#[path = "sql_bridge/write_errors.rs"]
mod write_errors;
#[path = "sql_bridge/manual_atomic.rs"]
mod manual_atomic;
use manual_atomic::run_manual_atomic_unit;
#[path = "sql_bridge/standalone_admission.rs"]
mod standalone_admission;
use standalone_admission::{
acquire_standalone_lease, acquire_unit_lease, admit_standalone_operation, StandaloneWriteError,
};
mod standalone_batch;
#[cfg(test)]
use standalone_batch::BatchPoisonReason;
use standalone_batch::{
run_standalone_batch, run_standalone_script, run_standalone_statement,
run_standalone_top_level, BatchFailure, PoisonedBatchError,
};
use std::any::Any;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::Instant;
use async_trait::async_trait;
use khive_storage::error::StorageError;
use khive_storage::types::{PageRequest, SqlColumn, SqlRow, SqlStatement, SqlValue};
use khive_storage::{AtomicUnitOp, StorageCapability, TopLevelMaintenance};
use tokio::sync::{OwnedSemaphorePermit, Semaphore};
use crate::error::SqliteError;
use crate::pool::{ConnectionPool, SharedReaderTransactionGuard, StandaloneReaderPurpose};
mod rows;
pub(crate) use rows::bind_params;
use rows::{
prepare_batch_statements, prepare_cached_sql_statement, prepare_sql_statement, row_to_sql_row,
AtomicEventRows, PreparedBatchStatement,
};
#[cfg(test)]
use rows::{COUNTED_EVENT_INSERT_LABELS, ROW_CONVERSIONS};
#[cfg(test)]
#[path = "atomic_event_usage_tests.rs"]
mod atomic_event_usage_tests;
fn execute_prepared_batch<'conn>(
conn: &'conn rusqlite::Connection,
prepared: Vec<PreparedBatchStatement<'conn>>,
statements: &[SqlStatement],
event_rows: Option<&AtomicEventRows>,
) -> Result<u64, rusqlite::Error> {
debug_assert_eq!(prepared.len(), statements.len());
let mut total = 0u64;
for (prepared, statement) in prepared.into_iter().zip(statements) {
let mut prepared = match prepared {
PreparedBatchStatement::Ready(prepared) => prepared,
PreparedBatchStatement::PrepareAtExecution => {
prepare_sql_statement(conn, &statement.sql)?
}
};
bind_params(&mut prepared, &statement.params)?;
let affected = prepared.raw_execute()? as u64;
if let Some(event_rows) = event_rows {
event_rows.observe(statement, affected);
}
total += affected;
}
Ok(total)
}
const TRANSACTION_CONTROL_KEYWORDS: [&str; 7] = [
"BEGIN",
"START",
"COMMIT",
"END",
"ROLLBACK",
"SAVEPOINT",
"RELEASE",
];
fn skip_sqlite_empty_prefix(mut rest: &[u8]) -> &[u8] {
loop {
let mut idx = 0;
while idx < rest.len() && rest[idx].is_ascii_whitespace() {
idx += 1;
}
rest = &rest[idx..];
if let Some(tail) = rest.strip_prefix(b"\xEF\xBB\xBF") {
rest = tail;
continue;
}
if let Some(tail) = rest.strip_prefix(b";") {
rest = tail;
continue;
}
if let Some(tail) = rest.strip_prefix(b"--") {
let mut idx = 0;
while idx < tail.len() && tail[idx] != b'\n' {
idx += 1;
}
rest = if idx < tail.len() {
&tail[idx + 1..]
} else {
&[]
};
continue;
}
if let Some(tail) = rest.strip_prefix(b"/*") {
let mut idx = 0;
while idx + 1 < tail.len() && !(tail[idx] == b'*' && tail[idx + 1] == b'/') {
idx += 1;
}
rest = if idx + 1 < tail.len() {
&tail[idx + 2..]
} else {
&[]
};
continue;
}
break;
}
rest
}
fn next_sqlite_token(mut rest: &[u8]) -> Option<(&[u8], &[u8])> {
loop {
let mut idx = 0;
while idx < rest.len() && rest[idx].is_ascii_whitespace() {
idx += 1;
}
rest = &rest[idx..];
if let Some(tail) = rest.strip_prefix(b"\xEF\xBB\xBF") {
rest = tail;
continue;
}
if let Some(tail) = rest.strip_prefix(b"--") {
let mut idx = 0;
while idx < tail.len() && tail[idx] != b'\n' {
idx += 1;
}
rest = if idx < tail.len() {
&tail[idx + 1..]
} else {
&[]
};
continue;
}
if let Some(tail) = rest.strip_prefix(b"/*") {
let mut idx = 0;
while idx + 1 < tail.len() && !(tail[idx] == b'*' && tail[idx + 1] == b'/') {
idx += 1;
}
rest = if idx + 1 < tail.len() {
&tail[idx + 2..]
} else {
&[]
};
continue;
}
break;
}
let len = rest
.iter()
.take_while(|byte| byte.is_ascii_alphanumeric() || **byte == b'_')
.count();
(len != 0).then_some((&rest[..len], &rest[len..]))
}
fn transaction_control_parts(sql: &str) -> Option<(&'static str, &[u8])> {
let rest = skip_sqlite_empty_prefix(sql.as_bytes());
TRANSACTION_CONTROL_KEYWORDS
.iter()
.copied()
.find_map(|keyword| {
let kw = keyword.as_bytes();
if rest.len() < kw.len() || !rest[..kw.len()].eq_ignore_ascii_case(kw) {
return None;
}
let boundary = match rest.get(kw.len()) {
Some(next) => !(next.is_ascii_alphanumeric() || *next == b'_'),
None => true,
};
boundary.then_some((keyword, &rest[kw.len()..]))
})
}
fn transaction_control_head(sql: &str) -> Option<&'static str> {
transaction_control_parts(sql).map(|(keyword, _)| keyword)
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum CachedReadTransactionControl {
BeginDeferred,
Finish(&'static str),
Unsupported(&'static str),
}
fn cached_read_transaction_control(sql: &str) -> Option<CachedReadTransactionControl> {
let (keyword, tail) = transaction_control_parts(sql)?;
match keyword {
"BEGIN" => {
let mut rest = tail;
let mut saw_deferred = false;
let mut saw_transaction = false;
while let Some((token, next)) = next_sqlite_token(rest) {
if !saw_deferred && !saw_transaction && token.eq_ignore_ascii_case(b"DEFERRED") {
saw_deferred = true;
} else if !saw_transaction && token.eq_ignore_ascii_case(b"TRANSACTION") {
saw_transaction = true;
} else {
return Some(CachedReadTransactionControl::Unsupported(keyword));
}
rest = next;
}
if !skip_sqlite_empty_prefix(rest).is_empty() {
return Some(CachedReadTransactionControl::Unsupported(keyword));
}
Some(CachedReadTransactionControl::BeginDeferred)
}
"COMMIT" | "END" => Some(CachedReadTransactionControl::Finish(keyword)),
"ROLLBACK" => {
let first = next_sqlite_token(tail);
let rollback_target = match first {
Some((token, rest)) if token.eq_ignore_ascii_case(b"TRANSACTION") => {
next_sqlite_token(rest).map(|(token, _)| token)
}
Some((token, _)) => Some(token),
None => None,
};
if rollback_target.is_some_and(|token| token.eq_ignore_ascii_case(b"TO")) {
Some(CachedReadTransactionControl::Unsupported(keyword))
} else {
Some(CachedReadTransactionControl::Finish(keyword))
}
}
_ => Some(CachedReadTransactionControl::Unsupported(keyword)),
}
}
fn reject_transaction_control_statements(
statements: &[SqlStatement],
operation: &'static str,
) -> khive_storage::types::StorageResult<()> {
for (index, statement) in statements.iter().enumerate() {
if let Some(keyword) = transaction_control_head(&statement.sql) {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message: format!(
"statement at index {index} is transaction control ({keyword}); \
execute_batch owns the BEGIN/COMMIT boundary for the whole \
batch — remove transaction-control statements from the batch"
),
});
}
}
Ok(())
}
fn settle_pooled_call<T>(
guard: &crate::pool::WriterGuard<'_>,
operation: &'static str,
result: khive_storage::types::StorageResult<T>,
) -> khive_storage::types::StorageResult<T> {
if guard.is_autocommit() {
return result;
}
if let Err(settlement) = guard.rollback_or_retire("pooled call left its transaction open") {
if let Err(error) = &result {
tracing::warn!(
operation,
%error,
"pooled call failed inside a transaction it opened, and its rollback \
could not prove autocommit; reporting the settlement failure"
);
}
return Err(settlement.into_storage_error(StorageCapability::Sql, operation));
}
result?;
Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message: "the call left a transaction open; it was rolled back before the pooled \
writer was released — use atomic_unit to run statements as one transaction"
.into(),
})
}
fn prepare_bound_statement<'conn>(
conn: &'conn rusqlite::Connection,
statement: &SqlStatement,
) -> Result<rusqlite::Statement<'conn>, rusqlite::Error> {
let mut stmt = prepare_sql_statement(conn, &statement.sql)?;
bind_params(&mut stmt, &statement.params)?;
Ok(stmt)
}
fn execute_prepared_query(
mut stmt: rusqlite::Statement<'_>,
) -> Result<Vec<SqlRow>, rusqlite::Error> {
let col_count = stmt.column_count();
let col_names: Vec<String> = (0..col_count)
.map(|i| stmt.column_name(i).unwrap_or("").to_string())
.collect();
let mut rows = Vec::new();
let mut raw_rows = stmt.raw_query();
while let Some(row) = raw_rows.next()? {
rows.push(row_to_sql_row(row, col_count, &col_names));
}
Ok(rows)
}
fn execute_prepared_query_row(
mut stmt: rusqlite::Statement<'_>,
) -> Result<Option<SqlRow>, rusqlite::Error> {
let col_count = stmt.column_count();
let col_names: Vec<String> = (0..col_count)
.map(|i| stmt.column_name(i).unwrap_or("").to_string())
.collect();
let mut raw_rows = stmt.raw_query();
Ok(raw_rows
.next()?
.map(|row| row_to_sql_row(row, col_count, &col_names)))
}
fn execute_prepared_query_page(
mut stmt: rusqlite::Statement<'_>,
page: &PageRequest,
) -> Result<Vec<SqlRow>, rusqlite::Error> {
if page.limit == 0 {
return Ok(Vec::new());
}
let col_count = stmt.column_count();
let col_names: Vec<String> = (0..col_count)
.map(|i| stmt.column_name(i).unwrap_or("").to_string())
.collect();
let mut rows = Vec::new();
let mut offset = page.offset;
let mut remaining = u64::from(page.limit);
let mut raw_rows = stmt.raw_query();
while remaining > 0 {
let Some(row) = raw_rows.next()? else {
break;
};
if offset > 0 {
offset -= 1;
continue;
}
rows.push(row_to_sql_row(row, col_count, &col_names));
remaining -= 1;
}
Ok(rows)
}
fn execute_query(
conn: &rusqlite::Connection,
statement: &SqlStatement,
) -> Result<Vec<SqlRow>, rusqlite::Error> {
execute_prepared_query(prepare_bound_statement(conn, statement)?)
}
fn execute_query_row(
conn: &rusqlite::Connection,
statement: &SqlStatement,
) -> Result<Option<SqlRow>, rusqlite::Error> {
execute_prepared_query_row(prepare_bound_statement(conn, statement)?)
}
fn execute_query_page(
conn: &rusqlite::Connection,
statement: &SqlStatement,
page: &PageRequest,
) -> Result<Vec<SqlRow>, rusqlite::Error> {
execute_prepared_query_page(prepare_bound_statement(conn, statement)?, page)
}
fn statement_is_cancellable_read(stmt: &rusqlite::Statement<'_>, sql: &str) -> bool {
stmt.readonly() && transaction_control_head(sql).is_none()
}
const READER_STRUCTURAL_PRAGMAS: [&str; 8] = [
"table_info",
"table_xinfo",
"table_list",
"index_list",
"index_info",
"index_xinfo",
"foreign_key_list",
"integrity_check",
];
const READER_SETTING_PRAGMAS: [&str; 10] = [
"database_list",
"collation_list",
"function_list",
"compile_options",
"page_count",
"freelist_count",
"user_version",
"schema_version",
"journal_mode",
"page_size",
];
fn skip_balanced_parens(rest: &[u8]) -> Option<&[u8]> {
debug_assert_eq!(rest.first(), Some(&b'('));
let mut depth: u32 = 0;
let mut idx = 0;
loop {
match *rest.get(idx)? {
b'(' => {
depth += 1;
idx += 1;
}
b')' => {
depth -= 1;
idx += 1;
if depth == 0 {
return Some(&rest[idx..]);
}
}
quote @ (b'\'' | b'"' | b'`') => {
idx += 1;
loop {
match *rest.get(idx)? {
byte if byte == quote => {
idx += 1;
if rest.get(idx) == Some("e) {
idx += 1; } else {
break;
}
}
_ => idx += 1,
}
}
}
b'[' => {
idx += 1;
while *rest.get(idx)? != b']' {
idx += 1;
}
idx += 1;
}
b'-' if rest.get(idx + 1) == Some(&b'-') => {
idx += 2;
while idx < rest.len() && rest[idx] != b'\n' {
idx += 1;
}
}
b'/' if rest.get(idx + 1) == Some(&b'*') => {
idx += 2;
while idx + 1 < rest.len() && !(rest[idx] == b'*' && rest[idx + 1] == b'/') {
idx += 1;
}
idx = (idx + 2).min(rest.len());
}
_ => idx += 1,
}
}
}
fn skip_sqlite_identifier(rest: &[u8]) -> Option<&[u8]> {
match *rest.first()? {
quote @ (b'"' | b'`') => {
let mut idx = 1;
loop {
match *rest.get(idx)? {
byte if byte == quote => {
idx += 1;
if rest.get(idx) == Some("e) {
idx += 1; } else {
break;
}
}
_ => idx += 1,
}
}
Some(&rest[idx..])
}
b'[' => {
let mut idx = 1;
while *rest.get(idx)? != b']' {
idx += 1;
}
Some(&rest[idx + 1..])
}
_ => next_sqlite_token(rest).map(|(_, next)| next),
}
}
fn skip_common_table_expressions(tail: &[u8]) -> Option<&[u8]> {
let mut rest = skip_sqlite_empty_prefix(tail);
if let Some((word, next)) = next_sqlite_token(rest) {
if word.eq_ignore_ascii_case(b"RECURSIVE") {
rest = skip_sqlite_empty_prefix(next);
}
}
loop {
let next = skip_sqlite_identifier(rest)?;
rest = skip_sqlite_empty_prefix(next);
if rest.first() == Some(&b'(') {
rest = skip_sqlite_empty_prefix(skip_balanced_parens(rest)?);
}
let (as_keyword, next) = next_sqlite_token(rest)?;
if !as_keyword.eq_ignore_ascii_case(b"AS") {
return None;
}
rest = skip_sqlite_empty_prefix(next);
if let Some((word, next)) = next_sqlite_token(rest) {
if word.eq_ignore_ascii_case(b"MATERIALIZED") {
rest = skip_sqlite_empty_prefix(next);
} else if word.eq_ignore_ascii_case(b"NOT") {
let (materialized, next) = next_sqlite_token(skip_sqlite_empty_prefix(next))?;
if !materialized.eq_ignore_ascii_case(b"MATERIALIZED") {
return None;
}
rest = skip_sqlite_empty_prefix(next);
}
}
if rest.first() != Some(&b'(') {
return None;
}
rest = skip_sqlite_empty_prefix(skip_balanced_parens(rest)?);
if rest.first() == Some(&b',') {
rest = skip_sqlite_empty_prefix(&rest[1..]);
continue;
}
return Some(rest);
}
}
pub(crate) fn reader_capability_admits(sql: &str) -> Result<(), String> {
let rest = skip_sqlite_empty_prefix(sql.as_bytes());
let Some((head, tail)) = next_sqlite_token(rest) else {
return Ok(());
};
if head.eq_ignore_ascii_case(b"SELECT") || head.eq_ignore_ascii_case(b"VALUES") {
return Ok(());
}
if head.eq_ignore_ascii_case(b"WITH") {
let Some(after_ctes) = skip_common_table_expressions(tail) else {
return Err(
"WITH statement's common-table-expression list could not be parsed; refusing \
to admit it through the reader capability"
.into(),
);
};
return match next_sqlite_token(after_ctes) {
Some((main_head, _))
if main_head.eq_ignore_ascii_case(b"SELECT")
|| main_head.eq_ignore_ascii_case(b"VALUES") =>
{
Ok(())
}
other => Err(format!(
"WITH ... {:?} is not admitted through the reader capability; only a \
read-only SELECT/VALUES body after the CTE list may run against a pooled \
reader connection",
other.map_or_else(
|| "<none>".to_string(),
|(main_head, _)| String::from_utf8_lossy(main_head).into_owned()
)
)),
};
}
if head.eq_ignore_ascii_case(b"EXPLAIN") {
let mut rest = skip_sqlite_empty_prefix(tail);
if let Some((query, next)) = next_sqlite_token(rest) {
if query.eq_ignore_ascii_case(b"QUERY") {
let after_query = skip_sqlite_empty_prefix(next);
match next_sqlite_token(after_query) {
Some((plan, next2)) if plan.eq_ignore_ascii_case(b"PLAN") => {
rest = skip_sqlite_empty_prefix(next2);
}
_ => {
return Err(
"EXPLAIN QUERY must be followed by PLAN through the reader capability"
.into(),
);
}
}
}
}
return reader_capability_admits(&String::from_utf8_lossy(rest));
}
if head.eq_ignore_ascii_case(b"PRAGMA") {
return reader_capability_admits_pragma(tail);
}
Err(format!(
"statement head {:?} is not admitted through the reader capability; only \
SELECT/WITH/VALUES/EXPLAIN and an allow-listed set of read-only PRAGMA forms \
may run against a pooled reader connection",
String::from_utf8_lossy(head)
))
}
fn reader_capability_admits_pragma(tail: &[u8]) -> Result<(), String> {
let rest = skip_sqlite_empty_prefix(tail);
let Some((mut name, mut after_name)) = next_sqlite_token(rest) else {
return Err("PRAGMA with no name is not admitted through the reader capability".into());
};
if after_name.first() == Some(&b'.') {
let (qualified_name, qualified_after) =
next_sqlite_token(&after_name[1..]).ok_or_else(|| {
"PRAGMA schema-qualifier with no pragma name is not admitted through the \
reader capability"
.to_string()
})?;
name = qualified_name;
after_name = qualified_after;
}
let after = skip_sqlite_empty_prefix(after_name);
let is_structural = READER_STRUCTURAL_PRAGMAS
.iter()
.any(|allowed| name.eq_ignore_ascii_case(allowed.as_bytes()));
let is_setting = READER_SETTING_PRAGMAS
.iter()
.any(|allowed| name.eq_ignore_ascii_case(allowed.as_bytes()));
if !is_structural && !is_setting {
return Err(format!(
"PRAGMA {:?} is not admitted through the reader capability",
String::from_utf8_lossy(name)
));
}
if after.first() == Some(&b'=') {
return Err(format!(
"PRAGMA {:?} may not be assigned through the reader capability",
String::from_utf8_lossy(name)
));
}
if after.first() == Some(&b'(') && !is_structural {
return Err(format!(
"PRAGMA {:?} may not carry an argument through the reader capability",
String::from_utf8_lossy(name)
));
}
Ok(())
}
fn admit_reader_capability_sql(
statement: &SqlStatement,
transaction_control: Option<CachedReadTransactionControl>,
operation: &'static str,
) -> khive_storage::types::StorageResult<()> {
if matches!(
transaction_control,
Some(CachedReadTransactionControl::BeginDeferred)
| Some(CachedReadTransactionControl::Finish(_))
) {
return Ok(());
}
reader_capability_admits(&statement.sql).map_err(|message| StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message,
})
}
fn execute_query_interruptibly(
scope: &crate::read_cancellation::InterruptibleReadScope,
conn: &rusqlite::Connection,
statement: &SqlStatement,
operation: &'static str,
rollback_interrupted_transaction: bool,
interruptible: bool,
) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
let stmt = prepare_bound_statement(conn, statement)
.map_err(|error| map_rusqlite_err(error, operation))?;
if interruptible && statement_is_cancellable_read(&stmt, &statement.sql) {
scope.run_with_interrupted_cleanup(
conn,
move || {
execute_prepared_query(stmt).map_err(|error| map_rusqlite_err(error, operation))
},
|| {
rollback_interrupted_read_transaction(
conn,
operation,
rollback_interrupted_transaction,
)
},
)
} else {
scope.mark_write_committed()?;
execute_prepared_query(stmt).map_err(|error| map_rusqlite_err(error, operation))
}
}
fn execute_query_row_interruptibly(
scope: &crate::read_cancellation::InterruptibleReadScope,
conn: &rusqlite::Connection,
statement: &SqlStatement,
operation: &'static str,
rollback_interrupted_transaction: bool,
interruptible: bool,
) -> khive_storage::types::StorageResult<Option<SqlRow>> {
let stmt = prepare_bound_statement(conn, statement)
.map_err(|error| map_rusqlite_err(error, operation))?;
if interruptible && statement_is_cancellable_read(&stmt, &statement.sql) {
scope.run_with_interrupted_cleanup(
conn,
move || {
execute_prepared_query_row(stmt).map_err(|error| map_rusqlite_err(error, operation))
},
|| {
rollback_interrupted_read_transaction(
conn,
operation,
rollback_interrupted_transaction,
)
},
)
} else {
scope.mark_write_committed()?;
execute_prepared_query_row(stmt).map_err(|error| map_rusqlite_err(error, operation))
}
}
fn execute_query_page_interruptibly(
scope: &crate::read_cancellation::InterruptibleReadScope,
conn: &rusqlite::Connection,
statement: &SqlStatement,
page: &PageRequest,
operation: &'static str,
rollback_interrupted_transaction: bool,
interruptible: bool,
) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
let stmt = prepare_bound_statement(conn, statement)
.map_err(|error| map_rusqlite_err(error, operation))?;
if interruptible && statement_is_cancellable_read(&stmt, &statement.sql) {
scope.run_with_interrupted_cleanup(
conn,
move || {
execute_prepared_query_page(stmt, page)
.map_err(|error| map_rusqlite_err(error, operation))
},
|| {
rollback_interrupted_read_transaction(
conn,
operation,
rollback_interrupted_transaction,
)
},
)
} else {
scope.mark_write_committed()?;
execute_prepared_query_page(stmt, page).map_err(|error| map_rusqlite_err(error, operation))
}
}
fn rollback_interrupted_read_transaction(
conn: &rusqlite::Connection,
operation: &'static str,
enabled: bool,
) -> khive_storage::types::StorageResult<()> {
if !enabled || conn.is_autocommit() {
return Ok(());
}
conn.execute_batch("ROLLBACK")
.map_err(|error| map_rusqlite_err(error, operation))?;
if conn.is_autocommit() {
Ok(())
} else {
Err(StorageError::Transaction {
operation: operation.into(),
message: "interrupted read transaction rollback did not restore autocommit".into(),
})
}
}
fn map_rusqlite_err(e: rusqlite::Error, op: &'static str) -> StorageError {
StorageError::driver(StorageCapability::Sql, op, e)
}
#[derive(Clone, Copy)]
enum SlotTimeoutClass {
Admission,
ReaderContract,
}
async fn acquire_reader_handle_slot(
pool: &ConnectionPool,
operation: &'static str,
class: SlotTimeoutClass,
) -> Result<OwnedSemaphorePermit, StorageError> {
let result = acquire_handle_slot(
pool.sql_bridge_reader_slots(),
pool.config().checkout_timeout,
operation,
class,
)
.await;
if matches!(
&result,
Err(StorageError::Timeout { .. } | StorageError::AdmissionTimeout { .. })
) {
pool.record_reader_admission_timeout();
}
result
}
pub(crate) async fn acquire_in_memory_write_unit(
pool: &ConnectionPool,
operation: &'static str,
) -> Result<OwnedSemaphorePermit, StorageError> {
acquire_handle_slot(
pool.sql_bridge_writer_slots(),
pool.config().checkout_timeout,
operation,
SlotTimeoutClass::Admission,
)
.await
}
async fn acquire_handle_slot(
slots: Arc<Semaphore>,
timeout: std::time::Duration,
operation: &'static str,
class: SlotTimeoutClass,
) -> Result<OwnedSemaphorePermit, StorageError> {
tokio::time::timeout(timeout, slots.acquire_owned())
.await
.map_err(|_| match class {
SlotTimeoutClass::Admission => StorageError::AdmissionTimeout {
operation: operation.into(),
timeout_ms: u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX),
pool_identity: None,
},
SlotTimeoutClass::ReaderContract => StorageError::Timeout {
operation: operation.into(),
},
})?
.map_err(|error| StorageError::Pool {
operation: operation.into(),
message: error.to_string(),
})
}
fn open_standalone_reader(pool: &ConnectionPool) -> Result<rusqlite::Connection, StorageError> {
pool.open_standalone_reader(StandaloneReaderPurpose::ExplicitSqlReadTransaction)
.map_err(|error| StorageError::driver(StorageCapability::Sql, "open_reader", error))
}
#[cfg(test)]
fn open_standalone_writer(pool: &ConnectionPool) -> Result<rusqlite::Connection, StorageError> {
let conn = pool
.open_standalone_writer()
.map_err(|e| e.into_storage_error(StorageCapability::Sql, "open_writer"))?;
configure_standalone_writer(pool, conn)
}
fn open_admitted_standalone_writer(
pool: &ConnectionPool,
) -> Result<rusqlite::Connection, StorageError> {
let conn = pool
.open_standalone_writer_for_admitted_operation()
.map_err(|e| e.into_storage_error(StorageCapability::Sql, "open_writer"))?;
configure_standalone_writer(pool, conn)
}
fn configure_standalone_writer(
pool: &ConnectionPool,
conn: rusqlite::Connection,
) -> Result<rusqlite::Connection, StorageError> {
let config = pool.config();
conn.busy_timeout(config.busy_timeout)
.map_err(|e| map_rusqlite_err(e, "open_writer"))?;
conn.pragma_update(None, "cache_size", "-65536")
.map_err(|e| map_rusqlite_err(e, "open_writer"))?;
conn.pragma_update(None, "mmap_size", "1073741824")
.map_err(|e| map_rusqlite_err(e, "open_writer"))?;
Ok(conn)
}
async fn open_standalone_on_blocking<F>(
pool: Arc<ConnectionPool>,
slot: OwnedSemaphorePermit,
operation: &'static str,
open: F,
) -> khive_storage::types::StorageResult<(rusqlite::Connection, OwnedSemaphorePermit)>
where
F: FnOnce(&ConnectionPool) -> Result<rusqlite::Connection, StorageError> + Send + 'static,
{
tokio::task::spawn_blocking(move || open(&pool).map(|conn| (conn, slot)))
.await
.map_err(|e| StorageError::driver(StorageCapability::Sql, operation, e))?
}
async fn open_standalone_reader_on_blocking(
pool: Arc<ConnectionPool>,
slot: OwnedSemaphorePermit,
) -> khive_storage::types::StorageResult<(rusqlite::Connection, OwnedSemaphorePermit)> {
open_standalone_on_blocking(pool, slot, "open_reader", open_standalone_reader).await
}
async fn open_standalone_writer_on_blocking(
pool: Arc<ConnectionPool>,
slot: OwnedSemaphorePermit,
) -> khive_storage::types::StorageResult<(rusqlite::Connection, OwnedSemaphorePermit)> {
open_standalone_on_blocking(pool, slot, "open_writer", open_admitted_standalone_writer).await
}
const CACHED_READ_TRANSACTION_LABEL: &str = "sql_bridge_cached_read_transaction";
struct CachedReadTransaction {
_slot: OwnedSemaphorePermit,
_tx_handle: khive_storage::tx_registry::TxHandle,
opened_at: Instant,
}
struct StandaloneHandle {
conn: rusqlite::Connection,
_retained_slot: Option<OwnedSemaphorePermit>,
read_transaction_slot: Option<CachedReadTransaction>,
}
impl StandaloneHandle {
fn is_cached_reader(&self) -> bool {
self._retained_slot.is_none()
}
fn has_read_transaction(&self) -> bool {
self.read_transaction_slot.is_some()
}
}
struct SqliteReader {
handle: Option<StandaloneHandle>,
pool: Arc<ConnectionPool>,
poisoned: bool,
}
async fn open_explicit_read_transaction_handle(
pool: Arc<ConnectionPool>,
) -> khive_storage::types::StorageResult<StandaloneHandle> {
let open_slot = crate::await_request_read_phase(
"sql_bridge.reader_open",
acquire_reader_handle_slot(
&pool,
"sql_bridge.reader_open",
SlotTimeoutClass::ReaderContract,
),
)
.await??;
let (conn, open_slot) = crate::await_request_read_phase(
"sql_bridge.reader_open",
open_standalone_reader_on_blocking(pool, open_slot),
)
.await??;
drop(open_slot);
Ok(StandaloneHandle {
conn,
_retained_slot: None,
read_transaction_slot: None,
})
}
impl SqliteReader {
async fn use_explicit_transaction_handle(
&mut self,
transaction_control: Option<CachedReadTransactionControl>,
operation: &'static str,
) -> khive_storage::types::StorageResult<bool> {
if self.poisoned {
return Err(StorageError::Pool {
operation: operation.into(),
message: "connection already consumed".into(),
});
}
if self.handle.is_some() {
return Ok(true);
}
match transaction_control {
None => Ok(false),
Some(CachedReadTransactionControl::BeginDeferred) => {
self.handle =
Some(open_explicit_read_transaction_handle(Arc::clone(&self.pool)).await?);
Ok(true)
}
Some(CachedReadTransactionControl::Finish(keyword))
| Some(CachedReadTransactionControl::Unsupported(keyword)) => {
Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message: format!(
"cached read-only handle has no admitted transaction for transaction \
control ({keyword})"
),
})
}
}
}
fn close_inactive_transaction_handle(&mut self) {
if self
.handle
.as_ref()
.is_some_and(|handle| handle.is_cached_reader() && !handle.has_read_transaction())
{
drop(self.handle.take());
}
}
}
async fn execute_standalone_read<R, F>(
handle: &mut Option<StandaloneHandle>,
pool: Arc<ConnectionPool>,
operation: &'static str,
transaction_control: Option<CachedReadTransactionControl>,
read: F,
) -> khive_storage::types::StorageResult<R>
where
R: Send + 'static,
F: FnOnce(
&crate::read_cancellation::InterruptibleReadScope,
&rusqlite::Connection,
bool,
bool,
) -> khive_storage::types::StorageResult<R>
+ Send
+ 'static,
{
if handle.is_none() {
return Err(StorageError::Pool {
operation: operation.into(),
message: "connection already consumed".into(),
});
}
let active_read_transaction = handle
.as_ref()
.is_some_and(|handle| handle.is_cached_reader() && handle.has_read_transaction());
let completion_preserving_writer_transaction = handle
.as_ref()
.is_some_and(|handle| !handle.is_cached_reader() && !handle.conn.is_autocommit());
let mut operation_slot = if active_read_transaction {
None
} else if completion_preserving_writer_transaction {
Some(acquire_reader_handle_slot(&pool, operation, SlotTimeoutClass::ReaderContract).await?)
} else {
Some(
crate::await_request_read_phase(
operation,
acquire_reader_handle_slot(&pool, operation, SlotTimeoutClass::ReaderContract),
)
.await??,
)
};
let Some(owned_handle) = handle.take() else {
return Err(StorageError::Pool {
operation: operation.into(),
message: "connection already consumed".into(),
});
};
let origin = pool.origin();
let read_tx_max_age = pool.config().read_tx_max_age;
let (owned_handle, result) = crate::read_cancellation::run_interruptible_read(
StorageCapability::Sql,
operation,
move |scope| {
let mut owned_handle = owned_handle;
let cached_reader = owned_handle.is_cached_reader();
let entered_with_transaction = owned_handle.has_read_transaction();
let entered_autocommit = owned_handle.conn.is_autocommit();
let mut restore_handle = true;
let mut result = if cached_reader && entered_with_transaction && entered_autocommit {
drop(owned_handle.read_transaction_slot.take());
Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message: "cached read-only handle retained transaction admission after SQLite \
had already returned to autocommit; the stale permit was released"
.into(),
})
} else if cached_reader && !entered_with_transaction && !entered_autocommit {
Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message: "cached read-only handle entered the operation outside autocommit; \
its transaction was rolled back before releasing the reader permit"
.into(),
})
} else if cached_reader
&& entered_with_transaction
&& owned_handle
.read_transaction_slot
.as_ref()
.is_some_and(|tx| tx.opened_at.elapsed() >= read_tx_max_age)
{
crate::checkpoint::note_read_tx_max_age_eviction();
match owned_handle.conn.execute_batch("ROLLBACK") {
Ok(()) if owned_handle.conn.is_autocommit() => {
drop(owned_handle.read_transaction_slot.take());
Err(StorageError::ReadTransactionAgeEvicted {
operation: operation.into(),
max_age_secs: read_tx_max_age.as_secs(),
})
}
Ok(()) => {
restore_handle = false;
Err(StorageError::ReadTransactionAgeEvictionCleanupFailed {
operation: operation.into(),
max_age_secs: read_tx_max_age.as_secs(),
message: "rollback did not restore autocommit".into(),
})
}
Err(error) => {
restore_handle = false;
Err(StorageError::ReadTransactionAgeEvictionCleanupFailed {
operation: operation.into(),
max_age_secs: read_tx_max_age.as_secs(),
message: format!("rollback failed: {error}"),
})
}
}
} else if cached_reader && entered_with_transaction {
match transaction_control {
None | Some(CachedReadTransactionControl::Finish(_)) => {
read(scope, &owned_handle.conn, true, true)
}
Some(CachedReadTransactionControl::BeginDeferred) => {
Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message: "cached read-only handle already owns an admitted read \
transaction; nested BEGIN is not supported"
.into(),
})
}
Some(CachedReadTransactionControl::Unsupported(keyword)) => {
Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message: format!(
"cached read-only transaction does not support nested or \
write-locking transaction control ({keyword})"
),
})
}
}
} else if cached_reader {
match transaction_control {
None | Some(CachedReadTransactionControl::BeginDeferred) => {
read(scope, &owned_handle.conn, false, true)
}
Some(CachedReadTransactionControl::Finish(keyword))
| Some(CachedReadTransactionControl::Unsupported(keyword)) => {
Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message: format!(
"cached read-only handle has no admitted transaction for \
transaction control ({keyword})"
),
})
}
}
} else {
read(scope, &owned_handle.conn, false, entered_autocommit)
};
if scope.cleanup_failed() {
restore_handle = false;
}
if cached_reader
&& matches!(result, Err(StorageError::Timeout { .. }))
&& !owned_handle.conn.is_autocommit()
{
match owned_handle.conn.execute_batch("ROLLBACK") {
Ok(()) if owned_handle.conn.is_autocommit() => {
drop(owned_handle.read_transaction_slot.take());
}
Ok(()) => {
restore_handle = false;
result = Err(StorageError::Transaction {
operation: operation.into(),
message:
"interrupted read transaction rollback did not restore autocommit; \
the connection was discarded"
.into(),
});
}
Err(error) => {
restore_handle = false;
result = Err(StorageError::Transaction {
operation: operation.into(),
message: format!(
"failed to roll back interrupted read transaction ({error}); \
the connection was discarded"
),
});
}
}
}
if cached_reader && entered_with_transaction {
if owned_handle.conn.is_autocommit() {
drop(owned_handle.read_transaction_slot.take());
if result.is_ok()
&& !matches!(
transaction_control,
Some(CachedReadTransactionControl::Finish(_))
)
{
result = Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message: "cached read-only operation unexpectedly ended its admitted \
transaction; reader admission was released after autocommit"
.into(),
});
}
} else if result.is_ok()
&& matches!(
transaction_control,
Some(CachedReadTransactionControl::Finish(_))
)
{
result = Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message: "transaction-ending control completed but the cached reader \
remained outside autocommit; its reader permit remains retained"
.into(),
});
}
} else if cached_reader
&& entered_autocommit
&& matches!(
transaction_control,
Some(CachedReadTransactionControl::BeginDeferred)
)
&& result.is_ok()
{
if owned_handle.conn.is_autocommit() {
result = Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message: "deferred BEGIN completed without opening a read transaction"
.into(),
});
} else {
match operation_slot.take() {
Some(slot) => {
let tx_handle = khive_storage::tx_registry::register_scoped(
Some(CACHED_READ_TRANSACTION_LABEL.to_string()),
origin.clone(),
);
owned_handle.read_transaction_slot = Some(CachedReadTransaction {
_slot: slot,
_tx_handle: tx_handle,
opened_at: Instant::now(),
});
}
None => {
result = Err(StorageError::Pool {
operation: operation.into(),
message: "successful cached-reader BEGIN had no operation permit; \
its transaction was rolled back before returning"
.into(),
});
}
}
}
}
if cached_reader
&& owned_handle.read_transaction_slot.is_none()
&& !owned_handle.conn.is_autocommit()
{
match owned_handle.conn.execute_batch("ROLLBACK") {
Ok(()) if owned_handle.conn.is_autocommit() => {
if result.is_ok() {
result = Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message: "cached read-only operation left the connection outside \
autocommit; its transaction was rolled back before \
releasing the reader permit"
.into(),
});
}
}
Ok(()) => {
restore_handle = false;
result = Err(StorageError::Transaction {
operation: operation.into(),
message: "ROLLBACK completed but the cached reader remained outside \
autocommit; the connection was discarded before releasing \
the reader permit"
.into(),
});
}
Err(error) => {
restore_handle = false;
result = Err(StorageError::Transaction {
operation: operation.into(),
message: format!(
"failed to roll back a cached reader outside autocommit ({error}); \
the connection was discarded before releasing the reader permit"
),
});
}
}
}
let owned_handle = if restore_handle {
Some(owned_handle)
} else {
drop(owned_handle);
None
};
drop(operation_slot);
Ok((owned_handle, result))
},
)
.await?;
*handle = owned_handle;
result
}
#[async_trait]
impl khive_storage::SqlReader for SqliteReader {
async fn query_row(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Option<SqlRow>> {
let transaction_control = cached_read_transaction_control(&statement.sql);
admit_reader_capability_sql(&statement, transaction_control, "query_row")?;
if !self
.use_explicit_transaction_handle(transaction_control, "query_row")
.await?
{
return run_pool_reader_query(
Arc::clone(&self.pool),
"query_row",
move |scope, conn| {
execute_query_row_interruptibly(
scope,
conn,
&statement,
"query_row",
false,
true,
)
},
)
.await;
}
let result = execute_standalone_read(
&mut self.handle,
Arc::clone(&self.pool),
"query_row",
transaction_control,
move |scope, conn, rollback, interruptible| {
execute_query_row_interruptibly(
scope,
conn,
&statement,
"query_row",
rollback,
interruptible,
)
},
)
.await;
if self.handle.is_none() {
self.poisoned = true;
}
self.close_inactive_transaction_handle();
result
}
async fn query_all(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
let transaction_control = cached_read_transaction_control(&statement.sql);
admit_reader_capability_sql(&statement, transaction_control, "query_all")?;
if !self
.use_explicit_transaction_handle(transaction_control, "query_all")
.await?
{
return run_pool_reader_query(
Arc::clone(&self.pool),
"query_all",
move |scope, conn| {
execute_query_interruptibly(scope, conn, &statement, "query_all", false, true)
},
)
.await;
}
let result = execute_standalone_read(
&mut self.handle,
Arc::clone(&self.pool),
"query_all",
transaction_control,
move |scope, conn, rollback, interruptible| {
execute_query_interruptibly(
scope,
conn,
&statement,
"query_all",
rollback,
interruptible,
)
},
)
.await;
if self.handle.is_none() {
self.poisoned = true;
}
self.close_inactive_transaction_handle();
result
}
async fn query_page(
&mut self,
statement: SqlStatement,
page: PageRequest,
) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
let transaction_control = cached_read_transaction_control(&statement.sql);
admit_reader_capability_sql(&statement, transaction_control, "query_page")?;
if !self
.use_explicit_transaction_handle(transaction_control, "query_page")
.await?
{
return run_pool_reader_query(
Arc::clone(&self.pool),
"query_page",
move |scope, conn| {
execute_query_page_interruptibly(
scope,
conn,
&statement,
&page,
"query_page",
false,
true,
)
},
)
.await;
}
let result = execute_standalone_read(
&mut self.handle,
Arc::clone(&self.pool),
"query_page",
transaction_control,
move |scope, conn, rollback, interruptible| {
execute_query_page_interruptibly(
scope,
conn,
&statement,
&page,
"query_page",
rollback,
interruptible,
)
},
)
.await;
if self.handle.is_none() {
self.poisoned = true;
}
self.close_inactive_transaction_handle();
result
}
async fn query_scalar(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Option<SqlValue>> {
let row = self.query_row(statement).await?;
Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
}
async fn explain(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
let explain_stmt = SqlStatement {
sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
params: statement.params,
label: statement.label,
};
self.query_all(explain_stmt).await
}
}
struct SqliteWriter {
observe_direct_errors: bool,
handle: Option<StandaloneHandle>,
writer_task: Option<crate::writer_task::WriterTaskHandle>,
origin: khive_storage::tx_registry::TxOrigin,
db: String,
pool: Arc<ConnectionPool>,
event_rows: Option<Arc<AtomicEventRows>>,
held_lease: Option<crate::disk_guard::DetachedVolumeLease>,
}
fn execute_top_level_maintenance(
pool: &ConnectionPool,
conn: &rusqlite::Connection,
maintenance: TopLevelMaintenance,
) -> rusqlite::Result<()> {
match maintenance {
TopLevelMaintenance::WalCheckpointTruncate => {
let result = conn.query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |row| {
Ok((
row.get::<_, i64>(0)?,
row.get::<_, i64>(1)?,
row.get::<_, i64>(2)?,
))
});
crate::checkpoint::record_checkpoint_run_result(pool, result.as_ref().ok().copied());
result.map(|_| ())
}
TopLevelMaintenance::Vacuum => conn.execute_batch(maintenance.as_sql()),
}
}
impl SqliteWriter {
fn require_standalone_handle(
&self,
operation: &'static str,
) -> khive_storage::types::StorageResult<()> {
if self.handle.is_none() {
return Err(StorageError::Pool {
operation: operation.into(),
message: "connection already consumed".into(),
});
}
Ok(())
}
async fn use_queue_read_transaction_handle(
&mut self,
transaction_control: Option<CachedReadTransactionControl>,
operation: &'static str,
) -> khive_storage::types::StorageResult<bool> {
if self.handle.is_some() {
return Ok(true);
}
match transaction_control {
None => Ok(false),
Some(CachedReadTransactionControl::BeginDeferred) => {
self.handle =
Some(open_explicit_read_transaction_handle(Arc::clone(&self.pool)).await?);
Ok(true)
}
Some(CachedReadTransactionControl::Finish(keyword))
| Some(CachedReadTransactionControl::Unsupported(keyword)) => {
Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message: format!(
"cached read-only handle has no admitted transaction for transaction \
control ({keyword})"
),
})
}
}
}
fn close_inactive_queue_read_transaction_handle(&mut self) {
if self
.handle
.as_ref()
.is_some_and(|handle| handle.is_cached_reader() && !handle.has_read_transaction())
{
drop(self.handle.take());
}
}
}
#[async_trait]
impl khive_storage::SqlReader for SqliteWriter {
async fn query_row(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Option<SqlRow>> {
if self.writer_task.is_some() {
let transaction_control = cached_read_transaction_control(&statement.sql);
if !self
.use_queue_read_transaction_handle(transaction_control, "writer.query_row")
.await?
{
admit_reader_capability_sql(&statement, transaction_control, "writer.query_row")?;
return run_pool_reader_query(
Arc::clone(&self.pool),
"writer.query_row",
move |scope, conn| {
execute_query_row_interruptibly(
scope,
conn,
&statement,
"writer.query_row",
false,
true,
)
},
)
.await;
}
let result = execute_standalone_read(
&mut self.handle,
Arc::clone(&self.pool),
"writer.query_row",
transaction_control,
move |scope, conn, rollback, interruptible| {
execute_query_row_interruptibly(
scope,
conn,
&statement,
"writer.query_row",
rollback,
interruptible,
)
},
)
.await;
self.close_inactive_queue_read_transaction_handle();
return result;
}
let transaction_control = cached_read_transaction_control(&statement.sql);
execute_standalone_read(
&mut self.handle,
Arc::clone(&self.pool),
"writer.query_row",
transaction_control,
move |scope, conn, rollback, interruptible| {
execute_query_row_interruptibly(
scope,
conn,
&statement,
"writer.query_row",
rollback,
interruptible,
)
},
)
.await
}
async fn query_all(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
if self.writer_task.is_some() {
let transaction_control = cached_read_transaction_control(&statement.sql);
if !self
.use_queue_read_transaction_handle(transaction_control, "writer.query_all")
.await?
{
admit_reader_capability_sql(&statement, transaction_control, "writer.query_all")?;
return run_pool_reader_query(
Arc::clone(&self.pool),
"writer.query_all",
move |scope, conn| {
execute_query_interruptibly(
scope,
conn,
&statement,
"writer.query_all",
false,
true,
)
},
)
.await;
}
let result = execute_standalone_read(
&mut self.handle,
Arc::clone(&self.pool),
"writer.query_all",
transaction_control,
move |scope, conn, rollback, interruptible| {
execute_query_interruptibly(
scope,
conn,
&statement,
"writer.query_all",
rollback,
interruptible,
)
},
)
.await;
self.close_inactive_queue_read_transaction_handle();
return result;
}
let transaction_control = cached_read_transaction_control(&statement.sql);
execute_standalone_read(
&mut self.handle,
Arc::clone(&self.pool),
"writer.query_all",
transaction_control,
move |scope, conn, rollback, interruptible| {
execute_query_interruptibly(
scope,
conn,
&statement,
"writer.query_all",
rollback,
interruptible,
)
},
)
.await
}
async fn query_page(
&mut self,
statement: SqlStatement,
page: PageRequest,
) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
if self.writer_task.is_some() {
let transaction_control = cached_read_transaction_control(&statement.sql);
if !self
.use_queue_read_transaction_handle(transaction_control, "writer.query_page")
.await?
{
admit_reader_capability_sql(&statement, transaction_control, "writer.query_page")?;
return run_pool_reader_query(
Arc::clone(&self.pool),
"writer.query_page",
move |scope, conn| {
execute_query_page_interruptibly(
scope,
conn,
&statement,
&page,
"writer.query_page",
false,
true,
)
},
)
.await;
}
let result = execute_standalone_read(
&mut self.handle,
Arc::clone(&self.pool),
"writer.query_page",
transaction_control,
move |scope, conn, rollback, interruptible| {
execute_query_page_interruptibly(
scope,
conn,
&statement,
&page,
"writer.query_page",
rollback,
interruptible,
)
},
)
.await;
self.close_inactive_queue_read_transaction_handle();
return result;
}
let transaction_control = cached_read_transaction_control(&statement.sql);
execute_standalone_read(
&mut self.handle,
Arc::clone(&self.pool),
"writer.query_page",
transaction_control,
move |scope, conn, rollback, interruptible| {
execute_query_page_interruptibly(
scope,
conn,
&statement,
&page,
"writer.query_page",
rollback,
interruptible,
)
},
)
.await
}
async fn query_scalar(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Option<SqlValue>> {
let row = khive_storage::SqlReader::query_row(self, statement).await?;
Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
}
async fn explain(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
let explain_stmt = SqlStatement {
sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
params: statement.params,
label: statement.label,
};
khive_storage::SqlReader::query_all(self, explain_stmt).await
}
}
#[async_trait]
impl khive_storage::SqlWriter for SqliteWriter {
async fn execute(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<u64> {
if let Some(writer_task) = self.writer_task.clone() {
let event_rows = self.event_rows.clone();
return writer_task
.send_bounded(move |conn| {
let mut stmt = prepare_cached_sql_statement(conn, &statement.sql)
.map_err(|e| map_rusqlite_err(e, "execute"))?;
bind_params(&mut stmt, &statement.params)
.map_err(|e| map_rusqlite_err(e, "execute"))?;
let affected = stmt
.raw_execute()
.map_err(|e| map_rusqlite_err(e, "execute"))?;
if let Some(event_rows) = event_rows.as_deref() {
event_rows.observe(&statement, affected as u64);
}
Ok(affected as u64)
})
.await;
}
self.require_standalone_handle("execute")?;
let unit_holds_lease = self.held_lease.is_some();
if !unit_holds_lease {
if let Some(keyword) = transaction_control_head(&statement.sql) {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: "execute".into(),
message: format!(
"statement is transaction control ({keyword}); a standalone \
statement holds the volume lease only for the call — use \
atomic_unit to run statements as one transaction"
),
});
}
}
let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
operation: "execute".into(),
message: "connection already consumed".into(),
})?;
let event_rows = self.event_rows.clone();
let pool = Arc::clone(&self.pool);
let (handle, result) = tokio::task::spawn_blocking(move || {
run_standalone_statement(handle, pool, unit_holds_lease, statement, event_rows)
})
.await
.map_err(|e| StorageError::driver(StorageCapability::Sql, "execute", e))?;
self.handle = Some(handle);
let affected = result.map_err(|failure| match failure {
StandaloneWriteError::Refused(error) => error,
StandaloneWriteError::Sql(error) => self.map_direct_error(error, "execute"),
})?;
Ok(affected as u64)
}
async fn execute_batch(
&mut self,
statements: Vec<SqlStatement>,
) -> khive_storage::types::StorageResult<u64> {
reject_transaction_control_statements(&statements, "execute_batch")?;
if let Some(writer_task) = self.writer_task.clone() {
let event_rows = self.event_rows.clone();
return writer_task
.send_bounded(move |conn| {
let prepared = prepare_batch_statements(conn, &statements)
.map_err(|e| map_rusqlite_err(e, "execute_batch"))?;
execute_prepared_batch(conn, prepared, &statements, event_rows.as_deref())
.map_err(|e| map_rusqlite_err(e, "execute_batch"))
})
.await;
}
self.require_standalone_handle("execute_batch")?;
let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
operation: "execute_batch".into(),
message: "connection already consumed".into(),
})?;
let origin = self.origin.clone();
let event_rows = self.event_rows.clone();
let pool = Arc::clone(&self.pool);
let unit_holds_lease = self.held_lease.is_some();
let (handle, result) = tokio::task::spawn_blocking(move || {
run_standalone_batch(
handle,
pool,
unit_holds_lease,
statements,
origin,
event_rows,
)
})
.await
.map_err(|e| StorageError::driver(StorageCapability::Sql, "execute_batch", e))?;
self.handle = handle;
result.map_err(|failure| match failure {
StandaloneWriteError::Refused(error) => error,
StandaloneWriteError::Sql(failure) => self.map_direct_batch_failure(failure),
})
}
async fn execute_script(&mut self, script: String) -> khive_storage::types::StorageResult<()> {
if let Some(writer_task) = self.writer_task.clone() {
return writer_task
.send_bounded(move |conn| {
conn.execute_batch(&script)
.map_err(|e| map_rusqlite_err(e, "execute_script"))
})
.await;
}
self.require_standalone_handle("execute_script")?;
let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
operation: "execute_script".into(),
message: "connection already consumed".into(),
})?;
let pool = Arc::clone(&self.pool);
let unit_holds_lease = self.held_lease.is_some();
let (handle, result) = tokio::task::spawn_blocking(move || {
run_standalone_script(handle, pool, unit_holds_lease, script)
})
.await
.map_err(|e| StorageError::driver(StorageCapability::Sql, "execute_script", e))?;
self.handle = handle;
result.map_err(|failure| match failure {
StandaloneWriteError::Refused(error) => error,
StandaloneWriteError::Sql(error) => self.map_direct_error(error, "execute_script"),
})
}
async fn execute_script_top_level(
&mut self,
maintenance: TopLevelMaintenance,
) -> khive_storage::types::StorageResult<()> {
if let Some(writer_task) = self.writer_task.clone() {
let pool = Arc::clone(&self.pool);
let execute = move |conn: &rusqlite::Connection| {
execute_top_level_maintenance(&pool, conn, maintenance)
.map_err(|e| map_rusqlite_err(e, "execute_script_top_level"))
};
return if maintenance == TopLevelMaintenance::WalCheckpointTruncate {
writer_task.send_checkpoint_bounded(execute).await
} else {
writer_task.send_vacuum_bounded(execute).await
};
}
let handle = self.handle.take().ok_or_else(|| StorageError::Pool {
operation: "execute_script_top_level".into(),
message: "connection already consumed".into(),
})?;
let pool = Arc::clone(&self.pool);
let unit_holds_lease = self.held_lease.is_some();
let (handle, result) = tokio::task::spawn_blocking(move || {
run_standalone_top_level(handle, pool, unit_holds_lease, maintenance)
})
.await
.map_err(|e| StorageError::driver(StorageCapability::Sql, "execute_script_top_level", e))?;
self.handle = Some(handle);
result.map_err(|failure| match failure {
StandaloneWriteError::Refused(error) => error,
StandaloneWriteError::Sql(error) => {
self.map_direct_error(error, "execute_script_top_level")
}
})
}
}
async fn run_pool_reader_query<T, F>(
pool: Arc<ConnectionPool>,
operation: &'static str,
query: F,
) -> khive_storage::types::StorageResult<T>
where
T: Send + 'static,
F: FnOnce(
&crate::read_cancellation::InterruptibleReadScope,
&rusqlite::Connection,
) -> khive_storage::types::StorageResult<T>
+ Send
+ 'static,
{
let admission = pool
.acquire_reader_admission(StorageCapability::Sql, operation)
.await?;
crate::read_cancellation::run_interruptible_read(
StorageCapability::Sql,
operation,
move |scope| {
let mut guard = pool.resolve_reader_checkout(
StorageCapability::Sql,
operation,
pool.reader_with_admission(admission, || scope.should_stop()),
)?;
guard.mark_dirty();
let result = scope.with_pooled_reader(&mut guard, |conn| query(scope, conn));
if let Err(error) = &result {
pool.record_reader_query_error(error);
}
result
},
)
.await
}
async fn run_pool_writer_query<T, F>(
pool: Arc<ConnectionPool>,
operation: &'static str,
query: F,
) -> khive_storage::types::StorageResult<T>
where
T: Send + 'static,
F: FnOnce(
&crate::read_cancellation::InterruptibleReadScope,
&rusqlite::Connection,
bool,
) -> khive_storage::types::StorageResult<T>
+ Send
+ 'static,
{
crate::read_cancellation::run_interruptible_read(
StorageCapability::Sql,
operation,
move |scope| {
let guard = pool.try_writer().map_err(|error: SqliteError| {
error.into_storage_error(StorageCapability::Sql, operation)
})?;
scope.with_pooled_writer(&pool, &guard, |conn| {
let interruptible = conn.is_autocommit();
query(scope, conn, interruptible)
})
},
)
.await
}
struct PoolBackedReader {
pool: Arc<ConnectionPool>,
transaction: Option<SharedReaderTransactionGuard>,
}
fn finish_pool_backed_reader_step<T>(
transaction: &mut Option<SharedReaderTransactionGuard>,
guard: SharedReaderTransactionGuard,
expect_open_after: bool,
operation: &'static str,
result: khive_storage::types::StorageResult<T>,
) -> khive_storage::types::StorageResult<T> {
let still_open = !guard.conn().is_autocommit();
if still_open == expect_open_after {
if still_open {
*transaction = Some(guard);
}
return result;
}
guard.poison();
let message = if expect_open_after {
"a read inside the pool-backed reader's admitted transaction unexpectedly ended it; \
the connection was discarded"
} else {
"transaction-ending control completed but the pool-backed reader's connection \
remained outside autocommit; the connection was discarded"
};
match result {
Err(error) => Err(error),
Ok(_) => Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message: message.into(),
}),
}
}
async fn open_pool_backed_reader_transaction<T, F>(
transaction: &mut Option<SharedReaderTransactionGuard>,
pool: Arc<ConnectionPool>,
operation: &'static str,
query: F,
) -> khive_storage::types::StorageResult<T>
where
T: Send + 'static,
F: FnOnce(
&crate::read_cancellation::InterruptibleReadScope,
&rusqlite::Connection,
bool,
bool,
) -> khive_storage::types::StorageResult<T>
+ Send
+ 'static,
{
let (guard, result) = crate::read_cancellation::run_interruptible_read(
StorageCapability::Sql,
operation,
move |scope| {
let Some(guard) = pool
.checkout_shared_reader_transaction(|| scope.should_stop())
.map_err(|error| StorageError::driver(StorageCapability::Sql, operation, error))?
else {
return Err(StorageError::Timeout {
operation: operation.into(),
});
};
let result = query(scope, guard.conn(), false, true);
if scope.cleanup_failed() {
guard.poison();
}
Ok((guard, result))
},
)
.await?;
finish_pool_backed_reader_step(transaction, guard, true, operation, result)
}
#[allow(clippy::too_many_lines)]
async fn run_pool_backed_reader_query<T, F>(
transaction: &mut Option<SharedReaderTransactionGuard>,
pool: Arc<ConnectionPool>,
operation: &'static str,
transaction_control: Option<CachedReadTransactionControl>,
query: F,
) -> khive_storage::types::StorageResult<T>
where
T: Send + 'static,
F: FnOnce(
&crate::read_cancellation::InterruptibleReadScope,
&rusqlite::Connection,
bool,
bool,
) -> khive_storage::types::StorageResult<T>
+ Send
+ 'static,
{
if transaction.is_none() {
return match transaction_control {
None => {
run_pool_reader_query(pool, operation, move |scope, conn| {
query(scope, conn, false, true)
})
.await
}
Some(CachedReadTransactionControl::Finish(keyword))
| Some(CachedReadTransactionControl::Unsupported(keyword)) => {
Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message: format!(
"pool-backed reader has no admitted transaction for transaction \
control ({keyword})"
),
})
}
Some(CachedReadTransactionControl::BeginDeferred) => {
open_pool_backed_reader_transaction(transaction, pool, operation, query).await
}
};
}
match transaction_control {
Some(CachedReadTransactionControl::BeginDeferred) => {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message: "pool-backed reader already owns an admitted read transaction; \
nested BEGIN is not supported"
.into(),
});
}
Some(CachedReadTransactionControl::Unsupported(keyword)) => {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: operation.into(),
message: format!(
"pool-backed reader's admitted read transaction does not support nested \
or write-locking transaction control ({keyword})"
),
});
}
None | Some(CachedReadTransactionControl::Finish(_)) => {}
}
let expect_open_after = transaction_control.is_none();
let guard = transaction.take().expect("checked Some above");
let (guard, result) = crate::read_cancellation::run_interruptible_read(
StorageCapability::Sql,
operation,
move |scope| {
let result = query(scope, guard.conn(), true, true);
if scope.cleanup_failed() {
guard.poison();
}
Ok((guard, result))
},
)
.await?;
finish_pool_backed_reader_step(transaction, guard, expect_open_after, operation, result)
}
#[async_trait]
impl khive_storage::SqlReader for PoolBackedReader {
async fn query_row(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Option<SqlRow>> {
let transaction_control = cached_read_transaction_control(&statement.sql);
admit_reader_capability_sql(&statement, transaction_control, "pool_reader.query_row")?;
let pool = Arc::clone(&self.pool);
run_pool_backed_reader_query(
&mut self.transaction,
pool,
"pool_reader.query_row",
transaction_control,
move |scope, conn, rollback, interruptible| {
execute_query_row_interruptibly(
scope,
conn,
&statement,
"pool_reader.query_row",
rollback,
interruptible,
)
},
)
.await
}
async fn query_all(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
let transaction_control = cached_read_transaction_control(&statement.sql);
admit_reader_capability_sql(&statement, transaction_control, "pool_reader.query_all")?;
let pool = Arc::clone(&self.pool);
run_pool_backed_reader_query(
&mut self.transaction,
pool,
"pool_reader.query_all",
transaction_control,
move |scope, conn, rollback, interruptible| {
execute_query_interruptibly(
scope,
conn,
&statement,
"pool_reader.query_all",
rollback,
interruptible,
)
},
)
.await
}
async fn query_page(
&mut self,
statement: SqlStatement,
page: PageRequest,
) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
let transaction_control = cached_read_transaction_control(&statement.sql);
admit_reader_capability_sql(&statement, transaction_control, "pool_reader.query_page")?;
let pool = Arc::clone(&self.pool);
run_pool_backed_reader_query(
&mut self.transaction,
pool,
"pool_reader.query_page",
transaction_control,
move |scope, conn, rollback, interruptible| {
execute_query_page_interruptibly(
scope,
conn,
&statement,
&page,
"pool_reader.query_page",
rollback,
interruptible,
)
},
)
.await
}
async fn query_scalar(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Option<SqlValue>> {
let row = self.query_row(statement).await?;
Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
}
async fn explain(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
let explain_stmt = SqlStatement {
sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
params: statement.params,
label: statement.label,
};
self.query_all(explain_stmt).await
}
}
struct PoolBackedWriter {
pool: Arc<ConnectionPool>,
}
#[async_trait]
impl khive_storage::SqlReader for PoolBackedWriter {
async fn query_row(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Option<SqlRow>> {
let pool = Arc::clone(&self.pool);
run_pool_writer_query(
pool,
"pool_writer.query_row",
move |scope, conn, interruptible| {
execute_query_row_interruptibly(
scope,
conn,
&statement,
"pool_writer.query_row",
false,
interruptible,
)
},
)
.await
}
async fn query_all(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
let pool = Arc::clone(&self.pool);
run_pool_writer_query(
pool,
"pool_writer.query_all",
move |scope, conn, interruptible| {
execute_query_interruptibly(
scope,
conn,
&statement,
"pool_writer.query_all",
false,
interruptible,
)
},
)
.await
}
async fn query_page(
&mut self,
statement: SqlStatement,
page: PageRequest,
) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
let pool = Arc::clone(&self.pool);
run_pool_writer_query(
pool,
"pool_writer.query_page",
move |scope, conn, interruptible| {
execute_query_page_interruptibly(
scope,
conn,
&statement,
&page,
"pool_writer.query_page",
false,
interruptible,
)
},
)
.await
}
async fn query_scalar(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Option<SqlValue>> {
let row = khive_storage::SqlReader::query_row(self, statement).await?;
Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
}
async fn explain(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
let explain_stmt = SqlStatement {
sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
params: statement.params,
label: statement.label,
};
khive_storage::SqlReader::query_all(self, explain_stmt).await
}
}
#[async_trait]
impl khive_storage::SqlWriter for PoolBackedWriter {
async fn execute(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<u64> {
if let Some(keyword) = transaction_control_head(&statement.sql) {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Sql,
operation: "pool_writer.execute".into(),
message: format!(
"statement is transaction control ({keyword}); a pooled writer is \
released after every call, so it cannot hold a transaction across \
calls — use atomic_unit to run statements as one transaction"
),
});
}
let pool = Arc::clone(&self.pool);
tokio::task::spawn_blocking(move || {
let guard = pool.try_writer().map_err(|e: SqliteError| {
StorageError::driver(StorageCapability::Sql, "pool_writer.execute", e)
})?;
let result = (|| {
let mut stmt = prepare_cached_sql_statement(&guard, &statement.sql)
.map_err(|e| map_rusqlite_err(e, "pool_writer.execute"))?;
bind_params(&mut stmt, &statement.params)
.map_err(|e| map_rusqlite_err(e, "pool_writer.execute"))?;
let rows = stmt
.raw_execute()
.map_err(|e| map_rusqlite_err(e, "pool_writer.execute"))?;
Ok(rows as u64)
})();
settle_pooled_call(&guard, "pool_writer.execute", result)
.inspect_err(|error| pool.record_direct_writer_error(error))
})
.await
.map_err(|e| StorageError::driver(StorageCapability::Sql, "pool_writer.execute", e))?
}
async fn execute_batch(
&mut self,
statements: Vec<SqlStatement>,
) -> khive_storage::types::StorageResult<u64> {
reject_transaction_control_statements(&statements, "pool_writer.execute_batch")?;
let pool = Arc::clone(&self.pool);
tokio::task::spawn_blocking(move || {
let guard = pool.try_writer().map_err(|e: SqliteError| {
StorageError::driver(StorageCapability::Sql, "pool_writer.execute_batch", e)
})?;
let result = (|| {
let prepared = prepare_batch_statements(&guard, &statements)
.map_err(|e| map_rusqlite_err(e, "pool_writer.execute_batch"))?;
guard
.execute_batch("BEGIN IMMEDIATE")
.map_err(|e| map_rusqlite_err(e, "pool_writer.execute_batch"))?;
let _tx_handle = khive_storage::tx_registry::register_scoped(
Some("pool_writer.execute_batch".to_string()),
pool.origin(),
);
let result = execute_prepared_batch(&guard, prepared, &statements, None)
.map_err(|e| map_rusqlite_err(e, "pool_writer.execute_batch"));
match result {
Ok(total) => {
if let Err(e) = guard.execute_batch("COMMIT") {
let _ = guard.execute_batch("ROLLBACK");
Err(map_rusqlite_err(e, "pool_writer.execute_batch"))
} else {
Ok(total)
}
}
Err(e) => {
let _ = guard.execute_batch("ROLLBACK");
Err(e)
}
}
})();
result.inspect_err(|error| pool.record_direct_writer_error(error))
})
.await
.map_err(|e| StorageError::driver(StorageCapability::Sql, "pool_writer.execute_batch", e))?
}
async fn execute_script(&mut self, script: String) -> khive_storage::types::StorageResult<()> {
let pool = Arc::clone(&self.pool);
tokio::task::spawn_blocking(move || {
let guard = pool.try_writer().map_err(|e: SqliteError| {
StorageError::driver(StorageCapability::Sql, "pool_writer.execute_script", e)
})?;
let result = guard
.execute_batch(&script)
.map_err(|e| map_rusqlite_err(e, "pool_writer.execute_script"));
settle_pooled_call(&guard, "pool_writer.execute_script", result)
.inspect_err(|error| pool.record_direct_writer_error(error))
})
.await
.map_err(|e| {
StorageError::driver(StorageCapability::Sql, "pool_writer.execute_script", e)
})?
}
}
struct InlineWriter {
event_rows: Option<Arc<AtomicEventRows>>,
conn: *const rusqlite::Connection,
}
unsafe impl Send for InlineWriter {}
impl InlineWriter {
fn conn(&self) -> &rusqlite::Connection {
unsafe { &*self.conn }
}
}
#[async_trait]
impl khive_storage::SqlReader for InlineWriter {
async fn query_row(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Option<SqlRow>> {
execute_query_row(self.conn(), &statement)
.map_err(|e| map_rusqlite_err(e, "inline.query_row"))
}
async fn query_all(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
execute_query(self.conn(), &statement).map_err(|e| map_rusqlite_err(e, "inline.query_all"))
}
async fn query_page(
&mut self,
statement: SqlStatement,
page: PageRequest,
) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
execute_query_page(self.conn(), &statement, &page)
.map_err(|e| map_rusqlite_err(e, "inline.query_page"))
}
async fn query_scalar(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Option<SqlValue>> {
let row = khive_storage::SqlReader::query_row(self, statement).await?;
Ok(row.and_then(|r| r.columns.into_iter().next().map(|c| c.value)))
}
async fn explain(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<Vec<SqlRow>> {
let explain_stmt = SqlStatement {
sql: format!("EXPLAIN QUERY PLAN {}", statement.sql),
params: statement.params,
label: statement.label,
};
khive_storage::SqlReader::query_all(self, explain_stmt).await
}
}
#[async_trait]
impl khive_storage::SqlWriter for InlineWriter {
async fn execute(
&mut self,
statement: SqlStatement,
) -> khive_storage::types::StorageResult<u64> {
let mut stmt = prepare_cached_sql_statement(self.conn(), &statement.sql)
.map_err(|e| map_rusqlite_err(e, "inline.execute"))?;
bind_params(&mut stmt, &statement.params)
.map_err(|e| map_rusqlite_err(e, "inline.execute"))?;
let affected = stmt
.raw_execute()
.map_err(|e| map_rusqlite_err(e, "inline.execute"))?;
if let Some(event_rows) = self.event_rows.as_deref() {
event_rows.observe(&statement, affected as u64);
}
Ok(affected as u64)
}
async fn execute_batch(
&mut self,
statements: Vec<SqlStatement>,
) -> khive_storage::types::StorageResult<u64> {
reject_transaction_control_statements(&statements, "inline.execute_batch")?;
let prepared = prepare_batch_statements(self.conn(), &statements)
.map_err(|e| map_rusqlite_err(e, "inline.execute_batch"))?;
execute_prepared_batch(
self.conn(),
prepared,
&statements,
self.event_rows.as_deref(),
)
.map_err(|e| map_rusqlite_err(e, "inline.execute_batch"))
}
async fn execute_script(&mut self, script: String) -> khive_storage::types::StorageResult<()> {
self.conn()
.execute_batch(&script)
.map_err(|e| map_rusqlite_err(e, "inline.execute_script"))
}
}
fn block_on_sync<F: std::future::Future>(fut: F) -> Result<F::Output, StorageError> {
use std::task::{Context, Poll, RawWaker, RawWakerVTable, Waker};
fn no_op(_: *const ()) {}
fn clone_waker(_: *const ()) -> RawWaker {
RawWaker::new(std::ptr::null(), &VTABLE)
}
static VTABLE: RawWakerVTable = RawWakerVTable::new(clone_waker, no_op, no_op, no_op);
let raw_waker = RawWaker::new(std::ptr::null(), &VTABLE);
let waker = unsafe { Waker::from_raw(raw_waker) };
let mut cx = Context::from_waker(&waker);
let mut fut = std::pin::pin!(fut);
match fut.as_mut().poll(&mut cx) {
Poll::Ready(v) => Ok(v),
Poll::Pending => {
tracing::error!(
"block_on_sync: atomic_unit future suspended on its first poll — \
the closure passed to SqlAccess::atomic_unit must be non-blocking \
(synchronous InlineWriter calls only, no real .await point)"
);
Err(StorageError::Internal(
"atomic_unit future suspended — closure must be non-blocking".to_string(),
))
}
}
}
pub struct SqlBridge {
pool: Arc<ConnectionPool>,
is_file_backed: bool,
}
impl SqlBridge {
pub fn new(pool: Arc<ConnectionPool>, _is_file_backed: bool) -> Self {
let is_file_backed = pool.canonical_path().is_some();
Self {
pool,
is_file_backed,
}
}
}
#[async_trait]
impl khive_storage::SqlAccess for SqlBridge {
fn database_path(&self) -> Option<std::path::PathBuf> {
self.pool.canonical_path().map(std::path::Path::to_path_buf)
}
async fn reader(
&self,
) -> khive_storage::types::StorageResult<Box<dyn khive_storage::SqlReader>> {
if self.is_file_backed {
Ok(Box::new(SqliteReader {
handle: None,
pool: Arc::clone(&self.pool),
poisoned: false,
}))
} else {
Ok(Box::new(PoolBackedReader {
pool: Arc::clone(&self.pool),
transaction: None,
}))
}
}
async fn writer(
&self,
) -> khive_storage::types::StorageResult<Box<dyn khive_storage::SqlWriter>> {
if self.is_file_backed {
if self.pool.config().read_only {
return Err(StorageError::Pool {
operation: "writer".into(),
message: "backend is read-only".into(),
});
}
let db = crate::timeout_sink::db_label(&self.pool);
let writer_task = match self.pool.writer_task_handle() {
Ok(handle) => handle,
Err(e) => {
if self.pool.config().write_routing_strict {
return Err(e);
}
tracing::warn!(
error = %e,
"KHIVE_WRITE_ROUTING is not strict; writer() degrades to the \
standalone-connection path"
);
None
}
};
if writer_task.is_none() && self.pool.config().write_routing_strict {
return Err(StorageError::Pool {
operation: "writer".into(),
message: "KHIVE_WRITE_ROUTING=strict but no writer-task handle is \
available; refusing to fall back to a direct connection"
.into(),
});
}
if writer_task.is_none() && self.pool.write_queue_active() {
crate::timeout_sink::emit_direct_route_violation(
&db,
crate::timeout_sink::Site::DirectRouteSqlBridgeWriter,
);
}
let handle = if writer_task.is_none() {
let handle_slot = acquire_handle_slot(
self.pool.sql_bridge_writer_slots(),
self.pool.config().checkout_timeout,
"sql_bridge.writer_handle",
SlotTimeoutClass::Admission,
)
.await?;
let (conn, handle_slot) =
open_standalone_writer_on_blocking(Arc::clone(&self.pool), handle_slot).await?;
Some(StandaloneHandle {
conn,
_retained_slot: Some(handle_slot),
read_transaction_slot: None,
})
} else {
None
};
Ok(Box::new(SqliteWriter {
observe_direct_errors: true,
event_rows: None,
handle,
writer_task,
origin: self.pool.origin(),
db,
pool: Arc::clone(&self.pool),
held_lease: None,
}))
} else {
Ok(Box::new(PoolBackedWriter {
pool: Arc::clone(&self.pool),
}))
}
}
async fn atomic_unit(
&self,
op: AtomicUnitOp,
) -> khive_storage::types::StorageResult<Box<dyn Any + Send>> {
let event_rows = Arc::new(AtomicEventRows::default());
let result = async {
if self.is_file_backed {
if self.pool.config().read_only {
return Err(StorageError::Pool {
operation: "atomic_unit".into(),
message: "backend is read-only".into(),
});
}
let handle = self.pool.writer_task_handle()?;
if handle.is_none() && self.pool.config().write_routing_strict {
return Err(StorageError::Pool {
operation: "atomic_unit".into(),
message: "KHIVE_WRITE_ROUTING=strict but no writer-task handle is \
available; refusing to fall back to a direct connection"
.into(),
});
}
if handle.is_none() && self.pool.write_queue_active() {
crate::timeout_sink::emit_direct_route_violation(
&crate::timeout_sink::db_label(&self.pool),
crate::timeout_sink::Site::DirectRouteAtomicUnit,
);
}
if let Some(writer_task) = handle {
let pending_event_rows = Arc::clone(&event_rows);
return writer_task
.send_bounded(move |conn| {
let mut inline = InlineWriter {
event_rows: Some(Arc::clone(&pending_event_rows)),
conn: conn as *const rusqlite::Connection,
};
match block_on_sync(op(&mut inline)) {
Ok(inner) => inner,
Err(e) => Err(e),
}
})
.await;
}
let handle_slot = acquire_handle_slot(
self.pool.sql_bridge_writer_slots(),
self.pool.config().checkout_timeout,
"sql_bridge.atomic_unit_handle",
SlotTimeoutClass::Admission,
)
.await?;
let unit_lease = acquire_unit_lease(Arc::clone(&self.pool)).await?;
let (conn, handle_slot) =
open_standalone_writer_on_blocking(Arc::clone(&self.pool), handle_slot).await?;
let mut writer = SqliteWriter {
observe_direct_errors: false,
event_rows: Some(Arc::clone(&event_rows)),
handle: Some(StandaloneHandle {
conn,
_retained_slot: Some(handle_slot),
read_transaction_slot: None,
}),
writer_task: None,
origin: self.pool.origin(),
db: crate::timeout_sink::db_label(&self.pool),
pool: Arc::clone(&self.pool),
held_lease: unit_lease,
};
run_manual_atomic_unit(&mut writer, op, self.pool.origin())
.await
.inspect_err(|error| self.pool.record_direct_writer_error(error))
} else {
let pool = Arc::clone(&self.pool);
let pending_event_rows = Arc::clone(&event_rows);
tokio::task::spawn_blocking(move || {
let guard = pool.try_writer().map_err(|error: SqliteError| {
StorageError::driver(StorageCapability::Sql, "atomic_unit", error)
})?;
let conn = guard.conn();
if !conn.is_autocommit() {
pool.retire_pooled_writer(conn);
return Err(StorageError::writer_task_terminated(
khive_storage::WriterTaskRequestState::SideEffectsUnknown,
));
}
if let Err(error) = conn.execute_batch("BEGIN IMMEDIATE") {
if !conn.is_autocommit() {
pool.retire_pooled_writer(conn);
return Err(StorageError::writer_task_terminated(
khive_storage::WriterTaskRequestState::SideEffectsUnknown,
));
}
return Err(map_rusqlite_err(error, "atomic_unit.begin"))
.inspect_err(|error| pool.record_direct_writer_error(error));
}
let _tx_handle = khive_storage::tx_registry::register_scoped(
Some("atomic_unit".to_string()),
pool.origin(),
);
let (result, terminal_state) = crate::writer_task::execute_wrapped_transaction(
conn,
"atomic_unit.commit",
|conn| {
let mut inline = InlineWriter {
event_rows: Some(Arc::clone(&pending_event_rows)),
conn: conn as *const rusqlite::Connection,
};
block_on_sync(op(&mut inline)).and_then(|result| result)
},
);
if terminal_state.is_some() {
pool.retire_pooled_writer(conn);
}
result.inspect_err(|error| pool.record_direct_writer_error(error))
})
.await
.map_err(|error| {
StorageError::driver(StorageCapability::Sql, "atomic_unit", error)
})?
}
}
.await;
khive_storage::usage::account_event_write(
result.as_ref().map(|_| event_rows.committed_rows()),
);
result
}
}
#[cfg(test)]
#[path = "sql_bridge_tests.rs"]
mod tests;
#[cfg(test)]
#[path = "sql_bridge/direct_busy_tests.rs"]
mod direct_busy_tests;
#[cfg(test)]
#[path = "sql_bridge/settlement_hazard_tests.rs"]
mod settlement_hazard_tests;