use std::path::{Path, PathBuf};
use leviath_agent_client::MAX_FRAME_BYTES;
use leviath_core::floor_char_boundary;
#[derive(Debug, Default)]
pub struct StageTail {
offsets: Vec<usize>,
}
impl StageTail {
pub fn new() -> Self {
Self::default()
}
pub fn pump(&mut self, runs_dir: &Path, run_id: &str) -> String {
let mut out = String::new();
let mut idx = 0;
loop {
let path = stage_output_path(runs_dir, run_id, idx);
let exists = path.exists();
if !exists && idx >= self.offsets.len() {
break;
}
let bytes = std::fs::read(&path).unwrap_or_default();
let seen = self.offsets.get(idx).copied().unwrap_or(0);
if bytes.len() > seen {
out.push_str(&String::from_utf8_lossy(&bytes[seen..]));
}
if idx < self.offsets.len() {
self.offsets[idx] = bytes.len();
} else {
self.offsets.push(bytes.len());
}
idx += 1;
}
out
}
}
fn stage_output_path(runs_dir: &Path, run_id: &str, idx: usize) -> PathBuf {
runs_dir
.join(run_id)
.join("stages")
.join(idx.to_string())
.join("output.log")
}
pub fn split_chunks(text: &str) -> Vec<&str> {
let mut chunks = Vec::new();
let mut rest = text;
while !rest.is_empty() {
let (chunk, tail) = rest.split_at(floor_char_boundary(rest, MAX_FRAME_BYTES));
chunks.push(chunk);
rest = tail;
}
chunks
}
#[cfg(test)]
mod tests {
use super::*;
fn write_stage(root: &Path, run_id: &str, idx: usize, text: &str) {
let path = stage_output_path(root, run_id, idx);
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(path, text).unwrap();
}
fn append_stage(root: &Path, run_id: &str, idx: usize, text: &str) {
use std::io::Write;
let path = stage_output_path(root, run_id, idx);
let mut f = std::fs::OpenOptions::new().append(true).open(path).unwrap();
f.write_all(text.as_bytes()).unwrap();
}
#[test]
fn pump_of_a_run_with_no_output_is_empty() {
let dir = tempfile::tempdir().unwrap();
let mut tail = StageTail::new();
assert_eq!(tail.pump(dir.path(), "run-1"), "");
}
#[test]
fn pump_returns_only_bytes_appended_since_last_call() {
let dir = tempfile::tempdir().unwrap();
let mut tail = StageTail::new();
write_stage(dir.path(), "run-1", 0, "hello\n");
assert_eq!(tail.pump(dir.path(), "run-1"), "hello\n");
assert_eq!(tail.pump(dir.path(), "run-1"), "");
append_stage(dir.path(), "run-1", 0, "world\n");
assert_eq!(tail.pump(dir.path(), "run-1"), "world\n");
}
#[test]
fn pump_streams_stages_in_order() {
let dir = tempfile::tempdir().unwrap();
let mut tail = StageTail::new();
write_stage(dir.path(), "run-1", 0, "stage zero\n");
assert_eq!(tail.pump(dir.path(), "run-1"), "stage zero\n");
write_stage(dir.path(), "run-1", 1, "stage one\n");
assert_eq!(tail.pump(dir.path(), "run-1"), "stage one\n");
}
#[test]
fn pump_drains_a_completed_stage_and_the_current_one_together() {
let dir = tempfile::tempdir().unwrap();
let mut tail = StageTail::new();
write_stage(dir.path(), "run-1", 0, "zero\n");
write_stage(dir.path(), "run-1", 1, "one\n");
assert_eq!(tail.pump(dir.path(), "run-1"), "zero\none\n");
}
#[test]
fn pump_after_a_stage_slot_is_tracked_rereads_a_late_appended_stage() {
let dir = tempfile::tempdir().unwrap();
let mut tail = StageTail::new();
write_stage(dir.path(), "run-1", 0, "zero\n");
assert_eq!(tail.pump(dir.path(), "run-1"), "zero\n");
append_stage(dir.path(), "run-1", 0, "zero-more\n");
write_stage(dir.path(), "run-1", 1, "one\n");
assert_eq!(tail.pump(dir.path(), "run-1"), "zero-more\none\n");
}
#[test]
fn pump_tolerates_invalid_utf8_without_panicking() {
let dir = tempfile::tempdir().unwrap();
let path = stage_output_path(dir.path(), "run-1", 0);
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(&path, [b'o', b'k', 0xff, b'\n']).unwrap();
let mut tail = StageTail::new();
assert!(tail.pump(dir.path(), "run-1").starts_with("ok"));
}
#[test]
fn split_of_empty_is_no_chunks() {
assert!(split_chunks("").is_empty());
}
#[test]
fn split_of_small_text_is_one_chunk() {
assert_eq!(split_chunks("hello world"), vec!["hello world"]);
}
#[test]
fn split_breaks_large_text_into_frame_sized_pieces() {
let big = "a".repeat(MAX_FRAME_BYTES * 2 + 5);
let chunks = split_chunks(&big);
assert_eq!(chunks.len(), 3);
assert_eq!(chunks[0].len(), MAX_FRAME_BYTES);
assert_eq!(chunks[1].len(), MAX_FRAME_BYTES);
assert_eq!(chunks[2].len(), 5);
assert_eq!(chunks.concat(), big);
}
#[test]
fn split_never_tears_a_multibyte_char() {
let s = "✓".repeat(MAX_FRAME_BYTES);
let chunks = split_chunks(&s);
assert!(chunks.iter().all(|c| c.len() <= MAX_FRAME_BYTES));
assert!(chunks.len() >= 2);
assert_eq!(chunks.concat(), s);
}
}