From c5f55d2a090b01fa21535d5e0954a801a7b04f8c Mon Sep 17 00:00:00 2001 From: Beinan Wang <> Date: Tue, 6 Oct 2026 06:08:24 +0000 Subject: [PATCH 1/2] fix: reconcile catchup admission deferrals without failure debt --- .../src/catchup/executor.rs | 1 + .../src/catchup/kubernetes.rs | 122 ++++++++- .../lance-context-master/src/catchup/mod.rs | 150 +++++++---- .../src/catchup/outcome.rs | 252 ++++++++++++++++++ .../lance-context-master/src/catchup/store.rs | 54 ++-- .../lance-context-master/src/catchup/tests.rs | 106 +++++++- docs/master-catchup.md | 32 +++ 7 files changed, 630 insertions(+), 87 deletions(-) create mode 100644 crates/lance-context-master/src/catchup/outcome.rs diff --git a/crates/lance-context-master/src/catchup/executor.rs b/crates/lance-context-master/src/catchup/executor.rs index 7cde32c6..aac1e0d1 100644 --- a/crates/lance-context-master/src/catchup/executor.rs +++ b/crates/lance-context-master/src/catchup/executor.rs @@ -51,6 +51,7 @@ pub async fn execute(mut config: MasterConfig, target: &str) -> Result Result<(Option, Option)> { + ) -> Result<(Option, Option)> { let (status, mut job) = self .request(Method::GET, &format!("{}/{}", self.base, record.job), None) .await?; @@ -219,7 +223,10 @@ impl Kubernetes { { return Ok((Some(uid), None)); } - Ok((Some(uid), Some(success))) + Ok(( + Some(uid), + Some(outcome::terminal_outcome(record, &job, items, success)?), + )) } } @@ -265,6 +272,8 @@ pub(super) fn render_job(config: &MasterConfig, record: &Record, mut spec: Value ), ("CATCHUP_TARGET", record.target.clone()), ("CATCHUP_JOB_NAME", record.job.clone()), + ("CATCHUP_ATTEMPT_ID", record.attempt.to_string()), + ("CATCHUP_OUTCOME_PATH", outcome::RECEIPT_PATH.into()), ("CATCHUP_SHARDS", shards.join(",")), ( "CATCHUP_MERGE_MAX_GENERATIONS", @@ -291,7 +300,13 @@ pub(super) fn render_job(config: &MasterConfig, record: &Record, mut spec: Value env.retain(|e| e["name"] != name); env.push(json!({"name": name, "value": value})); } + env.retain(|e| e["name"] != "CATCHUP_POD_UID"); + env.push( + json!({"name":"CATCHUP_POD_UID","valueFrom":{"fieldRef":{"fieldPath":"metadata.uid"}}}), + ); spec["containers"][0]["env"] = json!(env); + spec["containers"][0]["terminationMessagePath"] = json!(outcome::RECEIPT_PATH); + spec["containers"][0]["terminationMessagePolicy"] = json!("File"); spec["containers"][0]["args"] = json!([]); spec["restartPolicy"] = json!("Never"); spec["preemptionPolicy"] = json!("Never"); @@ -324,6 +339,7 @@ mod tests { patches: AtomicUsize, pods_stopped: AtomicBool, job: std::sync::Mutex>, + pod_override: std::sync::Mutex>, } async fn create( State(mock): State>, @@ -364,6 +380,9 @@ mod tests { (StatusCode::OK, Json(job.clone())) } async fn pods(State(mock): State>) -> Json { + if let Some(pods) = mock.pod_override.lock().unwrap().clone() { + return Json(pods); + } Json( json!({"items":[{"status":{"phase":if mock.pods_stopped.load(Ordering::SeqCst) {"Failed"} else {"Running"}}}]}), ) @@ -376,6 +395,7 @@ mod tests { patches: AtomicUsize::new(0), pods_stopped: AtomicBool::new(false), job: std::sync::Mutex::new(None), + pod_override: std::sync::Mutex::new(None), }); let app = Router::new() .route("/jobs", axum::routing::post(create)) @@ -444,10 +464,104 @@ mod tests { mock.pods_stopped.store(true, Ordering::SeqCst); assert_eq!( client.reconcile(&config, &record, false).await.unwrap().1, - Some(false) + Some(JobOutcome::Failed) ); server.abort(); } + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn terminal_deferred_job_releases_its_slot_without_touching_other_execution() { + use super::super::{reconcile_record, store::Inventory}; + let (_dir, state) = super::super::tests::fixture().await.unwrap(); + let inventory = Inventory::new(&state); + inventory.reserve("hot", "test", 1000, 0).await.unwrap(); + let record = inventory.get("hot").await.unwrap().unwrap(); + // A different operation arrived after admission but before our claim. + let coordinator = state.task_store.merge_coordinator(); + let other = lance_context_merge::Execution::new("hot", "master:other", "other-owner", 1); + state + .task_store + .etcd_client() + .clone() + .put( + lance_context_merge::execution_key(&state.config.etcd.etcd_prefix, "hot"), + serde_json::to_vec(&other).unwrap(), + None, + ) + .await + .unwrap(); + let mut job = render_job( + &state.config, + &record, + read_template(&state.config.catchup).unwrap(), + ); + job["metadata"]["uid"] = json!("job-uid"); + let message = json!({"version":1, "outcome":"deferred_before_claim", + "target":"hot", "job":record.job, "attempt":record.attempt.to_string(), + "pod_uid":"pod-uid"}) + .to_string(); + let pods = json!({"items":[{"metadata":{"uid":"pod-uid","ownerReferences":[{ + "uid":"job-uid","name":record.job,"kind":"Job","controller":true}]}, + "status":{"phase":"Failed","containerStatuses":[{"name":"catchup", "restartCount":0, + "state":{"terminated":{"exitCode":75,"message":message}}}]}}]}); + let mock = Arc::new(Mock { + creates: AtomicUsize::new(0), + terminal: AtomicBool::new(true), + patches: AtomicUsize::new(0), + pods_stopped: AtomicBool::new(true), + job: std::sync::Mutex::new(Some(job)), + pod_override: std::sync::Mutex::new(Some(pods)), + }); + let app = Router::new() + .route("/jobs/{name}", get(lookup).patch(patch)) + .route("/pods", get(self::pods)) + .with_state(mock.clone()); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + let server = tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + let token = _dir.path().join("kube-token"); + std::fs::write(&token, "test").unwrap(); + let kube = Kubernetes { + client: Client::new(), + base: format!("http://{addr}/jobs"), + pods_base: format!("http://{addr}/pods"), + token_file: token.to_string_lossy().into(), + }; + // A code 75 with lost terminal evidence remains unresolved even if + // another execution is healthy. No failure debt or ownership is reset. + let valid_pods = mock.pod_override.lock().unwrap().clone().unwrap(); + let mut missing_receipt = valid_pods.clone(); + missing_receipt["items"][0]["status"]["containerStatuses"][0]["state"]["terminated"] + ["message"] = json!(""); + *mock.pod_override.lock().unwrap() = Some(missing_receipt); + reconcile_record(&state, &inventory, &kube, record.clone()) + .await + .unwrap(); + let unresolved = inventory.get("hot").await.unwrap().unwrap(); + assert!(unresolved.active); + assert!(unresolved.needs_attention); + assert_eq!(unresolved.consecutive_failures, 0); + assert_eq!(inventory.active().await.unwrap().len(), 1); + *mock.pod_override.lock().unwrap() = Some(valid_pods); + reconcile_record(&state, &Inventory::new(&state), &kube, unresolved) + .await + .unwrap(); + let done = inventory.get("hot").await.unwrap().unwrap(); + assert!(!done.active); + assert_eq!(done.outcome.as_deref(), Some("deferred_before_claim")); + assert_eq!(done.consecutive_failures, 0); + assert!(inventory.active().await.unwrap().is_empty()); + assert_eq!( + serde_json::to_value(coordinator.get("hot").await.unwrap().unwrap()).unwrap(), + serde_json::to_value(other).unwrap() + ); + assert_eq!(mock.patches.load(Ordering::SeqCst), 0); + assert_eq!(mock.creates.load(Ordering::SeqCst), 0); + server.abort(); + } + #[test] fn invalid_or_unbounded_templates_are_rejected() { assert!(validate_template(&json!({"containers":[]})).is_err()); diff --git a/crates/lance-context-master/src/catchup/mod.rs b/crates/lance-context-master/src/catchup/mod.rs index c92ddd73..8f1dba9c 100644 --- a/crates/lance-context-master/src/catchup/mod.rs +++ b/crates/lance-context-master/src/catchup/mod.rs @@ -1,6 +1,7 @@ //! Master-owned, bounded Kubernetes catch-up jobs. No payload reads in admission. mod executor; mod kubernetes; +mod outcome; mod progress; pub(crate) mod store; #[cfg(test)] @@ -64,6 +65,13 @@ pub struct CatchupConfig { pub target: Option, #[arg(long, env = "CATCHUP_JOB_NAME", hide = true)] pub job_name: Option, + /// Controller Jobs opt in to identity-bound Kubernetes termination receipts. + #[arg(long, env = "CATCHUP_OUTCOME_PATH", hide = true)] + pub outcome_path: Option, + #[arg(long, env = "CATCHUP_ATTEMPT_ID", hide = true)] + pub attempt_id: Option, + #[arg(long, env = "CATCHUP_POD_UID", hide = true)] + pub pod_uid: Option, } impl Default for CatchupConfig { fn default() -> Self { @@ -84,6 +92,9 @@ impl Default for CatchupConfig { pipeline_enabled: true, target: None, job_name: None, + outcome_path: None, + attempt_id: None, + pod_uid: None, } } } @@ -329,6 +340,80 @@ pub fn spawn(state: &Arc) -> Option> { } })) } +async fn reconcile_record( + state: &Arc, + inventory: &Inventory, + kube: &Kubernetes, + mut record: Record, +) -> Result<()> { + // A busy executor may already have exited while a different maintenance + // execution owns the table. Reconcile its terminal receipt BEFORE sampling + // progress, which correctly refuses to observe/revoke that other owner. + match kube + .reconcile(&state.config, &record, record.termination_requested) + .await + { + Ok((uid, terminal)) => { + if record.job_uid.is_none() && uid.is_some() { + let mut updated = record.clone(); + updated.job_uid = uid; + if !inventory.update(&record, &updated).await? { + return Ok(()); + } + record = updated; + } + if let Some(outcome) = terminal { + return inventory + .complete(&record, outcome, chrono::Utc::now().timestamp_millis()) + .await; + } + } + Err(error) => { + tracing::warn!(target = %record.target, job = %record.job, %error, + "catch-up job unresolved; reservation retained"); + metrics::counter!("master_catchup_job_errors_total").increment(1); + return inventory.note_error(&record, &error).await; + } + } + if record.termination_requested { + return Ok(()); + } + let sample = match progress::sample(state, &record).await { + Ok(sample) => sample, + Err(error) => return inventory.note_error(&record, &error).await, + }; + let observed = progress::observe( + record.progress.as_ref(), + sample.clone(), + chrono::Utc::now().timestamp_millis(), + &state.config.catchup, + state.config.maintenance.maintenance_idle_timeout_secs, + ); + let mut updated = record.clone(); + updated.progress = Some(observed); + if !inventory.update(&record, &updated).await? { + return Ok(()); + } + record = updated; + // Recheck immediately before acting: an old observation cannot authorize + // termination of an execution that has since advanced or been replaced. + if record.progress.as_ref().is_some_and(|p| p.stalled) + && progress::revoke(state, &record, &sample).await? + { + let mut updated = record.clone(); + updated.termination_requested = true; + if !inventory.update(&record, &updated).await? { + return Ok(()); + } + if let Err(error) = kube.reconcile(&state.config, &updated, true).await { + inventory.note_error(&updated, &error).await?; + } + // Terminal results are collected on the next tick, retaining capacity + // until the Job and its Pods have all stopped. + } + Ok(()) +} + async fn tick(state: &Arc, cursor: &mut usize) -> Result<()> { let Some(_operation) = state.admission.try_admit() else { return Ok(()); @@ -346,63 +431,14 @@ async fn tick(state: &Arc, cursor: &mut usize) -> Result<()> { { use futures::{stream, StreamExt, TryStreamExt}; let result = async { - stream::iter(inventory.active().await?).map(|record| { - let inventory = &inventory; - let kube = &kube; - async move { - let mut record = record; - let mut stop = record.termination_requested; - if !stop { - match progress::sample(state, &record).await { - Ok(sample) => { - let now = chrono::Utc::now().timestamp_millis(); - let observed = progress::observe(record.progress.as_ref(), sample.clone(), now, - &state.config.catchup, state.config.maintenance.maintenance_idle_timeout_secs); - let mut updated = record.clone(); - updated.progress = Some(observed); - if !inventory.update(&record, &updated).await? { return Ok(()); } - record = updated; - // Recheck immediately before acting: a stale sample is not - // permission to terminate a now-progressing execution. - if record.progress.as_ref().is_some_and(|p| p.stalled) { - stop = progress::revoke(state, &record, &sample).await?; - if stop { - let mut updated = record.clone(); - updated.termination_requested = true; - if !inventory.update(&record, &updated).await? { return Ok(()); } - record = updated; - } - } - } - Err(error) => { - inventory.note_error(&record, &error).await?; - return Ok(()); - } - } - } - match kube.reconcile(&state.config, &record, stop).await { - Ok((uid, terminal)) => { - if record.job_uid.is_none() && uid.is_some() { - let mut updated = record.clone(); - updated.job_uid = uid; - if !inventory.update(&record, &updated).await? { return Ok(()); } - record = updated; - } - if let Some(success) = terminal { - inventory.complete(&record, success, chrono::Utc::now().timestamp_millis()).await?; - } - }, - Err(error) => { - tracing::warn!(target = %record.target, job = %record.job, %error, "catch-up job unresolved; reservation retained"); - metrics::counter!("master_catchup_job_errors_total").increment(1); - inventory.note_error(&record, &error).await?; - } - } - Ok::<_,String>(()) - } - }).buffer_unordered(8).try_collect::>().await?; - Ok::<_,String>(()) - }.await; + stream::iter(inventory.active().await?) + .map(|record| reconcile_record(state, &inventory, &kube, record)) + .buffer_unordered(8) + .try_collect::>() + .await?; + Ok::<_, String>(()) + } + .await; let released = state .task_store .release_coordination_lock(guard) diff --git a/crates/lance-context-master/src/catchup/outcome.rs b/crates/lance-context-master/src/catchup/outcome.rs new file mode 100644 index 00000000..fd8f1392 --- /dev/null +++ b/crates/lance-context-master/src/catchup/outcome.rs @@ -0,0 +1,252 @@ +//! Terminal evidence for controller-owned Jobs. A code 75 alone is not proof +//! that no claim was acquired: require the exact Pod's structured receipt. +use super::{store::Record, CatchupConfig, Result}; +use serde::{Deserialize, Serialize}; +use serde_json::Value; + +pub(super) const RECEIPT_PATH: &str = "/dev/termination-log"; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(super) enum JobOutcome { + Succeeded, + Deferred, + Failed, +} + +#[derive(Debug, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct DeferredReceipt { + version: u32, + outcome: String, + target: String, + job: String, + attempt: String, + pod_uid: String, +} + +/// Called only after every claim RPC returned a definitive None, before any +/// recovery or payload work. Native supervisors without this opt-in keep their +/// existing stdout contract. Failure to write the receipt remains a failure. +pub(super) fn report_deferred(config: &CatchupConfig, target: &str) -> Result<()> { + let Some(path) = &config.outcome_path else { + return Ok(()); + }; + let receipt = DeferredReceipt { + version: 1, + outcome: "deferred_before_claim".into(), + target: target.into(), + job: config.job_name.clone().ok_or("missing receipt Job")?, + attempt: config.attempt_id.clone().ok_or("missing receipt attempt")?, + pod_uid: config.pod_uid.clone().ok_or("missing receipt Pod UID")?, + }; + let bytes = serde_json::to_vec(&receipt).map_err(|e| e.to_string())?; + // Kubernetes retains at most 4096 bytes per termination message. Leave + // headroom and fail closed instead of depending on a truncated receipt. + if bytes.len() > 2048 || receipt.attempt.is_empty() || receipt.pod_uid.is_empty() { + return Err("invalid catch-up outcome identity or receipt length".into()); + } + std::fs::write(path, bytes).map_err(|e| format!("write catch-up outcome: {e}")) +} + +/// The caller already verified terminal Job state and all listed Pod phases. +/// Missing/ambiguous code-75 evidence retains the reservation for inspection; +/// ordinary failures keep their existing failure policy. +pub(super) fn terminal_outcome( + record: &Record, + job: &Value, + pods: &[Value], + success: bool, +) -> Result { + let has_busy_exit = pods.iter().any(|pod| { + pod["status"]["containerStatuses"] + .as_array() + .is_some_and(|statuses| { + statuses + .iter() + .any(|s| s["state"]["terminated"]["exitCode"] == 75) + }) + }); + if !has_busy_exit { + return Ok(if success { + JobOutcome::Succeeded + } else { + JobOutcome::Failed + }); + } + let invalid = || "unresolved catch-up deferred receipt; reservation retained".to_string(); + if success || pods.len() != 1 || record.termination_requested { + return Err(invalid()); + } + let pod = &pods[0]; + let job_uid = job["metadata"]["uid"].as_str().ok_or_else(invalid)?; + if !pod["metadata"]["ownerReferences"] + .as_array() + .is_some_and(|owners| { + owners.iter().any(|o| { + o["uid"] == job_uid + && o["name"] == record.job + && o["kind"] == "Job" + && o["controller"] == true + }) + }) + { + return Err(invalid()); + } + let statuses = pod["status"]["containerStatuses"] + .as_array() + .ok_or_else(invalid)?; + if statuses.len() != 1 { + return Err(invalid()); + } + let container = &statuses[0]; + let terminated = &container["state"]["terminated"]; + if pod["status"]["phase"] != "Failed" + || container["restartCount"] != 0 + || container["name"] != job["spec"]["template"]["spec"]["containers"][0]["name"] + || terminated["exitCode"] != 75 + || terminated["signal"] + .as_u64() + .is_some_and(|signal| signal != 0) + { + return Err(invalid()); + } + let message = terminated["message"].as_str().ok_or_else(invalid)?; + if message.len() > 2048 { + return Err(invalid()); + } + let receipt: DeferredReceipt = serde_json::from_str(message).map_err(|_| invalid())?; + if receipt.version != 1 + || receipt.outcome != "deferred_before_claim" + || receipt.target != record.target + || receipt.job != record.job + || receipt.attempt != record.attempt.to_string() + || receipt.pod_uid.is_empty() + || pod["metadata"]["uid"] != receipt.pod_uid + { + return Err(invalid()); + } + Ok(JobOutcome::Deferred) +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + + fn fixture() -> (Record, Value, Value) { + let record: Record = serde_json::from_value(json!({ + "target":"hot", "job":"lc-job", "job_uid":"job-uid", "slot":0, + "attempt":7, "active":true, "reason":"test", "pending_at_admission":1000, + "created_at_ms":0, "finished_at_ms":null, "consecutive_failures":2, + "needs_attention":true, "next_retry_ms":0, "outcome":null + })) + .unwrap(); + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("termination"); + report_deferred( + &CatchupConfig { + outcome_path: Some(path.to_string_lossy().into()), + job_name: Some(record.job.clone()), + attempt_id: Some(record.attempt.to_string()), + pod_uid: Some("pod-uid".into()), + ..Default::default() + }, + &record.target, + ) + .unwrap(); + let job = json!({"metadata":{"uid":"job-uid"}, + "spec":{"template":{"spec":{"containers":[{"name":"executor"}]}}}}); + let pod = json!({"metadata":{"uid":"pod-uid", "ownerReferences":[{ + "uid":"job-uid", "name":"lc-job", "kind":"Job", "controller":true}]}, + "status":{"phase":"Failed", "containerStatuses":[{"name":"executor", + "restartCount":0,"state":{"terminated":{"exitCode":75,"signal":0, + "message":std::fs::read_to_string(path).unwrap()}}}]}}); + (record, job, pod) + } + + #[test] + fn exact_receipt_is_deferred_only_after_terminal_failure() { + let (record, job, pod) = fixture(); + assert_eq!( + terminal_outcome(&record, &job, std::slice::from_ref(&pod), false).unwrap(), + JobOutcome::Deferred + ); + assert!(terminal_outcome(&record, &job, std::slice::from_ref(&pod), true).is_err()); + assert!(terminal_outcome(&record, &job, &[pod.clone(), pod], false).is_err()); + } + + #[test] + fn copied_truncated_and_unbound_receipts_never_enable_short_retry() { + let (record, job, pod) = fixture(); + for (pointer, value) in [ + ("/metadata/uid", json!("replacement-pod")), + ("/metadata/ownerReferences/0/uid", json!("other-job")), + ("/status/containerStatuses/0/name", json!("other-container")), + ("/status/containerStatuses/0/restartCount", json!(1)), + ( + "/status/containerStatuses/0/state/terminated/signal", + json!(9), + ), + ( + "/status/containerStatuses/0/state/terminated/message", + json!("{}"), + ), + ( + "/status/containerStatuses/0/state/terminated/message", + json!("busy before claim"), + ), + ] { + let mut changed = pod.clone(); + *changed.pointer_mut(pointer).unwrap() = value; + assert!( + terminal_outcome(&record, &job, &[changed], false).is_err(), + "{pointer}" + ); + } + for (key, value) in [ + ("version", json!(2)), + ("target", json!("other")), + ("job", json!("other")), + ("attempt", json!("6")), + ("pod_uid", json!("other")), + ("outcome", json!("failed_after_claim")), + ] { + let mut changed = pod.clone(); + let message = + &mut changed["status"]["containerStatuses"][0]["state"]["terminated"]["message"]; + let mut receipt: Value = serde_json::from_str(message.as_str().unwrap()).unwrap(); + receipt[key] = value; + *message = json!(receipt.to_string()); + assert!( + terminal_outcome(&record, &job, &[changed], false).is_err(), + "{key}" + ); + } + let mut stopped = record; + stopped.termination_requested = true; + assert!(terminal_outcome(&stopped, &job, &[pod], false).is_err()); + } + + #[test] + fn ordinary_terminal_results_and_legacy_executors_remain_supported() { + let (record, job, mut pod) = fixture(); + pod["status"]["containerStatuses"][0]["state"]["terminated"]["exitCode"] = json!(1); + assert_eq!( + terminal_outcome(&record, &job, &[pod], false).unwrap(), + JobOutcome::Failed + ); + assert_eq!( + terminal_outcome(&record, &job, &[], true).unwrap(), + JobOutcome::Succeeded + ); + assert!(report_deferred(&CatchupConfig::default(), "hot").is_ok()); + assert!(report_deferred( + &CatchupConfig { + outcome_path: Some("/nonexistent/catchup/receipt".into()), + ..Default::default() + }, + "hot" + ) + .is_err()); + } +} diff --git a/crates/lance-context-master/src/catchup/store.rs b/crates/lance-context-master/src/catchup/store.rs index dc3062bf..ccd65b8c 100644 --- a/crates/lance-context-master/src/catchup/store.rs +++ b/crates/lance-context-master/src/catchup/store.rs @@ -1,4 +1,4 @@ -use super::{Decision, Result}; +use super::{outcome::JobOutcome, Decision, Result}; use crate::state::MasterState; use etcd_client::{Client, Compare, CompareOp, GetOptions, Txn, TxnOp}; use serde::{Deserialize, Serialize}; @@ -197,8 +197,8 @@ impl Inventory { created_at_ms: now, finished_at_ms: None, consecutive_failures: old.as_ref().map_or(0, |r| r.consecutive_failures), - needs_attention: false, - next_retry_ms: 0, + needs_attention: old.as_ref().is_some_and(|r| r.needs_attention), + next_retry_ms: old.as_ref().map_or(0, |r| r.next_retry_ms), outcome: None, progress: None, }; @@ -250,30 +250,40 @@ impl Inventory { }, )) } - pub async fn complete(&self, record: &Record, success: bool, now: i64) -> Result<()> { + pub(super) async fn complete( + &self, + record: &Record, + outcome: JobOutcome, + now: i64, + ) -> Result<()> { let mut terminal = record.clone(); terminal.active = false; terminal.finished_at_ms = Some(now); - terminal.consecutive_failures = if success { - 0 - } else { - record.consecutive_failures.saturating_add(1) - }; - terminal.needs_attention = terminal.consecutive_failures >= 3; - let delay_secs = if success { - 60 - } else { - (60u64.saturating_mul(1u64 << terminal.consecutive_failures.min(6))).min(3600) + let delay_secs = match outcome { + JobOutcome::Succeeded => { + terminal.consecutive_failures = 0; + terminal.needs_attention = false; + terminal.outcome = Some("succeeded".into()); + 60 + } + JobOutcome::Deferred => { + // Admission acquired no work. Keep earned failure debt and + // attention, including deadlines surviving a controller restart. + terminal.outcome = Some("deferred_before_claim".into()); + // Stable per-Job jitter; no synchronized etcd retry burst. + 3 + record.job.bytes().fold(0u64, |n, b| (n + u64::from(b)) % 3) + } + JobOutcome::Failed => { + terminal.consecutive_failures = record.consecutive_failures.saturating_add(1); + terminal.needs_attention = terminal.consecutive_failures >= 3; + terminal.outcome = Some("failed; inspect Job status and executor logs".into()); + (60u64.saturating_mul(1u64 << terminal.consecutive_failures.min(6))).min(3600) + } }; terminal.next_retry_ms = now.saturating_add(delay_secs as i64 * 1000); - terminal.outcome = Some( - if success { - "succeeded" - } else { - "failed; inspect Job status and executor logs" - } - .into(), - ); + if outcome == JobOutcome::Deferred { + terminal.next_retry_ms = terminal.next_retry_ms.max(record.next_retry_ms); + } let old = serde_json::to_vec(record).map_err(|e| e.to_string())?; let new = serde_json::to_vec(&terminal).map_err(|e| e.to_string())?; let active = active_key(&self.prefix, &record.target); diff --git a/crates/lance-context-master/src/catchup/tests.rs b/crates/lance-context-master/src/catchup/tests.rs index 8a917687..b1727ffd 100644 --- a/crates/lance-context-master/src/catchup/tests.rs +++ b/crates/lance-context-master/src/catchup/tests.rs @@ -1,3 +1,4 @@ +use super::outcome::JobOutcome; use super::*; use clap::Parser; use lance_context_api::TaskKind; @@ -100,6 +101,8 @@ fn job_is_native_and_bounded_and_overrides_unsafe_inherited_mode() { for (name, expected) in [ ("CATCHUP_PIPELINE_ENABLED", "false"), ("CATCHUP_MERGE_MAX_GENERATIONS", "32"), + ("CATCHUP_ATTEMPT_ID", "1"), + ("CATCHUP_OUTCOME_PATH", "/dev/termination-log"), ] { let matches: Vec<_> = env.iter().filter(|e| e["name"] == name).collect(); assert_eq!(matches.len(), 1); @@ -127,6 +130,15 @@ fn job_is_native_and_bounded_and_overrides_unsafe_inherited_mode() { env.iter().find(|e| e["name"] == "CATCHUP_ENABLED").unwrap()["value"], "false" ); + assert_eq!( + env.iter().find(|e| e["name"] == "CATCHUP_POD_UID").unwrap()["valueFrom"]["fieldRef"] + ["fieldPath"], + "metadata.uid" + ); + assert_eq!( + job["spec"]["template"]["spec"]["containers"][0]["terminationMessagePolicy"], + "File" + ); let mut next = r.clone(); next.attempt = 2; let next_job = kubernetes::render_job(&c, &next, template()); @@ -138,7 +150,7 @@ fn job_is_native_and_bounded_and_overrides_unsafe_inherited_mode() { "worker-1,worker-0" ); } -async fn fixture() -> Option<(tempfile::TempDir, Arc)> { +pub(super) async fn fixture() -> Option<(tempfile::TempDir, Arc)> { let endpoints = std::env::var("ETCD_TEST_ENDPOINTS").expect("ETCD_TEST_ENDPOINTS is required"); let dir = tempfile::tempdir().unwrap(); let mut c = config(); @@ -179,7 +191,7 @@ async fn replicas_share_one_slot_and_dedupe_duplicate_requests() { ); assert_eq!(a.active().await.unwrap().len(), 1); let record = a.get("hot").await.unwrap().unwrap(); - b.complete(&record, false, 1000).await.unwrap(); + b.complete(&record, JobOutcome::Failed, 1000).await.unwrap(); let failed = a.get("hot").await.unwrap().unwrap(); assert_eq!(failed.consecutive_failures, 1); assert!(!failed.active); @@ -198,7 +210,9 @@ async fn replicas_share_one_slot_and_dedupe_duplicate_requests() { "reserved" ); // Replaying an old completion cannot release the new table's slot. - b.complete(&record, true, 2000).await.unwrap(); + b.complete(&record, JobOutcome::Succeeded, 2000) + .await + .unwrap(); assert_eq!(a.active().await.unwrap()[0].target, "other"); } #[tokio::test] @@ -423,7 +437,10 @@ async fn repeated_failed_jobs_keep_backoff_across_controller_restarts() { ); let record = inventory.get("hot").await.unwrap().unwrap(); assert_eq!(record.attempt, u64::from(attempt)); - inventory.complete(&record, false, now + 1).await.unwrap(); + inventory + .complete(&record, JobOutcome::Failed, now + 1) + .await + .unwrap(); let persisted = Inventory::new(&state).get("hot").await.unwrap().unwrap(); assert_eq!(persisted.consecutive_failures, attempt); assert_eq!(persisted.needs_attention, attempt >= 3); @@ -432,6 +449,73 @@ async fn repeated_failed_jobs_keep_backoff_across_controller_restarts() { } } +#[tokio::test] +#[ignore = "requires ETCD_TEST_ENDPOINTS"] +async fn deferred_jobs_preserve_failure_debt_and_fence_stale_completion() { + let (_dir, state) = fixture().await.unwrap(); + let inventory = Inventory::new(&state); + let mut now = 1000; + // Establish genuine failure debt, including the attention flag. + for _ in 0..3 { + inventory.reserve("hot", "test", 1000, now).await.unwrap(); + let record = inventory.get("hot").await.unwrap().unwrap(); + inventory + .complete(&record, JobOutcome::Failed, now + 1) + .await + .unwrap(); + now = inventory.get("hot").await.unwrap().unwrap().next_retry_ms; + } + for _ in 0..3 { + let inventory = Inventory::new(&state); // Simulate a fresh controller. + inventory.reserve("hot", "test", 1000, now).await.unwrap(); + let record = inventory.get("hot").await.unwrap().unwrap(); + assert!(record.needs_attention); + inventory + .complete(&record, JobOutcome::Deferred, now + 1) + .await + .unwrap(); + let deferred = inventory.get("hot").await.unwrap().unwrap(); + assert_eq!(deferred.consecutive_failures, 3); + assert!(deferred.needs_attention); + assert_eq!(deferred.outcome.as_deref(), Some("deferred_before_claim")); + assert!((3000..=5000).contains(&(deferred.next_retry_ms - now - 1))); + assert!(!deferred.active); + assert!(inventory.active().await.unwrap().is_empty()); + assert_ne!( + inventory + .reserve("hot", "too early", 1000, now + 2) + .await + .unwrap() + .decision, + "reserved" + ); + // An old controller cannot turn this deferred attempt into a failure. + inventory + .complete(&record, JobOutcome::Failed, now + 2) + .await + .unwrap(); + assert_eq!( + inventory + .get("hot") + .await + .unwrap() + .unwrap() + .consecutive_failures, + 3 + ); + now = deferred.next_retry_ms; + } + inventory.reserve("hot", "test", 1000, now).await.unwrap(); + let record = inventory.get("hot").await.unwrap().unwrap(); + inventory + .complete(&record, JobOutcome::Succeeded, now + 1) + .await + .unwrap(); + let recovered = inventory.get("hot").await.unwrap().unwrap(); + assert_eq!(recovered.consecutive_failures, 0); + assert!(!recovered.needs_attention); +} + #[tokio::test] #[ignore = "requires ETCD_TEST_ENDPOINTS"] async fn catchup_runtime_ceiling_does_not_cancel_work_within_idle_window() { @@ -605,12 +689,22 @@ async fn busy_executor_preserves_live_claim_then_merges_after_release() { .await .unwrap() .unwrap(); + let receipt_path = _dir.path().join("termination"); + cfg.catchup.outcome_path = Some(receipt_path.to_string_lossy().into()); + cfg.catchup.attempt_id = Some("1".into()); + cfg.catchup.pod_uid = Some("test-pod".into()); let started = std::time::Instant::now(); assert_eq!( execute(cfg.clone(), target).await.unwrap(), ExecuteOutcome::AdmissionBusy ); assert!(started.elapsed() >= std::time::Duration::from_secs(30)); + let receipt: serde_json::Value = + serde_json::from_slice(&std::fs::read(&receipt_path).unwrap()).unwrap(); + assert_eq!(receipt["outcome"], "deferred_before_claim"); + assert_eq!(receipt["target"], target); + assert_eq!(receipt["pod_uid"], "test-pod"); + std::fs::remove_file(&receipt_path).unwrap(); assert_eq!( state.task_store.get(&task.id).await.unwrap().unwrap().state, lance_context_api::TaskState::Running @@ -625,6 +719,10 @@ async fn busy_executor_preserves_live_claim_then_merges_after_release() { execute(cfg, target).await.unwrap(), ExecuteOutcome::Completed ); + assert!( + !receipt_path.exists(), + "claimed work must not emit a deferred receipt" + ); let reader = GenericStore::open_existing(&uri, GenericStoreOptions::default()) .await .unwrap(); diff --git a/docs/master-catchup.md b/docs/master-catchup.md index 68da0ee9..23064171 100644 --- a/docs/master-catchup.md +++ b/docs/master-catchup.md @@ -206,3 +206,35 @@ queue limits, cancellation and manifest-drain barriers remain. The legacy timeou configuration/wire fields remain for compatibility; older master/worker binaries may still enforce them during a mixed-version rollout. Deploy compatible versions throughout before relying on the absence of a total execution ceiling. + + +## Busy admission and terminal evidence + +Controller-created Jobs set a downward-API Pod UID, attempt identity and a +Kubernetes termination-message path. If every claim RPC definitively returned +no claim, the executor writes a versioned `deferred_before_claim` receipt before +exiting 75. Claim RPC errors (including accepted-but-lost responses), recovery, +payload failures and partial work do not produce this receipt. The standalone +native supervisor retains its separate stdout contract. + +The controller waits for terminal Job and Pod status, then checks the receipt +against the target, Job, attempt, actual Pod UID, container and Job owner +reference. Multiple Pods, container restarts, missing/truncated receipts or a +termination already requested by the controller remain unresolved; their +reservation is retained and `needs_attention` is exposed by the existing GET +catchup API. Exit 75 alone never authorizes a fast retry. + +A verified deferred result releases only this inventory reservation and records +`outcome=deferred_before_claim`. It preserves real failure counts, attention and +earned retry deadlines, with 3–5 seconds of per-Job jitter. Actual readmission +occurs on a subsequent configured controller tick (30 seconds by default), and +still checks current execution, target locks and durable failure cooldowns. +It is not a promise of a new Pod in 3–5 seconds. A completed Job is reconciled +before progress sampling so an unrelated maintenance owner cannot hide its +terminal outcome or be revoked by this controller. + +Deploy the updated executor image in the pinned Pod template together with the +controller. Existing successful/ordinary failed Jobs retain their handling; +older executors returning 75 without the receipt require reconciliation rather +than a guessed success or cleared failure history. This change does not adopt +externally managed publishers or enable production catchup automatically. From cbb69c11521395e4db39b8fd619c6fedb152b222 Mon Sep 17 00:00:00 2001 From: Beinan Wang <> Date: Tue, 6 Oct 2026 06:14:41 +0000 Subject: [PATCH 2/2] fix: label terminal metrics with typed outcome --- crates/lance-context-master/src/catchup/store.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/crates/lance-context-master/src/catchup/store.rs b/crates/lance-context-master/src/catchup/store.rs index ccd65b8c..8cfe0b49 100644 --- a/crates/lance-context-master/src/catchup/store.rs +++ b/crates/lance-context-master/src/catchup/store.rs @@ -311,8 +311,8 @@ impl Inventory { .map_err(|e| e.to_string())? .succeeded(); if changed { - tracing::info!(target = %record.target, job = %record.job, success, failures = terminal.consecutive_failures, "catch-up job terminal; capacity released"); - metrics::counter!("master_catchup_jobs_finished_total", "result" => if success { "ok" } else { "failed" }).increment(1); + tracing::info!(target = %record.target, job = %record.job, ?outcome, failures = terminal.consecutive_failures, "catch-up job terminal; capacity released"); + metrics::counter!("master_catchup_jobs_finished_total", "result" => match outcome { JobOutcome::Succeeded => "ok", JobOutcome::Deferred => "deferred", JobOutcome::Failed => "failed" }).increment(1); } Ok(()) }