This is an automated email from the ASF dual-hosted git repository.
mmodzelewski pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iggy.git
The following commit(s) were added to refs/heads/master by this push:
new 70903b836 fix(server-ng): stop redundant partition rebuilds resetting
offsets (#3846)
70903b836 is described below
commit 70903b8368340fb9a29b21bb1193702cea78ed65
Author: Grzegorz Koszyk <[email protected]>
AuthorDate: Mon Aug 10 10:26:19 2026 +0200
fix(server-ng): stop redundant partition rebuilds resetting offsets (#3846)
---
core/partitions/src/iggy_partition.rs | 19 ++-
core/partitions/src/iggy_partitions.rs | 7 +
core/server-ng/src/bootstrap.rs | 1 -
core/server-ng/src/partition_helpers.rs | 1 -
core/server-ng/src/partition_reconciler.rs | 239 ++++++++++++++++++++---------
core/shard/src/lib.rs | 62 ++++++--
core/shard/src/metrics.rs | 12 ++
7 files changed, 258 insertions(+), 83 deletions(-)
diff --git a/core/partitions/src/iggy_partition.rs
b/core/partitions/src/iggy_partition.rs
index 148c4078b..945c0739e 100644
--- a/core/partitions/src/iggy_partition.rs
+++ b/core/partitions/src/iggy_partition.rs
@@ -770,7 +770,24 @@ where
self.offset.store(recovered_end, Ordering::Release);
self.dirty_offset.store(recovered_end, Ordering::Relaxed);
self.should_increment_offset = true;
- self.stats.set_current_offset(recovered_end);
+ }
+
+ /// Copy this incarnation's offset counter into the shared
+ /// [`PartitionStats`], making it the value readers (offset validation,
+ /// `get_topic`, `get_stats`) see.
+ ///
+ /// Called from [`IggyPartitions::insert`](crate::IggyPartitions::insert)
+ /// only: when the instance BECOMES the addressable one, never while
+ /// building it. The stats registry keys on the namespace, not the
+ /// incarnation, so every build of a namespace holds the same `Arc` as
+ /// whatever is already serving it -- and a build is not guaranteed to be
+ /// adopted. Seeding from the build instead leaves a zeroed
`current_offset`
+ /// on the live incarnation, which then rejects every
+ /// `store_consumer_offset` above 0 with `InvalidOffset` until the next
send
+ /// re-seeds it.
+ pub(crate) fn publish_current_offset(&self) {
+ self.stats
+ .set_current_offset(self.offset.load(Ordering::Acquire));
}
/// The next message offset this replica will mint, `0` while the offset
diff --git a/core/partitions/src/iggy_partitions.rs
b/core/partitions/src/iggy_partitions.rs
index a25aed657..5fac54528 100644
--- a/core/partitions/src/iggy_partitions.rs
+++ b/core/partitions/src/iggy_partitions.rs
@@ -200,6 +200,12 @@ where
/// Insert a new partition and return its local index.
///
+ /// Insertion is the moment a build becomes the addressable incarnation,
+ /// so this is also where its offset counter is published into the shared
+ /// `PartitionStats` ([`IggyPartition::publish_current_offset`]). No
+ /// earlier point is safe: a build that is never adopted must leave the
+ /// live incarnation's counters alone.
+ ///
/// # Safety discipline (compiler cannot enforce)
///
/// Must only be called from the shard's pump task (i.e. inside
@@ -220,6 +226,7 @@ where
0,
"IggyPartitions::insert while a with_partition borrow is live"
);
+ partition.publish_current_offset();
// Safety: pump-only invariant, caller responsibility.
let partitions = unsafe { &mut *self.partitions.get() };
let local_idx = LocalIdx::new(partitions.len());
diff --git a/core/server-ng/src/bootstrap.rs b/core/server-ng/src/bootstrap.rs
index ad1671b20..b3bc66b35 100644
--- a/core/server-ng/src/bootstrap.rs
+++ b/core/server-ng/src/bootstrap.rs
@@ -2413,7 +2413,6 @@ async fn load_partition(
partition.offset.store(counter, Ordering::Release);
partition.dirty_offset.store(counter, Ordering::Relaxed);
partition.should_increment_offset = current_offset.is_some();
- partition.stats.set_current_offset(counter);
// The durable frontier is a LOWER BOUND on top of what the segments
proved:
// it is the only carrier left when the segments that named the frontier
are
// gone (an all-GC'd origin's install, a crash inside the swap window), and
diff --git a/core/server-ng/src/partition_helpers.rs
b/core/server-ng/src/partition_helpers.rs
index a7a291eb7..129ffb537 100644
--- a/core/server-ng/src/partition_helpers.rs
+++ b/core/server-ng/src/partition_helpers.rs
@@ -694,7 +694,6 @@ pub async fn build_partition_fresh(
partition.offset.store(0, Ordering::Release);
partition.dirty_offset.store(0, Ordering::Relaxed);
partition.should_increment_offset = false;
- partition.stats.set_current_offset(0);
debug_assert!(
!partition.log.has_segments(),
"fresh partition must not carry recovered segments"
diff --git a/core/server-ng/src/partition_reconciler.rs
b/core/server-ng/src/partition_reconciler.rs
index 63d6c2a5c..cf6c86de0 100644
--- a/core/server-ng/src/partition_reconciler.rs
+++ b/core/server-ng/src/partition_reconciler.rs
@@ -413,6 +413,10 @@ struct PassCounters {
/// tombstone and re-wakes us without bumping `Streams::revision`, so an
/// armed skip would swallow that wake and strand the rebuild forever.
deferred: usize,
+ /// Namespaces an earlier pass already built, whose `InsertOwned` the pump
+ /// has not applied yet. Counted so the pass does not arm the fast-skip
+ /// while work is in flight; applying it bumps no revision.
+ already_staged: usize,
}
impl PassCounters {
@@ -428,6 +432,7 @@ impl PassCounters {
+ self.purges_staged
+ self.deferred
+ self.parked_reclaimed
+ + self.already_staged
}
}
@@ -476,9 +481,9 @@ async fn reconcile_once(ctx: &ReconcilerCtx) -> bool {
let target_set: AHashSet<IggyNamespace> = target.iter().map(|(ns, _)|
*ns).collect();
let mut counters = PassCounters::default();
- let staged = reconcile_additions(ctx, target, &mut counters).await;
+ reconcile_additions(ctx, target, &mut counters).await;
reconcile_removals(ctx, &target_set, &mut counters).await;
- reconcile_parked_frames(ctx, &staged, &mut counters);
+ reconcile_parked_frames(ctx, &mut counters);
reconcile_consumer_group_offsets(ctx, &mut counters).await;
reconcile_segment_truncations(ctx, &mut counters);
reconcile_partition_purges(ctx, &mut counters);
@@ -506,6 +511,7 @@ async fn reconcile_once(ctx: &ReconcilerCtx) -> bool {
backoff_skipped = counters.backoff_skipped,
stale = counters.stale,
deferred = counters.deferred,
+ already_staged = counters.already_staged,
parked_reclaimed = counters.parked_reclaimed,
purges_staged = counters.purges_staged,
trims_pending = counters.trims_pending,
@@ -521,19 +527,14 @@ async fn reconcile_once(ctx: &ReconcilerCtx) -> bool {
true
}
-/// Returns the namespaces whose `ReconcileOp::InsertOwned` this pass staged.
The
-/// pump applies the op on its own task, so they are not in `IggyPartitions`
yet
-/// and [`reconcile_parked_frames`] would read them as un-materialised, aging
-/// their frames on the pass that built them.
async fn reconcile_additions(
ctx: &ReconcilerCtx,
target: Vec<(IggyNamespace, u64)>,
counters: &mut PassCounters,
-) -> AHashSet<IggyNamespace> {
+) {
let shard_id = ctx.shard.id;
let partitions = ctx.shard.plane.partitions();
let total_shards = u32::from(ctx.total_shards);
- let mut staged = AHashSet::new();
for (ns, epoch) in target {
if partitions.contains(&ns) {
@@ -583,7 +584,7 @@ async fn reconcile_additions(
// means the local partition is a prior incarnation carrying
// stale segments/offsets/log. Tear it down; the
// post-ConfirmRemove wake rebuilds it fresh next pass.
- if ctx.shard.shards_table().epoch_for(ns) == Some(epoch) {
+ if shards_table_has_epoch(ctx, ns, epoch) {
continue;
}
trace!(
@@ -599,10 +600,15 @@ async fn reconcile_additions(
let owning_shard = calculate_shard_assignment(&ns, total_shards);
if owning_shard != shard_id {
- // Compare the epoch, not just presence: a delete + recreate
recycles
- // the slab keys, so the row survives with the DEAD incarnation's
- // `created_revision`. A presence-only gate never refreshes it, and
- // nothing else writes a non-owner's row.
+ // Compare the epoch, not just presence: a delete + recreate
+ // recycles the slab keys, so the row survives with the DEAD
+ // incarnation's `created_revision`. A presence-only gate never
+ // refreshes it, and nothing else writes a non-owner's row.
+ //
+ // No mirror of the staged-`InsertOwned` guard below, deliberately:
+ // a lagging pump costs one duplicate `InsertRouted` per pass, and
+ // the apply is an idempotent row overwrite, while scanning the op
+ // queue per routed namespace would go quadratic.
if !shards_table_has_epoch(ctx, ns, epoch) {
ctx.shard.enqueue_reconcile_op(ReconcileOp::InsertRouted {
namespace: ns,
@@ -614,6 +620,17 @@ async fn reconcile_additions(
continue;
}
+ // An earlier pass already built this one and the pump has not applied
it
+ // yet, so the `contains` test above reads false for finished work.
+ // Rebuilding is not a wasted-effort question: the second build shares
+ // the namespace's `PartitionStats` with the queued sibling and
re-opens
+ // segment 0 with `file_exists = false`, truncating the file that
+ // sibling is about to serve.
+ if ctx.shard.has_staged_insert_owned(ns) {
+ counters.already_staged += 1;
+ continue;
+ }
+
let now = Instant::now();
if ctx.is_backed_off(ns, FailureCause::Add, now) {
counters.backoff_skipped += 1;
@@ -648,7 +665,6 @@ async fn reconcile_additions(
});
ctx.record_success(ns, FailureCause::Add);
counters.materialised += 1;
- staged.insert(ns);
}
Err(err) => {
ctx.record_failure(ns, FailureCause::Add, now);
@@ -664,8 +680,6 @@ async fn reconcile_additions(
}
}
}
-
- staged
}
/// Retire parked frames the shard cannot serve, age the ones it might.
@@ -705,18 +719,15 @@ async fn reconcile_additions(
/// timeout and no committed op dies on a local-convergence signal. Residency
/// only; see `ParkedFrame::passes`.
///
-/// `staged_this_pass` is exempt: its `InsertOwned` is queued but not applied,
so
-/// it reads as un-materialised here. Not a one-pass concession.
-/// `reconcile_additions` has no cross-pass guard against a
queued-but-unapplied
-/// op (it tests `partitions.contains`, false the whole time it sits in the
-/// queue), so it re-stages every pass until the pump drains. The exemption
-/// therefore covers arbitrary pump lag; dropping it ages frames on every
-/// commit-driven pass the pump falls behind.
-fn reconcile_parked_frames(
- ctx: &ReconcilerCtx,
- staged_this_pass: &AHashSet<IggyNamespace>,
- counters: &mut PassCounters,
-) {
+/// A namespace with a staged, unapplied `InsertOwned` is exempt: its partition
+/// is on the way but reads as un-materialised here. The queue is asked per
+/// parked namespace ([`shard::IggyShard::has_staged_insert_owned`]) rather
than
+/// carrying a set over from the additions pass, so the answer cannot go stale
+/// across `reconcile_removals`' awaits; `parked` is empty on the steady path,
+/// so the scan costs nothing there. The exemption spans passes, not just the
+/// one that built the namespace, covering arbitrary pump lag; dropping it ages
+/// frames on every commit-driven pass the pump falls behind.
+fn reconcile_parked_frames(ctx: &ReconcilerCtx, counters: &mut PassCounters) {
let parked = ctx.shard.parked_namespaces();
if parked.is_empty() {
return;
@@ -724,7 +735,7 @@ fn reconcile_parked_frames(
let partitions = ctx.shard.plane.partitions();
let total_shards = u32::from(ctx.total_shards);
for ns in parked {
- if staged_this_pass.contains(&ns) {
+ if ctx.shard.has_staged_insert_owned(ns) {
continue;
}
// Tombstoned namespaces are still in the map, so `contains` below
reads
@@ -1147,7 +1158,8 @@ pub fn install_tick_handler(shard: &Rc<ServerNgShard>,
wake_tx: WakeTx) {
#[cfg(test)]
mod tests {
use super::{
- FailureCause, FailureRecord, ReconcilerCtx,
delete_partitions_from_disk, reconcile_once,
+ FailureCause, FailureRecord, ReconcilerCtx, build_partition_fresh,
+ delete_partitions_from_disk, fetch_partition_stats, reconcile_once,
};
use configs::server_ng::{NgSystemConfig, ServerNgConfig};
use consensus::{MetadataHandle, PartitionsHandle};
@@ -1181,6 +1193,7 @@ mod tests {
use std::mem::size_of;
use std::rc::Rc;
use std::sync::Arc;
+ use std::sync::atomic::Ordering;
use std::time::Instant;
use tempfile::TempDir;
@@ -1595,61 +1608,149 @@ mod tests {
);
}
- /// Regression (deferred-apply window): the reconciler stages
- /// `ReconcileOp::InsertOwned` from a task separate from the pump that
- /// applies it, so under a commit burst it can run a second pass before
- /// the pump drains the first pass's staged ops. Both passes then
- /// observe `!contains(ns)` and build the same namespace. The pump's
- /// apply must be idempotent, else the second `insert` orphans the first
- /// partition (leaked VSR group + writers) and inflates `len`.
- /// `reconcile_pass` applies inline and cannot surface this, so here we
- /// run two passes and only then drain once.
+ /// The cross-pass guard: a pass must not rebuild a namespace an earlier
pass
+ /// already built and left queued. Rebuilding is not merely wasted work --
+ /// the second build shares the namespace's `PartitionStats` with the
queued
+ /// sibling and re-opens segment 0 truncating -- so a pass has to recognise
+ /// the staged op, not just `partitions.contains`.
#[compio::test]
- async fn deferred_apply_window_does_not_duplicate_owned_partition() {
+ async fn second_pass_does_not_rebuild_a_namespace_already_staged() {
let tmp = TempDir::new().expect("tempdir for system path");
let config = test_config(&tmp);
let mux = TestMux::default();
seed_stream(&mux, 1, "stream-a");
- seed_topic(
- &mux,
- 2,
- 0,
- "topic-a",
- vec![assignment(0, 1), assignment(1, 2)],
- );
+ seed_topic(&mux, 2, 0, "topic-a", vec![assignment(0, 1)]);
let shard = build_test_shard(0, &config, mux);
let ctx = make_ctx(Rc::clone(&shard), 1, Rc::new(config));
+ let ns = IggyNamespace::new(0, 0, 0);
- // Two passes with no pump drain in between: models the reconciler,
- // woken by a second commit tick, running pass N+1 before the pump
- // applies pass N's `InsertOwned`. Both passes see the namespaces as
- // unmaterialised and stage a build for each, so the queue holds two
- // `InsertOwned` per namespace when the pump finally drains.
reconcile_once(&ctx).await;
+ assert!(
+ ctx.shard.has_staged_insert_owned(ns),
+ "the first pass must leave an unapplied InsertOwned to guard
against"
+ );
+
+ // Second pass while that op is still queued: `partitions.contains(ns)`
+ // is false, so only the staged-op guard can stop the rebuild.
reconcile_once(&ctx).await;
ctx.shard.apply_reconcile_ops();
+ assert_eq!(
+ shard.plane.partitions().len(),
+ 1,
+ "the namespace must materialise exactly once"
+ );
- let partitions = shard.plane.partitions();
+ // Counting builds, not partitions: the pump discards the redundant
+ // `InsertOwned` either way, so `len` cannot tell one build from two.
+ // `ensure_initial_segment` plants exactly one segment per build and
+ // folds it into the namespace's shared stats, so this counter is the
+ // observable that separates them.
+ let stats = fetch_partition_stats(&ctx, ns).expect("materialised
namespace has stats");
assert_eq!(
- partitions.len(),
- 2,
- "deferred-apply window must not duplicate partitions: \
- each namespace materialises exactly once"
+ stats.segments_count_inconsistent(),
+ 1,
+ "a second build ran: its initial segment was folded into the \
+ namespace's shared stats on top of the live incarnation's"
+ );
+ }
+
+ /// The stats registry keys on the namespace, not the incarnation, so
+ /// `current_offset` moves on adoption only: the pump seeds it from the
+ /// incarnation it inserts, and a build that never becomes addressable
+ /// leaves it alone. Seeding from the build instead zeroed it under the
live
+ /// incarnation, after which the partition plane's admission check read an
+ /// empty offset space and answered every `store_consumer_offset` above 0
+ /// with `InvalidOffset` (error 4100) until the next send re-seeded it.
+ ///
+ /// Both incarnations are built and staged by hand: the adopted one so its
+ /// counter is non-zero BEFORE insertion (a reconciler build always adopts
+ /// at 0, where the publish is indistinguishable from a no-op), and the
+ /// redundant one because `has_staged_insert_owned` now stops a pass from
+ /// producing it.
+ #[compio::test]
+ async fn discarded_build_leaves_live_partition_offset_intact() {
+ const COMMITTED_OFFSET: u64 = 3;
+ const LIVE_EPOCH: u64 = 1;
+
+ let tmp = TempDir::new().expect("tempdir for system path");
+ let config = test_config(&tmp);
+ let mux = TestMux::default();
+ seed_stream(&mux, 1, "stream-a");
+ seed_topic(&mux, 2, 0, "topic-a", vec![assignment(0, 1)]);
+
+ let shard = build_test_shard(0, &config, mux);
+ let ctx = make_ctx(Rc::clone(&shard), 1, Rc::new(config.clone()));
+ let ns = IggyNamespace::new(0, 0, 0);
+
+ let stats = fetch_partition_stats(&ctx, ns).expect("committed
namespace has stats");
+ let live = build_partition_fresh(
+ &config,
+ ns,
+ Arc::clone(&stats),
+ LIVE_EPOCH,
+ CLUSTER_ID,
+ 0,
+ 1,
+ Rc::clone(&ctx.shard.bus),
+ )
+ .await
+ .expect("live build succeeds");
+ // What a recovery leaves behind: an incarnation whose own counter is
+ // ahead of the zeroed shared stats until adoption publishes it.
+ live.offset.store(COMMITTED_OFFSET, Ordering::Release);
+ ctx.shard.enqueue_reconcile_op(ReconcileOp::InsertOwned {
+ namespace: ns,
+ partition: Box::new(live),
+ epoch: LIVE_EPOCH,
+ });
+ ctx.shard.apply_reconcile_ops();
+
+ assert_eq!(
+ stats.current_offset(),
+ COMMITTED_OFFSET,
+ "adoption must publish the incarnation's offset into the shared
stats"
+ );
+
+ let redundant = build_partition_fresh(
+ &config,
+ ns,
+ Arc::clone(&stats),
+ LIVE_EPOCH + 1,
+ CLUSTER_ID,
+ 0,
+ 1,
+ Rc::clone(&ctx.shard.bus),
+ )
+ .await
+ .expect("redundant build succeeds over the live incarnation's path");
+ ctx.shard.enqueue_reconcile_op(ReconcileOp::InsertOwned {
+ namespace: ns,
+ partition: Box::new(redundant),
+ epoch: LIVE_EPOCH + 1,
+ });
+ ctx.shard.apply_reconcile_ops();
+
+ assert_eq!(
+ shard.plane.partitions().len(),
+ 1,
+ "the redundant build must be discarded, not adopted: a second
insert \
+ overwrites the ns -> idx entry and orphans the first partition, \
+ leaking its VSR group and segment writers"
+ );
+ // The epoch, not `shard_for`: on a single shard an adopt would write
+ // `ShardId::new(0)` too, but it would stamp the redundant op's epoch.
+ assert_eq!(
+ shard.shards_table().epoch_for(ns),
+ Some(LIVE_EPOCH),
+ "the discarded op must not rewrite the routing row"
+ );
+ assert_eq!(
+ stats.current_offset(),
+ COMMITTED_OFFSET,
+ "a discarded build must not reset the live incarnation's
current_offset"
);
- for partition_id in 0..2 {
- let ns = IggyNamespace::new(0, 0, partition_id);
- assert!(
- partitions.contains(&ns),
- "namespace {ns:?} must be addressable exactly once"
- );
- assert_eq!(
- shard.shards_table().shard_for(ns),
- Some(0),
- "shards_table must point at the owning shard"
- );
- }
}
/// Multi-shard scenario: only the partition whose hash maps to
diff --git a/core/shard/src/lib.rs b/core/shard/src/lib.rs
index 28537d1c5..9a0ca3fc9 100644
--- a/core/shard/src/lib.rs
+++ b/core/shard/src/lib.rs
@@ -1775,6 +1775,34 @@ where
let _ =
sender.try_send(ShardFrame::lifecycle(LifecycleFrame::ReconcileApply));
}
+ /// `true` when an `InsertOwned` for `namespace` is built and queued but
not
+ /// yet applied.
+ ///
+ /// The reconciler's own "already handled" test is
`IggyPartitions::contains`,
+ /// which only turns true once the pump applies, so without this a pass run
+ /// during that lag rebuilds a namespace an earlier pass already built. The
+ /// queue IS the record of that in-flight work, so asking it cannot drift
+ /// from reality the way a parallel set would: every op leaves the queue
+ /// through `apply_reconcile_ops`, which either inserts or discards.
+ ///
+ /// Deliberately blind to `epoch`. Matching it would let a delete +
recreate
+ /// landing inside the lag build a second incarnation over the queued one's
+ /// on-disk path, which is the case this exists to prevent; the recreate is
+ /// not lost, it costs one pass. The queued (dead-epoch) op applies, and
the
+ /// next pass reads the epoch mismatch off the routing row and takes the
+ /// stale-incarnation teardown into a clean rebuild.
+ pub fn has_staged_insert_owned(&self, namespace: IggyNamespace) -> bool {
+ self.reconcile_queue.borrow().iter().any(|op| {
+ matches!(
+ op,
+ ReconcileOp::InsertOwned {
+ namespace: staged_namespace,
+ ..
+ } if *staged_namespace == namespace
+ )
+ })
+ }
+
/// Stage a segment-cleaner pass for `namespace` on this shard's pump. The
/// timer task resolves retention config off-pump and stamps `now`; the
pump
/// is the single writer of partition state, so the deletion runs there,
@@ -1882,18 +1910,30 @@ where
epoch,
} => {
// Idempotent apply, mirroring `ConfirmRemove` (idempotent
- // via `remove`'s `None` early-return). The reconciler
- // stages this from a task separate from the pump, so under
- // a commit burst two passes can each observe
- // `!contains(ns)` and build the same namespace before
- // either drains here. A second unconditional `insert`
- // would push a duplicate partition and overwrite the
- // `ns -> idx` entry, orphaning the first (its VSR group +
- // segment writers leak and `len` inflates). The discarded
- // build is a fresh empty incarnation over the same on-disk
- // path the kept one owns, so dropping it just closes a few
- // fds.
+ // via `remove`'s `None` early-return). An unconditional
+ // `insert` over a live namespace would push a duplicate
+ // partition and overwrite the `ns -> idx` entry, orphaning
+ // the first: its VSR group + segment writers leak and
`len`
+ // inflates.
+ //
+ // A backstop, not the mechanism. `reconcile_additions`
+ // skips a namespace whose `InsertOwned` is already staged
+ // ([`Self::has_staged_insert_owned`]), so a second op for
a
+ // live namespace should not be built at all. Dropping one
+ // here is damage control rather than a free no-op: the
+ // build already planted its initial segment over the live
+ // incarnation's path and folded that into the namespace's
+ // shared stats.
if partitions.contains(&namespace) {
+ tracing::error!(
+ shard = self_shard_id,
+ ns_raw = namespace.inner(),
+ epoch,
+ "discarding duplicate InsertOwned for a live
namespace: the \
+ staged-op guard was bypassed and the build
re-planted segment 0 \
+ over the live incarnation's path"
+ );
+
self.metrics.record_duplicate_partition_build_discarded();
drop(partition);
continue;
}
diff --git a/core/shard/src/metrics.rs b/core/shard/src/metrics.rs
index 659550217..2df15c357 100644
--- a/core/shard/src/metrics.rs
+++ b/core/shard/src/metrics.rs
@@ -186,6 +186,7 @@ pub struct ShardMetrics {
partitions_materialised_total: Counter,
partitions_removed_total: Counter,
partitions_reconcile_failures_total: Counter,
+ partitions_duplicate_builds_discarded_total: Counter,
partition_transfer_refusals_total: Counter,
partition_frames_rejected_stale_total: Counter,
partition_frames_rejected_ahead_total: Counter,
@@ -220,6 +221,7 @@ impl ShardMetrics {
partitions_materialised_total: Counter::default(),
partitions_removed_total: Counter::default(),
partitions_reconcile_failures_total: Counter::default(),
+ partitions_duplicate_builds_discarded_total: Counter::default(),
partition_transfer_refusals_total: Counter::default(),
partition_frames_rejected_stale_total: Counter::default(),
partition_frames_rejected_ahead_total: Counter::default(),
@@ -259,6 +261,16 @@ impl ShardMetrics {
self.partitions_removed_total.inc();
}
+ /// Bumped when the pump discards a duplicate `InsertOwned` for a namespace
+ /// that is already live. The reconciler's staged-op guard should make this
+ /// unreachable, so a non-zero value is a caught correctness anomaly, not
+ /// routine churn: the discarded build re-planted segment 0 over the live
+ /// incarnation's path and folded its initial segment into the shared stats
+ /// before the pump caught it.
+ pub fn record_duplicate_partition_build_discarded(&self) {
+ self.partitions_duplicate_builds_discarded_total.inc();
+ }
+
/// Bumped each time `build_partition_fresh` or
/// `delete_partitions_from_disk` returns `Err`. The reconciler retries
/// next tick, but a sustained climb surfaces a stuck partition (disk