1#[path = "sql_bridge/write_errors.rs"]
13mod write_errors;
14
15#[path = "sql_bridge/manual_atomic.rs"]
16mod manual_atomic;
17use manual_atomic::run_manual_atomic_unit;
18#[path = "sql_bridge/standalone_admission.rs"]
19mod standalone_admission;
20use standalone_admission::{
21 acquire_standalone_lease, acquire_unit_lease, admit_standalone_operation, StandaloneWriteError,
22};
23mod standalone_batch;
24#[cfg(test)]
25use standalone_batch::BatchPoisonReason;
26use standalone_batch::{
27 run_standalone_batch, run_standalone_script, run_standalone_statement,
28 run_standalone_top_level, BatchFailure, PoisonedBatchError,
29};
30
31use std::any::Any;
32use std::sync::atomic::{AtomicU64, Ordering};
33use std::sync::Arc;
34use std::time::Instant;
35
36use async_trait::async_trait;
37
38use khive_storage::error::StorageError;
39use khive_storage::types::{PageRequest, SqlColumn, SqlRow, SqlStatement, SqlValue};
40use khive_storage::{AtomicUnitOp, StorageCapability, TopLevelMaintenance};
41use tokio::sync::{OwnedSemaphorePermit, Semaphore};
42
43use crate::error::SqliteError;
44use crate::pool::{ConnectionPool, SharedReaderTransactionGuard, StandaloneReaderPurpose};
45
46mod rows;
51
52pub(crate) use rows::bind_params;
53use rows::{
54 prepare_batch_statements, prepare_cached_sql_statement, prepare_sql_statement, row_to_sql_row,
55 AtomicEventRows, PreparedBatchStatement,
56};
57#[cfg(test)]
58use rows::{COUNTED_EVENT_INSERT_LABELS, ROW_CONVERSIONS};
59
60#[cfg(test)]
61#[path = "atomic_event_usage_tests.rs"]
62mod atomic_event_usage_tests;
63
64fn execute_prepared_batch<'conn>(
66 conn: &'conn rusqlite::Connection,
67 prepared: Vec<PreparedBatchStatement<'conn>>,
68 statements: &[SqlStatement],
69 event_rows: Option<&AtomicEventRows>,
70) -> Result<u64, rusqlite::Error> {
71 debug_assert_eq!(prepared.len(), statements.len());
72 let mut total = 0u64;
73 for (prepared, statement) in prepared.into_iter().zip(statements) {
74 let mut prepared = match prepared {
75 PreparedBatchStatement::Ready(prepared) => prepared,
76 PreparedBatchStatement::PrepareAtExecution => {
77 prepare_sql_statement(conn, &statement.sql)?
78 }
79 };
80 bind_params(&mut prepared, &statement.params)?;
81 let affected = prepared.raw_execute()? as u64;
82 if let Some(event_rows) = event_rows {
83 event_rows.observe(statement, affected);
84 }
85 total += affected;
86 }
87 Ok(total)
88}
89
90const TRANSACTION_CONTROL_KEYWORDS: [&str; 7] = [
101 "BEGIN",
102 "START",
103 "COMMIT",
104 "END",
105 "ROLLBACK",
106 "SAVEPOINT",
107 "RELEASE",
108];
109
110fn skip_sqlite_empty_prefix(mut rest: &[u8]) -> &[u8] {
113 loop {
114 let mut idx = 0;
115 while idx < rest.len() && rest[idx].is_ascii_whitespace() {
116 idx += 1;
117 }
118 rest = &rest[idx..];
119 if let Some(tail) = rest.strip_prefix(b"\xEF\xBB\xBF") {
120 rest = tail;
121 continue;
122 }
123 if let Some(tail) = rest.strip_prefix(b";") {
124 rest = tail;
125 continue;
126 }
127 if let Some(tail) = rest.strip_prefix(b"--") {
128 let mut idx = 0;
129 while idx < tail.len() && tail[idx] != b'\n' {
130 idx += 1;
131 }
132 rest = if idx < tail.len() {
133 &tail[idx + 1..]
134 } else {
135 &[]
136 };
137 continue;
138 }
139 if let Some(tail) = rest.strip_prefix(b"/*") {
140 let mut idx = 0;
141 while idx + 1 < tail.len() && !(tail[idx] == b'*' && tail[idx + 1] == b'/') {
142 idx += 1;
143 }
144 rest = if idx + 1 < tail.len() {
145 &tail[idx + 2..]
146 } else {
147 &[]
148 };
149 continue;
150 }
151 break;
152 }
153 rest
154}
155
156fn next_sqlite_token(mut rest: &[u8]) -> Option<(&[u8], &[u8])> {
160 loop {
161 let mut idx = 0;
162 while idx < rest.len() && rest[idx].is_ascii_whitespace() {
163 idx += 1;
164 }
165 rest = &rest[idx..];
166 if let Some(tail) = rest.strip_prefix(b"\xEF\xBB\xBF") {
167 rest = tail;
168 continue;
169 }
170 if let Some(tail) = rest.strip_prefix(b"--") {
171 let mut idx = 0;
172 while idx < tail.len() && tail[idx] != b'\n' {
173 idx += 1;
174 }
175 rest = if idx < tail.len() {
176 &tail[idx + 1..]
177 } else {
178 &[]
179 };
180 continue;
181 }
182 if let Some(tail) = rest.strip_prefix(b"/*") {
183 let mut idx = 0;
184 while idx + 1 < tail.len() && !(tail[idx] == b'*' && tail[idx + 1] == b'/') {
185 idx += 1;
186 }
187 rest = if idx + 1 < tail.len() {
188 &tail[idx + 2..]
189 } else {
190 &[]
191 };
192 continue;
193 }
194 break;
195 }
196
197 let len = rest
198 .iter()
199 .take_while(|byte| byte.is_ascii_alphanumeric() || **byte == b'_')
200 .count();
201 (len != 0).then_some((&rest[..len], &rest[len..]))
202}
203
204fn transaction_control_parts(sql: &str) -> Option<(&'static str, &[u8])> {
208 let rest = skip_sqlite_empty_prefix(sql.as_bytes());
209 TRANSACTION_CONTROL_KEYWORDS
210 .iter()
211 .copied()
212 .find_map(|keyword| {
213 let kw = keyword.as_bytes();
214 if rest.len() < kw.len() || !rest[..kw.len()].eq_ignore_ascii_case(kw) {
215 return None;
216 }
217 let boundary = match rest.get(kw.len()) {
218 Some(next) => !(next.is_ascii_alphanumeric() || *next == b'_'),
219 None => true,
220 };
221 boundary.then_some((keyword, &rest[kw.len()..]))
222 })
223}
224
225fn transaction_control_head(sql: &str) -> Option<&'static str> {
227 transaction_control_parts(sql).map(|(keyword, _)| keyword)
228}
229
230#[derive(Clone, Copy, Debug, PartialEq, Eq)]
231enum CachedReadTransactionControl {
232 BeginDeferred,
234 Finish(&'static str),
236 Unsupported(&'static str),
239}
240
241fn cached_read_transaction_control(sql: &str) -> Option<CachedReadTransactionControl> {
249 let (keyword, tail) = transaction_control_parts(sql)?;
250 match keyword {
251 "BEGIN" => {
252 let mut rest = tail;
269 let mut saw_deferred = false;
270 let mut saw_transaction = false;
271 while let Some((token, next)) = next_sqlite_token(rest) {
272 if !saw_deferred && !saw_transaction && token.eq_ignore_ascii_case(b"DEFERRED") {
273 saw_deferred = true;
274 } else if !saw_transaction && token.eq_ignore_ascii_case(b"TRANSACTION") {
275 saw_transaction = true;
276 } else {
277 return Some(CachedReadTransactionControl::Unsupported(keyword));
278 }
279 rest = next;
280 }
281 if !skip_sqlite_empty_prefix(rest).is_empty() {
282 return Some(CachedReadTransactionControl::Unsupported(keyword));
283 }
284 Some(CachedReadTransactionControl::BeginDeferred)
285 }
286 "COMMIT" | "END" => Some(CachedReadTransactionControl::Finish(keyword)),
287 "ROLLBACK" => {
288 let first = next_sqlite_token(tail);
289 let rollback_target = match first {
290 Some((token, rest)) if token.eq_ignore_ascii_case(b"TRANSACTION") => {
291 next_sqlite_token(rest).map(|(token, _)| token)
292 }
293 Some((token, _)) => Some(token),
294 None => None,
295 };
296 if rollback_target.is_some_and(|token| token.eq_ignore_ascii_case(b"TO")) {
297 Some(CachedReadTransactionControl::Unsupported(keyword))
298 } else {
299 Some(CachedReadTransactionControl::Finish(keyword))
300 }
301 }
302 _ => Some(CachedReadTransactionControl::Unsupported(keyword)),
303 }
304}
305
306fn reject_transaction_control_statements(
310 statements: &[SqlStatement],
311 operation: &'static str,
312) -> khive_storage::types::StorageResult<()> {
313 for (index, statement) in statements.iter().enumerate() {
314 if let Some(keyword) = transaction_control_head(&statement.sql) {
315 return Err(StorageError::InvalidInput {
316 capability: StorageCapability::Sql,
317 operation: operation.into(),
318 message: format!(
319 "statement at index {index} is transaction control ({keyword}); \
320 execute_batch owns the BEGIN/COMMIT boundary for the whole \
321 batch — remove transaction-control statements from the batch"
322 ),
323 });
324 }
325 }
326 Ok(())
327}
328
329fn settle_pooled_call<T>(
341 guard: &crate::pool::WriterGuard<'_>,
342 operation: &'static str,
343 result: khive_storage::types::StorageResult<T>,
344) -> khive_storage::types::StorageResult<T> {
345 if guard.is_autocommit() {
346 return result;
347 }
348 if let Err(settlement) = guard.rollback_or_retire("pooled call left its transaction open") {
349 if let Err(error) = &result {
350 tracing::warn!(
351 operation,
352 %error,
353 "pooled call failed inside a transaction it opened, and its rollback \
354 could not prove autocommit; reporting the settlement failure"
355 );
356 }
357 return Err(settlement.into_storage_error(StorageCapability::Sql, operation));
358 }
359 result?;
360 Err(StorageError::InvalidInput {
361 capability: StorageCapability::Sql,
362 operation: operation.into(),
363 message: "the call left a transaction open; it was rolled back before the pooled \
364 writer was released — use atomic_unit to run statements as one transaction"
365 .into(),
366 })
367}
368
369fn prepare_bound_statement<'conn>(
370 conn: &'conn rusqlite::Connection,
371 statement: &SqlStatement,
372) -> Result<rusqlite::Statement<'conn>, rusqlite::Error> {
373 let mut stmt = prepare_sql_statement(conn, &statement.sql)?;
374 bind_params(&mut stmt, &statement.params)?;
375 Ok(stmt)
376}
377
378fn execute_prepared_query(
379 mut stmt: rusqlite::Statement<'_>,
380) -> Result<Vec<SqlRow>, rusqlite::Error> {
381 let col_count = stmt.column_count();
382 let col_names: Vec<String> = (0..col_count)
383 .map(|i| stmt.column_name(i).unwrap_or("").to_string())
384 .collect();
385
386 let mut rows = Vec::new();
387 let mut raw_rows = stmt.raw_query();
388 while let Some(row) = raw_rows.next()? {
389 rows.push(row_to_sql_row(row, col_count, &col_names));
390 }
391 Ok(rows)
392}
393
394fn execute_prepared_query_row(
395 mut stmt: rusqlite::Statement<'_>,
396) -> Result<Option<SqlRow>, rusqlite::Error> {
397 let col_count = stmt.column_count();
398 let col_names: Vec<String> = (0..col_count)
399 .map(|i| stmt.column_name(i).unwrap_or("").to_string())
400 .collect();
401
402 let mut raw_rows = stmt.raw_query();
403 Ok(raw_rows
404 .next()?
405 .map(|row| row_to_sql_row(row, col_count, &col_names)))
406}
407
408fn execute_prepared_query_page(
409 mut stmt: rusqlite::Statement<'_>,
410 page: &PageRequest,
411) -> Result<Vec<SqlRow>, rusqlite::Error> {
412 if page.limit == 0 {
416 return Ok(Vec::new());
417 }
418
419 let col_count = stmt.column_count();
420 let col_names: Vec<String> = (0..col_count)
421 .map(|i| stmt.column_name(i).unwrap_or("").to_string())
422 .collect();
423
424 let mut rows = Vec::new();
425 let mut offset = page.offset;
426 let mut remaining = u64::from(page.limit);
427 let mut raw_rows = stmt.raw_query();
428 while remaining > 0 {
438 let Some(row) = raw_rows.next()? else {
439 break;
440 };
441 if offset > 0 {
442 offset -= 1;
443 continue;
444 }
445 rows.push(row_to_sql_row(row, col_count, &col_names));
446 remaining -= 1;
447 }
448 Ok(rows)
449}
450
451fn execute_query(
453 conn: &rusqlite::Connection,
454 statement: &SqlStatement,
455) -> Result<Vec<SqlRow>, rusqlite::Error> {
456 execute_prepared_query(prepare_bound_statement(conn, statement)?)
457}
458
459fn execute_query_row(
460 conn: &rusqlite::Connection,
461 statement: &SqlStatement,
462) -> Result<Option<SqlRow>, rusqlite::Error> {
463 execute_prepared_query_row(prepare_bound_statement(conn, statement)?)
464}
465
466fn execute_query_page(
467 conn: &rusqlite::Connection,
468 statement: &SqlStatement,
469 page: &PageRequest,
470) -> Result<Vec<SqlRow>, rusqlite::Error> {
471 execute_prepared_query_page(prepare_bound_statement(conn, statement)?, page)
472}
473
474fn statement_is_cancellable_read(stmt: &rusqlite::Statement<'_>, sql: &str) -> bool {
479 stmt.readonly() && transaction_control_head(sql).is_none()
480}
481
482const READER_STRUCTURAL_PRAGMAS: [&str; 8] = [
487 "table_info",
488 "table_xinfo",
489 "table_list",
490 "index_list",
491 "index_info",
492 "index_xinfo",
493 "foreign_key_list",
494 "integrity_check",
495];
496
497const READER_SETTING_PRAGMAS: [&str; 10] = [
505 "database_list",
506 "collation_list",
507 "function_list",
508 "compile_options",
509 "page_count",
510 "freelist_count",
511 "user_version",
512 "schema_version",
513 "journal_mode",
514 "page_size",
515];
516
517fn skip_balanced_parens(rest: &[u8]) -> Option<&[u8]> {
525 debug_assert_eq!(rest.first(), Some(&b'('));
526 let mut depth: u32 = 0;
527 let mut idx = 0;
528 loop {
529 match *rest.get(idx)? {
530 b'(' => {
531 depth += 1;
532 idx += 1;
533 }
534 b')' => {
535 depth -= 1;
536 idx += 1;
537 if depth == 0 {
538 return Some(&rest[idx..]);
539 }
540 }
541 quote @ (b'\'' | b'"' | b'`') => {
542 idx += 1;
543 loop {
544 match *rest.get(idx)? {
545 byte if byte == quote => {
546 idx += 1;
547 if rest.get(idx) == Some("e) {
548 idx += 1; } else {
550 break;
551 }
552 }
553 _ => idx += 1,
554 }
555 }
556 }
557 b'[' => {
558 idx += 1;
559 while *rest.get(idx)? != b']' {
560 idx += 1;
561 }
562 idx += 1;
563 }
564 b'-' if rest.get(idx + 1) == Some(&b'-') => {
565 idx += 2;
566 while idx < rest.len() && rest[idx] != b'\n' {
567 idx += 1;
568 }
569 }
570 b'/' if rest.get(idx + 1) == Some(&b'*') => {
571 idx += 2;
572 while idx + 1 < rest.len() && !(rest[idx] == b'*' && rest[idx + 1] == b'/') {
573 idx += 1;
574 }
575 idx = (idx + 2).min(rest.len());
576 }
577 _ => idx += 1,
578 }
579 }
580}
581
582fn skip_sqlite_identifier(rest: &[u8]) -> Option<&[u8]> {
590 match *rest.first()? {
591 quote @ (b'"' | b'`') => {
592 let mut idx = 1;
593 loop {
594 match *rest.get(idx)? {
595 byte if byte == quote => {
596 idx += 1;
597 if rest.get(idx) == Some("e) {
598 idx += 1; } else {
600 break;
601 }
602 }
603 _ => idx += 1,
604 }
605 }
606 Some(&rest[idx..])
607 }
608 b'[' => {
609 let mut idx = 1;
610 while *rest.get(idx)? != b']' {
611 idx += 1;
612 }
613 Some(&rest[idx + 1..])
614 }
615 _ => next_sqlite_token(rest).map(|(_, next)| next),
616 }
617}
618
619fn skip_common_table_expressions(tail: &[u8]) -> Option<&[u8]> {
627 let mut rest = skip_sqlite_empty_prefix(tail);
628 if let Some((word, next)) = next_sqlite_token(rest) {
629 if word.eq_ignore_ascii_case(b"RECURSIVE") {
630 rest = skip_sqlite_empty_prefix(next);
631 }
632 }
633 loop {
634 let next = skip_sqlite_identifier(rest)?;
636 rest = skip_sqlite_empty_prefix(next);
637 if rest.first() == Some(&b'(') {
639 rest = skip_sqlite_empty_prefix(skip_balanced_parens(rest)?);
640 }
641 let (as_keyword, next) = next_sqlite_token(rest)?;
642 if !as_keyword.eq_ignore_ascii_case(b"AS") {
643 return None;
644 }
645 rest = skip_sqlite_empty_prefix(next);
646 if let Some((word, next)) = next_sqlite_token(rest) {
648 if word.eq_ignore_ascii_case(b"MATERIALIZED") {
649 rest = skip_sqlite_empty_prefix(next);
650 } else if word.eq_ignore_ascii_case(b"NOT") {
651 let (materialized, next) = next_sqlite_token(skip_sqlite_empty_prefix(next))?;
652 if !materialized.eq_ignore_ascii_case(b"MATERIALIZED") {
653 return None;
654 }
655 rest = skip_sqlite_empty_prefix(next);
656 }
657 }
658 if rest.first() != Some(&b'(') {
660 return None;
661 }
662 rest = skip_sqlite_empty_prefix(skip_balanced_parens(rest)?);
663 if rest.first() == Some(&b',') {
664 rest = skip_sqlite_empty_prefix(&rest[1..]);
665 continue;
666 }
667 return Some(rest);
668 }
669}
670
671pub(crate) fn reader_capability_admits(sql: &str) -> Result<(), String> {
690 let rest = skip_sqlite_empty_prefix(sql.as_bytes());
691 let Some((head, tail)) = next_sqlite_token(rest) else {
692 return Ok(());
695 };
696 if head.eq_ignore_ascii_case(b"SELECT") || head.eq_ignore_ascii_case(b"VALUES") {
697 return Ok(());
698 }
699 if head.eq_ignore_ascii_case(b"WITH") {
700 let Some(after_ctes) = skip_common_table_expressions(tail) else {
701 return Err(
702 "WITH statement's common-table-expression list could not be parsed; refusing \
703 to admit it through the reader capability"
704 .into(),
705 );
706 };
707 return match next_sqlite_token(after_ctes) {
708 Some((main_head, _))
709 if main_head.eq_ignore_ascii_case(b"SELECT")
710 || main_head.eq_ignore_ascii_case(b"VALUES") =>
711 {
712 Ok(())
713 }
714 other => Err(format!(
715 "WITH ... {:?} is not admitted through the reader capability; only a \
716 read-only SELECT/VALUES body after the CTE list may run against a pooled \
717 reader connection",
718 other.map_or_else(
719 || "<none>".to_string(),
720 |(main_head, _)| String::from_utf8_lossy(main_head).into_owned()
721 )
722 )),
723 };
724 }
725 if head.eq_ignore_ascii_case(b"EXPLAIN") {
726 let mut rest = skip_sqlite_empty_prefix(tail);
727 if let Some((query, next)) = next_sqlite_token(rest) {
728 if query.eq_ignore_ascii_case(b"QUERY") {
729 let after_query = skip_sqlite_empty_prefix(next);
730 match next_sqlite_token(after_query) {
731 Some((plan, next2)) if plan.eq_ignore_ascii_case(b"PLAN") => {
732 rest = skip_sqlite_empty_prefix(next2);
733 }
734 _ => {
735 return Err(
736 "EXPLAIN QUERY must be followed by PLAN through the reader capability"
737 .into(),
738 );
739 }
740 }
741 }
742 }
743 return reader_capability_admits(&String::from_utf8_lossy(rest));
744 }
745 if head.eq_ignore_ascii_case(b"PRAGMA") {
746 return reader_capability_admits_pragma(tail);
747 }
748 Err(format!(
749 "statement head {:?} is not admitted through the reader capability; only \
750 SELECT/WITH/VALUES/EXPLAIN and an allow-listed set of read-only PRAGMA forms \
751 may run against a pooled reader connection",
752 String::from_utf8_lossy(head)
753 ))
754}
755
756fn reader_capability_admits_pragma(tail: &[u8]) -> Result<(), String> {
757 let rest = skip_sqlite_empty_prefix(tail);
758 let Some((mut name, mut after_name)) = next_sqlite_token(rest) else {
759 return Err("PRAGMA with no name is not admitted through the reader capability".into());
760 };
761 if after_name.first() == Some(&b'.') {
764 let (qualified_name, qualified_after) =
765 next_sqlite_token(&after_name[1..]).ok_or_else(|| {
766 "PRAGMA schema-qualifier with no pragma name is not admitted through the \
767 reader capability"
768 .to_string()
769 })?;
770 name = qualified_name;
771 after_name = qualified_after;
772 }
773 let after = skip_sqlite_empty_prefix(after_name);
774 let is_structural = READER_STRUCTURAL_PRAGMAS
775 .iter()
776 .any(|allowed| name.eq_ignore_ascii_case(allowed.as_bytes()));
777 let is_setting = READER_SETTING_PRAGMAS
778 .iter()
779 .any(|allowed| name.eq_ignore_ascii_case(allowed.as_bytes()));
780 if !is_structural && !is_setting {
781 return Err(format!(
782 "PRAGMA {:?} is not admitted through the reader capability",
783 String::from_utf8_lossy(name)
784 ));
785 }
786 if after.first() == Some(&b'=') {
787 return Err(format!(
788 "PRAGMA {:?} may not be assigned through the reader capability",
789 String::from_utf8_lossy(name)
790 ));
791 }
792 if after.first() == Some(&b'(') && !is_structural {
793 return Err(format!(
794 "PRAGMA {:?} may not carry an argument through the reader capability",
795 String::from_utf8_lossy(name)
796 ));
797 }
798 Ok(())
799}
800
801fn admit_reader_capability_sql(
806 statement: &SqlStatement,
807 transaction_control: Option<CachedReadTransactionControl>,
808 operation: &'static str,
809) -> khive_storage::types::StorageResult<()> {
810 if matches!(
816 transaction_control,
817 Some(CachedReadTransactionControl::BeginDeferred)
818 | Some(CachedReadTransactionControl::Finish(_))
819 ) {
820 return Ok(());
821 }
822 reader_capability_admits(&statement.sql).map_err(|message| StorageError::InvalidInput {
823 capability: StorageCapability::Sql,
824 operation: operation.into(),
825 message,
826 })
827}
828
829fn execute_query_interruptibly(
830 scope: &crate::read_cancellation::InterruptibleReadScope,
831 conn: &rusqlite::Connection,
832 statement: &SqlStatement,
833 operation: &'static str,
834 rollback_interrupted_transaction: bool,
835 interruptible: bool,
836) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
837 let stmt = prepare_bound_statement(conn, statement)
838 .map_err(|error| map_rusqlite_err(error, operation))?;
839 if interruptible && statement_is_cancellable_read(&stmt, &statement.sql) {
840 scope.run_with_interrupted_cleanup(
841 conn,
842 move || {
843 execute_prepared_query(stmt).map_err(|error| map_rusqlite_err(error, operation))
844 },
845 || {
846 rollback_interrupted_read_transaction(
847 conn,
848 operation,
849 rollback_interrupted_transaction,
850 )
851 },
852 )
853 } else {
854 scope.mark_write_committed()?;
855 execute_prepared_query(stmt).map_err(|error| map_rusqlite_err(error, operation))
856 }
857}
858
859fn execute_query_row_interruptibly(
860 scope: &crate::read_cancellation::InterruptibleReadScope,
861 conn: &rusqlite::Connection,
862 statement: &SqlStatement,
863 operation: &'static str,
864 rollback_interrupted_transaction: bool,
865 interruptible: bool,
866) -> khive_storage::types::StorageResult<Option<SqlRow>> {
867 let stmt = prepare_bound_statement(conn, statement)
868 .map_err(|error| map_rusqlite_err(error, operation))?;
869 if interruptible && statement_is_cancellable_read(&stmt, &statement.sql) {
870 scope.run_with_interrupted_cleanup(
871 conn,
872 move || {
873 execute_prepared_query_row(stmt).map_err(|error| map_rusqlite_err(error, operation))
874 },
875 || {
876 rollback_interrupted_read_transaction(
877 conn,
878 operation,
879 rollback_interrupted_transaction,
880 )
881 },
882 )
883 } else {
884 scope.mark_write_committed()?;
885 execute_prepared_query_row(stmt).map_err(|error| map_rusqlite_err(error, operation))
886 }
887}
888
889fn execute_query_page_interruptibly(
890 scope: &crate::read_cancellation::InterruptibleReadScope,
891 conn: &rusqlite::Connection,
892 statement: &SqlStatement,
893 page: &PageRequest,
894 operation: &'static str,
895 rollback_interrupted_transaction: bool,
896 interruptible: bool,
897) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
898 let stmt = prepare_bound_statement(conn, statement)
899 .map_err(|error| map_rusqlite_err(error, operation))?;
900 if interruptible && statement_is_cancellable_read(&stmt, &statement.sql) {
901 scope.run_with_interrupted_cleanup(
902 conn,
903 move || {
904 execute_prepared_query_page(stmt, page)
905 .map_err(|error| map_rusqlite_err(error, operation))
906 },
907 || {
908 rollback_interrupted_read_transaction(
909 conn,
910 operation,
911 rollback_interrupted_transaction,
912 )
913 },
914 )
915 } else {
916 scope.mark_write_committed()?;
917 execute_prepared_query_page(stmt, page).map_err(|error| map_rusqlite_err(error, operation))
918 }
919}
920
921fn rollback_interrupted_read_transaction(
922 conn: &rusqlite::Connection,
923 operation: &'static str,
924 enabled: bool,
925) -> khive_storage::types::StorageResult<()> {
926 if !enabled || conn.is_autocommit() {
927 return Ok(());
928 }
929 conn.execute_batch("ROLLBACK")
930 .map_err(|error| map_rusqlite_err(error, operation))?;
931 if conn.is_autocommit() {
932 Ok(())
933 } else {
934 Err(StorageError::Transaction {
935 operation: operation.into(),
936 message: "interrupted read transaction rollback did not restore autocommit".into(),
937 })
938 }
939}
940
941fn map_rusqlite_err(e: rusqlite::Error, op: &'static str) -> StorageError {
943 StorageError::driver(StorageCapability::Sql, op, e)
944}
945
946#[derive(Clone, Copy)]
957enum SlotTimeoutClass {
958 Admission,
959 ReaderContract,
960}
961
962async fn acquire_reader_handle_slot(
963 pool: &ConnectionPool,
964 operation: &'static str,
965 class: SlotTimeoutClass,
966) -> Result<OwnedSemaphorePermit, StorageError> {
967 let result = acquire_handle_slot(
968 pool.sql_bridge_reader_slots(),
969 pool.config().checkout_timeout,
970 operation,
971 class,
972 )
973 .await;
974 if matches!(
975 &result,
976 Err(StorageError::Timeout { .. } | StorageError::AdmissionTimeout { .. })
977 ) {
978 pool.record_reader_admission_timeout();
979 }
980 result
981}
982
983pub(crate) async fn acquire_in_memory_write_unit(
987 pool: &ConnectionPool,
988 operation: &'static str,
989) -> Result<OwnedSemaphorePermit, StorageError> {
990 acquire_handle_slot(
991 pool.sql_bridge_writer_slots(),
992 pool.config().checkout_timeout,
993 operation,
994 SlotTimeoutClass::Admission,
995 )
996 .await
997}
998
999async fn acquire_handle_slot(
1000 slots: Arc<Semaphore>,
1001 timeout: std::time::Duration,
1002 operation: &'static str,
1003 class: SlotTimeoutClass,
1004) -> Result<OwnedSemaphorePermit, StorageError> {
1005 tokio::time::timeout(timeout, slots.acquire_owned())
1006 .await
1007 .map_err(|_| match class {
1008 SlotTimeoutClass::Admission => StorageError::AdmissionTimeout {
1009 operation: operation.into(),
1010 timeout_ms: u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX),
1011 pool_identity: None,
1012 },
1013 SlotTimeoutClass::ReaderContract => StorageError::Timeout {
1014 operation: operation.into(),
1015 },
1016 })?
1017 .map_err(|error| StorageError::Pool {
1018 operation: operation.into(),
1019 message: error.to_string(),
1020 })
1021}
1022
1023fn open_standalone_reader(pool: &ConnectionPool) -> Result<rusqlite::Connection, StorageError> {
1028 pool.open_standalone_reader(StandaloneReaderPurpose::ExplicitSqlReadTransaction)
1029 .map_err(|error| StorageError::driver(StorageCapability::Sql, "open_reader", error))
1030}
1031
1032#[cfg(test)]
1033fn open_standalone_writer(pool: &ConnectionPool) -> Result<rusqlite::Connection, StorageError> {
1034 let conn = pool
1035 .open_standalone_writer()
1036 .map_err(|e| e.into_storage_error(StorageCapability::Sql, "open_writer"))?;
1037 configure_standalone_writer(pool, conn)
1038}
1039
1040fn open_admitted_standalone_writer(
1041 pool: &ConnectionPool,
1042) -> Result<rusqlite::Connection, StorageError> {
1043 let conn = pool
1044 .open_standalone_writer_for_admitted_operation()
1045 .map_err(|e| e.into_storage_error(StorageCapability::Sql, "open_writer"))?;
1046 configure_standalone_writer(pool, conn)
1047}
1048
1049fn configure_standalone_writer(
1050 pool: &ConnectionPool,
1051 conn: rusqlite::Connection,
1052) -> Result<rusqlite::Connection, StorageError> {
1053 let config = pool.config();
1054 conn.busy_timeout(config.busy_timeout)
1055 .map_err(|e| map_rusqlite_err(e, "open_writer"))?;
1056 conn.pragma_update(None, "cache_size", "-65536")
1057 .map_err(|e| map_rusqlite_err(e, "open_writer"))?;
1058 conn.pragma_update(None, "mmap_size", "1073741824")
1059 .map_err(|e| map_rusqlite_err(e, "open_writer"))?;
1060
1061 Ok(conn)
1062}
1063
1064async fn open_standalone_on_blocking<F>(
1071 pool: Arc<ConnectionPool>,
1072 slot: OwnedSemaphorePermit,
1073 operation: &'static str,
1074 open: F,
1075) -> khive_storage::types::StorageResult<(rusqlite::Connection, OwnedSemaphorePermit)>
1076where
1077 F: FnOnce(&ConnectionPool) -> Result<rusqlite::Connection, StorageError> + Send + 'static,
1078{
1079 tokio::task::spawn_blocking(move || open(&pool).map(|conn| (conn, slot)))
1080 .await
1081 .map_err(|e| StorageError::driver(StorageCapability::Sql, operation, e))?
1082}
1083
1084async fn open_standalone_reader_on_blocking(
1094 pool: Arc<ConnectionPool>,
1095 slot: OwnedSemaphorePermit,
1096) -> khive_storage::types::StorageResult<(rusqlite::Connection, OwnedSemaphorePermit)> {
1097 open_standalone_on_blocking(pool, slot, "open_reader", open_standalone_reader).await
1098}
1099
1100async fn open_standalone_writer_on_blocking(
1104 pool: Arc<ConnectionPool>,
1105 slot: OwnedSemaphorePermit,
1106) -> khive_storage::types::StorageResult<(rusqlite::Connection, OwnedSemaphorePermit)> {
1107 open_standalone_on_blocking(pool, slot, "open_writer", open_admitted_standalone_writer).await
1108}
1109
1110const CACHED_READ_TRANSACTION_LABEL: &str = "sql_bridge_cached_read_transaction";
1115
1116struct CachedReadTransaction {
1126 _slot: OwnedSemaphorePermit,
1127 _tx_handle: khive_storage::tx_registry::TxHandle,
1128 opened_at: Instant,
1134}
1135
1136struct StandaloneHandle {
1137 conn: rusqlite::Connection,
1138 _retained_slot: Option<OwnedSemaphorePermit>,
1143 read_transaction_slot: Option<CachedReadTransaction>,
1149}
1150
1151impl StandaloneHandle {
1152 fn is_cached_reader(&self) -> bool {
1155 self._retained_slot.is_none()
1156 }
1157
1158 fn has_read_transaction(&self) -> bool {
1159 self.read_transaction_slot.is_some()
1160 }
1161}
1162
1163struct SqliteReader {
1164 handle: Option<StandaloneHandle>,
1168 pool: Arc<ConnectionPool>,
1169 poisoned: bool,
1173}
1174
1175async fn open_explicit_read_transaction_handle(
1176 pool: Arc<ConnectionPool>,
1177) -> khive_storage::types::StorageResult<StandaloneHandle> {
1178 let open_slot = crate::await_request_read_phase(
1179 "sql_bridge.reader_open",
1180 acquire_reader_handle_slot(
1181 &pool,
1182 "sql_bridge.reader_open",
1183 SlotTimeoutClass::ReaderContract,
1184 ),
1185 )
1186 .await??;
1187 let (conn, open_slot) = crate::await_request_read_phase(
1188 "sql_bridge.reader_open",
1189 open_standalone_reader_on_blocking(pool, open_slot),
1190 )
1191 .await??;
1192 drop(open_slot);
1193 Ok(StandaloneHandle {
1194 conn,
1195 _retained_slot: None,
1196 read_transaction_slot: None,
1197 })
1198}
1199
1200impl SqliteReader {
1201 async fn use_explicit_transaction_handle(
1205 &mut self,
1206 transaction_control: Option<CachedReadTransactionControl>,
1207 operation: &'static str,
1208 ) -> khive_storage::types::StorageResult<bool> {
1209 if self.poisoned {
1210 return Err(StorageError::Pool {
1211 operation: operation.into(),
1212 message: "connection already consumed".into(),
1213 });
1214 }
1215 if self.handle.is_some() {
1216 return Ok(true);
1217 }
1218 match transaction_control {
1219 None => Ok(false),
1220 Some(CachedReadTransactionControl::BeginDeferred) => {
1221 self.handle =
1222 Some(open_explicit_read_transaction_handle(Arc::clone(&self.pool)).await?);
1223 Ok(true)
1224 }
1225 Some(CachedReadTransactionControl::Finish(keyword))
1226 | Some(CachedReadTransactionControl::Unsupported(keyword)) => {
1227 Err(StorageError::InvalidInput {
1228 capability: StorageCapability::Sql,
1229 operation: operation.into(),
1230 message: format!(
1231 "cached read-only handle has no admitted transaction for transaction \
1232 control ({keyword})"
1233 ),
1234 })
1235 }
1236 }
1237 }
1238
1239 fn close_inactive_transaction_handle(&mut self) {
1245 if self
1246 .handle
1247 .as_ref()
1248 .is_some_and(|handle| handle.is_cached_reader() && !handle.has_read_transaction())
1249 {
1250 drop(self.handle.take());
1251 }
1252 }
1253}
1254
1255async fn execute_standalone_read<R, F>(
1269 handle: &mut Option<StandaloneHandle>,
1270 pool: Arc<ConnectionPool>,
1271 operation: &'static str,
1272 transaction_control: Option<CachedReadTransactionControl>,
1273 read: F,
1274) -> khive_storage::types::StorageResult<R>
1275where
1276 R: Send + 'static,
1277 F: FnOnce(
1278 &crate::read_cancellation::InterruptibleReadScope,
1279 &rusqlite::Connection,
1280 bool,
1281 bool,
1282 ) -> khive_storage::types::StorageResult<R>
1283 + Send
1284 + 'static,
1285{
1286 if handle.is_none() {
1287 return Err(StorageError::Pool {
1288 operation: operation.into(),
1289 message: "connection already consumed".into(),
1290 });
1291 }
1292 let active_read_transaction = handle
1293 .as_ref()
1294 .is_some_and(|handle| handle.is_cached_reader() && handle.has_read_transaction());
1295 let completion_preserving_writer_transaction = handle
1296 .as_ref()
1297 .is_some_and(|handle| !handle.is_cached_reader() && !handle.conn.is_autocommit());
1298 let mut operation_slot = if active_read_transaction {
1299 None
1300 } else if completion_preserving_writer_transaction {
1301 Some(acquire_reader_handle_slot(&pool, operation, SlotTimeoutClass::ReaderContract).await?)
1305 } else {
1306 Some(
1307 crate::await_request_read_phase(
1308 operation,
1309 acquire_reader_handle_slot(&pool, operation, SlotTimeoutClass::ReaderContract),
1310 )
1311 .await??,
1312 )
1313 };
1314 let Some(owned_handle) = handle.take() else {
1315 return Err(StorageError::Pool {
1316 operation: operation.into(),
1317 message: "connection already consumed".into(),
1318 });
1319 };
1320 let origin = pool.origin();
1321 let read_tx_max_age = pool.config().read_tx_max_age;
1322 let (owned_handle, result) = crate::read_cancellation::run_interruptible_read(
1323 StorageCapability::Sql,
1324 operation,
1325 move |scope| {
1326 let mut owned_handle = owned_handle;
1327 let cached_reader = owned_handle.is_cached_reader();
1328 let entered_with_transaction = owned_handle.has_read_transaction();
1329 let entered_autocommit = owned_handle.conn.is_autocommit();
1330 let mut restore_handle = true;
1331 let mut result = if cached_reader && entered_with_transaction && entered_autocommit {
1332 drop(owned_handle.read_transaction_slot.take());
1335 Err(StorageError::InvalidInput {
1336 capability: StorageCapability::Sql,
1337 operation: operation.into(),
1338 message: "cached read-only handle retained transaction admission after SQLite \
1339 had already returned to autocommit; the stale permit was released"
1340 .into(),
1341 })
1342 } else if cached_reader && !entered_with_transaction && !entered_autocommit {
1343 Err(StorageError::InvalidInput {
1344 capability: StorageCapability::Sql,
1345 operation: operation.into(),
1346 message: "cached read-only handle entered the operation outside autocommit; \
1347 its transaction was rolled back before releasing the reader permit"
1348 .into(),
1349 })
1350 } else if cached_reader
1351 && entered_with_transaction
1352 && owned_handle
1353 .read_transaction_slot
1354 .as_ref()
1355 .is_some_and(|tx| tx.opened_at.elapsed() >= read_tx_max_age)
1356 {
1357 crate::checkpoint::note_read_tx_max_age_eviction();
1362 match owned_handle.conn.execute_batch("ROLLBACK") {
1363 Ok(()) if owned_handle.conn.is_autocommit() => {
1364 drop(owned_handle.read_transaction_slot.take());
1365 Err(StorageError::ReadTransactionAgeEvicted {
1366 operation: operation.into(),
1367 max_age_secs: read_tx_max_age.as_secs(),
1368 })
1369 }
1370 Ok(()) => {
1371 restore_handle = false;
1372 Err(StorageError::ReadTransactionAgeEvictionCleanupFailed {
1373 operation: operation.into(),
1374 max_age_secs: read_tx_max_age.as_secs(),
1375 message: "rollback did not restore autocommit".into(),
1376 })
1377 }
1378 Err(error) => {
1379 restore_handle = false;
1380 Err(StorageError::ReadTransactionAgeEvictionCleanupFailed {
1381 operation: operation.into(),
1382 max_age_secs: read_tx_max_age.as_secs(),
1383 message: format!("rollback failed: {error}"),
1384 })
1385 }
1386 }
1387 } else if cached_reader && entered_with_transaction {
1388 match transaction_control {
1389 None | Some(CachedReadTransactionControl::Finish(_)) => {
1390 read(scope, &owned_handle.conn, true, true)
1391 }
1392 Some(CachedReadTransactionControl::BeginDeferred) => {
1393 Err(StorageError::InvalidInput {
1394 capability: StorageCapability::Sql,
1395 operation: operation.into(),
1396 message: "cached read-only handle already owns an admitted read \
1397 transaction; nested BEGIN is not supported"
1398 .into(),
1399 })
1400 }
1401 Some(CachedReadTransactionControl::Unsupported(keyword)) => {
1402 Err(StorageError::InvalidInput {
1403 capability: StorageCapability::Sql,
1404 operation: operation.into(),
1405 message: format!(
1406 "cached read-only transaction does not support nested or \
1407 write-locking transaction control ({keyword})"
1408 ),
1409 })
1410 }
1411 }
1412 } else if cached_reader {
1413 match transaction_control {
1414 None | Some(CachedReadTransactionControl::BeginDeferred) => {
1415 read(scope, &owned_handle.conn, false, true)
1416 }
1417 Some(CachedReadTransactionControl::Finish(keyword))
1418 | Some(CachedReadTransactionControl::Unsupported(keyword)) => {
1419 Err(StorageError::InvalidInput {
1420 capability: StorageCapability::Sql,
1421 operation: operation.into(),
1422 message: format!(
1423 "cached read-only handle has no admitted transaction for \
1424 transaction control ({keyword})"
1425 ),
1426 })
1427 }
1428 }
1429 } else {
1430 read(scope, &owned_handle.conn, false, entered_autocommit)
1431 };
1432
1433 if scope.cleanup_failed() {
1434 restore_handle = false;
1439 }
1440
1441 if cached_reader
1447 && matches!(result, Err(StorageError::Timeout { .. }))
1448 && !owned_handle.conn.is_autocommit()
1449 {
1450 match owned_handle.conn.execute_batch("ROLLBACK") {
1451 Ok(()) if owned_handle.conn.is_autocommit() => {
1452 drop(owned_handle.read_transaction_slot.take());
1453 }
1454 Ok(()) => {
1455 restore_handle = false;
1456 result = Err(StorageError::Transaction {
1457 operation: operation.into(),
1458 message:
1459 "interrupted read transaction rollback did not restore autocommit; \
1460 the connection was discarded"
1461 .into(),
1462 });
1463 }
1464 Err(error) => {
1465 restore_handle = false;
1466 result = Err(StorageError::Transaction {
1467 operation: operation.into(),
1468 message: format!(
1469 "failed to roll back interrupted read transaction ({error}); \
1470 the connection was discarded"
1471 ),
1472 });
1473 }
1474 }
1475 }
1476
1477 if cached_reader && entered_with_transaction {
1478 if owned_handle.conn.is_autocommit() {
1479 drop(owned_handle.read_transaction_slot.take());
1483 if result.is_ok()
1484 && !matches!(
1485 transaction_control,
1486 Some(CachedReadTransactionControl::Finish(_))
1487 )
1488 {
1489 result = Err(StorageError::InvalidInput {
1490 capability: StorageCapability::Sql,
1491 operation: operation.into(),
1492 message: "cached read-only operation unexpectedly ended its admitted \
1493 transaction; reader admission was released after autocommit"
1494 .into(),
1495 });
1496 }
1497 } else if result.is_ok()
1498 && matches!(
1499 transaction_control,
1500 Some(CachedReadTransactionControl::Finish(_))
1501 )
1502 {
1503 result = Err(StorageError::InvalidInput {
1504 capability: StorageCapability::Sql,
1505 operation: operation.into(),
1506 message: "transaction-ending control completed but the cached reader \
1507 remained outside autocommit; its reader permit remains retained"
1508 .into(),
1509 });
1510 }
1511 } else if cached_reader
1512 && entered_autocommit
1513 && matches!(
1514 transaction_control,
1515 Some(CachedReadTransactionControl::BeginDeferred)
1516 )
1517 && result.is_ok()
1518 {
1519 if owned_handle.conn.is_autocommit() {
1520 result = Err(StorageError::InvalidInput {
1521 capability: StorageCapability::Sql,
1522 operation: operation.into(),
1523 message: "deferred BEGIN completed without opening a read transaction"
1524 .into(),
1525 });
1526 } else {
1527 match operation_slot.take() {
1528 Some(slot) => {
1529 let tx_handle = khive_storage::tx_registry::register_scoped(
1530 Some(CACHED_READ_TRANSACTION_LABEL.to_string()),
1531 origin.clone(),
1532 );
1533 owned_handle.read_transaction_slot = Some(CachedReadTransaction {
1534 _slot: slot,
1535 _tx_handle: tx_handle,
1536 opened_at: Instant::now(),
1537 });
1538 }
1539 None => {
1540 result = Err(StorageError::Pool {
1541 operation: operation.into(),
1542 message: "successful cached-reader BEGIN had no operation permit; \
1543 its transaction was rolled back before returning"
1544 .into(),
1545 });
1546 }
1547 }
1548 }
1549 }
1550
1551 if cached_reader
1554 && owned_handle.read_transaction_slot.is_none()
1555 && !owned_handle.conn.is_autocommit()
1556 {
1557 match owned_handle.conn.execute_batch("ROLLBACK") {
1558 Ok(()) if owned_handle.conn.is_autocommit() => {
1559 if result.is_ok() {
1560 result = Err(StorageError::InvalidInput {
1561 capability: StorageCapability::Sql,
1562 operation: operation.into(),
1563 message: "cached read-only operation left the connection outside \
1564 autocommit; its transaction was rolled back before \
1565 releasing the reader permit"
1566 .into(),
1567 });
1568 }
1569 }
1570 Ok(()) => {
1571 restore_handle = false;
1572 result = Err(StorageError::Transaction {
1573 operation: operation.into(),
1574 message: "ROLLBACK completed but the cached reader remained outside \
1575 autocommit; the connection was discarded before releasing \
1576 the reader permit"
1577 .into(),
1578 });
1579 }
1580 Err(error) => {
1581 restore_handle = false;
1582 result = Err(StorageError::Transaction {
1583 operation: operation.into(),
1584 message: format!(
1585 "failed to roll back a cached reader outside autocommit ({error}); \
1586 the connection was discarded before releasing the reader permit"
1587 ),
1588 });
1589 }
1590 }
1591 }
1592
1593 let owned_handle = if restore_handle {
1594 Some(owned_handle)
1595 } else {
1596 drop(owned_handle);
1599 None
1600 };
1601 drop(operation_slot);
1605 Ok((owned_handle, result))
1606 },
1607 )
1608 .await?;
1609 *handle = owned_handle;
1610 result
1611}
1612
1613#[async_trait]
1614impl khive_storage::SqlReader for SqliteReader {
1615 async fn query_row(
1616 &mut self,
1617 statement: SqlStatement,
1618 ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
1619 let transaction_control = cached_read_transaction_control(&statement.sql);
1620 admit_reader_capability_sql(&statement, transaction_control, "query_row")?;
1621 if !self
1622 .use_explicit_transaction_handle(transaction_control, "query_row")
1623 .await?
1624 {
1625 return run_pool_reader_query(
1626 Arc::clone(&self.pool),
1627 "query_row",
1628 move |scope, conn| {
1629 execute_query_row_interruptibly(
1630 scope,
1631 conn,
1632 &statement,
1633 "query_row",
1634 false,
1635 true,
1636 )
1637 },
1638 )
1639 .await;
1640 }
1641 let result = execute_standalone_read(
1642 &mut self.handle,
1643 Arc::clone(&self.pool),
1644 "query_row",
1645 transaction_control,
1646 move |scope, conn, rollback, interruptible| {
1647 execute_query_row_interruptibly(
1648 scope,
1649 conn,
1650 &statement,
1651 "query_row",
1652 rollback,
1653 interruptible,
1654 )
1655 },
1656 )
1657 .await;
1658 if self.handle.is_none() {
1659 self.poisoned = true;
1660 }
1661 self.close_inactive_transaction_handle();
1662 result
1663 }
1664
1665 async fn query_all(
1666 &mut self,
1667 statement: SqlStatement,
1668 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1669 let transaction_control = cached_read_transaction_control(&statement.sql);
1670 admit_reader_capability_sql(&statement, transaction_control, "query_all")?;
1671 if !self
1672 .use_explicit_transaction_handle(transaction_control, "query_all")
1673 .await?
1674 {
1675 return run_pool_reader_query(
1676 Arc::clone(&self.pool),
1677 "query_all",
1678 move |scope, conn| {
1679 execute_query_interruptibly(scope, conn, &statement, "query_all", false, true)
1680 },
1681 )
1682 .await;
1683 }
1684 let result = execute_standalone_read(
1685 &mut self.handle,
1686 Arc::clone(&self.pool),
1687 "query_all",
1688 transaction_control,
1689 move |scope, conn, rollback, interruptible| {
1690 execute_query_interruptibly(
1691 scope,
1692 conn,
1693 &statement,
1694 "query_all",
1695 rollback,
1696 interruptible,
1697 )
1698 },
1699 )
1700 .await;
1701 if self.handle.is_none() {
1702 self.poisoned = true;
1703 }
1704 self.close_inactive_transaction_handle();
1705 result
1706 }
1707
1708 async fn query_page(
1709 &mut self,
1710 statement: SqlStatement,
1711 page: PageRequest,
1712 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1713 let transaction_control = cached_read_transaction_control(&statement.sql);
1714 admit_reader_capability_sql(&statement, transaction_control, "query_page")?;
1715 if !self
1716 .use_explicit_transaction_handle(transaction_control, "query_page")
1717 .await?
1718 {
1719 return run_pool_reader_query(
1720 Arc::clone(&self.pool),
1721 "query_page",
1722 move |scope, conn| {
1723 execute_query_page_interruptibly(
1724 scope,
1725 conn,
1726 &statement,
1727 &page,
1728 "query_page",
1729 false,
1730 true,
1731 )
1732 },
1733 )
1734 .await;
1735 }
1736 let result = execute_standalone_read(
1737 &mut self.handle,
1738 Arc::clone(&self.pool),
1739 "query_page",
1740 transaction_control,
1741 move |scope, conn, rollback, interruptible| {
1742 execute_query_page_interruptibly(
1743 scope,
1744 conn,
1745 &statement,
1746 &page,
1747 "query_page",
1748 rollback,
1749 interruptible,
1750 )
1751 },
1752 )
1753 .await;
1754 if self.handle.is_none() {
1755 self.poisoned = true;
1756 }
1757 self.close_inactive_transaction_handle();
1758 result
1759 }
1760
1761 async fn query_scalar(
1762 &mut self,
1763 statement: SqlStatement,
1764 ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
1765 let row = self.query_row(statement).await?;
1766 Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
1767 }
1768
1769 async fn explain(
1770 &mut self,
1771 statement: SqlStatement,
1772 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1773 let explain_stmt = SqlStatement {
1774 sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
1775 params: statement.params,
1776 label: statement.label,
1777 };
1778 self.query_all(explain_stmt).await
1779 }
1780}
1781
1782struct SqliteWriter {
1787 observe_direct_errors: bool,
1789 handle: Option<StandaloneHandle>,
1798 writer_task: Option<crate::writer_task::WriterTaskHandle>,
1805 origin: khive_storage::tx_registry::TxOrigin,
1808 db: String,
1812 pool: Arc<ConnectionPool>,
1814 event_rows: Option<Arc<AtomicEventRows>>,
1815 held_lease: Option<crate::disk_guard::DetachedVolumeLease>,
1819}
1820
1821fn execute_top_level_maintenance(
1822 pool: &ConnectionPool,
1823 conn: &rusqlite::Connection,
1824 maintenance: TopLevelMaintenance,
1825) -> rusqlite::Result<()> {
1826 match maintenance {
1827 TopLevelMaintenance::WalCheckpointTruncate => {
1828 let result = conn.query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |row| {
1829 Ok((
1830 row.get::<_, i64>(0)?,
1831 row.get::<_, i64>(1)?,
1832 row.get::<_, i64>(2)?,
1833 ))
1834 });
1835 crate::checkpoint::record_checkpoint_run_result(pool, result.as_ref().ok().copied());
1836 result.map(|_| ())
1837 }
1838 TopLevelMaintenance::Vacuum => conn.execute_batch(maintenance.as_sql()),
1839 }
1840}
1841
1842impl SqliteWriter {
1843 fn require_standalone_handle(
1849 &self,
1850 operation: &'static str,
1851 ) -> khive_storage::types::StorageResult<()> {
1852 if self.handle.is_none() {
1853 return Err(StorageError::Pool {
1854 operation: operation.into(),
1855 message: "connection already consumed".into(),
1856 });
1857 }
1858 Ok(())
1859 }
1860
1861 async fn use_queue_read_transaction_handle(
1862 &mut self,
1863 transaction_control: Option<CachedReadTransactionControl>,
1864 operation: &'static str,
1865 ) -> khive_storage::types::StorageResult<bool> {
1866 if self.handle.is_some() {
1867 return Ok(true);
1868 }
1869 match transaction_control {
1870 None => Ok(false),
1871 Some(CachedReadTransactionControl::BeginDeferred) => {
1872 self.handle =
1873 Some(open_explicit_read_transaction_handle(Arc::clone(&self.pool)).await?);
1874 Ok(true)
1875 }
1876 Some(CachedReadTransactionControl::Finish(keyword))
1877 | Some(CachedReadTransactionControl::Unsupported(keyword)) => {
1878 Err(StorageError::InvalidInput {
1879 capability: StorageCapability::Sql,
1880 operation: operation.into(),
1881 message: format!(
1882 "cached read-only handle has no admitted transaction for transaction \
1883 control ({keyword})"
1884 ),
1885 })
1886 }
1887 }
1888 }
1889
1890 fn close_inactive_queue_read_transaction_handle(&mut self) {
1891 if self
1892 .handle
1893 .as_ref()
1894 .is_some_and(|handle| handle.is_cached_reader() && !handle.has_read_transaction())
1895 {
1896 drop(self.handle.take());
1897 }
1898 }
1899}
1900
1901#[async_trait]
1902impl khive_storage::SqlReader for SqliteWriter {
1903 async fn query_row(
1904 &mut self,
1905 statement: SqlStatement,
1906 ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
1907 if self.writer_task.is_some() {
1908 let transaction_control = cached_read_transaction_control(&statement.sql);
1909 if !self
1910 .use_queue_read_transaction_handle(transaction_control, "writer.query_row")
1911 .await?
1912 {
1913 admit_reader_capability_sql(&statement, transaction_control, "writer.query_row")?;
1914 return run_pool_reader_query(
1915 Arc::clone(&self.pool),
1916 "writer.query_row",
1917 move |scope, conn| {
1918 execute_query_row_interruptibly(
1919 scope,
1920 conn,
1921 &statement,
1922 "writer.query_row",
1923 false,
1924 true,
1925 )
1926 },
1927 )
1928 .await;
1929 }
1930 let result = execute_standalone_read(
1931 &mut self.handle,
1932 Arc::clone(&self.pool),
1933 "writer.query_row",
1934 transaction_control,
1935 move |scope, conn, rollback, interruptible| {
1936 execute_query_row_interruptibly(
1937 scope,
1938 conn,
1939 &statement,
1940 "writer.query_row",
1941 rollback,
1942 interruptible,
1943 )
1944 },
1945 )
1946 .await;
1947 self.close_inactive_queue_read_transaction_handle();
1948 return result;
1949 }
1950 let transaction_control = cached_read_transaction_control(&statement.sql);
1951 execute_standalone_read(
1952 &mut self.handle,
1953 Arc::clone(&self.pool),
1954 "writer.query_row",
1955 transaction_control,
1956 move |scope, conn, rollback, interruptible| {
1957 execute_query_row_interruptibly(
1958 scope,
1959 conn,
1960 &statement,
1961 "writer.query_row",
1962 rollback,
1963 interruptible,
1964 )
1965 },
1966 )
1967 .await
1968 }
1969
1970 async fn query_all(
1971 &mut self,
1972 statement: SqlStatement,
1973 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
1974 if self.writer_task.is_some() {
1975 let transaction_control = cached_read_transaction_control(&statement.sql);
1976 if !self
1977 .use_queue_read_transaction_handle(transaction_control, "writer.query_all")
1978 .await?
1979 {
1980 admit_reader_capability_sql(&statement, transaction_control, "writer.query_all")?;
1981 return run_pool_reader_query(
1982 Arc::clone(&self.pool),
1983 "writer.query_all",
1984 move |scope, conn| {
1985 execute_query_interruptibly(
1986 scope,
1987 conn,
1988 &statement,
1989 "writer.query_all",
1990 false,
1991 true,
1992 )
1993 },
1994 )
1995 .await;
1996 }
1997 let result = execute_standalone_read(
1998 &mut self.handle,
1999 Arc::clone(&self.pool),
2000 "writer.query_all",
2001 transaction_control,
2002 move |scope, conn, rollback, interruptible| {
2003 execute_query_interruptibly(
2004 scope,
2005 conn,
2006 &statement,
2007 "writer.query_all",
2008 rollback,
2009 interruptible,
2010 )
2011 },
2012 )
2013 .await;
2014 self.close_inactive_queue_read_transaction_handle();
2015 return result;
2016 }
2017 let transaction_control = cached_read_transaction_control(&statement.sql);
2018 execute_standalone_read(
2019 &mut self.handle,
2020 Arc::clone(&self.pool),
2021 "writer.query_all",
2022 transaction_control,
2023 move |scope, conn, rollback, interruptible| {
2024 execute_query_interruptibly(
2025 scope,
2026 conn,
2027 &statement,
2028 "writer.query_all",
2029 rollback,
2030 interruptible,
2031 )
2032 },
2033 )
2034 .await
2035 }
2036
2037 async fn query_page(
2038 &mut self,
2039 statement: SqlStatement,
2040 page: PageRequest,
2041 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2042 if self.writer_task.is_some() {
2043 let transaction_control = cached_read_transaction_control(&statement.sql);
2044 if !self
2045 .use_queue_read_transaction_handle(transaction_control, "writer.query_page")
2046 .await?
2047 {
2048 admit_reader_capability_sql(&statement, transaction_control, "writer.query_page")?;
2049 return run_pool_reader_query(
2050 Arc::clone(&self.pool),
2051 "writer.query_page",
2052 move |scope, conn| {
2053 execute_query_page_interruptibly(
2054 scope,
2055 conn,
2056 &statement,
2057 &page,
2058 "writer.query_page",
2059 false,
2060 true,
2061 )
2062 },
2063 )
2064 .await;
2065 }
2066 let result = execute_standalone_read(
2067 &mut self.handle,
2068 Arc::clone(&self.pool),
2069 "writer.query_page",
2070 transaction_control,
2071 move |scope, conn, rollback, interruptible| {
2072 execute_query_page_interruptibly(
2073 scope,
2074 conn,
2075 &statement,
2076 &page,
2077 "writer.query_page",
2078 rollback,
2079 interruptible,
2080 )
2081 },
2082 )
2083 .await;
2084 self.close_inactive_queue_read_transaction_handle();
2085 return result;
2086 }
2087 let transaction_control = cached_read_transaction_control(&statement.sql);
2088 execute_standalone_read(
2089 &mut self.handle,
2090 Arc::clone(&self.pool),
2091 "writer.query_page",
2092 transaction_control,
2093 move |scope, conn, rollback, interruptible| {
2094 execute_query_page_interruptibly(
2095 scope,
2096 conn,
2097 &statement,
2098 &page,
2099 "writer.query_page",
2100 rollback,
2101 interruptible,
2102 )
2103 },
2104 )
2105 .await
2106 }
2107
2108 async fn query_scalar(
2109 &mut self,
2110 statement: SqlStatement,
2111 ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
2112 let row = khive_storage::SqlReader::query_row(self, statement).await?;
2113 Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
2114 }
2115
2116 async fn explain(
2117 &mut self,
2118 statement: SqlStatement,
2119 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2120 let explain_stmt = SqlStatement {
2121 sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
2122 params: statement.params,
2123 label: statement.label,
2124 };
2125 khive_storage::SqlReader::query_all(self, explain_stmt).await
2126 }
2127}
2128
2129#[async_trait]
2130impl khive_storage::SqlWriter for SqliteWriter {
2131 async fn execute(
2132 &mut self,
2133 statement: SqlStatement,
2134 ) -> khive_storage::types::StorageResult<u64> {
2135 if let Some(writer_task) = self.writer_task.clone() {
2145 let event_rows = self.event_rows.clone();
2146 return writer_task
2147 .send_bounded(move |conn| {
2148 let mut stmt = prepare_cached_sql_statement(conn, &statement.sql)
2149 .map_err(|e| map_rusqlite_err(e, "execute"))?;
2150 bind_params(&mut stmt, &statement.params)
2151 .map_err(|e| map_rusqlite_err(e, "execute"))?;
2152 let affected = stmt
2153 .raw_execute()
2154 .map_err(|e| map_rusqlite_err(e, "execute"))?;
2155 if let Some(event_rows) = event_rows.as_deref() {
2156 event_rows.observe(&statement, affected as u64);
2157 }
2158 Ok(affected as u64)
2159 })
2160 .await;
2161 }
2162
2163 self.require_standalone_handle("execute")?;
2164 let unit_holds_lease = self.held_lease.is_some();
2165 if !unit_holds_lease {
2171 if let Some(keyword) = transaction_control_head(&statement.sql) {
2172 return Err(StorageError::InvalidInput {
2173 capability: StorageCapability::Sql,
2174 operation: "execute".into(),
2175 message: format!(
2176 "statement is transaction control ({keyword}); a standalone \
2177 statement holds the volume lease only for the call — use \
2178 atomic_unit to run statements as one transaction"
2179 ),
2180 });
2181 }
2182 }
2183 let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
2184 operation: "execute".into(),
2185 message: "connection already consumed".into(),
2186 })?;
2187 let event_rows = self.event_rows.clone();
2188 let pool = Arc::clone(&self.pool);
2189 let (handle, result) = tokio::task::spawn_blocking(move || {
2190 run_standalone_statement(handle, pool, unit_holds_lease, statement, event_rows)
2191 })
2192 .await
2193 .map_err(|e| StorageError::driver(StorageCapability::Sql, "execute", e))?;
2194 self.handle = Some(handle);
2195 let affected = result.map_err(|failure| match failure {
2196 StandaloneWriteError::Refused(error) => error,
2197 StandaloneWriteError::Sql(error) => self.map_direct_error(error, "execute"),
2198 })?;
2199 Ok(affected as u64)
2200 }
2201
2202 async fn execute_batch(
2203 &mut self,
2204 statements: Vec<SqlStatement>,
2205 ) -> khive_storage::types::StorageResult<u64> {
2206 reject_transaction_control_statements(&statements, "execute_batch")?;
2222 if let Some(writer_task) = self.writer_task.clone() {
2223 let event_rows = self.event_rows.clone();
2224 return writer_task
2225 .send_bounded(move |conn| {
2226 let prepared = prepare_batch_statements(conn, &statements)
2227 .map_err(|e| map_rusqlite_err(e, "execute_batch"))?;
2228 execute_prepared_batch(conn, prepared, &statements, event_rows.as_deref())
2229 .map_err(|e| map_rusqlite_err(e, "execute_batch"))
2230 })
2231 .await;
2232 }
2233
2234 self.require_standalone_handle("execute_batch")?;
2235 let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
2236 operation: "execute_batch".into(),
2237 message: "connection already consumed".into(),
2238 })?;
2239 let origin = self.origin.clone();
2240 let event_rows = self.event_rows.clone();
2241 let pool = Arc::clone(&self.pool);
2242 let unit_holds_lease = self.held_lease.is_some();
2243 let (handle, result) = tokio::task::spawn_blocking(move || {
2244 run_standalone_batch(
2245 handle,
2246 pool,
2247 unit_holds_lease,
2248 statements,
2249 origin,
2250 event_rows,
2251 )
2252 })
2253 .await
2254 .map_err(|e| StorageError::driver(StorageCapability::Sql, "execute_batch", e))?;
2255 self.handle = handle;
2256 result.map_err(|failure| match failure {
2257 StandaloneWriteError::Refused(error) => error,
2258 StandaloneWriteError::Sql(failure) => self.map_direct_batch_failure(failure),
2259 })
2260 }
2261
2262 async fn execute_script(&mut self, script: String) -> khive_storage::types::StorageResult<()> {
2263 if let Some(writer_task) = self.writer_task.clone() {
2276 return writer_task
2277 .send_bounded(move |conn| {
2278 conn.execute_batch(&script)
2279 .map_err(|e| map_rusqlite_err(e, "execute_script"))
2280 })
2281 .await;
2282 }
2283
2284 self.require_standalone_handle("execute_script")?;
2285 let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
2286 operation: "execute_script".into(),
2287 message: "connection already consumed".into(),
2288 })?;
2289 let pool = Arc::clone(&self.pool);
2290 let unit_holds_lease = self.held_lease.is_some();
2291 let (handle, result) = tokio::task::spawn_blocking(move || {
2292 run_standalone_script(handle, pool, unit_holds_lease, script)
2293 })
2294 .await
2295 .map_err(|e| StorageError::driver(StorageCapability::Sql, "execute_script", e))?;
2296 self.handle = handle;
2297 result.map_err(|failure| match failure {
2298 StandaloneWriteError::Refused(error) => error,
2299 StandaloneWriteError::Sql(error) => self.map_direct_error(error, "execute_script"),
2300 })
2301 }
2302
2303 async fn execute_script_top_level(
2304 &mut self,
2305 maintenance: TopLevelMaintenance,
2306 ) -> khive_storage::types::StorageResult<()> {
2307 if let Some(writer_task) = self.writer_task.clone() {
2317 let pool = Arc::clone(&self.pool);
2318 let execute = move |conn: &rusqlite::Connection| {
2319 execute_top_level_maintenance(&pool, conn, maintenance)
2320 .map_err(|e| map_rusqlite_err(e, "execute_script_top_level"))
2321 };
2322 return if maintenance == TopLevelMaintenance::WalCheckpointTruncate {
2323 writer_task.send_checkpoint_bounded(execute).await
2324 } else {
2325 writer_task.send_vacuum_bounded(execute).await
2326 };
2327 }
2328
2329 let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
2333 operation: "execute_script_top_level".into(),
2334 message: "connection already consumed".into(),
2335 })?;
2336 let pool = Arc::clone(&self.pool);
2337 let unit_holds_lease = self.held_lease.is_some();
2338 let (handle, result) = tokio::task::spawn_blocking(move || {
2339 run_standalone_top_level(handle, pool, unit_holds_lease, maintenance)
2340 })
2341 .await
2342 .map_err(|e| StorageError::driver(StorageCapability::Sql, "execute_script_top_level", e))?;
2343 self.handle = Some(handle);
2344 result.map_err(|failure| match failure {
2345 StandaloneWriteError::Refused(error) => error,
2346 StandaloneWriteError::Sql(error) => {
2347 self.map_direct_error(error, "execute_script_top_level")
2348 }
2349 })
2350 }
2351}
2352
2353async fn run_pool_reader_query<T, F>(
2358 pool: Arc<ConnectionPool>,
2359 operation: &'static str,
2360 query: F,
2361) -> khive_storage::types::StorageResult<T>
2362where
2363 T: Send + 'static,
2364 F: FnOnce(
2365 &crate::read_cancellation::InterruptibleReadScope,
2366 &rusqlite::Connection,
2367 ) -> khive_storage::types::StorageResult<T>
2368 + Send
2369 + 'static,
2370{
2371 let admission = pool
2373 .acquire_reader_admission(StorageCapability::Sql, operation)
2374 .await?;
2375 crate::read_cancellation::run_interruptible_read(
2376 StorageCapability::Sql,
2377 operation,
2378 move |scope| {
2379 let mut guard = pool.resolve_reader_checkout(
2383 StorageCapability::Sql,
2384 operation,
2385 pool.reader_with_admission(admission, || scope.should_stop()),
2386 )?;
2387 guard.mark_dirty();
2392 let result = scope.with_pooled_reader(&mut guard, |conn| query(scope, conn));
2393 if let Err(error) = &result {
2394 pool.record_reader_query_error(error);
2395 }
2396 result
2397 },
2398 )
2399 .await
2400}
2401
2402async fn run_pool_writer_query<T, F>(
2403 pool: Arc<ConnectionPool>,
2404 operation: &'static str,
2405 query: F,
2406) -> khive_storage::types::StorageResult<T>
2407where
2408 T: Send + 'static,
2409 F: FnOnce(
2410 &crate::read_cancellation::InterruptibleReadScope,
2411 &rusqlite::Connection,
2412 bool,
2413 ) -> khive_storage::types::StorageResult<T>
2414 + Send
2415 + 'static,
2416{
2417 crate::read_cancellation::run_interruptible_read(
2418 StorageCapability::Sql,
2419 operation,
2420 move |scope| {
2421 let guard = pool.try_writer().map_err(|error: SqliteError| {
2422 error.into_storage_error(StorageCapability::Sql, operation)
2423 })?;
2424 scope.with_pooled_writer(&pool, &guard, |conn| {
2425 let interruptible = conn.is_autocommit();
2426 query(scope, conn, interruptible)
2427 })
2428 },
2429 )
2430 .await
2431}
2432
2433struct PoolBackedReader {
2434 pool: Arc<ConnectionPool>,
2435 transaction: Option<SharedReaderTransactionGuard>,
2442}
2443
2444fn finish_pool_backed_reader_step<T>(
2455 transaction: &mut Option<SharedReaderTransactionGuard>,
2456 guard: SharedReaderTransactionGuard,
2457 expect_open_after: bool,
2458 operation: &'static str,
2459 result: khive_storage::types::StorageResult<T>,
2460) -> khive_storage::types::StorageResult<T> {
2461 let still_open = !guard.conn().is_autocommit();
2462 if still_open == expect_open_after {
2463 if still_open {
2464 *transaction = Some(guard);
2465 }
2466 return result;
2469 }
2470 guard.poison();
2471 let message = if expect_open_after {
2472 "a read inside the pool-backed reader's admitted transaction unexpectedly ended it; \
2473 the connection was discarded"
2474 } else {
2475 "transaction-ending control completed but the pool-backed reader's connection \
2476 remained outside autocommit; the connection was discarded"
2477 };
2478 match result {
2479 Err(error) => Err(error),
2480 Ok(_) => Err(StorageError::InvalidInput {
2481 capability: StorageCapability::Sql,
2482 operation: operation.into(),
2483 message: message.into(),
2484 }),
2485 }
2486}
2487
2488async fn open_pool_backed_reader_transaction<T, F>(
2492 transaction: &mut Option<SharedReaderTransactionGuard>,
2493 pool: Arc<ConnectionPool>,
2494 operation: &'static str,
2495 query: F,
2496) -> khive_storage::types::StorageResult<T>
2497where
2498 T: Send + 'static,
2499 F: FnOnce(
2500 &crate::read_cancellation::InterruptibleReadScope,
2501 &rusqlite::Connection,
2502 bool,
2503 bool,
2504 ) -> khive_storage::types::StorageResult<T>
2505 + Send
2506 + 'static,
2507{
2508 let (guard, result) = crate::read_cancellation::run_interruptible_read(
2509 StorageCapability::Sql,
2510 operation,
2511 move |scope| {
2512 let Some(guard) = pool
2513 .checkout_shared_reader_transaction(|| scope.should_stop())
2514 .map_err(|error| StorageError::driver(StorageCapability::Sql, operation, error))?
2515 else {
2516 return Err(StorageError::Timeout {
2517 operation: operation.into(),
2518 });
2519 };
2520 let result = query(scope, guard.conn(), false, true);
2521 if scope.cleanup_failed() {
2522 guard.poison();
2523 }
2524 Ok((guard, result))
2525 },
2526 )
2527 .await?;
2528 finish_pool_backed_reader_step(transaction, guard, true, operation, result)
2529}
2530
2531#[allow(clippy::too_many_lines)]
2537async fn run_pool_backed_reader_query<T, F>(
2538 transaction: &mut Option<SharedReaderTransactionGuard>,
2539 pool: Arc<ConnectionPool>,
2540 operation: &'static str,
2541 transaction_control: Option<CachedReadTransactionControl>,
2542 query: F,
2543) -> khive_storage::types::StorageResult<T>
2544where
2545 T: Send + 'static,
2546 F: FnOnce(
2547 &crate::read_cancellation::InterruptibleReadScope,
2548 &rusqlite::Connection,
2549 bool,
2550 bool,
2551 ) -> khive_storage::types::StorageResult<T>
2552 + Send
2553 + 'static,
2554{
2555 if transaction.is_none() {
2556 return match transaction_control {
2557 None => {
2558 run_pool_reader_query(pool, operation, move |scope, conn| {
2559 query(scope, conn, false, true)
2560 })
2561 .await
2562 }
2563 Some(CachedReadTransactionControl::Finish(keyword))
2564 | Some(CachedReadTransactionControl::Unsupported(keyword)) => {
2565 Err(StorageError::InvalidInput {
2566 capability: StorageCapability::Sql,
2567 operation: operation.into(),
2568 message: format!(
2569 "pool-backed reader has no admitted transaction for transaction \
2570 control ({keyword})"
2571 ),
2572 })
2573 }
2574 Some(CachedReadTransactionControl::BeginDeferred) => {
2575 open_pool_backed_reader_transaction(transaction, pool, operation, query).await
2576 }
2577 };
2578 }
2579
2580 match transaction_control {
2581 Some(CachedReadTransactionControl::BeginDeferred) => {
2582 return Err(StorageError::InvalidInput {
2583 capability: StorageCapability::Sql,
2584 operation: operation.into(),
2585 message: "pool-backed reader already owns an admitted read transaction; \
2586 nested BEGIN is not supported"
2587 .into(),
2588 });
2589 }
2590 Some(CachedReadTransactionControl::Unsupported(keyword)) => {
2591 return Err(StorageError::InvalidInput {
2592 capability: StorageCapability::Sql,
2593 operation: operation.into(),
2594 message: format!(
2595 "pool-backed reader's admitted read transaction does not support nested \
2596 or write-locking transaction control ({keyword})"
2597 ),
2598 });
2599 }
2600 None | Some(CachedReadTransactionControl::Finish(_)) => {}
2601 }
2602
2603 let expect_open_after = transaction_control.is_none();
2604 let guard = transaction.take().expect("checked Some above");
2605 let (guard, result) = crate::read_cancellation::run_interruptible_read(
2606 StorageCapability::Sql,
2607 operation,
2608 move |scope| {
2609 let result = query(scope, guard.conn(), true, true);
2610 if scope.cleanup_failed() {
2611 guard.poison();
2612 }
2613 Ok((guard, result))
2614 },
2615 )
2616 .await?;
2617 finish_pool_backed_reader_step(transaction, guard, expect_open_after, operation, result)
2618}
2619
2620#[async_trait]
2621impl khive_storage::SqlReader for PoolBackedReader {
2622 async fn query_row(
2623 &mut self,
2624 statement: SqlStatement,
2625 ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
2626 let transaction_control = cached_read_transaction_control(&statement.sql);
2627 admit_reader_capability_sql(&statement, transaction_control, "pool_reader.query_row")?;
2628 let pool = Arc::clone(&self.pool);
2629 run_pool_backed_reader_query(
2630 &mut self.transaction,
2631 pool,
2632 "pool_reader.query_row",
2633 transaction_control,
2634 move |scope, conn, rollback, interruptible| {
2635 execute_query_row_interruptibly(
2636 scope,
2637 conn,
2638 &statement,
2639 "pool_reader.query_row",
2640 rollback,
2641 interruptible,
2642 )
2643 },
2644 )
2645 .await
2646 }
2647
2648 async fn query_all(
2649 &mut self,
2650 statement: SqlStatement,
2651 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2652 let transaction_control = cached_read_transaction_control(&statement.sql);
2653 admit_reader_capability_sql(&statement, transaction_control, "pool_reader.query_all")?;
2654 let pool = Arc::clone(&self.pool);
2655 run_pool_backed_reader_query(
2656 &mut self.transaction,
2657 pool,
2658 "pool_reader.query_all",
2659 transaction_control,
2660 move |scope, conn, rollback, interruptible| {
2661 execute_query_interruptibly(
2662 scope,
2663 conn,
2664 &statement,
2665 "pool_reader.query_all",
2666 rollback,
2667 interruptible,
2668 )
2669 },
2670 )
2671 .await
2672 }
2673
2674 async fn query_page(
2675 &mut self,
2676 statement: SqlStatement,
2677 page: PageRequest,
2678 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2679 let transaction_control = cached_read_transaction_control(&statement.sql);
2680 admit_reader_capability_sql(&statement, transaction_control, "pool_reader.query_page")?;
2681 let pool = Arc::clone(&self.pool);
2682 run_pool_backed_reader_query(
2683 &mut self.transaction,
2684 pool,
2685 "pool_reader.query_page",
2686 transaction_control,
2687 move |scope, conn, rollback, interruptible| {
2688 execute_query_page_interruptibly(
2689 scope,
2690 conn,
2691 &statement,
2692 &page,
2693 "pool_reader.query_page",
2694 rollback,
2695 interruptible,
2696 )
2697 },
2698 )
2699 .await
2700 }
2701
2702 async fn query_scalar(
2703 &mut self,
2704 statement: SqlStatement,
2705 ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
2706 let row = self.query_row(statement).await?;
2707 Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
2708 }
2709
2710 async fn explain(
2711 &mut self,
2712 statement: SqlStatement,
2713 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2714 let explain_stmt = SqlStatement {
2715 sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
2716 params: statement.params,
2717 label: statement.label,
2718 };
2719 self.query_all(explain_stmt).await
2720 }
2721}
2722
2723struct PoolBackedWriter {
2724 pool: Arc<ConnectionPool>,
2725}
2726
2727#[async_trait]
2728impl khive_storage::SqlReader for PoolBackedWriter {
2729 async fn query_row(
2730 &mut self,
2731 statement: SqlStatement,
2732 ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
2733 let pool = Arc::clone(&self.pool);
2734 run_pool_writer_query(
2735 pool,
2736 "pool_writer.query_row",
2737 move |scope, conn, interruptible| {
2738 execute_query_row_interruptibly(
2739 scope,
2740 conn,
2741 &statement,
2742 "pool_writer.query_row",
2743 false,
2744 interruptible,
2745 )
2746 },
2747 )
2748 .await
2749 }
2750
2751 async fn query_all(
2752 &mut self,
2753 statement: SqlStatement,
2754 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2755 let pool = Arc::clone(&self.pool);
2756 run_pool_writer_query(
2757 pool,
2758 "pool_writer.query_all",
2759 move |scope, conn, interruptible| {
2760 execute_query_interruptibly(
2761 scope,
2762 conn,
2763 &statement,
2764 "pool_writer.query_all",
2765 false,
2766 interruptible,
2767 )
2768 },
2769 )
2770 .await
2771 }
2772
2773 async fn query_page(
2774 &mut self,
2775 statement: SqlStatement,
2776 page: PageRequest,
2777 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2778 let pool = Arc::clone(&self.pool);
2779 run_pool_writer_query(
2780 pool,
2781 "pool_writer.query_page",
2782 move |scope, conn, interruptible| {
2783 execute_query_page_interruptibly(
2784 scope,
2785 conn,
2786 &statement,
2787 &page,
2788 "pool_writer.query_page",
2789 false,
2790 interruptible,
2791 )
2792 },
2793 )
2794 .await
2795 }
2796
2797 async fn query_scalar(
2798 &mut self,
2799 statement: SqlStatement,
2800 ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
2801 let row = khive_storage::SqlReader::query_row(self, statement).await?;
2802 Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
2803 }
2804
2805 async fn explain(
2806 &mut self,
2807 statement: SqlStatement,
2808 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2809 let explain_stmt = SqlStatement {
2810 sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
2811 params: statement.params,
2812 label: statement.label,
2813 };
2814 khive_storage::SqlReader::query_all(self, explain_stmt).await
2815 }
2816}
2817
2818#[async_trait]
2819impl khive_storage::SqlWriter for PoolBackedWriter {
2820 async fn execute(
2821 &mut self,
2822 statement: SqlStatement,
2823 ) -> khive_storage::types::StorageResult<u64> {
2824 if let Some(keyword) = transaction_control_head(&statement.sql) {
2831 return Err(StorageError::InvalidInput {
2832 capability: StorageCapability::Sql,
2833 operation: "pool_writer.execute".into(),
2834 message: format!(
2835 "statement is transaction control ({keyword}); a pooled writer is \
2836 released after every call, so it cannot hold a transaction across \
2837 calls — use atomic_unit to run statements as one transaction"
2838 ),
2839 });
2840 }
2841 let pool = Arc::clone(&self.pool);
2842 tokio::task::spawn_blocking(move || {
2843 let guard = pool.try_writer().map_err(|e: SqliteError| {
2844 StorageError::driver(StorageCapability::Sql, "pool_writer.execute", e)
2845 })?;
2846 let result = (|| {
2847 let mut stmt = prepare_cached_sql_statement(&guard, &statement.sql)
2848 .map_err(|e| map_rusqlite_err(e, "pool_writer.execute"))?;
2849 bind_params(&mut stmt, &statement.params)
2850 .map_err(|e| map_rusqlite_err(e, "pool_writer.execute"))?;
2851 let rows = stmt
2852 .raw_execute()
2853 .map_err(|e| map_rusqlite_err(e, "pool_writer.execute"))?;
2854 Ok(rows as u64)
2855 })();
2856 settle_pooled_call(&guard, "pool_writer.execute", result)
2857 .inspect_err(|error| pool.record_direct_writer_error(error))
2858 })
2859 .await
2860 .map_err(|e| StorageError::driver(StorageCapability::Sql, "pool_writer.execute", e))?
2861 }
2862
2863 async fn execute_batch(
2864 &mut self,
2865 statements: Vec<SqlStatement>,
2866 ) -> khive_storage::types::StorageResult<u64> {
2867 reject_transaction_control_statements(&statements, "pool_writer.execute_batch")?;
2871 let pool = Arc::clone(&self.pool);
2872 tokio::task::spawn_blocking(move || {
2873 let guard = pool.try_writer().map_err(|e: SqliteError| {
2874 StorageError::driver(StorageCapability::Sql, "pool_writer.execute_batch", e)
2875 })?;
2876 let result = (|| {
2877 let prepared = prepare_batch_statements(&guard, &statements)
2878 .map_err(|e| map_rusqlite_err(e, "pool_writer.execute_batch"))?;
2879 guard
2880 .execute_batch("BEGIN IMMEDIATE")
2881 .map_err(|e| map_rusqlite_err(e, "pool_writer.execute_batch"))?;
2882 let _tx_handle = khive_storage::tx_registry::register_scoped(
2883 Some("pool_writer.execute_batch".to_string()),
2884 pool.origin(),
2885 );
2886 let result = execute_prepared_batch(&guard, prepared, &statements, None)
2887 .map_err(|e| map_rusqlite_err(e, "pool_writer.execute_batch"));
2888 match result {
2889 Ok(total) => {
2890 if let Err(e) = guard.execute_batch("COMMIT") {
2891 let _ = guard.execute_batch("ROLLBACK");
2892 Err(map_rusqlite_err(e, "pool_writer.execute_batch"))
2893 } else {
2894 Ok(total)
2895 }
2896 }
2897 Err(e) => {
2898 let _ = guard.execute_batch("ROLLBACK");
2899 Err(e)
2900 }
2901 }
2902 })();
2903 result.inspect_err(|error| pool.record_direct_writer_error(error))
2904 })
2905 .await
2906 .map_err(|e| StorageError::driver(StorageCapability::Sql, "pool_writer.execute_batch", e))?
2907 }
2908
2909 async fn execute_script(&mut self, script: String) -> khive_storage::types::StorageResult<()> {
2910 let pool = Arc::clone(&self.pool);
2915 tokio::task::spawn_blocking(move || {
2916 let guard = pool.try_writer().map_err(|e: SqliteError| {
2917 StorageError::driver(StorageCapability::Sql, "pool_writer.execute_script", e)
2918 })?;
2919 let result = guard
2920 .execute_batch(&script)
2921 .map_err(|e| map_rusqlite_err(e, "pool_writer.execute_script"));
2922 settle_pooled_call(&guard, "pool_writer.execute_script", result)
2923 .inspect_err(|error| pool.record_direct_writer_error(error))
2924 })
2925 .await
2926 .map_err(|e| {
2927 StorageError::driver(StorageCapability::Sql, "pool_writer.execute_script", e)
2928 })?
2929 }
2930}
2931
2932struct InlineWriter {
2958 event_rows: Option<Arc<AtomicEventRows>>,
2959 conn: *const rusqlite::Connection,
2960}
2961
2962unsafe impl Send for InlineWriter {}
2969
2970impl InlineWriter {
2971 fn conn(&self) -> &rusqlite::Connection {
2975 unsafe { &*self.conn }
2976 }
2977}
2978
2979#[async_trait]
2980impl khive_storage::SqlReader for InlineWriter {
2981 async fn query_row(
2982 &mut self,
2983 statement: SqlStatement,
2984 ) -> khive_storage::types::StorageResult<Option<SqlRow>> {
2985 execute_query_row(self.conn(), &statement)
2986 .map_err(|e| map_rusqlite_err(e, "inline.query_row"))
2987 }
2988
2989 async fn query_all(
2990 &mut self,
2991 statement: SqlStatement,
2992 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
2993 execute_query(self.conn(), &statement).map_err(|e| map_rusqlite_err(e, "inline.query_all"))
2994 }
2995
2996 async fn query_page(
2997 &mut self,
2998 statement: SqlStatement,
2999 page: PageRequest,
3000 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
3001 execute_query_page(self.conn(), &statement, &page)
3002 .map_err(|e| map_rusqlite_err(e, "inline.query_page"))
3003 }
3004
3005 async fn query_scalar(
3006 &mut self,
3007 statement: SqlStatement,
3008 ) -> khive_storage::types::StorageResult<Option<SqlValue>> {
3009 let row = khive_storage::SqlReader::query_row(self, statement).await?;
3010 Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
3011 }
3012
3013 async fn explain(
3014 &mut self,
3015 statement: SqlStatement,
3016 ) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
3017 let explain_stmt = SqlStatement {
3018 sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
3019 params: statement.params,
3020 label: statement.label,
3021 };
3022 khive_storage::SqlReader::query_all(self, explain_stmt).await
3023 }
3024}
3025
3026#[async_trait]
3027impl khive_storage::SqlWriter for InlineWriter {
3028 async fn execute(
3029 &mut self,
3030 statement: SqlStatement,
3031 ) -> khive_storage::types::StorageResult<u64> {
3032 let mut stmt = prepare_cached_sql_statement(self.conn(), &statement.sql)
3035 .map_err(|e| map_rusqlite_err(e, "inline.execute"))?;
3036 bind_params(&mut stmt, &statement.params)
3037 .map_err(|e| map_rusqlite_err(e, "inline.execute"))?;
3038 let affected = stmt
3039 .raw_execute()
3040 .map_err(|e| map_rusqlite_err(e, "inline.execute"))?;
3041 if let Some(event_rows) = self.event_rows.as_deref() {
3042 event_rows.observe(&statement, affected as u64);
3043 }
3044 Ok(affected as u64)
3045 }
3046
3047 async fn execute_batch(
3048 &mut self,
3049 statements: Vec<SqlStatement>,
3050 ) -> khive_storage::types::StorageResult<u64> {
3051 reject_transaction_control_statements(&statements, "inline.execute_batch")?;
3056 let prepared = prepare_batch_statements(self.conn(), &statements)
3057 .map_err(|e| map_rusqlite_err(e, "inline.execute_batch"))?;
3058 execute_prepared_batch(
3059 self.conn(),
3060 prepared,
3061 &statements,
3062 self.event_rows.as_deref(),
3063 )
3064 .map_err(|e| map_rusqlite_err(e, "inline.execute_batch"))
3065 }
3066
3067 async fn execute_script(&mut self, script: String) -> khive_storage::types::StorageResult<()> {
3068 self.conn()
3071 .execute_batch(&script)
3072 .map_err(|e| map_rusqlite_err(e, "inline.execute_script"))
3073 }
3074}
3075
3076fn block_on_sync<F: std::future::Future>(fut: F) -> Result<F::Output, StorageError> {
3099 use std::task::{Context, Poll, RawWaker, RawWakerVTable, Waker};
3100
3101 fn no_op(_: *const ()) {}
3102 fn clone_waker(_: *const ()) -> RawWaker {
3103 RawWaker::new(std::ptr::null(), &VTABLE)
3104 }
3105 static VTABLE: RawWakerVTable = RawWakerVTable::new(clone_waker, no_op, no_op, no_op);
3106
3107 let raw_waker = RawWaker::new(std::ptr::null(), &VTABLE);
3110 let waker = unsafe { Waker::from_raw(raw_waker) };
3111 let mut cx = Context::from_waker(&waker);
3112
3113 let mut fut = std::pin::pin!(fut);
3114 match fut.as_mut().poll(&mut cx) {
3115 Poll::Ready(v) => Ok(v),
3116 Poll::Pending => {
3117 tracing::error!(
3118 "block_on_sync: atomic_unit future suspended on its first poll — \
3119 the closure passed to SqlAccess::atomic_unit must be non-blocking \
3120 (synchronous InlineWriter calls only, no real .await point)"
3121 );
3122 Err(StorageError::Internal(
3123 "atomic_unit future suspended — closure must be non-blocking".to_string(),
3124 ))
3125 }
3126 }
3127}
3128
3129pub struct SqlBridge {
3142 pool: Arc<ConnectionPool>,
3143 is_file_backed: bool,
3144}
3145
3146impl SqlBridge {
3147 pub fn new(pool: Arc<ConnectionPool>, _is_file_backed: bool) -> Self {
3149 let is_file_backed = pool.canonical_path().is_some();
3153 Self {
3154 pool,
3155 is_file_backed,
3156 }
3157 }
3158}
3159
3160#[async_trait]
3161impl khive_storage::SqlAccess for SqlBridge {
3162 fn database_path(&self) -> Option<std::path::PathBuf> {
3163 self.pool.canonical_path().map(std::path::Path::to_path_buf)
3164 }
3165
3166 async fn reader(
3167 &self,
3168 ) -> khive_storage::types::StorageResult<Box<dyn khive_storage::SqlReader>> {
3169 if self.is_file_backed {
3170 Ok(Box::new(SqliteReader {
3171 handle: None,
3172 pool: Arc::clone(&self.pool),
3173 poisoned: false,
3174 }))
3175 } else {
3176 Ok(Box::new(PoolBackedReader {
3177 pool: Arc::clone(&self.pool),
3178 transaction: None,
3179 }))
3180 }
3181 }
3182
3183 async fn writer(
3184 &self,
3185 ) -> khive_storage::types::StorageResult<Box<dyn khive_storage::SqlWriter>> {
3186 if self.is_file_backed {
3187 if self.pool.config().read_only {
3188 return Err(StorageError::Pool {
3189 operation: "writer".into(),
3190 message: "backend is read-only".into(),
3191 });
3192 }
3193 let db = crate::timeout_sink::db_label(&self.pool);
3194 let writer_task = match self.pool.writer_task_handle() {
3200 Ok(handle) => handle,
3201 Err(e) => {
3202 if self.pool.config().write_routing_strict {
3203 return Err(e);
3204 }
3205 tracing::warn!(
3206 error = %e,
3207 "KHIVE_WRITE_ROUTING is not strict; writer() degrades to the \
3208 standalone-connection path"
3209 );
3210 None
3211 }
3212 };
3213 if writer_task.is_none() && self.pool.config().write_routing_strict {
3214 return Err(StorageError::Pool {
3215 operation: "writer".into(),
3216 message: "KHIVE_WRITE_ROUTING=strict but no writer-task handle is \
3217 available; refusing to fall back to a direct connection"
3218 .into(),
3219 });
3220 }
3221 if writer_task.is_none() && self.pool.write_queue_active() {
3222 crate::timeout_sink::emit_direct_route_violation(
3229 &db,
3230 crate::timeout_sink::Site::DirectRouteSqlBridgeWriter,
3231 );
3232 }
3233 let handle = if writer_task.is_none() {
3247 let handle_slot = acquire_handle_slot(
3248 self.pool.sql_bridge_writer_slots(),
3249 self.pool.config().checkout_timeout,
3250 "sql_bridge.writer_handle",
3251 SlotTimeoutClass::Admission,
3252 )
3253 .await?;
3254 let (conn, handle_slot) =
3255 open_standalone_writer_on_blocking(Arc::clone(&self.pool), handle_slot).await?;
3256 Some(StandaloneHandle {
3257 conn,
3258 _retained_slot: Some(handle_slot),
3259 read_transaction_slot: None,
3260 })
3261 } else {
3262 None
3263 };
3264 Ok(Box::new(SqliteWriter {
3265 observe_direct_errors: true,
3266 event_rows: None,
3267 handle,
3268 writer_task,
3269 origin: self.pool.origin(),
3270 db,
3271 pool: Arc::clone(&self.pool),
3272 held_lease: None,
3273 }))
3274 } else {
3275 Ok(Box::new(PoolBackedWriter {
3276 pool: Arc::clone(&self.pool),
3277 }))
3278 }
3279 }
3280
3281 async fn atomic_unit(
3292 &self,
3293 op: AtomicUnitOp,
3294 ) -> khive_storage::types::StorageResult<Box<dyn Any + Send>> {
3295 let event_rows = Arc::new(AtomicEventRows::default());
3296 let result = async {
3297 if self.is_file_backed {
3298 if self.pool.config().read_only {
3299 return Err(StorageError::Pool {
3300 operation: "atomic_unit".into(),
3301 message: "backend is read-only".into(),
3302 });
3303 }
3304 let handle = self.pool.writer_task_handle()?;
3310 if handle.is_none() && self.pool.config().write_routing_strict {
3311 return Err(StorageError::Pool {
3312 operation: "atomic_unit".into(),
3313 message: "KHIVE_WRITE_ROUTING=strict but no writer-task handle is \
3314 available; refusing to fall back to a direct connection"
3315 .into(),
3316 });
3317 }
3318 if handle.is_none() && self.pool.write_queue_active() {
3319 crate::timeout_sink::emit_direct_route_violation(
3320 &crate::timeout_sink::db_label(&self.pool),
3321 crate::timeout_sink::Site::DirectRouteAtomicUnit,
3322 );
3323 }
3324 if let Some(writer_task) = handle {
3325 let pending_event_rows = Arc::clone(&event_rows);
3331 return writer_task
3332 .send_bounded(move |conn| {
3333 let mut inline = InlineWriter {
3334 event_rows: Some(Arc::clone(&pending_event_rows)),
3335 conn: conn as *const rusqlite::Connection,
3336 };
3337 match block_on_sync(op(&mut inline)) {
3347 Ok(inner) => inner,
3348 Err(e) => Err(e),
3349 }
3350 })
3351 .await;
3352 }
3353 let handle_slot = acquire_handle_slot(
3367 self.pool.sql_bridge_writer_slots(),
3368 self.pool.config().checkout_timeout,
3369 "sql_bridge.atomic_unit_handle",
3370 SlotTimeoutClass::Admission,
3371 )
3372 .await?;
3373 let unit_lease = acquire_unit_lease(Arc::clone(&self.pool)).await?;
3376 let (conn, handle_slot) =
3377 open_standalone_writer_on_blocking(Arc::clone(&self.pool), handle_slot).await?;
3378 let mut writer = SqliteWriter {
3379 observe_direct_errors: false,
3380 event_rows: Some(Arc::clone(&event_rows)),
3381 handle: Some(StandaloneHandle {
3382 conn,
3383 _retained_slot: Some(handle_slot),
3384 read_transaction_slot: None,
3385 }),
3386 writer_task: None,
3387 origin: self.pool.origin(),
3388 db: crate::timeout_sink::db_label(&self.pool),
3389 pool: Arc::clone(&self.pool),
3390 held_lease: unit_lease,
3391 };
3392 run_manual_atomic_unit(&mut writer, op, self.pool.origin())
3393 .await
3394 .inspect_err(|error| self.pool.record_direct_writer_error(error))
3395 } else {
3396 let pool = Arc::clone(&self.pool);
3399 let pending_event_rows = Arc::clone(&event_rows);
3400 tokio::task::spawn_blocking(move || {
3401 let guard = pool.try_writer().map_err(|error: SqliteError| {
3402 StorageError::driver(StorageCapability::Sql, "atomic_unit", error)
3403 })?;
3404 let conn = guard.conn();
3405 if !conn.is_autocommit() {
3406 pool.retire_pooled_writer(conn);
3407 return Err(StorageError::WriterTaskTerminated {
3408 request_state:
3409 khive_storage::WriterTaskRequestState::SideEffectsUnknown,
3410 });
3411 }
3412 if let Err(error) = conn.execute_batch("BEGIN IMMEDIATE") {
3413 if !conn.is_autocommit() {
3414 pool.retire_pooled_writer(conn);
3415 return Err(StorageError::WriterTaskTerminated {
3416 request_state:
3417 khive_storage::WriterTaskRequestState::SideEffectsUnknown,
3418 });
3419 }
3420 return Err(map_rusqlite_err(error, "atomic_unit.begin"))
3421 .inspect_err(|error| pool.record_direct_writer_error(error));
3422 }
3423 let _tx_handle = khive_storage::tx_registry::register_scoped(
3424 Some("atomic_unit".to_string()),
3425 pool.origin(),
3426 );
3427 let (result, terminal_state) = crate::writer_task::execute_wrapped_transaction(
3428 conn,
3429 "atomic_unit.commit",
3430 |conn| {
3431 let mut inline = InlineWriter {
3432 event_rows: Some(Arc::clone(&pending_event_rows)),
3433 conn: conn as *const rusqlite::Connection,
3434 };
3435 block_on_sync(op(&mut inline)).and_then(|result| result)
3436 },
3437 );
3438 if terminal_state.is_some() {
3439 pool.retire_pooled_writer(conn);
3440 }
3441 result.inspect_err(|error| pool.record_direct_writer_error(error))
3442 })
3443 .await
3444 .map_err(|error| {
3445 StorageError::driver(StorageCapability::Sql, "atomic_unit", error)
3446 })?
3447 }
3448 }
3449 .await;
3450 khive_storage::usage::account_event_write(
3451 result.as_ref().map(|_| event_rows.committed_rows()),
3452 );
3453 result
3454 }
3455}
3456
3457#[cfg(test)]
3458mod tests {
3459 use super::*;
3460 use crate::pool::PoolConfig;
3461 use khive_storage::types::{SqlStatement, SqlValue};
3462 use khive_storage::{SqlAccess as _, SqlReader as _};
3463
3464 #[tokio::test]
3465 async fn top_level_wal_checkpoint_ends_the_active_pin_run() {
3466 let dir = tempfile::tempdir().unwrap();
3467 let pool = Arc::new(
3468 ConnectionPool::new(PoolConfig {
3469 path: Some(dir.path().join("top_level_checkpoint.db")),
3470 write_queue_enabled: Some(false),
3471 ..PoolConfig::for_test()
3472 })
3473 .unwrap(),
3474 );
3475 {
3476 let writer = pool.writer().unwrap();
3477 writer
3478 .conn()
3479 .execute_batch("CREATE TABLE t (value INTEGER); INSERT INTO t VALUES (1)")
3480 .unwrap();
3481 }
3482
3483 let _task_guard = crate::checkpoint::CheckpointRunTaskGuard::start(
3484 &pool,
3485 std::time::Duration::from_secs(60),
3486 );
3487 crate::checkpoint::record_checkpoint_run_result(&pool, Some((0, 20, 10)));
3488 let bridge = SqlBridge::new(Arc::clone(&pool), true);
3489 bridge
3490 .writer()
3491 .await
3492 .unwrap()
3493 .execute_script_top_level(TopLevelMaintenance::WalCheckpointTruncate)
3494 .await
3495 .unwrap();
3496
3497 assert_eq!(
3498 crate::checkpoint::checkpoint_run_status(&pool),
3499 crate::checkpoint::CheckpointRunStatus::NoObservation
3500 );
3501 }
3502
3503 include!("sql_bridge_atomic_serialization_tests.rs");
3504 include!("sql_bridge_write_admission_tests.rs");
3505
3506 fn database_tx_view(pool: &ConnectionPool) -> khive_storage::tx_registry::TxOriginFilter {
3507 match pool.origin() {
3508 khive_storage::tx_registry::TxOrigin::Database(identity) => {
3509 khive_storage::tx_registry::TxOriginFilter::Secondary(identity)
3510 }
3511 other => panic!("expected a file-backed database origin, got {other:?}"),
3512 }
3513 }
3514
3515 struct NotifyOnDrop(Arc<tokio::sync::Notify>);
3516
3517 impl Drop for NotifyOnDrop {
3518 fn drop(&mut self) {
3519 self.0.notify_one();
3520 }
3521 }
3522
3523 fn blocking_non_interrupting_progress_gate(
3524 conn: &rusqlite::Connection,
3525 ) -> (
3526 Arc<tokio::sync::Notify>,
3527 Arc<std::sync::Barrier>,
3528 Arc<tokio::sync::Notify>,
3529 ) {
3530 let entered = Arc::new(tokio::sync::Notify::new());
3531 let callback_entered = Arc::clone(&entered);
3532 let release = Arc::new(std::sync::Barrier::new(2));
3533 let callback_release = Arc::clone(&release);
3534 let completed = Arc::new(tokio::sync::Notify::new());
3535 let notify_on_drop = NotifyOnDrop(Arc::clone(&completed));
3536 let blocked_once = Arc::new(std::sync::atomic::AtomicBool::new(false));
3537 let callback_blocked_once = Arc::clone(&blocked_once);
3538 conn.progress_handler(
3539 1_000,
3540 Some(move || {
3541 let _keep_until_connection_drop = ¬ify_on_drop;
3542 if !callback_blocked_once.swap(true, std::sync::atomic::Ordering::SeqCst) {
3543 callback_entered.notify_one();
3544 callback_release.wait();
3545 return false;
3549 }
3550 false
3551 }),
3552 )
3553 .unwrap();
3554 (entered, release, completed)
3555 }
3556
3557 fn progress_gate_statement() -> SqlStatement {
3558 SqlStatement {
3559 sql: "WITH RECURSIVE rows(value) AS (\
3560 SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 999\
3561 ) SELECT SUM(value) FROM rows"
3562 .into(),
3563 params: vec![],
3564 label: None,
3565 }
3566 }
3567
3568 fn slow_insert_statement() -> SqlStatement {
3569 SqlStatement {
3570 sql: "INSERT INTO cancellation_write_probe(value) \
3571 WITH RECURSIVE rows(value) AS (\
3572 SELECT 1 UNION ALL SELECT value + 1 FROM rows WHERE value < 10000\
3573 ) SELECT value FROM rows"
3574 .into(),
3575 params: vec![],
3576 label: Some("non-interruptible-write-probe".into()),
3577 }
3578 }
3579
3580 fn passive_checkpoint(conn: &rusqlite::Connection) -> (i64, i64, i64) {
3581 conn.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
3582 Ok((row.get(0)?, row.get(1)?, row.get(2)?))
3583 })
3584 .unwrap()
3585 }
3586
3587 fn deliberately_slow_read_statement() -> SqlStatement {
3588 SqlStatement {
3589 sql: "WITH RECURSIVE numbers(value) AS (\
3590 SELECT 1 UNION ALL SELECT value + 1 FROM numbers WHERE value < 1000\
3591 ) SELECT SUM(a.value * b.value * c.value) \
3592 FROM numbers AS a CROSS JOIN numbers AS b CROSS JOIN numbers AS c"
3593 .into(),
3594 params: vec![],
3595 label: Some("read-cancellation-progress-probe".into()),
3596 }
3597 }
3598
3599 async fn wait_for_progress(probe: &std::sync::atomic::AtomicUsize) {
3600 tokio::time::timeout(std::time::Duration::from_secs(1), async {
3601 while probe.load(std::sync::atomic::Ordering::SeqCst) == 0 {
3602 tokio::task::yield_now().await;
3603 }
3604 })
3605 .await
3606 .expect("slow SQLite statement never reached its progress callback");
3607 }
3608
3609 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3610 async fn pooled_stats_count_cancellation_stops_scan_and_releases_reader() {
3611 let dir = tempfile::tempdir().unwrap();
3612 let pool = Arc::new(
3613 ConnectionPool::new(PoolConfig {
3614 path: Some(dir.path().join("stats-count-cancel.db")),
3615 max_readers: 1,
3616 ..PoolConfig::default()
3617 })
3618 .unwrap(),
3619 );
3620 pool.writer()
3623 .unwrap()
3624 .conn()
3625 .execute_batch(
3626 "CREATE TABLE count_fixture(n INTEGER PRIMARY KEY); \
3627 WITH RECURSIVE n(x) AS (SELECT 1 UNION ALL SELECT x+1 FROM n WHERE x<1000) \
3628 INSERT INTO count_fixture SELECT x FROM n; \
3629 CREATE VIEW events AS SELECT 'local' AS namespace, 'knowledge.learn' AS verb \
3630 FROM count_fixture a CROSS JOIN count_fixture b CROSS JOIN count_fixture c;",
3631 )
3632 .unwrap();
3633 let bridge = SqlBridge::new(Arc::clone(&pool), true);
3634 let mut reader = bridge.reader().await.unwrap();
3635 let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
3636 let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
3637 let query = tokio::spawn(crate::scope_test_read_progress(
3638 Arc::clone(&progress),
3639 crate::scope_request_read_cancellation(cancel_rx, async move {
3640 let result = reader.query_scalar(SqlStatement {
3641 sql: "SELECT COUNT(*) FROM events WHERE namespace = ?1 AND verb LIKE 'knowledge.%'".into(),
3642 params: vec![SqlValue::Text("local".into())],
3643 label: Some("knowledge.stats.event_count".into()),
3644 }).await;
3645 (reader, result)
3646 }),
3647 ));
3648 wait_for_progress(progress.as_ref()).await;
3649 assert!(
3650 !query.is_finished(),
3651 "COUNT must still be scanning before cancellation"
3652 );
3653 let started = std::time::Instant::now();
3654 let grace = crate::read_cancellation::sqlite_interrupt_grace_from_env();
3655 cancel_tx.send(true).unwrap();
3656 let (mut reader, result) = tokio::time::timeout(grace, query)
3657 .await
3658 .expect("COUNT did not settle within the interrupt grace")
3659 .unwrap();
3660 let elapsed = started.elapsed();
3661 assert!(
3662 matches!(result, Err(StorageError::Timeout { .. })),
3663 "{result:?}"
3664 );
3665 assert_eq!(
3666 pool.available_readers(),
3667 1,
3668 "COUNT retained the sole pooled reader"
3669 );
3670 let stopped = progress.load(std::sync::atomic::Ordering::SeqCst);
3671 let next = reader
3672 .query_scalar(SqlStatement {
3673 sql: "SELECT COUNT(*) FROM count_fixture".into(),
3674 params: vec![],
3675 label: None,
3676 })
3677 .await
3678 .unwrap();
3679 assert!(matches!(next, Some(SqlValue::Integer(1000))));
3680 assert_eq!(
3681 progress.load(std::sync::atomic::Ordering::SeqCst),
3682 stopped,
3683 "cancelled callback leaked into the next borrower"
3684 );
3685 eprintln!(
3686 "stats_count_cancel_ms={} grace_ms={}",
3687 elapsed.as_secs_f64() * 1000.0,
3688 grace.as_millis()
3689 );
3690 }
3691
3692 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3693 async fn cancellation_before_reader_checkout_is_prompt_and_executes_no_statement() {
3694 let dir = tempfile::tempdir().unwrap();
3695 let config = PoolConfig {
3696 path: Some(dir.path().join("sql_bridge_cancel_before_checkout.db")),
3697 max_readers: 1,
3698 checkout_timeout: std::time::Duration::from_secs(5),
3699 ..PoolConfig::default()
3700 };
3701 let pool = Arc::new(ConnectionPool::new(config).unwrap());
3702 pool.writer()
3703 .unwrap()
3704 .conn()
3705 .execute_batch(
3706 "CREATE TABLE checkout_cancel_probe(value INTEGER NOT NULL); \
3707 INSERT INTO checkout_cancel_probe VALUES (0);",
3708 )
3709 .unwrap();
3710 let held_reader = pool.reader().expect("hold the sole pooled reader");
3711 let bridge = SqlBridge::new(Arc::clone(&pool), true);
3719 let mut waiting_reader = bridge.reader().await.unwrap();
3720 let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
3721 let waiting = tokio::spawn(crate::scope_request_read_cancellation(
3722 cancel_rx,
3723 async move {
3724 waiting_reader
3725 .query_row(SqlStatement {
3726 sql: "SELECT value FROM checkout_cancel_probe".into(),
3727 params: vec![],
3728 label: Some("must-not-run-after-cancelled-checkout".into()),
3729 })
3730 .await
3731 },
3732 ));
3733
3734 tokio::task::yield_now().await;
3735 cancel_tx.send(true).unwrap();
3736 let result = tokio::time::timeout(std::time::Duration::from_millis(100), waiting)
3737 .await
3738 .expect("cancelled reader checkout waited for the five-second pool timeout")
3739 .expect("checkout task panicked");
3740 assert!(matches!(result, Err(StorageError::Timeout { .. })));
3744
3745 drop(held_reader);
3746 tokio::time::sleep(std::time::Duration::from_millis(25)).await;
3747 let value: i64 = pool
3748 .reader()
3749 .unwrap()
3750 .conn()
3751 .query_row("SELECT value FROM checkout_cancel_probe", [], |row| {
3752 row.get(0)
3753 })
3754 .unwrap();
3755 assert_eq!(
3756 value, 0,
3757 "the probe row must be untouched: nothing else in this test writes to it"
3758 );
3759 assert_eq!(
3760 pool.available_readers(),
3761 1,
3762 "reader checkout leaked a permit"
3763 );
3764 }
3765
3766 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3775 async fn probe_reader_capability_admission_of_state_mutating_statements() {
3776 let dir = tempfile::tempdir().unwrap();
3777 let config = PoolConfig {
3778 path: Some(dir.path().join("sql_bridge_reader_admission_probe.db")),
3779 max_readers: 2,
3780 ..PoolConfig::default()
3781 };
3782 let pool = Arc::new(ConnectionPool::new(config).unwrap());
3783 pool.writer()
3784 .unwrap()
3785 .conn()
3786 .execute_batch("CREATE TABLE reader_admission_probe(value INTEGER NOT NULL);")
3787 .unwrap();
3788 let bridge = SqlBridge::new(Arc::clone(&pool), true);
3789
3790 let probes: [&str; 5] = [
3791 "ATTACH DATABASE ':memory:' AS x",
3792 "PRAGMA writable_schema=ON",
3793 "PRAGMA busy_timeout=1",
3794 "PRAGMA cache_size(-64)",
3795 "CREATE TEMP TABLE t(x)",
3796 ];
3797 let mut admitted = Vec::new();
3798 for probe in probes {
3799 let mut reader = bridge.reader().await.unwrap();
3800 let result = reader
3801 .query_all(SqlStatement {
3802 sql: probe.into(),
3803 params: vec![],
3804 label: Some("reader-admission-probe".into()),
3805 })
3806 .await;
3807 admitted.push((probe, result.is_ok()));
3808 }
3809 eprintln!("reader capability admission per probe: {admitted:#?}");
3810 for (probe, was_admitted) in &admitted {
3811 assert!(
3812 !was_admitted,
3813 "reader capability must refuse {probe:?}; the pre-fix bridge wrongly admitted it"
3814 );
3815 }
3816
3817 let controls: [&str; 2] = [
3820 "SELECT value FROM reader_admission_probe",
3821 "PRAGMA table_info(reader_admission_probe)",
3822 ];
3823 for control in controls {
3824 let mut reader = bridge.reader().await.unwrap();
3825 let result = reader
3826 .query_all(SqlStatement {
3827 sql: control.into(),
3828 params: vec![],
3829 label: Some("reader-admission-control".into()),
3830 })
3831 .await;
3832 assert!(
3833 result.is_ok(),
3834 "reader capability must still admit {control:?}: {result:?}"
3835 );
3836 }
3837 }
3838
3839 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3845 async fn probe_reader_capability_admission_of_with_cte_dml_statements() {
3846 let dir = tempfile::tempdir().unwrap();
3847 let config = PoolConfig {
3848 path: Some(dir.path().join("sql_bridge_cte_dml_admission_probe.db")),
3849 max_readers: 2,
3850 ..PoolConfig::default()
3851 };
3852 let pool = Arc::new(ConnectionPool::new(config).unwrap());
3853 pool.writer()
3854 .unwrap()
3855 .conn()
3856 .execute_batch(
3857 "CREATE TABLE cte_dml_admission_probe(id INTEGER PRIMARY KEY, value INTEGER NOT NULL);",
3858 )
3859 .unwrap();
3860 let bridge = SqlBridge::new(Arc::clone(&pool), true);
3861
3862 let probes: [&str; 4] = [
3863 "WITH x(v) AS (SELECT 1) \
3864 INSERT INTO cte_dml_admission_probe(id, value) SELECT 1, v FROM x",
3865 "WITH x(v) AS (SELECT 1) \
3866 INSERT INTO cte_dml_admission_probe(id, value) SELECT 1, v FROM x RETURNING id",
3867 "WITH x(v) AS (SELECT 0) \
3868 UPDATE cte_dml_admission_probe SET value = value + (SELECT v FROM x) WHERE id = 1",
3869 "WITH x(v) AS (SELECT 1) \
3870 DELETE FROM cte_dml_admission_probe WHERE id = (SELECT v FROM x)",
3871 ];
3872 let mut admitted = Vec::new();
3873 for probe in probes {
3874 let mut reader = bridge.reader().await.unwrap();
3875 let result = reader
3876 .query_all(SqlStatement {
3877 sql: probe.into(),
3878 params: vec![],
3879 label: Some("cte-dml-admission-probe".into()),
3880 })
3881 .await;
3882 eprintln!("{probe:?} -> {result:?}");
3883 admitted.push((probe, result.is_ok()));
3884 }
3885 eprintln!("WITH-DML reader capability admission per probe: {admitted:#?}");
3886 for (probe, was_admitted) in &admitted {
3887 assert!(
3888 !was_admitted,
3889 "reader capability must refuse {probe:?}; it is a WITH-prefixed write, not a read"
3890 );
3891 }
3892
3893 let controls: [&str; 3] = [
3898 "WITH x(v) AS (SELECT 1) SELECT v FROM x",
3899 "WITH RECURSIVE n(v) AS (VALUES(0) UNION ALL SELECT v + 1 FROM n WHERE v < 3) \
3900 SELECT v FROM n",
3901 "WITH a(v) AS (SELECT 1), b(v) AS (SELECT v FROM a WHERE ',' || 'x' NOT LIKE '%,%') \
3902 SELECT v FROM b",
3903 ];
3904 for control in controls {
3905 let mut reader = bridge.reader().await.unwrap();
3906 let result = reader
3907 .query_all(SqlStatement {
3908 sql: control.into(),
3909 params: vec![],
3910 label: Some("cte-dml-admission-control".into()),
3911 })
3912 .await;
3913 assert!(
3914 result.is_ok(),
3915 "reader capability must still admit {control:?}: {result:?}"
3916 );
3917 }
3918 }
3919
3920 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3927 async fn probe_reader_capability_admission_of_quoted_cte_names() {
3928 let dir = tempfile::tempdir().unwrap();
3929 let config = PoolConfig {
3930 path: Some(dir.path().join("sql_bridge_quoted_cte_admission_probe.db")),
3931 max_readers: 2,
3932 ..PoolConfig::default()
3933 };
3934 let pool = Arc::new(ConnectionPool::new(config).unwrap());
3935 pool.writer()
3936 .unwrap()
3937 .conn()
3938 .execute_batch(
3939 "CREATE TABLE quoted_cte_admission_probe(id INTEGER PRIMARY KEY, value INTEGER NOT NULL);",
3940 )
3941 .unwrap();
3942 let bridge = SqlBridge::new(Arc::clone(&pool), true);
3943
3944 let controls: [&str; 3] = [
3945 "WITH \"my (cte)\"(v) AS (SELECT 1) SELECT v FROM \"my (cte)\"",
3946 "WITH `my (cte)`(v) AS (SELECT 1) SELECT v FROM `my (cte)`",
3947 "WITH [my (cte)](v) AS (SELECT 1) SELECT v FROM [my (cte)]",
3948 ];
3949 for control in controls {
3950 let mut reader = bridge.reader().await.unwrap();
3951 let result = reader
3952 .query_all(SqlStatement {
3953 sql: control.into(),
3954 params: vec![],
3955 label: Some("quoted-cte-admission-control".into()),
3956 })
3957 .await;
3958 assert!(
3959 result.is_ok(),
3960 "reader capability must admit a quoted CTE name in {control:?}: {result:?}"
3961 );
3962 }
3963
3964 let write_probe = "WITH \"my (cte)\"(v) AS (SELECT 1) \
3965 INSERT INTO quoted_cte_admission_probe(id, value) SELECT 1, v FROM \"my (cte)\"";
3966 let mut reader = bridge.reader().await.unwrap();
3967 let result = reader
3968 .query_all(SqlStatement {
3969 sql: write_probe.into(),
3970 params: vec![],
3971 label: Some("quoted-cte-admission-probe".into()),
3972 })
3973 .await;
3974 assert!(
3975 result.is_err(),
3976 "reader capability must refuse a write statement under a quoted CTE name: {result:?}"
3977 );
3978 }
3979
3980 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3992 async fn pool_backed_reader_admits_deferred_read_transaction_control() {
3993 let config = PoolConfig {
3994 path: None,
3995 ..PoolConfig::default()
3996 };
3997 let pool = Arc::new(ConnectionPool::new(config).unwrap());
3998 pool.writer()
3999 .unwrap()
4000 .conn()
4001 .execute_batch(
4002 "CREATE TABLE pool_backed_reader_probe(value INTEGER NOT NULL); \
4003 INSERT INTO pool_backed_reader_probe VALUES (1);",
4004 )
4005 .unwrap();
4006 let bridge = SqlBridge::new(Arc::clone(&pool), false);
4007 let mut reader = bridge.reader().await.unwrap();
4008
4009 reader
4010 .query_all(SqlStatement {
4011 sql: "BEGIN DEFERRED".into(),
4012 params: vec![],
4013 label: Some("pool-backed-reader-snapshot-begin".into()),
4014 })
4015 .await
4016 .expect("BEGIN DEFERRED must be admitted through the pool-backed reader");
4017 let rows = reader
4018 .query_all(SqlStatement {
4019 sql: "SELECT value FROM pool_backed_reader_probe".into(),
4020 params: vec![],
4021 label: Some("pool-backed-reader-snapshot-read".into()),
4022 })
4023 .await
4024 .expect("a read inside the admitted snapshot must succeed");
4025 assert_eq!(rows.len(), 1);
4026 reader
4027 .query_all(SqlStatement {
4028 sql: "COMMIT".into(),
4029 params: vec![],
4030 label: Some("pool-backed-reader-snapshot-commit".into()),
4031 })
4032 .await
4033 .expect("COMMIT must be admitted through the pool-backed reader");
4034
4035 let integrity = reader
4036 .query_scalar(SqlStatement {
4037 sql: "PRAGMA integrity_check".into(),
4038 params: vec![],
4039 label: Some("pool-backed-reader-integrity-check".into()),
4040 })
4041 .await
4042 .expect("PRAGMA integrity_check must be admitted through the pool-backed reader");
4043 assert!(matches!(integrity, Some(SqlValue::Text(ref s)) if s.eq_ignore_ascii_case("ok")));
4044
4045 for probe in [
4046 "ATTACH DATABASE ':memory:' AS x",
4047 "PRAGMA writable_schema=ON",
4048 "CREATE TEMP TABLE t(x)",
4049 "SAVEPOINT nested_snapshot",
4050 "WITH x(v) AS (SELECT 99) INSERT INTO pool_backed_reader_probe(value) SELECT v FROM x",
4055 "WITH x(v) AS (SELECT 99) \
4056 INSERT INTO pool_backed_reader_probe(value) SELECT v FROM x RETURNING value",
4057 "WITH x(v) AS (SELECT 0) \
4058 UPDATE pool_backed_reader_probe SET value = value + (SELECT v FROM x)",
4059 "WITH x(v) AS (SELECT 1) \
4060 DELETE FROM pool_backed_reader_probe WHERE value = (SELECT v FROM x)",
4061 ] {
4062 let mut reader = bridge.reader().await.unwrap();
4063 let result = reader
4064 .query_all(SqlStatement {
4065 sql: probe.into(),
4066 params: vec![],
4067 label: Some("pool-backed-reader-admission-probe".into()),
4068 })
4069 .await;
4070 assert!(
4071 result.is_err(),
4072 "pool-backed reader capability must refuse {probe:?}; got {result:?}"
4073 );
4074 }
4075
4076 let rows = bridge
4079 .reader()
4080 .await
4081 .unwrap()
4082 .query_all(SqlStatement {
4083 sql: "SELECT value FROM pool_backed_reader_probe".into(),
4084 params: vec![],
4085 label: Some("pool-backed-reader-post-probe-read".into()),
4086 })
4087 .await
4088 .expect("a plain read must still work after every probe above was refused");
4089 assert_eq!(
4090 rows.len(),
4091 1,
4092 "a refused WITH-DML probe must not have committed a row"
4093 );
4094 }
4095
4096 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4107 async fn abandoned_deferred_read_transaction_span_is_rolled_back_before_reuse() {
4108 let config = PoolConfig {
4109 path: None,
4110 ..PoolConfig::default()
4111 };
4112 let pool = Arc::new(ConnectionPool::new(config).unwrap());
4113 pool.writer()
4114 .unwrap()
4115 .conn()
4116 .execute_batch(
4117 "CREATE TABLE abandoned_span_probe(value INTEGER NOT NULL); \
4118 INSERT INTO abandoned_span_probe VALUES (1);",
4119 )
4120 .unwrap();
4121 let bridge = SqlBridge::new(Arc::clone(&pool), false);
4122
4123 {
4124 let mut reader = bridge.reader().await.unwrap();
4125 reader
4126 .query_all(SqlStatement {
4127 sql: "BEGIN DEFERRED".into(),
4128 params: vec![],
4129 label: Some("abandoned-span-begin".into()),
4130 })
4131 .await
4132 .expect("the span must open");
4133 let failed = reader
4134 .query_all(SqlStatement {
4135 sql: "SELECT value FROM abandoned_span_probe_missing_table".into(),
4136 params: vec![],
4137 label: Some("abandoned-span-failing-read".into()),
4138 })
4139 .await;
4140 assert!(
4141 failed.is_err(),
4142 "the probe read against a nonexistent table must fail"
4143 );
4144 }
4147
4148 let writer = pool.writer().unwrap();
4149 assert!(
4150 writer.conn().is_autocommit(),
4151 "an abandoned deferred-read span must be rolled back before its connection \
4152 returns to service"
4153 );
4154 drop(writer);
4155
4156 let mut reader = bridge.reader().await.unwrap();
4157 let rows = reader
4158 .query_all(SqlStatement {
4159 sql: "SELECT value FROM abandoned_span_probe".into(),
4160 params: vec![],
4161 label: Some("abandoned-span-post-recovery-read".into()),
4162 })
4163 .await
4164 .expect("a fresh checkout must read normally after the abandoned span");
4165 assert_eq!(
4166 rows.len(),
4167 1,
4168 "the original row must be intact; the abandoned span must not have committed \
4169 anything"
4170 );
4171 }
4172
4173 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4181 async fn pooled_reader_checkout_timeout_is_a_retryable_admission_timeout() {
4182 let dir = tempfile::tempdir().unwrap();
4183 let config = PoolConfig {
4184 path: Some(dir.path().join("sql_bridge_reader_admission_timeout.db")),
4185 max_readers: 1,
4186 checkout_timeout: std::time::Duration::from_millis(200),
4187 ..PoolConfig::default()
4188 };
4189 let pool = Arc::new(ConnectionPool::new(config).unwrap());
4190 pool.writer()
4191 .unwrap()
4192 .conn()
4193 .execute_batch(
4194 "CREATE TABLE reader_admission_probe(value INTEGER NOT NULL); \
4195 INSERT INTO reader_admission_probe VALUES (0);",
4196 )
4197 .unwrap();
4198 let held_reader = pool.reader().expect("hold the sole pooled reader");
4201
4202 let bridge = SqlBridge::new(Arc::clone(&pool), true);
4203 let mut contender = bridge.reader().await.unwrap();
4204 let blocked = contender
4205 .query_row(SqlStatement {
4206 sql: "SELECT value FROM reader_admission_probe".into(),
4207 params: vec![],
4208 label: Some("reader-admission-timeout-probe".into()),
4209 })
4210 .await;
4211 assert!(
4212 matches!(blocked, Err(StorageError::AdmissionTimeout { .. })),
4213 "an exhausted pooled-reader checkout must be a retryable AdmissionTimeout; got {blocked:?}"
4214 );
4215
4216 drop(held_reader);
4217 tokio::time::sleep(std::time::Duration::from_millis(25)).await;
4218 assert_eq!(
4219 pool.available_readers(),
4220 1,
4221 "reader checkout leaked a permit"
4222 );
4223 }
4224
4225 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4226 async fn abandoned_read_interrupts_sqlite_releases_permit_and_stops_work() {
4227 let dir = tempfile::tempdir().unwrap();
4228 let config = PoolConfig {
4229 path: Some(dir.path().join("sql_bridge_abandoned_read.db")),
4230 max_readers: 1,
4231 checkout_timeout: std::time::Duration::from_millis(500),
4232 ..PoolConfig::default()
4233 };
4234 let pool = Arc::new(ConnectionPool::new(config).unwrap());
4235 let bridge = SqlBridge::new(Arc::clone(&pool), true);
4236 let mut reader = SqliteReader {
4237 handle: Some(
4238 open_explicit_read_transaction_handle(Arc::clone(&pool))
4239 .await
4240 .unwrap(),
4241 ),
4242 pool: Arc::clone(&pool),
4243 poisoned: false,
4244 };
4245 let mut contender = bridge.reader().await.unwrap();
4246 let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
4247 let progress_in_scope = Arc::clone(&progress);
4248
4249 let query = tokio::spawn(crate::scope_test_read_progress(
4250 progress_in_scope,
4251 async move { reader.query_all(deliberately_slow_read_statement()).await },
4252 ));
4253 wait_for_progress(progress.as_ref()).await;
4254 query.abort();
4255 assert!(matches!(query.await, Err(error) if error.is_cancelled()));
4256
4257 tokio::time::timeout(
4258 std::time::Duration::from_millis(500),
4259 contender.query_row(SqlStatement {
4260 sql: "SELECT 1".into(),
4261 params: vec![],
4262 label: None,
4263 }),
4264 )
4265 .await
4266 .expect("abandoned SQLite statement did not return the sole reader promptly")
4267 .expect("reader probe failed after cancellation");
4268
4269 let stopped_at = progress.load(std::sync::atomic::Ordering::SeqCst);
4270 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
4271 assert_eq!(
4272 progress.load(std::sync::atomic::Ordering::SeqCst),
4273 stopped_at,
4274 "SQLite progress kept advancing after the abandoned request returned its reader"
4275 );
4276 }
4277
4278 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4279 async fn request_deadline_interrupts_statement_without_outer_timeout() {
4280 let dir = tempfile::tempdir().unwrap();
4281 let config = PoolConfig {
4282 path: Some(dir.path().join("sql_bridge_request_deadline.db")),
4283 max_readers: 1,
4284 checkout_timeout: std::time::Duration::from_millis(500),
4285 ..PoolConfig::default()
4286 };
4287 let pool = Arc::new(ConnectionPool::new(config).unwrap());
4288 let bridge = SqlBridge::new(Arc::clone(&pool), true);
4289 let mut reader = bridge.reader().await.unwrap();
4290 let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
4291
4292 let result = crate::scope_test_read_progress(
4293 Arc::clone(&progress),
4294 crate::scope_request_read_deadline(std::time::Duration::from_millis(25), async move {
4295 reader.query_all(deliberately_slow_read_statement()).await
4296 }),
4297 )
4298 .await;
4299 assert!(
4300 matches!(result, Err(StorageError::Timeout { .. })),
4301 "deadline must surface as a typed timeout, got {result:?}"
4302 );
4303
4304 let stopped_at = progress.load(std::sync::atomic::Ordering::SeqCst);
4305 assert!(stopped_at > 0, "deadline test never exercised SQLite work");
4306 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
4307 assert_eq!(
4308 progress.load(std::sync::atomic::Ordering::SeqCst),
4309 stopped_at,
4310 "deadline returned while SQLite kept consuming work"
4311 );
4312 }
4313
4314 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4315 async fn progress_handler_cleanup_failure_discards_pooled_connection() {
4316 let dir = tempfile::tempdir().unwrap();
4317 let config = PoolConfig {
4318 path: Some(dir.path().join("sql_bridge_cleanup_failure.db")),
4319 max_readers: 1,
4320 ..PoolConfig::default()
4321 };
4322 let pool = Arc::new(ConnectionPool::new(config).unwrap());
4323 let pool_for_read = Arc::clone(&pool);
4324 let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
4325 let result = crate::read_cancellation::scope_test_read_cleanup_failure(
4326 crate::scope_test_read_progress(
4327 Arc::clone(&progress),
4328 crate::read_cancellation::run_interruptible_read(
4329 StorageCapability::Sql,
4330 "cleanup_failure_probe",
4331 move |scope| {
4332 let mut guard = pool_for_read.reader().map_err(|error| {
4333 StorageError::driver(
4334 StorageCapability::Sql,
4335 "cleanup_failure_probe",
4336 error,
4337 )
4338 })?;
4339 scope.run_pooled_reader(&mut guard, |conn| {
4340 conn.query_row("SELECT 1", [], |row| row.get::<_, i64>(0))
4341 .map_err(|error| map_rusqlite_err(error, "cleanup_failure_probe"))
4342 })
4343 },
4344 ),
4345 ),
4346 )
4347 .await;
4348 assert!(
4349 matches!(result, Err(StorageError::Internal(ref message)) if message.contains("clear failure")),
4350 "injected cleanup failure must be surfaced; got {result:?}"
4351 );
4352 assert_eq!(
4353 pool.available_readers(),
4354 1,
4355 "discard must install a replacement"
4356 );
4357
4358 let calls_after_failed_read = progress.load(std::sync::atomic::Ordering::SeqCst);
4359 let guard = pool.reader().unwrap();
4360 let sum: i64 = guard
4361 .conn()
4362 .query_row(
4363 "WITH RECURSIVE n(x) AS (VALUES(0) UNION ALL SELECT x + 1 FROM n WHERE x < 10000) \
4364 SELECT sum(x) FROM n",
4365 [],
4366 |row| row.get(0),
4367 )
4368 .unwrap();
4369 assert_eq!(sum, 50_005_000);
4370 assert_eq!(
4371 progress.load(std::sync::atomic::Ordering::SeqCst),
4372 calls_after_failed_read,
4373 "a connection whose handler could not be cleared was reused"
4374 );
4375 }
4376
4377 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4378 async fn raw_pooled_reader_quarantines_cleanup_failure_during_unwind() {
4379 let dir = tempfile::tempdir().unwrap();
4380 let pool = Arc::new(
4381 ConnectionPool::new(PoolConfig {
4382 path: Some(dir.path().join("raw_reader_unwind_cleanup.db")),
4383 max_readers: 1,
4384 ..PoolConfig::default()
4385 })
4386 .unwrap(),
4387 );
4388 let worker_pool = Arc::clone(&pool);
4389 let result = crate::read_cancellation::scope_test_read_cleanup_failure(
4390 crate::read_cancellation::run_interruptible_read(
4391 StorageCapability::Sql,
4392 "raw_reader_unwind_cleanup",
4393 move |scope| {
4394 let mut guard = worker_pool.reader().map_err(|error| {
4395 StorageError::driver(
4396 StorageCapability::Sql,
4397 "raw_reader_unwind_cleanup",
4398 error,
4399 )
4400 })?;
4401 scope.with_pooled_reader(&mut guard, |conn| {
4402 scope.run(conn, || -> khive_storage::types::StorageResult<()> {
4403 panic!("injected raw reader panic after progress registration")
4404 })
4405 })
4406 },
4407 ),
4408 )
4409 .await;
4410 assert!(
4411 result.is_err(),
4412 "blocking panic must surface as a join error"
4413 );
4414 assert_eq!(
4415 pool.available_readers(),
4416 pool.max_readers(),
4417 "unwind cleanup failure must close and replace the raw pooled reader"
4418 );
4419 }
4420
4421 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4422 async fn raw_pooled_writer_retires_cleanup_failure_during_unwind() {
4423 let pool = Arc::new(ConnectionPool::new(PoolConfig::default()).unwrap());
4424 let worker_pool = Arc::clone(&pool);
4425 let result = crate::read_cancellation::scope_test_read_cleanup_failure(
4426 crate::read_cancellation::run_interruptible_read(
4427 StorageCapability::Sql,
4428 "raw_writer_unwind_cleanup",
4429 move |scope| {
4430 let guard = worker_pool.try_writer().map_err(|error| {
4431 StorageError::driver(
4432 StorageCapability::Sql,
4433 "raw_writer_unwind_cleanup",
4434 error,
4435 )
4436 })?;
4437 scope.with_pooled_writer(&worker_pool, &guard, |conn| {
4438 scope.run(conn, || -> khive_storage::types::StorageResult<()> {
4439 panic!("injected raw writer panic after progress registration")
4440 })
4441 })
4442 },
4443 ),
4444 )
4445 .await;
4446 assert!(
4447 result.is_err(),
4448 "blocking panic must surface as a join error"
4449 );
4450 assert!(
4451 pool.try_writer().is_err(),
4452 "unwind cleanup failure must retire the raw pooled writer"
4453 );
4454 }
4455
4456 #[test]
4457 fn query_row_converts_only_the_first_matching_row() {
4458 let conn = rusqlite::Connection::open_in_memory().unwrap();
4459 let statement = SqlStatement {
4460 sql: "WITH RECURSIVE rows(value) AS (\
4461 SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 99\
4462 ) SELECT value FROM rows ORDER BY value"
4463 .into(),
4464 params: vec![],
4465 label: None,
4466 };
4467
4468 ROW_CONVERSIONS.with(|count| count.set(0));
4469 let row = execute_query_row(&conn, &statement).unwrap().unwrap();
4470
4471 assert!(matches!(row.get("value"), Some(SqlValue::Integer(0))));
4472 ROW_CONVERSIONS.with(|count| assert_eq!(count.get(), 1));
4473 }
4474
4475 #[test]
4476 fn query_page_bounds_owned_rows_before_full_materialization() {
4477 let conn = rusqlite::Connection::open_in_memory().unwrap();
4478 let statement = SqlStatement {
4479 sql: "WITH RECURSIVE rows(value) AS (\
4480 SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 99\
4481 ) SELECT value FROM rows ORDER BY value"
4482 .into(),
4483 params: vec![],
4484 label: None,
4485 };
4486
4487 ROW_CONVERSIONS.with(|count| count.set(0));
4488 let rows = execute_query_page(
4489 &conn,
4490 &statement,
4491 &PageRequest {
4492 offset: 40,
4493 limit: 3,
4494 },
4495 )
4496 .unwrap();
4497
4498 assert_eq!(rows.len(), 3);
4499 assert!(matches!(rows[0].get("value"), Some(SqlValue::Integer(40))));
4500 assert!(matches!(rows[2].get("value"), Some(SqlValue::Integer(42))));
4501 ROW_CONVERSIONS.with(|count| assert_eq!(count.get(), 3));
4502 }
4503
4504 #[test]
4505 fn query_page_zero_limit_converts_no_rows_but_still_validates_sql() {
4506 let conn = rusqlite::Connection::open_in_memory().unwrap();
4507 let statement = SqlStatement {
4508 sql: "WITH RECURSIVE rows(value) AS (\
4509 SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 99\
4510 ) SELECT value FROM rows ORDER BY value"
4511 .into(),
4512 params: vec![],
4513 label: None,
4514 };
4515
4516 ROW_CONVERSIONS.with(|count| count.set(0));
4517 let rows = execute_query_page(
4518 &conn,
4519 &statement,
4520 &PageRequest {
4521 offset: 0,
4522 limit: 0,
4523 },
4524 )
4525 .unwrap();
4526
4527 assert!(rows.is_empty());
4528 ROW_CONVERSIONS.with(|count| assert_eq!(count.get(), 0));
4529
4530 let invalid = SqlStatement {
4531 sql: "SELECT FROM WHERE".into(),
4532 params: vec![],
4533 label: None,
4534 };
4535 assert!(
4536 execute_query_page(
4537 &conn,
4538 &invalid,
4539 &PageRequest {
4540 offset: 0,
4541 limit: 0
4542 }
4543 )
4544 .is_err(),
4545 "a zero-limit page must still fail on invalid SQL at prepare time"
4546 );
4547 }
4548
4549 #[test]
4550 fn cached_writer_prepare_preserves_the_single_statement_boundary() {
4551 let conn = rusqlite::Connection::open_in_memory().unwrap();
4552 assert!(matches!(
4553 prepare_cached_sql_statement(&conn, "SELECT 1; SELECT 2"),
4554 Err(rusqlite::Error::MultipleStatement)
4555 ));
4556 }
4557
4558 #[tokio::test]
4559 async fn queue_backed_execute_reuses_the_persistent_connection_statement_cache() {
4560 use rusqlite::hooks::{AuthAction, AuthContext, Authorization};
4561 use std::sync::atomic::{AtomicUsize, Ordering};
4562
4563 let dir = tempfile::tempdir().unwrap();
4564 let pool = Arc::new(
4565 ConnectionPool::new(PoolConfig {
4566 path: Some(dir.path().join("sql_bridge_writer_cache.db")),
4567 write_queue_enabled: Some(true),
4568 write_routing_strict: true,
4569 ..PoolConfig::for_test()
4570 })
4571 .unwrap(),
4572 );
4573 pool.writer()
4574 .unwrap()
4575 .conn()
4576 .execute_batch(
4577 "CREATE TABLE writer_cache_test (id INTEGER PRIMARY KEY, value TEXT NOT NULL)",
4578 )
4579 .unwrap();
4580
4581 let writer_task = pool
4582 .writer_task_handle()
4583 .unwrap()
4584 .expect("file-backed queue-enabled pool must expose its writer task");
4585 let prepare_count = Arc::new(AtomicUsize::new(0));
4586 let hook_count = Arc::clone(&prepare_count);
4587 writer_task
4588 .send_top_level(move |conn| {
4589 conn.authorizer(Some(move |context: AuthContext<'_>| {
4590 if matches!(
4591 context.action,
4592 AuthAction::Insert { table_name } if table_name == "writer_cache_test"
4593 ) {
4594 hook_count.fetch_add(1, Ordering::SeqCst);
4595 }
4596 Authorization::Allow
4597 }))
4598 .map_err(|error| map_rusqlite_err(error, "test.install_authorizer"))
4599 })
4600 .await
4601 .unwrap();
4602
4603 let bridge = SqlBridge::new(Arc::clone(&pool), true);
4604 let mut writer = bridge.writer().await.unwrap();
4605 for id in [1, 2] {
4606 khive_storage::SqlWriter::execute(
4607 &mut *writer,
4608 SqlStatement {
4609 sql: "INSERT INTO writer_cache_test (id, value) VALUES (?1, ?2)".into(),
4610 params: vec![SqlValue::Integer(id), SqlValue::Text(format!("value-{id}"))],
4611 label: None,
4612 },
4613 )
4614 .await
4615 .unwrap();
4616 }
4617
4618 assert_eq!(
4619 prepare_count.load(Ordering::SeqCst),
4620 1,
4621 "the second identical execute on the writer task's persistent connection must reuse \
4622 the cached SQLite statement instead of compiling it again"
4623 );
4624 writer_task
4625 .send_top_level(|conn| {
4626 conn.authorizer(None::<fn(AuthContext<'_>) -> Authorization>)
4627 .map_err(|error| map_rusqlite_err(error, "test.remove_authorizer"))
4628 })
4629 .await
4630 .unwrap();
4631 }
4632
4633 #[test]
4634 fn inline_execute_batch_prepares_each_statement_once() {
4635 use rusqlite::hooks::{AuthAction, AuthContext, Authorization};
4636 use std::sync::atomic::{AtomicUsize, Ordering};
4637
4638 let conn = rusqlite::Connection::open_in_memory().unwrap();
4639 conn.execute_batch(
4640 "CREATE TABLE single_prepare_test (id INTEGER PRIMARY KEY, value TEXT NOT NULL)",
4641 )
4642 .unwrap();
4643 let prepare_count = Arc::new(AtomicUsize::new(0));
4644 let hook_count = Arc::clone(&prepare_count);
4645 conn.authorizer(Some(move |context: AuthContext<'_>| {
4646 if matches!(
4647 context.action,
4648 AuthAction::Insert { table_name } if table_name == "single_prepare_test"
4649 ) {
4650 hook_count.fetch_add(1, Ordering::SeqCst);
4651 }
4652 Authorization::Allow
4653 }))
4654 .unwrap();
4655
4656 let mut writer = InlineWriter {
4657 event_rows: None,
4658 conn: &conn as *const rusqlite::Connection,
4659 };
4660 let affected = block_on_sync(khive_storage::SqlWriter::execute_batch(
4661 &mut writer,
4662 vec![SqlStatement {
4663 sql: "INSERT INTO single_prepare_test (id, value) VALUES (?1, ?2)".into(),
4664 params: vec![SqlValue::Integer(1), SqlValue::Text("once".into())],
4665 label: None,
4666 }],
4667 ))
4668 .expect("InlineWriter operations must resolve on their first poll")
4669 .expect("valid batch must execute");
4670
4671 assert_eq!(affected, 1);
4672 assert_eq!(
4673 prepare_count.load(Ordering::SeqCst),
4674 1,
4675 "classification and execution must share one prepared statement handle"
4676 );
4677 conn.authorizer(None::<fn(AuthContext<'_>) -> Authorization>)
4678 .unwrap();
4679 }
4680
4681 #[test]
4682 fn inline_execute_batch_preserves_schema_dependencies_between_statements() {
4683 let conn = rusqlite::Connection::open_in_memory().unwrap();
4684 let mut writer = InlineWriter {
4685 event_rows: None,
4686 conn: &conn as *const rusqlite::Connection,
4687 };
4688
4689 let affected = block_on_sync(khive_storage::SqlWriter::execute_batch(
4690 &mut writer,
4691 vec![
4692 SqlStatement {
4693 sql: "CREATE TABLE dependent_prepare_test (id INTEGER PRIMARY KEY)".into(),
4694 params: vec![],
4695 label: None,
4696 },
4697 SqlStatement {
4698 sql: "INSERT INTO dependent_prepare_test (id) VALUES (1)".into(),
4699 params: vec![],
4700 label: None,
4701 },
4702 ],
4703 ))
4704 .expect("InlineWriter operations must resolve on their first poll")
4705 .expect("a later statement must be prepared after its prerequisite schema change");
4706
4707 assert_eq!(affected, 1);
4708 let count: i64 = conn
4709 .query_row("SELECT COUNT(*) FROM dependent_prepare_test", [], |row| {
4710 row.get(0)
4711 })
4712 .unwrap();
4713 assert_eq!(count, 1);
4714 }
4715
4716 #[tokio::test]
4717 async fn pool_backed_query_page_beyond_result_set_returns_empty() {
4718 let config = PoolConfig {
4719 path: None,
4720 ..PoolConfig::default()
4721 };
4722 let pool = Arc::new(ConnectionPool::new(config).unwrap());
4723 {
4724 let writer = pool.writer().unwrap();
4725 writer
4726 .conn()
4727 .execute_batch(
4728 "CREATE TABLE page_test (id INTEGER PRIMARY KEY, val TEXT NOT NULL);\
4729 INSERT INTO page_test (id, val) VALUES (1, 'a'), (2, 'b'), (3, 'c');",
4730 )
4731 .unwrap();
4732 }
4733 let bridge = SqlBridge::new(Arc::clone(&pool), false);
4734
4735 let statement = || SqlStatement {
4736 sql: "SELECT val FROM page_test ORDER BY id".into(),
4737 params: vec![],
4738 label: None,
4739 };
4740
4741 let mut reader = bridge.reader().await.unwrap();
4742 let page = reader
4743 .query_page(
4744 statement(),
4745 PageRequest {
4746 offset: 1,
4747 limit: 2,
4748 },
4749 )
4750 .await
4751 .unwrap();
4752 assert_eq!(page.len(), 2);
4753 assert!(matches!(page[0].get("val"), Some(SqlValue::Text(v)) if v == "b"));
4754 assert!(matches!(page[1].get("val"), Some(SqlValue::Text(v)) if v == "c"));
4755
4756 let empty = reader
4757 .query_page(
4758 statement(),
4759 PageRequest {
4760 offset: 99,
4761 limit: 10,
4762 },
4763 )
4764 .await
4765 .unwrap();
4766 assert!(
4767 empty.is_empty(),
4768 "offset past the last row must return an empty page, got {empty:?}"
4769 );
4770 drop(reader);
4771
4772 let mut writer = bridge.writer().await.unwrap();
4773 let empty = writer
4774 .query_page(
4775 statement(),
4776 PageRequest {
4777 offset: 99,
4778 limit: 10,
4779 },
4780 )
4781 .await
4782 .unwrap();
4783 assert!(
4784 empty.is_empty(),
4785 "offset past the last row must return an empty page, got {empty:?}"
4786 );
4787 }
4788
4789 #[tokio::test]
4790 async fn file_bridge_scopes_reader_permits_to_operations_and_caps_writer_handles() {
4791 let dir = tempfile::tempdir().unwrap();
4792 let config = PoolConfig {
4793 path: Some(dir.path().join("sql_bridge_handle_cap.db")),
4794 write_queue_enabled: Some(false),
4795 max_readers: 2,
4796 checkout_timeout: std::time::Duration::from_millis(20),
4797 ..PoolConfig::default()
4798 };
4799 let pool = Arc::new(ConnectionPool::new(config).unwrap());
4800 let bridge = SqlBridge::new(Arc::clone(&pool), true);
4801 let second_bridge = SqlBridge::new(Arc::clone(&pool), true);
4802
4803 let mut retained_readers = Vec::new();
4804 for expected in 0..3 {
4805 let mut reader = second_bridge.reader().await.unwrap();
4806 let value = reader
4807 .query_scalar(SqlStatement {
4808 sql: format!("SELECT {expected}"),
4809 params: vec![],
4810 label: None,
4811 })
4812 .await
4813 .unwrap();
4814 assert!(matches!(value, Some(SqlValue::Integer(value)) if value == expected));
4815 retained_readers.push(reader);
4816 }
4817 assert_eq!(retained_readers.len(), 3);
4818
4819 let mut additional_reader = bridge.reader().await.unwrap();
4820 let page = additional_reader
4821 .query_page(
4822 SqlStatement {
4823 sql: "WITH RECURSIVE rows(value) AS (\
4824 SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 9\
4825 ) SELECT value FROM rows ORDER BY value"
4826 .into(),
4827 params: vec![],
4828 label: None,
4829 },
4830 PageRequest {
4831 offset: 7,
4832 limit: 2,
4833 },
4834 )
4835 .await
4836 .unwrap();
4837 assert_eq!(page.len(), 2);
4838 assert!(matches!(page[0].get("value"), Some(SqlValue::Integer(7))));
4839 assert!(matches!(page[1].get("value"), Some(SqlValue::Integer(8))));
4840 drop((additional_reader, retained_readers));
4841
4842 let writer = bridge.writer().await.unwrap();
4843 let writer_error = match second_bridge.writer().await {
4844 Ok(_) => panic!("a second live writer handle exceeded the one-handle cap"),
4845 Err(error) => error,
4846 };
4847 assert!(matches!(
4848 writer_error,
4849 StorageError::AdmissionTimeout { ref operation, .. }
4850 if operation.as_ref() == "sql_bridge.writer_handle"
4851 ));
4852 drop(writer);
4853 let writer_after_release = bridge.writer().await.unwrap();
4854 drop(writer_after_release);
4855 }
4856
4857 #[tokio::test]
4858 #[serial_test::serial(tx_registry)]
4859 async fn file_bridge_attributes_ordinary_pool_reads_and_explicit_transaction_exception() {
4860 let dir = tempfile::tempdir().unwrap();
4861 let pool = Arc::new(
4862 ConnectionPool::new(PoolConfig {
4863 path: Some(dir.path().join("sql_bridge_reader_routes.db")),
4864 max_readers: 1,
4865 ..PoolConfig::default()
4866 })
4867 .unwrap(),
4868 );
4869 let bridge = SqlBridge::new(Arc::clone(&pool), true);
4870 let mut reader = bridge.reader().await.unwrap();
4871
4872 assert_eq!(
4873 pool.reader_acquisition_snapshot(),
4874 crate::pool::ReaderAcquisitionSnapshot {
4875 reader_admission_capacity: 1,
4876 available_reader_admission_slots: 1,
4877 ..crate::pool::ReaderAcquisitionSnapshot::default()
4878 },
4879 "constructing or retaining an idle raw-SQL reader must open nothing"
4880 );
4881
4882 let value = reader
4883 .query_scalar(SqlStatement {
4884 sql: "SELECT 1".into(),
4885 params: vec![],
4886 label: None,
4887 })
4888 .await
4889 .unwrap();
4890 assert!(matches!(value, Some(SqlValue::Integer(1))));
4891 let ordinary = pool.reader_acquisition_snapshot();
4892 assert_eq!(ordinary.pooled_checkouts, 1);
4893 assert_eq!(ordinary.completed_pooled_checkouts, 1);
4894 assert_eq!(ordinary.standalone_opens, 0);
4895
4896 reader
4897 .query_all(SqlStatement {
4898 sql: "BEGIN DEFERRED".into(),
4899 params: vec![],
4900 label: None,
4901 })
4902 .await
4903 .expect("the documented explicit transaction exception opens");
4904 let begun = pool.reader_acquisition_snapshot();
4905 assert_eq!(begun.pooled_checkouts, 1);
4906 assert_eq!(begun.standalone_opens, 1);
4907
4908 reader
4909 .query_scalar(SqlStatement {
4910 sql: "SELECT 2".into(),
4911 params: vec![],
4912 label: None,
4913 })
4914 .await
4915 .expect("transaction query reuses its one exceptional connection");
4916 reader
4917 .query_all(SqlStatement {
4918 sql: "COMMIT".into(),
4919 params: vec![],
4920 label: None,
4921 })
4922 .await
4923 .expect("COMMIT closes the explicit transaction exception");
4924 let committed = pool.reader_acquisition_snapshot();
4925 assert_eq!(committed.reader_admission_capacity, 1);
4926 assert_eq!(committed.available_reader_admission_slots, 1);
4927 assert_eq!(committed.acquisitions, begun.acquisitions);
4928 assert_eq!(committed.pooled_checkouts, begun.pooled_checkouts);
4929 assert_eq!(committed.standalone_opens, begun.standalone_opens);
4930 assert_eq!(
4931 committed.completed_pooled_checkouts, begun.completed_pooled_checkouts,
4932 "queries and COMMIT inside one explicit transaction must not acquire another reader"
4933 );
4934
4935 reader
4936 .query_scalar(SqlStatement {
4937 sql: "SELECT 3".into(),
4938 params: vec![],
4939 label: None,
4940 })
4941 .await
4942 .expect("ordinary traffic returns to the reader pool after COMMIT");
4943 let after = pool.reader_acquisition_snapshot();
4944 assert_eq!(after.pooled_checkouts, 2);
4945 assert_eq!(after.completed_pooled_checkouts, 2);
4946 assert_eq!(after.standalone_opens, 1);
4947 assert_eq!(after.active_pooled_checkouts, 0);
4948 }
4949
4950 #[tokio::test]
4951 async fn explicit_reader_open_timeout_is_visible_without_a_standalone_fallback() {
4952 let dir = tempfile::tempdir().unwrap();
4953 let pool = Arc::new(
4954 ConnectionPool::new(PoolConfig {
4955 path: Some(dir.path().join("sql_bridge_reader_open_timeout.db")),
4956 max_readers: 1,
4957 checkout_timeout: std::time::Duration::from_millis(20),
4958 ..PoolConfig::default()
4959 })
4960 .unwrap(),
4961 );
4962 let bridge = SqlBridge::new(Arc::clone(&pool), true);
4963 let mut logical_reader = bridge.reader().await.unwrap();
4964 let held = pool.reader().expect("hold the shared reader budget");
4965
4966 let blocked = logical_reader
4967 .query_all(SqlStatement {
4968 sql: "BEGIN DEFERRED".into(),
4969 params: vec![],
4970 label: None,
4971 })
4972 .await;
4973 assert!(
4974 matches!(
4975 &blocked,
4976 Err(StorageError::Timeout { operation })
4977 if operation.as_ref() == "sql_bridge.reader_open"
4978 ),
4979 "the compatible explicit-transaction open phase must stay visible; got {blocked:?}"
4980 );
4981
4982 let snapshot = pool.reader_acquisition_snapshot();
4983 assert_eq!(snapshot.checkout_timeouts, 1);
4984 assert_eq!(snapshot.pooled_checkouts, 1);
4985 assert_eq!(snapshot.standalone_opens, 0);
4986 assert_eq!(snapshot.active_pooled_checkouts, 1);
4987 assert_eq!(snapshot.available_reader_admission_slots, 0);
4988 drop(held);
4989 }
4990
4991 #[tokio::test]
4996 async fn pooled_raw_sql_read_counts_busy_handler_timeouts() {
4997 let dir = tempfile::tempdir().unwrap();
4998 let path = dir.path().join("sql_bridge_busy_timeouts.db");
4999 let pool = Arc::new(
5000 ConnectionPool::new(PoolConfig {
5001 path: Some(path.clone()),
5002 wal_mode: false,
5003 write_queue_enabled: Some(false),
5004 busy_timeout: std::time::Duration::from_millis(50),
5005 ..PoolConfig::default()
5006 })
5007 .unwrap(),
5008 );
5009 pool.writer()
5010 .unwrap()
5011 .conn()
5012 .execute_batch("CREATE TABLE busy_fixture (id INTEGER PRIMARY KEY)")
5013 .unwrap();
5014 let bridge = SqlBridge::new(Arc::clone(&pool), true);
5015 let mut reader = bridge.reader().await.unwrap();
5016 let count = || SqlStatement {
5017 sql: "SELECT count(*) FROM busy_fixture".into(),
5018 params: vec![],
5019 label: None,
5020 };
5021
5022 let value = reader.query_scalar(count()).await.unwrap();
5024 assert!(matches!(value, Some(SqlValue::Integer(0))), "{value:?}");
5025 assert_eq!(pool.reader_acquisition_snapshot().busy_timeouts, 0);
5026
5027 let holder = rusqlite::Connection::open(&path).unwrap();
5028 holder.execute_batch("BEGIN EXCLUSIVE").unwrap();
5029 let refused = reader.query_scalar(count()).await.unwrap_err();
5030 assert!(
5031 matches!(
5032 &refused,
5033 StorageError::Driver { source, .. }
5034 if source
5035 .downcast_ref::<rusqlite::Error>()
5036 .and_then(|error| error.sqlite_error_code())
5037 == Some(rusqlite::ErrorCode::DatabaseBusy)
5038 ),
5039 "a read behind an exclusive lock must surface SQLITE_BUSY, got {refused:?}"
5040 );
5041 let snapshot = pool.reader_acquisition_snapshot();
5042 assert_eq!(snapshot.busy_timeouts, 1);
5043 assert_eq!(
5044 snapshot.checkout_timeouts, 0,
5045 "a busy-handler refusal after checkout is not a checkout timeout"
5046 );
5047 holder.execute_batch("ROLLBACK").unwrap();
5048 }
5049
5050 #[tokio::test]
5051 #[serial_test::serial(tx_registry)]
5052 async fn cached_read_transaction_retains_one_permit_until_commit_or_rollback() {
5053 let dir = tempfile::tempdir().unwrap();
5054 let config = PoolConfig {
5055 path: Some(dir.path().join("sql_bridge_reader_tx_control.db")),
5056 write_queue_enabled: Some(true),
5057 max_readers: 1,
5058 checkout_timeout: std::time::Duration::from_millis(20),
5059 ..PoolConfig::default()
5060 };
5061 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5062 let origin = pool.origin();
5063 let origin_view = database_tx_view(&pool);
5064 let unrelated_view = khive_storage::tx_registry::TxOriginFilter::Secondary(
5065 khive_storage::tx_registry::DbIdentity::new("unrelated-sql-bridge.db"),
5066 );
5067 let bridge = SqlBridge::new(Arc::clone(&pool), true);
5068 let mut reader = bridge.reader().await.unwrap();
5069 let mut contender = bridge.reader().await.unwrap();
5070
5071 assert!(
5072 khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
5073 "an idle cached reader must not register a transaction"
5074 );
5075
5076 reader
5077 .query_all(SqlStatement {
5078 sql: "BEGIN DEFERRED".into(),
5079 params: vec![],
5080 label: None,
5081 })
5082 .await
5083 .expect("BEGIN DEFERRED must open an admitted cached-reader snapshot");
5084 let opened = khive_storage::tx_registry::oldest_for(&origin_view)
5085 .expect("successful BEGIN must register the cached-reader transaction");
5086 assert_eq!(opened.label.as_deref(), Some(CACHED_READ_TRANSACTION_LABEL));
5087 assert_eq!(opened.origin, origin);
5088 assert!(
5089 khive_storage::tx_registry::oldest_for(&unrelated_view).is_none(),
5090 "the read transaction must be attributed only to its own backend"
5091 );
5092 assert_eq!(
5093 pool.sql_bridge_reader_slots().available_permits(),
5094 0,
5095 "the successful BEGIN must retain its operation permit"
5096 );
5097
5098 let value = reader
5099 .query_scalar(SqlStatement {
5100 sql: "SELECT 7".into(),
5101 params: vec![],
5102 label: None,
5103 })
5104 .await
5105 .expect("a query inside the admitted transaction must reuse its retained permit");
5106 assert!(matches!(value, Some(SqlValue::Integer(7))));
5107 assert_eq!(
5108 khive_storage::tx_registry::oldest_for(&origin_view)
5109 .expect("queries must retain the transaction registration")
5110 .id,
5111 opened.id,
5112 "queries inside the transaction must retain the original span"
5113 );
5114
5115 let blocked = contender
5116 .query_scalar(SqlStatement {
5117 sql: "SELECT 8".into(),
5118 params: vec![],
5119 label: None,
5120 })
5121 .await;
5122 assert!(
5123 matches!(
5124 &blocked,
5125 Err(StorageError::AdmissionTimeout { operation, .. })
5126 if operation.as_ref() == "query_row"
5127 ),
5128 "a second logical read must contend with the admitted transaction \
5129 and fail at the bounded pooled-admission stage; got {blocked:?}"
5130 );
5131
5132 reader
5133 .query_all(SqlStatement {
5134 sql: "COMMIT".into(),
5135 params: vec![],
5136 label: None,
5137 })
5138 .await
5139 .expect("COMMIT must close the admitted cached-reader snapshot");
5140 assert!(
5141 khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
5142 "COMMIT must deregister after SQLite returns to autocommit"
5143 );
5144 assert_eq!(
5145 pool.sql_bridge_reader_slots().available_permits(),
5146 1,
5147 "COMMIT may release the permit only after autocommit is restored"
5148 );
5149 let value = contender
5150 .query_scalar(SqlStatement {
5151 sql: "SELECT 8".into(),
5152 params: vec![],
5153 label: None,
5154 })
5155 .await
5156 .expect("the contender must run after COMMIT releases admission");
5157 assert!(matches!(value, Some(SqlValue::Integer(8))));
5158
5159 reader
5160 .query_all(SqlStatement {
5161 sql: "BEGIN TRANSACTION".into(),
5162 params: vec![],
5163 label: None,
5164 })
5165 .await
5166 .expect("plain deferred BEGIN TRANSACTION must also be admitted");
5167 let reopened = khive_storage::tx_registry::oldest_for(&origin_view)
5168 .expect("the second successful BEGIN must register a fresh span");
5169 assert_ne!(reopened.id, opened.id);
5170 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 0);
5171 let nested = reader
5172 .query_all(SqlStatement {
5173 sql: "ROLLBACK TO stale_snapshot".into(),
5174 params: vec![],
5175 label: None,
5176 })
5177 .await;
5178 assert!(
5179 matches!(&nested, Err(StorageError::InvalidInput { .. })),
5180 "ROLLBACK TO requires unsupported nested state; got {nested:?}"
5181 );
5182 assert_eq!(
5183 pool.sql_bridge_reader_slots().available_permits(),
5184 0,
5185 "rejected nested control must not release the still-live transaction admission"
5186 );
5187 assert_eq!(
5188 khive_storage::tx_registry::oldest_for(&origin_view)
5189 .expect("ROLLBACK TO rejection must retain the live span")
5190 .id,
5191 reopened.id
5192 );
5193 let savepoint = reader
5194 .query_all(SqlStatement {
5195 sql: "SAVEPOINT nested_snapshot".into(),
5196 params: vec![],
5197 label: None,
5198 })
5199 .await;
5200 assert!(
5201 matches!(&savepoint, Err(StorageError::InvalidInput { .. })),
5202 "SAVEPOINT must be rejected inside the admitted transaction; got {savepoint:?}"
5203 );
5204 assert_eq!(
5205 khive_storage::tx_registry::oldest_for(&origin_view)
5206 .expect("SAVEPOINT rejection must retain the live span")
5207 .id,
5208 reopened.id
5209 );
5210 reader
5211 .query_all(SqlStatement {
5212 sql: "ROLLBACK".into(),
5213 params: vec![],
5214 label: None,
5215 })
5216 .await
5217 .expect("ROLLBACK must close the admitted cached-reader snapshot");
5218 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
5219 assert!(
5220 khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
5221 "full ROLLBACK must deregister after SQLite returns to autocommit"
5222 );
5223 }
5224
5225 #[tokio::test]
5226 async fn failed_cached_reader_begin_does_not_register_a_transaction() {
5227 use rusqlite::hooks::{AuthAction, AuthContext, Authorization, TransactionOperation};
5228
5229 fn deny_begin(ctx: AuthContext<'_>) -> Authorization {
5230 match ctx.action {
5231 AuthAction::Transaction {
5232 operation: TransactionOperation::Begin,
5233 } => Authorization::Deny,
5234 _ => Authorization::Allow,
5235 }
5236 }
5237
5238 let dir = tempfile::tempdir().unwrap();
5239 let config = PoolConfig {
5240 path: Some(dir.path().join("sql_bridge_reader_failed_begin.db")),
5241 max_readers: 1,
5242 ..PoolConfig::default()
5243 };
5244 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5245 let origin_view = database_tx_view(&pool);
5246 let conn = open_standalone_reader(&pool).unwrap();
5247 conn.authorizer(Some(deny_begin)).unwrap();
5248 let mut reader = SqliteReader {
5249 handle: Some(StandaloneHandle {
5250 conn,
5251 _retained_slot: None,
5252 read_transaction_slot: None,
5253 }),
5254 pool: Arc::clone(&pool),
5255 poisoned: false,
5256 };
5257
5258 let begin = reader
5259 .query_all(SqlStatement {
5260 sql: "BEGIN DEFERRED".into(),
5261 params: vec![],
5262 label: None,
5263 })
5264 .await;
5265 assert!(begin.is_err(), "the authorizer must reject BEGIN");
5266 assert!(
5267 khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
5268 "a failed BEGIN must never enter the transaction registry"
5269 );
5270 assert_eq!(
5271 pool.sql_bridge_reader_slots().available_permits(),
5272 1,
5273 "a failed BEGIN must return the operation permit"
5274 );
5275 }
5276
5277 #[tokio::test]
5278 #[serial_test::serial(tx_registry)]
5279 async fn failed_cached_reader_rollback_deregisters_only_when_connection_is_discarded() {
5280 use rusqlite::hooks::{AuthAction, AuthContext, Authorization, TransactionOperation};
5281
5282 fn deny_rollback(ctx: AuthContext<'_>) -> Authorization {
5283 match ctx.action {
5284 AuthAction::Transaction {
5285 operation: TransactionOperation::Rollback,
5286 } => Authorization::Deny,
5287 _ => Authorization::Allow,
5288 }
5289 }
5290
5291 let dir = tempfile::tempdir().unwrap();
5292 let config = PoolConfig {
5293 path: Some(dir.path().join("sql_bridge_reader_failed_rollback.db")),
5294 max_readers: 1,
5295 ..PoolConfig::default()
5296 };
5297 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5298 let origin_view = database_tx_view(&pool);
5299 let conn = open_standalone_reader(&pool).unwrap();
5300 let mut reader = SqliteReader {
5301 handle: Some(StandaloneHandle {
5302 conn,
5303 _retained_slot: None,
5304 read_transaction_slot: None,
5305 }),
5306 pool: Arc::clone(&pool),
5307 poisoned: false,
5308 };
5309
5310 reader
5311 .query_all(SqlStatement {
5312 sql: "BEGIN DEFERRED".into(),
5313 params: vec![],
5314 label: None,
5315 })
5316 .await
5317 .expect("BEGIN must establish the registered transaction");
5318 let opened = khive_storage::tx_registry::oldest_for(&origin_view)
5319 .expect("the admitted transaction must be registered");
5320 reader
5321 .handle
5322 .as_ref()
5323 .expect("reader must retain its connection")
5324 .conn
5325 .authorizer(Some(deny_rollback))
5326 .unwrap();
5327
5328 let rollback = reader
5329 .query_all(SqlStatement {
5330 sql: "ROLLBACK".into(),
5331 params: vec![],
5332 label: None,
5333 })
5334 .await;
5335 assert!(rollback.is_err(), "the authorizer must reject ROLLBACK");
5336 assert_eq!(
5337 khive_storage::tx_registry::oldest_for(&origin_view)
5338 .expect("failed ROLLBACK must retain registry evidence")
5339 .id,
5340 opened.id
5341 );
5342 assert_eq!(
5343 pool.sql_bridge_reader_slots().available_permits(),
5344 0,
5345 "failed ROLLBACK must retain reader admission"
5346 );
5347
5348 drop(reader);
5349 assert!(
5350 khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
5351 "discarding the connection must not leak its registry entry"
5352 );
5353 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
5354 }
5355
5356 #[tokio::test]
5357 #[serial_test::serial(tx_registry)]
5358 async fn cached_read_only_handles_reject_unsupported_transaction_control_without_consumption() {
5359 let dir = tempfile::tempdir().unwrap();
5360 let config = PoolConfig {
5361 path: Some(
5362 dir.path()
5363 .join("sql_bridge_reader_unsupported_tx_control.db"),
5364 ),
5365 write_queue_enabled: Some(true),
5366 max_readers: 1,
5367 ..PoolConfig::default()
5368 };
5369 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5370 let bridge = SqlBridge::new(Arc::clone(&pool), true);
5371
5372 let mut reader = bridge.reader().await.unwrap();
5373 for (sql, keyword) in [
5374 ("BEGIN IMMEDIATE", "BEGIN"),
5375 ("BEGIN EXCLUSIVE", "BEGIN"),
5376 ("BEGIN TRANSACTION IMMEDIATE", "BEGIN"),
5381 ("BEGIN TRANSACTION EXCLUSIVE", "BEGIN"),
5382 ("BEGIN DEFERRED TRANSACTION trailing", "BEGIN"),
5383 ("BEGIN TRANSACTION \"IMMEDIATE\"", "BEGIN"),
5388 ("BEGIN TRANSACTION [IMMEDIATE]", "BEGIN"),
5389 ("BEGIN TRANSACTION `IMMEDIATE`", "BEGIN"),
5390 ("BEGIN TRANSACTION 'IMMEDIATE'", "BEGIN"),
5391 ("BEGIN \"DEFERRED\"", "BEGIN"),
5392 ("BEGIN; COMMIT", "BEGIN"),
5393 ("START TRANSACTION", "START"),
5394 ("COMMIT", "COMMIT"),
5395 ] {
5396 let rejected = reader
5397 .query_all(SqlStatement {
5398 sql: sql.into(),
5399 params: vec![],
5400 label: None,
5401 })
5402 .await;
5403 assert!(
5404 matches!(
5405 &rejected,
5406 Err(StorageError::InvalidInput { message, .. })
5407 if message.contains(keyword)
5408 ),
5409 "unsupported cached-reader control {sql:?} must fail closed; got {rejected:?}"
5410 );
5411 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
5412 }
5413
5414 let mut queue_backed_writer = bridge.writer().await.unwrap();
5415 let rejected = queue_backed_writer
5416 .query_all(SqlStatement {
5417 sql: "SAVEPOINT stale_snapshot".into(),
5418 params: vec![],
5419 label: None,
5420 })
5421 .await;
5422 assert!(
5423 matches!(
5424 &rejected,
5425 Err(StorageError::InvalidInput {
5426 operation,
5427 message,
5428 ..
5429 }) if operation.as_ref() == "writer.query_all"
5430 && message.contains("transaction control")
5431 && message.contains("SAVEPOINT")
5432 ),
5433 "a queue-backed writer without an explicit read transaction must reject nested \
5434 transaction control; got {rejected:?}"
5435 );
5436 assert_eq!(
5437 pool.sql_bridge_reader_slots().available_permits(),
5438 1,
5439 "queue-backed rejection must leave the operation permit available"
5440 );
5441 let value = queue_backed_writer
5442 .query_scalar(SqlStatement {
5443 sql: "SELECT 8".into(),
5444 params: vec![],
5445 label: None,
5446 })
5447 .await
5448 .expect("transaction-control rejection must not consume the queue-backed handle");
5449 assert!(matches!(value, Some(SqlValue::Integer(8))));
5450
5451 queue_backed_writer
5452 .query_all(SqlStatement {
5453 sql: "BEGIN DEFERRED".into(),
5454 params: vec![],
5455 label: None,
5456 })
5457 .await
5458 .expect("queue-backed cached reader must share explicit read-transaction admission");
5459 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 0);
5460 queue_backed_writer
5461 .query_all(SqlStatement {
5462 sql: "END".into(),
5463 params: vec![],
5464 label: None,
5465 })
5466 .await
5467 .expect("END must release queue-backed cached-reader admission");
5468 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
5469 }
5470
5471 #[tokio::test]
5472 #[serial_test::serial(tx_registry)]
5473 async fn dropping_cached_reader_transaction_closes_snapshot_before_releasing_permit() {
5474 let dir = tempfile::tempdir().unwrap();
5475 let config = PoolConfig {
5476 path: Some(dir.path().join("sql_bridge_reader_tx_drop.db")),
5477 max_readers: 1,
5478 checkout_timeout: std::time::Duration::from_millis(20),
5479 ..PoolConfig::default()
5480 };
5481 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5482 let origin_view = database_tx_view(&pool);
5483 let bridge = SqlBridge::new(Arc::clone(&pool), true);
5484 let mut reader = bridge.reader().await.unwrap();
5485 let mut contender = bridge.reader().await.unwrap();
5486
5487 reader
5488 .query_all(SqlStatement {
5489 sql: "BEGIN DEFERRED".into(),
5490 params: vec![],
5491 label: None,
5492 })
5493 .await
5494 .expect("begin admitted transaction");
5495 reader
5496 .query_all(SqlStatement {
5497 sql: "SELECT * FROM sqlite_schema".into(),
5498 params: vec![],
5499 label: None,
5500 })
5501 .await
5502 .expect("materialize read snapshot");
5503 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 0);
5504 assert!(
5505 khive_storage::tx_registry::oldest_for(&origin_view).is_some(),
5506 "the live snapshot must remain registered until handle drop"
5507 );
5508
5509 drop(reader);
5510 assert!(
5511 khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
5512 "handle drop must close SQLite before deregistering the snapshot"
5513 );
5514 assert_eq!(
5515 pool.sql_bridge_reader_slots().available_permits(),
5516 1,
5517 "dropping the handle must close its transaction before returning admission"
5518 );
5519 contender
5520 .query_all(SqlStatement {
5521 sql: "SELECT * FROM sqlite_schema".into(),
5522 params: vec![],
5523 label: None,
5524 })
5525 .await
5526 .expect("a new operation must run after the transactional handle drops");
5527 }
5528
5529 #[tokio::test]
5538 #[serial_test::serial(tx_registry)]
5539 async fn expired_cached_reader_transaction_is_rolled_back_on_reuse() {
5540 let dir = tempfile::tempdir().unwrap();
5541 let config = PoolConfig {
5542 path: Some(dir.path().join("sql_bridge_reader_tx_max_age.db")),
5543 max_readers: 1,
5544 checkout_timeout: std::time::Duration::from_millis(20),
5545 read_tx_max_age: std::time::Duration::from_millis(20),
5546 ..PoolConfig::default()
5547 };
5548 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5549 let origin_view = database_tx_view(&pool);
5550 let bridge = SqlBridge::new(Arc::clone(&pool), true);
5551 let mut reader = bridge.reader().await.unwrap();
5552
5553 reader
5554 .query_all(SqlStatement {
5555 sql: "BEGIN DEFERRED".into(),
5556 params: vec![],
5557 label: None,
5558 })
5559 .await
5560 .expect("begin admitted transaction");
5561 reader
5562 .query_all(SqlStatement {
5563 sql: "SELECT * FROM sqlite_schema".into(),
5564 params: vec![],
5565 label: None,
5566 })
5567 .await
5568 .expect("materialize read snapshot");
5569 assert!(
5570 khive_storage::tx_registry::oldest_for(&origin_view).is_some(),
5571 "the open transaction must be registered before it ages out"
5572 );
5573
5574 tokio::time::sleep(std::time::Duration::from_millis(40)).await;
5575
5576 let evictions_before = crate::checkpoint::read_tx_max_age_evictions();
5577 let error = reader
5578 .query_all(SqlStatement {
5579 sql: "SELECT * FROM sqlite_schema".into(),
5580 params: vec![],
5581 label: None,
5582 })
5583 .await
5584 .expect_err("reusing a transaction past read_tx_max_age must be refused");
5585 assert!(
5586 error.is_retryable(),
5587 "an evicted-transaction error must be retryable so the caller can open a fresh \
5588 snapshot: {error}"
5589 );
5590 match &error {
5591 StorageError::ReadTransactionAgeEvicted {
5592 operation,
5593 max_age_secs,
5594 } => {
5595 assert_eq!(operation.as_ref(), "query_all");
5596 assert_eq!(
5597 *max_age_secs, 0,
5598 "a 20ms read_tx_max_age truncates to 0 whole seconds"
5599 );
5600 }
5601 other => panic!(
5602 "a clean age-triggered rollback must surface the dedicated \
5603 ReadTransactionAgeEvicted variant, not a generic classification: {other:?}"
5604 ),
5605 }
5606 assert_eq!(
5607 crate::checkpoint::read_tx_max_age_evictions(),
5608 evictions_before + 1,
5609 "the eviction must be counted in the #1846 diagnostics gauge"
5610 );
5611 assert!(
5612 khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
5613 "the expired transaction must be rolled back and deregistered rather than \
5614 continuing to pin the WAL snapshot"
5615 );
5616
5617 reader
5618 .query_all(SqlStatement {
5619 sql: "SELECT * FROM sqlite_schema".into(),
5620 params: vec![],
5621 label: None,
5622 })
5623 .await
5624 .expect("the handle must remain usable for a fresh autocommit read after eviction");
5625 }
5626
5627 #[tokio::test]
5636 #[serial_test::serial(tx_registry)]
5637 async fn expired_cached_reader_transaction_rollback_denial_discards_connection_and_releases_admission(
5638 ) {
5639 use rusqlite::hooks::{AuthAction, AuthContext, Authorization, TransactionOperation};
5640
5641 fn deny_rollback(ctx: AuthContext<'_>) -> Authorization {
5642 match ctx.action {
5643 AuthAction::Transaction {
5644 operation: TransactionOperation::Rollback,
5645 } => Authorization::Deny,
5646 _ => Authorization::Allow,
5647 }
5648 }
5649
5650 let dir = tempfile::tempdir().unwrap();
5651 let config = PoolConfig {
5652 path: Some(
5653 dir.path()
5654 .join("sql_bridge_reader_tx_max_age_rollback_denied.db"),
5655 ),
5656 max_readers: 1,
5657 checkout_timeout: std::time::Duration::from_millis(20),
5658 read_tx_max_age: std::time::Duration::from_millis(20),
5659 ..PoolConfig::default()
5660 };
5661 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5662 let origin_view = database_tx_view(&pool);
5663 let conn = open_standalone_reader(&pool).unwrap();
5664 let mut reader = SqliteReader {
5665 handle: Some(StandaloneHandle {
5666 conn,
5667 _retained_slot: None,
5668 read_transaction_slot: None,
5669 }),
5670 pool: Arc::clone(&pool),
5671 poisoned: false,
5672 };
5673
5674 reader
5675 .query_all(SqlStatement {
5676 sql: "BEGIN DEFERRED".into(),
5677 params: vec![],
5678 label: None,
5679 })
5680 .await
5681 .expect("begin admitted transaction");
5682 reader
5683 .query_all(SqlStatement {
5684 sql: "SELECT * FROM sqlite_schema".into(),
5685 params: vec![],
5686 label: None,
5687 })
5688 .await
5689 .expect("materialize read snapshot");
5690 assert!(
5691 khive_storage::tx_registry::oldest_for(&origin_view).is_some(),
5692 "the open transaction must be registered before it ages out"
5693 );
5694
5695 reader
5696 .handle
5697 .as_ref()
5698 .expect("reader must retain its connection")
5699 .conn
5700 .authorizer(Some(deny_rollback))
5701 .unwrap();
5702
5703 tokio::time::sleep(std::time::Duration::from_millis(40)).await;
5704
5705 let evictions_before = crate::checkpoint::read_tx_max_age_evictions();
5706 let error = reader
5707 .query_all(SqlStatement {
5708 sql: "SELECT * FROM sqlite_schema".into(),
5709 params: vec![],
5710 label: None,
5711 })
5712 .await
5713 .expect_err("a denied rollback on an expired transaction must surface an error");
5714 assert!(
5715 error.is_retryable(),
5716 "even a failed cleanup rollback must remain classified retryable so callers open a \
5717 fresh handle: {error}"
5718 );
5719 match &error {
5720 StorageError::ReadTransactionAgeEvictionCleanupFailed {
5721 operation,
5722 max_age_secs,
5723 message,
5724 } => {
5725 assert_eq!(operation.as_ref(), "query_all");
5726 assert_eq!(
5727 *max_age_secs, 0,
5728 "a 20ms read_tx_max_age truncates to 0 whole seconds"
5729 );
5730 assert!(
5731 message.contains("rollback failed"),
5732 "the failure must be attributable to the denied ROLLBACK, not silent \
5733 success: {message}"
5734 );
5735 }
5736 other => panic!(
5737 "a denied cleanup rollback must surface the dedicated \
5738 ReadTransactionAgeEvictionCleanupFailed variant, not a generic Transaction \
5739 error the caller cannot machine-detect: {other:?}"
5740 ),
5741 }
5742 assert_eq!(
5743 crate::checkpoint::read_tx_max_age_evictions(),
5744 evictions_before + 1,
5745 "the eviction attempt must still be counted even though cleanup failed"
5746 );
5747 assert!(
5748 khive_storage::tx_registry::oldest_for(&origin_view).is_none(),
5749 "a denied rollback must discard the connection and deregister the expired \
5750 transaction span rather than leaking it"
5751 );
5752 assert_eq!(
5753 pool.sql_bridge_reader_slots().available_permits(),
5754 1,
5755 "discarding the poisoned connection must release the reader admission slot"
5756 );
5757
5758 let reuse = reader
5759 .query_all(SqlStatement {
5760 sql: "SELECT * FROM sqlite_schema".into(),
5761 params: vec![],
5762 label: None,
5763 })
5764 .await;
5765 let message = match reuse {
5766 Err(StorageError::Pool { message, .. }) => message,
5767 other => panic!(
5768 "reusing this discarded reader must fail loudly with 'connection already \
5769 consumed' rather than silently reopening; got {other:?}"
5770 ),
5771 };
5772 assert!(
5773 message.contains("connection already consumed"),
5774 "expected the discarded reader's reuse error to name the pinned failure; got \
5775 {message:?}"
5776 );
5777 }
5778
5779 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5780 #[serial_test::serial(tx_registry)]
5781 async fn cancelled_cached_reader_transaction_releases_guards_after_connection_closes() {
5782 let dir = tempfile::tempdir().unwrap();
5783 let config = PoolConfig {
5784 path: Some(dir.path().join("sql_bridge_reader_tx_drop_cancel.db")),
5785 max_readers: 1,
5786 checkout_timeout: std::time::Duration::from_millis(50),
5787 ..PoolConfig::default()
5788 };
5789 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5790 let origin_view = database_tx_view(&pool);
5791 let bridge = SqlBridge::new(Arc::clone(&pool), true);
5792 let mut reader = SqliteReader {
5793 handle: Some(
5794 open_explicit_read_transaction_handle(Arc::clone(&pool))
5795 .await
5796 .unwrap(),
5797 ),
5798 pool: Arc::clone(&pool),
5799 poisoned: false,
5800 };
5801 let mut contender = bridge.reader().await.unwrap();
5802
5803 reader
5804 .query_all(SqlStatement {
5805 sql: "BEGIN DEFERRED".into(),
5806 params: vec![],
5807 label: None,
5808 })
5809 .await
5810 .expect("begin admitted transaction");
5811 assert!(
5812 khive_storage::tx_registry::oldest_for(&origin_view).is_some(),
5813 "the explicit transaction must be registered before cancellation"
5814 );
5815 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 0);
5816
5817 let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
5818 let query = tokio::spawn(crate::scope_test_read_progress(
5819 Arc::clone(&progress),
5820 async move { reader.query_all(deliberately_slow_read_statement()).await },
5821 ));
5822 wait_for_progress(progress.as_ref()).await;
5823 query.abort();
5824 assert!(matches!(query.await, Err(error) if error.is_cancelled()));
5825 tokio::time::timeout(std::time::Duration::from_secs(1), async {
5826 while khive_storage::tx_registry::oldest_for(&origin_view).is_some()
5827 || pool.sql_bridge_reader_slots().available_permits() != 1
5828 {
5829 tokio::task::yield_now().await;
5830 }
5831 })
5832 .await
5833 .expect("connection cleanup leaked transaction evidence or reader admission");
5834 contender
5835 .query_all(SqlStatement {
5836 sql: "SELECT 1".into(),
5837 params: vec![],
5838 label: None,
5839 })
5840 .await
5841 .expect("admission must recover after the cancelled connection closes");
5842 }
5843
5844 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5845 #[serial_test::serial(tx_registry)]
5846 async fn cancelled_cached_reader_rolls_back_releases_wal_and_clears_handler() {
5847 let dir = tempfile::tempdir().unwrap();
5848 let config = PoolConfig {
5849 path: Some(dir.path().join("sql_bridge_reader_tx_request_cancel.db")),
5850 max_readers: 1,
5851 checkout_timeout: std::time::Duration::from_millis(500),
5852 ..PoolConfig::default()
5853 };
5854 let pool = Arc::new(ConnectionPool::new(config).unwrap());
5855 let bridge = SqlBridge::new(Arc::clone(&pool), true);
5856 let writer = open_standalone_writer(&pool).unwrap();
5857 writer
5858 .execute_batch(
5859 "CREATE TABLE snapshot_probe(id INTEGER PRIMARY KEY, value TEXT NOT NULL); \
5860 INSERT INTO snapshot_probe(value) VALUES ('seed');",
5861 )
5862 .unwrap();
5863 let mut reader = bridge.reader().await.unwrap();
5864 let mut contender = bridge.reader().await.unwrap();
5865
5866 reader
5867 .query_all(SqlStatement {
5868 sql: "BEGIN DEFERRED".into(),
5869 params: vec![],
5870 label: None,
5871 })
5872 .await
5873 .expect("begin admitted transaction");
5874 reader
5875 .query_all(SqlStatement {
5876 sql: "SELECT * FROM snapshot_probe".into(),
5877 params: vec![],
5878 label: None,
5879 })
5880 .await
5881 .expect("materialize a real WAL snapshot");
5882 writer
5883 .execute_batch(
5884 "WITH RECURSIVE rows(value) AS (\
5885 SELECT 1 UNION ALL SELECT value + 1 FROM rows WHERE value < 100\
5886 ) INSERT INTO snapshot_probe(value) SELECT printf('row-%d', value) FROM rows;",
5887 )
5888 .unwrap();
5889 let (_, log_before, checkpointed_before) = passive_checkpoint(&writer);
5890 assert!(
5891 log_before > checkpointed_before,
5892 "the explicit reader snapshot must pin WAL frames before cancellation"
5893 );
5894
5895 let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
5896 let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
5897 let query = tokio::spawn(crate::scope_test_read_progress(
5898 Arc::clone(&progress),
5899 crate::scope_request_read_cancellation(cancel_rx, async move {
5900 let result = reader.query_all(deliberately_slow_read_statement()).await;
5901 (reader, result)
5902 }),
5903 ));
5904 wait_for_progress(progress.as_ref()).await;
5905 cancel_tx.send(true).unwrap();
5906 let (mut reader, result) = tokio::time::timeout(std::time::Duration::from_secs(1), query)
5907 .await
5908 .expect("interrupted explicit read transaction did not stop promptly")
5909 .unwrap();
5910 assert!(
5911 matches!(result, Err(StorageError::Timeout { .. })),
5912 "request cancellation must surface as a typed timeout; got {result:?}"
5913 );
5914 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
5915
5916 let (_, log_after, checkpointed_after) = passive_checkpoint(&writer);
5917 assert_eq!(
5918 log_after, checkpointed_after,
5919 "cancellation must release the explicit reader's WAL snapshot"
5920 );
5921 contender
5922 .query_all(SqlStatement {
5923 sql: "SELECT 1".into(),
5924 params: vec![],
5925 label: None,
5926 })
5927 .await
5928 .expect("the sole reader permit must be reusable after rollback");
5929
5930 let stopped_at = progress.load(std::sync::atomic::Ordering::SeqCst);
5931 reader
5932 .query_all(SqlStatement {
5933 sql: "WITH RECURSIVE rows(value) AS (\
5934 SELECT 0 UNION ALL SELECT value + 1 FROM rows WHERE value < 10000\
5935 ) SELECT SUM(value) FROM rows"
5936 .into(),
5937 params: vec![],
5938 label: None,
5939 })
5940 .await
5941 .expect("same connection must remain usable after handler teardown");
5942 assert_eq!(
5943 progress.load(std::sync::atomic::Ordering::SeqCst),
5944 stopped_at,
5945 "the cancelled request's progress callback bled into the next borrower"
5946 );
5947 }
5948
5949 fn register_khive_test_slow_udf(
5957 conn: &rusqlite::Connection,
5958 sleep_ms: u64,
5959 started: Arc<std::sync::atomic::AtomicBool>,
5960 ) {
5961 conn.create_scalar_function(
5962 "khive_test_slow_udf",
5963 0,
5964 rusqlite::functions::FunctionFlags::SQLITE_UTF8,
5965 move |_| {
5966 started.store(true, std::sync::atomic::Ordering::Release);
5967 std::thread::sleep(std::time::Duration::from_millis(sleep_ms));
5968 Ok(0i64)
5969 },
5970 )
5971 .unwrap();
5972 }
5973
5974 async fn wait_for_flag(flag: &std::sync::atomic::AtomicBool) {
5975 tokio::time::timeout(std::time::Duration::from_secs(1), async {
5976 while !flag.load(std::sync::atomic::Ordering::Acquire) {
5977 tokio::task::yield_now().await;
5978 }
5979 })
5980 .await
5981 .expect("slow UDF never started");
5982 }
5983
5984 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5985 async fn abandoned_read_past_grace_recovers_admission_after_bounded_join() {
5986 let dir = tempfile::tempdir().unwrap();
5994 let config = PoolConfig {
5995 path: Some(dir.path().join("sql_bridge_grace_exceeded.db")),
5996 max_readers: 1,
5997 checkout_timeout: std::time::Duration::from_millis(2_000),
5998 ..PoolConfig::default()
5999 };
6000 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6001 let writer = open_standalone_writer(&pool).unwrap();
6002 writer
6003 .execute_batch(
6004 "CREATE TABLE grace_probe(id INTEGER PRIMARY KEY, value TEXT NOT NULL); \
6005 INSERT INTO grace_probe(value) VALUES ('seed');",
6006 )
6007 .unwrap();
6008
6009 let mut reader = SqliteReader {
6010 handle: Some(
6011 open_explicit_read_transaction_handle(Arc::clone(&pool))
6012 .await
6013 .unwrap(),
6014 ),
6015 pool: Arc::clone(&pool),
6016 poisoned: false,
6017 };
6018 let mut contender = SqliteReader {
6019 handle: Some(
6020 open_explicit_read_transaction_handle(Arc::clone(&pool))
6021 .await
6022 .unwrap(),
6023 ),
6024 pool: Arc::clone(&pool),
6025 poisoned: false,
6026 };
6027 let udf_started = Arc::new(std::sync::atomic::AtomicBool::new(false));
6028 register_khive_test_slow_udf(
6029 &reader.handle.as_ref().unwrap().conn,
6030 900,
6031 Arc::clone(&udf_started),
6032 );
6033
6034 reader
6035 .query_all(SqlStatement {
6036 sql: "BEGIN DEFERRED".into(),
6037 params: vec![],
6038 label: None,
6039 })
6040 .await
6041 .expect("begin admitted transaction");
6042 reader
6043 .query_all(SqlStatement {
6044 sql: "SELECT * FROM grace_probe".into(),
6045 params: vec![],
6046 label: None,
6047 })
6048 .await
6049 .expect("materialize a real WAL snapshot");
6050 writer
6051 .execute_batch(
6052 "WITH RECURSIVE rows(value) AS (\
6053 SELECT 1 UNION ALL SELECT value + 1 FROM rows WHERE value < 100\
6054 ) INSERT INTO grace_probe(value) SELECT printf('row-%d', value) FROM rows;",
6055 )
6056 .unwrap();
6057 let (_, log_before, checkpointed_before) = passive_checkpoint(&writer);
6058 assert!(
6059 log_before > checkpointed_before,
6060 "the explicit reader snapshot must pin WAL frames before cancellation"
6061 );
6062
6063 let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
6064 let query = tokio::spawn(crate::scope_request_read_cancellation(
6065 cancel_rx,
6066 async move {
6067 let result = reader
6068 .query_all(SqlStatement {
6069 sql: "SELECT khive_test_slow_udf()".into(),
6070 params: vec![],
6071 label: None,
6072 })
6073 .await;
6074 (reader, result)
6075 },
6076 ));
6077 wait_for_flag(udf_started.as_ref()).await;
6083 cancel_tx.send(true).unwrap();
6084
6085 let (mut reader, result) = tokio::time::timeout(std::time::Duration::from_secs(3), query)
6086 .await
6087 .expect(
6088 "a worker that settles within the grace+hard-cap bound must not hang the caller",
6089 )
6090 .unwrap();
6091 assert!(
6092 matches!(result, Err(StorageError::Timeout { .. })),
6093 "request cancellation must still surface as a typed timeout even after grace \
6094 was exceeded; got {result:?}"
6095 );
6096
6097 assert_eq!(
6103 pool.sql_bridge_reader_slots().available_permits(),
6104 1,
6105 "the sole reader permit must be visible again once the bounded join completes"
6106 );
6107
6108 let (_, log_after, checkpointed_after) = passive_checkpoint(&writer);
6109 assert_eq!(
6110 log_after, checkpointed_after,
6111 "the abandoned explicit read transaction must release its WAL snapshot by the \
6112 time the caller observes the timeout"
6113 );
6114
6115 contender
6116 .query_all(SqlStatement {
6117 sql: "SELECT 1".into(),
6118 params: vec![],
6119 label: None,
6120 })
6121 .await
6122 .expect("a fresh reader must be admitted once the zombie worker has settled");
6123
6124 reader
6128 .query_all(SqlStatement {
6129 sql: "SELECT 1".into(),
6130 params: vec![],
6131 label: None,
6132 })
6133 .await
6134 .expect("the interrupted connection must remain usable after settling");
6135 }
6136
6137 #[tokio::test]
6138 #[serial_test::serial(tx_registry)]
6139 async fn cached_reader_transaction_lifecycle_survives_sqlite_empty_prefixes() {
6140 let dir = tempfile::tempdir().unwrap();
6141 let config = PoolConfig {
6142 path: Some(dir.path().join("sql_bridge_reader_prefixed_tx_control.db")),
6143 write_queue_enabled: Some(true),
6144 max_readers: 1,
6145 ..PoolConfig::default()
6146 };
6147 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6148 let bridge = SqlBridge::new(Arc::clone(&pool), true);
6149 let mut reader = bridge.reader().await.unwrap();
6150
6151 reader
6152 .query_all(SqlStatement {
6153 sql: " ; -- empty statement\n /* leading comment */ \u{feff} BEGIN DEFERRED".into(),
6154 params: vec![],
6155 label: None,
6156 })
6157 .await
6158 .expect("prefixed BEGIN must enter the admitted transaction state");
6159 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 0);
6160 reader
6161 .query_all(SqlStatement {
6162 sql: " /* leading comment */ \u{feff} ; COMMIT".into(),
6163 params: vec![],
6164 label: None,
6165 })
6166 .await
6167 .expect("prefixed COMMIT must end the admitted transaction state");
6168 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
6169
6170 let rejected = reader
6171 .query_all(SqlStatement {
6172 sql: " ; /* no active transaction */ \u{feff} COMMIT".into(),
6173 params: vec![],
6174 label: None,
6175 })
6176 .await;
6177 assert!(
6178 matches!(
6179 &rejected,
6180 Err(StorageError::InvalidInput {
6181 operation,
6182 message,
6183 ..
6184 }) if operation.as_ref() == "query_all"
6185 && message.contains("transaction control")
6186 && message.contains("COMMIT")
6187 ),
6188 "a prefixed COMMIT without an admitted transaction must still fail closed; \
6189 got {rejected:?}"
6190 );
6191
6192 let mut queue_backed_writer = bridge.writer().await.unwrap();
6193 let rejected = queue_backed_writer
6194 .query_all(SqlStatement {
6195 sql: "-- leading comment\n \u{feff} ; /* empty */ SAVEPOINT pinned".into(),
6196 params: vec![],
6197 label: None,
6198 })
6199 .await;
6200 assert!(
6201 matches!(
6202 &rejected,
6203 Err(StorageError::InvalidInput {
6204 operation,
6205 message,
6206 ..
6207 }) if operation.as_ref() == "writer.query_all"
6208 && message.contains("transaction control")
6209 && message.contains("SAVEPOINT")
6210 ),
6211 "a queue-backed cached reader must classify transaction control through \
6212 comments, BOMs, and empty statements; got {rejected:?}"
6213 );
6214
6215 let value = reader
6216 .query_scalar(SqlStatement {
6217 sql: "SELECT 10".into(),
6218 params: vec![],
6219 label: None,
6220 })
6221 .await
6222 .expect("prefixed transaction lifecycle must preserve the cached reader");
6223 assert!(matches!(value, Some(SqlValue::Integer(10))));
6224 }
6225
6226 #[tokio::test]
6227 async fn cached_reader_restores_autocommit_before_releasing_its_operation_permit() {
6228 let dir = tempfile::tempdir().unwrap();
6229 let config = PoolConfig {
6230 path: Some(dir.path().join("sql_bridge_reader_autocommit.db")),
6231 max_readers: 1,
6232 ..PoolConfig::default()
6233 };
6234 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6235 let conn = open_standalone_reader(&pool).unwrap();
6236 conn.execute_batch("BEGIN DEFERRED; SELECT * FROM sqlite_schema")
6237 .unwrap();
6238 assert!(
6239 !conn.is_autocommit(),
6240 "the regression precondition needs a live read transaction"
6241 );
6242 let mut reader = SqliteReader {
6243 handle: Some(StandaloneHandle {
6244 conn,
6245 _retained_slot: None,
6246 read_transaction_slot: None,
6247 }),
6248 pool: Arc::clone(&pool),
6249 poisoned: false,
6250 };
6251
6252 let rejected = reader
6253 .query_all(SqlStatement {
6254 sql: "ROLLBACK".into(),
6258 params: vec![],
6259 label: None,
6260 })
6261 .await;
6262 assert!(
6263 matches!(
6264 &rejected,
6265 Err(StorageError::InvalidInput {
6266 operation,
6267 message,
6268 ..
6269 }) if operation.as_ref() == "query_all"
6270 && message.contains("outside autocommit")
6271 ),
6272 "a cached reader that reaches the boundary outside autocommit must fail closed; \
6273 got {rejected:?}"
6274 );
6275 assert_eq!(
6276 pool.sql_bridge_reader_slots().available_permits(),
6277 1,
6278 "the permit may be released only after the stale transaction is gone"
6279 );
6280 assert!(
6281 reader.handle.is_none(),
6282 "the restored connection must close instead of surviving as an idle standalone cache"
6283 );
6284
6285 let value = reader
6286 .query_scalar(SqlStatement {
6287 sql: "SELECT 9".into(),
6288 params: vec![],
6289 label: None,
6290 })
6291 .await
6292 .expect("the cleaned reader must remain usable through the pooled route");
6293 assert!(matches!(value, Some(SqlValue::Integer(9))));
6294 }
6295
6296 #[tokio::test]
6297 async fn standalone_writer_read_preserves_manual_atomic_transaction() {
6298 let dir = tempfile::tempdir().unwrap();
6299 let config = PoolConfig {
6300 path: Some(dir.path().join("sql_bridge_writer_atomic_read.db")),
6301 write_queue_enabled: Some(false),
6302 max_readers: 1,
6303 ..PoolConfig::default()
6304 };
6305 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6306 {
6307 let writer = pool.writer().unwrap();
6308 writer
6309 .conn()
6310 .execute_batch(
6311 "CREATE TABLE atomic_read_test \
6312 (id INTEGER PRIMARY KEY, value TEXT NOT NULL)",
6313 )
6314 .unwrap();
6315 }
6316 let bridge = SqlBridge::new(Arc::clone(&pool), true);
6317
6318 let observed = bridge
6319 .atomic_unit(Box::new(|writer| {
6320 Box::pin(async move {
6321 writer
6322 .execute(SqlStatement {
6323 sql: "INSERT INTO atomic_read_test (id, value) VALUES (1, 'pending')"
6324 .into(),
6325 params: vec![],
6326 label: None,
6327 })
6328 .await?;
6329 let count = writer
6330 .query_scalar(SqlStatement {
6331 sql: "SELECT COUNT(*) FROM atomic_read_test".into(),
6332 params: vec![],
6333 label: None,
6334 })
6335 .await?;
6336 Ok(Box::new(count) as Box<dyn std::any::Any + Send>)
6337 })
6338 }))
6339 .await
6340 .expect("manual atomic read must not be mistaken for an idle reader snapshot");
6341 let observed = match observed.downcast::<Option<SqlValue>>() {
6342 Ok(observed) => observed,
6343 Err(_) => panic!("unexpected atomic result type"),
6344 };
6345 assert!(matches!(*observed, Some(SqlValue::Integer(1))));
6346
6347 let mut reader = bridge.reader().await.unwrap();
6348 let committed = reader
6349 .query_scalar(SqlStatement {
6350 sql: "SELECT COUNT(*) FROM atomic_read_test".into(),
6351 params: vec![],
6352 label: None,
6353 })
6354 .await
6355 .unwrap();
6356 assert!(matches!(committed, Some(SqlValue::Integer(1))));
6357 }
6358
6359 #[tokio::test]
6360 async fn request_cancellation_preserves_file_backed_manual_atomic_read_and_commit() {
6361 let dir = tempfile::tempdir().unwrap();
6362 let pool = Arc::new(
6363 ConnectionPool::new(PoolConfig {
6364 path: Some(dir.path().join("sql_bridge_writer_tx_cancel.db")),
6365 write_queue_enabled: Some(false),
6366 ..PoolConfig::for_test()
6367 })
6368 .unwrap(),
6369 );
6370 pool.writer()
6371 .unwrap()
6372 .conn()
6373 .execute_batch(
6374 "CREATE TABLE writer_tx_cancel_probe(\
6375 id INTEGER PRIMARY KEY, value TEXT NOT NULL)",
6376 )
6377 .unwrap();
6378 let bridge = SqlBridge::new(Arc::clone(&pool), true);
6379 let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
6380
6381 let observed = crate::scope_request_read_cancellation(
6382 cancel_rx,
6383 bridge.atomic_unit(Box::new(move |writer| {
6384 Box::pin(async move {
6385 writer
6386 .execute(SqlStatement {
6387 sql: "INSERT INTO writer_tx_cancel_probe VALUES (1, 'before')".into(),
6388 params: vec![],
6389 label: None,
6390 })
6391 .await?;
6392 cancel_tx.send(true).unwrap();
6393 let count = writer
6394 .query_scalar(SqlStatement {
6395 sql: "SELECT COUNT(*) FROM writer_tx_cancel_probe".into(),
6396 params: vec![],
6397 label: None,
6398 })
6399 .await?;
6400 writer
6401 .execute(SqlStatement {
6402 sql: "INSERT INTO writer_tx_cancel_probe VALUES (2, 'after')".into(),
6403 params: vec![],
6404 label: None,
6405 })
6406 .await?;
6407 Ok(Box::new(count) as Box<dyn std::any::Any + Send>)
6408 })
6409 })),
6410 )
6411 .await
6412 .expect("request cancellation must not interrupt an admitted manual write transaction");
6413 let observed = match observed.downcast::<Option<SqlValue>>() {
6414 Ok(observed) => observed,
6415 Err(_) => panic!("unexpected atomic result type"),
6416 };
6417 assert!(matches!(*observed, Some(SqlValue::Integer(1))));
6418
6419 let reader = pool.reader().unwrap();
6420 let rows: i64 = reader
6421 .conn()
6422 .query_row("SELECT COUNT(*) FROM writer_tx_cancel_probe", [], |row| {
6423 row.get(0)
6424 })
6425 .unwrap();
6426 assert_eq!(rows, 2, "both writes around the SELECT must commit");
6427 }
6428
6429 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6430 async fn cancelled_standalone_writer_transaction_retains_active_reader_admission() {
6431 let dir = tempfile::tempdir().unwrap();
6432 let pool = Arc::new(
6433 ConnectionPool::new(PoolConfig {
6434 path: Some(dir.path().join("sql_bridge_writer_tx_admission.db")),
6435 write_queue_enabled: Some(false),
6436 max_readers: 1,
6437 checkout_timeout: std::time::Duration::from_millis(250),
6438 volume_lock_dir: Some(dir.path().join("volume-locks")),
6442 ..PoolConfig::default()
6443 })
6444 .unwrap(),
6445 );
6446 let writer_slot = pool
6447 .sql_bridge_writer_slots()
6448 .acquire_owned()
6449 .await
6450 .unwrap();
6451 let conn = open_standalone_writer(&pool).unwrap();
6452 let (entered, release, _completed) = blocking_non_interrupting_progress_gate(&conn);
6453 let unit_lease = acquire_unit_lease(Arc::clone(&pool)).await.unwrap();
6456 let mut writer = SqliteWriter {
6457 observe_direct_errors: true,
6458 event_rows: None,
6459 handle: Some(StandaloneHandle {
6460 conn,
6461 _retained_slot: Some(writer_slot),
6462 read_transaction_slot: None,
6463 }),
6464 writer_task: None,
6465 origin: pool.origin(),
6466 db: crate::timeout_sink::db_label(&pool),
6467 pool: Arc::clone(&pool),
6468 held_lease: unit_lease,
6469 };
6470 khive_storage::SqlWriter::execute(
6471 &mut writer,
6472 SqlStatement {
6473 sql: "BEGIN IMMEDIATE".into(),
6474 params: vec![],
6475 label: None,
6476 },
6477 )
6478 .await
6479 .unwrap();
6480 let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
6481 cancel_tx.send(true).unwrap();
6482
6483 let query = tokio::spawn(crate::scope_request_read_cancellation(
6484 cancel_rx,
6485 async move {
6486 let result =
6487 khive_storage::SqlReader::query_all(&mut writer, progress_gate_statement())
6488 .await;
6489 let rollback = khive_storage::SqlWriter::execute(
6490 &mut writer,
6491 SqlStatement {
6492 sql: "ROLLBACK".into(),
6493 params: vec![],
6494 label: None,
6495 },
6496 )
6497 .await;
6498 (result, rollback)
6499 },
6500 ));
6501 tokio::time::timeout(std::time::Duration::from_secs(1), entered.notified())
6502 .await
6503 .expect("cancelled writer-transaction SELECT never reached SQLite");
6504 assert_eq!(
6505 pool.sql_bridge_reader_slots().available_permits(),
6506 0,
6507 "a writer-supertrait SELECT must retain ordinary active-reader admission"
6508 );
6509
6510 tokio::task::spawn_blocking(move || release.wait())
6511 .await
6512 .unwrap();
6513 let (rows, rollback) = tokio::time::timeout(std::time::Duration::from_secs(2), query)
6514 .await
6515 .expect("writer-transaction SELECT did not finish after its gate opened")
6516 .unwrap();
6517 assert_eq!(
6518 rows.expect("request cancellation interrupted the admitted writer transaction")
6519 .len(),
6520 1
6521 );
6522 rollback.expect("writer transaction did not return to autocommit");
6523 assert_eq!(pool.sql_bridge_reader_slots().available_permits(), 1);
6524 }
6525
6526 #[tokio::test]
6527 async fn expired_deadline_preserves_pool_backed_manual_atomic_read_and_commit() {
6528 let pool = Arc::new(ConnectionPool::new(PoolConfig::default()).unwrap());
6529 pool.writer()
6530 .unwrap()
6531 .conn()
6532 .execute_batch(
6533 "CREATE TABLE pool_writer_tx_deadline_probe(\
6534 id INTEGER PRIMARY KEY, value TEXT NOT NULL)",
6535 )
6536 .unwrap();
6537 let bridge = SqlBridge::new(Arc::clone(&pool), false);
6538
6539 let observed = crate::scope_request_read_deadline(
6540 std::time::Duration::ZERO,
6541 bridge.atomic_unit(Box::new(|writer| {
6542 Box::pin(async move {
6543 writer
6544 .execute(SqlStatement {
6545 sql: "INSERT INTO pool_writer_tx_deadline_probe VALUES (1, 'before')"
6546 .into(),
6547 params: vec![],
6548 label: None,
6549 })
6550 .await?;
6551 let count = writer
6552 .query_scalar(SqlStatement {
6553 sql: "SELECT COUNT(*) FROM pool_writer_tx_deadline_probe".into(),
6554 params: vec![],
6555 label: None,
6556 })
6557 .await?;
6558 writer
6559 .execute(SqlStatement {
6560 sql: "INSERT INTO pool_writer_tx_deadline_probe VALUES (2, 'after')"
6561 .into(),
6562 params: vec![],
6563 label: None,
6564 })
6565 .await?;
6566 Ok(Box::new(count) as Box<dyn std::any::Any + Send>)
6567 })
6568 })),
6569 )
6570 .await
6571 .expect("an expired read deadline must not interrupt an admitted manual write transaction");
6572 let observed = match observed.downcast::<Option<SqlValue>>() {
6573 Ok(observed) => observed,
6574 Err(_) => panic!("unexpected atomic result type"),
6575 };
6576 assert!(matches!(*observed, Some(SqlValue::Integer(1))));
6577
6578 let reader = pool.reader().unwrap();
6579 let rows: i64 = reader
6580 .conn()
6581 .query_row(
6582 "SELECT COUNT(*) FROM pool_writer_tx_deadline_probe",
6583 [],
6584 |row| row.get(0),
6585 )
6586 .unwrap();
6587 assert_eq!(rows, 2, "both writes around the SELECT must commit");
6588 }
6589
6590 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6591 async fn cancelled_standalone_open_retains_slot_until_open_finishes() {
6592 let dir = tempfile::tempdir().unwrap();
6593 let config = PoolConfig {
6594 path: Some(dir.path().join("sql_bridge_cancelled_open.db")),
6595 max_readers: 1,
6596 checkout_timeout: std::time::Duration::from_millis(250),
6597 ..PoolConfig::default()
6598 };
6599 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6600 let slots = pool.sql_bridge_reader_slots();
6601 let slot = Arc::clone(&slots).acquire_owned().await.unwrap();
6602 assert_eq!(slots.available_permits(), 0);
6603
6604 let (entered_tx, entered_rx) = std::sync::mpsc::channel();
6605 let (release_tx, release_rx) = std::sync::mpsc::channel();
6606 let open = tokio::spawn(open_standalone_on_blocking(
6607 Arc::clone(&pool),
6608 slot,
6609 "test_open_reader",
6610 move |pool| {
6611 entered_tx.send(()).unwrap();
6612 release_rx.recv().unwrap();
6613 open_standalone_reader(pool)
6614 },
6615 ));
6616 tokio::task::spawn_blocking(move || entered_rx.recv())
6617 .await
6618 .unwrap()
6619 .unwrap();
6620
6621 open.abort();
6622 assert!(matches!(open.await, Err(error) if error.is_cancelled()));
6623 assert_eq!(
6624 slots.available_permits(),
6625 0,
6626 "the permit must remain in the detached open closure"
6627 );
6628 let contender = tokio::time::timeout(
6629 std::time::Duration::from_millis(50),
6630 Arc::clone(&slots).acquire_owned(),
6631 )
6632 .await;
6633 assert!(contender.is_err(), "an in-flight open must retain the cap");
6634
6635 release_tx.send(()).unwrap();
6636 let recovered = tokio::time::timeout(
6637 std::time::Duration::from_secs(1),
6638 Arc::clone(&slots).acquire_owned(),
6639 )
6640 .await
6641 .expect("the detached open did not release its permit")
6642 .unwrap();
6643 assert_eq!(slots.available_permits(), 0);
6644 drop(recovered);
6645 assert_eq!(slots.available_permits(), 1);
6646 }
6647
6648 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6649 async fn abandoned_writer_read_interrupts_and_releases_writer_handle() {
6650 let dir = tempfile::tempdir().unwrap();
6651 let config = PoolConfig {
6652 path: Some(dir.path().join("sql_bridge_cancelled_writer.db")),
6653 write_queue_enabled: Some(false),
6654 checkout_timeout: std::time::Duration::from_millis(250),
6655 ..PoolConfig::for_test()
6656 };
6657 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6658 let bridge = SqlBridge::new(Arc::clone(&pool), true);
6659
6660 let handle_slot = pool
6661 .sql_bridge_writer_slots()
6662 .acquire_owned()
6663 .await
6664 .unwrap();
6665 let conn = open_standalone_writer(&pool).unwrap();
6666 let mut writer = SqliteWriter {
6667 observe_direct_errors: true,
6668 event_rows: None,
6669 handle: Some(StandaloneHandle {
6670 conn,
6671 _retained_slot: Some(handle_slot),
6672 read_transaction_slot: None,
6673 }),
6674 writer_task: None,
6675 origin: pool.origin(),
6676 db: crate::timeout_sink::db_label(&pool),
6677 pool: Arc::clone(&pool),
6678 held_lease: None,
6679 };
6680 let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));
6681 let query = tokio::spawn(crate::scope_test_read_progress(
6682 Arc::clone(&progress),
6683 async move {
6684 khive_storage::SqlReader::query_all(&mut writer, deliberately_slow_read_statement())
6685 .await
6686 },
6687 ));
6688
6689 wait_for_progress(progress.as_ref()).await;
6690 query.abort();
6691 assert!(matches!(query.await, Err(error) if error.is_cancelled()));
6692 let writer_after =
6693 tokio::time::timeout(std::time::Duration::from_millis(500), bridge.writer())
6694 .await
6695 .expect("abandoned SQLite read did not release the writer handle promptly")
6696 .expect("writer handle remained unavailable after read interruption");
6697 drop(writer_after);
6698 let stopped_at = progress.load(std::sync::atomic::Ordering::SeqCst);
6699 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
6700 assert_eq!(
6701 progress.load(std::sync::atomic::Ordering::SeqCst),
6702 stopped_at,
6703 "writer-backed SQLite read kept consuming work after cancellation"
6704 );
6705 }
6706
6707 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6708 async fn request_cancellation_never_interrupts_admitted_execute_batch() {
6709 let dir = tempfile::tempdir().unwrap();
6710 let config = PoolConfig {
6711 path: Some(dir.path().join("sql_bridge_cancelled_writer_batch.db")),
6712 write_queue_enabled: Some(false),
6713 checkout_timeout: std::time::Duration::from_millis(250),
6714 ..PoolConfig::for_test()
6715 };
6716 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6717 let bridge = SqlBridge::new(Arc::clone(&pool), true);
6718 {
6719 let guard = pool.writer().unwrap();
6720 guard
6721 .conn()
6722 .execute_batch(
6723 "CREATE TABLE cancellation_write_probe(\
6724 id INTEGER PRIMARY KEY, value INTEGER NOT NULL)",
6725 )
6726 .unwrap();
6727 }
6728
6729 let handle_slot = pool
6730 .sql_bridge_writer_slots()
6731 .acquire_owned()
6732 .await
6733 .unwrap();
6734 let conn = open_standalone_writer(&pool).unwrap();
6735 let (entered, release, completed) = blocking_non_interrupting_progress_gate(&conn);
6736 let mut writer = SqliteWriter {
6737 observe_direct_errors: true,
6738 event_rows: None,
6739 handle: Some(StandaloneHandle {
6740 conn,
6741 _retained_slot: Some(handle_slot),
6742 read_transaction_slot: None,
6743 }),
6744 writer_task: None,
6745 origin: pool.origin(),
6746 db: crate::timeout_sink::db_label(&pool),
6747 pool: Arc::clone(&pool),
6748 held_lease: None,
6749 };
6750 let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
6751 let query = tokio::spawn(crate::scope_request_read_cancellation(
6752 cancel_rx,
6753 async move {
6754 khive_storage::SqlWriter::execute_batch(&mut writer, vec![slow_insert_statement()])
6755 .await
6756 },
6757 ));
6758
6759 tokio::time::timeout(std::time::Duration::from_secs(1), entered.notified())
6760 .await
6761 .expect("mutating execute_batch never reached SQLite VM work");
6762 cancel_tx.send(true).unwrap();
6763 tokio::time::sleep(std::time::Duration::from_millis(25)).await;
6764 assert!(
6765 !query.is_finished(),
6766 "request-read cancellation must not interrupt an admitted batch"
6767 );
6768
6769 let contender = bridge.writer().await;
6770 let retained_slot = matches!(
6771 &contender,
6772 Err(StorageError::AdmissionTimeout { operation, .. })
6773 if operation.as_ref() == "sql_bridge.writer_handle"
6774 );
6775 drop(contender);
6776
6777 tokio::task::spawn_blocking(move || release.wait())
6778 .await
6779 .unwrap();
6780 let affected = tokio::time::timeout(std::time::Duration::from_secs(2), query)
6781 .await
6782 .expect("admitted batch did not finish after its gate was released")
6783 .unwrap()
6784 .expect("request cancellation must preserve the batch result");
6785 assert_eq!(affected, 10_000);
6786 tokio::time::timeout(std::time::Duration::from_secs(1), completed.notified())
6787 .await
6788 .expect("completed batch did not release its connection");
6789 assert!(
6790 retained_slot,
6791 "request cancellation released the writer slot before the admitted batch stopped"
6792 );
6793 let reader = pool.reader().unwrap();
6794 let count: i64 = reader
6795 .conn()
6796 .query_row("SELECT COUNT(*) FROM cancellation_write_probe", [], |row| {
6797 row.get(0)
6798 })
6799 .unwrap();
6800 assert_eq!(count, 10_000, "the admitted batch must commit every row");
6801 }
6802
6803 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6804 async fn request_cancellation_never_interrupts_dml_returning_via_sql_reader() {
6805 let dir = tempfile::tempdir().unwrap();
6806 let config = PoolConfig {
6807 path: Some(dir.path().join("sql_bridge_dml_returning_cancel.db")),
6808 write_queue_enabled: Some(false),
6809 checkout_timeout: std::time::Duration::from_millis(250),
6810 ..PoolConfig::for_test()
6811 };
6812 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6813 {
6814 let guard = pool.writer().unwrap();
6815 guard
6816 .conn()
6817 .execute_batch(
6818 "CREATE TABLE returning_write_probe(\
6819 id INTEGER PRIMARY KEY, value INTEGER NOT NULL)",
6820 )
6821 .unwrap();
6822 }
6823
6824 let handle_slot = pool
6825 .sql_bridge_writer_slots()
6826 .acquire_owned()
6827 .await
6828 .unwrap();
6829 let conn = open_standalone_writer(&pool).unwrap();
6830 let (entered, release, completed) = blocking_non_interrupting_progress_gate(&conn);
6831 let mut writer = SqliteWriter {
6832 observe_direct_errors: true,
6833 event_rows: None,
6834 handle: Some(StandaloneHandle {
6835 conn,
6836 _retained_slot: Some(handle_slot),
6837 read_transaction_slot: None,
6838 }),
6839 writer_task: None,
6840 origin: pool.origin(),
6841 db: crate::timeout_sink::db_label(&pool),
6842 pool: Arc::clone(&pool),
6843 held_lease: None,
6844 };
6845 let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
6846 let query = tokio::spawn(crate::scope_request_read_cancellation(
6847 cancel_rx,
6848 async move {
6849 khive_storage::SqlReader::query_all(
6850 &mut writer,
6851 SqlStatement {
6852 sql: "INSERT INTO returning_write_probe(value) \
6853 WITH RECURSIVE rows(value) AS (\
6854 SELECT 1 UNION ALL SELECT value + 1 FROM rows WHERE value < 10000\
6855 ) SELECT value FROM rows RETURNING id"
6856 .into(),
6857 params: vec![],
6858 label: Some("non-interruptible-returning-probe".into()),
6859 },
6860 )
6861 .await
6862 },
6863 ));
6864
6865 tokio::time::timeout(std::time::Duration::from_secs(1), entered.notified())
6866 .await
6867 .expect("DML RETURNING never reached admitted SQLite work");
6868 cancel_tx.send(true).unwrap();
6869 tokio::time::sleep(std::time::Duration::from_millis(25)).await;
6870 assert!(
6871 !query.is_finished(),
6872 "request-read cancellation interrupted DML RETURNING"
6873 );
6874
6875 tokio::task::spawn_blocking(move || release.wait())
6876 .await
6877 .unwrap();
6878 let rows = tokio::time::timeout(std::time::Duration::from_secs(2), query)
6879 .await
6880 .expect("DML RETURNING did not finish after its gate was released")
6881 .unwrap()
6882 .expect("request cancellation must preserve DML RETURNING's result");
6883 assert_eq!(rows.len(), 10_000);
6884 tokio::time::timeout(std::time::Duration::from_secs(1), completed.notified())
6885 .await
6886 .expect("completed DML RETURNING did not release its connection");
6887
6888 let reader = pool.reader().unwrap();
6889 let count: i64 = reader
6890 .conn()
6891 .query_row("SELECT COUNT(*) FROM returning_write_probe", [], |row| {
6892 row.get(0)
6893 })
6894 .unwrap();
6895 assert_eq!(count, 10_000, "DML RETURNING must commit every row");
6896 }
6897
6898 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6906 async fn cancelled_call_invalidates_handle_reuse_fails_loud() {
6907 let dir = tempfile::tempdir().unwrap();
6908 let config = PoolConfig {
6909 path: Some(dir.path().join("sql_bridge_cancelled_reuse.db")),
6910 checkout_timeout: std::time::Duration::from_millis(250),
6911 ..PoolConfig::for_test()
6912 };
6913 let pool = Arc::new(ConnectionPool::new(config).unwrap());
6914
6915 let handle_slot = acquire_handle_slot(
6916 pool.sql_bridge_writer_slots(),
6917 pool.config().checkout_timeout,
6918 "sql_bridge.writer_handle",
6919 SlotTimeoutClass::Admission,
6920 )
6921 .await
6922 .unwrap();
6923 let conn = open_standalone_writer(&pool).unwrap();
6924 let (entered, release, completed) = blocking_non_interrupting_progress_gate(&conn);
6925 let writer = Arc::new(tokio::sync::Mutex::new(SqliteWriter {
6926 observe_direct_errors: true,
6927 event_rows: None,
6928 handle: Some(StandaloneHandle {
6929 conn,
6930 _retained_slot: Some(handle_slot),
6931 read_transaction_slot: None,
6932 }),
6933 writer_task: None,
6934 origin: pool.origin(),
6935 db: crate::timeout_sink::db_label(&pool),
6936 pool: Arc::clone(&pool),
6937 held_lease: None,
6938 }));
6939 let writer_clone = Arc::clone(&writer);
6940 let query = tokio::spawn(async move {
6941 khive_storage::SqlWriter::execute_batch(
6942 &mut *writer_clone.lock().await,
6943 vec![progress_gate_statement()],
6944 )
6945 .await
6946 });
6947
6948 entered.notified().await;
6949 query.abort();
6950 let cancelled = matches!(query.await, Err(error) if error.is_cancelled());
6951
6952 let reuse = khive_storage::SqlWriter::execute(
6953 &mut *writer.lock().await,
6954 SqlStatement {
6955 sql: "CREATE TABLE cancelled_reuse_probe (id INTEGER PRIMARY KEY)".into(),
6956 params: vec![],
6957 label: None,
6958 },
6959 )
6960 .await;
6961 let message = match reuse {
6962 Err(StorageError::Pool { message, .. }) => message,
6963 other => panic!(
6964 "reusing a cancelled writer handle must fail loudly with \
6965 'connection already consumed'; got {other:?}"
6966 ),
6967 };
6968 assert!(
6969 message.contains("connection already consumed"),
6970 "expected the cancelled handle's reuse error to name the pinned \
6971 failure; got {message:?}"
6972 );
6973
6974 tokio::task::spawn_blocking(move || release.wait())
6975 .await
6976 .unwrap();
6977 tokio::time::timeout(std::time::Duration::from_secs(1), completed.notified())
6978 .await
6979 .expect("cancelled writer's detached SQLite call did not finish");
6980 assert!(cancelled, "writer batch task did not report cancellation");
6981 }
6982
6983 #[tokio::test]
6992 async fn execute_batch_rejects_transaction_control_before_executing_anything() {
6993 let dir = tempfile::tempdir().unwrap();
6994 let config = PoolConfig {
6995 path: Some(dir.path().join("sql_bridge_tx_control_reject.db")),
6996 checkout_timeout: std::time::Duration::from_millis(250),
6997 write_queue_enabled: Some(false),
6998 ..PoolConfig::for_test()
6999 };
7000 let pool = Arc::new(ConnectionPool::new(config).unwrap());
7001 {
7002 let guard = pool.writer().unwrap();
7003 guard
7004 .conn()
7005 .execute_batch(
7006 "CREATE TABLE tx_reject_test (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
7007 )
7008 .unwrap();
7009 }
7010
7011 let handle_slot = acquire_handle_slot(
7012 pool.sql_bridge_writer_slots(),
7013 pool.config().checkout_timeout,
7014 "sql_bridge.writer_handle",
7015 SlotTimeoutClass::Admission,
7016 )
7017 .await
7018 .unwrap();
7019 let conn = open_standalone_writer(&pool).unwrap();
7020 let mut writer = SqliteWriter {
7021 observe_direct_errors: true,
7022 event_rows: None,
7023 handle: Some(StandaloneHandle {
7024 conn,
7025 _retained_slot: Some(handle_slot),
7026 read_transaction_slot: None,
7027 }),
7028 writer_task: None,
7029 origin: pool.origin(),
7030 db: crate::timeout_sink::db_label(&pool),
7031 pool: Arc::clone(&pool),
7032 held_lease: None,
7033 };
7034
7035 for tail in ["COMMIT", "BEGIN"] {
7036 let multi = khive_storage::SqlWriter::execute_batch(
7037 &mut writer,
7038 vec![SqlStatement {
7039 sql: format!(
7040 "INSERT INTO tx_reject_test (id, val) VALUES (10, 'tail'); {tail}"
7041 ),
7042 params: vec![],
7043 label: None,
7044 }],
7045 )
7046 .await;
7047 let message = multi
7048 .as_ref()
7049 .err()
7050 .map(ToString::to_string)
7051 .unwrap_or_default();
7052 assert!(
7053 message.contains("Multiple statements"),
7054 "a SqlStatement with trailing {tail} must be rejected before execution; got {message}"
7055 );
7056 }
7057
7058 let batch = khive_storage::SqlWriter::execute_batch(
7061 &mut writer,
7062 vec![
7063 SqlStatement {
7064 sql: "INSERT INTO tx_reject_test (id, val) VALUES (1, 'a')".into(),
7065 params: vec![],
7066 label: None,
7067 },
7068 SqlStatement {
7069 sql: "COMMIT".into(),
7070 params: vec![],
7071 label: None,
7072 },
7073 ],
7074 )
7075 .await;
7076 match &batch {
7077 Err(StorageError::InvalidInput {
7078 operation, message, ..
7079 }) => {
7080 assert_eq!(operation.as_ref(), "execute_batch");
7081 assert!(
7082 message.contains("transaction control") && message.contains("COMMIT"),
7083 "the rejection must name the offending statement head; got {message:?}"
7084 );
7085 }
7086 other => {
7087 panic!("a batch containing a bare COMMIT must be rejected up front; got {other:?}")
7088 }
7089 }
7090
7091 for sql in [
7094 "BEGIN IMMEDIATE",
7095 "START TRANSACTION",
7096 "commit",
7097 "End transaction",
7098 "ROLLBACK",
7099 "SAVEPOINT sp1",
7100 "RELEASE sp1",
7101 " -- leading comment\nCOMMIT",
7102 "/* block */ rollback to savepoint sp1",
7103 ] {
7104 let rejected = khive_storage::SqlWriter::execute_batch(
7105 &mut writer,
7106 vec![SqlStatement {
7107 sql: sql.into(),
7108 params: vec![],
7109 label: None,
7110 }],
7111 )
7112 .await;
7113 assert!(
7114 matches!(&rejected, Err(StorageError::InvalidInput { .. })),
7115 "transaction-control head {sql:?} must be rejected; got {rejected:?}"
7116 );
7117 }
7118
7119 let count: i64 = {
7123 let guard = pool.reader().unwrap();
7124 guard
7125 .conn()
7126 .query_row("SELECT COUNT(*) FROM tx_reject_test", [], |r| r.get(0))
7127 .unwrap()
7128 };
7129 assert_eq!(count, 0, "a rejected batch must not have executed anything");
7130
7131 let affected = khive_storage::SqlWriter::execute(
7132 &mut writer,
7133 SqlStatement {
7134 sql: "INSERT INTO tx_reject_test (id, val) VALUES (2, 'b')".into(),
7135 params: vec![],
7136 label: None,
7137 },
7138 )
7139 .await
7140 .expect("the handle must survive a rejected batch untouched");
7141 assert_eq!(affected, 1);
7142 }
7143
7144 #[tokio::test]
7145 async fn standalone_execute_batch_rejects_prefixed_commit_before_any_write() {
7146 let dir = tempfile::tempdir().unwrap();
7147 let config = PoolConfig {
7148 path: Some(dir.path().join("sql_bridge_prefixed_commit_standalone.db")),
7149 write_queue_enabled: Some(false),
7150 ..PoolConfig::for_test()
7151 };
7152 let pool = Arc::new(ConnectionPool::new(config).unwrap());
7153 {
7154 let guard = pool.writer().unwrap();
7155 guard
7156 .conn()
7157 .execute_batch("CREATE TABLE prefixed_commit (id INTEGER PRIMARY KEY)")
7158 .unwrap();
7159 }
7160 let bridge = SqlBridge::new(Arc::clone(&pool), true);
7161 let mut writer = bridge.writer().await.unwrap();
7162
7163 let rejected = writer
7164 .execute_batch(vec![
7165 SqlStatement {
7166 sql: "INSERT INTO prefixed_commit (id) VALUES (1)".into(),
7167 params: vec![],
7168 label: None,
7169 },
7170 SqlStatement {
7171 sql: " ; -- empty statement\n /* leading comment */ \u{feff} ; COMMIT".into(),
7172 params: vec![],
7173 label: None,
7174 },
7175 ])
7176 .await;
7177 assert!(
7178 matches!(
7179 &rejected,
7180 Err(StorageError::InvalidInput {
7181 operation,
7182 message,
7183 ..
7184 }) if operation.as_ref() == "execute_batch"
7185 && message.contains("transaction control")
7186 && message.contains("COMMIT")
7187 ),
7188 "standalone execute_batch must reject a prefixed COMMIT before the INSERT; \
7189 got {rejected:?}"
7190 );
7191
7192 let mut reader = bridge.reader().await.unwrap();
7193 let count = reader
7194 .query_scalar(SqlStatement {
7195 sql: "SELECT COUNT(*) FROM prefixed_commit".into(),
7196 params: vec![],
7197 label: None,
7198 })
7199 .await
7200 .unwrap();
7201 assert!(
7202 matches!(count, Some(SqlValue::Integer(0))),
7203 "prefixed COMMIT rejection must happen before the earlier INSERT; got {count:?}"
7204 );
7205
7206 let affected = writer
7207 .execute(SqlStatement {
7208 sql: "INSERT INTO prefixed_commit (id) VALUES (2)".into(),
7209 params: vec![],
7210 label: None,
7211 })
7212 .await
7213 .expect("prefixed COMMIT rejection must leave the standalone handle reusable");
7214 assert_eq!(affected, 1);
7215 }
7216
7217 #[tokio::test]
7218 async fn execute_batch_rejects_multi_statement_on_pool_backed_path() {
7219 let pool = Arc::new(ConnectionPool::new(PoolConfig::default()).unwrap());
7220 pool.writer()
7221 .unwrap()
7222 .conn()
7223 .execute_batch(
7224 "CREATE TABLE multi_statement_pool_test (id INTEGER PRIMARY KEY, val TEXT)",
7225 )
7226 .unwrap();
7227 let bridge = SqlBridge::new(Arc::clone(&pool), false);
7228 let mut writer = bridge.writer().await.unwrap();
7229
7230 let result = khive_storage::SqlWriter::execute_batch(
7231 &mut *writer,
7232 vec![SqlStatement {
7233 sql: "INSERT INTO multi_statement_pool_test (id, val) VALUES (1, 'x'); COMMIT"
7234 .into(),
7235 params: vec![],
7236 label: None,
7237 }],
7238 )
7239 .await;
7240 let message = result
7241 .as_ref()
7242 .err()
7243 .map(ToString::to_string)
7244 .unwrap_or_default();
7245 assert!(
7246 message.contains("Multiple statements"),
7247 "pool-backed execute_batch must reject a trailing COMMIT; got {message}"
7248 );
7249 let count: i64 = pool
7250 .reader()
7251 .unwrap()
7252 .conn()
7253 .query_row(
7254 "SELECT COUNT(*) FROM multi_statement_pool_test",
7255 [],
7256 |row| row.get(0),
7257 )
7258 .unwrap();
7259 assert_eq!(count, 0);
7260 }
7261
7262 #[tokio::test]
7263 async fn inline_execute_batch_rejects_multi_statement_sql() {
7264 let dir = tempfile::tempdir().unwrap();
7265 let pool = Arc::new(
7266 ConnectionPool::new(PoolConfig {
7267 path: Some(dir.path().join("sql_bridge_multi_statement_inline.db")),
7268 write_queue_enabled: Some(true),
7269 write_routing_strict: true,
7270 ..PoolConfig::for_test()
7271 })
7272 .unwrap(),
7273 );
7274 pool.writer()
7275 .unwrap()
7276 .conn()
7277 .execute_batch(
7278 "CREATE TABLE multi_statement_inline_test (id INTEGER PRIMARY KEY, val TEXT)",
7279 )
7280 .unwrap();
7281 let bridge = SqlBridge::new(Arc::clone(&pool), true);
7282
7283 let result = bridge
7284 .atomic_unit(Box::new(|writer| {
7285 Box::pin(async move {
7286 writer
7287 .execute_batch(vec![SqlStatement {
7288 sql: "INSERT INTO multi_statement_inline_test (id, val) VALUES (1, 'x'); BEGIN"
7289 .into(),
7290 params: vec![],
7291 label: None,
7292 }])
7293 .await
7294 .map(|_| Box::new(()) as Box<dyn Any + Send>)
7295 })
7296 }))
7297 .await;
7298 let message = result
7299 .as_ref()
7300 .err()
7301 .map(ToString::to_string)
7302 .unwrap_or_default();
7303 assert!(
7304 message.contains("Multiple statements"),
7305 "InlineWriter must reject a trailing BEGIN; got {message}"
7306 );
7307 let count: i64 = pool
7308 .reader()
7309 .unwrap()
7310 .conn()
7311 .query_row(
7312 "SELECT COUNT(*) FROM multi_statement_inline_test",
7313 [],
7314 |row| row.get(0),
7315 )
7316 .unwrap();
7317 assert_eq!(count, 0);
7318 }
7319
7320 #[test]
7325 fn transaction_control_head_classification_matrix() {
7326 for (sql, expected) in [
7327 ("BEGIN", Some("BEGIN")),
7328 ("begin immediate", Some("BEGIN")),
7329 ("START TRANSACTION", Some("START")),
7330 ("start transaction", Some("START")),
7331 ("COMMIT", Some("COMMIT")),
7332 ("commit;", Some("COMMIT")),
7333 ("END", Some("END")),
7334 ("end transaction", Some("END")),
7335 ("ROLLBACK", Some("ROLLBACK")),
7336 ("rollback to savepoint sp1", Some("ROLLBACK")),
7337 ("SAVEPOINT sp1", Some("SAVEPOINT")),
7338 ("RELEASE sp1", Some("RELEASE")),
7339 ("release savepoint sp1", Some("RELEASE")),
7340 (" \t COMMIT", Some("COMMIT")),
7341 ("\u{feff}BEGIN", Some("BEGIN")),
7342 (" ; BEGIN", Some("BEGIN")),
7343 (" ; ; -- empty\n /* comment */ COMMIT", Some("COMMIT")),
7344 (" \u{feff} SAVEPOINT sp1", Some("SAVEPOINT")),
7345 ("/* comment */ \u{feff} ; RELEASE sp1", Some("RELEASE")),
7346 ("\u{feff} ; \u{feff} -- empty\n ROLLBACK", Some("ROLLBACK")),
7347 ("-- a comment\nCOMMIT", Some("COMMIT")),
7348 ("/* /* nested? no */ */ COMMIT", None),
7351 ("-- one\n-- two\n /* x */ begin", Some("BEGIN")),
7352 ("INSERT INTO t VALUES (1)", None),
7353 ("UPDATE t SET x = 1", None),
7354 ("DELETE FROM t", None),
7355 ("SELECT * FROM commit_log", None),
7356 ("CREATE TABLE rollback_audit (id INTEGER)", None),
7357 ("/* comment only */", None),
7358 (" ; /* empty statements only */ ; ", None),
7359 (" ; SELECT 1", None),
7360 ("", None),
7361 ] {
7362 assert_eq!(
7363 transaction_control_head(sql),
7364 expected,
7365 "classification mismatch for {sql:?}"
7366 );
7367 }
7368 }
7369
7370 #[test]
7371 fn cached_read_transaction_control_classification_matrix() {
7372 use CachedReadTransactionControl::{BeginDeferred, Finish, Unsupported};
7373
7374 for (sql, expected) in [
7375 ("BEGIN", Some(BeginDeferred)),
7376 ("begin transaction", Some(BeginDeferred)),
7377 (
7378 "/* p */ \u{feff} ; BEGIN /* mode */ DEFERRED",
7379 Some(BeginDeferred),
7380 ),
7381 ("BEGIN DEFERRED TRANSACTION", Some(BeginDeferred)),
7382 ("BEGIN IMMEDIATE", Some(Unsupported("BEGIN"))),
7383 ("BEGIN /* lock */ EXCLUSIVE", Some(Unsupported("BEGIN"))),
7384 ("BEGIN TRANSACTION IMMEDIATE", Some(Unsupported("BEGIN"))),
7385 ("begin transaction exclusive", Some(Unsupported("BEGIN"))),
7386 ("BEGIN TRANSACTION DEFERRED", Some(Unsupported("BEGIN"))),
7387 ("BEGIN IMMEDIATE TRANSACTION", Some(Unsupported("BEGIN"))),
7388 ("BEGIN TRANSACTION named_txn", Some(Unsupported("BEGIN"))),
7389 (
7390 "BEGIN DEFERRED TRANSACTION trailing",
7391 Some(Unsupported("BEGIN")),
7392 ),
7393 ("BEGIN DEFERRED DEFERRED", Some(Unsupported("BEGIN"))),
7394 (
7398 "BEGIN TRANSACTION \"IMMEDIATE\"",
7399 Some(Unsupported("BEGIN")),
7400 ),
7401 ("BEGIN TRANSACTION [IMMEDIATE]", Some(Unsupported("BEGIN"))),
7402 ("BEGIN TRANSACTION `IMMEDIATE`", Some(Unsupported("BEGIN"))),
7403 ("BEGIN TRANSACTION 'IMMEDIATE'", Some(Unsupported("BEGIN"))),
7404 ("BEGIN \"DEFERRED\"", Some(Unsupported("BEGIN"))),
7405 ("BEGIN; COMMIT", Some(Unsupported("BEGIN"))),
7406 ("BEGIN;", Some(BeginDeferred)),
7408 ("BEGIN DEFERRED ; -- done", Some(BeginDeferred)),
7409 ("BEGIN TRANSACTION /* t */ ;;", Some(BeginDeferred)),
7410 ("START TRANSACTION", Some(Unsupported("START"))),
7411 ("COMMIT", Some(Finish("COMMIT"))),
7412 ("END TRANSACTION", Some(Finish("END"))),
7413 ("ROLLBACK", Some(Finish("ROLLBACK"))),
7414 ("ROLLBACK TRANSACTION", Some(Finish("ROLLBACK"))),
7415 ("ROLLBACK TO sp", Some(Unsupported("ROLLBACK"))),
7416 (
7417 "ROLLBACK /* nested */ TRANSACTION /* target */ TO sp",
7418 Some(Unsupported("ROLLBACK")),
7419 ),
7420 ("SAVEPOINT sp", Some(Unsupported("SAVEPOINT"))),
7421 ("SELECT 1", None),
7422 ] {
7423 assert_eq!(
7424 cached_read_transaction_control(sql),
7425 expected,
7426 "cached-reader transaction classification mismatch for {sql:?}"
7427 );
7428 }
7429 }
7430
7431 #[test]
7432 fn sqlite_accepts_utf8_bom_before_transaction_control() {
7433 let conn = rusqlite::Connection::open_in_memory().unwrap();
7434 conn.execute_batch("CREATE TABLE bom_transaction_test (id INTEGER)")
7435 .unwrap();
7436 conn.execute_batch("\u{feff}BEGIN IMMEDIATE").unwrap();
7437 conn.execute_batch("ROLLBACK").unwrap();
7438 }
7439
7440 #[tokio::test]
7447 async fn execute_batch_rejects_transaction_control_on_queue_backed_path() {
7448 let dir = tempfile::tempdir().unwrap();
7449 let config = PoolConfig {
7450 path: Some(dir.path().join("sql_bridge_tx_reject_queue.db")),
7451 checkout_timeout: std::time::Duration::from_millis(250),
7452 write_queue_enabled: Some(true),
7453 write_routing_strict: true,
7454 ..PoolConfig::for_test()
7455 };
7456 let pool = Arc::new(ConnectionPool::new(config).unwrap());
7457 {
7458 let guard = pool.writer().unwrap();
7459 guard
7460 .conn()
7461 .execute_batch(
7462 "CREATE TABLE tx_reject_queue_test (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
7463 )
7464 .unwrap();
7465 }
7466 let bridge = SqlBridge::new(Arc::clone(&pool), true);
7467 let mut writer = bridge.writer().await.unwrap();
7468
7469 let rejected = khive_storage::SqlWriter::execute_batch(
7470 &mut *writer,
7471 vec![
7472 SqlStatement {
7473 sql: "INSERT INTO tx_reject_queue_test (id, val) VALUES (1, 'a')".into(),
7474 params: vec![],
7475 label: None,
7476 },
7477 SqlStatement {
7478 sql: "COMMIT".into(),
7479 params: vec![],
7480 label: None,
7481 },
7482 ],
7483 )
7484 .await;
7485 assert!(
7486 matches!(&rejected, Err(StorageError::InvalidInput { .. })),
7487 "a bare COMMIT in a queue-backed batch must be rejected up front; got {rejected:?}"
7488 );
7489
7490 let prefixed = khive_storage::SqlWriter::execute_batch(
7491 &mut *writer,
7492 vec![
7493 SqlStatement {
7494 sql: "INSERT INTO tx_reject_queue_test (id, val) VALUES (3, 'prefixed')".into(),
7495 params: vec![],
7496 label: None,
7497 },
7498 SqlStatement {
7499 sql: "/* leading */ \u{feff} ; -- empty\n ; COMMIT".into(),
7500 params: vec![],
7501 label: None,
7502 },
7503 ],
7504 )
7505 .await;
7506 assert!(
7507 matches!(
7508 &prefixed,
7509 Err(StorageError::InvalidInput {
7510 operation,
7511 message,
7512 ..
7513 }) if operation.as_ref() == "execute_batch"
7514 && message.contains("transaction control")
7515 && message.contains("COMMIT")
7516 ),
7517 "a prefixed COMMIT must be rejected before touching the writer task; got {prefixed:?}"
7518 );
7519
7520 let affected = khive_storage::SqlWriter::execute_batch(
7521 &mut *writer,
7522 vec![SqlStatement {
7523 sql: "INSERT INTO tx_reject_queue_test (id, val) VALUES (2, 'b')".into(),
7524 params: vec![],
7525 label: None,
7526 }],
7527 )
7528 .await
7529 .expect("the writer task must survive the rejected batch");
7530 assert_eq!(affected, 1);
7531
7532 let count: i64 = {
7533 let guard = pool.reader().unwrap();
7534 guard
7535 .conn()
7536 .query_row("SELECT COUNT(*) FROM tx_reject_queue_test", [], |r| {
7537 r.get(0)
7538 })
7539 .unwrap()
7540 };
7541 assert_eq!(
7542 count, 1,
7543 "exactly the post-rejection batch's row may have landed"
7544 );
7545 }
7546
7547 #[tokio::test]
7560 async fn failed_rollback_poisons_handle_reuse_fails_loud() {
7561 use rusqlite::hooks::{AuthAction, AuthContext, Authorization, TransactionOperation};
7562
7563 fn deny_rollback(ctx: AuthContext<'_>) -> Authorization {
7564 match ctx.action {
7565 AuthAction::Transaction {
7566 operation: TransactionOperation::Rollback,
7567 } => Authorization::Deny,
7568 _ => Authorization::Allow,
7569 }
7570 }
7571
7572 let dir = tempfile::tempdir().unwrap();
7573 let config = PoolConfig {
7574 path: Some(dir.path().join("sql_bridge_rollback_poison.db")),
7575 checkout_timeout: std::time::Duration::from_millis(250),
7576 ..PoolConfig::for_test()
7577 };
7578 let pool = Arc::new(ConnectionPool::new(config).unwrap());
7579 {
7580 let guard = pool.writer().unwrap();
7581 guard
7582 .conn()
7583 .execute_batch(
7584 "CREATE TABLE rollback_poison_test (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
7585 )
7586 .unwrap();
7587 }
7588
7589 let handle_slot = acquire_handle_slot(
7590 pool.sql_bridge_writer_slots(),
7591 pool.config().checkout_timeout,
7592 "sql_bridge.writer_handle",
7593 SlotTimeoutClass::Admission,
7594 )
7595 .await
7596 .unwrap();
7597 let conn = open_standalone_writer(&pool).unwrap();
7598 conn.authorizer(Some(deny_rollback)).unwrap();
7599 let mut writer = SqliteWriter {
7600 observe_direct_errors: true,
7601 event_rows: None,
7602 handle: Some(StandaloneHandle {
7603 conn,
7604 _retained_slot: Some(handle_slot),
7605 read_transaction_slot: None,
7606 }),
7607 writer_task: None,
7608 origin: pool.origin(),
7609 db: crate::timeout_sink::db_label(&pool),
7610 pool: Arc::clone(&pool),
7611 held_lease: None,
7612 };
7613
7614 let batch = khive_storage::SqlWriter::execute_batch(
7615 &mut writer,
7616 vec![
7617 SqlStatement {
7618 sql: "INSERT INTO rollback_poison_test (id, val) VALUES (1, 'a')".into(),
7619 params: vec![],
7620 label: None,
7621 },
7622 SqlStatement {
7623 sql: "SELECT FROM WHERE".into(),
7624 params: vec![],
7625 label: None,
7626 },
7627 ],
7628 )
7629 .await;
7630 let batch_error = batch.expect_err("invalid second statement must fail the batch");
7631 let poison = match &batch_error {
7632 StorageError::Driver { source, .. } => source
7633 .downcast_ref::<PoisonedBatchError>()
7634 .expect("failed rollback must retain its typed poison wrapper"),
7635 other => panic!("failed rollback must return a driver error; got {other:?}"),
7636 };
7637 assert!(
7638 matches!(&poison.poison_reason, BatchPoisonReason::RollbackFailed(_)),
7639 "the poison cause must be compiler-checked as RollbackFailed; got {poison:?}"
7640 );
7641 let batch_message = batch_error.to_string();
7642 assert!(
7643 batch_message.contains("ROLLBACK after statement failure failed"),
7644 "the caller must see the poison context naming the failed \
7645 rollback; got {batch_message:?}"
7646 );
7647 assert!(
7648 batch_message.contains("original error"),
7649 "the original statement error must stay visible alongside the \
7650 poison context; got {batch_message:?}"
7651 );
7652
7653 let reuse = khive_storage::SqlWriter::execute(
7654 &mut writer,
7655 SqlStatement {
7656 sql: "CREATE TABLE rollback_poison_probe (id INTEGER PRIMARY KEY)".into(),
7657 params: vec![],
7658 label: None,
7659 },
7660 )
7661 .await;
7662 let message = match reuse {
7663 Err(StorageError::Pool { message, .. }) => message,
7664 other => panic!(
7665 "reusing a poisoned writer handle must fail loudly with \
7666 'connection already consumed'; got {other:?}"
7667 ),
7668 };
7669 assert!(
7670 message.contains("connection already consumed"),
7671 "expected the poisoned handle's reuse error to name the pinned \
7672 failure; got {message:?}"
7673 );
7674 }
7675
7676 #[tokio::test]
7682 async fn non_transient_begin_failure_poisons_handle() {
7683 let dir = tempfile::tempdir().unwrap();
7684 let config = PoolConfig {
7685 path: Some(dir.path().join("sql_bridge_begin_poison.db")),
7686 checkout_timeout: std::time::Duration::from_millis(250),
7687 ..PoolConfig::for_test()
7688 };
7689 let pool = Arc::new(ConnectionPool::new(config).unwrap());
7690
7691 let handle_slot = acquire_handle_slot(
7692 pool.sql_bridge_writer_slots(),
7693 pool.config().checkout_timeout,
7694 "sql_bridge.writer_handle",
7695 SlotTimeoutClass::Admission,
7696 )
7697 .await
7698 .unwrap();
7699 let conn = open_standalone_writer(&pool).unwrap();
7700 conn.execute_batch("BEGIN IMMEDIATE").unwrap();
7703 let mut writer = SqliteWriter {
7704 observe_direct_errors: true,
7705 event_rows: None,
7706 handle: Some(StandaloneHandle {
7707 conn,
7708 _retained_slot: Some(handle_slot),
7709 read_transaction_slot: None,
7710 }),
7711 writer_task: None,
7712 origin: pool.origin(),
7713 db: crate::timeout_sink::db_label(&pool),
7714 pool: Arc::clone(&pool),
7715 held_lease: None,
7716 };
7717
7718 let batch = khive_storage::SqlWriter::execute_batch(
7719 &mut writer,
7720 vec![SqlStatement {
7721 sql: "SELECT 1".into(),
7722 params: vec![],
7723 label: None,
7724 }],
7725 )
7726 .await;
7727 let batch_error = batch.expect_err("BEGIN inside an open transaction must fail");
7728 let poison = match &batch_error {
7729 StorageError::Driver { source, .. } => source
7730 .downcast_ref::<PoisonedBatchError>()
7731 .expect("failed BEGIN must retain its typed poison wrapper"),
7732 other => panic!("failed BEGIN must return a driver error; got {other:?}"),
7733 };
7734 assert!(
7735 matches!(&poison.poison_reason, BatchPoisonReason::BeginFailed),
7736 "the poison cause must be compiler-checked as BeginFailed; got {poison:?}"
7737 );
7738 let batch_message = batch_error.to_string();
7739 assert!(
7740 batch_message.contains("BEGIN IMMEDIATE failed non-transiently"),
7741 "a non-transient BEGIN failure must surface the poison context; \
7742 got {batch_message:?}"
7743 );
7744 assert!(
7745 batch_message.contains("cannot start a transaction within a transaction"),
7746 "the original BEGIN error must stay visible; got {batch_message:?}"
7747 );
7748
7749 let reuse = khive_storage::SqlWriter::execute(
7750 &mut writer,
7751 SqlStatement {
7752 sql: "CREATE TABLE begin_poison_probe (id INTEGER PRIMARY KEY)".into(),
7753 params: vec![],
7754 label: None,
7755 },
7756 )
7757 .await;
7758 assert!(
7759 matches!(
7760 &reuse,
7761 Err(StorageError::Pool { message, .. })
7762 if message.contains("connection already consumed")
7763 ),
7764 "a handle poisoned by a non-transient BEGIN failure must be \
7765 dropped, not restored; got {reuse:?}"
7766 );
7767 }
7768
7769 #[tokio::test]
7774 async fn busy_begin_failure_restores_handle_reusable() {
7775 let dir = tempfile::tempdir().unwrap();
7776 let config = PoolConfig {
7777 path: Some(dir.path().join("sql_bridge_begin_busy.db")),
7778 checkout_timeout: std::time::Duration::from_millis(250),
7779 busy_timeout: std::time::Duration::from_millis(100),
7780 ..PoolConfig::for_test()
7781 };
7782 let pool = Arc::new(ConnectionPool::new(config).unwrap());
7783 {
7784 let guard = pool.writer().unwrap();
7785 guard
7786 .conn()
7787 .execute_batch("CREATE TABLE begin_busy_test (id INTEGER PRIMARY KEY)")
7788 .unwrap();
7789 }
7790
7791 let lock_conn = pool.open_standalone_writer().unwrap();
7795 lock_conn.execute_batch("BEGIN IMMEDIATE").unwrap();
7796
7797 let handle_slot = acquire_handle_slot(
7798 pool.sql_bridge_writer_slots(),
7799 pool.config().checkout_timeout,
7800 "sql_bridge.writer_handle",
7801 SlotTimeoutClass::Admission,
7802 )
7803 .await
7804 .unwrap();
7805 let conn = open_standalone_writer(&pool).unwrap();
7806 let mut writer = SqliteWriter {
7807 observe_direct_errors: true,
7808 event_rows: None,
7809 handle: Some(StandaloneHandle {
7810 conn,
7811 _retained_slot: Some(handle_slot),
7812 read_transaction_slot: None,
7813 }),
7814 writer_task: None,
7815 origin: pool.origin(),
7816 db: crate::timeout_sink::db_label(&pool),
7817 pool: Arc::clone(&pool),
7818 held_lease: None,
7819 };
7820
7821 let batch = khive_storage::SqlWriter::execute_batch(
7822 &mut writer,
7823 vec![SqlStatement {
7824 sql: "INSERT INTO begin_busy_test (id) VALUES (1)".into(),
7825 params: vec![],
7826 label: None,
7827 }],
7828 )
7829 .await;
7830 let batch_error = batch.expect_err("BEGIN IMMEDIATE under a held write lock must fail");
7831 assert!(
7832 batch_error.to_string().contains("database is locked"),
7833 "the busy BEGIN failure must surface SQLite's busy error; got {batch_error:?}"
7834 );
7835
7836 lock_conn.execute_batch("ROLLBACK").unwrap();
7837 drop(lock_conn);
7838
7839 let affected = khive_storage::SqlWriter::execute(
7840 &mut writer,
7841 SqlStatement {
7842 sql: "INSERT INTO begin_busy_test (id) VALUES (2)".into(),
7843 params: vec![],
7844 label: None,
7845 },
7846 )
7847 .await
7848 .expect("a busy BEGIN failure must restore the handle as reusable");
7849 assert_eq!(affected, 1);
7850 }
7851
7852 #[tokio::test]
7857 async fn manual_atomic_unit_shares_writer_permit_budget_with_writer_handle() {
7858 let dir = tempfile::tempdir().unwrap();
7859 let config = PoolConfig {
7860 path: Some(dir.path().join("sql_bridge_atomic_unit_budget.db")),
7861 checkout_timeout: std::time::Duration::from_millis(50),
7862 write_queue_enabled: Some(false),
7863 ..PoolConfig::for_test()
7864 };
7865 let pool = Arc::new(ConnectionPool::new(config).unwrap());
7866 let bridge = SqlBridge::new(Arc::clone(&pool), true);
7867 {
7868 let guard = pool.writer().unwrap();
7869 guard
7870 .conn()
7871 .execute_batch(
7872 "CREATE TABLE IF NOT EXISTS atomic_unit_budget_test \
7873 (id INTEGER PRIMARY KEY, val INTEGER NOT NULL)",
7874 )
7875 .unwrap();
7876 }
7877
7878 fn insert_op(id: i64) -> AtomicUnitOp {
7879 Box::new(move |writer| {
7880 Box::pin(async move {
7881 writer
7882 .execute(SqlStatement {
7883 sql: "INSERT INTO atomic_unit_budget_test (id, val) VALUES (?1, ?2)"
7884 .into(),
7885 params: vec![SqlValue::Integer(id), SqlValue::Integer(id)],
7886 label: None,
7887 })
7888 .await
7889 .map_err(|e| {
7890 khive_storage::StorageError::driver(
7891 StorageCapability::Sql,
7892 "atomic_unit_budget_test_insert",
7893 e,
7894 )
7895 })?;
7896 Ok(Box::new(()) as Box<dyn std::any::Any + Send>)
7897 })
7898 })
7899 }
7900
7901 let writer_handle = bridge.writer().await.unwrap();
7902 let blocked = bridge.atomic_unit(insert_op(1)).await;
7903 assert!(
7904 matches!(
7905 &blocked,
7906 Err(StorageError::AdmissionTimeout { operation, .. })
7907 if operation.as_ref() == "sql_bridge.atomic_unit_handle"
7908 ),
7909 "atomic_unit must time out on the shared writer permit while a \
7910 writer handle is live; got {blocked:?}"
7911 );
7912
7913 drop(writer_handle);
7914 let unblocked = bridge.atomic_unit(insert_op(2)).await;
7915 assert!(
7916 unblocked.is_ok(),
7917 "atomic_unit must succeed once the writer handle releases the \
7918 shared writer permit; got {unblocked:?}"
7919 );
7920
7921 let mut reader = bridge.reader().await.unwrap();
7922 let count = reader
7923 .query_scalar(SqlStatement {
7924 sql: "SELECT COUNT(*) FROM atomic_unit_budget_test".into(),
7925 params: vec![],
7926 label: None,
7927 })
7928 .await
7929 .unwrap();
7930 assert!(
7931 matches!(count, Some(SqlValue::Integer(1))),
7932 "only the post-drop atomic_unit call may have committed; got {count:?}"
7933 );
7934 }
7935
7936 #[tokio::test]
7942 async fn execute_batch_routes_through_writer_task_when_flag_enabled() {
7943 let dir = tempfile::tempdir().unwrap();
7944 let path = dir.path().join("write_queue_execute_batch.db");
7945 let config = PoolConfig {
7946 path: Some(path.clone()),
7947 write_queue_enabled: Some(true),
7948 ..PoolConfig::for_test()
7949 };
7950 let pool = Arc::new(ConnectionPool::new(config).unwrap());
7951 {
7952 let guard = pool.writer().unwrap();
7953 guard
7954 .conn()
7955 .execute_batch(
7956 "CREATE TABLE IF NOT EXISTS write_queue_batch_test \
7957 (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
7958 )
7959 .unwrap();
7960 }
7961
7962 let bridge = SqlBridge::new(Arc::clone(&pool), true);
7963
7964 let mut writer = bridge.writer().await.unwrap();
7965 let affected = writer
7966 .execute_batch(vec![
7967 SqlStatement {
7968 sql: "INSERT INTO write_queue_batch_test (id, val) VALUES (?1, ?2)".into(),
7969 params: vec![SqlValue::Integer(1), SqlValue::Text("a".into())],
7970 label: None,
7971 },
7972 SqlStatement {
7973 sql: "INSERT INTO write_queue_batch_test (id, val) VALUES (?1, ?2)".into(),
7974 params: vec![SqlValue::Integer(2), SqlValue::Text("b".into())],
7975 label: None,
7976 },
7977 ])
7978 .await
7979 .unwrap();
7980 assert_eq!(affected, 2);
7981
7982 let mut reader = bridge.reader().await.unwrap();
7983 let count = reader
7984 .query_scalar(SqlStatement {
7985 sql: "SELECT COUNT(*) FROM write_queue_batch_test".into(),
7986 params: vec![],
7987 label: None,
7988 })
7989 .await
7990 .unwrap();
7991 assert!(
7992 matches!(count, Some(SqlValue::Integer(2))),
7993 "expected 2 rows, got {count:?}"
7994 );
7995 assert_eq!(
7996 pool.writer_task_spawn_count(),
7997 1,
7998 "the flag-ON path must actually spawn and use the writer task"
7999 );
8000 }
8001
8002 #[tokio::test]
8009 async fn execute_batch_rolls_back_atomically_on_mid_sequence_failure() {
8010 let dir = tempfile::tempdir().unwrap();
8011 let path = dir.path().join("write_queue_execute_batch_rollback.db");
8012 let config = PoolConfig {
8013 path: Some(path.clone()),
8014 write_queue_enabled: Some(true),
8015 ..PoolConfig::for_test()
8016 };
8017 let pool = Arc::new(ConnectionPool::new(config).unwrap());
8018 {
8019 let guard = pool.writer().unwrap();
8020 guard
8021 .conn()
8022 .execute_batch(
8023 "CREATE TABLE IF NOT EXISTS write_queue_rollback_test \
8024 (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
8025 )
8026 .unwrap();
8027 }
8028
8029 let bridge = SqlBridge::new(Arc::clone(&pool), true);
8030
8031 let mut writer = bridge.writer().await.unwrap();
8032 let result = writer
8033 .execute_batch(vec![
8034 SqlStatement {
8036 sql: "INSERT INTO write_queue_rollback_test (id, val) VALUES (?1, ?2)".into(),
8037 params: vec![SqlValue::Integer(1), SqlValue::Text("first".into())],
8038 label: None,
8039 },
8040 SqlStatement {
8042 sql: "INSERT INTO write_queue_rollback_test (id, val) VALUES (?1, ?2)".into(),
8043 params: vec![SqlValue::Integer(1), SqlValue::Text("duplicate".into())],
8044 label: None,
8045 },
8046 SqlStatement {
8048 sql: "INSERT INTO write_queue_rollback_test (id, val) VALUES (?1, ?2)".into(),
8049 params: vec![SqlValue::Integer(2), SqlValue::Text("third".into())],
8050 label: None,
8051 },
8052 ])
8053 .await;
8054 assert!(
8055 result.is_err(),
8056 "a batch with a mid-sequence PK conflict must return an error"
8057 );
8058
8059 let mut reader = bridge.reader().await.unwrap();
8060 let count = reader
8061 .query_scalar(SqlStatement {
8062 sql: "SELECT COUNT(*) FROM write_queue_rollback_test".into(),
8063 params: vec![],
8064 label: None,
8065 })
8066 .await
8067 .unwrap();
8068 assert!(
8069 matches!(count, Some(SqlValue::Integer(0))),
8070 "the whole request must roll back — including statement 1's \
8071 otherwise-successful INSERT — not just the failing statement; \
8072 got {count:?}"
8073 );
8074 }
8075
8076 #[tokio::test]
8077 async fn in_memory_atomic_unit_terminal_fault_retires_writer() {
8078 use khive_storage::WriterTaskRequestState;
8079 use rusqlite::hooks::{AuthAction, AuthContext, Authorization, TransactionOperation};
8080
8081 for (mode, deny_rollback, expected) in [
8084 ("error", true, WriterTaskRequestState::SideEffectsUnknown),
8085 ("commit", true, WriterTaskRequestState::SideEffectsUnknown),
8086 ("panic", true, WriterTaskRequestState::SideEffectsUnknown),
8087 (
8088 "panic",
8089 false,
8090 WriterTaskRequestState::TransactionRolledBack,
8091 ),
8092 ] {
8093 let pool = Arc::new(
8094 ConnectionPool::new(PoolConfig {
8095 path: None,
8096 ..PoolConfig::default()
8097 })
8098 .unwrap(),
8099 );
8100 {
8101 let guard = pool.writer().unwrap();
8102 guard
8103 .execute_batch("CREATE TABLE atomic_terminal_probe (id INTEGER PRIMARY KEY)")
8104 .unwrap();
8105 guard
8106 .authorizer(Some(move |ctx: AuthContext<'_>| match ctx.action {
8107 AuthAction::Transaction {
8108 operation: TransactionOperation::Rollback,
8109 } if deny_rollback => Authorization::Deny,
8110 AuthAction::Transaction {
8111 operation: TransactionOperation::Unknown,
8112 } if mode == "commit" => Authorization::Deny,
8113 _ => Authorization::Allow,
8114 }))
8115 .unwrap();
8116 }
8117 let bridge = SqlBridge::new(Arc::clone(&pool), false);
8118 let result = bridge
8119 .atomic_unit(Box::new(move |writer| {
8120 Box::pin(async move {
8121 writer
8122 .execute(SqlStatement {
8123 sql: "INSERT INTO atomic_terminal_probe VALUES (1)".into(),
8124 params: vec![],
8125 label: None,
8126 })
8127 .await?;
8128 match mode {
8129 "error" => Err(StorageError::Internal("terminal probe".into())),
8130 "panic" => panic!("terminal probe"),
8131 _ => Ok(Box::new(()) as Box<dyn Any + Send>),
8132 }
8133 })
8134 }))
8135 .await;
8136 assert!(
8137 matches!(result, Err(StorageError::WriterTaskTerminated { request_state })
8138 if request_state == expected),
8139 "{mode}: {result:?}"
8140 );
8141 assert!(
8142 pool.try_checkpoint_nowait().is_err(),
8143 "{mode}: writer was not retired"
8144 );
8145 let mut writer = bridge.writer().await.unwrap();
8146 assert!(
8147 writer
8148 .execute(SqlStatement {
8149 sql: "INSERT INTO atomic_terminal_probe VALUES (2)".into(),
8150 params: vec![],
8151 label: None,
8152 })
8153 .await
8154 .is_err(),
8155 "{mode}: ordinary write reused a terminal connection"
8156 );
8157 }
8158 }
8159
8160 #[tokio::test]
8161 async fn in_memory_atomic_unit_holds_writer_guard_through_rollback() {
8162 use std::sync::atomic::{AtomicBool, Ordering};
8163
8164 let pool = Arc::new(
8165 ConnectionPool::new(PoolConfig {
8166 path: None,
8167 ..PoolConfig::default()
8168 })
8169 .unwrap(),
8170 );
8171 pool.writer()
8172 .unwrap()
8173 .execute_batch("CREATE TABLE atomic_guard_probe (id INTEGER PRIMARY KEY)")
8174 .unwrap();
8175 let bridge = SqlBridge::new(Arc::clone(&pool), false);
8176 let excluded = Arc::new(AtomicBool::new(false));
8177 let observed = Arc::clone(&excluded);
8178 let probe_pool = Arc::clone(&pool);
8179 let result = bridge
8180 .atomic_unit(Box::new(move |writer| {
8181 Box::pin(async move {
8182 writer
8183 .execute(SqlStatement {
8184 sql: "INSERT INTO atomic_guard_probe VALUES (1)".into(),
8185 params: vec![],
8186 label: None,
8187 })
8188 .await?;
8189 observed.store(
8190 probe_pool.try_checkpoint_nowait().is_err(),
8191 Ordering::SeqCst,
8192 );
8193 Err(StorageError::Internal("rollback guard probe".into()))
8194 })
8195 }))
8196 .await;
8197 assert!(
8198 matches!(result, Err(StorageError::WriterTaskRequestFailed {
8199 request_state: khive_storage::WriterTaskRequestState::TransactionRolledBack,
8200 ref source,
8201 }) if matches!(source.as_ref(), StorageError::Internal(message)
8202 if message == "rollback guard probe")),
8203 "{result:?}"
8204 );
8205 assert!(
8206 excluded.load(Ordering::SeqCst),
8207 "atomic unit released its writer guard"
8208 );
8209 let mut writer = bridge.writer().await.unwrap();
8210 writer
8211 .execute(SqlStatement {
8212 sql: "INSERT INTO atomic_guard_probe VALUES (2)".into(),
8213 params: vec![],
8214 label: None,
8215 })
8216 .await
8217 .unwrap();
8218 let rows = writer
8219 .query_all(SqlStatement {
8220 sql: "SELECT id FROM atomic_guard_probe ORDER BY id".into(),
8221 params: vec![],
8222 label: None,
8223 })
8224 .await
8225 .unwrap();
8226 assert_eq!(rows.len(), 1, "{rows:?}");
8227 assert!(matches!(rows[0].columns[0].value, SqlValue::Integer(2)));
8228 }
8229
8230 #[tokio::test]
8231 async fn in_memory_atomic_unit_pending_op_rolls_back_and_releases_guard() {
8232 let pool = Arc::new(
8233 ConnectionPool::new(PoolConfig {
8234 path: None,
8235 ..PoolConfig::default()
8236 })
8237 .unwrap(),
8238 );
8239 pool.writer()
8240 .unwrap()
8241 .execute_batch("CREATE TABLE atomic_pending_probe (id INTEGER PRIMARY KEY)")
8242 .unwrap();
8243 let bridge = SqlBridge::new(Arc::clone(&pool), false);
8244 let result = tokio::time::timeout(
8245 std::time::Duration::from_secs(10),
8246 bridge.atomic_unit(Box::new(|writer| {
8247 Box::pin(async move {
8248 writer
8249 .execute(SqlStatement {
8250 sql: "INSERT INTO atomic_pending_probe VALUES (1)".into(),
8251 params: vec![],
8252 label: None,
8253 })
8254 .await?;
8255 std::future::pending::<
8256 khive_storage::types::StorageResult<Box<dyn Any + Send>>,
8257 >()
8258 .await
8259 })
8260 })),
8261 )
8262 .await
8263 .expect("suspending atomic unit must return promptly");
8264 assert!(
8265 matches!(result, Err(StorageError::WriterTaskRequestFailed {
8266 request_state: khive_storage::WriterTaskRequestState::TransactionRolledBack,
8267 ref source,
8268 }) if matches!(source.as_ref(), StorageError::Internal(message)
8269 if message.contains("future suspended"))),
8270 "{result:?}"
8271 );
8272 let mut writer = bridge.writer().await.unwrap();
8273 writer
8274 .execute(SqlStatement {
8275 sql: "INSERT INTO atomic_pending_probe VALUES (2)".into(),
8276 params: vec![],
8277 label: None,
8278 })
8279 .await
8280 .unwrap();
8281 let sum = writer
8282 .query_scalar(SqlStatement {
8283 sql: "SELECT SUM(id) FROM atomic_pending_probe".into(),
8284 params: vec![],
8285 label: None,
8286 })
8287 .await
8288 .unwrap();
8289 assert!(matches!(sum, Some(SqlValue::Integer(2))), "{sum:?}");
8290 assert!(pool.writer().unwrap().is_autocommit());
8291 }
8292
8293 #[tokio::test]
8310 async fn atomic_unit_pending_future_errors_without_killing_writer_task() {
8311 let dir = tempfile::tempdir().unwrap();
8312 let path = dir.path().join("atomic_unit_pending_future.db");
8313 let config = PoolConfig {
8314 path: Some(path.clone()),
8315 write_queue_enabled: Some(true),
8316 ..PoolConfig::for_test()
8317 };
8318 let pool = Arc::new(ConnectionPool::new(config).unwrap());
8319 {
8320 let guard = pool.writer().unwrap();
8321 guard
8322 .conn()
8323 .execute_batch(
8324 "CREATE TABLE IF NOT EXISTS atomic_unit_pending_test \
8325 (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
8326 )
8327 .unwrap();
8328 }
8329 assert!(
8330 pool.writer_task_handle().unwrap().is_some(),
8331 "writer task must be spawned with the flag on for a file-backed pool"
8332 );
8333
8334 let bridge = SqlBridge::new(Arc::clone(&pool), true);
8335
8336 let pending_op: AtomicUnitOp = Box::new(|_writer| {
8339 Box::pin(std::future::pending::<
8340 khive_storage::types::StorageResult<Box<dyn std::any::Any + Send>>,
8341 >())
8342 });
8343
8344 let pending_result = bridge.atomic_unit(pending_op).await;
8345 assert!(
8346 pending_result.is_err(),
8347 "a Pending-on-first-poll atomic_unit closure must return Err, \
8348 not panic; got {pending_result:?}"
8349 );
8350
8351 let ok_op: AtomicUnitOp = Box::new(|writer| {
8356 Box::pin(async move {
8357 writer
8358 .execute(SqlStatement {
8359 sql: "INSERT INTO atomic_unit_pending_test (id, val) VALUES (?1, ?2)"
8360 .into(),
8361 params: vec![SqlValue::Integer(1), SqlValue::Text("survived".into())],
8362 label: None,
8363 })
8364 .await
8365 .map_err(|e| {
8366 khive_storage::StorageError::driver(
8367 StorageCapability::Sql,
8368 "atomic_unit_pending_future_test_insert",
8369 e,
8370 )
8371 })?;
8372 Ok(Box::new(()) as Box<dyn std::any::Any + Send>)
8373 })
8374 });
8375 let ok_result = bridge.atomic_unit(ok_op).await;
8376 assert!(
8377 ok_result.is_ok(),
8378 "writer task must survive a Pending misuse and keep serving \
8379 subsequent well-behaved atomic_unit requests; got {ok_result:?}"
8380 );
8381
8382 let mut reader = bridge.reader().await.unwrap();
8383 let count = reader
8384 .query_scalar(SqlStatement {
8385 sql: "SELECT COUNT(*) FROM atomic_unit_pending_test".into(),
8386 params: vec![],
8387 label: None,
8388 })
8389 .await
8390 .unwrap();
8391 assert!(
8392 matches!(count, Some(SqlValue::Integer(1))),
8393 "the well-behaved atomic_unit call after the Pending misuse must \
8394 have actually committed its write; got {count:?}"
8395 );
8396 }
8397
8398 #[tokio::test]
8408 async fn writer_strict_routing_fails_closed_without_writer_task() {
8409 let dir = tempfile::tempdir().unwrap();
8410 let path = dir.path().join("strict_writer.db");
8411 let config = PoolConfig {
8412 path: Some(path),
8413 write_queue_enabled: Some(false),
8414 write_routing_strict: true,
8415 ..PoolConfig::for_test()
8416 };
8417 let pool = Arc::new(ConnectionPool::new(config).unwrap());
8418 let bridge = SqlBridge::new(Arc::clone(&pool), true);
8419
8420 let result = bridge.writer().await;
8421 let err = match result {
8422 Ok(_) => panic!(
8423 "KHIVE_WRITE_ROUTING=strict with no writer task must fail closed, not \
8424 silently degrade to a standalone connection"
8425 ),
8426 Err(e) => e,
8427 };
8428 assert!(
8429 err.to_string().contains("strict"),
8430 "error must name strict routing, got: {err}"
8431 );
8432 }
8433
8434 #[tokio::test]
8439 async fn atomic_unit_strict_routing_fails_closed_without_writer_task() {
8440 let dir = tempfile::tempdir().unwrap();
8441 let path = dir.path().join("strict_atomic_unit.db");
8442 let config = PoolConfig {
8443 path: Some(path),
8444 write_queue_enabled: Some(false),
8445 write_routing_strict: true,
8446 ..PoolConfig::for_test()
8447 };
8448 let pool = Arc::new(ConnectionPool::new(config).unwrap());
8449 let bridge = SqlBridge::new(Arc::clone(&pool), true);
8450
8451 let op: AtomicUnitOp = Box::new(|_writer| {
8452 Box::pin(async move { Ok(Box::new(()) as Box<dyn std::any::Any + Send>) })
8453 });
8454 let result = bridge.atomic_unit(op).await;
8455 assert!(
8456 result.is_err(),
8457 "KHIVE_WRITE_ROUTING=strict but the queue is off (no writer task handle) must \
8458 fail closed instead of falling back to a manual BEGIN IMMEDIATE; got {result:?}"
8459 );
8460 let msg = result.unwrap_err().to_string();
8461 assert!(
8462 msg.contains("strict"),
8463 "error must name strict routing, got: {msg}"
8464 );
8465 }
8466
8467 #[tokio::test]
8475 async fn writer_handle_supports_read_after_write_under_strict_queue() {
8476 let dir = tempfile::tempdir().unwrap();
8477 let path = dir.path().join("writer_read_after_write.db");
8478 let config = PoolConfig {
8479 path: Some(path),
8480 write_queue_enabled: Some(true),
8481 write_routing_strict: true,
8482 ..PoolConfig::for_test()
8483 };
8484 let pool = Arc::new(ConnectionPool::new(config).unwrap());
8485 {
8486 let guard = pool.writer().unwrap();
8487 guard
8488 .conn()
8489 .execute_batch(
8490 "CREATE TABLE IF NOT EXISTS writer_cursor_test \
8491 (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
8492 )
8493 .unwrap();
8494 }
8495
8496 let bridge = SqlBridge::new(Arc::clone(&pool), true);
8497
8498 let mut w = bridge.writer().await.unwrap();
8499 w.execute(SqlStatement {
8500 sql: "INSERT INTO writer_cursor_test (id, val) VALUES (?1, ?2)".into(),
8501 params: vec![SqlValue::Integer(1), SqlValue::Text("via-writer".into())],
8502 label: None,
8503 })
8504 .await
8505 .unwrap();
8506
8507 let row = w
8508 .query_row(SqlStatement {
8509 sql: "SELECT val FROM writer_cursor_test WHERE id = ?1".into(),
8510 params: vec![SqlValue::Integer(1)],
8511 label: None,
8512 })
8513 .await
8514 .unwrap()
8515 .expect("row inserted through the same writer handle must be visible to it");
8516 assert!(
8517 matches!(&row.columns[0].value, SqlValue::Text(v) if v == "via-writer"),
8518 "query_row through a queue-backed writer handle must see its own \
8519 committed write; got {:?}",
8520 row.columns[0].value
8521 );
8522 }
8523
8524 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
8530 async fn queue_backed_read_uses_pool_budget_and_remains_reusable_after_saturation() {
8531 let dir = tempfile::tempdir().unwrap();
8532 let path = dir.path().join("queue_backed_reader_budget.db");
8533 let config = PoolConfig {
8534 path: Some(path),
8535 write_queue_enabled: Some(true),
8536 write_routing_strict: true,
8537 max_readers: 1,
8538 checkout_timeout: std::time::Duration::from_millis(250),
8539 ..PoolConfig::default()
8540 };
8541 let pool = Arc::new(ConnectionPool::new(config).unwrap());
8542 {
8543 let guard = pool.writer().unwrap();
8544 guard
8545 .conn()
8546 .execute_batch(
8547 "CREATE TABLE IF NOT EXISTS reopen_test \
8548 (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
8549 )
8550 .unwrap();
8551 }
8552 let bridge = SqlBridge::new(Arc::clone(&pool), true);
8553
8554 let mut w = bridge.writer().await.unwrap();
8555 w.execute(SqlStatement {
8556 sql: "INSERT INTO reopen_test (id, val) VALUES (1, 'seed')".into(),
8557 params: vec![],
8558 label: None,
8559 })
8560 .await
8561 .unwrap();
8562
8563 let held = pool
8566 .sql_bridge_reader_slots()
8567 .acquire_owned()
8568 .await
8569 .unwrap();
8570 let starved = w
8571 .query_row(SqlStatement {
8572 sql: "SELECT val FROM reopen_test WHERE id = 1".into(),
8573 params: vec![],
8574 label: None,
8575 })
8576 .await;
8577 assert!(
8578 matches!(
8579 &starved,
8580 Err(StorageError::AdmissionTimeout { operation, .. })
8581 if operation.as_ref() == "writer.query_row"
8582 ),
8583 "queue-backed read with reader permits saturated must time out \
8584 at the shared pooled-reader admission stage; \
8585 got {starved:?}"
8586 );
8587 let saturated = pool.reader_acquisition_snapshot();
8588 assert_eq!(saturated.checkout_timeouts, 1);
8589 assert_eq!(saturated.standalone_opens, 0);
8590 drop(held);
8591
8592 let writer_task = pool
8598 .writer_task_handle()
8599 .expect("queue-enabled file pool must offer a writer task")
8600 .expect("writer task present under write_queue_enabled");
8601 let mut post_cancel = SqliteWriter {
8602 observe_direct_errors: true,
8603 event_rows: None,
8604 handle: None,
8605 writer_task: Some(writer_task),
8606 origin: pool.origin(),
8607 db: crate::timeout_sink::db_label(&pool),
8608 pool: Arc::clone(&pool),
8609 held_lease: None,
8610 };
8611 let row = post_cancel
8612 .query_row(SqlStatement {
8613 sql: "SELECT val FROM reopen_test WHERE id = 1".into(),
8614 params: vec![],
8615 label: None,
8616 })
8617 .await
8618 .expect("read on a queue-backed handle with no transaction connection must pool")
8619 .expect("seeded row must be visible");
8620 assert!(
8621 matches!(&row.columns[0].value, SqlValue::Text(v) if v == "seed"),
8622 "pooled read must return the seeded row; got {:?}",
8623 row.columns[0].value
8624 );
8625 }
8626
8627 #[tokio::test]
8633 async fn writer_query_row_rejects_dml_with_returning_on_queue_backed_handle() {
8634 let dir = tempfile::tempdir().unwrap();
8635 let path = dir.path().join("writer_readonly_returning.db");
8636 let config = PoolConfig {
8637 path: Some(path),
8638 write_queue_enabled: Some(true),
8639 write_routing_strict: true,
8640 ..PoolConfig::for_test()
8641 };
8642 let pool = Arc::new(ConnectionPool::new(config).unwrap());
8643 {
8644 let guard = pool.writer().unwrap();
8645 guard
8646 .conn()
8647 .execute_batch(
8648 "CREATE TABLE IF NOT EXISTS writer_returning_test \
8649 (id INTEGER PRIMARY KEY, val TEXT NOT NULL);
8650 INSERT INTO writer_returning_test (id, val) VALUES (1, 'original');",
8651 )
8652 .unwrap();
8653 }
8654
8655 let bridge = SqlBridge::new(Arc::clone(&pool), true);
8656
8657 let mut w = bridge.writer().await.unwrap();
8658 let result = w
8659 .query_row(SqlStatement {
8660 sql: "UPDATE writer_returning_test SET val = 'mutated' \
8661 WHERE id = ?1 RETURNING val"
8662 .into(),
8663 params: vec![SqlValue::Integer(1)],
8664 label: None,
8665 })
8666 .await;
8667 assert!(
8668 result.is_err(),
8669 "a DML-with-RETURNING statement through query_row on a \
8670 queue-backed writer handle must be rejected, not executed on \
8671 an untracked read-write connection; got {result:?}"
8672 );
8673
8674 let mut reader = bridge.reader().await.unwrap();
8675 let val = reader
8676 .query_scalar(SqlStatement {
8677 sql: "SELECT val FROM writer_returning_test WHERE id = ?1".into(),
8678 params: vec![SqlValue::Integer(1)],
8679 label: None,
8680 })
8681 .await
8682 .unwrap();
8683 assert!(
8684 matches!(&val, Some(SqlValue::Text(v)) if v == "original"),
8685 "the rejected UPDATE...RETURNING must not have altered the row; got {val:?}"
8686 );
8687 }
8688
8689 #[tokio::test]
8700 async fn writer_query_row_rejects_setting_pragma_on_queue_backed_reader_route() {
8701 let dir = tempfile::tempdir().unwrap();
8702 let path = dir.path().join("writer_reader_route_pragma_admission.db");
8703 let config = PoolConfig {
8704 path: Some(path),
8705 write_queue_enabled: Some(true),
8706 write_routing_strict: true,
8707 ..PoolConfig::for_test()
8708 };
8709 let pool = Arc::new(ConnectionPool::new(config).unwrap());
8710 let bridge = SqlBridge::new(Arc::clone(&pool), true);
8711
8712 let mut w = bridge.writer().await.unwrap();
8713 let result = w
8714 .query_row(SqlStatement {
8715 sql: "PRAGMA cache_size = -999".into(),
8716 params: vec![],
8717 label: None,
8718 })
8719 .await;
8720 assert!(
8721 result.is_err(),
8722 "a setting PRAGMA through the queue-backed writer's reader route must be \
8723 refused, exactly like it is through PoolBackedReader/SqliteReader; \
8724 got {result:?}"
8725 );
8726 }
8727
8728 #[tokio::test]
8748 async fn acceptance_five_op_batch_completes_under_concurrent_write_contention() {
8749 let dir = tempfile::tempdir().unwrap();
8750 let path = dir.path().join("acceptance_batch.db");
8751 let config = PoolConfig {
8752 path: Some(path),
8753 write_queue_enabled: Some(true),
8754 write_routing_strict: true,
8755 ..PoolConfig::for_test()
8756 };
8757 let pool = Arc::new(ConnectionPool::new(config).unwrap());
8758 {
8759 let guard = pool.writer().unwrap();
8760 guard
8761 .conn()
8762 .execute_batch(
8763 "CREATE TABLE IF NOT EXISTS acceptance_batch \
8764 (id INTEGER PRIMARY KEY, val TEXT NOT NULL);
8765 INSERT INTO acceptance_batch (id, val) VALUES \
8766 (200, 'seed-0'), (201, 'seed-1'), (202, 'seed-2'), (203, 'seed-3');",
8767 )
8768 .unwrap();
8769 }
8770
8771 let bridge = Arc::new(SqlBridge::new(Arc::clone(&pool), true));
8772
8773 let writer_task = pool
8774 .writer_task_handle()
8775 .unwrap()
8776 .expect("writer task must be spawned for a file-backed pool with the flag on");
8777
8778 let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
8784 let (release_tx, release_rx) = tokio::sync::oneshot::channel::<()>();
8785 let occupier = {
8786 let writer_task = writer_task.clone();
8787 tokio::spawn(async move {
8788 writer_task
8789 .send(move |_conn| {
8790 let _ = started_tx.send(());
8791 let _ = release_rx.blocking_recv();
8792 Ok::<(), StorageError>(())
8793 })
8794 .await
8795 })
8796 };
8797 started_rx
8798 .await
8799 .expect("occupier must signal it has started running inside the writer task");
8800 assert_eq!(
8801 writer_task.queue_depth(),
8802 0,
8803 "channel must start empty once the occupier has been dequeued and is running"
8804 );
8805
8806 let contenders: Vec<_> = (0..3)
8811 .map(|i| {
8812 let bridge = Arc::clone(&bridge);
8813 tokio::spawn(async move {
8814 let mut writer = bridge.writer().await?;
8815 writer
8816 .execute(SqlStatement {
8817 sql: "INSERT INTO acceptance_batch (id, val) VALUES (?1, ?2)".into(),
8818 params: vec![
8819 SqlValue::Integer(100 + i),
8820 SqlValue::Text(format!("contender-{i}")),
8821 ],
8822 label: None,
8823 })
8824 .await
8825 })
8826 })
8827 .collect();
8828
8829 let send = {
8831 let bridge = Arc::clone(&bridge);
8832 tokio::spawn(async move {
8833 let mut writer = bridge.writer().await?;
8834 writer
8835 .execute(SqlStatement {
8836 sql: "INSERT INTO acceptance_batch (id, val) VALUES (?1, ?2)".into(),
8837 params: vec![SqlValue::Integer(1), SqlValue::Text("send".into())],
8838 label: None,
8839 })
8840 .await
8841 })
8842 };
8843 let marks: Vec<_> = (0..4)
8844 .map(|i| {
8845 let bridge = Arc::clone(&bridge);
8846 tokio::spawn(async move {
8847 let mut writer = bridge.writer().await?;
8848 writer
8849 .execute(SqlStatement {
8850 sql: "UPDATE acceptance_batch SET val = ?2 WHERE id = ?1".into(),
8851 params: vec![
8852 SqlValue::Integer(200 + i),
8853 SqlValue::Text(format!("marked-{i}")),
8854 ],
8855 label: None,
8856 })
8857 .await
8858 })
8859 })
8860 .collect();
8861
8862 let mut saw_all_enqueued = false;
8865 for _ in 0..200 {
8866 if writer_task.queue_depth() >= 8 {
8867 saw_all_enqueued = true;
8868 break;
8869 }
8870 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
8871 }
8872 assert!(
8873 saw_all_enqueued,
8874 "not all 8 contending writes (3 contenders + send + 4 marks) reached \
8875 the writer task's channel while the occupier held the single drain \
8876 slot — got depth {}",
8877 writer_task.queue_depth()
8878 );
8879
8880 release_tx
8881 .send(())
8882 .expect("occupier must still be waiting on the release signal");
8883 occupier
8884 .await
8885 .expect("occupier task must not panic")
8886 .expect("occupier write must succeed");
8887
8888 for c in contenders {
8889 c.await
8890 .expect("contender task must not panic")
8891 .expect("contender write must complete without a checkout timeout");
8892 }
8893 send.await
8894 .expect("send task must not panic")
8895 .expect("send op must complete without a checkout timeout");
8896 for (i, m) in marks.into_iter().enumerate() {
8897 let affected = m
8898 .await
8899 .expect("mark task must not panic")
8900 .expect("mark op must complete without a checkout timeout — no starvation");
8901 assert_eq!(
8902 affected, 1,
8903 "mark {i} must have updated exactly its own row"
8904 );
8905 }
8906
8907 let mut reader = bridge.reader().await.unwrap();
8908 let count = reader
8909 .query_scalar(SqlStatement {
8910 sql: "SELECT COUNT(*) FROM acceptance_batch".into(),
8911 params: vec![],
8912 label: None,
8913 })
8914 .await
8915 .unwrap();
8916 assert!(
8917 matches!(count, Some(SqlValue::Integer(8))),
8918 "the 4 seeded mark rows plus 3 contenders plus the batch's own send \
8919 must all be present; got {count:?}"
8920 );
8921
8922 for i in 0..4i64 {
8923 let mut reader = bridge.reader().await.unwrap();
8924 let val = reader
8925 .query_scalar(SqlStatement {
8926 sql: "SELECT val FROM acceptance_batch WHERE id = ?1".into(),
8927 params: vec![SqlValue::Integer(200 + i)],
8928 label: None,
8929 })
8930 .await
8931 .unwrap();
8932 assert!(
8933 matches!(&val, Some(SqlValue::Text(v)) if *v == format!("marked-{i}")),
8934 "mark row {i} must reflect the persisted UPDATE after release; got {val:?}"
8935 );
8936 }
8937 }
8938
8939 #[tokio::test]
8940 async fn standalone_writer_handle_resamples_reserve_before_each_operation() {
8941 use std::sync::atomic::{AtomicUsize, Ordering};
8942
8943 for operation in [
8944 "execute",
8945 "execute_batch",
8946 "execute_script",
8947 "top_level_vacuum",
8948 ] {
8949 let dir = tempfile::tempdir().unwrap();
8950 let mut pool = ConnectionPool::new(PoolConfig {
8951 path: Some(dir.path().join(format!("reserve-{operation}.db"))),
8952 write_queue_enabled: Some(false),
8953 ..PoolConfig::for_test()
8954 })
8955 .unwrap();
8956 let samples = Arc::new(AtomicUsize::new(0));
8957 let observed = Arc::clone(&samples);
8958 pool.set_test_write_admission(100, move |_| {
8959 match observed.fetch_add(1, Ordering::SeqCst) {
8960 0 => Ok(101), 1 => Ok(100), extra => panic!("unexpected capacity sample {extra}"),
8963 }
8964 });
8965 let bridge = SqlBridge::new(Arc::new(pool), true);
8966 let mut writer = bridge.writer().await.unwrap();
8967 writer
8968 .execute(SqlStatement {
8969 sql: "CREATE TABLE reserve_test (id INTEGER PRIMARY KEY)".into(),
8970 params: vec![],
8971 label: None,
8972 })
8973 .await
8974 .unwrap();
8975
8976 let insert = || SqlStatement {
8977 sql: "INSERT INTO reserve_test (id) VALUES (1)".into(),
8978 params: vec![],
8979 label: None,
8980 };
8981 let error = match operation {
8982 "execute" => writer.execute(insert()).await.map(|_| ()),
8983 "execute_batch" => writer.execute_batch(vec![insert()]).await.map(|_| ()),
8984 "execute_script" => {
8985 writer
8986 .execute_script("INSERT INTO reserve_test (id) VALUES (1)".into())
8987 .await
8988 }
8989 "top_level_vacuum" => {
8990 writer
8991 .execute_script_top_level(TopLevelMaintenance::Vacuum)
8992 .await
8993 }
8994 _ => unreachable!(),
8995 }
8996 .expect_err("the second operation must see the new reserve sample");
8997 assert!(
8998 matches!(
8999 &error,
9000 StorageError::CapacityFloor {
9001 available_bytes,
9002 floor_bytes,
9003 ..
9004 } if *available_bytes == 100 && *floor_bytes == 100
9005 ),
9006 "{operation} must retain typed capacity-floor classification: {error:?}"
9007 );
9008 assert_eq!(samples.load(Ordering::SeqCst), 2, "{operation}");
9009 assert!(
9010 matches!(
9011 writer
9012 .query_scalar(SqlStatement {
9013 sql: "SELECT COUNT(*) FROM reserve_test".into(),
9014 params: vec![],
9015 label: None,
9016 })
9017 .await
9018 .unwrap(),
9019 Some(SqlValue::Integer(0))
9020 ),
9021 "{operation} must not write after the reserve refusal"
9022 );
9023 }
9024 }
9025
9026 #[tokio::test]
9027 async fn file_backed_bridge_counts_writer_and_flag_off_atomic_unit_acquisitions() {
9028 let dir = tempfile::tempdir().unwrap();
9029 let config = PoolConfig {
9030 path: Some(dir.path().join("bridge_writer_acquisitions.db")),
9031 write_queue_enabled: Some(false),
9032 ..PoolConfig::for_test()
9033 };
9034 let pool = Arc::new(ConnectionPool::new(config).unwrap());
9035 let bridge = SqlBridge::new(Arc::clone(&pool), true);
9036
9037 let before = pool.writer_acquisition_snapshot();
9038
9039 drop(bridge.writer().await.unwrap());
9040 let after_writer = pool.writer_acquisition_snapshot();
9041 assert_eq!(
9042 after_writer.standalone_acquisitions,
9043 before.standalone_acquisitions + 1
9044 );
9045 assert_eq!(after_writer.acquisitions, before.acquisitions + 1);
9046 assert_eq!(after_writer.pooled_acquisitions, before.pooled_acquisitions);
9047 assert_eq!(
9048 after_writer.writer_task_acquisitions,
9049 before.writer_task_acquisitions
9050 );
9051
9052 let op: AtomicUnitOp = Box::new(|_writer| {
9053 Box::pin(async { Ok(Box::new(()) as Box<dyn std::any::Any + Send>) })
9054 });
9055 bridge.atomic_unit(op).await.unwrap();
9056
9057 let after_atomic_unit = pool.writer_acquisition_snapshot();
9058 assert_eq!(
9059 after_atomic_unit.standalone_acquisitions,
9060 before.standalone_acquisitions + 2
9061 );
9062 assert_eq!(after_atomic_unit.acquisitions, before.acquisitions + 2);
9063 assert_eq!(
9064 after_atomic_unit.pooled_acquisitions,
9065 before.pooled_acquisitions
9066 );
9067 assert_eq!(
9068 after_atomic_unit.writer_task_acquisitions,
9069 before.writer_task_acquisitions
9070 );
9071 }
9072
9073 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
9086 async fn max_completed_hold_operation_names_the_caller_read_not_a_bridge_constant() {
9087 async fn recorded_hold_operation(through_writer: bool) -> (Option<&'static str>, u64) {
9088 let dir = tempfile::tempdir().unwrap();
9089 let path = dir.path().join("hold_attribution.db");
9090 let pool = Arc::new(
9091 ConnectionPool::new(PoolConfig {
9092 path: Some(path),
9093 write_queue_enabled: Some(true),
9094 write_routing_strict: true,
9095 ..PoolConfig::for_test()
9096 })
9097 .unwrap(),
9098 );
9099 {
9100 let guard = pool.writer().unwrap();
9101 guard
9102 .conn()
9103 .execute_batch("CREATE TABLE IF NOT EXISTS hold_attr (id INTEGER PRIMARY KEY)")
9104 .unwrap();
9105 }
9106 let bridge = SqlBridge::new(Arc::clone(&pool), true);
9107 let statement = SqlStatement {
9108 sql: "SELECT id FROM hold_attr".into(),
9109 params: vec![],
9110 label: None,
9111 };
9112 if through_writer {
9113 let mut w = bridge.writer().await.unwrap();
9114 w.query_row(statement).await.unwrap();
9115 } else {
9116 let mut r = bridge.reader().await.unwrap();
9117 r.query_all(statement).await.unwrap();
9118 }
9119 let snapshot = pool.reader_acquisition_snapshot();
9120 (
9121 snapshot.max_completed_hold_operation,
9122 snapshot.completed_pooled_checkouts,
9123 )
9124 }
9125
9126 let (read_side, read_completed) = recorded_hold_operation(false).await;
9127 let (write_side, write_completed) = recorded_hold_operation(true).await;
9128
9129 assert!(
9130 read_completed >= 1 && write_completed >= 1,
9131 "both arms must complete a pooled checkout, or the attribution \
9132 below is reading an empty population; got {read_completed} and \
9133 {write_completed}"
9134 );
9135 assert_ne!(
9141 read_side, write_side,
9142 "the recorded operation must tell two different reads apart; a \
9143 shared bridge constant makes these equal while still looking \
9144 like an answer"
9145 );
9146 assert_eq!(
9147 read_side,
9148 Some("query_all"),
9149 "a pooled read drawn through the reader must record its own operation"
9150 );
9151 assert_eq!(
9152 write_side,
9153 Some("writer.query_row"),
9154 "a pooled read drawn through the writer must record its own operation"
9155 );
9156 }
9157}
9158
9159#[cfg(test)]
9160#[path = "sql_bridge/direct_busy_tests.rs"]
9161mod direct_busy_tests;
9162
9163#[cfg(test)]
9164#[path = "sql_bridge/settlement_hazard_tests.rs"]
9165mod settlement_hazard_tests;