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
1 change: 1 addition & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -867,6 +867,7 @@ jobs:
--test diff_snapshot \
--test file_restore_shared_rss \
--test stock_fork_snapshot_compat \
--test pause_resume_vsock \
--test vmstate_only_snapshot \
--test substrate_kernel_capabilities \
--test substrate_uffd_base \
Expand Down
26 changes: 17 additions & 9 deletions crates/engram-coordinator/src/api/admin.rs
Original file line number Diff line number Diff line change
Expand Up @@ -222,9 +222,10 @@ pub(crate) async fn evacuate_session_core(
#[derive(Serialize)]
pub struct EvictIdleResponse {
pub session_id: SessionId,
/// "idle" — the session was paused, flushed, snapshotted, and its
/// local sandbox destroyed; the PG row is at `Idle` and a subsequent
/// `Resume` rebinds it.
/// The pipeline's actual outcome: "idle" (full suspend — paused,
/// flushed, snapshotted, sandbox destroyed, resume rebinds), the
/// ADR 0074 rung-2 "evicting (parked-paused, rung 2)" (VM paused in
/// place, un-parked by the next prompt), or a lease-race "skipped".
pub status: &'static str,
}

Expand Down Expand Up @@ -256,20 +257,27 @@ pub(crate) async fn evict_idle_core(
)));
};

crate::idle_evictor::evict_idle_session(state, session_id, sandbox_id)
let outcome = crate::idle_evictor::evict_idle_session(state, session_id, sandbox_id)
.await
.map_err(|e| ApiError::Internal(format!("idle-evict pipeline: {e}")))?;

// Report what actually happened — with ADR 0074 rung 2 the pipeline
// may PARK (pause in place, session Evicting + park_rung=2) instead
// of suspending to Idle, and a lease race skips entirely. The old
// hardcoded "idle" misreported both.
let status = match &outcome {
crate::idle_evictor::EvictOutcome::ParkedPaused => "evicting (parked-paused, rung 2)",
crate::idle_evictor::EvictOutcome::Skipped { .. } => "skipped (pipeline already in flight)",
_ => "idle",
};
tracing::info!(
%session_id,
%sandbox_id,
"admin evict-idle: session suspended to Idle via the idle-eviction primitive",
?outcome,
"admin evict-idle: idle-eviction primitive completed",
);

Ok(EvictIdleResponse {
session_id,
status: "idle",
})
Ok(EvictIdleResponse { session_id, status })
}

