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 indexer = project_index.zip(session_id);
187    let mut degraded = DegradedBuffer::new(MAX_BUFFERED);
188    let mut recovery_tick = tokio::time::interval(RECOVERY_INTERVAL);
189
190    loop {
191        tokio::select! {
192            biased;
193            _ = &mut stop_rx => {
194                while let Ok(event) = rx.try_recv() {
195                    handle_event(
196                        &mut file,
197                        &mut offset,
198                        event,
199                        &mut degraded,
200                        indexer.as_ref(),
201                        redactor.as_deref(),
202                    ).await;
203                }
204                while let Ok(waiter) = flush_rx.try_recv() {
205                    let _ = waiter.send(());
206                }
207                break;
208            }
209            _ = recovery_tick.tick() => {
210                retry_degraded_on_tick(
211                    &mut file,
212                    &mut offset,
213                    &mut degraded,
214                    indexer.as_ref(),
215                    redactor.as_deref(),
216                ).await;
217            }
218            maybe_event = rx.recv() => {
219                match maybe_event {
220                    Some(event) => {
221                        handle_event(
222                            &mut file,
223                            &mut offset,
224                            event,
225                            &mut degraded,
226                            indexer.as_ref(),
227                            redactor.as_deref(),
228                        ).await;
229                    }
230                    None => break,
231                }
232            }
233            maybe_flush = flush_rx.recv() => {
234                match maybe_flush {
235                    Some(waiter) => {
236                        while let Ok(event) = rx.try_recv() {
237                            handle_event(
238                                &mut file,
239                                &mut offset,
240                                event,
241                                &mut degraded,
242                                indexer.as_ref(),
243                                redactor.as_deref(),
244                            ).await;
245                        }
246                        retry_buffered(
247                            &mut file,
248                            &mut offset,
249                            &mut degraded,
250                            indexer.as_ref(),
251                            redactor.as_deref(),
252                        ).await;
253                        if let Err(e) = file.sync_data().await {
254                            crate::notify!(error, "event writer flush failed: {e}");
255                        }
256                        let _ = waiter.send(());
257                    }
258                    None => break,
259                }
260            }
261        }
262    }
263    retry_buffered(
264        &mut file,
265        &mut offset,
266        &mut degraded,
267        indexer.as_ref(),
268        redactor.as_deref(),
269    )
270    .await;
271    if !degraded.events.is_empty() {
272        crate::notify!(
273            error,
274            "event writer stopped with {} buffered event(s) not persisted (dropped {}): {}",
275            degraded.events.len(),
276            degraded.dropped,
277            degraded.error.as_deref().unwrap_or("write failed")
278        );
279    }
280    if let Err(e) = file.sync_data().await {
281        crate::notify!(error, "event writer final sync failed: {e}");
282    }
283    Ok(())
284}
285
286async fn handle_event(
287    file: &mut tokio::fs::File,
288    offset: &mut u64,
289    event: EventEnvelope,
290    degraded: &mut DegradedBuffer,
291    indexer: Option<&(Arc<AnchorIndex>, String)>,
292    redactor: Option<&Redactor>,
293) {
294    if degraded.is_degraded() {
295        degraded.buffer(event);
296        return;
297    }
298    if let Err(e) = write_event(file, offset, &event, indexer, redactor).await {
299        crate::notify!(
300            error,
301            "event writer write failed (seq={}): {e}; buffering events in memory",
302            event.seq
303        );
304        degraded.degrade(event, e.to_string());
305    }
306}
307
308async fn retry_degraded_on_tick(
309    file: &mut tokio::fs::File,
310    offset: &mut u64,
311    degraded: &mut DegradedBuffer,
312    indexer: Option<&(Arc<AnchorIndex>, String)>,
313    redactor: Option<&Redactor>,
314) {
315    if !degraded.is_degraded() || degraded.events.is_empty() {
316        return;
317    }
318
319    retry_buffered(file, offset, degraded, indexer, redactor).await;
320    if !degraded.is_degraded() {
321        crate::notify!(info, location = Status, "事件写入已恢复");
322    }
323}
324
325async fn retry_buffered(
326    file: &mut tokio::fs::File,
327    offset: &mut u64,
328    degraded: &mut DegradedBuffer,
329    indexer: Option<&(Arc<AnchorIndex>, String)>,
330    redactor: Option<&Redactor>,
331) {
332    while let Some(event) = degraded.events.pop_front() {
333        if let Err(e) = write_event(file, offset, &event, indexer, redactor).await {
334            crate::notify!(
335                error,
336                "event writer retry failed (seq={}): {e}; {} event(s) remain buffered",
337                event.seq,
338                degraded.events.len() + 1
339            );
340            degraded.events.push_front(event);
341            degraded.error = Some(e.to_string());
342            return;
343        }
344    }
345    degraded.recover();
346}
347
348async fn write_event(
349    file: &mut tokio::fs::File,
350    offset: &mut u64,
351    envelope: &EventEnvelope,
352    indexer: Option<&(Arc<AnchorIndex>, String)>,
353    redactor: Option<&Redactor>,
354) -> std::io::Result<()> {
355    let line = serialize_event(envelope, redactor);
356    let start = *offset;
357    file.seek(SeekFrom::Start(start)).await?;
358    if let Err(e) = file.write_all(line.as_bytes()).await {
359        *offset = start;
360        return Err(e);
361    }
362    if let Err(e) = file.write_all(b"\n").await {
363        *offset = start;
364        return Err(e);
365    }
366    let end = file.stream_position().await?;
367    if let Err(e) = file.sync_data().await {
368        *offset = start;
369        return Err(e);
370    }
371    *offset = end;
372    if let Some((idx, sid)) = indexer
373        && let Err(e) = insert_project_row(idx, sid, envelope, &line)
374    {
375        crate::notify!(
376            warn,
377            location = Log,
378            stack = merge_count("project_index.insert_failed", 60_000),
379            "project index insert failed (seq={}): {e}",
380            envelope.seq
381        );
382    }
383    Ok(())
384}
385
386fn serialize_event(envelope: &EventEnvelope, redactor: Option<&Redactor>) -> String {
387    let Some(r) = redactor else {
388        return serde_json::to_string(envelope).unwrap_or_else(|e| {
389            format!(
390                "{{\"type\":\"encode_error\",\"error\":{:?}}}",
391                e.to_string()
392            )
393        });
394    };
395    let mut value = match serde_json::to_value(envelope) {
396        Ok(v) => v,
397        Err(e) => {
398            return format!(
399                "{{\"type\":\"encode_error\",\"error\":{:?}}}",
400                e.to_string()
401            );
402        }
403    };
404    r.redact_json(&mut value);
405    serde_json::to_string(&value).unwrap_or_else(|e| {
406        format!(
407            "{{\"type\":\"encode_error\",\"error\":{:?}}}",
408            e.to_string()
409        )
410    })
411}
412
413fn insert_project_row(
414    index: &AnchorIndex,
415    session_id: &str,
416    envelope: &EventEnvelope,
417    payload_json: &str,
418) -> rusqlite::Result<()> {
419    let event = &envelope.event;
420    let ts = extract_ts(envelope);
421    let kind = event_kind(event);
422    let (turn_id, flow_run_id) = extract_anchors(event);
423    let text = extract_text_content(event).unwrap_or_default();
424    index.insert_project_event_raw(ProjectEventInsert {
425        session_id,
426        seq: envelope.seq as i64,
427        ts: &ts,
428        kind,
429        turn_id: turn_id.as_deref(),
430        flow_run_id: flow_run_id.as_deref(),
431        text_content: &text,
432        payload_json,
433    })?;
434    Ok(())
435}
436
437pub(crate) fn extract_ts(envelope: &EventEnvelope) -> String {
438    envelope.ts.to_rfc3339()
439}
440
441pub(crate) fn event_kind(event: &Event) -> &'static str {
442    match event {
443        Event::FlowStart { .. } => "flow_start",
444        Event::FlowEnd { .. } => "flow_end",
445        Event::LlmCall { .. } => "llm_call",
446        Event::TurnStart { .. } => "turn_start",
447        Event::TurnEnd { .. } => "turn_end",
448        Event::UserMsg { .. } => "user_msg",
449        Event::AssistantMsg { .. } => "assistant_msg",
450        Event::ToolResultMsg { .. } => "tool_result_msg",
451        Event::DiffPreview { .. } => "diff_preview",
452        Event::CompactionSummary { .. } => "compaction_summary",
453        Event::SystemMsg { .. } => "system_msg",
454        Event::UserInject { .. } => "user_inject",
455        Event::ContentFilterHit { .. } => "content_filter_hit",
456        Event::ContextCompact { .. } => "context_compact",
457        Event::Checkpoint { .. } => "checkpoint",
458        Event::ContextTruncated { .. } => "context_truncated",
459        Event::WatchWarn { .. } => "watch_warn",
460        Event::PendingPrompt { .. } => "pending_prompt",
461        Event::PromptResolved { .. } => "prompt_resolved",
462        Event::LlmPartialCall { .. } => "llm_partial_call",
463        Event::FlowGraph { .. } => "flow_graph",
464        Event::FlowNodeStart { .. } => "flow_node_start",
465        Event::FlowNodeEnd { .. } => "flow_node_end",
466        Event::ToolNode { .. } => "tool_node",
467        Event::AttachmentDegraded { .. } => "attachment_degraded",
468        Event::ToolPendingApproval { .. } => "tool_pending_approval",
469        Event::ToolApproved { .. } => "tool_approved",
470        Event::ToolDenied { .. } => "tool_denied",
471        Event::TerminalFinalState { .. } => "terminal_final_state",
472        Event::MermaidDiagram { .. } => "mermaid_diagram",
473    }
474}
475
476pub(crate) fn extract_anchors(event: &Event) -> (Option<String>, Option<String>) {
477    match event {
478        Event::FlowStart { run_id, .. } | Event::FlowEnd { run_id, .. } => {
479            (None, Some(run_id.0.to_string()))
480        }
481        Event::TurnStart { turn_id, .. } | Event::TurnEnd { turn_id, .. } => {
482            (Some(turn_id.0.to_string()), None)
483        }
484        Event::UserMsg { turn_id, .. }
485        | Event::SystemMsg { turn_id, .. }
486        | Event::UserInject { turn_id, .. } => (Some(turn_id.0.to_string()), None),
487        Event::AssistantMsg {
488            turn_id,
489            flow_run_id,
490            ..
491        }
492        | Event::ToolResultMsg {
493            turn_id,
494            flow_run_id,
495            ..
496        } => (
497            Some(turn_id.0.to_string()),
498            flow_run_id.as_ref().map(|r| r.0.to_string()),
499        ),
500        Event::DiffPreview {
501            turn_id,
502            flow_run_id,
503            ..
504        } => (
505            turn_id.as_ref().map(|t| t.0.to_string()),
506            flow_run_id.as_ref().map(|r| r.0.to_string()),
507        ),
508        Event::ContentFilterHit {
509            turn_id,
510            flow_run_id,
511            ..
512        }
513        | Event::ContextTruncated {
514            turn_id,
515            flow_run_id,
516            ..
517        }
518        | Event::WatchWarn {
519            turn_id,
520            flow_run_id,
521            ..
522        } => (
523            turn_id.as_ref().map(|t| t.0.to_string()),
524            flow_run_id.as_ref().map(|r| r.0.to_string()),
525        ),
526        Event::LlmPartialCall {
527            turn_id,
528            flow_run_id,
529            ..
530        } => (
531            turn_id.as_ref().map(|t| t.0.to_string()),
532            flow_run_id.as_ref().map(|r| r.0.to_string()),
533        ),
534        Event::FlowGraph { run_id, .. }
535        | Event::FlowNodeStart { run_id, .. }
536        | Event::FlowNodeEnd { run_id, .. }
537        | Event::ToolNode { run_id, .. }
538        | Event::ToolPendingApproval { run_id, .. }
539        | Event::ToolApproved { run_id, .. }
540        | Event::ToolDenied { run_id, .. } => (None, Some(run_id.0.to_string())),
541        Event::AttachmentDegraded {
542            turn_id,
543            flow_run_id,
544            ..
545        } => (
546            turn_id.as_ref().map(|t| t.0.to_string()),
547            flow_run_id.as_ref().map(|r| r.0.to_string()),
548        ),
549        Event::CompactionSummary { .. }
550        | Event::LlmCall { .. }
551        | Event::ContextCompact { .. }
552        | Event::Checkpoint { .. }
553        | Event::PendingPrompt { .. }
554        | Event::PromptResolved { .. }
555        | Event::TerminalFinalState { .. }
556        | Event::MermaidDiagram { .. } => (None, None),
557    }
558}
559
560pub(crate) fn extract_text_content(event: &Event) -> Option<String> {
561    match event {
562        Event::UserMsg { message, .. }
563        | Event::AssistantMsg { message, .. }
564        | Event::ToolResultMsg { message, .. }
565        | Event::SystemMsg { message, .. } => Some(message.text_concat()),
566        Event::WatchWarn { message, .. } => Some(message.clone()),
567        Event::CompactionSummary { summary, .. } => Some(summary.clone()),
568        Event::AttachmentDegraded {
569            file_basename,
570            reason,
571            ..
572        } => Some(format!("{file_basename} {reason}")),
573        _ => None,
574    }
575}
576
577#[cfg(test)]
578mod tests {
579    use super::*;
580    use crate::event::{Event, FlowRunId, FlowStatus};
581    use tempfile::TempDir;
582
583    fn flow_start(seq: u64) -> EventEnvelope {
584        EventEnvelope::new(
585            seq,
586            Event::FlowStart {
587                run_id: FlowRunId::now(),
588                flow_name: format!("flow_{seq}"),
589                parent_run_id: None,
590                parent_node_id: None,
591                spawned: false,
592            },
593        )
594    }
595
596    #[test]
597    fn degraded_buffer_keeps_later_events_and_drops_oldest_at_capacity() {
598        let mut degraded = DegradedBuffer::new(2);
599
600        degraded.degrade(flow_start(1), "disk full");
601        degraded.buffer(flow_start(2));
602
603        assert!(degraded.is_degraded());
604        assert_eq!(
605            degraded
606                .events
607                .iter()
608                .map(|event| event.seq)
609                .collect::<Vec<_>>(),
610            vec![1, 2]
611        );
612        assert_eq!(degraded.dropped, 0);
613
614        degraded.buffer(flow_start(3));
615
616        assert_eq!(
617            degraded
618                .events
619                .iter()
620                .map(|event| event.seq)
621                .collect::<Vec<_>>(),
622            vec![2, 3]
623        );
624        assert_eq!(degraded.dropped, 1);
625        assert_eq!(degraded.error.as_deref(), Some("disk full"));
626    }
627
628    #[tokio::test]
629    async fn recovery_tick_retries_buffered_events() {
630        let dir = TempDir::new().unwrap();
631        let path = dir.path().join("events.jsonl");
632        tokio::fs::write(&path, b"").await.unwrap();
633        let mut file = tokio::fs::File::open(&path).await.unwrap();
634        let mut offset = 0;
635        let mut degraded = DegradedBuffer::new(MAX_BUFFERED);
636
637        handle_event(
638            &mut file,
639            &mut offset,
640            flow_start(1),
641            &mut degraded,
642            None,
643            None,
644        )
645        .await;
646
647        assert!(degraded.is_degraded());
648        assert_eq!(degraded.events.len(), 1);
649
650        drop(file);
651        #[allow(clippy::suspicious_open_options)]
652        let mut file = tokio::fs::OpenOptions::new()
653            .read(true)
654            .write(true)
655            .open(&path)
656            .await
657            .unwrap();
658        let mut offset = 0;
659        retry_degraded_on_tick(&mut file, &mut offset, &mut degraded, None, None).await;
660
661        assert!(!degraded.is_degraded());
662        assert!(degraded.events.is_empty());
663        let contents = tokio::fs::read_to_string(path).await.unwrap();
664        let event: serde_json::Value = serde_json::from_str(contents.trim()).unwrap();
665        assert_eq!(event["seq"], 1);
666        assert_eq!(event["type"], "flow_start");
667    }
668
669    #[tokio::test]
670    async fn write_event_overwrites_partial_line() {
671        let dir = TempDir::new().unwrap();
672        let path = dir.path().join("events.jsonl");
673        let mut file = tokio::fs::OpenOptions::new()
674            .create(true)
675            .truncate(false)
676            .read(true)
677            .write(true)
678            .open(&path)
679            .await
680            .unwrap();
681        let event = flow_start(1);
682        let line = serialize_event(&event, None);
683        file.write_all(&line.as_bytes()[..line.len() / 2])
684            .await
685            .unwrap();
686        let mut offset = 0;
687
688        write_event(&mut file, &mut offset, &event, None, None)
689            .await
690            .unwrap();
691
692        let contents = tokio::fs::read_to_string(path).await.unwrap();
693        assert_eq!(contents.lines().count(), 1);
694        let _: serde_json::Value = serde_json::from_str(contents.trim()).unwrap();
695    }
696
697    #[tokio::test]
698    async fn writer_appends_events_as_jsonl() {
699        let dir = TempDir::new().unwrap();
700        let writer = EventWriter::spawn(dir.path()).unwrap();
701        let tx = writer.sender();
702        for i in 0..5 {
703            tx.send(EventEnvelope::new(
704                i as u64,
705                Event::FlowStart {
706                    run_id: FlowRunId::now(),
707                    flow_name: format!("flow_{i}"),
708                    parent_run_id: None,
709                    parent_node_id: None,
710                    spawned: false,
711                },
712            ))
713            .unwrap();
714        }
715        drop(tx);
716        writer.shutdown().await;
717        let path = dir.path().join("events.jsonl");
718        let contents = tokio::fs::read_to_string(&path).await.unwrap();
719        let lines: Vec<_> = contents.lines().collect();
720        assert_eq!(lines.len(), 5);
721        for line in lines {
722            let v: serde_json::Value = serde_json::from_str(line).unwrap();
723            assert_eq!(v["type"], "flow_start");
724            assert!(v["run_id"].is_string());
725            assert!(v["flow_name"].is_string());
726        }
727    }
728
729    #[tokio::test]
730    async fn writer_writes_to_project_index_with_session_id() {
731        let session_dir = TempDir::new().unwrap();
732        let project_dir = TempDir::new().unwrap();
733        let idx = Arc::new(AnchorIndex::open_project(project_dir.path()).unwrap());
734        let writer = EventWriter::spawn_full(
735            session_dir.path(),
736            None,
737            Some(idx.clone()),
738            Some("sess-x".into()),
739        )
740        .unwrap();
741        let tx = writer.sender();
742        for i in 0..3 {
743            tx.send(EventEnvelope::new(
744                (i + 1) as u64,
745                Event::FlowStart {
746                    run_id: FlowRunId::now(),
747                    flow_name: format!("flow_{i}"),
748                    parent_run_id: None,
749                    parent_node_id: None,
750                    spawned: false,
751                },
752            ))
753            .unwrap();
754        }
755        drop(tx);
756        writer.shutdown().await;
757
758        let jsonl_lines = tokio::fs::read_to_string(session_dir.path().join("events.jsonl"))
759            .await
760            .unwrap()
761            .lines()
762            .count();
763        assert_eq!(jsonl_lines, 3);
764
765        let conn = idx.conn();
766        let count: i64 = conn
767            .query_row(
768                "SELECT COUNT(*) FROM events WHERE session_id = ?",
769                rusqlite::params!["sess-x"],
770                |r| r.get(0),
771            )
772            .unwrap();
773        assert_eq!(count, 3);
774    }
775
776    #[tokio::test]
777    async fn writer_indexes_user_msg_text_for_project_fts() {
778        use crate::event::TurnId;
779        use crate::message::Message;
780
781        let session_dir = TempDir::new().unwrap();
782        let project_dir = TempDir::new().unwrap();
783        let idx = Arc::new(AnchorIndex::open_project(project_dir.path()).unwrap());
784        let writer = EventWriter::spawn_full(
785            session_dir.path(),
786            None,
787            Some(idx.clone()),
788            Some("sess-x".into()),
789        )
790        .unwrap();
791        let tid = TurnId::now();
792        writer
793            .sender()
794            .send(EventEnvelope::new(
795                1,
796                Event::UserMsg {
797                    turn_id: tid.clone(),
798                    flow_run_id: None,
799                    message: Message::user_text(tid, "sqlite fts full text search"),
800                },
801            ))
802            .unwrap();
803        writer.shutdown().await;
804
805        let hits = idx
806            .fts_search_project_events("sqlite", Some("sess-x"), 10)
807            .unwrap();
808        assert_eq!(hits.len(), 1);
809        assert_eq!(hits[0].session_id, "sess-x");
810        assert_eq!(hits[0].seq, 1);
811    }
812
813    #[tokio::test]
814    async fn writer_serializes_flow_end_with_status() {
815        let dir = TempDir::new().unwrap();
816        let writer = EventWriter::spawn(dir.path()).unwrap();
817        writer
818            .sender()
819            .send(EventEnvelope::new(
820                0,
821                Event::FlowEnd {
822                    run_id: FlowRunId::now(),
823                    flow_name: "t".into(),
824                    status: FlowStatus::Errored {
825                        message: "boom".into(),
826                    },
827                },
828            ))
829            .unwrap();
830        writer.shutdown().await;
831        let contents = tokio::fs::read_to_string(dir.path().join("events.jsonl"))
832            .await
833            .unwrap();
834        let v: serde_json::Value = serde_json::from_str(contents.trim()).unwrap();
835        assert_eq!(v["type"], "flow_end");
836        assert_eq!(v["status"]["kind"], "errored");
837    }
838}