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 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::PromptResolved { .. } => "prompt_resolved",
487 Event::LlmPartialCall { .. } => "llm_partial_call",
488 Event::FlowGraph { .. } => "flow_graph",
489 Event::FlowNodeStart { .. } => "flow_node_start",
490 Event::FlowNodeEnd { .. } => "flow_node_end",
491 Event::ToolNode { .. } => "tool_node",
492 Event::AttachmentDegraded { .. } => "attachment_degraded",
493 Event::ToolPendingApproval { .. } => "tool_pending_approval",
494 Event::ToolApproved { .. } => "tool_approved",
495 Event::ToolDenied { .. } => "tool_denied",
496 Event::PermissionRequestCreated { .. } => "permission_request_created",
497 Event::PermissionRequestTargeted { .. } => "permission_request_targeted",
498 Event::PermissionRequestDeferred { .. } => "permission_request_deferred",
499 Event::PermissionRequestApproved { .. } => "permission_request_approved",
500 Event::PermissionRequestDenied { .. } => "permission_request_denied",
501 Event::PermissionRequestCancelled { .. } => "permission_request_cancelled",
502 Event::PermissionGroupCreated { .. } => "permission_group_created",
503 Event::PermissionGroupUpdated { .. } => "permission_group_updated",
504 Event::PermissionGroupResolved { .. } => "permission_group_resolved",
505 Event::PermissionGrantCreated { .. } => "permission_grant_created",
506 Event::PermissionGrantExpired { .. } => "permission_grant_expired",
507 Event::UnrestrictedExecution { .. } => "unrestricted_execution",
508 Event::TerminalFinalState { .. } => "terminal_final_state",
509 Event::MermaidDiagram { .. } => "mermaid_diagram",
510 }
511}
512
513pub(crate) fn extract_anchors(event: &Event) -> (Option<String>, Option<String>) {
514 match event {
515 Event::FlowStart { run_id, .. }
516 | Event::FlowEnd { run_id, .. }
517 | Event::WorkspaceLifecycle { run_id, .. } => (None, Some(run_id.0.to_string())),
518 Event::TurnStart { turn_id, .. } | Event::TurnEnd { turn_id, .. } => {
519 (Some(turn_id.0.to_string()), None)
520 }
521 Event::UserInject { turn_id, .. } => (Some(turn_id.0.to_string()), None),
522 Event::UserMsg {
523 turn_id,
524 flow_run_id,
525 ..
526 }
527 | Event::AssistantMsg {
528 turn_id,
529 flow_run_id,
530 ..
531 }
532 | Event::ToolResultMsg {
533 turn_id,
534 flow_run_id,
535 ..
536 }
537 | Event::ToolResultMetrics {
538 turn_id,
539 flow_run_id,
540 ..
541 }
542 | Event::SystemMsg {
543 turn_id,
544 flow_run_id,
545 ..
546 } => (
547 Some(turn_id.0.to_string()),
548 flow_run_id.as_ref().map(|r| r.0.to_string()),
549 ),
550 Event::DiffPreview {
551 turn_id,
552 flow_run_id,
553 ..
554 } => (
555 turn_id.as_ref().map(|t| t.0.to_string()),
556 flow_run_id.as_ref().map(|r| r.0.to_string()),
557 ),
558 Event::FileEditApplied {
559 turn_id,
560 flow_run_id,
561 ..
562 } => (
563 turn_id.as_ref().map(|t| t.0.to_string()),
564 flow_run_id.as_ref().map(|r| r.0.to_string()),
565 ),
566 Event::ContentFilterHit {
567 turn_id,
568 flow_run_id,
569 ..
570 }
571 | Event::ContextTruncated {
572 turn_id,
573 flow_run_id,
574 ..
575 }
576 | Event::WatchWarn {
577 turn_id,
578 flow_run_id,
579 ..
580 } => (
581 turn_id.as_ref().map(|t| t.0.to_string()),
582 flow_run_id.as_ref().map(|r| r.0.to_string()),
583 ),
584 Event::LlmPartialCall {
585 turn_id,
586 flow_run_id,
587 ..
588 } => (
589 turn_id.as_ref().map(|t| t.0.to_string()),
590 flow_run_id.as_ref().map(|r| r.0.to_string()),
591 ),
592 Event::FlowGraph { run_id, .. }
593 | Event::FlowNodeStart { run_id, .. }
594 | Event::FlowNodeEnd { run_id, .. }
595 | Event::ToolNode { run_id, .. }
596 | Event::ToolPendingApproval { run_id, .. }
597 | Event::ToolApproved { run_id, .. }
598 | Event::ToolDenied { run_id, .. } => (None, Some(run_id.0.to_string())),
599 Event::PermissionRequestCreated { payload }
600 | Event::PermissionRequestTargeted { payload }
601 | Event::PermissionRequestDeferred { payload }
602 | Event::PermissionRequestApproved { payload }
603 | Event::PermissionRequestDenied { payload }
604 | Event::PermissionRequestCancelled { payload }
605 | Event::UnrestrictedExecution { payload } => {
606 (None, Some(payload.requesting_run_id.0.to_string()))
607 }
608 Event::PermissionGroupCreated { payload }
609 | Event::PermissionGroupUpdated { payload }
610 | Event::PermissionGroupResolved { payload } => {
611 let anchor = match &payload.owner {
612 crate::permission_audit::PermissionGroupAuditOwner::Flow { run_id } => {
613 run_id.to_string()
614 }
615 crate::permission_audit::PermissionGroupAuditOwner::User { session_id } => {
616 format!("user:{session_id}")
617 }
618 crate::permission_audit::PermissionGroupAuditOwner::System => "system".into(),
619 };
620 (None, Some(anchor))
621 }
622 Event::PermissionGrantCreated { payload } | Event::PermissionGrantExpired { payload } => {
623 (None, Some(payload.requesting_run_id.0.to_string()))
624 }
625 Event::AttachmentDegraded {
626 turn_id,
627 flow_run_id,
628 ..
629 } => (
630 turn_id.as_ref().map(|t| t.0.to_string()),
631 flow_run_id.as_ref().map(|r| r.0.to_string()),
632 ),
633 Event::CompactionSummary { flow_run_id, .. }
634 | Event::ContextCompact { flow_run_id, .. }
635 | Event::Checkpoint { flow_run_id, .. } => (
636 None,
637 flow_run_id.as_ref().map(|run_id| run_id.0.to_string()),
638 ),
639 Event::LlmCall { .. }
640 | Event::PendingPrompt { .. }
641 | Event::PromptResolved { .. }
642 | Event::TerminalFinalState { .. }
643 | Event::MermaidDiagram { .. } => (None, None),
644 }
645}
646
647pub(crate) fn extract_text_content(event: &Event) -> Option<String> {
648 match event {
649 Event::UserMsg { message, .. }
650 | Event::AssistantMsg { message, .. }
651 | Event::ToolResultMsg { message, .. }
652 | Event::SystemMsg { message, .. } => Some(message.text_concat()),
653 Event::WatchWarn { message, .. } => Some(message.clone()),
654 Event::CompactionSummary { summary, .. } => Some(summary.clone()),
655 Event::AttachmentDegraded {
656 file_basename,
657 reason,
658 ..
659 } => Some(format!("{file_basename} {reason}")),
660 _ => None,
661 }
662}
663
664#[cfg(test)]
665mod tests {
666 use super::*;
667 use crate::event::{Event, FlowRunId, FlowStatus};
668 use tempfile::TempDir;
669
670 fn flow_start(seq: u64) -> EventEnvelope {
671 EventEnvelope::new(
672 seq,
673 Event::FlowStart {
674 run_id: FlowRunId::now(),
675 flow_name: format!("flow_{seq}"),
676 parent_run_id: None,
677 parent_node_id: None,
678 spawned: false,
679 },
680 )
681 }
682
683 #[test]
684 fn every_owned_message_role_keeps_its_flow_anchor() {
685 let turn_id = crate::event::TurnId::now();
686 let flow_run_id = FlowRunId::now();
687 let user = crate::message::Message::user_text(turn_id.clone(), "delegated");
688 let system = crate::message::Message::system_text(turn_id.clone(), "handoff");
689
690 for event in [
691 Event::UserMsg {
692 turn_id: turn_id.clone(),
693 flow_run_id: Some(flow_run_id.clone()),
694 message: user,
695 },
696 Event::SystemMsg {
697 turn_id: turn_id.clone(),
698 flow_run_id: Some(flow_run_id.clone()),
699 message: system,
700 },
701 ] {
702 assert_eq!(
703 extract_anchors(&event),
704 (Some(turn_id.to_string()), Some(flow_run_id.to_string()))
705 );
706 }
707 }
708
709 #[test]
710 fn every_owned_compaction_event_keeps_its_flow_anchor() {
711 let flow_run_id = FlowRunId::now();
712 for event in [
713 Event::CompactionSummary {
714 session_id: "session".into(),
715 flow_run_id: Some(flow_run_id.clone()),
716 range_start: 0,
717 range_end: 1,
718 compacted_count: 2,
719 before_tokens: 100,
720 after_tokens: 10,
721 summary: "summary".into(),
722 },
723 Event::ContextCompact {
724 session_id: "session".into(),
725 flow_run_id: Some(flow_run_id.clone()),
726 before_tokens: 100,
727 after_tokens: 10,
728 compacted_range_start: 0,
729 compacted_range_end: 1,
730 summary_text: Some("summary".into()),
731 replacement_msg_seq: None,
732 },
733 Event::Checkpoint {
734 session_id: "session".into(),
735 flow_run_id: Some(flow_run_id.clone()),
736 messages: Vec::new(),
737 window_tokens: 10,
738 },
739 ] {
740 assert_eq!(
741 extract_anchors(&event),
742 (None, Some(flow_run_id.to_string()))
743 );
744 }
745 }
746
747 #[test]
748 fn degraded_buffer_keeps_later_events_and_drops_oldest_at_capacity() {
749 let mut degraded = DegradedBuffer::new(2);
750
751 degraded.degrade(flow_start(1), "disk full");
752 degraded.buffer(flow_start(2));
753
754 assert!(degraded.is_degraded());
755 assert_eq!(
756 degraded
757 .events
758 .iter()
759 .map(|event| event.seq)
760 .collect::<Vec<_>>(),
761 vec![1, 2]
762 );
763 assert_eq!(degraded.dropped, 0);
764
765 degraded.buffer(flow_start(3));
766
767 assert_eq!(
768 degraded
769 .events
770 .iter()
771 .map(|event| event.seq)
772 .collect::<Vec<_>>(),
773 vec![2, 3]
774 );
775 assert_eq!(degraded.dropped, 1);
776 assert_eq!(degraded.error.as_deref(), Some("disk full"));
777 }
778
779 #[tokio::test]
780 async fn recovery_tick_retries_buffered_events() {
781 let dir = TempDir::new().unwrap();
782 let path = dir.path().join("events.jsonl");
783 tokio::fs::write(&path, b"").await.unwrap();
784 let mut file = tokio::fs::File::open(&path).await.unwrap();
785 let mut offset = 0;
786 let mut degraded = DegradedBuffer::new(MAX_BUFFERED);
787
788 handle_event(
789 &mut file,
790 &mut offset,
791 flow_start(1),
792 &mut degraded,
793 None,
794 None,
795 dir.path(),
796 )
797 .await;
798
799 assert!(degraded.is_degraded());
800 assert_eq!(degraded.events.len(), 1);
801
802 drop(file);
803 #[allow(clippy::suspicious_open_options)]
804 let mut file = tokio::fs::OpenOptions::new()
805 .read(true)
806 .write(true)
807 .open(&path)
808 .await
809 .unwrap();
810 let mut offset = 0;
811 retry_degraded_on_tick(
812 &mut file,
813 &mut offset,
814 &mut degraded,
815 None,
816 None,
817 dir.path(),
818 )
819 .await;
820
821 assert!(!degraded.is_degraded());
822 assert!(degraded.events.is_empty());
823 let contents = tokio::fs::read_to_string(path).await.unwrap();
824 let event: serde_json::Value = serde_json::from_str(contents.trim()).unwrap();
825 assert_eq!(event["seq"], 1);
826 assert_eq!(event["type"], "flow_start");
827 }
828
829 #[tokio::test]
830 async fn write_event_overwrites_partial_line() {
831 let dir = TempDir::new().unwrap();
832 let path = dir.path().join("events.jsonl");
833 let mut file = tokio::fs::OpenOptions::new()
834 .create(true)
835 .truncate(false)
836 .read(true)
837 .write(true)
838 .open(&path)
839 .await
840 .unwrap();
841 let event = flow_start(1);
842 let line = serialize_event(&event, None);
843 file.write_all(&line.as_bytes()[..line.len() / 2])
844 .await
845 .unwrap();
846 let mut offset = 0;
847
848 write_event(&mut file, &mut offset, &event, None, None, dir.path())
849 .await
850 .unwrap();
851
852 let contents = tokio::fs::read_to_string(path).await.unwrap();
853 assert_eq!(contents.lines().count(), 1);
854 let _: serde_json::Value = serde_json::from_str(contents.trim()).unwrap();
855 }
856
857 #[tokio::test]
858 async fn writer_appends_events_as_jsonl() {
859 let dir = TempDir::new().unwrap();
860 let writer = EventWriter::spawn(dir.path()).unwrap();
861 let tx = writer.sender();
862 for i in 0..5 {
863 tx.send(EventEnvelope::new(
864 i as u64,
865 Event::FlowStart {
866 run_id: FlowRunId::now(),
867 flow_name: format!("flow_{i}"),
868 parent_run_id: None,
869 parent_node_id: None,
870 spawned: false,
871 },
872 ))
873 .unwrap();
874 }
875 drop(tx);
876 writer.shutdown().await;
877 let path = dir.path().join("events.jsonl");
878 let contents = tokio::fs::read_to_string(&path).await.unwrap();
879 let lines: Vec<_> = contents.lines().collect();
880 assert_eq!(lines.len(), 5);
881 for line in lines {
882 let v: serde_json::Value = serde_json::from_str(line).unwrap();
883 assert_eq!(v["type"], "flow_start");
884 assert!(v["run_id"].is_string());
885 assert!(v["flow_name"].is_string());
886 }
887 let stats = crate::session_meta::SessionStats::load_or_rebuild(dir.path()).unwrap();
888 assert_eq!(stats.event_count, 5);
889 assert_eq!(stats.message_count, 0);
890 assert_eq!(stats.event_bytes, std::fs::metadata(path).unwrap().len());
891 }
892
893 #[tokio::test]
894 async fn writer_writes_to_project_index_with_session_id() {
895 let session_dir = TempDir::new().unwrap();
896 let project_dir = TempDir::new().unwrap();
897 let idx = Arc::new(AnchorIndex::open_project(project_dir.path()).unwrap());
898 let writer = EventWriter::spawn_full(
899 session_dir.path(),
900 None,
901 Some(idx.clone()),
902 Some("sess-x".into()),
903 )
904 .unwrap();
905 let tx = writer.sender();
906 for i in 0..3 {
907 tx.send(EventEnvelope::new(
908 (i + 1) as u64,
909 Event::FlowStart {
910 run_id: FlowRunId::now(),
911 flow_name: format!("flow_{i}"),
912 parent_run_id: None,
913 parent_node_id: None,
914 spawned: false,
915 },
916 ))
917 .unwrap();
918 }
919 drop(tx);
920 writer.shutdown().await;
921
922 let jsonl_lines = tokio::fs::read_to_string(session_dir.path().join("events.jsonl"))
923 .await
924 .unwrap()
925 .lines()
926 .count();
927 assert_eq!(jsonl_lines, 3);
928
929 let conn = idx.conn();
930 let count: i64 = conn
931 .query_row(
932 "SELECT COUNT(*) FROM events WHERE session_id = ?",
933 rusqlite::params!["sess-x"],
934 |r| r.get(0),
935 )
936 .unwrap();
937 assert_eq!(count, 3);
938 }
939
940 #[tokio::test]
941 async fn writer_indexes_user_msg_text_for_project_fts() {
942 use crate::event::TurnId;
943 use crate::message::Message;
944
945 let session_dir = TempDir::new().unwrap();
946 let project_dir = TempDir::new().unwrap();
947 let idx = Arc::new(AnchorIndex::open_project(project_dir.path()).unwrap());
948 let writer = EventWriter::spawn_full(
949 session_dir.path(),
950 None,
951 Some(idx.clone()),
952 Some("sess-x".into()),
953 )
954 .unwrap();
955 let tid = TurnId::now();
956 writer
957 .sender()
958 .send(EventEnvelope::new(
959 1,
960 Event::UserMsg {
961 turn_id: tid.clone(),
962 flow_run_id: None,
963 message: Message::user_text(tid, "sqlite fts full text search"),
964 },
965 ))
966 .unwrap();
967 writer.shutdown().await;
968
969 let hits = idx
970 .fts_search_project_events("sqlite", Some("sess-x"), 10)
971 .unwrap();
972 assert_eq!(hits.len(), 1);
973 assert_eq!(hits[0].session_id, "sess-x");
974 assert_eq!(hits[0].seq, 1);
975 let stats = crate::session_meta::SessionStats::load_or_rebuild(session_dir.path()).unwrap();
976 assert_eq!(stats.event_count, 1);
977 assert_eq!(stats.user_message_count, 1);
978 assert_eq!(stats.message_count, 1);
979 }
980
981 #[tokio::test]
982 async fn writer_serializes_flow_end_with_status() {
983 let dir = TempDir::new().unwrap();
984 let writer = EventWriter::spawn(dir.path()).unwrap();
985 writer
986 .sender()
987 .send(EventEnvelope::new(
988 0,
989 Event::FlowEnd {
990 run_id: FlowRunId::now(),
991 flow_name: "t".into(),
992 status: FlowStatus::Errored {
993 message: "boom".into(),
994 },
995 },
996 ))
997 .unwrap();
998 writer.shutdown().await;
999 let contents = tokio::fs::read_to_string(dir.path().join("events.jsonl"))
1000 .await
1001 .unwrap();
1002 let v: serde_json::Value = serde_json::from_str(contents.trim()).unwrap();
1003 assert_eq!(v["type"], "flow_end");
1004 assert_eq!(v["status"]["kind"], "errored");
1005 }
1006}