fix(sdk): respect max_buffer_size when merging batches (#3933)
Merging adjacent producer batches released their byte permits before
the write completed, allowing buffered and in-flight data to exceed
max_buffer_size. Large permit totals could also overflow during the
merge and permanently reduce available capacity.
Retain each permit separately until send_internal completes. Validate
batch sizes before u32 conversion, split merged writes at the wire
limit, and handle zero buffer and in-flight limits without invalid
semaphores.
Update max_buffer_size documentation and add regression coverage for
permit lifetime, unlimited limits, and budgets above u32::MAX. Keep
the edge.6 package versions synchronized.
Fixes #3932
diff --git a/Cargo.lock b/Cargo.lock
index c6aa065..251f7f3 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -6659,7 +6659,7 @@
[[package]]
name = "iggy"
-version = "0.11.0-edge.5"
+version = "0.11.0-edge.6"
dependencies = [
"async-broadcast",
"async-dropper",
@@ -6693,7 +6693,7 @@
[[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 @@
[[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 @@
[[package]]
name = "iggy-connectors"
-version = "0.5.0-edge.5"
+version = "0.5.0-edge.6"
dependencies = [
"async-trait",
"axum",
@@ -6856,7 +6856,7 @@
[[package]]
name = "iggy-mcp"
-version = "0.5.0-edge.4"
+version = "0.5.0-edge.5"
dependencies = [
"axum",
"axum-server",
@@ -6890,7 +6890,7 @@
[[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 @@
[[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 ea2d476..7872f40 100644
--- a/Cargo.toml
+++ b/Cargo.toml
@@ -206,10 +206,10 @@
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-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 = { path = "core/sdk", version = "0.11.0-edge.6" }
+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" }
indexmap = "2.14.0"
integration = { path = "core/integration" }
diff --git a/bdd/python/uv.lock b/bdd/python/uv.lock
index b401a9a..9198d27 100644
--- a/bdd/python/uv.lock
+++ b/bdd/python/uv.lock
@@ -8,7 +8,7 @@
[[package]]
name = "apache-iggy"
-version = "0.9.0.dev5"
+version = "0.9.0.dev6"
source = { directory = "../../foreign/python" }
[package.metadata]
diff --git a/core/ai/mcp/Cargo.toml b/core/ai/mcp/Cargo.toml
index bab4769..5386c5c 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 60e521d..e2a3c8d 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 79e99c7..0a83a4c 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 bfecb9e..a2f8756 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 8055104..16aaa62 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 37b5351..d791d7d 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 0e3e504..09ef84d 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 1e49e90..f2d5948 100644
--- a/core/sdk/src/clients/producer_config.rs
+++ b/core/sdk/src/clients/producer_config.rs
@@ -105,7 +105,8 @@
/// 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 acb79c4..447a1b8 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_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 @@
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 @@
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 @@
};
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 @@
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, IggyError> {
+ 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 @@
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:?}");
}
}
@@ -238,6 +250,60 @@
}
#[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 668256f..619470b 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 @@
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,22 +115,36 @@
pub struct ShardMessageWithPermit {
pub inner: ShardMessage,
- _bytes_permit: Option<OwnedSemaphorePermit>,
+ size_bytes: u64,
+ bytes_permit: Option<OwnedSemaphorePermit>,
+ merged_bytes_permits: Vec<OwnedSemaphorePermit>,
}
impl ShardMessageWithPermit {
- pub fn new(msg: ShardMessage, permit_bytes: OwnedSemaphorePermit) -> Self {
+ pub fn new(msg: ShardMessage, bytes_permit: Option<OwnedSemaphorePermit>) -> Self {
+ let size_bytes = msg.get_size_bytes().as_bytes_u64();
Self {
inner: msg,
- _bytes_permit: Some(permit_bytes),
+ size_bytes,
+ bytes_permit,
+ merged_bytes_permits: Vec::new(),
}
}
+
+ fn merge(&mut self, other: Self) {
+ self.inner.messages.extend(other.inner.messages);
+ 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<ShardMessageWithPermit>,
closed: Arc<AtomicBool>,
- pub(crate) _handle: JoinHandle<()>,
+ pub(crate) handle: JoinHandle<()>,
}
impl Shard {
@@ -151,7 +170,7 @@
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(),
@@ -193,7 +212,7 @@
_ = 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() {
@@ -205,11 +224,26 @@
}
});
- Self {
- tx,
- closed,
- _handle: handle,
+ Self { tx, closed, handle }
+ }
+
+ /// Drains the buffer and combines adjacent messages with the same destination.
+ fn merge_batches(buffer: &mut Vec<ShardMessageWithPermit>) -> Vec<ShardMessageWithPermit> {
+ let mut merged_batches: Vec<ShardMessageWithPermit> = Vec::with_capacity(buffer.len());
+ for message in buffer.drain(..) {
+ if let Some(last) = merged_batches.last_mut()
+ && 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(message);
+ continue;
+ }
+ merged_batches.push(message);
}
+ merged_batches
}
async fn flush_buffer(
@@ -223,18 +257,7 @@
return;
}
- let mut merged_batches: Vec<ShardMessageWithPermit> = 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 +335,187 @@
.unwrap()
}
+ async fn charged_batch(
+ budget: &Arc<Semaphore>,
+ stream: Arc<Identifier>,
+ topic: Arc<Identifier>,
+ 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, Some(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);
+
+ 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);
+
+ 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;
+
+ 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);
+
+ 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,
+ "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");
+
+ stop_tx.send(()).unwrap();
+ shard.handle.await.unwrap();
+ }
+
#[tokio::test]
async fn test_shard_flushes_by_batch_length() {
let mut mock = MockProducerCoreBackend::new();
@@ -346,7 +550,7 @@
};
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();
}
@@ -387,11 +591,13 @@
};
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();
@@ -431,7 +637,7 @@
};
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();
@@ -472,7 +678,7 @@
};
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();
@@ -489,7 +695,7 @@
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));
@@ -502,7 +708,7 @@
};
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/examples/python/uv.lock b/examples/python/uv.lock
index 1f9c5b3..eee0abb 100644
--- a/examples/python/uv.lock
+++ b/examples/python/uv.lock
@@ -8,7 +8,7 @@
[[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 eefeaa1..ee723b4 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 <dev@iggy.apache.org>"]
license = "Apache-2.0"
@@ -37,7 +37,7 @@
[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 = [
diff --git a/foreign/python/pyproject.toml b/foreign/python/pyproject.toml
index 7a82108..d9623d4 100644
--- a/foreign/python/pyproject.toml
+++ b/foreign/python/pyproject.toml
@@ -22,7 +22,7 @@
[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 963b58a..214d702 100644
--- a/foreign/python/uv.lock
+++ b/foreign/python/uv.lock
@@ -8,7 +8,7 @@
[[package]]
name = "apache-iggy"
-version = "0.9.0.dev5"
+version = "0.9.0.dev6"
source = { editable = "." }
[package.optional-dependencies]