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


##########
core/server/src/partition_reconciler.rs:
##########
@@ -928,10 +926,21 @@ async fn tear_down_owned_partition(
     // on_replicate / on_ack frames that haven't observed the queued
     // tombstone yet. Idempotent on retry: already-tombstoned namespace
     // stays tombstoned; already-removed shards_table row is a no-op.
+    // Take the counters handle BEFORE the tombstone (every partition accessor
+    // is tombstone-gated) and settle it AFTER, once the tombstone has stopped
+    // new frames from resolving the partition. The metadata apply already
+    // rolled it out at commit time, so what is left here is whatever landed in
+    // the window between the two: in-flight appends adding to parents that
+    // outlive the partition, and a retention tick subtracting from counters 
the
+    // apply had already zeroed. A frame already past the gate can still land
+    // after this settle, which is why the rollback clamps rather than trusting
+    // the amount it is handed.
+    let live_stats = partitions.with_partition(&ns, |partition| 
Arc::clone(&partition.stats));
     if !partitions.is_tombstoned(&ns) {
         partitions.tombstone(ns);
     }
     shards_table.remove(&ns);
+    settle_partition_stats(ctx, ns, live_stats.as_ref());

Review Comment:
   warning: a pump handler suspended mid-append resumes after this settle, and 
its increment lands in the parent totals with nothing left to roll it back. 
settle off the removed partition in the `ConfirmRemove` arm instead.



##########
core/server/src/partition_reconciler.rs:
##########
@@ -1137,18 +1146,68 @@ 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's counters out of its parent topic and stream,
+/// and drop its registry entry.
+///
+/// Two handles, because they can differ. The registry entry is usually gone
+/// already (the metadata apply evicts it on the commit that acked the delete),
+/// but a stale incarnation being torn down after a slab-key reuse still has
+/// one, and it must not survive into the rebuild: `partition` is a
+/// get-or-create, so the rebuild would inherit the dead incarnation's counters
+/// and its `purged_generation`. `live` is the mounted partition's own handle,
+/// which is the only place the post-commit residue lives. When both name the
+/// same `Arc` the second zeroing rolls back 0.
+fn settle_partition_stats(
+    ctx: &ReconcilerCtx,
+    ns: IggyNamespace,
+    live: Option<&Arc<iggy_common::PartitionStats>>,
+) {
+    // Registry cloned out rather than mutated inside the read closure: the
+    // rollback cascades into parent totals, which the STM read has no part in.
+    let registry = ctx
+        .shard
+        .plane
+        .metadata()
+        .mux_stm
+        .streams()
+        .read(|inner| Arc::clone(&inner.stats_registry));
+    registry.remove_partition(ns.stream_id(), ns.topic_id(), 
ns.partition_id());

Review Comment:
   warning: evicting the whole entry drops `purged_generation`, so the rebuild 
mints 0 and a still-pending purge then wipes everything appended since. seed 
the gate from the committed `Partition::purge_generation` - keeping the entry 
would break the recreate test.



##########
core/metadata/src/stm/stream.rs:
##########
@@ -2360,6 +2553,10 @@ 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.
+        self.stats_registry.rebuild_parent_totals(&self.items);

Review Comment:
   warning: `retain_from_snapshot` keeps entries by slab key with no identity 
check, so behind a donor that recycled a key the receiver keeps the dead 
occupant's stats and this rebuild makes them authoritative. compare 
`created_revision` per key first.



##########
core/server/src/http/metrics.rs:
##########
@@ -45,6 +45,16 @@ 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. Alert on any increase.
+    ///
+    /// 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.
+    stats_rollup_underflows: Counter,

Review Comment:
   warning: this only advances inside the `/metrics` scrape, which needs both 
`http.enabled` and `http.metrics.enabled`. a tcp-only or quic-only deployment 
gets no counter at all, leaving one throttled warn line as the whole signal.



##########
core/common/src/types/streaming_stats.rs:
##########
@@ -19,6 +19,86 @@ 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 keeps the total merely low, which the

Review Comment:
   warning: "the next write corrects" is wrong - increments are `fetch_add`, so 
they stack on the clamped base and the shortfall sticks until a rebuild or 
restart. this is what tells an operator whether the alarm needs a response.



##########
core/integration/tests/server/scenarios/delete_stats_rollback_scenario.rs:
##########
@@ -0,0 +1,346 @@
+// 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 7 `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 — `metadata::stm::stream` for 
the
+//! rollback and eviction, `server::partition_reconciler` for the teardown
+//! settle and the membership gate — because 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.
+//!
+//! * 
`given_counted_partition_when_apply_delete_partitions_should_roll_it_out_of_the_parents`
+//! * 
`given_counted_topic_when_apply_delete_topic_should_roll_it_out_of_the_stream`
+//! * 
`given_many_counted_partitions_when_apply_delete_topic_should_roll_all_of_them_out`
+//! * 
`given_non_zero_base_partition_ids_when_apply_delete_partitions_should_keep_the_survivor`
+//! * 
`given_replayed_delete_partitions_when_applied_twice_should_not_double_roll_back`
+//! * 
`given_topic_residue_no_partition_entry_covers_when_apply_delete_topic_should_settle_it`
+//!
+//! 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;
+use bytes::Bytes;
+use iggy::prelude::*;
+use integration::harness::{TestHarness, assert_clean_system, login_root};
+use std::str::FromStr;
+use std::time::{Duration, Instant};
+use tokio::time::sleep;
+
+// Committed partition ops fold into the shared stats on the owning shard, so a
+// read can race the apply window. 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: the consts, `create_messages` and the three validators are 
copies of `stream_size_validation_scenario.rs` - `validate_system_stats` byte 
for byte. hoist them into `scenarios/mod.rs`, which already owns the shared 
consts and helpers.



##########
core/common/src/types/streaming_stats.rs:
##########
@@ -19,6 +19,86 @@ 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 keeps the total merely low, which the
+/// next write corrects; this counter plus the `warn!` below is what tells an
+/// operator the divergence happened at all.
+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.
+#[must_use]
+pub fn rollup_underflows() -> u64 {
+    ROLLUP_UNDERFLOWS.load(Ordering::Relaxed)
+}
+
+fn report_rollup_underflow(scope: &'static str, counter: &'static str, 
shortfall: u64) {
+    let total = ROLLUP_UNDERFLOWS.fetch_add(1, Ordering::Relaxed) + 1;
+    // One line for the first, then at powers of two. A skewed tree emits up to
+    // three of these per counter per rollback, and a bulk delete would turn a
+    // diagnostic into a log flood that buries the first one. The counter above
+    // stays exact, and it is what an alert reads.
+    if total.is_power_of_two() {
+        warn!(
+            scope,
+            counter,
+            shortfall,
+            total,
+            "rollup decrement exceeded the total it was subtracted from; 
clamped at zero"
+        );
+    }
+}
+
+/// 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.
+fn clamped_sub_u64(
+    counter: &AtomicU64,
+    amount: u64,
+    scope: &'static str,
+    name: &'static str,
+) -> u64 {
+    // The closure always yields `Some`, so `Err` carries the same previous

Review Comment:
   simplification: the `unwrap_or_else` arm can't run - the closure always 
returns `Some`, so `fetch_update` never gives `Err`. four lines of comment 
justify a dead branch, twice.



##########
core/common/src/types/streaming_stats.rs:
##########
@@ -19,6 +19,86 @@ 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 keeps the total merely low, which the
+/// next write corrects; this counter plus the `warn!` below is what tells an
+/// operator the divergence happened at all.
+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.
+#[must_use]
+pub fn rollup_underflows() -> u64 {
+    ROLLUP_UNDERFLOWS.load(Ordering::Relaxed)
+}
+
+fn report_rollup_underflow(scope: &'static str, counter: &'static str, 
shortfall: u64) {
+    let total = ROLLUP_UNDERFLOWS.fetch_add(1, Ordering::Relaxed) + 1;
+    // One line for the first, then at powers of two. A skewed tree emits up to
+    // three of these per counter per rollback, and a bulk delete would turn a
+    // diagnostic into a log flood that buries the first one. The counter above
+    // stays exact, and it is what an alert reads.
+    if total.is_power_of_two() {
+        warn!(

Review Comment:
   warning: no stream, topic or partition id here, so an operator alerting on 
the counter can't locate the divergence. return the shortfall and log where the 
ids are known - retention and purge are the likeliest source.



##########
core/integration/tests/server/scenarios/delete_stats_rollback_scenario.rs:
##########
@@ -0,0 +1,346 @@
+// 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 7 `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 — `metadata::stm::stream` for 
the
+//! rollback and eviction, `server::partition_reconciler` for the teardown
+//! settle and the membership gate — because 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.
+//!
+//! * 
`given_counted_partition_when_apply_delete_partitions_should_roll_it_out_of_the_parents`

Review Comment:
   nit: this hand-kept list already misses three of the nine tests this commit 
adds (`stream.rs:3542`, `:3656`, `:3759`), and it's split from the paragraph it 
belongs to. name the module instead.



##########
core/server/src/boot/mod.rs:
##########
@@ -552,9 +552,7 @@ async fn shard_main(
             // and decrement the parent `StreamStats` by it.
             let () = recovered.mux_stm.streams().read(|inner| {

Review Comment:
   warning: this runs after the wal replay, so a replayed `DeleteTopic` over a 
torn checkpoint clamps against the stream and trips the new underflow counter 
on every such boot. zero the aggregates before replay.



##########
core/server/src/http/metrics.rs:
##########
@@ -45,6 +45,16 @@ 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

Review Comment:
   warning: "alert on any increase" isn't actionable - a retention decrement 
against a settled partition trips it, and the new unit test drives exactly that 
case. restate as a bounded rate outside delete, purge, teardown and restore 
windows.



##########
core/server/src/partition_reconciler.rs:
##########
@@ -1137,18 +1146,68 @@ 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's counters out of its parent topic and stream,
+/// and drop its registry entry.
+///
+/// Two handles, because they can differ. The registry entry is usually gone
+/// already (the metadata apply evicts it on the commit that acked the delete),
+/// but a stale incarnation being torn down after a slab-key reuse still has
+/// one, and it must not survive into the rebuild: `partition` is a
+/// get-or-create, so the rebuild would inherit the dead incarnation's counters
+/// and its `purged_generation`. `live` is the mounted partition's own handle,
+/// which is the only place the post-commit residue lives. When both name the
+/// same `Arc` the second zeroing rolls back 0.
+fn settle_partition_stats(
+    ctx: &ReconcilerCtx,
+    ns: IggyNamespace,
+    live: Option<&Arc<iggy_common::PartitionStats>>,
+) {
+    // Registry cloned out rather than mutated inside the read closure: the
+    // rollback cascades into parent totals, which the STM read has no part in.
+    let registry = ctx
+        .shard
+        .plane
+        .metadata()
+        .mux_stm
+        .streams()
+        .read(|inner| Arc::clone(&inner.stats_registry));
+    registry.remove_partition(ns.stream_id(), ns.topic_id(), 
ns.partition_id());
+    if let Some(stats) = live {
+        stats.zero_out_all();
+    }
+}
+
+/// Everything one namespace's build needs out of committed metadata: the 
shared
+/// `Arc<PartitionStats>`, the topic's runtime options, and the committed
+/// [`Partition`] record the disk loader reads `created_at` / 
`created_revision`
+/// / `created_view` from.
+///
+/// One read for all three. The `Partition` lookup doubles as the membership
+/// check: the committed topic has to still list this partition, because a pass
+/// captures its targets once and then awaits disk work per namespace, so
+/// without it a target captured before a `DeletePartitions` re-creates the
+/// entry that delete just evicted -- and the apply on the second left-right
+/// buffer, deferred to a later publish, then zeroes a live partition's 
counters
+/// and drops its entry again.
+///
+/// `None` if the stream, the topic, or the partition vanished between the
+/// target snapshot and this read.
+fn fetch_partition_build_inputs(

Review Comment:
   nit: three references still name the old `fetch_partition_stats` - 
`stream.rs:411`, `responses.rs:1385`, and the dead intra-doc link at line 1107. 
that comment also still claims stats are fetched lazily, which this merge made 
false.



##########
core/server/src/boot/mod.rs:
##########
@@ -1121,4 +1140,55 @@ mod tests {
             "ops with no partition-shape effect must stay off the broadcast"
         );
     }
+
+    /// A checkpoint loads a stream's total and its topics' as separate reads
+    /// while the partition plane keeps counting, so the two can disagree in
+    /// either direction. Boot has to drop both, or the surviving residue sits
+    /// under a total nothing on disk backs.
+    #[test]
+    fn clear_snapshot_totals_drops_both_levels_of_a_torn_checkpoint() {
+        use iggy_common::{
+            CompressionAlgorithm, IggyExpiry, IggyTimestamp, MaxTopicSize, 
StreamStats,
+        };
+        use metadata::stm::stream::{Stream, Topic};
+
+        let stream_stats = Arc::new(StreamStats::default());
+        let mut stream = Stream::with_stats(
+            Arc::from("stream"),
+            IggyTimestamp::now(),
+            Arc::clone(&stream_stats),
+        );
+        let topic = Topic::new(
+            Arc::from("topic"),
+            IggyTimestamp::now(),
+            IggyExpiry::NeverExpire,
+            CompressionAlgorithm::default(),
+            MaxTopicSize::Unlimited,
+            Arc::clone(&stream_stats),
+        );
+        let topic_stats = Arc::clone(&topic.stats);
+        stream.topics.insert(topic);
+
+        // Topics summing above their stream is the torn shape: the stream was
+        // read first, and the topic kept counting before its own read.
+        stream_stats.store_from_snapshot(100, 1, 1);
+        topic_stats.store_from_snapshot(200, 2, 2);
+        let underflows_before = iggy_common::rollup_underflows();

Review Comment:
   nit: `rollup_underflows()` is a process-global static, so this 
exact-equality assert breaks at a distance under `cargo test -p server`, where 
the whole crate shares one process. assert on the totals instead.



##########
core/server/src/boot/mod.rs:
##########
@@ -1032,6 +1030,27 @@ async fn shard_main(
 /// and dropped (the periodic tick recovers). Installed via
 /// [`metadata::IggyMetadata::set_commit_notifier`] on shard 0 only, the
 /// sole writer of the metadata state machine.

Review Comment:
   nit: `clear_snapshot_totals` landed between this doc block and the fn it 
describes, so both blocks now attach to the stats helper and 
`make_metadata_commit_notifier` ships undocumented. move it above line 1024 or 
below the notifier.



##########
core/server/src/partition_reconciler.rs:
##########
@@ -1137,18 +1146,68 @@ 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's counters out of its parent topic and stream,
+/// and drop its registry entry.
+///
+/// Two handles, because they can differ. The registry entry is usually gone
+/// already (the metadata apply evicts it on the commit that acked the delete),
+/// but a stale incarnation being torn down after a slab-key reuse still has
+/// one, and it must not survive into the rebuild: `partition` is a
+/// get-or-create, so the rebuild would inherit the dead incarnation's counters
+/// and its `purged_generation`. `live` is the mounted partition's own handle,
+/// which is the only place the post-commit residue lives. When both name the
+/// same `Arc` the second zeroing rolls back 0.
+fn settle_partition_stats(
+    ctx: &ReconcilerCtx,
+    ns: IggyNamespace,
+    live: Option<&Arc<iggy_common::PartitionStats>>,

Review Comment:
   nit: `&Arc<PartitionStats>` buys nothing over `&PartitionStats` here - the 
body only calls `zero_out_all()`. the call site becomes `live_stats.as_deref()`.



##########
core/common/src/types/streaming_stats.rs:
##########
@@ -19,6 +19,86 @@ 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 keeps the total merely low, which the
+/// next write corrects; this counter plus the `warn!` below is what tells an
+/// operator the divergence happened at all.
+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.
+#[must_use]
+pub fn rollup_underflows() -> u64 {

Review Comment:
   nit: this reaches the root of the published `iggy_common` through the glob 
re-export, so a server-only diagnostic becomes semver-locked api. 
`#[doc(hidden)]` it.



##########
core/metadata/src/stm/stream.rs:
##########
@@ -384,10 +389,29 @@ 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
+///   current batch's `absorb_first` (0.11.8 `write.rs`), which is an
+///   implementation detail, not a contract.
+///
+/// Independently of both, `fetch_partition_stats` refuses to register a

Review Comment:
   warning: this doesn't close the live route - the check reads the left-right 
read side, but the eviction happens in `absorb_first` before `publish()`. a 
reader in that window recreates the entry the apply just dropped.



##########
core/server/src/boot/mod.rs:
##########
@@ -1121,4 +1140,55 @@ mod tests {
             "ops with no partition-shape effect must stay off the broadcast"
         );
     }
+
+    /// A checkpoint loads a stream's total and its topics' as separate reads
+    /// while the partition plane keeps counting, so the two can disagree in
+    /// either direction. Boot has to drop both, or the surviving residue sits
+    /// under a total nothing on disk backs.
+    #[test]
+    fn clear_snapshot_totals_drops_both_levels_of_a_torn_checkpoint() {
+        use iggy_common::{

Review Comment:
   nit: imports go at the top of the module - the sibling test above has none.



##########
core/common/src/types/streaming_stats.rs:
##########
@@ -19,6 +19,86 @@ 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 keeps the total merely low, which the
+/// next write corrects; this counter plus the `warn!` below is what tells an
+/// operator the divergence happened at all.
+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.
+#[must_use]
+pub fn rollup_underflows() -> u64 {
+    ROLLUP_UNDERFLOWS.load(Ordering::Relaxed)
+}
+
+fn report_rollup_underflow(scope: &'static str, counter: &'static str, 
shortfall: u64) {
+    let total = ROLLUP_UNDERFLOWS.fetch_add(1, Ordering::Relaxed) + 1;
+    // One line for the first, then at powers of two. A skewed tree emits up to
+    // three of these per counter per rollback, and a bulk delete would turn a
+    // diagnostic into a log flood that buries the first one. The counter above
+    // stays exact, and it is what an alert reads.
+    if total.is_power_of_two() {

Review Comment:
   nit: the throttle keys on one global count across all six scope/counter 
pairs, so a first-ever divergence in a fresh scope stays silent for up to 1023 
events. key it per pair, or on a time window.



##########
core/server/src/partition_reconciler.rs:
##########
@@ -1137,18 +1146,68 @@ 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's counters out of its parent topic and stream,
+/// and drop its registry entry.
+///
+/// Two handles, because they can differ. The registry entry is usually gone
+/// already (the metadata apply evicts it on the commit that acked the delete),
+/// but a stale incarnation being torn down after a slab-key reuse still has
+/// one, and it must not survive into the rebuild: `partition` is a
+/// get-or-create, so the rebuild would inherit the dead incarnation's counters
+/// and its `purged_generation`. `live` is the mounted partition's own handle,
+/// which is the only place the post-commit residue lives. When both name the
+/// same `Arc` the second zeroing rolls back 0.
+fn settle_partition_stats(
+    ctx: &ReconcilerCtx,
+    ns: IggyNamespace,
+    live: Option<&Arc<iggy_common::PartitionStats>>,
+) {
+    // Registry cloned out rather than mutated inside the read closure: the
+    // rollback cascades into parent totals, which the STM read has no part in.
+    let registry = ctx
+        .shard
+        .plane
+        .metadata()
+        .mux_stm
+        .streams()
+        .read(|inner| Arc::clone(&inner.stats_registry));
+    registry.remove_partition(ns.stream_id(), ns.topic_id(), 
ns.partition_id());
+    if let Some(stats) = live {
+        stats.zero_out_all();
+    }
+}
+
+/// Everything one namespace's build needs out of committed metadata: the 
shared
+/// `Arc<PartitionStats>`, the topic's runtime options, and the committed
+/// [`Partition`] record the disk loader reads `created_at` / 
`created_revision`
+/// / `created_view` from.
+///
+/// One read for all three. The `Partition` lookup doubles as the membership
+/// check: the committed topic has to still list this partition, because a pass
+/// captures its targets once and then awaits disk work per namespace, so
+/// without it a target captured before a `DeletePartitions` re-creates the
+/// entry that delete just evicted -- and the apply on the second left-right
+/// buffer, deferred to a later publish, then zeroes a live partition's 
counters
+/// and drops its entry again.
+///
+/// `None` if the stream, the topic, or the partition vanished between the
+/// target snapshot and this read.
+fn fetch_partition_build_inputs(
     ctx: &ReconcilerCtx,
     ns: IggyNamespace,
 ) -> Option<(
     Arc<iggy_common::PartitionStats>,
     iggy_common::TopicRuntimeOptions,
+    Partition,
 )> {
     ctx.shard.plane.metadata().mux_stm.streams().read(|inner| {
         let stream = inner.items.get(ns.stream_id())?;
         let topic = stream.topics.get(ns.topic_id())?;
+        let partition = topic
+            .partitions
+            .iter()
+            .find(|partition| partition.id == ns.partition_id())?

Review Comment:
   nit: this `find` now runs on every build, not just the disk-load branch, so 
materialising a k-partition topic is quadratic under the read guard. cheap next 
to the mkdir, but `iggy_partitions.rs:365` is the same shape.



##########
core/server/src/partition_reconciler.rs:
##########
@@ -928,10 +926,21 @@ async fn tear_down_owned_partition(
     // on_replicate / on_ack frames that haven't observed the queued
     // tombstone yet. Idempotent on retry: already-tombstoned namespace
     // stays tombstoned; already-removed shards_table row is a no-op.
+    // Take the counters handle BEFORE the tombstone (every partition accessor
+    // is tombstone-gated) and settle it AFTER, once the tombstone has stopped
+    // new frames from resolving the partition. The metadata apply already
+    // rolled it out at commit time, so what is left here is whatever landed in
+    // the window between the two: in-flight appends adding to parents that
+    // outlive the partition, and a retention tick subtracting from counters 
the
+    // apply had already zeroed. A frame already past the gate can still land
+    // after this settle, which is why the rollback clamps rather than trusting
+    // the amount it is handed.
+    let live_stats = partitions.with_partition(&ns, |partition| 
Arc::clone(&partition.stats));

Review Comment:
   nit: "take the counters handle BEFORE the tombstone" reads like the retry 
path is covered, but `with_partition` is tombstone-gated, so on retry this is 
always `None` and the settle is a no-op. say so.



##########
core/server/src/partition_reconciler.rs:
##########
@@ -1817,6 +1862,96 @@ 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() {

Review Comment:
   nit: this still passes with the `registry.remove_partition` at line 1174 
deleted, because the metadata apply already evicted the entry in this fixture. 
the stale-incarnation eviction has no test.



##########
core/server/src/partition_reconciler.rs:
##########
@@ -1137,18 +1146,68 @@ 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's counters out of its parent topic and stream,
+/// and drop its registry entry.
+///
+/// Two handles, because they can differ. The registry entry is usually gone
+/// already (the metadata apply evicts it on the commit that acked the delete),
+/// but a stale incarnation being torn down after a slab-key reuse still has
+/// one, and it must not survive into the rebuild: `partition` is a
+/// get-or-create, so the rebuild would inherit the dead incarnation's counters
+/// and its `purged_generation`. `live` is the mounted partition's own handle,
+/// which is the only place the post-commit residue lives. When both name the
+/// same `Arc` the second zeroing rolls back 0.
+fn settle_partition_stats(
+    ctx: &ReconcilerCtx,
+    ns: IggyNamespace,
+    live: Option<&Arc<iggy_common::PartitionStats>>,
+) {
+    // Registry cloned out rather than mutated inside the read closure: the
+    // rollback cascades into parent totals, which the STM read has no part in.
+    let registry = ctx

Review Comment:
   simplification: one full `streams().read()` per torn-down namespace just to 
clone the registry `Arc`. build it once into `ReconcilerCtx` - it's never 
re-minted after boot, and a topic delete opens one reader epoch per partition.



##########
core/metadata/src/stm/stream.rs:
##########
@@ -550,23 +578,109 @@ impl StatsRegistry {
     }
 
     fn remove_topic(&self, stream_id: usize, topic_id: usize) {
-        self.topics
+        // Partitions first: each one's rollback cascades through its parent
+        // topic into the stream, so zeroing the topic ahead of them would
+        // subtract the same bytes from the stream twice.
+        self.remove_partitions_where(|(sid, tid, _)| *sid == stream_id && *tid 
== topic_id);

Review Comment:
   simplification: this still uses the predicate sweep the comment at line 608 
argues against, and `remove_stream` does too. the ids are in hand at `:2234`, 
and every registry insert takes them from the committed topic, so keying is 
safe.



##########
core/integration/tests/server/scenarios/delete_stats_rollback_scenario.rs:
##########
@@ -0,0 +1,346 @@
+// 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

Review Comment:
   nit: the literal "7 unit tests fail" is a number nothing maintains - it goes 
stale on the next test added to that module. keep the conclusion, drop the 
count.



##########
core/integration/tests/server/scenarios/delete_stats_rollback_scenario.rs:
##########
@@ -0,0 +1,346 @@
+// 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 7 `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 — `metadata::stm::stream` for 
the
+//! rollback and eviction, `server::partition_reconciler` for the teardown
+//! settle and the membership gate — because 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.
+//!
+//! * 
`given_counted_partition_when_apply_delete_partitions_should_roll_it_out_of_the_parents`
+//! * 
`given_counted_topic_when_apply_delete_topic_should_roll_it_out_of_the_stream`
+//! * 
`given_many_counted_partitions_when_apply_delete_topic_should_roll_all_of_them_out`
+//! * 
`given_non_zero_base_partition_ids_when_apply_delete_partitions_should_keep_the_survivor`
+//! * 
`given_replayed_delete_partitions_when_applied_twice_should_not_double_roll_back`
+//! * 
`given_topic_residue_no_partition_entry_covers_when_apply_delete_topic_should_settle_it`
+//!
+//! 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;
+use bytes::Bytes;
+use iggy::prelude::*;
+use integration::harness::{TestHarness, assert_clean_system, login_root};
+use std::str::FromStr;
+use std::time::{Duration, Instant};
+use tokio::time::sleep;
+
+// Committed partition ops fold into the shared stats on the owning shard, so a
+// read can race the apply window. 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 STREAM_NAME: &str = "delete-stats-stream";
+const KEPT_TOPIC: &str = "kept-topic";
+const DELETED_TOPIC: &str = "deleted-topic";
+const MESSAGE_PAYLOAD_SIZE_BYTES: u64 = 57;
+const MSGS_COUNT: u64 = 17;
+
+// The server accounts the on-disk batch framing: one command header per append
+// pass plus a per-message header. Mirrors `stream_size_validation_scenario`.
+const NG_BATCH_HEADER_SIZE: u64 = 256;
+const NG_MESSAGE_HEADER_SIZE: u64 = 48;
+const MSGS_SIZE: u64 =
+    NG_BATCH_HEADER_SIZE + (NG_MESSAGE_HEADER_SIZE + 
MESSAGE_PAYLOAD_SIZE_BYTES) * MSGS_COUNT;
+
+pub async fn run(harness: &TestHarness) {
+    let client = harness
+        .new_client()
+        .await
+        .expect("Failed to create new client");
+    client.ping().await.unwrap();
+    login_root(&client).await.expect("login failed");
+
+    // Partition half first: its assertions are the tighter pair.
+    delete_partitions_rolls_back_topic_and_stream(&client).await;
+    delete_topic_rolls_back_the_stream(&client).await;
+
+    assert_clean_system(&client).await;
+}
+
+/// Deleting partitions that still hold messages must take their bytes out of
+/// both the topic and the stream. The retained partition's messages must stay
+/// counted, so a blanket zeroing fails here as loudly as no rollback at all.
+async fn delete_partitions_rolls_back_topic_and_stream(client: &IggyClient) {
+    create_stream(client, STREAM_NAME).await;
+    create_topic(client, STREAM_NAME, KEPT_TOPIC).await;
+
+    // One pass per partition, so the delete below removes a known share.
+    for partition_id in 0..PARTITIONS_COUNT {
+        send_one_pass(client, STREAM_NAME, KEPT_TOPIC, partition_id).await;
+    }
+    let all_partitions = MSGS_SIZE * u64::from(PARTITIONS_COUNT);
+    let all_messages = MSGS_COUNT * u64::from(PARTITIONS_COUNT);
+    validate_topic(
+        client,
+        STREAM_NAME,
+        KEPT_TOPIC,
+        all_partitions,
+        all_messages,
+    )
+    .await;
+    validate_stream(client, STREAM_NAME, all_partitions, all_messages).await;
+
+    client
+        .delete_partitions(
+            &Identifier::from_str(STREAM_NAME).unwrap(),
+            &Identifier::from_str(KEPT_TOPIC).unwrap(),
+            1,
+        )
+        .await
+        .unwrap();
+
+    // Single-shot again: the retained partitions' messages must still be
+    // counted, so this fails both on a missing rollback and on a blanket
+    // zeroing of the parents.
+    let retained = u64::from(PARTITIONS_COUNT - 1);
+    let topic = client
+        .get_topic(
+            &Identifier::from_str(STREAM_NAME).unwrap(),
+            &Identifier::from_str(KEPT_TOPIC).unwrap(),
+        )
+        .await
+        .unwrap()
+        .expect("Failed to get topic");
+    assert_eq!(
+        topic.size,
+        MSGS_SIZE * retained,
+        "the deleted partitions' bytes must leave the topic total on the 
delete's commit"
+    );
+    assert_eq!(topic.messages_count, MSGS_COUNT * retained);
+    let stream = client
+        .get_stream(&Identifier::from_str(STREAM_NAME).unwrap())
+        .await
+        .unwrap()
+        .expect("Failed to get stream");
+    assert_eq!(
+        stream.size,
+        MSGS_SIZE * retained,
+        "the deleted partitions' bytes must leave the stream total too"
+    );
+    assert_eq!(stream.messages_count, MSGS_COUNT * retained);
+
+    client
+        .delete_stream(&Identifier::from_str(STREAM_NAME).unwrap())
+        .await
+        .unwrap();
+}

Review Comment:
   nit: missing blank line before the next doc comment - every other item 
boundary in this file has one.



##########
core/common/src/types/streaming_stats.rs:
##########
@@ -19,6 +19,86 @@ 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 keeps the total merely low, which the
+/// next write corrects; this counter plus the `warn!` below is what tells an
+/// operator the divergence happened at all.
+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.
+#[must_use]
+pub fn rollup_underflows() -> u64 {
+    ROLLUP_UNDERFLOWS.load(Ordering::Relaxed)
+}
+
+fn report_rollup_underflow(scope: &'static str, counter: &'static str, 
shortfall: u64) {
+    let total = ROLLUP_UNDERFLOWS.fetch_add(1, Ordering::Relaxed) + 1;
+    // One line for the first, then at powers of two. A skewed tree emits up to
+    // three of these per counter per rollback, and a bulk delete would turn a
+    // diagnostic into a log flood that buries the first one. The counter above
+    // stays exact, and it is what an alert reads.
+    if total.is_power_of_two() {
+        warn!(
+            scope,
+            counter,
+            shortfall,
+            total,
+            "rollup decrement exceeded the total it was subtracted from; 
clamped at zero"
+        );
+    }
+}
+
+/// 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.
+fn clamped_sub_u64(

Review Comment:
   simplification: `clamped_sub_u64` and `clamped_sub_u32` are the same 16 
lines twice, differing only in the atomic type. one `macro_rules!` covers both 
- a sealed trait would be longer than the duplication.



##########
core/metadata/src/stm/stream.rs:
##########
@@ -550,23 +578,109 @@ impl StatsRegistry {
     }
 
     fn remove_topic(&self, stream_id: usize, topic_id: usize) {
-        self.topics
+        // Partitions first: each one's rollback cascades through its parent
+        // topic into the stream, so zeroing the topic ahead of them would
+        // subtract the same bytes from the stream twice.
+        self.remove_partitions_where(|(sid, tid, _)| *sid == stream_id && *tid 
== topic_id);
+        let topic = self
+            .topics
             .lock()
             .expect("stats registry mutex poisoned")
             .remove(&(stream_id, topic_id));
-        self.partitions
-            .lock()
-            .expect("stats registry mutex poisoned")
-            .retain(|(sid, tid, _), _| !(*sid == stream_id && *tid == 
topic_id));
+        // Whatever the topic still counts after the loop above is what its
+        // partitions did not account for: a snapshot restore that stored a
+        // total this node's partition entries never contributed, or bytes a
+        // partition kept adding after its own entry was evicted. Swapping it
+        // out settles that residue on the stream instead of stranding it.
+        if let Some(topic) = topic {
+            topic.zero_out_all();
+        }
     }
 
-    fn remove_partitions_from(&self, stream_id: usize, topic_id: usize, 
first_removed: usize) {
-        self.partitions
-            .lock()
-            .expect("stats registry mutex poisoned")
-            .retain(|(sid, tid, pid), _| {
-                !(*sid == stream_id && *tid == topic_id && *pid >= 
first_removed)
+    /// Roll the named partitions out of their parents and drop their entries.
+    ///
+    /// Ids, never positions. `DeletePartitions` truncates the tail of a `Vec`,
+    /// and a count-based predicate (`id >= retained`) only picks the same set
+    /// while ids happen to be dense; when they are not it evicts and zeroes a
+    /// SURVIVING partition, which strips its bytes from the parents and resets
+    /// the offset its clients store against.
+    fn remove_partitions(&self, stream_id: usize, topic_id: usize, 
removed_ids: &[usize]) {
+        // Keyed removes, not a predicate: the ids are known, and a predicate
+        // makes the map walk every entry on the node to drop a handful of 
them.
+        // Bulk deletes call this per topic, so that walk squares.
+        let dropped: Vec<Arc<PartitionStats>> = {
+            let mut entries = self
+                .partitions
+                .lock()
+                .expect("stats registry mutex poisoned");
+            removed_ids
+                .iter()
+                .filter_map(|partition_id| entries.remove(&(stream_id, 
topic_id, *partition_id)))
+                .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();
+        }
+    }
+
+    /// Drop every partition entry `should_drop` selects, rolling each one's
+    /// counters out of its parent topic and stream on the way.
+    ///
+    /// Both halves are load-bearing. A partition reports by incrementing its
+    /// parents through the `Arc`, so evicting alone strands what it 
contributed
+    /// in the topic and stream totals. And zeroing alone leaves an entry that
+    /// outlives its ids: a topic's slab key is recycled by the next
+    /// `create_topic`, and `DeletePartitions` truncates the tail so the next
+    /// `CreatePartitions` mints the freed ids again. The survivor would hand
+    /// its successor both its counters and its `purged_generation`, and that
+    /// stale generation makes the next purge's gate skip the reset.
+    ///
+    /// # Panics
+    /// If the registry mutex is poisoned.
+    fn remove_partitions_where(&self, should_drop: impl Fn(&(usize, usize, 
usize)) -> bool) {

Review Comment:
   simplification: generic over a predicate with exactly one caller. inline it 
as `remove_topic_partitions(stream_id, topic_id)` and drop the 
single-instantiation generic.



##########
core/integration/tests/server/scenarios/delete_stats_rollback_scenario.rs:
##########
@@ -0,0 +1,346 @@
+// 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 7 `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 — `metadata::stm::stream` for 
the
+//! rollback and eviction, `server::partition_reconciler` for the teardown
+//! settle and the membership gate — because 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.
+//!
+//! * 
`given_counted_partition_when_apply_delete_partitions_should_roll_it_out_of_the_parents`
+//! * 
`given_counted_topic_when_apply_delete_topic_should_roll_it_out_of_the_stream`
+//! * 
`given_many_counted_partitions_when_apply_delete_topic_should_roll_all_of_them_out`
+//! * 
`given_non_zero_base_partition_ids_when_apply_delete_partitions_should_keep_the_survivor`
+//! * 
`given_replayed_delete_partitions_when_applied_twice_should_not_double_roll_back`
+//! * 
`given_topic_residue_no_partition_entry_covers_when_apply_delete_topic_should_settle_it`
+//!
+//! 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;
+use bytes::Bytes;
+use iggy::prelude::*;
+use integration::harness::{TestHarness, assert_clean_system, login_root};
+use std::str::FromStr;

Review Comment:
   nit: `Identifier::named` is what the other scenarios use - `from_str` 
reinterprets a numeric-looking name as an id. no live bug with these names, 
just the wrong constructor.
   
   also at the other 15 call sites in this file.



##########
core/integration/tests/server/scenarios/delete_stats_rollback_scenario.rs:
##########
@@ -0,0 +1,346 @@
+// 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 7 `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 — `metadata::stm::stream` for 
the
+//! rollback and eviction, `server::partition_reconciler` for the teardown
+//! settle and the membership gate — because 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.
+//!
+//! * 
`given_counted_partition_when_apply_delete_partitions_should_roll_it_out_of_the_parents`
+//! * 
`given_counted_topic_when_apply_delete_topic_should_roll_it_out_of_the_stream`
+//! * 
`given_many_counted_partitions_when_apply_delete_topic_should_roll_all_of_them_out`
+//! * 
`given_non_zero_base_partition_ids_when_apply_delete_partitions_should_keep_the_survivor`
+//! * 
`given_replayed_delete_partitions_when_applied_twice_should_not_double_roll_back`
+//! * 
`given_topic_residue_no_partition_entry_covers_when_apply_delete_topic_should_settle_it`
+//!
+//! 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;
+use bytes::Bytes;
+use iggy::prelude::*;
+use integration::harness::{TestHarness, assert_clean_system, login_root};
+use std::str::FromStr;
+use std::time::{Duration, Instant};
+use tokio::time::sleep;
+
+// Committed partition ops fold into the shared stats on the owning shard, so a
+// read can race the apply window. 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 STREAM_NAME: &str = "delete-stats-stream";
+const KEPT_TOPIC: &str = "kept-topic";
+const DELETED_TOPIC: &str = "deleted-topic";
+const MESSAGE_PAYLOAD_SIZE_BYTES: u64 = 57;
+const MSGS_COUNT: u64 = 17;
+
+// The server accounts the on-disk batch framing: one command header per append
+// pass plus a per-message header. Mirrors `stream_size_validation_scenario`.
+const NG_BATCH_HEADER_SIZE: u64 = 256;
+const NG_MESSAGE_HEADER_SIZE: u64 = 48;
+const MSGS_SIZE: u64 =
+    NG_BATCH_HEADER_SIZE + (NG_MESSAGE_HEADER_SIZE + 
MESSAGE_PAYLOAD_SIZE_BYTES) * MSGS_COUNT;
+
+pub async fn run(harness: &TestHarness) {
+    let client = harness

Review Comment:
   simplification: this re-implements `scenarios::create_client` verbatim, 
panic string included. import it instead.



##########
core/integration/tests/server/scenarios/delete_stats_rollback_scenario.rs:
##########
@@ -0,0 +1,346 @@
+// 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 7 `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 — `metadata::stm::stream` for 
the
+//! rollback and eviction, `server::partition_reconciler` for the teardown
+//! settle and the membership gate — because 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.
+//!
+//! * 
`given_counted_partition_when_apply_delete_partitions_should_roll_it_out_of_the_parents`
+//! * 
`given_counted_topic_when_apply_delete_topic_should_roll_it_out_of_the_stream`
+//! * 
`given_many_counted_partitions_when_apply_delete_topic_should_roll_all_of_them_out`
+//! * 
`given_non_zero_base_partition_ids_when_apply_delete_partitions_should_keep_the_survivor`
+//! * 
`given_replayed_delete_partitions_when_applied_twice_should_not_double_roll_back`
+//! * 
`given_topic_residue_no_partition_entry_covers_when_apply_delete_topic_should_settle_it`
+//!
+//! 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;
+use bytes::Bytes;
+use iggy::prelude::*;
+use integration::harness::{TestHarness, assert_clean_system, login_root};
+use std::str::FromStr;
+use std::time::{Duration, Instant};
+use tokio::time::sleep;
+
+// Committed partition ops fold into the shared stats on the owning shard, so a
+// read can race the apply window. 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 STREAM_NAME: &str = "delete-stats-stream";
+const KEPT_TOPIC: &str = "kept-topic";
+const DELETED_TOPIC: &str = "deleted-topic";
+const MESSAGE_PAYLOAD_SIZE_BYTES: u64 = 57;
+const MSGS_COUNT: u64 = 17;
+
+// The server accounts the on-disk batch framing: one command header per append
+// pass plus a per-message header. Mirrors `stream_size_validation_scenario`.
+const NG_BATCH_HEADER_SIZE: u64 = 256;
+const NG_MESSAGE_HEADER_SIZE: u64 = 48;
+const MSGS_SIZE: u64 =
+    NG_BATCH_HEADER_SIZE + (NG_MESSAGE_HEADER_SIZE + 
MESSAGE_PAYLOAD_SIZE_BYTES) * MSGS_COUNT;
+
+pub async fn run(harness: &TestHarness) {
+    let client = harness
+        .new_client()
+        .await
+        .expect("Failed to create new client");
+    client.ping().await.unwrap();
+    login_root(&client).await.expect("login failed");
+
+    // Partition half first: its assertions are the tighter pair.
+    delete_partitions_rolls_back_topic_and_stream(&client).await;
+    delete_topic_rolls_back_the_stream(&client).await;
+
+    assert_clean_system(&client).await;
+}
+
+/// Deleting partitions that still hold messages must take their bytes out of
+/// both the topic and the stream. The retained partition's messages must stay
+/// counted, so a blanket zeroing fails here as loudly as no rollback at all.
+async fn delete_partitions_rolls_back_topic_and_stream(client: &IggyClient) {
+    create_stream(client, STREAM_NAME).await;
+    create_topic(client, STREAM_NAME, KEPT_TOPIC).await;
+
+    // One pass per partition, so the delete below removes a known share.
+    for partition_id in 0..PARTITIONS_COUNT {
+        send_one_pass(client, STREAM_NAME, KEPT_TOPIC, partition_id).await;
+    }
+    let all_partitions = MSGS_SIZE * u64::from(PARTITIONS_COUNT);
+    let all_messages = MSGS_COUNT * u64::from(PARTITIONS_COUNT);
+    validate_topic(
+        client,
+        STREAM_NAME,
+        KEPT_TOPIC,
+        all_partitions,
+        all_messages,
+    )
+    .await;
+    validate_stream(client, STREAM_NAME, all_partitions, all_messages).await;
+
+    client
+        .delete_partitions(
+            &Identifier::from_str(STREAM_NAME).unwrap(),
+            &Identifier::from_str(KEPT_TOPIC).unwrap(),
+            1,
+        )
+        .await
+        .unwrap();
+
+    // Single-shot again: the retained partitions' messages must still be
+    // counted, so this fails both on a missing rollback and on a blanket
+    // zeroing of the parents.
+    let retained = u64::from(PARTITIONS_COUNT - 1);
+    let topic = client
+        .get_topic(
+            &Identifier::from_str(STREAM_NAME).unwrap(),
+            &Identifier::from_str(KEPT_TOPIC).unwrap(),
+        )
+        .await
+        .unwrap()
+        .expect("Failed to get topic");
+    assert_eq!(
+        topic.size,
+        MSGS_SIZE * retained,
+        "the deleted partitions' bytes must leave the topic total on the 
delete's commit"
+    );
+    assert_eq!(topic.messages_count, MSGS_COUNT * retained);
+    let stream = client
+        .get_stream(&Identifier::from_str(STREAM_NAME).unwrap())
+        .await
+        .unwrap()
+        .expect("Failed to get stream");
+    assert_eq!(
+        stream.size,
+        MSGS_SIZE * retained,
+        "the deleted partitions' bytes must leave the stream total too"
+    );
+    assert_eq!(stream.messages_count, MSGS_COUNT * retained);
+
+    client
+        .delete_stream(&Identifier::from_str(STREAM_NAME).unwrap())
+        .await
+        .unwrap();
+}
+/// Deleting a topic that still holds messages must take its bytes out of the
+/// stream total, not just remove the topic.
+async fn delete_topic_rolls_back_the_stream(client: &IggyClient) {
+    create_stream(client, STREAM_NAME).await;
+    create_topic(client, STREAM_NAME, KEPT_TOPIC).await;
+    create_topic(client, STREAM_NAME, DELETED_TOPIC).await;
+
+    send_one_pass(client, STREAM_NAME, KEPT_TOPIC, 0).await;
+    send_one_pass(client, STREAM_NAME, DELETED_TOPIC, 0).await;
+    validate_stream(client, STREAM_NAME, MSGS_SIZE * 2, MSGS_COUNT * 2).await;
+
+    client
+        .delete_topic(
+            &Identifier::from_str(STREAM_NAME).unwrap(),
+            &Identifier::from_str(DELETED_TOPIC).unwrap(),
+        )
+        .await
+        .unwrap();
+
+    // Read once, with no convergence retry: the delete acks on the metadata
+    // commit, so the totals must already be right in the reply the client is
+    // holding. A retry here would wait for the reconciler's on-disk wipe and
+    // hide the window entirely.
+    let stream = client
+        .get_stream(&Identifier::from_str(STREAM_NAME).unwrap())
+        .await
+        .unwrap()
+        .expect("Failed to get stream");
+    assert_eq!(
+        stream.size, MSGS_SIZE,
+        "the deleted topic's bytes must leave the stream total on the commit 
that acked the delete"
+    );
+    assert_eq!(stream.messages_count, MSGS_COUNT);
+    validate_topic(client, STREAM_NAME, KEPT_TOPIC, MSGS_SIZE, 
MSGS_COUNT).await;
+    validate_system_stats(client, MSGS_SIZE, MSGS_COUNT).await;
+
+    client
+        .delete_stream(&Identifier::from_str(STREAM_NAME).unwrap())
+        .await
+        .unwrap();
+}
+
+async fn create_stream(client: &IggyClient, stream_name: &str) {

Review Comment:
   simplification: one-line passthrough with two call sites. inline it.



##########
core/metadata/src/stm/stream.rs:
##########
@@ -550,23 +578,109 @@ impl StatsRegistry {
     }
 
     fn remove_topic(&self, stream_id: usize, topic_id: usize) {
-        self.topics
+        // Partitions first: each one's rollback cascades through its parent
+        // topic into the stream, so zeroing the topic ahead of them would
+        // subtract the same bytes from the stream twice.
+        self.remove_partitions_where(|(sid, tid, _)| *sid == stream_id && *tid 
== topic_id);
+        let topic = self
+            .topics
             .lock()
             .expect("stats registry mutex poisoned")
             .remove(&(stream_id, topic_id));
-        self.partitions
-            .lock()
-            .expect("stats registry mutex poisoned")
-            .retain(|(sid, tid, _), _| !(*sid == stream_id && *tid == 
topic_id));
+        // Whatever the topic still counts after the loop above is what its
+        // partitions did not account for: a snapshot restore that stored a
+        // total this node's partition entries never contributed, or bytes a
+        // partition kept adding after its own entry was evicted. Swapping it
+        // out settles that residue on the stream instead of stranding it.
+        if let Some(topic) = topic {
+            topic.zero_out_all();
+        }
     }
 
-    fn remove_partitions_from(&self, stream_id: usize, topic_id: usize, 
first_removed: usize) {
-        self.partitions
-            .lock()
-            .expect("stats registry mutex poisoned")
-            .retain(|(sid, tid, pid), _| {
-                !(*sid == stream_id && *tid == topic_id && *pid >= 
first_removed)
+    /// Roll the named partitions out of their parents and drop their entries.
+    ///
+    /// Ids, never positions. `DeletePartitions` truncates the tail of a `Vec`,
+    /// and a count-based predicate (`id >= retained`) only picks the same set
+    /// while ids happen to be dense; when they are not it evicts and zeroes a
+    /// SURVIVING partition, which strips its bytes from the parents and resets
+    /// the offset its clients store against.
+    fn remove_partitions(&self, stream_id: usize, topic_id: usize, 
removed_ids: &[usize]) {
+        // Keyed removes, not a predicate: the ids are known, and a predicate
+        // makes the map walk every entry on the node to drop a handful of 
them.
+        // Bulk deletes call this per topic, so that walk squares.
+        let dropped: Vec<Arc<PartitionStats>> = {
+            let mut entries = self
+                .partitions
+                .lock()
+                .expect("stats registry mutex poisoned");
+            removed_ids
+                .iter()
+                .filter_map(|partition_id| entries.remove(&(stream_id, 
topic_id, *partition_id)))
+                .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();
+        }
+    }
+
+    /// Drop every partition entry `should_drop` selects, rolling each one's
+    /// counters out of its parent topic and stream on the way.
+    ///
+    /// Both halves are load-bearing. A partition reports by incrementing its
+    /// parents through the `Arc`, so evicting alone strands what it 
contributed
+    /// in the topic and stream totals. And zeroing alone leaves an entry that
+    /// outlives its ids: a topic's slab key is recycled by the next
+    /// `create_topic`, and `DeletePartitions` truncates the tail so the next
+    /// `CreatePartitions` mints the freed ids again. The survivor would hand
+    /// its successor both its counters and its `purged_generation`, and that
+    /// stale generation makes the next purge's gate skip the reset.
+    ///
+    /// # Panics
+    /// If the registry mutex is poisoned.
+    fn remove_partitions_where(&self, should_drop: impl Fn(&(usize, usize, 
usize)) -> bool) {
+        let dropped: Vec<Arc<PartitionStats>> = {
+            let mut entries = self
+                .partitions
+                .lock()
+                .expect("stats registry mutex poisoned");
+            let mut dropped = Vec::new();
+            entries.retain(|key, entry| {
+                if should_drop(key) {
+                    dropped.push(Arc::clone(&entry.stats));
+                    return false;
+                }
+                true
             });
+            dropped
+        };
+        // 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();
+        }
+    }
+
+    /// Roll one partition out of its parents and drop its entry.
+    ///
+    /// For the reconciler's owned-partition teardown. It tombstones the
+    /// namespace first, which stops any NEW frame from resolving the 
partition;
+    /// a frame already past that gate can still finish and touch the counters
+    /// afterwards, so this settles the bulk rather than closing the window. 
The
+    /// clamped rollback keeps that remainder from being taken out of a
+    /// sibling's data.
+    ///
+    /// The metadata apply has usually evicted the entry already, in which case
+    /// this is a no-op and the caller's own handle carries whatever landed
+    /// after; a stale incarnation being torn down after a slab-key reuse still
+    /// has its entry, and it must not survive into the rebuild.
+    ///
+    /// # Panics
+    /// If the registry mutex is poisoned.
+    pub fn remove_partition(&self, stream_id: usize, topic_id: usize, 
partition_id: usize) {

Review Comment:
   simplification: three-line wrapper over `remove_partitions` with one caller. 
call it directly with a one-element slice - `remove_partitions` has to go 
`pub`, so the public item count is unchanged.



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