mod reading;
use std::path::Path;
use super::*;
use crate::app::Snapshot;
use crate::boundary::tests::{agent, snapshot};
use crate::git_tree::{AgentState, Delta};
use crate::state::{SnapshotCell, new_snapshot_cell};
use tempfile::TempDir;
pub(super) const AGENT: &str = "c-1";
pub(super) fn text_delta(fragment: &str) -> String {
format!(
"{{\"type\":\"content_delta\",\"index\":0,\"delta\":{{\"text_delta\":\"{fragment}\"}}}}\n"
)
}
pub(super) fn thinking_delta(fragment: &str) -> String {
format!(
"{{\"type\":\"content_delta\",\"index\":0,\"delta\":{{\"thinking_delta\":\"{fragment}\"}}}}\n"
)
}
pub(super) fn response(ws: &Path, seq: u32) -> std::path::PathBuf {
let step = ws.join("steps").join(AGENT).join(format!("{seq:03}"));
std::fs::create_dir_all(&step).expect("step dir");
step.join("response.json")
}
pub(super) fn append(path: &Path, bytes: &str) {
use std::io::Write;
let mut file = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(path)
.expect("open response");
file.write_all(bytes.as_bytes()).expect("append");
}
pub(super) fn seated(ws: &Path, state: AgentState) -> SnapshotCell {
new_snapshot_cell(std::sync::Arc::new(snapshot(
ws,
"alba",
vec![agent(AGENT, state, 100)],
vec![],
)))
}
pub(super) fn flying() -> (TempDir, SnapshotCell, Follow) {
let dir = tempfile::tempdir().expect("tmp");
let cell = seated(dir.path(), AgentState::InFlight);
let follow = Follow::new(
std::sync::Arc::clone(&cell),
dir.path().to_path_buf(),
AGENT.to_owned(),
);
(dir, cell, follow)
}
pub(super) fn said(frame: Frame) -> Option<Stream> {
match frame {
Frame::Ready(stream) => Some(stream),
_ => None,
}
}
#[test]
pub(super) fn appended_bytes_are_on_a_frame_with_no_derivation_between() {
let (dir, cell, mut follow) = flying();
let file = response(dir.path(), 1);
let pinned = crate::state::latest_snapshot(&cell);
append(&file, &text_delta("the first "));
assert_eq!(
said(follow.poll()).and_then(|s| s.text).as_deref(),
Some("the first "),
"the first bytes are a frame"
);
append(&file, &text_delta("half."));
assert_eq!(
said(follow.poll()).and_then(|s| s.text).as_deref(),
Some("half."),
"every character that landed since is on the next frame, and nothing else"
);
assert!(
std::sync::Arc::ptr_eq(&pinned, &crate::state::latest_snapshot(&cell)),
"and no derivation ran to put it there — this is the whole claim"
);
}
#[test]
fn a_frame_says_exactly_what_the_pull_fold_says_of_the_same_bytes() {
let (dir, _cell, mut follow) = flying();
let file = response(dir.path(), 1);
append(&file, &thinking_delta("first I "));
append(&file, &text_delta("here goes"));
let framed = said(follow.poll()).expect("a frame");
let mut row = agent(AGENT, AgentState::InFlight, 100);
row.stream = crate::git_tree::stream_from_disk(dir.path(), AGENT);
let derived: Snapshot = snapshot(dir.path(), "alba", vec![row], vec![]);
let pulled = crate::boundary::answer::inspector::live_tail(&derived, dir.path(), AGENT)
.expect("in flight, so there is a tail");
assert_eq!(framed, pulled, "one fold, two ways of reaching it");
assert_eq!(framed.thinking.as_deref(), Some("first I "));
assert_eq!(framed.last_delta, Some(Delta::Text));
}
#[test]
fn reasoning_streams_and_carries_the_doing_split_with_it() {
let (dir, _cell, mut follow) = flying();
let file = response(dir.path(), 1);
append(&file, &thinking_delta("first I "));
let stream = said(follow.poll()).expect("a frame");
assert_eq!(stream.thinking.as_deref(), Some("first I "));
assert_eq!(stream.text, None, "no answer yet, and none invented");
assert_eq!(stream.last_delta, Some(Delta::Thinking));
append(&file, &text_delta("here goes"));
let stream = said(follow.poll()).expect("a frame");
assert_eq!(stream.text.as_deref(), Some("here goes"));
assert_eq!(stream.last_delta, Some(Delta::Text));
}
#[test]
fn the_frames_of_a_read_sum_to_the_answer_and_carry_it_once() {
let (dir, _cell, mut follow) = flying();
let file = response(dir.path(), 1);
let words = ["A sentence ", "that arrives ", "in pieces."];
let mut framed = Vec::new();
for word in words {
append(&file, &text_delta(word));
framed.push(
said(follow.poll())
.and_then(|s| s.text)
.unwrap_or_else(|| panic!("a frame for {word:?}")),
);
}
assert_eq!(framed, words, "each frame is exactly what landed for it");
assert_eq!(
framed.concat().len(),
words.concat().len(),
"and the answer crosses once, not once per frame"
);
}
#[test]
fn absorbing_every_frame_of_a_read_is_the_fold_of_the_same_bytes() {
let (dir, _cell, mut follow) = flying();
let file = response(dir.path(), 1);
let mut held = Stream::default();
for (thinking, text) in [("weighing ", "the answer "), ("it up. ", "so far.")] {
append(&file, &thinking_delta(thinking));
append(&file, &text_delta(text));
held.absorb(said(follow.poll()).expect("a frame"));
}
assert_eq!(
held,
crate::git_tree::stream_from_disk(dir.path(), AGENT),
"the seat's accumulation is the file's own fold"
);
assert_eq!(held.text.as_deref(), Some("the answer so far."));
assert_eq!(held.thinking.as_deref(), Some("weighing it up. "));
}
#[test]
fn a_read_that_starts_mid_answer_is_whole_on_its_first_frame() {
let (dir, cell, mut early) = flying();
let file = response(dir.path(), 1);
append(&file, &text_delta("already said. "));
assert!(matches!(early.poll(), Frame::Ready(_)));
drop(early);
append(&file, &text_delta("and more."));
let mut late = Follow::new(cell, dir.path().to_path_buf(), AGENT.to_owned());
assert_eq!(
said(late.poll()).and_then(|s| s.text).as_deref(),
Some("already said. and more."),
"a fresh read is answered from zero, holding nothing"
);
}