From 4ce7eb2aea50390b8268999d7abb77730a87f9a1 Mon Sep 17 00:00:00 2001 From: ImDanXie Date: Sat, 26 Sep 2026 01:59:11 +0800 Subject: [PATCH] fix(scheduler): base EMALong.update() on the decayed value get() applies time-based decay but never writes it back, while update() used the raw, undecayed prev.ema as its base. After an idle period the two methods therefore disagreed: get() reported congestion had subsided, yet the next update() revived the stale value at 0.9x its old magnitude and re-armed the decay delay - which can pin the batcher in permanent back-pressure ("release a few batches, drop for tens of seconds" loop). Extract the decay computation into a shared helper and use it in both get() and update(), so a fresh sample is blended onto the decayed value. Control experiment: updateAfterIdleUsesDecayedBase fails on the old code (18s instead of ~32ms revived) and passes with the fix; existing EMALong expectations unchanged (8/8). Fixes #301 --- .../apache/bifromq/basescheduler/EMALong.java | 13 ++++++++--- .../bifromq/basescheduler/EMALongTest.java | 22 +++++++++++++++++++ 2 files changed, 32 insertions(+), 3 deletions(-) diff --git a/base-scheduler/src/main/java/org/apache/bifromq/basescheduler/EMALong.java b/base-scheduler/src/main/java/org/apache/bifromq/basescheduler/EMALong.java index f412a559e..411a55be3 100644 --- a/base-scheduler/src/main/java/org/apache/bifromq/basescheduler/EMALong.java +++ b/base-scheduler/src/main/java/org/apache/bifromq/basescheduler/EMALong.java @@ -46,7 +46,12 @@ public void update(long newValue) { long now = nowSupplier.get(); while (true) { State prev = state.get(); - long newEma = (prev.ema == 0L) ? newValue : (long) Math.ceil(prev.ema * (1 - alpha) + newValue * alpha); + // Base the new sample on the *decayed* value, so update() and get() agree on the current + // EMA. Using the raw prev.ema here resurrects a stale congestion value after an idle period + // (get() reported it had decayed away) and re-arms the decay delay, which can pin the scheduler + // in permanent back-pressure. + long prevEma = decay(prev, now); + long newEma = (prevEma == 0L) ? newValue : (long) Math.ceil(prevEma * (1 - alpha) + newValue * alpha); State next = new State(newEma, now); if (state.compareAndSet(prev, next)) { return; @@ -55,8 +60,10 @@ public void update(long newValue) { } public long get() { - long now = nowSupplier.get(); - State s = state.get(); + return decay(state.get(), nowSupplier.get()); + } + + private long decay(State s, long now) { if (s.ema == 0L || s.lastTs == 0L) { return s.ema; } diff --git a/base-scheduler/src/test/java/org/apache/bifromq/basescheduler/EMALongTest.java b/base-scheduler/src/test/java/org/apache/bifromq/basescheduler/EMALongTest.java index 3ca03f620..39d1b98b9 100644 --- a/base-scheduler/src/test/java/org/apache/bifromq/basescheduler/EMALongTest.java +++ b/base-scheduler/src/test/java/org/apache/bifromq/basescheduler/EMALongTest.java @@ -113,4 +113,26 @@ void neverDecay() { // expected = 80 assertEquals(ema.get(), 80); } + + /** + * Regression: after an idle period a fresh sample must not resurrect the stale (decayed-away) + * congestion value. With the buggy base (undecayed prev.ema), update(0) after 60s idle would return + * 0.9 * 20s = 18s instead of ~32ms. + */ + @Test + void updateAfterIdleUsesDecayedBase() { + EMALong ema = new EMALong(nowSupplier, 0.1, 0.9, 5_000_000_000L); + fakeTime.set(1L); + ema.update(20_000_000_000L); // congestion spike: ema = 20s + // idle for 60s beyond the decay delay: get() reports ≈ 20s * 0.9^60 ≈ 36ms + fakeTime.set(5_000_000_000L + 60_000_000_000L + 1L); + long decayedBefore = ema.get(); + org.testng.Assert.assertTrue(decayedBefore < 1_000_000_000L, + "congestion should have decayed after 60s idle, but get()=" + decayedBefore); + // a fresh zero-latency sample must stay in the decayed range, not jump back to ~0.9 * 20s + ema.update(0L); + long after = ema.get(); + org.testng.Assert.assertTrue(after < 1_000_000_000L, + "update() resurrected the stale value after idle: " + after); + } }