1use 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
63const WRITER_BEGIN_RETRY_DELAYS: [Duration; 2] =
70 [Duration::from_millis(5), Duration::from_millis(10)];
71
72type WriteOp<R> = Box<dyn FnOnce(&Connection) -> Result<R, StorageError> + Send>;
82
83#[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
125pub 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
194pub struct WriteRequest<R: Send + 'static> {
209 op: WriteOp<R>,
210 reply: oneshot::Sender<Result<R, StorageError>>,
211 top_level: bool,
212 checkpoint_bypass: bool,
214 vacuum_copy_headroom: bool,
215 telemetry: WriteTelemetry,
216}
217
218mod sealed {
219 pub trait Sealed {
224 fn execute_and_reply_reporting_terminal(
225 self: Box<Self>,
226 conn: &rusqlite::Connection,
227 tx_span: Option<khive_storage::tx_registry::TxHandle>,
228 queue_wait: std::time::Duration,
229 transaction_acquire: std::time::Duration,
230 ) -> Option<khive_storage::error::WriterTaskRequestState>;
231
232 fn execute_and_reply_top_level_reporting_terminal(
233 self: Box<Self>,
234 conn: &rusqlite::Connection,
235 queue_wait: std::time::Duration,
236 ) -> Option<khive_storage::error::WriterTaskRequestState>;
237
238 fn reply_error_after_begin(
239 self: Box<Self>,
240 err: khive_storage::error::StorageError,
241 queue_wait: std::time::Duration,
242 transaction_acquire: std::time::Duration,
243 );
244 }
245}
246
247pub trait AnyWriteRequest: sealed::Sealed + Send {
253 fn execute_and_reply(self: Box<Self>, conn: &Connection);
266
267 fn execute_and_reply_top_level(self: Box<Self>, conn: &Connection);
277
278 fn reply_error(self: Box<Self>, err: StorageError);
290
291 fn is_top_level(&self) -> bool;
295
296 fn is_checkpoint_bypass(&self) -> bool;
297
298 fn needs_vacuum_headroom(&self) -> bool;
299
300 fn queue_wait(&self) -> Duration;
303}
304
305#[derive(Debug, Clone, Copy, PartialEq, Eq)]
306enum RollbackDisposition {
307 RolledBack,
308 SideEffectsUnknown,
309}
310
311fn rollback_after_failure(conn: &Connection, failure_context: &'static str) -> RollbackDisposition {
314 match conn.execute_batch("ROLLBACK") {
315 Ok(()) if conn.is_autocommit() => RollbackDisposition::RolledBack,
316 Ok(()) => {
317 tracing::error!(
318 failure_context,
319 "writer transaction: ROLLBACK returned success but the connection is still in a \
320 transaction; request side effects are unknown"
321 );
322 RollbackDisposition::SideEffectsUnknown
323 }
324 Err(rollback_error) => {
325 tracing::error!(
326 error = %rollback_error,
327 failure_context,
328 "writer transaction: rollback after request failure failed; request side effects are \
329 unknown"
330 );
331 RollbackDisposition::SideEffectsUnknown
332 }
333 }
334}
335
336pub(crate) fn execute_wrapped_transaction<R, F>(
339 conn: &Connection,
340 commit_operation: &'static str,
341 operation: F,
342) -> (Result<R, StorageError>, Option<WriterTaskRequestState>)
343where
344 F: FnOnce(&Connection) -> Result<R, StorageError>,
345{
346 let profiled = execute_wrapped_transaction_profiled(conn, commit_operation, None, operation);
347 (profiled.result, profiled.terminal_state)
348}
349
350struct ProfiledWrappedTransaction<R> {
351 result: Result<R, StorageError>,
352 terminal_state: Option<WriterTaskRequestState>,
353 body: Duration,
354 commit: Duration,
355}
356
357fn execute_wrapped_transaction_profiled<R, F>(
358 conn: &Connection,
359 commit_operation: &'static str,
360 db: Option<&str>,
361 operation: F,
362) -> ProfiledWrappedTransaction<R>
363where
364 F: FnOnce(&Connection) -> Result<R, StorageError>,
365{
366 let body_started = Instant::now();
367 let operation_outcome = catch_unwind(AssertUnwindSafe(|| operation(conn)));
368 let body = body_started.elapsed();
369
370 match operation_outcome {
371 Ok(Ok(value)) => {
372 let commit_started = Instant::now();
373 let commit_outcome = conn.execute_batch("COMMIT");
374 let commit = commit_started.elapsed();
375 match commit_outcome {
376 Ok(()) if conn.is_autocommit() => ProfiledWrappedTransaction {
377 result: Ok(value),
378 terminal_state: None,
379 body,
380 commit,
381 },
382 Ok(()) => {
383 tracing::error!(
384 "writer transaction: COMMIT returned success but the connection is still in \
385 a transaction; request side effects are unknown"
386 );
387 let request_state = WriterTaskRequestState::SideEffectsUnknown;
388 ProfiledWrappedTransaction {
389 result: Err(writer_task_terminated(request_state)),
390 terminal_state: Some(request_state),
391 body,
392 commit,
393 }
394 }
395 Err(commit_error) => {
396 if let Some(db) = db {
397 crate::timeout_sink::maybe_emit_sqlite_full(db, &commit_error);
398 }
399 match rollback_after_failure(conn, "commit failure") {
400 RollbackDisposition::RolledBack => ProfiledWrappedTransaction {
401 result: Err(StorageError::WriterTaskRequestFailed {
402 request_state: WriterTaskRequestState::TransactionRolledBack,
403 source: Box::new(StorageError::Pool {
404 operation: commit_operation.into(),
405 message: commit_error.to_string(),
406 }),
407 }),
408 terminal_state: None,
409 body,
410 commit,
411 },
412 RollbackDisposition::SideEffectsUnknown => {
413 let request_state = WriterTaskRequestState::SideEffectsUnknown;
414 ProfiledWrappedTransaction {
415 result: Err(writer_task_terminated_with_cause(
416 request_state,
417 &commit_error,
418 )),
419 terminal_state: Some(request_state),
420 body,
421 commit,
422 }
423 }
424 }
425 }
426 }
427 }
428 Ok(Err(operation_error)) => {
429 match rollback_after_failure(conn, "request operation failure") {
430 RollbackDisposition::RolledBack => ProfiledWrappedTransaction {
431 result: Err(StorageError::WriterTaskRequestFailed {
432 request_state: WriterTaskRequestState::TransactionRolledBack,
433 source: Box::new(operation_error),
434 }),
435 terminal_state: None,
436 body,
437 commit: Duration::ZERO,
438 },
439 RollbackDisposition::SideEffectsUnknown => {
440 if let Some(db) = db {
444 crate::timeout_sink::maybe_emit_sqlite_full(db, &operation_error);
445 }
446 let request_state = WriterTaskRequestState::SideEffectsUnknown;
447 ProfiledWrappedTransaction {
448 result: Err(writer_task_terminated_with_cause(
449 request_state,
450 &operation_error,
451 )),
452 terminal_state: Some(request_state),
453 body,
454 commit: Duration::ZERO,
455 }
456 }
457 }
458 }
459 Err(_panic_payload) => {
460 let request_state = match rollback_after_failure(conn, "request panic") {
461 RollbackDisposition::RolledBack => WriterTaskRequestState::TransactionRolledBack,
462 RollbackDisposition::SideEffectsUnknown => {
463 WriterTaskRequestState::SideEffectsUnknown
464 }
465 };
466 ProfiledWrappedTransaction {
467 result: Err(writer_task_terminated(request_state)),
468 terminal_state: Some(request_state),
469 body,
470 commit: Duration::ZERO,
471 }
472 }
473 }
474}
475
476impl<R: Send + 'static> sealed::Sealed for WriteRequest<R> {
477 fn execute_and_reply_reporting_terminal(
478 self: Box<Self>,
479 conn: &Connection,
480 tx_span: Option<khive_storage::tx_registry::TxHandle>,
481 queue_wait: Duration,
482 transaction_acquire: Duration,
483 ) -> Option<WriterTaskRequestState> {
484 let WriteRequest {
488 op,
489 reply,
490 telemetry,
491 ..
492 } = *self;
493 let profiled = execute_wrapped_transaction_profiled(
494 conn,
495 "writer_task_commit",
496 Some(&telemetry.db),
497 op,
498 );
499 if let Err(error) = &profiled.result {
500 crate::timeout_sink::maybe_emit_sqlite_full(&telemetry.db, error);
501 }
502 drop(tx_span);
507 telemetry.finish(
508 queue_wait,
509 transaction_acquire,
510 profiled.body,
511 profiled.commit,
512 );
513 let _ = reply.send(profiled.result);
516 profiled.terminal_state
517 }
518
519 fn execute_and_reply_top_level_reporting_terminal(
520 self: Box<Self>,
521 conn: &Connection,
522 queue_wait: Duration,
523 ) -> Option<WriterTaskRequestState> {
524 let WriteRequest {
525 op,
526 reply,
527 telemetry,
528 ..
529 } = *self;
530 let body_started = Instant::now();
531 let outcome = catch_unwind(AssertUnwindSafe(|| op(conn)));
532 let body = body_started.elapsed();
533 let telemetry_db = telemetry.db.clone();
534 telemetry.finish(queue_wait, Duration::ZERO, body, Duration::ZERO);
535 match outcome {
536 Ok(outcome) if conn.is_autocommit() => {
537 if let Err(error) = &outcome {
540 crate::timeout_sink::maybe_emit_sqlite_full(&telemetry_db, error);
541 }
542 let _ = reply.send(outcome);
543 None
544 }
545 Ok(_outcome) => {
546 tracing::error!(
547 "writer task: top-level request returned with an open transaction; request \
548 side effects are unknown"
549 );
550 let request_state = WriterTaskRequestState::SideEffectsUnknown;
551 let _ = reply.send(Err(writer_task_terminated(request_state)));
552 Some(request_state)
553 }
554 Err(_panic_payload) => {
555 let request_state = WriterTaskRequestState::SideEffectsUnknown;
559 let _ = reply.send(Err(writer_task_terminated(request_state)));
560 Some(request_state)
561 }
562 }
563 }
564
565 fn reply_error_after_begin(
566 self: Box<Self>,
567 err: StorageError,
568 queue_wait: Duration,
569 transaction_acquire: Duration,
570 ) {
571 let WriteRequest {
572 reply, telemetry, ..
573 } = *self;
574 telemetry.finish(
575 queue_wait,
576 transaction_acquire,
577 Duration::ZERO,
578 Duration::ZERO,
579 );
580 let _ = reply.send(Err(err));
581 }
582}
583
584impl<R: Send + 'static> AnyWriteRequest for WriteRequest<R> {
585 fn execute_and_reply(self: Box<Self>, conn: &Connection) {
586 let queue_wait = self.queue_wait();
587 let _ = sealed::Sealed::execute_and_reply_reporting_terminal(
588 self,
589 conn,
590 None,
591 queue_wait,
592 Duration::ZERO,
593 );
594 }
595
596 fn execute_and_reply_top_level(self: Box<Self>, conn: &Connection) {
597 let queue_wait = self.queue_wait();
598 let _ =
599 sealed::Sealed::execute_and_reply_top_level_reporting_terminal(self, conn, queue_wait);
600 }
601
602 fn reply_error(self: Box<Self>, err: StorageError) {
603 let queue_wait = self.queue_wait();
604 sealed::Sealed::reply_error_after_begin(self, err, queue_wait, Duration::ZERO);
605 }
606
607 fn is_top_level(&self) -> bool {
608 self.top_level
609 }
610
611 fn is_checkpoint_bypass(&self) -> bool {
612 self.checkpoint_bypass
613 }
614
615 fn needs_vacuum_headroom(&self) -> bool {
616 self.vacuum_copy_headroom
617 }
618
619 fn queue_wait(&self) -> Duration {
620 self.telemetry.queue_wait()
621 }
622}
623
624fn writer_task_terminated(request_state: WriterTaskRequestState) -> StorageError {
625 StorageError::writer_task_terminated(request_state)
626}
627
628fn writer_task_terminated_with_cause(
629 request_state: WriterTaskRequestState,
630 error: &(dyn std::error::Error + 'static),
631) -> StorageError {
632 let mut source = Some(error);
633 let mut sqlite_full_codes = None;
634 while let Some(current) = source {
635 if let Some(rusqlite::Error::SqliteFailure(code, _)) =
636 current.downcast_ref::<rusqlite::Error>()
637 {
638 if code.code == rusqlite::ErrorCode::DiskFull {
639 sqlite_full_codes = Some((code.extended_code & 0xff, code.extended_code));
640 break;
641 }
642 }
643 source = current.source();
644 }
645 StorageError::WriterTaskTerminated {
648 request_state,
649 sqlite_full_codes,
650 }
651}
652
653fn writer_task_begin_error(error: rusqlite::Error, busy_timeout: Duration) -> StorageError {
654 if crate::timeout_sink::is_busy_or_locked(&error) {
655 StorageError::WriterTaskBusy {
656 timeout_ms: u64::try_from(busy_timeout.as_millis()).unwrap_or(u64::MAX),
657 }
658 } else {
659 StorageError::Pool {
660 operation: "writer_task_begin".into(),
661 message: error.to_string(),
662 }
663 }
664}
665
666#[derive(Clone, Debug)]
670pub struct WriterTaskHandle {
671 tx: mpsc::Sender<Box<dyn AnyWriteRequest + Send>>,
672 backend_key: Option<PathBuf>,
675 db: String,
680 slow_write_threshold: Option<std::time::Duration>,
686 enqueue_timeout: std::time::Duration,
707}
708
709impl WriterTaskHandle {
710 async fn enqueue<R, F>(
724 &self,
725 op: F,
726 ) -> Result<oneshot::Receiver<Result<R, StorageError>>, StorageError>
727 where
728 R: Send + 'static,
729 F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
730 {
731 self.enqueue_inner(op, false, false, false).await
732 }
733
734 async fn enqueue_inner<R, F>(
738 &self,
739 op: F,
740 top_level: bool,
741 checkpoint_bypass: bool,
742 vacuum_copy_headroom: bool,
743 ) -> Result<oneshot::Receiver<Result<R, StorageError>>, StorageError>
744 where
745 R: Send + 'static,
746 F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
747 {
748 let (reply_tx, reply_rx) = oneshot::channel();
749 let telemetry = WriteTelemetry::new(
750 self.backend_key.clone(),
751 self.db.clone(),
752 self.queue_depth(),
753 self.slow_write_threshold,
754 );
755 let request = WriteRequest {
756 op: Box::new(op),
757 reply: reply_tx,
758 top_level,
759 checkpoint_bypass,
760 vacuum_copy_headroom,
761 telemetry,
762 };
763
764 self.tx
765 .send(Box::new(request))
766 .await
767 .map_err(|_| writer_task_terminated(WriterTaskRequestState::NotStarted))?;
768
769 Ok(reply_rx)
770 }
771
772 pub async fn send<R, F>(&self, op: F) -> Result<R, StorageError>
779 where
780 R: Send + 'static,
781 F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
782 {
783 let reply_rx = self.enqueue(op).await?;
784 reply_rx
785 .await
786 .map_err(|_| writer_task_terminated(WriterTaskRequestState::SideEffectsUnknown))?
787 }
788
789 pub async fn send_with_timeout<R, F>(
804 &self,
805 op: F,
806 timeout: std::time::Duration,
807 ) -> Result<R, StorageError>
808 where
809 R: Send + 'static,
810 F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
811 {
812 let reply_rx = match tokio::time::timeout(timeout, self.enqueue(op)).await {
813 Ok(Ok(reply_rx)) => reply_rx,
814 Ok(Err(e)) => return Err(e),
815 Err(_elapsed) => {
816 let timeout_ms = timeout.as_millis() as u64;
817 crate::timeout_sink::emit_queue_saturation(&self.db, timeout_ms);
818 return Err(StorageError::WriteQueueFull { timeout_ms });
819 }
820 };
821
822 reply_rx
823 .await
824 .map_err(|_| writer_task_terminated(WriterTaskRequestState::SideEffectsUnknown))?
825 }
826
827 pub async fn send_bounded<R, F>(&self, op: F) -> Result<R, StorageError>
838 where
839 R: Send + 'static,
840 F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
841 {
842 self.send_with_timeout(op, self.enqueue_timeout).await
843 }
844
845 pub async fn send_top_level<R, F>(&self, op: F) -> Result<R, StorageError>
857 where
858 R: Send + 'static,
859 F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
860 {
861 let reply_rx = self.enqueue_inner(op, true, false, false).await?;
862 reply_rx
863 .await
864 .map_err(|_| writer_task_terminated(WriterTaskRequestState::SideEffectsUnknown))?
865 }
866
867 pub async fn send_top_level_bounded<R, F>(&self, op: F) -> Result<R, StorageError>
872 where
873 R: Send + 'static,
874 F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
875 {
876 self.send_top_level_bounded_inner(op, false, false).await
877 }
878
879 pub(crate) async fn send_checkpoint_bounded<R, F>(&self, op: F) -> Result<R, StorageError>
882 where
883 R: Send + 'static,
884 F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
885 {
886 self.send_top_level_bounded_inner(op, true, false).await
887 }
888
889 pub(crate) async fn send_vacuum_bounded<R, F>(&self, op: F) -> Result<R, StorageError>
891 where
892 R: Send + 'static,
893 F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
894 {
895 self.send_top_level_bounded_inner(op, false, true).await
896 }
897
898 async fn send_top_level_bounded_inner<R, F>(
899 &self,
900 op: F,
901 checkpoint_bypass: bool,
902 vacuum_copy_headroom: bool,
903 ) -> Result<R, StorageError>
904 where
905 R: Send + 'static,
906 F: FnOnce(&Connection) -> Result<R, StorageError> + Send + 'static,
907 {
908 let reply_rx = match tokio::time::timeout(
909 self.enqueue_timeout,
910 self.enqueue_inner(op, true, checkpoint_bypass, vacuum_copy_headroom),
911 )
912 .await
913 {
914 Ok(Ok(reply_rx)) => reply_rx,
915 Ok(Err(e)) => return Err(e),
916 Err(_elapsed) => {
917 let timeout_ms = self.enqueue_timeout.as_millis() as u64;
918 crate::timeout_sink::emit_queue_saturation(&self.db, timeout_ms);
919 return Err(StorageError::WriteQueueFull { timeout_ms });
920 }
921 };
922
923 reply_rx
924 .await
925 .map_err(|_| writer_task_terminated(WriterTaskRequestState::SideEffectsUnknown))?
926 }
927
928 pub fn queue_depth(&self) -> usize {
937 self.tx.max_capacity() - self.tx.capacity()
938 }
939
940 pub fn capacity(&self) -> usize {
943 self.tx.max_capacity()
944 }
945}
946
947pub fn spawn(pool: &ConnectionPool, capacity: usize) -> Result<WriterTaskHandle, SqliteError> {
977 let conn = pool.open_standalone_writer_untracked()?;
981 let acquisition_counters = pool.writer_acquisition_counters();
982 let write_admission = pool.write_admission();
983 let busy_timeout = pool.config().busy_timeout;
984 let origin = pool.origin();
985 let backend_key = writer_db_key(pool);
986 let db = crate::timeout_sink::db_label(pool);
987 let (tx, rx) = mpsc::channel(capacity.max(1));
988 let join = tokio::spawn(run_writer_task(
989 conn,
990 rx,
991 origin,
992 db.clone(),
993 acquisition_counters,
994 write_admission,
995 busy_timeout,
996 ));
997 pool.set_writer_task_join(join);
1001 Ok(WriterTaskHandle {
1002 tx,
1003 backend_key,
1004 db,
1005 slow_write_threshold: crate::timeout_sink::slow_write_threshold(),
1006 enqueue_timeout: std::time::Duration::from_millis(
1007 pool.config().write_admission_deadline_ms,
1008 ),
1009 })
1010}
1011
1012async fn close_and_fail_queued_requests(rx: &mut mpsc::Receiver<Box<dyn AnyWriteRequest + Send>>) {
1020 rx.close();
1021 while let Some(request) = rx.recv().await {
1022 request.reply_error(writer_task_terminated(WriterTaskRequestState::NotStarted));
1023 }
1024}
1025
1026fn begin_immediate_with_retry(
1030 conn: &Connection,
1031 acquisition_counters: &WriterAcquisitionCounters,
1032 busy_timeout: Duration,
1033 mut set_busy_timeout: impl FnMut(&Connection, Duration) -> rusqlite::Result<()>,
1034) -> (rusqlite::Result<()>, Duration, u32) {
1035 let transaction_acquire_started = Instant::now();
1036 let mut begin_attempt = 1_u32;
1037 let mut retry_delays = WRITER_BEGIN_RETRY_DELAYS.into_iter();
1038 let mut busy_timeout_lowered = false;
1039 let begin_outcome = loop {
1040 match conn.execute_batch("BEGIN IMMEDIATE") {
1041 Ok(()) => break Ok(()),
1042 Err(error) if crate::timeout_sink::is_busy_or_locked(&error) => {
1043 acquisition_counters.record_writer_task_begin_busy();
1049 let Some(delay) = retry_delays.next() else {
1050 break Err(error);
1051 };
1052 let remaining_budget =
1060 busy_timeout.saturating_sub(transaction_acquire_started.elapsed());
1061 if remaining_budget.is_zero() {
1062 break Err(error);
1063 }
1064 if let Err(set_err) = set_busy_timeout(conn, remaining_budget) {
1065 tracing::warn!(
1066 error = %set_err,
1067 "writer task: failed to lower busy_timeout for BEGIN \
1068 retry; surfacing the original busy refusal"
1069 );
1070 break Err(error);
1072 }
1073 busy_timeout_lowered = true;
1074 acquisition_counters.record_writer_task_begin_busy_absorbed();
1077 tracing::debug!(
1078 attempt = begin_attempt,
1079 backoff_ms = delay.as_millis() as u64,
1080 budget_remaining_ms = remaining_budget.as_millis() as u64,
1081 "writer task: BEGIN IMMEDIATE refused busy; retrying before \
1082 request execution"
1083 );
1084 std::thread::sleep(delay);
1085 begin_attempt = begin_attempt.saturating_add(1);
1086 }
1087 Err(error) => break Err(error),
1088 }
1089 };
1090 let transaction_acquire = transaction_acquire_started.elapsed();
1091 if busy_timeout_lowered {
1097 if let Err(restore_err) = set_busy_timeout(conn, busy_timeout) {
1098 tracing::warn!(
1099 error = %restore_err,
1100 "writer task: failed to restore busy_timeout after a BEGIN retry \
1101 sequence"
1102 );
1103 }
1104 }
1105 (begin_outcome, transaction_acquire, begin_attempt)
1106}
1107
1108struct BlockingWriterConnection {
1109 conn: Option<Connection>,
1110 admission: Arc<crate::pool::WriteAdmission>,
1111 volume_lease: Option<crate::disk_guard::VolumeLease>,
1112}
1113
1114impl BlockingWriterConnection {
1115 fn new(conn: Connection, admission: Arc<crate::pool::WriteAdmission>) -> Self {
1116 Self {
1117 conn: Some(conn),
1118 admission,
1119 volume_lease: None,
1120 }
1121 }
1122
1123 fn finish(
1124 mut self,
1125 state: Option<WriterTaskRequestState>,
1126 ) -> (Option<Connection>, Option<WriterTaskRequestState>) {
1127 if state.is_some() {
1128 self.retire();
1129 (None, state)
1130 } else {
1131 (self.conn.take(), None)
1132 }
1133 }
1134
1135 fn retire(&mut self) {
1136 if let Some(conn) = self.conn.take() {
1137 let _ = self.admission.close_retired_connection(conn);
1140 }
1141 }
1142}
1143
1144impl Drop for BlockingWriterConnection {
1145 fn drop(&mut self) {
1146 self.retire();
1147 }
1148}
1149
1150impl std::ops::Deref for BlockingWriterConnection {
1151 type Target = Connection;
1152
1153 fn deref(&self) -> &Connection {
1154 self.conn
1155 .as_ref()
1156 .expect("blocking writer owns its connection until settlement")
1157 }
1158}
1159
1160async fn run_writer_task(
1174 mut conn: Connection,
1175 mut rx: mpsc::Receiver<Box<dyn AnyWriteRequest + Send>>,
1176 origin: khive_storage::tx_registry::TxOrigin,
1177 db: String,
1178 acquisition_counters: Arc<WriterAcquisitionCounters>,
1179 write_admission: Arc<crate::pool::WriteAdmission>,
1180 busy_timeout: Duration,
1181) {
1182 while let Some(request) = rx.recv().await {
1183 let queue_wait = request.queue_wait();
1187 let origin = origin.clone();
1188 let blocking_counters = Arc::clone(&acquisition_counters);
1189 let blocking_admission = Arc::clone(&write_admission);
1190 let blocking_db = db.clone();
1191 let outcome = tokio::task::spawn_blocking(move || {
1192 let mut conn = BlockingWriterConnection::new(conn, Arc::clone(&blocking_admission));
1193 let acquisition_counters = blocking_counters;
1194 if !conn.is_autocommit() {
1199 tracing::error!(
1200 "writer task: connection is not in autocommit mode before request dispatch; \
1201 retiring the poisoned writer without running the request"
1202 );
1203 let request_state = WriterTaskRequestState::NotStarted;
1204 request.reply_error(writer_task_terminated(request_state));
1205 return conn.finish(Some(request_state));
1206 }
1207
1208 conn.volume_lease = if request.is_checkpoint_bypass() {
1212 None
1213 } else {
1214 match blocking_admission.acquire() {
1215 Ok(lease) => lease,
1216 Err(error) => {
1217 request.reply_error(error.into_storage_error(
1218 khive_storage::StorageCapability::Sql,
1219 "writer_task_admission",
1220 ));
1221 return conn.finish(None);
1222 }
1223 }
1224 };
1225
1226 let terminal_state = if request.is_top_level() {
1227 if !request.is_checkpoint_bypass() {
1228 let admission = if request.needs_vacuum_headroom() {
1229 blocking_admission.check_for_vacuum()
1230 } else {
1231 blocking_admission.check()
1232 };
1233 if let Err(error) = admission {
1234 request.reply_error(error.into_storage_error(
1235 khive_storage::StorageCapability::Sql,
1236 "writer_task_admission",
1237 ));
1238 return conn.finish(None);
1239 }
1240 }
1241 acquisition_counters.record_writer_task_acquisition();
1249 sealed::Sealed::execute_and_reply_top_level_reporting_terminal(
1250 request, &conn, queue_wait,
1251 )
1252 } else {
1253 let tx_span = khive_storage::tx_registry::register_scoped(
1254 Some("writer_task_tx".to_string()),
1255 origin,
1256 );
1257 let (begin_outcome, transaction_acquire, begin_attempt) =
1258 begin_immediate_with_retry(
1259 &conn,
1260 &acquisition_counters,
1261 busy_timeout,
1262 Connection::busy_timeout,
1263 );
1264 match begin_outcome {
1265 Ok(()) => {
1266 if let Err(error) = blocking_admission.check() {
1267 let request_state =
1268 match rollback_after_failure(&conn, "capacity admission") {
1269 RollbackDisposition::RolledBack => None,
1270 RollbackDisposition::SideEffectsUnknown => {
1271 Some(WriterTaskRequestState::SideEffectsUnknown)
1272 }
1273 };
1274 drop(tx_span);
1275 let error = if let Some(state) = request_state {
1276 writer_task_terminated(state)
1277 } else {
1278 error.into_storage_error(
1279 khive_storage::StorageCapability::Sql,
1280 "writer_task_admission",
1281 )
1282 };
1283 sealed::Sealed::reply_error_after_begin(
1284 request,
1285 error,
1286 queue_wait,
1287 transaction_acquire,
1288 );
1289 return conn.finish(request_state);
1290 }
1291 acquisition_counters.record_writer_task_acquisition();
1292 sealed::Sealed::execute_and_reply_reporting_terminal(
1293 request,
1294 &conn,
1295 Some(tx_span),
1296 queue_wait,
1297 transaction_acquire,
1298 )
1299 }
1300 Err(e) => {
1301 crate::timeout_sink::maybe_emit_sqlite_full(&blocking_db, &e);
1307 tracing::warn!(
1308 error = %e,
1309 attempts = begin_attempt,
1310 "writer task: BEGIN IMMEDIATE failed; replying an \
1311 error without running the request's operation"
1312 );
1313 drop(tx_span);
1318 let begin_error = writer_task_begin_error(e, busy_timeout);
1325 if !matches!(&begin_error, StorageError::WriterTaskBusy { .. }) {
1326 acquisition_counters.record_writer_task_begin_error();
1327 }
1328 sealed::Sealed::reply_error_after_begin(
1329 request,
1330 begin_error,
1331 queue_wait,
1332 transaction_acquire,
1333 );
1334 None
1335 }
1336 }
1337 };
1338 conn.finish(terminal_state)
1339 })
1340 .await;
1341
1342 match outcome {
1343 Ok((Some(returned_conn), None)) => conn = returned_conn,
1344 Ok((_returned_conn, Some(request_state))) => {
1345 acquisition_counters.record_writer_task_request_failure();
1346 if request_state == WriterTaskRequestState::SideEffectsUnknown {
1347 acquisition_counters.record_writer_task_side_effects_unknown();
1348 }
1349 tracing::error!(
1350 request_state = %request_state,
1351 "writer task reached a terminal request or connection state; closing and \
1352 failing the queue without restarting"
1353 );
1354 crate::timeout_sink::emit_writer_task_retirement(
1355 &db,
1356 &format!("terminal request state: {request_state}"),
1357 );
1358 close_and_fail_queued_requests(&mut rx).await;
1359 return;
1360 }
1361 Ok((None, None)) => unreachable!("a reusable writer returns its connection"),
1362 Err(join_err) => {
1363 acquisition_counters.record_writer_task_request_failure();
1364 tracing::error!(
1365 error = %join_err,
1366 "writer task blocking closure failed outside the request \
1367 panic boundary; closing and failing the queue without restarting"
1368 );
1369 crate::timeout_sink::emit_writer_task_retirement(
1370 &db,
1371 &format!("blocking closure join failure: {join_err}"),
1372 );
1373 close_and_fail_queued_requests(&mut rx).await;
1374 return;
1375 }
1376 }
1377 }
1378}
1379
1380#[cfg(test)]
1381mod tests {
1382 use super::*;
1383 use crate::pool::PoolConfig;
1384 use rusqlite::hooks::{AuthAction, AuthContext, Authorization, TransactionOperation};
1385 use serial_test::serial;
1386
1387 #[test]
1388 fn begin_error_classification_is_code_based_and_narrow() {
1389 for code in [rusqlite::ffi::SQLITE_BUSY, rusqlite::ffi::SQLITE_LOCKED] {
1390 let error = rusqlite::Error::SqliteFailure(
1391 rusqlite::ffi::Error::new(code),
1392 Some("rendered text is irrelevant".to_string()),
1393 );
1394 assert!(matches!(
1395 writer_task_begin_error(error, Duration::from_millis(175)),
1396 StorageError::WriterTaskBusy { timeout_ms: 175 }
1397 ));
1398 }
1399
1400 let structural = rusqlite::Error::SqliteFailure(
1401 rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_CONSTRAINT),
1402 Some("database is locked".to_string()),
1403 );
1404 assert!(matches!(
1405 writer_task_begin_error(structural, Duration::from_millis(175)),
1406 StorageError::Pool { ref operation, .. } if operation == "writer_task_begin"
1407 ));
1408 }
1409 use std::future::Future;
1410 use std::pin::Pin;
1411 use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
1412 use std::sync::mpsc as std_mpsc;
1413 use std::sync::{Arc, Mutex};
1414 use std::task::{Context, Poll, Wake, Waker};
1415 use std::time::Duration;
1416
1417 fn file_pool(path: &std::path::Path) -> ConnectionPool {
1418 let cfg = PoolConfig {
1419 path: Some(path.to_path_buf()),
1420 ..PoolConfig::for_test()
1421 };
1422 ConnectionPool::new(cfg).expect("pool open")
1423 }
1424
1425 fn deny_commit_and_rollback(ctx: AuthContext<'_>) -> Authorization {
1426 match ctx.action {
1427 AuthAction::Transaction {
1430 operation: TransactionOperation::Unknown | TransactionOperation::Rollback,
1431 } => Authorization::Deny,
1432 _ => Authorization::Allow,
1433 }
1434 }
1435
1436 fn deny_commit(ctx: AuthContext<'_>) -> Authorization {
1437 match ctx.action {
1438 AuthAction::Transaction {
1439 operation: TransactionOperation::Unknown,
1440 } => Authorization::Deny,
1441 _ => Authorization::Allow,
1442 }
1443 }
1444
1445 fn deny_rollback(ctx: AuthContext<'_>) -> Authorization {
1446 match ctx.action {
1447 AuthAction::Transaction {
1448 operation: TransactionOperation::Rollback,
1449 } => Authorization::Deny,
1450 _ => Authorization::Allow,
1451 }
1452 }
1453
1454 fn assert_writer_task_terminal_state<T: std::fmt::Debug>(
1455 result: Result<T, StorageError>,
1456 expected: WriterTaskRequestState,
1457 ) {
1458 match result {
1459 Err(StorageError::WriterTaskTerminated { request_state, .. }) => {
1460 assert_eq!(request_state, expected)
1461 }
1462 other => panic!("expected WriterTaskTerminated({expected:?}), got {other:?}"),
1463 }
1464 }
1465
1466 #[test]
1467 fn terminal_native_full_evidence_uses_the_original_error_chain() {
1468 let extended = rusqlite::ffi::SQLITE_FULL | (3 << 8);
1469 let native = rusqlite::Error::SqliteFailure(
1470 rusqlite::ffi::Error::new(extended),
1471 Some("synthetic classifier control".into()),
1472 );
1473 let wrapped = StorageError::driver(khive_storage::StorageCapability::Sql, "nested", native);
1474 let error =
1475 writer_task_terminated_with_cause(WriterTaskRequestState::SideEffectsUnknown, &wrapped);
1476 assert!(matches!(error, StorageError::WriterTaskTerminated {
1477 request_state: WriterTaskRequestState::SideEffectsUnknown,
1478 sqlite_full_codes: Some((rusqlite::ffi::SQLITE_FULL, code)),
1479 } if code == extended));
1480 assert!(
1481 std::error::Error::source(&error).is_none(),
1482 "evidence must not replay the original sink cause"
1483 );
1484 assert!(!error.is_retryable());
1485 assert_eq!(error.capability(), None);
1486
1487 for code in [rusqlite::ffi::SQLITE_BUSY, rusqlite::ffi::SQLITE_IOERR] {
1488 let native = rusqlite::Error::SqliteFailure(
1489 rusqlite::ffi::Error::new(code),
1490 Some("database or disk is full".into()),
1491 );
1492 assert!(matches!(
1493 writer_task_terminated_with_cause(
1494 WriterTaskRequestState::SideEffectsUnknown,
1495 &native
1496 ),
1497 StorageError::WriterTaskTerminated {
1498 sqlite_full_codes: None,
1499 ..
1500 }
1501 ));
1502 }
1503 assert!(matches!(
1504 writer_task_terminated(WriterTaskRequestState::NotStarted),
1505 StorageError::WriterTaskTerminated {
1506 sqlite_full_codes: None,
1507 ..
1508 }
1509 ));
1510 }
1511
1512 struct ParkedWake {
1513 entered: std_mpsc::SyncSender<()>,
1514 release: Mutex<std_mpsc::Receiver<()>>,
1515 }
1516
1517 impl Wake for ParkedWake {
1518 fn wake(self: Arc<Self>) {
1519 self.entered
1520 .send(())
1521 .expect("reply sender must rendezvous with the test");
1522 self.release
1523 .lock()
1524 .unwrap_or_else(|poisoned| poisoned.into_inner())
1525 .recv()
1526 .expect("test must release the parked reply sender");
1527 }
1528 }
1529
1530 fn arm_parked_wake<F: Future>(
1531 mut future: Pin<&mut F>,
1532 ) -> (std_mpsc::Receiver<()>, std_mpsc::Sender<()>) {
1533 let (entered_tx, entered_rx) = std_mpsc::sync_channel(0);
1534 let (release_tx, release_rx) = std_mpsc::channel();
1535 let waker = Waker::from(Arc::new(ParkedWake {
1536 entered: entered_tx,
1537 release: Mutex::new(release_rx),
1538 }));
1539 let mut context = Context::from_waker(&waker);
1540 assert!(
1541 matches!(future.as_mut().poll(&mut context), Poll::Pending),
1542 "writer send must remain pending until its operation replies"
1543 );
1544 (entered_rx, release_tx)
1545 }
1546
1547 fn poll_ready<F: Future>(mut future: Pin<&mut F>) -> F::Output {
1548 let mut context = Context::from_waker(Waker::noop());
1549 match future.as_mut().poll(&mut context) {
1550 Poll::Ready(output) => output,
1551 Poll::Pending => panic!("reply wake must make the writer send ready"),
1552 }
1553 }
1554
1555 fn database_tx_view(pool: &ConnectionPool) -> khive_storage::tx_registry::TxOriginFilter {
1556 match pool.origin() {
1557 khive_storage::tx_registry::TxOrigin::Database(identity) => {
1558 khive_storage::tx_registry::TxOriginFilter::Secondary(identity)
1559 }
1560 other => panic!("expected a file-backed database origin, got {other:?}"),
1561 }
1562 }
1563
1564 async fn wait_for_writer_span_to_close(view: &khive_storage::tx_registry::TxOriginFilter) {
1565 tokio::time::timeout(Duration::from_secs(5), async {
1566 while khive_storage::tx_registry::any_open_labeled(view, "writer_task_tx") {
1567 tokio::task::yield_now().await;
1568 }
1569 })
1570 .await
1571 .expect("writer task transaction span must eventually close");
1572 }
1573
1574 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1575 #[serial(tx_registry)]
1576 async fn writer_task_connection_maintains_rfc3339_expression_indexes() {
1577 let dir = tempfile::tempdir().unwrap();
1578 let path = dir.path().join("writer_task_rfc3339_expression_index.db");
1579 let pool = file_pool(&path);
1580 {
1581 let writer = pool.writer().expect("pooled writer");
1582 writer
1583 .conn()
1584 .execute_batch(
1585 "CREATE TABLE deadlines(id INTEGER PRIMARY KEY, due TEXT);
1586 CREATE INDEX idx_deadlines_strict \
1587 ON deadlines(ifnull(khive_rfc3339_strict_key(due), x''));",
1588 )
1589 .expect("pooled writer registers the key function");
1590 }
1591 let handle = spawn(&pool, 8).expect("writer task spawn");
1592
1593 let inserted = handle
1594 .send(|conn| {
1595 conn.execute(
1596 "INSERT INTO deadlines(id, due) VALUES (1, '2026-01-01T00:00:00Z')",
1597 [],
1598 )
1599 .map_err(|error| StorageError::Pool {
1600 operation: "test_insert".into(),
1601 message: error.to_string(),
1602 })
1603 })
1604 .await
1605 .expect("the writer task's connection maintains the expression index");
1606 assert_eq!(inserted, 1);
1607 }
1608
1609 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1610 #[serial(tx_registry)]
1611 async fn successful_send_reply_waits_for_writer_tx_deregistration() {
1612 let dir = tempfile::tempdir().unwrap();
1613 let path = dir.path().join("writer_task_success_reply_lifecycle.db");
1614 let pool = file_pool(&path);
1615 let view = database_tx_view(&pool);
1616 let handle = spawn(&pool, 8).expect("writer task spawn");
1617 let (op_started_tx, op_started_rx) = std_mpsc::sync_channel(0);
1618 let (op_release_tx, op_release_rx) = std_mpsc::channel();
1619
1620 let send = handle.send(move |_conn| {
1621 op_started_tx
1622 .send(())
1623 .expect("operation must rendezvous with the test");
1624 op_release_rx
1625 .recv()
1626 .expect("test must release the operation");
1627 Ok::<_, StorageError>(())
1628 });
1629 tokio::pin!(send);
1630 let (reply_entered_rx, reply_release_tx) = arm_parked_wake(send.as_mut());
1631
1632 op_started_rx
1633 .recv_timeout(Duration::from_secs(5))
1634 .expect("writer operation must start");
1635 op_release_tx.send(()).expect("release writer operation");
1636 reply_entered_rx
1637 .recv_timeout(Duration::from_secs(5))
1638 .expect("reply sender must wake the waiting caller");
1639
1640 let reply = poll_ready(send.as_mut());
1641 let span_was_open_at_reply =
1642 khive_storage::tx_registry::any_open_labeled(&view, "writer_task_tx");
1643
1644 reply_release_tx
1645 .send(())
1646 .expect("release parked reply sender");
1647 wait_for_writer_span_to_close(&view).await;
1648
1649 reply.expect("committed operation reply");
1650 assert!(
1651 !span_was_open_at_reply,
1652 "a successful caller reply must not become observable while its committed writer_task_tx span remains registered"
1653 );
1654 }
1655
1656 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1657 #[serial(tx_registry)]
1658 async fn begin_failure_reply_waits_for_writer_tx_deregistration() {
1659 let dir = tempfile::tempdir().unwrap();
1660 let path = dir.path().join("writer_task_begin_reply_lifecycle.db");
1661 let cfg = PoolConfig {
1664 path: Some(path.clone()),
1665 volume_lock_dir: Some(dir.path().join("volume-locks")),
1666 busy_timeout: Duration::from_millis(150),
1667 ..PoolConfig::for_test()
1668 };
1669 let pool = ConnectionPool::new(cfg).unwrap();
1670 let view = database_tx_view(&pool);
1671 let handle = spawn(&pool, 8).expect("writer task spawn");
1672 let lock_holder = rusqlite::Connection::open(&path).expect("external writer");
1675 lock_holder
1676 .execute_batch("BEGIN IMMEDIATE")
1677 .expect("hold database write lock");
1678 let op_ran = Arc::new(AtomicBool::new(false));
1679 let op_ran_in_request = Arc::clone(&op_ran);
1680
1681 let send = handle.send(move |_conn| {
1682 op_ran_in_request.store(true, Ordering::SeqCst);
1683 Ok::<_, StorageError>(())
1684 });
1685 tokio::pin!(send);
1686 let (reply_entered_rx, reply_release_tx) = arm_parked_wake(send.as_mut());
1687 reply_entered_rx
1688 .recv_timeout(Duration::from_secs(5))
1689 .expect("BEGIN failure must wake the waiting caller");
1690
1691 let reply = poll_ready(send.as_mut());
1692 let span_was_open_at_reply =
1693 khive_storage::tx_registry::any_open_labeled(&view, "writer_task_tx");
1694
1695 reply_release_tx
1696 .send(())
1697 .expect("release parked reply sender");
1698 wait_for_writer_span_to_close(&view).await;
1699 lock_holder
1700 .execute_batch("ROLLBACK")
1701 .expect("release database write lock");
1702
1703 assert!(
1704 matches!(
1705 &reply,
1706 Err(StorageError::WriterTaskBusy { timeout_ms }) if *timeout_ms == 150
1707 ),
1708 "expected typed retryable writer-task contention, got {reply:?}"
1709 );
1710 assert!(!op_ran.load(Ordering::SeqCst));
1711 assert!(
1712 !span_was_open_at_reply,
1713 "a BEGIN-failure caller reply must not become observable while its writer_task_tx span remains registered"
1714 );
1715 }
1716
1717 #[tokio::test]
1724 #[serial(tx_registry)]
1725 async fn begin_immediate_failure_replies_error_without_running_op() {
1726 let dir = tempfile::tempdir().unwrap();
1731 let path = dir.path().join("writer_task_begin_failure.db");
1732 let cfg = PoolConfig {
1735 path: Some(path.clone()),
1736 volume_lock_dir: Some(dir.path().join("volume-locks")),
1737 busy_timeout: Duration::from_millis(150),
1738 ..PoolConfig::for_test()
1739 };
1740 let pool = ConnectionPool::new(cfg).unwrap();
1741 {
1742 let writer = pool.try_writer().unwrap();
1743 writer
1744 .conn()
1745 .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY, v TEXT)")
1746 .unwrap();
1747 }
1748
1749 let handle = spawn(&pool, 8).expect("writer task should spawn on a file-backed pool");
1750
1751 let lock_holder = rusqlite::Connection::open(&path).unwrap();
1752 lock_holder.execute_batch("BEGIN IMMEDIATE").unwrap();
1753
1754 let op_ran = Arc::new(AtomicBool::new(false));
1755 let op_ran_clone = Arc::clone(&op_ran);
1756 let result = handle
1757 .send(move |conn| {
1758 op_ran_clone.store(true, Ordering::SeqCst);
1759 conn.execute("INSERT INTO t (id, v) VALUES (99, 'should-not-land')", [])
1760 .map_err(|e| StorageError::Pool {
1761 operation: "test_insert".into(),
1762 message: e.to_string(),
1763 })
1764 })
1765 .await;
1766
1767 assert!(
1768 matches!(
1769 &result,
1770 Err(StorageError::WriterTaskBusy { timeout_ms }) if *timeout_ms == 150
1771 ),
1772 "expected a typed retryable error on contended BEGIN IMMEDIATE, got {result:?}"
1773 );
1774 assert!(
1775 !op_ran.load(Ordering::SeqCst),
1776 "the request's operation closure must never run when BEGIN \
1777 IMMEDIATE fails — running it would land a partial write in \
1778 autocommit mode for a request the caller is told failed"
1779 );
1780
1781 lock_holder.execute_batch("ROLLBACK").unwrap();
1784 drop(lock_holder);
1785
1786 handle
1787 .send(|conn| {
1788 conn.execute("INSERT INTO t (id, v) VALUES (100, 'next-request')", [])
1789 .map_err(|e| StorageError::Pool {
1790 operation: "test_insert_after_busy".into(),
1791 message: e.to_string(),
1792 })
1793 })
1794 .await
1795 .expect("transient contention must not retire the writer task");
1796
1797 let reader = pool.reader().expect("reader");
1798 let count: i64 = reader
1799 .conn()
1800 .query_row("SELECT COUNT(*) FROM t WHERE id IN (99, 100)", [], |row| {
1801 row.get(0)
1802 })
1803 .unwrap();
1804 assert_eq!(
1805 count, 1,
1806 "the failed request must not land, while the next request commits on the same task"
1807 );
1808 }
1809
1810 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1813 #[serial(tx_registry)]
1814 async fn transient_begin_contention_clears_within_budget_and_op_runs_once() {
1815 let dir = tempfile::tempdir().unwrap();
1824 let path = dir.path().join("writer_task_begin_transient_contention.db");
1825 let pool = ConnectionPool::new(PoolConfig {
1828 path: Some(path.clone()),
1829 volume_lock_dir: Some(dir.path().join("volume-locks")),
1830 busy_timeout: Duration::from_millis(500),
1831 ..PoolConfig::for_test()
1832 })
1833 .unwrap();
1834 {
1835 let writer = pool.try_writer().unwrap();
1836 writer
1837 .conn()
1838 .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY)")
1839 .unwrap();
1840 }
1841 let handle = spawn(&pool, 8).expect("writer task spawn");
1842 let lock_holder = Connection::open(&path).unwrap();
1846 lock_holder.execute_batch("BEGIN IMMEDIATE").unwrap();
1847
1848 let op_runs = Arc::new(AtomicUsize::new(0));
1849 let op_runs_in_request = Arc::clone(&op_runs);
1850 let send_future = handle.send(move |conn| {
1851 op_runs_in_request.fetch_add(1, Ordering::SeqCst);
1852 conn.execute("INSERT INTO t (id) VALUES (1)", [])
1853 .map_err(|error| StorageError::Pool {
1854 operation: "test_insert_after_transient_contention".into(),
1855 message: error.to_string(),
1856 })
1857 });
1858 let release_future = async {
1859 tokio::time::sleep(Duration::from_millis(50)).await;
1860 lock_holder.execute_batch("ROLLBACK").unwrap();
1861 };
1862 let (result, ()) = tokio::join!(send_future, release_future);
1863
1864 assert_eq!(
1865 result.expect("BEGIN IMMEDIATE succeeds once the transient lock clears"),
1866 1
1867 );
1868 assert_eq!(
1869 op_runs.load(Ordering::SeqCst),
1870 1,
1871 "the FnOnce request closure must execute exactly once"
1872 );
1873
1874 let settled = pool.writer_acquisition_snapshot();
1875 assert_eq!(
1876 settled.writer_task_begin_busy, 0,
1877 "contention absorbed inside SQLite's own busy_timeout wait must never \
1878 surface as a Rust-level refusal"
1879 );
1880 assert_eq!(settled.writer_task_begin_busy_absorbed, 0);
1881 let reader = pool.reader().unwrap();
1882 let rows: i64 = reader
1883 .conn()
1884 .query_row("SELECT COUNT(*) FROM t", [], |row| row.get(0))
1885 .unwrap();
1886 assert_eq!(rows, 1, "exactly one closure execution commits one row");
1887 }
1888
1889 include!("writer_task_admission_retry_tests.rs");
1890
1891 #[tokio::test]
1895 #[serial(tx_registry)]
1896 async fn writer_task_executes_op_and_commits() {
1897 let dir = tempfile::tempdir().unwrap();
1898 let path = dir.path().join("writer_task_commit.db");
1899 let pool = file_pool(&path);
1900 {
1901 let writer = pool.try_writer().unwrap();
1902 writer
1903 .conn()
1904 .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY, v TEXT)")
1905 .unwrap();
1906 }
1907
1908 let handle = spawn(&pool, 8).expect("writer task should spawn on a file-backed pool");
1909
1910 let affected = handle
1911 .send(|conn| {
1912 conn.execute("INSERT INTO t (id, v) VALUES (1, 'hello')", [])
1913 .map_err(|e| StorageError::Pool {
1914 operation: "test_insert".into(),
1915 message: e.to_string(),
1916 })
1917 })
1918 .await
1919 .expect("op should succeed");
1920 assert_eq!(affected, 1);
1921
1922 let reader = pool.reader().expect("reader");
1926 let v: String = reader
1927 .conn()
1928 .query_row("SELECT v FROM t WHERE id = 1", [], |row| row.get(0))
1929 .expect("row must be committed and visible to a reader");
1930 assert_eq!(v, "hello");
1931
1932 let counters = pool.writer_acquisition_snapshot();
1933 assert_eq!(counters.acquisitions, 2);
1934 assert_eq!(counters.pooled_acquisitions, 1);
1935 assert_eq!(counters.standalone_acquisitions, 0);
1936 assert_eq!(counters.writer_task_acquisitions, 1);
1937 assert_eq!(counters.timeouts, 0);
1938 }
1939
1940 #[tokio::test]
1941 async fn writer_task_connection_follows_checkpoint_ownership_claim() {
1942 let dir = tempfile::tempdir().unwrap();
1943 let path = dir.path().join("writer_task_autocheckpoint.db");
1944 let pool = file_pool(&path);
1945 let handle = pool
1949 .writer_task_handle()
1950 .expect("writer task should spawn")
1951 .expect("file-backed pool resolves the write queue on");
1952
1953 let read_pages = |handle: &WriterTaskHandle| {
1954 let handle = handle.clone();
1955 async move {
1956 handle
1957 .send_top_level(|conn| {
1958 conn.pragma_query_value(None, "wal_autocheckpoint", |row| {
1959 row.get::<_, u32>(0)
1960 })
1961 .map_err(|e| StorageError::Pool {
1962 operation: "test_wal_autocheckpoint".into(),
1963 message: e.to_string(),
1964 })
1965 })
1966 .await
1967 .expect("query writer-task connection pragma")
1968 }
1969 };
1970
1971 assert_eq!(
1974 read_pages(&handle).await,
1975 crate::pool::FALLBACK_WAL_AUTOCHECKPOINT_PAGES
1976 );
1977
1978 pool.claim_checkpoint_ownership().expect("claim ownership");
1981 pool.propagate_checkpoint_claim_to_writer_task()
1982 .await
1983 .expect("propagate claim to the running writer task");
1984 assert_eq!(read_pages(&handle).await, 0);
1985 }
1986
1987 #[test]
1988 fn spawn_fails_on_in_memory_pool() {
1989 let cfg = PoolConfig {
1995 path: None,
1996 ..PoolConfig::default()
1997 };
1998 let pool = ConnectionPool::new(cfg).unwrap();
1999 let result = spawn(&pool, 8);
2000 assert!(
2001 result.is_err(),
2002 "in-memory pools must reject spawn, not panic"
2003 );
2004 }
2005
2006 include!("writer_task_queue_capacity_tests.rs");
2007
2008 #[tokio::test]
2009 #[serial(tx_registry)]
2010 async fn operation_failure_with_successful_rollback_reports_finality_once_and_continues() {
2011 let dir = tempfile::tempdir().unwrap();
2012 let path = dir.path().join("writer_task_operation_rollback.db");
2013 let pool = file_pool(&path);
2014 {
2015 let writer = pool.try_writer().unwrap();
2016 writer
2017 .conn()
2018 .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY, v TEXT)")
2019 .unwrap();
2020 }
2021 let handle = spawn(&pool, 8).expect("writer task spawn");
2022
2023 let executions = Arc::new(AtomicUsize::new(0));
2024 let executions_in_op = Arc::clone(&executions);
2025 let original_error = handle
2026 .send(move |conn| -> Result<(), StorageError> {
2027 executions_in_op.fetch_add(1, Ordering::SeqCst);
2028 conn.execute("INSERT INTO t (id, v) VALUES (1, 'rolled-back')", [])
2029 .map_err(|e| StorageError::Pool {
2030 operation: "test_operation_error_insert".into(),
2031 message: e.to_string(),
2032 })?;
2033 Err(StorageError::Internal(
2034 "intentional operation failure".into(),
2035 ))
2036 })
2037 .await;
2038 match &original_error {
2039 Err(StorageError::WriterTaskRequestFailed {
2040 request_state: WriterTaskRequestState::TransactionRolledBack,
2041 source,
2042 }) => assert!(
2043 matches!(source.as_ref(), StorageError::Internal(message)
2044 if message == "intentional operation failure"),
2045 "the proven-rollback wrapper must retain the typed operation error: {source:?}"
2046 ),
2047 other => panic!(
2048 "a confirmed rollback must carry TransactionRolledBack and preserve the operation error, got {other:?}"
2049 ),
2050 }
2051 assert_eq!(
2052 executions.load(Ordering::SeqCst),
2053 1,
2054 "finality propagation must not replay the request closure"
2055 );
2056
2057 let affected = handle
2058 .send(|conn| {
2059 conn.execute("INSERT INTO t (id, v) VALUES (2, 'committed')", [])
2060 .map_err(|e| StorageError::Pool {
2061 operation: "test_operation_error_followup_insert".into(),
2062 message: e.to_string(),
2063 })
2064 })
2065 .await
2066 .expect("the writer must continue after a confirmed rollback");
2067 assert_eq!(affected, 1);
2068 assert_eq!(
2069 executions.load(Ordering::SeqCst),
2070 1,
2071 "serving a follow-up request must not replay the rolled-back closure"
2072 );
2073
2074 let reader = pool.reader().expect("reader");
2075 let rolled_back: i64 = reader
2076 .conn()
2077 .query_row("SELECT COUNT(*) FROM t WHERE id = 1", [], |row| row.get(0))
2078 .unwrap();
2079 let committed: i64 = reader
2080 .conn()
2081 .query_row("SELECT COUNT(*) FROM t WHERE id = 2", [], |row| row.get(0))
2082 .unwrap();
2083 assert_eq!(rolled_back, 0);
2084 assert_eq!(committed, 1);
2085 }
2086
2087 #[tokio::test]
2088 #[serial(tx_registry)]
2089 async fn commit_failure_with_successful_rollback_reports_finality_once_and_continues() {
2090 let dir = tempfile::tempdir().unwrap();
2091 let path = dir.path().join("writer_task_commit_rollback.db");
2092 let pool = file_pool(&path);
2093 {
2094 let writer = pool.try_writer().unwrap();
2095 writer
2096 .conn()
2097 .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY, v TEXT)")
2098 .unwrap();
2099 }
2100 let handle = spawn(&pool, 8).expect("writer task spawn");
2101
2102 let executions = Arc::new(AtomicUsize::new(0));
2103 let executions_in_op = Arc::clone(&executions);
2104 let commit_error = handle
2105 .send(move |conn| -> Result<usize, StorageError> {
2106 executions_in_op.fetch_add(1, Ordering::SeqCst);
2107 let affected = conn
2108 .execute("INSERT INTO t (id, v) VALUES (1, 'rolled-back')", [])
2109 .map_err(|e| StorageError::Pool {
2110 operation: "test_commit_error_insert".into(),
2111 message: e.to_string(),
2112 })?;
2113 conn.authorizer(Some(deny_commit))
2114 .map_err(|e| StorageError::Pool {
2115 operation: "test_install_authorizer".into(),
2116 message: e.to_string(),
2117 })?;
2118 Ok(affected)
2119 })
2120 .await;
2121 match &commit_error {
2122 Err(StorageError::WriterTaskRequestFailed {
2123 request_state: WriterTaskRequestState::TransactionRolledBack,
2124 source,
2125 }) => assert!(
2126 matches!(source.as_ref(), StorageError::Pool { operation, .. }
2127 if operation == "writer_task_commit"),
2128 "the proven-rollback wrapper must retain the typed COMMIT error: {source:?}"
2129 ),
2130 other => panic!(
2131 "a confirmed rollback must carry TransactionRolledBack and preserve the COMMIT error, got {other:?}"
2132 ),
2133 }
2134 assert!(
2135 commit_error
2136 .as_ref()
2137 .expect_err("COMMIT must be denied")
2138 .is_retryable(),
2139 "the existing retryable commit-error contract must remain unchanged after a \
2140 confirmed rollback"
2141 );
2142 assert_eq!(
2143 executions.load(Ordering::SeqCst),
2144 1,
2145 "finality propagation must not replay the request closure"
2146 );
2147
2148 let affected = handle
2149 .send(|conn| {
2150 conn.authorizer(None::<fn(AuthContext<'_>) -> Authorization>)
2151 .map_err(|e| StorageError::Pool {
2152 operation: "test_remove_authorizer".into(),
2153 message: e.to_string(),
2154 })?;
2155 conn.execute("INSERT INTO t (id, v) VALUES (2, 'committed')", [])
2156 .map_err(|e| StorageError::Pool {
2157 operation: "test_commit_error_followup_insert".into(),
2158 message: e.to_string(),
2159 })
2160 })
2161 .await
2162 .expect("the writer must continue after the failed COMMIT is rolled back");
2163 assert_eq!(affected, 1);
2164 assert_eq!(
2165 executions.load(Ordering::SeqCst),
2166 1,
2167 "serving a follow-up request must not replay the rolled-back closure"
2168 );
2169
2170 let reader = pool.reader().expect("reader");
2171 let rolled_back: i64 = reader
2172 .conn()
2173 .query_row("SELECT COUNT(*) FROM t WHERE id = 1", [], |row| row.get(0))
2174 .unwrap();
2175 let committed: i64 = reader
2176 .conn()
2177 .query_row("SELECT COUNT(*) FROM t WHERE id = 2", [], |row| row.get(0))
2178 .unwrap();
2179 assert_eq!(rolled_back, 0);
2180 assert_eq!(committed, 1);
2181 }
2182
2183 #[test]
2184 fn top_level_request_returning_with_open_transaction_reports_side_effects_unknown() {
2185 let conn = Connection::open_in_memory().expect("in-memory connection");
2186 conn.execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY)")
2187 .unwrap();
2188 let (reply_tx, mut reply_rx) = oneshot::channel();
2189 let request = WriteRequest {
2190 op: Box::new(|conn| -> Result<usize, StorageError> {
2191 conn.execute_batch("BEGIN IMMEDIATE")
2192 .map_err(|e| StorageError::Pool {
2193 operation: "test_top_level_begin".into(),
2194 message: e.to_string(),
2195 })?;
2196 conn.execute("INSERT INTO t (id) VALUES (1)", [])
2197 .map_err(|e| StorageError::Pool {
2198 operation: "test_top_level_insert".into(),
2199 message: e.to_string(),
2200 })
2201 }),
2202 reply: reply_tx,
2203 top_level: true,
2204 checkpoint_bypass: false,
2205 vacuum_copy_headroom: false,
2206 telemetry: WriteTelemetry::new(None, "test".to_string(), 0, None),
2207 };
2208
2209 let terminal_state = sealed::Sealed::execute_and_reply_top_level_reporting_terminal(
2210 Box::new(request),
2211 &conn,
2212 Duration::ZERO,
2213 );
2214 assert_eq!(
2215 terminal_state,
2216 Some(WriterTaskRequestState::SideEffectsUnknown)
2217 );
2218 let reply = reply_rx
2219 .try_recv()
2220 .expect("active request must receive a typed terminal reply");
2221 assert_writer_task_terminal_state(reply, WriterTaskRequestState::SideEffectsUnknown);
2222 assert!(
2223 !conn.is_autocommit(),
2224 "the fixture must prove the post-request autocommit check observed an open transaction"
2225 );
2226 }
2227
2228 #[test]
2229 fn commit_failure_with_failed_rollback_reports_side_effects_unknown() {
2230 let conn = Connection::open_in_memory().expect("in-memory connection");
2231 conn.execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY); BEGIN IMMEDIATE")
2232 .unwrap();
2233 let executions = Arc::new(AtomicUsize::new(0));
2234 let executions_in_op = Arc::clone(&executions);
2235 let (reply_tx, mut reply_rx) = oneshot::channel();
2236 let request = WriteRequest {
2237 op: Box::new(move |conn| -> Result<usize, StorageError> {
2238 executions_in_op.fetch_add(1, Ordering::SeqCst);
2239 let affected = conn
2240 .execute("INSERT INTO t (id) VALUES (1)", [])
2241 .map_err(|e| StorageError::Pool {
2242 operation: "test_insert_before_commit_failure".into(),
2243 message: e.to_string(),
2244 })?;
2245 conn.authorizer(Some(deny_commit_and_rollback))
2246 .map_err(|e| StorageError::Pool {
2247 operation: "test_install_authorizer".into(),
2248 message: e.to_string(),
2249 })?;
2250 Ok(affected)
2251 }),
2252 reply: reply_tx,
2253 top_level: false,
2254 checkpoint_bypass: false,
2255 vacuum_copy_headroom: false,
2256 telemetry: WriteTelemetry::new(None, "test".to_string(), 0, None),
2257 };
2258
2259 let terminal_state = sealed::Sealed::execute_and_reply_reporting_terminal(
2260 Box::new(request),
2261 &conn,
2262 None,
2263 Duration::ZERO,
2264 Duration::ZERO,
2265 );
2266 assert_eq!(
2267 terminal_state,
2268 Some(WriterTaskRequestState::SideEffectsUnknown)
2269 );
2270 let reply = reply_rx
2271 .try_recv()
2272 .expect("active request must receive a typed terminal reply");
2273 assert_writer_task_terminal_state(reply, WriterTaskRequestState::SideEffectsUnknown);
2274 assert_eq!(executions.load(Ordering::SeqCst), 1);
2275 assert!(
2276 !conn.is_autocommit(),
2277 "the denied COMMIT and ROLLBACK must leave the test connection poisoned"
2278 );
2279 }
2280
2281 #[tokio::test]
2282 #[serial(tx_registry)]
2283 async fn poisoned_connection_retires_before_queued_top_level_request() {
2284 let dir = tempfile::tempdir().unwrap();
2285 let path = dir.path().join("writer_task_rollback_poison.db");
2286 let pool = file_pool(&path);
2287 {
2288 let writer = pool.try_writer().unwrap();
2289 writer
2290 .conn()
2291 .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY, v TEXT)")
2292 .unwrap();
2293 }
2294 let handle = spawn(&pool, 8).expect("writer task spawn");
2295 let (started_tx, started_rx) = oneshot::channel::<()>();
2296 let (release_tx, release_rx) = std_mpsc::channel::<()>();
2297
2298 let active = tokio::spawn({
2299 let handle = handle.clone();
2300 async move {
2301 handle
2302 .send(move |conn| -> Result<usize, StorageError> {
2303 let affected = conn
2304 .execute("INSERT INTO t (id, v) VALUES (1, 'active')", [])
2305 .map_err(|e| StorageError::Pool {
2306 operation: "test_active_insert".into(),
2307 message: e.to_string(),
2308 })?;
2309 conn.authorizer(Some(deny_commit_and_rollback))
2310 .map_err(|e| StorageError::Pool {
2311 operation: "test_install_authorizer".into(),
2312 message: e.to_string(),
2313 })?;
2314 let _ = started_tx.send(());
2315 release_rx.recv().expect("test must release active op");
2316 Ok(affected)
2317 })
2318 .await
2319 }
2320 });
2321
2322 tokio::time::timeout(Duration::from_secs(5), started_rx)
2323 .await
2324 .expect("active request did not start")
2325 .expect("active request dropped its start signal");
2326
2327 let queued_ran = Arc::new(AtomicBool::new(false));
2328 let queued_ran_in_op = Arc::clone(&queued_ran);
2329 let queued_top_level = handle
2330 .enqueue_inner(
2331 move |conn| {
2332 queued_ran_in_op.store(true, Ordering::SeqCst);
2333 conn.execute("INSERT INTO t (id, v) VALUES (2, 'queued')", [])
2334 .map_err(|e| StorageError::Pool {
2335 operation: "test_queued_top_level_insert".into(),
2336 message: e.to_string(),
2337 })
2338 },
2339 true,
2340 false,
2341 false,
2342 )
2343 .await
2344 .expect("top-level request must queue behind active request");
2345 release_tx.send(()).expect("release active op");
2346
2347 let active_result = tokio::time::timeout(Duration::from_secs(5), active)
2348 .await
2349 .expect("active caller hung after rollback failure")
2350 .expect("active caller task join");
2351 assert_writer_task_terminal_state(
2352 active_result,
2353 WriterTaskRequestState::SideEffectsUnknown,
2354 );
2355
2356 let queued_result = tokio::time::timeout(Duration::from_secs(5), queued_top_level)
2357 .await
2358 .expect("queued top-level caller hung after terminal failure")
2359 .expect("terminal drain must preserve queued typed reply");
2360 assert_writer_task_terminal_state(queued_result, WriterTaskRequestState::NotStarted);
2361 assert!(
2362 !queued_ran.load(Ordering::SeqCst),
2363 "a top-level request must never run on the poisoned connection"
2364 );
2365
2366 let future_ran = Arc::new(AtomicBool::new(false));
2367 let future_ran_in_op = Arc::clone(&future_ran);
2368 let future_result = handle
2369 .send_top_level(move |_conn| {
2370 future_ran_in_op.store(true, Ordering::SeqCst);
2371 Ok::<(), StorageError>(())
2372 })
2373 .await;
2374 assert_writer_task_terminal_state(future_result, WriterTaskRequestState::NotStarted);
2375 assert!(!future_ran.load(Ordering::SeqCst));
2376 }
2377
2378 #[test]
2379 fn operation_failure_with_failed_rollback_reports_side_effects_unknown() {
2380 let conn = Connection::open_in_memory().expect("in-memory connection");
2381 conn.execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY); BEGIN IMMEDIATE")
2382 .unwrap();
2383 let executions = Arc::new(AtomicUsize::new(0));
2384 let executions_in_op = Arc::clone(&executions);
2385 let (reply_tx, mut reply_rx) = oneshot::channel();
2386 let request = WriteRequest {
2387 op: Box::new(move |conn| -> Result<(), StorageError> {
2388 executions_in_op.fetch_add(1, Ordering::SeqCst);
2389 conn.authorizer(Some(deny_rollback))
2390 .map_err(|e| StorageError::Pool {
2391 operation: "test_install_authorizer".into(),
2392 message: e.to_string(),
2393 })?;
2394 Err(StorageError::Internal(
2395 "intentional operation failure before denied rollback".into(),
2396 ))
2397 }),
2398 reply: reply_tx,
2399 top_level: false,
2400 checkpoint_bypass: false,
2401 vacuum_copy_headroom: false,
2402 telemetry: WriteTelemetry::new(None, "test".to_string(), 0, None),
2403 };
2404
2405 let terminal_state = sealed::Sealed::execute_and_reply_reporting_terminal(
2406 Box::new(request),
2407 &conn,
2408 None,
2409 Duration::ZERO,
2410 Duration::ZERO,
2411 );
2412 assert_eq!(
2413 terminal_state,
2414 Some(WriterTaskRequestState::SideEffectsUnknown)
2415 );
2416 let reply = reply_rx
2417 .try_recv()
2418 .expect("active request must receive a typed terminal reply");
2419 assert_writer_task_terminal_state(reply, WriterTaskRequestState::SideEffectsUnknown);
2420 assert_eq!(executions.load(Ordering::SeqCst), 1);
2421 assert!(
2422 !conn.is_autocommit(),
2423 "the denied ROLLBACK must leave the test connection poisoned"
2424 );
2425 }
2426
2427 #[test]
2428 fn wrapped_panic_with_failed_rollback_reports_side_effects_unknown() {
2429 let conn = Connection::open_in_memory().expect("in-memory connection");
2430 conn.execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY); BEGIN IMMEDIATE")
2431 .unwrap();
2432 let (reply_tx, mut reply_rx) = oneshot::channel();
2433 let request = WriteRequest {
2434 op: Box::new(|conn| -> Result<(), StorageError> {
2435 conn.execute_batch("INSERT INTO t (id) VALUES (1); COMMIT")
2440 .map_err(|e| StorageError::Pool {
2441 operation: "test_force_rollback_failure".into(),
2442 message: e.to_string(),
2443 })?;
2444 panic!("intentional panic after illicit commit");
2445 }),
2446 reply: reply_tx,
2447 top_level: false,
2448 checkpoint_bypass: false,
2449 vacuum_copy_headroom: false,
2450 telemetry: WriteTelemetry::new(None, "test".to_string(), 0, None),
2451 };
2452
2453 let terminal_state = sealed::Sealed::execute_and_reply_reporting_terminal(
2454 Box::new(request),
2455 &conn,
2456 None,
2457 Duration::ZERO,
2458 Duration::ZERO,
2459 );
2460 assert_eq!(
2461 terminal_state,
2462 Some(WriterTaskRequestState::SideEffectsUnknown)
2463 );
2464 let reply = reply_rx
2465 .try_recv()
2466 .expect("active request must receive a typed terminal reply");
2467 assert_writer_task_terminal_state(reply, WriterTaskRequestState::SideEffectsUnknown);
2468
2469 let count: i64 = conn
2470 .query_row("SELECT COUNT(*) FROM t", [], |row| row.get(0))
2471 .unwrap();
2472 assert_eq!(
2473 count, 1,
2474 "the fixture's committed side effect proves why the state must be unknown"
2475 );
2476 }
2477
2478 #[tokio::test]
2482 #[serial(tx_registry)]
2483 async fn wrapped_panic_rolls_back_and_terminally_fails_queue() {
2484 let dir = tempfile::tempdir().unwrap();
2485 let path = dir.path().join("writer_task_wrapped_panic.db");
2486 let cfg = PoolConfig {
2487 path: Some(path),
2488 write_queue_enabled: Some(true),
2489 write_queue_capacity: 8,
2490 ..PoolConfig::for_test()
2491 };
2492 let pool = ConnectionPool::new(cfg).unwrap();
2493 {
2494 let writer = pool.try_writer().unwrap();
2495 writer
2496 .conn()
2497 .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY, v TEXT)")
2498 .unwrap();
2499 }
2500
2501 let handle = pool
2504 .writer_task_handle()
2505 .expect("writer task lookup")
2506 .expect("file-backed queued pool must spawn its writer task");
2507 assert_eq!(pool.writer_task_spawn_count(), 1);
2508
2509 let (started_tx, started_rx) = oneshot::channel::<()>();
2510 let (release_tx, release_rx) = std_mpsc::channel::<()>();
2511 let active = tokio::spawn({
2512 let handle = handle.clone();
2513 async move {
2514 handle
2515 .send(move |conn| -> Result<usize, StorageError> {
2516 conn.execute("INSERT INTO t (id, v) VALUES (1, 'active')", [])
2517 .map_err(|e| StorageError::Pool {
2518 operation: "test_active_insert".into(),
2519 message: e.to_string(),
2520 })?;
2521 let _ = started_tx.send(());
2522 release_rx.recv().expect("test must release active op");
2523 panic!("intentional wrapped writer request panic");
2524 })
2525 .await
2526 }
2527 });
2528
2529 tokio::time::timeout(Duration::from_secs(5), started_rx)
2530 .await
2531 .expect("active request did not start")
2532 .expect("active request dropped its start signal");
2533
2534 let queued_one_ran = Arc::new(AtomicBool::new(false));
2535 let queued_one_ran_in_op = Arc::clone(&queued_one_ran);
2536 let queued_one = handle
2537 .enqueue(move |conn| {
2538 queued_one_ran_in_op.store(true, Ordering::SeqCst);
2539 conn.execute("INSERT INTO t (id, v) VALUES (2, 'queued-one')", [])
2540 .map_err(|e| StorageError::Pool {
2541 operation: "test_queued_one_insert".into(),
2542 message: e.to_string(),
2543 })
2544 })
2545 .await
2546 .expect("first queued request must be accepted");
2547
2548 let queued_two_ran = Arc::new(AtomicBool::new(false));
2549 let queued_two_ran_in_op = Arc::clone(&queued_two_ran);
2550 let queued_two = handle
2551 .enqueue(move |_conn| {
2552 queued_two_ran_in_op.store(true, Ordering::SeqCst);
2553 Ok::<String, StorageError>("queued-two-ran".to_string())
2554 })
2555 .await
2556 .expect("second queued request must be accepted");
2557
2558 assert_eq!(
2559 handle.queue_depth(),
2560 2,
2561 "both heterogeneous requests must be buffered behind the active op"
2562 );
2563 release_tx.send(()).expect("release active op");
2564
2565 let active_result = tokio::time::timeout(Duration::from_secs(5), active)
2566 .await
2567 .expect("active caller hung after panic")
2568 .expect("active caller task join");
2569 assert_writer_task_terminal_state(
2570 active_result,
2571 WriterTaskRequestState::TransactionRolledBack,
2572 );
2573
2574 let queued_one_result = tokio::time::timeout(Duration::from_secs(5), queued_one)
2575 .await
2576 .expect("first queued caller hung after terminal failure")
2577 .expect("terminal drain must preserve first typed reply");
2578 assert_writer_task_terminal_state(queued_one_result, WriterTaskRequestState::NotStarted);
2579
2580 let queued_two_result = tokio::time::timeout(Duration::from_secs(5), queued_two)
2581 .await
2582 .expect("second queued caller hung after terminal failure")
2583 .expect("terminal drain must preserve second typed reply");
2584 assert_writer_task_terminal_state(queued_two_result, WriterTaskRequestState::NotStarted);
2585 assert!(!queued_one_ran.load(Ordering::SeqCst));
2586 assert!(!queued_two_ran.load(Ordering::SeqCst));
2587
2588 let future_ran = Arc::new(AtomicBool::new(false));
2589 let future_ran_in_op = Arc::clone(&future_ran);
2590 let future_result = handle
2591 .send(move |_conn| {
2592 future_ran_in_op.store(true, Ordering::SeqCst);
2593 Ok::<(), StorageError>(())
2594 })
2595 .await;
2596 assert_writer_task_terminal_state(future_result, WriterTaskRequestState::NotStarted);
2597 assert!(!future_ran.load(Ordering::SeqCst));
2598
2599 let cached_after_failure = pool
2600 .writer_task_handle()
2601 .expect("cached writer task lookup")
2602 .expect("pool retains its terminal handle");
2603 assert_eq!(
2604 pool.writer_task_spawn_count(),
2605 1,
2606 "a terminal writer task must not be restarted behind callers' backs"
2607 );
2608 let cached_result = cached_after_failure
2609 .send(|_conn| Ok::<(), StorageError>(()))
2610 .await;
2611 assert_writer_task_terminal_state(cached_result, WriterTaskRequestState::NotStarted);
2612
2613 let reader = pool.reader().expect("reader");
2614 let count: i64 = reader
2615 .conn()
2616 .query_row("SELECT COUNT(*) FROM t", [], |row| row.get(0))
2617 .unwrap();
2618 assert_eq!(
2619 count, 0,
2620 "the active transaction must be rolled back and queued ops must never run"
2621 );
2622 }
2623
2624 #[tokio::test]
2625 async fn top_level_panic_reports_unknown_and_fails_queue_without_running_it() {
2626 let dir = tempfile::tempdir().unwrap();
2627 let path = dir.path().join("writer_task_top_level_panic.db");
2628 let pool = file_pool(&path);
2629 {
2630 let writer = pool.try_writer().unwrap();
2631 writer
2632 .conn()
2633 .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY, v TEXT)")
2634 .unwrap();
2635 }
2636 let handle = spawn(&pool, 8).expect("writer task spawn");
2637
2638 let (started_tx, started_rx) = oneshot::channel::<()>();
2639 let (release_tx, release_rx) = std_mpsc::channel::<()>();
2640 let active = tokio::spawn({
2641 let handle = handle.clone();
2642 async move {
2643 handle
2644 .send_top_level(move |conn| -> Result<usize, StorageError> {
2645 conn.execute("INSERT INTO t (id, v) VALUES (10, 'autocommitted')", [])
2646 .map_err(|e| StorageError::Pool {
2647 operation: "test_top_level_insert".into(),
2648 message: e.to_string(),
2649 })?;
2650 let _ = started_tx.send(());
2651 release_rx.recv().expect("test must release top-level op");
2652 panic!("intentional top-level writer request panic");
2653 })
2654 .await
2655 }
2656 });
2657
2658 tokio::time::timeout(Duration::from_secs(5), started_rx)
2659 .await
2660 .expect("top-level request did not start")
2661 .expect("top-level request dropped its start signal");
2662
2663 let queued_ran = Arc::new(AtomicBool::new(false));
2664 let queued_ran_in_op = Arc::clone(&queued_ran);
2665 let queued = handle
2666 .enqueue(move |conn| {
2667 queued_ran_in_op.store(true, Ordering::SeqCst);
2668 conn.execute("INSERT INTO t (id, v) VALUES (11, 'queued')", [])
2669 .map_err(|e| StorageError::Pool {
2670 operation: "test_top_level_queued_insert".into(),
2671 message: e.to_string(),
2672 })
2673 })
2674 .await
2675 .expect("queued request must be accepted");
2676 assert_eq!(handle.queue_depth(), 1);
2677 release_tx.send(()).expect("release top-level op");
2678
2679 let active_result = tokio::time::timeout(Duration::from_secs(5), active)
2680 .await
2681 .expect("top-level caller hung after panic")
2682 .expect("top-level caller task join");
2683 assert_writer_task_terminal_state(
2684 active_result,
2685 WriterTaskRequestState::SideEffectsUnknown,
2686 );
2687
2688 let queued_result = tokio::time::timeout(Duration::from_secs(5), queued)
2689 .await
2690 .expect("queued caller hung after top-level panic")
2691 .expect("terminal drain must preserve queued typed reply");
2692 assert_writer_task_terminal_state(queued_result, WriterTaskRequestState::NotStarted);
2693 assert!(!queued_ran.load(Ordering::SeqCst));
2694
2695 let reader = pool.reader().expect("reader");
2696 let active_count: i64 = reader
2697 .conn()
2698 .query_row("SELECT COUNT(*) FROM t WHERE id = 10", [], |row| row.get(0))
2699 .unwrap();
2700 let queued_count: i64 = reader
2701 .conn()
2702 .query_row("SELECT COUNT(*) FROM t WHERE id = 11", [], |row| row.get(0))
2703 .unwrap();
2704 assert_eq!(
2705 active_count, 1,
2706 "the completed top-level statement autocommits before the panic"
2707 );
2708 assert_eq!(queued_count, 0, "the queued request must never run");
2709 }
2710
2711 #[tokio::test]
2712 async fn closed_receiver_rejects_all_send_surfaces_as_not_started() {
2713 let (tx, rx) = mpsc::channel::<Box<dyn AnyWriteRequest + Send>>(4);
2717 drop(rx);
2718
2719 let handle = WriterTaskHandle {
2720 tx,
2721 backend_key: None,
2722 db: "test".to_string(),
2723 slow_write_threshold: None,
2724 enqueue_timeout: Duration::from_secs(5),
2725 };
2726 let send_result = handle.send(|_conn| Ok::<(), StorageError>(())).await;
2727 assert_writer_task_terminal_state(send_result, WriterTaskRequestState::NotStarted);
2728
2729 let timed_result = handle
2730 .send_with_timeout(|_conn| Ok::<(), StorageError>(()), Duration::from_secs(1))
2731 .await;
2732 assert_writer_task_terminal_state(timed_result, WriterTaskRequestState::NotStarted);
2733
2734 let top_level_result = handle
2735 .send_top_level(|_conn| Ok::<(), StorageError>(()))
2736 .await;
2737 assert_writer_task_terminal_state(top_level_result, WriterTaskRequestState::NotStarted);
2738 }
2739
2740 #[tokio::test]
2741 async fn accepted_request_lost_reply_is_side_effects_unknown() {
2742 let (tx, mut rx) = mpsc::channel::<Box<dyn AnyWriteRequest + Send>>(1);
2746 let handle = WriterTaskHandle {
2747 tx,
2748 backend_key: None,
2749 db: "test".to_string(),
2750 slow_write_threshold: None,
2751 enqueue_timeout: Duration::from_secs(5),
2752 };
2753 let request_ran = Arc::new(AtomicBool::new(false));
2754 let request_ran_in_op = Arc::clone(&request_ran);
2755
2756 let dropper = tokio::spawn(async move {
2757 let request = rx.recv().await.expect("request must be accepted");
2758 drop(request);
2759 });
2760 let result = tokio::time::timeout(
2761 Duration::from_secs(5),
2762 handle.send(move |_conn| {
2763 request_ran_in_op.store(true, Ordering::SeqCst);
2764 Ok::<(), StorageError>(())
2765 }),
2766 )
2767 .await
2768 .expect("caller hung after accepted request was dropped");
2769 dropper.await.expect("dropper task join");
2770
2771 assert_writer_task_terminal_state(result, WriterTaskRequestState::SideEffectsUnknown);
2772 assert!(!request_ran.load(Ordering::SeqCst));
2773 }
2774
2775 #[cfg(unix)]
2778 #[test]
2779 fn writer_stage_backend_key_preserves_non_utf8_path_bytes() {
2780 use std::ffi::OsString;
2781 use std::os::unix::ffi::OsStringExt;
2782
2783 let path_a =
2784 std::path::PathBuf::from(OsString::from_vec(b"/tmp/khive-writer-\x80.db".to_vec()));
2785 let path_b =
2786 std::path::PathBuf::from(OsString::from_vec(b"/tmp/khive-writer-\x81.db".to_vec()));
2787 assert_eq!(
2788 path_a.display().to_string(),
2789 path_b.display().to_string(),
2790 "fixture must reproduce the lossy display-label collision"
2791 );
2792 assert_ne!(
2793 writer_db_key_from_path(Some(&path_a)),
2794 writer_db_key_from_path(Some(&path_b)),
2795 "backend keys must retain the canonical path's exact OS bytes"
2796 );
2797 }
2798
2799 #[tokio::test]
2803 async fn writer_stage_sample_attributes_a_slow_body() {
2804 let dir = tempfile::tempdir().unwrap();
2805 let path = dir.path().join("writer_stage_sample.db");
2806 let pool = ConnectionPool::new(PoolConfig {
2809 path: Some(path.clone()),
2810 volume_lock_dir: Some(dir.path().join("volume-locks")),
2811 ..PoolConfig::for_test()
2812 })
2813 .expect("pool open");
2814 {
2815 let writer = pool.try_writer().unwrap();
2816 writer
2817 .conn()
2818 .execute_batch("CREATE TABLE t (id INTEGER PRIMARY KEY)")
2819 .unwrap();
2820 }
2821 let handle = spawn(&pool, 8).unwrap();
2822
2823 handle
2827 .send(|conn| {
2828 std::thread::sleep(Duration::from_millis(400));
2829 conn.execute("INSERT INTO t VALUES (1)", [])
2830 .map_err(|error| StorageError::Pool {
2831 operation: "writer_stage_sample".into(),
2832 message: error.to_string(),
2833 })
2834 })
2835 .await
2836 .unwrap();
2837
2838 let sample = last_writer_stage_observation(&pool).expect("writer stage sample");
2839 assert!(
2840 sample.body_micros >= 350_000,
2841 "the synthetic delay must land in the body stage: {sample:?}"
2842 );
2843 assert!(
2844 sample.body_micros > sample.queue_wait_micros,
2845 "fast queueing must not receive the body's delay: {sample:?}"
2846 );
2847 assert!(
2848 sample.body_micros > sample.transaction_acquire_micros,
2849 "an uncontended BEGIN must not receive the body's delay: {sample:?}"
2850 );
2851 assert!(
2852 sample.body_micros > sample.commit_micros,
2853 "a fast COMMIT must not receive the body's delay: {sample:?}"
2854 );
2855 assert!(sample.observed_at_unix_ms > 0);
2856 }
2857
2858 #[test]
2863 fn writer_queue_wait_excludes_blocking_pool_scheduling_delay() {
2864 let runtime = tokio::runtime::Builder::new_multi_thread()
2865 .worker_threads(1)
2866 .max_blocking_threads(1)
2867 .enable_all()
2868 .build()
2869 .expect("test runtime");
2870
2871 runtime.block_on(async {
2872 let dir = tempfile::tempdir().unwrap();
2873 let path = dir.path().join("writer_dequeue_boundary.db");
2874 let pool = file_pool(&path);
2875 let handle = spawn(&pool, 8).unwrap();
2876
2877 let (blocker_started_tx, blocker_started_rx) = std_mpsc::sync_channel(0);
2878 let (release_blocker_tx, release_blocker_rx) = std_mpsc::channel();
2879 let blocker = tokio::task::spawn_blocking(move || {
2880 blocker_started_tx.send(()).unwrap();
2881 release_blocker_rx.recv().unwrap();
2882 });
2883 blocker_started_rx
2884 .recv_timeout(Duration::from_secs(1))
2885 .expect("sole blocking worker must be occupied");
2886
2887 let reply = handle
2888 .enqueue(|_conn| Ok::<(), StorageError>(()))
2889 .await
2890 .expect("request must enter the bounded writer channel");
2891 let dequeue_deadline = Instant::now() + Duration::from_secs(1);
2892 while handle.queue_depth() != 0 {
2893 assert!(
2894 Instant::now() < dequeue_deadline,
2895 "writer drain never dequeued the accepted request"
2896 );
2897 tokio::task::yield_now().await;
2898 }
2899
2900 let scheduling_delay = Duration::from_millis(150);
2901 tokio::time::sleep(scheduling_delay).await;
2902 release_blocker_tx.send(()).unwrap();
2903 blocker.await.unwrap();
2904 reply.await.unwrap().unwrap();
2905
2906 let sample = last_writer_stage_observation(&pool).expect("writer stage sample");
2907 assert!(
2908 sample.total_micros.saturating_sub(sample.queue_wait_micros) >= 100_000,
2909 "the post-dequeue blocking-pool delay must not inflate queue_wait: {sample:?}"
2910 );
2911 });
2912 }
2913
2914 #[tokio::test]
2918 #[serial(tx_registry)]
2919 async fn writer_task_failure_counters_are_acquisition_site_exact() {
2920 {
2926 let dir = tempfile::tempdir().unwrap();
2927 let path = dir.path().join("writer_task_failure_counters_rollback.db");
2928 let pool = file_pool(&path);
2929 let handle = spawn(&pool, 8).expect("writer task should spawn");
2930
2931 let before = pool.writer_acquisition_snapshot();
2932 assert_eq!(before.writer_task_request_failures, 0);
2933 assert_eq!(before.writer_task_side_effects_unknown, 0);
2934
2935 let (started_tx, started_rx) = oneshot::channel::<()>();
2936 let (release_tx, release_rx) = std_mpsc::channel::<()>();
2937 let active = tokio::spawn({
2938 let handle = handle.clone();
2939 async move {
2940 handle
2941 .send(move |_conn| -> Result<(), StorageError> {
2942 let _ = started_tx.send(());
2943 release_rx.recv().expect("test must release active op");
2944 panic!("intentional rollback-clean panic for counter test");
2945 })
2946 .await
2947 }
2948 });
2949 tokio::time::timeout(Duration::from_secs(5), started_rx)
2950 .await
2951 .expect("active request did not start")
2952 .expect("active request dropped its start signal");
2953
2954 let queued = handle
2955 .enqueue(|_conn| Ok::<(), StorageError>(()))
2956 .await
2957 .expect("second request must queue behind the active one");
2958
2959 release_tx.send(()).expect("release active op");
2960 let active_result = tokio::time::timeout(Duration::from_secs(5), active)
2961 .await
2962 .expect("active caller hung after panic")
2963 .expect("active caller task join");
2964 assert_writer_task_terminal_state(
2965 active_result,
2966 WriterTaskRequestState::TransactionRolledBack,
2967 );
2968
2969 let queued_result = queued.await.expect("terminal drain must reply");
2970 assert_writer_task_terminal_state(queued_result, WriterTaskRequestState::NotStarted);
2971
2972 let after = pool.writer_acquisition_snapshot();
2973 assert_eq!(
2974 after.writer_task_request_failures, 1,
2975 "only the request that actually reached the seam counts, not the ones \
2976 failed by the queue-close drain"
2977 );
2978 assert_eq!(
2979 after.writer_task_side_effects_unknown, 0,
2980 "a clean rollback must not be counted as an unknown-side-effects outcome"
2981 );
2982 }
2983
2984 {
2988 let dir = tempfile::tempdir().unwrap();
2989 let path = dir.path().join("writer_task_failure_counters_unknown.db");
2990 let pool = file_pool(&path);
2991 let handle = spawn(&pool, 8).expect("writer task should spawn");
2992
2993 let before = pool.writer_acquisition_snapshot();
2994
2995 let active_result = handle
2996 .send(|conn| -> Result<(), StorageError> {
2997 conn.authorizer(Some(deny_rollback))
2998 .map_err(|e| StorageError::Pool {
2999 operation: "test_install_authorizer".into(),
3000 message: e.to_string(),
3001 })?;
3002 Err(StorageError::Internal(
3003 "intentional operation failure before denied rollback".into(),
3004 ))
3005 })
3006 .await;
3007 assert_writer_task_terminal_state(
3008 active_result,
3009 WriterTaskRequestState::SideEffectsUnknown,
3010 );
3011
3012 let sentinel_result = handle.send(|_conn| Ok::<(), StorageError>(())).await;
3019 assert_writer_task_terminal_state(sentinel_result, WriterTaskRequestState::NotStarted);
3020
3021 let after = pool.writer_acquisition_snapshot();
3022 assert_eq!(
3023 after.writer_task_request_failures - before.writer_task_request_failures,
3024 1
3025 );
3026 assert_eq!(
3027 after.writer_task_side_effects_unknown - before.writer_task_side_effects_unknown,
3028 1
3029 );
3030 }
3031 }
3032}
3033
3034#[cfg(test)]
3035#[path = "writer_task_lease_close_tests.rs"]
3036mod volume_lease_close_tests;
3037
3038#[cfg(all(test, any(unix, windows)))]
3039#[path = "writer_task_identity_tests.rs"]
3040mod identity_tests;