hubcio commented on code in PR #3975:
URL: https://github.com/apache/iggy/pull/3975#discussion_r3893568620
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -1812,14 +2037,18 @@ where
return Err(IggyError::CannotAppendMessage);
}
- let dirty_offset = if self.should_increment_offset {
+ let next = if self.should_increment_offset {
self.dirty_offset
.load(Ordering::Relaxed)
.checked_add(1)
.ok_or(IggyError::CannotAppendMessage)?
} else {
0
};
+ // Only here: this is the only path that mints. A backup re-stamps what
+ // the primary sends (`append_received_send_messages_to_journal`) and
+ // must follow it exactly, so raising ITS counter would fork the group.
+ let dirty_offset = next.max(self.mint_floor());
Review Comment:
`mint_floor()` clears the flag before the fence or the journal append can
fail, so the retry mints below the reservation and the fence fast path waves it
through - a confirmed offset re-minted. clear the flag next to
`dirty_offset.store` instead.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -4189,6 +4431,121 @@ where
Ok(())
}
+ /// Re-anchor the append point after boot re-seeded the offset counter
above
+ /// what the recovered segment chain holds.
+ ///
+ /// A hole INSIDE a segment is not survivable: `recover_segment_bounds`
walks
+ /// a segment from its FILENAME with a running `expected_offset`, so the
next
+ /// boot truncates the post-hole suffix and the counter falls back below
what
+ /// was already confirmed, undoing the fix on the second crash. On a
segment
+ /// BOUNDARY every reader copes -- absolute offsets in the index,
+ /// `disk_poll_start` walking on into later segments, a contiguity guard
that
+ /// only looks within one file.
+ ///
+ /// So an empty tail is unlinked (its name claims a range it does not hold)
+ /// and a sized tail, the only copy of its messages, is sealed with a fresh
+ /// segment planted at the frontier. An empty chain is left to the caller's
Review Comment:
`ensure_initial_segment` plants at `offset_frontier()`, still 0 here, while
the first mint lands at the floor - the hole inside a segment this doc calls
unsurvivable; the boot after the first flush refuses it. plant the empty chain
at `frontier` here too.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -975,6 +1097,25 @@ where
self.consensus.clock_realtime_micros() <
self.superblock_retry_after_micros.get()
}
+ /// Drop the reservation back onto the frontier, once a graceful flush has
+ /// made the segments account for every offset this replica confirmed.
+ ///
+ /// The reservation is there for the crash case, where they do not. Left
+ /// standing it would make every ordinary restart resume a lease block
+ /// higher and hole the offset space for nothing.
+ ///
+ /// Callers must have flushed FIRST, and must not call this when the flush
+ /// failed: the claim it makes is precisely that the flush succeeded.
+ #[allow(clippy::future_not_send)]
+ #[must_use = "the bool is the durability verdict; a failed collapse leaves
a gap"]
+ pub async fn collapse_offset_reservation(&self) -> bool {
+ let frontier = self.offset_frontier();
Review Comment:
after a crash-restart with no sends the counter is still below the armed
floor, so a clean stop writes `(0,0)` or a reservation under the planted
segment - next boot re-mints offset 0 or refuses the chain. collapse to
`offset_frontier().max(armed_mint_floor())`.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -2768,6 +2997,19 @@ where
.checked_add(u64::from(batch_messages_count) - 1)
.ok_or(IggyError::CannotAppendMessage)?;
+ // Past this line the offsets are in the journal, hence committable,
+ // pollable, confirmable and forwardable, and no later gate can take
+ // them back. See [`Self::reserve_offsets_through`].
+ //
+ // LOCK ORDER: the caller holds `write_lock` and this takes
+ // `superblock_lock` under it. The install path takes them in the
+ // reverse order, safely only because `reset_offset_frontier_at` drops
+ // `superblock_lock` before `try_install` takes `write_lock`. Never
hold
+ // `superblock_lock` across a `write_lock` acquire.
+ if !self.reserve_offsets_through(last_dirty_offset).await {
Review Comment:
on a solo primary a failed or backed-off write here makes `on_replicate`
drop the op after the sequencer took it, so nothing commits again until restart
and every produce times out. reserve in `on_request` before the pipeline push
and reply transient.
##########
core/partitions/src/state_transfer.rs:
##########
@@ -1994,7 +1994,13 @@ where
let purge_advances = offsets_wire.purge_generation >
committed_purge_generation
|| (self.applied_purge_generation < committed_purge_generation
&& offsets_wire.next_offset == 0);
- let local_next_offset = self.offset_frontier();
+ // What this replica HOLDS, not the append counter: the two diverge
once
+ // the counter is seeded from an offset reservation, which stands a
lease
+ // block above the last message anywhere and names none of them.
Against
+ // that, every legitimate offer inside the block reads as a rewind, and
+ // the replica that restarted with a reservation is exactly the one
+ // needing the transfer, so it would cycle refusal -> backoff forever.
+ let local_next_offset = self.held_offset_frontier();
Review Comment:
an empty chain installed at frontier N has `held_offset_frontier() == 0`, so
the guard is skipped and a stale offer below N rewinds the counter. max it with
`installed_frontier`.
##########
core/integration/tests/cluster/crash_offset_reuse.rs:
##########
@@ -97,8 +97,64 @@ async fn wait_until_serving(harness: &TestHarness, budget:
Duration) -> IggyClie
}
}
-// TODO(hubcio): fix this test
-#[ignore = "confirmed offsets re-minted after a crash; no durable offset
watermark"]
+/// Create the stream and its single-partition topic.
+async fn create_topic(client: &IggyClient, messages_required_to_save:
Option<u32>) {
+ client
+ .create_stream(STREAM_NAME)
+ .await
+ .expect("create stream");
+ client
+ .create_topic(
+ &Identifier::named(STREAM_NAME).unwrap(),
+ TOPIC_NAME,
+ &TopicCreateOptions {
+ partitions_count: Some(1),
+ message_expiry: Some(IggyExpiry::NeverExpire),
+ messages_required_to_save,
+ ..TopicCreateOptions::default()
+ },
+ )
+ .await
+ .expect("create topic");
+}
+
+/// 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:
+/// recovery tolerates a discontiguity inside a segment whenever the index
+/// survives, so a black-box offset assertion passes either way and the wrong
+/// shape sits there until an index is torn.
+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
+}
+
#[iggy_harness(cluster_nodes = 1)]
Review Comment:
this re-inlines `create_topic` and `crash_and_recover` from above.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -4189,6 +4431,121 @@ where
Ok(())
}
+ /// Re-anchor the append point after boot re-seeded the offset counter
above
+ /// what the recovered segment chain holds.
+ ///
+ /// A hole INSIDE a segment is not survivable: `recover_segment_bounds`
walks
+ /// a segment from its FILENAME with a running `expected_offset`, so the
next
+ /// boot truncates the post-hole suffix and the counter falls back below
what
+ /// was already confirmed, undoing the fix on the second crash. On a
segment
+ /// BOUNDARY every reader copes -- absolute offsets in the index,
+ /// `disk_poll_start` walking on into later segments, a contiguity guard
that
+ /// only looks within one file.
+ ///
+ /// So an empty tail is unlinked (its name claims a range it does not hold)
+ /// and a sized tail, the only copy of its messages, is sealed with a fresh
+ /// segment planted at the frontier. An empty chain is left to the caller's
+ /// `ensure_initial_segment`.
+ ///
+ /// # Errors
+ /// [`IggyError`] when the fresh segment cannot be created, leaving the
+ /// partition without a serviceable chain.
+ #[allow(clippy::future_not_send)]
+ pub async fn reanchor_to_offset_frontier(
+ &mut self,
+ config: &PartitionsConfig,
+ ) -> Result<(), IggyError> {
+ // Where the next append will land: the counter, or an armed mint floor
+ // above it. The floor is the whole reason a hole can appear, so
+ // anchoring to the counter alone would leave the chain as unprepared.
+ let frontier = self.offset_frontier().max(self.armed_mint_floor());
+ if frontier == 0 {
+ return Ok(());
+ }
+ let namespace = self.namespace();
+ let mut retired = 0usize;
+ while let Some(segment) = self.log.segments().last() {
+ if segment.size.as_bytes_u64() > 0 || segment.start_offset >=
frontier {
+ break;
+ }
+ let Some((segment, mut storage)) = self.log.retire_back() else {
+ break;
+ };
+ let (messages_path, index_path) =
storage.segment_and_index_paths();
+ let _ = storage.shutdown();
+ drop(storage);
+ for path in messages_path.into_iter().chain(index_path) {
+ match compio::fs::remove_file(&path).await {
+ Ok(()) => {}
+ Err(error) if error.kind() == std::io::ErrorKind::NotFound
=> {}
+ Err(error) => {
+ warn!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ path = %path,
+ %error,
+ "failed to unlink a stale empty segment during the
boot re-anchor"
+ );
+ }
+ }
+ }
+ tracing::info!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ start_offset = segment.start_offset,
+ offset_frontier = frontier,
+ "unlinked an empty segment named below the restored offset
frontier"
+ );
+ retired += 1;
+ }
+ // Durable before anything is planted beside them: a crash in between
+ // would boot the stale name back into the chain. No stats decrement to
+ // pair with the retire -- boot never counted the recovered chain, and
+ // `fetch_sub` does not saturate.
+ if retired > 0
+ && let Some(partition_dir) = self.partition_dir.clone()
+ && let Err(error) =
crate::state_transfer::fsync_dir(&partition_dir).await
+ {
+ warn!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ partition_dir,
+ %error,
+ "boot re-anchor could not fsync the partition dir after
unlinking"
+ );
+ }
+ // Only a SIZED tail: an empty one either just went, or is already
named
+ // at the frontier and can take the appends as it is.
+ let needs_plant = self.log.segments().last().is_some_and(|segment| {
+ segment.size.as_bytes_u64() > 0 && segment.end_offset + 1 <
frontier
Review Comment:
`saturating_add(1)`, like the sibling at 984.
##########
core/integration/tests/cluster/crash_offset_reuse.rs:
##########
@@ -97,8 +97,64 @@ async fn wait_until_serving(harness: &TestHarness, budget:
Duration) -> IggyClie
}
}
-// TODO(hubcio): fix this test
-#[ignore = "confirmed offsets re-minted after a crash; no durable offset
watermark"]
+/// Create the stream and its single-partition topic.
+async fn create_topic(client: &IggyClient, messages_required_to_save:
Option<u32>) {
+ client
+ .create_stream(STREAM_NAME)
+ .await
+ .expect("create stream");
+ client
+ .create_topic(
+ &Identifier::named(STREAM_NAME).unwrap(),
+ TOPIC_NAME,
+ &TopicCreateOptions {
+ partitions_count: Some(1),
+ message_expiry: Some(IggyExpiry::NeverExpire),
+ messages_required_to_save,
+ ..TopicCreateOptions::default()
+ },
+ )
+ .await
+ .expect("create topic");
+}
+
+/// 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:
+/// recovery tolerates a discontiguity inside a segment whenever the index
+/// survives, so a black-box offset assertion passes either way and the wrong
+/// shape sits there until an index is torn.
+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
+}
+
#[iggy_harness(cluster_nodes = 1)]
async fn
given_confirmed_sends_below_flush_threshold_when_a_solo_node_is_killed_should_not_remint_offsets(
Review Comment:
this builds the empty-chain shape but restarts only once, so it never sees
the boot that refuses it. add a send, a clean restart and a third boot, and
assert `segment_base_offsets`.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -998,6 +1139,80 @@ where
self.write_superblock(superblock.as_ref(), frontier).await
}
+ /// The lowest offset this replica may mint, given what its durable record
+ /// already permitted it to hand out. A floor, not a seed -- see
+ /// [`Self::restore_offset_frontier`]. It bites once, on the first mint
after
+ /// a crash that took acked-but-unflushed messages with it: segments and
+ /// counter come back short, and the reservation is the only surviving
+ /// witness that those offsets were confirmed.
+ ///
+ /// SOLO GROUPS ONLY. A backup validates an incoming prepare with
+ /// `base_offset == dirty_offset + 1`
+ /// ([`Self::append_received_send_messages_to_journal`]), so a primary
+ /// minting from a floor its peers do not share would have every one of
them
+ /// refuse the batch; relaxing that check deserves its own evidence. A
+ /// replicated group is also less exposed -- an ack means a quorum
journaled
+ /// the batch, so the hole there is a FULL-cluster crash.
+ fn mint_floor(&self) -> u64 {
+ let floor = self.armed_mint_floor();
+ self.mint_floor_pending.set(false);
+ floor
+ }
+
+ /// [`Self::mint_floor`] without spending it, for the boot re-anchor, which
+ /// has to know where the first mint will land before there is one.
+ const fn armed_mint_floor(&self) -> u64 {
+ if self.consensus.replica_count() > 1 ||
!self.mint_floor_pending.get() {
+ return 0;
+ }
+ self.durable_offset_reserved.get()
+ }
+
+ /// The append fence: make sure the durable record already permits every
+ /// offset up to and including `end_offset` before the caller lets them
+ /// exist.
+ ///
+ /// `SendMessagesResponse` hands clients concrete base offsets and the poll
+ /// path serves committed messages out of the resident journal, so an
offset
+ /// is client-visible long before the threshold-gated flush names it in a
+ /// segment, and a crash in between hands a second message an offset a
client
+ /// already holds. Fencing here rather than at commit puts it upstream of
+ /// every way an offset escapes -- the reply, the poll tier, the peers a
Review Comment:
`append_repaired_send_messages` journals without this fence, so 'every way
an offset escapes' overstates it; qualify the claim.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -4189,6 +4431,121 @@ where
Ok(())
}
+ /// Re-anchor the append point after boot re-seeded the offset counter
above
+ /// what the recovered segment chain holds.
+ ///
+ /// A hole INSIDE a segment is not survivable: `recover_segment_bounds`
walks
+ /// a segment from its FILENAME with a running `expected_offset`, so the
next
+ /// boot truncates the post-hole suffix and the counter falls back below
what
+ /// was already confirmed, undoing the fix on the second crash. On a
segment
+ /// BOUNDARY every reader copes -- absolute offsets in the index,
+ /// `disk_poll_start` walking on into later segments, a contiguity guard
that
+ /// only looks within one file.
+ ///
+ /// So an empty tail is unlinked (its name claims a range it does not hold)
+ /// and a sized tail, the only copy of its messages, is sealed with a fresh
+ /// segment planted at the frontier. An empty chain is left to the caller's
+ /// `ensure_initial_segment`.
+ ///
+ /// # Errors
+ /// [`IggyError`] when the fresh segment cannot be created, leaving the
+ /// partition without a serviceable chain.
+ #[allow(clippy::future_not_send)]
+ pub async fn reanchor_to_offset_frontier(
+ &mut self,
+ config: &PartitionsConfig,
+ ) -> Result<(), IggyError> {
+ // Where the next append will land: the counter, or an armed mint floor
+ // above it. The floor is the whole reason a hole can appear, so
+ // anchoring to the counter alone would leave the chain as unprepared.
+ let frontier = self.offset_frontier().max(self.armed_mint_floor());
+ if frontier == 0 {
+ return Ok(());
+ }
+ let namespace = self.namespace();
+ let mut retired = 0usize;
+ while let Some(segment) = self.log.segments().last() {
+ if segment.size.as_bytes_u64() > 0 || segment.start_offset >=
frontier {
+ break;
+ }
+ let Some((segment, mut storage)) = self.log.retire_back() else {
+ break;
+ };
+ let (messages_path, index_path) =
storage.segment_and_index_paths();
+ let _ = storage.shutdown();
+ drop(storage);
+ for path in messages_path.into_iter().chain(index_path) {
+ match compio::fs::remove_file(&path).await {
+ Ok(()) => {}
+ Err(error) if error.kind() == std::io::ErrorKind::NotFound
=> {}
+ Err(error) => {
+ warn!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ path = %path,
+ %error,
+ "failed to unlink a stale empty segment during the
boot re-anchor"
+ );
+ }
+ }
+ }
+ tracing::info!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ start_offset = segment.start_offset,
+ offset_frontier = frontier,
+ "unlinked an empty segment named below the restored offset
frontier"
+ );
+ retired += 1;
+ }
+ // Durable before anything is planted beside them: a crash in between
+ // would boot the stale name back into the chain. No stats decrement to
+ // pair with the retire -- boot never counted the recovered chain, and
+ // `fetch_sub` does not saturate.
+ if retired > 0
+ && let Some(partition_dir) = self.partition_dir.clone()
+ && let Err(error) =
crate::state_transfer::fsync_dir(&partition_dir).await
+ {
+ warn!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ partition_dir,
+ %error,
+ "boot re-anchor could not fsync the partition dir after
unlinking"
+ );
+ }
+ // Only a SIZED tail: an empty one either just went, or is already
named
+ // at the frontier and can take the appends as it is.
+ let needs_plant = self.log.segments().last().is_some_and(|segment| {
+ segment.size.as_bytes_u64() > 0 && segment.end_offset + 1 <
frontier
+ });
+ if !needs_plant {
+ return Ok(());
+ }
+ let sealed_index = self.log.segments().len() - 1;
+ let sealed_end = self.log.active_segment().end_offset;
+ self.log.active_segment_mut().sealed = true;
+ let sealed_storage = &mut self.log.storages_mut()[sealed_index];
Review Comment:
this seal-and-plant is `rotate_segment` minus `reset_read_state`. a
`rotate_segment_at(config, start)` shared by both keeps one seal path and its
create-then-teardown order.
##########
core/server/src/segment_recovery.rs:
##########
@@ -3018,6 +3051,52 @@ mod tests {
assert_eq!(bytes_of(&next_index_path), next_index);
}
+ #[compio::test]
Review Comment:
this pins the over-broad rule (any hole below the reservation); once the
guard only admits the last pair's forward gap, this test needs the tighter
shape, plus an overlap case and a mid-chain gap case.
##########
core/server/src/partition_helpers.rs:
##########
@@ -653,6 +653,7 @@ pub async fn build_partition_fresh(
config.partition.evicted_ring_capacity,
config.partition.evicted_ring_bytes_max.as_bytes_u64(),
);
+
partition.set_offset_reservation_lease(config.partition.offset_reservation_lease);
Review Comment:
same hole on this path: `set_superblock` arms the floor from the surviving
`offset_reserved`, `ensure_initial_segment` plants at the frontier, and the
re-anchor never runs here. the first mint lands a lease block into a segment
named below it, no crash needed.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -873,6 +965,34 @@ where
}
}
+ /// One past the highest offset this replica actually HOLDS: on disk in a
+ /// sized segment, or resident in the journal. `0` when it holds nothing.
+ ///
+ /// Not [`Self::offset_frontier`], which reads the append counter: after a
+ /// reservation-seeded boot that stands a lease block above the last byte
+ /// anywhere. Anything asking "would this destroy something I have?" has to
+ /// ask about messages. The journal arm is not redundant either, since the
+ /// threshold-gated flush routinely leaves committed messages unnamed by
any
+ /// segment.
+ #[must_use]
+ pub fn held_offset_frontier(&self) -> u64 {
+ let on_disk = self
+ .log
+ .segments()
+ .iter()
+ .filter(|segment| segment.size.as_bytes_u64() > 0)
+ .map(|segment| segment.end_offset.saturating_add(1))
+ .max()
+ .unwrap_or(0);
+ let journal = self.log.journal().info;
+ let resident = if journal.messages_count > 0 {
+ journal.current_offset.saturating_add(1)
Review Comment:
`current_offset` is the dirty tail, so the fence and view-change persists
record uncommitted offsets as the durable frontier. after a view change
truncates that tail, a restart seeds the counter above the group and this
replica refuses every prepare until state transfer. read the committed offset.
##########
core/server/src/segment_recovery.rs:
##########
@@ -567,7 +575,18 @@ fn ensure_contiguous_chain(
}
// `checked_add`, not `+`: an end offset at u64::MAX must read as a
// hole (no start offset can follow it), not overflow.
- if previous.end_offset.checked_add(1) != Some(next.start_offset) {
+ //
+ // A gap whose far side sits inside `offset_reserved` is the boot
+ // re-anchor's, not damage: `reanchor_to_offset_frontier` seals the
+ // recovered tail and plants the next segment at a frontier the
+ // superblock already claimed, so the skipped offsets were minted (or
+ // may have been) and no file was ever meant to hold them. Refusing it
+ // would tombstone the partition on the boot AFTER the one that fixed
+ // the re-mint. The ceiling is what bounds the leniency: a stray file
+ // splicing a gap ABOVE anything this replica claimed still refuses.
+ if previous.end_offset.checked_add(1) != Some(next.start_offset)
+ && next.start_offset > offset_reserved
Review Comment:
this also admits overlaps and missing middle segments, since every real
start sits below the reservation - after a clean stop, the whole chain. only
accept a forward gap on the last pair, and pass 0 for replicated groups.
##########
core/partitions/src/state_transfer.rs:
##########
@@ -1994,7 +1994,13 @@ where
let purge_advances = offsets_wire.purge_generation >
committed_purge_generation
|| (self.applied_purge_generation < committed_purge_generation
&& offsets_wire.next_offset == 0);
- let local_next_offset = self.offset_frontier();
+ // What this replica HOLDS, not the append counter: the two diverge
once
+ // the counter is seeded from an offset reservation, which stands a
lease
Review Comment:
same claim at iggy_partition.rs:734. the counter is never seeded from the
reservation - boot restores the frontier only and the floor bites at the first
mint - so this reasoning does not hold.
##########
core/consensus/src/vsr_state.rs:
##########
@@ -121,11 +136,15 @@ impl TryFrom<&[u8]> for VsrState {
type Error = VsrStateError;
fn try_from(bytes: &[u8]) -> Result<Self, Self::Error> {
- // Length-tolerant: a pre-`offset_frontier` record is padded out and
the
- // new field reads as 0, which is exactly "no recorded frontier" (the
- // read sites filter it). One length check up front then puts every
- // field slice below in bounds by construction, so the `try_into`s
- // cannot fail.
+ // Length-tolerant for ONE legacy layout: a pre-`offset_frontier`
record
+ // pads out and its trailing fields read as 0, which the read sites
+ // filter. The one length check then puts every field slice below in
+ // bounds by construction, so the `try_into`s cannot fail.
+ //
+ // `offset_reserved` gets no tolerance of its own: a silent 0 reads as
+ // "nothing reserved" and re-mints confirmed offsets, the exact defect
+ // the field closes, so a frontier-but-no-reservation record is
refused.
+ // Pre-production, so an older data directory is wiped, not migrated.
Review Comment:
server-0.9.0-edge.2 through edge.6 wrote 66-byte records on both planes, so
this refuses boot node-wide after an upgrade, with an error naming a 58-byte
layout no build ever wrote. accept 66 with zero-fill and drop the dead 58 arm.
##########
core/consensus/src/vsr_state.rs:
##########
@@ -31,8 +31,8 @@ use std::fmt;
/// Number of bytes [`VsrState::to_bytes`] produces: `cluster`(16) +
/// `replica_id`(1) + `replica_count`(1) + `view`(4) + `log_view`(4) +
/// `commit_max`(8) + `checkpoint_op`(8) + `checkpoint_checksum`(16) +
-/// `offset_frontier`(8).
-pub const ENCODED_LEN: usize = 66;
+/// `offset_frontier`(8) + `offset_reserved`(8).
+pub const ENCODED_LEN: usize = 74;
Review Comment:
a 74-byte record is `WrongLength` on the previous builds too, so rollback
needs a wipe as well - worth saying, or bump the record version.
##########
core/server/config.toml:
##########
@@ -973,6 +973,14 @@ clients_table_max = 8192
# u128 bitset, and this depth bounds that suffix.
prepare_queue_depth = 32
+# How many offsets a partition claims in its superblock ahead of the mint
+# counter before it will append, so a crash-restarted replica resumes above
+# every offset it confirmed to a client instead of re-minting it for a
different
+# message. One superblock write (two fsyncs) per block: lowering it raises the
+# fsync rate on the write path, raising it wastes at most one block of the u64
Review Comment:
same in core/configs/src/server_config/partition.rs:133. this protection
only applies to single-replica groups; say so here.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -998,6 +1139,80 @@ where
self.write_superblock(superblock.as_ref(), frontier).await
}
+ /// The lowest offset this replica may mint, given what its durable record
+ /// already permitted it to hand out. A floor, not a seed -- see
+ /// [`Self::restore_offset_frontier`]. It bites once, on the first mint
after
+ /// a crash that took acked-but-unflushed messages with it: segments and
+ /// counter come back short, and the reservation is the only surviving
+ /// witness that those offsets were confirmed.
+ ///
+ /// SOLO GROUPS ONLY. A backup validates an incoming prepare with
+ /// `base_offset == dirty_offset + 1`
+ /// ([`Self::append_received_send_messages_to_journal`]), so a primary
+ /// minting from a floor its peers do not share would have every one of
them
+ /// refuse the batch; relaxing that check deserves its own evidence. A
+ /// replicated group is also less exposed -- an ack means a quorum
journaled
+ /// the batch, so the hole there is a FULL-cluster crash.
+ fn mint_floor(&self) -> u64 {
+ let floor = self.armed_mint_floor();
+ self.mint_floor_pending.set(false);
+ floor
+ }
+
+ /// [`Self::mint_floor`] without spending it, for the boot re-anchor, which
+ /// has to know where the first mint will land before there is one.
+ const fn armed_mint_floor(&self) -> u64 {
+ if self.consensus.replica_count() > 1 ||
!self.mint_floor_pending.get() {
+ return 0;
+ }
+ self.durable_offset_reserved.get()
+ }
+
+ /// The append fence: make sure the durable record already permits every
+ /// offset up to and including `end_offset` before the caller lets them
+ /// exist.
+ ///
+ /// `SendMessagesResponse` hands clients concrete base offsets and the poll
+ /// path serves committed messages out of the resident journal, so an
offset
+ /// is client-visible long before the threshold-gated flush names it in a
+ /// segment, and a crash in between hands a second message an offset a
client
+ /// already holds. Fencing here rather than at commit puts it upstream of
+ /// every way an offset escapes -- the reply, the poll tier, the peers a
+ /// prepare reaches -- on the one path both a primary's mint and a backup's
+ /// re-stamp take. Claiming through `end_offset + 1 + lease` rather than
from
+ /// the live counter needs no special case for an oversized batch.
+ ///
+ /// `false` only when a write was attempted and failed, and the caller must
+ /// refuse the append. Fail-closed: the send is rejected with nothing
Review Comment:
the backoff arm returns false without attempting a write, and on the primary
the append error is swallowed in `on_replicate`, not rejected. reword once the
fence moves.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -998,6 +1139,80 @@ where
self.write_superblock(superblock.as_ref(), frontier).await
}
+ /// The lowest offset this replica may mint, given what its durable record
+ /// already permitted it to hand out. A floor, not a seed -- see
+ /// [`Self::restore_offset_frontier`]. It bites once, on the first mint
after
+ /// a crash that took acked-but-unflushed messages with it: segments and
+ /// counter come back short, and the reservation is the only surviving
+ /// witness that those offsets were confirmed.
+ ///
+ /// SOLO GROUPS ONLY. A backup validates an incoming prepare with
+ /// `base_offset == dirty_offset + 1`
+ /// ([`Self::append_received_send_messages_to_journal`]), so a primary
+ /// minting from a floor its peers do not share would have every one of
them
+ /// refuse the batch; relaxing that check deserves its own evidence. A
+ /// replicated group is also less exposed -- an ack means a quorum
journaled
+ /// the batch, so the hole there is a FULL-cluster crash.
+ fn mint_floor(&self) -> u64 {
+ let floor = self.armed_mint_floor();
+ self.mint_floor_pending.set(false);
+ floor
+ }
+
+ /// [`Self::mint_floor`] without spending it, for the boot re-anchor, which
+ /// has to know where the first mint will land before there is one.
+ const fn armed_mint_floor(&self) -> u64 {
+ if self.consensus.replica_count() > 1 ||
!self.mint_floor_pending.get() {
+ return 0;
+ }
+ self.durable_offset_reserved.get()
+ }
+
+ /// The append fence: make sure the durable record already permits every
+ /// offset up to and including `end_offset` before the caller lets them
+ /// exist.
+ ///
+ /// `SendMessagesResponse` hands clients concrete base offsets and the poll
+ /// path serves committed messages out of the resident journal, so an
offset
+ /// is client-visible long before the threshold-gated flush names it in a
+ /// segment, and a crash in between hands a second message an offset a
client
+ /// already holds. Fencing here rather than at commit puts it upstream of
+ /// every way an offset escapes -- the reply, the poll tier, the peers a
+ /// prepare reaches -- on the one path both a primary's mint and a backup's
+ /// re-stamp take. Claiming through `end_offset + 1 + lease` rather than
from
+ /// the live counter needs no special case for an oversized batch.
+ ///
+ /// `false` only when a write was attempted and failed, and the caller must
+ /// refuse the append. Fail-closed: the send is rejected with nothing
+ /// externalised, or a backup withholds its `PrepareOk` and the group
elects
+ /// around it, exactly as a failed view persist does.
+ #[allow(clippy::future_not_send)]
+ #[must_use = "the bool is the fence verdict; dropping it lets the append
escape unreserved"]
+ pub async fn reserve_offsets_through(&self, end_offset: u64) -> bool {
Review Comment:
nothing on a replicated group benefits from the reservation (the floor and
the re-anchor are solo-only), yet every replica pays this write per block and
the chain guard reads it to loosen itself. return early when `replica_count() >
1`.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -4189,6 +4431,121 @@ where
Ok(())
}
+ /// Re-anchor the append point after boot re-seeded the offset counter
above
+ /// what the recovered segment chain holds.
+ ///
+ /// A hole INSIDE a segment is not survivable: `recover_segment_bounds`
walks
+ /// a segment from its FILENAME with a running `expected_offset`, so the
next
+ /// boot truncates the post-hole suffix and the counter falls back below
what
+ /// was already confirmed, undoing the fix on the second crash. On a
segment
+ /// BOUNDARY every reader copes -- absolute offsets in the index,
+ /// `disk_poll_start` walking on into later segments, a contiguity guard
that
+ /// only looks within one file.
+ ///
+ /// So an empty tail is unlinked (its name claims a range it does not hold)
+ /// and a sized tail, the only copy of its messages, is sealed with a fresh
+ /// segment planted at the frontier. An empty chain is left to the caller's
+ /// `ensure_initial_segment`.
+ ///
+ /// # Errors
+ /// [`IggyError`] when the fresh segment cannot be created, leaving the
+ /// partition without a serviceable chain.
+ #[allow(clippy::future_not_send)]
+ pub async fn reanchor_to_offset_frontier(
+ &mut self,
+ config: &PartitionsConfig,
+ ) -> Result<(), IggyError> {
+ // Where the next append will land: the counter, or an armed mint floor
+ // above it. The floor is the whole reason a hole can appear, so
+ // anchoring to the counter alone would leave the chain as unprepared.
+ let frontier = self.offset_frontier().max(self.armed_mint_floor());
+ if frontier == 0 {
+ return Ok(());
+ }
+ let namespace = self.namespace();
+ let mut retired = 0usize;
+ while let Some(segment) = self.log.segments().last() {
+ if segment.size.as_bytes_u64() > 0 || segment.start_offset >=
frontier {
+ break;
+ }
+ let Some((segment, mut storage)) = self.log.retire_back() else {
+ break;
+ };
+ let (messages_path, index_path) =
storage.segment_and_index_paths();
+ let _ = storage.shutdown();
+ drop(storage);
+ for path in messages_path.into_iter().chain(index_path) {
+ match compio::fs::remove_file(&path).await {
+ Ok(()) => {}
+ Err(error) if error.kind() == std::io::ErrorKind::NotFound
=> {}
+ Err(error) => {
+ warn!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ path = %path,
+ %error,
Review Comment:
planting after a failed unlink leaves `[tail][stale empty][R]` on disk,
which the next boot refuses as an empty non-tail segment. return the error
instead (it aborts the node's boot, not a retry).
##########
core/shard/src/lib.rs:
##########
@@ -6762,6 +6762,20 @@ where
// after this flush). A partition already fenced by the commit
// path keeps its original fault.
partition.fence_flush_failure();
+ // The collapse below claims the segments account for every
+ // confirmed offset, which a failed flush is exactly the case
+ // against, so leave the reservation standing.
+ continue;
+ }
+ // The segments now prove where the offset space ends, so the
+ // reservation has nothing left to witness. Without the collapse
+ // every clean stop would leave a lease-block-wide hole.
+ if !partition.collapse_offset_reservation().await {
Review Comment:
one superblock write per active partition, serially, on every graceful stop.
fan these out with `join_all` - the lock and the store are per partition.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -4189,6 +4431,121 @@ where
Ok(())
}
+ /// Re-anchor the append point after boot re-seeded the offset counter
above
+ /// what the recovered segment chain holds.
+ ///
+ /// A hole INSIDE a segment is not survivable: `recover_segment_bounds`
walks
+ /// a segment from its FILENAME with a running `expected_offset`, so the
next
+ /// boot truncates the post-hole suffix and the counter falls back below
what
+ /// was already confirmed, undoing the fix on the second crash. On a
segment
+ /// BOUNDARY every reader copes -- absolute offsets in the index,
+ /// `disk_poll_start` walking on into later segments, a contiguity guard
that
+ /// only looks within one file.
+ ///
+ /// So an empty tail is unlinked (its name claims a range it does not hold)
+ /// and a sized tail, the only copy of its messages, is sealed with a fresh
+ /// segment planted at the frontier. An empty chain is left to the caller's
+ /// `ensure_initial_segment`.
+ ///
+ /// # Errors
+ /// [`IggyError`] when the fresh segment cannot be created, leaving the
+ /// partition without a serviceable chain.
+ #[allow(clippy::future_not_send)]
+ pub async fn reanchor_to_offset_frontier(
+ &mut self,
+ config: &PartitionsConfig,
+ ) -> Result<(), IggyError> {
+ // Where the next append will land: the counter, or an armed mint floor
+ // above it. The floor is the whole reason a hole can appear, so
+ // anchoring to the counter alone would leave the chain as unprepared.
+ let frontier = self.offset_frontier().max(self.armed_mint_floor());
+ if frontier == 0 {
+ return Ok(());
+ }
+ let namespace = self.namespace();
+ let mut retired = 0usize;
+ while let Some(segment) = self.log.segments().last() {
+ if segment.size.as_bytes_u64() > 0 || segment.start_offset >=
frontier {
+ break;
+ }
+ let Some((segment, mut storage)) = self.log.retire_back() else {
+ break;
+ };
+ let (messages_path, index_path) =
storage.segment_and_index_paths();
+ let _ = storage.shutdown();
+ drop(storage);
+ for path in messages_path.into_iter().chain(index_path) {
+ match compio::fs::remove_file(&path).await {
+ Ok(()) => {}
+ Err(error) if error.kind() == std::io::ErrorKind::NotFound
=> {}
+ Err(error) => {
+ warn!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ path = %path,
+ %error,
+ "failed to unlink a stale empty segment during the
boot re-anchor"
+ );
+ }
+ }
+ }
+ tracing::info!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ start_offset = segment.start_offset,
+ offset_frontier = frontier,
+ "unlinked an empty segment named below the restored offset
frontier"
+ );
+ retired += 1;
+ }
+ // Durable before anything is planted beside them: a crash in between
+ // would boot the stale name back into the chain. No stats decrement to
Review Comment:
boot does count the recovered chain - `load_persisted_segments` increments
per segment, empty tails included - so each retire here leaves `segments_count`
one too high on the wire. decrement per retire, like retention does.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -998,6 +1139,80 @@ where
self.write_superblock(superblock.as_ref(), frontier).await
}
+ /// The lowest offset this replica may mint, given what its durable record
+ /// already permitted it to hand out. A floor, not a seed -- see
+ /// [`Self::restore_offset_frontier`]. It bites once, on the first mint
after
+ /// a crash that took acked-but-unflushed messages with it: segments and
+ /// counter come back short, and the reservation is the only surviving
+ /// witness that those offsets were confirmed.
+ ///
+ /// SOLO GROUPS ONLY. A backup validates an incoming prepare with
+ /// `base_offset == dirty_offset + 1`
+ /// ([`Self::append_received_send_messages_to_journal`]), so a primary
+ /// minting from a floor its peers do not share would have every one of
them
+ /// refuse the batch; relaxing that check deserves its own evidence. A
+ /// replicated group is also less exposed -- an ack means a quorum
journaled
+ /// the batch, so the hole there is a FULL-cluster crash.
+ fn mint_floor(&self) -> u64 {
+ let floor = self.armed_mint_floor();
+ self.mint_floor_pending.set(false);
+ floor
+ }
+
+ /// [`Self::mint_floor`] without spending it, for the boot re-anchor, which
+ /// has to know where the first mint will land before there is one.
+ const fn armed_mint_floor(&self) -> u64 {
+ if self.consensus.replica_count() > 1 ||
!self.mint_floor_pending.get() {
+ return 0;
+ }
+ self.durable_offset_reserved.get()
+ }
+
+ /// The append fence: make sure the durable record already permits every
+ /// offset up to and including `end_offset` before the caller lets them
+ /// exist.
+ ///
+ /// `SendMessagesResponse` hands clients concrete base offsets and the poll
+ /// path serves committed messages out of the resident journal, so an
offset
+ /// is client-visible long before the threshold-gated flush names it in a
+ /// segment, and a crash in between hands a second message an offset a
client
+ /// already holds. Fencing here rather than at commit puts it upstream of
+ /// every way an offset escapes -- the reply, the poll tier, the peers a
+ /// prepare reaches -- on the one path both a primary's mint and a backup's
+ /// re-stamp take. Claiming through `end_offset + 1 + lease` rather than
from
+ /// the live counter needs no special case for an oversized batch.
+ ///
+ /// `false` only when a write was attempted and failed, and the caller must
+ /// refuse the append. Fail-closed: the send is rejected with nothing
+ /// externalised, or a backup withholds its `PrepareOk` and the group
elects
+ /// around it, exactly as a failed view persist does.
+ #[allow(clippy::future_not_send)]
+ #[must_use = "the bool is the fence verdict; dropping it lets the append
escape unreserved"]
+ pub async fn reserve_offsets_through(&self, end_offset: u64) -> bool {
+ // A frontier: offsets strictly below it are permitted, so covering
+ // `end_offset` needs a record strictly above it.
+ if self.durable_offset_reserved.get() > end_offset {
+ return true;
+ }
+ 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;
+ // A batch queued behind another append's write finds the block already
+ // extended.
+ if self.durable_offset_reserved.get() > end_offset {
+ return true;
+ }
+ let claim = end_offset
Review Comment:
two fsyncs awaited inline under `write_lock` inside the shard pump, so every
partition on the core stalls once per block - the first fsync on the default
append path. claim the next block off the append path, about half a block early.
##########
core/integration/tests/cluster/crash_offset_reuse.rs:
##########
@@ -141,3 +197,91 @@ async fn
given_confirmed_sends_below_flush_threshold_when_a_solo_node_is_killed_
silently different data"
);
}
+
+/// The fix has to survive its own side effect: the hole the reservation leaves
+/// between the recovered segments and the new append point truncates
everything
+/// past it on the next boot if it lands INSIDE a segment.
+///
+/// So the SECOND crash is the one that matters, and only if the run between
the
+/// two reaches disk, which is what the flush threshold is for.
+#[iggy_harness(cluster_nodes = 1)]
Review Comment:
none of these stops the node cleanly between crashes, and all three are
single-node - the shutdown collapse and the replicated fence path stay
untested. add a crash, restart, clean stop, restart, send case and a 3-node one.
##########
core/configs/src/server_config/partition.rs:
##########
@@ -177,6 +188,16 @@ impl Validatable<ConfigurationError> for PartitionConfig {
);
return Err(ConfigurationError::InvalidConfigurationValue);
}
+ if self.offset_reservation_lease == 0
Review Comment:
no zero or above-ceiling test for this knob; the `prepare_queue_depth` and
`evicted_ring` ones have both.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -4189,6 +4431,121 @@ where
Ok(())
}
+ /// Re-anchor the append point after boot re-seeded the offset counter
above
+ /// what the recovered segment chain holds.
+ ///
+ /// A hole INSIDE a segment is not survivable: `recover_segment_bounds`
walks
+ /// a segment from its FILENAME with a running `expected_offset`, so the
next
+ /// boot truncates the post-hole suffix and the counter falls back below
what
+ /// was already confirmed, undoing the fix on the second crash. On a
segment
Review Comment:
also at bootstrap.rs:2808 and crash_offset_reuse.rs:202, 252. neither
recovery walk truncates a hole any more - both refuse with
`OffsetDiscontinuity` and the solo arm tombstones, which is worse than this
says.
##########
core/partitions/src/log.rs:
##########
@@ -314,6 +314,28 @@ where
Some((segment, storage))
}
+ /// Retire the NEWEST segment, the mirror of [`Self::retire_front`].
+ ///
+ /// Boot re-anchor only, and only for an EMPTY tail named below the
re-seeded
+ /// counter, which would otherwise claim a range it does not hold. Nothing
+ /// else may take from the back: a sized tail is the only copy of its
+ /// messages.
+ pub fn retire_back(&mut self) -> Option<(Segment, SegmentStorage)> {
Review Comment:
the `segments_mut` doc at 194 lists `add_persisted_segment` and
`retire_front` as the only length mutators; add this one.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -4189,6 +4431,121 @@ where
Ok(())
}
+ /// Re-anchor the append point after boot re-seeded the offset counter
above
+ /// what the recovered segment chain holds.
+ ///
+ /// A hole INSIDE a segment is not survivable: `recover_segment_bounds`
walks
+ /// a segment from its FILENAME with a running `expected_offset`, so the
next
+ /// boot truncates the post-hole suffix and the counter falls back below
what
+ /// was already confirmed, undoing the fix on the second crash. On a
segment
+ /// BOUNDARY every reader copes -- absolute offsets in the index,
+ /// `disk_poll_start` walking on into later segments, a contiguity guard
that
+ /// only looks within one file.
+ ///
+ /// So an empty tail is unlinked (its name claims a range it does not hold)
+ /// and a sized tail, the only copy of its messages, is sealed with a fresh
+ /// segment planted at the frontier. An empty chain is left to the caller's
+ /// `ensure_initial_segment`.
+ ///
+ /// # Errors
+ /// [`IggyError`] when the fresh segment cannot be created, leaving the
+ /// partition without a serviceable chain.
+ #[allow(clippy::future_not_send)]
+ pub async fn reanchor_to_offset_frontier(
+ &mut self,
+ config: &PartitionsConfig,
+ ) -> Result<(), IggyError> {
+ // Where the next append will land: the counter, or an armed mint floor
+ // above it. The floor is the whole reason a hole can appear, so
+ // anchoring to the counter alone would leave the chain as unprepared.
+ let frontier = self.offset_frontier().max(self.armed_mint_floor());
+ if frontier == 0 {
+ return Ok(());
+ }
+ let namespace = self.namespace();
+ let mut retired = 0usize;
+ while let Some(segment) = self.log.segments().last() {
+ if segment.size.as_bytes_u64() > 0 || segment.start_offset >=
frontier {
+ break;
+ }
+ let Some((segment, mut storage)) = self.log.retire_back() else {
+ break;
+ };
+ let (messages_path, index_path) =
storage.segment_and_index_paths();
+ let _ = storage.shutdown();
+ drop(storage);
+ for path in messages_path.into_iter().chain(index_path) {
+ match compio::fs::remove_file(&path).await {
+ Ok(()) => {}
+ Err(error) if error.kind() == std::io::ErrorKind::NotFound
=> {}
+ Err(error) => {
+ warn!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ path = %path,
+ %error,
+ "failed to unlink a stale empty segment during the
boot re-anchor"
+ );
+ }
+ }
+ }
+ tracing::info!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ start_offset = segment.start_offset,
+ offset_frontier = frontier,
+ "unlinked an empty segment named below the restored offset
frontier"
+ );
+ retired += 1;
+ }
+ // Durable before anything is planted beside them: a crash in between
+ // would boot the stale name back into the chain. No stats decrement to
+ // pair with the retire -- boot never counted the recovered chain, and
+ // `fetch_sub` does not saturate.
+ if retired > 0
+ && let Some(partition_dir) = self.partition_dir.clone()
+ && let Err(error) =
crate::state_transfer::fsync_dir(&partition_dir).await
+ {
+ warn!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ partition_dir,
+ %error,
+ "boot re-anchor could not fsync the partition dir after
unlinking"
+ );
+ }
+ // Only a SIZED tail: an empty one either just went, or is already
named
+ // at the frontier and can take the appends as it is.
+ let needs_plant = self.log.segments().last().is_some_and(|segment| {
+ segment.size.as_bytes_u64() > 0 && segment.end_offset + 1 <
frontier
+ });
+ if !needs_plant {
+ return Ok(());
+ }
+ let sealed_index = self.log.segments().len() - 1;
+ let sealed_end = self.log.active_segment().end_offset;
+ self.log.active_segment_mut().sealed = true;
+ let sealed_storage = &mut self.log.storages_mut()[sealed_index];
+ let _ = sealed_storage.shutdown();
+ self.log.messages_writers_mut()[sealed_index] = None;
+ self.log.index_writers_mut()[sealed_index] = None;
+ self.log.indexes_mut()[sealed_index] = None;
+ self.install_empty_segment(config, frontier).await?;
Review Comment:
`purge` fsyncs the dir after its plant and says why; this one does not. a
line on why the floor makes it unnecessary would save the next reader the
question.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -752,14 +837,21 @@ where
/// proved.
///
/// The record is a lower bound, never a completeness claim: it exists
- /// because three paths leave a replica whose counter would otherwise
- /// restart at 0 while the group is at N (a transfer install of an all-GC'd
- /// origin, a crash inside the install's swap window, and the
- /// fence-and-rebuild path, which needs no crash at all). Restarting the
+ /// because four paths leave a replica whose counter would otherwise
+ /// restart below where the group already is (a transfer install of an
+ /// all-GC'd origin, a crash inside the install's swap window, the
+ /// fence-and-rebuild path, which needs no crash at all, and a crash while
+ /// acked messages were still resident in the journal). Restarting the
/// counter is not a lag -- replicas re-stamp `base_offset` from it and
/// recompute `batch_checksum` over the result, so the next replicated
/// prepare would persist different bytes here than on every peer,
silently.
///
+ /// The RESERVATION is deliberately not folded in here: a backup mints
Review Comment:
seed the solo counter from `max(offset_frontier, offset_reserved)` here
(gated on `replica_count() == 1`): `offset_frontier()` is then right for the
re-anchor, `ensure_initial_segment` and the collapse at once, and the
mint-floor state machine goes away. keep `held_offset_frontier` bytes-based so
a later 1 -> N topology change does not inherit an inflated frontier.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -1103,6 +1318,16 @@ where
.unwrap_or(config.messages_required_to_save)
}
+ /// Install the offset-reservation block size resolved from this node's
+ /// `PartitionsConfig`.
+ ///
+ /// Carried on the partition because the fence runs inside `on_request` /
+ /// `on_replicate`, which take no config. Floored at 1: a zero block
reserves
+ /// nothing and would write the superblock before every append.
+ pub const fn set_offset_reservation_lease(&mut self, lease: u32) {
Review Comment:
`u64::from(lease)` and drop `const` - the cast only exists to dodge
`cast_lossless`. keep the zero floor, unit tests and the simulator set this
directly.
--
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]