From 43f57e9710ee19d9278cdfc277068ea035a0ece5 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Thu, 16 Jul 2026 18:49:05 -0400 Subject: [PATCH 1/6] Bump --- tools/scripts/deps-download/init.sh | 2 +- tools/yscope-dev-utils | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/tools/scripts/deps-download/init.sh b/tools/scripts/deps-download/init.sh index a55db69867..c2f6c5fa7f 100755 --- a/tools/scripts/deps-download/init.sh +++ b/tools/scripts/deps-download/init.sh @@ -8,7 +8,7 @@ script_dir="$( cd "$( dirname "${BASH_SOURCE[0]}" )" &> /dev/null && pwd )" project_root_dir="$script_dir/../../../" download_dep_script="$script_dir/download-dep.py" -readonly YSCOPE_DEV_UTILS_COMMIT_SHA="2facad83f3ee492d704e650269ce04ee35f5f534" +readonly YSCOPE_DEV_UTILS_COMMIT_SHA="0c214c44acddff330a204201428ea145e641891d" python3 "${download_dep_script}" \ "https://github.com/y-scope/yscope-dev-utils/archive/${YSCOPE_DEV_UTILS_COMMIT_SHA}.zip" \ "yscope-dev-utils-${YSCOPE_DEV_UTILS_COMMIT_SHA}" \ diff --git a/tools/yscope-dev-utils b/tools/yscope-dev-utils index 2facad83f3..0c214c44ac 160000 --- a/tools/yscope-dev-utils +++ b/tools/yscope-dev-utils @@ -1 +1 @@ -Subproject commit 2facad83f3ee492d704e650269ce04ee35f5f534 +Subproject commit 0c214c44acddff330a204201428ea145e641891d From f0c1be15369a5472b1af89f97b08e1bf107f1cf6 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Thu, 16 Jul 2026 19:00:45 -0400 Subject: [PATCH 2/6] Rust reformat --- components/api-server/src/bin/api_server.rs | 3 +- components/api-server/src/client.rs | 32 ++++---- components/api-server/src/error.rs | 3 +- components/api-server/src/routes.rs | 75 +++++++++---------- .../src/clp_config/package/config.rs | 3 +- .../src/clp_config/s3_config.rs | 3 +- .../clp-rust-utils/src/database/mysql.rs | 6 +- .../src/job_config/clp_io_config.rs | 8 +- .../src/job_config/compression.rs | 6 +- .../src/job_config/ingestion.rs | 3 +- .../clp-rust-utils/src/job_config/search.rs | 6 +- components/clp-rust-utils/src/logging.rs | 11 +-- components/clp-rust-utils/src/s3/client.rs | 8 +- components/clp-rust-utils/src/sqs/client.rs | 8 +- components/clp-rust-utils/src/telemetry.rs | 6 +- .../clp-rust-utils/tests/clp_config_test.rs | 15 ++-- .../log-ingestor/src/aws_client_manager.rs | 3 +- .../log-ingestor/src/bin/log_ingestor.rs | 6 +- .../compression/compression_job_submitter.rs | 32 ++++---- .../log-ingestor/src/compression/listener.rs | 13 ++-- .../src/ingestion_job/s3_scanner.rs | 12 +-- .../src/ingestion_job/sqs_listener.rs | 25 +++---- .../log-ingestor/src/ingestion_job_manager.rs | 30 ++++---- .../ingestion_job_manager/clp_ingestion.rs | 65 ++++++++-------- components/log-ingestor/src/routes.rs | 42 +++++------ components/log-ingestor/tests/aws_config.rs | 3 +- .../tests/test_compression_listener.rs | 22 +++--- .../log-ingestor/tests/test_ingestion_job.rs | 54 +++++++------ components/log-ingestor/tests/test_scan.rs | 25 +++---- components/log-ingestor/tests/test_utils.rs | 6 +- 30 files changed, 268 insertions(+), 266 deletions(-) diff --git a/components/api-server/src/bin/api_server.rs b/components/api-server/src/bin/api_server.rs index 8976dcffb0..5f6603749a 100644 --- a/components/api-server/src/bin/api_server.rs +++ b/components/api-server/src/bin/api_server.rs @@ -1,6 +1,7 @@ use anyhow::Context; use clap::Parser; -use clp_rust_utils::{clp_config::package, serde::yaml}; +use clp_rust_utils::clp_config::package; +use clp_rust_utils::serde::yaml; #[derive(Parser)] #[command(version, about = "API Server for CLP.")] diff --git a/components/api-server/src/client.rs b/components/api-server/src/client.rs index 01d3ab6c3b..f8e2a3e01c 100644 --- a/components/api-server/src/client.rs +++ b/components/api-server/src/client.rs @@ -1,22 +1,28 @@ use std::pin::Pin; use async_stream::stream; -use chrono::{DateTime, TimeZone, Utc}; +use chrono::DateTime; +use chrono::TimeZone; +use chrono::Utc; +use clp_rust_utils::aws::AWS_DEFAULT_REGION; +use clp_rust_utils::clp_config::package::config::Config; +use clp_rust_utils::clp_config::package::config::StorageEngine; +use clp_rust_utils::clp_config::package::config::StreamOutputStorage; +use clp_rust_utils::clp_config::package::credentials::Credentials; +use clp_rust_utils::database::mysql::create_clp_db_mysql_pool; pub use clp_rust_utils::job_config::CompressionJobStatus; -use clp_rust_utils::{ - aws::AWS_DEFAULT_REGION, - clp_config::package::{ - config::{Config, StorageEngine, StreamOutputStorage}, - credentials::Credentials, - }, - database::mysql::create_clp_db_mysql_pool, - job_config::{QUERY_JOBS_TABLE_NAME, QueryJobStatus, QueryJobType, SearchJobConfig}, -}; -use futures::{Stream, StreamExt}; +use clp_rust_utils::job_config::QUERY_JOBS_TABLE_NAME; +use clp_rust_utils::job_config::QueryJobStatus; +use clp_rust_utils::job_config::QueryJobType; +use clp_rust_utils::job_config::SearchJobConfig; +use futures::Stream; +use futures::StreamExt; use pin_project_lite::pin_project; -use serde::{Deserialize, Serialize}; +use serde::Deserialize; +use serde::Serialize; use sqlx::Row; -use utoipa::{IntoParams, ToSchema}; +use utoipa::IntoParams; +use utoipa::ToSchema; pub use crate::error::ClientError; diff --git a/components/api-server/src/error.rs b/components/api-server/src/error.rs index 484b462b45..3447f3d616 100644 --- a/components/api-server/src/error.rs +++ b/components/api-server/src/error.rs @@ -1,4 +1,5 @@ -use aws_sdk_s3::{error::SdkError, primitives::ByteStreamError}; +use aws_sdk_s3::error::SdkError; +use aws_sdk_s3::primitives::ByteStreamError; use num_enum::TryFromPrimitive; use thiserror::Error; diff --git a/components/api-server/src/routes.rs b/components/api-server/src/routes.rs index c10ac175da..fae107f48a 100644 --- a/components/api-server/src/routes.rs +++ b/components/api-server/src/routes.rs @@ -1,29 +1,31 @@ -use axum::{ - Json, - extract::{Path, Query, State}, - http::StatusCode, - response::{ - IntoResponse, - Sse, - sse::{Event, KeepAlive}, - }, - routing::get, -}; -use futures::{Stream, StreamExt}; -use serde::{Deserialize, Serialize}; +use axum::Json; +use axum::extract::Path; +use axum::extract::Query; +use axum::extract::State; +use axum::http::StatusCode; +use axum::response::IntoResponse; +use axum::response::Sse; +use axum::response::sse::Event; +use axum::response::sse::KeepAlive; +use axum::routing::get; +use futures::Stream; +use futures::StreamExt; +use serde::Deserialize; +use serde::Serialize; use thiserror::Error; -use tower_http::cors::{Any, CorsLayer}; -use utoipa::{OpenApi, ToSchema}; -use utoipa_axum::{router::OpenApiRouter, routes}; - -use crate::client::{ - Client, - ClientError, - CompressionUsage, - CompressionUsageParams, - QueryConfig, - ValidatedCompressionUsageParams, -}; +use tower_http::cors::Any; +use tower_http::cors::CorsLayer; +use utoipa::OpenApi; +use utoipa::ToSchema; +use utoipa_axum::router::OpenApiRouter; +use utoipa_axum::routes; + +use crate::client::Client; +use crate::client::ClientError; +use crate::client::CompressionUsage; +use crate::client::CompressionUsageParams; +use crate::client::QueryConfig; +use crate::client::ValidatedCompressionUsageParams; /// Factory method to create an Axum router configured with all API routes. /// @@ -65,14 +67,12 @@ pub fn from_client(client: Client) -> Result { mod api_doc { // Using `super::...` can cause `super` to appear as a tag in the generated OpenAPI // documentation. Importing the paths directly prevents this issue. - use super::{ - __path_cancel_query, - __path_compression_usage, - __path_health, - __path_query, - __path_query_results, - CompressionUsage, - }; + use super::__path_cancel_query; + use super::__path_compression_usage; + use super::__path_health; + use super::__path_query; + use super::__path_query_results; + use super::CompressionUsage; use crate::client::CompressionJobStatus; #[derive(utoipa::OpenApi)] @@ -337,11 +337,10 @@ impl IntoResponse for HandlerError { #[cfg(test)] mod tests { - use axum::{ - body::Body, - http::{Request, StatusCode}, - routing::get, - }; + use axum::body::Body; + use axum::http::Request; + use axum::http::StatusCode; + use axum::routing::get; use http_body_util::BodyExt; use tower::ServiceExt; 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 98c62c55a5..8064f1e428 100644 --- a/components/clp-rust-utils/src/clp_config/package/config.rs +++ b/components/clp-rust-utils/src/clp_config/package/config.rs @@ -1,6 +1,7 @@ use serde::Deserialize; -use crate::clp_config::{AwsAuthentication, S3Config}; +use crate::clp_config::AwsAuthentication; +use crate::clp_config::S3Config; /// Mirror of `clp_py_utils.clp_config.ClpConfig`. /// diff --git a/components/clp-rust-utils/src/clp_config/s3_config.rs b/components/clp-rust-utils/src/clp_config/s3_config.rs index 0ceac30f91..50792094a1 100644 --- a/components/clp-rust-utils/src/clp_config/s3_config.rs +++ b/components/clp-rust-utils/src/clp_config/s3_config.rs @@ -1,5 +1,6 @@ use non_empty_string::NonEmptyString; -use serde::{Deserialize, Serialize}; +use serde::Deserialize; +use serde::Serialize; /// Represents the configuration for connecting to an S3 bucket. #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] diff --git a/components/clp-rust-utils/src/database/mysql.rs b/components/clp-rust-utils/src/database/mysql.rs index 1ea951184a..9a2aa127e3 100644 --- a/components/clp-rust-utils/src/database/mysql.rs +++ b/components/clp-rust-utils/src/database/mysql.rs @@ -1,10 +1,8 @@ use secrecy::ExposeSecret; use strum::IntoEnumIterator; -use crate::clp_config::package::{ - config::Database as DatabaseConfig, - credentials::Database as DatabaseCredentials, -}; +use crate::clp_config::package::config::Database as DatabaseConfig; +use crate::clp_config::package::credentials::Database as DatabaseCredentials; /// Implements [`sqlx::Type`] for `$ty` by delegating to `$delegate`. /// diff --git a/components/clp-rust-utils/src/job_config/clp_io_config.rs b/components/clp-rust-utils/src/job_config/clp_io_config.rs index 7f4e19cb51..312dac3e53 100644 --- a/components/clp-rust-utils/src/job_config/clp_io_config.rs +++ b/components/clp-rust-utils/src/job_config/clp_io_config.rs @@ -1,11 +1,9 @@ use non_empty_string::NonEmptyString; use serde::Serialize; -use crate::{ - clp_config::S3Config, - job_config::ingestion::JobId as IngestionJobId, - s3::S3ObjectMetadataId, -}; +use crate::clp_config::S3Config; +use crate::job_config::ingestion::JobId as IngestionJobId; +use crate::s3::S3ObjectMetadataId; /// Represents CLP IO config. #[derive(Debug, Clone, PartialEq, Eq, Serialize)] diff --git a/components/clp-rust-utils/src/job_config/compression.rs b/components/clp-rust-utils/src/job_config/compression.rs index ddebf2f880..e09ec10484 100644 --- a/components/clp-rust-utils/src/job_config/compression.rs +++ b/components/clp-rust-utils/src/job_config/compression.rs @@ -1,5 +1,7 @@ -use num_enum::{IntoPrimitive, TryFromPrimitive}; -use serde::{Deserialize, Serialize}; +use num_enum::IntoPrimitive; +use num_enum::TryFromPrimitive; +use serde::Deserialize; +use serde::Serialize; use strum::EnumString; use utoipa::ToSchema; diff --git a/components/clp-rust-utils/src/job_config/ingestion.rs b/components/clp-rust-utils/src/job_config/ingestion.rs index bc3f726ac2..43782e520d 100644 --- a/components/clp-rust-utils/src/job_config/ingestion.rs +++ b/components/clp-rust-utils/src/job_config/ingestion.rs @@ -1,6 +1,7 @@ pub mod s3 { use non_empty_string::NonEmptyString; - use serde::{Deserialize, Serialize}; + use serde::Deserialize; + use serde::Serialize; use thiserror::Error; use utoipa::ToSchema; diff --git a/components/clp-rust-utils/src/job_config/search.rs b/components/clp-rust-utils/src/job_config/search.rs index fcf25275b9..0ece3c2295 100644 --- a/components/clp-rust-utils/src/job_config/search.rs +++ b/components/clp-rust-utils/src/job_config/search.rs @@ -1,5 +1,7 @@ -use num_enum::{IntoPrimitive, TryFromPrimitive}; -use serde::{Deserialize, Serialize}; +use num_enum::IntoPrimitive; +use num_enum::TryFromPrimitive; +use serde::Deserialize; +use serde::Serialize; pub const QUERY_JOBS_TABLE_NAME: &str = "query_jobs"; diff --git a/components/clp-rust-utils/src/logging.rs b/components/clp-rust-utils/src/logging.rs index 07a8431760..bc130b08eb 100644 --- a/components/clp-rust-utils/src/logging.rs +++ b/components/clp-rust-utils/src/logging.rs @@ -1,8 +1,9 @@ -use tracing_appender::{ - non_blocking::{NonBlockingBuilder, WorkerGuard}, - rolling::{RollingFileAppender, Rotation}, -}; -use tracing_subscriber::{self, fmt::writer::MakeWriterExt}; +use tracing_appender::non_blocking::NonBlockingBuilder; +use tracing_appender::non_blocking::WorkerGuard; +use tracing_appender::rolling::RollingFileAppender; +use tracing_appender::rolling::Rotation; +use tracing_subscriber::fmt::writer::MakeWriterExt; +use tracing_subscriber::{self}; /// Opaque struct to hold the worker guards for the background log writers. /// These guards must be held for the lifetime of the program to ensure logs are flushed. diff --git a/components/clp-rust-utils/src/s3/client.rs b/components/clp-rust-utils/src/s3/client.rs index 9160f5d0f1..134b1645a9 100644 --- a/components/clp-rust-utils/src/s3/client.rs +++ b/components/clp-rust-utils/src/s3/client.rs @@ -1,8 +1,8 @@ use aws_config::BehaviorVersion; -use aws_sdk_s3::{ - Client, - config::{Builder, Credentials, Region}, -}; +use aws_sdk_s3::Client; +use aws_sdk_s3::config::Builder; +use aws_sdk_s3::config::Credentials; +use aws_sdk_s3::config::Region; use non_empty_string::NonEmptyString; use crate::clp_config::AwsAuthentication; diff --git a/components/clp-rust-utils/src/sqs/client.rs b/components/clp-rust-utils/src/sqs/client.rs index 4586837146..2d5be46171 100644 --- a/components/clp-rust-utils/src/sqs/client.rs +++ b/components/clp-rust-utils/src/sqs/client.rs @@ -1,8 +1,8 @@ use aws_config::BehaviorVersion; -use aws_sdk_sqs::{ - Client, - config::{Builder, Credentials, Region}, -}; +use aws_sdk_sqs::Client; +use aws_sdk_sqs::config::Builder; +use aws_sdk_sqs::config::Credentials; +use aws_sdk_sqs::config::Region; use non_empty_string::NonEmptyString; use crate::clp_config::AwsAuthentication; diff --git a/components/clp-rust-utils/src/telemetry.rs b/components/clp-rust-utils/src/telemetry.rs index 45d5ba5d32..ee237d12a1 100644 --- a/components/clp-rust-utils/src/telemetry.rs +++ b/components/clp-rust-utils/src/telemetry.rs @@ -2,9 +2,11 @@ use std::env; -use opentelemetry_sdk::metrics::{PeriodicReader, SdkMeterProvider}; +use opentelemetry_sdk::metrics::PeriodicReader; +use opentelemetry_sdk::metrics::SdkMeterProvider; -use crate::{Error, clp_config::package::config::Telemetry}; +use crate::Error; +use crate::clp_config::package::config::Telemetry; /// RAII guard that shuts down the meter provider and flushes pending metric exports when dropped. pub struct TelemetryGuard { diff --git a/components/clp-rust-utils/tests/clp_config_test.rs b/components/clp-rust-utils/tests/clp_config_test.rs index 41fd5edf11..adbfab2b21 100644 --- a/components/clp-rust-utils/tests/clp_config_test.rs +++ b/components/clp-rust-utils/tests/clp_config_test.rs @@ -1,9 +1,12 @@ -use clp_rust_utils::{ - clp_config::{AwsAuthentication, AwsCredentials, S3Config}, - job_config::{ClpIoConfig, InputConfig, OutputConfig, S3ObjectMetadataInputConfig}, - serde::BrotliMsgpack, - types::non_empty_string::ExpectedNonEmpty, -}; +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::job_config::ClpIoConfig; +use clp_rust_utils::job_config::InputConfig; +use clp_rust_utils::job_config::OutputConfig; +use clp_rust_utils::job_config::S3ObjectMetadataInputConfig; +use clp_rust_utils::serde::BrotliMsgpack; +use clp_rust_utils::types::non_empty_string::ExpectedNonEmpty; use non_empty_string::NonEmptyString; use serde_json::Value; diff --git a/components/log-ingestor/src/aws_client_manager.rs b/components/log-ingestor/src/aws_client_manager.rs index d19e420f7a..fd14131f26 100644 --- a/components/log-ingestor/src/aws_client_manager.rs +++ b/components/log-ingestor/src/aws_client_manager.rs @@ -2,7 +2,8 @@ use anyhow::Result; use async_trait::async_trait; use aws_sdk_s3::Client as S3Client; use aws_sdk_sqs::Client as SqsClient; -use clp_rust_utils::{aws::AWS_DEFAULT_REGION, clp_config::AwsAuthentication}; +use clp_rust_utils::aws::AWS_DEFAULT_REGION; +use clp_rust_utils::clp_config::AwsAuthentication; use non_empty_string::NonEmptyString; /// A marker trait for AWS client types. diff --git a/components/log-ingestor/src/bin/log_ingestor.rs b/components/log-ingestor/src/bin/log_ingestor.rs index a766c84156..c7436b39ac 100644 --- a/components/log-ingestor/src/bin/log_ingestor.rs +++ b/components/log-ingestor/src/bin/log_ingestor.rs @@ -1,7 +1,9 @@ use anyhow::Context; use clap::Parser; -use clp_rust_utils::{clp_config::package, serde::yaml}; -use log_ingestor::{ingestion_job_manager::IngestionJobManagerState, routes::create_router}; +use clp_rust_utils::clp_config::package; +use clp_rust_utils::serde::yaml; +use log_ingestor::ingestion_job_manager::IngestionJobManagerState; +use log_ingestor::routes::create_router; #[derive(Parser)] #[command(version, about = "log-ingestor for CLP.")] diff --git a/components/log-ingestor/src/compression/compression_job_submitter.rs b/components/log-ingestor/src/compression/compression_job_submitter.rs index b4f8f6287c..69beadb97c 100644 --- a/components/log-ingestor/src/compression/compression_job_submitter.rs +++ b/components/log-ingestor/src/compression/compression_job_submitter.rs @@ -1,24 +1,20 @@ use anyhow::Result; use async_trait::async_trait; -use clp_rust_utils::{ - clp_config::{ - AwsAuthentication, - S3Config, - package::{DEFAULT_DATASET_NAME, config::ArchiveOutput}, - }, - job_config::{ - ClpIoConfig, - CompressionJobId, - CompressionJobStatus, - InputConfig, - OutputConfig, - S3ObjectMetadataInputConfig, - ingestion::s3::BaseConfig, - }, - s3::S3ObjectMetadataId, -}; +use clp_rust_utils::clp_config::AwsAuthentication; +use clp_rust_utils::clp_config::S3Config; +use clp_rust_utils::clp_config::package::DEFAULT_DATASET_NAME; +use clp_rust_utils::clp_config::package::config::ArchiveOutput; +use clp_rust_utils::job_config::ClpIoConfig; +use clp_rust_utils::job_config::CompressionJobId; +use clp_rust_utils::job_config::CompressionJobStatus; +use clp_rust_utils::job_config::InputConfig; +use clp_rust_utils::job_config::OutputConfig; +use clp_rust_utils::job_config::S3ObjectMetadataInputConfig; +use clp_rust_utils::job_config::ingestion::s3::BaseConfig; +use clp_rust_utils::s3::S3ObjectMetadataId; -use crate::{compression::BufferSubmitter, ingestion_job_manager::ClpCompressionState}; +use crate::compression::BufferSubmitter; +use crate::ingestion_job_manager::ClpCompressionState; /// The CLP compression job table name. pub const CLP_COMPRESSION_JOB_TABLE_NAME: &str = "compression_jobs"; diff --git a/components/log-ingestor/src/compression/listener.rs b/components/log-ingestor/src/compression/listener.rs index bd951ee4c7..ee9c61f379 100644 --- a/components/log-ingestor/src/compression/listener.rs +++ b/components/log-ingestor/src/compression/listener.rs @@ -1,14 +1,15 @@ use std::time::Duration; use anyhow::Result; -use tokio::{ - select, - sync::mpsc, - time::{Instant, sleep_until}, -}; +use tokio::select; +use tokio::sync::mpsc; +use tokio::time::Instant; +use tokio::time::sleep_until; use tokio_util::sync::CancellationToken; -use crate::compression::{Buffer, BufferSubmitter, CompressionBufferEntry}; +use crate::compression::Buffer; +use crate::compression::BufferSubmitter; +use crate::compression::CompressionBufferEntry; /// Represents a listener task that buffers incoming [`CompressionBufferEntry`] values and submits /// when a certain size threshold is reached or on timeout. diff --git a/components/log-ingestor/src/ingestion_job/s3_scanner.rs b/components/log-ingestor/src/ingestion_job/s3_scanner.rs index 80f650cdf6..b2609ad050 100644 --- a/components/log-ingestor/src/ingestion_job/s3_scanner.rs +++ b/components/log-ingestor/src/ingestion_job/s3_scanner.rs @@ -2,15 +2,17 @@ use std::time::Duration; use anyhow::Result; use aws_sdk_s3::Client; -use clp_rust_utils::{job_config::ingestion::s3::S3ScannerConfig, s3::ObjectMetadata}; +use clp_rust_utils::job_config::ingestion::s3::S3ScannerConfig; +use clp_rust_utils::s3::ObjectMetadata; use non_empty_string::NonEmptyString; use tokio::select; use tokio_util::sync::CancellationToken; -use crate::{ - aws_client_manager::AwsClientManagerType, - ingestion_job::{IngestionJobId, IngestionJobState, S3ScannerState, scan_prefix}, -}; +use crate::aws_client_manager::AwsClientManagerType; +use crate::ingestion_job::IngestionJobId; +use crate::ingestion_job::IngestionJobState; +use crate::ingestion_job::S3ScannerState; +use crate::ingestion_job::scan_prefix; /// Represents a S3 scanner task that periodically scans a given prefix under the bucket to fetch /// object metadata for newly created objects. diff --git a/components/log-ingestor/src/ingestion_job/sqs_listener.rs b/components/log-ingestor/src/ingestion_job/sqs_listener.rs index 659d543142..84be1e7ba9 100644 --- a/components/log-ingestor/src/ingestion_job/sqs_listener.rs +++ b/components/log-ingestor/src/ingestion_job/sqs_listener.rs @@ -1,24 +1,21 @@ use std::cmp::min; use anyhow::Result; -use aws_sdk_sqs::{ - Client, - operation::receive_message::ReceiveMessageOutput, - types::DeleteMessageBatchRequestEntry, -}; -use clp_rust_utils::{ - job_config::ingestion::s3::ValidatedSqsListenerConfig, - s3::ObjectMetadata, - sqs::event::{Record, S3}, -}; +use aws_sdk_sqs::Client; +use aws_sdk_sqs::operation::receive_message::ReceiveMessageOutput; +use aws_sdk_sqs::types::DeleteMessageBatchRequestEntry; +use clp_rust_utils::job_config::ingestion::s3::ValidatedSqsListenerConfig; +use clp_rust_utils::s3::ObjectMetadata; +use clp_rust_utils::sqs::event::Record; +use clp_rust_utils::sqs::event::S3; use non_empty_string::NonEmptyString; use tokio::select; use tokio_util::sync::CancellationToken; -use crate::{ - aws_client_manager::AwsClientManagerType, - ingestion_job::{IngestionJobId, IngestionJobState, SqsListenerState}, -}; +use crate::aws_client_manager::AwsClientManagerType; +use crate::ingestion_job::IngestionJobId; +use crate::ingestion_job::IngestionJobState; +use crate::ingestion_job::SqsListenerState; type TaskId = usize; diff --git a/components/log-ingestor/src/ingestion_job_manager.rs b/components/log-ingestor/src/ingestion_job_manager.rs index 00639daf2c..cdabf63a39 100644 --- a/components/log-ingestor/src/ingestion_job_manager.rs +++ b/components/log-ingestor/src/ingestion_job_manager.rs @@ -1,26 +1,26 @@ mod clp_ingestion; -use std::{collections::HashMap, sync::Arc}; +use std::collections::HashMap; +use std::sync::Arc; pub use clp_ingestion::*; -use clp_rust_utils::{ - clp_config::{ - AwsAuthentication, - package::{ - config::{Config as ClpConfig, LogsInput}, - credentials::Credentials as ClpCredentials, - }, - }, - job_config::ingestion::s3::{ConfigError, S3IngestionJobConfig, ValidatedSqsListenerConfig}, -}; +use clp_rust_utils::clp_config::AwsAuthentication; +use clp_rust_utils::clp_config::package::config::Config as ClpConfig; +use clp_rust_utils::clp_config::package::config::LogsInput; +use clp_rust_utils::clp_config::package::credentials::Credentials as ClpCredentials; +use clp_rust_utils::job_config::ingestion::s3::ConfigError; +use clp_rust_utils::job_config::ingestion::s3::S3IngestionJobConfig; +use clp_rust_utils::job_config::ingestion::s3::ValidatedSqsListenerConfig; use serde::Serialize; use tokio::sync::Mutex; use utoipa::ToSchema; -use crate::{ - aws_client_manager::{S3ClientWrapper, SqsClientWrapper}, - ingestion_job::{IngestionJob, IngestionJobId, IngestionJobState, S3Scanner}, -}; +use crate::aws_client_manager::S3ClientWrapper; +use crate::aws_client_manager::SqsClientWrapper; +use crate::ingestion_job::IngestionJob; +use crate::ingestion_job::IngestionJobId; +use crate::ingestion_job::IngestionJobState; +use crate::ingestion_job::S3Scanner; /// Errors for ingestion job manager operations. #[derive(thiserror::Error, Debug)] diff --git a/components/log-ingestor/src/ingestion_job_manager/clp_ingestion.rs b/components/log-ingestor/src/ingestion_job_manager/clp_ingestion.rs index 7b65b5ae71..dc8aadb32b 100644 --- a/components/log-ingestor/src/ingestion_job_manager/clp_ingestion.rs +++ b/components/log-ingestor/src/ingestion_job_manager/clp_ingestion.rs @@ -1,43 +1,42 @@ use std::time::Duration; use async_trait::async_trait; -use clp_rust_utils::{ - clp_config::{ - AwsAuthentication, - package::{ - config::{ArchiveOutput, Config as ClpConfig, LogsInput}, - credentials::Credentials as ClpCredentials, - }, - }, - database::mysql::MySqlEnumFormat, - impl_sqlx_type, - job_config::{ - ClpIoConfig, - CompressionJobId, - CompressionJobStatus, - InputConfig, - ingestion::s3::{S3IngestionJobConfig, S3ScannerConfig}, - }, - s3::{ObjectMetadata, S3ObjectMetadataId}, -}; +use clp_rust_utils::clp_config::AwsAuthentication; +use clp_rust_utils::clp_config::package::config::ArchiveOutput; +use clp_rust_utils::clp_config::package::config::Config as ClpConfig; +use clp_rust_utils::clp_config::package::config::LogsInput; +use clp_rust_utils::clp_config::package::credentials::Credentials as ClpCredentials; +use clp_rust_utils::database::mysql::MySqlEnumFormat; +use clp_rust_utils::impl_sqlx_type; +use clp_rust_utils::job_config::ClpIoConfig; +use clp_rust_utils::job_config::CompressionJobId; +use clp_rust_utils::job_config::CompressionJobStatus; +use clp_rust_utils::job_config::InputConfig; +use clp_rust_utils::job_config::ingestion::s3::S3IngestionJobConfig; +use clp_rust_utils::job_config::ingestion::s3::S3ScannerConfig; +use clp_rust_utils::s3::ObjectMetadata; +use clp_rust_utils::s3::S3ObjectMetadataId; use const_format::formatcp; use non_empty_string::NonEmptyString; -use sqlx::{Connection, MySqlPool}; -use strum_macros::{AsRefStr, Display, EnumIter, EnumString}; +use sqlx::Connection; +use sqlx::MySqlPool; +use strum_macros::AsRefStr; +use strum_macros::Display; +use strum_macros::EnumIter; +use strum_macros::EnumString; use tokio::sync::mpsc; -use crate::{ - compression::{ - Buffer, - CLP_COMPRESSION_JOB_TABLE_NAME, - CompressionBufferEntry, - CompressionJobSubmitter, - Listener, - wait_for_compression_job_completion_and_update_metadata, - }, - ingestion_job::{IngestionJobState, S3ScannerState, SqsListenerState}, - ingestion_job_manager::{IngestionJobId, TerminalStatus}, -}; +use crate::compression::Buffer; +use crate::compression::CLP_COMPRESSION_JOB_TABLE_NAME; +use crate::compression::CompressionBufferEntry; +use crate::compression::CompressionJobSubmitter; +use crate::compression::Listener; +use crate::compression::wait_for_compression_job_completion_and_update_metadata; +use crate::ingestion_job::IngestionJobState; +use crate::ingestion_job::S3ScannerState; +use crate::ingestion_job::SqsListenerState; +use crate::ingestion_job_manager::IngestionJobId; +use crate::ingestion_job_manager::TerminalStatus; /// A bundle of objects for log-ingestor to recovery from a restart. pub struct LogIngestorRecoveryContext { diff --git a/components/log-ingestor/src/routes.rs b/components/log-ingestor/src/routes.rs index 5293a78a68..c718aaca1a 100644 --- a/components/log-ingestor/src/routes.rs +++ b/components/log-ingestor/src/routes.rs @@ -1,30 +1,26 @@ #![allow(clippy::needless_for_each)] -use axum::{ - Json, - Router, - extract::{Path, State}, - response::IntoResponse, - routing::get, -}; -use clp_rust_utils::job_config::ingestion::s3::{ - S3IngestionJobConfig, - S3ScannerConfig, - SqsListenerConfig, -}; +use axum::Json; +use axum::Router; +use axum::extract::Path; +use axum::extract::State; +use axum::response::IntoResponse; +use axum::routing::get; +use clp_rust_utils::job_config::ingestion::s3::S3IngestionJobConfig; +use clp_rust_utils::job_config::ingestion::s3::S3ScannerConfig; +use clp_rust_utils::job_config::ingestion::s3::SqsListenerConfig; use serde::Serialize; -use tower_http::cors::{Any, CorsLayer}; -use utoipa::{OpenApi, ToSchema}; -use utoipa_axum::{router::OpenApiRouter, routes}; +use tower_http::cors::Any; +use tower_http::cors::CorsLayer; +use utoipa::OpenApi; +use utoipa::ToSchema; +use utoipa_axum::router::OpenApiRouter; +use utoipa_axum::routes; -use crate::{ - ingestion_job::IngestionJobId, - ingestion_job_manager::{ - Error as IngestionJobManagerError, - IngestionJobManagerState, - TerminalStatus, - }, -}; +use crate::ingestion_job::IngestionJobId; +use crate::ingestion_job_manager::Error as IngestionJobManagerError; +use crate::ingestion_job_manager::IngestionJobManagerState; +use crate::ingestion_job_manager::TerminalStatus; #[derive(utoipa::OpenApi)] #[openapi( diff --git a/components/log-ingestor/tests/aws_config.rs b/components/log-ingestor/tests/aws_config.rs index 04968c9dec..fed487829d 100644 --- a/components/log-ingestor/tests/aws_config.rs +++ b/components/log-ingestor/tests/aws_config.rs @@ -1,4 +1,5 @@ -use anyhow::{Result, anyhow}; +use anyhow::Result; +use anyhow::anyhow; use non_empty_string::NonEmptyString; /// Default AWS configuration for local testing with `LocalStack`. diff --git a/components/log-ingestor/tests/test_compression_listener.rs b/components/log-ingestor/tests/test_compression_listener.rs index 7d7b1f2319..0319a2ca53 100644 --- a/components/log-ingestor/tests/test_compression_listener.rs +++ b/components/log-ingestor/tests/test_compression_listener.rs @@ -1,19 +1,17 @@ -use std::{sync::Arc, time::Duration}; +use std::sync::Arc; +use std::time::Duration; use anyhow::Result; use async_trait::async_trait; use clp_rust_utils::s3::S3ObjectMetadataId; -use log_ingestor::compression::{ - Buffer, - BufferSubmitter, - CompressionBufferEntry, - DEFAULT_LISTENER_CAPACITY, - Listener, -}; -use tokio::{ - sync::{Mutex, mpsc}, - time::Instant, -}; +use log_ingestor::compression::Buffer; +use log_ingestor::compression::BufferSubmitter; +use log_ingestor::compression::CompressionBufferEntry; +use log_ingestor::compression::DEFAULT_LISTENER_CAPACITY; +use log_ingestor::compression::Listener; +use tokio::sync::Mutex; +use tokio::sync::mpsc; +use tokio::time::Instant; const TEST_OBJECT_SIZE: u64 = 1024; diff --git a/components/log-ingestor/tests/test_ingestion_job.rs b/components/log-ingestor/tests/test_ingestion_job.rs index fad9a8d513..f9ad1a327e 100644 --- a/components/log-ingestor/tests/test_ingestion_job.rs +++ b/components/log-ingestor/tests/test_ingestion_job.rs @@ -1,38 +1,34 @@ -use std::{sync::Arc, time::Duration}; +use std::sync::Arc; +use std::time::Duration; -use anyhow::{Context, Result}; +use anyhow::Context; +use anyhow::Result; use async_trait::async_trait; -use clp_rust_utils::{ - clp_config::{AwsAuthentication, AwsCredentials}, - job_config::ingestion::s3::{ - BaseConfig, - BufferConfig, - S3ScannerConfig, - SqsListenerConfig, - ValidatedSqsListenerConfig, - }, - s3::ObjectMetadata, - types::non_empty_string::ExpectedNonEmpty, -}; -use log_ingestor::{ - aws_client_manager::{S3ClientWrapper, SqsClientWrapper}, - ingestion_job::{ - IngestionJobId, - IngestionJobState, - S3Scanner, - S3ScannerState, - SqsListener, - SqsListenerState, - }, -}; +use clp_rust_utils::clp_config::AwsAuthentication; +use clp_rust_utils::clp_config::AwsCredentials; +use clp_rust_utils::job_config::ingestion::s3::BaseConfig; +use clp_rust_utils::job_config::ingestion::s3::BufferConfig; +use clp_rust_utils::job_config::ingestion::s3::S3ScannerConfig; +use clp_rust_utils::job_config::ingestion::s3::SqsListenerConfig; +use clp_rust_utils::job_config::ingestion::s3::ValidatedSqsListenerConfig; +use clp_rust_utils::s3::ObjectMetadata; +use clp_rust_utils::types::non_empty_string::ExpectedNonEmpty; +use log_ingestor::aws_client_manager::S3ClientWrapper; +use log_ingestor::aws_client_manager::SqsClientWrapper; +use log_ingestor::ingestion_job::IngestionJobId; +use log_ingestor::ingestion_job::IngestionJobState; +use log_ingestor::ingestion_job::S3Scanner; +use log_ingestor::ingestion_job::S3ScannerState; +use log_ingestor::ingestion_job::SqsListener; +use log_ingestor::ingestion_job::SqsListenerState; use non_empty_string::NonEmptyString; use tokio::sync::Mutex; use uuid::Uuid; -use super::{ - aws_config::AwsConfig, - test_utils::{self, get_testing_prefix_as_non_empty_string, upload_test_objects}, -}; +use super::aws_config::AwsConfig; +use super::test_utils::get_testing_prefix_as_non_empty_string; +use super::test_utils::upload_test_objects; +use super::test_utils::{self}; const WAIT_FOR_INGESTED_OBJECTS_TIMEOUT_SEC: u64 = 30; const INGESTED_OBJECT_POLL_INTERVAL_MS: u64 = 100; diff --git a/components/log-ingestor/tests/test_scan.rs b/components/log-ingestor/tests/test_scan.rs index 8f9c7076bb..39f5900b66 100644 --- a/components/log-ingestor/tests/test_scan.rs +++ b/components/log-ingestor/tests/test_scan.rs @@ -1,25 +1,20 @@ -use std::{ - collections::HashSet, - sync::{ - Arc, - atomic::{AtomicUsize, Ordering}, - }, -}; +use std::collections::HashSet; +use std::sync::Arc; +use std::sync::atomic::AtomicUsize; +use std::sync::atomic::Ordering; use anyhow::Result; -use clp_rust_utils::{ - clp_config::{AwsAuthentication, AwsCredentials}, - s3::ObjectMetadata, -}; +use clp_rust_utils::clp_config::AwsAuthentication; +use clp_rust_utils::clp_config::AwsCredentials; +use clp_rust_utils::s3::ObjectMetadata; use log_ingestor::ingestion_job::scan_prefix; use non_empty_string::NonEmptyString; use tokio::sync::Mutex; use uuid::Uuid; -use super::{ - aws_config::AwsConfig, - test_utils::{get_testing_prefix_as_non_empty_string, upload_test_objects}, -}; +use super::aws_config::AwsConfig; +use super::test_utils::get_testing_prefix_as_non_empty_string; +use super::test_utils::upload_test_objects; /// `ListObjectsV2` returns at most 1,000 keys per response. const MAX_NUM_OBJECTS_PER_PAGE: usize = 1000; diff --git a/components/log-ingestor/tests/test_utils.rs b/components/log-ingestor/tests/test_utils.rs index 4ef11fa358..f442e4212c 100644 --- a/components/log-ingestor/tests/test_utils.rs +++ b/components/log-ingestor/tests/test_utils.rs @@ -1,5 +1,7 @@ -use anyhow::{Context, Result}; -use clp_rust_utils::{s3::ObjectMetadata, types::non_empty_string::ExpectedNonEmpty}; +use anyhow::Context; +use anyhow::Result; +use clp_rust_utils::s3::ObjectMetadata; +use clp_rust_utils::types::non_empty_string::ExpectedNonEmpty; use log_ingestor::ingestion_job::IngestionJobId; use non_empty_string::NonEmptyString; From 5c8f5e521f9712cc314eb18368dbc9cdf4dad367 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Tue, 21 Jul 2026 23:43:57 -0400 Subject: [PATCH 3/6] More reformat --- components/api-server/src/bin/api_server.rs | 3 +- components/api-server/src/client.rs | 32 ++++---- components/api-server/src/error.rs | 3 +- components/api-server/src/routes.rs | 75 ++++++++++--------- .../src/clp_config/package/config.rs | 3 +- .../src/clp_config/s3_config.rs | 3 +- .../clp-rust-utils/src/database/mysql.rs | 6 +- .../src/job_config/clp_io_config.rs | 8 +- .../src/job_config/compression.rs | 6 +- .../src/job_config/ingestion.rs | 3 +- .../clp-rust-utils/src/job_config/search.rs | 6 +- components/clp-rust-utils/src/logging.rs | 14 ++-- components/clp-rust-utils/src/s3/client.rs | 8 +- components/clp-rust-utils/src/sqs/client.rs | 8 +- components/clp-rust-utils/src/telemetry.rs | 6 +- .../clp-rust-utils/tests/clp_config_test.rs | 15 ++-- .../log-ingestor/src/aws_client_manager.rs | 3 +- .../log-ingestor/src/bin/log_ingestor.rs | 6 +- .../compression/compression_job_submitter.rs | 32 ++++---- .../log-ingestor/src/compression/listener.rs | 13 ++-- .../src/ingestion_job/s3_scanner.rs | 12 ++- .../src/ingestion_job/sqs_listener.rs | 25 ++++--- .../log-ingestor/src/ingestion_job_manager.rs | 30 ++++---- .../ingestion_job_manager/clp_ingestion.rs | 65 ++++++++-------- components/log-ingestor/src/routes.rs | 42 ++++++----- components/log-ingestor/tests/aws_config.rs | 3 +- .../tests/test_compression_listener.rs | 22 +++--- .../log-ingestor/tests/test_ingestion_job.rs | 58 +++++++------- components/log-ingestor/tests/test_scan.rs | 25 ++++--- components/log-ingestor/tests/test_utils.rs | 6 +- 30 files changed, 273 insertions(+), 268 deletions(-) diff --git a/components/api-server/src/bin/api_server.rs b/components/api-server/src/bin/api_server.rs index 5f6603749a..8976dcffb0 100644 --- a/components/api-server/src/bin/api_server.rs +++ b/components/api-server/src/bin/api_server.rs @@ -1,7 +1,6 @@ use anyhow::Context; use clap::Parser; -use clp_rust_utils::clp_config::package; -use clp_rust_utils::serde::yaml; +use clp_rust_utils::{clp_config::package, serde::yaml}; #[derive(Parser)] #[command(version, about = "API Server for CLP.")] diff --git a/components/api-server/src/client.rs b/components/api-server/src/client.rs index f8e2a3e01c..01d3ab6c3b 100644 --- a/components/api-server/src/client.rs +++ b/components/api-server/src/client.rs @@ -1,28 +1,22 @@ use std::pin::Pin; use async_stream::stream; -use chrono::DateTime; -use chrono::TimeZone; -use chrono::Utc; -use clp_rust_utils::aws::AWS_DEFAULT_REGION; -use clp_rust_utils::clp_config::package::config::Config; -use clp_rust_utils::clp_config::package::config::StorageEngine; -use clp_rust_utils::clp_config::package::config::StreamOutputStorage; -use clp_rust_utils::clp_config::package::credentials::Credentials; -use clp_rust_utils::database::mysql::create_clp_db_mysql_pool; +use chrono::{DateTime, TimeZone, Utc}; pub use clp_rust_utils::job_config::CompressionJobStatus; -use clp_rust_utils::job_config::QUERY_JOBS_TABLE_NAME; -use clp_rust_utils::job_config::QueryJobStatus; -use clp_rust_utils::job_config::QueryJobType; -use clp_rust_utils::job_config::SearchJobConfig; -use futures::Stream; -use futures::StreamExt; +use clp_rust_utils::{ + aws::AWS_DEFAULT_REGION, + clp_config::package::{ + config::{Config, StorageEngine, StreamOutputStorage}, + credentials::Credentials, + }, + database::mysql::create_clp_db_mysql_pool, + job_config::{QUERY_JOBS_TABLE_NAME, QueryJobStatus, QueryJobType, SearchJobConfig}, +}; +use futures::{Stream, StreamExt}; use pin_project_lite::pin_project; -use serde::Deserialize; -use serde::Serialize; +use serde::{Deserialize, Serialize}; use sqlx::Row; -use utoipa::IntoParams; -use utoipa::ToSchema; +use utoipa::{IntoParams, ToSchema}; pub use crate::error::ClientError; diff --git a/components/api-server/src/error.rs b/components/api-server/src/error.rs index 3447f3d616..484b462b45 100644 --- a/components/api-server/src/error.rs +++ b/components/api-server/src/error.rs @@ -1,5 +1,4 @@ -use aws_sdk_s3::error::SdkError; -use aws_sdk_s3::primitives::ByteStreamError; +use aws_sdk_s3::{error::SdkError, primitives::ByteStreamError}; use num_enum::TryFromPrimitive; use thiserror::Error; diff --git a/components/api-server/src/routes.rs b/components/api-server/src/routes.rs index fae107f48a..c10ac175da 100644 --- a/components/api-server/src/routes.rs +++ b/components/api-server/src/routes.rs @@ -1,31 +1,29 @@ -use axum::Json; -use axum::extract::Path; -use axum::extract::Query; -use axum::extract::State; -use axum::http::StatusCode; -use axum::response::IntoResponse; -use axum::response::Sse; -use axum::response::sse::Event; -use axum::response::sse::KeepAlive; -use axum::routing::get; -use futures::Stream; -use futures::StreamExt; -use serde::Deserialize; -use serde::Serialize; +use axum::{ + Json, + extract::{Path, Query, State}, + http::StatusCode, + response::{ + IntoResponse, + Sse, + sse::{Event, KeepAlive}, + }, + routing::get, +}; +use futures::{Stream, StreamExt}; +use serde::{Deserialize, Serialize}; use thiserror::Error; -use tower_http::cors::Any; -use tower_http::cors::CorsLayer; -use utoipa::OpenApi; -use utoipa::ToSchema; -use utoipa_axum::router::OpenApiRouter; -use utoipa_axum::routes; - -use crate::client::Client; -use crate::client::ClientError; -use crate::client::CompressionUsage; -use crate::client::CompressionUsageParams; -use crate::client::QueryConfig; -use crate::client::ValidatedCompressionUsageParams; +use tower_http::cors::{Any, CorsLayer}; +use utoipa::{OpenApi, ToSchema}; +use utoipa_axum::{router::OpenApiRouter, routes}; + +use crate::client::{ + Client, + ClientError, + CompressionUsage, + CompressionUsageParams, + QueryConfig, + ValidatedCompressionUsageParams, +}; /// Factory method to create an Axum router configured with all API routes. /// @@ -67,12 +65,14 @@ pub fn from_client(client: Client) -> Result { mod api_doc { // Using `super::...` can cause `super` to appear as a tag in the generated OpenAPI // documentation. Importing the paths directly prevents this issue. - use super::__path_cancel_query; - use super::__path_compression_usage; - use super::__path_health; - use super::__path_query; - use super::__path_query_results; - use super::CompressionUsage; + use super::{ + __path_cancel_query, + __path_compression_usage, + __path_health, + __path_query, + __path_query_results, + CompressionUsage, + }; use crate::client::CompressionJobStatus; #[derive(utoipa::OpenApi)] @@ -337,10 +337,11 @@ impl IntoResponse for HandlerError { #[cfg(test)] mod tests { - use axum::body::Body; - use axum::http::Request; - use axum::http::StatusCode; - use axum::routing::get; + use axum::{ + body::Body, + http::{Request, StatusCode}, + routing::get, + }; use http_body_util::BodyExt; use tower::ServiceExt; 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 80645c131e..29863686b0 100644 --- a/components/clp-rust-utils/src/clp_config/package/config.rs +++ b/components/clp-rust-utils/src/clp_config/package/config.rs @@ -1,7 +1,6 @@ use serde::Deserialize; -use crate::clp_config::AwsAuthentication; -use crate::clp_config::S3Config; +use crate::clp_config::{AwsAuthentication, S3Config}; /// Mirror of `clp_py_utils.clp_config.ClpConfig`. /// diff --git a/components/clp-rust-utils/src/clp_config/s3_config.rs b/components/clp-rust-utils/src/clp_config/s3_config.rs index 50792094a1..0ceac30f91 100644 --- a/components/clp-rust-utils/src/clp_config/s3_config.rs +++ b/components/clp-rust-utils/src/clp_config/s3_config.rs @@ -1,6 +1,5 @@ use non_empty_string::NonEmptyString; -use serde::Deserialize; -use serde::Serialize; +use serde::{Deserialize, Serialize}; /// Represents the configuration for connecting to an S3 bucket. #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] diff --git a/components/clp-rust-utils/src/database/mysql.rs b/components/clp-rust-utils/src/database/mysql.rs index 9a2aa127e3..1ea951184a 100644 --- a/components/clp-rust-utils/src/database/mysql.rs +++ b/components/clp-rust-utils/src/database/mysql.rs @@ -1,8 +1,10 @@ use secrecy::ExposeSecret; use strum::IntoEnumIterator; -use crate::clp_config::package::config::Database as DatabaseConfig; -use crate::clp_config::package::credentials::Database as DatabaseCredentials; +use crate::clp_config::package::{ + config::Database as DatabaseConfig, + credentials::Database as DatabaseCredentials, +}; /// Implements [`sqlx::Type`] for `$ty` by delegating to `$delegate`. /// diff --git a/components/clp-rust-utils/src/job_config/clp_io_config.rs b/components/clp-rust-utils/src/job_config/clp_io_config.rs index 312dac3e53..7f4e19cb51 100644 --- a/components/clp-rust-utils/src/job_config/clp_io_config.rs +++ b/components/clp-rust-utils/src/job_config/clp_io_config.rs @@ -1,9 +1,11 @@ use non_empty_string::NonEmptyString; use serde::Serialize; -use crate::clp_config::S3Config; -use crate::job_config::ingestion::JobId as IngestionJobId; -use crate::s3::S3ObjectMetadataId; +use crate::{ + clp_config::S3Config, + job_config::ingestion::JobId as IngestionJobId, + s3::S3ObjectMetadataId, +}; /// Represents CLP IO config. #[derive(Debug, Clone, PartialEq, Eq, Serialize)] diff --git a/components/clp-rust-utils/src/job_config/compression.rs b/components/clp-rust-utils/src/job_config/compression.rs index e09ec10484..ddebf2f880 100644 --- a/components/clp-rust-utils/src/job_config/compression.rs +++ b/components/clp-rust-utils/src/job_config/compression.rs @@ -1,7 +1,5 @@ -use num_enum::IntoPrimitive; -use num_enum::TryFromPrimitive; -use serde::Deserialize; -use serde::Serialize; +use num_enum::{IntoPrimitive, TryFromPrimitive}; +use serde::{Deserialize, Serialize}; use strum::EnumString; use utoipa::ToSchema; diff --git a/components/clp-rust-utils/src/job_config/ingestion.rs b/components/clp-rust-utils/src/job_config/ingestion.rs index 43782e520d..bc3f726ac2 100644 --- a/components/clp-rust-utils/src/job_config/ingestion.rs +++ b/components/clp-rust-utils/src/job_config/ingestion.rs @@ -1,7 +1,6 @@ pub mod s3 { use non_empty_string::NonEmptyString; - use serde::Deserialize; - use serde::Serialize; + use serde::{Deserialize, Serialize}; use thiserror::Error; use utoipa::ToSchema; diff --git a/components/clp-rust-utils/src/job_config/search.rs b/components/clp-rust-utils/src/job_config/search.rs index 0ece3c2295..fcf25275b9 100644 --- a/components/clp-rust-utils/src/job_config/search.rs +++ b/components/clp-rust-utils/src/job_config/search.rs @@ -1,7 +1,5 @@ -use num_enum::IntoPrimitive; -use num_enum::TryFromPrimitive; -use serde::Deserialize; -use serde::Serialize; +use num_enum::{IntoPrimitive, TryFromPrimitive}; +use serde::{Deserialize, Serialize}; pub const QUERY_JOBS_TABLE_NAME: &str = "query_jobs"; diff --git a/components/clp-rust-utils/src/logging.rs b/components/clp-rust-utils/src/logging.rs index bc130b08eb..ad3ef343b4 100644 --- a/components/clp-rust-utils/src/logging.rs +++ b/components/clp-rust-utils/src/logging.rs @@ -1,9 +1,11 @@ -use tracing_appender::non_blocking::NonBlockingBuilder; -use tracing_appender::non_blocking::WorkerGuard; -use tracing_appender::rolling::RollingFileAppender; -use tracing_appender::rolling::Rotation; -use tracing_subscriber::fmt::writer::MakeWriterExt; -use tracing_subscriber::{self}; +use tracing_appender::{ + non_blocking::{NonBlockingBuilder, WorkerGuard}, + rolling::{RollingFileAppender, Rotation}, +}; +use tracing_subscriber::{ + fmt::writer::MakeWriterExt, + {self}, +}; /// Opaque struct to hold the worker guards for the background log writers. /// These guards must be held for the lifetime of the program to ensure logs are flushed. diff --git a/components/clp-rust-utils/src/s3/client.rs b/components/clp-rust-utils/src/s3/client.rs index 134b1645a9..9160f5d0f1 100644 --- a/components/clp-rust-utils/src/s3/client.rs +++ b/components/clp-rust-utils/src/s3/client.rs @@ -1,8 +1,8 @@ use aws_config::BehaviorVersion; -use aws_sdk_s3::Client; -use aws_sdk_s3::config::Builder; -use aws_sdk_s3::config::Credentials; -use aws_sdk_s3::config::Region; +use aws_sdk_s3::{ + Client, + config::{Builder, Credentials, Region}, +}; use non_empty_string::NonEmptyString; use crate::clp_config::AwsAuthentication; diff --git a/components/clp-rust-utils/src/sqs/client.rs b/components/clp-rust-utils/src/sqs/client.rs index 2d5be46171..4586837146 100644 --- a/components/clp-rust-utils/src/sqs/client.rs +++ b/components/clp-rust-utils/src/sqs/client.rs @@ -1,8 +1,8 @@ use aws_config::BehaviorVersion; -use aws_sdk_sqs::Client; -use aws_sdk_sqs::config::Builder; -use aws_sdk_sqs::config::Credentials; -use aws_sdk_sqs::config::Region; +use aws_sdk_sqs::{ + Client, + config::{Builder, Credentials, Region}, +}; use non_empty_string::NonEmptyString; use crate::clp_config::AwsAuthentication; diff --git a/components/clp-rust-utils/src/telemetry.rs b/components/clp-rust-utils/src/telemetry.rs index ee237d12a1..45d5ba5d32 100644 --- a/components/clp-rust-utils/src/telemetry.rs +++ b/components/clp-rust-utils/src/telemetry.rs @@ -2,11 +2,9 @@ use std::env; -use opentelemetry_sdk::metrics::PeriodicReader; -use opentelemetry_sdk::metrics::SdkMeterProvider; +use opentelemetry_sdk::metrics::{PeriodicReader, SdkMeterProvider}; -use crate::Error; -use crate::clp_config::package::config::Telemetry; +use crate::{Error, clp_config::package::config::Telemetry}; /// RAII guard that shuts down the meter provider and flushes pending metric exports when dropped. pub struct TelemetryGuard { diff --git a/components/clp-rust-utils/tests/clp_config_test.rs b/components/clp-rust-utils/tests/clp_config_test.rs index adbfab2b21..41fd5edf11 100644 --- a/components/clp-rust-utils/tests/clp_config_test.rs +++ b/components/clp-rust-utils/tests/clp_config_test.rs @@ -1,12 +1,9 @@ -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::job_config::ClpIoConfig; -use clp_rust_utils::job_config::InputConfig; -use clp_rust_utils::job_config::OutputConfig; -use clp_rust_utils::job_config::S3ObjectMetadataInputConfig; -use clp_rust_utils::serde::BrotliMsgpack; -use clp_rust_utils::types::non_empty_string::ExpectedNonEmpty; +use clp_rust_utils::{ + clp_config::{AwsAuthentication, AwsCredentials, S3Config}, + job_config::{ClpIoConfig, InputConfig, OutputConfig, S3ObjectMetadataInputConfig}, + serde::BrotliMsgpack, + types::non_empty_string::ExpectedNonEmpty, +}; use non_empty_string::NonEmptyString; use serde_json::Value; diff --git a/components/log-ingestor/src/aws_client_manager.rs b/components/log-ingestor/src/aws_client_manager.rs index fd14131f26..d19e420f7a 100644 --- a/components/log-ingestor/src/aws_client_manager.rs +++ b/components/log-ingestor/src/aws_client_manager.rs @@ -2,8 +2,7 @@ use anyhow::Result; use async_trait::async_trait; use aws_sdk_s3::Client as S3Client; use aws_sdk_sqs::Client as SqsClient; -use clp_rust_utils::aws::AWS_DEFAULT_REGION; -use clp_rust_utils::clp_config::AwsAuthentication; +use clp_rust_utils::{aws::AWS_DEFAULT_REGION, clp_config::AwsAuthentication}; use non_empty_string::NonEmptyString; /// A marker trait for AWS client types. diff --git a/components/log-ingestor/src/bin/log_ingestor.rs b/components/log-ingestor/src/bin/log_ingestor.rs index c7436b39ac..a766c84156 100644 --- a/components/log-ingestor/src/bin/log_ingestor.rs +++ b/components/log-ingestor/src/bin/log_ingestor.rs @@ -1,9 +1,7 @@ use anyhow::Context; use clap::Parser; -use clp_rust_utils::clp_config::package; -use clp_rust_utils::serde::yaml; -use log_ingestor::ingestion_job_manager::IngestionJobManagerState; -use log_ingestor::routes::create_router; +use clp_rust_utils::{clp_config::package, serde::yaml}; +use log_ingestor::{ingestion_job_manager::IngestionJobManagerState, routes::create_router}; #[derive(Parser)] #[command(version, about = "log-ingestor for CLP.")] diff --git a/components/log-ingestor/src/compression/compression_job_submitter.rs b/components/log-ingestor/src/compression/compression_job_submitter.rs index 69beadb97c..b4f8f6287c 100644 --- a/components/log-ingestor/src/compression/compression_job_submitter.rs +++ b/components/log-ingestor/src/compression/compression_job_submitter.rs @@ -1,20 +1,24 @@ use anyhow::Result; use async_trait::async_trait; -use clp_rust_utils::clp_config::AwsAuthentication; -use clp_rust_utils::clp_config::S3Config; -use clp_rust_utils::clp_config::package::DEFAULT_DATASET_NAME; -use clp_rust_utils::clp_config::package::config::ArchiveOutput; -use clp_rust_utils::job_config::ClpIoConfig; -use clp_rust_utils::job_config::CompressionJobId; -use clp_rust_utils::job_config::CompressionJobStatus; -use clp_rust_utils::job_config::InputConfig; -use clp_rust_utils::job_config::OutputConfig; -use clp_rust_utils::job_config::S3ObjectMetadataInputConfig; -use clp_rust_utils::job_config::ingestion::s3::BaseConfig; -use clp_rust_utils::s3::S3ObjectMetadataId; +use clp_rust_utils::{ + clp_config::{ + AwsAuthentication, + S3Config, + package::{DEFAULT_DATASET_NAME, config::ArchiveOutput}, + }, + job_config::{ + ClpIoConfig, + CompressionJobId, + CompressionJobStatus, + InputConfig, + OutputConfig, + S3ObjectMetadataInputConfig, + ingestion::s3::BaseConfig, + }, + s3::S3ObjectMetadataId, +}; -use crate::compression::BufferSubmitter; -use crate::ingestion_job_manager::ClpCompressionState; +use crate::{compression::BufferSubmitter, ingestion_job_manager::ClpCompressionState}; /// The CLP compression job table name. pub const CLP_COMPRESSION_JOB_TABLE_NAME: &str = "compression_jobs"; diff --git a/components/log-ingestor/src/compression/listener.rs b/components/log-ingestor/src/compression/listener.rs index ee9c61f379..bd951ee4c7 100644 --- a/components/log-ingestor/src/compression/listener.rs +++ b/components/log-ingestor/src/compression/listener.rs @@ -1,15 +1,14 @@ use std::time::Duration; use anyhow::Result; -use tokio::select; -use tokio::sync::mpsc; -use tokio::time::Instant; -use tokio::time::sleep_until; +use tokio::{ + select, + sync::mpsc, + time::{Instant, sleep_until}, +}; use tokio_util::sync::CancellationToken; -use crate::compression::Buffer; -use crate::compression::BufferSubmitter; -use crate::compression::CompressionBufferEntry; +use crate::compression::{Buffer, BufferSubmitter, CompressionBufferEntry}; /// Represents a listener task that buffers incoming [`CompressionBufferEntry`] values and submits /// when a certain size threshold is reached or on timeout. diff --git a/components/log-ingestor/src/ingestion_job/s3_scanner.rs b/components/log-ingestor/src/ingestion_job/s3_scanner.rs index b2609ad050..80f650cdf6 100644 --- a/components/log-ingestor/src/ingestion_job/s3_scanner.rs +++ b/components/log-ingestor/src/ingestion_job/s3_scanner.rs @@ -2,17 +2,15 @@ use std::time::Duration; use anyhow::Result; use aws_sdk_s3::Client; -use clp_rust_utils::job_config::ingestion::s3::S3ScannerConfig; -use clp_rust_utils::s3::ObjectMetadata; +use clp_rust_utils::{job_config::ingestion::s3::S3ScannerConfig, s3::ObjectMetadata}; use non_empty_string::NonEmptyString; use tokio::select; use tokio_util::sync::CancellationToken; -use crate::aws_client_manager::AwsClientManagerType; -use crate::ingestion_job::IngestionJobId; -use crate::ingestion_job::IngestionJobState; -use crate::ingestion_job::S3ScannerState; -use crate::ingestion_job::scan_prefix; +use crate::{ + aws_client_manager::AwsClientManagerType, + ingestion_job::{IngestionJobId, IngestionJobState, S3ScannerState, scan_prefix}, +}; /// Represents a S3 scanner task that periodically scans a given prefix under the bucket to fetch /// object metadata for newly created objects. diff --git a/components/log-ingestor/src/ingestion_job/sqs_listener.rs b/components/log-ingestor/src/ingestion_job/sqs_listener.rs index 84be1e7ba9..659d543142 100644 --- a/components/log-ingestor/src/ingestion_job/sqs_listener.rs +++ b/components/log-ingestor/src/ingestion_job/sqs_listener.rs @@ -1,21 +1,24 @@ use std::cmp::min; use anyhow::Result; -use aws_sdk_sqs::Client; -use aws_sdk_sqs::operation::receive_message::ReceiveMessageOutput; -use aws_sdk_sqs::types::DeleteMessageBatchRequestEntry; -use clp_rust_utils::job_config::ingestion::s3::ValidatedSqsListenerConfig; -use clp_rust_utils::s3::ObjectMetadata; -use clp_rust_utils::sqs::event::Record; -use clp_rust_utils::sqs::event::S3; +use aws_sdk_sqs::{ + Client, + operation::receive_message::ReceiveMessageOutput, + types::DeleteMessageBatchRequestEntry, +}; +use clp_rust_utils::{ + job_config::ingestion::s3::ValidatedSqsListenerConfig, + s3::ObjectMetadata, + sqs::event::{Record, S3}, +}; use non_empty_string::NonEmptyString; use tokio::select; use tokio_util::sync::CancellationToken; -use crate::aws_client_manager::AwsClientManagerType; -use crate::ingestion_job::IngestionJobId; -use crate::ingestion_job::IngestionJobState; -use crate::ingestion_job::SqsListenerState; +use crate::{ + aws_client_manager::AwsClientManagerType, + ingestion_job::{IngestionJobId, IngestionJobState, SqsListenerState}, +}; type TaskId = usize; diff --git a/components/log-ingestor/src/ingestion_job_manager.rs b/components/log-ingestor/src/ingestion_job_manager.rs index cdabf63a39..00639daf2c 100644 --- a/components/log-ingestor/src/ingestion_job_manager.rs +++ b/components/log-ingestor/src/ingestion_job_manager.rs @@ -1,26 +1,26 @@ mod clp_ingestion; -use std::collections::HashMap; -use std::sync::Arc; +use std::{collections::HashMap, sync::Arc}; pub use clp_ingestion::*; -use clp_rust_utils::clp_config::AwsAuthentication; -use clp_rust_utils::clp_config::package::config::Config as ClpConfig; -use clp_rust_utils::clp_config::package::config::LogsInput; -use clp_rust_utils::clp_config::package::credentials::Credentials as ClpCredentials; -use clp_rust_utils::job_config::ingestion::s3::ConfigError; -use clp_rust_utils::job_config::ingestion::s3::S3IngestionJobConfig; -use clp_rust_utils::job_config::ingestion::s3::ValidatedSqsListenerConfig; +use clp_rust_utils::{ + clp_config::{ + AwsAuthentication, + package::{ + config::{Config as ClpConfig, LogsInput}, + credentials::Credentials as ClpCredentials, + }, + }, + job_config::ingestion::s3::{ConfigError, S3IngestionJobConfig, ValidatedSqsListenerConfig}, +}; use serde::Serialize; use tokio::sync::Mutex; use utoipa::ToSchema; -use crate::aws_client_manager::S3ClientWrapper; -use crate::aws_client_manager::SqsClientWrapper; -use crate::ingestion_job::IngestionJob; -use crate::ingestion_job::IngestionJobId; -use crate::ingestion_job::IngestionJobState; -use crate::ingestion_job::S3Scanner; +use crate::{ + aws_client_manager::{S3ClientWrapper, SqsClientWrapper}, + ingestion_job::{IngestionJob, IngestionJobId, IngestionJobState, S3Scanner}, +}; /// Errors for ingestion job manager operations. #[derive(thiserror::Error, Debug)] diff --git a/components/log-ingestor/src/ingestion_job_manager/clp_ingestion.rs b/components/log-ingestor/src/ingestion_job_manager/clp_ingestion.rs index dc8aadb32b..7b65b5ae71 100644 --- a/components/log-ingestor/src/ingestion_job_manager/clp_ingestion.rs +++ b/components/log-ingestor/src/ingestion_job_manager/clp_ingestion.rs @@ -1,42 +1,43 @@ use std::time::Duration; use async_trait::async_trait; -use clp_rust_utils::clp_config::AwsAuthentication; -use clp_rust_utils::clp_config::package::config::ArchiveOutput; -use clp_rust_utils::clp_config::package::config::Config as ClpConfig; -use clp_rust_utils::clp_config::package::config::LogsInput; -use clp_rust_utils::clp_config::package::credentials::Credentials as ClpCredentials; -use clp_rust_utils::database::mysql::MySqlEnumFormat; -use clp_rust_utils::impl_sqlx_type; -use clp_rust_utils::job_config::ClpIoConfig; -use clp_rust_utils::job_config::CompressionJobId; -use clp_rust_utils::job_config::CompressionJobStatus; -use clp_rust_utils::job_config::InputConfig; -use clp_rust_utils::job_config::ingestion::s3::S3IngestionJobConfig; -use clp_rust_utils::job_config::ingestion::s3::S3ScannerConfig; -use clp_rust_utils::s3::ObjectMetadata; -use clp_rust_utils::s3::S3ObjectMetadataId; +use clp_rust_utils::{ + clp_config::{ + AwsAuthentication, + package::{ + config::{ArchiveOutput, Config as ClpConfig, LogsInput}, + credentials::Credentials as ClpCredentials, + }, + }, + database::mysql::MySqlEnumFormat, + impl_sqlx_type, + job_config::{ + ClpIoConfig, + CompressionJobId, + CompressionJobStatus, + InputConfig, + ingestion::s3::{S3IngestionJobConfig, S3ScannerConfig}, + }, + s3::{ObjectMetadata, S3ObjectMetadataId}, +}; use const_format::formatcp; use non_empty_string::NonEmptyString; -use sqlx::Connection; -use sqlx::MySqlPool; -use strum_macros::AsRefStr; -use strum_macros::Display; -use strum_macros::EnumIter; -use strum_macros::EnumString; +use sqlx::{Connection, MySqlPool}; +use strum_macros::{AsRefStr, Display, EnumIter, EnumString}; use tokio::sync::mpsc; -use crate::compression::Buffer; -use crate::compression::CLP_COMPRESSION_JOB_TABLE_NAME; -use crate::compression::CompressionBufferEntry; -use crate::compression::CompressionJobSubmitter; -use crate::compression::Listener; -use crate::compression::wait_for_compression_job_completion_and_update_metadata; -use crate::ingestion_job::IngestionJobState; -use crate::ingestion_job::S3ScannerState; -use crate::ingestion_job::SqsListenerState; -use crate::ingestion_job_manager::IngestionJobId; -use crate::ingestion_job_manager::TerminalStatus; +use crate::{ + compression::{ + Buffer, + CLP_COMPRESSION_JOB_TABLE_NAME, + CompressionBufferEntry, + CompressionJobSubmitter, + Listener, + wait_for_compression_job_completion_and_update_metadata, + }, + ingestion_job::{IngestionJobState, S3ScannerState, SqsListenerState}, + ingestion_job_manager::{IngestionJobId, TerminalStatus}, +}; /// A bundle of objects for log-ingestor to recovery from a restart. pub struct LogIngestorRecoveryContext { diff --git a/components/log-ingestor/src/routes.rs b/components/log-ingestor/src/routes.rs index c718aaca1a..5293a78a68 100644 --- a/components/log-ingestor/src/routes.rs +++ b/components/log-ingestor/src/routes.rs @@ -1,26 +1,30 @@ #![allow(clippy::needless_for_each)] -use axum::Json; -use axum::Router; -use axum::extract::Path; -use axum::extract::State; -use axum::response::IntoResponse; -use axum::routing::get; -use clp_rust_utils::job_config::ingestion::s3::S3IngestionJobConfig; -use clp_rust_utils::job_config::ingestion::s3::S3ScannerConfig; -use clp_rust_utils::job_config::ingestion::s3::SqsListenerConfig; +use axum::{ + Json, + Router, + extract::{Path, State}, + response::IntoResponse, + routing::get, +}; +use clp_rust_utils::job_config::ingestion::s3::{ + S3IngestionJobConfig, + S3ScannerConfig, + SqsListenerConfig, +}; use serde::Serialize; -use tower_http::cors::Any; -use tower_http::cors::CorsLayer; -use utoipa::OpenApi; -use utoipa::ToSchema; -use utoipa_axum::router::OpenApiRouter; -use utoipa_axum::routes; +use tower_http::cors::{Any, CorsLayer}; +use utoipa::{OpenApi, ToSchema}; +use utoipa_axum::{router::OpenApiRouter, routes}; -use crate::ingestion_job::IngestionJobId; -use crate::ingestion_job_manager::Error as IngestionJobManagerError; -use crate::ingestion_job_manager::IngestionJobManagerState; -use crate::ingestion_job_manager::TerminalStatus; +use crate::{ + ingestion_job::IngestionJobId, + ingestion_job_manager::{ + Error as IngestionJobManagerError, + IngestionJobManagerState, + TerminalStatus, + }, +}; #[derive(utoipa::OpenApi)] #[openapi( diff --git a/components/log-ingestor/tests/aws_config.rs b/components/log-ingestor/tests/aws_config.rs index fed487829d..04968c9dec 100644 --- a/components/log-ingestor/tests/aws_config.rs +++ b/components/log-ingestor/tests/aws_config.rs @@ -1,5 +1,4 @@ -use anyhow::Result; -use anyhow::anyhow; +use anyhow::{Result, anyhow}; use non_empty_string::NonEmptyString; /// Default AWS configuration for local testing with `LocalStack`. diff --git a/components/log-ingestor/tests/test_compression_listener.rs b/components/log-ingestor/tests/test_compression_listener.rs index 0319a2ca53..7d7b1f2319 100644 --- a/components/log-ingestor/tests/test_compression_listener.rs +++ b/components/log-ingestor/tests/test_compression_listener.rs @@ -1,17 +1,19 @@ -use std::sync::Arc; -use std::time::Duration; +use std::{sync::Arc, time::Duration}; use anyhow::Result; use async_trait::async_trait; use clp_rust_utils::s3::S3ObjectMetadataId; -use log_ingestor::compression::Buffer; -use log_ingestor::compression::BufferSubmitter; -use log_ingestor::compression::CompressionBufferEntry; -use log_ingestor::compression::DEFAULT_LISTENER_CAPACITY; -use log_ingestor::compression::Listener; -use tokio::sync::Mutex; -use tokio::sync::mpsc; -use tokio::time::Instant; +use log_ingestor::compression::{ + Buffer, + BufferSubmitter, + CompressionBufferEntry, + DEFAULT_LISTENER_CAPACITY, + Listener, +}; +use tokio::{ + sync::{Mutex, mpsc}, + time::Instant, +}; const TEST_OBJECT_SIZE: u64 = 1024; diff --git a/components/log-ingestor/tests/test_ingestion_job.rs b/components/log-ingestor/tests/test_ingestion_job.rs index f9ad1a327e..3d8f125cf2 100644 --- a/components/log-ingestor/tests/test_ingestion_job.rs +++ b/components/log-ingestor/tests/test_ingestion_job.rs @@ -1,34 +1,42 @@ -use std::sync::Arc; -use std::time::Duration; +use std::{sync::Arc, time::Duration}; -use anyhow::Context; -use anyhow::Result; +use anyhow::{Context, Result}; use async_trait::async_trait; -use clp_rust_utils::clp_config::AwsAuthentication; -use clp_rust_utils::clp_config::AwsCredentials; -use clp_rust_utils::job_config::ingestion::s3::BaseConfig; -use clp_rust_utils::job_config::ingestion::s3::BufferConfig; -use clp_rust_utils::job_config::ingestion::s3::S3ScannerConfig; -use clp_rust_utils::job_config::ingestion::s3::SqsListenerConfig; -use clp_rust_utils::job_config::ingestion::s3::ValidatedSqsListenerConfig; -use clp_rust_utils::s3::ObjectMetadata; -use clp_rust_utils::types::non_empty_string::ExpectedNonEmpty; -use log_ingestor::aws_client_manager::S3ClientWrapper; -use log_ingestor::aws_client_manager::SqsClientWrapper; -use log_ingestor::ingestion_job::IngestionJobId; -use log_ingestor::ingestion_job::IngestionJobState; -use log_ingestor::ingestion_job::S3Scanner; -use log_ingestor::ingestion_job::S3ScannerState; -use log_ingestor::ingestion_job::SqsListener; -use log_ingestor::ingestion_job::SqsListenerState; +use clp_rust_utils::{ + clp_config::{AwsAuthentication, AwsCredentials}, + job_config::ingestion::s3::{ + BaseConfig, + BufferConfig, + S3ScannerConfig, + SqsListenerConfig, + ValidatedSqsListenerConfig, + }, + s3::ObjectMetadata, + types::non_empty_string::ExpectedNonEmpty, +}; +use log_ingestor::{ + aws_client_manager::{S3ClientWrapper, SqsClientWrapper}, + ingestion_job::{ + IngestionJobId, + IngestionJobState, + S3Scanner, + S3ScannerState, + SqsListener, + SqsListenerState, + }, +}; use non_empty_string::NonEmptyString; use tokio::sync::Mutex; use uuid::Uuid; -use super::aws_config::AwsConfig; -use super::test_utils::get_testing_prefix_as_non_empty_string; -use super::test_utils::upload_test_objects; -use super::test_utils::{self}; +use super::{ + aws_config::AwsConfig, + test_utils::{ + get_testing_prefix_as_non_empty_string, + upload_test_objects, + {self}, + }, +}; const WAIT_FOR_INGESTED_OBJECTS_TIMEOUT_SEC: u64 = 30; const INGESTED_OBJECT_POLL_INTERVAL_MS: u64 = 100; diff --git a/components/log-ingestor/tests/test_scan.rs b/components/log-ingestor/tests/test_scan.rs index 39f5900b66..8f9c7076bb 100644 --- a/components/log-ingestor/tests/test_scan.rs +++ b/components/log-ingestor/tests/test_scan.rs @@ -1,20 +1,25 @@ -use std::collections::HashSet; -use std::sync::Arc; -use std::sync::atomic::AtomicUsize; -use std::sync::atomic::Ordering; +use std::{ + collections::HashSet, + sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, + }, +}; use anyhow::Result; -use clp_rust_utils::clp_config::AwsAuthentication; -use clp_rust_utils::clp_config::AwsCredentials; -use clp_rust_utils::s3::ObjectMetadata; +use clp_rust_utils::{ + clp_config::{AwsAuthentication, AwsCredentials}, + s3::ObjectMetadata, +}; use log_ingestor::ingestion_job::scan_prefix; use non_empty_string::NonEmptyString; use tokio::sync::Mutex; use uuid::Uuid; -use super::aws_config::AwsConfig; -use super::test_utils::get_testing_prefix_as_non_empty_string; -use super::test_utils::upload_test_objects; +use super::{ + aws_config::AwsConfig, + test_utils::{get_testing_prefix_as_non_empty_string, upload_test_objects}, +}; /// `ListObjectsV2` returns at most 1,000 keys per response. const MAX_NUM_OBJECTS_PER_PAGE: usize = 1000; diff --git a/components/log-ingestor/tests/test_utils.rs b/components/log-ingestor/tests/test_utils.rs index f442e4212c..4ef11fa358 100644 --- a/components/log-ingestor/tests/test_utils.rs +++ b/components/log-ingestor/tests/test_utils.rs @@ -1,7 +1,5 @@ -use anyhow::Context; -use anyhow::Result; -use clp_rust_utils::s3::ObjectMetadata; -use clp_rust_utils::types::non_empty_string::ExpectedNonEmpty; +use anyhow::{Context, Result}; +use clp_rust_utils::{s3::ObjectMetadata, types::non_empty_string::ExpectedNonEmpty}; use log_ingestor::ingestion_job::IngestionJobId; use non_empty_string::NonEmptyString; From 0477b3160a87f4b211591c2bcc7561b1cc506997 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Tue, 21 Jul 2026 23:50:14 -0400 Subject: [PATCH 4/6] reformat --- components/api-server/src/bin/api_server.rs | 3 +- components/api-server/src/client.rs | 32 ++++---- components/api-server/src/error.rs | 3 +- components/api-server/src/routes.rs | 75 +++++++++---------- .../src/clp_config/package/config.rs | 3 +- .../src/clp_config/s3_config.rs | 3 +- .../clp-rust-utils/src/database/mysql.rs | 6 +- .../src/job_config/clp_io_config.rs | 8 +- .../src/job_config/compression.rs | 6 +- .../src/job_config/ingestion.rs | 3 +- .../clp-rust-utils/src/job_config/search.rs | 6 +- components/clp-rust-utils/src/logging.rs | 14 ++-- components/clp-rust-utils/src/s3/client.rs | 8 +- components/clp-rust-utils/src/sqs/client.rs | 8 +- .../clp-rust-utils/src/task_io/compression.rs | 3 +- components/clp-rust-utils/src/telemetry.rs | 6 +- .../clp-rust-utils/tests/clp_config_test.rs | 15 ++-- components/clp-tdl-package/src/common.rs | 7 +- .../src/task/compression/mod.rs | 12 +-- .../src/compression_job_submitter/mod.rs | 17 ++--- .../src/compression_job_submitter/spider.rs | 44 +++++------ .../compression-coordinator/src/partition.rs | 21 +++--- .../log-ingestor/src/aws_client_manager.rs | 3 +- .../log-ingestor/src/bin/log_ingestor.rs | 6 +- .../compression/compression_job_submitter.rs | 32 ++++---- .../log-ingestor/src/compression/listener.rs | 13 ++-- .../src/ingestion_job/s3_scanner.rs | 12 +-- .../src/ingestion_job/sqs_listener.rs | 25 +++---- .../log-ingestor/src/ingestion_job_manager.rs | 30 ++++---- .../ingestion_job_manager/clp_ingestion.rs | 65 ++++++++-------- components/log-ingestor/src/routes.rs | 42 +++++------ components/log-ingestor/tests/aws_config.rs | 3 +- .../tests/test_compression_listener.rs | 22 +++--- .../log-ingestor/tests/test_ingestion_job.rs | 58 +++++++------- components/log-ingestor/tests/test_scan.rs | 25 +++---- components/log-ingestor/tests/test_utils.rs | 6 +- 36 files changed, 316 insertions(+), 329 deletions(-) diff --git a/components/api-server/src/bin/api_server.rs b/components/api-server/src/bin/api_server.rs index 8976dcffb0..5f6603749a 100644 --- a/components/api-server/src/bin/api_server.rs +++ b/components/api-server/src/bin/api_server.rs @@ -1,6 +1,7 @@ use anyhow::Context; use clap::Parser; -use clp_rust_utils::{clp_config::package, serde::yaml}; +use clp_rust_utils::clp_config::package; +use clp_rust_utils::serde::yaml; #[derive(Parser)] #[command(version, about = "API Server for CLP.")] diff --git a/components/api-server/src/client.rs b/components/api-server/src/client.rs index 01d3ab6c3b..f8e2a3e01c 100644 --- a/components/api-server/src/client.rs +++ b/components/api-server/src/client.rs @@ -1,22 +1,28 @@ use std::pin::Pin; use async_stream::stream; -use chrono::{DateTime, TimeZone, Utc}; +use chrono::DateTime; +use chrono::TimeZone; +use chrono::Utc; +use clp_rust_utils::aws::AWS_DEFAULT_REGION; +use clp_rust_utils::clp_config::package::config::Config; +use clp_rust_utils::clp_config::package::config::StorageEngine; +use clp_rust_utils::clp_config::package::config::StreamOutputStorage; +use clp_rust_utils::clp_config::package::credentials::Credentials; +use clp_rust_utils::database::mysql::create_clp_db_mysql_pool; pub use clp_rust_utils::job_config::CompressionJobStatus; -use clp_rust_utils::{ - aws::AWS_DEFAULT_REGION, - clp_config::package::{ - config::{Config, StorageEngine, StreamOutputStorage}, - credentials::Credentials, - }, - database::mysql::create_clp_db_mysql_pool, - job_config::{QUERY_JOBS_TABLE_NAME, QueryJobStatus, QueryJobType, SearchJobConfig}, -}; -use futures::{Stream, StreamExt}; +use clp_rust_utils::job_config::QUERY_JOBS_TABLE_NAME; +use clp_rust_utils::job_config::QueryJobStatus; +use clp_rust_utils::job_config::QueryJobType; +use clp_rust_utils::job_config::SearchJobConfig; +use futures::Stream; +use futures::StreamExt; use pin_project_lite::pin_project; -use serde::{Deserialize, Serialize}; +use serde::Deserialize; +use serde::Serialize; use sqlx::Row; -use utoipa::{IntoParams, ToSchema}; +use utoipa::IntoParams; +use utoipa::ToSchema; pub use crate::error::ClientError; diff --git a/components/api-server/src/error.rs b/components/api-server/src/error.rs index 484b462b45..3447f3d616 100644 --- a/components/api-server/src/error.rs +++ b/components/api-server/src/error.rs @@ -1,4 +1,5 @@ -use aws_sdk_s3::{error::SdkError, primitives::ByteStreamError}; +use aws_sdk_s3::error::SdkError; +use aws_sdk_s3::primitives::ByteStreamError; use num_enum::TryFromPrimitive; use thiserror::Error; diff --git a/components/api-server/src/routes.rs b/components/api-server/src/routes.rs index c10ac175da..fae107f48a 100644 --- a/components/api-server/src/routes.rs +++ b/components/api-server/src/routes.rs @@ -1,29 +1,31 @@ -use axum::{ - Json, - extract::{Path, Query, State}, - http::StatusCode, - response::{ - IntoResponse, - Sse, - sse::{Event, KeepAlive}, - }, - routing::get, -}; -use futures::{Stream, StreamExt}; -use serde::{Deserialize, Serialize}; +use axum::Json; +use axum::extract::Path; +use axum::extract::Query; +use axum::extract::State; +use axum::http::StatusCode; +use axum::response::IntoResponse; +use axum::response::Sse; +use axum::response::sse::Event; +use axum::response::sse::KeepAlive; +use axum::routing::get; +use futures::Stream; +use futures::StreamExt; +use serde::Deserialize; +use serde::Serialize; use thiserror::Error; -use tower_http::cors::{Any, CorsLayer}; -use utoipa::{OpenApi, ToSchema}; -use utoipa_axum::{router::OpenApiRouter, routes}; - -use crate::client::{ - Client, - ClientError, - CompressionUsage, - CompressionUsageParams, - QueryConfig, - ValidatedCompressionUsageParams, -}; +use tower_http::cors::Any; +use tower_http::cors::CorsLayer; +use utoipa::OpenApi; +use utoipa::ToSchema; +use utoipa_axum::router::OpenApiRouter; +use utoipa_axum::routes; + +use crate::client::Client; +use crate::client::ClientError; +use crate::client::CompressionUsage; +use crate::client::CompressionUsageParams; +use crate::client::QueryConfig; +use crate::client::ValidatedCompressionUsageParams; /// Factory method to create an Axum router configured with all API routes. /// @@ -65,14 +67,12 @@ pub fn from_client(client: Client) -> Result { mod api_doc { // Using `super::...` can cause `super` to appear as a tag in the generated OpenAPI // documentation. Importing the paths directly prevents this issue. - use super::{ - __path_cancel_query, - __path_compression_usage, - __path_health, - __path_query, - __path_query_results, - CompressionUsage, - }; + use super::__path_cancel_query; + use super::__path_compression_usage; + use super::__path_health; + use super::__path_query; + use super::__path_query_results; + use super::CompressionUsage; use crate::client::CompressionJobStatus; #[derive(utoipa::OpenApi)] @@ -337,11 +337,10 @@ impl IntoResponse for HandlerError { #[cfg(test)] mod tests { - use axum::{ - body::Body, - http::{Request, StatusCode}, - routing::get, - }; + use axum::body::Body; + use axum::http::Request; + use axum::http::StatusCode; + use axum::routing::get; use http_body_util::BodyExt; use tower::ServiceExt; 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 29863686b0..80645c131e 100644 --- a/components/clp-rust-utils/src/clp_config/package/config.rs +++ b/components/clp-rust-utils/src/clp_config/package/config.rs @@ -1,6 +1,7 @@ use serde::Deserialize; -use crate::clp_config::{AwsAuthentication, S3Config}; +use crate::clp_config::AwsAuthentication; +use crate::clp_config::S3Config; /// Mirror of `clp_py_utils.clp_config.ClpConfig`. /// diff --git a/components/clp-rust-utils/src/clp_config/s3_config.rs b/components/clp-rust-utils/src/clp_config/s3_config.rs index 0ceac30f91..50792094a1 100644 --- a/components/clp-rust-utils/src/clp_config/s3_config.rs +++ b/components/clp-rust-utils/src/clp_config/s3_config.rs @@ -1,5 +1,6 @@ use non_empty_string::NonEmptyString; -use serde::{Deserialize, Serialize}; +use serde::Deserialize; +use serde::Serialize; /// Represents the configuration for connecting to an S3 bucket. #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] diff --git a/components/clp-rust-utils/src/database/mysql.rs b/components/clp-rust-utils/src/database/mysql.rs index 1ea951184a..9a2aa127e3 100644 --- a/components/clp-rust-utils/src/database/mysql.rs +++ b/components/clp-rust-utils/src/database/mysql.rs @@ -1,10 +1,8 @@ use secrecy::ExposeSecret; use strum::IntoEnumIterator; -use crate::clp_config::package::{ - config::Database as DatabaseConfig, - credentials::Database as DatabaseCredentials, -}; +use crate::clp_config::package::config::Database as DatabaseConfig; +use crate::clp_config::package::credentials::Database as DatabaseCredentials; /// Implements [`sqlx::Type`] for `$ty` by delegating to `$delegate`. /// diff --git a/components/clp-rust-utils/src/job_config/clp_io_config.rs b/components/clp-rust-utils/src/job_config/clp_io_config.rs index 7f4e19cb51..312dac3e53 100644 --- a/components/clp-rust-utils/src/job_config/clp_io_config.rs +++ b/components/clp-rust-utils/src/job_config/clp_io_config.rs @@ -1,11 +1,9 @@ use non_empty_string::NonEmptyString; use serde::Serialize; -use crate::{ - clp_config::S3Config, - job_config::ingestion::JobId as IngestionJobId, - s3::S3ObjectMetadataId, -}; +use crate::clp_config::S3Config; +use crate::job_config::ingestion::JobId as IngestionJobId; +use crate::s3::S3ObjectMetadataId; /// Represents CLP IO config. #[derive(Debug, Clone, PartialEq, Eq, Serialize)] diff --git a/components/clp-rust-utils/src/job_config/compression.rs b/components/clp-rust-utils/src/job_config/compression.rs index ddebf2f880..e09ec10484 100644 --- a/components/clp-rust-utils/src/job_config/compression.rs +++ b/components/clp-rust-utils/src/job_config/compression.rs @@ -1,5 +1,7 @@ -use num_enum::{IntoPrimitive, TryFromPrimitive}; -use serde::{Deserialize, Serialize}; +use num_enum::IntoPrimitive; +use num_enum::TryFromPrimitive; +use serde::Deserialize; +use serde::Serialize; use strum::EnumString; use utoipa::ToSchema; diff --git a/components/clp-rust-utils/src/job_config/ingestion.rs b/components/clp-rust-utils/src/job_config/ingestion.rs index bc3f726ac2..43782e520d 100644 --- a/components/clp-rust-utils/src/job_config/ingestion.rs +++ b/components/clp-rust-utils/src/job_config/ingestion.rs @@ -1,6 +1,7 @@ pub mod s3 { use non_empty_string::NonEmptyString; - use serde::{Deserialize, Serialize}; + use serde::Deserialize; + use serde::Serialize; use thiserror::Error; use utoipa::ToSchema; diff --git a/components/clp-rust-utils/src/job_config/search.rs b/components/clp-rust-utils/src/job_config/search.rs index fcf25275b9..0ece3c2295 100644 --- a/components/clp-rust-utils/src/job_config/search.rs +++ b/components/clp-rust-utils/src/job_config/search.rs @@ -1,5 +1,7 @@ -use num_enum::{IntoPrimitive, TryFromPrimitive}; -use serde::{Deserialize, Serialize}; +use num_enum::IntoPrimitive; +use num_enum::TryFromPrimitive; +use serde::Deserialize; +use serde::Serialize; pub const QUERY_JOBS_TABLE_NAME: &str = "query_jobs"; diff --git a/components/clp-rust-utils/src/logging.rs b/components/clp-rust-utils/src/logging.rs index ad3ef343b4..bc130b08eb 100644 --- a/components/clp-rust-utils/src/logging.rs +++ b/components/clp-rust-utils/src/logging.rs @@ -1,11 +1,9 @@ -use tracing_appender::{ - non_blocking::{NonBlockingBuilder, WorkerGuard}, - rolling::{RollingFileAppender, Rotation}, -}; -use tracing_subscriber::{ - fmt::writer::MakeWriterExt, - {self}, -}; +use tracing_appender::non_blocking::NonBlockingBuilder; +use tracing_appender::non_blocking::WorkerGuard; +use tracing_appender::rolling::RollingFileAppender; +use tracing_appender::rolling::Rotation; +use tracing_subscriber::fmt::writer::MakeWriterExt; +use tracing_subscriber::{self}; /// Opaque struct to hold the worker guards for the background log writers. /// These guards must be held for the lifetime of the program to ensure logs are flushed. diff --git a/components/clp-rust-utils/src/s3/client.rs b/components/clp-rust-utils/src/s3/client.rs index 9160f5d0f1..134b1645a9 100644 --- a/components/clp-rust-utils/src/s3/client.rs +++ b/components/clp-rust-utils/src/s3/client.rs @@ -1,8 +1,8 @@ use aws_config::BehaviorVersion; -use aws_sdk_s3::{ - Client, - config::{Builder, Credentials, Region}, -}; +use aws_sdk_s3::Client; +use aws_sdk_s3::config::Builder; +use aws_sdk_s3::config::Credentials; +use aws_sdk_s3::config::Region; use non_empty_string::NonEmptyString; use crate::clp_config::AwsAuthentication; diff --git a/components/clp-rust-utils/src/sqs/client.rs b/components/clp-rust-utils/src/sqs/client.rs index 4586837146..2d5be46171 100644 --- a/components/clp-rust-utils/src/sqs/client.rs +++ b/components/clp-rust-utils/src/sqs/client.rs @@ -1,8 +1,8 @@ use aws_config::BehaviorVersion; -use aws_sdk_sqs::{ - Client, - config::{Builder, Credentials, Region}, -}; +use aws_sdk_sqs::Client; +use aws_sdk_sqs::config::Builder; +use aws_sdk_sqs::config::Credentials; +use aws_sdk_sqs::config::Region; use non_empty_string::NonEmptyString; use crate::clp_config::AwsAuthentication; diff --git a/components/clp-rust-utils/src/task_io/compression.rs b/components/clp-rust-utils/src/task_io/compression.rs index db6c5ca627..3bbf62a760 100644 --- a/components/clp-rust-utils/src/task_io/compression.rs +++ b/components/clp-rust-utils/src/task_io/compression.rs @@ -1,7 +1,8 @@ //! Protocol types exchanged with the Spider (Huntsman) tasks that run CLP S3 compression jobs. use non_empty_string::NonEmptyString; -use serde::{Deserialize, Serialize}; +use serde::Deserialize; +use serde::Serialize; use crate::clp_config::AwsAuthentication; diff --git a/components/clp-rust-utils/src/telemetry.rs b/components/clp-rust-utils/src/telemetry.rs index 45d5ba5d32..ee237d12a1 100644 --- a/components/clp-rust-utils/src/telemetry.rs +++ b/components/clp-rust-utils/src/telemetry.rs @@ -2,9 +2,11 @@ use std::env; -use opentelemetry_sdk::metrics::{PeriodicReader, SdkMeterProvider}; +use opentelemetry_sdk::metrics::PeriodicReader; +use opentelemetry_sdk::metrics::SdkMeterProvider; -use crate::{Error, clp_config::package::config::Telemetry}; +use crate::Error; +use crate::clp_config::package::config::Telemetry; /// RAII guard that shuts down the meter provider and flushes pending metric exports when dropped. pub struct TelemetryGuard { diff --git a/components/clp-rust-utils/tests/clp_config_test.rs b/components/clp-rust-utils/tests/clp_config_test.rs index 41fd5edf11..adbfab2b21 100644 --- a/components/clp-rust-utils/tests/clp_config_test.rs +++ b/components/clp-rust-utils/tests/clp_config_test.rs @@ -1,9 +1,12 @@ -use clp_rust_utils::{ - clp_config::{AwsAuthentication, AwsCredentials, S3Config}, - job_config::{ClpIoConfig, InputConfig, OutputConfig, S3ObjectMetadataInputConfig}, - serde::BrotliMsgpack, - types::non_empty_string::ExpectedNonEmpty, -}; +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::job_config::ClpIoConfig; +use clp_rust_utils::job_config::InputConfig; +use clp_rust_utils::job_config::OutputConfig; +use clp_rust_utils::job_config::S3ObjectMetadataInputConfig; +use clp_rust_utils::serde::BrotliMsgpack; +use clp_rust_utils::types::non_empty_string::ExpectedNonEmpty; use non_empty_string::NonEmptyString; use serde_json::Value; diff --git a/components/clp-tdl-package/src/common.rs b/components/clp-tdl-package/src/common.rs index 557b40bc1a..11036b1fe0 100644 --- a/components/clp-tdl-package/src/common.rs +++ b/components/clp-tdl-package/src/common.rs @@ -1,10 +1,9 @@ //! Process-global state shared by this package's tasks: the Tokio runtime, the Spider task //! executor config, and this cdylib's `tracing` subscriber. -use std::{ - path::{Path, PathBuf}, - sync::OnceLock, -}; +use std::path::Path; +use std::path::PathBuf; +use std::sync::OnceLock; use anyhow::Context; use clp_rust_utils::clp_config::package::config::SpiderTaskExecutorConfig; diff --git a/components/clp-tdl-package/src/task/compression/mod.rs b/components/clp-tdl-package/src/task/compression/mod.rs index 8f08baf1d4..70cbdc86f8 100644 --- a/components/clp-tdl-package/src/task/compression/mod.rs +++ b/components/clp-tdl-package/src/task/compression/mod.rs @@ -1,11 +1,11 @@ //! The compression tasks: the `#[task]` wrappers Spider invokes and their implementations. -use clp_rust_utils::task_io::compression::{ - ClpSCompressionOption, - CompressionTaskOutput, - S3InputSource, -}; -use spider_tdl::{TaskContext, TdlError, task}; +use clp_rust_utils::task_io::compression::ClpSCompressionOption; +use clp_rust_utils::task_io::compression::CompressionTaskOutput; +use clp_rust_utils::task_io::compression::S3InputSource; +use spider_tdl::TaskContext; +use spider_tdl::TdlError; +use spider_tdl::task; mod commit; mod compress; diff --git a/components/compression-coordinator/src/compression_job_submitter/mod.rs b/components/compression-coordinator/src/compression_job_submitter/mod.rs index 012e3f11c4..31d32ec5cd 100644 --- a/components/compression-coordinator/src/compression_job_submitter/mod.rs +++ b/components/compression-coordinator/src/compression_job_submitter/mod.rs @@ -6,15 +6,14 @@ mod spider; use std::time::Duration; use async_trait::async_trait; -use clp_rust_utils::{ - job_config::CompressionJobId, - task_io::compression::{ClpSCompressionOption, S3InputSource}, -}; -use serde::{Deserialize, Serialize}; -use spider_core::{ - task::ExecutionPolicy, - types::id::{JobId, ResourceGroupId}, -}; +use clp_rust_utils::job_config::CompressionJobId; +use clp_rust_utils::task_io::compression::ClpSCompressionOption; +use clp_rust_utils::task_io::compression::S3InputSource; +use serde::Deserialize; +use serde::Serialize; +use spider_core::task::ExecutionPolicy; +use spider_core::types::id::JobId; +use spider_core::types::id::ResourceGroupId; use crate::error::Error; diff --git a/components/compression-coordinator/src/compression_job_submitter/spider.rs b/components/compression-coordinator/src/compression_job_submitter/spider.rs index b30fa6b3fb..e024053b8c 100644 --- a/components/compression-coordinator/src/compression_job_submitter/spider.rs +++ b/components/compression-coordinator/src/compression_job_submitter/spider.rs @@ -3,32 +3,26 @@ use std::time::Duration; use async_trait::async_trait; -use clp_rust_utils::{ - job_config::CompressionJobId, - task_io::compression::{ClpSCompressionOption, S3InputSource}, -}; -use spider_client::{SpiderClient, error::ClientError}; -use spider_core::{ - job::JobState, - task::{ - DataTypeDescriptor, - ExecutionPolicy, - TaskDescriptor, - TaskGraph, - TdlContext, - TerminationTaskDescriptor, - ValueTypeDescriptor, - }, - types::{ - id::{JobId, ResourceGroupId}, - io::TaskInput, - }, -}; +use clp_rust_utils::job_config::CompressionJobId; +use clp_rust_utils::task_io::compression::ClpSCompressionOption; +use clp_rust_utils::task_io::compression::S3InputSource; +use spider_client::SpiderClient; +use spider_client::error::ClientError; +use spider_core::job::JobState; +use spider_core::task::DataTypeDescriptor; +use spider_core::task::ExecutionPolicy; +use spider_core::task::TaskDescriptor; +use spider_core::task::TaskGraph; +use spider_core::task::TdlContext; +use spider_core::task::TerminationTaskDescriptor; +use spider_core::task::ValueTypeDescriptor; +use spider_core::types::id::JobId; +use spider_core::types::id::ResourceGroupId; +use spider_core::types::io::TaskInput; -use crate::{ - compression_job_submitter::{CompressionJobOutcome, S3CompressionJobSubmitter}, - error::Error, -}; +use crate::compression_job_submitter::CompressionJobOutcome; +use crate::compression_job_submitter::S3CompressionJobSubmitter; +use crate::error::Error; #[async_trait] impl S3CompressionJobSubmitter for SpiderClient { diff --git a/components/compression-coordinator/src/partition.rs b/components/compression-coordinator/src/partition.rs index 75f85707f1..a99a281887 100644 --- a/components/compression-coordinator/src/partition.rs +++ b/components/compression-coordinator/src/partition.rs @@ -1,13 +1,13 @@ //! Partitioning of S3 objects into compression-task inputs. -use std::{collections::VecDeque, path::Path}; +use std::collections::VecDeque; +use std::path::Path; -use clp_rust_utils::{ - clp_config::S3Config, - job_config::{ClpIoConfig, InputConfig}, - s3::ObjectMetadata, - task_io::compression::S3InputSource, -}; +use clp_rust_utils::clp_config::S3Config; +use clp_rust_utils::job_config::ClpIoConfig; +use clp_rust_utils::job_config::InputConfig; +use clp_rust_utils::s3::ObjectMetadata; +use clp_rust_utils::task_io::compression::S3InputSource; use crate::Error; @@ -292,10 +292,9 @@ fn group_files_by_similar_filenames(mut files: Vec) -> Vec Date: Thu, 23 Jul 2026 13:56:06 -0400 Subject: [PATCH 5/6] lint fix --- .../compression-coordinator/src/job_handle.rs | 27 ++++++++++--------- 1 file changed, 15 insertions(+), 12 deletions(-) diff --git a/components/compression-coordinator/src/job_handle.rs b/components/compression-coordinator/src/job_handle.rs index 8d47da8a41..3d355eaea9 100644 --- a/components/compression-coordinator/src/job_handle.rs +++ b/components/compression-coordinator/src/job_handle.rs @@ -1,20 +1,23 @@ //! Handle for driving a single S3 compression job to completion. -use std::{sync::Arc, time::Duration}; +use std::sync::Arc; +use std::time::Duration; -use clp_rust_utils::{ - clp_config::package::config::Database, - dataset::VALID_DATASET_NAME_REGEX, - job_config::{ClpIoConfig, CompressionJobId, CompressionJobStatus, InputConfig}, - task_io::compression::{ClpSCompressionOption, S3InputSource}, -}; -use spider_core::{ - task::ExecutionPolicy, - types::id::{JobId as SpiderJobId, ResourceGroupId}, -}; +use clp_rust_utils::clp_config::package::config::Database; +use clp_rust_utils::dataset::VALID_DATASET_NAME_REGEX; +use clp_rust_utils::job_config::ClpIoConfig; +use clp_rust_utils::job_config::CompressionJobId; +use clp_rust_utils::job_config::CompressionJobStatus; +use clp_rust_utils::job_config::InputConfig; +use clp_rust_utils::task_io::compression::ClpSCompressionOption; +use clp_rust_utils::task_io::compression::S3InputSource; +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, compression_job_submitter::S3CompressionJobSubmitter}; +use crate::Error; +use crate::compression_job_submitter::S3CompressionJobSubmitter; /// Options for a compression job running in Spider. pub struct SpiderOption { From a16de2d07a84c7d6deb9e7c3eb40549811db60c6 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Wed, 29 Jul 2026 15:28:21 -0400 Subject: [PATCH 6/6] Run rust linters --- .../src/clp_config/package/config.rs | 33 +++---- components/clp-rust-utils/src/dataset.rs | 3 +- .../src/job_config/clp_io_config.rs | 3 +- components/clp-rust-utils/src/s3/url.rs | 3 +- .../src/serde/brotli_msgpack.rs | 9 +- .../src/task/compression/commit.rs | 14 +-- .../src/task/compression/compress.rs | 98 +++++++++---------- .../src/bin/compression_coordinator.rs | 11 +-- .../src/coordination.rs | 41 ++++---- .../compression-coordinator/src/error.rs | 3 +- .../compression-coordinator/src/job_handle.rs | 39 ++++---- .../compression/compression_job_submitter.rs | 31 +++--- 12 files changed, 144 insertions(+), 144 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 f911728a27..badacf3804 100644 --- a/components/clp-rust-utils/src/clp_config/package/config.rs +++ b/components/clp-rust-utils/src/clp_config/package/config.rs @@ -1,15 +1,14 @@ -use std::{ - num::{NonZeroU32, NonZeroU64}, - path::{Path, PathBuf}, -}; +use std::num::NonZeroU32; +use std::num::NonZeroU64; +use std::path::Path; +use std::path::PathBuf; use non_empty_string::NonEmptyString; use serde::Deserialize; -use crate::{ - clp_config::{AwsAuthentication, S3Config}, - dataset::resolve_dataset_name, -}; +use crate::clp_config::AwsAuthentication; +use crate::clp_config::S3Config; +use crate::dataset::resolve_dataset_name; /// Mirror of `clp_py_utils.clp_config.ClpConfig`. /// @@ -562,13 +561,11 @@ fn default_archive_staging_directory() -> String { mod tests { use std::path::Path; - use super::{ - ArchiveOutput, - ArchiveOutputStorage, - Database, - LogsInput, - SpiderTaskExecutorConfig, - }; + use super::ArchiveOutput; + use super::ArchiveOutputStorage; + use super::Database; + use super::LogsInput; + use super::SpiderTaskExecutorConfig; #[test] fn deserialize_logs_input_s3_config() { @@ -671,7 +668,8 @@ mod tests { fn dataset_archive_storage_directory_s3() { use non_empty_string::NonEmptyString; - use crate::clp_config::{AwsAuthentication, S3Config}; + use crate::clp_config::AwsAuthentication; + use crate::clp_config::S3Config; let archive_output = ArchiveOutput { storage: ArchiveOutputStorage::S3 { @@ -769,7 +767,8 @@ mod tests { fn s3_config_with_staging_directory(staging_directory: &str) -> SpiderTaskExecutorConfig { use non_empty_string::NonEmptyString; - use crate::clp_config::{AwsAuthentication, S3Config}; + use crate::clp_config::AwsAuthentication; + use crate::clp_config::S3Config; SpiderTaskExecutorConfig { archive_output: ArchiveOutput { diff --git a/components/clp-rust-utils/src/dataset.rs b/components/clp-rust-utils/src/dataset.rs index 0c9dc9ecd2..4b460e0d36 100644 --- a/components/clp-rust-utils/src/dataset.rs +++ b/components/clp-rust-utils/src/dataset.rs @@ -18,7 +18,8 @@ pub fn resolve_dataset_name(dataset: Option<&str>) -> &str { #[cfg(test)] mod tests { - use super::{CLP_DEFAULT_DATASET_NAME, resolve_dataset_name}; + use super::CLP_DEFAULT_DATASET_NAME; + use super::resolve_dataset_name; #[test] fn resolve_dataset_name_passes_through_some() { diff --git a/components/clp-rust-utils/src/job_config/clp_io_config.rs b/components/clp-rust-utils/src/job_config/clp_io_config.rs index d0ff8d483c..84ae2e22c5 100644 --- a/components/clp-rust-utils/src/job_config/clp_io_config.rs +++ b/components/clp-rust-utils/src/job_config/clp_io_config.rs @@ -1,5 +1,6 @@ use non_empty_string::NonEmptyString; -use serde::{Deserialize, Serialize}; +use serde::Deserialize; +use serde::Serialize; use crate::clp_config::S3Config; use crate::job_config::ingestion::JobId as IngestionJobId; diff --git a/components/clp-rust-utils/src/s3/url.rs b/components/clp-rust-utils/src/s3/url.rs index 1b44b76db1..c9eb138903 100644 --- a/components/clp-rust-utils/src/s3/url.rs +++ b/components/clp-rust-utils/src/s3/url.rs @@ -82,7 +82,8 @@ mod tests { use non_empty_string::NonEmptyString; use super::generate_s3_url; - use crate::{Error, types::non_empty_string::ExpectedNonEmpty}; + use crate::Error; + use crate::types::non_empty_string::ExpectedNonEmpty; fn to_non_empty_string(value: &'static str) -> NonEmptyString { NonEmptyString::from_static_str(value) diff --git a/components/clp-rust-utils/src/serde/brotli_msgpack.rs b/components/clp-rust-utils/src/serde/brotli_msgpack.rs index 0ac1beae90..8ad301465c 100644 --- a/components/clp-rust-utils/src/serde/brotli_msgpack.rs +++ b/components/clp-rust-utils/src/serde/brotli_msgpack.rs @@ -1,7 +1,10 @@ -use std::io::{Read, Write}; +use std::io::Read; +use std::io::Write; -use brotli::{CompressorWriter, Decompressor}; -use serde::{Serialize, de::DeserializeOwned}; +use brotli::CompressorWriter; +use brotli::Decompressor; +use serde::Serialize; +use serde::de::DeserializeOwned; use crate::Error; diff --git a/components/clp-tdl-package/src/task/compression/commit.rs b/components/clp-tdl-package/src/task/compression/commit.rs index e32b0e2066..d325286655 100644 --- a/components/clp-tdl-package/src/task/compression/commit.rs +++ b/components/clp-tdl-package/src/task/compression/commit.rs @@ -1,13 +1,13 @@ //! The commit worker that publishes a compression job's archives to CLP's metadata store. use anyhow::Context; -use clp_rust_utils::{ - clp_config::package::credentials, - database::mysql::create_clp_db_mysql_pool, - dataset::{VALID_DATASET_NAME_REGEX, resolve_dataset_name}, - job_config::{CompressionJobId, CompressionJobStatus}, - task_io::compression::ArchiveMetadata, -}; +use clp_rust_utils::clp_config::package::credentials; +use clp_rust_utils::database::mysql::create_clp_db_mysql_pool; +use clp_rust_utils::dataset::VALID_DATASET_NAME_REGEX; +use clp_rust_utils::dataset::resolve_dataset_name; +use clp_rust_utils::job_config::CompressionJobId; +use clp_rust_utils::job_config::CompressionJobStatus; +use clp_rust_utils::task_io::compression::ArchiveMetadata; use secrecy::SecretString; use spider_core::types::id::JobId; diff --git a/components/clp-tdl-package/src/task/compression/compress.rs b/components/clp-tdl-package/src/task/compression/compress.rs index 9d8a205ede..af2eb94ff0 100644 --- a/components/clp-tdl-package/src/task/compression/compress.rs +++ b/components/clp-tdl-package/src/task/compression/compress.rs @@ -1,37 +1,33 @@ //! The `clp-s` compression worker that turns an S3 input source into archives. -use std::{ - ffi::OsString, - io::{BufRead, BufReader, Read}, - path::{Path, PathBuf}, - process::{Command, Stdio}, -}; +use std::ffi::OsString; +use std::io::BufRead; +use std::io::BufReader; +use std::io::Read; +use std::path::Path; +use std::path::PathBuf; +use std::process::Command; +use std::process::Stdio; use anyhow::Context; -use clp_rust_utils::{ - aws::AWS_DEFAULT_REGION, - clp_config::{ - AwsAuthentication, - S3Config, - package::config::{ - ArchiveOutput, - ArchiveOutputStorage, - Database, - SpiderTaskExecutorConfig, - }, - }, - dataset::resolve_dataset_name, - s3::{create_new_client, generate_s3_url}, - task_io::compression::{ - ArchiveMetadata, - ClpSCompressionOption, - CompressionTaskOutput, - S3InputSource, - }, -}; +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; +use clp_rust_utils::clp_config::package::config::Database; +use clp_rust_utils::clp_config::package::config::SpiderTaskExecutorConfig; +use clp_rust_utils::dataset::resolve_dataset_name; +use clp_rust_utils::s3::create_new_client; +use clp_rust_utils::s3::generate_s3_url; +use clp_rust_utils::task_io::compression::ArchiveMetadata; +use clp_rust_utils::task_io::compression::ClpSCompressionOption; +use clp_rust_utils::task_io::compression::CompressionTaskOutput; +use clp_rust_utils::task_io::compression::S3InputSource; use non_empty_string::NonEmptyString; -use crate::common::{clp_home, runtime}; +use crate::common::clp_home; +use crate::common::runtime; /// Compresses the given S3 objects into archives, uploads them to S3, and returns their metadata /// for the commit task. @@ -854,32 +850,30 @@ fn kill_clp_s_and_read_stderr( #[cfg(test)] mod tests { - use std::{ - ffi::OsString, - path::{Path, PathBuf}, - }; - - use clp_rust_utils::{ - clp_config::{ - AwsAuthentication, - AwsCredentials, - S3Config, - package::config::{ArchiveOutput, ArchiveOutputStorage, ClpDbNames, Database}, - }, - task_io::compression::{ArchiveMetadata, ClpSCompressionOption, S3InputSource}, - }; + use std::ffi::OsString; + use std::path::Path; + 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; + use clp_rust_utils::clp_config::package::config::ClpDbNames; + use clp_rust_utils::clp_config::package::config::Database; + use clp_rust_utils::task_io::compression::ArchiveMetadata; + use clp_rust_utils::task_io::compression::ClpSCompressionOption; + use clp_rust_utils::task_io::compression::S3InputSource; use non_empty_string::NonEmptyString; - use super::{ - ClpSInput, - build_clp_s_args, - build_indexer_args, - build_log_converter_args, - build_s3_logs_list, - create_archive_s3_key, - parse_archive_stats, - s3_credential_env, - }; + use super::ClpSInput; + use super::build_clp_s_args; + use super::build_indexer_args; + use super::build_log_converter_args; + 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<()> { diff --git a/components/compression-coordinator/src/bin/compression_coordinator.rs b/components/compression-coordinator/src/bin/compression_coordinator.rs index cf069567df..072960025c 100644 --- a/components/compression-coordinator/src/bin/compression_coordinator.rs +++ b/components/compression-coordinator/src/bin/compression_coordinator.rs @@ -1,11 +1,10 @@ -use std::{path::PathBuf, time::Duration}; +use std::path::PathBuf; +use std::time::Duration; use clap::Parser; -use clp_rust_utils::{ - clp_config::package::{self}, - database::mysql::create_clp_db_mysql_pool, - serde::yaml, -}; +use clp_rust_utils::clp_config::package::{self}; +use clp_rust_utils::database::mysql::create_clp_db_mysql_pool; +use clp_rust_utils::serde::yaml; /// Command-line arguments for the compression coordinator. #[derive(Debug, Parser)] diff --git a/components/compression-coordinator/src/coordination.rs b/components/compression-coordinator/src/coordination.rs index 0873e9d42d..762e62c38a 100644 --- a/components/compression-coordinator/src/coordination.rs +++ b/components/compression-coordinator/src/coordination.rs @@ -1,32 +1,31 @@ //! The coordinator poll loop that discovers pending CLP compression jobs and dispatches them to //! Spider. -use std::{sync::Arc, time::Duration}; - -use clp_rust_utils::{ - clp_config::package::config::{ - CompressionCoordinator as CoordinatorConfig, - Database as DatabaseConfig, - Spider as SpiderConfig, - SpiderResourceGroup, - }, - job_config::{ClpIoConfig, CompressionJobId, CompressionJobStatus}, - serde::BrotliMsgpack, -}; +use std::sync::Arc; +use std::time::Duration; + +use clp_rust_utils::clp_config::package::config::CompressionCoordinator as CoordinatorConfig; +use clp_rust_utils::clp_config::package::config::Database as DatabaseConfig; +use clp_rust_utils::clp_config::package::config::Spider as SpiderConfig; +use clp_rust_utils::clp_config::package::config::SpiderResourceGroup; +use clp_rust_utils::job_config::ClpIoConfig; +use clp_rust_utils::job_config::CompressionJobId; +use clp_rust_utils::job_config::CompressionJobStatus; +use clp_rust_utils::serde::BrotliMsgpack; use const_format::formatcp; use spider_client::SpiderClient; -use spider_core::{ - task::{ExecutionPolicy, TimeoutPolicy}, - types::id::{JobId as SpiderJobId, ResourceGroupId}, -}; -use tokio::{select, time::Instant}; +use spider_core::task::ExecutionPolicy; +use spider_core::task::TimeoutPolicy; +use spider_core::types::id::JobId as SpiderJobId; +use spider_core::types::id::ResourceGroupId; +use tokio::select; +use tokio::time::Instant; use tokio_util::sync::CancellationToken; use tonic::transport::Endpoint; -use crate::{ - Error, - job_handle::{S3CompressionJobHandle, SpiderOption}, -}; +use crate::Error; +use crate::job_handle::S3CompressionJobHandle; +use crate::job_handle::SpiderOption; /// Coordinator for fetching new compression jobs and submitting them to Spider. pub struct Coordinator { diff --git a/components/compression-coordinator/src/error.rs b/components/compression-coordinator/src/error.rs index b3da15ed29..5ceb703154 100644 --- a/components/compression-coordinator/src/error.rs +++ b/components/compression-coordinator/src/error.rs @@ -1,6 +1,7 @@ //! The crate-level error type for the compression coordinator. -use clp_rust_utils::{job_config::ingestion::JobId as IngestionJobId, s3::S3ObjectMetadataId}; +use clp_rust_utils::job_config::ingestion::JobId as IngestionJobId; +use clp_rust_utils::s3::S3ObjectMetadataId; /// Errors returned by the compression coordinator. #[derive(Debug, thiserror::Error)] diff --git a/components/compression-coordinator/src/job_handle.rs b/components/compression-coordinator/src/job_handle.rs index e3651eb44e..7358608b42 100644 --- a/components/compression-coordinator/src/job_handle.rs +++ b/components/compression-coordinator/src/job_handle.rs @@ -1,27 +1,30 @@ //! Handle for driving a single S3 compression job to completion. -use std::{sync::Arc, time::Duration}; - -use clp_rust_utils::{ - clp_config::package::config::Database, - dataset::VALID_DATASET_NAME_REGEX, - job_config::{ClpIoConfig, CompressionJobId, CompressionJobStatus, InputConfig}, - s3::{ObjectMetadata, S3ObjectMetadataId}, - task_io::compression::{ClpSCompressionOption, S3InputSource}, -}; +use std::sync::Arc; +use std::time::Duration; + +use clp_rust_utils::clp_config::package::config::Database; +use clp_rust_utils::dataset::VALID_DATASET_NAME_REGEX; +use clp_rust_utils::job_config::ClpIoConfig; +use clp_rust_utils::job_config::CompressionJobId; +use clp_rust_utils::job_config::CompressionJobStatus; +use clp_rust_utils::job_config::InputConfig; +use clp_rust_utils::s3::ObjectMetadata; +use clp_rust_utils::s3::S3ObjectMetadataId; +use clp_rust_utils::task_io::compression::ClpSCompressionOption; +use clp_rust_utils::task_io::compression::S3InputSource; use const_format::formatcp; use non_empty_string::NonEmptyString; -use spider_core::{ - task::{ExecutionPolicy, TimeoutPolicy}, - types::id::{JobId as SpiderJobId, ResourceGroupId}, -}; +use spider_core::task::ExecutionPolicy; +use spider_core::task::TimeoutPolicy; +use spider_core::types::id::JobId as SpiderJobId; +use spider_core::types::id::ResourceGroupId; use sqlx::MySqlPool; -use crate::{ - Error, - compression_job_submitter::{CompressionJobOutcome, S3CompressionJobSubmitter}, - partition::CompressionInputBuilder, -}; +use crate::Error; +use crate::compression_job_submitter::CompressionJobOutcome; +use crate::compression_job_submitter::S3CompressionJobSubmitter; +use crate::partition::CompressionInputBuilder; /// Options for a compression job running in Spider. pub struct SpiderOption { diff --git a/components/log-ingestor/src/compression/compression_job_submitter.rs b/components/log-ingestor/src/compression/compression_job_submitter.rs index ad1f8d220e..151672534a 100644 --- a/components/log-ingestor/src/compression/compression_job_submitter.rs +++ b/components/log-ingestor/src/compression/compression_job_submitter.rs @@ -1,23 +1,22 @@ use anyhow::Result; use async_trait::async_trait; -use clp_rust_utils::{ - clp_config::{AwsAuthentication, S3Config, package::config::ArchiveOutput}, - dataset::CLP_DEFAULT_DATASET_NAME, - job_config::{ - ClpIoConfig, - CompressionJobId, - CompressionJobStatus, - InputConfig, - OutputConfig, - S3ObjectMetadataInputConfig, - ingestion::s3::BaseConfig, - }, - s3::S3ObjectMetadataId, - types::non_empty_string::ExpectedNonEmpty, -}; +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::dataset::CLP_DEFAULT_DATASET_NAME; +use clp_rust_utils::job_config::ClpIoConfig; +use clp_rust_utils::job_config::CompressionJobId; +use clp_rust_utils::job_config::CompressionJobStatus; +use clp_rust_utils::job_config::InputConfig; +use clp_rust_utils::job_config::OutputConfig; +use clp_rust_utils::job_config::S3ObjectMetadataInputConfig; +use clp_rust_utils::job_config::ingestion::s3::BaseConfig; +use clp_rust_utils::s3::S3ObjectMetadataId; +use clp_rust_utils::types::non_empty_string::ExpectedNonEmpty; use non_empty_string::NonEmptyString; -use crate::{compression::BufferSubmitter, ingestion_job_manager::ClpCompressionState}; +use crate::compression::BufferSubmitter; +use crate::ingestion_job_manager::ClpCompressionState; /// The CLP compression job table name. pub const CLP_COMPRESSION_JOB_TABLE_NAME: &str = "compression_jobs";