diff --git a/plugins/README.md b/plugins/README.md index 21c48c4..62f0814 100644 --- a/plugins/README.md +++ b/plugins/README.md @@ -312,8 +312,10 @@ plugins/neoforge/gradlew -p plugins/neoforge build The first build of each mod downloads and remaps/decompiles Minecraft, so it takes a few minutes; subsequent builds are fast. Jars land in each module's `build/libs`. CI runs -both gates: `bash plugins/test.sh` (JDK 25 — the install-time plugins plus the -codec/invite/server-list tests) and `bash plugins/test-mods.sh` (JDK 17 — the three +both gates: `bash plugins/test.sh` (JDK 25 — the install-time plugins, the +codec/invite/server-list tests, and the proxy routing tests that drive ServerRegistry, +WaitingRouter and ControlChannel against the real velocity-api; run those alone with +`plugins/velocity/gradlew -p plugins/velocity routingTest`) and `bash plugins/test-mods.sh` (JDK 17 — the three loader mods, via the wrappers above). ### Dependency verification diff --git a/plugins/test.sh b/plugins/test.sh index f6d0d41..78c3e16 100644 --- a/plugins/test.sh +++ b/plugins/test.sh @@ -3,7 +3,7 @@ # # bash plugins/test.sh # -# Two gates, both runnable on any machine with a JDK 25 (Gradle comes from each +# Three gates, all runnable on any machine with a JDK 25 (Gradle comes from each # module's wrapper, sha256-pinned): # # 1. The hand-written, framework-free test mains under shared/test and @@ -37,6 +37,15 @@ # install-time surprise. limbo compiles against the API release the login gate # bundles: deploy/game-stack.lock's LIMBO_VERSION, which bootstrap passes too. # +# 3. The velocity routing self-tests (`./gradlew routingTest`): ServerRegistry, +# WaitingRouter and ControlChannel run against the real velocity-api with a +# fake proxy and a stub felis-api — a refresh registers, moves and drops +# backends and lets go of a renamed subdomain, the login gate and host routing +# admit only linked players, each wake refusal reaches the player as its own +# message, felis:control acts only for the connection's player and holds its +# frame budget. They ride the module's verified dependency set, which is why +# they live in Gradle rather than in the javac mains above. +# # No test framework: the mains are the same javac one-liners their javadocs document, # so a local run and CI run the same bytes. set -euo pipefail @@ -168,6 +177,9 @@ for module in velocity paper; do ( cd "plugins/$module" && ./gradlew --no-daemon build ) done +echo "==> plugins/velocity: ./gradlew --no-daemon routingTest" +( cd plugins/velocity && ./gradlew --no-daemon routingTest ) + limbo_version="$(sed -n 's/^LIMBO_VERSION=//p' deploy/game-stack.lock)" [ -n "$limbo_version" ] || { echo "deploy/game-stack.lock sets no LIMBO_VERSION" >&2; exit 1; } echo "==> plugins/limbo: ./gradlew --no-daemon -PlimboVersion=${limbo_version} build" diff --git a/plugins/velocity/build.gradle b/plugins/velocity/build.gradle index c27f7cc..b070143 100644 --- a/plugins/velocity/build.gradle +++ b/plugins/velocity/build.gradle @@ -68,6 +68,42 @@ tasks.withType(JavaCompile).configureEach { options.encoding = 'UTF-8' } +// Routing self-tests: ServerRegistry and WaitingRouter run unchanged against the real +// velocity-api types (faked with java.lang.reflect.Proxy, test/.../Fakes.java) and a +// stub felis-api on a loopback port. They reuse the compile classpath, which dependency +// verification has already checked, so they add no artifact. Kept out of `build` on +// purpose: bootstrap runs `gradle clean build` on every install, and a test there would +// only slow the install down. plugins/test.sh runs `./gradlew routingTest`. +sourceSets { + routingTest { + java { + srcDir 'test' + include 'best/lolicon/felis/velocity/Fakes.java' + include 'best/lolicon/felis/velocity/ServerRegistryTest.java' + include 'best/lolicon/felis/velocity/WaitingRouterTest.java' + include 'best/lolicon/felis/velocity/ControlChannelTest.java' + } + compileClasspath += sourceSets.main.output + configurations.compileClasspath + runtimeClasspath += output + compileClasspath + } +} + +def routingMains = ['ServerRegistryTest', 'WaitingRouterTest', 'ControlChannelTest'] +routingMains.each { name -> + tasks.register(name, JavaExec) { + group = 'verification' + description = "Runs the ${name} self-test main." + classpath = sourceSets.routingTest.runtimeClasspath + mainClass = "best.lolicon.felis.velocity.${name}" + } +} + +tasks.register('routingTest') { + group = 'verification' + description = 'Runs the routing self-tests (ServerRegistry, WaitingRouter, ControlChannel).' + dependsOn routingMains +} + // Byte-identical jars from identical sources. deploy/bootstrap.sh restarts the proxy only when // a file it runs changed (velocity_fingerprint), and build-time timestamps in the jar would // make every rebuild look like a change and disconnect every player on each rerun. diff --git a/plugins/velocity/src/main/java/best/lolicon/felis/velocity/ServerRegistry.java b/plugins/velocity/src/main/java/best/lolicon/felis/velocity/ServerRegistry.java index 2b66383..25fc0fb 100644 --- a/plugins/velocity/src/main/java/best/lolicon/felis/velocity/ServerRegistry.java +++ b/plugins/velocity/src/main/java/best/lolicon/felis/velocity/ServerRegistry.java @@ -10,6 +10,7 @@ import org.slf4j.Logger; import java.net.InetSocketAddress; import java.util.ArrayList; import java.util.Collection; +import java.util.HashMap; import java.util.HashSet; import java.util.Locale; import java.util.Map; @@ -55,6 +56,7 @@ final class ServerRegistry { /** refresh reconciles registrations against a freshly fetched server list. */ void refresh(Collection servers) { Set seen = new HashSet<>(); + Map subdomains = new HashMap<>(); for (ServerView v : servers) { String name = v.name(); if (name == null || name.isEmpty()) { @@ -64,7 +66,7 @@ final class ServerRegistry { byName.put(name, v); String sub = v.subdomain(); if (sub != null && !sub.isEmpty()) { - subdomainToName.put(sub.toLowerCase(Locale.ROOT), name); + subdomains.put(sub.toLowerCase(Locale.ROOT), name); } ensureRegistered(v); } @@ -75,7 +77,11 @@ final class ServerRegistry { deregister(name); } } - subdomainToName.values().removeIf(n -> !seen.contains(n)); + // The fetch is the whole truth for host routing too: a subdomain a server gave + // up (renamed, or gone with the server) stops routing, not only one another + // server took over. + subdomainToName.putAll(subdomains); + subdomainToName.keySet().retainAll(subdomains.keySet()); } private void ensureRegistered(ServerView v) { diff --git a/plugins/velocity/test/best/lolicon/felis/velocity/ControlChannelTest.java b/plugins/velocity/test/best/lolicon/felis/velocity/ControlChannelTest.java new file mode 100644 index 0000000..b512fd9 --- /dev/null +++ b/plugins/velocity/test/best/lolicon/felis/velocity/ControlChannelTest.java @@ -0,0 +1,307 @@ +package best.lolicon.felis.velocity; + +import best.lolicon.felis.link.Control; +import best.lolicon.felis.link.ControlFrame; +import best.lolicon.felis.link.FelisApiClient; +import best.lolicon.felis.link.LinkConfig; +import best.lolicon.felis.link.ServerView; + +import com.velocitypowered.api.event.connection.DisconnectEvent; +import com.velocitypowered.api.event.connection.PluginMessageEvent; +import com.velocitypowered.api.proxy.ServerConnection; +import com.velocitypowered.api.proxy.messages.ChannelIdentifier; +import com.velocitypowered.api.proxy.messages.MinecraftChannelIdentifier; +import com.velocitypowered.api.proxy.server.RegisteredServer; + +import java.net.InetAddress; +import java.net.ServerSocket; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.util.ArrayList; +import java.util.List; + +/** + * ControlChannelTest drives the real ControlChannel (with the real WaitingRouter, + * ServerRegistry, FelisApiClient and call pool behind it) through Velocity's + * PluginMessageEvent, against a fake proxy and a stub felis-api: which frames are + * consumed, which sources may act and as whom, what each lobby frame answers, + * that status answers are shared across players, that a claim refusal never wakes, + * and that the per-player frame budget holds. Framework free: a failed assertion + * throws. + * + *

