mod engine;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use serde_json::{Value, json};
use tempfile::TempDir;
use super::*;
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::server::{Answerer, Listener};
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())
}
}
struct Paced {
gate: Arc<AtomicUsize>,
frames: Vec<Value>,
}
impl Answerer for Paced {
fn answer(
&self,
_client: &crate::registry::Client,
_request: Value,
) -> Box<dyn Iterator<Item = Value>> {
let gate = Arc::clone(&self.gate);
let frames = self.frames.clone();
Box::new((0..=frames.len()).filter_map(move |i| {
while gate.load(Ordering::Relaxed) <= i {
std::thread::yield_now();
}
frames.get(i).cloned()
}))
}
}
pub(super) fn awaited(tail: &mut Tail, text: &str) -> bool {
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(20);
loop {
if declared(tail).and_then(|s| s.text).as_deref() == Some(text) {
return true;
}
if std::time::Instant::now() > deadline {
return false;
}
std::thread::yield_now();
}
}
pub(super) fn frame(text: &str) -> Value {
json!({"ok": true, "kind": "follow", "stream": {"text": text, "delta": "text"}})
}
pub(super) fn subject() -> Value {
crate::boundary::codec::encode(&crate::boundary::Gesture::Ask(
crate::boundary::Query::Follow {
workspace: "alba".to_owned(),
agent: "c-1".to_owned(),
},
))
}
fn wired(tmp: &TempDir, says: Vec<Value>) -> (Listener, Lane, Tail) {
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 (tail, end) = pair();
(listener, Lane::new(seat, end, Arc::new(NoRepaint)), tail)
}
pub(super) fn declared(tail: &mut Tail) -> Option<crate::git_tree::Stream> {
tail.settle();
tail.ask(&subject())
}
pub(super) fn watching(tail: &mut Tail) -> Option<crate::git_tree::Stream> {
declared(tail);
declared(tail)
}
#[test]
fn every_frame_the_engine_writes_reaches_the_seat_and_the_end_is_told() {
let tmp = TempDir::new().expect("tmp");
mint(tmp.path());
let gate = Arc::new(AtomicUsize::new(0));
let listener = Listener::bind(
&material(tmp.path(), Role::Server, EPHEMERAL),
Arc::new(Paced {
gate: Arc::clone(&gate),
frames: vec![frame("the first "), frame("the first half.")],
}),
Presence::default(),
)
.expect("bind");
let seat = crate::wire::client::Seat::open(&material(
tmp.path(),
Role::Window,
&crate::wire::loopback(&listener.address()),
))
.expect("seat");
let (mut tail, end) = pair();
let mut lane = Lane::new(seat, end, Arc::new(NoRepaint));
assert_eq!(watching(&mut tail), None, "nothing has crossed yet");
let held = std::thread::spawn(move || lane.turn());
gate.store(1, Ordering::Relaxed);
assert!(
awaited(&mut tail, "the first "),
"the first frame lands while the read is still open"
);
gate.store(2, Ordering::Relaxed);
assert!(
awaited(&mut tail, "the first half."),
"and so does the next — a frame carries the whole fold, so the newest wins"
);
gate.store(3, Ordering::Relaxed);
assert!(held.join().expect("the turn ends"), "frames landed");
assert_eq!(
declared(&mut tail),
None,
"the stream is over and the seat is back on the pull fold"
);
}
#[test]
fn a_lane_with_no_engine_lands_nothing_and_says_so() {
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 tail, end) = pair();
let mut lane = Lane::new(seat, end, Arc::new(NoRepaint));
watching(&mut tail);
assert!(!lane.turn(), "no engine, no frames");
assert_eq!(declared(&mut tail), None);
}
#[test]
fn a_lane_with_no_subject_asks_nothing() {
let tmp = TempDir::new().expect("tmp");
let (_listener, mut lane, _tail) = wired(&tmp, vec![frame("unread")]);
assert!(!lane.turn(), "nothing declared, nothing asked");
}
#[test]
fn an_answer_of_another_kind_ends_the_read() {
let tmp = TempDir::new().expect("tmp");
let (_listener, mut lane, mut tail) = wired(
&tmp,
vec![json!({"ok": false, "error": "no such conversation"})],
);
watching(&mut tail);
assert!(!lane.turn(), "nothing landed");
assert_eq!(declared(&mut tail), None);
}
#[test]
fn a_frame_for_the_conversation_just_left_lands_nowhere() {
let (mut tail, end) = pair();
let other = crate::boundary::codec::encode(&crate::boundary::Gesture::Ask(
crate::boundary::Query::Follow {
workspace: "alba".to_owned(),
agent: "c-2".to_owned(),
},
));
watching(&mut tail);
end.publish(
&subject().to_string(),
Some(crate::git_tree::Stream {
text: Some("c-1 is talking".to_owned()),
..crate::git_tree::Stream::default()
}),
);
assert!(
declared(&mut tail).is_some(),
"the fold is for this subject"
);
tail.settle();
tail.ask(&other);
end.publish(
&subject().to_string(),
Some(crate::git_tree::Stream {
text: Some("c-1 is still talking".to_owned()),
..crate::git_tree::Stream::default()
}),
);
tail.settle();
assert_eq!(tail.ask(&other), None, "and c-2 has said nothing yet");
}
#[test]
fn a_lane_whose_frame_end_is_gone_follows_nothing() {
let (mut tail, mut end) = pair();
watching(&mut tail);
assert!(end.standing().is_some(), "a live frame end has a subject");
drop(tail);
assert!(end.standing().is_none(), "and a dead one has none, forever");
}
#[test]
fn the_lane_thread_starts_and_stops() {
let tmp = TempDir::new().expect("tmp");
let (_listener, lane, _tail) = wired(&tmp, vec![frame("nobody asked")]);
drop(lane.start());
}