Skip to main content

atman_runtime/
event_writer.rs

1use std::collections::VecDeque;
2use std::path::{Path, PathBuf};
3use std::sync::Arc;
4use std::time::Duration;
5
6use tokio::io::{AsyncSeekExt, AsyncWriteExt, SeekFrom};
7use tokio::sync::{mpsc, oneshot};
8
9use crate::event::{Event, EventEnvelope};
10use crate::index::{AnchorIndex, ProjectEventInsert};
11use crate::redact::Redactor;
12
13pub struct EventWriter {
14    thread: Option<std::thread::JoinHandle<()>>,
15    tx: mpsc::UnboundedSender<EventEnvelope>,
16    flush_tx: mpsc::UnboundedSender<oneshot::Sender<()>>,
17    stop_tx: Option<oneshot::Sender<()>>,
18    events_path: PathBuf,
19}
20
21impl EventWriter {
22    pub fn spawn(session_dir: impl AsRef<Path>) -> std::io::Result<Self> {
23        Self::spawn_full(session_dir, None, None, None)
24    }
25
26    pub fn spawn_with(
27        session_dir: impl AsRef<Path>,
28        redactor: Option<Arc<Redactor>>,
29    ) -> std::io::Result<Self> {
30        Self::spawn_full(session_dir, redactor, None, None)
31    }
32
33    // Owns its own thread + rt so short-lived caller runtimes
34    // (spawn_blocking + throwaway current_thread rt) can't kill the loop.
35    pub fn spawn_full(
36        session_dir: impl AsRef<Path>,
37        redactor: Option<Arc<Redactor>>,
38        project_index: Option<Arc<AnchorIndex>>,
39        session_id: Option<String>,
40    ) -> std::io::Result<Self> {
41        let session_dir = session_dir.as_ref().to_path_buf();
42        let events_path = session_dir.join("events.jsonl");
43        std::fs::create_dir_all(&session_dir)?;
44        let (tx, rx) = mpsc::unbounded_channel::<EventEnvelope>();
45        let (flush_tx, flush_rx) = mpsc::unbounded_channel::<oneshot::Sender<()>>();
46        let (stop_tx, stop_rx) = oneshot::channel::<()>();
47        let file_path = events_path.clone();
48        let thread = std::thread::Builder::new()
49            .name("atman-event-writer".into())
50            .spawn(move || {
51                let rt = match tokio::runtime::Builder::new_current_thread()
52                    .enable_all()
53                    .build()
54                {
55                    Ok(rt) => rt,
56                    Err(e) => {
57                        crate::notify!(error, "event writer rt init failed: {e}");
58                        return;
59                    }
60                };
61                rt.block_on(async move {
62                    if let Err(e) = writer_loop(
63                        rx,
64                        flush_rx,
65                        stop_rx,
66                        &file_path,
67                        project_index,
68                        session_id,
69                        redactor,
70                    )
71                    .await
72                    {
73                        crate::notify!(error, "event writer failed: {e}");
74                    }
75                });
76            })?;
77        Ok(Self {
78            thread: Some(thread),
79            tx,
80            flush_tx,
81            stop_tx: Some(stop_tx),
82            events_path,
83        })
84    }
85
86    pub async fn flush(&self) {
87        let (tx, rx) = oneshot::channel::<()>();
88        if self.flush_tx.send(tx).is_err() {
89            return;
90        }
91        let _ = rx.await;
92    }
93
94    pub fn sender(&self) -> mpsc::UnboundedSender<EventEnvelope> {
95        self.tx.clone()
96    }
97
98    pub fn events_path(&self) -> &Path {
99        &self.events_path
100    }
101
102    pub async fn shutdown(mut self) {
103        if let Some(stop_tx) = self.stop_tx.take() {
104            let _ = stop_tx.send(());
105        }
106        if let Some(thread) = self.thread.take() {
107            let _ = tokio::task::spawn_blocking(move || {
108                let _ = thread.join();
109            })
110            .await;
111        }
112    }
113}
114
115impl Drop for EventWriter {
116    fn drop(&mut self) {
117        if let Some(stop_tx) = self.stop_tx.take() {
118            let _ = stop_tx.send(());
119        }
120        if let Some(thread) = self.thread.take() {
121            let _ = thread.join();
122        }
123    }
124}
125
126const MAX_BUFFERED: usize = 10_000;
127const RECOVERY_INTERVAL: Duration = Duration::from_secs(2);
128
129struct DegradedBuffer {
130    events: VecDeque<EventEnvelope>,
131    error: Option<String>,
132    dropped: u64,
133    max_buffered: usize,
134}
135
136impl DegradedBuffer {
137    fn new(max_buffered: usize) -> Self {
138        Self {
139            events: VecDeque::new(),
140            error: None,
141            dropped: 0,
142            max_buffered,
143        }
144    }
145
146    fn is_degraded(&self) -> bool {
147        self.error.is_some()
148    }
149
150    fn degrade(&mut self, event: EventEnvelope, error: impl Into<String>) {
151        self.error = Some(error.into());
152        self.buffer(event);
153    }
154
155    fn buffer(&mut self, event: EventEnvelope) {
156        if self.events.len() == self.max_buffered {
157            self.events.pop_front();
158            self.dropped += 1;
159        }
160        self.events.push_back(event);
161    }
162
163    fn recover(&mut self) {
164        self.error = None;
165    }
166}
167
168async fn writer_loop(
169    mut rx: mpsc::UnboundedReceiver<EventEnvelope>,
170    mut flush_rx: mpsc::UnboundedReceiver<oneshot::Sender<()>>,
171    mut stop_rx: oneshot::Receiver<()>,
172    path: &Path,
173    project_index: Option<Arc<AnchorIndex>>,
174    session_id: Option<String>,
175    redactor: Option<Arc<Redactor>>,
176) -> std::io::Result<()> {
177    // O_APPEND prevents seek-based overwrites required for idempotent retries.
178    #[allow(clippy::suspicious_open_options)]
179    let mut file = tokio::fs::OpenOptions::new()
180        .create(true)
181        .read(true)
182        .write(true)
183        .open(path)
184        .await?;
185    let mut offset = file.seek(SeekFrom::End(0)).await?;
186    let session_dir = path.parent().unwrap_or(path);
187    let indexer = project_index.zip(session_id);
188    let mut degraded = DegradedBuffer::new(MAX_BUFFERED);
189    let mut recovery_tick = tokio::time::interval(RECOVERY_INTERVAL);
190
191    loop {
192        tokio::select! {
193            biased;
194            _ = &mut stop_rx => {
195                while let Ok(event) = rx.try_recv() {
196                    handle_event(
197                        &mut file,
198                        &mut offset,
199                        event,
200                        &mut degraded,
201                        indexer.as_ref(),
202                        redactor.as_deref(),
203                        session_dir,
204                    ).await;
205                }
206                while let Ok(waiter) = flush_rx.try_recv() {
207                    let _ = waiter.send(());
208                }
209                break;
210            }
211            _ = recovery_tick.tick() => {
212                retry_degraded_on_tick(
213                    &mut file,
214                    &mut offset,
215                    &mut degraded,
216                    indexer.as_ref(),
217                    redactor.as_deref(),
218                    session_dir,
219                ).await;
220            }
221            maybe_event = rx.recv() => {
222                match maybe_event {
223                    Some(event) => {
224                        handle_event(
225                            &mut file,
226                            &mut offset,
227                            event,
228                            &mut degraded,
229                            indexer.as_ref(),
230                            redactor.as_deref(),
231                            session_dir,
232                        ).await;
233                    }
234                    None => break,
235                }
236            }
237            maybe_flush = flush_rx.recv() => {
238                match maybe_flush {
239                    Some(waiter) => {
240                        while let Ok(event) = rx.try_recv() {
241                            handle_event(
242                                &mut file,
243                                &mut offset,
244                                event,
245                                &mut degraded,
246                                indexer.as_ref(),
247                                redactor.as_deref(),
248                                session_dir,
249                            ).await;
250                        }
251                        retry_buffered(
252                            &mut file,
253                            &mut offset,
254                            &mut degraded,
255                            indexer.as_ref(),
256                            redactor.as_deref(),
257                            session_dir,
258                        ).await;
259                        if let Err(e) = file.sync_data().await {
260                            crate::notify!(error, "event writer flush failed: {e}");
261                        }
262                        let _ = waiter.send(());
263                    }
264                    None => break,
265                }
266            }
267        }
268    }
269    retry_buffered(
270        &mut file,
271        &mut offset,
272        &mut degraded,
273        indexer.as_ref(),
274        redactor.as_deref(),
275        session_dir,
276    )
277    .await;
278    if !degraded.events.is_empty() {
279        crate::notify!(
280            error,
281            "event writer stopped with {} buffered event(s) not persisted (dropped {}): {}",
282            degraded.events.len(),
283            degraded.dropped,
284            degraded.error.as_deref().unwrap_or("write failed")
285        );
286    }
287    if let Err(e) = file.sync_data().await {
288        crate::notify!(error, "event writer final sync failed: {e}");
289    }
290    Ok(())
291}
292
293async fn handle_event(
294    file: &mut tokio::fs::File,
295    offset: &mut u64,
296    event: EventEnvelope,
297    degraded: &mut DegradedBuffer,
298    indexer: Option<&(Arc<AnchorIndex>, String)>,
299    redactor: Option<&Redactor>,
300    session_dir: &Path,
301) {
302    if degraded.is_degraded() {
303        degraded.buffer(event);
304        return;
305    }
306    if let Err(e) = write_event(file, offset, &event, indexer, redactor, session_dir).await {
307        crate::notify!(
308            error,
309            "event writer write failed (seq={}): {e}; buffering events in memory",
310            event.seq
311        );
312        degraded.degrade(event, e.to_string());
313    }
314}
315
316async fn retry_degraded_on_tick(
317    file: &mut tokio::fs::File,
318    offset: &mut u64,
319    degraded: &mut DegradedBuffer,
320    indexer: Option<&(Arc<AnchorIndex>, String)>,
321    redactor: Option<&Redactor>,
322    session_dir: &Path,
323) {
324    if !degraded.is_degraded() || degraded.events.is_empty() {
325        return;
326    }
327
328    retry_buffered(file, offset, degraded, indexer, redactor, session_dir).await;
329    if !degraded.is_degraded() {
330        crate::notify!(info, location = Status, "事件写入已恢复");
331    }
332}
333
334async fn retry_buffered(
335    file: &mut tokio::fs::File,
336    offset: &mut u64,
337    degraded: &mut DegradedBuffer,
338    indexer: Option<&(Arc<AnchorIndex>, String)>,
339    redactor: Option<&Redactor>,
340    session_dir: &Path,
341) {
342    while let Some(event) = degraded.events.pop_front() {
343        if let Err(e) = write_event(file, offset, &event, indexer, redactor, session_dir).await {
344            crate::notify!(
345                error,
346                "event writer retry failed (seq={}): {e}; {} event(s) remain buffered",
347                event.seq,
348                degraded.events.len() + 1
349            );
350            degraded.events.push_front(event);
351            degraded.error = Some(e.to_string());
352            return;
353        }
354    }
355    degraded.recover();
356}
357
358async fn write_event(
359    file: &mut tokio::fs::File,
360    offset: &mut u64,
361    envelope: &EventEnvelope,
362    indexer: Option<&(Arc<AnchorIndex>, String)>,
363    redactor: Option<&Redactor>,
364    session_dir: &Path,
365) -> std::io::Result<()> {
366    let line = serialize_event(envelope, redactor);
367    let start = *offset;
368    file.seek(SeekFrom::Start(start)).await?;
369    if let Err(e) = file.write_all(line.as_bytes()).await {
370        *offset = start;
371        return Err(e);
372    }
373    if let Err(e) = file.write_all(b"\n").await {
374        *offset = start;
375        return Err(e);
376    }
377    let end = file.stream_position().await?;
378    if let Err(e) = file.sync_data().await {
379        *offset = start;
380        return Err(e);
381    }
382    *offset = end;
383    if let Err(error) =
384        crate::session_meta::SessionStats::record_persisted_event(session_dir, start, end, envelope)
385    {
386        crate::notify!(
387            warn,
388            location = Log,
389            stack = merge_count("session_stats.persist_failed", 60_000),
390            "session stats update failed (seq={}): {error}",
391            envelope.seq
392        );
393    }
394    if let Some((idx, sid)) = indexer
395        && let Err(e) = insert_project_row(idx, sid, envelope, &line)
396    {
397        crate::notify!(
398            warn,
399            location = Log,
400            stack = merge_count("project_index.insert_failed", 60_000),
401            "project index insert failed (seq={}): {e}",
402            envelope.seq
403        );
404    }
405    Ok(())
406}
407
408fn serialize_event(envelope: &EventEnvelope, redactor: Option<&Redactor>) -> String {
409    let Some(r) = redactor else {
410        return serde_json::to_string(envelope).unwrap_or_else(|e| {
411            format!(
412                "{{\"type\":\"encode_error\",\"error\":{:?}}}",
413                e.to_string()
414            )
415        });
416    };
417    let mut value = match serde_json::to_value(envelope) {
418        Ok(v) => v,
419        Err(e) => {
420            return format!(
421                "{{\"type\":\"encode_error\",\"error\":{:?}}}",
422                e.to_string()
423            );
424        }
425    };
426    r.redact_json(&mut value);
427    serde_json::to_string(&value).unwrap_or_else(|e| {
428        format!(
429            "{{\"type\":\"encode_error\",\"error\":{:?}}}",
430            e.to_string()
431        )
432    })
433}
434
435fn insert_project_row(
436    index: &AnchorIndex,
437    session_id: &str,
438    envelope: &EventEnvelope,
439    payload_json: &str,
440) -> rusqlite::Result<()> {
441    let event = &envelope.event;
442    let ts = extract_ts(envelope);
443    let kind = event_kind(event);
444    let (turn_id, flow_run_id) = extract_anchors(event);
445    let text = extract_text_content(event).unwrap_or_default();
446    index.insert_project_event_raw(ProjectEventInsert {
447        session_id,
448        seq: envelope.seq as i64,
449        ts: &ts,
450        kind,
451        turn_id: turn_id.as_deref(),
452        flow_run_id: flow_run_id.as_deref(),
453        text_content: &text,
454        payload_json,
455    })?;
456    Ok(())
457}
458
459pub(crate) fn extract_ts(envelope: &EventEnvelope) -> String {
460    envelope.ts.to_rfc3339()
461}
462
463pub(crate) fn event_kind(event: &Event) -> &'static str {
464    match event {
465        Event::FlowStart { .. } => "flow_start",
466        Event::FlowEnd { .. } => "flow_end",
467        Event::WorkspaceLifecycle { .. } => "workspace_lifecycle",
468        Event::LlmCall { .. } => "llm_call",
469        Event::TurnStart { .. } => "turn_start",
470        Event::TurnEnd { .. } => "turn_end",
471        Event::UserMsg { .. } => "user_msg",
472        Event::AssistantMsg { .. } => "assistant_msg",
473        Event::ToolResultMsg { .. } => "tool_result_msg",
474        Event::ToolResultMetrics { .. } => "tool_result_metrics",
475        Event::DiffPreview { .. } => "diff_preview",
476        Event::FileEditApplied { .. } => "file_edit_applied",
477        Event::CompactionSummary { .. } => "compaction_summary",
478        Event::SystemMsg { .. } => "system_msg",
479        Event::UserInject { .. } => "user_inject",
480        Event::ContentFilterHit { .. } => "content_filter_hit",
481        Event::ContextCompact { .. } => "context_compact",
482        Event::Checkpoint { .. } => "checkpoint",
483        Event::ContextTruncated { .. } => "context_truncated",
484        Event::WatchWarn { .. } => "watch_warn",
485        Event::PendingPrompt { .. } => "pending_prompt",
486        Event::PromptExpired { .. } => "prompt_expired",
487        Event::PromptResolved { .. } => "prompt_resolved",
488        Event::DeferredFormRecorded { .. } => "deferred_form_recorded",
489        Event::DeferredFormApplied { .. } => "deferred_form_applied",
490        Event::LlmPartialCall { .. } => "llm_partial_call",
491        Event::FlowGraph { .. } => "flow_graph",
492        Event::FlowNodeStart { .. } => "flow_node_start",
493        Event::FlowNodeEnd { .. } => "flow_node_end",
494        Event::ToolNode { .. } => "tool_node",
495        Event::AttachmentDegraded { .. } => "attachment_degraded",
496        Event::ToolPendingApproval { .. } => "tool_pending_approval",
497        Event::ToolApproved { .. } => "tool_approved",
498        Event::ToolDenied { .. } => "tool_denied",
499        Event::PermissionRequestCreated { .. } => "permission_request_created",
500        Event::PermissionRequestTargeted { .. } => "permission_request_targeted",
501        Event::PermissionRequestDeferred { .. } => "permission_request_deferred",
502        Event::PermissionRequestApproved { .. } => "permission_request_approved",
503        Event::PermissionRequestDenied { .. } => "permission_request_denied",
504        Event::PermissionRequestCancelled { .. } => "permission_request_cancelled",
505        Event::PermissionGroupCreated { .. } => "permission_group_created",
506        Event::PermissionGroupUpdated { .. } => "permission_group_updated",
507        Event::PermissionGroupResolved { .. } => "permission_group_resolved",
508        Event::PermissionGrantCreated { .. } => "permission_grant_created",
509        Event::PermissionGrantExpired { .. } => "permission_grant_expired",
510        Event::UnrestrictedExecution { .. } => "unrestricted_execution",
511        Event::TerminalFinalState { .. } => "terminal_final_state",
512        Event::MermaidDiagram { .. } => "mermaid_diagram",
513    }
514}
515
516pub(crate) fn extract_anchors(event: &Event) -> (Option<String>, Option<String>) {
517    match event {
518        Event::FlowStart { run_id, .. }
519        | Event::FlowEnd { run_id, .. }
520        | Event::WorkspaceLifecycle { run_id, .. } => (None, Some(run_id.0.to_string())),
521        Event::TurnStart { turn_id, .. } | Event::TurnEnd { turn_id, .. } => {
522            (Some(turn_id.0.to_string()), None)
523        }
524        Event::UserInject { turn_id, .. } => (Some(turn_id.0.to_string()), None),
525        Event::UserMsg {
526            turn_id,
527            flow_run_id,
528            ..
529        }
530        | Event::AssistantMsg {
531            turn_id,
532            flow_run_id,
533            ..
534        }
535        | Event::ToolResultMsg {
536            turn_id,
537            flow_run_id,
538            ..
539        }
540        | Event::ToolResultMetrics {
541            turn_id,
542            flow_run_id,
543            ..
544        }
545        | Event::SystemMsg {
546            turn_id,
547            flow_run_id,
548            ..
549        } => (
550            Some(turn_id.0.to_string()),
551            flow_run_id.as_ref().map(|r| r.0.to_string()),
552        ),
553        Event::DiffPreview {
554            turn_id,
555            flow_run_id,
556            ..
557        } => (
558            turn_id.as_ref().map(|t| t.0.to_string()),
559            flow_run_id.as_ref().map(|r| r.0.to_string()),
560        ),
561        Event::FileEditApplied {
562            turn_id,
563            flow_run_id,
564            ..
565        } => (
566            turn_id.as_ref().map(|t| t.0.to_string()),
567            flow_run_id.as_ref().map(|r| r.0.to_string()),
568        ),
569        Event::ContentFilterHit {
570            turn_id,
571            flow_run_id,
572            ..
573        }
574        | Event::ContextTruncated {
575            turn_id,
576            flow_run_id,
577            ..
578        }
579        | Event::WatchWarn {
580            turn_id,
581            flow_run_id,
582            ..
583        } => (
584            turn_id.as_ref().map(|t| t.0.to_string()),
585            flow_run_id.as_ref().map(|r| r.0.to_string()),
586        ),
587        Event::LlmPartialCall {
588            turn_id,
589            flow_run_id,
590            ..
591        } => (
592            turn_id.as_ref().map(|t| t.0.to_string()),
593            flow_run_id.as_ref().map(|r| r.0.to_string()),
594        ),
595        Event::FlowGraph { run_id, .. }
596        | Event::FlowNodeStart { run_id, .. }
597        | Event::FlowNodeEnd { run_id, .. }
598        | Event::ToolNode { run_id, .. }
599        | Event::ToolPendingApproval { run_id, .. }
600        | Event::ToolApproved { run_id, .. }
601        | Event::ToolDenied { run_id, .. } => (None, Some(run_id.0.to_string())),
602        Event::PermissionRequestCreated { payload }
603        | Event::PermissionRequestTargeted { payload }
604        | Event::PermissionRequestDeferred { payload }
605        | Event::PermissionRequestApproved { payload }
606        | Event::PermissionRequestDenied { payload }
607        | Event::PermissionRequestCancelled { payload }
608        | Event::UnrestrictedExecution { payload } => {
609            (None, Some(payload.requesting_run_id.0.to_string()))
610        }
611        Event::PermissionGroupCreated { payload }
612        | Event::PermissionGroupUpdated { payload }
613        | Event::PermissionGroupResolved { payload } => {
614            let anchor = match &payload.owner {
615                crate::permission_audit::PermissionGroupAuditOwner::Flow { run_id } => {
616                    run_id.to_string()
617                }
618                crate::permission_audit::PermissionGroupAuditOwner::User { session_id } => {
619                    format!("user:{session_id}")
620                }
621                crate::permission_audit::PermissionGroupAuditOwner::System => "system".into(),
622            };
623            (None, Some(anchor))
624        }
625        Event::PermissionGrantCreated { payload } | Event::PermissionGrantExpired { payload } => {
626            (None, Some(payload.requesting_run_id.0.to_string()))
627        }
628        Event::AttachmentDegraded {
629            turn_id,
630            flow_run_id,
631            ..
632        } => (
633            turn_id.as_ref().map(|t| t.0.to_string()),
634            flow_run_id.as_ref().map(|r| r.0.to_string()),
635        ),
636        Event::CompactionSummary { flow_run_id, .. }
637        | Event::ContextCompact { flow_run_id, .. }
638        | Event::Checkpoint { flow_run_id, .. } => (
639            None,
640            flow_run_id.as_ref().map(|run_id| run_id.0.to_string()),
641        ),
642        Event::LlmCall { .. }
643        | Event::PendingPrompt { .. }
644        | Event::PromptExpired { .. }
645        | Event::PromptResolved { .. }
646        | Event::DeferredFormRecorded { .. }
647        | Event::TerminalFinalState { .. }
648        | Event::MermaidDiagram { .. } => (None, None),
649        Event::DeferredFormApplied { message, .. } => (Some(message.turn_id.to_string()), None),
650    }
651}
652
653pub(crate) fn extract_text_content(event: &Event) -> Option<String> {
654    match event {
655        Event::UserMsg { message, .. }
656        | Event::AssistantMsg { message, .. }
657        | Event::ToolResultMsg { message, .. }
658        | Event::SystemMsg { message, .. } => Some(message.text_concat()),
659        Event::DeferredFormApplied { message, .. } => Some(message.text_concat()),
660        Event::WatchWarn { message, .. } => Some(message.clone()),
661        Event::CompactionSummary { summary, .. } => Some(summary.clone()),
662        Event::AttachmentDegraded {
663            file_basename,
664            reason,
665            ..
666        } => Some(format!("{file_basename} {reason}")),
667        _ => None,
668    }
669}
670
671#[cfg(test)]
672mod tests {
673    use super::*;
674    use crate::event::{Event, FlowRunId, FlowStatus};
675    use tempfile::TempDir;
676
677    fn flow_start(seq: u64) -> EventEnvelope {
678        EventEnvelope::new(
679            seq,
680            Event::FlowStart {
681                run_id: FlowRunId::now(),
682                flow_name: format!("flow_{seq}"),
683                parent_run_id: None,
684                parent_node_id: None,
685                spawned: false,
686            },
687        )
688    }
689
690    #[test]
691    fn every_owned_message_role_keeps_its_flow_anchor() {
692        let turn_id = crate::event::TurnId::now();
693        let flow_run_id = FlowRunId::now();
694        let user = crate::message::Message::user_text(turn_id.clone(), "delegated");
695        let system = crate::message::Message::system_text(turn_id.clone(), "handoff");
696
697        for event in [
698            Event::UserMsg {
699                turn_id: turn_id.clone(),
700                flow_run_id: Some(flow_run_id.clone()),
701                message: user,
702            },
703            Event::SystemMsg {
704                turn_id: turn_id.clone(),
705                flow_run_id: Some(flow_run_id.clone()),
706                message: system,
707            },
708        ] {
709            assert_eq!(
710                extract_anchors(&event),
711                (Some(turn_id.to_string()), Some(flow_run_id.to_string()))
712            );
713        }
714    }
715
716    #[test]
717    fn every_owned_compaction_event_keeps_its_flow_anchor() {
718        let flow_run_id = FlowRunId::now();
719        for event in [
720            Event::CompactionSummary {
721                session_id: "session".into(),
722                flow_run_id: Some(flow_run_id.clone()),
723                range_start: 0,
724                range_end: 1,
725                compacted_count: 2,
726                before_tokens: 100,
727                after_tokens: 10,
728                summary: "summary".into(),
729            },
730            Event::ContextCompact {
731                session_id: "session".into(),
732                flow_run_id: Some(flow_run_id.clone()),
733                before_tokens: 100,
734                after_tokens: 10,
735                compacted_range_start: 0,
736                compacted_range_end: 1,
737                summary_text: Some("summary".into()),
738                replacement_msg_seq: None,
739            },
740            Event::Checkpoint {
741                session_id: "session".into(),
742                flow_run_id: Some(flow_run_id.clone()),
743                messages: Vec::new(),
744                window_tokens: 10,
745            },
746        ] {
747            assert_eq!(
748                extract_anchors(&event),
749                (None, Some(flow_run_id.to_string()))
750            );
751        }
752    }
753
754    #[test]
755    fn degraded_buffer_keeps_later_events_and_drops_oldest_at_capacity() {
756        let mut degraded = DegradedBuffer::new(2);
757
758        degraded.degrade(flow_start(1), "disk full");
759        degraded.buffer(flow_start(2));
760
761        assert!(degraded.is_degraded());
762        assert_eq!(
763            degraded
764                .events
765                .iter()
766                .map(|event| event.seq)
767                .collect::<Vec<_>>(),
768            vec![1, 2]
769        );
770        assert_eq!(degraded.dropped, 0);
771
772        degraded.buffer(flow_start(3));
773
774        assert_eq!(
775            degraded
776                .events
777                .iter()
778                .map(|event| event.seq)
779                .collect::<Vec<_>>(),
780            vec![2, 3]
781        );
782        assert_eq!(degraded.dropped, 1);
783        assert_eq!(degraded.error.as_deref(), Some("disk full"));
784    }
785
786    #[tokio::test]
787    async fn recovery_tick_retries_buffered_events() {
788        let dir = TempDir::new().unwrap();
789        let path = dir.path().join("events.jsonl");
790        tokio::fs::write(&path, b"").await.unwrap();
791        let mut file = tokio::fs::File::open(&path).await.unwrap();
792        let mut offset = 0;
793        let mut degraded = DegradedBuffer::new(MAX_BUFFERED);
794
795        handle_event(
796            &mut file,
797            &mut offset,
798            flow_start(1),
799            &mut degraded,
800            None,
801            None,
802            dir.path(),
803        )
804        .await;
805
806        assert!(degraded.is_degraded());
807        assert_eq!(degraded.events.len(), 1);
808
809        drop(file);
810        #[allow(clippy::suspicious_open_options)]
811        let mut file = tokio::fs::OpenOptions::new()
812            .read(true)
813            .write(true)
814            .open(&path)
815            .await
816            .unwrap();
817        let mut offset = 0;
818        retry_degraded_on_tick(
819            &mut file,
820            &mut offset,
821            &mut degraded,
822            None,
823            None,
824            dir.path(),
825        )
826        .await;
827
828        assert!(!degraded.is_degraded());
829        assert!(degraded.events.is_empty());
830        let contents = tokio::fs::read_to_string(path).await.unwrap();
831        let event: serde_json::Value = serde_json::from_str(contents.trim()).unwrap();
832        assert_eq!(event["seq"], 1);
833        assert_eq!(event["type"], "flow_start");
834    }
835
836    #[tokio::test]
837    async fn write_event_overwrites_partial_line() {
838        let dir = TempDir::new().unwrap();
839        let path = dir.path().join("events.jsonl");
840        let mut file = tokio::fs::OpenOptions::new()
841            .create(true)
842            .truncate(false)
843            .read(true)
844            .write(true)
845            .open(&path)
846            .await
847            .unwrap();
848        let event = flow_start(1);
849        let line = serialize_event(&event, None);
850        file.write_all(&line.as_bytes()[..line.len() / 2])
851            .await
852            .unwrap();
853        let mut offset = 0;
854
855        write_event(&mut file, &mut offset, &event, None, None, dir.path())
856            .await
857            .unwrap();
858
859        let contents = tokio::fs::read_to_string(path).await.unwrap();
860        assert_eq!(contents.lines().count(), 1);
861        let _: serde_json::Value = serde_json::from_str(contents.trim()).unwrap();
862    }
863
864    #[tokio::test]
865    async fn writer_appends_events_as_jsonl() {
866        let dir = TempDir::new().unwrap();
867        let writer = EventWriter::spawn(dir.path()).unwrap();
868        let tx = writer.sender();
869        for i in 0..5 {
870            tx.send(EventEnvelope::new(
871                i as u64,
872                Event::FlowStart {
873                    run_id: FlowRunId::now(),
874                    flow_name: format!("flow_{i}"),
875                    parent_run_id: None,
876                    parent_node_id: None,
877                    spawned: false,
878                },
879            ))
880            .unwrap();
881        }
882        drop(tx);
883        writer.shutdown().await;
884        let path = dir.path().join("events.jsonl");
885        let contents = tokio::fs::read_to_string(&path).await.unwrap();
886        let lines: Vec<_> = contents.lines().collect();
887        assert_eq!(lines.len(), 5);
888        for line in lines {
889            let v: serde_json::Value = serde_json::from_str(line).unwrap();
890            assert_eq!(v["type"], "flow_start");
891            assert!(v["run_id"].is_string());
892            assert!(v["flow_name"].is_string());
893        }
894        let stats = crate::session_meta::SessionStats::load_or_rebuild(dir.path()).unwrap();
895        assert_eq!(stats.event_count, 5);
896        assert_eq!(stats.message_count, 0);
897        assert_eq!(stats.event_bytes, std::fs::metadata(path).unwrap().len());
898    }
899
900    #[tokio::test]
901    async fn writer_writes_to_project_index_with_session_id() {
902        let session_dir = TempDir::new().unwrap();
903        let project_dir = TempDir::new().unwrap();
904        let idx = Arc::new(AnchorIndex::open_project(project_dir.path()).unwrap());
905        let writer = EventWriter::spawn_full(
906            session_dir.path(),
907            None,
908            Some(idx.clone()),
909            Some("sess-x".into()),
910        )
911        .unwrap();
912        let tx = writer.sender();
913        for i in 0..3 {
914            tx.send(EventEnvelope::new(
915                (i + 1) as u64,
916                Event::FlowStart {
917                    run_id: FlowRunId::now(),
918                    flow_name: format!("flow_{i}"),
919                    parent_run_id: None,
920                    parent_node_id: None,
921                    spawned: false,
922                },
923            ))
924            .unwrap();
925        }
926        drop(tx);
927        writer.shutdown().await;
928
929        let jsonl_lines = tokio::fs::read_to_string(session_dir.path().join("events.jsonl"))
930            .await
931            .unwrap()
932            .lines()
933            .count();
934        assert_eq!(jsonl_lines, 3);
935
936        let conn = idx.conn();
937        let count: i64 = conn
938            .query_row(
939                "SELECT COUNT(*) FROM events WHERE session_id = ?",
940                rusqlite::params!["sess-x"],
941                |r| r.get(0),
942            )
943            .unwrap();
944        assert_eq!(count, 3);
945    }
946
947    #[tokio::test]
948    async fn writer_indexes_user_msg_text_for_project_fts() {
949        use crate::event::TurnId;
950        use crate::message::Message;
951
952        let session_dir = TempDir::new().unwrap();
953        let project_dir = TempDir::new().unwrap();
954        let idx = Arc::new(AnchorIndex::open_project(project_dir.path()).unwrap());
955        let writer = EventWriter::spawn_full(
956            session_dir.path(),
957            None,
958            Some(idx.clone()),
959            Some("sess-x".into()),
960        )
961        .unwrap();
962        let tid = TurnId::now();
963        writer
964            .sender()
965            .send(EventEnvelope::new(
966                1,
967                Event::UserMsg {
968                    turn_id: tid.clone(),
969                    flow_run_id: None,
970                    message: Message::user_text(tid, "sqlite fts full text search"),
971                },
972            ))
973            .unwrap();
974        writer.shutdown().await;
975
976        let hits = idx
977            .fts_search_project_events("sqlite", Some("sess-x"), 10)
978            .unwrap();
979        assert_eq!(hits.len(), 1);
980        assert_eq!(hits[0].session_id, "sess-x");
981        assert_eq!(hits[0].seq, 1);
982        let stats = crate::session_meta::SessionStats::load_or_rebuild(session_dir.path()).unwrap();
983        assert_eq!(stats.event_count, 1);
984        assert_eq!(stats.user_message_count, 1);
985        assert_eq!(stats.message_count, 1);
986    }
987
988    #[tokio::test]
989    async fn writer_serializes_flow_end_with_status() {
990        let dir = TempDir::new().unwrap();
991        let writer = EventWriter::spawn(dir.path()).unwrap();
992        writer
993            .sender()
994            .send(EventEnvelope::new(
995                0,
996                Event::FlowEnd {
997                    run_id: FlowRunId::now(),
998                    flow_name: "t".into(),
999                    status: FlowStatus::Errored {
1000                        message: "boom".into(),
1001                    },
1002                },
1003            ))
1004            .unwrap();
1005        writer.shutdown().await;
1006        let contents = tokio::fs::read_to_string(dir.path().join("events.jsonl"))
1007            .await
1008            .unwrap();
1009        let v: serde_json::Value = serde_json::from_str(contents.trim()).unwrap();
1010        assert_eq!(v["type"], "flow_end");
1011        assert_eq!(v["status"]["kind"], "errored");
1012    }
1013}