Skip to main content

khive_db/
writer_task.rs

1//! Single-writer task and bounded write queue (ADR-067 Component A).
2//!
3//! `WriterTask` (via `spawn` and the drain loop `run_writer_task`) owns a
4//! dedicated standalone writer `rusqlite::Connection` and is the only code
5//! path that issues `BEGIN IMMEDIATE` for write traffic routed through the
6//! channel it drains. Callers reach it exclusively through a
7//! [`WriterTaskHandle`], sending a typed closure and awaiting a typed
8//! oneshot reply so each store method's natural return type (e.g.
9//! `BatchWriteSummary`) survives the trip through the type-erased channel
10//! unmodified — a flat `Result<u64, StorageError>` reply would conflate
11//! `affected`/`failed` into one count and drop `first_error`.
12//!
13//! See `crates/khive-db/docs/api/writer-task.md` for migration-slice scope
14//! (which write paths currently route through this vs. the legacy
15//! pool-mutex path) and the ADR-067 component breakdown.
16//!
17//! ## ADR-136 D1 gate 5: writer classification
18//!
19//! Every connection that issues writes against a khive-db file falls into
20//! exactly one of the rows below. Request-path store writes, runtime
21//! merge/symmetric-edge transactions,
22//! and `khive_storage::SqlAccess` (`SqlBridge::writer`/`atomic_unit`) are subject
23//! to `PoolConfig::write_queue_enabled` / `write_routing_strict` routing.
24//! Everything else here is an
25//! explicit, intentional exemption — not an implicit gap left over from an
26//! incomplete migration.
27//!
28//! | Writer | Connection | Request-path write | Routing |
29//! | --- | --- | --- | --- |
30//! | Request-path writes (`create`, `update`, batch upserts, …) | `WriterTaskHandle` (queue-first) or a standalone/pool-mutex connection on degrade | Yes | Queue-first (ADR-136 D1 gate 1); strict routing fails closed on degrade |
31//! | Startup / schema migrations (`migrations::run_migrations`, `apply_schema_plan`) | The pool's own writer-mutex connection (`ConnectionPool::new`, before the `WriterTask` is ever spawned) | No | Exempt — runs once at boot, before any queue handle exists to route through |
32//! | Checkpointing (`checkpoint::CheckpointConnection`) | Dedicated standalone connection, opened once at task startup (ADR-091 Amendment 5, #1652) | No | Exempt by design — the whole point of that fix was moving checkpoint I/O OFF the pool writer's admission path; routing it back through the write queue would reintroduce the contention it removed |
33//! | Recovery (`walpin` beacon/sidecar bookkeeping) | None — file-level bookkeeping (PID beacons, heartbeats) alongside the database, not a SQL connection against it | No | N/A — never acquires a database writer connection |
34//! | Top-level maintenance (`VACUUM` via `execute_script_top_level`/`WriterTaskHandle::send_top_level`) | `WriterTaskHandle`, skipping the per-request `BEGIN IMMEDIATE`/`COMMIT` wrap only | Yes | Routed through the SAME single writer owner as every other queued request — only the transaction wrap is skipped, never the queue |
35//!
36//! Store constructors may run outside Tokio and temporarily cache no handle;
37//! every store request-path write refreshes through
38//! `ConnectionPool::writer_task_for_write` at execution time. Strict routing
39//! refuses a missing handle, while the explicit compatibility fallback emits
40//! its store-specific violation at the actual direct-writer seam. Runtime
41//! entity/note merges and symmetric edge updates use the same routing policy
42//! through a closed operation adapter. The strict default itself remains gated
43//! by ADR-135 F2 / ADR-136 D2 evidence.
44//!
45//! A `direct_route_violation` sink row (`timeout_sink::emit_direct_route_violation`,
46//! ADR-136 D1 gate 6c) is emitted ONLY for the first row's degrade path — a
47//! request-path write that bypassed an enabled queue. The other four
48//! rows are exempt by design and never emit that row.
49
50use rusqlite::Connection;
51use std::collections::HashMap;
52use std::panic::{catch_unwind, AssertUnwindSafe};
53use std::path::{Path, PathBuf};
54use std::sync::{Arc, Mutex, OnceLock};
55use std::time::{Duration, Instant};
56use tokio::sync::{mpsc, oneshot};
57
58use khive_storage::error::{StorageError, WriterTaskRequestState};
59
60use crate::error::SqliteError;
61use crate::pool::{ConnectionPool, WriterAcquisitionCounters};
62
63/// Retry delays after a busy/locked `BEGIN IMMEDIATE` refusal. The operation
64/// closure remains owned by the request until one BEGIN succeeds, so these
65/// are the only retries this path may perform. Persistent contention is
66/// bounded to three total BEGIN attempts and 15 ms of explicit backoff, and
67/// the attempts together share one configured `busy_timeout` acquisition
68/// budget rather than each waiting out a full window of their own.
69const WRITER_BEGIN_RETRY_DELAYS: [Duration; 2] =
70    [Duration::from_millis(5), Duration::from_millis(10)];
71
72/// Closure signature for a write operation executed against the writer
73/// task's dedicated connection.
74///
75/// `conn` is already inside a `BEGIN IMMEDIATE` transaction opened by
76/// `run_writer_task` when this runs. The closure must issue DML (and, in
77/// later slices, named `SAVEPOINT`s) only — never a bare `BEGIN` / `COMMIT`
78/// / `ROLLBACK` — a nested bare `BEGIN IMMEDIATE` would violate SQLite's
79/// nested-transaction rule and return `SQLITE_ERROR: cannot start a
80/// transaction within a transaction` (ADR-067 lines 271-276).
81type WriteOp<R> = Box<dyn FnOnce(&Connection) -> Result<R, StorageError> + Send>;
82
83/// Latest completed writer-task span for one database, decomposed at the
84/// boundaries an operator can act on (#1849). A zero value means that phase
85/// did not run (top-level/early-failure request) or completed below the
86/// microsecond clock resolution; `total_micros` retains the full span.
87#[derive(Debug, Clone, PartialEq, Eq)]
88pub struct WriterStageObservation {
89    pub queue_wait_micros: u64,
90    pub transaction_acquire_micros: u64,
91    pub body_micros: u64,
92    pub commit_micros: u64,
93    pub total_micros: u64,
94    pub queue_depth_at_entry: u64,
95    pub observed_at_unix_ms: u64,
96}
97
98static WRITER_STAGE_OBSERVATIONS: OnceLock<
99    Mutex<HashMap<Option<PathBuf>, WriterStageObservation>>,
100> = OnceLock::new();
101
102fn writer_stage_observations() -> &'static Mutex<HashMap<Option<PathBuf>, WriterStageObservation>> {
103    WRITER_STAGE_OBSERVATIONS.get_or_init(|| Mutex::new(HashMap::new()))
104}
105
106fn writer_db_key_from_path(path: Option<&Path>) -> Option<PathBuf> {
107    path.map(Path::to_path_buf)
108}
109
110fn writer_db_key(pool: &ConnectionPool) -> Option<PathBuf> {
111    writer_db_key_from_path(pool.canonical_path())
112}
113
114fn duration_micros(duration: Duration) -> u64 {
115    duration.as_micros().min(u128::from(u64::MAX)) as u64
116}
117
118fn observed_at_unix_ms() -> u64 {
119    std::time::SystemTime::now()
120        .duration_since(std::time::UNIX_EPOCH)
121        .map(|duration| duration.as_millis() as u64)
122        .unwrap_or(0)
123}
124
125/// Pure in-memory read of the most recently completed writer-task span for
126/// this exact backend. It acquires no SQLite connection and performs no I/O.
127pub fn last_writer_stage_observation(pool: &ConnectionPool) -> Option<WriterStageObservation> {
128    writer_stage_observations()
129        .lock()
130        .unwrap_or_else(std::sync::PoisonError::into_inner)
131        .get(&writer_db_key(pool))
132        .cloned()
133}
134
135struct WriteTelemetry {
136    backend_key: Option<PathBuf>,
137    db: String,
138    submitted_at: Instant,
139    queue_depth_at_entry: usize,
140    slow_write_threshold: Option<Duration>,
141}
142
143impl WriteTelemetry {
144    fn new(
145        backend_key: Option<PathBuf>,
146        db: String,
147        queue_depth_at_entry: usize,
148        slow_write_threshold: Option<Duration>,
149    ) -> Self {
150        Self {
151            backend_key,
152            db,
153            submitted_at: Instant::now(),
154            queue_depth_at_entry,
155            slow_write_threshold,
156        }
157    }
158
159    fn queue_wait(&self) -> Duration {
160        self.submitted_at.elapsed()
161    }
162
163    fn finish(
164        self,
165        queue_wait: Duration,
166        transaction_acquire: Duration,
167        body: Duration,
168        commit: Duration,
169    ) {
170        let total = self.submitted_at.elapsed();
171        let observation = WriterStageObservation {
172            queue_wait_micros: duration_micros(queue_wait),
173            transaction_acquire_micros: duration_micros(transaction_acquire),
174            body_micros: duration_micros(body),
175            commit_micros: duration_micros(commit),
176            total_micros: duration_micros(total),
177            queue_depth_at_entry: self.queue_depth_at_entry as u64,
178            observed_at_unix_ms: observed_at_unix_ms(),
179        };
180        writer_stage_observations()
181            .lock()
182            .unwrap_or_else(std::sync::PoisonError::into_inner)
183            .insert(self.backend_key, observation.clone());
184
185        if self
186            .slow_write_threshold
187            .is_some_and(|threshold| total >= threshold)
188        {
189            crate::timeout_sink::emit_slow_write(&self.db, &observation);
190        }
191    }
192}
193
194/// One write request awaiting execution by the writer task.
195///
196/// Carries a typed closure and a typed oneshot reply so that the concrete
197/// return type `R` (e.g. `BatchWriteSummary`) is preserved end to end,
198/// while [`AnyWriteRequest`] lets the drain loop hold heterogeneous
199/// requests in one homogeneous channel.
200///
201/// `top_level` (ADR-067 Component A): when `true`,
202/// the drain loop runs this request's operation WITHOUT wrapping it in a
203/// `BEGIN IMMEDIATE`/`COMMIT`/`ROLLBACK` — still serialized through the
204/// single writer owner (only one request drains at a time regardless of
205/// this flag), but with the transaction wrap skipped entirely. Exists for
206/// statements SQLite forbids inside any open transaction (e.g. `VACUUM`);
207/// see [`WriterTaskHandle::send_top_level`].
208pub struct WriteRequest<R: Send + 'static> {
209    op: WriteOp<R>,
210    reply: oneshot::Sender<Result<R, StorageError>>,
211    top_level: bool,
212    /// Only the typed checkpoint entry point may bypass disk admission.
213    checkpoint_bypass: bool,
214    vacuum_copy_headroom: bool,
215    telemetry: WriteTelemetry,
216}
217
218mod sealed {
219    /// Restricts [`super::AnyWriteRequest`] to implementations defined in
220    /// this module — only [`super::WriteRequest<R>`] implements it — and
221    /// carries the drain loop's internal terminal-state reporting methods
222    /// without changing the public execution-method signatures.
223    pub trait Sealed {
224        fn execute_and_reply_reporting_terminal(
225            self: Box<Self>,
226            conn: &rusqlite::Connection,
227            tx_span: Option<khive_storage::tx_registry::TxHandle>,
228            queue_wait: std::time::Duration,
229            transaction_acquire: std::time::Duration,
230        ) -> Option<khive_storage::error::WriterTaskRequestState>;
231
232        fn execute_and_reply_top_level_reporting_terminal(
233            self: Box<Self>,
234            conn: &rusqlite::Connection,
235            queue_wait: std::time::Duration,
236        ) -> Option<khive_storage::error::WriterTaskRequestState>;
237
238        fn reply_error_after_begin(
239            self: Box<Self>,
240            err: khive_storage::error::StorageError,
241            queue_wait: std::time::Duration,
242            transaction_acquire: std::time::Duration,
243        );
244    }
245}
246
247/// Type-erased write request the writer task's drain loop can hold in a
248/// homogeneous channel (`mpsc::Sender<Box<dyn AnyWriteRequest + Send>>`),
249/// while each concrete [`WriteRequest<R>`] still carries its own typed
250/// reply. Sealed: only this module may implement it (ADR-067 lines
251/// 210-212).
252pub trait AnyWriteRequest: sealed::Sealed + Send {
253    /// Runs this request's operation against `conn`, commits or rolls back
254    /// the enclosing transaction based on the outcome, and sends the
255    /// (possibly commit-failure-adjusted) result to the request's oneshot
256    /// reply channel.
257    ///
258    /// `conn` must already be inside a successfully-opened `BEGIN IMMEDIATE`
259    /// transaction opened by the caller (`run_writer_task`) — this method
260    /// issues only `COMMIT` / `ROLLBACK`, never `BEGIN`, so `run_writer_task`
261    /// remains the sole issuer of `BEGIN IMMEDIATE` (ADR-067 Component A).
262    /// Callers must use [`Self::reply_error`] instead when the enclosing
263    /// `BEGIN IMMEDIATE` itself failed — this method must not be invoked in
264    /// that case.
265    fn execute_and_reply(self: Box<Self>, conn: &Connection);
266
267    /// Runs this request's operation directly against `conn` — no
268    /// transaction wrap, no `COMMIT`/`ROLLBACK` — and sends the result to
269    /// the request's oneshot reply channel.
270    ///
271    /// Used only for [`Self::is_top_level`] requests: the drain loop calls
272    /// this INSTEAD of `execute_and_reply` for such requests, skipping
273    /// `BEGIN IMMEDIATE` entirely so a statement that must run outside any
274    /// transaction (e.g. `VACUUM`) can still be serialized through the
275    /// single writer owner.
276    fn execute_and_reply_top_level(self: Box<Self>, conn: &Connection);
277
278    /// Replies with `err` without running this request's operation or
279    /// touching `conn`.
280    ///
281    /// Used when the enclosing `BEGIN IMMEDIATE` failed (for example,
282    /// `SQLITE_BUSY` from lock contention with an unmigrated writer path
283    /// still holding the pool's writer mutex — reachable while only
284    /// `entity.rs` is routed through this channel). Running the operation
285    /// anyway would execute its DML against `conn` in autocommit mode,
286    /// landing partial writes for a request the caller is told failed.
287    /// Skipping the operation entirely keeps "the caller got an error" and
288    /// "no rows landed" true together.
289    fn reply_error(self: Box<Self>, err: StorageError);
290
291    /// `true` if the drain loop must run this request via
292    /// [`Self::execute_and_reply_top_level`] (no transaction wrap) instead
293    /// of [`Self::execute_and_reply`] (wrapped in `BEGIN IMMEDIATE`).
294    fn is_top_level(&self) -> bool;
295
296    fn is_checkpoint_bypass(&self) -> bool;
297
298    fn needs_vacuum_headroom(&self) -> bool;
299
300    /// Time from the caller constructing this request (before bounded-channel
301    /// admission) until the writer task dequeued it.
302    fn queue_wait(&self) -> Duration;
303}
304
305#[derive(Debug, Clone, Copy, PartialEq, Eq)]
306enum RollbackDisposition {
307    RolledBack,
308    SideEffectsUnknown,
309}
310
311/// Roll back a failed wrapped request and verify that the connection really
312/// returned to autocommit mode before it can serve another request.
313fn rollback_after_failure(conn: &Connection, failure_context: &'static str) -> RollbackDisposition {
314    match conn.execute_batch("ROLLBACK") {
315        Ok(()) if conn.is_autocommit() => RollbackDisposition::RolledBack,
316        Ok(()) => {
317            tracing::error!(
318                failure_context,
319                "writer transaction: ROLLBACK returned success but the connection is still in a \
320                 transaction; request side effects are unknown"
321            );
322            RollbackDisposition::SideEffectsUnknown
323        }
324        Err(rollback_error) => {
325            tracing::error!(
326                error = %rollback_error,
327                failure_context,
328                "writer transaction: rollback after request failure failed; request side effects are \
329                 unknown"
330            );
331            RollbackDisposition::SideEffectsUnknown
332        }
333    }
334}
335
336/// Execute an operation inside an already-open transaction and apply the
337/// shared commit, rollback, panic, and autocommit verification rules.
338pub(crate) fn execute_wrapped_transaction<R, F>(
339    conn: &Connection,
340    commit_operation: &'static str,
341    operation: F,
342) -> (Result<R, StorageError>, Option<WriterTaskRequestState>)
343where
344    F: FnOnce(&Connection) -> Result<R, StorageError>,
345{
346    let profiled = execute_wrapped_transaction_profiled(conn, commit_operation, None, operation);
347    (profiled.result, profiled.terminal_state)
348}
349
350struct ProfiledWrappedTransaction<R> {
351    result: Result<R, StorageError>,
352    terminal_state: Option<WriterTaskRequestState>,
353    body: Duration,
354    commit: Duration,
355}
356
357fn execute_wrapped_transaction_profiled<R, F>(
358    conn: &Connection,
359    commit_operation: &'static str,
360    db: Option<&str>,
361    operation: F,
362) -> ProfiledWrappedTransaction<R>
363where
364    F: FnOnce(&Connection) -> Result<R, StorageError>,
365{
366    let body_started = Instant::now();
367    let operation_outcome = catch_unwind(AssertUnwindSafe(|| operation(conn)));
368    let body = body_started.elapsed();
369
370    match operation_outcome {
371        Ok(Ok(value)) => {
372            let commit_started = Instant::now();
373            let commit_outcome = conn.execute_batch("COMMIT");
374            let commit = commit_started.elapsed();
375            match commit_outcome {
376                Ok(()) if conn.is_autocommit() => ProfiledWrappedTransaction {
377                    result: Ok(value),
378                    terminal_state: None,
379                    body,
380                    commit,
381                },
382                Ok(()) => {
383                    tracing::error!(
384                    "writer transaction: COMMIT returned success but the connection is still in \
385                     a transaction; request side effects are unknown"
386                );
387                    let request_state = WriterTaskRequestState::SideEffectsUnknown;
388                    ProfiledWrappedTransaction {
389                        result: Err(writer_task_terminated(request_state)),
390                        terminal_state: Some(request_state),
391                        body,
392                        commit,
393                    }
394                }
395                Err(commit_error) => {
396                    if let Some(db) = db {
397                        crate::timeout_sink::maybe_emit_sqlite_full(db, &commit_error);
398                    }
399                    match rollback_after_failure(conn, "commit failure") {
400                        RollbackDisposition::RolledBack => ProfiledWrappedTransaction {
401                            result: Err(StorageError::WriterTaskRequestFailed {
402                                request_state: WriterTaskRequestState::TransactionRolledBack,
403                                source: Box::new(StorageError::Pool {
404                                    operation: commit_operation.into(),
405                                    message: commit_error.to_string(),
406                                }),
407                            }),
408                            terminal_state: None,
409                            body,
410                            commit,
411                        },
412                        RollbackDisposition::SideEffectsUnknown => {
413                            let request_state = WriterTaskRequestState::SideEffectsUnknown;
414                            ProfiledWrappedTransaction {
415                                result: Err(writer_task_terminated_with_cause(
416                                    request_state,
417                                    &commit_error,
418                                )),
419                                terminal_state: Some(request_state),
420                                body,
421                                commit,
422                            }
423                        }
424                    }
425                }
426            }
427        }
428        Ok(Err(operation_error)) => {
429            match rollback_after_failure(conn, "request operation failure") {
430                RollbackDisposition::RolledBack => ProfiledWrappedTransaction {
431                    result: Err(StorageError::WriterTaskRequestFailed {
432                        request_state: WriterTaskRequestState::TransactionRolledBack,
433                        source: Box::new(operation_error),
434                    }),
435                    terminal_state: None,
436                    body,
437                    commit: Duration::ZERO,
438                },
439                RollbackDisposition::SideEffectsUnknown => {
440                    // SQLite may have already rolled back a failed transaction,
441                    // leaving no transaction for the explicit ROLLBACK.
442                    // Escalate the original cause once; retain native FULL evidence below.
443                    if let Some(db) = db {
444                        crate::timeout_sink::maybe_emit_sqlite_full(db, &operation_error);
445                    }
446                    let request_state = WriterTaskRequestState::SideEffectsUnknown;
447                    ProfiledWrappedTransaction {
448                        result: Err(writer_task_terminated_with_cause(
449                            request_state,
450                            &operation_error,
451                        )),
452                        terminal_state: Some(request_state),
453                        body,
454                        commit: Duration::ZERO,
455                    }
456                }
457            }
458        }
459        Err(_panic_payload) => {
460            let request_state = match rollback_after_failure(conn, "request panic") {
461                RollbackDisposition::RolledBack => WriterTaskRequestState::TransactionRolledBack,
462                RollbackDisposition::SideEffectsUnknown => {
463                    WriterTaskRequestState::SideEffectsUnknown
464                }
465            };
466            ProfiledWrappedTransaction {
467                result: Err(writer_task_terminated(request_state)),
468                terminal_state: Some(request_state),
469                body,
470                commit: Duration::ZERO,
471            }
472        }
473    }
474}
475
476impl<R: Send + 'static> sealed::Sealed for WriteRequest<R> {
477    fn execute_and_reply_reporting_terminal(
478        self: Box<Self>,
479        conn: &Connection,
480        tx_span: Option<khive_storage::tx_registry::TxHandle>,
481        queue_wait: Duration,
482        transaction_acquire: Duration,
483    ) -> Option<WriterTaskRequestState> {
484        // Keep the typed reply sender outside the unwind boundary. Calling
485        // `(self.op)(conn)` directly would drop `self.reply` while unwinding,
486        // leaving the active caller with only an untyped RecvError.
487        let WriteRequest {
488            op,
489            reply,
490            telemetry,
491            ..
492        } = *self;
493        let profiled = execute_wrapped_transaction_profiled(
494            conn,
495            "writer_task_commit",
496            Some(&telemetry.db),
497            op,
498        );
499        if let Err(error) = &profiled.result {
500            crate::timeout_sink::maybe_emit_sqlite_full(&telemetry.db, error);
501        }
502        // The reply wake can make the caller immediately observe the
503        // committed/error result. Deregister the transaction span first so
504        // tx_registry continues to mean "currently open SQL transaction",
505        // never "a transaction whose caller has already resumed" (#1790).
506        drop(tx_span);
507        telemetry.finish(
508            queue_wait,
509            transaction_acquire,
510            profiled.body,
511            profiled.commit,
512        );
513        // The receiver may already be gone (caller dropped its future) —
514        // that is not this task's problem to report.
515        let _ = reply.send(profiled.result);
516        profiled.terminal_state
517    }
518
519    fn execute_and_reply_top_level_reporting_terminal(
520        self: Box<Self>,
521        conn: &Connection,
522        queue_wait: Duration,
523    ) -> Option<WriterTaskRequestState> {
524        let WriteRequest {
525            op,
526            reply,
527            telemetry,
528            ..
529        } = *self;
530        let body_started = Instant::now();
531        let outcome = catch_unwind(AssertUnwindSafe(|| op(conn)));
532        let body = body_started.elapsed();
533        let telemetry_db = telemetry.db.clone();
534        telemetry.finish(queue_wait, Duration::ZERO, body, Duration::ZERO);
535        match outcome {
536            Ok(outcome) if conn.is_autocommit() => {
537                // No COMMIT/ROLLBACK here: this request explicitly did not
538                // open a transaction, so there is nothing to close.
539                if let Err(error) = &outcome {
540                    crate::timeout_sink::maybe_emit_sqlite_full(&telemetry_db, error);
541                }
542                let _ = reply.send(outcome);
543                None
544            }
545            Ok(_outcome) => {
546                tracing::error!(
547                    "writer task: top-level request returned with an open transaction; request \
548                     side effects are unknown"
549                );
550                let request_state = WriterTaskRequestState::SideEffectsUnknown;
551                let _ = reply.send(Err(writer_task_terminated(request_state)));
552                Some(request_state)
553            }
554            Err(_panic_payload) => {
555                // Statements completed before a top-level panic may already
556                // have autocommitted. Report the ambiguity; never invent a
557                // rollback for a request that opened no transaction.
558                let request_state = WriterTaskRequestState::SideEffectsUnknown;
559                let _ = reply.send(Err(writer_task_terminated(request_state)));
560                Some(request_state)
561            }
562        }
563    }
564
565    fn reply_error_after_begin(
566        self: Box<Self>,
567        err: StorageError,
568        queue_wait: Duration,
569        transaction_acquire: Duration,
570    ) {
571        let WriteRequest {
572            reply, telemetry, ..
573        } = *self;
574        telemetry.finish(
575            queue_wait,
576            transaction_acquire,
577            Duration::ZERO,
578            Duration::ZERO,
579        );
580        let _ = reply.send(Err(err));
581    }
582}
583
584impl<R: Send + 'static> AnyWriteRequest for WriteRequest<R> {
585    fn execute_and_reply(self: Box<Self>, conn: &Connection) {
586        let queue_wait = self.queue_wait();
587        let _ = sealed::Sealed::execute_and_reply_reporting_terminal(
588            self,
589            conn,
590            None,
591            queue_wait,
592            Duration::ZERO,
593        );
594    }
595
596    fn execute_and_reply_top_level(self: Box<Self>, conn: &Connection) {
597        let queue_wait = self.queue_wait();
598        let _ =
599            sealed::Sealed::execute_and_reply_top_level_reporting_terminal(self, conn, queue_wait);
600    }
601
602    fn reply_error(self: Box<Self>, err: StorageError) {
603        let queue_wait = self.queue_wait();
604        sealed::Sealed::reply_error_after_begin(self, err, queue_wait, Duration::ZERO);
605    }
606
607    fn is_top_level(&self) -> bool {
608        self.top_level
609    }
610
611    fn is_checkpoint_bypass(&self) -> bool {
612        self.checkpoint_bypass
613    }
614
615    fn needs_vacuum_headroom(&self) -> bool {
616        self.vacuum_copy_headroom
617    }
618
619    fn queue_wait(&self) -> Duration {
620        self.telemetry.queue_wait()
621    }
622}
623
624fn writer_task_terminated(request_state: WriterTaskRequestState) -> StorageError {
625    StorageError::writer_task_terminated(request_state)
626}
627
628fn writer_task_terminated_with_cause(
629    request_state: WriterTaskRequestState,
630    error: &(dyn std::error::Error + 'static),
631) -> StorageError {
632    let mut source = Some(error);
633    let mut sqlite_full_codes = None;
634    while let Some(current) = source {
635        if let Some(rusqlite::Error::SqliteFailure(code, _)) =
636            current.downcast_ref::<rusqlite::Error>()
637        {
638            if code.code == rusqlite::ErrorCode::DiskFull {
639                sqlite_full_codes = Some((code.extended_code & 0xff, code.extended_code));
640                break;
641            }
642        }
643        source = current.source();
644    }
645    // The original error was already escalated at the failing boundary.
646    // A source-free pair preserves caller evidence without duplicating that sink event.
647    StorageError::WriterTaskTerminated {
648        request_state,
649        sqlite_full_codes,
650    }
651}
652
653fn writer_task_begin_error(error: rusqlite::Error, busy_timeout: Duration) -> StorageError {
654    if crate::timeout_sink::is_busy_or_locked(&error) {
655        StorageError::WriterTaskBusy {
656            timeout_ms: u64::try_from(busy_timeout.as_millis()).unwrap_or(u64::MAX),
657        }
658    } else {
659        StorageError::Pool {
660            operation: "writer_task_begin".into(),
661            message: error.to_string(),
662        }
663    }
664}
665
666/// Sender half of the write queue. Cheaply cloneable (wraps an
667/// `mpsc::Sender`) — every migrated store that shares one writer task holds
668/// a clone of this handle.
669#[derive(Clone, Debug)]
670pub struct WriterTaskHandle {
671    tx: mpsc::Sender<Box<dyn AnyWriteRequest + Send>>,
672    /// Exact canonical backend identity used by the in-memory stage registry.
673    /// Kept separate from `db`, whose lossy display form is logging-only.
674    backend_key: Option<PathBuf>,
675    /// This handle's pool's writer-timeout sink identity (`timeout_sink::db_label`),
676    /// captured at spawn so `send_with_timeout`'s `WriteQueueFull` path can
677    /// report a `queue_saturation` sink row (ADR-136 D1 gate 6a) without
678    /// needing a `&ConnectionPool` reference.
679    db: String,
680    /// Slow-write latency bound (`timeout_sink::slow_write_threshold`),
681    /// resolved once at spawn. `None` disables slow-write rows. Captured
682    /// here rather than read per send so the caller path pays no env lookup
683    /// and tests get a deterministic value by setting the override before
684    /// spawning their own writer task.
685    slow_write_threshold: Option<std::time::Duration>,
686    /// Default enqueue-capacity deadline for [`Self::send_bounded`] /
687    /// [`Self::send_top_level_bounded`], captured from
688    /// `PoolConfig::write_admission_deadline_ms` at [`spawn`] (ADR-131
689    /// Decision 2). Bounds ONLY the wait for channel capacity, the same
690    /// boundary `send_with_timeout` already documents — never the reply wait
691    /// after acceptance (#1382).
692    ///
693    /// This is a dedicated admission authority, distinct from
694    /// `PoolConfig::checkout_timeout` (reader/pool checkout). The two used
695    /// to be conflated here before ADR-131 Decision 2 (#1382, #1643).
696    ///
697    /// ADR-131 Decision 2's outer-request-budget clamp ("when the caller's outer
698    /// request deadline leaves less time remaining than the configured
699    /// admission deadline, the admission deadline applied to that operation
700    /// is the remaining outer budget instead") is explicitly DEFERRED: no
701    /// outer request deadline is plumbed to `send_bounded` /
702    /// `send_top_level_bounded`'s call sites as of #1737, so there is no
703    /// budget to clamp against. When an outer deadline reaches this call
704    /// path, clamp here rather than faking a shorter deadline against a
705    /// nonexistent budget.
706    enqueue_timeout: std::time::Duration,
707}
708
709impl WriterTaskHandle {
710    /// Enqueue a write operation and return the oneshot receiver its reply
711    /// will arrive on, once the request has actually been accepted onto the
712    /// channel.
713    ///
714    /// Shared by [`Self::send`] and [`Self::send_with_timeout`] so that a
715    /// caller-supplied deadline (see `send_with_timeout`) can bound ONLY
716    /// this enqueue step — never the reply-wait that follows it. Once this
717    /// returns `Ok`, the request has been accepted by the writer task and
718    /// will either run to completion or receive a typed terminal error if an
719    /// earlier request kills the task before this operation begins. The
720    /// returned receiver must be awaited without a timeout; abandoning it
721    /// here would silently drop the request's eventual result, not cancel the
722    /// request itself.
723    async fn enqueue<R, F>(
724        &self,
725        op: F,
726    ) -> Result<oneshot::Receiver<Result<R, StorageError>>, StorageError>
727    where
728        R: Send + 'static,
729        F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
730    {
731        self.enqueue_inner(op, false, false, false).await
732    }
733
734    /// Shared enqueue path for both transaction-wrapped ([`Self::enqueue`])
735    /// and top-level ([`Self::send_top_level`]) requests — `top_level`
736    /// controls which [`AnyWriteRequest`] method the drain loop invokes.
737    async fn enqueue_inner<R, F>(
738        &self,
739        op: F,
740        top_level: bool,
741        checkpoint_bypass: bool,
742        vacuum_copy_headroom: bool,
743    ) -> Result<oneshot::Receiver<Result<R, StorageError>>, StorageError>
744    where
745        R: Send + 'static,
746        F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
747    {
748        let (reply_tx, reply_rx) = oneshot::channel();
749        let telemetry = WriteTelemetry::new(
750            self.backend_key.clone(),
751            self.db.clone(),
752            self.queue_depth(),
753            self.slow_write_threshold,
754        );
755        let request = WriteRequest {
756            op: Box::new(op),
757            reply: reply_tx,
758            top_level,
759            checkpoint_bypass,
760            vacuum_copy_headroom,
761            telemetry,
762        };
763
764        self.tx
765            .send(Box::new(request))
766            .await
767            .map_err(|_| writer_task_terminated(WriterTaskRequestState::NotStarted))?;
768
769        Ok(reply_rx)
770    }
771
772    /// Send a write operation to the writer task and await its typed reply.
773    ///
774    /// Backpressure: this suspends on the channel's `send().await` when the
775    /// bounded queue is full (ADR-067 "Channel capacity and queue-full
776    /// policy") — there is no `try_send` escape hatch. Callers that need a
777    /// deadline on that wait should use [`Self::send_with_timeout`] instead.
778    pub async fn send<R, F>(&self, op: F) -> Result<R, StorageError>
779    where
780        R: Send + 'static,
781        F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
782    {
783        let reply_rx = self.enqueue(op).await?;
784        reply_rx
785            .await
786            .map_err(|_| writer_task_terminated(WriterTaskRequestState::SideEffectsUnknown))?
787    }
788
789    /// Like [`Self::send`], but bounds the wait for the bounded channel to
790    /// free capacity with `timeout`.
791    ///
792    /// The timeout applies ONLY to enqueueing the request (the channel
793    /// `send().await` that can suspend on a full queue) — never to waiting
794    /// for the writer task's reply once the request has been accepted.
795    /// `StorageError::WriteQueueFull` means exactly "the bounded channel was
796    /// full and this request was never accepted"; it must never be returned
797    /// for a request that was accepted and is still executing (or already
798    /// committed) by the time `timeout` elapses — that would misreport a
799    /// slow op or a lock wait as a queue-capacity failure, and could tell a
800    /// caller a write failed when it actually landed. ADR-067's queue-full
801    /// policy has no immediate-error `try_send` path — only this caller-side
802    /// deadline on the enqueue step.
803    pub async fn send_with_timeout<R, F>(
804        &self,
805        op: F,
806        timeout: std::time::Duration,
807    ) -> Result<R, StorageError>
808    where
809        R: Send + 'static,
810        F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
811    {
812        let reply_rx = match tokio::time::timeout(timeout, self.enqueue(op)).await {
813            Ok(Ok(reply_rx)) => reply_rx,
814            Ok(Err(e)) => return Err(e),
815            Err(_elapsed) => {
816                let timeout_ms = timeout.as_millis() as u64;
817                crate::timeout_sink::emit_queue_saturation(&self.db, timeout_ms);
818                return Err(StorageError::WriteQueueFull { timeout_ms });
819            }
820        };
821
822        reply_rx
823            .await
824            .map_err(|_| writer_task_terminated(WriterTaskRequestState::SideEffectsUnknown))?
825    }
826
827    /// Like [`Self::send`], but bounds the enqueue-capacity wait with this
828    /// handle's configured `enqueue_timeout`
829    /// (`PoolConfig::write_admission_deadline_ms` at spawn time, ADR-131
830    /// Decision 2) instead of waiting indefinitely (#1382).
831    ///
832    /// Once the request is accepted onto the channel, the reply wait is
833    /// unbounded by this method, identical to [`Self::send`] — this bounds
834    /// only queue-capacity admission, never SQLite execution or reply
835    /// latency. Callers that need a non-default deadline should use
836    /// [`Self::send_with_timeout`] directly.
837    pub async fn send_bounded<R, F>(&self, op: F) -> Result<R, StorageError>
838    where
839        R: Send + 'static,
840        F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
841    {
842        self.send_with_timeout(op, self.enqueue_timeout).await
843    }
844
845    /// Send a write operation that MUST run outside any open transaction
846    /// (e.g. `VACUUM`, which SQLite forbids inside `BEGIN`/`COMMIT`) and
847    /// await its typed reply.
848    ///
849    /// Still serialized through the same single writer owner as
850    /// [`Self::send`] — the request goes through the identical bounded
851    /// channel and drain loop, one request at a time — but the drain loop
852    /// skips the per-request `BEGIN IMMEDIATE`/`COMMIT`/`ROLLBACK` wrap
853    /// entirely for this request (ADR-067 Component A). The single-writer
854    /// guarantee is preserved; only
855    /// the transaction wrap is skipped.
856    pub async fn send_top_level<R, F>(&self, op: F) -> Result<R, StorageError>
857    where
858        R: Send + 'static,
859        F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
860    {
861        let reply_rx = self.enqueue_inner(op, true, false, false).await?;
862        reply_rx
863            .await
864            .map_err(|_| writer_task_terminated(WriterTaskRequestState::SideEffectsUnknown))?
865    }
866
867    /// Like [`Self::send_top_level`], but bounds the enqueue-capacity wait
868    /// with this handle's configured `enqueue_timeout`, mirroring
869    /// [`Self::send_bounded`] for top-level (transaction-skipping) requests
870    /// (#1382).
871    pub async fn send_top_level_bounded<R, F>(&self, op: F) -> Result<R, StorageError>
872    where
873        R: Send + 'static,
874        F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
875    {
876        self.send_top_level_bounded_inner(op, false, false).await
877    }
878
879    /// Checkpoint-only top-level dispatch. It remains serialized by the
880    /// writer task but never waits for the capacity lease or floor probe.
881    pub(crate) async fn send_checkpoint_bounded<R, F>(&self, op: F) -> Result<R, StorageError>
882    where
883        R: Send + 'static,
884        F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
885    {
886        self.send_top_level_bounded_inner(op, true, false).await
887    }
888
889    /// Typed VACUUM dispatch with a copy-sized metadata estimate at dequeue.
890    pub(crate) async fn send_vacuum_bounded<R, F>(&self, op: F) -> Result<R, StorageError>
891    where
892        R: Send + 'static,
893        F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
894    {
895        self.send_top_level_bounded_inner(op, false, true).await
896    }
897
898    async fn send_top_level_bounded_inner<R, F>(
899        &self,
900        op: F,
901        checkpoint_bypass: bool,
902        vacuum_copy_headroom: bool,
903    ) -> Result<R, StorageError>
904    where
905        R: Send + 'static,
906        F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
907    {
908        let reply_rx = match tokio::time::timeout(
909            self.enqueue_timeout,
910            self.enqueue_inner(op, true, checkpoint_bypass, vacuum_copy_headroom),
911        )
912        .await
913        {
914            Ok(Ok(reply_rx)) => reply_rx,
915            Ok(Err(e)) => return Err(e),
916            Err(_elapsed) => {
917                let timeout_ms = self.enqueue_timeout.as_millis() as u64;
918                crate::timeout_sink::emit_queue_saturation(&self.db, timeout_ms);
919                return Err(StorageError::WriteQueueFull { timeout_ms });
920            }
921        };
922
923        reply_rx
924            .await
925            .map_err(|_| writer_task_terminated(WriterTaskRequestState::SideEffectsUnknown))?
926    }
927
928    /// Current write-queue backlog depth: requests enqueued but not yet
929    /// accepted by the writer task's drain loop.
930    ///
931    /// Reads `mpsc::Sender::max_capacity() - capacity()`, so it is a
932    /// point-in-time snapshot racy under concurrent senders/the drain loop
933    /// draining concurrently — acceptable for a monitoring gauge (the
934    /// load/perf harness metrics read-surface), never used for any correctness
935    /// decision.
936    pub fn queue_depth(&self) -> usize {
937        self.tx.max_capacity() - self.tx.capacity()
938    }
939
940    /// The bounded channel's configured capacity
941    /// (`PoolConfig::write_queue_capacity`).
942    pub fn capacity(&self) -> usize {
943        self.tx.max_capacity()
944    }
945}
946
947/// Spawn the write-owner task (ADR-067 Component A) on the current Tokio
948/// runtime.
949///
950/// Opens a dedicated standalone writer connection independent of the pool's
951/// Mutex-guarded `writer()` connection used by unmigrated paths. That one-time
952/// infrastructure open is uncounted; every dequeued top-level request or
953/// successful `BEGIN IMMEDIATE` increments the writer-task acquisition class.
954/// Returns the cloneable [`WriterTaskHandle`] sender half. The task normally
955/// runs until every handle clone is dropped and the channel closes; a request
956/// panic, failed rollback, or poisoned connection puts it into the permanent
957/// terminal state documented below.
958///
959/// # Terminal settlement
960/// Retiring an unsettled writer, including when its blocking owner unwinds,
961/// releases the volume lease only after cleanup restored autocommit or the
962/// connection closed. If neither is possible, path, volume, and errors are
963/// reported to stderr and tracing and the pool is poisoned: that request
964/// reports [`SqliteError::WriterSettlementUnknown`] and every later write is
965/// refused with [`SqliteError::WriterPoisoned`] before it starts.
966///
967/// `capacity` bounds the channel (`PoolConfig::write_queue_capacity` /
968/// `KHIVE_WRITE_QUEUE_CAPACITY`, ADR-067 recommends 256).
969///
970/// # Errors
971/// Must be called from within a Tokio runtime context (calls
972/// `tokio::spawn`). Returns an error if the pool cannot open a standalone
973/// writer connection (e.g. an in-memory pool has no standalone-connection
974/// support). See `crates/khive-db/docs/api/writer-task.md` for the
975/// migration-slice scope this commits per `BEGIN IMMEDIATE`.
976pub fn spawn(pool: &ConnectionPool, capacity: usize) -> Result<WriterTaskHandle, SqliteError> {
977    // The lifetime connection is infrastructure, not one acquisition per
978    // write. Each dequeued request is counted below at the task's actual
979    // ownership boundary instead.
980    let conn = pool.open_standalone_writer_untracked()?;
981    let acquisition_counters = pool.writer_acquisition_counters();
982    let write_admission = pool.write_admission();
983    let busy_timeout = pool.config().busy_timeout;
984    let origin = pool.origin();
985    let backend_key = writer_db_key(pool);
986    let db = crate::timeout_sink::db_label(pool);
987    let (tx, rx) = mpsc::channel(capacity.max(1));
988    let join = tokio::spawn(run_writer_task(
989        conn,
990        rx,
991        origin,
992        db.clone(),
993        acquisition_counters,
994        write_admission,
995        busy_timeout,
996    ));
997    // Stored on the pool (not returned) so the handle's clone-and-share
998    // contract stays untouched; see ConnectionPool::take_writer_task_join
999    // for who awaits it and why.
1000    pool.set_writer_task_join(join);
1001    Ok(WriterTaskHandle {
1002        tx,
1003        backend_key,
1004        db,
1005        slow_write_threshold: crate::timeout_sink::slow_write_threshold(),
1006        enqueue_timeout: std::time::Duration::from_millis(
1007            pool.config().write_admission_deadline_ms,
1008        ),
1009    })
1010}
1011
1012/// Permanently close admission, then reply to every request that was already
1013/// accepted into the bounded channel without invoking its operation closure.
1014///
1015/// Closing before draining is load-bearing: draining an open receiver would
1016/// wait forever while [`WriterTaskHandle`] clones still exist, while dropping
1017/// the receiver immediately would discard buffered requests and their typed
1018/// reply senders.
1019async fn close_and_fail_queued_requests(rx: &mut mpsc::Receiver<Box<dyn AnyWriteRequest + Send>>) {
1020    rx.close();
1021    while let Some(request) = rx.recv().await {
1022        request.reply_error(writer_task_terminated(WriterTaskRequestState::NotStarted));
1023    }
1024}
1025
1026/// Acquire the request transaction without multiplying its busy-timeout budget.
1027/// The setter is supplied so tests can fail SQLite's timeout update while still
1028/// executing real BEGIN statements; production uses `Connection::busy_timeout`.
1029fn begin_immediate_with_retry(
1030    conn: &Connection,
1031    acquisition_counters: &WriterAcquisitionCounters,
1032    busy_timeout: Duration,
1033    mut set_busy_timeout: impl FnMut(&Connection, Duration) -> rusqlite::Result<()>,
1034) -> (rusqlite::Result<()>, Duration, u32) {
1035    let transaction_acquire_started = Instant::now();
1036    let mut begin_attempt = 1_u32;
1037    let mut retry_delays = WRITER_BEGIN_RETRY_DELAYS.into_iter();
1038    let mut busy_timeout_lowered = false;
1039    let begin_outcome = loop {
1040        match conn.execute_batch("BEGIN IMMEDIATE") {
1041            Ok(()) => break Ok(()),
1042            Err(error) if crate::timeout_sink::is_busy_or_locked(&error) => {
1043                // Every busy/locked refusal is counted here,
1044                // whether or not it goes on to be retried: this
1045                // is the pre-PR meaning of the counter, and the
1046                // outer match below no longer counts the final
1047                // refusal a second time.
1048                acquisition_counters.record_writer_task_begin_busy();
1049                let Some(delay) = retry_delays.next() else {
1050                    break Err(error);
1051                };
1052                // The whole retry sequence shares one
1053                // `busy_timeout` budget: each SQLite attempt is
1054                // itself bounded by `busy_timeout`, so without a
1055                // shrinking budget three attempts could each
1056                // wait out a full window and multiply the
1057                // serialized writer's contention window instead
1058                // of bounding it.
1059                let remaining_budget =
1060                    busy_timeout.saturating_sub(transaction_acquire_started.elapsed());
1061                if remaining_budget.is_zero() {
1062                    break Err(error);
1063                }
1064                if let Err(set_err) = set_busy_timeout(conn, remaining_budget) {
1065                    tracing::warn!(
1066                        error = %set_err,
1067                        "writer task: failed to lower busy_timeout for BEGIN \
1068                         retry; surfacing the original busy refusal"
1069                    );
1070                    // A retry with the old timeout could exceed the shared budget.
1071                    break Err(error);
1072                }
1073                busy_timeout_lowered = true;
1074                // Count only refusals that will actually be retried. A failed
1075                // timeout reduction leaves this refusal visible to the caller.
1076                acquisition_counters.record_writer_task_begin_busy_absorbed();
1077                tracing::debug!(
1078                    attempt = begin_attempt,
1079                    backoff_ms = delay.as_millis() as u64,
1080                    budget_remaining_ms = remaining_budget.as_millis() as u64,
1081                    "writer task: BEGIN IMMEDIATE refused busy; retrying before \
1082                     request execution"
1083                );
1084                std::thread::sleep(delay);
1085                begin_attempt = begin_attempt.saturating_add(1);
1086            }
1087            Err(error) => break Err(error),
1088        }
1089    };
1090    let transaction_acquire = transaction_acquire_started.elapsed();
1091    // Restore the pool-configured busy_timeout only if a retry
1092    // actually lowered it, so the common uncontended request
1093    // never pays for an extra pragma write; the next request
1094    // dequeued on this same connection must still start from
1095    // the full configured budget.
1096    if busy_timeout_lowered {
1097        if let Err(restore_err) = set_busy_timeout(conn, busy_timeout) {
1098            tracing::warn!(
1099                error = %restore_err,
1100                "writer task: failed to restore busy_timeout after a BEGIN retry \
1101                 sequence"
1102            );
1103        }
1104    }
1105    (begin_outcome, transaction_acquire, begin_attempt)
1106}
1107
1108struct BlockingWriterConnection {
1109    conn: Option<Connection>,
1110    admission: Arc<crate::pool::WriteAdmission>,
1111    volume_lease: Option<crate::disk_guard::VolumeLease>,
1112}
1113
1114impl BlockingWriterConnection {
1115    fn new(conn: Connection, admission: Arc<crate::pool::WriteAdmission>) -> Self {
1116        Self {
1117            conn: Some(conn),
1118            admission,
1119            volume_lease: None,
1120        }
1121    }
1122
1123    fn finish(
1124        mut self,
1125        state: Option<WriterTaskRequestState>,
1126    ) -> (Option<Connection>, Option<WriterTaskRequestState>) {
1127        if state.is_some() {
1128            self.retire();
1129            (None, state)
1130        } else {
1131            (self.conn.take(), None)
1132        }
1133    }
1134
1135    fn retire(&mut self) {
1136        if let Some(conn) = self.conn.take() {
1137            // A terminal failure has poisoned the admission; the request that
1138            // retired this writer already carries its terminal state.
1139            let _ = self.admission.close_retired_connection(conn);
1140        }
1141    }
1142}
1143
1144impl Drop for BlockingWriterConnection {
1145    fn drop(&mut self) {
1146        self.retire();
1147    }
1148}
1149
1150impl std::ops::Deref for BlockingWriterConnection {
1151    type Target = Connection;
1152
1153    fn deref(&self) -> &Connection {
1154        self.conn
1155            .as_ref()
1156            .expect("blocking writer owns its connection until settlement")
1157    }
1158}
1159
1160/// Drain loop: the sole caller of `BEGIN IMMEDIATE` for write traffic routed
1161/// through the channel. Busy/locked `BEGIN IMMEDIATE` refusals receive the
1162/// bounded retry above before a final failure replies the request's error via
1163/// [`AnyWriteRequest::reply_error`] without invoking the request's closure.
1164///
1165/// A request-operation panic is contained inside the concrete
1166/// [`WriteRequest<R>`] so its typed reply survives. A panic, failed rollback,
1167/// or otherwise poisoned connection makes the task terminal. The task then
1168/// closes admission, explicitly fails every already-queued request as
1169/// [`WriterTaskRequestState::NotStarted`], and exits permanently. There is no
1170/// supervisor or connection restart; later sends observe the closed channel.
1171/// See
1172/// `crates/khive-db/docs/api/writer-task.md` for the full failure matrix.
1173async fn run_writer_task(
1174    mut conn: Connection,
1175    mut rx: mpsc::Receiver<Box<dyn AnyWriteRequest + Send>>,
1176    origin: khive_storage::tx_registry::TxOrigin,
1177    db: String,
1178    acquisition_counters: Arc<WriterAcquisitionCounters>,
1179    write_admission: Arc<crate::pool::WriteAdmission>,
1180    busy_timeout: Duration,
1181) {
1182    while let Some(request) = rx.recv().await {
1183        // The bounded-channel wait ends at this exact dequeue boundary.
1184        // Sampling inside the blocking closure would misattribute a saturated
1185        // Tokio blocking pool to writer-queue contention (#1849).
1186        let queue_wait = request.queue_wait();
1187        let origin = origin.clone();
1188        let blocking_counters = Arc::clone(&acquisition_counters);
1189        let blocking_admission = Arc::clone(&write_admission);
1190        let blocking_db = db.clone();
1191        let outcome = tokio::task::spawn_blocking(move || {
1192            let mut conn = BlockingWriterConnection::new(conn, Arc::clone(&blocking_admission));
1193            let acquisition_counters = blocking_counters;
1194            // A top-level request deliberately skips BEGIN, so it would
1195            // silently join any transaction leaked by an earlier request.
1196            // Refuse every request before dispatch if the connection is not
1197            // demonstrably clean, then retire the writer task.
1198            if !conn.is_autocommit() {
1199                tracing::error!(
1200                    "writer task: connection is not in autocommit mode before request dispatch; \
1201                     retiring the poisoned writer without running the request"
1202                );
1203                let request_state = WriterTaskRequestState::NotStarted;
1204                request.reply_error(writer_task_terminated(request_state));
1205                return conn.finish(Some(request_state));
1206            }
1207
1208            // The lease precedes SQLite writer acquisition and remains held
1209            // until the request has committed, rolled back, or returned to
1210            // autocommit. Checkpoint requests deliberately skip both steps.
1211            conn.volume_lease = if request.is_checkpoint_bypass() {
1212                None
1213            } else {
1214                match blocking_admission.acquire() {
1215                    Ok(lease) => lease,
1216                    Err(error) => {
1217                        request.reply_error(error.into_storage_error(
1218                            khive_storage::StorageCapability::Sql,
1219                            "writer_task_admission",
1220                        ));
1221                        return conn.finish(None);
1222                    }
1223                }
1224            };
1225
1226            let terminal_state = if request.is_top_level() {
1227                if !request.is_checkpoint_bypass() {
1228                    let admission = if request.needs_vacuum_headroom() {
1229                        blocking_admission.check_for_vacuum()
1230                    } else {
1231                        blocking_admission.check()
1232                    };
1233                    if let Err(error) = admission {
1234                        request.reply_error(error.into_storage_error(
1235                            khive_storage::StorageCapability::Sql,
1236                            "writer_task_admission",
1237                        ));
1238                        return conn.finish(None);
1239                    }
1240                }
1241                // ADR-067 Component A:
1242                // no BEGIN IMMEDIATE for this request — some statements
1243                // (e.g. VACUUM) are rejected by SQLite inside any open
1244                // transaction. Still runs on this task's dedicated
1245                // connection and still serialized one-request-at-a-time by
1246                // this same drain loop, so the single-writer guarantee
1247                // holds; only the transaction wrap is skipped.
1248                acquisition_counters.record_writer_task_acquisition();
1249                sealed::Sealed::execute_and_reply_top_level_reporting_terminal(
1250                    request, &conn, queue_wait,
1251                )
1252            } else {
1253                let tx_span = khive_storage::tx_registry::register_scoped(
1254                    Some("writer_task_tx".to_string()),
1255                    origin,
1256                );
1257                let (begin_outcome, transaction_acquire, begin_attempt) =
1258                    begin_immediate_with_retry(
1259                        &conn,
1260                        &acquisition_counters,
1261                        busy_timeout,
1262                        Connection::busy_timeout,
1263                    );
1264                match begin_outcome {
1265                    Ok(()) => {
1266                        if let Err(error) = blocking_admission.check() {
1267                            let request_state =
1268                                match rollback_after_failure(&conn, "capacity admission") {
1269                                    RollbackDisposition::RolledBack => None,
1270                                    RollbackDisposition::SideEffectsUnknown => {
1271                                        Some(WriterTaskRequestState::SideEffectsUnknown)
1272                                    }
1273                                };
1274                            drop(tx_span);
1275                            let error = if let Some(state) = request_state {
1276                                writer_task_terminated(state)
1277                            } else {
1278                                error.into_storage_error(
1279                                    khive_storage::StorageCapability::Sql,
1280                                    "writer_task_admission",
1281                                )
1282                            };
1283                            sealed::Sealed::reply_error_after_begin(
1284                                request,
1285                                error,
1286                                queue_wait,
1287                                transaction_acquire,
1288                            );
1289                            return conn.finish(request_state);
1290                        }
1291                        acquisition_counters.record_writer_task_acquisition();
1292                        sealed::Sealed::execute_and_reply_reporting_terminal(
1293                            request,
1294                            &conn,
1295                            Some(tx_span),
1296                            queue_wait,
1297                            transaction_acquire,
1298                        )
1299                    }
1300                    Err(e) => {
1301                        // Do NOT run the request's operation: `conn` never
1302                        // entered a transaction, so executing the op's DML
1303                        // here would run in autocommit mode and land partial
1304                        // writes for a request the caller is about to be told
1305                        // failed.
1306                        crate::timeout_sink::maybe_emit_sqlite_full(&blocking_db, &e);
1307                        tracing::warn!(
1308                            error = %e,
1309                            attempts = begin_attempt,
1310                            "writer task: BEGIN IMMEDIATE failed; replying an \
1311                             error without running the request's operation"
1312                        );
1313                        // Although BEGIN never opened a SQL transaction, the
1314                        // scoped attempt is observable in tx_registry. Drop
1315                        // it before waking the caller for the same reply-side
1316                        // lifecycle guarantee as successful requests (#1790).
1317                        drop(tx_span);
1318                        // A busy/locked refusal was already counted inside
1319                        // the retry loop above, including this final one —
1320                        // classifying and counting it again here would
1321                        // double it. Only a non-busy BEGIN failure (which
1322                        // never enters the retry loop's busy arm) still
1323                        // needs its counter recorded at this seam.
1324                        let begin_error = writer_task_begin_error(e, busy_timeout);
1325                        if !matches!(&begin_error, StorageError::WriterTaskBusy { .. }) {
1326                            acquisition_counters.record_writer_task_begin_error();
1327                        }
1328                        sealed::Sealed::reply_error_after_begin(
1329                            request,
1330                            begin_error,
1331                            queue_wait,
1332                            transaction_acquire,
1333                        );
1334                        None
1335                    }
1336                }
1337            };
1338            conn.finish(terminal_state)
1339        })
1340        .await;
1341
1342        match outcome {
1343            Ok((Some(returned_conn), None)) => conn = returned_conn,
1344            Ok((_returned_conn, Some(request_state))) => {
1345                acquisition_counters.record_writer_task_request_failure();
1346                if request_state == WriterTaskRequestState::SideEffectsUnknown {
1347                    acquisition_counters.record_writer_task_side_effects_unknown();
1348                }
1349                tracing::error!(
1350                    request_state = %request_state,
1351                    "writer task reached a terminal request or connection state; closing and \
1352                     failing the queue without restarting"
1353                );
1354                crate::timeout_sink::emit_writer_task_retirement(
1355                    &db,
1356                    &format!("terminal request state: {request_state}"),
1357                );
1358                close_and_fail_queued_requests(&mut rx).await;
1359                return;
1360            }
1361            Ok((None, None)) => unreachable!("a reusable writer returns its connection"),
1362            Err(join_err) => {
1363                acquisition_counters.record_writer_task_request_failure();
1364                tracing::error!(
1365                    error = %join_err,
1366                    "writer task blocking closure failed outside the request \
1367                     panic boundary; closing and failing the queue without restarting"
1368                );
1369                crate::timeout_sink::emit_writer_task_retirement(
1370                    &db,
1371                    &format!("blocking closure join failure: {join_err}"),
1372                );
1373                close_and_fail_queued_requests(&mut rx).await;
1374                return;
1375            }
1376        }
1377    }
1378}
1379
1380#[cfg(test)]
1381mod tests {
1382    use super::*;
1383    use crate::pool::PoolConfig;
1384    use rusqlite::hooks::{AuthAction, AuthContext, Authorization, TransactionOperation};
1385    use serial_test::serial;
1386
1387    #[test]
1388    fn begin_error_classification_is_code_based_and_narrow() {
1389        for code in [rusqlite::ffi::SQLITE_BUSY, rusqlite::ffi::SQLITE_LOCKED] {
1390            let error = rusqlite::Error::SqliteFailure(
1391                rusqlite::ffi::Error::new(code),
1392                Some("rendered text is irrelevant".to_string()),
1393            );
1394            assert!(matches!(
1395                writer_task_begin_error(error, Duration::from_millis(175)),
1396                StorageError::WriterTaskBusy { timeout_ms: 175 }
1397            ));
1398        }
1399
1400        let structural = rusqlite::Error::SqliteFailure(
1401            rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_CONSTRAINT),
1402            Some("database is locked".to_string()),
1403        );
1404        assert!(matches!(
1405            writer_task_begin_error(structural, Duration::from_millis(175)),
1406            StorageError::Pool { ref operation, .. } if operation == "writer_task_begin"
1407        ));
1408    }
1409    use std::future::Future;
1410    use std::pin::Pin;
1411    use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
1412    use std::sync::mpsc as std_mpsc;
1413    use std::sync::{Arc, Mutex};
1414    use std::task::{Context, Poll, Wake, Waker};
1415    use std::time::Duration;
1416
1417    fn file_pool(path: &std::path::Path) -> ConnectionPool {
1418        let cfg = PoolConfig {
1419            path: Some(path.to_path_buf()),
1420            ..PoolConfig::for_test()
1421        };
1422        ConnectionPool::new(cfg).expect("pool open")
1423    }
1424
1425    fn deny_commit_and_rollback(ctx: AuthContext<'_>) -> Authorization {
1426        match ctx.action {
1427            // SQLite reports COMMIT as the non-exhaustive Unknown transaction
1428            // operation in rusqlite 0.40; ROLLBACK has its own variant.
1429            AuthAction::Transaction {
1430                operation: TransactionOperation::Unknown | TransactionOperation::Rollback,
1431            } => Authorization::Deny,
1432            _ => Authorization::Allow,
1433        }
1434    }
1435
1436    fn deny_commit(ctx: AuthContext<'_>) -> Authorization {
1437        match ctx.action {
1438            AuthAction::Transaction {
1439                operation: TransactionOperation::Unknown,
1440            } => Authorization::Deny,
1441            _ => Authorization::Allow,
1442        }
1443    }
1444
1445    fn deny_rollback(ctx: AuthContext<'_>) -> Authorization {
1446        match ctx.action {
1447            AuthAction::Transaction {
1448                operation: TransactionOperation::Rollback,
1449            } => Authorization::Deny,
1450            _ => Authorization::Allow,
1451        }
1452    }
1453
1454    fn assert_writer_task_terminal_state<T: std::fmt::Debug>(
1455        result: Result<T, StorageError>,
1456        expected: WriterTaskRequestState,
1457    ) {
1458        match result {
1459            Err(StorageError::WriterTaskTerminated { request_state, .. }) => {
1460                assert_eq!(request_state, expected)
1461            }
1462            other => panic!("expected WriterTaskTerminated({expected:?}), got {other:?}"),
1463        }
1464    }
1465
1466    #[test]
1467    fn terminal_native_full_evidence_uses_the_original_error_chain() {
1468        let extended = rusqlite::ffi::SQLITE_FULL | (3 << 8);
1469        let native = rusqlite::Error::SqliteFailure(
1470            rusqlite::ffi::Error::new(extended),
1471            Some("synthetic classifier control".into()),
1472        );
1473        let wrapped = StorageError::driver(khive_storage::StorageCapability::Sql, "nested", native);
1474        let error =
1475            writer_task_terminated_with_cause(WriterTaskRequestState::SideEffectsUnknown, &wrapped);
1476        assert!(matches!(error, StorageError::WriterTaskTerminated {
1477            request_state: WriterTaskRequestState::SideEffectsUnknown,
1478            sqlite_full_codes: Some((rusqlite::ffi::SQLITE_FULL, code)),
1479        } if code == extended));
1480        assert!(
1481            std::error::Error::source(&error).is_none(),
1482            "evidence must not replay the original sink cause"
1483        );
1484        assert!(!error.is_retryable());
1485        assert_eq!(error.capability(), None);
1486
1487        for code in [rusqlite::ffi::SQLITE_BUSY, rusqlite::ffi::SQLITE_IOERR] {
1488            let native = rusqlite::Error::SqliteFailure(
1489                rusqlite::ffi::Error::new(code),
1490                Some("database or disk is full".into()),
1491            );
1492            assert!(matches!(
1493                writer_task_terminated_with_cause(
1494                    WriterTaskRequestState::SideEffectsUnknown,
1495                    &native
1496                ),
1497                StorageError::WriterTaskTerminated {
1498                    sqlite_full_codes: None,
1499                    ..
1500                }
1501            ));
1502        }
1503        assert!(matches!(
1504            writer_task_terminated(WriterTaskRequestState::NotStarted),
1505            StorageError::WriterTaskTerminated {
1506                sqlite_full_codes: None,
1507                ..
1508            }
1509        ));
1510    }
1511
1512    struct ParkedWake {
1513        entered: std_mpsc::SyncSender<()>,
1514        release: Mutex<std_mpsc::Receiver<()>>,
1515    }
1516
1517    impl Wake for ParkedWake {
1518        fn wake(self: Arc<Self>) {
1519            self.entered
1520                .send(())
1521                .expect("reply sender must rendezvous with the test");
1522            self.release
1523                .lock()
1524                .unwrap_or_else(|poisoned| poisoned.into_inner())
1525                .recv()
1526                .expect("test must release the parked reply sender");
1527        }
1528    }
1529
1530    fn arm_parked_wake<F: Future>(
1531        mut future: Pin<&mut F>,
1532    ) -> (std_mpsc::Receiver<()>, std_mpsc::Sender<()>) {
1533        let (entered_tx, entered_rx) = std_mpsc::sync_channel(0);
1534        let (release_tx, release_rx) = std_mpsc::channel();
1535        let waker = Waker::from(Arc::new(ParkedWake {
1536            entered: entered_tx,
1537            release: Mutex::new(release_rx),
1538        }));
1539        let mut context = Context::from_waker(&waker);
1540        assert!(
1541            matches!(future.as_mut().poll(&mut context), Poll::Pending),
1542            "writer send must remain pending until its operation replies"
1543        );
1544        (entered_rx, release_tx)
1545    }
1546
1547    fn poll_ready<F: Future>(mut future: Pin<&mut F>) -> F::Output {
1548        let mut context = Context::from_waker(Waker::noop());
1549        match future.as_mut().poll(&mut context) {
1550            Poll::Ready(output) => output,
1551            Poll::Pending => panic!("reply wake must make the writer send ready"),
1552        }
1553    }
1554
1555    fn database_tx_view(pool: &ConnectionPool) -> khive_storage::tx_registry::TxOriginFilter {
1556        match pool.origin() {
1557            khive_storage::tx_registry::TxOrigin::Database(identity) => {
1558                khive_storage::tx_registry::TxOriginFilter::Secondary(identity)
1559            }
1560            other => panic!("expected a file-backed database origin, got {other:?}"),
1561        }
1562    }
1563
1564    async fn wait_for_writer_span_to_close(view: &khive_storage::tx_registry::TxOriginFilter) {
1565        tokio::time::timeout(Duration::from_secs(5), async {
1566            while khive_storage::tx_registry::any_open_labeled(view, "writer_task_tx") {
1567                tokio::task::yield_now().await;
1568            }
1569        })
1570        .await
1571        .expect("writer task transaction span must eventually close");
1572    }
1573
1574    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1575    #[serial(tx_registry)]
1576    async fn writer_task_connection_maintains_rfc3339_expression_indexes() {
1577        let dir = tempfile::tempdir().unwrap();
1578        let path = dir.path().join("writer_task_rfc3339_expression_index.db");
1579        let pool = file_pool(&path);
1580        {
1581            let writer = pool.writer().expect("pooled writer");
1582            writer
1583                .conn()
1584                .execute_batch(
1585                    "CREATE TABLE deadlines(id INTEGER PRIMARY KEY, due TEXT);
1586                     CREATE INDEX idx_deadlines_strict \
1587                         ON deadlines(ifnull(khive_rfc3339_strict_key(due), x''));",
1588                )
1589                .expect("pooled writer registers the key function");
1590        }
1591        let handle = spawn(&pool, 8).expect("writer task spawn");
1592
1593        let inserted = handle
1594            .send(|conn| {
1595                conn.execute(
1596                    "INSERT INTO deadlines(id, due) VALUES (1, '2026-01-01T00:00:00Z')",
1597                    [],
1598                )
1599                .map_err(|error| StorageError::Pool {
1600                    operation: "test_insert".into(),
1601                    message: error.to_string(),
1602                })
1603            })
1604            .await
1605            .expect("the writer task's connection maintains the expression index");
1606        assert_eq!(inserted, 1);
1607    }
1608
1609    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1610    #[serial(tx_registry)]
1611    async fn successful_send_reply_waits_for_writer_tx_deregistration() {
1612        let dir = tempfile::tempdir().unwrap();
1613        let path = dir.path().join("writer_task_success_reply_lifecycle.db");
1614        let pool = file_pool(&path);
1615        let view = database_tx_view(&pool);
1616        let handle = spawn(&pool, 8).expect("writer task spawn");
1617        let (op_started_tx, op_started_rx) = std_mpsc::sync_channel(0);
1618        let (op_release_tx, op_release_rx) = std_mpsc::channel();
1619
1620        let send = handle.send(move |_conn| {
1621            op_started_tx
1622                .send(())
1623                .expect("operation must rendezvous with the test");
1624            op_release_rx
1625                .recv()
1626                .expect("test must release the operation");
1627            Ok::<_, StorageError>(())
1628        });
1629        tokio::pin!(send);
1630        let (reply_entered_rx, reply_release_tx) = arm_parked_wake(send.as_mut());
1631
1632        op_started_rx
1633            .recv_timeout(Duration::from_secs(5))
1634            .expect("writer operation must start");
1635        op_release_tx.send(()).expect("release writer operation");
1636        reply_entered_rx
1637            .recv_timeout(Duration::from_secs(5))
1638            .expect("reply sender must wake the waiting caller");
1639
1640        let reply = poll_ready(send.as_mut());
1641        let span_was_open_at_reply =
1642            khive_storage::tx_registry::any_open_labeled(&view, "writer_task_tx");
1643
1644        reply_release_tx
1645            .send(())
1646            .expect("release parked reply sender");
1647        wait_for_writer_span_to_close(&view).await;
1648
1649        reply.expect("committed operation reply");
1650        assert!(
1651            !span_was_open_at_reply,
1652            "a successful caller reply must not become observable while its committed writer_task_tx span remains registered"
1653        );
1654    }
1655
1656    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1657    #[serial(tx_registry)]
1658    async fn begin_failure_reply_waits_for_writer_tx_deregistration() {
1659        let dir = tempfile::tempdir().unwrap();
1660        let path = dir.path().join("writer_task_begin_reply_lifecycle.db");
1661        // A lock directory of its own keeps this deliberate writer hold off the
1662        // volume lease the other tests in this process share.
1663        let cfg = PoolConfig {
1664            path: Some(path.clone()),
1665            volume_lock_dir: Some(dir.path().join("volume-locks")),
1666            busy_timeout: Duration::from_millis(150),
1667            ..PoolConfig::for_test()
1668        };
1669        let pool = ConnectionPool::new(cfg).unwrap();
1670        let view = database_tx_view(&pool);
1671        let handle = spawn(&pool, 8).expect("writer task spawn");
1672        // An independent SQLite writer does not hold khive's cooperative
1673        // volume lease, so this reaches the intended BEGIN-busy branch.
1674        let lock_holder = rusqlite::Connection::open(&path).expect("external writer");
1675        lock_holder
1676            .execute_batch("BEGIN IMMEDIATE")
1677            .expect("hold database write lock");
1678        let op_ran = Arc::new(AtomicBool::new(false));
1679        let op_ran_in_request = Arc::clone(&op_ran);
1680
1681        let send = handle.send(move |_conn| {
1682            op_ran_in_request.store(true, Ordering::SeqCst);
1683            Ok::<_, StorageError>(())
1684        });
1685        tokio::pin!(send);
1686        let (reply_entered_rx, reply_release_tx) = arm_parked_wake(send.as_mut());
1687        reply_entered_rx
1688            .recv_timeout(Duration::from_secs(5))
1689            .expect("BEGIN failure must wake the waiting caller");
1690
1691        let reply = poll_ready(send.as_mut());
1692        let span_was_open_at_reply =
1693            khive_storage::tx_registry::any_open_labeled(&view, "writer_task_tx");
1694
1695        reply_release_tx
1696            .send(())
1697            .expect("release parked reply sender");
1698        wait_for_writer_span_to_close(&view).await;
1699        lock_holder
1700            .execute_batch("ROLLBACK")
1701            .expect("release database write lock");
1702
1703        assert!(
1704            matches!(
1705                &reply,
1706                Err(StorageError::WriterTaskBusy { timeout_ms }) if *timeout_ms == 150
1707            ),
1708            "expected typed retryable writer-task contention, got {reply:?}"
1709        );
1710        assert!(!op_ran.load(Ordering::SeqCst));
1711        assert!(
1712            !span_was_open_at_reply,
1713            "a BEGIN-failure caller reply must not become observable while its writer_task_tx span remains registered"
1714        );
1715    }
1716
1717    // `#[serial(tx_registry)]`: `run_writer_task` registers a `writer_task_tx`
1718    // handle in the process-wide `tx_registry` singleton for the life of each
1719    // `BEGIN IMMEDIATE`. Tests that observe the registry (the checkpoint
1720    // `tx_age_sweep_*` group) read `tx_registry::oldest()`; an un-serialized
1721    // spawning test here would leak a longer-lived `writer_task_tx` into that
1722    // read and make the sweep name the wrong transaction. Share the key.
1723    #[tokio::test]
1724    #[serial(tx_registry)]
1725    async fn begin_immediate_failure_replies_error_without_running_op() {
1726        // Real lock contention, not a simulation: an independent SQLite
1727        // connection holds the database write lock without taking khive's
1728        // cooperative volume lease. The writer task therefore reaches its
1729        // dedicated connection's `BEGIN IMMEDIATE` and receives SQLITE_BUSY.
1730        let dir = tempfile::tempdir().unwrap();
1731        let path = dir.path().join("writer_task_begin_failure.db");
1732        // A lock directory of its own keeps this deliberate writer hold off the
1733        // volume lease the other tests in this process share.
1734        let cfg = PoolConfig {
1735            path: Some(path.clone()),
1736            volume_lock_dir: Some(dir.path().join("volume-locks")),
1737            busy_timeout: Duration::from_millis(150),
1738            ..PoolConfig::for_test()
1739        };
1740        let pool = ConnectionPool::new(cfg).unwrap();
1741        {
1742            let writer = pool.try_writer().unwrap();
1743            writer
1744                .conn()
1745                .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY, v TEXT)")
1746                .unwrap();
1747        }
1748
1749        let handle = spawn(&pool, 8).expect("writer task should spawn on a file-backed pool");
1750
1751        let lock_holder = rusqlite::Connection::open(&path).unwrap();
1752        lock_holder.execute_batch("BEGIN IMMEDIATE").unwrap();
1753
1754        let op_ran = Arc::new(AtomicBool::new(false));
1755        let op_ran_clone = Arc::clone(&op_ran);
1756        let result = handle
1757            .send(move |conn| {
1758                op_ran_clone.store(true, Ordering::SeqCst);
1759                conn.execute("INSERT INTO t (id, v) VALUES (99, 'should-not-land')", [])
1760                    .map_err(|e| StorageError::Pool {
1761                        operation: "test_insert".into(),
1762                        message: e.to_string(),
1763                    })
1764            })
1765            .await;
1766
1767        assert!(
1768            matches!(
1769                &result,
1770                Err(StorageError::WriterTaskBusy { timeout_ms }) if *timeout_ms == 150
1771            ),
1772            "expected a typed retryable error on contended BEGIN IMMEDIATE, got {result:?}"
1773        );
1774        assert!(
1775            !op_ran.load(Ordering::SeqCst),
1776            "the request's operation closure must never run when BEGIN \
1777             IMMEDIATE fails — running it would land a partial write in \
1778             autocommit mode for a request the caller is told failed"
1779        );
1780
1781        // Release the contended lock, then verify no row landed from the
1782        // failed request.
1783        lock_holder.execute_batch("ROLLBACK").unwrap();
1784        drop(lock_holder);
1785
1786        handle
1787            .send(|conn| {
1788                conn.execute("INSERT INTO t (id, v) VALUES (100, 'next-request')", [])
1789                    .map_err(|e| StorageError::Pool {
1790                        operation: "test_insert_after_busy".into(),
1791                        message: e.to_string(),
1792                    })
1793            })
1794            .await
1795            .expect("transient contention must not retire the writer task");
1796
1797        let reader = pool.reader().expect("reader");
1798        let count: i64 = reader
1799            .conn()
1800            .query_row("SELECT COUNT(*) FROM t WHERE id IN (99, 100)", [], |row| {
1801                row.get(0)
1802            })
1803            .unwrap();
1804        assert_eq!(
1805            count, 1,
1806            "the failed request must not land, while the next request commits on the same task"
1807        );
1808    }
1809
1810    // `#[serial(tx_registry)]`: same rationale as
1811    // `begin_immediate_failure_replies_error_without_running_op`.
1812    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1813    #[serial(tx_registry)]
1814    async fn transient_begin_contention_clears_within_budget_and_op_runs_once() {
1815        // Hold a real SQLite write lock for a fixed, short delay — well
1816        // inside the configured busy_timeout — so SQLite's own internal
1817        // busy-handler wait absorbs the contention transparently within a
1818        // single `BEGIN IMMEDIATE` call. This is the common case the
1819        // shared retry budget (exercised by
1820        // `begin_retry_budget_makes_exactly_one_attempt_under_sustained_contention`
1821        // below) must not regress: recovery inside one busy_timeout window
1822        // should never touch the Rust-level retry loop or its counters.
1823        let dir = tempfile::tempdir().unwrap();
1824        let path = dir.path().join("writer_task_begin_transient_contention.db");
1825        // A lock directory of its own keeps this deliberate writer hold off the
1826        // volume lease the other tests in this process share.
1827        let pool = ConnectionPool::new(PoolConfig {
1828            path: Some(path.clone()),
1829            volume_lock_dir: Some(dir.path().join("volume-locks")),
1830            busy_timeout: Duration::from_millis(500),
1831            ..PoolConfig::for_test()
1832        })
1833        .unwrap();
1834        {
1835            let writer = pool.try_writer().unwrap();
1836            writer
1837                .conn()
1838                .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY)")
1839                .unwrap();
1840        }
1841        let handle = spawn(&pool, 8).expect("writer task spawn");
1842        // A second khive writer would contend on the cooperative volume lease
1843        // before SQLite. This fixture needs an independent SQLite writer so it
1844        // exercises the intended BEGIN busy-handler path.
1845        let lock_holder = Connection::open(&path).unwrap();
1846        lock_holder.execute_batch("BEGIN IMMEDIATE").unwrap();
1847
1848        let op_runs = Arc::new(AtomicUsize::new(0));
1849        let op_runs_in_request = Arc::clone(&op_runs);
1850        let send_future = handle.send(move |conn| {
1851            op_runs_in_request.fetch_add(1, Ordering::SeqCst);
1852            conn.execute("INSERT INTO t (id) VALUES (1)", [])
1853                .map_err(|error| StorageError::Pool {
1854                    operation: "test_insert_after_transient_contention".into(),
1855                    message: error.to_string(),
1856                })
1857        });
1858        let release_future = async {
1859            tokio::time::sleep(Duration::from_millis(50)).await;
1860            lock_holder.execute_batch("ROLLBACK").unwrap();
1861        };
1862        let (result, ()) = tokio::join!(send_future, release_future);
1863
1864        assert_eq!(
1865            result.expect("BEGIN IMMEDIATE succeeds once the transient lock clears"),
1866            1
1867        );
1868        assert_eq!(
1869            op_runs.load(Ordering::SeqCst),
1870            1,
1871            "the FnOnce request closure must execute exactly once"
1872        );
1873
1874        let settled = pool.writer_acquisition_snapshot();
1875        assert_eq!(
1876            settled.writer_task_begin_busy, 0,
1877            "contention absorbed inside SQLite's own busy_timeout wait must never \
1878             surface as a Rust-level refusal"
1879        );
1880        assert_eq!(settled.writer_task_begin_busy_absorbed, 0);
1881        let reader = pool.reader().unwrap();
1882        let rows: i64 = reader
1883            .conn()
1884            .query_row("SELECT COUNT(*) FROM t", [], |row| row.get(0))
1885            .unwrap();
1886        assert_eq!(rows, 1, "exactly one closure execution commits one row");
1887    }
1888
1889    include!("writer_task_admission_retry_tests.rs");
1890
1891    // `#[serial(tx_registry)]`: shares the key with the checkpoint
1892    // `tx_age_sweep_*` tests — see the note on
1893    // `begin_immediate_failure_replies_error_without_running_op`.
1894    #[tokio::test]
1895    #[serial(tx_registry)]
1896    async fn writer_task_executes_op_and_commits() {
1897        let dir = tempfile::tempdir().unwrap();
1898        let path = dir.path().join("writer_task_commit.db");
1899        let pool = file_pool(&path);
1900        {
1901            let writer = pool.try_writer().unwrap();
1902            writer
1903                .conn()
1904                .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY, v TEXT)")
1905                .unwrap();
1906        }
1907
1908        let handle = spawn(&pool, 8).expect("writer task should spawn on a file-backed pool");
1909
1910        let affected = handle
1911            .send(|conn| {
1912                conn.execute("INSERT INTO t (id, v) VALUES (1, 'hello')", [])
1913                    .map_err(|e| StorageError::Pool {
1914                        operation: "test_insert".into(),
1915                        message: e.to_string(),
1916                    })
1917            })
1918            .await
1919            .expect("op should succeed");
1920        assert_eq!(affected, 1);
1921
1922        // Verify the write actually committed to the shared file — read it
1923        // back via a fresh pooled reader connection, not the writer task's
1924        // own connection.
1925        let reader = pool.reader().expect("reader");
1926        let v: String = reader
1927            .conn()
1928            .query_row("SELECT v FROM t WHERE id = 1", [], |row| row.get(0))
1929            .expect("row must be committed and visible to a reader");
1930        assert_eq!(v, "hello");
1931
1932        let counters = pool.writer_acquisition_snapshot();
1933        assert_eq!(counters.acquisitions, 2);
1934        assert_eq!(counters.pooled_acquisitions, 1);
1935        assert_eq!(counters.standalone_acquisitions, 0);
1936        assert_eq!(counters.writer_task_acquisitions, 1);
1937        assert_eq!(counters.timeouts, 0);
1938    }
1939
1940    #[tokio::test]
1941    async fn writer_task_connection_follows_checkpoint_ownership_claim() {
1942        let dir = tempfile::tempdir().unwrap();
1943        let path = dir.path().join("writer_task_autocheckpoint.db");
1944        let pool = file_pool(&path);
1945        // The pool-registered handle (not a directly spawned side task):
1946        // `propagate_checkpoint_claim_to_writer_task` reaches exactly this
1947        // task's connection.
1948        let handle = pool
1949            .writer_task_handle()
1950            .expect("writer task should spawn")
1951            .expect("file-backed pool resolves the write queue on");
1952
1953        let read_pages = |handle: &WriterTaskHandle| {
1954            let handle = handle.clone();
1955            async move {
1956                handle
1957                    .send_top_level(|conn| {
1958                        conn.pragma_query_value(None, "wal_autocheckpoint", |row| {
1959                            row.get::<_, u32>(0)
1960                        })
1961                        .map_err(|e| StorageError::Pool {
1962                            operation: "test_wal_autocheckpoint".into(),
1963                            message: e.to_string(),
1964                        })
1965                    })
1966                    .await
1967                    .expect("query writer-task connection pragma")
1968            }
1969        };
1970
1971        // Spawned before any claim: the task's long-lived connection keeps
1972        // the bounded fallback.
1973        assert_eq!(
1974            read_pages(&handle).await,
1975            crate::pool::FALLBACK_WAL_AUTOCHECKPOINT_PAGES
1976        );
1977
1978        // After the claim propagates, the same connection is flipped to the
1979        // dedicated-owner setting without reopening.
1980        pool.claim_checkpoint_ownership().expect("claim ownership");
1981        pool.propagate_checkpoint_claim_to_writer_task()
1982            .await
1983            .expect("propagate claim to the running writer task");
1984        assert_eq!(read_pages(&handle).await, 0);
1985    }
1986
1987    #[test]
1988    fn spawn_fails_on_in_memory_pool() {
1989        // In-memory pools have no standalone-connection support
1990        // (the infrastructure-only standalone open) — `spawn` must surface
1991        // that as an error rather than panicking. Deliberately a plain
1992        // `#[test]` (no Tokio runtime): `spawn` fails before it ever reaches
1993        // `tokio::spawn`, so no runtime is required for this path.
1994        let cfg = PoolConfig {
1995            path: None,
1996            ..PoolConfig::default()
1997        };
1998        let pool = ConnectionPool::new(cfg).unwrap();
1999        let result = spawn(&pool, 8);
2000        assert!(
2001            result.is_err(),
2002            "in-memory pools must reject spawn, not panic"
2003        );
2004    }
2005
2006    include!("writer_task_queue_capacity_tests.rs");
2007
2008    #[tokio::test]
2009    #[serial(tx_registry)]
2010    async fn operation_failure_with_successful_rollback_reports_finality_once_and_continues() {
2011        let dir = tempfile::tempdir().unwrap();
2012        let path = dir.path().join("writer_task_operation_rollback.db");
2013        let pool = file_pool(&path);
2014        {
2015            let writer = pool.try_writer().unwrap();
2016            writer
2017                .conn()
2018                .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY, v TEXT)")
2019                .unwrap();
2020        }
2021        let handle = spawn(&pool, 8).expect("writer task spawn");
2022
2023        let executions = Arc::new(AtomicUsize::new(0));
2024        let executions_in_op = Arc::clone(&executions);
2025        let original_error = handle
2026            .send(move |conn| -> Result<(), StorageError> {
2027                executions_in_op.fetch_add(1, Ordering::SeqCst);
2028                conn.execute("INSERT INTO t (id, v) VALUES (1, 'rolled-back')", [])
2029                    .map_err(|e| StorageError::Pool {
2030                        operation: "test_operation_error_insert".into(),
2031                        message: e.to_string(),
2032                    })?;
2033                Err(StorageError::Internal(
2034                    "intentional operation failure".into(),
2035                ))
2036            })
2037            .await;
2038        match &original_error {
2039            Err(StorageError::WriterTaskRequestFailed {
2040                request_state: WriterTaskRequestState::TransactionRolledBack,
2041                source,
2042            }) => assert!(
2043                matches!(source.as_ref(), StorageError::Internal(message)
2044                    if message == "intentional operation failure"),
2045                "the proven-rollback wrapper must retain the typed operation error: {source:?}"
2046            ),
2047            other => panic!(
2048                "a confirmed rollback must carry TransactionRolledBack and preserve the operation error, got {other:?}"
2049            ),
2050        }
2051        assert_eq!(
2052            executions.load(Ordering::SeqCst),
2053            1,
2054            "finality propagation must not replay the request closure"
2055        );
2056
2057        let affected = handle
2058            .send(|conn| {
2059                conn.execute("INSERT INTO t (id, v) VALUES (2, 'committed')", [])
2060                    .map_err(|e| StorageError::Pool {
2061                        operation: "test_operation_error_followup_insert".into(),
2062                        message: e.to_string(),
2063                    })
2064            })
2065            .await
2066            .expect("the writer must continue after a confirmed rollback");
2067        assert_eq!(affected, 1);
2068        assert_eq!(
2069            executions.load(Ordering::SeqCst),
2070            1,
2071            "serving a follow-up request must not replay the rolled-back closure"
2072        );
2073
2074        let reader = pool.reader().expect("reader");
2075        let rolled_back: i64 = reader
2076            .conn()
2077            .query_row("SELECT COUNT(*) FROM t WHERE id = 1", [], |row| row.get(0))
2078            .unwrap();
2079        let committed: i64 = reader
2080            .conn()
2081            .query_row("SELECT COUNT(*) FROM t WHERE id = 2", [], |row| row.get(0))
2082            .unwrap();
2083        assert_eq!(rolled_back, 0);
2084        assert_eq!(committed, 1);
2085    }
2086
2087    #[tokio::test]
2088    #[serial(tx_registry)]
2089    async fn commit_failure_with_successful_rollback_reports_finality_once_and_continues() {
2090        let dir = tempfile::tempdir().unwrap();
2091        let path = dir.path().join("writer_task_commit_rollback.db");
2092        let pool = file_pool(&path);
2093        {
2094            let writer = pool.try_writer().unwrap();
2095            writer
2096                .conn()
2097                .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY, v TEXT)")
2098                .unwrap();
2099        }
2100        let handle = spawn(&pool, 8).expect("writer task spawn");
2101
2102        let executions = Arc::new(AtomicUsize::new(0));
2103        let executions_in_op = Arc::clone(&executions);
2104        let commit_error = handle
2105            .send(move |conn| -> Result<usize, StorageError> {
2106                executions_in_op.fetch_add(1, Ordering::SeqCst);
2107                let affected = conn
2108                    .execute("INSERT INTO t (id, v) VALUES (1, 'rolled-back')", [])
2109                    .map_err(|e| StorageError::Pool {
2110                        operation: "test_commit_error_insert".into(),
2111                        message: e.to_string(),
2112                    })?;
2113                conn.authorizer(Some(deny_commit))
2114                    .map_err(|e| StorageError::Pool {
2115                        operation: "test_install_authorizer".into(),
2116                        message: e.to_string(),
2117                    })?;
2118                Ok(affected)
2119            })
2120            .await;
2121        match &commit_error {
2122            Err(StorageError::WriterTaskRequestFailed {
2123                request_state: WriterTaskRequestState::TransactionRolledBack,
2124                source,
2125            }) => assert!(
2126                matches!(source.as_ref(), StorageError::Pool { operation, .. }
2127                    if operation == "writer_task_commit"),
2128                "the proven-rollback wrapper must retain the typed COMMIT error: {source:?}"
2129            ),
2130            other => panic!(
2131                "a confirmed rollback must carry TransactionRolledBack and preserve the COMMIT error, got {other:?}"
2132            ),
2133        }
2134        assert!(
2135            commit_error
2136                .as_ref()
2137                .expect_err("COMMIT must be denied")
2138                .is_retryable(),
2139            "the existing retryable commit-error contract must remain unchanged after a \
2140             confirmed rollback"
2141        );
2142        assert_eq!(
2143            executions.load(Ordering::SeqCst),
2144            1,
2145            "finality propagation must not replay the request closure"
2146        );
2147
2148        let affected = handle
2149            .send(|conn| {
2150                conn.authorizer(None::<fn(AuthContext<'_>) -> Authorization>)
2151                    .map_err(|e| StorageError::Pool {
2152                        operation: "test_remove_authorizer".into(),
2153                        message: e.to_string(),
2154                    })?;
2155                conn.execute("INSERT INTO t (id, v) VALUES (2, 'committed')", [])
2156                    .map_err(|e| StorageError::Pool {
2157                        operation: "test_commit_error_followup_insert".into(),
2158                        message: e.to_string(),
2159                    })
2160            })
2161            .await
2162            .expect("the writer must continue after the failed COMMIT is rolled back");
2163        assert_eq!(affected, 1);
2164        assert_eq!(
2165            executions.load(Ordering::SeqCst),
2166            1,
2167            "serving a follow-up request must not replay the rolled-back closure"
2168        );
2169
2170        let reader = pool.reader().expect("reader");
2171        let rolled_back: i64 = reader
2172            .conn()
2173            .query_row("SELECT COUNT(*) FROM t WHERE id = 1", [], |row| row.get(0))
2174            .unwrap();
2175        let committed: i64 = reader
2176            .conn()
2177            .query_row("SELECT COUNT(*) FROM t WHERE id = 2", [], |row| row.get(0))
2178            .unwrap();
2179        assert_eq!(rolled_back, 0);
2180        assert_eq!(committed, 1);
2181    }
2182
2183    #[test]
2184    fn top_level_request_returning_with_open_transaction_reports_side_effects_unknown() {
2185        let conn = Connection::open_in_memory().expect("in-memory connection");
2186        conn.execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY)")
2187            .unwrap();
2188        let (reply_tx, mut reply_rx) = oneshot::channel();
2189        let request = WriteRequest {
2190            op: Box::new(|conn| -> Result<usize, StorageError> {
2191                conn.execute_batch("BEGIN IMMEDIATE")
2192                    .map_err(|e| StorageError::Pool {
2193                        operation: "test_top_level_begin".into(),
2194                        message: e.to_string(),
2195                    })?;
2196                conn.execute("INSERT INTO t (id) VALUES (1)", [])
2197                    .map_err(|e| StorageError::Pool {
2198                        operation: "test_top_level_insert".into(),
2199                        message: e.to_string(),
2200                    })
2201            }),
2202            reply: reply_tx,
2203            top_level: true,
2204            checkpoint_bypass: false,
2205            vacuum_copy_headroom: false,
2206            telemetry: WriteTelemetry::new(None, "test".to_string(), 0, None),
2207        };
2208
2209        let terminal_state = sealed::Sealed::execute_and_reply_top_level_reporting_terminal(
2210            Box::new(request),
2211            &conn,
2212            Duration::ZERO,
2213        );
2214        assert_eq!(
2215            terminal_state,
2216            Some(WriterTaskRequestState::SideEffectsUnknown)
2217        );
2218        let reply = reply_rx
2219            .try_recv()
2220            .expect("active request must receive a typed terminal reply");
2221        assert_writer_task_terminal_state(reply, WriterTaskRequestState::SideEffectsUnknown);
2222        assert!(
2223            !conn.is_autocommit(),
2224            "the fixture must prove the post-request autocommit check observed an open transaction"
2225        );
2226    }
2227
2228    #[test]
2229    fn commit_failure_with_failed_rollback_reports_side_effects_unknown() {
2230        let conn = Connection::open_in_memory().expect("in-memory connection");
2231        conn.execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY); BEGIN IMMEDIATE")
2232            .unwrap();
2233        let executions = Arc::new(AtomicUsize::new(0));
2234        let executions_in_op = Arc::clone(&executions);
2235        let (reply_tx, mut reply_rx) = oneshot::channel();
2236        let request = WriteRequest {
2237            op: Box::new(move |conn| -> Result<usize, StorageError> {
2238                executions_in_op.fetch_add(1, Ordering::SeqCst);
2239                let affected = conn
2240                    .execute("INSERT INTO t (id) VALUES (1)", [])
2241                    .map_err(|e| StorageError::Pool {
2242                        operation: "test_insert_before_commit_failure".into(),
2243                        message: e.to_string(),
2244                    })?;
2245                conn.authorizer(Some(deny_commit_and_rollback))
2246                    .map_err(|e| StorageError::Pool {
2247                        operation: "test_install_authorizer".into(),
2248                        message: e.to_string(),
2249                    })?;
2250                Ok(affected)
2251            }),
2252            reply: reply_tx,
2253            top_level: false,
2254            checkpoint_bypass: false,
2255            vacuum_copy_headroom: false,
2256            telemetry: WriteTelemetry::new(None, "test".to_string(), 0, None),
2257        };
2258
2259        let terminal_state = sealed::Sealed::execute_and_reply_reporting_terminal(
2260            Box::new(request),
2261            &conn,
2262            None,
2263            Duration::ZERO,
2264            Duration::ZERO,
2265        );
2266        assert_eq!(
2267            terminal_state,
2268            Some(WriterTaskRequestState::SideEffectsUnknown)
2269        );
2270        let reply = reply_rx
2271            .try_recv()
2272            .expect("active request must receive a typed terminal reply");
2273        assert_writer_task_terminal_state(reply, WriterTaskRequestState::SideEffectsUnknown);
2274        assert_eq!(executions.load(Ordering::SeqCst), 1);
2275        assert!(
2276            !conn.is_autocommit(),
2277            "the denied COMMIT and ROLLBACK must leave the test connection poisoned"
2278        );
2279    }
2280
2281    #[tokio::test]
2282    #[serial(tx_registry)]
2283    async fn poisoned_connection_retires_before_queued_top_level_request() {
2284        let dir = tempfile::tempdir().unwrap();
2285        let path = dir.path().join("writer_task_rollback_poison.db");
2286        let pool = file_pool(&path);
2287        {
2288            let writer = pool.try_writer().unwrap();
2289            writer
2290                .conn()
2291                .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY, v TEXT)")
2292                .unwrap();
2293        }
2294        let handle = spawn(&pool, 8).expect("writer task spawn");
2295        let (started_tx, started_rx) = oneshot::channel::<()>();
2296        let (release_tx, release_rx) = std_mpsc::channel::<()>();
2297
2298        let active = tokio::spawn({
2299            let handle = handle.clone();
2300            async move {
2301                handle
2302                    .send(move |conn| -> Result<usize, StorageError> {
2303                        let affected = conn
2304                            .execute("INSERT INTO t (id, v) VALUES (1, 'active')", [])
2305                            .map_err(|e| StorageError::Pool {
2306                                operation: "test_active_insert".into(),
2307                                message: e.to_string(),
2308                            })?;
2309                        conn.authorizer(Some(deny_commit_and_rollback))
2310                            .map_err(|e| StorageError::Pool {
2311                                operation: "test_install_authorizer".into(),
2312                                message: e.to_string(),
2313                            })?;
2314                        let _ = started_tx.send(());
2315                        release_rx.recv().expect("test must release active op");
2316                        Ok(affected)
2317                    })
2318                    .await
2319            }
2320        });
2321
2322        tokio::time::timeout(Duration::from_secs(5), started_rx)
2323            .await
2324            .expect("active request did not start")
2325            .expect("active request dropped its start signal");
2326
2327        let queued_ran = Arc::new(AtomicBool::new(false));
2328        let queued_ran_in_op = Arc::clone(&queued_ran);
2329        let queued_top_level = handle
2330            .enqueue_inner(
2331                move |conn| {
2332                    queued_ran_in_op.store(true, Ordering::SeqCst);
2333                    conn.execute("INSERT INTO t (id, v) VALUES (2, 'queued')", [])
2334                        .map_err(|e| StorageError::Pool {
2335                            operation: "test_queued_top_level_insert".into(),
2336                            message: e.to_string(),
2337                        })
2338                },
2339                true,
2340                false,
2341                false,
2342            )
2343            .await
2344            .expect("top-level request must queue behind active request");
2345        release_tx.send(()).expect("release active op");
2346
2347        let active_result = tokio::time::timeout(Duration::from_secs(5), active)
2348            .await
2349            .expect("active caller hung after rollback failure")
2350            .expect("active caller task join");
2351        assert_writer_task_terminal_state(
2352            active_result,
2353            WriterTaskRequestState::SideEffectsUnknown,
2354        );
2355
2356        let queued_result = tokio::time::timeout(Duration::from_secs(5), queued_top_level)
2357            .await
2358            .expect("queued top-level caller hung after terminal failure")
2359            .expect("terminal drain must preserve queued typed reply");
2360        assert_writer_task_terminal_state(queued_result, WriterTaskRequestState::NotStarted);
2361        assert!(
2362            !queued_ran.load(Ordering::SeqCst),
2363            "a top-level request must never run on the poisoned connection"
2364        );
2365
2366        let future_ran = Arc::new(AtomicBool::new(false));
2367        let future_ran_in_op = Arc::clone(&future_ran);
2368        let future_result = handle
2369            .send_top_level(move |_conn| {
2370                future_ran_in_op.store(true, Ordering::SeqCst);
2371                Ok::<(), StorageError>(())
2372            })
2373            .await;
2374        assert_writer_task_terminal_state(future_result, WriterTaskRequestState::NotStarted);
2375        assert!(!future_ran.load(Ordering::SeqCst));
2376    }
2377
2378    #[test]
2379    fn operation_failure_with_failed_rollback_reports_side_effects_unknown() {
2380        let conn = Connection::open_in_memory().expect("in-memory connection");
2381        conn.execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY); BEGIN IMMEDIATE")
2382            .unwrap();
2383        let executions = Arc::new(AtomicUsize::new(0));
2384        let executions_in_op = Arc::clone(&executions);
2385        let (reply_tx, mut reply_rx) = oneshot::channel();
2386        let request = WriteRequest {
2387            op: Box::new(move |conn| -> Result<(), StorageError> {
2388                executions_in_op.fetch_add(1, Ordering::SeqCst);
2389                conn.authorizer(Some(deny_rollback))
2390                    .map_err(|e| StorageError::Pool {
2391                        operation: "test_install_authorizer".into(),
2392                        message: e.to_string(),
2393                    })?;
2394                Err(StorageError::Internal(
2395                    "intentional operation failure before denied rollback".into(),
2396                ))
2397            }),
2398            reply: reply_tx,
2399            top_level: false,
2400            checkpoint_bypass: false,
2401            vacuum_copy_headroom: false,
2402            telemetry: WriteTelemetry::new(None, "test".to_string(), 0, None),
2403        };
2404
2405        let terminal_state = sealed::Sealed::execute_and_reply_reporting_terminal(
2406            Box::new(request),
2407            &conn,
2408            None,
2409            Duration::ZERO,
2410            Duration::ZERO,
2411        );
2412        assert_eq!(
2413            terminal_state,
2414            Some(WriterTaskRequestState::SideEffectsUnknown)
2415        );
2416        let reply = reply_rx
2417            .try_recv()
2418            .expect("active request must receive a typed terminal reply");
2419        assert_writer_task_terminal_state(reply, WriterTaskRequestState::SideEffectsUnknown);
2420        assert_eq!(executions.load(Ordering::SeqCst), 1);
2421        assert!(
2422            !conn.is_autocommit(),
2423            "the denied ROLLBACK must leave the test connection poisoned"
2424        );
2425    }
2426
2427    #[test]
2428    fn wrapped_panic_with_failed_rollback_reports_side_effects_unknown() {
2429        let conn = Connection::open_in_memory().expect("in-memory connection");
2430        conn.execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY); BEGIN IMMEDIATE")
2431            .unwrap();
2432        let (reply_tx, mut reply_rx) = oneshot::channel();
2433        let request = WriteRequest {
2434            op: Box::new(|conn| -> Result<(), StorageError> {
2435                // Deliberately violate WriteOp's no-COMMIT contract to create
2436                // the otherwise rare but reachable rollback-failure state.
2437                // The row commits before the panic, so reporting a successful
2438                // rollback here would be dangerously false.
2439                conn.execute_batch("INSERT INTO t (id) VALUES (1); COMMIT")
2440                    .map_err(|e| StorageError::Pool {
2441                        operation: "test_force_rollback_failure".into(),
2442                        message: e.to_string(),
2443                    })?;
2444                panic!("intentional panic after illicit commit");
2445            }),
2446            reply: reply_tx,
2447            top_level: false,
2448            checkpoint_bypass: false,
2449            vacuum_copy_headroom: false,
2450            telemetry: WriteTelemetry::new(None, "test".to_string(), 0, None),
2451        };
2452
2453        let terminal_state = sealed::Sealed::execute_and_reply_reporting_terminal(
2454            Box::new(request),
2455            &conn,
2456            None,
2457            Duration::ZERO,
2458            Duration::ZERO,
2459        );
2460        assert_eq!(
2461            terminal_state,
2462            Some(WriterTaskRequestState::SideEffectsUnknown)
2463        );
2464        let reply = reply_rx
2465            .try_recv()
2466            .expect("active request must receive a typed terminal reply");
2467        assert_writer_task_terminal_state(reply, WriterTaskRequestState::SideEffectsUnknown);
2468
2469        let count: i64 = conn
2470            .query_row("SELECT COUNT(*) FROM t", [], |row| row.get(0))
2471            .unwrap();
2472        assert_eq!(
2473            count, 1,
2474            "the fixture's committed side effect proves why the state must be unknown"
2475        );
2476    }
2477
2478    // `#[serial(tx_registry)]`: the active request holds a registered
2479    // `writer_task_tx` until its panic is caught and rolled back. Share the
2480    // process-wide registry key with the other writer-task transaction tests.
2481    #[tokio::test]
2482    #[serial(tx_registry)]
2483    async fn wrapped_panic_rolls_back_and_terminally_fails_queue() {
2484        let dir = tempfile::tempdir().unwrap();
2485        let path = dir.path().join("writer_task_wrapped_panic.db");
2486        let cfg = PoolConfig {
2487            path: Some(path),
2488            write_queue_enabled: Some(true),
2489            write_queue_capacity: 8,
2490            ..PoolConfig::for_test()
2491        };
2492        let pool = ConnectionPool::new(cfg).unwrap();
2493        {
2494            let writer = pool.try_writer().unwrap();
2495            writer
2496                .conn()
2497                .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY, v TEXT)")
2498                .unwrap();
2499        }
2500
2501        // Use the pool-owned OnceLock path rather than `spawn` directly so
2502        // the test also proves a terminal task is never silently replaced.
2503        let handle = pool
2504            .writer_task_handle()
2505            .expect("writer task lookup")
2506            .expect("file-backed queued pool must spawn its writer task");
2507        assert_eq!(pool.writer_task_spawn_count(), 1);
2508
2509        let (started_tx, started_rx) = oneshot::channel::<()>();
2510        let (release_tx, release_rx) = std_mpsc::channel::<()>();
2511        let active = tokio::spawn({
2512            let handle = handle.clone();
2513            async move {
2514                handle
2515                    .send(move |conn| -> Result<usize, StorageError> {
2516                        conn.execute("INSERT INTO t (id, v) VALUES (1, 'active')", [])
2517                            .map_err(|e| StorageError::Pool {
2518                                operation: "test_active_insert".into(),
2519                                message: e.to_string(),
2520                            })?;
2521                        let _ = started_tx.send(());
2522                        release_rx.recv().expect("test must release active op");
2523                        panic!("intentional wrapped writer request panic");
2524                    })
2525                    .await
2526            }
2527        });
2528
2529        tokio::time::timeout(Duration::from_secs(5), started_rx)
2530            .await
2531            .expect("active request did not start")
2532            .expect("active request dropped its start signal");
2533
2534        let queued_one_ran = Arc::new(AtomicBool::new(false));
2535        let queued_one_ran_in_op = Arc::clone(&queued_one_ran);
2536        let queued_one = handle
2537            .enqueue(move |conn| {
2538                queued_one_ran_in_op.store(true, Ordering::SeqCst);
2539                conn.execute("INSERT INTO t (id, v) VALUES (2, 'queued-one')", [])
2540                    .map_err(|e| StorageError::Pool {
2541                        operation: "test_queued_one_insert".into(),
2542                        message: e.to_string(),
2543                    })
2544            })
2545            .await
2546            .expect("first queued request must be accepted");
2547
2548        let queued_two_ran = Arc::new(AtomicBool::new(false));
2549        let queued_two_ran_in_op = Arc::clone(&queued_two_ran);
2550        let queued_two = handle
2551            .enqueue(move |_conn| {
2552                queued_two_ran_in_op.store(true, Ordering::SeqCst);
2553                Ok::<String, StorageError>("queued-two-ran".to_string())
2554            })
2555            .await
2556            .expect("second queued request must be accepted");
2557
2558        assert_eq!(
2559            handle.queue_depth(),
2560            2,
2561            "both heterogeneous requests must be buffered behind the active op"
2562        );
2563        release_tx.send(()).expect("release active op");
2564
2565        let active_result = tokio::time::timeout(Duration::from_secs(5), active)
2566            .await
2567            .expect("active caller hung after panic")
2568            .expect("active caller task join");
2569        assert_writer_task_terminal_state(
2570            active_result,
2571            WriterTaskRequestState::TransactionRolledBack,
2572        );
2573
2574        let queued_one_result = tokio::time::timeout(Duration::from_secs(5), queued_one)
2575            .await
2576            .expect("first queued caller hung after terminal failure")
2577            .expect("terminal drain must preserve first typed reply");
2578        assert_writer_task_terminal_state(queued_one_result, WriterTaskRequestState::NotStarted);
2579
2580        let queued_two_result = tokio::time::timeout(Duration::from_secs(5), queued_two)
2581            .await
2582            .expect("second queued caller hung after terminal failure")
2583            .expect("terminal drain must preserve second typed reply");
2584        assert_writer_task_terminal_state(queued_two_result, WriterTaskRequestState::NotStarted);
2585        assert!(!queued_one_ran.load(Ordering::SeqCst));
2586        assert!(!queued_two_ran.load(Ordering::SeqCst));
2587
2588        let future_ran = Arc::new(AtomicBool::new(false));
2589        let future_ran_in_op = Arc::clone(&future_ran);
2590        let future_result = handle
2591            .send(move |_conn| {
2592                future_ran_in_op.store(true, Ordering::SeqCst);
2593                Ok::<(), StorageError>(())
2594            })
2595            .await;
2596        assert_writer_task_terminal_state(future_result, WriterTaskRequestState::NotStarted);
2597        assert!(!future_ran.load(Ordering::SeqCst));
2598
2599        let cached_after_failure = pool
2600            .writer_task_handle()
2601            .expect("cached writer task lookup")
2602            .expect("pool retains its terminal handle");
2603        assert_eq!(
2604            pool.writer_task_spawn_count(),
2605            1,
2606            "a terminal writer task must not be restarted behind callers' backs"
2607        );
2608        let cached_result = cached_after_failure
2609            .send(|_conn| Ok::<(), StorageError>(()))
2610            .await;
2611        assert_writer_task_terminal_state(cached_result, WriterTaskRequestState::NotStarted);
2612
2613        let reader = pool.reader().expect("reader");
2614        let count: i64 = reader
2615            .conn()
2616            .query_row("SELECT COUNT(*) FROM t", [], |row| row.get(0))
2617            .unwrap();
2618        assert_eq!(
2619            count, 0,
2620            "the active transaction must be rolled back and queued ops must never run"
2621        );
2622    }
2623
2624    #[tokio::test]
2625    async fn top_level_panic_reports_unknown_and_fails_queue_without_running_it() {
2626        let dir = tempfile::tempdir().unwrap();
2627        let path = dir.path().join("writer_task_top_level_panic.db");
2628        let pool = file_pool(&path);
2629        {
2630            let writer = pool.try_writer().unwrap();
2631            writer
2632                .conn()
2633                .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY, v TEXT)")
2634                .unwrap();
2635        }
2636        let handle = spawn(&pool, 8).expect("writer task spawn");
2637
2638        let (started_tx, started_rx) = oneshot::channel::<()>();
2639        let (release_tx, release_rx) = std_mpsc::channel::<()>();
2640        let active = tokio::spawn({
2641            let handle = handle.clone();
2642            async move {
2643                handle
2644                    .send_top_level(move |conn| -> Result<usize, StorageError> {
2645                        conn.execute("INSERT INTO t (id, v) VALUES (10, 'autocommitted')", [])
2646                            .map_err(|e| StorageError::Pool {
2647                                operation: "test_top_level_insert".into(),
2648                                message: e.to_string(),
2649                            })?;
2650                        let _ = started_tx.send(());
2651                        release_rx.recv().expect("test must release top-level op");
2652                        panic!("intentional top-level writer request panic");
2653                    })
2654                    .await
2655            }
2656        });
2657
2658        tokio::time::timeout(Duration::from_secs(5), started_rx)
2659            .await
2660            .expect("top-level request did not start")
2661            .expect("top-level request dropped its start signal");
2662
2663        let queued_ran = Arc::new(AtomicBool::new(false));
2664        let queued_ran_in_op = Arc::clone(&queued_ran);
2665        let queued = handle
2666            .enqueue(move |conn| {
2667                queued_ran_in_op.store(true, Ordering::SeqCst);
2668                conn.execute("INSERT INTO t (id, v) VALUES (11, 'queued')", [])
2669                    .map_err(|e| StorageError::Pool {
2670                        operation: "test_top_level_queued_insert".into(),
2671                        message: e.to_string(),
2672                    })
2673            })
2674            .await
2675            .expect("queued request must be accepted");
2676        assert_eq!(handle.queue_depth(), 1);
2677        release_tx.send(()).expect("release top-level op");
2678
2679        let active_result = tokio::time::timeout(Duration::from_secs(5), active)
2680            .await
2681            .expect("top-level caller hung after panic")
2682            .expect("top-level caller task join");
2683        assert_writer_task_terminal_state(
2684            active_result,
2685            WriterTaskRequestState::SideEffectsUnknown,
2686        );
2687
2688        let queued_result = tokio::time::timeout(Duration::from_secs(5), queued)
2689            .await
2690            .expect("queued caller hung after top-level panic")
2691            .expect("terminal drain must preserve queued typed reply");
2692        assert_writer_task_terminal_state(queued_result, WriterTaskRequestState::NotStarted);
2693        assert!(!queued_ran.load(Ordering::SeqCst));
2694
2695        let reader = pool.reader().expect("reader");
2696        let active_count: i64 = reader
2697            .conn()
2698            .query_row("SELECT COUNT(*) FROM t WHERE id = 10", [], |row| row.get(0))
2699            .unwrap();
2700        let queued_count: i64 = reader
2701            .conn()
2702            .query_row("SELECT COUNT(*) FROM t WHERE id = 11", [], |row| row.get(0))
2703            .unwrap();
2704        assert_eq!(
2705            active_count, 1,
2706            "the completed top-level statement autocommits before the panic"
2707        );
2708        assert_eq!(queued_count, 0, "the queued request must never run");
2709    }
2710
2711    #[tokio::test]
2712    async fn closed_receiver_rejects_all_send_surfaces_as_not_started() {
2713        // Simulates the writer task having terminated: its `rx` is gone, so
2714        // every send surface must fail deterministically without running the
2715        // supplied operation.
2716        let (tx, rx) = mpsc::channel::<Box<dyn AnyWriteRequest + Send>>(4);
2717        drop(rx);
2718
2719        let handle = WriterTaskHandle {
2720            tx,
2721            backend_key: None,
2722            db: "test".to_string(),
2723            slow_write_threshold: None,
2724            enqueue_timeout: Duration::from_secs(5),
2725        };
2726        let send_result = handle.send(|_conn| Ok::<(), StorageError>(())).await;
2727        assert_writer_task_terminal_state(send_result, WriterTaskRequestState::NotStarted);
2728
2729        let timed_result = handle
2730            .send_with_timeout(|_conn| Ok::<(), StorageError>(()), Duration::from_secs(1))
2731            .await;
2732        assert_writer_task_terminal_state(timed_result, WriterTaskRequestState::NotStarted);
2733
2734        let top_level_result = handle
2735            .send_top_level(|_conn| Ok::<(), StorageError>(()))
2736            .await;
2737        assert_writer_task_terminal_state(top_level_result, WriterTaskRequestState::NotStarted);
2738    }
2739
2740    #[tokio::test]
2741    async fn accepted_request_lost_reply_is_side_effects_unknown() {
2742        // A request can be accepted before an unexpected receiver/task loss.
2743        // The handle cannot prove whether its closure ran, so the oneshot
2744        // fallback must conservatively report SideEffectsUnknown.
2745        let (tx, mut rx) = mpsc::channel::<Box<dyn AnyWriteRequest + Send>>(1);
2746        let handle = WriterTaskHandle {
2747            tx,
2748            backend_key: None,
2749            db: "test".to_string(),
2750            slow_write_threshold: None,
2751            enqueue_timeout: Duration::from_secs(5),
2752        };
2753        let request_ran = Arc::new(AtomicBool::new(false));
2754        let request_ran_in_op = Arc::clone(&request_ran);
2755
2756        let dropper = tokio::spawn(async move {
2757            let request = rx.recv().await.expect("request must be accepted");
2758            drop(request);
2759        });
2760        let result = tokio::time::timeout(
2761            Duration::from_secs(5),
2762            handle.send(move |_conn| {
2763                request_ran_in_op.store(true, Ordering::SeqCst);
2764                Ok::<(), StorageError>(())
2765            }),
2766        )
2767        .await
2768        .expect("caller hung after accepted request was dropped");
2769        dropper.await.expect("dropper task join");
2770
2771        assert_writer_task_terminal_state(result, WriterTaskRequestState::SideEffectsUnknown);
2772        assert!(!request_ran.load(Ordering::SeqCst));
2773    }
2774
2775    /// #1849: writer telemetry uses the same canonical OS-path identity as
2776    /// the pool. A lossy display string would merge these distinct backends.
2777    #[cfg(unix)]
2778    #[test]
2779    fn writer_stage_backend_key_preserves_non_utf8_path_bytes() {
2780        use std::ffi::OsString;
2781        use std::os::unix::ffi::OsStringExt;
2782
2783        let path_a =
2784            std::path::PathBuf::from(OsString::from_vec(b"/tmp/khive-writer-\x80.db".to_vec()));
2785        let path_b =
2786            std::path::PathBuf::from(OsString::from_vec(b"/tmp/khive-writer-\x81.db".to_vec()));
2787        assert_eq!(
2788            path_a.display().to_string(),
2789            path_b.display().to_string(),
2790            "fixture must reproduce the lossy display-label collision"
2791        );
2792        assert_ne!(
2793            writer_db_key_from_path(Some(&path_a)),
2794            writer_db_key_from_path(Some(&path_b)),
2795            "backend keys must retain the canonical path's exact OS bytes"
2796        );
2797    }
2798
2799    /// #1849: a deliberately slow transaction body with no queue backlog or
2800    /// lock contention must be attributed to `body`, not flattened into one
2801    /// opaque send-to-reply duration.
2802    #[tokio::test]
2803    async fn writer_stage_sample_attributes_a_slow_body() {
2804        let dir = tempfile::tempdir().unwrap();
2805        let path = dir.path().join("writer_stage_sample.db");
2806        // A lock directory of its own keeps this deliberate writer hold off the
2807        // volume lease the other tests in this process share.
2808        let pool = ConnectionPool::new(PoolConfig {
2809            path: Some(path.clone()),
2810            volume_lock_dir: Some(dir.path().join("volume-locks")),
2811            ..PoolConfig::for_test()
2812        })
2813        .expect("pool open");
2814        {
2815            let writer = pool.try_writer().unwrap();
2816            writer
2817                .conn()
2818                .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY)")
2819                .unwrap();
2820        }
2821        let handle = spawn(&pool, 8).unwrap();
2822
2823        // The delay must dwarf a COMMIT on a slow disk (an fsync of a few
2824        // hundred milliseconds has been observed on shared runners), or the
2825        // commit-stage comparison below reads the disk instead of the body.
2826        handle
2827            .send(|conn| {
2828                std::thread::sleep(Duration::from_millis(400));
2829                conn.execute("INSERT INTO t VALUES (1)", [])
2830                    .map_err(|error| StorageError::Pool {
2831                        operation: "writer_stage_sample".into(),
2832                        message: error.to_string(),
2833                    })
2834            })
2835            .await
2836            .unwrap();
2837
2838        let sample = last_writer_stage_observation(&pool).expect("writer stage sample");
2839        assert!(
2840            sample.body_micros >= 350_000,
2841            "the synthetic delay must land in the body stage: {sample:?}"
2842        );
2843        assert!(
2844            sample.body_micros > sample.queue_wait_micros,
2845            "fast queueing must not receive the body's delay: {sample:?}"
2846        );
2847        assert!(
2848            sample.body_micros > sample.transaction_acquire_micros,
2849            "an uncontended BEGIN must not receive the body's delay: {sample:?}"
2850        );
2851        assert!(
2852            sample.body_micros > sample.commit_micros,
2853            "a fast COMMIT must not receive the body's delay: {sample:?}"
2854        );
2855        assert!(sample.observed_at_unix_ms > 0);
2856    }
2857
2858    /// #1849: bounded-channel queue wait ends when the drain loop receives
2859    /// the request. Saturation in Tokio's blocking pool happens after that
2860    /// boundary and must remain in `total - named_stages`, never be relabeled
2861    /// as writer-queue contention.
2862    #[test]
2863    fn writer_queue_wait_excludes_blocking_pool_scheduling_delay() {
2864        let runtime = tokio::runtime::Builder::new_multi_thread()
2865            .worker_threads(1)
2866            .max_blocking_threads(1)
2867            .enable_all()
2868            .build()
2869            .expect("test runtime");
2870
2871        runtime.block_on(async {
2872            let dir = tempfile::tempdir().unwrap();
2873            let path = dir.path().join("writer_dequeue_boundary.db");
2874            let pool = file_pool(&path);
2875            let handle = spawn(&pool, 8).unwrap();
2876
2877            let (blocker_started_tx, blocker_started_rx) = std_mpsc::sync_channel(0);
2878            let (release_blocker_tx, release_blocker_rx) = std_mpsc::channel();
2879            let blocker = tokio::task::spawn_blocking(move || {
2880                blocker_started_tx.send(()).unwrap();
2881                release_blocker_rx.recv().unwrap();
2882            });
2883            blocker_started_rx
2884                .recv_timeout(Duration::from_secs(1))
2885                .expect("sole blocking worker must be occupied");
2886
2887            let reply = handle
2888                .enqueue(|_conn| Ok::<(), StorageError>(()))
2889                .await
2890                .expect("request must enter the bounded writer channel");
2891            let dequeue_deadline = Instant::now() + Duration::from_secs(1);
2892            while handle.queue_depth() != 0 {
2893                assert!(
2894                    Instant::now() < dequeue_deadline,
2895                    "writer drain never dequeued the accepted request"
2896                );
2897                tokio::task::yield_now().await;
2898            }
2899
2900            let scheduling_delay = Duration::from_millis(150);
2901            tokio::time::sleep(scheduling_delay).await;
2902            release_blocker_tx.send(()).unwrap();
2903            blocker.await.unwrap();
2904            reply.await.unwrap().unwrap();
2905
2906            let sample = last_writer_stage_observation(&pool).expect("writer stage sample");
2907            assert!(
2908                sample.total_micros.saturating_sub(sample.queue_wait_micros) >= 100_000,
2909                "the post-dequeue blocking-pool delay must not inflate queue_wait: {sample:?}"
2910            );
2911        });
2912    }
2913
2914    // `#[serial(tx_registry)]`: the active request holds a registered
2915    // `writer_task_tx` before its terminal failure, sharing the process-wide
2916    // registry key with the other writer-task transaction tests.
2917    #[tokio::test]
2918    #[serial(tx_registry)]
2919    async fn writer_task_failure_counters_are_acquisition_site_exact() {
2920        // Case 1: a request that reaches the seam, panics, and rolls back
2921        // cleanly must move `writer_task_request_failures` by exactly one —
2922        // never `writer_task_side_effects_unknown` — and the requests that
2923        // never dequeued into the seam (failed only because the queue closed
2924        // behind them) must move neither counter.
2925        {
2926            let dir = tempfile::tempdir().unwrap();
2927            let path = dir.path().join("writer_task_failure_counters_rollback.db");
2928            let pool = file_pool(&path);
2929            let handle = spawn(&pool, 8).expect("writer task should spawn");
2930
2931            let before = pool.writer_acquisition_snapshot();
2932            assert_eq!(before.writer_task_request_failures, 0);
2933            assert_eq!(before.writer_task_side_effects_unknown, 0);
2934
2935            let (started_tx, started_rx) = oneshot::channel::<()>();
2936            let (release_tx, release_rx) = std_mpsc::channel::<()>();
2937            let active = tokio::spawn({
2938                let handle = handle.clone();
2939                async move {
2940                    handle
2941                        .send(move |_conn| -> Result<(), StorageError> {
2942                            let _ = started_tx.send(());
2943                            release_rx.recv().expect("test must release active op");
2944                            panic!("intentional rollback-clean panic for counter test");
2945                        })
2946                        .await
2947                }
2948            });
2949            tokio::time::timeout(Duration::from_secs(5), started_rx)
2950                .await
2951                .expect("active request did not start")
2952                .expect("active request dropped its start signal");
2953
2954            let queued = handle
2955                .enqueue(|_conn| Ok::<(), StorageError>(()))
2956                .await
2957                .expect("second request must queue behind the active one");
2958
2959            release_tx.send(()).expect("release active op");
2960            let active_result = tokio::time::timeout(Duration::from_secs(5), active)
2961                .await
2962                .expect("active caller hung after panic")
2963                .expect("active caller task join");
2964            assert_writer_task_terminal_state(
2965                active_result,
2966                WriterTaskRequestState::TransactionRolledBack,
2967            );
2968
2969            let queued_result = queued.await.expect("terminal drain must reply");
2970            assert_writer_task_terminal_state(queued_result, WriterTaskRequestState::NotStarted);
2971
2972            let after = pool.writer_acquisition_snapshot();
2973            assert_eq!(
2974                after.writer_task_request_failures, 1,
2975                "only the request that actually reached the seam counts, not the ones \
2976                 failed by the queue-close drain"
2977            );
2978            assert_eq!(
2979                after.writer_task_side_effects_unknown, 0,
2980                "a clean rollback must not be counted as an unknown-side-effects outcome"
2981            );
2982        }
2983
2984        // Case 2: a request whose rollback is itself refused (denied by an
2985        // installed authorizer) terminates `SideEffectsUnknown` and must move
2986        // both counters by exactly one.
2987        {
2988            let dir = tempfile::tempdir().unwrap();
2989            let path = dir.path().join("writer_task_failure_counters_unknown.db");
2990            let pool = file_pool(&path);
2991            let handle = spawn(&pool, 8).expect("writer task should spawn");
2992
2993            let before = pool.writer_acquisition_snapshot();
2994
2995            let active_result = handle
2996                .send(|conn| -> Result<(), StorageError> {
2997                    conn.authorizer(Some(deny_rollback))
2998                        .map_err(|e| StorageError::Pool {
2999                            operation: "test_install_authorizer".into(),
3000                            message: e.to_string(),
3001                        })?;
3002                    Err(StorageError::Internal(
3003                        "intentional operation failure before denied rollback".into(),
3004                    ))
3005                })
3006                .await;
3007            assert_writer_task_terminal_state(
3008                active_result,
3009                WriterTaskRequestState::SideEffectsUnknown,
3010            );
3011
3012            // The reply above races the counter increment: `send` observes
3013            // its reply before `run_writer_task` finishes its own await and
3014            // updates the pool counters. A trailing request against the same
3015            // (now-terminal) handle can only resolve once that update and the
3016            // subsequent queue-close drain have both happened, so awaiting it
3017            // orders the snapshot below strictly after the increment.
3018            let sentinel_result = handle.send(|_conn| Ok::<(), StorageError>(())).await;
3019            assert_writer_task_terminal_state(sentinel_result, WriterTaskRequestState::NotStarted);
3020
3021            let after = pool.writer_acquisition_snapshot();
3022            assert_eq!(
3023                after.writer_task_request_failures - before.writer_task_request_failures,
3024                1
3025            );
3026            assert_eq!(
3027                after.writer_task_side_effects_unknown - before.writer_task_side_effects_unknown,
3028                1
3029            );
3030        }
3031    }
3032}
3033
3034#[cfg(test)]
3035#[path = "writer_task_lease_close_tests.rs"]
3036mod volume_lease_close_tests;
3037
3038#[cfg(all(test, any(unix, windows)))]
3039#[path = "writer_task_identity_tests.rs"]
3040mod identity_tests;