use kanade_shared::kv::STDOUT_INLINE_THRESHOLD;
pub const MAX_CAPTURE_BYTES: usize = 32 * STDOUT_INLINE_THRESHOLD;
const _: () = assert!(MAX_CAPTURE_BYTES > STDOUT_INLINE_THRESHOLD);
const HEAD_FRACTION: usize = 3;
const fn split(cap: usize) -> (usize, usize) {
let head = cap / HEAD_FRACTION * (HEAD_FRACTION - 1);
(head, cap - head)
}
pub struct CappedOutput {
cap: usize,
head: Vec<u8>,
tail: std::collections::VecDeque<u8>,
total: u64,
}
impl CappedOutput {
pub fn new(cap: usize) -> Self {
let (head_cap, _) = split(cap);
Self {
cap,
head: Vec::with_capacity(head_cap.min(64 * 1024)),
tail: std::collections::VecDeque::new(),
total: 0,
}
}
pub fn total(&self) -> u64 {
self.total
}
pub fn truncated(&self) -> bool {
self.total > self.cap as u64
}
pub fn push(&mut self, bytes: &[u8]) {
self.total += bytes.len() as u64;
if self.cap == 0 {
return;
}
let (head_cap, tail_cap) = split(self.cap);
let mut rest = bytes;
if self.head.len() < head_cap {
let take = (head_cap - self.head.len()).min(rest.len());
self.head.extend_from_slice(&rest[..take]);
rest = &rest[take..];
}
if rest.is_empty() || tail_cap == 0 {
return;
}
let keep = &rest[rest.len().saturating_sub(tail_cap)..];
self.tail.extend(keep.iter().copied());
while self.tail.len() > tail_cap {
self.tail.pop_front();
}
}
pub fn finish(self) -> String {
let (a, b) = self.tail.as_slices();
let mut tail_bytes = Vec::with_capacity(a.len() + b.len());
tail_bytes.extend_from_slice(a);
tail_bytes.extend_from_slice(b);
if !self.truncated() {
let mut all = self.head;
all.extend_from_slice(&tail_bytes);
return String::from_utf8_lossy(&all).into_owned();
}
let head = String::from_utf8_lossy(&self.head).into_owned();
let tail = String::from_utf8_lossy(&tail_bytes).into_owned();
let dropped = self.total - (self.head.len() + tail_bytes.len()) as u64;
format!(
"{head}\n\
[kanade] output truncated: {total} bytes written, {dropped} dropped \
(kept the first {head_len} and last {tail_len}).\n\
[kanade] stdout is the control path for a run, not a bulk transfer \
channel — use a `collect:` job to ship large files.\n",
total = self.total,
head_len = self.head.len(),
tail_len = self.tail.len(),
) + &tail
}
}
#[cfg(test)]
mod tests {
use super::*;
fn feed(cap: usize, chunks: &[&[u8]]) -> String {
let mut c = CappedOutput::new(cap);
for ch in chunks {
c.push(ch);
}
c.finish()
}
#[test]
fn the_split_always_accounts_for_the_whole_budget() {
for cap in [0usize, 1, 2, 3, 4, 5, 10, 11, 99, 100, MAX_CAPTURE_BYTES] {
let (h, t) = split(cap);
assert_eq!(h + t, cap, "cap {cap} split into {h} + {t}");
}
}
#[test]
fn output_under_the_cap_is_returned_verbatim() {
let out = feed(100, &[b"hello ", b"world"]);
assert_eq!(out, "hello world");
assert!(!out.contains("truncated"));
}
#[test]
fn output_at_exactly_the_cap_is_not_truncated() {
let out = feed(10, &[b"0123456789"]);
assert_eq!(out, "0123456789");
}
#[test]
fn truncation_keeps_both_ends_and_says_what_it_dropped() {
let body: Vec<u8> = (0..1000u32).map(|i| b'a' + (i % 26) as u8).collect();
let out = feed(120, &[&body]);
assert!(out.starts_with("abcdefghij"), "out: {}", &out[..40]);
assert!(out.ends_with(std::str::from_utf8(&body[body.len() - 10..]).unwrap()));
assert!(out.contains("1000 bytes written"));
assert!(out.contains("880 dropped"));
}
#[test]
fn the_tail_is_the_last_bytes_across_many_chunks() {
let mut c = CappedOutput::new(30);
for i in 0..100u8 {
c.push(&[b'0' + (i % 10)]);
}
let out = c.finish();
assert!(
out.ends_with("6789"),
"out ends: {:?}",
&out[out.len() - 8..]
);
}
#[test]
fn one_huge_chunk_is_bounded_the_same_as_many_small_ones() {
let big = vec![b'x'; 10 * 1024 * 1024];
let mut c = CappedOutput::new(1024);
c.push(&big);
assert_eq!(c.total(), 10 * 1024 * 1024);
assert!(c.head.len() + c.tail.len() <= 1024);
}
#[test]
fn a_zero_cap_keeps_nothing_but_still_counts() {
let mut c = CappedOutput::new(0);
c.push(b"discarded");
assert_eq!(c.total(), 9);
assert!(c.truncated());
assert!(c.finish().contains("9 bytes written"));
}
#[test]
fn a_codepoint_split_by_the_boundary_does_not_lose_the_rest() {
let s = "あいうえお".repeat(50);
let out = feed(64, &[s.as_bytes()]);
assert!(out.contains("truncated"));
assert!(out.starts_with('あ'), "out: {out}");
}
#[test]
fn the_cap_leaves_room_for_the_object_store_path() {
let out = feed(
MAX_CAPTURE_BYTES,
&[&vec![b'x'; STDOUT_INLINE_THRESHOLD * 4]],
);
assert!(!out.contains("truncated"));
}
}