From ff72dba01394404c810e9e706fab7907efb029f4 Mon Sep 17 00:00:00 2001 From: Prajwal-banakar Date: Sat, 21 Mar 2026 05:24:49 +0000 Subject: [PATCH 1/4] Add validation for numeric config fields --- crates/fluss/src/client/connection.rs | 4 +- crates/fluss/src/client/table/remote_log.rs | 2 +- crates/fluss/src/config.rs | 178 ++++++++++++++++++-- 3 files changed, 168 insertions(+), 16 deletions(-) diff --git a/crates/fluss/src/client/connection.rs b/crates/fluss/src/client/connection.rs index a3ffd755..9727c8d7 100644 --- a/crates/fluss/src/client/connection.rs +++ b/crates/fluss/src/client/connection.rs @@ -41,7 +41,9 @@ impl FlussConnection { pub async fn new(arg: Config) -> Result { arg.validate_security() .map_err(|msg| Error::IllegalArgument { message: msg })?; - arg.validate_scanner_fetch() + arg.validate_scanner() + .map_err(|msg| Error::IllegalArgument { message: msg })?; + arg.validate_writer() .map_err(|msg| Error::IllegalArgument { message: msg })?; let timeout = Duration::from_millis(arg.connect_timeout_ms); diff --git a/crates/fluss/src/client/table/remote_log.rs b/crates/fluss/src/client/table/remote_log.rs index 6bc95512..4d96ce96 100644 --- a/crates/fluss/src/client/table/remote_log.rs +++ b/crates/fluss/src/client/table/remote_log.rs @@ -778,7 +778,7 @@ impl RemoteLogDownloader { let fetcher = Arc::new(ProductionFetcher { credentials_rx, local_log_dir: Arc::new(local_log_dir), - remote_log_read_concurrency: remote_log_read_concurrency.max(1), + remote_log_read_concurrency, }); Self::new_with_fetcher(fetcher, max_prefetch_segments, max_concurrent_downloads) diff --git a/crates/fluss/src/config.rs b/crates/fluss/src/config.rs index 2900e2f4..71c04e54 100644 --- a/crates/fluss/src/config.rs +++ b/crates/fluss/src/config.rs @@ -361,7 +361,16 @@ impl Config { } Ok(()) } - pub fn validate_scanner_fetch(&self) -> Result<(), String> { + pub fn validate_scanner(&self) -> Result<(), String> { + if self.scanner_remote_log_prefetch_num == 0 { + return Err("scanner_remote_log_prefetch_num must be > 0".to_string()); + } + if self.scanner_remote_log_read_concurrency == 0 { + return Err("scanner_remote_log_read_concurrency must be > 0".to_string()); + } + if self.remote_file_download_thread_num == 0 { + return Err("remote_file_download_thread_num must be > 0".to_string()); + } if self.scanner_log_fetch_min_bytes <= 0 { return Err("scanner_log_fetch_min_bytes must be > 0".to_string()); } @@ -387,6 +396,57 @@ impl Config { } Ok(()) } + + pub fn validate_writer(&self) -> Result<(), String> { + if self.writer_request_max_size <= 0 { + return Err("writer_request_max_size must be > 0".to_string()); + } + if self.writer_batch_size <= 0 { + return Err("writer_batch_size must be > 0".to_string()); + } + if self.writer_batch_timeout_ms < 0 { + return Err("writer_batch_timeout_ms must be >= 0".to_string()); + } + if self.writer_max_inflight_requests_per_bucket == 0 { + return Err("writer_max_inflight_requests_per_bucket must be > 0".to_string()); + } + if self.writer_buffer_memory_size == 0 { + return Err("writer_buffer_memory_size must be > 0".to_string()); + } + if self.writer_batch_size > self.writer_request_max_size { + return Err("writer_batch_size must be <= writer_request_max_size".to_string()); + } + if self.writer_batch_size as usize > self.writer_buffer_memory_size { + return Err("writer_batch_size must be <= writer_buffer_memory_size".to_string()); + } + // idempotence checks + if !self.writer_enable_idempotence { + return Ok(()); + } + let acks_is_all = self.writer_acks.eq_ignore_ascii_case("all") || self.writer_acks == "-1"; + if !acks_is_all { + return Err(format!( + "Idempotent writes require acks='all' (-1), but got acks='{}'", + self.writer_acks + )); + } + if self.writer_retries <= 0 { + return Err(format!( + "Idempotent writes require retries > 0, but got retries={}", + self.writer_retries + )); + } + if self.writer_max_inflight_requests_per_bucket + > MAX_IN_FLIGHT_REQUESTS_PER_BUCKET_FOR_IDEMPOTENCE + { + return Err(format!( + "Idempotent writes require max-inflight-requests-per-bucket <= {}, but got {}", + MAX_IN_FLIGHT_REQUESTS_PER_BUCKET_FOR_IDEMPOTENCE, + self.writer_max_inflight_requests_per_bucket + )); + } + Ok(()) + } } #[cfg(test)] @@ -456,13 +516,38 @@ mod tests { }; assert!(config.validate_security().is_err()); } + #[test] - fn test_scanner_fetch_defaults_valid() { + fn test_scanner_defaults_valid() { let config = Config::default(); - assert!(config.validate_scanner_fetch().is_ok()); - assert_eq!(config.scanner_log_fetch_max_bytes, 16 * 1024 * 1024); - assert_eq!(config.scanner_log_fetch_min_bytes, 1); - assert_eq!(config.scanner_log_fetch_wait_max_time_ms, 500); + assert!(config.validate_scanner().is_ok()); + } + + #[test] + fn test_scanner_remote_log_prefetch_num_zero() { + let config = Config { + scanner_remote_log_prefetch_num: 0, + ..Config::default() + }; + assert!(config.validate_scanner().is_err()); + } + + #[test] + fn test_scanner_remote_log_read_concurrency_zero() { + let config = Config { + scanner_remote_log_read_concurrency: 0, + ..Config::default() + }; + assert!(config.validate_scanner().is_err()); + } + + #[test] + fn test_remote_file_download_thread_num_zero() { + let config = Config { + remote_file_download_thread_num: 0, + ..Config::default() + }; + assert!(config.validate_scanner().is_err()); } #[test] @@ -472,7 +557,7 @@ mod tests { scanner_log_fetch_max_bytes: 1, ..Config::default() }; - assert!(config.validate_scanner_fetch().is_err()); + assert!(config.validate_scanner().is_err()); } #[test] @@ -481,13 +566,78 @@ mod tests { scanner_log_fetch_wait_max_time_ms: -1, ..Config::default() }; - assert!(config.validate_scanner_fetch().is_err()); + assert!(config.validate_scanner().is_err()); } #[test] - fn test_idempotence_default_is_valid() { + fn test_writer_defaults_valid() { let config = Config::default(); - assert!(config.validate_idempotence().is_ok()); + assert!(config.validate_writer().is_ok()); + } + + #[test] + fn test_writer_request_max_size_zero() { + let config = Config { + writer_request_max_size: 0, + ..Config::default() + }; + assert!(config.validate_writer().is_err()); + } + + #[test] + fn test_writer_batch_size_zero() { + let config = Config { + writer_batch_size: 0, + ..Config::default() + }; + assert!(config.validate_writer().is_err()); + } + + #[test] + fn test_writer_batch_timeout_negative() { + let config = Config { + writer_batch_timeout_ms: -1, + ..Config::default() + }; + assert!(config.validate_writer().is_err()); + } + + #[test] + fn test_writer_max_inflight_requests_per_bucket_zero() { + let config = Config { + writer_max_inflight_requests_per_bucket: 0, + ..Config::default() + }; + assert!(config.validate_writer().is_err()); + } + + #[test] + fn test_writer_buffer_memory_size_zero() { + let config = Config { + writer_buffer_memory_size: 0, + ..Config::default() + }; + assert!(config.validate_writer().is_err()); + } + + #[test] + fn test_writer_batch_size_exceeds_request_max_size() { + let config = Config { + writer_batch_size: 20 * 1024 * 1024, + writer_request_max_size: 10 * 1024 * 1024, + ..Config::default() + }; + assert!(config.validate_writer().is_err()); + } + + #[test] + fn test_writer_batch_size_exceeds_buffer_memory_size() { + let config = Config { + writer_batch_size: 128 * 1024 * 1024, + writer_buffer_memory_size: 64 * 1024 * 1024, + ..Config::default() + }; + assert!(config.validate_writer().is_err()); } #[test] @@ -499,7 +649,7 @@ mod tests { writer_max_inflight_requests_per_bucket: 100, ..Config::default() }; - assert!(config.validate_idempotence().is_ok()); + assert!(config.validate_writer().is_ok()); } #[test] @@ -509,7 +659,7 @@ mod tests { writer_acks: "1".to_string(), ..Config::default() }; - assert!(config.validate_idempotence().is_err()); + assert!(config.validate_writer().is_err()); } #[test] @@ -519,7 +669,7 @@ mod tests { writer_retries: 0, ..Config::default() }; - assert!(config.validate_idempotence().is_err()); + assert!(config.validate_writer().is_err()); } #[test] @@ -529,6 +679,6 @@ mod tests { writer_max_inflight_requests_per_bucket: 10, ..Config::default() }; - assert!(config.validate_idempotence().is_err()); + assert!(config.validate_writer().is_err()); } } From 6c324883d304a968e218fae74f0be68bb4a03947 Mon Sep 17 00:00:00 2001 From: Prajwal Banakar Date: Thu, 2 Apr 2026 15:36:09 +0000 Subject: [PATCH 2/4] improved --- .../fluss/src/client/write/writer_client.rs | 4 --- crates/fluss/src/config.rs | 29 +------------------ 2 files changed, 1 insertion(+), 32 deletions(-) diff --git a/crates/fluss/src/client/write/writer_client.rs b/crates/fluss/src/client/write/writer_client.rs index aee6bcd9..ffdf96b1 100644 --- a/crates/fluss/src/client/write/writer_client.rs +++ b/crates/fluss/src/client/write/writer_client.rs @@ -54,10 +54,6 @@ impl WriterClient { pub fn new(config: Config, metadata: Arc) -> Result { let ack = Self::get_ack(&config)?; - config - .validate_idempotence() - .map_err(|message| Error::IllegalArgument { message })?; - let idempotence_manager = Arc::new(IdempotenceManager::new( config.writer_enable_idempotence, config.writer_max_inflight_requests_per_bucket, diff --git a/crates/fluss/src/config.rs b/crates/fluss/src/config.rs index 71c04e54..aca37654 100644 --- a/crates/fluss/src/config.rs +++ b/crates/fluss/src/config.rs @@ -307,34 +307,7 @@ impl Config { /// Validates idempotence configuration. Returns `Ok(())` when the config is /// consistent, or an error message when idempotence is enabled but other /// settings are incompatible. - pub fn validate_idempotence(&self) -> Result<(), String> { - if !self.writer_enable_idempotence { - return Ok(()); - } - let acks_is_all = self.writer_acks.eq_ignore_ascii_case("all") || self.writer_acks == "-1"; - if !acks_is_all { - return Err(format!( - "Idempotent writes require acks='all' (-1), but got acks='{}'", - self.writer_acks - )); - } - if self.writer_retries <= 0 { - return Err(format!( - "Idempotent writes require retries > 0, but got retries={}", - self.writer_retries - )); - } - if self.writer_max_inflight_requests_per_bucket - > MAX_IN_FLIGHT_REQUESTS_PER_BUCKET_FOR_IDEMPOTENCE - { - return Err(format!( - "Idempotent writes require max-inflight-requests-per-bucket <= {}, but got {}", - MAX_IN_FLIGHT_REQUESTS_PER_BUCKET_FOR_IDEMPOTENCE, - self.writer_max_inflight_requests_per_bucket - )); - } - Ok(()) - } + /// Validates security configuration. Returns `Ok(())` when the config is /// consistent, or an error message when SASL is enabled but the config is From 7a808b5b6037a4f3ad5d0492af82f5e76e46f51a Mon Sep 17 00:00:00 2001 From: Prajwal Banakar Date: Tue, 14 Apr 2026 04:59:37 +0000 Subject: [PATCH 3/4] Added issue comments --- crates/fluss/src/client/connection.rs | 2 ++ crates/fluss/src/config.rs | 3 ++- 2 files changed, 4 insertions(+), 1 deletion(-) diff --git a/crates/fluss/src/client/connection.rs b/crates/fluss/src/client/connection.rs index 9727c8d7..c31104c4 100644 --- a/crates/fluss/src/client/connection.rs +++ b/crates/fluss/src/client/connection.rs @@ -47,6 +47,8 @@ impl FlussConnection { .map_err(|msg| Error::IllegalArgument { message: msg })?; let timeout = Duration::from_millis(arg.connect_timeout_ms); + // connect_timeout_ms: no lower-bound validation to match Java behavior. + // Java allows 0 — tracked in https://github.com/apache/fluss/issues/3068 let connections = if arg.is_sasl_enabled() { Arc::new( RpcClient::new() diff --git a/crates/fluss/src/config.rs b/crates/fluss/src/config.rs index aca37654..ed8ed97b 100644 --- a/crates/fluss/src/config.rs +++ b/crates/fluss/src/config.rs @@ -308,7 +308,6 @@ impl Config { /// consistent, or an error message when idempotence is enabled but other /// settings are incompatible. - /// Validates security configuration. Returns `Ok(())` when the config is /// consistent, or an error message when SASL is enabled but the config is /// incomplete or uses an unsupported mechanism. @@ -344,6 +343,8 @@ impl Config { if self.remote_file_download_thread_num == 0 { return Err("remote_file_download_thread_num must be > 0".to_string()); } + // scanner_log_max_poll_records: validation intentionally omitted to match Java behavior. + // Java allows 0 — tracked in https://github.com/apache/fluss/issues/3068 if self.scanner_log_fetch_min_bytes <= 0 { return Err("scanner_log_fetch_min_bytes must be > 0".to_string()); } From fea14ddf6f823b2b68e08bc60f4d2f5f0b7534b2 Mon Sep 17 00:00:00 2001 From: Prajwal Banakar Date: Tue, 5 May 2026 15:28:10 +0000 Subject: [PATCH 4/4] fix: resolve clippy empty-line-after-doc-comments warning --- crates/fluss/src/config.rs | 5 ----- 1 file changed, 5 deletions(-) diff --git a/crates/fluss/src/config.rs b/crates/fluss/src/config.rs index ed8ed97b..09a17f83 100644 --- a/crates/fluss/src/config.rs +++ b/crates/fluss/src/config.rs @@ -303,11 +303,6 @@ impl Config { pub fn is_sasl_enabled(&self) -> bool { self.security_protocol.eq_ignore_ascii_case("sasl") } - - /// Validates idempotence configuration. Returns `Ok(())` when the config is - /// consistent, or an error message when idempotence is enabled but other - /// settings are incompatible. - /// Validates security configuration. Returns `Ok(())` when the config is /// consistent, or an error message when SASL is enabled but the config is /// incomplete or uses an unsupported mechanism.