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::writer_task_terminated(
3408 khive_storage::WriterTaskRequestState::SideEffectsUnknown,
3409 ));
3410 }
3411 if let Err(error) = conn.execute_batch("BEGIN IMMEDIATE") {
3412 if !conn.is_autocommit() {
3413 pool.retire_pooled_writer(conn);
3414 return Err(StorageError::writer_task_terminated(
3415 khive_storage::WriterTaskRequestState::SideEffectsUnknown,
3416 ));
3417 }
3418 return Err(map_rusqlite_err(error, "atomic_unit.begin"))
3419 .inspect_err(|error| pool.record_direct_writer_error(error));
3420 }
3421 let _tx_handle = khive_storage::tx_registry::register_scoped(
3422 Some("atomic_unit".to_string()),
3423 pool.origin(),
3424 );
3425 let (result, terminal_state) = crate::writer_task::execute_wrapped_transaction(
3426 conn,
3427 "atomic_unit.commit",
3428 |conn| {
3429 let mut inline = InlineWriter {
3430 event_rows: Some(Arc::clone(&pending_event_rows)),
3431 conn: conn as *const rusqlite::Connection,
3432 };
3433 block_on_sync(op(&mut inline)).and_then(|result| result)
3434 },
3435 );
3436 if terminal_state.is_some() {
3437 pool.retire_pooled_writer(conn);
3438 }
3439 result.inspect_err(|error| pool.record_direct_writer_error(error))
3440 })
3441 .await
3442 .map_err(|error| {
3443 StorageError::driver(StorageCapability::Sql, "atomic_unit", error)
3444 })?
3445 }
3446 }
3447 .await;
3448 khive_storage::usage::account_event_write(
3449 result.as_ref().map(|_| event_rows.committed_rows()),
3450 );
3451 result
3452 }
3453}
3454
3455#[cfg(test)]
3456#[path = "sql_bridge_tests.rs"]
3457mod tests;
3458
3459#[cfg(test)]
3460#[path = "sql_bridge/direct_busy_tests.rs"]
3461mod direct_busy_tests;
3462
3463#[cfg(test)]
3464#[path = "sql_bridge/settlement_hazard_tests.rs"]
3465mod settlement_hazard_tests;