mirror of
https://github.com/lancedb/lancedb.git
synced 2026-08-18 12:08:35 +00:00
client: jobs surface consolidated onto the platform API
Every connection-level job op now rides /v1/jobs (or the relocated /v1/errors); the legacy /v1/job wire structs are deleted: - list_jobs -> POST /v1/jobs/list with include_status: rows map to JobInfo with the subtype as the job label and units/rows/error from the status payload. State vocabulary unifies on running / finished / failed / cancelled. - get_job -> resolve + describe: a point snapshot keyed by the caller's submission id. - cancel_job -> resolve + platform cancel; false when the job never registered (the legacy best-effort contract). - job_history -> describe + POST /v1/jobs/query_events: one summary row with the lifecycle event timeline for a given id; the no-id form is a registry listing, timeline-free. - errors -> GET /v1/errors (relocated; outside the jobs API). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
+256
-166
@@ -134,30 +134,6 @@ struct RemoteListMaterializedViewsResponse {
|
||||
views: Vec<RemoteMaterializedViewEntry>,
|
||||
}
|
||||
|
||||
#[derive(serde::Deserialize)]
|
||||
struct RemoteJobEntry {
|
||||
table: String,
|
||||
job_id: String,
|
||||
job_type: String,
|
||||
state: String,
|
||||
#[serde(default)]
|
||||
column: Option<String>,
|
||||
#[serde(default)]
|
||||
age_seconds: Option<i64>,
|
||||
#[serde(default)]
|
||||
command: Option<String>,
|
||||
#[serde(default)]
|
||||
units_done: Option<i64>,
|
||||
#[serde(default)]
|
||||
units_total: Option<i64>,
|
||||
#[serde(default)]
|
||||
committed: bool,
|
||||
#[serde(default)]
|
||||
rows_skipped: u64,
|
||||
#[serde(default)]
|
||||
error: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(serde::Deserialize)]
|
||||
struct RemoteDescribePlatformJobResponse {
|
||||
job_id: String,
|
||||
@@ -180,87 +156,49 @@ struct RemoteListPlatformJobsResponse {
|
||||
#[derive(serde::Deserialize)]
|
||||
struct RemotePlatformJobRow {
|
||||
job_id: String,
|
||||
}
|
||||
|
||||
#[derive(serde::Deserialize)]
|
||||
struct RemoteListJobsResponse {
|
||||
jobs: Vec<RemoteJobEntry>,
|
||||
}
|
||||
|
||||
#[derive(serde::Deserialize)]
|
||||
struct RemoteGetJobResponse {
|
||||
#[serde(default)]
|
||||
job: Option<RemoteJobEntry>,
|
||||
}
|
||||
|
||||
#[derive(serde::Deserialize)]
|
||||
struct RemoteCancelJobResponse {
|
||||
cancelled: bool,
|
||||
}
|
||||
|
||||
impl From<RemoteJobEntry> for JobInfo {
|
||||
fn from(j: RemoteJobEntry) -> Self {
|
||||
JobInfo {
|
||||
table: j.table,
|
||||
job_id: j.job_id,
|
||||
job_type: j.job_type,
|
||||
state: j.state,
|
||||
column: j.column,
|
||||
age_seconds: j.age_seconds,
|
||||
command: j.command,
|
||||
units_done: j.units_done,
|
||||
units_total: j.units_total,
|
||||
committed: j.committed,
|
||||
rows_skipped: j.rows_skipped,
|
||||
error: j.error,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(serde::Deserialize)]
|
||||
struct RemoteJobHistoryEntry {
|
||||
table: String,
|
||||
job_id: String,
|
||||
job_type: String,
|
||||
#[serde(default)]
|
||||
job_subtype: String,
|
||||
#[serde(default)]
|
||||
state: String,
|
||||
#[serde(default)]
|
||||
column: Option<String>,
|
||||
created_ms: i64,
|
||||
updated_ms: i64,
|
||||
created_at_millis: i64,
|
||||
#[serde(default)]
|
||||
completed_ms: Option<i64>,
|
||||
#[serde(default)]
|
||||
rows_processed: Option<i64>,
|
||||
#[serde(default)]
|
||||
rows_skipped: Option<i64>,
|
||||
#[serde(default)]
|
||||
error: Option<String>,
|
||||
#[serde(default)]
|
||||
events: Option<String>,
|
||||
status: serde_json::Value,
|
||||
}
|
||||
|
||||
#[derive(serde::Deserialize)]
|
||||
struct RemoteJobHistoryResponse {
|
||||
jobs: Vec<RemoteJobHistoryEntry>,
|
||||
}
|
||||
|
||||
impl From<RemoteJobHistoryEntry> for JobHistoryInfo {
|
||||
fn from(j: RemoteJobHistoryEntry) -> Self {
|
||||
JobHistoryInfo {
|
||||
table: j.table,
|
||||
job_id: j.job_id,
|
||||
job_type: j.job_type,
|
||||
state: j.state,
|
||||
column: j.column,
|
||||
created_ms: j.created_ms,
|
||||
updated_ms: j.updated_ms,
|
||||
completed_ms: j.completed_ms,
|
||||
rows_processed: j.rows_processed,
|
||||
rows_skipped: j.rows_skipped,
|
||||
error: j.error,
|
||||
events: j.events,
|
||||
}
|
||||
/// Platform list-row state -> the client's job vocabulary.
|
||||
fn platform_state_to_client(state: &str) -> String {
|
||||
match state {
|
||||
"in_progress" => "running",
|
||||
"done" => "finished",
|
||||
other => other,
|
||||
}
|
||||
.to_string()
|
||||
}
|
||||
|
||||
/// Describe job_state -> the client's job vocabulary.
|
||||
fn describe_state_to_client(state: &str) -> String {
|
||||
match state {
|
||||
"IN_PROGRESS" => "running",
|
||||
"DONE" => "finished",
|
||||
"FAILED" => "failed",
|
||||
"CANCELLED" => "cancelled",
|
||||
other => other,
|
||||
}
|
||||
.to_string()
|
||||
}
|
||||
|
||||
fn payload_i64(status: &serde_json::Value, key: &str) -> Option<i64> {
|
||||
status.get(key).and_then(serde_json::Value::as_i64)
|
||||
}
|
||||
|
||||
fn payload_error(status: &serde_json::Value) -> Option<String> {
|
||||
status
|
||||
.get("error")
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.map(str::to_string)
|
||||
}
|
||||
|
||||
#[derive(serde::Deserialize)]
|
||||
@@ -1078,24 +1016,63 @@ impl<S: HttpSend> Database for RemoteDatabase<S> {
|
||||
}
|
||||
|
||||
async fn list_jobs(&self) -> Result<Vec<JobInfo>> {
|
||||
let req = self.client.get("/v1/job/list");
|
||||
let req = self
|
||||
.client
|
||||
.post("/v1/jobs/list")
|
||||
.json(&serde_json::json!({ "include_status": true }));
|
||||
let (request_id, rsp) = self.client.send(req).await?;
|
||||
let rsp = self.client.check_response(&request_id, rsp).await?;
|
||||
let body: RemoteListJobsResponse = rsp.json().await.err_to_http(request_id)?;
|
||||
Ok(body.jobs.into_iter().map(JobInfo::from).collect())
|
||||
let body: RemoteListPlatformJobsResponse = rsp.json().await.err_to_http(request_id)?;
|
||||
let now_ms = std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.map(|d| d.as_millis() as i64)
|
||||
.unwrap_or(0);
|
||||
Ok(body
|
||||
.jobs
|
||||
.into_iter()
|
||||
.map(|row| JobInfo {
|
||||
table: row.table,
|
||||
job_id: row.job_id,
|
||||
// The platform job_type is always "indexer"; the subtype
|
||||
// (udf / mv_refresh / compaction / ...) is the useful label.
|
||||
job_type: row.job_subtype,
|
||||
state: platform_state_to_client(&row.state),
|
||||
column: None,
|
||||
age_seconds: (row.created_at_millis > 0)
|
||||
.then(|| (now_ms - row.created_at_millis) / 1000),
|
||||
command: None,
|
||||
units_done: payload_i64(&row.status, "units_done"),
|
||||
units_total: payload_i64(&row.status, "units_total"),
|
||||
committed: row.state == "done",
|
||||
rows_skipped: payload_i64(&row.status, "rows_skipped").unwrap_or(0) as u64,
|
||||
error: payload_error(&row.status),
|
||||
})
|
||||
.collect())
|
||||
}
|
||||
|
||||
async fn get_job(&self, job_id: &str, table: Option<&str>) -> Result<Option<JobInfo>> {
|
||||
// Point-access poll path: GET /v1/job/{id}, with the table as the O(1)
|
||||
// hint when known. `query` handles URL-encoding the table name.
|
||||
let mut req = self.client.get(&format!("/v1/job/{job_id}"));
|
||||
if let Some(t) = table {
|
||||
req = req.query(&[("table", t)]);
|
||||
}
|
||||
let (request_id, rsp) = self.client.send(req).await?;
|
||||
let rsp = self.client.check_response(&request_id, rsp).await?;
|
||||
let body: RemoteGetJobResponse = rsp.json().await.err_to_http(request_id)?;
|
||||
Ok(body.job.map(JobInfo::from))
|
||||
// A point snapshot from the platform API: resolve the submission id,
|
||||
// then describe. The snapshot keeps the caller's id.
|
||||
let Some(platform_id) = self.resolve_platform_job_id(job_id, table).await? else {
|
||||
return Ok(None);
|
||||
};
|
||||
let Some(described) = self.describe_platform_job(&platform_id).await? else {
|
||||
return Ok(None);
|
||||
};
|
||||
Ok(Some(JobInfo {
|
||||
table: table.unwrap_or_default().to_string(),
|
||||
job_id: job_id.to_string(),
|
||||
job_type: described.job_subtype,
|
||||
state: describe_state_to_client(&described.job_state),
|
||||
column: None,
|
||||
age_seconds: None,
|
||||
command: None,
|
||||
units_done: payload_i64(&described.status, "units_done"),
|
||||
units_total: payload_i64(&described.status, "units_total"),
|
||||
committed: described.job_state == "DONE",
|
||||
rows_skipped: payload_i64(&described.status, "rows_skipped").unwrap_or(0) as u64,
|
||||
error: payload_error(&described.status),
|
||||
}))
|
||||
}
|
||||
|
||||
async fn describe_platform_job(
|
||||
@@ -1149,26 +1126,135 @@ impl<S: HttpSend> Database for RemoteDatabase<S> {
|
||||
}
|
||||
|
||||
async fn cancel_job(&self, job_id: &str) -> Result<bool> {
|
||||
let req = self.client.post(&format!("/v1/job/{}/cancel", job_id));
|
||||
let (request_id, rsp) = self.client.send(req).await?;
|
||||
let rsp = self.client.check_response(&request_id, rsp).await?;
|
||||
let body: RemoteCancelJobResponse = rsp.json().await.err_to_http(request_id)?;
|
||||
Ok(body.cancelled)
|
||||
// Resolve the submission id and cancel through the platform API.
|
||||
// False when no matching job has registered (the legacy best-effort
|
||||
// contract).
|
||||
let Some(platform_id) = self.resolve_platform_job_id(job_id, None).await? else {
|
||||
return Ok(false);
|
||||
};
|
||||
self.cancel_platform_job(&platform_id).await?;
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
async fn job_history(&self, job_id: Option<&str>) -> Result<Vec<JobHistoryInfo>> {
|
||||
let mut req = self.client.get("/v1/job/history");
|
||||
if let Some(j) = job_id {
|
||||
req = req.query(&[("job", j)]);
|
||||
}
|
||||
// One job: describe (identity) plus query_events (timeline). No id:
|
||||
// a registry listing, timeline-free -- pass an id for the event log.
|
||||
let Some(caller_id) = job_id else {
|
||||
let rows = self.list_jobs().await?;
|
||||
let now_ms = std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.map(|d| d.as_millis() as i64)
|
||||
.unwrap_or(0);
|
||||
return Ok(rows
|
||||
.into_iter()
|
||||
.map(|j| JobHistoryInfo {
|
||||
table: j.table,
|
||||
job_id: j.job_id,
|
||||
job_type: j.job_type,
|
||||
state: j.state,
|
||||
column: j.column,
|
||||
created_ms: j.age_seconds.map(|a| now_ms - a * 1000).unwrap_or_default(),
|
||||
updated_ms: 0,
|
||||
completed_ms: None,
|
||||
rows_processed: None,
|
||||
rows_skipped: j.rows_skipped.try_into().ok(),
|
||||
error: j.error,
|
||||
events: None,
|
||||
})
|
||||
.collect());
|
||||
};
|
||||
// Accept either a platform id or a submission id.
|
||||
let platform_id = match self.describe_platform_job(caller_id).await? {
|
||||
Some(_) => caller_id.to_string(),
|
||||
None => match self.resolve_platform_job_id(caller_id, None).await? {
|
||||
Some(id) => id,
|
||||
None => return Ok(Vec::new()),
|
||||
},
|
||||
};
|
||||
let Some(described) = self.describe_platform_job(&platform_id).await? else {
|
||||
return Ok(Vec::new());
|
||||
};
|
||||
let req = self
|
||||
.client
|
||||
.post("/v1/jobs/query_events")
|
||||
.json(&serde_json::json!({ "job_id": platform_id }));
|
||||
let (request_id, rsp) = self.client.send(req).await?;
|
||||
let rsp = self.client.check_response(&request_id, rsp).await?;
|
||||
let body: RemoteJobHistoryResponse = rsp.json().await.err_to_http(request_id)?;
|
||||
Ok(body.jobs.into_iter().map(JobHistoryInfo::from).collect())
|
||||
let body = rsp.bytes().await.err_to_http(request_id.clone())?;
|
||||
let reader = arrow_ipc::reader::StreamReader::try_new(std::io::Cursor::new(body), None)
|
||||
.map_err(|e| Error::Http {
|
||||
source: format!("failed to read job-events IPC stream: {e}").into(),
|
||||
request_id: request_id.clone(),
|
||||
status_code: None,
|
||||
})?;
|
||||
|
||||
let mut created_ms = i64::MAX;
|
||||
let mut updated_ms = 0i64;
|
||||
let mut completed_ms = None;
|
||||
let mut last_error = None;
|
||||
let mut events = Vec::new();
|
||||
for batch in reader {
|
||||
let batch = batch.map_err(|e| Error::Http {
|
||||
source: format!("failed to decode job-events batch: {e}").into(),
|
||||
request_id: request_id.clone(),
|
||||
status_code: None,
|
||||
})?;
|
||||
let states = batch
|
||||
.column_by_name("state")
|
||||
.and_then(|c| c.as_any().downcast_ref::<arrow_array::StringArray>());
|
||||
let times = batch
|
||||
.column_by_name("updated_at_millis")
|
||||
.and_then(|c| c.as_any().downcast_ref::<arrow_array::Int64Array>());
|
||||
let payloads = batch
|
||||
.column_by_name("payload")
|
||||
.and_then(|c| c.as_any().downcast_ref::<arrow_array::StringArray>());
|
||||
let (Some(states), Some(times)) = (states, times) else {
|
||||
continue;
|
||||
};
|
||||
for i in 0..batch.num_rows() {
|
||||
let state = states.value(i);
|
||||
let ts = times.value(i);
|
||||
created_ms = created_ms.min(ts);
|
||||
updated_ms = updated_ms.max(ts);
|
||||
if matches!(state, "succeeded" | "failed" | "timed_out" | "canceled") {
|
||||
completed_ms = Some(ts);
|
||||
}
|
||||
if let Some(payloads) = payloads {
|
||||
if !arrow_array::Array::is_null(payloads, i) {
|
||||
if let Ok(payload) =
|
||||
serde_json::from_str::<serde_json::Value>(payloads.value(i))
|
||||
{
|
||||
if let Some(e) = payload_error(&payload) {
|
||||
last_error = Some(e);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
events.push(format!("{state} {ts}"));
|
||||
}
|
||||
}
|
||||
Ok(vec![JobHistoryInfo {
|
||||
table: String::new(),
|
||||
job_id: caller_id.to_string(),
|
||||
job_type: described.job_subtype,
|
||||
state: describe_state_to_client(&described.job_state),
|
||||
column: None,
|
||||
created_ms: if created_ms == i64::MAX {
|
||||
described.creation_ms
|
||||
} else {
|
||||
created_ms
|
||||
},
|
||||
updated_ms,
|
||||
completed_ms,
|
||||
rows_processed: payload_i64(&described.status, "rows_committed"),
|
||||
rows_skipped: payload_i64(&described.status, "rows_skipped"),
|
||||
error: last_error.or_else(|| payload_error(&described.status)),
|
||||
events: (!events.is_empty()).then(|| events.join("\n")),
|
||||
}])
|
||||
}
|
||||
|
||||
async fn errors(&self, job_id: Option<&str>, table: Option<&str>) -> Result<Vec<JobErrorInfo>> {
|
||||
let mut req = self.client.get("/v1/job/errors");
|
||||
let mut req = self.client.get("/v1/errors");
|
||||
if let Some(j) = job_id {
|
||||
req = req.query(&[("job", j)]);
|
||||
}
|
||||
@@ -2253,76 +2339,80 @@ mod tests {
|
||||
assert_eq!(views[0].source_table, "docs");
|
||||
assert!(views[0].auto_refresh);
|
||||
|
||||
// list_jobs
|
||||
// list_jobs: platform listing with status payloads
|
||||
let conn = Connection::new_with_handler(|request| {
|
||||
assert_eq!(request.method(), &reqwest::Method::GET);
|
||||
assert_eq!(request.url().path(), "/v1/job/list");
|
||||
assert_eq!(request.method(), &reqwest::Method::POST);
|
||||
assert_eq!(request.url().path(), "/v1/jobs/list");
|
||||
http::Response::builder()
|
||||
.status(200)
|
||||
.body(
|
||||
r#"{"jobs":[{"table":"docs","job_id":"j-3","job_type":"udf_virtual_column_backfill","state":"running","column":"vec","age_seconds":4,"command":null,"units_done":1,"units_total":2,"committed":false,"rows_skipped":0,"error":null}]}"#,
|
||||
r#"{"jobs":[{"job_id":"plat-3","table":"docs","job_type":"indexer","job_subtype":"udf","state":"in_progress","created_at_millis":1000,"status":{"units_done":1,"units_total":2}}]}"#,
|
||||
)
|
||||
.unwrap()
|
||||
});
|
||||
let jobs = conn.list_jobs().await.unwrap();
|
||||
assert_eq!(jobs.len(), 1);
|
||||
assert_eq!(jobs[0].state, "running");
|
||||
assert_eq!(jobs[0].job_type, "udf");
|
||||
assert_eq!(jobs[0].units_total, Some(2));
|
||||
|
||||
// cancel_job
|
||||
let conn = Connection::new_with_handler(|request| {
|
||||
assert_eq!(request.method(), &reqwest::Method::POST);
|
||||
assert_eq!(request.url().path(), "/v1/job/j-3/cancel");
|
||||
http::Response::builder()
|
||||
// cancel_job: resolve via the manifest-id list filter, then cancel
|
||||
let conn = Connection::new_with_handler(|request| match request.url().path() {
|
||||
"/v1/jobs/list" => http::Response::builder()
|
||||
.status(200)
|
||||
.body(r#"{"cancelled":true}"#)
|
||||
.unwrap()
|
||||
.body(r#"{"jobs":[{"job_id":"plat-3","state":"in_progress"}]}"#)
|
||||
.unwrap(),
|
||||
"/v1/jobs/cancel" => {
|
||||
assert_eq!(request.method(), &reqwest::Method::POST);
|
||||
http::Response::builder()
|
||||
.status(200)
|
||||
.body(r#"{"job_id":"plat-3"}"#)
|
||||
.unwrap()
|
||||
}
|
||||
other => panic!("unexpected path {other}"),
|
||||
});
|
||||
assert!(conn.cancel_job("j-3").await.unwrap());
|
||||
|
||||
// cancel_job: no such inflight job -> false, not an error
|
||||
// cancel_job: never registered -> false, and no cancel request
|
||||
let conn = Connection::new_with_handler(|request| {
|
||||
assert_eq!(request.url().path(), "/v1/job/gone/cancel");
|
||||
http::Response::builder()
|
||||
.status(200)
|
||||
.body(r#"{"cancelled":false}"#)
|
||||
.unwrap()
|
||||
});
|
||||
assert!(!conn.cancel_job("gone").await.unwrap());
|
||||
|
||||
// job_history: GET /v1/job/history, no filter
|
||||
let conn = Connection::new_with_handler(|request| {
|
||||
assert_eq!(request.method(), &reqwest::Method::GET);
|
||||
assert_eq!(request.url().path(), "/v1/job/history");
|
||||
assert!(request.url().query().is_none());
|
||||
http::Response::builder()
|
||||
.status(200)
|
||||
.body(
|
||||
r#"{"jobs":[{"table":"docs","job_id":"j-1","job_type":"udf_virtual_column_backfill","state":"done","column":"vec","created_ms":1000,"updated_ms":2000,"completed_ms":2000,"rows_processed":42,"rows_skipped":3,"error":null,"events":"created\ndone"}]}"#,
|
||||
)
|
||||
.unwrap()
|
||||
});
|
||||
let hist = conn.job_history(None).await.unwrap();
|
||||
assert_eq!(hist.len(), 1);
|
||||
assert_eq!(hist[0].state, "done");
|
||||
assert_eq!(hist[0].rows_processed, Some(42));
|
||||
assert_eq!(hist[0].events.as_deref(), Some("created\ndone"));
|
||||
|
||||
// job_history: ?job= narrows to one job
|
||||
let conn = Connection::new_with_handler(|request| {
|
||||
assert_eq!(request.url().path(), "/v1/job/history");
|
||||
assert_eq!(request.url().query(), Some("job=j-1"));
|
||||
assert_eq!(request.url().path(), "/v1/jobs/list");
|
||||
http::Response::builder()
|
||||
.status(200)
|
||||
.body(r#"{"jobs":[]}"#)
|
||||
.unwrap()
|
||||
});
|
||||
assert!(!conn.cancel_job("gone").await.unwrap());
|
||||
|
||||
// job_history(None): a registry listing, timeline-free
|
||||
let conn = Connection::new_with_handler(|request| {
|
||||
assert_eq!(request.url().path(), "/v1/jobs/list");
|
||||
http::Response::builder()
|
||||
.status(200)
|
||||
.body(
|
||||
r#"{"jobs":[{"job_id":"plat-1","table":"docs","job_type":"indexer","job_subtype":"udf","state":"done","created_at_millis":1000,"status":{"rows_skipped":3}}]}"#,
|
||||
)
|
||||
.unwrap()
|
||||
});
|
||||
let hist = conn.job_history(None).await.unwrap();
|
||||
assert_eq!(hist.len(), 1);
|
||||
assert_eq!(hist[0].state, "finished");
|
||||
assert_eq!(hist[0].rows_skipped, Some(3));
|
||||
|
||||
// job_history(id): unknown everywhere -> empty
|
||||
let conn = Connection::new_with_handler(|request| match request.url().path() {
|
||||
"/v1/jobs/describe" => http::Response::builder().status(404).body("").unwrap(),
|
||||
"/v1/jobs/list" => http::Response::builder()
|
||||
.status(200)
|
||||
.body(r#"{"jobs":[]}"#)
|
||||
.unwrap(),
|
||||
other => panic!("unexpected path {other}"),
|
||||
});
|
||||
assert!(conn.job_history(Some("j-1")).await.unwrap().is_empty());
|
||||
|
||||
// errors: GET /v1/job/errors with job + table filters
|
||||
// errors: GET /v1/errors with job + table filters
|
||||
let conn = Connection::new_with_handler(|request| {
|
||||
assert_eq!(request.method(), &reqwest::Method::GET);
|
||||
assert_eq!(request.url().path(), "/v1/job/errors");
|
||||
assert_eq!(request.url().path(), "/v1/errors");
|
||||
assert_eq!(request.url().query(), Some("job=j-1&table=docs"));
|
||||
http::Response::builder()
|
||||
.status(200)
|
||||
|
||||
Reference in New Issue
Block a user