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.