hubcio commented on code in PR #3975:
URL: https://github.com/apache/iggy/pull/3975#discussion_r3905498066
##########
core/server/src/bootstrap.rs:
##########
@@ -2601,7 +2618,6 @@ async fn load_partition(
.max()
.filter(|&start| sized_end.is_none() && start > 0);
let current_offset = sized_end.or_else(|| empty_frontier.map(|start| start
- 1));
Review Comment:
an empty chain's name used to be a committed frontier; the re-anchor and
`ensure_initial_segment` now plant it at `mint_frontier()`, which is a
reservation. two crashes below the flush threshold then store committed offset
65537 for a partition holding nothing, and `store_consumer_offset` admits the
whole hole. bound the promotion by `recovered_state.offset_frontier`, which a
real install writes at the group frontier and a re-anchor plant leaves far
below.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -752,36 +829,78 @@ 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 two counters are seeded separately, because the record carries two
+ /// bounds. `offset_frontier` is what messages reached, so it seeds the
+ /// COMMITTED counter; `offset_reserved` is what may have been handed to a
+ /// client, so it seeds the APPEND counter and nothing else. Folding the
+ /// reservation into the committed one would publish a `current_offset`
over
+ /// a lease-block hole, and `store_consumer_offset` would then admit
offsets
+ /// inside it.
+ ///
+ /// The reservation is SOLO ONLY. A backup mints nothing: it re-stamps what
+ /// the primary sends and rejects anything that does not continue its own
+ /// counter, so an append point a lease block above its group would have
every
+ /// peer refuse the batch. A replicated group is also less exposed, since
an
+ /// ack there means a quorum journaled the batch and the hole needs a
+ /// FULL-cluster crash.
+ ///
/// Lives HERE rather than in the server crate so the boot paths and the
/// simulator share one implementation. A copy in the harness was a copy of
/// the max rule that had lost the max, in the one place built to catch
/// violations of it.
pub fn restore_offset_frontier(&mut self, recovered:
Option<&consensus::VsrState>) {
- let Some(frontier) = recovered
- .map(|state| state.offset_frontier)
- .filter(|&f| f > 0)
- else {
+ let Some(state) = recovered else {
return;
};
- let recovered_end = frontier - 1;
- if self.should_increment_offset && self.offset.load(Ordering::Acquire)
>= recovered_end {
+ let frontier = state.offset_frontier;
+ let reserved = if self.consensus.replica_count() == 1 {
+ state.offset_reserved
+ } else {
+ 0
+ };
+ // NOT `frontier > 0`: the shape a crash before the first flush leaves
is
+ // a zero frontier and a nonzero reservation, because the append fence
+ // runs before the journal append, so the first record a partition ever
+ // writes names no data at all.
+ let append_point = frontier.max(reserved);
+ if append_point == 0 {
+ return;
+ }
+ let seeded = self.should_increment_offset;
+ // Each counter takes its own max, since the record's two bounds move
+ // independently: a graceful stop collapses the reservation onto the
+ // append point while the frontier stays where the data ended.
+ let committed_restored =
frontier.checked_sub(1).is_some_and(|committed_end| {
+ let raise = !seeded || self.offset.load(Ordering::Acquire) <
committed_end;
+ if raise {
+ self.offset.store(committed_end, Ordering::Release);
+ }
+ raise
+ });
+ let append_end = append_point - 1;
+ let append_restored = !seeded ||
self.dirty_offset.load(Ordering::Relaxed) < append_end;
+ if append_restored {
+ self.dirty_offset.store(append_end, Ordering::Relaxed);
+ }
+ if !committed_restored && !append_restored {
return;
}
tracing::info!(
namespace_raw = self.consensus().group(),
offset_frontier = frontier,
- "restored partition offset frontier from its superblock"
+ offset_reserved = reserved,
+ append_point,
+ "restored partition offset counters from its superblock"
);
- self.offset.store(recovered_end, Ordering::Release);
- self.dirty_offset.store(recovered_end, Ordering::Relaxed);
self.should_increment_offset = true;
Review Comment:
`should_increment_offset` now means the append counter is live, but
`offset_frontier()` still reads it as the committed counter naming data, so it
returns 1 for a partition holding nothing. a separate bit for the committed
seed stops that; it does not help the boot path above, which rebuilds from the
file name.
##########
core/partitions/src/segment_anchor.rs:
##########
@@ -0,0 +1,205 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! The record that makes a gap in the segment chain legitimate.
+//!
+//! Recovery derives each segment's bounds from its own bytes, so a gap has no
+//! author: the boot re-anchor's planted gap and a lost segment look identical.
+//! The re-anchor writes its intent down instead, and the chain guard admits a
+//! forward gap only when the far side carries an anchor naming exactly the
near
+//! side.
+//!
+//! Written and directory-fsynced BEFORE that segment is created. A crash in
the
+//! window then leaves an anchor with no segment, which the boot sweep
collects;
+//! the other order leaves a planted segment with no anchor, which the guard
+//! reads as damage on an intact chain.
+
+use compio::io::AsyncWriteAtExt;
+use consensus::state_artifact_checksum;
+use std::io;
+
+/// File extension for an anchor record, `{start_offset:020}.anchor` beside the
+/// `{start_offset:020}.log` it belongs to.
+pub const ANCHOR_EXTENSION: &str = "anchor";
+
+/// Leading bytes of an anchor record, so a file that is not one (a truncated
+/// write, an operator's copy) is refused rather than decoded.
+const ANCHOR_MAGIC: u64 = u64::from_le_bytes(*b"IGGYANCH");
+
+/// `magic`(8) + `sealed_start`(8) + `sealed_end`(8) + `checksum`(8).
+pub const ANCHOR_ENCODED_LEN: usize = 32;
+
+/// The gap one planted segment is allowed to leave behind it.
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub struct SegmentAnchor {
+ /// Start offset of the segment that was sealed, i.e. the near side of the
+ /// gap. Names WHICH segment, so an anchor cannot be satisfied by a
+ /// different file that happens to end where this one expects.
+ pub sealed_start: u64,
+ /// End offset the sealed segment held when it was sealed. The gap runs
from
+ /// here to the planted segment's own start offset.
+ pub sealed_end: u64,
+}
+
+impl SegmentAnchor {
+ /// Encode to the fixed little-endian on-disk layout.
+ #[must_use]
+ pub fn to_bytes(&self) -> [u8; ANCHOR_ENCODED_LEN] {
+ let mut out = [0u8; ANCHOR_ENCODED_LEN];
+ out[0..8].copy_from_slice(&ANCHOR_MAGIC.to_le_bytes());
+ out[8..16].copy_from_slice(&self.sealed_start.to_le_bytes());
+ out[16..24].copy_from_slice(&self.sealed_end.to_le_bytes());
+ let checksum = state_artifact_checksum(&out[0..24]);
+ out[24..32].copy_from_slice(&checksum.to_le_bytes());
+ out
+ }
+
+ /// Decode a record, returning `None` for anything this build did not
write:
+ /// a wrong length, a wrong magic, or a checksum that does not match.
+ ///
+ /// A `None` is never treated as "no gap was intended". It means the record
+ /// proves nothing, so the gap it would have covered stays damage.
+ #[must_use]
+ pub fn from_bytes(bytes: &[u8]) -> Option<Self> {
+ if bytes.len() != ANCHOR_ENCODED_LEN {
+ return None;
+ }
+ let field = |at: usize| -> u64 {
+ let mut raw = [0u8; 8];
+ raw.copy_from_slice(&bytes[at..at + 8]);
+ u64::from_le_bytes(raw)
+ };
+ if field(0) != ANCHOR_MAGIC || field(24) !=
state_artifact_checksum(&bytes[0..24]) {
+ return None;
+ }
+ Some(Self {
+ sealed_start: field(8),
+ sealed_end: field(16),
+ })
+ }
+
+ /// Whether this anchor legitimises the gap between the segment starting at
+ /// `sealed_start` / ending at `sealed_end` and the segment it sits beside.
+ ///
+ /// BOTH bounds must match. Matching only the end offset would let an
anchor
+ /// left by an earlier incarnation of the chain cover a gap it never saw.
+ #[must_use]
+ pub const fn covers(&self, sealed_start: u64, sealed_end: u64) -> bool {
+ self.sealed_start == sealed_start && self.sealed_end == sealed_end
+ }
+}
+
+/// Path of the anchor record beside the segment starting at `start_offset`.
+#[must_use]
+pub fn anchor_path(partition_dir: &str, start_offset: u64) -> String {
+ format!("{partition_dir}/{start_offset:0>20}.{ANCHOR_EXTENSION}")
+}
+
+/// Write the anchor for a segment about to be planted at `start_offset`, then
+/// fsync the directory so the record cannot arrive after the segment it
+/// describes.
+///
+/// # Errors
+///
+/// Any I/O failure. The caller must NOT plant the segment: a planted segment
+/// whose anchor is missing reads as damage on the next boot.
+pub async fn write_anchor(
+ partition_dir: &str,
+ start_offset: u64,
+ anchor: SegmentAnchor,
+) -> io::Result<()> {
+ let path = anchor_path(partition_dir, start_offset);
+ let mut file = compio::fs::File::create(&path).await?;
+ let (result, _buf) = file
+ .write_all_at(anchor.to_bytes().to_vec(), 0)
+ .await
+ .into();
+ result?;
+ file.sync_all().await?;
+ crate::state_transfer::fsync_dir(partition_dir).await
+}
+
+/// Read the anchor beside the segment starting at `start_offset`.
+///
+/// `None` when the file is absent or does not decode; both mean the same thing
+/// to the guard, so the caller needs no distinction.
+pub async fn read_anchor(partition_dir: &str, start_offset: u64) ->
Option<SegmentAnchor> {
+ let path = anchor_path(partition_dir, start_offset);
+ let bytes = compio::fs::read(&path).await.ok()?;
Review Comment:
`.ok()?` turns an `EACCES` or `EIO` into no gap intended, so a healthy
planted chain is refused as `Hole` - a solo group goes dark for the life of the
process and every boot while the error lasts, a replicated one re-transfers the
whole partition over one bad read. `file_len` fail-stops on exactly this class;
do the same and treat `NotFound` separately.
##########
core/shard/src/lib.rs:
##########
@@ -6554,6 +6573,26 @@ where
}
continue;
}
+ // Same bound the metadata plane fail-stops on, applied per group.
A
+ // partition whose superblock keeps refusing withholds every
+ // view-scoped send AND refuses every append, so it serves nothing
+ // while the process still reports healthy.
+ let superblock_failures = partition.superblock_write_failures();
+ if superblock_wedged(
Review Comment:
one unwritable partition dir now exits the node and crashloops on restart,
while `iggy_partition.rs:692` still says only that group is fenced and the rest
keeps serving. `refuses every append` is also wrong above one replica -
`reserve_offsets_through` returns true there - and the timeout's operator docs
still scope it to the metadata superblock.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -975,6 +1145,46 @@ 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.
+ ///
+ /// Collapses onto [`Self::mint_frontier`], not [`Self::offset_frontier`]:
the
+ /// append point is what the next boot has to resume at, and on a boot that
+ /// consumed a reservation without appending it is the reservation itself,
so
+ /// reading the committed frontier here would write a record BELOW what an
+ /// earlier life already confirmed to a client. A clean stop is the runbook
+ /// answer to an incident, which would make it the one action that undoes
the
+ /// protection.
+ ///
+ /// The frontier field still records only what is held: a graceful stop
+ /// flushes the committed prefix, but the journal can hold an uncommitted
tail
+ /// that the next view legitimately truncates.
+ ///
+ /// 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 append_point = self.mint_frontier();
+ if self.durable_offset_reserved.get() <= append_point {
+ return true;
+ }
+ let Some(superblock) = self.superblock.as_ref().map(Rc::clone) else {
+ return true;
+ };
+ if self.superblock_write_is_backed_off() {
Review Comment:
a superblock error just before a clean stop leaves the reservation standing,
so the next boot seeds the append counter a lease block above the data and
holes the offset space for nothing - what the doc at :1152 warns about.
`record_frontier_before_quarantine` skips the backoff gate for exactly this
reason; do the same here.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -4189,41 +4613,245 @@ where
Ok(())
}
- /// Record the purge's frontier reset BEFORE the purge touches anything.
+ /// Re-anchor the append point after boot re-seeded the offset counter
above
+ /// what the recovered segment chain holds.
///
- /// The unlinks are made durable by their own directory fsync, so a crash
- /// between them and a reset written afterwards boots a purged directory
- /// whose record still names the pre-purge offset space:
- /// `restore_offset_frontier` re-seeds the counter to it while every peer
- /// restarted at 0, and the first append stamps a `base_offset` and
- /// `batch_checksum` no peer shares. Writing 0 first inverts the window
into
- /// a harmless one -- the record under-claims while the segments still
- /// exist, and boot takes the max of the record and what the segments
prove.
+ /// A hole INSIDE a segment is not survivable: `recover_segment_bounds`
walks
+ /// a segment from its FILENAME with a running `expected_offset` and
REFUSES
+ /// at the first offset that does not continue it (`OffsetDiscontinuity`),
+ /// which on a solo group tombstones the partition. A surviving index does
not
+ /// help: a first entry that is not the file-name offset makes recovery
+ /// discard the index and walk from byte 0, reaching the same refusal. On a
+ /// segment BOUNDARY every reader copes -- absolute offsets in the index,
+ /// `disk_poll_start` walking on into later segments, and a chain guard
that
+ /// admits a forward gap the reservation covers.
///
- /// Spelled out rather than read off the counter, which still holds the
- /// pre-purge frontier at this point.
+ /// So an empty tail is unlinked (its name claims a range it does not
hold), a
+ /// sized tail (the only copy of its messages) is sealed with a fresh
segment
+ /// planted at the append point, and a chain the unlinks emptied is planted
+ /// directly -- `ensure_initial_segment` names its segment for the
COMMITTED
+ /// frontier and would put the first mint inside it.
///
/// # Errors
- /// [`PurgeError::FrontierNotRecorded`]. Refused rather than logged:
nothing
- /// has been mutated yet, and a purge that cannot record its reset must not
- /// be the one that erases the data proving the old frontier. The caller
- /// RETRIES; it must not fence, since the chain is still whole and the live
- /// counter still names the pre-purge space.
+ /// [`IggyError`] when the fresh segment cannot be created, leaving the
+ /// partition without a serviceable chain.
#[allow(clippy::future_not_send)]
- async fn record_purge_frontier_reset(&mut self, generation: u64) ->
Result<(), PurgeError> {
- if self.reset_offset_frontier_at(0).await {
- self.purge_deferred = false;
+ 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.mint_frontier();
+ if frontier == 0 {
return Ok(());
}
- self.purge_deferred = true;
- // The ONLY operator-visible signal for the withhold: `send_prepare_ok`
- // returns silently, correctly, since it runs per prepare. So this line
- // has to say that the replica is now out of quorum for this group, or
- // the symptom reads as a network fault. The consecutive count
- // correlates it with the superblock writer's own error log, which
- // carries the `ENOSPC` / `EIO` cause but is rate-limited to
- // power-of-two failures, while this deferral repeats per reconciler
- // pass.
+ 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) => {
+ // Refused, not logged. The segment is already out of
the
+ // in-memory chain, so a file left behind becomes a
+ // non-tail empty segment as soon as the plant lands --
+ // `[sized][stale empty][planted]` -- which the next
boot
+ // refuses outright as `EmptyNonTailSegment`. Failing
boot
+ // here says so while the directory is still readable.
+ error!(
+ 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; refusing to plant beside it"
+ );
+ return Err(IggyError::CannotDeleteFile);
+ }
+ }
+ }
+ 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"
+ );
+ // Boot DOES count the recovered chain -- `load_persisted_segments`
+ // increments per segment before it looks at the size, so empty
tails
+ // are in the total -- and retention pairs its own retire with a
+ // decrement. Without this the count stays one high on the wire for
+ // the life of the process.
+ self.stats.decrement_segments_count(1);
+ retired += 1;
+ }
+ // Durable before anything is planted beside them: a crash in between
+ // would boot the stale name back into the chain.
+ 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"
+ );
+ }
+ // Bounds copied out: the plant below takes `&mut self`, so the borrow
on
+ // the chain cannot still be live.
+ let tail = self
+ .log
+ .segments()
+ .last()
+ .map(|segment| (segment.start_offset, segment.end_offset,
segment.size));
+ match tail {
+ // An EMPTIED chain still needs the plant, and it cannot be left to
+ // the caller's `ensure_initial_segment`, which names the segment
for
+ // the COMMITTED frontier and knows nothing of the append point. On
+ // the shape a crash before the first flush leaves -- committed
+ // frontier 0, append point a lease block up -- that plants
+ // `0.log` and then mints inside it, the hole this function exists
to
+ // prevent. The index does not save it either: a first entry that
is
+ // not the file-name offset makes recovery discard the index and
walk
+ // from byte 0, where the discontinuity tombstones the partition.
+ //
+ // No anchor: with nothing before it the plant leaves no gap, so
the
+ // chain guard has no pair to judge.
+ None => {
+ self.install_empty_segment(config, frontier).await?;
+ self.stats.increment_segments_count(1);
+ tracing::info!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ offset_frontier = frontier,
+ "planted a fresh segment at the restored offset frontier
over an \
+ empty recovered chain"
+ );
+ }
+ // 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.
+ Some((sealed_start, sealed_end, size))
+ if size.as_bytes_u64() > 0 && sealed_end.saturating_add(1) <
frontier =>
+ {
+ self.record_reanchor_gap(frontier, sealed_start, sealed_end)
+ .await?;
+ self.rotate_segment_at(config, frontier).await?;
+ tracing::info!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ sealed_end,
+ offset_frontier = frontier,
+ "sealed the recovered tail and planted a fresh segment at
the \
+ restored offset frontier"
+ );
+ }
+ Some(_) => {}
+ }
+ Ok(())
+ }
+
+ /// Write the anchor that makes the gap a plant at `frontier` leaves
+ /// legitimate, and make it durable before the segment exists.
+ ///
+ /// The chain guard admits a forward gap only when the far side carries an
+ /// anchor naming exactly the near side, so this record is what separates
the
+ /// re-anchor's own gap from a lost segment. Ordering is load-bearing in
one
+ /// direction only: an anchor with no segment is swept at the next boot,
+ /// while a segment with no anchor reads as damage.
+ ///
+ /// # Errors
+ /// [`IggyError::CannotCreateSegmentLogFile`] naming the anchor path. The
+ /// caller must not plant.
+ #[allow(clippy::future_not_send)]
+ async fn record_reanchor_gap(
+ &self,
+ frontier: u64,
+ sealed_start: u64,
+ sealed_end: u64,
+ ) -> Result<(), IggyError> {
+ // No directory means an in-memory partition, whose chain no boot
reads.
+ let Some(partition_dir) = self.partition_dir() else {
+ return Ok(());
+ };
+ let anchor = crate::segment_anchor::SegmentAnchor {
+ sealed_start,
+ sealed_end,
+ };
+ if let Err(error) =
+ crate::segment_anchor::write_anchor(&partition_dir, frontier,
anchor).await
Review Comment:
anchors have no lifecycle - every unlink site omits `.anchor`, including the
quarantine and converge sweeps that filter a directory listing to
`.log`/`.index`/`.staging`. purge resets the offset space to 0 and leaves the
anchor there, where two of `covers`'s three checks match for free and only the
far-side offset has to coincide. add `ANCHOR_EXTENSION` to those sets.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -998,6 +1208,158 @@ where
self.write_superblock(superblock.as_ref(), frontier).await
}
+ /// Whether the offset reservation is close enough to being consumed that
it
+ /// should be extended NOW, off the append path.
+ ///
+ /// The append fence is correct but badly placed: it writes the superblock
+ /// inline in the shard's frame pump, where the consensus tick is a sibling
+ /// arm, so its two fsyncs delay heartbeats for every group on the core.
The
+ /// fix is to make the fence's fast path
+ /// (`durable_offset_reserved > end_offset`) the only path it ever takes
under
+ /// load, by extending from the tick instead.
+ ///
+ /// HALF a block of headroom, which is a wide margin on purpose:
over-claiming
+ /// costs nothing but offset space, while arriving late puts the write
back on
+ /// the append path. A partition that has never minted is skipped -- its
first
+ /// append legitimately pays for the first claim, and extending every idle
+ /// partition at boot would write a superblock per partition for nothing.
+ #[must_use]
+ pub fn needs_offset_reservation_extension(&self) -> bool {
+ if self.consensus.replica_count() > 1
+ || self.superblock.is_none()
+ || !self.should_increment_offset
+ {
+ return false;
+ }
+ let headroom = self
+ .durable_offset_reserved
+ .get()
+ .saturating_sub(self.mint_frontier());
+ headroom < self.offset_reservation_lease / 2
+ }
+
+ /// Extend the reservation a full block past the current append point.
+ ///
+ /// Pairs with [`Self::needs_offset_reservation_extension`]; the caller is
the
+ /// shard tick, so this write is off the append path. A failure needs no
+ /// handling beyond the writer's own logging and backoff: the fence at the
+ /// mint is still there, and it is what refuses the append if the ceiling
+ /// never caught up.
+ #[allow(clippy::future_not_send)]
+ pub async fn extend_offset_reservation(&self) -> bool {
+ let Some(superblock) = self.superblock.as_ref().map(Rc::clone) else {
+ return true;
+ };
+ if self.superblock_write_is_backed_off() {
+ return false;
+ }
+ let _superblock_guard = self.superblock_lock.acquire().await;
+ self.write_offset_claim(superblock.as_ref(), self.mint_frontier())
+ .await
+ }
+
+ /// Upper bound on the offsets a pending `SendMessages` request will mint,
for
+ /// fencing it BEFORE it enters the pipeline.
+ ///
+ /// `project` assigns an op, not a base offset, so the exact range is
unknown
+ /// until the mint runs under `write_lock`. This is deliberately loose: a
+ /// request pipelined behind others can land above it, and the fence at the
+ /// mint stays as the exact check. It does not need to be tight -- the
claim
+ /// runs a whole lease block past whatever it is handed, so one of these
+ /// covers every batch in flight unless a run of them crosses a block
+ /// boundary.
+ ///
+ /// `None` when the body is not one canonical batch, which
+ /// `convert_request_message` has already rejected by the time this runs.
+ fn request_mint_ceiling(&self, message: &Message<RoutedRequestHeader>) ->
Option<u64> {
+ let body = message
+ .as_slice()
+
.get(std::mem::size_of::<RoutedRequestHeader>()..message.header().size as
usize)?;
+ let count = decode_batch_slice(body).ok()?.message_count();
+ Some(self.mint_frontier().saturating_add(u64::from(count)))
+ }
+
+ /// 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 a NEWLY minted 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. Journal repair
+ /// (`append_repaired_send_messages`) is the exception and needs none: it
+ /// re-journals offsets a peer already minted and fenced, so there is
nothing
+ /// new to claim. Claiming through `end_offset + 1 + lease` rather than
from
+ /// the live counter needs no special case for an oversized batch.
+ ///
+ /// SOLO ONLY, like everything the reservation feeds:
`restore_offset_frontier`
+ /// seeds no counter from it above one replica, and the boot re-anchor that
+ /// shapes the chain around it never runs there either. A replicated group
+ /// paying a superblock write per block would buy nothing -- and an ack
there
+ /// already means a quorum journaled the batch, so re-minting needs a
+ /// FULL-cluster crash.
+ ///
+ /// `false` when a write was attempted and failed, and also when an earlier
+ /// failure's backoff window is still open, in which case no write is
+ /// attempted at all -- that cell is shared with every other superblock
writer
+ /// on the partition, so this can refuse with no fault on the append path
+ /// itself. Either way the caller must refuse the append. Fail-closed: the
+ /// send is rejected with nothing externalised, exactly as a failed view
+ /// persist withholds its sends.
+ ///
+ /// WHERE it is called decides how much a refusal costs. Ahead of the
pipeline
+ /// (`on_request`) the client gets a retryable transient and the group
keeps
+ /// serving. At the mint the op already has its number and its ack is
already
+ /// skipped, so `commit_max` can never pass it and nothing later can commit
+ /// either: `on_replicate` fences the partition there and takes the node
down,
+ /// because a one-second backoff must not cost a partition the rest of the
+ /// process's life in the dark.
+ #[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 {
+ if self.consensus.replica_count() > 1 {
+ return true;
+ }
+ // 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() {
Review Comment:
one failed superblock write by any writer on this partition - a purge
frontier reset, the tick's own extension - opens a 20 ms window, up to 1 s
after repeats, where an already-admitted send whose mint crosses the lease
boundary sets `fatal` and takes the node down. rare at the default lease,
routine at a small one. give it the grace `superblock_wedged` grants.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -4189,41 +4613,245 @@ where
Ok(())
}
- /// Record the purge's frontier reset BEFORE the purge touches anything.
+ /// Re-anchor the append point after boot re-seeded the offset counter
above
+ /// what the recovered segment chain holds.
///
- /// The unlinks are made durable by their own directory fsync, so a crash
- /// between them and a reset written afterwards boots a purged directory
- /// whose record still names the pre-purge offset space:
- /// `restore_offset_frontier` re-seeds the counter to it while every peer
- /// restarted at 0, and the first append stamps a `base_offset` and
- /// `batch_checksum` no peer shares. Writing 0 first inverts the window
into
- /// a harmless one -- the record under-claims while the segments still
- /// exist, and boot takes the max of the record and what the segments
prove.
+ /// A hole INSIDE a segment is not survivable: `recover_segment_bounds`
walks
+ /// a segment from its FILENAME with a running `expected_offset` and
REFUSES
+ /// at the first offset that does not continue it (`OffsetDiscontinuity`),
+ /// which on a solo group tombstones the partition. A surviving index does
not
+ /// help: a first entry that is not the file-name offset makes recovery
+ /// discard the index and walk from byte 0, reaching the same refusal. On a
+ /// segment BOUNDARY every reader copes -- absolute offsets in the index,
+ /// `disk_poll_start` walking on into later segments, and a chain guard
that
+ /// admits a forward gap the reservation covers.
///
- /// Spelled out rather than read off the counter, which still holds the
- /// pre-purge frontier at this point.
+ /// So an empty tail is unlinked (its name claims a range it does not
hold), a
+ /// sized tail (the only copy of its messages) is sealed with a fresh
segment
+ /// planted at the append point, and a chain the unlinks emptied is planted
+ /// directly -- `ensure_initial_segment` names its segment for the
COMMITTED
+ /// frontier and would put the first mint inside it.
///
/// # Errors
- /// [`PurgeError::FrontierNotRecorded`]. Refused rather than logged:
nothing
- /// has been mutated yet, and a purge that cannot record its reset must not
- /// be the one that erases the data proving the old frontier. The caller
- /// RETRIES; it must not fence, since the chain is still whole and the live
- /// counter still names the pre-purge space.
+ /// [`IggyError`] when the fresh segment cannot be created, leaving the
+ /// partition without a serviceable chain.
#[allow(clippy::future_not_send)]
- async fn record_purge_frontier_reset(&mut self, generation: u64) ->
Result<(), PurgeError> {
- if self.reset_offset_frontier_at(0).await {
- self.purge_deferred = false;
+ 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.mint_frontier();
+ if frontier == 0 {
return Ok(());
}
- self.purge_deferred = true;
- // The ONLY operator-visible signal for the withhold: `send_prepare_ok`
- // returns silently, correctly, since it runs per prepare. So this line
- // has to say that the replica is now out of quorum for this group, or
- // the symptom reads as a network fault. The consecutive count
- // correlates it with the superblock writer's own error log, which
- // carries the `ENOSPC` / `EIO` cause but is rate-limited to
- // power-of-two failures, while this deferral repeats per reconciler
- // pass.
+ 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) => {
+ // Refused, not logged. The segment is already out of
the
+ // in-memory chain, so a file left behind becomes a
+ // non-tail empty segment as soon as the plant lands --
+ // `[sized][stale empty][planted]` -- which the next
boot
+ // refuses outright as `EmptyNonTailSegment`. Failing
boot
+ // here says so while the directory is still readable.
+ error!(
+ 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; refusing to plant beside it"
+ );
+ return Err(IggyError::CannotDeleteFile);
+ }
+ }
+ }
+ 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"
+ );
+ // Boot DOES count the recovered chain -- `load_persisted_segments`
+ // increments per segment before it looks at the size, so empty
tails
+ // are in the total -- and retention pairs its own retire with a
+ // decrement. Without this the count stays one high on the wire for
+ // the life of the process.
+ self.stats.decrement_segments_count(1);
+ retired += 1;
+ }
+ // Durable before anything is planted beside them: a crash in between
Review Comment:
the comment promises durability before anything is planted, but the fsync
error is only logged. the sized-tail arm is saved by `write_anchor`'s own
directory fsync; the emptied-chain arm goes through `install_empty_segment`,
which has neither - so the promise holds by accident.
##########
core/partitions/src/segment_anchor.rs:
##########
@@ -0,0 +1,205 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! The record that makes a gap in the segment chain legitimate.
+//!
+//! Recovery derives each segment's bounds from its own bytes, so a gap has no
+//! author: the boot re-anchor's planted gap and a lost segment look identical.
+//! The re-anchor writes its intent down instead, and the chain guard admits a
+//! forward gap only when the far side carries an anchor naming exactly the
near
+//! side.
+//!
+//! Written and directory-fsynced BEFORE that segment is created. A crash in
the
+//! window then leaves an anchor with no segment, which the boot sweep
collects;
+//! the other order leaves a planted segment with no anchor, which the guard
+//! reads as damage on an intact chain.
+
+use compio::io::AsyncWriteAtExt;
+use consensus::state_artifact_checksum;
+use std::io;
+
+/// File extension for an anchor record, `{start_offset:020}.anchor` beside the
+/// `{start_offset:020}.log` it belongs to.
+pub const ANCHOR_EXTENSION: &str = "anchor";
+
+/// Leading bytes of an anchor record, so a file that is not one (a truncated
+/// write, an operator's copy) is refused rather than decoded.
+const ANCHOR_MAGIC: u64 = u64::from_le_bytes(*b"IGGYANCH");
+
+/// `magic`(8) + `sealed_start`(8) + `sealed_end`(8) + `checksum`(8).
+pub const ANCHOR_ENCODED_LEN: usize = 32;
+
+/// The gap one planted segment is allowed to leave behind it.
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub struct SegmentAnchor {
+ /// Start offset of the segment that was sealed, i.e. the near side of the
+ /// gap. Names WHICH segment, so an anchor cannot be satisfied by a
+ /// different file that happens to end where this one expects.
+ pub sealed_start: u64,
+ /// End offset the sealed segment held when it was sealed. The gap runs
from
+ /// here to the planted segment's own start offset.
+ pub sealed_end: u64,
+}
+
+impl SegmentAnchor {
+ /// Encode to the fixed little-endian on-disk layout.
+ #[must_use]
+ pub fn to_bytes(&self) -> [u8; ANCHOR_ENCODED_LEN] {
+ let mut out = [0u8; ANCHOR_ENCODED_LEN];
+ out[0..8].copy_from_slice(&ANCHOR_MAGIC.to_le_bytes());
+ out[8..16].copy_from_slice(&self.sealed_start.to_le_bytes());
+ out[16..24].copy_from_slice(&self.sealed_end.to_le_bytes());
+ let checksum = state_artifact_checksum(&out[0..24]);
+ out[24..32].copy_from_slice(&checksum.to_le_bytes());
+ out
+ }
+
+ /// Decode a record, returning `None` for anything this build did not
write:
+ /// a wrong length, a wrong magic, or a checksum that does not match.
+ ///
+ /// A `None` is never treated as "no gap was intended". It means the record
+ /// proves nothing, so the gap it would have covered stays damage.
+ #[must_use]
+ pub fn from_bytes(bytes: &[u8]) -> Option<Self> {
+ if bytes.len() != ANCHOR_ENCODED_LEN {
+ return None;
+ }
+ let field = |at: usize| -> u64 {
+ let mut raw = [0u8; 8];
+ raw.copy_from_slice(&bytes[at..at + 8]);
+ u64::from_le_bytes(raw)
+ };
+ if field(0) != ANCHOR_MAGIC || field(24) !=
state_artifact_checksum(&bytes[0..24]) {
+ return None;
+ }
+ Some(Self {
+ sealed_start: field(8),
+ sealed_end: field(16),
+ })
+ }
+
+ /// Whether this anchor legitimises the gap between the segment starting at
+ /// `sealed_start` / ending at `sealed_end` and the segment it sits beside.
+ ///
+ /// BOTH bounds must match. Matching only the end offset would let an
anchor
+ /// left by an earlier incarnation of the chain cover a gap it never saw.
+ #[must_use]
+ pub const fn covers(&self, sealed_start: u64, sealed_end: u64) -> bool {
+ self.sealed_start == sealed_start && self.sealed_end == sealed_end
+ }
+}
+
+/// Path of the anchor record beside the segment starting at `start_offset`.
+#[must_use]
+pub fn anchor_path(partition_dir: &str, start_offset: u64) -> String {
+ format!("{partition_dir}/{start_offset:0>20}.{ANCHOR_EXTENSION}")
+}
+
+/// Write the anchor for a segment about to be planted at `start_offset`, then
+/// fsync the directory so the record cannot arrive after the segment it
+/// describes.
+///
+/// # Errors
+///
+/// Any I/O failure. The caller must NOT plant the segment: a planted segment
+/// whose anchor is missing reads as damage on the next boot.
+pub async fn write_anchor(
+ partition_dir: &str,
+ start_offset: u64,
+ anchor: SegmentAnchor,
+) -> io::Result<()> {
+ let path = anchor_path(partition_dir, start_offset);
+ let mut file = compio::fs::File::create(&path).await?;
Review Comment:
`File::create` truncates in place with no second slot, unlike the superblock
in this same directory (tmp, `sync_all`, rename, dir fsync). unreachable today
because the boot sweep drops any anchor whose `.log` is missing, and a live
`.log` cannot be re-selected by `sealed_end + 1 < frontier`.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -1812,14 +2184,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;
Review Comment:
dead rebinding left over from `next.max(self.mint_floor())`. rename the
binding above back to `dirty_offset` and drop this line.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -677,22 +692,78 @@ where
/// THIS partition's group is fenced; the rest of the node keeps serving.
#[allow(clippy::future_not_send)]
async fn write_superblock(&self, superblock: &SB, offset_frontier: u64) ->
bool {
- // ADVANCE direction: never below what this replica has already minted,
- // and never below what the record ALREADY holds. Both bounds are
- // needed and neither implies the other -- a failed install leaves the
- // counter behind the record it wrote before the swap, so maxing
against
- // the counter alone lets the fence that follows lower the durable
- // frontier. The reset direction goes through `write_superblock_inner`.
- let advanced = offset_frontier
- .max(self.offset_frontier())
- .max(self.durable_offset_frontier.get());
- self.write_superblock_inner(superblock, advanced).await
+ self.write_superblock_advancing(superblock, offset_frontier, 0)
+ .await
+ }
+
+ /// [`Self::write_superblock`] for a caller that also has a reservation to
+ /// claim. Both fields advance, neither can regress.
+ #[allow(clippy::future_not_send)]
+ async fn write_superblock_advancing(
+ &self,
+ superblock: &SB,
+ offset_frontier: u64,
+ offset_reserved: u64,
+ ) -> bool {
+ // ADVANCE direction; the reset direction goes through
+ // `write_superblock_inner`. Both bounds inside `advanced_frontier` are
+ // needed: a failed install leaves the chain behind the record it wrote
+ // before the swap, so maxing against the data alone would let the
fence
+ // that follows lower the durable frontier.
+ let advanced = self.advanced_frontier(offset_frontier);
+ // Nothing but the record witnesses a reservation, so a caller with no
+ // claim of its own (every view-change write) passes 0 and carries the
+ // recorded one forward; dropping it would let the next boot seed the
+ // counter below what an earlier append already fenced.
+ let reserved = offset_reserved.max(self.durable_offset_reserved.get());
+ self.write_superblock_inner(superblock, advanced, reserved)
+ .await
+ }
+
+ /// The advance rule for the frontier, shared by every writer that claims
+ /// one: never below what this replica holds, never below what the record
+ /// already says. Held messages, NOT [`Self::mint_frontier`], which stands
a
+ /// lease block above them after a reservation-seeded boot and names none
of
+ /// them.
+ fn advanced_frontier(&self, claim: u64) -> u64 {
+ claim
+ .max(self.held_offset_frontier())
+ .max(self.durable_offset_frontier.get())
+ }
+
+ /// Record an incoming state-transfer frontier, advancing the frontier but
+ /// SETTING the reservation to it.
Review Comment:
not what happens - `write_superblock_inner` clamps the reservation to
`max(offset_frontier)`, so it lands on `advanced_frontier(frontier)`, above the
offer whenever an earlier over-claiming write left the durable frontier higher.
say to the frontier this write records.
##########
core/consensus/src/vsr_state.rs:
##########
@@ -121,16 +142,26 @@ 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, the pre-`offset_reserved`
record
+ // every tagged release wrote. The one length check then puts every
field
+ // slice below in bounds by construction, so the `try_into`s cannot
fail.
+ //
+ // `offset_reserved` is filled from `offset_frontier`, not zeroed. The
two
+ // agree on a record written before the reservation existed: the
frontier
+ // is what that build's data proved, and the write side clamps the
+ // reservation up to it anyway, so this is the same value a first write
+ // under this build would record. A 0 would instead claim "nothing
+ // reserved" for offsets the frontier says exist.
+ //
+ // It cannot recover what the old build never wrote down -- offsets
acked
+ // out of RAM above the frontier are gone with the process either way
--
+ // but refusing the record recovers nothing and costs the node its
boot.
let mut padded = [0u8; ENCODED_LEN];
match bytes.len() {
ENCODED_LEN => padded.copy_from_slice(bytes),
- ENCODED_LEN_WITHOUT_FRONTIER => {
- padded[..ENCODED_LEN_WITHOUT_FRONTIER].copy_from_slice(bytes);
+ ENCODED_LEN_WITHOUT_RESERVATION => {
+
padded[..ENCODED_LEN_WITHOUT_RESERVATION].copy_from_slice(bytes);
+ padded[66..74].copy_from_slice(&bytes[58..66]);
Review Comment:
the `66` here is a literal two lines after `ENCODED_LEN_WITHOUT_RESERVATION`
names it. `padded[ENCODED_LEN_WITHOUT_RESERVATION..ENCODED_LEN]` is the whole
fix - the `58` has no constant and numeric offsets are this file's convention.
##########
core/consensus/src/vsr_state.rs:
##########
@@ -31,20 +31,26 @@ 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).
+///
+/// Growing this is one-way: a record of this length is
[`VsrStateError::WrongLength`]
Review Comment:
a rollback needing the data directory wiped is an operator instruction
living only in rustdoc on a const. put it where someone planning a downgrade
would look.
##########
core/partitions/src/state_transfer.rs:
##########
@@ -2044,7 +2051,7 @@ where
self.reset_offset_frontier_at(offsets_wire.next_offset)
.await
} else {
- self.persist_offset_frontier_at(offsets_wire.next_offset)
+ self.install_offset_frontier_at(offsets_wire.next_offset)
Review Comment:
does this reset buy anything today? the fence never runs above one replica,
so on every group that can receive an offer the reservation already equals the
frontier and `advanced_frontier(offer)` is never below it - the write is
identical to the advancing form's.
##########
core/configs/src/server_config/partition.rs:
##########
@@ -177,6 +191,16 @@ impl Validatable<ConfigurationError> for PartitionConfig {
);
return Err(ConfigurationError::InvalidConfigurationValue);
}
+ if self.offset_reservation_lease == 0
Review Comment:
the validation landed but `mod tests` has no zero or above-ceiling case for
this knob, while `prepare_queue_depth` and `evicted_ring_capacity` each keep
theirs.
##########
core/integration/tests/cluster/crash_offset_reuse.rs:
##########
@@ -115,20 +110,63 @@ async fn
given_confirmed_sends_below_flush_threshold_when_a_solo_node_is_killed_
&TopicCreateOptions {
partitions_count: Some(1),
message_expiry: Some(IggyExpiry::NeverExpire),
+ messages_required_to_save,
..TopicCreateOptions::default()
},
)
.await
.expect("create topic");
+}
- let acked = produce_acked(&client, "pre-crash", PRE_CRASH_SENDS).await;
- let highest_confirmed = *acked.last().expect("confirmed sends");
- drop(client);
+/// Base offsets of every segment file under `root`, from the file names, which
+/// are the on-disk claim about where each range begins.
+///
+/// Reading them is the only way to assert the shape the re-anchor produces. A
+/// black-box offset assertion passes either way on the boot that WRITES the
+/// wrong shape; the cost only lands on the boot that reads it back, where the
+/// walk refuses and the solo arm tombstones the partition.
+fn segment_base_offsets(root: &Path) -> Vec<u64> {
+ let mut offsets = Vec::new();
+ let mut stack = vec![root.to_path_buf()];
+ while let Some(dir) = stack.pop() {
+ let Ok(entries) = std::fs::read_dir(&dir) else {
+ continue;
+ };
+ for entry in entries.flatten() {
+ let path = entry.path();
+ if path.is_dir() {
+ stack.push(path);
+ } else if path.extension().is_some_and(|extension| extension ==
"log")
+ && let Some(stem) = path.file_stem().and_then(|stem|
stem.to_str())
+ && let Ok(offset) = stem.parse::<u64>()
+ {
+ offsets.push(offset);
+ }
+ }
+ }
+ offsets.sort_unstable();
+ offsets
+}
+/// Kill the node, bring it back, and return a client onto the restarted one.
+async fn crash_and_recover(harness: &mut TestHarness) -> IggyClient {
harness.kill_node(0).expect("SIGKILL the only node");
harness.restart_node(0).expect("restart it");
+ wait_until_serving(harness, SERVE_TIMEOUT).await
+}
- let client = wait_until_serving(harness, SERVE_TIMEOUT).await;
+#[iggy_harness(cluster_nodes = 1)]
Review Comment:
no test combines a flush with a graceful stop - `:220` stops cleanly but
takes no traffic between lives, and `:299` flushes then re-enters through
SIGKILL. `ensure_initial_segment` also never runs with a non-zero frontier: the
only superblock fixture hardcodes `offset_reserved: 0`.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -3846,98 +4319,49 @@ where
}
async fn rotate_segment(&mut self, config: &PartitionsConfig) ->
Result<(), IggyError> {
- let namespace = self.namespace();
- let old_segment_index = self.log.segments().len() - 1;
- let active_segment = self.log.active_segment_mut();
- active_segment.sealed = true;
- let start_offset = active_segment.end_offset + 1;
-
- let segment_size = self.effective_segment_size(config);
- let enforce_fsync = self.effective_enforce_fsync(config);
- let preallocate_segments = self.effective_preallocate_segments(config);
- let segment = Segment::new(start_offset, segment_size);
- // Prefer the active writer's location: a per-topic path override or a
- // config change after the initial segment was created must not scatter
- // one partition's segments across two directories. The config layout
- // only decides for a partition with no writer yet.
- let (messages_path, index_path) = self.partition_dir().map_or_else(
- || {
- (
- config.get_messages_path(
- namespace.stream_id(),
- namespace.topic_id(),
- namespace.partition_id(),
- start_offset,
- ),
- config.get_index_path(
- namespace.stream_id(),
- namespace.topic_id(),
- namespace.partition_id(),
- start_offset,
- ),
- )
- },
- |dir| {
- (
- format!("{dir}/{start_offset:0>20}.log"),
- format!("{dir}/{start_offset:0>20}.index"),
- )
- },
- );
+ let start_offset = self.log.active_segment().end_offset + 1;
+ self.rotate_segment_at(config, start_offset).await
+ }
- let storage = SegmentStorage::new(&messages_path, &index_path, 0, 0,
false)
- .await
- .map_err(|_|
IggyError::CannotCreateSegmentLogFile(messages_path.clone()))?;
- let messages_size_bytes = storage
- .messages_writer
- .as_ref()
- .ok_or_else(||
IggyError::CannotCreateSegmentLogFile(messages_path.clone()))?
- .size_counter();
- let messages_writer = Rc::new(
- MessagesWriter::new(
- &messages_path,
- messages_size_bytes,
- enforce_fsync,
- false,
- preallocate_segments.then_some(segment_size),
- )
- .await
- .map_err(|_|
IggyError::CannotCreateSegmentLogFile(messages_path.clone()))?,
- );
- let index_size_bytes = storage
- .index_writer
- .as_ref()
- .ok_or_else(||
IggyError::CannotCreateSegmentIndexFile(index_path.clone()))?
- .size_counter();
- let index_writer = Rc::new(
- IggyIndexWriter::new(&index_path, index_size_bytes, enforce_fsync,
false)
- .await
- .map_err(|_|
IggyError::CannotCreateSegmentIndexFile(index_path.clone()))?,
- );
+ /// Seal the active segment and plant a fresh empty one at `start_offset`.
+ ///
+ /// Shared by the size-driven roll, which plants at `end_offset + 1`, and
the
+ /// boot re-anchor, which plants at the append point the reservation moved
the
+ /// counter to. One seal path, and one order: the new segment's files are
+ /// created BEFORE the sealed segment's writers are torn down, so a failed
+ /// create leaves the chain serviceable.
+ async fn rotate_segment_at(
Review Comment:
nothing here ties `start_offset` to the sealed tail or to an anchor having
been written. both callers are correct today, but the obligation lives only in
prose at :4754 - a `debug_assert!` would stop a third caller inheriting it
silently.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -4189,41 +4613,245 @@ where
Ok(())
}
- /// Record the purge's frontier reset BEFORE the purge touches anything.
+ /// Re-anchor the append point after boot re-seeded the offset counter
above
+ /// what the recovered segment chain holds.
///
- /// The unlinks are made durable by their own directory fsync, so a crash
- /// between them and a reset written afterwards boots a purged directory
- /// whose record still names the pre-purge offset space:
- /// `restore_offset_frontier` re-seeds the counter to it while every peer
- /// restarted at 0, and the first append stamps a `base_offset` and
- /// `batch_checksum` no peer shares. Writing 0 first inverts the window
into
- /// a harmless one -- the record under-claims while the segments still
- /// exist, and boot takes the max of the record and what the segments
prove.
+ /// A hole INSIDE a segment is not survivable: `recover_segment_bounds`
walks
+ /// a segment from its FILENAME with a running `expected_offset` and
REFUSES
+ /// at the first offset that does not continue it (`OffsetDiscontinuity`),
+ /// which on a solo group tombstones the partition. A surviving index does
not
+ /// help: a first entry that is not the file-name offset makes recovery
+ /// discard the index and walk from byte 0, reaching the same refusal. On a
+ /// segment BOUNDARY every reader copes -- absolute offsets in the index,
+ /// `disk_poll_start` walking on into later segments, and a chain guard
that
+ /// admits a forward gap the reservation covers.
///
- /// Spelled out rather than read off the counter, which still holds the
- /// pre-purge frontier at this point.
+ /// So an empty tail is unlinked (its name claims a range it does not
hold), a
+ /// sized tail (the only copy of its messages) is sealed with a fresh
segment
+ /// planted at the append point, and a chain the unlinks emptied is planted
+ /// directly -- `ensure_initial_segment` names its segment for the
COMMITTED
+ /// frontier and would put the first mint inside it.
///
/// # Errors
- /// [`PurgeError::FrontierNotRecorded`]. Refused rather than logged:
nothing
- /// has been mutated yet, and a purge that cannot record its reset must not
- /// be the one that erases the data proving the old frontier. The caller
- /// RETRIES; it must not fence, since the chain is still whole and the live
- /// counter still names the pre-purge space.
+ /// [`IggyError`] when the fresh segment cannot be created, leaving the
+ /// partition without a serviceable chain.
#[allow(clippy::future_not_send)]
- async fn record_purge_frontier_reset(&mut self, generation: u64) ->
Result<(), PurgeError> {
- if self.reset_offset_frontier_at(0).await {
- self.purge_deferred = false;
+ pub async fn reanchor_to_offset_frontier(
Review Comment:
a data dir written by this branch's earlier push has no anchor, so this
build boots it as `Hole` and the solo arm tombstones it. master plants no gaps
and edge builds come only from master, so nothing deployed regresses - but
anyone who ran the earlier push needs a line in the PR description telling them
to wipe.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -998,6 +1208,158 @@ where
self.write_superblock(superblock.as_ref(), frontier).await
}
+ /// Whether the offset reservation is close enough to being consumed that
it
+ /// should be extended NOW, off the append path.
+ ///
+ /// The append fence is correct but badly placed: it writes the superblock
+ /// inline in the shard's frame pump, where the consensus tick is a sibling
+ /// arm, so its two fsyncs delay heartbeats for every group on the core.
The
+ /// fix is to make the fence's fast path
+ /// (`durable_offset_reserved > end_offset`) the only path it ever takes
under
+ /// load, by extending from the tick instead.
+ ///
+ /// HALF a block of headroom, which is a wide margin on purpose:
over-claiming
+ /// costs nothing but offset space, while arriving late puts the write
back on
+ /// the append path. A partition that has never minted is skipped -- its
first
+ /// append legitimately pays for the first claim, and extending every idle
+ /// partition at boot would write a superblock per partition for nothing.
+ #[must_use]
+ pub fn needs_offset_reservation_extension(&self) -> bool {
+ if self.consensus.replica_count() > 1
+ || self.superblock.is_none()
+ || !self.should_increment_offset
+ {
+ return false;
+ }
+ let headroom = self
+ .durable_offset_reserved
+ .get()
+ .saturating_sub(self.mint_frontier());
+ headroom < self.offset_reservation_lease / 2
Review Comment:
`lease / 2` is 0 when the lease is 1, which validation allows, so the
extension never fires and every append pays the inline claim - two fsyncs on
every multi-message batch. use `(lease / 2).max(1)`, or floor the config at 2.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -954,7 +1124,7 @@ where
let _superblock_guard = self.superblock_lock.acquire().await;
match intended {
Some(frontier) => {
- self.write_superblock_inner(superblock.as_ref(), frontier)
+ self.write_superblock_inner(superblock.as_ref(), frontier,
frontier)
Review Comment:
this writes `(f, f)` with no max against `durable_offset_reserved`, which
contradicts the monotone ceiling `vsr_state.rs:106` claims. inert today - the
only `Some` caller is the replicated `ConvergeFailed` arm, where the fence
never ran and `reserved == frontier`. carry the max anyway?
--
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]