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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<u64, RetrievalPortError> {
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)
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
51 changes: 51 additions & 0 deletions crates/tracedecay-runtime-core/src/resident_memory.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1304,6 +1304,23 @@ pub struct ResidentMemorySnapshotV1 {
pub process_shared_charges: Vec<ProcessSharedMemoryChargeV1>,
}

/// 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,
Expand Down Expand Up @@ -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()
Expand Down
67 changes: 67 additions & 0 deletions crates/tracedecay-runtime-core/src/resident_memory/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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));
}
Loading