From 20efb3f1307b4a33717a671e668a9c1cec7c7e13 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Sun, 27 Sep 2026 11:04:58 +0000 Subject: [PATCH] fix(code-index): serve the first index without the whole decode A first index at a 6 GB cap parked forever: after the text build, the sealed graph build (charged the decode's 2.82 GB retained bytes) and then the serving decode were measured against VmRSS, which counts clean mapped container pages, and ready waited on a decoded seat. - Admission reads RssAnon + RssShmem through the pressure cell's sampler (injectable), and a refused decode or text build returns the pool workers' freed allocator pages before re-measuring, whether or not an owner was shed. - A publication seats its graph head from the text owner's mapped store (verified-head recovery), so ready no longer needs the decode. With no decoded seat, search serves the text owner, so an empty seat is terminal; the released-seat bookkeeping is gone. - Project-open owners (query authority, advisory, TypeScript producer), the LSP census and feedback identity read the sealed manifest and no longer demand the whole decode. - A graph build refused while its own text build holds the memory retries at the text join instead of parking; a refused decode after the graph serves does not park convergence. - A decode records its measured peak growth and its next decode is charged that, not the retained figure. Cross-file edges resolve before per-file copies into one exact-capacity vector. --- .../src/code_index_scheduler/memory_tests.rs | 85 +++++++- .../code_index_scheduler/publication_store.rs | 97 ++++++--- .../src/code_index_scheduler/query_runtime.rs | 8 +- .../src/code_index_scheduler/registry.rs | 26 +-- .../code_index_scheduler/registry/mount.rs | 198 +++++++++++++++++- .../registry/serving_readiness_tests.rs | 4 +- .../registry/serving_reads.rs | 12 -- .../src/code_index_scheduler/residency.rs | 20 +- .../src/code_index_scheduler/serving.rs | 16 +- .../code_index_scheduler/tests/residency.rs | 12 +- .../src/production/helpers.rs | 22 +- .../src/production/mod.rs | 15 ++ .../src/production/partitioned_codec.rs | 41 ++-- .../src/production/sealed_codec.rs | 42 +++- .../staging_footprint.rs | 2 +- .../src/resident_memory.rs | 113 ++++++++-- .../src/resident_memory/owners.rs | 1 + .../src/resident_memory/tests.rs | 20 ++ ...de_index_runtime_graph_activation_tests.rs | 158 +++++++++++++- crates/tracedecay/src/daemon/maintenance.rs | 16 +- .../src/daemon/production_harness.rs | 28 +-- .../code_index_activation.rs | 11 - .../src/daemon/project_open_owners.rs | 7 +- .../advisory_runtime/deferred.rs | 40 +--- .../compiler_diagnostics_producer.rs | 3 - 25 files changed, 754 insertions(+), 243 deletions(-) diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/memory_tests.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/memory_tests.rs index 0bac49c733..22001a2bd9 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/memory_tests.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/memory_tests.rs @@ -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; @@ -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::>(), + ["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" + ); +} 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 5c866ddfbf..9376043b2b 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 @@ -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::{ @@ -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 { @@ -503,7 +521,7 @@ pub struct DaemonCodeIndexPublicationStoreV1 { active_encoded_bytes: Arc, /// Resident bytes of the active generation as last installed, kept /// after the decode is released: what materializing it again costs. - active_decode_charge: Arc>>, + active_decode_charge: Arc>>, decode_admission: Arc>>, decoded_content: SharedDecodedContentPoolV1, pub(super) seal_encoded_segment_bytes: Arc, @@ -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))) } @@ -2317,28 +2339,31 @@ impl DaemonCodeIndexPublicationStoreV1 { fn admit_active_decode( &self, ) -> Result, 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, 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, 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() @@ -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!( @@ -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() )) @@ -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)?; @@ -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 { if self.cache.lock_state()?.active.is_some() { return Ok(ActiveGenerationDecodeChargeV1::Decoded); @@ -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, diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/query_runtime.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/query_runtime.rs index 11d899555c..3b53d978d9 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/query_runtime.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/query_runtime.rs @@ -109,16 +109,12 @@ pub async fn retry_deferred_query_authority_until_serving( 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) diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry.rs index b9fed3438a..73b3ce68da 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry.rs @@ -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, } } @@ -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, @@ -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()), ) } @@ -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)) 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 1955937fd7..63efe0114c 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 @@ -208,6 +208,105 @@ impl CodeIndexSchedulerRegistryV1 { .await } + /// Seat the graph head a publication just wrote on its text owner, the + /// way a restart recovers it, so graph reads serve from the mapped store + /// without the whole-generation decode. `false` leaves the decode and its + /// activation to seat the graph. + async fn recover_published_graph_head( + scheduler: &Arc>, + shutting_down: &Arc, + passes: &Arc, + graph_activation: &CodeGraphActivationAuthorityV1, + (project_id, repository_id, worktree_id): ( + &ProjectId, + &tracedecay_domain::RepositoryId, + &tracedecay_domain::WorktreeId, + ), + text: &LatestCodeTextGenerationV1, + ) -> bool { + let generation_id = text.metadata().manifest().generation_id.clone(); + let binding_scheduler = Arc::clone(scheduler); + let binding_shutting_down = Arc::clone(shutting_down); + let binding_passes = Arc::clone(passes); + let binding = tokio::task::spawn_blocking(move || { + Self::lock_scheduler_for_graph_step( + &binding_scheduler, + &binding_shutting_down, + &binding_passes, + )? + .1 + .code_graph_replay_binding(&generation_id) + }) + .await; + let binding = match binding { + Ok(Ok(binding)) => binding, + Ok(Err(error)) => { + tracing::warn!( + event = "code_index_published_graph_head_binding_unavailable", + error = %error, + "the published graph head has no replay binding; the serving decode seats it" + ); + return false; + } + Err(error) => { + tracing::warn!( + event = "code_index_published_graph_head_binding_task_failed", + error = %error, + "the published graph head binding task failed; the serving decode seats it" + ); + return false; + } + }; + match graph_activation + .recover_verified_head( + project_id, + repository_id, + worktree_id, + text.clone(), + binding, + Arc::clone(shutting_down), + ) + .await + { + Ok(recovered) => recovered, + Err(error) => { + tracing::warn!( + event = "code_index_published_graph_head_recovery_failed", + error = %error, + "the published graph head did not seat from the text owner; the serving \ + decode seats it" + ); + false + } + } + } + + /// Drop a seat that no longer names the advertised generation. Search + /// serves the text owner while the seat is empty, and holding the + /// predecessor's decode would only keep a corpus-sized generation alive. + fn release_superseded_serving_seat( + serving_generation: &ServingGenerationSlot, + serving_generation_epoch: &AtomicU64, + serving_source_witness: &RwLock>, + serving_seats: &tokio::sync::watch::Sender, + serving_generation_changed: &tokio::sync::watch::Sender<()>, + ) { + let displaced = serving_generation + .write() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .take(); + if displaced.is_none() { + return; + } + serving_generation_epoch.fetch_add(1, Ordering::AcqRel); + *serving_source_witness + .write() + .unwrap_or_else(std::sync::PoisonError::into_inner) = None; + drop(displaced); + Self::record_serving_seat(serving_seats); + serving_generation_changed.send_replace(()); + } + /// A same-root remount keeps the incumbent owner: it may only refresh the /// graph activation policy, and it wakes the worker so a policy change or /// an edit that raced the remount is picked up by the next pass. @@ -1490,8 +1589,12 @@ impl CodeIndexSchedulerRegistryV1 { prepare_graph = serving_empty && worker_complete_generation_requested.load(Ordering::Acquire); } + // Recovery installs a fresh graph store on the owner; an owner + // that already serves one (a publication seated it) keeps it + // and its warm catalog. if prepare_graph && !published_pass + && !graph_already_serves && !retained_graph_head_recovery_attempted && let Some(retained) = graph_text .as_ref() @@ -1681,6 +1784,11 @@ impl CodeIndexSchedulerRegistryV1 { ); } } + // A publication's graph build refused while its own text + // projection still holds the build memory is sequencing, not + // a stall: it runs again as soon as that projection joins. + let mut graph_waits_for_text = false; + let text_projection_running = published_text_projection.is_some(); let mut result = match source_result { Ok(mut outcome) if prepare_graph => { // Publish the graph head from the sealed segments @@ -1690,7 +1798,10 @@ impl CodeIndexSchedulerRegistryV1 { // together; the activation after the decode recovers // the head published here. let mut graph_publish_refusal = None; - if let Some(text) = graph_text.as_ref() { + let mut graph_head_published = false; + // An owner that already serves this generation's graph + // needs no second build: only the decode is demanded. + if let Some(text) = graph_text.as_ref().filter(|_| !graph_already_serves) { let generation_id = text.metadata().manifest().generation_id.clone(); let binding_scheduler = Arc::clone(&worker_scheduler); let shutting_down = Arc::clone(&worker_shutting_down); @@ -1740,7 +1851,7 @@ impl CodeIndexSchedulerRegistryV1 { .await; drop(reservation); match published { - Ok(_) => {} + Ok(published) => graph_head_published = published, Err(error) if error.is_resident_memory_graph_refusal() => { graph_publish_refusal = Some(error.to_string()); } @@ -1766,6 +1877,46 @@ impl CodeIndexSchedulerRegistryV1 { ), } } + // The head just published serves graph reads from the + // text owner's mapped store, exactly as a restart's + // verified-head recovery does, so the whole-generation + // decode below runs only for a reader that demanded it. + let graph_serves_from_text = match graph_text.as_ref() { + Some(text) if graph_head_published => { + Self::recover_published_graph_head( + &worker_scheduler, + &worker_shutting_down, + &worker_reconcile_in_progress, + &worker_graph_activation, + ( + &worker_project_id, + &worker_repository_id, + &worker_worktree_id, + ), + text, + ) + .await + } + _ => false, + }; + let defer_serving_decode = graph_serves_from_text + && !worker_complete_generation_requested.load(Ordering::Acquire); + if defer_serving_decode { + clear_graph_resident_memory_park(&worker_convergence_park); + worker_memory_retry.reset(); + Self::release_superseded_serving_seat( + &worker_serving_generation, + &worker_serving_generation_epoch, + &worker_serving_source_witness, + &worker_serving_seats, + &worker_serving_generation_changed, + ); + tracing::info!( + event = "code_index_serving_decode_deferred", + "the published graph head serves from the text owner; the \ + whole-generation decode waits for a reader that needs it" + ); + } let graph_scheduler = Arc::clone(&worker_scheduler); let graph_text = graph_text.clone(); let shutting_down = Arc::clone(&worker_shutting_down); @@ -1779,6 +1930,9 @@ impl CodeIndexSchedulerRegistryV1 { if let Some(detail) = graph_publish_refusal { return Ok((None, None, false, Some(detail))); } + if defer_serving_decode { + return Ok((None, None, false, None)); + } let decoder = Self::lock_scheduler_for_graph_step( &graph_scheduler, &shutting_down, @@ -1880,17 +2034,32 @@ impl CodeIndexSchedulerRegistryV1 { .await { Ok(Ok((_, _, _, Some(detail)))) => { - park_convergence( - &worker_convergence_park, - detail.clone(), - CONVERGENCE_PARK_GRAPH_RESIDENT_MEMORY_REMEDIATION_V1, - Some(CodeIndexBuildBlockedReasonV1::ResidentMemory), - true, - ); - worker_memory_retry.schedule(&worker_pending_wake, &worker_wake); + // Once the text owner serves the graph, the + // generation has converged: only the reader + // that demanded the whole decode waits for + // memory, and freshness is not parked for it. + let converged = graph_serves_from_text || graph_already_serves; + // A build that waits for this pass's own text + // projection is rescheduled at its join below. + graph_waits_for_text = published_pass && text_projection_running; + if !graph_waits_for_text { + if !converged { + park_convergence( + &worker_convergence_park, + detail.clone(), + CONVERGENCE_PARK_GRAPH_RESIDENT_MEMORY_REMEDIATION_V1, + Some(CodeIndexBuildBlockedReasonV1::ResidentMemory), + true, + ); + } + worker_memory_retry + .schedule(&worker_pending_wake, &worker_wake); + } tracing::warn!( event = "code_index_graph_prepare_decode_refused", published_pass, + converged, + graph_waits_for_text, detail = %detail, "the sealed generation waits to decode until memory is \ given back; text serving is unaffected" @@ -1898,7 +2067,7 @@ impl CodeIndexSchedulerRegistryV1 { Ok((outcome, None, None)) } Ok(Ok((latest, replay_binding, roster_refusal_rebuild, None))) => { - if latest.is_none() { + if latest.is_none() && !defer_serving_decode { tracing::warn!( event = "code_index_graph_prepare_no_servable_generation", published_pass, @@ -2103,6 +2272,9 @@ impl CodeIndexSchedulerRegistryV1 { } }); } + if graph_waits_for_text { + Self::note_worker_continuation(&worker_pending_wake, &worker_wake); + } if let Some(outcome) = published_text_projection_outcome.take() { // A clone-fingerprint successor is still `Unfinished` work // after exact and lexical owners are ready. That must not @@ -2184,6 +2356,9 @@ impl CodeIndexSchedulerRegistryV1 { ), } } + // The ready text owner now serves this generation, + // with or without a decoded seat after it. + worker_serving_generation_changed.send_replace(()); } PublishedTextProjectionOutcomeV1::Shutdown => { tracing::info!( @@ -2765,6 +2940,7 @@ impl CodeIndexSchedulerRegistryV1 { }; if matches!(outcome, PublishedTextProjectionOutcomeV1::Finished) { worker_memory_retry.reset(); + worker_serving_generation_changed.send_replace(()); } match outcome { PublishedTextProjectionOutcomeV1::Finished diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/serving_readiness_tests.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/serving_readiness_tests.rs index 2c6fa644fb..f5c71f8dbc 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/serving_readiness_tests.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/serving_readiness_tests.rs @@ -73,8 +73,8 @@ fn graph_rebuild_split_is_not_terminal_freshness() { "search served the predecessor while status advertised the head" ); assert!( - !serving_seat_matches_advertised_generation(true, None, Some("generation.head")), - "a text owner installed before the serving swap is not yet the served generation" + serving_seat_matches_advertised_generation(true, None, Some("generation.head")), + "with no decoded seat search serves the advertised text owner itself" ); assert!(serving_seat_matches_advertised_generation( true, diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/serving_reads.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/serving_reads.rs index 741e9f1537..651fd33201 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/serving_reads.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/serving_reads.rs @@ -280,7 +280,6 @@ impl CodeIndexSchedulerRegistryV1 { pending_wake, source_freshness, graph_activation_enabled, - residency, ) = { let mounted = self.mounted.lock().await; let Some(worktree) = mounted.get(&canonical_root) else { @@ -299,7 +298,6 @@ impl CodeIndexSchedulerRegistryV1 { Arc::clone(&worktree.pending_wake), worktree.source_freshness.clone(), worktree.graph_activation.policy().is_enabled(), - Arc::clone(&worktree.residency), ) }; tokio::task::spawn_blocking(move || { @@ -352,10 +350,6 @@ impl CodeIndexSchedulerRegistryV1 { .read() .unwrap_or_else(std::sync::PoisonError::into_inner) .clone(); - let released_seat = latest - .is_none() - .then(|| residency.released_seat()) - .flatten(); let text = text_generation .read() .unwrap_or_else(std::sync::PoisonError::into_inner) @@ -391,7 +385,6 @@ impl CodeIndexSchedulerRegistryV1 { }); let ready = dashboard_terminal_status( latest.as_ref(), - released_seat.as_ref().map(CodeGenerationId::as_str), text.as_ref(), text_ready, graph_activation_enabled, @@ -448,10 +441,6 @@ impl CodeIndexSchedulerRegistryV1 { .read() .unwrap_or_else(std::sync::PoisonError::into_inner) .clone(); - let released_seat = latest - .is_none() - .then(|| residency.released_seat()) - .flatten(); let text = text_generation .read() .unwrap_or_else(std::sync::PoisonError::into_inner) @@ -479,7 +468,6 @@ impl CodeIndexSchedulerRegistryV1 { }); let ready = dashboard_terminal_status( latest.as_ref(), - released_seat.as_ref().map(CodeGenerationId::as_str), text.as_ref(), text_ready, graph_activation_enabled, diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/residency.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/residency.rs index a1bae9f06d..2ed25b3bb7 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/residency.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/residency.rs @@ -21,7 +21,6 @@ use tracedecay_code_index::graph_projection::{ use tracedecay_code_index::production::{ CodeIndexPublishedGenerationV1, DecodedGenerationContentV1, }; -use tracedecay_domain::CodeGenerationId; use tracedecay_runtime_core::resident_memory::{ ResidentHoldingV1, ResidentOwnerBytesV1, ResidentOwnerKindV1, ResidentOwnerRegistrationV1, ResidentOwnerReleaseV1, ResidentOwnerSampleV1, ResidentOwnerScopeV1, ResidentOwnerV1, @@ -41,10 +40,6 @@ pub(super) struct WorktreeResidencyV1 { publication: DaemonCodeIndexPublicationStoreV1, text_generation: Arc>>, last_used: Mutex, - /// The generation the seat held when the inventory released it. The - /// worktree still serves it from disk, so freshness keeps reporting it - /// as seated until a later generation takes the seat. - released_seat: Mutex>, } pub(super) struct WorktreeResidencyPartsV1 { @@ -68,7 +63,6 @@ impl WorktreeResidencyV1 { publication: parts.publication, text_generation: parts.text_generation, last_used: Mutex::new(Instant::now()), - released_seat: Mutex::new(None), } } @@ -87,13 +81,6 @@ impl WorktreeResidencyV1 { .unwrap_or_else(PoisonError::into_inner) } - pub(super) fn released_seat(&self) -> Option { - self.released_seat - .lock() - .unwrap_or_else(PoisonError::into_inner) - .clone() - } - fn seated(&self) -> Option> { self.serving_generation .read() @@ -258,12 +245,7 @@ impl ResidentOwnerV1 for ServingDecodeOwnerV1 { residency .complete_generation_requested .store(false, Ordering::Release); - if let Some(seated) = &seated { - *residency - .released_seat - .lock() - .unwrap_or_else(PoisonError::into_inner) = - Some(seated.manifest().generation_id.clone()); + if seated.is_some() { residency .serving_generation_epoch .fetch_add(1, Ordering::AcqRel); diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/serving.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/serving.rs index 6ccca6676f..3e09ccde67 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/serving.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/serving.rs @@ -46,7 +46,7 @@ use tracedecay_private_fs::{open_private_file, validate_private_directory}; 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::{ @@ -1189,13 +1189,7 @@ impl DaemonCodeTextArtifactStoreV1 { minimum: NonZeroU64, ) -> Result<(u64, u64, u64, u64), RetrievalPortError> { let snapshot = self.resident_memory.snapshot(); - let observed_bytes = sampled_process_resident_bytes_v1().map_or(0, |observed| { - self.resident_memory - .pressure() - .publish_observed_resident_bytes(observed) - .observed_bytes() - .unwrap_or(observed) - }); + let observed_bytes = self.resident_memory.pressure().measure_admission_bytes(); let unmodeled_live_bytes = observed_bytes.saturating_sub(snapshot.used_bytes); let admission_watermark = self .resident_memory @@ -1250,17 +1244,17 @@ impl DaemonCodeTextArtifactStoreV1 { let released = self .resident_owners .shed(minimum.get(), std::time::Instant::now()); - if released.is_empty() { - return Err(RetrievalPortError::ResidentMemoryRefused(detail)); - } for release in &released { log_resident_owner_release_v1(release); } + // Freed pages sit in the pool workers' allocator heaps + // whether or not an owner was shed. let trim = release_process_allocator_memory_v1(); tracing::info!( event = "code_text_artifact_build_shed_retained_state", released_owners = released.len(), trimmed_bytes = trim.released_bytes(), + refused = %detail, "text-artifact build shed retained state before re-measuring headroom" ); self.text_artifact_admission(preferred, minimum) 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 d38c32bb8a..b1106002bf 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 @@ -212,17 +212,17 @@ async fn linked_worktrees_on_identical_content_hold_one_decoded_generation() { worktrees.sort(); assert_eq!( decoded(&alone), - [(vec![primary.worktree_id.clone()], true, Some(1_462_697))] + [(vec![primary.worktree_id.clone()], true, Some(1_458_681))] ); - assert_eq!(alone.measured_bytes, 1_462_697); - // Two copies would hold 2,925,394 bytes; the linked worktree adds only + assert_eq!(alone.measured_bytes, 1_458_681); + // Two copies would hold 2,917,362 bytes; the linked worktree adds only // the manifest, lineage, and projection evidence it sealed itself. assert_eq!( decoded(&both), - [(worktrees, true, Some(1_854_381))], + [(worktrees, true, Some(1_846_349))], "one decode row for one content, naming both worktrees" ); - assert_eq!(both.measured_bytes, 1_854_381); + assert_eq!(both.measured_bytes, 1_846_349); let later = Instant::now() + IDLE_WINDOW; let released = owners.release_idle(later); @@ -247,7 +247,7 @@ async fn linked_worktrees_on_identical_content_hold_one_decoded_generation() { .iter() .map(|release| release.bytes.measured().unwrap_or(0)) .sum::(), - 1_854_381, + 1_846_349, "the two releases give back exactly what the shared row held" ); let idle = owners.report(later); diff --git a/crates/tracedecay-code-index/src/production/helpers.rs b/crates/tracedecay-code-index/src/production/helpers.rs index d0e7ffa49a..f7d398ca28 100644 --- a/crates/tracedecay-code-index/src/production/helpers.rs +++ b/crates/tracedecay-code-index/src/production/helpers.rs @@ -366,11 +366,19 @@ pub(crate) fn collect_edge_evidence( where T: AsRef + Sync, { - let mut edges = files + // Resolution's whole-set indexes are gone before the per-file copies are + // made, and the one exact-capacity vector never regrows, so the peak is + // the edges this returns plus one sort buffer. + let cross_file = resolve_cross_file_references(files)?; + let per_file = files .iter() - .flat_map(|file| file.as_ref().artifacts.edges.clone()) - .collect::>(); - edges.extend(resolve_cross_file_references(files)?); + .map(|file| file.as_ref().artifacts.edges.len()) + .sum::(); + let mut edges = Vec::with_capacity(per_file.saturating_add(cross_file.len())); + for file in files { + edges.extend(file.as_ref().artifacts.edges.iter().cloned()); + } + edges.extend(cross_file); edges.sort_by(edge_order); let mut abstentions = files .iter() @@ -443,7 +451,11 @@ where index, ) })?; - let mut edges = per_file.into_iter().flatten().collect::>(); + drop((by_simple_name, rust_files, typescript_modules)); + let mut edges = Vec::with_capacity(per_file.iter().map(Vec::len).sum()); + for file_edges in per_file { + edges.extend(file_edges); + } hotpath::measure_block!("code_index.seal.edge_materialization", { edges.sort_by(edge_order); edges.dedup(); diff --git a/crates/tracedecay-code-index/src/production/mod.rs b/crates/tracedecay-code-index/src/production/mod.rs index 12c2dc0fd7..833715f1d3 100644 --- a/crates/tracedecay-code-index/src/production/mod.rs +++ b/crates/tracedecay-code-index/src/production/mod.rs @@ -801,6 +801,11 @@ pub struct CodeIndexPublishedGenerationV1 { chunk_policy: OnceLock, /// [`Self::retained_bytes`] of the immutable decode, measured once. retained_bytes: OnceLock, + /// How far this process's unreclaimable resident bytes rose above their + /// starting point while this generation was decoded, sampled at every + /// decode pass boundary. `None` for a generation built in memory, or when + /// the kernel reports no resident set. + decode_peak_growth_bytes: Option, } /// The chunk policy-revision census of one immutable generation: no chunks at @@ -999,6 +1004,15 @@ impl CodeIndexPublishedGenerationV1 { decode.saturating_add(attribution) } + /// What decoding this generation from its sealed bytes cost at its peak, + /// as measured when this copy was decoded; `None` for a generation built + /// in memory. Admission charges this, not [`Self::retained_bytes`], for + /// the next decode of the same generation. + #[must_use] + pub fn decode_peak_growth_bytes(&self) -> Option { + self.decode_peak_growth_bytes + } + fn measure_decode_bytes(&self) -> usize { self.measure_resident_bytes() } @@ -2196,6 +2210,7 @@ where attribution: OnceLock::new(), chunk_policy: OnceLock::new(), retained_bytes: OnceLock::new(), + decode_peak_growth_bytes: None, }; hotpath::measure_block!( "code_index.build.assemble.validate", diff --git a/crates/tracedecay-code-index/src/production/partitioned_codec.rs b/crates/tracedecay-code-index/src/production/partitioned_codec.rs index b17e4f71e4..f240de2319 100644 --- a/crates/tracedecay-code-index/src/production/partitioned_codec.rs +++ b/crates/tracedecay-code-index/src/production/partitioned_codec.rs @@ -62,10 +62,10 @@ use super::projection_rows::{ PersistedProjectionRequestV1, chunk_roster, }; use super::sealed_codec::{ - FileScopeIdentityV1, PersistedFileGenerationArtifactsRefV2, PersistedFileGenerationArtifactsV1, - PersistedFileGenerationArtifactsV2, SEALED_GENERATION_FORMAT_REVISION_V1, - StreamingPersistedPublishedGenerationV1, assemble_published_generation, restore_file_pages, - superseded_sealed_generation_revision, + DecodePeakProbeV1, FileScopeIdentityV1, PersistedFileGenerationArtifactsRefV2, + PersistedFileGenerationArtifactsV1, PersistedFileGenerationArtifactsV2, + SEALED_GENERATION_FORMAT_REVISION_V1, StreamingPersistedPublishedGenerationV1, + assemble_published_generation, restore_file_pages, superseded_sealed_generation_revision, }; use super::*; @@ -2138,6 +2138,7 @@ impl CodeIndexPublishedGenerationV1 { &mut Vec, ) -> Result<(), CodeIndexProductionErrorV1>, ) -> Result { + let mut probe = DecodePeakProbeV1::start(); let generation = hotpath::measure_block!( "code_index.restore.manifest", parse_partitioned_manifest(bytes) @@ -2151,10 +2152,12 @@ impl CodeIndexPublishedGenerationV1 { decode_file_pages(&generation, &scope, &mut read_segment)?, )), }; + probe.sample(); let evidence = hotpath::measure_block!( "code_index.restore.generation_evidence", decode_generation_evidence(&generation.generation_evidence, read_segment) )?; + probe.sample(); let files = &content.files; let (lineage, projection_request, projection_receipt) = hotpath::measure_block!("code_index.restore.evidence_expand", { @@ -2172,19 +2175,23 @@ impl CodeIndexPublishedGenerationV1 { .expand(&generation.manifest.generation_id, &symbols)?; Ok::<_, CodeIndexProductionErrorV1>((lineage, request, receipt)) })?; - assemble_published_generation(StreamingPersistedPublishedGenerationV1 { - manifest: generation.manifest, - snapshot: generation.snapshot, - repository_parse_identity: generation.repository_parse_identity, - ignored_source_admissions: generation.ignored_source_admissions, - ignored_source_admissions_digest: generation.ignored_source_admissions_digest, - content, - lineage, - coverage: generation.coverage, - capability: generation.capability, - projection_request, - projection_receipt, - }) + probe.sample(); + assemble_published_generation( + StreamingPersistedPublishedGenerationV1 { + manifest: generation.manifest, + snapshot: generation.snapshot, + repository_parse_identity: generation.repository_parse_identity, + ignored_source_admissions: generation.ignored_source_admissions, + ignored_source_admissions_digest: generation.ignored_source_admissions_digest, + content, + lineage, + coverage: generation.coverage, + capability: generation.capability, + projection_request, + projection_receipt, + }, + probe, + ) } /// Authenticate only the tiny partitioned manifest and return the metadata diff --git a/crates/tracedecay-code-index/src/production/sealed_codec.rs b/crates/tracedecay-code-index/src/production/sealed_codec.rs index 7ed4de3fb8..f989680161 100644 --- a/crates/tracedecay-code-index/src/production/sealed_codec.rs +++ b/crates/tracedecay-code-index/src/production/sealed_codec.rs @@ -1087,8 +1087,41 @@ pub(super) fn restore_file_pages( .collect()) } +/// Unreclaimable process bytes sampled at a decode's pass boundaries. Pages a +/// pass frees stay with the allocator until it is collected, so the sample +/// after each pass reads that pass's high-water mark. +pub(super) struct DecodePeakProbeV1 { + start_bytes: Option, + peak_bytes: u64, +} + +impl DecodePeakProbeV1 { + pub(super) fn start() -> Self { + let start_bytes = + tracedecay_runtime_core::resident_memory::sampled_process_resident_bytes_v1(); + Self { + start_bytes, + peak_bytes: start_bytes.unwrap_or(0), + } + } + + pub(super) fn sample(&mut self) { + if let Some(bytes) = + tracedecay_runtime_core::resident_memory::sampled_process_resident_bytes_v1() + { + self.peak_bytes = self.peak_bytes.max(bytes); + } + } + + fn growth_bytes(&self) -> Option { + self.start_bytes + .map(|start| self.peak_bytes.saturating_sub(start)) + } +} + pub(super) fn assemble_published_generation( generation: StreamingPersistedPublishedGenerationV1, + mut probe: DecodePeakProbeV1, ) -> Result { let StreamingPersistedPublishedGenerationV1 { manifest, @@ -1136,10 +1169,12 @@ pub(super) fn assemble_published_generation( }) })? .map_err(CodeIndexProductionErrorV1::Increment)?; + probe.sample(); let symbols = hotpath::measure_block!("code_index.sealed_decode.symbol_index", { GenerationSymbolIndexV1::new(manifest.generation_id.clone(), symbol_rows) }) .map_err(CodeIndexProductionErrorV1::Lineage)?; + probe.sample(); let imports = hotpath::measure_block!("code_index.sealed_decode.import_evidence", { derive_import_evidence(&files) }); @@ -1147,6 +1182,7 @@ pub(super) fn assemble_published_generation( hotpath::measure_block!("code_index.sealed_decode.edge_evidence", { collect_edge_evidence(&files) })?; + probe.sample(); let projection = hotpath::measure_block!("code_index.sealed_decode.projection_handoff", { ProjectionPublicationHandoffV1::restore(projection_request, projection_receipt) @@ -1162,7 +1198,8 @@ pub(super) fn assemble_published_generation( projection, )) })?; - let published = CodeIndexPublishedGenerationV1 { + probe.sample(); + let mut published = CodeIndexPublishedGenerationV1 { statistics: hotpath::measure_block!( "code_index.sealed_decode.statistics", CodeIndexGenerationStatisticsV1::from_generation_parts( @@ -1194,11 +1231,14 @@ pub(super) fn assemble_published_generation( attribution: OnceLock::new(), chunk_policy: OnceLock::new(), retained_bytes: OnceLock::new(), + decode_peak_growth_bytes: None, }; hotpath::measure_block!( "code_index.sealed_decode.corpus_validation", published.validate_fresh() )?; + probe.sample(); + published.decode_peak_growth_bytes = probe.growth_bytes(); Ok(published) } diff --git a/crates/tracedecay-graph-db/tests/graph_db_suite/verified_generation_contract/staging_footprint.rs b/crates/tracedecay-graph-db/tests/graph_db_suite/verified_generation_contract/staging_footprint.rs index c51969135f..ee05ef1406 100644 --- a/crates/tracedecay-graph-db/tests/graph_db_suite/verified_generation_contract/staging_footprint.rs +++ b/crates/tracedecay-graph-db/tests/graph_db_suite/verified_generation_contract/staging_footprint.rs @@ -1439,7 +1439,7 @@ fn corrupt_sealed_containers(root: &Path) -> usize { /// recovered digest its relational head names. /// /// This is the graph-db half of the daemon's -/// `corrupt_graph_restart_repairs_through_canonical_serialized_activation`: +/// `corrupt_graph_restart_rebuilds_from_the_sealed_segments_without_decoding`: /// publish a sealed code generation, restart, and destroy the sealed /// container. With no staging rows (the direct-seal shape) recovery has /// nothing to serve from, so the canonical serialized activation path diff --git a/crates/tracedecay-runtime-core/src/resident_memory.rs b/crates/tracedecay-runtime-core/src/resident_memory.rs index a593c40caa..f6d27c1e1d 100644 --- a/crates/tracedecay-runtime-core/src/resident_memory.rs +++ b/crates/tracedecay-runtime-core/src/resident_memory.rs @@ -291,27 +291,51 @@ pub fn resident_memory_watermark_bytes_v1(limit_bytes: NonZeroU64, permille: u64 u64::try_from(scaled).unwrap_or(u64::MAX) } -/// Sample this process's resident set size directly from the kernel. +/// One kernel reading of this process's resident set, split by whether the +/// kernel can take the pages back without swapping. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct ProcessResidentSampleV1 { + /// Every resident page (`VmRSS`), clean file-backed mappings included. + pub resident_bytes: u64, + /// Anonymous and shared-memory pages (`RssAnon + RssShmem`). Clean mapped + /// file pages (the sealed container, graph stores) are dropped by the + /// kernel on demand before the cgroup kill line, so admission counts only + /// these. + pub unreclaimable_bytes: u64, +} + +fn status_kib_field_bytes(status: &str, field: &str) -> Option { + status + .lines() + .find_map(|line| line.strip_prefix(field)?.strip_prefix(':'))? + .split_whitespace() + .next()? + .parse::() + .ok()? + .checked_mul(1_024) +} + +fn process_resident_sample_from_status_v1(status: &str) -> Option { + let anon = status_kib_field_bytes(status, "RssAnon")?; + let shmem = status_kib_field_bytes(status, "RssShmem")?; + Some(ProcessResidentSampleV1 { + resident_bytes: status_kib_field_bytes(status, "VmRSS")?, + unreclaimable_bytes: anon.checked_add(shmem)?, + }) +} + +/// Sample this process's resident set directly from the kernel. /// -/// The one `/proc/self/status` `VmRSS` parser in the workspace: the daemon's -/// dedicated resident-memory sampler publishes its samples into -/// [`process_resident_memory_pressure_v1`], and load-scoped watchdogs (the -/// semantic session pool's cold-load resident bound) sample it directly. -/// Returns `None` where the kernel surface is unavailable (non-Linux hosts), -/// which callers must treat as unobserved, never as zero. +/// The one `/proc/self/status` parser in the workspace: the daemon's +/// dedicated resident-memory sampler and every admission re-measure read it +/// through [`ResidentMemoryPressureV1::sample_and_publish`]. Returns `None` +/// where the kernel surface is unavailable (non-Linux hosts), which callers +/// must treat as unobserved, never as zero. #[must_use] -pub fn sampled_process_resident_bytes_v1() -> Option { +pub fn sampled_process_resident_v1() -> Option { #[cfg(target_os = "linux")] { - let status = std::fs::read_to_string("/proc/self/status").ok()?; - let kib = status - .lines() - .find_map(|line| line.strip_prefix("VmRSS:"))? - .split_whitespace() - .next()? - .parse::() - .ok()?; - kib.checked_mul(1_024) + process_resident_sample_from_status_v1(&std::fs::read_to_string("/proc/self/status").ok()?) } #[cfg(not(target_os = "linux"))] { @@ -319,6 +343,15 @@ pub fn sampled_process_resident_bytes_v1() -> Option { } } +/// The bytes admission trusts: [`ProcessResidentSampleV1::unreclaimable_bytes`]. +#[must_use] +pub fn sampled_process_resident_bytes_v1() -> Option { + sampled_process_resident_v1().map(|sample| sample.unreclaimable_bytes) +} + +/// Where a pressure cell reads the process's resident set. +pub type ProcessResidentSamplerV1 = dyn Fn() -> Option + Send + Sync; + /// Share of the last ten seconds, in percent, that some task in this process's /// cgroup stalled on memory, at or above which the daemon sheds retained state /// even while RSS is under its watermark. Under `MemoryHigh` the kernel @@ -427,10 +460,10 @@ pub struct ResidentMemoryPressureRegistrationFailureV1; /// The measured side of the memory accounting loop. /// -/// One dedicated reader samples real RSS (`/proc/self/status` `VmRSS` on -/// Linux), publishes the `daemon.process.resident_bytes` gauge, and feeds this -/// cell. Admission reads the same canonical observation; there is no second -/// parser or publisher. +/// One dedicated reader samples the process (`/proc/self/status` on Linux), +/// publishes the `daemon.process.resident_bytes` gauge, and feeds this cell +/// the unreclaimable bytes. Admission re-measures through the same sampler; +/// there is no second parser or publisher. pub struct ResidentMemoryPressureV1 { limit_bytes: NonZeroU64, high_watermark_bytes: u64, @@ -439,6 +472,7 @@ pub struct ResidentMemoryPressureV1 { observed: AtomicBool, over_budget: AtomicBool, state: ProfiledMutex, + sampler: Arc, } impl fmt::Debug for ResidentMemoryPressureV1 { @@ -456,14 +490,24 @@ impl fmt::Debug for ResidentMemoryPressureV1 { impl ResidentMemoryPressureV1 { #[must_use] pub fn new(limit_bytes: NonZeroU64) -> Self { - Self::with_reclaim_line(limit_bytes, None) + Self::with_reclaim_line(limit_bytes, None, Arc::new(sampled_process_resident_v1)) + } + + /// A cell that reads the process through `sampler` instead of the kernel. + #[must_use] + pub fn with_sampler(limit_bytes: NonZeroU64, sampler: Arc) -> Self { + Self::with_reclaim_line(limit_bytes, None, sampler) } /// `reclaim_watermark_bytes` is a cgroup `memory.high` that sits strictly /// below `limit_bytes`. It replaces the percentage high watermark so the /// operator's band down to `memory.max` is not discounted again. Absent, /// zero, or not strictly below the ceiling, the percentage watermarks stand. - fn with_reclaim_line(limit_bytes: NonZeroU64, reclaim_watermark_bytes: Option) -> Self { + fn with_reclaim_line( + limit_bytes: NonZeroU64, + reclaim_watermark_bytes: Option, + sampler: Arc, + ) -> Self { let percentage_high = resident_memory_watermark_bytes_v1( limit_bytes, RESIDENT_MEMORY_PRESSURE_HIGH_WATERMARK_PERMILLE_V1, @@ -497,9 +541,31 @@ impl ResidentMemoryPressureV1 { Mutex::new(ResidentMemoryPressureReclaimerStateV1::default()), label = "runtime_core.resident.pressure" ), + sampler, } } + /// Read the process and publish its unreclaimable bytes as the admission + /// observation. `None` when the process cannot be read, which leaves the + /// last observation standing. + pub fn sample_and_publish( + &self, + ) -> Option<(ProcessResidentSampleV1, ResidentMemoryPressureStateV1)> { + let sample = (self.sampler)()?; + Some(( + sample, + self.publish_observed_resident_bytes(sample.unreclaimable_bytes), + )) + } + + /// [`Self::sample_and_publish`] reduced to the admission bytes: the + /// post-reclaim observation, or zero when the process cannot be read. + pub fn measure_admission_bytes(&self) -> u64 { + self.sample_and_publish().map_or(0, |(sample, state)| { + state.observed_bytes().unwrap_or(sample.unreclaimable_bytes) + }) + } + #[must_use] #[hotpath::skip] pub const fn limit_bytes(&self) -> u64 { @@ -685,6 +751,7 @@ pub fn process_resident_memory_pressure_v1() -> &'static Arc 0, - "a corrupt retained graph must defer its typed direct-recovery failure to the \ - canonical serialized activation path" + assert_eq!( + sealed_decode_count, 0, + "a corrupt retained graph is rebuilt from its sealed segments and seated from the \ + text owner, without decoding the generation" ); } else if !dirty_before_restart { assert_eq!( @@ -1495,6 +1495,154 @@ async fn restart_seats_the_retained_graph_while_its_text_owner_still_projects() .expect("join graph reconciliation tasks"); } +/// A first index reaches `ready` (fresh, graph serving) with the graph head +/// seated from the sealed manifest and the mapped graph store alone: no +/// reader asked for the whole decoded generation, so none is decoded, and +/// graph reads and the generation census answer from the text owner. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn first_index_serves_graph_reads_without_decoding_the_generation() { + let fixture = GitFixture::new(ALPHA_LIB_V1); + let store = TempDir::new().expect("store root"); + let profile = TempDir::new().expect("profile root"); + let profile_root = profile.path().join("profile"); + let project_id = test_project_id(); + tracedecay_runtime_core::storage::pin_fixture_repository_identity( + fixture.path(), + project_id.as_str(), + ) + .expect("project enrollment"); + let identity = tracedecay_daemon_identity::profile_identity::load_or_create(&profile_root) + .expect("profile identity"); + let _database_scope = tracedecay_runtime_core::db::enter_daemon_database_scope( + &profile_root, + 97, + "first index graph serving without decode", + ) + .expect("daemon database scope"); + let graph_runtime = Arc::new( + DaemonSessionRuntimeRegistryV1::open(identity) + .await + .expect("graph runtime registry"), + ); + let project_database = graph_runtime + .project_memory(project_id.clone(), [fixture.path().to_path_buf()]) + .await + .expect("writable project database"); + tracedecay_project::test_support::host_admission::await_bound_graph_runtime( + &project_database, + "bind first index graph runtime", + ) + .await + .expect("bound project graph runtime"); + + let registry = CodeIndexSchedulerRegistryV1::with_background_reconcile_permits(1, 1); + registry + .mount_worktree_with_graph_runtime( + project_id.clone(), + fixture.path(), + store.path().to_path_buf(), + graph_runtime.code_graph_seat_port(), + project_database, + CodeGraphActivationPolicyV1::Enabled, + ) + .await + .expect("mount first index"); + let reached = registry + .wait_for_readiness( + fixture.path(), + tracedecay_contracts::code_index_freshness::CodeIndexReadinessTargetV1::Ready, + Duration::from_secs(30), + ) + .await + .expect("readiness read"); + assert!( + matches!( + reached, + tracedecay_contracts::code_index_freshness::CodeIndexReadinessWaitReadV1::Reached + ), + "the first index must reach fresh with graph serving: {reached:?}" + ); + let scheduler = registry + .scheduler_handle(fixture.path()) + .await + .expect("mounted scheduler"); + assert_eq!( + scheduler + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .sealed_decode_count(), + 0, + "no reader demanded the whole generation, so it is never decoded" + ); + + let text = registry + .retained_text_owner_for_root(fixture.path()) + .await + .expect("text owner"); + let snapshot = text.metadata().snapshot(); + let scope = ResolvedScope::new( + project_id, + snapshot.repository.clone(), + snapshot.worktree.clone().expect("worktree identity"), + snapshot.reference.clone(), + ) + .expect("resolved scope"); + let generation_id = text.metadata().manifest().generation_id.clone(); + let statistics = text.metadata().generation_statistics().clone(); + let port = project_code_graph_projection_read_port( + registry.clone(), + fixture.path().to_path_buf(), + scope.clone(), + ); + let context = graph_request_context(scope.clone(), "first-index"); + let read = port + .open(CodeGraphReadRequest::from_context(&context, now_micros())) + .await + .expect("graph read on the first index"); + assert_eq!(read.freshness(), CodeGraphReadFreshnessV1::Current); + assert_eq!(read.generation(), &generation_id); + let page = read + .reader(&context, now_micros()) + .expect("admitted graph reader") + .symbols_page(None, 16, request_graph_cancellation(&context)) + .expect("bounded symbol query"); + assert_eq!( + page.symbols + .iter() + .filter_map(|symbol| symbol.metadata.as_ref()) + .map(|metadata| metadata.simple_name.as_str()) + .collect::>(), + ["alpha"] + ); + let census = + project_code_index_generation_census_reader(registry.clone(), fixture.path().into(), scope); + assert_eq!( + census().await, + GenerationCensusSnapshot::Observed { + generation_id: generation_id.as_str().to_owned(), + freshness: GenerationCensusServingFreshness::Current, + statistics: tracedecay_runtime_core::runtime_telemetry::GenerationCensusStatistics { + source_total_bytes: 28, + symbol_count: statistics.symbol_count, + edge_count: statistics.edge_count, + }, + } + ); + assert_eq!( + scheduler + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .sealed_decode_count(), + 0, + "graph reads and the census do not decode the generation" + ); + registry.shutdown().await; + graph_runtime + .shutdown_memory_graph_reconciliation_tasks() + .await + .expect("join graph reconciliation tasks"); +} + fn graph_request_context(scope: ResolvedScope, suffix: &str) -> RequestContext { let capability = CapabilityId::new("capability.code-graph-status-projection").expect("capability"); diff --git a/crates/tracedecay/src/daemon/maintenance.rs b/crates/tracedecay/src/daemon/maintenance.rs index 1219c41343..6c2da71ba5 100644 --- a/crates/tracedecay/src/daemon/maintenance.rs +++ b/crates/tracedecay/src/daemon/maintenance.rs @@ -841,8 +841,8 @@ impl MaintenanceCoordinator { } } -/// Samples this process's current resident set size, republishes it as a -/// Hotpath gauge, and feeds it to the resident-memory admission authority. +/// Samples this process's resident set, republishes it as Hotpath gauges, and +/// feeds its unreclaimable bytes to the resident-memory admission authority. /// /// A 20G RSS overrun past the admission limit was visible only to `ps` during /// a 2026-08 incident; the dedicated sampler closes that gap on a short @@ -850,19 +850,19 @@ impl MaintenanceCoordinator { /// [`process_resident_memory_pressure_v1`](tracedecay_runtime_core::resident_memory::process_resident_memory_pressure_v1) /// closes the loop: admission stops trusting its reservation model once the /// measurement says the process is over budget. The post-reclaim observation -/// returned by the pressure cell is the authority for the gauge and logs. +/// returned by the pressure cell is the authority for the unreclaimable gauge +/// and logs. #[cfg(target_os = "linux")] fn record_process_resident_memory_gauge(log: &std::sync::Mutex) { use tracedecay_runtime_core::resident_memory::ResidentMemoryPressureStateV1; - let Some(bytes) = tracedecay_runtime_core::resident_memory::sampled_process_resident_bytes_v1() - else { + let pressure = tracedecay_runtime_core::resident_memory::process_resident_memory_pressure_v1(); + let Some((sample, state)) = pressure.sample_and_publish() else { return; }; - let pressure = tracedecay_runtime_core::resident_memory::process_resident_memory_pressure_v1(); - let state = pressure.publish_observed_resident_bytes(bytes); + hotpath::gauge!("daemon.process.resident_bytes").set(sample.resident_bytes); if let Some(observed_bytes) = state.observed_bytes() { - hotpath::gauge!("daemon.process.resident_bytes").set(observed_bytes); + hotpath::gauge!("daemon.process.unreclaimable_resident_bytes").set(observed_bytes); } let over_budget = matches!(state, ResidentMemoryPressureStateV1::OverBudget { .. }); let transition = { diff --git a/crates/tracedecay/src/daemon/production_harness.rs b/crates/tracedecay/src/daemon/production_harness.rs index 1e4de073ca..646a77bc3a 100644 --- a/crates/tracedecay/src/daemon/production_harness.rs +++ b/crates/tracedecay/src/daemon/production_harness.rs @@ -1048,26 +1048,20 @@ async fn wait_for_production_composition_code_index( .latest_complete_ready_for_scope(scope) .await .is_some(); - // A clean restart whose retained revision-7 head recovered serves - // every read through the text projection and deliberately leaves - // the sealed seat empty, because replaying the partitions to seat - // a second copy of what already serves is the cost that recovery - // exists to avoid. No publication edge follows a quiet checkout, - // so waiting for that seat would always exhaust the timeout. - // - // The native graph is the discriminator, not text readiness - // alone: a *publishing* pass activates the graph only after the - // serving swap, so an already-serving graph over an empty seat is - // the recovered restart and nothing else. Accepting bare text - // readiness here would let the first open race ahead of its own - // seat, and every consumer that needs the decoded generation - // would then find nothing seated. - let recovered_text_ready = invocation + // A publication and a clean restart both seat the graph head on + // the text owner and leave the decoded seat empty until a reader + // demands it, so a text owner serving the native graph is ready. + // Graph serving installs before its interactive catalog finishes + // warming in the background; the composition is handed over only + // once catalog-dependent reads answer instead of reporting warming. + let text_graph_catalog_warm = invocation .code_index_schedulers .latest_text_serving_for_scope(scope) .await - .is_some_and(|text| text.interactive_graph_store().is_ok()); - if (generation_ready || recovered_text_ready) + .and_then(|text| text.interactive_graph_store().ok()) + .map(|store| store.interactive_catalog_is_warm().unwrap_or(false)); + if (generation_ready || text_graph_catalog_warm.is_some()) + && text_graph_catalog_warm != Some(false) && invocation .code_index_schedulers .query_authority_for_scope(scope) diff --git a/crates/tracedecay/src/daemon/project_composition/code_index_activation.rs b/crates/tracedecay/src/daemon/project_composition/code_index_activation.rs index dd1545388c..d521c861e4 100644 --- a/crates/tracedecay/src/daemon/project_composition/code_index_activation.rs +++ b/crates/tracedecay/src/daemon/project_composition/code_index_activation.rs @@ -148,17 +148,6 @@ fn spawn_query_authority_when_generation_ready(inputs: QueryAuthorityWaitInputs) let mut seats = authority_invocation .code_index_schedulers .subscribe_serving_seats(); - // A restored generation seats its graph only for a reader that - // demands the complete generation. The query authority is that - // reader: a route whose full upgrade is refused (a reset-required - // session store) has no advisory owner to make the demand. - tokio::select! { - biased; - () = authority_cancellation.cancelled() => return, - latest = authority_invocation - .code_index_schedulers - .latest_complete_ready_for_scope(&authority_scope) => drop(latest), - } let generation_ready = loop { if authority_invocation .code_index_schedulers diff --git a/crates/tracedecay/src/daemon/project_open_owners.rs b/crates/tracedecay/src/daemon/project_open_owners.rs index 6d2ee9fe9a..65537b40f9 100644 --- a/crates/tracedecay/src/daemon/project_open_owners.rs +++ b/crates/tracedecay/src/daemon/project_open_owners.rs @@ -617,18 +617,17 @@ pub(super) async fn register_project_open_production_owners( let mut mounted_providers = Vec::new(); let mut lsp_session_factory = None; let diagnostic_broker = server.diagnostics_lsp(); - // Bounds the orchestration wait around the scheduler's generation decode; - // the decode itself is instrumented inside the code-index subsystem. + // The census is the sealed manifest's file list; it never decodes. let indexed_generation = hotpath::future!( invocation .code_index_schedulers - .latest_complete_ready_decoded_for_root_scope(project_root, &scope), + .latest_feedback_generation_for_scope(project_root, &scope), label = "daemon.project.open.owners.lsp_census" ) .await; if let Some(generation) = indexed_generation { let mut indexed_files = generation - .generation() + .metadata() .snapshot() .files .iter() diff --git a/crates/tracedecay/src/daemon/project_open_owners/advisory_runtime/deferred.rs b/crates/tracedecay/src/daemon/project_open_owners/advisory_runtime/deferred.rs index 5c89f011ed..ebede9f2c6 100644 --- a/crates/tracedecay/src/daemon/project_open_owners/advisory_runtime/deferred.rs +++ b/crates/tracedecay/src/daemon/project_open_owners/advisory_runtime/deferred.rs @@ -60,20 +60,14 @@ pub(super) fn spawn( let mut serving_changes = None; let mut partial_publication_retried = false; loop { - // Sealing announces durable source before the complete serving - // owner is installed. Subscribe before probing that owner so a - // later serving swap can finish this mount without another edit. + // Sealing announces durable source before its text owner is + // installed. Subscribe before probing that owner so a later + // installation can finish this mount without another edit. if serving_changes.is_none() { serving_changes = invocation .code_index_schedulers .subscribe_serving_generation_changes(&project_root) .await; - if serving_changes.is_some() { - let _ = invocation - .code_index_schedulers - .request_complete_generation(&project_root) - .await; - } } match try_mount(&invocation, &project_root, &mut state).await { Attempt::Terminal => return, @@ -324,32 +318,14 @@ async fn classify_failure( Attempt::RetryPartialPublication } else if invocation .code_index_schedulers - .latest_complete_ready_for_scope(&state.scope) - .await - .is_none() - && invocation - .code_index_schedulers - .latest_text_serving_for_scope(&state.scope) - .await - .is_none() - { - Attempt::AwaitNextPublication - } else if invocation - .code_index_schedulers - .latest_complete_ready(project_root) + .latest_feedback_generation_for_scope(project_root, &state.scope) .await .is_none() { - // `try_mount` also admits the recovered text-serving level, but the - // feedback cycle it then composes mints its provider identity through - // `ProductionFeedbackDocumentIdentityPort`, which serves only - // `latest_complete_ready` for the exact root. When the text projection - // is ahead of that authority the composition fails with "project-open - // provider code-index identity is inconsistent with the application - // contract", earliness, not a missing composition. Classifying it - // terminal abandoned the upgrade for the daemon's whole life: the - // project kept the typed-unavailable feedback cycle and the warming - // LSP owner that advertises no analyzer method at all. + // The feedback cycle mints its provider identity from the same + // selection. While it answers nothing for the exact root the + // composition failed early, not for want of a composition; a + // terminal verdict here abandoned the upgrade for the daemon's life. Attempt::AwaitNextPublication } else { // A serving generation exists and no feedback cycle was published, so diff --git a/crates/tracedecay/src/daemon/project_open_owners/compiler_diagnostics_producer.rs b/crates/tracedecay/src/daemon/project_open_owners/compiler_diagnostics_producer.rs index f6462583a4..90f4cc2429 100644 --- a/crates/tracedecay/src/daemon/project_open_owners/compiler_diagnostics_producer.rs +++ b/crates/tracedecay/src/daemon/project_open_owners/compiler_diagnostics_producer.rs @@ -86,9 +86,6 @@ pub(super) fn spawn_typescript_diagnostics_producer( serving_changes = schedulers .subscribe_serving_generation_changes(&project_root) .await; - if serving_changes.is_some() { - let _ = schedulers.request_complete_generation(&project_root).await; - } } // The sealed generation, not the serving one: tsc reads only // the source the seal proved, so the check starts at the seal