Run: {@code ./gradlew routingTest} in plugins/velocity. + */ +public final class ControlChannelTest { + + private static final String SERVERS = "/api/v1/internal/servers/"; + + private static int checks; + private static int nextPlayer = 1; + + private static Fakes.Net net; + private static Fakes.Log log; + private static Fakes.Api api; + private static WaitingRouter router; + private static ControlChannel channel; + private static RegisteredServer login; + private static RegisteredServer lobby; + private static RegisteredServer beta; + + public static void main(String[] args) throws Exception { + net = new Fakes.Net(); + log = new Fakes.Log(); + api = new Fakes.Api(); + try { + FelisApiClient client = new FelisApiClient(new LinkConfig(api.url(), "t")); + ServerRegistry reg = new ServerRegistry(net.proxy, log.logger, "mc.test"); + // Deliberately unsorted: the tile list must come back sorted. + reg.refresh(List.of( + view("gamma", false, "10.43.0.5:25565"), + view("login", true, "10.43.0.1:25565"), + view("beta", true, "10.43.0.4:25565"), + view("lobby", true, "10.43.0.2:25565"), + view("alpha", false, "10.43.0.3:25565"))); + FelisVelocityPlugin plugin = + new FelisVelocityPlugin(net.proxy, log.logger, Files.createTempDirectory("felis-control")); + router = new WaitingRouter(net.proxy, log.logger, client, reg, plugin, "login", "lobby"); + channel = new ControlChannel(net.proxy, log.logger, client, router, reg, plugin, "login", "lobby"); + login = net.proxy.getServer("login").orElseThrow(); + lobby = net.proxy.getServer("lobby").orElseThrow(); + beta = net.proxy.getServer("beta").orElseThrow(); + + channel.register(); + assertEq("register opens felis:control", List.of("felis:control"), List.copyOf(net.channels)); + + consumption(); + sources(); + statusQueries(reg, plugin); + claims(); + wakeAndTransfer(); + budget(); + } finally { + api.close(); + } + System.out.println("ControlChannelTest OK (" + checks + " checks)"); + } + + private static void consumption() { + Fakes.FakePlayer p = player(true); + PluginMessageEvent other = message(p.on(lobby), MinecraftChannelIdentifier.create("other", "x"), + Control.encode(ControlFrame.statusQuery("alpha"))); + assertEq("another channel is left alone", true, other.getResult().isAllowed()); + + PluginMessageEvent fromClient = new PluginMessageEvent(p.player, p.on(lobby), ControlChannel.CHANNEL, + Control.encode(ControlFrame.statusQuery("alpha"))); + channel.onPluginMessage(fromClient); + assertEq("a client's frame is consumed", false, fromClient.getResult().isAllowed()); + + PluginMessageEvent junk = message(p.on(lobby), ControlChannel.CHANNEL, + "not a frame".getBytes(StandardCharsets.UTF_8)); + assertEq("a malformed frame is consumed", false, junk.getResult().isAllowed()); + assertEq("... and neither is answered", 0, p.pluginMessages.size()); + assertEq("... nor reaches felis-api", 0, api.count("GET " + SERVERS + "alpha/menu")); + } + + private static void sources() { + // A user backend is running its owner's plugins: refused, logged once a minute. + Fakes.FakePlayer p = player(true); + PluginMessageEvent e = message(p.on(beta), ControlChannel.CHANNEL, + Control.encode(ControlFrame.wakeRequest(p.name, "alpha"))); + assertEq("user backend: consumed", false, e.getResult().isAllowed()); + message(p.on(beta), ControlChannel.CHANNEL, Control.encode(ControlFrame.wakeRequest(p.name, "alpha"))); + assertEq("user backend: logged once", 1, log.count("WARN", "refused felis:control 'WakeRequest' from backend beta")); + assertEq("user backend: not woken", 0, api.count("POST " + SERVERS + "alpha/wake")); + + // LoginRelease belongs to the login gate alone. + message(p.on(lobby), ControlChannel.CHANNEL, Control.encode(ControlFrame.loginRelease(p.name))); + assertEq("lobby's LoginRelease: nothing moved", 0, p.connects.size()); + assertEq("lobby's LoginRelease: logged", 1, log.count("WARN", "refused felis:control 'LoginRelease' from backend lobby")); + message(p.on(login), ControlChannel.CHANNEL, Control.encode(ControlFrame.loginRelease(p.name))); + assertEq("login's LoginRelease: moved to the lobby", List.of("lobby"), List.copyOf(p.connects)); + + // The lobby may only name managed user servers; a StatusQuery gets an answer + // anyway so the tile stops loading, without echoing a malformed name. + Fakes.FakePlayer q = player(true); + message(q.on(lobby), ControlChannel.CHANNEL, Control.encode(ControlFrame.statusQuery("nope"))); + message(q.on(lobby), ControlChannel.CHANNEL, Control.encode(ControlFrame.statusQuery("lobby"))); + message(q.on(lobby), ControlChannel.CHANNEL, Control.encode(ControlFrame.statusQuery("../x"))); + message(q.on(lobby), ControlChannel.CHANNEL, Control.encode(ControlFrame.wakeRequest(q.name, "nope"))); + List r = frames(q); + assertEq("bad server: three answers, none for the wake", 3, r.size()); + assertEq("unknown server", "Error not_found nope", brief(r.get(0))); + assertEq("system server", "Error not_found lobby", brief(r.get(1))); + assertEq("malformed name not echoed", "Error not_found null", brief(r.get(2))); + assertEq("bad server: no felis-api call", 0, api.count("POST " + SERVERS + "nope/wake")); + + // The tile list: user servers only, sorted. + message(q.on(lobby), ControlChannel.CHANNEL, Control.encode(ControlFrame.listRequest())); + List all = frames(q); + ControlFrame list = all.get(all.size() - 1); + assertEq("list type", ControlFrame.LIST_UPDATE, list.type()); + assertEq("list: user servers, sorted", List.of("alpha", "beta", "gamma"), list.servers()); + } + + private static void statusQueries(ServerRegistry reg, FelisVelocityPlugin plugin) throws Exception { + api.claimable.add("alpha"); + Fakes.FakePlayer p = player(true); + message(p.on(lobby), ControlChannel.CHANNEL, Control.encode(ControlFrame.statusQuery("alpha"))); + Fakes.await("status answered", () -> p.pluginMessages.size() == 1); + ControlFrame s = frames(p).get(0); + assertEq("status type", ControlFrame.STATUS_UPDATE, s.type()); + assertEq("status fields", "alpha Stopped false 3/20 claimable=true", + s.server() + " " + s.phase() + " " + s.ready() + " " + s.playersOnline() + "/" + s.playersMax() + + " claimable=" + s.claimable()); + + // A second player within the TTL is answered from the shared projection, at once. + Fakes.FakePlayer q = player(true); + message(q.on(lobby), ControlChannel.CHANNEL, Control.encode(ControlFrame.statusQuery("alpha"))); + assertEq("shared answer is immediate", 1, q.pluginMessages.size()); + assertEq("shared answer is the same", "alpha Stopped false", brief3(frames(q).get(0))); + assertEq("felis-api read once", 1, api.count("GET " + SERVERS + "alpha/menu")); + + // felis-api's own refusal comes back as its code and message. + api.menuError.put("gamma", "404 not_found"); + message(p.on(lobby), ControlChannel.CHANNEL, Control.encode(ControlFrame.statusQuery("gamma"))); + Fakes.await("refusal answered", () -> p.pluginMessages.size() == 2); + ControlFrame refused = frames(p).get(1); + assertEq("refusal", "Error not_found gamma", brief(refused)); + assertEq("refusal message", "stub not_found", refused.message()); + + // A transport failure never shows its internals (host, port) to the player. + int closed; + try (ServerSocket probe = new ServerSocket(0, 1, InetAddress.getLoopbackAddress())) { + closed = probe.getLocalPort(); + } + FelisApiClient dead = new FelisApiClient(new LinkConfig("http://127.0.0.1:" + closed, "t")); + ControlChannel cut = new ControlChannel(net.proxy, log.logger, dead, router, reg, plugin, "login", "lobby"); + Fakes.FakePlayer t = player(true); + cut.onPluginMessage(new PluginMessageEvent(t.on(lobby), t.player, ControlChannel.CHANNEL, + Control.encode(ControlFrame.statusQuery("beta")))); + Fakes.await("transport failure answered", () -> t.pluginMessages.size() == 1); + ControlFrame down = frames(t).get(0); + assertEq("transport failure code", "Error transport_error beta", brief(down)); + assertEq("transport failure message", + "Felis 暂时不可用,请稍后再试 / Felis is temporarily unavailable — please try again.", down.message()); + } + + private static void claims() { + // A refused claim answers with the refusal and never wakes. + api.claimError.put("gamma", "403 quota_exceeded"); + Fakes.FakePlayer p = player(true); + p.current = lobby; + message(p.on(lobby), ControlChannel.CHANNEL, Control.encode(ControlFrame.claimRequest("someone-else", "gamma"))); + Fakes.await("claim refusal answered", () -> p.pluginMessages.size() == 1); + assertEq("claim refused", "Error quota_exceeded gamma", brief(frames(p).get(0))); + assertEq("claim refused: not woken", 0, api.count("POST " + SERVERS + "gamma/wake")); + assertEq("the claim acted as the connection's player", + true, api.bodies.get("POST " + SERVERS + "gamma/claim").contains(p.id.toString())); + + // A granted claim wakes and parks, and drops the stale "claimable" answer. Warm + // the shared answer right before, so only the claim can have dropped it. + Fakes.FakePlayer w = player(true); + int reads = api.count("GET " + SERVERS + "alpha/menu"); + message(w.on(lobby), ControlChannel.CHANNEL, Control.encode(ControlFrame.statusQuery("alpha"))); + Fakes.await("warm status", () -> w.pluginMessages.size() == 1); + int warm = api.count("GET " + SERVERS + "alpha/menu"); + assertEq("warmed (read at most once more)", true, warm == reads || warm == reads + 1); + message(p.on(lobby), ControlChannel.CHANNEL, Control.encode(ControlFrame.claimRequest(p.name, "alpha"))); + Fakes.await("claimed server woken", () -> api.count("POST " + SERVERS + "alpha/wake") == 1); + Fakes.await("claimer queued", () -> router.waitingCount() == 1); + assertEq("the wake acted as the connection's player", + true, api.bodies.get("POST " + SERVERS + "alpha/wake").contains(p.id.toString())); + Fakes.FakePlayer q = player(true); + message(q.on(lobby), ControlChannel.CHANNEL, Control.encode(ControlFrame.statusQuery("alpha"))); + Fakes.await("fresh status after the claim", () -> q.pluginMessages.size() == 1); + assertEq("status re-read after the claim", warm + 1, api.count("GET " + SERVERS + "alpha/menu")); + } + + private static void wakeAndTransfer() { + // A WakeRequest parks the player; when the backend is ready the lobby hears + // TransferReady just before the proxy moves them. + Fakes.FakePlayer p = player(true); + p.current = lobby; + message(p.on(lobby), ControlChannel.CHANNEL, Control.encode(ControlFrame.wakeRequest("forged-name", "beta"))); + Fakes.await("beta woken", () -> api.count("POST " + SERVERS + "beta/wake") == 1); + Fakes.await("waker queued", () -> router.waitingCount() == 2); + api.ready.put("beta", true); + router.tick(); + List r = frames(p); + assertEq("one frame to the lobby", 1, r.size()); + assertEq("transfer ready", ControlFrame.TRANSFER_READY + " " + p.name + " beta", + r.get(0).type() + " " + r.get(0).player() + " " + r.get(0).server()); + assertEq("then moved", List.of("beta"), List.copyOf(p.connects)); + } + + private static void budget() { + // 96 frames in a burst, then 10 a second: a flood is cut off, the connection + // is not. ListRequest is answered inline, so every accepted frame shows at once. + Fakes.FakePlayer p = player(true); + for (int i = 0; i < 150; i++) { + message(p.on(lobby), ControlChannel.CHANNEL, Control.encode(ControlFrame.listRequest())); + } + int answered = p.pluginMessages.size(); + assertEq("flood cut at the burst (plus what refilled meanwhile): " + answered, + true, answered >= 96 && answered <= 99); + + // Another player has a budget of their own. + Fakes.FakePlayer q = player(true); + message(q.on(lobby), ControlChannel.CHANNEL, Control.encode(ControlFrame.listRequest())); + assertEq("other player unaffected", 1, q.pluginMessages.size()); + + // Leaving the proxy forgets the budget: a fresh burst, far more than a refill. + channel.onDisconnect(new DisconnectEvent(p.player, DisconnectEvent.LoginStatus.SUCCESSFUL_LOGIN)); + p.pluginMessages.clear(); + for (int i = 0; i < 50; i++) { + message(p.on(lobby), ControlChannel.CHANNEL, Control.encode(ControlFrame.listRequest())); + } + assertEq("budget forgotten on disconnect", 50, p.pluginMessages.size()); + } + + // ---- helpers ---- + + private static Fakes.FakePlayer player(boolean linked) { + int n = nextPlayer++; + Fakes.FakePlayer p = new Fakes.FakePlayer(net, "player" + n, n); + if (linked) { + api.linked.add(p.id); + } + return p; + } + + // message delivers a frame a backend sent up its connection toward its player. + private static PluginMessageEvent message(ServerConnection source, ChannelIdentifier id, byte[] data) { + PluginMessageEvent e = new PluginMessageEvent(source, source.getPlayer(), id, data); + channel.onPluginMessage(e); + return e; + } + + private static List frames(Fakes.FakePlayer p) { + List out = new ArrayList<>(); + synchronized (p.pluginMessages) { + for (byte[] b : p.pluginMessages) { + out.add(Control.decode(b)); + } + } + return out; + } + + private static String brief(ControlFrame f) { + return f.type() + " " + f.code() + " " + f.server(); + } + + private static String brief3(ControlFrame f) { + return f.server() + " " + f.phase() + " " + f.ready(); + } + + private static ServerView view(String name, boolean ready, String addr) { + return new ServerView(name, name, ready ? "Running" : "Stopped", ready, "ownerOnly", + "Running", "ClusterIP", addr, 0, 20); + } + + private static void assertEq(String what, Object want, Object got) { + if (want == null ? got != null : !want.equals(got)) { + throw new AssertionError(what + ": got " + got + ", want " + want); + } + checks++; + } +} diff --git a/plugins/velocity/test/best/lolicon/felis/velocity/Fakes.java b/plugins/velocity/test/best/lolicon/felis/velocity/Fakes.java new file mode 100644 index 0000000..bcd35ff --- /dev/null +++ b/plugins/velocity/test/best/lolicon/felis/velocity/Fakes.java @@ -0,0 +1,542 @@ +package best.lolicon.felis.velocity; + +import com.sun.net.httpserver.HttpExchange; +import com.sun.net.httpserver.HttpServer; +import com.velocitypowered.api.proxy.ConnectionRequestBuilder; +import com.velocitypowered.api.proxy.Player; +import com.velocitypowered.api.proxy.ProxyServer; +import com.velocitypowered.api.proxy.ServerConnection; +import com.velocitypowered.api.proxy.messages.ChannelIdentifier; +import com.velocitypowered.api.proxy.messages.ChannelRegistrar; +import com.velocitypowered.api.proxy.player.PlayerSettings; +import com.velocitypowered.api.proxy.server.RegisteredServer; +import com.velocitypowered.api.proxy.server.ServerInfo; +import net.kyori.adventure.text.Component; +import net.kyori.adventure.text.TextComponent; +import org.slf4j.Logger; + +import java.io.IOException; +import java.io.OutputStream; +import java.lang.reflect.InvocationHandler; +import java.lang.reflect.Proxy; +import java.net.InetSocketAddress; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.Optional; +import java.util.Set; +import java.util.UUID; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentHashMap; +import java.util.function.BooleanSupplier; + +/** + * Fakes are the test doubles the routing tests run the real proxy classes against: + * the velocity-api interfaces they touch (ProxyServer, Player, RegisteredServer, the + * slf4j Logger) built with {@link Proxy}, and a stub felis-api on a loopback port. + * They behave the way the real thing does where the code under test relies on it: + * server names are case-insensitive, registering a taken name at another address + * throws, and unregistering needs the exact ServerInfo that was registered. + */ +final class Fakes { + private Fakes() { + } + + /** Returned by an {@link Answers} for a method it does not answer. */ + static final Object UNANSWERED = new Object(); + + interface Answers { + Object answer(String method, Object[] args) throws Throwable; + } + + /** + * fake builds an instance of an interface. Methods the answers leave unanswered run + * the interface's default body when it has one and return a zero value otherwise + * (Optional.empty() for an Optional). + */ + static T fake(Class type, Answers answers) { + InvocationHandler h = (proxy, m, args) -> { + Object[] a = args == null ? new Object[0] : args; + if (m.getDeclaringClass() == Object.class) { + switch (m.getName()) { + case "equals": + return proxy == a[0]; + case "hashCode": + return System.identityHashCode(proxy); + default: + return type.getSimpleName() + "@fake"; + } + } + Object r = answers.answer(m.getName(), a); + if (r != UNANSWERED) { + return r; + } + if (m.isDefault()) { + return InvocationHandler.invokeDefault(proxy, m, args); + } + return zero(m.getReturnType()); + }; + return type.cast(Proxy.newProxyInstance(type.getClassLoader(), new Class[]{type}, h)); + } + + private static Object zero(Class t) { + if (t == boolean.class) { + return false; + } + if (t == int.class || t == short.class || t == byte.class) { + return 0; + } + if (t == long.class) { + return 0L; + } + if (t == double.class || t == float.class) { + return 0.0; + } + if (t == Optional.class) { + return Optional.empty(); + } + return null; + } + + /** text is the plain text of a component the code under test built with Component.text. */ + static String text(Component c) { + StringBuilder sb = new StringBuilder(); + if (c instanceof TextComponent) { + sb.append(((TextComponent) c).content()); + } + for (Component child : c.children()) { + sb.append(text(child)); + } + return sb.toString(); + } + + /** await polls until the condition holds, failing after two seconds (async work). */ + static void await(String what, BooleanSupplier condition) { + long end = System.nanoTime() + 2_000_000_000L; + while (!condition.getAsBoolean()) { + if (System.nanoTime() > end) { + throw new AssertionError("timed out waiting for: " + what); + } + try { + Thread.sleep(5); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new AssertionError("interrupted waiting for: " + what); + } + } + } + + // ---- the proxy ---- + + static final class Net { + private final Map servers = new ConcurrentHashMap<>(); + final Map players = new ConcurrentHashMap<>(); + /** registrations logs every register ("+name host:port") and unregister ("-name"). */ + final List registrations = Collections.synchronizedList(new ArrayList<>()); + /** channels lists every plugin-message channel id registered. */ + final List channels = Collections.synchronizedList(new ArrayList<>()); + + private final ChannelRegistrar registrar = fake(ChannelRegistrar.class, (m, a) -> { + if ("register".equals(m)) { + for (Object id : (Object[]) a[0]) { + channels.add(((ChannelIdentifier) id).getId()); + } + return null; + } + return UNANSWERED; + }); + + final ProxyServer proxy = fake(ProxyServer.class, (m, a) -> { + switch (m) { + case "getServer": + return Optional.ofNullable(servers.get(key((String) a[0]))); + case "registerServer": + return register((ServerInfo) a[0]); + case "unregisterServer": + unregister((ServerInfo) a[0]); + return null; + case "getPlayer": + return a[0] instanceof UUID + ? Optional.ofNullable(players.get((UUID) a[0])) + : Optional.empty(); + case "getAllServers": + return new ArrayList<>(servers.values()); + case "getChannelRegistrar": + return registrar; + default: + return UNANSWERED; + } + }); + + private static String key(String name) { + return name.toLowerCase(Locale.ROOT); + } + + private RegisteredServer register(ServerInfo info) { + RegisteredServer existing = servers.get(key(info.getName())); + if (existing != null) { + if (existing.getServerInfo().equals(info)) { + return existing; + } + throw new IllegalArgumentException("server " + info.getName() + " is already registered"); + } + RegisteredServer rs = server(info); + servers.put(key(info.getName()), rs); + InetSocketAddress addr = info.getAddress(); + registrations.add("+" + info.getName() + " " + addr.getHostString() + ":" + addr.getPort()); + return rs; + } + + private void unregister(ServerInfo info) { + RegisteredServer existing = servers.get(key(info.getName())); + if (existing == null || !existing.getServerInfo().equals(info)) { + throw new IllegalArgumentException("server " + info.getName() + " is not registered with that info"); + } + servers.remove(key(info.getName())); + registrations.add("-" + info.getName()); + } + + /** add registers a server the way velocity.toml would (not logged in registrations). */ + RegisteredServer add(String name, String host, int port) { + RegisteredServer rs = server(new ServerInfo(name, InetSocketAddress.createUnresolved(host, port))); + servers.put(key(name), rs); + return rs; + } + + void remove(String name) { + servers.remove(key(name)); + } + + /** address is where a registered server points, or null. */ + String address(String name) { + RegisteredServer rs = servers.get(key(name)); + if (rs == null) { + return null; + } + InetSocketAddress a = rs.getServerInfo().getAddress(); + return a.getHostString() + ":" + a.getPort(); + } + } + + static RegisteredServer server(ServerInfo info) { + return fake(RegisteredServer.class, (m, a) -> "getServerInfo".equals(m) ? info : UNANSWERED); + } + + // ---- players ---- + + static final class FakePlayer { + final UUID id; + final String name; + volatile Locale locale = Locale.ENGLISH; + volatile String virtualHost; + volatile RegisteredServer current; + volatile boolean connectSucceeds = true; + volatile String disconnectedWith; + final List messages = Collections.synchronizedList(new ArrayList<>()); + /** connects lists every server a connection request was sent to, by name. */ + final List connects = Collections.synchronizedList(new ArrayList<>()); + /** pluginMessages lists every plugin message sent down any of this player's server connections. */ + final List pluginMessages = Collections.synchronizedList(new ArrayList<>()); + + final Player player; + + FakePlayer(Net net, String name, int n) { + this.id = UUID.fromString(String.format("00000000-0000-0000-0000-%012d", n)); + this.name = name; + PlayerSettings settings = fake(PlayerSettings.class, + (m, a) -> "getLocale".equals(m) ? locale : UNANSWERED); + this.player = fake(Player.class, (m, a) -> { + switch (m) { + case "getUniqueId": + return id; + case "getUsername": + return name; + case "getPlayerSettings": + return settings; + case "sendMessage": + for (Object o : a) { + if (o instanceof Component) { + messages.add(text((Component) o)); + } + } + return null; + case "disconnect": + disconnectedWith = text((Component) a[0]); + return null; + case "getVirtualHost": + String host = virtualHost; + return host == null + ? Optional.empty() + : Optional.of(InetSocketAddress.createUnresolved(host, 25565)); + case "getCurrentServer": + RegisteredServer on = current; + return on == null ? Optional.empty() : Optional.of(on(on)); + case "createConnectionRequest": + return request((RegisteredServer) a[0]); + default: + return UNANSWERED; + } + }); + net.players.put(id, player); + } + + /** on is this player's connection to a backend, as a plugin message arrives on it. */ + ServerConnection on(RegisteredServer on) { + return fake(ServerConnection.class, (m, a) -> { + switch (m) { + case "getServer": + return on; + case "getServerInfo": + return on.getServerInfo(); + case "getPlayer": + return player; + case "sendPluginMessage": + pluginMessages.add((byte[]) a[1]); + return true; + default: + return UNANSWERED; + } + }); + } + + private ConnectionRequestBuilder request(RegisteredServer target) { + return fake(ConnectionRequestBuilder.class, (m, a) -> { + switch (m) { + case "getServer": + return target; + case "connect": + connects.add(target.getServerInfo().getName()); + boolean ok = connectSucceeds; + ConnectionRequestBuilder.Status status = ok + ? ConnectionRequestBuilder.Status.SUCCESS + : ConnectionRequestBuilder.Status.SERVER_DISCONNECTED; + return CompletableFuture.completedFuture(fake(ConnectionRequestBuilder.Result.class, + (m2, a2) -> { + switch (m2) { + case "getStatus": + return status; + case "getAttemptedConnection": + return target; + default: + return UNANSWERED; + } + })); + default: + return UNANSWERED; + } + }); + } + + /** said reports whether any chat line so far contains the text. */ + boolean said(String part) { + synchronized (messages) { + for (String m : messages) { + if (m.contains(part)) { + return true; + } + } + } + return false; + } + } + + // ---- logging ---- + + static final class Log { + final List lines = Collections.synchronizedList(new ArrayList<>()); + + final Logger logger = fake(Logger.class, (m, a) -> { + switch (m) { + case "trace": + case "debug": + case "info": + case "warn": + case "error": + if (a.length > 0 && a[0] instanceof String) { + lines.add(m.toUpperCase(Locale.ROOT) + " " + format((String) a[0], a)); + } + return null; + case "isTraceEnabled": + case "isDebugEnabled": + case "isInfoEnabled": + case "isWarnEnabled": + case "isErrorEnabled": + return true; + default: + return UNANSWERED; + } + }); + + // format fills slf4j {} placeholders from the rest of the call's arguments + // (the varargs form arrives as one Object[]). + private static String format(String pattern, Object[] a) { + List args = new ArrayList<>(); + for (int i = 1; i < a.length; i++) { + if (a[i] instanceof Object[]) { + Collections.addAll(args, (Object[]) a[i]); + } else { + args.add(a[i]); + } + } + StringBuilder sb = new StringBuilder(); + int from = 0; + int next = 0; + int at; + while ((at = pattern.indexOf("{}", from)) >= 0) { + sb.append(pattern, from, at).append(next < args.size() ? args.get(next++) : "{}"); + from = at + 2; + } + return sb.append(pattern.substring(from)).toString(); + } + + /** count is how many lines at this level contain the text. */ + int count(String level, String part) { + int n = 0; + synchronized (lines) { + for (String l : lines) { + if (l.startsWith(level + " ") && l.contains(part)) { + n++; + } + } + } + return n; + } + } + + // ---- felis-api ---- + + /** + * Api is a stub felis-api internal face: link status from {@link #linked}, server + * status from {@link #ready}, wake answering 202 unless {@link #wakeError} holds an + * answer for the server, and join-events recorded. {@link #linkDown} makes the + * link-status route answer 500. + */ + static final class Api implements AutoCloseable { + final Set linked = ConcurrentHashMap.newKeySet(); + volatile boolean linkDown; + final Map ready = new ConcurrentHashMap<>(); + /** wakeError maps a server to "status code" (e.g. "403 forbidden"). */ + final Map wakeError = new ConcurrentHashMap<>(); + volatile int joinStatus = 204; + /** claimable names the servers the menu route reports as claimable. */ + final Set claimable = ConcurrentHashMap.newKeySet(); + /** menuError and claimError map a server to "status code", like wakeError. */ + final Map menuError = new ConcurrentHashMap<>(); + final Map claimError = 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. */ + final List calls = Collections.synchronizedList(new ArrayList<>()); + + private final HttpServer http; + + Api() throws IOException { + http = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); + http.createContext("/", this::handle); + http.start(); + } + + String url() { + return "http://127.0.0.1:" + http.getAddress().getPort(); + } + + private void handle(HttpExchange ex) throws IOException { + String method = ex.getRequestMethod(); + String path = ex.getRequestURI().getPath(); + bodies.put(method + " " + path, new String(ex.getRequestBody().readAllBytes(), StandardCharsets.UTF_8)); + calls.add(method + " " + path); + String status = "/api/v1/internal/account/link/status/"; + String servers = "/api/v1/internal/servers/"; + if (path.startsWith(status)) { + if (linkDown) { + reply(ex, 500, "{\"error\":{\"code\":\"internal\",\"message\":\"down\"}}"); + return; + } + boolean yes = linked.contains(UUID.fromString(path.substring(status.length()))); + reply(ex, 200, "{\"linked\":" + yes + "}"); + return; + } + if (path.startsWith(servers)) { + String[] parts = path.substring(servers.length()).split("/"); + String name = parts[0]; + String action = parts.length > 1 ? parts[1] : ""; + switch (method + " " + action) { + case "GET status": + reply(ex, 200, view(name, ready.getOrDefault(name, false))); + return; + case "POST wake": + if (!error(ex, wakeError.get(name))) { + reply(ex, 202, view(name, false)); + } + return; + case "GET menu": + if (!error(ex, menuError.get(name))) { + reply(ex, 200, "{\"name\":\"" + name + "\",\"phase\":\"Stopped\",\"ready\":false," + + "\"playersOnline\":3,\"playersMax\":20,\"claimable\":" + claimable.contains(name) + "}"); + } + return; + case "POST claim": + if (!error(ex, claimError.get(name))) { + reply(ex, 200, "{\"claimed\":true}"); + } + return; + case "POST join-event": + reply(ex, joinStatus, joinStatus == 204 ? "" : "{\"error\":{\"code\":\"internal\"}}"); + return; + default: + break; + } + } + reply(ex, 404, "{\"error\":{\"code\":\"not_found\",\"message\":\"stub\"}}"); + } + + // error answers with a {"error":{code,message}} envelope when one is configured. + private static boolean error(HttpExchange ex, String statusAndCode) throws IOException { + if (statusAndCode == null) { + return false; + } + String[] sc = statusAndCode.split(" ", 2); + reply(ex, Integer.parseInt(sc[0]), "{\"error\":{\"code\":\"" + sc[1] + "\",\"message\":\"stub " + sc[1] + "\"}}"); + return true; + } + + private static String view(String name, boolean ready) { + return "{\"name\":\"" + name + "\",\"subdomain\":\"" + name + "\",\"ready\":" + ready + "}"; + } + + private static void reply(HttpExchange ex, int status, String body) throws IOException { + byte[] b = body.getBytes(StandardCharsets.UTF_8); + if (b.length == 0) { + ex.sendResponseHeaders(status, -1); + ex.close(); + return; + } + ex.getResponseHeaders().set("Content-Type", "application/json"); + ex.sendResponseHeaders(status, b.length); + try (OutputStream os = ex.getResponseBody()) { + os.write(b); + } + } + + /** count is how many requests so far match "METHOD path". */ + int count(String call) { + int n = 0; + synchronized (calls) { + for (String c : calls) { + if (c.equals(call)) { + n++; + } + } + } + return n; + } + + @Override + public void close() { + http.stop(0); + } + } +} diff --git a/plugins/velocity/test/best/lolicon/felis/velocity/ServerRegistryTest.java b/plugins/velocity/test/best/lolicon/felis/velocity/ServerRegistryTest.java new file mode 100644 index 0000000..ee57f93 --- /dev/null +++ b/plugins/velocity/test/best/lolicon/felis/velocity/ServerRegistryTest.java @@ -0,0 +1,128 @@ +package best.lolicon.felis.velocity; + +import best.lolicon.felis.link.ServerView; + +import java.net.InetSocketAddress; +import java.util.List; +import java.util.Optional; + +/** + * ServerRegistryTest drives the real ServerRegistry against a fake ProxyServer: a + * refresh registers every server that has a backend address, leaves one registered at + * the same address alone, moves one whose address changed, drops one that vanished + * from a successful fetch, and host routing resolves only {@code .}. + * Framework free: a failed assertion throws. + * + *

Run: {@code ./gradlew routingTest} in plugins/velocity (it needs the velocity-api + * classes the plugin compiles against). + */ +public final class ServerRegistryTest { + + private static int checks; + + public static void main(String[] args) { + Fakes.Net net = new Fakes.Net(); + Fakes.Log log = new Fakes.Log(); + // velocity.toml's dead login placeholder, which the first refresh must replace. + net.add("login", "127.0.0.1", 1); + ServerRegistry reg = new ServerRegistry(net.proxy, log.logger, "MC.Example.test"); + + reg.refresh(List.of( + view("login", "login", true, "10.43.0.1:25565"), + view("lobby", "lobby", true, "10.43.0.2:25565"), + view("alpha", "alpha", false, "10.43.0.3:25566"), + view("fresh", "fresh", false, null), + view("", "blank", true, "10.43.0.9:25565"))); + assertEq("first refresh registrations", + List.of("-login", "+login 10.43.0.1:25565", "+lobby 10.43.0.2:25565", "+alpha 10.43.0.3:25566"), + List.copyOf(net.registrations)); + assertEq("login placeholder replaced", "10.43.0.1:25565", net.address("login")); + assertEq("a server with no address is known", true, reg.isManaged("fresh")); + assertEq("... but not registered", null, net.address("fresh")); + assertEq("a nameless entry is skipped", false, reg.isManaged("")); + assertEq("managed count", 4, reg.all().size()); + assertEq("registration logged", 1, log.count("INFO", "registered backend alpha -> 10.43.0.3:25566")); + + net.registrations.clear(); + reg.refresh(List.of( + view("login", "login", true, "10.43.0.1:25565"), + view("lobby", "lobby", true, "10.43.0.2:25565"), + view("alpha", "alpha", true, "10.43.0.7:25566"), + view("fresh", "fresh", true, "10.43.0.8:25565"))); + assertEq("second refresh: only the moved and the new", + List.of("-alpha", "+alpha 10.43.0.7:25566", "+fresh 10.43.0.8:25565"), + List.copyOf(net.registrations)); + assertEq("view follows the fetch", true, reg.view("alpha").ready()); + + assertEq("host routing", "alpha", name(reg.resolveByHost("alpha.mc.example.test"))); + assertEq("host routing ignores case", "alpha", name(reg.resolveByHost("ALPHA.Mc.Example.Test"))); + assertEq("the root itself routes nowhere", null, name(reg.resolveByHost("mc.example.test"))); + assertEq("a look-alike suffix routes nowhere", null, name(reg.resolveByHost("alphamc.example.test"))); + assertEq("another domain routes nowhere", null, name(reg.resolveByHost("alpha.mc.example.test.evil"))); + assertEq("an unknown subdomain routes nowhere", null, name(reg.resolveByHost("nope.mc.example.test"))); + assertEq("no host routes nowhere", null, name(reg.resolveByHost(null))); + + net.registrations.clear(); + reg.refresh(List.of( + view("login", "login", true, "10.43.0.1:25565"), + view("lobby", "lobby", true, "10.43.0.2:25565"), + view("fresh", "fresh", true, "10.43.0.8:25565"))); + assertEq("a vanished server is dropped", List.of("-alpha"), List.copyOf(net.registrations)); + assertEq("... from the views", false, reg.isManaged("alpha")); + assertEq("... and from host routing", null, name(reg.resolveByHost("alpha.mc.example.test"))); + assertEq("deregistration logged", 1, log.count("INFO", "deregistered backend alpha")); + + // A server whose subdomain changed answers on the new one only. + reg.refresh(List.of( + view("login", "login", true, "10.43.0.1:25565"), + view("lobby", "lobby", true, "10.43.0.2:25565"), + view("fresh", "renamed", true, "10.43.0.8:25565"))); + assertEq("new subdomain", "fresh", name(reg.resolveByHost("renamed.mc.example.test"))); + assertEq("old subdomain let go", null, name(reg.resolveByHost("fresh.mc.example.test"))); + + // A subdomain handed to another server routes to the new holder. + reg.refresh(List.of( + view("login", "login", true, "10.43.0.1:25565"), + view("lobby", "lobby", true, "10.43.0.2:25565"), + view("other", "renamed", true, "10.43.0.5:25565"), + view("fresh", "fresh", true, "10.43.0.8:25565"))); + assertEq("subdomain moved", "other", name(reg.resolveByHost("renamed.mc.example.test"))); + assertEq("subdomain taken back", "fresh", name(reg.resolveByHost("fresh.mc.example.test"))); + + // An empty successful fetch really does empty the proxy. + net.registrations.clear(); + reg.refresh(List.of()); + assertEq("everything deregistered", 4, net.registrations.size()); + assertEq("nothing managed", 0, reg.all().size()); + + assertAddr("host:port", "10.43.0.1", 25570, ServerRegistry.parseAddress("10.43.0.1:25570")); + assertAddr("bare host takes the Minecraft port", "backend.svc", 25565, + ServerRegistry.parseAddress("backend.svc")); + assertAddr("a non-numeric port is not a port", "backend:x", 25565, + ServerRegistry.parseAddress("backend:x")); + + System.out.println("ServerRegistryTest OK (" + checks + " checks)"); + } + + private static ServerView view(String name, String sub, boolean ready, String addr) { + return new ServerView(name, sub, ready ? "Running" : "Stopped", ready, "ownerOnly", + "Running", "ClusterIP", addr, 0, 20); + } + + private static String name(Optional v) { + return v.map(ServerView::name).orElse(null); + } + + private static void assertAddr(String what, String host, int port, InetSocketAddress got) { + assertEq(what + " host", host, got.getHostString()); + assertEq(what + " port", port, got.getPort()); + assertEq(what + " stays unresolved", true, got.isUnresolved()); + } + + private static void assertEq(String what, Object want, Object got) { + if (want == null ? got != null : !want.equals(got)) { + throw new AssertionError(what + ": got " + got + ", want " + want); + } + checks++; + } +} diff --git a/plugins/velocity/test/best/lolicon/felis/velocity/WaitingRouterTest.java b/plugins/velocity/test/best/lolicon/felis/velocity/WaitingRouterTest.java new file mode 100644 index 0000000..b26c44d --- /dev/null +++ b/plugins/velocity/test/best/lolicon/felis/velocity/WaitingRouterTest.java @@ -0,0 +1,488 @@ +package best.lolicon.felis.velocity; + +import best.lolicon.felis.link.FelisApiClient; +import best.lolicon.felis.link.LinkConfig; +import best.lolicon.felis.link.ServerView; + +import com.velocitypowered.api.event.Continuation; +import com.velocitypowered.api.event.EventTask; +import com.velocitypowered.api.event.connection.DisconnectEvent; +import com.velocitypowered.api.event.player.PlayerChooseInitialServerEvent; +import com.velocitypowered.api.event.player.ServerConnectedEvent; +import com.velocitypowered.api.event.player.ServerPreConnectEvent; +import com.velocitypowered.api.proxy.server.RegisteredServer; + +import java.nio.file.Files; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Locale; + +/** + * WaitingRouterTest drives the real WaitingRouter, ServerRegistry, FelisApiClient and + * FelisVelocityPlugin call pool through Velocity's own event objects, against a fake + * proxy and a stub felis-api: where a connecting host sends a player, what the login + * gate may release them to and when, how each wake refusal is told apart, how the + * waiting queue is drained (one status poll per server per tick, only linked players + * moved), and which joins are reported. Framework free: a failed assertion throws. + * + *

Not covered here: the 120 s queue timeout (the router reads the wall clock) and a + * real client's timing, which needs an online-mode account on a live stack. + * + *

Run: {@code ./gradlew routingTest} in plugins/velocity. + */ +public final class WaitingRouterTest { + + private static final String LINK = "GET /api/v1/internal/account/link/status/"; + private static final String SERVERS = "/api/v1/internal/servers/"; + + private static int checks; + private static int nextPlayer = 1; + + private static Fakes.Net net; + private static Fakes.Log log; + private static Fakes.Api api; + private static ServerRegistry reg; + private static FelisVelocityPlugin plugin; + private static WaitingRouter router; + private static RegisteredServer login; + private static RegisteredServer lobby; + private static RegisteredServer beta; + + public static void main(String[] args) throws Exception { + net = new Fakes.Net(); + log = new Fakes.Log(); + api = new Fakes.Api(); + try { + FelisApiClient client = new FelisApiClient(new LinkConfig(api.url(), "t")); + reg = new ServerRegistry(net.proxy, log.logger, "mc.test"); + reg.refresh(servers(true)); + plugin = new FelisVelocityPlugin(net.proxy, log.logger, Files.createTempDirectory("felis-routing")); + router = new WaitingRouter(net.proxy, log.logger, client, reg, plugin, "login", "lobby"); + login = net.proxy.getServer("login").orElseThrow(); + lobby = net.proxy.getServer("lobby").orElseThrow(); + beta = net.proxy.getServer("beta").orElseThrow(); + + hostRouting(); + loginGate(); + wakeRefusals(); + queue(); + menuAndCommands(); + joins(); + disconnectAndRelease(); + } finally { + api.close(); + } + System.out.println("WaitingRouterTest OK (" + checks + " checks)"); + } + + // The server list the stub control plane serves. "gone" is dropped by a later + // refresh; "fresh" has no backend address yet. + private static List servers(boolean withGone) { + List list = new ArrayList<>(List.of( + view("login", true, "10.43.0.1:25565"), + view("lobby", true, "10.43.0.2:25565"), + view("alpha", false, "10.43.0.3:25565"), + view("beta", true, "10.43.0.4:25565"), + view("gamma", false, "10.43.0.5:25565"), + view("delta", false, "10.43.0.6:25565"), + view("epsilon", false, "10.43.0.7:25565"), + view("zeta", false, "10.43.0.8:25565"), + view("fresh", false, null))); + if (withGone) { + list.add(view("gone", false, "10.43.0.9:25565")); + } + return list; + } + + private static void hostRouting() { + // A sleeping server: the player lands on login, and the linked release goes to + // the lobby while the server is woken and the player queued. + Fakes.FakePlayer p1 = player("alpha.mc.test", true); + assertEq("fresh connection goes to login", "login", choose(p1)); + ServerPreConnectEvent e = release(p1); + assertEq("sleeping target: released to the lobby", "lobby", allowedTo(e)); + assertEq("sleeping target: woken once", 1, api.count("POST " + SERVERS + "alpha/wake")); + assertEq("sleeping target: told", true, p1.said("Starting « alpha »")); + assertEq("sleeping target: queued", 1, router.waitingCount()); + assertEq("the release checked the link", 1, api.count(LINK + p1.id)); + + // A ready server: the release goes straight to it, no wake; until the connect + // lands, a repeated release (the gate retries) goes there again. + Fakes.FakePlayer p2 = player("BETA.mc.test", true); + choose(p2); + assertEq("ready target: released to it", "beta", allowedTo(release(p2))); + assertEq("ready target: retried release still goes to it", "beta", allowedTo(release(p2))); + assertEq("ready target: no wake", 0, api.count("POST " + SERVERS + "beta/wake")); + router.onServerConnected(new ServerConnectedEvent(p2.player, beta, login)); + assertEq("after the connect the target is spent", "lobby", allowedTo(release(p2))); + + // login. is the gate, never a destination: remembering it would loop. + Fakes.FakePlayer p3 = player("login.mc.test", true); + choose(p3); + assertEq("login host: released to the lobby", "lobby", allowedTo(release(p3))); + assertEq("login host: nothing woken", 0, api.count("POST " + SERVERS + "login/wake")); + + // lobby. is where the release goes anyway; even a lobby felis-api reports + // as not ready is never woken through the queue. + List lobbyDown = servers(true); + lobbyDown.set(1, view("lobby", false, "10.43.0.2:25565")); + reg.refresh(lobbyDown); + Fakes.FakePlayer p4 = player("lobby.mc.test", true); + choose(p4); + assertEq("lobby host: released to the lobby", "lobby", allowedTo(release(p4))); + assertEq("lobby host: not woken", 0, api.count("POST " + SERVERS + "lobby/wake")); + assertEq("lobby host: not queued", 1, router.waitingCount()); + reg.refresh(servers(true)); + + // Hosts that name no server fall through to the lobby. + for (String host : new String[]{"nope.mc.test", "alpha.evil.test", "mc.test"}) { + Fakes.FakePlayer p = player(host, true); + assertEq(host + ": still login first", "login", choose(p)); + assertEq(host + ": released to the lobby", "lobby", allowedTo(release(p))); + } + assertEq("no stray wake from other hosts", 1, api.count("POST " + SERVERS + "alpha/wake")); + + // A reconnect never inherits the host of an earlier connection. + Fakes.FakePlayer p5 = player("alpha.mc.test", true); + choose(p5); + p5.virtualHost = null; + choose(p5); + assertEq("reconnect: released to the lobby", "lobby", allowedTo(release(p5))); + assertEq("reconnect: nothing woken", 1, api.count("POST " + SERVERS + "alpha/wake")); + + // A target that vanished between connect and release. + Fakes.FakePlayer p6 = player("gone.mc.test", true); + choose(p6); + reg.refresh(servers(false)); + assertEq("vanished target: released to the lobby", "lobby", allowedTo(release(p6))); + assertEq("vanished target: told", true, p6.said("« gone » is no longer available")); + assertEq("vanished target: not woken", 0, api.count("POST " + SERVERS + "gone/wake")); + } + + private static void loginGate() { + Fakes.FakePlayer p = player(null, true); + + // Only a move out of login is gated. + ServerPreConnectEvent first = new ServerPreConnectEvent(p.player, beta); + assertEq("initial connect is not gated", null, router.onServerPreConnect(first)); + assertEq("... and keeps its target", "beta", allowedTo(first)); + ServerPreConnectEvent fromLobby = new ServerPreConnectEvent(p.player, beta, lobby); + assertEq("a move from the lobby is not gated", null, router.onServerPreConnect(fromLobby)); + assertEq("... and keeps its target", "beta", allowedTo(fromLobby)); + + // The gate may only release to the lobby, linked or not. + Fakes.FakePlayer zh = player(null, true); + zh.locale = Locale.SIMPLIFIED_CHINESE; + ServerPreConnectEvent elsewhere = new ServerPreConnectEvent(zh.player, beta, login); + assertEq("login to a user server: decided at once", null, router.onServerPreConnect(elsewhere)); + assertEq("login to a user server: denied", false, elsewhere.getResult().isAllowed()); + assertEq("login to a user server: told (zh)", true, zh.said("登录网关只能把玩家放行到大厅。")); + assertEq("login to a user server: logged", 1, log.count("WARN", "denied login-gate transfer for " + zh.id)); + assertEq("login to a user server: no link lookup", 0, api.count(LINK + zh.id)); + + // Unlinked: denied, and the gate's retries do not repeat the line. + Fakes.FakePlayer unlinked = player(null, false); + assertEq("unlinked: denied", false, release(unlinked).getResult().isAllowed()); + assertEq("unlinked: denied again", false, release(unlinked).getResult().isAllowed()); + assertEq("unlinked: told once", 1, count(unlinked.messages, "Finish signing in before leaving")); + assertEq("unlinked: both asked felis-api", 2, api.count(LINK + unlinked.id)); + + // felis-api down and no recent confirmation: fail closed. + Fakes.FakePlayer unknown = player(null, true); + api.linkDown = true; + assertEq("api down, never confirmed: denied", false, release(unknown).getResult().isAllowed()); + assertEq("api down, never confirmed: told", true, unknown.said("temporarily unavailable")); + assertEq("api down, never confirmed: logged", 1, + log.count("WARN", "could not verify login release for " + unknown.id)); + api.linkDown = false; + + // felis-api down after a confirmation minutes ago: admitted, and said so. + Fakes.FakePlayer known = player(null, true); + assertEq("confirmed: released", "lobby", allowedTo(release(known))); + api.linkDown = true; + assertEq("api down, recently confirmed: released", "lobby", allowedTo(release(known))); + assertEq("api down, recently confirmed: logged", 1, log.count("WARN", "admitting " + known.id)); + api.linkDown = false; + } + + private static void wakeRefusals() { + String[][] cases = { + {"403 forbidden", "You're not allowed to start « gamma »", null}, + {"409 maintenance_in_progress", "« gamma » is under maintenance", null}, + {"409 conflict", "Couldn't start « gamma » right now", "wake gamma failed (status=409)"}, + {"503 at_capacity", "The cluster is at capacity right now", null}, + {"503 unavailable", "Couldn't start « gamma » right now", "wake gamma failed (status=503)"}, + {"500 internal", "Couldn't start « gamma » right now", "wake gamma failed (status=500)"}, + }; + for (String[] c : cases) { + api.wakeError.put("gamma", c[0]); + int before = router.waitingCount(); + int warned = c[2] == null ? 0 : log.count("WARN", c[2]); + Fakes.FakePlayer p = player("gamma.mc.test", true); + choose(p); + assertEq(c[0] + ": released to the lobby", "lobby", allowedTo(release(p))); + assertEq(c[0] + ": told", true, p.said(c[1])); + assertEq(c[0] + ": not queued", before, router.waitingCount()); + assertEq(c[0] + ": no promise of a start", false, p.said("Starting « gamma »")); + if (c[2] != null) { + assertEq(c[0] + ": logged", warned + 1, log.count("WARN", c[2])); + } + } + // 429: a wake is already in flight, so join the wait. + api.wakeError.put("gamma", "429 cooldown"); + int before = router.waitingCount(); + Fakes.FakePlayer p = player("gamma.mc.test", true); + choose(p); + release(p); + assertEq("429: queued", before + 1, router.waitingCount()); + assertEq("429: told it is starting", true, p.said("Starting « gamma »")); + api.wakeError.remove("gamma"); + } + + private static void queue() { + // Queue now: one player on alpha (hostRouting), one on gamma (the 429). + Fakes.FakePlayer onAlpha = player("alpha.mc.test", true); + choose(onAlpha); + release(onAlpha); + assertEq("queue holds three", 3, router.waitingCount()); + List alphaWaiters = List.of(onAlpha); + + int alphaPolls = api.count("GET " + SERVERS + "alpha/status"); + int gammaPolls = api.count("GET " + SERVERS + "gamma/status"); + router.tick(); + assertEq("one status poll per server, however many wait on it", alphaPolls + 1, + api.count("GET " + SERVERS + "alpha/status")); + assertEq("gamma polled too", gammaPolls + 1, api.count("GET " + SERVERS + "gamma/status")); + assertEq("nothing ready: nobody moved", 0, onAlpha.connects.size()); + assertEq("nothing ready: all still queued", 3, router.waitingCount()); + + api.ready.put("alpha", true); + router.tick(); + assertEq("alpha ready: its waiters moved", List.of("alpha"), List.copyOf(onAlpha.connects)); + assertEq("alpha ready: told", true, onAlpha.said("« alpha » is ready — moving you in")); + assertEq("alpha ready: one poll again", alphaPolls + 2, api.count("GET " + SERVERS + "alpha/status")); + assertEq("alpha ready: only gamma's waiter left", 1, router.waitingCount()); + for (Fakes.FakePlayer w : alphaWaiters) { + assertEq("moved player checked linked again", true, api.count(LINK + w.id) >= 2); + } + + // A waiter who left the proxy is dropped without polling for them. + net.players.clear(); + int polls = api.count("GET " + SERVERS + "gamma/status"); + router.tick(); + assertEq("departed waiter dropped", 0, router.waitingCount()); + assertEq("departed waiter: no poll", polls, api.count("GET " + SERVERS + "gamma/status")); + + // Unlinked by the time the server is ready: dropped and told, not moved. + Fakes.FakePlayer lapsed = player("delta.mc.test", true); + choose(lapsed); + release(lapsed); + assertEq("lapsed: queued", 1, router.waitingCount()); + api.linked.remove(lapsed.id); + api.ready.put("delta", true); + router.tick(); + assertEq("lapsed: dropped", 0, router.waitingCount()); + assertEq("lapsed: not moved", 0, lapsed.connects.size()); + assertEq("lapsed: told", true, lapsed.said("no longer linked")); + + // Ready but with no backend registered yet: wait for the next refresh. + Fakes.FakePlayer early = player("fresh.mc.test", true); + choose(early); + release(early); + api.ready.put("fresh", true); + router.tick(); + assertEq("unregistered: still queued", 1, router.waitingCount()); + assertEq("unregistered: not moved", 0, early.connects.size()); + List list = servers(false); + list.set(list.size() - 1, view("fresh", true, "10.43.0.10:25565")); + reg.refresh(list); + router.tick(); + assertEq("registered: moved", List.of("fresh"), List.copyOf(early.connects)); + assertEq("registered: queue empty", 0, router.waitingCount()); + } + + private static void menuAndCommands() { + List notified = Collections.synchronizedList(new ArrayList<>()); + router.setMenuTransferListener((player, server) -> notified.add(player.getUsername() + "@" + server)); + + Fakes.FakePlayer menu = player(null, true); + menu.current = lobby; + Fakes.FakePlayer command = player(null, true); + command.current = lobby; + router.enqueueFromMenu(menu.player, "epsilon"); + router.enqueueFromCommand(command.player, "epsilon"); + Fakes.await("both queued", () -> router.waitingCount() == 2); + assertEq("each entry woke epsilon", 2, api.count("POST " + SERVERS + "epsilon/wake")); + api.ready.put("epsilon", true); + router.tick(); + assertEq("both moved (menu)", List.of("epsilon"), List.copyOf(menu.connects)); + assertEq("both moved (command)", List.of("epsilon"), List.copyOf(command.connects)); + assertEq("only the menu entry tells the lobby", List.of(menu.name + "@epsilon"), List.copyOf(notified)); + + // Asking for the server you stand on is answered at once, case-insensitively. + Fakes.FakePlayer there = player(null, true); + there.current = beta; + router.enqueueFromCommand(there.player, "BETA"); + assertEq("already there: told", true, there.said("You're already on « BETA »")); + assertEq("already there: nothing asked", 0, api.count(LINK + there.id)); + + // An invite to a running server joins it; to a stopped one it goes through the + // policy-gated wake like any other entry. + Fakes.FakePlayer invited = player(null, true); + invited.current = lobby; + router.enqueueFromInvite(invited.player, "beta"); + Fakes.await("invitee moved", () -> invited.connects.size() == 1); + assertEq("invite to a running server: joined", List.of("beta"), List.copyOf(invited.connects)); + assertEq("invite to a running server: no wake", 0, api.count("POST " + SERVERS + "beta/wake")); + + Fakes.FakePlayer invitedStopped = player(null, true); + invitedStopped.current = lobby; + router.enqueueFromInvite(invitedStopped.player, "zeta"); + Fakes.await("zeta woken", () -> api.count("POST " + SERVERS + "zeta/wake") == 1); + Fakes.await("invitee queued", () -> router.waitingCount() == 1); + assertEq("invite to a stopped server: not moved yet", 0, invitedStopped.connects.size()); + + Fakes.FakePlayer unlinked = player(null, false); + unlinked.current = lobby; + router.enqueueFromCommand(unlinked.player, "zeta"); + Fakes.await("unlinked told", () -> unlinked.said("Finish signing in before joining a server")); + assertEq("unlinked command: no wake", 1, api.count("POST " + SERVERS + "zeta/wake")); + assertEq("unlinked command: not queued", 1, router.waitingCount()); + } + + private static void joins() throws InterruptedException { + Fakes.FakePlayer p = player(null, true); + RegisteredServer alpha = net.proxy.getServer("alpha").orElseThrow(); + RegisteredServer stat = net.add("static", "10.0.0.50", 25565); + router.onServerConnected(new ServerConnectedEvent(p.player, lobby, null)); + router.onServerConnected(new ServerConnectedEvent(p.player, login, null)); + router.onServerConnected(new ServerConnectedEvent(p.player, stat, lobby)); + router.onServerConnected(new ServerConnectedEvent(p.player, alpha, lobby)); + Fakes.await("alpha join reported", () -> api.count("POST " + SERVERS + "alpha/join-event") == 1); + Thread.sleep(200); // let anything wrongly submitted for the others land too + assertEq("lobby join not reported", 0, api.count("POST " + SERVERS + "lobby/join-event")); + assertEq("login join not reported", 0, api.count("POST " + SERVERS + "login/join-event")); + assertEq("unmanaged join not reported", 0, api.count("POST " + SERVERS + "static/join-event")); + assertEq("no failure counted", 0L, plugin.stats().total(ProxyStats.Event.JOIN_EVENT_FAILED)); + + api.joinStatus = 500; + router.onServerConnected(new ServerConnectedEvent(p.player, beta, lobby)); + Fakes.await("failed join counted", () -> plugin.stats().total(ProxyStats.Event.JOIN_EVENT_FAILED) == 1); + Fakes.await("failed join logged", + () -> log.count("WARN", "join-event for " + p.id + " on beta failed (status=500)") == 1); + api.joinStatus = 204; + + // A transfer the backend refuses is counted, logged and told. + Fakes.FakePlayer refused = player(null, true); + refused.current = lobby; + refused.connectSucceeds = false; + router.enqueueFromInvite(refused.player, "beta"); + Fakes.await("failed transfer counted", () -> plugin.stats().total(ProxyStats.Event.TRANSFER_FAILED) == 1); + assertEq("failed transfer: told", true, refused.said("Couldn't connect you to « beta »")); + assertEq("failed transfer: logged", 1, log.count("WARN", "transfer of " + refused.id + " to beta failed")); + } + + private static void disconnectAndRelease() { + // Leaving the proxy drops the queue entry. + Fakes.FakePlayer leaver = player(null, true); + leaver.current = lobby; + int before = router.waitingCount(); + router.enqueueFromCommand(leaver.player, "zeta"); + Fakes.await("leaver queued", () -> router.waitingCount() == before + 1); + router.onDisconnect(new DisconnectEvent(leaver.player, DisconnectEvent.LoginStatus.SUCCESSFUL_LOGIN)); + assertEq("disconnect drops the entry", before, router.waitingCount()); + + // The login gate's release asks Velocity for the lobby. + Fakes.FakePlayer released = player(null, true); + router.releaseFromLogin(released.player); + assertEq("release connects to the lobby", List.of("lobby"), List.copyOf(released.connects)); + + net.remove("lobby"); + Fakes.FakePlayer early = player(null, true); + router.releaseFromLogin(early.player); + assertEq("no lobby yet: nothing sent", 0, early.connects.size()); + assertEq("no lobby yet: logged", 1, log.count("WARN", "login release for " + early.id + " but the lobby")); + + // With no login gate registered a host-routed player is turned away. + net.remove("login"); + Fakes.FakePlayer stranded = player("beta.mc.test", true); + PlayerChooseInitialServerEvent e = new PlayerChooseInitialServerEvent(stranded.player, null); + router.onChooseInitialServer(e); + assertEq("no login: no initial server", false, e.getInitialServer().isPresent()); + assertEq("no login: disconnected with a reason", "The Felis login gate is unavailable. Please reconnect shortly.", + stranded.disconnectedWith); + } + + // ---- helpers ---- + + private static Fakes.FakePlayer player(String host, boolean linked) { + int n = nextPlayer++; + Fakes.FakePlayer p = new Fakes.FakePlayer(net, "player" + n, n); + p.virtualHost = host; + if (linked) { + api.linked.add(p.id); + } + return p; + } + + // choose runs PlayerChooseInitialServerEvent the way Velocity does (the first + // try server, login, preset) and returns the server the player will land on. + private static String choose(Fakes.FakePlayer p) { + PlayerChooseInitialServerEvent e = new PlayerChooseInitialServerEvent(p.player, login); + router.onChooseInitialServer(e); + return e.getInitialServer().map(s -> s.getServerInfo().getName()).orElse(null); + } + + // release is the login gate asking to move the player to the lobby, with the + // router's async part run to completion the way Velocity's event manager would. + private static ServerPreConnectEvent release(Fakes.FakePlayer p) { + p.current = login; + ServerPreConnectEvent e = new ServerPreConnectEvent(p.player, lobby, login); + EventTask task = router.onServerPreConnect(e); + if (task != null) { + task.execute(new Continuation() { + @Override + public void resume() { + } + + @Override + public void resumeWithException(Throwable t) { + throw new AssertionError("pre-connect task failed", t); + } + }); + } + return e; + } + + private static String allowedTo(ServerPreConnectEvent e) { + if (!e.getResult().isAllowed()) { + return "denied"; + } + return e.getResult().getServer().map(s -> s.getServerInfo().getName()).orElse(null); + } + + private static int count(List lines, String part) { + int n = 0; + synchronized (lines) { + for (String l : lines) { + if (l.contains(part)) { + n++; + } + } + } + return n; + } + + private static ServerView view(String name, boolean ready, String addr) { + return new ServerView(name, name, ready ? "Running" : "Stopped", ready, "ownerOnly", + "Running", "ClusterIP", addr, 0, 20); + } + + private static void assertEq(String what, Object want, Object got) { + if (want == null ? got != null : !want.equals(got)) { + throw new AssertionError(what + ": got " + got + ", want " + want); + } + checks++; + } +}