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