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
12 changes: 11 additions & 1 deletion cli/src/gen/engram/app/v1/fleet_pb.ts

Large diffs are not rendered by default.

4 changes: 4 additions & 0 deletions crates/engram-coordinator/src/api/hosts.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,8 @@ pub(crate) async fn enrichment_for_session(

#[derive(Serialize)]
pub struct HostView {
pub nbd_slots_total: u32,
pub nbd_slots_in_use: u32,
pub id: HostId,
pub hostname: String,
pub status: &'static str,
Expand Down Expand Up @@ -163,6 +165,8 @@ impl HostView {
util_base_shm_mib: row.utilization.base_shm_mib,
util_parked_pss_mib: row.utilization.parked_pss_mib,
util_running_pss_mib: row.utilization.running_pss_mib,
nbd_slots_total: row.utilization.nbd_slots_total,
nbd_slots_in_use: row.utilization.nbd_slots_in_use,
allocatable_mib,
reserved_mib,
free_mib: allocatable_mib.saturating_sub(reserved_mib),
Expand Down
4 changes: 4 additions & 0 deletions crates/engram-coordinator/src/api/sessions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -640,6 +640,7 @@ async fn boot_prepared(
oauth_credential: Option<engram_core::types::oauth::OAuthCredentialKey>,
) -> Result<CreateSessionResponse, ApiError> {
let crate::session_boot::PreparedBoot {
nbd_slot_need,
inputs,
memory_mib,
cpu_budget_vcpus,
Expand All @@ -666,6 +667,7 @@ async fn boot_prepared(
// a capacity miss does. A straggler host that hasn't staged yet simply
// isn't in the ranked pool; its next heartbeat un-gates it.
let ctx = crate::placement::ScheduleContext {
nbd_slot_need,
repo: &image_repo,
image_version: &image_tag,
// ADR 0078: a base snapshot is fleet-wide (prewarmed on many
Expand Down Expand Up @@ -744,6 +746,7 @@ async fn boot_prepared(
};

let write_set = engram_core::traits::SessionCreateWriteSet {
nbd_slot_need,
session_id,
spec: inputs.spec.clone(),
mem_budget_mib: memory_mib as i64,
Expand Down Expand Up @@ -1494,6 +1497,7 @@ async fn prepare_inner(
let manifest_digest = bundle.enabled.manifest_digest.clone();

Ok(crate::session_boot::PreparedBoot {
nbd_slot_need: 1 + u32::from(config.resolved_swap_mib() > 0),
inputs: crate::session_boot::BootInputs {
session_id,
spec,
Expand Down
2 changes: 2 additions & 0 deletions crates/engram-coordinator/src/api/snapshot.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1165,6 +1165,7 @@ async fn resume_disk_only_cold_boot(
spec.rootfs_manifest = session.live_disk_manifest;
let (repo, tag) = engram_core::types::session::split_image_ref(&session.image);
let context = crate::placement::ScheduleContext {
nbd_slot_need: 1 + u32::from(spec.swap_mib.unwrap_or(0) > 0),
repo,
image_version: tag,
snapshot_host: None,
Expand Down Expand Up @@ -1678,6 +1679,7 @@ async fn resume_from_fc_snapshot(
let resume_budget =
crate::boot_materializer::resolve_resume_budget(&state.services.meta, &session).await;
let ctx = ScheduleContext {
nbd_slot_need: 1 + u32::from(record.swap_manifest.is_some()),
repo: image_repo,
image_version: image_tag,
// ADR 0078: authoritative affinity — the host that holds this
Expand Down
4 changes: 4 additions & 0 deletions crates/engram-coordinator/src/grpc_app/convert.rs
Original file line number Diff line number Diff line change
Expand Up @@ -418,6 +418,8 @@ pub(crate) fn host_view_to_proto(v: &crate::api::hosts::HostView) -> app::HostVi
cordoned,
cordon_owner,
retirement,
nbd_slots_total,
nbd_slots_in_use,
allocatable_mib,
reserved_mib,
free_mib,
Expand Down Expand Up @@ -457,6 +459,8 @@ pub(crate) fn host_view_to_proto(v: &crate::api::hosts::HostView) -> app::HostVi
.map(|o| o.as_str().to_owned())
.unwrap_or_default(),
retirement: retirement.as_ref().map(retirement_to_proto),
nbd_slots_total: *nbd_slots_total,
nbd_slots_in_use: *nbd_slots_in_use,
allocatable_mib: *allocatable_mib,
reserved_mib: *reserved_mib,
free_mib: *free_mib,
Expand Down
86 changes: 83 additions & 3 deletions crates/engram-coordinator/src/placement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ use crate::host_registry::HostRegistry;
/// the backing store moved from the in-memory mirror to the hosts rows.)
#[derive(Clone, Debug)]
pub struct ScheduleContext<'a> {
pub nbd_slot_need: u32,
pub repo: &'a str,
pub image_version: &'a str,
/// ADR 0078 (GCS-free resume) tier-0 authoritative affinity: the host
Expand Down Expand Up @@ -543,7 +544,16 @@ pub fn pick_from(
pick_from_2d(hosts, reserved, ctx, now, ttl, false)
}

/// The ranked 2D pick, with an explicit `require_fit` knob.
fn nbd_slots_fit(h: &HostRecord, reserved: &HashMap<HostId, ReservedBudget>, need: u32) -> bool {
h.utilization.nbd_slots_total == 0
|| i64::from(h.utilization.nbd_slots_total)
- i64::from(h.utilization.nbd_slots_in_use)
- reserved.get(&h.id).map(|r| r.nbd_slots).unwrap_or(0)
>= i64::from(need)
}

/// The ranked RAM/CPU pick, with an explicit `require_fit` knob.
/// NBD slot capacity is a hard gate in every mode.
///
/// `require_fit=false` (the historical `pick_from` behavior) keeps the
/// ADR 0046 capacity-SOFT last-resort fallback: when nothing fits both
Expand Down Expand Up @@ -574,6 +584,17 @@ pub fn pick_from_2d(
None => Err(PickError::NoCapacity),
};
}
// Slots are a hard gate, including affinity and unmeasured-RAM hosts.
let slot_hosts: Vec<_> = hosts
.iter()
.filter(|h| nbd_slots_fit(h, reserved, ctx.nbd_slot_need))
.cloned()
.collect();
let hosts = slot_hosts.as_slice();
let ranked = rank_hosts(hosts, ctx, now, ttl);
if ranked.hosts.is_empty() {
return Err(PickError::NoCapacity);
}
let need_mib = ctx.memory_mib.unwrap_or(0) as i64;
let need_vcpus = ctx.cpu_budget_vcpus.unwrap_or(0) as i64;
// Free RAM for `id`: None ⇒ unmeasured (treated as "fits" softly).
Expand Down Expand Up @@ -803,6 +824,9 @@ pub async fn placement_preview(
return false;
};
let alloc = h.utilization.allocatable_mib as i64;
if !nbd_slots_fit(h, &reserved, ctx.nbd_slot_need) {
return false;
}
if alloc <= 0 {
return true; // unmeasured → soft fallback fits
}
Expand Down Expand Up @@ -907,7 +931,12 @@ pub async fn log_reserve_no_fit(
) -> Vec<(HostId, String)> {
let mut summary: Vec<(HostId, String)> = Vec::new();
match meta
.placement_no_fit_details(candidates, mem_budget_mib, cpu_budget_vcpus)
.placement_no_fit_details(
candidates,
mem_budget_mib,
cpu_budget_vcpus,
ctx.nbd_slot_need,
)
.await
{
Ok(details) => {
Expand Down Expand Up @@ -1020,6 +1049,7 @@ pub async fn pick_specific_host(
registry: &HostRegistry,
host_id: HostId,
exclude_host: Option<HostId>,
nbd_slot_need: u32,
now: DateTime<Utc>,
) -> Result<(HostId, Arc<dyn HostClient>), PickError> {
if Some(host_id) == exclude_host {
Expand All @@ -1041,6 +1071,9 @@ pub async fn pick_specific_host(
if host_meets_capabilities(h, &CapabilityRequirements::default()).is_err() {
return Err(PickError::NoCapacity);
}
if !nbd_slots_fit(h, &reserved, nbd_slot_need) {
return Err(PickError::NoCapacity);
}
let alloc = h.utilization.allocatable_mib as i64;
let reserved_mib = reserved.get(&host_id).map(|r| r.mem_mib).unwrap_or(0);
if alloc > 0 && alloc - reserved_mib <= 0 {
Expand Down Expand Up @@ -1728,8 +1761,44 @@ mod tests {
}
}

#[test]
fn nbd_slots_are_hard_even_for_affinity_and_soft_fallback() {
let mut h = host(1);
h.utilization.nbd_slots_total = 2;
h.utilization.nbd_slots_in_use = 1;
let mut c = ctx();
c.nbd_slot_need = 2;
c.snapshot_host = Some(h.id);
c.prefer_host = Some(h.id);
let mut reserved = HashMap::new();
for require_fit in [false, true] {
assert!(
pick_from_2d(&[h.clone()], &reserved, &c, Utc::now(), TTL, require_fit).is_err()
);
}
c.nbd_slot_need = 1;
assert_eq!(
pick_from_2d(&[h.clone()], &reserved, &c, Utc::now(), TTL, false).unwrap(),
h.id
);
reserved.insert(
h.id,
ReservedBudget {
nbd_slots: 1,
..Default::default()
},
);
assert!(pick_from_2d(&[h.clone()], &reserved, &c, Utc::now(), TTL, false).is_err());
h.utilization.nbd_slots_total = 0;
assert_eq!(
pick_from_2d(&[h.clone()], &reserved, &c, Utc::now(), TTL, false).unwrap(),
h.id
);
}

fn ctx<'a>() -> ScheduleContext<'a> {
ScheduleContext {
nbd_slot_need: 1,
repo: "r",
image_version: "v",
snapshot_host: None,
Expand All @@ -1748,7 +1817,16 @@ mod tests {
fn mem_reserved(entries: &[(HostId, i64)]) -> HashMap<HostId, ReservedBudget> {
entries
.iter()
.map(|&(id, mem_mib)| (id, ReservedBudget { mem_mib, vcpus: 0 }))
.map(|&(id, mem_mib)| {
(
id,
ReservedBudget {
mem_mib,
vcpus: 0,
nbd_slots: 0,
},
)
})
.collect()
}

Expand Down Expand Up @@ -1977,13 +2055,15 @@ mod tests {
(
hid(1),
ReservedBudget {
nbd_slots: 0,
mem_mib: 0,
vcpus: 32,
},
),
(
hid(2),
ReservedBudget {
nbd_slots: 0,
mem_mib: 0,
vcpus: 0,
},
Expand Down
22 changes: 20 additions & 2 deletions crates/engram-coordinator/src/queue_scanner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@
//! ## Per-fit-class FIFO (queue-fairness follow-up to ADR 0048)
//!
//! [`FitClass`] partitions the FIFO-ordered queue purely by the 2D
//! `(mem_budget_mib, cpu_budget_vcpus)` pair ([`partition_queue`]) and each
//! `(mem_budget_mib, cpu_budget_vcpus, nbd_slot_need)` tuple ([`partition_queue`]) and each
//! class sweeps strict FIFO independently. Classes are attempted
//! oldest-head-first (the most-starved class gets first claim on freed
//! capacity this tick); a create that hits `NoCapacity` stops only its
Expand Down Expand Up @@ -149,6 +149,7 @@ fn env_secs(name: &str, default: u64) -> u64 {
/// blocking. See the module doc.
#[derive(Clone, Copy, Debug, Eq, PartialEq, Hash)]
struct FitClass {
nbd_slot_need: u32,
mem_budget_mib: i64,
cpu_budget_vcpus: i32,
}
Expand All @@ -158,6 +159,7 @@ impl FitClass {
Self {
mem_budget_mib: q.mem_budget_mib,
cpu_budget_vcpus: q.cpu_budget_vcpus,
nbd_slot_need: q.nbd_slot_need,
}
}
}
Expand Down Expand Up @@ -482,6 +484,7 @@ async fn place_create(state: &SharedState, q: &QueuedSession) -> PlaceOutcome {
.map(|e| e.base_snapshot_memory_manifest.is_some())
.unwrap_or(false);
let ctx = crate::placement::ScheduleContext {
nbd_slot_need: q.nbd_slot_need,
repo,
image_version: tag,
snapshot_host: None,
Expand Down Expand Up @@ -598,6 +601,7 @@ async fn resume_has_capacity(state: &SharedState, q: &QueuedSession) -> ResumeCa
.ok()
.flatten();
let ctx = crate::placement::ScheduleContext {
nbd_slot_need: q.nbd_slot_need,
repo,
image_version: tag,
snapshot_host: None,
Expand Down Expand Up @@ -836,7 +840,12 @@ async fn time_out_session(state: &SharedState, q: &QueuedSession) {
if let Ok(details) = state
.services
.meta
.placement_no_fit_details(&ids, q.mem_budget_mib, q.cpu_budget_vcpus)
.placement_no_fit_details(
&ids,
q.mem_budget_mib,
q.cpu_budget_vcpus,
q.nbd_slot_need,
)
.await
{
payload["hosts"] = details
Expand Down Expand Up @@ -928,9 +937,18 @@ mod tests {
use super::*;
use engram_core::types::session::{Session, SessionMode};

#[test]
fn slot_need_separates_queue_fit_classes() {
let root = queued(1, 1024, 1, 10);
let mut swap = root.clone();
swap.nbd_slot_need = 2;
assert_eq!(partition_queue(vec![swap, root]).len(), 2);
}

fn queued(id_seed: u8, mem: i64, cpu: i32, secs_ago: i64) -> QueuedSession {
let id = SessionId::new();
QueuedSession {
nbd_slot_need: 1,
session: Session {
id,
status: SessionState::Queued,
Expand Down
1 change: 1 addition & 0 deletions crates/engram-coordinator/src/session_boot.rs
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,7 @@ pub(crate) struct BootInputs {
/// its response). Built by `sessions::prepare_from_request` (from a live
/// request) or `sessions::prepare_from_row` (from a durable queued row).
pub(crate) struct PreparedBoot {
pub nbd_slot_need: u32,
pub inputs: BootInputs,
pub memory_mib: u32,
pub cpu_budget_vcpus: u32,
Expand Down
1 change: 1 addition & 0 deletions crates/engram-coordinator/src/teleport.rs
Original file line number Diff line number Diff line change
Expand Up @@ -674,6 +674,7 @@ mod steps {
};
let candidates_for = |caps: crate::placement::CapabilityRequirements| {
let context = crate::placement::ScheduleContext {
nbd_slot_need: 1,
repo: &session.image,
image_version: "",
snapshot_host: None,
Expand Down
2 changes: 2 additions & 0 deletions crates/engram-coordinator/tests/admin_evac_live_pg.rs
Original file line number Diff line number Diff line change
Expand Up @@ -422,6 +422,7 @@ async fn durable_cordon_excludes_host_from_placement_on_every_replica() {

// Both hosts available → either can be picked.
let ctx = ScheduleContext {
nbd_slot_need: 1,
repo: "test/img",
image_version: "v1",
snapshot_host: None,
Expand Down Expand Up @@ -497,6 +498,7 @@ async fn durable_cordon_excludes_host_from_placement_on_every_replica() {
.await
.expect("uncordon");
let exclude_healthy_ctx = ScheduleContext {
nbd_slot_need: 1,
exclude_host: Some(healthy),
..ctx.clone()
};
Expand Down
3 changes: 3 additions & 0 deletions crates/engram-coordinator/tests/capture_jobs_live_pg.rs
Original file line number Diff line number Diff line change
Expand Up @@ -235,6 +235,7 @@ async fn record_capture_job_report_is_fenced_by_epoch() {
epoch: job.epoch + 1,
stage: CaptureJobStage::Booting,
progress: Some(CaptureJobProgress {
sandbox_id: None,
detail: Some("should never land".into()),
log_tail: None,
warm_stages: Vec::new(),
Expand Down Expand Up @@ -262,6 +263,7 @@ async fn record_capture_job_report_is_fenced_by_epoch() {
epoch: job.epoch,
stage: CaptureJobStage::Booting,
progress: Some(CaptureJobProgress {
sandbox_id: None,
detail: Some("cold boot".into()),
log_tail: None,
warm_stages: Vec::new(),
Expand Down Expand Up @@ -377,6 +379,7 @@ async fn reassign_bumps_epoch_and_attempts_and_resets_stage() {
epoch: job.epoch,
stage: CaptureJobStage::Booting,
progress: Some(CaptureJobProgress {
sandbox_id: None,
detail: Some("booting".into()),
log_tail: None,
warm_stages: Vec::new(),
Expand Down
1 change: 1 addition & 0 deletions crates/engram-coordinator/tests/ha_listener.rs
Original file line number Diff line number Diff line change
Expand Up @@ -556,6 +556,7 @@ async fn cross_replica_scheduling_pins_and_tokens() {
registry_b.register(h2, backend.clone());

let ctx = ScheduleContext {
nbd_slot_need: 1,
repo: "r",
image_version: "v",
snapshot_host: None,
Expand Down
Loading
Loading