diff --git a/plugins/velocity/src/main/java/best/lolicon/felis/velocity/FelisVelocityPlugin.java b/plugins/velocity/src/main/java/best/lolicon/felis/velocity/FelisVelocityPlugin.java index 312cbb4..137411d 100644 --- a/plugins/velocity/src/main/java/best/lolicon/felis/velocity/FelisVelocityPlugin.java +++ b/plugins/velocity/src/main/java/best/lolicon/felis/velocity/FelisVelocityPlugin.java @@ -279,13 +279,15 @@ public final class FelisVelocityPlugin { } /** 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)); + boolean async(Player player, Runnable task) { + if (async(task)) { + return true; } + player.sendMessage(Component.text( + zh(player) ? "服务器现在很忙,请过一会儿再试。" + : "The network is busy right now — try again in a moment.", + NamedTextColor.YELLOW)); + return false; } // A run that falls due while the previous one is still going is skipped (SkipIfRunning). diff --git a/plugins/velocity/src/main/java/best/lolicon/felis/velocity/WaitingRouter.java b/plugins/velocity/src/main/java/best/lolicon/felis/velocity/WaitingRouter.java index 5dcdf5a..582bd57 100644 --- a/plugins/velocity/src/main/java/best/lolicon/felis/velocity/WaitingRouter.java +++ b/plugins/velocity/src/main/java/best/lolicon/felis/velocity/WaitingRouter.java @@ -24,6 +24,7 @@ import java.util.HashMap; import java.util.Locale; import java.util.Map; import java.util.Optional; +import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicBoolean; @@ -96,6 +97,11 @@ public final class WaitingRouter { private final Map waiting = new ConcurrentHashMap<>(); private final Map pendingTargets = new ConcurrentHashMap<>(); private final Map lastGateNotice = new ConcurrentHashMap<>(); + // Players whose login release or queue entry is being authorized right now. The + // gate re-sends its release every 1 to 8 s and a player can click a tile again; + // while felis-api was slow each of those ran its own link check and wake, so one + // player woke the server, and was told it was starting, several times over. + private final Set authorizing = ConcurrentHashMap.newKeySet(); // tick is scheduled at a fixed rate and makes blocking calls; when felis-api is // slow a run can outlast the interval, and overlapping runs would multiply the // load on the thing that is already slow. @@ -332,7 +338,17 @@ public final class WaitingRouter { return null; } - return EventTask.async(() -> authorizeLoginRelease(event)); + UUID id = player.getUniqueId(); + if (!authorizing.add(id)) { + return null; // the release already being authorized answers; the gate asks again after + } + return EventTask.async(() -> { + try { + authorizeLoginRelease(event); + } finally { + authorizing.remove(id); + } + }); } private void authorizeLoginRelease(ServerPreConnectEvent event) { @@ -593,28 +609,56 @@ public final class WaitingRouter { : "You're already on « " + serverName + " ».", NamedTextColor.YELLOW)); return; } - plugin.async(player, () -> { + // A wait that is already on for this server answers a second click itself: going + // round again would only wake the server once more and restart the wait. + Waiter queued = waiting.get(id); + if (queued != null && queued.serverName.equalsIgnoreCase(serverName)) { + player.sendMessage(Component.text( + zh ? "你已在排队等待「" + queued.serverName + "」,就绪后会自动把你传送过去。" + : "You're already waiting for « " + queued.serverName + " » — you'll be moved in when it's ready.", + NamedTextColor.GRAY)); + return; + } + if (!authorizing.add(id)) { + player.sendMessage(Component.text( + zh ? "上一个请求还在处理,请稍候。" : "Still working on your last request — hold on a moment.", + NamedTextColor.GRAY)); + return; + } + boolean taken = plugin.async(player, () -> { try { - if (!linked(id)) { - player.sendMessage(Component.text( - zh ? "请先完成登录,再加入服务器。" - : "Finish signing in before joining a server.", NamedTextColor.YELLOW)); - return; - } - } catch (LinkException e) { - log.warn("Felis: could not verify queue entry for {} (status={}): {}", - id, e.statusCode(), e.getMessage()); - player.sendMessage(Component.text( - zh ? "登录验证暂时不可用,请稍后重试。" - : "Login verification is temporarily unavailable. Please try again shortly.", - NamedTextColor.RED)); - return; + authorizeEntry(player, zh, serverName, fromMenu); + } finally { + authorizing.remove(id); } - if (joinIfRunning(player, serverName, fromMenu)) { - return; - } - wakeAndWaitLinked(player, serverName, fromMenu); }); + if (!taken) { + authorizing.remove(id); + } + } + + private void authorizeEntry(Player player, boolean zh, String serverName, boolean fromMenu) { + UUID id = player.getUniqueId(); + try { + if (!linked(id)) { + player.sendMessage(Component.text( + zh ? "请先完成登录,再加入服务器。" + : "Finish signing in before joining a server.", NamedTextColor.YELLOW)); + return; + } + } catch (LinkException e) { + log.warn("Felis: could not verify queue entry for {} (status={}): {}", + id, e.statusCode(), e.getMessage()); + player.sendMessage(Component.text( + zh ? "登录验证暂时不可用,请稍后重试。" + : "Login verification is temporarily unavailable. Please try again shortly.", + NamedTextColor.RED)); + return; + } + if (joinIfRunning(player, serverName, fromMenu)) { + return; + } + wakeAndWaitLinked(player, serverName, fromMenu); } /** diff --git a/plugins/velocity/test/best/lolicon/felis/velocity/Fakes.java b/plugins/velocity/test/best/lolicon/felis/velocity/Fakes.java index e8f9d36..42cc525 100644 --- a/plugins/velocity/test/best/lolicon/felis/velocity/Fakes.java +++ b/plugins/velocity/test/best/lolicon/felis/velocity/Fakes.java @@ -500,6 +500,11 @@ final class Fakes { volatile String menuAccessError; /** menuAccessHold, while set, holds every menu-access request until it opens. */ volatile CountDownLatch menuAccessHold; + /** + * holds maps a "METHOD path" prefix to a latch: a matching request waits for it to + * open (5 s at most), the way a slow felis-api keeps a call hanging. + */ + final Map holds = new ConcurrentHashMap<>(); /** bodies holds the last request body per "METHOD path". */ final Map bodies = new ConcurrentHashMap<>(); /** calls lists every request as "METHOD path", in arrival order. */ @@ -522,6 +527,15 @@ final class Fakes { String path = ex.getRequestURI().getPath(); bodies.put(method + " " + path, new String(ex.getRequestBody().readAllBytes(), StandardCharsets.UTF_8)); calls.add(method + " " + path); + for (Map.Entry h : holds.entrySet()) { + if ((method + " " + path).startsWith(h.getKey())) { + try { + h.getValue().await(5, TimeUnit.SECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + } String status = "/api/v1/internal/account/link/status/"; String servers = "/api/v1/internal/servers/"; String menuAccess = "/api/v1/internal/player/menu-access/"; diff --git a/plugins/velocity/test/best/lolicon/felis/velocity/WaitingRouterTest.java b/plugins/velocity/test/best/lolicon/felis/velocity/WaitingRouterTest.java index 19815d2..3779798 100644 --- a/plugins/velocity/test/best/lolicon/felis/velocity/WaitingRouterTest.java +++ b/plugins/velocity/test/best/lolicon/felis/velocity/WaitingRouterTest.java @@ -17,7 +17,10 @@ import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.Locale; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicReference; /** * WaitingRouterTest drives the real WaitingRouter, ServerRegistry, FelisApiClient and @@ -75,6 +78,7 @@ public final class WaitingRouterTest { longStart(); menuAndCommands(); joins(); + slowApi(); backToLobby(); disconnectAndRelease(); } finally { @@ -102,6 +106,8 @@ public final class WaitingRouterTest { view("lambda", false, "10.43.0.25:25565"), view("omicron", false, "10.43.0.26:25565"), view("sigma", false, "10.43.0.27:25565"), + view("upsilon", false, "10.43.0.28:25565"), + view("omega", false, "10.43.0.29:25565"), view("fresh", false, null))); if (withGone) { list.add(view("gone", false, "10.43.0.9:25565")); @@ -669,6 +675,87 @@ public final class WaitingRouterTest { assertEq("refused transfer: dialled once", 1, full.connects.size()); } + // A slow felis-api still gets one authorization per player at a time: the login + // gate's re-sent release and a second click wait for the first instead of each + // running its own link check and wake. + private static void slowApi() throws InterruptedException { + // The gate sends its release again while the first is still at the link check. + Fakes.FakePlayer slow = player("upsilon.mc.test", true); + assertEq("slow release: routed to login", "login", choose(slow)); + CountDownLatch linkHeld = new CountDownLatch(1); + api.holds.put(LINK + slow.id, linkHeld); + AtomicReference first = new AtomicReference<>(); + Thread gate = new Thread(() -> first.set(release(slow))); + gate.start(); + Fakes.await("slow release: the first at the link check", () -> api.count(LINK + slow.id) == 1); + ServerPreConnectEvent again = release(slow); + assertEq("slow release: the re-sent one denied at once", "denied", allowedTo(again)); + linkHeld.countDown(); + gate.join(); + api.holds.remove(LINK + slow.id); + assertEq("slow release: the first goes to the lobby", "lobby", allowedTo(first.get())); + assertEq("slow release: one link check", 1, api.count(LINK + slow.id)); + assertEq("slow release: woken once", 1, api.count("POST " + SERVERS + "upsilon/wake")); + assertEq("slow release: told once", 1, count(slow.messages, "Starting « upsilon »")); + assertEq("slow release: the next release authorized again", "lobby", allowedTo(release(slow))); + router.onDisconnect(new DisconnectEvent(slow.player, DisconnectEvent.LoginStatus.SUCCESSFUL_LOGIN)); + + // A second entry while the first is still waking the server is told to hold on. + int queued = router.waitingCount(); + Fakes.FakePlayer eager = player(null, true); + eager.current = lobby; + String omegaWake = "POST " + SERVERS + "omega/wake"; + CountDownLatch wakeHeld = new CountDownLatch(1); + api.holds.put(omegaWake, wakeHeld); + router.enqueueFromMenu(eager.player, "omega"); + Fakes.await("eager: the first entry waking", () -> api.count(omegaWake) == 1); + router.enqueueFromCommand(eager.player, "omega"); + assertEq("eager: the second told to hold on", true, eager.said("Still working on your last request")); + wakeHeld.countDown(); + api.holds.remove(omegaWake); + Fakes.await("eager: queued", () -> router.waitingCount() == queued + 1); + Thread.sleep(200); // let anything wrongly submitted land too + assertEq("eager: woken once", 1, api.count(omegaWake)); + assertEq("eager: one link check", 1, api.count(LINK + eager.id)); + + // Queued, a click for the same server is answered from the queue. + router.enqueueFromMenu(eager.player, "OMEGA"); + assertEq("eager: told it is waiting", true, eager.said("You're already waiting for « omega »")); + Thread.sleep(200); + assertEq("eager: still one link check", 1, api.count(LINK + eager.id)); + assertEq("eager: still woken once", 1, api.count(omegaWake)); + assertEq("eager: told it is starting once", 1, count(eager.messages, "Starting « omega »")); + router.onDisconnect(new DisconnectEvent(eager.player, DisconnectEvent.LoginStatus.SUCCESSFUL_LOGIN)); + + // An entry the full call pool turns away leaves nothing behind to block the next. + CountDownLatch busy = new CountDownLatch(1); + Runnable parked = () -> { + try { + busy.await(5, TimeUnit.SECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + }; + // Fill every thread and queue slot; a thread that starts late frees a queue slot, + // so fill again until the pool still refuses after settling. + do { + while (plugin.async(parked)) { + // next slot + } + Thread.sleep(50); + } while (plugin.async(parked)); + Fakes.FakePlayer turned = player(null, true); + turned.current = lobby; + router.enqueueFromCommand(turned.player, "omega"); + assertEq("pool full: told", true, turned.said("The network is busy right now")); + busy.countDown(); + Fakes.await("pool drained", () -> plugin.async(() -> { })); + router.enqueueFromCommand(turned.player, "omega"); + Fakes.await("pool drained: queued", () -> router.waitingCount() == queued + 1); + assertEq("pool drained: not told to hold on", false, turned.said("Still working")); + router.onDisconnect(new DisconnectEvent(turned.player, DisconnectEvent.LoginStatus.SUCCESSFUL_LOGIN)); + } + // /felis lobby and /felis go lobby: the lobby is a system server, so it never // appeared among the servers /felis go knows. private static void backToLobby() {