This is an automated email from the ASF dual-hosted git repository.

hubcio pushed a commit to branch benchmark/deferred-polling-task-per-partition
in repository https://gitbox.apache.org/repos/asf/iggy.git

commit 391fb66436289e1cff005ecf50265c3f77ecc7c4
Author: Hubert Gruszecki <[email protected]>
AuthorDate: Thu Sep 17 18:51:51 2026 +0200

    fix(shard): charge a deferred poll for its reply, not its borrows
    
    A deferred poll was refused with InvalidSizeBytes whenever its result
    borrowed more memory than the read had reserved. The charge counted a
    shared allocation once per fragment slicing it, and it ran before the
    reply was trimmed to the caller's byte limit, so a serveable poll failed
    on how the journal happened to allocate rather than on what it returned.
    Under load every poll on a busy partition failed, the client treated the
    refusal as a lost assignment, and the consumer fell behind for good.
    
    Count each backing allocation once, and take the charge after the trim.
    When the reply still pins more than the reservation, copy it into one
    exact-size allocation instead of refusing: fragment boundaries survive,
    so batch framing is unchanged, and the charge drops to the reply's own
    size. A poll that can return fewer messages must never fail.
    
    The check itself stays. It bounds real pinned memory, which the shard
    aggregates through max_inflight_bytes.
---
 core/partitions/src/poll_plan.rs | 101 ++++++++++++++++++++++++++++++++++++---
 core/partitions/src/types.rs     |   7 +++
 core/shard/src/poll/deferred.rs  |  14 +++---
 3 files changed, 109 insertions(+), 13 deletions(-)

diff --git a/core/partitions/src/poll_plan.rs b/core/partitions/src/poll_plan.rs
index a65c5b9b2..525e41df6 100644
--- a/core/partitions/src/poll_plan.rs
+++ b/core/partitions/src/poll_plan.rs
@@ -271,15 +271,50 @@ impl PollReadResult {
         Ok((self, limited))
     }
 
-    /// Conservative charge: shared allocations count once for each fragment.
+    /// Memory this result keeps alive, counting every backing allocation once
+    /// however many fragments slice it.
     #[must_use]
     pub fn retained_bytes(&self) -> usize {
-        self.fragments.iter().fold(
-            self.fragments
-                .capacity()
-                .saturating_mul(size_of::<crate::Fragment>()),
-            |bytes, fragment| 
bytes.saturating_add(fragment.allocation_bytes()),
-        )
+        let mut bytes = self
+            .fragments
+            .capacity()
+            .saturating_mul(size_of::<crate::Fragment>());
+        for (index, fragment) in self.fragments.iter().enumerate() {
+            if !self.fragments[..index]
+                .iter()
+                .any(|earlier| earlier.shares_allocation_with(fragment))
+            {
+                bytes = bytes.saturating_add(fragment.allocation_bytes());
+            }
+        }
+        bytes
+    }
+
+    /// Copy the selection into one exact-size allocation and drop the borrowed
+    /// buffers. A reply that slices a few bytes out of segment-sized journal
+    /// buffers pins all of them, so the copy trades one pass over the reply 
for
+    /// a charge that tracks the reply itself. Fragment boundaries survive, 
which
+    /// keeps batch framing intact.
+    #[must_use]
+    pub fn compacted(mut self) -> Self {
+        let total: usize = self.fragments.iter().map(Fragment::len).sum();
+        let mut buffer = Owned::with_capacity(total);
+        for fragment in &self.fragments {
+            buffer.extend_from_slice(fragment.as_slice());
+        }
+        let compacted = Frozen::from(buffer);
+        let mut start = 0;
+        self.fragments = self
+            .fragments
+            .iter()
+            .map(|fragment| {
+                let end = start + fragment.len();
+                let piece = Fragment::slice(compacted.clone(), start, end);
+                start = end;
+                piece
+            })
+            .collect();
+        self
     }
 
     #[must_use]
@@ -1713,4 +1748,56 @@ mod tests {
         assert_eq!(walk.matched, 0);
         assert!(walk.fragments.is_empty());
     }
