hubcio commented on code in PR #4046:
URL: https://github.com/apache/iggy/pull/4046#discussion_r3952415105


##########
core/metadata/src/stm/stream.rs:
##########
@@ -603,7 +719,82 @@ impl StatsRegistry {
         self.partitions
             .lock()
             .expect("stats registry mutex poisoned")
-            .retain(|key, _| live_partitions.contains(key));
+            .retain(|key, entry| {

Review Comment:
   critical: this eviction doesn't zero the entry's stats, and the 
still-mounted stale incarnation holds the same `Arc` - its later 
`ConfirmRemove` subtracts those bytes again out of live siblings. zero the 
dropped entries here, like `remove_partitions` does.



##########
core/common/src/types/streaming_stats.rs:
##########
@@ -19,6 +19,110 @@ use std::sync::{
     Arc,
     atomic::{AtomicU32, AtomicU64, Ordering},
 };
+use tracing::warn;
+
+/// Number of rollup decrements that could not be covered by the counter they
+/// were subtracted from.
+///
+/// A decrement bigger than the total it targets means the tree lost
+/// `parent >= sum(children)` somewhere upstream. The counters are unsigned, so
+/// without a clamp that subtraction wraps to ~1.8e19 and every reader
+/// (`get_stream`, `get_topic`, `/stats`, `/metrics`) serves the wrapped value
+/// until the process restarts.
+///
+/// Clamping trades that for a total that is merely low, and low is where it
+/// stays: increments are `fetch_add`, so every later write stacks on the
+/// clamped base and carries the shortfall with it. Only a rebuild
+/// (`rebuild_parent_totals`) or a restart puts the total back. This counter
+/// plus the `warn!` below is what tells an operator the divergence happened,
+/// and that the aggregate needs a rebuild rather than time.
+static ROLLUP_UNDERFLOWS: AtomicU64 = AtomicU64::new(0);
+
+/// Monotonic count of clamped rollup decrements. Surfaced on `/metrics` so the
+/// clamp is alertable rather than log-grep-able.
+///
+/// Hidden from the docs: `iggy_common` re-exports this module at its root, and
+/// a server-only diagnostic has no business in the published SDK surface.
+#[doc(hidden)]
+#[must_use]
+pub fn rollup_underflows() -> u64 {

Review Comment:
   warning: `doc(hidden)` hides this from rustdoc but not from the crate's 
public API, and the glob in `lib.rs` puts it at the `iggy_common` root with 
nothing in CI to catch it. rename to `stats_rollup_underflows()`, or reach it 
through a named `pub mod`.



##########
core/shard/src/lib.rs:
##########
@@ -2537,6 +2537,17 @@ where
                         // the floor with the partition value.
                         self.metrics
                             
.record_partition_prepare_gap_drops(partition.take_prepare_gap_drops());
+                        // Roll whatever this partition still counts out of its
+                        // parent topic and stream. Here and not in the
+                        // reconciler's teardown: a handler suspended 
mid-append
+                        // holds its own `Arc` past the tombstone, and an 
earlier
+                        // settle leaves its increment in the parents with the
+                        // partition already gone. This is the drop point, so
+                        // nothing can add through that handle afterwards. The
+                        // rollback clamps, so the usual case -- the metadata
+                        // apply already zeroed these counters at commit -- 
takes
+                        // nothing.
+                        partition.stats.zero_out_all();

Review Comment:
   nit: this drop point rolls its counters back, but the two discard arms above 
at 2469 and 2493 just `drop(partition)` and leave what the build already folded 
into the shared stats. both are unreachable today - fold at adoption in 
`IggyPartitions::insert` instead.



##########
core/server/src/http/metrics.rs:
##########
@@ -203,6 +245,20 @@ mod tests {
         );
     }
 
+    #[test]
+    fn rollup_underflow_total_tracks_the_sampled_count() {
+        let metrics = HttpMetrics::init(&[]);
+        metrics.observe_rollup_underflows(3);
+        metrics.observe_rollup_underflows(5);
+        // A restart resets the source; the counter must not run backwards.

Review Comment:
   nit: a restart can't produce this state - the static and `HttpMetrics` are 
recreated together, so `recorded` never exceeds `total`. say the call guards a 
non-monotonic source defensively, or drop it.



##########
core/metadata/src/stm/stream.rs:
##########
@@ -439,22 +474,30 @@ impl StatsRegistry {
     /// reads the same `Arc`, so partition-plane counters are visible
     /// cross-shard without a gather.
     ///
+    /// Takes the committed record rather than a bare id, because a fresh entry

Review Comment:
   nit: the doc says committed record, but not that it has to be the record for 
this `(stream_id, topic_id)` key, or that `parent` must be the committed 
topic's `Arc`. neither is checked, and a fresh entry inherits both.



##########
core/common/src/types/streaming_stats.rs:
##########
@@ -19,6 +19,110 @@ use std::sync::{
     Arc,
     atomic::{AtomicU32, AtomicU64, Ordering},
 };
+use tracing::warn;
+
+/// Number of rollup decrements that could not be covered by the counter they
+/// were subtracted from.
+///
+/// A decrement bigger than the total it targets means the tree lost
+/// `parent >= sum(children)` somewhere upstream. The counters are unsigned, so
+/// without a clamp that subtraction wraps to ~1.8e19 and every reader
+/// (`get_stream`, `get_topic`, `/stats`, `/metrics`) serves the wrapped value
+/// until the process restarts.
+///
+/// Clamping trades that for a total that is merely low, and low is where it
+/// stays: increments are `fetch_add`, so every later write stacks on the
+/// clamped base and carries the shortfall with it. Only a rebuild
+/// (`rebuild_parent_totals`) or a restart puts the total back. This counter
+/// plus the `warn!` below is what tells an operator the divergence happened,
+/// and that the aggregate needs a rebuild rather than time.
+static ROLLUP_UNDERFLOWS: AtomicU64 = AtomicU64::new(0);
+
+/// Monotonic count of clamped rollup decrements. Surfaced on `/metrics` so the
+/// clamp is alertable rather than log-grep-able.
+///
+/// Hidden from the docs: `iggy_common` re-exports this module at its root, and
+/// a server-only diagnostic has no business in the published SDK surface.
+#[doc(hidden)]
+#[must_use]
+pub fn rollup_underflows() -> u64 {
+    ROLLUP_UNDERFLOWS.load(Ordering::Relaxed)
+}
+
+/// One scope/counter pair's clamp state: the labels an operator reads, plus 
the
+/// count the log throttle keys on.
+struct UnderflowSite {
+    scope: &'static str,
+    counter: &'static str,
+    clamped: AtomicU64,
+}
+
+impl UnderflowSite {
+    const fn new(scope: &'static str, counter: &'static str) -> Self {
+        Self {
+            scope,
+            counter,
+            clamped: AtomicU64::new(0),
+        }
+    }
+
+    /// Count the clamp globally, then log this pair's first and every power of
+    /// two after it.
+    ///
+    /// The throttle is per pair, not global: a skewed tree emits up to three
+    /// lines per counter per rollback and a bulk delete turns that into a
+    /// flood, so a shared count would hold a quiet pair's first-ever
+    /// divergence back for up to 1023 events behind a noisy one.
+    fn report(&self, shortfall: u64) {
+        ROLLUP_UNDERFLOWS.fetch_add(1, Ordering::Relaxed);
+        let clamped = self.clamped.fetch_add(1, Ordering::Relaxed) + 1;

Review Comment:
   nit: `clamped` is the per-pair count but the metric it bumps is 
process-global and carries no scope or counter label, so the two numbers can't 
be reconciled. log both, or label the metric.



##########
core/common/src/types/streaming_stats.rs:
##########
@@ -282,21 +399,39 @@ impl PartitionStats {
         self.parent.increment_segments_count(segments_count);
     }
 
-    pub fn decrement_size_bytes(&self, size_bytes: u64) {
-        self.size_bytes.fetch_sub(size_bytes, Ordering::AcqRel);
-        self.decrement_parent_size_bytes(size_bytes);
-    }
-
-    pub fn decrement_messages_count(&self, messages_count: u64) {
-        self.messages_count
-            .fetch_sub(messages_count, Ordering::AcqRel);
-        self.decrement_parent_messages_count(messages_count);
-    }
-
-    pub fn decrement_segments_count(&self, segments_count: u32) {
-        self.segments_count
-            .fetch_sub(segments_count, Ordering::AcqRel);
-        self.decrement_parent_segments_count(segments_count);
+    // Forward what this level actually gave up, not what was asked for. A
+    // decrement bigger than the counter holds means the bytes were never in
+    // this subtree, so the ancestors do not hold them either -- passing the
+    // full amount up would take them out of a sibling's live data instead.
+    // Under `parent == sum(children)` the two are equal and nothing changes.
+    //
+    // Each returns this level's shortfall -- how much of the amount it did not
+    // hold, zero when the decrement was covered -- so a caller that knows the
+    // ids can name where the divergence is.
+    pub fn decrement_size_bytes(&self, size_bytes: u64) -> u64 {
+        let taken = clamped_sub_u64(&self.size_bytes, size_bytes, 
&PARTITION_SIZE_BYTES);
+        self.decrement_parent_size_bytes(taken);
+        size_bytes - taken
+    }
+
+    pub fn decrement_messages_count(&self, messages_count: u64) -> u64 {
+        let taken = clamped_sub_u64(
+            &self.messages_count,
+            messages_count,
+            &PARTITION_MESSAGES_COUNT,
+        );
+        self.decrement_parent_messages_count(taken);
+        messages_count - taken
+    }
+
+    pub fn decrement_segments_count(&self, segments_count: u32) -> u32 {
+        let taken = clamped_sub_u32(
+            &self.segments_count,
+            segments_count,
+            &PARTITION_SEGMENTS_COUNT,
+        );
+        self.decrement_parent_segments_count(taken);
+        segments_count - taken
     }
 
     pub fn decrement_parent_size_bytes(&self, size_bytes: u64) {

Review Comment:
   simplification: these twelve `increment_parent_*` and `decrement_parent_*` 
wrappers have no caller outside this file, and they're public on a published 
crate. inline them - `self.parent.increment_size_bytes(x)` is the same one line.



##########
core/integration/tests/server/scenarios/mod.rs:
##########
@@ -72,6 +74,133 @@ const USERNAME_3: &str = "user3";
 const CONSUMER_KIND: ConsumerKind = ConsumerKind::Consumer;
 const MESSAGES_COUNT: u32 = 1337;
 
+// The partition plane applies committed ops asynchronously on the owning shard
+// (sends fold into the shared stats at commit-apply; purge/delete zero them
+// when the reconciler drives the wipe), so a read racing that window can see a
+// pre-apply value. Retry until the expectation holds, then make the terminal
+// assertion for a real mismatch.
+const STATS_CONVERGENCE_TIMEOUT: Duration = Duration::from_secs(10);

Review Comment:
   simplification: same 10s and 100ms as `POLL_CONVERGENCE_TIMEOUT` and 
`POLL_RETRY_INTERVAL` at 65, and both wait out the same convergence. keep one 
neutral pair and update the doc links at 114 and 205.



##########
core/server/src/partition_reconciler.rs:
##########
@@ -694,11 +707,12 @@ async fn reconcile_additions(
             continue;
         }
 
-        // Resolve the shared stats `Arc` only for namespaces actually
-        // built, not once per committed partition every pass. A topic that
-        // vanished between the target snapshot and this read defers to the
-        // next pass.
-        let Some((partition_stats, topic_runtime)) = 
fetch_partition_stats(ctx, ns) else {
+        // Resolved only for namespaces actually built, not once per committed
+        // partition every pass. A stream, topic or partition that vanished
+        // between the target snapshot and this read defers to the next pass.
+        let Some((partition_stats, topic_runtime, partition_metadata)) =
+            fetch_partition_build_inputs(ctx, ns)

Review Comment:
   warning: `created_revision` is in hand here but the build and row still get 
the snapshot `epoch`, so a mid-pass recreate lands a dead-epoch row. skip on 
mismatch only when the dir is absent - a bare `continue` lets the next pass 
adopt the old segments.



##########
core/metadata/src/stm/stream.rs:
##########
@@ -603,7 +719,82 @@ impl StatsRegistry {
         self.partitions
             .lock()
             .expect("stats registry mutex poisoned")
-            .retain(|key, _| live_partitions.contains(key));
+            .retain(|key, entry| {
+                live_partitions
+                    .get(key)
+                    .is_some_and(|created_revision| *created_revision == 
entry.created_revision)
+            });
+    }
+
+    /// Recompute every topic and stream total as the sum of the partition
+    /// entries this node holds.
+    ///
+    /// A snapshot carries the DONOR's topic and stream totals, but partition
+    /// counters are node-local and never snapshotted: they keep moving on the
+    /// receiver while the snapshot is captured, shipped and installed. Storing
+    /// the donor's totals over the receiver's own children leaves
+    /// `topic < sum(partitions)` as the ordinary post-transfer shape, and the
+    /// next delete then subtracts a child from a parent that never counted it.
+    /// Every rollback in this file assumes `parent == sum(children)`; this is
+    /// where that gets re-established.
+    ///
+    /// Partitions this node has not materialized yet contribute nothing, the
+    /// same convergence boot relies on: the reconciler registers each one and
+    /// folds its on-disk delta in as it goes.
+    ///
+    /// The `partitions` guard fences the MAP, not the counters, which the data
+    /// plane reaches through the `Arc`. An append landing between a child read
+    /// and its parent store is folded into the parent by `fetch_add` and then
+    /// overwritten, so the invariant is restored modulo whatever arrives 
during
+    /// the walk. That residue is bounded by the walk and by one append, and 
the
+    /// saturating rollback absorbs it; quiescing the data plane for a metadata
+    /// install would cost far more than it buys.
+    ///
+    /// # Panics
+    /// If the registry mutex is poisoned.
+    // One acquisition for the whole walk: the stores it makes are plain atomic
+    // writes on the topic and stream `Arc`s, with no parent cascade and no way
+    // back into the map, so there is nothing for a tighter guard to protect.
+    #[allow(clippy::significant_drop_tightening)]
+    fn rebuild_parent_totals(&self, streams: &IdSlab<Stream>) {

Review Comment:
   warning: this holds the `partitions` mutex across a walk of every partition 
on the node, and other threads take that same lock through the get-or-create. 
collect under the guard, drop it, then store.



##########
core/server/src/partition_reconciler.rs:
##########
@@ -1244,47 +1261,69 @@ fn current_revision(ctx: &ReconcilerCtx) -> u64 {
         .read(|inner| inner.revision)
 }
 
-/// Clone the parent topic's `Arc<TopicStats>` for a single namespace.
-/// `None` if the topic vanished between the target snapshot and this read.
-fn fetch_partition_stats(
+/// Roll a torn-down partition out of its parent topic and stream by dropping
+/// its registry entry.
+///
+/// Usually a no-op: the metadata apply evicts the entry on the commit that
+/// acked the delete. What is left for this call is a stale incarnation being
+/// torn down after a slab-key reuse, whose entry the apply never named. It 
must
+/// not survive into the rebuild -- `StatsRegistry::partition` is a
+/// get-or-create, so the rebuild would inherit the dead incarnation's 
counters.
+/// (Its `purged_generation` no longer rides along either way: a fresh entry
+/// seeds that gate from the committed partition.)
+///
+/// The mounted partition's own handle is settled later, on `ConfirmRemove`.
+fn settle_partition_stats(ctx: &ReconcilerCtx, ns: IggyNamespace) {

Review Comment:
   warning: on a successful teardown this eviction leaves the partition 
reporting the registry-miss shape, and the parents low by its whole size, until 
the rebuild re-folds the disk delta a pass later. document the window, or carry 
a materialisation signal.



##########
core/server/src/partition_reconciler.rs:
##########
@@ -966,6 +977,12 @@ async fn tear_down_owned_partition(
         partitions.tombstone(ns);
     }
     shards_table.remove(&ns);
+    // Registry entry only. The mounted partition's OWN counters are settled by
+    // the `ConfirmRemove` arm, which is the point where the partition value is
+    // dropped: a handler suspended mid-append resumes and increments through
+    // its cached handle, and anything settled before that drop leaves the
+    // increment in the parent totals with nothing left to roll it back.
+    settle_partition_stats(ctx, ns);

Review Comment:
   warning: settle runs before the disk delete, so a failed delete returns with 
no `ConfirmRemove` and the still-mounted partition then reports one empty 
segment at offset 0 while its files remain. move it down next to the enqueue.



##########
core/common/src/types/streaming_stats.rs:
##########
@@ -19,6 +19,110 @@ use std::sync::{
     Arc,
     atomic::{AtomicU32, AtomicU64, Ordering},
 };
+use tracing::warn;
+
+/// Number of rollup decrements that could not be covered by the counter they
+/// were subtracted from.
+///
+/// A decrement bigger than the total it targets means the tree lost
+/// `parent >= sum(children)` somewhere upstream. The counters are unsigned, so
+/// without a clamp that subtraction wraps to ~1.8e19 and every reader
+/// (`get_stream`, `get_topic`, `/stats`, `/metrics`) serves the wrapped value
+/// until the process restarts.
+///
+/// Clamping trades that for a total that is merely low, and low is where it
+/// stays: increments are `fetch_add`, so every later write stacks on the
+/// clamped base and carries the shortfall with it. Only a rebuild
+/// (`rebuild_parent_totals`) or a restart puts the total back. This counter
+/// plus the `warn!` below is what tells an operator the divergence happened,
+/// and that the aggregate needs a rebuild rather than time.
+static ROLLUP_UNDERFLOWS: AtomicU64 = AtomicU64::new(0);

Review Comment:
   nit: this static is process-global, which the docs never say here. worth 
stating, since the throttle right below is deliberately split per pair to stop 
a noisy counter hiding a quiet one's first divergence.



##########
core/metadata/src/stm/stream.rs:
##########
@@ -384,10 +389,32 @@ impl Stream {
 /// read.
 ///
 /// Shared across buffers and reader shards via `Arc` (a `StreamsInner` clone
-/// shares it). Only shard 0's writer mutates the maps, under the
-/// single-threaded Absorb; the `Mutex` is for `Sync` (uncontended), not
-/// concurrency. Ids are deterministic across replicas (same op order), so both
-/// buffers resolve the same key.
+/// shares it). Shard 0's writer owns the map mutations that ride the op-log,
+/// but it is NOT the only writer: [`Self::partition`] is a get-or-create the
+/// owning shard's reconciler calls when it materializes a namespace, and boot
+/// recovery calls it from every shard at once. The `Mutex` is load-bearing for
+/// that concurrency, not merely for `Sync`. Ids are deterministic across
+/// replicas (same op order), so both buffers resolve the same key.
+///
+/// Every eviction here therefore runs twice, once per buffer, and the second
+/// run is deferred. Two properties keep that from doing damage, and neither is
+/// guaranteed by the `Absorb` trait:
+///
+/// * `WriteCell::apply` is `append(cmd).publish()`, so a batch is one op. 
Batch
+///   two and a delete could be absorbed on the second buffer AFTER a create
+///   that re-minted the same id, and the delete would then evict the new
+///   partition's entry.
+/// * `left_right` drains the previous batch's `absorb_second` before the

Review Comment:
   nit: the keyed evictions lean on this drain order, and `left-right = "0.11"` 
in the manifest floats any 0.11.x. pin `=0.11.8`, or add an `Absorb` impl that 
records the call sequence.



##########
core/partitions/src/iggy_partition.rs:
##########
@@ -7263,6 +7260,39 @@ pub struct SegmentRemoval {
     pub budget_spent: bool,
 }
 
+/// Roll one cleaned-up segment out of the partition counters, naming the
+/// partition when the rollback could not be covered.
+///
+/// Retention is the likeliest source of a clamped rollback: a partition on its
+/// way out keeps serving cleanup passes after the delete already settled its
+/// counters into the parents. The `stats_rollup_underflows` metric says a
+/// divergence happened and carries no ids; this says which partition and which
+/// segment, which is what an operator needs to act on it.
+fn settle_cleaned_segment(
+    stats: &PartitionStats,
+    namespace: IggyNamespace,
+    segment: &Segment,
+    messages_count: u64,
+) {
+    let size_shortfall = 
stats.decrement_size_bytes(segment.size.as_bytes_u64());
+    let segments_shortfall = stats.decrement_segments_count(1);
+    let messages_shortfall = stats.decrement_messages_count(messages_count);
+    if size_shortfall == 0 && segments_shortfall == 0 && messages_shortfall == 
0 {
+        return;
+    }
+    warn!(

Review Comment:
   nit: this warn fires per segment with no throttle, and 
`UnderflowSite::report` already logs the same event power-of-two throttled plus 
counts it. emit one line per `clean_*` call with the removed offset range.



##########
core/server/src/http/metrics.rs:
##########
@@ -45,6 +45,30 @@ pub(in crate::http) struct HttpMetrics {
     pub(in crate::http) messages: Gauge,
     pub(in crate::http) users: Gauge,
     pub(in crate::http) clients: Gauge,
+    /// Not a legacy-parity metric: the running count of rollup decrements that
+    /// had to be clamped at zero. Non-zero means a topic or stream total was
+    /// asked to give back more than it held, so the aggregate it now reports 
is
+    /// low, and stays low until a rebuild or a restart.
+    ///
+    /// Not "alert on any increase". A delete, a purge, a partition teardown 
and
+    /// a snapshot restore each open a window where a retention pass hands back
+    /// bytes the parents have already given up, and the clamp is the intended
+    /// outcome there -- the 
`given_a_late_decrement_on_a_rolled_back_partition`

Review Comment:
   nit: the test is 
`given_a_late_decrement_on_a_rolled_back_partition_should_leave_siblings_alone`,
 and `iggy_common::streaming_stats` isn't a reachable path since the module is 
`pub(crate)`. cite the full name plus the file.



##########
core/server/src/http/metrics.rs:
##########
@@ -45,6 +45,30 @@ pub(in crate::http) struct HttpMetrics {
     pub(in crate::http) messages: Gauge,
     pub(in crate::http) users: Gauge,
     pub(in crate::http) clients: Gauge,
+    /// Not a legacy-parity metric: the running count of rollup decrements that
+    /// had to be clamped at zero. Non-zero means a topic or stream total was

Review Comment:
   nit: says "a topic or stream total", but the three partition sites feed the 
same counter - including the example the doc cites, where the partition clamps 
and no parent does. say "partition, topic or stream".



##########
core/integration/tests/data_integrity/verify_after_server_restart.rs:
##########
@@ -520,3 +557,17 @@ async fn should_handle_resource_deletion_and_restart() {
         );
     }
 }
+
+/// Payload for `should_handle_resource_deletion_and_restart`: small, fixed and
+/// non-empty, so a scope that gets deleted has bytes worth rolling back.
+fn deletion_test_messages() -> Vec<IggyMessage> {
+    (0..8u128)
+        .map(|offset| {
+            IggyMessage::builder()
+                .id(offset + 1)
+                .payload(bytes::Bytes::from_static(b"deletion-stats-payload"))

Review Comment:
   nit: `Bytes` isn't in scope anywhere in this file and the rest of it imports 
at the top. add `use bytes::Bytes;` - this PR added the same import to 
`scenarios/mod.rs`.



##########
core/common/src/types/streaming_stats.rs:
##########
@@ -19,6 +19,110 @@ use std::sync::{
     Arc,
     atomic::{AtomicU32, AtomicU64, Ordering},
 };
+use tracing::warn;
+
+/// Number of rollup decrements that could not be covered by the counter they
+/// were subtracted from.
+///
+/// A decrement bigger than the total it targets means the tree lost
+/// `parent >= sum(children)` somewhere upstream. The counters are unsigned, so
+/// without a clamp that subtraction wraps to ~1.8e19 and every reader
+/// (`get_stream`, `get_topic`, `/stats`, `/metrics`) serves the wrapped value
+/// until the process restarts.
+///
+/// Clamping trades that for a total that is merely low, and low is where it
+/// stays: increments are `fetch_add`, so every later write stacks on the
+/// clamped base and carries the shortfall with it. Only a rebuild
+/// (`rebuild_parent_totals`) or a restart puts the total back. This counter
+/// plus the `warn!` below is what tells an operator the divergence happened,
+/// and that the aggregate needs a rebuild rather than time.
+static ROLLUP_UNDERFLOWS: AtomicU64 = AtomicU64::new(0);
+
+/// Monotonic count of clamped rollup decrements. Surfaced on `/metrics` so the
+/// clamp is alertable rather than log-grep-able.
+///
+/// Hidden from the docs: `iggy_common` re-exports this module at its root, and
+/// a server-only diagnostic has no business in the published SDK surface.
+#[doc(hidden)]
+#[must_use]
+pub fn rollup_underflows() -> u64 {
+    ROLLUP_UNDERFLOWS.load(Ordering::Relaxed)
+}
+
+/// One scope/counter pair's clamp state: the labels an operator reads, plus 
the
+/// count the log throttle keys on.
+struct UnderflowSite {
+    scope: &'static str,
+    counter: &'static str,
+    clamped: AtomicU64,
+}
+
+impl UnderflowSite {
+    const fn new(scope: &'static str, counter: &'static str) -> Self {
+        Self {
+            scope,
+            counter,
+            clamped: AtomicU64::new(0),
+        }
+    }
+
+    /// Count the clamp globally, then log this pair's first and every power of
+    /// two after it.
+    ///
+    /// The throttle is per pair, not global: a skewed tree emits up to three
+    /// lines per counter per rollback and a bulk delete turns that into a
+    /// flood, so a shared count would hold a quiet pair's first-ever
+    /// divergence back for up to 1023 events behind a noisy one.
+    fn report(&self, shortfall: u64) {
+        ROLLUP_UNDERFLOWS.fetch_add(1, Ordering::Relaxed);
+        let clamped = self.clamped.fetch_add(1, Ordering::Relaxed) + 1;
+        if clamped.is_power_of_two() {
+            warn!(

Review Comment:
   nit: no `target:` here, so this can't be filtered with the other diag 
streams. a scope-neutral target like `iggy.stats.diag` would do - the site logs 
stream and topic scopes too, so the partitions target is wrong.



##########
core/common/src/types/streaming_stats.rs:
##########
@@ -357,3 +492,115 @@ impl PartitionStats {
         self.zero_out_current_offset();
     }
 }
+
+#[cfg(test)]
+mod tests {
+    use super::*;
+
+    fn tree() -> (Arc<StreamStats>, Arc<TopicStats>, Arc<PartitionStats>) {
+        let stream = Arc::new(StreamStats::default());
+        let topic = Arc::new(TopicStats::new(stream.clone()));
+        let partition = Arc::new(PartitionStats::new(topic.clone()));
+        (stream, topic, partition)
+    }
+
+    #[test]
+    fn 
given_a_parent_short_of_its_child_when_rolling_back_should_clamp_instead_of_wrapping()
 {
+        let (stream, topic, partition) = tree();
+        partition.increment_size_bytes(512);
+        // A snapshot restore stores absolute parent totals, so a parent can 
end
+        // up holding less than its children do.
+        topic.store_from_snapshot(0, 0, 0);
+        stream.store_from_snapshot(0, 0, 0);
+
+        partition.zero_out_all();
+
+        assert_eq!(topic.size_bytes_inconsistent(), 0);
+        assert_eq!(stream.size_bytes_inconsistent(), 0);
+    }
+
+    /// A deleted partition keeps its handle until the reconciler tears it 
down,
+    /// so a retention sweep can decrement counters the delete already rolled
+    /// back. Those bytes are not in the parents any more, and taking them 
again
+    /// would take a live sibling's instead.
+    #[test]
+    fn 
given_a_late_decrement_on_a_rolled_back_partition_should_leave_siblings_alone() 
{
+        let stream = Arc::new(StreamStats::default());
+        let topic = Arc::new(TopicStats::new(stream.clone()));
+        let survivor = Arc::new(PartitionStats::new(topic.clone()));
+        let deleted = Arc::new(PartitionStats::new(topic.clone()));
+        survivor.increment_size_bytes(488);
+        deleted.increment_size_bytes(512);
+
+        deleted.zero_out_all();
+        assert_eq!(topic.size_bytes_inconsistent(), 488);
+
+        // Retention retiring a segment of the partition that is on its way 
out.
+        deleted.decrement_size_bytes(100);
+
+        assert_eq!(
+            topic.size_bytes_inconsistent(),
+            488,
+            "the survivor's bytes are the only ones left; they are not the 
deleted partition's"
+        );
+        assert_eq!(stream.size_bytes_inconsistent(), 488);
+        assert_eq!(survivor.size_bytes_inconsistent(), 488);
+    }
+
+    /// `delete_topic` settles the topic's residue on the stream, and the
+    /// reconciler later settles the same partition's own handle. The second
+    /// pass must find nothing left to take, or it takes it from another topic.
+    #[test]
+    fn 
given_a_topic_already_settled_when_its_partition_settles_should_leave_other_topics_alone()
 {
+        let stream = Arc::new(StreamStats::default());
+        let doomed_topic = Arc::new(TopicStats::new(stream.clone()));
+        let live_topic = Arc::new(TopicStats::new(stream.clone()));
+        let doomed_partition = 
Arc::new(PartitionStats::new(doomed_topic.clone()));
+        let live_partition = Arc::new(PartitionStats::new(live_topic.clone()));
+        live_partition.increment_size_bytes(900);
+        // Landed through the cached handle after the delete evicted its entry.
+        doomed_partition.increment_size_bytes(100);
+        assert_eq!(stream.size_bytes_inconsistent(), 1000);
+
+        doomed_topic.zero_out_all();
+        assert_eq!(stream.size_bytes_inconsistent(), 900);
+
+        doomed_partition.zero_out_all();

Review Comment:
   nit: both widths expand from the same `clamped_sub!` body using 
`saturating_sub`, so nothing modulus-specific runs - say it covers the u32 
wiring instead. same wording at `stream.rs:3480`.



##########
core/integration/tests/server/general.rs:
##########
@@ -88,6 +88,17 @@ async fn stream_size_validation(harness: &TestHarness) {
     stream_size_validation_scenario::run(harness).await;
 }
 
+#[iggy_harness(

Review Comment:
   simplification: the scenario's own doc calls it transport-irrelevant and 
says it guards neither half, and every assertion goes through `get_topic`, 
`get_stream` or `get_stats`, all already covered on four transports. `[Tcp]` 
drops 9 server processes.



##########
core/server/src/partition_reconciler.rs:
##########
@@ -1926,6 +1965,148 @@ mod tests {
         );
     }
 
+    /// A pass captures its targets once and then awaits disk work per
+    /// namespace, so a delete committed mid-pass leaves a stale target behind.
+    /// Resolving stats for it would get-or-CREATE the registry entry the 
delete
+    /// just evicted, and the apply on the second left-right buffer would then
+    /// zero and drop a live partition's counters.
+    #[compio::test]
+    async fn 
stats_are_refused_for_a_partition_the_committed_topic_no_longer_lists() {
+        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));
+
+        let listed = IggyNamespace::new(0, 0, 0);
+        assert!(
+            fetch_partition_build_inputs(&ctx, listed).is_some(),
+            "a partition the topic lists must still resolve"
+        );
+
+        let unlisted = IggyNamespace::new(0, 0, 7);
+        assert!(
+            fetch_partition_build_inputs(&ctx, unlisted).is_none(),
+            "a partition the committed topic does not list must not resolve"
+        );
+        assert!(
+            registry(&ctx).partition_get(0, 0, 7).is_none(),
+            "and the refusal must not have created its entry on the way"
+        );
+    }
+
+    /// The metadata apply rolls a deleted partition out of its parents at
+    /// commit, but the partition stays mounted until the reconciler tears it
+    /// down, and appends in that window land in parents with no entry left to
+    /// account for them. Teardown is where that residue gets settled.
+    #[compio::test]
+    async fn 
teardown_settles_the_residue_the_metadata_delete_could_not_reach() {
+        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));
+        let ns = IggyNamespace::new(0, 0, 0);
+
+        reconcile_once(&ctx).await;
+        ctx.shard.apply_reconcile_ops();
+        let (stats, ..) =
+            fetch_partition_build_inputs(&ctx, ns).expect("materialised 
namespace has stats");
+        stats.increment_size_bytes(512);
+
+        seed_delete_topic(&ctx.shard.plane.metadata().mux_stm, 3, 0, 0);
+        // The apply evicted the entry and rolled back what it could see. This
+        // is the append that beat the teardown, through the handle the mounted
+        // partition still holds.
+        stats.increment_size_bytes(100);
+        assert_eq!(stream_size(&ctx), 100);
+
+        reconcile_once(&ctx).await;
+        ctx.shard.apply_reconcile_ops();
+
+        assert_eq!(
+            stream_size(&ctx),
+            0,
+            "teardown must settle what landed after the commit that acked the 
delete"

Review Comment:
   nit: the metadata apply already evicted the entry, so this passes on 
`zero_out_all()` at the drop point alone - delete `settle_partition_stats` and 
it stays green. the sibling test is the one guarding the settle, so rename this 
for the drop point.



##########
core/common/src/types/streaming_stats.rs:
##########
@@ -357,3 +492,115 @@ impl PartitionStats {
         self.zero_out_current_offset();
     }
 }
+
+#[cfg(test)]
+mod tests {
+    use super::*;
+
+    fn tree() -> (Arc<StreamStats>, Arc<TopicStats>, Arc<PartitionStats>) {
+        let stream = Arc::new(StreamStats::default());
+        let topic = Arc::new(TopicStats::new(stream.clone()));
+        let partition = Arc::new(PartitionStats::new(topic.clone()));
+        (stream, topic, partition)
+    }
+
+    #[test]
+    fn 
given_a_parent_short_of_its_child_when_rolling_back_should_clamp_instead_of_wrapping()
 {
+        let (stream, topic, partition) = tree();
+        partition.increment_size_bytes(512);
+        // A snapshot restore stores absolute parent totals, so a parent can 
end
+        // up holding less than its children do.
+        topic.store_from_snapshot(0, 0, 0);
+        stream.store_from_snapshot(0, 0, 0);
+
+        partition.zero_out_all();
+
+        assert_eq!(topic.size_bytes_inconsistent(), 0);
+        assert_eq!(stream.size_bytes_inconsistent(), 0);
+    }
+
+    /// A deleted partition keeps its handle until the reconciler tears it 
down,
+    /// so a retention sweep can decrement counters the delete already rolled
+    /// back. Those bytes are not in the parents any more, and taking them 
again
+    /// would take a live sibling's instead.
+    #[test]
+    fn 
given_a_late_decrement_on_a_rolled_back_partition_should_leave_siblings_alone() 
{

Review Comment:
   nit: no `when_` in this name, while the other four new tests here have one.



##########
core/integration/tests/server/scenarios/delete_stats_rollback_scenario.rs:
##########
@@ -0,0 +1,221 @@
+// 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.
+
+//! Deleting data that is still counted must leave the parent totals correct.
+//!
+//! `stream_size_validation_scenario` covers delete too, but it purges each
+//! topic first, so every delete it performs removes an already-empty scope and
+//! cannot observe a rollback that never happened. These scenarios delete 
scopes
+//! that still hold messages, which is what the parent totals are wrong about.
+//!
+//! What this guards, and what it does not. Two mechanisms roll a deleted scope
+//! out of its parents: the metadata STM evicts the registry entries at commit
+//! and zeroes them into their parents, and the reconciler's partition teardown
+//! settles whatever landed after that. The STM half runs inside the apply that
+//! produces the ack, on counters shared across every shard and both left-right
+//! buffers, so it alone satisfies every assertion here. The reconciler runs as
+//! a separate task, but a delete commit wakes it (`signal_reconcile_wake`), so
+//! it too finishes before a client can complete another round trip.
+//!
+//! Measured, not reasoned: with only the STM rollback disabled this scenario
+//! PASSES (while the `metadata::stm::stream` unit tests fail), and with only
+//! the reconciler settle disabled it also passes. It fails only when BOTH are
+//! disabled. Either mechanism alone keeps a client's view correct, so this is 
a
+//! contract test for what a client observes and not a regression guard for
+//! either half.
+//!
+//! Each half is guarded by unit tests instead: the `metadata::stm::stream`
+//! module for the rollback and eviction, `server::partition_reconciler` for 
the
+//! teardown settle and the membership gate. The integration layer cannot hold
+//! the reconciler off by configuration -- the tick is capped at 30s, cannot be
+//! set to zero, and the commit wake bypasses it regardless.
+//!
+//! Verifying any of this needs `cargo build --bin iggy-server` first. The
+//! harness spawns the built binary, so a source-only mutation runs against the
+//! previous build and reports a green that means nothing.
+//!
+//! The reads after each delete are deliberately single-shot. Retrying until 
the
+//! numbers converge, the way the pre-delete assertions do for the async send
+//! folding, would also wait out a broken rollback and pass regardless.
+
+use crate::server::scenarios::{
+    PARTITIONS_COUNT, batch_size, create_client, create_messages, 
validate_stream,

Review Comment:
   nit: the scenario never checks server-wide totals go back to zero after the 
deletes, which is the clearest statement of what this fixes. 
`assert_clean_system` reads no stats - add `validate_system_stats(client, 0, 
0)` before it.



##########
core/server/src/http/metrics.rs:
##########
@@ -106,6 +137,17 @@ impl HttpMetrics {
         self.http_requests.clone()
     }
 
+    /// Advance the rollup-underflow counter to `total`, the process-wide count

Review Comment:
   simplification: the get, `checked_sub`, `inc_by` dance exists only because 
`Counter` has no setter, and it guards a state the doc calls impossible. a 
~12-line `Collector` with `ConstCounter` emits the same `_total` series - add a 
test for it, `PARITY_METRIC_NAMES` skips this one.



##########
core/partitions/src/iggy_partition.rs:
##########
@@ -7263,6 +7260,39 @@ pub struct SegmentRemoval {
     pub budget_spent: bool,
 }
 
+/// Roll one cleaned-up segment out of the partition counters, naming the
+/// partition when the rollback could not be covered.
+///
+/// Retention is the likeliest source of a clamped rollback: a partition on its
+/// way out keeps serving cleanup passes after the delete already settled its
+/// counters into the parents. The `stats_rollup_underflows` metric says a
+/// divergence happened and carries no ids; this says which partition and which
+/// segment, which is what an operator needs to act on it.
+fn settle_cleaned_segment(

Review Comment:
   nit: `messages_count` is derived from the same `&segment` one line above the 
call. derive it inside and return it, then feed `removal.messages` from the 
return.



##########
core/common/src/types/streaming_stats.rs:
##########
@@ -42,18 +146,21 @@ impl StreamStats {
             .fetch_add(segments_count, Ordering::AcqRel);
     }
 
-    pub fn decrement_size_bytes(&self, size_bytes: u64) {
-        self.size_bytes.fetch_sub(size_bytes, Ordering::AcqRel);
+    /// Returns the shortfall -- how much of `size_bytes` this counter did not
+    /// hold -- so a caller that knows the ids can name where the divergence 
is.
+    /// Zero whenever the decrement was covered.
+    pub fn decrement_size_bytes(&self, size_bytes: u64) -> u64 {
+        size_bytes - clamped_sub_u64(&self.size_bytes, size_bytes, 
&STREAM_SIZE_BYTES)
     }
 
-    pub fn decrement_messages_count(&self, messages_count: u64) {
-        self.messages_count
-            .fetch_sub(messages_count, Ordering::AcqRel);
+    pub fn decrement_messages_count(&self, messages_count: u64) -> u64 {

Review Comment:
   simplification: nothing reads the stream-level or topic-level shortfall 
anywhere - each level computes its own from `taken`. return `()` at both levels 
and keep the three on `PartitionStats`, which also cuts the public signature 
churn from nine to three.



##########
core/common/src/types/streaming_stats.rs:
##########
@@ -19,6 +19,110 @@ use std::sync::{
     Arc,
     atomic::{AtomicU32, AtomicU64, Ordering},
 };
+use tracing::warn;
+
+/// Number of rollup decrements that could not be covered by the counter they
+/// were subtracted from.
+///
+/// A decrement bigger than the total it targets means the tree lost
+/// `parent >= sum(children)` somewhere upstream. The counters are unsigned, so
+/// without a clamp that subtraction wraps to ~1.8e19 and every reader
+/// (`get_stream`, `get_topic`, `/stats`, `/metrics`) serves the wrapped value
+/// until the process restarts.
+///
+/// Clamping trades that for a total that is merely low, and low is where it
+/// stays: increments are `fetch_add`, so every later write stacks on the
+/// clamped base and carries the shortfall with it. Only a rebuild
+/// (`rebuild_parent_totals`) or a restart puts the total back. This counter
+/// plus the `warn!` below is what tells an operator the divergence happened,
+/// and that the aggregate needs a rebuild rather than time.
+static ROLLUP_UNDERFLOWS: AtomicU64 = AtomicU64::new(0);
+
+/// Monotonic count of clamped rollup decrements. Surfaced on `/metrics` so the
+/// clamp is alertable rather than log-grep-able.
+///
+/// Hidden from the docs: `iggy_common` re-exports this module at its root, and
+/// a server-only diagnostic has no business in the published SDK surface.
+#[doc(hidden)]
+#[must_use]
+pub fn rollup_underflows() -> u64 {
+    ROLLUP_UNDERFLOWS.load(Ordering::Relaxed)
+}
+
+/// One scope/counter pair's clamp state: the labels an operator reads, plus 
the
+/// count the log throttle keys on.
+struct UnderflowSite {
+    scope: &'static str,
+    counter: &'static str,
+    clamped: AtomicU64,
+}
+
+impl UnderflowSite {
+    const fn new(scope: &'static str, counter: &'static str) -> Self {
+        Self {
+            scope,
+            counter,
+            clamped: AtomicU64::new(0),
+        }
+    }
+
+    /// Count the clamp globally, then log this pair's first and every power of
+    /// two after it.
+    ///
+    /// The throttle is per pair, not global: a skewed tree emits up to three
+    /// lines per counter per rollback and a bulk delete turns that into a
+    /// flood, so a shared count would hold a quiet pair's first-ever
+    /// divergence back for up to 1023 events behind a noisy one.
+    fn report(&self, shortfall: u64) {
+        ROLLUP_UNDERFLOWS.fetch_add(1, Ordering::Relaxed);
+        let clamped = self.clamped.fetch_add(1, Ordering::Relaxed) + 1;
+        if clamped.is_power_of_two() {
+            warn!(
+                scope = self.scope,
+                counter = self.counter,
+                shortfall,
+                clamped,
+                "rollup decrement exceeded the total it was subtracted from; 
clamped at zero"
+            );
+        }
+    }
+}
+
+static STREAM_SIZE_BYTES: UnderflowSite = UnderflowSite::new("stream", 
"size_bytes");
+static STREAM_MESSAGES_COUNT: UnderflowSite = UnderflowSite::new("stream", 
"messages_count");
+static STREAM_SEGMENTS_COUNT: UnderflowSite = UnderflowSite::new("stream", 
"segments_count");

Review Comment:
   simplification: nine per-pair statics where three per-scope sites plus a 
`counter` argument would do, about 9 lines. costs you the property the doc 
right above defends - a noisy counter would then swallow a quiet one's first 
divergence.



##########
core/metadata/src/stm/stream.rs:
##########
@@ -2338,7 +2547,22 @@ impl Snapshotable for Streams {
         // Boot: no live registry exists yet, so mint one. Safe because
         // `new_from_empty` clones this single inner onto the other left-right
         // buffer rather than building a second one.
-        Ok(StreamsInner::inner_from_snapshot(snapshot, 
Arc::new(StatsRegistry::default())).into())
+        let inner = StreamsInner::inner_from_snapshot(snapshot, 
Arc::new(StatsRegistry::default()));
+        // Drop the checkpoint's aggregates here, BEFORE journal replay runs 
over
+        // this state machine. Boot rebuilds every counter from disk (each 
shard
+        // folds its `load_partition` deltas in), so the stored totals are 
never
+        // authoritative -- and a checkpoint reads a stream's total and its
+        // topics' as separate loads while the partition plane keeps counting, 
so
+        // they can disagree in either direction. Left in place, a replayed
+        // `DeleteTopic` over that torn shape rolls the topic back against a
+        // stream that never held it, clamps, and raises the rollup-underflow
+        // alarm on every boot with no real divergence behind it.
+        //
+        // The registry is empty here, so this stores zero at both levels; it 
is
+        // the same walk state transfer uses, which is why there is no second
+        // spelling of it.
+        inner.stats_registry.rebuild_parent_totals(&inner.items);

Review Comment:
   simplification: this rebuild overwrites the two `store_from_snapshot` calls 
in `inner_from_snapshot` on every non-test path, and over a fresh registry it's 
storing 0 over 0. drop the stores and this call, about 25 lines - the tests at 
4110 and 3573 go vacuous.



##########
core/common/src/types/streaming_stats.rs:
##########
@@ -19,6 +19,110 @@ use std::sync::{
     Arc,
     atomic::{AtomicU32, AtomicU64, Ordering},
 };
+use tracing::warn;
+
+/// Number of rollup decrements that could not be covered by the counter they
+/// were subtracted from.
+///
+/// A decrement bigger than the total it targets means the tree lost
+/// `parent >= sum(children)` somewhere upstream. The counters are unsigned, so
+/// without a clamp that subtraction wraps to ~1.8e19 and every reader
+/// (`get_stream`, `get_topic`, `/stats`, `/metrics`) serves the wrapped value
+/// until the process restarts.
+///
+/// Clamping trades that for a total that is merely low, and low is where it
+/// stays: increments are `fetch_add`, so every later write stacks on the
+/// clamped base and carries the shortfall with it. Only a rebuild
+/// (`rebuild_parent_totals`) or a restart puts the total back. This counter
+/// plus the `warn!` below is what tells an operator the divergence happened,
+/// and that the aggregate needs a rebuild rather than time.
+static ROLLUP_UNDERFLOWS: AtomicU64 = AtomicU64::new(0);
+
+/// Monotonic count of clamped rollup decrements. Surfaced on `/metrics` so the
+/// clamp is alertable rather than log-grep-able.
+///
+/// Hidden from the docs: `iggy_common` re-exports this module at its root, and
+/// a server-only diagnostic has no business in the published SDK surface.
+#[doc(hidden)]
+#[must_use]
+pub fn rollup_underflows() -> u64 {
+    ROLLUP_UNDERFLOWS.load(Ordering::Relaxed)
+}
+
+/// One scope/counter pair's clamp state: the labels an operator reads, plus 
the
+/// count the log throttle keys on.
+struct UnderflowSite {
+    scope: &'static str,
+    counter: &'static str,
+    clamped: AtomicU64,
+}
+
+impl UnderflowSite {
+    const fn new(scope: &'static str, counter: &'static str) -> Self {
+        Self {
+            scope,
+            counter,
+            clamped: AtomicU64::new(0),
+        }
+    }
+
+    /// Count the clamp globally, then log this pair's first and every power of
+    /// two after it.
+    ///
+    /// The throttle is per pair, not global: a skewed tree emits up to three
+    /// lines per counter per rollback and a bulk delete turns that into a
+    /// flood, so a shared count would hold a quiet pair's first-ever
+    /// divergence back for up to 1023 events behind a noisy one.
+    fn report(&self, shortfall: u64) {
+        ROLLUP_UNDERFLOWS.fetch_add(1, Ordering::Relaxed);
+        let clamped = self.clamped.fetch_add(1, Ordering::Relaxed) + 1;
+        if clamped.is_power_of_two() {
+            warn!(
+                scope = self.scope,
+                counter = self.counter,
+                shortfall,
+                clamped,
+                "rollup decrement exceeded the total it was subtracted from; 
clamped at zero"
+            );
+        }
+    }
+}
+
+static STREAM_SIZE_BYTES: UnderflowSite = UnderflowSite::new("stream", 
"size_bytes");
+static STREAM_MESSAGES_COUNT: UnderflowSite = UnderflowSite::new("stream", 
"messages_count");
+static STREAM_SEGMENTS_COUNT: UnderflowSite = UnderflowSite::new("stream", 
"segments_count");
+static TOPIC_SIZE_BYTES: UnderflowSite = UnderflowSite::new("topic", 
"size_bytes");
+static TOPIC_MESSAGES_COUNT: UnderflowSite = UnderflowSite::new("topic", 
"messages_count");
+static TOPIC_SEGMENTS_COUNT: UnderflowSite = UnderflowSite::new("topic", 
"segments_count");
+static PARTITION_SIZE_BYTES: UnderflowSite = UnderflowSite::new("partition", 
"size_bytes");
+static PARTITION_MESSAGES_COUNT: UnderflowSite = 
UnderflowSite::new("partition", "messages_count");
+static PARTITION_SEGMENTS_COUNT: UnderflowSite = 
UnderflowSite::new("partition", "segments_count");
+
+/// Subtract `amount`, clamping at zero instead of wrapping. Returns what was
+/// actually taken, which is what the caller passes on to its parent.
+///
+/// `fetch_update` rather than a load followed by a subtract: the counters are
+/// written from the metadata shard and from whichever shard owns the 
partition,
+/// so a separate load leaves a window where the clamp reads one value and
+/// subtracts from another.
+macro_rules! clamped_sub {
+    ($name:ident, $counter:ty, $amount:ty) => {
+        fn $name(counter: &$counter, amount: $amount, site: &'static 
UnderflowSite) -> $amount {
+            let previous = counter
+                .fetch_update(Ordering::AcqRel, Ordering::Acquire, |current| {
+                    Some(current.saturating_sub(amount))

Review Comment:
   simplification: the closure always returns `Some`, so the `Err` arm can't 
happen. `(current != 0).then(|| current.saturating_sub(amount))` gives it a 
meaning and skips the CAS on an already-zero counter.



-- 
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