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