Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
71 changes: 70 additions & 1 deletion crates/ironclaw_filesystem/src/backend.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,10 @@
use async_trait::async_trait;
use ironclaw_host_api::VirtualPath;

use crate::{CasExpectation, Entry, FilesystemError, RecordVersion, SeqNo, VersionedEntry};
use crate::{
CasExpectation, Entry, FilesystemError, FilesystemOperation, RecordVersion, SeqNo,
VersionedEntry,
};

/// Multi-key transactional handle returned by [`RootFilesystem::begin`].
///
Expand All @@ -40,6 +43,13 @@ pub trait StorageTxn: Send {

async fn delete(&mut self, path: &VirtualPath) -> Result<(), FilesystemError>;

async fn reserve_sequence(&mut self, path: &VirtualPath) -> Result<SeqNo, FilesystemError> {
Err(FilesystemError::Unsupported {
path: path.clone(),
operation: FilesystemOperation::ReserveSeq,
})
}

async fn commit(self: Box<Self>) -> Result<(), FilesystemError>;

async fn rollback(self: Box<Self>);
Expand All @@ -60,3 +70,62 @@ impl EventRecord {
Self { seq, payload }
}
}

#[cfg(test)]
mod tests {
use super::*;

struct DummyTxn;

#[async_trait]
impl StorageTxn for DummyTxn {
async fn put(
&mut self,
path: &VirtualPath,
_entry: Entry,
_cas: CasExpectation,
) -> Result<RecordVersion, FilesystemError> {
Err(FilesystemError::Unsupported {
path: path.clone(),
operation: FilesystemOperation::WriteFile,
})
}

async fn get(
&mut self,
path: &VirtualPath,
) -> Result<Option<VersionedEntry>, FilesystemError> {
Err(FilesystemError::Unsupported {
path: path.clone(),
operation: FilesystemOperation::ReadFile,
})
}

async fn delete(&mut self, path: &VirtualPath) -> Result<(), FilesystemError> {
Err(FilesystemError::Unsupported {
path: path.clone(),
operation: FilesystemOperation::Delete,
})
}

async fn commit(self: Box<Self>) -> Result<(), FilesystemError> {
Ok(())
}

async fn rollback(self: Box<Self>) {}
}

#[tokio::test]
async fn storage_txn_reserve_sequence_fails_closed_by_default() {
let path = VirtualPath::new("/events/log").unwrap();
let mut txn = DummyTxn;

let err = txn.reserve_sequence(&path).await.unwrap_err();

assert!(matches!(
err,
FilesystemError::Unsupported { path: actual, operation: FilesystemOperation::ReserveSeq }
if actual == path
));
}
}
70 changes: 60 additions & 10 deletions crates/ironclaw_filesystem/src/libsql.rs
Original file line number Diff line number Diff line change
Expand Up @@ -141,13 +141,25 @@ impl LibSqlRootFilesystem {
}

#[cfg(feature = "libsql")]
async fn connect_with_retry<F>(mut open: F) -> Result<libsql::Connection, FilesystemError>
async fn connect_with_retry<F>(open: F) -> Result<libsql::Connection, FilesystemError>
where
F: FnMut() -> Result<libsql::Connection, libsql::Error>,
{
connect_with_retry_and_pragmas(open, |_| LIBSQL_CONNECTION_PRAGMAS).await
}

#[cfg(feature = "libsql")]
async fn connect_with_retry_and_pragmas<F, P>(
mut open: F,
mut pragmas_for_attempt: P,
) -> Result<libsql::Connection, FilesystemError>
where
F: FnMut() -> Result<libsql::Connection, libsql::Error>,
P: FnMut(u32) -> &'static str,
{
// Match the legacy libSQL backend's connection policy: every
// operation gets its own connection, concurrent writers wait on
// SQLite locks, and transient file-open races get a short retry
// SQLite locks, and transient file-open/setup races get a short retry
// budget before surfacing as infrastructure errors.
let mut last_error = None;
for attempt in 0..LIBSQL_CONNECT_ATTEMPTS {
Expand All @@ -157,12 +169,15 @@ where
// `execute_batch` runs each statement and discards the rows
// PRAGMAs like `busy_timeout` return, which is exactly what
// we want — we only care about the side effect.
conn.execute_batch(LIBSQL_CONNECTION_PRAGMAS)
.await
.map_err(|error| {
infrastructure_libsql_error(FilesystemOperation::Stat, error)
})?;
return Ok(conn);
match conn.execute_batch(pragmas_for_attempt(attempt)).await {
Ok(_) => return Ok(conn),
Err(error) => {
last_error = Some(error);
if attempt + 1 < LIBSQL_CONNECT_ATTEMPTS {
tokio::time::sleep(connect_backoff(attempt)).await;
}
}
}
}
Err(error) => {
last_error = Some(error);
Expand All @@ -176,11 +191,13 @@ where
let reason = match last_error {
Some(error) => {
format!(
"failed to create libSQL connection after {LIBSQL_CONNECT_ATTEMPTS} attempts: {error}"
"failed to create or initialize libSQL connection after {LIBSQL_CONNECT_ATTEMPTS} attempts: {error}"
)
}
None => {
format!("failed to create libSQL connection after {LIBSQL_CONNECT_ATTEMPTS} attempts")
format!(
"failed to create or initialize libSQL connection after {LIBSQL_CONNECT_ATTEMPTS} attempts"
)
}
};
Err(crate::db::infrastructure_error(
Expand Down Expand Up @@ -2137,6 +2154,39 @@ mod tests {
assert_eq!(timeout, 5000);
}

#[tokio::test]
async fn connect_retries_transient_pragma_failures_before_succeeding() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("connect-retry-pragma-test.db");
let db = libsql::Builder::new_local(db_path).build().await.unwrap();
let mut opens = 0;
let mut initializers = 0;

let conn = connect_with_retry_and_pragmas(
|| {
opens += 1;
db.connect()
},
|_| {
initializers += 1;
if initializers == 1 {
"THIS IS NOT SQL"
} else {
LIBSQL_CONNECTION_PRAGMAS
}
},
)
.await
.unwrap();

assert_eq!(opens, 2);
assert_eq!(initializers, 2);
let mut rows = conn.query("PRAGMA busy_timeout", ()).await.unwrap();
let row = rows.next().await.unwrap().unwrap();
let timeout: i64 = row.get(0).unwrap();
assert_eq!(timeout, 5000);
}

/// `run_migrations` must switch the database into WAL journaling, which
/// is the property that lets readers run concurrently with the single
/// writer instead of serialising behind a whole-file EXCLUSIVE lock.
Expand Down
Loading
Loading