use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use serde_json::{Value, json};
use tempfile::TempDir;
use super::super::{Answerer, Listener, Peer};
use crate::registry::presence::Presence;
use crate::test_support::seat::Seat;
use crate::test_support::wire::{material, mint};
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, _peer: &Peer, _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::test_support::seat::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");
let mut first = 0usize;
seat.followed(&json!({"op": "ops"}), &mut |_| {
first += 1;
false
})
.expect("the reader's own end is a clean end");
assert_eq!(first, 1, "the read ended on the frame that said so");
}