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