From a374ac5fc44dd835c13ce51e63685ee352ff57cb Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Sat, 26 Sep 2026 22:26:37 +0000 Subject: [PATCH] perf(graph-db): stream the sealed store build one column at a time --- .cargo/config.toml | 4 +- Cargo.lock | 10 +-- Cargo.toml | 25 ++++--- crates/tracedecay-graph-db/Cargo.toml | 4 +- crates/tracedecay-graph-db/src/lib.rs | 2 + .../tracedecay-graph-db/src/sealed_store.rs | 66 +++++++++++++--- .../src/thread_allocation.rs | 75 +++++++++++++++++++ 7 files changed, 155 insertions(+), 31 deletions(-) create mode 100644 crates/tracedecay-graph-db/src/thread_allocation.rs diff --git a/.cargo/config.toml b/.cargo/config.toml index eb2d621eb6..4186abb838 100644 --- a/.cargo/config.toml +++ b/.cargo/config.toml @@ -134,9 +134,9 @@ directory = ".pnpm/crates/crates-io" [source.pnpm-git] directory = ".pnpm/crates/git" -[source."git+https://github.com/ScriptedAlchemy/grafeo.git?rev=4dc04db299bb63269a2c216c66aa94fab40d9819"] +[source."git+https://github.com/ScriptedAlchemy/grafeo.git?rev=2027d1093c19d9d631618b9b4d1e575120c12fb2"] git = "https://github.com/ScriptedAlchemy/grafeo.git" -rev = "4dc04db299bb63269a2c216c66aa94fab40d9819" +rev = "2027d1093c19d9d631618b9b4d1e575120c12fb2" replace-with = "pnpm-git" [source."git+https://github.com/ScriptedAlchemy/tokensave-large-treesitters?rev=102d0103b6fd093c1e5989cceede1ed9f08fe6a9"] diff --git a/Cargo.lock b/Cargo.lock index 8bac0a425a..044dc16fce 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2580,7 +2580,7 @@ dependencies = [ [[package]] name = "grafeo-adapters" version = "0.5.42" -source = "git+https://github.com/ScriptedAlchemy/grafeo.git?rev=4dc04db299bb63269a2c216c66aa94fab40d9819#4dc04db299bb63269a2c216c66aa94fab40d9819" +source = "git+https://github.com/ScriptedAlchemy/grafeo.git?rev=2027d1093c19d9d631618b9b4d1e575120c12fb2#2027d1093c19d9d631618b9b4d1e575120c12fb2" dependencies = [ "bincode", "grafeo-common", @@ -2595,7 +2595,7 @@ dependencies = [ [[package]] name = "grafeo-common" version = "0.5.42" -source = "git+https://github.com/ScriptedAlchemy/grafeo.git?rev=4dc04db299bb63269a2c216c66aa94fab40d9819#4dc04db299bb63269a2c216c66aa94fab40d9819" +source = "git+https://github.com/ScriptedAlchemy/grafeo.git?rev=2027d1093c19d9d631618b9b4d1e575120c12fb2#2027d1093c19d9d631618b9b4d1e575120c12fb2" dependencies = [ "arcstr", "bincode", @@ -2615,7 +2615,7 @@ dependencies = [ [[package]] name = "grafeo-core" version = "0.5.42" -source = "git+https://github.com/ScriptedAlchemy/grafeo.git?rev=4dc04db299bb63269a2c216c66aa94fab40d9819#4dc04db299bb63269a2c216c66aa94fab40d9819" +source = "git+https://github.com/ScriptedAlchemy/grafeo.git?rev=2027d1093c19d9d631618b9b4d1e575120c12fb2#2027d1093c19d9d631618b9b4d1e575120c12fb2" dependencies = [ "arc-swap", "arcstr", @@ -2638,7 +2638,7 @@ dependencies = [ [[package]] name = "grafeo-engine" version = "0.5.42" -source = "git+https://github.com/ScriptedAlchemy/grafeo.git?rev=4dc04db299bb63269a2c216c66aa94fab40d9819#4dc04db299bb63269a2c216c66aa94fab40d9819" +source = "git+https://github.com/ScriptedAlchemy/grafeo.git?rev=2027d1093c19d9d631618b9b4d1e575120c12fb2#2027d1093c19d9d631618b9b4d1e575120c12fb2" dependencies = [ "arcstr", "bincode", @@ -2660,7 +2660,7 @@ dependencies = [ [[package]] name = "grafeo-storage" version = "0.5.42" -source = "git+https://github.com/ScriptedAlchemy/grafeo.git?rev=4dc04db299bb63269a2c216c66aa94fab40d9819#4dc04db299bb63269a2c216c66aa94fab40d9819" +source = "git+https://github.com/ScriptedAlchemy/grafeo.git?rev=2027d1093c19d9d631618b9b4d1e575120c12fb2#2027d1093c19d9d631618b9b4d1e575120c12fb2" dependencies = [ "bincode", "byteorder", diff --git a/Cargo.toml b/Cargo.toml index 7b8a3fa279..b4d3c31dc4 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -239,14 +239,17 @@ tokensave-large-treesitters = { git = "https://github.com/ScriptedAlchemy/tokens # exhaustion errors. The fork also carries graph correctness and bounded WAL # group-write fixes, and (ScriptedAlchemy/grafeo#4) the direct compact build: # `IncrementalCompactStoreBuilder` + `GrafeoDB::write_compact_container` -# turn a verified row stream into a compacted container with one columnar -# copy in memory and one container write, with no staging LPG, WAL, or -# in-memory `compact()` copy, which is how sealed generation stores are -# built. Its workspace resolves sibling crates by path, so every -# grafeo crate must come from the same revision. A mixed registry/git graph -# would build type-incompatible instances of the patched crates. -grafeo-adapters = { git = "https://github.com/ScriptedAlchemy/grafeo.git", rev = "4dc04db299bb63269a2c216c66aa94fab40d9819" } -grafeo-common = { git = "https://github.com/ScriptedAlchemy/grafeo.git", rev = "4dc04db299bb63269a2c216c66aa94fab40d9819" } -grafeo-core = { git = "https://github.com/ScriptedAlchemy/grafeo.git", rev = "4dc04db299bb63269a2c216c66aa94fab40d9819" } -grafeo-engine = { git = "https://github.com/ScriptedAlchemy/grafeo.git", rev = "4dc04db299bb63269a2c216c66aa94fab40d9819" } -grafeo-storage = { git = "https://github.com/ScriptedAlchemy/grafeo.git", rev = "4dc04db299bb63269a2c216c66aa94fab40d9819" } +# turn a verified row stream into a compacted container with one container +# write, with no staging LPG, WAL, or in-memory `compact()` copy, which is +# how sealed generation stores are built. Since ScriptedAlchemy/grafeo#5 the +# builder spools pushed values to a file and the container write encodes +# them one column at a time, so the build never holds every row's values +# or the finished columnar store. Its workspace resolves sibling crates by +# path, so every grafeo crate must come from the same revision. A mixed +# registry/git graph would build type-incompatible instances of the patched +# crates. +grafeo-adapters = { git = "https://github.com/ScriptedAlchemy/grafeo.git", rev = "2027d1093c19d9d631618b9b4d1e575120c12fb2" } +grafeo-common = { git = "https://github.com/ScriptedAlchemy/grafeo.git", rev = "2027d1093c19d9d631618b9b4d1e575120c12fb2" } +grafeo-core = { git = "https://github.com/ScriptedAlchemy/grafeo.git", rev = "2027d1093c19d9d631618b9b4d1e575120c12fb2" } +grafeo-engine = { git = "https://github.com/ScriptedAlchemy/grafeo.git", rev = "2027d1093c19d9d631618b9b4d1e575120c12fb2" } +grafeo-storage = { git = "https://github.com/ScriptedAlchemy/grafeo.git", rev = "2027d1093c19d9d631618b9b4d1e575120c12fb2" } diff --git a/crates/tracedecay-graph-db/Cargo.toml b/crates/tracedecay-graph-db/Cargo.toml index 1636ce47ee..b260d152da 100644 --- a/crates/tracedecay-graph-db/Cargo.toml +++ b/crates/tracedecay-graph-db/Cargo.toml @@ -82,11 +82,13 @@ tracedecay-store = { path = "../tracedecay-store" } # is least able to notice from the inside: it is silent, correct, and # slow. The fallback paths say so out loud. tracing = "0.1" +# The sealed build's column spool: an anonymous file in the staging +# directory that the OS removes even when the process is killed. +tempfile = "3" [dev-dependencies] criterion = "0.5" rusqlite = { version = "0.40.1", default-features = false, features = ["backup"] } -tempfile = "3" tracedecay-rusqlite-runtime = { path = "../tracedecay-rusqlite-runtime" } [[bench]] diff --git a/crates/tracedecay-graph-db/src/lib.rs b/crates/tracedecay-graph-db/src/lib.rs index b91718386f..a2ae41aac9 100644 --- a/crates/tracedecay-graph-db/src/lib.rs +++ b/crates/tracedecay-graph-db/src/lib.rs @@ -22,6 +22,8 @@ mod runtime; mod schema; mod sealed_store; mod state; +#[cfg(test)] +mod thread_allocation; mod traversal; mod verified_marker; diff --git a/crates/tracedecay-graph-db/src/sealed_store.rs b/crates/tracedecay-graph-db/src/sealed_store.rs index 98d38f76f1..81edf686f4 100644 --- a/crates/tracedecay-graph-db/src/sealed_store.rs +++ b/crates/tracedecay-graph-db/src/sealed_store.rs @@ -647,12 +647,16 @@ struct PreparedSealedRelation { } impl SealedCompactRows { - fn new() -> Self { - Self { - builder: IncrementalCompactStoreBuilder::new(), + /// Spools pushed column values to an anonymous file in `staging`, so + /// the build holds the row topology rather than every value. + fn new(staging: &Path) -> Result { + let spool = tempfile::tempfile_in(staging) + .map_err(|error| sealed_store_io_failure("column spool", error))?; + Ok(Self { + builder: IncrementalCompactStoreBuilder::spooling_to(spool), next_node: 0, next_edge: 0, - } + }) } fn push_node( @@ -800,17 +804,16 @@ impl SealedCompactRows { } /// Encodes every pushed row and writes the sealed container at `path` - /// in one durable pass. The container holds the compact store, an empty - /// LPG overlay whose id allocators start past the written ids, and a - /// catalog naming the unique-key property indexes. + /// in one durable pass, one column at a time. The container holds the + /// compact store, an empty LPG overlay whose id allocators start past + /// the written ids, and a catalog naming the unique-key property + /// indexes. fn write_container(self, path: &Path) -> Result<(), GraphDbError> { - let store = hotpath::measure_block!("code_index.seal.write.compact", self.builder.finish()) - .map_err(|error| sealed_build_failure("encode", error))?; hotpath::measure_block!( "code_index.seal.write.container", GrafeoDB::write_compact_container( path, - Arc::new(store), + self.builder, INDEXED_PROPERTIES .iter() .map(|property| (*property).to_owned()), @@ -1508,7 +1511,7 @@ fn build_sealed_container( ) -> Result<(usize, usize), GraphDbError> { hotpath::gauge!("code_index.seal.encode.effective_workers").set(1); let physical_namespace = identity.physical_namespace()?; - let mut sealed = SealedCompactRows::new(); + let mut sealed = SealedCompactRows::new(staging)?; let (entity_count, relation_count, dependency_namespaces_written) = hotpath::measure_block!("code_index.seal.encode", { match rows { @@ -2140,7 +2143,7 @@ mod build_tests { use super::{ SEALED_STORE_DATABASE_FILE, SealedRowSource, build_or_open_sealed_store, - sealed_generation_directory, sealed_store_root, + build_sealed_container, sealed_generation_directory, sealed_store_root, }; use crate::{ GraphDbError, GraphDbLocation, GraphDbOpenOptions, GraphDbOwner, GraphDurability, @@ -2218,6 +2221,45 @@ mod build_tests { owner.issue_lease().unwrap() } + /// The direct build holds the rows' topology and one column at a time, + /// not every pushed value: its peak above the resident manifest stays + /// within a fixed budget for a 30,000-entity, 45,000-relation generation, + /// and the column spool it used is gone once the container is written. + /// Holding every value until the columnar store was encoded peaked at + /// 89,678,855 bytes on this generation. + #[test] + fn direct_build_holds_one_column_not_every_value() { + const PEAK_BUDGET_BYTES: usize = 30_000_000; + let check: &dyn Fn() -> Result<(), GraphDbError> = &|| Ok(()); + let manifest = manifest(30_000, 45_000); + let identity = manifest.identity(); + let staging = tempfile::tempdir().unwrap(); + + let (built, peak) = crate::thread_allocation::peak_above_start(|| { + build_sealed_container( + SealedRowSource::Manifest(&manifest), + &identity, + staging.path(), + check, + ) + }); + eprintln!("SEALED BUILD peak {peak}"); + + assert_eq!(built.unwrap(), (30_000, 45_000)); + assert!( + peak <= PEAK_BUDGET_BYTES, + "the sealed build held {peak} bytes at peak, over its {PEAK_BUDGET_BYTES}-byte budget" + ); + let left: Vec<_> = std::fs::read_dir(staging.path()) + .unwrap() + .map(|entry| entry.unwrap().file_name()) + .collect(); + assert_eq!( + left, + vec![std::ffi::OsString::from(SEALED_STORE_DATABASE_FILE)] + ); + } + /// The sealed container is written once, complete, after every row has /// been encoded: no engine, WAL, or partial container exists under the /// staging directory while rows stream. A build cancelled mid-copy leaves diff --git a/crates/tracedecay-graph-db/src/thread_allocation.rs b/crates/tracedecay-graph-db/src/thread_allocation.rs new file mode 100644 index 0000000000..ff50e19377 --- /dev/null +++ b/crates/tracedecay-graph-db/src/thread_allocation.rs @@ -0,0 +1,75 @@ +//! Per-thread allocation accounting for the lib test binary. +//! +//! Tests run concurrently in one process, so a process-wide counter would +//! charge one test for another's allocations. Each thread counts its own +//! live bytes instead; a measurement is only meaningful for work that runs +//! entirely on the measuring thread. + +use std::alloc::{GlobalAlloc, Layout, System}; +use std::cell::Cell; + +thread_local! { + static LIVE: Cell = const { Cell::new(0) }; + static PEAK: Cell = const { Cell::new(0) }; +} + +struct ThreadCountingAllocator; + +fn charge(bytes: isize) { + let _ = LIVE.try_with(|live| { + let now = live.get() + bytes; + live.set(now); + let _ = PEAK.try_with(|peak| peak.set(peak.get().max(now))); + }); +} + +// SAFETY: every call delegates to `System` with the caller's layout; the +// accounting touches only const-initialized thread-locals, which never +// allocate. +unsafe impl GlobalAlloc for ThreadCountingAllocator { + unsafe fn alloc(&self, layout: Layout) -> *mut u8 { + // SAFETY: forwarded unchanged. + let pointer = unsafe { System.alloc(layout) }; + if !pointer.is_null() { + charge(layout.size().cast_signed()); + } + pointer + } + + unsafe fn alloc_zeroed(&self, layout: Layout) -> *mut u8 { + // SAFETY: forwarded unchanged. + let pointer = unsafe { System.alloc_zeroed(layout) }; + if !pointer.is_null() { + charge(layout.size().cast_signed()); + } + pointer + } + + unsafe fn dealloc(&self, pointer: *mut u8, layout: Layout) { + // SAFETY: forwarded unchanged. + unsafe { System.dealloc(pointer, layout) }; + charge(-layout.size().cast_signed()); + } + + unsafe fn realloc(&self, pointer: *mut u8, layout: Layout, new_size: usize) -> *mut u8 { + // SAFETY: forwarded unchanged. + let resized = unsafe { System.realloc(pointer, layout, new_size) }; + if !resized.is_null() { + charge(new_size.cast_signed() - layout.size().cast_signed()); + } + resized + } +} + +#[global_allocator] +static ALLOCATOR: ThreadCountingAllocator = ThreadCountingAllocator; + +/// Runs `work` and returns its result with the most bytes this thread held +/// above what it held when `work` started. +pub(crate) fn peak_above_start(work: impl FnOnce() -> R) -> (R, usize) { + let start = LIVE.with(Cell::get); + PEAK.with(|peak| peak.set(start)); + let result = work(); + let peak = PEAK.with(Cell::get) - start; + (result, usize::try_from(peak).unwrap_or(0)) +}