From 635eed51568cee990511d5d50e814c823a8ff6b3 Mon Sep 17 00:00:00 2001 From: haubur Date: Thu, 20 Aug 2026 19:22:23 +0200 Subject: [PATCH 1/3] fix: respect max_buffer_size when merging batches --- core/sdk/src/clients/producer_sharding.rs | 159 ++++++++++++++++++++-- 1 file changed, 145 insertions(+), 14 deletions(-) diff --git a/core/sdk/src/clients/producer_sharding.rs b/core/sdk/src/clients/producer_sharding.rs index 668256f83d..c2db25b42e 100644 --- a/core/sdk/src/clients/producer_sharding.rs +++ b/core/sdk/src/clients/producer_sharding.rs @@ -110,16 +110,23 @@ impl Sizeable for ShardMessage { pub struct ShardMessageWithPermit { pub inner: ShardMessage, - _bytes_permit: Option, + bytes_permit: OwnedSemaphorePermit, } impl ShardMessageWithPermit { pub fn new(msg: ShardMessage, permit_bytes: OwnedSemaphorePermit) -> Self { Self { inner: msg, - _bytes_permit: Some(permit_bytes), + bytes_permit: permit_bytes, } } + + /// Takes over `other`'s messages together with its byte permit, so the buffered bytes stay + /// charged against the `max_buffer_size` budget until the merged batch has been written. + fn merge(&mut self, other: Self) { + self.inner.messages.extend(other.inner.messages); + self.bytes_permit.merge(other.bytes_permit); + } } pub struct Shard { @@ -212,6 +219,27 @@ impl Shard { } } + /// Drains the buffer into batches. + /// + /// If adjacent ShardMessages have the same destination (stream, topic, partition combination), + /// merge to avoid multiple sends. + /// The merge happens on both, the inner [`IggyMessage`] and the [`OwnedSemaphorePermit`]. + fn merge_batches(buffer: &mut Vec) -> Vec { + let mut merged_batches: Vec = Vec::with_capacity(buffer.len()); + // Since buffer is a mutable reference, the drain leaves the buffer intact but empty, + // such that it can be filled up again. + for msg in buffer.drain(..) { + if let Some(last) = merged_batches.last_mut() + && Self::same_destination(&last.inner, &msg.inner) + { + last.merge(msg); + continue; + } + merged_batches.push(msg); + } + merged_batches + } + async fn flush_buffer( core: &Arc, slots_permit: &Arc, @@ -223,18 +251,7 @@ impl Shard { return; } - let mut merged_batches: Vec = Vec::new(); - for msg in buffer.drain(..) { - if let Some(last) = merged_batches.last_mut() - && Self::same_destination(&last.inner, &msg.inner) - { - last.inner.messages.extend(msg.inner.messages); - continue; - } - merged_batches.push(msg); - } - - for msg in merged_batches { + for msg in Self::merge_batches(buffer) { let _slot_permit = slots_permit.acquire().await; let result = core @@ -312,6 +329,120 @@ mod tests { .unwrap() } + async fn charged_batch( + budget: &Arc, + stream: Arc, + topic: Arc, + payload_size: usize, + ) -> ShardMessageWithPermit { + let message = ShardMessage { + stream, + topic, + messages: vec![dummy_message(payload_size)], + partitioning: None, + }; + let permit = budget + .clone() + .acquire_many_owned(message.get_size_bytes().as_bytes_u32()) + .await + .unwrap(); + ShardMessageWithPermit::new(message, permit) + } + + #[tokio::test] + async fn test_merge_batches_keeps_permits_of_merged_batches_charged() { + let budget = Arc::new(Semaphore::new(10_000)); + let stream = dummy_identifier(); + let topic = dummy_identifier(); + + let mut buffer = Vec::new(); + for _ in 0..3 { + buffer.push(charged_batch(&budget, stream.clone(), topic.clone(), 10).await); + } + let charged = 10_000 - budget.available_permits(); + let merged = Shard::merge_batches(&mut buffer); + + // The original buffer should be drained. + assert!(buffer.is_empty()); + assert_eq!(merged.len(), 1); + assert_eq!(merged[0].inner.messages.len(), 3); + assert_eq!(budget.available_permits(), 10_000 - charged); + + // Dropping merged gives back to the semaphore. + drop(merged); + assert_eq!(budget.available_permits(), 10_000); + } + + #[tokio::test] + async fn test_shard_keeps_budget_charged_until_merged_batch_is_written() { + const BUDGET: usize = 10_000; + + // The two channels simulate the timing of a real write, which `flush_buffer` awaits at + // `core.send_internal`. Receiving on `write_started_rx` means the worker has entered that + // call, so the batch is on the wire and its bytes are still buffered from the producer's + // point of view. Sending on `release_write_tx` simulates that the call + // returns, the loop iteration ends and the batch drops and frees its permit. + let (write_started_tx, write_started_rx) = flume::unbounded::<()>(); + let (release_write_tx, release_write_rx) = flume::unbounded::<()>(); + + let mut mock = MockProducerCoreBackend::new(); + mock.expect_send_internal() + .times(1) + .returning(move |_, _, _, _| { + let write_started_tx = write_started_tx.clone(); + let release_write_rx = release_write_rx.clone(); + Box::pin(async move { + write_started_tx.send_async(()).await.unwrap(); + release_write_rx.recv_async().await.unwrap(); + Ok(no_confirmations()) + }) + }); + + let bb = BackgroundConfig::builder() + .batch_length(3) + .batch_size(0) + .linger_time(IggyDuration::new_from_secs(60)); + let config = Arc::new(bb.build()); + + let budget = Arc::new(Semaphore::new(BUDGET)); + let slots_permit = Arc::new(Semaphore::new(100)); + + let (_stop_tx, stop_rx) = broadcast::channel(1); + let shard = Shard::new( + Arc::new(mock), + config, + slots_permit, + flume::unbounded().0, + stop_rx, + ); + + let stream = dummy_identifier(); + let topic = dummy_identifier(); + for _ in 0..3 { + let batch = charged_batch(&budget, stream.clone(), topic.clone(), 100).await; + shard.send(batch).await.unwrap(); + } + let charged = BUDGET - budget.available_permits(); + assert!(charged > 0); + + // The three batches share a destination, so the flush merges them into a single write. + write_started_rx.recv_async().await.unwrap(); + assert_eq!( + budget.available_permits(), + BUDGET - charged, + "the merged batch must hold every permit it absorbed until the write completes" + ); + + release_write_tx.send_async(()).await.unwrap(); + tokio::time::timeout(Duration::from_secs(1), async { + while budget.available_permits() != BUDGET { + sleep(Duration::from_millis(5)).await; + } + }) + .await + .expect("the written batch must give its permits back"); + } + #[tokio::test] async fn test_shard_flushes_by_batch_length() { let mut mock = MockProducerCoreBackend::new(); From 9ff6faa3a0fc001dcd417ca35b0acbcbc5142461 Mon Sep 17 00:00:00 2001 From: spetz Date: Mon, 31 Aug 2026 13:00:07 +0200 Subject: [PATCH 2/3] improvements --- Cargo.lock | 14 +- Cargo.toml | 6 +- core/ai/mcp/Cargo.toml | 2 +- core/bench/Cargo.toml | 2 +- core/binary_protocol/Cargo.toml | 2 +- core/cli/Cargo.toml | 2 +- core/common/Cargo.toml | 2 +- core/connectors/runtime/Cargo.toml | 2 +- core/sdk/Cargo.toml | 2 +- core/sdk/src/clients/producer_config.rs | 3 +- core/sdk/src/clients/producer_dispatcher.rs | 146 +++++++++++----- core/sdk/src/clients/producer_sharding.rs | 179 ++++++++++++++------ foreign/python/Cargo.toml | 2 +- 13 files changed, 253 insertions(+), 111 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index c6aa065741..251f7f309a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6659,7 +6659,7 @@ checksum = "cd62e6b5e86ea8eeeb8db1de02880a6abc01a397b2ebb64b5d74ac255318f5cb" [[package]] name = "iggy" -version = "0.11.0-edge.5" +version = "0.11.0-edge.6" dependencies = [ "async-broadcast", "async-dropper", @@ -6693,7 +6693,7 @@ dependencies = [ [[package]] name = "iggy-bench" -version = "0.6.0-edge.5" +version = "0.6.0-edge.6" dependencies = [ "async-trait", "bench-report", @@ -6750,7 +6750,7 @@ dependencies = [ [[package]] name = "iggy-cli" -version = "0.14.0-edge.5" +version = "0.14.0-edge.6" dependencies = [ "anyhow", "apple-native-keyring-store", @@ -6784,7 +6784,7 @@ dependencies = [ [[package]] name = "iggy-connectors" -version = "0.5.0-edge.5" +version = "0.5.0-edge.6" dependencies = [ "async-trait", "axum", @@ -6856,7 +6856,7 @@ dependencies = [ [[package]] name = "iggy-mcp" -version = "0.5.0-edge.4" +version = "0.5.0-edge.5" dependencies = [ "axum", "axum-server", @@ -6890,7 +6890,7 @@ dependencies = [ [[package]] name = "iggy_binary_protocol" -version = "0.11.0-edge.5" +version = "0.11.0-edge.6" dependencies = [ "aligned-vec", "bytemuck", @@ -6903,7 +6903,7 @@ dependencies = [ [[package]] name = "iggy_common" -version = "0.11.0-edge.5" +version = "0.11.0-edge.6" dependencies = [ "aes-gcm", "async-broadcast", diff --git a/Cargo.toml b/Cargo.toml index ea2d476acb..7e69e39c7b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -206,10 +206,10 @@ hyper-util = { version = "0.1.20", features = ["server-auto", "service"] } iceberg = "0.9.1" iceberg-catalog-rest = "0.9.1" iceberg-storage-opendal = "0.9.1" -iggy = { path = "core/sdk", version = "0.11.0-edge.5" } +iggy = { path = "core/sdk", version = "0.11.0-edge.6" } iggy-cli = { path = "core/cli", version = "0.14.0-edge.5" } -iggy_binary_protocol = { path = "core/binary_protocol", version = "0.11.0-edge.5" } -iggy_common = { path = "core/common", version = "0.11.0-edge.5" } +iggy_binary_protocol = { path = "core/binary_protocol", version = "0.11.0-edge.6" } +iggy_common = { path = "core/common", version = "0.11.0-edge.6" } iggy_connector_sdk = { path = "core/connectors/sdk", version = "0.4.0-edge.3" } indexmap = "2.14.0" integration = { path = "core/integration" } diff --git a/core/ai/mcp/Cargo.toml b/core/ai/mcp/Cargo.toml index bab4769b3e..5386c5c4a0 100644 --- a/core/ai/mcp/Cargo.toml +++ b/core/ai/mcp/Cargo.toml @@ -17,7 +17,7 @@ [package] name = "iggy-mcp" -version = "0.5.0-edge.4" +version = "0.5.0-edge.5" description = "MCP Server for Iggy message streaming platform" edition = "2024" license = "Apache-2.0" diff --git a/core/bench/Cargo.toml b/core/bench/Cargo.toml index 60e521d910..e2a3c8d8d0 100644 --- a/core/bench/Cargo.toml +++ b/core/bench/Cargo.toml @@ -17,7 +17,7 @@ [package] name = "iggy-bench" -version = "0.6.0-edge.5" +version = "0.6.0-edge.6" edition = "2024" license = "Apache-2.0" repository = "https://github.com/apache/iggy" diff --git a/core/binary_protocol/Cargo.toml b/core/binary_protocol/Cargo.toml index 79e99c768f..0a83a4c010 100644 --- a/core/binary_protocol/Cargo.toml +++ b/core/binary_protocol/Cargo.toml @@ -17,7 +17,7 @@ [package] name = "iggy_binary_protocol" -version = "0.11.0-edge.5" +version = "0.11.0-edge.6" description = "Wire protocol types and codec for the Iggy binary protocol. Shared between server and SDK." edition = "2024" rust-version.workspace = true diff --git a/core/cli/Cargo.toml b/core/cli/Cargo.toml index bfecb9ef32..a2f8756bb7 100644 --- a/core/cli/Cargo.toml +++ b/core/cli/Cargo.toml @@ -17,7 +17,7 @@ [package] name = "iggy-cli" -version = "0.14.0-edge.5" +version = "0.14.0-edge.6" edition = "2024" rust-version.workspace = true authors = ["bartosz.ciesla@gmail.com"] diff --git a/core/common/Cargo.toml b/core/common/Cargo.toml index 8055104fe7..16aaa621c5 100644 --- a/core/common/Cargo.toml +++ b/core/common/Cargo.toml @@ -17,7 +17,7 @@ [package] name = "iggy_common" -version = "0.11.0-edge.5" +version = "0.11.0-edge.6" description = "Iggy is the persistent message streaming platform written in Rust, supporting QUIC, TCP and HTTP transport protocols, capable of processing millions of messages per second." edition = "2024" rust-version.workspace = true diff --git a/core/connectors/runtime/Cargo.toml b/core/connectors/runtime/Cargo.toml index 37b5351d5b..d791d7d26a 100644 --- a/core/connectors/runtime/Cargo.toml +++ b/core/connectors/runtime/Cargo.toml @@ -17,7 +17,7 @@ [package] name = "iggy-connectors" -version = "0.5.0-edge.5" +version = "0.5.0-edge.6" description = "Connectors runtime for Iggy message streaming platform" edition = "2024" license = "Apache-2.0" diff --git a/core/sdk/Cargo.toml b/core/sdk/Cargo.toml index 0e3e504b51..09ef84d03c 100644 --- a/core/sdk/Cargo.toml +++ b/core/sdk/Cargo.toml @@ -17,7 +17,7 @@ [package] name = "iggy" -version = "0.11.0-edge.5" +version = "0.11.0-edge.6" description = "Iggy is the persistent message streaming platform written in Rust, supporting QUIC, TCP and HTTP transport protocols, capable of processing millions of messages per second." edition = "2024" rust-version.workspace = true diff --git a/core/sdk/src/clients/producer_config.rs b/core/sdk/src/clients/producer_config.rs index 1e49e90f26..f2d5948527 100644 --- a/core/sdk/src/clients/producer_config.rs +++ b/core/sdk/src/clients/producer_config.rs @@ -105,7 +105,8 @@ pub struct BackgroundConfig { /// Action to apply when back-pressure limits are reached #[builder(default = BackpressureMode::Block)] pub failure_mode: BackpressureMode, - /// Upper bound for the **bytes held in memory** across *all* shards. + /// Upper bound for the **bytes buffered or in flight** across *all* shards. + /// Bytes remain charged until the corresponding write completes. /// `IggyByteSize::from(0)` ⇒ unlimited. #[builder(default = IggyByteSize::from(32 * MIB as u64))] pub max_buffer_size: IggyByteSize, diff --git a/core/sdk/src/clients/producer_dispatcher.rs b/core/sdk/src/clients/producer_dispatcher.rs index acb79c45bf..447a1b86cd 100644 --- a/core/sdk/src/clients/producer_dispatcher.rs +++ b/core/sdk/src/clients/producer_dispatcher.rs @@ -20,7 +20,7 @@ use crate::clients::producer_config::{BackgroundConfig, BackpressureMode}; use crate::clients::producer_error_callback::ErrorCtx; use crate::clients::producer_sharding::{Shard, ShardMessage, ShardMessageWithPermit}; use futures::FutureExt; -use iggy_common::{Identifier, IggyError, IggyMessage, Partitioning, Sizeable}; +use iggy_common::{Identifier, IggyByteSize, IggyError, IggyMessage, Partitioning, Sizeable}; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; use tokio::sync::{Semaphore, broadcast}; @@ -61,13 +61,16 @@ impl ProducerDispatcher { tracing::debug!("error-callback worker finished"); }); - let bytes_permit = { - let bytes = config.max_buffer_size.as_bytes_usize(); - if bytes == 0 { usize::MAX } else { bytes } - }; + let max_buffer_size = config.max_buffer_size.as_bytes_u64(); + assert!( + max_buffer_size == 0 || max_buffer_size <= Semaphore::MAX_PERMITS as u64, + "max_buffer_size cannot exceed {} bytes on this platform", + Semaphore::MAX_PERMITS + ); + let bytes_permit = Arc::new(Semaphore::new(max_buffer_size as usize)); let slots_permit = Arc::new(Semaphore::new(if config.max_in_flight == 0 { - usize::MAX + Semaphore::MAX_PERMITS } else { config.max_in_flight })); @@ -87,7 +90,7 @@ impl ProducerDispatcher { shards, config, closed: AtomicBool::new(false), - bytes_permit: Arc::new(Semaphore::new(bytes_permit)), + bytes_permit, stop_tx, join_handle: handle, } @@ -112,41 +115,45 @@ impl ProducerDispatcher { }; let batch_bytes = shard_message.get_size_bytes(); - if batch_bytes > self.config.max_buffer_size { + if self.config.max_buffer_size != 0 && batch_bytes > self.config.max_buffer_size { return Err(IggyError::BackgroundSendBufferOverflow); } - let permit_bytes = match self - .bytes_permit - .clone() - .try_acquire_many_owned(batch_bytes.as_bytes_u32()) - { - Ok(perm) => perm, - Err(_) => match self.config.failure_mode { - BackpressureMode::FailImmediately => { - return Err(IggyError::BackgroundSendBufferOverflow); - } - BackpressureMode::Block => self - .bytes_permit - .clone() - .acquire_many_owned(batch_bytes.as_bytes_u32()) - .await - .map_err(|_| IggyError::BackgroundSendError)?, - BackpressureMode::BlockWithTimeout(timeout_dur) => { - match tokio::time::timeout( - timeout_dur.get_duration(), - self.bytes_permit - .clone() - .acquire_many_owned(batch_bytes.as_bytes_u32()), - ) - .await - { - Ok(Ok(perm)) => perm, - Ok(Err(_)) => return Err(IggyError::BackgroundSendError), - Err(_) => return Err(IggyError::BackgroundSendTimeout), + let permit_count = Self::permit_count(batch_bytes)?; + let bytes_permit = if self.config.max_buffer_size == 0 { + None + } else { + let permit = match self + .bytes_permit + .clone() + .try_acquire_many_owned(permit_count) + { + Ok(permit) => permit, + Err(_) => match &self.config.failure_mode { + BackpressureMode::FailImmediately => { + return Err(IggyError::BackgroundSendBufferOverflow); } - } - }, + BackpressureMode::Block => self + .bytes_permit + .clone() + .acquire_many_owned(permit_count) + .await + .map_err(|_| IggyError::BackgroundSendError)?, + BackpressureMode::BlockWithTimeout(timeout_duration) => { + match tokio::time::timeout( + timeout_duration.get_duration(), + self.bytes_permit.clone().acquire_many_owned(permit_count), + ) + .await + { + Ok(Ok(permit)) => permit, + Ok(Err(_)) => return Err(IggyError::BackgroundSendError), + Err(_) => return Err(IggyError::BackgroundSendTimeout), + } + } + }, + }; + Some(permit) }; let shard_ix = self.config.sharding.pick_shard( @@ -159,10 +166,15 @@ impl ProducerDispatcher { let shard = &self.shards[shard_ix]; shard - .send(ShardMessageWithPermit::new(shard_message, permit_bytes)) + .send(ShardMessageWithPermit::new(shard_message, bytes_permit)) .await } + fn permit_count(batch_size: IggyByteSize) -> Result { + u32::try_from(batch_size.as_bytes_u64()) + .map_err(|_| IggyError::BackgroundSendBufferOverflow) + } + /// Flushes each shard's buffer and stops its worker. Dropping the /// dispatcher instead of calling this silently discards any buffered, /// not-yet-sent messages. @@ -174,7 +186,7 @@ impl ProducerDispatcher { let _ = self.stop_tx.send(()); for shard in self.shards.drain(..) { - if let Err(e) = shard._handle.await { + if let Err(e) = shard.handle.await { tracing::error!("shard panicked: {e:?}"); } } @@ -237,6 +249,60 @@ mod tests { assert!(result.is_ok()); } + #[tokio::test] + async fn test_dispatch_succeeds_with_unlimited_buffer_and_in_flight_requests() { + let mut mock = MockProducerCoreBackend::new(); + mock.expect_send_internal() + .times(1) + .returning(|_, _, _, _| Box::pin(async { Ok(no_confirmations()) })); + + let config = BackgroundConfig::builder() + .max_buffer_size(0.into()) + .max_in_flight(0) + .batch_length(1) + .build(); + let dispatcher = ProducerDispatcher::new(Arc::new(mock), config); + + assert_eq!(dispatcher.bytes_permit.available_permits(), 0); + dispatcher + .dispatch( + vec![dummy_message(5)], + dummy_identifier(), + dummy_identifier(), + None, + ) + .await + .unwrap(); + dispatcher.shutdown().await; + } + + #[cfg(target_pointer_width = "64")] + #[tokio::test] + async fn test_dispatcher_supports_buffer_budget_above_u32_max() { + let mock = MockProducerCoreBackend::new(); + let budget_size = u32::MAX as u64 + 1; + let config = BackgroundConfig::builder() + .max_buffer_size(budget_size.into()) + .build(); + let dispatcher = ProducerDispatcher::new(Arc::new(mock), config); + + assert_eq!( + dispatcher.bytes_permit.available_permits(), + budget_size as usize + ); + dispatcher.shutdown().await; + } + + #[test] + fn test_permit_count_rejects_batch_above_u32_max() { + let result = ProducerDispatcher::permit_count(IggyByteSize::from(u32::MAX as u64 + 1)); + + assert!(matches!( + result, + Err(IggyError::BackgroundSendBufferOverflow) + )); + } + #[tokio::test] async fn test_dispatch_fails_on_buffer_overflow_immediate() { let mock = MockProducerCoreBackend::new(); diff --git a/core/sdk/src/clients/producer_sharding.rs b/core/sdk/src/clients/producer_sharding.rs index c2db25b42e..619470b31d 100644 --- a/core/sdk/src/clients/producer_sharding.rs +++ b/core/sdk/src/clients/producer_sharding.rs @@ -15,18 +15,20 @@ // specific language governing permissions and limitations // under the License. -use crate::clients::producer::ProducerCoreBackend; -use crate::clients::producer_config::BackgroundConfig; -use crate::clients::producer_error_callback::ErrorCtx; -use iggy_common::{Identifier, IggyByteSize, IggyError, IggyMessage, Partitioning, Sizeable}; use std::hash::DefaultHasher; use std::hash::{Hash, Hasher}; use std::sync::Arc; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; + +use iggy_common::{Identifier, IggyByteSize, IggyError, IggyMessage, Partitioning, Sizeable}; use tokio::sync::{OwnedSemaphorePermit, Semaphore, broadcast}; use tokio::task::JoinHandle; use tracing::{debug, error}; +use crate::clients::producer::ProducerCoreBackend; +use crate::clients::producer_config::BackgroundConfig; +use crate::clients::producer_error_callback::ErrorCtx; + /// A strategy for distributing messages across shards. /// /// Implementors of this trait define how to choose a shard for a given batch of messages. @@ -101,8 +103,11 @@ impl Sizeable for ShardMessage { let mut total = IggyByteSize::new(0); total += self.stream.get_size_bytes(); total += self.topic.get_size_bytes(); - for msg in &self.messages { - total += msg.get_size_bytes(); + if let Some(partitioning) = &self.partitioning { + total += partitioning.get_size_bytes(); + } + for message in &self.messages { + total += message.get_size_bytes(); } total } @@ -110,29 +115,36 @@ impl Sizeable for ShardMessage { pub struct ShardMessageWithPermit { pub inner: ShardMessage, - bytes_permit: OwnedSemaphorePermit, + size_bytes: u64, + bytes_permit: Option, + merged_bytes_permits: Vec, } impl ShardMessageWithPermit { - pub fn new(msg: ShardMessage, permit_bytes: OwnedSemaphorePermit) -> Self { + pub fn new(msg: ShardMessage, bytes_permit: Option) -> Self { + let size_bytes = msg.get_size_bytes().as_bytes_u64(); Self { inner: msg, - bytes_permit: permit_bytes, + size_bytes, + bytes_permit, + merged_bytes_permits: Vec::new(), } } - /// Takes over `other`'s messages together with its byte permit, so the buffered bytes stay - /// charged against the `max_buffer_size` budget until the merged batch has been written. fn merge(&mut self, other: Self) { self.inner.messages.extend(other.inner.messages); - self.bytes_permit.merge(other.bytes_permit); + self.size_bytes += other.size_bytes; + // Tokio stores a merged permit count in a u32, so retain permits separately to avoid + // overflowing when the buffer budget exceeds u32::MAX. + self.merged_bytes_permits.extend(other.bytes_permit); + self.merged_bytes_permits.extend(other.merged_bytes_permits); } } pub struct Shard { tx: flume::Sender, closed: Arc, - pub(crate) _handle: JoinHandle<()>, + pub(crate) handle: JoinHandle<()>, } impl Shard { @@ -158,7 +170,7 @@ impl Shard { maybe_msg = rx.recv_async() => { match maybe_msg { Ok(msg) => { - buffer_bytes += msg.inner.get_size_bytes().as_bytes_usize(); + buffer_bytes += msg.size_bytes as usize; buffer.push(msg); debug!( buffer_len = buffer.len(), @@ -200,7 +212,7 @@ impl Shard { _ = stop_rx.recv() => { closed_clone.store(true, Ordering::Release); while let Ok(msg) = rx.try_recv() { - buffer_bytes += msg.inner.get_size_bytes().as_bytes_usize(); + buffer_bytes += msg.size_bytes as usize; buffer.push(msg); } if !buffer.is_empty() { @@ -212,30 +224,24 @@ impl Shard { } }); - Self { - tx, - closed, - _handle: handle, - } + Self { tx, closed, handle } } - /// Drains the buffer into batches. - /// - /// If adjacent ShardMessages have the same destination (stream, topic, partition combination), - /// merge to avoid multiple sends. - /// The merge happens on both, the inner [`IggyMessage`] and the [`OwnedSemaphorePermit`]. + /// Drains the buffer and combines adjacent messages with the same destination. fn merge_batches(buffer: &mut Vec) -> Vec { let mut merged_batches: Vec = Vec::with_capacity(buffer.len()); - // Since buffer is a mutable reference, the drain leaves the buffer intact but empty, - // such that it can be filled up again. - for msg in buffer.drain(..) { + for message in buffer.drain(..) { if let Some(last) = merged_batches.last_mut() - && Self::same_destination(&last.inner, &msg.inner) + && Self::same_destination(&last.inner, &message.inner) + && last + .size_bytes + .checked_add(message.size_bytes) + .is_some_and(|size_bytes| size_bytes <= u32::MAX as u64) { - last.merge(msg); + last.merge(message); continue; } - merged_batches.push(msg); + merged_batches.push(message); } merged_batches } @@ -346,7 +352,7 @@ mod tests { .acquire_many_owned(message.get_size_bytes().as_bytes_u32()) .await .unwrap(); - ShardMessageWithPermit::new(message, permit) + ShardMessageWithPermit::new(message, Some(permit)) } #[tokio::test] @@ -362,26 +368,88 @@ mod tests { let charged = 10_000 - budget.available_permits(); let merged = Shard::merge_batches(&mut buffer); - // The original buffer should be drained. assert!(buffer.is_empty()); assert_eq!(merged.len(), 1); assert_eq!(merged[0].inner.messages.len(), 3); assert_eq!(budget.available_permits(), 10_000 - charged); - // Dropping merged gives back to the semaphore. drop(merged); assert_eq!(budget.available_permits(), 10_000); } + #[cfg(target_pointer_width = "64")] + #[tokio::test] + async fn test_merge_batches_keeps_more_than_u32_max_permits_charged() { + let budget_size = u32::MAX as usize + 1; + let budget = Arc::new(Semaphore::new(budget_size)); + let stream = dummy_identifier(); + let topic = dummy_identifier(); + let first_permit = budget.clone().acquire_many_owned(u32::MAX).await.unwrap(); + let second_permit = budget.clone().acquire_owned().await.unwrap(); + let mut buffer = vec![ + ShardMessageWithPermit::new( + ShardMessage { + stream: stream.clone(), + topic: topic.clone(), + messages: vec![dummy_message(1)], + partitioning: None, + }, + Some(first_permit), + ), + ShardMessageWithPermit::new( + ShardMessage { + stream, + topic, + messages: vec![dummy_message(1)], + partitioning: None, + }, + Some(second_permit), + ), + ]; + + let merged = Shard::merge_batches(&mut buffer); + + assert_eq!(merged.len(), 1); + assert_eq!(budget.available_permits(), 0); + drop(merged); + assert_eq!(budget.available_permits(), budget_size); + } + + #[test] + fn test_merge_batches_splits_batches_above_u32_max_bytes() { + let stream = dummy_identifier(); + let topic = dummy_identifier(); + let mut first = ShardMessageWithPermit::new( + ShardMessage { + stream: stream.clone(), + topic: topic.clone(), + messages: vec![dummy_message(1)], + partitioning: None, + }, + None, + ); + first.size_bytes = u32::MAX as u64; + let mut second = ShardMessageWithPermit::new( + ShardMessage { + stream, + topic, + messages: vec![dummy_message(1)], + partitioning: None, + }, + None, + ); + second.size_bytes = 1; + let mut buffer = vec![first, second]; + + let merged = Shard::merge_batches(&mut buffer); + + assert_eq!(merged.len(), 2); + } + #[tokio::test] async fn test_shard_keeps_budget_charged_until_merged_batch_is_written() { const BUDGET: usize = 10_000; - // The two channels simulate the timing of a real write, which `flush_buffer` awaits at - // `core.send_internal`. Receiving on `write_started_rx` means the worker has entered that - // call, so the batch is on the wire and its bytes are still buffered from the producer's - // point of view. Sending on `release_write_tx` simulates that the call - // returns, the loop iteration ends and the batch drops and frees its permit. let (write_started_tx, write_started_rx) = flume::unbounded::<()>(); let (release_write_tx, release_write_rx) = flume::unbounded::<()>(); @@ -407,7 +475,7 @@ mod tests { let budget = Arc::new(Semaphore::new(BUDGET)); let slots_permit = Arc::new(Semaphore::new(100)); - let (_stop_tx, stop_rx) = broadcast::channel(1); + let (stop_tx, stop_rx) = broadcast::channel(1); let shard = Shard::new( Arc::new(mock), config, @@ -425,8 +493,10 @@ mod tests { let charged = BUDGET - budget.available_permits(); assert!(charged > 0); - // The three batches share a destination, so the flush merges them into a single write. - write_started_rx.recv_async().await.unwrap(); + tokio::time::timeout(Duration::from_secs(1), write_started_rx.recv_async()) + .await + .expect("the merged write must start") + .unwrap(); assert_eq!( budget.available_permits(), BUDGET - charged, @@ -441,6 +511,9 @@ mod tests { }) .await .expect("the written batch must give its permits back"); + + stop_tx.send(()).unwrap(); + shard.handle.await.unwrap(); } #[tokio::test] @@ -477,7 +550,7 @@ mod tests { }; let wrapped = ShardMessageWithPermit::new( message, - permit_bytes.clone().acquire_many_owned(1).await.unwrap(), + Some(permit_bytes.clone().acquire_many_owned(1).await.unwrap()), ); shard.send(wrapped).await.unwrap(); } @@ -518,11 +591,13 @@ mod tests { }; let wrapped = ShardMessageWithPermit::new( message, - permit_bytes - .clone() - .acquire_many_owned(10_000) - .await - .unwrap(), + Some( + permit_bytes + .clone() + .acquire_many_owned(10_000) + .await + .unwrap(), + ), ); shard.send(wrapped).await.unwrap(); @@ -562,7 +637,7 @@ mod tests { }; let wrapped = ShardMessageWithPermit::new( message, - permit_bytes.clone().acquire_many_owned(1).await.unwrap(), + Some(permit_bytes.clone().acquire_many_owned(1).await.unwrap()), ); shard.send(wrapped).await.unwrap(); @@ -603,7 +678,7 @@ mod tests { }; let wrapped = ShardMessageWithPermit::new( message, - permit_bytes.clone().acquire_many_owned(1).await.unwrap(), + Some(permit_bytes.clone().acquire_many_owned(1).await.unwrap()), ); shard.send(wrapped).await.unwrap(); @@ -620,7 +695,7 @@ mod tests { let shard = Shard { tx, closed: Arc::new(AtomicBool::new(false)), - _handle: tokio::spawn(async {}), + handle: tokio::spawn(async {}), }; let permit_bytes = Arc::new(Semaphore::new(10_000)); @@ -633,7 +708,7 @@ mod tests { }; let wrapped = ShardMessageWithPermit::new( message, - permit_bytes.clone().acquire_many_owned(1).await.unwrap(), + Some(permit_bytes.clone().acquire_many_owned(1).await.unwrap()), ); let result = shard.send(wrapped).await; diff --git a/foreign/python/Cargo.toml b/foreign/python/Cargo.toml index eefeaa1c6f..7d5f2210a0 100644 --- a/foreign/python/Cargo.toml +++ b/foreign/python/Cargo.toml @@ -37,7 +37,7 @@ doc = false [dependencies] bytes = "1.12.1" futures = "0.3.33" -iggy = { path = "../../core/sdk", version = "0.11.0-edge.5" } +iggy = { path = "../../core/sdk", version = "0.11.0-edge.6" } paste = "1" pyo3 = "0.29.0" pyo3-async-runtimes = { version = "0.29.0", features = [ From 2fef029efd046255c6994e0cb7cbc3ba08bdd1ca Mon Sep 17 00:00:00 2001 From: spetz Date: Mon, 31 Aug 2026 13:26:26 +0200 Subject: [PATCH 3/3] fix versions --- Cargo.toml | 2 +- bdd/python/uv.lock | 2 +- examples/python/uv.lock | 2 +- foreign/python/Cargo.toml | 2 +- foreign/python/pyproject.toml | 2 +- foreign/python/uv.lock | 2 +- 6 files changed, 6 insertions(+), 6 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 7e69e39c7b..7872f40937 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -207,7 +207,7 @@ iceberg = "0.9.1" iceberg-catalog-rest = "0.9.1" iceberg-storage-opendal = "0.9.1" iggy = { path = "core/sdk", version = "0.11.0-edge.6" } -iggy-cli = { path = "core/cli", version = "0.14.0-edge.5" } +iggy-cli = { path = "core/cli", version = "0.14.0-edge.6" } iggy_binary_protocol = { path = "core/binary_protocol", version = "0.11.0-edge.6" } iggy_common = { path = "core/common", version = "0.11.0-edge.6" } iggy_connector_sdk = { path = "core/connectors/sdk", version = "0.4.0-edge.3" } diff --git a/bdd/python/uv.lock b/bdd/python/uv.lock index b401a9ae2c..9198d276fa 100644 --- a/bdd/python/uv.lock +++ b/bdd/python/uv.lock @@ -8,7 +8,7 @@ exclude-newer-span = "P7D" [[package]] name = "apache-iggy" -version = "0.9.0.dev5" +version = "0.9.0.dev6" source = { directory = "../../foreign/python" } [package.metadata] diff --git a/examples/python/uv.lock b/examples/python/uv.lock index 1f9c5b3132..eee0abbdd9 100644 --- a/examples/python/uv.lock +++ b/examples/python/uv.lock @@ -8,7 +8,7 @@ exclude-newer-span = "P7D" [[package]] name = "apache-iggy" -version = "0.9.0.dev5" +version = "0.9.0.dev6" source = { directory = "../../foreign/python" } [package.metadata] diff --git a/foreign/python/Cargo.toml b/foreign/python/Cargo.toml index 7d5f2210a0..ee723b4888 100644 --- a/foreign/python/Cargo.toml +++ b/foreign/python/Cargo.toml @@ -17,7 +17,7 @@ [package] name = "apache-iggy" -version = "0.9.0-dev5" +version = "0.9.0-dev6" edition = "2024" authors = ["Iggy Committers "] license = "Apache-2.0" diff --git a/foreign/python/pyproject.toml b/foreign/python/pyproject.toml index 7a82108b6c..d9623d4312 100644 --- a/foreign/python/pyproject.toml +++ b/foreign/python/pyproject.toml @@ -22,7 +22,7 @@ build-backend = "maturin" [project] name = "apache-iggy" requires-python = ">=3.10" -version = "0.9.0.dev5" +version = "0.9.0.dev6" description = "Apache Iggy is the persistent message streaming platform written in Rust, supporting QUIC, TCP and HTTP transport protocols, capable of processing millions of messages per second." readme = "README.md" license = { file = "LICENSE" } diff --git a/foreign/python/uv.lock b/foreign/python/uv.lock index 963b58acf9..214d7023b5 100644 --- a/foreign/python/uv.lock +++ b/foreign/python/uv.lock @@ -8,7 +8,7 @@ exclude-newer-span = "P7D" [[package]] name = "apache-iggy" -version = "0.9.0.dev5" +version = "0.9.0.dev6" source = { editable = "." } [package.optional-dependencies]