use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use serde_json::{Value, json};
use tempfile::TempDir;
use super::super::{Answerer, Client, Listener};
use crate::registry::presence::Presence;
use crate::test_support::wire::{material, mint};
use crate::wire::client::Seat;
use crate::wire::material::Role;
#[test]
fn a_frame_is_written_as_it_is_produced_rather_than_when_the_answer_ends() {
struct Paced(Arc<AtomicUsize>);
impl Answerer for Paced {
fn answer(&self, _client: &Client, _request: Value) -> Box<dyn Iterator<Item = Value>> {
let gate = Arc::clone(&self.0);
Box::new((0..2).map(move |n| {
while gate.load(Ordering::Relaxed) <= n {
std::thread::yield_now();
}
json!({"ok": true, "kind": "echo", "seq": n})
}))
}
}
let tmp = TempDir::new().expect("tmp");
mint(tmp.path());
let gate = Arc::new(AtomicUsize::new(1));
let listener = Listener::bind(
&material(
tmp.path(),
Role::Server,
crate::test_support::wire::EPHEMERAL,
),
Arc::new(Paced(Arc::clone(&gate))),
Presence::default(),
)
.expect("bind");
let seat = Seat::open(&material(
tmp.path(),
Role::Window,
&crate::wire::loopback(&listener.address()),
))
.expect("seat");
let mut seen = 0usize;
seat.followed(&json!({"op": "ops"}), &mut |landed| {
assert!(landed.is_err(), "the fixture's echo is nobody's Reply");
seen += 1;
gate.store(2, Ordering::Relaxed);
true
})
.expect("the stream ends cleanly");
assert_eq!(seen, 2, "and both frames arrive, in order");
}