From 0b03851f3161d84b09dd08419bfc46c81614dba5 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Thu, 27 Aug 2026 13:17:16 -0400 Subject: [PATCH 01/44] Init --- components/clp-rust-utils/src/task_io.rs | 1 + .../clp-rust-utils/src/task_io/query.rs | 39 +++++++++++++++++++ components/clp-tdl-package/src/lib.rs | 8 +++- components/clp-tdl-package/src/task/mod.rs | 1 + .../clp-tdl-package/src/task/query/mod.rs | 20 ++++++++++ 5 files changed, 67 insertions(+), 2 deletions(-) create mode 100644 components/clp-rust-utils/src/task_io/query.rs create mode 100644 components/clp-tdl-package/src/task/query/mod.rs diff --git a/components/clp-rust-utils/src/task_io.rs b/components/clp-rust-utils/src/task_io.rs index 376ef623d5..09f7d4812e 100644 --- a/components/clp-rust-utils/src/task_io.rs +++ b/components/clp-rust-utils/src/task_io.rs @@ -1 +1,2 @@ pub mod compression; +pub mod query; diff --git a/components/clp-rust-utils/src/task_io/query.rs b/components/clp-rust-utils/src/task_io/query.rs new file mode 100644 index 0000000000..b25ea44c1b --- /dev/null +++ b/components/clp-rust-utils/src/task_io/query.rs @@ -0,0 +1,39 @@ +//! Protocol types exchanged with the Spider tasks that run CLP-S query jobs. + +use std::num::NonZeroU32; + +use serde::Deserialize; +use serde::Serialize; + +/// The job-wide CLP-S options shared by every archive-search task in a query job. +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +pub struct ClpSQueryOption { + /// The query string passed positionally to `clp-s`. + pub query_string: String, + + /// The maximum number of results retained by one archive-search invocation. + pub max_num_results: NonZeroU32, + + /// The inclusive lower timestamp bound (`--tge`), in Unix epoch microseconds. + pub begin_timestamp: Option, + + /// The inclusive upper timestamp bound (`--tle`), in Unix epoch microseconds. + pub end_timestamp: Option, + + /// Whether `clp-s` performs a case-insensitive search. + pub ignore_case: bool, +} + +/// The graph output of one successfully completed archive-search task. +/// +/// Search results are written directly to the results cache and are not returned through Spider. +/// This output only identifies the archive whose search invocation completed; it does not +/// finalize the query job. +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +pub struct QueryTaskOutput { + /// The resolved dataset containing the searched archive. + pub dataset: String, + + /// The identifier of the searched archive. + pub archive_id: String, +} diff --git a/components/clp-tdl-package/src/lib.rs b/components/clp-tdl-package/src/lib.rs index 42aa104fb8..d0a5454a03 100644 --- a/components/clp-tdl-package/src/lib.rs +++ b/components/clp-tdl-package/src/lib.rs @@ -1,4 +1,4 @@ -//! Spider TDL task package `clp`: the CLP compression tasks the Spider task executor loads. +//! Spider TDL task package `clp`: the CLP tasks the Spider task executor loads. pub mod common; mod task; @@ -28,5 +28,9 @@ fn package_init() -> Result<(), TdlError> { spider_tdl::register_tdl_package! { package_name: "clp", init: package_init, - tasks: [task::compression::s3_compress_task, task::compression::commit_task], + tasks: [ + task::compression::s3_compress_task, + task::compression::commit_task, + task::query::clp_s_search_to_results_cache_task, + ], } diff --git a/components/clp-tdl-package/src/task/mod.rs b/components/clp-tdl-package/src/task/mod.rs index f672b3ee36..25acd8a5ff 100644 --- a/components/clp-tdl-package/src/task/mod.rs +++ b/components/clp-tdl-package/src/task/mod.rs @@ -1,3 +1,4 @@ //! The task implementations this package registers with Spider. pub mod compression; +pub mod query; diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs new file mode 100644 index 0000000000..7b4eb900ee --- /dev/null +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -0,0 +1,20 @@ +//! The query-task signatures registered with Spider. + +use clp_rust_utils::task_io::query::ClpSQueryOption; +use clp_rust_utils::task_io::query::QueryTaskOutput; +use spider_tdl::TaskContext; +use spider_tdl::TdlError; +use spider_tdl::task; + +/// Searches one CLP-S archive and writes its matches directly to the results cache. +#[task(name = "search::clp_s_search_to_results_cache")] +pub(crate) fn clp_s_search_to_results_cache_task( + ctx: TaskContext, + query_job_id: i32, + clp_s_query_option: ClpSQueryOption, + dataset: String, + archive_id: String, +) -> Result { + let _ = (ctx, query_job_id, clp_s_query_option, dataset, archive_id); + todo!("Implement the CLP-S results-cache query task") +} From dbc198630d21b7be9954bf87e5b12607061c7404 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Thu, 27 Aug 2026 18:18:59 -0400 Subject: [PATCH 02/44] Change search to query; Add msgpack tests --- .../clp-rust-utils/src/task_io/query.rs | 70 +++++++++++++++++-- components/clp-tdl-package/src/lib.rs | 2 +- .../clp-tdl-package/src/task/query/mod.rs | 6 +- 3 files changed, 67 insertions(+), 11 deletions(-) diff --git a/components/clp-rust-utils/src/task_io/query.rs b/components/clp-rust-utils/src/task_io/query.rs index b25ea44c1b..7293006480 100644 --- a/components/clp-rust-utils/src/task_io/query.rs +++ b/components/clp-rust-utils/src/task_io/query.rs @@ -5,13 +5,13 @@ use std::num::NonZeroU32; use serde::Deserialize; use serde::Serialize; -/// The job-wide CLP-S options shared by every archive-search task in a query job. +/// The job-wide CLP-S options shared by every archive-query task in a query job. #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] pub struct ClpSQueryOption { /// The query string passed positionally to `clp-s`. pub query_string: String, - /// The maximum number of results retained by one archive-search invocation. + /// The maximum number of results retained by one archive-query invocation. pub max_num_results: NonZeroU32, /// The inclusive lower timestamp bound (`--tge`), in Unix epoch microseconds. @@ -24,16 +24,72 @@ pub struct ClpSQueryOption { pub ignore_case: bool, } -/// The graph output of one successfully completed archive-search task. +/// The graph output of one successfully completed archive-query task. /// -/// Search results are written directly to the results cache and are not returned through Spider. -/// This output only identifies the archive whose search invocation completed; it does not +/// Query results are written directly to the results cache and are not returned through Spider. +/// This output only identifies the archive whose query invocation completed; it does not /// finalize the query job. #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] pub struct QueryTaskOutput { - /// The resolved dataset containing the searched archive. + /// The resolved dataset containing the queried archive. pub dataset: String, - /// The identifier of the searched archive. + /// The identifier of the queried archive. pub archive_id: String, } + +#[cfg(test)] +mod tests { + use std::num::NonZeroU32; + + use super::ClpSQueryOption; + use super::QueryTaskOutput; + + #[test] + fn clp_s_query_option_with_timestamp_bounds_round_trips_through_msgpack() { + let expected = ClpSQueryOption { + query_string: "level:error".to_owned(), + max_num_results: NonZeroU32::new(1_000).expect("1,000 is nonzero"), + begin_timestamp: Some(1_700_000_000_000_001), + end_timestamp: Some(1_700_000_000_000_999), + ignore_case: true, + }; + + let serialized = rmp_serde::to_vec(&expected).expect("query options should serialize"); + let actual: ClpSQueryOption = + rmp_serde::from_slice(&serialized).expect("query options should deserialize"); + + assert_eq!(expected, actual); + } + + #[test] + fn clp_s_query_option_without_timestamp_bounds_round_trips_through_msgpack() { + let expected = ClpSQueryOption { + query_string: "*".to_owned(), + max_num_results: NonZeroU32::new(1).expect("1 is nonzero"), + begin_timestamp: None, + end_timestamp: None, + ignore_case: false, + }; + + let serialized = rmp_serde::to_vec(&expected).expect("query options should serialize"); + let actual: ClpSQueryOption = + rmp_serde::from_slice(&serialized).expect("query options should deserialize"); + + assert_eq!(expected, actual); + } + + #[test] + fn query_task_output_round_trips_through_msgpack() { + let expected = QueryTaskOutput { + dataset: "default".to_owned(), + archive_id: "archive-id".to_owned(), + }; + + let serialized = rmp_serde::to_vec(&expected).expect("task output should serialize"); + let actual: QueryTaskOutput = + rmp_serde::from_slice(&serialized).expect("task output should deserialize"); + + assert_eq!(expected, actual); + } +} diff --git a/components/clp-tdl-package/src/lib.rs b/components/clp-tdl-package/src/lib.rs index d0a5454a03..71a94cbe6c 100644 --- a/components/clp-tdl-package/src/lib.rs +++ b/components/clp-tdl-package/src/lib.rs @@ -31,6 +31,6 @@ spider_tdl::register_tdl_package! { tasks: [ task::compression::s3_compress_task, task::compression::commit_task, - task::query::clp_s_search_to_results_cache_task, + task::query::clp_s_query_to_results_cache_task, ], } diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs index 7b4eb900ee..06d1b6b283 100644 --- a/components/clp-tdl-package/src/task/query/mod.rs +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -6,9 +6,9 @@ use spider_tdl::TaskContext; use spider_tdl::TdlError; use spider_tdl::task; -/// Searches one CLP-S archive and writes its matches directly to the results cache. -#[task(name = "search::clp_s_search_to_results_cache")] -pub(crate) fn clp_s_search_to_results_cache_task( +/// Queries one CLP-S archive and writes its matches directly to the results cache. +#[task(name = "query::clp_s_query_to_results_cache")] +pub(crate) fn clp_s_query_to_results_cache_task( ctx: TaskContext, query_job_id: i32, clp_s_query_option: ClpSQueryOption, From 4b2acb6f5709903afd7061c5600ddae5bd54d764 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Thu, 27 Aug 2026 18:38:00 -0400 Subject: [PATCH 03/44] Create and submitter skeleton --- Cargo.lock | 11 ++++++ Cargo.toml | 3 +- .../clp-rust-utils/src/job_config/search.rs | 2 ++ components/query-coordinator/Cargo.toml | 11 ++++++ components/query-coordinator/src/error.rs | 8 +++++ components/query-coordinator/src/lib.rs | 6 ++++ .../src/query_job_submitter/mod.rs | 30 ++++++++++++++++ .../src/query_job_submitter/spider.rs | 36 +++++++++++++++++++ 8 files changed, 106 insertions(+), 1 deletion(-) create mode 100644 components/query-coordinator/Cargo.toml create mode 100644 components/query-coordinator/src/error.rs create mode 100644 components/query-coordinator/src/lib.rs create mode 100644 components/query-coordinator/src/query_job_submitter/mod.rs create mode 100644 components/query-coordinator/src/query_job_submitter/spider.rs diff --git a/Cargo.lock b/Cargo.lock index 3c490b08db..9c64c562eb 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3271,6 +3271,17 @@ dependencies = [ "pulldown-cmark", ] +[[package]] +name = "query-coordinator" +version = "0.13.1-dev" +dependencies = [ + "async-trait", + "clp-rust-utils", + "spider-client", + "spider-core", + "thiserror", +] + [[package]] name = "quote" version = "1.0.47" diff --git a/Cargo.toml b/Cargo.toml index d493d3e5fd..3f1f57c0c0 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -4,7 +4,8 @@ members = [ "components/clp-rust-utils", "components/clp-tdl-package", "components/compression-coordinator", - "components/log-ingestor" + "components/log-ingestor", + "components/query-coordinator" ] resolver = "3" diff --git a/components/clp-rust-utils/src/job_config/search.rs b/components/clp-rust-utils/src/job_config/search.rs index dded4e7c39..09c4503f2a 100644 --- a/components/clp-rust-utils/src/job_config/search.rs +++ b/components/clp-rust-utils/src/job_config/search.rs @@ -3,6 +3,8 @@ use num_enum::TryFromPrimitive; use serde::Deserialize; use serde::Serialize; +pub type QueryJobId = i32; + pub const QUERY_JOBS_TABLE_NAME: &str = "query_jobs"; /// Mirror of `job_orchestration.scheduler.job_config.AggregationConfig`. Must be kept in sync. diff --git a/components/query-coordinator/Cargo.toml b/components/query-coordinator/Cargo.toml new file mode 100644 index 0000000000..e28483e6e4 --- /dev/null +++ b/components/query-coordinator/Cargo.toml @@ -0,0 +1,11 @@ +[package] +name = "query-coordinator" +version = { workspace = true } +edition = { workspace = true } + +[dependencies] +async-trait = { workspace = true } +clp-rust-utils = { workspace = true } +spider-client = { workspace = true } +spider-core = { workspace = true } +thiserror = { workspace = true } diff --git a/components/query-coordinator/src/error.rs b/components/query-coordinator/src/error.rs new file mode 100644 index 0000000000..2d0be81fd4 --- /dev/null +++ b/components/query-coordinator/src/error.rs @@ -0,0 +1,8 @@ +//! The crate-level error type for the query coordinator. + +/// Errors returned by the query coordinator. +#[derive(Debug, thiserror::Error)] +pub enum Error { + #[error("spider request failure: {0}")] + SpiderClient(#[from] spider_client::error::ClientError), +} diff --git a/components/query-coordinator/src/lib.rs b/components/query-coordinator/src/lib.rs new file mode 100644 index 0000000000..27431a1b20 --- /dev/null +++ b/components/query-coordinator/src/lib.rs @@ -0,0 +1,6 @@ +//! Coordination for CLP query jobs. + +mod error; +pub mod query_job_submitter; + +pub use error::Error; diff --git a/components/query-coordinator/src/query_job_submitter/mod.rs b/components/query-coordinator/src/query_job_submitter/mod.rs new file mode 100644 index 0000000000..7ba30eb4fa --- /dev/null +++ b/components/query-coordinator/src/query_job_submitter/mod.rs @@ -0,0 +1,30 @@ +//! The query-job submission interface. + +mod spider; + +use async_trait::async_trait; +use clp_rust_utils::job_config::QueryJobId; +use clp_rust_utils::task_io::query::ClpSQueryOption; +use spider_core::task::ExecutionPolicy; +use spider_core::types::id::JobId; +use spider_core::types::id::ResourceGroupId; + +use crate::Error; + +/// Registers CLP-S query jobs with a distributed task scheduler. +#[async_trait] +pub trait QueryJobSubmitter: Clone + Send + Sync { + /// Registers, but does not start, one query task per `(dataset, archive_id)` pair. + /// + /// # Errors + /// + /// Implementations must document their error conditions. + async fn submit_query_job( + &self, + query_job_id: QueryJobId, + resource_group_id: ResourceGroupId, + clp_s_query_option: ClpSQueryOption, + archives: Vec<(String, String)>, + query_task_execution_policy: ExecutionPolicy, + ) -> Result; +} diff --git a/components/query-coordinator/src/query_job_submitter/spider.rs b/components/query-coordinator/src/query_job_submitter/spider.rs new file mode 100644 index 0000000000..51f7ddf817 --- /dev/null +++ b/components/query-coordinator/src/query_job_submitter/spider.rs @@ -0,0 +1,36 @@ +//! [`QueryJobSubmitter`] skeleton for [`spider_client::SpiderClient`]. + +use async_trait::async_trait; +use clp_rust_utils::job_config::QueryJobId; +use clp_rust_utils::task_io::query::ClpSQueryOption; +use spider_client::SpiderClient; +use spider_core::task::ExecutionPolicy; +use spider_core::types::id::JobId; +use spider_core::types::id::ResourceGroupId; + +use crate::Error; +use crate::query_job_submitter::QueryJobSubmitter; + +#[async_trait] +impl QueryJobSubmitter for SpiderClient { + /// # Errors + /// + /// Task-graph construction and submission are not implemented yet. + async fn submit_query_job( + &self, + query_job_id: QueryJobId, + resource_group_id: ResourceGroupId, + clp_s_query_option: ClpSQueryOption, + archives: Vec<(String, String)>, + query_task_execution_policy: ExecutionPolicy, + ) -> Result { + let _ = ( + query_job_id, + resource_group_id, + clp_s_query_option, + archives, + query_task_execution_policy, + ); + todo!("Construct and submit the CLP-S query task graph") + } +} From d13c6e63bb83ee30d53e9bbb88cee34e348463bb Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Thu, 27 Aug 2026 18:42:15 -0400 Subject: [PATCH 04/44] docs(clp-tdl-package): Clarify package task scope --- components/clp-tdl-package/src/lib.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/components/clp-tdl-package/src/lib.rs b/components/clp-tdl-package/src/lib.rs index 71a94cbe6c..cd742490c1 100644 --- a/components/clp-tdl-package/src/lib.rs +++ b/components/clp-tdl-package/src/lib.rs @@ -1,4 +1,4 @@ -//! Spider TDL task package `clp`: the CLP tasks the Spider task executor loads. +//! Spider TDL package `clp`, providing CLP compression and query tasks for Spider task executors. pub mod common; mod task; From 0ba6fe54f8579dde07b9dd058c9fe6b3c20c32f6 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Thu, 27 Aug 2026 18:43:30 -0400 Subject: [PATCH 05/44] docs(clp-tdl-package): Clarify package task scope --- components/clp-tdl-package/src/lib.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/components/clp-tdl-package/src/lib.rs b/components/clp-tdl-package/src/lib.rs index 71a94cbe6c..cd742490c1 100644 --- a/components/clp-tdl-package/src/lib.rs +++ b/components/clp-tdl-package/src/lib.rs @@ -1,4 +1,4 @@ -//! Spider TDL task package `clp`: the CLP tasks the Spider task executor loads. +//! Spider TDL package `clp`, providing CLP compression and query tasks for Spider task executors. pub mod common; mod task; From 6298f54afd405020a092cf16e1b46044ced44827 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Fri, 28 Aug 2026 13:20:22 -0400 Subject: [PATCH 06/44] polish --- components/clp-rust-utils/src/task_io/query.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/components/clp-rust-utils/src/task_io/query.rs b/components/clp-rust-utils/src/task_io/query.rs index 7293006480..e0c02a65a6 100644 --- a/components/clp-rust-utils/src/task_io/query.rs +++ b/components/clp-rust-utils/src/task_io/query.rs @@ -1,11 +1,11 @@ -//! Protocol types exchanged with the Spider tasks that run CLP-S query jobs. +//! Protocol types exchanged with the Spider (Huntsman) tasks that run CLP query jobs. use std::num::NonZeroU32; use serde::Deserialize; use serde::Serialize; -/// The job-wide CLP-S options shared by every archive-query task in a query job. +/// `clp-s` tuning and engine options for a query job. #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] pub struct ClpSQueryOption { /// The query string passed positionally to `clp-s`. From ff5bc49b4c1f268d428a6902efe857a502508396 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Sat, 29 Aug 2026 14:36:39 -0400 Subject: [PATCH 07/44] Update components/clp-rust-utils/src/task_io/query.rs Co-authored-by: Lin Zhihao <59785146+LinZhihao-723@users.noreply.github.com> --- components/clp-rust-utils/src/task_io/query.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/components/clp-rust-utils/src/task_io/query.rs b/components/clp-rust-utils/src/task_io/query.rs index e0c02a65a6..8844b4b1bd 100644 --- a/components/clp-rust-utils/src/task_io/query.rs +++ b/components/clp-rust-utils/src/task_io/query.rs @@ -5,7 +5,7 @@ use std::num::NonZeroU32; use serde::Deserialize; use serde::Serialize; -/// `clp-s` tuning and engine options for a query job. +/// `clp-s` options for a query job. #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] pub struct ClpSQueryOption { /// The query string passed positionally to `clp-s`. From 37396cde35b2851b0107253221301fcd68617cbf Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Sat, 29 Aug 2026 14:44:04 -0400 Subject: [PATCH 08/44] Update components/clp-rust-utils/src/task_io/query.rs Co-authored-by: Lin Zhihao <59785146+LinZhihao-723@users.noreply.github.com> --- .../clp-rust-utils/src/task_io/query.rs | 55 ------------------- 1 file changed, 55 deletions(-) diff --git a/components/clp-rust-utils/src/task_io/query.rs b/components/clp-rust-utils/src/task_io/query.rs index 8844b4b1bd..3496d15beb 100644 --- a/components/clp-rust-utils/src/task_io/query.rs +++ b/components/clp-rust-utils/src/task_io/query.rs @@ -38,58 +38,3 @@ pub struct QueryTaskOutput { pub archive_id: String, } -#[cfg(test)] -mod tests { - use std::num::NonZeroU32; - - use super::ClpSQueryOption; - use super::QueryTaskOutput; - - #[test] - fn clp_s_query_option_with_timestamp_bounds_round_trips_through_msgpack() { - let expected = ClpSQueryOption { - query_string: "level:error".to_owned(), - max_num_results: NonZeroU32::new(1_000).expect("1,000 is nonzero"), - begin_timestamp: Some(1_700_000_000_000_001), - end_timestamp: Some(1_700_000_000_000_999), - ignore_case: true, - }; - - let serialized = rmp_serde::to_vec(&expected).expect("query options should serialize"); - let actual: ClpSQueryOption = - rmp_serde::from_slice(&serialized).expect("query options should deserialize"); - - assert_eq!(expected, actual); - } - - #[test] - fn clp_s_query_option_without_timestamp_bounds_round_trips_through_msgpack() { - let expected = ClpSQueryOption { - query_string: "*".to_owned(), - max_num_results: NonZeroU32::new(1).expect("1 is nonzero"), - begin_timestamp: None, - end_timestamp: None, - ignore_case: false, - }; - - let serialized = rmp_serde::to_vec(&expected).expect("query options should serialize"); - let actual: ClpSQueryOption = - rmp_serde::from_slice(&serialized).expect("query options should deserialize"); - - assert_eq!(expected, actual); - } - - #[test] - fn query_task_output_round_trips_through_msgpack() { - let expected = QueryTaskOutput { - dataset: "default".to_owned(), - archive_id: "archive-id".to_owned(), - }; - - let serialized = rmp_serde::to_vec(&expected).expect("task output should serialize"); - let actual: QueryTaskOutput = - rmp_serde::from_slice(&serialized).expect("task output should deserialize"); - - assert_eq!(expected, actual); - } -} From da02988cf2d6060bf5090ad98527493e6f763301 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Sun, 30 Aug 2026 19:33:55 -0400 Subject: [PATCH 09/44] Change to milliseconds --- components/clp-rust-utils/src/task_io/query.rs | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/components/clp-rust-utils/src/task_io/query.rs b/components/clp-rust-utils/src/task_io/query.rs index 3496d15beb..ca2afd4a08 100644 --- a/components/clp-rust-utils/src/task_io/query.rs +++ b/components/clp-rust-utils/src/task_io/query.rs @@ -14,11 +14,11 @@ pub struct ClpSQueryOption { /// The maximum number of results retained by one archive-query invocation. pub max_num_results: NonZeroU32, - /// The inclusive lower timestamp bound (`--tge`), in Unix epoch microseconds. - pub begin_timestamp: Option, + /// Inclusive `--tge` bound in Unix epoch milliseconds. + pub begin_timestamp_millisecs: Option, - /// The inclusive upper timestamp bound (`--tle`), in Unix epoch microseconds. - pub end_timestamp: Option, + /// Inclusive `--tle` bound in Unix epoch milliseconds. + pub end_timestamp_millisecs: Option, /// Whether `clp-s` performs a case-insensitive search. pub ignore_case: bool, @@ -37,4 +37,3 @@ pub struct QueryTaskOutput { /// The identifier of the queried archive. pub archive_id: String, } - From bf2ebc043ba8c3ebd0d4a10ab10f53ffa00358fa Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Sun, 30 Aug 2026 19:35:57 -0400 Subject: [PATCH 10/44] Update components/clp-tdl-package/src/task/query/mod.rs Co-authored-by: Lin Zhihao <59785146+LinZhihao-723@users.noreply.github.com> --- components/clp-tdl-package/src/task/query/mod.rs | 13 ++++++------- 1 file changed, 6 insertions(+), 7 deletions(-) diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs index 06d1b6b283..17a2bac1b3 100644 --- a/components/clp-tdl-package/src/task/query/mod.rs +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -9,12 +9,11 @@ use spider_tdl::task; /// Queries one CLP-S archive and writes its matches directly to the results cache. #[task(name = "query::clp_s_query_to_results_cache")] pub(crate) fn clp_s_query_to_results_cache_task( - ctx: TaskContext, - query_job_id: i32, - clp_s_query_option: ClpSQueryOption, - dataset: String, - archive_id: String, -) -> Result { - let _ = (ctx, query_job_id, clp_s_query_option, dataset, archive_id); + _ctx: TaskContext, + _query_job_id: i32, + _clp_s_query_option: ClpSQueryOption, + _dataset: String, + _archive_id: String, +) -> Result<(), TdlError> { todo!("Implement the CLP-S results-cache query task") } From 2ce88821329571fd4e138757abb3ec5262957829 Mon Sep 17 00:00:00 2001 From: LinZhihao-723 Date: Mon, 31 Aug 2026 11:12:17 -0400 Subject: [PATCH 11/44] refactor(clp-tdl-package): Move `clp_binary_path` and `s3_credential_env` into a shared `task::utils` module. --- .../src/task/compression/compress.rs | 106 +--------------- components/clp-tdl-package/src/task/mod.rs | 1 + components/clp-tdl-package/src/task/utils.rs | 116 ++++++++++++++++++ 3 files changed, 119 insertions(+), 104 deletions(-) create mode 100644 components/clp-tdl-package/src/task/utils.rs diff --git a/components/clp-tdl-package/src/task/compression/compress.rs b/components/clp-tdl-package/src/task/compression/compress.rs index a1ccc35b72..0cdf035603 100644 --- a/components/clp-tdl-package/src/task/compression/compress.rs +++ b/components/clp-tdl-package/src/task/compression/compress.rs @@ -10,10 +10,7 @@ use std::process::Command; use std::process::Stdio; use anyhow::Context; -use aws_config::BehaviorVersion; -use aws_sdk_s3::config::ProvideCredentials; use clp_rust_utils::aws::AWS_DEFAULT_REGION; -use clp_rust_utils::clp_config::AwsAuthentication; use clp_rust_utils::clp_config::S3Config; use clp_rust_utils::clp_config::package::config::ArchiveOutput; use clp_rust_utils::clp_config::package::config::ArchiveOutputStorage; @@ -30,6 +27,8 @@ use non_empty_string::NonEmptyString; use crate::common::clp_home; use crate::common::runtime; +use crate::task::utils::clp_binary_path; +use crate::task::utils::s3_credential_env; /// Compresses the given S3 objects into archives, uploads them to S3, and returns their metadata /// for the commit task. @@ -345,74 +344,6 @@ fn build_s3_logs_list(input_source: &S3InputSource) -> anyhow::Result { Ok(list) } -/// Resolves the AWS credential env vars clp-s needs to access the S3 objects. -/// -/// # Returns -/// -/// The env-var name-value pairs with the following environment variables set: -/// -/// * `AWS_ACCESS_KEY_ID` -/// * `AWS_SECRET_ACCESS_KEY` -/// * `AWS_SESSION_TOKEN` (if any) -/// -/// # Errors -/// -/// Returns an error if: -/// -/// * The default AWS SDK credential provider chain has no provider. -/// * Forwards [`ProvideCredentials::provide_credentials`]'s return values on failure. -fn s3_credential_env( - runtime: &tokio::runtime::Handle, - region: &str, - auth: &AwsAuthentication, -) -> anyhow::Result> { - /// The env var holding the AWS access key ID. - const AWS_ACCESS_KEY_ID_ENV_VAR: &str = "AWS_ACCESS_KEY_ID"; - - /// The env var holding the AWS secret access key. - const AWS_SECRET_ACCESS_KEY_ENV_VAR: &str = "AWS_SECRET_ACCESS_KEY"; - - /// The env var holding the AWS session token. - const AWS_SESSION_TOKEN_ENV_VAR: &str = "AWS_SESSION_TOKEN"; - - let (access_key_id, secret_access_key, session_token) = match auth { - AwsAuthentication::Credentials { credentials } => ( - credentials.access_key_id.clone(), - credentials.secret_access_key.clone(), - credentials.session_token.clone(), - ), - AwsAuthentication::Default => { - let sdk_config = runtime.block_on( - aws_config::defaults(BehaviorVersion::latest()) - .region(aws_sdk_s3::config::Region::new(region.to_string())) - .load(), - ); - let provider = sdk_config - .credentials_provider() - .context("default AWS SDK credential provider is unavailable")?; - let credentials = runtime - .block_on(provider.provide_credentials()) - .context("failed to resolve credentials from the default AWS SDK provider chain")?; - ( - credentials.access_key_id().to_string(), - credentials.secret_access_key().to_string(), - credentials - .session_token() - .map(std::string::ToString::to_string), - ) - } - }; - - let mut env = vec![ - (AWS_ACCESS_KEY_ID_ENV_VAR, access_key_id), - (AWS_SECRET_ACCESS_KEY_ENV_VAR, secret_access_key), - ]; - if let Some(session_token) = session_token { - env.push((AWS_SESSION_TOKEN_ENV_VAR, session_token)); - } - Ok(env) -} - /// Parses a single clp-s `--print-archive-stats` stdout line into an [`ArchiveMetadata`]. /// /// NOTE: clp-s emits a superset of [`ArchiveMetadata`]'s fields per line; unknown fields are @@ -585,15 +516,6 @@ fn build_log_converter_args(output_dir: &Path, inputs_from_path: &Path) -> Vec PathBuf { - clp_home.join("bin").join(binary) -} - /// Resolves the S3 config the archives are uploaded to from `config`. /// /// # Returns @@ -905,7 +827,6 @@ mod tests { use std::path::PathBuf; use clp_rust_utils::clp_config::AwsAuthentication; - use clp_rust_utils::clp_config::AwsCredentials; use clp_rust_utils::clp_config::S3Config; use clp_rust_utils::clp_config::package::config::ArchiveOutput; use clp_rust_utils::clp_config::package::config::ArchiveOutputStorage; @@ -923,7 +844,6 @@ mod tests { use super::build_s3_logs_list; use super::create_archive_s3_key; use super::parse_archive_stats; - use super::s3_credential_env; #[test] fn build_s3_logs_list_default_endpoint() -> anyhow::Result<()> { @@ -947,28 +867,6 @@ mod tests { Ok(()) } - #[test] - fn s3_credential_env_credentials() { - let runtime = tokio::runtime::Runtime::new().expect("failed to create Tokio runtime"); - let auth = AwsAuthentication::Credentials { - credentials: AwsCredentials { - access_key_id: "the-access-key".to_string(), - secret_access_key: "the-secret-key".to_string(), - session_token: Some("the-session-token".to_string()), - }, - }; - - assert_eq!( - s3_credential_env(runtime.handle(), "us-east-1", &auth) - .expect("failed to resolve credentials"), - vec![ - ("AWS_ACCESS_KEY_ID", "the-access-key".to_string()), - ("AWS_SECRET_ACCESS_KEY", "the-secret-key".to_string()), - ("AWS_SESSION_TOKEN", "the-session-token".to_string()), - ] - ); - } - #[test] fn parse_archive_stats_ignores_extra_keys() { let line = concat!( diff --git a/components/clp-tdl-package/src/task/mod.rs b/components/clp-tdl-package/src/task/mod.rs index f672b3ee36..32fbf2c96f 100644 --- a/components/clp-tdl-package/src/task/mod.rs +++ b/components/clp-tdl-package/src/task/mod.rs @@ -1,3 +1,4 @@ //! The task implementations this package registers with Spider. pub mod compression; +pub mod utils; diff --git a/components/clp-tdl-package/src/task/utils.rs b/components/clp-tdl-package/src/task/utils.rs new file mode 100644 index 0000000000..197eed0c2e --- /dev/null +++ b/components/clp-tdl-package/src/task/utils.rs @@ -0,0 +1,116 @@ +//! Helpers shared by the tasks that invoke CLP's core binaries. + +use std::path::Path; +use std::path::PathBuf; + +use anyhow::Context; +use aws_config::BehaviorVersion; +use aws_sdk_s3::config::ProvideCredentials; +use clp_rust_utils::clp_config::AwsAuthentication; + +/// Resolves the path of a CLP binary under `clp_home`, joining `bin/{binary}`. +/// +/// # Returns +/// +/// The path to the named binary under the CLP installation. +pub(super) fn clp_binary_path(clp_home: &Path, binary: &str) -> PathBuf { + clp_home.join("bin").join(binary) +} + +/// Resolves the AWS credential env vars clp-s needs to access the S3 objects. +/// +/// # Returns +/// +/// The env-var name-value pairs with the following environment variables set: +/// +/// * `AWS_ACCESS_KEY_ID` +/// * `AWS_SECRET_ACCESS_KEY` +/// * `AWS_SESSION_TOKEN` (if any) +/// +/// # Errors +/// +/// Returns an error if: +/// +/// * The default AWS SDK credential provider chain has no provider. +/// * Forwards [`ProvideCredentials::provide_credentials`]'s return values on failure. +pub(super) fn s3_credential_env( + runtime: &tokio::runtime::Handle, + region: &str, + auth: &AwsAuthentication, +) -> anyhow::Result> { + /// The env var holding the AWS access key ID. + const AWS_ACCESS_KEY_ID_ENV_VAR: &str = "AWS_ACCESS_KEY_ID"; + + /// The env var holding the AWS secret access key. + const AWS_SECRET_ACCESS_KEY_ENV_VAR: &str = "AWS_SECRET_ACCESS_KEY"; + + /// The env var holding the AWS session token. + const AWS_SESSION_TOKEN_ENV_VAR: &str = "AWS_SESSION_TOKEN"; + + let (access_key_id, secret_access_key, session_token) = match auth { + AwsAuthentication::Credentials { credentials } => ( + credentials.access_key_id.clone(), + credentials.secret_access_key.clone(), + credentials.session_token.clone(), + ), + AwsAuthentication::Default => { + let sdk_config = runtime.block_on( + aws_config::defaults(BehaviorVersion::latest()) + .region(aws_sdk_s3::config::Region::new(region.to_string())) + .load(), + ); + let provider = sdk_config + .credentials_provider() + .context("default AWS SDK credential provider is unavailable")?; + let credentials = runtime + .block_on(provider.provide_credentials()) + .context("failed to resolve credentials from the default AWS SDK provider chain")?; + ( + credentials.access_key_id().to_string(), + credentials.secret_access_key().to_string(), + credentials + .session_token() + .map(std::string::ToString::to_string), + ) + } + }; + + let mut env = vec![ + (AWS_ACCESS_KEY_ID_ENV_VAR, access_key_id), + (AWS_SECRET_ACCESS_KEY_ENV_VAR, secret_access_key), + ]; + if let Some(session_token) = session_token { + env.push((AWS_SESSION_TOKEN_ENV_VAR, session_token)); + } + Ok(env) +} + +#[cfg(test)] +mod tests { + use clp_rust_utils::clp_config::AwsAuthentication; + use clp_rust_utils::clp_config::AwsCredentials; + + use super::s3_credential_env; + + #[test] + fn s3_credential_env_credentials() { + let runtime = tokio::runtime::Runtime::new().expect("failed to create Tokio runtime"); + let auth = AwsAuthentication::Credentials { + credentials: AwsCredentials { + access_key_id: "the-access-key".to_string(), + secret_access_key: "the-secret-key".to_string(), + session_token: Some("the-session-token".to_string()), + }, + }; + + assert_eq!( + s3_credential_env(runtime.handle(), "us-east-1", &auth) + .expect("failed to resolve credentials"), + vec![ + ("AWS_ACCESS_KEY_ID", "the-access-key".to_string()), + ("AWS_SECRET_ACCESS_KEY", "the-secret-key".to_string()), + ("AWS_SESSION_TOKEN", "the-session-token".to_string()), + ] + ); + } +} From 0296f0c5d1152678890e1d06d6944a9ce2962c4a Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Mon, 31 Aug 2026 16:52:17 -0400 Subject: [PATCH 12/44] Remove task redundant description --- components/clp-tdl-package/src/task/query/mod.rs | 1 - 1 file changed, 1 deletion(-) diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs index 17a2bac1b3..6e404b23ed 100644 --- a/components/clp-tdl-package/src/task/query/mod.rs +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -6,7 +6,6 @@ use spider_tdl::TaskContext; use spider_tdl::TdlError; use spider_tdl::task; -/// Queries one CLP-S archive and writes its matches directly to the results cache. #[task(name = "query::clp_s_query_to_results_cache")] pub(crate) fn clp_s_query_to_results_cache_task( _ctx: TaskContext, From 2096e9b42e74b25d0d03881a492c387c0da6a181 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Mon, 31 Aug 2026 18:29:49 -0400 Subject: [PATCH 13/44] Add query job ID type --- components/clp-rust-utils/src/job_config/search.rs | 2 ++ components/clp-tdl-package/src/task/query/mod.rs | 3 ++- 2 files changed, 4 insertions(+), 1 deletion(-) diff --git a/components/clp-rust-utils/src/job_config/search.rs b/components/clp-rust-utils/src/job_config/search.rs index dded4e7c39..09c4503f2a 100644 --- a/components/clp-rust-utils/src/job_config/search.rs +++ b/components/clp-rust-utils/src/job_config/search.rs @@ -3,6 +3,8 @@ use num_enum::TryFromPrimitive; use serde::Deserialize; use serde::Serialize; +pub type QueryJobId = i32; + pub const QUERY_JOBS_TABLE_NAME: &str = "query_jobs"; /// Mirror of `job_orchestration.scheduler.job_config.AggregationConfig`. Must be kept in sync. diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs index 6e404b23ed..3cab0d0df9 100644 --- a/components/clp-tdl-package/src/task/query/mod.rs +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -1,5 +1,6 @@ //! The query-task signatures registered with Spider. +use clp_rust_utils::job_config::QueryJobId; use clp_rust_utils::task_io::query::ClpSQueryOption; use clp_rust_utils::task_io::query::QueryTaskOutput; use spider_tdl::TaskContext; @@ -9,7 +10,7 @@ use spider_tdl::task; #[task(name = "query::clp_s_query_to_results_cache")] pub(crate) fn clp_s_query_to_results_cache_task( _ctx: TaskContext, - _query_job_id: i32, + _query_job_id: QueryJobId, _clp_s_query_option: ClpSQueryOption, _dataset: String, _archive_id: String, From cb87c6eaf72fabe725892a042223030e54c5ddba Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Mon, 31 Aug 2026 18:51:34 -0400 Subject: [PATCH 14/44] Make query task dataset optional --- components/clp-tdl-package/src/task/query/mod.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs index 3cab0d0df9..3b5712d8eb 100644 --- a/components/clp-tdl-package/src/task/query/mod.rs +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -12,7 +12,7 @@ pub(crate) fn clp_s_query_to_results_cache_task( _ctx: TaskContext, _query_job_id: QueryJobId, _clp_s_query_option: ClpSQueryOption, - _dataset: String, + _dataset: Option, _archive_id: String, ) -> Result<(), TdlError> { todo!("Implement the CLP-S results-cache query task") From f8a68a6ebc36f94d24ad4af8f7a8f7ffd3c7cf12 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Mon, 31 Aug 2026 19:07:58 -0400 Subject: [PATCH 15/44] Require non-empty query task strings --- components/clp-rust-utils/src/task_io/query.rs | 7 ++++--- components/clp-tdl-package/src/task/query/mod.rs | 5 +++-- 2 files changed, 7 insertions(+), 5 deletions(-) diff --git a/components/clp-rust-utils/src/task_io/query.rs b/components/clp-rust-utils/src/task_io/query.rs index ca2afd4a08..fb005aa9a6 100644 --- a/components/clp-rust-utils/src/task_io/query.rs +++ b/components/clp-rust-utils/src/task_io/query.rs @@ -2,6 +2,7 @@ use std::num::NonZeroU32; +use non_empty_string::NonEmptyString; use serde::Deserialize; use serde::Serialize; @@ -9,7 +10,7 @@ use serde::Serialize; #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] pub struct ClpSQueryOption { /// The query string passed positionally to `clp-s`. - pub query_string: String, + pub query_string: NonEmptyString, /// The maximum number of results retained by one archive-query invocation. pub max_num_results: NonZeroU32, @@ -32,8 +33,8 @@ pub struct ClpSQueryOption { #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] pub struct QueryTaskOutput { /// The resolved dataset containing the queried archive. - pub dataset: String, + pub dataset: NonEmptyString, /// The identifier of the queried archive. - pub archive_id: String, + pub archive_id: NonEmptyString, } diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs index 3b5712d8eb..919425cd5c 100644 --- a/components/clp-tdl-package/src/task/query/mod.rs +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -3,6 +3,7 @@ use clp_rust_utils::job_config::QueryJobId; use clp_rust_utils::task_io::query::ClpSQueryOption; use clp_rust_utils::task_io::query::QueryTaskOutput; +use non_empty_string::NonEmptyString; use spider_tdl::TaskContext; use spider_tdl::TdlError; use spider_tdl::task; @@ -12,8 +13,8 @@ pub(crate) fn clp_s_query_to_results_cache_task( _ctx: TaskContext, _query_job_id: QueryJobId, _clp_s_query_option: ClpSQueryOption, - _dataset: Option, - _archive_id: String, + _dataset: Option, + _archive_id: NonEmptyString, ) -> Result<(), TdlError> { todo!("Implement the CLP-S results-cache query task") } From d5c0450794a2261ff2541258f7696b5666105da8 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Mon, 31 Aug 2026 19:20:54 -0400 Subject: [PATCH 16/44] Use shared default for query result limit --- components/clp-rust-utils/src/clp_config/package/config.rs | 5 ++++- components/clp-rust-utils/src/task_io/query.rs | 5 +++-- 2 files changed, 7 insertions(+), 3 deletions(-) diff --git a/components/clp-rust-utils/src/clp_config/package/config.rs b/components/clp-rust-utils/src/clp_config/package/config.rs index 1c9fa72d71..787e25acd9 100644 --- a/components/clp-rust-utils/src/clp_config/package/config.rs +++ b/components/clp-rust-utils/src/clp_config/package/config.rs @@ -196,6 +196,9 @@ impl Database { } } +/// The default maximum number of results retained for a query. +pub const DEFAULT_MAX_NUM_QUERY_RESULTS: u32 = 1000; + #[derive(Clone, Debug, Deserialize, Eq, PartialEq)] #[serde(default)] pub struct ApiServer { @@ -211,7 +214,7 @@ impl Default for ApiServer { host: "localhost".to_owned(), port: 3001, query_job_polling: QueryJobPollingConfig::default(), - default_max_num_query_results: 1000, + default_max_num_query_results: DEFAULT_MAX_NUM_QUERY_RESULTS, } } } diff --git a/components/clp-rust-utils/src/task_io/query.rs b/components/clp-rust-utils/src/task_io/query.rs index fb005aa9a6..b425b24eb8 100644 --- a/components/clp-rust-utils/src/task_io/query.rs +++ b/components/clp-rust-utils/src/task_io/query.rs @@ -12,8 +12,9 @@ pub struct ClpSQueryOption { /// The query string passed positionally to `clp-s`. pub query_string: NonEmptyString, - /// The maximum number of results retained by one archive-query invocation. - pub max_num_results: NonZeroU32, + /// The per-archive result limit. When absent, + /// [`crate::clp_config::package::config::DEFAULT_MAX_NUM_QUERY_RESULTS`] is used. + pub max_num_results: Option, /// Inclusive `--tge` bound in Unix epoch milliseconds. pub begin_timestamp_millisecs: Option, From ca2ebe22bb57f6a6ea3f8b6ed54fe65a097052e9 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Mon, 31 Aug 2026 19:22:56 -0400 Subject: [PATCH 17/44] Lint fix remove unused --- components/clp-tdl-package/src/task/query/mod.rs | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs index 919425cd5c..6b7492374c 100644 --- a/components/clp-tdl-package/src/task/query/mod.rs +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -2,10 +2,8 @@ use clp_rust_utils::job_config::QueryJobId; use clp_rust_utils::task_io::query::ClpSQueryOption; -use clp_rust_utils::task_io::query::QueryTaskOutput; use non_empty_string::NonEmptyString; use spider_tdl::TaskContext; -use spider_tdl::TdlError; use spider_tdl::task; #[task(name = "query::clp_s_query_to_results_cache")] @@ -15,6 +13,6 @@ pub(crate) fn clp_s_query_to_results_cache_task( _clp_s_query_option: ClpSQueryOption, _dataset: Option, _archive_id: NonEmptyString, -) -> Result<(), TdlError> { +) -> Result<(), spider_tdl::TdlError> { todo!("Implement the CLP-S results-cache query task") } From 94293962e04fd02a08b5347d55c2387972ff038f Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Mon, 31 Aug 2026 19:26:41 -0400 Subject: [PATCH 18/44] Remove unused query task output --- components/clp-rust-utils/src/task_io/query.rs | 14 -------------- 1 file changed, 14 deletions(-) diff --git a/components/clp-rust-utils/src/task_io/query.rs b/components/clp-rust-utils/src/task_io/query.rs index b425b24eb8..e516bc6f9d 100644 --- a/components/clp-rust-utils/src/task_io/query.rs +++ b/components/clp-rust-utils/src/task_io/query.rs @@ -25,17 +25,3 @@ pub struct ClpSQueryOption { /// Whether `clp-s` performs a case-insensitive search. pub ignore_case: bool, } - -/// The graph output of one successfully completed archive-query task. -/// -/// Query results are written directly to the results cache and are not returned through Spider. -/// This output only identifies the archive whose query invocation completed; it does not -/// finalize the query job. -#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] -pub struct QueryTaskOutput { - /// The resolved dataset containing the queried archive. - pub dataset: NonEmptyString, - - /// The identifier of the queried archive. - pub archive_id: NonEmptyString, -} From 323acfc9acf0776f15342703a77553d0f677542a Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Mon, 31 Aug 2026 19:37:17 -0400 Subject: [PATCH 19/44] Rename task --- components/clp-tdl-package/src/lib.rs | 2 +- components/clp-tdl-package/src/task/query/mod.rs | 6 +++--- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/components/clp-tdl-package/src/lib.rs b/components/clp-tdl-package/src/lib.rs index cd742490c1..750047df6a 100644 --- a/components/clp-tdl-package/src/lib.rs +++ b/components/clp-tdl-package/src/lib.rs @@ -31,6 +31,6 @@ spider_tdl::register_tdl_package! { tasks: [ task::compression::s3_compress_task, task::compression::commit_task, - task::query::clp_s_query_to_results_cache_task, + task::query::clp_s_search_task, ], } diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs index 6b7492374c..7527aca29b 100644 --- a/components/clp-tdl-package/src/task/query/mod.rs +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -6,13 +6,13 @@ use non_empty_string::NonEmptyString; use spider_tdl::TaskContext; use spider_tdl::task; -#[task(name = "query::clp_s_query_to_results_cache")] -pub(crate) fn clp_s_query_to_results_cache_task( +#[task(name = "query::clp_s_search")] +pub(crate) fn clp_s_search_task( _ctx: TaskContext, _query_job_id: QueryJobId, _clp_s_query_option: ClpSQueryOption, _dataset: Option, _archive_id: NonEmptyString, ) -> Result<(), spider_tdl::TdlError> { - todo!("Implement the CLP-S results-cache query task") + todo!("Implement the CLP-S search task") } From 8f509ef5ee0cd9517ccc53d45bfc50bc15acfb6a Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Mon, 31 Aug 2026 19:45:37 -0400 Subject: [PATCH 20/44] Add query task output handle --- components/clp-rust-utils/src/task_io/query.rs | 7 +++++++ components/clp-tdl-package/src/task/query/mod.rs | 2 ++ 2 files changed, 9 insertions(+) diff --git a/components/clp-rust-utils/src/task_io/query.rs b/components/clp-rust-utils/src/task_io/query.rs index e516bc6f9d..c5c25a1f74 100644 --- a/components/clp-rust-utils/src/task_io/query.rs +++ b/components/clp-rust-utils/src/task_io/query.rs @@ -25,3 +25,10 @@ pub struct ClpSQueryOption { /// Whether `clp-s` performs a case-insensitive search. pub ignore_case: bool, } + +/// The output handler that `clp-s` writes a query task's results to. +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +pub enum QueryOutputHandle { + /// Writes results to the results-cache collection named after the query job's ID. + ResultsCache, +} diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs index 7527aca29b..0210565a3e 100644 --- a/components/clp-tdl-package/src/task/query/mod.rs +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -2,6 +2,7 @@ use clp_rust_utils::job_config::QueryJobId; use clp_rust_utils::task_io::query::ClpSQueryOption; +use clp_rust_utils::task_io::query::QueryOutputHandle; use non_empty_string::NonEmptyString; use spider_tdl::TaskContext; use spider_tdl::task; @@ -13,6 +14,7 @@ pub(crate) fn clp_s_search_task( _clp_s_query_option: ClpSQueryOption, _dataset: Option, _archive_id: NonEmptyString, + _output_handle: QueryOutputHandle, ) -> Result<(), spider_tdl::TdlError> { todo!("Implement the CLP-S search task") } From b67cb7008f2da6ba94fc3b84f50fe10e320809d6 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Mon, 31 Aug 2026 19:51:22 -0400 Subject: [PATCH 21/44] Fix todo uppercase --- components/clp-tdl-package/src/task/query/mod.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs index 0210565a3e..ba9e140d9b 100644 --- a/components/clp-tdl-package/src/task/query/mod.rs +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -16,5 +16,5 @@ pub(crate) fn clp_s_search_task( _archive_id: NonEmptyString, _output_handle: QueryOutputHandle, ) -> Result<(), spider_tdl::TdlError> { - todo!("Implement the CLP-S search task") + todo!("query task is not implemented") } From 3bee52b05fdd39ef2d8f9ee8d10c7cd53ac48ef4 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Mon, 31 Aug 2026 19:58:09 -0400 Subject: [PATCH 22/44] Document Python mirror for query result limit default --- components/clp-rust-utils/src/clp_config/package/config.rs | 3 +++ 1 file changed, 3 insertions(+) diff --git a/components/clp-rust-utils/src/clp_config/package/config.rs b/components/clp-rust-utils/src/clp_config/package/config.rs index 787e25acd9..722a92fee4 100644 --- a/components/clp-rust-utils/src/clp_config/package/config.rs +++ b/components/clp-rust-utils/src/clp_config/package/config.rs @@ -197,6 +197,9 @@ impl Database { } /// The default maximum number of results retained for a query. +/// +/// Mirror of `clp_py_utils.clp_config.ApiServer.default_max_num_query_results`. Must be kept in +/// sync. pub const DEFAULT_MAX_NUM_QUERY_RESULTS: u32 = 1000; #[derive(Clone, Debug, Deserialize, Eq, PartialEq)] From c39b0865305c1a488144ad8f62ba7a66148f14b6 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Mon, 31 Aug 2026 20:04:18 -0400 Subject: [PATCH 23/44] Name query task output handle parameter after its type --- components/clp-tdl-package/src/task/query/mod.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs index ba9e140d9b..d3a507be97 100644 --- a/components/clp-tdl-package/src/task/query/mod.rs +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -14,7 +14,7 @@ pub(crate) fn clp_s_search_task( _clp_s_query_option: ClpSQueryOption, _dataset: Option, _archive_id: NonEmptyString, - _output_handle: QueryOutputHandle, + _query_output_handle: QueryOutputHandle, ) -> Result<(), spider_tdl::TdlError> { todo!("query task is not implemented") } From 3cee9144490bcb4ea2ea33c49a4a5a712ba17837 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Mon, 31 Aug 2026 22:26:50 -0400 Subject: [PATCH 24/44] Fix symbol ordering --- .../clp-rust-utils/src/clp_config/package/config.rs | 12 ++++++------ components/clp-rust-utils/src/job_config/search.rs | 4 ++-- 2 files changed, 8 insertions(+), 8 deletions(-) diff --git a/components/clp-rust-utils/src/clp_config/package/config.rs b/components/clp-rust-utils/src/clp_config/package/config.rs index 722a92fee4..bc5e234471 100644 --- a/components/clp-rust-utils/src/clp_config/package/config.rs +++ b/components/clp-rust-utils/src/clp_config/package/config.rs @@ -11,6 +11,12 @@ use crate::clp_config::AwsAuthentication; use crate::clp_config::S3Config; use crate::dataset::resolve_dataset_name; +/// The default maximum number of results retained for a query. +/// +/// Mirror of `clp_py_utils.clp_config.ApiServer.default_max_num_query_results`. Must be kept in +/// sync. +pub const DEFAULT_MAX_NUM_QUERY_RESULTS: u32 = 1000; + /// Mirror of `clp_py_utils.clp_config.ClpConfig`. /// /// # NOTE @@ -196,12 +202,6 @@ impl Database { } } -/// The default maximum number of results retained for a query. -/// -/// Mirror of `clp_py_utils.clp_config.ApiServer.default_max_num_query_results`. Must be kept in -/// sync. -pub const DEFAULT_MAX_NUM_QUERY_RESULTS: u32 = 1000; - #[derive(Clone, Debug, Deserialize, Eq, PartialEq)] #[serde(default)] pub struct ApiServer { diff --git a/components/clp-rust-utils/src/job_config/search.rs b/components/clp-rust-utils/src/job_config/search.rs index 09c4503f2a..e6eda35af4 100644 --- a/components/clp-rust-utils/src/job_config/search.rs +++ b/components/clp-rust-utils/src/job_config/search.rs @@ -3,10 +3,10 @@ use num_enum::TryFromPrimitive; use serde::Deserialize; use serde::Serialize; -pub type QueryJobId = i32; - pub const QUERY_JOBS_TABLE_NAME: &str = "query_jobs"; +pub type QueryJobId = i32; + /// Mirror of `job_orchestration.scheduler.job_config.AggregationConfig`. Must be kept in sync. #[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)] #[serde(default)] From 789f671d248849bc46cf4032659f5921244e77ae Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Mon, 31 Aug 2026 22:57:26 -0400 Subject: [PATCH 25/44] Use clp-s default query result limit --- .../clp-rust-utils/src/clp_config/package/config.rs | 8 +------- components/clp-rust-utils/src/task_io/query.rs | 4 ++-- 2 files changed, 3 insertions(+), 9 deletions(-) diff --git a/components/clp-rust-utils/src/clp_config/package/config.rs b/components/clp-rust-utils/src/clp_config/package/config.rs index bc5e234471..1c9fa72d71 100644 --- a/components/clp-rust-utils/src/clp_config/package/config.rs +++ b/components/clp-rust-utils/src/clp_config/package/config.rs @@ -11,12 +11,6 @@ use crate::clp_config::AwsAuthentication; use crate::clp_config::S3Config; use crate::dataset::resolve_dataset_name; -/// The default maximum number of results retained for a query. -/// -/// Mirror of `clp_py_utils.clp_config.ApiServer.default_max_num_query_results`. Must be kept in -/// sync. -pub const DEFAULT_MAX_NUM_QUERY_RESULTS: u32 = 1000; - /// Mirror of `clp_py_utils.clp_config.ClpConfig`. /// /// # NOTE @@ -217,7 +211,7 @@ impl Default for ApiServer { host: "localhost".to_owned(), port: 3001, query_job_polling: QueryJobPollingConfig::default(), - default_max_num_query_results: DEFAULT_MAX_NUM_QUERY_RESULTS, + default_max_num_query_results: 1000, } } } diff --git a/components/clp-rust-utils/src/task_io/query.rs b/components/clp-rust-utils/src/task_io/query.rs index c5c25a1f74..412790fe6f 100644 --- a/components/clp-rust-utils/src/task_io/query.rs +++ b/components/clp-rust-utils/src/task_io/query.rs @@ -12,8 +12,8 @@ pub struct ClpSQueryOption { /// The query string passed positionally to `clp-s`. pub query_string: NonEmptyString, - /// The per-archive result limit. When absent, - /// [`crate::clp_config::package::config::DEFAULT_MAX_NUM_QUERY_RESULTS`] is used. + /// The per-archive result limit. When absent, the task omits `--max-num-results` and uses the + /// `clp-s` default. pub max_num_results: Option, /// Inclusive `--tge` bound in Unix epoch milliseconds. From 184d94029fb25c17c9f91e96228a60e48f9c635e Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Tue, 1 Sep 2026 10:40:21 -0400 Subject: [PATCH 26/44] Apply batched suggestions from code review Co-authored-by: Lin Zhihao <59785146+LinZhihao-723@users.noreply.github.com> --- components/clp-rust-utils/src/task_io/query.rs | 2 -- components/clp-tdl-package/src/task/query/mod.rs | 2 +- 2 files changed, 1 insertion(+), 3 deletions(-) diff --git a/components/clp-rust-utils/src/task_io/query.rs b/components/clp-rust-utils/src/task_io/query.rs index 412790fe6f..bd44d50819 100644 --- a/components/clp-rust-utils/src/task_io/query.rs +++ b/components/clp-rust-utils/src/task_io/query.rs @@ -29,6 +29,4 @@ pub struct ClpSQueryOption { /// The output handler that `clp-s` writes a query task's results to. #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] pub enum QueryOutputHandle { - /// Writes results to the results-cache collection named after the query job's ID. - ResultsCache, } diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs index d3a507be97..42aa3c0717 100644 --- a/components/clp-tdl-package/src/task/query/mod.rs +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -16,5 +16,5 @@ pub(crate) fn clp_s_search_task( _archive_id: NonEmptyString, _query_output_handle: QueryOutputHandle, ) -> Result<(), spider_tdl::TdlError> { - todo!("query task is not implemented") + todo!("clp-s search task is not implemented") } From c479fade00c78e4dff9f4199db60bfba13cda12c Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Tue, 1 Sep 2026 10:45:22 -0400 Subject: [PATCH 27/44] Propagate rename --- components/clp-rust-utils/src/task_io/query.rs | 2 +- components/clp-tdl-package/src/task/query/mod.rs | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/components/clp-rust-utils/src/task_io/query.rs b/components/clp-rust-utils/src/task_io/query.rs index bd44d50819..cf38fcdf7c 100644 --- a/components/clp-rust-utils/src/task_io/query.rs +++ b/components/clp-rust-utils/src/task_io/query.rs @@ -28,5 +28,5 @@ pub struct ClpSQueryOption { /// The output handler that `clp-s` writes a query task's results to. #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] -pub enum QueryOutputHandle { +pub enum OutputHandle { } diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs index 42aa3c0717..9326b20631 100644 --- a/components/clp-tdl-package/src/task/query/mod.rs +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -2,7 +2,7 @@ use clp_rust_utils::job_config::QueryJobId; use clp_rust_utils::task_io::query::ClpSQueryOption; -use clp_rust_utils::task_io::query::QueryOutputHandle; +use clp_rust_utils::task_io::query::OutputHandle; use non_empty_string::NonEmptyString; use spider_tdl::TaskContext; use spider_tdl::task; @@ -14,7 +14,7 @@ pub(crate) fn clp_s_search_task( _clp_s_query_option: ClpSQueryOption, _dataset: Option, _archive_id: NonEmptyString, - _query_output_handle: QueryOutputHandle, + _output_handle: OutputHandle, ) -> Result<(), spider_tdl::TdlError> { todo!("clp-s search task is not implemented") } From cfa0a97693fa4bc8013c62a5f54cd13a16e31588 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Tue, 1 Sep 2026 12:14:43 -0400 Subject: [PATCH 28/44] lint fix --- components/clp-rust-utils/src/task_io/query.rs | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/components/clp-rust-utils/src/task_io/query.rs b/components/clp-rust-utils/src/task_io/query.rs index cf38fcdf7c..b5f53a9c4b 100644 --- a/components/clp-rust-utils/src/task_io/query.rs +++ b/components/clp-rust-utils/src/task_io/query.rs @@ -28,5 +28,4 @@ pub struct ClpSQueryOption { /// The output handler that `clp-s` writes a query task's results to. #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] -pub enum OutputHandle { -} +pub enum OutputHandle {} From 89ca7b57522f5657e8c5e84b7c5ccc47bc7cf342 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Wed, 2 Sep 2026 04:20:28 -0400 Subject: [PATCH 29/44] Fix inconsistencies --- Cargo.lock | 1 + components/clp-rust-utils/src/job_config/search.rs | 2 -- components/query-coordinator/Cargo.toml | 1 + components/query-coordinator/src/query_job_submitter/mod.rs | 3 ++- components/query-coordinator/src/query_job_submitter/spider.rs | 3 ++- 5 files changed, 6 insertions(+), 4 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 9c64c562eb..e64950616f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3277,6 +3277,7 @@ version = "0.13.1-dev" dependencies = [ "async-trait", "clp-rust-utils", + "non-empty-string", "spider-client", "spider-core", "thiserror", diff --git a/components/clp-rust-utils/src/job_config/search.rs b/components/clp-rust-utils/src/job_config/search.rs index 1659a00ab9..e6eda35af4 100644 --- a/components/clp-rust-utils/src/job_config/search.rs +++ b/components/clp-rust-utils/src/job_config/search.rs @@ -3,8 +3,6 @@ use num_enum::TryFromPrimitive; use serde::Deserialize; use serde::Serialize; -pub type QueryJobId = i32; - pub const QUERY_JOBS_TABLE_NAME: &str = "query_jobs"; pub type QueryJobId = i32; diff --git a/components/query-coordinator/Cargo.toml b/components/query-coordinator/Cargo.toml index e28483e6e4..23c89d0442 100644 --- a/components/query-coordinator/Cargo.toml +++ b/components/query-coordinator/Cargo.toml @@ -6,6 +6,7 @@ edition = { workspace = true } [dependencies] async-trait = { workspace = true } clp-rust-utils = { workspace = true } +non-empty-string = { workspace = true } spider-client = { workspace = true } spider-core = { workspace = true } thiserror = { workspace = true } diff --git a/components/query-coordinator/src/query_job_submitter/mod.rs b/components/query-coordinator/src/query_job_submitter/mod.rs index 7ba30eb4fa..6e71f6e19c 100644 --- a/components/query-coordinator/src/query_job_submitter/mod.rs +++ b/components/query-coordinator/src/query_job_submitter/mod.rs @@ -5,6 +5,7 @@ mod spider; use async_trait::async_trait; use clp_rust_utils::job_config::QueryJobId; use clp_rust_utils::task_io::query::ClpSQueryOption; +use non_empty_string::NonEmptyString; use spider_core::task::ExecutionPolicy; use spider_core::types::id::JobId; use spider_core::types::id::ResourceGroupId; @@ -24,7 +25,7 @@ pub trait QueryJobSubmitter: Clone + Send + Sync { query_job_id: QueryJobId, resource_group_id: ResourceGroupId, clp_s_query_option: ClpSQueryOption, - archives: Vec<(String, String)>, + archives: Vec<(Option, NonEmptyString)>, query_task_execution_policy: ExecutionPolicy, ) -> Result; } diff --git a/components/query-coordinator/src/query_job_submitter/spider.rs b/components/query-coordinator/src/query_job_submitter/spider.rs index 51f7ddf817..325fe74d95 100644 --- a/components/query-coordinator/src/query_job_submitter/spider.rs +++ b/components/query-coordinator/src/query_job_submitter/spider.rs @@ -3,6 +3,7 @@ use async_trait::async_trait; use clp_rust_utils::job_config::QueryJobId; use clp_rust_utils::task_io::query::ClpSQueryOption; +use non_empty_string::NonEmptyString; use spider_client::SpiderClient; use spider_core::task::ExecutionPolicy; use spider_core::types::id::JobId; @@ -21,7 +22,7 @@ impl QueryJobSubmitter for SpiderClient { query_job_id: QueryJobId, resource_group_id: ResourceGroupId, clp_s_query_option: ClpSQueryOption, - archives: Vec<(String, String)>, + archives: Vec<(Option, NonEmptyString)>, query_task_execution_policy: ExecutionPolicy, ) -> Result { let _ = ( From 4e063fb45c860fe7cfb0f462d2d122ce4f34b066 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Wed, 2 Sep 2026 11:34:34 -0400 Subject: [PATCH 30/44] Add missing output handle arg --- components/query-coordinator/src/query_job_submitter/mod.rs | 2 ++ components/query-coordinator/src/query_job_submitter/spider.rs | 3 +++ 2 files changed, 5 insertions(+) diff --git a/components/query-coordinator/src/query_job_submitter/mod.rs b/components/query-coordinator/src/query_job_submitter/mod.rs index 6e71f6e19c..835e1f0b19 100644 --- a/components/query-coordinator/src/query_job_submitter/mod.rs +++ b/components/query-coordinator/src/query_job_submitter/mod.rs @@ -5,6 +5,7 @@ mod spider; use async_trait::async_trait; use clp_rust_utils::job_config::QueryJobId; use clp_rust_utils::task_io::query::ClpSQueryOption; +use clp_rust_utils::task_io::query::OutputHandle; use non_empty_string::NonEmptyString; use spider_core::task::ExecutionPolicy; use spider_core::types::id::JobId; @@ -25,6 +26,7 @@ pub trait QueryJobSubmitter: Clone + Send + Sync { query_job_id: QueryJobId, resource_group_id: ResourceGroupId, clp_s_query_option: ClpSQueryOption, + output_handle: OutputHandle, archives: Vec<(Option, NonEmptyString)>, query_task_execution_policy: ExecutionPolicy, ) -> Result; diff --git a/components/query-coordinator/src/query_job_submitter/spider.rs b/components/query-coordinator/src/query_job_submitter/spider.rs index 325fe74d95..d78a19d880 100644 --- a/components/query-coordinator/src/query_job_submitter/spider.rs +++ b/components/query-coordinator/src/query_job_submitter/spider.rs @@ -3,6 +3,7 @@ use async_trait::async_trait; use clp_rust_utils::job_config::QueryJobId; use clp_rust_utils::task_io::query::ClpSQueryOption; +use clp_rust_utils::task_io::query::OutputHandle; use non_empty_string::NonEmptyString; use spider_client::SpiderClient; use spider_core::task::ExecutionPolicy; @@ -22,6 +23,7 @@ impl QueryJobSubmitter for SpiderClient { query_job_id: QueryJobId, resource_group_id: ResourceGroupId, clp_s_query_option: ClpSQueryOption, + output_handle: OutputHandle, archives: Vec<(Option, NonEmptyString)>, query_task_execution_policy: ExecutionPolicy, ) -> Result { @@ -29,6 +31,7 @@ impl QueryJobSubmitter for SpiderClient { query_job_id, resource_group_id, clp_s_query_option, + output_handle, archives, query_task_execution_policy, ); From f5860dedd096b7511ca628e991a269e495bab639 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Wed, 2 Sep 2026 12:00:15 -0400 Subject: [PATCH 31/44] Add query job lifecycle handling --- Cargo.lock | 3 + .../initialize-orchestration-db.py | 38 ++- components/query-coordinator/Cargo.toml | 3 + components/query-coordinator/src/error.rs | 9 + .../query-coordinator/src/job_handle.rs | 258 ++++++++++++++++++ components/query-coordinator/src/lib.rs | 1 + .../src/query_job_submitter/mod.rs | 27 ++ .../src/query_job_submitter/spider.rs | 49 ++++ 8 files changed, 386 insertions(+), 2 deletions(-) create mode 100644 components/query-coordinator/src/job_handle.rs diff --git a/Cargo.lock b/Cargo.lock index e64950616f..85032a8a0c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3280,7 +3280,10 @@ dependencies = [ "non-empty-string", "spider-client", "spider-core", + "sqlx", "thiserror", + "tokio", + "tracing", ] [[package]] diff --git a/components/clp-py-utils/clp_py_utils/initialize-orchestration-db.py b/components/clp-py-utils/clp_py_utils/initialize-orchestration-db.py index 5af48908a6..1b33f52a84 100644 --- a/components/clp-py-utils/clp_py_utils/initialize-orchestration-db.py +++ b/components/clp-py-utils/clp_py_utils/initialize-orchestration-db.py @@ -11,7 +11,7 @@ QueryJobStatus, QueryTaskStatus, ) -from mysql.connector.errorcode import ER_DUP_KEYNAME +from mysql.connector.errorcode import ER_DUP_FIELDNAME, ER_DUP_KEYNAME from pydantic import ValidationError from clp_py_utils.clp_config import ( @@ -134,19 +134,53 @@ def main(argv): `id` INT NOT NULL AUTO_INCREMENT, `type` INT NOT NULL, `status` INT NOT NULL DEFAULT '{QueryJobStatus.PENDING}', + `status_msg` VARCHAR(512) NOT NULL DEFAULT '', `creation_time` DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3), `num_tasks` INT NOT NULL DEFAULT '0', `num_tasks_completed` INT NOT NULL DEFAULT '0', `start_time` DATETIME(3) NULL DEFAULT NULL, `duration` FLOAT NULL DEFAULT NULL, `job_config` MEDIUMBLOB NOT NULL, + `spider_id` BIGINT UNSIGNED NULL DEFAULT NULL, PRIMARY KEY (`id`) USING BTREE, INDEX `CREATION_TIME` (`creation_time`) USING BTREE, - INDEX `JOB_STATUS` (`status`) USING BTREE + INDEX `JOB_STATUS` (`status`) USING BTREE, + INDEX `JOB_SPIDER_ID` (`spider_id`) USING BTREE ) ROW_FORMAT=DYNAMIC """ ) + # Upgrade query-job tables created before Spider lifecycle support was added. + query_jobs_table_upgrades = ( + ( + f""" + ALTER TABLE `{QUERY_JOBS_TABLE_NAME}` + ADD COLUMN `status_msg` VARCHAR(512) NOT NULL DEFAULT '' AFTER `status` + """, + ER_DUP_FIELDNAME, + ), + ( + f""" + ALTER TABLE `{QUERY_JOBS_TABLE_NAME}` + ADD COLUMN `spider_id` BIGINT UNSIGNED NULL DEFAULT NULL AFTER `job_config` + """, + ER_DUP_FIELDNAME, + ), + ( + f""" + ALTER TABLE `{QUERY_JOBS_TABLE_NAME}` + ADD INDEX `JOB_SPIDER_ID` (`spider_id`) USING BTREE + """, + ER_DUP_KEYNAME, + ), + ) + for upgrade_query, duplicate_error_code in query_jobs_table_upgrades: + try: + scheduling_db_cursor.execute(upgrade_query) + except Exception as err: + if not (hasattr(err, "errno") and err.errno == duplicate_error_code): + raise + scheduling_db_cursor.execute( f""" CREATE TABLE IF NOT EXISTS `{QUERY_TASKS_TABLE_NAME}` ( diff --git a/components/query-coordinator/Cargo.toml b/components/query-coordinator/Cargo.toml index 23c89d0442..c5bb002a41 100644 --- a/components/query-coordinator/Cargo.toml +++ b/components/query-coordinator/Cargo.toml @@ -9,4 +9,7 @@ clp-rust-utils = { workspace = true } non-empty-string = { workspace = true } spider-client = { workspace = true } spider-core = { workspace = true } +sqlx = { workspace = true } thiserror = { workspace = true } +tokio = { workspace = true } +tracing = { workspace = true } diff --git a/components/query-coordinator/src/error.rs b/components/query-coordinator/src/error.rs index 2d0be81fd4..bd7b31135c 100644 --- a/components/query-coordinator/src/error.rs +++ b/components/query-coordinator/src/error.rs @@ -3,6 +3,15 @@ /// Errors returned by the query coordinator. #[derive(Debug, thiserror::Error)] pub enum Error { + #[error("query job {0} is no longer pending")] + JobNotPending(clp_rust_utils::job_config::QueryJobId), + #[error("spider request failure: {0}")] SpiderClient(#[from] spider_client::error::ClientError), + + #[error("sqlx error: {0}")] + Sqlx(#[from] sqlx::Error), + + #[error("number of query tasks {0} exceeds `i32::MAX`")] + TooManyQueryTasks(usize), } diff --git a/components/query-coordinator/src/job_handle.rs b/components/query-coordinator/src/job_handle.rs new file mode 100644 index 0000000000..8e21208e25 --- /dev/null +++ b/components/query-coordinator/src/job_handle.rs @@ -0,0 +1,258 @@ +//! Lifecycle management for one coordinator-planned query job. + +use std::sync::Arc; +use std::time::Duration; + +use clp_rust_utils::job_config::QUERY_JOBS_TABLE_NAME; +use clp_rust_utils::job_config::QueryJobId; +use clp_rust_utils::job_config::QueryJobStatus; +use clp_rust_utils::task_io::query::ClpSQueryOption; +use clp_rust_utils::task_io::query::OutputHandle; +use non_empty_string::NonEmptyString; +use spider_core::task::ExecutionPolicy; +use spider_core::types::id::JobId as SpiderJobId; +use spider_core::types::id::ResourceGroupId; +use sqlx::MySqlPool; + +use crate::Error; +use crate::query_job_submitter::QueryJobOutcome; +use crate::query_job_submitter::QueryJobSubmitter; + +/// The coordinator-prepared inputs for one Spider query graph. +pub struct QueryPlan { + /// Query behavior shared by every archive task. + pub clp_s_query_option: ClpSQueryOption, + + /// Result destination shared by every archive task. + pub output_handle: OutputHandle, + + /// One optional dataset and non-empty archive ID pair per task. + pub archives: Vec<(Option, NonEmptyString)>, + + /// Execution policy applied to every archive task. + pub query_task_execution_policy: ExecutionPolicy, +} + +/// Spider polling options shared by query-job handles. +pub struct SpiderOption { + /// Delay before the first Spider job-state poll. + pub initial_poll_backoff: Duration, + + /// Maximum delay between Spider job-state polls. + pub max_poll_backoff: Duration, +} + +/// Drives one already-planned query job through submission and terminal persistence. +pub struct QueryJobHandle { + db_pool: MySqlPool, + query_job_id: QueryJobId, + job_submitter: SubmitterType, + resource_group_id: ResourceGroupId, + query_plan: QueryPlan, + spider_option: Arc, +} + +impl QueryJobHandle { + /// Constructs a handle for an already-planned query job. + pub fn new( + db_pool: MySqlPool, + query_job_id: QueryJobId, + job_submitter: SubmitterType, + resource_group_id: ResourceGroupId, + query_plan: QueryPlan, + spider_option: Arc, + ) -> Self { + Self { + db_pool, + query_job_id, + job_submitter, + resource_group_id, + query_plan, + spider_option, + } + } + + /// Submits the prepared graph and drives the query job to a terminal state. + /// + /// On an orchestration failure, this method makes a best-effort attempt to mark the CLP query + /// job as failed before returning the original error. + /// + /// # Errors + /// + /// Returns an error if submission, submission persistence, polling, or terminal persistence + /// fails. + pub async fn run(self) -> Result<(), Error> { + tracing::info!(query_job_id = %self.query_job_id, "Starting query job."); + + let result = self.submit_and_wait().await; + if let Err(error) = &result { + self.report_failure(error).await; + } + result + } + + /// Resumes a query job that was already submitted to Spider. + /// + /// The caller must ensure `spider_job_id` belongs to this CLP query job. + /// + /// # Errors + /// + /// Returns an error if polling or terminal persistence fails. + pub async fn recover(self, spider_job_id: SpiderJobId) -> Result<(), Error> { + tracing::info!( + query_job_id = %self.query_job_id, + spider_job_id = %spider_job_id, + "Recovering query job.", + ); + + let result = self.to_completion(spider_job_id).await; + if let Err(error) = &result { + self.report_failure(error).await; + } + result + } + + async fn submit_and_wait(&self) -> Result<(), Error> { + let num_tasks = self.query_plan.archives.len(); + let persisted_num_tasks = + i32::try_from(num_tasks).map_err(|_| Error::TooManyQueryTasks(num_tasks))?; + let spider_job_id = self + .job_submitter + .submit_query_job( + self.query_job_id, + self.resource_group_id, + self.query_plan.clp_s_query_option.clone(), + self.query_plan.output_handle.clone(), + self.query_plan.archives.clone(), + self.query_plan.query_task_execution_policy.clone(), + ) + .await?; + + tracing::info!( + query_job_id = %self.query_job_id, + spider_job_id = %spider_job_id, + num_tasks, + "Query job submitted.", + ); + + self.persist_submission(spider_job_id, persisted_num_tasks) + .await?; + self.to_completion(spider_job_id).await + } + + async fn persist_submission( + &self, + spider_job_id: SpiderJobId, + num_tasks: i32, + ) -> Result<(), Error> { + let query = format!( + "UPDATE `{QUERY_JOBS_TABLE_NAME}` SET `spider_id` = ?, `status` = ?, `num_tasks` = ?, \ + `start_time` = CURRENT_TIMESTAMP(3) WHERE `id` = ? AND `status` = ?" + ); + let result = sqlx::query(&query) + .bind(spider_job_id.get()) + .bind(i32::from(QueryJobStatus::Running)) + .bind(num_tasks) + .bind(self.query_job_id) + .bind(i32::from(QueryJobStatus::Pending)) + .execute(&self.db_pool) + .await?; + + if 1 != result.rows_affected() { + return Err(Error::JobNotPending(self.query_job_id)); + } + Ok(()) + } + + async fn to_completion(&self, spider_job_id: SpiderJobId) -> Result<(), Error> { + let outcome = self + .job_submitter + .run_query_job_to_completion( + spider_job_id, + self.spider_option.initial_poll_backoff, + self.spider_option.max_poll_backoff, + ) + .await?; + + tracing::info!( + query_job_id = %self.query_job_id, + spider_job_id = %spider_job_id, + outcome = ?outcome, + "Query job reached a terminal Spider state.", + ); + + match outcome { + QueryJobOutcome::Succeeded => { + self.update_terminal_status(QueryJobStatus::Succeeded, "", false) + .await + } + QueryJobOutcome::Failed { error_message } => { + self.update_terminal_status( + QueryJobStatus::Failed, + &format!("The Spider query job failed: {error_message}"), + false, + ) + .await + } + QueryJobOutcome::UnexpectedlyCancelled => { + self.update_terminal_status( + QueryJobStatus::Failed, + "The Spider query job was unexpectedly cancelled.", + false, + ) + .await + } + } + } + + async fn report_failure(&self, error: &Error) { + tracing::error!( + query_job_id = %self.query_job_id, + error = %error, + "Query-job orchestration failed.", + ); + + if let Err(status_error) = self + .update_terminal_status( + QueryJobStatus::Failed, + &format!("Query-job orchestration failed: {error}"), + true, + ) + .await + { + tracing::error!( + query_job_id = %self.query_job_id, + error = %status_error, + "Failed to persist the query-job failure.", + ); + } + } + + /// Updates a non-terminal query job while preserving every existing terminal or cancellation + /// state. When `allow_pending` is false, only a running job may transition. + async fn update_terminal_status( + &self, + status: QueryJobStatus, + status_message: &str, + allow_pending: bool, + ) -> Result<(), Error> { + let eligible_statuses = if allow_pending { "?, ?" } else { "?" }; + let query = format!( + "UPDATE `{QUERY_JOBS_TABLE_NAME}` SET `status` = ?, `status_msg` = LEFT(?, 512), \ + `duration` = \ + CASE WHEN `start_time` IS NULL THEN 0 ELSE TIMESTAMPDIFF(MICROSECOND, `start_time`, \ + CURRENT_TIMESTAMP(3)) / 1000000.0 END WHERE `id` = ? AND `status` IN \ + ({eligible_statuses})" + ); + let mut query = sqlx::query(&query) + .bind(i32::from(status)) + .bind(status_message) + .bind(self.query_job_id) + .bind(i32::from(QueryJobStatus::Running)); + if allow_pending { + query = query.bind(i32::from(QueryJobStatus::Pending)); + } + query.execute(&self.db_pool).await?; + Ok(()) + } +} diff --git a/components/query-coordinator/src/lib.rs b/components/query-coordinator/src/lib.rs index 27431a1b20..32869c80d1 100644 --- a/components/query-coordinator/src/lib.rs +++ b/components/query-coordinator/src/lib.rs @@ -1,6 +1,7 @@ //! Coordination for CLP query jobs. mod error; +pub mod job_handle; pub mod query_job_submitter; pub use error::Error; diff --git a/components/query-coordinator/src/query_job_submitter/mod.rs b/components/query-coordinator/src/query_job_submitter/mod.rs index 835e1f0b19..d26a7cf525 100644 --- a/components/query-coordinator/src/query_job_submitter/mod.rs +++ b/components/query-coordinator/src/query_job_submitter/mod.rs @@ -2,6 +2,8 @@ mod spider; +use std::time::Duration; + use async_trait::async_trait; use clp_rust_utils::job_config::QueryJobId; use clp_rust_utils::task_io::query::ClpSQueryOption; @@ -13,6 +15,19 @@ use spider_core::types::id::ResourceGroupId; use crate::Error; +/// The terminal outcome of a query job. +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum QueryJobOutcome { + /// Every archive query completed successfully. + Succeeded, + + /// At least one archive query failed. + Failed { error_message: String }, + + /// Spider cancelled the job unexpectedly. User-requested cancellation is outside the MVP. + UnexpectedlyCancelled, +} + /// Registers CLP-S query jobs with a distributed task scheduler. #[async_trait] pub trait QueryJobSubmitter: Clone + Send + Sync { @@ -30,4 +45,16 @@ pub trait QueryJobSubmitter: Clone + Send + Sync { archives: Vec<(Option, NonEmptyString)>, query_task_execution_policy: ExecutionPolicy, ) -> Result; + + /// Idempotently starts `spider_job_id` and waits for it to reach a terminal state. + /// + /// # Errors + /// + /// Implementations must document their error conditions. + async fn run_query_job_to_completion( + &self, + spider_job_id: JobId, + initial_poll_backoff: Duration, + max_poll_backoff: Duration, + ) -> Result; } diff --git a/components/query-coordinator/src/query_job_submitter/spider.rs b/components/query-coordinator/src/query_job_submitter/spider.rs index d78a19d880..46b6537581 100644 --- a/components/query-coordinator/src/query_job_submitter/spider.rs +++ b/components/query-coordinator/src/query_job_submitter/spider.rs @@ -1,16 +1,21 @@ //! [`QueryJobSubmitter`] skeleton for [`spider_client::SpiderClient`]. +use std::time::Duration; + use async_trait::async_trait; use clp_rust_utils::job_config::QueryJobId; use clp_rust_utils::task_io::query::ClpSQueryOption; use clp_rust_utils::task_io::query::OutputHandle; use non_empty_string::NonEmptyString; use spider_client::SpiderClient; +use spider_client::error::ClientError; +use spider_core::job::JobState; use spider_core::task::ExecutionPolicy; use spider_core::types::id::JobId; use spider_core::types::id::ResourceGroupId; use crate::Error; +use crate::query_job_submitter::QueryJobOutcome; use crate::query_job_submitter::QueryJobSubmitter; #[async_trait] @@ -37,4 +42,48 @@ impl QueryJobSubmitter for SpiderClient { ); todo!("Construct and submit the CLP-S query task graph") } + + /// # Errors + /// + /// Returns an error if: + /// + /// * Starting the job fails for a reason other than it already having been started. + /// * Fetching the Spider job state fails. + async fn run_query_job_to_completion( + &self, + spider_job_id: JobId, + initial_poll_backoff: Duration, + max_poll_backoff: Duration, + ) -> Result { + const POLL_BACKOFF_FACTOR: u32 = 2; + + match self.start_job(spider_job_id).await { + Ok(_) | Err(ClientError::InvalidJobState(_)) => {} + Err(error) => return Err(error.into()), + } + + let mut backoff = initial_poll_backoff.min(max_poll_backoff); + let terminal_state = loop { + let state = self.get_job_state(spider_job_id).await?; + if state.is_terminal() { + break state; + } + tokio::time::sleep(backoff).await; + backoff = backoff + .saturating_mul(POLL_BACKOFF_FACTOR) + .min(max_poll_backoff); + }; + + Ok(match terminal_state { + JobState::Succeeded => QueryJobOutcome::Succeeded, + JobState::Failed => QueryJobOutcome::Failed { + error_message: self + .get_job_error(spider_job_id) + .await + .unwrap_or_else(|error| format!("")), + }, + JobState::Cancelled => QueryJobOutcome::UnexpectedlyCancelled, + _ => unreachable!("a terminal Spider state must have a terminal outcome"), + }) + } } From 842c93c584234e0b999355ca3938e64b04567f4e Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Wed, 2 Sep 2026 12:09:47 -0400 Subject: [PATCH 32/44] Address review comments --- .../clp-rust-utils/src/job_config/search.rs | 3 +++ .../clp-tdl-package/src/task/query/mod.rs | 3 ++- .../src/query_job_submitter/mod.rs | 19 ++++++++++++++--- .../src/query_job_submitter/spider.rs | 21 ++++++------------- 4 files changed, 27 insertions(+), 19 deletions(-) diff --git a/components/clp-rust-utils/src/job_config/search.rs b/components/clp-rust-utils/src/job_config/search.rs index e6eda35af4..ea2d6651d9 100644 --- a/components/clp-rust-utils/src/job_config/search.rs +++ b/components/clp-rust-utils/src/job_config/search.rs @@ -1,3 +1,4 @@ +use non_empty_string::NonEmptyString; use num_enum::IntoPrimitive; use num_enum::TryFromPrimitive; use serde::Deserialize; @@ -5,6 +6,8 @@ use serde::Serialize; pub const QUERY_JOBS_TABLE_NAME: &str = "query_jobs"; +pub type ArchiveId = NonEmptyString; + pub type QueryJobId = i32; /// Mirror of `job_orchestration.scheduler.job_config.AggregationConfig`. Must be kept in sync. diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs index 9326b20631..76697069a5 100644 --- a/components/clp-tdl-package/src/task/query/mod.rs +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -1,5 +1,6 @@ //! The query-task signatures registered with Spider. +use clp_rust_utils::job_config::ArchiveId; use clp_rust_utils::job_config::QueryJobId; use clp_rust_utils::task_io::query::ClpSQueryOption; use clp_rust_utils::task_io::query::OutputHandle; @@ -13,7 +14,7 @@ pub(crate) fn clp_s_search_task( _query_job_id: QueryJobId, _clp_s_query_option: ClpSQueryOption, _dataset: Option, - _archive_id: NonEmptyString, + _archive_id: ArchiveId, _output_handle: OutputHandle, ) -> Result<(), spider_tdl::TdlError> { todo!("clp-s search task is not implemented") diff --git a/components/query-coordinator/src/query_job_submitter/mod.rs b/components/query-coordinator/src/query_job_submitter/mod.rs index 835e1f0b19..c7ba7d5499 100644 --- a/components/query-coordinator/src/query_job_submitter/mod.rs +++ b/components/query-coordinator/src/query_job_submitter/mod.rs @@ -3,6 +3,7 @@ mod spider; use async_trait::async_trait; +use clp_rust_utils::job_config::ArchiveId; use clp_rust_utils::job_config::QueryJobId; use clp_rust_utils::task_io::query::ClpSQueryOption; use clp_rust_utils::task_io::query::OutputHandle; @@ -13,10 +14,23 @@ use spider_core::types::id::ResourceGroupId; use crate::Error; +/// Coordinator-side metadata for an archive query task. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct ArchiveMetadata { + /// The archive's ID. + pub id: ArchiveId, + + /// The archive's dataset, or `None` for the default dataset. + pub dataset: Option, + + /// The archive's compressed size in bytes. + pub size: u64, +} + /// Registers CLP-S query jobs with a distributed task scheduler. #[async_trait] pub trait QueryJobSubmitter: Clone + Send + Sync { - /// Registers, but does not start, one query task per `(dataset, archive_id)` pair. + /// Registers, but does not start, one query task per archive. /// /// # Errors /// @@ -27,7 +41,6 @@ pub trait QueryJobSubmitter: Clone + Send + Sync { resource_group_id: ResourceGroupId, clp_s_query_option: ClpSQueryOption, output_handle: OutputHandle, - archives: Vec<(Option, NonEmptyString)>, - query_task_execution_policy: ExecutionPolicy, + archives_to_search: Vec<(ArchiveMetadata, ExecutionPolicy)>, ) -> Result; } diff --git a/components/query-coordinator/src/query_job_submitter/spider.rs b/components/query-coordinator/src/query_job_submitter/spider.rs index d78a19d880..84f0b78a4b 100644 --- a/components/query-coordinator/src/query_job_submitter/spider.rs +++ b/components/query-coordinator/src/query_job_submitter/spider.rs @@ -4,13 +4,13 @@ use async_trait::async_trait; use clp_rust_utils::job_config::QueryJobId; use clp_rust_utils::task_io::query::ClpSQueryOption; use clp_rust_utils::task_io::query::OutputHandle; -use non_empty_string::NonEmptyString; use spider_client::SpiderClient; use spider_core::task::ExecutionPolicy; use spider_core::types::id::JobId; use spider_core::types::id::ResourceGroupId; use crate::Error; +use crate::query_job_submitter::ArchiveMetadata; use crate::query_job_submitter::QueryJobSubmitter; #[async_trait] @@ -20,21 +20,12 @@ impl QueryJobSubmitter for SpiderClient { /// Task-graph construction and submission are not implemented yet. async fn submit_query_job( &self, - query_job_id: QueryJobId, - resource_group_id: ResourceGroupId, - clp_s_query_option: ClpSQueryOption, - output_handle: OutputHandle, - archives: Vec<(Option, NonEmptyString)>, - query_task_execution_policy: ExecutionPolicy, + _query_job_id: QueryJobId, + _resource_group_id: ResourceGroupId, + _clp_s_query_option: ClpSQueryOption, + _output_handle: OutputHandle, + _archives_to_search: Vec<(ArchiveMetadata, ExecutionPolicy)>, ) -> Result { - let _ = ( - query_job_id, - resource_group_id, - clp_s_query_option, - output_handle, - archives, - query_task_execution_policy, - ); todo!("Construct and submit the CLP-S query task graph") } } From 2e689fe52a918fced5ec505a71d9fb2bae791cd0 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Wed, 2 Sep 2026 17:42:08 -0400 Subject: [PATCH 33/44] Add trailing comma --- Cargo.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Cargo.toml b/Cargo.toml index 3f1f57c0c0..5b1de8aa1c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -5,7 +5,7 @@ members = [ "components/clp-tdl-package", "components/compression-coordinator", "components/log-ingestor", - "components/query-coordinator" + "components/query-coordinator", ] resolver = "3" From b2fed1f539006c6d759690efaff8664ee4edbc38 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Wed, 2 Sep 2026 17:44:36 -0400 Subject: [PATCH 34/44] Apply batched suggestions from code review Co-authored-by: Lin Zhihao <59785146+LinZhihao-723@users.noreply.github.com> --- components/query-coordinator/src/query_job_submitter/mod.rs | 4 ++-- .../query-coordinator/src/query_job_submitter/spider.rs | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/components/query-coordinator/src/query_job_submitter/mod.rs b/components/query-coordinator/src/query_job_submitter/mod.rs index c7ba7d5499..3a78a63dd1 100644 --- a/components/query-coordinator/src/query_job_submitter/mod.rs +++ b/components/query-coordinator/src/query_job_submitter/mod.rs @@ -14,7 +14,7 @@ use spider_core::types::id::ResourceGroupId; use crate::Error; -/// Coordinator-side metadata for an archive query task. +/// Identifies an archive handled by query tasks. #[derive(Clone, Debug, Eq, PartialEq)] pub struct ArchiveMetadata { /// The archive's ID. @@ -27,7 +27,7 @@ pub struct ArchiveMetadata { pub size: u64, } -/// Registers CLP-S query jobs with a distributed task scheduler. +/// Drives CLP query jobs on a Spider (Huntsman) cluster. #[async_trait] pub trait QueryJobSubmitter: Clone + Send + Sync { /// Registers, but does not start, one query task per archive. diff --git a/components/query-coordinator/src/query_job_submitter/spider.rs b/components/query-coordinator/src/query_job_submitter/spider.rs index 84f0b78a4b..3171f42bd1 100644 --- a/components/query-coordinator/src/query_job_submitter/spider.rs +++ b/components/query-coordinator/src/query_job_submitter/spider.rs @@ -1,4 +1,4 @@ -//! [`QueryJobSubmitter`] skeleton for [`spider_client::SpiderClient`]. +//! [`QueryJobSubmitter`] implementation for [`spider_client::SpiderClient`]. use async_trait::async_trait; use clp_rust_utils::job_config::QueryJobId; From e5f0501df662a6382e6651230b3d94ae741eb544 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Wed, 2 Sep 2026 17:46:10 -0400 Subject: [PATCH 35/44] Apply suggestion from @LinZhihao-723 Co-authored-by: Lin Zhihao <59785146+LinZhihao-723@users.noreply.github.com> --- .../src/query_job_submitter/mod.rs | 16 +++++++++++++++- 1 file changed, 15 insertions(+), 1 deletion(-) diff --git a/components/query-coordinator/src/query_job_submitter/mod.rs b/components/query-coordinator/src/query_job_submitter/mod.rs index 3a78a63dd1..e02a57b4ce 100644 --- a/components/query-coordinator/src/query_job_submitter/mod.rs +++ b/components/query-coordinator/src/query_job_submitter/mod.rs @@ -30,7 +30,21 @@ pub struct ArchiveMetadata { /// Drives CLP query jobs on a Spider (Huntsman) cluster. #[async_trait] pub trait QueryJobSubmitter: Clone + Send + Sync { - /// Registers, but does not start, one query task per archive. + /// Builds the query task graph for the given archives and registers it with Spider, without + /// starting it. + /// + /// # Parameters + /// + /// * `query_job_id` - The unique ID of the CLP query job. + /// * `resource_group_id` - The Spider resource group to register the job under. + /// * `clp_s_query_options` - `clp-s` query options shared by every task in the job. + /// * `output_handle` - The output handle selecting how the query outputs are returned. + /// * `archives_to_search` - The archives to search, each represents a query task paired with + /// the task execution policy. + /// + /// # Returns + /// + /// The job ID issued by Spider on success. /// /// # Errors /// From e0566a733580a20dde806fae5085afb3a54cdffd Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Wed, 2 Sep 2026 18:08:03 -0400 Subject: [PATCH 36/44] docs(query-coordinator): Fix deferred API documentation --- components/clp-tdl-package/src/task/query/mod.rs | 3 +++ components/query-coordinator/src/query_job_submitter/mod.rs | 2 +- .../query-coordinator/src/query_job_submitter/spider.rs | 6 +++--- 3 files changed, 7 insertions(+), 4 deletions(-) diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs index 76697069a5..df35c84415 100644 --- a/components/clp-tdl-package/src/task/query/mod.rs +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -9,6 +9,9 @@ use spider_tdl::TaskContext; use spider_tdl::task; #[task(name = "query::clp_s_search")] +/// # Panics +/// +/// Panics because `clp-s` search tasks are not implemented yet. pub(crate) fn clp_s_search_task( _ctx: TaskContext, _query_job_id: QueryJobId, diff --git a/components/query-coordinator/src/query_job_submitter/mod.rs b/components/query-coordinator/src/query_job_submitter/mod.rs index e02a57b4ce..ba0aaf53d0 100644 --- a/components/query-coordinator/src/query_job_submitter/mod.rs +++ b/components/query-coordinator/src/query_job_submitter/mod.rs @@ -37,7 +37,7 @@ pub trait QueryJobSubmitter: Clone + Send + Sync { /// /// * `query_job_id` - The unique ID of the CLP query job. /// * `resource_group_id` - The Spider resource group to register the job under. - /// * `clp_s_query_options` - `clp-s` query options shared by every task in the job. + /// * `clp_s_query_option` - `clp-s` query options shared by every task in the job. /// * `output_handle` - The output handle selecting how the query outputs are returned. /// * `archives_to_search` - The archives to search, each represents a query task paired with /// the task execution policy. diff --git a/components/query-coordinator/src/query_job_submitter/spider.rs b/components/query-coordinator/src/query_job_submitter/spider.rs index 3171f42bd1..cdbb2e3e71 100644 --- a/components/query-coordinator/src/query_job_submitter/spider.rs +++ b/components/query-coordinator/src/query_job_submitter/spider.rs @@ -15,9 +15,9 @@ use crate::query_job_submitter::QueryJobSubmitter; #[async_trait] impl QueryJobSubmitter for SpiderClient { - /// # Errors + /// # Panics /// - /// Task-graph construction and submission are not implemented yet. + /// Panics because task-graph construction and submission are not implemented yet. async fn submit_query_job( &self, _query_job_id: QueryJobId, @@ -26,6 +26,6 @@ impl QueryJobSubmitter for SpiderClient { _output_handle: OutputHandle, _archives_to_search: Vec<(ArchiveMetadata, ExecutionPolicy)>, ) -> Result { - todo!("Construct and submit the CLP-S query task graph") + todo!("construct and submit the clp-s query task graph") } } From a714f6a45a136543987cb4d04f77d245252444f3 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Wed, 2 Sep 2026 19:10:48 -0400 Subject: [PATCH 37/44] Fix according to guidelines --- .../initialize-orchestration-db.py | 33 +--- components/query-coordinator/src/error.rs | 10 + .../query-coordinator/src/job_handle.rs | 172 +++++++++++------- .../src/query_job_submitter/mod.rs | 16 +- .../src/query_job_submitter/spider.rs | 29 ++- 5 files changed, 155 insertions(+), 105 deletions(-) diff --git a/components/clp-py-utils/clp_py_utils/initialize-orchestration-db.py b/components/clp-py-utils/clp_py_utils/initialize-orchestration-db.py index 1b33f52a84..1dd800cd99 100644 --- a/components/clp-py-utils/clp_py_utils/initialize-orchestration-db.py +++ b/components/clp-py-utils/clp_py_utils/initialize-orchestration-db.py @@ -11,7 +11,7 @@ QueryJobStatus, QueryTaskStatus, ) -from mysql.connector.errorcode import ER_DUP_FIELDNAME, ER_DUP_KEYNAME +from mysql.connector.errorcode import ER_DUP_KEYNAME from pydantic import ValidationError from clp_py_utils.clp_config import ( @@ -150,37 +150,6 @@ def main(argv): """ ) - # Upgrade query-job tables created before Spider lifecycle support was added. - query_jobs_table_upgrades = ( - ( - f""" - ALTER TABLE `{QUERY_JOBS_TABLE_NAME}` - ADD COLUMN `status_msg` VARCHAR(512) NOT NULL DEFAULT '' AFTER `status` - """, - ER_DUP_FIELDNAME, - ), - ( - f""" - ALTER TABLE `{QUERY_JOBS_TABLE_NAME}` - ADD COLUMN `spider_id` BIGINT UNSIGNED NULL DEFAULT NULL AFTER `job_config` - """, - ER_DUP_FIELDNAME, - ), - ( - f""" - ALTER TABLE `{QUERY_JOBS_TABLE_NAME}` - ADD INDEX `JOB_SPIDER_ID` (`spider_id`) USING BTREE - """, - ER_DUP_KEYNAME, - ), - ) - for upgrade_query, duplicate_error_code in query_jobs_table_upgrades: - try: - scheduling_db_cursor.execute(upgrade_query) - except Exception as err: - if not (hasattr(err, "errno") and err.errno == duplicate_error_code): - raise - scheduling_db_cursor.execute( f""" CREATE TABLE IF NOT EXISTS `{QUERY_TASKS_TABLE_NAME}` ( diff --git a/components/query-coordinator/src/error.rs b/components/query-coordinator/src/error.rs index bd7b31135c..c4cb07e2ad 100644 --- a/components/query-coordinator/src/error.rs +++ b/components/query-coordinator/src/error.rs @@ -12,6 +12,16 @@ pub enum Error { #[error("sqlx error: {0}")] Sqlx(#[from] sqlx::Error), + #[error("failed to persist terminal status for query job {query_job_id}: {source}")] + TerminalStatusPersistence { + /// The query job whose terminal state could not be persisted. + query_job_id: clp_rust_utils::job_config::QueryJobId, + + /// The persistence failure. + #[source] + source: sqlx::Error, + }, + #[error("number of query tasks {0} exceeds `i32::MAX`")] TooManyQueryTasks(usize), } diff --git a/components/query-coordinator/src/job_handle.rs b/components/query-coordinator/src/job_handle.rs index 8e21208e25..67947118c1 100644 --- a/components/query-coordinator/src/job_handle.rs +++ b/components/query-coordinator/src/job_handle.rs @@ -8,13 +8,13 @@ use clp_rust_utils::job_config::QueryJobId; use clp_rust_utils::job_config::QueryJobStatus; use clp_rust_utils::task_io::query::ClpSQueryOption; use clp_rust_utils::task_io::query::OutputHandle; -use non_empty_string::NonEmptyString; use spider_core::task::ExecutionPolicy; use spider_core::types::id::JobId as SpiderJobId; use spider_core::types::id::ResourceGroupId; use sqlx::MySqlPool; use crate::Error; +use crate::query_job_submitter::ArchiveMetadata; use crate::query_job_submitter::QueryJobOutcome; use crate::query_job_submitter::QueryJobSubmitter; @@ -26,16 +26,13 @@ pub struct QueryPlan { /// Result destination shared by every archive task. pub output_handle: OutputHandle, - /// One optional dataset and non-empty archive ID pair per task. - pub archives: Vec<(Option, NonEmptyString)>, - - /// Execution policy applied to every archive task. - pub query_task_execution_policy: ExecutionPolicy, + /// The archives to query, each paired with its task execution policy. + pub archives_to_search: Vec<(ArchiveMetadata, ExecutionPolicy)>, } /// Spider polling options shared by query-job handles. pub struct SpiderOption { - /// Delay before the first Spider job-state poll. + /// Initial delay after a non-terminal Spider job-state poll. pub initial_poll_backoff: Duration, /// Maximum delay between Spider job-state polls. @@ -43,6 +40,10 @@ pub struct SpiderOption { } /// Drives one already-planned query job through submission and terminal persistence. +/// +/// # Type Parameters +/// +/// * `SubmitterType` - The type of the job submitter for Spider job submission. pub struct QueryJobHandle { db_pool: MySqlPool, query_job_id: QueryJobId, @@ -53,7 +54,11 @@ pub struct QueryJobHandle { } impl QueryJobHandle { - /// Constructs a handle for an already-planned query job. + /// Factory function. + /// + /// # Returns + /// + /// A newly created [`QueryJobHandle`] for the given already-planned query job. pub fn new( db_pool: MySqlPool, query_job_id: QueryJobId, @@ -74,21 +79,29 @@ impl QueryJobHandle { /// Submits the prepared graph and drives the query job to a terminal state. /// - /// On an orchestration failure, this method makes a best-effort attempt to mark the CLP query - /// job as failed before returning the original error. + /// On a submission failure, this method makes a best-effort attempt to mark the CLP query job + /// as failed before returning the original error. After the job is durably running, monitoring + /// and terminal-persistence failures leave it running so recovery can reattach to Spider. /// /// # Errors /// - /// Returns an error if submission, submission persistence, polling, or terminal persistence - /// fails. + /// Returns an error if: + /// + /// * Forwards [`Self::submit`]'s return values on failure. + /// * Forwards [`Self::to_completion`]'s return values on failure. pub async fn run(self) -> Result<(), Error> { - tracing::info!(query_job_id = %self.query_job_id, "Starting query job."); + tracing::info!(query_job_id = % self.query_job_id, "Starting query job."); - let result = self.submit_and_wait().await; - if let Err(error) = &result { - self.report_failure(error).await; - } - result + let spider_job_id = match self.submit().await { + Ok(spider_job_id) => spider_job_id, + Err(error) => { + if !matches!(error, Error::JobNotPending(_)) { + self.report_failure(&error).await; + } + return Err(error); + } + }; + self.to_completion(spider_job_id).await } /// Resumes a query job that was already submitted to Spider. @@ -97,23 +110,34 @@ impl QueryJobHandle { /// /// # Errors /// - /// Returns an error if polling or terminal persistence fails. + /// Returns an error if: + /// + /// * Forwards [`Self::to_completion`]'s return values on failure. pub async fn recover(self, spider_job_id: SpiderJobId) -> Result<(), Error> { tracing::info!( - query_job_id = %self.query_job_id, - spider_job_id = %spider_job_id, + query_job_id = % self.query_job_id, + spider_job_id = % spider_job_id, "Recovering query job.", ); - let result = self.to_completion(spider_job_id).await; - if let Err(error) = &result { - self.report_failure(error).await; - } - result + self.to_completion(spider_job_id).await } - async fn submit_and_wait(&self) -> Result<(), Error> { - let num_tasks = self.query_plan.archives.len(); + /// Submits the query job to Spider and persists its running state. + /// + /// # Returns + /// + /// The submitted Spider job ID on success. + /// + /// # Errors + /// + /// Returns an error if: + /// + /// * [`Error::TooManyQueryTasks`] if the number of query tasks exceeds `i32`'s range. + /// * Forwards [`QueryJobSubmitter::submit_query_job`]'s return values on failure. + /// * Forwards [`Self::persist_submission`]'s return values on failure. + async fn submit(&self) -> Result { + let num_tasks = self.query_plan.archives_to_search.len(); let persisted_num_tasks = i32::try_from(num_tasks).map_err(|_| Error::TooManyQueryTasks(num_tasks))?; let spider_job_id = self @@ -123,23 +147,30 @@ impl QueryJobHandle { self.resource_group_id, self.query_plan.clp_s_query_option.clone(), self.query_plan.output_handle.clone(), - self.query_plan.archives.clone(), - self.query_plan.query_task_execution_policy.clone(), + self.query_plan.archives_to_search.clone(), ) .await?; tracing::info!( - query_job_id = %self.query_job_id, - spider_job_id = %spider_job_id, + query_job_id = % self.query_job_id, + spider_job_id = % spider_job_id, num_tasks, "Query job submitted.", ); self.persist_submission(spider_job_id, persisted_num_tasks) .await?; - self.to_completion(spider_job_id).await + Ok(spider_job_id) } + /// Persists the Spider job ID and marks the query job as running. + /// + /// # Errors + /// + /// Returns an error if: + /// + /// * [`Error::JobNotPending`] if the query job is no longer pending. + /// * Forwards [`sqlx::query::Query::execute`]'s return values on failure. async fn persist_submission( &self, spider_job_id: SpiderJobId, @@ -164,6 +195,14 @@ impl QueryJobHandle { Ok(()) } + /// Waits for the associated Spider job to complete and finalizes the query job. + /// + /// # Errors + /// + /// Returns an error if: + /// + /// * [`Error::TerminalStatusPersistence`] if the terminal query-job status cannot be persisted. + /// * Forwards [`QueryJobSubmitter::run_query_job_to_completion`]'s return values on failure. async fn to_completion(&self, spider_job_id: SpiderJobId) -> Result<(), Error> { let outcome = self .job_submitter @@ -175,40 +214,39 @@ impl QueryJobHandle { .await?; tracing::info!( - query_job_id = %self.query_job_id, - spider_job_id = %spider_job_id, - outcome = ?outcome, + query_job_id = % self.query_job_id, + spider_job_id = % spider_job_id, + outcome = ? outcome, "Query job reached a terminal Spider state.", ); - match outcome { - QueryJobOutcome::Succeeded => { - self.update_terminal_status(QueryJobStatus::Succeeded, "", false) - .await - } - QueryJobOutcome::Failed { error_message } => { - self.update_terminal_status( - QueryJobStatus::Failed, - &format!("The Spider query job failed: {error_message}"), - false, - ) - .await - } - QueryJobOutcome::UnexpectedlyCancelled => { - self.update_terminal_status( - QueryJobStatus::Failed, - "The Spider query job was unexpectedly cancelled.", - false, - ) - .await - } - } + let (status, status_message) = match outcome { + QueryJobOutcome::Succeeded => (QueryJobStatus::Succeeded, String::new()), + QueryJobOutcome::Failed { error_message } => ( + QueryJobStatus::Failed, + format!("The Spider query job failed: {error_message}"), + ), + QueryJobOutcome::UnexpectedlyCancelled => ( + QueryJobStatus::Failed, + "The Spider query job was unexpectedly cancelled.".to_string(), + ), + }; + self.update_terminal_status(status, &status_message, false) + .await + .map_err(|source| Error::TerminalStatusPersistence { + query_job_id: self.query_job_id, + source, + }) } + /// Reports a query-job orchestration failure. + /// + /// Logs the original error and makes a best-effort attempt to mark the query job as failed. If + /// terminal-status persistence fails, the status-update error is logged and otherwise ignored. async fn report_failure(&self, error: &Error) { tracing::error!( - query_job_id = %self.query_job_id, - error = %error, + query_job_id = % self.query_job_id, + error = % error, "Query-job orchestration failed.", ); @@ -221,8 +259,8 @@ impl QueryJobHandle { .await { tracing::error!( - query_job_id = %self.query_job_id, - error = %status_error, + query_job_id = % self.query_job_id, + error = % status_error, "Failed to persist the query-job failure.", ); } @@ -230,12 +268,20 @@ impl QueryJobHandle { /// Updates a non-terminal query job while preserving every existing terminal or cancellation /// state. When `allow_pending` is false, only a running job may transition. + /// A zero-row update is treated as success so an ineligible or missing job row is left + /// unchanged. + /// + /// # Errors + /// + /// Returns an error if: + /// + /// * Forwards [`sqlx::query::Query::execute`]'s return values on failure. async fn update_terminal_status( &self, status: QueryJobStatus, status_message: &str, allow_pending: bool, - ) -> Result<(), Error> { + ) -> Result<(), sqlx::Error> { let eligible_statuses = if allow_pending { "?, ?" } else { "?" }; let query = format!( "UPDATE `{QUERY_JOBS_TABLE_NAME}` SET `status` = ?, `status_msg` = LEFT(?, 512), \ diff --git a/components/query-coordinator/src/query_job_submitter/mod.rs b/components/query-coordinator/src/query_job_submitter/mod.rs index 6aaeb12278..8b6679ca2c 100644 --- a/components/query-coordinator/src/query_job_submitter/mod.rs +++ b/components/query-coordinator/src/query_job_submitter/mod.rs @@ -36,14 +36,16 @@ pub enum QueryJobOutcome { Succeeded, /// At least one archive query failed. - Failed { error_message: String }, + Failed { + /// The error reported by Spider. + error_message: String, + }, /// Spider cancelled the job unexpectedly. User-requested cancellation is outside the MVP. UnexpectedlyCancelled, } /// Drives CLP query jobs on a Spider (Huntsman) cluster. ->>>>>>> query-coordinator/crate #[async_trait] pub trait QueryJobSubmitter: Clone + Send + Sync { /// Builds the query task graph for the given archives and registers it with Spider, without @@ -76,6 +78,16 @@ pub trait QueryJobSubmitter: Clone + Send + Sync { /// Idempotently starts `spider_job_id` and waits for it to reach a terminal state. /// + /// # Parameters + /// + /// * `spider_job_id` - The ID of the Spider job to start and monitor. + /// * `initial_poll_backoff` - The initial delay after a non-terminal job-state poll. + /// * `max_poll_backoff` - The maximum delay between job-state polls. + /// + /// # Returns + /// + /// The terminal query-job outcome on success. + /// /// # Errors /// /// Implementations must document their error conditions. diff --git a/components/query-coordinator/src/query_job_submitter/spider.rs b/components/query-coordinator/src/query_job_submitter/spider.rs index 9260e0489c..393de8e13a 100644 --- a/components/query-coordinator/src/query_job_submitter/spider.rs +++ b/components/query-coordinator/src/query_job_submitter/spider.rs @@ -38,8 +38,13 @@ impl QueryJobSubmitter for SpiderClient { /// /// Returns an error if: /// - /// * Starting the job fails for a reason other than it already having been started. - /// * Fetching the Spider job state fails. + /// * Forwards [`SpiderClient::start_job`]'s return values on failure, except + /// [`ClientError::InvalidJobState`]. + /// * Forwards [`SpiderClient::get_job_state`]'s return values on failure. + /// + /// # Panics + /// + /// Panics if Spider returns a terminal state without a corresponding [`QueryJobOutcome`]. async fn run_query_job_to_completion( &self, spider_job_id: JobId, @@ -67,12 +72,20 @@ impl QueryJobSubmitter for SpiderClient { Ok(match terminal_state { JobState::Succeeded => QueryJobOutcome::Succeeded, - JobState::Failed => QueryJobOutcome::Failed { - error_message: self - .get_job_error(spider_job_id) - .await - .unwrap_or_else(|error| format!("")), - }, + JobState::Failed => { + let error_message = match self.get_job_error(spider_job_id).await { + Ok(error_message) => error_message, + Err(error) => { + tracing::warn!( + spider_job_id = % spider_job_id, + error = % error, + "Failed to fetch the Spider job error.", + ); + format!("") + } + }; + QueryJobOutcome::Failed { error_message } + } JobState::Cancelled => QueryJobOutcome::UnexpectedlyCancelled, _ => unreachable!("a terminal Spider state must have a terminal outcome"), }) From 012c8b2b3a0f742cc668f53a48f62fed82955bb5 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Wed, 2 Sep 2026 19:12:24 -0400 Subject: [PATCH 38/44] Lint fix --- components/query-coordinator/src/job_handle.rs | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/components/query-coordinator/src/job_handle.rs b/components/query-coordinator/src/job_handle.rs index 67947118c1..0164beaaf9 100644 --- a/components/query-coordinator/src/job_handle.rs +++ b/components/query-coordinator/src/job_handle.rs @@ -285,9 +285,8 @@ impl QueryJobHandle { let eligible_statuses = if allow_pending { "?, ?" } else { "?" }; let query = format!( "UPDATE `{QUERY_JOBS_TABLE_NAME}` SET `status` = ?, `status_msg` = LEFT(?, 512), \ - `duration` = \ - CASE WHEN `start_time` IS NULL THEN 0 ELSE TIMESTAMPDIFF(MICROSECOND, `start_time`, \ - CURRENT_TIMESTAMP(3)) / 1000000.0 END WHERE `id` = ? AND `status` IN \ + `duration` = CASE WHEN `start_time` IS NULL THEN 0 ELSE TIMESTAMPDIFF(MICROSECOND, \ + `start_time`, CURRENT_TIMESTAMP(3)) / 1000000.0 END WHERE `id` = ? AND `status` IN \ ({eligible_statuses})" ); let mut query = sqlx::query(&query) From 4891cf858d43b25918d96c3ec4bbd0bced9ab7df Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Thu, 3 Sep 2026 13:45:47 -0400 Subject: [PATCH 39/44] fix(query-coordinator): Return errors from staging paths --- components/clp-tdl-package/src/task/query/mod.rs | 7 +++---- components/query-coordinator/src/error.rs | 3 +++ .../query-coordinator/src/query_job_submitter/spider.rs | 7 ++++--- 3 files changed, 10 insertions(+), 7 deletions(-) diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs index df35c84415..849ad54589 100644 --- a/components/clp-tdl-package/src/task/query/mod.rs +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -9,9 +9,6 @@ use spider_tdl::TaskContext; use spider_tdl::task; #[task(name = "query::clp_s_search")] -/// # Panics -/// -/// Panics because `clp-s` search tasks are not implemented yet. pub(crate) fn clp_s_search_task( _ctx: TaskContext, _query_job_id: QueryJobId, @@ -20,5 +17,7 @@ pub(crate) fn clp_s_search_task( _archive_id: ArchiveId, _output_handle: OutputHandle, ) -> Result<(), spider_tdl::TdlError> { - todo!("clp-s search task is not implemented") + Err(spider_tdl::TdlError::ExecutionError( + "clp-s search task is not implemented".to_owned(), + )) } diff --git a/components/query-coordinator/src/error.rs b/components/query-coordinator/src/error.rs index 2d0be81fd4..afc5700af3 100644 --- a/components/query-coordinator/src/error.rs +++ b/components/query-coordinator/src/error.rs @@ -3,6 +3,9 @@ /// Errors returned by the query coordinator. #[derive(Debug, thiserror::Error)] pub enum Error { + #[error("query job submission is not implemented")] + QueryJobSubmissionNotImplemented, + #[error("spider request failure: {0}")] SpiderClient(#[from] spider_client::error::ClientError), } diff --git a/components/query-coordinator/src/query_job_submitter/spider.rs b/components/query-coordinator/src/query_job_submitter/spider.rs index cdbb2e3e71..bd92e3b101 100644 --- a/components/query-coordinator/src/query_job_submitter/spider.rs +++ b/components/query-coordinator/src/query_job_submitter/spider.rs @@ -15,9 +15,10 @@ use crate::query_job_submitter::QueryJobSubmitter; #[async_trait] impl QueryJobSubmitter for SpiderClient { - /// # Panics + /// # Errors /// - /// Panics because task-graph construction and submission are not implemented yet. + /// Returns [`Error::QueryJobSubmissionNotImplemented`] until task-graph construction and + /// submission are implemented. async fn submit_query_job( &self, _query_job_id: QueryJobId, @@ -26,6 +27,6 @@ impl QueryJobSubmitter for SpiderClient { _output_handle: OutputHandle, _archives_to_search: Vec<(ArchiveMetadata, ExecutionPolicy)>, ) -> Result { - todo!("construct and submit the clp-s query task graph") + Err(Error::QueryJobSubmissionNotImplemented) } } From 4a45d04248e98f90f18fe994d2b5e3508d4cd90c Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Thu, 3 Sep 2026 13:47:07 -0400 Subject: [PATCH 40/44] Change back to todo --- components/clp-tdl-package/src/task/query/mod.rs | 4 +--- components/query-coordinator/src/error.rs | 3 --- .../query-coordinator/src/query_job_submitter/spider.rs | 6 +----- 3 files changed, 2 insertions(+), 11 deletions(-) diff --git a/components/clp-tdl-package/src/task/query/mod.rs b/components/clp-tdl-package/src/task/query/mod.rs index 849ad54589..76697069a5 100644 --- a/components/clp-tdl-package/src/task/query/mod.rs +++ b/components/clp-tdl-package/src/task/query/mod.rs @@ -17,7 +17,5 @@ pub(crate) fn clp_s_search_task( _archive_id: ArchiveId, _output_handle: OutputHandle, ) -> Result<(), spider_tdl::TdlError> { - Err(spider_tdl::TdlError::ExecutionError( - "clp-s search task is not implemented".to_owned(), - )) + todo!("clp-s search task is not implemented") } diff --git a/components/query-coordinator/src/error.rs b/components/query-coordinator/src/error.rs index afc5700af3..2d0be81fd4 100644 --- a/components/query-coordinator/src/error.rs +++ b/components/query-coordinator/src/error.rs @@ -3,9 +3,6 @@ /// Errors returned by the query coordinator. #[derive(Debug, thiserror::Error)] pub enum Error { - #[error("query job submission is not implemented")] - QueryJobSubmissionNotImplemented, - #[error("spider request failure: {0}")] SpiderClient(#[from] spider_client::error::ClientError), } diff --git a/components/query-coordinator/src/query_job_submitter/spider.rs b/components/query-coordinator/src/query_job_submitter/spider.rs index bd92e3b101..c5c33f596d 100644 --- a/components/query-coordinator/src/query_job_submitter/spider.rs +++ b/components/query-coordinator/src/query_job_submitter/spider.rs @@ -15,10 +15,6 @@ use crate::query_job_submitter::QueryJobSubmitter; #[async_trait] impl QueryJobSubmitter for SpiderClient { - /// # Errors - /// - /// Returns [`Error::QueryJobSubmissionNotImplemented`] until task-graph construction and - /// submission are implemented. async fn submit_query_job( &self, _query_job_id: QueryJobId, @@ -27,6 +23,6 @@ impl QueryJobSubmitter for SpiderClient { _output_handle: OutputHandle, _archives_to_search: Vec<(ArchiveMetadata, ExecutionPolicy)>, ) -> Result { - Err(Error::QueryJobSubmissionNotImplemented) + todo!("construct and submit the clp-s query task graph") } } From 61414eb48be6d51b82bebdc214e7ac77f4cabdfa Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Thu, 3 Sep 2026 13:53:38 -0400 Subject: [PATCH 41/44] lint fix --- components/query-coordinator/src/job_handle.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/components/query-coordinator/src/job_handle.rs b/components/query-coordinator/src/job_handle.rs index 0164beaaf9..c8bac666b4 100644 --- a/components/query-coordinator/src/job_handle.rs +++ b/components/query-coordinator/src/job_handle.rs @@ -59,7 +59,7 @@ impl QueryJobHandle { /// # Returns /// /// A newly created [`QueryJobHandle`] for the given already-planned query job. - pub fn new( + pub const fn new( db_pool: MySqlPool, query_job_id: QueryJobId, job_submitter: SubmitterType, From 5e84b4f7253a4fd4e25c7cf9c1834de9c0008676 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Thu, 3 Sep 2026 13:57:17 -0400 Subject: [PATCH 42/44] Lint fix --- components/query-coordinator/src/job_handle.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/components/query-coordinator/src/job_handle.rs b/components/query-coordinator/src/job_handle.rs index c8bac666b4..c591de88e8 100644 --- a/components/query-coordinator/src/job_handle.rs +++ b/components/query-coordinator/src/job_handle.rs @@ -138,7 +138,7 @@ impl QueryJobHandle { /// * Forwards [`Self::persist_submission`]'s return values on failure. async fn submit(&self) -> Result { let num_tasks = self.query_plan.archives_to_search.len(); - let persisted_num_tasks = + let _persisted_num_tasks = i32::try_from(num_tasks).map_err(|_| Error::TooManyQueryTasks(num_tasks))?; let spider_job_id = self .job_submitter @@ -158,7 +158,7 @@ impl QueryJobHandle { "Query job submitted.", ); - self.persist_submission(spider_job_id, persisted_num_tasks) + self.persist_submission(spider_job_id, _persisted_num_tasks) .await?; Ok(spider_job_id) } From bd8d4b4fc67e652f2b604d77d29dcc1c11caa8e4 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Thu, 3 Sep 2026 14:02:25 -0400 Subject: [PATCH 43/44] Add a submit task and fix unused variable --- .../query-coordinator/src/job_handle.rs | 40 +++++++++++++------ 1 file changed, 27 insertions(+), 13 deletions(-) diff --git a/components/query-coordinator/src/job_handle.rs b/components/query-coordinator/src/job_handle.rs index c591de88e8..358276acba 100644 --- a/components/query-coordinator/src/job_handle.rs +++ b/components/query-coordinator/src/job_handle.rs @@ -134,22 +134,13 @@ impl QueryJobHandle { /// Returns an error if: /// /// * [`Error::TooManyQueryTasks`] if the number of query tasks exceeds `i32`'s range. - /// * Forwards [`QueryJobSubmitter::submit_query_job`]'s return values on failure. + /// * Forwards [`Self::submit_to_spider`]'s return values on failure. /// * Forwards [`Self::persist_submission`]'s return values on failure. async fn submit(&self) -> Result { let num_tasks = self.query_plan.archives_to_search.len(); - let _persisted_num_tasks = + let persisted_num_tasks = i32::try_from(num_tasks).map_err(|_| Error::TooManyQueryTasks(num_tasks))?; - let spider_job_id = self - .job_submitter - .submit_query_job( - self.query_job_id, - self.resource_group_id, - self.query_plan.clp_s_query_option.clone(), - self.query_plan.output_handle.clone(), - self.query_plan.archives_to_search.clone(), - ) - .await?; + let spider_job_id = self.submit_to_spider().await?; tracing::info!( query_job_id = % self.query_job_id, @@ -158,11 +149,34 @@ impl QueryJobHandle { "Query job submitted.", ); - self.persist_submission(spider_job_id, _persisted_num_tasks) + self.persist_submission(spider_job_id, persisted_num_tasks) .await?; Ok(spider_job_id) } + /// Submits the prepared query graph to Spider. + /// + /// # Returns + /// + /// The submitted Spider job ID on success. + /// + /// # Errors + /// + /// Returns an error if: + /// + /// * Forwards [`QueryJobSubmitter::submit_query_job`]'s return values on failure. + async fn submit_to_spider(&self) -> Result { + self.job_submitter + .submit_query_job( + self.query_job_id, + self.resource_group_id, + self.query_plan.clp_s_query_option.clone(), + self.query_plan.output_handle.clone(), + self.query_plan.archives_to_search.clone(), + ) + .await + } + /// Persists the Spider job ID and marks the query job as running. /// /// # Errors From 401f9882442ab332212645364582e9202f4cac07 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Tue, 8 Sep 2026 11:15:30 -0400 Subject: [PATCH 44/44] Remove redundant QueryPlan and rename SpiderOption --- components/query-coordinator/src/error.rs | 3 + .../query-coordinator/src/job_handle.rs | 55 ++++++++----------- 2 files changed, 26 insertions(+), 32 deletions(-) diff --git a/components/query-coordinator/src/error.rs b/components/query-coordinator/src/error.rs index c4cb07e2ad..af25857ab9 100644 --- a/components/query-coordinator/src/error.rs +++ b/components/query-coordinator/src/error.rs @@ -24,4 +24,7 @@ pub enum Error { #[error("number of query tasks {0} exceeds `i32::MAX`")] TooManyQueryTasks(usize), + + #[error("no archives were selected for the query job")] + NoArchivesToSearch, } diff --git a/components/query-coordinator/src/job_handle.rs b/components/query-coordinator/src/job_handle.rs index 358276acba..02a8dd6965 100644 --- a/components/query-coordinator/src/job_handle.rs +++ b/components/query-coordinator/src/job_handle.rs @@ -18,18 +18,6 @@ use crate::query_job_submitter::ArchiveMetadata; use crate::query_job_submitter::QueryJobOutcome; use crate::query_job_submitter::QueryJobSubmitter; -/// The coordinator-prepared inputs for one Spider query graph. -pub struct QueryPlan { - /// Query behavior shared by every archive task. - pub clp_s_query_option: ClpSQueryOption, - - /// Result destination shared by every archive task. - pub output_handle: OutputHandle, - - /// The archives to query, each paired with its task execution policy. - pub archives_to_search: Vec<(ArchiveMetadata, ExecutionPolicy)>, -} - /// Spider polling options shared by query-job handles. pub struct SpiderOption { /// Initial delay after a non-terminal Spider job-state poll. @@ -49,7 +37,9 @@ pub struct QueryJobHandle { query_job_id: QueryJobId, job_submitter: SubmitterType, resource_group_id: ResourceGroupId, - query_plan: QueryPlan, + clp_s_query_option: ClpSQueryOption, + output_handle: OutputHandle, + archives_to_search: Vec<(ArchiveMetadata, ExecutionPolicy)>, spider_option: Arc, } @@ -64,7 +54,9 @@ impl QueryJobHandle { query_job_id: QueryJobId, job_submitter: SubmitterType, resource_group_id: ResourceGroupId, - query_plan: QueryPlan, + clp_s_query_option: ClpSQueryOption, + output_handle: OutputHandle, + archives_to_search: Vec<(ArchiveMetadata, ExecutionPolicy)>, spider_option: Arc, ) -> Self { Self { @@ -72,7 +64,9 @@ impl QueryJobHandle { query_job_id, job_submitter, resource_group_id, - query_plan, + clp_s_query_option, + output_handle, + archives_to_search, spider_option, } } @@ -137,7 +131,10 @@ impl QueryJobHandle { /// * Forwards [`Self::submit_to_spider`]'s return values on failure. /// * Forwards [`Self::persist_submission`]'s return values on failure. async fn submit(&self) -> Result { - let num_tasks = self.query_plan.archives_to_search.len(); + let num_tasks = self.archives_to_search.len(); + if num_tasks == 0 { + return Err(Error::NoArchivesToSearch); + } let persisted_num_tasks = i32::try_from(num_tasks).map_err(|_| Error::TooManyQueryTasks(num_tasks))?; let spider_job_id = self.submit_to_spider().await?; @@ -170,9 +167,9 @@ impl QueryJobHandle { .submit_query_job( self.query_job_id, self.resource_group_id, - self.query_plan.clp_s_query_option.clone(), - self.query_plan.output_handle.clone(), - self.query_plan.archives_to_search.clone(), + self.clp_s_query_option.clone(), + self.output_handle.clone(), + self.archives_to_search.clone(), ) .await } @@ -245,7 +242,7 @@ impl QueryJobHandle { "The Spider query job was unexpectedly cancelled.".to_string(), ), }; - self.update_terminal_status(status, &status_message, false) + self.update_terminal_status(status, &status_message, QueryJobStatus::Running) .await .map_err(|source| Error::TerminalStatusPersistence { query_job_id: self.query_job_id, @@ -268,7 +265,7 @@ impl QueryJobHandle { .update_terminal_status( QueryJobStatus::Failed, &format!("Query-job orchestration failed: {error}"), - true, + QueryJobStatus::Pending, ) .await { @@ -280,8 +277,7 @@ impl QueryJobHandle { } } - /// Updates a non-terminal query job while preserving every existing terminal or cancellation - /// state. When `allow_pending` is false, only a running job may transition. + /// Updates a query job only when it has the expected non-terminal status. /// A zero-row update is treated as success so an ineligible or missing job row is left /// unchanged. /// @@ -294,23 +290,18 @@ impl QueryJobHandle { &self, status: QueryJobStatus, status_message: &str, - allow_pending: bool, + expected_status: QueryJobStatus, ) -> Result<(), sqlx::Error> { - let eligible_statuses = if allow_pending { "?, ?" } else { "?" }; let query = format!( "UPDATE `{QUERY_JOBS_TABLE_NAME}` SET `status` = ?, `status_msg` = LEFT(?, 512), \ `duration` = CASE WHEN `start_time` IS NULL THEN 0 ELSE TIMESTAMPDIFF(MICROSECOND, \ - `start_time`, CURRENT_TIMESTAMP(3)) / 1000000.0 END WHERE `id` = ? AND `status` IN \ - ({eligible_statuses})" + `start_time`, CURRENT_TIMESTAMP(3)) / 1000000.0 END WHERE `id` = ? AND `status` = ?" ); - let mut query = sqlx::query(&query) + let query = sqlx::query(&query) .bind(i32::from(status)) .bind(status_message) .bind(self.query_job_id) - .bind(i32::from(QueryJobStatus::Running)); - if allow_pending { - query = query.bind(i32::from(QueryJobStatus::Pending)); - } + .bind(i32::from(expected_status)); query.execute(&self.db_pool).await?; Ok(()) }