From b2f0b2f03ee75ff7b431c5adb232e41043772270 Mon Sep 17 00:00:00 2001 From: Minekube AI Engineer <1535738075139801220+minekube-ai@users.noreply.github.com> Date: Wed, 19 Aug 2026 15:43:46 +0000 Subject: [PATCH] fix(connect): recover closed admission coordinator on plugin reload A plugin reload (disable -> enable on the same plugin instance) left the parent-scoped BedrockAdmissionCoordinator permanently closed: disable() closes it, while WatcherRegister/WatchClient are recreated per enable in a child injector, so the new watch routed every session proposal into the closed coordinator, throwing IllegalStateException 'Bedroom admission coordinator is closed' which failed the WebSocket and put the watch into an endless reconnect loop until JVM restart. - ConnectPlatform.enable() now resets the shared coordinator before the new cycle's watcher binds: pending admissions are dropped (stale across cycles), the identity registry is reopened, and the cleanup executor that close() shut down is recreated. - WatchClient now rejects a proposal over the wire (SessionRejection) when the coordinator is closed instead of letting the ISE escape and kill the WebSocket stream. Regression tests (RED on unfixed main, GREEN with fix): - ConnectPlatformReloadRecoveryTest: enable() after disable() recovers the coordinator; full enable -> proposal -> disable -> enable -> proposal reload cycle delivers the second proposal instead of the ISE. - WatchClientTest.closedCoordinatorRejectsProposalInsteadOfFailingTheWatch: a proposal into a closed coordinator is rejected on the wire, the watch stream stays up. --- .../com/minekube/connect/ConnectPlatform.java | 7 + .../bedrock/BedrockAdmissionCoordinator.java | 24 +- .../VerifiedBedrockIdentityRegistry.java | 12 + .../minekube/connect/watch/WatchClient.java | 21 +- .../ConnectPlatformReloadRecoveryTest.java | 221 ++++++++++++++++++ .../connect/watch/WatchClientTest.java | 61 +++++ 6 files changed, 342 insertions(+), 4 deletions(-) create mode 100644 core/src/test/java/com/minekube/connect/ConnectPlatformReloadRecoveryTest.java diff --git a/core/src/main/java/com/minekube/connect/ConnectPlatform.java b/core/src/main/java/com/minekube/connect/ConnectPlatform.java index 2e07a643a..36c8796ea 100644 --- a/core/src/main/java/com/minekube/connect/ConnectPlatform.java +++ b/core/src/main/java/com/minekube/connect/ConnectPlatform.java @@ -128,6 +128,13 @@ public boolean enable(Module... postInitializeModules) { return false; } + // The admission coordinator is a parent-scoped singleton shared by every enable() cycle. + // A previous disable() (plugin reload) closed it; reopen it with fresh state before the new + // cycle's watcher binds, otherwise every session proposal throws ISE and kills the watch. + if (admissionCoordinator != null) { + admissionCoordinator.reset(); + } + try { if (!injector.inject()) { // TODO && !bootstrap.getGeyserConfig().isUseDirectConnection() logger.error("Failed to inject the packet listener!"); diff --git a/core/src/main/java/com/minekube/connect/bedrock/BedrockAdmissionCoordinator.java b/core/src/main/java/com/minekube/connect/bedrock/BedrockAdmissionCoordinator.java index 8201d27c8..f7fdf2f40 100644 --- a/core/src/main/java/com/minekube/connect/bedrock/BedrockAdmissionCoordinator.java +++ b/core/src/main/java/com/minekube/connect/bedrock/BedrockAdmissionCoordinator.java @@ -32,7 +32,7 @@ public final class BedrockAdmissionCoordinator implements AutoCloseable { private final VerifiedBedrockIdentityRegistry identities; private final BedrockPrincipalConsumer principalConsumer; - private final ScheduledExecutorService cleanupExecutor; + private ScheduledExecutorService cleanupExecutor; private final Map admissions = new HashMap<>(); private final Map latestBySession = new HashMap<>(); private final Map players = new IdentityHashMap<>(); @@ -204,6 +204,28 @@ public synchronized void close() { cleanupExecutor.shutdownNow(); } + /** + * Reopens the coordinator for a fresh enable() cycle after {@link #close()}. + * + *

The coordinator is a parent-scoped singleton shared by every enable() cycle, so a plugin + * reload (disable → enable without platform reconstruction) must not leave it permanently + * closed: the next cycle's watcher would route every session proposal into the closed + * coordinator, throwing {@code IllegalStateException: ... coordinator is closed} and killing + * the watch. Clears all pending admissions, reopens the identity registry, and recreates the + * cleanup executor that {@link #close()} shut down. + */ + public synchronized void reset() { + new ArrayList<>(admissions.values()).forEach(this::removeAdmission); + admissions.clear(); + latestBySession.clear(); + players.clear(); + closed = false; + identities.reset(); + if (cleanupExecutor.isShutdown()) { + cleanupExecutor = newCleanupExecutor(); + } + } + private synchronized void expire(AdmissionToken token, Admission expected) { if (admissions.get(token) == expected) { removeAdmission(expected); diff --git a/core/src/main/java/com/minekube/connect/bedrock/VerifiedBedrockIdentityRegistry.java b/core/src/main/java/com/minekube/connect/bedrock/VerifiedBedrockIdentityRegistry.java index 91d9b8125..0e3fc9578 100644 --- a/core/src/main/java/com/minekube/connect/bedrock/VerifiedBedrockIdentityRegistry.java +++ b/core/src/main/java/com/minekube/connect/bedrock/VerifiedBedrockIdentityRegistry.java @@ -84,6 +84,18 @@ public synchronized void close() { identities.clear(); } + /** + * Clears recorded identities and reopens the registry after {@link #close()}. + * + *

The registry is a parent-scoped singleton shared by every enable() cycle, so a plugin + * reload (disable → enable on the same platform) must be able to reuse it. Entries from a + * previous cycle can never match a new session's players, so clearing is safe. + */ + public synchronized void reset() { + closed = false; + identities.clear(); + } + private void ensureOpen() { if (closed) { throw new IllegalStateException("Bedrock identity registry is closed"); diff --git a/core/src/main/java/com/minekube/connect/watch/WatchClient.java b/core/src/main/java/com/minekube/connect/watch/WatchClient.java index a3d03553f..ba05b4050 100644 --- a/core/src/main/java/com/minekube/connect/watch/WatchClient.java +++ b/core/src/main/java/com/minekube/connect/watch/WatchClient.java @@ -28,6 +28,8 @@ import com.google.inject.Inject; import com.google.inject.name.Named; import com.google.protobuf.InvalidProtocolBufferException; +import com.google.rpc.Code; +import com.google.rpc.Status; import com.minekube.connect.bedrock.BedrockIdentityKeyProvider; import com.minekube.connect.bedrock.BedrockIdentityReadiness; import com.minekube.connect.bedrock.BedrockIdentityReadiness.Transport; @@ -186,9 +188,22 @@ public void onMessage(@NotNull WebSocket webSocket, @NotNull ByteString bytes) { .toByteArray() )); }; - SessionProposal prop = admissionCoordinator == null - ? new SessionProposal(res.getSession(), rejectProposal) - : admissionCoordinator.proposal(res.getSession(), rejectProposal, "", ""); + SessionProposal prop; + try { + prop = admissionCoordinator == null + ? new SessionProposal(res.getSession(), rejectProposal) + : admissionCoordinator.proposal(res.getSession(), rejectProposal, "", ""); + } catch (IllegalStateException e) { + // The coordinator is only closed when a disable() is tearing the plugin down or + // a reload failed to recover it. Reject the proposal over the wire instead of + // letting the exception escape and fail the WebSocket, which would put the watch + // into an endless reconnect loop; the next enable() recovers the coordinator. + rejectProposal.accept(Status.newBuilder() + .setCode(Code.INTERNAL_VALUE) + .setMessage("Bedrock admission coordinator is closed") + .build()); + return; + } watcher.onProposal(prop); } diff --git a/core/src/test/java/com/minekube/connect/ConnectPlatformReloadRecoveryTest.java b/core/src/test/java/com/minekube/connect/ConnectPlatformReloadRecoveryTest.java new file mode 100644 index 000000000..b86acdd2d --- /dev/null +++ b/core/src/test/java/com/minekube/connect/ConnectPlatformReloadRecoveryTest.java @@ -0,0 +1,221 @@ +package com.minekube.connect; + +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.google.inject.AbstractModule; +import com.google.inject.Guice; +import com.google.inject.Injector; +import com.google.inject.Module; +import com.google.inject.name.Names; +import com.minekube.connect.api.ConnectApi; +import com.minekube.connect.api.SimpleConnectApi; +import com.minekube.connect.api.inject.PlatformInjector; +import com.minekube.connect.api.logger.ConnectLogger; +import com.minekube.connect.bedrock.BedrockAdmissionCoordinator; +import com.minekube.connect.bedrock.BedrockIdentityKeyProvider; +import com.minekube.connect.bedrock.BedrockIdentityReadiness; +import com.minekube.connect.bedrock.BedrockPrincipalReadiness; +import com.minekube.connect.bedrock.VerifiedBedrockIdentityRegistry; +import com.minekube.connect.config.ConnectConfig; +import com.minekube.connect.inject.CommonPlatformInjector; +import com.minekube.connect.module.WatcherModule; +import com.minekube.connect.platform.util.PlatformUtils; +import com.minekube.connect.tunnel.Tunneler; +import com.minekube.connect.tunnel.p2p.Libp2pEndpoint; +import com.minekube.connect.util.Metrics; +import com.minekube.connect.util.UpdateChecker; +import com.minekube.connect.watch.SessionProposal; +import java.lang.reflect.Field; +import java.net.InetSocketAddress; +import minekube.connect.v1alpha1.WatchServiceOuterClass.GameProfile; +import minekube.connect.v1alpha1.WatchServiceOuterClass.Player; +import minekube.connect.v1alpha1.WatchServiceOuterClass.Session; +import minekube.connect.v1alpha1.WatchServiceOuterClass.WatchResponse; +import okhttp3.OkHttpClient; +import okhttp3.Request; +import okhttp3.WebSocket; +import okhttp3.WebSocketListener; +import okio.ByteString; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; + +/** + * Regression guard: a plugin reload (disable → enable on the same platform instance) must not leave + * the parent-scoped {@link BedrockAdmissionCoordinator} closed, otherwise every session proposal + * throws {@code IllegalStateException: ... coordinator is closed} and the watch dies. + * + *

{@code ConnectPlatform.disable()} closes the coordinator (a parent-scoped singleton shared by + * every enable() cycle), while {@code WatcherRegister}/{@code WatchClient} are recreated per enable + * in a child injector. {@code ConnectPlatform.enable()} must recover the coordinator before the new + * cycle's watcher binds, or the first proposal after a reload kills the WebSocket and the watch + * reconnects forever (public report: "Connection error with WatchService: ... coordinator is + * closed"). + */ +class ConnectPlatformReloadRecoveryTest { + + static { + // bstats MetricsBase refuses to construct when its classes are not shaded/relocated + // (release-jar-only). The documented test escape hatch; the test config disables metrics + // anyway, so Metrics stays inert. + System.setProperty("bstats.relocatecheck", "false"); + } + + /** The enable() wiring alone: a coordinator closed by the previous disable() is usable again. */ + @Test + void enableRecoversCoordinatorClosedByPreviousDisable() throws Exception { + VerifiedBedrockIdentityRegistry registry = new VerifiedBedrockIdentityRegistry(); + BedrockAdmissionCoordinator coordinator = new BedrockAdmissionCoordinator(registry); + PlatformInjector platformInjector = mock(PlatformInjector.class); + when(platformInjector.inject()).thenReturn(true); + Injector guice = mock(Injector.class); + when(guice.createChildInjector(any(Module[].class))).thenReturn(guice); + when(guice.getInstance(UpdateChecker.class)).thenReturn(mock(UpdateChecker.class)); + ConnectPlatform platform = new ConnectPlatform( + mock(ConnectApi.class), platformInjector, mock(ConnectLogger.class), guice, coordinator); + setConfig(platform, new ConnectConfig()); + + coordinator.close(); // simulate the previous disable() cycle + assertTrue(platform.enable()); + + // The shared coordinator must accept proposals again; a closed one throws ISE here. + Session session = session("session-recovered"); + SessionProposal proposal = assertDoesNotThrow( + () -> coordinator.proposal(session, reason -> { }, "", "")); + coordinator.discard(proposal); + } + + /** + * End-to-end reload cycle: enable → deliver proposal (join accepted) → disable (watcher stopped, + * coordinator closed) → enable with a new child injector → deliver proposal again. The second + * proposal must be accepted, not fail with the "coordinator is closed" ISE that kills the watch. + */ + @Test + void reloadCycleDeliversSessionProposalThroughRecoveredCoordinator() throws Exception { + ConnectConfig config = new ConnectConfig(); + disableMetrics(config); + OkHttpClient watchHttpClient = mock(OkHttpClient.class); + PlatformInjector platformInjector = mock(PlatformInjector.class); + when(platformInjector.inject()).thenReturn(true); + when(platformInjector.getServerSocketAddress()) + .thenReturn(new InetSocketAddress("127.0.0.1", 25565)); + Tunneler tunneler = mock(Tunneler.class); + VerifiedBedrockIdentityRegistry registry = new VerifiedBedrockIdentityRegistry(); + BedrockAdmissionCoordinator coordinator = new BedrockAdmissionCoordinator(registry); + ConnectLogger logger = mock(ConnectLogger.class); + + Injector parent = Guice.createInjector(testModule( + config, watchHttpClient, platformInjector, tunneler, coordinator, logger)); + ConnectPlatform platform = new ConnectPlatform( + mock(ConnectApi.class), platformInjector, logger, parent, coordinator); + setConfig(platform, config); + + // Cycle 1: enable, watch connects, proposal is accepted. + platform.enable(new WatcherModule()); + WebSocketListener listener1 = captureListener(watchHttpClient, 1); + Session session1 = session("session-1"); + listener1.onMessage(mock(WebSocket.class), bytes(session1)); + verify(tunneler).prepare(session1); + + // Reload teardown: watcher stopped, parent coordinator closed. + platform.disable(); + + // Cycle 2: the reload re-instantiates the platform on the SAME parent injector (new child + // injector, same parent coordinator singleton, which disable() closed). + ConnectPlatform reloaded = new ConnectPlatform( + mock(ConnectApi.class), platformInjector, logger, parent, coordinator); + setConfig(reloaded, config); + reloaded.enable(new WatcherModule()); + WebSocketListener listener2 = captureListener(watchHttpClient, 2); + Session session2 = session("session-2"); + // RED before the fix: proposal() throws ISE (closed coordinator) and the watch would die. + assertDoesNotThrow( + () -> listener2.onMessage(mock(WebSocket.class), bytes(session2)), + "a proposal after reload must be accepted, not fail with 'coordinator is closed'"); + verify(tunneler, times(2)).prepare(any(Session.class)); + + reloaded.disable(); // stop the cycle-2 watcher's scheduler and coordinator executor + } + + private static Module testModule( + ConnectConfig config, + OkHttpClient watchHttpClient, + PlatformInjector platformInjector, + Tunneler tunneler, + BedrockAdmissionCoordinator coordinator, + ConnectLogger logger) { + return new AbstractModule() { + @Override + protected void configure() { + bind(ConnectConfig.class).toInstance(config); + bind(ConnectApi.class).toInstance(mock(ConnectApi.class)); + bind(SimpleConnectApi.class).toInstance(new SimpleConnectApi(logger)); + bind(PlatformInjector.class).toInstance(platformInjector); + bind(ConnectLogger.class).toInstance(logger); + bind(OkHttpClient.class) + .annotatedWith(Names.named("watchHttpClient")) + .toInstance(watchHttpClient); + bind(BedrockIdentityReadiness.class).toInstance(new BedrockIdentityReadiness( + config, new BedrockIdentityKeyProvider(config, new OkHttpClient()))); + bind(BedrockPrincipalReadiness.class).toInstance(new BedrockPrincipalReadiness(config)); + bind(BedrockAdmissionCoordinator.class).toInstance(coordinator); + bind(PlatformUtils.class).toInstance(mock(PlatformUtils.class)); + bind(String.class) + .annotatedWith(Names.named("platformName")) + .toInstance("test"); + bind(UpdateChecker.class).toInstance(mock(UpdateChecker.class)); + bind(Libp2pEndpoint.class).toInstance(mock(Libp2pEndpoint.class)); + bind(Tunneler.class).toInstance(tunneler); + bind(CommonPlatformInjector.class).toInstance(mock(CommonPlatformInjector.class)); + } + }; + } + + private static WebSocketListener captureListener(OkHttpClient httpClient, int invocation) + throws Exception { + ArgumentCaptor listener = ArgumentCaptor.forClass(WebSocketListener.class); + verify(httpClient, times(invocation)).newWebSocket(any(Request.class), listener.capture()); + return listener.getValue(); + } + + private static ByteString bytes(Session session) { + return ByteString.of(WatchResponse.newBuilder().setSession(session).build().toByteArray()); + } + + private static Session session(String id) { + return Session.newBuilder() + .setId(id) + .setTunnelServiceAddr("wss://tunnel.example") + .setPlayer(Player.newBuilder() + .setAddr("127.0.0.1") + .setProfile(GameProfile.newBuilder() + .setId("00000000-0000-0000-0000-000000000001") + .setName("Player"))) + .build(); + } + + private static void setConfig(ConnectPlatform platform, ConnectConfig config) throws Exception { + Field field = ConnectPlatform.class.getDeclaredField("config"); + field.setAccessible(true); + field.set(platform, config); + } + + /** bstats must stay inert in tests: it would otherwise start a daemon scheduler thread. */ + private static void disableMetrics(ConnectConfig config) throws Exception { + Field metricsField = ConnectConfig.class.getDeclaredField("metrics"); + metricsField.setAccessible(true); + Object metrics = metricsField.get(config); + if (metrics == null) { + metrics = new ConnectConfig.MetricsConfig(); + metricsField.set(config, metrics); + } + Field disabled = metrics.getClass().getDeclaredField("disabled"); + disabled.setAccessible(true); + disabled.set(metrics, true); + } +} diff --git a/core/src/test/java/com/minekube/connect/watch/WatchClientTest.java b/core/src/test/java/com/minekube/connect/watch/WatchClientTest.java index e4ea81a99..8957b875c 100644 --- a/core/src/test/java/com/minekube/connect/watch/WatchClientTest.java +++ b/core/src/test/java/com/minekube/connect/watch/WatchClientTest.java @@ -1,7 +1,10 @@ package com.minekube.connect.watch; +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.mock; @@ -90,6 +93,64 @@ void privateWatchSessionUsesCoordinatorTokenAndSanitizedProposal() { } } + @Test + void closedCoordinatorRejectsProposalInsteadOfFailingTheWatch() throws Exception { + ConnectConfig config = new ConnectConfig(); + OkHttpClient httpClient = mock(OkHttpClient.class); + VerifiedBedrockIdentityRegistry registry = new VerifiedBedrockIdentityRegistry(); + BedrockAdmissionCoordinator coordinator = new BedrockAdmissionCoordinator(registry); + WatchClient client = new WatchClient( + httpClient, + config, + new BedrockIdentityReadiness( + config, + new BedrockIdentityKeyProvider(config, new OkHttpClient())), + coordinator); + java.util.concurrent.atomic.AtomicReference proposalRef = + new java.util.concurrent.atomic.AtomicReference<>(); + java.util.concurrent.atomic.AtomicReference errorRef = + new java.util.concurrent.atomic.AtomicReference<>(); + Watcher watcher = new Watcher() { + @Override public void onOpen(WatchBootstrap bootstrap) { } + @Override public void onProposal(SessionProposal proposal) { proposalRef.set(proposal); } + @Override public void onCompleted() { } + @Override public void onError(Throwable throwable) { errorRef.set(throwable); } + }; + + try { + coordinator.close(); // disable() closed the shared coordinator + client.watch(watcher); + ArgumentCaptor listener = + ArgumentCaptor.forClass(WebSocketListener.class); + verify(httpClient).newWebSocket(any(Request.class), listener.capture()); + WebSocket socket = mock(WebSocket.class); + Session raw = Session.newBuilder() + .setId("session-closed-coordinator") + .setPlayer(Player.newBuilder() + .setAddr("127.0.0.1") + .setProfile(GameProfile.newBuilder() + .setId("f912bf90-8349-565f-9dc0-9891923c0cc3") + .setName("Player"))) + .build(); + + // A proposal into a closed coordinator must be rejected over the wire, not allowed to + // escape and fail the WebSocket (which would put the watch into a reconnect loop). + assertDoesNotThrow(() -> listener.getValue().onMessage( + socket, ByteString.of(WatchResponse.newBuilder() + .setSession(raw).build().toByteArray()))); + + assertNull(proposalRef.get(), "a closed coordinator must not accept the proposal"); + assertNull(errorRef.get(), "a closed coordinator must not fail the watch stream"); + ArgumentCaptor sent = ArgumentCaptor.forClass(ByteString.class); + verify(socket).send(sent.capture()); + WatchRequest request = WatchRequest.parseFrom(sent.getValue().toByteArray()); + assertTrue(request.hasSessionRejection()); + assertEquals("session-closed-coordinator", request.getSessionRejection().getId()); + } finally { + coordinator.close(); + } + } + @Test void defaultDisabledConfigurationDoesNotAdvertiseBedrockIdentity() { OkHttpClient httpClient = mock(OkHttpClient.class);