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 crates/lance-context-master/src/catchup/executor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@ pub async fn execute(mut config: MasterConfig, target: &str) -> Result<ExecuteOu
)
.await?
else {
super::outcome::report_deferred(&state.config.catchup, target)?;
return Ok(ExecuteOutcome::AdmissionBusy);
};
tracing::info!(%target, admission_seconds = admission_started.elapsed().as_secs_f64(),
Expand Down
122 changes: 118 additions & 4 deletions crates/lance-context-master/src/catchup/kubernetes.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,8 @@
use super::{store::Record, CatchupConfig, Result};
use super::{
outcome::{self, JobOutcome},
store::Record,
CatchupConfig, Result,
};
use crate::config::MasterConfig;
use reqwest::{Client, Method, StatusCode};
use serde_json::{json, Value};
Expand Down Expand Up @@ -111,7 +115,7 @@ impl Kubernetes {
config: &MasterConfig,
record: &Record,
stop: bool,
) -> Result<(Option<String>, Option<bool>)> {
) -> Result<(Option<String>, Option<JobOutcome>)> {
let (status, mut job) = self
.request(Method::GET, &format!("{}/{}", self.base, record.job), None)
.await?;
Expand Down Expand Up @@ -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)?),
))
}
}

Expand Down Expand Up @@ -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",
Expand All @@ -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");
Expand Down Expand Up @@ -324,6 +339,7 @@ mod tests {
patches: AtomicUsize,
pods_stopped: AtomicBool,
job: std::sync::Mutex<Option<Value>>,
pod_override: std::sync::Mutex<Option<Value>>,
}
async fn create(
State(mock): State<Arc<Mock>>,
Expand Down Expand Up @@ -364,6 +380,9 @@ mod tests {
(StatusCode::OK, Json(job.clone()))
}
async fn pods(State(mock): State<Arc<Mock>>) -> Json<Value> {
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"}}}]}),
)
Expand All @@ -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))
Expand Down Expand Up @@ -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());
Expand Down
150 changes: 93 additions & 57 deletions crates/lance-context-master/src/catchup/mod.rs
Original file line number Diff line number Diff line change
@@ -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)]
Expand Down Expand Up @@ -64,6 +65,13 @@ pub struct CatchupConfig {
pub target: Option<String>,
#[arg(long, env = "CATCHUP_JOB_NAME", hide = true)]
pub job_name: Option<String>,
/// Controller Jobs opt in to identity-bound Kubernetes termination receipts.
#[arg(long, env = "CATCHUP_OUTCOME_PATH", hide = true)]
pub outcome_path: Option<String>,
#[arg(long, env = "CATCHUP_ATTEMPT_ID", hide = true)]
pub attempt_id: Option<String>,
#[arg(long, env = "CATCHUP_POD_UID", hide = true)]
pub pod_uid: Option<String>,
}
impl Default for CatchupConfig {
fn default() -> Self {
Expand All @@ -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,
}
}
}
Expand Down Expand Up @@ -329,6 +340,80 @@ pub fn spawn(state: &Arc<MasterState>) -> Option<tokio::task::JoinHandle<()>> {
}
}))
}
async fn reconcile_record(
state: &Arc<MasterState>,
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<MasterState>, cursor: &mut usize) -> Result<()> {
let Some(_operation) = state.admission.try_admit() else {
return Ok(());
Expand All @@ -346,63 +431,14 @@ async fn tick(state: &Arc<MasterState>, 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::<Vec<_>>().await?;
Ok::<_,String>(())
}.await;
stream::iter(inventory.active().await?)
.map(|record| reconcile_record(state, &inventory, &kube, record))
.buffer_unordered(8)
.try_collect::<Vec<_>>()
.await?;
Ok::<_, String>(())
}
.await;
let released = state
.task_store
.release_coordination_lock(guard)
Expand Down
Loading
Loading