Repository navigation
perf(storage): recover the WAL without blocking the async runtime - #2011
Conversation
aaj3f
left a comment
There was a problem hiding this comment.
@zonotope this is nice and a nice follow-up to #2001 & #2003. I think the only item from the review worth pointing out up-top (though in reality it feels quite minor) is that the structural guarantee put in place with the !Send marker (_not_send: std::marker::PhantomData<*const ()>), is only a guarantee insofar as a later refactor doesn't remove it. Which is to say, all tests pass green if you delete that field from OperationHold. Some kind of static assertion may be nice to better ensure (rather than merely with comments) that this is durable and not accidentally refactored away. Claude provided some more details below inline. Really nothing other than some nit items though (full review below):
Having recover_wal_locked take the write guard as a parameter means the replay can't run without the gate, moving that guard into the blocking task keeps it off async tasks entirely, and sharing open_file_storage/build_file between build and build_async (with the same split in the connection crate) leaves one copy of each startup path instead of two. I probed the parts that are easy to get subtly wrong and they hold: a panic in the replay comes back to the caller as the original panic and frees the root gate; a recover_wal_async dropped after its first poll still finishes the replay before anyone else gets the gate; and one aborted while it waits replays nothing. The !Send marker also does what the body says — the pre-#2001 compare-and-swap shape and a hold across .await in a Send future both fail with the *const () error.
Notes inline: one 🟠 to pin the !Send marker with a static assertion, one ❓ about S3 address identifiers on file/memory clients, and two docs 🟡s. One more with no line to sit on:
- This is more of a question than a suggestion. The title feeds the release notes, and
build_asyncis new public API that embedders have to switch to before they see any of this (build_client*and the connection paths switch on their own). Something like "perf(storage): recover the WAL without blocking the async runtime; add FlureeBuilder::build_async" would tell them.
Adherence to repo commitments:
- Patterns/abstractions: ✔ splits recovery into a locked core with sync and async front doors, shares the setup between them in each crate, and extracts
run_blockingfrom the joincompare_and_swapalready had; no parallel construct. - Performance (speed first, memory second): ✔ startup-only.
compare_and_swapkeeps its single blocking hop with the same panic-resume, andrecover_wal_asynctrades a parked runtime worker for one blocking-pool hop per root at startup. No performance-degradation risk. - Deployment targets: ✔
recover_wal_asyncandrun_blockinglive in thenative, non-wasm32file.rs, the connection crate's file arms carry the same gate, and CI's wasm32 clippy command is clean locally.fluree-db-serverbuilds throughbuild_client, so it picks this up as is; the standalone solo binary opens its engine withbuild_client/build_client_with_nameservicetoo, so it gets async recovery at its next db pin with no code change there, and solo's Lambdas use S3. (Small behavior note: file-backedconnect_async/build_clientnow need a Tokio runtime context to recover, which everyFileStorageI/O call already needed.) - Testing: ✔ six new tests across core, connection and api, all run by name here; the runtime test goes red with a
block_onswapped back in, and the async connection andbuild_clienttests go red without their recovery calls. - Conventions:
⚠️ commits are subject-only (the body carries the rationale); fmt and clippy clean; docs updated, with the stragglers inline.
Verified locally at 6c336279f — CI hasn't run on this PR since it's stacked, so I ran what it would have for these crates: cargo fmt --check; clippy -D warnings --all-targets on core, connection, api and cli, with and without aws; CI's wasm32 clippy command; the fluree-db-core storage tests (137), the two connection tests with and without aws, and the two grp_misc startup-recovery tests — all green, with the new tests present by name.
Before merging. Since this is stacked on #2003's branch, CI hasn't run on it beyond the release workflow's plan job: no test or clippy run exists for 6c336279f. It's also based on #2003 as of ab3fe3d7f, before your two newer commits there; merging fix/cas-deadlock in again is clean, and on that merge all 138 core storage tests pass, with #2003's new panic test now pinning run_blocking. Once #2003 merges, I'd retarget this to main and push, because a retarget alone may not start CI (it arrives as an edited event, which isn't among the default activity types ci.yml's pull_request trigger uses), so those jobs actually run on this head before it lands. I also checked it against the other open PRs and there's nothing to coordinate, except one line against #2017 in fluree-db-api/src/lib.rs: your build_file refactor drops the Ok( wrapper, and #2017 renames attachment_provider_cell to ledger_manager_cell. That's trivial whichever of the two lands second.
Approving so you can merge when ready, but maybe worth folding in the marker assertion and the doc stragglers first — they're small, and I'd rather see them in this PR than lost in the backlog.
| log: Option<Arc<Wal>>, | ||
| durability: Durability, | ||
| _root_gate: tokio::sync::OwnedRwLockReadGuard<()>, | ||
| _not_send: std::marker::PhantomData<*const ()>, |
There was a problem hiding this comment.
🟠 Should address — nothing pins this marker, so the rule the body says the compiler now enforces can be refactored away without a red build.
I confirmed the marker does what the body says: putting the pre-#2001 shape back (returning locked_read_blocking's result out of a spawn_blocking) fails with "*const () cannot be sent between threads safely", and so does holding an OperationHold across .await in a Send future. But it all hangs on this one field. I deleted _not_send and its initializer, and fluree-db-core and its tests compiled as before — so a cleanup pass that drops an unused-looking underscore field, or a later split of this struct, would quietly let the #2001 shape compile again.
The workspace doesn't depend on static_assertions, but its assert_not_impl_any! expands to something small enough to inline. I added this next to the struct: it compiles at this head, and with the marker removed it fails with E0283 ("type annotations needed"):
// Stops compiling if `OperationHold` becomes `Send`: both impls would then apply.
const _: fn() = || {
trait AmbiguousIfSend<A> {
fn some_item() {}
}
impl<T: ?Sized> AmbiguousIfSend<()> for T {}
impl<T: ?Sized + Send> AmbiguousIfSend<u8> for T {}
let _ = <OperationHold as AmbiguousIfSend<_>>::some_item;
};Smaller, fwiw: !Send can't see a non-Send future. An OperationHold held across .await inside a LocalSet task (the kind of context #2001's test ran in) still compiles — I checked. So "Its release never waits on an async task to be polled" holds because every holder lives in this file on a blocking thread, not because of the type. It may be worth the doc comment saying so, so whoever adds the next holder knows the type only covers half of it.
Since the body leans on this as the enforcement, I'd rather see it pinned in this PR than lost in the backlog.
There was a problem hiding this comment.
Right on both counts. I added the assertion and updated the docstring in e1ac529
| for (identifier, storage_config) in addr_ids { | ||
| let id_storage: Arc<dyn Storage> = match &storage_config.storage_type { | ||
| #[cfg(feature = "aws")] | ||
| StorageType::S3(_) => build_s3_storage_from_config(storage_config).await?, |
There was a problem hiding this comment.
❓ Question — is the S3 arm for file- and memory-backed clients a deliberate change?
Before this PR, build_client_file and build_client_memory went through the local wrap_address_identifiers, which sent every identifier to build_local_storage_from_config — so an S3 address identifier on a file- or memory-backed client failed with "S3 storage in addressIdentifiers is only supported with 'aws' feature", even with aws on. With the two wrappers merged, an aws build now opens an S3 storage for that identifier instead.
I think that's what the old error message implied should happen, so this may well be intended. But the body only describes the no-aws side ("gets the same configuration error as before"), and I don't see a test with a local base and an S3 identifier. If it's deliberate, a sentence in the body and a small test would make it visible; if not, the local paths probably want to keep rejecting it.
There was a problem hiding this comment.
It was not purely intentional in that this wasn't a design goal from the beginning, but the constraint was only a byproduct of the old implementation. I chose to leave the behavior change in place because restricting it now would just be artificial. This relaxes a requirement, so nothing that previously worked stopped working. I also chose not to document for users or test it so we don't have to commit to this behavior until we have a specific use case and the extra information that brings.
| use fluree_db_indexer::IndexerConfig; | ||
|
|
||
| let fluree = FlureeBuilder::file("/path/to/data").build().await?; | ||
| let fluree = FlureeBuilder::file("/path/to/data").build_async().await?; |
There was a problem hiding this comment.
🟡 Optional — the same build().await shape is in two more pages, and two file-backed build()? examples were missed.
docs/graph-sources/overview.md:193 and docs/graph-sources/r2rml.md:29 both have let fluree = FlureeBuilder::default().build().await?;. Like the reindex examples, .await on a Result doesn't compile — and default() has no storage path, so even a sync build() there returns "File storage requires a path". FlureeBuilder::file("/path/to/data").build_async().await? would match this page.
Two more still call build()? from what reads as async code: the without_ledger_caching() example at docs/getting-started/rust-api.md:772-774, and the remote-connection example at docs/operations/configuration.md:1129-1132 (the same example this PR converted at rust-api.md:1402-1408).
Four small edits, so if you agree, I'd rather see them in this PR than lost in the backlog.
Commenting here because those files are not in this diff.
| ``` | ||
|
|
||
| `build_async()` replays the storage's write-ahead log without blocking the async runtime. | ||
| `build()` builds the same instance from synchronous code, replaying the log on the calling thread. |
There was a problem hiding this comment.
🟡 Optional — "from synchronous code" can read as "without a runtime".
FlureeBuilder's own docs (fluree-db-api/src/lib.rs:1464-1466) say every build* method, the sync ones included, must be called inside a Tokio runtime, since the default ledger cache spawns its listener. I called FlureeBuilder::file(dir).build() from a plain thread to check, and it panics with "there is no reactor running, must be called from the context of a Tokio 1.x runtime". Someone reading this line could reasonably try it from a plain fn main. Maybe something like:
| `build()` builds the same instance from synchronous code, replaying the log on the calling thread. | |
| `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. |
It's a one-line edit to a line this PR adds, so I'd rather it land here.
Follow-up: #2001. Builds on #2003.
Summary
FileStorage::recover_walreplays a root's write-ahead log at startup. It takes the root gate withfutures::executor::block_onand replays with synchronous file I/O, all on the calling thread. Async startup paths called it from a runtime thread, so that thread stalled until every in-flight file operation on the root finished, and then for the whole replay.This PR adds an async recovery path and uses it wherever recovery runs in async code. It also makes the compiler enforce part of the rule the #2001 fix depends on: only blocking-pool threads hold the root gate.
Enforcing the root-gate rule
OperationHoldholds a read guard on the root gate. It is now!Send. Returning one from aspawn_blockingtask, or holding one across an.awaitin aSendfuture, no longer compiles. Both are ways an async task could hold the root gate whilerecover_walblocks the thread that task needs. The shape of the #2001 bug, a guard returned through aspawn_blockingresult, fails with "*const ()cannot be sent between threads safely".A static assertion next to the struct pins the marker. If
OperationHoldbecomesSend, for example because the marker field is removed, the crate fails to compile with E0283.The type does not cover a
!Sendfuture, such as aLocalSettask, holding one across an.await. The struct's docstring says so, and requires holders to take and release it on a blocking thread, which every current holder does.Async recovery
recover_walis split.recover_wal_lockeddoes the replay and takes the root gate's write guard as a parameter, so it cannot run without the gate held. The syncrecover_waltakes the gate withblock_onand calls it.FileStorage::recover_wal_asyncawaits the root gate and moves the write guard into a blocking task, which replays throughrecover_wal_lockedand then releases the gate. No async task holds the gate. Dropping the future before it takes the gate replays nothing. Dropping it afterward leaves the replay to finish on the blocking pool.run_blockingruns a closure on the blocking pool, resumes a panic on the caller, and returns only cancellation as an error.recover_wal_asyncandcompare_and_swapboth use it.recover_wal's docstring now states that it blocks the calling thread until in-flight file operations on the root finish, and that the replay runs on that thread.Call sites
Each call site uses the mechanism that fits its context.
create_sync_connectionconnectrecover_walcreate_async_connection, with and withoutawsrecover_wal_asyncFlureeBuilder::build_clientandbuild_client_with_nameservice, used by the serverrecover_wal_async, for the main root and file-backed address identifiersFlureeBuilder::build,build_encrypted,build_encrypted_from_config,fluree_filerecover_walFlureeBuilder::build_async(new)recover_wal_asyncbuild_flureebuild_asyncfluree-bench-virtualopenbuildWithout the
awsfeature,create_async_connectionused to fall back to the sync path for file storage. It now recovers asynchronously too. The connection crate's two copies of the file arm areopen_file_storage,file_connection, andfile_connection_async.File-backed
connect_asyncandbuild_clientnow need a Tokio runtime context to recover, sincerecover_wal_asyncusesspawn_blocking.FileStorage's async I/O already needed one.In the API crate, the file client, the memory client, and
build_local_storage_from_configare async. The local and AWS address-identifier wrappers are one asyncwrap_address_identifiers, with its S3 arm behind theawsfeature. Withoutaws, an S3 identifier gets the same configuration error as before. Withaws, file- and memory-backed clients now open S3 storage for an S3 identifier, as S3-backed clients do. They used to reject it, because the local wrapper sent every identifier tobuild_local_storage_from_config. This is deliberate but not documented for users, so it is not a commitment yet.FlureeBuilder::build_asyncbuilds the same file-backed instance asbuild. Both areopen_file_storage, then recovery, thenbuild_file.The CLI's
build_flureeis async and awaitsbuild_async. Its 45 call sites gained an.await, andbuild_store, its one sync caller, is async too.Docs
Examples that build a file-backed instance inside async code now use
build_async().await. This covers the Rust API guide, the BM25, running, benches, encryption, and configuration docs, and the crate-level Quick Start.Four examples did not compile: the two reindex examples called
build().await, and the graph-source overview and R2RML pages calledFlureeBuilder::default().build().await, which also has no storage path. They now build a file-backed instance withbuild_async().await.The Rust API guide notes when to use
build_asyncandbuild, including thatbuildneeds a Tokio runtime even from synchronous code. The encryption guide listsbuild_asyncamong the build methods that apply a configured key. Thebuild_encrypted*examples stay sync, since those methods have no async counterpart.API changes
New public API:
FileStorage::recover_wal_asyncandFlureeBuilder::build_async.FileStorage::crash_with_a_logged_write_for_testis a new#[doc(hidden)]test hook beside the existing ones.One public signature changed:
fluree_db_cli::context::build_flureeis nowasync. Nothing outside the CLI calls it.Testing
New tests, each shown to fail with its protection removed:
recover_wal_async_waits_for_the_root_gate_without_blocking_the_runtime(core). Another thread holds an operation on the root, and a timer on the same current-thread runtime must fire while recovery waits. With the syncrecover_walin its place, it times out.recover_wal_async_restores_acknowledged_writes_after_a_crash(core). It replays a crashed log and checks the rebuilt files on disk, sinceread_bytesreplays on a miss by itself. Without the recovery call, it fails at the first disk read.file_connection_recovers_the_walandasync_file_connection_recovers_the_wal(connection). Each fails without its connection path's recovery call. Both pass with and withoutaws, which covers bothcreate_async_connectionvariants.build_async_recovers_the_walandbuild_client_recovers_the_wal(api,it_file_startup_recovery). Each fails without its build path's recovery call.Restoring the pre-#2001 compare-and-swap shape against the
!Sendmarker fails to compile, and removing the marker fails the static assertion.