Skip to main content

vtcode_core/core/agent/events/
mod.rs

1//! Event recording utilities for the agent runner.
2
3mod 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
41/// Callback type alias for streaming structured events.
42pub type EventSink = Arc<Mutex<Box<dyn FnMut(&ThreadEvent) + Send>>>;
43
44/// Ordered validation of canonical evidence identities within the current task.
45pub 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
185/// Owns the session persistence drain and exposes its final health result.
186///
187/// The event callback remains synchronous for compatibility, while the runner
188/// retains this handle so a completed task cannot report success before its
189/// authoritative event queue has drained.
190pub(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
220/// Authoritative event sink for one canonical session store.
221///
222/// Events are accepted in order through one bounded, non-blocking queue owned
223/// by a blocking drain actor. The queue is bounded both by event count and by
224/// an estimated serialized payload budget, so a burst of large tool outputs
225/// cannot retain unbounded memory. If the queue or store fails, the sink fails
226/// closed and [`Self::close`] reports the loss. Call [`Self::close`] exactly
227/// once after the terminal event and propagate its error before reporting a
228/// successful run.
229pub 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    /// Ordered, acknowledged matrix writes and replay through this session drain.
237    pub fn matrix_persistence(&self) -> crate::subagents::matrix::MatrixPersistence {
238        matrix_persistence(Arc::clone(&self.state))
239    }
240    /// Validate task-owned evidence after all previously accepted events.
241    pub fn decision_validator(&self) -> DecisionEvidenceValidator {
242        decision_validator(Arc::clone(&self.state))
243    }
244    /// Open the canonical session store and start its bounded drain.
245    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    /// Enqueue one event for canonical persistence.
255    pub fn emit(&self, event: &ThreadEvent) -> Result<()> {
256        enqueue_session_event(&self.state, event)
257    }
258
259    /// Return the callback form used by the core event recorder.
260    pub fn event_sink(&self) -> EventSink {
261        event_sink_for_state(Arc::clone(&self.state))
262    }
263
264    /// Ordered query barrier; accepted events precede this retained snapshot.
265    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    /// Return only a bounded page from the ordered cached projection.
275    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    /// Resolve a bounded evidence page through the same ordered barrier.
286    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    /// Drain and close the canonical persistence task.
310    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
342/// Build an event sink that persists every recorded event to the unified
343/// per-session store ([`vtcode_memory`]), making it the canonical source of
344/// truth for session state/history.
345///
346/// The sink hands events to one bounded, non-blocking queue and a blocking
347/// background drain. If the queue cannot accept an event, the sink records a
348/// fatal persistence failure and `close()` fails the run. Disk I/O and
349/// manifest writes remain off the Tokio runtime worker through `spawn_blocking`.
350pub 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    // Two scopes, invalidated on append; oversized projections remain uncached.
431    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                    // Cap rewrites can retain checkpoints after evicting the
488                    // command records. Their evidence remains historical.
489                    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
782/// Combine two optional event sinks into one that fans out to both.
783pub 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/// Helper responsible for recording execution events and relaying them to optional sinks.
816#[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    /// Attach a session result to its original verifier's canonical output.
879    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                    // A failed begin_turn means a previous turn was never
925                    // finished (or a concurrent turn is in flight). Surface
926                    // it instead of silently dropping the submission id:
927                    // without it, every event of this turn loses its
928                    // submission context on retried attempts.
929                    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    /// Start a task with its original public request.
938    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    /// Record the terminal lifecycle events for an execution failure.
998    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    /// Emit a tool latency harness event with recorded duration.
1307    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    /// Emit the single terminal harness observation for one invocation.
1312    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    /// Emit an `ErrorRecovered` harness event, recording that the agent
1343    /// successfully recovered from a transient error after retries.
1344    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    /// Emit a `ToolRetryAttempted` harness event, recording that a transient
1357    /// tool failure triggered an automatic retry.
1358    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    /// Drain the collected events while keeping the recorder available for a
1371    /// terminal failure path that may run after the normal result assembly.
1372    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}