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]