This is an automated email from the ASF dual-hosted git repository. hubcio pushed a commit to branch fix/docs-audit-core in repository https://gitbox.apache.org/repos/asf/iggy.git
commit ac4955d4046f9eaea0e817f41c762e47edca5111 Author: Hubert Gruszecki <[email protected]> AuthorDate: Sat Sep 12 20:06:33 2026 +0200 fix(cli): reject invalid message keys before connecting Invalid keys reached connection and login before local validation. Use the SDK validator during clap parsing to keep the byte limits consistent and reject malformed keys before creating the client. --- core/cli/src/args/message.rs | 34 +++++++++++++++++++++- .../tests/cli/message/test_message_send_command.rs | 22 ++++++++------ 2 files changed, 46 insertions(+), 10 deletions(-) diff --git a/core/cli/src/args/message.rs b/core/cli/src/args/message.rs index 2c7dde6a9..f435fd610 100644 --- a/core/cli/src/args/message.rs +++ b/core/cli/src/args/message.rs @@ -83,7 +83,7 @@ pub(crate) struct SendMessagesArgs { /// /// The key must contain 1 to 255 bytes. Binary clients resolve the partition ID; HTTP resolves it on the server. #[clap(verbatim_doc_comment)] - #[clap(short, long, group = "partitioning")] + #[clap(short, long, value_parser = parse_message_key, group = "partitioning")] pub(crate) message_key: Option<String>, /// Messages to be sent /// @@ -116,6 +116,10 @@ pub(crate) struct SendMessagesArgs { pub(crate) input_file: Option<String>, } +fn parse_message_key(value: &str) -> Result<String, IggyError> { + Partitioning::messages_key_str(value).map(|_| value.to_owned()) +} + /// Parse Header Key, Kind and Value from the string separated by a ':' fn parse_key_val(s: &str) -> Result<(HeaderKey, HeaderValue), IggyError> { let parts = s.splitn(3, ':').collect::<Vec<_>>(); @@ -288,8 +292,36 @@ pub(crate) struct FlushMessagesArgs { #[cfg(test)] mod tests { use super::*; + use crate::args::{Command, IggyConsoleArgs}; + use clap::Parser; use std::str::FromStr; + #[test] + fn given_valid_message_key_when_sending_should_preserve_key_bytes() { + let max_key_bytes = usize::from(u8::MAX); + for key in [ + "x".to_owned(), + "x".repeat(max_key_bytes), + format!("{}x", "é".repeat(max_key_bytes / "é".len())), + ] { + let parsed = IggyConsoleArgs::try_parse_from([ + "iggy", + "message", + "send", + "--message-key", + &key, + "stream", + "topic", + "payload", + ]) + .unwrap(); + let Some(Command::Message(MessageAction::Send(args))) = parsed.command else { + panic!("Expected the message send command"); + }; + assert_eq!(args.message_key.as_deref(), Some(key.as_str())); + } + } + #[test] fn parse_key_val_should_parse_string() { let expected_value: &str = "value"; diff --git a/core/integration/tests/cli/message/test_message_send_command.rs b/core/integration/tests/cli/message/test_message_send_command.rs index fc2701eac..abde36aef 100644 --- a/core/integration/tests/cli/message/test_message_send_command.rs +++ b/core/integration/tests/cli/message/test_message_send_command.rs @@ -22,11 +22,11 @@ use assert_cmd::assert::Assert; use async_trait::async_trait; use iggy::prelude::defaults::{DEFAULT_ROOT_PASSWORD, DEFAULT_ROOT_USERNAME}; use iggy::prelude::*; -use integration::iggy_harness; use predicates::str::{contains, diff}; use serial_test::parallel; use std::collections::BTreeMap; use std::str::from_utf8; +use std::time::Duration; use twox_hash::XxHash32; #[derive(Debug)] @@ -464,22 +464,24 @@ Options: .await; } -#[iggy_harness] -async fn given_invalid_message_key_when_sending_should_return_error_without_panicking( - harness: &TestHarness, -) { - let server_address = harness.server().raw_tcp_addr().unwrap(); +#[test] +#[parallel] +fn given_invalid_message_key_when_sending_should_reject_before_connecting() { + const UNREACHABLE_SERVER_ADDRESS: &str = "127.0.0.1:0"; + const ARGUMENT_PARSE_TIMEOUT: Duration = Duration::from_secs(5); + let cli_home = tempfile::tempdir().unwrap(); let oversized_key = "x".repeat(usize::from(u8::MAX) + 1); + let oversized_unicode_key = "é".repeat(usize::from(u8::MAX) / "é".len() + 1); - for key in ["", oversized_key.as_str()] { + for key in ["", oversized_key.as_str(), oversized_unicode_key.as_str()] { #[allow(deprecated)] let mut command = assert_cmd::Command::cargo_bin("iggy").unwrap(); command .env("IGGY_HOME", cli_home.path()) .args([ "--tcp-server-address", - &server_address, + UNREACHABLE_SERVER_ADDRESS, "-u", DEFAULT_ROOT_USERNAME, "-p", @@ -492,8 +494,10 @@ async fn given_invalid_message_key_when_sending_should_return_error_without_pani "topic", "payload", ]) + .timeout(ARGUMENT_PARSE_TIMEOUT) .assert() - .code(1) + .code(2) + .stderr(contains("--message-key <MESSAGE_KEY>")) .stderr(contains("Invalid command")); } }
