From a6f9de9581e0cae5d826c8f71ada64cb2a4ce812 Mon Sep 17 00:00:00 2001 From: Wouter Gritter Date: Wed, 16 Sep 2026 13:41:57 +0200 Subject: [PATCH] Bound IntervalledCounter size under packet floods The counter stored one entry per addTime call, so a client sending tiny (even empty) packets could grow it without limit while staying under the bytes-per-second limit, eventually OOMing the proxy. Merge data points within 1ms of the newest one so the entry count is bounded by the window size (~7000 for the default 7s window) rather than by packet rate. --- .../proxy/util/IntervalledCounter.java | 40 ++++-- .../proxy/util/IntervalledCounterTest.java | 115 ++++++++++++++++++ 2 files changed, 146 insertions(+), 9 deletions(-) create mode 100644 proxy/src/test/java/com/velocitypowered/proxy/util/IntervalledCounterTest.java diff --git a/proxy/src/main/java/com/velocitypowered/proxy/util/IntervalledCounter.java b/proxy/src/main/java/com/velocitypowered/proxy/util/IntervalledCounter.java index 9b986dca..7e0bf4c8 100644 --- a/proxy/src/main/java/com/velocitypowered/proxy/util/IntervalledCounter.java +++ b/proxy/src/main/java/com/velocitypowered/proxy/util/IntervalledCounter.java @@ -30,40 +30,51 @@ package com.velocitypowered.proxy.util; *

This class is not thread-safe. If multiple threads access an instance concurrently, * external synchronization is required.

