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]

Reply via email to