+
+    fn resident(fragments: PollFragments) -> PollReadResult {
+        PollPlan {
+            context: PollContext {
+                history: PollHistoryId::default(),
+                consumer: PollingConsumer::Consumer(0, 0),
+                auto_commit: true,
+            },
+            commit_offset: 0,
+            tier: PollTier::Resident {
+                fragments,
+                last_matching_offset: None,
+                message_count: 0,
+            },
+        }
+        .execute_resident()
+    }
+
+    fn windows(source: &Frozen<4096>, count: usize) -> PollFragments {
+        (0..count)
+            .map(|window| Fragment::slice(source.clone(), window * 4096, 
window * 4096 + 512))
+            .collect()
+    }
+
+    #[test]
+    fn 
retained_bytes_counts_one_allocation_once_however_many_slices_borrow_it() {
+        let source: Frozen<4096> = Owned::copy_from_slice(&vec![7; 1 << 
20]).into();
+        let one = resident(windows(&source, 1)).retained_bytes();
+        let four = resident(windows(&source, 4)).retained_bytes();
+        assert_eq!(one, four);
+        assert!(four < 2 * source.allocation_bytes());
+    }
+
+    #[test]
+    fn 
compacting_releases_the_borrowed_allocation_and_preserves_the_selection() {
+        let source: Frozen<4096> = Owned::copy_from_slice(&vec![7; 1 << 
20]).into();
+        let borrowed = resident(windows(&source, 4));
+        let selection: Vec<Vec<u8>> = borrowed
+            .fragments
+            .iter()
+            .map(|fragment| fragment.as_slice().to_vec())
+            .collect();
+        let charge = borrowed.retained_bytes();
+        let compacted = borrowed.compacted();
+        assert!(compacted.retained_bytes() < charge);
+        let kept: Vec<Vec<u8>> = compacted
+            .fragments
+            .iter()
+            .map(|fragment| fragment.as_slice().to_vec())
+            .collect();
+        assert_eq!(kept, selection);
+    }
 }
diff --git a/core/partitions/src/types.rs b/core/partitions/src/types.rs
index 7619552c7..50a08570e 100644
--- a/core/partitions/src/types.rs
+++ b/core/partitions/src/types.rs
@@ -102,6 +102,13 @@ impl<const ALIGN: usize> Fragment<ALIGN> {
     pub fn borrows_from(&self, source: &Frozen<ALIGN>) -> bool {
         self.source.shares_allocation(source)
     }
+
+    /// Whether both fragments keep the same allocation alive, so a byte charge
+    /// counts it once instead of once per slice.
+    #[must_use]
+    pub fn shares_allocation_with(&self, other: &Self) -> bool {
+        self.source.shares_allocation(&other.source)
+    }
 }
 
 /// Arguments for polling messages from a partition.
diff --git a/core/shard/src/poll/deferred.rs b/core/shard/src/poll/deferred.rs
index d9c51e75d..4f7e96d98 100644
--- a/core/shard/src/poll/deferred.rs
+++ b/core/shard/src/poll/deferred.rs
@@ -626,12 +626,7 @@ where
                 return;
             }
         };
-        let config = *self.deferred_polls.config.borrow();
-        if result.retained_bytes() > config.max_read_bytes {
-            self.reject_deferred_poll(id, &waiter, 
IggyError::InvalidSizeBytes);
-            return;
-        }
-        let (result, byte_limited) =
+        let (mut result, byte_limited) =
             match result.limit_bytes(waiter.request.options.max_bytes as 
usize) {
                 Ok(result) => result,
                 Err(error) => {
@@ -639,6 +634,13 @@ where
                     return;
                 }
             };
+        // Charge what the reply carries, not what the read walked over. A
+        // selection that borrows journal buffers can still pin more than the
+        // reservation, so copy it out instead of refusing a serveable poll.
+        let config = *self.deferred_polls.config.borrow();
+        if result.retained_bytes() > config.max_read_bytes {
+            result = result.compacted();
+        }
         if result.message_count() >= waiter.request.options.min_count
             || byte_limited
             || waiter.terminal

Reply via email to