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 c4e8f56be91c9208d51655158937ffdf3386a235
Author: Hubert Gruszecki <[email protected]>
AuthorDate: Thu Sep 17 17:19:07 2026 +0200

    fix(partitions): admit deferred polls regardless of journal residency
    
    A deferred poll was refused with InvalidSizeBytes on any partition holding
    more than about 1.3 MiB of resident journal, while an immediate poll on the
    same partition at the same instant succeeded. The admission check charged 
the
    whole resident journal against the read budget instead of the offset range 
the
    poll would read, so the refusal did not depend on the request at all. Under
    load a partition crossed that line within half a second of production and 
then
    refused every deferred poll from that point on.
    
    The byte budget is already enforced where the bytes exist: the served 
result is
    rejected above the configured max_read_bytes, and the reply is trimmed to 
the
    caller's max_bytes with the trailing batch re-framed. Without the check the
    bounded plan and snapshot builders were identical to the plain ones, and the
    journal scan behind it, which decoded every resident batch on the shard pump
    for every probe, has no caller left.
    
    (cherry picked from commit c5ef0d6fde6b7dbd7c1a6efc685182bc8334b310)
---
 core/partitions/src/iggy_partition.rs  | 29 -----------------------------
 core/partitions/src/iggy_partitions.rs | 20 --------------------
 core/partitions/src/journal.rs         | 19 -------------------
 core/shard/src/poll/deferred.rs        | 16 ++++------------
 core/shard/src/poll/deferred_tests.rs  | 10 ++++++++++
 core/shard/src/poll/test_support.rs    |  8 +++++++-
 6 files changed, 21 insertions(+), 81 deletions(-)

diff --git a/core/partitions/src/iggy_partition.rs 
b/core/partitions/src/iggy_partition.rs
index 71436e49f..90d05f3d6 100644
--- a/core/partitions/src/iggy_partition.rs
+++ b/core/partitions/src/iggy_partition.rs
@@ -4598,35 +4598,6 @@ where
         }
     }
 
-    pub(crate) fn build_bounded_poll_plan(
-        &mut self,
-        consumer: PollingConsumer,
-        args: &PollingArgs,
-        validate_checksum: bool,
-        max_bytes: usize,
-    ) -> Result<PollPlan, IggyError> {
-        // Source snapshots, sparse copies and fragment descriptors can 
coexist.
-        const RESIDENT_ALLOCATION_FACTOR: usize = 4;
-        // Covers file paths, read handles and fixed IO/selection bookkeeping.
-        const READ_METADATA_BYTES: usize = 64 * 1024;
-        let resident_bytes = 
self.log.journal().inner.resident_message_bytes()?;
-        let segment_bytes = self
-            .log
-            .segments()
-            .len()
-            .saturating_mul(size_of::<DiskSegment>());
-        // Selection may copy sparse fragments while the source snapshot 
remains live.
-        if resident_bytes
-            .saturating_mul(RESIDENT_ALLOCATION_FACTOR)
-            .saturating_add(segment_bytes)
-            .saturating_add(READ_METADATA_BYTES)
-            > max_bytes
-        {
-            return Err(IggyError::InvalidSizeBytes);
-        }
-        Ok(self.build_poll_plan(consumer, args, validate_checksum))
-    }
-
     /// Snapshot read resources synchronously on the partition owner.
     /// Only the owned snapshot crosses a suspension during disk I/O.
     #[allow(clippy::too_many_lines)]
diff --git a/core/partitions/src/iggy_partitions.rs 
b/core/partitions/src/iggy_partitions.rs
index 15fa26e34..ef0fb4ca2 100644
--- a/core/partitions/src/iggy_partitions.rs
+++ b/core/partitions/src/iggy_partitions.rs
@@ -587,26 +587,6 @@ where
         Some(partition.build_poll_plan(consumer, args, validate_checksum))
     }
 
-    /// Snapshot a deferred read after checking its resident allocation budget.
-    ///
-    /// # Errors
-    /// Returns a size error when the snapshot exceeds the budget, or a read
-    /// error when a resident journal entry cannot be decoded.
-    pub fn build_bounded_poll_snapshot(
-        &self,
-        namespace: &IggyNamespace,
-        consumer: PollingConsumer,
-        args: &PollingArgs,
-        max_bytes: usize,
-    ) -> Result<Option<PollPlan>, IggyError> {
-        let validate_checksum = self.config.validate_checksum;
-        self.get_mut_by_ns(namespace)
-            .map(|partition| {
-                partition.build_bounded_poll_plan(consumer, args, 
validate_checksum, max_bytes)
-            })
-            .transpose()
-    }
-
     /// Validate and accept a poll synchronously on the owning pump.
     /// Attempt the reply synchronously, then immediately await any returned
     /// continuation through [`Self::replicate_poll_completion`] on the same 
pump.
diff --git a/core/partitions/src/journal.rs b/core/partitions/src/journal.rs
index d2595bc3d..6cf64d215 100644
--- a/core/partitions/src/journal.rs
+++ b/core/partitions/src/journal.rs
@@ -16,7 +16,6 @@
 // under the License.
 
 use iggy_binary_protocol::{Operation, PrepareHeader};
