Development Documentation (main branch) - For stable release docs, see docs.rs/eidetica
Skip to main content

eidetica/backend/database/sql/
mod.rs

1//! SQL-based backend implementations for Eidetica storage.
2//!
3//! This module provides SQL database backends that implement the `BackendImpl` trait,
4//! allowing Eidetica entries to be stored in relational databases.
5//!
6//! ## Available Backends
7//!
8//! - **SQLite** (feature: `sqlite`): Embedded database
9//! - **PostgreSQL** (feature: `postgres`): PostgreSQL database
10//!
11//! ## Architecture
12//!
13//! The SQL backend uses sqlx with `AnyPool` for multi-database support.
14//! All methods are async to match the async `BackendImpl` trait.
15//!
16//! ## Schema and Migrations
17//!
18//! The database schema is defined in the [`schema`] module and automatically
19//! initialized when connecting. Migrations are handled via code-based functions
20//! rather than SQL files to support dialect differences between SQLite and PostgreSQL.
21//!
22//! See [`schema`] module documentation for details on adding migrations.
23
24mod storage;
25mod traversal;
26
27/// Schema definition and migration system.
28pub mod schema;
29
30use std::any::Any;
31#[cfg(feature = "sqlite")]
32use std::fs::{File, OpenOptions, TryLockError};
33#[cfg(feature = "sqlite")]
34use std::io;
35#[cfg(feature = "sqlite")]
36use std::path::{Path, PathBuf};
37#[cfg(feature = "sqlite")]
38use std::str::FromStr;
39#[cfg(any(feature = "sqlite", feature = "postgres"))]
40use std::time::Duration;
41
42use async_trait::async_trait;
43#[cfg(feature = "postgres")]
44use sqlx::AnyConnection;
45#[cfg(feature = "postgres")]
46use sqlx::Connection;
47use sqlx::any::AnyPoolOptions;
48#[cfg(feature = "sqlite")]
49use sqlx::sqlite::SqliteConnectOptions;
50use sqlx::{AnyPool, Executor};
51
52use crate::Result;
53use crate::backend::errors::BackendError;
54use crate::backend::{
55    BackendImpl, InstanceMetadata, InstanceSecrets, RecordMutations, RecordPage, RecordRange,
56    RecordView, StagingToken, StoreStateRequest, VerificationStatus,
57};
58use crate::entry::{Entry, ID};
59use crate::snapshot::Snapshot;
60
61/// Extension trait for sqlx Result types to simplify error handling.
62///
63/// Similar to `anyhow::Context`, this trait adds a method to convert
64/// sqlx errors to `BackendError::SqlxError` with a context message.
65pub(crate) trait SqlxResultExt<T> {
66    /// Convert sqlx error to BackendError with context message.
67    fn sql_context(self, context: &str) -> Result<T>;
68}
69
70impl<T> SqlxResultExt<T> for std::result::Result<T, sqlx::Error> {
71    fn sql_context(self, context: &str) -> Result<T> {
72        self.map_err(|e| {
73            BackendError::SqlxError {
74                reason: format!("{context}: {e}"),
75                source: Some(e),
76            }
77            .into()
78        })
79    }
80}
81
82/// Database backend kind for SQL dialect selection.
83#[derive(Debug, Clone, Copy, PartialEq, Eq)]
84pub enum DbKind {
85    /// SQLite database
86    Sqlite,
87    /// PostgreSQL database
88    Postgres,
89}
90
91/// SQL-based backend implementing `BackendImpl` using sqlx.
92///
93/// This backend supports both SQLite and PostgreSQL through sqlx's `AnyPool`.
94///
95/// # Concurrency
96///
97/// `SqlxBackend` is `Send + Sync` as required by `BackendImpl`. The underlying
98/// sqlx pool handles connection pooling and thread safety. Each backend owns its
99/// persistent storage namespace exclusively for its lifetime. Share one backend
100/// through the Eidetica service rather than opening the same storage directly.
101///
102/// # Test Isolation
103///
104/// For PostgreSQL, each backend instance can use its own schema for test isolation.
105/// Use `connect_postgres_isolated()` to create an isolated backend for testing.
106///
107/// ```compile_fail
108/// use eidetica::backend::database::SqlxBackend;
109///
110/// fn cannot_escape_pool(backend: SqlxBackend) {
111///     let _ = backend.pool();
112/// }
113/// ```
114pub struct SqlxBackend {
115    pool: Option<AnyPool>,
116    kind: DbKind,
117    _owner: Option<StorageOwner>,
118    #[cfg(all(feature = "postgres", feature = "testing"))]
119    postgres_token: Option<String>,
120}
121
122impl Drop for SqlxBackend {
123    fn drop(&mut self) {
124        // Mark every pool handle closed before releasing the ownership lock. Creating this
125        // future performs the close transition; waiting is unnecessary for the fence.
126        // Checked-out SQLx connections remain valid until returned, so PostgreSQL retains
127        // their shared advisory locks until then.
128        if let Some(pool) = self.pool.take() {
129            drop(pool.close());
130        }
131    }
132}
133
134enum StorageOwner {
135    #[cfg(feature = "sqlite")]
136    Sqlite { _lock: File },
137    #[cfg(feature = "postgres")]
138    Postgres {
139        _connection: tokio::sync::Mutex<AnyConnection>,
140    },
141}
142
143#[cfg(feature = "sqlite")]
144fn prepare_sqlite(url: &str) -> Result<Option<StorageOwner>> {
145    let normalized_url = normalize_sqlite_url(url);
146    let options = SqliteConnectOptions::from_str(&normalized_url).map_err(|error| {
147        BackendError::SqlxError {
148            reason: format!("Failed to parse SQLite connection URL: {error}"),
149            source: Some(error),
150        }
151    })?;
152    if sqlite_is_in_memory(&normalized_url) {
153        return Ok(None);
154    }
155
156    let database_path = canonical_database_path(options.get_filename()).map_err(|error| {
157        BackendError::SqlxError {
158            reason: format!(
159                "Failed to identify SQLite database `{}`: {error}",
160                options.get_filename().display()
161            ),
162            source: None,
163        }
164    })?;
165    let lock_path = sqlite_owner_lock_path(&database_path);
166    let lock = OpenOptions::new()
167        .read(true)
168        .write(true)
169        .create(true)
170        .truncate(false)
171        .open(&lock_path)
172        .map_err(|error| BackendError::SqlxError {
173            reason: format!(
174                "Failed to open SQLite ownership sidecar `{}`: {error}",
175                lock_path.display()
176            ),
177            source: None,
178        })?;
179    match lock.try_lock() {
180        Ok(()) => Ok(Some(StorageOwner::Sqlite { _lock: lock })),
181        Err(TryLockError::WouldBlock) => Err(BackendError::StorageAlreadyOwned {
182            namespace: database_path.display().to_string(),
183        }
184        .into()),
185        Err(TryLockError::Error(error)) => Err(BackendError::SqlxError {
186            reason: format!(
187                "Failed to claim SQLite database ownership `{}`: {error}",
188                database_path.display()
189            ),
190            source: None,
191        }
192        .into()),
193    }
194}
195
196#[cfg(feature = "sqlite")]
197fn sqlite_owner_lock_path(database_path: &Path) -> PathBuf {
198    let mut path = database_path.as_os_str().to_owned();
199    path.push(".eidetica-owner");
200    PathBuf::from(path)
201}
202
203#[cfg(feature = "sqlite")]
204fn normalize_sqlite_url(url: &str) -> String {
205    let Some(rest) = url.strip_prefix("sqlite:file:") else {
206        return url.to_owned();
207    };
208    if rest == ":memory:" || rest.starts_with(":memory:?") {
209        return url.to_owned();
210    }
211    format!("sqlite:{rest}")
212}
213
214#[cfg(feature = "sqlite")]
215fn sqlite_is_in_memory(url: &str) -> bool {
216    let url = url
217        .trim_start_matches("sqlite://")
218        .trim_start_matches("sqlite:");
219    let (database, _) = url.split_once('?').unwrap_or((url, ""));
220    database == ":memory:" || database == "file::memory:" || sqlite_file_mode(url)
221}
222
223#[cfg(feature = "sqlite")]
224fn sqlite_file_mode(url: &str) -> bool {
225    let query = url.split_once('?').map_or("", |(_, query)| query);
226    url::form_urlencoded::parse(query.as_bytes())
227        .any(|(key, value)| key == "mode" && value == "memory")
228}
229
230#[cfg(feature = "sqlite")]
231fn canonical_database_path(path: &Path) -> io::Result<PathBuf> {
232    let absolute = if path.is_absolute() {
233        path.to_owned()
234    } else {
235        std::env::current_dir()?.join(path)
236    };
237
238    match absolute.symlink_metadata() {
239        Ok(_) => return absolute.canonicalize(),
240        Err(error) if error.kind() == io::ErrorKind::NotFound => {}
241        Err(error) => return Err(error),
242    }
243
244    let file_name = absolute.file_name().ok_or_else(|| {
245        io::Error::new(
246            io::ErrorKind::InvalidInput,
247            "database path has no file name",
248        )
249    })?;
250    let parent = absolute.parent().ok_or_else(|| {
251        io::Error::new(io::ErrorKind::InvalidInput, "database path has no parent")
252    })?;
253    Ok(parent.canonicalize()?.join(file_name))
254}
255
256#[cfg(feature = "sqlite")]
257fn sqlite_path_url(path: &Path) -> Result<String> {
258    let absolute = if path.is_absolute() {
259        path.to_owned()
260    } else {
261        std::env::current_dir()
262            .map_err(|error| BackendError::SqlxError {
263                reason: format!("Failed to resolve SQLite database path: {error}"),
264                source: None,
265            })?
266            .join(path)
267    };
268    let file_url = url::Url::from_file_path(&absolute).map_err(|()| BackendError::SqlxError {
269        reason: format!(
270            "Failed to encode SQLite database path `{}` as a file URL",
271            absolute.display()
272        ),
273        source: None,
274    })?;
275    Ok(format!(
276        "sqlite:{}?mode=rwc",
277        &file_url.as_str()["file:".len()..]
278    ))
279}
280
281impl SqlxBackend {
282    /// Get a reference to the underlying pool for SQL backend modules.
283    pub(crate) fn pool(&self) -> &AnyPool {
284        self.pool.as_ref().expect("SQL pool must exist until drop")
285    }
286
287    #[cfg(all(feature = "postgres", feature = "testing"))]
288    #[doc(hidden)]
289    pub async fn test_postgres_checked_out_connection(
290        &self,
291    ) -> Result<sqlx::pool::PoolConnection<sqlx::Any>> {
292        self.pool()
293            .acquire()
294            .await
295            .sql_context("Failed to acquire PostgreSQL test connection")
296    }
297
298    /// Get the database kind.
299    pub fn kind(&self) -> DbKind {
300        self.kind
301    }
302
303    /// Check if this backend is using SQLite.
304    pub fn is_sqlite(&self) -> bool {
305        self.kind == DbKind::Sqlite
306    }
307
308    /// Check if this backend is using PostgreSQL.
309    pub fn is_postgres(&self) -> bool {
310        self.kind == DbKind::Postgres
311    }
312
313    #[cfg(all(feature = "postgres", feature = "testing"))]
314    #[doc(hidden)]
315    pub async fn test_postgres_owner_pid(&self) -> Result<i32> {
316        let _connection = match self
317            ._owner
318            .as_ref()
319            .expect("PostgreSQL backend must hold an ownership connection")
320        {
321            StorageOwner::Postgres { _connection } => _connection,
322            #[cfg(feature = "sqlite")]
323            StorageOwner::Sqlite { .. } => unreachable!("only PostgreSQL ownership is queried"),
324        };
325        let mut connection = _connection.lock().await;
326        let (pid,): (i32,) = sqlx::query_as("SELECT pg_backend_pid()")
327            .fetch_one(&mut *connection)
328            .await
329            .sql_context("Failed to read PostgreSQL ownership session")?;
330        Ok(pid)
331    }
332
333    #[cfg(all(feature = "postgres", feature = "testing"))]
334    #[doc(hidden)]
335    pub async fn test_postgres_pool_pids(&self) -> Result<Vec<i32>> {
336        let mut connections = Vec::with_capacity(2);
337        for _ in 0..2 {
338            connections.push(
339                self.pool()
340                    .acquire()
341                    .await
342                    .sql_context("Failed to acquire PostgreSQL test connection")?,
343            );
344        }
345        let mut pids = Vec::with_capacity(connections.len());
346        for connection in &mut connections {
347            let (pid,): (i32,) = sqlx::query_as("SELECT pg_backend_pid()")
348                .fetch_one(&mut **connection)
349                .await
350                .sql_context("Failed to read PostgreSQL pool session")?;
351            pids.push(pid);
352        }
353        Ok(pids)
354    }
355
356    #[cfg(all(feature = "postgres", feature = "testing"))]
357    #[doc(hidden)]
358    pub fn test_postgres_token(&self) -> &str {
359        self.postgres_token
360            .as_deref()
361            .expect("PostgreSQL backend must hold an ownership token")
362    }
363}
364
365// Test-only Store-state stage pause hook.
366//
367// A regression test for same-token stage-vs-publish interleavings needs to
368// park a `stage_store_state_records` call after token validation (advisory
369// locks held) but before any record writes, run a competing publish, then
370// release the stage. Production call paths cannot pause mid-transaction, and
371// a sleep-based race would be flaky by construction, so this narrow hook
372// exists behind the `testing` feature only: it is compiled out of every
373// production build, adds no trait surface, and is a no-op (one map miss)
374// unless a test registered a gate for the exact staging namespace. Gates are
375// one-shot and keyed by namespace UUID, so parallel tests cannot observe or
376// disturb each other.
377#[cfg(feature = "testing")]
378#[derive(Debug)]
379pub struct StoreStateStagePause {
380    validated: tokio::sync::Notify,
381    release: tokio::sync::Notify,
382}
383
384#[cfg(feature = "testing")]
385impl StoreStateStagePause {
386    fn new() -> Self {
387        Self {
388            validated: tokio::sync::Notify::new(),
389            release: tokio::sync::Notify::new(),
390        }
391    }
392
393    /// Wait until the staged call reaches the pause point (bounded).
394    ///
395    /// # Panics
396    ///
397    /// Panics after 15 seconds: a timeout means the test never drove a stage
398    /// through the gate, i.e. a harness bug, not a backend result.
399    pub async fn wait_validated(&self) {
400        if tokio::time::timeout(
401            std::time::Duration::from_secs(15),
402            self.validated.notified(),
403        )
404        .await
405        .is_err()
406        {
407            panic!("Store-state stage pause gate never reached: test harness bug");
408        }
409    }
410
411    /// Let the paused stage proceed to its writes.
412    pub fn release(&self) {
413        self.release.notify_one();
414    }
415}
416
417#[cfg(feature = "testing")]
418static STORE_STATE_STAGE_PAUSES: std::sync::OnceLock<
419    tokio::sync::Mutex<std::collections::HashMap<String, std::sync::Arc<StoreStateStagePause>>>,
420> = std::sync::OnceLock::new();
421
422#[cfg(feature = "testing")]
423fn store_state_stage_pauses() -> &'static tokio::sync::Mutex<
424    std::collections::HashMap<String, std::sync::Arc<StoreStateStagePause>>,
425> {
426    STORE_STATE_STAGE_PAUSES
427        .get_or_init(|| tokio::sync::Mutex::new(std::collections::HashMap::new()))
428}
429
430/// Register a one-shot pause gate for the given staging namespace.
431///
432/// The next `stage_store_state_records` call for this namespace signals the
433/// gate after validating its token and waits (bounded) for [`StoreStateStagePause::release`]
434/// before writing any records. Returns the gate the test drives.
435#[cfg(feature = "testing")]
436impl SqlxBackend {
437    pub async fn testing_register_stage_pause(
438        namespace_id: &str,
439    ) -> std::sync::Arc<StoreStateStagePause> {
440        let gate = std::sync::Arc::new(StoreStateStagePause::new());
441        store_state_stage_pauses()
442            .lock()
443            .await
444            .insert(namespace_id.to_string(), gate.clone());
445        gate
446    }
447}
448
449/// Fire the pause gate for a namespace, if a test registered one.
450///
451/// Called from `stage_store_state_records` after token validation, before
452/// record writes. One-shot: the gate is removed before signalling, so a late
453/// duplicate stage for the same namespace proceeds unpaused.
454#[cfg(feature = "testing")]
455pub(crate) async fn fire_store_state_stage_pause(namespace_id: &str) {
456    let gate = store_state_stage_pauses().lock().await.remove(namespace_id);
457    let Some(gate) = gate else { return };
458    gate.validated.notify_one();
459    if tokio::time::timeout(std::time::Duration::from_secs(15), gate.release.notified())
460        .await
461        .is_err()
462    {
463        panic!(
464            "Store-state stage pause gate for namespace {namespace_id} was never released: test harness bug"
465        );
466    }
467}
468
469// SQLite-specific implementations
470#[cfg(feature = "sqlite")]
471impl SqlxBackend {
472    /// Open a SQLite database at the given path.
473    ///
474    /// Creates the database file and schema if they don't exist.
475    ///
476    /// # Arguments
477    ///
478    /// * `path` - Path to the SQLite database file
479    ///
480    /// # Example
481    ///
482    /// ```ignore
483    /// use eidetica::backend::database::sql::SqlxBackend;
484    ///
485    /// #[tokio::main]
486    /// async fn main() {
487    ///     let backend = SqlxBackend::open_sqlite("my_database.db").await.unwrap();
488    /// }
489    /// ```
490    pub async fn open_sqlite<P: AsRef<std::path::Path>>(path: P) -> Result<Self> {
491        // mode=rwc: read-write-create (create file if it doesn't exist)
492        let url = sqlite_path_url(path.as_ref())?;
493        Self::connect_sqlite(&url).await
494    }
495
496    /// Connect to a SQLite database using a connection URL.
497    ///
498    /// # Arguments
499    ///
500    /// * `url` - SQLite connection URL (e.g., "sqlite:./my.db")
501    pub async fn connect_sqlite(url: &str) -> Result<Self> {
502        // Install any driver support
503        sqlx::any::install_default_drivers();
504
505        let storage = prepare_sqlite(url)?;
506        let is_in_memory = storage.is_none();
507        let normalized_url = normalize_sqlite_url(url);
508
509        // For SQLite in-memory databases with shared cache, we must prevent
510        // all connections from being closed. When the last connection closes,
511        // the in-memory database is destroyed and all data is lost.
512        //
513        // IMPORTANT: SQLite pragmas like busy_timeout and synchronous are per-connection
514        // settings. We use after_connect to ensure every connection in the pool has
515        // these configured, not just one.
516        let pool = if is_in_memory {
517            AnyPoolOptions::new()
518                .max_connections(1)
519                .min_connections(1)
520                .idle_timeout(None)
521                .max_lifetime(None)
522                .after_connect(|conn, _meta| {
523                    Box::pin(async move {
524                        // In-memory databases don't need WAL mode (all in RAM)
525                        // but still need busy_timeout for lock contention
526                        conn.execute("PRAGMA busy_timeout = 5000;").await?;
527                        Ok(())
528                    })
529                })
530                .connect(&normalized_url)
531                .await
532                .sql_context("Failed to connect to SQLite")?
533        } else {
534            AnyPoolOptions::new()
535                .max_connections(5)
536                .after_connect(|conn, _meta| {
537                    Box::pin(async move {
538                        // File-based SQLite per-connection settings:
539                        // - synchronous=NORMAL: Balanced durability (safe with WAL)
540                        // - busy_timeout=5000: Wait up to 5s for locks before failing
541                        //
542                        // Note: journal_mode=WAL is a database-level setting that persists,
543                        // so we only set it once after pool creation, not per-connection.
544                        conn.execute("PRAGMA synchronous = NORMAL; PRAGMA busy_timeout = 5000;")
545                            .await?;
546                        Ok(())
547                    })
548                })
549                .connect(&normalized_url)
550                .await
551                .sql_context("Failed to connect to SQLite")?
552        };
553
554        // Set WAL mode once (database-level setting that persists in the file)
555        if !is_in_memory {
556            sqlx::query("PRAGMA journal_mode = WAL;")
557                .execute(&pool)
558                .await
559                .sql_context("Failed to set SQLite WAL mode")?;
560        }
561
562        let backend = Self {
563            pool: Some(pool),
564            kind: DbKind::Sqlite,
565            _owner: storage,
566            #[cfg(all(feature = "postgres", feature = "testing"))]
567            postgres_token: None,
568        };
569
570        // Initialize schema
571        schema::initialize(&backend).await?;
572
573        Ok(backend)
574    }
575
576    /// Create an in-memory SQLite database (async).
577    ///
578    /// The database exists only for the lifetime of this backend instance.
579    /// Useful for testing.
580    ///
581    /// # Example
582    ///
583    /// ```ignore
584    /// use eidetica::backend::database::sql::SqlxBackend;
585    ///
586    /// #[tokio::main]
587    /// async fn main() {
588    ///     let backend = SqlxBackend::sqlite_in_memory().await.unwrap();
589    /// }
590    /// ```
591    pub async fn sqlite_in_memory() -> Result<Self> {
592        // Use shared cache mode for in-memory SQLite so all connections in the pool
593        // share the same database. Without this, each connection gets its own
594        // isolated in-memory database.
595        // Use a unique name per instance to avoid sharing between tests.
596        let unique_id = uuid::Uuid::new_v4();
597        let url = format!("sqlite:file:mem_{unique_id}?mode=memory&cache=shared");
598        Self::connect_sqlite(&url).await
599    }
600}
601
602// PostgreSQL-specific implementations
603#[cfg(feature = "postgres")]
604const POSTGRES_OWNERSHIP_TABLE: &str = "_eidetica_storage_owner";
605
606#[cfg(feature = "postgres")]
607const POSTGRES_NAMESPACE_LOCK: &str = "hashtextextended(format('eidetica-storage-v1:%s/%s:%s/%s', octet_length(current_database()), current_database(), octet_length(current_schema()), current_schema()), 0)";
608
609#[cfg(feature = "postgres")]
610impl SqlxBackend {
611    /// Connect to a PostgreSQL database using a connection URL.
612    ///
613    /// This connects to the default (public) schema. For test isolation,
614    /// use `connect_postgres_isolated()` instead.
615    ///
616    /// # Arguments
617    ///
618    /// * `url` - PostgreSQL connection URL (e.g., "postgres://user:pass@localhost/dbname")
619    ///
620    /// # Example
621    ///
622    /// ```ignore
623    /// use eidetica::backend::database::sql::SqlxBackend;
624    ///
625    /// let backend = SqlxBackend::connect_postgres("postgres://localhost/eidetica").await.unwrap();
626    /// ```
627    pub async fn connect_postgres(url: &str) -> Result<Self> {
628        Self::connect_postgres_with_schema(url, None).await
629    }
630
631    /// Connect to a PostgreSQL database with a specific schema for isolation.
632    ///
633    /// Creates a unique schema if `schema_name` is provided, providing test isolation.
634    /// Each test can use its own schema so they don't interfere with each other.
635    ///
636    /// # Arguments
637    ///
638    /// * `url` - PostgreSQL connection URL
639    /// * `schema_name` - Optional schema name. If None, uses the default (public) schema.
640    async fn connect_postgres_with_schema(url: &str, schema_name: Option<String>) -> Result<Self> {
641        // Install any driver support
642        sqlx::any::install_default_drivers();
643
644        // If schema_name is provided, first create the schema, then use after_connect
645        // to set search_path on each connection. This is more reliable than URL options
646        // which don't work consistently across all network configurations.
647        if let Some(ref schema) = schema_name {
648            // First connect to create the schema if needed
649            let temp_pool = AnyPoolOptions::new()
650                .max_connections(1)
651                .connect(url)
652                .await
653                .sql_context("Failed to connect to PostgreSQL")?;
654
655            // Create schema if it doesn't exist
656            let create_schema = format!("CREATE SCHEMA IF NOT EXISTS {schema}");
657            sqlx::query(&create_schema)
658                .execute(&temp_pool)
659                .await
660                .sql_context(&format!("Failed to create schema {schema}"))?;
661
662            temp_pool.close().await;
663        }
664
665        // Build pool with after_connect hook to set search_path on each connection
666        // For isolated (test) connections, use smaller pool to avoid exhausting
667        // PostgreSQL's max_connections when running many tests in parallel.
668        let mut owner = AnyConnection::connect(url)
669            .await
670            .sql_context("Failed to connect to PostgreSQL")?;
671        if let Some(ref schema) = schema_name {
672            let set_path = format!("SET search_path TO {schema}");
673            owner
674                .execute(set_path.as_str())
675                .await
676                .sql_context("Failed to select PostgreSQL storage namespace")?;
677        }
678
679        let (database, schema, acquired): (String, String, bool) = sqlx::query_as(&format!(
680            "SELECT current_database()::text, current_schema()::text, pg_try_advisory_lock({POSTGRES_NAMESPACE_LOCK})"
681        ))
682        .fetch_one(&mut owner)
683        .await
684        .sql_context("Failed to claim PostgreSQL storage ownership")?;
685        if !acquired {
686            return Err(BackendError::StorageAlreadyOwned {
687                namespace: format!("PostgreSQL database `{database}` schema `{schema}`"),
688            }
689            .into());
690        }
691
692        owner
693            .execute(format!(
694                "CREATE TABLE IF NOT EXISTS {POSTGRES_OWNERSHIP_TABLE} (id SMALLINT PRIMARY KEY CHECK (id = 1), token TEXT NOT NULL)"
695            ).as_str())
696            .await
697            .sql_context("Failed to initialize PostgreSQL storage ownership metadata")?;
698        let token = uuid::Uuid::new_v4().to_string();
699        sqlx::query(&format!(
700            "INSERT INTO {POSTGRES_OWNERSHIP_TABLE} (id, token) VALUES (1, $1) ON CONFLICT (id) DO UPDATE SET token = EXCLUDED.token"
701        ))
702        .bind(&token)
703        .execute(&mut owner)
704        .await
705        .sql_context("Failed to publish PostgreSQL storage ownership token")?;
706        owner
707            .execute(format!("SELECT pg_advisory_lock_shared({POSTGRES_NAMESPACE_LOCK})").as_str())
708            .await
709            .sql_context("Failed to fence PostgreSQL storage ownership")?;
710        owner
711            .execute(format!("SELECT pg_advisory_unlock({POSTGRES_NAMESPACE_LOCK})").as_str())
712            .await
713            .sql_context("Failed to finish PostgreSQL storage ownership claim")?;
714
715        let schema_for_hook = schema_name.clone();
716        let is_isolated = schema_name.is_some();
717        let mut pool_options = AnyPoolOptions::new();
718
719        if is_isolated {
720            // Test isolation: 2 connections is enough, with longer timeout to wait
721            // rather than fail when many tests run in parallel
722            pool_options = pool_options
723                .max_connections(2)
724                .acquire_timeout(Duration::from_secs(30));
725        } else {
726            // Production: 5 connections for real concurrency needs
727            pool_options = pool_options.max_connections(5);
728        }
729
730        let token_for_hook = token.clone();
731        let pool = pool_options
732            .after_connect(move |conn, _meta| {
733                let schema = schema_for_hook.clone();
734                let token = token_for_hook.clone();
735                Box::pin(async move {
736                    if let Some(ref schema) = schema {
737                        let set_path = format!("SET search_path TO {schema}");
738                        conn.execute(set_path.as_str()).await?;
739                    }
740                    conn.execute(
741                        format!("SELECT pg_advisory_lock_shared({POSTGRES_NAMESPACE_LOCK})")
742                            .as_str(),
743                    )
744                    .await?;
745                    let valid: (bool,) = sqlx::query_as(&format!(
746                        "SELECT token = $1 FROM {POSTGRES_OWNERSHIP_TABLE} WHERE id = 1"
747                    ))
748                    .bind(token)
749                    .fetch_one(&mut *conn)
750                    .await?;
751                    if !valid.0 {
752                        return Err(sqlx::Error::Protocol(
753                            "PostgreSQL storage ownership changed".to_string(),
754                        ));
755                    }
756                    Ok(())
757                })
758            })
759            .connect(url)
760            .await
761            .sql_context("Failed to connect to PostgreSQL")?;
762
763        let backend = Self {
764            pool: Some(pool),
765            kind: DbKind::Postgres,
766            _owner: Some(StorageOwner::Postgres {
767                _connection: tokio::sync::Mutex::new(owner),
768            }),
769            #[cfg(feature = "testing")]
770            postgres_token: Some(token),
771        };
772
773        // Initialize schema (tables will be created in the current search_path)
774        schema::initialize(&backend).await?;
775
776        Ok(backend)
777    }
778
779    /// Connect to a PostgreSQL database with test isolation.
780    ///
781    /// Creates a unique schema for this backend instance, ensuring tests
782    /// don't interfere with each other when run in parallel.
783    ///
784    /// # Arguments
785    ///
786    /// * `url` - PostgreSQL connection URL (e.g., "postgres://user:pass@localhost/dbname")
787    ///
788    /// # Example
789    ///
790    /// ```ignore
791    /// use eidetica::backend::database::sql::SqlxBackend;
792    ///
793    /// let backend = SqlxBackend::connect_postgres_isolated("postgres://localhost/eidetica").await.unwrap();
794    /// // This backend uses its own isolated schema
795    /// ```
796    pub async fn connect_postgres_isolated(url: &str) -> Result<Self> {
797        // Generate a unique schema name using UUID
798        // PostgreSQL schema names must start with a letter and be lowercase
799        let unique_id = uuid::Uuid::new_v4().simple().to_string();
800        let schema_name = format!("test_{unique_id}");
801        Self::connect_postgres_with_schema(url, Some(schema_name)).await
802    }
803
804    #[cfg(feature = "testing")]
805    #[doc(hidden)]
806    pub async fn test_connect_postgres_schema(url: &str, schema: String) -> Result<Self> {
807        Self::connect_postgres_with_schema(url, Some(schema)).await
808    }
809}
810
811#[async_trait]
812impl BackendImpl for SqlxBackend {
813    async fn resolve_store_state(&self, request: &StoreStateRequest) -> Result<Option<RecordView>> {
814        storage::resolve_store_state(self, request).await
815    }
816
817    async fn begin_store_state_staging(&self, request: StoreStateRequest) -> Result<StagingToken> {
818        storage::begin_store_state_staging(self, request).await
819    }
820
821    async fn stage_store_state_records(
822        &self,
823        token: &StagingToken,
824        records: RecordMutations,
825    ) -> Result<()> {
826        storage::stage_store_state_records(self, token, records).await
827    }
828
829    async fn publish_store_state(&self, token: StagingToken) -> Result<RecordView> {
830        storage::publish_store_state(self, token).await
831    }
832
833    async fn abort_store_state(&self, token: StagingToken) -> Result<()> {
834        storage::abort_store_state(self, token).await
835    }
836
837    async fn store_state_record_get(
838        &self,
839        view: &RecordView,
840        key: &[u8],
841    ) -> Result<Option<Vec<u8>>> {
842        storage::store_state_record_get(self, view, key).await
843    }
844
845    async fn store_state_record_scan(
846        &self,
847        view: &RecordView,
848        range: &RecordRange,
849        after: Option<&[u8]>,
850        limit: usize,
851    ) -> Result<RecordPage> {
852        storage::store_state_record_scan(self, view, range, after, limit).await
853    }
854
855    async fn clear_derived_store_state(&self) -> Result<()> {
856        storage::clear_derived_store_state(self).await
857    }
858
859    async fn reset_local_verification(&self) -> Result<()> {
860        storage::reset_local_verification(self).await
861    }
862    async fn get(&self, id: &ID) -> Result<Entry> {
863        storage::get(self, id).await
864    }
865
866    async fn get_verification_status(&self, id: &ID) -> Result<VerificationStatus> {
867        storage::get_verification_status(self, id).await
868    }
869
870    async fn put(&self, entry: Entry) -> Result<()> {
871        storage::put(self, entry).await
872    }
873
874    async fn update_verification_status(
875        &self,
876        id: &ID,
877        verification_status: VerificationStatus,
878    ) -> Result<()> {
879        storage::update_verification_status(self, id, verification_status).await
880    }
881
882    async fn get_entries_by_verification_status(
883        &self,
884        status: VerificationStatus,
885    ) -> Result<Vec<ID>> {
886        storage::get_entries_by_verification_status(self, status).await
887    }
888
889    async fn snapshot(&self, tree: &ID) -> Result<Snapshot> {
890        traversal::snapshot(self, tree).await.map(Snapshot::new)
891    }
892
893    async fn store_snapshot(&self, tree: &ID, store: &str) -> Result<Snapshot> {
894        traversal::store_snapshot(self, tree, store)
895            .await
896            .map(Snapshot::new)
897    }
898
899    async fn store_snapshot_at(
900        &self,
901        tree: &ID,
902        store: &str,
903        main_snapshot: &Snapshot,
904    ) -> Result<Snapshot> {
905        traversal::store_snapshot_at(self, tree, store, main_snapshot.tips())
906            .await
907            .map(Snapshot::new)
908    }
909
910    async fn all_roots(&self) -> Result<Vec<ID>> {
911        storage::all_roots(self).await
912    }
913
914    async fn find_merge_base(
915        &self,
916        tree: &ID,
917        store: &str,
918        entry_ids: &[ID],
919    ) -> Result<Option<ID>> {
920        traversal::find_merge_base(self, tree, store, entry_ids).await
921    }
922
923    fn as_any(&self) -> &dyn Any {
924        self
925    }
926
927    async fn get_tree(&self, tree: &ID) -> Result<Vec<Entry>> {
928        storage::get_tree(self, tree).await
929    }
930
931    async fn get_store(&self, tree: &ID, store: &str) -> Result<Vec<Entry>> {
932        storage::get_store(self, tree, store).await
933    }
934
935    async fn get_tree_from_tips(&self, tree: &ID, tips: &[ID]) -> Result<Vec<Entry>> {
936        traversal::get_tree_from_tips(self, tree, tips).await
937    }
938
939    async fn store_at(&self, tree: &ID, store: &str, snapshot: &Snapshot) -> Result<Vec<Entry>> {
940        traversal::store_at(self, tree, store, snapshot.tips()).await
941    }
942
943    async fn get_sorted_store_parents(
944        &self,
945        tree_id: &ID,
946        entry_id: &ID,
947        store: &str,
948    ) -> Result<Vec<ID>> {
949        traversal::get_sorted_store_parents(self, tree_id, entry_id, store).await
950    }
951
952    async fn get_path_from_to(
953        &self,
954        tree_id: &ID,
955        store: &str,
956        from_id: Option<&ID>,
957        to_ids: &[ID],
958    ) -> Result<Vec<ID>> {
959        traversal::get_path_from_to(self, tree_id, store, from_id, to_ids).await
960    }
961
962    async fn get_instance_metadata(&self) -> Result<Option<InstanceMetadata>> {
963        storage::get_instance_metadata(self).await
964    }
965
966    async fn set_instance_metadata(&self, metadata: &InstanceMetadata) -> Result<()> {
967        storage::set_instance_metadata(self, metadata).await
968    }
969
970    async fn get_instance_secrets(&self) -> Result<Option<InstanceSecrets>> {
971        storage::get_instance_secrets(self).await
972    }
973
974    async fn set_instance_secrets(&self, secrets: &InstanceSecrets) -> Result<()> {
975        storage::set_instance_secrets(self, secrets).await
976    }
977}
978
979/// Namespace for SQLite database constructors.
980///
981/// Provides ergonomic factory methods for creating SQLite-backed storage.
982/// All methods return `SqlxBackend` which implements `BackendImpl`.
983///
984/// # Example
985///
986/// ```ignore
987/// use eidetica::backend::database::Sqlite;
988///
989/// // File-based storage
990/// let backend = Sqlite::open("my_data.db").await?;
991///
992/// // In-memory (for testing)
993/// let backend = Sqlite::in_memory().await?;
994/// ```
995#[cfg(feature = "sqlite")]
996pub struct Sqlite;
997
998#[cfg(feature = "sqlite")]
999impl Sqlite {
1000    /// Open a SQLite database at the given path.
1001    ///
1002    /// Creates the database file and schema if they don't exist.
1003    ///
1004    /// # Arguments
1005    ///
1006    /// * `path` - Path to the SQLite database file
1007    pub async fn open<P: AsRef<std::path::Path>>(path: P) -> Result<SqlxBackend> {
1008        SqlxBackend::open_sqlite(path).await
1009    }
1010
1011    /// Create an in-memory SQLite database.
1012    ///
1013    /// The database exists only for the lifetime of the returned backend.
1014    /// Useful for testing.
1015    pub async fn in_memory() -> Result<SqlxBackend> {
1016        SqlxBackend::sqlite_in_memory().await
1017    }
1018
1019    /// Connect to a SQLite database using a connection URL.
1020    ///
1021    /// # Arguments
1022    ///
1023    /// * `url` - SQLite connection URL (e.g., "sqlite:./my.db")
1024    pub async fn connect(url: &str) -> Result<SqlxBackend> {
1025        SqlxBackend::connect_sqlite(url).await
1026    }
1027}
1028
1029/// Namespace for PostgreSQL database constructors.
1030///
1031/// Provides ergonomic factory methods for creating PostgreSQL-backed storage.
1032/// All methods return `SqlxBackend` which implements `BackendImpl`.
1033///
1034/// # Example
1035///
1036/// ```ignore
1037/// use eidetica::backend::database::Postgres;
1038///
1039/// // Connect to PostgreSQL
1040/// let backend = Postgres::connect("postgres://user:pass@localhost/mydb").await?;
1041///
1042/// // With test isolation (unique schema per instance)
1043/// let backend = Postgres::connect_isolated("postgres://localhost/test").await?;
1044/// ```
1045#[cfg(feature = "postgres")]
1046pub struct Postgres;
1047
1048#[cfg(feature = "postgres")]
1049impl Postgres {
1050    /// Connect to a PostgreSQL database using a connection URL.
1051    ///
1052    /// This connects to the default (public) schema. For test isolation,
1053    /// use `connect_isolated()` instead.
1054    ///
1055    /// # Arguments
1056    ///
1057    /// * `url` - PostgreSQL connection URL (e.g., "postgres://user:pass@localhost/dbname")
1058    pub async fn connect(url: &str) -> Result<SqlxBackend> {
1059        SqlxBackend::connect_postgres(url).await
1060    }
1061
1062    /// Connect to a PostgreSQL database with test isolation.
1063    ///
1064    /// Creates a unique schema for this backend instance, ensuring tests
1065    /// don't interfere with each other when run in parallel.
1066    ///
1067    /// # Arguments
1068    ///
1069    /// * `url` - PostgreSQL connection URL
1070    pub async fn connect_isolated(url: &str) -> Result<SqlxBackend> {
1071        SqlxBackend::connect_postgres_isolated(url).await
1072    }
1073}
1074
1075#[cfg(all(test, feature = "postgres"))]
1076mod tests {
1077    use super::*;
1078
1079    #[tokio::test]
1080    async fn failed_postgres_initialization_releases_ownership() {
1081        if std::env::var("TEST_BACKEND").as_deref() != Ok("postgres") {
1082            return;
1083        }
1084
1085        let url = std::env::var("TEST_POSTGRES_URL")
1086            .unwrap_or_else(|_| "postgres://localhost/eidetica_test".to_string());
1087        sqlx::any::install_default_drivers();
1088        let schema = format!("test_{}", uuid::Uuid::new_v4().simple());
1089        let setup = AnyPoolOptions::new()
1090            .max_connections(1)
1091            .connect(&url)
1092            .await
1093            .unwrap();
1094        sqlx::query(&format!("CREATE SCHEMA {schema}"))
1095            .execute(&setup)
1096            .await
1097            .unwrap();
1098        sqlx::query(&format!(
1099            "CREATE VIEW {schema}.entries AS SELECT 1 AS value"
1100        ))
1101        .execute(&setup)
1102        .await
1103        .unwrap();
1104
1105        assert!(
1106            SqlxBackend::connect_postgres_with_schema(&url, Some(schema.clone()))
1107                .await
1108                .is_err(),
1109            "the conflicting view must make schema initialization fail"
1110        );
1111        sqlx::query(&format!("DROP VIEW {schema}.entries"))
1112            .execute(&setup)
1113            .await
1114            .unwrap();
1115
1116        SqlxBackend::connect_postgres_with_schema(&url, Some(schema.clone()))
1117            .await
1118            .expect("failed initialization must release PostgreSQL ownership");
1119        sqlx::query(&format!("DROP SCHEMA {schema} CASCADE"))
1120            .execute(&setup)
1121            .await
1122            .unwrap();
1123        setup.close().await;
1124    }
1125}
1126
1127/// A publish carrying a token whose target does not match the namespace it
1128/// names must never disturb a ready namespace.
1129///
1130/// Unlike `InMemory` (which resolves by target and adopts the ready winner),
1131/// the SQL publish names the namespace: the `UPDATE` only flips
1132/// lifecycle/status and keeps the begin-time identity, using `target` for
1133/// locking and failure-path winner adoption. Either way the ready snapshot
1134/// and its records always survive a malformed clone.
1135#[cfg(all(test, feature = "sqlite"))]
1136mod store_state_token_tests {
1137    use std::collections::BTreeMap;
1138
1139    use super::Sqlite;
1140    use crate::backend::{
1141        BackendImpl, CacheScope, ProjectionDescriptor, StagingToken, StoreStateLifecycle,
1142        StoreStateRequest,
1143    };
1144    use crate::entry::ID;
1145
1146    fn request(store: &str) -> StoreStateRequest {
1147        StoreStateRequest {
1148            database: ID::from_bytes("db"),
1149            store: store.to_string(),
1150            lifecycle: StoreStateLifecycle::Derived,
1151            scope: CacheScope::Shared,
1152            projection: ProjectionDescriptor {
1153                name: "test/opaque".to_string(),
1154                version: 0,
1155            },
1156            source_key: b"snapshot".to_vec(),
1157        }
1158    }
1159
1160    #[tokio::test]
1161    async fn mismatched_target_publish_preserves_ready_namespace() {
1162        let backend = Sqlite::in_memory().await.unwrap();
1163
1164        // Ready namespace for target A.
1165        let request_a = request("store-a");
1166        let token_a = backend
1167            .begin_store_state_staging(request_a.clone())
1168            .await
1169            .unwrap();
1170        backend
1171            .stage_store_state_records(
1172                &token_a,
1173                BTreeMap::from([(b"key".to_vec(), Some(b"value-a".to_vec()))]),
1174            )
1175            .await
1176            .unwrap();
1177        let view_a = backend.publish_store_state(token_a).await.unwrap();
1178
1179        // Staging namespace for target B.
1180        let request_b = request("store-b");
1181        let token_b = backend
1182            .begin_store_state_staging(request_b.clone())
1183            .await
1184            .unwrap();
1185        backend
1186            .stage_store_state_records(
1187                &token_b,
1188                BTreeMap::from([(b"key".to_vec(), Some(b"value-b".to_vec()))]),
1189            )
1190            .await
1191            .unwrap();
1192
1193        // Malformed clone: B's namespace id, A's target. The SQL publish names
1194        // the namespace (the UPDATE only flips lifecycle/status and keeps the
1195        // begin-time identity; `target` drives locking and winner adoption),
1196        // so B becomes ready under its own identity while ready A is
1197        // untouched: no ready snapshot is ever modified by a mismatched token.
1198        let bad = StagingToken {
1199            namespace_id: token_b.namespace_id.clone(),
1200            target: request_a.clone(),
1201        };
1202        let published_b = backend.publish_store_state(bad).await.unwrap();
1203        assert_eq!(published_b.namespace_id, token_b.namespace_id);
1204        assert_eq!(
1205            backend.resolve_store_state(&request_b).await.unwrap(),
1206            Some(published_b.clone())
1207        );
1208        assert_eq!(
1209            backend
1210                .store_state_record_get(&published_b, b"key")
1211                .await
1212                .unwrap(),
1213            Some(b"value-b".to_vec())
1214        );
1215        assert_eq!(
1216            backend.resolve_store_state(&request_a).await.unwrap(),
1217            Some(view_a.clone())
1218        );
1219        assert_eq!(
1220            backend
1221                .store_state_record_get(&view_a, b"key")
1222                .await
1223                .unwrap(),
1224            Some(b"value-a".to_vec())
1225        );
1226
1227        // Malformed clone: A's (ready) namespace id with a target that
1228        // resolves nowhere. The publish fails and the ready row survives the
1229        // guarded discard with its records intact.
1230        let nowhere = request("store-nowhere");
1231        let bad_ready = StagingToken {
1232            namespace_id: view_a.namespace_id.clone(),
1233            target: nowhere.clone(),
1234        };
1235        assert!(backend.publish_store_state(bad_ready).await.is_err());
1236        assert_eq!(
1237            backend.resolve_store_state(&request_a).await.unwrap(),
1238            Some(view_a.clone())
1239        );
1240        assert_eq!(
1241            backend
1242                .store_state_record_get(&view_a, b"key")
1243                .await
1244                .unwrap(),
1245            Some(b"value-a".to_vec())
1246        );
1247        assert_eq!(backend.resolve_store_state(&nowhere).await.unwrap(), None);
1248    }
1249}