*/ -@SuppressWarnings("checkstyle:WhitespaceAfter") // Not our class public final class IntervalledCounter { private static final int INITIAL_SIZE = 8; + /** + * Data points within this many nanoseconds of the newest one are merged into it, bounding the + * number of stored data points to roughly {@code interval / COALESCE_INTERVAL}. + */ + private static final long COALESCE_INTERVAL = 1_000_000L; // 1ms + /** * Ring buffer holding the timestamp (in nanoseconds) for each data point. */ - protected long[] times; + private long[] times; + /** * Ring buffer holding the count associated with each timestamp. */ - protected long[] counts; + private long[] counts; + /** * The sliding window size in nanoseconds. Only entries with time >= (currentTime - interval) * are considered part of the window. */ - protected final long interval; + private final long interval; + /** * Cached lower bound of the window (in nanoseconds) after the last update. */ - protected long minTime; + private long minTime; + /** * Running sum of all counts currently within the window. */ - protected long sum; + private long sum; + /** * Head index (inclusive) of the ring buffer. */ - protected int head; // inclusive + private int head; // inclusive + /** * Tail index (exclusive) of the ring buffer. */ - protected int tail; // exclusive + private int tail; // exclusive /** * Creates a new counter with the specified interval. @@ -131,6 +142,8 @@ public final class IntervalledCounter { /** * Adds {@code count} units at the specified timestamp, assuming the timestamp is within the * current window. If the timestamp is older than {@code minTime}, the value is ignored. + * If the timestamp is within {@link #COALESCE_INTERVAL} of the newest stored data point, the + * count is merged into that data point instead of creating a new one. * This method does not automatically advance the window; callers should invoke * {@link #updateCurrentTime()} or {@link #updateCurrentTime(long)} beforehand. * @@ -142,6 +155,15 @@ public final class IntervalledCounter { if (currTime - this.minTime < 0) { return; } + if (this.head != this.tail) { + final int last = this.tail == 0 ? this.times.length - 1 : this.tail - 1; + // guard against overflow by using subtraction + if (currTime - this.times[last] < COALESCE_INTERVAL) { + this.counts[last] += count; + this.sum += count; + return; + } + } int nextTail = (this.tail + 1) % this.times.length; if (nextTail == this.head) { this.resize(); @@ -219,7 +241,7 @@ public final class IntervalledCounter { * @return the rate in units per second for the current window */ public double getRate() { - return (double)this.sum / ((double)this.interval * 1.0E-9); + return (double) this.sum / ((double) this.interval * 1.0E-9); } /** diff --git a/proxy/src/test/java/com/velocitypowered/proxy/util/IntervalledCounterTest.java b/proxy/src/test/java/com/velocitypowered/proxy/util/IntervalledCounterTest.java new file mode 100644 index 00000000..f187ff2c --- /dev/null +++ b/proxy/src/test/java/com/velocitypowered/proxy/util/IntervalledCounterTest.java @@ -0,0 +1,115 @@ +/* + * Copyright (C) 2026 Velocity Contributors + * + * This program is free software: you can redistribute it and/or modify + * it under the terms of the GNU General Public License as published by + * the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU General Public License for more details. + * + * You should have received a copy of the GNU General Public License + * along with this program. If not, see . + */ + +package com.velocitypowered.proxy.util; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.concurrent.TimeUnit; +import org.junit.jupiter.api.Test; + +class IntervalledCounterTest { + + private static final long INTERVAL = TimeUnit.SECONDS.toNanos(7); + + @Test + void sumAndRateTrackAddedCounts() { + IntervalledCounter counter = new IntervalledCounter(INTERVAL); + long now = 0; + + counter.updateAndAdd(700, now); + assertEquals(700, counter.getSum()); + assertEquals(100.0, counter.getRate(), 1e-9); + + now += TimeUnit.SECONDS.toNanos(1); + counter.updateAndAdd(1400, now); + assertEquals(2100, counter.getSum()); + assertEquals(300.0, counter.getRate(), 1e-9); + } + + @Test + void dataPointsOutsideTheWindowAreEvicted() { + IntervalledCounter counter = new IntervalledCounter(INTERVAL); + long now = 0; + + counter.updateAndAdd(10, now); + now += TimeUnit.SECONDS.toNanos(3); + counter.updateAndAdd(20, now); + assertEquals(30, counter.getSum()); + assertEquals(2, counter.totalDataPoints()); + + // 7.5s after the first point: only the second point remains + now = TimeUnit.MILLISECONDS.toNanos(7500); + counter.updateCurrentTime(now); + assertEquals(20, counter.getSum()); + assertEquals(1, counter.totalDataPoints()); + + // 10.5s: everything has expired + now = TimeUnit.MILLISECONDS.toNanos(10500); + counter.updateCurrentTime(now); + assertEquals(0, counter.getSum()); + assertEquals(0, counter.totalDataPoints()); + } + + @Test + void storedDataPointsAreBoundedByTimeNotByCallCount() { + IntervalledCounter counter = new IntervalledCounter(INTERVAL); + + // Simulate a flood of tiny (even zero-sized) packets arriving far faster than one per + // millisecond for the whole window. Before coalescing was introduced, each call stored a + // separate data point, so an attacker could grow the ring buffer without bound while staying + // under a bytes-per-second limit. + final long stepNanos = 100; + final long calls = INTERVAL / stepNanos; + long now = 0; + for (long i = 0; i < calls; i++) { + counter.updateAndAdd(i % 2, now); + now += stepNanos; + } + + long expectedSum = calls / 2; + assertEquals(expectedSum, counter.getSum()); + // 7s window at 1ms coalescing granularity is ~7000 points; leave some slack. + assertTrue(counter.totalDataPoints() <= 7100, + "expected at most ~7000 data points, got " + counter.totalDataPoints()); + assertTrue(counter.totalDataPoints() >= 7000, + "expected at least 7000 data points, got " + counter.totalDataPoints()); + } + + @Test + void coalescedDataPointsExpireTogether() { + IntervalledCounter counter = new IntervalledCounter(INTERVAL); + + long now = 0; + counter.updateAndAdd(5, now); + // Within the same 1ms bucket: merged into the previous point + counter.updateAndAdd(7, now + 500_000L); + assertEquals(12, counter.getSum()); + assertEquals(1, counter.totalDataPoints()); + + // A new bucket starts a new point + counter.updateAndAdd(1, now + 1_000_000L); + assertEquals(13, counter.getSum()); + assertEquals(2, counter.totalDataPoints()); + + // The merged bucket carries the timestamp of its first point and expires with it + counter.updateCurrentTime(now + INTERVAL + 1); + assertEquals(1, counter.getSum()); + assertEquals(1, counter.totalDataPoints()); + } +}