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 @@ -2,14 +2,14 @@ use std::fs;
use std::num::NonZeroU64;
use std::path::Path;
use std::process::Command;
use std::sync::Arc;
use std::sync::{Arc, Mutex};

use tempfile::TempDir;
use tracedecay_code_index::parallelism::CodeIndexWorkerRuntimeV1;
use tracedecay_domain::{ProjectId, configuration::CodeIndexWorkerSelectionV1};
use tracedecay_runtime_core::resident_memory::{
DEFAULT_PROCESS_RESIDENT_MEMORY_LIMIT_V1, ProcessResidentMemoryV1, ResidentMemoryComponentIdV1,
ResidentMemoryPressureV1,
DEFAULT_PROCESS_RESIDENT_MEMORY_LIMIT_V1, ProcessResidentMemoryV1, ProcessResidentSampleV1,
ResidentMemoryComponentIdV1, ResidentMemoryPressureV1,
};

use crate::code_index::production::CodeIndexPublicationStoreErrorV1;
Expand Down Expand Up @@ -550,3 +550,82 @@ fn a_released_generation_decodes_again_only_once_its_bytes_fit_the_budget() {
"the decode's charge is released once it completes"
);
}

/// Admission counts only what the kernel cannot take back without swapping.
/// Here clean file pages (the mapped sealed container) fill the resident set
/// to one byte under the admission watermark while anonymous memory leaves
/// room, so the decode runs. Once the unreclaimable bytes themselves fill the
/// headroom, the same decode is refused and nothing is decoded.
#[test]
fn a_decode_is_admitted_against_unreclaimable_bytes_not_clean_file_pages() {
const GIB: u64 = 1024 * 1024 * 1024;
let project = fixture();
let store = TempDir::new().expect("store root");
let mut scheduler = CodeIndexWorktreeSchedulerV1::open(
ProjectId::new("project.code-index-decode-unreclaimable").expect("valid project"),
project.path(),
store.path().to_path_buf(),
Arc::new(SharedCodeIndexBytePoolV1::default()),
)
.expect("open scheduler");
let limit = NonZeroU64::new(16 * GIB).expect("limit");
let view = Arc::new(Mutex::new(ProcessResidentSampleV1 {
resident_bytes: GIB,
unreclaimable_bytes: GIB,
}));
let sampled = Arc::clone(&view);
let pressure = Arc::new(ResidentMemoryPressureV1::with_sampler(
limit,
Arc::new(move || Some(*sampled.lock().expect("view"))),
));
let high_watermark = pressure.high_watermark_bytes();
let authority = Arc::new(ProcessResidentMemoryV1::with_pressure(limit, pressure));
scheduler.bind_resident_memory(Arc::clone(&authority));
assert!(matches!(
scheduler.reconcile_now().expect("publish generation"),
CodeIndexReconcileOutcomeV1::Published(_)
));
scheduler
.publication
.release_decoded_active_after_seal()
.expect("release the sealed decode");
let decodes_before = scheduler.sealed_decode_count();

*view.lock().expect("view") = ProcessResidentSampleV1 {
resident_bytes: high_watermark - 1,
unreclaimable_bytes: 2 * GIB,
};
let decoded = scheduler
.latest_complete()
.expect("clean file pages do not refuse a decode that fits in anonymous headroom");
assert_eq!(scheduler.sealed_decode_count(), decodes_before + 1);
assert_eq!(
decoded
.generation
.symbols()
.symbols
.iter()
.map(|symbol| symbol.qualified_name.as_str())
.collect::<Vec<_>>(),
["src/lib.rs::retained_generation"]
);
drop(decoded);
scheduler
.publication
.release_decoded_active_after_seal()
.expect("release the decode again");

*view.lock().expect("view") = ProcessResidentSampleV1 {
resident_bytes: high_watermark - 1,
unreclaimable_bytes: high_watermark - 1,
};
assert!(matches!(
scheduler.publication.load_active_shared(),
Err(CodeIndexPublicationStoreErrorV1::ResidentMemoryRefused(_))
));
assert_eq!(
scheduler.sealed_decode_count(),
decodes_before + 1,
"a decode refused on unreclaimable bytes decodes nothing"
);
}
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ use tracedecay_private_fs::framed_log::{DirectorySyncPolicy, DurableFileBatch};
use tracedecay_runtime_core::resident_memory::{
ProcessResidentMemoryV1, ResidentMemoryComponentIdV1, ResidentMemoryKeyV1,
ResidentMemoryReservationV1, ResidentOwnersV1, log_resident_owner_release_v1,
release_process_allocator_memory_v1, sampled_process_resident_bytes_v1,
release_process_allocator_memory_v1,
};

