From 9d3d0d06401d42ada4493eec971df0c01415a6fa Mon Sep 17 00:00:00 2001 From: Xuanwo Date: Wed, 12 Aug 2026 01:39:51 +0800 Subject: [PATCH] feat: decode remote job results --- rust/lancedb/src/remote/job.rs | 592 ++++++++++++++++++++++++++++++++- 1 file changed, 584 insertions(+), 8 deletions(-) diff --git a/rust/lancedb/src/remote/job.rs b/rust/lancedb/src/remote/job.rs index 8bd387f96..93ee845d9 100644 --- a/rust/lancedb/src/remote/job.rs +++ b/rust/lancedb/src/remote/job.rs @@ -63,6 +63,42 @@ impl From<&str> for JobState { } } +/// Closed current no-result job_type vocabulary. +const KNOWN_NO_RESULT_JOB_TYPES: &[&str] = &[ + "create_index", + "reindex_ivf_pq", + "reindex_ivf_flat", + "reindex_ivf_rq", + "reindex_ivf_hnsw_sq", + "reindex_btree", + "reindex_fts", + "reindex_fm", + "reindex_bitmap", + "reindex_label_list", + "reindex_zonemap", + "reindex_ngram", + "reindex_bloom_filter", + "reindex_rtree", + "compact", + "cleanup", + "index_remap", + "spfresh_merge", + "prewarm_page_cache", + "prewarm_index_cache", + "job_registry_archive", +]; + +const REGISTER_FUNCTION_JOB_TYPE: &str = "register_function"; + +/// Closed success-result expectation for a describe `job_type`. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum ExpectedSuccessResult { + Function, + None, + /// Historical empty/missing job_type: only missing/null wire result. + AbsentResultOnly, +} + /// The server's account of why a job failed. Absent from older servers, which /// report only the terminal state. #[derive(Deserialize)] @@ -78,13 +114,83 @@ struct ReportedFailure { retryable: Option, } +/// Outer `/v1/jobs/describe` envelope. Unknown informational fields are ignored; +/// `result` stays raw so lifecycle reads do not depend on JobResult decoding. #[derive(Deserialize)] struct DescribeJobResponse { job_state: JobState, #[serde(default)] + job_type: String, + #[serde(default)] + result: Option, + #[serde(default)] failure: Option, } +fn protocol_http(request_id: String, message: impl Into) -> Error { + Error::Http { + source: message.into().into(), + request_id, + status_code: None, + } +} + +fn expected_success_result(job_type: &str) -> Option { + if job_type == REGISTER_FUNCTION_JOB_TYPE { + return Some(ExpectedSuccessResult::Function); + } + if job_type.is_empty() { + return Some(ExpectedSuccessResult::AbsentResultOnly); + } + if KNOWN_NO_RESULT_JOB_TYPES.contains(&job_type) { + return Some(ExpectedSuccessResult::None); + } + None +} + +/// Strict DONE projection of a describe payload into [`JobResult`]. +fn project_done_job_result( + job_type: &str, + raw_result: Option<&serde_json::Value>, + request_id: String, +) -> Result { + let Some(expected) = expected_success_result(job_type) else { + return Err(protocol_http( + request_id, + "unknown job_type for success result projection", + )); + }; + + match (expected, raw_result) { + (ExpectedSuccessResult::Function, None) => Err(protocol_http( + request_id, + "register_function DONE response missing Function result", + )), + (ExpectedSuccessResult::None, None) | (ExpectedSuccessResult::AbsentResultOnly, None) => { + Ok(JobResult::None) + } + (ExpectedSuccessResult::AbsentResultOnly, Some(_)) => Err(protocol_http( + request_id, + "historical empty job_type cannot carry an explicit success result", + )), + (ExpectedSuccessResult::Function, Some(raw)) | (ExpectedSuccessResult::None, Some(raw)) => { + let Ok(decoded) = serde_json::from_value::(raw.clone()) else { + return Err(protocol_http(request_id, "failed to decode job result")); + }; + match (expected, decoded) { + (ExpectedSuccessResult::Function, JobResult::Function(function)) => { + Ok(JobResult::Function(function)) + } + (ExpectedSuccessResult::None, JobResult::None) => Ok(JobResult::None), + _ => Err(protocol_http( + request_id, + "job result kind does not match job_type expectation", + )), + } + } + } +} + pub struct RemoteJob { client: RestfulLanceDbClient, job_id: String, @@ -96,7 +202,7 @@ impl RemoteJob { } /// One `/v1/jobs/describe` round trip. - async fn describe(&self) -> Result { + async fn describe(&self) -> Result<(String, DescribeJobResponse)> { let request = self .client .post("/v1/jobs/describe") @@ -107,10 +213,10 @@ impl RemoteJob { let description: DescribeJobResponse = serde_json::from_str(&body).map_err(|e| Error::Http { source: format!("failed to parse job description: {}", e).into(), - request_id, + request_id: request_id.clone(), status_code: None, })?; - Ok(description) + Ok((request_id, description)) } } @@ -121,17 +227,27 @@ impl JobHandle for RemoteJob { } async fn status(&self) -> Result { - Ok(self.describe().await?.job_state.client_label()) + Ok(self.describe().await?.1.job_state.client_label()) } async fn wait(&self) -> Result { let mut interval = INITIAL_POLL_INTERVAL; loop { - let description = self.describe().await?; + 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", + )); + } match description.job_state { - // Existing DONE responses have no success-result field; map to - // None until strict result decoding is added. - JobState::Done => return Ok(JobResult::None), + JobState::Done => { + return project_done_job_result( + &description.job_type, + description.result.as_ref(), + request_id, + ); + } JobState::Failed => { return Err(Error::JobFailed { job_id: Some(self.job_id.clone()), @@ -179,8 +295,13 @@ impl JobHandle for RemoteJob { mod tests { use super::*; use crate::error::FunctionErrorCode; + use crate::function::{ + Function, FunctionId, FunctionOutput, FunctionParameter, FunctionSignature, + }; use crate::job::JobResult; use crate::remote::client::test_utils::client_with_handler; + use arrow_schema::DataType; + use serde_json::{Value, json}; /// Terminal DONE with no result field projects as [`JobResult::None`]. #[tokio::test] @@ -272,4 +393,459 @@ mod tests { other => panic!("expected Error::JobFailed, got {other:?}"), } } + + /// Closed current no-result job_type vocabulary. + const KNOWN_NO_RESULT_JOB_TYPES: &[&str] = &[ + "create_index", + "reindex_ivf_pq", + "reindex_ivf_flat", + "reindex_ivf_rq", + "reindex_ivf_hnsw_sq", + "reindex_btree", + "reindex_fts", + "reindex_fm", + "reindex_bitmap", + "reindex_label_list", + "reindex_zonemap", + "reindex_ngram", + "reindex_bloom_filter", + "reindex_rtree", + "compact", + "cleanup", + "index_remap", + "spfresh_merge", + "prewarm_page_cache", + "prewarm_index_cache", + "job_registry_archive", + ]; + + #[derive(Clone)] + enum JsonField { + Absent, + Null, + Present(Value), + } + + fn sample_remote_function() -> Function { + let id = FunctionId::try_new("fn.exact.remote-job-result").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 assert_http_protocol_err(err: Error) { + match err { + Error::Http { .. } => {} + other => panic!("expected Error::Http protocol failure, got {other:?}"), + } + } + + 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, + }); + 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 remote_job_with_describe_body(job_id: &str, body: String) -> RemoteJob { + let client = client_with_handler(move |_| { + http::Response::builder() + .status(200) + .body(body.clone()) + .unwrap() + }); + RemoteJob::new(client, job_id.into()) + } + + async fn wait_expect_none(job: &RemoteJob) { + let result = job + .wait() + .await + .expect("expected successful None projection"); + assert_eq!(result, JobResult::None); + } + + async fn wait_expect_http(job: &RemoteJob) { + let err = job + .wait() + .await + .expect_err("expected remote protocol Error::Http"); + assert_http_protocol_err(err); + } + + /// register_function DONE with a Function result returns the exact ID and signature. + #[tokio::test] + async fn remote_job_result_register_function_done_returns_exact_function() { + let expected = sample_remote_function(); + let body = describe_body( + "job-register", + "DONE", + JsonField::Present(Value::String("register_function".into())), + JsonField::Present(job_result_function_wire(&expected)), + ); + let job = remote_job_with_describe_body("job-register", body); + let result = job + .wait() + .await + .expect("register_function DONE with Function must succeed"); + match result { + JobResult::Function(function) => assert_exact_function(&function, &expected), + JobResult::None => panic!("register_function success must not project as None"), + } + } + + /// register_function DONE with missing, null, or explicit None is protocol Http. + #[tokio::test] + async fn remote_job_result_register_function_missing_null_or_explicit_none_is_http() { + let cases = [ + ("absent", JsonField::Absent), + ("null", JsonField::Null), + ("explicit_none", JsonField::Present(job_result_none_wire())), + ]; + let mut unexpected = Vec::new(); + for (label, result_field) in cases { + let body = describe_body( + "job-register-missing", + "DONE", + JsonField::Present(Value::String("register_function".into())), + result_field, + ); + let job = remote_job_with_describe_body("job-register-missing", body); + match job.wait().await { + Err(Error::Http { .. }) => {} + Ok(value) => unexpected.push(format!("{label}: Ok({value:?})")), + Err(other) => unexpected.push(format!("{label}: Err({other:?})")), + } + } + assert!( + unexpected.is_empty(), + "register_function without Function must be Error::Http for every case: {unexpected:?}" + ); + } + + /// create_index accepts explicit None or missing; every known no-result type accepts missing. + #[tokio::test] + async fn remote_job_result_known_no_result_types_project_none() { + let create_index_cases = [ + JsonField::Absent, + JsonField::Present(job_result_none_wire()), + ]; + for result_field in create_index_cases { + let body = describe_body( + "job-create-index", + "DONE", + JsonField::Present(Value::String("create_index".into())), + result_field, + ); + let job = remote_job_with_describe_body("job-create-index", body); + wait_expect_none(&job).await; + } + + for job_type in KNOWN_NO_RESULT_JOB_TYPES { + let body = describe_body( + "job-no-result", + "DONE", + JsonField::Present(Value::String((*job_type).into())), + JsonField::Absent, + ); + let job = remote_job_with_describe_body("job-no-result", body); + wait_expect_none(&job).await; + } + } + + /// A known no-result job_type carrying Function is protocol Http. + #[tokio::test] + async fn remote_job_result_known_no_result_with_function_is_http() { + let function = sample_remote_function(); + let body = describe_body( + "job-create-index-function", + "DONE", + JsonField::Present(Value::String("create_index".into())), + JsonField::Present(job_result_function_wire(&function)), + ); + let job = remote_job_with_describe_body("job-create-index-function", body); + wait_expect_http(&job).await; + } + + /// Unknown job_type fails wait; empty historical job_type succeeds only without an explicit result. + #[tokio::test] + async fn remote_job_result_unknown_or_empty_job_type_expectation() { + // Historical empty job_type without an explicit result remains None. + for result_field in [JsonField::Absent, JsonField::Null] { + let body = describe_body( + "job-empty-type", + "DONE", + JsonField::Present(Value::String("".into())), + result_field, + ); + let job = remote_job_with_describe_body("job-empty-type", body); + wait_expect_none(&job).await; + } + + let mut unexpected = Vec::new(); + let reject_cases = [ + ( + "unknown_missing_result", + JsonField::Present(Value::String("future_job_type_xyz".into())), + JsonField::Absent, + ), + ( + "unknown_explicit_none", + JsonField::Present(Value::String("future_job_type_xyz".into())), + JsonField::Present(job_result_none_wire()), + ), + ( + "empty_explicit_none", + JsonField::Present(Value::String("".into())), + JsonField::Present(job_result_none_wire()), + ), + ( + "empty_explicit_function", + JsonField::Present(Value::String("".into())), + JsonField::Present(job_result_function_wire(&sample_remote_function())), + ), + ]; + for (label, job_type, result_field) in reject_cases { + let body = describe_body("job-type-expectation", "DONE", job_type, result_field); + let job = remote_job_with_describe_body("job-type-expectation", body); + match job.wait().await { + Err(Error::Http { .. }) => {} + Ok(value) => unexpected.push(format!("{label}: Ok({value:?})")), + Err(other) => unexpected.push(format!("{label}: Err({other:?})")), + } + } + assert!( + unexpected.is_empty(), + "unknown/empty job_type expectation mismatches must be Error::Http: {unexpected:?}" + ); + } + + /// Unknown or malformed result kind, version, outer field, or nested Function fails wait. + #[tokio::test] + async fn remote_job_result_malformed_or_unknown_result_is_http() { + let function = sample_remote_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 mut unknown_nested_field = function_wire.clone(); + unknown_nested_field["function"] + .as_object_mut() + .unwrap() + .insert("unexpected_field".into(), Value::Bool(true)); + + let mut malformed_nested_version = function_wire.clone(); + malformed_nested_version["function"]["format_version"] = Value::from(2); + + let malformed_results = [ + ("unknown_kind", unknown_kind), + ("unknown_version", unknown_version), + ("unknown_outer_field", unknown_outer_field), + ("empty_function_id", empty_function_id), + ("unknown_nested_field", unknown_nested_field), + ("malformed_nested_version", malformed_nested_version), + ]; + let mut unexpected = Vec::new(); + for (label, result) in malformed_results { + let body = describe_body( + "job-bad-result", + "DONE", + JsonField::Present(Value::String("register_function".into())), + JsonField::Present(result), + ); + let job = remote_job_with_describe_body("job-bad-result", body); + match job.wait().await { + Err(Error::Http { .. }) => {} + Ok(value) => unexpected.push(format!("{label}: Ok({value:?})")), + Err(other) => unexpected.push(format!("{label}: Err({other:?})")), + } + } + assert!( + unexpected.is_empty(), + "malformed/unknown result must be Error::Http for every case: {unexpected:?}" + ); + } + + /// For unknown result or job_type, status reports finished while wait fails on a separate handle. + #[tokio::test] + async fn remote_job_result_status_observes_finished_while_wait_rejects_unknown() { + let cases = [ + ( + "job-unknown-result", + JsonField::Present(Value::String("create_index".into())), + JsonField::Present(json!({ + "format_version": 1, + "kind": "future_result_kind", + "raw": {"keep": true} + })), + ), + ( + "job-unknown-type", + JsonField::Present(Value::String("future_job_type_xyz".into())), + JsonField::Present(job_result_none_wire()), + ), + ]; + + let mut unexpected = Vec::new(); + for (job_id, job_type, result) in cases { + let body = describe_body(job_id, "DONE", job_type, result); + + let status_job = remote_job_with_describe_body(job_id, body.clone()); + let status = status_job + .status() + .await + .expect("status must observe terminal DONE"); + assert_eq!(status, "finished"); + + let wait_job = remote_job_with_describe_body(job_id, body); + match wait_job.wait().await { + Err(Error::Http { .. }) => {} + Ok(value) => unexpected.push(format!("{job_id}: Ok({value:?})")), + Err(other) => unexpected.push(format!("{job_id}: Err({other:?})")), + } + } + assert!( + unexpected.is_empty(), + "wait must be Error::Http while status stayed finished: {unexpected:?}" + ); + } + + /// non-DONE FAILED or CANCELLED carrying a success result is protocol Http. + #[tokio::test] + async fn remote_job_result_non_done_carrying_success_result_is_http() { + let function_wire = job_result_function_wire(&sample_remote_function()); + let none_wire = job_result_none_wire(); + let cases = [ + ( + "job-failed-function", + "FAILED", + JsonField::Present(Value::String("register_function".into())), + JsonField::Present(function_wire), + ), + ( + "job-cancelled-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) in cases { + let body = describe_body(job_id, job_state, job_type, result); + let job = remote_job_with_describe_body(job_id, body); + match job.wait().await { + Err(Error::Http { .. }) => {} + Err(Error::JobFailed { .. }) => { + unexpected.push(format!("{job_id}: JobFailed")); + } + Err(Error::JobCancelled { .. }) => { + unexpected.push(format!("{job_id}: JobCancelled")); + } + Ok(value) => unexpected.push(format!("{job_id}: Ok({value:?})")), + Err(other) => unexpected.push(format!("{job_id}: Err({other:?})")), + } + } + assert!( + unexpected.is_empty(), + "non-DONE success result must be Error::Http, not lifecycle errors: {unexpected:?}" + ); + } + + /// Unknown informational fields on the describe outer envelope are tolerated. + #[tokio::test] + async fn remote_job_result_unknown_outer_envelope_fields_tolerated() { + let create_index_body = json!({ + "job_id": "job-create-index-extra", + "job_state": "DONE", + "job_type": "create_index", + "creation_ms": 7, + "server_note": "informational-only", + }); + let create_index_job = + remote_job_with_describe_body("job-create-index-extra", create_index_body.to_string()); + wait_expect_none(&create_index_job).await; + + let expected = sample_remote_function(); + let register_body = json!({ + "job_id": "job-register-extra", + "job_state": "DONE", + "job_type": "register_function", + "result": job_result_function_wire(&expected), + "creation_ms": 42, + "server_note": "informational-only", + "spec": {"ignored": true}, + }); + let register_job = + remote_job_with_describe_body("job-register-extra", register_body.to_string()); + let result = register_job + .wait() + .await + .expect("unknown outer fields must not block Function projection"); + match result { + JobResult::Function(function) => assert_exact_function(&function, &expected), + JobResult::None => panic!("register_function success must not project as None"), + } + } }