Skip to main content

khive_runtime/
audit_batch.rs

1//! ADR-133 Slice 1: the audit-batch seam.
2//!
3//! Incidental audit writes (gate denials, dispatch outcomes, config-lock
4//! rows, git.digest receipts, and pure-observability rows like
5//! `RecallExecuted`) no longer take one writer-task acquisition per row on
6//! the request hot path. Concurrent submissions arriving while a generation
7//! is committing share the *next* generation instead of each taking their
8//! own writer acquisition, so N concurrent producers collapse to one
9//! [`khive_storage::EventStore::append_events_idempotent`] call per
10//! generation rather than N.
11//!
12//! [`AuditBatch`] owns this seam. A lazily-spawned supervisor task drains
13//! pending rows into generations and drives each through the store; the
14//! supervisor's own `JoinHandle` is retained (never discarded) so an
15//! abnormal exit — panic, cancellation, a lost child join, or a driver that
16//! returns `Ok` while state is not terminally consistent — is observed and
17//! converted into a `Failed` transition with all accepted waiters resolved,
18//! per owner ruling R1/R4 (`.khive/OWNER_RULING_adr133_gate.md`).
19
20#[cfg(any(test, feature = "fault-injection"))]
21use std::sync::atomic::Ordering;
22use std::sync::Arc;
23use std::time::Duration;
24
25use parking_lot::Mutex;
26use tokio::sync::oneshot;
27use tokio::task::JoinHandle;
28
29use khive_storage::event::EventAppendDisposition;
30use khive_storage::{Event, EventStore, StorageError, WriterTaskRequestState};
31
32/// Coarse durability classification for an [`AuditProducer`]. See
33/// [`classify`] for the exhaustive mapping.
34#[derive(Debug, Clone, Copy, PartialEq, Eq)]
35pub(crate) enum AuditProductionClass {
36    /// The dispatch this row audits carries an obligation: the caller-visible
37    /// outcome must not silently diverge from what was durably recorded.
38    DispatchObligation,
39    /// The row is a best-effort observability signal; a non-commit degrades
40    /// gracefully rather than blocking or failing the dispatch it audits.
41    PureObservability,
42}
43
44/// Every call site that can submit a row through [`AuditBatchControl`].
45/// Adding a variant here without extending the crate-private `classify`
46/// function's match is a compile error — there is no wildcard arm.
47#[derive(Debug, Clone, Copy, PartialEq, Eq)]
48pub enum AuditProducer {
49    /// The gate denied a dispatch; the denial itself is audited.
50    GateDenied,
51    /// The gate backend was unreachable and the dispatch failed closed; the
52    /// refusal is audited best-effort: the typed `GateUnavailable` refusal
53    /// already fails the dispatch, so a lost row degrades diagnostics only.
54    GateUnavailable,
55    /// A pack dispatch returned a successful result.
56    DispatchSucceeded,
57    /// A pack dispatch returned an error result.
58    DispatchFailed,
59    /// A redirected KG read consulted the gate on its effective entity id.
60    EffectiveTargetCheck,
61    /// The gate allowed a verb no pack owns.
62    UnknownVerb,
63    /// The strict `git.digest` success receipt (schema v2).
64    GitDigestReceipt,
65    /// A drained process-lifetime `OnceLock` config-lock row.
66    ConfigLocked,
67    /// A `memory.recall` execution's pure-observability audit row.
68    ///
69    /// Classified and ready for routing, but not yet wired to a live call
70    /// site in this slice: `khive-pack-memory`'s recall handler reaches only
71    /// `KhiveRuntime` (`crates/khive-runtime/src/runtime.rs`), which does not
72    /// hold this batch seam and is outside this change's file ownership
73    /// (`final_file_ownership_r2.md` assigns `runtime.rs` to the D8 author).
74    /// Wiring this variant to `emit_recall_executed_event` needs either an
75    /// ownership-map amendment granting `runtime.rs` a narrow accessor, or
76    /// threading a second seam onto `KhiveRuntime` — both out of scope here.
77    #[allow(dead_code)]
78    RecallExecuted,
79}
80
81/// Classify an [`AuditProducer`] into its [`AuditProductionClass`]. One
82/// exhaustive match, no wildcard arm.
83pub(crate) const fn classify(producer: AuditProducer) -> AuditProductionClass {
84    match producer {
85        AuditProducer::GateDenied
86        | AuditProducer::DispatchSucceeded
87        | AuditProducer::DispatchFailed
88        | AuditProducer::EffectiveTargetCheck
89        | AuditProducer::UnknownVerb
90        | AuditProducer::GitDigestReceipt => AuditProductionClass::DispatchObligation,
91        AuditProducer::ConfigLocked
92        | AuditProducer::RecallExecuted
93        | AuditProducer::GateUnavailable => AuditProductionClass::PureObservability,
94    }
95}
96
97/// Exhaustive terminal reasons an [`AuditBatchControl::submit`],
98/// [`AuditBatchControl::quiesce`], or [`AuditBatchControl::close_and_drain`]
99/// call can resolve to. Never mapped through a wildcard arm anywhere in this
100/// module (R4).
101#[derive(Debug, Clone, Copy, PartialEq, Eq)]
102pub enum AuditTerminalReason {
103    /// `EventStore::preflight_event` rejected the row before it was ever
104    /// enqueued. The row was never counted as submitted.
105    PreflightRejected,
106    /// The batch is `Closing` or `Closed`; new admission is refused.
107    AdmissionClosed,
108    /// `AuditBatchConfig::max_pending_rows` was reached before this row was
109    /// enqueued. The row was never counted as submitted and never shared a
110    /// generation with anyone — safe to retry, and doing so applies the
111    /// obligation at most once.
112    QueueAdmissionExhausted,
113    /// The row was already enqueued (counted in `submitted_rows`) when this
114    /// caller's `AuditBatchConfig::admission_deadline` elapsed waiting for
115    /// its generation's outcome. Unlike [`Self::QueueAdmissionExhausted`],
116    /// the row was not refused: by the moment the deadline fires it may
117    /// still be sitting in `state.pending`, or the driver may have already
118    /// drained it into an in-flight generation — either way it remains
119    /// enqueued and unresolved, and is committed (or terminally failed) by
120    /// the generation driver independently of this caller's timeout, so the
121    /// caller cannot tell from this reason alone whether the row eventually
122    /// committed, or even which of those two states it was in when the
123    /// deadline elapsed. Retrying is only safe for an idempotent caller —
124    /// the prior submission may still land. When this reason degrades an
125    /// admission-degrade-safe read's own audit obligation, it is counted
126    /// separately from [`Self::QueueAdmissionExhausted`] — see
127    /// `pack::audit_admission_unresolved_obligation_count` — precisely
128    /// because a row counted here may still commit, unlike one refused
129    /// before enqueue.
130    AdmissionDeadlineExpired,
131    /// Reached only through `AuditBatch::submit_until_resolved`: the row's
132    /// [`Self::AdmissionDeadlineExpired`] wait had already elapsed, and this
133    /// caller's own `AuditBatchConfig::resolution_deadline` then also
134    /// elapsed still waiting for the row's real generation outcome. The
135    /// caller reaches this only after its domain effect has already
136    /// committed, so the effect is never retried and the row is never
137    /// re-enqueued — but unlike an ordinary commit, the caller now learns
138    /// the effect committed while the audit outcome itself is unresolved.
139    /// Same as [`Self::AdmissionDeadlineExpired`], the row is left exactly
140    /// where the driver holds it (`state.pending`, or already mid-generation)
141    /// for the driver to resolve independently — this reason performs no
142    /// removal. Kept distinct from [`Self::AdmissionDeadlineExpired`] so a
143    /// caller, and diagnostics reading this reason, can tell a merely-slow
144    /// admission wait apart from a resolution wait that gave up entirely.
145    ResolutionDeadlineExpired,
146    /// A row shared this generation's id with a previously stored row whose
147    /// columns or observation projection did not match exactly.
148    IdentityConflict,
149    /// `classify_store_error` judged the store's error non-retryable, and the
150    /// generation stopped on that attempt. Usually that is the first attempt,
151    /// but not necessarily: a generation whose earlier attempts failed
152    /// retryably and whose next one returns a non-retryable error reports
153    /// this reason too, because the attempt that decided the outcome is the
154    /// non-retryable one. Distinct from [`Self::RetryExhausted`], which the
155    /// classifier judged safe to retry on every attempt and which failed
156    /// anyway once they ran out — this reason carries no such hope. Whatever
157    /// the store returned on the deciding attempt, it is not the kind of
158    /// failure `AuditBatchConfig::max_commit_attempts` exists to ride out,
159    /// so a caller or an automated retry policy reading this reason should
160    /// not schedule a bare retry of the same call and should instead treat
161    /// it as a storage fault needing attention.
162    StoreFailure,
163    /// The store returned a `classify_store_error`-retryable error on every
164    /// one of `AuditBatchConfig::max_commit_attempts` attempts for this
165    /// generation, and the last attempt still failed retryable. Kept
166    /// distinct from [`Self::StoreFailure`] — the same way
167    /// [`Self::ResolutionDeadlineExpired`] is kept distinct from
168    /// [`Self::AdmissionDeadlineExpired`] — so a caller, and diagnostics
169    /// reading this reason, can tell a store call the classifier judged
170    /// hopeless apart from one that kept failing a condition (write-queue or
171    /// writer-task pressure, pool or timeout) the classifier judged
172    /// transient. The underlying condition may still be
173    /// transient at the moment attempts run out (a daemon restart, pool
174    /// pressure outlasting the configured backoff), so an operator or an
175    /// automated retry policy sitting above this batch can choose to wait
176    /// longer and try again rather than treating it identically to
177    /// [`Self::StoreFailure`]. This reason changes no tolerance, deadline,
178    /// retry count, or backoff on its own — it only names which of the two
179    /// causes produced the generation's failure.
180    RetryExhausted,
181    /// The configured `EventStore` backend does not implement
182    /// `append_events_idempotent`.
183    IdempotencyUnsupported,
184    /// The generation driver task panicked, or the supervisor awaiting it
185    /// unwound while armed.
186    DriverPanicked,
187    /// The generation driver task was cancelled/aborted, or the supervisor
188    /// awaiting it was dropped mid-await (shutdown abort) while armed.
189    DriverCancelled,
190    /// The child driver's `JoinHandle` was lost — dropped without ever being
191    /// inspected — so its outcome could not be classified.
192    DriverJoinLost,
193    /// The driver returned `Ok` but locked batch state was not proved
194    /// terminally consistent afterward (`in_flight` still set, or the
195    /// generation's rows were never resolved).
196    DriverExitedInconsistent,
197    /// The driver's own bound on a single generation's
198    /// `EventStore::append_events_idempotent()` call
199    /// (`supervisor_loop`'s `driver_append_deadline`) elapsed while that
200    /// call was still in flight. Unlike [`Self::AdmissionDeadlineExpired`]
201    /// and [`Self::ResolutionDeadlineExpired`], which bound only how long a
202    /// *caller* keeps waiting while the row stays wherever the driver holds
203    /// it, this bounds the driver's own hold: every waiter on the abandoned
204    /// generation is resolved with this reason and the generation is
205    /// removed from the driver's in-flight state, so a stalled store call
206    /// cannot pin `pending` at `max_pending_rows` and starve all later
207    /// admission (khive#2331). The underlying append is not cancelled — it
208    /// may not be safely cancellable mid-flight — so it is left to run to
209    /// completion in the background and its eventual result, whatever it
210    /// is, is discarded; no waiter is still listening for it. The caller's
211    /// domain effect (if any already committed before this row was
212    /// enqueued) is never retried, and the row is never re-enqueued.
213    DriverAppendAbandoned,
214    /// Before spawning a generation's child, the driver found
215    /// `AuditBatchConfig::max_abandoned_appends` detached appends already
216    /// outstanding from prior [`Self::DriverAppendAbandoned`] generations
217    /// (khive#2331). The store is treated as wedged: this generation's rows
218    /// are shed without ever attempting an append — no child task is
219    /// spawned, no store call is made, and every waiter resolves with this
220    /// reason immediately, well inside `driver_append_deadline`. Same
221    /// non-retry contract as [`Self::DriverAppendAbandoned`]: any
222    /// already-committed domain effect is not retried and the row is never
223    /// re-enqueued, and pure-observability producers record degradation.
224    /// The driver reattempts an append on the next generation as soon as an
225    /// outstanding append returns (commit or failure) and the count drops
226    /// back below the cap — recovery needs no timer of its own.
227    StoreWedged,
228}
229
230/// One row accepted for batching: the immutable event identity plus the
231/// producer that minted it, used for classification and, on failure,
232/// degradation accounting.
233pub struct PreparedAuditRow {
234    pub event: Event,
235    pub producer: AuditProducer,
236}
237
238/// What [`AuditBatchControl::submit`] resolves to on a non-error outcome.
239#[derive(Debug, Clone, Copy, PartialEq, Eq)]
240pub enum AuditCommitOutcome {
241    /// The row was freshly inserted.
242    Committed,
243    /// A prior row with the same identity already matched exactly (internal
244    /// retry replaying the same producer-minted identity).
245    AlreadyPresentIdentical,
246}
247
248/// Production-visible snapshot of [`AuditBatch::health_metrics`]. See there
249/// for field semantics. `khive_db::diagnostics::RuntimeAuditBatchMetrics` carries
250/// these three fields plus `admission_refused_obligations` and
251/// `admission_unresolved_obligations`, which are sourced from process-wide
252/// counters outside `AuditBatch` rather than from this struct — see
253/// `VerbRegistry::audit_batch_metrics`.
254#[derive(Debug, Clone, Copy, PartialEq, Eq)]
255pub struct AuditBatchHealthMetrics {
256    pub flush_failures: u64,
257    pub degraded_rows: u64,
258    pub degraded: bool,
259    /// Detached [`AuditTerminalReason::DriverAppendAbandoned`] appends that
260    /// eventually returned a commit, after their generation had already
261    /// resolved every waiter with `DriverAppendAbandoned`. Lets an operator
262    /// see that a store recorded as wedged later drained.
263    pub late_append_commits: u64,
264    /// Detached [`AuditTerminalReason::DriverAppendAbandoned`] appends that
265    /// eventually returned a failure (store error, panic, or cancellation),
266    /// after their generation had already resolved every waiter with
267    /// `DriverAppendAbandoned`.
268    pub late_append_failures: u64,
269}
270
271/// Tunables for the batch seam. Defaults are conservative; every field is
272/// exercised by at least one mechanism test.
273#[derive(Debug, Clone)]
274pub struct AuditBatchConfig {
275    pub max_pending_rows: std::num::NonZeroUsize,
276    pub max_rows_per_generation: std::num::NonZeroUsize,
277    pub max_commit_attempts: std::num::NonZeroU8,
278    pub retry_backoff: Duration,
279    pub admission_deadline: Duration,
280    /// Caps how much longer `AuditBatch::submit_until_resolved` keeps
281    /// waiting on a row's real generation outcome once `admission_deadline`
282    /// has already elapsed on it. Without this bound a generation stuck on a
283    /// stalled `EventStore::append_events_idempotent` call retains the
284    /// completed write's caller, its request slot, and its audit-lane waiter
285    /// forever, exhausting both request and audit capacity (khive#2331).
286    /// Must be at least `admission_deadline` — validated (debug-only) in
287    /// [`AuditBatch::new`]. Defaults to 6x `admission_deadline`.
288    ///
289    /// Also the single source value the driver's own per-generation append
290    /// bound (`supervisor_loop`'s `driver_append_deadline`) is derived from
291    /// — that bound exists so a stalled append cannot keep the driver from
292    /// ever draining `pending` again, which is what actually exhausts
293    /// admission for every later caller, not just the one already waiting
294    /// on the stuck row. See [`AuditTerminalReason::DriverAppendAbandoned`].
295    pub resolution_deadline: Duration,
296    /// Caps how many [`AuditTerminalReason::DriverAppendAbandoned`] appends
297    /// may be outstanding (handed to a detached task, still running) at
298    /// once. Each outstanding append retains up to
299    /// `max_rows_per_generation` events until it finally returns, so the
300    /// retained-buffer bound this places on a wedged store is
301    /// `max_abandoned_appends * max_rows_per_generation` rows — without it,
302    /// a store whose append never returns mints one such task per
303    /// `driver_append_deadline` with no cap, and both the live task count
304    /// and the retained event batches inside them grow until the process
305    /// exits (khive#2331). Once the cap is reached, `supervisor_loop` sheds
306    /// further generations with [`AuditTerminalReason::StoreWedged`] instead
307    /// of attempting another append; it resumes attempting appends as soon
308    /// as an outstanding one returns and the count drops back below the
309    /// cap. Defaults to 4.
310    pub max_abandoned_appends: std::num::NonZeroUsize,
311}
312
313impl Default for AuditBatchConfig {
314    fn default() -> Self {
315        let admission_deadline = Duration::from_secs(5);
316        Self {
317            max_pending_rows: std::num::NonZeroUsize::new(4096).unwrap(),
318            max_rows_per_generation: std::num::NonZeroUsize::new(256).unwrap(),
319            max_commit_attempts: std::num::NonZeroU8::new(3).unwrap(),
320            retry_backoff: Duration::from_millis(20),
321            admission_deadline,
322            resolution_deadline: admission_deadline * 6,
323            max_abandoned_appends: std::num::NonZeroUsize::new(4).unwrap(),
324        }
325    }
326}
327
328#[async_trait::async_trait]
329pub trait AuditBatchControl: Send + Sync {
330    async fn submit(
331        &self,
332        row: PreparedAuditRow,
333    ) -> Result<AuditCommitOutcome, AuditTerminalReason>;
334    async fn quiesce(&self) -> Result<(), AuditTerminalReason>;
335    async fn close_and_drain(&self) -> Result<(), AuditTerminalReason>;
336}
337
338#[derive(Debug, Clone, PartialEq, Eq)]
339enum Lifecycle {
340    Open,
341    Closing,
342    Closed,
343    Failed(AuditTerminalReason),
344}
345
346struct Waiting {
347    event: Event,
348    producer: AuditProducer,
349    responder: oneshot::Sender<Result<AuditCommitOutcome, AuditTerminalReason>>,
350}
351
352/// A committed/failed generation's accounting, retained for the lifetime of
353/// the process (bounded in practice by process lifetime and generation
354/// volume; this slice makes no attempt to prune history).
355#[derive(Debug, Clone, PartialEq, Eq)]
356pub struct AuditGenerationSnapshot {
357    pub generation_id: u64,
358    pub submitted_rows: u64,
359    pub committed_rows: u64,
360    pub store_batch_calls: u64,
361    pub terminal_reason: Option<AuditTerminalReason>,
362}
363
364struct State {
365    lifecycle: Lifecycle,
366    pending: Vec<Waiting>,
367    in_flight_generation: Option<u64>,
368    driver_active: bool,
369    next_generation_id: u64,
370    submitted_rows: u64,
371    committed_rows: u64,
372    store_batch_calls: u64,
373    generations: Vec<AuditGenerationSnapshot>,
374    flush_failures: u64,
375    degraded_rows: u64,
376    degraded: bool,
377    /// Detached `DriverAppendAbandoned` appends currently still running,
378    /// bounded by `AuditBatchConfig::max_abandoned_appends`. Incremented
379    /// when `supervisor_loop` abandons a generation and hands its child to a
380    /// detached task; decremented by that task once the child finally
381    /// returns.
382    outstanding_abandoned_appends: usize,
383    late_append_commits: u64,
384    late_append_failures: u64,
385}
386
387impl State {
388    fn new() -> Self {
389        Self {
390            lifecycle: Lifecycle::Open,
391            pending: Vec::new(),
392            in_flight_generation: None,
393            driver_active: false,
394            next_generation_id: 0,
395            submitted_rows: 0,
396            committed_rows: 0,
397            store_batch_calls: 0,
398            generations: Vec::new(),
399            flush_failures: 0,
400            degraded_rows: 0,
401            degraded: false,
402            outstanding_abandoned_appends: 0,
403            late_append_commits: 0,
404            late_append_failures: 0,
405        }
406    }
407
408    fn is_idle(&self) -> bool {
409        self.pending.is_empty() && self.in_flight_generation.is_none() && !self.driver_active
410    }
411}
412
413struct Inner {
414    state: Mutex<State>,
415}
416
417impl Inner {
418    /// Wins once: the first abnormal driver/supervisor exit sets `Failed`
419    /// and drains every accepted waiter (this generation's, plus anything
420    /// still queued for a future one) with the same typed reason. A later
421    /// abnormal path observes the existing terminal state and does not
422    /// double-count (R1).
423    fn fail_driver(&self, reason: AuditTerminalReason, in_flight_waiters: Vec<Waiting>) {
424        let mut state = self.state.lock();
425        let already_failed = matches!(state.lifecycle, Lifecycle::Failed(_));
426        if !already_failed {
427            state.lifecycle = Lifecycle::Failed(reason);
428        }
429        state.driver_active = false;
430        state.in_flight_generation = None;
431        let drained: Vec<Waiting> = std::mem::take(&mut state.pending);
432        if !already_failed {
433            state.flush_failures += 1;
434            state.degraded = true;
435            let generation_id = state.next_generation_id;
436            state.generations.push(AuditGenerationSnapshot {
437                generation_id,
438                submitted_rows: (in_flight_waiters.len() + drained.len()) as u64,
439                committed_rows: 0,
440                store_batch_calls: 0,
441                terminal_reason: Some(reason),
442            });
443            state.next_generation_id += 1;
444        }
445        drop(state);
446        for waiting in in_flight_waiters.into_iter().chain(drained) {
447            record_degradation_if_pure(self, waiting.producer);
448            let _ = waiting.responder.send(Err(reason));
449        }
450    }
451
452    fn record_degradation(&self) {
453        let mut state = self.state.lock();
454        state.degraded_rows += 1;
455        state.degraded = true;
456    }
457}
458
459fn record_degradation_if_pure(inner: &Inner, producer: AuditProducer) {
460    if classify(producer) == AuditProductionClass::PureObservability {
461        inner.record_degradation();
462    }
463}
464
465/// Armed for the lifetime of one generation's supervision. If dropped while
466/// still armed — the supervisor's own frame unwinding (panic) or being
467/// dropped mid-await (cancellation/shutdown abort) — the guard fails the
468/// generation before returning control to whatever tore it down, so the
469/// caller never observes a background-task-count restoration that raced
470/// ahead of the failure broadcast (R1).
471struct SupervisorGuard<'a> {
472    inner: &'a Arc<Inner>,
473    waiting: Option<Vec<Waiting>>,
474}
475
476impl<'a> SupervisorGuard<'a> {
477    fn armed(inner: &'a Arc<Inner>, waiting: Vec<Waiting>) -> Self {
478        Self {
479            inner,
480            waiting: Some(waiting),
481        }
482    }
483
484    /// Clean disarm: the caller has fully classified the outcome and taken
485    /// ownership of the waiters to resolve them itself.
486    fn disarm(mut self) -> Vec<Waiting> {
487        self.waiting.take().unwrap_or_default()
488    }
489
490    /// Explicit failure classification reached without unwinding/cancelling
491    /// this frame (e.g. a `JoinError` was returned normally). Equivalent to
492    /// what `Drop` does when armed, but callable inline.
493    fn fail(mut self, reason: AuditTerminalReason) {
494        if let Some(waiting) = self.waiting.take() {
495            self.inner.fail_driver(reason, waiting);
496        }
497    }
498}
499
500impl Drop for SupervisorGuard<'_> {
501    fn drop(&mut self) {
502        let Some(waiting) = self.waiting.take() else {
503            return;
504        };
505        let reason = if std::thread::panicking() {
506            AuditTerminalReason::DriverPanicked
507        } else {
508            AuditTerminalReason::DriverCancelled
509        };
510        self.inner.fail_driver(reason, waiting);
511    }
512}
513
514enum RetryDecision {
515    Retry,
516    Terminal(AuditTerminalReason),
517}
518
519fn classify_store_error(err: &StorageError) -> RetryDecision {
520    match err {
521        StorageError::WriteQueueFull { .. } | StorageError::WriterTaskBusy { .. } => {
522            RetryDecision::Retry
523        }
524        // Transient availability conditions, not judgments on the batch. The
525        // events-daemon forwarding lane (ADR-170) reports an unreachable or
526        // stalled daemon as `Pool`/`Timeout`; the direct SQL path reports
527        // acquisition pressure the same way. Both are exactly what the
528        // configured bounded retries exist for — treating them as terminal
529        // would abandon a generation on the first blip of a daemon restart.
530        StorageError::Pool { .. } | StorageError::Timeout { .. } => RetryDecision::Retry,
531        StorageError::WriterTaskRequestFailed { request_state, .. }
532        | StorageError::WriterTaskTerminated { request_state } => match request_state {
533            WriterTaskRequestState::NotStarted | WriterTaskRequestState::TransactionRolledBack => {
534                RetryDecision::Retry
535            }
536            WriterTaskRequestState::SideEffectsUnknown => RetryDecision::Retry,
537        },
538        StorageError::Unsupported { operation, .. }
539            if operation.as_ref() == "append_events_idempotent" =>
540        {
541            RetryDecision::Terminal(AuditTerminalReason::IdempotencyUnsupported)
542        }
543        _ => RetryDecision::Terminal(AuditTerminalReason::StoreFailure),
544    }
545}
546
547enum GenerationResult {
548    Committed(Vec<EventAppendDisposition>),
549    Failed(AuditTerminalReason),
550    /// Fault-injection only: proves the supervisor's post-`Ok` consistency
551    /// check actually runs.
552    #[cfg_attr(not(any(test, feature = "fault-injection")), allow(dead_code))]
553    FakedInconsistent,
554}
555
556/// The driver's own bound on how long `supervisor_loop` waits for one
557/// generation's `run_generation` child task before abandoning it
558/// (khive#2331; see [`AuditTerminalReason::DriverAppendAbandoned`]).
559///
560/// Derived from `resolution_deadline` rather than a second config field, at
561/// 3x it. That margin is load-bearing, not arbitrary: `AuditBatch::new`
562/// already enforces `resolution_deadline >= admission_deadline`, so the
563/// worst-case *caller*-side wait for a `submit_until_resolved` row —
564/// `admission_deadline` then `resolution_deadline` in sequence — is always
565/// `<= 2 * resolution_deadline`. Bounding the driver at `3 *
566/// resolution_deadline` guarantees it can never resolve a still-waiting
567/// caller's row with `DriverAppendAbandoned` ahead of that caller's own,
568/// more specific `ResolutionDeadlineExpired` reason.
569///
570/// This bound alone only stops the driver from holding one stalled
571/// generation forever — a store whose append never returns still mints one
572/// detached append per `driver_append_deadline` with nothing capping how
573/// many run at once. `AuditBatchConfig::max_abandoned_appends` closes that:
574/// once that many detached appends are outstanding, further generations are
575/// shed with [`AuditTerminalReason::StoreWedged`] instead of attempting
576/// another append. Combined, the two bounds guarantee at most
577/// `max_abandoned_appends` live detached appends at any time, at most that
578/// many retained event batches, and a shed generation costs no store work
579/// at all — it resolves inside this function's caller without ever calling
580/// `run_generation`.
581fn driver_append_deadline(config: &AuditBatchConfig) -> Duration {
582    config.resolution_deadline.saturating_mul(3)
583}
584
585async fn run_generation(
586    store: Arc<dyn EventStore>,
587    events: Vec<Event>,
588    config: Arc<AuditBatchConfig>,
589) -> GenerationResult {
590    #[cfg(any(test, feature = "fault-injection"))]
591    if fault::CHILD_PANIC.swap(false, Ordering::SeqCst) {
592        panic!("adr133 fault injection: audit_batch child_panic");
593    }
594    #[cfg(any(test, feature = "fault-injection"))]
595    if fault::INCONSISTENT_EXIT.swap(false, Ordering::SeqCst) {
596        return GenerationResult::FakedInconsistent;
597    }
598
599    let mut attempt: u8 = 0;
600    loop {
601        attempt += 1;
602        match store.append_events_idempotent(events.clone()).await {
603            Ok(result) => return GenerationResult::Committed(result.rows),
604            Err(err) => match classify_store_error(&err) {
605                RetryDecision::Retry if attempt < config.max_commit_attempts.get() => {
606                    tokio::time::sleep(config.retry_backoff).await;
607                }
608                RetryDecision::Retry => {
609                    return GenerationResult::Failed(AuditTerminalReason::RetryExhausted)
610                }
611                RetryDecision::Terminal(reason) => return GenerationResult::Failed(reason),
612            },
613        }
614    }
615}
616
617/// The batch owner. Constructed once per configured `EventStore`; every
618/// dispatch-audit call site routes its row through [`AuditBatch::submit`]
619/// instead of taking its own writer-task acquisition.
620pub struct AuditBatch {
621    inner: Arc<Inner>,
622    store: Arc<dyn EventStore>,
623    config: Arc<AuditBatchConfig>,
624    supervisor: Mutex<Option<JoinHandle<()>>>,
625}
626
627impl AuditBatch {
628    pub fn new(store: Arc<dyn EventStore>, config: AuditBatchConfig) -> Arc<Self> {
629        debug_assert!(
630            config.resolution_deadline >= config.admission_deadline,
631            "resolution_deadline ({:?}) must be at least admission_deadline ({:?})",
632            config.resolution_deadline,
633            config.admission_deadline
634        );
635        Arc::new(Self {
636            inner: Arc::new(Inner {
637                state: Mutex::new(State::new()),
638            }),
639            store,
640            config: Arc::new(config),
641            supervisor: Mutex::new(None),
642        })
643    }
644
645    /// Process-lifetime audit-batch health counters, for the registry/
646    /// runtime owner to feed into `db_diagnostics` (D8's operator surface).
647    /// Unlike `test_internals::AuditBatchSnapshot::metrics_snapshot`,
648    /// which is test-only (the module is cfg-gated and invisible to the
649    /// default doc build, so an intra-doc link cannot resolve), this is
650    /// always available.
651    pub fn health_metrics(&self) -> AuditBatchHealthMetrics {
652        let state = self.inner.state.lock();
653        AuditBatchHealthMetrics {
654            flush_failures: state.flush_failures,
655            degraded_rows: state.degraded_rows,
656            degraded: state.degraded,
657            late_append_commits: state.late_append_commits,
658            late_append_failures: state.late_append_failures,
659        }
660    }
661
662    fn spawn_supervisor_if_idle(&self) {
663        let inner = self.inner.clone();
664        let store = self.store.clone();
665        let config = self.config.clone();
666        let handle = tokio::spawn(async move {
667            supervisor_loop(inner, store, config).await;
668        });
669        *self.supervisor.lock() = Some(handle);
670    }
671
672    /// Enqueue one row and keep waiting for its real generation outcome past
673    /// the ordinary admission wait deadline, up to `resolution_deadline`.
674    ///
675    /// This is the narrow khive#2256 seam for a successful operation whose
676    /// domain effect has already committed. Returning
677    /// [`AuditTerminalReason::AdmissionDeadlineExpired`] there would report a
678    /// false operation failure and invite an unsafe retry while the same
679    /// audit row remains enqueued. Pre-enqueue refusal and genuine terminal
680    /// generation failures still return normally. The post-admission wait is
681    /// itself bounded by `resolution_deadline`
682    /// ([`AuditTerminalReason::ResolutionDeadlineExpired`]) so a stalled
683    /// store cannot retain this caller, its request slot, and its audit-lane
684    /// waiter forever (khive#2331).
685    ///
686    /// The two deadlines therefore mean different things to the caller, and
687    /// callers key on the difference. Admission expiry means the row is
688    /// enqueued and this seam keeps waiting; the generation commits it
689    /// independently, so a dispatch that reaches it reports its committed
690    /// result and counts the row as unresolved. Resolution expiry means the
691    /// caller waited for the real outcome and never received one: the commit
692    /// is unconfirmed, not proven absent. That is returned as a terminal
693    /// reason and the dispatch propagates it as a structured
694    /// committed-outcome error carrying the domain result, with no retryable
695    /// context, for every verb rather than only for the receipt verb. Never
696    /// replay the handler or re-enqueue the row on it; the original driver
697    /// work continues on its own.
698    pub(crate) async fn submit_until_resolved(
699        &self,
700        row: PreparedAuditRow,
701    ) -> Result<AuditCommitOutcome, AuditTerminalReason> {
702        self.submit_with_wait_policy(row, true).await
703    }
704
705    async fn submit_with_wait_policy(
706        &self,
707        row: PreparedAuditRow,
708        wait_until_resolved: bool,
709    ) -> Result<AuditCommitOutcome, AuditTerminalReason> {
710        // Pre-enqueue validation (invariant 3): a malformed row is rejected
711        // before it can share a generation with anyone else's.
712        if self.store.preflight_event(&row.event).is_err() {
713            return Err(AuditTerminalReason::PreflightRejected);
714        }
715
716        let producer = row.producer;
717        let (tx, mut rx) = oneshot::channel();
718        let need_spawn = {
719            let mut state = self.inner.state.lock();
720            match state.lifecycle {
721                Lifecycle::Closed | Lifecycle::Closing => {
722                    return Err(AuditTerminalReason::AdmissionClosed)
723                }
724                Lifecycle::Failed(reason) => return Err(reason),
725                Lifecycle::Open => {}
726            }
727            if state.pending.len() >= self.config.max_pending_rows.get() {
728                return Err(AuditTerminalReason::QueueAdmissionExhausted);
729            }
730            state.pending.push(Waiting {
731                event: row.event,
732                producer,
733                responder: tx,
734            });
735            state.submitted_rows += 1;
736            let need_spawn = !state.driver_active;
737            if need_spawn {
738                state.driver_active = true;
739            }
740            need_spawn
741        };
742        if need_spawn {
743            self.spawn_supervisor_if_idle();
744        }
745
746        if wait_until_resolved {
747            match tokio::time::timeout(self.config.admission_deadline, &mut rx).await {
748                Ok(Ok(result)) => return result,
749                Ok(Err(_recv_error)) => return Err(AuditTerminalReason::DriverJoinLost),
750                Err(_elapsed) => tracing::warn!(
751                    ?producer,
752                    "strict audit obligation remains enqueued after the admission wait deadline; \
753                     waiting up to the resolution deadline for its real terminal outcome"
754                ),
755            }
756            return match tokio::time::timeout(self.config.resolution_deadline, rx).await {
757                Ok(Ok(result)) => result,
758                Ok(Err(_recv_error)) => Err(AuditTerminalReason::DriverJoinLost),
759                // Same non-removal contract as the `AdmissionDeadlineExpired`
760                // arm below: the row is left exactly where the driver holds
761                // it. The caller's domain effect already committed by the
762                // time it reached this wait, so it is never retried here;
763                // the audit outcome itself is what remains unresolved.
764                Err(_elapsed) => {
765                    tracing::warn!(
766                        ?producer,
767                        "strict audit obligation remains unresolved after the resolution \
768                         deadline; giving up on this caller's wait — the committed effect is \
769                         not retried and the row is not re-enqueued"
770                    );
771                    Err(AuditTerminalReason::ResolutionDeadlineExpired)
772                }
773            };
774        }
775
776        match tokio::time::timeout(self.config.admission_deadline, rx).await {
777            Ok(Ok(result)) => result,
778            Ok(Err(_recv_error)) => Err(AuditTerminalReason::DriverJoinLost),
779            // The row was already pushed onto `state.pending` above (and
780            // `submitted_rows` incremented) before this wait began — this is
781            // a deadline elapsing on an enqueued row, not a queue-full
782            // refusal, so it gets its own terminal reason (khive#2117,
783            // khive#2208). The row is left in place for the driver to drain;
784            // this arm performs no removal.
785            Err(_elapsed) => Err(AuditTerminalReason::AdmissionDeadlineExpired),
786        }
787    }
788}
789
790async fn supervisor_loop(
791    inner: Arc<Inner>,
792    store: Arc<dyn EventStore>,
793    config: Arc<AuditBatchConfig>,
794) {
795    loop {
796        let (waiting, wedged_generation_id) = {
797            let mut state = inner.state.lock();
798            if matches!(state.lifecycle, Lifecycle::Failed(_)) {
799                state.driver_active = false;
800                break;
801            }
802            if state.pending.is_empty() {
803                state.driver_active = false;
804                break;
805            }
806            let take_n = state
807                .pending
808                .len()
809                .min(config.max_rows_per_generation.get());
810            let waiting: Vec<Waiting> = state.pending.drain(..take_n).collect();
811            let generation_id = state.next_generation_id;
812            state.next_generation_id += 1;
813            // khive#2331: a store whose append never returns must not mint
814            // an unbounded number of detached appends. Before committing to
815            // spawning this generation's child, check whether the cap is
816            // already saturated by prior abandoned generations still
817            // running in the background — if so, this generation is shed
818            // below instead of ever calling the store.
819            if state.outstanding_abandoned_appends >= config.max_abandoned_appends.get() {
820                (waiting, Some(generation_id))
821            } else {
822                state.in_flight_generation = Some(generation_id);
823                (waiting, None)
824            }
825        };
826
827        if let Some(generation_id) = wedged_generation_id {
828            tracing::warn!(
829                generation_id,
830                max_abandoned_appends = config.max_abandoned_appends.get(),
831                "audit store treated as wedged: max_abandoned_appends detached appends are \
832                 already outstanding; shedding this generation without attempting an append"
833            );
834            let submitted = waiting.len() as u64;
835            {
836                let mut state = inner.state.lock();
837                state.flush_failures += 1;
838                state.generations.push(AuditGenerationSnapshot {
839                    generation_id,
840                    submitted_rows: submitted,
841                    committed_rows: 0,
842                    store_batch_calls: 0,
843                    terminal_reason: Some(AuditTerminalReason::StoreWedged),
844                });
845            }
846            for w in waiting {
847                record_degradation_if_pure(&inner, w.producer);
848                let _ = w.responder.send(Err(AuditTerminalReason::StoreWedged));
849            }
850            continue;
851        }
852
853        #[cfg(any(test, feature = "fault-injection"))]
854        if fault::SUPERVISOR_PANIC.swap(false, Ordering::SeqCst) {
855            let _guard = SupervisorGuard::armed(&inner, waiting);
856            panic!("adr133 fault injection: audit_batch supervisor_panic");
857        }
858        #[cfg(any(test, feature = "fault-injection"))]
859        if fault::SUPERVISOR_SLEEP_BEFORE_SPAWN.swap(false, Ordering::SeqCst) {
860            let guard = SupervisorGuard::armed(&inner, waiting);
861            tokio::time::sleep(Duration::from_secs(3600)).await;
862            drop(guard);
863            continue;
864        }
865
866        let guard = SupervisorGuard::armed(&inner, waiting);
867        let events: Vec<Event> = guard
868            .waiting
869            .as_ref()
870            .expect("guard freshly armed")
871            .iter()
872            .map(|w| w.event.clone())
873            .collect();
874
875        let mut child: JoinHandle<GenerationResult> =
876            tokio::spawn(run_generation(store.clone(), events, config.clone()));
877
878        #[cfg(any(test, feature = "fault-injection"))]
879        if fault::CHILD_CANCEL.swap(false, Ordering::SeqCst) {
880            child.abort();
881        }
882        #[cfg(any(test, feature = "fault-injection"))]
883        if fault::JOIN_LOST.swap(false, Ordering::SeqCst) {
884            drop(child);
885            guard.fail(AuditTerminalReason::DriverJoinLost);
886            continue;
887        }
888
889        let append_deadline = driver_append_deadline(&config);
890        let join_result = match tokio::time::timeout(append_deadline, &mut child).await {
891            Ok(join_result) => join_result,
892            Err(_elapsed) => {
893                // The append is still in flight and may not be safely
894                // cancellable mid-write (a blocking storage call cannot be
895                // forced to stop), so `child` is not aborted here — it is
896                // handed to a detached reaper that drives it to completion
897                // and discards whatever it eventually returns. Every waiter
898                // on this generation is resolved now, and the driver loops
899                // back to `pending` immediately instead of staying pinned on
900                // this one stalled call, which is what actually starves
901                // later admission (khive#2331).
902                tracing::warn!(
903                    ?append_deadline,
904                    "audit generation append exceeded the driver's own bound; abandoning \
905                     this generation so admission capacity recovers for later callers"
906                );
907                let waiting = guard.disarm();
908                let submitted = waiting.len() as u64;
909                {
910                    let mut state = inner.state.lock();
911                    state.in_flight_generation = None;
912                    state.flush_failures += 1;
913                    state.outstanding_abandoned_appends += 1;
914                    let generation_id = state.next_generation_id.saturating_sub(1);
915                    state.generations.push(AuditGenerationSnapshot {
916                        generation_id,
917                        submitted_rows: submitted,
918                        committed_rows: 0,
919                        store_batch_calls: 0,
920                        terminal_reason: Some(AuditTerminalReason::DriverAppendAbandoned),
921                    });
922                }
923                for w in waiting {
924                    record_degradation_if_pure(&inner, w.producer);
925                    let _ = w
926                        .responder
927                        .send(Err(AuditTerminalReason::DriverAppendAbandoned));
928                }
929                let reaper_inner = inner.clone();
930                tokio::spawn(async move {
931                    let result = child.await;
932                    let mut state = reaper_inner.state.lock();
933                    state.outstanding_abandoned_appends =
934                        state.outstanding_abandoned_appends.saturating_sub(1);
935                    match result {
936                        Ok(GenerationResult::Committed(_)) => state.late_append_commits += 1,
937                        Ok(GenerationResult::Failed(_) | GenerationResult::FakedInconsistent)
938                        | Err(_) => state.late_append_failures += 1,
939                    }
940                });
941                continue;
942            }
943        };
944        match join_result {
945            Err(join_err) => {
946                let reason = if join_err.is_panic() {
947                    AuditTerminalReason::DriverPanicked
948                } else {
949                    AuditTerminalReason::DriverCancelled
950                };
951                guard.fail(reason);
952            }
953            Ok(GenerationResult::FakedInconsistent) => {
954                guard.fail(AuditTerminalReason::DriverExitedInconsistent);
955            }
956            Ok(GenerationResult::Failed(reason)) => {
957                let waiting = guard.disarm();
958                let submitted = waiting.len() as u64;
959                {
960                    let mut state = inner.state.lock();
961                    state.in_flight_generation = None;
962                    state.store_batch_calls += 1;
963                    state.flush_failures += 1;
964                    state.degraded = true;
965                    let generation_id = state.next_generation_id.saturating_sub(1);
966                    state.generations.push(AuditGenerationSnapshot {
967                        generation_id,
968                        submitted_rows: submitted,
969                        committed_rows: 0,
970                        store_batch_calls: 1,
971                        terminal_reason: Some(reason),
972                    });
973                }
974                for w in waiting {
975                    record_degradation_if_pure(&inner, w.producer);
976                    let _ = w.responder.send(Err(reason));
977                }
978            }
979            Ok(GenerationResult::Committed(dispositions)) => {
980                let waiting = guard.disarm();
981                if dispositions.len() != waiting.len() {
982                    // Defensive: the store must preserve input order/length.
983                    // Treat a mismatch as the driver having exited in a
984                    // state that cannot be reconciled with the accepted
985                    // waiters.
986                    {
987                        let mut state = inner.state.lock();
988                        state.in_flight_generation = None;
989                    }
990                    inner.fail_driver(AuditTerminalReason::DriverExitedInconsistent, waiting);
991                    continue;
992                }
993                let mut committed_n = 0u64;
994                for d in &dispositions {
995                    if !matches!(d, EventAppendDisposition::IdentityConflict) {
996                        committed_n += 1;
997                    }
998                }
999                let submitted = dispositions.len() as u64;
1000                {
1001                    let mut state = inner.state.lock();
1002                    state.in_flight_generation = None;
1003                    state.store_batch_calls += 1;
1004                    state.committed_rows += committed_n;
1005                    let generation_id = state.next_generation_id.saturating_sub(1);
1006                    state.generations.push(AuditGenerationSnapshot {
1007                        generation_id,
1008                        submitted_rows: submitted,
1009                        committed_rows: committed_n,
1010                        store_batch_calls: 1,
1011                        terminal_reason: None,
1012                    });
1013                }
1014                for (w, disposition) in waiting.into_iter().zip(dispositions) {
1015                    let result = match disposition {
1016                        EventAppendDisposition::Inserted => Ok(AuditCommitOutcome::Committed),
1017                        EventAppendDisposition::AlreadyPresentIdentical => {
1018                            Ok(AuditCommitOutcome::AlreadyPresentIdentical)
1019                        }
1020                        EventAppendDisposition::IdentityConflict => {
1021                            record_degradation_if_pure(&inner, w.producer);
1022                            Err(AuditTerminalReason::IdentityConflict)
1023                        }
1024                    };
1025                    let _ = w.responder.send(result);
1026                }
1027            }
1028        }
1029    }
1030}
1031
1032#[async_trait::async_trait]
1033impl AuditBatchControl for AuditBatch {
1034    async fn submit(
1035        &self,
1036        row: PreparedAuditRow,
1037    ) -> Result<AuditCommitOutcome, AuditTerminalReason> {
1038        self.submit_with_wait_policy(row, false).await
1039    }
1040
1041    async fn quiesce(&self) -> Result<(), AuditTerminalReason> {
1042        loop {
1043            let (lifecycle, idle) = {
1044                let state = self.inner.state.lock();
1045                (state.lifecycle.clone(), state.is_idle())
1046            };
1047            match lifecycle {
1048                Lifecycle::Failed(reason) => return Err(reason),
1049                Lifecycle::Closed => return Ok(()),
1050                _ if idle => return Ok(()),
1051                _ => tokio::time::sleep(Duration::from_millis(2)).await,
1052            }
1053        }
1054    }
1055
1056    async fn close_and_drain(&self) -> Result<(), AuditTerminalReason> {
1057        {
1058            let mut state = self.inner.state.lock();
1059            if matches!(state.lifecycle, Lifecycle::Open) {
1060                state.lifecycle = Lifecycle::Closing;
1061            }
1062        }
1063        let result = AuditBatchControl::quiesce(self).await;
1064        {
1065            let mut state = self.inner.state.lock();
1066            if result.is_ok() && matches!(state.lifecycle, Lifecycle::Closing) {
1067                state.lifecycle = Lifecycle::Closed;
1068            }
1069        }
1070        let handle = self.supervisor.lock().take();
1071        if let Some(handle) = handle {
1072            let _ = handle.await;
1073        }
1074        result
1075    }
1076}
1077
1078#[cfg(any(test, feature = "fault-injection"))]
1079mod fault {
1080    use std::sync::atomic::AtomicBool;
1081
1082    pub(super) static CHILD_PANIC: AtomicBool = AtomicBool::new(false);
1083    pub(super) static CHILD_CANCEL: AtomicBool = AtomicBool::new(false);
1084    pub(super) static SUPERVISOR_PANIC: AtomicBool = AtomicBool::new(false);
1085    pub(super) static SUPERVISOR_SLEEP_BEFORE_SPAWN: AtomicBool = AtomicBool::new(false);
1086    pub(super) static JOIN_LOST: AtomicBool = AtomicBool::new(false);
1087    pub(super) static INCONSISTENT_EXIT: AtomicBool = AtomicBool::new(false);
1088}
1089
1090/// Deterministic fault-injection arms consumed by exactly one subsequent
1091/// generation each. Test-only surface, gated identically to the rest of this
1092/// crate's `fault-injection` fixtures (see `crate::operations`).
1093#[cfg(any(test, feature = "fault-injection"))]
1094pub mod fault_injection {
1095    use std::sync::atomic::Ordering;
1096
1097    pub fn arm_child_panic() {
1098        super::fault::CHILD_PANIC.store(true, Ordering::SeqCst);
1099    }
1100    pub fn arm_child_cancel() {
1101        super::fault::CHILD_CANCEL.store(true, Ordering::SeqCst);
1102    }
1103    pub fn arm_supervisor_panic() {
1104        super::fault::SUPERVISOR_PANIC.store(true, Ordering::SeqCst);
1105    }
1106    pub fn arm_supervisor_sleep_before_spawn() {
1107        super::fault::SUPERVISOR_SLEEP_BEFORE_SPAWN.store(true, Ordering::SeqCst);
1108    }
1109    pub fn arm_join_lost() {
1110        super::fault::JOIN_LOST.store(true, Ordering::SeqCst);
1111    }
1112    pub fn arm_inconsistent_exit() {
1113        super::fault::INCONSISTENT_EXIT.store(true, Ordering::SeqCst);
1114    }
1115}
1116
1117// ── Test-attribution surface (R2) ────────────────────────────────────────
1118//
1119// Additive only: production attribution still lands in the public
1120// `db_diagnostics` counters via `metrics_snapshot`-shaped data supplied at
1121// that seam. This surface never substitutes for it.
1122#[cfg(any(test, feature = "test-internals"))]
1123mod test_internals {
1124    use super::*;
1125
1126    /// A point-in-time view of the batch's counters and generation history,
1127    /// used by mechanism tests to compute a delta across a measured
1128    /// operation.
1129    #[derive(Debug, Clone)]
1130    pub struct AuditBatchSnapshot {
1131        pub pending_rows: usize,
1132        pub in_flight_generation: Option<u64>,
1133        pub driver_active: bool,
1134        pub next_generation_id: u64,
1135        pub submitted_rows: u64,
1136        pub committed_rows: u64,
1137        pub store_batch_calls: u64,
1138        pub per_generation: Vec<AuditGenerationSnapshot>,
1139        /// Live count of detached `DriverAppendAbandoned` appends still
1140        /// running, bounded by `AuditBatchConfig::max_abandoned_appends`.
1141        pub outstanding_abandoned_appends: usize,
1142    }
1143
1144    impl AuditBatchSnapshot {
1145        pub fn is_idle(&self) -> bool {
1146            self.pending_rows == 0 && self.in_flight_generation.is_none() && !self.driver_active
1147        }
1148    }
1149
1150    #[derive(Debug, Clone, Copy)]
1151    pub struct AuditBatchMetricsSnapshot {
1152        pub flush_failures: u64,
1153        pub degraded_rows: u64,
1154        pub degraded: bool,
1155        pub late_append_commits: u64,
1156        pub late_append_failures: u64,
1157    }
1158
1159    #[derive(Debug, Clone, PartialEq)]
1160    pub struct AuditBatchDelta {
1161        pub submitted_rows: u64,
1162        pub committed_rows: u64,
1163        pub store_batch_calls: u64,
1164        pub per_generation: Vec<AuditGenerationSnapshot>,
1165    }
1166
1167    #[derive(Debug, Clone, Copy, PartialEq, Eq)]
1168    pub enum AuditSnapshotError {
1169        CounterRegressed,
1170        GenerationHistoryRegressed,
1171    }
1172
1173    impl AuditBatch {
1174        pub fn test_snapshot(&self) -> AuditBatchSnapshot {
1175            let state = self.inner.state.lock();
1176            AuditBatchSnapshot {
1177                pending_rows: state.pending.len(),
1178                in_flight_generation: state.in_flight_generation,
1179                driver_active: state.driver_active,
1180                next_generation_id: state.next_generation_id,
1181                submitted_rows: state.submitted_rows,
1182                committed_rows: state.committed_rows,
1183                store_batch_calls: state.store_batch_calls,
1184                per_generation: state.generations.clone(),
1185                outstanding_abandoned_appends: state.outstanding_abandoned_appends,
1186            }
1187        }
1188
1189        pub fn metrics_snapshot(&self) -> AuditBatchMetricsSnapshot {
1190            let state = self.inner.state.lock();
1191            AuditBatchMetricsSnapshot {
1192                flush_failures: state.flush_failures,
1193                degraded_rows: state.degraded_rows,
1194                degraded: state.degraded,
1195                late_append_commits: state.late_append_commits,
1196                late_append_failures: state.late_append_failures,
1197            }
1198        }
1199
1200        /// Abort the currently-retained supervisor `JoinHandle`, if any,
1201        /// simulating a shutdown abort landing mid-generation. Returns
1202        /// whether a handle was found and aborted. Test-only: exercises the
1203        /// R1 supervisor-cancellation path without reaching into a private
1204        /// field from outside this module.
1205        pub fn test_abort_supervisor(&self) -> bool {
1206            let handle = self.supervisor.lock().take();
1207            match handle {
1208                Some(handle) => {
1209                    handle.abort();
1210                    true
1211                }
1212                None => false,
1213            }
1214        }
1215    }
1216
1217    /// Checked monotonic subtraction. Rejects a regressed counter or a
1218    /// generation history that is not an append-only extension of `before`.
1219    pub fn audit_delta(
1220        before: &AuditBatchSnapshot,
1221        after: &AuditBatchSnapshot,
1222    ) -> Result<AuditBatchDelta, AuditSnapshotError> {
1223        let submitted_rows = after
1224            .submitted_rows
1225            .checked_sub(before.submitted_rows)
1226            .ok_or(AuditSnapshotError::CounterRegressed)?;
1227        let committed_rows = after
1228            .committed_rows
1229            .checked_sub(before.committed_rows)
1230            .ok_or(AuditSnapshotError::CounterRegressed)?;
1231        let store_batch_calls = after
1232            .store_batch_calls
1233            .checked_sub(before.store_batch_calls)
1234            .ok_or(AuditSnapshotError::CounterRegressed)?;
1235        if after.per_generation.len() < before.per_generation.len() {
1236            return Err(AuditSnapshotError::GenerationHistoryRegressed);
1237        }
1238        if after.per_generation[..before.per_generation.len()] != before.per_generation[..] {
1239            return Err(AuditSnapshotError::GenerationHistoryRegressed);
1240        }
1241        let per_generation = after.per_generation[before.per_generation.len()..].to_vec();
1242        Ok(AuditBatchDelta {
1243            submitted_rows,
1244            committed_rows,
1245            store_batch_calls,
1246            per_generation,
1247        })
1248    }
1249
1250    /// Exhaustive, non-wildcard producer classification (see
1251    /// [`super::AuditProducer`] and [`super::classify`], the real
1252    /// crate-private definitions this mirrors one-for-one). The doctest
1253    /// below proves the general property those definitions rely on: a match
1254    /// over a non-`#[non_exhaustive]` enum that omits a variant, with no
1255    /// wildcard arm to silently absorb it, fails to compile rather than
1256    /// passing an incomplete classification.
1257    ///
1258    /// ```compile_fail
1259    /// enum AuditProducer {
1260    ///     GateDenied,
1261    ///     DispatchSucceeded,
1262    ///     DispatchFailed,
1263    ///     UnknownVerb,
1264    ///     GitDigestReceipt,
1265    ///     ConfigLocked,
1266    ///     RecallExecuted,
1267    /// }
1268    ///
1269    /// fn describe(p: AuditProducer) -> &'static str {
1270    ///     match p {
1271    ///         AuditProducer::GateDenied => "obligation",
1272    ///         AuditProducer::DispatchSucceeded => "obligation",
1273    ///         AuditProducer::DispatchFailed => "obligation",
1274    ///         AuditProducer::UnknownVerb => "obligation",
1275    ///         AuditProducer::GitDigestReceipt => "obligation",
1276    ///         AuditProducer::ConfigLocked => "observability",
1277    ///         // RecallExecuted intentionally omitted: a non-exhaustive
1278    ///         // match must fail to compile, proving no variant can
1279    ///         // silently fall through an absent wildcard arm.
1280    ///     }
1281    /// }
1282    /// ```
1283    #[allow(dead_code)]
1284    struct DoctestAnchor;
1285}
1286
1287#[cfg(any(test, feature = "test-internals"))]
1288pub use test_internals::*;