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 4fa64fa4ac..8058a53825 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 @@ -2422,14 +2422,12 @@ impl DaemonCodeIndexPublicationStoreV1 { None => return Ok(None), }, }; - let pressure = admission.resident_memory.pressure(); - let watermark = pressure - .high_watermark_bytes() - .min(admission.resident_memory.snapshot().limit_bytes); + let watermark = admission.resident_memory.admission_watermark_bytes(); let admissible = || -> Result<(), String> { - let used = admission.resident_memory.snapshot().used_bytes; - let observed = pressure.measure_admission_bytes(); - let available = watermark.saturating_sub(used.max(observed)); + let available = admission + .resident_memory + .headroom_below(watermark) + .available_bytes; if requested.get() <= available { Ok(()) } else { diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/reconcile.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/reconcile.rs index 749920c331..79e0c0831e 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/reconcile.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/reconcile.rs @@ -1431,20 +1431,13 @@ impl CodeIndexWorktreeSchedulerV1 { .load_active_shared() .map_err(CodeIndexProductionErrorV1::Publication)?; let planned_workers = tracedecay_code_index::parallelism::indexing_workers(); - let snapshot = self.resident_memory.snapshot(); - // Admission refuses a request larger than the measured headroom, so - // the slab is planned against that too: the ledger alone does not see - // live state no owner charges, and a width planned from it asks for - // more than the process has left and is refused on every retry. The - // decoded parent is such state, so the measurement is taken now. - let pressure = self.resident_memory.pressure(); - let measured_remaining = pressure - .limit_bytes() - .saturating_sub(pressure.measure_admission_bytes()); - let remaining = snapshot - .limit_bytes - .saturating_sub(snapshot.used_bytes) - .min(measured_remaining); + // Planned from the view admission refuses against, taken after the + // parent decoded: a width the process cannot hold is refused on every + // retry. + let remaining = self + .resident_memory + .headroom_below(u64::MAX) + .available_bytes; // The process-global worker plan may have been installed against a // larger authority (standalone seed using detected host RAM). This // scheduler's remaining bytes are a different authority: the 6 GiB 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 3f77f4bae2..2940fff981 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 @@ -972,26 +972,18 @@ pub(super) fn text_artifact_resident_memory_charges( pub(super) fn text_artifact_admitted_build_budget( preferred_bytes: u64, minimum_bytes: u64, - limit_bytes: u64, - used_bytes: u64, - observed_bytes: u64, - watermark_headroom: u64, + available_bytes: u64, ) -> Result { if minimum_bytes == 0 || preferred_bytes < minimum_bytes { return Err(RetrievalPortError::Contract( "text-artifact build budget bounds are invalid".to_owned(), )); } - let unmodeled_live_bytes = observed_bytes.saturating_sub(used_bytes); - let available_for_growth = limit_bytes - .saturating_sub(used_bytes) - .saturating_sub(unmodeled_live_bytes) - .saturating_sub(watermark_headroom); - let admitted_bytes = preferred_bytes.min(available_for_growth); + let admitted_bytes = preferred_bytes.min(available_bytes); if admitted_bytes < minimum_bytes { return Err(RetrievalPortError::ResidentMemoryRefused(format!( "text-artifact build needs at least {minimum_bytes} bytes; \ - {available_for_growth} bytes are available below the resident-memory watermark" + {available_bytes} bytes are available below the resident-memory watermark" ))); } Ok(admitted_bytes) @@ -1199,22 +1191,19 @@ impl DaemonCodeTextArtifactStoreV1 { preferred: NonZeroU64, minimum: NonZeroU64, ) -> Result<(u64, u64, u64, u64), RetrievalPortError> { - let snapshot = self.resident_memory.snapshot(); - 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 + let admission_watermark = self.resident_memory.admission_watermark_bytes(); + let headroom = self.resident_memory.headroom_below(admission_watermark); + let observed_bytes = headroom.observed_bytes; + let unmodeled_live_bytes = observed_bytes.saturating_sub(headroom.used_bytes); + let watermark_headroom = self .resident_memory - .pressure() - .high_watermark_bytes() - .min(snapshot.limit_bytes); - let watermark_headroom = snapshot.limit_bytes.saturating_sub(admission_watermark); + .snapshot() + .limit_bytes + .saturating_sub(admission_watermark); let admitted_bytes = text_artifact_admitted_build_budget( preferred.get(), minimum.get(), - snapshot.limit_bytes, - snapshot.used_bytes, - observed_bytes, - watermark_headroom, + headroom.available_bytes, ) .inspect_err(|_| { self.resident_memory diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/serving.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/serving.rs index 83cf903f00..21be518d38 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/serving.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/serving.rs @@ -2097,38 +2097,23 @@ fn overlapping_text_builds_share_one_admission_watermark_headroom() { fn text_build_budget_shrinks_to_available_headroom_without_dropping_below_its_floor() { const GIB: u64 = 1024 * 1024 * 1024; const MIB: u64 = 1024 * 1024; - let limit = 26 * GIB; - let preferred = limit / 8; + let preferred = 26 * GIB / 8; let minimum = 1536 * MIB; - let watermark_headroom = limit - (limit * 900 / 1000); - let observed = 21 * GIB; - let available = limit - observed - watermark_headroom; assert_eq!( - super::super::text_artifact_admitted_build_budget( - preferred, - minimum, - limit, - 0, - observed, - watermark_headroom, - ), - Ok(available), + super::super::text_artifact_admitted_build_budget(preferred, minimum, 2 * GIB), + Ok(2 * GIB), "a replacement build must use the supported smaller budget instead of deadlocking behind the stale graph" ); assert_eq!( - super::super::text_artifact_admitted_build_budget( - preferred, - minimum, - limit, - 0, - 22 * GIB, - watermark_headroom, - ), + super::super::text_artifact_admitted_build_budget(preferred, minimum, 8 * GIB), + Ok(preferred), + ); + assert_eq!( + super::super::text_artifact_admitted_build_budget(preferred, minimum, GIB), Err( tracedecay_query::retrieval::RetrievalPortError::ResidentMemoryRefused(format!( - "text-artifact build needs at least {minimum} bytes; {} bytes are available below the resident-memory watermark", - limit - 22 * GIB - watermark_headroom + "text-artifact build needs at least {minimum} bytes; {GIB} bytes are available below the resident-memory watermark" )) ), "less than the builder's supported floor must remain a typed capacity refusal" diff --git a/crates/tracedecay-runtime-core/src/resident_memory.rs b/crates/tracedecay-runtime-core/src/resident_memory.rs index 192fe73bce..df7db21c75 100644 --- a/crates/tracedecay-runtime-core/src/resident_memory.rs +++ b/crates/tracedecay-runtime-core/src/resident_memory.rs @@ -1304,6 +1304,23 @@ pub struct ResidentMemorySnapshotV1 { pub process_shared_charges: Vec, } +/// One reading of the process's room for new work. +/// +/// The ledger charges the retained owners and the work in flight; it does not +/// see the runtime's own allocations or the allocator's retained pages, which +/// are not yet small and fixed enough to charge as a constant. The measured +/// process covers them, so new work has to fit beside both. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct ResidentMemoryHeadroomV1 { + /// Bytes the ledger charges. + pub used_bytes: u64, + /// Bytes measured for the process; zero when it cannot be read. + pub observed_bytes: u64, + /// What fits below the ceiling beside both; zero while the pressure + /// latch holds. + pub available_bytes: u64, +} + #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub struct ProcessSharedMemoryChargeV1 { pub component: ResidentMemoryComponentIdV1, @@ -1552,6 +1569,40 @@ impl ProcessResidentMemoryV1 { } } + /// The process's room for new work below `ceiling`, from a fresh + /// measurement: what every admission sizes from. The ledger is held to + /// this authority's limit and the measurement to the pressure cell's, as + /// [`Self::reserve`] holds them. + pub fn headroom_below(&self, ceiling: u64) -> ResidentMemoryHeadroomV1 { + let observed_bytes = self.pressure.measure_admission_bytes(); + let used_bytes = self.lock_state().used_bytes; + let available_bytes = if self.pressure.state().is_over_budget() { + 0 + } else { + let ledger_room = ceiling + .min(self.limit_bytes.get()) + .saturating_sub(used_bytes); + let measured_room = ceiling + .min(self.pressure.limit_bytes()) + .saturating_sub(observed_bytes); + ledger_room.min(measured_room) + }; + ResidentMemoryHeadroomV1 { + used_bytes, + observed_bytes, + available_bytes, + } + } + + /// The ceiling decode and artifact builds size below: the pressure high + /// watermark, leaving the band above it for work already admitted. + #[must_use] + pub fn admission_watermark_bytes(&self) -> u64 { + self.pressure + .high_watermark_bytes() + .min(self.limit_bytes.get()) + } + fn lock_state(&self) -> ProfiledMutexGuard<'_, ResidentMemoryStateV1> { self.state .lock() diff --git a/crates/tracedecay-runtime-core/src/resident_memory/tests.rs b/crates/tracedecay-runtime-core/src/resident_memory/tests.rs index 9792efc0e0..d38c7682ef 100644 --- a/crates/tracedecay-runtime-core/src/resident_memory/tests.rs +++ b/crates/tracedecay-runtime-core/src/resident_memory/tests.rs @@ -1442,3 +1442,70 @@ fn checkpoints_read_the_process_once_per_interval_and_then_see_growth() { ); assert_eq!(reads.load(Ordering::Acquire), reads_taken + 1); } + +/// Every admission sizes from one view: the ceiling less the larger of the +/// ledger and the measured process, and nothing while the pressure latch +/// holds, so a width planned from it is one admission accepts. +#[test] +fn headroom_is_the_ceiling_less_the_larger_of_ledger_and_measurement() { + const MIB: u64 = 1024 * 1024; + let measured = Arc::new(AtomicU64::new(300 * MIB)); + let sampled = Arc::clone(&measured); + let authority = Arc::new(ProcessResidentMemoryV1::with_pressure( + bytes(1000 * MIB), + Arc::new(ResidentMemoryPressureV1::with_sampler( + bytes(1000 * MIB), + Arc::new(move || { + let observed = sampled.load(Ordering::SeqCst); + Some(ProcessResidentSampleV1 { + resident_bytes: observed, + unreclaimable_bytes: observed, + swapped_bytes: 0, + cgroup_committed_bytes: None, + }) + }), + )), + )); + let owner = key("project-a", "worktree-a", "generation-a", "reader"); + let charged = authority + .reserve(owner.clone(), bytes(100 * MIB)) + .expect("charged owner"); + + let headroom = authority.headroom_below(u64::MAX); + assert_eq!( + ( + headroom.used_bytes, + headroom.observed_bytes, + headroom.available_bytes + ), + (100 * MIB, 300 * MIB, 700 * MIB), + "state no owner charges counts through the measurement" + ); + assert_eq!( + authority.headroom_below(500 * MIB).available_bytes, + 200 * MIB + ); + assert!( + authority.reserve(owner.clone(), bytes(701 * MIB)).is_err(), + "admission refuses past the same view" + ); + + let in_flight = authority + .reserve(owner, bytes(300 * MIB)) + .expect("in-flight work"); + assert_eq!( + authority.headroom_below(u64::MAX).available_bytes, + 600 * MIB, + "a ledger above the measurement counts in full" + ); + + measured.store(990 * MIB, Ordering::SeqCst); + assert_eq!(authority.headroom_below(u64::MAX).available_bytes, 0); + measured.store(800 * MIB, Ordering::SeqCst); + assert_eq!( + authority.headroom_below(u64::MAX).available_bytes, + 0, + "the latch holds until the low watermark" + ); + drop((charged, in_flight)); +}