numinnex commented on code in PR #4046:
URL: https://github.com/apache/iggy/pull/4046#discussion_r3955158489
##########
core/shard/src/lib.rs:
##########
@@ -2537,6 +2537,17 @@ where
// the floor with the partition value.
self.metrics
.record_partition_prepare_gap_drops(partition.take_prepare_gap_drops());
+ // Roll whatever this partition still counts out of its
+ // parent topic and stream. Here and not in the
+ // reconciler's teardown: a handler suspended
mid-append
+ // holds its own `Arc` past the tombstone, and an
earlier
+ // settle leaves its increment in the parents with the
+ // partition already gone. This is the drop point, so
+ // nothing can add through that handle afterwards. The
+ // rollback clamps, so the usual case -- the metadata
+ // apply already zeroed these counters at commit --
takes
+ // nothing.
+ partition.stats.zero_out_all();
Review Comment:
Leaving this one, with a reason rather than a fix.
The discarded build shares the namespace's `PartitionStats` Arc with the
live incarnation, so neither arm can roll back what it folded: zeroing there is
exactly the clobber that already cost a live partition its `current_offset`
once.
Folding at adoption is the right shape, but it means the loader and the
builder stop touching the shared Arc at all and hand `insert` a local delta
instead. `load_partition` counts as it walks the chain, and
`load_partition_or_fence` zeroes the shared Arc on a refusal, so that is a
boot-path change rather than a local one. Both arms are unreachable today and
already log plus count as damage control, so the refactor is out of proportion
to this PR. Worth a separate issue.
##########
core/common/src/types/streaming_stats.rs:
##########
@@ -19,6 +19,110 @@ use std::sync::{
Arc,
atomic::{AtomicU32, AtomicU64, Ordering},
};
+use tracing::warn;
+
+/// Number of rollup decrements that could not be covered by the counter they
+/// were subtracted from.
+///
+/// A decrement bigger than the total it targets means the tree lost
+/// `parent >= sum(children)` somewhere upstream. The counters are unsigned, so
+/// without a clamp that subtraction wraps to ~1.8e19 and every reader
+/// (`get_stream`, `get_topic`, `/stats`, `/metrics`) serves the wrapped value
+/// until the process restarts.
+///
+/// Clamping trades that for a total that is merely low, and low is where it
+/// stays: increments are `fetch_add`, so every later write stacks on the
+/// clamped base and carries the shortfall with it. Only a rebuild
+/// (`rebuild_parent_totals`) or a restart puts the total back. This counter
+/// plus the `warn!` below is what tells an operator the divergence happened,
+/// and that the aggregate needs a rebuild rather than time.
+static ROLLUP_UNDERFLOWS: AtomicU64 = AtomicU64::new(0);
+
+/// Monotonic count of clamped rollup decrements. Surfaced on `/metrics` so the
+/// clamp is alertable rather than log-grep-able.
+///
+/// Hidden from the docs: `iggy_common` re-exports this module at its root, and
+/// a server-only diagnostic has no business in the published SDK surface.
+#[doc(hidden)]
+#[must_use]
+pub fn rollup_underflows() -> u64 {
+ ROLLUP_UNDERFLOWS.load(Ordering::Relaxed)
+}
+
+/// One scope/counter pair's clamp state: the labels an operator reads, plus
the
+/// count the log throttle keys on.
+struct UnderflowSite {
+ scope: &'static str,
+ counter: &'static str,
+ clamped: AtomicU64,
+}
+
+impl UnderflowSite {
+ const fn new(scope: &'static str, counter: &'static str) -> Self {
+ Self {
+ scope,
+ counter,
+ clamped: AtomicU64::new(0),
+ }
+ }
+
+ /// Count the clamp globally, then log this pair's first and every power of
+ /// two after it.
+ ///
+ /// The throttle is per pair, not global: a skewed tree emits up to three
+ /// lines per counter per rollback and a bulk delete turns that into a
+ /// flood, so a shared count would hold a quiet pair's first-ever
+ /// divergence back for up to 1023 events behind a noisy one.
+ fn report(&self, shortfall: u64) {
+ ROLLUP_UNDERFLOWS.fetch_add(1, Ordering::Relaxed);
+ let clamped = self.clamped.fetch_add(1, Ordering::Relaxed) + 1;
+ if clamped.is_power_of_two() {
+ warn!(
+ scope = self.scope,
+ counter = self.counter,
+ shortfall,
+ clamped,
+ "rollup decrement exceeded the total it was subtracted from;
clamped at zero"
+ );
+ }
+ }
+}
+
+static STREAM_SIZE_BYTES: UnderflowSite = UnderflowSite::new("stream",
"size_bytes");
+static STREAM_MESSAGES_COUNT: UnderflowSite = UnderflowSite::new("stream",
"messages_count");
+static STREAM_SEGMENTS_COUNT: UnderflowSite = UnderflowSite::new("stream",
"segments_count");
Review Comment:
Keeping the nine statics.
The trade you name is the deciding one: three sites plus a `counter`
argument drops the per-pair throttle, and that property is why the throttle is
not global already. A bulk delete floods one pair, and a shared count would
hold a quiet pair's first-ever divergence back behind it for up to 1023 events.
Nine flat statics is the cheapest shape that keeps it. An array per scope
needs a counter enum plus index plumbing at every call site, which is more code
than the nine lines it saves.
--
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]