use super::*;
use crate::registry::presence::Presence;
use crate::test_support::wire::{material, mint};
use crate::wire::client::Seat;
use crate::wire::material::Role;
use serde_json::json;
use std::sync::atomic::AtomicUsize;
use tempfile::TempDir;
mod lazy;
mod presence;
struct Echo {
asked: Arc<AtomicUsize>,
chunks: usize,
}
impl Answerer for Echo {
fn answer(&self, client: &Client, request: Value) -> Box<dyn Iterator<Item = Value>> {
self.asked.fetch_add(1, Ordering::Relaxed);
let name = client.name();
Box::new((0..self.chunks).map(move |n| {
json!({"ok": true, "kind": "echo", "seq": n,
"asked": request, "client": name})
}))
}
}
fn wired(chunks: usize) -> (TempDir, Listener, Seat, Arc<AtomicUsize>) {
let tmp = TempDir::new().expect("tmp");
mint(tmp.path());
let asked = Arc::new(AtomicUsize::new(0));
let listener = Listener::bind(
&material(
tmp.path(),
Role::Server,
crate::test_support::wire::EPHEMERAL,
),
Arc::new(Echo {
asked: Arc::clone(&asked),
chunks,
}),
Presence::default(),
)
.expect("bind");
let seat = Seat::open(&material(tmp.path(), Role::Client, &listener.address())).expect("seat");
(tmp, listener, seat, asked)
}
#[test]
fn a_certificate_bearing_seat_is_answered() {
let (_tmp, listener, seat, asked) = wired(1);
assert!(listener.address().starts_with("127.0.0.1:"));
let stream = seat.ask(&json!({"op": "workspaces"})).expect("answered");
assert_eq!(stream.len(), 1);
assert_eq!(stream[0]["asked"], json!({"op": "workspaces"}));
assert_eq!(asked.load(Ordering::Relaxed), 1);
seat.ask(&json!({"op": "board"})).expect("answered");
assert_eq!(asked.load(Ordering::Relaxed), 2);
}
#[test]
fn a_many_frame_answer_is_the_same_shape() {
let (_tmp, _listener, seat, _asked) = wired(3);
let stream = seat.ask(&json!({"op": "ops"})).expect("answered");
assert_eq!(stream.len(), 3);
assert_eq!(stream[2]["seq"], json!(2));
}
#[test]
fn an_uncertificated_peer_is_refused_by_tls() {
let tmp = TempDir::new().expect("tmp");
mint(tmp.path());
let asked = Arc::new(AtomicUsize::new(0));
let listener = Listener::bind(
&material(
tmp.path(),
Role::Server,
crate::test_support::wire::EPHEMERAL,
),
Arc::new(Echo {
asked: Arc::clone(&asked),
chunks: 1,
}),
Presence::default(),
)
.expect("bind");
let mut plain = std::net::TcpStream::connect(listener.address()).expect("connect");
let _ = crate::wire::frame::write_value(&mut plain, &json!({"op": "workspaces"}));
plain
.set_read_timeout(Some(Duration::from_secs(5)))
.expect("timeout");
let mut buf = [0u8; 1];
let _ = std::io::Read::read(&mut plain, &mut buf);
assert_eq!(asked.load(Ordering::Relaxed), 0);
}
#[test]
fn a_foreign_certificate_is_refused() {
let tmp = TempDir::new().expect("tmp");
mint(tmp.path());
let stranger = TempDir::new().expect("tmp");
mint(stranger.path());
let asked = Arc::new(AtomicUsize::new(0));
let listener = Listener::bind(
&material(
tmp.path(),
Role::Server,
crate::test_support::wire::EPHEMERAL,
),
Arc::new(Echo {
asked: Arc::clone(&asked),
chunks: 1,
}),
Presence::default(),
)
.expect("bind");
std::fs::copy(tmp.path().join("ca.pem"), stranger.path().join("ca.pem")).expect("copy");
let seat = Seat::open(&material(
stranger.path(),
Role::Client,
&listener.address(),
))
.expect("seat");
seat.ask(&json!({"op": "workspaces"})).expect_err("refused");
assert_eq!(asked.load(Ordering::Relaxed), 0);
}
#[test]
fn a_silent_connection_ends_at_eof() {
let (_tmp, listener, _seat, asked) = wired(1);
drop(std::net::TcpStream::connect(listener.address()).expect("connect"));
assert_eq!(asked.load(Ordering::Relaxed), 0);
}
#[test]
fn an_unusable_config_drops_the_connection() {
let tmp = TempDir::new().expect("tmp");
mint(tmp.path());
let mut config = Arc::into_inner(
crate::wire::tls::server_config(&material(tmp.path(), Role::Server, "127.0.0.1:0"))
.expect("config"),
)
.expect("sole owner");
config.max_fragment_size = Some(1);
let listener = TcpListener::bind("127.0.0.1:0").expect("bind");
let address = listener.local_addr().expect("addr").to_string();
let peer = std::thread::spawn(move || {
let (stream, _) = listener.accept().expect("accept");
serve(
stream,
&Arc::new(config),
&Echo {
asked: Arc::new(AtomicUsize::new(0)),
chunks: 1,
},
&Presence::default(),
);
});
drop(TcpStream::connect(address).expect("connect"));
peer.join().expect("served");
}
#[test]
fn dropping_the_listener_stops_it() {
let (_tmp, listener, seat, _asked) = wired(1);
seat.ask(&json!({"op": "workspaces"})).expect("answered");
drop(listener);
seat.ask(&json!({"op": "workspaces"})).expect_err("stopped");
}
#[test]
fn an_unbindable_address_refuses() {
let tmp = TempDir::new().expect("tmp");
mint(tmp.path());
let refusal = Listener::bind(
&material(tmp.path(), Role::Server, "256.256.256.256:1"),
Arc::new(Echo {
asked: Arc::new(AtomicUsize::new(0)),
chunks: 1,
}),
Presence::default(),
)
.err()
.expect("refused");
assert!(refusal.contains("256.256.256.256:1"), "{refusal}");
}
#[test]
fn unusable_material_refuses_before_binding() {
let tmp = TempDir::new().expect("tmp");
mint(tmp.path());
std::fs::write(tmp.path().join("ca.pem"), "").expect("write");
assert!(
Listener::bind(
&material(
tmp.path(),
Role::Server,
crate::test_support::wire::EPHEMERAL
),
Arc::new(Echo {
asked: Arc::new(AtomicUsize::new(0)),
chunks: 1,
}),
Presence::default(),
)
.is_err()
);
}
#[test]
fn the_answerer_is_handed_the_certificates_leaf_name() {
let (_tmp, _listener, seat, _asked) = wired(1);
let stream = seat.ask(&json!({"op": "workspaces"})).expect("answered");
assert_eq!(stream[0]["client"], "yog-client");
}
#[test]
fn a_chain_naming_no_usable_identity_is_nobody() {
assert!(peer_client(None).is_none());
assert!(peer_client(Some(&[])).is_none());
let nameless = CertificateDer::from(vec![0x30, 0x00]);
assert!(peer_client(Some(&[nameless])).is_none());
}