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::*;