use std::io::{self, ErrorKind};
use std::net::TcpStream;
use std::time::{Duration, Instant};
use rustls::{ClientConnection, StreamOwned};
use serde_json::{Value, json};
use super::frame;
pub(crate) const GONE: Duration = Duration::from_mins(2);
pub(crate) const PING: Duration = Duration::from_secs(25);
pub(crate) type Tls = StreamOwned<ClientConnection, TcpStream>;
pub(crate) struct Line {
pub(crate) tls: Tls,
pub(crate) spelled: u32,
pub(crate) fresh: bool,
pub(crate) kept: bool,
}
impl Line {
pub(crate) fn from_held(held: Held) -> Line {
Line {
tls: held.tls,
spelled: held.spelled,
fresh: false,
kept: true,
}
}
pub(crate) fn give_back(self, key: &str, now: Instant) {
if !self.kept {
return;
}
let held = Held {
tls: self.tls,
spelled: self.spelled,
heard: now,
};
crate::state::worked(key, |w| w.held.push(held));
}
}
pub(crate) struct Held {
tls: Tls,
spelled: u32,
heard: Instant,
}
impl Held {
pub(crate) fn gone(&mut self, now: Instant) -> Option<&'static str> {
if now.saturating_duration_since(self.heard) + PING >= GONE {
Some("past the silence bound")
} else if open(&mut self.tls) {
None
} else {
Some("closed at the far end")
}
}
}
fn open(tls: &mut Tls) -> bool {
let _ = tls.sock.set_nonblocking(true);
let verdict = loop {
match tls.conn.process_new_packets() {
Ok(state) if state.peer_has_closed() => break false,
Ok(_) => {}
Err(_) => break false,
}
if !tls.conn.wants_read() {
break true;
}
match tls.conn.read_tls(&mut tls.sock) {
Ok(0) => break false,
Ok(_) => {}
Err(e) if e.kind() == ErrorKind::WouldBlock => break true,
Err(_) => break false,
}
};
let _ = tls.sock.set_nonblocking(false);
verdict
}
pub(crate) fn read(r: &mut dyn io::Read, pings: &mut usize) -> io::Result<Option<Value>> {
loop {
match frame::read_value(r)? {
Some(frame) if frame == ping() => *pings += 1,
other => return Ok(other),
}
}
}
pub(crate) fn ping() -> Value {
json!({"ping": true})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_bounds_are_the_engines() {
assert_eq!(GONE, Duration::from_mins(2));
assert_eq!(PING, Duration::from_secs(25));
assert_eq!(ping(), json!({"ping": true}));
}
#[test]
fn a_read_skips_every_ping_and_hands_back_the_frame_or_the_end() {
let mut bytes = Vec::new();
frame::write_value(&mut bytes, &ping()).unwrap();
frame::write_value(&mut bytes, &ping()).unwrap();
frame::write_value(&mut bytes, &json!({"n": 1})).unwrap();
frame::write_value(&mut bytes, &ping()).unwrap();
frame::write_end(&mut bytes).unwrap();
let mut cursor = std::io::Cursor::new(bytes);
let mut pings = 0;
assert_eq!(
read(&mut cursor, &mut pings).unwrap(),
Some(json!({"n": 1}))
);
assert_eq!(pings, 2);
assert_eq!(read(&mut cursor, &mut pings).unwrap(), None);
assert_eq!(pings, 3);
}
}