1mod context;
4mod lifecycle;
5pub use context::{ExecutionContextTracker, decision_completed_event, validate_decision_input};
6pub use lifecycle::{
7 SharedLifecycleEmitter, ToolOutputPayload, error_item_completed_event, file_change_completed_event,
8 tool_invocation_completed_event, tool_output_completed_event, tool_output_item_id, tool_output_payload_from_value,
9 tool_output_started_event, tool_output_updated_event, tool_started_event,
10};
11
12use crate::core::threads::{SubmissionId, ThreadRuntimeHandle};
13use crate::exec::events::{
14 CommandExecutionItem, CommandExecutionStatus, CompactionMode, CompactionTrigger, EVENT_SCHEMA_VERSION, ErrorItem,
15 HarnessEventItem, HarnessEventKind, ItemCompletedEvent, ItemStartedEvent, ThreadCompactBoundaryEvent,
16 ThreadCompletedEvent, ThreadCompletionSubtype, ThreadEvent, ThreadItem, ThreadItemDetails, ThreadStartedEvent,
17 ToolOutcome, TurnBlockedEvent, TurnCompletedEvent, TurnFailedEvent, TurnStartedEvent, Usage, VersionedThreadEvent,
18 tool_outcome_from_status,
19};
20use anyhow::{Context, Result, anyhow};
21use parking_lot::Mutex;
22use serde::Serialize;
23use serde_json::Value;
24use std::io::{self, Write};
25use std::mem::size_of;
26use std::path::Path;
27use std::sync::Arc;
28use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
29use std::sync::mpsc::{self, Receiver, TrySendError};
30use tokio::sync::OnceCell;
31use tokio::task::JoinHandle;
32use tokio::task::spawn_blocking;
33use uuid::Uuid;
34
35use vtcode_memory::event_log::DEFAULT_MAX_EVENTS;
36
37const SESSION_STORE_DRAIN_CAPACITY: usize = 8192;
38const SESSION_STORE_DRAIN_MAX_BYTES: usize = 16 * 1024 * 1024;
39const EXPLANATION_CACHE_MAX_BYTES: usize = 32 * 1024 * 1024;
40
41pub type EventSink = Arc<Mutex<Box<dyn FnMut(&ThreadEvent) + Send>>>;
43
44pub type DecisionEvidenceValidator =
46 Arc<dyn Fn(String, Vec<String>) -> futures::future::BoxFuture<'static, Result<()>> + Send + Sync>;
47
48fn decision_validator(state: Arc<SessionStoreSinkState>) -> DecisionEvidenceValidator {
49 Arc::new(move |task, ids| {
50 let state = Arc::clone(&state);
51 Box::pin(async move {
52 anyhow::ensure!(
53 !state.health.failed.load(Ordering::Acquire),
54 "canonical persistence failed; decision evidence unavailable"
55 );
56 let (reply, receiver) = tokio::sync::oneshot::channel();
57 state
58 .sender
59 .lock()
60 .as_ref()
61 .context("canonical event sink is closed")?
62 .try_send(SessionStoreRequest::ValidateDecision(task, ids, reply))
63 .map_err(|error| anyhow!("canonical event queue unavailable; retry decision recording: {error}"))?;
64 receiver.await.context("decision evidence drain unavailable")?
65 })
66 })
67}
68
69fn matrix_persistence(state: Arc<SessionStoreSinkState>) -> crate::subagents::matrix::MatrixPersistence {
70 let write_state = Arc::clone(&state);
71 crate::subagents::matrix::MatrixPersistence {
72 persist: Arc::new(move |snapshot| {
73 let state = Arc::clone(&write_state);
74 Box::pin(async move {
75 anyhow::ensure!(!state.health.failed.load(Ordering::Acquire), "canonical persistence failed");
76 let event = ThreadEvent::MatrixUpdated(Box::new(snapshot));
77 let queued = prepare_queued_session_event(&state, &event)?;
78 let reserved_bytes = queued.reserved_bytes;
79 let (reply, receiver) = tokio::sync::oneshot::channel();
80 let sent = state
81 .sender
82 .lock()
83 .as_ref()
84 .context("canonical event sink is closed")
85 .and_then(|sender| {
86 sender
87 .try_send(SessionStoreRequest::MatrixPersist(queued, reply))
88 .map_err(|error| anyhow!("canonical matrix queue unavailable: {error}"))
89 });
90 if let Err(error) = sent {
91 release_reserved_bytes(&state, reserved_bytes);
92 state.health.failed.store(true, Ordering::Release);
93 return Err(error);
94 }
95 state.health.accepted_events.fetch_add(1, Ordering::Relaxed);
96 receiver.await.context("matrix persistence barrier unavailable")?
97 })
98 }),
99 load: Arc::new(move || {
100 let state = Arc::clone(&state);
101 Box::pin(async move {
102 anyhow::ensure!(!state.health.failed.load(Ordering::Acquire), "canonical persistence failed");
103 let (reply, receiver) = tokio::sync::oneshot::channel();
104 state
105 .sender
106 .lock()
107 .as_ref()
108 .context("canonical event sink is closed")?
109 .try_send(SessionStoreRequest::MatrixLoad(reply))
110 .map_err(|error| anyhow!("canonical matrix queue unavailable: {error}"))?;
111 receiver.await.context("matrix replay barrier unavailable")?
112 })
113 }),
114 }
115}
116
117#[derive(Debug, Default)]
118struct SessionStoreSinkHealth {
119 accepted_events: AtomicU64,
120 persisted_events: AtomicU64,
121 append_failures: AtomicU64,
122 serialization_failures: AtomicU64,
123 channel_failures: AtomicU64,
124 failed: AtomicBool,
125}
126
127#[cfg(test)]
128#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
129struct SessionStoreSinkHealthSnapshot {
130 accepted_events: u64,
131 persisted_events: u64,
132 append_failures: u64,
133 serialization_failures: u64,
134 channel_failures: u64,
135 failed: bool,
136}
137
138impl SessionStoreSinkHealth {
139 #[cfg(test)]
140 fn snapshot(&self) -> SessionStoreSinkHealthSnapshot {
141 SessionStoreSinkHealthSnapshot {
142 accepted_events: self.accepted_events.load(Ordering::Relaxed),
143 persisted_events: self.persisted_events.load(Ordering::Relaxed),
144 append_failures: self.append_failures.load(Ordering::Relaxed),
145 serialization_failures: self.serialization_failures.load(Ordering::Relaxed),
146 channel_failures: self.channel_failures.load(Ordering::Relaxed),
147 failed: self.failed.load(Ordering::Relaxed),
148 }
149 }
150}
151
152enum SessionStoreRequest {
153 MatrixPersist(QueuedSessionEvent, tokio::sync::oneshot::Sender<Result<()>>),
154 MatrixLoad(tokio::sync::oneshot::Sender<Result<Vec<crate::exec::events::matrix::MatrixSnapshot>>>),
155 Event(QueuedSessionEvent),
156 ValidateDecision(String, Vec<String>, tokio::sync::oneshot::Sender<Result<()>>),
157 Explanation(
158 vtcode_memory::explanation::ExplanationScope,
159 tokio::sync::oneshot::Sender<Result<vtcode_memory::explanation::ExplanationModel>>,
160 ),
161 ExplanationPage(
162 vtcode_memory::explanation::ExplanationScope,
163 usize,
164 tokio::sync::oneshot::Sender<Result<vtcode_memory::explanation::ExplanationPage>>,
165 ),
166 Evidence(
167 vtcode_memory::explanation::EvidenceRef,
168 usize,
169 tokio::sync::oneshot::Sender<Result<vtcode_memory::explanation::EvidencePage>>,
170 ),
171}
172
173struct QueuedSessionEvent {
174 event: ThreadEvent,
175 reserved_bytes: usize,
176}
177
178struct SessionStoreSinkState {
179 sender: Mutex<Option<mpsc::SyncSender<SessionStoreRequest>>>,
180 reserved_bytes: AtomicU64,
181 max_bytes: usize,
182 health: Arc<SessionStoreSinkHealth>,
183}
184
185pub(crate) struct SessionStoreSinkHandle {
191 state: Arc<SessionStoreSinkState>,
192 drain: Option<JoinHandle<()>>,
193}
194
195impl SessionStoreSinkHandle {
196 pub(crate) fn matrix_persistence(&self) -> crate::subagents::matrix::MatrixPersistence {
197 matrix_persistence(Arc::clone(&self.state))
198 }
199 pub(crate) fn decision_validator(&self) -> DecisionEvidenceValidator {
200 decision_validator(Arc::clone(&self.state))
201 }
202 pub(crate) async fn close(mut self) -> Result<()> {
203 self.state.sender.lock().take();
204 if let Some(drain) = self.drain.take() {
205 drain.await.context("session event drain task failed")?;
206 }
207 if self.state.health.failed.load(Ordering::Acquire) {
208 return Err(anyhow!("session event persistence failed"));
209 }
210 Ok(())
211 }
212}
213
214impl Drop for SessionStoreSinkHandle {
215 fn drop(&mut self) {
216 self.state.sender.lock().take();
217 }
218}
219
220pub struct SessionStoreSink {
230 state: Arc<SessionStoreSinkState>,
231 handle: Arc<Mutex<Option<SessionStoreSinkHandle>>>,
232 close_result: Arc<OnceCell<std::result::Result<(), String>>>,
233}
234
235impl SessionStoreSink {
236 pub fn matrix_persistence(&self) -> crate::subagents::matrix::MatrixPersistence {
238 matrix_persistence(Arc::clone(&self.state))
239 }
240 pub fn decision_validator(&self) -> DecisionEvidenceValidator {
242 decision_validator(Arc::clone(&self.state))
243 }
244 pub async fn open(workspace: &Path, session_id: &str) -> Result<Self> {
246 let (state, handle) = open_session_store_sink(workspace, session_id, SESSION_STORE_DRAIN_CAPACITY).await?;
247 Ok(Self {
248 state,
249 handle: Arc::new(Mutex::new(Some(handle))),
250 close_result: Arc::new(OnceCell::new()),
251 })
252 }
253
254 pub fn emit(&self, event: &ThreadEvent) -> Result<()> {
256 enqueue_session_event(&self.state, event)
257 }
258
259 pub fn event_sink(&self) -> EventSink {
261 event_sink_for_state(Arc::clone(&self.state))
262 }
263
264 pub async fn explanation(
266 &self,
267 scope: vtcode_memory::explanation::ExplanationScope,
268 ) -> Result<vtcode_memory::explanation::ExplanationModel> {
269 let (sender, receiver) = tokio::sync::oneshot::channel();
270 self.enqueue_query(SessionStoreRequest::Explanation(scope, sender))?;
271 receiver.await.context("explanation snapshot drain unavailable")?
272 }
273
274 pub async fn explanation_page(
276 &self,
277 scope: vtcode_memory::explanation::ExplanationScope,
278 offset: usize,
279 ) -> Result<vtcode_memory::explanation::ExplanationPage> {
280 let (sender, receiver) = tokio::sync::oneshot::channel();
281 self.enqueue_query(SessionStoreRequest::ExplanationPage(scope, offset, sender))?;
282 receiver.await.context("explanation page drain unavailable")?
283 }
284
285 pub async fn evidence(
287 &self,
288 reference: vtcode_memory::explanation::EvidenceRef,
289 offset: usize,
290 ) -> Result<vtcode_memory::explanation::EvidencePage> {
291 let (sender, receiver) = tokio::sync::oneshot::channel();
292 self.enqueue_query(SessionStoreRequest::Evidence(reference, offset, sender))?;
293 receiver.await.context("explanation evidence drain unavailable")?
294 }
295
296 fn enqueue_query(&self, request: SessionStoreRequest) -> Result<()> {
297 if self.state.health.failed.load(Ordering::Acquire) {
298 return Err(anyhow!("canonical persistence failed; explanation unavailable"));
299 }
300 self.state
301 .sender
302 .lock()
303 .as_ref()
304 .context("canonical event sink is closed")?
305 .try_send(request)
306 .map_err(|error| anyhow!("canonical event queue unavailable; retry explanation query: {error}"))
307 }
308
309 pub async fn close(&self) -> Result<()> {
311 let result = self
312 .close_result
313 .get_or_init(|| async {
314 let handle = self.handle.lock().take();
315 match handle {
316 Some(handle) => handle.close().await.map_err(|error| error.to_string()),
317 None => Ok(()),
318 }
319 })
320 .await;
321 match result {
322 Ok(()) => Ok(()),
323 Err(error) => Err(anyhow!(error.clone())),
324 }
325 }
326}
327
328impl Drop for SessionStoreSink {
329 fn drop(&mut self) {
330 self.state.sender.lock().take();
331 }
332}
333
334#[doc(hidden)]
335pub fn event_sink<F>(callback: F) -> EventSink
336where
337 F: FnMut(&ThreadEvent) + Send + 'static,
338{
339 Arc::new(Mutex::new(Box::new(callback)))
340}
341
342pub async fn session_store_sink(workspace: &Path, session_id: &str) -> Result<SessionStoreSink> {
351 SessionStoreSink::open(workspace, session_id).await
352}
353
354pub(crate) async fn session_store_sink_with_handle(
355 workspace: &Path,
356 session_id: &str,
357) -> Result<(EventSink, SessionStoreSinkHandle)> {
358 let (state, handle) = open_session_store_sink(workspace, session_id, SESSION_STORE_DRAIN_CAPACITY).await?;
359 Ok((event_sink_for_state(state), handle))
360}
361
362async fn open_session_store_sink(
363 workspace: &Path,
364 session_id: &str,
365 capacity: usize,
366) -> Result<(Arc<SessionStoreSinkState>, SessionStoreSinkHandle)> {
367 open_session_store_sink_with_limits(workspace, session_id, capacity, SESSION_STORE_DRAIN_MAX_BYTES).await
368}
369
370async fn open_session_store_sink_with_limits(
371 workspace: &Path,
372 session_id: &str,
373 capacity: usize,
374 max_bytes: usize,
375) -> Result<(Arc<SessionStoreSinkState>, SessionStoreSinkHandle)> {
376 let workspace = workspace.to_path_buf();
377 let session_id_owned = session_id.to_string();
378 let log = spawn_blocking(move || vtcode_memory::open(&workspace, &session_id_owned, DEFAULT_MAX_EVENTS))
379 .await
380 .context("canonical session store open task failed")??;
381
382 let queue_capacity = capacity.max(1);
383 let (sender, receiver) = mpsc::sync_channel::<SessionStoreRequest>(queue_capacity);
384 let health = Arc::new(SessionStoreSinkHealth::default());
385 let state = Arc::new(SessionStoreSinkState {
386 sender: Mutex::new(Some(sender)),
387 reserved_bytes: AtomicU64::new(0),
388 max_bytes: max_bytes.max(1),
389 health: Arc::clone(&health),
390 });
391 let drain_state = Arc::clone(&state);
392 let drain_session_id = session_id.to_string();
393 let drain = spawn_blocking(move || drain_session_events(receiver, log, drain_session_id, drain_state));
394
395 Ok((state.clone(), SessionStoreSinkHandle { state, drain: Some(drain) }))
396}
397
398#[cfg(test)]
399async fn session_store_sink_with_capacity_handle(
400 workspace: &Path,
401 session_id: &str,
402 capacity: usize,
403) -> Result<(EventSink, SessionStoreSinkHandle)> {
404 let (state, handle) = open_session_store_sink(workspace, session_id, capacity).await?;
405 Ok((event_sink_for_state(state), handle))
406}
407
408fn event_sink_for_state(state: Arc<SessionStoreSinkState>) -> EventSink {
409 event_sink(move |event: &ThreadEvent| {
410 if let Err(error) = enqueue_session_event(&state, event) {
411 tracing::error!(error = %error, "canonical session event was not accepted");
412 }
413 })
414}
415
416struct CachedExplanation {
417 scope: vtcode_memory::explanation::ExplanationScope,
418 model: Arc<vtcode_memory::explanation::ExplanationModel>,
419}
420
421fn cached_explanation(
422 log: &vtcode_memory::SessionEventLog,
423 scope: vtcode_memory::explanation::ExplanationScope,
424 explanations: &mut Vec<CachedExplanation>,
425) -> Result<Arc<vtcode_memory::explanation::ExplanationModel>> {
426 if let Some(cached) = explanations.iter().find(|cached| cached.scope == scope) {
427 return Ok(Arc::clone(&cached.model));
428 }
429 let model = Arc::new(vtcode_memory::explanation::query_explanation(log, scope)?);
430 let mut counter = LimitedByteCounter { bytes: 0, limit: EXPLANATION_CACHE_MAX_BYTES };
432 if serde_json::to_writer(&mut counter, model.as_ref()).is_ok() {
433 explanations.push(CachedExplanation { scope, model: Arc::clone(&model) });
434 }
435 Ok(model)
436}
437
438fn matrix_command_completed_event(
439 snapshot: &crate::exec::events::matrix::MatrixSnapshot,
440 task: &crate::exec::events::matrix::MatrixTaskState,
441 evidence: &crate::exec::events::matrix::MatrixCommandEvidence,
442 session_id: &str,
443) -> ThreadEvent {
444 ThreadEvent::ItemCompleted(ItemCompletedEvent {
445 item: ThreadItem {
446 id: evidence.event_id.clone(),
447 context: Some(Box::new(crate::exec::events::ItemContext {
448 task_id: format!("{}/{}", snapshot.spec.id, task.id),
449 turn_id: evidence.attempt_id.clone(),
450 actor_id: evidence.worker_id.clone(),
451 parent_actor_id: Some(session_id.to_owned()),
452 timestamp: chrono::Utc::now().to_rfc3339(),
453 activity: Some(crate::exec::events::CommandActivity::Verification),
454 })),
455 details: ThreadItemDetails::CommandExecution(Box::new(CommandExecutionItem {
456 command: evidence.command.clone(),
457 arguments: Some(serde_json::json!({
458 "matrix_id":snapshot.spec.id, "task_id":task.id,
459 "attempt_id":evidence.attempt_id, "worker_id":evidence.worker_id,
460 "generation":evidence.generation, "cancelled":evidence.cancelled,
461 })),
462 aggregated_output: String::new(),
463 exit_code: evidence.exit_code,
464 status: if evidence.exit_code == Some(0) && !evidence.cancelled {
465 CommandExecutionStatus::Completed
466 } else {
467 CommandExecutionStatus::Failed
468 },
469 })),
470 },
471 })
472}
473
474fn restore_matrix_evidence_ids(
475 log: &vtcode_memory::SessionEventLog,
476 ids: &mut std::collections::HashSet<String>,
477) -> Result<()> {
478 log.visit_snapshot(|_, bytes| {
479 if let Ok(versioned) = serde_json::from_slice::<VersionedThreadEvent>(bytes) {
480 match versioned.into_event() {
481 ThreadEvent::ItemCompleted(event)
482 if matches!(event.item.details, ThreadItemDetails::CommandExecution(_)) =>
483 {
484 ids.insert(event.item.id);
485 }
486 ThreadEvent::MatrixUpdated(snapshot) => {
487 ids.extend(
490 snapshot
491 .tasks
492 .into_iter()
493 .flat_map(|task| task.attempts)
494 .flat_map(|attempt| attempt.evidence)
495 .map(|evidence| evidence.event_id),
496 );
497 }
498 _ => {}
499 }
500 }
501 })
502 .context("restore canonical matrix command identities")?;
503 Ok(())
504}
505
506fn drain_session_events(
507 rx: Receiver<SessionStoreRequest>,
508 log: vtcode_memory::SessionEventLog,
509 session_id: String,
510 state: Arc<SessionStoreSinkState>,
511) {
512 let mut explanations = Vec::new();
513 let mut matrix_evidence_ids = std::collections::HashSet::new();
514 let mut matrix_evidence_restored = false;
515 while let Ok(request) = rx.recv() {
516 let queued = match request {
517 SessionStoreRequest::MatrixPersist(queued, reply) => {
518 let result = (|| -> Result<()> {
519 if let ThreadEvent::MatrixUpdated(snapshot) = &queued.event {
520 if !matrix_evidence_restored {
521 restore_matrix_evidence_ids(&log, &mut matrix_evidence_ids)?;
522 matrix_evidence_restored = true;
523 }
524 for task in &snapshot.tasks {
525 for attempt in &task.attempts {
526 for evidence in &attempt.evidence {
527 if matrix_evidence_ids.insert(evidence.event_id.clone()) {
528 log.append(&matrix_command_completed_event(
529 snapshot,
530 task,
531 evidence,
532 &session_id,
533 ))?;
534 }
535 }
536 }
537 }
538 }
539 log.append(&queued.event)?;
540 log.flush()?;
541 Ok(())
542 })();
543 release_reserved_bytes(&state, queued.reserved_bytes);
544 let failed = result.is_err();
545 if failed {
546 state.health.append_failures.fetch_add(1, Ordering::Relaxed);
547 state.health.failed.store(true, Ordering::Release);
548 } else {
549 state.health.persisted_events.fetch_add(1, Ordering::Relaxed);
550 }
551 explanations.clear();
552 let _ = reply.send(result);
553 if failed {
554 break;
555 }
556 continue;
557 }
558 SessionStoreRequest::MatrixLoad(reply) => {
559 let result = (|| -> Result<_> {
560 let mut latest = std::collections::BTreeMap::<String, vtcode_memory::matrix::MatrixState>::new();
561 let mut replay_error = None;
562 log.visit_snapshot(|_, bytes| {
563 if replay_error.is_some() {
564 return;
565 }
566 if let Ok(versioned) = serde_json::from_slice::<VersionedThreadEvent>(bytes)
567 && let ThreadEvent::MatrixUpdated(snapshot) = versioned.into_event()
568 {
569 let previous = latest.remove(&snapshot.spec.id);
570 let events = previous
571 .into_iter()
572 .map(|state| state.event())
573 .chain([ThreadEvent::MatrixUpdated(snapshot)]);
574 match vtcode_memory::matrix::replay(events) {
575 Ok(replayed) => latest.extend(replayed),
576 Err(error) => replay_error = Some(error),
577 }
578 }
579 })?;
580 if let Some(error) = replay_error {
581 return Err(error.into());
582 }
583 Ok(latest.into_values().map(|state| state.snapshot().clone()).collect())
584 })();
585 let _ = reply.send(result);
586 continue;
587 }
588 SessionStoreRequest::Event(event) => event,
589 SessionStoreRequest::ValidateDecision(task_id, ids, reply) => {
590 let result = (|| -> Result<()> {
591 let mut found = std::collections::HashSet::new();
592 log.visit_snapshot(|_, bytes| {
593 if let Ok(versioned) = serde_json::from_slice::<VersionedThreadEvent>(bytes) {
594 let item = match versioned.into_event() {
595 ThreadEvent::ItemStarted(e) => Some(e.item),
596 ThreadEvent::ItemUpdated(e) => Some(e.item),
597 ThreadEvent::ItemCompleted(e) => Some(e.item),
598 _ => None,
599 };
600 if let Some(item) = item
601 && item.context.as_ref().is_some_and(|c| c.task_id == task_id)
602 && ids.contains(&item.id)
603 {
604 found.insert(item.id);
605 }
606 }
607 })?;
608 anyhow::ensure!(
609 ids.iter().all(|id| found.contains(id)),
610 "decision evidence is unavailable or does not belong to the current task"
611 );
612 Ok(())
613 })();
614 let _ = reply.send(result);
615 continue;
616 }
617 SessionStoreRequest::Explanation(scope, reply) => {
618 let result = cached_explanation(&log, scope, &mut explanations).map(|model| model.as_ref().clone());
619 let _ = reply.send(result);
620 continue;
621 }
622 SessionStoreRequest::ExplanationPage(scope, offset, reply) => {
623 let result = cached_explanation(&log, scope, &mut explanations)
624 .map(|model| vtcode_memory::explanation::page_explanation(&model, offset));
625 let _ = reply.send(result);
626 continue;
627 }
628 SessionStoreRequest::Evidence(reference, offset, reply) => {
629 let result = vtcode_memory::explanation::query_evidence(&log, &reference, offset, 32 * 1024)
630 .map_err(anyhow::Error::from);
631 let _ = reply.send(result);
632 continue;
633 }
634 };
635 explanations.clear();
636 let reserved_bytes = queued.reserved_bytes;
637 match log.append(&queued.event) {
638 Ok(()) => {
639 state.health.persisted_events.fetch_add(1, Ordering::Relaxed);
640 }
641 Err(err) => {
642 state.health.append_failures.fetch_add(1, Ordering::Relaxed);
643 state.health.failed.store(true, Ordering::Release);
644 tracing::error!(
645 session_id = %session_id,
646 error = %err,
647 "failed to persist session event; stopping authoritative drain"
648 );
649 release_reserved_bytes(&state, reserved_bytes);
650 while let Ok(request) = rx.try_recv() {
651 if let SessionStoreRequest::Event(queued) | SessionStoreRequest::MatrixPersist(queued, _) = request
652 {
653 release_reserved_bytes(&state, queued.reserved_bytes);
654 }
655 }
656 break;
657 }
658 }
659 release_reserved_bytes(&state, reserved_bytes);
660 }
661
662 if let Err(err) = log.flush() {
663 state.health.append_failures.fetch_add(1, Ordering::Relaxed);
664 state.health.failed.store(true, Ordering::Release);
665 tracing::error!(
666 session_id = %session_id,
667 error = %err,
668 "failed to flush session event log during drain shutdown"
669 );
670 }
671}
672
673fn prepare_queued_session_event(state: &SessionStoreSinkState, event: &ThreadEvent) -> Result<QueuedSessionEvent> {
674 if state.health.failed.load(Ordering::Acquire) {
675 state.health.channel_failures.fetch_add(1, Ordering::Relaxed);
676 return Err(anyhow!("canonical session event sink has failed"));
677 }
678
679 let serialized_bytes = match serialized_event_size(event, state.max_bytes) {
680 Ok(bytes) => bytes,
681 Err(error) => {
682 state.health.serialization_failures.fetch_add(1, Ordering::Relaxed);
683 state.health.failed.store(true, Ordering::Release);
684 return Err(error.context("failed to serialize canonical session event for bounded handoff"));
685 }
686 };
687 let reserved_bytes = serialized_bytes.saturating_mul(2).saturating_add(size_of::<ThreadEvent>());
688 if reserved_bytes > state.max_bytes {
689 state.health.serialization_failures.fetch_add(1, Ordering::Relaxed);
690 state.health.failed.store(true, Ordering::Release);
691 return Err(anyhow!(
692 "canonical session event exceeds the bounded persistence budget ({reserved_bytes} > {})",
693 state.max_bytes
694 ));
695 }
696 if !reserve_bytes(&state.reserved_bytes, reserved_bytes, state.max_bytes) {
697 state.health.failed.store(true, Ordering::Release);
698 state.health.channel_failures.fetch_add(1, Ordering::Relaxed);
699 return Err(anyhow!("canonical session event queue reached its bounded byte capacity"));
700 }
701
702 Ok(QueuedSessionEvent { event: event.clone(), reserved_bytes })
703}
704
705fn enqueue_session_event(state: &SessionStoreSinkState, event: &ThreadEvent) -> Result<()> {
706 let queued = prepare_queued_session_event(state, event)?;
707 let reserved_bytes = queued.reserved_bytes;
708 let sender = state.sender.lock();
709 let Some(sender) = sender.as_ref() else {
710 release_reserved_bytes(state, reserved_bytes);
711 state.health.failed.store(true, Ordering::Release);
712 state.health.channel_failures.fetch_add(1, Ordering::Relaxed);
713 return Err(anyhow!("canonical session event sink is closed"));
714 };
715 match sender.try_send(SessionStoreRequest::Event(queued)) {
716 Ok(()) => {
717 state.health.accepted_events.fetch_add(1, Ordering::Relaxed);
718 Ok(())
719 }
720 Err(TrySendError::Full(request) | TrySendError::Disconnected(request)) => {
721 if let SessionStoreRequest::Event(queued) | SessionStoreRequest::MatrixPersist(queued, _) = request {
722 release_reserved_bytes(state, queued.reserved_bytes);
723 }
724 state.health.failed.store(true, Ordering::Release);
725 state.health.channel_failures.fetch_add(1, Ordering::Relaxed);
726 Err(anyhow!("canonical session event queue could not accept the event"))
727 }
728 }
729}
730
731fn reserve_bytes(reserved_bytes: &AtomicU64, bytes: usize, max_bytes: usize) -> bool {
732 let bytes = u64::try_from(bytes).unwrap_or(u64::MAX);
733 let max_bytes = u64::try_from(max_bytes).unwrap_or(u64::MAX);
734 reserved_bytes
735 .fetch_update(Ordering::AcqRel, Ordering::Acquire, |current| {
736 current.checked_add(bytes).filter(|next| *next <= max_bytes)
737 })
738 .is_ok()
739}
740
741fn release_reserved_bytes(state: &SessionStoreSinkState, bytes: usize) {
742 let bytes = u64::try_from(bytes).unwrap_or(u64::MAX);
743 state.reserved_bytes.fetch_sub(bytes, Ordering::AcqRel);
744}
745
746#[derive(Serialize)]
747struct BorrowedVersionedThreadEvent<'a> {
748 schema_version: &'static str,
749 event: &'a ThreadEvent,
750}
751
752struct LimitedByteCounter {
753 bytes: usize,
754 limit: usize,
755}
756
757impl Write for LimitedByteCounter {
758 fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
759 let next = self
760 .bytes
761 .checked_add(bytes.len())
762 .ok_or_else(|| io::Error::other("canonical event size overflow"))?;
763 if next > self.limit {
764 return Err(io::Error::other("canonical event exceeds bounded size"));
765 }
766 self.bytes = next;
767 Ok(bytes.len())
768 }
769
770 fn flush(&mut self) -> io::Result<()> {
771 Ok(())
772 }
773}
774
775fn serialized_event_size(event: &ThreadEvent, limit: usize) -> Result<usize> {
776 let mut counter = LimitedByteCounter { bytes: 0, limit };
777 serde_json::to_writer(&mut counter, &BorrowedVersionedThreadEvent { schema_version: EVENT_SCHEMA_VERSION, event })
778 .context("canonical event serialization failed")?;
779 counter.bytes.checked_add(1).context("canonical event size overflow")
780}
781
782pub fn combine_event_sinks(a: Option<EventSink>, b: Option<EventSink>) -> Option<EventSink> {
784 match (a, b) {
785 (None, None) => None,
786 (Some(s), None) | (None, Some(s)) => Some(s),
787 (Some(a), Some(b)) => Some(event_sink(move |e: &ThreadEvent| {
788 a.lock()(e);
789 b.lock()(e);
790 })),
791 }
792}
793
794#[derive(Debug, Clone)]
795pub struct ActiveCommandHandle {
796 id: String,
797 command: String,
798}
799
800#[derive(Debug, Clone)]
801pub struct ActiveToolHandle {
802 id: String,
803 tool_name: String,
804 arguments: Option<Value>,
805 tool_call_id: Option<String>,
806}
807
808impl ActiveToolHandle {
809 #[must_use]
810 pub fn item_id(&self) -> &str {
811 &self.id
812 }
813}
814
815#[derive(Default)]
817pub struct ExecEventRecorder {
818 context: ExecutionContextTracker,
819 thread_id: String,
820 events: Vec<ThreadEvent>,
821 event_sink: Option<EventSink>,
822 thread_handle: Option<ThreadRuntimeHandle>,
823 active_submission_id: Option<SubmissionId>,
824 active_turn_id: Option<String>,
825 lifecycle: SharedLifecycleEmitter,
826}
827
828impl ExecEventRecorder {
829 pub fn new(
830 thread_id: impl Into<String>,
831 event_sink: Option<EventSink>,
832 thread_handle: Option<ThreadRuntimeHandle>,
833 ) -> Self {
834 let thread_id = thread_id.into();
835 let mut recorder = Self {
836 context: ExecutionContextTracker::default(),
837 thread_id: thread_id.clone(),
838 events: Vec::new(),
839 event_sink,
840 thread_handle,
841 active_submission_id: None,
842 active_turn_id: None,
843 lifecycle: SharedLifecycleEmitter::default(),
844 };
845 recorder.record_with_context(None, None, ThreadEvent::ThreadStarted(ThreadStartedEvent { thread_id }));
846 recorder
847 }
848
849 fn record(&mut self, event: ThreadEvent) {
850 self.record_with_context(self.active_submission_id.clone(), self.active_turn_id.clone(), event);
851 }
852
853 fn record_with_context(
854 &mut self,
855 submission_id: Option<SubmissionId>,
856 turn_id: Option<String>,
857 mut event: ThreadEvent,
858 ) {
859 self.context.annotate(&mut event);
860 let decision = decision_completed_event(&event);
861 let exec_completion = self.context.background_exec_output_event(&mut event);
862 if let Some(sink) = &self.event_sink {
863 let mut callback = sink.lock();
864 callback(&event);
865 }
866 if let Some(handle) = &self.thread_handle {
867 handle.record_event(submission_id, turn_id, event.clone());
868 }
869 self.events.push(event);
870 if let Some(decision) = decision {
871 self.record(decision);
872 }
873 if let Some(completion) = exec_completion {
874 self.record(completion);
875 }
876 }
877
878 pub fn record_exec_session_output(&mut self, call_item_id: &str, tool_name: &str, args: &Value, output: &Value) {
880 if let Some(event) = self.context.exec_session_output_event(call_item_id, tool_name, args, output) {
881 self.record(event);
882 }
883 if let Some(event) = self.context.delegation_status_event(call_item_id, tool_name, output) {
884 self.record(event);
885 }
886 }
887
888 pub fn record_thread_event(&mut self, event: ThreadEvent) {
889 self.record(event);
890 }
891
892 pub fn record_thread_events<I>(&mut self, events: I)
893 where
894 I: IntoIterator<Item = ThreadEvent>,
895 {
896 for event in events {
897 self.record(event);
898 }
899 }
900
901 fn record_pending_lifecycle_events(&mut self) {
902 for event in self.lifecycle.drain_events() {
903 self.record(event);
904 }
905 }
906
907 fn next_item_id(&mut self) -> String {
908 self.lifecycle.next_item_id()
909 }
910
911 pub fn turn_started(&mut self) {
912 let turn_id = format!("turn-{}", Uuid::new_v4());
913 self.context
914 .begin(&self.thread_id, &turn_id, "", crate::exec::events::InputOrigin::Continuation);
915 self.active_turn_id = Some(turn_id);
916 self.begin_turn_submission();
917 }
918
919 fn begin_turn_submission(&mut self) {
920 if let Some(handle) = &self.thread_handle {
921 match handle.begin_turn() {
922 Ok(submission_id) => self.active_submission_id = Some(submission_id),
923 Err(err) => {
924 tracing::warn!(error = %err, "failed to begin turn; submission id unavailable");
930 self.active_submission_id = None;
931 }
932 }
933 }
934 self.record(ThreadEvent::TurnStarted(TurnStartedEvent::default()));
935 }
936
937 pub fn task_started(&mut self, goal: &str) -> String {
939 let context = self.context.begin(
940 &self.thread_id,
941 &format!("turn-{}", Uuid::new_v4()),
942 goal,
943 crate::exec::events::InputOrigin::User,
944 );
945 self.active_turn_id = Some(context.turn_id);
946 self.begin_turn_submission();
947 context.task_id
948 }
949
950 pub fn turn_completed(&mut self, usage: Usage) {
951 self.record(ThreadEvent::TurnCompleted(TurnCompletedEvent {
952 completed_at: None,
953 usage,
954 in_progress_exec_sessions: Vec::new(),
955 }));
956 self.finish_turn();
957 }
958
959 pub fn turn_failed(&mut self, message: &str, usage: Option<Usage>) {
960 self.record(ThreadEvent::TurnFailed(TurnFailedEvent {
961 completed_at: None,
962 message: message.to_string(),
963 usage,
964 }));
965 self.finish_turn();
966 }
967
968 pub fn turn_blocked(&mut self, event: TurnBlockedEvent) {
969 self.record(ThreadEvent::TurnBlocked(Box::new(event)));
970 }
971
972 pub fn thread_completed(
973 &mut self,
974 session_id: &str,
975 subtype: ThreadCompletionSubtype,
976 outcome_code: &str,
977 result: Option<&str>,
978 stop_reason: Option<&str>,
979 usage: Usage,
980 total_cost_usd: Option<serde_json::Number>,
981 num_turns: usize,
982 ) {
983 self.record(ThreadEvent::ThreadCompleted(Box::new(ThreadCompletedEvent {
984 completed_at: None,
985 thread_id: self.thread_id.clone(),
986 session_id: session_id.to_string(),
987 subtype,
988 outcome_code: outcome_code.to_string(),
989 result: result.map(str::to_string),
990 stop_reason: stop_reason.map(str::to_string),
991 usage,
992 total_cost_usd,
993 num_turns,
994 })));
995 }
996
997 pub fn thread_failed(&mut self, session_id: &str, message: &str, num_turns: usize) {
999 self.turn_failed(message, None);
1000 self.thread_completed(
1001 session_id,
1002 ThreadCompletionSubtype::ErrorDuringExecution,
1003 "error",
1004 None,
1005 Some(message),
1006 Usage::default(),
1007 None,
1008 num_turns,
1009 );
1010 }
1011
1012 pub fn compact_boundary(
1013 &mut self,
1014 trigger: CompactionTrigger,
1015 mode: CompactionMode,
1016 original_message_count: usize,
1017 compacted_message_count: usize,
1018 history_artifact_path: Option<&str>,
1019 previous_envelope: Option<&crate::core::agent::request_envelope::SessionRequestEnvelope>,
1020 new_envelope: Option<&crate::core::agent::request_envelope::SessionRequestEnvelope>,
1021 ) {
1022 let previous_segment_id = previous_envelope.map(|envelope| envelope.segment_id().to_owned());
1023 let new_segment_id = new_envelope.map(|envelope| envelope.segment_id().to_owned());
1024 let previous_prefix_hash = previous_envelope.map(|envelope| format!("{:016x}", envelope.prefix_hash()));
1025 let new_prefix_hash = new_envelope.map(|envelope| format!("{:016x}", envelope.prefix_hash()));
1026 let previous_catalog_hash = previous_envelope
1027 .and_then(crate::core::agent::request_envelope::SessionRequestEnvelope::catalog_hash)
1028 .map(|hash| format!("{hash:016x}"));
1029 let new_catalog_hash = new_envelope
1030 .and_then(crate::core::agent::request_envelope::SessionRequestEnvelope::catalog_hash)
1031 .map(|hash| format!("{hash:016x}"));
1032 self.record(ThreadEvent::ThreadCompactBoundary(Box::new(ThreadCompactBoundaryEvent {
1033 thread_id: self.thread_id.clone(),
1034 trigger,
1035 mode,
1036 original_message_count,
1037 compacted_message_count,
1038 history_artifact_path: history_artifact_path.map(str::to_string),
1039 previous_segment_id,
1040 new_segment_id,
1041 previous_prefix_hash,
1042 new_prefix_hash,
1043 previous_catalog_hash,
1044 new_catalog_hash,
1045 })));
1046 }
1047
1048 fn finish_turn(&mut self) {
1049 if let Some(handle) = &self.thread_handle {
1050 handle.finish_turn();
1051 }
1052 self.active_submission_id = None;
1053 self.active_turn_id = None;
1054 }
1055
1056 pub fn agent_message(&mut self, text: &str) {
1057 self.lifecycle.emit_completed_agent_message(text);
1058 self.record_pending_lifecycle_events();
1059 }
1060
1061 pub fn agent_message_stream_update(&mut self, text: &str) -> bool {
1062 if text.trim().is_empty() || !self.lifecycle.replace_assistant_text(text) {
1063 return false;
1064 }
1065 let emitted = self.lifecycle.emit_assistant_snapshot(None);
1066 self.record_pending_lifecycle_events();
1067 emitted
1068 }
1069
1070 pub fn agent_message_stream_complete(&mut self) {
1071 let _ = self.lifecycle.complete_assistant_stream();
1072 self.record_pending_lifecycle_events();
1073 }
1074
1075 pub fn reasoning(&mut self, text: &str) {
1076 self.lifecycle.emit_completed_reasoning(text);
1077 self.record_pending_lifecycle_events();
1078 }
1079
1080 pub fn set_reasoning_stage(&mut self, stage: &str) {
1081 if !self.lifecycle.set_reasoning_stage(Some(stage.to_string())) {
1082 return;
1083 }
1084 let _ = self.lifecycle.emit_reasoning_stage_update();
1085 self.record_pending_lifecycle_events();
1086 }
1087
1088 pub fn reasoning_stream_update(&mut self, text: &str) -> bool {
1089 if text.trim().is_empty() || !self.lifecycle.replace_reasoning_text(text) {
1090 return false;
1091 }
1092 let emitted = self.lifecycle.emit_reasoning_snapshot(None);
1093 self.record_pending_lifecycle_events();
1094 emitted
1095 }
1096
1097 pub fn reasoning_stream_complete(&mut self) {
1098 let _ = self.lifecycle.complete_reasoning_stream();
1099 self.record_pending_lifecycle_events();
1100 }
1101
1102 pub fn tool_started(
1103 &mut self,
1104 tool_name: &str,
1105 arguments: Option<&Value>,
1106 tool_call_id: Option<&str>,
1107 ) -> ActiveToolHandle {
1108 let handle = ActiveToolHandle {
1109 id: self.next_item_id(),
1110 tool_name: tool_name.to_string(),
1111 arguments: arguments.cloned(),
1112 tool_call_id: tool_call_id.map(str::to_string),
1113 };
1114 self.record(tool_started_event(
1115 handle.id.clone(),
1116 &handle.tool_name,
1117 handle.arguments.as_ref(),
1118 handle.tool_call_id.as_deref(),
1119 ));
1120 handle
1121 }
1122
1123 pub fn tool_finished(
1124 &mut self,
1125 handle: &ActiveToolHandle,
1126 status: crate::exec::events::ToolCallStatus,
1127 exit_code: Option<i32>,
1128 aggregated_output: &str,
1129 spool_path: Option<&str>,
1130 ) {
1131 let outcome = tool_outcome_from_status(&status);
1132 self.record(tool_invocation_completed_event(
1133 handle.id.clone(),
1134 &handle.tool_name,
1135 handle.arguments.as_ref(),
1136 handle.tool_call_id.as_deref(),
1137 status.clone(),
1138 outcome,
1139 ));
1140 self.record(tool_output_completed_event(
1141 handle.id.clone(),
1142 handle.tool_call_id.as_deref(),
1143 status,
1144 exit_code,
1145 spool_path,
1146 aggregated_output,
1147 ));
1148 }
1149
1150 pub fn tool_output_started(&mut self, call_item_id: &str, tool_call_id: Option<&str>) {
1151 self.record(tool_output_started_event(call_item_id.to_string(), tool_call_id));
1152 }
1153
1154 pub fn tool_output_updated(&mut self, call_item_id: &str, tool_call_id: Option<&str>, output: &str) {
1155 self.record(tool_output_updated_event(call_item_id.to_string(), tool_call_id, output));
1156 }
1157
1158 pub fn tool_output_finished(
1159 &mut self,
1160 call_item_id: &str,
1161 tool_call_id: Option<&str>,
1162 status: crate::exec::events::ToolCallStatus,
1163 exit_code: Option<i32>,
1164 aggregated_output: &str,
1165 spool_path: Option<&str>,
1166 ) {
1167 self.record(tool_output_completed_event(
1168 call_item_id.to_string(),
1169 tool_call_id,
1170 status,
1171 exit_code,
1172 spool_path,
1173 aggregated_output,
1174 ));
1175 }
1176
1177 pub fn tool_rejected(
1178 &mut self,
1179 tool_name: &str,
1180 arguments: Option<&Value>,
1181 tool_call_id: Option<&str>,
1182 detail: &str,
1183 ) {
1184 let handle = self.tool_started(tool_name, arguments, tool_call_id);
1185 let call_item_id = handle.id.clone();
1186 self.record(tool_invocation_completed_event(
1187 call_item_id.clone(),
1188 tool_name,
1189 arguments,
1190 tool_call_id,
1191 crate::exec::events::ToolCallStatus::Failed,
1192 ToolOutcome::HookDenied,
1193 ));
1194 self.record(tool_output_started_event(call_item_id.clone(), tool_call_id));
1195 self.record(tool_output_completed_event(
1196 call_item_id,
1197 tool_call_id,
1198 crate::exec::events::ToolCallStatus::Failed,
1199 None,
1200 None,
1201 detail,
1202 ));
1203 let error_item_id = self.next_item_id();
1204 self.record(error_item_completed_event(error_item_id, detail.to_string()));
1205 }
1206
1207 pub fn permission_requested(&mut self, tool_name: &str) {
1208 self.record(ThreadEvent::PermissionRequested(crate::exec::events::PermissionRequestedEvent {
1209 tool_name: tool_name.to_string(),
1210 }));
1211 }
1212
1213 pub fn permission_resolved(
1214 &mut self,
1215 tool_name: &str,
1216 decision: crate::exec::events::PermissionDecision,
1217 wait_ms: u64,
1218 ) {
1219 self.record(ThreadEvent::PermissionResolved(crate::exec::events::PermissionResolvedEvent {
1220 tool_name: tool_name.to_string(),
1221 decision,
1222 wait_ms,
1223 }));
1224 }
1225
1226 pub fn command_started(&mut self, command: &str) -> ActiveCommandHandle {
1227 let id = self.next_item_id();
1228 let item = ThreadItem {
1229 context: None,
1230 id: id.clone(),
1231 details: ThreadItemDetails::CommandExecution(Box::new(CommandExecutionItem {
1232 command: command.to_string(),
1233 arguments: None,
1234 aggregated_output: String::new(),
1235 exit_code: None,
1236 status: CommandExecutionStatus::InProgress,
1237 })),
1238 };
1239 self.record(ThreadEvent::ItemStarted(ItemStartedEvent { item }));
1240 ActiveCommandHandle { id, command: command.to_string() }
1241 }
1242
1243 pub fn command_finished(
1244 &mut self,
1245 handle: &ActiveCommandHandle,
1246 status: CommandExecutionStatus,
1247 exit_code: Option<i32>,
1248 aggregated_output: &str,
1249 ) {
1250 let item = ThreadItem {
1251 context: None,
1252 id: handle.id.clone(),
1253 details: ThreadItemDetails::CommandExecution(Box::new(CommandExecutionItem {
1254 command: handle.command.clone(),
1255 arguments: None,
1256 aggregated_output: aggregated_output.to_string(),
1257 exit_code,
1258 status,
1259 })),
1260 };
1261 self.record(ThreadEvent::ItemCompleted(ItemCompletedEvent { item }));
1262 }
1263
1264 pub fn warning(&mut self, message: &str) {
1265 let item = ThreadItem {
1266 context: None,
1267 id: self.next_item_id(),
1268 details: ThreadItemDetails::Error(ErrorItem { message: message.to_string() }),
1269 };
1270 self.record(ThreadEvent::ItemCompleted(ItemCompletedEvent { item }));
1271 }
1272
1273 pub fn harness_event(
1274 &mut self,
1275 event: HarnessEventKind,
1276 message: Option<String>,
1277 command: Option<String>,
1278 path: Option<String>,
1279 exit_code: Option<i32>,
1280 attempt: Option<u32>,
1281 error_category: Option<String>,
1282 ) {
1283 let item = ThreadItem {
1284 context: None,
1285 id: self.next_item_id(),
1286 details: ThreadItemDetails::Harness(Box::new(HarnessEventItem {
1287 event,
1288 message,
1289 command,
1290 path,
1291 exit_code,
1292 attempt,
1293 error_category,
1294 duration_ms: None,
1295 task_id: None,
1296 session_id: None,
1297 exec_session_id: None,
1298 status: None,
1299 transcript_path: None,
1300 archive_path: None,
1301 })),
1302 };
1303 self.record(ThreadEvent::ItemCompleted(ItemCompletedEvent { item }));
1304 }
1305
1306 pub fn record_tool_latency(&mut self, tool_name: &str, duration_ms: u64) {
1308 self.record_tool_outcome(tool_name, 1, duration_ms, None);
1309 }
1310
1311 pub fn record_tool_outcome(
1313 &mut self,
1314 tool_name: &str,
1315 attempts: u32,
1316 duration_ms: u64,
1317 error_category: Option<&str>,
1318 ) {
1319 let item = ThreadItem {
1320 context: None,
1321 id: self.next_item_id(),
1322 details: ThreadItemDetails::Harness(Box::new(HarnessEventItem {
1323 event: HarnessEventKind::ToolLatencyRecorded,
1324 message: Some(format!("{tool_name} completed in {duration_ms}ms")),
1325 command: None,
1326 path: None,
1327 exit_code: None,
1328 attempt: Some(attempts),
1329 error_category: error_category.map(str::to_string),
1330 duration_ms: Some(duration_ms),
1331 task_id: None,
1332 session_id: None,
1333 exec_session_id: None,
1334 status: None,
1335 transcript_path: None,
1336 archive_path: None,
1337 })),
1338 };
1339 self.record(ThreadEvent::ItemCompleted(ItemCompletedEvent { item }));
1340 }
1341
1342 pub fn error_recovered(&mut self, tool_name: &str, attempt: u32, error_category: &str) {
1345 self.harness_event(
1346 HarnessEventKind::ErrorRecovered,
1347 Some(format!("{tool_name} recovered after {attempt} retries")),
1348 None,
1349 None,
1350 None,
1351 Some(attempt),
1352 Some(error_category.to_string()),
1353 );
1354 }
1355
1356 pub fn tool_retry_attempted(&mut self, tool_name: &str, attempt: u32, error_category: &str, delay_ms: u64) {
1359 self.harness_event(
1360 HarnessEventKind::ToolRetryAttempted,
1361 Some(format!("{tool_name}: retry {attempt} after {delay_ms}ms")),
1362 None,
1363 None,
1364 None,
1365 Some(attempt),
1366 Some(error_category.to_string()),
1367 );
1368 }
1369
1370 pub fn take_events(&mut self) -> Vec<ThreadEvent> {
1373 self.lifecycle.complete_open_items();
1374 self.record_pending_lifecycle_events();
1375 std::mem::take(&mut self.events)
1376 }
1377
1378 pub fn into_events(mut self) -> Vec<ThreadEvent> {
1379 self.take_events()
1380 }
1381}
1382
1383#[cfg(test)]
1384mod tests {
1385 fn pending_verifier(recorder: &mut ExecEventRecorder, session: &str) -> String {
1386 let args = serde_json::json!({"cmd":"cargo check --locked"});
1387 let launch = recorder.tool_started("exec_command", Some(&args), Some(session));
1388 recorder.tool_finished(&launch, crate::exec::events::ToolCallStatus::Completed, None, "still running", None);
1389 recorder.record_exec_session_output(
1390 &launch.id,
1391 "exec_command",
1392 &args,
1393 &serde_json::json!({"session_id":session,"output":"still running"}),
1394 );
1395 launch.id
1396 }
1397
1398 #[tokio::test]
1399 async fn exec_session_verification_tracks_terminal_output_and_preserves_launch_task() {
1400 use vtcode_memory::explanation::ExplanationScope;
1401 let workspace = tempfile::tempdir().unwrap();
1402 let sink = SessionStoreSink::open(workspace.path(), "exec-verification").await.unwrap();
1403 let mut recorder = ExecEventRecorder::new("exec-verification", Some(sink.event_sink()), None);
1404 let launch_task = recorder.task_started("Verify changes");
1405 for session in ["successful", "failed", "pending"] {
1406 pending_verifier(&mut recorder, session);
1407 }
1408 let pending = sink.explanation(ExplanationScope::Task).await.unwrap();
1409 assert_eq!(pending.verification.len(), 3);
1410 assert!(
1411 pending
1412 .verification
1413 .iter()
1414 .all(|entry| entry.fact.status == "pending" && !entry.fresh)
1415 );
1416 for (requested, returned) in [("unrelated", "unrelated"), ("successful", "failed")] {
1417 recorder.record_exec_session_output(
1418 "poll",
1419 "write_stdin",
1420 &serde_json::json!({"session_id":requested,"chars":""}),
1421 &serde_json::json!({"session_id":returned,"exit_code":0}),
1422 );
1423 }
1424 assert!(
1425 sink.explanation(ExplanationScope::Task)
1426 .await
1427 .unwrap()
1428 .verification
1429 .iter()
1430 .all(|entry| entry.fact.status == "pending")
1431 );
1432 for (session, exit) in [("successful", Some(0)), ("failed", Some(7)), ("pending", None)] {
1433 recorder.record_exec_session_output(
1434 "poll",
1435 "write_stdin",
1436 &serde_json::json!({"session_id":session,"chars":""}),
1437 &serde_json::json!({"session_id":session,"exit_code":exit,"output":"terminal observation"}),
1438 );
1439 }
1440 let terminal = sink.explanation(ExplanationScope::Task).await.unwrap();
1441 assert_eq!(
1442 terminal
1443 .verification
1444 .iter()
1445 .map(|entry| (entry.fact.status.as_str(), entry.exit_code, entry.fresh))
1446 .collect::<Vec<_>>(),
1447 vec![
1448 ("passed", Some(0), true),
1449 ("failed or denied", Some(7), false),
1450 ("pending", None, false)
1451 ]
1452 );
1453 let new_task = recorder.task_started("A different task");
1454 assert_ne!(launch_task, new_task);
1455 recorder.record_exec_session_output(
1456 "poll-new-task",
1457 "write_stdin",
1458 &serde_json::json!({"session_id":"pending","chars":""}),
1459 &serde_json::json!({"session_id":"pending","exit_code":0}),
1460 );
1461 assert!(sink.explanation(ExplanationScope::Task).await.unwrap().verification.is_empty());
1462 let all = sink.explanation(ExplanationScope::Session).await.unwrap();
1463 assert_eq!(all.verification.len(), 3);
1464 assert!(
1465 all.verification
1466 .iter()
1467 .all(|entry| entry.fact.task_id.as_deref() == Some(launch_task.as_str()))
1468 );
1469 assert_eq!(all.verification[2].fact.status, "passed");
1470 sink.close().await.unwrap();
1471 }
1472
1473 #[tokio::test]
1474 async fn exec_session_verification_launched_before_a_mutation_remains_stale() {
1475 let workspace = tempfile::tempdir().unwrap();
1476 let sink = SessionStoreSink::open(workspace.path(), "exec-freshness").await.unwrap();
1477 let mut recorder = ExecEventRecorder::new("exec-freshness", Some(sink.event_sink()), None);
1478 recorder.task_started("Verify changes");
1479 pending_verifier(&mut recorder, "before-edit");
1480 recorder.record_thread_event(ThreadEvent::ItemCompleted(ItemCompletedEvent {
1481 item: ThreadItem {
1482 id: "edit".into(),
1483 context: None,
1484 details: ThreadItemDetails::FileChange(Box::new(crate::exec::events::FileChangeItem {
1485 changes: vec![crate::exec::events::FileUpdateChange {
1486 path: "parser.rs".into(),
1487 kind: crate::exec::events::PatchChangeKind::Update,
1488 }],
1489 status: crate::exec::events::PatchApplyStatus::Completed,
1490 unified_diff: None,
1491 diff_incomplete: None,
1492 additions: Some(1),
1493 deletions: Some(0),
1494 })),
1495 },
1496 }));
1497 recorder.record_exec_session_output(
1498 "poll",
1499 "write_stdin",
1500 &serde_json::json!({"session_id":"before-edit","chars":""}),
1501 &serde_json::json!({"session_id":"before-edit","exit_code":0}),
1502 );
1503 let model = sink
1504 .explanation(vtcode_memory::explanation::ExplanationScope::Task)
1505 .await
1506 .unwrap();
1507 assert_eq!(model.verification[0].fact.status, "passed");
1508 assert!(!model.verification[0].fresh);
1509 pending_verifier(&mut recorder, "after-edit");
1510 recorder.record_exec_session_output(
1511 "poll",
1512 "write_stdin",
1513 &serde_json::json!({"session_id":"after-edit","chars":""}),
1514 &serde_json::json!({"session_id":"after-edit","exit_code":0}),
1515 );
1516 assert!(
1517 sink.explanation(vtcode_memory::explanation::ExplanationScope::Task)
1518 .await
1519 .unwrap()
1520 .verification[1]
1521 .fresh
1522 );
1523 sink.close().await.unwrap();
1524 }
1525
1526 #[tokio::test]
1527 async fn explanation_barrier_and_decision_validation_include_accepted_events() {
1528 let workspace = tempfile::tempdir().unwrap();
1529 let sink = SessionStoreSink::open(workspace.path(), "explanation-barrier").await.unwrap();
1530 let mut recorder = ExecEventRecorder::new("explanation-barrier", Some(sink.event_sink()), None);
1531 let task_id = recorder.task_started("Fix parser");
1532 recorder.record_thread_event(ThreadEvent::ItemCompleted(ItemCompletedEvent {
1533 item: ThreadItem {
1534 id: "parser-change".into(),
1535 context: None,
1536 details: ThreadItemDetails::Decision(Box::new(crate::exec::events::DecisionItem {
1537 summary: "Use existing parser".into(),
1538 rationale: "Preserve validated inputs".into(),
1539 alternatives: vec![],
1540 evidence_ids: vec![],
1541 })),
1542 },
1543 }));
1544 let validate = sink.decision_validator();
1545 validate(task_id.clone(), vec!["parser-change".into()]).await.unwrap();
1546 assert!(validate("other-task".into(), vec!["parser-change".into()]).await.is_err());
1547 assert!(validate(task_id.clone(), vec!["missing".into()]).await.is_err());
1548 let model = sink
1549 .explanation(vtcode_memory::explanation::ExplanationScope::Task)
1550 .await
1551 .unwrap();
1552 assert_eq!(model.task_id.as_deref(), Some(task_id.as_str()));
1553 assert_eq!(model.decisions.len(), 1);
1554 let projected_page = sink
1555 .explanation_page(vtcode_memory::explanation::ExplanationScope::Task, 0)
1556 .await
1557 .unwrap();
1558 assert_eq!(projected_page.model.revision, model.revision);
1559 assert_eq!(projected_page.model.decisions, model.decisions);
1560 let reference = model.decisions[0].fact.evidence.clone();
1561 let page = sink.evidence(reference, 0).await.unwrap();
1562 assert!(page.text.contains("Preserve validated inputs"));
1563 recorder.turn_completed(Usage::default());
1564 recorder.turn_started();
1565 let updated_page = sink
1566 .explanation_page(vtcode_memory::explanation::ExplanationScope::Task, 0)
1567 .await
1568 .unwrap();
1569 assert_ne!(updated_page.model.revision, projected_page.model.revision);
1570 assert_eq!(updated_page.model.decisions, model.decisions);
1571 let starts: Vec<_> = recorder
1572 .events
1573 .iter()
1574 .filter_map(|e| {
1575 if let ThreadEvent::TurnStarted(t) = e {
1576 t.context.as_deref()
1577 } else {
1578 None
1579 }
1580 })
1581 .collect();
1582 assert_eq!(starts.len(), 2);
1583 assert_eq!(starts[0].task_id, starts[1].task_id);
1584 assert_ne!(starts[0].turn_id, starts[1].turn_id);
1585 sink.close().await.unwrap();
1586 assert!(validate(task_id, vec!["parser-change".into()]).await.is_err());
1587 }
1588 #[test]
1589 fn ten_thousand_events_share_cached_projections_for_bounded_pages() {
1590 use vtcode_memory::explanation::{ExplanationScope, page_explanation};
1591
1592 let workspace = tempfile::tempdir().unwrap();
1593 let log = vtcode_memory::open(workspace.path(), "large-explanation", DEFAULT_MAX_EVENTS).unwrap();
1594 let mut context = ExecutionContextTracker::default();
1595 context.begin("root", "turn", "Verify retained work", crate::exec::events::InputOrigin::User);
1596 let mut started = ThreadEvent::TurnStarted(TurnStartedEvent::default());
1597 context.annotate(&mut started);
1598 log.append(&started).unwrap();
1599 for index in 0..DEFAULT_MAX_EVENTS - 1 {
1600 let mut event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1601 item: ThreadItem {
1602 id: format!("check-{index}"),
1603 context: None,
1604 details: ThreadItemDetails::CommandExecution(Box::new(CommandExecutionItem {
1605 command: format!("cargo check --locked # {}", "x".repeat(256)),
1606 arguments: None,
1607 aggregated_output: "discarded output".repeat(64),
1608 exit_code: Some(0),
1609 status: CommandExecutionStatus::Completed,
1610 })),
1611 },
1612 });
1613 context.annotate(&mut event);
1614 log.append(&event).unwrap();
1615 }
1616 let mut cache = Vec::new();
1617 let original = cached_explanation(&log, ExplanationScope::Task, &mut cache).unwrap();
1618 assert_eq!(original.actions.len(), DEFAULT_MAX_EVENTS - 1);
1619 let mut counter = LimitedByteCounter { bytes: 0, limit: EXPLANATION_CACHE_MAX_BYTES };
1620 serde_json::to_writer(&mut counter, original.as_ref()).unwrap();
1621 assert!(counter.bytes > 8 * 1024 * 1024, "fixture exceeds the old cache ceiling");
1622 let mut offset = 0;
1623 let mut seen = 0;
1624 loop {
1625 let shared = cached_explanation(&log, ExplanationScope::Task, &mut cache).unwrap();
1626 assert!(Arc::ptr_eq(&original, &shared), "pages reuse the same reduced snapshot");
1627 let page = page_explanation(&shared, offset);
1628 assert!(serde_json::to_vec(&page).unwrap().len() <= 64 * 1024);
1629 assert_eq!(page.model.revision, original.revision);
1630 seen += page.model.actions.len();
1631 match page.next_offset {
1632 Some(next) => offset = next,
1633 None => break,
1634 }
1635 }
1636 assert_eq!(seen, original.actions.len());
1637 assert_eq!(cache.len(), 1);
1638 }
1639
1640 use super::*;
1641 use crate::core::threads::{ThreadBootstrap, ThreadManager};
1642 use std::fs;
1643 use std::time::Duration;
1644 use tempfile::TempDir;
1645 use tokio::sync::oneshot;
1646 use tokio::time::timeout;
1647
1648 fn make_recorder() -> ExecEventRecorder {
1649 ExecEventRecorder::new("thread", None, None)
1650 }
1651
1652 #[test]
1653 fn compact_boundary_reports_the_installed_segment_identity() {
1654 let previous = crate::core::agent::request_envelope::SessionRequestEnvelope::with_prefix_hash(
1655 "segment-1",
1656 "prompt",
1657 vec![crate::llm::provider::ToolDefinition::function(
1658 "read_file".to_string(),
1659 "Read a file".to_string(),
1660 serde_json::json!({"type": "object"}),
1661 )],
1662 11,
1663 22,
1664 );
1665 let next = previous.begin_segment("segment-2");
1666 let mut recorder = make_recorder();
1667
1668 recorder.compact_boundary(
1669 CompactionTrigger::Auto,
1670 CompactionMode::Local,
1671 20,
1672 8,
1673 None,
1674 Some(&previous),
1675 Some(&next),
1676 );
1677
1678 let events = recorder.into_events();
1679 let Some(ThreadEvent::ThreadCompactBoundary(boundary)) = events.last() else {
1680 panic!("expected compact boundary");
1681 };
1682 assert_eq!(boundary.previous_segment_id.as_deref(), Some(previous.segment_id()));
1683 assert_eq!(boundary.new_segment_id.as_deref(), Some(next.segment_id()));
1684 assert_eq!(boundary.previous_prefix_hash, boundary.new_prefix_hash);
1685 assert_eq!(boundary.previous_catalog_hash, boundary.new_catalog_hash);
1686 assert!(boundary.previous_catalog_hash.is_some());
1687 }
1688
1689 #[tokio::test]
1690 async fn session_store_sink_preserves_order_when_queue_is_small() {
1691 let workspace = TempDir::new().expect("workspace");
1692 let (sink, handle) = session_store_sink_with_capacity_handle(workspace.path(), "session", 1)
1693 .await
1694 .expect("session sink");
1695 let health = Arc::clone(&handle.state.health);
1696 let events = vec![
1697 ThreadEvent::ThreadStarted(ThreadStartedEvent { thread_id: "thread".to_string() }),
1698 ThreadEvent::TurnStarted(TurnStartedEvent::default()),
1699 ThreadEvent::TurnCompleted(TurnCompletedEvent {
1700 completed_at: None,
1701 usage: Usage::default(),
1702 in_progress_exec_sessions: Vec::new(),
1703 }),
1704 ];
1705
1706 for (index, event) in events.iter().enumerate() {
1707 if index > 0 {
1708 timeout(Duration::from_secs(5), async {
1709 while health.snapshot().persisted_events < index as u64 {
1710 tokio::task::yield_now().await;
1711 }
1712 })
1713 .await
1714 .expect("bounded queue should keep draining");
1715 }
1716 let mut callback = sink.lock();
1717 callback(event);
1718 assert_eq!(health.snapshot().accepted_events, index as u64 + 1);
1719 }
1720 drop(sink);
1721
1722 timeout(Duration::from_secs(5), handle.close())
1723 .await
1724 .expect("session sink should drain before timeout")
1725 .expect("session sink should close successfully");
1726
1727 assert_eq!(health.snapshot().accepted_events, events.len() as u64);
1728 let log = vtcode_memory::open(workspace.path(), "session", DEFAULT_MAX_EVENTS).expect("reopen session");
1729 assert_eq!(log.event_count(), events.len() as u64);
1730 let event_path = workspace.path().join(".vtcode/sessions/session/events.jsonl");
1731 let persisted = fs::read_to_string(event_path)
1732 .expect("read persisted events")
1733 .lines()
1734 .map(|line| serde_json::from_str::<VersionedThreadEvent>(line).expect("decode event"))
1735 .map(VersionedThreadEvent::into_event)
1736 .collect::<Vec<_>>();
1737 assert_eq!(persisted, events);
1738 assert_eq!(health.snapshot().append_failures, 0);
1739 assert_eq!(health.snapshot().channel_failures, 0);
1740 assert!(!health.snapshot().failed);
1741 }
1742
1743 #[tokio::test]
1744 async fn session_store_sink_open_failure_is_propagated() {
1745 let temp_dir = TempDir::new().expect("workspace");
1746 let workspace_file = temp_dir.path().join("workspace-file");
1747 fs::write(&workspace_file, "not a directory").expect("workspace file");
1748
1749 let result = SessionStoreSink::open(&workspace_file, "session").await;
1750 assert!(result.is_err(), "canonical store setup must fail closed");
1751 }
1752
1753 #[tokio::test]
1754 async fn session_store_sink_drain_failure_is_propagated() {
1755 let health = Arc::new(SessionStoreSinkHealth::default());
1756 health.failed.store(true, Ordering::Release);
1757 let state = Arc::new(SessionStoreSinkState {
1758 sender: Mutex::new(None),
1759 reserved_bytes: AtomicU64::new(0),
1760 max_bytes: SESSION_STORE_DRAIN_MAX_BYTES,
1761 health,
1762 });
1763 let drain = tokio::spawn(async {});
1764 let handle = SessionStoreSinkHandle { state, drain: Some(drain) };
1765
1766 let result = handle.close().await;
1767 assert!(result.is_err(), "canonical drain failures must fail closed");
1768 }
1769
1770 #[tokio::test]
1771 async fn session_store_sink_concurrent_close_waits_for_shared_completion() {
1772 let health = Arc::new(SessionStoreSinkHealth::default());
1773 let state = Arc::new(SessionStoreSinkState {
1774 sender: Mutex::new(None),
1775 reserved_bytes: AtomicU64::new(0),
1776 max_bytes: SESSION_STORE_DRAIN_MAX_BYTES,
1777 health,
1778 });
1779 let (started_sender, started_receiver) = oneshot::channel();
1780 let (release_sender, release_receiver) = oneshot::channel();
1781 let drain = tokio::spawn(async move {
1782 started_sender.send(()).expect("notify drain start");
1783 release_receiver.await.expect("release drain");
1784 });
1785 let handle = SessionStoreSinkHandle { state: Arc::clone(&state), drain: Some(drain) };
1786 let sink = Arc::new(SessionStoreSink {
1787 state,
1788 handle: Arc::new(Mutex::new(Some(handle))),
1789 close_result: Arc::new(OnceCell::new()),
1790 });
1791
1792 let first = {
1793 let sink = Arc::clone(&sink);
1794 tokio::spawn(async move { sink.close().await })
1795 };
1796 started_receiver.await.expect("drain should start");
1797 let second = {
1798 let sink = Arc::clone(&sink);
1799 tokio::spawn(async move { sink.close().await })
1800 };
1801 tokio::task::yield_now().await;
1802 assert!(!second.is_finished(), "concurrent close must wait for the first drain");
1803
1804 release_sender.send(()).expect("release drain");
1805 assert!(first.await.expect("first close task").is_ok());
1806 assert!(second.await.expect("second close task").is_ok());
1807 }
1808
1809 #[tokio::test]
1810 async fn session_store_sink_rejects_events_over_byte_budget() {
1811 let workspace = TempDir::new().expect("workspace");
1812 let (state, handle) = open_session_store_sink_with_limits(workspace.path(), "session", 4, 128)
1813 .await
1814 .expect("session sink");
1815 let health = Arc::clone(&state.health);
1816 let event = ThreadEvent::ThreadStarted(ThreadStartedEvent { thread_id: "x".repeat(256) });
1817
1818 let result = enqueue_session_event(&state, &event);
1819 assert!(result.is_err(), "oversized events must fail closed");
1820 handle.close().await.expect_err("failed sink must not close successfully");
1821
1822 let snapshot = health.snapshot();
1823 assert_eq!(snapshot.accepted_events, 0);
1824 assert_eq!(snapshot.serialization_failures, 1);
1825 assert!(snapshot.failed);
1826 }
1827
1828 #[test]
1829 fn session_store_sink_saturation_fails_closed() {
1830 let (sender, receiver) = mpsc::sync_channel(1);
1831 let state = SessionStoreSinkState {
1832 sender: Mutex::new(Some(sender)),
1833 reserved_bytes: AtomicU64::new(0),
1834 max_bytes: SESSION_STORE_DRAIN_MAX_BYTES,
1835 health: Arc::new(SessionStoreSinkHealth::default()),
1836 };
1837 let first = ThreadEvent::TurnStarted(TurnStartedEvent::default());
1838 let second = ThreadEvent::TurnCompleted(TurnCompletedEvent {
1839 completed_at: None,
1840 usage: Usage::default(),
1841 in_progress_exec_sessions: Vec::new(),
1842 });
1843
1844 enqueue_session_event(&state, &first).expect("first event should fit");
1845 assert!(enqueue_session_event(&state, &second).is_err(), "queue saturation must fail closed");
1846
1847 let snapshot = state.health.snapshot();
1848 assert_eq!(snapshot.accepted_events, 1);
1849 assert_eq!(snapshot.channel_failures, 1);
1850 assert!(snapshot.failed);
1851 drop(receiver);
1852 }
1853
1854 #[test]
1855 fn closed_session_store_channel_is_observable() {
1856 let (sender, receiver) = mpsc::sync_channel(1);
1857 drop(receiver);
1858 let health = SessionStoreSinkHealth::default();
1859 let event = ThreadEvent::TurnStarted(TurnStartedEvent::default());
1860
1861 if sender.send(event).is_err() {
1862 health.channel_failures.fetch_add(1, Ordering::Relaxed);
1863 }
1864
1865 assert_eq!(health.snapshot().channel_failures, 1);
1866 }
1867
1868 #[test]
1869 fn streaming_events_flush_on_completion() {
1870 let mut recorder = make_recorder();
1871 recorder.turn_started();
1872 assert!(recorder.agent_message_stream_update("partial"));
1873 recorder.agent_message_stream_complete();
1874 let events = recorder.into_events();
1875 assert!(events.iter().any(|event| matches!(event, ThreadEvent::ItemCompleted(_))));
1876 }
1877
1878 #[test]
1879 fn thread_failure_emits_terminal_lifecycle_events() {
1880 let mut recorder = make_recorder();
1881 recorder.turn_started();
1882 recorder.thread_failed("session", "setup failed", 1);
1883 let events = recorder.into_events();
1884
1885 assert!(events.iter().any(|event| matches!(event, ThreadEvent::TurnFailed(_))));
1886 assert!(events.iter().any(|event| {
1887 matches!(
1888 event,
1889 ThreadEvent::ThreadCompleted(item)
1890 if item.subtype == ThreadCompletionSubtype::ErrorDuringExecution
1891 && item.outcome_code == "error"
1892 )
1893 }));
1894 }
1895
1896 #[test]
1897 fn command_events_capture_status() {
1898 let mut recorder = make_recorder();
1899 let handle = recorder.command_started("git status");
1900 recorder.command_finished(&handle, CommandExecutionStatus::Completed, Some(0), "");
1901 let events = recorder.into_events();
1902 let command = events
1903 .into_iter()
1904 .filter_map(|event| match event {
1905 ThreadEvent::ItemCompleted(event) => Some(event.item),
1906 _ => None,
1907 })
1908 .find(|item| matches!(item.details, ThreadItemDetails::CommandExecution(_)))
1909 .expect("command event should be emitted");
1910
1911 match command.details {
1912 ThreadItemDetails::CommandExecution(details) => {
1913 assert_eq!(details.command, "git status");
1914 assert_eq!(details.status, CommandExecutionStatus::Completed);
1915 }
1916 _ => panic!("unexpected event variant"),
1917 }
1918 }
1919
1920 #[test]
1921 fn rejected_tool_call_emits_failed_tool_output_item() {
1922 let mut recorder = make_recorder();
1923 recorder.tool_rejected("read_file", None, Some("call_1"), "Tool permission denied");
1924
1925 let events = recorder.into_events();
1926 let tool_outputs = events
1927 .iter()
1928 .filter_map(|event| match event {
1929 ThreadEvent::ItemCompleted(ItemCompletedEvent { item }) => match &item.details {
1930 ThreadItemDetails::ToolOutput(details) => Some(details),
1931 _ => None,
1932 },
1933 _ => None,
1934 })
1935 .collect::<Vec<_>>();
1936
1937 assert_eq!(tool_outputs.len(), 1);
1938 assert_eq!(tool_outputs[0].tool_call_id.as_deref(), Some("call_1"));
1939 assert_eq!(tool_outputs[0].status, crate::exec::events::ToolCallStatus::Failed);
1940 assert_eq!(tool_outputs[0].output, "Tool permission denied");
1941 }
1942
1943 #[test]
1944 fn thread_backed_recorder_reuses_submission_id_within_turn() {
1945 let handle = ThreadManager::new().start_thread_with_identifier("thread", ThreadBootstrap::new(None));
1946 let mut recorder = ExecEventRecorder::new("thread", None, Some(handle.clone()));
1947
1948 recorder.turn_started();
1949 recorder.agent_message("hello");
1950 recorder.turn_completed(Usage::default());
1951
1952 let records = handle.replay_recent();
1953 let submission_ids: std::collections::BTreeSet<String> = records
1954 .iter()
1955 .filter_map(|record| record.submission_id.as_ref().map(|id| id.as_str().to_string()))
1956 .collect();
1957
1958 assert_eq!(submission_ids.len(), 1);
1959 assert!(
1960 records
1961 .iter()
1962 .any(|record| matches!(record.event, ThreadEvent::TurnStarted(_)) && record.submission_id.is_some())
1963 );
1964 assert!(
1965 records
1966 .iter()
1967 .any(|record| matches!(record.event, ThreadEvent::TurnCompleted(_)) && record.submission_id.is_some())
1968 );
1969 }
1970
1971 #[test]
1972 fn thread_backed_recorder_keeps_full_event_history_beyond_thread_buffer() {
1973 let handle = ThreadManager::with_event_buffer_capacity(2)
1974 .start_thread_with_identifier("thread", ThreadBootstrap::new(None));
1975 let mut recorder = ExecEventRecorder::new("thread", None, Some(handle.clone()));
1976
1977 recorder.turn_started();
1978 recorder.agent_message("first");
1979 recorder.agent_message("second");
1980 recorder.turn_completed(Usage::default());
1981
1982 let full_events = recorder.into_events();
1983 let buffered_events = handle.recent_events();
1984
1985 assert_eq!(buffered_events.len(), 2);
1986 assert!(full_events.len() > buffered_events.len());
1987 assert_eq!(
1988 full_events
1989 .iter()
1990 .filter(|event| matches!(event, ThreadEvent::ItemCompleted(_)))
1991 .count(),
1992 2
1993 );
1994 }
1995}