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
144 changes: 131 additions & 13 deletions crates/lance-context-master/src/catchup/kubernetes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand Down Expand Up @@ -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<Vec<Value>> {
if uid.is_empty() || !uid.bytes().all(|b| b.is_ascii_alphanumeric() || b == b'-') {
return Err("invalid Job UID".into());
}
let url = format!(
Expand All @@ -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()
Expand Down Expand Up @@ -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]
Expand Down Expand Up @@ -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());
Expand Down
22 changes: 21 additions & 1 deletion crates/lance-context-master/src/catchup/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<String>,
#[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)]
Expand Down Expand Up @@ -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,
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down
12 changes: 10 additions & 2 deletions crates/lance-context-master/src/catchup/store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
90 changes: 90 additions & 0 deletions crates/lance-context-master/src/catchup/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
34 changes: 34 additions & 0 deletions docs/master-catchup.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Loading