This is an automated email from the ASF dual-hosted git repository.

numinnex pushed a commit to branch stats_update_on_delete_partitions
in repository https://gitbox.apache.org/repos/asf/iggy.git


The following commit(s) were added to 
refs/heads/stats_update_on_delete_partitions by this push:
     new cdd326a19 address review
cdd326a19 is described below

commit cdd326a19bd6ca0d4394eb475553a2b941578f8e
Author: Grzegorz Koszyk <[email protected]>
AuthorDate: Tue Sep 8 08:57:06 2026 +0200

    address review
---
 Cargo.toml                                         |   6 +-
 core/common/src/types/streaming_stats.rs           | 144 +++++--------
 .../data_integrity/verify_after_server_restart.rs  |   3 +-
 core/integration/tests/server/general.rs           |  11 +-
 .../scenarios/delete_stats_rollback_scenario.rs    |   4 +
 core/integration/tests/server/scenarios/mod.rs     |  38 ++--
 .../server/scenarios/purge_delete_scenario.rs      |  26 +--
 core/metadata/src/stm/stream.rs                    | 225 +++++++++++++--------
 core/partitions/src/iggy_partition.rs              | 110 ++++++----
 core/server/src/http/handlers.rs                   |   1 -
 core/server/src/http/metrics.rs                    | 113 ++++++-----
 core/server/src/partition_reconciler.rs            | 106 +++++++---
 12 files changed, 458 insertions(+), 329 deletions(-)

diff --git a/Cargo.toml b/Cargo.toml
index beb92e80f..43747aae0 100644
--- a/Cargo.toml
+++ b/Cargo.toml
@@ -214,7 +214,11 @@ jsonwebtoken = { version = "11.0.0", features = 
["rust_crypto"] }
 kafka-protocol = { version = "0.18.0", default-features = false, features = 
["broker"] }
 keyring-core = "1.0.0"
 lazy_static = "1.5.0"
-left-right = "0.11"
+# Pinned exactly: `metadata::stm::stream`'s keyed stats evictions lean on
+# left-right draining the previous batch's `absorb_second` before the current
+# batch's `absorb_first`, which is an implementation detail of 0.11.8's
+# `write.rs`, not a documented contract.
+left-right = "=0.11.8"
 libc = "0.2.189"
 log = "0.4.34"
 lz4_flex = "0.14.0"
diff --git a/core/common/src/types/streaming_stats.rs 
b/core/common/src/types/streaming_stats.rs
index 9ac283a1b..5d64c632c 100644
--- a/core/common/src/types/streaming_stats.rs
+++ b/core/common/src/types/streaming_stats.rs
@@ -21,8 +21,9 @@ use std::sync::{
 };
 use tracing::warn;
 
-/// Number of rollup decrements that could not be covered by the counter they
-/// were subtracted from.
+/// Process-wide number of rollup decrements that could not be covered by the
+/// counter they were subtracted from. One static for the whole node, summed
+/// across every stream, topic and partition it holds.
 ///
 /// A decrement bigger than the total it targets means the tree lost
 /// `parent >= sum(children)` somewhere upstream. The counters are unsigned, so
@@ -38,14 +39,15 @@ use tracing::warn;
 /// 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.
+/// Monotonic process-wide 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)]
+/// Named for the root it is reached from, not for this module: `iggy_common`
+/// globs the module into its crate root, so a bare `rollup_underflows` would
+/// sit in the published SDK surface next to the stream and topic types saying
+/// nothing about which rollup it counts.
 #[must_use]
