1use std::any::Any;
10use std::sync::Arc;
11use std::time::Instant;
12
13use async_trait::async_trait;
14
15use khive_storage::error::StorageError;
16use khive_storage::types::{PageRequest, SqlColumn, SqlRow, SqlStatement, SqlValue};
17use khive_storage::{AtomicUnitOp, StorageCapability};
18use tokio::sync::{OwnedSemaphorePermit, Semaphore};
19
20use crate::error::SqliteError;
21use crate::pool::ConnectionPool;
22
23fn row_to_sql_row(row: &rusqlite::Row<'_>, col_count: usize, col_names: &[String]) -> SqlRow {
29 #[cfg(test)]
30 ROW_CONVERSIONS.with(|count| count.set(count.get() + 1));
31
32 let mut columns = Vec::with_capacity(col_count);
33 for i in 0..col_count {
34 let value = match row.get_ref(i) {
35 Ok(rusqlite::types::ValueRef::Null) => SqlValue::Null,
36 Ok(rusqlite::types::ValueRef::Integer(v)) => SqlValue::Integer(v),
37 Ok(rusqlite::types::ValueRef::Real(v)) => SqlValue::Float(v),
38 Ok(rusqlite::types::ValueRef::Text(bytes)) => {
39 SqlValue::Text(String::from_utf8_lossy(bytes).into_owned())
40 }
41 Ok(rusqlite::types::ValueRef::Blob(bytes)) => SqlValue::Blob(bytes.to_vec()),
42 Err(_) => SqlValue::Null,
43 };
44 columns.push(SqlColumn {
45 name: col_names.get(i).cloned().unwrap_or_default(),
46 value,
47 });
48 }
49 SqlRow { columns }
50}
51
52#[cfg(test)]
53thread_local! {
54 static ROW_CONVERSIONS: std::cell::Cell<usize> = const { std::cell::Cell::new(0) };
55}
56
57pub(crate) fn bind_params(
64 stmt: &mut rusqlite::Statement<'_>,
65 params: &[SqlValue],
66) -> Result<(), rusqlite::Error> {
67 for (i, param) in params.iter().enumerate() {
68 let idx = i + 1; match param {
70 SqlValue::Null => stmt.raw_bind_parameter(idx, rusqlite::types::Null)?,
71 SqlValue::Bool(v) => stmt.raw_bind_parameter(idx, *v as i64)?,
72 SqlValue::Integer(v) => stmt.raw_bind_parameter(idx, *v)?,
73 SqlValue::Float(v) => stmt.raw_bind_parameter(idx, *v)?,
74 SqlValue::Text(v) => stmt.raw_bind_parameter(idx, v.as_str())?,
75 SqlValue::Blob(v) => stmt.raw_bind_parameter(idx, v.as_slice())?,
76 SqlValue::Json(v) => {
77 let s = serde_json::to_string(v).unwrap_or_default();
78 stmt.raw_bind_parameter(idx, s.as_str())?;
79 }
80 SqlValue::Uuid(v) => stmt.raw_bind_parameter(idx, v.to_string().as_str())?,
81 SqlValue::Timestamp(v) => {
82 stmt.raw_bind_parameter(idx, v.timestamp_micros())?;
83 }
84 }
85 }
86 Ok(())
87}
88
89fn prepare_sql_statement<'conn>(
97 conn: &'conn rusqlite::Connection,
98 sql: &str,
99) -> Result<rusqlite::Statement<'conn>, rusqlite::Error> {
100 conn.prepare(sql)
101}
102
103fn prepare_cached_sql_statement<'conn>(
107 conn: &'conn rusqlite::Connection,
108 sql: &str,
109) -> Result<rusqlite::CachedStatement<'conn>, rusqlite::Error> {
110 conn.prepare_cached(sql)
111}
112
113enum PreparedBatchStatement<'conn> {
119 Ready(rusqlite::Statement<'conn>),
120 PrepareAtExecution,
121}
122
123fn prepare_batch_statements<'conn>(
126 conn: &'conn rusqlite::Connection,
127 statements: &[SqlStatement],
128) -> Result<Vec<PreparedBatchStatement<'conn>>, rusqlite::Error> {
129 let mut prepared = Vec::with_capacity(statements.len());
130 for statement in statements {
131 match prepare_sql_statement(conn, &statement.sql) {
132 Ok(statement) => prepared.push(PreparedBatchStatement::Ready(statement)),
133 Err(error @ rusqlite::Error::MultipleStatement) => return Err(error),
134 Err(_) => prepared.push(PreparedBatchStatement::PrepareAtExecution),
135 }
136 }
137 Ok(prepared)
138}
139
140fn execute_prepared_batch<'conn>(
142 conn: &'conn rusqlite::Connection,
143 prepared: Vec<PreparedBatchStatement<'conn>>,
144 statements: &[SqlStatement],
145) -> Result<u64, rusqlite::Error> {
146 debug_assert_eq!(prepared.len(), statements.len());
147 let mut total = 0u64;
148 for (prepared, statement) in prepared.into_iter().zip(statements) {
149 let mut prepared = match prepared {
150 PreparedBatchStatement::Ready(prepared) => prepared,
151 PreparedBatchStatement::PrepareAtExecution => {
152 prepare_sql_statement(conn, &statement.sql)?
153 }
154 };
155 bind_params(&mut prepared, &statement.params)?;
156 total += prepared.raw_execute()? as u64;
157 }
158 Ok(total)
159}
160
161const TRANSACTION_CONTROL_KEYWORDS: [&str; 7] = [
172 "BEGIN",
173 "START",
174 "COMMIT",
175 "END",
176 "ROLLBACK",
177 "SAVEPOINT",
178 "RELEASE",
179];
180
181fn skip_sqlite_empty_prefix(mut rest: &[u8]) -> &[u8] {
184 loop {
185 let mut idx = 0;
186 while idx < rest.len() && rest[idx].is_ascii_whitespace() {
187 idx += 1;
188 }
189 rest = &rest[idx..];
190 if let Some(tail) = rest.strip_prefix(b"\xEF\xBB\xBF") {
191 rest = tail;
192 continue;
193 }
194 if let Some(tail) = rest.strip_prefix(b";") {
195 rest = tail;
196 continue;
197 }
198 if let Some(tail) = rest.strip_prefix(b"--") {
199 let mut idx = 0;
200 while idx < tail.len() && tail[idx] != b'\n' {
201 idx += 1;
202 }
203 rest = if idx < tail.len() {
204 &tail[idx + 1..]
205 } else {
206 &[]
207 };
208 continue;
209 }
210 if let Some(tail) = rest.strip_prefix(b"/*") {
211 let mut idx = 0;
212 while idx + 1 < tail.len() && !(tail[idx] == b'*' && tail[idx + 1] == b'/') {
213 idx += 1;
214 }
215 rest = if idx + 1 < tail.len() {
216 &tail[idx + 2..]
217 } else {
218 &[]
219 };
220 continue;
221 }
222 break;
223 }
224 rest
225}
226
227fn next_sqlite_token(mut rest: &[u8]) -> Option<(&[u8], &[u8])> {
231 loop {
232 let mut idx = 0;
233 while idx < rest.len() && rest[idx].is_ascii_whitespace() {
234 idx += 1;
235 }
236 rest = &rest[idx..];
237 if let Some(tail) = rest.strip_prefix(b"\xEF\xBB\xBF") {
238 rest = tail;
239 continue;
240 }
241 if let Some(tail) = rest.strip_prefix(b"--") {
242 let mut idx = 0;
243 while idx < tail.len() && tail[idx] != b'\n' {
244 idx += 1;
245 }
246 rest = if idx < tail.len() {
247 &tail[idx + 1..]
248 } else {
249 &[]
250 };
251 continue;
252 }
253 if let Some(tail) = rest.strip_prefix(b"/*") {
254 let mut idx = 0;
255 while idx + 1 < tail.len() && !(tail[idx] == b'*' && tail[idx + 1] == b'/') {
256 idx += 1;
257 }
258 rest = if idx + 1 < tail.len() {
259 &tail[idx + 2..]
260 } else {
261 &[]
262 };
263 continue;
264 }
265 break;
266 }
267
268 let len = rest
269 .iter()
270 .take_while(|byte| byte.is_ascii_alphanumeric() || **byte == b'_')
271 .count();
272 (len != 0).then_some((&rest[..len], &rest[len..]))
273}
274
275fn transaction_control_parts(sql: &str) -> Option<(&'static str, &[u8])> {
279 let rest = skip_sqlite_empty_prefix(sql.as_bytes());
280 TRANSACTION_CONTROL_KEYWORDS
281 .iter()
282 .copied()
283 .find_map(|keyword| {
284 let kw = keyword.as_bytes();
285 if rest.len() < kw.len() || !rest[..kw.len()].eq_ignore_ascii_case(kw) {
286 return None;
287 }
288 let boundary = match rest.get(kw.len()) {
289 Some(next) => !(next.is_ascii_alphanumeric() || *next == b'_'),
290 None => true,
291 };
292 boundary.then_some((keyword, &rest[kw.len()..]))
293 })
294}
295
296fn transaction_control_head(sql: &str) -> Option<&'static str> {
298 transaction_control_parts(sql).map(|(keyword, _)| keyword)
299}
300
301#[derive(Clone, Copy, Debug, PartialEq, Eq)]
302enum CachedReadTransactionControl {
303 BeginDeferred,
305 Finish(&'static str),
307 Unsupported(&'static str),
310}
311
312fn cached_read_transaction_control(sql: &str) -> Option<CachedReadTransactionControl> {
320 let (keyword, tail) = transaction_control_parts(sql)?;
321 match keyword {
322 "BEGIN" => {
323 let mut rest = tail;
340 let mut saw_deferred = false;
341 let mut saw_transaction = false;
342 while let Some((token, next)) = next_sqlite_token(rest) {
343 if !saw_deferred && !saw_transaction && token.eq_ignore_ascii_case(b"DEFERRED") {
344 saw_deferred = true;
345 } else if !saw_transaction && token.eq_ignore_ascii_case(b"TRANSACTION") {
346 saw_transaction = true;
347 } else {
348 return Some(CachedReadTransactionControl::Unsupported(keyword));
349 }
350 rest = next;
351 }
352 if !skip_sqlite_empty_prefix(rest).is_empty() {
353 return Some(CachedReadTransactionControl::Unsupported(keyword));
354 }
355 Some(CachedReadTransactionControl::BeginDeferred)
356 }
357 "COMMIT" | "END" => Some(CachedReadTransactionControl::Finish(keyword)),
358 "ROLLBACK" => {
359 let first = next_sqlite_token(tail);
360 let rollback_target = match first {
361 Some((token, rest)) if token.eq_ignore_ascii_case(b"TRANSACTION") => {
362 next_sqlite_token(rest).map(|(token, _)| token)
363 }
364 Some((token, _)) => Some(token),
365 None => None,
366 };
367 if rollback_target.is_some_and(|token| token.eq_ignore_ascii_case(b"TO")) {
368 Some(CachedReadTransactionControl::Unsupported(keyword))
369 } else {
370 Some(CachedReadTransactionControl::Finish(keyword))
371 }
372 }
373 _ => Some(CachedReadTransactionControl::Unsupported(keyword)),
374 }
375}
376
377fn reject_transaction_control_statements(
381 statements: &[SqlStatement],
382 operation: &'static str,
383) -> khive_storage::types::StorageResult<()> {
384 for (index, statement) in statements.iter().enumerate() {
385 if let Some(keyword) = transaction_control_head(&statement.sql) {
386 return Err(StorageError::InvalidInput {
387 capability: StorageCapability::Sql,
388 operation: operation.into(),
389 message: format!(
390 "statement at index {index} is transaction control ({keyword}); \
391 execute_batch owns the BEGIN/COMMIT boundary for the whole \
392 batch — remove transaction-control statements from the batch"
393 ),
394 });
395 }
396 }
397 Ok(())
398}
399
400struct BatchFailure {
403 error: rusqlite::Error,
404 poison_reason: Option<BatchPoisonReason>,
405}
406
407#[derive(Clone, Copy, Debug, PartialEq, Eq)]
408enum BatchHandleDisposition {
409 Retain,
410 Poison,
411}
412
413fn execute_standalone_batch(
418 conn: &rusqlite::Connection,
419 statements: &[SqlStatement],
420 origin: khive_storage::tx_registry::TxOrigin,
421) -> (BatchHandleDisposition, Result<u64, BatchFailure>) {
422 let prepared = match prepare_batch_statements(conn, statements) {
423 Ok(prepared) => prepared,
424 Err(error) => {
425 return (
426 BatchHandleDisposition::Retain,
427 Err(BatchFailure {
428 error,
429 poison_reason: None,
430 }),
431 );
432 }
433 };
434 if let Err(begin_error) = conn.execute_batch("BEGIN IMMEDIATE") {
435 drop(prepared);
440 let (disposition, poison_reason) = if crate::timeout_sink::is_busy_or_locked(&begin_error) {
441 (BatchHandleDisposition::Retain, None)
442 } else {
443 tracing::warn!(
444 %begin_error,
445 "execute_batch: BEGIN IMMEDIATE failed non-transiently; \
446 poisoning the standalone connection — the handle is \
447 dropped and must be re-acquired"
448 );
449 (
450 BatchHandleDisposition::Poison,
451 Some(BatchPoisonReason::BeginFailed),
452 )
453 };
454 return (
455 disposition,
456 Err(BatchFailure {
457 error: begin_error,
458 poison_reason,
459 }),
460 );
461 }
462
463 let _tx_handle =
466 khive_storage::tx_registry::register_scoped(Some("execute_batch".to_string()), origin);
467 let result = (|| -> Result<u64, rusqlite::Error> {
468 let total = execute_prepared_batch(conn, prepared, statements)?;
469 conn.execute_batch("COMMIT")?;
470 Ok(total)
471 })();
472
473 let mut disposition = BatchHandleDisposition::Retain;
474 let mut poison_reason = None;
475 if let Err(error) = &result {
476 if let Err(rollback_error) = conn.execute_batch("ROLLBACK") {
477 tracing::warn!(
481 %error,
482 %rollback_error,
483 "execute_batch: ROLLBACK after statement failure failed; \
484 poisoning the standalone connection — the handle is \
485 dropped and must be re-acquired"
486 );
487 disposition = BatchHandleDisposition::Poison;
488 poison_reason = Some(BatchPoisonReason::RollbackFailed(rollback_error));
489 }
490 }
491
492 (
493 disposition,
494 result.map_err(|error| BatchFailure {
495 error,
496 poison_reason,
497 }),
498 )
499}
500
501#[derive(Debug)]
502enum BatchPoisonReason {
503 BeginFailed,
504 RollbackFailed(rusqlite::Error),
505}
506
507impl std::fmt::Display for BatchPoisonReason {
508 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
509 match self {
510 Self::BeginFailed => f.write_str(
511 "BEGIN IMMEDIATE failed non-transiently; connection transaction state is suspect",
512 ),
513 Self::RollbackFailed(error) => {
514 write!(f, "ROLLBACK after statement failure failed: {error}")
515 }
516 }
517 }
518}
519
520#[derive(Debug)]
525struct PoisonedBatchError {
526 original: rusqlite::Error,
527 poison_reason: BatchPoisonReason,
528}
529
530impl std::fmt::Display for PoisonedBatchError {
531 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
532 write!(
533 f,
534 "{}; original error: {}",
535 self.poison_reason, self.original
536 )
537 }
538}
539
540impl std::error::Error for PoisonedBatchError {
541 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
542 Some(&self.original)
543 }
544}
545
546fn prepare_bound_statement<'conn>(
547 conn: &'conn rusqlite::Connection,
548 statement: &SqlStatement,
549) -> Result<rusqlite::Statement<'conn>, rusqlite::Error> {
550 let mut stmt = prepare_sql_statement(conn, &statement.sql)?;
551 bind_params(&mut stmt, &statement.params)?;
552 Ok(stmt)
553}
554
555fn execute_prepared_query(
556 mut stmt: rusqlite::Statement<'_>,
557) -> Result<Vec<SqlRow>, rusqlite::Error> {
558 let col_count = stmt.column_count();
559 let col_names: Vec<String> = (0..col_count)
560 .map(|i| stmt.column_name(i).unwrap_or("").to_string())
561 .collect();
562
563 let mut rows = Vec::new();
564 let mut raw_rows = stmt.raw_query();
565 while let Some(row) = raw_rows.next()? {
566 rows.push(row_to_sql_row(row, col_count, &col_names));
567 }
568 Ok(rows)
569}
570
571fn execute_prepared_query_row(
572 mut stmt: rusqlite::Statement<'_>,
573) -> Result<Option<SqlRow>, rusqlite::Error> {
574 let col_count = stmt.column_count();
575 let col_names: Vec<String> = (0..col_count)
576 .map(|i| stmt.column_name(i).unwrap_or("").to_string())
577 .collect();
578
579 let mut raw_rows = stmt.raw_query();
580 Ok(raw_rows
581 .next()?
582 .map(|row| row_to_sql_row(row, col_count, &col_names)))
583}
584
585fn execute_prepared_query_page(
586 mut stmt: rusqlite::Statement<'_>,
587 page: &PageRequest,
588) -> Result<Vec<SqlRow>, rusqlite::Error> {
589 if page.limit == 0 {
593 return Ok(Vec::new());
594 }
595
596 let col_count = stmt.column_count();
597 let col_names: Vec<String> = (0..col_count)
598 .map(|i| stmt.column_name(i).unwrap_or("").to_string())
599 .collect();
600
601 let mut rows = Vec::new();
602 let mut offset = page.offset;
603 let mut remaining = u64::from(page.limit);
604 let mut raw_rows = stmt.raw_query();
605 while remaining > 0 {
615 let Some(row) = raw_rows.next()? else {
616 break;
617 };
618 if offset > 0 {
619 offset -= 1;
620 continue;
621 }
622 rows.push(row_to_sql_row(row, col_count, &col_names));
623 remaining -= 1;
624 }
625 Ok(rows)
626}
627
628fn execute_query(
630 conn: &rusqlite::Connection,
631 statement: &SqlStatement,
632) -> Result<Vec<SqlRow>, rusqlite::Error> {
633 execute_prepared_query(prepare_bound_statement(conn, statement)?)
634}
635
636fn execute_query_row(
637 conn: &rusqlite::Connection,
638 statement: &SqlStatement,
639) -> Result<Option<SqlRow>, rusqlite::Error> {
640 execute_prepared_query_row(prepare_bound_statement(conn, statement)?)
641}
642
643fn execute_query_page(
644 conn: &rusqlite::Connection,
645 statement: &SqlStatement,
646 page: &PageRequest,
647) -> Result<Vec<SqlRow>, rusqlite::Error> {
648 execute_prepared_query_page(prepare_bound_statement(conn, statement)?, page)
649}
650
651fn statement_is_cancellable_read(stmt: &rusqlite::Statement<'_>, sql: &str) -> bool {
656 stmt.readonly() && transaction_control_head(sql).is_none()
657}
658
659fn execute_query_interruptibly(
660 scope: &crate::read_cancellation::InterruptibleReadScope,
661 conn: &rusqlite::Connection,
662 statement: &SqlStatement,
663 operation: &'static str,
664 rollback_interrupted_transaction: bool,
665 interruptible: bool,
666) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
667 let stmt = prepare_bound_statement(conn, statement)
668 .map_err(|error| map_rusqlite_err(error, operation))?;
669 if interruptible && statement_is_cancellable_read(&stmt, &statement.sql) {
670 scope.run_with_interrupted_cleanup(
671 conn,
672 move || {
673 execute_prepared_query(stmt).map_err(|error| map_rusqlite_err(error, operation))
674 },
675 || {
676 rollback_interrupted_read_transaction(
677 conn,
678 operation,
679 rollback_interrupted_transaction,
680 )
681 },
682 )
683 } else {
684 scope.mark_write_committed()?;
685 execute_prepared_query(stmt).map_err(|error| map_rusqlite_err(error, operation))
686 }
687}
688
689fn execute_query_row_interruptibly(
690 scope: &crate::read_cancellation::InterruptibleReadScope,
691 conn: &rusqlite::Connection,
692 statement: &SqlStatement,
693 operation: &'static str,
694 rollback_interrupted_transaction: bool,
695 interruptible: bool,
696) -> khive_storage::types::StorageResult<Option<SqlRow>> {
697 let stmt = prepare_bound_statement(conn, statement)
698 .map_err(|error| map_rusqlite_err(error, operation))?;
699 if interruptible && statement_is_cancellable_read(&stmt, &statement.sql) {
700 scope.run_with_interrupted_cleanup(
701 conn,
702 move || {
703 execute_prepared_query_row(stmt).map_err(|error| map_rusqlite_err(error, operation))
704 },
705 || {
706 rollback_interrupted_read_transaction(
707 conn,
708 operation,
709 rollback_interrupted_transaction,
710 )
711 },
712 )
713 } else {
714 scope.mark_write_committed()?;
715 execute_prepared_query_row(stmt).map_err(|error| map_rusqlite_err(error, operation))
716 }
717}
718
719fn execute_query_page_interruptibly(
720 scope: &crate::read_cancellation::InterruptibleReadScope,
721 conn: &rusqlite::Connection,
722 statement: &SqlStatement,
723 page: &PageRequest,
724 operation: &'static str,
725 rollback_interrupted_transaction: bool,
726 interruptible: bool,
727) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
728 let stmt = prepare_bound_statement(conn, statement)
729 .map_err(|error| map_rusqlite_err(error, operation))?;
730 if interruptible && statement_is_cancellable_read(&stmt, &statement.sql) {
731 scope.run_with_interrupted_cleanup(
732 conn,
733 move || {
734 execute_prepared_query_page(stmt, page)
735 .map_err(|error| map_rusqlite_err(error, operation))
736 },
737 || {
738 rollback_interrupted_read_transaction(
739 conn,
740 operation,
741 rollback_interrupted_transaction,
742 )
743 },
744 )
745 } else {
746 scope.mark_write_committed()?;
747 execute_prepared_query_page(stmt, page).map_err(|error| map_rusqlite_err(error, operation))
748 }
749}
750
751fn rollback_interrupted_read_transaction(
752 conn: &rusqlite::Connection,
753 operation: &'static str,
754 enabled: bool,
755) -> khive_storage::types::StorageResult<()> {
756 if !enabled || conn.is_autocommit() {
757 return Ok(());
758 }
759 conn.execute_batch("ROLLBACK")
760 .map_err(|error| map_rusqlite_err(error, operation))?;
761 if conn.is_autocommit() {
762 Ok(())
763 } else {
764 Err(StorageError::Transaction {
765 operation: operation.into(),
766 message: "interrupted read transaction rollback did not restore autocommit".into(),
767 })
768 }
769}
770
771fn map_rusqlite_err(e: rusqlite::Error, op: &'static str) -> StorageError {
773 StorageError::driver(StorageCapability::Sql, op, e)
774}
775
776#[derive(Clone, Copy)]
781enum SlotTimeoutClass {
782 Admission,
783 ReaderContract,
784}
785
786async fn acquire_handle_slot(
787 slots: Arc<Semaphore>,
788 timeout: std::time::Duration,
789 operation: &'static str,
790 class: SlotTimeoutClass,
791) -> Result<OwnedSemaphorePermit, StorageError> {
792 tokio::time::timeout(timeout, slots.acquire_owned())
793 .await
794 .map_err(|_| match class {
795 SlotTimeoutClass::Admission => StorageError::AdmissionTimeout {
796 operation: operation.into(),
797 timeout_ms: u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX),
798 },
799 SlotTimeoutClass::ReaderContract => StorageError::Timeout {
800 operation: operation.into(),
801 },
802 })?
803 .map_err(|error| StorageError::Pool {
804 operation: operation.into(),
805 message: error.to_string(),
806 })
807}
808
809fn open_standalone_reader(pool: &ConnectionPool) -> Result<rusqlite::Connection, StorageError> {
814 pool.open_standalone_reader()
815 .map_err(|error| StorageError::driver(StorageCapability::Sql, "open_reader", error))
816}
817
818fn open_standalone_writer(pool: &ConnectionPool) -> Result<rusqlite::Connection, StorageError> {
819 let config = pool.config();
820 let conn = pool
821 .open_standalone_writer()
822 .map_err(|e| StorageError::driver(StorageCapability::Sql, "open_writer", e))?;
823
824 conn.busy_timeout(config.busy_timeout)
825 .map_err(|e| map_rusqlite_err(e, "open_writer"))?;
826 conn.pragma_update(None, "cache_size", "-65536")
827 .map_err(|e| map_rusqlite_err(e, "open_writer"))?;
828 conn.pragma_update(None, "mmap_size", "1073741824")
829 .map_err(|e| map_rusqlite_err(e, "open_writer"))?;
830
831 Ok(conn)
832}
833
834async fn open_standalone_on_blocking<F>(
841 pool: Arc<ConnectionPool>,
842 slot: OwnedSemaphorePermit,
843 operation: &'static str,
844 open: F,
845) -> khive_storage::types::StorageResult<(rusqlite::Connection, OwnedSemaphorePermit)>
846where
847 F: FnOnce(&ConnectionPool) -> Result<rusqlite::Connection, StorageError> + Send + 'static,
848{
849 tokio::task::spawn_blocking(move || open(&pool).map(|conn| (conn, slot)))
850 .await
851 .map_err(|e| StorageError::driver(StorageCapability::Sql, operation, e))?
852}
853
854async fn open_standalone_reader_on_blocking(
864 pool: Arc<ConnectionPool>,
865 slot: OwnedSemaphorePermit,
866) -> khive_storage::types::StorageResult<(rusqlite::Connection, OwnedSemaphorePermit)> {
867 open_standalone_on_blocking(pool, slot, "open_reader", open_standalone_reader).await
868}
869
870async fn open_standalone_writer_on_blocking(
873 pool: Arc<ConnectionPool>,
874 slot: OwnedSemaphorePermit,
875) -> khive_storage::types::StorageResult<(rusqlite::Connection, OwnedSemaphorePermit)> {
876 open_standalone_on_blocking(pool, slot, "open_writer", open_standalone_writer).await
877}
878
879const CACHED_READ_TRANSACTION_LABEL: &str = "sql_bridge_cached_read_transaction";
884
885struct CachedReadTransaction {
895 _slot: OwnedSemaphorePermit,
896 _tx_handle: khive_storage::tx_registry::TxHandle,
897 opened_at: Instant,
903}
904
905struct StandaloneHandle {
906 conn: rusqlite::Connection,
907 _retained_slot: Option<OwnedSemaphorePermit>,
913 read_transaction_slot: Option<CachedReadTransaction>,
919}
920
921impl StandaloneHandle {
922 fn is_cached_reader(&self) -> bool {
925 self._retained_slot.is_none()
926 }
927
928 fn has_read_transaction(&self) -> bool {
929 self.read_transaction_slot.is_some()
930 }
931}
932
933struct SqliteReader {
934 handle: Option<StandaloneHandle>,
935 pool: Arc<ConnectionPool>,
936}
937
938async fn open_cached_reader_handle(
939 pool: Arc<ConnectionPool>,
940) -> khive_storage::types::StorageResult<StandaloneHandle> {
941 let open_slot = crate::await_request_read_phase(
942 "sql_bridge.reader_open",
943 acquire_handle_slot(
944 pool.sql_bridge_reader_slots(),
945 pool.config().checkout_timeout,
946 "sql_bridge.reader_open",
947 SlotTimeoutClass::ReaderContract,
948 ),
949 )
950 .await??;
951 let (conn, open_slot) = crate::await_request_read_phase(
952 "sql_bridge.reader_open",
953 open_standalone_reader_on_blocking(pool, open_slot),
954 )
955 .await??;
956 drop(open_slot);
957 Ok(StandaloneHandle {
958 conn,
959 _retained_slot: None,
960 read_transaction_slot: None,
961 })
962}
963
964async fn execute_standalone_read<R, F>(
979 handle: &mut Option<StandaloneHandle>,
980 pool: Arc<ConnectionPool>,
981 operation: &'static str,
982 transaction_control: Option<CachedReadTransactionControl>,
983 read: F,
984) -> khive_storage::types::StorageResult<R>
985where
986 R: Send + 'static,
987 F: FnOnce(
988 &crate::read_cancellation::InterruptibleReadScope,
989 &rusqlite::Connection,
990 bool,
991 bool,
992 ) -> khive_storage::types::StorageResult<R>
993 + Send
994 + 'static,
995{
996 if handle.is_none() {
997 return Err(StorageError::Pool {
998 operation: operation.into(),
999 message: "connection already consumed".into(),
1000 });
1001 }
1002 let active_read_transaction = handle
1003 .as_ref()
1004 .is_some_and(|handle| handle.is_cached_reader() && handle.has_read_transaction());
1005 let completion_preserving_writer_transaction = handle
1006 .as_ref()
1007 .is_some_and(|handle| !handle.is_cached_reader() && !handle.conn.is_autocommit());
1008 let mut operation_slot = if active_read_transaction {
1009 None
1010 } else if completion_preserving_writer_transaction {
1011 Some(
1015 acquire_handle_slot(
1016 pool.sql_bridge_reader_slots(),
1017 pool.config().checkout_timeout,
1018 "sql_bridge.reader_operation",
1019 SlotTimeoutClass::ReaderContract,
1020 )
1021 .await?,
1022 )
1023 } else {
1024 Some(
1025 crate::await_request_read_phase(
1026 "sql_bridge.reader_operation",
1027 acquire_handle_slot(
1028 pool.sql_bridge_reader_slots(),
1029 pool.config().checkout_timeout,
1030 "sql_bridge.reader_operation",
1031 SlotTimeoutClass::ReaderContract,
1032 ),
1033 )
1034 .await??,
1035 )
1036 };
1037 let Some(owned_handle) = handle.take() else {
1038 return Err(StorageError::Pool {
1039 operation: operation.into(),
1040 message: "connection already consumed".into(),
1041 });
1042 };
1043 let origin = pool.origin();
1044 let read_tx_max_age = pool.config().read_tx_max_age;
1045 let (owned_handle, result) = crate::read_cancellation::run_interruptible_read(
1046 StorageCapability::Sql,
1047 operation,
1048 move |scope| {
1049 let mut owned_handle = owned_handle;
1050 let cached_reader = owned_handle.is_cached_reader();
1051 let entered_with_transaction = owned_handle.has_read_transaction();
1052 let entered_autocommit = owned_handle.conn.is_autocommit();
1053 let mut restore_handle = true;
1054 let mut result = if cached_reader && entered_with_transaction && entered_autocommit {
1055 drop(owned_handle.read_transaction_slot.take());
1058 Err(StorageError::InvalidInput {
1059 capability: StorageCapability::Sql,
1060 operation: operation.into(),
1061 message: "cached read-only handle retained transaction admission after SQLite \
1062 had already returned to autocommit; the stale permit was released"
1063 .into(),
1064 })
1065 } else if cached_reader && !entered_with_transaction && !entered_autocommit {
1066 Err(StorageError::InvalidInput {
1067 capability: StorageCapability::Sql,
1068 operation: operation.into(),
1069 message: "cached read-only handle entered the operation outside autocommit; \
1070 its transaction was rolled back before releasing the reader permit"
1071 .into(),
1072 })
1073 } else if cached_reader
1074 && entered_with_transaction
1075 && owned_handle
1076 .read_transaction_slot
1077 .as_ref()
1078 .is_some_and(|tx| tx.opened_at.elapsed() >= read_tx_max_age)
1079 {
1080 crate::checkpoint::note_read_tx_max_age_eviction();
1085 match owned_handle.conn.execute_batch("ROLLBACK") {
1086 Ok(()) if owned_handle.conn.is_autocommit() => {
1087 drop(owned_handle.read_transaction_slot.take());
1088 Err(StorageError::ReadTransactionAgeEvicted {
1089 operation: operation.into(),
1090 max_age_secs: read_tx_max_age.as_secs(),
1091 })
1092 }
1093 Ok(()) => {
1094 restore_handle = false;
1095 Err(StorageError::ReadTransactionAgeEvictionCleanupFailed {
1096 operation: operation.into(),
1097 max_age_secs: read_tx_max_age.as_secs(),
1098 message: "rollback did not restore autocommit".into(),
1099 })
1100 }
1101 Err(error) => {
1102 restore_handle = false;
1103 Err(StorageError::ReadTransactionAgeEvictionCleanupFailed {
1104 operation: operation.into(),
1105 max_age_secs: read_tx_max_age.as_secs(),
1106 message: format!("rollback failed: {error}"),
1107 })
1108 }
1109 }
1110 } else if cached_reader && entered_with_transaction {
1111 match transaction_control {
1112 None | Some(CachedReadTransactionControl::Finish(_)) => {
1113 read(scope, &owned_handle.conn, true, true)
1114 }
1115 Some(CachedReadTransactionControl::BeginDeferred) => {
1116 Err(StorageError::InvalidInput {
1117 capability: StorageCapability::Sql,
1118 operation: operation.into(),
1119 message: "cached read-only handle already owns an admitted read \
1120 transaction; nested BEGIN is not supported"
1121 .into(),
1122 })
1123 }
1124 Some(CachedReadTransactionControl::Unsupported(keyword)) => {
1125 Err(StorageError::InvalidInput {
1126 capability: StorageCapability::Sql,
1127 operation: operation.into(),
1128 message: format!(
1129 "cached read-only transaction does not support nested or \
1130 write-locking transaction control ({keyword})"
1131 ),
1132 })
1133 }
1134 }
1135 } else if cached_reader {
1136 match transaction_control {
1137 None | Some(CachedReadTransactionControl::BeginDeferred) => {
1138 read(scope, &owned_handle.conn, false, true)
1139 }
1140 Some(CachedReadTransactionControl::Finish(keyword))
1141 | Some(CachedReadTransactionControl::Unsupported(keyword)) => {
1142 Err(StorageError::InvalidInput {
1143 capability: StorageCapability::Sql,
1144 operation: operation.into(),
1145 message: format!(
1146 "cached read-only handle has no admitted transaction for \
1147 transaction control ({keyword})"
1148 ),
1149 })
1150 }
1151 }
1152 } else {
1153 read(scope, &owned_handle.conn, false, entered_autocommit)
1154 };
1155
1156 if scope.cleanup_failed() {
1157 restore_handle = false;
1162 }
1163
1164 if cached_reader
1170 && matches!(result, Err(StorageError::Timeout { .. }))
1171 && !owned_handle.conn.is_autocommit()
1172 {
1173 match owned_handle.conn.execute_batch("ROLLBACK") {
1174 Ok(()) if owned_handle.conn.is_autocommit() => {
1175 drop(owned_handle.read_transaction_slot.take());
1176 }
1177 Ok(()) => {
1178 restore_handle = false;
1179 result = Err(StorageError::Transaction {
1180 operation: operation.into(),
1181 message:
1182 "interrupted read transaction rollback did not restore autocommit; \
1183 the connection was discarded"
1184 .into(),
1185 });
1186 }
1187 Err(error) => {
1188 restore_handle = false;
1189 result = Err(StorageError::Transaction {
1190 operation: operation.into(),
1191 message: format!(
1192 "failed to roll back interrupted read transaction ({error}); \
1193 the connection was discarded"
1194 ),
1195 });
1196 }
1197 }
1198 }
1199
1200 if cached_reader && entered_with_transaction {
1201 if owned_handle.conn.is_autocommit() {
1202 drop(owned_handle.read_transaction_slot.take());
1206 if result.is_ok()
1207 && !matches!(
1208 transaction_control,
1209 Some(CachedReadTransactionControl::Finish(_))
1210 )
1211 {
1212 result = Err(StorageError::InvalidInput {
1213 capability: StorageCapability::Sql,
1214 operation: operation.into(),
1215 message: "cached read-only operation unexpectedly ended its admitted \
1216 transaction; reader admission was released after autocommit"
1217 .into(),
1218 });
1219 }
1220 } else if result.is_ok()
1221 && matches!(
1222 transaction_control,
1223 Some(CachedReadTransactionControl::Finish(_))
1224 )
1225 {
1226 result = Err(StorageError::InvalidInput {
1227 capability: StorageCapability::Sql,
1228 operation: operation.into(),
1229 message: "transaction-ending control completed but the cached reader \
1230 remained outside autocommit; its reader permit remains retained"
1231 .into(),
1232 });
1233 }
1234 } else if cached_reader
1235 && entered_autocommit
1236 && matches!(
1237 transaction_control,
1238 Some(CachedReadTransactionControl::BeginDeferred)
1239 )
1240 && result.is_ok()
1241 {
1242 if owned_handle.conn.is_autocommit() {
1243 result = Err(StorageError::InvalidInput {
1244 capability: StorageCapability::Sql,
1245 operation: operation.into(),
1246 message: "deferred BEGIN completed without opening a read transaction"
1247 .into(),
1248 });
1249 } else {
1250 match operation_slot.take() {
1251 Some(slot) => {
1252 let tx_handle = khive_storage::tx_registry::register_scoped(
1253 Some(CACHED_READ_TRANSACTION_LABEL.to_string()),
1254 origin.clone(),
1255 );
1256 owned_handle.read_transaction_slot = Some(CachedReadTransaction {
1257 _slot: slot,
1258 _tx_handle: tx_handle,
1259 opened_at: Instant::now(),
1260 });
1261 }
1262 None => {
1263 result = Err(StorageError::Pool {
1264 operation: operation.into(),
1265 message: "successful cached-reader BEGIN had no operation permit; \
1266 its transaction was rolled back before returning"
1267 .into(),
1268 });
1269 }
1270 }
1271 }
1272 }
1273
1274 if cached_reader
1277 && owned_handle.read_transaction_slot.is_none()
1278 && !owned_handle.conn.is_autocommit()
1279 {
1280 match owned_handle.conn.execute_batch("ROLLBACK") {
1281 Ok(()) if owned_handle.conn.is_autocommit() => {
1282 if result.is_ok() {
1283 result = Err(StorageError::InvalidInput {
1284 capability: StorageCapability::Sql,
1285 operation: operation.into(),
1286 message: "cached read-only operation left the connection outside \
1287 autocommit; its transaction was rolled back before \
1288 releasing the reader permit"
1289 .into(),
1290 });
1291 }
1292 }
1293 Ok(()) => {
1294 restore_handle = false;
1295 result = Err(StorageError::Transaction {
1296 operation: operation.into(),
1297 message: "ROLLBACK completed but the cached reader remained outside \
1298 autocommit; the connection was discarded before releasing \
1299 the reader permit"
1300 .into(),
1301 });
1302 }
1303 Err(error) => {
1304 restore_handle = false;
1305 result = Err(StorageError::Transaction {
1306 operation: operation.into(),
1307 message: format!(
1308 "failed to roll back a cached reader outside autocommit ({error}); \
1309 the connection was discarded before releasing the reader permit"
1310 ),
1311 });
1312 }
1313 }
1314 }
1315
1316 let owned_handle = if restore_handle {
1317 Some(owned_handle)
1318 } else {
1319 drop(owned_handle);
1322 None
1323 };
1324 drop(operation_slot);
1328 Ok((owned_handle, result))
1329 },
1330 )
1331 .await?;
1332 *handle = owned_handle;
1333 result
1334}
1335
1336#[async_trait]
1337impl khive_storage::SqlReader for SqliteReader {
1338 async fn query_row(
1339 &mut self,
1340 statement: SqlStatement,
1341 ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
1342 let transaction_control = cached_read_transaction_control(&statement.sql);
1343 execute_standalone_read(
1344 &mut self.handle,
1345 Arc::clone(&self.pool),
1346 "query_row",
1347 transaction_control,
1348 move |scope, conn, rollback, interruptible| {
1349 execute_query_row_interruptibly(
1350 scope,
1351 conn,
1352 &statement,
1353 "query_row",
1354 rollback,
1355 interruptible,
1356 )
1357 },
1358 )
1359 .await
1360 }
1361
1362 async fn query_all(
1363 &mut self,
1364 statement: SqlStatement,
1365 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1366 let transaction_control = cached_read_transaction_control(&statement.sql);
1367 execute_standalone_read(
1368 &mut self.handle,
1369 Arc::clone(&self.pool),
1370 "query_all",
1371 transaction_control,
1372 move |scope, conn, rollback, interruptible| {
1373 execute_query_interruptibly(
1374 scope,
1375 conn,
1376 &statement,
1377 "query_all",
1378 rollback,
1379 interruptible,
1380 )
1381 },
1382 )
1383 .await
1384 }
1385
1386 async fn query_page(
1387 &mut self,
1388 statement: SqlStatement,
1389 page: PageRequest,
1390 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1391 let transaction_control = cached_read_transaction_control(&statement.sql);
1392 execute_standalone_read(
1393 &mut self.handle,
1394 Arc::clone(&self.pool),
1395 "query_page",
1396 transaction_control,
1397 move |scope, conn, rollback, interruptible| {
1398 execute_query_page_interruptibly(
1399 scope,
1400 conn,
1401 &statement,
1402 &page,
1403 "query_page",
1404 rollback,
1405 interruptible,
1406 )
1407 },
1408 )
1409 .await
1410 }
1411
1412 async fn query_scalar(
1413 &mut self,
1414 statement: SqlStatement,
1415 ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
1416 let row = self.query_row(statement).await?;
1417 Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
1418 }
1419
1420 async fn explain(
1421 &mut self,
1422 statement: SqlStatement,
1423 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1424 let explain_stmt = SqlStatement {
1425 sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
1426 params: statement.params,
1427 label: statement.label,
1428 };
1429 self.query_all(explain_stmt).await
1430 }
1431}
1432
1433struct SqliteWriter {
1438 handle: Option<StandaloneHandle>,
1452 writer_task: Option<crate::writer_task::WriterTaskHandle>,
1459 origin: khive_storage::tx_registry::TxOrigin,
1462 db: String,
1466 pool: Arc<ConnectionPool>,
1470}
1471
1472impl SqliteWriter {
1473 async fn ensure_conn(
1498 pool: Arc<ConnectionPool>,
1499 ) -> khive_storage::types::StorageResult<StandaloneHandle> {
1500 open_cached_reader_handle(pool).await
1501 }
1502}
1503
1504#[async_trait]
1505impl khive_storage::SqlReader for SqliteWriter {
1506 async fn query_row(
1507 &mut self,
1508 statement: SqlStatement,
1509 ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
1510 if self.handle.is_none() && self.writer_task.is_some() {
1511 self.handle = Some(Self::ensure_conn(Arc::clone(&self.pool)).await?);
1512 }
1513 let transaction_control = cached_read_transaction_control(&statement.sql);
1514 execute_standalone_read(
1515 &mut self.handle,
1516 Arc::clone(&self.pool),
1517 "writer.query_row",
1518 transaction_control,
1519 move |scope, conn, rollback, interruptible| {
1520 execute_query_row_interruptibly(
1521 scope,
1522 conn,
1523 &statement,
1524 "writer.query_row",
1525 rollback,
1526 interruptible,
1527 )
1528 },
1529 )
1530 .await
1531 }
1532
1533 async fn query_all(
1534 &mut self,
1535 statement: SqlStatement,
1536 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1537 if self.handle.is_none() && self.writer_task.is_some() {
1538 self.handle = Some(Self::ensure_conn(Arc::clone(&self.pool)).await?);
1539 }
1540 let transaction_control = cached_read_transaction_control(&statement.sql);
1541 execute_standalone_read(
1542 &mut self.handle,
1543 Arc::clone(&self.pool),
1544 "writer.query_all",
1545 transaction_control,
1546 move |scope, conn, rollback, interruptible| {
1547 execute_query_interruptibly(
1548 scope,
1549 conn,
1550 &statement,
1551 "writer.query_all",
1552 rollback,
1553 interruptible,
1554 )
1555 },
1556 )
1557 .await
1558 }
1559
1560 async fn query_page(
1561 &mut self,
1562 statement: SqlStatement,
1563 page: PageRequest,
1564 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1565 if self.handle.is_none() && self.writer_task.is_some() {
1566 self.handle = Some(Self::ensure_conn(Arc::clone(&self.pool)).await?);
1567 }
1568 let transaction_control = cached_read_transaction_control(&statement.sql);
1569 execute_standalone_read(
1570 &mut self.handle,
1571 Arc::clone(&self.pool),
1572 "writer.query_page",
1573 transaction_control,
1574 move |scope, conn, rollback, interruptible| {
1575 execute_query_page_interruptibly(
1576 scope,
1577 conn,
1578 &statement,
1579 &page,
1580 "writer.query_page",
1581 rollback,
1582 interruptible,
1583 )
1584 },
1585 )
1586 .await
1587 }
1588
1589 async fn query_scalar(
1590 &mut self,
1591 statement: SqlStatement,
1592 ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
1593 let row = khive_storage::SqlReader::query_row(self, statement).await?;
1594 Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
1595 }
1596
1597 async fn explain(
1598 &mut self,
1599 statement: SqlStatement,
1600 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1601 let explain_stmt = SqlStatement {
1602 sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
1603 params: statement.params,
1604 label: statement.label,
1605 };
1606 khive_storage::SqlReader::query_all(self, explain_stmt).await
1607 }
1608}
1609
1610#[async_trait]
1611impl khive_storage::SqlWriter for SqliteWriter {
1612 async fn execute(
1613 &mut self,
1614 statement: SqlStatement,
1615 ) -> khive_storage::types::StorageResult<u64> {
1616 if let Some(writer_task) = self.writer_task.clone() {
1624 return writer_task
1625 .send_bounded(move |conn| {
1626 let mut stmt = prepare_cached_sql_statement(conn, &statement.sql)
1627 .map_err(|e| map_rusqlite_err(e, "execute"))?;
1628 bind_params(&mut stmt, &statement.params)
1629 .map_err(|e| map_rusqlite_err(e, "execute"))?;
1630 let affected = stmt
1631 .raw_execute()
1632 .map_err(|e| map_rusqlite_err(e, "execute"))?;
1633 Ok(affected as u64)
1634 })
1635 .await;
1636 }
1637
1638 let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
1639 operation: "execute".into(),
1640 message: "connection already consumed".into(),
1641 })?;
1642 let (handle, result) = tokio::task::spawn_blocking(move || {
1643 let res = (|| -> Result<usize, rusqlite::Error> {
1644 let mut stmt = prepare_cached_sql_statement(&handle.conn, &statement.sql)?;
1645 bind_params(&mut stmt, &statement.params)?;
1646 stmt.raw_execute()
1647 })();
1648 (handle, res)
1649 })
1650 .await
1651 .map_err(|e| StorageError::driver(StorageCapability::Sql, "execute", e))?;
1652 self.handle = Some(handle);
1653 let affected = result.map_err(|e| {
1654 crate::timeout_sink::maybe_emit_busy(
1655 &self.db,
1656 crate::timeout_sink::Site::StandaloneSqlBridge,
1657 &e,
1658 );
1659 map_rusqlite_err(e, "execute")
1660 })?;
1661 Ok(affected as u64)
1662 }
1663
1664 async fn execute_batch(
1665 &mut self,
1666 statements: Vec<SqlStatement>,
1667 ) -> khive_storage::types::StorageResult<u64> {
1668 reject_transaction_control_statements(&statements, "execute_batch")?;
1684 if let Some(writer_task) = self.writer_task.clone() {
1685 return writer_task
1686 .send_bounded(move |conn| {
1687 let prepared = prepare_batch_statements(conn, &statements)
1688 .map_err(|e| map_rusqlite_err(e, "execute_batch"))?;
1689 execute_prepared_batch(conn, prepared, &statements)
1690 .map_err(|e| map_rusqlite_err(e, "execute_batch"))
1691 })
1692 .await;
1693 }
1694
1695 let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
1696 operation: "execute_batch".into(),
1697 message: "connection already consumed".into(),
1698 })?;
1699 let origin = self.origin.clone();
1700 let (handle, result) = tokio::task::spawn_blocking(move || {
1701 let (disposition, result) = execute_standalone_batch(&handle.conn, &statements, origin);
1702 let retained = match disposition {
1703 BatchHandleDisposition::Retain => Some(handle),
1704 BatchHandleDisposition::Poison => None,
1705 };
1706 (retained, result)
1707 })
1708 .await
1709 .map_err(|e| StorageError::driver(StorageCapability::Sql, "execute_batch", e))?;
1710 self.handle = handle;
1711 result.map_err(|failure| {
1712 crate::timeout_sink::maybe_emit_busy(
1713 &self.db,
1714 crate::timeout_sink::Site::StandaloneSqlBridge,
1715 &failure.error,
1716 );
1717 match failure.poison_reason {
1718 Some(poison_reason) => StorageError::driver(
1719 StorageCapability::Sql,
1720 "execute_batch",
1721 PoisonedBatchError {
1722 original: failure.error,
1723 poison_reason,
1724 },
1725 ),
1726 None => map_rusqlite_err(failure.error, "execute_batch"),
1727 }
1728 })
1729 }
1730
1731 async fn execute_script(&mut self, script: String) -> khive_storage::types::StorageResult<()> {
1732 if let Some(writer_task) = self.writer_task.clone() {
1745 return writer_task
1746 .send_bounded(move |conn| {
1747 conn.execute_batch(&script)
1748 .map_err(|e| map_rusqlite_err(e, "execute_script"))
1749 })
1750 .await;
1751 }
1752
1753 let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
1754 operation: "execute_script".into(),
1755 message: "connection already consumed".into(),
1756 })?;
1757 let (handle, result) = tokio::task::spawn_blocking(move || {
1758 let res = handle.conn.execute_batch(&script);
1759 (handle, res)
1760 })
1761 .await
1762 .map_err(|e| StorageError::driver(StorageCapability::Sql, "execute_script", e))?;
1763 self.handle = Some(handle);
1764 result.map_err(|e| {
1765 crate::timeout_sink::maybe_emit_busy(
1766 &self.db,
1767 crate::timeout_sink::Site::StandaloneSqlBridge,
1768 &e,
1769 );
1770 map_rusqlite_err(e, "execute_script")
1771 })
1772 }
1773
1774 async fn execute_script_top_level(
1775 &mut self,
1776 script: String,
1777 ) -> khive_storage::types::StorageResult<()> {
1778 if let Some(writer_task) = self.writer_task.clone() {
1788 return writer_task
1789 .send_top_level_bounded(move |conn| {
1790 conn.execute_batch(&script)
1791 .map_err(|e| map_rusqlite_err(e, "execute_script_top_level"))
1792 })
1793 .await;
1794 }
1795
1796 let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
1800 operation: "execute_script_top_level".into(),
1801 message: "connection already consumed".into(),
1802 })?;
1803 let (handle, result) = tokio::task::spawn_blocking(move || {
1804 let res = handle.conn.execute_batch(&script);
1805 (handle, res)
1806 })
1807 .await
1808 .map_err(|e| StorageError::driver(StorageCapability::Sql, "execute_script_top_level", e))?;
1809 self.handle = Some(handle);
1810 result.map_err(|e| {
1811 crate::timeout_sink::maybe_emit_busy(
1812 &self.db,
1813 crate::timeout_sink::Site::StandaloneSqlBridge,
1814 &e,
1815 );
1816 map_rusqlite_err(e, "execute_script_top_level")
1817 })
1818 }
1819}
1820
1821async fn run_pool_reader_query<T, F>(
1826 pool: Arc<ConnectionPool>,
1827 operation: &'static str,
1828 query: F,
1829) -> khive_storage::types::StorageResult<T>
1830where
1831 T: Send + 'static,
1832 F: FnOnce(
1833 &crate::read_cancellation::InterruptibleReadScope,
1834 &rusqlite::Connection,
1835 ) -> khive_storage::types::StorageResult<T>
1836 + Send
1837 + 'static,
1838{
1839 crate::read_cancellation::run_interruptible_read(
1840 StorageCapability::Sql,
1841 operation,
1842 move |scope| {
1843 let mut guard = pool.resolve_reader_checkout(
1847 StorageCapability::Sql,
1848 operation,
1849 pool.reader_until(|| scope.should_stop()),
1850 )?;
1851 scope.with_pooled_reader(&mut guard, |conn| query(scope, conn))
1852 },
1853 )
1854 .await
1855}
1856
1857async fn run_pool_writer_query<T, F>(
1858 pool: Arc<ConnectionPool>,
1859 operation: &'static str,
1860 query: F,
1861) -> khive_storage::types::StorageResult<T>
1862where
1863 T: Send + 'static,
1864 F: FnOnce(
1865 &crate::read_cancellation::InterruptibleReadScope,
1866 &rusqlite::Connection,
1867 bool,
1868 ) -> khive_storage::types::StorageResult<T>
1869 + Send
1870 + 'static,
1871{
1872 crate::read_cancellation::run_interruptible_read(
1873 StorageCapability::Sql,
1874 operation,
1875 move |scope| {
1876 let guard = pool.try_writer().map_err(|error: SqliteError| {
1877 StorageError::driver(StorageCapability::Sql, operation, error)
1878 })?;
1879 scope.with_pooled_writer(&pool, &guard, |conn| {
1880 let interruptible = conn.is_autocommit();
1881 query(scope, conn, interruptible)
1882 })
1883 },
1884 )
1885 .await
1886}
1887
1888struct PoolBackedReader {
1889 pool: Arc<ConnectionPool>,
1890}
1891
1892#[async_trait]
1893impl khive_storage::SqlReader for PoolBackedReader {
1894 async fn query_row(
1895 &mut self,
1896 statement: SqlStatement,
1897 ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
1898 let pool = Arc::clone(&self.pool);
1899 run_pool_reader_query(pool, "pool_reader.query_row", move |scope, conn| {
1900 execute_query_row_interruptibly(
1901 scope,
1902 conn,
1903 &statement,
1904 "pool_reader.query_row",
1905 false,
1906 true,
1907 )
1908 })
1909 .await
1910 }
1911
1912 async fn query_all(
1913 &mut self,
1914 statement: SqlStatement,
1915 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1916 let pool = Arc::clone(&self.pool);
1917 run_pool_reader_query(pool, "pool_reader.query_all", move |scope, conn| {
1918 execute_query_interruptibly(
1919 scope,
1920 conn,
1921 &statement,
1922 "pool_reader.query_all",
1923 false,
1924 true,
1925 )
1926 })
1927 .await
1928 }
1929
1930 async fn query_page(
1931 &mut self,
1932 statement: SqlStatement,
1933 page: PageRequest,
1934 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1935 let pool = Arc::clone(&self.pool);
1936 run_pool_reader_query(pool, "pool_reader.query_page", move |scope, conn| {
1937 execute_query_page_interruptibly(
1938 scope,
1939 conn,
1940 &statement,
1941 &page,
1942 "pool_reader.query_page",
1943 false,
1944 true,
1945 )
1946 })
1947 .await
1948 }
1949
1950 async fn query_scalar(
1951 &mut self,
1952 statement: SqlStatement,
1953 ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
1954 let row = self.query_row(statement).await?;
1955 Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
1956 }
1957
1958 async fn explain(
1959 &mut self,
1960 statement: SqlStatement,
1961 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1962 let explain_stmt = SqlStatement {
1963 sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
1964 params: statement.params,
1965 label: statement.label,
1966 };
1967 self.query_all(explain_stmt).await
1968 }
1969}
1970
1971struct PoolBackedWriter {
1972 pool: Arc<ConnectionPool>,
1973}
1974
1975#[async_trait]
1976impl khive_storage::SqlReader for PoolBackedWriter {
1977 async fn query_row(
1978 &mut self,
1979 statement: SqlStatement,
1980 ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
1981 let pool = Arc::clone(&self.pool);
1982 run_pool_writer_query(
1983 pool,
1984 "pool_writer.query_row",
1985 move |scope, conn, interruptible| {
1986 execute_query_row_interruptibly(
1987 scope,
1988 conn,
1989 &statement,
1990 "pool_writer.query_row",
1991 false,
1992 interruptible,
1993 )
1994 },
1995 )
1996 .await
1997 }
1998
1999 async fn query_all(
2000 &mut self,
2001 statement: SqlStatement,
2002 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2003 let pool = Arc::clone(&self.pool);
2004 run_pool_writer_query(
2005 pool,
2006 "pool_writer.query_all",
2007 move |scope, conn, interruptible| {
2008 execute_query_interruptibly(
2009 scope,
2010 conn,
2011 &statement,
2012 "pool_writer.query_all",
2013 false,
2014 interruptible,
2015 )
2016 },
2017 )
2018 .await
2019 }
2020
2021 async fn query_page(
2022 &mut self,
2023 statement: SqlStatement,
2024 page: PageRequest,
2025 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2026 let pool = Arc::clone(&self.pool);
2027 run_pool_writer_query(
2028 pool,
2029 "pool_writer.query_page",
2030 move |scope, conn, interruptible| {
2031 execute_query_page_interruptibly(
2032 scope,
2033 conn,
2034 &statement,
2035 &page,
2036 "pool_writer.query_page",
2037 false,
2038 interruptible,
2039 )
2040 },
2041 )
2042 .await
2043 }
2044
2045 async fn query_scalar(
2046 &mut self,
2047 statement: SqlStatement,
2048 ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
2049 let row = khive_storage::SqlReader::query_row(self, statement).await?;
2050 Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
2051 }
2052
2053 async fn explain(
2054 &mut self,
2055 statement: SqlStatement,
2056 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2057 let explain_stmt = SqlStatement {
2058 sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
2059 params: statement.params,
2060 label: statement.label,
2061 };
2062 khive_storage::SqlReader::query_all(self, explain_stmt).await
2063 }
2064}
2065
2066#[async_trait]
2067impl khive_storage::SqlWriter for PoolBackedWriter {
2068 async fn execute(
2069 &mut self,
2070 statement: SqlStatement,
2071 ) -> khive_storage::types::StorageResult<u64> {
2072 let pool = Arc::clone(&self.pool);
2077 tokio::task::spawn_blocking(move || {
2078 let guard = pool.try_writer().map_err(|e: SqliteError| {
2079 StorageError::driver(StorageCapability::Sql, "pool_writer.execute", e)
2080 })?;
2081 let mut stmt = prepare_cached_sql_statement(&guard, &statement.sql)
2082 .map_err(|e| map_rusqlite_err(e, "pool_writer.execute"))?;
2083 bind_params(&mut stmt, &statement.params)
2084 .map_err(|e| map_rusqlite_err(e, "pool_writer.execute"))?;
2085 let rows = stmt
2086 .raw_execute()
2087 .map_err(|e| map_rusqlite_err(e, "pool_writer.execute"))?;
2088 Ok(rows as u64)
2089 })
2090 .await
2091 .map_err(|e| StorageError::driver(StorageCapability::Sql, "pool_writer.execute", e))?
2092 }
2093
2094 async fn execute_batch(
2095 &mut self,
2096 statements: Vec<SqlStatement>,
2097 ) -> khive_storage::types::StorageResult<u64> {
2098 reject_transaction_control_statements(&statements, "pool_writer.execute_batch")?;
2102 let pool = Arc::clone(&self.pool);
2103 tokio::task::spawn_blocking(move || {
2104 let guard = pool.try_writer().map_err(|e: SqliteError| {
2105 StorageError::driver(StorageCapability::Sql, "pool_writer.execute_batch", e)
2106 })?;
2107 let prepared = prepare_batch_statements(&guard, &statements)
2108 .map_err(|e| map_rusqlite_err(e, "pool_writer.execute_batch"))?;
2109 guard
2110 .execute_batch("BEGIN IMMEDIATE")
2111 .map_err(|e| map_rusqlite_err(e, "pool_writer.execute_batch"))?;
2112 let _tx_handle = khive_storage::tx_registry::register_scoped(
2113 Some("pool_writer.execute_batch".to_string()),
2114 pool.origin(),
2115 );
2116 let result = execute_prepared_batch(&guard, prepared, &statements)
2117 .map_err(|e| map_rusqlite_err(e, "pool_writer.execute_batch"));
2118 match result {
2119 Ok(total) => {
2120 if let Err(e) = guard.execute_batch("COMMIT") {
2121 let _ = guard.execute_batch("ROLLBACK");
2122 Err(map_rusqlite_err(e, "pool_writer.execute_batch"))
2123 } else {
2124 Ok(total)
2125 }
2126 }
2127 Err(e) => {
2128 let _ = guard.execute_batch("ROLLBACK");
2129 Err(e)
2130 }
2131 }
2132 })
2133 .await
2134 .map_err(|e| StorageError::driver(StorageCapability::Sql, "pool_writer.execute_batch", e))?
2135 }
2136
2137 async fn execute_script(&mut self, script: String) -> khive_storage::types::StorageResult<()> {
2138 let pool = Arc::clone(&self.pool);
2141 tokio::task::spawn_blocking(move || {
2142 let guard = pool.try_writer().map_err(|e: SqliteError| {
2143 StorageError::driver(StorageCapability::Sql, "pool_writer.execute_script", e)
2144 })?;
2145 guard
2146 .execute_batch(&script)
2147 .map_err(|e| map_rusqlite_err(e, "pool_writer.execute_script"))
2148 })
2149 .await
2150 .map_err(|e| {
2151 StorageError::driver(StorageCapability::Sql, "pool_writer.execute_script", e)
2152 })?
2153 }
2154}
2155
2156struct InlineWriter {
2182 conn: *const rusqlite::Connection,
2183}
2184
2185unsafe impl Send for InlineWriter {}
2192
2193impl InlineWriter {
2194 fn conn(&self) -> &rusqlite::Connection {
2198 unsafe { &*self.conn }
2199 }
2200}
2201
2202#[async_trait]
2203impl khive_storage::SqlReader for InlineWriter {
2204 async fn query_row(
2205 &mut self,
2206 statement: SqlStatement,
2207 ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
2208 execute_query_row(self.conn(), &statement)
2209 .map_err(|e| map_rusqlite_err(e, "inline.query_row"))
2210 }
2211
2212 async fn query_all(
2213 &mut self,
2214 statement: SqlStatement,
2215 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2216 execute_query(self.conn(), &statement).map_err(|e| map_rusqlite_err(e, "inline.query_all"))
2217 }
2218
2219 async fn query_page(
2220 &mut self,
2221 statement: SqlStatement,
2222 page: PageRequest,
2223 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2224 execute_query_page(self.conn(), &statement, &page)
2225 .map_err(|e| map_rusqlite_err(e, "inline.query_page"))
2226 }
2227
2228 async fn query_scalar(
2229 &mut self,
2230 statement: SqlStatement,
2231 ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
2232 let row = khive_storage::SqlReader::query_row(self, statement).await?;
2233 Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
2234 }
2235
2236 async fn explain(
2237 &mut self,
2238 statement: SqlStatement,
2239 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2240 let explain_stmt = SqlStatement {
2241 sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
2242 params: statement.params,
2243 label: statement.label,
2244 };
2245 khive_storage::SqlReader::query_all(self, explain_stmt).await
2246 }
2247}
2248
2249#[async_trait]
2250impl khive_storage::SqlWriter for InlineWriter {
2251 async fn execute(
2252 &mut self,
2253 statement: SqlStatement,
2254 ) -> khive_storage::types::StorageResult<u64> {
2255 let mut stmt = prepare_cached_sql_statement(self.conn(), &statement.sql)
2258 .map_err(|e| map_rusqlite_err(e, "inline.execute"))?;
2259 bind_params(&mut stmt, &statement.params)
2260 .map_err(|e| map_rusqlite_err(e, "inline.execute"))?;
2261 let affected = stmt
2262 .raw_execute()
2263 .map_err(|e| map_rusqlite_err(e, "inline.execute"))?;
2264 Ok(affected as u64)
2265 }
2266
2267 async fn execute_batch(
2268 &mut self,
2269 statements: Vec<SqlStatement>,
2270 ) -> khive_storage::types::StorageResult<u64> {
2271 reject_transaction_control_statements(&statements, "inline.execute_batch")?;
2276 let prepared = prepare_batch_statements(self.conn(), &statements)
2277 .map_err(|e| map_rusqlite_err(e, "inline.execute_batch"))?;
2278 execute_prepared_batch(self.conn(), prepared, &statements)
2279 .map_err(|e| map_rusqlite_err(e, "inline.execute_batch"))
2280 }
2281
2282 async fn execute_script(&mut self, script: String) -> khive_storage::types::StorageResult<()> {
2283 self.conn()
2286 .execute_batch(&script)
2287 .map_err(|e| map_rusqlite_err(e, "inline.execute_script"))
2288 }
2289}
2290
2291fn block_on_sync<F: std::future::Future>(fut: F) -> Result<F::Output, StorageError> {
2314 use std::task::{Context, Poll, RawWaker, RawWakerVTable, Waker};
2315
2316 fn no_op(_: *const ()) {}
2317 fn clone_waker(_: *const ()) -> RawWaker {
2318 RawWaker::new(std::ptr::null(), &VTABLE)
2319 }
2320 static VTABLE: RawWakerVTable = RawWakerVTable::new(clone_waker, no_op, no_op, no_op);
2321
2322 let raw_waker = RawWaker::new(std::ptr::null(), &VTABLE);
2325 let waker = unsafe { Waker::from_raw(raw_waker) };
2326 let mut cx = Context::from_waker(&waker);
2327
2328 let mut fut = std::pin::pin!(fut);
2329 match fut.as_mut().poll(&mut cx) {
2330 Poll::Ready(v) => Ok(v),
2331 Poll::Pending => {
2332 tracing::error!(
2333 "block_on_sync: atomic_unit future suspended on its first poll — \
2334 the closure passed to SqlAccess::atomic_unit must be non-blocking \
2335 (synchronous InlineWriter calls only, no real .await point)"
2336 );
2337 Err(StorageError::Internal(
2338 "atomic_unit future suspended — closure must be non-blocking".to_string(),
2339 ))
2340 }
2341 }
2342}
2343
2344async fn run_manual_atomic_unit(
2349 writer: &mut dyn khive_storage::SqlWriter,
2350 op: AtomicUnitOp,
2351 origin: khive_storage::tx_registry::TxOrigin,
2352) -> khive_storage::types::StorageResult<Box<dyn Any + Send>> {
2353 fn tx_stmt(sql: &str, label: &str) -> SqlStatement {
2354 SqlStatement {
2355 sql: sql.to_string(),
2356 params: vec![],
2357 label: Some(label.to_string()),
2358 }
2359 }
2360 khive_storage::SqlWriter::execute(writer, tx_stmt("BEGIN IMMEDIATE", "begin")).await?;
2361 let _tx_handle =
2362 khive_storage::tx_registry::register_scoped(Some("atomic_unit".to_string()), origin);
2363
2364 let result = op(writer).await;
2365
2366 match result {
2367 Ok(value) => {
2368 match khive_storage::SqlWriter::execute(writer, tx_stmt("COMMIT", "commit")).await {
2369 Ok(_) => Ok(value),
2370 Err(e) => {
2371 let _ =
2372 khive_storage::SqlWriter::execute(writer, tx_stmt("ROLLBACK", "rollback"))
2373 .await;
2374 Err(e)
2375 }
2376 }
2377 }
2378 Err(e) => {
2379 let _ =
2380 khive_storage::SqlWriter::execute(writer, tx_stmt("ROLLBACK", "rollback")).await;
2381 Err(e)
2382 }
2383 }
2384}
2385
2386pub struct SqlBridge {
2399 pool: Arc<ConnectionPool>,
2400 is_file_backed: bool,
2401}
2402
2403impl SqlBridge {
2404 pub fn new(pool: Arc<ConnectionPool>, is_file_backed: bool) -> Self {
2406 Self {
2407 pool,
2408 is_file_backed,
2409 }
2410 }
2411}
2412
2413#[async_trait]
2414impl khive_storage::SqlAccess for SqlBridge {
2415 fn database_path(&self) -> Option<std::path::PathBuf> {
2416 self.pool.canonical_path().map(std::path::Path::to_path_buf)
2417 }
2418
2419 async fn reader(
2420 &self,
2421 ) -> khive_storage::types::StorageResult<Box<dyn khive_storage::SqlReader>> {
2422 if self.is_file_backed {
2423 Ok(Box::new(SqliteReader {
2424 handle: Some(open_cached_reader_handle(Arc::clone(&self.pool)).await?),
2425 pool: Arc::clone(&self.pool),
2426 }))
2427 } else {
2428 Ok(Box::new(PoolBackedReader {
2429 pool: Arc::clone(&self.pool),
2430 }))
2431 }
2432 }
2433
2434 async fn writer(
2435 &self,
2436 ) -> khive_storage::types::StorageResult<Box<dyn khive_storage::SqlWriter>> {
2437 if self.is_file_backed {
2438 if self.pool.config().read_only {
2439 return Err(StorageError::Pool {
2440 operation: "writer".into(),
2441 message: "backend is read-only".into(),
2442 });
2443 }
2444 let db = crate::timeout_sink::db_label(&self.pool);
2445 let writer_task = match self.pool.writer_task_handle() {
2451 Ok(handle) => handle,
2452 Err(e) => {
2453 if self.pool.config().write_routing_strict {
2454 return Err(e);
2455 }
2456 tracing::warn!(
2457 error = %e,
2458 "KHIVE_WRITE_ROUTING is not strict; writer() degrades to the \
2459 standalone-connection path"
2460 );
2461 None
2462 }
2463 };
2464 if writer_task.is_none() && self.pool.config().write_routing_strict {
2465 return Err(StorageError::Pool {
2466 operation: "writer".into(),
2467 message: "KHIVE_WRITE_ROUTING=strict but no writer-task handle is \
2468 available; refusing to fall back to a direct connection"
2469 .into(),
2470 });
2471 }
2472 if writer_task.is_none() && self.pool.write_queue_active() {
2473 crate::timeout_sink::emit_direct_route_violation(
2480 &db,
2481 crate::timeout_sink::Site::DirectRouteSqlBridgeWriter,
2482 );
2483 }
2484 let handle = if writer_task.is_none() {
2497 let handle_slot = acquire_handle_slot(
2498 self.pool.sql_bridge_writer_slots(),
2499 self.pool.config().checkout_timeout,
2500 "sql_bridge.writer_handle",
2501 SlotTimeoutClass::Admission,
2502 )
2503 .await?;
2504 let (conn, handle_slot) =
2505 open_standalone_writer_on_blocking(Arc::clone(&self.pool), handle_slot).await?;
2506 Some(StandaloneHandle {
2507 conn,
2508 _retained_slot: Some(handle_slot),
2509 read_transaction_slot: None,
2510 })
2511 } else {
2512 None
2513 };
2514 Ok(Box::new(SqliteWriter {
2515 handle,
2516 writer_task,
2517 origin: self.pool.origin(),
2518 db,
2519 pool: Arc::clone(&self.pool),
2520 }))
2521 } else {
2522 Ok(Box::new(PoolBackedWriter {
2523 pool: Arc::clone(&self.pool),
2524 }))
2525 }
2526 }
2527
2528 async fn atomic_unit(
2539 &self,
2540 op: AtomicUnitOp,
2541 ) -> khive_storage::types::StorageResult<Box<dyn Any + Send>> {
2542 if self.is_file_backed {
2543 if self.pool.config().read_only {
2544 return Err(StorageError::Pool {
2545 operation: "atomic_unit".into(),
2546 message: "backend is read-only".into(),
2547 });
2548 }
2549 let handle = self.pool.writer_task_handle()?;
2555 if handle.is_none() && self.pool.config().write_routing_strict {
2556 return Err(StorageError::Pool {
2557 operation: "atomic_unit".into(),
2558 message: "KHIVE_WRITE_ROUTING=strict but no writer-task handle is \
2559 available; refusing to fall back to a direct connection"
2560 .into(),
2561 });
2562 }
2563 if handle.is_none() && self.pool.write_queue_active() {
2564 crate::timeout_sink::emit_direct_route_violation(
2565 &crate::timeout_sink::db_label(&self.pool),
2566 crate::timeout_sink::Site::DirectRouteAtomicUnit,
2567 );
2568 }
2569 if let Some(writer_task) = handle {
2570 return writer_task
2576 .send_bounded(move |conn| {
2577 let mut inline = InlineWriter {
2578 conn: conn as *const rusqlite::Connection,
2579 };
2580 match block_on_sync(op(&mut inline)) {
2590 Ok(inner) => inner,
2591 Err(e) => Err(e),
2592 }
2593 })
2594 .await;
2595 }
2596 let handle_slot = acquire_handle_slot(
2610 self.pool.sql_bridge_writer_slots(),
2611 self.pool.config().checkout_timeout,
2612 "sql_bridge.atomic_unit_handle",
2613 SlotTimeoutClass::Admission,
2614 )
2615 .await?;
2616 let (conn, handle_slot) =
2617 open_standalone_writer_on_blocking(Arc::clone(&self.pool), handle_slot).await?;
2618 let mut writer = SqliteWriter {
2619 handle: Some(StandaloneHandle {
2620 conn,
2621 _retained_slot: Some(handle_slot),
2622 read_transaction_slot: None,
2623 }),
2624 writer_task: None,
2625 origin: self.pool.origin(),
2626 db: crate::timeout_sink::db_label(&self.pool),
2627 pool: Arc::clone(&self.pool),
2628 };
2629 run_manual_atomic_unit(&mut writer, op, self.pool.origin()).await
2630 } else {
2631 let mut writer = PoolBackedWriter {
2635 pool: Arc::clone(&self.pool),
2636 };
2637 run_manual_atomic_unit(&mut writer, op, self.pool.origin()).await
2638 }
2639 }
2640}
2641
2642#[cfg(test)]
2643mod tests {
2644 use super::*;
2645 use crate::pool::PoolConfig;
2646 use khive_storage::types::{SqlStatement, SqlValue};
2647 use khive_storage::{SqlAccess as _, SqlReader as _};
2648
2649 fn database_tx_view(pool: &ConnectionPool) -> khive_storage::tx_registry::TxOriginFilter {
2650 match pool.origin() {
2651 khive_storage::tx_registry::TxOrigin::Database(identity) => {
2652 khive_storage::tx_registry::TxOriginFilter::Secondary(identity)
2653 }
2654 other => panic!("expected a file-backed database origin, got {other:?}"),
2655 }
2656 }
2657
2658 struct NotifyOnDrop(Arc<tokio::sync::Notify>);
2659
2660 impl Drop for NotifyOnDrop {
2661 fn drop(&mut self) {
2662 self.0.notify_one();
2663 }
2664 }
2665
2666 fn blocking_non_interrupting_progress_gate(
2667 conn: &rusqlite::Connection,
2668 ) -> (
2669 Arc<tokio::sync::Notify>,
2670 Arc<std::sync::Barrier>,
2671 Arc<tokio::sync::Notify>,
2672 ) {
2673 let entered = Arc::new(tokio::sync::Notify::new());
2674 let callback_entered = Arc::clone(&entered);
2675 let release = Arc::new(std::sync::Barrier::new(2));
2676 let callback_release = Arc::clone(&release);
2677 let completed = Arc::new(tokio::sync::Notify::new());
2678 let notify_on_drop = NotifyOnDrop(Arc::clone(&completed));
2679 let blocked_once = Arc::new(std::sync::atomic::AtomicBool::new(false));
2680 let callback_blocked_once = Arc::clone(&blocked_once);
2681 conn.progress_handler(
2682 1_000,
2683 Some(move || {
2684 let _keep_until_connection_drop = ¬ify_on_drop;
2685 if !callback_blocked_once.swap(true, std::sync::atomic::Ordering::SeqCst) {
2686 callback_entered.notify_one();
2687 callback_release.wait();
2688 return false;
2692 }
2693 false
2694 }),
2695 )
2696 .unwrap();
2697 (entered, release, completed)
2698 }
2699
2700 fn progress_gate_statement() -> SqlStatement {
2701 SqlStatement {
2702 sql: "WITH RECURSIVE rows(value) AS (\
2703 SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 999\
2704 ) SELECT SUM(value) FROM rows"
2705 .into(),
2706 params: vec![],
2707 label: None,
2708 }
2709 }
2710
2711 fn slow_insert_statement() -> SqlStatement {
2712 SqlStatement {
2713 sql: "INSERT INTO cancellation_write_probe(value) \
2714 WITH RECURSIVE rows(value) AS (\
2715 SELECT 1 UNION ALL SELECT value + 1 FROM rows WHERE value < 10000\
2716 ) SELECT value FROM rows"
2717 .into(),
2718 params: vec![],
2719 label: Some("non-interruptible-write-probe".into()),
2720 }
2721 }
2722
2723 fn passive_checkpoint(conn: &rusqlite::Connection) -> (i64, i64, i64) {
2724 conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
2725 Ok((row.get(0)?, row.get(1)?, row.get(2)?))
2726 })
2727 .unwrap()
2728 }
2729
2730 fn deliberately_slow_read_statement() -> SqlStatement {
2731 SqlStatement {
2732 sql: "WITH RECURSIVE numbers(value) AS (\
2733 SELECT 1 UNION ALL SELECT value + 1 FROM numbers WHERE value < 1000\
2734 ) SELECT SUM(a.value * b.value * c.value) \
2735 FROM numbers AS a CROSS JOIN numbers AS b CROSS JOIN numbers AS c"
2736 .into(),
2737 params: vec![],
2738 label: Some("read-cancellation-progress-probe".into()),
2739 }
2740 }
2741
2742 async fn wait_for_progress(probe: &std::sync::atomic::AtomicUsize) {
2743 tokio::time::timeout(std::time::Duration::from_secs(1), async {
2744 while probe.load(std::sync::atomic::Ordering::SeqCst) == 0 {
2745 tokio::task::yield_now().await;
2746 }
2747 })
2748 .await
2749 .expect("slow SQLite statement never reached its progress callback");
2750 }
2751
2752 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2753 async fn cancellation_before_reader_checkout_is_prompt_and_executes_no_statement() {
2754 let dir = tempfile::tempdir().unwrap();
2755 let config = PoolConfig {
2756 path: Some(dir.path().join("sql_bridge_cancel_before_checkout.db")),
2757 max_readers: 1,
2758 checkout_timeout: std::time::Duration::from_secs(5),
2759 ..PoolConfig::default()
2760 };
2761 let pool = Arc::new(ConnectionPool::new(config).unwrap());
2762 pool.writer()
2763 .unwrap()
2764 .conn()
2765 .execute_batch(
2766 "CREATE TABLE checkout_cancel_probe(value INTEGER NOT NULL); \
2767 INSERT INTO checkout_cancel_probe VALUES (0);",
2768 )
2769 .unwrap();
2770 let held_reader = pool.reader().expect("hold the sole pooled reader");
2771 let bridge = SqlBridge::new(Arc::clone(&pool), false);
2775 let mut waiting_reader = bridge.reader().await.unwrap();
2776 let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
2777 let waiting = tokio::spawn(crate::scope_request_read_cancellation(
2778 cancel_rx,
2779 async move {
2780 waiting_reader
2781 .query_row(SqlStatement {
2782 sql: "UPDATE checkout_cancel_probe SET value = value + 1 RETURNING value"
2783 .into(),
2784 params: vec![],
2785 label: Some("must-not-run-after-cancelled-checkout".into()),
2786 })
2787 .await
2788 },
2789 ));
2790
2791 tokio::task::yield_now().await;
2792 cancel_tx.send(true).unwrap();
2793 let result = tokio::time::timeout(std::time::Duration::from_millis(100), waiting)
2794 .await
2795 .expect("cancelled reader checkout waited for the five-second pool timeout")
2796 .expect("checkout task panicked");
2797 assert!(matches!(result, Err(StorageError::Timeout { .. })));
2801
2802 drop(held_reader);
2803 tokio::time::sleep(std::time::Duration::from_millis(25)).await;
2804 let value: i64 = pool
2805 .reader()
2806 .unwrap()
2807 .conn()
2808 .query_row("SELECT value FROM checkout_cancel_probe", [], |row| {
2809 row.get(0)
2810 })
2811 .unwrap();
2812 assert_eq!(
2813 value, 0,
2814 "a DML statement started after its pre-admission checkout was cancelled"
2815 );
2816 assert_eq!(
2817 pool.available_readers(),
2818 1,
2819 "reader checkout leaked a permit"
2820 );
2821 }
2822
2823 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2831 async fn pooled_reader_checkout_timeout_is_a_retryable_admission_timeout() {
2832 let dir = tempfile::tempdir().unwrap();
2833 let config = PoolConfig {
2834 path: Some(dir.path().join("sql_bridge_reader_admission_timeout.db")),
2835 max_readers: 1,
2836 checkout_timeout: std::time::Duration::from_millis(200),
2837 ..PoolConfig::default()
2838 };
2839 let pool = Arc::new(ConnectionPool::new(config).unwrap());
2840 pool.writer()
2841 .unwrap()
2842 .conn()
2843 .execute_batch(
2844 "CREATE TABLE reader_admission_probe(value INTEGER NOT NULL); \
2845 INSERT INTO reader_admission_probe VALUES (0);",
2846 )
2847 .unwrap();
2848 let held_reader = pool.reader().expect("hold the sole pooled reader");
2851
2852 let bridge = SqlBridge::new(Arc::clone(&pool), false);
2853 let mut contender = bridge.reader().await.unwrap();
2854 let blocked = contender
2855 .query_row(SqlStatement {
2856 sql: "SELECT value FROM reader_admission_probe".into(),
2857 params: vec![],
2858 label: Some("reader-admission-timeout-probe".into()),
2859 })
2860 .await;
2861 assert!(
2862 matches!(blocked, Err(StorageError::AdmissionTimeout { .. })),
2863 "an exhausted pooled-reader checkout must be a retryable AdmissionTimeout; got {blocked:?}"
2864 );
2865
2866 drop(held_reader);
2867 tokio::time::sleep(std::time::Duration::from_millis(25)).await;
2868 assert_eq!(
2869 pool.available_readers(),
2870 1,
2871 "reader checkout leaked a permit"
2872 );
2873 }
2874
2875 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2876 async fn abandoned_read_interrupts_sqlite_releases_permit_and_stops_work() {
2877 let dir = tempfile::tempdir().unwrap();
2878 let config = PoolConfig {
2879 path: Some(dir.path().join("sql_bridge_abandoned_read.db")),
2880 max_readers: 1,
2881 checkout_timeout: std::time::Duration::from_millis(500),
2882 ..PoolConfig::default()
2883 };
2884 let pool = Arc::new(ConnectionPool::new(config).unwrap());
2885 let bridge = SqlBridge::new(Arc::clone(&pool), true);
2886 let mut reader = SqliteReader {
2887 handle: Some(open_cached_reader_handle(Arc::clone(&pool)).await.unwrap()),
2888 pool: Arc::clone(&pool),
2889 };
2890 let mut contender = bridge.reader().await.unwrap();
2891 let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
2892 let progress_in_scope = Arc::clone(&progress);
2893
2894 let query = tokio::spawn(crate::scope_test_read_progress(
2895 progress_in_scope,
2896 async move { reader.query_all(deliberately_slow_read_statement()).await },
2897 ));
2898 wait_for_progress(progress.as_ref()).await;
2899 query.abort();
2900 assert!(matches!(query.await, Err(error) if error.is_cancelled()));
2901
2902 tokio::time::timeout(
2903 std::time::Duration::from_millis(500),
2904 contender.query_row(SqlStatement {
2905 sql: "SELECT 1".into(),
2906 params: vec![],
2907 label: None,
2908 }),
2909 )
2910 .await
2911 .expect("abandoned SQLite statement did not return the sole reader promptly")
2912 .expect("reader probe failed after cancellation");
2913
2914 let stopped_at = progress.load(std::sync::atomic::Ordering::SeqCst);
2915 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
2916 assert_eq!(
2917 progress.load(std::sync::atomic::Ordering::SeqCst),
2918 stopped_at,
2919 "SQLite progress kept advancing after the abandoned request returned its reader"
2920 );
2921 }
2922
2923 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2924 async fn request_deadline_interrupts_statement_without_outer_timeout() {
2925 let dir = tempfile::tempdir().unwrap();
2926 let config = PoolConfig {
2927 path: Some(dir.path().join("sql_bridge_request_deadline.db")),
2928 max_readers: 1,
2929 checkout_timeout: std::time::Duration::from_millis(500),
2930 ..PoolConfig::default()
2931 };
2932 let pool = Arc::new(ConnectionPool::new(config).unwrap());
2933 let bridge = SqlBridge::new(Arc::clone(&pool), true);
2934 let mut reader = bridge.reader().await.unwrap();
2935 let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
2936
2937 let result = crate::scope_test_read_progress(
2938 Arc::clone(&progress),
2939 crate::scope_request_read_deadline(std::time::Duration::from_millis(25), async move {
2940 reader.query_all(deliberately_slow_read_statement()).await
2941 }),
2942 )
2943 .await;
2944 assert!(
2945 matches!(result, Err(StorageError::Timeout { .. })),
2946 "deadline must surface as a typed timeout, got {result:?}"
2947 );
2948
2949 let stopped_at = progress.load(std::sync::atomic::Ordering::SeqCst);
2950 assert!(stopped_at > 0, "deadline test never exercised SQLite work");
2951 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
2952 assert_eq!(
2953 progress.load(std::sync::atomic::Ordering::SeqCst),
2954 stopped_at,
2955 "deadline returned while SQLite kept consuming work"
2956 );
2957 }
2958
2959 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2960 async fn progress_handler_cleanup_failure_discards_pooled_connection() {
2961 let dir = tempfile::tempdir().unwrap();
2962 let config = PoolConfig {
2963 path: Some(dir.path().join("sql_bridge_cleanup_failure.db")),
2964 max_readers: 1,
2965 ..PoolConfig::default()
2966 };
2967 let pool = Arc::new(ConnectionPool::new(config).unwrap());
2968 let pool_for_read = Arc::clone(&pool);
2969 let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
2970 let result = crate::read_cancellation::scope_test_read_cleanup_failure(
2971 crate::scope_test_read_progress(
2972 Arc::clone(&progress),
2973 crate::read_cancellation::run_interruptible_read(
2974 StorageCapability::Sql,
2975 "cleanup_failure_probe",
2976 move |scope| {
2977 let mut guard = pool_for_read.reader().map_err(|error| {
2978 StorageError::driver(
2979 StorageCapability::Sql,
2980 "cleanup_failure_probe",
2981 error,
2982 )
2983 })?;
2984 scope.run_pooled_reader(&mut guard, |conn| {
2985 conn.query_row("SELECT 1", [], |row| row.get::<_, i64>(0))
2986 .map_err(|error| map_rusqlite_err(error, "cleanup_failure_probe"))
2987 })
2988 },
2989 ),
2990 ),
2991 )
2992 .await;
2993 assert!(
2994 matches!(result, Err(StorageError::Internal(ref message)) if message.contains("clear failure")),
2995 "injected cleanup failure must be surfaced; got {result:?}"
2996 );
2997 assert_eq!(
2998 pool.available_readers(),
2999 1,
3000 "discard must install a replacement"
3001 );
3002
3003 let calls_after_failed_read = progress.load(std::sync::atomic::Ordering::SeqCst);
3004 let guard = pool.reader().unwrap();
3005 let sum: i64 = guard
3006 .conn()
3007 .query_row(
3008 "WITH RECURSIVE n(x) AS (VALUES(0) UNION ALL SELECT x + 1 FROM n WHERE x < 10000) \
3009 SELECT sum(x) FROM n",
3010 [],
3011 |row| row.get(0),
3012 )
3013 .unwrap();
3014 assert_eq!(sum, 50_005_000);
3015 assert_eq!(
3016 progress.load(std::sync::atomic::Ordering::SeqCst),
3017 calls_after_failed_read,
3018 "a connection whose handler could not be cleared was reused"
3019 );
3020 }
3021
3022 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3023 async fn raw_pooled_reader_quarantines_cleanup_failure_during_unwind() {
3024 let dir = tempfile::tempdir().unwrap();
3025 let pool = Arc::new(
3026 ConnectionPool::new(PoolConfig {
3027 path: Some(dir.path().join("raw_reader_unwind_cleanup.db")),
3028 max_readers: 1,
3029 ..PoolConfig::default()
3030 })
3031 .unwrap(),
3032 );
3033 let worker_pool = Arc::clone(&pool);
3034 let result = crate::read_cancellation::scope_test_read_cleanup_failure(
3035 crate::read_cancellation::run_interruptible_read(
3036 StorageCapability::Sql,
3037 "raw_reader_unwind_cleanup",
3038 move |scope| {
3039 let mut guard = worker_pool.reader().map_err(|error| {
3040 StorageError::driver(
3041 StorageCapability::Sql,
3042 "raw_reader_unwind_cleanup",
3043 error,
3044 )
3045 })?;
3046 scope.with_pooled_reader(&mut guard, |conn| {
3047 scope.run(conn, || -> khive_storage::types::StorageResult<()> {
3048 panic!("injected raw reader panic after progress registration")
3049 })
3050 })
3051 },
3052 ),
3053 )
3054 .await;
3055 assert!(
3056 result.is_err(),
3057 "blocking panic must surface as a join error"
3058 );
3059 assert_eq!(
3060 pool.available_readers(),
3061 pool.max_readers(),
3062 "unwind cleanup failure must close and replace the raw pooled reader"
3063 );
3064 }
3065
3066 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3067 async fn raw_pooled_writer_retires_cleanup_failure_during_unwind() {
3068 let pool = Arc::new(ConnectionPool::new(PoolConfig::default()).unwrap());
3069 let worker_pool = Arc::clone(&pool);
3070 let result = crate::read_cancellation::scope_test_read_cleanup_failure(
3071 crate::read_cancellation::run_interruptible_read(
3072 StorageCapability::Sql,
3073 "raw_writer_unwind_cleanup",
3074 move |scope| {
3075 let guard = worker_pool.try_writer().map_err(|error| {
3076 StorageError::driver(
3077 StorageCapability::Sql,
3078 "raw_writer_unwind_cleanup",
3079 error,
3080 )
3081 })?;
3082 scope.with_pooled_writer(&worker_pool, &guard, |conn| {
3083 scope.run(conn, || -> khive_storage::types::StorageResult<()> {
3084 panic!("injected raw writer panic after progress registration")
3085 })
3086 })
3087 },
3088 ),
3089 )
3090 .await;
3091 assert!(
3092 result.is_err(),
3093 "blocking panic must surface as a join error"
3094 );
3095 assert!(
3096 pool.try_writer().is_err(),
3097 "unwind cleanup failure must retire the raw pooled writer"
3098 );
3099 }
3100
3101 #[test]
3102 fn query_row_converts_only_the_first_matching_row() {
3103 let conn = rusqlite::Connection::open_in_memory().unwrap();
3104 let statement = SqlStatement {
3105 sql: "WITH RECURSIVE rows(value) AS (\
3106 SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 99\
3107 ) SELECT value FROM rows ORDER BY value"
3108 .into(),
3109 params: vec![],
3110 label: None,
3111 };
3112
3113 ROW_CONVERSIONS.with(|count| count.set(0));
3114 let row = execute_query_row(&conn, &statement).unwrap().unwrap();
3115
3116 assert!(matches!(row.get("value"), Some(SqlValue::Integer(0))));
3117 ROW_CONVERSIONS.with(|count| assert_eq!(count.get(), 1));
3118 }
3119
3120 #[test]
3121 fn query_page_bounds_owned_rows_before_full_materialization() {
3122 let conn = rusqlite::Connection::open_in_memory().unwrap();
3123 let statement = SqlStatement {
3124 sql: "WITH RECURSIVE rows(value) AS (\
3125 SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 99\
3126 ) SELECT value FROM rows ORDER BY value"
3127 .into(),
3128 params: vec![],
3129 label: None,
3130 };
3131
3132 ROW_CONVERSIONS.with(|count| count.set(0));
3133 let rows = execute_query_page(
3134 &conn,
3135 &statement,
3136 &PageRequest {
3137 offset: 40,
3138 limit: 3,
3139 },
3140 )
3141 .unwrap();
3142
3143 assert_eq!(rows.len(), 3);
3144 assert!(matches!(rows[0].get("value"), Some(SqlValue::Integer(40))));
3145 assert!(matches!(rows[2].get("value"), Some(SqlValue::Integer(42))));
3146 ROW_CONVERSIONS.with(|count| assert_eq!(count.get(), 3));
3147 }
3148
3149 #[test]
3150 fn query_page_zero_limit_converts_no_rows_but_still_validates_sql() {
3151 let conn = rusqlite::Connection::open_in_memory().unwrap();
3152 let statement = SqlStatement {
3153 sql: "WITH RECURSIVE rows(value) AS (\
3154 SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 99\
3155 ) SELECT value FROM rows ORDER BY value"
3156 .into(),
3157 params: vec![],
3158 label: None,
3159 };
3160
3161 ROW_CONVERSIONS.with(|count| count.set(0));
3162 let rows = execute_query_page(
3163 &conn,
3164 &statement,
3165 &PageRequest {
3166 offset: 0,
3167 limit: 0,
3168 },
3169 )
3170 .unwrap();
3171
3172 assert!(rows.is_empty());
3173 ROW_CONVERSIONS.with(|count| assert_eq!(count.get(), 0));
3174
3175 let invalid = SqlStatement {
3176 sql: "SELECT FROM WHERE".into(),
3177 params: vec![],
3178 label: None,
3179 };
3180 assert!(
3181 execute_query_page(
3182 &conn,
3183 &invalid,
3184 &PageRequest {
3185 offset: 0,
3186 limit: 0
3187 }
3188 )
3189 .is_err(),
3190 "a zero-limit page must still fail on invalid SQL at prepare time"
3191 );
3192 }
3193
3194 #[test]
3195 fn cached_writer_prepare_preserves_the_single_statement_boundary() {
3196 let conn = rusqlite::Connection::open_in_memory().unwrap();
3197 assert!(matches!(
3198 prepare_cached_sql_statement(&conn, "SELECT 1; SELECT 2"),
3199 Err(rusqlite::Error::MultipleStatement)
3200 ));
3201 }
3202
3203 #[tokio::test]
3204 async fn queue_backed_execute_reuses_the_persistent_connection_statement_cache() {
3205 use rusqlite::hooks::{AuthAction, AuthContext, Authorization};
3206 use std::sync::atomic::{AtomicUsize, Ordering};
3207
3208 let dir = tempfile::tempdir().unwrap();
3209 let pool = Arc::new(
3210 ConnectionPool::new(PoolConfig {
3211 path: Some(dir.path().join("sql_bridge_writer_cache.db")),
3212 write_queue_enabled: Some(true),
3213 write_routing_strict: true,
3214 ..PoolConfig::default()
3215 })
3216 .unwrap(),
3217 );
3218 pool.writer()
3219 .unwrap()
3220 .conn()
3221 .execute_batch(
3222 "CREATE TABLE writer_cache_test (id INTEGER PRIMARY KEY, value TEXT NOT NULL)",
3223 )
3224 .unwrap();
3225
3226 let writer_task = pool
3227 .writer_task_handle()
3228 .unwrap()
3229 .expect("file-backed queue-enabled pool must expose its writer task");
3230 let prepare_count = Arc::new(AtomicUsize::new(0));
3231 let hook_count = Arc::clone(&prepare_count);
3232 writer_task
3233 .send_top_level(move |conn| {
3234 conn.authorizer(Some(move |context: AuthContext<'_>| {
3235 if matches!(
3236 context.action,
3237 AuthAction::Insert { table_name } if table_name == "writer_cache_test"
3238 ) {
3239 hook_count.fetch_add(1, Ordering::SeqCst);
3240 }
3241 Authorization::Allow
3242 }))
3243 .map_err(|error| map_rusqlite_err(error, "test.install_authorizer"))
3244 })
3245 .await
3246 .unwrap();
3247
3248 let bridge = SqlBridge::new(Arc::clone(&pool), true);
3249 let mut writer = bridge.writer().await.unwrap();
3250 for id in [1, 2] {
3251 khive_storage::SqlWriter::execute(
3252 &mut *writer,
3253 SqlStatement {
3254 sql: "INSERT INTO writer_cache_test (id, value) VALUES (?1, ?2)".into(),
3255 params: vec![SqlValue::Integer(id), SqlValue::Text(format!("value-{id}"))],
3256 label: None,
3257 },
3258 )
3259 .await
3260 .unwrap();
3261 }
3262
3263 assert_eq!(
3264 prepare_count.load(Ordering::SeqCst),
3265 1,
3266 "the second identical execute on the writer task's persistent connection must reuse \
3267 the cached SQLite statement instead of compiling it again"
3268 );
3269 writer_task
3270 .send_top_level(|conn| {
3271 conn.authorizer(None::<fn(AuthContext<'_>) -> Authorization>)
3272 .map_err(|error| map_rusqlite_err(error, "test.remove_authorizer"))
3273 })
3274 .await
3275 .unwrap();
3276 }
3277
3278 #[test]
3279 fn inline_execute_batch_prepares_each_statement_once() {
3280 use rusqlite::hooks::{AuthAction, AuthContext, Authorization};
3281 use std::sync::atomic::{AtomicUsize, Ordering};
3282
3283 let conn = rusqlite::Connection::open_in_memory().unwrap();
3284 conn.execute_batch(
3285 "CREATE TABLE single_prepare_test (id INTEGER PRIMARY KEY, value TEXT NOT NULL)",
3286 )
3287 .unwrap();
3288 let prepare_count = Arc::new(AtomicUsize::new(0));
3289 let hook_count = Arc::clone(&prepare_count);
3290 conn.authorizer(Some(move |context: AuthContext<'_>| {
3291 if matches!(
3292 context.action,
3293 AuthAction::Insert { table_name } if table_name == "single_prepare_test"
3294 ) {
3295 hook_count.fetch_add(1, Ordering::SeqCst);
3296 }
3297 Authorization::Allow
3298 }))
3299 .unwrap();
3300
3301 let mut writer = InlineWriter {
3302 conn: &conn as *const rusqlite::Connection,
3303 };
3304 let affected = block_on_sync(khive_storage::SqlWriter::execute_batch(
3305 &mut writer,
3306 vec![SqlStatement {
3307 sql: "INSERT INTO single_prepare_test (id, value) VALUES (?1, ?2)".into(),
3308 params: vec![SqlValue::Integer(1), SqlValue::Text("once".into())],
3309 label: None,
3310 }],
3311 ))
3312 .expect("InlineWriter operations must resolve on their first poll")
3313 .expect("valid batch must execute");
3314
3315 assert_eq!(affected, 1);
3316 assert_eq!(
3317 prepare_count.load(Ordering::SeqCst),
3318 1,
3319 "classification and execution must share one prepared statement handle"
3320 );
3321 conn.authorizer(None::<fn(AuthContext<'_>) -> Authorization>)
3322 .unwrap();
3323 }
3324
3325 #[test]
3326 fn inline_execute_batch_preserves_schema_dependencies_between_statements() {
3327 let conn = rusqlite::Connection::open_in_memory().unwrap();
3328 let mut writer = InlineWriter {
3329 conn: &conn as *const rusqlite::Connection,
3330 };
3331
3332 let affected = block_on_sync(khive_storage::SqlWriter::execute_batch(
3333 &mut writer,
3334 vec![
3335 SqlStatement {
3336 sql: "CREATE TABLE dependent_prepare_test (id INTEGER PRIMARY KEY)".into(),
3337 params: vec![],
3338 label: None,
3339 },
3340 SqlStatement {
3341 sql: "INSERT INTO dependent_prepare_test (id) VALUES (1)".into(),
3342 params: vec![],
3343 label: None,
3344 },
3345 ],
3346 ))
3347 .expect("InlineWriter operations must resolve on their first poll")
3348 .expect("a later statement must be prepared after its prerequisite schema change");
3349
3350 assert_eq!(affected, 1);
3351 let count: i64 = conn
3352 .query_row("SELECT COUNT(*) FROM dependent_prepare_test", [], |row| {
3353 row.get(0)
3354 })
3355 .unwrap();
3356 assert_eq!(count, 1);
3357 }
3358
3359 #[tokio::test]
3360 async fn pool_backed_query_page_beyond_result_set_returns_empty() {
3361 let config = PoolConfig {
3362 path: None,
3363 ..PoolConfig::default()
3364 };
3365 let pool = Arc::new(ConnectionPool::new(config).unwrap());
3366 {
3367 let writer = pool.writer().unwrap();
3368 writer
3369 .conn()
3370 .execute_batch(
3371 "CREATE TABLE page_test (id INTEGER PRIMARY KEY, val TEXT NOT NULL);\
3372 INSERT INTO page_test (id, val) VALUES (1, 'a'), (2, 'b'), (3, 'c');",
3373 )
3374 .unwrap();
3375 }
3376 let bridge = SqlBridge::new(Arc::clone(&pool), false);
3377
3378 let statement = || SqlStatement {
3379 sql: "SELECT val FROM page_test ORDER BY id".into(),
3380 params: vec![],
3381 label: None,
3382 };
3383
3384 let mut reader = bridge.reader().await.unwrap();
3385 let page = reader
3386 .query_page(
3387 statement(),
3388 PageRequest {
3389 offset: 1,
3390 limit: 2,
3391 },
3392 )
3393 .await
3394 .unwrap();
3395 assert_eq!(page.len(), 2);
3396 assert!(matches!(page[0].get("val"), Some(SqlValue::Text(v)) if v == "b"));
3397 assert!(matches!(page[1].get("val"), Some(SqlValue::Text(v)) if v == "c"));
3398
3399 let empty = reader
3400 .query_page(
3401 statement(),
3402 PageRequest {
3403 offset: 99,
3404 limit: 10,
3405 },
3406 )
3407 .await
3408 .unwrap();
3409 assert!(
3410 empty.is_empty(),
3411 "offset past the last row must return an empty page, got {empty:?}"
3412 );
3413 drop(reader);
3414
3415 let mut writer = bridge.writer().await.unwrap();
3416 let empty = writer
3417 .query_page(
3418 statement(),
3419 PageRequest {
3420 offset: 99,
3421 limit: 10,
3422 },
3423 )
3424 .await
3425 .unwrap();
3426 assert!(
3427 empty.is_empty(),
3428 "offset past the last row must return an empty page, got {empty:?}"
3429 );
3430 }
3431
3432 #[tokio::test]
3433 async fn file_bridge_scopes_reader_permits_to_operations_and_caps_writer_handles() {
3434 let dir = tempfile::tempdir().unwrap();
3435 let config = PoolConfig {
3436 path: Some(dir.path().join("sql_bridge_handle_cap.db")),
3437 write_queue_enabled: Some(false),
3438 max_readers: 2,
3439 checkout_timeout: std::time::Duration::from_millis(20),
3440 ..PoolConfig::default()
3441 };
3442 let pool = Arc::new(ConnectionPool::new(config).unwrap());
3443 let bridge = SqlBridge::new(Arc::clone(&pool), true);
3444 let second_bridge = SqlBridge::new(Arc::clone(&pool), true);
3445
3446 let mut retained_readers = Vec::new();
3447 for expected in 0..3 {
3448 let mut reader = second_bridge.reader().await.unwrap();
3449 let value = reader
3450 .query_scalar(SqlStatement {
3451 sql: format!("SELECT {expected}"),
3452 params: vec![],
3453 label: None,
3454 })
3455 .await
3456 .unwrap();
3457 assert!(matches!(value, Some(SqlValue::Integer(value)) if value == expected));
3458 retained_readers.push(reader);
3459 }
3460 assert_eq!(retained_readers.len(), 3);
3461
3462 let mut additional_reader = bridge.reader().await.unwrap();
3463 let page = additional_reader
3464 .query_page(
3465 SqlStatement {
3466 sql: "WITH RECURSIVE rows(value) AS (\
3467 SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 9\
3468 ) SELECT value FROM rows ORDER BY value"
3469 .into(),
3470 params: vec![],
3471 label: None,
3472 },
3473 PageRequest {
3474 offset: 7,
3475 limit: 2,
3476 },
3477 )
3478 .await
3479 .unwrap();
3480 assert_eq!(page.len(), 2);
3481 assert!(matches!(page[0].get("value"), Some(SqlValue::Integer(7))));
3482 assert!(matches!(page[1].get("value"), Some(SqlValue::Integer(8))));
3483 drop((additional_reader, retained_readers));
3484
3485 let writer = bridge.writer().await.unwrap();
3486 let writer_error = match second_bridge.writer().await {
3487 Ok(_) => panic!("a second live writer handle exceeded the one-handle cap"),
3488 Err(error) => error,
3489 };
3490 assert!(matches!(
3491 writer_error,
3492 StorageError::AdmissionTimeout { ref operation, .. }
3493 if operation.as_ref() == "sql_bridge.writer_handle"
3494 ));
3495 drop(writer);
3496 let writer_after_release = bridge.writer().await.unwrap();
3497 drop(writer_after_release);
3498 }
3499
3500 #[tokio::test]
3501 #[serial_test::serial(tx_registry)]
3502 async fn cached_read_transaction_retains_one_permit_until_commit_or_rollback() {
3503 let dir = tempfile::tempdir().unwrap();
3504 let config = PoolConfig {
3505 path: Some(dir.path().join("sql_bridge_reader_tx_control.db")),
3506 write_queue_enabled: Some(true),
3507 max_readers: 1,
3508 checkout_timeout: std::time::Duration::from_millis(20),
3509 ..PoolConfig::default()
3510 };
3511 let pool = Arc::new(ConnectionPool::new(config).unwrap());
3512 let origin = pool.origin();
3513 let origin_view = database_tx_view(&pool);
3514 let unrelated_view = khive_storage::tx_registry::TxOriginFilter::Secondary(
3515 khive_storage::tx_registry::DbIdentity::new("unrelated-sql-bridge.db"),
3516 );
3517 let bridge = SqlBridge::new(Arc::clone(&pool), true);
3518 let mut reader = bridge.reader().await.unwrap();
3519 let mut contender = bridge.reader().await.unwrap();
3520
3521 assert!(
3522 khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
3523 "an idle cached reader must not register a transaction"
3524 );
3525
3526 reader
3527 .query_all(SqlStatement {
3528 sql: "BEGIN DEFERRED".into(),
3529 params: vec![],
3530 label: None,
3531 })
3532 .await
3533 .expect("BEGIN DEFERRED must open an admitted cached-reader snapshot");
3534 let opened = khive_storage::tx_registry::oldest_for(&origin_view)
3535 .expect("successful BEGIN must register the cached-reader transaction");
3536 assert_eq!(opened.label.as_deref(), Some(CACHED_READ_TRANSACTION_LABEL));
3537 assert_eq!(opened.origin, origin);
3538 assert!(
3539 khive_storage::tx_registry::oldest_for(&unrelated_view).is_none(),
3540 "the read transaction must be attributed only to its own backend"
3541 );
3542 assert_eq!(
3543 pool.sql_bridge_reader_slots().available_permits(),
3544 0,
3545 "the successful BEGIN must retain its operation permit"
3546 );
3547
3548 let value = reader
3549 .query_scalar(SqlStatement {
3550 sql: "SELECT 7".into(),
3551 params: vec![],
3552 label: None,
3553 })
3554 .await
3555 .expect("a query inside the admitted transaction must reuse its retained permit");
3556 assert!(matches!(value, Some(SqlValue::Integer(7))));
3557 assert_eq!(
3558 khive_storage::tx_registry::oldest_for(&origin_view)
3559 .expect("queries must retain the transaction registration")
3560 .id,
3561 opened.id,
3562 "queries inside the transaction must retain the original span"
3563 );
3564
3565 let blocked = contender
3566 .query_scalar(SqlStatement {
3567 sql: "SELECT 8".into(),
3568 params: vec![],
3569 label: None,
3570 })
3571 .await;
3572 assert!(
3573 matches!(
3574 &blocked,
3575 Err(StorageError::Timeout { operation })
3576 if operation.as_ref() == "sql_bridge.reader_operation"
3577 ),
3578 "a second logical read must contend with the admitted transaction \
3579 and time out with the ADR-005 reader contract error; got {blocked:?}"
3580 );
3581
3582 reader
3583 .query_all(SqlStatement {
3584 sql: "COMMIT".into(),
3585 params: vec![],
3586 label: None,
3587 })
3588 .await
3589 .expect("COMMIT must close the admitted cached-reader snapshot");
3590 assert!(
3591 khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
3592 "COMMIT must deregister after SQLite returns to autocommit"
3593 );
3594 assert_eq!(
3595 pool.sql_bridge_reader_slots().available_permits(),
3596 1,
3597 "COMMIT may release the permit only after autocommit is restored"
3598 );
3599 let value = contender
3600 .query_scalar(SqlStatement {
3601 sql: "SELECT 8".into(),
3602 params: vec![],
3603 label: None,
3604 })
3605 .await
3606 .expect("the contender must run after COMMIT releases admission");
3607 assert!(matches!(value, Some(SqlValue::Integer(8))));
3608
3609 reader
3610 .query_all(SqlStatement {
3611 sql: "BEGIN TRANSACTION".into(),
3612 params: vec![],
3613 label: None,
3614 })
3615 .await
3616 .expect("plain deferred BEGIN TRANSACTION must also be admitted");
3617 let reopened = khive_storage::tx_registry::oldest_for(&origin_view)
3618 .expect("the second successful BEGIN must register a fresh span");
3619 assert_ne!(reopened.id, opened.id);
3620 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 0);
3621 let nested = reader
3622 .query_all(SqlStatement {
3623 sql: "ROLLBACK TO stale_snapshot".into(),
3624 params: vec![],
3625 label: None,
3626 })
3627 .await;
3628 assert!(
3629 matches!(&nested, Err(StorageError::InvalidInput { .. })),
3630 "ROLLBACK TO requires unsupported nested state; got {nested:?}"
3631 );
3632 assert_eq!(
3633 pool.sql_bridge_reader_slots().available_permits(),
3634 0,
3635 "rejected nested control must not release the still-live transaction admission"
3636 );
3637 assert_eq!(
3638 khive_storage::tx_registry::oldest_for(&origin_view)
3639 .expect("ROLLBACK TO rejection must retain the live span")
3640 .id,
3641 reopened.id
3642 );
3643 let savepoint = reader
3644 .query_all(SqlStatement {
3645 sql: "SAVEPOINT nested_snapshot".into(),
3646 params: vec![],
3647 label: None,
3648 })
3649 .await;
3650 assert!(
3651 matches!(&savepoint, Err(StorageError::InvalidInput { .. })),
3652 "SAVEPOINT must be rejected inside the admitted transaction; got {savepoint:?}"
3653 );
3654 assert_eq!(
3655 khive_storage::tx_registry::oldest_for(&origin_view)
3656 .expect("SAVEPOINT rejection must retain the live span")
3657 .id,
3658 reopened.id
3659 );
3660 reader
3661 .query_all(SqlStatement {
3662 sql: "ROLLBACK".into(),
3663 params: vec![],
3664 label: None,
3665 })
3666 .await
3667 .expect("ROLLBACK must close the admitted cached-reader snapshot");
3668 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
3669 assert!(
3670 khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
3671 "full ROLLBACK must deregister after SQLite returns to autocommit"
3672 );
3673 }
3674
3675 #[tokio::test]
3676 async fn failed_cached_reader_begin_does_not_register_a_transaction() {
3677 use rusqlite::hooks::{AuthAction, AuthContext, Authorization, TransactionOperation};
3678
3679 fn deny_begin(ctx: AuthContext<'_>) -> Authorization {
3680 match ctx.action {
3681 AuthAction::Transaction {
3682 operation: TransactionOperation::Begin,
3683 } => Authorization::Deny,
3684 _ => Authorization::Allow,
3685 }
3686 }
3687
3688 let dir = tempfile::tempdir().unwrap();
3689 let config = PoolConfig {
3690 path: Some(dir.path().join("sql_bridge_reader_failed_begin.db")),
3691 max_readers: 1,
3692 ..PoolConfig::default()
3693 };
3694 let pool = Arc::new(ConnectionPool::new(config).unwrap());
3695 let origin_view = database_tx_view(&pool);
3696 let conn = open_standalone_reader(&pool).unwrap();
3697 conn.authorizer(Some(deny_begin)).unwrap();
3698 let mut reader = SqliteReader {
3699 handle: Some(StandaloneHandle {
3700 conn,
3701 _retained_slot: None,
3702 read_transaction_slot: None,
3703 }),
3704 pool: Arc::clone(&pool),
3705 };
3706
3707 let begin = reader
3708 .query_all(SqlStatement {
3709 sql: "BEGIN DEFERRED".into(),
3710 params: vec![],
3711 label: None,
3712 })
3713 .await;
3714 assert!(begin.is_err(), "the authorizer must reject BEGIN");
3715 assert!(
3716 khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
3717 "a failed BEGIN must never enter the transaction registry"
3718 );
3719 assert_eq!(
3720 pool.sql_bridge_reader_slots().available_permits(),
3721 1,
3722 "a failed BEGIN must return the operation permit"
3723 );
3724 }
3725
3726 #[tokio::test]
3727 #[serial_test::serial(tx_registry)]
3728 async fn failed_cached_reader_rollback_deregisters_only_when_connection_is_discarded() {
3729 use rusqlite::hooks::{AuthAction, AuthContext, Authorization, TransactionOperation};
3730
3731 fn deny_rollback(ctx: AuthContext<'_>) -> Authorization {
3732 match ctx.action {
3733 AuthAction::Transaction {
3734 operation: TransactionOperation::Rollback,
3735 } => Authorization::Deny,
3736 _ => Authorization::Allow,
3737 }
3738 }
3739
3740 let dir = tempfile::tempdir().unwrap();
3741 let config = PoolConfig {
3742 path: Some(dir.path().join("sql_bridge_reader_failed_rollback.db")),
3743 max_readers: 1,
3744 ..PoolConfig::default()
3745 };
3746 let pool = Arc::new(ConnectionPool::new(config).unwrap());
3747 let origin_view = database_tx_view(&pool);
3748 let conn = open_standalone_reader(&pool).unwrap();
3749 let mut reader = SqliteReader {
3750 handle: Some(StandaloneHandle {
3751 conn,
3752 _retained_slot: None,
3753 read_transaction_slot: None,
3754 }),
3755 pool: Arc::clone(&pool),
3756 };
3757
3758 reader
3759 .query_all(SqlStatement {
3760 sql: "BEGIN DEFERRED".into(),
3761 params: vec![],
3762 label: None,
3763 })
3764 .await
3765 .expect("BEGIN must establish the registered transaction");
3766 let opened = khive_storage::tx_registry::oldest_for(&origin_view)
3767 .expect("the admitted transaction must be registered");
3768 reader
3769 .handle
3770 .as_ref()
3771 .expect("reader must retain its connection")
3772 .conn
3773 .authorizer(Some(deny_rollback))
3774 .unwrap();
3775
3776 let rollback = reader
3777 .query_all(SqlStatement {
3778 sql: "ROLLBACK".into(),
3779 params: vec![],
3780 label: None,
3781 })
3782 .await;
3783 assert!(rollback.is_err(), "the authorizer must reject ROLLBACK");
3784 assert_eq!(
3785 khive_storage::tx_registry::oldest_for(&origin_view)
3786 .expect("failed ROLLBACK must retain registry evidence")
3787 .id,
3788 opened.id
3789 );
3790 assert_eq!(
3791 pool.sql_bridge_reader_slots().available_permits(),
3792 0,
3793 "failed ROLLBACK must retain reader admission"
3794 );
3795
3796 drop(reader);
3797 assert!(
3798 khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
3799 "discarding the connection must not leak its registry entry"
3800 );
3801 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
3802 }
3803
3804 #[tokio::test]
3805 #[serial_test::serial(tx_registry)]
3806 async fn cached_read_only_handles_reject_unsupported_transaction_control_without_consumption() {
3807 let dir = tempfile::tempdir().unwrap();
3808 let config = PoolConfig {
3809 path: Some(
3810 dir.path()
3811 .join("sql_bridge_reader_unsupported_tx_control.db"),
3812 ),
3813 write_queue_enabled: Some(true),
3814 max_readers: 1,
3815 ..PoolConfig::default()
3816 };
3817 let pool = Arc::new(ConnectionPool::new(config).unwrap());
3818 let bridge = SqlBridge::new(Arc::clone(&pool), true);
3819
3820 let mut reader = bridge.reader().await.unwrap();
3821 for (sql, keyword) in [
3822 ("BEGIN IMMEDIATE", "BEGIN"),
3823 ("BEGIN EXCLUSIVE", "BEGIN"),
3824 ("BEGIN TRANSACTION IMMEDIATE", "BEGIN"),
3829 ("BEGIN TRANSACTION EXCLUSIVE", "BEGIN"),
3830 ("BEGIN DEFERRED TRANSACTION trailing", "BEGIN"),
3831 ("BEGIN TRANSACTION \"IMMEDIATE\"", "BEGIN"),
3836 ("BEGIN TRANSACTION [IMMEDIATE]", "BEGIN"),
3837 ("BEGIN TRANSACTION `IMMEDIATE`", "BEGIN"),
3838 ("BEGIN TRANSACTION 'IMMEDIATE'", "BEGIN"),
3839 ("BEGIN \"DEFERRED\"", "BEGIN"),
3840 ("BEGIN; COMMIT", "BEGIN"),
3841 ("START TRANSACTION", "START"),
3842 ("COMMIT", "COMMIT"),
3843 ] {
3844 let rejected = reader
3845 .query_all(SqlStatement {
3846 sql: sql.into(),
3847 params: vec![],
3848 label: None,
3849 })
3850 .await;
3851 assert!(
3852 matches!(
3853 &rejected,
3854 Err(StorageError::InvalidInput { message, .. })
3855 if message.contains(keyword)
3856 ),
3857 "unsupported cached-reader control {sql:?} must fail closed; got {rejected:?}"
3858 );
3859 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
3860 }
3861
3862 let mut queue_backed_writer = bridge.writer().await.unwrap();
3863 let rejected = queue_backed_writer
3864 .query_all(SqlStatement {
3865 sql: "SAVEPOINT stale_snapshot".into(),
3866 params: vec![],
3867 label: None,
3868 })
3869 .await;
3870 assert!(
3871 matches!(
3872 &rejected,
3873 Err(StorageError::InvalidInput {
3874 operation,
3875 message,
3876 ..
3877 }) if operation.as_ref() == "writer.query_all"
3878 && message.contains("transaction control")
3879 && message.contains("SAVEPOINT")
3880 ),
3881 "a queue-backed writer's cached read-only connection must reject transaction \
3882 control; got {rejected:?}"
3883 );
3884 assert_eq!(
3885 pool.sql_bridge_reader_slots().available_permits(),
3886 1,
3887 "queue-backed rejection must leave the operation permit available"
3888 );
3889 let value = queue_backed_writer
3890 .query_scalar(SqlStatement {
3891 sql: "SELECT 8".into(),
3892 params: vec![],
3893 label: None,
3894 })
3895 .await
3896 .expect("transaction-control rejection must not consume the queue-backed handle");
3897 assert!(matches!(value, Some(SqlValue::Integer(8))));
3898
3899 queue_backed_writer
3900 .query_all(SqlStatement {
3901 sql: "BEGIN DEFERRED".into(),
3902 params: vec![],
3903 label: None,
3904 })
3905 .await
3906 .expect("queue-backed cached reader must share explicit read-transaction admission");
3907 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 0);
3908 queue_backed_writer
3909 .query_all(SqlStatement {
3910 sql: "END".into(),
3911 params: vec![],
3912 label: None,
3913 })
3914 .await
3915 .expect("END must release queue-backed cached-reader admission");
3916 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
3917 }
3918
3919 #[tokio::test]
3920 #[serial_test::serial(tx_registry)]
3921 async fn dropping_cached_reader_transaction_closes_snapshot_before_releasing_permit() {
3922 let dir = tempfile::tempdir().unwrap();
3923 let config = PoolConfig {
3924 path: Some(dir.path().join("sql_bridge_reader_tx_drop.db")),
3925 max_readers: 1,
3926 checkout_timeout: std::time::Duration::from_millis(20),
3927 ..PoolConfig::default()
3928 };
3929 let pool = Arc::new(ConnectionPool::new(config).unwrap());
3930 let origin_view = database_tx_view(&pool);
3931 let bridge = SqlBridge::new(Arc::clone(&pool), true);
3932 let mut reader = bridge.reader().await.unwrap();
3933 let mut contender = bridge.reader().await.unwrap();
3934
3935 reader
3936 .query_all(SqlStatement {
3937 sql: "BEGIN DEFERRED".into(),
3938 params: vec![],
3939 label: None,
3940 })
3941 .await
3942 .expect("begin admitted transaction");
3943 reader
3944 .query_all(SqlStatement {
3945 sql: "SELECT * FROM sqlite_schema".into(),
3946 params: vec![],
3947 label: None,
3948 })
3949 .await
3950 .expect("materialize read snapshot");
3951 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 0);
3952 assert!(
3953 khive_storage::tx_registry::oldest_for(&origin_view).is_some(),
3954 "the live snapshot must remain registered until handle drop"
3955 );
3956
3957 drop(reader);
3958 assert!(
3959 khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
3960 "handle drop must close SQLite before deregistering the snapshot"
3961 );
3962 assert_eq!(
3963 pool.sql_bridge_reader_slots().available_permits(),
3964 1,
3965 "dropping the handle must close its transaction before returning admission"
3966 );
3967 contender
3968 .query_all(SqlStatement {
3969 sql: "SELECT * FROM sqlite_schema".into(),
3970 params: vec![],
3971 label: None,
3972 })
3973 .await
3974 .expect("a new operation must run after the transactional handle drops");
3975 }
3976
3977 #[tokio::test]
3986 #[serial_test::serial(tx_registry)]
3987 async fn expired_cached_reader_transaction_is_rolled_back_on_reuse() {
3988 let dir = tempfile::tempdir().unwrap();
3989 let config = PoolConfig {
3990 path: Some(dir.path().join("sql_bridge_reader_tx_max_age.db")),
3991 max_readers: 1,
3992 checkout_timeout: std::time::Duration::from_millis(20),
3993 read_tx_max_age: std::time::Duration::from_millis(20),
3994 ..PoolConfig::default()
3995 };
3996 let pool = Arc::new(ConnectionPool::new(config).unwrap());
3997 let origin_view = database_tx_view(&pool);
3998 let bridge = SqlBridge::new(Arc::clone(&pool), true);
3999 let mut reader = bridge.reader().await.unwrap();
4000
4001 reader
4002 .query_all(SqlStatement {
4003 sql: "BEGIN DEFERRED".into(),
4004 params: vec![],
4005 label: None,
4006 })
4007 .await
4008 .expect("begin admitted transaction");
4009 reader
4010 .query_all(SqlStatement {
4011 sql: "SELECT * FROM sqlite_schema".into(),
4012 params: vec![],
4013 label: None,
4014 })
4015 .await
4016 .expect("materialize read snapshot");
4017 assert!(
4018 khive_storage::tx_registry::oldest_for(&origin_view).is_some(),
4019 "the open transaction must be registered before it ages out"
4020 );
4021
4022 tokio::time::sleep(std::time::Duration::from_millis(40)).await;
4023
4024 let evictions_before = crate::checkpoint::read_tx_max_age_evictions();
4025 let error = reader
4026 .query_all(SqlStatement {
4027 sql: "SELECT * FROM sqlite_schema".into(),
4028 params: vec![],
4029 label: None,
4030 })
4031 .await
4032 .expect_err("reusing a transaction past read_tx_max_age must be refused");
4033 assert!(
4034 error.is_retryable(),
4035 "an evicted-transaction error must be retryable so the caller can open a fresh \
4036 snapshot: {error}"
4037 );
4038 match &error {
4039 StorageError::ReadTransactionAgeEvicted {
4040 operation,
4041 max_age_secs,
4042 } => {
4043 assert_eq!(operation.as_ref(), "query_all");
4044 assert_eq!(
4045 *max_age_secs, 0,
4046 "a 20ms read_tx_max_age truncates to 0 whole seconds"
4047 );
4048 }
4049 other => panic!(
4050 "a clean age-triggered rollback must surface the dedicated \
4051 ReadTransactionAgeEvicted variant, not a generic classification: {other:?}"
4052 ),
4053 }
4054 assert_eq!(
4055 crate::checkpoint::read_tx_max_age_evictions(),
4056 evictions_before + 1,
4057 "the eviction must be counted in the #1846 diagnostics gauge"
4058 );
4059 assert!(
4060 khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
4061 "the expired transaction must be rolled back and deregistered rather than \
4062 continuing to pin the WAL snapshot"
4063 );
4064
4065 reader
4066 .query_all(SqlStatement {
4067 sql: "SELECT * FROM sqlite_schema".into(),
4068 params: vec![],
4069 label: None,
4070 })
4071 .await
4072 .expect("the handle must remain usable for a fresh autocommit read after eviction");
4073 }
4074
4075 #[tokio::test]
4084 #[serial_test::serial(tx_registry)]
4085 async fn expired_cached_reader_transaction_rollback_denial_discards_connection_and_releases_admission(
4086 ) {
4087 use rusqlite::hooks::{AuthAction, AuthContext, Authorization, TransactionOperation};
4088
4089 fn deny_rollback(ctx: AuthContext<'_>) -> Authorization {
4090 match ctx.action {
4091 AuthAction::Transaction {
4092 operation: TransactionOperation::Rollback,
4093 } => Authorization::Deny,
4094 _ => Authorization::Allow,
4095 }
4096 }
4097
4098 let dir = tempfile::tempdir().unwrap();
4099 let config = PoolConfig {
4100 path: Some(
4101 dir.path()
4102 .join("sql_bridge_reader_tx_max_age_rollback_denied.db"),
4103 ),
4104 max_readers: 1,
4105 checkout_timeout: std::time::Duration::from_millis(20),
4106 read_tx_max_age: std::time::Duration::from_millis(20),
4107 ..PoolConfig::default()
4108 };
4109 let pool = Arc::new(ConnectionPool::new(config).unwrap());
4110 let origin_view = database_tx_view(&pool);
4111 let conn = open_standalone_reader(&pool).unwrap();
4112 let mut reader = SqliteReader {
4113 handle: Some(StandaloneHandle {
4114 conn,
4115 _retained_slot: None,
4116 read_transaction_slot: None,
4117 }),
4118 pool: Arc::clone(&pool),
4119 };
4120
4121 reader
4122 .query_all(SqlStatement {
4123 sql: "BEGIN DEFERRED".into(),
4124 params: vec![],
4125 label: None,
4126 })
4127 .await
4128 .expect("begin admitted transaction");
4129 reader
4130 .query_all(SqlStatement {
4131 sql: "SELECT * FROM sqlite_schema".into(),
4132 params: vec![],
4133 label: None,
4134 })
4135 .await
4136 .expect("materialize read snapshot");
4137 assert!(
4138 khive_storage::tx_registry::oldest_for(&origin_view).is_some(),
4139 "the open transaction must be registered before it ages out"
4140 );
4141
4142 reader
4143 .handle
4144 .as_ref()
4145 .expect("reader must retain its connection")
4146 .conn
4147 .authorizer(Some(deny_rollback))
4148 .unwrap();
4149
4150 tokio::time::sleep(std::time::Duration::from_millis(40)).await;
4151
4152 let evictions_before = crate::checkpoint::read_tx_max_age_evictions();
4153 let error = reader
4154 .query_all(SqlStatement {
4155 sql: "SELECT * FROM sqlite_schema".into(),
4156 params: vec![],
4157 label: None,
4158 })
4159 .await
4160 .expect_err("a denied rollback on an expired transaction must surface an error");
4161 assert!(
4162 error.is_retryable(),
4163 "even a failed cleanup rollback must remain classified retryable so callers open a \
4164 fresh handle: {error}"
4165 );
4166 match &error {
4167 StorageError::ReadTransactionAgeEvictionCleanupFailed {
4168 operation,
4169 max_age_secs,
4170 message,
4171 } => {
4172 assert_eq!(operation.as_ref(), "query_all");
4173 assert_eq!(
4174 *max_age_secs, 0,
4175 "a 20ms read_tx_max_age truncates to 0 whole seconds"
4176 );
4177 assert!(
4178 message.contains("rollback failed"),
4179 "the failure must be attributable to the denied ROLLBACK, not silent \
4180 success: {message}"
4181 );
4182 }
4183 other => panic!(
4184 "a denied cleanup rollback must surface the dedicated \
4185 ReadTransactionAgeEvictionCleanupFailed variant, not a generic Transaction \
4186 error the caller cannot machine-detect: {other:?}"
4187 ),
4188 }
4189 assert_eq!(
4190 crate::checkpoint::read_tx_max_age_evictions(),
4191 evictions_before + 1,
4192 "the eviction attempt must still be counted even though cleanup failed"
4193 );
4194 assert!(
4195 khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
4196 "a denied rollback must discard the connection and deregister the expired \
4197 transaction span rather than leaking it"
4198 );
4199 assert_eq!(
4200 pool.sql_bridge_reader_slots().available_permits(),
4201 1,
4202 "discarding the poisoned connection must release the reader admission slot"
4203 );
4204
4205 let reuse = reader
4206 .query_all(SqlStatement {
4207 sql: "SELECT * FROM sqlite_schema".into(),
4208 params: vec![],
4209 label: None,
4210 })
4211 .await;
4212 let message = match reuse {
4213 Err(StorageError::Pool { message, .. }) => message,
4214 other => panic!(
4215 "reusing this discarded reader must fail loudly with 'connection already \
4216 consumed' rather than silently reopening; got {other:?}"
4217 ),
4218 };
4219 assert!(
4220 message.contains("connection already consumed"),
4221 "expected the discarded reader's reuse error to name the pinned failure; got \
4222 {message:?}"
4223 );
4224 }
4225
4226 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4227 #[serial_test::serial(tx_registry)]
4228 async fn cancelled_cached_reader_transaction_releases_guards_after_connection_closes() {
4229 let dir = tempfile::tempdir().unwrap();
4230 let config = PoolConfig {
4231 path: Some(dir.path().join("sql_bridge_reader_tx_drop_cancel.db")),
4232 max_readers: 1,
4233 checkout_timeout: std::time::Duration::from_millis(50),
4234 ..PoolConfig::default()
4235 };
4236 let pool = Arc::new(ConnectionPool::new(config).unwrap());
4237 let origin_view = database_tx_view(&pool);
4238 let bridge = SqlBridge::new(Arc::clone(&pool), true);
4239 let mut reader = SqliteReader {
4240 handle: Some(open_cached_reader_handle(Arc::clone(&pool)).await.unwrap()),
4241 pool: Arc::clone(&pool),
4242 };
4243 let mut contender = bridge.reader().await.unwrap();
4244
4245 reader
4246 .query_all(SqlStatement {
4247 sql: "BEGIN DEFERRED".into(),
4248 params: vec![],
4249 label: None,
4250 })
4251 .await
4252 .expect("begin admitted transaction");
4253 assert!(
4254 khive_storage::tx_registry::oldest_for(&origin_view).is_some(),
4255 "the explicit transaction must be registered before cancellation"
4256 );
4257 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 0);
4258
4259 let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
4260 let query = tokio::spawn(crate::scope_test_read_progress(
4261 Arc::clone(&progress),
4262 async move { reader.query_all(deliberately_slow_read_statement()).await },
4263 ));
4264 wait_for_progress(progress.as_ref()).await;
4265 query.abort();
4266 assert!(matches!(query.await, Err(error) if error.is_cancelled()));
4267 tokio::time::timeout(std::time::Duration::from_secs(1), async {
4268 while khive_storage::tx_registry::oldest_for(&origin_view).is_some()
4269 || pool.sql_bridge_reader_slots().available_permits() != 1
4270 {
4271 tokio::task::yield_now().await;
4272 }
4273 })
4274 .await
4275 .expect("connection cleanup leaked transaction evidence or reader admission");
4276 contender
4277 .query_all(SqlStatement {
4278 sql: "SELECT 1".into(),
4279 params: vec![],
4280 label: None,
4281 })
4282 .await
4283 .expect("admission must recover after the cancelled connection closes");
4284 }
4285
4286 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4287 #[serial_test::serial(tx_registry)]
4288 async fn cancelled_cached_reader_rolls_back_releases_wal_and_clears_handler() {
4289 let dir = tempfile::tempdir().unwrap();
4290 let config = PoolConfig {
4291 path: Some(dir.path().join("sql_bridge_reader_tx_request_cancel.db")),
4292 max_readers: 1,
4293 checkout_timeout: std::time::Duration::from_millis(500),
4294 ..PoolConfig::default()
4295 };
4296 let pool = Arc::new(ConnectionPool::new(config).unwrap());
4297 let bridge = SqlBridge::new(Arc::clone(&pool), true);
4298 let writer = open_standalone_writer(&pool).unwrap();
4299 writer
4300 .execute_batch(
4301 "CREATE TABLE snapshot_probe(id INTEGER PRIMARY KEY, value TEXT NOT NULL); \
4302 INSERT INTO snapshot_probe(value) VALUES ('seed');",
4303 )
4304 .unwrap();
4305 let mut reader = bridge.reader().await.unwrap();
4306 let mut contender = bridge.reader().await.unwrap();
4307
4308 reader
4309 .query_all(SqlStatement {
4310 sql: "BEGIN DEFERRED".into(),
4311 params: vec![],
4312 label: None,
4313 })
4314 .await
4315 .expect("begin admitted transaction");
4316 reader
4317 .query_all(SqlStatement {
4318 sql: "SELECT * FROM snapshot_probe".into(),
4319 params: vec![],
4320 label: None,
4321 })
4322 .await
4323 .expect("materialize a real WAL snapshot");
4324 writer
4325 .execute_batch(
4326 "WITH RECURSIVE rows(value) AS (\
4327 SELECT 1 UNION ALL SELECT value + 1 FROM rows WHERE value < 100\
4328 ) INSERT INTO snapshot_probe(value) SELECT printf('row-%d', value) FROM rows;",
4329 )
4330 .unwrap();
4331 let (_, log_before, checkpointed_before) = passive_checkpoint(&writer);
4332 assert!(
4333 log_before > checkpointed_before,
4334 "the explicit reader snapshot must pin WAL frames before cancellation"
4335 );
4336
4337 let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
4338 let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
4339 let query = tokio::spawn(crate::scope_test_read_progress(
4340 Arc::clone(&progress),
4341 crate::scope_request_read_cancellation(cancel_rx, async move {
4342 let result = reader.query_all(deliberately_slow_read_statement()).await;
4343 (reader, result)
4344 }),
4345 ));
4346 wait_for_progress(progress.as_ref()).await;
4347 cancel_tx.send(true).unwrap();
4348 let (mut reader, result) = tokio::time::timeout(std::time::Duration::from_secs(1), query)
4349 .await
4350 .expect("interrupted explicit read transaction did not stop promptly")
4351 .unwrap();
4352 assert!(
4353 matches!(result, Err(StorageError::Timeout { .. })),
4354 "request cancellation must surface as a typed timeout; got {result:?}"
4355 );
4356 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
4357
4358 let (_, log_after, checkpointed_after) = passive_checkpoint(&writer);
4359 assert_eq!(
4360 log_after, checkpointed_after,
4361 "cancellation must release the explicit reader's WAL snapshot"
4362 );
4363 contender
4364 .query_all(SqlStatement {
4365 sql: "SELECT 1".into(),
4366 params: vec![],
4367 label: None,
4368 })
4369 .await
4370 .expect("the sole reader permit must be reusable after rollback");
4371
4372 let stopped_at = progress.load(std::sync::atomic::Ordering::SeqCst);
4373 reader
4374 .query_all(SqlStatement {
4375 sql: "WITH RECURSIVE rows(value) AS (\
4376 SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 10000\
4377 ) SELECT SUM(value) FROM rows"
4378 .into(),
4379 params: vec![],
4380 label: None,
4381 })
4382 .await
4383 .expect("same connection must remain usable after handler teardown");
4384 assert_eq!(
4385 progress.load(std::sync::atomic::Ordering::SeqCst),
4386 stopped_at,
4387 "the cancelled request's progress callback bled into the next borrower"
4388 );
4389 }
4390
4391 fn register_khive_test_slow_udf(
4399 conn: &rusqlite::Connection,
4400 sleep_ms: u64,
4401 started: Arc<std::sync::atomic::AtomicBool>,
4402 ) {
4403 conn.create_scalar_function(
4404 "khive_test_slow_udf",
4405 0,
4406 rusqlite::functions::FunctionFlags::SQLITE_UTF8,
4407 move |_| {
4408 started.store(true, std::sync::atomic::Ordering::Release);
4409 std::thread::sleep(std::time::Duration::from_millis(sleep_ms));
4410 Ok(0i64)
4411 },
4412 )
4413 .unwrap();
4414 }
4415
4416 async fn wait_for_flag(flag: &std::sync::atomic::AtomicBool) {
4417 tokio::time::timeout(std::time::Duration::from_secs(1), async {
4418 while !flag.load(std::sync::atomic::Ordering::Acquire) {
4419 tokio::task::yield_now().await;
4420 }
4421 })
4422 .await
4423 .expect("slow UDF never started");
4424 }
4425
4426 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4427 async fn abandoned_read_past_grace_recovers_admission_after_bounded_join() {
4428 let dir = tempfile::tempdir().unwrap();
4436 let config = PoolConfig {
4437 path: Some(dir.path().join("sql_bridge_grace_exceeded.db")),
4438 max_readers: 1,
4439 checkout_timeout: std::time::Duration::from_millis(2_000),
4440 ..PoolConfig::default()
4441 };
4442 let pool = Arc::new(ConnectionPool::new(config).unwrap());
4443 let writer = open_standalone_writer(&pool).unwrap();
4444 writer
4445 .execute_batch(
4446 "CREATE TABLE grace_probe(id INTEGER PRIMARY KEY, value TEXT NOT NULL); \
4447 INSERT INTO grace_probe(value) VALUES ('seed');",
4448 )
4449 .unwrap();
4450
4451 let mut reader = SqliteReader {
4452 handle: Some(open_cached_reader_handle(Arc::clone(&pool)).await.unwrap()),
4453 pool: Arc::clone(&pool),
4454 };
4455 let mut contender = SqliteReader {
4456 handle: Some(open_cached_reader_handle(Arc::clone(&pool)).await.unwrap()),
4457 pool: Arc::clone(&pool),
4458 };
4459 let udf_started = Arc::new(std::sync::atomic::AtomicBool::new(false));
4460 register_khive_test_slow_udf(
4461 &reader.handle.as_ref().unwrap().conn,
4462 900,
4463 Arc::clone(&udf_started),
4464 );
4465
4466 reader
4467 .query_all(SqlStatement {
4468 sql: "BEGIN DEFERRED".into(),
4469 params: vec![],
4470 label: None,
4471 })
4472 .await
4473 .expect("begin admitted transaction");
4474 reader
4475 .query_all(SqlStatement {
4476 sql: "SELECT * FROM grace_probe".into(),
4477 params: vec![],
4478 label: None,
4479 })
4480 .await
4481 .expect("materialize a real WAL snapshot");
4482 writer
4483 .execute_batch(
4484 "WITH RECURSIVE rows(value) AS (\
4485 SELECT 1 UNION ALL SELECT value + 1 FROM rows WHERE value < 100\
4486 ) INSERT INTO grace_probe(value) SELECT printf('row-%d', value) FROM rows;",
4487 )
4488 .unwrap();
4489 let (_, log_before, checkpointed_before) = passive_checkpoint(&writer);
4490 assert!(
4491 log_before > checkpointed_before,
4492 "the explicit reader snapshot must pin WAL frames before cancellation"
4493 );
4494
4495 let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
4496 let query = tokio::spawn(crate::scope_request_read_cancellation(
4497 cancel_rx,
4498 async move {
4499 let result = reader
4500 .query_all(SqlStatement {
4501 sql: "SELECT khive_test_slow_udf()".into(),
4502 params: vec![],
4503 label: None,
4504 })
4505 .await;
4506 (reader, result)
4507 },
4508 ));
4509 wait_for_flag(udf_started.as_ref()).await;
4515 cancel_tx.send(true).unwrap();
4516
4517 let (mut reader, result) = tokio::time::timeout(std::time::Duration::from_secs(3), query)
4518 .await
4519 .expect(
4520 "a worker that settles within the grace+hard-cap bound must not hang the caller",
4521 )
4522 .unwrap();
4523 assert!(
4524 matches!(result, Err(StorageError::Timeout { .. })),
4525 "request cancellation must still surface as a typed timeout even after grace \
4526 was exceeded; got {result:?}"
4527 );
4528
4529 assert_eq!(
4535 pool.sql_bridge_reader_slots().available_permits(),
4536 1,
4537 "the sole reader permit must be visible again once the bounded join completes"
4538 );
4539
4540 let (_, log_after, checkpointed_after) = passive_checkpoint(&writer);
4541 assert_eq!(
4542 log_after, checkpointed_after,
4543 "the abandoned explicit read transaction must release its WAL snapshot by the \
4544 time the caller observes the timeout"
4545 );
4546
4547 contender
4548 .query_all(SqlStatement {
4549 sql: "SELECT 1".into(),
4550 params: vec![],
4551 label: None,
4552 })
4553 .await
4554 .expect("a fresh reader must be admitted once the zombie worker has settled");
4555
4556 reader
4560 .query_all(SqlStatement {
4561 sql: "SELECT 1".into(),
4562 params: vec![],
4563 label: None,
4564 })
4565 .await
4566 .expect("the interrupted connection must remain usable after settling");
4567 }
4568
4569 #[tokio::test]
4570 #[serial_test::serial(tx_registry)]
4571 async fn cached_reader_transaction_lifecycle_survives_sqlite_empty_prefixes() {
4572 let dir = tempfile::tempdir().unwrap();
4573 let config = PoolConfig {
4574 path: Some(dir.path().join("sql_bridge_reader_prefixed_tx_control.db")),
4575 write_queue_enabled: Some(true),
4576 max_readers: 1,
4577 ..PoolConfig::default()
4578 };
4579 let pool = Arc::new(ConnectionPool::new(config).unwrap());
4580 let bridge = SqlBridge::new(Arc::clone(&pool), true);
4581 let mut reader = bridge.reader().await.unwrap();
4582
4583 reader
4584 .query_all(SqlStatement {
4585 sql: " ; -- empty statement\n /* leading comment */ \u{feff} BEGIN DEFERRED".into(),
4586 params: vec![],
4587 label: None,
4588 })
4589 .await
4590 .expect("prefixed BEGIN must enter the admitted transaction state");
4591 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 0);
4592 reader
4593 .query_all(SqlStatement {
4594 sql: " /* leading comment */ \u{feff} ; COMMIT".into(),
4595 params: vec![],
4596 label: None,
4597 })
4598 .await
4599 .expect("prefixed COMMIT must end the admitted transaction state");
4600 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
4601
4602 let rejected = reader
4603 .query_all(SqlStatement {
4604 sql: " ; /* no active transaction */ \u{feff} COMMIT".into(),
4605 params: vec![],
4606 label: None,
4607 })
4608 .await;
4609 assert!(
4610 matches!(
4611 &rejected,
4612 Err(StorageError::InvalidInput {
4613 operation,
4614 message,
4615 ..
4616 }) if operation.as_ref() == "query_all"
4617 && message.contains("transaction control")
4618 && message.contains("COMMIT")
4619 ),
4620 "a prefixed COMMIT without an admitted transaction must still fail closed; \
4621 got {rejected:?}"
4622 );
4623
4624 let mut queue_backed_writer = bridge.writer().await.unwrap();
4625 let rejected = queue_backed_writer
4626 .query_all(SqlStatement {
4627 sql: "-- leading comment\n \u{feff} ; /* empty */ SAVEPOINT pinned".into(),
4628 params: vec![],
4629 label: None,
4630 })
4631 .await;
4632 assert!(
4633 matches!(
4634 &rejected,
4635 Err(StorageError::InvalidInput {
4636 operation,
4637 message,
4638 ..
4639 }) if operation.as_ref() == "writer.query_all"
4640 && message.contains("transaction control")
4641 && message.contains("SAVEPOINT")
4642 ),
4643 "a queue-backed cached reader must classify transaction control through \
4644 comments, BOMs, and empty statements; got {rejected:?}"
4645 );
4646
4647 let value = reader
4648 .query_scalar(SqlStatement {
4649 sql: "SELECT 10".into(),
4650 params: vec![],
4651 label: None,
4652 })
4653 .await
4654 .expect("prefixed transaction lifecycle must preserve the cached reader");
4655 assert!(matches!(value, Some(SqlValue::Integer(10))));
4656 }
4657
4658 #[tokio::test]
4659 async fn cached_reader_restores_autocommit_before_releasing_its_operation_permit() {
4660 let dir = tempfile::tempdir().unwrap();
4661 let config = PoolConfig {
4662 path: Some(dir.path().join("sql_bridge_reader_autocommit.db")),
4663 max_readers: 1,
4664 ..PoolConfig::default()
4665 };
4666 let pool = Arc::new(ConnectionPool::new(config).unwrap());
4667 let conn = open_standalone_reader(&pool).unwrap();
4668 conn.execute_batch("BEGIN DEFERRED; SELECT * FROM sqlite_schema")
4669 .unwrap();
4670 assert!(
4671 !conn.is_autocommit(),
4672 "the regression precondition needs a live read transaction"
4673 );
4674 let mut reader = SqliteReader {
4675 handle: Some(StandaloneHandle {
4676 conn,
4677 _retained_slot: None,
4678 read_transaction_slot: None,
4679 }),
4680 pool: Arc::clone(&pool),
4681 };
4682
4683 let rejected = reader
4684 .query_all(SqlStatement {
4685 sql: "ROLLBACK".into(),
4689 params: vec![],
4690 label: None,
4691 })
4692 .await;
4693 assert!(
4694 matches!(
4695 &rejected,
4696 Err(StorageError::InvalidInput {
4697 operation,
4698 message,
4699 ..
4700 }) if operation.as_ref() == "query_all"
4701 && message.contains("outside autocommit")
4702 ),
4703 "a cached reader that reaches the boundary outside autocommit must fail closed; \
4704 got {rejected:?}"
4705 );
4706 assert_eq!(
4707 pool.sql_bridge_reader_slots().available_permits(),
4708 1,
4709 "the permit may be released only after the stale transaction is gone"
4710 );
4711 assert!(
4712 reader
4713 .handle
4714 .as_ref()
4715 .expect("successful rollback should preserve the cached handle")
4716 .conn
4717 .is_autocommit(),
4718 "the restored idle connection must not retain a WAL snapshot"
4719 );
4720
4721 let value = reader
4722 .query_scalar(SqlStatement {
4723 sql: "SELECT 9".into(),
4724 params: vec![],
4725 label: None,
4726 })
4727 .await
4728 .expect("the cleaned cached reader must remain usable");
4729 assert!(matches!(value, Some(SqlValue::Integer(9))));
4730 }
4731
4732 #[tokio::test]
4733 async fn standalone_writer_read_preserves_manual_atomic_transaction() {
4734 let dir = tempfile::tempdir().unwrap();
4735 let config = PoolConfig {
4736 path: Some(dir.path().join("sql_bridge_writer_atomic_read.db")),
4737 write_queue_enabled: Some(false),
4738 max_readers: 1,
4739 ..PoolConfig::default()
4740 };
4741 let pool = Arc::new(ConnectionPool::new(config).unwrap());
4742 {
4743 let writer = pool.writer().unwrap();
4744 writer
4745 .conn()
4746 .execute_batch(
4747 "CREATE TABLE atomic_read_test \
4748 (id INTEGER PRIMARY KEY, value TEXT NOT NULL)",
4749 )
4750 .unwrap();
4751 }
4752 let bridge = SqlBridge::new(Arc::clone(&pool), true);
4753
4754 let observed = bridge
4755 .atomic_unit(Box::new(|writer| {
4756 Box::pin(async move {
4757 writer
4758 .execute(SqlStatement {
4759 sql: "INSERT INTO atomic_read_test (id, value) VALUES (1, 'pending')"
4760 .into(),
4761 params: vec![],
4762 label: None,
4763 })
4764 .await?;
4765 let count = writer
4766 .query_scalar(SqlStatement {
4767 sql: "SELECT COUNT(*) FROM atomic_read_test".into(),
4768 params: vec![],
4769 label: None,
4770 })
4771 .await?;
4772 Ok(Box::new(count) as Box<dyn std::any::Any + Send>)
4773 })
4774 }))
4775 .await
4776 .expect("manual atomic read must not be mistaken for an idle reader snapshot");
4777 let observed = match observed.downcast::<Option<SqlValue>>() {
4778 Ok(observed) => observed,
4779 Err(_) => panic!("unexpected atomic result type"),
4780 };
4781 assert!(matches!(*observed, Some(SqlValue::Integer(1))));
4782
4783 let mut reader = bridge.reader().await.unwrap();
4784 let committed = reader
4785 .query_scalar(SqlStatement {
4786 sql: "SELECT COUNT(*) FROM atomic_read_test".into(),
4787 params: vec![],
4788 label: None,
4789 })
4790 .await
4791 .unwrap();
4792 assert!(matches!(committed, Some(SqlValue::Integer(1))));
4793 }
4794
4795 #[tokio::test]
4796 async fn request_cancellation_preserves_file_backed_manual_atomic_read_and_commit() {
4797 let dir = tempfile::tempdir().unwrap();
4798 let pool = Arc::new(
4799 ConnectionPool::new(PoolConfig {
4800 path: Some(dir.path().join("sql_bridge_writer_tx_cancel.db")),
4801 write_queue_enabled: Some(false),
4802 ..PoolConfig::default()
4803 })
4804 .unwrap(),
4805 );
4806 pool.writer()
4807 .unwrap()
4808 .conn()
4809 .execute_batch(
4810 "CREATE TABLE writer_tx_cancel_probe(\
4811 id INTEGER PRIMARY KEY, value TEXT NOT NULL)",
4812 )
4813 .unwrap();
4814 let bridge = SqlBridge::new(Arc::clone(&pool), true);
4815 let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
4816
4817 let observed = crate::scope_request_read_cancellation(
4818 cancel_rx,
4819 bridge.atomic_unit(Box::new(move |writer| {
4820 Box::pin(async move {
4821 writer
4822 .execute(SqlStatement {
4823 sql: "INSERT INTO writer_tx_cancel_probe VALUES (1, 'before')".into(),
4824 params: vec![],
4825 label: None,
4826 })
4827 .await?;
4828 cancel_tx.send(true).unwrap();
4829 let count = writer
4830 .query_scalar(SqlStatement {
4831 sql: "SELECT COUNT(*) FROM writer_tx_cancel_probe".into(),
4832 params: vec![],
4833 label: None,
4834 })
4835 .await?;
4836 writer
4837 .execute(SqlStatement {
4838 sql: "INSERT INTO writer_tx_cancel_probe VALUES (2, 'after')".into(),
4839 params: vec![],
4840 label: None,
4841 })
4842 .await?;
4843 Ok(Box::new(count) as Box<dyn std::any::Any + Send>)
4844 })
4845 })),
4846 )
4847 .await
4848 .expect("request cancellation must not interrupt an admitted manual write transaction");
4849 let observed = match observed.downcast::<Option<SqlValue>>() {
4850 Ok(observed) => observed,
4851 Err(_) => panic!("unexpected atomic result type"),
4852 };
4853 assert!(matches!(*observed, Some(SqlValue::Integer(1))));
4854
4855 let reader = pool.reader().unwrap();
4856 let rows: i64 = reader
4857 .conn()
4858 .query_row("SELECT COUNT(*) FROM writer_tx_cancel_probe", [], |row| {
4859 row.get(0)
4860 })
4861 .unwrap();
4862 assert_eq!(rows, 2, "both writes around the SELECT must commit");
4863 }
4864
4865 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4866 async fn cancelled_standalone_writer_transaction_retains_active_reader_admission() {
4867 let dir = tempfile::tempdir().unwrap();
4868 let pool = Arc::new(
4869 ConnectionPool::new(PoolConfig {
4870 path: Some(dir.path().join("sql_bridge_writer_tx_admission.db")),
4871 write_queue_enabled: Some(false),
4872 max_readers: 1,
4873 checkout_timeout: std::time::Duration::from_millis(250),
4874 ..PoolConfig::default()
4875 })
4876 .unwrap(),
4877 );
4878 let writer_slot = pool
4879 .sql_bridge_writer_slots()
4880 .acquire_owned()
4881 .await
4882 .unwrap();
4883 let conn = open_standalone_writer(&pool).unwrap();
4884 let (entered, release, _completed) = blocking_non_interrupting_progress_gate(&conn);
4885 let mut writer = SqliteWriter {
4886 handle: Some(StandaloneHandle {
4887 conn,
4888 _retained_slot: Some(writer_slot),
4889 read_transaction_slot: None,
4890 }),
4891 writer_task: None,
4892 origin: pool.origin(),
4893 db: crate::timeout_sink::db_label(&pool),
4894 pool: Arc::clone(&pool),
4895 };
4896 khive_storage::SqlWriter::execute(
4897 &mut writer,
4898 SqlStatement {
4899 sql: "BEGIN IMMEDIATE".into(),
4900 params: vec![],
4901 label: None,
4902 },
4903 )
4904 .await
4905 .unwrap();
4906 let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
4907 cancel_tx.send(true).unwrap();
4908
4909 let query = tokio::spawn(crate::scope_request_read_cancellation(
4910 cancel_rx,
4911 async move {
4912 let result =
4913 khive_storage::SqlReader::query_all(&mut writer, progress_gate_statement())
4914 .await;
4915 let rollback = khive_storage::SqlWriter::execute(
4916 &mut writer,
4917 SqlStatement {
4918 sql: "ROLLBACK".into(),
4919 params: vec![],
4920 label: None,
4921 },
4922 )
4923 .await;
4924 (result, rollback)
4925 },
4926 ));
4927 tokio::time::timeout(std::time::Duration::from_secs(1), entered.notified())
4928 .await
4929 .expect("cancelled writer-transaction SELECT never reached SQLite");
4930 assert_eq!(
4931 pool.sql_bridge_reader_slots().available_permits(),
4932 0,
4933 "a writer-supertrait SELECT must retain ordinary active-reader admission"
4934 );
4935
4936 tokio::task::spawn_blocking(move || release.wait())
4937 .await
4938 .unwrap();
4939 let (rows, rollback) = tokio::time::timeout(std::time::Duration::from_secs(2), query)
4940 .await
4941 .expect("writer-transaction SELECT did not finish after its gate opened")
4942 .unwrap();
4943 assert_eq!(
4944 rows.expect("request cancellation interrupted the admitted writer transaction")
4945 .len(),
4946 1
4947 );
4948 rollback.expect("writer transaction did not return to autocommit");
4949 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
4950 }
4951
4952 #[tokio::test]
4953 async fn expired_deadline_preserves_pool_backed_manual_atomic_read_and_commit() {
4954 let pool = Arc::new(ConnectionPool::new(PoolConfig::default()).unwrap());
4955 pool.writer()
4956 .unwrap()
4957 .conn()
4958 .execute_batch(
4959 "CREATE TABLE pool_writer_tx_deadline_probe(\
4960 id INTEGER PRIMARY KEY, value TEXT NOT NULL)",
4961 )
4962 .unwrap();
4963 let bridge = SqlBridge::new(Arc::clone(&pool), false);
4964
4965 let observed = crate::scope_request_read_deadline(
4966 std::time::Duration::ZERO,
4967 bridge.atomic_unit(Box::new(|writer| {
4968 Box::pin(async move {
4969 writer
4970 .execute(SqlStatement {
4971 sql: "INSERT INTO pool_writer_tx_deadline_probe VALUES (1, 'before')"
4972 .into(),
4973 params: vec![],
4974 label: None,
4975 })
4976 .await?;
4977 let count = writer
4978 .query_scalar(SqlStatement {
4979 sql: "SELECT COUNT(*) FROM pool_writer_tx_deadline_probe".into(),
4980 params: vec![],
4981 label: None,
4982 })
4983 .await?;
4984 writer
4985 .execute(SqlStatement {
4986 sql: "INSERT INTO pool_writer_tx_deadline_probe VALUES (2, 'after')"
4987 .into(),
4988 params: vec![],
4989 label: None,
4990 })
4991 .await?;
4992 Ok(Box::new(count) as Box<dyn std::any::Any + Send>)
4993 })
4994 })),
4995 )
4996 .await
4997 .expect("an expired read deadline must not interrupt an admitted manual write transaction");
4998 let observed = match observed.downcast::<Option<SqlValue>>() {
4999 Ok(observed) => observed,
5000 Err(_) => panic!("unexpected atomic result type"),
5001 };
5002 assert!(matches!(*observed, Some(SqlValue::Integer(1))));
5003
5004 let reader = pool.reader().unwrap();
5005 let rows: i64 = reader
5006 .conn()
5007 .query_row(
5008 "SELECT COUNT(*) FROM pool_writer_tx_deadline_probe",
5009 [],
5010 |row| row.get(0),
5011 )
5012 .unwrap();
5013 assert_eq!(rows, 2, "both writes around the SELECT must commit");
5014 }
5015
5016 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5017 async fn cancelled_standalone_open_retains_slot_until_open_finishes() {
5018 let dir = tempfile::tempdir().unwrap();
5019 let config = PoolConfig {
5020 path: Some(dir.path().join("sql_bridge_cancelled_open.db")),
5021 max_readers: 1,
5022 checkout_timeout: std::time::Duration::from_millis(250),
5023 ..PoolConfig::default()
5024 };
5025 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5026 let slots = pool.sql_bridge_reader_slots();
5027 let slot = Arc::clone(&slots).acquire_owned().await.unwrap();
5028 assert_eq!(slots.available_permits(), 0);
5029
5030 let (entered_tx, entered_rx) = std::sync::mpsc::channel();
5031 let (release_tx, release_rx) = std::sync::mpsc::channel();
5032 let open = tokio::spawn(open_standalone_on_blocking(
5033 Arc::clone(&pool),
5034 slot,
5035 "test_open_reader",
5036 move |pool| {
5037 entered_tx.send(()).unwrap();
5038 release_rx.recv().unwrap();
5039 open_standalone_reader(pool)
5040 },
5041 ));
5042 tokio::task::spawn_blocking(move || entered_rx.recv())
5043 .await
5044 .unwrap()
5045 .unwrap();
5046
5047 open.abort();
5048 assert!(matches!(open.await, Err(error) if error.is_cancelled()));
5049 assert_eq!(
5050 slots.available_permits(),
5051 0,
5052 "the permit must remain in the detached open closure"
5053 );
5054 let contender = tokio::time::timeout(
5055 std::time::Duration::from_millis(50),
5056 Arc::clone(&slots).acquire_owned(),
5057 )
5058 .await;
5059 assert!(contender.is_err(), "an in-flight open must retain the cap");
5060
5061 release_tx.send(()).unwrap();
5062 let recovered = tokio::time::timeout(
5063 std::time::Duration::from_secs(1),
5064 Arc::clone(&slots).acquire_owned(),
5065 )
5066 .await
5067 .expect("the detached open did not release its permit")
5068 .unwrap();
5069 assert_eq!(slots.available_permits(), 0);
5070 drop(recovered);
5071 assert_eq!(slots.available_permits(), 1);
5072 }
5073
5074 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5075 async fn abandoned_writer_read_interrupts_and_releases_writer_handle() {
5076 let dir = tempfile::tempdir().unwrap();
5077 let config = PoolConfig {
5078 path: Some(dir.path().join("sql_bridge_cancelled_writer.db")),
5079 write_queue_enabled: Some(false),
5080 checkout_timeout: std::time::Duration::from_millis(250),
5081 ..PoolConfig::default()
5082 };
5083 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5084 let bridge = SqlBridge::new(Arc::clone(&pool), true);
5085
5086 let handle_slot = pool
5087 .sql_bridge_writer_slots()
5088 .acquire_owned()
5089 .await
5090 .unwrap();
5091 let conn = open_standalone_writer(&pool).unwrap();
5092 let mut writer = SqliteWriter {
5093 handle: Some(StandaloneHandle {
5094 conn,
5095 _retained_slot: Some(handle_slot),
5096 read_transaction_slot: None,
5097 }),
5098 writer_task: None,
5099 origin: pool.origin(),
5100 db: crate::timeout_sink::db_label(&pool),
5101 pool: Arc::clone(&pool),
5102 };
5103 let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
5104 let query = tokio::spawn(crate::scope_test_read_progress(
5105 Arc::clone(&progress),
5106 async move {
5107 khive_storage::SqlReader::query_all(&mut writer, deliberately_slow_read_statement())
5108 .await
5109 },
5110 ));
5111
5112 wait_for_progress(progress.as_ref()).await;
5113 query.abort();
5114 assert!(matches!(query.await, Err(error) if error.is_cancelled()));
5115 let writer_after =
5116 tokio::time::timeout(std::time::Duration::from_millis(500), bridge.writer())
5117 .await
5118 .expect("abandoned SQLite read did not release the writer handle promptly")
5119 .expect("writer handle remained unavailable after read interruption");
5120 drop(writer_after);
5121 let stopped_at = progress.load(std::sync::atomic::Ordering::SeqCst);
5122 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
5123 assert_eq!(
5124 progress.load(std::sync::atomic::Ordering::SeqCst),
5125 stopped_at,
5126 "writer-backed SQLite read kept consuming work after cancellation"
5127 );
5128 }
5129
5130 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5131 async fn request_cancellation_never_interrupts_admitted_execute_batch() {
5132 let dir = tempfile::tempdir().unwrap();
5133 let config = PoolConfig {
5134 path: Some(dir.path().join("sql_bridge_cancelled_writer_batch.db")),
5135 write_queue_enabled: Some(false),
5136 checkout_timeout: std::time::Duration::from_millis(250),
5137 ..PoolConfig::default()
5138 };
5139 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5140 let bridge = SqlBridge::new(Arc::clone(&pool), true);
5141 {
5142 let guard = pool.writer().unwrap();
5143 guard
5144 .conn()
5145 .execute_batch(
5146 "CREATE TABLE cancellation_write_probe(\
5147 id INTEGER PRIMARY KEY, value INTEGER NOT NULL)",
5148 )
5149 .unwrap();
5150 }
5151
5152 let handle_slot = pool
5153 .sql_bridge_writer_slots()
5154 .acquire_owned()
5155 .await
5156 .unwrap();
5157 let conn = open_standalone_writer(&pool).unwrap();
5158 let (entered, release, completed) = blocking_non_interrupting_progress_gate(&conn);
5159 let mut writer = SqliteWriter {
5160 handle: Some(StandaloneHandle {
5161 conn,
5162 _retained_slot: Some(handle_slot),
5163 read_transaction_slot: None,
5164 }),
5165 writer_task: None,
5166 origin: pool.origin(),
5167 db: crate::timeout_sink::db_label(&pool),
5168 pool: Arc::clone(&pool),
5169 };
5170 let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
5171 let query = tokio::spawn(crate::scope_request_read_cancellation(
5172 cancel_rx,
5173 async move {
5174 khive_storage::SqlWriter::execute_batch(&mut writer, vec![slow_insert_statement()])
5175 .await
5176 },
5177 ));
5178
5179 tokio::time::timeout(std::time::Duration::from_secs(1), entered.notified())
5180 .await
5181 .expect("mutating execute_batch never reached SQLite VM work");
5182 cancel_tx.send(true).unwrap();
5183 tokio::time::sleep(std::time::Duration::from_millis(25)).await;
5184 assert!(
5185 !query.is_finished(),
5186 "request-read cancellation must not interrupt an admitted batch"
5187 );
5188
5189 let contender = bridge.writer().await;
5190 let retained_slot = matches!(
5191 &contender,
5192 Err(StorageError::AdmissionTimeout { operation, .. })
5193 if operation.as_ref() == "sql_bridge.writer_handle"
5194 );
5195 drop(contender);
5196
5197 tokio::task::spawn_blocking(move || release.wait())
5198 .await
5199 .unwrap();
5200 let affected = tokio::time::timeout(std::time::Duration::from_secs(2), query)
5201 .await
5202 .expect("admitted batch did not finish after its gate was released")
5203 .unwrap()
5204 .expect("request cancellation must preserve the batch result");
5205 assert_eq!(affected, 10_000);
5206 tokio::time::timeout(std::time::Duration::from_secs(1), completed.notified())
5207 .await
5208 .expect("completed batch did not release its connection");
5209 assert!(
5210 retained_slot,
5211 "request cancellation released the writer slot before the admitted batch stopped"
5212 );
5213 let reader = pool.reader().unwrap();
5214 let count: i64 = reader
5215 .conn()
5216 .query_row("SELECT COUNT(*) FROM cancellation_write_probe", [], |row| {
5217 row.get(0)
5218 })
5219 .unwrap();
5220 assert_eq!(count, 10_000, "the admitted batch must commit every row");
5221 }
5222
5223 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5224 async fn request_cancellation_never_interrupts_dml_returning_via_sql_reader() {
5225 let dir = tempfile::tempdir().unwrap();
5226 let config = PoolConfig {
5227 path: Some(dir.path().join("sql_bridge_dml_returning_cancel.db")),
5228 write_queue_enabled: Some(false),
5229 checkout_timeout: std::time::Duration::from_millis(250),
5230 ..PoolConfig::default()
5231 };
5232 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5233 {
5234 let guard = pool.writer().unwrap();
5235 guard
5236 .conn()
5237 .execute_batch(
5238 "CREATE TABLE returning_write_probe(\
5239 id INTEGER PRIMARY KEY, value INTEGER NOT NULL)",
5240 )
5241 .unwrap();
5242 }
5243
5244 let handle_slot = pool
5245 .sql_bridge_writer_slots()
5246 .acquire_owned()
5247 .await
5248 .unwrap();
5249 let conn = open_standalone_writer(&pool).unwrap();
5250 let (entered, release, completed) = blocking_non_interrupting_progress_gate(&conn);
5251 let mut writer = SqliteWriter {
5252 handle: Some(StandaloneHandle {
5253 conn,
5254 _retained_slot: Some(handle_slot),
5255 read_transaction_slot: None,
5256 }),
5257 writer_task: None,
5258 origin: pool.origin(),
5259 db: crate::timeout_sink::db_label(&pool),
5260 pool: Arc::clone(&pool),
5261 };
5262 let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
5263 let query = tokio::spawn(crate::scope_request_read_cancellation(
5264 cancel_rx,
5265 async move {
5266 khive_storage::SqlReader::query_all(
5267 &mut writer,
5268 SqlStatement {
5269 sql: "INSERT INTO returning_write_probe(value) \
5270 WITH RECURSIVE rows(value) AS (\
5271 SELECT 1 UNION ALL SELECT value + 1 FROM rows WHERE value < 10000\
5272 ) SELECT value FROM rows RETURNING id"
5273 .into(),
5274 params: vec![],
5275 label: Some("non-interruptible-returning-probe".into()),
5276 },
5277 )
5278 .await
5279 },
5280 ));
5281
5282 tokio::time::timeout(std::time::Duration::from_secs(1), entered.notified())
5283 .await
5284 .expect("DML RETURNING never reached admitted SQLite work");
5285 cancel_tx.send(true).unwrap();
5286 tokio::time::sleep(std::time::Duration::from_millis(25)).await;
5287 assert!(
5288 !query.is_finished(),
5289 "request-read cancellation interrupted DML RETURNING"
5290 );
5291
5292 tokio::task::spawn_blocking(move || release.wait())
5293 .await
5294 .unwrap();
5295 let rows = tokio::time::timeout(std::time::Duration::from_secs(2), query)
5296 .await
5297 .expect("DML RETURNING did not finish after its gate was released")
5298 .unwrap()
5299 .expect("request cancellation must preserve DML RETURNING's result");
5300 assert_eq!(rows.len(), 10_000);
5301 tokio::time::timeout(std::time::Duration::from_secs(1), completed.notified())
5302 .await
5303 .expect("completed DML RETURNING did not release its connection");
5304
5305 let reader = pool.reader().unwrap();
5306 let count: i64 = reader
5307 .conn()
5308 .query_row("SELECT COUNT(*) FROM returning_write_probe", [], |row| {
5309 row.get(0)
5310 })
5311 .unwrap();
5312 assert_eq!(count, 10_000, "DML RETURNING must commit every row");
5313 }
5314
5315 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5323 async fn cancelled_call_invalidates_handle_reuse_fails_loud() {
5324 let dir = tempfile::tempdir().unwrap();
5325 let config = PoolConfig {
5326 path: Some(dir.path().join("sql_bridge_cancelled_reuse.db")),
5327 checkout_timeout: std::time::Duration::from_millis(250),
5328 ..PoolConfig::default()
5329 };
5330 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5331
5332 let handle_slot = acquire_handle_slot(
5333 pool.sql_bridge_writer_slots(),
5334 pool.config().checkout_timeout,
5335 "sql_bridge.writer_handle",
5336 SlotTimeoutClass::Admission,
5337 )
5338 .await
5339 .unwrap();
5340 let conn = open_standalone_writer(&pool).unwrap();
5341 let (entered, release, completed) = blocking_non_interrupting_progress_gate(&conn);
5342 let writer = Arc::new(tokio::sync::Mutex::new(SqliteWriter {
5343 handle: Some(StandaloneHandle {
5344 conn,
5345 _retained_slot: Some(handle_slot),
5346 read_transaction_slot: None,
5347 }),
5348 writer_task: None,
5349 origin: pool.origin(),
5350 db: crate::timeout_sink::db_label(&pool),
5351 pool: Arc::clone(&pool),
5352 }));
5353 let writer_clone = Arc::clone(&writer);
5354 let query = tokio::spawn(async move {
5355 khive_storage::SqlWriter::execute_batch(
5356 &mut *writer_clone.lock().await,
5357 vec![progress_gate_statement()],
5358 )
5359 .await
5360 });
5361
5362 entered.notified().await;
5363 query.abort();
5364 let cancelled = matches!(query.await, Err(error) if error.is_cancelled());
5365
5366 let reuse = khive_storage::SqlWriter::execute(
5367 &mut *writer.lock().await,
5368 SqlStatement {
5369 sql: "CREATE TABLE cancelled_reuse_probe (id INTEGER PRIMARY KEY)".into(),
5370 params: vec![],
5371 label: None,
5372 },
5373 )
5374 .await;
5375 let message = match reuse {
5376 Err(StorageError::Pool { message, .. }) => message,
5377 other => panic!(
5378 "reusing a cancelled writer handle must fail loudly with \
5379 'connection already consumed'; got {other:?}"
5380 ),
5381 };
5382 assert!(
5383 message.contains("connection already consumed"),
5384 "expected the cancelled handle's reuse error to name the pinned \
5385 failure; got {message:?}"
5386 );
5387
5388 tokio::task::spawn_blocking(move || release.wait())
5389 .await
5390 .unwrap();
5391 tokio::time::timeout(std::time::Duration::from_secs(1), completed.notified())
5392 .await
5393 .expect("cancelled writer's detached SQLite call did not finish");
5394 assert!(cancelled, "writer batch task did not report cancellation");
5395 }
5396
5397 #[tokio::test]
5406 async fn execute_batch_rejects_transaction_control_before_executing_anything() {
5407 let dir = tempfile::tempdir().unwrap();
5408 let config = PoolConfig {
5409 path: Some(dir.path().join("sql_bridge_tx_control_reject.db")),
5410 checkout_timeout: std::time::Duration::from_millis(250),
5411 write_queue_enabled: Some(false),
5412 ..PoolConfig::default()
5413 };
5414 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5415 {
5416 let guard = pool.writer().unwrap();
5417 guard
5418 .conn()
5419 .execute_batch(
5420 "CREATE TABLE tx_reject_test (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
5421 )
5422 .unwrap();
5423 }
5424
5425 let handle_slot = acquire_handle_slot(
5426 pool.sql_bridge_writer_slots(),
5427 pool.config().checkout_timeout,
5428 "sql_bridge.writer_handle",
5429 SlotTimeoutClass::Admission,
5430 )
5431 .await
5432 .unwrap();
5433 let conn = open_standalone_writer(&pool).unwrap();
5434 let mut writer = SqliteWriter {
5435 handle: Some(StandaloneHandle {
5436 conn,
5437 _retained_slot: Some(handle_slot),
5438 read_transaction_slot: None,
5439 }),
5440 writer_task: None,
5441 origin: pool.origin(),
5442 db: crate::timeout_sink::db_label(&pool),
5443 pool: Arc::clone(&pool),
5444 };
5445
5446 for tail in ["COMMIT", "BEGIN"] {
5447 let multi = khive_storage::SqlWriter::execute_batch(
5448 &mut writer,
5449 vec![SqlStatement {
5450 sql: format!(
5451 "INSERT INTO tx_reject_test (id, val) VALUES (10, 'tail'); {tail}"
5452 ),
5453 params: vec![],
5454 label: None,
5455 }],
5456 )
5457 .await;
5458 let message = multi
5459 .as_ref()
5460 .err()
5461 .map(ToString::to_string)
5462 .unwrap_or_default();
5463 assert!(
5464 message.contains("Multiple statements"),
5465 "a SqlStatement with trailing {tail} must be rejected before execution; got {message}"
5466 );
5467 }
5468
5469 let batch = khive_storage::SqlWriter::execute_batch(
5472 &mut writer,
5473 vec![
5474 SqlStatement {
5475 sql: "INSERT INTO tx_reject_test (id, val) VALUES (1, 'a')".into(),
5476 params: vec![],
5477 label: None,
5478 },
5479 SqlStatement {
5480 sql: "COMMIT".into(),
5481 params: vec![],
5482 label: None,
5483 },
5484 ],
5485 )
5486 .await;
5487 match &batch {
5488 Err(StorageError::InvalidInput {
5489 operation, message, ..
5490 }) => {
5491 assert_eq!(operation.as_ref(), "execute_batch");
5492 assert!(
5493 message.contains("transaction control") && message.contains("COMMIT"),
5494 "the rejection must name the offending statement head; got {message:?}"
5495 );
5496 }
5497 other => {
5498 panic!("a batch containing a bare COMMIT must be rejected up front; got {other:?}")
5499 }
5500 }
5501
5502 for sql in [
5505 "BEGIN IMMEDIATE",
5506 "START TRANSACTION",
5507 "commit",
5508 "End transaction",
5509 "ROLLBACK",
5510 "SAVEPOINT sp1",
5511 "RELEASE sp1",
5512 " -- leading comment\nCOMMIT",
5513 "/* block */ rollback to savepoint sp1",
5514 ] {
5515 let rejected = khive_storage::SqlWriter::execute_batch(
5516 &mut writer,
5517 vec![SqlStatement {
5518 sql: sql.into(),
5519 params: vec![],
5520 label: None,
5521 }],
5522 )
5523 .await;
5524 assert!(
5525 matches!(&rejected, Err(StorageError::InvalidInput { .. })),
5526 "transaction-control head {sql:?} must be rejected; got {rejected:?}"
5527 );
5528 }
5529
5530 let count: i64 = {
5534 let guard = pool.reader().unwrap();
5535 guard
5536 .conn()
5537 .query_row("SELECT COUNT(*) FROM tx_reject_test", [], |r| r.get(0))
5538 .unwrap()
5539 };
5540 assert_eq!(count, 0, "a rejected batch must not have executed anything");
5541
5542 let affected = khive_storage::SqlWriter::execute(
5543 &mut writer,
5544 SqlStatement {
5545 sql: "INSERT INTO tx_reject_test (id, val) VALUES (2, 'b')".into(),
5546 params: vec![],
5547 label: None,
5548 },
5549 )
5550 .await
5551 .expect("the handle must survive a rejected batch untouched");
5552 assert_eq!(affected, 1);
5553 }
5554
5555 #[tokio::test]
5556 async fn standalone_execute_batch_rejects_prefixed_commit_before_any_write() {
5557 let dir = tempfile::tempdir().unwrap();
5558 let config = PoolConfig {
5559 path: Some(dir.path().join("sql_bridge_prefixed_commit_standalone.db")),
5560 write_queue_enabled: Some(false),
5561 ..PoolConfig::default()
5562 };
5563 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5564 {
5565 let guard = pool.writer().unwrap();
5566 guard
5567 .conn()
5568 .execute_batch("CREATE TABLE prefixed_commit (id INTEGER PRIMARY KEY)")
5569 .unwrap();
5570 }
5571 let bridge = SqlBridge::new(Arc::clone(&pool), true);
5572 let mut writer = bridge.writer().await.unwrap();
5573
5574 let rejected = writer
5575 .execute_batch(vec![
5576 SqlStatement {
5577 sql: "INSERT INTO prefixed_commit (id) VALUES (1)".into(),
5578 params: vec![],
5579 label: None,
5580 },
5581 SqlStatement {
5582 sql: " ; -- empty statement\n /* leading comment */ \u{feff} ; COMMIT".into(),
5583 params: vec![],
5584 label: None,
5585 },
5586 ])
5587 .await;
5588 assert!(
5589 matches!(
5590 &rejected,
5591 Err(StorageError::InvalidInput {
5592 operation,
5593 message,
5594 ..
5595 }) if operation.as_ref() == "execute_batch"
5596 && message.contains("transaction control")
5597 && message.contains("COMMIT")
5598 ),
5599 "standalone execute_batch must reject a prefixed COMMIT before the INSERT; \
5600 got {rejected:?}"
5601 );
5602
5603 let mut reader = bridge.reader().await.unwrap();
5604 let count = reader
5605 .query_scalar(SqlStatement {
5606 sql: "SELECT COUNT(*) FROM prefixed_commit".into(),
5607 params: vec![],
5608 label: None,
5609 })
5610 .await
5611 .unwrap();
5612 assert!(
5613 matches!(count, Some(SqlValue::Integer(0))),
5614 "prefixed COMMIT rejection must happen before the earlier INSERT; got {count:?}"
5615 );
5616
5617 let affected = writer
5618 .execute(SqlStatement {
5619 sql: "INSERT INTO prefixed_commit (id) VALUES (2)".into(),
5620 params: vec![],
5621 label: None,
5622 })
5623 .await
5624 .expect("prefixed COMMIT rejection must leave the standalone handle reusable");
5625 assert_eq!(affected, 1);
5626 }
5627
5628 #[tokio::test]
5629 async fn execute_batch_rejects_multi_statement_on_pool_backed_path() {
5630 let pool = Arc::new(ConnectionPool::new(PoolConfig::default()).unwrap());
5631 pool.writer()
5632 .unwrap()
5633 .conn()
5634 .execute_batch(
5635 "CREATE TABLE multi_statement_pool_test (id INTEGER PRIMARY KEY, val TEXT)",
5636 )
5637 .unwrap();
5638 let bridge = SqlBridge::new(Arc::clone(&pool), false);
5639 let mut writer = bridge.writer().await.unwrap();
5640
5641 let result = khive_storage::SqlWriter::execute_batch(
5642 &mut *writer,
5643 vec![SqlStatement {
5644 sql: "INSERT INTO multi_statement_pool_test (id, val) VALUES (1, 'x'); COMMIT"
5645 .into(),
5646 params: vec![],
5647 label: None,
5648 }],
5649 )
5650 .await;
5651 let message = result
5652 .as_ref()
5653 .err()
5654 .map(ToString::to_string)
5655 .unwrap_or_default();
5656 assert!(
5657 message.contains("Multiple statements"),
5658 "pool-backed execute_batch must reject a trailing COMMIT; got {message}"
5659 );
5660 let count: i64 = pool
5661 .reader()
5662 .unwrap()
5663 .conn()
5664 .query_row(
5665 "SELECT COUNT(*) FROM multi_statement_pool_test",
5666 [],
5667 |row| row.get(0),
5668 )
5669 .unwrap();
5670 assert_eq!(count, 0);
5671 }
5672
5673 #[tokio::test]
5674 async fn inline_execute_batch_rejects_multi_statement_sql() {
5675 let dir = tempfile::tempdir().unwrap();
5676 let pool = Arc::new(
5677 ConnectionPool::new(PoolConfig {
5678 path: Some(dir.path().join("sql_bridge_multi_statement_inline.db")),
5679 write_queue_enabled: Some(true),
5680 write_routing_strict: true,
5681 ..PoolConfig::default()
5682 })
5683 .unwrap(),
5684 );
5685 pool.writer()
5686 .unwrap()
5687 .conn()
5688 .execute_batch(
5689 "CREATE TABLE multi_statement_inline_test (id INTEGER PRIMARY KEY, val TEXT)",
5690 )
5691 .unwrap();
5692 let bridge = SqlBridge::new(Arc::clone(&pool), true);
5693
5694 let result = bridge
5695 .atomic_unit(Box::new(|writer| {
5696 Box::pin(async move {
5697 writer
5698 .execute_batch(vec![SqlStatement {
5699 sql: "INSERT INTO multi_statement_inline_test (id, val) VALUES (1, 'x'); BEGIN"
5700 .into(),
5701 params: vec![],
5702 label: None,
5703 }])
5704 .await
5705 .map(|_| Box::new(()) as Box<dyn Any + Send>)
5706 })
5707 }))
5708 .await;
5709 let message = result
5710 .as_ref()
5711 .err()
5712 .map(ToString::to_string)
5713 .unwrap_or_default();
5714 assert!(
5715 message.contains("Multiple statements"),
5716 "InlineWriter must reject a trailing BEGIN; got {message}"
5717 );
5718 let count: i64 = pool
5719 .reader()
5720 .unwrap()
5721 .conn()
5722 .query_row(
5723 "SELECT COUNT(*) FROM multi_statement_inline_test",
5724 [],
5725 |row| row.get(0),
5726 )
5727 .unwrap();
5728 assert_eq!(count, 0);
5729 }
5730
5731 #[test]
5736 fn transaction_control_head_classification_matrix() {
5737 for (sql, expected) in [
5738 ("BEGIN", Some("BEGIN")),
5739 ("begin immediate", Some("BEGIN")),
5740 ("START TRANSACTION", Some("START")),
5741 ("start transaction", Some("START")),
5742 ("COMMIT", Some("COMMIT")),
5743 ("commit;", Some("COMMIT")),
5744 ("END", Some("END")),
5745 ("end transaction", Some("END")),
5746 ("ROLLBACK", Some("ROLLBACK")),
5747 ("rollback to savepoint sp1", Some("ROLLBACK")),
5748 ("SAVEPOINT sp1", Some("SAVEPOINT")),
5749 ("RELEASE sp1", Some("RELEASE")),
5750 ("release savepoint sp1", Some("RELEASE")),
5751 (" \t COMMIT", Some("COMMIT")),
5752 ("\u{feff}BEGIN", Some("BEGIN")),
5753 (" ; BEGIN", Some("BEGIN")),
5754 (" ; ; -- empty\n /* comment */ COMMIT", Some("COMMIT")),
5755 (" \u{feff} SAVEPOINT sp1", Some("SAVEPOINT")),
5756 ("/* comment */ \u{feff} ; RELEASE sp1", Some("RELEASE")),
5757 ("\u{feff} ; \u{feff} -- empty\n ROLLBACK", Some("ROLLBACK")),
5758 ("-- a comment\nCOMMIT", Some("COMMIT")),
5759 ("/* /* nested? no */ */ COMMIT", None),
5762 ("-- one\n-- two\n /* x */ begin", Some("BEGIN")),
5763 ("INSERT INTO t VALUES (1)", None),
5764 ("UPDATE t SET x = 1", None),
5765 ("DELETE FROM t", None),
5766 ("SELECT * FROM commit_log", None),
5767 ("CREATE TABLE rollback_audit (id INTEGER)", None),
5768 ("/* comment only */", None),
5769 (" ; /* empty statements only */ ; ", None),
5770 (" ; SELECT 1", None),
5771 ("", None),
5772 ] {
5773 assert_eq!(
5774 transaction_control_head(sql),
5775 expected,
5776 "classification mismatch for {sql:?}"
5777 );
5778 }
5779 }
5780
5781 #[test]
5782 fn cached_read_transaction_control_classification_matrix() {
5783 use CachedReadTransactionControl::{BeginDeferred, Finish, Unsupported};
5784
5785 for (sql, expected) in [
5786 ("BEGIN", Some(BeginDeferred)),
5787 ("begin transaction", Some(BeginDeferred)),
5788 (
5789 "/* p */ \u{feff} ; BEGIN /* mode */ DEFERRED",
5790 Some(BeginDeferred),
5791 ),
5792 ("BEGIN DEFERRED TRANSACTION", Some(BeginDeferred)),
5793 ("BEGIN IMMEDIATE", Some(Unsupported("BEGIN"))),
5794 ("BEGIN /* lock */ EXCLUSIVE", Some(Unsupported("BEGIN"))),
5795 ("BEGIN TRANSACTION IMMEDIATE", Some(Unsupported("BEGIN"))),
5796 ("begin transaction exclusive", Some(Unsupported("BEGIN"))),
5797 ("BEGIN TRANSACTION DEFERRED", Some(Unsupported("BEGIN"))),
5798 ("BEGIN IMMEDIATE TRANSACTION", Some(Unsupported("BEGIN"))),
5799 ("BEGIN TRANSACTION named_txn", Some(Unsupported("BEGIN"))),
5800 (
5801 "BEGIN DEFERRED TRANSACTION trailing",
5802 Some(Unsupported("BEGIN")),
5803 ),
5804 ("BEGIN DEFERRED DEFERRED", Some(Unsupported("BEGIN"))),
5805 (
5809 "BEGIN TRANSACTION \"IMMEDIATE\"",
5810 Some(Unsupported("BEGIN")),
5811 ),
5812 ("BEGIN TRANSACTION [IMMEDIATE]", Some(Unsupported("BEGIN"))),
5813 ("BEGIN TRANSACTION `IMMEDIATE`", Some(Unsupported("BEGIN"))),
5814 ("BEGIN TRANSACTION 'IMMEDIATE'", Some(Unsupported("BEGIN"))),
5815 ("BEGIN \"DEFERRED\"", Some(Unsupported("BEGIN"))),
5816 ("BEGIN; COMMIT", Some(Unsupported("BEGIN"))),
5817 ("BEGIN;", Some(BeginDeferred)),
5819 ("BEGIN DEFERRED ; -- done", Some(BeginDeferred)),
5820 ("BEGIN TRANSACTION /* t */ ;;", Some(BeginDeferred)),
5821 ("START TRANSACTION", Some(Unsupported("START"))),
5822 ("COMMIT", Some(Finish("COMMIT"))),
5823 ("END TRANSACTION", Some(Finish("END"))),
5824 ("ROLLBACK", Some(Finish("ROLLBACK"))),
5825 ("ROLLBACK TRANSACTION", Some(Finish("ROLLBACK"))),
5826 ("ROLLBACK TO sp", Some(Unsupported("ROLLBACK"))),
5827 (
5828 "ROLLBACK /* nested */ TRANSACTION /* target */ TO sp",
5829 Some(Unsupported("ROLLBACK")),
5830 ),
5831 ("SAVEPOINT sp", Some(Unsupported("SAVEPOINT"))),
5832 ("SELECT 1", None),
5833 ] {
5834 assert_eq!(
5835 cached_read_transaction_control(sql),
5836 expected,
5837 "cached-reader transaction classification mismatch for {sql:?}"
5838 );
5839 }
5840 }
5841
5842 #[test]
5843 fn sqlite_accepts_utf8_bom_before_transaction_control() {
5844 let conn = rusqlite::Connection::open_in_memory().unwrap();
5845 conn.execute_batch("CREATE TABLE bom_transaction_test (id INTEGER)")
5846 .unwrap();
5847 conn.execute_batch("\u{feff}BEGIN IMMEDIATE").unwrap();
5848 conn.execute_batch("ROLLBACK").unwrap();
5849 }
5850
5851 #[tokio::test]
5858 async fn execute_batch_rejects_transaction_control_on_queue_backed_path() {
5859 let dir = tempfile::tempdir().unwrap();
5860 let config = PoolConfig {
5861 path: Some(dir.path().join("sql_bridge_tx_reject_queue.db")),
5862 checkout_timeout: std::time::Duration::from_millis(250),
5863 write_queue_enabled: Some(true),
5864 write_routing_strict: true,
5865 ..PoolConfig::default()
5866 };
5867 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5868 {
5869 let guard = pool.writer().unwrap();
5870 guard
5871 .conn()
5872 .execute_batch(
5873 "CREATE TABLE tx_reject_queue_test (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
5874 )
5875 .unwrap();
5876 }
5877 let bridge = SqlBridge::new(Arc::clone(&pool), true);
5878 let mut writer = bridge.writer().await.unwrap();
5879
5880 let rejected = khive_storage::SqlWriter::execute_batch(
5881 &mut *writer,
5882 vec![
5883 SqlStatement {
5884 sql: "INSERT INTO tx_reject_queue_test (id, val) VALUES (1, 'a')".into(),
5885 params: vec![],
5886 label: None,
5887 },
5888 SqlStatement {
5889 sql: "COMMIT".into(),
5890 params: vec![],
5891 label: None,
5892 },
5893 ],
5894 )
5895 .await;
5896 assert!(
5897 matches!(&rejected, Err(StorageError::InvalidInput { .. })),
5898 "a bare COMMIT in a queue-backed batch must be rejected up front; got {rejected:?}"
5899 );
5900
5901 let prefixed = khive_storage::SqlWriter::execute_batch(
5902 &mut *writer,
5903 vec![
5904 SqlStatement {
5905 sql: "INSERT INTO tx_reject_queue_test (id, val) VALUES (3, 'prefixed')".into(),
5906 params: vec![],
5907 label: None,
5908 },
5909 SqlStatement {
5910 sql: "/* leading */ \u{feff} ; -- empty\n ; COMMIT".into(),
5911 params: vec![],
5912 label: None,
5913 },
5914 ],
5915 )
5916 .await;
5917 assert!(
5918 matches!(
5919 &prefixed,
5920 Err(StorageError::InvalidInput {
5921 operation,
5922 message,
5923 ..
5924 }) if operation.as_ref() == "execute_batch"
5925 && message.contains("transaction control")
5926 && message.contains("COMMIT")
5927 ),
5928 "a prefixed COMMIT must be rejected before touching the writer task; got {prefixed:?}"
5929 );
5930
5931 let affected = khive_storage::SqlWriter::execute_batch(
5932 &mut *writer,
5933 vec![SqlStatement {
5934 sql: "INSERT INTO tx_reject_queue_test (id, val) VALUES (2, 'b')".into(),
5935 params: vec![],
5936 label: None,
5937 }],
5938 )
5939 .await
5940 .expect("the writer task must survive the rejected batch");
5941 assert_eq!(affected, 1);
5942
5943 let count: i64 = {
5944 let guard = pool.reader().unwrap();
5945 guard
5946 .conn()
5947 .query_row("SELECT COUNT(*) FROM tx_reject_queue_test", [], |r| {
5948 r.get(0)
5949 })
5950 .unwrap()
5951 };
5952 assert_eq!(
5953 count, 1,
5954 "exactly the post-rejection batch's row may have landed"
5955 );
5956 }
5957
5958 #[tokio::test]
5971 async fn failed_rollback_poisons_handle_reuse_fails_loud() {
5972 use rusqlite::hooks::{AuthAction, AuthContext, Authorization, TransactionOperation};
5973
5974 fn deny_rollback(ctx: AuthContext<'_>) -> Authorization {
5975 match ctx.action {
5976 AuthAction::Transaction {
5977 operation: TransactionOperation::Rollback,
5978 } => Authorization::Deny,
5979 _ => Authorization::Allow,
5980 }
5981 }
5982
5983 let dir = tempfile::tempdir().unwrap();
5984 let config = PoolConfig {
5985 path: Some(dir.path().join("sql_bridge_rollback_poison.db")),
5986 checkout_timeout: std::time::Duration::from_millis(250),
5987 ..PoolConfig::default()
5988 };
5989 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5990 {
5991 let guard = pool.writer().unwrap();
5992 guard
5993 .conn()
5994 .execute_batch(
5995 "CREATE TABLE rollback_poison_test (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
5996 )
5997 .unwrap();
5998 }
5999
6000 let handle_slot = acquire_handle_slot(
6001 pool.sql_bridge_writer_slots(),
6002 pool.config().checkout_timeout,
6003 "sql_bridge.writer_handle",
6004 SlotTimeoutClass::Admission,
6005 )
6006 .await
6007 .unwrap();
6008 let conn = open_standalone_writer(&pool).unwrap();
6009 conn.authorizer(Some(deny_rollback)).unwrap();
6010 let mut writer = SqliteWriter {
6011 handle: Some(StandaloneHandle {
6012 conn,
6013 _retained_slot: Some(handle_slot),
6014 read_transaction_slot: None,
6015 }),
6016 writer_task: None,
6017 origin: pool.origin(),
6018 db: crate::timeout_sink::db_label(&pool),
6019 pool: Arc::clone(&pool),
6020 };
6021
6022 let batch = khive_storage::SqlWriter::execute_batch(
6023 &mut writer,
6024 vec![
6025 SqlStatement {
6026 sql: "INSERT INTO rollback_poison_test (id, val) VALUES (1, 'a')".into(),
6027 params: vec![],
6028 label: None,
6029 },
6030 SqlStatement {
6031 sql: "SELECT FROM WHERE".into(),
6032 params: vec![],
6033 label: None,
6034 },
6035 ],
6036 )
6037 .await;
6038 let batch_error = batch.expect_err("invalid second statement must fail the batch");
6039 let poison = match &batch_error {
6040 StorageError::Driver { source, .. } => source
6041 .downcast_ref::<PoisonedBatchError>()
6042 .expect("failed rollback must retain its typed poison wrapper"),
6043 other => panic!("failed rollback must return a driver error; got {other:?}"),
6044 };
6045 assert!(
6046 matches!(&poison.poison_reason, BatchPoisonReason::RollbackFailed(_)),
6047 "the poison cause must be compiler-checked as RollbackFailed; got {poison:?}"
6048 );
6049 let batch_message = batch_error.to_string();
6050 assert!(
6051 batch_message.contains("ROLLBACK after statement failure failed"),
6052 "the caller must see the poison context naming the failed \
6053 rollback; got {batch_message:?}"
6054 );
6055 assert!(
6056 batch_message.contains("original error"),
6057 "the original statement error must stay visible alongside the \
6058 poison context; got {batch_message:?}"
6059 );
6060
6061 let reuse = khive_storage::SqlWriter::execute(
6062 &mut writer,
6063 SqlStatement {
6064 sql: "CREATE TABLE rollback_poison_probe (id INTEGER PRIMARY KEY)".into(),
6065 params: vec![],
6066 label: None,
6067 },
6068 )
6069 .await;
6070 let message = match reuse {
6071 Err(StorageError::Pool { message, .. }) => message,
6072 other => panic!(
6073 "reusing a poisoned writer handle must fail loudly with \
6074 'connection already consumed'; got {other:?}"
6075 ),
6076 };
6077 assert!(
6078 message.contains("connection already consumed"),
6079 "expected the poisoned handle's reuse error to name the pinned \
6080 failure; got {message:?}"
6081 );
6082 }
6083
6084 #[tokio::test]
6090 async fn non_transient_begin_failure_poisons_handle() {
6091 let dir = tempfile::tempdir().unwrap();
6092 let config = PoolConfig {
6093 path: Some(dir.path().join("sql_bridge_begin_poison.db")),
6094 checkout_timeout: std::time::Duration::from_millis(250),
6095 ..PoolConfig::default()
6096 };
6097 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6098
6099 let handle_slot = acquire_handle_slot(
6100 pool.sql_bridge_writer_slots(),
6101 pool.config().checkout_timeout,
6102 "sql_bridge.writer_handle",
6103 SlotTimeoutClass::Admission,
6104 )
6105 .await
6106 .unwrap();
6107 let conn = open_standalone_writer(&pool).unwrap();
6108 conn.execute_batch("BEGIN IMMEDIATE").unwrap();
6111 let mut writer = SqliteWriter {
6112 handle: Some(StandaloneHandle {
6113 conn,
6114 _retained_slot: Some(handle_slot),
6115 read_transaction_slot: None,
6116 }),
6117 writer_task: None,
6118 origin: pool.origin(),
6119 db: crate::timeout_sink::db_label(&pool),
6120 pool: Arc::clone(&pool),
6121 };
6122
6123 let batch = khive_storage::SqlWriter::execute_batch(
6124 &mut writer,
6125 vec![SqlStatement {
6126 sql: "SELECT 1".into(),
6127 params: vec![],
6128 label: None,
6129 }],
6130 )
6131 .await;
6132 let batch_error = batch.expect_err("BEGIN inside an open transaction must fail");
6133 let poison = match &batch_error {
6134 StorageError::Driver { source, .. } => source
6135 .downcast_ref::<PoisonedBatchError>()
6136 .expect("failed BEGIN must retain its typed poison wrapper"),
6137 other => panic!("failed BEGIN must return a driver error; got {other:?}"),
6138 };
6139 assert!(
6140 matches!(&poison.poison_reason, BatchPoisonReason::BeginFailed),
6141 "the poison cause must be compiler-checked as BeginFailed; got {poison:?}"
6142 );
6143 let batch_message = batch_error.to_string();
6144 assert!(
6145 batch_message.contains("BEGIN IMMEDIATE failed non-transiently"),
6146 "a non-transient BEGIN failure must surface the poison context; \
6147 got {batch_message:?}"
6148 );
6149 assert!(
6150 batch_message.contains("cannot start a transaction within a transaction"),
6151 "the original BEGIN error must stay visible; got {batch_message:?}"
6152 );
6153
6154 let reuse = khive_storage::SqlWriter::execute(
6155 &mut writer,
6156 SqlStatement {
6157 sql: "CREATE TABLE begin_poison_probe (id INTEGER PRIMARY KEY)".into(),
6158 params: vec![],
6159 label: None,
6160 },
6161 )
6162 .await;
6163 assert!(
6164 matches!(
6165 &reuse,
6166 Err(StorageError::Pool { message, .. })
6167 if message.contains("connection already consumed")
6168 ),
6169 "a handle poisoned by a non-transient BEGIN failure must be \
6170 dropped, not restored; got {reuse:?}"
6171 );
6172 }
6173
6174 #[tokio::test]
6179 async fn busy_begin_failure_restores_handle_reusable() {
6180 let dir = tempfile::tempdir().unwrap();
6181 let config = PoolConfig {
6182 path: Some(dir.path().join("sql_bridge_begin_busy.db")),
6183 checkout_timeout: std::time::Duration::from_millis(250),
6184 busy_timeout: std::time::Duration::from_millis(100),
6185 ..PoolConfig::default()
6186 };
6187 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6188 {
6189 let guard = pool.writer().unwrap();
6190 guard
6191 .conn()
6192 .execute_batch("CREATE TABLE begin_busy_test (id INTEGER PRIMARY KEY)")
6193 .unwrap();
6194 }
6195
6196 let lock_conn = pool.open_standalone_writer().unwrap();
6200 lock_conn.execute_batch("BEGIN IMMEDIATE").unwrap();
6201
6202 let handle_slot = acquire_handle_slot(
6203 pool.sql_bridge_writer_slots(),
6204 pool.config().checkout_timeout,
6205 "sql_bridge.writer_handle",
6206 SlotTimeoutClass::Admission,
6207 )
6208 .await
6209 .unwrap();
6210 let conn = open_standalone_writer(&pool).unwrap();
6211 let mut writer = SqliteWriter {
6212 handle: Some(StandaloneHandle {
6213 conn,
6214 _retained_slot: Some(handle_slot),
6215 read_transaction_slot: None,
6216 }),
6217 writer_task: None,
6218 origin: pool.origin(),
6219 db: crate::timeout_sink::db_label(&pool),
6220 pool: Arc::clone(&pool),
6221 };
6222
6223 let batch = khive_storage::SqlWriter::execute_batch(
6224 &mut writer,
6225 vec![SqlStatement {
6226 sql: "INSERT INTO begin_busy_test (id) VALUES (1)".into(),
6227 params: vec![],
6228 label: None,
6229 }],
6230 )
6231 .await;
6232 let batch_error = batch.expect_err("BEGIN IMMEDIATE under a held write lock must fail");
6233 assert!(
6234 batch_error.to_string().contains("database is locked"),
6235 "the busy BEGIN failure must surface SQLite's busy error; got {batch_error:?}"
6236 );
6237
6238 lock_conn.execute_batch("ROLLBACK").unwrap();
6239 drop(lock_conn);
6240
6241 let affected = khive_storage::SqlWriter::execute(
6242 &mut writer,
6243 SqlStatement {
6244 sql: "INSERT INTO begin_busy_test (id) VALUES (2)".into(),
6245 params: vec![],
6246 label: None,
6247 },
6248 )
6249 .await
6250 .expect("a busy BEGIN failure must restore the handle as reusable");
6251 assert_eq!(affected, 1);
6252 }
6253
6254 #[tokio::test]
6259 async fn manual_atomic_unit_shares_writer_permit_budget_with_writer_handle() {
6260 let dir = tempfile::tempdir().unwrap();
6261 let config = PoolConfig {
6262 path: Some(dir.path().join("sql_bridge_atomic_unit_budget.db")),
6263 checkout_timeout: std::time::Duration::from_millis(50),
6264 write_queue_enabled: Some(false),
6265 ..PoolConfig::default()
6266 };
6267 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6268 let bridge = SqlBridge::new(Arc::clone(&pool), true);
6269 {
6270 let guard = pool.writer().unwrap();
6271 guard
6272 .conn()
6273 .execute_batch(
6274 "CREATE TABLE IF NOT EXISTS atomic_unit_budget_test \
6275 (id INTEGER PRIMARY KEY, val INTEGER NOT NULL)",
6276 )
6277 .unwrap();
6278 }
6279
6280 fn insert_op(id: i64) -> AtomicUnitOp {
6281 Box::new(move |writer| {
6282 Box::pin(async move {
6283 writer
6284 .execute(SqlStatement {
6285 sql: "INSERT INTO atomic_unit_budget_test (id, val) VALUES (?1, ?2)"
6286 .into(),
6287 params: vec![SqlValue::Integer(id), SqlValue::Integer(id)],
6288 label: None,
6289 })
6290 .await
6291 .map_err(|e| {
6292 khive_storage::StorageError::driver(
6293 StorageCapability::Sql,
6294 "atomic_unit_budget_test_insert",
6295 e,
6296 )
6297 })?;
6298 Ok(Box::new(()) as Box<dyn std::any::Any + Send>)
6299 })
6300 })
6301 }
6302
6303 let writer_handle = bridge.writer().await.unwrap();
6304 let blocked = bridge.atomic_unit(insert_op(1)).await;
6305 assert!(
6306 matches!(
6307 &blocked,
6308 Err(StorageError::AdmissionTimeout { operation, .. })
6309 if operation.as_ref() == "sql_bridge.atomic_unit_handle"
6310 ),
6311 "atomic_unit must time out on the shared writer permit while a \
6312 writer handle is live; got {blocked:?}"
6313 );
6314
6315 drop(writer_handle);
6316 let unblocked = bridge.atomic_unit(insert_op(2)).await;
6317 assert!(
6318 unblocked.is_ok(),
6319 "atomic_unit must succeed once the writer handle releases the \
6320 shared writer permit; got {unblocked:?}"
6321 );
6322
6323 let mut reader = bridge.reader().await.unwrap();
6324 let count = reader
6325 .query_scalar(SqlStatement {
6326 sql: "SELECT COUNT(*) FROM atomic_unit_budget_test".into(),
6327 params: vec![],
6328 label: None,
6329 })
6330 .await
6331 .unwrap();
6332 assert!(
6333 matches!(count, Some(SqlValue::Integer(1))),
6334 "only the post-drop atomic_unit call may have committed; got {count:?}"
6335 );
6336 }
6337
6338 #[tokio::test]
6344 async fn execute_batch_routes_through_writer_task_when_flag_enabled() {
6345 let dir = tempfile::tempdir().unwrap();
6346 let path = dir.path().join("write_queue_execute_batch.db");
6347 let config = PoolConfig {
6348 path: Some(path.clone()),
6349 write_queue_enabled: Some(true),
6350 ..PoolConfig::default()
6351 };
6352 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6353 {
6354 let guard = pool.writer().unwrap();
6355 guard
6356 .conn()
6357 .execute_batch(
6358 "CREATE TABLE IF NOT EXISTS write_queue_batch_test \
6359 (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
6360 )
6361 .unwrap();
6362 }
6363
6364 let bridge = SqlBridge::new(Arc::clone(&pool), true);
6365
6366 let mut writer = bridge.writer().await.unwrap();
6367 let affected = writer
6368 .execute_batch(vec![
6369 SqlStatement {
6370 sql: "INSERT INTO write_queue_batch_test (id, val) VALUES (?1, ?2)".into(),
6371 params: vec![SqlValue::Integer(1), SqlValue::Text("a".into())],
6372 label: None,
6373 },
6374 SqlStatement {
6375 sql: "INSERT INTO write_queue_batch_test (id, val) VALUES (?1, ?2)".into(),
6376 params: vec![SqlValue::Integer(2), SqlValue::Text("b".into())],
6377 label: None,
6378 },
6379 ])
6380 .await
6381 .unwrap();
6382 assert_eq!(affected, 2);
6383
6384 let mut reader = bridge.reader().await.unwrap();
6385 let count = reader
6386 .query_scalar(SqlStatement {
6387 sql: "SELECT COUNT(*) FROM write_queue_batch_test".into(),
6388 params: vec![],
6389 label: None,
6390 })
6391 .await
6392 .unwrap();
6393 assert!(
6394 matches!(count, Some(SqlValue::Integer(2))),
6395 "expected 2 rows, got {count:?}"
6396 );
6397 assert_eq!(
6398 pool.writer_task_spawn_count(),
6399 1,
6400 "the flag-ON path must actually spawn and use the writer task"
6401 );
6402 }
6403
6404 #[tokio::test]
6411 async fn execute_batch_rolls_back_atomically_on_mid_sequence_failure() {
6412 let dir = tempfile::tempdir().unwrap();
6413 let path = dir.path().join("write_queue_execute_batch_rollback.db");
6414 let config = PoolConfig {
6415 path: Some(path.clone()),
6416 write_queue_enabled: Some(true),
6417 ..PoolConfig::default()
6418 };
6419 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6420 {
6421 let guard = pool.writer().unwrap();
6422 guard
6423 .conn()
6424 .execute_batch(
6425 "CREATE TABLE IF NOT EXISTS write_queue_rollback_test \
6426 (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
6427 )
6428 .unwrap();
6429 }
6430
6431 let bridge = SqlBridge::new(Arc::clone(&pool), true);
6432
6433 let mut writer = bridge.writer().await.unwrap();
6434 let result = writer
6435 .execute_batch(vec![
6436 SqlStatement {
6438 sql: "INSERT INTO write_queue_rollback_test (id, val) VALUES (?1, ?2)".into(),
6439 params: vec![SqlValue::Integer(1), SqlValue::Text("first".into())],
6440 label: None,
6441 },
6442 SqlStatement {
6444 sql: "INSERT INTO write_queue_rollback_test (id, val) VALUES (?1, ?2)".into(),
6445 params: vec![SqlValue::Integer(1), SqlValue::Text("duplicate".into())],
6446 label: None,
6447 },
6448 SqlStatement {
6450 sql: "INSERT INTO write_queue_rollback_test (id, val) VALUES (?1, ?2)".into(),
6451 params: vec![SqlValue::Integer(2), SqlValue::Text("third".into())],
6452 label: None,
6453 },
6454 ])
6455 .await;
6456 assert!(
6457 result.is_err(),
6458 "a batch with a mid-sequence PK conflict must return an error"
6459 );
6460
6461 let mut reader = bridge.reader().await.unwrap();
6462 let count = reader
6463 .query_scalar(SqlStatement {
6464 sql: "SELECT COUNT(*) FROM write_queue_rollback_test".into(),
6465 params: vec![],
6466 label: None,
6467 })
6468 .await
6469 .unwrap();
6470 assert!(
6471 matches!(count, Some(SqlValue::Integer(0))),
6472 "the whole request must roll back — including statement 1's \
6473 otherwise-successful INSERT — not just the failing statement; \
6474 got {count:?}"
6475 );
6476 }
6477
6478 #[tokio::test]
6495 async fn atomic_unit_pending_future_errors_without_killing_writer_task() {
6496 let dir = tempfile::tempdir().unwrap();
6497 let path = dir.path().join("atomic_unit_pending_future.db");
6498 let config = PoolConfig {
6499 path: Some(path.clone()),
6500 write_queue_enabled: Some(true),
6501 ..PoolConfig::default()
6502 };
6503 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6504 {
6505 let guard = pool.writer().unwrap();
6506 guard
6507 .conn()
6508 .execute_batch(
6509 "CREATE TABLE IF NOT EXISTS atomic_unit_pending_test \
6510 (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
6511 )
6512 .unwrap();
6513 }
6514 assert!(
6515 pool.writer_task_handle().unwrap().is_some(),
6516 "writer task must be spawned with the flag on for a file-backed pool"
6517 );
6518
6519 let bridge = SqlBridge::new(Arc::clone(&pool), true);
6520
6521 let pending_op: AtomicUnitOp = Box::new(|_writer| {
6524 Box::pin(std::future::pending::<
6525 khive_storage::types::StorageResult<Box<dyn std::any::Any + Send>>,
6526 >())
6527 });
6528
6529 let pending_result = bridge.atomic_unit(pending_op).await;
6530 assert!(
6531 pending_result.is_err(),
6532 "a Pending-on-first-poll atomic_unit closure must return Err, \
6533 not panic; got {pending_result:?}"
6534 );
6535
6536 let ok_op: AtomicUnitOp = Box::new(|writer| {
6541 Box::pin(async move {
6542 writer
6543 .execute(SqlStatement {
6544 sql: "INSERT INTO atomic_unit_pending_test (id, val) VALUES (?1, ?2)"
6545 .into(),
6546 params: vec![SqlValue::Integer(1), SqlValue::Text("survived".into())],
6547 label: None,
6548 })
6549 .await
6550 .map_err(|e| {
6551 khive_storage::StorageError::driver(
6552 StorageCapability::Sql,
6553 "atomic_unit_pending_future_test_insert",
6554 e,
6555 )
6556 })?;
6557 Ok(Box::new(()) as Box<dyn std::any::Any + Send>)
6558 })
6559 });
6560 let ok_result = bridge.atomic_unit(ok_op).await;
6561 assert!(
6562 ok_result.is_ok(),
6563 "writer task must survive a Pending misuse and keep serving \
6564 subsequent well-behaved atomic_unit requests; got {ok_result:?}"
6565 );
6566
6567 let mut reader = bridge.reader().await.unwrap();
6568 let count = reader
6569 .query_scalar(SqlStatement {
6570 sql: "SELECT COUNT(*) FROM atomic_unit_pending_test".into(),
6571 params: vec![],
6572 label: None,
6573 })
6574 .await
6575 .unwrap();
6576 assert!(
6577 matches!(count, Some(SqlValue::Integer(1))),
6578 "the well-behaved atomic_unit call after the Pending misuse must \
6579 have actually committed its write; got {count:?}"
6580 );
6581 }
6582
6583 #[tokio::test]
6593 async fn writer_strict_routing_fails_closed_without_writer_task() {
6594 let dir = tempfile::tempdir().unwrap();
6595 let path = dir.path().join("strict_writer.db");
6596 let config = PoolConfig {
6597 path: Some(path),
6598 write_queue_enabled: Some(false),
6599 write_routing_strict: true,
6600 ..PoolConfig::default()
6601 };
6602 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6603 let bridge = SqlBridge::new(Arc::clone(&pool), true);
6604
6605 let result = bridge.writer().await;
6606 let err = match result {
6607 Ok(_) => panic!(
6608 "KHIVE_WRITE_ROUTING=strict with no writer task must fail closed, not \
6609 silently degrade to a standalone connection"
6610 ),
6611 Err(e) => e,
6612 };
6613 assert!(
6614 err.to_string().contains("strict"),
6615 "error must name strict routing, got: {err}"
6616 );
6617 }
6618
6619 #[tokio::test]
6624 async fn atomic_unit_strict_routing_fails_closed_without_writer_task() {
6625 let dir = tempfile::tempdir().unwrap();
6626 let path = dir.path().join("strict_atomic_unit.db");
6627 let config = PoolConfig {
6628 path: Some(path),
6629 write_queue_enabled: Some(false),
6630 write_routing_strict: true,
6631 ..PoolConfig::default()
6632 };
6633 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6634 let bridge = SqlBridge::new(Arc::clone(&pool), true);
6635
6636 let op: AtomicUnitOp = Box::new(|_writer| {
6637 Box::pin(async move { Ok(Box::new(()) as Box<dyn std::any::Any + Send>) })
6638 });
6639 let result = bridge.atomic_unit(op).await;
6640 assert!(
6641 result.is_err(),
6642 "KHIVE_WRITE_ROUTING=strict but the queue is off (no writer task handle) must \
6643 fail closed instead of falling back to a manual BEGIN IMMEDIATE; got {result:?}"
6644 );
6645 let msg = result.unwrap_err().to_string();
6646 assert!(
6647 msg.contains("strict"),
6648 "error must name strict routing, got: {msg}"
6649 );
6650 }
6651
6652 #[tokio::test]
6660 async fn writer_handle_supports_read_after_write_under_strict_queue() {
6661 let dir = tempfile::tempdir().unwrap();
6662 let path = dir.path().join("writer_read_after_write.db");
6663 let config = PoolConfig {
6664 path: Some(path),
6665 write_queue_enabled: Some(true),
6666 write_routing_strict: true,
6667 ..PoolConfig::default()
6668 };
6669 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6670 {
6671 let guard = pool.writer().unwrap();
6672 guard
6673 .conn()
6674 .execute_batch(
6675 "CREATE TABLE IF NOT EXISTS writer_cursor_test \
6676 (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
6677 )
6678 .unwrap();
6679 }
6680
6681 let bridge = SqlBridge::new(Arc::clone(&pool), true);
6682
6683 let mut w = bridge.writer().await.unwrap();
6684 w.execute(SqlStatement {
6685 sql: "INSERT INTO writer_cursor_test (id, val) VALUES (?1, ?2)".into(),
6686 params: vec![SqlValue::Integer(1), SqlValue::Text("via-writer".into())],
6687 label: None,
6688 })
6689 .await
6690 .unwrap();
6691
6692 let row = w
6693 .query_row(SqlStatement {
6694 sql: "SELECT val FROM writer_cursor_test WHERE id = ?1".into(),
6695 params: vec![SqlValue::Integer(1)],
6696 label: None,
6697 })
6698 .await
6699 .unwrap()
6700 .expect("row inserted through the same writer handle must be visible to it");
6701 assert!(
6702 matches!(&row.columns[0].value, SqlValue::Text(v) if v == "via-writer"),
6703 "query_row through a queue-backed writer handle must see its own \
6704 committed write; got {:?}",
6705 row.columns[0].value
6706 );
6707 }
6708
6709 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6718 async fn queue_backed_read_uses_reader_budget_and_reopens_after_cancel() {
6719 let dir = tempfile::tempdir().unwrap();
6720 let path = dir.path().join("queue_backed_reader_budget.db");
6721 let config = PoolConfig {
6722 path: Some(path),
6723 write_queue_enabled: Some(true),
6724 write_routing_strict: true,
6725 max_readers: 1,
6726 checkout_timeout: std::time::Duration::from_millis(250),
6727 ..PoolConfig::default()
6728 };
6729 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6730 {
6731 let guard = pool.writer().unwrap();
6732 guard
6733 .conn()
6734 .execute_batch(
6735 "CREATE TABLE IF NOT EXISTS reopen_test \
6736 (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
6737 )
6738 .unwrap();
6739 }
6740 let bridge = SqlBridge::new(Arc::clone(&pool), true);
6741
6742 let mut w = bridge.writer().await.unwrap();
6743 w.execute(SqlStatement {
6744 sql: "INSERT INTO reopen_test (id, val) VALUES (1, 'seed')".into(),
6745 params: vec![],
6746 label: None,
6747 })
6748 .await
6749 .unwrap();
6750
6751 let held = pool
6754 .sql_bridge_reader_slots()
6755 .acquire_owned()
6756 .await
6757 .unwrap();
6758 let starved = w
6759 .query_row(SqlStatement {
6760 sql: "SELECT val FROM reopen_test WHERE id = 1".into(),
6761 params: vec![],
6762 label: None,
6763 })
6764 .await;
6765 assert!(
6766 matches!(
6767 &starved,
6768 Err(StorageError::Timeout { operation })
6769 if operation.as_ref() == "sql_bridge.reader_open"
6770 ),
6771 "queue-backed read with reader permits saturated must time out \
6772 on the reader budget with the ADR-005 reader contract error; \
6773 got {starved:?}"
6774 );
6775 drop(held);
6776
6777 let writer_task = pool
6785 .writer_task_handle()
6786 .expect("queue-enabled file pool must offer a writer task")
6787 .expect("writer task present under write_queue_enabled");
6788 let mut post_cancel = SqliteWriter {
6789 handle: None,
6790 writer_task: Some(writer_task),
6791 origin: pool.origin(),
6792 db: crate::timeout_sink::db_label(&pool),
6793 pool: Arc::clone(&pool),
6794 };
6795 let row = post_cancel
6796 .query_row(SqlStatement {
6797 sql: "SELECT val FROM reopen_test WHERE id = 1".into(),
6798 params: vec![],
6799 label: None,
6800 })
6801 .await
6802 .expect("read on a queue-backed handle with no resident connection must reopen")
6803 .expect("seeded row must be visible");
6804 assert!(
6805 matches!(&row.columns[0].value, SqlValue::Text(v) if v == "seed"),
6806 "reopened read must return the seeded row; got {:?}",
6807 row.columns[0].value
6808 );
6809 }
6810
6811 #[tokio::test]
6822 async fn writer_query_row_rejects_dml_with_returning_on_queue_backed_handle() {
6823 let dir = tempfile::tempdir().unwrap();
6824 let path = dir.path().join("writer_readonly_returning.db");
6825 let config = PoolConfig {
6826 path: Some(path),
6827 write_queue_enabled: Some(true),
6828 write_routing_strict: true,
6829 ..PoolConfig::default()
6830 };
6831 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6832 {
6833 let guard = pool.writer().unwrap();
6834 guard
6835 .conn()
6836 .execute_batch(
6837 "CREATE TABLE IF NOT EXISTS writer_returning_test \
6838 (id INTEGER PRIMARY KEY, val TEXT NOT NULL);
6839 INSERT INTO writer_returning_test (id, val) VALUES (1, 'original');",
6840 )
6841 .unwrap();
6842 }
6843
6844 let bridge = SqlBridge::new(Arc::clone(&pool), true);
6845
6846 let mut w = bridge.writer().await.unwrap();
6847 let result = w
6848 .query_row(SqlStatement {
6849 sql: "UPDATE writer_returning_test SET val = 'mutated' \
6850 WHERE id = ?1 RETURNING val"
6851 .into(),
6852 params: vec![SqlValue::Integer(1)],
6853 label: None,
6854 })
6855 .await;
6856 assert!(
6857 result.is_err(),
6858 "a DML-with-RETURNING statement through query_row on a \
6859 queue-backed writer handle must be rejected, not executed on \
6860 an untracked read-write connection; got {result:?}"
6861 );
6862
6863 let mut reader = bridge.reader().await.unwrap();
6864 let val = reader
6865 .query_scalar(SqlStatement {
6866 sql: "SELECT val FROM writer_returning_test WHERE id = ?1".into(),
6867 params: vec![SqlValue::Integer(1)],
6868 label: None,
6869 })
6870 .await
6871 .unwrap();
6872 assert!(
6873 matches!(&val, Some(SqlValue::Text(v)) if v == "original"),
6874 "the rejected UPDATE...RETURNING must not have altered the row; got {val:?}"
6875 );
6876 }
6877
6878 #[tokio::test]
6898 async fn acceptance_five_op_batch_completes_under_concurrent_write_contention() {
6899 let dir = tempfile::tempdir().unwrap();
6900 let path = dir.path().join("acceptance_batch.db");
6901 let config = PoolConfig {
6902 path: Some(path),
6903 write_queue_enabled: Some(true),
6904 write_routing_strict: true,
6905 ..PoolConfig::default()
6906 };
6907 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6908 {
6909 let guard = pool.writer().unwrap();
6910 guard
6911 .conn()
6912 .execute_batch(
6913 "CREATE TABLE IF NOT EXISTS acceptance_batch \
6914 (id INTEGER PRIMARY KEY, val TEXT NOT NULL);
6915 INSERT INTO acceptance_batch (id, val) VALUES \
6916 (200, 'seed-0'), (201, 'seed-1'), (202, 'seed-2'), (203, 'seed-3');",
6917 )
6918 .unwrap();
6919 }
6920
6921 let bridge = Arc::new(SqlBridge::new(Arc::clone(&pool), true));
6922
6923 let writer_task = pool
6924 .writer_task_handle()
6925 .unwrap()
6926 .expect("writer task must be spawned for a file-backed pool with the flag on");
6927
6928 let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
6934 let (release_tx, release_rx) = tokio::sync::oneshot::channel::<()>();
6935 let occupier = {
6936 let writer_task = writer_task.clone();
6937 tokio::spawn(async move {
6938 writer_task
6939 .send(move |_conn| {
6940 let _ = started_tx.send(());
6941 let _ = release_rx.blocking_recv();
6942 Ok::<(), StorageError>(())
6943 })
6944 .await
6945 })
6946 };
6947 started_rx
6948 .await
6949 .expect("occupier must signal it has started running inside the writer task");
6950 assert_eq!(
6951 writer_task.queue_depth(),
6952 0,
6953 "channel must start empty once the occupier has been dequeued and is running"
6954 );
6955
6956 let contenders: Vec<_> = (0..3)
6961 .map(|i| {
6962 let bridge = Arc::clone(&bridge);
6963 tokio::spawn(async move {
6964 let mut writer = bridge.writer().await?;
6965 writer
6966 .execute(SqlStatement {
6967 sql: "INSERT INTO acceptance_batch (id, val) VALUES (?1, ?2)".into(),
6968 params: vec![
6969 SqlValue::Integer(100 + i),
6970 SqlValue::Text(format!("contender-{i}")),
6971 ],
6972 label: None,
6973 })
6974 .await
6975 })
6976 })
6977 .collect();
6978
6979 let send = {
6981 let bridge = Arc::clone(&bridge);
6982 tokio::spawn(async move {
6983 let mut writer = bridge.writer().await?;
6984 writer
6985 .execute(SqlStatement {
6986 sql: "INSERT INTO acceptance_batch (id, val) VALUES (?1, ?2)".into(),
6987 params: vec![SqlValue::Integer(1), SqlValue::Text("send".into())],
6988 label: None,
6989 })
6990 .await
6991 })
6992 };
6993 let marks: Vec<_> = (0..4)
6994 .map(|i| {
6995 let bridge = Arc::clone(&bridge);
6996 tokio::spawn(async move {
6997 let mut writer = bridge.writer().await?;
6998 writer
6999 .execute(SqlStatement {
7000 sql: "UPDATE acceptance_batch SET val = ?2 WHERE id = ?1".into(),
7001 params: vec![
7002 SqlValue::Integer(200 + i),
7003 SqlValue::Text(format!("marked-{i}")),
7004 ],
7005 label: None,
7006 })
7007 .await
7008 })
7009 })
7010 .collect();
7011
7012 let mut saw_all_enqueued = false;
7015 for _ in 0..200 {
7016 if writer_task.queue_depth() >= 8 {
7017 saw_all_enqueued = true;
7018 break;
7019 }
7020 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
7021 }
7022 assert!(
7023 saw_all_enqueued,
7024 "not all 8 contending writes (3 contenders + send + 4 marks) reached \
7025 the writer task's channel while the occupier held the single drain \
7026 slot — got depth {}",
7027 writer_task.queue_depth()
7028 );
7029
7030 release_tx
7031 .send(())
7032 .expect("occupier must still be waiting on the release signal");
7033 occupier
7034 .await
7035 .expect("occupier task must not panic")
7036 .expect("occupier write must succeed");
7037
7038 for c in contenders {
7039 c.await
7040 .expect("contender task must not panic")
7041 .expect("contender write must complete without a checkout timeout");
7042 }
7043 send.await
7044 .expect("send task must not panic")
7045 .expect("send op must complete without a checkout timeout");
7046 for (i, m) in marks.into_iter().enumerate() {
7047 let affected = m
7048 .await
7049 .expect("mark task must not panic")
7050 .expect("mark op must complete without a checkout timeout — no starvation");
7051 assert_eq!(
7052 affected, 1,
7053 "mark {i} must have updated exactly its own row"
7054 );
7055 }
7056
7057 let mut reader = bridge.reader().await.unwrap();
7058 let count = reader
7059 .query_scalar(SqlStatement {
7060 sql: "SELECT COUNT(*) FROM acceptance_batch".into(),
7061 params: vec![],
7062 label: None,
7063 })
7064 .await
7065 .unwrap();
7066 assert!(
7067 matches!(count, Some(SqlValue::Integer(8))),
7068 "the 4 seeded mark rows plus 3 contenders plus the batch's own send \
7069 must all be present; got {count:?}"
7070 );
7071
7072 for i in 0..4i64 {
7073 let mut reader = bridge.reader().await.unwrap();
7074 let val = reader
7075 .query_scalar(SqlStatement {
7076 sql: "SELECT val FROM acceptance_batch WHERE id = ?1".into(),
7077 params: vec![SqlValue::Integer(200 + i)],
7078 label: None,
7079 })
7080 .await
7081 .unwrap();
7082 assert!(
7083 matches!(&val, Some(SqlValue::Text(v)) if *v == format!("marked-{i}")),
7084 "mark row {i} must reflect the persisted UPDATE after release; got {val:?}"
7085 );
7086 }
7087 }
7088
7089 #[tokio::test]
7090 async fn file_backed_bridge_counts_writer_and_flag_off_atomic_unit_acquisitions() {
7091 let dir = tempfile::tempdir().unwrap();
7092 let config = PoolConfig {
7093 path: Some(dir.path().join("bridge_writer_acquisitions.db")),
7094 write_queue_enabled: Some(false),
7095 ..PoolConfig::default()
7096 };
7097 let pool = Arc::new(ConnectionPool::new(config).unwrap());
7098 let bridge = SqlBridge::new(Arc::clone(&pool), true);
7099
7100 let before = pool.writer_acquisition_snapshot();
7101
7102 drop(bridge.writer().await.unwrap());
7103 let after_writer = pool.writer_acquisition_snapshot();
7104 assert_eq!(
7105 after_writer.standalone_acquisitions,
7106 before.standalone_acquisitions + 1
7107 );
7108 assert_eq!(after_writer.acquisitions, before.acquisitions + 1);
7109 assert_eq!(after_writer.pooled_acquisitions, before.pooled_acquisitions);
7110 assert_eq!(
7111 after_writer.writer_task_acquisitions,
7112 before.writer_task_acquisitions
7113 );
7114
7115 let op: AtomicUnitOp = Box::new(|_writer| {
7116 Box::pin(async { Ok(Box::new(()) as Box<dyn std::any::Any + Send>) })
7117 });
7118 bridge.atomic_unit(op).await.unwrap();
7119
7120 let after_atomic_unit = pool.writer_acquisition_snapshot();
7121 assert_eq!(
7122 after_atomic_unit.standalone_acquisitions,
7123 before.standalone_acquisitions + 2
7124 );
7125 assert_eq!(after_atomic_unit.acquisitions, before.acquisitions + 2);
7126 assert_eq!(
7127 after_atomic_unit.pooled_acquisitions,
7128 before.pooled_acquisitions
7129 );
7130 assert_eq!(
7131 after_atomic_unit.writer_task_acquisitions,
7132 before.writer_task_acquisitions
7133 );
7134 }
7135}