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 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 #[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}