From c370b2068b570a539d64a32fba1c50ecde3fdec3 Mon Sep 17 00:00:00 2001 From: Frederick Baier Date: Sun, 6 Sep 2026 19:49:22 +0200 Subject: [PATCH] fix: Propagate player disconnect and server-switch failures --- .../integration/player/PlayerIntegration.java | 41 +++- .../player/PlayerIntegrationTest.java | 197 ++++++++++++++++++ 2 files changed, 232 insertions(+), 6 deletions(-) create mode 100644 api/src/test/java/app/simplecloud/api/internal/integration/player/PlayerIntegrationTest.java diff --git a/api/src/main/java/app/simplecloud/api/internal/integration/player/PlayerIntegration.java b/api/src/main/java/app/simplecloud/api/internal/integration/player/PlayerIntegration.java index 3c6b2c0..b0a7ac0 100644 --- a/api/src/main/java/app/simplecloud/api/internal/integration/player/PlayerIntegration.java +++ b/api/src/main/java/app/simplecloud/api/internal/integration/player/PlayerIntegration.java @@ -9,6 +9,8 @@ import java.time.Duration; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; +import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.BiFunction; @@ -29,8 +31,12 @@ public class PlayerIntegration { private BiFunction> connectHandler; public PlayerIntegration(CloudApiImpl cloudApi) { - this.natsConnection = cloudApi.getNatsConnection(); - this.networkId = cloudApi.getNetworkId(); + this(cloudApi.getNatsConnection(), cloudApi.getNetworkId()); + } + + PlayerIntegration(Connection natsConnection, String networkId) { + this.natsConnection = natsConnection; + this.networkId = networkId; } /** @@ -82,6 +88,7 @@ public CompletableFuture login( /** * Notifies the controller that a player disconnected. + * Completes exceptionally if the controller rejects the request or no valid reply is received. */ public CompletableFuture disconnect(String playerId) { return CompletableFuture.runAsync(() -> { @@ -91,14 +98,25 @@ public CompletableFuture disconnect(String playerId) { .build(); String subject = networkId + ".player.disconnect"; - natsConnection.request(subject, request.toByteArray(), REQUEST_TIMEOUT); - } catch (Exception ignored) { + Message response = natsConnection.request(subject, request.toByteArray(), REQUEST_TIMEOUT); + if (response == null) { + throw new TimeoutException("No controller reply within " + REQUEST_TIMEOUT.toSeconds() + " seconds"); + } + if (!PlayerDisconnectResponse.parseFrom(response.getData()).getSuccess()) { + throw new IllegalStateException("Controller rejected the disconnect request"); + } + } catch (Exception e) { + if (e instanceof InterruptedException) { + Thread.currentThread().interrupt(); + } + throw new CompletionException("Failed to disconnect player " + playerId + ": " + e.getMessage(), e); } }); } /** * Notifies the controller that a player switched servers. + * Completes exceptionally if the controller rejects the request or no valid reply is received. */ public CompletableFuture serverSwitch(String playerId, String newServerName) { return CompletableFuture.runAsync(() -> { @@ -109,8 +127,19 @@ public CompletableFuture serverSwitch(String playerId, String newServerNam .build(); String subject = networkId + ".player.switch"; - natsConnection.request(subject, request.toByteArray(), REQUEST_TIMEOUT); - } catch (Exception ignored) { + Message response = natsConnection.request(subject, request.toByteArray(), REQUEST_TIMEOUT); + if (response == null) { + throw new TimeoutException("No controller reply within " + REQUEST_TIMEOUT.toSeconds() + " seconds"); + } + if (!PlayerServerSwitchResponse.parseFrom(response.getData()).getSuccess()) { + throw new IllegalStateException("Controller rejected the server-switch request"); + } + } catch (Exception e) { + if (e instanceof InterruptedException) { + Thread.currentThread().interrupt(); + } + throw new CompletionException("Failed to switch player " + playerId + " to server " + + newServerName + ": " + e.getMessage(), e); } }); } diff --git a/api/src/test/java/app/simplecloud/api/internal/integration/player/PlayerIntegrationTest.java b/api/src/test/java/app/simplecloud/api/internal/integration/player/PlayerIntegrationTest.java new file mode 100644 index 0000000..836807b --- /dev/null +++ b/api/src/test/java/app/simplecloud/api/internal/integration/player/PlayerIntegrationTest.java @@ -0,0 +1,197 @@ +package app.simplecloud.api.internal.integration.player; + +import build.buf.gen.simplecloud.player.v2.PlayerDisconnectRequest; +import build.buf.gen.simplecloud.player.v2.PlayerDisconnectResponse; +import build.buf.gen.simplecloud.player.v2.PlayerLoginResponse; +import build.buf.gen.simplecloud.player.v2.PlayerServerSwitchRequest; +import build.buf.gen.simplecloud.player.v2.PlayerServerSwitchResponse; +import com.google.protobuf.InvalidProtocolBufferException; +import io.nats.client.Connection; +import io.nats.client.Message; +import io.nats.client.impl.NatsMessage; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; + +import java.lang.reflect.Proxy; +import java.time.Duration; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.jupiter.api.Assertions.*; + +@Timeout(5) +class PlayerIntegrationTest { + + private static final String NETWORK_ID = "test-network"; + private static final String PLAYER_ID = "11111111-1111-1111-1111-111111111111"; + + enum Operation { + DISCONNECT("disconnect"), SERVER_SWITCH("switch"); + + final String subjectSuffix; + + Operation(String subjectSuffix) { + this.subjectSuffix = subjectSuffix; + } + + CompletableFuture invoke(PlayerIntegration integration) { + return this == DISCONNECT + ? integration.disconnect(PLAYER_ID) + : integration.serverSwitch(PLAYER_ID, "lobby-1"); + } + + byte[] response(boolean success) { + return this == DISCONNECT + ? PlayerDisconnectResponse.newBuilder().setSuccess(success) + .setSessionDurationSeconds(success ? 42 : 0).build().toByteArray() + : PlayerServerSwitchResponse.newBuilder().setSuccess(success).build().toByteArray(); + } + } + + @ParameterizedTest + @EnumSource(Operation.class) + void successfulReplyCompletesNormallyAndSendsExpectedRequest(Operation operation) { + AtomicInteger requests = new AtomicInteger(); + PlayerIntegration integration = integration((subject, data, timeout) -> { + requests.incrementAndGet(); + assertEquals(NETWORK_ID + ".player." + operation.subjectSuffix, subject); + assertEquals(Duration.ofSeconds(30), timeout); + if (operation == Operation.DISCONNECT) { + assertEquals(PLAYER_ID, PlayerDisconnectRequest.parseFrom(data).getPlayerId()); + } else { + PlayerServerSwitchRequest request = PlayerServerSwitchRequest.parseFrom(data); + assertEquals(PLAYER_ID, request.getPlayerId()); + assertEquals("lobby-1", request.getNewServerName()); + } + return message(operation.response(true)); + }); + + assertNull(operation.invoke(integration).join()); + assertEquals(1, requests.get()); + } + + @Test + void alreadyOfflineDisconnectCompletesNormally() { + // Controller's idempotent reply has success=true and no session duration. + PlayerIntegration integration = integration((subject, data, timeout) -> message( + PlayerDisconnectResponse.newBuilder().setSuccess(true).build().toByteArray() + )); + + assertNull(integration.disconnect(PLAYER_ID).join()); + } + + @ParameterizedTest + @EnumSource(Operation.class) + void unsuccessfulReplyReachesTheProxyErrorCallback(Operation operation) { + AtomicInteger requests = new AtomicInteger(); + PlayerIntegration integration = integration((subject, data, timeout) -> { + requests.incrementAndGet(); + return message(operation.response(false)); + }); + AtomicReference reportedFailure = new AtomicReference<>(); + CompletableFuture future = operation.invoke(integration); + + // Both proxy listeners attach this callback to the returned future. + future.exceptionally(error -> { + reportedFailure.set(error); + return null; + }).join(); + + assertNotNull(reportedFailure.get()); + CompletionException failure = assertFailure(operation, future); + assertInstanceOf(IllegalStateException.class, failure.getCause()); + assertTrue(failure.getMessage().contains("rejected")); + assertEquals(1, requests.get()); + } + + @ParameterizedTest + @EnumSource(Operation.class) + void missingReplyFailsAsTimeout(Operation operation) { + PlayerIntegration integration = integration((subject, data, timeout) -> null); + + CompletionException failure = assertFailure(operation, operation.invoke(integration)); + + assertInstanceOf(TimeoutException.class, failure.getCause()); + } + + @ParameterizedTest + @EnumSource(Operation.class) + void malformedReplyPreservesDecodeFailure(Operation operation) { + PlayerIntegration integration = integration((subject, data, timeout) -> message(new byte[]{(byte) 0xff})); + + CompletionException failure = assertFailure(operation, operation.invoke(integration)); + + assertInstanceOf(InvalidProtocolBufferException.class, failure.getCause()); + } + + @ParameterizedTest + @EnumSource(Operation.class) + void transportFailurePreservesCause(Operation operation) { + IllegalStateException transportFailure = new IllegalStateException("connection closed"); + PlayerIntegration integration = integration((subject, data, timeout) -> { + throw transportFailure; + }); + + CompletionException failure = assertFailure(operation, operation.invoke(integration)); + + assertSame(transportFailure, failure.getCause()); + } + + @ParameterizedTest + @EnumSource(Operation.class) + void interruptionPreservesCause(Operation operation) { + InterruptedException interruption = new InterruptedException("request interrupted"); + PlayerIntegration integration = integration((subject, data, timeout) -> { + throw interruption; + }); + + CompletionException failure = assertFailure(operation, operation.invoke(integration)); + + assertSame(interruption, failure.getCause()); + } + + @Test + void loginStillReturnsControllerFailureAsLoginResult() { + PlayerIntegration integration = integration((subject, data, timeout) -> message( + PlayerLoginResponse.newBuilder().setSuccess(false).setErrorMessage("login rejected").build().toByteArray() + )); + + LoginResult result = integration.login(PLAYER_ID, "Player", "Player", "proxy-1", "hash", + "en_US", 1, true, null).join(); + + assertFalse(result.isSuccess()); + assertEquals("login rejected", result.getErrorMessage()); + } + + private static CompletionException assertFailure(Operation operation, CompletableFuture future) { + CompletionException failure = assertThrows(CompletionException.class, future::join); + assertTrue(failure.getMessage().contains(PLAYER_ID)); + assertTrue(failure.getMessage().contains(operation.subjectSuffix)); + return failure; + } + + private static Message message(byte[] data) { + return NatsMessage.builder().subject("reply").data(data).build(); + } + + private static PlayerIntegration integration(RequestHandler handler) { + Connection connection = (Connection) Proxy.newProxyInstance( + Connection.class.getClassLoader(), new Class[]{Connection.class}, (proxy, method, args) -> { + if (method.getName().equals("request") && args.length == 3) { + return handler.request((String) args[0], (byte[]) args[1], (Duration) args[2]); + } + throw new AssertionError("Unexpected connection call: " + method.getName()); + }); + return new PlayerIntegration(connection, NETWORK_ID); + } + + @FunctionalInterface + private interface RequestHandler { + Message request(String subject, byte[] data, Duration timeout) throws Exception; + } +}