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.
This commit is contained in:
@@ -30,40 +30,51 @@ package com.velocitypowered.proxy.util;
|
||||
* <p>This class is not thread-safe. If multiple threads access an instance concurrently,
|
||||
* external synchronization is required.</p>
|
||||
*/
|
||||
@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);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -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 <https://www.gnu.org/licenses/>.
|
||||
*/
|
||||
|
||||
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());
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user