krishvishal commented on code in PR #3975:
URL: https://github.com/apache/iggy/pull/3975#discussion_r3907466505
##########
core/configs/src/server_config/partition.rs:
##########
@@ -177,6 +191,16 @@ impl Validatable<ConfigurationError> for PartitionConfig {
);
return Err(ConfigurationError::InvalidConfigurationValue);
}
+ if self.offset_reservation_lease == 0
Review Comment:
Fixed. Zero, above-ceiling and at-ceiling cases added, mirroring the other
two knobs.
##########
core/integration/tests/cluster/crash_offset_reuse.rs:
##########
@@ -115,20 +110,63 @@ async fn
given_confirmed_sends_below_flush_threshold_when_a_solo_node_is_killed_
&TopicCreateOptions {
partitions_count: Some(1),
message_expiry: Some(IggyExpiry::NeverExpire),
+ messages_required_to_save,
..TopicCreateOptions::default()
},
)
.await
.expect("create topic");
+}
- let acked = produce_acked(&client, "pre-crash", PRE_CRASH_SENDS).await;
- let highest_confirmed = *acked.last().expect("confirmed sends");
- drop(client);
+/// Base offsets of every segment file under `root`, from the file names, which
+/// are the on-disk claim about where each range begins.
+///
+/// Reading them is the only way to assert the shape the re-anchor produces. A
+/// black-box offset assertion passes either way on the boot that WRITES the
+/// wrong shape; the cost only lands on the boot that reads it back, where the
+/// walk refuses and the solo arm tombstones the partition.
+fn segment_base_offsets(root: &Path) -> Vec<u64> {
+ let mut offsets = Vec::new();
+ let mut stack = vec![root.to_path_buf()];
+ while let Some(dir) = stack.pop() {
+ let Ok(entries) = std::fs::read_dir(&dir) else {
+ continue;
+ };
+ for entry in entries.flatten() {
+ let path = entry.path();
+ if path.is_dir() {
+ stack.push(path);
+ } else if path.extension().is_some_and(|extension| extension ==
"log")
+ && let Some(stem) = path.file_stem().and_then(|stem|
stem.to_str())
+ && let Ok(offset) = stem.parse::<u64>()
+ {
+ offsets.push(offset);
+ }
+ }
+ }
+ offsets.sort_unstable();
+ offsets
+}
+/// Kill the node, bring it back, and return a client onto the restarted one.
+async fn crash_and_recover(harness: &mut TestHarness) -> IggyClient {
harness.kill_node(0).expect("SIGKILL the only node");
harness.restart_node(0).expect("restart it");
+ wait_until_serving(harness, SERVE_TIMEOUT).await
+}
- let client = wait_until_serving(harness, SERVE_TIMEOUT).await;
+#[iggy_harness(cluster_nodes = 1)]
Review Comment:
Fixed. Added a flush-then-graceful-stop cluster case here, and
`ensure_initial_segment` over a non-zero reservation in `partition_helpers`.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -2132,6 +2517,40 @@ where
offset,
}
} else {
+ // Fence AHEAD of the pipeline for a mint, not only at the
mint.
+ // A refusal at the mint arrives after the sequencer took the
op,
+ // where the only honest answer left is to fence the partition
and
+ // take the node down (`on_replicate`). Here the request has
+ // entered nothing, so a transient disk fault costs the client
one
+ // retry instead of costing the process its life.
+ //
+ // `TransientNotAccepted`, per its contract: nothing was
admitted,
+ // so the client may re-issue anywhere without double-apply
risk.
+ // It does make the SDK recheck the leader and walk the roster,
+ // which finds no better node when the fault is this one's
disk --
+ // wasteful, but the weaker code would claim an unknown outcome
+ // for a request that provably has none.
+ if message.header().operation == Operation::SendMessages
+ && let Some(ceiling) = self.request_mint_ceiling(&message)
Review Comment:
Fixed, and this was the one that mattered. `request_mint_ceiling` reads the
batch header instead of the verifying decode, so the fence fires on ordinary
sends again; covered by a test that pins the premise.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -998,6 +1208,158 @@ where
self.write_superblock(superblock.as_ref(), frontier).await
}
+ /// Whether the offset reservation is close enough to being consumed that
it
+ /// should be extended NOW, off the append path.
+ ///
+ /// The append fence is correct but badly placed: it writes the superblock
+ /// inline in the shard's frame pump, where the consensus tick is a sibling
+ /// arm, so its two fsyncs delay heartbeats for every group on the core.
The
+ /// fix is to make the fence's fast path
+ /// (`durable_offset_reserved > end_offset`) the only path it ever takes
under
+ /// load, by extending from the tick instead.
+ ///
+ /// HALF a block of headroom, which is a wide margin on purpose:
over-claiming
+ /// costs nothing but offset space, while arriving late puts the write
back on
+ /// the append path. A partition that has never minted is skipped -- its
first
+ /// append legitimately pays for the first claim, and extending every idle
+ /// partition at boot would write a superblock per partition for nothing.
+ #[must_use]
+ pub fn needs_offset_reservation_extension(&self) -> bool {
+ if self.consensus.replica_count() > 1
+ || self.superblock.is_none()
+ || !self.should_increment_offset
+ {
+ return false;
+ }
+ let headroom = self
+ .durable_offset_reserved
+ .get()
+ .saturating_sub(self.mint_frontier());
+ headroom < self.offset_reservation_lease / 2
+ }
+
+ /// Extend the reservation a full block past the current append point.
+ ///
+ /// Pairs with [`Self::needs_offset_reservation_extension`]; the caller is
the
+ /// shard tick, so this write is off the append path. A failure needs no
+ /// handling beyond the writer's own logging and backoff: the fence at the
+ /// mint is still there, and it is what refuses the append if the ceiling
+ /// never caught up.
+ #[allow(clippy::future_not_send)]
+ pub async fn extend_offset_reservation(&self) -> bool {
+ let Some(superblock) = self.superblock.as_ref().map(Rc::clone) else {
+ return true;
+ };
+ if self.superblock_write_is_backed_off() {
+ return false;
+ }
+ let _superblock_guard = self.superblock_lock.acquire().await;
+ self.write_offset_claim(superblock.as_ref(), self.mint_frontier())
+ .await
+ }
+
+ /// Upper bound on the offsets a pending `SendMessages` request will mint,
for
+ /// fencing it BEFORE it enters the pipeline.
+ ///
+ /// `project` assigns an op, not a base offset, so the exact range is
unknown
+ /// until the mint runs under `write_lock`. This is deliberately loose: a
+ /// request pipelined behind others can land above it, and the fence at the
+ /// mint stays as the exact check. It does not need to be tight -- the
claim
+ /// runs a whole lease block past whatever it is handed, so one of these
+ /// covers every batch in flight unless a run of them crosses a block
+ /// boundary.
+ ///
+ /// `None` when the body is not one canonical batch, which
+ /// `convert_request_message` has already rejected by the time this runs.
+ fn request_mint_ceiling(&self, message: &Message<RoutedRequestHeader>) ->
Option<u64> {
+ let body = message
+ .as_slice()
+
.get(std::mem::size_of::<RoutedRequestHeader>()..message.header().size as
usize)?;
+ let count = decode_batch_slice(body).ok()?.message_count();
Review Comment:
Fixed. `BatchHeader::decode` for `message_count`, and an early
`replica_count() > 1` return before the body is touched.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]