use std::collections::{HashMap, VecDeque};
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Duration;
use tokio::sync::mpsc;
use zeph_sanitizer::pii::PiiFilter;
use zeph_sanitizer::secret_mask::SecretMaskRegistry;
use zeph_sanitizer::secret_shape::scrub_secret_shapes;
use zeph_sanitizer::{ContentSanitizer, ContentSource, ContentSourceKind};
use crate::state::SubAgentState;
const FORWARD_CHANNEL_CAPACITY: usize = 128;
const FORWARD_RING_CAPACITY: usize = 200;
const FORWARD_BUFFER_GRACE: Duration = Duration::from_secs(5);
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct ForwardSurfaces {
pub tui: bool,
pub bare: bool,
}
impl ForwardSurfaces {
#[must_use]
pub fn any(self) -> bool {
self.tui || self.bare
}
}
#[derive(Debug, Clone)]
pub(crate) struct RawChunk {
kind: ForwardChunkKind,
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub(crate) enum ForwardChunkKind {
Text(String),
Thinking(String),
Terminal(SubAgentState),
}
#[derive(Debug, Clone)]
pub(crate) struct SanitizedChunk {
pub(crate) task_id: Arc<str>,
pub(crate) def_name: Arc<str>,
pub(crate) seq: u64,
pub(crate) kind: SanitizedChunkKind,
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub(crate) enum SanitizedChunkKind {
Text(String),
Thinking(String),
Terminal(SubAgentState),
}
pub(crate) struct SanitizeLayers {
pub(crate) sanitizer: ContentSanitizer,
pub(crate) secret_registry: Option<Arc<SecretMaskRegistry>>,
pub(crate) pii_filter: Option<PiiFilter>,
}
fn sanitize_text(raw_text: &str, def_name: &str, layers: &SanitizeLayers) -> String {
let source = ContentSource::new(ContentSourceKind::ToolResult).with_identifier(def_name);
let mut body = layers.sanitizer.sanitize(raw_text, source).body;
if let Some(registry) = &layers.secret_registry {
body = registry.mask(&body);
}
body = scrub_secret_shapes(&body).into_owned();
if let Some(filter) = &layers.pii_filter {
body = filter.scrub(&body).into_owned();
}
body
}
const SANITIZE_HOLDBACK_BYTES: usize = 256;
const PEM_HOLDBACK_CAP_BYTES: usize = zeph_common::secrets::PEM_BODY_CAP;
#[derive(Default)]
struct PendingSanitizeBuffers {
text: String,
thinking: String,
}
fn last_header_before(buf: &str, before: usize) -> Option<usize> {
let region = &buf[..before];
[region.rfind("-----BEGIN"), region.rfind("---- BEGIN")]
.into_iter()
.flatten()
.max()
}
fn footer_end_after(buf: &str, marker_idx: usize) -> Option<usize> {
let tail = &buf[marker_idx..];
let end_offset = [tail.find("-----END"), tail.find("---- END")]
.into_iter()
.flatten()
.min()?;
Some(marker_idx + end_offset + 8)
}
fn pem_safe_flush_target(buf: &str, natural_target: usize) -> usize {
let mut candidate = natural_target;
loop {
let Some(marker_idx) = last_header_before(buf, candidate) else {
return candidate;
};
match footer_end_after(buf, marker_idx) {
Some(block_end) if block_end <= candidate => return candidate,
Some(_) => candidate = marker_idx,
None => {
let capped = buf.len().saturating_sub(PEM_HOLDBACK_CAP_BYTES);
return marker_idx.max(capped);
}
}
}
}
fn split_off_safe_prefix(buf: &mut String, holdback: usize) -> Option<String> {
if buf.is_empty() {
return None;
}
let target = if holdback == 0 {
buf.len()
} else {
let natural_target = buf.floor_char_boundary(buf.len().saturating_sub(holdback));
pem_safe_flush_target(buf, natural_target)
};
let boundary = buf.floor_char_boundary(target.min(buf.len()));
if boundary == 0 {
return None;
}
let prefix = buf[..boundary].to_owned();
buf.drain(..boundary);
Some(prefix)
}
fn try_flush_kind(
buf: &mut String,
holdback: usize,
def_name: &str,
layers: &SanitizeLayers,
wrap_kind: fn(String) -> SanitizedChunkKind,
) -> Option<SanitizedChunkKind> {
let safe = split_off_safe_prefix(buf, holdback)?;
Some(wrap_kind(sanitize_text(&safe, def_name, layers)))
}
fn make_sanitized_chunk(
task_id: &Arc<str>,
def_name: &Arc<str>,
seq: u64,
kind: SanitizedChunkKind,
) -> SanitizedChunk {
SanitizedChunk {
task_id: Arc::clone(task_id),
def_name: Arc::clone(def_name),
seq,
kind,
}
}
#[allow(clippy::too_many_arguments)]
fn flush_all_pending(
pending: &mut PendingSanitizeBuffers,
task_id: &Arc<str>,
def_name: &Arc<str>,
layers: &SanitizeLayers,
surfaces: ForwardSurfaces,
buffer: &ForwardBuffer,
dispatch: &mut impl FnMut(&SanitizedChunk, ForwardSurfaces, &ForwardBuffer),
emit_seq: &mut u64,
) {
if let Some(kind) = try_flush_kind(
&mut pending.text,
0,
def_name.as_ref(),
layers,
SanitizedChunkKind::Text,
) {
dispatch(
&make_sanitized_chunk(task_id, def_name, *emit_seq, kind),
surfaces,
buffer,
);
*emit_seq += 1;
}
if let Some(kind) = try_flush_kind(
&mut pending.thinking,
0,
def_name.as_ref(),
layers,
SanitizedChunkKind::Thinking,
) {
dispatch(
&make_sanitized_chunk(task_id, def_name, *emit_seq, kind),
surfaces,
buffer,
);
*emit_seq += 1;
}
}
pub(crate) struct ForwardSender {
tx: mpsc::Sender<RawChunk>,
task_id: Arc<str>,
def_name: Arc<str>,
seq: AtomicU64,
dropped: AtomicU64,
}
impl ForwardSender {
pub(crate) fn new(tx: mpsc::Sender<RawChunk>, task_id: Arc<str>, def_name: Arc<str>) -> Self {
Self {
tx,
task_id,
def_name,
seq: AtomicU64::new(0),
dropped: AtomicU64::new(0),
}
}
fn try_send(&self, kind: ForwardChunkKind) {
let seq = self.seq.fetch_add(1, Ordering::Relaxed);
let chunk = RawChunk { kind };
if self.tx.try_send(chunk).is_ok() {
tracing::debug!(
task_id = %self.task_id,
def_name = %self.def_name,
seq,
"subagent.forward.emit"
);
} else {
let dropped = self.dropped.fetch_add(1, Ordering::Relaxed) + 1;
tracing::warn!(
task_id = %self.task_id,
def_name = %self.def_name,
seq,
dropped,
"subagent.forward.drop: ingress channel full, chunk dropped"
);
}
}
pub(crate) fn send_text(&self, text: &str) {
if text.is_empty() {
return;
}
self.try_send(ForwardChunkKind::Text(text.to_owned()));
}
pub(crate) fn send_thinking(&self, text: &str) {
if text.is_empty() {
return;
}
self.try_send(ForwardChunkKind::Thinking(text.to_owned()));
}
pub(crate) fn send_terminal(&self, state: SubAgentState) {
tracing::debug!(task_id = %self.task_id, ?state, "subagent.forward.terminal");
self.try_send(ForwardChunkKind::Terminal(state));
}
}
pub(crate) type ForwardBuffer = std::sync::Mutex<HashMap<String, VecDeque<String>>>;
fn display_line(kind: &SanitizedChunkKind) -> Option<String> {
match kind {
SanitizedChunkKind::Text(t) => Some(t.clone()),
SanitizedChunkKind::Thinking(t) => Some(format!("[thinking] {t}")),
SanitizedChunkKind::Terminal(_) => None,
}
}
fn state_str(state: SubAgentState) -> &'static str {
match state {
SubAgentState::Submitted => "submitted",
SubAgentState::Working => "working",
SubAgentState::Completed => "completed",
SubAgentState::Failed => "failed",
SubAgentState::Canceled => "canceled",
}
}
fn emit_bare_line(chunk: &SanitizedChunk) {
#[derive(serde::Serialize)]
struct BareForwardEvent<'a> {
task_id: &'a str,
def_name: &'a str,
seq: u64,
kind: &'static str,
#[serde(skip_serializing_if = "Option::is_none")]
content: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
state: Option<&'static str>,
}
let (kind, content, state) = match &chunk.kind {
SanitizedChunkKind::Text(t) => ("text", Some(t.as_str()), None),
SanitizedChunkKind::Thinking(t) => ("thinking", Some(t.as_str()), None),
SanitizedChunkKind::Terminal(s) => ("terminal", None, Some(state_str(*s))),
};
let event = BareForwardEvent {
task_id: &chunk.task_id,
def_name: &chunk.def_name,
seq: chunk.seq,
kind,
content,
state,
};
if let Ok(line) = serde_json::to_string(&event) {
println!("{line}");
}
}
fn dispatch_chunk(chunk: &SanitizedChunk, surfaces: ForwardSurfaces, buffer: &ForwardBuffer) {
if surfaces.tui
&& let Some(line) = display_line(&chunk.kind)
{
let mut guard = buffer
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let ring = guard.entry(chunk.task_id.to_string()).or_default();
ring.push_back(line);
while ring.len() > FORWARD_RING_CAPACITY {
ring.pop_front();
}
}
if surfaces.bare {
emit_bare_line(chunk);
}
}
pub(crate) fn new_channel(
task_id: Arc<str>,
def_name: Arc<str>,
) -> (ForwardSender, mpsc::Receiver<RawChunk>) {
let (tx, rx) = mpsc::channel(FORWARD_CHANNEL_CAPACITY);
(ForwardSender::new(tx, task_id, def_name), rx)
}
pub(crate) async fn run_forward_drain(
task_id: Arc<str>,
def_name: Arc<str>,
rx: mpsc::Receiver<RawChunk>,
layers: SanitizeLayers,
surfaces: ForwardSurfaces,
buffer: Arc<ForwardBuffer>,
) {
run_forward_drain_with(
task_id,
def_name,
rx,
layers,
surfaces,
buffer,
dispatch_chunk,
)
.await;
}
async fn run_forward_drain_with(
task_id: Arc<str>,
def_name: Arc<str>,
mut rx: mpsc::Receiver<RawChunk>,
layers: SanitizeLayers,
surfaces: ForwardSurfaces,
buffer: Arc<ForwardBuffer>,
mut dispatch: impl FnMut(&SanitizedChunk, ForwardSurfaces, &ForwardBuffer),
) {
let mut pending = PendingSanitizeBuffers::default();
let mut emit_seq: u64 = 0;
loop {
if let Some(raw) = rx.recv().await {
match raw.kind {
ForwardChunkKind::Text(delta) => {
pending.text.push_str(&delta);
if let Some(kind) = try_flush_kind(
&mut pending.text,
SANITIZE_HOLDBACK_BYTES,
def_name.as_ref(),
&layers,
SanitizedChunkKind::Text,
) {
dispatch(
&make_sanitized_chunk(&task_id, &def_name, emit_seq, kind),
surfaces,
&buffer,
);
emit_seq += 1;
}
}
ForwardChunkKind::Thinking(delta) => {
pending.thinking.push_str(&delta);
if let Some(kind) = try_flush_kind(
&mut pending.thinking,
SANITIZE_HOLDBACK_BYTES,
def_name.as_ref(),
&layers,
SanitizedChunkKind::Thinking,
) {
dispatch(
&make_sanitized_chunk(&task_id, &def_name, emit_seq, kind),
surfaces,
&buffer,
);
emit_seq += 1;
}
}
ForwardChunkKind::Terminal(state) => {
flush_all_pending(
&mut pending,
&task_id,
&def_name,
&layers,
surfaces,
&buffer,
&mut dispatch,
&mut emit_seq,
);
let chunk = make_sanitized_chunk(
&task_id,
&def_name,
emit_seq,
SanitizedChunkKind::Terminal(state),
);
dispatch(&chunk, surfaces, &buffer);
break;
}
}
} else {
tracing::warn!(
task_id = %task_id,
"subagent.forward.terminal: ingress channel closed without an explicit \
terminal chunk — synthesizing hard-abort backstop"
);
flush_all_pending(
&mut pending,
&task_id,
&def_name,
&layers,
surfaces,
&buffer,
&mut dispatch,
&mut emit_seq,
);
let synthesized = make_sanitized_chunk(
&task_id,
&def_name,
emit_seq,
SanitizedChunkKind::Terminal(SubAgentState::Canceled),
);
dispatch(&synthesized, surfaces, &buffer);
break;
}
}
tokio::time::sleep(FORWARD_BUFFER_GRACE).await;
buffer
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(task_id.as_ref());
}
pub(crate) fn forwarded_tail(buffer: &ForwardBuffer, task_id: &str, n: usize) -> Vec<String> {
let guard = buffer
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
guard.get(task_id).map_or_else(Vec::new, |ring| {
ring.iter().rev().take(n).rev().cloned().collect()
})
}
pub(crate) fn new_buffer() -> Arc<ForwardBuffer> {
Arc::new(std::sync::Mutex::new(HashMap::new()))
}
#[cfg(test)]
mod tests {
use std::sync::atomic::AtomicUsize;
use zeph_config::sanitizer::PiiFilterConfig;
use super::*;
fn layers() -> SanitizeLayers {
SanitizeLayers {
sanitizer: ContentSanitizer::new(&zeph_config::ContentIsolationConfig::default()),
secret_registry: None,
pii_filter: None,
}
}
async fn run_and_count_terminals(
task_id: Arc<str>,
def_name: Arc<str>,
rx: mpsc::Receiver<RawChunk>,
surfaces: ForwardSurfaces,
buffer: Arc<ForwardBuffer>,
) -> usize {
let terminal_dispatches = Arc::new(AtomicUsize::new(0));
let counter = Arc::clone(&terminal_dispatches);
run_forward_drain_with(
task_id,
def_name,
rx,
layers(),
surfaces,
buffer,
move |chunk, surfaces, buffer| {
if matches!(chunk.kind, SanitizedChunkKind::Terminal(_)) {
counter.fetch_add(1, Ordering::SeqCst);
}
dispatch_chunk(chunk, surfaces, buffer);
},
)
.await;
terminal_dispatches.load(Ordering::SeqCst)
}
#[tokio::test(start_paused = true)]
async fn happy_path_emits_no_spurious_second_terminal() {
let task_id: Arc<str> = Arc::from("task-1");
let def_name: Arc<str> = Arc::from("agent-1");
let (sender, rx) = new_channel(Arc::clone(&task_id), Arc::clone(&def_name));
let buffer = new_buffer();
sender.send_text("hello");
sender.send_terminal(SubAgentState::Completed);
drop(sender);
let terminal_count = run_and_count_terminals(
Arc::clone(&task_id),
def_name,
rx,
ForwardSurfaces {
tui: true,
bare: false,
},
Arc::clone(&buffer),
)
.await;
assert_eq!(
terminal_count, 1,
"exactly one terminal chunk must be dispatched — a second would mean the drain \
looped back to recv() after the explicit terminal (C-new-1 regression)"
);
let tail = forwarded_tail(&buffer, &task_id, 10);
assert!(
tail.is_empty(),
"buffer entry must be evicted after grace window"
);
}
#[tokio::test(start_paused = true)]
async fn hard_abort_without_explicit_terminal_synthesizes_backstop() {
let task_id: Arc<str> = Arc::from("task-2");
let def_name: Arc<str> = Arc::from("agent-2");
let (sender, rx) = new_channel(Arc::clone(&task_id), Arc::clone(&def_name));
let buffer = new_buffer();
sender.send_text("partial output");
drop(sender);
let terminal_count = run_and_count_terminals(
Arc::clone(&task_id),
def_name,
rx,
ForwardSurfaces {
tui: true,
bare: false,
},
buffer,
)
.await;
assert_eq!(
terminal_count, 1,
"exactly one synthesized backstop terminal must be dispatched on hard abort"
);
}
#[tokio::test(start_paused = true)]
async fn zero_consumer_surfaces_still_drains_without_panicking() {
let task_id: Arc<str> = Arc::from("task-3");
let def_name: Arc<str> = Arc::from("agent-3");
let (sender, rx) = new_channel(Arc::clone(&task_id), Arc::clone(&def_name));
let buffer = new_buffer();
sender.send_text("no one is listening");
sender.send_terminal(SubAgentState::Completed);
drop(sender);
run_forward_drain(
task_id,
def_name,
rx,
layers(),
ForwardSurfaces::default(),
buffer,
)
.await;
}
#[tokio::test(start_paused = true)]
async fn secret_registry_masks_forwarded_text_and_thinking() {
use zeph_sanitizer::secret_mask::{SecretCategory, SecretMaskRegistry};
let registry = Arc::new(SecretMaskRegistry::new());
registry.register(
"MY_KEY",
"sk-live-topsecretvalue123",
SecretCategory::ApiKey,
);
let task_id: Arc<str> = Arc::from("task-secret");
let def_name: Arc<str> = Arc::from("agent-secret");
let (sender, rx) = new_channel(Arc::clone(&task_id), Arc::clone(&def_name));
let buffer = new_buffer();
sender.send_text("the key is sk-live-topsecretvalue123, use it wisely");
sender.send_thinking("I will use sk-live-topsecretvalue123 to authenticate");
sender.send_terminal(SubAgentState::Completed);
drop(sender);
let seen = Arc::new(std::sync::Mutex::new(Vec::<SanitizedChunk>::new()));
let collected = Arc::clone(&seen);
let layers = SanitizeLayers {
sanitizer: ContentSanitizer::new(&zeph_config::ContentIsolationConfig::default()),
secret_registry: Some(registry),
pii_filter: None,
};
run_forward_drain_with(
task_id,
def_name,
rx,
layers,
ForwardSurfaces {
tui: true,
bare: false,
},
buffer,
move |chunk, surfaces, buffer| {
collected.lock().unwrap().push(chunk.clone());
dispatch_chunk(chunk, surfaces, buffer);
},
)
.await;
let chunks = seen.lock().unwrap();
for chunk in chunks.iter() {
match &chunk.kind {
SanitizedChunkKind::Text(t) | SanitizedChunkKind::Thinking(t) => {
assert!(
!t.contains("sk-live-topsecretvalue123"),
"forwarded content must not contain the raw secret: {t}"
);
}
SanitizedChunkKind::Terminal(_) => {}
}
}
}
#[tokio::test(start_paused = true)]
async fn registered_secret_that_also_matches_a_shape_gets_typed_placeholder_not_generic() {
use zeph_sanitizer::secret_mask::{SecretCategory, SecretMaskRegistry};
let registry = Arc::new(SecretMaskRegistry::new());
registry.register(
"MY_KEY",
"sk-live-topsecretvalue123",
SecretCategory::ApiKey,
);
let task_id: Arc<str> = Arc::from("task-secret-typed");
let def_name: Arc<str> = Arc::from("agent-secret-typed");
let (sender, rx) = new_channel(Arc::clone(&task_id), Arc::clone(&def_name));
let buffer = new_buffer();
sender.send_text("the key is sk-live-topsecretvalue123, use it wisely");
sender.send_terminal(SubAgentState::Completed);
drop(sender);
let seen = Arc::new(std::sync::Mutex::new(Vec::<SanitizedChunk>::new()));
let collected = Arc::clone(&seen);
let layers = SanitizeLayers {
sanitizer: ContentSanitizer::new(&zeph_config::ContentIsolationConfig::default()),
secret_registry: Some(registry),
pii_filter: None,
};
run_forward_drain_with(
task_id,
def_name,
rx,
layers,
ForwardSurfaces {
tui: true,
bare: false,
},
buffer,
move |chunk, surfaces, buffer| {
collected.lock().unwrap().push(chunk.clone());
dispatch_chunk(chunk, surfaces, buffer);
},
)
.await;
let combined = collect_forwarded_text(&seen.lock().unwrap());
assert!(
!combined.contains("sk-live-topsecretvalue123"),
"raw secret must not survive the pipeline: {combined}"
);
assert!(
combined.contains("<SECRET:api_key:"),
"registry masking must run first and produce its typed placeholder: {combined}"
);
assert!(
!combined.contains("[REDACTED]"),
"shape scrub must not double-mask the registry's own placeholder output: {combined}"
);
}
#[tokio::test(start_paused = true)]
async fn generic_secret_shape_masked_without_registration() {
let task_id: Arc<str> = Arc::from("task-shape");
let def_name: Arc<str> = Arc::from("agent-shape");
let (sender, rx) = new_channel(Arc::clone(&task_id), Arc::clone(&def_name));
let buffer = new_buffer();
sender.send_text("here is a key: sk-test-abc123def456, use it wisely");
sender.send_terminal(SubAgentState::Completed);
drop(sender);
let seen = Arc::new(std::sync::Mutex::new(Vec::<SanitizedChunk>::new()));
let collected = Arc::clone(&seen);
run_forward_drain_with(
task_id,
def_name,
rx,
layers(),
ForwardSurfaces {
tui: true,
bare: false,
},
buffer,
move |chunk, surfaces, buffer| {
collected.lock().unwrap().push(chunk.clone());
dispatch_chunk(chunk, surfaces, buffer);
},
)
.await;
let combined = collect_forwarded_text(&seen.lock().unwrap());
assert!(
!combined.contains("sk-test-abc123def456"),
"generic secret-shaped string must be masked without prior registration: {combined}"
);
assert!(
combined.contains("[REDACTED]"),
"masked placeholder must be present in the combined forwarded text: {combined}"
);
}
#[tokio::test(start_paused = true)]
async fn pii_filter_scrubs_forwarded_email() {
let task_id: Arc<str> = Arc::from("task-pii");
let def_name: Arc<str> = Arc::from("agent-pii");
let (sender, rx) = new_channel(Arc::clone(&task_id), Arc::clone(&def_name));
let buffer = new_buffer();
sender.send_text("contact me at victim@example.com for details");
sender.send_terminal(SubAgentState::Completed);
drop(sender);
let seen = Arc::new(std::sync::Mutex::new(Vec::<SanitizedChunk>::new()));
let collected = Arc::clone(&seen);
let layers = SanitizeLayers {
sanitizer: ContentSanitizer::new(&zeph_config::ContentIsolationConfig::default()),
secret_registry: None,
pii_filter: Some(PiiFilter::new(PiiFilterConfig::default())),
};
run_forward_drain_with(
task_id,
def_name,
rx,
layers,
ForwardSurfaces {
tui: true,
bare: false,
},
buffer,
move |chunk, surfaces, buffer| {
collected.lock().unwrap().push(chunk.clone());
dispatch_chunk(chunk, surfaces, buffer);
},
)
.await;
let chunks = seen.lock().unwrap();
let text_chunk = chunks
.iter()
.find(|c| matches!(c.kind, SanitizedChunkKind::Text(_)))
.expect("one text chunk must have been dispatched");
let SanitizedChunkKind::Text(ref t) = text_chunk.kind else {
unreachable!()
};
assert!(
!t.contains("victim@example.com"),
"forwarded content must not contain the raw email address: {t}"
);
}
fn collect_forwarded_text(chunks: &[SanitizedChunk]) -> String {
chunks
.iter()
.filter_map(|c| match &c.kind {
SanitizedChunkKind::Text(t) => Some(t.as_str()),
_ => None,
})
.collect()
}
#[tokio::test(start_paused = true)]
async fn secret_split_across_two_deltas_is_still_masked() {
use zeph_sanitizer::secret_mask::{SecretCategory, SecretMaskRegistry};
let secret_value = "sk-live-topsecretvalue123456789";
let registry = Arc::new(SecretMaskRegistry::new());
registry.register("MY_KEY", secret_value, SecretCategory::ApiKey);
let (first_half, second_half) = secret_value.split_at(secret_value.len() / 2);
let task_id: Arc<str> = Arc::from("task-split");
let def_name: Arc<str> = Arc::from("agent-split");
let (sender, rx) = new_channel(Arc::clone(&task_id), Arc::clone(&def_name));
let buffer = new_buffer();
sender.send_text(&format!("the key is {first_half}"));
sender.send_text(&format!("{second_half}, use it wisely"));
sender.send_terminal(SubAgentState::Completed);
drop(sender);
let seen = Arc::new(std::sync::Mutex::new(Vec::<SanitizedChunk>::new()));
let collected = Arc::clone(&seen);
let layers = SanitizeLayers {
sanitizer: ContentSanitizer::new(&zeph_config::ContentIsolationConfig::default()),
secret_registry: Some(registry),
pii_filter: None,
};
run_forward_drain_with(
task_id,
def_name,
rx,
layers,
ForwardSurfaces {
tui: true,
bare: false,
},
buffer,
move |chunk, surfaces, buffer| {
collected.lock().unwrap().push(chunk.clone());
dispatch_chunk(chunk, surfaces, buffer);
},
)
.await;
let combined = collect_forwarded_text(&seen.lock().unwrap());
assert!(
!combined.contains(secret_value),
"secret split across two forwarded deltas must still be masked: {combined}"
);
assert!(
combined.contains("<SECRET:api_key:"),
"masked placeholder must be present in the combined forwarded text: {combined}"
);
}
#[tokio::test(start_paused = true)]
async fn email_split_across_two_deltas_is_still_scrubbed() {
let email = "victim@example.com";
let (first_half, second_half) = email.split_at(email.len() / 2);
let task_id: Arc<str> = Arc::from("task-split-pii");
let def_name: Arc<str> = Arc::from("agent-split-pii");
let (sender, rx) = new_channel(Arc::clone(&task_id), Arc::clone(&def_name));
let buffer = new_buffer();
sender.send_text(&format!("contact me at {first_half}"));
sender.send_text(&format!("{second_half} for details"));
sender.send_terminal(SubAgentState::Completed);
drop(sender);
let seen = Arc::new(std::sync::Mutex::new(Vec::<SanitizedChunk>::new()));
let collected = Arc::clone(&seen);
let layers = SanitizeLayers {
sanitizer: ContentSanitizer::new(&zeph_config::ContentIsolationConfig::default()),
secret_registry: None,
pii_filter: Some(PiiFilter::new(PiiFilterConfig::default())),
};
run_forward_drain_with(
task_id,
def_name,
rx,
layers,
ForwardSurfaces {
tui: true,
bare: false,
},
buffer,
move |chunk, surfaces, buffer| {
collected.lock().unwrap().push(chunk.clone());
dispatch_chunk(chunk, surfaces, buffer);
},
)
.await;
let combined = collect_forwarded_text(&seen.lock().unwrap());
assert!(
!combined.contains(email),
"email split across two forwarded deltas must still be scrubbed: {combined}"
);
}
#[tokio::test(start_paused = true)]
async fn secret_split_across_progressive_flush_boundary_is_still_masked() {
use zeph_sanitizer::secret_mask::{SecretCategory, SecretMaskRegistry};
let secret_value = "sk-live-anothersecretvalue987654321";
let registry = Arc::new(SecretMaskRegistry::new());
registry.register("MY_KEY", secret_value, SecretCategory::ApiKey);
let (first_half, second_half) = secret_value.split_at(secret_value.len() / 2);
let task_id: Arc<str> = Arc::from("task-window");
let def_name: Arc<str> = Arc::from("agent-window");
let (sender, rx) = new_channel(Arc::clone(&task_id), Arc::clone(&def_name));
let buffer = new_buffer();
for i in 0..40 {
sender.send_text(&format!("filler-chunk-{i:03} "));
}
sender.send_text(first_half);
sender.send_text(second_half);
sender.send_terminal(SubAgentState::Completed);
drop(sender);
let seen = Arc::new(std::sync::Mutex::new(Vec::<SanitizedChunk>::new()));
let collected = Arc::clone(&seen);
let layers = SanitizeLayers {
sanitizer: ContentSanitizer::new(&zeph_config::ContentIsolationConfig::default()),
secret_registry: Some(registry),
pii_filter: None,
};
run_forward_drain_with(
task_id,
def_name,
rx,
layers,
ForwardSurfaces {
tui: true,
bare: false,
},
buffer,
move |chunk, surfaces, buffer| {
collected.lock().unwrap().push(chunk.clone());
dispatch_chunk(chunk, surfaces, buffer);
},
)
.await;
let seen = seen.lock().unwrap();
let text_chunk_count = seen
.iter()
.filter(|c| matches!(c.kind, SanitizedChunkKind::Text(_)))
.count();
assert!(
text_chunk_count > 1,
"filler well over the holdback window must have produced at least one \
progressive flush before the terminal-triggered final flush, got \
{text_chunk_count} text chunk(s)"
);
let combined = collect_forwarded_text(&seen);
assert!(
!combined.contains(secret_value),
"secret split across the streaming boundary must still be masked: {combined}"
);
}
#[tokio::test(start_paused = true)]
async fn pem_block_split_across_chunk_boundary_is_still_fully_masked() {
let task_id: Arc<str> = Arc::from("task-pem-chunked");
let def_name: Arc<str> = Arc::from("agent-pem-chunked");
let (sender, rx) = new_channel(Arc::clone(&task_id), Arc::clone(&def_name));
let buffer = new_buffer();
let body = "X".repeat(300);
sender.send_text("intro text before the key ");
sender.send_text(&format!("-----BEGIN RSA PRIVATE KEY-----\n{body}"));
sender.send_text("\n-----END RSA PRIVATE KEY-----\nfollowing text after the key");
sender.send_terminal(SubAgentState::Completed);
drop(sender);
let seen = Arc::new(std::sync::Mutex::new(Vec::<SanitizedChunk>::new()));
let collected = Arc::clone(&seen);
run_forward_drain_with(
task_id,
def_name,
rx,
layers(),
ForwardSurfaces {
tui: true,
bare: false,
},
buffer,
move |chunk, surfaces, buffer| {
collected.lock().unwrap().push(chunk.clone());
dispatch_chunk(chunk, surfaces, buffer);
},
)
.await;
let combined = collect_forwarded_text(&seen.lock().unwrap());
assert!(
!combined.contains(&body),
"PEM body must not survive split across a chunk boundary: {combined}"
);
assert!(
!combined.contains('X'),
"no raw PEM body fragment may leak through an isolated flush of a middle slice \
that itself contains neither BEGIN nor END: {combined}"
);
assert!(
combined.contains("[REDACTED_PEM_KEY]"),
"PEM placeholder must be present in the combined forwarded text: {combined}"
);
assert!(
combined.contains("intro text before the key"),
"text preceding the PEM block must still be forwarded: {combined}"
);
assert!(
combined.contains("following text after the key"),
"text following the PEM block must still be forwarded: {combined}"
);
}
#[tokio::test(start_paused = true)]
async fn complete_key_immediately_followed_by_different_pem_armor_leaks_nothing() {
let task_id: Arc<str> = Arc::from("task-pem-bundle");
let def_name: Arc<str> = Arc::from("agent-pem-bundle");
let (sender, rx) = new_channel(Arc::clone(&task_id), Arc::clone(&def_name));
let buffer = new_buffer();
let key_body = "Z".repeat(400);
sender.send_text(&format!(
"-----BEGIN RSA PRIVATE KEY-----\n{key_body}\n-----END RSA PRIVATE KEY-----\n\
-----BEGIN CERTIFICATE-----\nMIIBcertbody"
));
sender.send_text("\n-----END CERTIFICATE-----\nbundle complete");
sender.send_terminal(SubAgentState::Completed);
drop(sender);
let seen = Arc::new(std::sync::Mutex::new(Vec::<SanitizedChunk>::new()));
let collected = Arc::clone(&seen);
run_forward_drain_with(
task_id,
def_name,
rx,
layers(),
ForwardSurfaces {
tui: true,
bare: false,
},
buffer,
move |chunk, surfaces, buffer| {
collected.lock().unwrap().push(chunk.clone());
dispatch_chunk(chunk, surfaces, buffer);
},
)
.await;
let combined = collect_forwarded_text(&seen.lock().unwrap());
assert!(
!combined.contains('Z'),
"no fragment of the RSA key body may leak when immediately followed by a \
different PEM armor type in the same delta: {combined}"
);
assert!(
combined.contains("[REDACTED_PEM_KEY]"),
"PEM placeholder must be present for the private key: {combined}"
);
assert!(
combined.contains("bundle complete"),
"text following the bundle must still be forwarded: {combined}"
);
}
#[tokio::test(start_paused = true)]
async fn complete_key_followed_by_bare_trailing_begin_leaks_nothing() {
let task_id: Arc<str> = Arc::from("task-pem-trailing-begin");
let def_name: Arc<str> = Arc::from("agent-pem-trailing-begin");
let (sender, rx) = new_channel(Arc::clone(&task_id), Arc::clone(&def_name));
let buffer = new_buffer();
let key_body = "Z".repeat(400);
sender.send_text(&format!(
"-----BEGIN RSA PRIVATE KEY-----\n{key_body}\n-----END RSA PRIVATE KEY-----\n-----BEGIN"
));
sender.send_text(" EC PRIVATE KEY-----\nsecondbody\n-----END EC PRIVATE KEY-----\ndone");
sender.send_terminal(SubAgentState::Completed);
drop(sender);
let seen = Arc::new(std::sync::Mutex::new(Vec::<SanitizedChunk>::new()));
let collected = Arc::clone(&seen);
run_forward_drain_with(
task_id,
def_name,
rx,
layers(),
ForwardSurfaces {
tui: true,
bare: false,
},
buffer,
move |chunk, surfaces, buffer| {
collected.lock().unwrap().push(chunk.clone());
dispatch_chunk(chunk, surfaces, buffer);
},
)
.await;
let combined = collect_forwarded_text(&seen.lock().unwrap());
assert!(
!combined.contains('Z'),
"no fragment of the first key body may leak when a bare trailing -----BEGIN \
follows it in the same buffer: {combined}"
);
assert!(
!combined.contains("secondbody"),
"no fragment of the second key body may leak either: {combined}"
);
assert_eq!(
combined.matches("[REDACTED_PEM_KEY]").count(),
2,
"both blocks must be redacted independently: {combined}"
);
assert!(
combined.contains("done"),
"trailing text must survive: {combined}"
);
}
#[tokio::test(start_paused = true)]
async fn pem_holdback_boundary_computation_does_not_panic_on_multibyte_text() {
let task_id: Arc<str> = Arc::from("task-pem-multibyte");
let def_name: Arc<str> = Arc::from("agent-pem-multibyte");
let (sender, rx) = new_channel(Arc::clone(&task_id), Arc::clone(&def_name));
let buffer = new_buffer();
let cjk_body: String = std::iter::repeat_n('中', 300).collect();
sender.send_text(&format!("-----BEGIN RSA PRIVATE KEY-----\n{cjk_body}"));
sender.send_text("\n-----END RSA PRIVATE KEY-----\ndone");
sender.send_terminal(SubAgentState::Completed);
drop(sender);
let seen = Arc::new(std::sync::Mutex::new(Vec::<SanitizedChunk>::new()));
let collected = Arc::clone(&seen);
run_forward_drain_with(
task_id,
def_name,
rx,
layers(),
ForwardSurfaces {
tui: true,
bare: false,
},
buffer,
move |chunk, surfaces, buffer| {
collected.lock().unwrap().push(chunk.clone());
dispatch_chunk(chunk, surfaces, buffer);
},
)
.await;
let combined = collect_forwarded_text(&seen.lock().unwrap());
assert!(
!combined.contains('中'),
"CJK key body must not leak: {combined}"
);
assert!(combined.contains("[REDACTED_PEM_KEY]"));
assert!(
combined.contains("done"),
"trailing text must survive: {combined}"
);
}
#[tokio::test(start_paused = true)]
async fn buffer_entry_survives_during_grace_window_then_evicted() {
let task_id: Arc<str> = Arc::from("task-grace");
let def_name: Arc<str> = Arc::from("agent-grace");
let (sender, rx) = new_channel(Arc::clone(&task_id), Arc::clone(&def_name));
let buffer = new_buffer();
sender.send_text("visible during the grace window");
sender.send_terminal(SubAgentState::Completed);
drop(sender);
let drain_buffer = Arc::clone(&buffer);
let drain_task_id = Arc::clone(&task_id);
let handle = tokio::spawn(run_forward_drain(
drain_task_id,
def_name,
rx,
layers(),
ForwardSurfaces {
tui: true,
bare: false,
},
drain_buffer,
));
tokio::time::advance(Duration::from_millis(1)).await;
tokio::task::yield_now().await;
let mid_window_tail = forwarded_tail(&buffer, &task_id, 10);
assert_eq!(
mid_window_tail.len(),
1,
"exactly one forwarded line expected"
);
assert!(
mid_window_tail[0].contains("visible during the grace window"),
"the transcript must still be visible during the grace window, got: {:?}",
mid_window_tail[0]
);
tokio::time::advance(FORWARD_BUFFER_GRACE + Duration::from_millis(1)).await;
handle.await.expect("drain task must not panic");
let post_eviction_tail = forwarded_tail(&buffer, &task_id, 10);
assert!(
post_eviction_tail.is_empty(),
"buffer entry must be evicted once the grace window elapses"
);
}
#[test]
fn empty_text_is_not_sent() {
let task_id: Arc<str> = Arc::from("task-4");
let def_name: Arc<str> = Arc::from("agent-4");
let (sender, mut rx) = new_channel(task_id, def_name);
sender.send_text("");
sender.send_thinking("");
drop(sender);
assert!(
rx.try_recv().is_err(),
"empty text/thinking must not be sent onto the ingress channel"
);
}
#[test]
fn channel_full_increments_drop_counter_and_does_not_panic() {
let task_id: Arc<str> = Arc::from("task-5");
let def_name: Arc<str> = Arc::from("agent-5");
let (sender, mut rx) = new_channel(task_id, def_name);
for i in 0..FORWARD_CHANNEL_CAPACITY + 10 {
sender.send_text(&format!("chunk {i}"));
}
let mut received = 0;
while rx.try_recv().is_ok() {
received += 1;
}
assert!(
received > 0,
"at least some chunks must have been delivered"
);
assert!(
received <= FORWARD_CHANNEL_CAPACITY,
"received must never exceed channel capacity"
);
}
#[test]
fn forward_surfaces_any() {
assert!(!ForwardSurfaces::default().any());
assert!(
ForwardSurfaces {
tui: true,
bare: false
}
.any()
);
assert!(
ForwardSurfaces {
tui: false,
bare: true
}
.any()
);
}
}