diff --git a/src/builder.rs b/src/builder.rs index f0f38783fb..d68cf4db33 100644 --- a/src/builder.rs +++ b/src/builder.rs @@ -60,6 +60,7 @@ use crate::event::EventQueue; use crate::fee_estimator::OnchainFeeEstimator; use crate::gossip::GossipSource; use crate::io::sqlite_store::SqliteStore; +use crate::io::tier_store::{setup_index_store, TierStore}; use crate::io::utils::{ open_or_migrate_fs_store, read_all_objects, read_event_queue, read_external_pathfinding_scores_from_cache, read_n_objects, read_network_graph, @@ -158,6 +159,12 @@ impl std::fmt::Debug for LogWriterConfig { } } +#[derive(Default, Debug)] +struct TierStoreConfig { + ephemeral_storage_dir_path: Option, + backup_storage_dir_path: Option, +} + /// An error encountered during building a [`Node`]. /// /// [`Node`]: crate::Node @@ -311,6 +318,7 @@ pub struct NodeBuilder { liquidity_source_config: Option, log_writer_config: Option, async_payments_role: Option, + tier_store_config: Option, runtime_handle: Option, pathfinding_scores_sync_config: Option, probing_config: Option, @@ -329,6 +337,7 @@ impl NodeBuilder { let gossip_source_config = None; let liquidity_source_config = None; let log_writer_config = None; + let tier_store_config = None; let runtime_handle = None; let pathfinding_scores_sync_config = None; let probing_config = None; @@ -338,6 +347,7 @@ impl NodeBuilder { gossip_source_config, liquidity_source_config, log_writer_config, + tier_store_config, runtime_handle, async_payments_role: None, pathfinding_scores_sync_config, @@ -663,6 +673,41 @@ impl NodeBuilder { self } + /// Configures a local SQLite backup store for disaster recovery. + /// + /// When building with tiered storage, a SQLite store will be created at the + /// given directory path using [`SQLITE_BACKUP_DB_FILE_NAME`] as its database + /// file name. It receives a second durable copy of data written to the + /// primary store. + /// + /// Writes and removals for primary-backed data only succeed once both the + /// primary and backup SQLite stores complete successfully. + /// + /// If not set, durable data will be stored only in the primary store. + /// + /// [`SQLITE_BACKUP_DB_FILE_NAME`]: crate::io::sqlite_store::SQLITE_BACKUP_DB_FILE_NAME + #[cfg(not(feature = "uniffi"))] + pub fn set_backup_storage_dir_path(&mut self, backup_storage_dir_path: String) -> &mut Self { + let tier_store_config = self.tier_store_config.get_or_insert(TierStoreConfig::default()); + tier_store_config.backup_storage_dir_path = Some(backup_storage_dir_path.into()); + self + } + + /// Configures the ephemeral storage directory path for non-critical, frequently-accessed data. + /// + /// When set, a local SQLite store is created at this path for ephemeral data like + /// the network graph and scorer. Data stored here can be rebuilt if lost. + /// + /// If not set, non-critical data will be stored in the primary store. + #[cfg(not(feature = "uniffi"))] + pub fn set_ephemeral_storage_dir_path( + &mut self, ephemeral_storage_dir_path: String, + ) -> &mut Self { + let tier_store_config = self.tier_store_config.get_or_insert(TierStoreConfig::default()); + tier_store_config.ephemeral_storage_dir_path = Some(ephemeral_storage_dir_path.into()); + self + } + /// Builds a [`Node`] instance with a [`SqliteStore`] backend and according to the options /// previously configured. pub fn build(&self, node_entropy: NodeEntropy) -> Result { @@ -872,11 +917,18 @@ impl NodeBuilder { } /// Builds a [`Node`] instance according to the options previously configured. + /// + /// The provided `kv_store` will be used as the primary storage backend. Optionally, + /// an ephemeral store for frequently-accessed non-critical data (e.g., network graph, scorer) + /// and a local SQLite backup store for disaster recovery can be configured via + /// [`set_ephemeral_storage_dir_path`] and [`set_backup_storage_dir_path`]. + /// + /// [`set_ephemeral_storage_dir_path`]: Self::set_ephemeral_storage_dir_path + /// [`set_backup_storage_dir_path`]: Self::set_backup_storage_dir_path pub fn build_with_store( &self, node_entropy: NodeEntropy, kv_store: S, ) -> Result { let logger = setup_logger(&self.log_writer_config, &self.config)?; - self.build_with_store_and_logger(node_entropy, kv_store, logger) } @@ -901,6 +953,46 @@ impl NodeBuilder { fn build_with_store_runtime_and_logger( &self, node_entropy: NodeEntropy, kv_store: S, runtime: Arc, logger: Arc, ) -> Result { + let ts_config = self.tier_store_config.as_ref(); + let primary_store = Arc::new(DynStoreWrapper(kv_store)); + let mut tier_store = TierStore::new(primary_store, Arc::clone(&logger)); + if let Some(config) = ts_config { + if let Some(ephemeral_storage_dir_path) = config.ephemeral_storage_dir_path.as_ref() { + let index_store = runtime + .block_on(setup_index_store(self.config.storage_dir_path.clone().into())) + .map_err(|e| { + log_error!(logger, "Failed to setup tier-store index: {}", e); + BuildError::KVStoreSetupFailed + })?; + let ephemeral_store = SqliteStore::new( + ephemeral_storage_dir_path.clone(), + Some(io::sqlite_store::SQLITE_EPHEMERAL_DB_FILE_NAME.to_string()), + Some(io::sqlite_store::KV_TABLE_NAME.to_string()), + ) + .map_err(|e| { + log_error!(logger, "Failed to setup ephemeral SQLite store: {}", e); + BuildError::KVStoreSetupFailed + })?; + let ephemeral_store: Arc = Arc::new(DynStoreWrapper(ephemeral_store)); + tier_store.set_index_store(index_store); + tier_store.set_ephemeral_store(ephemeral_store); + } + + if let Some(backup_storage_dir_path) = config.backup_storage_dir_path.as_ref() { + let backup_store = SqliteStore::new( + backup_storage_dir_path.clone(), + Some(io::sqlite_store::SQLITE_BACKUP_DB_FILE_NAME.to_string()), + Some(io::sqlite_store::KV_TABLE_NAME.to_string()), + ) + .map_err(|e| { + log_error!(logger, "Failed to setup backup SQLite store: {}", e); + BuildError::KVStoreSetupFailed + })?; + let backup_store: Arc = Arc::new(DynStoreWrapper(backup_store)); + tier_store.set_backup_store(backup_store); + } + } + let seed_bytes = node_entropy.to_seed_bytes(); let config = Arc::new(self.config.clone()); @@ -915,7 +1007,7 @@ impl NodeBuilder { seed_bytes, runtime, logger, - Arc::new(DynStoreWrapper(kv_store)), + Arc::new(DynStoreWrapper(tier_store)), ) } } diff --git a/src/io/mod.rs b/src/io/mod.rs index c70c68d96e..b7856e1b8a 100644 --- a/src/io/mod.rs +++ b/src/io/mod.rs @@ -12,6 +12,7 @@ pub mod postgres_store; pub mod sqlite_store; #[cfg(test)] pub(crate) mod test_utils; +pub(crate) mod tier_store; pub(crate) mod utils; pub mod vss_store; diff --git a/src/io/sqlite_store/mod.rs b/src/io/sqlite_store/mod.rs index 2587220598..6702f842de 100644 --- a/src/io/sqlite_store/mod.rs +++ b/src/io/sqlite_store/mod.rs @@ -12,12 +12,14 @@ use std::future::Future; use std::path::PathBuf; use std::sync::atomic::{AtomicI64, AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; +use std::time::Duration; use lightning::io; use lightning::util::persist::{ KVStore, MigratableKVStore, PageToken, PaginatedKVStore, PaginatedListResponse, }; use lightning_types::string::PrintableString; +use rusqlite::ffi::ErrorCode; use rusqlite::{named_params, Connection}; use crate::io::utils::check_namespace_key_validity; @@ -26,6 +28,12 @@ mod migrations; /// LDK Node's database file name. pub const SQLITE_DB_FILE_NAME: &str = "ldk_node_data.sqlite"; +/// LDK Node's internal tier-store index database file name. +pub(crate) const SQLITE_TIER_INDEX_DB_FILE_NAME: &str = "ldk_node_tier_index.sqlite"; +/// LDK Node's backup database file name. +pub const SQLITE_BACKUP_DB_FILE_NAME: &str = "ldk_node_data_backup.sqlite"; +/// LDK Node's ephemeral database file name. +pub const SQLITE_EPHEMERAL_DB_FILE_NAME: &str = "ldk_node_data_ephemeral.sqlite"; /// LDK Node's table in which we store all data. pub const KV_TABLE_NAME: &str = "ldk_node_data"; @@ -41,6 +49,17 @@ const SCHEMA_USER_VERSION: u16 = 3; // The number of entries returned per page in paginated list operations. const PAGE_SIZE: usize = 50; +fn exclusive_lock_error_kind(error: &rusqlite::Error) -> io::ErrorKind { + match error { + rusqlite::Error::SqliteFailure(error, _) + if matches!(error.code, ErrorCode::DatabaseBusy | ErrorCode::DatabaseLocked) => + { + io::ErrorKind::AlreadyExists + }, + _ => io::ErrorKind::Other, + } +} + /// A [`KVStore`] implementation that writes to and reads from an [SQLite] database. /// /// [SQLite]: https://sqlite.org @@ -62,7 +81,22 @@ impl SqliteStore { pub fn new( data_dir: PathBuf, db_file_name: Option, kv_table_name: Option, ) -> io::Result { - let inner = Arc::new(SqliteStoreInner::new(data_dir, db_file_name, kv_table_name)?); + Self::new_internal(data_dir, db_file_name, kv_table_name, false) + } + + /// Constructs a new [`SqliteStore`] which exclusively owns its database for its lifetime. + pub(crate) fn new_exclusive( + data_dir: PathBuf, db_file_name: Option, kv_table_name: Option, + ) -> io::Result { + Self::new_internal(data_dir, db_file_name, kv_table_name, true) + } + + fn new_internal( + data_dir: PathBuf, db_file_name: Option, kv_table_name: Option, + exclusive: bool, + ) -> io::Result { + let inner = + Arc::new(SqliteStoreInner::new(data_dir, db_file_name, kv_table_name, exclusive)?); let next_write_version = AtomicU64::new(1); Ok(Self { inner, next_write_version }) @@ -230,6 +264,7 @@ struct SqliteStoreInner { impl SqliteStoreInner { fn new( data_dir: PathBuf, db_file_name: Option, kv_table_name: Option, + exclusive: bool, ) -> io::Result { let db_file_name = db_file_name.unwrap_or(DEFAULT_SQLITE_DB_FILE_NAME.to_string()); let kv_table_name = kv_table_name.unwrap_or(DEFAULT_KV_TABLE_NAME.to_string()); @@ -251,6 +286,33 @@ impl SqliteStoreInner { io::Error::new(io::ErrorKind::Other, msg) })?; + if exclusive { + connection.busy_timeout(Duration::ZERO).map_err(|e| { + let msg = format!( + "Failed to configure exclusive database lock timeout for {}: {}", + db_file_path.display(), + e + ); + io::Error::new(io::ErrorKind::Other, msg) + })?; + connection.pragma_update(None, "locking_mode", "EXCLUSIVE").map_err(|e| { + let msg = format!( + "Failed to exclusively lock database file {}: {}", + db_file_path.display(), + e + ); + io::Error::new(io::ErrorKind::Other, msg) + })?; + connection.execute_batch("BEGIN EXCLUSIVE; COMMIT;").map_err(|e| { + let msg = format!( + "Failed to exclusively lock database file {}: {}", + db_file_path.display(), + e + ); + io::Error::new(exclusive_lock_error_kind(&e), msg) + })?; + } + let sql = format!("SELECT user_version FROM pragma_user_version"); let version_res: u16 = connection.query_row(&sql, [], |row| row.get(0)).map_err(|e| { let msg = format!("Failed to read PRAGMA user_version: {}", e); @@ -700,6 +762,21 @@ mod tests { } } + #[test] + fn exclusive_lock_error_kind_distinguishes_contention_from_other_failures() { + for result_code in [rusqlite::ffi::SQLITE_BUSY, rusqlite::ffi::SQLITE_LOCKED] { + let error = + rusqlite::Error::SqliteFailure(rusqlite::ffi::Error::new(result_code), None); + assert_eq!(exclusive_lock_error_kind(&error), io::ErrorKind::AlreadyExists); + } + + let io_error = rusqlite::Error::SqliteFailure( + rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_IOERR), + None, + ); + assert_eq!(exclusive_lock_error_kind(&io_error), io::ErrorKind::Other); + } + #[tokio::test] async fn read_write_remove_list_persist() { let mut temp_path = random_storage_path(); diff --git a/src/io/tier_store.rs b/src/io/tier_store.rs new file mode 100644 index 0000000000..fb63f032b1 --- /dev/null +++ b/src/io/tier_store.rs @@ -0,0 +1,3177 @@ +// This file is Copyright its original authors, visible in version control history. +// +// This file is licensed under the Apache License, Version 2.0 or the MIT license , at your option. You may not use this file except in +// accordance with one or both of these licenses. + +use std::collections::HashMap; +use std::future::Future; +use std::path::PathBuf; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::{Arc, Mutex}; + +use bitcoin::hashes::{sha256, Hash, HashEngine}; +use bitcoin::hex::{DisplayHex, FromHex}; +use lightning::util::persist::{ + KVStore, PageToken, PaginatedKVStore, PaginatedListResponse, NETWORK_GRAPH_PERSISTENCE_KEY, + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, SCORER_PERSISTENCE_KEY, + SCORER_PERSISTENCE_PRIMARY_NAMESPACE, +}; +use lightning::util::ser::{Readable, Writeable}; +use lightning::{impl_writeable_tlv_based, impl_writeable_tlv_based_enum, io, log_error}; +use tokio::sync::Mutex as TokioMutex; + +use crate::io::sqlite_store::{SqliteStore, KV_TABLE_NAME, SQLITE_TIER_INDEX_DB_FILE_NAME}; +use crate::io::utils::{check_namespace_key_validity, EXTERNAL_PATHFINDING_SCORES_CACHE_KEY}; +use crate::logger::{LdkLogger, Logger}; +use crate::types::{DynStore, DynStoreWrapper}; + +const INDEX_DATABASE_ID_LEN: usize = 16; +const PAGE_TOKEN_FORMAT_VERSION: u8 = 1; +const INDEX_ENTRIES_PRIMARY_NAMESPACE: &str = "_tier_store_entries"; +const INDEX_JOURNAL_PRIMARY_NAMESPACE: &str = "_tier_store_journal"; +const INDEX_METADATA_PRIMARY_NAMESPACE: &str = "_tier_store_metadata"; +const INDEX_DATABASE_ID_KEY: &str = "index_database_id"; +const INDEX_NAMESPACE_READY_KEY_PREFIX: &str = "ready_"; +const INDEX_CACHE_READY_KEY_PREFIX: &str = "cache_ready_"; +const INDEX_ENTRY_VALUE: &[u8] = &[1]; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum ValueTier { + Primary, + Ephemeral, +} + +impl_writeable_tlv_based_enum!(ValueTier, + (0, Primary) => {}, + (2, Ephemeral) => {}, +); + +#[derive(Debug, PartialEq, Eq)] +enum JournalOperation { + Create { + /// The complete intended value, retained locally so recovery does not depend on which remote + /// or local value-store write completed before interruption. + value: Vec, + }, + Remove { + lazy: bool, + }, +} + +impl_writeable_tlv_based_enum!(JournalOperation, + (0, Create) => { + (0, value, required), + }, + (2, Remove) => { + (0, lazy, required), + }, +); + +#[derive(Debug, PartialEq, Eq)] +struct JournalEntry { + primary_namespace: String, + secondary_namespace: String, + key: String, + tier: ValueTier, + requires_backup: bool, + operation: JournalOperation, +} + +impl_writeable_tlv_based!(JournalEntry, { + (0, primary_namespace, required), + (2, secondary_namespace, required), + (4, key, required), + (6, tier, required), + (8, requires_backup, required), + (10, operation, required), +}); + +impl JournalEntry { + /// Encodes a pending operation for persistence in the local index database. + fn serialize(&self) -> Vec { + Writeable::encode(self) + } + + /// Decodes and validates a pending operation from the local index database. + fn deserialize(encoded: &[u8]) -> io::Result { + Readable::read(&mut &*encoded) + .map_err(|_| io::Error::new(io::ErrorKind::InvalidData, "Invalid journal entry")) + } +} + +struct TierStorePageToken { + format_version: u8, + index_database_id: Vec, + namespace_id: String, + index_page_token: String, +} + +impl_writeable_tlv_based!(TierStorePageToken, { + (0, format_version, required), + (2, index_database_id, required), + (4, namespace_id, required), + (6, index_page_token, required), +}); + +impl TierStorePageToken { + /// Wraps an index-store page token with the context required to validate its later use. + fn encode( + index_database_id: &[u8; INDEX_DATABASE_ID_LEN], namespace_id: String, + index_page_token: PageToken, + ) -> PageToken { + let token = Self { + format_version: PAGE_TOKEN_FORMAT_VERSION, + index_database_id: index_database_id.to_vec(), + namespace_id, + index_page_token: index_page_token.to_string(), + }; + PageToken::new(Writeable::encode(&token).to_lower_hex_string()) + } + + /// Decodes a TierStore token and rejects tokens issued for another index or namespace. + fn decode( + token: PageToken, expected_database_id: &[u8; INDEX_DATABASE_ID_LEN], + expected_namespace_id: &str, + ) -> io::Result { + let encoded = Vec::from_hex(token.as_str()).map_err(|_| { + io::Error::new(io::ErrorKind::InvalidInput, "Invalid TierStore page token") + })?; + let mut reader = &*encoded; + let token: Self = Readable::read(&mut reader).map_err(|_| { + io::Error::new(io::ErrorKind::InvalidInput, "Invalid TierStore page token") + })?; + if !reader.is_empty() + || token.format_version != PAGE_TOKEN_FORMAT_VERSION + || token.index_database_id != expected_database_id + || token.namespace_id != expected_namespace_id + { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + "TierStore page token does not belong to this index namespace", + )); + } + Ok(PageToken::new(token.index_page_token)) + } +} + +pub(crate) struct TierStoreIndex { + // Holding the store keeps its exclusive SQLite lock for the lifetime of the tier store. + store: Arc, + database_id: [u8; INDEX_DATABASE_ID_LEN], +} + +impl TierStoreIndex { + /// Opens the internal SQLite index and ensures that it has a persistent database identity. + async fn new(data_dir: PathBuf) -> io::Result { + let store = SqliteStore::new_exclusive( + data_dir, + Some(SQLITE_TIER_INDEX_DB_FILE_NAME.to_string()), + Some(KV_TABLE_NAME.to_string()), + )?; + let store: Arc = Arc::new(DynStoreWrapper(store)); + let database_id = Self::read_or_create_database_id(store.as_ref()).await?; + Ok(Self { store, database_id }) + } + + /// Constructs an index over the supplied store for tests that do not need SQLite persistence. + #[cfg(test)] + fn from_store(store: Arc) -> Self { + Self { store, database_id: [1; INDEX_DATABASE_ID_LEN] } + } + + /// Constructs a test index with a specified database identity. + #[cfg(test)] + fn from_store_with_database_id( + store: Arc, database_id: [u8; INDEX_DATABASE_ID_LEN], + ) -> Self { + Self { store, database_id } + } + + /// Derives the internal secondary namespace for a logical namespace pair. + /// + /// Length-prefixing keeps distinct pairs from producing the same hash input. The original pair + /// is also stored as namespace metadata so that a hash collision can be detected. + fn namespace_id(primary_namespace: &str, secondary_namespace: &str) -> String { + let mut engine = sha256::Hash::engine(); + engine.input(&(primary_namespace.len() as u64).to_be_bytes()); + engine.input(primary_namespace.as_bytes()); + engine.input(&(secondary_namespace.len() as u64).to_be_bytes()); + engine.input(secondary_namespace.as_bytes()); + sha256::Hash::from_engine(engine).to_string() + } + + /// Derives the metadata key that records whether a logical namespace is index-backed. + fn namespace_ready_key(primary_namespace: &str, secondary_namespace: &str) -> String { + format!( + "{}{}", + INDEX_NAMESPACE_READY_KEY_PREFIX, + Self::namespace_id(primary_namespace, secondary_namespace) + ) + } + + /// Derives the metadata key recording that indexed cache values occupy the ephemeral tier. + fn cache_ready_key(primary_namespace: &str, secondary_namespace: &str) -> String { + format!( + "{}{}", + INDEX_CACHE_READY_KEY_PREFIX, + Self::namespace_id(primary_namespace, secondary_namespace) + ) + } + + /// Encodes the original logical namespace pair for collision detection. + fn namespace_metadata(primary_namespace: &str, secondary_namespace: &str) -> Vec { + // The fixed five-byte overhead is one format-version byte plus two big-endian u16 + // namespace-length prefixes: 1 + 2 + 2 = 5. + let mut metadata = + Vec::with_capacity(5 + primary_namespace.len() + secondary_namespace.len()); + metadata.push(1); + metadata.extend_from_slice(&(primary_namespace.len() as u16).to_be_bytes()); + metadata.extend_from_slice(primary_namespace.as_bytes()); + metadata.extend_from_slice(&(secondary_namespace.len() as u16).to_be_bytes()); + metadata.extend_from_slice(secondary_namespace.as_bytes()); + metadata + } + + /// Returns whether the namespace's index is authoritative for listing. + /// + /// Returns an error if the stored namespace metadata does not match the requested namespace, + /// which indicates either corrupt metadata or a namespace-ID collision. + async fn is_namespace_ready( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> io::Result { + self.is_namespace_marker_set( + &Self::namespace_ready_key(primary_namespace, secondary_namespace), + primary_namespace, + secondary_namespace, + ) + .await + } + + /// Returns whether indexed cache values have been reconciled into ephemeral storage. + async fn is_cache_ready( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> io::Result { + self.is_namespace_marker_set( + &Self::cache_ready_key(primary_namespace, secondary_namespace), + primary_namespace, + secondary_namespace, + ) + .await + } + + async fn is_namespace_marker_set( + &self, marker_key: &str, primary_namespace: &str, secondary_namespace: &str, + ) -> io::Result { + match KVStore::read(self.store.as_ref(), INDEX_METADATA_PRIMARY_NAMESPACE, "", marker_key) + .await + { + Ok(metadata) + if metadata == Self::namespace_metadata(primary_namespace, secondary_namespace) => + { + Ok(true) + }, + Ok(_) => Err(io::Error::new( + io::ErrorKind::InvalidData, + "Tier-store index namespace collision", + )), + Err(e) if e.kind() == io::ErrorKind::NotFound => Ok(false), + Err(e) => Err(e), + } + } + + /// Marks the namespace's index as authoritative by persisting its original namespace pair. + async fn mark_namespace_ready( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> io::Result<()> { + self.set_namespace_marker( + &Self::namespace_ready_key(primary_namespace, secondary_namespace), + primary_namespace, + secondary_namespace, + ) + .await + } + + /// Marks indexed cache values as reconciled into ephemeral storage. + async fn mark_cache_ready( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> io::Result<()> { + self.set_namespace_marker( + &Self::cache_ready_key(primary_namespace, secondary_namespace), + primary_namespace, + secondary_namespace, + ) + .await + } + + async fn set_namespace_marker( + &self, marker_key: &str, primary_namespace: &str, secondary_namespace: &str, + ) -> io::Result<()> { + KVStore::write( + self.store.as_ref(), + INDEX_METADATA_PRIMARY_NAMESPACE, + "", + marker_key, + Self::namespace_metadata(primary_namespace, secondary_namespace), + ) + .await + } + + /// Adds a logical key to the namespace's ordered listing index. + /// + /// Rewriting an existing key preserves its original index-store creation order. + async fn write_entry( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, + ) -> io::Result<()> { + KVStore::write( + self.store.as_ref(), + INDEX_ENTRIES_PRIMARY_NAMESPACE, + &Self::namespace_id(primary_namespace, secondary_namespace), + key, + INDEX_ENTRY_VALUE.to_vec(), + ) + .await + } + + /// Returns whether the namespace's listing index contains the logical key. + async fn contains_entry( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, + ) -> io::Result { + match KVStore::read( + self.store.as_ref(), + INDEX_ENTRIES_PRIMARY_NAMESPACE, + &Self::namespace_id(primary_namespace, secondary_namespace), + key, + ) + .await + { + Ok(value) if value == INDEX_ENTRY_VALUE => Ok(true), + Ok(_) => Err(io::Error::new(io::ErrorKind::InvalidData, "Invalid index entry")), + Err(e) if e.kind() == io::ErrorKind::NotFound => Ok(false), + Err(e) => Err(e), + } + } + + /// Removes a logical key from the namespace's listing index. + async fn remove_entry( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, lazy: bool, + ) -> io::Result<()> { + KVStore::remove( + self.store.as_ref(), + INDEX_ENTRIES_PRIMARY_NAMESPACE, + &Self::namespace_id(primary_namespace, secondary_namespace), + key, + lazy, + ) + .await + } + + /// Lists all logical keys recorded in the namespace's index. + async fn list( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> io::Result> { + KVStore::list( + self.store.as_ref(), + INDEX_ENTRIES_PRIMARY_NAMESPACE, + &Self::namespace_id(primary_namespace, secondary_namespace), + ) + .await + } + + /// Lists logical keys in the index store's creation order using a namespace-bound token. + async fn list_paginated( + &self, primary_namespace: &str, secondary_namespace: &str, page_token: Option, + ) -> io::Result { + let namespace_id = Self::namespace_id(primary_namespace, secondary_namespace); + let index_page_token = page_token + .map(|token| TierStorePageToken::decode(token, &self.database_id, &namespace_id)) + .transpose()?; + let response = PaginatedKVStore::list_paginated( + self.store.as_ref(), + INDEX_ENTRIES_PRIMARY_NAMESPACE, + &namespace_id, + index_page_token, + ) + .await?; + let next_page_token = response + .next_page_token + .map(|token| TierStorePageToken::encode(&self.database_id, namespace_id, token)); + Ok(PaginatedListResponse { keys: response.keys, next_page_token }) + } + + /// Persists a pending operation before its value-store effects begin. + async fn write_journal_entry(&self, entry: &JournalEntry) -> io::Result<()> { + KVStore::write( + self.store.as_ref(), + INDEX_JOURNAL_PRIMARY_NAMESPACE, + &Self::namespace_id(&entry.primary_namespace, &entry.secondary_namespace), + &entry.key, + entry.serialize(), + ) + .await + } + + /// Reads a pending operation for a logical key. + async fn read_journal_entry( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, + ) -> io::Result { + let encoded = KVStore::read( + self.store.as_ref(), + INDEX_JOURNAL_PRIMARY_NAMESPACE, + &Self::namespace_id(primary_namespace, secondary_namespace), + key, + ) + .await?; + let entry = JournalEntry::deserialize(&encoded)?; + if entry.primary_namespace != primary_namespace + || entry.secondary_namespace != secondary_namespace + || entry.key != key + { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "Journal entry identity does not match its location", + )); + } + Ok(entry) + } + + /// Lists pending-operation keys from oldest to newest journal insertion. + async fn list_journal_entries( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> io::Result> { + let secondary_namespace = Self::namespace_id(primary_namespace, secondary_namespace); + let mut keys = Vec::new(); + let mut page_token = None; + loop { + let page = PaginatedKVStore::list_paginated( + self.store.as_ref(), + INDEX_JOURNAL_PRIMARY_NAMESPACE, + &secondary_namespace, + page_token, + ) + .await?; + keys.extend(page.keys); + match page.next_page_token { + Some(next_page_token) => page_token = Some(next_page_token), + None => break, + } + } + keys.reverse(); + Ok(keys) + } + + /// Clears a pending operation after all of its intended effects are durable. + async fn remove_journal_entry( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, + ) -> io::Result<()> { + KVStore::remove( + self.store.as_ref(), + INDEX_JOURNAL_PRIMARY_NAMESPACE, + &Self::namespace_id(primary_namespace, secondary_namespace), + key, + false, + ) + .await + } + + /// Reads the index database identity, creating and persisting one when it is absent. + async fn read_or_create_database_id( + store: &DynStore, + ) -> io::Result<[u8; INDEX_DATABASE_ID_LEN]> { + match KVStore::read(store, INDEX_METADATA_PRIMARY_NAMESPACE, "", INDEX_DATABASE_ID_KEY) + .await + { + Ok(bytes) => bytes.try_into().map_err(|_| { + io::Error::new(io::ErrorKind::InvalidData, "Invalid tier-store index database ID") + }), + Err(e) if e.kind() == io::ErrorKind::NotFound => { + let mut database_id = [0; INDEX_DATABASE_ID_LEN]; + getrandom::fill(&mut database_id).map_err(|e| { + io::Error::new( + io::ErrorKind::Other, + format!("Failed to generate tier-store index database ID: {e}"), + ) + })?; + KVStore::write( + store, + INDEX_METADATA_PRIMARY_NAMESPACE, + "", + INDEX_DATABASE_ID_KEY, + database_id.to_vec(), + ) + .await?; + Ok(database_id) + }, + Err(e) => Err(e), + } + } +} + +/// A 3-tiered [`KVStore`] implementation that routes data across +/// storage backends that may be local or remote: +/// - a primary store for durable, authoritative persistence, +/// - an optional backup store that maintains an additional durable copy of +/// primary-backed data, and +/// - an optional ephemeral store for non-critical, rebuildable cached data. +/// +/// When a backup store is configured, writes and removals for primary-backed data +/// are issued to the primary and backup stores concurrently and only succeed once +/// both stores complete successfully. +/// +/// Reads and lists do not consult the backup store during normal operation. +/// Ephemeral data is read from and written to the ephemeral store when configured. +/// Namespaces are indexed locally so unpaginated and paginated listings expose the +/// same logical contents in cross-tier creation order. Existing primary-store keys +/// are imported into the index before a namespace is first accessed. Cache values +/// imported from primary storage are then moved to ephemeral storage without changing +/// their index positions. +/// +/// Note that dual-store writes and removals are not atomic across the primary and +/// backup stores. If one store succeeds and the other fails, the operation +/// returns an error even though one store may already reflect the change. +pub(crate) struct TierStore { + inner: Arc, +} + +impl TierStore { + pub fn new(primary_store: Arc, logger: Arc) -> Self { + let inner = Arc::new(TierStoreInner::new(primary_store, Arc::clone(&logger))); + + Self { inner } + } + + /// Configures a backup store for primary-backed data. + /// + /// Once set, writes and removals targeting the primary tier succeed only if both + /// the primary and backup stores succeed. The two operations are issued + /// concurrently, and any failure is returned to the caller. + /// + /// Note: dual-store writes/removals are not atomic. An error may be returned + /// after the primary store has already been updated if the backup store fails. + /// + /// The backup store is not consulted for normal reads or lists. + pub fn set_backup_store(&mut self, backup: Arc) { + debug_assert_eq!(Arc::strong_count(&self.inner), 1); + + let inner = Arc::get_mut(&mut self.inner).expect( + "TierStore should not be shared during configuration. No other references should exist", + ); + + inner.backup_store = Some(backup); + } + + /// Configures the ephemeral store for non-critical, rebuildable data. + /// + /// When configured, selected cache-like data is routed to this store instead of + /// the primary store. + pub fn set_ephemeral_store(&mut self, ephemeral: Arc) { + debug_assert_eq!(Arc::strong_count(&self.inner), 1); + + let inner = Arc::get_mut(&mut self.inner).expect( + "TierStore should not be shared during configuration. No other references should exist", + ); + + inner.ephemeral_store = Some(ephemeral); + } + + pub(crate) fn set_index_store(&mut self, index: TierStoreIndex) { + debug_assert_eq!(Arc::strong_count(&self.inner), 1); + + let inner = Arc::get_mut(&mut self.inner).expect( + "TierStore should not be shared during configuration. No other references should exist", + ); + + inner.index = Some(index); + } +} + +pub(crate) async fn setup_index_store(data_dir: PathBuf) -> io::Result { + TierStoreIndex::new(data_dir).await +} + +impl KVStore for TierStore { + fn read( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, + ) -> impl Future, io::Error>> + 'static + Send { + let inner = Arc::clone(&self.inner); + + let primary_namespace = primary_namespace.to_string(); + let secondary_namespace = secondary_namespace.to_string(); + let key = key.to_string(); + + async move { inner.read_internal(primary_namespace, secondary_namespace, key).await } + } + + fn write( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, buf: Vec, + ) -> impl Future> + 'static + Send { + let inner = Arc::clone(&self.inner); + let locking_key = inner.build_locking_key(primary_namespace, secondary_namespace, key); + let (lock_ref, version) = inner.get_new_version_and_lock_ref(locking_key.clone()); + + let primary_namespace = primary_namespace.to_string(); + let secondary_namespace = secondary_namespace.to_string(); + let key = key.to_string(); + + async move { + inner + .write_internal( + primary_namespace, + secondary_namespace, + key, + buf, + lock_ref, + locking_key, + version, + ) + .await + } + } + + fn remove( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, lazy: bool, + ) -> impl Future> + 'static + Send { + let inner = Arc::clone(&self.inner); + let locking_key = inner.build_locking_key(primary_namespace, secondary_namespace, key); + let (lock_ref, version) = inner.get_new_version_and_lock_ref(locking_key.clone()); + + let primary_namespace = primary_namespace.to_string(); + let secondary_namespace = secondary_namespace.to_string(); + let key = key.to_string(); + + async move { + inner + .remove_internal( + primary_namespace, + secondary_namespace, + key, + lazy, + lock_ref, + locking_key, + version, + ) + .await + } + } + + fn list( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> impl Future, io::Error>> + 'static + Send { + let inner = Arc::clone(&self.inner); + + let primary_namespace = primary_namespace.to_string(); + let secondary_namespace = secondary_namespace.to_string(); + + async move { inner.list_internal(primary_namespace, secondary_namespace).await } + } +} + +impl PaginatedKVStore for TierStore { + fn list_paginated( + &self, primary_namespace: &str, secondary_namespace: &str, page_token: Option, + ) -> impl Future> + 'static + Send { + let inner = Arc::clone(&self.inner); + + let primary_namespace = primary_namespace.to_string(); + let secondary_namespace = secondary_namespace.to_string(); + + async move { + inner.list_paginated_internal(primary_namespace, secondary_namespace, page_token).await + } + } +} + +struct TierStoreInner { + /// The authoritative store for durable data. + primary_store: Arc, + /// The store used for non-critical, rebuildable cached data. + ephemeral_store: Option>, + /// An optional second durable store for primary-backed data. + backup_store: Option>, + /// The local store used to index the logical contents across tiers. + index: Option, + /// Per-namespace locks for serializing first-use index initialization. + index_initialization_locks: Mutex>>>, + /// Per-key locks for serializing primary+backup operations and skipping stale writes. + locks: Mutex>>>, + next_write_version: AtomicU64, + logger: Arc, +} + +impl TierStoreInner { + /// Creates a tier store with the primary data store. + pub fn new(primary_store: Arc, logger: Arc) -> Self { + Self { + primary_store, + ephemeral_store: None, + backup_store: None, + index: None, + index_initialization_locks: Mutex::new(HashMap::new()), + locks: Mutex::new(HashMap::new()), + next_write_version: AtomicU64::new(1), + logger, + } + } + + fn get_new_version_and_lock_ref(&self, locking_key: String) -> (Arc>, u64) { + let version = self.next_write_version.fetch_add(1, Ordering::Relaxed); + if version == u64::MAX { + panic!("TierStore version counter overflowed"); + } + + (self.get_lock_ref(locking_key), version) + } + + /// Returns the lock that serializes operations for one logical key. + fn get_lock_ref(&self, locking_key: String) -> Arc> { + let mut locks = self.locks.lock().expect("lock"); + Arc::clone(locks.entry(locking_key).or_insert_with(|| Arc::new(TokioMutex::new(0)))) + } + + fn clean_locks(&self, lock_ref: &Arc>, locking_key: String) { + let mut locks = self.locks.lock().expect("lock"); + let strong_count = Arc::strong_count(lock_ref); + debug_assert!(strong_count >= 2, "Unexpected TierStore lock strong count"); + if strong_count == 2 { + locks.remove(&locking_key); + } + } + + /// Returns the lock that serializes initialization of the given logical namespace. + fn get_index_initialization_lock( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> Arc> { + let namespace_id = TierStoreIndex::namespace_id(primary_namespace, secondary_namespace); + let mut locks = self.index_initialization_locks.lock().expect("lock"); + Arc::clone(locks.entry(namespace_id).or_insert_with(|| Arc::new(TokioMutex::new(())))) + } + + /// Removes an initialization lock after its final active user releases it. + fn clean_index_initialization_locks( + &self, lock_ref: &Arc>, primary_namespace: &str, secondary_namespace: &str, + ) { + let namespace_id = TierStoreIndex::namespace_id(primary_namespace, secondary_namespace); + let mut locks = self.index_initialization_locks.lock().expect("lock"); + if Arc::strong_count(lock_ref) == 2 { + locks.remove(&namespace_id); + } + } + + fn build_locking_key( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, + ) -> String { + if primary_namespace.is_empty() { + key.to_owned() + } else { + format!("{}#{}#{}", primary_namespace, secondary_namespace, key) + } + } + + /// Reads from the primary data store. + async fn read_primary( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, + ) -> io::Result> { + match KVStore::read( + self.primary_store.as_ref(), + primary_namespace, + secondary_namespace, + key, + ) + .await + { + Ok(data) => Ok(data), + Err(e) => Err(e), + } + } + + /// Lists keys from the primary data store. + async fn list_primary( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> io::Result> { + match KVStore::list(self.primary_store.as_ref(), primary_namespace, secondary_namespace) + .await + { + Ok(keys) => Ok(keys), + Err(e) => { + log_error!( + self.logger, + "Failed to list from primary store for namespace {}/{}: {}.", + primary_namespace, + secondary_namespace, + e + ); + Err(e) + }, + } + } + + async fn write_primary_backup_async( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, buf: Vec, + ) -> io::Result<()> { + if let Some(backup_store) = self.backup_store.as_ref() { + let primary_fut = KVStore::write( + self.primary_store.as_ref(), + primary_namespace, + secondary_namespace, + key, + buf.clone(), + ); + + let backup_fut = KVStore::write( + backup_store.as_ref(), + primary_namespace, + secondary_namespace, + key, + buf, + ); + + let (primary_res, backup_res) = tokio::join!(primary_fut, backup_fut); + + self.handle_primary_backup_results( + "write", + primary_namespace, + secondary_namespace, + key, + primary_res, + backup_res, + ) + } else { + KVStore::write( + self.primary_store.as_ref(), + primary_namespace, + secondary_namespace, + key, + buf, + ) + .await + } + } + + async fn remove_primary_backup_async( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, lazy: bool, + ) -> io::Result<()> { + let primary_fut = KVStore::remove( + self.primary_store.as_ref(), + primary_namespace, + secondary_namespace, + key, + lazy, + ); + + if let Some(backup_store) = self.backup_store.as_ref() { + let backup_fut = KVStore::remove( + backup_store.as_ref(), + primary_namespace, + secondary_namespace, + key, + lazy, + ); + + let (primary_res, backup_res) = tokio::join!(primary_fut, backup_fut); + + self.handle_primary_backup_results( + "removal", + primary_namespace, + secondary_namespace, + key, + primary_res, + backup_res, + ) + } else { + primary_fut.await + } + } + + async fn execute_locked_write( + &self, lock_ref: Arc>, locking_key: String, version: u64, callback: F, + ) -> io::Result<()> + where + F: FnOnce() -> Fut, + Fut: Future>, + { + let res = { + let mut last_written_version = lock_ref.lock().await; + + if version <= *last_written_version { + Ok(()) + } else { + let res = callback().await; + // A failed multi-store operation may still have updated one of its stores. We record + // the attempted version regardless so an older operation cannot overwrite newer state. + *last_written_version = version; + res + } + }; + + self.clean_locks(&lock_ref, locking_key); + res + } + + async fn read_internal( + &self, primary_namespace: String, secondary_namespace: String, key: String, + ) -> io::Result> { + check_namespace_key_validity( + primary_namespace.as_str(), + secondary_namespace.as_str(), + Some(key.as_str()), + "read", + )?; + self.prepare_namespace(&primary_namespace, &secondary_namespace).await?; + + if is_ephemeral_cached_key(&primary_namespace, &secondary_namespace, &key) { + if let Some(eph_store) = self.ephemeral_store.as_ref() { + // We don't retry ephemeral-store reads here. Local failures are treated as + // terminal for this access path rather than falling back to another store. + return KVStore::read( + eph_store.as_ref(), + &primary_namespace, + &secondary_namespace, + &key, + ) + .await; + } + } + + self.read_primary(&primary_namespace, &secondary_namespace, &key).await + } + + async fn write_internal( + &self, primary_namespace: String, secondary_namespace: String, key: String, buf: Vec, + lock_ref: Arc>, locking_key: String, version: u64, + ) -> io::Result<()> { + check_namespace_key_validity( + primary_namespace.as_str(), + secondary_namespace.as_str(), + Some(key.as_str()), + "write", + )?; + self.prepare_namespace(&primary_namespace, &secondary_namespace).await?; + + self.execute_locked_write(lock_ref, locking_key, version, || async move { + self.write_locked(&primary_namespace, &secondary_namespace, &key, buf).await + }) + .await + } + + /// Writes one key after recovering any pending operation while its per-key lock is held. + async fn write_locked( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, value: Vec, + ) -> io::Result<()> { + self.recover_key_locked(primary_namespace, secondary_namespace, key).await?; + let tier = self.value_tier(primary_namespace, secondary_namespace, key); + let Some(index) = self.index.as_ref() else { + return self + .write_value(tier, primary_namespace, secondary_namespace, key, value) + .await; + }; + + if index.contains_entry(primary_namespace, secondary_namespace, key).await? { + return self + .write_value(tier, primary_namespace, secondary_namespace, key, value) + .await; + } + if self.value_exists(tier, primary_namespace, secondary_namespace, key).await? { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "Value exists without a tier-store index entry", + )); + } + + let journal_entry = JournalEntry { + primary_namespace: primary_namespace.to_string(), + secondary_namespace: secondary_namespace.to_string(), + key: key.to_string(), + tier, + requires_backup: tier == ValueTier::Primary && self.backup_store.is_some(), + operation: JournalOperation::Create { value }, + }; + index.write_journal_entry(&journal_entry).await?; + self.apply_pending_create(&journal_entry).await?; + index.remove_journal_entry(primary_namespace, secondary_namespace, key).await + } + + async fn remove_internal( + &self, primary_namespace: String, secondary_namespace: String, key: String, lazy: bool, + lock_ref: Arc>, locking_key: String, version: u64, + ) -> io::Result<()> { + check_namespace_key_validity( + primary_namespace.as_str(), + secondary_namespace.as_str(), + Some(key.as_str()), + "remove", + )?; + self.prepare_namespace(&primary_namespace, &secondary_namespace).await?; + + self.execute_locked_write(lock_ref, locking_key, version, || async move { + self.remove_locked(&primary_namespace, &secondary_namespace, &key, lazy).await + }) + .await + } + + /// Removes one key after recovering any pending operation while its per-key lock is held. + async fn remove_locked( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, lazy: bool, + ) -> io::Result<()> { + self.recover_key_locked(primary_namespace, secondary_namespace, key).await?; + let tier = self.value_tier(primary_namespace, secondary_namespace, key); + let Some(index) = self.index.as_ref() else { + return self + .remove_value(tier, primary_namespace, secondary_namespace, key, lazy) + .await; + }; + + let journal_entry = JournalEntry { + primary_namespace: primary_namespace.to_string(), + secondary_namespace: secondary_namespace.to_string(), + key: key.to_string(), + tier, + requires_backup: tier == ValueTier::Primary && self.backup_store.is_some(), + operation: JournalOperation::Remove { lazy }, + }; + index.write_journal_entry(&journal_entry).await?; + self.apply_pending_remove(&journal_entry).await?; + index.remove_journal_entry(primary_namespace, secondary_namespace, key).await + } + + /// Selects the authoritative value tier for a logical key. + fn value_tier( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, + ) -> ValueTier { + if self.ephemeral_store.is_some() + && is_ephemeral_cached_key(primary_namespace, secondary_namespace, key) + { + ValueTier::Ephemeral + } else { + ValueTier::Primary + } + } + + /// Checks for an existing value before classifying an unindexed key as a new creation. + async fn value_exists( + &self, tier: ValueTier, primary_namespace: &str, secondary_namespace: &str, key: &str, + ) -> io::Result { + let store = match tier { + ValueTier::Primary => &self.primary_store, + ValueTier::Ephemeral => self.ephemeral_store.as_ref().ok_or_else(|| { + io::Error::new(io::ErrorKind::NotFound, "Ephemeral store is unavailable") + })?, + }; + match KVStore::read(store.as_ref(), primary_namespace, secondary_namespace, key).await { + Ok(_) => Ok(true), + Err(e) if e.kind() == io::ErrorKind::NotFound => { + if tier == ValueTier::Primary { + if let Some(backup_store) = self.backup_store.as_ref() { + return match KVStore::read( + backup_store.as_ref(), + primary_namespace, + secondary_namespace, + key, + ) + .await + { + Ok(_) => Err(io::Error::new( + io::ErrorKind::InvalidData, + "Backup value exists without a primary or index entry", + )), + Err(e) if e.kind() == io::ErrorKind::NotFound => Ok(false), + Err(e) => Err(e), + }; + } + } + Ok(false) + }, + Err(e) => Err(e), + } + } + + /// Writes a value to its authoritative tier and any required backup. + async fn write_value( + &self, tier: ValueTier, primary_namespace: &str, secondary_namespace: &str, key: &str, + value: Vec, + ) -> io::Result<()> { + match tier { + ValueTier::Primary => { + self.write_primary_backup_async(primary_namespace, secondary_namespace, key, value) + .await + }, + ValueTier::Ephemeral => { + let store = self.ephemeral_store.as_ref().ok_or_else(|| { + io::Error::new(io::ErrorKind::NotFound, "Ephemeral store is unavailable") + })?; + KVStore::write(store.as_ref(), primary_namespace, secondary_namespace, key, value) + .await + }, + } + } + + /// Removes a value from its authoritative tier and any required backup. + async fn remove_value( + &self, tier: ValueTier, primary_namespace: &str, secondary_namespace: &str, key: &str, + lazy: bool, + ) -> io::Result<()> { + match tier { + ValueTier::Primary => { + self.remove_primary_backup_async(primary_namespace, secondary_namespace, key, lazy) + .await + }, + ValueTier::Ephemeral => { + let store = self.ephemeral_store.as_ref().ok_or_else(|| { + io::Error::new(io::ErrorKind::NotFound, "Ephemeral store is unavailable") + })?; + KVStore::remove(store.as_ref(), primary_namespace, secondary_namespace, key, lazy) + .await + }, + } + } + + /// Completes a journaled create and makes it visible in the listing index. + async fn apply_pending_create(&self, entry: &JournalEntry) -> io::Result<()> { + let JournalOperation::Create { value } = &entry.operation else { + return Err(io::Error::new(io::ErrorKind::InvalidData, "Expected pending create")); + }; + if entry.requires_backup && self.backup_store.is_none() { + return Err(io::Error::new( + io::ErrorKind::NotFound, + "Required backup store is unavailable", + )); + } + self.write_value( + entry.tier, + &entry.primary_namespace, + &entry.secondary_namespace, + &entry.key, + value.clone(), + ) + .await?; + self.index + .as_ref() + .expect("pending operations require an index") + .write_entry(&entry.primary_namespace, &entry.secondary_namespace, &entry.key) + .await + } + + /// Completes a journaled removal, hiding the key before deleting its value copies. + async fn apply_pending_remove(&self, entry: &JournalEntry) -> io::Result<()> { + let JournalOperation::Remove { lazy } = &entry.operation else { + return Err(io::Error::new(io::ErrorKind::InvalidData, "Expected pending removal")); + }; + if entry.requires_backup && self.backup_store.is_none() { + return Err(io::Error::new( + io::ErrorKind::NotFound, + "Required backup store is unavailable", + )); + } + self.index + .as_ref() + .expect("pending operations require an index") + .remove_entry(&entry.primary_namespace, &entry.secondary_namespace, &entry.key, *lazy) + .await?; + self.remove_value( + entry.tier, + &entry.primary_namespace, + &entry.secondary_namespace, + &entry.key, + *lazy, + ) + .await + } + + async fn list_internal( + &self, primary_namespace: String, secondary_namespace: String, + ) -> io::Result> { + check_namespace_key_validity( + primary_namespace.as_str(), + secondary_namespace.as_str(), + None, + "list", + )?; + + self.prepare_namespace(&primary_namespace, &secondary_namespace).await?; + if let Some(index) = self.index.as_ref() { + return index.list(&primary_namespace, &secondary_namespace).await; + } + + self.list_value_stores(&primary_namespace, &secondary_namespace).await + } + + /// Imports a namespace's existing primary-store keys and makes its index authoritative. + /// + /// Primary pagination returns keys from newest to oldest, so the complete result is reversed + /// before insertion to preserve that order in the local index. Values already present in the + /// ephemeral store are discarded because their positions relative to primary keys cannot be + /// reconstructed without the index. Cache values still in the primary store are retained and + /// subsequently moved to ephemeral storage without changing their imported index positions. + async fn ensure_namespace_indexed( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> io::Result<()> { + let Some(index) = self.index.as_ref() else { + return Ok(()); + }; + if index.is_namespace_ready(primary_namespace, secondary_namespace).await? { + return Ok(()); + } + + let lock_ref = self.get_index_initialization_lock(primary_namespace, secondary_namespace); + let result: io::Result<()> = async { + let _guard = lock_ref.lock().await; + if index.is_namespace_ready(primary_namespace, secondary_namespace).await? { + Ok(()) + } else { + self.discard_unindexed_ephemeral_cache(primary_namespace, secondary_namespace) + .await?; + let mut keys = Vec::new(); + let mut page_token = None; + loop { + let page = PaginatedKVStore::list_paginated( + self.primary_store.as_ref(), + primary_namespace, + secondary_namespace, + page_token, + ) + .await?; + keys.extend(page.keys); + match page.next_page_token { + Some(next_page_token) => page_token = Some(next_page_token), + None => break, + } + } + + for key in keys.into_iter().rev() { + index.write_entry(primary_namespace, secondary_namespace, &key).await?; + } + index.mark_namespace_ready(primary_namespace, secondary_namespace).await + } + } + .await; + self.clean_index_initialization_locks(&lock_ref, primary_namespace, secondary_namespace); + result + } + + /// Initializes a namespace, recovers pending membership changes, and reconciles cache placement. + async fn prepare_namespace( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> io::Result<()> { + self.ensure_namespace_indexed(primary_namespace, secondary_namespace).await?; + self.recover_namespace(primary_namespace, secondary_namespace).await?; + self.ensure_ephemeral_cache_reconciled(primary_namespace, secondary_namespace).await + } + + /// Discards cache values whose ordering cannot be recovered from a missing index. + async fn discard_unindexed_ephemeral_cache( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> io::Result<()> { + let Some(ephemeral_store) = self.ephemeral_store.as_ref() else { + return Ok(()); + }; + for key in ephemeral_cache_keys(primary_namespace) { + KVStore::remove( + ephemeral_store.as_ref(), + primary_namespace, + secondary_namespace, + key, + false, + ) + .await?; + } + Ok(()) + } + + /// Reconciles cache placement once and records completion for subsequent namespace accesses. + async fn ensure_ephemeral_cache_reconciled( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> io::Result<()> { + let (Some(index), Some(_)) = (self.index.as_ref(), self.ephemeral_store.as_ref()) else { + return Ok(()); + }; + if ephemeral_cache_keys(primary_namespace).next().is_none() + || index.is_cache_ready(primary_namespace, secondary_namespace).await? + { + return Ok(()); + } + + let lock_ref = self.get_index_initialization_lock(primary_namespace, secondary_namespace); + let result: io::Result<()> = async { + let _guard = lock_ref.lock().await; + if index.is_cache_ready(primary_namespace, secondary_namespace).await? { + Ok(()) + } else { + self.reconcile_ephemeral_cache(primary_namespace, secondary_namespace).await?; + index.mark_cache_ready(primary_namespace, secondary_namespace).await + } + } + .await; + self.clean_index_initialization_locks(&lock_ref, primary_namespace, secondary_namespace); + result + } + + /// Moves indexed cache values to ephemeral storage without changing their index positions. + /// + /// Destination writes precede source removals. Repeating this after interruption either copies + /// the primary value again or removes a stale primary/backup copy after finding the destination. + async fn reconcile_ephemeral_cache( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> io::Result<()> { + let (Some(index), Some(ephemeral_store)) = + (self.index.as_ref(), self.ephemeral_store.as_ref()) + else { + return Ok(()); + }; + + for key in ephemeral_cache_keys(primary_namespace) { + let locking_key = self.build_locking_key(primary_namespace, secondary_namespace, key); + let lock_ref = self.get_lock_ref(locking_key.clone()); + let result: io::Result<()> = async { + let _guard = lock_ref.lock().await; + self.recover_key_locked(primary_namespace, secondary_namespace, key).await?; + if !index.contains_entry(primary_namespace, secondary_namespace, key).await? { + return Ok(()); + } + + match KVStore::read( + ephemeral_store.as_ref(), + primary_namespace, + secondary_namespace, + key, + ) + .await + { + Ok(_) => {}, + Err(e) if e.kind() == io::ErrorKind::NotFound => { + match self.read_primary(primary_namespace, secondary_namespace, key).await { + Ok(value) => { + KVStore::write( + ephemeral_store.as_ref(), + primary_namespace, + secondary_namespace, + key, + value, + ) + .await?; + }, + Err(e) if e.kind() == io::ErrorKind::NotFound => { + self.remove_primary_backup_async( + primary_namespace, + secondary_namespace, + key, + false, + ) + .await?; + return index + .remove_entry( + primary_namespace, + secondary_namespace, + key, + false, + ) + .await; + }, + Err(e) => return Err(e), + } + }, + Err(e) => return Err(e), + } + + self.remove_primary_backup_async(primary_namespace, secondary_namespace, key, false) + .await + } + .await; + self.clean_locks(&lock_ref, locking_key); + result?; + } + Ok(()) + } + + /// Completes journaled creates and removals before exposing a namespace. + async fn recover_namespace( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> io::Result<()> { + let Some(index) = self.index.as_ref() else { + return Ok(()); + }; + for key in index.list_journal_entries(primary_namespace, secondary_namespace).await? { + let locking_key = self.build_locking_key(primary_namespace, secondary_namespace, &key); + let lock_ref = self.get_lock_ref(locking_key.clone()); + let result: io::Result<()> = async { + let _guard = lock_ref.lock().await; + self.recover_key_locked(primary_namespace, secondary_namespace, &key).await + } + .await; + self.clean_locks(&lock_ref, locking_key); + result?; + } + Ok(()) + } + + /// Recovers one key while its per-key operation lock is held by the caller. + async fn recover_key_locked( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, + ) -> io::Result<()> { + let Some(index) = self.index.as_ref() else { + return Ok(()); + }; + let entry = + match index.read_journal_entry(primary_namespace, secondary_namespace, key).await { + Ok(entry) => entry, + Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(()), + Err(e) => return Err(e), + }; + match &entry.operation { + JournalOperation::Create { .. } => self.apply_pending_create(&entry).await?, + JournalOperation::Remove { .. } => self.apply_pending_remove(&entry).await?, + } + index.remove_journal_entry(primary_namespace, secondary_namespace, key).await + } + + /// Lists the authoritative logical keys directly from the primary and ephemeral value stores. + /// + /// Ephemeral-routed keys are taken only from the ephemeral store, excluding stale primary + /// copies. This path is used when no local index store is configured. + async fn list_value_stores( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> io::Result> { + let mut keys = self.list_primary(primary_namespace, secondary_namespace).await?; + + let Some(ephemeral_store) = self.ephemeral_store.as_ref() else { + return Ok(keys); + }; + + if primary_namespace != NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE + && primary_namespace != SCORER_PERSISTENCE_PRIMARY_NAMESPACE + { + return Ok(keys); + } + + // The ephemeral store is authoritative for keys routed there. Exclude stale + // primary copies, then add only routed keys from the ephemeral store. + keys.retain(|key| !is_ephemeral_cached_key(primary_namespace, secondary_namespace, key)); + + let ephemeral_keys = + KVStore::list(ephemeral_store.as_ref(), primary_namespace, secondary_namespace).await?; + + for key in ephemeral_keys { + if is_ephemeral_cached_key(primary_namespace, secondary_namespace, &key) + && !keys.contains(&key) + { + keys.push(key); + } + } + + Ok(keys) + } + + async fn list_paginated_internal( + &self, primary_namespace: String, secondary_namespace: String, + page_token: Option, + ) -> io::Result { + check_namespace_key_validity( + primary_namespace.as_str(), + secondary_namespace.as_str(), + None, + "list_paginated", + )?; + + self.prepare_namespace(&primary_namespace, &secondary_namespace).await?; + if let Some(index) = self.index.as_ref() { + return index + .list_paginated(&primary_namespace, &secondary_namespace, page_token) + .await; + } + + PaginatedKVStore::list_paginated( + self.primary_store.as_ref(), + &primary_namespace, + &secondary_namespace, + page_token, + ) + .await + } + + fn handle_primary_backup_results( + &self, op: &str, primary_namespace: &str, secondary_namespace: &str, key: &str, + primary_res: io::Result<()>, backup_res: io::Result<()>, + ) -> io::Result<()> { + match (primary_res, backup_res) { + (Ok(()), Ok(())) => Ok(()), + (Err(primary_err), Ok(())) => { + log_error!( + self.logger, + "Primary {} failed after backup {} succeeded for key {}/{}/{}; primary and backup may have diverged: {}", + op, + op, + primary_namespace, + secondary_namespace, + key, + primary_err + ); + Err(primary_err) + }, + (Ok(()), Err(backup_err)) => { + log_error!( + self.logger, + "Backup {} failed after primary {} succeeded for key {}/{}/{}; primary and backup may have diverged: {}", + op, + op, + primary_namespace, + secondary_namespace, + key, + backup_err + ); + Err(backup_err) + }, + (Err(primary_err), Err(backup_err)) => { + log_error!( + self.logger, + "Primary and backup {}s both failed for key {}/{}/{}: primary={}, backup={}", + op, + primary_namespace, + secondary_namespace, + key, + primary_err, + backup_err + ); + Err(primary_err) + }, + } + } +} + +fn is_ephemeral_cached_key(pn: &str, _sn: &str, key: &str) -> bool { + ephemeral_cache_keys(pn).any(|cache_key| cache_key == key) +} + +fn ephemeral_cache_keys(primary_namespace: &str) -> impl Iterator + '_ { + [ + (NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, NETWORK_GRAPH_PERSISTENCE_KEY), + (SCORER_PERSISTENCE_PRIMARY_NAMESPACE, SCORER_PERSISTENCE_KEY), + (SCORER_PERSISTENCE_PRIMARY_NAMESPACE, EXTERNAL_PATHFINDING_SCORES_CACHE_KEY), + ] + .into_iter() + .filter_map(move |(cache_namespace, cache_key)| { + (primary_namespace == cache_namespace).then_some(cache_key) + }) +} + +#[cfg(test)] +mod tests { + use std::future::Future; + use std::panic::RefUnwindSafe; + use std::path::PathBuf; + use std::sync::atomic::{AtomicUsize, Ordering}; + use std::sync::{Arc, Mutex}; + + use lightning::util::logger::Level; + use lightning::util::persist::{ + CHANNEL_MANAGER_PERSISTENCE_KEY, CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + CHANNEL_MONITOR_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MONITOR_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, SCORER_PERSISTENCE_SECONDARY_NAMESPACE, + }; + use lightning_persister::fs_store::v2::FilesystemStoreV2; + use tokio::sync::oneshot; + + use super::*; + use crate::io::test_utils::{ + do_read_write_remove_list_persist, random_storage_path, InMemoryStore, + }; + use crate::io::tier_store::TierStore; + use crate::logger::Logger; + use crate::types::{DynStore, DynStoreWrapper}; + + impl RefUnwindSafe for TierStore {} + + struct CleanupDir(PathBuf); + impl Drop for CleanupDir { + fn drop(&mut self) { + let _ = std::fs::remove_dir_all(&self.0); + } + } + + fn setup_tier_store(primary_store: Arc, logger: Arc) -> TierStore { + TierStore::new(primary_store, logger) + } + + fn set_test_index_store(tier: &mut TierStore) { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + tier.set_index_store(TierStoreIndex::from_store(store)); + } + + #[test] + fn journal_entry_roundtrips() { + for operation in [ + JournalOperation::Create { value: vec![0, 1, 2, 255] }, + JournalOperation::Remove { lazy: true }, + ] { + let entry = JournalEntry { + primary_namespace: "primary".to_string(), + secondary_namespace: "secondary".to_string(), + key: "key".to_string(), + tier: ValueTier::Primary, + requires_backup: true, + operation, + }; + assert_eq!(JournalEntry::deserialize(&entry.serialize()).unwrap(), entry); + } + } + + #[test] + fn page_token_roundtrips_and_validates_context() { + let database_id = [2; INDEX_DATABASE_ID_LEN]; + let namespace_id = TierStoreIndex::namespace_id("primary", "secondary"); + let token = TierStorePageToken::encode( + &database_id, + namespace_id.clone(), + PageToken::new("opaque:index-token".to_string()), + ); + + let decoded = + TierStorePageToken::decode(token.clone(), &database_id, &namespace_id).unwrap(); + assert_eq!(decoded.as_str(), "opaque:index-token"); + assert_eq!( + TierStorePageToken::decode(token.clone(), &[3; INDEX_DATABASE_ID_LEN], &namespace_id) + .unwrap_err() + .kind(), + io::ErrorKind::InvalidInput + ); + assert_eq!( + TierStorePageToken::decode(token, &database_id, "another-namespace") + .unwrap_err() + .kind(), + io::ErrorKind::InvalidInput + ); + assert_eq!( + TierStorePageToken::decode( + PageToken::new("not-a-tier-store-token".to_string()), + &database_id, + &namespace_id, + ) + .unwrap_err() + .kind(), + io::ErrorKind::InvalidInput + ); + let unsupported_version = TierStorePageToken { + format_version: PAGE_TOKEN_FORMAT_VERSION + 1, + index_database_id: database_id.to_vec(), + namespace_id: namespace_id.clone(), + index_page_token: "opaque:index-token".to_string(), + }; + let unsupported_version = + PageToken::new(Writeable::encode(&unsupported_version).to_lower_hex_string()); + assert_eq!( + TierStorePageToken::decode(unsupported_version, &database_id, &namespace_id) + .unwrap_err() + .kind(), + io::ErrorKind::InvalidInput + ); + } + + #[tokio::test] + async fn index_store_is_internal_persistent_sqlite_store() { + let base_dir = random_storage_path(); + let _cleanup = CleanupDir(base_dir.clone()); + + let index = setup_index_store(base_dir.clone()).await.unwrap(); + assert!(base_dir.join(SQLITE_TIER_INDEX_DB_FILE_NAME).exists()); + + let database_id = + TierStoreIndex::read_or_create_database_id(index.store.as_ref()).await.unwrap(); + let persisted_database_id = + TierStoreIndex::read_or_create_database_id(index.store.as_ref()).await.unwrap(); + assert_ne!(database_id, [0; INDEX_DATABASE_ID_LEN]); + assert_eq!(persisted_database_id, database_id); + } + + #[tokio::test] + async fn index_store_rejects_second_owner() { + let base_dir = random_storage_path(); + let _cleanup = CleanupDir(base_dir.clone()); + + let _index = setup_index_store(base_dir.clone()).await.unwrap(); + let error = match setup_index_store(base_dir).await { + Ok(_) => panic!("a second index-store owner must be rejected"), + Err(e) => e, + }; + assert_eq!(error.kind(), io::ErrorKind::AlreadyExists); + } + + #[tokio::test] + async fn indexed_listing_orders_keys_across_primary_and_ephemeral_stores() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + let _cleanup = CleanupDir(base_dir); + + let primary_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let ephemeral_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let mut tier = setup_tier_store(primary_store, logger); + set_test_index_store(&mut tier); + tier.set_ephemeral_store(ephemeral_store); + + for key in ["primary-a", NETWORK_GRAPH_PERSISTENCE_KEY, "primary-b"] { + tier.write( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + key, + vec![1], + ) + .await + .unwrap(); + } + + let page = PaginatedKVStore::list_paginated( + &tier, + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + None, + ) + .await + .unwrap(); + assert_eq!( + page.keys, + vec![ + "primary-b".to_string(), + NETWORK_GRAPH_PERSISTENCE_KEY.to_string(), + "primary-a".to_string(), + ] + ); + + let mut listed = KVStore::list( + &tier, + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + ) + .await + .unwrap(); + listed.sort(); + assert_eq!( + listed, + vec![ + NETWORK_GRAPH_PERSISTENCE_KEY.to_string(), + "primary-a".to_string(), + "primary-b".to_string(), + ] + ); + } + + #[tokio::test] + async fn indexed_listing_preserves_updates_and_reorders_recreated_keys() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + let _cleanup = CleanupDir(base_dir); + + let primary_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let mut tier = setup_tier_store(primary_store, logger); + set_test_index_store(&mut tier); + + for key in ["a", "b"] { + tier.write("namespace", "", key, vec![1]).await.unwrap(); + } + tier.write("namespace", "", "a", vec![2]).await.unwrap(); + + let page = PaginatedKVStore::list_paginated(&tier, "namespace", "", None).await.unwrap(); + assert_eq!(page.keys, vec!["b".to_string(), "a".to_string()]); + + tier.remove("namespace", "", "a", false).await.unwrap(); + tier.write("namespace", "", "a", vec![3]).await.unwrap(); + let page = PaginatedKVStore::list_paginated(&tier, "namespace", "", None).await.unwrap(); + assert_eq!(page.keys, vec!["a".to_string(), "b".to_string()]); + } + + #[tokio::test] + async fn adopting_ephemeral_storage_migrates_cache_keys_without_reordering_them() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + let _cleanup = CleanupDir(base_dir); + + let primary_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let backup_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let ephemeral_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + for key in ["old", NETWORK_GRAPH_PERSISTENCE_KEY, "new"] { + primary_store + .write( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + key, + key.as_bytes().to_vec(), + ) + .await + .unwrap(); + backup_store + .write( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + key, + key.as_bytes().to_vec(), + ) + .await + .unwrap(); + } + let mut tier = setup_tier_store(Arc::clone(&primary_store), logger); + set_test_index_store(&mut tier); + tier.set_backup_store(Arc::clone(&backup_store)); + tier.set_ephemeral_store(Arc::clone(&ephemeral_store)); + + let page = PaginatedKVStore::list_paginated( + &tier, + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + None, + ) + .await + .unwrap(); + assert_eq!( + page.keys, + vec!["new".to_string(), NETWORK_GRAPH_PERSISTENCE_KEY.to_string(), "old".to_string(),] + ); + assert_eq!( + ephemeral_store + .read( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + ) + .await + .unwrap(), + NETWORK_GRAPH_PERSISTENCE_KEY.as_bytes() + ); + assert!(primary_store + .read( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + ) + .await + .is_err()); + assert!(backup_store + .read( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + ) + .await + .is_err()); + } + + #[tokio::test] + async fn missing_index_discards_ephemeral_only_cache_data_before_reading() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + let _cleanup = CleanupDir(base_dir); + + let primary_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let ephemeral_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + ephemeral_store + .write( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + vec![1], + ) + .await + .unwrap(); + let mut tier = setup_tier_store(primary_store, logger); + set_test_index_store(&mut tier); + tier.set_ephemeral_store(Arc::clone(&ephemeral_store)); + + assert!(tier + .read( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + ) + .await + .is_err()); + assert!(ephemeral_store + .read( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + ) + .await + .is_err()); + assert!(KVStore::list( + &tier, + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + ) + .await + .unwrap() + .is_empty()); + } + + #[tokio::test] + async fn missing_index_rebuilds_cache_from_primary_instead_of_ephemeral() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + let _cleanup = CleanupDir(base_dir); + + let primary_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let ephemeral_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + primary_store + .write( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + vec![2], + ) + .await + .unwrap(); + ephemeral_store + .write( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + vec![1], + ) + .await + .unwrap(); + let mut tier = setup_tier_store(Arc::clone(&primary_store), logger); + set_test_index_store(&mut tier); + tier.set_ephemeral_store(Arc::clone(&ephemeral_store)); + + assert_eq!( + tier.read( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + ) + .await + .unwrap(), + vec![2] + ); + assert!(primary_store + .read( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + ) + .await + .is_err()); + assert!(tier + .inner + .index + .as_ref() + .unwrap() + .contains_entry( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + ) + .await + .unwrap()); + } + + #[tokio::test] + async fn namespace_preparation_finishes_interrupted_cache_migration() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + let _cleanup = CleanupDir(base_dir); + + let primary_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let ephemeral_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let mut tier = setup_tier_store(Arc::clone(&primary_store), logger); + set_test_index_store(&mut tier); + tier.set_ephemeral_store(Arc::clone(&ephemeral_store)); + let index = tier.inner.index.as_ref().unwrap(); + for key in ["old", NETWORK_GRAPH_PERSISTENCE_KEY, "new"] { + index + .write_entry( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + key, + ) + .await + .unwrap(); + } + index + .mark_namespace_ready( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + ) + .await + .unwrap(); + primary_store + .write( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + vec![1], + ) + .await + .unwrap(); + ephemeral_store + .write( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + vec![1], + ) + .await + .unwrap(); + + let page = PaginatedKVStore::list_paginated( + &tier, + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + None, + ) + .await + .unwrap(); + assert_eq!( + page.keys, + vec!["new".to_string(), NETWORK_GRAPH_PERSISTENCE_KEY.to_string(), "old".to_string(),] + ); + assert!(primary_store + .read( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + ) + .await + .is_err()); + } + + #[tokio::test] + async fn namespace_initialization_preserves_existing_primary_order_across_pages() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + let _cleanup = CleanupDir(base_dir); + + let primary_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + for i in 0..55 { + primary_store + .write("namespace", "", &format!("existing-{i:02}"), vec![1]) + .await + .unwrap(); + } + let mut tier = setup_tier_store(primary_store, logger); + set_test_index_store(&mut tier); + + let mut actual = Vec::new(); + let mut page_token = None; + loop { + let page = + PaginatedKVStore::list_paginated(&tier, "namespace", "", page_token).await.unwrap(); + actual.extend(page.keys); + match page.next_page_token { + Some(next_page_token) => page_token = Some(next_page_token), + None => break, + } + } + + let expected = (0..55).rev().map(|i| format!("existing-{i:02}")).collect::>(); + assert_eq!(actual, expected); + assert_eq!(KVStore::list(&tier, "namespace", "").await.unwrap().len(), 55); + } + + #[tokio::test] + async fn paginated_listing_rejects_tokens_from_another_namespace_or_index() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + let _cleanup = CleanupDir(base_dir); + + let primary_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let mut tier = setup_tier_store(Arc::clone(&primary_store), Arc::clone(&logger)); + let index_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + tier.set_index_store(TierStoreIndex::from_store_with_database_id( + index_store, + [1; INDEX_DATABASE_ID_LEN], + )); + for i in 0..51 { + tier.write("namespace", "", &format!("key-{i:02}"), vec![1]).await.unwrap(); + } + let token = PaginatedKVStore::list_paginated(&tier, "namespace", "", None) + .await + .unwrap() + .next_page_token + .unwrap(); + + let namespace_error = + PaginatedKVStore::list_paginated(&tier, "other-namespace", "", Some(token.clone())) + .await + .unwrap_err(); + assert_eq!(namespace_error.kind(), io::ErrorKind::InvalidInput); + + let replacement_index_store: Arc = + Arc::new(DynStoreWrapper(InMemoryStore::new())); + tier.set_index_store(TierStoreIndex::from_store_with_database_id( + replacement_index_store, + [2; INDEX_DATABASE_ID_LEN], + )); + let index_error = PaginatedKVStore::list_paginated(&tier, "namespace", "", Some(token)) + .await + .unwrap_err(); + assert_eq!(index_error.kind(), io::ErrorKind::InvalidInput); + } + + #[tokio::test] + async fn pending_creates_roll_forward_to_primary_and_backup_in_journal_order() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + let _cleanup = CleanupDir(base_dir); + + let primary_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let backup_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let index_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let mut tier = setup_tier_store(Arc::clone(&primary_store), Arc::clone(&logger)); + tier.set_backup_store(Arc::clone(&backup_store)); + tier.set_index_store(TierStoreIndex::from_store(Arc::clone(&index_store))); + tier.inner.ensure_namespace_indexed("namespace", "").await.unwrap(); + + for (key, value) in [("first", vec![1]), ("second", vec![2])] { + let entry = JournalEntry { + primary_namespace: "namespace".to_string(), + secondary_namespace: String::new(), + key: key.to_string(), + tier: ValueTier::Primary, + requires_backup: true, + operation: JournalOperation::Create { value: value.clone() }, + }; + tier.inner.index.as_ref().unwrap().write_journal_entry(&entry).await.unwrap(); + // Simulate the backup write winning the race before interruption. + backup_store.write("namespace", "", key, value).await.unwrap(); + } + + drop(tier); + let mut recovered_tier = setup_tier_store(Arc::clone(&primary_store), logger); + recovered_tier.set_backup_store(Arc::clone(&backup_store)); + recovered_tier.set_index_store(TierStoreIndex::from_store(index_store)); + + let page = + PaginatedKVStore::list_paginated(&recovered_tier, "namespace", "", None).await.unwrap(); + assert_eq!(page.keys, vec!["second".to_string(), "first".to_string()]); + for (key, value) in [("first", vec![1]), ("second", vec![2])] { + assert_eq!(primary_store.read("namespace", "", key).await.unwrap(), value); + assert_eq!(backup_store.read("namespace", "", key).await.unwrap(), value); + assert!(recovered_tier + .inner + .index + .as_ref() + .unwrap() + .read_journal_entry("namespace", "", key) + .await + .is_err()); + } + } + + #[tokio::test] + async fn pending_removal_hides_index_entry_before_removing_value_copies() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + let _cleanup = CleanupDir(base_dir); + + let primary_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let backup_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let mut tier = setup_tier_store(Arc::clone(&primary_store), logger); + tier.set_backup_store(Arc::clone(&backup_store)); + set_test_index_store(&mut tier); + tier.write("namespace", "", "key", vec![1]).await.unwrap(); + + let entry = JournalEntry { + primary_namespace: "namespace".to_string(), + secondary_namespace: String::new(), + key: "key".to_string(), + tier: ValueTier::Primary, + requires_backup: true, + operation: JournalOperation::Remove { lazy: false }, + }; + let index = tier.inner.index.as_ref().unwrap(); + index.write_journal_entry(&entry).await.unwrap(); + index.remove_entry("namespace", "", "key", false).await.unwrap(); + assert!(!index.contains_entry("namespace", "", "key").await.unwrap()); + assert!(primary_store.read("namespace", "", "key").await.is_ok()); + assert!(backup_store.read("namespace", "", "key").await.is_ok()); + assert!(index.read_journal_entry("namespace", "", "key").await.is_ok()); + + assert!(KVStore::list(&tier, "namespace", "").await.unwrap().is_empty()); + assert!(primary_store.read("namespace", "", "key").await.is_err()); + assert!(backup_store.read("namespace", "", "key").await.is_err()); + assert!(index.read_journal_entry("namespace", "", "key").await.is_err()); + } + + #[tokio::test] + async fn failed_create_retains_journal_without_exposing_index_entry() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + let _cleanup = CleanupDir(base_dir); + + let primary_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let mut tier = setup_tier_store(Arc::clone(&primary_store), logger); + let attempts = Arc::new(AtomicUsize::new(0)); + let backup_store: Arc = + Arc::new(DynStoreWrapper(InstrumentedStore::new(StoreBehavior::FailWrite { + attempts: Arc::clone(&attempts), + }))); + tier.set_backup_store(backup_store); + set_test_index_store(&mut tier); + + assert!(tier.write("namespace", "", "key", vec![1]).await.is_err()); + let index = tier.inner.index.as_ref().unwrap(); + assert!(primary_store.read("namespace", "", "key").await.is_ok()); + assert!(!index.contains_entry("namespace", "", "key").await.unwrap()); + assert!(matches!( + index.read_journal_entry("namespace", "", "key").await.unwrap().operation, + JournalOperation::Create { .. } + )); + assert_eq!(attempts.load(Ordering::Relaxed), 1); + } + + #[tokio::test] + async fn queued_write_recovers_pending_create_under_key_lock_before_classifying() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + let _cleanup = CleanupDir(base_dir); + + let primary_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let (snapshot_taken_tx, snapshot_taken_rx) = oneshot::channel(); + let (allow_return_tx, allow_return_rx) = oneshot::channel(); + let index_store: Arc = + Arc::new(DynStoreWrapper(InstrumentedStore::new(StoreBehavior::GateJournalList { + snapshot_taken: Mutex::new(Some(snapshot_taken_tx)), + allow_return: Mutex::new(Some(allow_return_rx)), + }))); + let mut tier = setup_tier_store(Arc::clone(&primary_store), logger); + tier.set_index_store(TierStoreIndex::from_store(index_store)); + tier.inner.ensure_namespace_indexed("namespace", "").await.unwrap(); + + let tier = Arc::new(tier); + let locking_key = tier.inner.build_locking_key("namespace", "", "key"); + let lock_ref = tier.inner.get_lock_ref(locking_key); + let guard = lock_ref.lock().await; + let write_task = { + let tier = Arc::clone(&tier); + tokio::spawn(async move { tier.write("namespace", "", "key", vec![2]).await }) + }; + snapshot_taken_rx.await.unwrap(); + + let pending = JournalEntry { + primary_namespace: "namespace".to_string(), + secondary_namespace: String::new(), + key: "key".to_string(), + tier: ValueTier::Primary, + requires_backup: false, + operation: JournalOperation::Create { value: vec![1] }, + }; + let index = tier.inner.index.as_ref().unwrap(); + index.write_journal_entry(&pending).await.unwrap(); + primary_store.write("namespace", "", "key", vec![1]).await.unwrap(); + allow_return_tx.send(()).unwrap(); + drop(guard); + drop(lock_ref); + + write_task.await.unwrap().unwrap(); + assert_eq!(primary_store.read("namespace", "", "key").await.unwrap(), vec![2]); + assert!(index.contains_entry("namespace", "", "key").await.unwrap()); + assert!(index.read_journal_entry("namespace", "", "key").await.is_err()); + } + + #[tokio::test] + async fn update_rejects_value_without_index_entry() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + let _cleanup = CleanupDir(base_dir); + + let primary_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let mut tier = setup_tier_store(Arc::clone(&primary_store), logger); + set_test_index_store(&mut tier); + tier.write("namespace", "", "key", vec![1]).await.unwrap(); + tier.inner + .index + .as_ref() + .unwrap() + .remove_entry("namespace", "", "key", false) + .await + .unwrap(); + + let error = tier.write("namespace", "", "key", vec![2]).await.unwrap_err(); + assert_eq!(error.kind(), io::ErrorKind::InvalidData); + assert_eq!(primary_store.read("namespace", "", "key").await.unwrap(), vec![1]); + } + + #[tokio::test] + async fn failed_removal_retains_journal_and_hides_index_entry() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + let _cleanup = CleanupDir(base_dir); + + let primary_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let backup_store: Arc = + Arc::new(DynStoreWrapper(InstrumentedStore::new(StoreBehavior::FailRemove))); + let mut tier = setup_tier_store(primary_store, logger); + tier.set_backup_store(backup_store); + set_test_index_store(&mut tier); + tier.write("namespace", "", "key", vec![1]).await.unwrap(); + + assert!(tier.remove("namespace", "", "key", false).await.is_err()); + let index = tier.inner.index.as_ref().unwrap(); + assert!(!index.contains_entry("namespace", "", "key").await.unwrap()); + assert!(matches!( + index.read_journal_entry("namespace", "", "key").await.unwrap().operation, + JournalOperation::Remove { .. } + )); + assert!(KVStore::list(&tier, "namespace", "").await.is_err()); + } + + enum StoreBehavior { + FailList, + FailWrite { + attempts: Arc, + }, + FailRemove, + GateJournalList { + snapshot_taken: Mutex>>, + allow_return: Mutex>>, + }, + } + + /// A store that injects selected failures or synchronization points while delegating other + /// operations to an inner [`InMemoryStore`]. + struct InstrumentedStore { + inner: InMemoryStore, + behavior: StoreBehavior, + } + + impl InstrumentedStore { + fn new(behavior: StoreBehavior) -> Self { + Self { inner: InMemoryStore::new(), behavior } + } + } + + impl KVStore for InstrumentedStore { + fn read( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, + ) -> impl Future, io::Error>> + 'static + Send { + KVStore::read(&self.inner, primary_namespace, secondary_namespace, key) + } + fn write( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, buf: Vec, + ) -> impl Future> + 'static + Send { + let write = if let StoreBehavior::FailWrite { attempts } = &self.behavior { + attempts.fetch_add(1, Ordering::Relaxed); + None + } else { + Some(KVStore::write(&self.inner, primary_namespace, secondary_namespace, key, buf)) + }; + async move { + match write { + Some(write) => write.await, + None => Err(io::Error::new(io::ErrorKind::Other, "write failed")), + } + } + } + fn remove( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, lazy: bool, + ) -> impl Future> + 'static + Send { + let remove = match &self.behavior { + StoreBehavior::FailRemove => None, + _ => Some(KVStore::remove( + &self.inner, + primary_namespace, + secondary_namespace, + key, + lazy, + )), + }; + async move { + match remove { + Some(remove) => remove.await, + None => Err(io::Error::new(io::ErrorKind::Other, "remove failed")), + } + } + } + fn list( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> impl Future, io::Error>> + 'static + Send { + let list = match &self.behavior { + StoreBehavior::FailList => None, + StoreBehavior::FailWrite { .. } + | StoreBehavior::FailRemove + | StoreBehavior::GateJournalList { .. } => { + Some(KVStore::list(&self.inner, primary_namespace, secondary_namespace)) + }, + }; + async move { + match list { + Some(list) => list.await, + None => Err(io::Error::new(io::ErrorKind::Other, "list failed")), + } + } + } + } + + impl PaginatedKVStore for InstrumentedStore { + fn list_paginated( + &self, primary_namespace: &str, secondary_namespace: &str, + page_token: Option, + ) -> impl Future> + 'static + Send { + let fails = matches!(&self.behavior, StoreBehavior::FailList); + let gate = match &self.behavior { + StoreBehavior::GateJournalList { snapshot_taken, allow_return } => { + if primary_namespace == INDEX_JOURNAL_PRIMARY_NAMESPACE { + let snapshot_taken = snapshot_taken.lock().unwrap().take(); + let allow_return = allow_return.lock().unwrap().take(); + match (snapshot_taken, allow_return) { + (Some(snapshot_taken), Some(allow_return)) => { + Some((snapshot_taken, allow_return)) + }, + _ => None, + } + } else { + None + } + }, + _ => None, + }; + let list = PaginatedKVStore::list_paginated( + &self.inner, + primary_namespace, + secondary_namespace, + page_token, + ); + async move { + if fails { + return Err(io::Error::new(io::ErrorKind::Other, "list_paginated failed")); + } + let response = list.await?; + if let Some((snapshot_taken, allow_return)) = gate { + let _ = snapshot_taken.send(()); + allow_return.await.map_err(|_| { + io::Error::new(io::ErrorKind::Other, "journal-list gate was dropped") + })?; + }; + Ok(response) + } + } + } + + #[tokio::test] + async fn write_read_list_remove() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + + let _cleanup = CleanupDir(base_dir.clone()); + + let primary_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("primary")).unwrap())); + let tier = setup_tier_store(primary_store, logger); + + do_read_write_remove_list_persist(&tier).await; + } + + #[tokio::test] + async fn ephemeral_routing() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + + let _cleanup = CleanupDir(base_dir.clone()); + + let primary_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("primary")).unwrap())); + let mut tier = setup_tier_store(Arc::clone(&primary_store), logger); + + let ephemeral_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("ephemeral")).unwrap())); + tier.set_ephemeral_store(Arc::clone(&ephemeral_store)); + + let data = vec![42u8; 32]; + + tier.write( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + data.clone(), + ) + .await + .unwrap(); + + tier.write( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_KEY, + data.clone(), + ) + .await + .unwrap(); + + let primary_read_ng = primary_store + .read( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + ) + .await; + let ephemeral_read_ng = ephemeral_store + .read( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + ) + .await; + + let primary_read_cm = primary_store + .read( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_KEY, + ) + .await; + let ephemeral_read_cm = ephemeral_store + .read( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_KEY, + ) + .await; + + assert!(primary_read_ng.is_err()); + assert_eq!(ephemeral_read_ng.unwrap(), data); + + assert!(ephemeral_read_cm.is_err()); + assert_eq!(primary_read_cm.unwrap(), data); + } + + #[tokio::test] + async fn external_pathfinding_scores_cache_routes_to_ephemeral_store() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + + let _cleanup = CleanupDir(base_dir.clone()); + + let primary_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("primary")).unwrap())); + let mut tier = setup_tier_store(Arc::clone(&primary_store), logger); + + let ephemeral_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("ephemeral")).unwrap())); + tier.set_ephemeral_store(Arc::clone(&ephemeral_store)); + + let data = vec![42u8; 32]; + tier.write( + SCORER_PERSISTENCE_PRIMARY_NAMESPACE, + SCORER_PERSISTENCE_SECONDARY_NAMESPACE, + EXTERNAL_PATHFINDING_SCORES_CACHE_KEY, + data.clone(), + ) + .await + .unwrap(); + + assert!(primary_store + .read( + SCORER_PERSISTENCE_PRIMARY_NAMESPACE, + SCORER_PERSISTENCE_SECONDARY_NAMESPACE, + EXTERNAL_PATHFINDING_SCORES_CACHE_KEY, + ) + .await + .is_err()); + assert_eq!( + tier.read( + SCORER_PERSISTENCE_PRIMARY_NAMESPACE, + SCORER_PERSISTENCE_SECONDARY_NAMESPACE, + EXTERNAL_PATHFINDING_SCORES_CACHE_KEY, + ) + .await + .unwrap(), + data + ); + + tier.remove( + SCORER_PERSISTENCE_PRIMARY_NAMESPACE, + SCORER_PERSISTENCE_SECONDARY_NAMESPACE, + EXTERNAL_PATHFINDING_SCORES_CACHE_KEY, + false, + ) + .await + .unwrap(); + assert!(ephemeral_store + .read( + SCORER_PERSISTENCE_PRIMARY_NAMESPACE, + SCORER_PERSISTENCE_SECONDARY_NAMESPACE, + EXTERNAL_PATHFINDING_SCORES_CACHE_KEY, + ) + .await + .is_err()); + } + + #[tokio::test] + async fn list_exposes_primary_and_routed_ephemeral_keys() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + + let _cleanup = CleanupDir(base_dir.clone()); + + let primary_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("primary")).unwrap())); + let mut tier = setup_tier_store(Arc::clone(&primary_store), logger); + + let ephemeral_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("ephemeral")).unwrap())); + tier.set_ephemeral_store(Arc::clone(&ephemeral_store)); + + // A durable root-namespace key, routed to primary since it isn't ephemeral-cached. + tier.write( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_KEY, + vec![1u8; 32], + ) + .await + .unwrap(); + + // The ephemeral-cached key, routed to the ephemeral store. + tier.write( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + vec![2u8; 32], + ) + .await + .unwrap(); + + // A decoy sitting in the ephemeral store under an unrelated namespace. This must + // never leak into a listing for that namespace just because an ephemeral + // store happens to be configured. + ephemeral_store + .write( + CHANNEL_MONITOR_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MONITOR_PERSISTENCE_SECONDARY_NAMESPACE, + "ephemeral-decoy", + vec![3u8; 32], + ) + .await + .unwrap(); + ephemeral_store + .write( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + "ephemeral-root-decoy", + vec![4u8; 32], + ) + .await + .unwrap(); + + // This is `list("", "")`: CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE and + // NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE are the same empty string, so both + // keys live in the exact namespace. + let root_keys = KVStore::list( + &tier, + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + ) + .await + .unwrap(); + + // Unpaginated listing exposes the logical view across both tiers without leaking + // unrelated keys from the ephemeral store. + assert!(root_keys.contains(&CHANNEL_MANAGER_PERSISTENCE_KEY.to_string())); + assert!(root_keys.contains(&NETWORK_GRAPH_PERSISTENCE_KEY.to_string())); + assert!(!root_keys.contains(&"ephemeral-root-decoy".to_string())); + + let monitor_keys = KVStore::list( + &tier, + CHANNEL_MONITOR_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MONITOR_PERSISTENCE_SECONDARY_NAMESPACE, + ) + .await + .unwrap(); + + // The unrelated-namespace decoy sitting in the ephemeral store must not leak + // into a listing for a namespace it was never routed to. + assert!(!monitor_keys.contains(&"ephemeral-decoy".to_string())); + } + + #[tokio::test] + async fn list_paginated_only_exposes_primary_keys() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + + let _cleanup = CleanupDir(base_dir.clone()); + + let primary_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("primary")).unwrap())); + let mut tier = setup_tier_store(Arc::clone(&primary_store), logger); + + let ephemeral_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("ephemeral")).unwrap())); + tier.set_ephemeral_store(Arc::clone(&ephemeral_store)); + + tier.write( + CHANNEL_MONITOR_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MONITOR_PERSISTENCE_SECONDARY_NAMESPACE, + "monitor-key", + vec![1u8; 32], + ) + .await + .unwrap(); + + // This decoy uses the same namespace but the opposite physical store, so it + // would show up if paginated listing routed to the wrong tier. + ephemeral_store + .write( + CHANNEL_MONITOR_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MONITOR_PERSISTENCE_SECONDARY_NAMESPACE, + "ephemeral-decoy", + vec![2u8; 32], + ) + .await + .unwrap(); + + // This key shares the network graph's namespace tuple ("", "") but is not + // itself an ephemeral-cached key, standing in for durable root-namespace data + // such as `manager`/`output_sweeper`/`peers`. It must still be listed even + // though the ephemeral store is configured and authoritative for + // `network_graph`/`scorer` specifically. + primary_store + .write( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + "other-root-namespace-key", + vec![3u8; 32], + ) + .await + .unwrap(); + + tier.write( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + vec![4u8; 32], + ) + .await + .unwrap(); + + let primary_response = PaginatedKVStore::list_paginated( + &tier, + CHANNEL_MONITOR_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MONITOR_PERSISTENCE_SECONDARY_NAMESPACE, + None, + ) + .await + .unwrap(); + assert_eq!(primary_response.keys, vec!["monitor-key".to_string()]); + + let root_response = PaginatedKVStore::list_paginated( + &tier, + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + None, + ) + .await + .unwrap(); + + assert_eq!(root_response.keys, vec!["other-root-namespace-key".to_string()]); + } + + #[tokio::test] + async fn listings_only_consult_ephemeral_store_for_routed_namespaces() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + + let _cleanup = CleanupDir(base_dir.clone()); + + let primary_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let mut tier = setup_tier_store(Arc::clone(&primary_store), logger); + + // An ephemeral store whose `list`/`list_paginated` always fail. + let ephemeral_store: Arc = + Arc::new(DynStoreWrapper(InstrumentedStore::new(StoreBehavior::FailList))); + tier.set_ephemeral_store(Arc::clone(&ephemeral_store)); + + // A durable key in a namespace that can never hold an ephemeral-cached key. + tier.write( + CHANNEL_MONITOR_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MONITOR_PERSISTENCE_SECONDARY_NAMESPACE, + "monitor-key", + vec![1u8; 32], + ) + .await + .unwrap(); + + // Listing that namespace must not consult (or depend on) the ephemeral store, so it + // succeeds even though the ephemeral list would fail. + let monitor_keys = KVStore::list( + &tier, + CHANNEL_MONITOR_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MONITOR_PERSISTENCE_SECONDARY_NAMESPACE, + ) + .await + .unwrap(); + assert_eq!(monitor_keys, vec!["monitor-key".to_string()]); + + // The paginated path always exposes only the primary store. + let monitor_page = PaginatedKVStore::list_paginated( + &tier, + CHANNEL_MONITOR_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MONITOR_PERSISTENCE_SECONDARY_NAMESPACE, + None, + ) + .await + .unwrap(); + assert_eq!(monitor_page.keys, vec!["monitor-key".to_string()]); + + // An unpaginated root listing must consult the authoritative ephemeral store. + assert!(KVStore::list( + &tier, + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + ) + .await + .is_err()); + } + + #[tokio::test] + async fn list_hides_stale_primary_copy_when_ephemeral_key_is_missing() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + + let _cleanup = CleanupDir(base_dir.clone()); + + let primary_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let mut tier = setup_tier_store(Arc::clone(&primary_store), logger); + + let ephemeral_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + tier.set_ephemeral_store(Arc::clone(&ephemeral_store)); + + // A durable root-namespace key that must always be discoverable. + tier.write( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_KEY, + vec![1u8; 32], + ) + .await + .unwrap(); + + // A stale copy of `network_graph` sitting in primary as if it had been persisted there + // before the ephemeral store was configured. The ephemeral store holds no copy, so + // `read` routes to ephemeral and would fail for this key. + primary_store + .write( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + vec![2u8; 32], + ) + .await + .unwrap(); + + let root_keys = KVStore::list( + &tier, + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + ) + .await + .unwrap(); + + // The ephemeral store is authoritative for routed keys, so its missing entry hides + // the stale primary copy. + assert!(root_keys.contains(&CHANNEL_MANAGER_PERSISTENCE_KEY.to_string())); + assert!(!root_keys.contains(&NETWORK_GRAPH_PERSISTENCE_KEY.to_string())); + } + + #[tokio::test] + async fn primary_backed_writes_preserve_latest_call_order() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + + let _cleanup = CleanupDir(base_dir.clone()); + + let primary_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("primary")).unwrap())); + let tier = setup_tier_store(primary_store, logger); + + let old_data = vec![1u8; 32]; + let new_data = vec![2u8; 32]; + + let old_write = tier.write( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_KEY, + old_data, + ); + let new_write = tier.write( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_KEY, + new_data.clone(), + ); + + new_write.await.unwrap(); + old_write.await.unwrap(); + + // Stale data doesn't overwrite latest + let persisted = tier + .read( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_KEY, + ) + .await + .unwrap(); + assert_eq!(persisted, new_data); + } + + #[tokio::test] + async fn failed_newer_backup_write_still_supersedes_older_write() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + + let _cleanup = CleanupDir(base_dir); + + let primary_store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let mut tier = setup_tier_store(Arc::clone(&primary_store), logger); + + let backup_write_attempts = Arc::new(AtomicUsize::new(0)); + let backup_store: Arc = + Arc::new(DynStoreWrapper(InstrumentedStore::new(StoreBehavior::FailWrite { + attempts: Arc::clone(&backup_write_attempts), + }))); + tier.set_backup_store(backup_store); + + let old_data = vec![1u8; 32]; + let new_data = vec![2u8; 32]; + let old_write = tier.write( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_KEY, + old_data, + ); + let new_write = tier.write( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_KEY, + new_data.clone(), + ); + + // The primary write succeeds, but the same newer write fails on the backup. + assert!(new_write.await.is_err()); + // The older operation must be treated as stale even though the newer operation failed. + old_write.await.unwrap(); + + let persisted = primary_store + .read( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_KEY, + ) + .await + .unwrap(); + assert_eq!(persisted, new_data); + assert_eq!(backup_write_attempts.load(Ordering::Relaxed), 1); + } + + #[tokio::test] + async fn ephemeral_writes_preserve_latest_call_order() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + + let _cleanup = CleanupDir(base_dir.clone()); + + let primary_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("primary")).unwrap())); + let mut tier = setup_tier_store(primary_store, logger); + + let ephemeral_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("ephemeral")).unwrap())); + tier.set_ephemeral_store(ephemeral_store); + + let old_data = vec![1u8; 32]; + let new_data = vec![2u8; 32]; + + let old_write = tier.write( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + old_data, + ); + let new_write = tier.write( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + new_data.clone(), + ); + + new_write.await.unwrap(); + old_write.await.unwrap(); + + let persisted = tier + .read( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + ) + .await + .unwrap(); + assert_eq!(persisted, new_data); + } + + #[tokio::test] + async fn ephemeral_removes_preserve_latest_call_order() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + + let _cleanup = CleanupDir(base_dir.clone()); + + let primary_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("primary")).unwrap())); + let mut tier = setup_tier_store(primary_store, logger); + + let ephemeral_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("ephemeral")).unwrap())); + tier.set_ephemeral_store(ephemeral_store); + + let data = vec![2u8; 32]; + + let stale_remove = tier.remove( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + true, + ); + let new_write = tier.write( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + data.clone(), + ); + + new_write.await.unwrap(); + stale_remove.await.unwrap(); + + let persisted = tier + .read( + NETWORK_GRAPH_PERSISTENCE_PRIMARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_SECONDARY_NAMESPACE, + NETWORK_GRAPH_PERSISTENCE_KEY, + ) + .await + .unwrap(); + assert_eq!(persisted, data); + } + + #[tokio::test] + async fn backup_write_is_part_of_success_path() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + + let _cleanup = CleanupDir(base_dir.clone()); + + let primary_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("primary")).unwrap())); + let mut tier = setup_tier_store(Arc::clone(&primary_store), logger); + + let backup_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("backup")).unwrap())); + tier.set_backup_store(Arc::clone(&backup_store)); + + let data = vec![42u8; 32]; + + tier.write( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_KEY, + data.clone(), + ) + .await + .unwrap(); + + let primary_read = primary_store + .read( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_KEY, + ) + .await; + let backup_read = backup_store + .read( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_KEY, + ) + .await; + + assert_eq!(primary_read.unwrap(), data); + assert_eq!(backup_read.unwrap(), data); + } + + #[tokio::test] + async fn backup_remove_is_part_of_success_path() { + let base_dir = random_storage_path(); + let log_path = base_dir.join("tier_store_test.log").to_string_lossy().into_owned(); + let logger = Arc::new(Logger::new_fs_writer(log_path, Level::Trace).unwrap()); + + let _cleanup = CleanupDir(base_dir.clone()); + + let primary_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("primary")).unwrap())); + let mut tier = setup_tier_store(Arc::clone(&primary_store), logger); + + let backup_store: Arc = + Arc::new(DynStoreWrapper(FilesystemStoreV2::new(base_dir.join("backup")).unwrap())); + tier.set_backup_store(Arc::clone(&backup_store)); + + let data = vec![42u8; 32]; + let key = CHANNEL_MANAGER_PERSISTENCE_KEY; + + tier.write( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + key, + data, + ) + .await + .unwrap(); + + tier.remove( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + key, + true, + ) + .await + .unwrap(); + + let primary_read = primary_store + .read( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + key, + ) + .await; + let backup_read = backup_store + .read( + CHANNEL_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, + CHANNEL_MANAGER_PERSISTENCE_SECONDARY_NAMESPACE, + key, + ) + .await; + + assert!(primary_read.is_err()); + assert!(backup_read.is_err()); + } +} diff --git a/tests/common/mod.rs b/tests/common/mod.rs index 85f618c958..2ba52846c9 100644 --- a/tests/common/mod.rs +++ b/tests/common/mod.rs @@ -56,7 +56,7 @@ use lightning::ln::msgs::SocketAddress; use lightning::routing::gossip::NodeAlias; use lightning::util::persist::{KVStore, PageToken, PaginatedKVStore, PaginatedListResponse}; use lightning_invoice::{Bolt11InvoiceDescription, Description}; -use lightning_persister::fs_store::v1::FilesystemStore; +use lightning_persister::fs_store::v2::FilesystemStoreV2; use lightning_types::payment::{PaymentHash, PaymentPreimage}; use logging::TestLogWriter; use rand::distr::Alphanumeric; @@ -1901,7 +1901,7 @@ impl PaginatedKVStore for TestSyncStore { struct TestSyncStoreInner { serializer: tokio::sync::RwLock<()>, test_store: InMemoryStore, - fs_store: FilesystemStore, + fs_store: FilesystemStoreV2, sqlite_store: SqliteStore, } @@ -1910,7 +1910,7 @@ impl TestSyncStoreInner { let serializer = tokio::sync::RwLock::new(()); let mut fs_dir = dest_dir.clone(); fs_dir.push("fs_store"); - let fs_store = FilesystemStore::new(fs_dir); + let fs_store = FilesystemStoreV2::new(fs_dir).unwrap(); let mut sql_dir = dest_dir.clone(); sql_dir.push("sqlite_store"); let sqlite_store = SqliteStore::new( diff --git a/tests/integration_tests_rust.rs b/tests/integration_tests_rust.rs index fd247f74cf..79a3162e70 100644 --- a/tests/integration_tests_rust.rs +++ b/tests/integration_tests_rust.rs @@ -38,6 +38,8 @@ use ldk_node::config::{ AsyncPaymentsRole, EsploraSyncConfig, ADDRESS_POOL_SIZE, DEFAULT_FULL_SCAN_STOP_GAP, }; use ldk_node::entropy::NodeEntropy; +#[cfg(not(feature = "uniffi"))] +use ldk_node::io::sqlite_store::SqliteStore; use ldk_node::liquidity::LSPS2ServiceConfig; use ldk_node::payment::{ ConfirmationStatus, PayerProofOptions, PaymentDetails, PaymentDirection, PaymentKind, @@ -4925,3 +4927,74 @@ async fn do_lsps2_multi_lsp_picks_cheapest(reverse_order: bool) { cheap.stop().unwrap(); expensive.stop().unwrap(); } + +// Builder backup-store configuration is not yet exposed via FFI (see #871) +#[cfg(not(feature = "uniffi"))] +#[tokio::test(flavor = "multi_thread", worker_threads = 1)] +async fn builder_configures_sqlite_backup_store() { + let (bitcoind, electrsd) = setup_bitcoind_and_electrsd(); + let chain_source = random_chain_source(&bitcoind, &electrsd); + + let mut config_a = random_config(); + config_a.store_type = TestStoreType::Sqlite; + let primary_dir = config_a.node_config.storage_dir_path.clone(); + let backup_dir = common::random_storage_path(); + + // Build node_a with backup storage configured + setup_builder!(builder_a, config_a.node_config.clone()); + builder_a.set_chain_source_esplora( + format!("http://{}", electrsd.esplora_url.as_ref().unwrap()), + None, + ); + builder_a.set_filesystem_logger(None, None); + builder_a.set_backup_storage_dir_path(backup_dir.to_str().unwrap().to_owned()); + + let node_a = builder_a.build(config_a.node_entropy.into()).unwrap(); + node_a.start().unwrap(); + assert!(node_a.status().is_running); + assert!(node_a.status().latest_fee_rate_cache_update_timestamp.is_some()); + + let mut config_b = random_config(); + config_b.node_config.manually_handle_unknown_bolt11_payments = true; + let node_b = setup_node(&chain_source, config_b); + + do_channel_full_cycle( + node_a, + node_b, + &bitcoind.client, + &electrsd.client, + false, + true, + true, + false, + ) + .await; + + let primary_store = SqliteStore::new( + primary_dir.into(), + Some(ldk_node::io::sqlite_store::SQLITE_DB_FILE_NAME.to_string()), + Some(ldk_node::io::sqlite_store::KV_TABLE_NAME.to_string()), + ) + .unwrap(); + + let backup_store = SqliteStore::new( + backup_dir, + Some(ldk_node::io::sqlite_store::SQLITE_BACKUP_DB_FILE_NAME.to_string()), + Some(ldk_node::io::sqlite_store::KV_TABLE_NAME.to_string()), + ) + .unwrap(); + + for (pn, sn, key) in [ + ("bdk_wallet", "", "descriptor"), + ("bdk_wallet", "", "change_descriptor"), + ("bdk_wallet", "", "network"), + ("", "", "node_metrics"), + ("", "", "events"), + ("", "", "peers"), + ] { + let primary = primary_store.read(pn, sn, key).await.unwrap(); + let backup = backup_store.read(pn, sn, key).await.unwrap(); + + assert_eq!(backup, primary, "backup mismatch for {pn}/{sn}/{key}"); + } +}