Skip to main content

car_eventlog/
lib.rs

1//! Event log with JSONL persistence for Common Agent Runtime.
2//!
3//! Append-only event log. Every runtime operation is recorded here.
4//! Supports optional JSONL journal persistence for replay and audit.
5
6pub 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/// Event kinds matching the Python EventKind enum.
44#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize, schemars::JsonSchema)]
45#[serde(rename_all = "snake_case")]
46pub enum EventKind {
47    /// The authenticated `runs.start` bracket reached its durable boundary.
48    RunStarted,
49    /// A body-free authenticated run cancellation request became durable.
50    RunCancellationRequested,
51    /// A deterministic run cancellation receipt became durable.
52    RunCancellationResult,
53    ProposalReceived,
54    /// A proposal reached its deterministic terminal result. Emitted by the
55    /// CAR server after the runtime has emitted every action transition.
56    ProposalCompleted,
57    ActionValidated,
58    ActionRejected,
59    ActionExecuting,
60    ActionSucceeded,
61    ActionFailed,
62    ActionSkipped,
63    ActionRetrying,
64    ActionDeduplicated,
65    PolicyViolation,
66    StateChanged,
67    StateSnapshot,
68    /// Proposal-level aggregate of observed state mutations that survived the
69    /// transaction boundary. Distinct from provisional per-action
70    /// `state_changed` rows and from declared expected effects.
71    StateCommitted,
72    StateRollback,
73    // Skill lifecycle events (SkillRL-inspired)
74    SkillDistilled,
75    SkillEvolved,
76    SkillDeprecated,
77    EvolutionTriggered,
78    /// A provisional skill candidate passed the validation gate and was promoted
79    /// to Active, superseding its incumbent (SkillOpt-inspired — see
80    /// `docs/solutions/gated-skill-optimization.md`).
81    CandidatePromoted,
82    /// A provisional skill candidate failed the validation gate and was rejected
83    /// (recorded in the rejected-edit buffer so it isn't regenerated).
84    CandidateRejected,
85    // Memory consolidation ("dream") events
86    Consolidated,
87    // Proactive memory intervention events (arXiv 2607.08716-inspired):
88    // Phase 1 maintenance records compact bank edits derived from recent
89    // trajectory telemetry; Phase 2 records whether the selector injected a
90    // grounded reminder or explicitly remained silent.
91    ProactiveMemoryMaintained,
92    ProactiveMemoryIntervention,
93    // Replanning events
94    ReplanAttempted,
95    ReplanProposalReceived,
96    ReplanRejected,
97    ReplanExhausted,
98    // Voice turn telemetry — emitted by car-engine's voice_turn dispatch
99    // and the orchestrator. `data` carries `turn_id` (u64) plus
100    // event-specific fields like `text_len`, `error`, `timeout_ms`.
101    VoiceFastTurnStarted,
102    VoiceFastTurnEnded,
103    VoiceSidecarResolved,
104    VoiceSidecarFailed,
105    VoiceSidecarTimedOut,
106    VoiceTurnCancelled,
107    VoiceBridgePlayed,
108    // Foreman merge-verify gate (verified-parallel-coding-orchestrator).
109    // Emitted by car-multi's foreman gate when a farmed-out worktree is
110    // verified before integration. `data` carries `subtask`, `changed_symbols`,
111    // `containment_violations`, `semantic_conflicts`, and `build_test`. This is
112    // the audit trail that makes the gate policy-aware rather than a bare merge.
113    GateAccepted,
114    GateRejected,
115    // The inference chain served a call on a model other than the preferred
116    // candidate, mid-run. `data` carries `from`, `to` and `reason` (one of
117    // `credential_rejected` / `credential_absent` / `rate_limited` /
118    // `timed_out` / `failed`).
119    //
120    // car#1333 recorded WHO wrote each turn (`models_served`); this records
121    // WHY the backbone changed, which is a different fact. Without it a run
122    // that degraded mid-session had nothing on disk explaining a surprising
123    // result, and mining would attribute it to the code under test rather than
124    // to the model swap (car#1351).
125    //
126    // One row per HOP of a fallback chain. The live `model_fallback` event is
127    // latched once per phase so the stream does not narrate every routing
128    // decision; the journal wants each distinct transition, and not the same
129    // one re-stated on all fifty turns of a run whose credential stayed dead.
130    // Repeats are therefore collapsed CONSECUTIVELY, not globally — a lane
131    // that fails, recovers, and fails again later is a new episode, and
132    // suppressing it would leave a journal that cannot tell "changed once,
133    // early" from "flapped all run".
134    //
135    // A candidate the ROUTER never offered uses the alternate payload
136    // `{lane, reason, until}`: `rate_limited` after a prior 429, or
137    // `circuit_open` while its breaker blocks dispatch. `until` is the Unix
138    // cooldown end when known and null for a session-long exclusion or an
139    // in-flight half-open probe. Those persistent conditions are deduplicated
140    // once per lane per coder session rather than consecutively.
141    ModelFallback,
142    // Per-execution caller / tenant scope (Parslee-ai/car#187 phase 3).
143    // Emitted by Runtime::execute_scoped* once per proposal when the
144    // RuntimeScope carries any identity. `data` carries `caller_id`,
145    // `tenant_id`, and `claims` — exact set depends on what the
146    // dispatcher forwarded. Audit / log analysis correlates actions
147    // back to the caller / tenant that triggered them.
148    SessionScope,
149    // Permission-tier gate decisions (survey "Code as Agent Harness"
150    // §3.4.3, §5.2.5 — the harness as safety governor). Emitted by
151    // car-engine's TierPermissionHandler when the permission gate
152    // evaluates an action. `data` carries `gate_decision` (allow /
153    // needs_approval / deny), `required_tier`, `granted_tier`, and (for
154    // escalation/deny) `fingerprint` + `reason`. The audit trail that
155    // makes permission tiers inspectable rather than implicit.
156    //
157    // `data` also carries `reversibility` (reversible / compensable /
158    // irreversible) on every variant — the SECOND, independent axis, from
159    // car_policy::classify_reversibility. The tier answers "who may
160    // authorize this?" and says nothing about whether the effect can be
161    // undone: a `git push` and a charged card are both full_access /
162    // needs_approval and have different rollback contracts. The gate does
163    // not act on this field; it is recorded so an audit can tell those two
164    // rows apart without re-deriving the classification later.
165    PermissionDecision,
166    // A durable human-in-the-loop approval/rejection was recorded
167    // (§5.2.5 — "approvals should be auditable state transitions").
168    // `data` carries `fingerprint`, `approval` (approved / rejected),
169    // `required_tier`, `reviewer`, `reason`, and optional `evidence`.
170    // The auditable counterpart to the ApprovalLedger's durable record.
171    ApprovalRecorded,
172    // Deep-telemetry breadcrumbs (survey §3.5.1 — deep telemetry as the
173    // optimization substrate; "decision-tree traces show where the agent
174    // repeatedly chooses unproductive paths"). A BranchDecision records a
175    // fork the harness took and why; `data` carries `branch` (the chosen
176    // path), `reason`, and any decision-specific context. The substrate an
177    // Evolution Agent (§3.5.2) replays to find where the loop wastes work.
178    BranchDecision,
179    // An alternative the harness considered and discarded — a failed
180    // attempt superseded by a retry/replan, a candidate not selected.
181    // `data` carries `alternative` (what was rejected) and `reason`.
182    // Without this, telemetry shows only the path taken, not the paths
183    // pruned, which is exactly what failure-mode diagnosis needs.
184    AlternativeRejected,
185    // An inference call's token/cost telemetry (§3.5.1). Carries the
186    // standardized metric keys (`tokens_in`, `tokens_out`, `cost_usd`) via
187    // `append_metered`. A dedicated kind so model cost feeds
188    // `metrics_totals` without inflating action-success counts.
189    InferenceMetered,
190    // A transactional conflict the harness detected before executing a
191    // proposal against the versioned shared state (survey §4.3/§5.2.4).
192    // Emitted by the executor's pre-execution transaction check. `data`
193    // carries `kind` (write_write / read_write / stale_assumption), `key`,
194    // `actions`, `explanation`, and `resolution`. Under strict mode the
195    // proposal is rejected; under warn mode it is only recorded.
196    TransactionConflict,
197    // A proposal-admission gate decision (EPIC A / task A1 — the
198    // executor's pre-execution safety seam). Emitted once per registered
199    // `AdmissionGate` that runs during proposal admission. `data` carries
200    // `gate` (the gate name, e.g. information_flow / concurrency / policy),
201    // `decision` (allow / reject / needs_approval), and — when the gate
202    // objects — `reason`, `blocked` (the offending action ids), and an
203    // optional `fingerprint` for approval escalations. The audit trail
204    // that makes the verified safety checks inspectable as live
205    // enforcement rather than dormant library functions.
206    AdmissionGateDecision,
207    // A tool-use hallucination caught by cross-checking the model's claims
208    // against the runtime's own execution receipts (EPIC A / A6 — arXiv
209    // 2603.10060). `data` carries `count` and `hallucinations` (each with
210    // kind/tool/explanation). Deterministic and zero-inference: the runtime
211    // ran the tools, so it holds unforgeable ground truth.
212    ToolReceiptHallucination,
213    // Deterministic goal-loop verifier pass. Emitted after the runtime gathers
214    // ground truth and `car-verify` evaluates a goal condition. `data` carries
215    // `iteration`, `met`, `grounded`, `reason`, and `model_id`/`model_tier`
216    // (local|cloud|unknown) so `/goal` outcomes — and which model tier produced
217    // an ungrounded completion — are queryable from the same append-only
218    // journal as the tool receipts they depend on.
219    GoalEvaluated,
220    // Terminal decision of an assistant/coder turn loop. `data` carries
221    // `decision` ("empty_tool_calls" | "max_turns" | "stalled"), `stop_reason`,
222    // `was_truncated`, and `turns`. The default (goal-less) loop declares success
223    // the instant the model emits no tool calls, with no truncation/outcome
224    // check — so this makes "why did the loop stop" (a clean finish vs a
225    // truncated, turn-capped, or stalled one) queryable from the journal, on the
226    // ungrounded default path that emits no `GoalEvaluated`.
227    TurnCompleted,
228    /// The active `runs.start` bracket reached one terminal state. This is
229    /// distinct from `ProposalCompleted`: one run can contain many proposals.
230    RunCompleted,
231}
232
233/// Status of a trace span.
234#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
235#[serde(rename_all = "snake_case")]
236pub enum SpanStatus {
237    Ok,
238    Error,
239    Unset,
240}
241
242/// A trace span representing a unit of work.
243#[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
255/// Standardized `Event.data` keys for cross-cutting telemetry metrics, so
256/// every emit site records them under the same name and aggregation can
257/// rely on it (survey §3.5.1: deep telemetry "records the decision process
258/// in greater detail: token usage and cost, model/tool latency …").
259pub mod metric_keys {
260    /// Wall-clock duration of the unit of work, milliseconds (f64).
261    pub const DURATION_MS: &str = "duration_ms";
262    /// Input/prompt tokens consumed (u64).
263    pub const TOKENS_IN: &str = "tokens_in";
264    /// Output/completion tokens produced (u64).
265    pub const TOKENS_OUT: &str = "tokens_out";
266    /// Estimated cost in USD (f64).
267    pub const COST_USD: &str = "cost_usd";
268}
269
270/// Cross-cutting telemetry metrics attachable to any event. All optional —
271/// a tool call has latency but no tokens; an inference has all four. Merged
272/// into `Event.data` under [`metric_keys`] by [`EventLog::append_metered`],
273/// and read back via the `Event` accessors, so downstream aggregation
274/// (harness-level metrics, the Evolution Agent) has a uniform source.
275#[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    /// Latency-only metrics (the common tool/action case).
289    pub fn latency(duration_ms: f64) -> Self {
290        Self {
291            duration_ms: Some(duration_ms),
292            ..Default::default()
293        }
294    }
295
296    /// Token + cost metrics for an inference call.
297    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    /// Merge these metrics into an event `data` map under [`metric_keys`].
312    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/// A single event in the log.
329#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, schemars::JsonSchema)]
330pub struct Event {
331    pub kind: EventKind,
332    /// Authenticated active-run identity stamped by CAR at append time.
333    /// Historical journals omit this field and replay as `None`.
334    #[serde(default, skip_serializing_if = "Option::is_none")]
335    pub run_id: Option<String>,
336    /// WebSocket client identity that opened `run_id` via `runs.start`.
337    /// Never reconstructed from the journal filename during replay.
338    #[serde(default, skip_serializing_if = "Option::is_none")]
339    pub client_id: Option<String>,
340    /// CAR-minted policy-session identity used for this proposal, when the
341    /// caller selected a live `session.policy.open` session. Unvalidated
342    /// caller labels are never copied here.
343    #[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    /// Hash of the previous event in the chain (EPIC A / A9 tamper-
354    /// evidence). `None` when hash chaining is disabled (the default) —
355    /// the field is skipped in serialization, so logs without chaining are
356    /// byte-identical to before this was added.
357    #[serde(default, skip_serializing_if = "Option::is_none")]
358    pub prev_hash: Option<String>,
359    /// This event's own content hash, computed over its fields plus
360    /// `prev_hash`. Present only when hash chaining is enabled.
361    #[serde(default, skip_serializing_if = "Option::is_none")]
362    pub hash: Option<String>,
363}
364
365impl Event {
366    /// Wall-clock duration recorded on this event, if any.
367    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    /// Input tokens recorded on this event, if any.
374    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    /// Output tokens recorded on this event, if any.
381    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    /// Estimated cost (USD) recorded on this event, if any.
388    pub fn cost_usd(&self) -> Option<f64> {
389        self.data.get(metric_keys::COST_USD).and_then(Value::as_f64)
390    }
391
392    /// All metrics carried on this event, gathered into a [`Metrics`].
393    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/// Summed telemetry metrics across a set of events — the trajectory-level
404/// totals harness-level evaluation (§5.2.1) and the Evolution Agent
405/// (§3.5.2) reason over. `tokens` is the sum of in + out.
406#[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    /// Number of events that carried at least one metric.
414    pub metered_events: usize,
415}
416
417/// Sum the telemetry metrics across a slice of events. The single
418/// implementation behind both [`EventLog::metrics_totals`] and the
419/// harness-metrics computation, so the two can never drift on the metric
420/// contract (neo review: avoid a duplicated copy).
421pub 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/// Per-agent cost/token attribution (EPIC G / G3). Folded from
451/// `InferenceMetered` events that carry an `agent` field in `data`.
452#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
453pub struct AgentCost {
454    pub agent: String,
455    /// Number of metered inference events attributed to this agent.
456    pub calls: u64,
457    pub tokens_in: u64,
458    pub tokens_out: u64,
459    pub cost_usd: f64,
460}
461
462/// Attribute token/cost totals per agent by folding `InferenceMetered` events
463/// grouped by their `data["agent"]` field (EPIC G / G3). Events with no `agent`
464/// field are grouped under `"unknown"`. Ordered by agent name (BTreeMap) so the
465/// report is deterministic. This is how a multi-agent run reports cost per agent
466/// (e.g. Researcher $2, Coordinator $0.5) and, joined with the `tools`/`workflow`
467/// provenance fields the emit sites stamp, how a tool call is traceable to its
468/// agent.
469pub 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
494/// Background JSONL journal writer. `EventLog::append` hands a serialized event
495/// line to this over a channel; a dedicated thread owns the file and does the
496/// actual write. So `append` never does file I/O while a caller holds the log
497/// mutex — the head-of-line blocking that bites when many concurrent tasks
498/// (e.g. Foreman gate verifications running under one shared, journaled session
499/// log) each re-opened and wrote the file under the lock.
500///
501/// Best-effort, like the journal it replaces: an open/write failure drops the
502/// line (the in-memory event vec is unaffected) — but unlike the old silent
503/// journal, the hard failures (can't spawn the thread, can't open the file) are
504/// surfaced via `tracing::warn!`, since this carries the gate audit trail and a
505/// silently-broken audit log is worse than a noisy one.
506///
507/// The channel is unbounded so a burst never blocks the hot path. This relies on
508/// an envelope: low per-session journal volume and a writer that keeps up, so the
509/// backlog stays small. It is not a *new* unbounded-growth risk — the in-memory
510/// `events` vec already grows without bound under the same pathological
511/// hot-loop-`append` workload, so the channel is not the first thing to OOM.
512enum JournalMessage {
513    Async(String),
514    Critical {
515        line: String,
516        known_existing: bool,
517        ack: JournalAcknowledgement,
518    },
519    #[cfg(test)]
520    Shutdown,
521}
522
523/// Hard upper bound for one asynchronous critical-journal acknowledgement.
524/// Callers may choose a shorter deadline, but never create a timer entry that
525/// lives longer than this process-wide contract.
526pub const MAX_CRITICAL_ACKNOWLEDGEMENT_TIMEOUT: Duration = Duration::from_secs(30);
527
528const DEFAULT_CRITICAL_ACKNOWLEDGEMENT_CAPACITY: usize = 64;
529
530/// Deterministic failure seam for durable-journal tests. Production callers
531/// use the default empty queue; embedders may inject one failure at an exact
532/// write/flush/fsync boundary without replacing the filesystem.
533#[derive(Debug, Clone, Copy, PartialEq, Eq)]
534pub enum JournalFailurePoint {
535    AsyncWrite,
536    Write,
537    Flush,
538    Fsync,
539    /// Accept and durably write a critical row, but retain its acknowledgement.
540    /// This models the ambiguous boundary where a caller cannot know whether
541    /// the writer completed before its acknowledgement deadline.
542    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    /// Number of acknowledgements retained by the explicit
580    /// [`JournalFailurePoint::HoldAcknowledgement`] test seam.
581    #[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    /// Release acknowledgements retained by the explicit stalled-writer test
590    /// seam. Production code never arms that seam.
591    #[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/// Failure from a bounded critical append. A durability-unknown result retains
1075/// the exact serialized row and is safe to retry with the same lifecycle
1076/// identity and data.
1077#[derive(Debug, Clone, PartialEq, Eq)]
1078pub enum CriticalAppendError {
1079    /// The append was rejected before an acknowledgement wait could begin.
1080    Rejected { reason: String },
1081    /// The writer accepted the request, but durable completion was not known
1082    /// before the acknowledgement bound. The exact row remains pending.
1083    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    /// `None` only if the writer thread could not be spawned (journaling then
1110    /// silently disabled — still best-effort).
1111    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            // Drop tx (rx dies with it); journaling becomes a no-op.
1163            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            // Best-effort: if the writer thread has gone, drop the line.
1177            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        // Close the channel so the writer drains its backlog, flushes, and
1252        // exits; join so buffered lines are durable by the time the log is gone.
1253        self.tx.take();
1254        if let Some(handle) = self.handle.take() {
1255            let _ = handle.join();
1256        }
1257    }
1258}
1259
1260/// The journal thread's body: own the file, write each line, flush when the
1261/// channel goes momentarily idle (batches bursts, keeps durability prompt).
1262///
1263/// The file is opened **lazily on the first line to write**, not at thread
1264/// start. A session that never appends an event — health checks, heartbeats,
1265/// `agents.list` polls, and every other no-op connection — then leaves no
1266/// journal behind. Opening eagerly created a 0-byte `<client_id>.jsonl` per
1267/// connection that accumulated without bound (177K empties observed on a
1268/// long-lived daemon). Sessions that DO log are unaffected: the file is created
1269/// on their first event exactly as before.
1270fn 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    // Rows whose first asynchronous write failed remain ahead of the next
1280    // critical boundary. Rows received after a failed critical boundary stay
1281    // behind that exact row. Keeping the queues separate preserves the
1282    // producer order across retries.
1283    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                // A critical acknowledgement is an ordered durability
1347                // barrier. Repair every earlier asynchronous row first; if
1348                // any row is still unwritable, reserve this exact critical
1349                // line and fail closed instead of creating an audit gap.
1350                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)); // ENOSPC
1453        }
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            // `write_all` may have written a prefix before returning an
1468            // error. Restore the exact pre-row boundary so retry cannot leave
1469            // a truncated JSON fragment or duplicate suffix in the journal.
1470            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    // Windows has no portable directory-fsync primitive: FlushFileBuffers on
1505    // a directory handle returns ERROR_ACCESS_DENIED on supported NTFS
1506    // runners. open_private_append already completes the platform-specific
1507    // parent metadata boundary before returning the newly-created file.
1508    Ok(())
1509}
1510
1511/// Append-only event log with optional JSONL journal.
1512pub struct EventLog {
1513    events: Vec<Event>,
1514    spans: Vec<Span>,
1515    journal: Option<JournalWriter>,
1516    /// When true, each appended event is hash-chained to its predecessor
1517    /// (EPIC A / A9). Off by default — enabling it is opt-in so existing
1518    /// JSONL output stays byte-identical for consumers that don't need
1519    /// tamper-evidence.
1520    hash_chaining: bool,
1521    /// The hash of the most recently appended event, threaded into the
1522    /// next event's `prev_hash`. The genesis link uses the empty string.
1523    last_hash: Option<String>,
1524    /// Auto-retention policy (EPIC G / G2). When set, `append` caps the
1525    /// in-memory event count at `max_events` (dropping oldest) so the log
1526    /// can't grow unbounded. Age-based trimming is applied by
1527    /// `enforce_retention`. `None` = keep everything (unchanged default).
1528    retention: Option<RetentionPolicy>,
1529    /// Path of the JSONL journal, kept so retention trims can compact
1530    /// (rewrite) the file — the background [`JournalWriter`] only appends.
1531    journal_path: Option<PathBuf>,
1532    /// Approximate number of event lines currently in the journal file:
1533    /// incremented per journaled append, seeded from the parsed event count
1534    /// on [`EventLog::load`], reset to the retained count after a
1535    /// compaction. Drives the compaction throttle.
1536    journal_lines: usize,
1537    /// Total events ever dropped from the in-memory log (retention trims,
1538    /// manual truncation, `clear`). Monotonic. Lets consumers that project
1539    /// over `events()` — e.g. the tool-receipt verifier (A6) — know the
1540    /// retained window is incomplete instead of mistaking an evicted event
1541    /// for one that never happened.
1542    trimmed_events: u64,
1543    /// Monotonic cumulative cost (USD) across every event ever appended
1544    /// (EPIC G / G1). Updated at append time and **never** decremented by
1545    /// retention trims, truncation, or `clear`, so a cumulative budget check
1546    /// can't slide backward when old events are evicted. Seeded from the
1547    /// journal on [`EventLog::load`].
1548    cumulative_cost_usd: f64,
1549    /// Live producer binding for CAR-owned journal appends. This is runtime
1550    /// state, not replay state: loading historical rows never fabricates an
1551    /// active run from whatever happens to be at the journal tail.
1552    active_binding: Option<EventBinding>,
1553    /// Exact serialized critical events whose first durability attempt failed.
1554    /// A retry reuses the same timestamp/bytes and asks the writer to finish
1555    /// flush+fsync instead of minting a conflicting duplicate terminal.
1556    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
1573/// Journal-compaction throttle floor (G2): a retention trim only triggers a
1574/// journal rewrite once the journal holds at least this many more lines than
1575/// the retained set (and the rewrite would shrink it by ≥25% — see
1576/// [`EventLog::maybe_compact_journal`]). Keeps frequent small trims from
1577/// rewriting the file on every append.
1578const JOURNAL_COMPACT_MIN_EXCESS: usize = 1024;
1579
1580/// Compute the content hash of an event for the tamper-evidence chain.
1581///
1582/// Hashes `prev_hash` plus a canonical rendering of the event's content
1583/// (kind, action/proposal ids, sorted `data`, timestamp). The top-level
1584/// `data` map is sorted by key so the digest is stable across a
1585/// serialize/deserialize round-trip (serde_json already emits nested object
1586/// keys in sorted order). Any after-the-fact edit to a chained event — or an
1587/// interior deletion/reordering — breaks the chain from that point on. The
1588/// chain has no anchored head hash, so truncation at either end (dropping a
1589/// prefix or a suffix of the log wholesale) is NOT detectable; see
1590/// [`EventLog::verify_chain`] for the precise guarantee.
1591fn 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    // Preserve historical hashes byte-for-byte when every binding field is
1616    // absent. Bound v0.51 events add one domain-separated identity segment.
1617    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/// Retention policy for an [`EventLog`] (EPIC G / G2). Bounds the log by
1638/// **size** (`max_events`, enforced automatically on append — oldest dropped)
1639/// and by **age** (`max_age_secs`, applied by [`EventLog::enforce_retention`]).
1640/// Both `None` = keep everything.
1641#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
1642pub struct RetentionPolicy {
1643    /// Cap the in-memory event count; on overflow the oldest are dropped.
1644    #[serde(default)]
1645    pub max_events: Option<usize>,
1646    /// Drop events older than this many seconds when `enforce_retention` runs.
1647    #[serde(default)]
1648    pub max_age_secs: Option<i64>,
1649}
1650
1651/// A structured audit query over the event log (EPIC G / G2). Every field is
1652/// an AND-conjoined filter; empty/`None` fields don't constrain. Answers
1653/// "who ran what tool when, and which approvals applied" by filtering the
1654/// `SessionScope` / `PermissionDecision` / `ApprovalRecorded` / action trail.
1655#[derive(Debug, Clone, Default, Serialize, Deserialize)]
1656pub struct EventQuery {
1657    /// Restrict to these event kinds (empty = any kind).
1658    #[serde(default)]
1659    pub kinds: Vec<EventKind>,
1660    /// Exact match on `action_id`.
1661    #[serde(default)]
1662    pub action_id: Option<String>,
1663    /// Exact match on `proposal_id`.
1664    #[serde(default)]
1665    pub proposal_id: Option<String>,
1666    /// Inclusive lower time bound.
1667    #[serde(default)]
1668    pub since: Option<DateTime<Utc>>,
1669    /// Exclusive upper time bound.
1670    #[serde(default)]
1671    pub until: Option<DateTime<Utc>>,
1672    /// Match events whose `data` contains ALL these key→value pairs (compared
1673    /// as strings). Covers caller/tenant/tool/gate/decision, which live in
1674    /// `data` on the audit events.
1675    #[serde(default)]
1676    pub data_matches: std::collections::HashMap<String, String>,
1677    /// Cap the number of results (most-recent-first). `None`/0 = unlimited.
1678    #[serde(default)]
1679    pub limit: Option<usize>,
1680}
1681
1682/// Does a JSON `data` value equal the query string? Compares strings directly
1683/// and stringifies scalars so `{"count": 3}` matches `"3"`.
1684fn 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    /// Does `e` satisfy every constraint in this query?
1694    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    /// Test/embedder seam for deterministic journal write/flush/fsync faults.
1764    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    /// Test/embedder seam for deterministic first-use directory-entry faults.
1787    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    /// Bind this log to one authenticated active run. Exact repeat binding is
1799    /// idempotent; a different run/client is rejected instead of silently
1800    /// re-attributing later action events.
1801    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    /// Attach a CAR-minted policy session to the currently bound proposal.
1823    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    /// Enable tamper-evident hash chaining for events appended from now on
1882    /// (EPIC A / A9). The chain continues from the last already-appended
1883    /// event's hash if one exists (re-enabling after a load), else from the
1884    /// genesis link. Returns `self` for builder-style use.
1885    pub fn with_hash_chaining(mut self) -> Self {
1886        self.enable_hash_chaining();
1887        self
1888    }
1889
1890    /// Turn on hash chaining in place. Idempotent.
1891    pub fn enable_hash_chaining(&mut self) {
1892        self.hash_chaining = true;
1893        // Continue the chain from whatever the last event already carries.
1894        if self.last_hash.is_none() {
1895            self.last_hash = self.events.last().and_then(|e| e.hash.clone());
1896        }
1897    }
1898
1899    /// Whether hash chaining is currently enabled.
1900    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                &timestamp,
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        // Hand the serialized line to the background writer — no file I/O here,
1948        // so a caller holding the log mutex is never blocked on disk.
1949        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        // Monotonic cumulative cost (G1): fold cost in at append time so a
1957        // budget check survives retention trims of the underlying events.
1958        if let Some(c) = event.cost_usd() {
1959            self.cumulative_cost_usd += c;
1960        }
1961
1962        self.events.push(event);
1963        // Auto-retention (EPIC G / G2): cap the in-memory log at max_events so
1964        // it can't grow unbounded. Cheap — a bounded pop from the front only
1965        // when over the cap. Age-based trimming is on-demand via
1966        // enforce_retention (walking every event on each append would be O(n)).
1967        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                // The journal keeps the dropped events until the (throttled)
1972                // compaction rewrites it to the retained set.
1973                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                &timestamp,
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    /// Append a lifecycle-critical event and return only after its exact JSONL
2072    /// row and all earlier queued rows have been flushed and fsynced. If a
2073    /// write/flush/fsync attempt fails, the exact event remains pending in
2074    /// memory so an identical retry finishes the same row rather than minting
2075    /// a second terminal with a new timestamp.
2076    ///
2077    /// This compatibility API performs an unbounded blocking acknowledgement
2078    /// wait and is intended only for genuinely synchronous callers. Async
2079    /// callers must use [`Self::append_critical_async`] so a stalled filesystem
2080    /// cannot occupy an executor worker indefinitely.
2081    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    /// Synchronous critical append with a hard acknowledgement bound.
2122    ///
2123    /// This is the startup-thread counterpart of [`Self::append_critical_async`].
2124    /// The exact serialized row is installed in `critical_pending` before the
2125    /// bounded wait begins, so timeout cannot claim success or authorize a
2126    /// different lifecycle event. An identical later startup replay safely
2127    /// reconciles the row whether or not the writer completed before timeout.
2128    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    /// Append a lifecycle-critical event without blocking an async executor
2181    /// worker on filesystem acknowledgement.
2182    ///
2183    /// The acknowledgement wait is bounded by `acknowledgement_timeout`. Once
2184    /// the writer accepts the message, the exact serialized row is marked
2185    /// pending before this future can yield. Timeout, cancellation by an outer
2186    /// request deadline, writer failure, and acknowledgement-coordinator shutdown
2187    /// therefore leave an identical retry safe: it reuses the original event
2188    /// timestamp/hash and the writer suppresses a duplicate row. Until that
2189    /// exact retry reconciles the pending row, a different critical event is
2190    /// rejected before enqueue.
2191    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        // This happens before the first `.await`, so dropping the future at an
2223        // outer Tokio timeout cannot lose the exact retry identity.
2224        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    /// Verify the tamper-evidence hash chain over the currently-loaded
2237    /// events (EPIC A / A9). Walks every event that carries a `hash`,
2238    /// recomputing it from its content + the running `prev_hash` and
2239    /// checking the links join up. Returns `Ok(n)` with the number of
2240    /// chained events verified, or `Err(index)` naming the first event
2241    /// whose hash or linkage doesn't match — i.e. the point at which a
2242    /// chained event was edited, or an interior event was deleted or
2243    /// reordered.
2244    ///
2245    /// **Scope of the guarantee:** the chain detects *interior*
2246    /// edits/reorderings/deletions only. It cannot detect truncation at
2247    /// either end: there is no anchored head hash, so the first chained
2248    /// event's `prev_hash` is taken on trust (dropping a prefix goes
2249    /// unnoticed), and nothing pins the tail (dropping a suffix goes
2250    /// unnoticed). Detecting head/tail truncation requires anchoring the
2251    /// chain head (and a trusted latest-hash witness), which is out of
2252    /// scope until that anchor exists.
2253    ///
2254    /// Events without a `hash` (appended before chaining was enabled) are
2255    /// skipped, so a partially-chained log verifies its chained suffix.
2256    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                // Once the chain has started, a gap is a break.
2263                if chain_started {
2264                    return Err(i);
2265                }
2266                continue;
2267            };
2268            // The recorded prev_hash must match the running hash.
2269            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    /// Append an event with cross-cutting [`Metrics`] (duration, tokens,
2295    /// cost) merged into its `data` under [`metric_keys`]. Use this for any
2296    /// event whose latency or token cost should feed trajectory-level
2297    /// aggregation (`metrics_totals`) — the deep-telemetry substrate of
2298    /// §3.5.1. Metric keys present in both `data` and `metrics` take the
2299    /// `metrics` value (the metrics argument wins).
2300    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    /// Sum the telemetry metrics across every event in the log — the
2313    /// trajectory-level totals (tokens, cost, wall-clock) that harness-level
2314    /// evaluation (§5.2.1) and the Evolution Agent (§3.5.2) reason over.
2315    ///
2316    /// Contract: this sums **every** event carrying a [`metric_keys`] value,
2317    /// regardless of which append path emitted it. A duration recorded once
2318    /// per action (e.g. `ActionSucceeded`) is counted once; the standardized
2319    /// keys mean there is a single value per metric per event, so there is no
2320    /// double-count as long as each unit of work meters itself once. Token
2321    /// metrics from `InferenceMetered` and latency from action events sum
2322    /// into the same totals — that is intended (total cost = model + tools).
2323    pub fn metrics_totals(&self) -> MetricsTotals {
2324        metrics_totals_of(&self.events)
2325    }
2326
2327    /// Per-agent cost/token report (EPIC G / G3) — see [`cost_by_agent_of`].
2328    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    /// Drop every retained event and span, releasing their memory. The
2371    /// JSONL journal is left untouched (it is the audit trail); the
2372    /// monotonic counters (`trimmed_events`, `cumulative_cost_usd`) are
2373    /// preserved — `clear` frees memory, it doesn't reset the log's history.
2374    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    /// Total events ever dropped from the in-memory log (retention trims,
2385    /// manual truncation, `clear`). Monotonic; `> 0` means the retained
2386    /// window is incomplete — consumers projecting over [`Self::events`]
2387    /// (e.g. the A6 tool-receipt verifier) must treat an absent event as
2388    /// possibly-evicted, not as never-happened.
2389    pub fn trimmed_events(&self) -> u64 {
2390        self.trimmed_events
2391    }
2392
2393    /// Monotonic cumulative cost (USD) across every event ever appended
2394    /// (EPIC G / G1). Unlike folding `cost_usd` over [`Self::events`] — which
2395    /// slides backward when retention trims metered events — this counter
2396    /// only grows, so it is the correct denominator for a cumulative budget
2397    /// (`AlertThresholds::max_cost_usd`). Seeded from the journal on
2398    /// [`Self::load`]; survives trims and [`Self::clear`].
2399    pub fn cumulative_cost_usd(&self) -> f64 {
2400        self.cumulative_cost_usd
2401    }
2402
2403    /// Current size of the JSONL journal file in bytes, if a journal is
2404    /// configured and stat-able. The background writer batches, so this may
2405    /// momentarily lag the last few appends.
2406    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    /// Journal-compaction throttle (G2): rewrite only when the journal holds
2412    /// at least [`JOURNAL_COMPACT_MIN_EXCESS`] more lines than the retained
2413    /// set AND the rewrite would shrink it by ≥25%. Frequent small trims
2414    /// therefore cost nothing; each compaction rewrites at most the retained
2415    /// set and is amortized O(1) per append.
2416    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    /// Rewrite the JSONL journal to contain exactly the currently-retained
2427    /// events (G2 journal compaction — before this, retention trimmed the
2428    /// in-memory log only and the journal grew unbounded). Atomic: writes a
2429    /// sibling temp file and renames it over the journal. The background
2430    /// writer is joined first (draining its backlog and closing its handle —
2431    /// renaming under a live append-mode handle would orphan subsequent
2432    /// writes to the old inode), then respawned on the compacted file.
2433    ///
2434    /// Hash chaining (A9) survives: [`Self::verify_chain`] anchors the first
2435    /// hashed event on its *stored* `prev_hash`, so the retained tail of a
2436    /// chained log still verifies after a compact + reload. Corollary: a
2437    /// head-trim by retention is indistinguishable from compaction — tamper
2438    /// evidence covers the retained tail only.
2439    ///
2440    /// Returns `true` if the journal was rewritten. Failure is best-effort
2441    /// like the journal itself: a warning is logged, the old (uncompacted)
2442    /// journal stays in place, and appending resumes against it.
2443    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        // Join the writer so pending lines are flushed and its handle closed.
2453        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                // Compaction fsynced these exact rows before replacing the
2477                // journal. Reconcile producer-side retry state before the
2478                // writer respawns so an identical critical retry is treated
2479                // as an already-written durability barrier, not appended a
2480                // second time by the fresh writer's empty identity cache.
2481                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    /// Run a structured audit [`EventQuery`], returning matching events
2499    /// most-recent-first, capped at `query.limit` (EPIC G / G2).
2500    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(); // most recent first for audit review
2503        if let Some(limit) = query.limit.filter(|l| *l > 0) {
2504            out.truncate(limit);
2505        }
2506        out
2507    }
2508
2509    /// Install an auto-retention policy (EPIC G / G2). `max_events` is then
2510    /// enforced on every `append`; call [`Self::enforce_retention`] to also
2511    /// apply the age bound.
2512    pub fn set_retention(&mut self, policy: Option<RetentionPolicy>) {
2513        self.retention = policy;
2514    }
2515
2516    /// The active retention policy, if any.
2517    pub fn retention(&self) -> Option<&RetentionPolicy> {
2518        self.retention.as_ref()
2519    }
2520
2521    /// Apply a retention policy now: drop events older than `max_age_secs`
2522    /// and cap the count at `max_events` (keeping the most recent). Returns
2523    /// the number of events removed. Independent of the installed policy, so a
2524    /// caller can run a one-off sweep. When a journal is configured, a trim
2525    /// also triggers the throttled journal compaction (see
2526    /// [`Self::compact_journal`]) so the JSONL file tracks retention instead
2527    /// of growing unbounded.
2528    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    /// Begin a new trace span. Returns the generated span_id.
2565    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    /// End an open span by setting its status and end time.
2588    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    /// Return all spans.
2596    pub fn spans(&self) -> Vec<Span> {
2597        self.spans.clone()
2598    }
2599
2600    /// Export traces as OTLP-compatible JSON.
2601    pub fn export_traces(&self) -> String {
2602        // Group spans by trace_id
2603        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    /// Load an event log from a JSONL journal file.
2670    pub fn load(path: &Path) -> std::io::Result<Self> {
2671        Self::load_with_writer(path, JournalWriter::spawn(path.to_path_buf()))
2672    }
2673
2674    /// Load and validate an event log without attaching a writer or modifying
2675    /// the source journal. An unterminated final row is reported as
2676    /// [`std::io::ErrorKind::UnexpectedEof`]; callers that own live append
2677    /// recovery should use [`EventLog::load`] instead.
2678    pub fn load_read_only(path: &Path) -> std::io::Result<Self> {
2679        Self::load_from_journal(path, None, false)
2680    }
2681
2682    /// Load a journal while installing the deterministic durability-failure
2683    /// seam for subsequent appends.
2684    #[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                    // A process can crash after writing only a prefix of its
2736                    // final JSONL row. Nothing follows this unterminated tail,
2737                    // so discard it before any resumed append. A terminated
2738                    // bad row is authoritative corruption and fails below.
2739                    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        // If the loaded tail is chained, keep chaining ENABLED and continue
2752        // from the last hash. Restoring `last_hash` but leaving chaining off
2753        // (the pre-fix behaviour) permanently broke the chain: one unchained
2754        // append before a manual re-enable left a gap that made every future
2755        // `verify_chain` report tampering, with no repair path (review C-9b).
2756        let last_hash = events.last().and_then(|e| e.hash.clone());
2757        let hash_chaining = last_hash.is_some();
2758        // Seed the monotonic counters from what the journal preserved: the
2759        // cumulative cost restarts from the journaled spend (G1), and the
2760        // journal line count from the validated events. A recoverable torn
2761        // tail is rewritten to exactly this set before loading completes.
2762        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            // Live loads journal subsequent appends back to the same file;
2768            // read-only loads intentionally carry no writer.
2769            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        // Filter by kind.
2949        let q = EventQuery {
2950            kinds: vec![EventKind::PermissionDecision],
2951            ..Default::default()
2952        };
2953        assert_eq!(log.query(&q).len(), 2);
2954
2955        // Filter by a data field (who ran the tool).
2956        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        // Combined tool + kind.
2965        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        // Most-recent-first + limit → the bob decision.
2972        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        // BTreeMap order: coordinator, researcher.
3017        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        // Only the last 3 survive.
3042        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        // Two events; backdate the first well past the age bound.
3051        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            // Trims dropped 3 metered events; the counter never slid back.
3122            assert_eq!(log.len(), 1);
3123            assert!((log.cumulative_cost_usd() - 10.0).abs() < 1e-9);
3124        }
3125        // Reload seeds the counter from what the journal preserved (here the
3126        // journal was never compacted, so the full spend survives).
3127        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        // Real journal compaction (review G2): once the excess clears the
3134        // throttle (≥ JOURNAL_COMPACT_MIN_EXCESS lines AND ≥25% shrink), the
3135        // retention trim rewrites the JSONL file to exactly the retained
3136        // events instead of letting it grow forever.
3137        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        } // drop joins the writer → file settled.
3158
3159        // The journal holds only the retained tail, not all `total` lines.
3160        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        // The newest events survived, contiguously up to the last append.
3167        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        // A9 × G2: verify_chain anchors the first hashed event on its stored
3176        // prev_hash, so a compacted journal's retained tail must still verify
3177        // after reload even though the chain's head was dropped.
3178        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            // Trim to the last 5 and force an (unthrottled) compaction.
3191            log.truncate_events_keep_last(5);
3192            assert!(log.compact_journal(), "compaction must succeed");
3193            // Appends after compaction land in the compacted file and keep
3194            // chaining from the retained tail.
3195            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        // verify_chain over an unchained log is vacuously ok (0 verified).
3230        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        // Every event hashed, links join.
3245        assert!(log.events().iter().all(|e| e.hash.is_some()));
3246        assert_eq!(log.verify_chain(), Ok(5));
3247        // First event's prev_hash is the genesis (empty) link.
3248        assert_eq!(log.events()[0].prev_hash.as_deref(), Some(""));
3249        // Each subsequent prev_hash equals the prior event's hash.
3250        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        // Tamper with event #2's data after the fact.
3268        log.events[2].data.insert("v".to_string(), Value::from(999));
3269        // The chain breaks exactly at the edited event.
3270        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        // Remove the second event — the next event's prev_hash no longer
3285        // matches the running hash.
3286        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        // Serialize each event to JSON and back, then re-verify — the
3306        // digest must be stable across the round-trip.
3307        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        // Regression (review C-9b): `load` used to restore `last_hash` but
3322        // hard-set `hash_chaining: false`, so the first append after a load
3323        // produced an unchained event mid-chain — a permanent, unrepairable
3324        // verify_chain failure. Loading a chained tail must keep chaining on.
3325        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        } // drop joins the writer thread → lines flushed.
3333
3334        {
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        // And the whole thing still verifies after a second reload.
3345        let reloaded = EventLog::load(&journal).unwrap();
3346        assert_eq!(reloaded.len(), 3);
3347        assert_eq!(reloaded.verify_chain(), Ok(3));
3348
3349        // An UNCHAINED journal must not turn chaining on.
3350        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        // Original data preserved; metrics merged under standardized keys.
3371        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        // An un-metered event must not affect totals.
3396        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        // Contract: metrics_totals sums any event carrying a metric key,
3410        // regardless of append path. A legacy raw `append` that puts
3411        // "duration_ms" in data must still be counted (locks the contract
3412        // documented on metrics_totals).
3413        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        // The new kinds must round-trip as snake_case for the JSON wire.
3428        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        // A session that never logs an event must leave no journal file behind.
3792        // Eager open created a 0-byte file per connection that accumulated
3793        // without bound; the writer now opens lazily on the first line.
3794        let dir = tempfile::tempdir().unwrap();
3795        let journal = dir.path().join("no-events.jsonl");
3796        {
3797            let _log = EventLog::with_journal(journal.clone());
3798            // No append. Drop joins the writer thread, which never opened the
3799            // file because no line was ever sent.
3800        }
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        // The background writer must not lose or reorder events under a tight
3810        // append burst; drop-join guarantees the backlog is flushed before the
3811        // log is gone.
3812        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        } // drop joins the writer thread → all 500 lines flushed.
3825
3826        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        // The whole "best-effort" promise rests on this branch: a journal path
3840        // that can't be opened (here: the path IS an existing directory) must not
3841        // panic or block append — the in-memory log keeps working.
3842        let dir = tempfile::tempdir().unwrap();
3843        let journal = dir.path().join("a-directory");
3844        fs::create_dir(&journal).unwrap(); // open(append) on a dir fails
3845
3846        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        // Drop must still terminate cleanly (writer thread drained and joined).
3860    }
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        // Resume: load, append more, drop → both old and new are on disk.
3876        {
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        // Verify OTLP structure
4009        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        // Verify the child has parentSpanId
4019        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        // The first failure drops the asynchronous attempt. The second makes
4720        // the first critical barrier prove that it cannot repair the gap yet.
4721        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}