1pub mod harness_adapt;
7pub mod harness_metrics;
8pub mod observability;
9pub mod tool_receipts;
10
11pub use observability::{
12 evaluate_alerts, summarize, summarize_log, Alert, AlertKind, AlertThresholds, MetricsSummary,
13};
14
15use car_secrets::{
16 atomic_replace_private_file, create_private_file, open_private_append, revalidate_private_file,
17 revalidate_private_path,
18};
19use chrono::{DateTime, Utc};
20use serde::{Deserialize, Serialize};
21use serde_json::Value;
22use std::collections::{HashMap, HashSet, VecDeque};
23use std::fs;
24use std::future::Future;
25use std::io::{BufRead, BufReader, BufWriter, Read, Seek, SeekFrom, Write};
26use std::path::{Path, PathBuf};
27use std::pin::Pin;
28use std::sync::{mpsc, Arc, Condvar, Mutex, Weak};
29use std::task::{Context, Poll, Waker};
30use std::thread;
31use std::time::{Duration, Instant};
32use uuid::Uuid;
33
34#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
35#[serde(rename_all = "camelCase")]
36pub struct EventLogStats {
37 pub events: usize,
38 pub spans: usize,
39 pub approx_event_bytes: usize,
40 pub approx_span_bytes: usize,
41}
42
43#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize, schemars::JsonSchema)]
45#[serde(rename_all = "snake_case")]
46pub enum EventKind {
47 RunStarted,
49 RunCancellationRequested,
51 RunCancellationResult,
53 ProposalReceived,
54 ProposalCompleted,
57 ActionValidated,
58 ActionRejected,
59 ActionExecuting,
60 ActionSucceeded,
61 ActionFailed,
62 ActionSkipped,
63 ActionRetrying,
64 ActionDeduplicated,
65 PolicyViolation,
66 StateChanged,
67 StateSnapshot,
68 StateCommitted,
72 StateRollback,
73 SkillDistilled,
75 SkillEvolved,
76 SkillDeprecated,
77 EvolutionTriggered,
78 CandidatePromoted,
82 CandidateRejected,
85 Consolidated,
87 ProactiveMemoryMaintained,
92 ProactiveMemoryIntervention,
93 ReplanAttempted,
95 ReplanProposalReceived,
96 ReplanRejected,
97 ReplanExhausted,
98 VoiceFastTurnStarted,
102 VoiceFastTurnEnded,
103 VoiceSidecarResolved,
104 VoiceSidecarFailed,
105 VoiceSidecarTimedOut,
106 VoiceTurnCancelled,
107 VoiceBridgePlayed,
108 GateAccepted,
114 GateRejected,
115 ModelFallback,
142 SessionScope,
149 PermissionDecision,
166 ApprovalRecorded,
172 BranchDecision,
179 AlternativeRejected,
185 InferenceMetered,
190 TransactionConflict,
197 AdmissionGateDecision,
207 ToolReceiptHallucination,
213 GoalEvaluated,
220 TurnCompleted,
228 RunCompleted,
231}
232
233#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
235#[serde(rename_all = "snake_case")]
236pub enum SpanStatus {
237 Ok,
238 Error,
239 Unset,
240}
241
242#[derive(Debug, Clone, Serialize, Deserialize)]
244pub struct Span {
245 pub trace_id: String,
246 pub span_id: String,
247 pub parent_span_id: Option<String>,
248 pub name: String,
249 pub start_time: DateTime<Utc>,
250 pub end_time: Option<DateTime<Utc>>,
251 pub status: SpanStatus,
252 pub attributes: HashMap<String, Value>,
253}
254
255pub mod metric_keys {
260 pub const DURATION_MS: &str = "duration_ms";
262 pub const TOKENS_IN: &str = "tokens_in";
264 pub const TOKENS_OUT: &str = "tokens_out";
266 pub const COST_USD: &str = "cost_usd";
268}
269
270#[derive(Debug, Clone, Copy, Default, PartialEq, Serialize, Deserialize)]
276pub struct Metrics {
277 #[serde(default, skip_serializing_if = "Option::is_none")]
278 pub duration_ms: Option<f64>,
279 #[serde(default, skip_serializing_if = "Option::is_none")]
280 pub tokens_in: Option<u64>,
281 #[serde(default, skip_serializing_if = "Option::is_none")]
282 pub tokens_out: Option<u64>,
283 #[serde(default, skip_serializing_if = "Option::is_none")]
284 pub cost_usd: Option<f64>,
285}
286
287impl Metrics {
288 pub fn latency(duration_ms: f64) -> Self {
290 Self {
291 duration_ms: Some(duration_ms),
292 ..Default::default()
293 }
294 }
295
296 pub fn inference(tokens_in: u64, tokens_out: u64, cost_usd: Option<f64>) -> Self {
298 Self {
299 duration_ms: None,
300 tokens_in: Some(tokens_in),
301 tokens_out: Some(tokens_out),
302 cost_usd,
303 }
304 }
305
306 pub fn with_duration(mut self, duration_ms: f64) -> Self {
307 self.duration_ms = Some(duration_ms);
308 self
309 }
310
311 fn merge_into(&self, data: &mut HashMap<String, Value>) {
313 if let Some(d) = self.duration_ms {
314 data.insert(metric_keys::DURATION_MS.into(), Value::from(d));
315 }
316 if let Some(t) = self.tokens_in {
317 data.insert(metric_keys::TOKENS_IN.into(), Value::from(t));
318 }
319 if let Some(t) = self.tokens_out {
320 data.insert(metric_keys::TOKENS_OUT.into(), Value::from(t));
321 }
322 if let Some(c) = self.cost_usd {
323 data.insert(metric_keys::COST_USD.into(), Value::from(c));
324 }
325 }
326}
327
328#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, schemars::JsonSchema)]
330pub struct Event {
331 pub kind: EventKind,
332 #[serde(default, skip_serializing_if = "Option::is_none")]
335 pub run_id: Option<String>,
336 #[serde(default, skip_serializing_if = "Option::is_none")]
339 pub client_id: Option<String>,
340 #[serde(default, skip_serializing_if = "Option::is_none")]
344 pub policy_session_id: Option<String>,
345 #[serde(default, skip_serializing_if = "Option::is_none")]
346 pub action_id: Option<String>,
347 #[serde(default, skip_serializing_if = "Option::is_none")]
348 pub proposal_id: Option<String>,
349 #[serde(default)]
350 pub data: HashMap<String, Value>,
351 #[serde(default = "Utc::now")]
352 pub timestamp: DateTime<Utc>,
353 #[serde(default, skip_serializing_if = "Option::is_none")]
358 pub prev_hash: Option<String>,
359 #[serde(default, skip_serializing_if = "Option::is_none")]
362 pub hash: Option<String>,
363}
364
365impl Event {
366 pub fn duration_ms(&self) -> Option<f64> {
368 self.data
369 .get(metric_keys::DURATION_MS)
370 .and_then(Value::as_f64)
371 }
372
373 pub fn tokens_in(&self) -> Option<u64> {
375 self.data
376 .get(metric_keys::TOKENS_IN)
377 .and_then(Value::as_u64)
378 }
379
380 pub fn tokens_out(&self) -> Option<u64> {
382 self.data
383 .get(metric_keys::TOKENS_OUT)
384 .and_then(Value::as_u64)
385 }
386
387 pub fn cost_usd(&self) -> Option<f64> {
389 self.data.get(metric_keys::COST_USD).and_then(Value::as_f64)
390 }
391
392 pub fn metrics(&self) -> Metrics {
394 Metrics {
395 duration_ms: self.duration_ms(),
396 tokens_in: self.tokens_in(),
397 tokens_out: self.tokens_out(),
398 cost_usd: self.cost_usd(),
399 }
400 }
401}
402
403#[derive(Debug, Clone, Copy, Default, PartialEq, Serialize, Deserialize)]
407pub struct MetricsTotals {
408 pub duration_ms: f64,
409 pub tokens_in: u64,
410 pub tokens_out: u64,
411 pub tokens: u64,
412 pub cost_usd: f64,
413 pub metered_events: usize,
415}
416
417pub fn metrics_totals_of(events: &[Event]) -> MetricsTotals {
422 let mut totals = MetricsTotals::default();
423 for ev in events {
424 let m = ev.metrics();
425 let mut metered = false;
426 if let Some(d) = m.duration_ms {
427 totals.duration_ms += d;
428 metered = true;
429 }
430 if let Some(t) = m.tokens_in {
431 totals.tokens_in = totals.tokens_in.saturating_add(t);
432 metered = true;
433 }
434 if let Some(t) = m.tokens_out {
435 totals.tokens_out = totals.tokens_out.saturating_add(t);
436 metered = true;
437 }
438 if let Some(c) = m.cost_usd {
439 totals.cost_usd += c;
440 metered = true;
441 }
442 if metered {
443 totals.metered_events += 1;
444 }
445 }
446 totals.tokens = totals.tokens_in.saturating_add(totals.tokens_out);
447 totals
448}
449
450#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
453pub struct AgentCost {
454 pub agent: String,
455 pub calls: u64,
457 pub tokens_in: u64,
458 pub tokens_out: u64,
459 pub cost_usd: f64,
460}
461
462pub fn cost_by_agent_of(events: &[Event]) -> Vec<AgentCost> {
470 use std::collections::BTreeMap;
471 let mut map: BTreeMap<String, AgentCost> = BTreeMap::new();
472 for e in events {
473 if e.kind != EventKind::InferenceMetered {
474 continue;
475 }
476 let agent = e
477 .data
478 .get("agent")
479 .and_then(|v| v.as_str())
480 .unwrap_or("unknown")
481 .to_string();
482 let entry = map.entry(agent.clone()).or_insert_with(|| AgentCost {
483 agent,
484 ..Default::default()
485 });
486 entry.calls += 1;
487 entry.tokens_in = entry.tokens_in.saturating_add(e.tokens_in().unwrap_or(0));
488 entry.tokens_out = entry.tokens_out.saturating_add(e.tokens_out().unwrap_or(0));
489 entry.cost_usd += e.cost_usd().unwrap_or(0.0);
490 }
491 map.into_values().collect()
492}
493
494enum JournalMessage {
513 Async(String),
514 Critical {
515 line: String,
516 known_existing: bool,
517 ack: JournalAcknowledgement,
518 },
519 #[cfg(test)]
520 Shutdown,
521}
522
523pub const MAX_CRITICAL_ACKNOWLEDGEMENT_TIMEOUT: Duration = Duration::from_secs(30);
527
528const DEFAULT_CRITICAL_ACKNOWLEDGEMENT_CAPACITY: usize = 64;
529
530#[derive(Debug, Clone, Copy, PartialEq, Eq)]
534pub enum JournalFailurePoint {
535 AsyncWrite,
536 Write,
537 Flush,
538 Fsync,
539 HoldAcknowledgement,
543}
544
545#[derive(Debug, Clone, Default)]
546pub struct JournalFailureInjector {
547 failures: Arc<Mutex<VecDeque<JournalFailurePoint>>>,
548 held_acknowledgements: Arc<Mutex<Vec<JournalAcknowledgement>>>,
549}
550
551impl JournalFailureInjector {
552 pub fn fail_next(&self, point: JournalFailurePoint) {
553 self.failures
554 .lock()
555 .expect("journal failure injector mutex poisoned")
556 .push_back(point);
557 }
558
559 fn take(&self, point: JournalFailurePoint) -> bool {
560 let mut failures = self
561 .failures
562 .lock()
563 .expect("journal failure injector mutex poisoned");
564 if failures.front() == Some(&point) {
565 failures.pop_front();
566 true
567 } else {
568 false
569 }
570 }
571
572 fn hold_acknowledgement(&self, ack: JournalAcknowledgement) {
573 self.held_acknowledgements
574 .lock()
575 .expect("journal held-acknowledgement mutex poisoned")
576 .push(ack);
577 }
578
579 #[doc(hidden)]
582 pub fn held_acknowledgement_count(&self) -> usize {
583 self.held_acknowledgements
584 .lock()
585 .expect("journal held-acknowledgement mutex poisoned")
586 .len()
587 }
588
589 #[doc(hidden)]
592 pub fn release_held_acknowledgements(&self) {
593 let acknowledgements: Vec<_> = self
594 .held_acknowledgements
595 .lock()
596 .expect("journal held-acknowledgement mutex poisoned")
597 .drain(..)
598 .collect();
599 for acknowledgement in acknowledgements {
600 acknowledgement.send(Ok(()));
601 }
602 }
603}
604
605#[derive(Debug)]
606enum CriticalPreAcceptanceError {
607 WriterUnavailable,
608 WriterStopped,
609 CoordinatorUnavailable(String),
610 CapacityExhausted { capacity: usize },
611 InvalidAcknowledgementTimeout { requested: Duration },
612}
613
614impl std::fmt::Display for CriticalPreAcceptanceError {
615 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
616 match self {
617 Self::WriterUnavailable => write!(formatter, "journal writer thread is unavailable"),
618 Self::WriterStopped => {
619 write!(
620 formatter,
621 "journal writer stopped before accepting critical append"
622 )
623 }
624 Self::CoordinatorUnavailable(reason) => write!(
625 formatter,
626 "critical acknowledgement coordinator is unavailable: {reason}"
627 ),
628 Self::CapacityExhausted { capacity } => write!(
629 formatter,
630 "critical acknowledgement capacity is exhausted ({capacity} in flight)"
631 ),
632 Self::InvalidAcknowledgementTimeout { requested } => write!(
633 formatter,
634 "critical acknowledgement timeout must be between 1ns and {}ms, got {}ms",
635 MAX_CRITICAL_ACKNOWLEDGEMENT_TIMEOUT.as_millis(),
636 requested.as_millis()
637 ),
638 }
639 }
640}
641
642#[derive(Debug)]
643enum CriticalPostAcceptanceError {
644 DurabilityFailure(String),
645 AcknowledgementTimedOut { timeout: Duration },
646 CoordinatorStopped,
647}
648
649impl std::fmt::Display for CriticalPostAcceptanceError {
650 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
651 match self {
652 Self::DurabilityFailure(reason) => write!(formatter, "{reason}"),
653 Self::AcknowledgementTimedOut { timeout } => write!(
654 formatter,
655 "journal writer did not acknowledge within {}ms",
656 timeout.as_millis()
657 ),
658 Self::CoordinatorStopped => {
659 write!(
660 formatter,
661 "acknowledgement coordinator stopped after enqueue"
662 )
663 }
664 }
665 }
666}
667
668struct AsyncAcknowledgementState {
669 terminal: bool,
670 result: Option<Result<(), CriticalPostAcceptanceError>>,
671 waker: Option<Waker>,
672}
673
674struct AsyncAcknowledgementEntry {
675 id: u64,
676 deadline: Instant,
677 timeout: Duration,
678 manager: Weak<AsyncAcknowledgementManagerInner>,
679 state: Mutex<AsyncAcknowledgementState>,
680}
681
682impl AsyncAcknowledgementEntry {
683 fn complete(&self, result: Result<(), CriticalPostAcceptanceError>) {
684 self.complete_deciding(|| result);
685 }
686
687 fn complete_writer(&self, result: std::io::Result<()>) {
688 self.complete_deciding(|| {
689 if Instant::now() >= self.deadline {
690 Err(CriticalPostAcceptanceError::AcknowledgementTimedOut {
691 timeout: self.timeout,
692 })
693 } else {
694 result.map_err(|error| {
695 CriticalPostAcceptanceError::DurabilityFailure(error.to_string())
696 })
697 }
698 });
699 }
700
701 fn complete_deciding(&self, decide: impl FnOnce() -> Result<(), CriticalPostAcceptanceError>) {
702 let waker = {
703 let mut state = self
704 .state
705 .lock()
706 .expect("journal async-acknowledgement mutex poisoned");
707 if state.terminal {
708 return;
709 }
710 state.terminal = true;
711 state.result = Some(decide());
712 state.waker.take()
713 };
714 if let Some(manager) = self.manager.upgrade() {
715 manager.remove(self.id);
716 }
717 if let Some(waker) = waker {
718 waker.wake();
719 }
720 }
721
722 fn cancel_preacceptance(&self) {
723 {
724 let mut state = self
725 .state
726 .lock()
727 .expect("journal async-acknowledgement mutex poisoned");
728 if state.terminal {
729 return;
730 }
731 state.terminal = true;
732 }
733 if let Some(manager) = self.manager.upgrade() {
734 manager.remove(self.id);
735 }
736 }
737}
738
739struct AsyncAcknowledgement {
740 entry: Arc<AsyncAcknowledgementEntry>,
741}
742
743impl Future for AsyncAcknowledgement {
744 type Output = Result<(), CriticalPostAcceptanceError>;
745
746 fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
747 let mut state = self
748 .entry
749 .state
750 .lock()
751 .expect("journal async-acknowledgement mutex poisoned");
752 match state.result.take() {
753 Some(result) => Poll::Ready(result),
754 None => {
755 state.waker = Some(context.waker().clone());
756 Poll::Pending
757 }
758 }
759 }
760}
761
762struct AsyncAcknowledgementManagerState {
763 entries: HashMap<u64, Arc<AsyncAcknowledgementEntry>>,
764 next_id: u64,
765 shutting_down: bool,
766}
767
768struct AsyncAcknowledgementManagerInner {
769 capacity: usize,
770 state: Mutex<AsyncAcknowledgementManagerState>,
771 changed: Condvar,
772 #[cfg(test)]
773 expiry_barrier: Mutex<Option<AsyncAcknowledgementExpiryBarrier>>,
774}
775
776#[cfg(test)]
777#[derive(Clone)]
778struct AsyncAcknowledgementExpiryBarrier {
779 removed: Arc<std::sync::Barrier>,
780 release: Arc<std::sync::Barrier>,
781}
782
783#[cfg(test)]
784impl AsyncAcknowledgementExpiryBarrier {
785 fn new() -> Self {
786 Self {
787 removed: Arc::new(std::sync::Barrier::new(2)),
788 release: Arc::new(std::sync::Barrier::new(2)),
789 }
790 }
791
792 fn pause_after_removal(&self) {
793 self.removed.wait();
794 self.release.wait();
795 }
796
797 fn wait_until_removed(&self) {
798 self.removed.wait();
799 }
800
801 fn allow_timeout_completion(&self) {
802 self.release.wait();
803 }
804}
805
806impl AsyncAcknowledgementManagerInner {
807 fn remove(&self, id: u64) {
808 let removed = self
809 .state
810 .lock()
811 .expect("journal acknowledgement-manager mutex poisoned")
812 .entries
813 .remove(&id)
814 .is_some();
815 if removed {
816 self.changed.notify_all();
817 }
818 }
819}
820
821struct AsyncAcknowledgementManager {
822 inner: Arc<AsyncAcknowledgementManagerInner>,
823 worker: Mutex<Option<thread::JoinHandle<()>>>,
824}
825
826impl AsyncAcknowledgementManager {
827 fn new(capacity: usize) -> Self {
828 Self {
829 inner: Arc::new(AsyncAcknowledgementManagerInner {
830 capacity,
831 state: Mutex::new(AsyncAcknowledgementManagerState {
832 entries: HashMap::new(),
833 next_id: 0,
834 shutting_down: false,
835 }),
836 changed: Condvar::new(),
837 #[cfg(test)]
838 expiry_barrier: Mutex::new(None),
839 }),
840 worker: Mutex::new(None),
841 }
842 }
843
844 #[cfg(test)]
845 fn pause_next_expiry_after_removal(&self, barrier: AsyncAcknowledgementExpiryBarrier) {
846 *self
847 .inner
848 .expiry_barrier
849 .lock()
850 .expect("journal acknowledgement expiry-barrier mutex poisoned") = Some(barrier);
851 }
852
853 fn ensure_worker(&self) -> Result<(), CriticalPreAcceptanceError> {
854 let mut worker = self
855 .worker
856 .lock()
857 .expect("journal acknowledgement-worker mutex poisoned");
858 if worker.is_some() {
859 return Ok(());
860 }
861 let inner = self.inner.clone();
862 let handle = thread::Builder::new()
863 .name("car-eventlog-critical-ack".into())
864 .spawn(move || async_acknowledgement_timer_loop(inner))
865 .map_err(|error| {
866 CriticalPreAcceptanceError::CoordinatorUnavailable(error.to_string())
867 })?;
868 *worker = Some(handle);
869 Ok(())
870 }
871
872 fn reserve(
873 &self,
874 timeout: Duration,
875 ) -> Result<AsyncAcknowledgementReservation, CriticalPreAcceptanceError> {
876 if timeout.is_zero() || timeout > MAX_CRITICAL_ACKNOWLEDGEMENT_TIMEOUT {
877 return Err(CriticalPreAcceptanceError::InvalidAcknowledgementTimeout {
878 requested: timeout,
879 });
880 }
881 self.ensure_worker()?;
882 let mut state = self
883 .inner
884 .state
885 .lock()
886 .expect("journal acknowledgement-manager mutex poisoned");
887 if state.shutting_down {
888 return Err(CriticalPreAcceptanceError::CoordinatorUnavailable(
889 "coordinator is shutting down".to_string(),
890 ));
891 }
892 if state.entries.len() >= self.inner.capacity {
893 return Err(CriticalPreAcceptanceError::CapacityExhausted {
894 capacity: self.inner.capacity,
895 });
896 }
897 let id = loop {
898 let candidate = state.next_id;
899 state.next_id = state.next_id.wrapping_add(1);
900 if !state.entries.contains_key(&candidate) {
901 break candidate;
902 }
903 };
904 let entry = Arc::new(AsyncAcknowledgementEntry {
905 id,
906 deadline: Instant::now() + timeout,
907 timeout,
908 manager: Arc::downgrade(&self.inner),
909 state: Mutex::new(AsyncAcknowledgementState {
910 terminal: false,
911 result: None,
912 waker: None,
913 }),
914 });
915 state.entries.insert(id, entry.clone());
916 drop(state);
917 self.inner.changed.notify_all();
918 Ok(AsyncAcknowledgementReservation {
919 entry,
920 preaccepted: true,
921 })
922 }
923
924 fn shutdown(&self) {
925 let entries = {
926 let mut state = self
927 .inner
928 .state
929 .lock()
930 .expect("journal acknowledgement-manager mutex poisoned");
931 state.shutting_down = true;
932 let entries = state
933 .entries
934 .drain()
935 .map(|(_, entry)| entry)
936 .collect::<Vec<_>>();
937 self.inner.changed.notify_all();
938 entries
939 };
940 for entry in entries {
941 entry.complete(Err(CriticalPostAcceptanceError::CoordinatorStopped));
942 }
943 if let Some(worker) = self
944 .worker
945 .lock()
946 .expect("journal acknowledgement-worker mutex poisoned")
947 .take()
948 {
949 let _ = worker.join();
950 }
951 }
952}
953
954fn async_acknowledgement_timer_loop(inner: Arc<AsyncAcknowledgementManagerInner>) {
955 loop {
956 let expired = {
957 let mut state = inner
958 .state
959 .lock()
960 .expect("journal acknowledgement-manager mutex poisoned");
961 loop {
962 if state.shutting_down {
963 return;
964 }
965 let now = Instant::now();
966 let expired_ids: Vec<_> = state
967 .entries
968 .iter()
969 .filter_map(|(id, entry)| (entry.deadline <= now).then_some(*id))
970 .collect();
971 if !expired_ids.is_empty() {
972 break expired_ids
973 .into_iter()
974 .filter_map(|id| state.entries.remove(&id))
975 .collect::<Vec<_>>();
976 }
977 if let Some(deadline) = state.entries.values().map(|entry| entry.deadline).min() {
978 let wait = deadline.saturating_duration_since(now);
979 let (next, _) = inner
980 .changed
981 .wait_timeout(state, wait)
982 .expect("journal acknowledgement-manager mutex poisoned");
983 state = next;
984 } else {
985 state = inner
986 .changed
987 .wait(state)
988 .expect("journal acknowledgement-manager mutex poisoned");
989 }
990 }
991 };
992 #[cfg(test)]
993 if !expired.is_empty() {
994 if let Some(barrier) = inner
995 .expiry_barrier
996 .lock()
997 .expect("journal acknowledgement expiry-barrier mutex poisoned")
998 .take()
999 {
1000 barrier.pause_after_removal();
1001 }
1002 }
1003 for entry in expired {
1004 entry.complete(Err(CriticalPostAcceptanceError::AcknowledgementTimedOut {
1005 timeout: entry.timeout,
1006 }));
1007 }
1008 }
1009}
1010
1011struct AsyncAcknowledgementReservation {
1012 entry: Arc<AsyncAcknowledgementEntry>,
1013 preaccepted: bool,
1014}
1015
1016impl AsyncAcknowledgementReservation {
1017 fn sender(&self) -> AsyncAcknowledgementSender {
1018 AsyncAcknowledgementSender {
1019 entry: self.entry.clone(),
1020 }
1021 }
1022
1023 fn into_future(mut self) -> AsyncAcknowledgement {
1024 self.preaccepted = false;
1025 AsyncAcknowledgement {
1026 entry: self.entry.clone(),
1027 }
1028 }
1029}
1030
1031impl Drop for AsyncAcknowledgementReservation {
1032 fn drop(&mut self) {
1033 if self.preaccepted {
1034 self.entry.cancel_preacceptance();
1035 }
1036 }
1037}
1038
1039struct AsyncAcknowledgementSender {
1040 entry: Arc<AsyncAcknowledgementEntry>,
1041}
1042
1043impl AsyncAcknowledgementSender {
1044 fn send(&self, result: std::io::Result<()>) {
1045 self.entry.complete_writer(result);
1046 }
1047}
1048
1049enum JournalAcknowledgement {
1050 Sync(mpsc::SyncSender<std::io::Result<()>>),
1051 Async(AsyncAcknowledgementSender),
1052}
1053
1054impl JournalAcknowledgement {
1055 fn send(&self, result: std::io::Result<()>) {
1056 match self {
1057 Self::Sync(sender) => {
1058 let _ = sender.send(result);
1059 }
1060 Self::Async(sender) => sender.send(result),
1061 }
1062 }
1063}
1064
1065impl std::fmt::Debug for JournalAcknowledgement {
1066 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1067 match self {
1068 Self::Sync(_) => formatter.write_str("JournalAcknowledgement::Sync"),
1069 Self::Async(_) => formatter.write_str("JournalAcknowledgement::Async"),
1070 }
1071 }
1072}
1073
1074#[derive(Debug, Clone, PartialEq, Eq)]
1078pub enum CriticalAppendError {
1079 Rejected { reason: String },
1081 DurabilityUnknown { reason: String },
1084}
1085
1086impl CriticalAppendError {
1087 pub fn is_retry_safe(&self) -> bool {
1088 matches!(self, Self::DurabilityUnknown { .. })
1089 }
1090}
1091
1092impl std::fmt::Display for CriticalAppendError {
1093 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1094 match self {
1095 Self::Rejected { reason } => {
1096 write!(formatter, "critical journal append rejected: {reason}")
1097 }
1098 Self::DurabilityUnknown { reason } => write!(
1099 formatter,
1100 "critical journal durability is unknown; retry the exact event safely: {reason}"
1101 ),
1102 }
1103 }
1104}
1105
1106impl std::error::Error for CriticalAppendError {}
1107
1108struct JournalWriter {
1109 tx: Option<mpsc::Sender<JournalMessage>>,
1112 handle: Option<thread::JoinHandle<()>>,
1113 acknowledgements: AsyncAcknowledgementManager,
1114}
1115
1116impl JournalWriter {
1117 fn spawn(path: PathBuf) -> Self {
1118 Self::spawn_with_injectors(path, JournalFailureInjector::default(), None)
1119 }
1120
1121 fn spawn_with_injector(path: PathBuf, failures: JournalFailureInjector) -> Self {
1122 Self::spawn_with_injectors(path, failures, None)
1123 }
1124
1125 fn spawn_with_private_path_injector(
1126 path: PathBuf,
1127 failures: car_secrets::PrivatePathDurabilityFailureInjector,
1128 ) -> Self {
1129 Self::spawn_with_injectors(path, JournalFailureInjector::default(), Some(failures))
1130 }
1131
1132 fn spawn_with_injectors(
1133 path: PathBuf,
1134 failures: JournalFailureInjector,
1135 private_path_failures: Option<car_secrets::PrivatePathDurabilityFailureInjector>,
1136 ) -> Self {
1137 Self::spawn_with_injectors_and_ack_capacity(
1138 path,
1139 failures,
1140 private_path_failures,
1141 DEFAULT_CRITICAL_ACKNOWLEDGEMENT_CAPACITY,
1142 )
1143 }
1144
1145 fn spawn_with_injectors_and_ack_capacity(
1146 path: PathBuf,
1147 failures: JournalFailureInjector,
1148 private_path_failures: Option<car_secrets::PrivatePathDurabilityFailureInjector>,
1149 acknowledgement_capacity: usize,
1150 ) -> Self {
1151 let (tx, rx) = mpsc::channel::<JournalMessage>();
1152 let acknowledgements = AsyncAcknowledgementManager::new(acknowledgement_capacity);
1153 match thread::Builder::new()
1154 .name("car-eventlog-journal".into())
1155 .spawn(move || journal_loop(path, rx, failures, private_path_failures))
1156 {
1157 Ok(handle) => Self {
1158 tx: Some(tx),
1159 handle: Some(handle),
1160 acknowledgements,
1161 },
1162 Err(e) => {
1164 tracing::warn!(error = %e, "car-eventlog: failed to spawn journal writer thread — journaling disabled for this log");
1165 Self {
1166 tx: None,
1167 handle: None,
1168 acknowledgements,
1169 }
1170 }
1171 }
1172 }
1173
1174 fn send(&self, line: String) {
1175 if let Some(tx) = &self.tx {
1176 let _ = tx.send(JournalMessage::Async(line));
1178 }
1179 }
1180
1181 fn enqueue_critical_sync(
1182 &self,
1183 line: String,
1184 known_existing: bool,
1185 ) -> Result<mpsc::Receiver<std::io::Result<()>>, CriticalPreAcceptanceError> {
1186 let tx = self
1187 .tx
1188 .as_ref()
1189 .ok_or(CriticalPreAcceptanceError::WriterUnavailable)?;
1190 let (ack_tx, ack_rx) = mpsc::sync_channel(0);
1191 tx.send(JournalMessage::Critical {
1192 line,
1193 known_existing,
1194 ack: JournalAcknowledgement::Sync(ack_tx),
1195 })
1196 .map_err(|_| CriticalPreAcceptanceError::WriterStopped)?;
1197 Ok(ack_rx)
1198 }
1199
1200 fn reserve_async_acknowledgement(
1201 &self,
1202 acknowledgement_timeout: Duration,
1203 ) -> Result<AsyncAcknowledgementReservation, CriticalPreAcceptanceError> {
1204 if self.tx.is_none() {
1205 return Err(CriticalPreAcceptanceError::WriterUnavailable);
1206 }
1207 self.acknowledgements.reserve(acknowledgement_timeout)
1208 }
1209
1210 fn enqueue_critical_async(
1211 &self,
1212 line: String,
1213 known_existing: bool,
1214 reservation: AsyncAcknowledgementReservation,
1215 ) -> Result<AsyncAcknowledgement, CriticalPreAcceptanceError> {
1216 let tx = self
1217 .tx
1218 .as_ref()
1219 .ok_or(CriticalPreAcceptanceError::WriterUnavailable)?;
1220 tx.send(JournalMessage::Critical {
1221 line,
1222 known_existing,
1223 ack: JournalAcknowledgement::Async(reservation.sender()),
1224 })
1225 .map_err(|_| CriticalPreAcceptanceError::WriterStopped)?;
1226 Ok(reservation.into_future())
1227 }
1228
1229 #[cfg(test)]
1230 fn remove_sender_for_test(&mut self) {
1231 self.tx.take();
1232 if let Some(handle) = self.handle.take() {
1233 let _ = handle.join();
1234 }
1235 }
1236
1237 #[cfg(test)]
1238 fn stop_receiver_for_test(&mut self) {
1239 if let Some(tx) = &self.tx {
1240 let _ = tx.send(JournalMessage::Shutdown);
1241 }
1242 if let Some(handle) = self.handle.take() {
1243 let _ = handle.join();
1244 }
1245 }
1246}
1247
1248impl Drop for JournalWriter {
1249 fn drop(&mut self) {
1250 self.acknowledgements.shutdown();
1251 self.tx.take();
1254 if let Some(handle) = self.handle.take() {
1255 let _ = handle.join();
1256 }
1257 }
1258}
1259
1260fn journal_loop(
1271 path: PathBuf,
1272 rx: mpsc::Receiver<JournalMessage>,
1273 failures: JournalFailureInjector,
1274 private_path_failures: Option<car_secrets::PrivatePathDurabilityFailureInjector>,
1275) {
1276 let mut writer: Option<std::fs::File> = None;
1277 let mut written_critical_lines = HashSet::new();
1278 let mut blocked_critical: Option<String> = None;
1279 let mut failed_async = VecDeque::new();
1284 let mut after_blocked_critical = VecDeque::new();
1285 while let Ok(message) = rx.recv() {
1286 let existed_before_open = path.exists();
1287 let open = || match private_path_failures.as_ref() {
1288 Some(failures) => {
1289 car_secrets::open_private_append_with_failure_injector(&path, failures)
1290 }
1291 None => open_private_append(&path),
1292 };
1293 match message {
1294 #[cfg(test)]
1295 JournalMessage::Shutdown => break,
1296 JournalMessage::Async(line) => {
1297 if blocked_critical.is_some() {
1298 after_blocked_critical.push_back(line);
1299 continue;
1300 }
1301 if !failed_async.is_empty() {
1302 failed_async.push_back(line);
1303 continue;
1304 }
1305 if writer.is_none() {
1306 writer = open().ok();
1307 }
1308 let result = match writer.as_mut() {
1309 Some(file) => append_journal_line(
1310 &path,
1311 file,
1312 &line,
1313 false,
1314 false,
1315 existed_before_open,
1316 &failures,
1317 &mut written_critical_lines,
1318 ),
1319 None => Err(std::io::Error::other("cannot open journal file")),
1320 };
1321 if let Err(error) = result {
1322 failed_async.push_back(line);
1323 tracing::warn!(path = %path.display(), %error, "car-eventlog: asynchronous journal append failed and is awaiting ordered retry");
1324 }
1325 continue;
1326 }
1327 JournalMessage::Critical {
1328 line,
1329 known_existing,
1330 ack,
1331 } => {
1332 if blocked_critical
1333 .as_deref()
1334 .is_some_and(|pending| pending != line)
1335 {
1336 ack.send(Err(std::io::Error::new(
1337 std::io::ErrorKind::WouldBlock,
1338 "another critical journal row is awaiting durability",
1339 )));
1340 continue;
1341 }
1342 if writer.is_none() {
1343 writer = open().ok();
1344 }
1345
1346 while let Some(pending) = failed_async.front() {
1351 let replay = match writer.as_mut() {
1352 Some(file) => append_journal_line(
1353 &path,
1354 file,
1355 pending,
1356 false,
1357 false,
1358 existed_before_open,
1359 &failures,
1360 &mut written_critical_lines,
1361 ),
1362 None => Err(std::io::Error::other("cannot open journal file")),
1363 };
1364 match replay {
1365 Ok(()) => {
1366 failed_async.pop_front();
1367 }
1368 Err(error) => {
1369 blocked_critical = Some(line.clone());
1370 ack.send(Err(std::io::Error::new(
1371 error.kind(),
1372 format!("prior asynchronous journal row is not durable: {error}"),
1373 )));
1374 break;
1375 }
1376 }
1377 }
1378 if !failed_async.is_empty() {
1379 continue;
1380 }
1381 let result = match writer.as_mut() {
1382 Some(file) => append_journal_line(
1383 &path,
1384 file,
1385 &line,
1386 true,
1387 known_existing,
1388 existed_before_open,
1389 &failures,
1390 &mut written_critical_lines,
1391 ),
1392 None => Err(std::io::Error::other("cannot open journal file")),
1393 };
1394 match result {
1395 Ok(()) => {
1396 blocked_critical = None;
1397 while let Some(queued) = after_blocked_critical.pop_front() {
1398 if let Some(file) = writer.as_mut() {
1399 if let Err(error) = append_journal_line(
1400 &path,
1401 file,
1402 &queued,
1403 false,
1404 false,
1405 true,
1406 &failures,
1407 &mut written_critical_lines,
1408 ) {
1409 failed_async.push_back(queued);
1410 failed_async.append(&mut after_blocked_critical);
1411 tracing::warn!(path = %path.display(), %error, "car-eventlog: queued asynchronous append failed after critical recovery and is awaiting ordered retry");
1412 break;
1413 }
1414 }
1415 }
1416 if failures.take(JournalFailurePoint::HoldAcknowledgement) {
1417 failures.hold_acknowledgement(ack);
1418 } else {
1419 ack.send(Ok(()));
1420 }
1421 }
1422 Err(error) => {
1423 blocked_critical = Some(line);
1424 ack.send(Err(error));
1425 }
1426 }
1427 continue;
1428 }
1429 };
1430 }
1431 if let Some(mut writer) = writer {
1432 let _ = writer.flush();
1433 }
1434}
1435
1436fn append_journal_line(
1437 path: &Path,
1438 file: &mut std::fs::File,
1439 line: &str,
1440 critical: bool,
1441 known_existing: bool,
1442 existed_before_open: bool,
1443 failures: &JournalFailureInjector,
1444 written_critical_lines: &mut HashSet<String>,
1445) -> std::io::Result<()> {
1446 revalidate_private_path(path, file)?;
1447 let already_written = known_existing || written_critical_lines.contains(line);
1448 if !already_written {
1449 if (!critical && failures.take(JournalFailurePoint::AsyncWrite))
1450 || (critical && failures.take(JournalFailurePoint::Write))
1451 {
1452 return Err(std::io::Error::from_raw_os_error(28)); }
1454 let mut bytes = Vec::with_capacity(line.len() + 2);
1455 let len = file.seek(SeekFrom::End(0))?;
1456 if len > 0 {
1457 file.seek(SeekFrom::End(-1))?;
1458 let mut tail = [0u8; 1];
1459 file.read_exact(&mut tail)?;
1460 if tail[0] != b'\n' {
1461 bytes.push(b'\n');
1462 }
1463 }
1464 bytes.extend_from_slice(line.as_bytes());
1465 bytes.push(b'\n');
1466 if let Err(error) = file.write_all(&bytes) {
1467 let _ = file.set_len(len);
1471 let _ = file.seek(SeekFrom::End(0));
1472 return Err(error);
1473 }
1474 if critical {
1475 written_critical_lines.insert(line.to_string());
1476 }
1477 }
1478 if critical && failures.take(JournalFailurePoint::Flush) {
1479 return Err(std::io::Error::other("injected journal flush failure"));
1480 }
1481 file.flush()?;
1482 if critical {
1483 if failures.take(JournalFailurePoint::Fsync) {
1484 return Err(std::io::Error::other("injected journal fsync failure"));
1485 }
1486 file.sync_all()?;
1487 if !existed_before_open {
1488 sync_journal_parent(path)?;
1489 }
1490 }
1491 revalidate_private_path(path, file)
1492}
1493
1494#[cfg(not(target_os = "windows"))]
1495fn sync_journal_parent(path: &Path) -> std::io::Result<()> {
1496 if let Some(parent) = path.parent() {
1497 std::fs::File::open(parent)?.sync_all()?;
1498 }
1499 Ok(())
1500}
1501
1502#[cfg(target_os = "windows")]
1503fn sync_journal_parent(_path: &Path) -> std::io::Result<()> {
1504 Ok(())
1509}
1510
1511pub struct EventLog {
1513 events: Vec<Event>,
1514 spans: Vec<Span>,
1515 journal: Option<JournalWriter>,
1516 hash_chaining: bool,
1521 last_hash: Option<String>,
1524 retention: Option<RetentionPolicy>,
1529 journal_path: Option<PathBuf>,
1532 journal_lines: usize,
1537 trimmed_events: u64,
1543 cumulative_cost_usd: f64,
1549 active_binding: Option<EventBinding>,
1553 critical_pending: HashSet<String>,
1557}
1558
1559#[derive(Debug, Clone, PartialEq, Eq)]
1560struct EventBinding {
1561 run_id: String,
1562 client_id: String,
1563 policy_session_id: Option<String>,
1564}
1565
1566struct PreparedCriticalAppend {
1567 existing_index: Option<usize>,
1568 event: Option<Event>,
1569 line: String,
1570 known_existing: bool,
1571}
1572
1573const JOURNAL_COMPACT_MIN_EXCESS: usize = 1024;
1579
1580fn event_digest(
1592 prev_hash: &str,
1593 kind: &EventKind,
1594 run_id: Option<&str>,
1595 client_id: Option<&str>,
1596 policy_session_id: Option<&str>,
1597 action_id: Option<&str>,
1598 proposal_id: Option<&str>,
1599 data: &HashMap<String, Value>,
1600 timestamp: &DateTime<Utc>,
1601) -> String {
1602 use sha2::{Digest, Sha256};
1603 let mut sorted: Vec<(&String, &Value)> = data.iter().collect();
1604 sorted.sort_by(|a, b| a.0.cmp(b.0));
1605 let data_canon: String = sorted
1606 .iter()
1607 .map(|(k, v)| format!("{k}={}", v))
1608 .collect::<Vec<_>>()
1609 .join("\u{1f}");
1610 let kind_str = serde_json::to_string(kind).unwrap_or_default();
1611 let mut hasher = Sha256::new();
1612 hasher.update(prev_hash.as_bytes());
1613 hasher.update(b"\x1e");
1614 hasher.update(kind_str.as_bytes());
1615 if run_id.is_some() || client_id.is_some() || policy_session_id.is_some() {
1618 hasher.update(b"\x1d");
1619 hasher.update(run_id.unwrap_or("").as_bytes());
1620 hasher.update(b"\x1f");
1621 hasher.update(client_id.unwrap_or("").as_bytes());
1622 hasher.update(b"\x1f");
1623 hasher.update(policy_session_id.unwrap_or("").as_bytes());
1624 }
1625 hasher.update(b"\x1e");
1626 hasher.update(action_id.unwrap_or("").as_bytes());
1627 hasher.update(b"\x1e");
1628 hasher.update(proposal_id.unwrap_or("").as_bytes());
1629 hasher.update(b"\x1e");
1630 hasher.update(data_canon.as_bytes());
1631 hasher.update(b"\x1e");
1632 hasher.update(timestamp.to_rfc3339().as_bytes());
1633 let digest = hasher.finalize();
1634 digest.iter().map(|b| format!("{b:02x}")).collect()
1635}
1636
1637#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
1642pub struct RetentionPolicy {
1643 #[serde(default)]
1645 pub max_events: Option<usize>,
1646 #[serde(default)]
1648 pub max_age_secs: Option<i64>,
1649}
1650
1651#[derive(Debug, Clone, Default, Serialize, Deserialize)]
1656pub struct EventQuery {
1657 #[serde(default)]
1659 pub kinds: Vec<EventKind>,
1660 #[serde(default)]
1662 pub action_id: Option<String>,
1663 #[serde(default)]
1665 pub proposal_id: Option<String>,
1666 #[serde(default)]
1668 pub since: Option<DateTime<Utc>>,
1669 #[serde(default)]
1671 pub until: Option<DateTime<Utc>>,
1672 #[serde(default)]
1676 pub data_matches: std::collections::HashMap<String, String>,
1677 #[serde(default)]
1679 pub limit: Option<usize>,
1680}
1681
1682fn data_value_matches(v: &Value, want: &str) -> bool {
1685 match v {
1686 Value::String(s) => s == want,
1687 Value::Null => false,
1688 other => *other == want,
1689 }
1690}
1691
1692impl EventQuery {
1693 pub fn matches(&self, e: &Event) -> bool {
1695 if !self.kinds.is_empty() && !self.kinds.contains(&e.kind) {
1696 return false;
1697 }
1698 if let Some(aid) = &self.action_id {
1699 if e.action_id.as_deref() != Some(aid.as_str()) {
1700 return false;
1701 }
1702 }
1703 if let Some(pid) = &self.proposal_id {
1704 if e.proposal_id.as_deref() != Some(pid.as_str()) {
1705 return false;
1706 }
1707 }
1708 if let Some(since) = self.since {
1709 if e.timestamp < since {
1710 return false;
1711 }
1712 }
1713 if let Some(until) = self.until {
1714 if e.timestamp >= until {
1715 return false;
1716 }
1717 }
1718 for (k, want) in &self.data_matches {
1719 match e.data.get(k) {
1720 Some(v) if data_value_matches(v, want) => {}
1721 _ => return false,
1722 }
1723 }
1724 true
1725 }
1726}
1727
1728impl EventLog {
1729 pub fn new() -> Self {
1730 Self {
1731 events: Vec::new(),
1732 spans: Vec::new(),
1733 journal: None,
1734 hash_chaining: false,
1735 last_hash: None,
1736 retention: None,
1737 journal_path: None,
1738 journal_lines: 0,
1739 trimmed_events: 0,
1740 cumulative_cost_usd: 0.0,
1741 active_binding: None,
1742 critical_pending: HashSet::new(),
1743 }
1744 }
1745
1746 pub fn with_journal(path: PathBuf) -> Self {
1747 Self {
1748 events: Vec::new(),
1749 spans: Vec::new(),
1750 journal: Some(JournalWriter::spawn(path.clone())),
1751 hash_chaining: false,
1752 last_hash: None,
1753 retention: None,
1754 journal_path: Some(path),
1755 journal_lines: 0,
1756 trimmed_events: 0,
1757 cumulative_cost_usd: 0.0,
1758 active_binding: None,
1759 critical_pending: HashSet::new(),
1760 }
1761 }
1762
1763 pub fn with_journal_failure_injector(path: PathBuf, failures: JournalFailureInjector) -> Self {
1765 let mut log = Self::with_journal(path.clone());
1766 log.journal = Some(JournalWriter::spawn_with_injector(path, failures));
1767 log
1768 }
1769
1770 #[cfg(test)]
1771 fn with_journal_failure_injector_and_ack_capacity(
1772 path: PathBuf,
1773 failures: JournalFailureInjector,
1774 acknowledgement_capacity: usize,
1775 ) -> Self {
1776 let mut log = Self::with_journal(path.clone());
1777 log.journal = Some(JournalWriter::spawn_with_injectors_and_ack_capacity(
1778 path,
1779 failures,
1780 None,
1781 acknowledgement_capacity,
1782 ));
1783 log
1784 }
1785
1786 pub fn with_private_path_failure_injector(
1788 path: PathBuf,
1789 failures: car_secrets::PrivatePathDurabilityFailureInjector,
1790 ) -> Self {
1791 let mut log = Self::with_journal(path.clone());
1792 log.journal = Some(JournalWriter::spawn_with_private_path_injector(
1793 path, failures,
1794 ));
1795 log
1796 }
1797
1798 pub fn bind_run(&mut self, run_id: &str, client_id: &str) -> Result<(), String> {
1802 if run_id.is_empty() || client_id.is_empty() {
1803 return Err("active journal binding requires non-empty run_id and client_id".into());
1804 }
1805 match &self.active_binding {
1806 Some(binding) if binding.run_id == run_id && binding.client_id == client_id => Ok(()),
1807 Some(binding) => Err(format!(
1808 "journal is already bound to run_id `{}` and client_id `{}`",
1809 binding.run_id, binding.client_id
1810 )),
1811 None => {
1812 self.active_binding = Some(EventBinding {
1813 run_id: run_id.to_string(),
1814 client_id: client_id.to_string(),
1815 policy_session_id: None,
1816 });
1817 Ok(())
1818 }
1819 }
1820 }
1821
1822 pub fn bind_policy_session(&mut self, policy_session_id: &str) -> Result<(), String> {
1824 if policy_session_id.is_empty() {
1825 return Err("policy_session_id must be non-empty".into());
1826 }
1827 let binding = self
1828 .active_binding
1829 .as_mut()
1830 .ok_or_else(|| "cannot bind a policy session without an active run".to_string())?;
1831 match binding.policy_session_id.as_deref() {
1832 Some(existing) if existing != policy_session_id => Err(format!(
1833 "journal proposal is already bound to policy_session_id `{existing}`"
1834 )),
1835 _ => {
1836 binding.policy_session_id = Some(policy_session_id.to_string());
1837 Ok(())
1838 }
1839 }
1840 }
1841
1842 pub fn clear_policy_session(&mut self, policy_session_id: &str) -> Result<(), String> {
1843 let binding = self
1844 .active_binding
1845 .as_mut()
1846 .ok_or_else(|| "cannot clear a policy session without an active run".to_string())?;
1847 if binding.policy_session_id.as_deref() != Some(policy_session_id) {
1848 return Err("policy_session_id does not match the active journal binding".into());
1849 }
1850 binding.policy_session_id = None;
1851 Ok(())
1852 }
1853
1854 pub fn clear_run_binding(&mut self, run_id: &str, client_id: &str) -> Result<(), String> {
1855 let binding = self
1856 .active_binding
1857 .as_ref()
1858 .ok_or_else(|| "journal has no active run binding".to_string())?;
1859 if binding.run_id != run_id || binding.client_id != client_id {
1860 return Err("run_id/client_id does not match the active journal binding".into());
1861 }
1862 if binding.policy_session_id.is_some() {
1863 return Err(
1864 "cannot clear an active run while a proposal policy session is bound".into(),
1865 );
1866 }
1867 self.active_binding = None;
1868 Ok(())
1869 }
1870
1871 pub fn active_run_binding(&self) -> Option<(&str, &str, Option<&str>)> {
1872 self.active_binding.as_ref().map(|binding| {
1873 (
1874 binding.run_id.as_str(),
1875 binding.client_id.as_str(),
1876 binding.policy_session_id.as_deref(),
1877 )
1878 })
1879 }
1880
1881 pub fn with_hash_chaining(mut self) -> Self {
1886 self.enable_hash_chaining();
1887 self
1888 }
1889
1890 pub fn enable_hash_chaining(&mut self) {
1892 self.hash_chaining = true;
1893 if self.last_hash.is_none() {
1895 self.last_hash = self.events.last().and_then(|e| e.hash.clone());
1896 }
1897 }
1898
1899 pub fn hash_chaining_enabled(&self) -> bool {
1901 self.hash_chaining
1902 }
1903
1904 pub fn append(
1905 &mut self,
1906 kind: EventKind,
1907 action_id: Option<&str>,
1908 proposal_id: Option<&str>,
1909 data: HashMap<String, Value>,
1910 ) -> &Event {
1911 let timestamp = Utc::now();
1912 let (prev_hash, hash) = if self.hash_chaining {
1913 let prev = self.last_hash.clone().unwrap_or_default();
1914 let binding = self.active_binding.as_ref();
1915 let h = event_digest(
1916 &prev,
1917 &kind,
1918 binding.map(|b| b.run_id.as_str()),
1919 binding.map(|b| b.client_id.as_str()),
1920 binding.and_then(|b| b.policy_session_id.as_deref()),
1921 action_id,
1922 proposal_id,
1923 &data,
1924 ×tamp,
1925 );
1926 self.last_hash = Some(h.clone());
1927 (Some(prev), Some(h))
1928 } else {
1929 (None, None)
1930 };
1931 let event = Event {
1932 kind,
1933 run_id: self.active_binding.as_ref().map(|b| b.run_id.clone()),
1934 client_id: self.active_binding.as_ref().map(|b| b.client_id.clone()),
1935 policy_session_id: self
1936 .active_binding
1937 .as_ref()
1938 .and_then(|b| b.policy_session_id.clone()),
1939 action_id: action_id.map(|s| s.to_string()),
1940 proposal_id: proposal_id.map(|s| s.to_string()),
1941 data,
1942 timestamp,
1943 prev_hash,
1944 hash,
1945 };
1946
1947 if let Some(journal) = &self.journal {
1950 if let Ok(json) = serde_json::to_string(&event) {
1951 journal.send(json);
1952 self.journal_lines += 1;
1953 }
1954 }
1955
1956 if let Some(c) = event.cost_usd() {
1959 self.cumulative_cost_usd += c;
1960 }
1961
1962 self.events.push(event);
1963 if let Some(max) = self.retention.as_ref().and_then(|p| p.max_events) {
1968 if self.events.len() > max {
1969 let removed = truncate_vec_keep_last(&mut self.events, max);
1970 self.trimmed_events += removed as u64;
1971 self.maybe_compact_journal();
1974 }
1975 }
1976 self.events.last().unwrap()
1977 }
1978
1979 fn prepare_critical_append(
1980 &mut self,
1981 kind: EventKind,
1982 action_id: Option<&str>,
1983 proposal_id: Option<&str>,
1984 data: HashMap<String, Value>,
1985 ) -> Result<PreparedCriticalAppend, String> {
1986 let binding = self.active_binding.as_ref().ok_or_else(|| {
1987 "critical lifecycle event requires an authenticated run binding".to_string()
1988 })?;
1989 let existing = self.events.iter().position(|event| {
1990 event.kind == kind
1991 && event.run_id.as_deref() == Some(binding.run_id.as_str())
1992 && event.client_id.as_deref() == Some(binding.client_id.as_str())
1993 && event.policy_session_id.as_deref() == binding.policy_session_id.as_deref()
1994 && event.action_id.as_deref() == action_id
1995 && event.proposal_id.as_deref() == proposal_id
1996 && event.data == data
1997 });
1998 if existing.is_none() && !self.critical_pending.is_empty() {
1999 return Err(
2000 "another critical lifecycle event is awaiting an exact durability retry".into(),
2001 );
2002 }
2003
2004 if let Some(index) = existing {
2005 let line = serde_json::to_string(&self.events[index]).map_err(|e| e.to_string())?;
2006 return Ok(PreparedCriticalAppend {
2007 existing_index: Some(index),
2008 event: None,
2009 known_existing: !self.critical_pending.contains(&line),
2010 line,
2011 });
2012 }
2013
2014 let timestamp = Utc::now();
2015 let (prev_hash, hash) = if self.hash_chaining {
2016 let prev = self.last_hash.clone().unwrap_or_default();
2017 let hash = event_digest(
2018 &prev,
2019 &kind,
2020 Some(binding.run_id.as_str()),
2021 Some(binding.client_id.as_str()),
2022 binding.policy_session_id.as_deref(),
2023 action_id,
2024 proposal_id,
2025 &data,
2026 ×tamp,
2027 );
2028 (Some(prev), Some(hash))
2029 } else {
2030 (None, None)
2031 };
2032 let event = Event {
2033 kind,
2034 run_id: Some(binding.run_id.clone()),
2035 client_id: Some(binding.client_id.clone()),
2036 policy_session_id: binding.policy_session_id.clone(),
2037 action_id: action_id.map(str::to_string),
2038 proposal_id: proposal_id.map(str::to_string),
2039 data,
2040 timestamp,
2041 prev_hash,
2042 hash,
2043 };
2044 let line = serde_json::to_string(&event).map_err(|error| error.to_string())?;
2045 Ok(PreparedCriticalAppend {
2046 existing_index: None,
2047 event: Some(event),
2048 line,
2049 known_existing: false,
2050 })
2051 }
2052
2053 fn commit_prepared_critical(&mut self, prepared: PreparedCriticalAppend) -> (usize, String) {
2054 let index = match prepared.existing_index {
2055 Some(index) => index,
2056 None => {
2057 let event = prepared
2058 .event
2059 .expect("new critical append must carry its prepared event");
2060 if self.hash_chaining {
2061 self.last_hash = event.hash.clone();
2062 }
2063 self.events.push(event);
2064 self.journal_lines += 1;
2065 self.events.len() - 1
2066 }
2067 };
2068 (index, prepared.line)
2069 }
2070
2071 pub fn append_critical(
2082 &mut self,
2083 kind: EventKind,
2084 action_id: Option<&str>,
2085 proposal_id: Option<&str>,
2086 data: HashMap<String, Value>,
2087 ) -> Result<&Event, String> {
2088 if self.journal.is_none() {
2089 return Err("critical lifecycle event requires an enabled journal".to_string());
2090 }
2091 let prepared = self.prepare_critical_append(kind, action_id, proposal_id, data)?;
2092 let acknowledgement = self
2093 .journal
2094 .as_ref()
2095 .expect("journal presence checked above")
2096 .enqueue_critical_sync(prepared.line.clone(), prepared.known_existing)
2097 .map_err(|error| error.to_string())?;
2098 let (index, line) = self.commit_prepared_critical(prepared);
2099 self.critical_pending.insert(line.clone());
2100 let result = acknowledgement
2101 .recv()
2102 .map_err(|_| {
2103 std::io::Error::new(
2104 std::io::ErrorKind::BrokenPipe,
2105 "journal writer stopped before critical acknowledgement",
2106 )
2107 })
2108 .and_then(|result| result);
2109 match result {
2110 Ok(()) => {
2111 self.critical_pending.remove(&line);
2112 Ok(&self.events[index])
2113 }
2114 Err(error) => {
2115 self.critical_pending.insert(line);
2116 Err(format!("critical journal append was not durable: {error}"))
2117 }
2118 }
2119 }
2120
2121 pub fn append_critical_bounded(
2129 &mut self,
2130 kind: EventKind,
2131 action_id: Option<&str>,
2132 proposal_id: Option<&str>,
2133 data: HashMap<String, Value>,
2134 acknowledgement_timeout: Duration,
2135 ) -> Result<&Event, CriticalAppendError> {
2136 if acknowledgement_timeout.is_zero()
2137 || acknowledgement_timeout > MAX_CRITICAL_ACKNOWLEDGEMENT_TIMEOUT
2138 {
2139 return Err(CriticalAppendError::Rejected {
2140 reason: CriticalPreAcceptanceError::InvalidAcknowledgementTimeout {
2141 requested: acknowledgement_timeout,
2142 }
2143 .to_string(),
2144 });
2145 }
2146 let prepared = self
2147 .prepare_critical_append(kind, action_id, proposal_id, data)
2148 .map_err(|reason| CriticalAppendError::Rejected { reason })?;
2149 let acknowledgement = self
2150 .journal
2151 .as_ref()
2152 .ok_or_else(|| CriticalAppendError::Rejected {
2153 reason: "critical lifecycle event requires an enabled journal".to_string(),
2154 })?
2155 .enqueue_critical_sync(prepared.line.clone(), prepared.known_existing)
2156 .map_err(|error| CriticalAppendError::Rejected {
2157 reason: error.to_string(),
2158 })?;
2159 let (index, line) = self.commit_prepared_critical(prepared);
2160 self.critical_pending.insert(line.clone());
2161 let result = match acknowledgement.recv_timeout(acknowledgement_timeout) {
2162 Ok(result) => result.map_err(|error| error.to_string()),
2163 Err(mpsc::RecvTimeoutError::Timeout) => Err(format!(
2164 "journal writer did not acknowledge within {}ms",
2165 acknowledgement_timeout.as_millis()
2166 )),
2167 Err(mpsc::RecvTimeoutError::Disconnected) => {
2168 Err("journal writer stopped before critical acknowledgement".to_string())
2169 }
2170 };
2171 match result {
2172 Ok(()) => {
2173 self.critical_pending.remove(&line);
2174 Ok(&self.events[index])
2175 }
2176 Err(reason) => Err(CriticalAppendError::DurabilityUnknown { reason }),
2177 }
2178 }
2179
2180 pub async fn append_critical_async(
2192 &mut self,
2193 kind: EventKind,
2194 action_id: Option<&str>,
2195 proposal_id: Option<&str>,
2196 data: HashMap<String, Value>,
2197 acknowledgement_timeout: Duration,
2198 ) -> Result<&Event, CriticalAppendError> {
2199 let reservation = self
2200 .journal
2201 .as_ref()
2202 .ok_or_else(|| CriticalAppendError::Rejected {
2203 reason: "critical lifecycle event requires an enabled journal".to_string(),
2204 })?
2205 .reserve_async_acknowledgement(acknowledgement_timeout)
2206 .map_err(|error| CriticalAppendError::Rejected {
2207 reason: error.to_string(),
2208 })?;
2209 let prepared = self
2210 .prepare_critical_append(kind, action_id, proposal_id, data)
2211 .map_err(|reason| CriticalAppendError::Rejected { reason })?;
2212 let acknowledgement = self
2213 .journal
2214 .as_ref()
2215 .expect("journal presence checked before reservation")
2216 .enqueue_critical_async(prepared.line.clone(), prepared.known_existing, reservation)
2217 .map_err(|error| CriticalAppendError::Rejected {
2218 reason: error.to_string(),
2219 })?;
2220 let (index, line) = self.commit_prepared_critical(prepared);
2221
2222 self.critical_pending.insert(line.clone());
2225 match acknowledgement.await {
2226 Ok(()) => {
2227 self.critical_pending.remove(&line);
2228 Ok(&self.events[index])
2229 }
2230 Err(error) => Err(CriticalAppendError::DurabilityUnknown {
2231 reason: error.to_string(),
2232 }),
2233 }
2234 }
2235
2236 pub fn verify_chain(&self) -> Result<usize, usize> {
2257 let mut prev = String::new();
2258 let mut verified = 0usize;
2259 let mut chain_started = false;
2260 for (i, ev) in self.events.iter().enumerate() {
2261 let Some(stored) = &ev.hash else {
2262 if chain_started {
2264 return Err(i);
2265 }
2266 continue;
2267 };
2268 let recorded_prev = ev.prev_hash.clone().unwrap_or_default();
2270 if chain_started && recorded_prev != prev {
2271 return Err(i);
2272 }
2273 let recomputed = event_digest(
2274 &recorded_prev,
2275 &ev.kind,
2276 ev.run_id.as_deref(),
2277 ev.client_id.as_deref(),
2278 ev.policy_session_id.as_deref(),
2279 ev.action_id.as_deref(),
2280 ev.proposal_id.as_deref(),
2281 &ev.data,
2282 &ev.timestamp,
2283 );
2284 if &recomputed != stored {
2285 return Err(i);
2286 }
2287 prev = stored.clone();
2288 chain_started = true;
2289 verified += 1;
2290 }
2291 Ok(verified)
2292 }
2293
2294 pub fn append_metered(
2301 &mut self,
2302 kind: EventKind,
2303 action_id: Option<&str>,
2304 proposal_id: Option<&str>,
2305 mut data: HashMap<String, Value>,
2306 metrics: Metrics,
2307 ) -> &Event {
2308 metrics.merge_into(&mut data);
2309 self.append(kind, action_id, proposal_id, data)
2310 }
2311
2312 pub fn metrics_totals(&self) -> MetricsTotals {
2324 metrics_totals_of(&self.events)
2325 }
2326
2327 pub fn cost_by_agent(&self) -> Vec<AgentCost> {
2329 cost_by_agent_of(&self.events)
2330 }
2331
2332 pub fn events(&self) -> &[Event] {
2333 &self.events
2334 }
2335
2336 pub fn len(&self) -> usize {
2337 self.events.len()
2338 }
2339
2340 pub fn span_len(&self) -> usize {
2341 self.spans.len()
2342 }
2343
2344 pub fn is_empty(&self) -> bool {
2345 self.events.is_empty()
2346 }
2347
2348 pub fn stats(&self) -> EventLogStats {
2349 EventLogStats {
2350 events: self.events.len(),
2351 spans: self.spans.len(),
2352 approx_event_bytes: approx_json_bytes(&self.events),
2353 approx_span_bytes: approx_json_bytes(&self.spans),
2354 }
2355 }
2356
2357 pub fn truncate_events_keep_last(&mut self, keep_last: usize) -> usize {
2358 let removed = truncate_vec_keep_last(&mut self.events, keep_last);
2359 self.trimmed_events += removed as u64;
2360 if removed > 0 {
2361 self.maybe_compact_journal();
2362 }
2363 removed
2364 }
2365
2366 pub fn truncate_spans_keep_last(&mut self, keep_last: usize) -> usize {
2367 truncate_vec_keep_last(&mut self.spans, keep_last)
2368 }
2369
2370 pub fn clear(&mut self) -> EventLogStats {
2375 let removed = self.stats();
2376 self.trimmed_events += removed.events as u64;
2377 self.events.clear();
2378 self.events.shrink_to_fit();
2379 self.spans.clear();
2380 self.spans.shrink_to_fit();
2381 removed
2382 }
2383
2384 pub fn trimmed_events(&self) -> u64 {
2390 self.trimmed_events
2391 }
2392
2393 pub fn cumulative_cost_usd(&self) -> f64 {
2400 self.cumulative_cost_usd
2401 }
2402
2403 pub fn journal_size_bytes(&self) -> Option<u64> {
2407 let path = self.journal_path.as_ref()?;
2408 fs::metadata(path).ok().map(|m| m.len())
2409 }
2410
2411 fn maybe_compact_journal(&mut self) {
2417 if self.journal_path.is_none() {
2418 return;
2419 }
2420 let excess = self.journal_lines.saturating_sub(self.events.len());
2421 if excess >= JOURNAL_COMPACT_MIN_EXCESS && excess.saturating_mul(4) >= self.journal_lines {
2422 self.compact_journal();
2423 }
2424 }
2425
2426 pub fn compact_journal(&mut self) -> bool {
2444 let Some(path) = self.journal_path.clone() else {
2445 return false;
2446 };
2447 let compacted_lines: HashSet<String> = self
2448 .events
2449 .iter()
2450 .filter_map(|event| serde_json::to_string(event).ok())
2451 .collect();
2452 self.journal = None;
2454 let file_name = path
2455 .file_name()
2456 .and_then(|name| name.to_str())
2457 .unwrap_or("journal");
2458 let tmp = path.with_file_name(format!(".{file_name}.compact-{}.tmp", Uuid::new_v4()));
2459 let rewrite = (|| -> std::io::Result<()> {
2460 let file = create_private_file(&tmp)?;
2461 let mut writer = BufWriter::new(file);
2462 for ev in &self.events {
2463 let line = serde_json::to_string(ev).map_err(std::io::Error::other)?;
2464 writeln!(writer, "{line}")?;
2465 }
2466 writer.flush()?;
2467 let file = writer.into_inner().map_err(|error| error.into_error())?;
2468 file.sync_all()?;
2469 revalidate_private_file(&file)?;
2470 drop(file);
2471 atomic_replace_private_file(&tmp, &path)
2472 })();
2473 let ok = match rewrite {
2474 Ok(()) => {
2475 self.journal_lines = self.events.len();
2476 self.critical_pending
2482 .retain(|line| !compacted_lines.contains(line));
2483 true
2484 }
2485 Err(e) => {
2486 let _ = fs::remove_file(&tmp);
2487 tracing::warn!(
2488 path = %path.display(), error = %e,
2489 "car-eventlog: journal compaction failed — journal keeps growing until the next successful compaction"
2490 );
2491 false
2492 }
2493 };
2494 self.journal = Some(JournalWriter::spawn(path));
2495 ok
2496 }
2497
2498 pub fn query(&self, query: &EventQuery) -> Vec<&Event> {
2501 let mut out: Vec<&Event> = self.events.iter().filter(|e| query.matches(e)).collect();
2502 out.reverse(); if let Some(limit) = query.limit.filter(|l| *l > 0) {
2504 out.truncate(limit);
2505 }
2506 out
2507 }
2508
2509 pub fn set_retention(&mut self, policy: Option<RetentionPolicy>) {
2513 self.retention = policy;
2514 }
2515
2516 pub fn retention(&self) -> Option<&RetentionPolicy> {
2518 self.retention.as_ref()
2519 }
2520
2521 pub fn enforce_retention(&mut self, policy: &RetentionPolicy, now: DateTime<Utc>) -> usize {
2529 let before = self.events.len();
2530 if let Some(age) = policy.max_age_secs {
2531 let cutoff = now - chrono::Duration::seconds(age);
2532 self.events.retain(|e| e.timestamp >= cutoff);
2533 }
2534 if let Some(max) = policy.max_events {
2535 truncate_vec_keep_last(&mut self.events, max);
2536 }
2537 let removed = before.saturating_sub(self.events.len());
2538 self.trimmed_events += removed as u64;
2539 if removed > 0 {
2540 self.maybe_compact_journal();
2541 }
2542 removed
2543 }
2544
2545 pub fn filter(&self, kind: Option<&EventKind>, action_id: Option<&str>) -> Vec<&Event> {
2546 self.events
2547 .iter()
2548 .filter(|e| {
2549 if let Some(k) = kind {
2550 if &e.kind != k {
2551 return false;
2552 }
2553 }
2554 if let Some(aid) = action_id {
2555 if e.action_id.as_deref() != Some(aid) {
2556 return false;
2557 }
2558 }
2559 true
2560 })
2561 .collect()
2562 }
2563
2564 pub fn begin_span(
2566 &mut self,
2567 name: &str,
2568 trace_id: &str,
2569 parent_span_id: Option<&str>,
2570 attributes: HashMap<String, Value>,
2571 ) -> String {
2572 let span_id = Uuid::new_v4().to_string();
2573 let span = Span {
2574 trace_id: trace_id.to_string(),
2575 span_id: span_id.clone(),
2576 parent_span_id: parent_span_id.map(|s| s.to_string()),
2577 name: name.to_string(),
2578 start_time: Utc::now(),
2579 end_time: None,
2580 status: SpanStatus::Unset,
2581 attributes,
2582 };
2583 self.spans.push(span);
2584 span_id
2585 }
2586
2587 pub fn end_span(&mut self, span_id: &str, status: SpanStatus) {
2589 if let Some(span) = self.spans.iter_mut().find(|s| s.span_id == span_id) {
2590 span.end_time = Some(Utc::now());
2591 span.status = status;
2592 }
2593 }
2594
2595 pub fn spans(&self) -> Vec<Span> {
2597 self.spans.clone()
2598 }
2599
2600 pub fn export_traces(&self) -> String {
2602 let mut traces: HashMap<&str, Vec<&Span>> = HashMap::new();
2604 for span in &self.spans {
2605 traces.entry(span.trace_id.as_str()).or_default().push(span);
2606 }
2607
2608 let resource_spans: Vec<Value> = traces.into_values().map(|spans| {
2609 let scope_spans = spans
2610 .iter()
2611 .map(|s| {
2612 let mut span_obj = serde_json::json!({
2613 "traceId": s.trace_id,
2614 "spanId": s.span_id,
2615 "name": s.name,
2616 "startTimeUnixNano": s.start_time.timestamp_nanos_opt().unwrap_or(0).to_string(),
2617 "status": {
2618 "code": match s.status {
2619 SpanStatus::Ok => 1,
2620 SpanStatus::Error => 2,
2621 SpanStatus::Unset => 0,
2622 }
2623 },
2624 "attributes": s.attributes.iter().map(|(k, v)| {
2625 serde_json::json!({
2626 "key": k,
2627 "value": { "stringValue": v.to_string() }
2628 })
2629 }).collect::<Vec<_>>(),
2630 });
2631
2632 if let Some(ref parent) = s.parent_span_id {
2633 span_obj.as_object_mut().unwrap().insert(
2634 "parentSpanId".to_string(),
2635 Value::from(parent.as_str()),
2636 );
2637 }
2638 if let Some(end) = s.end_time {
2639 span_obj.as_object_mut().unwrap().insert(
2640 "endTimeUnixNano".to_string(),
2641 Value::from(end.timestamp_nanos_opt().unwrap_or(0).to_string()),
2642 );
2643 }
2644
2645 span_obj
2646 })
2647 .collect::<Vec<_>>();
2648
2649 serde_json::json!({
2650 "resource": {
2651 "attributes": [
2652 { "key": "service.name", "value": { "stringValue": "car-runtime" } }
2653 ]
2654 },
2655 "scopeSpans": [{
2656 "scope": { "name": "car-eventlog" },
2657 "spans": scope_spans
2658 }]
2659 })
2660 })
2661 .collect();
2662
2663 serde_json::to_string(&serde_json::json!({
2664 "resourceSpans": resource_spans
2665 }))
2666 .unwrap_or_else(|_| "{}".to_string())
2667 }
2668
2669 pub fn load(path: &Path) -> std::io::Result<Self> {
2671 Self::load_with_writer(path, JournalWriter::spawn(path.to_path_buf()))
2672 }
2673
2674 pub fn load_read_only(path: &Path) -> std::io::Result<Self> {
2679 Self::load_from_journal(path, None, false)
2680 }
2681
2682 #[doc(hidden)]
2685 pub fn load_with_journal_failure_injector(
2686 path: &Path,
2687 failures: JournalFailureInjector,
2688 ) -> std::io::Result<Self> {
2689 Self::load_with_writer(
2690 path,
2691 JournalWriter::spawn_with_injector(path.to_path_buf(), failures),
2692 )
2693 }
2694
2695 fn load_with_writer(path: &Path, writer: JournalWriter) -> std::io::Result<Self> {
2696 Self::load_from_journal(path, Some(writer), true)
2697 }
2698
2699 fn load_from_journal(
2700 path: &Path,
2701 writer: Option<JournalWriter>,
2702 repair_torn_tail: bool,
2703 ) -> std::io::Result<Self> {
2704 let file = fs::File::open(path)?;
2705 let mut reader = BufReader::new(file);
2706 let mut events = Vec::new();
2707 let mut event_lines = Vec::new();
2708 let mut line_bytes = Vec::new();
2709 let mut line_number = 0usize;
2710 let mut torn_tail_line = None;
2711
2712 loop {
2713 line_bytes.clear();
2714 let bytes_read = reader.read_until(b'\n', &mut line_bytes)?;
2715 if bytes_read == 0 {
2716 break;
2717 }
2718 line_number += 1;
2719 let terminated = line_bytes.ends_with(b"\n");
2720 if line_bytes.iter().all(|byte| byte.is_ascii_whitespace()) {
2721 if !terminated {
2722 torn_tail_line = Some(line_number);
2723 }
2724 continue;
2725 }
2726 match serde_json::from_slice::<Event>(&line_bytes) {
2727 Ok(event) => {
2728 events.push(event);
2729 event_lines.push(line_number);
2730 if !terminated {
2731 torn_tail_line = Some(line_number);
2732 }
2733 }
2734 Err(_) if !terminated => {
2735 torn_tail_line = Some(line_number);
2740 }
2741 Err(error) => {
2742 return Err(invalid_journal_data(path, line_number, error));
2743 }
2744 }
2745 if !terminated {
2746 break;
2747 }
2748 }
2749 drop(reader);
2750
2751 let last_hash = events.last().and_then(|e| e.hash.clone());
2757 let hash_chaining = last_hash.is_some();
2758 let cumulative_cost_usd = events.iter().filter_map(Event::cost_usd).sum();
2763 let journal_lines = events.len();
2764 let loaded = Self {
2765 events,
2766 spans: Vec::new(),
2767 journal: writer,
2770 hash_chaining,
2771 last_hash,
2772 retention: None,
2773 journal_path: repair_torn_tail.then(|| path.to_path_buf()),
2774 journal_lines,
2775 trimmed_events: 0,
2776 cumulative_cost_usd,
2777 active_binding: None,
2778 critical_pending: HashSet::new(),
2779 };
2780 if let Err(index) = loaded.verify_chain() {
2781 let source_line = event_lines.get(index).copied().unwrap_or(index + 1);
2782 return Err(invalid_journal_data(
2783 path,
2784 source_line,
2785 "hash chain integrity check failed",
2786 ));
2787 }
2788 if let Some(line_number) = torn_tail_line {
2789 if !repair_torn_tail {
2790 return Err(torn_journal_tail(path, line_number));
2791 }
2792 rewrite_loaded_journal(path, &loaded.events)?;
2793 }
2794 Ok(loaded)
2795 }
2796}
2797
2798fn torn_journal_tail(path: &Path, line_number: usize) -> std::io::Error {
2799 std::io::Error::new(
2800 std::io::ErrorKind::UnexpectedEof,
2801 format!(
2802 "event journal torn tail: path={} line={} reason=unterminated final record",
2803 path.display(),
2804 line_number
2805 ),
2806 )
2807}
2808
2809fn invalid_journal_data(
2810 path: &Path,
2811 line_number: usize,
2812 reason: impl std::fmt::Display,
2813) -> std::io::Error {
2814 std::io::Error::new(
2815 std::io::ErrorKind::InvalidData,
2816 format!(
2817 "event journal corruption: path={} line={} reason={reason}",
2818 path.display(),
2819 line_number
2820 ),
2821 )
2822}
2823
2824fn rewrite_loaded_journal(path: &Path, events: &[Event]) -> std::io::Result<()> {
2825 let file_name = path
2826 .file_name()
2827 .and_then(|name| name.to_str())
2828 .unwrap_or("journal");
2829 let temp = path.with_file_name(format!(".{file_name}.recover-{}.tmp", Uuid::new_v4()));
2830 let result = (|| {
2831 let file = create_private_file(&temp)?;
2832 let mut output = BufWriter::new(file);
2833 for event in events {
2834 let line = serde_json::to_string(event).map_err(std::io::Error::other)?;
2835 writeln!(output, "{line}")?;
2836 }
2837 output.flush()?;
2838 let file = output.into_inner().map_err(|error| error.into_error())?;
2839 file.sync_all()?;
2840 revalidate_private_file(&file)?;
2841 drop(file);
2842 atomic_replace_private_file(&temp, path)
2843 })();
2844 if result.is_err() {
2845 let _ = fs::remove_file(&temp);
2846 }
2847 result
2848}
2849
2850fn approx_json_bytes<T: Serialize>(value: &T) -> usize {
2851 serde_json::to_vec(value)
2852 .map(|bytes| bytes.len())
2853 .unwrap_or(0)
2854}
2855
2856fn truncate_vec_keep_last<T>(items: &mut Vec<T>, keep_last: usize) -> usize {
2857 let len = items.len();
2858 if len <= keep_last {
2859 return 0;
2860 }
2861 let removed = len - keep_last;
2862 items.drain(..removed);
2863 items.shrink_to_fit();
2864 removed
2865}
2866
2867impl Default for EventLog {
2868 fn default() -> Self {
2869 Self::new()
2870 }
2871}
2872
2873#[cfg(test)]
2874mod tests {
2875 use super::*;
2876
2877 struct ThreadWake(std::thread::Thread);
2878
2879 impl std::task::Wake for ThreadWake {
2880 fn wake(self: Arc<Self>) {
2881 self.0.unpark();
2882 }
2883
2884 fn wake_by_ref(self: &Arc<Self>) {
2885 self.0.unpark();
2886 }
2887 }
2888
2889 fn test_waker() -> std::task::Waker {
2890 std::task::Waker::from(Arc::new(ThreadWake(std::thread::current())))
2891 }
2892
2893 fn block_on_test_future<F: std::future::Future>(future: F) -> F::Output {
2894 let waker = test_waker();
2895 let mut context = std::task::Context::from_waker(&waker);
2896 let mut future = Box::pin(future);
2897 loop {
2898 match future.as_mut().poll(&mut context) {
2899 std::task::Poll::Ready(output) => return output,
2900 std::task::Poll::Pending => std::thread::park(),
2901 }
2902 }
2903 }
2904
2905 #[test]
2906 fn append_and_read() {
2907 let mut log = EventLog::new();
2908 log.append(
2909 EventKind::ProposalReceived,
2910 None,
2911 Some("p1"),
2912 [("source".to_string(), Value::from("test"))].into(),
2913 );
2914 assert_eq!(log.len(), 1);
2915 assert_eq!(log.events()[0].kind, EventKind::ProposalReceived);
2916 }
2917
2918 #[test]
2919 fn query_filters_by_kind_data_and_time() {
2920 let mut log = EventLog::new();
2921 log.append(
2922 EventKind::PermissionDecision,
2923 Some("a1"),
2924 None,
2925 [
2926 ("caller".to_string(), Value::from("alice")),
2927 ("tool".to_string(), Value::from("shell")),
2928 ]
2929 .into(),
2930 );
2931 log.append(
2932 EventKind::PermissionDecision,
2933 Some("a2"),
2934 None,
2935 [
2936 ("caller".to_string(), Value::from("bob")),
2937 ("tool".to_string(), Value::from("shell")),
2938 ]
2939 .into(),
2940 );
2941 log.append(
2942 EventKind::StateChanged,
2943 Some("a3"),
2944 None,
2945 Default::default(),
2946 );
2947
2948 let q = EventQuery {
2950 kinds: vec![EventKind::PermissionDecision],
2951 ..Default::default()
2952 };
2953 assert_eq!(log.query(&q).len(), 2);
2954
2955 let q = EventQuery {
2957 data_matches: [("caller".to_string(), "alice".to_string())].into(),
2958 ..Default::default()
2959 };
2960 let hits = log.query(&q);
2961 assert_eq!(hits.len(), 1);
2962 assert_eq!(hits[0].action_id.as_deref(), Some("a1"));
2963
2964 let q = EventQuery {
2966 kinds: vec![EventKind::PermissionDecision],
2967 data_matches: [("tool".to_string(), "shell".to_string())].into(),
2968 limit: Some(1),
2969 ..Default::default()
2970 };
2971 let hits = log.query(&q);
2973 assert_eq!(hits.len(), 1);
2974 assert_eq!(hits[0].action_id.as_deref(), Some("a2"));
2975 }
2976
2977 #[test]
2978 fn cost_by_agent_folds_metered_events() {
2979 let mut log = EventLog::new();
2980 log.append_metered(
2981 EventKind::InferenceMetered,
2982 None,
2983 None,
2984 [("agent".to_string(), Value::from("researcher"))].into(),
2985 Metrics {
2986 tokens_in: Some(100),
2987 tokens_out: Some(50),
2988 cost_usd: Some(2.0),
2989 ..Default::default()
2990 },
2991 );
2992 log.append_metered(
2993 EventKind::InferenceMetered,
2994 None,
2995 None,
2996 [("agent".to_string(), Value::from("researcher"))].into(),
2997 Metrics {
2998 tokens_in: Some(10),
2999 tokens_out: Some(5),
3000 cost_usd: Some(0.2),
3001 ..Default::default()
3002 },
3003 );
3004 log.append_metered(
3005 EventKind::InferenceMetered,
3006 None,
3007 None,
3008 [("agent".to_string(), Value::from("coordinator"))].into(),
3009 Metrics {
3010 cost_usd: Some(0.5),
3011 ..Default::default()
3012 },
3013 );
3014 let report = log.cost_by_agent();
3015 assert_eq!(report.len(), 2);
3016 assert_eq!(report[0].agent, "coordinator");
3018 assert_eq!(report[0].cost_usd, 0.5);
3019 assert_eq!(report[1].agent, "researcher");
3020 assert_eq!(report[1].calls, 2);
3021 assert_eq!(report[1].tokens_in, 110);
3022 assert_eq!(report[1].tokens_out, 55);
3023 assert!((report[1].cost_usd - 2.2).abs() < 1e-9);
3024 }
3025
3026 #[test]
3027 fn auto_retention_caps_event_count() {
3028 let mut log = EventLog::new();
3029 log.set_retention(Some(RetentionPolicy {
3030 max_events: Some(3),
3031 max_age_secs: None,
3032 }));
3033 for i in 0..10 {
3034 log.append(
3035 EventKind::StateChanged,
3036 Some(&format!("a{i}")),
3037 None,
3038 Default::default(),
3039 );
3040 }
3041 assert_eq!(log.len(), 3);
3043 assert_eq!(log.events()[0].action_id.as_deref(), Some("a7"));
3044 assert_eq!(log.events()[2].action_id.as_deref(), Some("a9"));
3045 }
3046
3047 #[test]
3048 fn enforce_retention_drops_old_by_age() {
3049 let mut log = EventLog::new();
3050 log.append(
3052 EventKind::StateChanged,
3053 Some("old"),
3054 None,
3055 Default::default(),
3056 );
3057 log.events[0].timestamp = Utc::now() - chrono::Duration::seconds(3600);
3058 log.append(
3059 EventKind::StateChanged,
3060 Some("fresh"),
3061 None,
3062 Default::default(),
3063 );
3064
3065 let removed = log.enforce_retention(
3066 &RetentionPolicy {
3067 max_events: None,
3068 max_age_secs: Some(60),
3069 },
3070 Utc::now(),
3071 );
3072 assert_eq!(removed, 1);
3073 assert_eq!(log.len(), 1);
3074 assert_eq!(log.events()[0].action_id.as_deref(), Some("fresh"));
3075 }
3076
3077 #[test]
3078 fn retention_trims_are_counted() {
3079 let mut log = EventLog::new();
3080 log.set_retention(Some(RetentionPolicy {
3081 max_events: Some(2),
3082 max_age_secs: None,
3083 }));
3084 for i in 0..5 {
3085 log.append(
3086 EventKind::StateChanged,
3087 Some(&format!("a{i}")),
3088 None,
3089 Default::default(),
3090 );
3091 }
3092 assert_eq!(log.trimmed_events(), 3);
3093 assert_eq!(log.truncate_events_keep_last(1), 1);
3094 assert_eq!(log.trimmed_events(), 4);
3095 log.clear();
3096 assert_eq!(log.trimmed_events(), 5);
3097 }
3098
3099 #[test]
3100 fn cumulative_cost_is_monotonic_across_trims_and_reload() {
3101 let dir = tempfile::tempdir().unwrap();
3102 let journal = dir.path().join("cost.jsonl");
3103 {
3104 let mut log = EventLog::with_journal(journal.clone());
3105 log.set_retention(Some(RetentionPolicy {
3106 max_events: Some(1),
3107 max_age_secs: None,
3108 }));
3109 for _ in 0..4 {
3110 log.append_metered(
3111 EventKind::InferenceMetered,
3112 None,
3113 None,
3114 Default::default(),
3115 Metrics {
3116 cost_usd: Some(2.5),
3117 ..Default::default()
3118 },
3119 );
3120 }
3121 assert_eq!(log.len(), 1);
3123 assert!((log.cumulative_cost_usd() - 10.0).abs() < 1e-9);
3124 }
3125 let reloaded = EventLog::load(&journal).unwrap();
3128 assert!((reloaded.cumulative_cost_usd() - 10.0).abs() < 1e-9);
3129 }
3130
3131 #[test]
3132 fn journal_compaction_rewrites_to_retained_set() {
3133 let dir = tempfile::tempdir().unwrap();
3138 let journal = dir.path().join("compact.jsonl");
3139 let keep = 16usize;
3140 let total = keep + JOURNAL_COMPACT_MIN_EXCESS + 8;
3141 {
3142 let mut log = EventLog::with_journal(journal.clone());
3143 log.set_retention(Some(RetentionPolicy {
3144 max_events: Some(keep),
3145 max_age_secs: None,
3146 }));
3147 for i in 0..total {
3148 log.append(
3149 EventKind::ActionSucceeded,
3150 Some(&format!("a{i}")),
3151 None,
3152 HashMap::new(),
3153 );
3154 }
3155 assert_eq!(log.len(), keep);
3156 assert!(log.journal_size_bytes().unwrap_or(0) > 0);
3157 } let reloaded = EventLog::load(&journal).unwrap();
3161 assert!(
3162 reloaded.len() < total,
3163 "journal must have been compacted (got {} lines)",
3164 reloaded.len()
3165 );
3166 assert_eq!(
3168 reloaded.events().last().unwrap().action_id.as_deref(),
3169 Some(format!("a{}", total - 1).as_str())
3170 );
3171 }
3172
3173 #[test]
3174 fn compact_journal_preserves_hash_chain_of_retained_tail() {
3175 let dir = tempfile::tempdir().unwrap();
3179 let journal = dir.path().join("chained.jsonl");
3180 {
3181 let mut log = EventLog::with_journal(journal.clone()).with_hash_chaining();
3182 for i in 0..20 {
3183 log.append(
3184 EventKind::ActionSucceeded,
3185 Some(&format!("a{i}")),
3186 None,
3187 HashMap::new(),
3188 );
3189 }
3190 log.truncate_events_keep_last(5);
3192 assert!(log.compact_journal(), "compaction must succeed");
3193 log.append(
3196 EventKind::ActionSucceeded,
3197 Some("post"),
3198 None,
3199 HashMap::new(),
3200 );
3201 }
3202 let reloaded = EventLog::load(&journal).unwrap();
3203 assert_eq!(reloaded.len(), 6);
3204 assert_eq!(reloaded.verify_chain(), Ok(6), "retained tail must verify");
3205 assert_eq!(reloaded.events()[0].action_id.as_deref(), Some("a15"));
3206 assert_eq!(reloaded.events()[5].action_id.as_deref(), Some("post"));
3207 }
3208
3209 #[test]
3210 fn compact_journal_without_journal_is_noop() {
3211 let mut log = EventLog::new();
3212 log.append(EventKind::StateChanged, Some("a"), None, Default::default());
3213 assert!(!log.compact_journal());
3214 assert_eq!(log.journal_size_bytes(), None);
3215 }
3216
3217 #[test]
3218 fn chaining_off_by_default_no_hashes() {
3219 let mut log = EventLog::new();
3220 log.append(
3221 EventKind::ActionSucceeded,
3222 Some("a1"),
3223 Some("p1"),
3224 HashMap::new(),
3225 );
3226 assert!(!log.hash_chaining_enabled());
3227 assert!(log.events()[0].hash.is_none());
3228 assert!(log.events()[0].prev_hash.is_none());
3229 assert_eq!(log.verify_chain(), Ok(0));
3231 }
3232
3233 #[test]
3234 fn hash_chain_verifies_clean_log() {
3235 let mut log = EventLog::new().with_hash_chaining();
3236 for i in 0..5 {
3237 log.append(
3238 EventKind::ActionSucceeded,
3239 Some(&format!("a{i}")),
3240 Some("p"),
3241 [("i".to_string(), Value::from(i))].into(),
3242 );
3243 }
3244 assert!(log.events().iter().all(|e| e.hash.is_some()));
3246 assert_eq!(log.verify_chain(), Ok(5));
3247 assert_eq!(log.events()[0].prev_hash.as_deref(), Some(""));
3249 for w in log.events().windows(2) {
3251 assert_eq!(w[1].prev_hash, w[0].hash);
3252 }
3253 }
3254
3255 #[test]
3256 fn tampering_with_data_breaks_chain() {
3257 let mut log = EventLog::new().with_hash_chaining();
3258 for i in 0..4 {
3259 log.append(
3260 EventKind::ActionSucceeded,
3261 Some(&format!("a{i}")),
3262 Some("p"),
3263 [("v".to_string(), Value::from(i))].into(),
3264 );
3265 }
3266 assert_eq!(log.verify_chain(), Ok(4));
3267 log.events[2].data.insert("v".to_string(), Value::from(999));
3269 assert_eq!(log.verify_chain(), Err(2));
3271 }
3272
3273 #[test]
3274 fn deleting_an_event_breaks_chain() {
3275 let mut log = EventLog::new().with_hash_chaining();
3276 for i in 0..4 {
3277 log.append(
3278 EventKind::ActionSucceeded,
3279 Some(&format!("a{i}")),
3280 Some("p"),
3281 HashMap::new(),
3282 );
3283 }
3284 log.events.remove(1);
3287 assert_eq!(log.verify_chain(), Err(1));
3288 }
3289
3290 #[test]
3291 fn chain_survives_serialize_roundtrip() {
3292 let mut log = EventLog::new().with_hash_chaining();
3293 for i in 0..3 {
3294 log.append(
3295 EventKind::PermissionDecision,
3296 Some(&format!("a{i}")),
3297 Some("p"),
3298 [
3299 ("decision".to_string(), Value::from("allow")),
3300 ("nested".to_string(), serde_json::json!({"z": 1, "a": 2})),
3301 ]
3302 .into(),
3303 );
3304 }
3305 let lines: Vec<String> = log
3308 .events()
3309 .iter()
3310 .map(|e| serde_json::to_string(e).unwrap())
3311 .collect();
3312 let mut rebuilt = EventLog::new();
3313 for line in &lines {
3314 rebuilt.events.push(serde_json::from_str(line).unwrap());
3315 }
3316 assert_eq!(rebuilt.verify_chain(), Ok(3));
3317 }
3318
3319 #[test]
3320 fn chain_survives_journal_load_and_append() {
3321 let dir = tempfile::tempdir().unwrap();
3326 let journal = dir.path().join("chain.jsonl");
3327 {
3328 let mut log = EventLog::with_journal(journal.clone());
3329 log.enable_hash_chaining();
3330 log.append(EventKind::ActionSucceeded, Some("a1"), None, HashMap::new());
3331 log.append(EventKind::ActionSucceeded, Some("a2"), None, HashMap::new());
3332 } {
3335 let mut log = EventLog::load(&journal).unwrap();
3336 assert!(
3337 log.hash_chaining_enabled(),
3338 "loading a chained tail re-enables chaining"
3339 );
3340 log.append(EventKind::ActionSucceeded, Some("a3"), None, HashMap::new());
3341 assert_eq!(log.verify_chain(), Ok(3), "post-load append stays chained");
3342 }
3343
3344 let reloaded = EventLog::load(&journal).unwrap();
3346 assert_eq!(reloaded.len(), 3);
3347 assert_eq!(reloaded.verify_chain(), Ok(3));
3348
3349 let plain = dir.path().join("plain.jsonl");
3351 {
3352 let mut log = EventLog::with_journal(plain.clone());
3353 log.append(EventKind::ActionSucceeded, Some("a1"), None, HashMap::new());
3354 }
3355 let loaded = EventLog::load(&plain).unwrap();
3356 assert!(!loaded.hash_chaining_enabled(), "unchained tail stays off");
3357 }
3358
3359 #[test]
3360 fn metered_event_carries_metrics_in_data() {
3361 let mut log = EventLog::new();
3362 log.append_metered(
3363 EventKind::ActionSucceeded,
3364 Some("a1"),
3365 Some("p1"),
3366 [("tool".to_string(), Value::from("search"))].into(),
3367 Metrics::inference(120, 45, Some(0.0012)).with_duration(83.0),
3368 );
3369 let ev = &log.events()[0];
3370 assert_eq!(ev.data.get("tool").unwrap(), "search");
3372 assert_eq!(ev.duration_ms(), Some(83.0));
3373 assert_eq!(ev.tokens_in(), Some(120));
3374 assert_eq!(ev.tokens_out(), Some(45));
3375 assert_eq!(ev.cost_usd(), Some(0.0012));
3376 }
3377
3378 #[test]
3379 fn metrics_totals_sum_across_events() {
3380 let mut log = EventLog::new();
3381 log.append_metered(
3382 EventKind::ActionSucceeded,
3383 Some("a1"),
3384 None,
3385 HashMap::new(),
3386 Metrics::latency(50.0),
3387 );
3388 log.append_metered(
3389 EventKind::ActionSucceeded,
3390 Some("a2"),
3391 None,
3392 HashMap::new(),
3393 Metrics::inference(100, 20, Some(0.5)).with_duration(70.0),
3394 );
3395 log.append(EventKind::ProposalReceived, None, None, HashMap::new());
3397
3398 let t = log.metrics_totals();
3399 assert_eq!(t.duration_ms, 120.0);
3400 assert_eq!(t.tokens_in, 100);
3401 assert_eq!(t.tokens_out, 20);
3402 assert_eq!(t.tokens, 120);
3403 assert_eq!(t.cost_usd, 0.5);
3404 assert_eq!(t.metered_events, 2);
3405 }
3406
3407 #[test]
3408 fn metrics_totals_counts_raw_appended_duration_key() {
3409 let mut log = EventLog::new();
3414 log.append(
3415 EventKind::ActionSucceeded,
3416 Some("a1"),
3417 None,
3418 [(metric_keys::DURATION_MS.to_string(), Value::from(42.0))].into(),
3419 );
3420 let t = log.metrics_totals();
3421 assert_eq!(t.duration_ms, 42.0);
3422 assert_eq!(t.metered_events, 1);
3423 }
3424
3425 #[test]
3426 fn new_telemetry_event_kinds_serialize_snake_case() {
3427 let json = serde_json::to_string(&EventKind::BranchDecision).unwrap();
3429 assert_eq!(json, "\"branch_decision\"");
3430 let json = serde_json::to_string(&EventKind::AlternativeRejected).unwrap();
3431 assert_eq!(json, "\"alternative_rejected\"");
3432 let json = serde_json::to_string(&EventKind::InferenceMetered).unwrap();
3433 assert_eq!(json, "\"inference_metered\"");
3434 }
3435
3436 #[test]
3437 fn filter_by_kind() {
3438 let mut log = EventLog::new();
3439 log.append(
3440 EventKind::ProposalReceived,
3441 None,
3442 Some("p1"),
3443 HashMap::new(),
3444 );
3445 log.append(
3446 EventKind::ActionValidated,
3447 Some("a1"),
3448 Some("p1"),
3449 HashMap::new(),
3450 );
3451 log.append(
3452 EventKind::ActionSucceeded,
3453 Some("a1"),
3454 Some("p1"),
3455 HashMap::new(),
3456 );
3457
3458 let validated = log.filter(Some(&EventKind::ActionValidated), None);
3459 assert_eq!(validated.len(), 1);
3460 }
3461
3462 #[test]
3463 fn filter_by_action_id() {
3464 let mut log = EventLog::new();
3465 log.append(EventKind::ActionValidated, Some("a1"), None, HashMap::new());
3466 log.append(EventKind::ActionValidated, Some("a2"), None, HashMap::new());
3467
3468 let a1_events = log.filter(None, Some("a1"));
3469 assert_eq!(a1_events.len(), 1);
3470 }
3471
3472 #[test]
3473 fn journal_write_and_reload() {
3474 let dir = tempfile::tempdir().unwrap();
3475 let journal = dir.path().join("events.jsonl");
3476
3477 {
3478 let mut log = EventLog::with_journal(journal.clone());
3479 log.append(
3480 EventKind::ProposalReceived,
3481 None,
3482 Some("p1"),
3483 HashMap::new(),
3484 );
3485 log.append(
3486 EventKind::ActionSucceeded,
3487 Some("a1"),
3488 Some("p1"),
3489 HashMap::new(),
3490 );
3491 }
3492
3493 assert!(journal.exists());
3494
3495 let reloaded = EventLog::load(&journal).unwrap();
3496 assert_eq!(reloaded.len(), 2);
3497 assert_eq!(reloaded.events()[0].kind, EventKind::ProposalReceived);
3498 assert_eq!(reloaded.events()[1].kind, EventKind::ActionSucceeded);
3499 }
3500
3501 #[test]
3502 fn load_rejects_newline_terminated_corrupt_middle_before_later_terminal() {
3503 let dir = tempfile::tempdir().unwrap();
3504 let journal = dir.path().join("corrupt-middle.jsonl");
3505 let started = serde_json::json!({
3506 "kind": "run_started",
3507 "run_id": "run-corrupt",
3508 "client_id": "client-1",
3509 "data": {"agent_id": "daily-continuity-newsroom"},
3510 "timestamp": "2026-08-30T09:30:00Z"
3511 });
3512 let completed = serde_json::json!({
3513 "kind": "run_completed",
3514 "run_id": "run-corrupt",
3515 "client_id": "client-1",
3516 "data": {
3517 "completion_digest": "must-not-be-trusted",
3518 "termination": {"kind": "outcome", "status": "success", "outcome": {}}
3519 },
3520 "timestamp": "2026-08-30T10:00:00Z"
3521 });
3522 fs::write(
3523 &journal,
3524 format!("{started}\n{{this-is-not-json}}\n{completed}\n"),
3525 )
3526 .unwrap();
3527
3528 let error = match EventLog::load(&journal) {
3529 Ok(_) => panic!("newline-terminated middle corruption must fail closed"),
3530 Err(error) => error,
3531 };
3532 assert_eq!(error.kind(), std::io::ErrorKind::InvalidData);
3533 assert!(error.to_string().contains("line=2"), "{error}");
3534 }
3535
3536 #[test]
3537 fn load_repairs_crash_torn_final_record_before_append() {
3538 let dir = tempfile::tempdir().unwrap();
3539 let journal = dir.path().join("torn-tail.jsonl");
3540 let started = serde_json::json!({
3541 "kind": "run_started",
3542 "run_id": "run-torn",
3543 "client_id": "client-1",
3544 "data": {"agent_id": "daily-continuity-newsroom"},
3545 "timestamp": "2026-08-30T09:30:00Z"
3546 });
3547 let mut journal_file = create_private_file(&journal).unwrap();
3548 write!(journal_file, "{started}\n{{\"kind\":\"run_completed\"").unwrap();
3549 drop(journal_file);
3550
3551 {
3552 let mut loaded = EventLog::load(&journal).expect("torn final row is recoverable");
3553 assert_eq!(loaded.len(), 1);
3554 loaded.append(
3555 EventKind::ProposalReceived,
3556 None,
3557 Some("proposal-after-recovery"),
3558 HashMap::new(),
3559 );
3560 }
3561
3562 let bytes = fs::read_to_string(&journal).unwrap();
3563 assert!(!bytes.contains("{\"kind\":\"run_completed\""), "{bytes}");
3564 assert!(bytes.ends_with('\n'));
3565 let reloaded = EventLog::load(&journal).unwrap();
3566 assert_eq!(reloaded.len(), 2);
3567 assert_eq!(reloaded.events()[0].kind, EventKind::RunStarted);
3568 assert_eq!(reloaded.events()[1].kind, EventKind::ProposalReceived);
3569 }
3570
3571 #[test]
3572 fn load_read_only_rejects_a_torn_tail_without_modifying_the_journal() {
3573 let dir = tempfile::tempdir().unwrap();
3574 let journal = dir.path().join("read-only-torn-tail.jsonl");
3575 let started = serde_json::json!({
3576 "kind": "run_started",
3577 "run_id": "run-torn",
3578 "client_id": "client-1",
3579 "data": {"agent_id": "daily-continuity-newsroom"},
3580 "timestamp": "2026-08-30T09:30:00Z"
3581 });
3582 fs::write(&journal, format!("{started}\n{{\"kind\":\"run_completed\"")).unwrap();
3583 let before = fs::read(&journal).unwrap();
3584
3585 let error = match EventLog::load_read_only(&journal) {
3586 Ok(_) => panic!("read-only loading must expose a torn final row"),
3587 Err(error) => error,
3588 };
3589
3590 assert_eq!(error.kind(), std::io::ErrorKind::UnexpectedEof);
3591 assert!(error.to_string().contains("event journal torn tail"));
3592 assert!(error.to_string().contains("line=2"));
3593 assert_eq!(fs::read(&journal).unwrap(), before);
3594 }
3595
3596 #[test]
3597 fn load_rejects_existing_hash_chain_tampering() {
3598 let dir = tempfile::tempdir().unwrap();
3599 let journal = dir.path().join("tampered-chain.jsonl");
3600 {
3601 let mut log = EventLog::with_journal(journal.clone()).with_hash_chaining();
3602 log.append(
3603 EventKind::RunStarted,
3604 None,
3605 None,
3606 [(
3607 "agent_id".to_string(),
3608 Value::from("daily-continuity-newsroom"),
3609 )]
3610 .into(),
3611 );
3612 log.append(EventKind::RunCompleted, None, None, HashMap::new());
3613 }
3614 let mut rows: Vec<Value> = fs::read_to_string(&journal)
3615 .unwrap()
3616 .lines()
3617 .map(|line| serde_json::from_str(line).unwrap())
3618 .collect();
3619 rows[0]["data"]["agent_id"] = Value::from("tampered-agent");
3620 fs::write(
3621 &journal,
3622 rows.iter()
3623 .map(Value::to_string)
3624 .collect::<Vec<_>>()
3625 .join("\n")
3626 + "\n",
3627 )
3628 .unwrap();
3629
3630 let error = match EventLog::load(&journal) {
3631 Ok(_) => panic!("hash-chain tampering must fail closed during load"),
3632 Err(error) => error,
3633 };
3634 assert_eq!(error.kind(), std::io::ErrorKind::InvalidData);
3635 assert!(error.to_string().contains("hash chain"), "{error}");
3636 assert!(error.to_string().contains("line=1"), "{error}");
3637 }
3638
3639 #[cfg(unix)]
3640 fn unix_mode(path: &Path) -> u32 {
3641 use std::os::unix::fs::PermissionsExt;
3642 fs::symlink_metadata(path).unwrap().permissions().mode() & 0o777
3643 }
3644
3645 #[cfg(unix)]
3646 #[test]
3647 fn car_owned_journal_and_created_parents_are_private() {
3648 let root = tempfile::tempdir().unwrap();
3649 let parent = root.path().join("eventlogs").join("session");
3650 let journal = parent.join("events.jsonl");
3651 {
3652 let mut log = EventLog::with_journal(journal.clone());
3653 log.append(EventKind::StateChanged, Some("a1"), None, HashMap::new());
3654 }
3655
3656 assert_eq!(unix_mode(&root.path().join("eventlogs")), 0o700);
3657 assert_eq!(unix_mode(&parent), 0o700);
3658 assert_eq!(unix_mode(&journal), 0o600);
3659 }
3660
3661 #[cfg(unix)]
3662 #[test]
3663 fn append_hardens_preexisting_owned_permissive_journal() {
3664 use std::os::unix::fs::PermissionsExt;
3665
3666 let dir = tempfile::tempdir().unwrap();
3667 let journal = dir.path().join("events.jsonl");
3668 fs::write(&journal, b"").unwrap();
3669 fs::set_permissions(&journal, fs::Permissions::from_mode(0o644)).unwrap();
3670 {
3671 let mut log = EventLog::with_journal(journal.clone());
3672 log.append(EventKind::StateChanged, Some("a1"), None, HashMap::new());
3673 }
3674 assert_eq!(unix_mode(&journal), 0o600);
3675 }
3676
3677 #[cfg(unix)]
3678 #[test]
3679 fn journal_refuses_symlink_and_hardlink_destinations() {
3680 use std::os::unix::fs::symlink;
3681
3682 let dir = tempfile::tempdir().unwrap();
3683 let victim = dir.path().join("victim");
3684 fs::write(&victim, b"unchanged").unwrap();
3685
3686 for journal in [dir.path().join("symlink"), dir.path().join("hardlink")] {
3687 if journal.ends_with("symlink") {
3688 symlink(&victim, &journal).unwrap();
3689 } else {
3690 fs::hard_link(&victim, &journal).unwrap();
3691 }
3692 {
3693 let mut log = EventLog::with_journal(journal);
3694 log.append(EventKind::StateChanged, Some("a1"), None, HashMap::new());
3695 }
3696 assert_eq!(fs::read(&victim).unwrap(), b"unchanged");
3697 }
3698 }
3699
3700 #[cfg(unix)]
3701 #[test]
3702 fn journal_stops_if_the_opened_path_is_substituted() {
3703 let dir = tempfile::tempdir().unwrap();
3704 let journal = dir.path().join("events.jsonl");
3705 let moved = dir.path().join("moved.jsonl");
3706 let mut log = EventLog::with_journal(journal.clone());
3707 log.append(EventKind::StateChanged, Some("first"), None, HashMap::new());
3708
3709 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
3710 while fs::metadata(&journal).map_or(true, |metadata| metadata.len() == 0) {
3711 assert!(
3712 std::time::Instant::now() < deadline,
3713 "first event was not persisted"
3714 );
3715 std::thread::yield_now();
3716 }
3717 fs::rename(&journal, &moved).unwrap();
3718 let _substitute = create_private_file(&journal).unwrap();
3719 log.append(
3720 EventKind::StateChanged,
3721 Some("second"),
3722 None,
3723 HashMap::new(),
3724 );
3725 drop(log);
3726
3727 assert_eq!(EventLog::load(&moved).unwrap().len(), 1);
3728 assert_eq!(fs::metadata(&journal).unwrap().len(), 0);
3729 }
3730
3731 #[cfg(unix)]
3732 #[test]
3733 fn compaction_preserves_private_mode_and_leaves_no_temp_name() {
3734 let dir = tempfile::tempdir().unwrap();
3735 let journal = dir.path().join("events.jsonl");
3736 let mut log = EventLog::with_journal(journal.clone());
3737 for index in 0..4 {
3738 log.append(
3739 EventKind::StateChanged,
3740 Some(&format!("a{index}")),
3741 None,
3742 HashMap::new(),
3743 );
3744 }
3745 log.truncate_events_keep_last(2);
3746 assert!(log.compact_journal());
3747 drop(log);
3748
3749 assert_eq!(unix_mode(&journal), 0o600);
3750 let names: Vec<_> = fs::read_dir(dir.path())
3751 .unwrap()
3752 .map(|entry| entry.unwrap().file_name())
3753 .collect();
3754 assert_eq!(names, vec![journal.file_name().unwrap()]);
3755 }
3756
3757 #[cfg(unix)]
3758 #[test]
3759 fn historical_world_readable_journal_can_be_loaded_without_mutation() {
3760 use std::os::unix::fs::PermissionsExt;
3761
3762 let dir = tempfile::tempdir().unwrap();
3763 let journal = dir.path().join("historical.jsonl");
3764 let event = Event {
3765 kind: EventKind::StateChanged,
3766 run_id: None,
3767 client_id: None,
3768 policy_session_id: None,
3769 action_id: Some("historical".into()),
3770 proposal_id: None,
3771 data: HashMap::new(),
3772 timestamp: Utc::now(),
3773 prev_hash: None,
3774 hash: None,
3775 };
3776 fs::write(
3777 &journal,
3778 format!("{}\n", serde_json::to_string(&event).unwrap()),
3779 )
3780 .unwrap();
3781 fs::set_permissions(&journal, fs::Permissions::from_mode(0o644)).unwrap();
3782
3783 let loaded = EventLog::load(&journal).unwrap();
3784 assert_eq!(loaded.len(), 1);
3785 drop(loaded);
3786 assert_eq!(unix_mode(&journal), 0o644);
3787 }
3788
3789 #[test]
3790 fn journal_not_created_without_appends() {
3791 let dir = tempfile::tempdir().unwrap();
3795 let journal = dir.path().join("no-events.jsonl");
3796 {
3797 let _log = EventLog::with_journal(journal.clone());
3798 }
3801 assert!(
3802 !journal.exists(),
3803 "journal file must not be created when nothing is appended"
3804 );
3805 }
3806
3807 #[test]
3808 fn journal_preserves_order_and_count_under_burst() {
3809 let dir = tempfile::tempdir().unwrap();
3813 let journal = dir.path().join("burst.jsonl");
3814 {
3815 let mut log = EventLog::with_journal(journal.clone());
3816 for i in 0..500 {
3817 log.append(
3818 EventKind::ActionSucceeded,
3819 Some(&format!("a{i}")),
3820 None,
3821 HashMap::new(),
3822 );
3823 }
3824 } let reloaded = EventLog::load(&journal).unwrap();
3827 assert_eq!(reloaded.len(), 500, "no events lost");
3828 for (i, event) in reloaded.events().iter().enumerate() {
3829 assert_eq!(
3830 event.action_id.as_deref(),
3831 Some(format!("a{i}").as_str()),
3832 "order preserved at {i}"
3833 );
3834 }
3835 }
3836
3837 #[test]
3838 fn unopenable_journal_is_best_effort_not_fatal() {
3839 let dir = tempfile::tempdir().unwrap();
3843 let journal = dir.path().join("a-directory");
3844 fs::create_dir(&journal).unwrap(); let mut log = EventLog::with_journal(journal);
3847 log.append(
3848 EventKind::ProposalReceived,
3849 None,
3850 Some("p1"),
3851 HashMap::new(),
3852 );
3853 log.append(EventKind::ActionSucceeded, Some("a1"), None, HashMap::new());
3854 assert_eq!(
3855 log.len(),
3856 2,
3857 "in-memory log unaffected by an unwritable journal"
3858 );
3859 }
3861
3862 #[test]
3863 fn load_then_append_preserves_existing_and_adds() {
3864 let dir = tempfile::tempdir().unwrap();
3865 let journal = dir.path().join("resume.jsonl");
3866 {
3867 let mut log = EventLog::with_journal(journal.clone());
3868 log.append(
3869 EventKind::ProposalReceived,
3870 None,
3871 Some("p1"),
3872 HashMap::new(),
3873 );
3874 }
3875 {
3877 let mut log = EventLog::load(&journal).unwrap();
3878 assert_eq!(log.len(), 1);
3879 log.append(EventKind::ActionSucceeded, Some("a1"), None, HashMap::new());
3880 }
3881 let reloaded = EventLog::load(&journal).unwrap();
3882 assert_eq!(reloaded.len(), 2, "append-mode preserved the loaded line");
3883 assert_eq!(reloaded.events()[0].kind, EventKind::ProposalReceived);
3884 assert_eq!(reloaded.events()[1].kind, EventKind::ActionSucceeded);
3885 }
3886
3887 #[test]
3888 fn event_kind_serializes_snake_case() {
3889 assert_eq!(
3890 serde_json::to_string(&EventKind::ProposalReceived).unwrap(),
3891 "\"proposal_received\""
3892 );
3893 assert_eq!(
3894 serde_json::to_string(&EventKind::StateSnapshot).unwrap(),
3895 "\"state_snapshot\""
3896 );
3897 }
3898
3899 #[test]
3900 fn stats_truncate_and_clear_release_retained_entries() {
3901 let mut log = EventLog::new();
3902 for idx in 0..5 {
3903 log.append(
3904 EventKind::ActionSucceeded,
3905 Some(&format!("a{idx}")),
3906 Some("p1"),
3907 [("payload".to_string(), Value::from("x".repeat(16)))].into(),
3908 );
3909 log.begin_span("action.tool_call", "trace", None, HashMap::new());
3910 }
3911
3912 let stats = log.stats();
3913 assert_eq!(stats.events, 5);
3914 assert_eq!(stats.spans, 5);
3915 assert!(stats.approx_event_bytes > 0);
3916 assert!(stats.approx_span_bytes > 0);
3917
3918 assert_eq!(log.truncate_events_keep_last(2), 3);
3919 assert_eq!(log.truncate_spans_keep_last(1), 4);
3920 assert_eq!(log.len(), 2);
3921 assert_eq!(log.span_len(), 1);
3922 assert_eq!(log.events()[0].action_id.as_deref(), Some("a3"));
3923
3924 let removed = log.clear();
3925 assert_eq!(removed.events, 2);
3926 assert_eq!(removed.spans, 1);
3927 assert_eq!(log.len(), 0);
3928 assert_eq!(log.span_len(), 0);
3929 }
3930
3931 #[test]
3932 fn span_begin_end_lifecycle() {
3933 let mut log = EventLog::new();
3934 let trace_id = "trace-1".to_string();
3935
3936 let span_id = log.begin_span(
3937 "test.operation",
3938 &trace_id,
3939 None,
3940 [("key".to_string(), Value::from("value"))].into(),
3941 );
3942
3943 let spans = log.spans();
3944 assert_eq!(spans.len(), 1);
3945 assert_eq!(spans[0].name, "test.operation");
3946 assert_eq!(spans[0].trace_id, "trace-1");
3947 assert!(spans[0].parent_span_id.is_none());
3948 assert!(spans[0].end_time.is_none());
3949 assert_eq!(spans[0].status, SpanStatus::Unset);
3950
3951 log.end_span(&span_id, SpanStatus::Ok);
3952
3953 let spans = log.spans();
3954 assert!(spans[0].end_time.is_some());
3955 assert_eq!(spans[0].status, SpanStatus::Ok);
3956 }
3957
3958 #[test]
3959 fn span_parent_child_relationship() {
3960 let mut log = EventLog::new();
3961 let trace_id = "trace-2".to_string();
3962
3963 let parent_id = log.begin_span("parent.op", &trace_id, None, HashMap::new());
3964 let child_id = log.begin_span("child.op", &trace_id, Some(&parent_id), HashMap::new());
3965
3966 let spans = log.spans();
3967 assert_eq!(spans.len(), 2);
3968
3969 let child = spans.iter().find(|s| s.span_id == child_id).unwrap();
3970 assert_eq!(child.parent_span_id.as_deref(), Some(parent_id.as_str()));
3971 assert_eq!(child.trace_id, trace_id);
3972
3973 let parent = spans.iter().find(|s| s.span_id == parent_id).unwrap();
3974 assert!(parent.parent_span_id.is_none());
3975 }
3976
3977 #[test]
3978 fn export_traces_produces_valid_json() {
3979 let mut log = EventLog::new();
3980 let trace_id = "trace-3".to_string();
3981
3982 let root = log.begin_span(
3983 "proposal.execute",
3984 &trace_id,
3985 None,
3986 [("proposal_id".to_string(), Value::from("p1"))].into(),
3987 );
3988 let child = log.begin_span(
3989 "action.tool_call",
3990 &trace_id,
3991 Some(&root),
3992 [("tool".to_string(), Value::from("read_file"))].into(),
3993 );
3994 log.end_span(&child, SpanStatus::Ok);
3995 log.end_span(&root, SpanStatus::Ok);
3996
3997 let json_str = log.export_traces();
3998 let parsed: Value =
3999 serde_json::from_str(&json_str).expect("export_traces must produce valid JSON");
4000
4001 let resource_spans = parsed["resourceSpans"].as_array().unwrap();
4002 assert_eq!(resource_spans.len(), 1);
4003
4004 let scope_spans = &resource_spans[0]["scopeSpans"][0]["spans"];
4005 let spans_arr = scope_spans.as_array().unwrap();
4006 assert_eq!(spans_arr.len(), 2);
4007
4008 for span in spans_arr {
4010 assert!(span.get("traceId").is_some());
4011 assert!(span.get("spanId").is_some());
4012 assert!(span.get("name").is_some());
4013 assert!(span.get("startTimeUnixNano").is_some());
4014 assert!(span.get("endTimeUnixNano").is_some());
4015 assert!(span.get("status").is_some());
4016 }
4017
4018 let child_span = spans_arr
4020 .iter()
4021 .find(|s| s["name"] == "action.tool_call")
4022 .unwrap();
4023 assert!(child_span.get("parentSpanId").is_some());
4024 }
4025
4026 #[test]
4027 fn span_status_set_on_error() {
4028 let mut log = EventLog::new();
4029 let trace_id = "trace-4".to_string();
4030
4031 let span_id = log.begin_span("failing.op", &trace_id, None, HashMap::new());
4032 log.end_span(&span_id, SpanStatus::Error);
4033
4034 let spans = log.spans();
4035 assert_eq!(spans[0].status, SpanStatus::Error);
4036 assert!(spans[0].end_time.is_some());
4037 }
4038
4039 #[test]
4040 fn active_run_binding_stamps_every_new_event_and_rejects_conflicts() {
4041 let mut log = EventLog::new();
4042 log.bind_run("run-a", "client-a")
4043 .expect("first active run binds");
4044 log.bind_policy_session("policy-session-a")
4045 .expect("CAR-minted policy session binds inside the run");
4046
4047 log.append(
4048 EventKind::ProposalReceived,
4049 None,
4050 Some("same-proposal"),
4051 HashMap::new(),
4052 );
4053 log.append(
4054 EventKind::ActionSucceeded,
4055 Some("action-a"),
4056 Some("same-proposal"),
4057 HashMap::new(),
4058 );
4059
4060 for event in log.events() {
4061 assert_eq!(event.run_id.as_deref(), Some("run-a"));
4062 assert_eq!(event.client_id.as_deref(), Some("client-a"));
4063 assert_eq!(event.policy_session_id.as_deref(), Some("policy-session-a"));
4064 }
4065 assert!(log.bind_run("run-b", "client-a").is_err());
4066 assert!(log.bind_run("run-a", "client-b").is_err());
4067 assert!(log.clear_run_binding("run-b", "client-a").is_err());
4068 assert_eq!(
4069 log.active_run_binding(),
4070 Some(("run-a", "client-a", Some("policy-session-a")))
4071 );
4072
4073 log.clear_policy_session("policy-session-a")
4074 .expect("exact policy session clears");
4075 log.clear_run_binding("run-a", "client-a")
4076 .expect("exact active run clears");
4077 log.append(
4078 EventKind::ProposalReceived,
4079 None,
4080 Some("unbound-legacy"),
4081 HashMap::new(),
4082 );
4083 let legacy = log.events().last().unwrap();
4084 assert!(legacy.run_id.is_none());
4085 assert!(legacy.client_id.is_none());
4086 assert!(legacy.policy_session_id.is_none());
4087 }
4088
4089 #[test]
4090 fn historical_event_without_binding_fields_still_deserializes() {
4091 let historical = r#"{"kind":"proposal_received","proposal_id":"p-old","data":{},"timestamp":"2026-01-02T03:04:05Z"}"#;
4092 let event: Event = serde_json::from_str(historical).expect("historical event replays");
4093 assert!(event.run_id.is_none());
4094 assert!(event.client_id.is_none());
4095 assert!(event.policy_session_id.is_none());
4096 }
4097
4098 #[test]
4099 fn async_acknowledgement_timeout_wins_after_expiry_removal() {
4100 let manager = AsyncAcknowledgementManager::new(1);
4101 let barrier = AsyncAcknowledgementExpiryBarrier::new();
4102 manager.pause_next_expiry_after_removal(barrier.clone());
4103 let reservation = manager.reserve(Duration::from_millis(10)).unwrap();
4104 let sender = reservation.sender();
4105 let acknowledgement = reservation.into_future();
4106
4107 barrier.wait_until_removed();
4108 sender.send(Ok(()));
4109 let result = block_on_test_future(acknowledgement);
4110 barrier.allow_timeout_completion();
4111 manager.shutdown();
4112
4113 assert!(matches!(
4114 result,
4115 Err(CriticalPostAcceptanceError::AcknowledgementTimedOut { .. })
4116 ));
4117 }
4118
4119 #[test]
4120 fn async_acknowledgement_capacity_is_atomic_under_concurrent_reservation() {
4121 let manager = Arc::new(AsyncAcknowledgementManager::new(1));
4122 let barrier = Arc::new(std::sync::Barrier::new(3));
4123 let handles: Vec<_> = (0..2)
4124 .map(|_| {
4125 let manager = manager.clone();
4126 let barrier = barrier.clone();
4127 std::thread::spawn(move || {
4128 let reservation = manager.reserve(MAX_CRITICAL_ACKNOWLEDGEMENT_TIMEOUT);
4129 barrier.wait();
4130 reservation
4131 })
4132 })
4133 .collect();
4134
4135 barrier.wait();
4136 let results: Vec<_> = handles
4137 .into_iter()
4138 .map(|handle| handle.join().unwrap())
4139 .collect();
4140 assert_eq!(results.iter().filter(|result| result.is_ok()).count(), 1);
4141 assert_eq!(
4142 results
4143 .iter()
4144 .filter(|result| {
4145 matches!(
4146 result,
4147 Err(CriticalPreAcceptanceError::CapacityExhausted { capacity: 1 })
4148 )
4149 })
4150 .count(),
4151 1
4152 );
4153 drop(results);
4154 manager.shutdown();
4155 }
4156
4157 #[test]
4158 fn async_acknowledgement_completed_before_future_construction_is_observed() {
4159 let manager = AsyncAcknowledgementManager::new(1);
4160 let reservation = manager.reserve(Duration::from_secs(1)).unwrap();
4161 reservation.sender().send(Ok(()));
4162 let acknowledgement = reservation.into_future();
4163
4164 assert!(block_on_test_future(acknowledgement).is_ok());
4165 manager.shutdown();
4166 }
4167
4168 #[test]
4169 fn journal_writer_drop_completes_pending_async_acknowledgement() {
4170 let dir = tempfile::tempdir().unwrap();
4171 let failures = JournalFailureInjector::default();
4172 failures.fail_next(JournalFailurePoint::HoldAcknowledgement);
4173 let writer = JournalWriter::spawn_with_injector(
4174 dir.path().join("pending-ack-shutdown.jsonl"),
4175 failures.clone(),
4176 );
4177 let reservation = writer
4178 .reserve_async_acknowledgement(MAX_CRITICAL_ACKNOWLEDGEMENT_TIMEOUT)
4179 .unwrap();
4180 let acknowledgement = writer
4181 .enqueue_critical_async("{}".to_string(), false, reservation)
4182 .unwrap();
4183 let deadline = Instant::now() + Duration::from_secs(2);
4184 while failures.held_acknowledgement_count() == 0 {
4185 assert!(
4186 Instant::now() < deadline,
4187 "writer never retained the pending acknowledgement"
4188 );
4189 std::thread::yield_now();
4190 }
4191
4192 drop(writer);
4193 let result = block_on_test_future(acknowledgement);
4194 assert!(matches!(
4195 result,
4196 Err(CriticalPostAcceptanceError::CoordinatorStopped)
4197 ));
4198 failures.release_held_acknowledgements();
4199 }
4200
4201 #[test]
4202 fn prepare_failure_releases_async_acknowledgement_capacity() {
4203 let dir = tempfile::tempdir().unwrap();
4204 let mut log = EventLog::with_journal_failure_injector_and_ack_capacity(
4205 dir.path().join("prepare-failure-capacity.jsonl"),
4206 JournalFailureInjector::default(),
4207 1,
4208 );
4209 let data = HashMap::from([("completion_digest".to_string(), Value::from("7".repeat(64)))]);
4210
4211 let error = block_on_test_future(log.append_critical_async(
4212 EventKind::RunCompleted,
4213 None,
4214 None,
4215 data.clone(),
4216 Duration::from_secs(1),
4217 ))
4218 .expect_err("an unbound critical event must fail during preparation");
4219 assert!(matches!(error, CriticalAppendError::Rejected { .. }));
4220
4221 log.bind_run("run-after-prepare-failure", "client-after-prepare-failure")
4222 .unwrap();
4223 block_on_test_future(log.append_critical_async(
4224 EventKind::RunCompleted,
4225 None,
4226 None,
4227 data,
4228 Duration::from_secs(1),
4229 ))
4230 .expect("the failed preparation must release the only acknowledgement slot");
4231 }
4232
4233 #[test]
4234 fn async_critical_preacceptance_failures_reject_without_fabricating_events() {
4235 let data = HashMap::from([("completion_digest".to_string(), Value::from("f".repeat(64)))]);
4236
4237 let mut no_journal = EventLog::new();
4238 no_journal
4239 .bind_run("run-no-writer", "client-no-writer")
4240 .unwrap();
4241 let error = block_on_test_future(no_journal.append_critical_async(
4242 EventKind::RunCompleted,
4243 None,
4244 None,
4245 data.clone(),
4246 Duration::from_millis(100),
4247 ))
4248 .expect_err("a missing writer must reject before acceptance");
4249 assert!(matches!(error, CriticalAppendError::Rejected { .. }));
4250 assert!(!error.is_retry_safe());
4251 assert!(no_journal.events().is_empty());
4252 assert!(no_journal.critical_pending.is_empty());
4253
4254 let dir = tempfile::tempdir().unwrap();
4255 let unavailable_path = dir.path().join("sender-missing.jsonl");
4256 let mut unavailable = EventLog::with_journal(unavailable_path.clone());
4257 unavailable
4258 .bind_run("run-sender-missing", "client-sender-missing")
4259 .unwrap();
4260 unavailable
4261 .journal
4262 .as_mut()
4263 .unwrap()
4264 .remove_sender_for_test();
4265 let error = block_on_test_future(unavailable.append_critical_async(
4266 EventKind::RunCompleted,
4267 None,
4268 None,
4269 data.clone(),
4270 Duration::from_millis(100),
4271 ))
4272 .expect_err("a missing writer sender must reject before event construction");
4273 assert!(matches!(error, CriticalAppendError::Rejected { .. }));
4274 assert!(unavailable.events().is_empty());
4275 assert!(unavailable.critical_pending.is_empty());
4276 assert!(!unavailable_path.exists());
4277
4278 let stopped_path = dir.path().join("receiver-stopped.jsonl");
4279 let mut stopped = EventLog::with_journal(stopped_path.clone());
4280 stopped
4281 .bind_run("run-receiver-stopped", "client-receiver-stopped")
4282 .unwrap();
4283 stopped.journal.as_mut().unwrap().stop_receiver_for_test();
4284 let error = block_on_test_future(stopped.append_critical_async(
4285 EventKind::RunCompleted,
4286 None,
4287 None,
4288 data,
4289 Duration::from_millis(100),
4290 ))
4291 .expect_err("a stopped writer receiver must reject a failed enqueue");
4292 assert!(matches!(error, CriticalAppendError::Rejected { .. }));
4293 assert!(stopped.events().is_empty());
4294 assert!(stopped.critical_pending.is_empty());
4295 assert!(!stopped_path.exists());
4296 }
4297
4298 #[test]
4299 fn async_critical_rejects_excessive_acknowledgement_duration_before_enqueue() {
4300 let dir = tempfile::tempdir().unwrap();
4301 let path = dir.path().join("excessive-ack-duration.jsonl");
4302 let mut log = EventLog::with_journal(path.clone());
4303 log.bind_run("run-excessive-ack", "client-excessive-ack")
4304 .unwrap();
4305
4306 let error = block_on_test_future(log.append_critical_async(
4307 EventKind::RunCompleted,
4308 None,
4309 None,
4310 HashMap::new(),
4311 MAX_CRITICAL_ACKNOWLEDGEMENT_TIMEOUT + Duration::from_millis(1),
4312 ))
4313 .expect_err("an excessive acknowledgement duration must reject");
4314
4315 assert!(matches!(error, CriticalAppendError::Rejected { .. }));
4316 assert!(log.events().is_empty());
4317 assert!(log.critical_pending.is_empty());
4318 assert!(!path.exists());
4319 }
4320
4321 #[test]
4322 fn async_critical_capacity_exhaustion_rejects_before_enqueue_and_exact_retry_unblocks() {
4323 let dir = tempfile::tempdir().unwrap();
4324 let path = dir.path().join("critical-ack-capacity.jsonl");
4325 let failures = JournalFailureInjector::default();
4326 failures.fail_next(JournalFailurePoint::HoldAcknowledgement);
4327 let mut log = EventLog::with_journal_failure_injector_and_ack_capacity(
4328 path.clone(),
4329 failures.clone(),
4330 1,
4331 );
4332 log.bind_run("run-capacity", "client-capacity").unwrap();
4333 let first_data =
4334 HashMap::from([("completion_digest".to_string(), Value::from("1".repeat(64)))]);
4335 let second_data =
4336 HashMap::from([("completion_digest".to_string(), Value::from("2".repeat(64)))]);
4337
4338 let mut first = Box::pin(log.append_critical_async(
4339 EventKind::RunCompleted,
4340 None,
4341 None,
4342 first_data.clone(),
4343 MAX_CRITICAL_ACKNOWLEDGEMENT_TIMEOUT,
4344 ));
4345 let waker = test_waker();
4346 let mut context = std::task::Context::from_waker(&waker);
4347 assert!(matches!(
4348 first.as_mut().poll(&mut context),
4349 std::task::Poll::Pending
4350 ));
4351 drop(first);
4352 assert_eq!(log.events().len(), 1);
4353 assert_eq!(log.critical_pending.len(), 1);
4354
4355 for attempt in 0..100 {
4356 let error = block_on_test_future(log.append_critical_async(
4357 EventKind::ProposalCompleted,
4358 None,
4359 Some("capacity-rejected"),
4360 second_data.clone(),
4361 Duration::from_millis(100),
4362 ))
4363 .expect_err("exhausted acknowledgement capacity must reject");
4364 let CriticalAppendError::Rejected { reason } = error else {
4365 panic!("attempt {attempt} was not a pre-enqueue rejection");
4366 };
4367 assert!(
4368 reason.contains("acknowledgement capacity is exhausted"),
4369 "attempt {attempt} bypassed capacity admission: {reason}"
4370 );
4371 assert_eq!(log.events().len(), 1);
4372 assert_eq!(log.critical_pending.len(), 1);
4373 }
4374
4375 let hold_deadline = std::time::Instant::now() + Duration::from_secs(2);
4376 while failures.held_acknowledgement_count() == 0 {
4377 assert!(
4378 std::time::Instant::now() < hold_deadline,
4379 "writer never reached the held acknowledgement"
4380 );
4381 std::thread::yield_now();
4382 }
4383 let pending_line = serde_json::to_string(&log.events()[0]).unwrap();
4384 assert_eq!(
4385 fs::read_to_string(&path).unwrap(),
4386 format!("{pending_line}\n")
4387 );
4388 failures.release_held_acknowledgements();
4389
4390 let original_timestamp = log.events()[0].timestamp;
4391 let retried = block_on_test_future(log.append_critical_async(
4392 EventKind::RunCompleted,
4393 None,
4394 None,
4395 first_data,
4396 Duration::from_millis(500),
4397 ))
4398 .expect("the exact cancelled row must reconcile pending state");
4399 assert_eq!(retried.timestamp, original_timestamp);
4400 assert!(log.critical_pending.is_empty());
4401
4402 block_on_test_future(log.append_critical_async(
4403 EventKind::ProposalCompleted,
4404 None,
4405 Some("capacity-rejected"),
4406 second_data,
4407 Duration::from_millis(500),
4408 ))
4409 .expect("a distinct row is allowed after exact reconciliation");
4410 drop(log);
4411
4412 let loaded = EventLog::load(&path).unwrap();
4413 assert_eq!(loaded.events().len(), 2);
4414 assert_eq!(loaded.events()[0].kind, EventKind::RunCompleted);
4415 assert_eq!(loaded.events()[1].kind, EventKind::ProposalCompleted);
4416 }
4417
4418 #[test]
4419 fn async_critical_never_acknowledged_row_remains_exactly_retryable() {
4420 let dir = tempfile::tempdir().unwrap();
4421 let path = dir.path().join("critical-never-acknowledged.jsonl");
4422 let failures = JournalFailureInjector::default();
4423 failures.fail_next(JournalFailurePoint::HoldAcknowledgement);
4424 let mut log = EventLog::with_journal_failure_injector(path.clone(), failures.clone());
4425 log.bind_run("run-never-ack", "client-never-ack").unwrap();
4426 let data = HashMap::from([("completion_digest".to_string(), Value::from("9".repeat(64)))]);
4427
4428 let error = block_on_test_future(log.append_critical_async(
4429 EventKind::RunCompleted,
4430 None,
4431 None,
4432 data.clone(),
4433 Duration::from_millis(20),
4434 ))
4435 .expect_err("the first acknowledgement is retained forever");
4436 assert!(matches!(
4437 error,
4438 CriticalAppendError::DurabilityUnknown { .. }
4439 ));
4440 let original = serde_json::to_string(&log.events()[0]).unwrap();
4441 assert!(log.critical_pending.contains(&original));
4442
4443 let retried = block_on_test_future(log.append_critical_async(
4444 EventKind::RunCompleted,
4445 None,
4446 None,
4447 data,
4448 Duration::from_millis(500),
4449 ))
4450 .expect("the exact row must reconcile without the first acknowledgement");
4451 assert_eq!(serde_json::to_string(retried).unwrap(), original);
4452 assert!(log.critical_pending.is_empty());
4453 assert_eq!(
4454 failures.held_acknowledgement_count(),
4455 1,
4456 "the original acknowledgement must remain unsent"
4457 );
4458 drop(log);
4459
4460 assert_eq!(fs::read_to_string(&path).unwrap(), format!("{original}\n"));
4461 assert_eq!(EventLog::load(&path).unwrap().events().len(), 1);
4462 }
4463
4464 #[test]
4465 fn bounded_sync_critical_timeout_is_exactly_retryable() {
4466 let dir = tempfile::tempdir().unwrap();
4467 let path = dir.path().join("critical-bounded-sync-timeout.jsonl");
4468 let failures = JournalFailureInjector::default();
4469 failures.fail_next(JournalFailurePoint::HoldAcknowledgement);
4470 let mut log = EventLog::with_journal_failure_injector(path.clone(), failures.clone());
4471 log.bind_run("run-bounded-sync", "client-bounded-sync")
4472 .unwrap();
4473 let data = HashMap::from([("completion_digest".to_string(), Value::from("7".repeat(64)))]);
4474
4475 let error = log
4476 .append_critical_bounded(
4477 EventKind::RunCompleted,
4478 None,
4479 None,
4480 data.clone(),
4481 Duration::from_millis(20),
4482 )
4483 .expect_err("the retained acknowledgement must hit the exact sync bound");
4484 assert_eq!(
4485 error,
4486 CriticalAppendError::DurabilityUnknown {
4487 reason: "journal writer did not acknowledge within 20ms".to_string()
4488 }
4489 );
4490 assert!(error.is_retry_safe());
4491 let pending = log
4492 .critical_pending
4493 .iter()
4494 .next()
4495 .expect("the exact timed-out row remains pending")
4496 .clone();
4497 assert_eq!(pending, serde_json::to_string(&log.events()[0]).unwrap());
4498
4499 let retried = log
4500 .append_critical_bounded(
4501 EventKind::RunCompleted,
4502 None,
4503 None,
4504 data,
4505 Duration::from_millis(500),
4506 )
4507 .expect("an exact retry must reconcile without the held acknowledgement");
4508 assert_eq!(serde_json::to_string(retried).unwrap(), pending);
4509 assert!(log.critical_pending.is_empty());
4510 assert_eq!(failures.held_acknowledgement_count(), 1);
4511 drop(log);
4512 assert_eq!(fs::read_to_string(&path).unwrap(), format!("{pending}\n"));
4513 }
4514
4515 #[test]
4516 fn async_critical_append_bounds_unknown_ack_and_preserves_exact_retry() {
4517 let dir = tempfile::tempdir().unwrap();
4518 let path = dir.path().join("critical-unknown-ack.jsonl");
4519 let failures = JournalFailureInjector::default();
4520 failures.fail_next(JournalFailurePoint::HoldAcknowledgement);
4521 let mut log = EventLog::with_journal_failure_injector(path.clone(), failures.clone());
4522 log.bind_run("run-unknown-ack", "client-unknown-ack")
4523 .unwrap();
4524 let data = HashMap::from([("completion_digest".to_string(), Value::from("e".repeat(64)))]);
4525
4526 let started = std::time::Instant::now();
4527 let error = block_on_test_future(log.append_critical_async(
4528 EventKind::RunCompleted,
4529 None,
4530 None,
4531 data.clone(),
4532 std::time::Duration::from_millis(40),
4533 ))
4534 .expect_err("a retained acknowledgement must become durability-unknown");
4535 assert!(
4536 started.elapsed() < std::time::Duration::from_millis(500),
4537 "the async acknowledgement wait exceeded its bounded allowance"
4538 );
4539 assert!(matches!(
4540 error,
4541 CriticalAppendError::DurabilityUnknown { .. }
4542 ));
4543 assert!(error.is_retry_safe());
4544 assert!(error
4545 .to_string()
4546 .contains("did not acknowledge within 40ms"));
4547
4548 let pending = log
4549 .critical_pending
4550 .iter()
4551 .next()
4552 .expect("the exact unacknowledged row remains pending")
4553 .clone();
4554 assert_eq!(log.critical_pending.len(), 1);
4555 assert_eq!(pending, serde_json::to_string(&log.events()[0]).unwrap());
4556 let disk_row = format!("{pending}\n");
4557 let writer_deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
4558 while fs::read_to_string(&path).unwrap_or_default() != disk_row {
4559 assert!(
4560 std::time::Instant::now() < writer_deadline,
4561 "the writer never durably accepted the held-acknowledgement row"
4562 );
4563 std::thread::yield_now();
4564 }
4565 while failures.held_acknowledgement_count() == 0 {
4566 assert!(
4567 std::time::Instant::now() < writer_deadline,
4568 "the writer never retained the late acknowledgement"
4569 );
4570 std::thread::yield_now();
4571 }
4572 assert_eq!(failures.held_acknowledgement_count(), 1);
4573 failures.release_held_acknowledgements();
4574 let original_timestamp = log.events()[0].timestamp;
4575
4576 let distinct_error = block_on_test_future(log.append_critical_async(
4577 EventKind::ProposalCompleted,
4578 None,
4579 Some("different-terminal"),
4580 HashMap::new(),
4581 Duration::from_millis(100),
4582 ))
4583 .expect_err("a different critical row must not bypass exact retry");
4584 assert!(matches!(
4585 distinct_error,
4586 CriticalAppendError::Rejected { .. }
4587 ));
4588 assert_eq!(log.events().len(), 1);
4589
4590 let retried = block_on_test_future(log.append_critical_async(
4591 EventKind::RunCompleted,
4592 None,
4593 None,
4594 data,
4595 std::time::Duration::from_secs(1),
4596 ))
4597 .expect("an identical retry must finish the retained row");
4598 assert_eq!(retried.timestamp, original_timestamp);
4599 assert!(log.critical_pending.is_empty());
4600
4601 block_on_test_future(log.append_critical_async(
4602 EventKind::ProposalCompleted,
4603 None,
4604 Some("different-terminal"),
4605 HashMap::new(),
4606 Duration::from_millis(500),
4607 ))
4608 .expect("a different row is allowed after exact retry reconciliation");
4609 drop(log);
4610
4611 let loaded = EventLog::load(&path).unwrap();
4612 assert_eq!(loaded.events().len(), 2);
4613 assert_eq!(serde_json::to_string(&loaded.events()[0]).unwrap(), pending);
4614 }
4615
4616 #[test]
4617 fn critical_append_failures_retry_same_row_and_fsync_prior_async_events() {
4618 for point in [
4619 JournalFailurePoint::Write,
4620 JournalFailurePoint::Flush,
4621 JournalFailurePoint::Fsync,
4622 ] {
4623 let dir = tempfile::tempdir().unwrap();
4624 let path = dir.path().join(format!("critical-{point:?}.jsonl"));
4625 let failures = JournalFailureInjector::default();
4626 failures.fail_next(point);
4627 let mut log = EventLog::with_journal_failure_injector(path.clone(), failures);
4628 log.bind_run("run-critical", "client-critical").unwrap();
4629 log.append(
4630 EventKind::ProposalReceived,
4631 None,
4632 Some("proposal-critical"),
4633 HashMap::new(),
4634 );
4635 let data =
4636 HashMap::from([("completion_digest".to_string(), Value::from("a".repeat(64)))]);
4637 assert!(
4638 log.append_critical(EventKind::RunCompleted, None, None, data.clone())
4639 .is_err(),
4640 "{point:?} failure must not acknowledge"
4641 );
4642 log.append(
4643 EventKind::ActionSucceeded,
4644 Some("after-pending-terminal"),
4645 Some("proposal-critical"),
4646 HashMap::new(),
4647 );
4648 log.append_critical(EventKind::RunCompleted, None, None, data)
4649 .expect("retry finishes the exact critical row");
4650 drop(log);
4651
4652 let loaded = EventLog::load(&path).unwrap();
4653 let events = loaded.events();
4654 assert_eq!(events[0].kind, EventKind::ProposalReceived);
4655 assert_eq!(events[1].kind, EventKind::RunCompleted);
4656 assert_eq!(events[2].kind, EventKind::ActionSucceeded);
4657 assert_eq!(
4658 events
4659 .iter()
4660 .filter(|event| event.kind == EventKind::RunCompleted)
4661 .count(),
4662 1,
4663 "{point:?} retry must not duplicate a terminal"
4664 );
4665 }
4666 }
4667
4668 #[test]
4669 fn compacting_a_failed_critical_row_makes_its_retry_idempotent() {
4670 let dir = tempfile::tempdir().unwrap();
4671 let path = dir.path().join("critical-compact-retry.jsonl");
4672 let failures = JournalFailureInjector::default();
4673 failures.fail_next(JournalFailurePoint::Fsync);
4674 let mut log =
4675 EventLog::with_journal_failure_injector(path.clone(), failures).with_hash_chaining();
4676 log.bind_run("run-critical-compact", "client-critical-compact")
4677 .unwrap();
4678 log.append(
4679 EventKind::ProposalReceived,
4680 None,
4681 Some("proposal-critical-compact"),
4682 HashMap::new(),
4683 );
4684 let data = HashMap::from([("completion_digest".to_string(), Value::from("c".repeat(64)))]);
4685
4686 assert!(
4687 log.append_critical(EventKind::RunCompleted, None, None, data.clone())
4688 .is_err(),
4689 "the injected fsync failure must leave the exact terminal pending"
4690 );
4691 assert!(
4692 log.compact_journal(),
4693 "compaction persists the in-memory pending terminal"
4694 );
4695 log.append_critical(EventKind::RunCompleted, None, None, data)
4696 .expect("identical retry recognizes the compacted terminal as durable");
4697 drop(log);
4698
4699 let loaded = EventLog::load(&path).unwrap();
4700 assert_eq!(
4701 loaded
4702 .events()
4703 .iter()
4704 .filter(|event| event.kind == EventKind::RunCompleted)
4705 .count(),
4706 1,
4707 "writer respawn must not duplicate the compacted terminal"
4708 );
4709 assert_eq!(loaded.events()[0].kind, EventKind::ProposalReceived);
4710 assert_eq!(loaded.events()[1].kind, EventKind::RunCompleted);
4711 assert_eq!(loaded.verify_chain(), Ok(2));
4712 }
4713
4714 #[test]
4715 fn critical_append_cannot_ack_until_failed_prior_async_row_is_replayed() {
4716 let dir = tempfile::tempdir().unwrap();
4717 let path = dir.path().join("prior-async-failure.jsonl");
4718 let failures = JournalFailureInjector::default();
4719 failures.fail_next(JournalFailurePoint::AsyncWrite);
4722 failures.fail_next(JournalFailurePoint::AsyncWrite);
4723 let mut log = EventLog::with_journal_failure_injector(path.clone(), failures);
4724 log.bind_run("run-ordered", "client-ordered").unwrap();
4725 log.append(
4726 EventKind::ActionSucceeded,
4727 Some("action-ordered"),
4728 Some("proposal-ordered"),
4729 HashMap::new(),
4730 );
4731 let data = HashMap::from([("completion_digest".to_string(), Value::from("b".repeat(64)))]);
4732
4733 assert!(
4734 log.append_critical(
4735 EventKind::ProposalCompleted,
4736 None,
4737 Some("proposal-ordered"),
4738 data.clone()
4739 )
4740 .is_err(),
4741 "a terminal must not acknowledge while an earlier row is still missing"
4742 );
4743 log.append_critical(
4744 EventKind::ProposalCompleted,
4745 None,
4746 Some("proposal-ordered"),
4747 data,
4748 )
4749 .expect("retry repairs the prior row before acknowledging the terminal");
4750 drop(log);
4751
4752 let loaded = EventLog::load(&path).unwrap();
4753 assert_eq!(loaded.events().len(), 2);
4754 assert_eq!(loaded.events()[0].kind, EventKind::ActionSucceeded);
4755 assert_eq!(loaded.events()[1].kind, EventKind::ProposalCompleted);
4756 }
4757
4758 #[test]
4759 fn first_use_parent_sync_failure_blocks_terminal_until_prior_async_replays() {
4760 let dir = tempfile::tempdir().unwrap();
4761 let path = dir.path().join("nested").join("first-use.jsonl");
4762 let failures = car_secrets::PrivatePathDurabilityFailureInjector::default();
4763 failures.fail_next(car_secrets::PrivatePathDurabilityFailurePoint::ParentDirectorySync);
4764 failures.fail_next(car_secrets::PrivatePathDurabilityFailurePoint::ParentDirectorySync);
4765 let mut log = EventLog::with_private_path_failure_injector(path.clone(), failures);
4766 log.bind_run("run-first-use", "client-first-use").unwrap();
4767 log.append(
4768 EventKind::ActionSucceeded,
4769 Some("action-first-use"),
4770 Some("proposal-first-use"),
4771 HashMap::new(),
4772 );
4773 let data = HashMap::from([("completion_digest".to_string(), Value::from("d".repeat(64)))]);
4774
4775 assert!(log
4776 .append_critical(
4777 EventKind::ProposalCompleted,
4778 None,
4779 Some("proposal-first-use"),
4780 data.clone(),
4781 )
4782 .is_err());
4783 log.append_critical(
4784 EventKind::ProposalCompleted,
4785 None,
4786 Some("proposal-first-use"),
4787 data,
4788 )
4789 .expect("retry durably replays async row then exact terminal");
4790 drop(log);
4791
4792 let loaded = EventLog::load(&path).unwrap();
4793 assert_eq!(loaded.events().len(), 2);
4794 assert_eq!(loaded.events()[0].kind, EventKind::ActionSucceeded);
4795 assert_eq!(loaded.events()[1].kind, EventKind::ProposalCompleted);
4796 }
4797}