From 1e5d254575729fec0961aa0b63a66c53d5a23a52 Mon Sep 17 00:00:00 2001 From: Beinan Wang <> Date: Tue, 6 Oct 2026 06:18:57 +0000 Subject: [PATCH] feat: keep selected catchup tables covered across Job lifecycles --- .../src/catchup/kubernetes.rs | 144 ++++++++++++++++-- .../lance-context-master/src/catchup/mod.rs | 22 ++- .../lance-context-master/src/catchup/store.rs | 12 +- .../lance-context-master/src/catchup/tests.rs | 90 +++++++++++ docs/master-catchup.md | 34 +++++ 5 files changed, 286 insertions(+), 16 deletions(-) diff --git a/crates/lance-context-master/src/catchup/kubernetes.rs b/crates/lance-context-master/src/catchup/kubernetes.rs index 2916ea48..272f6910 100644 --- a/crates/lance-context-master/src/catchup/kubernetes.rs +++ b/crates/lance-context-master/src/catchup/kubernetes.rs @@ -120,7 +120,22 @@ impl Kubernetes { .request(Method::GET, &format!("{}/{}", self.base, record.job), None) .await?; if status == StatusCode::NOT_FOUND { - if stop || record.job_uid.is_some() { + if let Some(uid) = &record.job_uid { + // Deletion/preemption can remove the Job before its terminal + // status is collected. Only exact, stopped child Pods establish + // resource retirement; absence alone proves nothing about an + // old process on a partitioned node. Storage recovery is still + // required by normal admission before any successor can write. + let pods = self.job_pods(uid).await?; + if confirmed_retired_children(record, uid, &pods) { + return Ok((Some(uid.clone()), Some(JobOutcome::Failed))); + } + return Err( + "catch-up Job missing; no exact terminal Pod evidence; reservation retained" + .into(), + ); + } + if stop { return Err( "catch-up Job missing after confirmation or startup stall; reservation retained" .into(), @@ -201,7 +216,21 @@ impl Kubernetes { }; // Older Kubernetes versions publish Failed while Pods still terminate. // Check Job UID and every Pod phase before giving the resource slot back. - if !uid.bytes().all(|b| b.is_ascii_alphanumeric() || b == b'-') { + let items = self.job_pods(&uid).await?; + if items + .iter() + .any(|p| !matches!(p["status"]["phase"].as_str(), Some("Succeeded" | "Failed"))) + { + return Ok((Some(uid), None)); + } + Ok(( + Some(uid), + Some(outcome::terminal_outcome(record, &job, &items, success)?), + )) + } + + async fn job_pods(&self, uid: &str) -> Result> { + if uid.is_empty() || !uid.bytes().all(|b| b.is_ascii_alphanumeric() || b == b'-') { return Err("invalid Job UID".into()); } let url = format!( @@ -216,20 +245,53 @@ impl Kubernetes { { return Err("cannot confirm catch-up Pods terminated".into()); } - let items = pods["items"].as_array().ok_or("invalid Pod list")?; - if items - .iter() - .any(|p| !matches!(p["status"]["phase"].as_str(), Some("Succeeded" | "Failed"))) - { - return Ok((Some(uid), None)); - } - Ok(( - Some(uid), - Some(outcome::terminal_outcome(record, &job, items, success)?), - )) + pods["items"] + .as_array() + .cloned() + .ok_or("invalid Pod list".into()) } } +fn confirmed_retired_children(record: &Record, uid: &str, pods: &[Value]) -> bool { + !pods.is_empty() + && pods.iter().all(|pod| { + pod["metadata"]["uid"] + .as_str() + .is_some_and(|uid| !uid.is_empty()) + && pod["metadata"]["ownerReferences"] + .as_array() + .is_some_and(|owners| { + owners.iter().any(|o| { + o["uid"] == uid + && o["name"] == record.job + && o["kind"] == "Job" + && o["controller"] == true + }) + }) + && matches!( + pod["status"]["phase"].as_str(), + Some("Failed" | "Succeeded") + ) + && pod["spec"]["restartPolicy"] == "Never" + && pod["spec"]["containers"] + .as_array() + .is_some_and(|c| c.len() == 1) + && pod["status"]["containerStatuses"] + .as_array() + .is_some_and(|statuses| { + statuses.len() == 1 + && statuses[0]["restartCount"] == 0 + && statuses[0]["state"]["terminated"]["finishedAt"] + .as_str() + .is_some_and(|s| !s.is_empty()) + && statuses[0]["state"]["terminated"]["exitCode"] + .as_i64() + .is_some() + && statuses[0]["name"] == pod["spec"]["containers"][0]["name"] + }) + }) +} + pub(super) fn render_job(config: &MasterConfig, record: &Record, mut spec: Value) -> Value { let mut env = spec["containers"][0]["env"] .as_array() @@ -466,6 +528,25 @@ mod tests { client.reconcile(&config, &record, false).await.unwrap().1, Some(JobOutcome::Failed) ); + // Simulate Job deletion before the controller collected terminal status. + *mock.job.lock().unwrap() = None; + assert!(client.reconcile(&config, &record, false).await.is_err()); + *mock.pod_override.lock().unwrap() = Some(json!({"items":[{ + "metadata":{"uid":"pod-uid","ownerReferences":[{"uid":"test-uid", + "name":record.job,"kind":"Job","controller":true}]}, + "spec":{"restartPolicy":"Never","containers":[{"name":"executor"}]}, + "status":{"phase":"Failed","containerStatuses":[{"name":"executor","restartCount":0, + "state":{"terminated":{"exitCode":137,"finishedAt":"2026-10-06T00:00:00Z"}}}]} + }]})); + assert_eq!( + client.reconcile(&config, &record, false).await.unwrap().1, + Some(JobOutcome::Failed) + ); + assert_eq!( + mock.creates.load(Ordering::SeqCst), + 1, + "must not recreate a missing confirmed Job under the old name" + ); server.abort(); } #[tokio::test] @@ -562,6 +643,43 @@ mod tests { server.abort(); } + #[test] + fn missing_job_requires_exact_terminal_children_not_absence_or_pod_phase_alone() { + let record: Record = serde_json::from_value(json!({ + "target":"hot", "job":"job", "job_uid":"job-uid", "slot":0, + "attempt":1, "active":true, "reason":"test", "pending_at_admission":1000, + "created_at_ms":0, "finished_at_ms":null, "consecutive_failures":0, + "needs_attention":false, "next_retry_ms":0, "outcome":null + })) + .unwrap(); + let pod = json!({"metadata":{"uid":"pod", "ownerReferences":[{ + "uid":"job-uid", "name":"job", "kind":"Job", "controller":true}]}, + "spec":{"restartPolicy":"Never", "containers":[{"name":"executor"}]}, + "status":{"phase":"Failed", "containerStatuses":[{"name":"executor", "restartCount":0, + "state":{"terminated":{"exitCode":137,"finishedAt":"2026-10-06T00:00:00Z"}}}]}}); + assert!(confirmed_retired_children( + &record, + "job-uid", + std::slice::from_ref(&pod) + )); + assert!(!confirmed_retired_children(&record, "job-uid", &[])); + for (pointer, value) in [ + ("/metadata/uid", json!("")), + ("/metadata/ownerReferences/0/uid", json!("other")), + ("/status/phase", json!("Running")), + ("/status/containerStatuses/0/restartCount", json!(1)), + ("/status/containerStatuses/0/state/terminated", json!(null)), + ("/spec/restartPolicy", json!("Always")), + ] { + let mut changed = pod.clone(); + *changed.pointer_mut(pointer).unwrap() = value; + assert!( + !confirmed_retired_children(&record, "job-uid", &[changed]), + "{pointer}" + ); + } + } + #[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 8f1dba9c..edee8f36 100644 --- a/crates/lance-context-master/src/catchup/mod.rs +++ b/crates/lance-context-master/src/catchup/mod.rs @@ -34,6 +34,10 @@ pub struct CatchupConfig { pub max_jobs: usize, #[arg(long, env = "CATCHUP_MIN_PENDING", default_value_t = 256)] pub min_pending: i64, + /// Keep servicing any positive WAL count on these explicitly owned targets. + /// Existing external publishers require a separate qualified ownership handoff. + #[arg(long, env = "CATCHUP_CONTINUOUS_TARGETS", value_delimiter = ',')] + pub continuous_targets: Vec, #[arg(long, env = "CATCHUP_STATS_MAX_AGE_SECS", default_value_t = 900)] pub stats_max_age_secs: u64, #[arg(long, env = "CATCHUP_INTERVAL_SECS", default_value_t = 30)] @@ -81,6 +85,7 @@ impl Default for CatchupConfig { namespace: "default".into(), max_jobs: 4, min_pending: 256, + continuous_targets: Vec::new(), stats_max_age_secs: 900, interval_secs: 30, slice_secs: 1800, @@ -106,6 +111,11 @@ impl CatchupConfig { if self.max_jobs == 0 || self.max_jobs > 256 || self.min_pending < 1 + || self.continuous_targets.len() > 256 + || self + .continuous_targets + .iter() + .any(|t| t.is_empty() || t.len() > 512) || self.interval_secs == 0 || self.stats_max_age_secs == 0 || self.startup_timeout_secs < 120 @@ -236,7 +246,17 @@ fn eligibility( { return Some("stale_stats"); } - if row.pending_wal_generations < config.catchup.min_pending { + let threshold = if config + .catchup + .continuous_targets + .iter() + .any(|t| t == target) + { + 1 + } else { + config.catchup.min_pending + }; + if row.pending_wal_generations < threshold { return Some("below_threshold"); } None diff --git a/crates/lance-context-master/src/catchup/store.rs b/crates/lance-context-master/src/catchup/store.rs index 8cfe0b49..abfd4786 100644 --- a/crates/lance-context-master/src/catchup/store.rs +++ b/crates/lance-context-master/src/catchup/store.rs @@ -65,8 +65,16 @@ impl Inventory { pub async fn ensure_policy(&self, master: &crate::config::MasterConfig) -> Result<()> { let config = &master.catchup; let template = super::kubernetes::read_template(config)?; - let policy = serde_json::json!({"namespace":config.namespace,"max_jobs":config.max_jobs,"template":template, - "shards":config.shards,"merge_max_bytes":config.merge_max_bytes,"merge_memory_bytes":config.merge_memory_bytes,"slice":config.slice_secs,"startup_timeout":config.startup_timeout_secs,"idle_timeout":master.maintenance.maintenance_idle_timeout_secs}).to_string(); + let mut policy = serde_json::json!({"namespace":config.namespace,"max_jobs":config.max_jobs,"template":template, + "shards":config.shards,"merge_max_bytes":config.merge_max_bytes,"merge_memory_bytes":config.merge_memory_bytes,"slice":config.slice_secs,"startup_timeout":config.startup_timeout_secs,"idle_timeout":master.maintenance.maintenance_idle_timeout_secs}); + // Preserve the existing policy bytes when continuous service is off. + if !config.continuous_targets.is_empty() { + let mut targets = config.continuous_targets.clone(); + targets.sort(); + targets.dedup(); + policy["continuous_targets"] = serde_json::json!(targets); + } + let policy = policy.to_string(); let key = format!("{}/catchup-policy", self.prefix); let mut client = self.client.clone(); client diff --git a/crates/lance-context-master/src/catchup/tests.rs b/crates/lance-context-master/src/catchup/tests.rs index b1727ffd..408c0e68 100644 --- a/crates/lance-context-master/src/catchup/tests.rs +++ b/crates/lance-context-master/src/catchup/tests.rs @@ -51,6 +51,96 @@ fn admission_requires_fresh_owned_pressure_even_for_manual_requests() { c.catchup.enabled = false; assert_eq!(eligibility(&c, Some(&r), "hot", now), Some("disabled")); } + +#[test] +fn continuous_coverage_services_small_tails_without_relaxing_ownership_or_freshness() { + let mut c = config(); + c.catchup.continuous_targets = vec!["hot".into(), "legacy".into()]; + let now = 10_000_000; + let mut r = row(now); + r.pending_wal_generations = 1; + assert!(eligibility(&c, Some(&r), "hot", now).is_none()); + assert_eq!( + eligibility(&c, Some(&r), "other", now), + Some("below_threshold") + ); + assert_eq!( + eligibility(&c, Some(&r), "legacy", now), + Some("requires_owned_target") + ); + r.pending_wal_generations = 0; + assert_eq!( + eligibility(&c, Some(&r), "hot", now), + Some("below_threshold") + ); + r.pending_wal_generations = 1; + r.scanned_at = now - 901_000; + assert_eq!(eligibility(&c, Some(&r), "hot", now), Some("stale_stats")); + r.scanned_at = now; + c.merge_rollout.drain_targets.push("hot".into()); + assert_eq!( + eligibility(&c, Some(&r), "hot", now), + Some("requires_owned_target") + ); +} + +#[tokio::test] +#[ignore = "requires ETCD_TEST_ENDPOINTS"] +async fn continuous_coverage_is_durable_and_resumes_only_after_fresh_success_stats() { + let (_dir, state) = fixture().await.unwrap(); + let mut config = state.config.clone(); + config.catchup.continuous_targets = vec!["hot".into()]; + let first = MasterState::new(config.clone()).await.unwrap(); + let inventory = Inventory::new(&first); + inventory.ensure_policy(&config).await.unwrap(); + // A restarted master must declare the same desired coverage. + assert!(Inventory::new(&state) + .ensure_policy(&state.config) + .await + .is_err()); + let now = chrono::Utc::now().timestamp_millis(); + let mut r = row(now); + r.pending_wal_generations = 1; + let request = Trigger { + target: "hot".into(), + reason: "continuous".into(), + dry_run: false, + }; + assert_eq!( + admit(&first, Some(&r), &request).await.unwrap().decision, + "reserved" + ); + let active = inventory.get("hot").await.unwrap().unwrap(); + // An older successful pass with expired cooldown, but no newer stats yet. + inventory + .complete(&active, JobOutcome::Succeeded, now - 120_000) + .await + .unwrap(); + r.scanned_at = now - 120_001; + let successor = MasterState::new(config).await.unwrap(); + assert_eq!( + admit(&successor, Some(&r), &request) + .await + .unwrap() + .decision, + "awaiting_fresh_stats" + ); + r.scanned_at = now; + assert_eq!( + admit(&successor, Some(&r), &request) + .await + .unwrap() + .decision, + "reserved" + ); + let next = Inventory::new(&successor) + .get("hot") + .await + .unwrap() + .unwrap(); + assert_eq!(next.attempt, active.attempt + 1); + assert_ne!(next.job, active.job); +} #[test] fn job_is_native_and_bounded_and_overrides_unsafe_inherited_mode() { let mut c = config(); diff --git a/docs/master-catchup.md b/docs/master-catchup.md index 23064171..f6517d94 100644 --- a/docs/master-catchup.md +++ b/docs/master-catchup.md @@ -238,3 +238,37 @@ 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. + + +## Continuous coverage for selected tables + +`CATCHUP_CONTINUOUS_TARGETS=table_a,table_b` makes the existing controller service +any positive fresh WAL count for these tables, including small tails after a +successful slice. Other tables retain `CATCHUP_MIN_PENDING`. Zero pending does +not create an idle Job. Continuous targets must still be explicitly owned and +non-draining, with the same storage fences, failure deadlines, global capacity +and fresh-after-success stats requirement as ordinary catchup. + +The continuous target set is stored in the shared catchup policy, canonicalized +across replicas. A master with a different set refuses autonomous reconciliation +instead of silently abandoning desired coverage after failover. Use the existing +policy-change procedure to change it. Empty configuration preserves existing +policy compatibility. The existing GET/POST catchup API observes/triggers the +same admission, so manual requests cannot bypass the ownership or memory bounds. + +If a confirmed Job disappears, the controller now checks children belonging to +that exact Job UID. It only releases the inventory slot when nonempty child +inventory proves terminal containers with restartPolicy Never and matching owner +references. This records a failure, preserving retry policy. It does not clear +execution, claim or target-lock keys: the normal fenced recovery path must resolve +any old storage operation before successor work. Missing children, running Pods, +restarts, paginated/incomplete listings and identity mismatches retain the +reservation and report an actionable unresolved state. Job absence is never a +proof that a process on a partitioned node stopped. + +This provides continuous lifecycle coverage for controller-managed Jobs. It does +not adopt external persistent native publishers merely by listing their names. +Migrate one external owner at a time after qualified join/fence evidence, preserving +its attempt/failure debt and legacy admission exclusion. Until that handoff is +implemented and verified for the actual old runtime, keep the external owner +active. Do not delete its catchup-active or operator keys to make admission pass.