use std::time::Duration;
use super::super::*;
use super::{AGENT, append, flying, response, said, seated, text_delta};
use crate::boundary::tests::{agent, snapshot};
use crate::git_tree::AgentState;
#[test]
fn a_step_that_has_opened_no_response_file_is_waited_on() {
let (dir, _cell, mut follow) = flying();
let file = response(dir.path(), 1);
assert!(!file.exists(), "the step dir is there and the file is not");
assert!(matches!(follow.poll(), Frame::Waiting));
append(&file, &text_delta("and then it speaks"));
assert!(matches!(follow.poll(), Frame::Ready(..)));
}
#[test]
fn a_half_written_line_waits_for_its_newline() {
let (dir, _cell, mut follow) = flying();
let file = response(dir.path(), 1);
let whole = text_delta("atomic");
let (head, tail) = whole.split_at(whole.len() / 2);
append(&file, head);
assert!(said(follow.poll()).is_none(), "nothing whole arrived");
append(&file, tail);
assert_eq!(
said(follow.poll()).and_then(|s| s.text).as_deref(),
Some("atomic")
);
}
#[test]
fn an_event_the_operator_cannot_see_is_not_a_frame() {
let (dir, _cell, mut follow) = flying();
let file = response(dir.path(), 1);
append(&file, "{\"type\":\"message_start\"}\n");
assert!(matches!(follow.poll(), Frame::Waiting), "nothing to say");
append(&file, &text_delta("now something"));
assert!(matches!(follow.poll(), Frame::Ready(..)));
assert!(
matches!(follow.poll(), Frame::Waiting),
"and an unchanged file is not news either"
);
}
#[test]
fn a_file_that_shrank_is_read_again_from_the_start() {
let (dir, _cell, mut follow) = flying();
let file = response(dir.path(), 1);
append(&file, &text_delta("the long first answer"));
assert!(matches!(follow.poll(), Frame::Ready(..)));
std::fs::write(&file, text_delta("short")).expect("truncate");
assert_eq!(
said(follow.poll()).and_then(|s| s.text).as_deref(),
Some("short"),
"the replacement file is read whole and appended, so a seat's \
accumulation is now wrong by the bytes that vanished — the follower's \
own ruling, kept: the next derivation is what corrects it"
);
}
#[test]
fn a_settled_step_is_not_followed() {
let dir = tempfile::tempdir().expect("tmp");
let cell = seated(dir.path(), AgentState::Quiescent);
append(&response(dir.path(), 1), &text_delta("already committed"));
let mut follow = Follow::new(cell, dir.path().to_path_buf(), AGENT.to_owned());
assert!(
matches!(follow.poll(), Frame::Waiting),
"a hold, not an answer — and never a frame"
);
}
#[test]
fn the_stream_ends_when_the_call_does_and_the_last_bytes_come_out_first() {
let (dir, cell, mut follow) = flying();
let file = response(dir.path(), 1);
append(&file, &text_delta("the model begins"));
assert!(
matches!(follow.poll(), Frame::Ready(..)),
"the stream is open"
);
append(&file, &text_delta(" and ends"));
crate::state::publish_snapshot(
&cell,
std::sync::Arc::new(snapshot(
dir.path(),
"alba",
vec![agent(AGENT, AgentState::Quiescent, 100)],
vec![],
)),
);
assert_eq!(
said(follow.poll()).and_then(|s| s.text).as_deref(),
Some(" and ends"),
"nothing written is dropped by the close"
);
assert!(matches!(follow.poll(), Frame::Over), "then the terminator");
}
#[test]
fn a_step_advancing_ends_the_stream() {
let (dir, _cell, mut follow) = flying();
append(&response(dir.path(), 1), &text_delta("step one"));
assert!(matches!(follow.poll(), Frame::Ready(..)));
append(&response(dir.path(), 2), &text_delta("step two"));
assert!(matches!(follow.poll(), Frame::Over));
}
#[test]
fn a_tree_that_went_away_ends_the_stream() {
let (dir, _cell, mut follow) = flying();
append(&response(dir.path(), 1), &text_delta("here"));
assert!(matches!(follow.poll(), Frame::Ready(..)));
std::fs::remove_dir_all(dir.path().join("steps")).expect("delete the steps");
assert!(matches!(follow.poll(), Frame::Over));
}
#[test]
fn a_quiet_hold_expires_and_a_frame_resets_it() {
let (dir, cell, _drop) = flying();
let ws = dir.path().to_path_buf();
let mut quiet = Follow::holding(
std::sync::Arc::clone(&cell),
ws.clone(),
AGENT.to_owned(),
2,
Duration::ZERO,
);
assert!(
quiet.next().is_none(),
"two quiet looks and the hold is over"
);
let mut talking = Follow::holding(cell, ws.clone(), AGENT.to_owned(), 2, Duration::ZERO);
let file = response(&ws, 1);
append(&file, &text_delta("one"));
assert!(talking.next().is_some());
append(&file, &text_delta(" two"));
assert!(
talking.next().is_some(),
"the quiet look before this one did not count against a hold a frame reset"
);
}
#[test]
fn a_stream_that_is_over_ends_the_hold_at_once() {
let (dir, cell, _drop) = flying();
let ws = dir.path().to_path_buf();
let file = response(&ws, 1);
let mut held = Follow::holding(
cell,
ws.clone(),
AGENT.to_owned(),
u32::MAX,
Duration::ZERO,
);
append(&file, &text_delta("one"));
assert!(held.next().is_some(), "the open step's bytes");
let next_step = response(&ws, 2);
append(&next_step, &text_delta("two"));
assert!(
held.next().is_none(),
"the stream is over, and the hold ends with it"
);
}