This is an automated email from the ASF dual-hosted git repository.
hubcio pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iggy.git
The following commit(s) were added to refs/heads/master by this push:
new 9c6451d62 fix(consensus): keep a late retry's reply instead of
refusing it (#3945)
9c6451d62 is described below
commit 9c6451d62f483dea8ebb4c8e08f0c91ecde8e6f2
Author: Grzegorz Koszyk <[email protected]>
AuthorDate: Thu Aug 27 14:27:04 2026 +0200
fix(consensus): keep a late retry's reply instead of refusing it (#3945)
---
core/consensus/src/client_table.rs | 254 ++++++++++++++++++---
core/consensus/src/metadata_helpers.rs | 8 +-
.../tests/cluster/client_table_adversarial.rs | 56 +++--
3 files changed, 260 insertions(+), 58 deletions(-)
diff --git a/core/consensus/src/client_table.rs
b/core/consensus/src/client_table.rs
index 46c9d9451..67cb42e25 100644
--- a/core/consensus/src/client_table.rs
+++ b/core/consensus/src/client_table.rs
@@ -78,6 +78,11 @@ impl CachedReply {
fn as_bytes(&self) -> &[u8] {
self.bytes.as_slice()
}
+
+ /// Wire size of this reply, the unit [`REPLY_RING_RETENTION_BYTES`]
budgets.
+ fn byte_len(&self) -> usize {
+ self.bytes.len()
+ }
}
/// Reserved request number for
[`Operation::Register`](iggy_binary_protocol::Operation::Register).
@@ -133,10 +138,64 @@ pub const CLIENTS_TABLE_SLOT_MAX: usize = 1 << 16;
/// live client can be waiting for is its latest (`request == watermark`).
/// Older entries answer old retransmits and post-rebind stragglers with the
/// original bytes instead of a bare "already applied"; losing one
-/// degrades the answer, never correctness. In-memory only: ring contents are
-/// refcount bumps and are never persisted or transferred.
+/// degrades the answer, never correctness.
+///
+/// This many replies are retained unconditionally, whatever they weigh, which
+/// is what bounds the memory a client holding megabyte replies can pin.
+/// Retention past it is governed by [`REPLY_RING_RETENTION_BYTES`]: a client
+/// sending small operations -- the common case -- keeps a far deeper replay
+/// history for the same memory, so a retry that arrives late still replays its
+/// original bytes instead of drawing a bare "already applied".
+///
+/// # What the deeper retention is worth
+///
+/// It is a live-memory property of the replica that served the request, and it
+/// survives neither rebuild path. `transferable_replies` ships at most this
+/// many replies in a state transfer, and [`ClientTable::from_snapshot`]
+/// restores only each entry's latest reply from a checkpoint, so a retry that
+/// would have replayed on the serving replica draws
+/// [`RequestStatus::AlreadyApplied`] once the receiving replica installs a
+/// transfer, or once this one restarts. Bounded cache depth, not a durable
+/// guarantee; `state_transfer_cuts_retention_back_to_the_floor` and
+/// `snapshot_drops_stale_ring_replies_but_keeps_at_most_once` pin each path.
+///
+/// What is durable is at-most-once itself: the watermark rides both paths
+/// intact, so a lost reply costs the caller its result bytes and never
+/// re-executes the operation.
pub const REPLY_RING_CAPACITY: usize = 5;
+/// Byte budget for the replies retained past [`REPLY_RING_CAPACITY`].
+///
+/// Deep retention exists for the slow retrier: a request whose reply aged out
+/// can only be answered "already applied, reply gone", which tells the caller
+/// its operation succeeded but hands back no result. Budgeting in bytes rather
+/// than in replies puts the depth where it is cheapest: a session sending
+/// metadata operations keeps a long history, one pulling large batches keeps
+/// none past the floor.
+///
+/// # Depth
+///
+/// A reply is never shorter than its 256-byte [`ReplyHeader`], so this budget
+/// is also the only thing bounding the ring's length: 8 KiB / 256 B = 32
+/// replies at the deepest, against a floor of [`REPLY_RING_CAPACITY`].
+///
+/// The 32 is a chosen bound, not a measured one. The SDK holds one request in
+/// flight per session, so what has to fit is the number of newer requests the
+/// same session commits between a reply going unacknowledged and its retry
+/// landing, and nothing in the tree measures that today.
+///
+/// # Cost
+///
+/// Not a wash on the common case. The common metadata reply is header-only, so
+/// per-slot retention goes from 5 x 256 B = 1.25 KiB to 8 KiB, a 6.4x rise: a
+/// saturated table at the default `clients_table_max` of 8192 goes from 10 MiB
+/// to 64 MiB, and from 41k live [`Frozen`] buffers to 262k. At the
+/// [`CLIENTS_TABLE_SLOT_MAX`] ceiling it is 512 MiB. Replies carrying a
payload
+/// exhaust the budget sooner and cost proportionally less; above roughly
+/// 1.6 KiB apiece the [`REPLY_RING_CAPACITY`] floor dominates and this budget
+/// adds nothing at all.
+pub const REPLY_RING_RETENTION_BYTES: usize = 8 * 1024;
+
/// What eviction and disconnect cleanup keep after reclaiming an entry's slot.
///
/// At-most-once needs only the fence: the watermark says which request numbers
@@ -1338,9 +1397,10 @@ pub enum ClientTableWireError {
/// slot occupied but unindexed, which desynchronizes the capacity check in
/// [`ClientTable::commit_register`] from the actual occupancy.
DuplicateClientId { slot: usize, client_id: u128 },
- /// A reply ring longer than [`REPLY_RING_CAPACITY`]. `push_latest` only
- /// evicts on equality, so an over-capacity ring grows without bound, and
- /// `encode` writes its length as a `u8`.
+ /// A reply ring longer than [`REPLY_RING_CAPACITY`], which is every reply
+ /// `transferable_replies` ever writes. `encode` writes the length as a
+ /// `u8`, and a peer that sends more replies than this crate transfers is
+ /// reporting state this one cannot have produced.
RingTooLong { slot: usize, len: u8, max: usize },
}
@@ -1443,8 +1503,7 @@ impl ClientTable {
.map(|entry| {
ENCODED_ENTRY_FIXED_LEN
+ entry
- .ring
- .iter()
+ .transferable_replies()
.map(|reply| size_of::<u32>() + reply.bytes.len())
.sum::<usize>()
})
@@ -1470,8 +1529,8 @@ impl ClientTable {
out.extend_from_slice(&entry.user_id.to_le_bytes());
out.extend_from_slice(&entry.watermark.to_le_bytes());
out.extend_from_slice(&entry.watermark_checksum.to_le_bytes());
- out.push(entry.ring.len() as u8);
- for reply in &entry.ring {
+ out.push(entry.transferable_replies().count() as u8);
+ for reply in entry.transferable_replies() {
let bytes = reply.bytes.as_slice();
out.extend_from_slice(&(bytes.len() as u32).to_le_bytes());
out.extend_from_slice(bytes);
@@ -1556,10 +1615,10 @@ impl ClientTable {
return Err(ClientTableWireError::EmptyRing);
}
// The artifact checksum only proves the bytes survived transit; it
- // says nothing about the peer that computed them. An over-capacity
- // ring is admitted forever after (`push_latest` evicts only on
- // equality) and eventually wraps `encode`'s `u8` length, making
- // this table permanently un-transferable onward.
+ // says nothing about the peer that computed them. Bounding at what
+ // `transferable_replies` emits keeps an installed ring inside the
+ // unconditional floor, so it satisfies `trim_ring`'s byte budget
on
+ // arrival and needs no trimming of its own.
if usize::from(ring_len) > REPLY_RING_CAPACITY {
return Err(ClientTableWireError::RingTooLong {
slot: slot_idx,
@@ -1679,6 +1738,20 @@ impl ClientEntry {
.expect("ring is never empty after registration")
}
+ /// The newest replies a state transfer carries, oldest first.
+ ///
+ /// Capped at [`REPLY_RING_CAPACITY`] rather than shipping whatever
+ /// retention holds locally: an artifact a recovering node has to fetch is
+ /// worth keeping small, and the deeper history rebuilds itself from the
+ /// receiver's own commits. The cost is a late retrier's result bytes on
the
+ /// transferred node, never its at-most-once fence, which is the bound
+ /// [`REPLY_RING_CAPACITY`] documents.
+ fn transferable_replies(&self) -> impl Iterator<Item = &CachedReply> {
+ self.ring
+ .iter()
+ .skip(self.ring.len().saturating_sub(REPLY_RING_CAPACITY))
+ }
+
/// Cached reply whose `request` matches (scan order is irrelevant
/// because request numbers in the ring are unique).
fn find_cached(&self, request: u64) -> Option<&CachedReply> {
@@ -1687,14 +1760,28 @@ impl ClientEntry {
.find(|cached| cached.header().request == request)
}
- /// Push the newest committed reply, evicting the oldest when full, and
- /// refresh the denormalized `latest_commit`.
+ /// Push the newest committed reply, drop the oldest ones retention no
+ /// longer covers, and refresh the denormalized `latest_commit`.
fn push_latest(&mut self, cached: CachedReply) {
self.latest_commit = cached.header().commit;
- if self.ring.len() == REPLY_RING_CAPACITY {
- self.ring.pop_front();
- }
self.ring.push_back(cached);
+ self.trim_ring();
+ }
+
+ /// Drop the oldest replies once the entry holds more than
+ /// [`REPLY_RING_CAPACITY`] and exceeds [`REPLY_RING_RETENTION_BYTES`].
+ ///
+ /// The total is summed here rather than denormalized onto the entry: no
+ /// reply is shorter than a header, so the budget holds the ring to 32
+ /// entries and this is a bounded walk of buffer lengths with no header
+ /// casts, and it leaves no running total for the sites that write the ring
+ /// to drift out of sync with.
+ fn trim_ring(&mut self) {
+ let mut bytes: usize =
self.ring.iter().map(CachedReply::byte_len).sum();
+ while self.ring.len() > REPLY_RING_CAPACITY && bytes >
REPLY_RING_RETENTION_BYTES {
+ let dropped = self.ring.pop_front().expect("length checked above");
+ bytes -= dropped.byte_len();
+ }
}
}
@@ -2205,6 +2292,29 @@ mod tests {
make_reply_with_checksum(client, request, commit, 0)
}
+ /// A reply heavy enough that two of them exhaust
+ /// [`REPLY_RING_RETENTION_BYTES`], so only the floor keeps it cached.
+ #[allow(clippy::cast_possible_truncation)]
+ fn make_big_reply(client: u128, request: u64, commit: u64) ->
Message<ReplyHeader> {
+ let header_size = std::mem::size_of::<ReplyHeader>();
+ let size = header_size + REPLY_RING_RETENTION_BYTES;
+ let mut msg = Message::<ReplyHeader>::new(size);
+ let header = bytemuck::checked::try_from_bytes_mut::<ReplyHeader>(
+ &mut msg.as_mut_slice()[..header_size],
+ )
+ .expect("zeroed bytes are valid");
+ *header = ReplyHeader {
+ client,
+ request,
+ commit,
+ size: size as u32,
+ command: Command::Reply,
+ operation: Operation::SendMessages,
+ ..ReplyHeader::default()
+ };
+ msg
+ }
+
#[allow(clippy::cast_possible_truncation)]
fn make_reply_with_checksum(
client: u128,
@@ -2600,28 +2710,75 @@ mod tests {
// Below-watermark duplicate whose reply aged out of the ring is refused
// execution with nothing to replay.
#[test]
- fn check_request_below_watermark_past_ring_is_already_applied() {
+ fn check_request_below_watermark_past_retention_is_already_applied() {
let (mut table, epoch) = table_with_client();
- // Requests 1..=6: request 1's reply is displaced beyond the ring
- // (capacity 5 holds 2..=6 once 6 commits; the register reply and
- // request 1 aged out first).
- for request in 1..=6u64 {
+ // Enough small replies to exhaust the byte budget several times over,
+ // so the oldest are certain to have been dropped.
+ let requests = (REPLY_RING_RETENTION_BYTES / size_of::<ReplyHeader>()
+ 8) as u64;
+ for request in 1..=requests {
table.commit_reply(1, TEST_USER_ID, make_reply_for(1, request, 10
+ request));
}
match table.check_request(1, epoch, 1, 0) {
RequestStatus::AlreadyApplied { request, watermark } => {
assert_eq!(request, 1);
- assert_eq!(watermark, 6);
+ assert_eq!(watermark, requests);
}
other => panic!("expected AlreadyApplied, got {other:?}"),
}
- // The oldest retained entry still answers.
- match table.check_request(1, epoch, 2, 0) {
- RequestStatus::Duplicate(cached) =>
assert_eq!(cached.header().request, 2),
+ // The newest still answers with its own bytes.
+ match table.check_request(1, epoch, requests, 0) {
+ RequestStatus::Duplicate(cached) =>
assert_eq!(cached.header().request, requests),
other => panic!("expected Duplicate, got {other:?}"),
}
}
+ // A retry that arrives after more commits than the floor holds still gets
+ // its original bytes back: retention past the floor is budgeted in bytes,
+ // and small replies are what a late retrier usually has outstanding.
+ #[test]
+ fn a_late_retry_replays_while_the_retention_budget_holds_it() {
+ let (mut table, epoch) = table_with_client();
+ let requests = REPLY_RING_CAPACITY as u64 + 2;
+ for request in 1..=requests {
+ table.commit_reply(1, TEST_USER_ID, make_reply_for(1, request, 10
+ request));
+ }
+
+ match table.check_request(1, epoch, 1, 0) {
+ RequestStatus::Duplicate(cached) => {
+ assert_eq!(cached.header().request, 1);
+ assert_eq!(
+ cached.header().commit,
+ 11,
+ "the original reply, not a re-run"
+ );
+ }
+ other => panic!("expected the original reply to replay, got
{other:?}"),
+ }
+ }
+
+ // Heavy replies stay bounded by the floor, so the deeper retention cannot
+ // be turned into a memory amplifier by a client polling large batches.
+ #[test]
+ fn heavy_replies_are_retained_only_to_the_floor() {
+ let (mut table, epoch) = table_with_client();
+ let requests = REPLY_RING_CAPACITY as u64 + 2;
+ for request in 1..=requests {
+ table.commit_reply(1, TEST_USER_ID, make_big_reply(1, request, 10
+ request));
+ }
+
+ // The floor counts the register reply out: it aged out first, leaving
+ // the last REPLY_RING_CAPACITY app replies.
+ let oldest_retained = requests - REPLY_RING_CAPACITY as u64 + 1;
+ assert!(matches!(
+ table.check_request(1, epoch, oldest_retained - 1, 0),
+ RequestStatus::AlreadyApplied { .. }
+ ));
+ assert!(matches!(
+ table.check_request(1, epoch, oldest_retained, 0),
+ RequestStatus::Duplicate(_)
+ ));
+ }
+
// Dedup across view change. Backup inherits client_table via
// commit_journal; on failover, retry must return ORIGINAL cached reply
// (same request, same commit op), no re-execution. Pipeline state is
@@ -2949,6 +3106,43 @@ mod tests {
assert_eq!(decoded.encode(), encoded);
}
+ // Deep retention is a live-memory property, not a transferred one: the
+ // artifact carries the floor, so a retry the serving replica would have
+ // replayed loses its bytes on the receiver. The fence still rides along,
so
+ // the answer degrades and the operation is still never re-executed.
+ #[test]
+ fn state_transfer_cuts_retention_back_to_the_floor() {
+ let (mut table, epoch) = table_with_client();
+ let requests = REPLY_RING_CAPACITY as u64 + 3;
+ for request in 1..=requests {
+ table.commit_reply(1, TEST_USER_ID, make_reply_for(1, request, 10
+ request));
+ }
+ // The whole run is still cached locally: retention is byte-budgeted
and
+ // these replies are headers.
+ assert!(matches!(
+ table.check_request(1, epoch, 1, 0),
+ RequestStatus::Duplicate(_)
+ ));
+
+ let decoded = ClientTable::decode(&table.encode(),
10).expect("roundtrip decodes");
+
+ assert_eq!(decoded.get_watermark(1), Some(requests));
+ let oldest_transferred = requests - REPLY_RING_CAPACITY as u64 + 1;
+ match decoded.check_request(1, epoch, oldest_transferred - 1, 0) {
+ RequestStatus::AlreadyApplied { request, watermark } => {
+ assert_eq!(request, oldest_transferred - 1);
+ assert_eq!(watermark, requests, "the fence survives the
transfer");
+ }
+ other => panic!("expected AlreadyApplied past the transferred
floor, got {other:?}"),
+ }
+ match decoded.check_request(1, epoch, oldest_transferred, 0) {
+ RequestStatus::Duplicate(cached) => {
+ assert_eq!(cached.header().request, oldest_transferred);
+ }
+ other => panic!("expected Duplicate, got {other:?}"),
+ }
+ }
+
// The denormalized `latest_commit` is rebuilt from the ring on decode, so
// eviction ranks a transferred table exactly like the original.
#[test]
@@ -3079,9 +3273,9 @@ mod tests {
));
}
- // `push_latest` evicts only on equality, so an over-capacity ring is
- // admitted permanently and grows on every later reply until `encode`'s u8
- // length wraps and the table stops being transferable at all.
+ // A ring longer than `transferable_replies` emits comes from a peer this
+ // one cannot model; admitting it would eventually wrap `encode`'s u8
length
+ // and leave the table untransferable onward.
#[test]
fn decode_rejects_a_ring_longer_than_capacity() {
let mut table = ClientTable::new(1);
diff --git a/core/consensus/src/metadata_helpers.rs
b/core/consensus/src/metadata_helpers.rs
index c7fad442c..43af22a87 100644
--- a/core/consensus/src/metadata_helpers.rs
+++ b/core/consensus/src/metadata_helpers.rs
@@ -468,7 +468,7 @@ pub async fn send_eviction_to_client<B, P>(
#[cfg(test)]
mod tests {
use super::*;
- use crate::client_table::REGISTER_REQUEST_ID;
+ use crate::client_table::{REGISTER_REQUEST_ID, REPLY_RING_RETENTION_BYTES};
use crate::{CLIENTS_TABLE_MAX, LocalPipeline};
use iggy_binary_protocol::{Command, Operation, ReplyHeader};
use message_bus::SendError;
@@ -740,8 +740,10 @@ mod tests {
client_table
.borrow_mut()
.commit_register(client_id, ACTING_USER_ID, initial_reply);
- // Ring capacity is 5, so request 1's reply is displaced once 6
commits.
- for request in 1..=6u64 {
+ // Enough replies to exhaust the retention budget, so request 1's is
+ // certain to have been dropped.
+ let requests = (REPLY_RING_RETENTION_BYTES / size_of::<ReplyHeader>()
+ 8) as u64;
+ for request in 1..=requests {
let reply =
synthesize_send_messages_reply(&consensus, client_id, request,
100 + request);
client_table
diff --git a/core/integration/tests/cluster/client_table_adversarial.rs
b/core/integration/tests/cluster/client_table_adversarial.rs
index 90c1ae245..5181c8fea 100644
--- a/core/integration/tests/cluster/client_table_adversarial.rs
+++ b/core/integration/tests/cluster/client_table_adversarial.rs
@@ -18,17 +18,22 @@
//! Adversarial specs against the VSR client table's at-most-once guarantees.
//!
//! Both assert the dedup contract a retrying client needs at the table's two
-//! resource edges. The capacity one now passes; the reply-ring one is still a
-//! RED SPEC, expected to FAIL.
+//! resource edges.
//!
//! 1. Capacity: a full table evicts the entry with the oldest commit. Eviction
//! keeps that client's request watermark (and the watermark's reply when
the
//! ring still held it), so a client that was merely quiet re-registers and
//! its retry of an already-committed request id is answered, not
re-executed.
-//! 2. Reply ring: each entry retains only its `REPLY_RING_CAPACITY` most
-//! recent committed replies. A retry of a request whose reply aged out is
-//! refused with the terminal `RequestAlreadyApplied` and no result payload,
-//! indistinguishable from a rejection to the caller.
+//! 2. Reply retention: an entry keeps `REPLY_RING_CAPACITY` replies whatever
+//! they weigh and older ones for as long as they fit
+//! `REPLY_RING_RETENTION_BYTES`, so a retry arriving more commits late than
+//! the floor alone would hold still replays its original bytes. The depth
is
+//! bounded, not unlimited: a retry that outlives the byte budget is still
+//! answered with the terminal `RequestAlreadyApplied` and no result
payload,
+//! which the caller cannot tell apart from a rejection. That edge, and the
+//! fact that neither a state transfer nor a restart carries the deeper
+//! history, are pinned by the consensus crate's unit tests; this file pins
+//! the depth the budget buys on a live server.
//!
//! The frames are hand-crafted on raw TCP sockets, same technique and frame
//! builders as the clients-table restart tests, because the churn needs
@@ -145,19 +150,21 @@ async fn
given_a_low_client_table_cap_when_connects_churn_should_keep_a_live_ded
}
}
-/// RED SPEC, expected to FAIL: a retry that falls off the reply ring must not
-/// be terminally rejected without a result.
+/// A retry that arrives after more commits than the unconditional floor holds
+/// must still replay its original result while the byte budget covers it.
///
-/// Request 1 commits, then later requests on the same session push its reply
-/// out of the ring. The retry of request 1 is then answered with the terminal
-/// `RequestAlreadyApplied` code and no result payload. The
-/// caller cannot distinguish "your operation succeeded, the reply aged out"
-/// from "your operation was rejected", so a slow retrier is forced to treat a
-/// success as a failure.
-// TODO(hubcio): fix this test
-#[ignore = "replay past the reply ring draws terminal RequestAlreadyApplied,
no result payload"]
+/// Request 1 commits, then more requests than `REPLY_RING_CAPACITY` commit on
+/// the same session. Retention past the floor is budgeted in bytes and these
+/// replies are small, so request 1's is still cached and its retry answers
with
+/// the original bytes. Answering `RequestAlreadyApplied` instead tells the
+/// caller its operation succeeded while handing back no result, which a slow
+/// retrier cannot tell apart from a rejection.
+///
+/// Scoped to what the budget covers. A retry late enough to outlive it, and
one
+/// that lands on a replica rebuilt by transfer or restart, still draw the
+/// terminal code; the consensus crate's unit tests own those cases.
#[iggy_harness(cluster_nodes = 1)]
-async fn
given_replays_past_the_reply_ring_when_a_request_id_falls_off_should_not_terminally_reject(
+async fn
given_a_retry_past_the_reply_floor_when_the_retention_budget_still_holds_it_should_replay_from_cache(
harness: &mut TestHarness,
) {
let addr = tcp_addr(harness);
@@ -165,8 +172,8 @@ async fn
given_replays_past_the_reply_ring_when_a_request_id_falls_off_should_no
let aged_payload = create_stream_payload("adv-l-aged");
let committed = commit_request(&mut stream, CLIENT_A, session, 1,
&aged_payload).await;
- // One more commit than the ring holds, so request 1's reply is evicted
- // even though the register reply seeded a slot of its own.
+ // More commits than the unconditional floor holds, so passing depends on
+ // the byte-budgeted retention rather than on the floor alone.
let later_requests = REPLY_RING_CAPACITY as u64 + 1;
for index in 0..later_requests {
let payload = create_stream_payload(&format!("adv-l-filler-{index}"));
@@ -180,12 +187,11 @@ async fn
given_replays_past_the_reply_ring_when_a_request_id_falls_off_should_no
assert_replayed_from_cache(&committed, &replayed, 1);
}
other => panic!(
- "a committed request replayed past the reply ring is terminally
rejected: \
- request 1 was applied and confirmed, but after {later_requests}
newer commits \
- its reply aged out of the {REPLY_RING_CAPACITY}-deep ring and the
server \
- answered the terminal RequestAlreadyApplied code with no result
payload, which the \
- client cannot distinguish from a rejection of an operation that
in fact \
- succeeded; got {other:?}"
+ "a committed request replayed past the reply floor lost its
result: request 1 \
+ was applied and confirmed, and {later_requests} newer commits of
this size fit \
+ the retention budget, so its reply must still replay; a terminal \
+ RequestAlreadyApplied tells the client its operation succeeded
while handing \
+ back no result, which it cannot tell apart from a rejection; got
{other:?}"
),
}
}