diff --git a/deploy/limbo/README.md b/deploy/limbo/README.md index d3829e8..2b549bf 100644 --- a/deploy/limbo/README.md +++ b/deploy/limbo/README.md @@ -58,6 +58,8 @@ Configuration (deployment inputs, never compiled in; env wins over a | `FELIS_LOBBY_SERVER` | Velocity server name to transfer to | `lobby` | | `FELIS_LOGIN_TIMEOUT_SECONDS` | login window (clamped 30–3600) | `600` | | `FELIS_HEALTH_PORT` | readiness port | `8080` | +| `FELIS_API_CONNECT_TIMEOUT_SECONDS` | felis-api connect timeout (1–120) | `10` | +| `FELIS_API_REQUEST_TIMEOUT_SECONDS` | felis-api call timeout (1–120) | `10` | If the API config **or** the root domain is absent the login flow stays **OFF** and the plugin runs readiness-only (the same "load un-crippled" fail-safe the other diff --git a/plugins/README.md b/plugins/README.md index a16144b..828d3ea 100644 --- a/plugins/README.md +++ b/plugins/README.md @@ -120,7 +120,15 @@ so there is nothing to shade and each jar is self-contained. - `LinkConfigLoader` — reads `FELIS_API_BASE_URL` / `FELIS_SERVICE_TOKEN` (env wins) or a `felis-link.properties` file written as a commented template on first run. **The API URL and service token are deployment inputs and are never - compiled in.** + compiled in.** Two optional keys bound each call, in whole seconds from 1 to 120 + (default 10): `connect-timeout-seconds` / `FELIS_API_CONNECT_TIMEOUT_SECONDS` + for opening the connection, `request-timeout-seconds` / + `FELIS_API_REQUEST_TIMEOUT_SECONDS` for the whole call. Anything else fails the + load with the key named. +- `FelisApiClient` — the routing client. A GET that failed on a dropped + connection or a 502/503/504 is tried once more after a 100–400 ms jittered + pause; a GET that timed out, and every POST (wake, claim, join-event, + approvals), is never repeated. Threading: the command runs on the server thread; the HTTP call is dispatched to a daemon single-thread executor and the reply is hopped back onto the server @@ -175,6 +183,15 @@ Velocity-only config keys (read from the same `felis-link.properties` / env as | `login-server` | `FELIS_LOGIN_SERVER` | The system auth gate every fresh connection must pass. Defaults to `login`. | | `lobby-server` | `FELIS_LOBBY_SERVER` | The distinct post-auth holding server used while a backend wakes. Defaults to `lobby`; it must not equal `login-server`. | +Load bounds on the proxy: felis-api calls run on a pool of 8 threads with 64 +waiting slots. Past that a call is refused at once — the player reads "busy", a +lobby frame gets a `busy` error, and a dropped join-event is logged at warn — +rather than piling up threads while felis-api is slow. The registration refresh +(15 s) and the waiting-queue poll (2 s) skip a run that falls due while the last +one is still going. Acting commands (`/link`, `/felis claim`, migrate, op +approve) share a per-player budget of 5 then one per 5 s; felis:control frames +from the lobby are metered per player and menu status answers are cached. + ## Lobby menu (§12) The `paper/` module is the lobby's player-facing face for §27 scenario 10 diff --git a/plugins/paper/src/main/java/best/lolicon/felis/paper/FelisPaperPlugin.java b/plugins/paper/src/main/java/best/lolicon/felis/paper/FelisPaperPlugin.java index ae6eab1..b687f61 100644 --- a/plugins/paper/src/main/java/best/lolicon/felis/paper/FelisPaperPlugin.java +++ b/plugins/paper/src/main/java/best/lolicon/felis/paper/FelisPaperPlugin.java @@ -384,6 +384,9 @@ public final class FelisPaperPlugin extends JavaPlugin implements Listener, Plug case "invalid_server_name": return zh ? "这台服务器已不存在。" : "That server no longer exists."; + case "busy": + return zh ? "Felis 现在很忙,请过一会儿再试。" + : "Felis is busy right now — try again in a moment."; case "transport_error": case "interrupted": return zh ? "Felis 暂时不可用,请稍后再试。" diff --git a/plugins/shared/src/main/java/best/lolicon/felis/link/FelisApiClient.java b/plugins/shared/src/main/java/best/lolicon/felis/link/FelisApiClient.java index 37e3082..5cbfd23 100644 --- a/plugins/shared/src/main/java/best/lolicon/felis/link/FelisApiClient.java +++ b/plugins/shared/src/main/java/best/lolicon/felis/link/FelisApiClient.java @@ -6,11 +6,14 @@ import java.net.URLEncoder; import java.net.http.HttpClient; import java.net.http.HttpRequest; import java.net.http.HttpResponse; +import java.net.http.HttpTimeoutException; import java.nio.charset.StandardCharsets; import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.Set; import java.util.UUID; +import java.util.concurrent.ThreadLocalRandom; import java.util.regex.Pattern; /** @@ -43,6 +46,11 @@ import java.util.regex.Pattern; public final class FelisApiClient { // Mirrors internal/naming.serverNameRE plus its no-leading/trailing-dash rule. private static final Pattern SERVER_NAME = Pattern.compile("^[a-z0-9][a-z0-9-]{1,30}[a-z0-9]$"); + // Answers a second GET can outlive: felis-api restarting behind its Service, or a + // gateway with no ready endpoint. + private static final Set RETRY_STATUSES = Set.of(502, 503, 504); + private static final int RETRY_PAUSE_MIN_MILLIS = 100; + private static final int RETRY_PAUSE_JITTER_MILLIS = 300; private final LinkConfig config; private final HttpClient http; @@ -50,7 +58,7 @@ public final class FelisApiClient { public FelisApiClient(LinkConfig config) { this.config = Objects.requireNonNull(config, "config"); this.http = HttpClient.newBuilder() - .connectTimeout(config.timeout()) + .connectTimeout(config.connectTimeout()) .build(); } @@ -261,8 +269,34 @@ public final class FelisApiClient { return URLEncoder.encode(value, StandardCharsets.UTF_8).replace("+", "%20"); } + // getObject tries a GET a second time, after a jittered 100-400 ms pause, when the + // first attempt failed in a way a retry can fix: the connection was refused or + // reset (felis-api restarting, its pod moving) or the answer was 502/503/504. GETs + // only read, so repeating one is safe; POSTs (wake, claim, join-event, approvals) + // are never repeated, because the first attempt may have landed. A GET that timed + // out is not repeated either: felis-api is then slow rather than gone, and a second + // full wait would hold the caller twice as long while adding load to an API that + // is already behind. The jitter keeps a proxy's queued callers from retrying in + // one burst. private Map getObject(String path, int expect) throws LinkException { HttpRequest req = base(path).GET().build(); + try { + HttpResponse res = exchange(req); + if (!RETRY_STATUSES.contains(res.statusCode())) { + return expectObject(res, expect); + } + } catch (HttpTimeoutException e) { + throw transportError(e); + } catch (IOException e) { + // refused or reset: retried below + } catch (InterruptedException e) { + throw interrupted(e); + } + try { + Thread.sleep(RETRY_PAUSE_MIN_MILLIS + ThreadLocalRandom.current().nextInt(RETRY_PAUSE_JITTER_MILLIS + 1)); + } catch (InterruptedException e) { + throw interrupted(e); + } return expectObject(send(req), expect); } @@ -289,23 +323,37 @@ public final class FelisApiClient { } return HttpRequest.newBuilder() .uri(uri) - .timeout(config.timeout()) + .timeout(config.requestTimeout()) .header("Authorization", "Bearer " + config.serviceToken()) .header("Accept", "application/json"); } private HttpResponse send(HttpRequest req) throws LinkException { try { - return http.send(req, HttpResponse.BodyHandlers.ofString()); + return exchange(req); } catch (IOException e) { - throw new LinkException(0, "transport_error", - "could not reach felis-api: " + e.getMessage(), e); + throw transportError(e); } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - throw new LinkException(0, "interrupted", "felis-api request interrupted", e); + throw interrupted(e); } } + private HttpResponse exchange(HttpRequest req) throws IOException, InterruptedException { + return http.send(req, HttpResponse.BodyHandlers.ofString()); + } + + // A refused connection arrives as a ConnectException with no message; the class + // name keeps the log line from reading "could not reach felis-api: null". + private static LinkException transportError(IOException e) { + String why = e.getMessage() == null ? e.getClass().getSimpleName() : e.getMessage(); + return new LinkException(0, "transport_error", "could not reach felis-api: " + why, e); + } + + private static LinkException interrupted(InterruptedException e) { + Thread.currentThread().interrupt(); + return new LinkException(0, "interrupted", "felis-api request interrupted", e); + } + private Map expectObject(HttpResponse res, int expect) throws LinkException { int status = res.statusCode(); if (status != expect) { diff --git a/plugins/shared/src/main/java/best/lolicon/felis/link/LinkClient.java b/plugins/shared/src/main/java/best/lolicon/felis/link/LinkClient.java index 6b248d8..e51dfa5 100644 --- a/plugins/shared/src/main/java/best/lolicon/felis/link/LinkClient.java +++ b/plugins/shared/src/main/java/best/lolicon/felis/link/LinkClient.java @@ -40,7 +40,7 @@ public final class LinkClient { public LinkClient(LinkConfig config) { this.config = Objects.requireNonNull(config, "config"); this.http = HttpClient.newBuilder() - .connectTimeout(config.timeout()) + .connectTimeout(config.connectTimeout()) .build(); } @@ -50,7 +50,7 @@ public final class LinkClient { String body = "{\"mc_uuid\":\"" + mcUuid + "\"}"; HttpRequest req = HttpRequest.newBuilder() .uri(URI.create(config.apiBaseUrl() + PATH)) - .timeout(config.timeout()) + .timeout(config.requestTimeout()) .header("Authorization", "Bearer " + config.serviceToken()) .header("Content-Type", "application/json") .header("Accept", "application/json") diff --git a/plugins/shared/src/main/java/best/lolicon/felis/link/LinkConfig.java b/plugins/shared/src/main/java/best/lolicon/felis/link/LinkConfig.java index 7cbe81d..3ab990b 100644 --- a/plugins/shared/src/main/java/best/lolicon/felis/link/LinkConfig.java +++ b/plugins/shared/src/main/java/best/lolicon/felis/link/LinkConfig.java @@ -10,29 +10,39 @@ import java.util.Objects; * and are never compiled in. Keeping them out of source is what lets the * tree stay domain- and credential-free; each platform's config loader is * responsible for sourcing them. + * + *

Two timeouts bound each call: {@code connectTimeout} for opening the TCP + * connection and {@code requestTimeout} for the whole exchange once it is open. A + * felis-api that is gone fails fast on the first; one that is slow is cut off by the + * second, so a blocked call cannot hold a plugin thread indefinitely. */ public final class LinkConfig { private final String apiBaseUrl; private final String serviceToken; - private final Duration timeout; + private final Duration connectTimeout; + private final Duration requestTimeout; - public LinkConfig(String apiBaseUrl, String serviceToken, Duration timeout) { + public LinkConfig(String apiBaseUrl, String serviceToken, Duration connectTimeout, Duration requestTimeout) { this.apiBaseUrl = stripTrailingSlash(Objects.requireNonNull(apiBaseUrl, "apiBaseUrl")); this.serviceToken = Objects.requireNonNull(serviceToken, "serviceToken"); - this.timeout = Objects.requireNonNull(timeout, "timeout"); + this.connectTimeout = Objects.requireNonNull(connectTimeout, "connectTimeout"); + this.requestTimeout = Objects.requireNonNull(requestTimeout, "requestTimeout"); if (this.apiBaseUrl.isEmpty()) { throw new IllegalArgumentException("apiBaseUrl is empty"); } if (this.serviceToken.isEmpty()) { throw new IllegalArgumentException("serviceToken is empty"); } - if (this.timeout.isZero() || this.timeout.isNegative()) { - throw new IllegalArgumentException("timeout must be positive"); + if (this.connectTimeout.isZero() || this.connectTimeout.isNegative()) { + throw new IllegalArgumentException("connectTimeout must be positive"); + } + if (this.requestTimeout.isZero() || this.requestTimeout.isNegative()) { + throw new IllegalArgumentException("requestTimeout must be positive"); } } public LinkConfig(String apiBaseUrl, String serviceToken) { - this(apiBaseUrl, serviceToken, Duration.ofSeconds(10)); + this(apiBaseUrl, serviceToken, Duration.ofSeconds(10), Duration.ofSeconds(10)); } public String apiBaseUrl() { @@ -43,8 +53,12 @@ public final class LinkConfig { return serviceToken; } - public Duration timeout() { - return timeout; + public Duration connectTimeout() { + return connectTimeout; + } + + public Duration requestTimeout() { + return requestTimeout; } private static String stripTrailingSlash(String u) { diff --git a/plugins/shared/src/main/java/best/lolicon/felis/link/LinkConfigLoader.java b/plugins/shared/src/main/java/best/lolicon/felis/link/LinkConfigLoader.java index 62f5b4e..b1432af 100644 --- a/plugins/shared/src/main/java/best/lolicon/felis/link/LinkConfigLoader.java +++ b/plugins/shared/src/main/java/best/lolicon/felis/link/LinkConfigLoader.java @@ -8,6 +8,7 @@ import java.nio.file.Files; import java.nio.file.Path; import java.time.Duration; import java.util.Properties; +import java.util.function.UnaryOperator; /** * LinkConfigLoader resolves a {@link LinkConfig} the same way on every platform: @@ -18,12 +19,25 @@ import java.util.Properties; * the source tree stays domain- and credential-free. On first run it writes a * commented template and then reports the values as missing, so an operator gets * a file to fill in rather than a silent half-configured plugin. + * + *

The two call timeouts ({@code connect-timeout-seconds}, + * {@code request-timeout-seconds}, or {@code FELIS_API_CONNECT_TIMEOUT_SECONDS} / + * {@code FELIS_API_REQUEST_TIMEOUT_SECONDS}) are optional and default to + * {@value #DEFAULT_TIMEOUT_SECONDS} s. A value that is not a whole number of seconds + * from 1 to {@value #MAX_TIMEOUT_SECONDS} fails the load: a typo quietly falling back + * to the default would hide the setting the operator meant to change. */ public final class LinkConfigLoader { public static final String ENV_URL = "FELIS_API_BASE_URL"; public static final String ENV_TOKEN = "FELIS_SERVICE_TOKEN"; private static final String KEY_URL = "api-base-url"; + public static final String ENV_CONNECT_TIMEOUT = "FELIS_API_CONNECT_TIMEOUT_SECONDS"; + public static final String ENV_REQUEST_TIMEOUT = "FELIS_API_REQUEST_TIMEOUT_SECONDS"; + static final int DEFAULT_TIMEOUT_SECONDS = 10; + static final int MAX_TIMEOUT_SECONDS = 120; private static final String KEY_TOKEN = "service-token"; + private static final String KEY_CONNECT_TIMEOUT = "connect-timeout-seconds"; + private static final String KEY_REQUEST_TIMEOUT = "request-timeout-seconds"; private LinkConfigLoader() { } @@ -33,9 +47,15 @@ public final class LinkConfigLoader { * if the file does not yet exist. * * @throws IOException if the file cannot be read/created, or if neither the - * environment nor the file supplies both required values. + * environment nor the file supplies both required values, or if a + * timeout is not a whole number of seconds in range. */ public static LinkConfig load(Path propertiesFile) throws IOException { + return load(propertiesFile, System::getenv); + } + + // load with the environment passed in, so the precedence rules are testable. + static LinkConfig load(Path propertiesFile, UnaryOperator env) throws IOException { Properties props = new Properties(); if (Files.exists(propertiesFile)) { try (InputStream in = Files.newInputStream(propertiesFile)) { @@ -45,14 +65,36 @@ public final class LinkConfigLoader { writeTemplate(propertiesFile); } - String url = firstNonBlank(System.getenv(ENV_URL), props.getProperty(KEY_URL)); - String token = firstNonBlank(System.getenv(ENV_TOKEN), props.getProperty(KEY_TOKEN)); + String url = firstNonBlank(env.apply(ENV_URL), props.getProperty(KEY_URL)); + String token = firstNonBlank(env.apply(ENV_TOKEN), props.getProperty(KEY_TOKEN)); if (isBlank(url) || isBlank(token)) { throw new IOException("set " + ENV_URL + "/" + ENV_TOKEN + " or fill in " + propertiesFile + " (" + KEY_URL + ", " + KEY_TOKEN + ")"); } - return new LinkConfig(url, token, Duration.ofSeconds(10)); + Duration connect = seconds(KEY_CONNECT_TIMEOUT, ENV_CONNECT_TIMEOUT, + firstNonBlank(env.apply(ENV_CONNECT_TIMEOUT), props.getProperty(KEY_CONNECT_TIMEOUT))); + Duration request = seconds(KEY_REQUEST_TIMEOUT, ENV_REQUEST_TIMEOUT, + firstNonBlank(env.apply(ENV_REQUEST_TIMEOUT), props.getProperty(KEY_REQUEST_TIMEOUT))); + return new LinkConfig(url, token, connect, request); + } + + private static Duration seconds(String key, String envName, String raw) throws IOException { + if (isBlank(raw)) { + return Duration.ofSeconds(DEFAULT_TIMEOUT_SECONDS); + } + String v = raw.trim(); + int n; + try { + n = Integer.parseInt(v); + } catch (NumberFormatException e) { + n = 0; + } + if (n < 1 || n > MAX_TIMEOUT_SECONDS) { + throw new IOException(key + " (" + envName + ") must be a whole number of seconds from 1 to " + + MAX_TIMEOUT_SECONDS + ", got \"" + v + "\""); + } + return Duration.ofSeconds(n); } private static void writeTemplate(Path file) throws IOException { @@ -72,7 +114,14 @@ public final class LinkConfigLoader { + KEY_URL + "=\n" + "#\n" + "# " + KEY_TOKEN + ": the internal service token (keep this secret).\n" - + KEY_TOKEN + "=\n"; + + KEY_TOKEN + "=\n" + + "#\n" + + "# Optional call timeouts in whole seconds (1-" + MAX_TIMEOUT_SECONDS + ", default " + + DEFAULT_TIMEOUT_SECONDS + "): " + KEY_CONNECT_TIMEOUT + " bounds opening the\n" + + "# connection, " + KEY_REQUEST_TIMEOUT + " the whole call once it is open\n" + + "# (" + ENV_CONNECT_TIMEOUT + ", " + ENV_REQUEST_TIMEOUT + ").\n" + + "#" + KEY_CONNECT_TIMEOUT + "=" + DEFAULT_TIMEOUT_SECONDS + "\n" + + "#" + KEY_REQUEST_TIMEOUT + "=" + DEFAULT_TIMEOUT_SECONDS + "\n"; try (OutputStream out = Files.newOutputStream(file)) { out.write(template.getBytes(StandardCharsets.UTF_8)); } diff --git a/plugins/shared/test/best/lolicon/felis/link/FelisApiClientRetryTest.java b/plugins/shared/test/best/lolicon/felis/link/FelisApiClientRetryTest.java new file mode 100644 index 0000000..7fd89a4 --- /dev/null +++ b/plugins/shared/test/best/lolicon/felis/link/FelisApiClientRetryTest.java @@ -0,0 +1,191 @@ +package best.lolicon.felis.link; + +import com.sun.net.httpserver.HttpExchange; +import com.sun.net.httpserver.HttpServer; + +import java.io.IOException; +import java.io.OutputStream; +import java.net.InetSocketAddress; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.atomic.AtomicInteger; + +/** + * FelisApiClientRetryTest pins which failures FelisApiClient repeats, against a stub + * felis-api that counts the requests each path receives: a GET is tried once more + * after a connection that dropped mid-reply or a 502/503/504, never a third time; a POST is + * never repeated (the first attempt may have landed); a GET that ran past the + * request timeout is cut off at that timeout and not repeated. + * + *

Run: {@code javac -d shared/src/main/java/best/lolicon/felis/link/*.java + * shared/test/best/lolicon/felis/link/FelisApiClientRetryTest.java && java -cp + * best.lolicon.felis.link.FelisApiClientRetryTest}. + */ +public final class FelisApiClientRetryTest { + + private static final String READY = "{\"name\":\"alpha\",\"phase\":\"Running\",\"ready\":true}"; + private static final Map hits = new ConcurrentHashMap<>(); + private static int checks; + + public static void main(String[] args) throws Exception { + HttpServer stub = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); + // The slow handler sleeps; the others must not queue behind it. + ExecutorService pool = Executors.newCachedThreadPool(); + stub.setExecutor(pool); + stub.createContext("/", FelisApiClientRetryTest::handle); + stub.start(); + try { + String base = "http://127.0.0.1:" + stub.getAddress().getPort(); + FelisApiClient api = new FelisApiClient(new LinkConfig( + base, "test-token", Duration.ofSeconds(2), Duration.ofMillis(400))); + getRetriesOnceAfter503(api); + getRetriesOnceAfterACutReply(api); + getGivesUpAfterTheSecondTry(api); + getIsNotRetriedOnAClientError(api); + postIsNeverRetried(api); + timedOutGetIsCutOffAndNotRetried(api); + } finally { + stub.stop(0); + pool.shutdownNow(); + } + System.out.println("FelisApiClientRetryTest OK (" + checks + " checks)"); + } + + // Each server name is one scripted behaviour; the n-th request to it picks the reply. + private static void handle(HttpExchange ex) throws IOException { + String path = ex.getRequestURI().getRawPath(); + ex.getRequestBody().readAllBytes(); + int n = hits.computeIfAbsent(ex.getRequestMethod() + " " + path, k -> new AtomicInteger()).incrementAndGet(); + switch (path) { + case "/api/v1/internal/servers/flaky/status": + reply(ex, n == 1 ? 503 : 200, n == 1 ? error("unavailable") : READY); + return; + case "/api/v1/internal/servers/cut/status": + if (n == 1) { + // Headers promise 1000 bytes, the connection drops after 10: felis-api + // dying mid-reply. (A drop before any response byte is not a usable + // probe: the JDK client resends such a GET on its own.) + ex.getResponseHeaders().set("Content-Type", "application/json"); + ex.sendResponseHeaders(200, 1000); + ex.getResponseBody().write("{\"name\":\"a".getBytes(StandardCharsets.UTF_8)); + ex.getResponseBody().flush(); + ex.close(); + return; + } + reply(ex, 200, READY); + return; + case "/api/v1/internal/servers/down/status": + reply(ex, 502, error("bad_gateway")); + return; + case "/api/v1/internal/servers/gone/status": + reply(ex, 404, error("not_found")); + return; + case "/api/v1/internal/servers/busy/wake": + reply(ex, 503, error("at_capacity")); + return; + case "/api/v1/internal/servers/slow/status": + try { + Thread.sleep(1500); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + reply(ex, 200, READY); + return; + default: + reply(ex, 500, error("unexpected")); + } + } + + private static void getRetriesOnceAfter503(FelisApiClient api) throws LinkException { + ServerView v = api.serverStatus("flaky"); + assertEq("flaky phase", "Running", v.phase()); + assertEq("flaky requests", 2, count("GET /api/v1/internal/servers/flaky/status")); + } + + private static void getRetriesOnceAfterACutReply(FelisApiClient api) throws LinkException { + ServerView v = api.serverStatus("cut"); + assertEq("cut phase", "Running", v.phase()); + assertEq("cut requests", 2, count("GET /api/v1/internal/servers/cut/status")); + } + + private static void getGivesUpAfterTheSecondTry(FelisApiClient api) { + LinkException e = expectFailure("down", () -> api.serverStatus("down")); + assertEq("down status", 502, e.statusCode()); + assertEq("down code", "bad_gateway", e.errorCode()); + assertEq("down requests", 2, count("GET /api/v1/internal/servers/down/status")); + } + + private static void getIsNotRetriedOnAClientError(FelisApiClient api) { + LinkException e = expectFailure("gone", () -> api.serverStatus("gone")); + assertEq("gone status", 404, e.statusCode()); + assertEq("gone requests", 1, count("GET /api/v1/internal/servers/gone/status")); + } + + private static void postIsNeverRetried(FelisApiClient api) { + UUID id = UUID.fromString("00000000-0000-0000-0000-00000000000a"); + LinkException e = expectFailure("busy", () -> api.wake("busy", id)); + assertEq("busy status", 503, e.statusCode()); + assertEq("busy code", "at_capacity", e.errorCode()); + assertEq("busy requests", 1, count("POST /api/v1/internal/servers/busy/wake")); + } + + private static void timedOutGetIsCutOffAndNotRetried(FelisApiClient api) { + long start = System.nanoTime(); + LinkException e = expectFailure("slow", () -> api.serverStatus("slow")); + long tookMillis = (System.nanoTime() - start) / 1_000_000; + assertEq("slow status", 0, e.statusCode()); + assertEq("slow code", "transport_error", e.errorCode()); + // 400 ms request timeout against a 1500 ms handler: well under the handler's + // time, which a missing or ignored request timeout would have to wait out. + if (tookMillis >= 1200) { + throw new AssertionError("slow: took " + tookMillis + " ms, want the 400 ms request timeout"); + } + checks++; + assertEq("slow requests", 1, count("GET /api/v1/internal/servers/slow/status")); + } + + // ---- harness ---- + + interface Call { + void run() throws LinkException; + } + + private static LinkException expectFailure(String what, Call call) { + try { + call.run(); + } catch (LinkException e) { + return e; + } + throw new AssertionError(what + ": call succeeded, want a LinkException"); + } + + private static int count(String key) { + AtomicInteger n = hits.get(key); + return n == null ? 0 : n.get(); + } + + private static String error(String code) { + return "{\"error\":{\"code\":\"" + code + "\",\"message\":\"stub\"}}"; + } + + private static void reply(HttpExchange ex, int status, String body) throws IOException { + byte[] b = body.getBytes(StandardCharsets.UTF_8); + ex.getResponseHeaders().set("Content-Type", "application/json"); + ex.sendResponseHeaders(status, b.length); + try (OutputStream os = ex.getResponseBody()) { + os.write(b); + } + } + + private static void assertEq(String what, Object want, Object got) { + if (!want.equals(got)) { + throw new AssertionError(what + ": got " + got + ", want " + want); + } + checks++; + } +} diff --git a/plugins/shared/test/best/lolicon/felis/link/LinkConfigLoaderTest.java b/plugins/shared/test/best/lolicon/felis/link/LinkConfigLoaderTest.java new file mode 100644 index 0000000..3df8f4e --- /dev/null +++ b/plugins/shared/test/best/lolicon/felis/link/LinkConfigLoaderTest.java @@ -0,0 +1,166 @@ +package best.lolicon.felis.link; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.time.Duration; +import java.util.Comparator; +import java.util.HashMap; +import java.util.Map; +import java.util.stream.Stream; + +/** + * LinkConfigLoaderTest checks how felis-link.properties and the environment combine: + * the environment wins, both required values must come from somewhere, a first run + * leaves a template and still refuses to start, and the call timeouts default to + * 10 s, take the operator's value, and refuse anything that is not a whole number + * of seconds from 1 to 120. + * + *

Run: {@code javac -d shared/src/main/java/best/lolicon/felis/link/*.java + * shared/test/best/lolicon/felis/link/LinkConfigLoaderTest.java && java -cp + * best.lolicon.felis.link.LinkConfigLoaderTest}. + */ +public final class LinkConfigLoaderTest { + + private static int checks; + private static Path dir; + + public static void main(String[] args) throws Exception { + dir = Files.createTempDirectory("felis-link-test"); + try { + fileAloneIsEnough(); + environmentWinsOverTheFile(); + missingTokenIsRefused(); + firstRunWritesATemplateAndRefuses(); + timeoutsComeFromFileOrEnvironment(); + badTimeoutsAreRefused(); + } finally { + try (Stream walk = Files.walk(dir)) { + walk.sorted(Comparator.reverseOrder()).forEach(p -> p.toFile().delete()); + } + } + System.out.println("LinkConfigLoaderTest OK (" + checks + " checks)"); + } + + private static void fileAloneIsEnough() throws IOException { + Path f = write("plain.properties", + "api-base-url=http://10.43.0.10:8081/\nservice-token=file-token\n"); + LinkConfig c = LinkConfigLoader.load(f, env()); + assertEq("file url (trailing slash stripped)", "http://10.43.0.10:8081", c.apiBaseUrl()); + assertEq("file token", "file-token", c.serviceToken()); + assertEq("default connect timeout", Duration.ofSeconds(10), c.connectTimeout()); + assertEq("default request timeout", Duration.ofSeconds(10), c.requestTimeout()); + } + + private static void environmentWinsOverTheFile() throws IOException { + Path f = write("both.properties", + "api-base-url=http://file:8081\nservice-token=file-token\n"); + LinkConfig c = LinkConfigLoader.load(f, env( + "FELIS_API_BASE_URL", "http://env:8081", + "FELIS_SERVICE_TOKEN", "env-token")); + assertEq("env url", "http://env:8081", c.apiBaseUrl()); + assertEq("env token", "env-token", c.serviceToken()); + + // A blank variable is unset, not an override to empty. + LinkConfig blank = LinkConfigLoader.load(f, env("FELIS_SERVICE_TOKEN", " ")); + assertEq("blank env token falls back", "file-token", blank.serviceToken()); + } + + private static void missingTokenIsRefused() throws IOException { + Path f = write("no-token.properties", "api-base-url=http://10.43.0.10:8081\n"); + IOException e = expectRefused("no token", f, env()); + assertContains("no token message", e.getMessage(), "service-token"); + } + + private static void firstRunWritesATemplateAndRefuses() throws IOException { + Path f = dir.resolve("fresh/felis-link.properties"); + expectRefused("first run", f, env()); + String template = new String(Files.readAllBytes(f), StandardCharsets.UTF_8); + assertContains("template url key", template, "\napi-base-url=\n"); + assertContains("template token key", template, "\nservice-token=\n"); + assertContains("template connect timeout", template, "\n#connect-timeout-seconds=10\n"); + assertContains("template request timeout", template, "\n#request-timeout-seconds=10\n"); + // The template's commented timeouts leave the defaults in force once filled in. + Files.write(f, (template.replace("\napi-base-url=\n", "\napi-base-url=http://x:8081\n") + .replace("\nservice-token=\n", "\nservice-token=t\n")).getBytes(StandardCharsets.UTF_8)); + LinkConfig c = LinkConfigLoader.load(f, env()); + assertEq("filled template connect timeout", Duration.ofSeconds(10), c.connectTimeout()); + assertEq("filled template request timeout", Duration.ofSeconds(10), c.requestTimeout()); + } + + private static void timeoutsComeFromFileOrEnvironment() throws IOException { + Path f = write("timeouts.properties", "api-base-url=http://x:8081\nservice-token=t\n" + + "connect-timeout-seconds=3\nrequest-timeout-seconds= 25 \n"); + LinkConfig c = LinkConfigLoader.load(f, env()); + assertEq("file connect timeout", Duration.ofSeconds(3), c.connectTimeout()); + assertEq("file request timeout", Duration.ofSeconds(25), c.requestTimeout()); + + LinkConfig e = LinkConfigLoader.load(f, env( + "FELIS_API_CONNECT_TIMEOUT_SECONDS", "1", + "FELIS_API_REQUEST_TIMEOUT_SECONDS", "120")); + assertEq("env connect timeout", Duration.ofSeconds(1), e.connectTimeout()); + assertEq("env request timeout", Duration.ofSeconds(120), e.requestTimeout()); + } + + private static void badTimeoutsAreRefused() throws IOException { + String base = "api-base-url=http://x:8081\nservice-token=t\n"; + String[][] bad = { + {"connect-timeout-seconds", "0"}, + {"connect-timeout-seconds", "-5"}, + {"connect-timeout-seconds", "ten"}, + {"request-timeout-seconds", "121"}, + {"request-timeout-seconds", "2.5"}, + {"request-timeout-seconds", "10s"}, + }; + for (String[] kv : bad) { + Path f = write("bad.properties", base + kv[0] + "=" + kv[1] + "\n"); + IOException e = expectRefused(kv[0] + "=" + kv[1], f, env()); + assertContains(kv[0] + "=" + kv[1] + " names the key", e.getMessage(), kv[0]); + assertContains(kv[0] + "=" + kv[1] + " quotes the value", e.getMessage(), "\"" + kv[1] + "\""); + } + Path f = write("env-bad.properties", base); + IOException e = expectRefused("env request timeout 0", f, env("FELIS_API_REQUEST_TIMEOUT_SECONDS", "0")); + assertContains("env bad names the variable", e.getMessage(), "FELIS_API_REQUEST_TIMEOUT_SECONDS"); + } + + // ---- harness ---- + + private static java.util.function.UnaryOperator env(String... kv) { + Map m = new HashMap<>(); + for (int i = 0; i < kv.length; i += 2) { + m.put(kv[i], kv[i + 1]); + } + return m::get; + } + + private static Path write(String name, String body) throws IOException { + Path f = dir.resolve(name); + Files.write(f, body.getBytes(StandardCharsets.UTF_8)); + return f; + } + + private static IOException expectRefused(String what, Path f, java.util.function.UnaryOperator env) { + try { + LinkConfigLoader.load(f, env); + } catch (IOException e) { + checks++; + return e; + } + throw new AssertionError(what + ": loaded, want an IOException"); + } + + private static void assertContains(String what, String got, String want) { + if (got == null || !got.contains(want)) { + throw new AssertionError(what + ": " + got + " does not contain " + want); + } + checks++; + } + + private static void assertEq(String what, Object want, Object got) { + if (!want.equals(got)) { + throw new AssertionError(what + ": got " + got + ", want " + want); + } + checks++; + } +} diff --git a/plugins/test.sh b/plugins/test.sh index bba9f4b..1d338f0 100644 --- a/plugins/test.sh +++ b/plugins/test.sh @@ -10,7 +10,11 @@ # velocity/test. They check what "compiles" cannot: the felis:control codec # round-trips every frame kind (spec §12), no server name can re-aim a # felis-api request path, only the lobby and the login gate may drive -# felis:control (and each only with its own frames), the link-status outage +# felis:control (and each only with its own frames), a GET is retried once +# after a dropped reply or a 502/503/504 and a POST never, the request timeout +# cuts a slow call off, felis-link.properties timeouts are validated, the +# felis-api call pool refuses instead of growing and a repeating task never +# overlaps itself, the link-status outage # fallback fails closed outside its window, a proxy restarted during an API # outage routes on the last saved server list (and only until a fetch # succeeds), /invite prompts cannot double-fire @@ -71,6 +75,18 @@ javac -d "$work/shared-classes" \ plugins/shared/test/best/lolicon/felis/link/FelisApiClientTest.java java -cp "$work/shared-classes" best.lolicon.felis.link.FelisApiClientTest +echo "==> FelisApiClientRetryTest (which failures are retried, request timeout, shared)" +javac -d "$work/shared-classes" \ + plugins/shared/src/main/java/best/lolicon/felis/link/*.java \ + plugins/shared/test/best/lolicon/felis/link/FelisApiClientRetryTest.java +java -cp "$work/shared-classes" best.lolicon.felis.link.FelisApiClientRetryTest + +echo "==> LinkConfigLoaderTest (env/file precedence, template, timeouts, shared)" +javac -d "$work/shared-classes" \ + plugins/shared/src/main/java/best/lolicon/felis/link/*.java \ + plugins/shared/test/best/lolicon/felis/link/LinkConfigLoaderTest.java +java -cp "$work/shared-classes" best.lolicon.felis.link.LinkConfigLoaderTest + echo "==> ControlPolicyTest (who may send what on felis:control, velocity)" mkdir -p "$work/policy-classes" javac -d "$work/policy-classes" \ @@ -93,6 +109,14 @@ javac -d "$work/list-classes" \ plugins/velocity/test/best/lolicon/felis/velocity/ServerListSourceTest.java java -cp "$work/list-classes" best.lolicon.felis.velocity.ServerListSourceTest +echo "==> ApiPoolTest (bounded felis-api pool, non-overlapping repeats, velocity)" +mkdir -p "$work/pool-classes" +javac -d "$work/pool-classes" \ + plugins/velocity/src/main/java/best/lolicon/felis/velocity/BoundedExecutor.java \ + plugins/velocity/src/main/java/best/lolicon/felis/velocity/SkipIfRunning.java \ + plugins/velocity/test/best/lolicon/felis/velocity/ApiPoolTest.java +java -cp "$work/pool-classes" best.lolicon.felis.velocity.ApiPoolTest + echo "==> InviteBookTest (/invite prompt store, velocity)" mkdir -p "$work/velocity-classes" javac -d "$work/velocity-classes" \ diff --git a/plugins/velocity/src/main/java/best/lolicon/felis/velocity/BoundedExecutor.java b/plugins/velocity/src/main/java/best/lolicon/felis/velocity/BoundedExecutor.java new file mode 100644 index 0000000..4fbb9f1 --- /dev/null +++ b/plugins/velocity/src/main/java/best/lolicon/felis/velocity/BoundedExecutor.java @@ -0,0 +1,68 @@ +package best.lolicon.felis.velocity; + +import java.util.concurrent.ArrayBlockingQueue; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.Consumer; + +/** + * BoundedExecutor runs the plugin's felis-api calls on a fixed number of threads with + * a bounded wait queue. Velocity's scheduler starts a thread for every task it is + * handed, so while felis-api was slow a burst of joins, menu clicks and commands + * stacked blocking calls without limit, and each of them landed on the API the moment + * it recovered. Here at most {@code threads} calls run and at most {@code queue} wait; + * past that {@link #submit} refuses at once, so the caller can say "busy" instead of + * leaving the player waiting on a call that would run minutes later. + * + *

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

Run: {@code javac -d velocity/src/main/java/best/lolicon/felis/velocity/{BoundedExecutor,SkipIfRunning}.java + * velocity/test/best/lolicon/felis/velocity/ApiPoolTest.java && java -cp + * best.lolicon.felis.velocity.ApiPoolTest}. + */ +public final class ApiPoolTest { + + private static int checks; + + public static void main(String[] args) throws Exception { + fullPoolRefusesInsteadOfGrowing(); + throwingTaskIsReportedAndThePoolLives(); + shutdownRefuses(); + repeatingTaskNeverOverlaps(); + throwingRunDoesNotWedgeTheGuard(); + System.out.println("ApiPoolTest OK (" + checks + " checks)"); + } + + // Two threads, two queue slots: four tasks are taken (two running, two waiting), the + // fifth is refused; once the blockers finish, every accepted task runs. + private static void fullPoolRefusesInsteadOfGrowing() throws Exception { + BoundedExecutor pool = new BoundedExecutor("t", 2, 2, t -> { }); + CountDownLatch started = new CountDownLatch(2); + CountDownLatch release = new CountDownLatch(1); + AtomicInteger ran = new AtomicInteger(); + Runnable blocker = () -> { + started.countDown(); + await(release); + ran.incrementAndGet(); + }; + assertEq("first taken", true, pool.submit(blocker)); + assertEq("second taken", true, pool.submit(blocker)); + assertEq("both blockers running", true, started.await(5, TimeUnit.SECONDS)); + assertEq("running while blocked", 2, pool.running()); + assertEq("third queued", true, pool.submit(ran::incrementAndGet)); + assertEq("fourth queued", true, pool.submit(ran::incrementAndGet)); + assertEq("waiting while blocked", 2, pool.waiting()); + assertEq("fifth refused", false, pool.submit(ran::incrementAndGet)); + assertEq("nothing finished yet", 0, ran.get()); + + release.countDown(); + CountDownLatch drained = new CountDownLatch(1); + waitFor(() -> ran.get() == 4, "accepted tasks all run"); + assertEq("refused task never ran", 4, ran.get()); + // Room again once drained. + assertEq("taken after drain", true, pool.submit(drained::countDown)); + assertEq("post-drain task ran", true, drained.await(5, TimeUnit.SECONDS)); + pool.shutdown(); + } + + private static void throwingTaskIsReportedAndThePoolLives() throws Exception { + List errors = new CopyOnWriteArrayList<>(); + BoundedExecutor pool = new BoundedExecutor("t", 1, 1, errors::add); + pool.submit(() -> { + throw new IllegalStateException("boom"); + }); + CountDownLatch next = new CountDownLatch(1); + waitFor(() -> errors.size() == 1, "error reported"); + assertEq("reported message", "boom", errors.get(0).getMessage()); + assertEq("taken after a throw", true, pool.submit(next::countDown)); + assertEq("task after a throw ran", true, next.await(5, TimeUnit.SECONDS)); + pool.shutdown(); + } + + private static void shutdownRefuses() { + BoundedExecutor pool = new BoundedExecutor("t", 1, 1, t -> { }); + pool.shutdown(); + assertEq("refused after shutdown", false, pool.submit(() -> { })); + } + + // A second run that falls due while the first is inside the task is skipped, not + // started beside it; the next run after the first finishes goes ahead. + private static void repeatingTaskNeverOverlaps() throws Exception { + CountDownLatch inside = new CountDownLatch(1); + CountDownLatch release = new CountDownLatch(1); + AtomicInteger entered = new AtomicInteger(); + AtomicInteger concurrent = new AtomicInteger(); + AtomicInteger maxConcurrent = new AtomicInteger(); + SkipIfRunning guarded = new SkipIfRunning(() -> { + int now = concurrent.incrementAndGet(); + maxConcurrent.accumulateAndGet(now, Math::max); + entered.incrementAndGet(); + if (entered.get() == 1) { + inside.countDown(); + await(release); + } + concurrent.decrementAndGet(); + }); + Thread first = new Thread(guarded); + first.start(); + assertEq("first run inside", true, inside.await(5, TimeUnit.SECONDS)); + guarded.run(); // due while the first is still running + guarded.run(); + assertEq("overlapping runs skipped", 1, entered.get()); + release.countDown(); + first.join(5000); + guarded.run(); + assertEq("run after the first finished", 2, entered.get()); + assertEq("never two at once", 1, maxConcurrent.get()); + } + + private static void throwingRunDoesNotWedgeTheGuard() { + AtomicInteger calls = new AtomicInteger(); + SkipIfRunning guarded = new SkipIfRunning(() -> { + if (calls.incrementAndGet() == 1) { + throw new IllegalStateException("felis-api down"); + } + }); + try { + guarded.run(); + throw new AssertionError("the first run's exception was swallowed"); + } catch (IllegalStateException e) { + assertEq("first run's exception propagates", "felis-api down", e.getMessage()); + } + guarded.run(); + assertEq("runs after a throw", 2, calls.get()); + } + + // ---- harness ---- + + interface Cond { + boolean ok(); + } + + private static void waitFor(Cond cond, String what) throws InterruptedException { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5); + while (!cond.ok()) { + if (System.nanoTime() > deadline) { + throw new AssertionError(what + ": timed out"); + } + Thread.sleep(5); + } + checks++; + } + + private static void await(CountDownLatch l) { + try { + l.await(5, TimeUnit.SECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + + private static void assertEq(String what, Object want, Object got) { + if (!want.equals(got)) { + throw new AssertionError(what + ": got " + got + ", want " + want); + } + checks++; + } +}