-use iggy_common::IggyError;
 use journal::{Journal, Storage};
 use server_common::{
     iobuf::{Frozen, Owned},
@@ -758,24 +757,6 @@ impl PartitionJournal<PartitionJournalMemStorage> {
         entries
     }
 
-    /// Conservative snapshot charge without allocating an entry vector.
-    pub fn resident_message_bytes(&self) -> Result<usize, IggyError> {
-        let offset_to_op = unsafe { &*self.offset_to_op.get() };
-        let op_to_storage_offset = unsafe { &*self.op_to_storage_offset.get() 
};
-        let inner = unsafe { &*self.inner.get() };
-        offset_to_op.values().try_fold(0usize, |bytes, op| {
-            let entry = op_to_storage_offset
-                .get(op)
-                .and_then(|offset| inner.storage.read_at_sync(*offset))
-                .ok_or(IggyError::CannotReadMessage)?;
-            decode_prepare_slice_trusted(entry.as_slice())
-                .map_err(|_| IggyError::CannotReadMessage)?;
-            Ok(bytes
-                .saturating_add(entry.allocation_bytes())
-                .saturating_add(size_of::<JournalBuffer>()))
-        })
-    }
-
     /// Owned, append-ordered clones of every resident entry above the purge
     /// floor, control ops included; one `Frozen` refcount bump each. Linear in
     /// the resident journal, so not for the poll path, which takes
diff --git a/core/shard/src/poll/deferred.rs b/core/shard/src/poll/deferred.rs
index 30104ef88..d9c51e75d 100644
--- a/core/shard/src/poll/deferred.rs
+++ b/core/shard/src/poll/deferred.rs
@@ -554,21 +554,13 @@ where
         };
         // Snapshot, disk selection and index loading can coexist.
         let max_bytes = reservation.bytes / READ_BUDGET_PARTS;
-        let plan = match self.plane.partitions().build_bounded_poll_snapshot(
+        let Some(plan) = self.plane.partitions().build_poll_snapshot(
             &waiter.namespace,
             waiter.request.consumer,
             &waiter.request.args,
-            max_bytes,
-        ) {
-            Ok(Some(plan)) => plan,
-            Ok(None) => {
-                self.reject_deferred_poll(id, &waiter, 
IggyError::TransientNotAccepted);
-                return;
-            }
-            Err(error) => {
-                self.reject_deferred_poll(id, &waiter, error);
-                return;
-            }
+        ) else {
+            self.reject_deferred_poll(id, &waiter, 
IggyError::TransientNotAccepted);
+            return;
         };
         waiter.last_visibility = Some(visibility);
         if plan.needs_off_pump_io() {
diff --git a/core/shard/src/poll/deferred_tests.rs 
b/core/shard/src/poll/deferred_tests.rs
index 0d99cba9e..cd5f82481 100644
--- a/core/shard/src/poll/deferred_tests.rs
+++ b/core/shard/src/poll/deferred_tests.rs
@@ -416,3 +416,13 @@ async fn 
read_crossing_readiness_deadline_finishes_within_request_budget() {
         );
     }
 }
+
+#[compio::test]
+async fn poll_admission_is_independent_of_resident_journal_size() {
+    // 1.5 MiB resident, one message polled: the read budget bounds the reply,
+    // not how much the partition holds.
+    let payload = "x".repeat(768 * 1024);
+    let owner = owner(&[payload.as_str(); 2]).await;
+    let replies = submit(&owner, request(&owner, 1, false, false)).await;
+    assert_eq!(offsets(&replies), [0]);
+}
diff --git a/core/shard/src/poll/test_support.rs 
b/core/shard/src/poll/test_support.rs
index 6664288ed..d8d13dfc4 100644
--- a/core/shard/src/poll/test_support.rs
+++ b/core/shard/src/poll/test_support.rs
@@ -33,6 +33,8 @@ use server_common::sharding::IggyNamespace;
 
 pub(super) type PollTestMetadata = MuxStateMachine<variadic!(Users, Streams)>;
 
+const SEGMENT_FLOOR: u64 = 1024 * 1024;
+
 /// Commit one batch starting at offset zero and keep it resident. Each call
 /// creates an independent partition history, even for the same namespace.
 #[allow(clippy::future_not_send)]
@@ -44,7 +46,11 @@ pub(super) async fn partition_with_messages<B: MessageBus + 
Clone>(
     let cluster_id = 1;
     let replica_id = 0;
     let replica_count = 3;
-    let segment_size = IggyByteSize::from(1_048_576_u64);
+    // The batch must fit one segment or it never materialises, and it must 
stay
+    // under the flush threshold or it stops being resident. Doubling the 
payload
+    // total covers framing; the floor keeps small fixtures where they were.
+    let payload_bytes: u64 = payloads.iter().map(|payload| payload.len() as 
u64).sum();
+    let segment_size = 
IggyByteSize::from(SEGMENT_FLOOR.max(payload_bytes.saturating_mul(2)));
     let config = PartitionsConfig {
         messages_required_to_save: 100,
         size_of_messages_required_to_save: segment_size,

Reply via email to