use super::*;
use crate::registry::presence::Presence;
use crate::wire::frame;
use rustls::ClientConnection;
use rustls::pki_types::ServerName;
use std::net::{Ipv4Addr, TcpStream};
use std::path::Path;
use std::time::Instant;
struct Watcher(Presence);
impl Answerer for Watcher {
fn answer(&self, _peer: &Peer, _request: Value) -> Box<dyn Iterator<Item = Value>> {
Box::new(std::iter::once(
json!({"live": self.0.live().into_iter().collect::<Vec<String>>()}),
))
}
}
#[test]
fn a_live_connection_is_present_and_a_closed_one_is_not() {
let tmp = TempDir::new().expect("tmp");
mint(tmp.path());
let presence = Presence::default();
let listener = Listener::bind(
&material(
tmp.path(),
Role::Server,
crate::test_support::wire::EPHEMERAL,
),
Arc::new(Watcher(presence.clone())),
presence.clone(),
)
.expect("bind");
let seat = Seat::open(&material(tmp.path(), Role::Client, &listener.address())).expect("seat");
let stream = seat.ask(&json!({"op": "workspaces"})).expect("answered");
assert_eq!(stream[0]["live"], json!(["yog-client"]));
let deadline = Instant::now() + Duration::from_secs(5);
while !presence.live().is_empty() && Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(10));
}
assert!(presence.live().is_empty(), "the connection's guard dropped");
}
#[test]
fn a_peer_that_goes_silent_is_reaped() {
let tmp = TempDir::new().expect("tmp");
mint(tmp.path());
let presence = Presence::default();
let config = crate::wire::tls::server_config(&material(
tmp.path(),
Role::Server,
crate::test_support::wire::EPHEMERAL,
))
.expect("config");
let tcp = TcpListener::bind("127.0.0.1:0").expect("bind");
let address = tcp.local_addr().expect("addr").to_string();
let watcher = Watcher(presence.clone());
let held = presence.clone();
let served = std::thread::spawn(move || {
let (stream, _) = tcp.accept().expect("accept");
serve(stream, &config, &watcher, &held, Duration::from_millis(50));
});
let mut tls = client(tmp.path(), &address);
frame::write_value(&mut tls, &json!({"protocol": crate::wire::hello::PROTOCOL}))
.expect("preface");
frame::write_value(&mut tls, &json!({"op": "workspaces"})).expect("ask");
frame::read_value(&mut tls).expect("read").expect("preface");
let answer = frame::read_value(&mut tls).expect("read").expect("answer");
assert_eq!(answer["live"], json!(["yog-client"]));
let deadline = Instant::now() + Duration::from_secs(5);
while !presence.live().is_empty() && Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(10));
}
assert!(presence.live().is_empty(), "the silent peer was reaped");
drop(tls);
served.join().expect("the serve thread returned");
}
fn client(dir: &Path, address: &str) -> StreamOwned<ClientConnection, TcpStream> {
let config =
crate::wire::tls::client_config(&material(dir, Role::Client, address)).expect("config");
let conn = ClientConnection::new(config, ServerName::IpAddress(Ipv4Addr::LOCALHOST.into()))
.expect("tls");
let tcp = TcpStream::connect(address).expect("connect");
tcp.set_read_timeout(Some(Duration::from_secs(10)))
.expect("timeout");
StreamOwned::new(conn, tcp)
}