From 203f6536a6df56ccaee508f96b2cb92994477e3a Mon Sep 17 00:00:00 2001 From: Xuanwo Date: Wed, 12 Aug 2026 02:09:33 +0800 Subject: [PATCH] feat: expose typed remote job results --- rust/lancedb/src/database.rs | 6 + rust/lancedb/src/remote/db.rs | 545 +++++++++++++++++++++++++++++---- rust/lancedb/src/remote/job.rs | 211 +++++++++++-- 3 files changed, 677 insertions(+), 85 deletions(-) diff --git a/rust/lancedb/src/database.rs b/rust/lancedb/src/database.rs index f99f6e12a..a69c623a2 100644 --- a/rust/lancedb/src/database.rs +++ b/rust/lancedb/src/database.rs @@ -230,6 +230,12 @@ pub struct JobDescription { pub creation_ms: i64, /// The job-type-specific specification. Null when the server omits it. pub spec: serde_json::Value, + /// Explicit success result from the describe envelope, when present. + /// + /// Missing or JSON `null` wire `result` is [`None`]. An explicit + /// [`crate::JobResult::None`] object is `Some(JobResult::None)`. An exact + /// Function result is `Some(JobResult::Function(...))`. + pub result: Option, /// Why the job failed, when the job is failed and the server reports a /// reason. pub failure: Option, diff --git a/rust/lancedb/src/remote/db.rs b/rust/lancedb/src/remote/db.rs index 262b7d211..f3042ec0e 100644 --- a/rust/lancedb/src/remote/db.rs +++ b/rust/lancedb/src/remote/db.rs @@ -453,51 +453,6 @@ struct RemoteListJobsResponse { page_token: Option, } -/// The server's account of why a job failed. Absent from older servers, -/// which report only the terminal state. -#[derive(serde::Deserialize)] -struct RemoteReportedFailure { - /// Stable Function error category when the server supplied one. - #[serde(default)] - error_code: Option, - #[serde(default)] - phase: Option, - #[serde(default)] - message: Option, - #[serde(default)] - retryable: Option, -} - -#[derive(serde::Deserialize)] -struct RemoteDescribeJobResponse { - job_id: String, - #[serde(default)] - job_type: String, - job_state: String, - #[serde(default)] - creation_ms: i64, - #[serde(default)] - spec: serde_json::Value, - #[serde(default)] - failure: Option, -} - -/// Server job states -> the client vocabulary ("running" / "finished" / -/// "failed" / "cancelled"). Covers both the describe enum (IN_PROGRESS / -/// DONE / FAILED / CANCELLED) and the registry's lowercase list-row states -/// (in_progress / succeeded / failed / canceled / timed_out). States this -/// client version does not know (e.g. created, queued) pass through as-is. -fn job_state_to_client(state: &str) -> String { - match state { - "IN_PROGRESS" | "in_progress" => "running", - "DONE" | "done" | "succeeded" => "finished", - "FAILED" | "failed" | "TIMED_OUT" | "timed_out" => "failed", - "CANCELLED" | "cancelled" | "canceled" => "cancelled", - other => other, - } - .to_string() -} - /// Bound on `list_jobs` page walking; a warning is logged when the listing /// is truncated at this many pages. const MAX_LIST_JOBS_PAGES: usize = 100; @@ -538,7 +493,7 @@ impl Database for RemoteDatabase { job_id: row.job_id, table: row.table, job_type: row.job_type, - state: job_state_to_client(&row.state), + state: super::job::job_state_to_client(&row.state), created_at_millis: row.created_at_millis, })); page_token = body.page_token; @@ -570,20 +525,23 @@ impl Database for RemoteDatabase { }) => return Ok(None), Err(err) => return Err(err), }; - let body: RemoteDescribeJobResponse = rsp.json().await.err_to_http(request_id)?; + let body: super::job::DescribeJobResponse = + rsp.json().await.err_to_http(request_id.clone())?; + let job_id = body.require_job_id(request_id.clone())?; + let result = body + .project_success_result(request_id)? + .into_description_result(); + let state = body.client_state(); Ok(Some(JobDescription { - job_id: body.job_id, + job_id, job_type: body.job_type, - state: job_state_to_client(&body.job_state), + state, creation_ms: body.creation_ms, spec: body.spec, - failure: body.failure.map(|reported| crate::error::JobFailure { - error_code: reported.error_code, - phase: reported.phase, - message: reported.message, - retryable: reported.retryable, - source: None, - }), + result, + failure: body + .failure + .map(super::job::ReportedFailure::into_job_failure), })) } @@ -1118,8 +1076,12 @@ mod tests { use crate::{ Connection, Error, database::CreateTableMode, + error::FunctionErrorCode, + function::{Function, FunctionId, FunctionOutput, FunctionParameter, FunctionSignature}, + job::JobResult, remote::{ARROW_STREAM_CONTENT_TYPE, ClientConfig, HeaderProvider, JSON_CONTENT_TYPE}, }; + use serde_json::{Value, json}; #[test] fn test_cache_key_security() { @@ -2336,6 +2298,67 @@ mod tests { assert_eq!(failure.retryable, Some(true)); } + /// Typed get_job keeps requiring an echoed response job_id even when the + /// rest of a create_index DONE describe looks valid. + #[tokio::test] + async fn test_get_job_omitted_job_id_is_http() { + let conn = Connection::new_with_handler(|request| { + assert_eq!(request.url().path(), "/v1/jobs/describe"); + http::Response::builder() + .status(200) + .body(r#"{"job_type":"create_index","job_state":"DONE","creation_ms":1,"spec":{}}"#) + .unwrap() + }); + let err = conn + .get_job("job-1") + .await + .expect_err("get_job must require response job_id"); + match err { + Error::Http { .. } => {} + other => panic!("expected Error::Http, got {other:?}"), + } + } + + /// Documented lowercase describe/list aliases stay stable on get_job. + #[tokio::test] + async fn test_get_job_normalizes_documented_state_aliases() { + let cases = [ + ("in_progress", "running"), + ("done", "finished"), + ("succeeded", "finished"), + ("failed", "failed"), + ("timed_out", "failed"), + ("cancelled", "cancelled"), + ("canceled", "cancelled"), + ("IN_PROGRESS", "running"), + ("DONE", "finished"), + ("FAILED", "failed"), + ("CANCELLED", "cancelled"), + ("TIMED_OUT", "failed"), + ]; + for (wire_state, expected_client_state) in cases { + let body = format!( + r#"{{"job_id":"job-alias","job_type":"create_index","job_state":"{wire_state}","creation_ms":1,"spec":{{}}}}"# + ); + let conn = Connection::new_with_handler(move |_| { + http::Response::builder() + .status(200) + .body(body.clone()) + .unwrap() + }); + let job = conn + .get_job("job-alias") + .await + .unwrap_or_else(|err| panic!("alias {wire_state} must describe: {err:?}")) + .expect("job must exist"); + assert_eq!( + job.state, expected_client_state, + "get_job alias {wire_state} must normalize to {expected_client_state}" + ); + assert_eq!(job.job_id, "job-alias"); + } + } + #[tokio::test] async fn test_get_job_missing_is_none() { let conn = Connection::new_with_handler(|_| { @@ -2529,4 +2552,414 @@ mod tests { other => panic!("expected Error::JobFailed, got {other:?}"), } } + + /// Optional JSON object field for deterministic `/v1/jobs/describe` fixtures. + #[derive(Clone)] + enum JsonField { + Absent, + Null, + Present(Value), + } + + fn sample_description_function() -> Function { + let id = FunctionId::try_new("fn.exact.remote-job-description").expect("valid FunctionId"); + let signature = FunctionSignature::try_new( + vec![ + FunctionParameter::new("x", DataType::Int32), + FunctionParameter::new("label", DataType::Utf8), + ], + FunctionOutput::new(DataType::Int32, true), + ) + .expect("valid FunctionSignature"); + Function::new(id, signature) + } + + fn assert_exact_function(actual: &Function, expected: &Function) { + assert_eq!(actual.id(), expected.id()); + assert_eq!(actual.signature(), expected.signature()); + } + + fn job_result_none_wire() -> Value { + serde_json::to_value(JobResult::None).expect("serialize JobResult::None wire") + } + + fn job_result_function_wire(function: &Function) -> Value { + serde_json::to_value(JobResult::Function(function.clone())) + .expect("serialize JobResult::Function wire") + } + + fn describe_body( + job_id: &str, + job_state: &str, + job_type: JsonField, + result: JsonField, + ) -> String { + let mut body = json!({ + "job_id": job_id, + "job_state": job_state, + "creation_ms": 1, + "spec": {}, + }); + let object = body + .as_object_mut() + .expect("describe body must be a JSON object"); + match job_type { + JsonField::Absent => {} + JsonField::Null => { + object.insert("job_type".into(), Value::Null); + } + JsonField::Present(value) => { + object.insert("job_type".into(), value); + } + } + match result { + JsonField::Absent => {} + JsonField::Null => { + object.insert("result".into(), Value::Null); + } + JsonField::Present(value) => { + object.insert("result".into(), value); + } + } + body.to_string() + } + + fn conn_with_describe_body(body: String) -> Connection { + Connection::new_with_handler(move |request| { + assert_eq!(request.method(), &reqwest::Method::POST); + assert_eq!(request.url().path(), "/v1/jobs/describe"); + http::Response::builder() + .status(200) + .body(body.clone()) + .unwrap() + }) + } + + /// register_function DONE Function is shared by get_job and job(id).wait(). + #[tokio::test] + async fn remote_job_description_result_register_function_done_matches_wait() { + let expected = sample_description_function(); + let body = describe_body( + "job-register", + "DONE", + JsonField::Present(Value::String("register_function".into())), + JsonField::Present(job_result_function_wire(&expected)), + ); + let conn = conn_with_describe_body(body); + + let description = conn + .get_job("job-register") + .await + .expect("valid register_function describe must succeed") + .expect("job must exist"); + let from_get = description + .result + .as_ref() + .and_then(JobResult::function) + .expect("register_function success must be Some(Function)"); + assert_exact_function(from_get, &expected); + + let waited = conn + .job("job-register") + .expect("job handle") + .wait() + .await + .expect("wait over the same fixture must succeed"); + let from_wait = waited + .function() + .expect("wait must return the exact Function"); + assert_exact_function(from_wait, &expected); + } + + /// Known no-result DONE: omitted/null stay None; explicit None is Some(None). + #[tokio::test] + async fn remote_job_description_result_known_no_result_omission_vs_explicit_none() { + for result_field in [JsonField::Absent, JsonField::Null] { + let body = describe_body( + "job-create-index", + "DONE", + JsonField::Present(Value::String("create_index".into())), + result_field, + ); + let description = conn_with_describe_body(body) + .get_job("job-create-index") + .await + .expect("known no-result describe must succeed") + .expect("job must exist"); + assert_eq!(description.result, None); + } + + let body = describe_body( + "job-create-index-explicit", + "DONE", + JsonField::Present(Value::String("create_index".into())), + JsonField::Present(job_result_none_wire()), + ); + let description = conn_with_describe_body(body) + .get_job("job-create-index-explicit") + .await + .expect("explicit None describe must succeed") + .expect("job must exist"); + assert_eq!(description.result, Some(JobResult::None)); + } + + /// Unknown nonterminal and ordinary FAILED/CANCELLED stay describable with no result. + #[tokio::test] + async fn remote_job_description_result_nonterminal_unknown_and_failed_cancelled_without_result() + { + let unknown_running = describe_body( + "job-future-running", + "IN_PROGRESS", + JsonField::Present(Value::String("future_job_type_xyz".into())), + JsonField::Absent, + ); + let description = conn_with_describe_body(unknown_running) + .get_job("job-future-running") + .await + .expect("unknown nonterminal without result must remain describable") + .expect("job must exist"); + assert_eq!(description.job_type, "future_job_type_xyz"); + assert_eq!(description.state, "running"); + assert_eq!(description.result, None); + + for (job_id, job_state, client_state) in [ + ("job-failed", "FAILED", "failed"), + ("job-cancelled", "CANCELLED", "cancelled"), + ] { + let body = describe_body( + job_id, + job_state, + JsonField::Present(Value::String("create_index".into())), + JsonField::Absent, + ); + let description = conn_with_describe_body(body) + .get_job(job_id) + .await + .expect("FAILED/CANCELLED without result must remain describable") + .expect("job must exist"); + assert_eq!(description.state, client_state); + assert_eq!( + description.result, None, + "must not invent a result for {job_state}" + ); + } + } + + /// get_job rejects the same strict invalid describe shapes as RemoteJob::wait. + #[tokio::test] + async fn remote_job_description_result_strict_invalid_cases_are_http() { + let function = sample_description_function(); + let function_wire = job_result_function_wire(&function); + let none_wire = job_result_none_wire(); + + let mut unknown_kind = none_wire.clone(); + unknown_kind["kind"] = Value::String("artifact".into()); + + let mut unknown_version = none_wire.clone(); + unknown_version["format_version"] = Value::from(2); + + let mut unknown_outer_field = none_wire.clone(); + unknown_outer_field + .as_object_mut() + .unwrap() + .insert("unexpected_field".into(), Value::Bool(true)); + + let mut empty_function_id = function_wire.clone(); + empty_function_id["function"]["id"] = Value::String("".into()); + + let cases = [ + ( + "register_missing", + "DONE", + JsonField::Present(Value::String("register_function".into())), + JsonField::Absent, + ), + ( + "register_null", + "DONE", + JsonField::Present(Value::String("register_function".into())), + JsonField::Null, + ), + ( + "register_explicit_none", + "DONE", + JsonField::Present(Value::String("register_function".into())), + JsonField::Present(none_wire.clone()), + ), + ( + "known_no_result_with_function", + "DONE", + JsonField::Present(Value::String("create_index".into())), + JsonField::Present(function_wire.clone()), + ), + ( + "unknown_done_missing", + "DONE", + JsonField::Present(Value::String("future_job_type_xyz".into())), + JsonField::Absent, + ), + ( + "unknown_done_explicit_none", + "DONE", + JsonField::Present(Value::String("future_job_type_xyz".into())), + JsonField::Present(none_wire.clone()), + ), + ( + "unknown_done_explicit_function", + "DONE", + JsonField::Present(Value::String("future_job_type_xyz".into())), + JsonField::Present(function_wire.clone()), + ), + ( + "malformed_unknown_kind", + "DONE", + JsonField::Present(Value::String("register_function".into())), + JsonField::Present(unknown_kind), + ), + ( + "malformed_unknown_version", + "DONE", + JsonField::Present(Value::String("register_function".into())), + JsonField::Present(unknown_version), + ), + ( + "malformed_unknown_outer_field", + "DONE", + JsonField::Present(Value::String("register_function".into())), + JsonField::Present(unknown_outer_field), + ), + ( + "malformed_empty_function_id", + "DONE", + JsonField::Present(Value::String("register_function".into())), + JsonField::Present(empty_function_id), + ), + ( + "failed_carrying_function", + "FAILED", + JsonField::Present(Value::String("register_function".into())), + JsonField::Present(function_wire), + ), + ( + "cancelled_carrying_none", + "CANCELLED", + JsonField::Present(Value::String("create_index".into())), + JsonField::Present(none_wire), + ), + ]; + + let mut unexpected = Vec::new(); + for (job_id, job_state, job_type, result_field) in cases { + let body = describe_body(job_id, job_state, job_type, result_field); + match conn_with_describe_body(body).get_job(job_id).await { + Err(Error::Http { .. }) => {} + Ok(Some(description)) => unexpected.push(format!( + "{job_id}: Ok(Some(result={:?}))", + description.result + )), + Ok(None) => unexpected.push(format!("{job_id}: Ok(None)")), + Err(other) => unexpected.push(format!("{job_id}: Err({other:?})")), + } + } + assert!( + unexpected.is_empty(), + "strict invalid get_job cases must be Error::Http: {unexpected:?}" + ); + } + + /// Unknown informational outer fields are tolerated for a valid Function description. + #[tokio::test] + async fn remote_job_description_result_unknown_outer_fields_tolerated() { + let expected = sample_description_function(); + let body = json!({ + "job_id": "job-register-extra", + "job_state": "DONE", + "job_type": "register_function", + "creation_ms": 42, + "spec": {"ignored": true}, + "result": job_result_function_wire(&expected), + "server_note": "informational-only", + "extra_admin_field": 7, + }) + .to_string(); + let description = conn_with_describe_body(body) + .get_job("job-register-extra") + .await + .expect("unknown outer fields must not block description") + .expect("job must exist"); + assert_eq!(description.job_id, "job-register-extra"); + assert_eq!(description.creation_ms, 42); + let function = description + .result + .as_ref() + .and_then(JobResult::function) + .expect("expected Some(Function)"); + assert_exact_function(function, &expected); + } + + /// Existing description fields and stable failure error_code stay intact with absent result. + #[tokio::test] + async fn remote_job_description_result_existing_fields_intact_when_result_absent() { + let body = json!({ + "job_id": "job-1", + "job_type": "create_index", + "job_state": "FAILED", + "creation_ms": 1000, + "spec": {"column": "vec"}, + "failure": { + "error_code": "name_or_function_not_found", + "phase": "validate", + "message": "looks like definition_validation_failure", + "retryable": false + } + }) + .to_string(); + let description = conn_with_describe_body(body) + .get_job("job-1") + .await + .expect("failure description without result must succeed") + .expect("job must exist"); + assert_eq!(description.job_id, "job-1"); + assert_eq!(description.job_type, "create_index"); + assert_eq!(description.state, "failed"); + assert_eq!(description.creation_ms, 1000); + assert_eq!(description.spec["column"], "vec"); + assert_eq!(description.result, None); + let failure = description.failure.expect("failure payload present"); + match &failure.error_code { + Some(code) => { + assert_eq!(code, &FunctionErrorCode::NameOrFunctionNotFound); + assert_ne!(code, &FunctionErrorCode::DefinitionValidationFailure); + } + None => panic!("known error_code must remain decoded while result is absent"), + } + assert_eq!(failure.phase.as_deref(), Some("validate")); + assert_eq!(failure.retryable, Some(false)); + } + + /// Fixture wires round-trip through the pinned public JobResult serde. + #[test] + fn remote_job_description_result_fixture_wires_match_job_result_serde() { + let none_wire = job_result_none_wire(); + let none: JobResult = + serde_json::from_value(none_wire.clone()).expect("None wire must decode"); + assert_eq!(none, JobResult::None); + assert_eq!( + serde_json::to_value(JobResult::None).expect("serialize None"), + none_wire + ); + + let expected = sample_description_function(); + let function_wire = job_result_function_wire(&expected); + let decoded: JobResult = + serde_json::from_value(function_wire.clone()).expect("Function wire must decode"); + match decoded { + JobResult::Function(function) => assert_exact_function(&function, &expected), + JobResult::None => panic!("Function fixture must not decode as None"), + } + assert_eq!( + serde_json::to_value(JobResult::Function(expected)).expect("serialize Function"), + function_wire + ); + } } diff --git a/rust/lancedb/src/remote/job.rs b/rust/lancedb/src/remote/job.rs index 93ee845d9..06736e73f 100644 --- a/rust/lancedb/src/remote/job.rs +++ b/rust/lancedb/src/remote/job.rs @@ -36,7 +36,9 @@ impl<'de> Deserialize<'de> for JobState { } impl JobState { - /// The client vocabulary label for this state. + /// Forward-observing label for [`RemoteJob::status`]: known describe + /// variants map to the client vocabulary; unrecognized/lowercase wire + /// labels pass through unchanged. fn client_label(&self) -> String { match self { Self::InProgress => "running".to_string(), @@ -63,6 +65,22 @@ impl From<&str> for JobState { } } +/// Server job states -> the client vocabulary ("running" / "finished" / +/// "failed" / "cancelled"). Covers both the describe enum (IN_PROGRESS / +/// DONE / FAILED / CANCELLED) and the registry's lowercase list-row states +/// (in_progress / succeeded / failed / canceled / timed_out). States this +/// client version does not know (e.g. created, queued) pass through as-is. +pub(super) fn job_state_to_client(state: &str) -> String { + match state { + "IN_PROGRESS" | "in_progress" => "running", + "DONE" | "done" | "succeeded" => "finished", + "FAILED" | "failed" | "TIMED_OUT" | "timed_out" => "failed", + "CANCELLED" | "cancelled" | "canceled" => "cancelled", + other => other, + } + .to_string() +} + /// Closed current no-result job_type vocabulary. const KNOWN_NO_RESULT_JOB_TYPES: &[&str] = &[ "create_index", @@ -102,7 +120,7 @@ enum ExpectedSuccessResult { /// The server's account of why a job failed. Absent from older servers, which /// report only the terminal state. #[derive(Deserialize)] -struct ReportedFailure { +pub(super) struct ReportedFailure { /// Stable Function error category when the server supplied one. #[serde(default)] error_code: Option, @@ -114,17 +132,98 @@ struct ReportedFailure { retryable: Option, } +impl ReportedFailure { + pub(super) fn into_job_failure(self) -> JobFailure { + JobFailure { + error_code: self.error_code, + phase: self.phase, + message: self.message, + retryable: self.retryable, + source: None, + } + } +} + /// Outer `/v1/jobs/describe` envelope. Unknown informational fields are ignored; /// `result` stays raw so lifecycle reads do not depend on JobResult decoding. +/// +/// Response `job_id` is optional so [`RemoteJob`] status/wait can parse +/// state-only envelopes. Typed [`crate::database::JobDescription`] paths must +/// call [`DescribeJobResponse::require_job_id`]. #[derive(Deserialize)] -struct DescribeJobResponse { +pub(super) struct DescribeJobResponse { + #[serde(default)] + job_id: Option, job_state: JobState, #[serde(default)] - job_type: String, + pub(super) job_type: String, + #[serde(default)] + pub(super) creation_ms: i64, + #[serde(default)] + pub(super) spec: serde_json::Value, #[serde(default)] result: Option, #[serde(default)] - failure: Option, + pub(super) failure: Option, +} + +impl DescribeJobResponse { + /// Require a non-empty echoed response `job_id` for typed get_job. + pub(super) fn require_job_id(&self, request_id: String) -> Result { + match self.job_id.as_deref() { + Some(job_id) if !job_id.is_empty() => Ok(job_id.to_string()), + _ => Err(protocol_http( + request_id, + "describe response missing job_id", + )), + } + } + + /// Description/list-aligned client label, including documented lowercase + /// aliases. Distinct from [`JobState::client_label`] used by status. + pub(super) fn client_state(&self) -> String { + match &self.job_state { + JobState::Other(raw) => job_state_to_client(raw), + state => state.client_label(), + } + } + + /// Strict success-result projection for describe: Missing versus Present. + pub(super) fn project_success_result(&self, request_id: String) -> Result { + project_describe_success_result( + &self.job_state, + &self.job_type, + self.result.as_ref(), + request_id, + ) + } +} + +/// Wire success-result presence after strict job-type expectation checks. +/// +/// Distinguishes a missing/null wire `result` from an explicit decoded +/// [`JobResult`] object so [`crate::database::JobDescription`] can keep +/// `None` versus `Some(JobResult::None)`. +#[derive(Debug, Clone, PartialEq, Eq)] +pub(super) enum RawSuccessResult { + Missing, + Present(JobResult), +} + +impl RawSuccessResult { + pub(super) fn into_description_result(self) -> Option { + match self { + Self::Missing => None, + Self::Present(result) => Some(result), + } + } + + fn into_wait_result(self) -> JobResult { + match self { + Self::Missing => JobResult::None, + Self::Present(result) => result, + } + } } fn protocol_http(request_id: String, message: impl Into) -> Error { @@ -148,12 +247,12 @@ fn expected_success_result(job_type: &str) -> Option { None } -/// Strict DONE projection of a describe payload into [`JobResult`]. -fn project_done_job_result( +/// Strict DONE projection that preserves Missing versus Present([`JobResult`]). +fn project_done_raw_result( job_type: &str, raw_result: Option<&serde_json::Value>, request_id: String, -) -> Result { +) -> Result { let Some(expected) = expected_success_result(job_type) else { return Err(protocol_http( request_id, @@ -167,7 +266,7 @@ fn project_done_job_result( "register_function DONE response missing Function result", )), (ExpectedSuccessResult::None, None) | (ExpectedSuccessResult::AbsentResultOnly, None) => { - Ok(JobResult::None) + Ok(RawSuccessResult::Missing) } (ExpectedSuccessResult::AbsentResultOnly, Some(_)) => Err(protocol_http( request_id, @@ -179,9 +278,11 @@ fn project_done_job_result( }; match (expected, decoded) { (ExpectedSuccessResult::Function, JobResult::Function(function)) => { - Ok(JobResult::Function(function)) + Ok(RawSuccessResult::Present(JobResult::Function(function))) + } + (ExpectedSuccessResult::None, JobResult::None) => { + Ok(RawSuccessResult::Present(JobResult::None)) } - (ExpectedSuccessResult::None, JobResult::None) => Ok(JobResult::None), _ => Err(protocol_http( request_id, "job result kind does not match job_type expectation", @@ -191,6 +292,24 @@ fn project_done_job_result( } } +fn project_describe_success_result( + job_state: &JobState, + job_type: &str, + raw_result: Option<&serde_json::Value>, + request_id: String, +) -> Result { + if !matches!(job_state, JobState::Done) { + if raw_result.is_some() { + return Err(protocol_http( + request_id, + "non-DONE job describe response carried a success result", + )); + } + return Ok(RawSuccessResult::Missing); + } + project_done_raw_result(job_type, raw_result, request_id) +} + pub struct RemoteJob { client: RestfulLanceDbClient, job_id: String, @@ -234,32 +353,17 @@ impl JobHandle for RemoteJob { let mut interval = INITIAL_POLL_INTERVAL; loop { let (request_id, description) = self.describe().await?; - if !matches!(description.job_state, JobState::Done) && description.result.is_some() { - return Err(protocol_http( - request_id, - "non-DONE job describe response carried a success result", - )); - } + let projected = description.project_success_result(request_id)?; match description.job_state { JobState::Done => { - return project_done_job_result( - &description.job_type, - description.result.as_ref(), - request_id, - ); + return Ok(projected.into_wait_result()); } JobState::Failed => { return Err(Error::JobFailed { job_id: Some(self.job_id.clone()), failure: description .failure - .map(|reported| JobFailure { - error_code: reported.error_code, - phase: reported.phase, - message: reported.message, - retryable: reported.retryable, - source: None, - }) + .map(ReportedFailure::into_job_failure) .unwrap_or_default(), }); } @@ -317,6 +421,55 @@ mod tests { assert_eq!(result, JobResult::None); } + /// Lifecycle reads accept a state-only DONE envelope: no echoed job_id, + /// job_type, creation_ms, spec, failure, or result. + #[tokio::test] + async fn remote_job_lifecycle_state_only_done_without_job_id() { + let client = client_with_handler(|_| { + http::Response::builder() + .status(200) + .body(r#"{"job_state":"DONE"}"#) + .unwrap() + }); + let status_job = RemoteJob::new(client, "job-state-only".into()); + let status = status_job + .status() + .await + .expect("state-only DONE must not require response job_id"); + assert_eq!(status, "finished"); + + let client = client_with_handler(|_| { + http::Response::builder() + .status(200) + .body(r#"{"job_state":"DONE"}"#) + .unwrap() + }); + let wait_job = RemoteJob::new(client, "job-state-only".into()); + let result = wait_job + .wait() + .await + .expect("historical missing job_type/result DONE must project None"); + assert_eq!(result, JobResult::None); + } + + /// status keeps forward-observing lowercase/unrecognized labels; it does + /// not apply get_job's documented alias normalization. + #[tokio::test] + async fn remote_job_lifecycle_status_preserves_lowercase_state_label() { + let client = client_with_handler(|_| { + http::Response::builder() + .status(200) + .body(r#"{"job_id":"job-lower","job_state":"done"}"#) + .unwrap() + }); + let job = RemoteJob::new(client, "job-lower".into()); + let status = job + .status() + .await + .expect("lowercase describe state must remain readable"); + assert_eq!(status, "done"); + } + #[tokio::test] async fn wait_decodes_known_failure_error_code() { let client = client_with_handler(|_| {