diff --git a/docs/contributing/benches.md b/docs/contributing/benches.md index be7ff50740..1db0321b0d 100644 --- a/docs/contributing/benches.md +++ b/docs/contributing/benches.md @@ -362,12 +362,13 @@ The fix is to wrap setup work that touches the runtime in `rt.block_on`: let rt = bench_runtime(); b.iter_batched( - // setup — wrap in block_on so the reactor is alive while - // FlureeBuilder::file(...).build() runs. + // Setup runs in block_on, so the reactor is alive while + // FlureeBuilder::file(...).build_async() runs. || rt.block_on(async { let dir = tempfile::tempdir().unwrap(); let fluree = FlureeBuilder::file(dir.path().to_string_lossy().to_string()) - .build() + .build_async() + .await .unwrap(); (dir, fluree) }), diff --git a/docs/getting-started/rust-api.md b/docs/getting-started/rust-api.md index 8e3ad830c9..436fb7e993 100644 --- a/docs/getting-started/rust-api.md +++ b/docs/getting-started/rust-api.md @@ -72,7 +72,7 @@ use fluree_db_api::{FlureeBuilder, Result}; #[tokio::main] async fn main() -> Result<()> { // Use file-backed storage for persistence - let fluree = FlureeBuilder::file("./data").build()?; + let fluree = FlureeBuilder::file("./data").build_async().await?; // Create a new ledger (or load an existing one) let ledger = fluree.create_ledger("mydb").await?; @@ -84,6 +84,9 @@ async fn main() -> Result<()> { } ``` +`build_async()` replays the storage's write-ahead log without blocking the async runtime. +`build()` builds the same instance from synchronous code that has entered a Tokio runtime (for example with `Runtime::enter`), replaying the log on the calling thread. + ### Bulk import (high throughput) For initial ledger bootstraps (large Turtle or JSON-LD datasets), Fluree exposes a bulk import @@ -94,7 +97,7 @@ use fluree_db_api::{FlureeBuilder, Result}; #[tokio::main] async fn main() -> Result<()> { - let fluree = FlureeBuilder::file("./data").build()?; + let fluree = FlureeBuilder::file("./data").build_async().await?; // `chunks_dir` can be: // - a directory containing *.ttl, *.trig, or *.jsonld files (sorted lexicographically), OR @@ -423,7 +426,7 @@ use tokio::sync::mpsc; #[tokio::main] async fn main() -> Result<()> { - let fluree = Arc::new(FlureeBuilder::file("./data").build()?); + let fluree = Arc::new(FlureeBuilder::file("./data").build_async().await?); // Plan against a borrowed GraphDb, then move the owned LedgerState into the // spawned producer (GraphDb borrows the state, so plan first). @@ -540,7 +543,7 @@ use fluree_db_api::{ #[tokio::main] async fn main() -> Result<()> { - let fluree = FlureeBuilder::file("./data").build()?; + let fluree = FlureeBuilder::file("./data").build_async().await?; // Get a cached ledger handle let handle = fluree.ledger_cached("mydb:main").await?; @@ -660,7 +663,7 @@ use std::fs::File; #[tokio::main] async fn main() -> Result<()> { - let fluree = FlureeBuilder::file("./data").build()?; + let fluree = FlureeBuilder::file("./data").build_async().await?; // Export as Turtle to a file let file = File::create("backup.ttl").unwrap(); @@ -751,7 +754,7 @@ use fluree_db_api::{FlureeBuilder, Result}; #[tokio::main] async fn main() -> Result<()> { // Caching is on by default — no extra call needed - let fluree = FlureeBuilder::file("./data").build()?; + let fluree = FlureeBuilder::file("./data").build_async().await?; // First call loads from storage let ledger = fluree.ledger("mydb:main").await?; @@ -768,7 +771,8 @@ To **disable** caching (e.g., for a CLI tool that runs once and exits): ```rust let fluree = FlureeBuilder::file("./data") .without_ledger_caching() - .build()?; + .build_async() + .await?; ``` #### Disconnecting Ledgers @@ -780,7 +784,7 @@ use fluree_db_api::{FlureeBuilder, Result}; #[tokio::main] async fn main() -> Result<()> { - let fluree = FlureeBuilder::file("./data").build()?; + let fluree = FlureeBuilder::file("./data").build_async().await?; // Load and use ledger let ledger = fluree.ledger("mydb:main").await?; @@ -814,7 +818,7 @@ use fluree_db_api::{FlureeBuilder, Result}; #[tokio::main] async fn main() -> Result<()> { - let fluree = FlureeBuilder::file("./data").build()?; + let fluree = FlureeBuilder::file("./data").build_async().await?; // Check if ledger exists (lightweight nameservice lookup) if fluree.ledger_exists("mydb:main").await? { @@ -850,7 +854,7 @@ use fluree_db_api::{FlureeBuilder, DropMode, DropStatus, Result}; #[tokio::main] async fn main() -> Result<()> { - let fluree = FlureeBuilder::file("./data").build()?; + let fluree = FlureeBuilder::file("./data").build_async().await?; // Soft drop: retract every branch in the nameservice, preserve artifacts let report = fluree.drop_ledger("mydb", DropMode::Soft).await?; @@ -949,7 +953,7 @@ use fluree_db_api::{FlureeBuilder, NotifyResult, RefreshOpts, Result}; #[tokio::main] async fn main() -> Result<()> { - let fluree = FlureeBuilder::file("./data").build()?; + let fluree = FlureeBuilder::file("./data").build_async().await?; // Load ledger into cache let _ledger = fluree.ledger_cached("mydb:main").await?; @@ -1022,7 +1026,7 @@ use serde_json::json; #[tokio::main] async fn main() -> Result<()> { - let fluree = FlureeBuilder::file("./data").build()?; + let fluree = FlureeBuilder::file("./data").build_async().await?; let handle = fluree.ledger_cached("mydb:main").await?; // Transaction returns the commit's t value @@ -1399,7 +1403,7 @@ async fn main() -> Result<()> { "https://acme-fluree.example.com", Some("eyJhbG...".to_string()), ) - .build()?; + .build_async().await?; let db = fluree.view("local-ledger:main").await?; @@ -1468,7 +1472,7 @@ use tokio::time::{sleep, Duration}; #[tokio::main] async fn main() -> Result<()> { - let fluree = Arc::new(FlureeBuilder::file("./data").build()?); + let fluree = Arc::new(FlureeBuilder::file("./data").build_async().await?); // Start background indexer let indexer = BackgroundIndexerWorker::new( @@ -1718,7 +1722,7 @@ async fn test_persistence() -> Result<()> { // Create ledger and write data { - let fluree = FlureeBuilder::file(path).build()?; + let fluree = FlureeBuilder::file(path).build_async().await?; let ledger = fluree.create_ledger("test").await?; let data = json!({"@context": {}, "@graph": [{"@id": "ex:test"}]}); @@ -1731,7 +1735,7 @@ async fn test_persistence() -> Result<()> { // Verify persistence by reopening { - let fluree = FlureeBuilder::file(path).build()?; + let fluree = FlureeBuilder::file(path).build_async().await?; let ledger = fluree.ledger("test:main").await?; assert!(ledger.t() > 0); @@ -2021,7 +2025,7 @@ use serde_json::json; #[tokio::main] async fn main() -> Result<()> { // Caching is on by default (required for stage) - let fluree = FlureeBuilder::file("./data").build()?; + let fluree = FlureeBuilder::file("./data").build_async().await?; // Get a cached handle let handle = fluree.ledger_cached("mydb:main").await?; @@ -2150,7 +2154,7 @@ use serde_json::json; #[tokio::main] async fn main() -> Result<()> { - let fluree = FlureeBuilder::file("./data").build()?; + let fluree = FlureeBuilder::file("./data").build_async().await?; // Get ledger info with optional context for IRI compaction let context = json!({ @@ -2225,7 +2229,7 @@ use serde_json::json; #[tokio::main] async fn main() -> Result<()> { - let fluree = FlureeBuilder::file("./data").build()?; + let fluree = FlureeBuilder::file("./data").build_async().await?; // Find all ledgers on main branch let query = json!({ diff --git a/docs/graph-sources/overview.md b/docs/graph-sources/overview.md index 72b9db604d..61363f3616 100644 --- a/docs/graph-sources/overview.md +++ b/docs/graph-sources/overview.md @@ -190,7 +190,7 @@ Graph sources are created and registered via the `fluree-db-api` Rust API, which ```rust use fluree_db_api::{FlureeBuilder, R2rmlCreateConfig}; -let fluree = FlureeBuilder::default().build().await?; +let fluree = FlureeBuilder::file("/path/to/data").build_async().await?; let config = R2rmlCreateConfig::new_direct( "execution-log", diff --git a/docs/graph-sources/r2rml.md b/docs/graph-sources/r2rml.md index f3e806fa6b..107f221c87 100644 --- a/docs/graph-sources/r2rml.md +++ b/docs/graph-sources/r2rml.md @@ -26,7 +26,7 @@ If you use **Direct S3** mode, Fluree resolves the current Iceberg metadata by r ```rust use fluree_db_api::{FlureeBuilder, R2rmlCreateConfig}; -let fluree = FlureeBuilder::default().build().await?; +let fluree = FlureeBuilder::file("/path/to/data").build_async().await?; let config = R2rmlCreateConfig::new_direct( "airlines-rdf", diff --git a/docs/indexing-and-search/bm25.md b/docs/indexing-and-search/bm25.md index e973868ab7..953d72196f 100644 --- a/docs/indexing-and-search/bm25.md +++ b/docs/indexing-and-search/bm25.md @@ -31,7 +31,7 @@ The CLI (`fluree bm25 create/list/sync/drop`) drives a server when one is reacha use fluree_db_api::{Bm25CreateConfig, FlureeBuilder}; use serde_json::json; -let fluree = FlureeBuilder::file("/path/to/data").build()?; +let fluree = FlureeBuilder::file("/path/to/data").build_async().await?; // Create a ledger and insert some data let ledger = fluree.create_ledger("docs:main").await?; diff --git a/docs/indexing-and-search/reindex.md b/docs/indexing-and-search/reindex.md index 30bd991d76..19024577cb 100644 --- a/docs/indexing-and-search/reindex.md +++ b/docs/indexing-and-search/reindex.md @@ -39,7 +39,7 @@ use fluree_db_api::{FlureeBuilder, ReindexOptions, ReindexResult}; // Create Fluree instance let fluree = FlureeBuilder::file("/path/to/data") - .build() + .build_async() .await?; // Reindex with default options @@ -55,7 +55,7 @@ println!("Root ID: {}", result.root_id); use fluree_db_api::{FlureeBuilder, ReindexOptions}; use fluree_db_indexer::IndexerConfig; -let fluree = FlureeBuilder::file("/path/to/data").build().await?; +let fluree = FlureeBuilder::file("/path/to/data").build_async().await?; let result = fluree.reindex("mydb:main", ReindexOptions::default() // Use custom index node sizes diff --git a/docs/operations/configuration.md b/docs/operations/configuration.md index 3fbe2c9a0f..5a215ed0d2 100644 --- a/docs/operations/configuration.md +++ b/docs/operations/configuration.md @@ -1129,7 +1129,8 @@ Register remote connections on the `FlureeBuilder`: let fluree = FlureeBuilder::file("./data") .remote_connection("acme", "https://acme-fluree.example.com", Some(token)) .remote_connection("partner", "https://partner.example.com", None) - .build()?; + .build_async() + .await?; ``` Each call registers a named connection. The name is used in SPARQL queries: diff --git a/docs/operations/running-fluree.md b/docs/operations/running-fluree.md index 505d031f32..d79951d7be 100644 --- a/docs/operations/running-fluree.md +++ b/docs/operations/running-fluree.md @@ -315,7 +315,7 @@ use serde_json::json; let fluree = FlureeBuilder::memory().build_memory(); // File-based persistence — typed -let fluree = FlureeBuilder::file("./data").build()?; +let fluree = FlureeBuilder::file("./data").build_async().await?; // AWS S3 (requires `aws` feature) — typed let fluree = FlureeBuilder::s3("my-bucket", "https://s3.us-east-1.amazonaws.com") @@ -337,7 +337,7 @@ Quick setup with typed builders: use fluree_db_api::FlureeBuilder; let fluree = FlureeBuilder::memory().build_memory(); // In-memory -let fluree = FlureeBuilder::file("./data").build()?; // File-based +let fluree = FlureeBuilder::file("./data").build_async().await?; // File-based let fluree = FlureeBuilder::s3("bucket", "endpoint").build_client().await?; // S3 ``` diff --git a/docs/security/encryption.md b/docs/security/encryption.md index 3951bcafd0..a2edd082fc 100644 --- a/docs/security/encryption.md +++ b/docs/security/encryption.md @@ -58,11 +58,12 @@ let client = FlureeBuilder::from_json_ld(&config)? // and the id that encrypts new writes. let fluree = FlureeBuilder::file("/data/fluree") .with_encryption_keys(vec![(1, old_key), (2, new_key)], 2)? - .build()?; + .build_async() + .await?; ``` A key set on the builder (`with_encryption_key*()`, `with_encryption_keys()`, -or `AES256Key` / `AES256Keys` in JSON-LD) is applied by every terminal build method — `build()`, `build_memory()`, `build_s3()`, +or `AES256Key` / `AES256Keys` in JSON-LD) is applied by every terminal build method: `build()`, `build_async()`, `build_memory()`, `build_s3()`, `build_client()` and friends — on every backend. The `build_*_encrypted()` methods remain for callers that want the key to be an explicit argument; `build_encrypted(key)` replaces any configured key set with that one key, as diff --git a/fluree-db-api/src/lib.rs b/fluree-db-api/src/lib.rs index d3f80dfe9d..afc6e8e871 100644 --- a/fluree-db-api/src/lib.rs +++ b/fluree-db-api/src/lib.rs @@ -15,7 +15,7 @@ //! use fluree_db_api::{FlureeBuilder, GraphDb}; //! //! // Create a file-backed Fluree instance -//! let fluree = FlureeBuilder::file("/data/fluree").build()?; +//! let fluree = FlureeBuilder::file("/data/fluree").build_async().await?; //! //! // Create a new ledger //! let ledger = fluree.create_ledger("mydb").await?; @@ -1334,7 +1334,7 @@ async fn build_s3_storage_from_config( /// Build a local (memory/file) storage instance from a StorageConfig. #[cfg(feature = "native")] -fn build_local_storage_from_config( +async fn build_local_storage_from_config( storage_config: &fluree_db_connection::config::StorageConfig, ) -> Result> { use fluree_db_connection::config::StorageType; @@ -1353,7 +1353,7 @@ fn build_local_storage_from_config( // is startup, so the startup sweep of crash-orphaned staging // files is taken here explicitly. storage.sweep_orphaned_staging(); - storage.recover_wal()?; + storage.recover_wal_async().await?; encrypt_storage_from_config(Arc::new(storage), storage_config) } StorageType::S3(_) => Err(ApiError::config( @@ -1367,7 +1367,7 @@ fn build_local_storage_from_config( /// Build a memory storage instance from a StorageConfig (non-native fallback). #[cfg(not(feature = "native"))] -fn build_local_storage_from_config( +async fn build_local_storage_from_config( storage_config: &fluree_db_connection::config::StorageConfig, ) -> Result> { use fluree_db_connection::config::StorageType; @@ -2477,22 +2477,46 @@ impl FlureeBuilder { /// appropriate). When indexing is enabled, a `BackgroundIndexerWorker` is /// spawned on the tokio runtime, so `build()` must be called within a /// tokio context. + /// + /// It recovers the storage root's WAL on the calling thread. + /// From async code, use [`Self::build_async`]. #[cfg(feature = "native")] pub fn build(mut self) -> Result { + let storage = self.open_file_storage()?; + storage.recover_wal()?; + Ok(self.build_file(storage)) + } + + /// Build a file-backed Fluree instance, recovering the storage root's WAL + /// without blocking the caller. + /// + /// It builds the same instance as [`Self::build`]. + #[cfg(feature = "native")] + pub async fn build_async(mut self) -> Result { + let storage = self.open_file_storage()?; + storage.recover_wal_async().await?; + Ok(self.build_file(storage)) + } + + /// Takes the builder's storage path and opens its file storage. + /// + /// Building an instance is startup, so this starts the sweep of staging + /// files a crash left behind. The sweep runs once per base path per + /// process. The file nameservice shares the tree and needs no sweep of its own. + #[cfg(feature = "native")] + fn open_file_storage(&mut self) -> Result { let path = self .storage_path .take() .ok_or_else(|| ApiError::config("File storage requires a path"))?; - let storage = self.file_storage(&path); - // Building the instance is startup: reclaim staging files a crash - // left behind. Explicit here rather than a side effect of `new`, and - // once per base path per process — the nameservice below shares this - // tree and needs no sweep of its own. storage.sweep_orphaned_staging(); - // Likewise the WAL: acknowledged writes a crash left unflushed - // are applied before anything reads this tree. - storage.recover_wal()?; + Ok(storage) + } + + /// Assembles a file-backed instance over `storage`, whose WAL is already recovered. + #[cfg(feature = "native")] + fn build_file(self, storage: FileStorage) -> Fluree { let nameservice = FileNameService::with_storage(storage.clone()); let event_bus = self.resolve_event_bus(); let notifying = @@ -2503,7 +2527,7 @@ impl FlureeBuilder { let attachment_provider_cell = Self::new_attachment_provider_cell(); let indexing_mode = self.start_background_indexing(&backend, ¬ifying, &attachment_provider_cell); - Ok(Self::finalize_with_backend( + Self::finalize_with_backend( self.ledger_cache_config, self.config, RuntimeParts { @@ -2518,7 +2542,7 @@ impl FlureeBuilder { self.remote_mounts, #[cfg(feature = "iceberg")] self.secret_resolver, - )) + ) } /// Build a Fluree instance with custom storage and nameservice. @@ -3217,8 +3241,8 @@ impl FlureeBuilder { // --- Local (memory/filesystem) --- match &self.config.index_storage.storage_type { - StorageType::Memory => self.build_client_memory(nameservice), - StorageType::File => self.build_client_file(nameservice), + StorageType::Memory => self.build_client_memory(nameservice).await, + StorageType::File => self.build_client_file(nameservice).await, StorageType::S3(_) => Err(ApiError::config( "S3 storage requires the 'aws' feature on fluree-db-api", )), @@ -3232,11 +3256,16 @@ impl FlureeBuilder { /// is `Some`, it replaces the default `MemoryNameService` (and /// the notifying wrapper + background indexer that ride along /// with it). - fn build_client_memory(self, nameservice: Option) -> Result { + async fn build_client_memory( + self, + nameservice: Option, + ) -> Result { let base_storage = self.encrypt_if_configured(Arc::new(MemoryStorage::new())); // Wrap with address identifier routing if configured - let storage = self.wrap_address_identifiers(base_storage)?; + let storage = self + .wrap_address_identifiers(base_storage, &self.config) + .await?; let backend = StorageBackend::Managed(storage); let event_bus = self.resolve_event_bus(); let index_config = self.derive_indexing(); @@ -3278,7 +3307,7 @@ impl FlureeBuilder { /// is `Some`, it replaces the default `FileNameService` (and /// the notifying wrapper + background indexer that ride along /// with it). - fn build_client_file(self, nameservice: Option) -> Result { + async fn build_client_file(self, nameservice: Option) -> Result { #[cfg(not(feature = "native"))] { let _ = nameservice; @@ -3308,12 +3337,14 @@ impl FlureeBuilder { // Client build is startup: take the explicit sweep of // crash-orphaned staging files here, where startup is known. file_storage.sweep_orphaned_staging(); - file_storage.recover_wal()?; + file_storage.recover_wal_async().await?; let ns_storage = file_storage.clone(); let base_storage = self.encrypt_if_configured(Arc::new(file_storage)); // Wrap with address identifier routing if configured - let storage = self.wrap_address_identifiers(base_storage)?; + let storage = self + .wrap_address_identifiers(base_storage, &self.config) + .await?; let backend = StorageBackend::Managed(storage); let event_bus = self.resolve_event_bus(); let index_config = self.derive_indexing(); @@ -3385,7 +3416,7 @@ impl FlureeBuilder { // Wrap with address identifier routing if configured let storage = self - .wrap_address_identifiers_aws(base_storage, aws_handle.config()) + .wrap_address_identifiers(base_storage, aws_handle.config()) .await?; let backend = StorageBackend::Managed(storage); let event_bus = self.resolve_event_bus(); @@ -3421,26 +3452,11 @@ impl FlureeBuilder { )) } - /// Wrap base storage with address identifier routing for local backends. - fn wrap_address_identifiers(&self, base_storage: Arc) -> Result> { - if let Some(addr_ids) = &self.config.address_identifiers { - let mut identifier_map = std::collections::HashMap::new(); - for (identifier, storage_config) in addr_ids { - let id_storage = build_local_storage_from_config(storage_config)?; - identifier_map.insert(identifier.to_string(), id_storage); - } - Ok(Arc::new(AddressIdentifierResolverStorage::new( - base_storage, - identifier_map, - ))) - } else { - Ok(base_storage) - } - } - - /// Wrap base storage with address identifier routing for AWS backends. - #[cfg(feature = "aws")] - async fn wrap_address_identifiers_aws( + /// Wrap base storage with routing for `config`'s address identifiers, if it has any. + /// + /// An S3 identifier needs the `aws` feature. + /// Without it, the identifier is rejected with a configuration error. + async fn wrap_address_identifiers( &self, base_storage: Arc, config: &ConnectionConfig, @@ -3449,8 +3465,9 @@ impl FlureeBuilder { let mut identifier_map = std::collections::HashMap::new(); for (identifier, storage_config) in addr_ids { let id_storage: Arc = match &storage_config.storage_type { + #[cfg(feature = "aws")] StorageType::S3(_) => build_s3_storage_from_config(storage_config).await?, - _ => build_local_storage_from_config(storage_config)?, + _ => build_local_storage_from_config(storage_config).await?, }; identifier_map.insert(identifier.to_string(), id_storage); } diff --git a/fluree-db-api/tests/grp_misc.rs b/fluree-db-api/tests/grp_misc.rs index 4e97469360..7ae4cf774a 100644 --- a/fluree-db-api/tests/grp_misc.rs +++ b/fluree-db-api/tests/grp_misc.rs @@ -21,6 +21,8 @@ mod it_edge_annotations_parse; mod it_fast_group_count; #[path = "it_file_backed.rs"] mod it_file_backed; +#[path = "it_file_startup_recovery.rs"] +mod it_file_startup_recovery; #[path = "it_file_storage_jsonld.rs"] mod it_file_storage_jsonld; #[path = "it_fuel_floor.rs"] diff --git a/fluree-db-api/tests/it_file_startup_recovery.rs b/fluree-db-api/tests/it_file_startup_recovery.rs new file mode 100644 index 0000000000..d68c7feca1 --- /dev/null +++ b/fluree-db-api/tests/it_file_startup_recovery.rs @@ -0,0 +1,45 @@ +//! Building a file-backed instance replays the WAL a crash left behind. + +#![cfg(feature = "native")] + +use fluree_db_api::FlureeBuilder; +use fluree_db_core::FileStorage; +use std::path::{Path, PathBuf}; + +async fn crash_with_a_logged_write(root: &Path) -> PathBuf { + FileStorage::crash_with_a_logged_write_for_test(root, "fluree:file://a.bin", b"logged") + .await + .unwrap() +} + +#[tokio::test] +async fn build_async_recovers_the_wal() { + let dir = tempfile::tempdir().unwrap(); + let path = crash_with_a_logged_write(dir.path()).await; + + let fluree = FlureeBuilder::file(dir.path().to_string_lossy().to_string()) + .without_indexing() + .build_async() + .await + .unwrap(); + + assert_eq!(std::fs::read(&path).unwrap(), b"logged"); + fluree + .create_ledger("startup/recovery:main") + .await + .expect("the recovered instance writes"); +} + +#[tokio::test] +async fn build_client_recovers_the_wal() { + let dir = tempfile::tempdir().unwrap(); + let path = crash_with_a_logged_write(dir.path()).await; + + let _client = FlureeBuilder::file(dir.path().to_string_lossy().to_string()) + .without_indexing() + .build_client() + .await + .unwrap(); + + assert_eq!(std::fs::read(&path).unwrap(), b"logged"); +} diff --git a/fluree-db-cli/src/commands/bm25.rs b/fluree-db-cli/src/commands/bm25.rs index 509d3d15f3..232c4253c5 100644 --- a/fluree-db-cli/src/commands/bm25.rs +++ b/fluree-db-cli/src/commands/bm25.rs @@ -182,7 +182,7 @@ async fn run_list( } async fn local_index_rows(dirs: &FlureeDir) -> CliResult> { - let fluree = build_fluree(dirs)?; + let fluree = build_fluree(dirs).await?; let ledgers = fluree.nameservice().all_records().await?; let sources = fluree.nameservice().all_graph_source_records().await?; @@ -393,7 +393,7 @@ async fn run_create( config = config.with_b(b); } - let fluree = build_fluree(dirs)?; + let fluree = build_fluree(dirs).await?; let result = fluree .create_full_text_index(config) .await @@ -502,7 +502,7 @@ async fn run_sync( return Ok(()); } - let fluree = build_fluree(dirs)?; + let fluree = build_fluree(dirs).await?; let result = match target_t { Some(t) => fluree.sync_bm25_index_to(index, t, None).await, None => fluree.sync_bm25_index(index).await, @@ -568,7 +568,7 @@ async fn run_drop( return report_remote_drop(index, &response); } - let fluree = build_fluree(dirs)?; + let fluree = build_fluree(dirs).await?; let result = fluree .drop_full_text_index(index) .await diff --git a/fluree-db-cli/src/commands/config_cmd.rs b/fluree-db-cli/src/commands/config_cmd.rs index 46c2db9e42..2b49832f0b 100644 --- a/fluree-db-cli/src/commands/config_cmd.rs +++ b/fluree-db-cli/src/commands/config_cmd.rs @@ -388,7 +388,7 @@ pub async fn run_set_origins(ledger: &str, file: &Path, dirs: &FlureeDir) -> Cli .map_err(|e| CliError::Config(format!("invalid origins config: {e}")))?; let ledger_id = context::to_ledger_id(ledger)?; - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; // Serialize to canonical bytes and store in CAS. let canonical_bytes = config.to_bytes(); diff --git a/fluree-db-cli/src/commands/context_cmd.rs b/fluree-db-cli/src/commands/context_cmd.rs index 369ff1085a..ba331d7f9f 100644 --- a/fluree-db-cli/src/commands/context_cmd.rs +++ b/fluree-db-cli/src/commands/context_cmd.rs @@ -31,7 +31,7 @@ pub async fn get( } } - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; match fluree.get_default_context(&ledger_id).await? { Some(ctx) => { println!( @@ -138,7 +138,7 @@ pub async fn set( } } - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; match fluree.set_default_context(&ledger_id, &ctx_value).await? { fluree_db_api::SetContextResult::Updated => { eprintln!("Default context updated for '{alias}'."); diff --git a/fluree-db-cli/src/commands/create.rs b/fluree-db-cli/src/commands/create.rs index bde3bdcdbc..17c7ca894e 100644 --- a/fluree-db-cli/src/commands/create.rs +++ b/fluree-db-cli/src/commands/create.rs @@ -412,7 +412,7 @@ pub async fn run( ))); } - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; match from { Some(path) if is_flpack_path(path) => { @@ -1354,7 +1354,7 @@ pub async fn run_memory_import( let include_user = !no_user; let commits = git_memory_commits(&repo_root, include_user)?; - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; // Create ledger + transact memory schema fluree diff --git a/fluree-db-cli/src/commands/delta.rs b/fluree-db-cli/src/commands/delta.rs index 6b7aeb9780..b198ffcdcc 100644 --- a/fluree-db-cli/src/commands/delta.rs +++ b/fluree-db-cli/src/commands/delta.rs @@ -255,7 +255,7 @@ fn create_config( #[cfg(feature = "delta")] async fn run_delta_map_local(args: DeltaMapArgs, dirs: &FlureeDir) -> CliResult<()> { - let fluree = crate::context::build_fluree(dirs)?; + let fluree = crate::context::build_fluree(dirs).await?; let config = fluree_db_api::DeltaCreateConfig { branch: args.branch.clone(), model: args.model.clone(), @@ -364,7 +364,7 @@ impl Ask<'_> { #[cfg(feature = "delta")] async fn local(&self, dirs: &FlureeDir) -> CliResult { - let fluree = crate::context::build_fluree(dirs)?; + let fluree = crate::context::build_fluree(dirs).await?; let unity = match self.unity() { Some(unity) => unity_config(unity)?, None => None, diff --git a/fluree-db-cli/src/commands/doc.rs b/fluree-db-cli/src/commands/doc.rs index 5ed51a4d8e..c97a5a5cc3 100644 --- a/fluree-db-cli/src/commands/doc.rs +++ b/fluree-db-cli/src/commands/doc.rs @@ -522,7 +522,7 @@ async fn run_ingest(args: DocIngestArgs, dirs: &FlureeDir) -> CliResult<()> { let fluree = if args.dry_run { None } else { - let fluree = build_fluree(dirs)?; + let fluree = build_fluree(dirs).await?; if !fluree.ledger_exists(&ledger_id).await? { fluree.create_ledger(&ledger_id).await?; eprintln!("{} created ledger {alias}", "→".dimmed()); @@ -1106,7 +1106,7 @@ async fn graph_source_present(fluree: &Fluree, id: &str) -> CliResult { async fn run_search(args: DocSearchArgs, dirs: &FlureeDir) -> CliResult<()> { let alias = context::resolve_ledger(args.ledger.as_deref(), dirs)?; - let fluree = build_fluree(dirs)?; + let fluree = build_fluree(dirs).await?; let (text_id, _) = text_index_id(&alias); // The vector lane needs embeddings on the chunks, not an index: it diff --git a/fluree-db-cli/src/commands/doc_sources.rs b/fluree-db-cli/src/commands/doc_sources.rs index 4d8dbdb038..3409111757 100644 --- a/fluree-db-cli/src/commands/doc_sources.rs +++ b/fluree-db-cli/src/commands/doc_sources.rs @@ -80,7 +80,7 @@ impl Opened { pub async fn open(source: &Source, dirs: &FlureeDir) -> CliResult { match &source.kind { SourceKind::Ledger(alias) => { - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let ledger_id = context::to_ledger_id(alias)?; if !fluree.ledger_exists(&ledger_id).await? { return Err(CliError::NotFound(format!( diff --git a/fluree-db-cli/src/commands/drop.rs b/fluree-db-cli/src/commands/drop.rs index 5ef5ede8d1..e3844b4fc3 100644 --- a/fluree-db-cli/src/commands/drop.rs +++ b/fluree-db-cli/src/commands/drop.rs @@ -94,7 +94,7 @@ async fn run_remote(name: &str, client: &RemoteLedgerClient) -> CliResult<()> { } async fn run_local(name: &str, dirs: &FlureeDir) -> CliResult<()> { - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; // Try dropping as a ledger first let report = fluree diff --git a/fluree-db-cli/src/commands/export.rs b/fluree-db-cli/src/commands/export.rs index b38578a0cc..a0ce75aed2 100644 --- a/fluree-db-cli/src/commands/export.rs +++ b/fluree-db-cli/src/commands/export.rs @@ -226,7 +226,7 @@ async fn run_ledger_archive( )); } - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; match output { Some(path) => { @@ -481,7 +481,7 @@ async fn run_local_rdf( )); } - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let format = parse_rdf_format(format_str)?; let mut builder = fluree.export(alias).format(format); diff --git a/fluree-db-cli/src/commands/graphql.rs b/fluree-db-cli/src/commands/graphql.rs index 1183019d08..df403b74d3 100644 --- a/fluree-db-cli/src/commands/graphql.rs +++ b/fluree-db-cli/src/commands/graphql.rs @@ -28,7 +28,7 @@ pub async fn run( dirs: &FlureeDir, ) -> CliResult<()> { let alias = context::resolve_ledger(explicit_ledger, dirs)?; - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let ledger_id = context::to_ledger_id(&alias)?; // The default context decides the GraphQL names and the form `id` values take, // so it is not optional here the way it is for a raw JSON-LD query. diff --git a/fluree-db-cli/src/commands/history.rs b/fluree-db-cli/src/commands/history.rs index 72736e2784..00c875837e 100644 --- a/fluree-db-cli/src/commands/history.rs +++ b/fluree-db-cli/src/commands/history.rs @@ -79,7 +79,7 @@ pub async fn run( )); } - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let ledger_view = fluree.ledger(&alias).await?; let result = fluree.query_connection(&query).await?; let json = result.to_jsonld(&ledger_view.snapshot)?; diff --git a/fluree-db-cli/src/commands/iceberg.rs b/fluree-db-cli/src/commands/iceberg.rs index cc334d999b..45cbc3132c 100644 --- a/fluree-db-cli/src/commands/iceberg.rs +++ b/fluree-db-cli/src/commands/iceberg.rs @@ -68,7 +68,7 @@ pub async fn run_iceberg_list( } } - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let gs_records = fluree.nameservice().all_graph_source_records().await?; let mut entries: Vec<_> = gs_records .into_iter() @@ -133,7 +133,7 @@ pub async fn run_iceberg_info( } } - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let gs_id = context::to_ledger_id(name)?; let gs = fluree .nameservice() @@ -202,7 +202,7 @@ pub async fn run_iceberg_drop( } } - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let gs_id = context::to_ledger_id(name)?; let gs = fluree .nameservice() @@ -581,7 +581,7 @@ fn print_remote_drop_response(response: &serde_json::Value) -> CliResult<()> { #[cfg(feature = "iceberg")] async fn run_iceberg_map_local(args: IcebergMapArgs, dirs: &FlureeDir) -> CliResult<()> { - let fluree = crate::context::build_fluree(dirs)?; + let fluree = crate::context::build_fluree(dirs).await?; let iceberg_config = build_iceberg_config(&args)?; if let Some(ref r2rml_path) = args.r2rml { diff --git a/fluree-db-cli/src/commands/index.rs b/fluree-db-cli/src/commands/index.rs index ca4877f586..c64da1beb3 100644 --- a/fluree-db-cli/src/commands/index.rs +++ b/fluree-db-cli/src/commands/index.rs @@ -18,7 +18,7 @@ pub struct IndexOutcome { /// a full rebuild when incremental isn't possible. pub async fn run_index(ledger: Option<&str>, dirs: &FlureeDir) -> CliResult<()> { let alias = context::resolve_ledger(ledger, dirs)?; - let fluree = build_fluree(dirs)?; + let fluree = build_fluree(dirs).await?; let ledger_id = context::to_ledger_id(&alias)?; // Verify ledger exists diff --git a/fluree-db-cli/src/commands/info.rs b/fluree-db-cli/src/commands/info.rs index 3de9dd6a55..845debd0f6 100644 --- a/fluree-db-cli/src/commands/info.rs +++ b/fluree-db-cli/src/commands/info.rs @@ -36,7 +36,7 @@ pub async fn run( Err(CliError::NotFound(_)) => { // Ledger not found — try graph source lookup let alias = context::resolve_ledger(ledger, dirs)?; - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let gs_id = context::to_ledger_id(&alias)?; if let Some(gs) = fluree.nameservice().lookup_graph_source(&gs_id).await? { if graph.is_some() { diff --git a/fluree-db-cli/src/commands/list.rs b/fluree-db-cli/src/commands/list.rs index 63d1e33eae..ce5e4b0845 100644 --- a/fluree-db-cli/src/commands/list.rs +++ b/fluree-db-cli/src/commands/list.rs @@ -18,7 +18,7 @@ pub async fn run(dirs: &FlureeDir, remote_flag: Option<&str>, direct: bool) -> C } } - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let active = config::read_active_ledger(dirs.data_dir()); let records = fluree.nameservice().all_records().await?; let gs_records = fluree.nameservice().all_graph_source_records().await?; diff --git a/fluree-db-cli/src/commands/log.rs b/fluree-db-cli/src/commands/log.rs index cbd4d2e20b..c79d9bf380 100644 --- a/fluree-db-cli/src/commands/log.rs +++ b/fluree-db-cli/src/commands/log.rs @@ -167,7 +167,7 @@ async fn run_local( )); } - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let ledger_id = context::to_ledger_id(&alias)?; // The same call the server's `GET /v1/fluree/log` makes, so `--direct` and diff --git a/fluree-db-cli/src/commands/materialize.rs b/fluree-db-cli/src/commands/materialize.rs index 30a2b21973..b73e1773d1 100644 --- a/fluree-db-cli/src/commands/materialize.rs +++ b/fluree-db-cli/src/commands/materialize.rs @@ -114,7 +114,7 @@ pub async fn run(dirs: &FlureeDir, params: &MaterializeParams<'_>) -> CliResult< // import producer thread. Reclaimed by the OS at process exit. Keeping ONE // provider (hence one catalog session) across the build is required for the // snapshot pin + watermark capture. - let fluree: &'static Fluree = Box::leak(Box::new(context::build_fluree(dirs)?)); + let fluree: &'static Fluree = Box::leak(Box::new(context::build_fluree(dirs).await?)); let provider = Arc::new(FlureeR2rmlProvider::new(fluree)); // CRITICAL-1 (#1529 review): a failed parity gate drops the WHOLE ledger NAME diff --git a/fluree-db-cli/src/commands/memory.rs b/fluree-db-cli/src/commands/memory.rs index 7d59ac1f52..a9b40e0709 100644 --- a/fluree-db-cli/src/commands/memory.rs +++ b/fluree-db-cli/src/commands/memory.rs @@ -75,13 +75,13 @@ pub async fn run(action: MemoryAction, dirs: &FlureeDir) -> CliResult<()> { } } -fn build_store(dirs: &FlureeDir) -> CliResult { +async fn build_store(dirs: &FlureeDir) -> CliResult { // Short-lived CLI commands keep a persistent (file-backed) ledger so that // `import` and the `init` legacy-ledger migration work and repeated // invocations don't rebuild from scratch. The long-lived `mcp serve` path // uses an ephemeral in-memory ledger instead (see `mcp_serve`), which is // what makes many concurrent MCP processes safe. - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; // Determine memory_dir: use .fluree-memory/ at the project root. // In unified (local) mode, data_dir is .fluree/ so its parent is the project root. @@ -100,7 +100,7 @@ fn build_store(dirs: &FlureeDir) -> CliResult { /// any `.ttl` files on disk. Use from every memory subcommand except /// `init`, which intentionally constructs an empty store before sync. async fn build_synced_store(dirs: &FlureeDir) -> CliResult { - let store = build_store(dirs)?; + let store = build_store(dirs).await?; store.ensure_synced().await.map_err(memory_err)?; Ok(store) } diff --git a/fluree-db-cli/src/commands/multi_query.rs b/fluree-db-cli/src/commands/multi_query.rs index 38914a2d5c..8f01266eab 100644 --- a/fluree-db-cli/src/commands/multi_query.rs +++ b/fluree-db-cli/src/commands/multi_query.rs @@ -229,7 +229,7 @@ async fn run_in_process( })?; inject_policy_into_envelope(&mut envelope, policy)?; - let fluree = Arc::new(context::build_fluree(dirs)?); + let fluree = Arc::new(context::build_fluree(dirs).await?); let mut builder = fluree.multi_query().envelope(envelope); if let Some(cfg) = formatter_config { builder = builder.format(cfg); diff --git a/fluree-db-cli/src/commands/query.rs b/fluree-db-cli/src/commands/query.rs index fd11aafaa8..e6ad3286f9 100644 --- a/fluree-db-cli/src/commands/query.rs +++ b/fluree-db-cli/src/commands/query.rs @@ -403,7 +403,7 @@ pub async fn run( Ok(gs) => gs, Err(CliError::NoActiveLedger) if force_connection => { context::QueryTarget::Ledger(LedgerMode::Local { - fluree: Box::new(context::build_fluree(dirs)?), + fluree: Box::new(context::build_fluree(dirs).await?), alias: String::new(), }) } diff --git a/fluree-db-cli/src/commands/sql.rs b/fluree-db-cli/src/commands/sql.rs index 93c6003c5f..0e09beea0d 100644 --- a/fluree-db-cli/src/commands/sql.rs +++ b/fluree-db-cli/src/commands/sql.rs @@ -172,7 +172,7 @@ pub async fn run_sql_check(name: &str, dirs: &FlureeDir) -> CliResult<()> { #[cfg(feature = "sql")] async fn run_sql_check_local(name: &str, dirs: &FlureeDir) -> CliResult<()> { - let fluree = crate::context::build_fluree(dirs)?; + let fluree = crate::context::build_fluree(dirs).await?; let id = if name.contains(':') { name.to_string() } else { @@ -208,7 +208,7 @@ async fn run_sql_check_local(_name: &str, _dirs: &FlureeDir) -> CliResult<()> { async fn run_sql_map_local(args: SqlMapArgs, dirs: &FlureeDir) -> CliResult<()> { use fluree_db_api::{SqlAuthConfig, SqlConfigValue, SqlDialect, WireProtocol}; - let fluree = crate::context::build_fluree(dirs)?; + let fluree = crate::context::build_fluree(dirs).await?; let mapping = read_mapping(&args)?; let mut config = fluree_db_api::SqlCreateConfig::new(&args.name, &args.endpoint, mapping); config.branch = args.branch.clone(); diff --git a/fluree-db-cli/src/commands/sync.rs b/fluree-db-cli/src/commands/sync.rs index dd4493d286..3fed182093 100644 --- a/fluree-db-cli/src/commands/sync.rs +++ b/fluree-db-cli/src/commands/sync.rs @@ -116,7 +116,7 @@ fn map_sync_auth_error(remote: &str, err: &str) -> Option { /// Build a SyncDriver with all configured remotes async fn build_sync_driver(dirs: &FlureeDir) -> CliResult<(SyncDriver, Arc)> { - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let config_store = Arc::new(TomlSyncConfigStore::new(dirs.config_dir().to_path_buf())); // Get the nameservice as RefPublisher @@ -280,7 +280,7 @@ pub async fn run_pull(ledger: Option<&str>, no_indexes: bool, dirs: &FlureeDir) .ok_or_else(|| CliError::Config("remote ledger-info response missing 't'".into()))?; // Resolve local head. - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let storage = fluree .backend() .admin_storage_cloned() @@ -836,7 +836,7 @@ pub async fn run_push(ledger: Option<&str>, dirs: &FlureeDir) -> CliResult<()> { .and_then(|s| s.parse().ok()); // Resolve local head. - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let local_ref = fluree .nameservice_mode() .get_ref(&ledger_id, RefKind::CommitHead) @@ -955,7 +955,7 @@ pub async fn run_publish( ); // Resolve local head. - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let local_ref = fluree .nameservice_mode() .get_ref(&ledger_id, RefKind::CommitHead) @@ -1133,7 +1133,7 @@ pub async fn run_clone( } // Create the local ledger. - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let storage = fluree .backend() .admin_storage_cloned() @@ -1451,7 +1451,7 @@ pub async fn run_clone_origin( } // 4. Create the local ledger. - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let storage = fluree .backend() .admin_storage_cloned() @@ -1789,7 +1789,7 @@ async fn run_pull_via_origins( no_indexes: bool, dirs: &FlureeDir, ) -> CliResult<()> { - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let storage = fluree .backend() .admin_storage_cloned() diff --git a/fluree-db-cli/src/commands/track.rs b/fluree-db-cli/src/commands/track.rs index ac9d3fe879..a931778029 100644 --- a/fluree-db-cli/src/commands/track.rs +++ b/fluree-db-cli/src/commands/track.rs @@ -105,7 +105,7 @@ async fn run_add( let effective_remote_alias = crate::context::to_ledger_id(remote_alias.unwrap_or(ledger))?; // Check mutual exclusion: refuse if local ledger exists - let fluree = crate::context::build_fluree(dirs)?; + let fluree = crate::context::build_fluree(dirs).await?; let local_ledger_id = &local_alias; if fluree.ledger_exists(local_ledger_id).await.unwrap_or(false) { return Err(CliError::Config(format!( diff --git a/fluree-db-cli/src/commands/use_cmd.rs b/fluree-db-cli/src/commands/use_cmd.rs index 3c009a06fb..d81c297b36 100644 --- a/fluree-db-cli/src/commands/use_cmd.rs +++ b/fluree-db-cli/src/commands/use_cmd.rs @@ -4,7 +4,7 @@ use crate::error::{CliError, CliResult}; use fluree_db_api::server_defaults::FlureeDir; pub async fn run(ledger: &str, dirs: &FlureeDir) -> CliResult<()> { - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let ledger_id = context::to_ledger_id(ledger)?; // Check if it's a local ledger diff --git a/fluree-db-cli/src/commands/verify.rs b/fluree-db-cli/src/commands/verify.rs index f848aa5942..8752bb9d24 100644 --- a/fluree-db-cli/src/commands/verify.rs +++ b/fluree-db-cli/src/commands/verify.rs @@ -10,7 +10,7 @@ pub async fn run( dirs: &FlureeDir, ) -> CliResult<()> { let alias = context::resolve_ledger(ledger, dirs)?; - let fluree = context::build_fluree(dirs)?; + let fluree = context::build_fluree(dirs).await?; let ledger_id = context::to_ledger_id(&alias)?; let report = fluree.verify_ledger(&ledger_id, limit).await.map_err(|e| { diff --git a/fluree-db-cli/src/context.rs b/fluree-db-cli/src/context.rs index b3d64b1394..9344d25d1c 100644 --- a/fluree-db-cli/src/context.rs +++ b/fluree-db-cli/src/context.rs @@ -119,7 +119,7 @@ pub async fn resolve_query_target( return Ok(QueryTarget::Ledger(mode)); } - let fluree = build_fluree(dirs)?; + let fluree = build_fluree(dirs).await?; // Check if local ledger exists (local wins) let ledger_id = to_ledger_id(ledger_part)?; @@ -539,7 +539,7 @@ pub fn resolve_ledger(explicit: Option<&str>, dirs: &FlureeDir) -> CliResult CliResult { +pub async fn build_fluree(dirs: &FlureeDir) -> CliResult { let storage = config::resolve_storage_path(dirs); let storage_str = storage.to_string_lossy().to_string(); let mut builder = FlureeBuilder::file(storage_str).without_ledger_caching(); @@ -562,7 +562,8 @@ pub fn build_fluree(dirs: &FlureeDir) -> CliResult { .with_novelty_thresholds(min_bytes, max_bytes); builder - .build() + .build_async() + .await .map_err(|e| CliError::Config(format!("failed to initialize Fluree: {e}"))) } @@ -799,7 +800,7 @@ mod tests { /// `fluree iceberg map` (which needs a live catalog) so the resolution /// branch can be tested in isolation. async fn register_graph_source(dirs: &FlureeDir, name: &str) { - let fluree = build_fluree(dirs).unwrap(); + let fluree = build_fluree(dirs).await.unwrap(); fluree .publisher() .unwrap() diff --git a/fluree-db-connection/src/lib.rs b/fluree-db-connection/src/lib.rs index 821f3d56dc..83d59c4ab9 100644 --- a/fluree-db-connection/src/lib.rs +++ b/fluree-db-connection/src/lib.rs @@ -184,30 +184,7 @@ pub async fn connect_from_config(config: ConnectionConfig) -> Result Result { match &config.index_storage.storage_type { - StorageType::File => { - #[cfg(not(all(feature = "native", not(target_arch = "wasm32"))))] - { - Err(ConnectionError::unsupported_component( - "https://ns.flur.ee/system#filePath (native feature disabled)", - )) - } - #[cfg(all(feature = "native", not(target_arch = "wasm32")))] - { - let path = - config.index_storage.path.as_ref().ok_or_else(|| { - ConnectionError::invalid_config("File storage requires path") - })?; - let storage = FileStorage::new(path.as_ref()) - .with_durability(Durability::resolve(config.index_storage.durability)); - // Opening a connection is startup, and startup is the layer - // that knows it: reclaiming staging files a crash left behind - // is an explicit action here, not a side effect of holding a - // storage handle. - storage.sweep_orphaned_staging(); - storage.recover_wal()?; - Ok(ConnectionHandle::File { config, storage }) - } - } + StorageType::File => file_connection(config), StorageType::Memory => { let storage = MemoryStorage::new(); Ok(ConnectionHandle::Memory { config, storage }) @@ -226,30 +203,7 @@ fn create_sync_connection(config: ConnectionConfig) -> Result async fn create_async_connection(config: ConnectionConfig) -> Result { match &config.index_storage.storage_type { StorageType::S3(s3_config) => create_aws_connection(config.clone(), s3_config).await, - StorageType::File => { - #[cfg(not(all(feature = "native", not(target_arch = "wasm32"))))] - { - Err(ConnectionError::unsupported_component( - "https://ns.flur.ee/system#filePath (native feature disabled)", - )) - } - #[cfg(all(feature = "native", not(target_arch = "wasm32")))] - { - let path = - config.index_storage.path.as_ref().ok_or_else(|| { - ConnectionError::invalid_config("File storage requires path") - })?; - let storage = FileStorage::new(path.as_ref()) - .with_durability(Durability::resolve(config.index_storage.durability)); - // Opening a connection is startup, and startup is the layer - // that knows it: reclaiming staging files a crash left behind - // is an explicit action here, not a side effect of holding a - // storage handle. - storage.sweep_orphaned_staging(); - storage.recover_wal()?; - Ok(ConnectionHandle::File { config, storage }) - } - } + StorageType::File => file_connection_async(config).await, StorageType::Memory => { let storage = MemoryStorage::new(); Ok(ConnectionHandle::Memory { config, storage }) @@ -267,10 +221,65 @@ async fn create_async_connection(config: ConnectionConfig) -> Result Err(ConnectionError::unsupported_component( "S3 storage requires the 'aws' feature to be enabled", )), + StorageType::File => file_connection_async(config).await, _ => create_sync_connection(config), } } +/// Opens the file storage `config` names and recovers its WAL on the calling thread. +fn file_connection(config: ConnectionConfig) -> Result { + #[cfg(not(all(feature = "native", not(target_arch = "wasm32"))))] + { + let _ = config; + Err(file_storage_unsupported()) + } + #[cfg(all(feature = "native", not(target_arch = "wasm32")))] + { + let storage = open_file_storage(&config)?; + storage.recover_wal()?; + Ok(ConnectionHandle::File { config, storage }) + } +} + +/// Opens the file storage `config` names and recovers its WAL without blocking the caller. +async fn file_connection_async(config: ConnectionConfig) -> Result { + #[cfg(not(all(feature = "native", not(target_arch = "wasm32"))))] + { + let _ = config; + Err(file_storage_unsupported()) + } + #[cfg(all(feature = "native", not(target_arch = "wasm32")))] + { + let storage = open_file_storage(&config)?; + storage.recover_wal_async().await?; + Ok(ConnectionHandle::File { config, storage }) + } +} + +/// Opens the file storage `config` names and starts the sweep of crash-orphaned staging files. +/// +/// Opening a connection is startup, so the sweep runs here. +/// Holding a storage handle does not start it. +#[cfg(all(feature = "native", not(target_arch = "wasm32")))] +fn open_file_storage(config: &ConnectionConfig) -> Result { + let path = config + .index_storage + .path + .as_ref() + .ok_or_else(|| ConnectionError::invalid_config("File storage requires path"))?; + let storage = FileStorage::new(path.as_ref()) + .with_durability(Durability::resolve(config.index_storage.durability)); + storage.sweep_orphaned_staging(); + Ok(storage) +} + +#[cfg(not(all(feature = "native", not(target_arch = "wasm32"))))] +fn file_storage_unsupported() -> ConnectionError { + ConnectionError::unsupported_component( + "https://ns.flur.ee/system#filePath (native feature disabled)", + ) +} + /// Create AWS connection from parsed JSON-LD config /// /// Uses StorageRegistry for storage sharing when the same @id is referenced @@ -448,4 +457,38 @@ mod tests { "opening a file connection did not sweep a stale staging orphan" ); } + + #[cfg(all(feature = "native", not(target_arch = "wasm32")))] + async fn crash_with_a_logged_write(root: &std::path::Path) -> std::path::PathBuf { + FileStorage::crash_with_a_logged_write_for_test(root, "fluree:file://a.bin", b"logged") + .await + .unwrap() + } + + /// Opening a file connection replays the WAL a crash left behind. + #[tokio::test] + #[cfg(all(feature = "native", not(target_arch = "wasm32")))] + async fn file_connection_recovers_the_wal() { + let dir = tempfile::tempdir().unwrap(); + let path = crash_with_a_logged_write(dir.path()).await; + + let _handle = + create_sync_connection(ConnectionConfig::file(dir.path().to_str().unwrap())).unwrap(); + + assert_eq!(std::fs::read(&path).unwrap(), b"logged"); + } + + /// Opening a file connection from async code replays the WAL a crash left behind. + #[tokio::test] + #[cfg(all(feature = "native", not(target_arch = "wasm32")))] + async fn async_file_connection_recovers_the_wal() { + let dir = tempfile::tempdir().unwrap(); + let path = crash_with_a_logged_write(dir.path()).await; + + let _handle = connect_from_config(ConnectionConfig::file(dir.path().to_str().unwrap())) + .await + .unwrap(); + + assert_eq!(std::fs::read(&path).unwrap(), b"logged"); + } } diff --git a/fluree-db-core/src/storage/file.rs b/fluree-db-core/src/storage/file.rs index 2b31ca5e5a..18e531c377 100644 --- a/fluree-db-core/src/storage/file.rs +++ b/fluree-db-core/src/storage/file.rs @@ -496,6 +496,23 @@ fn create_new_in_place(path: &Path, bytes: &[u8], policy: &WritePolicy) -> std:: Ok(true) } +/// Runs `f` on the blocking pool and returns its result. +/// +/// A panic in `f` resumes on the caller. +/// The only error is the task's cancellation, as during runtime shutdown. +async fn run_blocking(f: F) -> std::result::Result +where + F: FnOnce() -> T + Send + 'static, + T: Send + 'static, +{ + tokio::task::spawn_blocking(f) + .await + .or_else(|e| match e.try_into_panic() { + Ok(payload) => std::panic::resume_unwind(payload), + Err(e) => Err(e), + }) +} + /// File-based storage backed by `tokio::fs`. #[derive(Debug, Clone)] pub struct FileStorage { @@ -570,13 +587,31 @@ const WAL_RETRY_INTERVAL: std::time::Duration = std::time::Duration::from_secs(2 /// /// `durability` and `log` stay fixed while it is held. Fields drop in /// declaration order, so the key stripe is released before the root gate. +/// +/// It is `!Send`, so a `Send` future cannot hold it across an `.await`. +/// A `spawn_blocking` task cannot return one either. +/// The type does not stop a `!Send` future, such as a `LocalSet` task, from +/// holding it across an `.await`. Holders must take and release it on a +/// blocking thread. struct OperationHold { _key_stripe: Option, log: Option>, durability: Durability, _root_gate: tokio::sync::OwnedRwLockReadGuard<()>, + _not_send: std::marker::PhantomData<*const ()>, } +// Stops compiling if `OperationHold` becomes `Send`. +// Both impls would then apply, so the `some_item` path below is ambiguous. +const _: fn() = || { + trait AmbiguousIfSend { + fn some_item() {} + } + impl AmbiguousIfSend<()> for T {} + impl AmbiguousIfSend for T {} + let _ = >::some_item; +}; + impl FileStorage { /// Create a new file storage with the given base path /// @@ -846,25 +881,55 @@ impl FileStorage { } } - /// Replay the WAL an earlier run left under this root, if any, so - /// state acknowledged before a crash is on disk before the first read. + /// Replays the WAL an earlier run left under this root, if any. + /// + /// State acknowledged before a crash is then on disk before the first read. + /// The connection and builder paths call it at startup, like + /// [`Self::sweep_orphaned_staging`]. Constructing a handle does not call it. + /// + /// Blocking. It takes the root gate exclusively, so the calling thread + /// waits until every in-flight file operation on this root has finished. + /// The replay itself is synchronous file I/O on the calling thread. /// - /// A startup action like [`Self::sweep_orphaned_staging`]: the connection - /// and builder paths call it, constructing a handle never does. Blocking. - /// Replays under any durability setting, so an operator who switched back + /// It replays under any durability setting. An operator who switched back /// to per-write flushing after a crash still sees the acknowledged tail. - /// Leaves no trace on a root that never journaled. + /// A root that never journaled is left untouched. pub fn recover_wal(&self) -> Result<()> { - let _root = futures::executor::block_on(self.root_gate()?.write_owned()); + let root_gate = futures::executor::block_on(self.root_gate()?.write_owned()); + self.recover_wal_locked(&root_gate) + } + + /// Replays the WAL an earlier run left under this root, if any, without blocking the caller. + /// + /// It replays what [`Self::recover_wal`] replays. + /// It waits for the root gate asynchronously. + /// The replay then runs on the blocking pool, which holds the root gate until it finishes. + /// Dropping this future before it takes the root gate replays nothing. + /// Dropping it afterward leaves the replay to finish on the blocking pool. + pub async fn recover_wal_async(&self) -> Result<()> { + let root_gate = self.root_gate()?.write_owned().await; + let storage = self.clone(); + run_blocking(move || storage.recover_wal_locked(&root_gate)) + .await + .map_err(|e| crate::error::Error::io(format!("WAL recovery join: {e}")))? + } + + /// Replays the WAL under this root while `_root_gate` holds its root gate exclusively. + /// + /// Blocking: the replay is synchronous file I/O. + fn recover_wal_locked( + &self, + _root_gate: &tokio::sync::OwnedRwLockWriteGuard<()>, + ) -> Result<()> { if self.durability == Durability::Wal { self.attach_wal_locked(false)?; } else { - // Replay and let go: this handle is not going to journal. + // This handle does not journal, so the log is replayed and then released. Wal::acquire(&self.base_path, self.wal_owner.as_deref(), false) .map(drop) .map_err(|e| Self::recovery_error(&self.base_path, e))?; } - // A root several processes journal: apply what a stopped one left. + // Several processes can journal one root. Apply what stopped owners left. wal::replay_unowned(&self.base_path) .map(drop) .map_err(|e| Self::recovery_error(&self.base_path, e)) @@ -1005,6 +1070,7 @@ impl FileStorage { log, durability, _root_gate: root, + _not_send: std::marker::PhantomData, }) } @@ -1074,6 +1140,26 @@ impl FileStorage { Ok(()) } + /// Leave `root` with a write its WAL logged but a crash kept off disk. + /// + /// The write puts `bytes` at `address`. Returns the path of the file the crash lost. + #[doc(hidden)] + pub async fn crash_with_a_logged_write_for_test( + root: impl AsRef, + address: &str, + bytes: &[u8], + ) -> Result { + let storage = Self::new(root.as_ref()).with_durability(Durability::Wal); + storage.hold_wal_segments_for_test()?; + storage.write_bytes(address, bytes).await?; + storage.sync().await?; + storage.simulate_crash_for_test(); + let path = storage.resolve_path(address)?; + std::fs::remove_file(&path) + .map_err(|e| crate::error::Error::io(format!("remove {}: {e}", path.display())))?; + Ok(path) + } + /// Run `f` on the writing thread between an oversized write's checkpoint /// and its write. Shared across clones. #[cfg(test)] @@ -1948,10 +2034,10 @@ impl StorageCas for FileStorage { // An async task holding the root gate could need that thread to run. // Both would then wait forever. let storage = self.clone(); - // `spawn_blocking` does not inherit the caller's span. + // The blocking task does not inherit the caller's span. // Entering it keeps the closure's events and "cas phases" under the caller's span. let span = tracing::Span::current(); - tokio::task::spawn_blocking(move || { + run_blocking(move || { let _entered = span.enter(); // Phase 1: take the key lock and read. let phase = std::time::Instant::now(); @@ -1976,11 +2062,7 @@ impl StorageCas for FileStorage { } }) .await - .unwrap_or_else(|e| match e.try_into_panic() { - // A panic in `f` belongs to the caller, so it resumes on the caller's task. - Ok(payload) => std::panic::resume_unwind(payload), - Err(e) => Err(StorageExtError::io(format!("spawn_blocking join: {e}"))), - }) + .map_err(|e| StorageExtError::io(format!("spawn_blocking join: {e}")))? } } @@ -3177,6 +3259,31 @@ mod wal_tests { assert_eq!(wal_entries(dir.path()), ["LOCK"], "replay retires the log"); } + /// The async path replays the same records and retires the log. + /// It reads the files from disk, since `read_bytes` replays on a miss by itself. + #[tokio::test] + async fn recover_wal_async_restores_acknowledged_writes_after_a_crash() { + let dir = tempfile::tempdir().unwrap(); + let storage = with_wal(dir.path()); + storage.hold_wal_segments_for_test().unwrap(); + commit(&storage, 1).await; + commit(&storage, 2).await; + storage.simulate_crash_for_test(); + drop(storage); + let resolver = FileStorage::new(dir.path()); + for address in [TXN, COMMIT, HEAD] { + std::fs::remove_file(resolver.resolve_path(address).unwrap()).unwrap(); + } + + let storage = FileStorage::new(dir.path()).with_durability(Durability::Wal); + storage.recover_wal_async().await.unwrap(); + let on_disk = |address| std::fs::read(resolver.resolve_path(address).unwrap()).unwrap(); + assert_eq!(on_disk(TXN), b"r\x02"); + assert_eq!(on_disk(COMMIT), b"c\x02"); + assert_eq!(on_disk(HEAD), b"h\x02"); + assert_eq!(wal_entries(dir.path()), ["LOCK"], "replay retires the log"); + } + /// Replay applies records in order, so a delete cannot resurrect what it /// removed, and a later write after the delete wins. #[tokio::test] @@ -3753,6 +3860,53 @@ mod wal_tests { .expect("recover_wal deadlocked against an in-flight CAS"); } + /// `recover_wal_async` waits for the root gate without blocking its runtime thread. + /// + /// Another thread holds an operation on the root, so recovery has to wait. + /// A timer on the same current-thread runtime must still fire during that wait. + /// The test releases the operation only after the timer fires. + #[test] + fn recover_wal_async_waits_for_the_root_gate_without_blocking_the_runtime() { + use std::time::Duration; + let (done_tx, done_rx) = std::sync::mpsc::channel(); + // A separate thread lets a blocked runtime fail the test instead of hanging it. + std::thread::spawn(move || { + let dir = tempfile::tempdir().unwrap(); + let storage = FileStorage::new(dir.path()); + let (held_tx, held_rx) = std::sync::mpsc::channel(); + let (release_tx, release_rx) = std::sync::mpsc::channel::<()>(); + let holder = { + let storage = storage.clone(); + std::thread::spawn(move || { + let _operation = storage.begin_operation(storage.durability, "k").unwrap(); + held_tx.send(()).unwrap(); + release_rx.recv().unwrap(); + }) + }; + held_rx.recv().unwrap(); + + let rt = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .unwrap(); + rt.block_on(async { + let recovery = tokio::spawn(async move { storage.recover_wal_async().await }); + tokio::time::sleep(Duration::from_millis(50)).await; + assert!( + !recovery.is_finished(), + "recovery ran without the root gate" + ); + release_tx.send(()).unwrap(); + recovery.await.unwrap().unwrap(); + }); + holder.join().unwrap(); + done_tx.send(()).unwrap(); + }); + done_rx + .recv_timeout(Duration::from_secs(10)) + .expect("waiting for the root gate blocked the runtime thread"); + } + #[tokio::test] async fn retry_does_not_replay_crash_history_over_completed_fallback_writes() { let dir = tempfile::tempdir().unwrap();