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        // The writer task wraps any request operation that failed and was rolled
532        // back, whatever the cause, so the wrapper alone does not say whether a
533        // retry can succeed: a missing column fails the same way on every
534        // attempt. Judge the cause it carries. A request left in an unknown
535        // state is replayed regardless, because the append is idempotent.
536        StorageError::WriterTaskRequestFailed {
537            request_state,
538            source,
539        } => match request_state {
540            WriterTaskRequestState::NotStarted | WriterTaskRequestState::TransactionRolledBack => {
541                if is_sqlite_busy_or_locked(source) {
542                    RetryDecision::Retry
543                } else {
544                    classify_store_error(source)
545                }
546            }
547            WriterTaskRequestState::SideEffectsUnknown => RetryDecision::Retry,
548        },
549        StorageError::WriterTaskTerminated { request_state } => match request_state {
550            WriterTaskRequestState::NotStarted | WriterTaskRequestState::TransactionRolledBack => {
551                RetryDecision::Retry
552            }
553            WriterTaskRequestState::SideEffectsUnknown => RetryDecision::Retry,
554        },
555        StorageError::Unsupported { operation, .. }
556            if operation.as_ref() == "append_events_idempotent" =>
557        {
558            RetryDecision::Terminal(AuditTerminalReason::IdempotencyUnsupported)
559        }
560        _ => RetryDecision::Terminal(AuditTerminalReason::StoreFailure),
561    }
562}
563
564/// Whether a failed writer request's preserved cause is SQLite contention. The
565/// writer keeps the request body's driver error as the cause, and a driver
566/// error is not one of the typed transient variants `classify_store_error`
567/// retries on its own.
568fn is_sqlite_busy_or_locked(err: &StorageError) -> bool {
569    let StorageError::Driver { source, .. } = err else {
570        return false;
571    };
572    matches!(
573        source
574            .downcast_ref::<rusqlite::Error>()
575            .and_then(|error| error.sqlite_error_code()),
576        Some(rusqlite::ErrorCode::DatabaseBusy | rusqlite::ErrorCode::DatabaseLocked)
577    )
578}
579
580enum GenerationResult {
581    Committed(Vec<EventAppendDisposition>),
582    Failed(AuditTerminalReason),
583    /// Fault-injection only: proves the supervisor's post-`Ok` consistency
584    /// check actually runs.
585    #[cfg_attr(not(any(test, feature = "fault-injection")), allow(dead_code))]
586    FakedInconsistent,
587}
588
589/// The driver's own bound on how long `supervisor_loop` waits for one
590/// generation's `run_generation` child task before abandoning it
591/// (khive#2331; see [`AuditTerminalReason::DriverAppendAbandoned`]).
592///
593/// Derived from `resolution_deadline` rather than a second config field, at
594/// 3x it. That margin is load-bearing, not arbitrary: `AuditBatch::new`
595/// already enforces `resolution_deadline >= admission_deadline`, so the
596/// worst-case *caller*-side wait for a `submit_until_resolved` row —
597/// `admission_deadline` then `resolution_deadline` in sequence — is always
598/// `<= 2 * resolution_deadline`. Bounding the driver at `3 *
599/// resolution_deadline` guarantees it can never resolve a still-waiting
600/// caller's row with `DriverAppendAbandoned` ahead of that caller's own,
601/// more specific `ResolutionDeadlineExpired` reason.
602///
603/// This bound alone only stops the driver from holding one stalled
604/// generation forever — a store whose append never returns still mints one
605/// detached append per `driver_append_deadline` with nothing capping how
606/// many run at once. `AuditBatchConfig::max_abandoned_appends` closes that:
607/// once that many detached appends are outstanding, further generations are
608/// shed with [`AuditTerminalReason::StoreWedged`] instead of attempting
609/// another append. Combined, the two bounds guarantee at most
610/// `max_abandoned_appends` live detached appends at any time, at most that
611/// many retained event batches, and a shed generation costs no store work
612/// at all — it resolves inside this function's caller without ever calling
613/// `run_generation`.
614fn driver_append_deadline(config: &AuditBatchConfig) -> Duration {
615    config.resolution_deadline.saturating_mul(3)
616}
617
618async fn run_generation(
619    store: Arc<dyn EventStore>,
620    events: Vec<Event>,
621    config: Arc<AuditBatchConfig>,
622) -> GenerationResult {
623    #[cfg(any(test, feature = "fault-injection"))]
624    if fault::CHILD_PANIC.swap(false, Ordering::SeqCst) {
625        panic!("adr133 fault injection: audit_batch child_panic");
626    }
627    #[cfg(any(test, feature = "fault-injection"))]
628    if fault::INCONSISTENT_EXIT.swap(false, Ordering::SeqCst) {
629        return GenerationResult::FakedInconsistent;
630    }
631
632    let mut attempt: u8 = 0;
633    loop {
634        attempt += 1;
635        match store.append_events_idempotent(events.clone()).await {
636            Ok(result) => return GenerationResult::Committed(result.rows),
637            Err(err) => {
638                let reason = match classify_store_error(&err) {
639                    RetryDecision::Retry if attempt < config.max_commit_attempts.get() => {
640                        tokio::time::sleep(config.retry_backoff).await;
641                        continue;
642                    }
643                    RetryDecision::Retry => AuditTerminalReason::RetryExhausted,
644                    RetryDecision::Terminal(reason) => reason,
645                };
646                tracing::warn!(
647                    error = %err,
648                    attempts = attempt,
649                    ?reason,
650                    "audit generation failed; its rows were not committed"
651                );
652                return GenerationResult::Failed(reason);
653            }
654        }
655    }
656}
657
658/// The batch owner. Constructed once per configured `EventStore`; every
659/// dispatch-audit call site routes its row through [`AuditBatch::submit`]
660/// instead of taking its own writer-task acquisition.
661pub struct AuditBatch {
662    inner: Arc<Inner>,
663    store: Arc<dyn EventStore>,
664    config: Arc<AuditBatchConfig>,
665    supervisor: Mutex<Option<JoinHandle<()>>>,
666}
667
668impl AuditBatch {
669    pub fn new(store: Arc<dyn EventStore>, config: AuditBatchConfig) -> Arc<Self> {
670        debug_assert!(
671            config.resolution_deadline >= config.admission_deadline,
672            "resolution_deadline ({:?}) must be at least admission_deadline ({:?})",
673            config.resolution_deadline,
674            config.admission_deadline
675        );
676        Arc::new(Self {
677            inner: Arc::new(Inner {
678                state: Mutex::new(State::new()),
679            }),
680            store,
681            config: Arc::new(config),
682            supervisor: Mutex::new(None),
683        })
684    }
685
686    /// Process-lifetime audit-batch health counters, for the registry/
687    /// runtime owner to feed into `db_diagnostics` (D8's operator surface).
688    /// Unlike `test_internals::AuditBatchSnapshot::metrics_snapshot`,
689    /// which is test-only (the module is cfg-gated and invisible to the
690    /// default doc build, so an intra-doc link cannot resolve), this is
691    /// always available.
692    pub fn health_metrics(&self) -> AuditBatchHealthMetrics {
693        let state = self.inner.state.lock();
694        AuditBatchHealthMetrics {
695            flush_failures: state.flush_failures,
696            degraded_rows: state.degraded_rows,
697            degraded: state.degraded,
698            late_append_commits: state.late_append_commits,
699            late_append_failures: state.late_append_failures,
700        }
701    }
702
703    fn spawn_supervisor_if_idle(&self) {
704        let inner = self.inner.clone();
705        let store = self.store.clone();
706        let config = self.config.clone();
707        let handle = tokio::spawn(async move {
708            supervisor_loop(inner, store, config).await;
709        });
710        *self.supervisor.lock() = Some(handle);
711    }
712
713    /// Enqueue one row and keep waiting for its real generation outcome past
714    /// the ordinary admission wait deadline, up to `resolution_deadline`.
715    ///
716    /// This is the narrow khive#2256 seam for a successful operation whose
717    /// domain effect has already committed. Returning
718    /// [`AuditTerminalReason::AdmissionDeadlineExpired`] there would report a
719    /// false operation failure and invite an unsafe retry while the same
720    /// audit row remains enqueued. Pre-enqueue refusal and genuine terminal
721    /// generation failures still return normally. The post-admission wait is
722    /// itself bounded by `resolution_deadline`
723    /// ([`AuditTerminalReason::ResolutionDeadlineExpired`]) so a stalled
724    /// store cannot retain this caller, its request slot, and its audit-lane
725    /// waiter forever (khive#2331).
726    ///
727    /// The two deadlines therefore mean different things to the caller, and
728    /// callers key on the difference. Admission expiry means the row is
729    /// enqueued and this seam keeps waiting; the generation commits it
730    /// independently, so a dispatch that reaches it reports its committed
731    /// result and counts the row as unresolved. Resolution expiry means the
732    /// caller waited for the real outcome and never received one: the commit
733    /// is unconfirmed, not proven absent. That is returned as a terminal
734    /// reason and the dispatch propagates it as a structured
735    /// committed-outcome error carrying the domain result, with no retryable
736    /// context, for every verb rather than only for the receipt verb. Never
737    /// replay the handler or re-enqueue the row on it; the original driver
738    /// work continues on its own.
739    pub(crate) async fn submit_until_resolved(
740        &self,
741        row: PreparedAuditRow,
742    ) -> Result<AuditCommitOutcome, AuditTerminalReason> {
743        self.submit_with_wait_policy(row, true).await
744    }
745
746    async fn submit_with_wait_policy(
747        &self,
748        row: PreparedAuditRow,
749        wait_until_resolved: bool,
750    ) -> Result<AuditCommitOutcome, AuditTerminalReason> {
751        // Pre-enqueue validation (invariant 3): a malformed row is rejected
752        // before it can share a generation with anyone else's.
753        if self.store.preflight_event(&row.event).is_err() {
754            return Err(AuditTerminalReason::PreflightRejected);
755        }
756
757        let producer = row.producer;
758        let (tx, mut rx) = oneshot::channel();
759        let need_spawn = {
760            let mut state = self.inner.state.lock();
761            match state.lifecycle {
762                Lifecycle::Closed | Lifecycle::Closing => {
763                    return Err(AuditTerminalReason::AdmissionClosed)
764                }
765                Lifecycle::Failed(reason) => return Err(reason),
766                Lifecycle::Open => {}
767            }
768            if state.pending.len() >= self.config.max_pending_rows.get() {
769                return Err(AuditTerminalReason::QueueAdmissionExhausted);
770            }
771            state.pending.push(Waiting {
772                event: row.event,
773                producer,
774                responder: tx,
775            });
776            state.submitted_rows += 1;
777            let need_spawn = !state.driver_active;
778            if need_spawn {
779                state.driver_active = true;
780            }
781            need_spawn
782        };
783        if need_spawn {
784            self.spawn_supervisor_if_idle();
785        }
786
787        if wait_until_resolved {
788            match tokio::time::timeout(self.config.admission_deadline, &mut rx).await {
789                Ok(Ok(result)) => return result,
790                Ok(Err(_recv_error)) => return Err(AuditTerminalReason::DriverJoinLost),
791                Err(_elapsed) => tracing::warn!(
792                    ?producer,
793                    "strict audit obligation remains enqueued after the admission wait deadline; \
794                     waiting up to the resolution deadline for its real terminal outcome"
795                ),
796            }
797            return match tokio::time::timeout(self.config.resolution_deadline, rx).await {
798                Ok(Ok(result)) => result,
799                Ok(Err(_recv_error)) => Err(AuditTerminalReason::DriverJoinLost),
800                // Same non-removal contract as the `AdmissionDeadlineExpired`
801                // arm below: the row is left exactly where the driver holds
802                // it. The caller's domain effect already committed by the
803                // time it reached this wait, so it is never retried here;
804                // the audit outcome itself is what remains unresolved.
805                Err(_elapsed) => {
806                    tracing::warn!(
807                        ?producer,
808                        "strict audit obligation remains unresolved after the resolution \
809                         deadline; giving up on this caller's wait — the committed effect is \
810                         not retried and the row is not re-enqueued"
811                    );
812                    Err(AuditTerminalReason::ResolutionDeadlineExpired)
813                }
814            };
815        }
816
817        match tokio::time::timeout(self.config.admission_deadline, rx).await {
818            Ok(Ok(result)) => result,
819            Ok(Err(_recv_error)) => Err(AuditTerminalReason::DriverJoinLost),
820            // The row was already pushed onto `state.pending` above (and
821            // `submitted_rows` incremented) before this wait began — this is
822            // a deadline elapsing on an enqueued row, not a queue-full
823            // refusal, so it gets its own terminal reason (khive#2117,
824            // khive#2208). The row is left in place for the driver to drain;
825            // this arm performs no removal.
826            Err(_elapsed) => Err(AuditTerminalReason::AdmissionDeadlineExpired),
827        }
828    }
829}
830
831async fn supervisor_loop(
832    inner: Arc<Inner>,
833    store: Arc<dyn EventStore>,
834    config: Arc<AuditBatchConfig>,
835) {
836    loop {
837        let (waiting, wedged_generation_id) = {
838            let mut state = inner.state.lock();
839            if matches!(state.lifecycle, Lifecycle::Failed(_)) {
840                state.driver_active = false;
841                break;
842            }
843            if state.pending.is_empty() {
844                state.driver_active = false;
845                break;
846            }
847            let take_n = state
848                .pending
849                .len()
850                .min(config.max_rows_per_generation.get());
851            let waiting: Vec<Waiting> = state.pending.drain(..take_n).collect();
852            let generation_id = state.next_generation_id;
853            state.next_generation_id += 1;
854            // khive#2331: a store whose append never returns must not mint
855            // an unbounded number of detached appends. Before committing to
856            // spawning this generation's child, check whether the cap is
857            // already saturated by prior abandoned generations still
858            // running in the background — if so, this generation is shed
859            // below instead of ever calling the store.
860            if state.outstanding_abandoned_appends >= config.max_abandoned_appends.get() {
861                (waiting, Some(generation_id))
862            } else {
863                state.in_flight_generation = Some(generation_id);
864                (waiting, None)
865            }
866        };
867
868        if let Some(generation_id) = wedged_generation_id {
869            tracing::warn!(
870                generation_id,
871                max_abandoned_appends = config.max_abandoned_appends.get(),
872                "audit store treated as wedged: max_abandoned_appends detached appends are \
873                 already outstanding; shedding this generation without attempting an append"
874            );
875            let submitted = waiting.len() as u64;
876            {
877                let mut state = inner.state.lock();
878                state.flush_failures += 1;
879                state.generations.push(AuditGenerationSnapshot {
880                    generation_id,
881                    submitted_rows: submitted,
882                    committed_rows: 0,
883                    store_batch_calls: 0,
884                    terminal_reason: Some(AuditTerminalReason::StoreWedged),
885                });
886            }
887            for w in waiting {
888                record_degradation_if_pure(&inner, w.producer);
889                let _ = w.responder.send(Err(AuditTerminalReason::StoreWedged));
890            }
891            continue;
892        }
893
894        #[cfg(any(test, feature = "fault-injection"))]
895        if fault::SUPERVISOR_PANIC.swap(false, Ordering::SeqCst) {
896            let _guard = SupervisorGuard::armed(&inner, waiting);
897            panic!("adr133 fault injection: audit_batch supervisor_panic");
898        }
899        #[cfg(any(test, feature = "fault-injection"))]
900        if fault::SUPERVISOR_SLEEP_BEFORE_SPAWN.swap(false, Ordering::SeqCst) {
901            let guard = SupervisorGuard::armed(&inner, waiting);
902            tokio::time::sleep(Duration::from_secs(3600)).await;
903            drop(guard);
904            continue;
905        }
906
907        let guard = SupervisorGuard::armed(&inner, waiting);
908        let events: Vec<Event> = guard
909            .waiting
910            .as_ref()
911            .expect("guard freshly armed")
912            .iter()
913            .map(|w| w.event.clone())
914            .collect();
915
916        let mut child: JoinHandle<GenerationResult> =
917            tokio::spawn(run_generation(store.clone(), events, config.clone()));
918
919        #[cfg(any(test, feature = "fault-injection"))]
920        if fault::CHILD_CANCEL.swap(false, Ordering::SeqCst) {
921            child.abort();
922        }
923        #[cfg(any(test, feature = "fault-injection"))]
924        if fault::JOIN_LOST.swap(false, Ordering::SeqCst) {
925            drop(child);
926            guard.fail(AuditTerminalReason::DriverJoinLost);
927            continue;
928        }
929
930        let append_deadline = driver_append_deadline(&config);
931        let join_result = match tokio::time::timeout(append_deadline, &mut child).await {
932            Ok(join_result) => join_result,
933            Err(_elapsed) => {
934                // The append is still in flight and may not be safely
935                // cancellable mid-write (a blocking storage call cannot be
936                // forced to stop), so `child` is not aborted here — it is
937                // handed to a detached reaper that drives it to completion
938                // and discards whatever it eventually returns. Every waiter
939                // on this generation is resolved now, and the driver loops
940                // back to `pending` immediately instead of staying pinned on
941                // this one stalled call, which is what actually starves
942                // later admission (khive#2331).
943                tracing::warn!(
944                    ?append_deadline,
945                    "audit generation append exceeded the driver's own bound; abandoning \
946                     this generation so admission capacity recovers for later callers"
947                );
948                let waiting = guard.disarm();
949                let submitted = waiting.len() as u64;
950                {
951                    let mut state = inner.state.lock();
952                    state.in_flight_generation = None;
953                    state.flush_failures += 1;
954                    state.outstanding_abandoned_appends += 1;
955                    let generation_id = state.next_generation_id.saturating_sub(1);
956                    state.generations.push(AuditGenerationSnapshot {
957                        generation_id,
958                        submitted_rows: submitted,
959                        committed_rows: 0,
960                        store_batch_calls: 0,
961                        terminal_reason: Some(AuditTerminalReason::DriverAppendAbandoned),
962                    });
963                }
964                for w in waiting {
965                    record_degradation_if_pure(&inner, w.producer);
966                    let _ = w
967                        .responder
968                        .send(Err(AuditTerminalReason::DriverAppendAbandoned));
969                }
970                let reaper_inner = inner.clone();
971                tokio::spawn(async move {
972                    let result = child.await;
973                    let mut state = reaper_inner.state.lock();
974                    state.outstanding_abandoned_appends =
975                        state.outstanding_abandoned_appends.saturating_sub(1);
976                    match result {
977                        Ok(GenerationResult::Committed(_)) => state.late_append_commits += 1,
978                        Ok(GenerationResult::Failed(_) | GenerationResult::FakedInconsistent)
979                        | Err(_) => state.late_append_failures += 1,
980                    }
981                });
982                continue;
983            }
984        };
985        match join_result {
986            Err(join_err) => {
987                let reason = if join_err.is_panic() {
988                    AuditTerminalReason::DriverPanicked
989                } else {
990                    AuditTerminalReason::DriverCancelled
991                };
992                guard.fail(reason);
993            }
994            Ok(GenerationResult::FakedInconsistent) => {
995                guard.fail(AuditTerminalReason::DriverExitedInconsistent);
996            }
997            Ok(GenerationResult::Failed(reason)) => {
998                let waiting = guard.disarm();
999                let submitted = waiting.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.flush_failures += 1;
1005                    state.degraded = true;
1006                    let generation_id = state.next_generation_id.saturating_sub(1);
1007                    state.generations.push(AuditGenerationSnapshot {
1008                        generation_id,
1009                        submitted_rows: submitted,
1010                        committed_rows: 0,
1011                        store_batch_calls: 1,
1012                        terminal_reason: Some(reason),
1013                    });
1014                }
1015                for w in waiting {
1016                    record_degradation_if_pure(&inner, w.producer);
1017                    let _ = w.responder.send(Err(reason));
1018                }
1019            }
1020            Ok(GenerationResult::Committed(dispositions)) => {
1021                let waiting = guard.disarm();
1022                if dispositions.len() != waiting.len() {
1023                    // Defensive: the store must preserve input order/length.
1024                    // Treat a mismatch as the driver having exited in a
1025                    // state that cannot be reconciled with the accepted
1026                    // waiters.
1027                    {
1028                        let mut state = inner.state.lock();
1029                        state.in_flight_generation = None;
1030                    }
1031                    inner.fail_driver(AuditTerminalReason::DriverExitedInconsistent, waiting);
1032                    continue;
1033                }
1034                let mut committed_n = 0u64;
1035                for d in &dispositions {
1036                    if !matches!(d, EventAppendDisposition::IdentityConflict) {
1037                        committed_n += 1;
1038                    }
1039                }
1040                let submitted = dispositions.len() as u64;
1041                {
1042                    let mut state = inner.state.lock();
1043                    state.in_flight_generation = None;
1044                    state.store_batch_calls += 1;
1045                    state.committed_rows += committed_n;
1046                    let generation_id = state.next_generation_id.saturating_sub(1);
1047                    state.generations.push(AuditGenerationSnapshot {
1048                        generation_id,
1049                        submitted_rows: submitted,
1050                        committed_rows: committed_n,
1051                        store_batch_calls: 1,
1052                        terminal_reason: None,
1053                    });
1054                }
1055                for (w, disposition) in waiting.into_iter().zip(dispositions) {
1056                    let result = match disposition {
1057                        EventAppendDisposition::Inserted => Ok(AuditCommitOutcome::Committed),
1058                        EventAppendDisposition::AlreadyPresentIdentical => {
1059                            Ok(AuditCommitOutcome::AlreadyPresentIdentical)
1060                        }
1061                        EventAppendDisposition::IdentityConflict => {
1062                            record_degradation_if_pure(&inner, w.producer);
1063                            Err(AuditTerminalReason::IdentityConflict)
1064                        }
1065                    };
1066                    let _ = w.responder.send(result);
1067                }
1068            }
1069        }
1070    }
1071}
1072
1073#[async_trait::async_trait]
1074impl AuditBatchControl for AuditBatch {
1075    async fn submit(
1076        &self,
1077        row: PreparedAuditRow,
1078    ) -> Result<AuditCommitOutcome, AuditTerminalReason> {
1079        self.submit_with_wait_policy(row, false).await
1080    }
1081
1082    async fn quiesce(&self) -> Result<(), AuditTerminalReason> {
1083        loop {
1084            let (lifecycle, idle) = {
1085                let state = self.inner.state.lock();
1086                (state.lifecycle.clone(), state.is_idle())
1087            };
1088            match lifecycle {
1089                Lifecycle::Failed(reason) => return Err(reason),
1090                Lifecycle::Closed => return Ok(()),
1091                _ if idle => return Ok(()),
1092                _ => tokio::time::sleep(Duration::from_millis(2)).await,
1093            }
1094        }
1095    }
1096
1097    async fn close_and_drain(&self) -> Result<(), AuditTerminalReason> {
1098        {
1099            let mut state = self.inner.state.lock();
1100            if matches!(state.lifecycle, Lifecycle::Open) {
1101                state.lifecycle = Lifecycle::Closing;
1102            }
1103        }
1104        let result = AuditBatchControl::quiesce(self).await;
1105        {
1106            let mut state = self.inner.state.lock();
1107            if result.is_ok() && matches!(state.lifecycle, Lifecycle::Closing) {
1108                state.lifecycle = Lifecycle::Closed;
1109            }
1110        }
1111        let handle = self.supervisor.lock().take();
1112        if let Some(handle) = handle {
1113            let _ = handle.await;
1114        }
1115        result
1116    }
1117}
1118
1119#[cfg(any(test, feature = "fault-injection"))]
1120mod fault {
1121    use std::sync::atomic::AtomicBool;
1122
1123    pub(super) static CHILD_PANIC: AtomicBool = AtomicBool::new(false);
1124    pub(super) static CHILD_CANCEL: AtomicBool = AtomicBool::new(false);
1125    pub(super) static SUPERVISOR_PANIC: AtomicBool = AtomicBool::new(false);
1126    pub(super) static SUPERVISOR_SLEEP_BEFORE_SPAWN: AtomicBool = AtomicBool::new(false);
1127    pub(super) static JOIN_LOST: AtomicBool = AtomicBool::new(false);
1128    pub(super) static INCONSISTENT_EXIT: AtomicBool = AtomicBool::new(false);
1129}
1130
1131/// Deterministic fault-injection arms consumed by exactly one subsequent
1132/// generation each. Test-only surface, gated identically to the rest of this
1133/// crate's `fault-injection` fixtures (see `crate::operations`).
1134#[cfg(any(test, feature = "fault-injection"))]
1135pub mod fault_injection {
1136    use std::sync::atomic::Ordering;
1137
1138    pub fn arm_child_panic() {
1139        super::fault::CHILD_PANIC.store(true, Ordering::SeqCst);
1140    }
1141    pub fn arm_child_cancel() {
1142        super::fault::CHILD_CANCEL.store(true, Ordering::SeqCst);
1143    }
1144    pub fn arm_supervisor_panic() {
1145        super::fault::SUPERVISOR_PANIC.store(true, Ordering::SeqCst);
1146    }
1147    pub fn arm_supervisor_sleep_before_spawn() {
1148        super::fault::SUPERVISOR_SLEEP_BEFORE_SPAWN.store(true, Ordering::SeqCst);
1149    }
1150    pub fn arm_join_lost() {
1151        super::fault::JOIN_LOST.store(true, Ordering::SeqCst);
1152    }
1153    pub fn arm_inconsistent_exit() {
1154        super::fault::INCONSISTENT_EXIT.store(true, Ordering::SeqCst);
1155    }
1156}
1157
1158// ── Test-attribution surface (R2) ────────────────────────────────────────
1159//
1160// Additive only: production attribution still lands in the public
1161// `db_diagnostics` counters via `metrics_snapshot`-shaped data supplied at
1162// that seam. This surface never substitutes for it.
1163#[cfg(any(test, feature = "test-internals"))]
1164mod test_internals {
1165    use super::*;
1166
1167    /// A point-in-time view of the batch's counters and generation history,
1168    /// used by mechanism tests to compute a delta across a measured
1169    /// operation.
1170    #[derive(Debug, Clone)]
1171    pub struct AuditBatchSnapshot {
1172        pub pending_rows: usize,
1173        pub in_flight_generation: Option<u64>,
1174        pub driver_active: bool,
1175        pub next_generation_id: u64,
1176        pub submitted_rows: u64,
1177        pub committed_rows: u64,
1178        pub store_batch_calls: u64,
1179        pub per_generation: Vec<AuditGenerationSnapshot>,
1180        /// Live count of detached `DriverAppendAbandoned` appends still
1181        /// running, bounded by `AuditBatchConfig::max_abandoned_appends`.
1182        pub outstanding_abandoned_appends: usize,
1183    }
1184
1185    impl AuditBatchSnapshot {
1186        pub fn is_idle(&self) -> bool {
1187            self.pending_rows == 0 && self.in_flight_generation.is_none() && !self.driver_active
1188        }
1189    }
1190
1191    #[derive(Debug, Clone, Copy)]
1192    pub struct AuditBatchMetricsSnapshot {
1193        pub flush_failures: u64,
1194        pub degraded_rows: u64,
1195        pub degraded: bool,
1196        pub late_append_commits: u64,
1197        pub late_append_failures: u64,
1198    }
1199
1200    #[derive(Debug, Clone, PartialEq)]
1201    pub struct AuditBatchDelta {
1202        pub submitted_rows: u64,
1203        pub committed_rows: u64,
1204        pub store_batch_calls: u64,
1205        pub per_generation: Vec<AuditGenerationSnapshot>,
1206    }
1207
1208    #[derive(Debug, Clone, Copy, PartialEq, Eq)]
1209    pub enum AuditSnapshotError {
1210        CounterRegressed,
1211        GenerationHistoryRegressed,
1212    }
1213
1214    impl AuditBatch {
1215        pub fn test_snapshot(&self) -> AuditBatchSnapshot {
1216            let state = self.inner.state.lock();
1217            AuditBatchSnapshot {
1218                pending_rows: state.pending.len(),
1219                in_flight_generation: state.in_flight_generation,
1220                driver_active: state.driver_active,
1221                next_generation_id: state.next_generation_id,
1222                submitted_rows: state.submitted_rows,
1223                committed_rows: state.committed_rows,
1224                store_batch_calls: state.store_batch_calls,
1225                per_generation: state.generations.clone(),
1226                outstanding_abandoned_appends: state.outstanding_abandoned_appends,
1227            }
1228        }
1229
1230        pub fn metrics_snapshot(&self) -> AuditBatchMetricsSnapshot {
1231            let state = self.inner.state.lock();
1232            AuditBatchMetricsSnapshot {
1233                flush_failures: state.flush_failures,
1234                degraded_rows: state.degraded_rows,
1235                degraded: state.degraded,
1236                late_append_commits: state.late_append_commits,
1237                late_append_failures: state.late_append_failures,
1238            }
1239        }
1240
1241        /// Abort the currently-retained supervisor `JoinHandle`, if any,
1242        /// simulating a shutdown abort landing mid-generation. Returns
1243        /// whether a handle was found and aborted. Test-only: exercises the
1244        /// R1 supervisor-cancellation path without reaching into a private
1245        /// field from outside this module.
1246        pub fn test_abort_supervisor(&self) -> bool {
1247            let handle = self.supervisor.lock().take();
1248            match handle {
1249                Some(handle) => {
1250                    handle.abort();
1251                    true
1252                }
1253                None => false,
1254            }
1255        }
1256    }
1257
1258    /// Checked monotonic subtraction. Rejects a regressed counter or a
1259    /// generation history that is not an append-only extension of `before`.
1260    pub fn audit_delta(
1261        before: &AuditBatchSnapshot,
1262        after: &AuditBatchSnapshot,
1263    ) -> Result<AuditBatchDelta, AuditSnapshotError> {
1264        let submitted_rows = after
1265            .submitted_rows
1266            .checked_sub(before.submitted_rows)
1267            .ok_or(AuditSnapshotError::CounterRegressed)?;
1268        let committed_rows = after
1269            .committed_rows
1270            .checked_sub(before.committed_rows)
1271            .ok_or(AuditSnapshotError::CounterRegressed)?;
1272        let store_batch_calls = after
1273            .store_batch_calls
1274            .checked_sub(before.store_batch_calls)
1275            .ok_or(AuditSnapshotError::CounterRegressed)?;
1276        if after.per_generation.len() < before.per_generation.len() {
1277            return Err(AuditSnapshotError::GenerationHistoryRegressed);
1278        }
1279        if after.per_generation[..before.per_generation.len()] != before.per_generation[..] {
1280            return Err(AuditSnapshotError::GenerationHistoryRegressed);
1281        }
1282        let per_generation = after.per_generation[before.per_generation.len()..].to_vec();
1283        Ok(AuditBatchDelta {
1284            submitted_rows,
1285            committed_rows,
1286            store_batch_calls,
1287            per_generation,
1288        })
1289    }
1290
1291    /// Exhaustive, non-wildcard producer classification (see
1292    /// [`super::AuditProducer`] and [`super::classify`], the real
1293    /// crate-private definitions this mirrors one-for-one). The doctest
1294    /// below proves the general property those definitions rely on: a match
1295    /// over a non-`#[non_exhaustive]` enum that omits a variant, with no
1296    /// wildcard arm to silently absorb it, fails to compile rather than
1297    /// passing an incomplete classification.
1298    ///
1299    /// ```compile_fail
1300    /// enum AuditProducer {
1301    ///     GateDenied,
1302    ///     DispatchSucceeded,
1303    ///     DispatchFailed,
1304    ///     UnknownVerb,
1305    ///     GitDigestReceipt,
1306    ///     ConfigLocked,
1307    ///     RecallExecuted,
1308    /// }
1309    ///
1310    /// fn describe(p: AuditProducer) -> &'static str {
1311    ///     match p {
1312    ///         AuditProducer::GateDenied => "obligation",
1313    ///         AuditProducer::DispatchSucceeded => "obligation",
1314    ///         AuditProducer::DispatchFailed => "obligation",
1315    ///         AuditProducer::UnknownVerb => "obligation",
1316    ///         AuditProducer::GitDigestReceipt => "obligation",
1317    ///         AuditProducer::ConfigLocked => "observability",
1318    ///         // RecallExecuted intentionally omitted: a non-exhaustive
1319    ///         // match must fail to compile, proving no variant can
1320    ///         // silently fall through an absent wildcard arm.
1321    ///     }
1322    /// }
1323    /// ```
1324    #[allow(dead_code)]
1325    struct DoctestAnchor;
1326}
1327
1328#[cfg(any(test, feature = "test-internals"))]
1329pub use test_internals::*;