use super::*;
use crate::boundary::reply::Reply;
use crate::registry::presence::Presence;
use crate::test_support::wire::{EPHEMERAL, NO_LISTENER, material, mint};
use crate::watch::NoRepaint;
use crate::wire::material::Role;
use crate::wire::post::{Post, pair};
use crate::wire::server::{Answerer, Listener};
use serde_json::{Value, json};
use tempfile::TempDir;
struct Says(Vec<Value>);
impl Answerer for Says {
fn answer(
&self,
_client: &crate::registry::Client,
_request: Value,
) -> Box<dyn Iterator<Item = Value>> {
Box::new(self.0.clone().into_iter())
}
}
fn wired(tmp: &TempDir, says: Vec<Value>) -> (Listener, Poster, Post) {
mint(tmp.path());
let listener = Listener::bind(
&material(tmp.path(), Role::Server, EPHEMERAL),
Arc::new(Says(says)),
Presence::default(),
)
.expect("bind");
let seat = crate::wire::client::Seat::open(&material(
tmp.path(),
Role::Window,
&crate::wire::loopback(&listener.address()),
))
.expect("seat");
let (post, outbox) = pair();
(
listener,
Poster::new(seat, outbox, Arc::new(NoRepaint)),
post,
)
}
#[test]
fn a_posted_act_crosses_the_wire_and_its_receipt_lands_under_its_ticket() {
let tmp = TempDir::new().expect("tmp");
let (_listener, mut poster, mut post) = wired(&tmp, vec![json!({"ok": true, "kind": "acked"})]);
let ticket = post.send(&json!({"op": "seen", "workspace": "home", "agent": "c-1"}));
assert!(post.settle().is_empty(), "nothing has crossed yet");
assert!(poster.pass(), "one act, sent and answered");
assert_eq!(post.settle(), vec![ticket]);
assert_eq!(post.receipt(ticket), Some(Ok(Reply::Acked)));
}
#[test]
fn a_dead_engine_lands_a_sentence_rather_than_blocking() {
let tmp = TempDir::new().expect("tmp");
mint(tmp.path());
let seat = crate::wire::client::Seat::open(&material(tmp.path(), Role::Window, NO_LISTENER))
.expect("seat");
let (mut post, outbox) = pair();
let mut poster = Poster::new(seat, outbox, Arc::new(NoRepaint));
let ticket = post.send(&json!({"op": "scan", "workspace": "home"}));
assert!(poster.pass());
post.settle();
let Some(Err(said)) = post.receipt(ticket) else {
panic!("a refusal");
};
assert!(said.starts_with("connect "), "{said}");
}
#[test]
fn the_thread_runs_until_the_window_drops() {
let tmp = TempDir::new().expect("tmp");
let (_listener, poster, mut post) = wired(&tmp, vec![json!({"ok": true, "kind": "acked"})]);
let ticket = post.send(&json!({"op": "seen", "workspace": "home", "agent": "c-1"}));
let thread = poster.start();
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10);
loop {
post.settle();
if post.receipt(ticket).is_some() {
break;
}
assert!(
std::time::Instant::now() < deadline,
"the act never crossed"
);
std::thread::sleep(std::time::Duration::from_millis(10));
}
drop(post);
thread.join().expect("the poster ends when nobody can post");
}