numinnex commented on code in PR #3975:
URL: https://github.com/apache/iggy/pull/3975#discussion_r3881964667
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -910,6 +1032,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:
`offset_frontier()` is the append counter, so a boot mint floor that was
never consumed is erased here.
Reproduced on this branch: crash, restart, clean stop, restart, produce —
offset 4 was ACKed before the SIGKILL and the first send after the clean
restart is confirmed at offset 0 again. The graceful stop is what destroys the
protection; a second SIGKILL would have preserved it, so the runbook-correct
incident response is the one action that defeats the fix. Nothing logs the
moment.
Suggest:
```rust
let frontier = self
.offset_frontier()
.max(self.armed_mint_floor())
.max(self.durable_offset_frontier.get());
```
The `armed_mint_floor` term covers the case above. The
`durable_offset_frontier` term is a separate fix: `reset_offset_frontier_at` is
the non-monotone writer, so in the failed-install window this currently
regresses the recorded frontier.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -933,6 +1074,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);
Review Comment:
The flag is spent at mint time, but every failure path after the mint
discards the mint itself: the `count == 0` return, the `checked_add`,
`stamp_prepare_for_persistence`'s `?`, the reservation fence, and the journal
append all return before `dirty_offset` is stored.
The next append then mints `dirty_offset + 1` — the pre-crash value — while
`reserve_offsets_through` short-circuits on `durable_offset_reserved >
end_offset` and writes nothing. So one transient superblock error on the first
append after a crash-restart re-opens the exact defect this PR closes, silently.
Suggest disarming where `dirty_offset` is stored instead. That point is past
every early return, and it also covers the `?` in
`stamp_prepare_for_persistence`, which sits after the mint and is not covered
today.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -2666,6 +2895,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:
This two-fsync `atomic_replace` runs inline in the shard's sequential frame
pump, and the consensus tick is a `select_biased!` arm on that same pump — its
own comment bounds tick delay to one main frame body's longest `.await`. A slow
superblock write therefore withholds heartbeats for every group on the shard,
not only this partition, and peers can elect around the node. It also sits
ahead of `send_prepare_ok` on the backup path, so it is on the primary's
commit-latency path too.
The LOCK ORDER note above is correct about lock ordering but understates the
blast radius.
Worth saying that the throughput cost is negligible and the arithmetic in
the description checks out: roughly 3 fsyncs/s at 100k msg/s with the 64Ki
lease, about 15 ns/message amortized. The concern is only where the stall
lands. Extending the reservation off the critical path at a watermark, keeping
this fence as the backstop, would preserve the guarantee without putting the
write inside the pump.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -4033,6 +4275,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 {
Review Comment:
When the unlink loop empties the chain, `last()` is `None`, `needs_plant` is
false, and this returns without planting. `ensure_initial_segment` then plants
at `offset_frontier()` and knows nothing about `armed_mint_floor()`, so the
first mint lands at the reservation inside a segment named for the frontier —
the hole-inside-a-segment shape the re-anchor exists to prevent.
It is tolerated while the index survives, so it is invisible in the common
case. With a torn index (`enforce_fsync = false` is the shipped default) the
index-less walk hits `base_offset != expected_offset` and raises
`OffsetDiscontinuity`, tombstoning the partition.
This is the shape
`given_confirmed_sends_below_flush_threshold_when_a_solo_node_is_killed_should_not_remint_offsets`
produces; its assertion only checks the offset moved forward, so it passes
over the segment shape. `build_partition_fresh` has the same gap — it arms the
floor via `set_superblock` but never re-anchors before `ensure_initial_segment`.
Suggest planting here when the chain is empty and `frontier > 0`, or passing
the floor into `ensure_initial_segment`.
##########
core/server/src/segment_recovery.rs:
##########
@@ -507,7 +513,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:
`write_superblock_inner` clamps `offset_reserved` up to `offset_frontier`,
and `advanced_frontier` floors the frontier at `held_offset_frontier()`.
`offset_reserved` therefore sits at or above the end of every legitimate
segment, which makes `next.start_offset > offset_reserved` false for any real
gap — so `Hole` can no longer fire for a partition that has ever appended. A
genuinely lost middle segment now recovers silently as a holed log instead of
refusing loudly, and the admitted arm logs nothing.
Narrowing the predicate is harder than it first looks. The gaps accumulate,
one per crash. Probed on this branch, segment base offsets across successive
crash cycles:
```
[0]
[0, 65537]
[0, 65537, 131074]
[0, 65537, 131074, 196611]
```
So an "at most one gap" rule would tombstone a healthy partition by the
second crash, and
`given_a_crash_restarted_node_when_it_flushes_and_crashes_again_should_still_not_remint_offsets`
already reaches the two-gap state. Pinning `next.start_offset ==
offset_reserved` is unstable for a different reason: the first append after the
plant extends the reservation past it, so the same pair fails on the following
boot.
A single monotone scalar may not be able to separate re-anchor gaps from
damage across repeated cycles. Recording the anchor points durably, or making
the planted segment self-describing (carrying its predecessor's end offset),
would give the guard something stable to check against.
--
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]