From 3a3cd1d6b48eecade9b8a4f2342fd72b0eafb22b Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Thu, 1 Oct 2026 00:49:07 +0000 Subject: [PATCH] fix(code-index): free captured sources once a pass lets them go The byte pool and the physical artifact pool index through Weak, and a Weak keeps its allocation, so every source buffer and artifact header a pass dropped stayed allocated. The worker now drops dead entries once a pass lets go of its capture and build. --- .../code_index_scheduler/publication_store.rs | 51 ++++++++++++------- .../code_index_scheduler/registry/mount.rs | 2 + .../code_index_scheduler/tests/residency.rs | 26 ++++++++++ .../src/production/mod.rs | 19 +++++++ .../tests/resident_accounting.rs | 40 ++++++++++++++- 5 files changed, 120 insertions(+), 18 deletions(-) diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/publication_store.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/publication_store.rs index 46adb57d4e..4fa64fa4ac 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/publication_store.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/publication_store.rs @@ -8,11 +8,14 @@ use std::{ path::{Path, PathBuf}, sync::{ Arc, Condvar, Mutex, MutexGuard, PoisonError, Weak, - atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering}, + atomic::{AtomicBool, AtomicU64, Ordering}, }, time::{Duration, SystemTime}, }; +#[cfg(any(test, feature = "hotpath"))] +use std::sync::atomic::AtomicUsize; + use same_file::Handle; use sha2::{Digest, Sha256}; use tracedecay_application::code_index::DaemonCodeIndexControlV1; @@ -67,19 +70,23 @@ static CODE_INDEX_GENERATION_DECODE_WAITERS: AtomicUsize = AtomicUsize::new(0); pub struct CodeIndexBytePoolStatsV1 { pub parse_chunk_inserted: u64, pub parse_chunk_reused: u64, + /// Source allocations the pool keeps, live or not. + pub source_allocations: usize, + /// Those some capture or generation still holds. + pub live_sources: usize, } +/// Captured sources by content, shared by every capture in the process. +/// +/// Entries are weak, and a `Weak<[u8]>` keeps its whole allocation, bytes +/// included, after the last `Arc` dropped. Each pass therefore drops the +/// entries its capture let go of, see [`Self::release_dead_entries`]. pub struct SharedCodeIndexBytePoolV1 { bytes: ProfiledStdMutex>>, pub(super) physical_artifacts: SharedPhysicalCodeArtifactPoolV1, /// Decoded generation pages by content, so linked worktrees that sealed /// identical trees hold one decode between them. pub(super) decoded_content: SharedDecodedContentPoolV1, - /// Map length recorded after the last dead-entry prune. Weak entries whose - /// `Arc` dropped are never removed by lookups, so `intern` prunes them once - /// the map doubles past this baseline, bounding growth over the daemon - /// lifetime at amortized O(1) per insert. - last_prune_len: AtomicUsize, } impl Default for SharedCodeIndexBytePoolV1 { @@ -91,7 +98,6 @@ impl Default for SharedCodeIndexBytePoolV1 { ), physical_artifacts: SharedPhysicalCodeArtifactPoolV1::default(), decoded_content: SharedDecodedContentPoolV1::default(), - last_prune_len: AtomicUsize::new(0), } } } @@ -114,25 +120,36 @@ impl SharedCodeIndexBytePoolV1 { } let shared: Arc<[u8]> = Arc::from(bytes); pool.insert(digest.clone(), Arc::downgrade(&shared)); - if pool.len() - > self - .last_prune_len - .load(Ordering::Relaxed) - .saturating_mul(2) - { - pool.retain(|_, entry| entry.strong_count() > 0); - self.last_prune_len - .store(pool.len().max(1), Ordering::Relaxed); - } (digest, shared) } + /// Drop the entries no capture or generation holds any more, which frees + /// their sources, and the physical artifact pool's likewise. A dead entry + /// can never be reused, so this gives up nothing. Runs once a pass let go + /// of its capture and build. + pub(super) fn release_dead_entries(&self) { + self.bytes + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .retain(|_, entry| entry.strong_count() > 0); + self.physical_artifacts.release_dead_entries(); + } + #[cfg(test)] pub(super) fn stats(&self) -> CodeIndexBytePoolStatsV1 { let physical_artifacts = self.physical_artifacts.stats(); + let bytes = self + .bytes + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); CodeIndexBytePoolStatsV1 { parse_chunk_inserted: physical_artifacts.inserted, parse_chunk_reused: physical_artifacts.reused, + source_allocations: bytes.len(), + live_sources: bytes + .values() + .filter(|entry| entry.strong_count() > 0) + .count(), } } } diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/mount.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/mount.rs index 73827faa7f..d8de55afee 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/mount.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/mount.rs @@ -641,6 +641,7 @@ impl CodeIndexSchedulerRegistryV1 { )); let worker_residency = Arc::clone(&residency); let worker_resident_owners = Arc::clone(&self.resident_owners); + let worker_byte_pool = Arc::clone(&self.byte_pool); let mut worker_owner_headroom = self.resident_owners.subscribe_headroom(); let mut worker_admission_headroom = self.resident_memory.pressure().subscribe_headroom(); // Boxed at definition on purpose: this worker's state machine is the @@ -2922,6 +2923,7 @@ impl CodeIndexSchedulerRegistryV1 { ); } drop(reconcile_pass.take()); + worker_byte_pool.release_dead_entries(); let refused_for_memory = matches!( &result, Ok((Err(error), _, _)) if error.is_resident_memory_refusal() diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/residency.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/residency.rs index 7aff943542..e945387f61 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/residency.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/residency.rs @@ -393,6 +393,32 @@ async fn an_increment_reports_its_retained_parses_and_the_idle_window_releases_t registry.shutdown().await; } +/// A pass lets go of the sources it captured once its build is sealed, and +/// the shared byte pool then frees them instead of keeping each one +/// allocated behind a dead weak entry. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_settled_index_keeps_no_captured_source_allocated() { + let fixture = GitFixture::new(&[ + ("src/main.rs", "fn main() { helper(); }\n"), + ("src/helper.rs", "pub fn helper() {}\n"), + ("src/other.rs", "pub fn other() {}\n"), + ]); + let store = TempDir::new().expect("store root"); + let (registry, _) = + mounted_core_query_worktree_in(CodeIndexSchedulerRegistryV1::new(1), &fixture, &store) + .await; + wait_for_settled_owner(®istry, fixture.path()).await; + + let stats = registry.byte_pool_stats(); + assert_eq!( + (stats.source_allocations, stats.live_sources), + (0, 0), + "{stats:?}" + ); + + registry.shutdown().await; +} + async fn next_receipt_trigger( receipts: &mut tokio::sync::watch::Receiver, ) -> CodeIndexCadenceTriggerV1 { diff --git a/crates/tracedecay-code-index/src/production/mod.rs b/crates/tracedecay-code-index/src/production/mod.rs index 8bd63adb16..169e803eb0 100644 --- a/crates/tracedecay-code-index/src/production/mod.rs +++ b/crates/tracedecay-code-index/src/production/mod.rs @@ -598,6 +598,25 @@ impl SharedPhysicalCodeArtifactPoolV1 { }) } + /// Drop the index entries whose artifact no generation owns any more. A + /// `Weak` keeps its allocation, so a dead entry still pins the artifact's + /// header, and it can never be reused. + pub fn release_dead_entries(&self) { + let mut state = self + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + state + .artifacts + .retain(|_, artifact| artifact.strong_count() > 0); + let PhysicalCodeArtifactPoolStateV1 { + artifacts, + insertion_order, + .. + } = &mut *state; + insertion_order.retain(|key| artifacts.contains_key(key)); + } + fn record_clone_payloads(&self, reused: u64, computed: u64) { let mut state = self .state diff --git a/crates/tracedecay-code-index/tests/resident_accounting.rs b/crates/tracedecay-code-index/tests/resident_accounting.rs index ff4a19f418..69d36b851c 100644 --- a/crates/tracedecay-code-index/tests/resident_accounting.rs +++ b/crates/tracedecay-code-index/tests/resident_accounting.rs @@ -20,7 +20,7 @@ use tracedecay_code_index::production::{ CodeIndexProductionErrorV1, CodeIndexProductionOwnerV1, CodeIndexPublicationStoreErrorV1, CodeIndexPublishedGenerationV1, CodeIndexRepositoryParseIdentityV1, SealedGenerationFileWindowsV1, SealedGenerationSegmentPublicationV1, - SealedGenerationSegmentReadV1, SharedDecodedContentPoolV1, + SealedGenerationSegmentReadV1, SharedDecodedContentPoolV1, SharedPhysicalCodeArtifactPoolV1, }; use tracedecay_code_index::projection::{ ChunkProjectionDecisionV1, CodeChunkProjectionSink, ProjectionReceiptBuilderV1, @@ -369,6 +369,44 @@ fn a_full_build_leaves_live_only_what_its_generation_charges() { ); } +/// The daemon keeps one physical artifact pool for its lifetime. A `Weak` +/// keeps its allocation, so every index entry whose generation is gone pins +/// that artifact's allocation until the pool drops the entry. +#[test] +fn a_dropped_generation_leaves_nothing_once_the_pool_drops_its_dead_entries() { + let _measurement = MEASUREMENT.lock().unwrap_or_else(PoisonError::into_inner); + drop( + CodeIndexProductionOwnerV1::new(config(), Publication, Projection) + .expect("owner") + .build_and_publish(request(300), &Active) + .expect("warm build"), + ); + let pool = SharedPhysicalCodeArtifactPoolV1::default(); + let mut owner = CodeIndexProductionOwnerV1::new(config(), Publication, Projection) + .expect("owner") + .with_physical_artifact_pool(pool.clone()); + + let before = LIVE.load(Ordering::Relaxed); + drop( + owner + .build_and_publish(request(300), &Active) + .expect("build"), + ); + let pinned = LIVE.load(Ordering::Relaxed).saturating_sub(before); + pool.release_dead_entries(); + let left = LIVE.load(Ordering::Relaxed).saturating_sub(before); + eprintln!("ACCOUNTING dead entries pinned {pinned} left {left}"); + assert_eq!(pool.stats().resident, 0); + assert!( + pinned >= 300 * 512, + "300 dead entries pin their artifacts' allocations: {pinned} bytes" + ); + assert!( + left <= pinned / 20, + "dropping the dead entries left {left} of {pinned} bytes live" + ); +} + #[test] fn retained_bytes_account_for_what_a_decode_leaves_live() { let _measurement = MEASUREMENT.lock().unwrap_or_else(PoisonError::into_inner);