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 telemetry: WriteTelemetry,
213}
214
215mod sealed {
216 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
244pub trait AnyWriteRequest: sealed::Sealed + Send {
250 fn execute_and_reply(self: Box<Self>, conn: &Connection);
263
264 fn execute_and_reply_top_level(self: Box<Self>, conn: &Connection);
274
275 fn reply_error(self: Box<Self>, err: StorageError);
287
288 fn is_top_level(&self) -> bool;
292
293 fn queue_wait(&self) -> Duration;
296}
297
298#[derive(Debug, Clone, Copy, PartialEq, Eq)]
299enum RollbackDisposition {
300 RolledBack,
301 SideEffectsUnknown,
302}
303
304fn 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
329pub(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 let WriteRequest {
463 op,
464 reply,
465 telemetry,
466 ..
467 } = *self;
468 let profiled = execute_wrapped_transaction_profiled(conn, "writer_task_commit", op);
469 drop(tx_span);
474 telemetry.finish(
475 queue_wait,
476 transaction_acquire,
477 profiled.body,
478 profiled.commit,
479 );
480 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 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 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#[derive(Clone, Debug)]
600pub struct WriterTaskHandle {
601 tx: mpsc::Sender<Box<dyn AnyWriteRequest + Send>>,
602 backend_key: Option<PathBuf>,
605 db: String,
610 slow_write_threshold: Option<std::time::Duration>,
616 enqueue_timeout: std::time::Duration,
637}
638
639impl WriterTaskHandle {
640 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 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 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 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 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 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 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 pub fn queue_depth(&self) -> usize {
827 self.tx.max_capacity() - self.tx.capacity()
828 }
829
830 pub fn capacity(&self) -> usize {
833 self.tx.max_capacity()
834 }
835}
836
837pub fn spawn(pool: &ConnectionPool, capacity: usize) -> Result<WriterTaskHandle, SqliteError> {
859 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 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
892async 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
906fn 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 acquisition_counters.record_writer_task_begin_busy();
929 let Some(delay) = retry_delays.next() else {
930 break Err(error);
931 };
932 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 break Err(error);
952 }
953 busy_timeout_lowered = true;
954 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 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
988async 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 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 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 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 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 drop(tx_span);
1084 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 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 #[tokio::test]
1405 #[serial(tx_registry)]
1406 async fn begin_immediate_failure_replies_error_without_running_op() {
1407 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 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 #[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 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 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1565 #[serial(tx_registry)]
1566 async fn transient_begin_refusal_retries_once_and_restores_timeout() {
1567 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 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 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 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 #[tokio::test]
1767 #[serial(tx_registry)]
1768 async fn contended_begin_exhaustion_separates_absorbed_and_surfaced_refusals() {
1769 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 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 #[tokio::test]
1883 #[serial(tx_registry)]
1884 async fn begin_retry_budget_makes_exactly_one_attempt_under_sustained_contention() {
1885 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 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 #[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 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 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 assert_eq!(
2031 read_pages(&handle).await,
2032 crate::pool::FALLBACK_WAL_AUTOCHECKPOINT_PAGES
2033 );
2034
2035 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 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 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 let first = tokio::spawn({
2081 let handle = handle.clone();
2082 async move {
2083 let _ = handle.send(|_conn| Ok::<(), StorageError>(())).await;
2084 }
2085 });
2086
2087 tokio::time::sleep(Duration::from_millis(20)).await;
2089
2090 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 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 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 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 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_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 #[tokio::test]
2225 #[serial(tx_registry)]
2226 async fn send_with_timeout_returns_op_result_when_op_outlives_the_timeout() {
2227 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 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 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 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 #[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 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 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 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 #[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 #[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 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 #[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 #[tokio::test]
3171 #[serial(tx_registry)]
3172 async fn writer_task_failure_counters_are_acquisition_site_exact() {
3173 {
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 {
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 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}