use crate::code_index::{
Expand Down Expand Up @@ -483,6 +483,24 @@ pub(super) struct GenerationDecodeBudgetV1 {
const GENERATION_DECODE_RESIDENT_COMPONENT_V1: &str = "code-index-generation-decode-v1";
const SEALED_GRAPH_BUILD_RESIDENT_COMPONENT_V1: &str = "code-graph-sealed-build-v1";

/// The two corpus-sized passes over the active generation admission charges.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) enum ActiveGenerationWorkV1 {
/// Materializing the whole generation.
Decode,
/// Projecting its code graph from the sealed segments.
SealedGraphBuild,
}

/// What each pass over one generation costs, measured on a copy this process
/// held.
#[derive(Clone, Debug, PartialEq, Eq)]
struct ActiveGenerationChargesV1 {
generation_id: CodeGenerationId,
decode_bytes: u64,
graph_build_bytes: u64,
}

/// The resident cost of materializing the active generation.
#[derive(Clone, Debug, PartialEq, Eq)]
pub(super) enum ActiveGenerationDecodeChargeV1 {
Expand All @@ -503,7 +521,7 @@ pub struct DaemonCodeIndexPublicationStoreV1 {
active_encoded_bytes: Arc<AtomicU64>,
/// Resident bytes of the active generation as last installed, kept
/// after the decode is released: what materializing it again costs.
active_decode_charge: Arc<Mutex<Option<(CodeGenerationId, u64)>>>,
active_decode_charge: Arc<Mutex<Option<ActiveGenerationChargesV1>>>,
decode_admission: Arc<Mutex<Option<GenerationDecodeBudgetV1>>>,
decoded_content: SharedDecodedContentPoolV1,
pub(super) seal_encoded_segment_bytes: Arc<AtomicU64>,
Expand Down Expand Up @@ -2297,6 +2315,10 @@ impl DaemonCodeIndexPublicationStoreV1 {
self.active_encoded_bytes
.store(encoded_bytes, Ordering::Release);
hotpath::gauge!("daemon.code_index.generation.decode.bytes").set(encoded_bytes);
if let Some(peak_growth) = generation.decode_peak_growth_bytes() {
hotpath::gauge!("daemon.code_index.generation.decode.peak_growth_bytes")
.set(peak_growth);
}
Ok(Some(Arc::new(generation)))
}

Expand All @@ -2317,28 +2339,31 @@ impl DaemonCodeIndexPublicationStoreV1 {
fn admit_active_decode(
&self,
) -> Result<Option<ResidentMemoryReservationV1>, CodeIndexPublicationStoreErrorV1> {
self.admit_active_generation_work(GENERATION_DECODE_RESIDENT_COMPONENT_V1, "decoding")
self.admit_active_generation_work(ActiveGenerationWorkV1::Decode)
}

/// Charge building the active generation's code graph from its sealed
/// segments, the same way and the same bytes as decoding it: the build
/// holds the generation's cross-file resolution inputs and then its
/// segments the way a decode is charged, with the generation's retained
/// bytes: the build holds its cross-file resolution inputs and then its
/// compact graph store, both bounded by the generation it projects. The
/// caller holds the reservation for the build.
pub(super) fn admit_sealed_graph_build(
&self,
) -> Result<Option<ResidentMemoryReservationV1>, CodeIndexPublicationStoreErrorV1> {
self.admit_active_generation_work(
SEALED_GRAPH_BUILD_RESIDENT_COMPONENT_V1,
"building the code graph of",
)
self.admit_active_generation_work(ActiveGenerationWorkV1::SealedGraphBuild)
}

fn admit_active_generation_work(
&self,
component: &'static str,
work: &str,
work: ActiveGenerationWorkV1,
) -> Result<Option<ResidentMemoryReservationV1>, CodeIndexPublicationStoreErrorV1> {
let (component, work_label) = match work {
ActiveGenerationWorkV1::Decode => (GENERATION_DECODE_RESIDENT_COMPONENT_V1, "decoding"),
ActiveGenerationWorkV1::SealedGraphBuild => (
SEALED_GRAPH_BUILD_RESIDENT_COMPONENT_V1,
"building the code graph of",
),
};
let Some(admission) = self
.decode_admission
.lock()
Expand All @@ -2347,7 +2372,7 @@ impl DaemonCodeIndexPublicationStoreV1 {
else {
return Ok(None);
};
let (generation_id, requested) = match self.active_decode_charge()? {
let (generation_id, requested) = match self.active_generation_charge(work)? {
ActiveGenerationDecodeChargeV1::Decoded => return Ok(None),
ActiveGenerationDecodeChargeV1::Unmeasured => {
tracing::debug!(
Expand All @@ -2367,19 +2392,14 @@ impl DaemonCodeIndexPublicationStoreV1 {
let admissible = || -> Result<(), String> {
let snapshot = admission.resident_memory.snapshot();
let pressure = admission.resident_memory.pressure();
let observed = sampled_process_resident_bytes_v1().map_or(0, |observed| {
pressure
.publish_observed_resident_bytes(observed)
.observed_bytes()
.unwrap_or(observed)
});
let observed = pressure.measure_admission_bytes();
let watermark = pressure.high_watermark_bytes().min(snapshot.limit_bytes);
let available = watermark.saturating_sub(snapshot.used_bytes.max(observed));
if requested.get() <= available {
Ok(())
} else {
Err(format!(
"{work} generation {generation_id} needs {} resident bytes; {available} are \
"{work_label} generation {generation_id} needs {} resident bytes; {available} are \
available below the {watermark}-byte admission watermark",
requested.get()
))
Expand All @@ -2392,16 +2412,15 @@ impl DaemonCodeIndexPublicationStoreV1 {
for release in &released {
log_resident_owner_release_v1(release);
}
if released.is_empty() {
return Err(CodeIndexPublicationStoreErrorV1::ResidentMemoryRefused(
detail,
));
}
// Freed pages sit in the pool workers' allocator heaps whether or
// not an owner was shed, so they are returned before the
// re-measure that decides the refusal.
let trim = release_process_allocator_memory_v1();
tracing::info!(
event = "code_index_generation_decode_shed_retained_state",
released_owners = released.len(),
trimmed_bytes = trim.released_bytes(),
refused = %detail,
"generation decode shed retained state before re-measuring headroom"
);
admissible().map_err(CodeIndexPublicationStoreErrorV1::ResidentMemoryRefused)?;
Expand All @@ -2424,21 +2443,29 @@ impl DaemonCodeIndexPublicationStoreV1 {
})
}

/// A decoded copy charges its next decode what this decode measured at
/// its peak, which covers the pass transients above what it retains; a
/// generation built in memory has only its retained bytes to go on.
fn record_active_decode_charge(&self, generation: &CodeIndexPublishedGenerationV1) {
let retained = generation.retained_bytes();
*self
.active_decode_charge
.lock()
.unwrap_or_else(PoisonError::into_inner) = Some((
generation.manifest().generation_id.clone(),
generation.retained_bytes(),
));
.unwrap_or_else(PoisonError::into_inner) = Some(ActiveGenerationChargesV1 {
generation_id: generation.manifest().generation_id.clone(),
decode_bytes: generation
.decode_peak_growth_bytes()
.map_or(retained, |peak| peak.max(retained)),
graph_build_bytes: retained,
});
}

/// What reading the whole active generation into memory costs now:
/// nothing when it is already decoded, its measured resident bytes when
/// this process has held it before, unmeasured otherwise.
pub(super) fn active_decode_charge(
pub(super) fn active_generation_charge(
&self,
work: ActiveGenerationWorkV1,
) -> Result<ActiveGenerationDecodeChargeV1, CodeIndexPublicationStoreErrorV1> {
if self.cache.lock_state()?.active.is_some() {
return Ok(ActiveGenerationDecodeChargeV1::Decoded);
Expand All @@ -2451,8 +2478,16 @@ impl DaemonCodeIndexPublicationStoreV1 {
.lock()
.unwrap_or_else(PoisonError::into_inner)
.as_ref()
.filter(|(generation, _)| generation.as_str() == pointer.generation_id)
.map(|(generation, bytes)| (generation.clone(), *bytes));
.filter(|charges| charges.generation_id.as_str() == pointer.generation_id)
.map(|charges| {
(
charges.generation_id.clone(),
match work {
ActiveGenerationWorkV1::Decode => charges.decode_bytes,
ActiveGenerationWorkV1::SealedGraphBuild => charges.graph_build_bytes,
},
)
});
Ok(match charge {
Some((generation_id, bytes)) => ActiveGenerationDecodeChargeV1::Measured {
generation_id,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -109,16 +109,12 @@ pub async fn retry_deferred_query_authority_until_serving<F, Fut>(
let mut serving_changes = None;
loop {
// Subscribe before probing so a seat that lands between subscribe and
// the ready check remains observable. Demand a complete generation once
// the watch exists so a late subscribe after a silent Noop still gets a
// follow-up wake (same contract as the deferred advisory owner).
// the ready check remains observable. The mount needs only the text
// owner, so it demands no decoded generation.
if serving_changes.is_none() {
serving_changes = registry
.subscribe_serving_generation_changes(&project_root)
.await;
if serving_changes.is_some() {
let _ = registry.request_complete_generation(&project_root).await;
}
}
if registry
.retained_text_owner_for_root(&project_root)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1217,19 +1217,19 @@ fn dashboard_generation_is_ready(

/// Whether the generation status advertises is the one search will serve.
///
/// Status identity is taken from the text owner, but a publication installs
/// the replacement text owner before the serving swap (`registry/mount.rs`),
/// and a graph-on search answers from the decoded seat. Until the seat catches
/// up, search serves the predecessor and the split is not terminal. A graph-off
/// mount serves the text owner directly and deliberately never seats the
/// decoded generation, so there is no seat to compare.
/// Status identity is taken from the text owner. A graph-on search answers
/// from a decoded seat when one is held; while the seat still names the
/// predecessor, search serves it and the split is not terminal. With no seat,
/// search serves the text owner itself (exact, lexical, and the graph its
/// mapped store serves), which is the advertised generation. A graph-off
/// mount never seats the decoded generation, so there is no seat to compare.
pub(super) fn serving_seat_matches_advertised_generation(
graph_activation_enabled: bool,
serving_generation_id: Option<&str>,
advertised_text_generation_id: Option<&str>,
) -> bool {
match advertised_text_generation_id {
Some(advertised) if graph_activation_enabled => serving_generation_id == Some(advertised),
match (advertised_text_generation_id, serving_generation_id) {
(Some(advertised), Some(seated)) if graph_activation_enabled => seated == advertised,
_ => true,
}
}
Expand All @@ -1238,7 +1238,6 @@ pub(super) fn serving_seat_matches_advertised_generation(
/// status advertises.
pub(super) fn dashboard_terminal_status(
latest: Option<&LatestCompleteCodeIndexV1>,
released_seat: Option<&str>,
text: Option<&LatestCodeTextGenerationV1>,
text_ready: bool,
graph_activation_enabled: bool,
Expand All @@ -1251,9 +1250,7 @@ pub(super) fn dashboard_terminal_status(
code_graph_serving,
) && serving_seat_matches_advertised_generation(
graph_activation_enabled,
latest
.map(|latest| latest.generation().manifest().generation_id.as_str())
.or(released_seat),
latest.map(|latest| latest.generation().manifest().generation_id.as_str()),
text.map(|text| text.metadata().manifest().generation_id.as_str()),
)
}
Expand Down Expand Up @@ -3537,15 +3534,14 @@ impl CodeIndexSchedulerRegistryV1 {
if mounted_root != project_root {
return None;
}
// Feedback identity reads only the sealed manifest's snapshot, so an
// undecoded generation's text owner answers as well as a decoded one.
if let Some(generation) = self
.latest_complete_ready_decoded_for_root_scope(&project_root, scope)
.await
{
return Some(generation.text_generation_handle());
}
if let Some(generation) = self.latest_complete_ready_for_scope(scope).await {
return Some(generation.text_generation_handle());
}
self.latest_text_serving_freshness_for_scope(scope)
.await
.and_then(|(generation, current)| current.then_some(generation))
Expand Down
Loading
Loading