diff --git a/crates/ironclaw_filesystem/src/backend.rs b/crates/ironclaw_filesystem/src/backend.rs index cc44c863b3d..ea61ff75bb6 100644 --- a/crates/ironclaw_filesystem/src/backend.rs +++ b/crates/ironclaw_filesystem/src/backend.rs @@ -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`]. /// @@ -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 { + Err(FilesystemError::Unsupported { + path: path.clone(), + operation: FilesystemOperation::ReserveSeq, + }) + } + async fn commit(self: Box) -> Result<(), FilesystemError>; async fn rollback(self: Box); @@ -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 { + Err(FilesystemError::Unsupported { + path: path.clone(), + operation: FilesystemOperation::WriteFile, + }) + } + + async fn get( + &mut self, + path: &VirtualPath, + ) -> Result, 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) -> Result<(), FilesystemError> { + Ok(()) + } + + async fn rollback(self: Box) {} + } + + #[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 + )); + } +} diff --git a/crates/ironclaw_filesystem/src/libsql.rs b/crates/ironclaw_filesystem/src/libsql.rs index 9e5b046a0e8..80fe6fce2b1 100644 --- a/crates/ironclaw_filesystem/src/libsql.rs +++ b/crates/ironclaw_filesystem/src/libsql.rs @@ -141,13 +141,25 @@ impl LibSqlRootFilesystem { } #[cfg(feature = "libsql")] -async fn connect_with_retry(mut open: F) -> Result +async fn connect_with_retry(open: F) -> Result where F: FnMut() -> Result, +{ + connect_with_retry_and_pragmas(open, |_| LIBSQL_CONNECTION_PRAGMAS).await +} + +#[cfg(feature = "libsql")] +async fn connect_with_retry_and_pragmas( + mut open: F, + mut pragmas_for_attempt: P, +) -> Result +where + F: FnMut() -> Result, + 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 { @@ -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); @@ -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( @@ -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. diff --git a/crates/ironclaw_filesystem/src/postgres.rs b/crates/ironclaw_filesystem/src/postgres.rs index 4604e9e7d26..898f2d49f74 100644 --- a/crates/ironclaw_filesystem/src/postgres.rs +++ b/crates/ironclaw_filesystem/src/postgres.rs @@ -1,4 +1,6 @@ use std::{collections::BTreeMap, error::Error, time::Duration}; +#[cfg(feature = "postgres")] +use std::{collections::HashSet, sync::OnceLock}; use async_trait::async_trait; use ironclaw_host_api::VirtualPath; @@ -33,6 +35,11 @@ const POSTGRES_MIGRATION_CONNECT_DEFAULT_MAX_WAIT: Duration = Duration::from_sec const POSTGRES_MIGRATION_CONNECT_INITIAL_BACKOFF: Duration = Duration::from_millis(250); #[cfg(feature = "postgres")] const POSTGRES_MIGRATION_CONNECT_MAX_BACKOFF: Duration = Duration::from_secs(10); +#[cfg(feature = "postgres")] +const POSTGRES_ROOT_FILESYSTEM_MIGRATION_ADVISORY_LOCK: i64 = 824_917_203; +#[cfg(feature = "postgres")] +static POSTGRES_ROOT_FILESYSTEM_MIGRATED_SCHEMAS: OnceLock>> = + OnceLock::new(); #[cfg(feature = "postgres")] impl PostgresRootFilesystem { @@ -41,11 +48,36 @@ impl PostgresRootFilesystem { } pub async fn run_migrations(&self) -> Result<(), FilesystemError> { - let client = self.migration_client_with_retry().await?; - client + let mut client = self.migration_client_with_retry().await?; + let migration_key = postgres_root_filesystem_migration_key(&client).await?; + let registry = POSTGRES_ROOT_FILESYSTEM_MIGRATED_SCHEMAS + .get_or_init(|| tokio::sync::Mutex::new(HashSet::new())); + let mut migrated_schemas = registry.lock().await; + if migrated_schemas.contains(&migration_key) { + return Ok(()); + } + let transaction = client + .transaction() + .await + .map_err(|error| infrastructure_pg_error(FilesystemOperation::CreateDirAll, error))?; + transaction + .execute( + "SELECT pg_advisory_xact_lock($1)", + &[&POSTGRES_ROOT_FILESYSTEM_MIGRATION_ADVISORY_LOCK], + ) + .await + .map_err(|error| infrastructure_pg_error(FilesystemOperation::CreateDirAll, error))?; + transaction .batch_execute(POSTGRES_ROOT_FILESYSTEM_SCHEMA) .await - .map_err(|error| infrastructure_pg_error(FilesystemOperation::CreateDirAll, error)) + .map_err(|error| infrastructure_pg_error(FilesystemOperation::CreateDirAll, error))?; + drop_legacy_projection_indexes(&transaction).await?; + transaction + .commit() + .await + .map_err(|error| infrastructure_pg_error(FilesystemOperation::CreateDirAll, error))?; + migrated_schemas.insert(migration_key); + Ok(()) } async fn migration_client_with_retry( @@ -96,6 +128,61 @@ impl PostgresRootFilesystem { } } +#[cfg(feature = "postgres")] +async fn drop_legacy_projection_indexes( + transaction: &tokio_postgres::Transaction<'_>, +) -> Result<(), FilesystemError> { + let rows = transaction + .query( + "SELECT indexname \ + FROM pg_indexes \ + WHERE schemaname = current_schema() \ + AND tablename = 'root_filesystem_entries' \ + AND indexname LIKE 'idx_rfs_%' \ + AND indexname NOT LIKE 'idx_rfs_shared_%' \ + AND indexdef LIKE '%btree (((indexed ->>%'", + &[], + ) + .await + .map_err(|error| infrastructure_pg_error(FilesystemOperation::CreateDirAll, error))?; + for row in rows { + let index_name: String = row.get(0); + let quoted = quote_postgres_identifier(&index_name); + transaction + .batch_execute(&format!("DROP INDEX IF EXISTS {quoted}")) + .await + .map_err(|error| infrastructure_pg_error(FilesystemOperation::CreateDirAll, error))?; + } + Ok(()) +} + +#[cfg(feature = "postgres")] +fn quote_postgres_identifier(identifier: &str) -> String { + format!("\"{}\"", identifier.replace('"', "\"\"")) +} + +#[cfg(feature = "postgres")] +async fn postgres_root_filesystem_migration_key( + client: &deadpool_postgres::Object, +) -> Result { + let row = client + .query_one( + "SELECT \ + current_database(), \ + current_schema(), \ + COALESCE(inet_server_addr()::text, 'local'), \ + COALESCE(inet_server_port()::text, 'local')", + &[], + ) + .await + .map_err(|error| infrastructure_pg_error(FilesystemOperation::Connect, error))?; + let database: String = row.get(0); + let schema: String = row.get(1); + let host: String = row.get(2); + let port: String = row.get(3); + Ok(format!("{host}:{port}/{database}/{schema}")) +} + #[cfg(feature = "postgres")] fn postgres_migration_connect_backoff(attempt: u32) -> Duration { POSTGRES_MIGRATION_CONNECT_INITIAL_BACKOFF @@ -248,14 +335,17 @@ impl RootFilesystem for PostgresRootFilesystem { let index_name = sql_index_name(path.as_str(), spec.name.as_str()); match &spec.kind { IndexKind::Exact | IndexKind::Prefix => { + let index_name = postgres_shared_projection_index_name(spec); let expressions: Vec = spec .keys .iter() .map(|k| format!("((indexed->>'{}'))", k.as_str())) .collect(); + let mut columns = expressions; + columns.push("path".to_string()); let ddl = format!( "CREATE INDEX IF NOT EXISTS {index_name} ON root_filesystem_entries ({})", - expressions.join(", ") + columns.join(", ") ); client.batch_execute(&ddl).await.map_err(|error| { db_error(path.clone(), FilesystemOperation::EnsureIndex, error) @@ -373,8 +463,12 @@ impl RootFilesystem for PostgresRootFilesystem { let client = self.client().await?; let params_ref: Vec<&(dyn tokio_postgres::types::ToSql + Sync)> = params.iter().map(|p| p.as_ref() as _).collect(); + let statement = client + .prepare_cached(sql.as_str()) + .await + .map_err(|error| db_error(path.clone(), FilesystemOperation::Query, error))?; let rows = client - .query(sql.as_str(), ¶ms_ref[..]) + .query(&statement, ¶ms_ref[..]) .await .map_err(|error| db_error(path.clone(), FilesystemOperation::Query, error))?; rows.into_iter() @@ -564,27 +658,7 @@ impl RootFilesystem for PostgresRootFilesystem { async fn stat(&self, path: &VirtualPath) -> Result { let client = self.client().await?; - if let Some((len, file_type, modified)) = - self.exact_entry_with_client(&client, path).await? - { - return Ok(FileStat { - path: path.clone(), - file_type, - len, - modified, - sensitive: false, - }); - } - if self.has_child_entry_with_client(&client, path).await? { - return Ok(FileStat { - path: path.clone(), - file_type: FileType::Directory, - len: 0, - modified: None, - sensitive: false, - }); - } - Err(not_found(path.clone(), FilesystemOperation::Stat)) + postgres_stat_with_client(&client, path).await } async fn delete(&self, path: &VirtualPath) -> Result<(), FilesystemError> { @@ -750,22 +824,7 @@ impl RootFilesystem for PostgresRootFilesystem { async fn reserve_sequence(&self, path: &VirtualPath) -> Result { let client = self.client().await?; - let row = cached_query_one( - &client, - r#" - INSERT INTO root_filesystem_sequences (path, next_seq, updated_at) - VALUES ($1, 2, NOW()) - ON CONFLICT (path) DO UPDATE SET - next_seq = root_filesystem_sequences.next_seq + 1, - updated_at = NOW() - RETURNING next_seq - 1 AS reserved - "#, - &[&path.as_str()], - ) - .await - .map_err(|error| db_error(path.clone(), FilesystemOperation::ReserveSeq, error))?; - let reserved: i64 = row.get("reserved"); - seq_no_from_i64(path, reserved, FilesystemOperation::ReserveSeq) + postgres_reserve_sequence_with_client(&client, path).await } async fn create_dir_all(&self, path: &VirtualPath) -> Result<(), FilesystemError> { @@ -1052,6 +1111,11 @@ impl StorageTxn for PostgresStorageTxn { postgres_delete_with_client(self.client()?, path).await } + async fn reserve_sequence(&mut self, path: &VirtualPath) -> Result { + self.check_path(path)?; + postgres_reserve_sequence_with_client(self.client()?, path).await + } + async fn commit(mut self: Box) -> Result<(), FilesystemError> { let client = self.client.take().ok_or_else(|| FilesystemError::Backend { path: self.prefix.clone(), @@ -1082,6 +1146,29 @@ impl StorageTxn for PostgresStorageTxn { } } +#[cfg(feature = "postgres")] +async fn postgres_reserve_sequence_with_client( + client: &deadpool_postgres::Object, + path: &VirtualPath, +) -> Result { + let row = cached_query_one( + client, + r#" + INSERT INTO root_filesystem_sequences (path, next_seq, updated_at) + VALUES ($1, 2, NOW()) + ON CONFLICT (path) DO UPDATE SET + next_seq = root_filesystem_sequences.next_seq + 1, + updated_at = NOW() + RETURNING next_seq - 1 AS reserved + "#, + &[&path.as_str()], + ) + .await + .map_err(|error| db_error(path.clone(), FilesystemOperation::ReserveSeq, error))?; + let reserved: i64 = row.get("reserved"); + seq_no_from_i64(path, reserved, FilesystemOperation::ReserveSeq) +} + #[cfg(feature = "postgres")] impl Drop for PostgresStorageTxn { fn drop(&mut self) { @@ -1105,10 +1192,12 @@ impl Drop for PostgresStorageTxn { /// what keeps the small hosted pool from starving the heartbeat/webui and /// wedging the runner lease. /// -/// Only use these with *static* SQL — dynamic SQL would grow the cache -/// unbounded, so the filter `query` and index DDL paths stay on the uncached -/// `tokio_postgres` calls. The error type stays `tokio_postgres::Error` so -/// existing `db_error` mapping at call sites is unchanged. +/// Prefer these with static or low-cardinality query-shape SQL. Filesystem +/// `query` SQL is generated from the filter shape and indexed key names while +/// paths and values stay bound parameters, so it uses `prepare_cached` directly +/// at the call site. Index DDL remains uncached. The error type stays +/// `tokio_postgres::Error` so existing `db_error` mapping at call sites is +/// unchanged. #[cfg(feature = "postgres")] async fn cached_query_opt( client: &deadpool_postgres::Object, @@ -1414,6 +1503,74 @@ async fn postgres_get_with_client( })) } +#[cfg(feature = "postgres")] +async fn postgres_stat_with_client( + client: &deadpool_postgres::Object, + path: &VirtualPath, +) -> Result { + let (prefix_lower, prefix_upper) = descendant_path_range(path); + let row = cached_query_opt( + client, + r#" + WITH exact AS ( + SELECT OCTET_LENGTH(contents)::bigint AS len, + is_dir, + EXTRACT(EPOCH FROM updated_at)::bigint AS updated_at_epoch + FROM root_filesystem_entries + WHERE path = $1 + ), + child AS ( + SELECT 1 + FROM root_filesystem_entries + WHERE path >= $2 AND path < $3 + LIMIT 1 + ) + SELECT len, is_dir, updated_at_epoch, TRUE AS exact_match + FROM exact + UNION ALL + SELECT NULL::bigint AS len, TRUE AS is_dir, NULL::bigint AS updated_at_epoch, + FALSE AS exact_match + FROM child + WHERE NOT EXISTS (SELECT 1 FROM exact) + LIMIT 1 + "#, + &[&path.as_str(), &prefix_lower, &prefix_upper], + ) + .await + .map_err(|error| db_error(path.clone(), FilesystemOperation::Stat, error))?; + let Some(row) = row else { + return Err(not_found(path.clone(), FilesystemOperation::Stat)); + }; + let exact_match: bool = row.get("exact_match"); + if !exact_match { + return Ok(FileStat { + path: path.clone(), + file_type: FileType::Directory, + len: 0, + modified: None, + sensitive: false, + }); + } + let len: Option = row.get("len"); + let is_dir: bool = row.get("is_dir"); + let updated_at_epoch: Option = row.get("updated_at_epoch"); + Ok(FileStat { + path: path.clone(), + file_type: if is_dir { + FileType::Directory + } else { + FileType::File + }, + len: if is_dir { + 0 + } else { + len.unwrap_or_default().max(0) as u64 + }, + modified: updated_at_epoch.and_then(system_time_from_unix_seconds), + sensitive: false, + }) +} + #[cfg(feature = "postgres")] async fn postgres_delete_with_client( client: &deadpool_postgres::Object, @@ -1518,6 +1675,26 @@ fn descendant_path_range(path: &VirtualPath) -> (String, String) { (format!("{prefix}/"), format!("{prefix}0")) } +#[cfg(feature = "postgres")] +fn postgres_shared_projection_index_name(spec: &IndexSpec) -> String { + let kind = match &spec.kind { + IndexKind::Exact => "exact", + IndexKind::Prefix => "prefix", + IndexKind::Fts => "fts", + IndexKind::Vector { .. } => "vector", + }; + let mut keys = postgres_projection_index_component(kind); + for key in &spec.keys { + keys.push_str(&postgres_projection_index_component(key.as_str())); + } + sql_index_name(&format!("/shared/{kind}/{keys}"), spec.name.as_str()) +} + +#[cfg(feature = "postgres")] +fn postgres_projection_index_component(value: &str) -> String { + format!("{}:{value}", value.len()) +} + /// Translate a [`Filter`] tree into a postgres WHERE-clause fragment. /// Bound parameters use `$N` placeholders sized from `params.len() + 1`. /// @@ -1879,4 +2056,77 @@ mod tests { assert!("/secrets/a/b" < lower); // the path itself is excluded assert!("/secrets/a/bb" >= upper); // prefix-sharing sibling excluded } + + #[test] + fn shared_projection_index_name_ignores_prefix_specific_declarations() { + let spec = IndexSpec::new( + crate::IndexName::new("bucket_exact").unwrap(), + vec![crate::IndexKey::new("bucket").unwrap()], + IndexKind::Exact, + ); + let first = postgres_shared_projection_index_name(&spec); + let second = postgres_shared_projection_index_name(&spec); + assert_eq!(first, second); + assert!(first.contains("shared")); + assert!(first.contains("bucket_exact")); + assert!(first.len() <= 62); + } + + #[test] + fn shared_projection_index_name_separates_kind_and_keys() { + let exact_bucket = IndexSpec::new( + crate::IndexName::new("bucket_exact").unwrap(), + vec![crate::IndexKey::new("bucket").unwrap()], + IndexKind::Exact, + ); + let prefix_bucket = IndexSpec::new( + crate::IndexName::new("bucket_exact").unwrap(), + vec![crate::IndexKey::new("bucket").unwrap()], + IndexKind::Prefix, + ); + let exact_tenant = IndexSpec::new( + crate::IndexName::new("bucket_exact").unwrap(), + vec![crate::IndexKey::new("tenant_id").unwrap()], + IndexKind::Exact, + ); + assert_ne!( + postgres_shared_projection_index_name(&exact_bucket), + postgres_shared_projection_index_name(&prefix_bucket) + ); + assert_ne!( + postgres_shared_projection_index_name(&exact_bucket), + postgres_shared_projection_index_name(&exact_tenant) + ); + } + + #[test] + fn shared_projection_index_name_uses_injective_key_encoding() { + let split_keys = IndexSpec::new( + crate::IndexName::new("bucket_exact").unwrap(), + vec![ + crate::IndexKey::new("a").unwrap(), + crate::IndexKey::new("b").unwrap(), + ], + IndexKind::Exact, + ); + let joined_key = IndexSpec::new( + crate::IndexName::new("bucket_exact").unwrap(), + vec![crate::IndexKey::new("a_b").unwrap()], + IndexKind::Exact, + ); + + assert_ne!( + postgres_shared_projection_index_name(&split_keys), + postgres_shared_projection_index_name(&joined_key) + ); + } + + #[test] + fn quote_postgres_identifier_doubles_embedded_quotes() { + assert_eq!(quote_postgres_identifier("idx_simple"), "\"idx_simple\""); + assert_eq!( + quote_postgres_identifier("idx_\"quoted\""), + "\"idx_\"\"quoted\"\"\"" + ); + } } diff --git a/crates/ironclaw_filesystem/src/scoped.rs b/crates/ironclaw_filesystem/src/scoped.rs index 855fea8a080..c0dfd2f874f 100644 --- a/crates/ironclaw_filesystem/src/scoped.rs +++ b/crates/ironclaw_filesystem/src/scoped.rs @@ -703,6 +703,12 @@ impl StorageTxn for ScopedStorageTxn { self.inner.delete(path).await } + async fn reserve_sequence(&mut self, path: &VirtualPath) -> Result { + self.check(FilesystemOperation::ReserveSeq)?; + self.check_path(path)?; + self.inner.reserve_sequence(path).await + } + async fn commit(self: Box) -> Result<(), FilesystemError> { self.inner.commit().await } diff --git a/crates/ironclaw_reborn_event_store/src/lib.rs b/crates/ironclaw_reborn_event_store/src/lib.rs index 581fa4233d7..b5a21dc6f3e 100644 --- a/crates/ironclaw_reborn_event_store/src/lib.rs +++ b/crates/ironclaw_reborn_event_store/src/lib.rs @@ -291,6 +291,19 @@ pub async fn build_reborn_event_stores( /// `/events`. Production composition reuses this on top of a libSQL / /// PostgreSQL `RootFilesystem` so the backend choice is a property of the /// filesystem rather than of the durable-log impl. +/// +/// The caller must run any backend schema migrations before calling this +/// helper. Config-based builders perform their own migration step. +#[cfg(any(feature = "libsql", feature = "postgres"))] +pub fn build_reborn_event_stores_from_root_filesystem( + root: Arc, +) -> Result +where + F: RootFilesystem + Send + Sync + 'static, +{ + wrap_root_filesystem_as_event_stores(root) +} + #[cfg(any(feature = "libsql", feature = "postgres"))] fn wrap_root_filesystem_as_event_stores( root: Arc,