feat(velocity): felis-api 调用改走有界线程池、定时任务防重叠,超时可配置,幂等 GET 带抖动重试一次

This commit is contained in:
Lemon-miaow committed 2026-09-26 00:03:19 +08:00
1 parent 001f060027
commit fcf5c305ea
16 files changed
+872 -37

No files matched your search

@@ -0,0 +1,68 @@
package best.lolicon.felis.velocity;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Consumer;
/**
* BoundedExecutor runs the plugin's felis-api calls on a fixed number of threads with
* a bounded wait queue. Velocity's scheduler starts a thread for every task it is
* handed, so while felis-api was slow a burst of joins, menu clicks and commands
* stacked blocking calls without limit, and each of them landed on the API the moment
* it recovered. Here at most {@code threads} calls run and at most {@code queue} wait;
* past that {@link #submit} refuses at once, so the caller can say "busy" instead of
* leaving the player waiting on a call that would run minutes later.
*
* <p>A task that throws is handed to {@code onError} and the worker carries on, so one
* bad task never shrinks the pool or disappears without a log line.
*/
final class BoundedExecutor {
private final ThreadPoolExecutor pool;
private final Consumer<Throwable> onError;
BoundedExecutor(String name, int threads, int queue, Consumer<Throwable> onError) {
this.onError = onError;
AtomicInteger n = new AtomicInteger();
this.pool = new ThreadPoolExecutor(threads, threads, 0L, TimeUnit.MILLISECONDS,
new ArrayBlockingQueue<>(queue),
r -> {
Thread t = new Thread(r, name + "-" + n.incrementAndGet());
t.setDaemon(true);
return t;
},
new ThreadPoolExecutor.AbortPolicy());
}
/** submit queues the task, or returns false when every thread is busy and the queue is full, or after shutdown. */
boolean submit(Runnable task) {
try {
pool.execute(() -> {
try {
task.run();
} catch (Throwable t) {
onError.accept(t);
}
});
return true;
} catch (RejectedExecutionException e) {
return false;
}
}
/** running is how many tasks are executing now; waiting how many are queued. For the busy log line. */
int running() {
return pool.getActiveCount();
}
int waiting() {
return pool.getQueue().size();
}
/** shutdown stops taking tasks and interrupts the running ones (the proxy is going down). */
void shutdown() {
pool.shutdownNow();
}
}
@@ -237,7 +237,7 @@ public final class ControlChannel implements WaitingRouter.MenuTransferListener
send(source, cached.frame);
return;
}
plugin.async(() -> {
boolean taken = plugin.async(() -> {
try {
MenuStatus s = api.menuStatus(server);
ControlFrame frame = ControlFrame.statusUpdate(
@@ -248,6 +248,9 @@ public final class ControlChannel implements WaitingRouter.MenuTransferListener
send(source, errorFrame(e, server));
}
});
if (!taken) {
send(source, busyFrame(server));
}
}
// A WakeRequest is the menu's Join/Wake button on a server the player owns: wake
@@ -262,7 +265,7 @@ public final class ControlChannel implements WaitingRouter.MenuTransferListener
// autostartPolicy gate). A claim refusal answers with Error and never wakes.
private void handleClaim(ServerConnection source, Player player, String server) {
UUID id = player.getUniqueId();
plugin.async(() -> {
boolean taken = plugin.async(() -> {
try {
api.claim(server, id);
} catch (LinkException e) {
@@ -274,6 +277,9 @@ public final class ControlChannel implements WaitingRouter.MenuTransferListener
// hop for the wake, which is fine from here.
router.enqueueFromMenu(player, server);
});
if (!taken) {
send(source, busyFrame(server));
}
}
/**
@@ -305,6 +311,12 @@ public final class ControlChannel implements WaitingRouter.MenuTransferListener
return ControlFrame.error(code, message, server);
}
// busyFrame answers a frame whose felis-api call the bounded pool refused.
private static ControlFrame busyFrame(String server) {
return ControlFrame.error("busy",
"Felis 现在很忙,请稍后再试 / Felis is busy right now — please try again in a moment.", server);
}
private static final class CachedStatus {
final ControlFrame frame;
final long atMillis;
@@ -18,6 +18,7 @@ import com.velocitypowered.api.command.CommandSource;
import com.velocitypowered.api.event.Subscribe;
import com.velocitypowered.api.event.connection.DisconnectEvent;
import com.velocitypowered.api.event.proxy.ProxyInitializeEvent;
import com.velocitypowered.api.event.proxy.ProxyShutdownEvent;
import com.velocitypowered.api.plugin.Plugin;
import com.velocitypowered.api.plugin.annotation.DataDirectory;
import com.velocitypowered.api.proxy.Player;
@@ -34,6 +35,7 @@ import java.util.List;
import java.util.Locale;
import java.util.Optional;
import java.util.UUID;
import java.util.concurrent.atomic.AtomicLong;
import java.util.regex.Pattern;
/**
@@ -78,6 +80,13 @@ public final class FelisVelocityPlugin {
// friends around one after another is not. Sub-TTL on purpose: a sender may hold several
// live invites, they just cannot post them all in one breath.
private static final Duration INVITE_COOLDOWN = Duration.ofSeconds(30);
// felis-api calls in flight, and waiting, at most. A healthy call takes milliseconds,
// so this is far above what one proxy's players produce; it only binds while
// felis-api is slow, where 8 x the 10 s request timeout already means the 64th
// waiter hears back after about 80 s. Refusing past that beats a reply minutes late.
private static final int API_THREADS = 8;
private static final int API_QUEUE = 64;
private static final long BUSY_LOG_INTERVAL_MILLIS = 60_000L;
private final ProxyServer proxy;
private final Logger logger;
@@ -88,6 +97,8 @@ public final class FelisVelocityPlugin {
// typing, and a bound on a client macro that would otherwise mint link codes or
// fire claims as fast as it can send chat.
private final FrameBudget commandBudget = new FrameBudget(5, 0.2, System::currentTimeMillis);
private final BoundedExecutor apiCalls;
private final AtomicLong lastBusyLog = new AtomicLong();
private FelisVelocityConfig config;
private LinkClient linkClient;
@@ -103,6 +114,8 @@ public final class FelisVelocityPlugin {
this.proxy = proxy;
this.logger = logger;
this.dataDirectory = dataDirectory;
this.apiCalls = new BoundedExecutor("felis-api", API_THREADS, API_QUEUE,
t -> logger.error("Felis: a felis-api task failed", t));
}
@Subscribe
@@ -167,6 +180,11 @@ public final class FelisVelocityPlugin {
commandBudget.forget(event.getPlayer().getUniqueId());
}
@Subscribe
public void onProxyShutdown(ProxyShutdownEvent event) {
apiCalls.shutdown();
}
// withinBudget spends one command token for the player, or tells them to slow down.
private boolean withinBudget(Player player) {
if (commandBudget.tryTake(player.getUniqueId())) {
@@ -178,13 +196,37 @@ public final class FelisVelocityPlugin {
return false;
}
/** async runs a task on Velocity's scheduler so felis-api I/O never blocks the proxy thread. */
void async(Runnable task) {
proxy.getScheduler().buildTask(this, task).schedule();
/**
* async runs a felis-api task off the proxy thread, on the bounded pool. It returns
* false, without running the task, when API_THREADS calls are already in flight and
* API_QUEUE more are waiting; the caller then answers "busy" in its own terms.
*/
boolean async(Runnable task) {
if (apiCalls.submit(task)) {
return true;
}
long now = System.currentTimeMillis();
long last = lastBusyLog.get();
if (now - last >= BUSY_LOG_INTERVAL_MILLIS && lastBusyLog.compareAndSet(last, now)) {
logger.warn("Felis: felis-api is not keeping up ({} calls running, {} waiting); refusing new calls "
+ "until the queue drains.", apiCalls.running(), apiCalls.waiting());
}
return false;
}
/** async for a task a player is waiting on: a refusal is told to them at once. */
void async(Player player, Runnable task) {
if (!async(task)) {
player.sendMessage(Component.text(
zh(player) ? "服务器现在很忙,请过一会儿再试。"
: "The network is busy right now — try again in a moment.",
NamedTextColor.YELLOW));
}
}
// A run that falls due while the previous one is still going is skipped (SkipIfRunning).
private void repeating(Duration interval, Runnable task) {
proxy.getScheduler().buildTask(this, task).delay(interval).repeat(interval).schedule();
proxy.getScheduler().buildTask(this, new SkipIfRunning(task)).delay(interval).repeat(interval).schedule();
}
/**
@@ -263,7 +305,7 @@ public final class FelisVelocityPlugin {
boolean zh = zh(player);
player.sendMessage(Component.text(
zh ? "正在获取绑定码……" : "Requesting a link code…", NamedTextColor.GRAY));
async(() -> {
async(player, () -> {
try {
LinkCode code = linkClient.requestCode(player.getUniqueId());
player.sendMessage(Component.text(
@@ -583,7 +625,7 @@ public final class FelisVelocityPlugin {
UUID uuid = player.getUniqueId();
player.sendMessage(Component.text(
zh ? "正在认领「" + name + "」……" : "Claiming « " + name + " »…", NamedTextColor.GRAY));
async(() -> {
async(player, () -> {
try {
apiClient.claim(name, uuid);
player.sendMessage(Component.text(
@@ -619,7 +661,7 @@ public final class FelisVelocityPlugin {
String who = player.getUsername();
player.sendMessage(Component.text(
zh ? "正在发起账户迁移……" : "Starting account migration…", NamedTextColor.GRAY));
async(() -> {
async(player, () -> {
try {
apiClient.migrateStart(uuid);
String panelHost = config.panelHostname();
@@ -718,7 +760,7 @@ public final class FelisVelocityPlugin {
UUID approver = player.getUniqueId();
String who = player.getUsername();
if (accountArg == null) {
async(() -> {
async(player, () -> {
try {
OpLoginView req = apiClient.opLoginShow(code, approver);
long age = req.createdAt() == null ? -1
@@ -735,7 +777,7 @@ public final class FelisVelocityPlugin {
String account = accountArg.trim();
player.sendMessage(Component.text(
zh ? "正在批准管理员登录……" : "Approving operator sign-in…", NamedTextColor.GRAY));
async(() -> {
async(player, () -> {
try {
OpLoginView done = apiClient.opLoginApprove(code, approver, account);
player.sendMessage(Component.text(
@@ -0,0 +1,32 @@
package best.lolicon.felis.velocity;
import java.util.concurrent.atomic.AtomicBoolean;
/**
* SkipIfRunning wraps a repeating task so that a run falling due while the previous
* one is still going is skipped rather than started beside it. Velocity's
* {@code repeat()} fires on schedule whatever the last run is doing, so a felis-api
* call slower than the interval (a 10 s request timeout against the 15 s registration
* refresh, or several calls in one pass) stacked runs that each held a thread and
* sent the same requests again, adding load exactly when the API was struggling.
*/
final class SkipIfRunning implements Runnable {
private final Runnable task;
private final AtomicBoolean running = new AtomicBoolean();
SkipIfRunning(Runnable task) {
this.task = task;
}
@Override
public void run() {
if (!running.compareAndSet(false, true)) {
return;
}
try {
task.run();
} finally {
running.set(false);
}
}
}
@@ -343,7 +343,7 @@ public final class WaitingRouter {
|| name.equalsIgnoreCase(lobbyServer)) {
return; // system/static servers do not affect user-server activity
}
plugin.async(() -> {
boolean taken = plugin.async(() -> {
try {
api.reportJoin(name, id);
} catch (LinkException e) {
@@ -353,6 +353,9 @@ public final class WaitingRouter {
id, name, e.statusCode(), e.getMessage());
}
});
if (!taken) {
log.warn("Felis: join-event for {} on {} dropped: the felis-api call queue is full", id, name);
}
}
/** tick drains the waiting queue; the plugin schedules it on the async pool. */
@@ -451,7 +454,7 @@ public final class WaitingRouter {
: "You're already on « " + serverName + " ».", NamedTextColor.YELLOW));
return;
}
plugin.async(() -> {
plugin.async(player, () -> {
try {
if (!linked(id)) {
player.sendMessage(Component.text(
@@ -0,0 +1,164 @@
package best.lolicon.felis.velocity;
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
/**
* ApiPoolTest covers the two guards around the proxy's felis-api calls:
* {@link BoundedExecutor} (a fixed number of calls in flight, a bounded queue, refusal
* past it, a throwing task logged and survived) and {@link SkipIfRunning} (a repeating
* task never overlaps itself, and a run that throws does not block later runs).
* Framework free: a failed assertion throws and the process exits non-zero.
*
* <p>Run: {@code javac -d <out> velocity/src/main/java/best/lolicon/felis/velocity/{BoundedExecutor,SkipIfRunning}.java
* velocity/test/best/lolicon/felis/velocity/ApiPoolTest.java && java -cp <out>
* best.lolicon.felis.velocity.ApiPoolTest}.
*/
public final class ApiPoolTest {
private static int checks;
public static void main(String[] args) throws Exception {
fullPoolRefusesInsteadOfGrowing();
throwingTaskIsReportedAndThePoolLives();
shutdownRefuses();
repeatingTaskNeverOverlaps();
throwingRunDoesNotWedgeTheGuard();
System.out.println("ApiPoolTest OK (" + checks + " checks)");
}
// Two threads, two queue slots: four tasks are taken (two running, two waiting), the
// fifth is refused; once the blockers finish, every accepted task runs.
private static void fullPoolRefusesInsteadOfGrowing() throws Exception {
BoundedExecutor pool = new BoundedExecutor("t", 2, 2, t -> { });
CountDownLatch started = new CountDownLatch(2);
CountDownLatch release = new CountDownLatch(1);
AtomicInteger ran = new AtomicInteger();
Runnable blocker = () -> {
started.countDown();
await(release);
ran.incrementAndGet();
};
assertEq("first taken", true, pool.submit(blocker));
assertEq("second taken", true, pool.submit(blocker));
assertEq("both blockers running", true, started.await(5, TimeUnit.SECONDS));
assertEq("running while blocked", 2, pool.running());
assertEq("third queued", true, pool.submit(ran::incrementAndGet));
assertEq("fourth queued", true, pool.submit(ran::incrementAndGet));
assertEq("waiting while blocked", 2, pool.waiting());
assertEq("fifth refused", false, pool.submit(ran::incrementAndGet));
assertEq("nothing finished yet", 0, ran.get());
release.countDown();
CountDownLatch drained = new CountDownLatch(1);
waitFor(() -> ran.get() == 4, "accepted tasks all run");
assertEq("refused task never ran", 4, ran.get());
// Room again once drained.
assertEq("taken after drain", true, pool.submit(drained::countDown));
assertEq("post-drain task ran", true, drained.await(5, TimeUnit.SECONDS));
pool.shutdown();
}
private static void throwingTaskIsReportedAndThePoolLives() throws Exception {
List<Throwable> errors = new CopyOnWriteArrayList<>();
BoundedExecutor pool = new BoundedExecutor("t", 1, 1, errors::add);
pool.submit(() -> {
throw new IllegalStateException("boom");
});
CountDownLatch next = new CountDownLatch(1);
waitFor(() -> errors.size() == 1, "error reported");
assertEq("reported message", "boom", errors.get(0).getMessage());
assertEq("taken after a throw", true, pool.submit(next::countDown));
assertEq("task after a throw ran", true, next.await(5, TimeUnit.SECONDS));
pool.shutdown();
}
private static void shutdownRefuses() {
BoundedExecutor pool = new BoundedExecutor("t", 1, 1, t -> { });
pool.shutdown();
assertEq("refused after shutdown", false, pool.submit(() -> { }));
}
// A second run that falls due while the first is inside the task is skipped, not
// started beside it; the next run after the first finishes goes ahead.
private static void repeatingTaskNeverOverlaps() throws Exception {
CountDownLatch inside = new CountDownLatch(1);
CountDownLatch release = new CountDownLatch(1);
AtomicInteger entered = new AtomicInteger();
AtomicInteger concurrent = new AtomicInteger();
AtomicInteger maxConcurrent = new AtomicInteger();
SkipIfRunning guarded = new SkipIfRunning(() -> {
int now = concurrent.incrementAndGet();
maxConcurrent.accumulateAndGet(now, Math::max);
entered.incrementAndGet();
if (entered.get() == 1) {
inside.countDown();
await(release);
}
concurrent.decrementAndGet();
});
Thread first = new Thread(guarded);
first.start();
assertEq("first run inside", true, inside.await(5, TimeUnit.SECONDS));
guarded.run(); // due while the first is still running
guarded.run();
assertEq("overlapping runs skipped", 1, entered.get());
release.countDown();
first.join(5000);
guarded.run();
assertEq("run after the first finished", 2, entered.get());
assertEq("never two at once", 1, maxConcurrent.get());
}
private static void throwingRunDoesNotWedgeTheGuard() {
AtomicInteger calls = new AtomicInteger();
SkipIfRunning guarded = new SkipIfRunning(() -> {
if (calls.incrementAndGet() == 1) {
throw new IllegalStateException("felis-api down");
}
});
try {
guarded.run();
throw new AssertionError("the first run's exception was swallowed");
} catch (IllegalStateException e) {
assertEq("first run's exception propagates", "felis-api down", e.getMessage());
}
guarded.run();
assertEq("runs after a throw", 2, calls.get());
}
// ---- harness ----
interface Cond {
boolean ok();
}
private static void waitFor(Cond cond, String what) throws InterruptedException {
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5);
while (!cond.ok()) {
if (System.nanoTime() > deadline) {
throw new AssertionError(what + ": timed out");
}
Thread.sleep(5);
}
checks++;
}
private static void await(CountDownLatch l) {
try {
l.await(5, TimeUnit.SECONDS);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
private static void assertEq(String what, Object want, Object got) {
if (!want.equals(got)) {
throw new AssertionError(what + ": got " + got + ", want " + want);
}
checks++;
}
}