Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<BTreeMap<ContentDigest, Weak<[u8]>>>,
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 {
Expand All @@ -91,7 +98,6 @@ impl Default for SharedCodeIndexBytePoolV1 {
),
physical_artifacts: SharedPhysicalCodeArtifactPoolV1::default(),
decoded_content: SharedDecodedContentPoolV1::default(),
last_prune_len: AtomicUsize::new(0),
}
}
}
Expand All @@ -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(),
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(&registry, 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<CodeIndexCadenceTelemetryV1>,
) -> CodeIndexCadenceTriggerV1 {
Expand Down
19 changes: 19 additions & 0 deletions crates/tracedecay-code-index/src/production/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
40 changes: 39 additions & 1 deletion crates/tracedecay-code-index/tests/resident_accounting.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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);
Expand Down
Loading