diff --git a/rust/lancedb/src/remote/db.rs b/rust/lancedb/src/remote/db.rs index e92d34ce3..1c03b4691 100644 --- a/rust/lancedb/src/remote/db.rs +++ b/rust/lancedb/src/remote/db.rs @@ -134,30 +134,6 @@ struct RemoteListMaterializedViewsResponse { views: Vec, } -#[derive(serde::Deserialize)] -struct RemoteJobEntry { - table: String, - job_id: String, - job_type: String, - state: String, - #[serde(default)] - column: Option, - #[serde(default)] - age_seconds: Option, - #[serde(default)] - command: Option, - #[serde(default)] - units_done: Option, - #[serde(default)] - units_total: Option, - #[serde(default)] - committed: bool, - #[serde(default)] - rows_skipped: u64, - #[serde(default)] - error: Option, -} - #[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, -} - -#[derive(serde::Deserialize)] -struct RemoteGetJobResponse { #[serde(default)] - job: Option, -} - -#[derive(serde::Deserialize)] -struct RemoteCancelJobResponse { - cancelled: bool, -} - -impl From 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, - created_ms: i64, - updated_ms: i64, + created_at_millis: i64, #[serde(default)] - completed_ms: Option, - #[serde(default)] - rows_processed: Option, - #[serde(default)] - rows_skipped: Option, - #[serde(default)] - error: Option, - #[serde(default)] - events: Option, + status: serde_json::Value, } -#[derive(serde::Deserialize)] -struct RemoteJobHistoryResponse { - jobs: Vec, -} - -impl From 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 { + status.get(key).and_then(serde_json::Value::as_i64) +} + +fn payload_error(status: &serde_json::Value) -> Option { + status + .get("error") + .and_then(serde_json::Value::as_str) + .map(str::to_string) } #[derive(serde::Deserialize)] @@ -1078,24 +1016,63 @@ impl Database for RemoteDatabase { } async fn list_jobs(&self) -> Result> { - 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> { - // 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 Database for RemoteDatabase { } async fn cancel_job(&self, job_id: &str) -> Result { - 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> { - 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::()); + let times = batch + .column_by_name("updated_at_millis") + .and_then(|c| c.as_any().downcast_ref::()); + let payloads = batch + .column_by_name("payload") + .and_then(|c| c.as_any().downcast_ref::()); + 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::(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> { - 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)