-pub fn rollup_underflows() -> u64 {
+pub fn stats_rollup_underflows() -> u64 {
     ROLLUP_UNDERFLOWS.load(Ordering::Relaxed)
 }
 
@@ -73,15 +75,22 @@ impl UnderflowSite {
     /// 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.
+    ///
+    /// Both counts go on the line. `clamped` is this pair's, `total` is the
+    /// process-wide one [`stats_rollup_underflows`] exports, which carries no
+    /// scope or counter label -- without both an operator cannot tell which
+    /// pair moved the metric.
     fn report(&self, shortfall: u64) {
-        ROLLUP_UNDERFLOWS.fetch_add(1, Ordering::Relaxed);
+        let total = ROLLUP_UNDERFLOWS.fetch_add(1, Ordering::Relaxed) + 1;
         let clamped = self.clamped.fetch_add(1, Ordering::Relaxed) + 1;
         if clamped.is_power_of_two() {
             warn!(
+                target: "iggy.stats.diag",
                 scope = self.scope,
                 counter = self.counter,
                 shortfall,
                 clamped,
+                total,
                 "rollup decrement exceeded the total it was subtracted from; 
clamped at zero"
             );
         }
@@ -110,7 +119,11 @@ macro_rules! clamped_sub {
         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))
+                    // `None` on an already-zero counter: no store, and the
+                    // `Err` it returns carries that same zero, so the clamp
+                    // reports identically either way. That is the steady shape
+                    // for a scope whose rollback has already run.
+                    (current != 0).then(|| current.saturating_sub(amount))
                 })
                 .unwrap_or_else(|previous| previous);
             if previous < amount {
@@ -146,21 +159,20 @@ impl StreamStats {
             .fetch_add(segments_count, 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)
+    // A stream is the root, so a clamp here has nowhere left to be forwarded
+    // and no caller that could name a scope the `warn!` in `UnderflowSite`
+    // does not already carry. Only `PartitionStats` returns its shortfall,
+    // where the caller holds the namespace.
+    pub fn decrement_size_bytes(&self, size_bytes: u64) {
+        clamped_sub_u64(&self.size_bytes, size_bytes, &STREAM_SIZE_BYTES);
     }
 
-    pub fn decrement_messages_count(&self, messages_count: u64) -> u64 {
-        messages_count
-            - clamped_sub_u64(&self.messages_count, messages_count, 
&STREAM_MESSAGES_COUNT)
+    pub fn decrement_messages_count(&self, messages_count: u64) {
+        clamped_sub_u64(&self.messages_count, messages_count, 
&STREAM_MESSAGES_COUNT);
     }
 
-    pub fn decrement_segments_count(&self, segments_count: u32) -> u32 {
-        segments_count
-            - clamped_sub_u32(&self.segments_count, segments_count, 
&STREAM_SEGMENTS_COUNT)
+    pub fn decrement_segments_count(&self, segments_count: u32) {
+        clamped_sub_u32(&self.segments_count, segments_count, 
&STREAM_SEGMENTS_COUNT);
     }
 
     pub fn size_bytes_inconsistent(&self) -> u64 {
@@ -230,45 +242,21 @@ impl TopicStats {
         self.parent.clone()
     }
 
-    pub fn increment_parent_size_bytes(&self, size_bytes: u64) {
-        self.parent.increment_size_bytes(size_bytes);
-    }
-
-    pub fn increment_parent_messages_count(&self, messages_count: u64) {
-        self.parent.increment_messages_count(messages_count);
-    }
-
-    pub fn increment_parent_segments_count(&self, segments_count: u32) {
-        self.parent.increment_segments_count(segments_count);
-    }
-
     pub fn increment_size_bytes(&self, size_bytes: u64) {
         self.size_bytes.fetch_add(size_bytes, Ordering::AcqRel);
-        self.increment_parent_size_bytes(size_bytes);
+        self.parent.increment_size_bytes(size_bytes);
     }
 
     pub fn increment_messages_count(&self, messages_count: u64) {
         self.messages_count
             .fetch_add(messages_count, Ordering::AcqRel);
-        self.increment_parent_messages_count(messages_count);
+        self.parent.increment_messages_count(messages_count);
     }
 
     pub fn increment_segments_count(&self, segments_count: u32) {
         self.segments_count
             .fetch_add(segments_count, Ordering::AcqRel);
-        self.increment_parent_segments_count(segments_count);
-    }
-
-    pub fn decrement_parent_size_bytes(&self, size_bytes: u64) {
-        self.parent.decrement_size_bytes(size_bytes);
-    }
-
-    pub fn decrement_parent_messages_count(&self, messages_count: u64) {
-        self.parent.decrement_messages_count(messages_count);
-    }
-
-    pub fn decrement_parent_segments_count(&self, segments_count: u32) {
-        self.parent.decrement_segments_count(segments_count);
+        self.parent.increment_segments_count(segments_count);
     }
 
     // Forward what this level actually gave up, not what was asked for. A
@@ -276,26 +264,19 @@ impl TopicStats {
     // 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 {
+    pub fn decrement_size_bytes(&self, size_bytes: u64) {
         let taken = clamped_sub_u64(&self.size_bytes, size_bytes, 
&TOPIC_SIZE_BYTES);
-        self.decrement_parent_size_bytes(taken);
-        size_bytes - taken
+        self.parent.decrement_size_bytes(taken);
     }
 
-    pub fn decrement_messages_count(&self, messages_count: u64) -> u64 {
+    pub fn decrement_messages_count(&self, messages_count: u64) {
         let taken = clamped_sub_u64(&self.messages_count, messages_count, 
&TOPIC_MESSAGES_COUNT);
-        self.decrement_parent_messages_count(taken);
-        messages_count - taken
+        self.parent.decrement_messages_count(taken);
     }
 
-    pub fn decrement_segments_count(&self, segments_count: u32) -> u32 {
+    pub fn decrement_segments_count(&self, segments_count: u32) {
         let taken = clamped_sub_u32(&self.segments_count, segments_count, 
&TOPIC_SEGMENTS_COUNT);
-        self.decrement_parent_segments_count(taken);
-        segments_count - taken
+        self.parent.decrement_segments_count(taken);
     }
 
     pub fn size_bytes_inconsistent(&self) -> u64 {
@@ -372,30 +353,18 @@ impl PartitionStats {
 
     pub fn increment_size_bytes(&self, size_bytes: u64) {
         self.size_bytes.fetch_add(size_bytes, Ordering::AcqRel);
-        self.increment_parent_size_bytes(size_bytes);
+        self.parent.increment_size_bytes(size_bytes);
     }
 
     pub fn increment_messages_count(&self, messages_count: u64) {
         self.messages_count
             .fetch_add(messages_count, Ordering::AcqRel);
-        self.increment_parent_messages_count(messages_count);
+        self.parent.increment_messages_count(messages_count);
     }
 
     pub fn increment_segments_count(&self, segments_count: u32) {
         self.segments_count
             .fetch_add(segments_count, Ordering::AcqRel);
-        self.increment_parent_segments_count(segments_count);
-    }
-
-    pub fn increment_parent_size_bytes(&self, size_bytes: u64) {
-        self.parent.increment_size_bytes(size_bytes);
-    }
-
-    pub fn increment_parent_messages_count(&self, messages_count: u64) {
-        self.parent.increment_messages_count(messages_count);
-    }
-
-    pub fn increment_parent_segments_count(&self, segments_count: u32) {
         self.parent.increment_segments_count(segments_count);
     }
 
@@ -410,7 +379,7 @@ impl PartitionStats {
     // 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);
+        self.parent.decrement_size_bytes(taken);
         size_bytes - taken
     }
 
@@ -420,7 +389,7 @@ impl PartitionStats {
             messages_count,
             &PARTITION_MESSAGES_COUNT,
         );
-        self.decrement_parent_messages_count(taken);
+        self.parent.decrement_messages_count(taken);
         messages_count - taken
     }
 
@@ -430,22 +399,10 @@ impl PartitionStats {
             segments_count,
             &PARTITION_SEGMENTS_COUNT,
         );
-        self.decrement_parent_segments_count(taken);
+        self.parent.decrement_segments_count(taken);
         segments_count - taken
     }
 
-    pub fn decrement_parent_size_bytes(&self, size_bytes: u64) {
-        self.parent.decrement_size_bytes(size_bytes);
-    }
-
-    pub fn decrement_parent_messages_count(&self, messages_count: u64) {
-        self.parent.decrement_messages_count(messages_count);
-    }
-
-    pub fn decrement_parent_segments_count(&self, segments_count: u32) {
-        self.parent.decrement_segments_count(segments_count);
-    }
-
     pub fn size_bytes_inconsistent(&self) -> u64 {
         self.size_bytes.load(Ordering::Relaxed)
     }
@@ -524,7 +481,7 @@ mod tests {
     /// 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() 
{
+    fn 
given_a_rolled_back_partition_when_a_late_decrement_arrives_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()));
@@ -584,8 +541,9 @@ mod tests {
 
         partition.zero_out_all();
 
-        // u32 wraps at a different modulus than the two u64 counters, so it
-        // needs its own coverage.
+        // `clamped_sub!` expands separately per width, so the u32 counter is
+        // its own wiring: same `saturating_sub` body, its own `UnderflowSite`
+        // and its own call sites at all three levels.
         assert_eq!(topic.segments_count_inconsistent(), 0);
         assert_eq!(stream.segments_count_inconsistent(), 0);
     }
diff --git 
a/core/integration/tests/data_integrity/verify_after_server_restart.rs 
b/core/integration/tests/data_integrity/verify_after_server_restart.rs
index 3fea2e006..a9dece200 100644
--- a/core/integration/tests/data_integrity/verify_after_server_restart.rs
+++ b/core/integration/tests/data_integrity/verify_after_server_restart.rs
@@ -15,6 +15,7 @@
 // specific language governing permissions and limitations
 // under the License.
 
+use bytes::Bytes;
 use iggy::prelude::*;
 use integration::bench_utils::run_bench_and_wait_for_finish;
 use integration::harness::{TestHarness, TestServerConfig};
@@ -565,7 +566,7 @@ fn deletion_test_messages() -> Vec<IggyMessage> {
         .map(|offset| {
             IggyMessage::builder()
                 .id(offset + 1)
-                .payload(bytes::Bytes::from_static(b"deletion-stats-payload"))
+                .payload(Bytes::from_static(b"deletion-stats-payload"))
                 .build()
                 .expect("Failed to build message")
         })
diff --git a/core/integration/tests/server/general.rs 
b/core/integration/tests/server/general.rs
index e96ebc233..f1dd8382b 100644
--- a/core/integration/tests/server/general.rs
+++ b/core/integration/tests/server/general.rs
@@ -88,13 +88,10 @@ async fn stream_size_validation(harness: &TestHarness) {
     stream_size_validation_scenario::run(harness).await;
 }
 
-#[iggy_harness(
-    test_client_transport = [Tcp, Http, Quic, WebSocket],
-    server(
-        quic.max_idle_timeout = "500s",
-        quic.keep_alive_interval = "15s"
-    )
-)]
+// One transport: the scenario asserts server-side rollback arithmetic, and
+// every call it makes (`get_stream`, `get_topic`, `get_stats`) is already
+// exercised on all four transports by the scenarios above.
+#[iggy_harness(test_client_transport = [Tcp])]
 async fn delete_stats_rollback(harness: &TestHarness) {
     delete_stats_rollback_scenario::run(harness).await;
 }
diff --git 
a/core/integration/tests/server/scenarios/delete_stats_rollback_scenario.rs 
b/core/integration/tests/server/scenarios/delete_stats_rollback_scenario.rs
index 0624a6815..653806445 100644
--- a/core/integration/tests/server/scenarios/delete_stats_rollback_scenario.rs
+++ b/core/integration/tests/server/scenarios/delete_stats_rollback_scenario.rs
@@ -74,6 +74,10 @@ pub async fn run(harness: &TestHarness) {
     delete_partitions_rolls_back_topic_and_stream(&client).await;
     delete_topic_rolls_back_the_stream(&client).await;
 
+    // The plainest statement of what this guards: every scope the scenario
+    // created is gone, so the server-wide totals owe nothing. 
`assert_clean_system`
+    // reads the entity lists, never the stats.
+    validate_system_stats(&client, 0, 0).await;
     assert_clean_system(&client).await;
 }
 
diff --git a/core/integration/tests/server/scenarios/mod.rs 
b/core/integration/tests/server/scenarios/mod.rs
index 1d4ef848a..c6451d782 100644
--- a/core/integration/tests/server/scenarios/mod.rs
+++ b/core/integration/tests/server/scenarios/mod.rs
@@ -62,8 +62,14 @@ use std::time::{Duration, Instant};
 use tokio::time::sleep;
 
 const PARTITION_ID: u32 = 0;
-const POLL_CONVERGENCE_TIMEOUT: Duration = Duration::from_secs(10);
-const POLL_RETRY_INTERVAL: Duration = Duration::from_millis(100);
+// One pair for every wait in these scenarios, because they all wait out the
+// same thing: the partition plane applies committed ops asynchronously on the
+// owning shard (a send folds into the shared stats and into the servable log 
at
+// commit-apply; purge and delete zero them when the reconciler drives the 
wipe),
+// so a read racing that window sees a pre-apply value. Retry until the
+// expectation holds, then make the terminal assertion for a real mismatch.
+const CONVERGENCE_TIMEOUT: Duration = Duration::from_secs(10);
+const RETRY_INTERVAL: Duration = Duration::from_millis(100);
 const STREAM_NAME: &str = "test-stream";
 const TOPIC_NAME: &str = "test-topic";
 const PARTITIONS_COUNT: u32 = 3;
@@ -74,14 +80,6 @@ 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);
-const STATS_RETRY_INTERVAL: Duration = Duration::from_millis(100);
-
 const MESSAGE_PAYLOAD_SIZE_BYTES: u64 = 57;
 // The server accounts the actual on-disk batch framing: one 256-byte
 // `SendMessages` command header per append pass plus a 48-byte per-message
@@ -111,7 +109,7 @@ fn create_messages(messages_count: u64) -> Vec<IggyMessage> 
{
         .collect()
 }
 
-/// Fetch the stream until its totals match or [`STATS_CONVERGENCE_TIMEOUT`]
+/// Fetch the stream until its totals match or [`CONVERGENCE_TIMEOUT`]
 /// expires, then assert on the last read.
 async fn validate_stream(
     client: &IggyClient,
@@ -119,7 +117,7 @@ async fn validate_stream(
     expected_size: u64,
     expected_messages_count: u64,
 ) {
-    let deadline = Instant::now() + STATS_CONVERGENCE_TIMEOUT;
+    let deadline = Instant::now() + CONVERGENCE_TIMEOUT;
     let stream = loop {
         let stream = client
             .get_stream(&Identifier::named(stream_name).unwrap())
@@ -131,7 +129,7 @@ async fn validate_stream(
         {
             break stream;
         }
-        sleep(STATS_RETRY_INTERVAL).await;
+        sleep(RETRY_INTERVAL).await;
     };
     assert_eq!(stream.size, expected_size, "stream size mismatch");
     assert_eq!(
@@ -148,7 +146,7 @@ async fn validate_topic(
     expected_size: u64,
     expected_messages_count: u64,
 ) {
-    let deadline = Instant::now() + STATS_CONVERGENCE_TIMEOUT;
+    let deadline = Instant::now() + CONVERGENCE_TIMEOUT;
     let topic = loop {
         let topic = client
             .get_topic(
@@ -163,7 +161,7 @@ async fn validate_topic(
         {
             break topic;
         }
-        sleep(STATS_RETRY_INTERVAL).await;
+        sleep(RETRY_INTERVAL).await;
     };
     assert_eq!(topic.size, expected_size, "topic size mismatch");
     assert_eq!(
@@ -179,7 +177,7 @@ async fn validate_system_stats(
     expected_size: u64,
     expected_messages_count: u64,
 ) {
-    let deadline = Instant::now() + STATS_CONVERGENCE_TIMEOUT;
+    let deadline = Instant::now() + CONVERGENCE_TIMEOUT;
     let stats = loop {
         let stats = client.get_stats().await.unwrap();
         if (stats.messages_count == expected_messages_count
@@ -188,7 +186,7 @@ async fn validate_system_stats(
         {
             break stats;
         }
-        sleep(STATS_RETRY_INTERVAL).await;
+        sleep(RETRY_INTERVAL).await;
     };
     assert_eq!(
         stats.messages_count, expected_messages_count,
@@ -202,7 +200,7 @@ async fn validate_system_stats(
 }
 
 /// Poll until the partition serves `expected_count` messages or
-/// [`POLL_CONVERGENCE_TIMEOUT`] expires, returning the last poll result.
+/// [`CONVERGENCE_TIMEOUT`] expires, returning the last poll result.
 ///
 /// `send_messages` acks at consensus commit while the owning shard applies
 /// the batch asynchronously (see the materialisation race note at the top
@@ -217,7 +215,7 @@ async fn poll_until_expected_count(
     strategy: &PollingStrategy,
     expected_count: u32,
 ) -> PolledMessages {
-    let deadline = Instant::now() + POLL_CONVERGENCE_TIMEOUT;
+    let deadline = Instant::now() + CONVERGENCE_TIMEOUT;
     loop {
         let polled = client
             .poll_messages(
@@ -234,7 +232,7 @@ async fn poll_until_expected_count(
         if polled.messages.len() as u32 == expected_count || Instant::now() >= 
deadline {
             return polled;
         }
-        sleep(POLL_RETRY_INTERVAL).await;
+        sleep(RETRY_INTERVAL).await;
     }
 }
 
diff --git a/core/integration/tests/server/scenarios/purge_delete_scenario.rs 
b/core/integration/tests/server/scenarios/purge_delete_scenario.rs
index ce222ec62..b29130e2a 100644
--- a/core/integration/tests/server/scenarios/purge_delete_scenario.rs
+++ b/core/integration/tests/server/scenarios/purge_delete_scenario.rs
@@ -15,7 +15,7 @@
 // specific language governing permissions and limitations
 // under the License.
 
-use super::{POLL_CONVERGENCE_TIMEOUT, POLL_RETRY_INTERVAL};
+use super::{CONVERGENCE_TIMEOUT, RETRY_INTERVAL};
 use bytes::Bytes;
 use iggy::prelude::*;
 use iggy_common::Credentials;
@@ -1203,14 +1203,14 @@ pub async fn run_resident_purge_no_resurface(harness: 
&mut TestHarness) {
 /// Poll from offset 0 with headroom (count 100) until exactly `expected`
 /// messages are served, so an extra resurfaced message fails the count
 /// instead of being cropped by the poll size. Panics after
-/// [`POLL_CONVERGENCE_TIMEOUT`] with the last observed count.
+/// [`CONVERGENCE_TIMEOUT`] with the last observed count.
 async fn poll_exactly(
     client: &IggyClient,
     stream_ident: &Identifier,
     topic_ident: &Identifier,
     expected: usize,
 ) -> PolledMessages {
-    let deadline = std::time::Instant::now() + POLL_CONVERGENCE_TIMEOUT;
+    let deadline = std::time::Instant::now() + CONVERGENCE_TIMEOUT;
     loop {
         let polled = client
             .poll_messages(
@@ -1232,7 +1232,7 @@ async fn poll_exactly(
             "poll did not converge to {expected} messages, last saw {}",
             polled.messages.len()
         );
-        tokio::time::sleep(POLL_RETRY_INTERVAL).await;
+        tokio::time::sleep(RETRY_INTERVAL).await;
     }
 }
 
@@ -1318,7 +1318,7 @@ async fn maybe_restart(harness: &mut TestHarness, client: 
&IggyClient, restart_s
     let _ = client.disconnect().await;
     harness.restart_server().await.unwrap();
 
-    let deadline = tokio::time::Instant::now() + POLL_CONVERGENCE_TIMEOUT;
+    let deadline = tokio::time::Instant::now() + CONVERGENCE_TIMEOUT;
     loop {
         // `connect` re-authenticates from the embedded credentials, so a
         // successful ping means the shards are serving, not merely listening.
@@ -1327,12 +1327,12 @@ async fn maybe_restart(harness: &mut TestHarness, 
client: &IggyClient, restart_s
         }
         assert!(
             tokio::time::Instant::now() < deadline,
-            "server did not serve again within {POLL_CONVERGENCE_TIMEOUT:?} of 
restart"
+            "server did not serve again within {CONVERGENCE_TIMEOUT:?} of 
restart"
         );
         // Back to Disconnected, else the next `connect` no-ops on a half-open
         // connection and the ping keeps failing until the deadline.
         let _ = client.disconnect().await;
-        tokio::time::sleep(POLL_RETRY_INTERVAL).await;
+        tokio::time::sleep(RETRY_INTERVAL).await;
     }
 }
 
@@ -1501,17 +1501,17 @@ fn get_sorted_segment_offsets(partition_path: &str) -> 
Vec<u64> {
 }
 
 /// Wait until `dir` contains at least one entry, panicking with `context`
-/// once [`POLL_CONVERGENCE_TIMEOUT`] expires.
+/// once [`CONVERGENCE_TIMEOUT`] expires.
 ///
 /// A stored consumer offset is served from memory as soon as the store is
 /// acked, while the offset file is created asynchronously, so a single-shot
 /// existence check can run ahead of the flush. An offset that is never
 /// flushed still fails once the deadline expires.
 async fn await_dir_not_empty(dir: &str, context: &str) {
-    let deadline = std::time::Instant::now() + POLL_CONVERGENCE_TIMEOUT;
+    let deadline = std::time::Instant::now() + CONVERGENCE_TIMEOUT;
     while is_dir_empty(dir) {
         assert!(std::time::Instant::now() < deadline, "{context}");
-        tokio::time::sleep(POLL_RETRY_INTERVAL).await;
+        tokio::time::sleep(RETRY_INTERVAL).await;
     }
 }
 
@@ -1550,7 +1550,7 @@ async fn assert_fresh_empty_partition(partition_path: 
&str) {
 }
 
 /// Asserts no orphaned segment files remain after deletion, polling until the
-/// counts converge or [`POLL_CONVERGENCE_TIMEOUT`] expires.
+/// counts converge or [`CONVERGENCE_TIMEOUT`] expires.
 ///
 /// `get_sorted_segment_offsets` only checks .log files -- this additionally
 /// verifies that the .index file count matches, catching stale .index files
@@ -1558,7 +1558,7 @@ async fn assert_fresh_empty_partition(partition_path: 
&str) {
 /// separate awaits, so a layout that already converged on .log files can
 /// transiently show one extra .index file.
 async fn assert_no_orphaned_segment_files(partition_path: &str, 
expected_count: usize) {
-    let deadline = std::time::Instant::now() + POLL_CONVERGENCE_TIMEOUT;
+    let deadline = std::time::Instant::now() + CONVERGENCE_TIMEOUT;
     loop {
         let log_count = count_files_with_ext(partition_path, LOG_EXTENSION);
         let index_count = count_files_with_ext(partition_path, 
INDEX_EXTENSION);
@@ -1569,7 +1569,7 @@ async fn assert_no_orphaned_segment_files(partition_path: 
&str, expected_count:
             std::time::Instant::now() < deadline,
             "Expected {expected_count} .log and .index files, found 
{log_count} .log and {index_count} .index"
         );
-        tokio::time::sleep(POLL_RETRY_INTERVAL).await;
+        tokio::time::sleep(RETRY_INTERVAL).await;
     }
 }
 
diff --git a/core/metadata/src/stm/stream.rs b/core/metadata/src/stm/stream.rs
index b33ca0c17..f2707ea0f 100644
--- a/core/metadata/src/stm/stream.rs
+++ b/core/metadata/src/stm/stream.rs
@@ -406,7 +406,8 @@ impl Stream {
 ///   partition's entry.
 /// * `left_right` drains the previous batch's `absorb_second` before the
 ///   current batch's `absorb_first` (0.11.8 `write.rs`), which is an
-///   implementation detail, not a contract.
+///   implementation detail, not a contract -- which is why the workspace
+///   manifest pins `left-right` to `=0.11.8` rather than floating 0.11.x.
 ///
 /// Narrowing both, `fetch_partition_build_inputs` refuses to register a
 /// partition the committed topic does not list. It does not close the window:
@@ -481,6 +482,16 @@ impl StatsRegistry {
     /// partition was torn down and rebuilt still counts as pending, so its
     /// deferred second-buffer apply wipes everything appended since.
     ///
+    /// Caller contract, unchecked either way: `partition` must be the 
committed
+    /// record listed under `(stream_id, topic_id)`, and `parent` the committed
+    /// topic's own `Arc`. A fresh entry inherits both without comparing them 
to
+    /// anything, so a record from another topic seeds the identity gate with a
+    /// revision that never matches, and a parent from another topic sends this
+    /// partition's increments into a stranger's totals for the entry's whole
+    /// life. Read all three out of one `streams().read` (see
+    /// `fetch_partition_build_inputs` in the reconciler) and both hold by
+    /// construction.
+    ///
     /// # Panics
     /// If the registry mutex is poisoned.
     pub fn partition(
@@ -716,14 +727,31 @@ impl StatsRegistry {
             .lock()
             .expect("stats registry mutex poisoned")
             .retain(|key, _| live_topics.contains(key));
-        self.partitions
-            .lock()
-            .expect("stats registry mutex poisoned")
-            .retain(|key, entry| {
-                live_partitions
-                    .get(key)
-                    .is_some_and(|created_revision| *created_revision == 
entry.created_revision)
-            });
+        // Zeroing is the other half of the eviction, exactly as in
+        // [`Self::remove_partitions`]: a stale incarnation still mounted holds
+        // the same `Arc`, and its `ConfirmRemove` rolls those counters back
+        // through it. `rebuild_parent_totals` runs right after this and counts
+        // only survivors, so an entry dropped full leaves that later rollback
+        // to come out of a live sibling's totals.
+        let dropped: Vec<Arc<PartitionStats>> = {
+            let mut entries = self
+                .partitions
+                .lock()
+                .expect("stats registry mutex poisoned");
+            entries
+                .extract_if(|key, entry| {
+                    live_partitions
+                        .get(key)
+                        .is_none_or(|created_revision| *created_revision != 
entry.created_revision)
+                })
+                .map(|(_, entry)| entry.stats)
+                .collect()
+        };
+        // Guard released first: the rollback cascades into parent totals, 
which
+        // the partition map has no part in.
+        for stats in dropped {
+            stats.zero_out_all();
+        }
     }
 
     /// Recompute every topic and stream total as the sum of the partition
@@ -742,25 +770,41 @@ impl StatsRegistry {
     /// 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.
+    /// The counters are sampled under the guard and summed without it: the 
walk
+    /// visits every partition in the tree, and the same mutex is on the
+    /// get-or-create path the reconciler and every shard's boot recovery take.
+    /// Holding it across the walk would serialize them behind it.
+    ///
+    /// The guard fences the MAP, never the counters, which the data plane
+    /// reaches through the `Arc` regardless. An append landing between the
+    /// sample and the parent store is folded into the parent by `fetch_add` 
and
+    /// then overwritten, so the invariant is restored modulo whatever arrives
+    /// in between. 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>) {
-        let entries = self
-            .partitions
-            .lock()
-            .expect("stats registry mutex poisoned");
+        let entries: AHashMap<(usize, usize, usize), (u64, u64, u32)> = {
+            let partitions = self
+                .partitions
+                .lock()
+                .expect("stats registry mutex poisoned");
+            partitions
+                .iter()
+                .map(|(key, entry)| {
+                    (
+                        *key,
+                        (
+                            entry.stats.size_bytes_inconsistent(),
+                            entry.stats.messages_count_inconsistent(),
+                            entry.stats.segments_count_inconsistent(),
+                        ),
+                    )
+                })
+                .collect()
+        };
         for (stream_key, stream) in streams {
             let mut stream_size_bytes = 0u64;
             let mut stream_messages_count = 0u64;
@@ -770,15 +814,14 @@ impl StatsRegistry {
                 let mut topic_messages_count = 0u64;
                 let mut topic_segments_count = 0u32;
                 for partition in &topic.partitions {
-                    let Some(entry) = entries.get(&(stream_key, topic_key, 
partition.id)) else {
+                    let Some((size_bytes, messages_count, segments_count)) =
+                        entries.get(&(stream_key, topic_key, partition.id))
+                    else {
                         continue;
                     };
-                    topic_size_bytes =
-                        
topic_size_bytes.saturating_add(entry.stats.size_bytes_inconsistent());
-                    topic_messages_count = topic_messages_count
-                        
.saturating_add(entry.stats.messages_count_inconsistent());
-                    topic_segments_count = topic_segments_count
-                        
.saturating_add(entry.stats.segments_count_inconsistent());
+                    topic_size_bytes = 
topic_size_bytes.saturating_add(*size_bytes);
+                    topic_messages_count = 
topic_messages_count.saturating_add(*messages_count);
+                    topic_segments_count = 
topic_segments_count.saturating_add(*segments_count);
                 }
                 topic.stats.store_from_snapshot(
                     topic_size_bytes,
@@ -2547,21 +2590,11 @@ 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.
-        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);
+        // The checkpoint's topic and stream aggregates are not adopted (see
+        // `inner_from_snapshot`), so this starts every total at zero and boot
+        // folds the real ones back in from disk.
+        let inner = StreamsInner::inner_from_snapshot(snapshot, 
Arc::new(StatsRegistry::default()));
         Ok(inner.into())
     }
 }
@@ -2584,9 +2617,8 @@ impl StreamsInner {
         // that slot next.
         registry.retain_from_snapshot(&snapshot);
         *self = Self::inner_from_snapshot(snapshot, registry);
-        // Last, over the totals `inner_from_snapshot` just stored: the donor's
-        // aggregates describe the donor's partitions, and this node kept its
-        // own the whole time.
+        // The donor's aggregates describe the donor's partitions; this node 
kept
+        // its own the whole time, so the totals come from the surviving 
entries.
         self.stats_registry.rebuild_parent_totals(&self.items);
     }
 
@@ -2594,6 +2626,18 @@ impl StreamsInner {
     /// `stats_registry`. Shared by wrapper construction
     /// ([`Snapshotable::from_snapshot`]) and the in-place restore command
     /// (state transfer), which absorbs it on both left-right buffers.
+    ///
+    /// [`StatsSnapshot`] is decoded but never stored. Both callers derive the
+    /// topic and stream totals from this node's own partition entries instead:
+    /// state transfer through [`StatsRegistry::rebuild_parent_totals`], boot 
by
+    /// folding each shard's `load_partition` delta in as it materializes.
+    /// Adopting them would be actively wrong at both ends. 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, and a
+    /// replayed `DeleteTopic` over that torn shape rolls a 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. A snapshot's totals 
are
+    /// the DONOR's, over children the receiver kept the whole time.
     pub(crate) fn inner_from_snapshot(
         snapshot: StreamsSnapshot,
         stats_registry: Arc<StatsRegistry>,
@@ -2603,11 +2647,6 @@ impl StreamsInner {
 
         for (slab_key, stream_snap) in snapshot.items {
             let stream_stats = stats_registry.stream(slab_key);
-            stream_stats.store_from_snapshot(
-                stream_snap.stats.size_bytes,
-                stream_snap.stats.messages_count,
-                stream_snap.stats.segments_count,
-            );
 
             let mut topic_index: AHashMap<Arc<str>, usize> = AHashMap::new();
             let mut topic_entries: Vec<(usize, Topic)> = Vec::new();
@@ -2615,11 +2654,6 @@ impl StreamsInner {
             for (topic_slab_key, topic_snap) in stream_snap.topics {
                 let topic_stats =
                     stats_registry.topic(slab_key, topic_slab_key, 
stream_stats.clone());
-                topic_stats.store_from_snapshot(
-                    topic_snap.stats.size_bytes,
-                    topic_snap.stats.messages_count,
-                    topic_snap.stats.segments_count,
-                );
                 let topic_name: Arc<str> = Arc::from(topic_snap.name.as_str());
                 let topic = Topic {
                     id: topic_snap.id,
@@ -3563,14 +3597,14 @@ mod tests {
         assert_eq!(stream_stats.messages_count_inconsistent(), 0);
     }
 
-    /// A state transfer stores the DONOR's topic and stream totals over
-    /// partition counters the receiver kept the whole time, so
-    /// `topic < sum(partitions)` is the ordinary post-transfer shape. Rolling 
a
-    /// partition out of a parent that never counted it subtracts past zero, 
and
-    /// the totals are unsigned: `get_stream` / `get_topic` / `/stats` would
-    /// serve ~1.8e19 until the process restarted.
+    /// A restore replaces the whole tree while the receiver's partition
+    /// counters keep running, and it adopts no aggregate of its own: without
+    /// the rebuild every topic and stream would come back at zero over live
+    /// children. Rolling a partition out of a parent that never counted it
+    /// subtracts past zero, and the totals are unsigned: `get_stream` /
+    /// `get_topic` / `/stats` would serve ~1.8e19 until the process restarted.
     #[test]
-    fn 
given_snapshot_totals_below_the_live_partitions_when_restoring_should_not_wrap_on_delete()
 {
+    fn 
given_a_restore_over_live_partition_counters_when_deleting_should_not_wrap() {
         let mut inner = inner_with_registered_partition();
         let stats = inner.stats_registry.partition_get(0, 0, 
0).expect("stats");
         stats.increment_segments_count(3);
@@ -3578,14 +3612,7 @@ mod tests {
         stats.increment_size_bytes(512);
 
         // The donor captured this tree before any of that landed.
-        let mut snapshot = Streams::from(inner.clone()).to_snapshot();
-        let donor_totals = StatsSnapshot {
-            size_bytes: 0,
-            messages_count: 0,
-            segments_count: 0,
-        };
-        snapshot.items[0].1.stats = donor_totals.clone();
-        snapshot.items[0].1.topics[0].1.stats = donor_totals;
+        let snapshot = Streams::from(inner.clone()).to_snapshot();
         inner.restore_in_place(snapshot);
 
         // The rebuild has to put the receiver's own children back into the
@@ -3611,8 +3638,8 @@ mod tests {
         let stream_stats = &inner.items[0].stats;
         assert_eq!(stream_stats.size_bytes_inconsistent(), 0);
         assert_eq!(stream_stats.messages_count_inconsistent(), 0);
-        // u32, so it wraps at a different modulus than the two above and needs
-        // its own assertion.
+        // Its own `clamped_sub!` expansion and its own `UnderflowSite`, so the
+        // u32 wiring needs an assertion of its own.
         assert_eq!(stream_stats.segments_count_inconsistent(), 0);
     }
 
@@ -4102,12 +4129,12 @@ mod tests {
 
     /// A checkpoint reads a stream's total and each of its topics' as separate
     /// loads while the partition plane keeps counting, so the two can disagree
-    /// in either direction. The boot restore drops both levels, and drops them
-    /// before journal replay runs: a replayed `DeleteTopic` over the surviving
-    /// residue rolls a topic back against a stream that never held it, clamps,
-    /// and raises the underflow alarm on every boot after.
+    /// in either direction. The boot restore adopts neither level, and journal
+    /// replay runs over what it produces: a replayed `DeleteTopic` over an
+    /// adopted torn shape would roll a topic back against a stream that never
+    /// held it, clamp, and raise the underflow alarm on every boot after.
     #[test]
-    fn 
given_a_torn_checkpoint_when_restoring_at_boot_should_drop_both_levels() {
+    fn 
given_a_torn_checkpoint_when_restoring_at_boot_should_adopt_neither_level() {
         let mut snapshot = 
Streams::from(inner_with_registered_partition()).to_snapshot();
         // Topics summing above their stream is the torn shape: the stream was
         // read first, and the topic kept counting before its own read.
@@ -4209,6 +4236,44 @@ mod tests {
         );
     }
 
+    /// The eviction has to empty what it drops. A partition the snapshot does
+    /// not carry stays mounted until the reconciler reaches it, holding the
+    /// same `Arc`, and its `ConfirmRemove` rolls those counters back through 
it.
+    /// The rebuild counts survivors only, so an entry dropped full leaves that
+    /// rollback to come out of a live sibling's totals.
+    #[test]
+    fn 
given_a_pruned_entry_when_its_partition_is_torn_down_later_should_leave_survivors_alone()
 {
+        let mut inner = inner_with_registered_partition();
+        let doomed = inner.stats_registry.partition_get(0, 0, 
0).expect("stats");
+        doomed.increment_size_bytes(512);
+
+        // A donor tree with the same shape but a partition this node's entry
+        // cannot be, so the retain drops it and the restore registers nothing
+        // in its place.
+        let mut snapshot = Streams::from(inner.clone()).to_snapshot();
+        snapshot.items[0].1.topics[0].1.partitions[0].created_revision += 1;
+        inner.restore_in_place(snapshot);
+        assert!(inner.stats_registry.partition_get(0, 0, 0).is_none());
+
+        let survivor = inner.stats_registry.partition(
+            0,
+            0,
+            &committed_partition(&inner, 0, 0, 0),
+            inner.items[0].topics[0].stats.clone(),
+        );
+        survivor.increment_size_bytes(900);
+
+        // The mounted stale incarnation, reaching its drop point.
+        doomed.zero_out_all();
+
+        assert_eq!(
+            inner.items[0].topics[0].stats.size_bytes_inconsistent(),
+            900,
+            "the survivor's bytes must not pay for the pruned entry's rollback"
+        );
+        assert_eq!(inner.items[0].stats.size_bytes_inconsistent(), 900);
+    }
+
     /// Admission is all that stands between a client and a namespace 
collision.
     /// `IggyNamespace::new` used to mask, so slab key `MAX_TOPICS` packed
     /// byte-identically to key 0: same shard, consensus group, directory and
diff --git a/core/partitions/src/iggy_partition.rs 
b/core/partitions/src/iggy_partition.rs
index 228551fe2..5d363ef24 100644
--- a/core/partitions/src/iggy_partition.rs
+++ b/core/partitions/src/iggy_partition.rs
@@ -5947,6 +5947,8 @@ where
             budget_spent,
             ..SegmentRemoval::default()
         };
+        let mut shortfall = CleanupShortfall::default();
+        let mut removed_offsets: Option<(u64, u64)> = None;
         for _ in 0..removable {
             // The removable run is always a prefix (oldest first), so the next
             // victim is the front once the previous one is gone.
@@ -5979,13 +5981,13 @@ where
                 }
             }
 
-            // The removal loop above only reaches sealed segments, which 
always
-            // hold at least one message, so the count is inclusive 
end..=start.
-            // A one-message sealed segment has `start_offset == end_offset`, 
so
-            // the `+ 1` is required (a `start == end -> 0` special case would
-            // undercount it).
-            let messages_in_segment = segment.end_offset - 
segment.start_offset + 1;
-            settle_cleaned_segment(&self.stats, namespace, &segment, 
messages_in_segment);
+            let (messages_in_segment, segment_shortfall) =
+                settle_cleaned_segment(&self.stats, &segment);
+            shortfall.absorb(segment_shortfall);
+            removed_offsets = Some(match removed_offsets {
+                Some((from, _)) => (from, segment.end_offset),
+                None => (segment.start_offset, segment.end_offset),
+            });
 
             removal.segments += 1;
             removal.messages += messages_in_segment;
@@ -6000,6 +6002,8 @@ where
             );
         }
 
+        shortfall.report(namespace, removed_offsets);
+
         removal
     }
 
@@ -7260,37 +7264,69 @@ 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.
+/// What a cleanup rollback could not take out of the partition counters,
+/// summed over one [`IggyPartition::remove_sealed_segments_up_to`] call.
 ///
-/// 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!(
-        target: "iggy.partitions.diag",
-        plane = "partitions",
-        namespace_raw = namespace.inner(),
-        start_offset = segment.start_offset,
-        size_shortfall,
-        segments_shortfall,
-        messages_shortfall,
-        "segment cleanup gave back more than the partition counters held; the 
parent totals \
-         are now low by the shortfall until a rebuild or a restart"
-    );
+/// Accumulated rather than reported per segment: `UnderflowSite::report` in
+/// `iggy_common` already counts every clamp and power-of-two throttles its own
+/// line, so a per-segment `warn!` here is an unthrottled second copy of it. 
One
+/// line per call, carrying the offset range the call removed, says which
+/// partition and which segments without that.
+#[derive(Debug, Default, Clone, Copy)]
+struct CleanupShortfall {
+    size_bytes: u64,
+    segments: u32,
+    messages: u64,
+}
+
+impl CleanupShortfall {
+    const fn absorb(&mut self, other: Self) {
+        self.size_bytes += other.size_bytes;
+        self.segments += other.segments;
+        self.messages += other.messages;
+    }
+
+    /// One line for the whole call, naming the partition and the offset range
+    /// it removed. Silent when the counters covered every rollback.
+    fn report(self, namespace: IggyNamespace, removed_offsets: Option<(u64, 
u64)>) {
+        if self.size_bytes == 0 && self.segments == 0 && self.messages == 0 {
+            return;
+        }
+        let (removed_from, removed_to) = removed_offsets.unwrap_or_default();
+        warn!(
+            target: "iggy.partitions.diag",
+            plane = "partitions",
+            namespace_raw = namespace.inner(),
+            removed_from,
+            removed_to,
+            size_shortfall = self.size_bytes,
+            segments_shortfall = self.segments,
+            messages_shortfall = self.messages,
+            "segment cleanup gave back more than the partition counters held; 
the parent \
+             totals are now low by the shortfall until a rebuild or a restart"
+        );
+    }
+}
+
+/// Roll one cleaned-up segment out of the partition counters.
+///
+/// Returns the messages the segment held, which is what the caller reports as
+/// removed, plus whatever the rollback could not cover. 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.
+fn settle_cleaned_segment(stats: &PartitionStats, segment: &Segment) -> (u64, 
CleanupShortfall) {
+    // The removal loop only reaches sealed segments, which always hold at 
least
+    // one message, so the count is inclusive start..=end. A one-message sealed
+    // segment has `start_offset == end_offset`, so the `+ 1` is required (a
+    // `start == end -> 0` special case would undercount it).
+    let messages = segment.end_offset - segment.start_offset + 1;
+    let shortfall = CleanupShortfall {
+        size_bytes: stats.decrement_size_bytes(segment.size.as_bytes_u64()),
+        segments: stats.decrement_segments_count(1),
+        messages: stats.decrement_messages_count(messages),
+    };
+    (messages, shortfall)
 }
 
 /// Highest `end_offset` among the leading run of expired sealed segments, or
diff --git a/core/server/src/http/handlers.rs b/core/server/src/http/handlers.rs
index 64bcbd8a8..cdd8f638c 100644
--- a/core/server/src/http/handlers.rs
+++ b/core/server/src/http/handlers.rs
@@ -670,7 +670,6 @@ pub(in crate::http) async fn get_metrics(
     metrics.messages.set(gauge_value(messages_count));
     metrics.users.set(gauge_value(users_count));
     metrics.clients.set(gauge_value(clients_count));
-    metrics.observe_rollup_underflows(iggy_common::rollup_underflows());
     metrics.formatted_output()
 }
 
diff --git a/core/server/src/http/metrics.rs b/core/server/src/http/metrics.rs
index 45c34309d..c409392ff 100644
--- a/core/server/src/http/metrics.rs
+++ b/core/server/src/http/metrics.rs
@@ -21,13 +21,58 @@
 //! other route handlers so this leaf never imports the state hub.
 
 use configs::http::HttpMetricsConfig;
-use iggy_common::IggyError;
+use iggy_common::{IggyError, stats_rollup_underflows};
+use prometheus_client::collector::Collector;
 use prometheus_client::encoding::text::encode;
-use prometheus_client::metrics::counter::Counter;
+use prometheus_client::encoding::{DescriptorEncoder, EncodeMetric};
+use prometheus_client::metrics::counter::{ConstCounter, Counter};
 use prometheus_client::metrics::gauge::Gauge;
 use prometheus_client::registry::Registry;
 use tracing::error;
 
+/// Exports the process-wide clamped-rollup count as a counter series, read at
+/// encode time rather than mirrored into one.
+///
+/// Non-zero means a partition, topic or stream total was asked to give back
+/// more than it held, so what that scope now reports is low, and stays low
+/// until a rebuild or a restart. All three levels clamp and all three feed 
this
+/// one counter, which carries no scope label -- the `warn!` in `iggy_common`
+/// names the scope and counter that moved it.
+///
+/// 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 -- 
`given_a_rolled_back_partition_when_a_late_decrement_arrives_should_leave_siblings_alone`
+/// in `core/common/src/types/streaming_stats.rs` drives exactly that. What is
+/// worth paging on is a bounded rate OUTSIDE those windows: that is the shape
+/// that says the tree is diverging rather than settling.
+///
+/// A counter, not a gauge: the source only ever climbs within a process, so
+/// `rate()` and `increase()` are the queries an operator wants, and both are
+/// counter-only. A restart resets the source and the exposition together,
+/// which is the counter reset Prometheus expects.
+///
+/// Read only by the `/metrics` scrape, so it needs `http.enabled` and
+/// `http.metrics.enabled` -- as does every other series in this registry,
+/// which is the only Prometheus surface the server has. A TCP-only or
+/// QUIC-only deployment exports nothing, and the per-scope `warn!` in
+/// `iggy_common` is the whole signal there.
+#[derive(Debug)]
+struct StatsRollupUnderflows;
+
+impl Collector for StatsRollupUnderflows {
+    fn encode(&self, mut encoder: DescriptorEncoder) -> Result<(), 
std::fmt::Error> {
+        let counter = ConstCounter::new(stats_rollup_underflows());
+        let metric_encoder = encoder.encode_descriptor(
+            "stats_rollup_underflows",
+            "total count of aggregate stats decrements clamped at zero",
+            None,
+            counter.metric_type(),
+        )?;
+        counter.encode(metric_encoder)
+    }
+}
+
 /// The legacy server's metric set, registered under the same names and help
 /// texts so existing dashboards and alerts keep working unchanged.
 ///
@@ -45,30 +90,6 @@ 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`
-    /// unit test in `iggy_common::streaming_stats` drives exactly that. What 
is
-    /// worth paging on is a bounded rate OUTSIDE those windows: that is the
-    /// shape that says the tree is diverging rather than settling.
-    ///
-    /// A `Counter`, not a `Gauge`: it only ever climbs within a process, so
-    /// `rate()` and `increase()` are the queries an operator wants, and both 
are
-    /// counter-only. The scrape handler samples the source and advances this 
by
-    /// the difference, which is why the source is monotonic per process.
-    ///
-    /// Advanced only by the `/metrics` scrape, so it needs `http.enabled` and
-    /// `http.metrics.enabled` -- as does every other series in this registry,
-    /// which is the only Prometheus surface the server has. A TCP-only or
-    /// QUIC-only deployment exports nothing, and the per-scope `warn!` in
-    /// `iggy_common` is the whole signal there.
-    stats_rollup_underflows: Counter,
 }
 
 impl HttpMetrics {
@@ -82,7 +103,6 @@ impl HttpMetrics {
         let messages = Gauge::default();
         let users = Gauge::default();
         let clients = Gauge::default();
-        let stats_rollup_underflows = Counter::default();
         registry.register(
             "http_requests",
             "total count of http_requests",
@@ -99,11 +119,9 @@ impl HttpMetrics {
         registry.register("messages", "total count of messages", 
messages.clone());
         registry.register("users", "total count of users", users.clone());
         registry.register("clients", "total count of clients", 
clients.clone());
-        registry.register(
-            "stats_rollup_underflows",
-            "total count of aggregate stats decrements clamped at zero",
-            stats_rollup_underflows.clone(),
-        );
+        // Not a legacy-parity metric, and not a mirrored one: the source is a
+        // process-wide static the scrape reads directly.
+        registry.register_collector(Box::new(StatsRollupUnderflows));
         // Every shard's drop / reconcile / partition counters, one
         // `shard`-labelled sub-registry per shard so series stay per-shard
         // without a `shard_id` label in the counter label sets (see
@@ -127,7 +145,6 @@ impl HttpMetrics {
             messages,
             users,
             clients,
-            stats_rollup_underflows,
         }
     }
 
@@ -137,17 +154,6 @@ impl HttpMetrics {
         self.http_requests.clone()
     }
 
-    /// Advance the rollup-underflow counter to `total`, the process-wide count
-    /// sampled at scrape time. A counter cannot be set, so this adds the
-    /// difference; `total` only climbs within a process, and a restart resets
-    /// both sides together, which is the counter reset Prometheus expects.
-    pub(in crate::http) fn observe_rollup_underflows(&self, total: u64) {
-        let recorded = self.stats_rollup_underflows.get();
-        if let Some(delta) = total.checked_sub(recorded) {
-            self.stats_rollup_underflows.inc_by(delta);
-        }
-    }
-
     pub(in crate::http) fn formatted_output(&self) -> String {
         let mut buffer = String::new();
         if let Err(error) = encode(&mut buffer, &self.registry) {
@@ -245,17 +251,22 @@ mod tests {
         );
     }
 
+    /// Outside `PARITY_METRIC_NAMES`: the legacy server had no such series, so
+    /// it gets its own assertion rather than a row in the parity list.
+    ///
+    /// The value is a process-wide static shared with every other test in this
+    /// binary, so this asserts the series and its type, not a number.
     #[test]
-    fn rollup_underflow_total_tracks_the_sampled_count() {
+    fn rollup_underflow_collector_lands_in_the_exposition() {
         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.
-        metrics.observe_rollup_underflows(0);
         let output = metrics.formatted_output();
         assert!(
-            output.contains("\nstats_rollup_underflows_total 5\n"),
-            "expected the sampled total in the exposition:\n{output}"
+            output.contains("# TYPE stats_rollup_underflows counter\n"),
+            "expected the collector's series in the exposition:\n{output}"
+        );
+        assert!(
+            output.contains("\nstats_rollup_underflows_total "),
+            "expected the collector's sample in the exposition:\n{output}"
         );
     }
 
diff --git a/core/server/src/partition_reconciler.rs 
b/core/server/src/partition_reconciler.rs
index 8a0c61f9c..10be2b704 100644
--- a/core/server/src/partition_reconciler.rs
+++ b/core/server/src/partition_reconciler.rs
@@ -580,12 +580,7 @@ async fn reconcile_additions(
     let partitions = ctx.shard.plane.partitions();
     let total_shards = u32::from(ctx.total_shards);
 
-    for TargetPartition {
-        ns,
-        epoch,
-        created_view,
-    } in target
-    {
+    for TargetPartition { ns, epoch } in target {
         if partitions.contains(&ns) {
             // Tombstoned but still in the map. Two cases, told apart by
             // whether teardown's disk delete succeeded:
@@ -660,6 +655,20 @@ async fn reconcile_additions(
         // bumps `Streams::revision`, which forces the next pass past the
         // fast-skip.
         if partitions.is_tombstoned(&ns) {
+            // Unless a teardown put it there. The mid-pass-recreate arm below
+            // tears down a prior life this shard never mounted, and a disk
+            // delete that failed leaves the tombstone standing with no
+            // `ConfirmRemove` behind it -- the same permanent fence the in-map
+            // branch above escapes, told apart by the same signal.
+            if ctx.has_pending_delete_failure(ns) {
+                trace!(
+                    shard = shard_id,
+                    ns_raw = ns.inner(),
+                    "additions: ns tombstoned before materialisation with a 
failed disk delete; re-driving teardown"
+                );
+                tear_down_owned_partition(ctx, ns, counters).await;
+                continue;
+            }
             trace!(
                 shard = shard_id,
                 ns_raw = ns.inner(),
@@ -727,7 +736,36 @@ async fn reconcile_additions(
             ctx.config
                 .system
                 .get_partition_path(ns.stream_id(), ns.topic_id(), 
ns.partition_id());
-        let built = if std::fs::metadata(&partition_dir).is_ok() {
+        let prior_life_on_disk = std::fs::metadata(&partition_dir).is_ok();
+
+        // The target was snapshotted before this read, so a delete plus a
+        // recreate of the same slab keys can commit in between. Everything
+        // below has to describe the incarnation the row will carry, so the
+        // committed record wins over the snapshot on both counts.
+        let created_revision = partition_metadata.created_revision;
+        let created_view = partition_metadata.created_view;
+        if created_revision != epoch {
+            trace!(
+                shard = shard_id,
+                ns_raw = ns.inner(),
+                target_epoch = epoch,
+                created_revision,
+                prior_life_on_disk,
+                "additions: recreate committed mid-pass; deferring the build 
to the next pass"
+            );
+            counters.deferred += 1;
+            // Skipping alone only settles it when nothing is on disk. With a
+            // directory there, the next pass takes the loader arm below and
+            // hydrates the NEW incarnation out of the OLD segments, so the
+            // prior life has to go first. Teardown tombstones until its
+            // `ConfirmRemove` lands, which is what holds that pass off.
+            if prior_life_on_disk {
+                tear_down_owned_partition(ctx, ns, counters).await;
+            }
+            continue;
+        }
+
+        let built = if prior_life_on_disk {
             load_partition_or_fence(
                 ctx.config.as_ref(),
                 ns,
@@ -746,7 +784,7 @@ async fn reconcile_additions(
                 ctx.config.as_ref(),
                 ns,
                 partition_stats,
-                epoch,
+                created_revision,
                 topic_runtime,
                 ctx.cluster_id,
                 ctx.self_replica_id,
@@ -762,7 +800,7 @@ async fn reconcile_additions(
                 ctx.shard.enqueue_reconcile_op(ReconcileOp::InsertOwned {
                     namespace: ns,
                     partition: Box::new(partition),
-                    epoch,
+                    epoch: created_revision,
                 });
                 ctx.record_success(ns, FailureCause::Add);
                 counters.materialised += 1;
@@ -977,12 +1015,6 @@ 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);
 
     if let Err(err) = delete_partitions_from_disk(
         ns.stream_id(),
@@ -1005,6 +1037,17 @@ async fn tear_down_owned_partition(
         return;
     }
 
+    // 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.
+    //
+    // Paired with the enqueue, not with the fence above it: a failed disk
+    // delete returns without a `ConfirmRemove` behind it, and an entry already
+    // evicted would leave the still-mounted partition serving the 
registry-miss
+    // shape while its files are still there.
+    settle_partition_stats(ctx, ns);
     ctx.shard
         .enqueue_reconcile_op(ReconcileOp::ConfirmRemove { namespace: ns });
     ctx.record_success(ns, FailureCause::Delete);
@@ -1222,9 +1265,6 @@ struct TargetPartition {
     /// partition: [`fetch_partition_build_inputs`] re-reads committed metadata
     /// for the namespaces actually built, and takes the stats there.
     epoch: u64,
-    /// The view a fresh materialisation seeds its consensus group with; see
-    /// `build_partition_fresh`.
-    created_view: u32,
 }
 
 /// Every committed partition, as the additions pass needs it.
@@ -1241,7 +1281,6 @@ fn snapshot_target_namespaces(ctx: &ReconcilerCtx) -> 
Vec<TargetPartition> {
                     entries.push(TargetPartition {
                         ns: IggyNamespace::new(stream.id, topic_id, 
partition.id),
                         epoch: partition.created_revision,
-                        created_view: partition.created_view,
                     });
                 }
             }
@@ -1273,6 +1312,16 @@ fn current_revision(ctx: &ReconcilerCtx) -> u64 {
 /// seeds that gate from the committed partition.)
 ///
 /// The mounted partition's own handle is settled later, on `ConfirmRemove`.
+///
+/// The window this opens, on the teardown-for-rebuild path only. Between this
+/// eviction and the rebuild's get-or-create a pass later, `partition_get`
+/// answers `None`, so every reply builder serves the registry-miss shape for
+/// the namespace (no segments, offset 0) and its parents read low by whatever
+/// the dead incarnation held. Deliberate: the alternative is carrying a
+/// materialisation signal through the registry so readers could tell "not
+/// mounted here" from "empty", and the numbers are wrong either way while the
+/// rebuild is pending. The rebuild folds the on-disk delta back in and the
+/// window closes on its own; a delete has no rebuild and so no window.
 fn settle_partition_stats(ctx: &ReconcilerCtx, ns: IggyNamespace) {
     ctx.stats_registry
         .remove_partitions(ns.stream_id(), ns.topic_id(), 
&[ns.partition_id()]);
@@ -2001,9 +2050,16 @@ mod tests {
     /// 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.
+    /// account for them. The `ConfirmRemove` arm is where that residue gets
+    /// settled: it drops the partition value, so it is the last moment a
+    /// suspended handler can still increment through its cached handle.
+    ///
+    /// Named for the drop point because that is what it holds: the apply had
+    /// already evicted the entry here, so [`settle_partition_stats`] finds
+    /// nothing and the rollback comes from the mounted partition's own handle
+    /// in `shard`'s pump. The eviction is guarded by the sibling test below.
     #[compio::test]
-    async fn 
teardown_settles_the_residue_the_metadata_delete_could_not_reach() {
+    async fn 
confirm_remove_rolls_back_what_the_metadata_delete_could_not_reach() {
         let tmp = TempDir::new().expect("tempdir for system path");
         let config = test_config(&tmp);
         let mux = TestMux::default();
@@ -2033,13 +2089,13 @@ mod tests {
         assert_eq!(
             stream_size(&ctx),
             0,
-            "teardown must settle what landed after the commit that acked the 
delete"
+            "the drop point must roll back what landed after the commit that 
acked the delete"
         );
     }
 
-    /// The eviction half of the teardown settle, which the test above cannot
-    /// reach: there the metadata apply had already dropped the entry, so the
-    /// rollback came from the mounted partition's own handle. A stale
+    /// [`settle_partition_stats`], the half the test above cannot reach: there
+    /// the metadata apply had already dropped the entry, so the rollback came
+    /// from the mounted partition's own handle at the drop point. A stale
     /// incarnation left by a slab-key reuse still holds an entry the apply
     /// never named, and it has to go, or the rebuild's get-or-create inherits
     /// the dead incarnation's counters.

Reply via email to