fix(velocity): 同一玩家同时只跑一个授权和唤醒,排队中再点同一服直接告知
This commit is contained in:
4 files changed
+173
-26
No files matched your search
@@ -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. */
|
/** async for a task a player is waiting on: a refusal is told to them at once. */
|
||||||
void async(Player player, Runnable task) {
|
boolean async(Player player, Runnable task) {
|
||||||
if (!async(task)) {
|
if (async(task)) {
|
||||||
player.sendMessage(Component.text(
|
return true;
|
||||||
zh(player) ? "服务器现在很忙,请过一会儿再试。"
|
|
||||||
: "The network is busy right now — try again in a moment.",
|
|
||||||
NamedTextColor.YELLOW));
|
|
||||||
}
|
}
|
||||||
|
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).
|
// A run that falls due while the previous one is still going is skipped (SkipIfRunning).
|
||||||
|
|||||||
@@ -24,6 +24,7 @@ import java.util.HashMap;
|
|||||||
import java.util.Locale;
|
import java.util.Locale;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.Optional;
|
import java.util.Optional;
|
||||||
|
import java.util.Set;
|
||||||
import java.util.UUID;
|
import java.util.UUID;
|
||||||
import java.util.concurrent.ConcurrentHashMap;
|
import java.util.concurrent.ConcurrentHashMap;
|
||||||
import java.util.concurrent.atomic.AtomicBoolean;
|
import java.util.concurrent.atomic.AtomicBoolean;
|
||||||
@@ -96,6 +97,11 @@ public final class WaitingRouter {
|
|||||||
private final Map<UUID, Waiter> waiting = new ConcurrentHashMap<>();
|
private final Map<UUID, Waiter> waiting = new ConcurrentHashMap<>();
|
||||||
private final Map<UUID, String> pendingTargets = new ConcurrentHashMap<>();
|
private final Map<UUID, String> pendingTargets = new ConcurrentHashMap<>();
|
||||||
private final Map<UUID, Long> lastGateNotice = new ConcurrentHashMap<>();
|
private final Map<UUID, Long> 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<UUID> authorizing = ConcurrentHashMap.newKeySet();
|
||||||
// tick is scheduled at a fixed rate and makes blocking calls; when felis-api is
|
// 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
|
// slow a run can outlast the interval, and overlapping runs would multiply the
|
||||||
// load on the thing that is already slow.
|
// load on the thing that is already slow.
|
||||||
@@ -332,7 +338,17 @@ public final class WaitingRouter {
|
|||||||
return null;
|
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) {
|
private void authorizeLoginRelease(ServerPreConnectEvent event) {
|
||||||
@@ -593,28 +609,56 @@ public final class WaitingRouter {
|
|||||||
: "You're already on « " + serverName + " ».", NamedTextColor.YELLOW));
|
: "You're already on « " + serverName + " ».", NamedTextColor.YELLOW));
|
||||||
return;
|
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 {
|
try {
|
||||||
if (!linked(id)) {
|
authorizeEntry(player, zh, serverName, fromMenu);
|
||||||
player.sendMessage(Component.text(
|
} finally {
|
||||||
zh ? "请先完成登录,再加入服务器。"
|
authorizing.remove(id);
|
||||||
: "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);
|
|
||||||
});
|
});
|
||||||
|
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);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -500,6 +500,11 @@ final class Fakes {
|
|||||||
volatile String menuAccessError;
|
volatile String menuAccessError;
|
||||||
/** menuAccessHold, while set, holds every menu-access request until it opens. */
|
/** menuAccessHold, while set, holds every menu-access request until it opens. */
|
||||||
volatile CountDownLatch menuAccessHold;
|
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<String, CountDownLatch> holds = new ConcurrentHashMap<>();
|
||||||
/** bodies holds the last request body per "METHOD path". */
|
/** bodies holds the last request body per "METHOD path". */
|
||||||
final Map<String, String> bodies = new ConcurrentHashMap<>();
|
final Map<String, String> bodies = new ConcurrentHashMap<>();
|
||||||
/** calls lists every request as "METHOD path", in arrival order. */
|
/** calls lists every request as "METHOD path", in arrival order. */
|
||||||
@@ -522,6 +527,15 @@ final class Fakes {
|
|||||||
String path = ex.getRequestURI().getPath();
|
String path = ex.getRequestURI().getPath();
|
||||||
bodies.put(method + " " + path, new String(ex.getRequestBody().readAllBytes(), StandardCharsets.UTF_8));
|
bodies.put(method + " " + path, new String(ex.getRequestBody().readAllBytes(), StandardCharsets.UTF_8));
|
||||||
calls.add(method + " " + path);
|
calls.add(method + " " + path);
|
||||||
|
for (Map.Entry<String, CountDownLatch> 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 status = "/api/v1/internal/account/link/status/";
|
||||||
String servers = "/api/v1/internal/servers/";
|
String servers = "/api/v1/internal/servers/";
|
||||||
String menuAccess = "/api/v1/internal/player/menu-access/";
|
String menuAccess = "/api/v1/internal/player/menu-access/";
|
||||||
|
|||||||
@@ -17,7 +17,10 @@ import java.util.ArrayList;
|
|||||||
import java.util.Collections;
|
import java.util.Collections;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Locale;
|
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.AtomicLong;
|
||||||
|
import java.util.concurrent.atomic.AtomicReference;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* WaitingRouterTest drives the real WaitingRouter, ServerRegistry, FelisApiClient and
|
* WaitingRouterTest drives the real WaitingRouter, ServerRegistry, FelisApiClient and
|
||||||
@@ -75,6 +78,7 @@ public final class WaitingRouterTest {
|
|||||||
longStart();
|
longStart();
|
||||||
menuAndCommands();
|
menuAndCommands();
|
||||||
joins();
|
joins();
|
||||||
|
slowApi();
|
||||||
backToLobby();
|
backToLobby();
|
||||||
disconnectAndRelease();
|
disconnectAndRelease();
|
||||||
} finally {
|
} finally {
|
||||||
@@ -102,6 +106,8 @@ public final class WaitingRouterTest {
|
|||||||
view("lambda", false, "10.43.0.25:25565"),
|
view("lambda", false, "10.43.0.25:25565"),
|
||||||
view("omicron", false, "10.43.0.26:25565"),
|
view("omicron", false, "10.43.0.26:25565"),
|
||||||
view("sigma", false, "10.43.0.27: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)));
|
view("fresh", false, null)));
|
||||||
if (withGone) {
|
if (withGone) {
|
||||||
list.add(view("gone", false, "10.43.0.9:25565"));
|
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());
|
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<ServerPreConnectEvent> 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
|
// /felis lobby and /felis go lobby: the lobby is a system server, so it never
|
||||||
// appeared among the servers /felis go knows.
|
// appeared among the servers /felis go knows.
|
||||||
private static void backToLobby() {
|
private static void backToLobby() {
|
||||||
|
|||||||
Reference in new issue
Block a user