Mercurial > templog
view rust/src/params.rs @ 628:e1b5938de122 rust
futures sink? unsure what this was
author | Matt Johnston <matt@ucc.asn.au> |
---|---|
date | Fri, 19 Apr 2019 13:57:40 +0800 |
parents | d5075136442f |
children | 3e5e52d50af5 |
line wrap: on
line source
extern crate tokio_core; extern crate futures_await as futures; extern crate rand; extern crate serde_json; extern crate base64; extern crate atomicwrites; use std::time::Duration; use std::io; use std::str; use std::rc::Rc; use std::sync::{Arc,Mutex}; use std::error::Error; use std::cell::{Cell,RefCell}; use std::fs::File; use std::io::Read; use futures::prelude::*; use tokio_core::reactor::Interval; use tokio_core::reactor::Handle; use futures::{Stream,Future,future}; use self::rand::Rng; use std::str::FromStr; use hyper; use hyper::header::{Headers, ETag, EntityTag}; use hyper::client::Client; use types::*; use ::Config; #[derive(Deserialize, Serialize, Debug, Clone)] pub struct Params { pub fridge_setpoint: f32, pub fridge_difference: f32, pub overshoot_delay: u64, pub overshoot_factor: f32, pub disabled: bool, pub nowort: bool, pub fridge_range_lower: f32, pub fridge_range_upper: f32, } impl Params { pub fn defaults() -> Params { Params { fridge_setpoint: 16.0, fridge_difference: 0.2, overshoot_delay: 720, // 12 minutes overshoot_factor: 1.0, disabled: false, nowort: false, fridge_range_lower: 3.0, fridge_range_upper: 3.0, } } fn try_load(filename: &str) -> Result<Params, TemplogError> { let mut s = String::new(); File::open(filename)?.read_to_string(&mut s)?; Ok(serde_json::from_str(&s)?) } pub fn load(config: &Config) -> Params { Self::try_load(&config.PARAMS_FILE) .unwrap_or_else(|_| Params::defaults()) } } #[derive(Deserialize, Debug)] struct Response { // sent as an opaque etag: header. Has format "epoch-nonce", // responses where the epoch do not match ParamWaiter::epoch are dropped etag: String, params: Params, } #[derive(Clone)] pub struct ParamWaiter { limitlog: NotTooOften, // last_etag is used for long-polling. last_etag: RefCell<String>, epoch: String, config: Config, handle: Handle, } const LOG_MINUTES: u64 = 15; const MAX_RESPONSE_SIZE: usize = 10000; const TIMEOUT_MINUTES: u64 = 5; impl ParamWaiter { fn make_request(&self) -> hyper::client::FutureResponse { let uri = self.config.SETTINGS_URL.parse().expect("Bad SETTINGS_URL in config"); let mut req = hyper::client::Request::new(hyper::Method::Get, uri); { let headers = req.headers_mut(); headers.set(hyper::header::ETag(EntityTag::new(false, self.last_etag.borrow().clone()))); } hyper::client::Client::new(&self.handle).request(req) // XXX how do we do this? // req.timeout(Duration::new(TIMEOUT_MINUTES * 60, 0)).unwrap(); } fn handle_response(&self, buf : hyper::Chunk, status: hyper::StatusCode) -> Result<Params, TemplogError> { let text = String::from_utf8_lossy(buf.as_ref()); match status { hyper::StatusCode::Ok => { // new params let r: Response = serde_json::from_str(&text)?; let mut le = self.last_etag.borrow_mut(); *le = r.etag; // update params if the epoch is correct if let Some(e) = le.split('-').next() { if e == &self.epoch { self.write_params(&r.params); return Ok(r.params); } } Err(TemplogError::new(&format!("Bad epoch from server '{}' expected '{}'", *le, self.epoch))) } hyper::StatusCode::NotModified => { // XXX this isn't really an error Err(TemplogError::new("304 unmodified (long polling timeout at the server)")) }, _ => { Err(TemplogError::new(&format!("Wrong server response code {}: {}", status.as_u16(), text))) }, } } fn write_params(&self, params: &Params) { let p = atomicwrites::AtomicFile::new(&self.config.PARAMS_FILE, atomicwrites::AllowOverwrite); p.write(|f| { serde_json::to_writer(f, params) }); } #[async_stream(item = Params)] pub fn stream(config: Config, handle: Handle) -> Result<(), TemplogError> { let rcself = Rc::new(ParamWaiter::new(config, handle)); let dur = Duration::from_millis(4000); #[async] for _ in Interval::new(dur, &rcself.handle).unwrap() { // fetch params // TODO - skip if inflight. let r = await!(rcself.make_request()).map_err(|e| TemplogError::new_hyper("response", e))?; let status = r.status(); let b = await!(r.body().concat2()).map_err(|e| TemplogError::new_hyper("body", e))?; if let Ok(params) = rcself.handle_response(b, status) { stream_yield!(params); } } Ok(()) } fn new(config: Config, handle: Handle) -> Self { let mut b = [0u8; 15]; // 15 bytes -> 20 characters base64 rand::OsRng::new().expect("Opening system RNG failed").fill_bytes(&mut b); let epoch = base64::encode(&b); ParamWaiter { limitlog: NotTooOften::new(LOG_MINUTES*60), last_etag: RefCell::new(String::new()), epoch: epoch, config: config, handle: handle, } } }