use super::{WireEvent, approvals::Approvals};
use crate::{
agent::AgentOutputSink,
cancellation::AgentCancellation,
output::{ActivityEvent, ActivityKind, ActivitySender, ActivityStatus, OutputEvent},
sessions::{Session, usage::SessionUsageLedger, usage_recorder::SessionUsageRecorder},
};
use crossbeam_channel::{SendTimeoutError, Sender};
use serde_json::{Value, json};
use std::{
collections::HashMap,
sync::{
Arc, Mutex,
atomic::{AtomicBool, Ordering},
},
time::{Duration, Instant},
};
const MAX_CONTENT: usize = 8 * 1024 * 1024;
const MAX_SEGMENTS: usize = 512;
const TEXT_PART: usize = 4000;
#[derive(Clone)]
pub(super) struct Events {
sender: Sender<WireEvent>,
pub(super) cancellation: AgentCancellation,
cancel: Arc<AtomicBool>,
session_id: String,
run_id: String,
pub(super) stopped: Arc<AtomicBool>,
pub(super) failed: Arc<AtomicBool>,
}
impl Events {
pub(super) fn new(
sender: Sender<WireEvent>,
cancel: Arc<AtomicBool>,
session_id: String,
run_id: String,
) -> Self {
Self {
sender,
cancellation: AgentCancellation::new(Arc::clone(&cancel)),
cancel,
session_id,
run_id,
stopped: Arc::new(AtomicBool::new(false)),
failed: Arc::new(AtomicBool::new(false)),
}
}
pub(super) fn cancel(&self) {
self.cancel.store(true, Ordering::SeqCst);
}
pub(super) fn emit(
&self,
name: &'static str,
mut data: Value,
critical: bool,
) -> anyhow::Result<()> {
if self.stopped.load(Ordering::SeqCst) {
return Ok(());
}
if self.failed.load(Ordering::SeqCst) {
return Err(super::RunFailure("output_unavailable").into());
}
data["session_id"] = json!(self.session_id);
data["run_id"] = json!(self.run_id);
let mut event = WireEvent { name, data };
if !critical {
let _ = self.sender.try_send(event);
return Ok(());
}
let deadline = Instant::now() + Duration::from_secs(2);
loop {
match self.sender.send_timeout(event, Duration::from_millis(20)) {
Ok(()) => return Ok(()),
Err(SendTimeoutError::Timeout(returned)) if Instant::now() < deadline => {
event = returned
}
Err(_) => {
self.failed.store(true, Ordering::SeqCst);
self.cancel();
return Err(super::RunFailure("output_unavailable").into());
}
}
}
}
pub(super) fn diagnostic(&self, code: &'static str, summary: &str) -> anyhow::Result<()> {
self.emit(
"diagnostic",
json!({"code":code,"summary":summary_text(summary)}),
true,
)
}
}
pub(super) fn summary_text(text: &str) -> String {
let text = crate::output::sanitize_display_text(text);
let end = text
.char_indices()
.find(|(i, c)| i + c.len_utf8() > 1024)
.map(|(i, _)| i)
.unwrap_or(text.len());
text[..end].to_owned()
}
fn parts(text: &str) -> impl Iterator<Item = &str> {
let mut remaining = text;
std::iter::from_fn(move || {
if remaining.is_empty() {
return None;
}
let mut end = remaining.len().min(TEXT_PART);
while !remaining.is_char_boundary(end) {
end -= 1;
}
let (part, rest) = remaining.split_at(end);
remaining = rest;
Some(part)
})
}
struct Segment {
id: String,
text: String,
delta_part: usize,
redactor: crate::output::StreamingRedactor,
}
#[derive(Default)]
struct ActivityProjection {
entries: HashMap<String, (String, &'static str, String)>,
seen: std::collections::HashSet<String>,
closed: bool,
}
impl ActivityProjection {
fn admit(&mut self, event: &ActivityEvent) -> anyhow::Result<bool> {
if self.closed {
return Ok(false);
}
let id = match event {
ActivityEvent::Started { id, .. }
| ActivityEvent::UsageSnapshot { id, .. }
| ActivityEvent::UsageUpdate { id, .. } => id.as_str(),
_ => return Ok(true),
};
if id.len() > 1024 || (!self.seen.contains(id) && self.seen.len() == 4096) {
return Err(super::RunFailure("content_limit").into());
}
self.seen.insert(id.to_owned());
Ok(true)
}
fn project(&mut self, events: &Events, event: ActivityEvent) -> anyhow::Result<()> {
match event {
ActivityEvent::Started {
id, kind, metadata, ..
} => {
let kind = match kind {
ActivityKind::Tool => "tool",
ActivityKind::SubagentBatch | ActivityKind::SubagentTask => "subagent",
ActivityKind::Compaction => "compaction",
_ => return Ok(()),
};
if id.as_str().len() > 1024 || self.entries.len() >= 4096 {
return Err(super::RunFailure("content_limit").into());
}
let entry = self.entries.entry(id.0).or_insert_with(|| {
(
uuid::Uuid::new_v4().to_string(),
kind,
summary_text(&metadata.label),
)
});
events.emit("activity", json!({"activity_id":entry.0,"kind":entry.1,"state":"started","name":entry.2,"summary":"started"}), true)
}
ActivityEvent::Finished { id, status, .. } => {
if let Some((id, kind, name)) = self.entries.remove(id.as_str()) {
let state = if matches!(status, ActivityStatus::Success) {
"completed"
} else {
"failed"
};
events.emit("activity", json!({"activity_id":id,"kind":kind,"state":state,"name":name,"summary":state}), true)?;
}
Ok(())
}
_ => Ok(()), }
}
}
pub(super) struct Sink {
pub(super) events: Events,
session: Session,
approvals: Arc<Approvals>,
usage: Arc<SessionUsageRecorder>,
activity: Arc<Mutex<ActivityProjection>>,
pub(super) degraded: Arc<AtomicBool>,
pub(super) persisted: bool,
segments: Vec<Segment>,
current: Option<usize>,
prior: String,
content_bytes: usize,
compaction_id: Option<String>,
}
impl Sink {
pub(super) fn new(
events: Events,
session: Session,
cwd: std::path::PathBuf,
approvals: Arc<Approvals>,
) -> Self {
let usage = Arc::new(SessionUsageRecorder::new(
Some(session.clone()),
cwd,
Arc::new(Mutex::new(SessionUsageLedger::default())),
));
Self {
events,
session,
approvals,
usage,
activity: Default::default(),
degraded: Arc::new(AtomicBool::new(false)),
persisted: false,
segments: Vec::new(),
current: None,
prior: String::new(),
content_bytes: 0,
compaction_id: None,
}
}
fn segment(&mut self) -> anyhow::Result<usize> {
if let Some(index) = self.current {
return Ok(index);
}
if self.segments.len() == MAX_SEGMENTS {
return Err(super::RunFailure("content_limit").into());
}
let index = self.segments.len();
self.segments.push(Segment {
id: uuid::Uuid::new_v4().to_string(),
text: String::new(),
delta_part: 0,
redactor: Default::default(),
});
self.current = Some(index);
Ok(index)
}
fn boundary(&mut self) -> anyhow::Result<()> {
if let Some(index) = self.current.take() {
self.reconcile(index, true)?;
self.prior.push_str(&self.segments[index].text);
}
Ok(())
}
fn reconcile(&self, index: usize, partial: bool) -> anyhow::Result<()> {
let segment = &self.segments[index];
let text = crate::output::sanitize_display_text(&segment.text);
if text.is_empty() {
return self.events.emit("message.reconcile", json!({"message_id":segment.id,"part":0,"text":"","last_part":true,"partial":partial}), true);
}
let mut chunks = parts(&text).peekable();
let mut part = 0;
while let Some(text) = chunks.next() {
self.events.emit("message.reconcile", json!({"message_id":segment.id,"part":part,"text":text,"last_part":chunks.peek().is_none(),"partial":partial}), true)?;
part += 1;
}
Ok(())
}
pub(super) fn finish(&self, partial: bool) -> anyhow::Result<()> {
self.activity
.lock()
.unwrap_or_else(|e| e.into_inner())
.closed = true;
for index in 0..self.segments.len() {
self.reconcile(index, partial)?;
}
Ok(())
}
fn record_usage(&self, event: &ActivityEvent) {
if self.usage.record_activity(event).is_err() {
self.degraded.store(true, Ordering::SeqCst);
}
}
}
impl AgentOutputSink for Sink {
fn user_input_persisted(&mut self) -> anyhow::Result<()> {
if self.persisted {
return Ok(());
}
self.persisted = true;
let snapshot =
super::super::sessions::Snapshot::read(&self.session).map_err(super::RunFailure)?;
let (message_id, revision) = snapshot
.input_identity()
.ok_or(super::RunFailure("session_invalid"))?;
self.events.emit(
"input.persisted",
json!({"message_id":message_id,"durable":true,"revision":revision}),
true,
)?;
Ok(())
}
fn persistence_degraded(&mut self) {
self.degraded.store(true, Ordering::SeqCst);
}
fn assistant_delta(&mut self, text: &str) -> anyhow::Result<()> {
let text = crate::output::sanitize_display_controls(text);
self.content_bytes = self.content_bytes.saturating_add(text.len());
if self.content_bytes > MAX_CONTENT {
return Err(super::RunFailure("content_limit").into());
}
let index = self.segment()?;
let segment = &mut self.segments[index];
segment.text.push_str(&text);
let text = segment.redactor.push(&text);
for text in parts(&text) {
let part = segment.delta_part;
segment.delta_part += 1;
self.events.emit(
"message.delta",
json!({"message_id":segment.id,"part":part,"text":text}),
false,
)?;
}
Ok(())
}
fn tool_block(&mut self, _block: &str) -> anyhow::Result<()> {
Ok(())
}
fn activity_event(&mut self, event: ActivityEvent) -> anyhow::Result<()> {
let mut activity = self.activity.lock().unwrap_or_else(|e| e.into_inner());
if !activity.admit(&event)? {
return Ok(());
}
self.record_usage(&event);
activity.project(&self.events, event)
}
fn activity_sender(&self) -> Option<ActivitySender> {
let events = self.events.clone();
let activity = Arc::clone(&self.activity);
let usage = Arc::clone(&self.usage);
let degraded = Arc::clone(&self.degraded);
Some(Arc::new(move |event| {
let mut activity = activity.lock().unwrap_or_else(|e| e.into_inner());
let result = (|| {
if !activity.admit(&event)? {
return Ok(());
}
if usage.record_activity(&event).is_err() {
degraded.store(true, Ordering::SeqCst);
}
activity.project(&events, event)
})();
if result.is_err() {
events.failed.store(true, Ordering::SeqCst);
events.cancel();
}
}))
}
fn request_bash_approval(
&mut self,
request: crate::protection::bash::BashApprovalRequest,
cancellation: &AgentCancellation,
) -> anyhow::Result<bool> {
self.approvals.request(cancellation, |id| self.events.emit("approval.requested", json!({"approval_id":id,"summary":summary_text(&request.command),"expires_in_ms":super::approvals::APPROVAL_TIMEOUT.as_millis()}), true))
}
fn output_event(&mut self, event: OutputEvent) -> anyhow::Result<()> {
match event {
OutputEvent::AssistantDelta { text } => self.assistant_delta(&text),
OutputEvent::AssistantComplete { text } => {
let current = self
.current
.map(|i| self.segments[i].text.as_str())
.unwrap_or_default();
let output = crate::sessions::chat::reconciled_assistant_output(
&crate::output::sanitize_display_controls(&text),
&self.prior,
current,
);
if output.is_empty() && self.current.is_none() {
return Ok(());
}
let index = self.segment()?;
self.content_bytes = self
.content_bytes
.saturating_sub(self.segments[index].text.len())
.saturating_add(output.len());
if self.content_bytes > MAX_CONTENT {
return Err(super::RunFailure("content_limit").into());
}
self.segments[index].text = output;
self.reconcile(index, true)
}
OutputEvent::ToolStarted { .. } => self.boundary(),
OutputEvent::Diagnostic { message, .. } => {
self.events.diagnostic("execution_notice", &message)
}
OutputEvent::HookDiagnostic { diagnostic } => self
.events
.diagnostic("hook_notice", &diagnostic.sanitized_message()),
OutputEvent::UsageSnapshot {
usage,
request_sequence,
final_usage,
} => {
if self
.usage
.record("main", usage, request_sequence, final_usage)
.is_err()
{
self.persistence_degraded();
}
Ok(())
}
OutputEvent::ContextUsage {
request_sequence, ..
} => {
if self
.usage
.request_started("main", request_sequence)
.is_err()
{
self.persistence_degraded();
}
Ok(())
}
OutputEvent::CompactionStarted => {
self.boundary()?;
let id = uuid::Uuid::new_v4().to_string();
self.compaction_id = Some(id.clone());
self.events.emit("activity", json!({"activity_id":id,"kind":"compaction","state":"started","name":"compaction","summary":"started"}), true)
}
OutputEvent::CompactionCompleted { .. } | OutputEvent::CompactionFailed { .. } => {
let state = if matches!(event, OutputEvent::CompactionCompleted { .. }) {
"completed"
} else {
"failed"
};
let id = self
.compaction_id
.take()
.unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
self.prior.clear(); self.events.emit("activity", json!({"activity_id":id,"kind":"compaction","state":state,"name":"compaction","summary":state}), true)
}
_ => Ok(()),
}
}
}