// ---------------------------------------------------------------------
Expand Down
41 changes: 26 additions & 15 deletions crates/engram-coordinator/src/api/snapshot.rs
Original file line number Diff line number Diff line change
Expand Up @@ -489,15 +489,14 @@ pub async fn ensure_active(state: &SharedState, id: SessionId) -> Result<(), Api
let session = state.services.meta.get_session(id).await?;
match session.status {
SessionState::Active => {
// ADR 0074 rung-2: a session can be Active AND parked-paused
// (`park_rung == 2`, VM frozen) — the admin `EvictIdle` path parks
// from Active, and `evict_session_to_state`'s park branch returns
// WITHOUT a status transition, so it inherits `Active`. Returning Ok
// would advertise a live session over a PAUSED VM and stall the next
// prompt on a frozen harness. Un-pause it first (the natural
// nomination path parks from `Evicting`, routing through the
// `try_cancel` ascent; this covers the Active/paused shape). No-op
// for the common not-parked Active session.
// ADR 0074 rung-2 backstop: parked-paused now uniformly means
// `Evicting` (the park branch transitions the admin path's
// Active entry too), so Active + `park_rung == 2` only occurs
// in the crash window between the host `pause` landing and the
// park's PG bookkeeping committing. Returning Ok would
// advertise a live session over a PAUSED VM and stall the next
// prompt on a frozen harness — un-pause it first. No-op for
// the common not-parked Active session.
if session.park_rung == 2 {
let _ = try_cancel_nominated_eviction(state, id).await?;
}
Expand Down Expand Up @@ -574,7 +573,13 @@ pub async fn ensure_active(state: &SharedState, id: SessionId) -> Result<(), Api
SessionState::Dead => Err(ApiError::Gone(
"session is dead — chunked manifests are gone or never existed".into(),
)),
SessionState::Failed | SessionState::Completed => Err(ApiError::Conflict(format!(
// Terminal is GONE, not Conflict: a 409 reads as "retry later" to
// every caller, and the outbox driver in particular deferred a
// completed session's un-acked rows every backoff-tick forever
// (prod: 3 rows spinning the driver for hours). Gone maps to the
// driver's Terminal arm — the row is dropped — and to an honest
// 410 for exec/upload/relay callers.
SessionState::Failed | SessionState::Completed => Err(ApiError::Gone(format!(
"session is {} (terminal); no work to dispatch",
session.status.as_str()
))),
Expand Down Expand Up @@ -626,6 +631,14 @@ pub(crate) async fn try_cancel_nominated_eviction(
return Ok(false);
}
}
// Clear the park stamp the moment the un-pause lands, NOT in the
// transition's Ok arm: the `ensure_active` Active-arm backstop
// un-parks a session that is already Active (the pause-then-crash
// window before the park's Evicting transition), and Active→Active
// below is a same-state Conflict — tying the clear to the Ok arm
// left `park_rung=2` advertised forever over a running VM.
let _ = state.services.meta.set_session_park_rung(id, 0, None).await;
::metrics::counter!(crate::metrics::EVICTION_UNPARKED_PAUSED_TOTAL).increment(1);
}
let result = match state
.services
Expand All @@ -634,11 +647,7 @@ pub(crate) async fn try_cancel_nominated_eviction(
.await
{
Ok(prev) => {
let _ = state.services.meta.set_session_park_rung(id, 0, None).await;
::metrics::counter!(crate::metrics::EVICTION_CANCELLED_TOTAL).increment(1);
if parked_paused {
::metrics::counter!(crate::metrics::EVICTION_UNPARKED_PAUSED_TOTAL).increment(1);
}
tracing::info!(
session_id = %id,
park_rung = if parked_paused { 2 } else { 1 },
Expand Down Expand Up @@ -2683,7 +2692,9 @@ mod evicting_gate_tests {
.expect_err("Completed session has no work to dispatch");
flipper.await.unwrap();

assert_eq!(err.status(), axum::http::StatusCode::CONFLICT);
// Terminal is GONE (410), not a retryable 409: the outbox driver
// maps Gone to its Terminal drop arm instead of deferring forever.
assert_eq!(err.status(), axum::http::StatusCode::GONE);
assert!(
err.to_string().contains("terminal"),
"must surface the terminal-state error, not the mid-eviction 409, got: {err}",
Expand Down
93 changes: 72 additions & 21 deletions crates/engram-coordinator/src/idle_evictor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,45 @@ pub enum EvictOutcome {
/// matching the historical idle-eviction shape. Operator-driven
/// drains (ADR 0018 commit 12) use [`evict_session_to_state`]
/// directly with `target_state = Evacuating`.
/// ADR 0074 rung 2: does `sandbox_id`'s host have memory headroom to
/// ADR 0074 rung-2 park bookkeeping, run after a successful `pause`:
/// stamp `park_rung=2`/`parked_at`, and make the STATUS say what the VM
/// is doing. The natural path enters already `Evicting` (nominated); the
/// admin `EvictIdle` path enters `Active`, and returning `ParkedPaused`
/// without a transition used to leave an ACTIVE session over a frozen VM
/// — `park_rung == 2` uniformly means `Evicting` now. (The `ensure_active`
/// Active-arm un-park stays as the backstop for the crash window between
/// the pause and this transition.)
async fn park_paused_bookkeeping(
state: &SharedState,
session_id: SessionId,
entry_status: SessionState,
) -> Result<(), engram_core::MetaError> {
state
.services
.meta
.set_session_park_rung(session_id, 2, Some(Utc::now()))
.await?;
if entry_status == SessionState::Active {
state
.services
.meta
.transition_session(session_id, SessionState::Evicting)
.await?;
let _ = state
.emit(
session_id,
crate::state::SessionEvent::StatusChanged {
from: SessionState::Active,
to: SessionState::Evicting,
at: Utc::now(),
},
)
.await;
}
Ok(())
}

/// ADR 0074 rung 2: does the session's host have memory headroom to
/// keep a VM PAUSED (rung 2) rather than fully evicting it? Reads the
/// heartbeat-persisted `hosts.utilization`. Fails CLOSED (no headroom →
/// full eviction) on a telemetry gap: a paused VM frees no RAM, so
Expand All @@ -85,8 +123,17 @@ pub enum EvictOutcome {
/// Headroom threshold reuses the detector's mem floor (default 15% free)
/// plus a margin, so a host that is not "under pressure" for eviction
/// purposes has room to hold a paused VM.
async fn host_has_memory_headroom(state: &SharedState, sandbox_id: SandboxId) -> bool {
let Some(host_id) = state.host_registry.host_of(sandbox_id) else {
async fn host_has_memory_headroom(state: &SharedState, session_id: SessionId) -> bool {
// Resolve the host through PG, NOT the in-memory `host_registry`: this
// runs on whichever coord replica took the RPC / scanner tick, and a
// replica that never cached the sandbox→host bind would fail closed
// here — silently degrading every park into a full eviction on that
// pod (a coin flip in a 2-replica deployment).
let host_id = match state.services.meta.get_session(session_id).await {
Ok(s) => s.host_id,
Err(_) => return false,
};
let Some(host_id) = host_id else {
return false;
};
let hosts = match state.services.meta.list_active_hosts().await {
Expand Down Expand Up @@ -218,7 +265,7 @@ pub async fn evict_session_to_state(
// binding; this catches the case that wedged a session when an idle-evict
// completed exactly as a drain dispatched it — already `Idle`, but the
// drain still drove it to `Evacuating` on a destroyed sandbox.
match state.services.meta.get_session(session_id).await {
let entry_status = match state.services.meta.get_session(session_id).await {
Ok(s) if !matches!(s.status, SessionState::Active | SessionState::Evicting) => {
tracing::info!(
session_id = %session_id,
Expand All @@ -229,13 +276,13 @@ pub async fn evict_session_to_state(
reason: "session no longer evictable (a concurrent eviction won the lease first)",
});
}
Ok(_) => {}
Ok(s) => s.status,
Err(e) => {
return Err(EvictError::Meta(format!(
"evict: re-read session state after lease: {e}"
)));
}
}
};

// ADR 0016 A.1.1: entry log. Was silent before — a coord pod
// running the pipeline repeatedly (e.g. retry storm, post-roll
Expand Down Expand Up @@ -298,26 +345,27 @@ pub async fn evict_session_to_state(
// idle-evict path parks; drain/evac (Evacuating) always captures.
// `allow_park == false` is the reaper's DESCENT path (already
// parked, now forcing the full capture) — it must never re-park.
if allow_park && host_has_memory_headroom(state, sandbox_id).await {
if allow_park && host_has_memory_headroom(state, session_id).await {
match state.services.host.pause(sandbox_id).await {
Ok(()) => {
if let Err(e) = state
.services
.meta
.set_session_park_rung(session_id, 2, Some(Utc::now()))
.await
{
tracing::warn!(session_id = %session_id, error = %e,
"rung-2 park: park_rung stamp failed; un-pausing to avoid a stuck paused VM");
let _ = state.services.host.resume(sandbox_id).await;
} else {
Ok(()) => match park_paused_bookkeeping(state, session_id, entry_status).await {
Ok(()) => {
::metrics::counter!(crate::metrics::EVICTION_PARKED_PAUSED_TOTAL)
.increment(1);
tracing::info!(session_id = %session_id, %sandbox_id,
"rung-2 park: VM paused in place (host has memory headroom)");
"rung-2 park: VM paused in place (host has memory headroom)");
return Ok(EvictOutcome::ParkedPaused);
}
}
Err(e) => {
tracing::warn!(session_id = %session_id, error = %e,
"rung-2 park: bookkeeping failed; un-pausing and falling through to full eviction");
let _ = state.services.host.resume(sandbox_id).await;
let _ = state
.services
.meta
.set_session_park_rung(session_id, 0, None)
.await;
}
},
Err(engram_core::SandboxError::InvalidSpec(_)) => {
// Backend can't pause (VZ/Process) — fall through to
// the full eviction below.
Expand Down Expand Up @@ -1233,7 +1281,7 @@ async fn park_reaper_advance_one(
.parked_at
.map(|at| Utc::now().signed_duration_since(at) >= park_dwell_cap())
.unwrap_or(true); // no stamp → treat as long-parked (descend)
let has_headroom = host_has_memory_headroom(state, sandbox_id).await;
let has_headroom = host_has_memory_headroom(state, session_id).await;
let reason = if !has_headroom {
"pressure"
} else if dwell_exceeded {
Expand Down Expand Up @@ -3217,6 +3265,7 @@ mod tests {
.host_registry
.host_of(sandbox_id)
.expect("sandbox routed");
meta.session.lock().host_id = Some(host_id); // headroom resolves via PG now
seed_host_with_free_ram(&meta, host_id, 60_000); // ~91% free ≥ 30 floor

let outcome = evict_idle_session(&state, session_id, sandbox_id)
Expand Down Expand Up @@ -3288,6 +3337,7 @@ mod tests {
.host_registry
.host_of(sandbox_id)
.expect("sandbox routed");
meta.session.lock().host_id = Some(host_id); // headroom resolves via PG now
seed_host_with_free_ram(&meta, host_id, 60_000);

// Park first.
Expand Down Expand Up @@ -3361,6 +3411,7 @@ mod tests {
.await
.unwrap();
let host_id = state.host_registry.host_of(sandbox_id).expect("routed");
meta.session.lock().host_id = Some(host_id); // headroom resolves via PG now
seed_host_with_free_ram(&meta, host_id, 60_000);
// Freshly parked (within dwell) with headroom → the reaper holds.
state
Expand Down
46 changes: 9 additions & 37 deletions crates/engram-coordinator/src/outbox_delivery.rs
Original file line number Diff line number Diff line change
Expand Up @@ -46,15 +46,6 @@ const RESCAN_INTERVAL: Duration = Duration::from_secs(2);
/// after a cold resume; the cost of redelivering early is nil (dedup).
const ACK_TIMEOUT: Duration = Duration::from_secs(30);

/// ADR 0074 rung-2: on a `NotFound` (VM alive, no harness bound), how many
/// plain-retry attempts to wait for the in-guest harness to SELF-reattach
/// before falling back to `start_agent`. A parked-paused un-pause leaves the
/// harness alive in RAM; it re-dials the hub on its own within a beat, and a
/// premature `start_agent` would kill it and pay a full ~40s handshake. Over
/// `failure_backoff`'s 2s/4s/6s schedule this is a ~12s self-reattach window —
/// ample for an in-guest redial, a small delay for the rare genuine desync.
const HARNESS_SELF_REATTACH_ATTEMPTS: i32 = 3;

/// Backoff for rows whose delivery attempt FAILED (resume error, host
/// error). Grows linearly with attempts, capped — an unresumable
/// session shouldn't spin the driver, and there is deliberately no
Expand Down Expand Up @@ -220,34 +211,15 @@ async fn deliver_one(state: &SharedState, row: &OutboxRow) -> Result<(), Deliver
match forward().await {
Ok(()) => Ok(()),
Err(SandboxError::NotFound) => {
// The VM is alive but no harness is attached for delivery. There are
// TWO causes, and they want OPPOSITE responses:
//
// 1. ADR 0074 rung-2 parked-paused un-pause: the VM was PAUSED with
// the harness alive in RAM, then un-paused by the ascent. The
// harness's ADR 0073 self-auth loop re-dials the hub on its own
// (its vsock dropped across the pause) and re-binds with its
// still-valid epoch within a beat. `start_agent` here is not just
// wasteful — it KILLS that self-reattaching harness and pays a
// full ~40s agent_handshake (respawn + resume prefault), turning
// a returning-user un-pause (should be sub-second) into WORSE
// than a plain evict+resume. This was the rung-2 regression.
//
// 2. Genuine harness-unbound desync (harness process gone, VM alive):
// the harness will NOT come back on its own — `start_agent` (the
// e35ed1fa self-heal) is required.
//
// We can't synchronously distinguish them (no per-sandbox harness-
// attach RPC), so give the self-reattach a bounded window of plain
// retries FIRST; only fall back to `start_agent` once the harness
// clearly isn't self-reattaching. The window (`failure_backoff`
// cumulative over `HARNESS_SELF_REATTACH_ATTEMPTS`) is a few seconds
// — ample for an in-guest redial, a small delay for the rare desync.
if row.attempts < HARNESS_SELF_REATTACH_ATTEMPTS {
return Err(DeliverError::Retry(
"harness not attached — awaiting in-guest self-reattach".into(),
));
}
// The VM is alive but no harness is attached: the genuine
// harness-unbound desync (harness process gone, VM alive — the
// e35ed1fa class). The harness will NOT come back on its own —
// its connection loop only re-dials when its established link
// drops or agentd SIGUSR1s it — so `start_agent` is the remedy,
// immediately. (A #594-era "wait for in-guest self-reattach"
// window here was wrong twice over: an ADR 0074 rung-2 un-pause
// keeps the vsock connection INTACT — delivery succeeds and this
// arm never runs — and a real desync has nothing to wait for.)
match crate::api::snapshot::reattach_harness_in_place(state, row.session_id, sandbox_id)
.await
{
Expand Down
13 changes: 13 additions & 0 deletions crates/engram-coordinator/tests/outbox_live_pg.rs
Original file line number Diff line number Diff line change
Expand Up @@ -135,6 +135,19 @@ async fn delivered_row_rearms_after_ack_timeout() {
.await
.expect("defer");
assert!(meta.outbox_next_due(sid).await.unwrap().is_none());

// Defers count as attempts too — `failure_backoff(attempts)` only
// grows if a row failing BEFORE the forward (ensure_active error,
// NotFound) bumps the counter; it used to sit at the floor backoff
// forever and read as attempts=0 in every investigation.
meta.outbox_defer(&row.prompt_id, Duration::ZERO)
.await
.expect("re-defer to now");
let re = meta.outbox_next_due(sid).await.unwrap().expect("due again");
assert_eq!(
re.attempts, 3,
"each defer bumps attempts (1 delivery + 2 defers)"
);
}

#[tokio::test]
Expand Down
Loading
Loading