This is an automated email from the ASF dual-hosted git repository. numinnex pushed a commit to branch rework_topic_commands in repository https://gitbox.apache.org/repos/asf/iggy.git
commit 6057a1e106aa9048c22e4763d90c67900a75c4ee Author: Grzegorz Koszyk <[email protected]> AuthorDate: Fri Aug 14 10:40:56 2026 +0200 further refinements --- core/ai/mcp/src/service/mod.rs | 25 +-- .../src/requests/topics/update_topic.rs | 38 +---- .../src/responses/streams/get_stream.rs | 10 +- .../src/responses/topics/get_topic.rs | 6 +- .../cli/src/commands/binary_topics/update_topic.rs | 21 ++- core/common/src/http/topics/update_topic.rs | 29 ++-- core/common/src/traits/binary_impls/topics.rs | 11 +- core/common/src/traits/topic_client.rs | 8 +- core/common/src/types/options/mod.rs | 183 +++++++++++++-------- core/integration/tests/sdk/options.rs | 132 ++++++++++++++- .../server/scenarios/authentication_scenario.rs | 3 - .../tests/server/scenarios/permissions_scenario.rs | 9 - .../tests/server/scenarios/system_scenario.rs | 10 +- .../tests/server/topic_admission_vsr.rs | 24 +-- core/metadata/src/impls/metadata.rs | 18 +- core/metadata/src/stm/stream.rs | 34 ++-- core/partitions/src/messages_writer.rs | 26 ++- .../sdk/src/client_wrappers/binary_topic_client.rs | 56 +------ core/sdk/src/clients/binary_topics.rs | 16 +- core/sdk/src/http/topics.rs | 11 +- core/server/src/http/handlers.rs | 8 +- core/simulator/src/client.rs | 3 - foreign/cpp/src/client.rs | 27 +-- .../csharp/Iggy_SDK/Contracts/Tcp/TcpContracts.cs | 35 ++-- foreign/csharp/Iggy_SDK/Mappers/BinaryMapper.cs | 85 ++++++++-- foreign/go/internal/command/topic.go | 31 ++-- foreign/go/internal/command/topic_test.go | 13 +- .../iggy/client/async/tcp/TopicsTcpClient.java | 12 +- .../node/src/wire/topic/update-topic.command.ts | 29 +++- foreign/python/apache_iggy.pyi | 5 +- foreign/python/src/client.rs | 32 ++-- 31 files changed, 573 insertions(+), 377 deletions(-) diff --git a/core/ai/mcp/src/service/mod.rs b/core/ai/mcp/src/service/mod.rs index 716c7b52c..e93e94224 100644 --- a/core/ai/mcp/src/service/mod.rs +++ b/core/ai/mcp/src/service/mod.rs @@ -200,24 +200,17 @@ impl IggyService { }): Parameters<UpdateTopic>, ) -> Result<CallToolResult, ErrorData> { self.permissions.ensure_update()?; - let compression_algorithm = compression_algorithm - .and_then(|ca| ca.parse().ok()) - .unwrap_or_default(); - let message_expiry = message_expiry - .and_then(|me| me.parse().ok()) - .unwrap_or_default(); - let max_size = max_size.and_then(|ms| ms.parse().ok()).unwrap_or_default(); + // Absent means "leave alone" now, so an unparseable value must not + // silently become a reset to the server default. + let options = TopicUpdateOptions { + compression_algorithm: compression_algorithm.and_then(|ca| ca.parse().ok()), + message_expiry: message_expiry.and_then(|me| me.parse().ok()), + max_topic_size: max_size.and_then(|ms| ms.parse().ok()), + ..TopicUpdateOptions::default() + }; request( self.client - .update_topic( - &id(&stream_id)?, - &id(&topic_id)?, - &name, - compression_algorithm, - message_expiry, - max_size, - &TopicUpdateOptions::default(), - ) + .update_topic(&id(&stream_id)?, &id(&topic_id)?, &name, &options) .await, ) } diff --git a/core/binary_protocol/src/requests/topics/update_topic.rs b/core/binary_protocol/src/requests/topics/update_topic.rs index 44077a052..902971b37 100644 --- a/core/binary_protocol/src/requests/topics/update_topic.rs +++ b/core/binary_protocol/src/requests/topics/update_topic.rs @@ -17,19 +17,21 @@ use crate::WireError; use crate::WireIdentifier; -use crate::codec::{WireDecode, WireEncode, read_u8, read_u64_le}; +use crate::codec::{WireDecode, WireEncode}; use crate::primitives::identifier::WireName; use crate::primitives::options::WireOptions; -use bytes::{BufMut, BytesMut}; +use bytes::BytesMut; /// `UpdateTopic` request. /// /// Wire format: -/// `[stream_id:WireIdentifier][topic_id:WireIdentifier][compression_algorithm:u8] -/// [message_expiry:u64_le][max_topic_size:u64_le][name_len:u8][name:N][options TLV to end]` +/// `[stream_id:WireIdentifier][topic_id:WireIdentifier][name_len:u8][name:N] +/// [options TLV to end]` /// -/// The options block mirrors `CreateTopic`'s and carries the same catalog, so -/// a knob added there is updatable here without another layout change. +/// Identity and the new name are the only fixed fields; every SETTING rides the +/// options block, which mirrors `CreateTopic`'s and carries the same catalog. +/// A knob added there is updatable here without another layout change, and no +/// setting has two homes to disagree between. /// /// Keys absent from the block are LEFT ALONE rather than reset to their /// defaults. A client built before a key existed cannot send it, so treating @@ -40,20 +42,14 @@ use bytes::{BufMut, BytesMut}; pub struct UpdateTopicRequest { pub stream_id: WireIdentifier, pub topic_id: WireIdentifier, - pub compression_algorithm: u8, - pub message_expiry: u64, - pub max_topic_size: u64, pub name: WireName, pub options: WireOptions, } -const FIXED_FIELDS_SIZE: usize = 1 + 8 + 8; // 17 bytes - impl WireEncode for UpdateTopicRequest { fn encoded_size(&self) -> usize { self.stream_id.encoded_size() + self.topic_id.encoded_size() - + FIXED_FIELDS_SIZE + self.name.encoded_size() + self.options.encoded_size() } @@ -61,9 +57,6 @@ impl WireEncode for UpdateTopicRequest { fn encode(&self, buf: &mut BytesMut) { self.stream_id.encode(buf); self.topic_id.encode(buf); - buf.put_u8(self.compression_algorithm); - buf.put_u64_le(self.message_expiry); - buf.put_u64_le(self.max_topic_size); self.name.encode(buf); self.options.encode(buf); } @@ -74,12 +67,6 @@ impl WireDecode for UpdateTopicRequest { let (stream_id, mut pos) = WireIdentifier::decode(buf)?; let (topic_id, consumed) = WireIdentifier::decode(&buf[pos..])?; pos += consumed; - let compression_algorithm = read_u8(buf, pos)?; - pos += 1; - let message_expiry = read_u64_le(buf, pos)?; - pos += 8; - let max_topic_size = read_u64_le(buf, pos)?; - pos += 8; let (name, name_consumed) = WireName::decode(&buf[pos..])?; pos += name_consumed; let options = WireOptions::from_slice(&buf[pos..])?; @@ -87,9 +74,6 @@ impl WireDecode for UpdateTopicRequest { Self { stream_id, topic_id, - compression_algorithm, - message_expiry, - max_topic_size, name, options, }, @@ -113,9 +97,6 @@ mod tests { UpdateTopicRequest { stream_id: WireIdentifier::numeric(1), topic_id: WireIdentifier::numeric(2), - compression_algorithm: 1, - message_expiry: 7200, - max_topic_size: 500_000, name: WireName::new("updated-topic").unwrap(), options: sample_options(), } @@ -135,9 +116,6 @@ mod tests { let req = UpdateTopicRequest { stream_id: WireIdentifier::named("stream-a").unwrap(), topic_id: WireIdentifier::named("topic-b").unwrap(), - compression_algorithm: 0, - message_expiry: 0, - max_topic_size: u64::MAX, name: WireName::new("new-name").unwrap(), options: WireOptions::empty(), }; diff --git a/core/binary_protocol/src/responses/streams/get_stream.rs b/core/binary_protocol/src/responses/streams/get_stream.rs index 6e6d53ed0..2b09afe71 100644 --- a/core/binary_protocol/src/responses/streams/get_stream.rs +++ b/core/binary_protocol/src/responses/streams/get_stream.rs @@ -149,12 +149,20 @@ impl WireEncode for GetStreamResponse { } } +/// Smallest a topic header can encode as: the fixed ids, timestamps, sizes and +/// counts, plus a name length and two length-prefixed option blocks. +const MIN_TOPIC_HEADER_SIZE: usize = 45; + impl WireDecode for GetStreamResponse { fn decode(buf: &[u8]) -> Result<(Self, usize), WireError> { let (stream, mut pos) = StreamResponse::decode(buf)?; // Count-driven: a topic element carries variable-length options // blocks, so "consume until the buffer ends" no longer delimits it. - let mut topics = Vec::with_capacity(stream.topics_count as usize); + let mut topics = Vec::with_capacity(crate::codec::bounded_capacity( + stream.topics_count as usize, + buf.len().saturating_sub(pos), + MIN_TOPIC_HEADER_SIZE, + )); for _ in 0..stream.topics_count { let (topic, consumed) = TopicHeader::decode(&buf[pos..])?; pos += consumed; diff --git a/core/binary_protocol/src/responses/topics/get_topic.rs b/core/binary_protocol/src/responses/topics/get_topic.rs index e567e0957..793800a3c 100644 --- a/core/binary_protocol/src/responses/topics/get_topic.rs +++ b/core/binary_protocol/src/responses/topics/get_topic.rs @@ -116,7 +116,11 @@ impl WireDecode for GetTopicResponse { // Count-driven so the element stays delimited even when embedded in a // larger payload; the header's variable-length options blocks removed // the old "everything after the header is partitions" property. - let mut partitions = Vec::with_capacity(topic.partitions_count as usize); + let mut partitions = Vec::with_capacity(crate::codec::bounded_capacity( + topic.partitions_count as usize, + buf.len().saturating_sub(pos), + PartitionResponse::FIXED_SIZE, + )); for _ in 0..topic.partitions_count { let (partition, consumed) = PartitionResponse::decode(&buf[pos..])?; pos += consumed; diff --git a/core/cli/src/commands/binary_topics/update_topic.rs b/core/cli/src/commands/binary_topics/update_topic.rs index 2afe74a41..ab9ce8c9b 100644 --- a/core/cli/src/commands/binary_topics/update_topic.rs +++ b/core/cli/src/commands/binary_topics/update_topic.rs @@ -27,6 +27,7 @@ use tracing::{Level, event}; pub struct UpdateTopicCmd { update_topic: UpdateTopic, + compression_algorithm: CompressionAlgorithm, message_expiry: IggyExpiry, max_topic_size: MaxTopicSize, options: TopicUpdateOptions, @@ -47,14 +48,20 @@ impl UpdateTopicCmd { stream_id, topic_id, name, - compression_algorithm, - message_expiry, - max_topic_size, + compression_algorithm: Some(compression_algorithm), + message_expiry: Some(message_expiry), + max_topic_size: Some(max_topic_size), options: BTreeMap::new(), }, + compression_algorithm, message_expiry, max_topic_size, - options: TopicUpdateOptions::default(), + options: TopicUpdateOptions { + compression_algorithm: Some(compression_algorithm), + message_expiry: Some(message_expiry), + max_topic_size: Some(max_topic_size), + ..TopicUpdateOptions::default() + }, } } } @@ -67,7 +74,7 @@ impl CliCommand for UpdateTopicCmd { async fn execute_cmd(&mut self, client: &dyn Client) -> anyhow::Result<(), anyhow::Error> { client - .update_topic(&self.update_topic.stream_id, &self.update_topic.topic_id, &self.update_topic.name, self.update_topic.compression_algorithm, self.message_expiry, self.max_topic_size, &self.options) + .update_topic(&self.update_topic.stream_id, &self.update_topic.topic_id, &self.update_topic.name, &self.options) .await .with_context(|| { format!( @@ -84,7 +91,7 @@ impl CliCommand for UpdateTopicCmd { self.update_topic.topic_id, self.update_topic.name, self.message_expiry, - self.update_topic.compression_algorithm, + self.compression_algorithm, self.max_topic_size, self.update_topic.stream_id, ); @@ -97,7 +104,7 @@ impl fmt::Display for UpdateTopicCmd { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { let topic_id = &self.update_topic.topic_id; let topic_name = &self.update_topic.name; - let compression_algorithm = &self.update_topic.compression_algorithm; + let compression_algorithm = &self.compression_algorithm; let message_expiry = &self.message_expiry; let max_topic_size = &self.max_topic_size; let stream_id = &self.update_topic.stream_id; diff --git a/core/common/src/http/topics/update_topic.rs b/core/common/src/http/topics/update_topic.rs index 5f457cb28..7208c88b9 100644 --- a/core/common/src/http/topics/update_topic.rs +++ b/core/common/src/http/topics/update_topic.rs @@ -29,10 +29,11 @@ use std::collections::BTreeMap; /// It has additional payload: /// - `stream_id` - unique stream ID (numeric or name). /// - `topic_id` - unique topic ID (numeric or name). -/// - `message_expiry` - message expiry, if `NeverExpire` then messages will never expire. -/// - `max_topic_size` - maximum size of the topic in bytes, if `Unlimited` then topic size is unlimited. -/// Can't be lower than segment size in the config. /// - `name` - unique topic name, max length is 255 characters. +/// - `compression_algorithm`, `message_expiry`, `max_topic_size` - omit a field +/// to leave the topic's current value alone. Named here for REST ergonomics; +/// the server folds them into the same option keys the binary protocol uses, +/// so there is still one source per setting. /// - `options` - additional option keys as strings; only keys the update path /// accepts are allowed. #[derive(Debug, Serialize, Deserialize, PartialEq, Clone)] @@ -43,13 +44,15 @@ pub struct UpdateTopic { /// Unique topic ID (numeric or name). #[serde(skip)] pub topic_id: Identifier, - /// Compression algorithm for the topic. - pub compression_algorithm: CompressionAlgorithm, - /// Message expiry, if `NeverExpire` then messages will never expire. - pub message_expiry: IggyExpiry, - /// Max topic size, if `Unlimited` then topic size is unlimited. - /// Can't be lower than segment size in the config. - pub max_topic_size: MaxTopicSize, + /// Compression algorithm; omit to leave the current one alone. + #[serde(default)] + pub compression_algorithm: Option<CompressionAlgorithm>, + /// Message expiry; omit to leave the current one alone. + #[serde(default)] + pub message_expiry: Option<IggyExpiry>, + /// Max topic size; omit to leave the current one alone. + #[serde(default)] + pub max_topic_size: Option<MaxTopicSize>, /// Unique topic name, max length is 255 characters. pub name: String, /// Additional topic options as string key-values. Restricted to the keys @@ -63,9 +66,9 @@ impl Default for UpdateTopic { UpdateTopic { stream_id: Identifier::default(), topic_id: Identifier::default(), - compression_algorithm: Default::default(), - message_expiry: IggyExpiry::NeverExpire, - max_topic_size: MaxTopicSize::ServerDefault, + compression_algorithm: None, + message_expiry: None, + max_topic_size: None, name: "topic".to_string(), options: BTreeMap::new(), } diff --git a/core/common/src/traits/binary_impls/topics.rs b/core/common/src/traits/binary_impls/topics.rs index 9c9e5882f..976c8b9d1 100644 --- a/core/common/src/traits/binary_impls/topics.rs +++ b/core/common/src/traits/binary_impls/topics.rs @@ -18,9 +18,8 @@ use crate::traits::binary_auth::fail_if_not_authenticated; use crate::wire_conversions::{identifier_to_wire, topics_from_wire}; use crate::{ - BinaryClient, CompressionAlgorithm, DEFAULT_PARTITIONS_COUNT, Identifier, IggyError, - IggyExpiry, MaxTopicSize, Topic, TopicClient, TopicCreateOptions, TopicDetails, - TopicUpdateOptions, + BinaryClient, DEFAULT_PARTITIONS_COUNT, Identifier, IggyError, Topic, TopicClient, + TopicCreateOptions, TopicDetails, TopicUpdateOptions, }; use iggy_binary_protocol::WireName; use iggy_binary_protocol::codec::WireEncode; @@ -111,9 +110,6 @@ impl<B: BinaryClient> TopicClient for B { stream_id: &Identifier, topic_id: &Identifier, name: &str, - compression_algorithm: CompressionAlgorithm, - message_expiry: IggyExpiry, - max_topic_size: MaxTopicSize, options: &TopicUpdateOptions, ) -> Result<(), IggyError> { fail_if_not_authenticated(self).await?; @@ -125,9 +121,6 @@ impl<B: BinaryClient> TopicClient for B { UpdateTopicRequest { stream_id: wire_stream_id, topic_id: wire_topic_id, - compression_algorithm: compression_algorithm.as_code(), - message_expiry: u64::from(message_expiry), - max_topic_size: u64::from(max_topic_size), name: wire_name, options: options.to_wire()?, } diff --git a/core/common/src/traits/topic_client.rs b/core/common/src/traits/topic_client.rs index 54399fcb7..930fa1ead 100644 --- a/core/common/src/traits/topic_client.rs +++ b/core/common/src/traits/topic_client.rs @@ -15,10 +15,7 @@ // specific language governing permissions and limitations // under the License. -use crate::{ - CompressionAlgorithm, Identifier, IggyError, IggyExpiry, MaxTopicSize, Topic, - TopicCreateOptions, TopicDetails, TopicUpdateOptions, -}; +use crate::{Identifier, IggyError, Topic, TopicCreateOptions, TopicDetails, TopicUpdateOptions}; use async_trait::async_trait; /// This trait defines the methods to interact with the topic module. @@ -57,9 +54,6 @@ pub trait TopicClient { stream_id: &Identifier, topic_id: &Identifier, name: &str, - compression_algorithm: CompressionAlgorithm, - message_expiry: IggyExpiry, - max_topic_size: MaxTopicSize, options: &TopicUpdateOptions, ) -> Result<(), IggyError>; /// Delete a topic by unique ID or name. diff --git a/core/common/src/types/options/mod.rs b/core/common/src/types/options/mod.rs index 852004ad8..45428723b 100644 --- a/core/common/src/types/options/mod.rs +++ b/core/common/src/types/options/mod.rs @@ -334,19 +334,23 @@ pub const TOPIC_OPTION_KEYS: &[&str] = &[ /// The subset of [`TOPIC_OPTION_KEYS`] an `UpdateTopic` options block may /// carry. /// -/// Empty. Two separate reasons keep every key out: +/// These three used to be fixed fields of the update command, which meant one +/// setting had two homes and an update always rewrote all three whether the +/// caller meant to or not. As options they are patched: a key the client did +/// not send keeps its current value. /// -/// * `compression_algorithm`, `message_expiry` and `max_topic_size` have -/// dedicated fixed fields on the update layout. Accepting them in the block -/// too would mean two sources for one setting and a precedence rule nobody -/// can guess, so sending them here is an error rather than a silent -/// last-writer-wins. -/// * The partition runtime knobs (`segment_size`, `enforce_fsync`, both flush -/// thresholds, `preallocate_segments`) are pushed to partitions when the -/// topic is built. Nothing re-pushes them on update, so accepting one would -/// store a value the partitions never see -- a knob that reads as applied -/// and is not. They stay create-only until that propagation exists. -pub const UPDATABLE_TOPIC_OPTION_KEYS: &[&str] = &[]; +/// The partition runtime knobs (`segment_size`, `enforce_fsync`, both flush +/// thresholds, `preallocate_segments`) stay out, and not only because nothing +/// re-pushes them to a live partition. They describe how a partition's storage +/// was laid down: changing `segment_size` mid-segment leaves one segment sized +/// by the old cap and the next by the new one, and `preallocate_segments` can +/// only act on a file not yet opened. A topic gets them at creation and keeps +/// them, so its segments stay uniform. +pub const UPDATABLE_TOPIC_OPTION_KEYS: &[&str] = &[ + topic_option_keys::COMPRESSION_ALGORITHM, + topic_option_keys::MESSAGE_EXPIRY, + topic_option_keys::MAX_TOPIC_SIZE, +]; /// Keys an `UpdateStream` options block may carry. Empty because streams have /// no catalog keys yet: the block exists so the first one costs a catalog @@ -408,8 +412,10 @@ fn raw_options_to_wire(raw: &BTreeMap<String, String>) -> Result<WireOptions, Ig /// the option map, it does not replace it. See [`TopicUpdateOptions`] for why. #[derive(Debug, Clone, Default, PartialEq)] pub struct StreamUpdateOptions { - /// Keys sent as `String` values. Checked against - /// [`UPDATABLE_STREAM_OPTION_KEYS`] server-side. + /// Keys sent as `String` values, checked against + /// [`UPDATABLE_STREAM_OPTION_KEYS`] server-side. That list is empty, so + /// every key is currently refused by name; the field exists so the first + /// updatable stream key costs a catalog entry, not a wire change. pub raw: BTreeMap<String, String>, } @@ -428,8 +434,10 @@ impl StreamUpdateOptions { /// [`StreamUpdateOptions`]. #[derive(Debug, Clone, Default, PartialEq)] pub struct UserUpdateOptions { - /// Keys sent as `String` values. Checked against - /// [`UPDATABLE_USER_OPTION_KEYS`] server-side. + /// Keys sent as `String` values, checked against + /// [`UPDATABLE_USER_OPTION_KEYS`] server-side. That list is empty, so every + /// key is currently refused by name; the field exists so the first + /// updatable user key costs a catalog entry, not a wire change. pub raw: BTreeMap<String, String>, } @@ -497,19 +505,54 @@ impl TopicRuntimeOptions { /// value alone -- an update patches the option map, it does not replace it. #[derive(Debug, Clone, Default, PartialEq)] pub struct TopicUpdateOptions { - /// Keys sent as `String` values, for setting a key this build does not - /// know yet. Still checked against the updatable set server-side. + /// `None` leaves the topic's current algorithm alone. + pub compression_algorithm: Option<CompressionAlgorithm>, + /// `None` leaves the topic's current expiry alone. + pub message_expiry: Option<IggyExpiry>, + /// `None` leaves the topic's current cap alone. + pub max_topic_size: Option<MaxTopicSize>, + /// Keys sent as `String` values, checked against + /// [`UPDATABLE_TOPIC_OPTION_KEYS`] server-side. Lets a client reach an + /// updatable key added to the catalog after this build shipped. pub raw: BTreeMap<String, String>, } impl TopicUpdateOptions { - /// Encode the present keys into an options block. + /// Encode the present keys into an options block, in canonical kinds. + /// + /// A typed field is inserted after the raw entries, so it wins on collision + /// and keeps its canonical kind. /// /// # Errors /// - /// See [`raw_options_to_wire`]. + /// See [`raw_options_map`]. pub fn to_wire(&self) -> Result<WireOptions, IggyError> { - raw_options_to_wire(&self.raw) + let mut options = raw_options_map(&self.raw)?; + if let Some(compression_algorithm) = self.compression_algorithm { + options.insert( + HeaderKey::from_str(topic_option_keys::COMPRESSION_ALGORITHM) + .expect("catalog key is a valid header key"), + OptionValue::explicit( + HeaderValue::try_from(compression_algorithm.to_string().as_str()) + .expect("compression name fits a header value"), + ), + ); + } + if let Some(message_expiry) = self.message_expiry { + options.insert( + HeaderKey::from_str(topic_option_keys::MESSAGE_EXPIRY) + .expect("catalog key is a valid header key"), + OptionValue::explicit(HeaderValue::from(u64::from(message_expiry))), + ); + } + if let Some(max_topic_size) = self.max_topic_size { + options.insert( + HeaderKey::from_str(topic_option_keys::MAX_TOPIC_SIZE) + .expect("catalog key is a valid header key"), + OptionValue::explicit(HeaderValue::from(u64::from(max_topic_size))), + ); + } + crate::wire_conversions::resource_options_to_wire(&options, OptionsProvenance::All) } } @@ -896,41 +939,6 @@ impl TopicCreateOptions { } } -/// The option entries mirroring the three settings `UpdateTopic` carries as -/// fixed fields of the command rather than as option keys. -/// -/// An update writes both, so the stored map has to be rewritten alongside the -/// typed fields; leaving it alone makes `GetTopic` report the new expiry in -/// its field and the create-time one under `options`, with no way to tell -/// which is in force. -#[must_use] -pub fn topic_field_options( - compression_algorithm: CompressionAlgorithm, - message_expiry: IggyExpiry, - max_topic_size: MaxTopicSize, -) -> ResourceOptions { - ResourceOptions::from([ - ( - HeaderKey::from_str(topic_option_keys::COMPRESSION_ALGORITHM) - .expect("catalog key is a valid header key"), - OptionValue::explicit( - HeaderValue::try_from(compression_algorithm.to_string().as_str()) - .expect("compression name fits a header value"), - ), - ), - ( - HeaderKey::from_str(topic_option_keys::MESSAGE_EXPIRY) - .expect("catalog key is a valid header key"), - OptionValue::explicit(HeaderValue::from(u64::from(message_expiry))), - ), - ( - HeaderKey::from_str(topic_option_keys::MAX_TOPIC_SIZE) - .expect("catalog key is a valid header key"), - OptionValue::explicit(HeaderValue::from(u64::from(max_topic_size))), - ), - ]) -} - /// What a parse does with an entry this build cannot interpret: a key outside /// [`TOPIC_OPTION_KEYS`], or a catalog key whose value kind or payload does /// not parse. @@ -998,10 +1006,15 @@ fn parse_byte_size(entry: &WireUserHeaderEntry<'_>, key: &str) -> Result<u64, Ig fn parse_bool(entry: &WireUserHeaderEntry<'_>, key: &str) -> Result<bool, IggyError> { if entry.value_kind.0 == HeaderKind::Bool.as_code() { - return match entry.value.first() { - Some(0) => Ok(false), - Some(_) => Ok(true), - None => Err(IggyError::InvalidOptionValue(key.to_string())), + // Exactly one byte of 0 or 1. Anything looser admits a value that + // `HeaderValue::as_bool` later refuses to read, so the stored map would + // hold an entry its own public accessor rejects; a multi-byte payload + // would pass this gate and only fail at apply, turning a client + // mistake into a committed state-machine rejection. + return match entry.value { + [0] => Ok(false), + [1] => Ok(true), + _ => Err(IggyError::InvalidOptionValue(key.to_string())), }; } if entry.value_kind.0 == HeaderKind::String.as_code() { @@ -1018,12 +1031,11 @@ fn parse_compression( key: &str, ) -> Result<CompressionAlgorithm, IggyError> { if entry.value_kind.0 == HeaderKind::Uint8.as_code() { - let code = entry - .value - .first() - .copied() - .ok_or_else(|| IggyError::InvalidOptionValue(key.to_string()))?; - return CompressionAlgorithm::from_code(code) + // Exactly one byte, for the same reason as `parse_bool`. + let [code] = entry.value else { + return Err(IggyError::InvalidOptionValue(key.to_string())); + }; + return CompressionAlgorithm::from_code(*code) .map_err(|_| IggyError::InvalidOptionValue(key.to_string())); } if entry.value_kind.0 == HeaderKind::String.as_code() { @@ -1126,6 +1138,45 @@ mod tests { assert_eq!(parsed.segment_size, None); } + #[test] + fn sentinel_zeros_do_not_survive_re_encoding() { + // A client may send 0 to mean "resolve the default". Parsing normalizes + // it to absent, so admission puts the resolved value in the derived + // block -- but apply merges with explicit winning. If the literal 0 + // were still in the explicit block it would land back on top as the + // stored effective value, and a restart would re-parse it to absent and + // fall back to the node default rather than the value resolved at + // creation. Re-encoding from the parse is what drops it. + let sent = TopicCreateOptions { + message_expiry: Some(IggyExpiry::from(0u64)), + max_topic_size: Some(MaxTopicSize::from(0u64)), + segment_size: Some(IggyByteSize::from(0u64)), + size_of_messages_required_to_save: Some(IggyByteSize::from(0u64)), + enforce_fsync: Some(true), + ..TopicCreateOptions::default() + }; + let parsed = TopicCreateOptions::parse(&sent.to_wire().unwrap()).unwrap(); + let re_encoded = parsed.to_wire().unwrap(); + let stored = + crate::wire_conversions::resource_options_from_wire(&re_encoded, true).unwrap(); + + for key in [ + topic_option_keys::MESSAGE_EXPIRY, + topic_option_keys::MAX_TOPIC_SIZE, + topic_option_keys::SEGMENT_SIZE, + topic_option_keys::SIZE_OF_MESSAGES_REQUIRED_TO_SAVE, + ] { + assert!( + !stored.contains_key(&HeaderKey::from_str(key).unwrap()), + "{key} sentinel must not be persisted as an explicit value" + ); + } + // A non-sentinel key alongside them still rides through untouched. + assert!( + stored.contains_key(&HeaderKey::from_str(topic_option_keys::ENFORCE_FSYNC).unwrap()) + ); + } + #[test] fn runtime_options_derive_from_the_persisted_map() { let options = TopicCreateOptions { diff --git a/core/integration/tests/sdk/options.rs b/core/integration/tests/sdk/options.rs index 5ca6c5869..16b46fad5 100644 --- a/core/integration/tests/sdk/options.rs +++ b/core/integration/tests/sdk/options.rs @@ -291,14 +291,12 @@ async fn given_update_options_when_updating_topic_should_patch_not_replace(harne &stream, &topic, "update-topic", - CompressionAlgorithm::None, - IggyExpiry::NeverExpire, - MaxTopicSize::ServerDefault, &TopicUpdateOptions { raw: BTreeMap::from([( topic_option_keys::SEGMENT_SIZE.to_string(), "2MiB".to_string(), )]), + ..TopicUpdateOptions::default() }, ) .await; @@ -309,9 +307,6 @@ async fn given_update_options_when_updating_topic_should_patch_not_replace(harne &stream, &topic, "update-topic", - CompressionAlgorithm::None, - IggyExpiry::NeverExpire, - MaxTopicSize::ServerDefault, &TopicUpdateOptions::default(), ) .await @@ -399,3 +394,128 @@ async fn given_unknown_key_when_updating_stream_or_user_should_reject(harness: & .await .unwrap(); } + +#[iggy_harness] +async fn given_sentinel_zeros_when_creating_topic_should_report_resolved_defaults( + harness: &TestHarness, +) { + let client = harness.root_client().await.unwrap(); + client.create_stream("sentinel-stream").await.unwrap(); + let stream = Identifier::named("sentinel-stream").unwrap(); + + // 0 means "resolve the server default", not "expire immediately" / "no + // space". Admission normalizes it away, so the stored map must report the + // resolved value as derived -- never a literal 0 marked explicit, which is + // what a client would then read back as the effective configuration. + client + .create_topic( + &stream, + "sentinel-topic", + &TopicCreateOptions { + partitions_count: Some(1), + message_expiry: Some(IggyExpiry::from(0u64)), + max_topic_size: Some(MaxTopicSize::from(0u64)), + segment_size: Some(IggyByteSize::from(0u64)), + ..TopicCreateOptions::default() + }, + ) + .await + .unwrap(); + + let details = client + .get_topic(&stream, &Identifier::named("sentinel-topic").unwrap()) + .await + .unwrap() + .expect("topic exists"); + + for key in [ + topic_option_keys::MESSAGE_EXPIRY, + topic_option_keys::MAX_TOPIC_SIZE, + topic_option_keys::SEGMENT_SIZE, + ] { + let option = details + .options + .get(&HeaderKey::from_str(key).unwrap()) + .unwrap_or_else(|| panic!("{key} resolves to a default rather than vanishing")); + assert!( + !option.explicit, + "{key} sentinel must resolve to a derived default, not persist as explicit" + ); + assert_ne!( + option.value.as_bytes(), + 0u64.to_le_bytes(), + "{key} must report the resolved value, not the sentinel" + ); + } +} + +#[iggy_harness] +async fn given_rename_only_when_updating_topic_should_leave_settings_alone(harness: &TestHarness) { + let client = harness.root_client().await.unwrap(); + client.create_stream("patch-stream").await.unwrap(); + let stream = Identifier::named("patch-stream").unwrap(); + + client + .create_topic( + &stream, + "patch-topic", + &TopicCreateOptions { + partitions_count: Some(1), + compression_algorithm: Some(CompressionAlgorithm::Gzip), + message_expiry: Some(IggyExpiry::from(5_000_000u64)), + ..TopicCreateOptions::default() + }, + ) + .await + .unwrap(); + let topic = Identifier::named("patch-topic").unwrap(); + + // Settings live only in the options block, so an update that carries none + // of them is a pure rename. While they were fixed fields of the command + // this was impossible: every update rewrote all three whether or not the + // caller meant to. + client + .update_topic( + &stream, + &topic, + "patch-renamed", + &TopicUpdateOptions::default(), + ) + .await + .unwrap(); + + let details = client + .get_topic(&stream, &Identifier::named("patch-renamed").unwrap()) + .await + .unwrap() + .expect("topic exists"); + assert_eq!(details.name, "patch-renamed"); + assert_eq!(details.compression_algorithm, CompressionAlgorithm::Gzip); + assert_eq!(details.message_expiry, IggyExpiry::from(5_000_000u64)); + + // Sending one key changes that key and leaves the other alone. + client + .update_topic( + &stream, + &Identifier::named("patch-renamed").unwrap(), + "patch-renamed", + &TopicUpdateOptions { + message_expiry: Some(IggyExpiry::from(9_000_000u64)), + ..TopicUpdateOptions::default() + }, + ) + .await + .unwrap(); + + let details = client + .get_topic(&stream, &Identifier::named("patch-renamed").unwrap()) + .await + .unwrap() + .expect("topic exists"); + assert_eq!(details.message_expiry, IggyExpiry::from(9_000_000u64)); + assert_eq!( + details.compression_algorithm, + CompressionAlgorithm::Gzip, + "a key the update did not carry keeps its value" + ); +} diff --git a/core/integration/tests/server/scenarios/authentication_scenario.rs b/core/integration/tests/server/scenarios/authentication_scenario.rs index 53f173911..2736ea56a 100644 --- a/core/integration/tests/server/scenarios/authentication_scenario.rs +++ b/core/integration/tests/server/scenarios/authentication_scenario.rs @@ -234,9 +234,6 @@ async fn test_all_commands_require_auth(client: &IggyClient) { &ctx.stream_id, &ctx.topic_id, "x", - CompressionAlgorithm::None, - IggyExpiry::NeverExpire, - MaxTopicSize::ServerDefault, &TopicUpdateOptions::default(), ) .await diff --git a/core/integration/tests/server/scenarios/permissions_scenario.rs b/core/integration/tests/server/scenarios/permissions_scenario.rs index 798dcfb9b..6c27ff21c 100644 --- a/core/integration/tests/server/scenarios/permissions_scenario.rs +++ b/core/integration/tests/server/scenarios/permissions_scenario.rs @@ -578,9 +578,6 @@ async fn test_topic_permissions(harness: &TestHarness, root_client: &IggyClient) &stream_id, &topic_id, "new-name", - CompressionAlgorithm::None, - IggyExpiry::NeverExpire, - MaxTopicSize::ServerDefault, &TopicUpdateOptions::default(), ) .await, @@ -626,9 +623,6 @@ async fn test_topic_permissions(harness: &TestHarness, root_client: &IggyClient) &stream_id, &Identifier::named("temp-topic").unwrap(), "temp-topic-v2", - CompressionAlgorithm::None, - IggyExpiry::NeverExpire, - MaxTopicSize::ServerDefault, &TopicUpdateOptions::default(), ) .await @@ -1560,9 +1554,6 @@ async fn test_stream_permission_inheritance(harness: &TestHarness, root_client: &stream_id, &Identifier::named("temp-manage-stream-topic").unwrap(), "temp-manage-stream-topic-v2", - CompressionAlgorithm::None, - IggyExpiry::NeverExpire, - MaxTopicSize::ServerDefault, &TopicUpdateOptions::default(), ) .await diff --git a/core/integration/tests/server/scenarios/system_scenario.rs b/core/integration/tests/server/scenarios/system_scenario.rs index ab6d3d9b9..d8264af11 100644 --- a/core/integration/tests/server/scenarios/system_scenario.rs +++ b/core/integration/tests/server/scenarios/system_scenario.rs @@ -601,10 +601,12 @@ pub async fn run(harness: &TestHarness) { &Identifier::named(STREAM_NAME).unwrap(), &Identifier::named(TOPIC_NAME).unwrap(), &updated_topic_name, - CompressionAlgorithm::Gzip, - IggyExpiry::ExpireDuration(message_expiry_duration), - updated_max_topic_size, - &TopicUpdateOptions::default(), + &TopicUpdateOptions { + compression_algorithm: Some(CompressionAlgorithm::Gzip), + message_expiry: Some(IggyExpiry::ExpireDuration(message_expiry_duration)), + max_topic_size: Some(updated_max_topic_size), + ..TopicUpdateOptions::default() + }, ) .await .unwrap(); diff --git a/core/integration/tests/server/topic_admission_vsr.rs b/core/integration/tests/server/topic_admission_vsr.rs index 97e84a750..80f1ed9c4 100644 --- a/core/integration/tests/server/topic_admission_vsr.rs +++ b/core/integration/tests/server/topic_admission_vsr.rs @@ -147,10 +147,11 @@ async fn given_updated_topic_when_getting_topic_should_echo_stored_values(harnes stream_id, topic_id, "echo-topic", - CompressionAlgorithm::None, - message_expiry, - max_topic_size, - &TopicUpdateOptions::default(), + &TopicUpdateOptions { + message_expiry: Some(message_expiry), + max_topic_size: Some(max_topic_size), + ..TopicUpdateOptions::default() + }, ) .await .expect("update topic"); @@ -163,15 +164,16 @@ async fn given_updated_topic_when_getting_topic_should_echo_stored_values(harnes } }; - // The server echoes both stored sentinels as wire 0 (legacy parity). The - // SDK decodes a topic-response size 0 as `ServerDefault` but an expiry 0 - // as `NeverExpire` (`wire_conversions`), so that is the legacy-identical - // client-visible read-back; the node default must NOT leak into either. + // Settings ride the options block and 0 is its "resolve the default" + // sentinel, so a `ServerDefault` on update carries no key at all: the topic + // keeps what it already had. Resetting a setting back to the node default + // is deliberately not expressible -- an update states the values it wants, + // and everything it omits survives. + let created_size = MaxTopicSize::Custom(IggyByteSize::from_str("2GiB").expect("byte size")); assert_eq!( update_topic(MaxTopicSize::ServerDefault, IggyExpiry::ServerDefault).await, - (MaxTopicSize::ServerDefault, IggyExpiry::NeverExpire), - "an update to ServerDefault must echo the stored sentinel, \ - not the node default frozen at update time" + (created_size, IggyExpiry::NeverExpire), + "a sentinel carries no key, so the value set at creation survives" ); let custom_size = MaxTopicSize::Custom(IggyByteSize::from_str("3GiB").expect("byte size")); let custom_expiry = IggyExpiry::ExpireDuration(IggyDuration::from_str("5s").expect("duration")); diff --git a/core/metadata/src/impls/metadata.rs b/core/metadata/src/impls/metadata.rs index 0a4af6ada..632af893a 100644 --- a/core/metadata/src/impls/metadata.rs +++ b/core/metadata/src/impls/metadata.rs @@ -3404,15 +3404,25 @@ where match header.operation { Operation::CreateTopic => { - let request = WireCreateTopicRequest::decode_from(body) + let mut request = WireCreateTopicRequest::decode_from(body) .map_err(|_| IggyError::InvalidCommand)?; // Resolve every absent catalog key against server config here, // at primary admission, so the replicated payload carries // concrete values and every replica commits the same state - // regardless of local config. The client's explicit block is - // forwarded verbatim; resolved defaults ride a separate + // regardless of local config. Resolved defaults ride a separate // derived block, preserving per-key provenance for `GetTopic`. let explicit = TopicCreateOptions::parse(&request.options)?; + // Re-encode the explicit block from the parse rather than + // forwarding the client's bytes. Parsing normalizes a zero + // sentinel to "absent", so the resolved value goes in the + // derived block -- but apply merges with explicit winning, so a + // forwarded literal `0` would land back on top as the stored + // effective value. `GetTopic` would report 0, and a restart + // would re-parse that 0 to absent and fall back to whatever the + // node default is by then, not the value resolved at creation. + // Re-encoding also canonicalizes kinds (a `"128MiB"` string + // becomes `Uint64`), so the stored map reads back uniformly. + request.options = explicit.to_wire()?; let resolved_segment_size = explicit .segment_size .unwrap_or_else(|| IggyByteSize::from(self.default_segment_size.get())); @@ -4544,7 +4554,7 @@ mod tests { .expect("create topic with assignments prepare must decode"); assert!( persisted.request.options.is_empty(), - "the client's explicit block is forwarded verbatim" + "a client that sent no options gets an empty explicit block" ); let derived = iggy_common::TopicCreateOptions::parse(&persisted.derived_options) .expect("derived block parses against the catalog"); diff --git a/core/metadata/src/stm/stream.rs b/core/metadata/src/stm/stream.rs index 22f038af2..db0d4e02b 100644 --- a/core/metadata/src/stm/stream.rs +++ b/core/metadata/src/stm/stream.rs @@ -63,7 +63,7 @@ use iggy_binary_protocol::{WireIdentifier, WireName}; use iggy_common::wire_conversions::{resource_options_from_wire, resource_options_to_wire_split}; use iggy_common::{ CompressionAlgorithm, IggyExpiry, IggyTimestamp, MaxTopicSize, PartitionStats, ResourceOptions, - StreamStats, TopicCreateOptions, TopicStats, topic_field_options, + StreamStats, TopicCreateOptions, TopicStats, }; use serde::{Deserialize, Serialize}; use server_common::sharding::IggyNamespace; @@ -1846,23 +1846,27 @@ impl StateHandler for UpdateTopicRequest { let Ok(updated_options) = resource_options_from_wire(&self.options, true) else { return ApplyReply::err(UpdateTopicResult::InvalidOptionValue); }; + // Read leniently, like every other committed op: a key this build does + // not know is skipped rather than failing an operation its peers + // accepted. + let updated = TopicCreateOptions::parse_committed(&self.options); stream.topic_index.remove(&topic.name); topic.name = new_name_arc.clone(); - topic.compression_algorithm = - CompressionAlgorithm::from_code(self.compression_algorithm).unwrap_or_default(); - topic.message_expiry = IggyExpiry::from(self.message_expiry); - topic.max_topic_size = MaxTopicSize::from(self.max_topic_size); - // The three settings the command carries as fixed fields are mirrored - // into the map, so a reader of `options` never sees a value the typed - // field has already moved past. - topic.options.extend(topic_field_options( - topic.compression_algorithm, - topic.message_expiry, - topic.max_topic_size, - )); - // Patch, never replace: keys the client did not send keep their - // current value, so a client that predates a key cannot erase it. + // Settings arrive only through the options block now, so the typed + // fields are a projection of it and cannot drift. Absent means absent: + // a client that sends just a rename leaves every setting alone, and one + // built before a key existed cannot erase it. + if let Some(compression_algorithm) = updated.compression_algorithm { + topic.compression_algorithm = compression_algorithm; + } + if let Some(message_expiry) = updated.message_expiry { + topic.message_expiry = message_expiry; + } + if let Some(max_topic_size) = updated.max_topic_size { + topic.max_topic_size = max_topic_size; + } + // Patch, never replace, for the stored map too. topic.options.extend(updated_options); stream.topic_index.insert(new_name_arc, topic_id); ApplyReply::ok(Bytes::new()) diff --git a/core/partitions/src/messages_writer.rs b/core/partitions/src/messages_writer.rs index 5956d7ca0..4ecaadf50 100644 --- a/core/partitions/src/messages_writer.rs +++ b/core/partitions/src/messages_writer.rs @@ -168,15 +168,23 @@ fn preallocate_file(file: &File, file_path: &str, len: u64) { return; }; - // Runs INLINE on the shard thread, deliberately. Shard executors disable - // the blocking fallback pool (`thread_pool_limit(0)` in - // `server_common::executor`), so `spawn_blocking` does not park a task -- - // compio panics the shard outright with "the thread pool is needed but no - // worker thread is running". `FALLOC_FL_KEEP_SIZE` is a metadata-only - // extent reservation, microseconds on a local filesystem, which is the - // deployment this reservation exists for. A network filesystem can make it - // block; there it is better to turn the topic's `preallocate_segments` - // option off than to reintroduce a pool the runtime does not have. + // Runs INLINE on the shard thread, deliberately. `server_common::executor` + // sets `thread_pool_limit(0)` on the shard proactor, so `spawn_blocking` + // has no worker to park a task on and compio panics the shard outright with + // "the thread pool is needed but no worker thread is running". (That limit + // is skipped on macOS/aarch64, where the pool does exist -- see the FIXME + // there -- so the panic is Linux-and-most-targets, not universal. This arm + // is Linux-only regardless.) + // + // The cost is acceptable only because of what this call is: a metadata-only + // extent reservation, microseconds on the local filesystems this option + // exists for, and an immediate `EOPNOTSUPP` where the filesystem cannot do + // it. Where it can genuinely block -- NFSv4.2 `ALLOCATE`, FUSE, a badly + // fragmented extent tree forcing a journal commit -- it stalls the whole + // core, not one partition, because nothing here yields. Preallocation is + // opt-in per topic at creation for that reason; on such a deployment, + // create topics without `preallocate_segments` rather than reintroducing a + // pool the shard runtime does not have. if let Err(error) = fallocate(file, FallocateFlags::FALLOC_FL_KEEP_SIZE, 0, len) { warn!( target: "iggy.partitions.storage", diff --git a/core/sdk/src/client_wrappers/binary_topic_client.rs b/core/sdk/src/client_wrappers/binary_topic_client.rs index 133ce727f..52c12a71f 100644 --- a/core/sdk/src/client_wrappers/binary_topic_client.rs +++ b/core/sdk/src/client_wrappers/binary_topic_client.rs @@ -19,8 +19,7 @@ use crate::client_wrappers::client_wrapper::ClientWrapper; use async_trait::async_trait; use iggy_common::TopicClient; use iggy_common::{ - CompressionAlgorithm, Identifier, IggyError, IggyExpiry, MaxTopicSize, Topic, - TopicCreateOptions, TopicDetails, TopicUpdateOptions, + Identifier, IggyError, Topic, TopicCreateOptions, TopicDetails, TopicUpdateOptions, }; #[async_trait] @@ -69,75 +68,32 @@ impl TopicClient for ClientWrapper { stream_id: &Identifier, topic_id: &Identifier, name: &str, - compression_algorithm: CompressionAlgorithm, - message_expiry: IggyExpiry, - max_topic_size: MaxTopicSize, options: &TopicUpdateOptions, ) -> Result<(), IggyError> { match self { ClientWrapper::Iggy(client) => { client - .update_topic( - stream_id, - topic_id, - name, - compression_algorithm, - message_expiry, - max_topic_size, - options, - ) + .update_topic(stream_id, topic_id, name, options) .await } ClientWrapper::Http(client) => { client - .update_topic( - stream_id, - topic_id, - name, - compression_algorithm, - message_expiry, - max_topic_size, - options, - ) + .update_topic(stream_id, topic_id, name, options) .await } ClientWrapper::Tcp(client) => { client - .update_topic( - stream_id, - topic_id, - name, - compression_algorithm, - message_expiry, - max_topic_size, - options, - ) + .update_topic(stream_id, topic_id, name, options) .await } ClientWrapper::Quic(client) => { client - .update_topic( - stream_id, - topic_id, - name, - compression_algorithm, - message_expiry, - max_topic_size, - options, - ) + .update_topic(stream_id, topic_id, name, options) .await } ClientWrapper::WebSocket(client) => { client - .update_topic( - stream_id, - topic_id, - name, - compression_algorithm, - message_expiry, - max_topic_size, - options, - ) + .update_topic(stream_id, topic_id, name, options) .await } } diff --git a/core/sdk/src/clients/binary_topics.rs b/core/sdk/src/clients/binary_topics.rs index 190445c5a..60182eb60 100644 --- a/core/sdk/src/clients/binary_topics.rs +++ b/core/sdk/src/clients/binary_topics.rs @@ -20,8 +20,7 @@ use async_trait::async_trait; use iggy_common::TopicClient; use iggy_common::locking::IggyRwLockFn; use iggy_common::{ - CompressionAlgorithm, Identifier, IggyError, IggyExpiry, MaxTopicSize, Topic, - TopicCreateOptions, TopicDetails, TopicUpdateOptions, + Identifier, IggyError, Topic, TopicCreateOptions, TopicDetails, TopicUpdateOptions, }; #[async_trait] @@ -60,23 +59,12 @@ impl TopicClient for IggyClient { stream_id: &Identifier, topic_id: &Identifier, name: &str, - compression_algorithm: CompressionAlgorithm, - message_expiry: IggyExpiry, - max_topic_size: MaxTopicSize, options: &TopicUpdateOptions, ) -> Result<(), IggyError> { self.client .read() .await - .update_topic( - stream_id, - topic_id, - name, - compression_algorithm, - message_expiry, - max_topic_size, - options, - ) + .update_topic(stream_id, topic_id, name, options) .await } diff --git a/core/sdk/src/http/topics.rs b/core/sdk/src/http/topics.rs index 1d5e62470..365707b44 100644 --- a/core/sdk/src/http/topics.rs +++ b/core/sdk/src/http/topics.rs @@ -17,7 +17,7 @@ use crate::http::http_client::HttpClient; use crate::http::http_transport::HttpTransport; -use crate::prelude::{CompressionAlgorithm, Identifier, IggyError, IggyExpiry, MaxTopicSize}; +use crate::prelude::{Identifier, IggyError, IggyExpiry, MaxTopicSize}; use async_trait::async_trait; use iggy_common::TopicClient; use iggy_common::create_topic::CreateTopic; @@ -101,9 +101,6 @@ impl TopicClient for HttpClient { stream_id: &Identifier, topic_id: &Identifier, name: &str, - compression_algorithm: CompressionAlgorithm, - message_expiry: IggyExpiry, - max_topic_size: MaxTopicSize, options: &TopicUpdateOptions, ) -> Result<(), IggyError> { self.put( @@ -112,9 +109,9 @@ impl TopicClient for HttpClient { stream_id: stream_id.clone(), topic_id: topic_id.clone(), name: name.to_string(), - compression_algorithm, - message_expiry, - max_topic_size, + compression_algorithm: options.compression_algorithm, + message_expiry: options.message_expiry, + max_topic_size: options.max_topic_size, options: options.raw.clone(), }, ) diff --git a/core/server/src/http/handlers.rs b/core/server/src/http/handlers.rs index 7ebc9dc33..0668ac32c 100644 --- a/core/server/src/http/handlers.rs +++ b/core/server/src/http/handlers.rs @@ -910,7 +910,12 @@ pub(in crate::http) async fn update_topic( let stream_id = Identifier::from_str_value(&stream_id).map_err(WriteError::Rejected)?; let topic_id = Identifier::from_str_value(&topic_id).map_err(WriteError::Rejected)?; command.validate().map_err(WriteError::Rejected)?; + // The named JSON fields fold into the same option keys the binary protocol + // uses, so REST keeps its ergonomics without giving a setting two homes. let wire_options = TopicUpdateOptions { + compression_algorithm: command.compression_algorithm, + message_expiry: command.message_expiry, + max_topic_size: command.max_topic_size, raw: command.options, } .to_wire() @@ -922,9 +927,6 @@ pub(in crate::http) async fn update_topic( let request = UpdateTopicRequest { stream_id: identifier_to_wire(&stream_id).map_err(WriteError::Rejected)?, topic_id: identifier_to_wire(&topic_id).map_err(WriteError::Rejected)?, - compression_algorithm: command.compression_algorithm.as_code(), - message_expiry: command.message_expiry.into(), - max_topic_size: command.max_topic_size.into(), name: WireName::new(command.name) .map_err(|_| WriteError::Rejected(IggyError::InvalidTopicName))?, options: wire_options, diff --git a/core/simulator/src/client.rs b/core/simulator/src/client.rs index 20df2a069..db73cfb51 100644 --- a/core/simulator/src/client.rs +++ b/core/simulator/src/client.rs @@ -303,9 +303,6 @@ impl SimClient { let wire = UpdateTopicRequest { stream_id: WireIdentifier::named(stream).expect("stream name must be valid"), topic_id: WireIdentifier::named(topic).expect("topic name must be valid"), - compression_algorithm: 0, - message_expiry: 0, - max_topic_size: 0, name: WireName::new(new_name).expect("topic name must be valid"), options: WireOptions::empty(), }; diff --git a/foreign/cpp/src/client.rs b/foreign/cpp/src/client.rs index 39281cfbd..e53384590 100644 --- a/foreign/cpp/src/client.rs +++ b/foreign/cpp/src/client.rs @@ -533,10 +533,6 @@ impl Client { )); } }; - // Replication factor rides the options block now, not a fixed field. - let update_options = TopicUpdateOptions { - ..TopicUpdateOptions::default() - }; let rust_max_topic_size = match max_topic_size.as_str() { "" | "server_default" | "0" => RustMaxTopicSize::ServerDefault, _ => RustMaxTopicSize::from_str(&max_topic_size).map_err(|error| { @@ -546,17 +542,22 @@ impl Client { })?, }; + // Settings ride the options block; a server-default sentinel means the + // caller did not set the key, so the topic keeps its current value. + let update_options = TopicUpdateOptions { + compression_algorithm: (rust_compression_algorithm + != RustCompressionAlgorithm::default()) + .then_some(rust_compression_algorithm), + message_expiry: (rust_message_expiry != RustIggyExpiry::ServerDefault) + .then_some(rust_message_expiry), + max_topic_size: (rust_max_topic_size != RustMaxTopicSize::ServerDefault) + .then_some(rust_max_topic_size), + ..TopicUpdateOptions::default() + }; + RUNTIME.block_on(async { self.inner - .update_topic( - &rust_stream_id, - &rust_topic_id, - &topic_name, - rust_compression_algorithm, - rust_message_expiry, - rust_max_topic_size, - &update_options, - ) + .update_topic(&rust_stream_id, &rust_topic_id, &topic_name, &update_options) .await .map_err(|error| { format!( diff --git a/foreign/csharp/Iggy_SDK/Contracts/Tcp/TcpContracts.cs b/foreign/csharp/Iggy_SDK/Contracts/Tcp/TcpContracts.cs index 435c6f496..2da5395d3 100644 --- a/foreign/csharp/Iggy_SDK/Contracts/Tcp/TcpContracts.cs +++ b/foreign/csharp/Iggy_SDK/Contracts/Tcp/TcpContracts.cs @@ -642,23 +642,34 @@ internal static class TcpContracts internal static byte[] UpdateTopic(Identifier streamId, Identifier topicId, string name, CompressionAlgorithm compressionAlgorithm, ulong maxTopicSize, ulong messageExpiry) { - // No key may ride an update yet; the block exists so the first one - // costs a catalog entry rather than another wire change. + // Settings ride the options block. A default value means the caller did + // not set the key, so it is omitted and the server leaves the topic's + // current value alone. var options = new Dictionary<HeaderKey, HeaderValue>(); + if (compressionAlgorithm != CompressionAlgorithm.None) + { + options[HeaderKey.FromString("compression_algorithm")] + = HeaderValue.FromString(compressionAlgorithm.ToString().ToLowerInvariant()); + } + + if (messageExpiry != 0) + { + options[HeaderKey.FromString("message_expiry")] = HeaderValue.FromUInt64(messageExpiry); + } + + if (maxTopicSize != 0) + { + options[HeaderKey.FromString("max_topic_size")] = HeaderValue.FromUInt64(maxTopicSize); + } + var optionsLength = HeadersByteLength(options); Span<byte> bytes = - stackalloc byte[4 + streamId.Length + topicId.Length + 18 + name.Length + optionsLength]; + stackalloc byte[4 + streamId.Length + topicId.Length + 1 + name.Length + optionsLength]; bytes.WriteBytesFromStreamAndTopicIdentifiers(streamId, topicId); var position = 4 + streamId.Length + topicId.Length; - bytes[position] = (byte)compressionAlgorithm; - position += 1; - BinaryPrimitives.WriteUInt64LittleEndian(bytes[position..(position + 8)], - messageExpiry); - BinaryPrimitives.WriteUInt64LittleEndian(bytes[(position + 8)..(position + 16)], - maxTopicSize); - bytes[position + 16] = (byte)name.Length; - Encoding.UTF8.GetBytes(name, bytes[(position + 17)..(position + 17 + name.Length)]); - WriteHeadersTo(bytes[(position + 17 + name.Length)..], options); + bytes[position] = (byte)name.Length; + Encoding.UTF8.GetBytes(name, bytes[(position + 1)..(position + 1 + name.Length)]); + WriteHeadersTo(bytes[(position + 1 + name.Length)..], options); return bytes.ToArray(); } diff --git a/foreign/csharp/Iggy_SDK/Mappers/BinaryMapper.cs b/foreign/csharp/Iggy_SDK/Mappers/BinaryMapper.cs index 0b79df005..914c1af9a 100644 --- a/foreign/csharp/Iggy_SDK/Mappers/BinaryMapper.cs +++ b/foreign/csharp/Iggy_SDK/Mappers/BinaryMapper.cs @@ -748,27 +748,32 @@ internal static class BinaryMapper private static Dictionary<HeaderKey, HeaderValue> MapOptions(ReadOnlySpan<byte> payload, int position, out int readBytes) { - var optionsLength = (int)BinaryPrimitives.ReadUInt32LittleEndian(payload[position..(position + 4)]); - readBytes = 4 + optionsLength; + // Every length here is server-controlled. Read the block length as long + // so a value above int.MaxValue cannot wrap negative, and bound each + // entry against the block before slicing: an entry that overruns `end` + // would otherwise be accepted and silently consume the response bytes + // that follow the block. + var optionsLength = BinaryPrimitives.ReadUInt32LittleEndian(payload[position..(position + 4)]); + var available = (long)payload.Length - (position + 4); + if (optionsLength > available) + { + throw new MalformedResponseException( + $"Malformed options block at byte {position}: declared length {optionsLength} exceeds the " + + $"{available} bytes remaining in the payload."); + } + + readBytes = 4 + (int)optionsLength; var options = new Dictionary<HeaderKey, HeaderValue>(); var cursor = position + 4; - var end = position + 4 + optionsLength; + var end = cursor + (int)optionsLength; while (cursor < end) { - var keyKind = MapHeaderKind(payload[cursor]); - cursor += 1; - var keyLength = BinaryPrimitives.ReadInt32LittleEndian(payload[cursor..(cursor + 4)]); - cursor += 4; - var key = payload[cursor..(cursor + keyLength)].ToArray(); - cursor += keyLength; - - var valueKind = MapHeaderKind(payload[cursor]); - cursor += 1; - var valueLength = BinaryPrimitives.ReadInt32LittleEndian(payload[cursor..(cursor + 4)]); - cursor += 4; - var value = payload[cursor..(cursor + valueLength)].ToArray(); - cursor += valueLength; + var keyKind = MapHeaderKind(ReadOptionByte(payload, ref cursor, end, position)); + var key = ReadOptionField(payload, ref cursor, end, position, "key"); + + var valueKind = MapHeaderKind(ReadOptionByte(payload, ref cursor, end, position)); + var value = ReadOptionField(payload, ref cursor, end, position, "value"); options[new HeaderKey { @@ -781,9 +786,57 @@ internal static class BinaryMapper }; } + if (cursor != end) + { + throw new MalformedResponseException( + $"Malformed options block at byte {position}: entries ended at {cursor}, block ends at {end}."); + } + return options; } + private static byte ReadOptionByte(ReadOnlySpan<byte> payload, ref int cursor, int end, int blockStart) + { + if (cursor + 1 > end) + { + throw new MalformedResponseException( + $"Malformed options block at byte {blockStart}: entry kind runs past the end of the block."); + } + + var value = payload[cursor]; + cursor += 1; + return value; + } + + private static byte[] ReadOptionField(ReadOnlySpan<byte> payload, ref int cursor, int end, int blockStart, + string field) + { + if (cursor + 4 > end) + { + throw new MalformedResponseException( + $"Malformed options block at byte {blockStart}: {field} length runs past the end of the block."); + } + + var length = BinaryPrimitives.ReadUInt32LittleEndian(payload[cursor..(cursor + 4)]); + cursor += 4; + if (length is < 1 or > 255) + { + throw new MalformedResponseException( + $"Malformed options block at byte {blockStart}: {field} length {length} is outside 1..=255."); + } + + if (cursor + (int)length > end) + { + throw new MalformedResponseException( + $"Malformed options block at byte {blockStart}: {field} of {length} bytes runs past the end of " + + "the block."); + } + + var bytes = payload[cursor..(cursor + (int)length)].ToArray(); + cursor += (int)length; + return bytes; + } + internal static IReadOnlyList<StreamResponse> MapStreams(ReadOnlySpan<byte> payload) { List<StreamResponse> streams = new(); diff --git a/foreign/go/internal/command/topic.go b/foreign/go/internal/command/topic.go index 7eb4917c3..45140c9b7 100644 --- a/foreign/go/internal/command/topic.go +++ b/foreign/go/internal/command/topic.go @@ -133,12 +133,14 @@ func (d *DeleteTopic) MarshalBinary() ([]byte, error) { } type UpdateTopic struct { - StreamId iggcon.Identifier `json:"streamId"` - TopicId iggcon.Identifier `json:"topicId"` + StreamId iggcon.Identifier `json:"streamId"` + TopicId iggcon.Identifier `json:"topicId"` + Name string `json:"name"` + // Settings ride the options block. A zero value means "leave it alone", + // so a rename does not silently reset the rest. CompressionAlgorithm iggcon.CompressionAlgorithm `json:"compressionAlgorithm"` MessageExpiry iggcon.Duration `json:"messageExpiry"` MaxTopicSize uint64 `json:"maxTopicSize"` - Name string `json:"name"` } func (u *UpdateTopic) Code() Code { @@ -147,8 +149,20 @@ func (u *UpdateTopic) Code() Code { // options builds the trailing options block. Only keys an update may change // are allowed; the server rejects the create-time knobs by name. +// options builds the trailing options block. A zero value means the caller did +// not set the key, so it is omitted and the server leaves the topic's current +// value alone. func (u *UpdateTopic) options() []iggcon.HeaderEntry { var options []iggcon.HeaderEntry + if name := u.CompressionAlgorithm.String(); name != "none" { + options = append(options, stringOption(topicOptionCompressionAlgorithm, name)) + } + if u.MessageExpiry != 0 { + options = append(options, uint64Option(topicOptionMessageExpiry, uint64(u.MessageExpiry))) + } + if u.MaxTopicSize != 0 { + options = append(options, uint64Option(topicOptionMaxTopicSize, u.MaxTopicSize)) + } return options } @@ -163,22 +177,13 @@ func (u *UpdateTopic) MarshalBinary() ([]byte, error) { } optionsBytes := iggcon.GetHeadersBytes(u.options()) - buffer := make([]byte, 18+len(streamIdBytes)+len(topicIdBytes)+len(u.Name)+len(optionsBytes)) + buffer := make([]byte, 1+len(streamIdBytes)+len(topicIdBytes)+len(u.Name)+len(optionsBytes)) offset := 0 offset += copy(buffer[offset:], streamIdBytes) offset += copy(buffer[offset:], topicIdBytes) - buffer[offset] = byte(u.CompressionAlgorithm) - offset++ - - binary.LittleEndian.PutUint64(buffer[offset:], uint64(u.MessageExpiry)) - offset += 8 - - binary.LittleEndian.PutUint64(buffer[offset:], u.MaxTopicSize) - offset += 8 - buffer[offset] = uint8(len(u.Name)) offset++ diff --git a/foreign/go/internal/command/topic_test.go b/foreign/go/internal/command/topic_test.go index 439d9c89f..3363ff743 100644 --- a/foreign/go/internal/command/topic_test.go +++ b/foreign/go/internal/command/topic_test.go @@ -133,12 +133,17 @@ func TestSerialize_UpdateTopic(t *testing.T) { 0x01, // TopicId Kind (NumericId) 0x04, // TopicId Length (4) 0x01, 0x00, 0x00, 0x00, // TopicId Value (1) - 0x00, // compression algorithm - 0x64, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, // Message Expiry (100) - 0x64, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, // Max Topic Size (100) 0x0C, // Name Length (12) 0x75, 0x70, 0x64, 0x61, 0x74, 0x65, 0x5F, 0x74, 0x6F, 0x70, 0x69, 0x63, // Name ("update_topic") - // No options block: this update sets no option keys. + // Settings ride the options block; compression is None so it is omitted. + 0x02, 0x0E, 0x00, 0x00, 0x00, // key kind (String), key length (14) + 0x6D, 0x65, 0x73, 0x73, 0x61, 0x67, 0x65, 0x5F, 0x65, 0x78, 0x70, 0x69, 0x72, 0x79, // "message_expiry" + 0x0C, 0x08, 0x00, 0x00, 0x00, // value kind (Uint64), value length (8) + 0x64, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, // 100 + 0x02, 0x0E, 0x00, 0x00, 0x00, // key kind (String), key length (14) + 0x6D, 0x61, 0x78, 0x5F, 0x74, 0x6F, 0x70, 0x69, 0x63, 0x5F, 0x73, 0x69, 0x7A, 0x65, // "max_topic_size" + 0x0C, 0x08, 0x00, 0x00, 0x00, // value kind (Uint64), value length (8) + 0x64, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, // 100 } if !bytes.Equal(serialized1, expected) { diff --git a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/TopicsTcpClient.java b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/TopicsTcpClient.java index be4d0bcda..f2abaead0 100644 --- a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/TopicsTcpClient.java +++ b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/TopicsTcpClient.java @@ -175,14 +175,12 @@ public class TopicsTcpClient implements TopicsClient { var payload = Unpooled.buffer(); payload.writeBytes(toBytes(streamId)); payload.writeBytes(toBytes(topicId)); - payload.writeByte(compressionAlgorithm.asCode()); - payload.writeBytes(toBytesAsU64(messageExpiry)); - payload.writeBytes(toBytesAsU64(maxTopicSize)); payload.writeBytes(BytesSerializer.toBytes(name)); - // Only the updatable subset may ride an update; the server rejects the - // create-time knobs by name. - Map<HeaderKey, HeaderValue> options = new LinkedHashMap<>(); - payload.writeBytes(BytesSerializer.toBytes(options)); + // Settings ride the options block. A default value means the caller did + // not set the key, so it is omitted and the server leaves the topic's + // current value alone. + payload.writeBytes(BytesSerializer.toBytes( + createTopicOptions(compressionAlgorithm, messageExpiry, maxTopicSize))); return connection() .send(CommandCode.Topic.UPDATE.getValue(), payload) diff --git a/foreign/node/src/wire/topic/update-topic.command.ts b/foreign/node/src/wire/topic/update-topic.command.ts index 0ba0aa03a..955018e30 100644 --- a/foreign/node/src/wire/topic/update-topic.command.ts +++ b/foreign/node/src/wire/topic/update-topic.command.ts @@ -25,6 +25,7 @@ import { COMMAND_CODE } from '../command.code.js'; import { type CompressionAlgorithm as CompressionAlgorithmT, CompressionAlgorithm, + compressionAlgorithmName, isValidCompressionAlgorithm } from './topic.utils.js'; @@ -69,17 +70,29 @@ export const UPDATE_TOPIC = { if (bName.length < 1 || bName.length > 255) throw new Error('Topic name should be between 1 and 255 bytes'); if(!isValidCompressionAlgorithm(compressionAlgorithm)) - throw new Error(`createTopic: invalid compressionAlgorithm (${compressionAlgorithm})`); + throw new Error(`updateTopic: invalid compressionAlgorithm (${compressionAlgorithm})`); - // No key may ride an update yet; the block exists so the first one costs - // a catalog entry rather than another wire change. + // Settings ride the options block. A default value means the caller did not + // set the key, so it is omitted and the server leaves the current value be. const options: OptionEntry[] = []; + if (compressionAlgorithm !== CompressionAlgorithm.None) + options.push({ + key: 'compression_algorithm', + value: HeaderValue.String(compressionAlgorithmName(compressionAlgorithm)) + }); + if (messageExpiry !== 0n) + options.push({ + key: 'message_expiry', + value: HeaderValue.Uint64(messageExpiry) + }); + if (maxTopicSize !== 0n) + options.push({ + key: 'max_topic_size', + value: HeaderValue.Uint64(maxTopicSize) + }); - const b = Buffer.allocUnsafe(8 + 8 + 1 + 1); - b.writeUInt8(compressionAlgorithm, 0); - b.writeBigUInt64LE(messageExpiry, 1); // 0 is unlimited ??? - b.writeBigUInt64LE(maxTopicSize, 9); // optional, 0 is null - b.writeUInt8(bName.length, 17); + const b = Buffer.allocUnsafe(1); + b.writeUInt8(bName.length, 0); return Buffer.concat([ streamIdentifier, diff --git a/foreign/python/apache_iggy.pyi b/foreign/python/apache_iggy.pyi index 0a560c440..88bfbb5ec 100644 --- a/foreign/python/apache_iggy.pyi +++ b/foreign/python/apache_iggy.pyi @@ -1083,8 +1083,9 @@ class IggyClient: r""" Update an existing topic. - This is a full replacement: any optional parameter left unset is reset to - its server default rather than preserved. + A patch, not a replacement: every setting rides the options block, so a + field left unset keeps the topic's current value rather than resetting + it to a server default. Args: stream_id: Stream identifier as `str | int`. diff --git a/foreign/python/src/client.rs b/foreign/python/src/client.rs index dcd0b1cac..0ffc7822d 100644 --- a/foreign/python/src/client.rs +++ b/foreign/python/src/client.rs @@ -609,8 +609,9 @@ impl IggyClient { /// Update an existing topic. /// - /// This is a full replacement: any optional parameter left unset is reset to - /// its server default rather than preserved. + /// A patch, not a replacement: every setting rides the options block, so a + /// field left unset keeps the topic's current value rather than resetting + /// it to a server default. /// /// Args: /// stream_id: Stream identifier as `str | int`. @@ -647,25 +648,28 @@ impl IggyClient { &MaxTopicSize, >, ) -> PyResult<Bound<'a, PyAny>> { - let (compression_algorithm, expiry, max_size) = - resolve_topic_params(compression_algorithm, message_expiry, max_topic_size)?; + // Absent stays absent: a key the caller did not pass is left alone + // server-side rather than reset to a default. + let compression_algorithm = compression_algorithm + .map(|algo| { + CompressionAlgorithm::from_str(&algo) + .map_err(|e| PyErr::new::<pyo3::exceptions::PyRuntimeError, _>(e.to_string())) + }) + .transpose()?; + let update_options = TopicUpdateOptions { + compression_algorithm, + message_expiry: message_expiry.map(RustIggyExpiry::try_from).transpose()?, + max_topic_size: max_topic_size.map(RustMaxTopicSize::try_from).transpose()?, + ..TopicUpdateOptions::default() + }; let stream_id = Identifier::try_from(stream_id)?; let topic_id = Identifier::try_from(topic_id)?; let inner = self.inner.clone(); - let update_options = TopicUpdateOptions::default(); future_into_py(py, async move { inner - .update_topic( - &stream_id, - &topic_id, - &name, - compression_algorithm, - expiry, - max_size, - &update_options, - ) + .update_topic(&stream_id, &topic_id, &name, &update_options) .await .map_err(|e| PyErr::new::<pyo3::exceptions::PyRuntimeError, _>(e.to_string()))?; Ok(())
