use anyhow::Context;
use chrono::{DateTime, Utc};
use futures_util::FutureExt;
use std::future::Future;
use std::panic::AssertUnwindSafe;
use std::path::Path;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use tracing::{debug, error, info, warn};
use turso::Builder;
pub(crate) use turso::{IntoParams, Row, Value, params};
pub(crate) mod cdc;
pub mod checkpoint;
pub(crate) mod checkpoint_cause;
pub mod debug;
pub(crate) mod failure_record;
pub mod ipc;
pub(crate) mod migrations;
pub(crate) mod shrink_gate;
#[cfg(all(unix, test))]
mod store_lock_check;
pub mod wal_guard;
#[must_use]
pub fn now() -> String {
Utc::now().to_rfc3339()
}
pub(crate) fn parse_utc_timestamp(s: &str) -> Result<DateTime<Utc>, chrono::ParseError> {
DateTime::parse_from_rfc3339(s).map(|dt| dt.with_timezone(&Utc))
}
#[cfg_attr(
not(test),
expect(
dead_code,
reason = "Referenced only by assertion tests; kept for documentation"
)
)]
pub(crate) const EXPERIMENTAL_FEATURES: &[&str] = &["index_method"];
pub(crate) const CONSOLIDATED_DB_NAME: &str = "core";
const DOMAIN_STORE_NAMES: &[&str] = &[
"board",
"sessions",
"workspaces",
"users",
"config",
"chat_history",
"alarms",
];
pub(crate) const LOG_DB_NAME: &str = "logs";
pub(crate) static DOMAIN_CONN: tokio::sync::OnceCell<Connection> =
tokio::sync::OnceCell::const_new();
pub(crate) fn iter_checkpoint_stores()
-> impl Iterator<Item = (&'static str, Option<&'static crate::db::Connection>)> {
[
(CONSOLIDATED_DB_NAME, DOMAIN_CONN.get()),
(LOG_DB_NAME, crate::logs::LOG_STORE.get().map(|s| &s.conn)),
]
.into_iter()
}
pub(crate) fn store_names() -> Vec<&'static str> {
let mut names: Vec<&'static str> = DOMAIN_STORE_NAMES.to_vec();
names.push(LOG_DB_NAME);
names
}
pub(crate) fn debug_db_names() -> Vec<&'static str> {
let mut names = store_names();
names.push(CONSOLIDATED_DB_NAME);
names
}
pub async fn init_all_stores() -> anyhow::Result<()> {
if DOMAIN_CONN.get().is_some() {
return Ok(());
}
let root = crate::config::CONFIG.global_storage_root();
let conn = open_consolidated_store(&root).await?;
let board = crate::pipeline::board::BoardStore { conn: conn.clone() };
board.after_open(&root).await?;
let users = crate::users::UserStore { conn: conn.clone() };
users.ensure_admin_user().await?;
let session = crate::session::SessionStore { conn: conn.clone() };
let workspace = crate::workspace::WorkspaceStore { conn: conn.clone() };
let config = crate::config_db::ConfigStore { conn: conn.clone() };
let chat_history = crate::channels::chat_history::ChatHistoryStore { conn: conn.clone() };
let alarms = crate::alarms::AlarmStore { conn: conn.clone() };
let _ = DOMAIN_CONN.set(conn.clone());
crate::db::cdc::spawn_drainer(conn.clone());
init_cell(&crate::pipeline::board::BOARD, board, "BOARD")?;
init_cell(&crate::session::SESSIONS, session, "SESSIONS")?;
init_cell(&crate::workspace::WORKSPACES, workspace, "WORKSPACES")?;
init_cell(&crate::users::USER_STORE, users, "USER_STORE")?;
init_cell(&crate::config_db::CONFIG_STORE, config, "CONFIG_STORE")?;
init_cell(
&crate::channels::chat_history::CHAT_HISTORY,
chat_history,
"CHAT_HISTORY",
)?;
init_cell(&crate::alarms::ALARMS, alarms, "ALARMS")?;
Ok(())
}
fn init_cell<T>(cell: &tokio::sync::OnceCell<T>, value: T, name: &str) -> anyhow::Result<()> {
cell.set(value)
.map_err(|_| anyhow::anyhow!("{name} already initialized"))
}
#[must_use]
pub(crate) fn experimental_database_opts() -> turso::core::DatabaseOpts {
turso::core::DatabaseOpts::new().with_index_method(true)
}
static TANTIVY_SPECIAL: &[char] = &[
'+', '^', '~', ':', '{', '}', '"', '\'', '`', '[', ']', '(', ')', '\\', '*', '-', ];
#[must_use]
pub(crate) fn sanitize_fts_query(query: &str) -> String {
query
.split(|c: char| c.is_whitespace() || TANTIVY_SPECIAL.contains(&c))
.map(|word| word.trim_start_matches('/'))
.filter(|word| !word.is_empty())
.collect::<Vec<_>>()
.join(" ")
}
#[must_use]
pub(crate) fn sql_in_placeholders(count: usize) -> String {
vec!["?"; count].join(", ")
}
#[derive(Clone, Debug)]
pub(crate) struct Connection {
conn: Arc<tokio::sync::Mutex<turso::Connection>>,
has_dangling_tx: Arc<AtomicBool>,
db_path: Arc<std::path::Path>,
}
pub(crate) struct ReadonlyRows {
pub columns: Vec<String>,
pub rows: Vec<Vec<Value>>,
pub truncated: bool,
}
const KNOWN_FTS_DIR_COUNT_FALSE_POSITIVE: &str =
"wrong # of entries in index __turso_internal_fts_dir_";
const FTS_INTERNAL_INDEX_PREFIX: &str = "__turso_internal_fts_dir_";
pub(crate) const USER_OBJECT_FILTER: &str = "name NOT LIKE 'sqlite_%' AND name NOT LIKE '__turso_internal_%' \
AND name NOT IN ('turso_cdc', 'turso_cdc_version')";
fn remove_rebuild_temp(temp: &Path) {
let _ = std::fs::remove_file(temp);
let _ = std::fs::remove_file(wal_path(temp));
}
struct TempCleanup<'a>(&'a Path);
impl Drop for TempCleanup<'_> {
fn drop(&mut self) {
remove_rebuild_temp(self.0);
}
}
#[must_use]
fn known_fts_dir_false_positive(message: &str) -> bool {
message.contains(FTS_INTERNAL_INDEX_PREFIX)
}
fn map_rows<T, E>(
rows: &[Row],
mut map: impl FnMut(&Row) -> std::result::Result<T, E>,
) -> Vec<turso::Result<T>>
where
E: std::fmt::Display,
{
rows.iter()
.map(|row| map(row).map_err(|e| turso::Error::Error(e.to_string())))
.collect()
}
fn optional_row<T>(r: turso::Result<T>) -> anyhow::Result<Option<T>> {
match r {
Ok(val) => Ok(Some(val)),
Err(::turso::Error::QueryReturnedNoRows) => Ok(None),
Err(e) => Err(e.into()),
}
}
async fn map_first_row<T, E>(
mut rows: turso::Rows,
map: impl FnOnce(&Row) -> std::result::Result<T, E>,
) -> turso::Result<T>
where
E: std::fmt::Display,
{
let row = rows
.next()
.await?
.ok_or(turso::Error::QueryReturnedNoRows)?;
map(&row).map_err(|e| turso::Error::Error(e.to_string()))
}
fn strict_collect<T>(rows: Vec<turso::Result<T>>) -> anyhow::Result<Vec<T>> {
rows.into_iter()
.collect::<std::result::Result<Vec<_>, _>>()
.map_err(Into::into)
}
async fn drain_rows(mut rows: turso::Rows) -> turso::Result<Vec<Row>> {
let mut result = Vec::new();
while let Some(row) = rows.next().await? {
result.push(row);
}
Ok(result)
}
async fn run_readonly_query(
conn: &turso::Connection,
sql: &str,
params: impl IntoParams + Send + 'static,
row_limit: usize,
) -> turso::Result<ReadonlyRows> {
let mut rows = conn.query(sql, params).await?;
let columns = rows.column_names();
let col_count = columns.len();
let mut result_rows = Vec::new();
let mut count = 0usize;
let mut truncated = false;
while let Some(row) = rows.next().await? {
if count >= row_limit {
truncated = true;
break;
}
let mut values = Vec::with_capacity(col_count);
for idx in 0..col_count {
values.push(row.get_value(idx)?);
}
result_rows.push(values);
count += 1;
}
Ok(ReadonlyRows {
columns,
rows: result_rows,
truncated,
})
}
const OPEN_LOCK_RETRY_ATTEMPTS: usize = 10;
const OPEN_LOCK_RETRY_DELAY: Duration = Duration::from_millis(500);
const LOCK_STEP_PREFIX: &str = "Locking error: Failed locking";
const LOCK_CONTENTION_MARKER: &str = "File is locked by another process";
const WINDOWS_LOCK_VIOLATION_CODE: &str = "(os error 33)";
fn is_open_time_lock_error(err: &turso::Error) -> bool {
let turso::Error::Error(msg) = err else {
return false;
};
msg.contains(LOCK_STEP_PREFIX)
&& (msg.contains(LOCK_CONTENTION_MARKER) || msg.contains(WINDOWS_LOCK_VIOLATION_CODE))
}
pub(crate) fn is_store_lock_error(e: &anyhow::Error) -> bool {
e.chain().any(|cause| {
cause
.downcast_ref::<turso::Error>()
.is_some_and(is_open_time_lock_error)
})
}
#[derive(Debug)]
pub(crate) struct StoreRefusal {
pub store: &'static str,
pub db_path: std::path::PathBuf,
pub reason: String,
pub environment: bool,
}
impl StoreRefusal {
#[must_use]
pub(crate) fn new(store: &'static str, db_path: &Path, reason: impl Into<String>) -> Self {
Self {
store,
db_path: db_path.to_path_buf(),
reason: reason.into(),
environment: false,
}
}
#[must_use]
pub(crate) fn with_environment(mut self, environment: bool) -> Self {
self.environment = environment;
self
}
}
impl std::fmt::Display for StoreRefusal {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"refusing to start: store '{}' is not usable — {}",
self.store, self.reason
)
}
}
impl std::error::Error for StoreRefusal {}
async fn with_open_lock_retry<T, F, Fut>(
attempts: usize,
delay: Duration,
mut attempt: F,
) -> Result<T, turso::Error>
where
F: FnMut() -> Fut,
Fut: Future<Output = Result<T, turso::Error>>,
{
let mut remaining = attempts;
loop {
match attempt().await {
Ok(value) => return Ok(value),
Err(err) if is_open_time_lock_error(&err) => {
warn!(
attempt = attempts - remaining + 1,
total_attempts = attempts + 1,
error = ?err,
"database open blocked by another process (open-time store lock)"
);
if remaining == 0 {
return Err(err);
}
tokio::time::sleep(delay).await;
remaining -= 1;
}
Err(err) => return Err(err),
}
}
}
impl Connection {
pub async fn open(path: &Path) -> anyhow::Result<Self> {
let path_str = path
.to_str()
.with_context(|| format!("database path must be UTF-8: {}", path.display()))?;
let opts = experimental_database_opts();
let db = with_open_lock_retry(OPEN_LOCK_RETRY_ATTEMPTS, OPEN_LOCK_RETRY_DELAY, || async {
Builder::new_local(path_str)
.experimental_index_method(opts.enable_index_method)
.build()
.await
})
.await
.context("failed to open local database")?;
let conn = db.connect()?;
conn.busy_timeout(Duration::from_mins(1))?;
conn.execute("PRAGMA temp_store = MEMORY;", ())
.await
.context("failed to set in-memory temp storage (PRAGMA temp_store = MEMORY)")?;
Ok(Self {
conn: Arc::new(tokio::sync::Mutex::new(conn)),
has_dangling_tx: Arc::new(AtomicBool::new(false)),
db_path: Arc::from(path),
})
}
pub(crate) fn db_path(&self) -> &std::path::Path {
&self.db_path
}
async fn lock_and_cleanup(&self) -> tokio::sync::MutexGuard<'_, turso::Connection> {
Self::lock_and_cleanup_owned(&self.conn, &self.has_dangling_tx).await
}
async fn lock_and_cleanup_owned<'a>(
conn: &'a tokio::sync::Mutex<turso::Connection>,
has_dangling_tx: &'a AtomicBool,
) -> tokio::sync::MutexGuard<'a, turso::Connection> {
let guard = conn.lock().await;
let armed = has_dangling_tx.swap(false, Ordering::SeqCst);
if !matches!(guard.is_autocommit(), Ok(true)) {
match guard.execute("ROLLBACK", ()).await {
Ok(_) => tracing::warn!(
armed,
"rolled back a transaction left open on the shared connection"
),
Err(e) => {
has_dangling_tx.store(true, Ordering::SeqCst);
tracing::warn!(
error = %e,
armed,
"failed to roll back a transaction left open on the shared connection — \
the next access retries it"
);
}
}
}
guard
}
async fn run_detached<T, F>(&self, op: F) -> turso::Result<T>
where
T: Send + 'static,
F: for<'a> FnOnce(
&'a tokio::sync::MutexGuard<'a, turso::Connection>,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = turso::Result<T>> + Send + 'a>,
> + Send
+ 'static,
{
self.run_detached_capturing(op, None).await
}
async fn run_detached_capturing<T, F>(
&self,
op: F,
cause: Option<&checkpoint_cause::CauseSink>,
) -> turso::Result<T>
where
T: Send + 'static,
F: for<'a> FnOnce(
&'a tokio::sync::MutexGuard<'a, turso::Connection>,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = turso::Result<T>> + Send + 'a>,
> + Send
+ 'static,
{
let conn = self.conn.clone();
let dangling = self.has_dangling_tx.clone();
let task = async move {
let guard = Self::lock_and_cleanup_owned(&conn, &dangling).await;
op(&guard).await
};
let joined = match cause.cloned() {
Some(cause) => tokio::spawn(cause.scoped(task)).await,
None => tokio::spawn(task).await,
};
match joined {
Ok(inner) => inner,
Err(joined) => match joined.try_into_panic() {
Ok(payload) => std::panic::resume_unwind(payload),
Err(_) => Err(turso::Error::Error(
"database op cancelled (runtime shutdown)".to_string(),
)),
},
}
}
pub async fn upsert_row<U, I>(
&self,
update_sql: &str,
update_params: impl Fn() -> U,
insert_sql: &str,
insert_params: I,
) -> turso::Result<()>
where
U: IntoParams + Send + 'static,
I: IntoParams + Send + 'static,
{
if self.execute(update_sql, update_params()).await? > 0 {
return Ok(());
}
if self.execute(insert_sql, insert_params).await? == 0 {
self.execute(update_sql, update_params()).await?;
}
Ok(())
}
pub async fn execute(
&self,
sql: &str,
params: impl IntoParams + Send + 'static,
) -> turso::Result<u64> {
let sql = sql.to_owned();
self.run_detached(move |conn| Box::pin(async move { conn.execute(&sql, params).await }))
.await
}
pub(crate) async fn execute_cached(
&self,
sql: &str,
params: impl IntoParams + Send + 'static,
) -> turso::Result<u64> {
let sql = sql.to_owned();
self.run_detached(move |conn| {
Box::pin(async move {
let mut stmt = conn.prepare_cached(&sql).await?;
stmt.execute(params).await
})
})
.await
}
pub(crate) async fn execute_batch(&self, sql: &str) -> turso::Result<()> {
let sql = sql.to_owned();
self.run_detached(move |conn| Box::pin(async move { conn.execute_batch(&sql).await }))
.await
}
pub async fn begin_tx(&self) -> turso::Result<TxGuard<'_>> {
let conn = self.lock_and_cleanup().await;
self.has_dangling_tx.store(true, Ordering::SeqCst);
conn.execute("BEGIN", ()).await?;
self.has_dangling_tx.store(false, Ordering::SeqCst);
Ok(TxGuard {
conn,
has_dangling_tx: Some(self.has_dangling_tx.clone()),
})
}
pub async fn query(
&self,
sql: &str,
params: impl IntoParams + Send + 'static,
) -> turso::Result<Vec<Row>> {
self.query_with_cause(sql, params, None).await
}
async fn query_with_cause(
&self,
sql: &str,
params: impl IntoParams + Send + 'static,
cause: Option<&checkpoint_cause::CauseSink>,
) -> turso::Result<Vec<Row>> {
let sql = sql.to_owned();
self.run_detached_capturing(
move |conn| Box::pin(async move { Self::query_impl(conn, &sql, params).await }),
cause,
)
.await
}
pub(crate) async fn query_readonly(
&self,
sql: &str,
params: impl IntoParams + Send + 'static,
row_limit: usize,
) -> turso::Result<ReadonlyRows> {
let sql = sql.to_owned();
let dangling = self.has_dangling_tx.clone();
self.run_detached(move |conn| {
Box::pin(async move {
conn.execute("PRAGMA query_only = 1;", ()).await?;
let result =
AssertUnwindSafe(run_readonly_query(conn, &sql, params, row_limit))
.catch_unwind()
.await;
let (reset, tx_reset) = match AssertUnwindSafe(async {
let reset = conn.execute("PRAGMA query_only = 0;", ()).await;
let tx_reset = match conn.is_autocommit() {
Ok(true) => Ok(()),
_ => conn.execute("ROLLBACK;", ()).await.map(|_| ()),
};
(reset, tx_reset)
})
.catch_unwind()
.await
{
Ok(ok) => ok,
Err(payload) => {
dangling.store(true, Ordering::SeqCst);
return Err(turso::Error::Error(format!(
"the database engine panicked while resetting the read-only query state: {}",
crate::util::panic_message(&*payload)
)));
}
};
if let Err(e) = &reset {
tracing::error!(
error = %e,
"query_readonly: failed to reset PRAGMA query_only — instance writes may \
be left disabled on the shared connection"
);
}
if let Err(e) = &tx_reset {
dangling.store(true, Ordering::SeqCst);
tracing::error!(
error = %e,
"query_readonly: failed to roll back a leaked transaction — instance writes may \
be left disabled on the shared connection"
);
}
reset?;
tx_reset?;
match result {
Ok(Ok(rows)) => Ok(rows),
Ok(Err(e)) => Err(e),
Err(payload) => Err(turso::Error::Error(format!(
"the database engine panicked while executing the query: {}",
crate::util::panic_message(&*payload)
))),
}
})
})
.await
}
async fn query_impl(
conn: &turso::Connection,
sql: &str,
params: impl IntoParams + Send + 'static,
) -> turso::Result<Vec<Row>> {
let rows = conn.query(sql, params).await?;
drain_rows(rows).await
}
pub async fn query_map<T, E>(
&self,
sql: &str,
params: impl IntoParams + Send + 'static,
map: impl FnMut(&Row) -> std::result::Result<T, E> + Send + 'static,
) -> turso::Result<Vec<turso::Result<T>>>
where
T: Send + 'static,
E: std::fmt::Display + Send + Sync + 'static,
{
let rows = self.query(sql, params).await?;
Ok(map_rows(&rows, map))
}
pub async fn query_map_strict<T, E>(
&self,
sql: &str,
params: impl IntoParams + Send + 'static,
map: impl FnMut(&Row) -> std::result::Result<T, E> + Send + 'static,
) -> anyhow::Result<Vec<T>>
where
T: Send + 'static,
E: std::fmt::Display + Send + Sync + 'static,
{
strict_collect(self.query_map(sql, params, map).await?)
}
pub async fn query_row<T, E>(
&self,
sql: &str,
params: impl IntoParams + Send + 'static,
map: impl FnOnce(&Row) -> std::result::Result<T, E> + Send + 'static,
) -> turso::Result<T>
where
T: Send + 'static,
E: std::fmt::Display + Send + Sync + 'static,
{
let sql = sql.to_owned();
self.run_detached(move |conn| {
Box::pin(async move { Self::query_row_impl(conn, &sql, params, map).await })
})
.await
}
pub async fn query_optional<T, E>(
&self,
sql: &str,
params: impl IntoParams + Send + 'static,
map: impl FnOnce(&Row) -> std::result::Result<T, E> + Send + 'static,
) -> anyhow::Result<Option<T>>
where
T: Send + 'static,
E: std::fmt::Display + Send + Sync + 'static,
{
optional_row(self.query_row(sql, params, map).await)
}
async fn query_row_impl<T, E>(
conn: &turso::Connection,
sql: &str,
params: impl IntoParams + Send + 'static,
map: impl FnOnce(&Row) -> std::result::Result<T, E> + Send + 'static,
) -> turso::Result<T>
where
T: Send + 'static,
E: std::fmt::Display + Send + Sync + 'static,
{
let rows = conn.query(sql, params).await?;
map_first_row(rows, map).await
}
pub async fn query_cached(
&self,
sql: &str,
params: impl IntoParams + Send + 'static,
) -> turso::Result<Vec<Row>> {
let sql = sql.to_owned();
self.run_detached(move |conn| {
Box::pin(async move { Self::query_cached_impl(conn, &sql, params).await })
})
.await
}
async fn query_cached_impl(
conn: &turso::Connection,
sql: &str,
params: impl IntoParams + Send + 'static,
) -> turso::Result<Vec<Row>> {
let mut stmt = conn.prepare_cached(sql).await?;
let rows = stmt.query(params).await?;
drain_rows(rows).await
}
pub async fn query_map_strict_cached<T, E>(
&self,
sql: &str,
params: impl IntoParams + Send + 'static,
map: impl FnMut(&Row) -> std::result::Result<T, E> + Send + 'static,
) -> anyhow::Result<Vec<T>>
where
T: Send + 'static,
E: std::fmt::Display + Send + Sync + 'static,
{
let rows = self.query_cached(sql, params).await?;
strict_collect(map_rows(&rows, map))
}
async fn query_row_cached<T, E>(
&self,
sql: &str,
params: impl IntoParams + Send + 'static,
map: impl FnOnce(&Row) -> std::result::Result<T, E> + Send + 'static,
) -> turso::Result<T>
where
T: Send + 'static,
E: std::fmt::Display + Send + Sync + 'static,
{
let sql = sql.to_owned();
self.run_detached(move |conn| {
Box::pin(async move {
let mut stmt = conn.prepare_cached(&sql).await?;
map_first_row(stmt.query(params).await?, map).await
})
})
.await
}
pub async fn query_optional_cached<T, E>(
&self,
sql: &str,
params: impl IntoParams + Send + 'static,
map: impl FnOnce(&Row) -> std::result::Result<T, E> + Send + 'static,
) -> anyhow::Result<Option<T>>
where
T: Send + 'static,
E: std::fmt::Display + Send + Sync + 'static,
{
optional_row(self.query_row_cached(sql, params, map).await)
}
pub(crate) async fn checkpoint_ungated(&self) -> anyhow::Result<CheckpointOutcome> {
self.run_checkpoint(CheckpointMode::Truncate).await
}
async fn checkpoint_passive(&self) -> anyhow::Result<CheckpointOutcome> {
self.run_checkpoint(CheckpointMode::Passive).await
}
pub(crate) async fn checkpoint_mode(
&self,
truncate: bool,
) -> anyhow::Result<CheckpointOutcome> {
if truncate {
self.checkpoint_ungated().await
} else {
self.checkpoint_passive().await
}
}
async fn run_checkpoint(&self, mode: CheckpointMode) -> anyhow::Result<CheckpointOutcome> {
let cause = checkpoint_cause::CauseSink::new();
let sql = format!("PRAGMA wal_checkpoint({});", mode.label());
let rows = self
.query_with_cause(&sql, (), Some(&cause))
.await
.map_err(|e| cause.attach(anyhow::Error::from(e)))
.context("Failed to checkpoint WAL")?;
let row = rows
.first()
.context("PRAGMA wal_checkpoint returned no result row")?;
parse_checkpoint_row(row).map_err(|e| cause.attach(e))
}
pub async fn quick_check(&self) -> anyhow::Result<()> {
if let Some(problem) = self.quick_check_problems().await?.into_iter().next() {
anyhow::bail!("Database integrity check failed: {problem}");
}
Ok(())
}
pub(crate) async fn quick_check_problems(&self) -> anyhow::Result<Vec<String>> {
let rows = self
.query("PRAGMA quick_check;", ())
.await
.context("Failed to execute PRAGMA quick_check")?;
scan_integrity_rows(&rows)
}
}
fn scan_integrity_rows(rows: &[Row]) -> anyhow::Result<Vec<String>> {
let mut problems: Vec<String> = Vec::new();
for row in rows {
match row.get_value(0)? {
Value::Text(s) if s == "ok" => {}
Value::Text(s) if s.contains(KNOWN_FTS_DIR_COUNT_FALSE_POSITIVE) => {}
Value::Text(s) => problems.push(s),
_ => anyhow::bail!("Unexpected result from PRAGMA quick_check"),
}
}
Ok(problems)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct CheckpointOutcome {
pub busy: bool,
pub log_frames: i64,
pub checkpointed_frames: i64,
}
impl CheckpointOutcome {
#[must_use]
pub fn is_complete(&self) -> bool {
!self.busy && self.log_frames <= self.checkpointed_frames
}
}
fn int_column(row: &Row, idx: usize) -> anyhow::Result<i64> {
match row.get_value(idx)? {
Value::Integer(n) => Ok(n),
_ => anyhow::bail!("Unexpected result from PRAGMA wal_checkpoint"),
}
}
fn parse_checkpoint_row(row: &Row) -> anyhow::Result<CheckpointOutcome> {
Ok(CheckpointOutcome {
busy: int_column(row, 0)? != 0,
log_frames: int_column(row, 1)?,
checkpointed_frames: int_column(row, 2)?,
})
}
pub(crate) struct TxGuard<'a> {
conn: tokio::sync::MutexGuard<'a, turso::Connection>,
has_dangling_tx: Option<Arc<AtomicBool>>,
}
impl TxGuard<'_> {
pub async fn upsert_row<U, I>(
&self,
update_sql: &str,
update_params: impl Fn() -> U,
insert_sql: &str,
insert_params: I,
) -> turso::Result<()>
where
U: IntoParams + Send + 'static,
I: IntoParams + Send + 'static,
{
if self.execute(update_sql, update_params()).await? > 0 {
return Ok(());
}
if self.execute(insert_sql, insert_params).await? == 0 {
self.execute(update_sql, update_params()).await?;
}
Ok(())
}
pub async fn execute(
&self,
sql: &str,
params: impl IntoParams + Send + 'static,
) -> turso::Result<u64> {
self.conn.execute(sql, params).await
}
pub async fn execute_batch(&self, sql: &str) -> turso::Result<()> {
self.conn.execute_batch(sql).await
}
pub async fn query_row<T, E>(
&self,
sql: &str,
params: impl IntoParams + Send + 'static,
map: impl FnOnce(&Row) -> std::result::Result<T, E> + Send + 'static,
) -> turso::Result<T>
where
T: Send + 'static,
E: std::fmt::Display + Send + Sync + 'static,
{
Connection::query_row_impl(&self.conn, sql, params, map).await
}
pub async fn query(
&self,
sql: &str,
params: impl IntoParams + Send + 'static,
) -> turso::Result<Vec<Row>> {
Connection::query_impl(&self.conn, sql, params).await
}
pub async fn query_map_strict<T, E>(
&self,
sql: &str,
params: impl IntoParams + Send + 'static,
map: impl FnMut(&Row) -> std::result::Result<T, E> + Send + 'static,
) -> anyhow::Result<Vec<T>>
where
T: Send + 'static,
E: std::fmt::Display + Send + Sync + 'static,
{
let rows = self.query(sql, params).await?;
strict_collect(map_rows(&rows, map))
}
pub async fn commit(mut self) -> turso::Result<()> {
self.conn.execute("COMMIT", ()).await?;
self.has_dangling_tx = None;
Ok(())
}
pub async fn rollback(mut self) -> turso::Result<()> {
self.conn.execute("ROLLBACK", ()).await?;
self.has_dangling_tx = None;
Ok(())
}
}
impl Drop for TxGuard<'_> {
fn drop(&mut self) {
if let Some(flag) = &self.has_dangling_tx {
flag.store(true, Ordering::SeqCst);
}
}
}
pub(crate) async fn ensure_fts_index(
conn: &Connection,
index_name: &str,
tokenizer: &str,
ddl: &str,
) -> anyhow::Result<()> {
let existing_sql: Option<String> = conn
.query_optional(
"SELECT sql FROM sqlite_master WHERE type='index' AND name=?1 LIMIT 1",
params![index_name],
|row| match row.get_value(0)? {
Value::Text(s) => Ok::<_, ::turso::Error>(s),
_ => Ok::<_, ::turso::Error>(String::new()),
},
)
.await?
.filter(|s| !s.is_empty());
let needs_rebuild = existing_sql
.as_deref()
.is_none_or(|sql| !sql.to_lowercase().contains(&tokenizer.to_lowercase()));
if needs_rebuild {
conn.execute(&format!("DROP INDEX IF EXISTS {index_name}"), ())
.await?;
conn.execute(ddl, ()).await?;
}
Ok(())
}
const FTS_PROBE_TERM: &str = "mahbot";
fn names_ticket_title_fts(problem: &str, index: &str) -> bool {
problem.contains(index) || problem.contains(FTS_INTERNAL_INDEX_PREFIX)
}
async fn is_fts_index(conn: &Connection, index: &str) -> bool {
conn.query_optional(
"SELECT sql FROM sqlite_master WHERE type = 'index' AND name = ?",
(index.to_string(),),
|row| row.get::<String>(0),
)
.await
.ok()
.flatten()
.is_some_and(|sql| sql.to_lowercase().contains("using fts"))
}
pub(crate) const TICKETS_FTS_INDEX_NAME: &str = "idx_tickets_title_fts";
pub(crate) const TICKETS_FTS_INDEX_DDL: &str = "CREATE INDEX IF NOT EXISTS idx_tickets_title_fts ON tickets \
USING fts (title) WITH (tokenizer = 'ngram')";
async fn title_fts_match_probe(conn: &Connection) -> anyhow::Result<()> {
conn.query_optional(
"SELECT id FROM tickets WHERE title MATCH ?1 LIMIT 1",
params![FTS_PROBE_TERM.to_string()],
|_| Ok::<_, ::turso::Error>(()),
)
.await
.map(|_: Option<()>| ())
}
async fn detect_ticket_title_fts_corruption(
conn: &Connection,
index: &str,
) -> anyhow::Result<Option<String>> {
let problems = conn.quick_check_problems().await?;
let fts_problems: Vec<&str> = problems
.iter()
.map(String::as_str)
.filter(|p| names_ticket_title_fts(p, index))
.collect();
if !fts_problems.is_empty() {
return Ok(Some(fts_problems.join("; ")));
}
if let Err(e) = title_fts_match_probe(conn).await {
return Ok(Some(format!("MATCH probe failed: {e}")));
}
Ok(None)
}
async fn verify_rebuilt_ticket_title_fts(conn: &Connection, index: &str) -> anyhow::Result<()> {
let problems = conn.quick_check_problems().await?;
if !problems.is_empty() {
warn!(
index,
?problems,
"ticket-title FTS post-rebuild quick_check still reports problems"
);
}
let present = conn
.query_optional(
"SELECT 1 FROM sqlite_master WHERE type = 'index' AND name = ?1",
params![index.to_string()],
|_| Ok::<_, ::turso::Error>(()),
)
.await?;
anyhow::ensure!(
present.is_some(),
"rebuilt index '{index}' not found in sqlite_master"
);
let titles: Vec<String> = conn
.query_map(
"SELECT title FROM tickets WHERE title IS NOT NULL AND title != '' LIMIT 16",
(),
|row| row.get::<String>(0),
)
.await?
.into_iter()
.collect::<std::result::Result<Vec<_>, _>>()?;
let Some((title, term)) = titles.iter().find_map(|title| {
let sanitized = sanitize_fts_query(title);
let term = sanitized.split(' ').find(|tok| tok.chars().count() >= 3)?;
Some((title.clone(), term.to_string()))
}) else {
debug!(
index,
"no ticket title with an ngram-matchable token — MATCH smoke skipped"
);
return Ok(());
};
let hit = conn
.query_optional(
"SELECT id FROM tickets WHERE title MATCH ?1 LIMIT 1",
params![term.clone()],
|_| Ok::<_, ::turso::Error>(()),
)
.await?;
anyhow::ensure!(
hit.is_some(),
"MATCH smoke failed: token '{term}' from title '{title}' not found via MATCH"
);
info!(
index,
title, "ticket-title FTS MATCH smoke assertion passed"
);
Ok(())
}
struct TicketTitleFtsRebuildOutcome {
drop_create: anyhow::Result<()>,
verified: anyhow::Result<()>,
}
async fn rebuild_ticket_title_fts(
conn: &Connection,
index: &str,
ddl: &str,
reclaim_after_rebuild: bool,
) -> TicketTitleFtsRebuildOutcome {
let quoted = index.replace('"', "\"\"");
let drop_create = async {
let tx = conn.begin_tx().await?;
tx.execute_batch(&format!("DROP INDEX IF EXISTS \"{quoted}\"; {ddl};"))
.await?;
tx.commit().await
}
.await
.map_err(anyhow::Error::from);
let verified = verify_rebuilt_ticket_title_fts(conn, index).await;
match (&drop_create, &verified) {
(Ok(()), Ok(())) => {
info!(index, "ticket-title FTS index rebuilt and verified");
}
(drop_create, verified) => {
warn!(
index,
repair_error = ?drop_create.as_ref().err(),
verify_error = ?verified.as_ref().err(),
"ticket-title FTS rebuild finished with warnings — check earlier log lines"
);
}
}
if reclaim_after_rebuild {
match conn.checkpoint_ungated().await {
Ok(o) if o.is_complete() => {
info!(
index,
checkpointed_frames = o.checkpointed_frames,
"post-repair WAL checkpoint complete — corruption cleared"
);
}
Ok(o) => {
warn!(
index,
busy = o.busy,
log_frames = o.log_frames,
checkpointed_frames = o.checkpointed_frames,
"post-repair WAL checkpoint incomplete"
);
}
Err(e) => {
warn!(
index,
error = %e,
"post-repair WAL checkpoint failed"
);
}
}
}
TicketTitleFtsRebuildOutcome {
drop_create,
verified,
}
}
pub(crate) async fn repair_ticket_title_fts_if_corrupt(
conn: &Connection,
index: &str,
ddl: &str,
store: &'static str,
root: Option<&Path>,
) {
let detected = AssertUnwindSafe(detect_ticket_title_fts_corruption(conn, index))
.catch_unwind()
.await;
let evidence = match detected {
Ok(Ok(evidence)) => evidence,
Ok(Err(e)) => {
error!(
index,
error = %e,
"ticket-title FTS boot probe failed — boot continues; index left as-is for operator review"
);
return;
}
Err(panic) => {
warn!(
index,
panic = %crate::util::panic_message(&panic),
"ticket-title FTS boot probe panicked — treating as corruption evidence, attempting the rebuild"
);
Some("detection probe panicked".to_string())
}
};
let Some(evidence) = evidence else {
info!(
index,
"ticket-title FTS boot probe: index healthy — no repair"
);
return;
};
warn!(
index,
evidence = %evidence,
"ticket-title FTS corruption signature detected at boot — rebuilding the index (ticket rows untouched)"
);
let reclaim_after_rebuild = crate::db::shrink_gate::shrink_allowed(conn, store, root).await;
let rebuilt = AssertUnwindSafe(rebuild_ticket_title_fts(
conn,
index,
ddl,
reclaim_after_rebuild,
))
.catch_unwind()
.await;
if let Err(panic) = rebuilt {
error!(
index,
panic = %crate::util::panic_message(&panic),
"ticket-title FTS boot repair panicked — boot continues; index left as-is for operator review"
);
}
}
#[derive(Debug)]
pub(crate) enum TicketTitleFtsRuntimeRepair {
NotApplicable,
Healthy,
Rebuilt(String),
Failed(String),
}
impl TicketTitleFtsRuntimeRepair {
#[must_use]
pub(crate) fn summary(&self) -> String {
match self {
Self::NotApplicable => "not applicable (no ticket-title FTS index)".to_string(),
Self::Healthy => "healthy (index present, probe clean)".to_string(),
Self::Rebuilt(evidence) => format!("rebuilt after detection — {evidence}"),
Self::Failed(reason) => format!("failed — {reason}"),
}
}
}
pub(crate) async fn repair_ticket_title_fts_runtime(
conn: &Connection,
) -> TicketTitleFtsRuntimeRepair {
let present = AssertUnwindSafe(conn.query_optional(
"SELECT 1 FROM sqlite_master WHERE type = 'index' AND name = ?",
params![TICKETS_FTS_INDEX_NAME.to_string()],
|_| Ok::<_, ::turso::Error>(()),
))
.catch_unwind()
.await;
let present = match present {
Ok(Ok(present)) => present,
Ok(Err(e)) => {
return TicketTitleFtsRuntimeRepair::Failed(format!(
"index presence probe failed: {e}"
));
}
Err(panic) => {
return TicketTitleFtsRuntimeRepair::Failed(format!(
"index presence probe panicked: {}",
crate::util::panic_message(&panic)
));
}
};
if present.is_none() {
return TicketTitleFtsRuntimeRepair::NotApplicable;
}
let detected = AssertUnwindSafe(detect_ticket_title_fts_corruption(
conn,
TICKETS_FTS_INDEX_NAME,
))
.catch_unwind()
.await;
let evidence = match detected {
Ok(Ok(None)) => return TicketTitleFtsRuntimeRepair::Healthy,
Ok(Ok(Some(evidence))) => evidence,
Ok(Err(e)) => {
return TicketTitleFtsRuntimeRepair::Failed(format!("detection probe failed: {e}"));
}
Err(panic) => format!(
"detection probe panicked: {}",
crate::util::panic_message(&panic)
),
};
let rebuilt = AssertUnwindSafe(rebuild_ticket_title_fts(
conn,
TICKETS_FTS_INDEX_NAME,
TICKETS_FTS_INDEX_DDL,
false,
))
.catch_unwind()
.await;
let TicketTitleFtsRebuildOutcome {
drop_create,
verified,
} = match rebuilt {
Ok(outcome) => outcome,
Err(panic) => {
return TicketTitleFtsRuntimeRepair::Failed(format!(
"rebuild panicked: {}; evidence: {evidence}",
crate::util::panic_message(&panic)
));
}
};
if drop_create.is_err() || verified.is_err() {
return TicketTitleFtsRuntimeRepair::Failed(format!(
"rebuild incomplete: drop_create error: {:?}; verify error: {:?}; evidence: {evidence}",
drop_create.as_ref().err(),
verified.as_ref().err(),
));
}
TicketTitleFtsRuntimeRepair::Rebuilt(format!("corruption detected: {evidence}"))
}
pub(crate) async fn open_with_schema(db_path: &Path, schema: &str) -> anyhow::Result<Connection> {
if let Some(parent) = db_path.parent() {
std::fs::create_dir_all(parent)
.with_context(|| format!("Failed to create directory: {}", parent.display()))?;
}
let conn = Connection::open(db_path)
.await
.with_context(|| format!("Failed to open database: {}", db_path.display()))?;
conn.execute("PRAGMA foreign_keys = ON;", ())
.await
.context("Failed to enable foreign key enforcement")?;
execute_schema_ddl(&conn, schema).await?;
Ok(conn)
}
async fn execute_schema_ddl(conn: &Connection, sql: &str) -> anyhow::Result<()> {
if sql.trim().is_empty() {
return Ok(());
}
conn.execute_batch(sql)
.await
.context("Failed to run schema DDL")?;
Ok(())
}
#[must_use]
pub(crate) fn store_db_path(root: &Path, name: &str) -> std::path::PathBuf {
let stem = if name == LOG_DB_NAME {
LOG_DB_NAME
} else {
CONSOLIDATED_DB_NAME
};
root.join("db").join(format!("{stem}.db"))
}
#[must_use]
pub(crate) fn wal_path(db_path: &Path) -> std::path::PathBuf {
std::path::PathBuf::from(format!("{}-wal", db_path.display()))
}
const RESOURCE_SIGNAL_KEYWORDS: [&str; 4] = [
"no space left on device",
"too many open files",
"out of memory",
"permission denied",
];
fn has_resource_signal(lower: &str) -> bool {
RESOURCE_SIGNAL_KEYWORDS.iter().any(|k| lower.contains(k))
}
const LOCK_SIGNAL_KEYWORDS: [&str; 3] = ["locking", "locked", "busy"];
fn has_lock_signal(lower: &str) -> bool {
LOCK_SIGNAL_KEYWORDS.iter().any(|k| lower.contains(k))
}
pub(crate) fn is_actionable_signal(e: &anyhow::Error) -> bool {
text_is_actionable_signal(&format!("{e:#}"))
}
pub(crate) fn text_is_actionable_signal(text: &str) -> bool {
let lower = text.to_lowercase();
has_resource_signal(&lower) || has_lock_signal(&lower)
}
pub(crate) async fn open_store(
root: &Path,
name: &'static str,
schema: &str,
) -> anyhow::Result<Connection> {
let db_path = store_db_path(root, name);
match crate::db::wal_guard::classify_store_shape(&db_path) {
crate::db::wal_guard::StoreShape::Absent => open_with_schema(&db_path, schema).await,
crate::db::wal_guard::StoreShape::Unusable(defect) => {
Err(StoreRefusal::new(name, &db_path, defect.reason())
.with_environment(defect.is_environment_caused())
.into())
}
crate::db::wal_guard::StoreShape::Present => {
let opened = AssertUnwindSafe(open_present_store(root, &db_path, name, schema))
.catch_unwind()
.await;
match opened {
Ok(result) => result,
Err(payload) => Err(StoreRefusal::new(
name,
&db_path,
format!("open panicked: {}", crate::util::panic_message(&*payload)),
)
.into()),
}
}
}
}
async fn open_present_store(
root: &Path,
db_path: &Path,
name: &'static str,
schema: &str,
) -> anyhow::Result<Connection> {
let conn = match open_with_schema(db_path, schema).await {
Ok(conn) => conn,
Err(e) if is_actionable_signal(&e) => return Err(e),
Err(e) => {
return Err(
StoreRefusal::new(name, db_path, format!("it could not be opened: {e:#}")).into(),
);
}
};
match has_product_schema(&conn).await {
Ok(true) => {}
Ok(false) => {
return Err(StoreRefusal::new(
name,
db_path,
"it opens but does not carry the product's own schema",
)
.into());
}
Err(e) if is_actionable_signal(&e) => return Err(e),
Err(e) => {
return Err(StoreRefusal::new(
name,
db_path,
format!("the schema probe failed: {e:#}"),
)
.into());
}
}
run_repairs(conn, root, db_path, name, schema).await
}
async fn has_product_schema(conn: &Connection) -> anyhow::Result<bool> {
let table = crate::db::migrations::MIGRATIONS_TABLE;
let ledger: Option<i64> = conn
.query_optional(
&format!(
"SELECT COUNT(*) FROM sqlite_master WHERE type = 'table' AND name = '{table}'"
),
(),
|row| row.get::<i64>(0),
)
.await?;
if ledger.unwrap_or(0) == 0 {
return Ok(false);
}
let recorded: i64 = conn
.query_row(&format!("SELECT COUNT(*) FROM {table}"), (), |row| {
row.get::<i64>(0)
})
.await?;
Ok(recorded > 0)
}
async fn run_repairs(
conn: Connection,
root: &Path,
db_path: &Path,
name: &'static str,
schema: &str,
) -> anyhow::Result<Connection> {
match repair_btree_index_if_desynced(conn, root, db_path, name, schema).await {
Ok((RepairOutcome::Unreadable, _conn)) => {
Err(StoreRefusal::new(name, db_path, "its integrity check cannot read a table").into())
}
Ok((_, conn)) => Ok(conn),
Err(e) if is_actionable_signal(&e) => Err(e),
Err(e) => Err(StoreRefusal::new(
name,
db_path,
format!("the data-preserving rebuild could not be completed: {e:#}"),
)
.into()),
}
}
pub(crate) async fn open_consolidated_store(root: &Path) -> anyhow::Result<Connection> {
let conn = open_store(root, CONSOLIDATED_DB_NAME, "").await?;
migrations::run_migrations(&conn, migrations::TargetDb::Core).await?;
cdc::enable_capture(&conn).await?;
Ok(conn)
}
enum RepairOutcome {
NoRepair,
Repaired,
Migrated,
Unreadable,
}
#[cfg(test)]
async fn open_and_repair_for_test(
root: &Path,
db_path: &Path,
name: &'static str,
schema: &str,
) -> anyhow::Result<Connection> {
let conn = open_with_schema(db_path, schema).await?;
run_repairs(conn, root, db_path, name, schema).await
}
const BAKED_OVERFLOW_ALIASING_READ: &str = "short read on page 167772160";
enum RepairTarget {
Unreadable,
OverflowAliasing,
Index(String),
Unknown,
}
fn classify_repair_target(problems: &[String]) -> RepairTarget {
if problems.iter().any(|p| p.contains("Invalid page type")) {
RepairTarget::Unreadable
} else if problems.iter().any(|p| {
p.contains("referenced multiple times") || p.contains(BAKED_OVERFLOW_ALIASING_READ)
}) {
RepairTarget::OverflowAliasing
} else if problems.iter().any(|p| p.contains("short read")) {
RepairTarget::Unreadable
} else {
match problems.iter().find_map(|p| desynced_index_name(p)) {
Some(index) => RepairTarget::Index(index.to_string()),
None => RepairTarget::Unknown,
}
}
}
#[must_use]
pub(crate) fn is_tolerated_integrity_report(problems: &[String]) -> bool {
matches!(classify_repair_target(problems), RepairTarget::Unknown)
}
#[must_use]
fn pre_reindex_snapshot_path(db_path: &Path) -> std::path::PathBuf {
std::path::PathBuf::from(format!(
"{}.pre-reindex-{}",
db_path.display(),
family_stamp()
))
}
async fn snapshot_store_via_engine(
conn: &Connection,
db_path: &Path,
) -> anyhow::Result<std::path::PathBuf> {
let dest = pre_reindex_snapshot_path(db_path);
let escaped = dest.display().to_string().replace('\'', "''");
match AssertUnwindSafe(conn.execute(&format!("VACUUM INTO '{escaped}'"), ()))
.catch_unwind()
.await
{
Ok(Ok(_)) => Ok(dest),
Ok(Err(e)) => {
Err(anyhow::Error::from(e).context(format!("VACUUM INTO {}", dest.display())))
}
Err(payload) => Err(anyhow::anyhow!(
"VACUUM INTO {} panicked: {}",
dest.display(),
crate::util::panic_message(&*payload)
)),
}
}
async fn repair_btree_index_if_desynced(
conn: Connection,
root: &Path,
db_path: &Path,
name: &'static str,
schema: &str,
) -> anyhow::Result<(RepairOutcome, Connection)> {
let problems = match conn.quick_check_problems().await {
Ok(p) if p.is_empty() => return Ok((RepairOutcome::NoRepair, conn)),
Ok(p) => p,
Err(e) => {
if is_actionable_signal(&e) {
return Err(e);
}
crate::boot::boot_diagnostic(format!(
"store '{name}' quick_check could not run: {e} — unreadable table; the \
store is refused",
));
return Ok((RepairOutcome::Unreadable, conn));
}
};
let index = match classify_repair_target(&problems) {
RepairTarget::Unreadable => {
crate::boot::boot_diagnostic(format!(
"store '{name}' quick_check cannot scan a table ({}) — unreadable table; \
the store is refused",
problems.join("; "),
));
return Ok((RepairOutcome::Unreadable, conn));
}
RepairTarget::OverflowAliasing => {
crate::boot::boot_diagnostic(format!(
"store '{name}' quick_check reports overflow-aliasing ({}) — rebuilding \
the store data-preservingly in a fresh file",
problems.join("; "),
));
return migrate_overflow_aliased_store(conn, root, db_path, name, schema).await;
}
RepairTarget::Unknown => {
crate::boot::boot_diagnostic(format!(
"store '{name}' quick_check flagged an unknown condition: {}",
problems.join("; "),
));
return Ok((RepairOutcome::NoRepair, conn));
}
RepairTarget::Index(index) => {
if known_fts_dir_false_positive(&index) {
crate::boot::boot_diagnostic(format!(
"store '{name}' quick_check flagged the known FTS false positive — \
not repairing",
));
return Ok((RepairOutcome::NoRepair, conn));
}
if is_fts_index(&conn, &index).await {
crate::boot::boot_diagnostic(format!(
"store '{name}' quick_check flagged FTS index '{index}' — in-place repair \
out of scope; the after_open FTS repair owns its rebuild",
));
return Ok((RepairOutcome::NoRepair, conn));
}
index
}
};
if let Err(e) = snapshot_store_via_engine(&conn, db_path).await {
warn!("Failed to write pre-reindex snapshot: {e:#}");
}
let quoted = index.replace('"', "\"\"");
if conn
.execute_batch(&format!("REINDEX \"{quoted}\";"))
.await
.is_ok()
&& conn.quick_check().await.is_ok()
{
info!(
db = %name,
index = %index,
"class-B btree index desync repaired in place (REINDEX)",
);
return Ok((RepairOutcome::Repaired, conn));
}
crate::boot::boot_diagnostic(format!(
"store '{name}' REINDEX of '{index}' did not clear the quick_check desync — \
falling back to DROP+CREATE",
));
Ok((
drop_create_index_fallback(&conn, name, &index, "ed).await,
conn,
))
}
async fn count_table_rows(conn: &Connection, table: &str) -> turso::Result<i64> {
let quoted = table.replace('"', "\"\"");
conn.query_row(&format!("SELECT COUNT(*) FROM \"{quoted}\""), (), |r| {
r.get::<i64>(0)
})
.await
}
#[expect(clippy::too_many_lines)] async fn migrate_overflow_aliased_store(
conn: Connection,
root: &Path,
db_path: &Path,
name: &'static str,
schema: &str,
) -> anyhow::Result<(RepairOutcome, Connection)> {
let mut counts: Vec<(String, i64)> = Vec::new();
let tables = match conn
.query(
&format!(
"SELECT name FROM sqlite_master WHERE type='table' \
AND {USER_OBJECT_FILTER} ORDER BY rowid"
),
(),
)
.await
{
Ok(t) => t,
Err(e) => {
let err = anyhow::anyhow!(e);
if is_actionable_signal(&err) {
return Err(err);
}
crate::boot::boot_diagnostic(format!(
"store '{name}' cannot enumerate tables for the rebuild: {err} — left \
for operator review",
));
return Ok((RepairOutcome::NoRepair, conn));
}
};
for t in &tables {
let tbl = match t.get::<String>(0) {
Ok(tbl) => tbl,
Err(e) => {
crate::boot::boot_diagnostic(format!(
"store '{name}' table name is not text ({e}) — rebuild aborted; \
left for operator review",
));
return Ok((RepairOutcome::NoRepair, conn));
}
};
match count_table_rows(&conn, &tbl).await {
Ok(count) => counts.push((tbl, count)),
Err(e) => {
let err = anyhow::anyhow!(e);
return Err(if is_actionable_signal(&err) {
err
} else {
err.context(format!("table '{tbl}' cannot be read"))
});
}
}
}
let temp =
std::path::PathBuf::from(format!("{}.rebuild-{}", db_path.display(), family_stamp()));
let _temp_guard = TempCleanup(&temp);
let fresh = match Connection::open(&temp).await {
Ok(f) => f,
Err(e) if is_actionable_signal(&e) => return Err(e),
Err(e) => {
crate::boot::boot_diagnostic(format!(
"store '{name}' rebuild store open failed: {e} — left for operator review",
));
return Ok((RepairOutcome::NoRepair, conn));
}
};
let migrated = migrate_schema_and_data(&conn, &fresh, &counts).await;
let outcome = match migrated {
Ok(()) => match fresh.quick_check().await {
Ok(()) => checkpoint_rebuilt_before_swap(&fresh, name, root).await,
Err(e) => Err(classify_migrate(e)),
},
Err(f) => Err(f),
};
drop(fresh); if let Err(failure) = outcome {
return match failure {
MigrateFailure::Actionable(e) => Err(e),
MigrateFailure::Finding(msg) => {
crate::boot::boot_diagnostic(format!(
"store '{name}' rebuild aborted: {msg} — original store preserved \
for operator review (no data changed)",
));
Ok((RepairOutcome::NoRepair, conn))
}
};
}
drop(conn); if !quarantine_store_artifacts(db_path) {
crate::boot::boot_diagnostic(format!(
"store '{name}' original family could not be fully quarantined — the \
rebuild swap would clobber it without a forensic record; rebuild aborted, \
migrated data discarded",
));
return if db_path.exists() {
crate::boot::boot_diagnostic(format!(
"store '{name}' original family remains in place — left for operator \
review",
));
open_with_schema(db_path, schema)
.await
.map(|reopened| (RepairOutcome::NoRepair, reopened))
} else {
crate::boot::boot_diagnostic(format!(
"store '{name}' original family is in the quarantine and the store \
path is empty — boot aborted; recover from the quarantine and retry",
));
Err(anyhow::anyhow!(
"the rebuild aborted with the original family partially quarantined and the \
main file already moved — the store path was not rebuilt"
))
};
}
let wal = wal_path(db_path);
let temp_wal = wal_path(&temp);
if let Err(e) = std::fs::rename(&temp, db_path) {
return Err(anyhow::anyhow!(e).context(
"the rebuild swap could not move the rebuilt store into place — migrated temp \
family discarded; original family quarantined for recovery",
));
}
if temp_wal.exists()
&& let Err(e) = std::fs::rename(&temp_wal, &wal)
{
warn!(
error = %e,
from = %temp_wal.display(),
to = %wal.display(),
"rebuild swap sidecar rename failed",
);
}
let reopened = match open_with_schema(db_path, schema).await {
Ok(reopened) => reopened,
Err(e) if is_actionable_signal(&e) => return Err(e),
Err(e) => {
crate::boot::boot_diagnostic(format!(
"store '{name}' rebuilt store reopen failed: {e} — original family is \
quarantined for recovery",
));
return Err(e);
}
};
let mut verified = true;
let mut finding_reported = false;
match reopened.quick_check().await {
Ok(()) => {}
Err(e) if is_actionable_signal(&e) => {
warn!(error = %e, db = %name, "post-swap verification quick_check hit an actionable signal");
verified = false;
}
Err(e) => {
crate::boot::boot_diagnostic(format!(
"store '{name}' rebuilt store failed post-swap quick_check ({e}) — \
original family is quarantined for recovery",
));
verified = false;
finding_reported = true;
}
}
for (tbl, expected) in &counts {
match count_table_rows(&reopened, tbl).await {
Err(e) => {
let err = anyhow::anyhow!(e);
if is_actionable_signal(&err) {
verified = false;
warn!(
error = %err,
db = %name,
table = %tbl,
"post-swap counts verification query hit an actionable signal",
);
} else {
crate::boot::boot_diagnostic(format!(
"store '{name}' post-swap verification counts query for table \
'{tbl}' failed ({err}) — verification incomplete; original \
family is quarantined for recovery",
));
verified = false;
finding_reported = true;
}
}
Ok(c) if c != *expected => {
crate::boot::boot_diagnostic(format!(
"store '{name}' post-swap verification failed for table '{tbl}' \
(expected {expected} rows, found {c}) — the swap lost data; \
original family is quarantined for recovery",
));
verified = false;
finding_reported = true;
}
Ok(_) => {}
}
}
if verified {
info!(db = %name, "overflow-aliasing repaired: store rebuilt data-preservingly");
Ok((RepairOutcome::Migrated, reopened))
} else {
if !finding_reported {
warn!(db = %name, "post-swap verification incomplete (transient read errors) — store is in place; re-verified on the next boot");
}
Ok((RepairOutcome::NoRepair, reopened))
}
}
async fn checkpoint_rebuilt_before_swap(
conn: &Connection,
name: &'static str,
root: &Path,
) -> Result<(), MigrateFailure> {
if !crate::db::shrink_gate::shrink_allowed(conn, name, Some(root)).await {
return Ok(());
}
match conn.checkpoint_ungated().await {
Ok(o) if o.is_complete() => Ok(()),
Ok(o) => Err(MigrateFailure::Finding(format!(
"rebuilt store checkpoint incomplete (busy={}, {} of {} frames)",
o.busy, o.checkpointed_frames, o.log_frames,
))),
Err(e) => Err(classify_migrate(e)),
}
}
enum MigrateFailure {
Actionable(anyhow::Error),
Finding(String),
}
fn classify_migrate<E: Into<anyhow::Error>>(e: E) -> MigrateFailure {
let err = e.into();
if is_actionable_signal(&err) {
MigrateFailure::Actionable(err)
} else {
MigrateFailure::Finding(format!("{err:#}"))
}
}
fn row_insert_sql(table_ref: &str, row: &Row) -> Result<(String, Vec<Value>), MigrateFailure> {
let ncols = row.column_count();
let mut vals = Vec::with_capacity(ncols);
for c in 0..ncols {
vals.push(row.get_value(c).map_err(classify_migrate)?);
}
Ok((
format!(
"INSERT INTO \"{table_ref}\" VALUES ({})",
sql_in_placeholders(ncols)
),
vals,
))
}
async fn migrate_schema_and_data(
old: &Connection,
fresh: &Connection,
counts: &[(String, i64)],
) -> Result<(), MigrateFailure> {
let ddl_rows = old
.query(
&format!(
"SELECT sql FROM sqlite_master \
WHERE sql IS NOT NULL \
AND {USER_OBJECT_FILTER} \
ORDER BY rowid"
),
(),
)
.await
.map_err(classify_migrate)?;
for row in &ddl_rows {
let ddl = row
.get::<String>(0)
.map_err(|e| MigrateFailure::Finding(format!("DDL replay row is not text ({e})")))?;
fresh.execute(&ddl, ()).await.map_err(classify_migrate)?;
}
let has_sequence: i64 = old
.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='sqlite_sequence'",
(),
|r| r.get::<i64>(0),
)
.await
.map_err(classify_migrate)?;
if has_sequence > 0 {
let seq_rows = old
.query("SELECT * FROM sqlite_sequence", ())
.await
.map_err(classify_migrate)?;
for row in &seq_rows {
let (sql, vals) = row_insert_sql("sqlite_sequence", row)?;
fresh.execute(&sql, vals).await.map_err(classify_migrate)?;
}
}
for (tbl, expected) in counts {
let quoted = tbl.replace('"', "\"\"");
let rows = old
.query(&format!("SELECT * FROM \"{quoted}\""), ())
.await
.map_err(classify_migrate)?;
let tx = fresh.begin_tx().await.map_err(classify_migrate)?;
let mut copied = 0i64;
for row in &rows {
let (sql, vals) = row_insert_sql("ed, row)?;
tx.execute(&sql, vals).await.map_err(classify_migrate)?;
copied += 1;
}
tx.commit().await.map_err(classify_migrate)?;
if copied != *expected {
return Err(MigrateFailure::Finding(format!(
"table '{tbl}' copy count {copied} != pre-migration {expected}",
)));
}
}
Ok(())
}
async fn drop_create_index_fallback(
conn: &Connection,
name: &str,
index: &str,
quoted: &str,
) -> RepairOutcome {
let ddl = conn
.query_optional(
"SELECT sql FROM sqlite_master WHERE type = 'index' AND name = ?",
(index.to_string(),),
|row| row.get::<String>(0),
)
.await
.ok()
.flatten();
let Some(ddl) = ddl else {
crate::boot::boot_diagnostic(format!(
"store '{name}' index '{index}' DDL not found in sqlite_master — cannot \
DROP+CREATE; left for operator review",
));
return RepairOutcome::NoRepair;
};
let outcome = async {
let tx = conn.begin_tx().await?;
tx.execute_batch(&format!("DROP INDEX \"{quoted}\"; {ddl};"))
.await?;
tx.commit().await
}
.await;
match outcome {
Ok(()) if conn.quick_check().await.is_ok() => {
info!(
db = %name,
index = %index,
"class-B btree index desync repaired in place (DROP+CREATE)",
);
RepairOutcome::Repaired
}
Ok(()) => {
crate::boot::boot_diagnostic(format!(
"store '{name}' DROP+CREATE of '{index}' post-rebuild quick_check \
not clean — left for operator review",
));
RepairOutcome::NoRepair
}
Err(repair_err) => {
crate::boot::boot_diagnostic(format!(
"store '{name}' DROP+CREATE of '{index}' failed: {repair_err} — DROP \
rolled back on the next connection op, original index preserved; left \
for operator review",
));
RepairOutcome::NoRepair
}
}
}
fn desynced_index_name(msg: &str) -> Option<&str> {
const PREFIX: &str = "wrong # of entries in index ";
let rest = msg.strip_prefix(PREFIX)?;
let end = rest.find(['\n', ';', '\r']).unwrap_or(rest.len());
let name = &rest[..end];
(!name.is_empty()).then_some(name)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum CheckpointMode {
Passive,
Truncate,
}
impl CheckpointMode {
fn label(self) -> &'static str {
match self {
Self::Passive => "PASSIVE",
Self::Truncate => "TRUNCATE",
}
}
}
#[must_use]
pub(crate) fn quarantine_store_artifacts(db_path: &Path) -> bool {
quarantine_family(db_path, &[(db_path, ""), (&wal_path(db_path), "-wal")])
}
#[must_use]
fn quarantine_family(db_path: &Path, sources: &[(&Path, &str)]) -> bool {
static QUARANTINE_SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
let stamp = family_stamp();
let seq = QUARANTINE_SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let base = if seq == 0 {
format!(
"{}.quarantine-{stamp}",
db_path.file_name().unwrap_or_default().to_string_lossy()
)
} else {
format!(
"{}.quarantine-{stamp}-{seq}",
db_path.file_name().unwrap_or_default().to_string_lossy()
)
};
let mut complete = true;
for (src, suffix) in sources {
if !src.exists() {
continue;
}
let dst = db_path.with_file_name(format!("{base}{suffix}"));
if let Err(e) = std::fs::rename(src, &dst) {
warn!(
error = %e,
from = %src.display(),
to = %dst.display(),
"Failed to quarantine store artifact",
);
complete = false;
}
}
complete
}
pub(crate) async fn with_tx(
conn: &Connection,
ticket_id: &str,
action_label: &str,
work: impl AsyncFnOnce(&TxGuard<'_>) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
with_tx_outcome(conn, ticket_id, action_label, async |tx| {
work(tx).await?;
Ok(true)
})
.await
.map(|_| ())
}
pub(crate) async fn with_tx_outcome(
conn: &Connection,
ticket_id: &str,
action_label: &str,
work: impl AsyncFnOnce(&TxGuard<'_>) -> anyhow::Result<bool>,
) -> anyhow::Result<bool> {
let tx = conn
.begin_tx()
.await
.map_err(|e| {
warn!(
ticket = %ticket_id,
error = %e,
"Failed to begin transaction for {action_label}",
);
e
})
.with_context(|| format!("Failed to begin transaction for {action_label}"))?;
match work(&tx).await {
Ok(true) => {
tx.commit()
.await
.map_err(|e| {
warn!(
ticket = %ticket_id,
error = %e,
"Failed to commit transaction for {action_label}",
);
e
})
.with_context(|| format!("Failed to commit transaction for {action_label}"))?;
Ok(true)
}
Ok(false) => {
tx.rollback()
.await
.map_err(|e| {
warn!(
ticket = %ticket_id,
error = %e,
"Transaction rollback also failed for {action_label}",
);
e
})
.with_context(|| format!("Failed to roll back transaction for {action_label}"))?;
Ok(false)
}
Err(e) => {
if let Err(rollback_err) = tx.rollback().await {
warn!(
ticket = %ticket_id,
error = %rollback_err,
"Transaction rollback also failed for {action_label}",
);
}
warn!(
ticket = %ticket_id,
error = %e,
"{action_label}: transaction rolled back",
);
Err(e.context(format!("{action_label}: transaction rolled back")))
}
}
}
#[must_use]
fn family_stamp() -> String {
format!(
"{}-{}",
Utc::now().format("%Y%m%dT%H%M%SZ"),
std::process::id()
)
}
#[cfg(test)]
pub(crate) mod test_support {
use std::path::Path;
use std::process::Command;
use anyhow::{Context, Result, ensure};
use super::wal_guard::WAL_HEADER_BYTES;
use super::{Connection, TICKETS_FTS_INDEX_NAME, now, params};
pub(crate) async fn insert_fts_ticket(conn: &Connection, id: &str, title: &str) {
conn.execute(
"INSERT INTO tickets (id, title, description, workspace_name, created_at, updated_at) \
VALUES (?1, ?2, 'desc', 'ws', ?3, ?3)",
params![id.to_string(), title.to_string(), now()],
)
.await
.unwrap();
}
pub(crate) fn fts_corruption_ddl() -> String {
format!(
"DROP INDEX {TICKETS_FTS_INDEX_NAME}; \
CREATE INDEX {TICKETS_FTS_INDEX_NAME} ON tickets(title);"
)
}
pub(crate) fn assert_not_quarantined(dir: &std::path::Path, what: &str) {
let quarantined: Vec<String> = std::fs::read_dir(dir)
.expect("read the store dir")
.filter_map(std::result::Result::ok)
.map(|e| e.file_name().to_string_lossy().into_owned())
.filter(|n| n.contains("quarantine-"))
.collect();
assert!(
quarantined.is_empty(),
"{what} must not be quarantined, found: {quarantined:?}"
);
}
const PAGE_1_DATABASE_SIZE_OFFSET: u64 = 28;
const WAL_FRAME_HEADER_BYTES: usize = 24;
const FIXTURE_SCHEMA: &str = "CREATE TABLE fixture (id INTEGER PRIMARY KEY, v TEXT);";
const FIXTURE_ROW: &str = "INSERT INTO fixture (v) VALUES ('data')";
async fn create_fixture_store(db_path: &Path) {
let conn = super::open_with_schema(db_path, FIXTURE_SCHEMA)
.await
.expect("create the shrink-gate fixture store");
conn.execute(FIXTURE_ROW, ())
.await
.expect("insert the shrink-gate fixture row");
drop(conn);
}
pub(crate) async fn build_zero_page_count_store(db_path: &Path) {
use std::io::{Seek, SeekFrom, Write};
let conn = super::open_with_schema(db_path, FIXTURE_SCHEMA)
.await
.expect("create the shrink-gate fixture store");
conn.execute(FIXTURE_ROW, ())
.await
.expect("insert the shrink-gate fixture row");
conn.checkpoint_ungated()
.await
.expect("truncate-checkpoint the shrink-gate fixture store");
drop(conn);
let mut file = std::fs::OpenOptions::new()
.write(true)
.open(db_path)
.expect("open the shrink-gate fixture store to patch its page-1 header");
file.seek(SeekFrom::Start(PAGE_1_DATABASE_SIZE_OFFSET))
.expect("seek to the page-1 database-size field");
file.write_all(&0u32.to_be_bytes())
.expect("zero the page-1 database-size field");
}
pub(crate) async fn build_empty_main_file_store(db_path: &Path) {
create_fixture_store(db_path).await;
let wal = super::wal_path(db_path);
let wal_bytes = std::fs::metadata(&wal)
.expect("the fixture store must leave a journal")
.len();
assert!(
wal_bytes > WAL_HEADER_BYTES,
"the fixture must leave committed frames in its journal, found {wal_bytes} bytes",
);
std::fs::OpenOptions::new()
.write(true)
.truncate(true)
.open(db_path)
.expect("truncate the fixture's main file");
}
pub(crate) async fn build_zero_page_count_wal_store(db_path: &Path) {
create_fixture_store(db_path).await;
zero_wal_page_1_database_size(db_path);
}
fn zero_wal_page_1_database_size(db_path: &Path) {
let wal = super::wal_path(db_path);
let mut bytes = std::fs::read(&wal).expect("read the fixture journal");
let header_len = usize::try_from(WAL_HEADER_BYTES).expect("the WAL header fits a usize");
let use_native = cfg!(target_endian = "big") == (be_u32(&bytes, 0) & 1 == 1);
let page_size = usize::try_from(be_u32(&bytes, 8)).expect("the WAL page size fits a usize");
let frame_size = WAL_FRAME_HEADER_BYTES + page_size;
let mut seed = (be_u32(&bytes, 24), be_u32(&bytes, 28));
let mut newest_page_1 = None;
let mut offset = header_len;
while offset + frame_size <= bytes.len() && be_u32(&bytes, offset) != 0 {
let computed = frame_checksum(&bytes, offset, frame_size, seed, use_native);
assert_eq!(
computed,
(be_u32(&bytes, offset + 16), be_u32(&bytes, offset + 20)),
"the fixture journal's frame at offset {offset} does not carry the cumulative \
checksum of its bytes — the fixture's checksum algorithm does not match the \
engine's",
);
seed = computed;
if be_u32(&bytes, offset) == 1 {
newest_page_1 = Some(offset);
}
offset += frame_size;
}
let patch_at = newest_page_1.expect("the fixture journal must hold a page-1 image");
let field = patch_at
+ WAL_FRAME_HEADER_BYTES
+ usize::try_from(PAGE_1_DATABASE_SIZE_OFFSET)
.expect("the page-1 database-size offset fits a usize");
bytes[field..field + 4].fill(0);
let mut seed = (be_u32(&bytes, 24), be_u32(&bytes, 28));
let mut offset = header_len;
while offset + frame_size <= bytes.len() && be_u32(&bytes, offset) != 0 {
let computed = frame_checksum(&bytes, offset, frame_size, seed, use_native);
bytes[offset + 16..offset + 20].copy_from_slice(&computed.0.to_be_bytes());
bytes[offset + 20..offset + 24].copy_from_slice(&computed.1.to_be_bytes());
seed = computed;
offset += frame_size;
}
std::fs::write(&wal, &bytes).expect("write the patched fixture journal");
}
fn frame_checksum(
bytes: &[u8],
offset: usize,
frame_size: usize,
seed: (u32, u32),
use_native: bool,
) -> (u32, u32) {
wal_checksum(
&bytes[offset + WAL_FRAME_HEADER_BYTES..offset + frame_size],
wal_checksum(&bytes[offset..offset + 8], seed, use_native),
use_native,
)
}
fn wal_checksum(data: &[u8], seed: (u32, u32), use_native: bool) -> (u32, u32) {
let word = |bytes: &[u8]| {
let mut word = [0u8; 4];
word.copy_from_slice(bytes);
if use_native {
u32::from_ne_bytes(word)
} else {
u32::from_ne_bytes(word).swap_bytes()
}
};
let (mut s1, mut s2) = seed;
let (groups, remainder) = data.as_chunks::<8>();
assert!(
remainder.is_empty(),
"the cumulative WAL checksum is defined over 8-byte groups only",
);
for group in groups {
let v0 = word(&group[0..4]);
let v1 = word(&group[4..8]);
s1 = s1.wrapping_add(v0.wrapping_add(s2));
s2 = s2.wrapping_add(v1.wrapping_add(s1));
}
(s1, s2)
}
fn be_u32(bytes: &[u8], offset: usize) -> u32 {
u32::from_be_bytes(
bytes[offset..offset + 4]
.try_into()
.expect("the field is four bytes"),
)
}
pub(crate) async fn fixture_rows(conn: &Connection) -> i64 {
conn.query_row("SELECT COUNT(*) FROM fixture", (), |r| r.get::<i64>(0))
.await
.expect("count the fixture rows")
}
pub(crate) async fn page_count(conn: &Connection) -> i64 {
conn.query_row("PRAGMA page_count;", (), |r| r.get::<i64>(0))
.await
.expect("ask the fixture store for its page count")
}
pub(crate) const PROBES: [(&str, &str); 2] = [
("read", "SELECT count(*) FROM sqlite_schema;"),
(
"write",
"CREATE TABLE lock_probe (x); INSERT INTO lock_probe VALUES (1);",
),
];
pub(crate) fn sqlite3_available() -> bool {
Command::new("sqlite3")
.arg("--version")
.output()
.is_ok_and(|output| output.status.success())
}
pub(crate) fn probe_store(store: &Path, sql: &str) -> Result<String> {
let output = Command::new("sqlite3")
.arg(store)
.arg(sql)
.output()
.with_context(|| format!("spawning sqlite3 to probe {}", store.display()))?;
let refusal = format!(
"{}{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr),
);
ensure!(
!output.status.success(),
"sqlite3 was NOT refused by {} (exit {}) — the store is unlocked: {}",
store.display(),
output.status,
refusal.trim(),
);
let lower = refusal.to_lowercase();
ensure!(
lower.contains("locked") || lower.contains("busy"),
"sqlite3 was refused by {} without a lock/busy error — the failure is not the lock: {}",
store.display(),
refusal.trim(),
);
Ok(refusal
.lines()
.next()
.unwrap_or_default()
.trim()
.to_string())
}
pub(crate) fn foreign_client_is_refused(store: &Path) {
if !sqlite3_available() {
println!(
"foreign-client probes SKIPPED for {}: no stock sqlite3 on PATH",
store.display(),
);
return;
}
for (kind, sql) in PROBES {
let refusal = probe_store(store, sql).unwrap_or_else(|e| {
panic!("the foreign {kind} probe must be refused as locked: {e:#}")
});
println!(
"foreign {kind} probe on {} refused: {refusal}",
store.display()
);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::AtomicUsize;
use test_support::{assert_not_quarantined, fts_corruption_ddl, insert_fts_ticket};
#[test]
fn experimental_features_are_consistent() {
let opts = experimental_database_opts();
for feature in EXPERIMENTAL_FEATURES {
match *feature {
"index_method" => assert!(
opts.enable_index_method,
"index_method should be enabled per EXPERIMENTAL_FEATURES"
),
other => panic!("unknown experimental feature: {other}"),
}
}
assert!(
!opts.enable_views,
"views is not an active experimental feature"
);
assert!(
!opts.enable_custom_types,
"custom_types is not an active experimental feature"
);
assert!(
!opts.enable_encryption,
"encryption is not an active experimental feature"
);
assert!(
!opts.enable_autovacuum,
"autovacuum is not an active experimental feature"
);
assert!(
!opts.enable_vacuum,
"vacuum is not an active experimental feature"
);
assert!(
!opts.enable_attach,
"attach is not an active experimental feature"
);
assert!(
!opts.enable_generated_columns,
"generated_columns is not an active experimental feature"
);
assert!(
!opts.enable_without_rowid,
"without_rowid is not an active experimental feature"
);
assert!(
!opts.unsafe_testing,
"unsafe_testing is not an active experimental feature"
);
let fam = experimental_database_opts();
assert!(
fam.enable_index_method,
"family reads need index_method (stores are created with it)"
);
}
#[tokio::test]
async fn existing_store_without_product_schema_is_refused() {
let tmp = tempfile::TempDir::new().unwrap();
let root = tmp.path();
let db_path = store_db_path(root, "board");
std::fs::create_dir_all(db_path.parent().unwrap()).unwrap();
{
let _conn = open_with_schema(
&db_path,
"CREATE TABLE IF NOT EXISTS t (id INTEGER PRIMARY KEY);",
)
.await
.unwrap();
}
let before = std::fs::read(&db_path).unwrap();
let err = open_store(root, "board", "")
.await
.expect_err("a store without the product schema must be refused");
let refusal = err
.downcast_ref::<StoreRefusal>()
.unwrap_or_else(|| panic!("expected a StoreRefusal, got: {err:#}"));
assert!(
refusal.reason.contains("product's own schema"),
"the reason must name the missing product schema: {}",
refusal.reason
);
assert_not_quarantined(
db_path.parent().unwrap(),
"a store without the product schema",
);
assert_eq!(
std::fs::read(&db_path).unwrap(),
before,
"the main file must be unchanged byte-for-byte"
);
}
#[test]
fn family_writer_names_parse_via_debug_parser() {
let dir = tempfile::TempDir::new().unwrap();
let db_path = dir.path().join("board.db");
let wal = wal_path(&db_path);
for f in [&db_path, &wal] {
std::fs::write(f, b"x").unwrap();
}
assert!(
quarantine_store_artifacts(&db_path),
"quarantine writer must move every existing source"
);
let moved: Vec<String> = std::fs::read_dir(dir.path())
.unwrap()
.map(|e| e.unwrap().file_name().to_string_lossy().into_owned())
.filter(|n| n.contains("board.db.quarantine-"))
.collect();
assert_eq!(moved.len(), 2, "db + wal moved: {moved:?}");
for name in moved {
let base = name.strip_suffix("-wal").unwrap_or(&name);
let meta = crate::db::debug::parse_family_name(base).unwrap_or_else(|| {
panic!("quarantine writer name must parse as a family id: {name}")
});
assert_eq!(meta.store, "board");
assert_eq!(meta.kind, crate::db::debug::FamilyKind::Quarantine);
}
let snap = pre_reindex_snapshot_path(&db_path);
let meta = crate::db::debug::parse_family_name(snap.file_name().unwrap().to_str().unwrap())
.expect("pre-reindex snapshot name must parse as a family id");
assert_eq!(meta.kind, crate::db::debug::FamilyKind::PreReindex);
}
#[test]
fn builder_mapping_matches_experimental_features() {
let opts = experimental_database_opts();
let mut mapped: Vec<&str> = Vec::new();
if opts.enable_index_method {
mapped.push("index_method");
}
mapped.sort_unstable();
let mut expected: Vec<&str> = EXPERIMENTAL_FEATURES.to_vec();
expected.sort_unstable();
assert_eq!(
mapped, expected,
"Connection::open builder mapping enables features that differ from \
EXPERIMENTAL_FEATURES.\n\
If you added a feature: add the experimental_*() guard above AND \
add it to Connection::open.\n\
If you removed a feature: remove it from both places.\n\
See EXPERIMENTAL_FEATURES docs for naming asymmetries."
);
}
#[test]
fn raw_builder_usage_is_confined_to_persistence_and_debug() {
const PATTERNS: [&str; 2] = ["turso::Builder", "Builder::new_local"];
const ALLOWED: [&str; 2] = ["src/db/mod.rs", "src/db/debug.rs"];
let manifest_dir = std::path::Path::new(env!("CARGO_MANIFEST_DIR"));
let files = crate::util::test::rs_files_under(&manifest_dir.join("src"));
let mut violations: Vec<(String, &'static str)> = Vec::new();
for file in files {
let rel = crate::util::test::rel_source_path(manifest_dir, &file);
if ALLOWED.contains(&rel.as_str()) {
continue;
}
let content = std::fs::read_to_string(&file).expect("read source file");
let code_only: String = content
.lines()
.filter(|line| !line.trim_start().starts_with("//"))
.collect::<Vec<_>>()
.join("\n");
for pattern in PATTERNS {
if code_only.contains(pattern) {
violations.push((rel.clone(), pattern));
}
}
}
assert!(
violations.is_empty(),
"raw turso::Builder usage outside the persistence module and debug CLI.\n\
All database access must go through crate::db::Connection (db/mod.rs);\n\
the debug CLI (db/debug.rs) is the only documented exception.\n\
Violations: {violations:#?}"
);
}
#[test]
fn in_memory_temp_store_applied_on_both_opening_paths() {
const PRAGMA: &str = "PRAGMA temp_store = MEMORY";
const EXPECTED: [(&str, &str); 2] = [
("src/db/mod.rs", "impl Connection"),
("src/db/debug.rs", "fn connect_readonly"),
];
let manifest_dir = std::path::Path::new(env!("CARGO_MANIFEST_DIR"));
let mut violations: Vec<String> = Vec::new();
for (file, carrier) in EXPECTED {
let path = manifest_dir.join(file);
let content =
std::fs::read_to_string(&path).unwrap_or_else(|e| panic!("read {file}: {e}"));
let code_only: String = content
.lines()
.filter(|line| !line.trim_start().starts_with("//"))
.collect::<Vec<_>>()
.join("\n");
if !code_only.contains(PRAGMA) {
violations.push(format!(
"{file} ({carrier}) no longer applies the in-memory temp-store \
setting (missing `{PRAGMA}` in code).\n\
The missing-temp-root guarantee requires EVERY turso connection \
the application opens — service factory AND debug CLI — to run \
with in-memory temp storage."
));
}
}
assert!(
violations.is_empty(),
"in-memory temp-store regression guard failed:\n{}",
violations.join("\n\n")
);
}
#[test]
fn test_sanitize_fts_query() {
let cases = [
("hello world", "hello world"),
("+-~", ""),
("`Hello ${name}`", "Hello $ name"),
(
"contact user@example.com now",
"contact user@example.com now",
),
(
"hello, world! How's it going?",
"hello, world! How s it going?",
),
("!@#$%", "!@#$%"),
("my_function", "my_function"),
("#381", "#381"),
("v1.2.3", "v1.2.3"),
("feature/x", "feature/x"),
("/something", "something"),
("-hello", "hello"),
("hello-world", "hello world"),
("+term", "term"),
("don't", "don t"),
];
for (input, expected) in cases {
assert_eq!(sanitize_fts_query(input), expected, "input: {input:?}");
}
}
#[test]
fn test_parse_utc_timestamp() {
let valid_cases = [
("2024-01-15T10:30:00Z", "2024-01-15T10:30:00+00:00"),
("2024-06-15T14:30:00+05:00", "2024-06-15T09:30:00+00:00"),
("2024-12-25T20:00:00-08:00", "2024-12-26T04:00:00+00:00"),
];
for (input, expected) in valid_cases {
let ts = parse_utc_timestamp(input)
.unwrap_or_else(|e| panic!("parse_utc_timestamp({input:?}) failed: {e}"));
assert_eq!(ts.to_rfc3339(), expected, "input: {input:?}");
}
for invalid in ["garbage", "", "2024-01-15"] {
assert!(
parse_utc_timestamp(invalid).is_err(),
"expected error for: {invalid:?}",
);
}
}
#[tokio::test]
async fn test_checkpoint_reports_complete_outcome() {
let tmp = tempfile::TempDir::new().expect("temp dir for test");
let conn = Connection::open(tmp.path().join("test.db").as_path())
.await
.expect("open test database");
conn.execute(
"CREATE TABLE IF NOT EXISTS _test (id INTEGER PRIMARY KEY, val TEXT NOT NULL)",
(),
)
.await
.expect("create test table");
conn.execute("INSERT INTO _test (id, val) VALUES (1, 'hello')", ())
.await
.expect("insert test row");
for mode in ["checkpoint", "checkpoint_passive"] {
let outcome = conn
.checkpoint_mode(mode == "checkpoint")
.await
.expect("checkpoint should succeed on a healthy database");
assert!(
outcome.is_complete(),
"{mode} outcome must be complete: {outcome:?}"
);
}
}
#[test]
fn checkpoint_outcome_completeness_predicate() {
let complete = CheckpointOutcome {
busy: false,
log_frames: 0,
checkpointed_frames: 0,
};
assert!(complete.is_complete());
let busy = CheckpointOutcome {
busy: true,
log_frames: 0,
checkpointed_frames: 0,
};
assert!(!busy.is_complete(), "busy outcome must be incomplete");
let partial = CheckpointOutcome {
busy: false,
log_frames: 12,
checkpointed_frames: 5,
};
assert!(
!partial.is_complete(),
"uncheckpointed WAL frames must be incomplete"
);
}
#[tokio::test]
async fn test_quick_check_passes_on_healthy_db() {
let tmp = tempfile::TempDir::new().expect("temp dir for test");
let conn = Connection::open(tmp.path().join("test.db").as_path())
.await
.expect("open test database");
conn.quick_check()
.await
.expect("quick_check should pass on a healthy empty database");
}
#[test]
fn known_fts_dir_false_positive_classification() {
let fp = "wrong # of entries in index __turso_internal_fts_dir_idx_tickets_title_fts_key";
assert!(fp.contains(KNOWN_FTS_DIR_COUNT_FALSE_POSITIVE));
let genuine_missing =
"row 5 missing from index __turso_internal_fts_dir_idx_tickets_title_fts_key";
let genuine_unique =
"non-unique entry in index __turso_internal_fts_dir_idx_tickets_title_fts_key";
assert!(!genuine_missing.contains(KNOWN_FTS_DIR_COUNT_FALSE_POSITIVE));
assert!(!genuine_unique.contains(KNOWN_FTS_DIR_COUNT_FALSE_POSITIVE));
}
#[tokio::test]
async fn open_store_refuses_unusable_store_files() {
let mut bad_page = vec![0u8; 4096];
bad_page[..16].copy_from_slice(b"SQLite format 3\0");
let cases = [
("0 bytes", Vec::new(), true),
("64 bytes", vec![0u8; 64], false),
("garbage (bad magic)", vec![0x42; 128], false),
("valid magic, page size 0", bad_page, false),
];
for (name, bytes, with_wal) in cases {
let tmp = tempfile::TempDir::new().expect("temp dir for test");
let root = tmp.path();
let db_path = store_db_path(root, "board");
std::fs::create_dir_all(db_path.parent().unwrap()).unwrap();
std::fs::write(&db_path, &bytes).unwrap();
let wal_bytes = [0x7au8; 512];
if with_wal {
std::fs::write(wal_path(&db_path), wal_bytes).unwrap();
}
let err = open_store(root, "board", "")
.await
.expect_err("an unusable store file must refuse the boot");
assert!(
err.downcast_ref::<StoreRefusal>().is_some(),
"{name}: expected a StoreRefusal, got: {err:#}"
);
assert_eq!(
std::fs::read(&db_path).unwrap(),
bytes,
"{name}: the main file must be unchanged"
);
if with_wal {
assert_eq!(
std::fs::read(wal_path(&db_path)).unwrap(),
wal_bytes,
"{name}: the -wal must be unchanged"
);
}
assert_not_quarantined(db_path.parent().unwrap(), name);
}
let tmp = tempfile::TempDir::new().expect("temp dir for test");
let root = tmp.path();
let db_path = store_db_path(root, "board");
let _conn = open_store(root, "board", "CREATE TABLE IF NOT EXISTS t (id INTEGER);")
.await
.expect("an absent store file is a genuine first launch");
assert!(db_path.exists(), "first launch must create the store file");
}
#[tokio::test]
async fn existing_corrupt_store_is_refused_not_recreated() {
let tmp = tempfile::TempDir::new().expect("temp dir for test");
let root = tmp.path();
{
let conn = open_consolidated_store(root)
.await
.expect("build a real product store");
conn.checkpoint_ungated()
.await
.expect("checkpoint the fixture store");
}
let db_path = store_db_path(root, CONSOLIDATED_DB_NAME);
let bytes = std::fs::read(&db_path).unwrap();
assert!(bytes.len() > 8192, "fixture needs a multi-page core.db");
let mut corrupted = bytes.clone();
corrupted[4096..8192].fill(0); std::fs::write(&db_path, &corrupted).unwrap();
let err = open_consolidated_store(root)
.await
.expect_err("a content-damaged store must be refused");
assert!(
err.downcast_ref::<StoreRefusal>().is_some(),
"expected a StoreRefusal, got: {err:#}"
);
assert_eq!(
std::fs::read(&db_path).unwrap(),
corrupted,
"the refused store's main file must be unchanged"
);
assert_not_quarantined(db_path.parent().unwrap(), "a refused store");
}
#[test]
fn desynced_index_name_parses_quick_check_shapes() {
assert_eq!(
desynced_index_name("wrong # of entries in index idx_tickets_phase"),
Some("idx_tickets_phase")
);
assert_eq!(
desynced_index_name("wrong # of entries in index a; row 5 missing"),
Some("a")
);
assert!(
desynced_index_name(
"wrong # of entries in index __turso_internal_fts_dir_idx_tickets_title_fts_key"
)
.is_some()
);
assert!(desynced_index_name("Page referenced multiple times: page 167772160").is_none());
assert!(desynced_index_name("short read on page 12").is_none());
assert!(desynced_index_name("row 5 missing from index idx_x").is_none());
assert!(desynced_index_name("ok").is_none());
}
#[test]
fn overflow_aliasing_vetoes_reindex_across_rows() {
assert!(matches!(
classify_repair_target(&[
"wrong # of entries in index idx_phase".to_string(),
"Page referenced multiple times: page 167772160".to_string(),
]),
RepairTarget::OverflowAliasing
));
assert!(matches!(
classify_repair_target(&[
"*** in database main ***\nPage 871 referenced multiple times \
(references=[3, 3], page_category=Normal)"
.to_string(),
"wrong # of entries in index idx_t_v".to_string(),
]),
RepairTarget::OverflowAliasing
));
assert!(matches!(
classify_repair_target(&[
"Page referenced multiple times: page 5".to_string(),
"wrong # of entries in index idx_phase".to_string(),
]),
RepairTarget::OverflowAliasing
));
assert!(matches!(
classify_repair_target(&["short read on page 167772160".to_string()]),
RepairTarget::OverflowAliasing
));
assert!(matches!(
classify_repair_target(&[
"wrong # of entries in index idx_phase".to_string(),
"short read on page 12".to_string(),
]),
RepairTarget::Unreadable
));
assert!(matches!(
classify_repair_target(&["short read on page 12".to_string()]),
RepairTarget::Unreadable
));
assert!(matches!(
classify_repair_target(&["wrong # of entries in index idx_phase".to_string()]),
RepairTarget::Index(name) if name == "idx_phase"
));
assert!(matches!(
classify_repair_target(&["row 5 missing from index idx_x".to_string()]),
RepairTarget::Unknown
));
}
#[test]
fn only_unrecognised_integrity_signatures_are_tolerated() {
assert!(is_tolerated_integrity_report(&[
"row 5 missing from index idx_x".to_string(),
]));
assert!(!is_tolerated_integrity_report(&[
"wrong # of entries in index idx_phase".to_string(),
]));
assert!(!is_tolerated_integrity_report(&[
"Invalid page type: page 7".to_string(),
]));
assert!(is_tolerated_integrity_report(&[]));
}
#[tokio::test]
async fn transactional_ddl_rollback_preserves_dropped_index() {
let tmp = tempfile::TempDir::new().unwrap();
let conn = open_with_schema(
&tmp.path().join("t.db"),
"CREATE TABLE t (a TEXT, b TEXT); \
INSERT INTO t VALUES ('1', 'x'), ('1', 'y'); \
CREATE UNIQUE INDEX u ON t(b);",
)
.await
.expect("open test store");
let tx = conn.begin_tx().await.expect("begin tx");
let err = tx
.execute_batch("DROP INDEX u; CREATE UNIQUE INDEX u2 ON t(a);")
.await
.expect_err("CREATE UNIQUE on duplicated data must fail");
assert!(
format!("{err}").to_lowercase().contains("unique"),
"expected a constraint violation, got: {err}"
);
tx.rollback().await.expect("rollback");
{
let tx = conn.begin_tx().await.expect("begin tx");
let err = tx
.execute_batch("DROP INDEX u; CREATE UNIQUE INDEX u2 ON t(a);")
.await
.expect_err("CREATE UNIQUE on duplicated data must fail");
assert!(
format!("{err}").to_lowercase().contains("unique"),
"expected a constraint violation, got: {err}"
);
}
let names: Vec<String> = conn
.query_map_strict(
"SELECT name FROM sqlite_master WHERE type = 'index' AND name IN ('u', 'u2')",
(),
|row| row.get::<String>(0),
)
.await
.expect("query sqlite_master");
assert_eq!(names, vec!["u".to_string()], "original index must survive");
conn.quick_check()
.await
.expect("store must be consistent after the rollback");
let n: i64 = conn
.query_row("SELECT COUNT(*) FROM t INDEXED BY u", (), |row| {
row.get::<i64>(0)
})
.await
.expect("INDEXED BY must resolve");
assert_eq!(n, 2);
}
async fn open_orphan_tx_store() -> (tempfile::TempDir, Connection) {
let tmp = tempfile::TempDir::new().unwrap();
let conn = open_with_schema(&tmp.path().join("t.db"), "CREATE TABLE t (a INTEGER);")
.await
.expect("open test store");
(tmp, conn)
}
async fn orphan_tx_rows(conn: &Connection) -> Vec<i64> {
conn.query_map_strict("SELECT a FROM t ORDER BY a", (), |row| row.get::<i64>(0))
.await
.expect("query t")
}
async fn orphan_tx_autocommit(conn: &Connection) -> bool {
conn.conn
.lock()
.await
.is_autocommit()
.expect("probe autocommit")
}
#[tokio::test]
async fn an_abandoned_transaction_is_rolled_back_before_the_connection_is_reused() {
let (_tmp, conn) = open_orphan_tx_store().await;
{
let tx = conn.begin_tx().await.expect("begin tx");
tx.execute("INSERT INTO t VALUES (1)", ())
.await
.expect("insert in the abandoned tx");
}
assert!(
conn.has_dangling_tx.load(Ordering::SeqCst),
"an abandoned transaction must leave the deferred arm set"
);
conn.execute("INSERT INTO t VALUES (2)", ())
.await
.expect("write after the abandon");
assert_eq!(
orphan_tx_rows(&conn).await,
vec![2],
"the abandoned insert must be rolled back, not committed into the new one"
);
assert!(
orphan_tx_autocommit(&conn).await,
"the connection must be in autocommit before it is reused"
);
}
#[tokio::test]
async fn an_open_transaction_the_arm_never_saw_is_rolled_back_before_reuse() {
let (_tmp, conn) = open_orphan_tx_store().await;
{
let raw = conn.conn.lock().await;
raw.execute("BEGIN", ()).await.expect("raw BEGIN");
raw.execute("INSERT INTO t VALUES (1)", ())
.await
.expect("insert in the orphan tx");
}
assert!(
!conn.has_dangling_tx.load(Ordering::SeqCst),
"the arm must be false — that is the missed window"
);
conn.execute("INSERT INTO t VALUES (2)", ())
.await
.expect("write after the orphan tx");
assert_eq!(
orphan_tx_rows(&conn).await,
vec![2],
"the orphan transaction must be rolled back before the connection is reused"
);
assert!(
orphan_tx_autocommit(&conn).await,
"the connection must be in autocommit before it is reused"
);
}
#[tokio::test]
async fn a_stale_arm_without_an_open_transaction_is_cleared() {
let (_tmp, conn) = open_orphan_tx_store().await;
conn.has_dangling_tx.store(true, Ordering::SeqCst);
conn.execute("INSERT INTO t VALUES (1)", ())
.await
.expect("write after a stale arm");
assert!(
!conn.has_dangling_tx.load(Ordering::SeqCst),
"nothing was open, so the arm must be cleared rather than kept"
);
assert_eq!(orphan_tx_rows(&conn).await, vec![1]);
}
fn synthesize_overflow_aliasing(db_path: &Path) -> u32 {
let file = std::fs::read(db_path).unwrap();
let page_size = u16::from_be_bytes([file[16], file[17]]) as usize;
let local = ((page_size - 12) * 32 / 255) - 23;
let find = |file: &[u8], prefix: &[u8], rowid_serial: u8| -> (usize, u32) {
for i in 0..file.len().saturating_sub(prefix.len() + 12) {
if &file[i..i + prefix.len()] == prefix
&& file[i - 4..i] == [0x04, 0x9f, 0x3b, rowid_serial]
{
let ptr =
u32::from_be_bytes(file[i + local - 4..i + local].try_into().unwrap());
return (i, ptr);
}
}
panic!(
"index cell for {:?} not found",
String::from_utf8_lossy(prefix)
);
};
let (_, shared) = find(&file, b"v04000-", 0x02); let (patch_at, _) = find(&file, b"v00000-", 0x09); let mut file = file;
file[patch_at + local - 4..patch_at + local].copy_from_slice(&shared.to_be_bytes());
std::fs::write(db_path, &file).unwrap();
shared
}
async fn build_aliasing_candidate(
db_path: &Path,
schema: &str,
extra_row: Option<(&str, i64)>,
) {
{
let conn = open_with_schema(db_path, schema).await.unwrap();
for i in 0..5000i32 {
let mut big = format!("v{i:05}-");
big.push_str(&"q".repeat(2000));
conn.execute(
"INSERT INTO t (v, n) VALUES (?1, 1);",
turso::params![big.clone()],
)
.await
.unwrap();
}
if let Some((v, n)) = extra_row {
conn.execute("PRAGMA ignore_check_constraints = ON;", ())
.await
.unwrap();
conn.execute(
"INSERT INTO t (v, n) VALUES (?1, ?2);",
turso::params![v, n],
)
.await
.unwrap();
conn.execute("PRAGMA ignore_check_constraints = OFF;", ())
.await
.unwrap();
}
drop(conn);
}
{
let conn = Connection::open(db_path).await.unwrap();
conn.checkpoint_ungated().await.unwrap();
drop(conn);
}
}
#[ignore = "performs real overflow-page DB surgery on a multi-MB fixture (~4 s); runs only when explicitly invoked"]
#[tokio::test]
async fn overflow_aliasing_rebuild_preserves_fts_store() {
let tmp = tempfile::TempDir::new().unwrap();
let root = tmp.path();
let db_path = store_db_path(root, "board");
std::fs::create_dir_all(db_path.parent().unwrap()).unwrap();
let schema = "CREATE TABLE IF NOT EXISTS t (id INTEGER PRIMARY KEY, v TEXT NOT NULL, \
n INTEGER NOT NULL DEFAULT 1); \
CREATE INDEX IF NOT EXISTS idx_t_v ON t(v); \
CREATE TABLE IF NOT EXISTS ft (title TEXT NOT NULL); \
CREATE INDEX IF NOT EXISTS idx_ft_fts ON ft USING fts (title) \
WITH (tokenizer = 'ngram');";
build_aliasing_candidate(&db_path, schema, None).await;
{
let conn = open_with_schema(&db_path, schema).await.unwrap();
for i in 0..5 {
conn.execute(
"INSERT INTO ft (title) VALUES (?1);",
turso::params![format!("ticket title {i}")],
)
.await
.unwrap();
}
drop(conn);
}
{
let conn = Connection::open(&db_path).await.unwrap();
conn.checkpoint_ungated().await.unwrap();
drop(conn);
}
let shared = synthesize_overflow_aliasing(&db_path);
assert!(shared > 0, "surgery must reference a real overflow page");
let conn = open_and_repair_for_test(root, &db_path, "board", schema)
.await
.expect("overflow-aliasing repair must succeed on an FTS store");
conn.quick_check()
.await
.expect("rebuilt store must pass quick_check");
let count: i64 = conn
.query_row("SELECT COUNT(*) FROM t", (), |r| r.get::<i64>(0))
.await
.unwrap();
assert_eq!(count, 5000, "all rows must survive the rebuild");
let idx: i64 = conn
.query_row("SELECT COUNT(*) FROM t INDEXED BY idx_t_v", (), |r| {
r.get::<i64>(0)
})
.await
.unwrap();
assert_eq!(idx, 5000, "rebuilt index must be valid and complete");
let ft: i64 = conn
.query_row("SELECT COUNT(*) FROM ft", (), |r| r.get::<i64>(0))
.await
.unwrap();
assert_eq!(ft, 5, "FTS table data must survive the rebuild");
let matched: i64 = conn
.query_row(
"SELECT COUNT(*) FROM ft WHERE title MATCH 'title'",
(),
|r| r.get::<i64>(0),
)
.await
.unwrap();
assert_eq!(matched, 5, "rebuilt FTS index must answer MATCH queries");
let quarantined = std::fs::read_dir(db_path.parent().unwrap())
.unwrap()
.filter_map(std::result::Result::ok)
.any(|e| e.file_name().to_string_lossy().contains("quarantine-"));
assert!(
quarantined,
"original family must be quarantined (forensic record)"
);
conn.execute("INSERT INTO t (v, n) VALUES ('post-rebuild', 1);", ())
.await
.unwrap();
conn.quick_check()
.await
.expect("rebuilt store must stay clean after a write");
}
#[ignore = "performs real overflow-page DB surgery on a multi-MB fixture (~4 s); runs only when explicitly invoked"]
#[tokio::test]
async fn overflow_aliasing_rebuild_constraint_finding_aborts() {
let tmp = tempfile::TempDir::new().unwrap();
let root = tmp.path();
let db_path = store_db_path(root, "board");
std::fs::create_dir_all(db_path.parent().unwrap()).unwrap();
let schema = "CREATE TABLE IF NOT EXISTS t (id INTEGER PRIMARY KEY, v TEXT NOT NULL, \
n INTEGER NOT NULL CHECK (n < 50)); \
CREATE INDEX IF NOT EXISTS idx_t_v ON t(v);";
build_aliasing_candidate(&db_path, schema, Some(("v09999x", 99))).await;
synthesize_overflow_aliasing(&db_path);
let conn = open_and_repair_for_test(root, &db_path, "board", schema)
.await
.expect("aborted rebuild must still open the store (report-only)");
let count: i64 = conn
.query_row("SELECT COUNT(*) FROM t", (), |r| r.get::<i64>(0))
.await
.unwrap();
assert_eq!(
count, 5001,
"original data must be preserved, not recreated"
);
assert_not_quarantined(
db_path.parent().unwrap(),
"a constraint finding (the rebuild abort)",
);
}
#[ignore = "performs real overflow-page DB surgery on a multi-MB fixture (~4 s); runs only when explicitly invoked"]
#[tokio::test]
async fn overflow_aliasing_rebuild_preserves_autoincrement_store() {
let tmp = tempfile::TempDir::new().unwrap();
let root = tmp.path();
let db_path = store_db_path(root, "board");
std::fs::create_dir_all(db_path.parent().unwrap()).unwrap();
let schema = "CREATE TABLE IF NOT EXISTS t (id INTEGER PRIMARY KEY AUTOINCREMENT, \
v TEXT NOT NULL, n INTEGER NOT NULL DEFAULT 1); \
CREATE INDEX IF NOT EXISTS idx_t_v ON t(v);";
build_aliasing_candidate(&db_path, schema, None).await;
let shared = synthesize_overflow_aliasing(&db_path);
assert!(shared > 0, "surgery must reference a real overflow page");
{
let conn = Connection::open(&db_path).await.unwrap();
conn.execute("DELETE FROM t WHERE id > 4000", ())
.await
.unwrap();
drop(conn);
}
let conn = open_and_repair_for_test(root, &db_path, "board", schema)
.await
.expect("overflow-aliasing repair must succeed on an AUTOINCREMENT store");
conn.quick_check()
.await
.expect("rebuilt store must pass quick_check");
let count: i64 = conn
.query_row("SELECT COUNT(*) FROM t", (), |r| r.get::<i64>(0))
.await
.unwrap();
assert_eq!(count, 4000, "all surviving rows must be preserved");
let idx: i64 = conn
.query_row("SELECT COUNT(*) FROM t INDEXED BY idx_t_v", (), |r| {
r.get::<i64>(0)
})
.await
.unwrap();
assert_eq!(idx, 4000, "rebuilt index must be valid and complete");
conn.execute("INSERT INTO t (v, n) VALUES ('auto-next', 1);", ())
.await
.unwrap();
let max_id: i64 = conn
.query_row("SELECT MAX(id) FROM t", (), |r| r.get::<i64>(0))
.await
.unwrap();
assert_eq!(
max_id, 5001,
"AUTOINCREMENT must advance past the old watermark, never re-issue ids"
);
let seq: i64 = conn
.query_row(
"SELECT seq FROM sqlite_sequence WHERE name = 't'",
(),
|r| r.get::<i64>(0),
)
.await
.unwrap();
assert_eq!(
seq, 5001,
"sqlite_sequence watermark must be carried over and advanced"
);
conn.quick_check()
.await
.expect("rebuilt store must stay clean after a write");
let quarantined = std::fs::read_dir(db_path.parent().unwrap())
.unwrap()
.filter_map(std::result::Result::ok)
.any(|e| e.file_name().to_string_lossy().contains("quarantine-"));
assert!(
quarantined,
"original family must be quarantined (forensic record)"
);
}
#[test]
fn open_time_lock_error_predicate() {
let lock = |msg: &str| turso::Error::Error(msg.to_string());
assert!(is_open_time_lock_error(&lock(
"Locking error: Failed locking file '/x/logs.db-wal'. File is locked by another process"
)));
assert!(is_open_time_lock_error(&lock(
"Locking error: Failed locking file. File is locked by another process"
)));
assert!(is_open_time_lock_error(&lock(
"Locking error: Failed locking shared WAL coordination file. File is locked by another process"
)));
assert!(is_open_time_lock_error(&lock(
"Locking error: Failed locking file, The process cannot access the file because another process has locked a portion of the file. (os error 33)"
)));
assert!(!is_open_time_lock_error(&turso::Error::Busy(
"database is locked".into()
)));
assert!(!is_open_time_lock_error(&lock(
"Runtime error: database table is locked"
)));
assert!(!is_open_time_lock_error(&lock(
"Locking error: mmap shared WAL coordination file failed: boom (offset=0)"
)));
assert!(!is_open_time_lock_error(&lock(
"Locking error: Failed to release file lock: not locked"
)));
assert!(!is_open_time_lock_error(&lock(
"Locking error: Failed locking file '/x/core.db', Too many open files (os error 24)"
)));
assert!(!is_open_time_lock_error(&turso::Error::IoError(
std::io::ErrorKind::WouldBlock,
"read"
)));
assert!(!is_open_time_lock_error(&lock(
"Database file is corrupted"
)));
}
#[test]
fn lock_and_resource_signals_are_actionable() {
for msg in [
"Locking error: File is locked by another process",
"Runtime error: database table is locked",
"Database is busy",
"no space left on device",
] {
let err = anyhow::Error::new(turso::Error::Error(msg.to_string()));
assert!(
is_actionable_signal(&err),
"an external condition must propagate, not recreate: {msg}"
);
}
let corrupt = anyhow::Error::new(turso::Error::Error("Invalid page type".to_string()));
assert!(
!is_actionable_signal(&corrupt),
"corruption must not classify as an actionable signal"
);
}
#[tokio::test]
async fn engine_snapshot_is_queryable_and_leaves_the_source_usable() {
let tmp = tempfile::TempDir::new().unwrap();
let db_path = tmp.path().join("core.db");
let conn = open_with_schema(&db_path, "CREATE TABLE t (id INTEGER PRIMARY KEY, v TEXT);")
.await
.unwrap();
conn.execute("INSERT INTO t (v) VALUES ('kept')", ())
.await
.unwrap();
let snap = snapshot_store_via_engine(&conn, &db_path).await.unwrap();
assert!(
snap.exists(),
"the snapshot must exist at {}",
snap.display()
);
let wal = std::path::PathBuf::from(format!("{}-wal", snap.display()));
assert!(
wal.exists(),
"VACUUM INTO must leave {} beside the snapshot",
wal.display()
);
let reopened = open_with_schema(&snap, "").await.unwrap();
let v: String = reopened
.query_row("SELECT v FROM t WHERE id = 1", (), |r| r.get(0))
.await
.unwrap();
assert_eq!(v, "kept", "the snapshot must carry the committed row");
drop(reopened);
let rows: i64 = conn
.query_row("SELECT COUNT(*) FROM t", (), |r| r.get(0))
.await
.unwrap();
assert_eq!(rows, 1, "the source store must stay queryable");
}
#[tokio::test]
async fn open_lock_retry_succeeds_after_transient_lock() {
let calls = AtomicUsize::new(0);
let result = with_open_lock_retry(3, Duration::from_millis(1), || {
let n = calls.fetch_add(1, Ordering::SeqCst);
async move {
if n < 2 {
Err(turso::Error::Error(
"Locking error: Failed locking file '/x/logs.db-wal'. File is locked by another process".to_string(),
))
} else {
Ok("opened")
}
}
})
.await;
assert!(
result.is_ok(),
"open must succeed after the transient lock clears"
);
assert_eq!(
calls.load(Ordering::SeqCst),
3,
"closure must be called twice before success"
);
}
#[tokio::test]
async fn open_lock_retry_exhausts_and_preserves_error() {
let calls = AtomicUsize::new(0);
let err = with_open_lock_retry(2, Duration::from_millis(1), || {
calls.fetch_add(1, Ordering::SeqCst);
async {
Err::<(), turso::Error>(turso::Error::Error(
"Locking error: Failed locking file '/x/logs.db-wal'. File is locked by another process".to_string(),
))
}
})
.await
.expect_err("a persistent open-time lock must eventually surface");
assert_eq!(
calls.load(Ordering::SeqCst),
3,
"max total tries must be attempts + 1"
);
assert!(
is_open_time_lock_error(&err),
"the original lock error must be preserved"
);
}
#[tokio::test]
async fn open_lock_no_retry_on_non_lock_error() {
let calls = AtomicUsize::new(0);
let err = with_open_lock_retry(2, Duration::from_millis(1), || {
calls.fetch_add(1, Ordering::SeqCst);
async { Err::<(), turso::Error>(turso::Error::Error("Disk I/O error".to_string())) }
})
.await
.expect_err("a non-lock open failure must not be retried");
assert_eq!(
calls.load(Ordering::SeqCst),
1,
"non-lock errors must not trigger retries"
);
assert!(
matches!(&err, turso::Error::Error(msg) if msg == "Disk I/O error"),
"the original non-lock error must be returned"
);
}
#[test]
fn names_ticket_title_fts_classification() {
let index = "idx_tickets_title_fts";
assert!(names_ticket_title_fts(
"wrong # of entries in index idx_tickets_title_fts",
index
));
assert!(names_ticket_title_fts(
"row 5 missing from index __turso_internal_fts_dir_3_key",
index
));
assert!(names_ticket_title_fts(
"wrong # of entries in index __turso_internal_fts_dir_3_key",
index
));
assert!(!names_ticket_title_fts(
"wrong # of entries in index idx_sessions_created",
index
));
assert!(!names_ticket_title_fts(
"cell_index_read_payload_ptr called on non-index page",
index
));
}
#[tokio::test]
async fn is_fts_index_classifies_index_method() {
let tmp = tempfile::TempDir::new().unwrap();
let db_path = tmp.path().join("test.db");
let conn = open_with_schema(
&db_path,
"CREATE TABLE IF NOT EXISTS ft (title TEXT NOT NULL); \
CREATE INDEX IF NOT EXISTS idx_ft_fts ON ft USING fts (title) \
WITH (tokenizer = 'ngram'); \
CREATE TABLE IF NOT EXISTS t (x INTEGER); \
CREATE INDEX IF NOT EXISTS idx_plain ON t(x);",
)
.await
.unwrap();
assert!(is_fts_index(&conn, "idx_ft_fts").await);
assert!(!is_fts_index(&conn, "idx_plain").await);
assert!(!is_fts_index(&conn, "does_not_exist").await);
}
#[tokio::test]
async fn detect_healthy_title_fts_is_none() {
let tmp = tempfile::TempDir::new().unwrap();
let conn = open_consolidated_store(tmp.path()).await.unwrap();
insert_fts_ticket(&conn, "t-1", "Important bug fix").await;
let evidence = detect_ticket_title_fts_corruption(&conn, "idx_tickets_title_fts")
.await
.unwrap();
assert!(
evidence.is_none(),
"healthy FTS index must not report corruption: {evidence:?}"
);
}
#[tokio::test]
async fn rebuild_ticket_title_fts_repairs_and_verifies() {
let tmp = tempfile::TempDir::new().unwrap();
let conn = open_consolidated_store(tmp.path()).await.unwrap();
insert_fts_ticket(&conn, "t-1", "Important bug fix one").await;
insert_fts_ticket(&conn, "t-2", "Another relevant thing").await;
let before: i64 = conn
.query_row("SELECT COUNT(*) FROM tickets", (), |r| r.get::<i64>(0))
.await
.unwrap();
assert_eq!(before, 2);
rebuild_ticket_title_fts(
&conn,
"idx_tickets_title_fts",
"CREATE INDEX IF NOT EXISTS idx_tickets_title_fts ON tickets \
USING fts (title) WITH (tokenizer = 'ngram')",
true,
)
.await;
let evidence = detect_ticket_title_fts_corruption(&conn, "idx_tickets_title_fts")
.await
.unwrap();
assert!(
evidence.is_none(),
"rebuilt FTS index must be healthy: {evidence:?}"
);
let sanitized = sanitize_fts_query("Important bug fix one");
let matched: String = conn
.query_row(
"SELECT id FROM tickets WHERE title MATCH ?1 LIMIT 1",
params![sanitized],
|r| r.get::<String>(0),
)
.await
.unwrap();
assert_eq!(matched, "t-1");
let after: i64 = conn
.query_row("SELECT COUNT(*) FROM tickets", (), |r| r.get::<i64>(0))
.await
.unwrap();
assert_eq!(after, 2, "ticket rows must be untouched by the FTS rebuild");
let present = conn
.query_optional(
"SELECT 1 FROM sqlite_master WHERE type = 'index' AND name = ?1",
params!["idx_tickets_title_fts".to_string()],
|_| Ok::<_, ::turso::Error>(()),
)
.await
.unwrap();
assert!(present.is_some());
}
#[tokio::test]
async fn broken_title_fts_store_is_detected_and_repaired_end_to_end() {
let tmp = tempfile::TempDir::new().unwrap();
let conn = open_consolidated_store(tmp.path()).await.unwrap();
insert_fts_ticket(&conn, "t-1", "Important bug fix one").await;
conn.execute_batch(&fts_corruption_ddl()).await.unwrap();
assert!(!is_fts_index(&conn, "idx_tickets_title_fts").await);
repair_ticket_title_fts_if_corrupt(
&conn,
"idx_tickets_title_fts",
"CREATE INDEX IF NOT EXISTS idx_tickets_title_fts ON tickets \
USING fts (title) WITH (tokenizer = 'ngram')",
"core",
Some(tmp.path()),
)
.await;
assert!(is_fts_index(&conn, "idx_tickets_title_fts").await);
let matched: String = conn
.query_row(
"SELECT id FROM tickets WHERE title MATCH ?1 LIMIT 1",
params![sanitize_fts_query("Important bug fix one")],
|r| r.get::<String>(0),
)
.await
.unwrap();
assert_eq!(matched, "t-1");
assert!(
detect_ticket_title_fts_corruption(&conn, "idx_tickets_title_fts")
.await
.unwrap()
.is_none()
);
let count: i64 = conn
.query_row("SELECT COUNT(*) FROM tickets", (), |r| r.get::<i64>(0))
.await
.unwrap();
assert_eq!(count, 1);
}
#[tokio::test]
async fn title_fts_match_probe_ok_on_healthy_store() {
let tmp = tempfile::TempDir::new().unwrap();
let conn = open_consolidated_store(tmp.path()).await.unwrap();
title_fts_match_probe(&conn)
.await
.expect("probe must pass on an empty FTS index");
insert_fts_ticket(&conn, "t-1", "Important bug fix").await;
title_fts_match_probe(&conn)
.await
.expect("probe must pass on a healthy populated FTS index");
}
#[tokio::test]
async fn runtime_repair_classifies_missing_and_healthy_index() {
let tmp = tempfile::TempDir::new().unwrap();
let plain_path = tmp.path().join("plain.db");
let plain = open_with_schema(&plain_path, "CREATE TABLE plain (id INTEGER PRIMARY KEY);")
.await
.unwrap();
assert!(matches!(
repair_ticket_title_fts_runtime(&plain).await,
TicketTitleFtsRuntimeRepair::NotApplicable
));
let conn = open_consolidated_store(tmp.path()).await.unwrap();
insert_fts_ticket(&conn, "t-1", "Important bug fix").await;
assert!(matches!(
repair_ticket_title_fts_runtime(&conn).await,
TicketTitleFtsRuntimeRepair::Healthy
));
}
}