diff --git a/.claude/rules/database.md b/.claude/rules/database.md index 44bd3baa720..005d4134cdc 100644 --- a/.claude/rules/database.md +++ b/.claude/rules/database.md @@ -6,7 +6,34 @@ paths: --- # Database Rules -Dual-backend persistence: PostgreSQL + libSQL/Turso. **All new persistence features must support both backends.** +## Status & Direction + +The repo is migrating off per-crate `Store`/`Repository` traits onto a +single universal `RootFilesystem` mount table (`crates/ironclaw_filesystem/`). +Under the new model, every persistence concern is a mount path +(`/system/secrets`, `/system/processes`, `/engine/threads`, …) backed by +exactly one `RootFilesystem` implementation — typed stores become thin +wrappers around `ScopedFilesystem` and own no backend dispatch of their +own. See `crates/ironclaw_filesystem/CLAUDE.md` and the +`2026-05-14-universal-fs-dispatch.md` plan/ADR. + +**New persistence features go on `ScopedFilesystem`, not into `src/db/`.** +The rules below cover the *legacy* per-crate dual-backend pattern that +still exists in `src/db/`, `src/history/`, and `migrations/`. Touch them +only when fixing or extending code that already lives there; do not add +new sub-traits or per-domain backends. + +This file is `paths`-scoped to those legacy directories so the rule +loads when (and only when) you're inside them. New code under +`crates/ironclaw_filesystem/`, consumer crates routing through it, or +any new mount-backed store should follow the unified-surface contract +in `crates/ironclaw_filesystem/CLAUDE.md` instead. + +--- + +## Legacy: Dual-Backend Per-Crate Pattern + +Dual-backend persistence: PostgreSQL + libSQL/Turso. **All new persistence features must support both backends.** *(Applies only inside the legacy directories scoped above. For new crates, mount through `RootFilesystem` and let the wiring layer pick the backend.)* See `src/db/CLAUDE.md` for full schema, dialect differences, and libSQL limitations. diff --git a/Cargo.lock b/Cargo.lock index c2dfff6030b..af2117167c6 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4487,7 +4487,6 @@ dependencies = [ "ironclaw_events", "ironclaw_filesystem", "ironclaw_host_api", - "ironclaw_storage", "ironclaw_turns", "libsql", "serde", @@ -4801,18 +4800,6 @@ dependencies = [ "urlencoding", ] -[[package]] -name = "ironclaw_storage" -version = "0.1.0" -dependencies = [ - "async-trait", - "serde", - "serde_json", - "thiserror 2.0.18", - "tokio", - "tracing", -] - [[package]] name = "ironclaw_telegram_v2_adapter" version = "0.1.0" diff --git a/Cargo.toml b/Cargo.toml index 298bf6667f3..6ecd77baad3 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,5 +1,5 @@ [workspace] -members = [".", "crates/ironclaw_common", "crates/ironclaw_host_api", "crates/ironclaw_storage", "crates/ironclaw_filesystem", "crates/ironclaw_memory", "crates/ironclaw_events", "crates/ironclaw_event_projections", "crates/ironclaw_reborn_event_store", "crates/ironclaw_extensions", "crates/ironclaw_processes", "crates/ironclaw_dispatcher", "crates/ironclaw_scripts", "crates/ironclaw_mcp", "crates/ironclaw_wasm", "crates/ironclaw_capabilities", "crates/ironclaw_secrets", "crates/ironclaw_network", "crates/ironclaw_host_runtime", "crates/ironclaw_runtime_policy", "crates/ironclaw_authorization", "crates/ironclaw_run_state", "crates/ironclaw_approvals", "crates/ironclaw_resources", "crates/ironclaw_trust", "crates/ironclaw_turns", "crates/ironclaw_agent_loop", "crates/ironclaw_threads", "crates/ironclaw_loop_support", "crates/ironclaw_reborn", "crates/ironclaw_reborn_config", "crates/ironclaw_reborn_composition", "crates/ironclaw_reborn_cli", "crates/ironclaw_conversations", "crates/ironclaw_product_adapters", "crates/ironclaw_product_workflow", "crates/ironclaw_product_adapter_registry", "crates/ironclaw_wasm_product_adapters", "crates/ironclaw_telegram_v2_adapter", "crates/ironclaw_outbound", "crates/ironclaw_architecture", "crates/ironclaw_safety", "crates/ironclaw_skills", "crates/ironclaw_llm", "crates/ironclaw_engine", "crates/ironclaw_gateway", "crates/ironclaw_tui"] +members = [".", "crates/ironclaw_common", "crates/ironclaw_host_api", "crates/ironclaw_filesystem", "crates/ironclaw_memory", "crates/ironclaw_events", "crates/ironclaw_event_projections", "crates/ironclaw_reborn_event_store", "crates/ironclaw_extensions", "crates/ironclaw_processes", "crates/ironclaw_dispatcher", "crates/ironclaw_scripts", "crates/ironclaw_mcp", "crates/ironclaw_wasm", "crates/ironclaw_capabilities", "crates/ironclaw_secrets", "crates/ironclaw_network", "crates/ironclaw_host_runtime", "crates/ironclaw_runtime_policy", "crates/ironclaw_authorization", "crates/ironclaw_run_state", "crates/ironclaw_approvals", "crates/ironclaw_resources", "crates/ironclaw_trust", "crates/ironclaw_turns", "crates/ironclaw_agent_loop", "crates/ironclaw_threads", "crates/ironclaw_loop_support", "crates/ironclaw_reborn", "crates/ironclaw_reborn_config", "crates/ironclaw_reborn_composition", "crates/ironclaw_reborn_cli", "crates/ironclaw_conversations", "crates/ironclaw_product_adapters", "crates/ironclaw_product_workflow", "crates/ironclaw_product_adapter_registry", "crates/ironclaw_wasm_product_adapters", "crates/ironclaw_telegram_v2_adapter", "crates/ironclaw_outbound", "crates/ironclaw_architecture", "crates/ironclaw_safety", "crates/ironclaw_skills", "crates/ironclaw_llm", "crates/ironclaw_engine", "crates/ironclaw_gateway", "crates/ironclaw_tui"] exclude = [ "channels-src/discord", "channels-src/feishu", diff --git a/crates/ironclaw_architecture/tests/reborn_dependency_boundaries.rs b/crates/ironclaw_architecture/tests/reborn_dependency_boundaries.rs index 75c982edc37..b9919fc3910 100644 --- a/crates/ironclaw_architecture/tests/reborn_dependency_boundaries.rs +++ b/crates/ironclaw_architecture/tests/reborn_dependency_boundaries.rs @@ -766,51 +766,6 @@ fn boundary_rules() -> Vec { "ironclaw_wasm_product_adapters", ], }, - BoundaryRule { - crate_name: "ironclaw_storage", - forbidden: vec![ - "ironclaw", - "ironclaw_approvals", - "ironclaw_architecture", - "ironclaw_authorization", - "ironclaw_capabilities", - "ironclaw_common", - "ironclaw_conversations", - "ironclaw_dispatcher", - "ironclaw_engine", - "ironclaw_event_projections", - "ironclaw_events", - "ironclaw_extensions", - "ironclaw_filesystem", - "ironclaw_gateway", - "ironclaw_host_api", - "ironclaw_host_runtime", - "ironclaw_llm", - "ironclaw_loop_support", - "ironclaw_mcp", - "ironclaw_memory", - "ironclaw_network", - "ironclaw_outbound", - "ironclaw_processes", - "ironclaw_product_adapters", - "ironclaw_reborn", - "ironclaw_reborn_cli", - "ironclaw_reborn_config", - "ironclaw_reborn_event_store", - "ironclaw_resources", - "ironclaw_run_state", - "ironclaw_runtime_policy", - "ironclaw_safety", - "ironclaw_scripts", - "ironclaw_secrets", - "ironclaw_skills", - "ironclaw_threads", - "ironclaw_trust", - "ironclaw_tui", - "ironclaw_turns", - "ironclaw_wasm", - ], - }, BoundaryRule { crate_name: "ironclaw_reborn_config", forbidden: vec![ @@ -1012,7 +967,10 @@ fn boundary_rules() -> Vec { "ironclaw_conversations", "ironclaw_dispatcher", "ironclaw_extensions", - "ironclaw_filesystem", + // ironclaw_filesystem is permitted: FilesystemOutboundStateStore + // routes outbound persistence through ScopedFilesystem under + // the universal-fs-dispatch rework (plan + // 2026-05-14-universal-fs-dispatch). "ironclaw_gateway", "ironclaw_host_runtime", "ironclaw_mcp", diff --git a/crates/ironclaw_filesystem/src/hsm.rs b/crates/ironclaw_filesystem/src/hsm.rs new file mode 100644 index 00000000000..c168d4dac62 --- /dev/null +++ b/crates/ironclaw_filesystem/src/hsm.rs @@ -0,0 +1,264 @@ +//! Placeholder HSM-style backend demonstrating the universal dispatch seam. +//! +//! `HsmBackend` is intentionally minimal: it shows what a backend looks like +//! when its storage primitives sit behind an external boundary (a hardware +//! security module, a TEE-resident KMS, an OS keychain). A production HSM +//! implementation would replace [`HsmBackend::new`] with constructors that +//! accept an HSM session handle and would route [`put`](RootFilesystem::put) / +//! [`get`](RootFilesystem::get) / [`delete`](RootFilesystem::delete) through +//! the HSM's encrypt/decrypt API. +//! +//! The point of having it in-tree is to prove that adding a new backend is a +//! single-file change. The seam this demonstrates: +//! +//! 1. **One trait.** `HsmBackend` implements `RootFilesystem` and nothing +//! else; the composite dispatcher routes through it like any other mount. +//! 2. **Declared capabilities.** The HSM exposes only the encrypted-bytes +//! surface (`Read` / `Write` / `Stat` / `Delete`) — no records, no query, +//! no index, no events, no transactions. `BackendCapabilities` advertises +//! that up front; `CompositeRootFilesystem::mount_dyn` then refuses any +//! `MountDescriptor` that claims more than the HSM delivers +//! (`FilesystemError::DescriptorOverclaims`). Consumers cannot accidentally +//! attach a query-requiring store onto an HSM mount. +//! 3. **Swap by wiring, not by consumer edits.** A consumer that holds a +//! `ScopedFilesystem` bound to `/system/secrets` does not change when the +//! mount swaps from `LibSqlRootFilesystem` to `HsmBackend` — only the +//! `mount()` call at startup changes. +//! +//! The placeholder stores ciphertext in process memory so the trait can be +//! exercised end-to-end in tests. It is not a security boundary. Real HSM +//! backends are sealed behind external infrastructure. + +use async_trait::async_trait; +use ironclaw_host_api::VirtualPath; + +use crate::in_memory::InMemoryBackend; +use crate::{ + BackendCapabilities, Capability, CasExpectation, DirEntry, Entry, FileStat, FilesystemError, + FilesystemOperation, RecordVersion, RootFilesystem, TxnCapability, VersionedEntry, +}; + +/// Placeholder HSM-style backend. See the module-level docs. +pub struct HsmBackend { + inner: InMemoryBackend, +} + +impl HsmBackend { + pub fn new() -> Self { + Self { + inner: InMemoryBackend::new(), + } + } + + fn declared_capabilities() -> BackendCapabilities { + BackendCapabilities::empty() + .with(Capability::Read) + .with(Capability::Write) + .with(Capability::Stat) + .with(Capability::Delete) + .with_txn(TxnCapability::Cas) + } +} + +impl Default for HsmBackend { + fn default() -> Self { + Self::new() + } +} + +#[async_trait] +impl RootFilesystem for HsmBackend { + fn capabilities(&self) -> BackendCapabilities { + Self::declared_capabilities() + } + + async fn put( + &self, + path: &VirtualPath, + entry: Entry, + cas: CasExpectation, + ) -> Result { + if entry.kind.is_some() || !entry.indexed.is_empty() { + return Err(FilesystemError::Unsupported { + path: path.clone(), + operation: FilesystemOperation::WriteFile, + }); + } + self.inner.put(path, entry, cas).await + } + + async fn get(&self, path: &VirtualPath) -> Result, FilesystemError> { + self.inner.get(path).await + } + + async fn list_dir(&self, path: &VirtualPath) -> Result, FilesystemError> { + self.inner.list_dir(path).await + } + + async fn stat(&self, path: &VirtualPath) -> Result { + self.inner.stat(path).await + } + + async fn delete(&self, path: &VirtualPath) -> Result<(), FilesystemError> { + self.inner.delete(path).await + } +} + +#[cfg(test)] +mod tests { + use std::sync::Arc; + + use ironclaw_host_api::VirtualPath; + + use crate::{ + BackendCapabilities, BackendId, BackendKind, Capability, CasExpectation, + CompositeRootFilesystem, ContentKind, Entry, FilesystemError, IndexKey, IndexKind, + IndexName, IndexPolicy, IndexSpec, IndexValue, MountDescriptor, RootFilesystem, + StorageClass, TxnCapability, + }; + + use super::HsmBackend; + + fn vpath(value: &str) -> VirtualPath { + VirtualPath::new(value).unwrap() + } + + fn hsm_descriptor(capabilities: BackendCapabilities) -> MountDescriptor { + MountDescriptor { + virtual_root: vpath("/secrets"), + backend_id: BackendId::new("hsm-secrets").unwrap(), + backend_kind: BackendKind::Custom("hsm".into()), + storage_class: StorageClass::FileContent, + content_kind: ContentKind::SystemState, + index_policy: IndexPolicy::NotIndexed, + capabilities, + } + } + + #[tokio::test] + async fn hsm_supports_encrypted_bytes_round_trip() { + let hsm = HsmBackend::new(); + let path = vpath("/secrets/account/api-key"); + + let version = hsm + .put( + &path, + Entry::bytes(b"ciphertext-blob".to_vec()), + CasExpectation::Absent, + ) + .await + .unwrap(); + + let read_back = hsm.get(&path).await.unwrap().expect("entry present"); + assert_eq!(read_back.entry.body, b"ciphertext-blob"); + assert_eq!(read_back.version, version); + } + + #[tokio::test] + async fn hsm_rejects_structured_records() { + let hsm = HsmBackend::new(); + let path = vpath("/secrets/account/api-key"); + let entry = Entry::record( + crate::RecordKind::new("credential_lease").unwrap(), + &serde_json::json!({"scope": "team"}), + ) + .unwrap(); + + let err = hsm + .put(&path, entry, CasExpectation::Any) + .await + .unwrap_err(); + assert!(matches!(err, FilesystemError::Unsupported { .. })); + } + + #[tokio::test] + async fn hsm_rejects_query_and_index_ops() { + let hsm = HsmBackend::new(); + let path = vpath("/secrets"); + + let query_err = hsm + .query(&path, &crate::Filter::All, crate::Page::new(0, 10)) + .await + .unwrap_err(); + assert!(matches!(query_err, FilesystemError::Unsupported { .. })); + + let spec = IndexSpec::new( + IndexName::new("by_scope").unwrap(), + vec![IndexKey::new("scope").unwrap()], + IndexKind::Exact, + ); + let index_err = hsm.ensure_index(&path, &spec).await.unwrap_err(); + assert!(matches!(index_err, FilesystemError::Unsupported { .. })); + } + + #[tokio::test] + async fn composite_rejects_overclaimed_hsm_descriptor() { + let mut composite = CompositeRootFilesystem::new(); + let over_claimed = BackendCapabilities::empty() + .with(Capability::Read) + .with(Capability::Write) + .with(Capability::Stat) + .with(Capability::Delete) + .with(Capability::Query) + .with(Capability::IndexExact); + + let err = composite + .mount_dyn(hsm_descriptor(over_claimed), Arc::new(HsmBackend::new())) + .unwrap_err(); + + match err { + FilesystemError::DescriptorOverclaims { missing, .. } => { + assert!(missing.contains(&Capability::Query)); + assert!(missing.contains(&Capability::IndexExact)); + } + other => panic!("expected DescriptorOverclaims, got {other:?}"), + } + } + + #[tokio::test] + async fn composite_routes_to_hsm_under_secrets_mount() { + // Demonstrates the swap-by-wiring acceptance gate: a /secrets mount + // can point at HsmBackend with no consumer changes. Consumer code + // sees the same RootFilesystem trait — capability validation at + // mount time ensures the descriptor never over-claims. + let mut composite = CompositeRootFilesystem::new(); + let honest_descriptor = hsm_descriptor( + BackendCapabilities::empty() + .with(Capability::Read) + .with(Capability::Write) + .with(Capability::Stat) + .with(Capability::Delete) + .with_txn(TxnCapability::Cas), + ); + composite + .mount_dyn(honest_descriptor, Arc::new(HsmBackend::new())) + .unwrap(); + + let path = vpath("/secrets/account/key"); + composite + .put( + &path, + Entry::bytes(b"ciphertext".to_vec()), + CasExpectation::Absent, + ) + .await + .unwrap(); + + // Read back via the composite — same surface the consumer would use. + let read = composite.get(&path).await.unwrap().expect("entry present"); + assert_eq!(read.entry.body, b"ciphertext"); + + // Indexed values are still rejected because the HSM declared no + // index/query capability — even though the consumer code path + // didn't change. + let with_index = Entry::bytes(b"more".to_vec()).with_indexed( + IndexKey::new("scope").unwrap(), + IndexValue::Text("team".into()), + ); + let err = composite + .put(&path, with_index, CasExpectation::Any) + .await + .unwrap_err(); + assert!(matches!(err, FilesystemError::Unsupported { .. })); + } +} diff --git a/crates/ironclaw_filesystem/src/lib.rs b/crates/ironclaw_filesystem/src/lib.rs index 5e4835c21f8..2dea7dc2d52 100644 --- a/crates/ironclaw_filesystem/src/lib.rs +++ b/crates/ironclaw_filesystem/src/lib.rs @@ -19,6 +19,7 @@ mod backend; mod catalog; #[cfg(any(feature = "postgres", feature = "libsql"))] mod db; +mod hsm; mod in_memory; mod index; #[cfg(feature = "libsql")] @@ -33,6 +34,7 @@ mod types; pub use backend::{EventRecord, StorageTxn}; pub use catalog::{CompositeRootFilesystem, FilesystemCatalog, MountDescriptor, PathPlacement}; +pub use hsm::HsmBackend; pub use in_memory::InMemoryBackend; pub use index::{Filter, IndexKey, IndexKind, IndexName, IndexSpec, IndexValue, Page}; #[cfg(feature = "libsql")] diff --git a/crates/ironclaw_outbound/Cargo.toml b/crates/ironclaw_outbound/Cargo.toml index ccc54efc2c4..28639eff4fb 100644 --- a/crates/ironclaw_outbound/Cargo.toml +++ b/crates/ironclaw_outbound/Cargo.toml @@ -18,7 +18,6 @@ hex = "0.4" ironclaw_event_projections = { path = "../ironclaw_event_projections" } ironclaw_filesystem = { path = "../ironclaw_filesystem" } ironclaw_host_api = { path = "../ironclaw_host_api" } -ironclaw_storage = { path = "../ironclaw_storage" } ironclaw_turns = { path = "../ironclaw_turns" } libsql = { version = "0.6", optional = true, default-features = false, features = ["core", "replication", "remote", "tls"] } serde = { version = "1", features = ["derive"] } diff --git a/crates/ironclaw_outbound/src/db.rs b/crates/ironclaw_outbound/src/db.rs index 1356d966dcf..b835c9725ce 100644 --- a/crates/ironclaw_outbound/src/db.rs +++ b/crates/ironclaw_outbound/src/db.rs @@ -20,24 +20,27 @@ struct DeliveryIdentity<'a> { #[cfg(any(feature = "libsql", feature = "postgres"))] pub(crate) fn to_json(value: &T) -> Result { - ironclaw_storage::encode_json(value).map_err(|_| OutboundError::Serialization) + serde_json::to_string(value).map_err(|_| OutboundError::Serialization) } #[cfg(any(feature = "libsql", feature = "postgres"))] pub(crate) fn from_json(value: &str) -> Result { - ironclaw_storage::decode_json(value).map_err(|_| OutboundError::Serialization) + serde_json::from_str(value).map_err(|_| OutboundError::Serialization) } +// Backend errors are logged with raw detail for operator diagnostics, then +// collapsed into a payload-free variant — outbound CLAUDE.md forbids +// returning raw backend error details to callers. #[cfg(any(feature = "libsql", feature = "postgres"))] pub(crate) fn db_error(error: impl std::fmt::Display) -> OutboundError { - tracing::debug!(error = %&error, "outbound storage backend error"); - let redacted = ironclaw_storage::redacted_backend_error(error); - debug_assert_eq!(redacted, ironclaw_storage::StorageError::Backend); + tracing::error!(%error, "outbound storage backend error"); OutboundError::Backend } +// Sentinel value for optional composite-key components stored as columns. +// Inlined from the dissolved ironclaw_storage crate (only consumer). #[cfg(any(feature = "libsql", feature = "postgres"))] -const ABSENT_SCOPE_ID: &str = ironclaw_storage::ABSENT_SCOPE_COMPONENT; +const ABSENT_SCOPE_ID: &str = ""; #[cfg(any(feature = "libsql", feature = "postgres"))] pub(crate) fn scope_agent_db_value(scope: &TurnScope) -> &str { diff --git a/crates/ironclaw_storage/AGENTS.md b/crates/ironclaw_storage/AGENTS.md deleted file mode 100644 index 35521fa8859..00000000000 --- a/crates/ironclaw_storage/AGENTS.md +++ /dev/null @@ -1,11 +0,0 @@ -# ironclaw_storage - -Shared storage substrate primitives for Reborn persistence adapters. - -## Invariants - -- Own storage mechanics only: backend identity, redacted errors, migration descriptors, pagination, serialization helpers, primitive `BlobStore`/`RecordStore` traits, and future append-log/lock/transaction/encrypted-blob primitives. -- Do **not** own domain semantics. No turn/thread/outbound/secret-specific operations or schemas here; domain traits stay in their owning crates. -- Do **not** depend on other IronClaw domain crates. Domain adapters may depend on this crate, not the reverse. -- Errors must stay redacted: no raw SQL/backend messages, host paths, secrets, payload snippets, or provider/runtime details. -- Filesystem remains file-shaped/path authority; this crate is not a universal filesystem replacement. diff --git a/crates/ironclaw_storage/CLAUDE.md b/crates/ironclaw_storage/CLAUDE.md deleted file mode 100644 index fee59945f74..00000000000 --- a/crates/ironclaw_storage/CLAUDE.md +++ /dev/null @@ -1,11 +0,0 @@ -# ironclaw_storage - -Shared storage substrate primitives for Reborn persistence adapters. - -## Invariants - -- Own DB/storage mechanics only: backend identity, redacted errors, migration descriptors, pagination, serialization helpers, primitive `BlobStore`/`RecordStore` traits, and future append-log/lock/transaction/encrypted-blob primitives. -- Do **not** own domain semantics. No turn/thread/outbound/secret-specific operations or schemas here; domain traits stay in their owning crates. -- Do **not** depend on other IronClaw domain crates. Domain adapters may depend on this crate, not the reverse. -- Errors must stay redacted: no raw SQL/backend messages, host paths, secrets, payload snippets, or provider/runtime details. -- Filesystem remains file-shaped/path authority; this crate is not a universal filesystem replacement. diff --git a/crates/ironclaw_storage/Cargo.toml b/crates/ironclaw_storage/Cargo.toml deleted file mode 100644 index 729dc764c81..00000000000 --- a/crates/ironclaw_storage/Cargo.toml +++ /dev/null @@ -1,17 +0,0 @@ -[package] -name = "ironclaw_storage" -version = "0.1.0" -edition = "2024" -rust-version = "1.92" -description = "Shared storage substrate primitives for IronClaw Reborn persistence adapters" -publish = false - -[dependencies] -async-trait = "0.1" -serde = { version = "1", features = ["derive"] } -serde_json = "1" -thiserror = "2" -tracing = "0.1" - -[dev-dependencies] -tokio = { version = "1", features = ["macros", "rt"] } diff --git a/crates/ironclaw_storage/src/lib.rs b/crates/ironclaw_storage/src/lib.rs deleted file mode 100644 index bda997daabd..00000000000 --- a/crates/ironclaw_storage/src/lib.rs +++ /dev/null @@ -1,660 +0,0 @@ -//! Shared storage substrate primitives for Reborn persistence adapters. -//! -//! This crate owns reusable persistence mechanics only: backend identity, -//! redacted storage errors, JSON serialization helpers, pagination limits, -//! and migration descriptors. Domain crates still own their store traits, -//! schemas, validation, and query semantics. - -use async_trait::async_trait; -use serde::{ - Deserialize, Deserializer, Serialize, - de::{DeserializeOwned, IgnoredAny}, -}; -use thiserror::Error; - -/// Sentinel value for optional scope components stored in composite SQL keys. -/// -/// Domain stores should continue to own which scope fields participate in a -/// key. This constant only keeps the storage representation consistent across -/// adapters when an optional scope component is absent. -pub const ABSENT_SCOPE_COMPONENT: &str = ""; - -/// Supported durable backend families known to the shared substrate. -/// -/// This is an identity/support marker, not an authority grant and not a domain -/// repository selector. Composition code may use it for diagnostics and -/// migration routing after the owning domain has selected a store. -/// `Filesystem` identifies a blob/record store implementation backed by -/// filesystem mechanics; it does not grant file-shaped path authority or turn -/// this crate into a filesystem abstraction. -#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] -pub enum StorageBackendKind { - Memory, - LibSql, - Postgres, - Filesystem, - Object, -} - -impl StorageBackendKind { - pub fn as_str(self) -> &'static str { - match self { - Self::Memory => "memory", - Self::LibSql => "libsql", - Self::Postgres => "postgres", - Self::Filesystem => "filesystem", - Self::Object => "object", - } - } -} - -/// Redacted storage-substrate error. -/// -/// Variants intentionally avoid carrying raw backend messages. Domain stores can -/// map these into their own error types without leaking SQL details, host paths, -/// secret material, or provider/runtime payloads. -// Copy is intentional: variants must remain payload-free so redaction stays trivial. -#[derive(Debug, Clone, Copy, PartialEq, Eq, Error)] -pub enum StorageError { - #[error("storage backend operation failed")] - Backend, - #[error("storage serialization failed")] - Serialization, - #[error("storage migration failed")] - Migration, - #[error("storage record validation failed")] - Validation, - #[error("storage write conflict")] - Conflict, - #[error("storage operation is unsupported by backend")] - Unsupported, -} - -/// Convert any backend error into a redacted storage error. -/// -/// The input is accepted so callers can pass concrete DB/client errors at the -/// call site, but it is deliberately discarded to prevent accidental leakage. -/// Callers that need operational diagnostics must log the raw error before -/// passing it here; this function is the redaction boundary. -pub fn redacted_backend_error(error: impl std::fmt::Display) -> StorageError { - tracing::error!(%error, "storage backend operation failed"); - StorageError::Backend -} - -/// Serialize a structured payload without exposing serializer internals. -pub fn encode_json(value: &T) -> Result { - serde_json::to_string(value).map_err(|_| StorageError::Serialization) -} - -/// Deserialize a structured payload without exposing raw payload snippets. -pub fn decode_json(value: &str) -> Result { - serde_json::from_str(value).map_err(|_| StorageError::Serialization) -} - -/// Return a stable SQL value for an optional scoped identifier. -pub fn optional_scope_component(value: Option<&str>) -> &str { - value.unwrap_or(ABSENT_SCOPE_COMPONENT) -} - -/// Backend-neutral storage key for primitive stores. -/// -/// Keys may be path-like, but they are not authority-bearing filesystem paths. -/// Domain stores own key grammar and scope semantics before constructing one. -#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize)] -#[serde(transparent)] -pub struct StorageKey(String); - -impl StorageKey { - pub fn new(value: impl Into) -> Result { - let value = value.into(); - validate_storage_key(&value)?; - Ok(Self(value)) - } - - pub fn as_str(&self) -> &str { - &self.0 - } -} - -impl<'de> Deserialize<'de> for StorageKey { - fn deserialize(deserializer: D) -> Result - where - D: Deserializer<'de>, - { - let value = String::deserialize(deserializer)?; - Self::new(value).map_err(|_| serde::de::Error::custom("invalid storage key")) - } -} - -impl AsRef for StorageKey { - fn as_ref(&self) -> &str { - self.as_str() - } -} - -impl std::fmt::Display for StorageKey { - fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - formatter.write_str(self.as_str()) - } -} - -/// Opaque backend version/fencing value for primitive compare-and-swap writes. -#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize)] -#[serde(transparent)] -pub struct StorageVersion(String); - -impl StorageVersion { - pub fn new(value: impl Into) -> Result { - let value = value.into(); - validate_storage_token(&value, "storage version", 128)?; - Ok(Self(value)) - } - - pub fn as_str(&self) -> &str { - &self.0 - } -} - -impl<'de> Deserialize<'de> for StorageVersion { - fn deserialize(deserializer: D) -> Result - where - D: Deserializer<'de>, - { - let value = String::deserialize(deserializer)?; - Self::new(value).map_err(|_| serde::de::Error::custom("invalid storage version")) - } -} - -impl AsRef for StorageVersion { - fn as_ref(&self) -> &str { - self.as_str() - } -} - -impl std::fmt::Display for StorageVersion { - fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - formatter.write_str(self.as_str()) - } -} - -/// Primitive write precondition shared by record/blob-like backends. -#[derive(Debug, Clone, Default, PartialEq, Eq)] -pub enum PutCondition { - #[default] - Any, - IfAbsent, - IfVersion(StorageVersion), -} - -impl PutCondition { - pub fn allows(&self, current: Option<&StorageVersion>) -> bool { - match self { - Self::Any => true, - Self::IfAbsent => current.is_none(), - Self::IfVersion(expected) => current == Some(expected), - } - } -} - -/// Blob payload returned by [`BlobStore`]. -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct StoredBlob { - pub key: StorageKey, - pub bytes: Vec, - pub version: StorageVersion, -} - -/// Validated structured JSON payload for [`RecordStore`] values. -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct RecordPayloadJson(String); - -impl RecordPayloadJson { - pub fn new(value: impl Into) -> Result { - let value = value.into(); - let _: IgnoredAny = - serde_json::from_str(&value).map_err(|_| StorageError::Serialization)?; - Ok(Self(value)) - } - - pub fn as_str(&self) -> &str { - &self.0 - } - - pub fn into_inner(self) -> String { - self.0 - } -} - -impl AsRef for RecordPayloadJson { - fn as_ref(&self) -> &str { - self.as_str() - } -} - -impl std::fmt::Display for RecordPayloadJson { - fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - formatter.write_str(self.as_str()) - } -} - -/// Structured JSON record returned by [`RecordStore`]. -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct StoredRecord { - pub key: StorageKey, - pub payload_json: RecordPayloadJson, - pub version: StorageVersion, -} - -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct PutBlobRequest { - pub key: StorageKey, - pub bytes: Vec, - pub condition: PutCondition, -} - -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct PutRecordRequest { - pub key: StorageKey, - pub payload_json: RecordPayloadJson, - pub condition: PutCondition, -} - -/// Primitive binary/object storage. Domain semantics live above this trait. -#[async_trait] -pub trait BlobStore: Send + Sync { - async fn put_blob(&self, request: PutBlobRequest) -> Result; - - async fn get_blob(&self, key: &StorageKey) -> Result, StorageError>; - - async fn delete_blob( - &self, - key: &StorageKey, - condition: PutCondition, - ) -> Result<(), StorageError>; -} - -/// Primitive keyed structured-record storage. Domain stores own schemas. -#[async_trait] -pub trait RecordStore: Send + Sync { - async fn put_record(&self, request: PutRecordRequest) -> Result; - - async fn get_record(&self, key: &StorageKey) -> Result, StorageError>; - - async fn delete_record( - &self, - key: &StorageKey, - condition: PutCondition, - ) -> Result<(), StorageError>; -} - -fn validate_storage_key(value: &str) -> Result<(), StorageError> { - validate_storage_token(value, "storage key", 512)?; - if value.starts_with('/') - || value.starts_with('\\') - || value.contains('\\') - || value.split('/').any(|segment| segment == "..") - || looks_like_windows_absolute_path(value) - { - return Err(StorageError::Validation); - } - Ok(()) -} - -fn validate_storage_token( - value: &str, - _label: &'static str, - max_bytes: usize, -) -> Result<(), StorageError> { - if value.is_empty() || value.len() > max_bytes { - return Err(StorageError::Validation); - } - if value - .chars() - .any(|character| character == '\0' || character.is_control()) - { - return Err(StorageError::Validation); - } - Ok(()) -} - -fn looks_like_windows_absolute_path(value: &str) -> bool { - let bytes = value.as_bytes(); - bytes.len() >= 3 - && bytes[1] == b':' - && (bytes[2] == b'\\' || bytes[2] == b'/') - && bytes[0].is_ascii_alphabetic() -} - -/// Bounded pagination helper for storage reads. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub struct PageLimit { - value: usize, -} - -impl PageLimit { - pub fn new(requested: usize, default: usize, max: usize) -> Self { - let max = max.max(1); - let default = default.clamp(1, max); - let value = if requested == 0 { - default - } else { - requested.min(max) - }; - Self { value } - } - - pub fn get(self) -> usize { - self.value - } -} - -/// Static migration descriptor used by storage adapters and composition code. -/// -/// The SQL remains owned by the domain adapter; this descriptor gives shared -/// diagnostics and migration registries a common shape. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub struct StorageMigration { - pub id: &'static str, - pub description: &'static str, - pub backend: StorageBackendKind, - pub sql: Option<&'static str>, -} - -#[cfg(test)] -mod tests { - use std::{ - collections::HashMap, - sync::{ - Mutex, - atomic::{AtomicU64, Ordering}, - }, - }; - - use super::*; - - #[derive(Debug, PartialEq, Eq, serde::Serialize, serde::Deserialize)] - struct Payload { - value: String, - } - - #[derive(Default)] - struct InMemoryPrimitiveStore { - blobs: Mutex>, - records: Mutex>, - next_version: AtomicU64, - } - - impl InMemoryPrimitiveStore { - fn next_storage_version(&self) -> StorageVersion { - let version = self.next_version.fetch_add(1, Ordering::Relaxed) + 1; - StorageVersion::new(format!("memory-{version}")) - .expect("generated storage version is valid") - } - } - - #[async_trait::async_trait] - impl BlobStore for InMemoryPrimitiveStore { - async fn put_blob(&self, request: PutBlobRequest) -> Result { - let mut blobs = self.blobs.lock().map_err(|_| StorageError::Backend)?; - let allowed = request - .condition - .allows(blobs.get(&request.key).map(|blob| &blob.version)); - if !allowed { - return Err(StorageError::Conflict); - } - let stored = StoredBlob { - key: request.key, - bytes: request.bytes, - version: self.next_storage_version(), - }; - blobs.insert(stored.key.clone(), stored.clone()); - Ok(stored) - } - - async fn get_blob(&self, key: &StorageKey) -> Result, StorageError> { - let blobs = self.blobs.lock().map_err(|_| StorageError::Backend)?; - Ok(blobs.get(key).cloned()) - } - - async fn delete_blob( - &self, - key: &StorageKey, - condition: PutCondition, - ) -> Result<(), StorageError> { - let mut blobs = self.blobs.lock().map_err(|_| StorageError::Backend)?; - if !condition.allows(blobs.get(key).map(|blob| &blob.version)) { - return Err(StorageError::Conflict); - } - blobs.remove(key); - Ok(()) - } - } - - #[async_trait::async_trait] - impl RecordStore for InMemoryPrimitiveStore { - async fn put_record( - &self, - request: PutRecordRequest, - ) -> Result { - let mut records = self.records.lock().map_err(|_| StorageError::Backend)?; - let allowed = request - .condition - .allows(records.get(&request.key).map(|record| &record.version)); - if !allowed { - return Err(StorageError::Conflict); - } - let stored = StoredRecord { - key: request.key, - payload_json: request.payload_json, - version: self.next_storage_version(), - }; - records.insert(stored.key.clone(), stored.clone()); - Ok(stored) - } - - async fn get_record(&self, key: &StorageKey) -> Result, StorageError> { - let records = self.records.lock().map_err(|_| StorageError::Backend)?; - Ok(records.get(key).cloned()) - } - - async fn delete_record( - &self, - key: &StorageKey, - condition: PutCondition, - ) -> Result<(), StorageError> { - let mut records = self.records.lock().map_err(|_| StorageError::Backend)?; - if !condition.allows(records.get(key).map(|record| &record.version)) { - return Err(StorageError::Conflict); - } - records.remove(key); - Ok(()) - } - } - - #[test] - fn json_helpers_round_trip_without_domain_semantics() { - let payload = Payload { - value: "hello".to_string(), - }; - - let encoded = encode_json(&payload).expect("test payload serializes"); - let decoded: Payload = decode_json(&encoded).expect("test payload deserializes"); - - assert_eq!(decoded, payload); - } - - #[test] - fn json_decode_error_is_redacted() { - let error = decode_json::("{RAW_SECRET").unwrap_err(); - - assert_eq!(error, StorageError::Serialization); - assert!(!format!("{error:?}").contains("RAW_SECRET")); - assert!(!error.to_string().contains("RAW_SECRET")); - } - - #[test] - fn backend_error_discards_raw_detail() { - let error = redacted_backend_error("host path /tmp/secret.db failed"); - - assert_eq!(error, StorageError::Backend); - assert!(!format!("{error:?}").contains("/tmp/secret.db")); - assert!(!error.to_string().contains("/tmp/secret.db")); - } - - #[test] - fn storage_key_and_version_reject_empty_control_and_oversized_values() { - assert_eq!( - StorageKey::new("thread/message").unwrap().as_str(), - "thread/message" - ); - assert_eq!(StorageKey::new("").unwrap_err(), StorageError::Validation); - assert_eq!( - StorageKey::new("bad\nkey").unwrap_err(), - StorageError::Validation - ); - assert_eq!( - StorageKey::new("x".repeat(513)).unwrap_err(), - StorageError::Validation - ); - assert_eq!(StorageVersion::new("v1").unwrap().as_str(), "v1"); - assert_eq!( - StorageVersion::new("v".repeat(129)).unwrap_err(), - StorageError::Validation - ); - } - - #[test] - fn storage_key_and_version_reject_invalid_deserialized_values() { - assert!(serde_json::from_str::("\"\"").is_err()); - assert!(serde_json::from_str::(&format!("\"{}\"", "x".repeat(513))).is_err()); - assert!(serde_json::from_str::("\"bad\\nkey\"").is_err()); - assert!(serde_json::from_str::("\"\"").is_err()); - assert!( - serde_json::from_str::(&format!("\"{}\"", "v".repeat(129))).is_err() - ); - } - - #[test] - fn storage_version_serializes_and_deserializes_transparently() { - let version = StorageVersion::new("opaque-v1").unwrap(); - - let encoded = serde_json::to_string(&version).unwrap(); - let decoded: StorageVersion = serde_json::from_str(&encoded).unwrap(); - - assert_eq!(encoded, "\"opaque-v1\""); - assert_eq!(decoded, version); - } - - #[test] - fn storage_key_rejects_path_traversal_and_platform_absolute_forms() { - for invalid in ["../x", "/x", "a/../../b", "a\\..\\b", "C:\\secrets"] { - assert_eq!( - StorageKey::new(invalid).unwrap_err(), - StorageError::Validation - ); - } - assert_eq!( - StorageKey::new("version..2").unwrap().as_str(), - "version..2" - ); - } - - #[test] - fn record_payload_json_rejects_malformed_payloads() { - assert!(RecordPayloadJson::new("{\"ok\":true}").is_ok()); - assert_eq!( - RecordPayloadJson::new("not json").unwrap_err(), - StorageError::Serialization - ); - } - - #[test] - fn put_conditions_encode_cas_without_domain_semantics() { - let current = StorageVersion::new("v1").unwrap(); - let other = StorageVersion::new("v2").unwrap(); - - assert!(PutCondition::Any.allows(None)); - assert!(PutCondition::Any.allows(Some(¤t))); - assert!(PutCondition::IfAbsent.allows(None)); - assert!(!PutCondition::IfAbsent.allows(Some(¤t))); - assert!(PutCondition::IfVersion(current.clone()).allows(Some(¤t))); - assert!(!PutCondition::IfVersion(other).allows(Some(¤t))); - assert!(!PutCondition::IfVersion(current).allows(None)); - } - - #[tokio::test] - async fn in_memory_store_exercises_blob_and_record_traits() { - let store = InMemoryPrimitiveStore::default(); - let blob_key = StorageKey::new("blobs/example").unwrap(); - - assert!(store.get_blob(&blob_key).await.unwrap().is_none()); - let first_blob = store - .put_blob(PutBlobRequest { - key: blob_key.clone(), - bytes: b"first".to_vec(), - condition: PutCondition::IfAbsent, - }) - .await - .unwrap(); - assert_eq!(first_blob.bytes, b"first"); - assert_eq!( - store - .put_blob(PutBlobRequest { - key: blob_key.clone(), - bytes: b"conflict".to_vec(), - condition: PutCondition::IfAbsent, - }) - .await - .unwrap_err(), - StorageError::Conflict - ); - let updated_blob = store - .put_blob(PutBlobRequest { - key: blob_key.clone(), - bytes: b"second".to_vec(), - condition: PutCondition::IfVersion(first_blob.version.clone()), - }) - .await - .unwrap(); - assert_eq!(updated_blob.bytes, b"second"); - assert_ne!(updated_blob.version, first_blob.version); - assert_eq!( - store.get_blob(&blob_key).await.unwrap().unwrap(), - updated_blob - ); - store - .delete_blob(&blob_key, PutCondition::Any) - .await - .unwrap(); - assert!(store.get_blob(&blob_key).await.unwrap().is_none()); - - let record_key = StorageKey::new("records/example").unwrap(); - let payload = RecordPayloadJson::new(r#"{"value":"hello"}"#).unwrap(); - let stored_record = store - .put_record(PutRecordRequest { - key: record_key.clone(), - payload_json: payload.clone(), - condition: PutCondition::IfAbsent, - }) - .await - .unwrap(); - assert_eq!(stored_record.payload_json, payload); - assert_eq!( - store.get_record(&record_key).await.unwrap().unwrap(), - stored_record - ); - store - .delete_record(&record_key, PutCondition::Any) - .await - .unwrap(); - assert!(store.get_record(&record_key).await.unwrap().is_none()); - } - - #[test] - fn page_limit_applies_default_and_max_bounds() { - assert_eq!(PageLimit::new(0, 50, 100).get(), 50); - assert_eq!(PageLimit::new(500, 50, 100).get(), 100); - assert_eq!(PageLimit::new(10, 50, 100).get(), 10); - assert_eq!(PageLimit::new(0, 0, 0).get(), 1); - } -}