fix(velocity): 停服的 fallback 名不再注册成后端,唤醒就绪时按轮询到的地址立即注册再传送

This commit is contained in:
Lemon-miaow committed 2026-09-26 21:18:08 +08:00
1 parent 4b386c5c2e
commit 3843082153
6 files changed
+158 -41

No files matched your search

@@ -36,9 +36,18 @@ import java.util.concurrent.ConcurrentHashMap;
* first successful refresh replaces that placeholder with the live ClusterIP. A * first successful refresh replaces that placeholder with the live ClusterIP. A
* proxy that starts while felis-api is down refreshes once from the last list the * proxy that starts while felis-api is down refreshes once from the last list the
* API answered with instead ({@link ServerListSource}). * API answered with instead ({@link ServerListSource}).
*
* <p>Only a {@code direct} endpoint carries a backend address. A server that is not
* up reports {@code fallback}, with the NAME of its fallback server ("login") in the
* address field: dialled as a hostname that resolves nowhere, so registering it
* turned every wake-and-transfer into "couldn't connect you". A fallback report leaves
* the last direct registration in place — the ClusterIP belongs to the Service and
* outlives the pod — and {@link #observe(ServerView)} lets the waiting queue apply the
* address a fresh status poll reports before it transfers anyone.
*/ */
final class ServerRegistry { final class ServerRegistry {
private static final int DEFAULT_PORT = 25565; private static final int DEFAULT_PORT = 25565;
private static final String ENDPOINT_DIRECT = "direct";
private final ProxyServer proxy; private final ProxyServer proxy;
private final Logger log; private final Logger log;
@@ -84,10 +93,23 @@ final class ServerRegistry {
subdomainToName.keySet().retainAll(subdomains.keySet()); subdomainToName.keySet().retainAll(subdomains.keySet());
} }
/**
* observe applies one server's freshly polled view between refreshes: the view
* routing reads, and its registration when the poll reports a direct endpoint. A
* name the last refresh did not list is left to the next one.
*/
void observe(ServerView v) {
String name = v.name();
if (name == null || byName.replace(name, v) == null) {
return;
}
ensureRegistered(v);
}
private void ensureRegistered(ServerView v) { private void ensureRegistered(ServerView v) {
String addr = v.endpointAddress(); String addr = v.endpointAddress();
if (addr == null || addr.isEmpty()) { if (!ENDPOINT_DIRECT.equalsIgnoreCase(v.endpointMode()) || addr == null || addr.isEmpty()) {
return; // no backend address yet (server never started) → nothing to register return; // not up (the address is a fallback server's name) or never started
} }
InetSocketAddress target = parseAddress(addr); InetSocketAddress target = parseAddress(addr);
Optional<RegisteredServer> existing = proxy.getServer(v.name()); Optional<RegisteredServer> existing = proxy.getServer(v.name());
@@ -401,7 +401,14 @@ public final class WaitingRouter {
Boolean ready = readyCache.get(w.serverName); Boolean ready = readyCache.get(w.serverName);
if (ready == null) { if (ready == null) {
try { try {
ready = api.serverStatus(w.serverName).ready(); ServerView status = api.serverStatus(w.serverName);
ready = status.ready();
// The registry refreshes every 15 s; a server that just came up may
// still be registered at its old address, or not at all. Register
// what this poll reports before transferring anyone to it.
if (ready) {
registry.observe(status);
}
} catch (LinkException ex) { } catch (LinkException ex) {
ready = Boolean.FALSE; // transient → keep waiting until the deadline ready = Boolean.FALSE; // transient → keep waiting until the deadline
} }
@@ -293,9 +293,11 @@ public final class ControlChannelTest {
return f.server() + " " + f.phase() + " " + f.ready(); return f.server() + " " + f.phase() + " " + f.ready();
} }
// Shaped like the operator's status: up is direct at addr; down is the fallback,
// whose name ("login") sits in the address field.
private static ServerView view(String name, boolean ready, String addr) { private static ServerView view(String name, boolean ready, String addr) {
return new ServerView(name, name, ready ? "Running" : "Stopped", ready, "ownerOnly", return new ServerView(name, name, ready ? "Running" : "Stopped", ready, "ownerOnly",
"Running", "ClusterIP", addr, 0, 20); "Running", ready ? "direct" : "fallback", ready ? addr : "login", 0, 20);
} }
private static void assertEq(String what, Object want, Object got) { private static void assertEq(String what, Object want, Object got) {
@@ -413,11 +413,17 @@ final class Fakes {
* status from {@link #ready}, wake answering 202 unless {@link #wakeError} holds an * status from {@link #ready}, wake answering 202 unless {@link #wakeError} holds an
* answer for the server, and join-events recorded. {@link #linkDown} makes the * answer for the server, and join-events recorded. {@link #linkDown} makes the
* link-status route answer 500. * link-status route answer 500.
*
* <p>Status replies carry the endpoint the way the operator writes it: a server that
* is up reports {@code direct} and its {@link #address}; one that is not reports
* {@code fallback} with the fallback server's NAME, "login", in the address field.
*/ */
static final class Api implements AutoCloseable { static final class Api implements AutoCloseable {
final Set<UUID> linked = ConcurrentHashMap.newKeySet(); final Set<UUID> linked = ConcurrentHashMap.newKeySet();
volatile boolean linkDown; volatile boolean linkDown;
final Map<String, Boolean> ready = new ConcurrentHashMap<>(); final Map<String, Boolean> ready = new ConcurrentHashMap<>();
/** address is the direct endpoint a server reports while it is ready. */
final Map<String, String> address = new ConcurrentHashMap<>();
/** wakeError maps a server to "status code" (e.g. "403 forbidden"). */ /** wakeError maps a server to "status code" (e.g. "403 forbidden"). */
final Map<String, String> wakeError = new ConcurrentHashMap<>(); final Map<String, String> wakeError = new ConcurrentHashMap<>();
volatile int joinStatus = 204; volatile int joinStatus = 204;
@@ -465,11 +471,13 @@ final class Fakes {
String action = parts.length > 1 ? parts[1] : ""; String action = parts.length > 1 ? parts[1] : "";
switch (method + " " + action) { switch (method + " " + action) {
case "GET status": case "GET status":
reply(ex, 200, view(name, ready.getOrDefault(name, false))); reply(ex, 200, status(name, ready.getOrDefault(name, false)));
return; return;
case "POST wake": case "POST wake":
if (!error(ex, wakeError.get(name))) { if (!error(ex, wakeError.get(name))) {
reply(ex, 202, view(name, false)); // The real wake reply is this subset of the view: no endpoint.
reply(ex, 202, "{\"name\":\"" + name + "\",\"desiredState\":\"Running\","
+ "\"phase\":\"Stopped\",\"ready\":false}");
} }
return; return;
case "GET menu": case "GET menu":
@@ -503,8 +511,13 @@ final class Fakes {
return true; return true;
} }
private static String view(String name, boolean ready) { private String status(String name, boolean up) {
return "{\"name\":\"" + name + "\",\"subdomain\":\"" + name + "\",\"ready\":" + ready + "}"; String addr = address.get(name);
String endpoint = up
? "\"endpointMode\":\"direct\"" + (addr == null ? "" : ",\"endpointAddress\":\"" + addr + "\"")
: "\"endpointMode\":\"fallback\",\"endpointAddress\":\"login\"";
return "{\"name\":\"" + name + "\",\"subdomain\":\"" + name + "\",\"phase\":\""
+ (up ? "Running" : "Stopped") + "\",\"ready\":" + up + "," + endpoint + "}";
} }
private static void reply(HttpExchange ex, int status, String body) throws IOException { private static void reply(HttpExchange ex, int status, String body) throws IOException {
@@ -8,10 +8,15 @@ import java.util.Optional;
/** /**
* ServerRegistryTest drives the real ServerRegistry against a fake ProxyServer: a * ServerRegistryTest drives the real ServerRegistry against a fake ProxyServer: a
* refresh registers every server that has a backend address, leaves one registered at * refresh registers every server that is up at a direct endpoint, leaves one
* the same address alone, moves one whose address changed, drops one that vanished * registered at the same address alone, moves one whose address changed, keeps the
* from a successful fetch, and host routing resolves only {@code <subdomain>.<root>}. * last registration of one that went down, drops one that vanished from a successful
* Framework free: a failed assertion throws. * fetch, a single polled view re-registers between refreshes, and host routing
* resolves only {@code <subdomain>.<root>}. Framework free: a failed assertion throws.
*
* <p>The views are shaped the way the operator writes status: up is {@code direct}
* plus the Service address; down is {@code fallback} with the fallback server's NAME
* in the address field; a server never reconciled has neither.
* *
* <p>Run: {@code ./gradlew routingTest} in plugins/velocity (it needs the velocity-api * <p>Run: {@code ./gradlew routingTest} in plugins/velocity (it needs the velocity-api
* classes the plugin compiles against). * classes the plugin compiles against).
@@ -28,32 +33,65 @@ public final class ServerRegistryTest {
ServerRegistry reg = new ServerRegistry(net.proxy, log.logger, "MC.Example.test"); ServerRegistry reg = new ServerRegistry(net.proxy, log.logger, "MC.Example.test");
reg.refresh(List.of( reg.refresh(List.of(
view("login", "login", true, "10.43.0.1:25565"), up("login", "login", "10.43.0.1:25565"),
view("lobby", "lobby", true, "10.43.0.2:25565"), up("lobby", "lobby", "10.43.0.2:25565"),
view("alpha", "alpha", false, "10.43.0.3:25566"), down("alpha", "alpha"),
view("fresh", "fresh", false, null), unstarted("fresh", "fresh"),
view("", "blank", true, "10.43.0.9:25565"))); up("", "blank", "10.43.0.9:25565")));
assertEq("first refresh registrations", 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.of("-login", "+login 10.43.0.1:25565", "+lobby 10.43.0.2:25565"),
List.copyOf(net.registrations)); List.copyOf(net.registrations));
assertEq("login placeholder replaced", "10.43.0.1:25565", net.address("login")); 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("a stopped server is known", true, reg.isManaged("alpha"));
assertEq("... but its fallback's name is not dialled", null, net.address("alpha"));
assertEq("a server with no endpoint is known", true, reg.isManaged("fresh"));
assertEq("... but not registered", null, net.address("fresh")); assertEq("... but not registered", null, net.address("fresh"));
assertEq("a nameless entry is skipped", false, reg.isManaged("")); assertEq("a nameless entry is skipped", false, reg.isManaged(""));
assertEq("managed count", 4, reg.all().size()); assertEq("managed count", 4, reg.all().size());
assertEq("registration logged", 1, log.count("INFO", "registered backend alpha -> 10.43.0.3:25566")); assertEq("registration logged", 1, log.count("INFO", "registered backend lobby -> 10.43.0.2:25565"));
net.registrations.clear(); net.registrations.clear();
reg.refresh(List.of( reg.refresh(List.of(
view("login", "login", true, "10.43.0.1:25565"), up("login", "login", "10.43.0.1:25565"),
view("lobby", "lobby", true, "10.43.0.2:25565"), up("lobby", "lobby", "10.43.0.2:25565"),
view("alpha", "alpha", true, "10.43.0.7:25566"), up("alpha", "alpha", "10.43.0.3:25566"),
view("fresh", "fresh", true, "10.43.0.8:25565"))); up("fresh", "fresh", "10.43.0.8:25565")));
assertEq("second refresh: only the moved and the new", assertEq("second refresh: only the servers that came up",
List.of("-alpha", "+alpha 10.43.0.7:25566", "+fresh 10.43.0.8:25565"), List.of("+alpha 10.43.0.3:25566", "+fresh 10.43.0.8:25565"),
List.copyOf(net.registrations)); List.copyOf(net.registrations));
assertEq("view follows the fetch", true, reg.view("alpha").ready()); assertEq("view follows the fetch", true, reg.view("alpha").ready());
// alpha goes down: its Service keeps the ClusterIP, so the registration stays.
net.registrations.clear();
reg.refresh(List.of(
up("login", "login", "10.43.0.1:25565"),
up("lobby", "lobby", "10.43.0.2:25565"),
down("alpha", "alpha"),
up("fresh", "fresh", "10.43.0.8:25565")));
assertEq("down: registrations untouched", List.of(), List.copyOf(net.registrations));
assertEq("down: last direct address kept", "10.43.0.3:25566", net.address("alpha"));
assertEq("down: the view says so", false, reg.view("alpha").ready());
// A polled view between refreshes: up at a new address re-registers at once.
reg.observe(up("alpha", "alpha", "10.43.0.7:25566"));
assertEq("observed: moved", List.of("-alpha", "+alpha 10.43.0.7:25566"), List.copyOf(net.registrations));
assertEq("observed: the view follows", true, reg.view("alpha").ready());
net.registrations.clear();
reg.observe(down("alpha", "alpha"));
assertEq("observed down: registration kept", "10.43.0.7:25566", net.address("alpha"));
reg.observe(up("stranger", "stranger", "10.43.0.99:25565"));
assertEq("observed a name no refresh listed: left to the next refresh", false, reg.isManaged("stranger"));
assertEq("... and not registered", null, net.address("stranger"));
assertEq("observing registered nothing", List.of(), List.copyOf(net.registrations));
net.registrations.clear();
reg.refresh(List.of(
up("login", "login", "10.43.0.1:25565"),
up("lobby", "lobby", "10.43.0.2:25565"),
up("alpha", "alpha", "10.43.0.7:25566"),
up("fresh", "fresh", "10.43.0.8:25565")));
assertEq("an unchanged refresh changes nothing", List.of(), List.copyOf(net.registrations));
assertEq("host routing", "alpha", name(reg.resolveByHost("alpha.mc.example.test"))); 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("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("the root itself routes nowhere", null, name(reg.resolveByHost("mc.example.test")));
@@ -64,9 +102,9 @@ public final class ServerRegistryTest {
net.registrations.clear(); net.registrations.clear();
reg.refresh(List.of( reg.refresh(List.of(
view("login", "login", true, "10.43.0.1:25565"), up("login", "login", "10.43.0.1:25565"),
view("lobby", "lobby", true, "10.43.0.2:25565"), up("lobby", "lobby", "10.43.0.2:25565"),
view("fresh", "fresh", true, "10.43.0.8:25565"))); up("fresh", "fresh", "10.43.0.8:25565")));
assertEq("a vanished server is dropped", List.of("-alpha"), List.copyOf(net.registrations)); assertEq("a vanished server is dropped", List.of("-alpha"), List.copyOf(net.registrations));
assertEq("... from the views", false, reg.isManaged("alpha")); assertEq("... from the views", false, reg.isManaged("alpha"));
assertEq("... and from host routing", null, name(reg.resolveByHost("alpha.mc.example.test"))); assertEq("... and from host routing", null, name(reg.resolveByHost("alpha.mc.example.test")));
@@ -74,18 +112,18 @@ public final class ServerRegistryTest {
// A server whose subdomain changed answers on the new one only. // A server whose subdomain changed answers on the new one only.
reg.refresh(List.of( reg.refresh(List.of(
view("login", "login", true, "10.43.0.1:25565"), up("login", "login", "10.43.0.1:25565"),
view("lobby", "lobby", true, "10.43.0.2:25565"), up("lobby", "lobby", "10.43.0.2:25565"),
view("fresh", "renamed", true, "10.43.0.8:25565"))); up("fresh", "renamed", "10.43.0.8:25565")));
assertEq("new subdomain", "fresh", name(reg.resolveByHost("renamed.mc.example.test"))); assertEq("new subdomain", "fresh", name(reg.resolveByHost("renamed.mc.example.test")));
assertEq("old subdomain let go", null, name(reg.resolveByHost("fresh.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. // A subdomain handed to another server routes to the new holder.
reg.refresh(List.of( reg.refresh(List.of(
view("login", "login", true, "10.43.0.1:25565"), up("login", "login", "10.43.0.1:25565"),
view("lobby", "lobby", true, "10.43.0.2:25565"), up("lobby", "lobby", "10.43.0.2:25565"),
view("other", "renamed", true, "10.43.0.5:25565"), up("other", "renamed", "10.43.0.5:25565"),
view("fresh", "fresh", true, "10.43.0.8:25565"))); up("fresh", "fresh", "10.43.0.8:25565")));
assertEq("subdomain moved", "other", name(reg.resolveByHost("renamed.mc.example.test"))); assertEq("subdomain moved", "other", name(reg.resolveByHost("renamed.mc.example.test")));
assertEq("subdomain taken back", "fresh", name(reg.resolveByHost("fresh.mc.example.test"))); assertEq("subdomain taken back", "fresh", name(reg.resolveByHost("fresh.mc.example.test")));
@@ -104,9 +142,16 @@ public final class ServerRegistryTest {
System.out.println("ServerRegistryTest OK (" + checks + " checks)"); System.out.println("ServerRegistryTest OK (" + checks + " checks)");
} }
private static ServerView view(String name, String sub, boolean ready, String addr) { private static ServerView up(String name, String sub, String addr) {
return new ServerView(name, sub, ready ? "Running" : "Stopped", ready, "ownerOnly", return new ServerView(name, sub, "Running", true, "ownerOnly", "Running", "direct", addr, 0, 20);
"Running", "ClusterIP", addr, 0, 20); }
private static ServerView down(String name, String sub) {
return new ServerView(name, sub, "Stopped", false, "ownerOnly", "Stopped", "fallback", "login", 0, 0);
}
private static ServerView unstarted(String name, String sub) {
return new ServerView(name, sub, "", false, "ownerOnly", "Stopped", null, null, 0, 0);
} }
private static String name(Optional<ServerView> v) { private static String name(Optional<ServerView> v) {
@@ -62,6 +62,9 @@ public final class WaitingRouterTest {
login = net.proxy.getServer("login").orElseThrow(); login = net.proxy.getServer("login").orElseThrow();
lobby = net.proxy.getServer("lobby").orElseThrow(); lobby = net.proxy.getServer("lobby").orElseThrow();
beta = net.proxy.getServer("beta").orElseThrow(); beta = net.proxy.getServer("beta").orElseThrow();
// A stopped server's address field names its fallback; dialled, "login"
// resolves nowhere.
assertEq("a stopped server is not registered at its fallback's name", null, net.address("alpha"));
hostRouting(); hostRouting();
loginGate(); loginGate();
@@ -259,6 +262,8 @@ public final class WaitingRouterTest {
api.ready.put("alpha", true); api.ready.put("alpha", true);
router.tick(); router.tick();
assertEq("alpha ready: registered from the poll, not the next refresh", "10.43.0.3:25565",
net.address("alpha"));
assertEq("alpha ready: its waiters moved", List.of("alpha"), List.copyOf(onAlpha.connects)); 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: 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: one poll again", alphaPolls + 2, api.count("GET " + SERVERS + "alpha/status"));
@@ -286,7 +291,24 @@ public final class WaitingRouterTest {
assertEq("lapsed: not moved", 0, lapsed.connects.size()); assertEq("lapsed: not moved", 0, lapsed.connects.size());
assertEq("lapsed: told", true, lapsed.said("no longer linked")); assertEq("lapsed: told", true, lapsed.said("no longer linked"));
// Ready but with no backend registered yet: wait for the next refresh. // alpha stops: the fallback report leaves its last registration alone. It comes
// back behind a new address, and the waiter goes there, not to the old one.
api.ready.put("alpha", false);
reg.refresh(servers(false));
assertEq("stopped: last direct registration kept", "10.43.0.3:25565", net.address("alpha"));
Fakes.FakePlayer back = player("alpha.mc.test", true);
choose(back);
release(back);
assertEq("stopped again: queued", 1, router.waitingCount());
api.address.put("alpha", "10.43.0.13:25565");
api.ready.put("alpha", true);
router.tick();
assertEq("new address: re-registered from the poll", "10.43.0.13:25565", net.address("alpha"));
assertEq("new address: moved", List.of("alpha"), List.copyOf(back.connects));
assertEq("new address: queue empty", 0, router.waitingCount());
// Ready, but neither the registry nor the poll has an address for it yet: wait
// for the refresh that brings one.
Fakes.FakePlayer early = player("fresh.mc.test", true); Fakes.FakePlayer early = player("fresh.mc.test", true);
choose(early); choose(early);
release(early); release(early);
@@ -474,9 +496,15 @@ public final class WaitingRouterTest {
return n; return n;
} }
// view is a server as GET /servers lists it: up at its direct address, or down on
// the fallback, where the operator writes the fallback server's name ("login")
// into the address. addr is also what the stub API reports once the server is up.
private static ServerView view(String name, boolean ready, String addr) { private static ServerView view(String name, boolean ready, String addr) {
if (addr != null) {
api.address.put(name, addr);
}
return new ServerView(name, name, ready ? "Running" : "Stopped", ready, "ownerOnly", return new ServerView(name, name, ready ? "Running" : "Stopped", ready, "ownerOnly",
"Running", "ClusterIP", addr, 0, 20); "Running", ready ? "direct" : "fallback", ready ? addr : "login", 0, 20);
} }
private static void assertEq(String what, Object want, Object got) { private static void assertEq(String what, Object want, Object got) {