diff --git a/companion/src/main/java/com/github/minecraft_ta/totalDebugCompanion/CompanionApplication.java b/companion/src/main/java/com/github/minecraft_ta/totalDebugCompanion/CompanionApplication.java index 56450a6e6..f5fd6fd36 100644 --- a/companion/src/main/java/com/github/minecraft_ta/totalDebugCompanion/CompanionApplication.java +++ b/companion/src/main/java/com/github/minecraft_ta/totalDebugCompanion/CompanionApplication.java @@ -33,10 +33,6 @@ import com.github.minecraft_ta.totalDebugCompanion.script.ScriptCompilationService; import com.github.minecraft_ta.totalDebugCompanion.script.ScriptExecutionService; import com.github.minecraft_ta.totaldebug.protocol.scnet.InspectSubjectMessage; -import com.github.minecraft_ta.totaldebug.protocol.scnet.ChangeResultMessage; -import com.github.minecraft_ta.totaldebug.protocol.scnet.DatapacksMessage; -import com.github.minecraft_ta.totaldebug.protocol.scnet.PackStackMessage; -import com.github.minecraft_ta.totaldebug.protocol.scnet.ReloadResultMessage; import com.github.minecraft_ta.totalDebugCompanion.inspection.ItemIconService; import com.github.minecraft_ta.totaldebug.protocol.scnet.StopScriptMessage; import com.github.minecraft_ta.totalDebugCompanion.mcp.CompanionMcpServer; @@ -94,6 +90,8 @@ public final class CompanionApplication implements AutoCloseable, ProjectControl private final CompanionLaunchConfiguration launchConfiguration; private final Object lifecycleLock = new Object(); private volatile ProjectScope current; + /** Removes the current project's routes of the game's messages; under the lifecycle lock. */ + private Runnable currentMessages = () -> { }; private final InstanceState emptyState = InstanceState.inMemory(); private final CodeInsightService codeInsightService = new CodeInsightService( () -> { throw new IllegalStateException("Runtime class index is not ready"); }, RuntimeSourceCatalog.empty()); @@ -216,37 +214,12 @@ else if (gameStatus == null || gameStatus.state() != ServiceStatus.State.FAILED) if (pendingLaunch != null) queueLaunch(pendingLaunch); } - @Override - public void changeResult(ChangeResultMessage message) { - ProjectScope scope = current; - if (scope != null) scope.pipeline().answered(message.payload()); - } - - @Override - public void packStack(PackStackMessage message) { - ProjectScope scope = current; - if (scope != null) scope.packs().named(message.payload()); - } - - @Override - public void datapacks(DatapacksMessage message) { - ProjectScope scope = current; - if (scope != null) scope.packs().datapacks(message.world(), message.payload(), message.refusal()); - } - @Override public void playing(PlayingMessage message) { - ProjectScope scope = current; - if (scope != null) scope.location().playing(message.payload()); + // The current project's game location takes it too, through its own route. requestServerScripts(message.payload()); } - @Override - public void reloadResult(ReloadResultMessage message) { - ProjectScope scope = current; - if (scope != null) scope.pipeline().reloads().answered(message.payload()); - } - @Override public void failed(String detail, ClientHelloMessage hello) { synchronized (lifecycleLock) { if (closed || switching) return; @@ -267,14 +240,7 @@ public void relayFailed(RelayFailedMessage message) { case CompanionProtocol.RUN_SCRIPT, CompanionProtocol.STOP_SCRIPT -> { if (executionRuns != null) executionRuns.relayFailed(message.correlation(), message.reason()); } - case CompanionProtocol.RELOAD -> { - ProjectScope scope = current; - if (scope != null) scope.pipeline().reloads().relayFailed(message.correlation(), message.reason()); - } - case CompanionProtocol.CHANGE -> { - ProjectScope scope = current; - if (scope != null) scope.pipeline().relayFailed(message.correlation(), message.reason()); - } + // A change's or reload's reaches the current project's pipeline through its own route. default -> { } } } @@ -293,6 +259,8 @@ public void debugTarget(DebugTargetMessage message) { handleDebugTarget(message); } }); + // The project reopened above came before the session: its messages reach it from now on. + synchronized (lifecycleLock) { makeCurrent(current); } scriptExecutions = new ScriptExecutionService(session, scriptCompiler, this::isConnected); executionRuns = new ExecutionRuns(session, scriptExecutions); editorRuns = new EditorScriptRunService(executionRuns, notifications); @@ -361,7 +329,7 @@ public void start() throws IOException { synchronized (lifecycleLock) { scope = current; if (scope != null) scope.retire(); - current = null; + makeCurrent(null); } if (scope != null) { if (scope.runtime() != null) CompanionClassIndex.clear(); @@ -725,7 +693,7 @@ private void switchProject(CompanionProfile requested) throws IOException { } synchronized (lifecycleLock) { if (old != null) old.retire(); - current = null; + makeCurrent(null); if (runtimeIndexService != null) runtimeIndexService.clear(); } // Retirement is terminal. Attempt every detach and install the prepared replacement even if @@ -742,7 +710,7 @@ private void switchProject(CompanionProfile requested) throws IOException { try { old.close(); } catch (IOException | RuntimeException failure) { reportCleanupFailure("Close retired project", failure); } } - synchronized (lifecycleLock) { current = replacement; } + synchronized (lifecycleLock) { makeCurrent(replacement); } installed = true; restoreCatalog(replacement); runCleanup("Restore debugger preferences", () -> restoreProjectState(replacement)); @@ -793,6 +761,17 @@ private void runCleanup(String description, Runnable action) { catch (RuntimeException failure) { reportCleanupFailure(description, failure); } } + /** + * Makes {@code scope} the current project, or none: the game's messages for a project reach only the current one's + * owners. Under the lifecycle lock. + */ + private void makeCurrent(ProjectScope scope) { + this.currentMessages.run(); + this.currentMessages = () -> { }; + current = scope; + if (scope != null && session != null) this.currentMessages = scope.listen(session); + } + private void reportCleanupFailure(String description, Exception failure) { System.getLogger(CompanionApplication.class.getName()).log(System.Logger.Level.WARNING, description + " failed", failure); @@ -801,7 +780,7 @@ private void reportCleanupFailure(String description, Exception failure) { private void activateProfile(CompanionProfile requested) throws IOException { validateProfile(requested); ProjectScope replacement = ProjectScope.open(lifecycleLock, requested); - synchronized (lifecycleLock) { current = replacement; } + synchronized (lifecycleLock) { makeCurrent(replacement); } restoreCatalog(replacement); restoreProjectState(replacement); if (runtimeIndexService != null) runtimeIndexService.restore(requested.dataDirectory(), requested.workspaceDirectory()); diff --git a/companion/src/main/java/com/github/minecraft_ta/totalDebugCompanion/project/ProjectScope.java b/companion/src/main/java/com/github/minecraft_ta/totalDebugCompanion/project/ProjectScope.java index 9fe18a276..262f0f347 100644 --- a/companion/src/main/java/com/github/minecraft_ta/totalDebugCompanion/project/ProjectScope.java +++ b/companion/src/main/java/com/github/minecraft_ta/totalDebugCompanion/project/ProjectScope.java @@ -1,5 +1,13 @@ package com.github.minecraft_ta.totalDebugCompanion.project; +import com.github.minecraft_ta.totaldebug.protocol.scnet.RelayFailedMessage; +import com.github.minecraft_ta.totaldebug.protocol.scnet.PlayingMessage; +import com.github.minecraft_ta.totaldebug.protocol.scnet.DatapacksMessage; +import com.github.minecraft_ta.totaldebug.protocol.scnet.PackStackMessage; +import com.github.minecraft_ta.totaldebug.protocol.scnet.ReloadResultMessage; +import com.github.minecraft_ta.totaldebug.protocol.scnet.ChangeResultMessage; +import com.github.minecraft_ta.totaldebug.protocol.CompanionProtocol; +import com.github.minecraft_ta.totalDebugCompanion.session.MessageRoutes; import com.github.minecraft_ta.totalDebugCompanion.catalog.ConfigLabels; import com.github.minecraft_ta.totalDebugCompanion.catalog.ConfigSettings; import com.github.minecraft_ta.totalDebugCompanion.catalog.KeyBindingLabels; @@ -115,6 +123,28 @@ public ProjectScope(Object lock, CompanionProfile profile, InstanceState state, this.packSelections = new PackSelections(this.resources); } + /** + * Takes the game's messages this project's owners handle, until the returned removal runs, as when another project + * becomes the current one: a project that is not current receives nothing. + */ + public Runnable listen(MessageRoutes routes) { + List removals = List.of( + routes.on(ChangeResultMessage.class, message -> this.pipeline.answered(message.payload())), + routes.on(ReloadResultMessage.class, message -> this.pipeline.reloads().answered(message.payload())), + routes.on(PackStackMessage.class, message -> this.packs.named(message.payload())), + routes.on(DatapacksMessage.class, message -> this.packs.datapacks(message.world(), message.payload(), message.refusal())), + routes.on(PlayingMessage.class, message -> this.location.playing(message.payload())), + // The refused message and its correlation name the request together. + routes.on(RelayFailedMessage.class, message -> { + switch (message.messageId()) { + case CompanionProtocol.CHANGE -> this.pipeline.relayFailed(message.correlation(), message.reason()); + case CompanionProtocol.RELOAD -> this.pipeline.reloads().relayFailed(message.correlation(), message.reason()); + default -> { } + } + })); + return () -> removals.forEach(Runnable::run); + } + public static ProjectScope open(Object lock, CompanionProfile profile) throws IOException { var sources = RuntimeSourceCatalog.empty(); try { sources = new RuntimeSourceCatalog(LocalModSources.discover(profile.workspaceDirectory())); } diff --git a/companion/src/main/java/com/github/minecraft_ta/totalDebugCompanion/session/CompanionSession.java b/companion/src/main/java/com/github/minecraft_ta/totalDebugCompanion/session/CompanionSession.java index 05be42083..eabcca2a2 100644 --- a/companion/src/main/java/com/github/minecraft_ta/totalDebugCompanion/session/CompanionSession.java +++ b/companion/src/main/java/com/github/minecraft_ta/totalDebugCompanion/session/CompanionSession.java @@ -5,6 +5,7 @@ import com.github.minecraft_ta.totaldebug.protocol.relay.RelayedMessages; import com.github.minecraft_ta.totaldebug.protocol.scnet.ExecutionResultMessage; import java.util.function.BiConsumer; +import java.util.function.Consumer; import com.github.minecraft_ta.totaldebug.protocol.scnet.FromServerMessage; import com.github.minecraft_ta.totaldebug.protocol.scnet.ProtocolBindings; import com.github.minecraft_ta.totaldebug.protocol.scnet.RelayFailedMessage; @@ -14,10 +15,6 @@ import com.github.minecraft_ta.totaldebug.protocol.scnet.FocusWindowMessage; import com.github.minecraft_ta.totaldebug.protocol.scnet.ReadyMessage; import com.github.minecraft_ta.totaldebug.protocol.scnet.InspectSubjectMessage; -import com.github.minecraft_ta.totaldebug.protocol.scnet.ChangeResultMessage; -import com.github.minecraft_ta.totaldebug.protocol.scnet.DatapacksMessage; -import com.github.minecraft_ta.totaldebug.protocol.scnet.PackStackMessage; -import com.github.minecraft_ta.totaldebug.protocol.scnet.ReloadResultMessage; import com.github.minecraft_ta.totaldebug.protocol.scnet.DebugTargetMessage; import com.github.minecraft_ta.totaldebug.protocol.scnet.ClientHelloMessage; import com.github.minecraft_ta.totaldebug.protocol.scnet.ServerHelloMessage; @@ -41,10 +38,11 @@ import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; -public final class CompanionSession implements AutoCloseable { +public final class CompanionSession implements AutoCloseable, MessageRoutes { private final List> serverResultListeners = new CopyOnWriteArrayList<>(); private enum State { WAITING_FOR_HELLO, @@ -63,16 +61,6 @@ public interface Listener { default void inspectSubject(InspectSubjectMessage message) { } - /** The game answered a change of values it keeps. */ - default void changeResult(ChangeResultMessage message) { } - - default void packStack(PackStackMessage message) { } - - /** The server of the world the game plays named its datapacks. */ - default void datapacks(DatapacksMessage message) { } - - default void reloadResult(ReloadResultMessage message) { } - default void focusWindow() { } default void connecting() { @@ -164,6 +152,21 @@ public void removeExecutionResultListener(BiConsumer Runnable on(Class type, Consumer handler) { + // A key of its own, so removing it never removes another registration of the same handler. + Object key = new Object(); + // A message the bus is handing out while the route is removed is not handed to it any more. + AtomicBoolean active = new AtomicBoolean(true); + this.server.getMessageBus().listenAlways(type, key, message -> { + if (active.get()) handler.accept(message); + }); + return () -> { + active.set(false); + this.server.getMessageBus().unregister(type, key); + }; + } + public void setProjectSelectionHandler(AttachmentHandler handler) { if (this.projectSelections != null) throw new IllegalStateException("Session is already published"); this.projectSelectionHandler = Objects.requireNonNull(handler); @@ -278,10 +281,6 @@ private void registerHandlers() { this.server.getMessageBus().listenAlways(PlayingMessage.class, this.listener::playing); this.server.getMessageBus().listenAlways(DebugTargetMessage.class, this.listener::debugTarget); this.server.getMessageBus().listenAlways(InspectSubjectMessage.class, this.listener::inspectSubject); - this.server.getMessageBus().listenAlways(ChangeResultMessage.class, this.listener::changeResult); - this.server.getMessageBus().listenAlways(PackStackMessage.class, this.listener::packStack); - this.server.getMessageBus().listenAlways(DatapacksMessage.class, this.listener::datapacks); - this.server.getMessageBus().listenAlways(ReloadResultMessage.class, this.listener::reloadResult); this.server.getMessageBus().listenAlways(FocusWindowMessage.class, message -> SwingUtilities.invokeLater(this.listener::focusWindow)); this.server.addConnectionListener(new IConnectionListener() { diff --git a/companion/src/main/java/com/github/minecraft_ta/totalDebugCompanion/session/MessageRoutes.java b/companion/src/main/java/com/github/minecraft_ta/totalDebugCompanion/session/MessageRoutes.java new file mode 100644 index 000000000..eeb51ec6c --- /dev/null +++ b/companion/src/main/java/com/github/minecraft_ta/totalDebugCompanion/session/MessageRoutes.java @@ -0,0 +1,17 @@ +package com.github.minecraft_ta.totalDebugCompanion.session; + +import com.github.tth05.scnet.message.AbstractMessage; + +import java.util.function.Consumer; + +/** + * Where the game's messages are delivered (docs/SYSTEMS.md, section 6): an owner registers the messages it handles, and + * removes them when it closes, as a project that is no longer the current one. + */ +public interface MessageRoutes { + /** + * Runs {@code handler} for each message of {@code type} from the game, also one the server sent through it, on the + * connection's thread in arrival order; returns what removes it. + */ + Runnable on(Class type, Consumer handler); +} diff --git a/companion/src/test/java/com/github/minecraft_ta/totalDebugCompanion/ApplicationNavigationTest.java b/companion/src/test/java/com/github/minecraft_ta/totalDebugCompanion/ApplicationNavigationTest.java index 336a0c4dc..741436ed9 100644 --- a/companion/src/test/java/com/github/minecraft_ta/totalDebugCompanion/ApplicationNavigationTest.java +++ b/companion/src/test/java/com/github/minecraft_ta/totalDebugCompanion/ApplicationNavigationTest.java @@ -1,5 +1,8 @@ package com.github.minecraft_ta.totalDebugCompanion; +import com.github.minecraft_ta.totaldebug.protocol.scnet.PackStackMessage; +import com.github.minecraft_ta.totaldebug.protocol.message.PackStackPayload; +import com.github.minecraft_ta.totaldebug.protocol.message.ClientPacksPayload; import com.github.minecraft_ta.totalDebugCompanion.ui.HtmlText; import com.github.minecraft_ta.totalDebugCompanion.testui.UiTest; import com.github.minecraft_ta.totalDebugCompanion.bytecode.RuntimeSnapshotBytecodeSource; @@ -19,6 +22,7 @@ import com.github.minecraft_ta.totalDebugCompanion.resource.LocalFileSource; import com.github.minecraft_ta.totalDebugCompanion.ui.views.MainWindow; import org.junit.jupiter.api.io.TempDir; +import java.util.List; import java.nio.file.Files; import java.net.Socket; import java.awt.Container; @@ -232,6 +236,36 @@ private static void advanceRuntimeFollowUpBeforePublication(CompanionApplication } } + @Test void theProjectReopenedAtStartupTakesTheGamesMessages() throws Exception { + var configuration = new CompanionLaunchConfiguration(directory); + var profile = CompanionProfile.forGame(Files.createDirectories(directory.resolve("game"))); + try (var first = new CompanionApplication(configuration, "test-token")) { + first.openProject(profile).get(3, TimeUnit.SECONDS); + } + var game = new Client(); + try (var app = new CompanionApplication(configuration, "test-token")) { + assertEquals(profile.id(), app.currentProject().id(), "Companion reopens the project it had open"); + app.session().bindAndPublish(configuration); + ProtocolBindings.registerMod(game.getMessageProcessor()); + var ready = new CompletableFuture(); + game.getMessageBus().listenAlways(ReadyMessage.class, message -> ready.complete(null)); + int port = CompanionSessionDescriptor.read(configuration.descriptorFile(), CompanionProtocol.VERSION).port(); + assertTrue(game.connect(new InetSocketAddress("127.0.0.1", port))); + game.getMessageProcessor().enqueueMessage(new ClientHelloMessage(CompanionProtocol.VERSION, "test-token", profile.id(), + profile.dataDirectory().toString(), profile.workspaceDirectory().toString())); + ready.get(3, TimeUnit.SECONDS); + + // The game names its resource packs once it connected. + var stack = new PackStackPayload(34, List.of(new PackStackPayload.Pack("file/Faithful", "Faithful", ""))); + game.getMessageProcessor().enqueueMessage(new PackStackMessage(new ClientPacksPayload(stack, 48))); + long until = System.nanoTime() + TimeUnit.SECONDS.toNanos(3); + while (app.currentScope().packs().resourcePacks() == null && System.nanoTime() < until) Thread.sleep(10); + assertEquals(stack, app.currentScope().packs().resourcePacks(), "the reopened project takes the game's messages"); + } finally { + game.close(); + } + } + @Test void reconnectPendingKeepsRunsUntilTheTransportActuallyDisconnects() throws Exception { var configuration = new CompanionLaunchConfiguration(directory); var game = new Client(); diff --git a/companion/src/test/java/com/github/minecraft_ta/totalDebugCompanion/catalog/KeyAssignmentsTest.java b/companion/src/test/java/com/github/minecraft_ta/totalDebugCompanion/catalog/KeyAssignmentsTest.java index 000bee49b..135a013b2 100644 --- a/companion/src/test/java/com/github/minecraft_ta/totalDebugCompanion/catalog/KeyAssignmentsTest.java +++ b/companion/src/test/java/com/github/minecraft_ta/totalDebugCompanion/catalog/KeyAssignmentsTest.java @@ -1,5 +1,6 @@ package com.github.minecraft_ta.totalDebugCompanion.catalog; +import java.nio.file.AccessDeniedException; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; @@ -77,11 +78,23 @@ void aPageThatReadBeforeTheOwnerStillHearsOfTheNextChange() throws Exception { } } - /** Writes {@code text} beside {@code file} and moves it over the file, as the game and Companion save it. */ + /** + * Writes {@code text} beside {@code file} and moves it over the file, as an editor may save it. Windows refuses the move + * while the file is being read, as by the watch itself, so it is tried again for a moment. + */ private static void replace(Path file, String text) throws Exception { Path staged = file.resolveSibling(file.getFileName() + ".tmp"); Files.writeString(staged, text); - Files.move(staged, file, StandardCopyOption.REPLACE_EXISTING, StandardCopyOption.ATOMIC_MOVE); + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(2); + while (true) { + try { + Files.move(staged, file, StandardCopyOption.REPLACE_EXISTING, StandardCopyOption.ATOMIC_MOVE); + return; + } catch (AccessDeniedException reading) { + if (System.nanoTime() > deadline) throw reading; + Thread.sleep(20); + } + } } private static void await(IntSupplier count, int expected) throws InterruptedException { diff --git a/companion/src/test/java/com/github/minecraft_ta/totalDebugCompanion/project/ProjectScopeMessagesTest.java b/companion/src/test/java/com/github/minecraft_ta/totalDebugCompanion/project/ProjectScopeMessagesTest.java new file mode 100644 index 000000000..a9d672563 --- /dev/null +++ b/companion/src/test/java/com/github/minecraft_ta/totalDebugCompanion/project/ProjectScopeMessagesTest.java @@ -0,0 +1,68 @@ +package com.github.minecraft_ta.totalDebugCompanion.project; + +import com.github.minecraft_ta.totalDebugCompanion.session.CompanionProfile; +import com.github.minecraft_ta.totalDebugCompanion.session.MessageRoutes; +import com.github.minecraft_ta.totaldebug.protocol.message.ClientPacksPayload; +import com.github.minecraft_ta.totaldebug.protocol.message.PackStackPayload; +import com.github.minecraft_ta.totaldebug.protocol.scnet.PackStackMessage; +import com.github.tth05.scnet.message.AbstractMessage; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.List; +import java.util.function.Consumer; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +/** A project's owners take the game's messages through their own routes, and only while the project is the current one. */ +class ProjectScopeMessagesTest { + @TempDir Path directory; + + @Test + void aProjectTakesItsMessagesUntilAnotherBecomesCurrent() throws Exception { + Routes routes = new Routes(); + ProjectScope scope = ProjectScope.open(new Object(), CompanionProfile.forGame(this.directory)); + try { + Runnable stop = scope.listen(routes); + PackStackPayload faithful = stack("file/Faithful"); + routes.deliver(new PackStackMessage(new ClientPacksPayload(faithful, 48))); + assertEquals(faithful, scope.packs().resourcePacks(), "the packs the game named reach the project's packs"); + + // Another project becomes the current one. + stop.run(); + routes.deliver(new PackStackMessage(new ClientPacksPayload(stack("file/Other"), 48))); + assertEquals(faithful, scope.packs().resourcePacks(), "a project that is not current takes nothing"); + assertEquals(0, routes.handlers.size(), "and leaves no route behind"); + } finally { + scope.retire(); + scope.close(); + } + } + + private static PackStackPayload stack(String pack) { + return new PackStackPayload(34, List.of(new PackStackPayload.Pack(pack, pack, ""))); + } + + /** Routes that deliver what the test hands them, as the connection delivers what the game sends. */ + private static final class Routes implements MessageRoutes { + private record Handler(Class type, Consumer handler) { + } + + final List handlers = new ArrayList<>(); + + @Override + public Runnable on(Class type, Consumer handler) { + Handler added = new Handler(type, message -> handler.accept(type.cast(message))); + this.handlers.add(added); + return () -> this.handlers.remove(added); + } + + void deliver(AbstractMessage message) { + for (Handler handler : List.copyOf(this.handlers)) { + if (handler.type().isInstance(message)) handler.handler().accept(message); + } + } + } +} diff --git a/companion/src/test/java/com/github/minecraft_ta/totalDebugCompanion/session/CompanionSessionRoutesTest.java b/companion/src/test/java/com/github/minecraft_ta/totalDebugCompanion/session/CompanionSessionRoutesTest.java new file mode 100644 index 000000000..dfcab251a --- /dev/null +++ b/companion/src/test/java/com/github/minecraft_ta/totalDebugCompanion/session/CompanionSessionRoutesTest.java @@ -0,0 +1,71 @@ +package com.github.minecraft_ta.totalDebugCompanion.session; + +import com.github.minecraft_ta.totaldebug.protocol.CompanionProtocol; +import com.github.minecraft_ta.totaldebug.protocol.scnet.ProtocolBindings; +import com.github.minecraft_ta.totaldebug.protocol.message.ClientPacksPayload; +import com.github.minecraft_ta.totaldebug.protocol.message.PackStackPayload; +import com.github.minecraft_ta.totaldebug.protocol.scnet.ClientHelloMessage; +import com.github.minecraft_ta.totaldebug.protocol.scnet.PackStackMessage; +import com.github.minecraft_ta.totaldebug.protocol.scnet.ReadyMessage; +import com.github.minecraft_ta.totaldebug.storage.CompanionSessionDescriptor; +import com.github.tth05.scnet.Client; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.nio.file.Path; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** The game's messages reach the routes registered for them, and a removed route no longer, even mid-delivery. */ +class CompanionSessionRoutesTest { + private static final String TOKEN = "correct-token-value-1234567890abcdef"; + + @TempDir Path directory; + + @Test + void aRouteRemovedWhileAMessageIsHandedOutDoesNotTakeIt() throws Exception { + CompanionLaunchConfiguration configuration = new CompanionLaunchConfiguration(this.directory); + Client game = new Client(); + try (CompanionSession session = new CompanionSession(TOKEN)) { + CountDownLatch entered = new CountDownLatch(1); + CountDownLatch release = new CountDownLatch(1); + AtomicInteger late = new AtomicInteger(); + session.on(PackStackMessage.class, message -> { + entered.countDown(); + try { + release.await(5, TimeUnit.SECONDS); + } catch (InterruptedException interrupted) { + Thread.currentThread().interrupt(); + } + }); + Runnable removed = session.on(PackStackMessage.class, message -> late.incrementAndGet()); + session.bindAndPublish(configuration); + + ProtocolBindings.registerMod(game.getMessageProcessor()); + CompletableFuture ready = new CompletableFuture<>(); + game.getMessageBus().listenAlways(ReadyMessage.class, message -> ready.complete(null)); + int port = CompanionSessionDescriptor.read(configuration.descriptorFile(), CompanionProtocol.VERSION).port(); + assertTrue(game.connect(CompanionSession.sessionAddress(port))); + game.getMessageProcessor().enqueueMessage(new ClientHelloMessage(CompanionProtocol.VERSION, TOKEN, "profile", + this.directory.toString(), this.directory.toString())); + ready.get(5, TimeUnit.SECONDS); + + game.getMessageProcessor().enqueueMessage(new PackStackMessage(new ClientPacksPayload( + new PackStackPayload(34, List.of(new PackStackPayload.Pack("file/Faithful", "Faithful", ""))), 48))); + assertTrue(entered.await(5, TimeUnit.SECONDS), "the first route takes the message"); + // As the project the second route belongs to stops being the current one while the message is handed out. + removed.run(); + release.countDown(); + Thread.sleep(200); + assertEquals(0, late.get(), "a route removed meanwhile does not take it"); + } finally { + game.close(); + } + } +} diff --git a/mod/src/main/java/com/github/minecraft_ta/totaldebug/client/TotalDebugClient.java b/mod/src/main/java/com/github/minecraft_ta/totaldebug/client/TotalDebugClient.java index 5bca3a54a..0cf0b5ac1 100644 --- a/mod/src/main/java/com/github/minecraft_ta/totaldebug/client/TotalDebugClient.java +++ b/mod/src/main/java/com/github/minecraft_ta/totaldebug/client/TotalDebugClient.java @@ -1,5 +1,10 @@ package com.github.minecraft_ta.totaldebug.client; +import com.github.minecraft_ta.totaldebug.protocol.scnet.ToServerMessage; +import com.github.minecraft_ta.totaldebug.protocol.scnet.StopScriptMessage; +import com.github.minecraft_ta.totaldebug.protocol.scnet.RunScriptMessage; +import com.github.minecraft_ta.totaldebug.protocol.scnet.ReloadMessage; +import com.github.minecraft_ta.totaldebug.protocol.scnet.ChangeMessage; import com.github.minecraft_ta.totaldebug.change.ChangeTable; import com.github.minecraft_ta.totaldebug.client.catalog.KeyBindingEdits; import com.github.minecraft_ta.totaldebug.client.catalog.PackCatalogCapture; @@ -99,9 +104,9 @@ private TotalDebugClient(Path gameDirectory) { }); ChangeTable changes = new ChangeTable(Map.of(KeyBindingEdits.CATEGORY, new KeyBindingEdits(), ResourcePackEdits.CATEGORY, new ResourcePackEdits()), Minecraft.getInstance()); - companionApp.setChangeHandler((message, companion) -> Minecraft.getInstance().execute(() -> changes.apply(message.payload()) + companionApp.on(ChangeMessage.class, (message, companion) -> Minecraft.getInstance().execute(() -> changes.apply(message.payload()) .thenAccept(result -> companionApp.sendChangeResult(companion, new ChangeResultMessage(result))))); - companionApp.setReloadHandler((message, companion) -> Minecraft.getInstance().execute(() -> + companionApp.on(ReloadMessage.class, (message, companion) -> Minecraft.getInstance().execute(() -> ResourceReloads.reload(message.payload(), result -> answerReload(companion, result)))); this.codeView = new CodeViewOperation(new CodeViewOperation.Actions() { @Override @@ -127,10 +132,11 @@ public void focusCompanion() { this.keptStacks); ClientRelay relay = new ClientRelay(companionApp, () -> this.playing.identity()); this.relay = relay; - companionApp.setToServerHandler((message, companion) -> Minecraft.getInstance().execute(() -> relay.toServer(companion, message))); + companionApp.on(ToServerMessage.class, (message, companion) -> Minecraft.getInstance().execute(() -> + relay.toServer(companion, message.payload()))); TotalDebug.get().network().setCompanionReceiver(relay::fromServer); - companionApp.setScriptRequestHandler(this.scripts::handleRunRequest); - companionApp.setStopScriptHandler(this.scripts::stopScript); + companionApp.on(RunScriptMessage.class, (message, companion) -> this.scripts.handleRunRequest(message)); + companionApp.on(StopScriptMessage.class, (message, companion) -> this.scripts.stopScript(message.scriptId())); companionApp.setSessionClosedHandler(companion -> { this.scripts.close(); if (companion != 0) Minecraft.getInstance().execute(() -> relay.companionLeft(companion)); diff --git a/mod/src/main/java/com/github/minecraft_ta/totaldebug/client/companion/CompanionAppClient.java b/mod/src/main/java/com/github/minecraft_ta/totaldebug/client/companion/CompanionAppClient.java index 67357a245..452c073f5 100644 --- a/mod/src/main/java/com/github/minecraft_ta/totaldebug/client/companion/CompanionAppClient.java +++ b/mod/src/main/java/com/github/minecraft_ta/totaldebug/client/companion/CompanionAppClient.java @@ -62,6 +62,7 @@ import java.util.Map; import java.util.Objects; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -118,21 +119,29 @@ private static final class Connection { private volatile CompanionDiscovery discovery; private volatile BooleanSupplier companionEnabled = () -> true; - private volatile Consumer scriptRequestHandler = message -> TotalDebug.LOGGER.warn( - "Ignoring companion script request {} because no handler is installed", - message.scriptId() - ); - private volatile IntConsumer stopScriptHandler = scriptId -> TotalDebug.LOGGER.warn( - "Ignoring companion stop-script request {} because no handler is installed", - scriptId - ); + /** + * What handles each of Companion's requests, with the number of the connection that sent it: the answers below until + * the game installs its own ({@link #on}). + */ + private final Map, ObjIntConsumer> requests = new ConcurrentHashMap<>(); + + { + on(RunScriptMessage.class, (message, companion) -> TotalDebug.LOGGER.warn( + "Ignoring companion script request {} because no handler is installed", message.scriptId())); + on(StopScriptMessage.class, (message, companion) -> TotalDebug.LOGGER.warn( + "Ignoring companion stop-script request {} because no handler is installed", message.scriptId())); + on(ChangeMessage.class, (message, companion) -> sendChangeResult(companion, new ChangeResultMessage( + ChangeResultPayload.refused(message.payload().requestId(), "The game is not ready to change anything yet")))); + on(ReloadMessage.class, (message, companion) -> sendReloadResult(companion, new ReloadResultMessage( + new ReloadResultPayload(message.payload().requestId(), 0, List.of(), "The game is not ready to reload resources yet")))); + on(ToServerMessage.class, (message, companion) -> sendRelayFailed(companion, message.payload(), + "The game's relay to the server is not ready")); + on(RetryRuntimeInventoryMessage.class, (message, companion) -> startRuntimeInventoryPreparation(true)); + } + private volatile IntConsumer sessionClosedHandler = number -> { }; private volatile Runnable sessionOpenedHandler = () -> { }; private volatile BiConsumer> packCatalogHandler = (inventoryId, modules) -> { }; - private volatile ObjIntConsumer changeHandler = (message, companion) -> sendChangeResult(companion, - new ChangeResultMessage(ChangeResultPayload.refused(message.payload().requestId(), "The game is not ready to change anything yet"))); - private volatile ObjIntConsumer reloadHandler = (message, companion) -> sendReloadResult(companion, - new ReloadResultMessage(new ReloadResultPayload(message.payload().requestId(), 0, List.of(), "The game is not ready to reload resources yet"))); private volatile RuntimeInventoryPublisher.PublishedInventory publishedInventory; private volatile Consumer progressListener = progress -> { }; private volatile boolean closing; @@ -243,12 +252,13 @@ private synchronized CompanionDiscovery.Result tryAutomaticConnection() { } } - public void setScriptRequestHandler(Consumer handler) { - this.scriptRequestHandler = Objects.requireNonNull(handler, "handler"); - } - - public void setStopScriptHandler(IntConsumer handler) { - this.stopScriptHandler = Objects.requireNonNull(handler, "handler"); + /** + * Handles Companion's requests of {@code type}, with the number of the connection that sent them, on the connection + * thread, replacing what handled them before. A request that comes before its connection's handshake completed fails + * that session instead, whatever its type (docs/SYSTEMS.md, section 6). + */ + public void on(Class type, ObjIntConsumer handler) { + this.requests.put(Objects.requireNonNull(type, "type"), Objects.requireNonNull(handler, "handler")); } /** Runs once for each connection that ended, with its number, or 0 when it never authenticated. */ @@ -278,24 +288,11 @@ public void sendPreparedFile(PreparedFilePayload file) { } } - /** - * Receives Companion's changes of values the game keeps, such as key bindings, with its connection's number; runs on - * the connection thread. - */ - public void setChangeHandler(ObjIntConsumer handler) { - this.changeHandler = Objects.requireNonNull(handler, "handler"); - } - /** Answers a change of Companion connection {@code companion}, and no connection opened since. */ public void sendChangeResult(int companion, ChangeResultMessage message) { send(companion, message); } - /** Receives Companion's requests to reload resources, with its connection's number; runs on the connection thread. */ - public void setReloadHandler(ObjIntConsumer handler) { - this.reloadHandler = Objects.requireNonNull(handler, "handler"); - } - /** Answers a reload of Companion connection {@code companion}, and no connection opened since. */ public void sendReloadResult(int companion, ReloadResultMessage message) { send(companion, message); @@ -313,14 +310,6 @@ public void setProgressListener(Consumer listener) { this.progressListener = Objects.requireNonNull(listener, "listener"); } - private volatile ObjIntConsumer toServerHandler = (message, companion) -> sendRelayFailed( - companion, message, "The game's relay to the server is not ready"); - - /** Receives the messages Companion addressed to the server with its connection's number; runs on the connection thread. */ - public void setToServerHandler(ObjIntConsumer handler) { - this.toServerHandler = Objects.requireNonNull(handler, "handler"); - } - /** Hands Companion connection {@code companion} a message the server sent it, unread. */ public void sendFromServer(int companion, RelayedMessage message) { send(companion, new FromServerMessage(message)); @@ -378,6 +367,31 @@ private void transferForeground(Runnable beforeTransfer, Runnable sendRequest) t this.foregroundHandoff.transfer(descriptor.processId(), beforeTransfer, sendRequest); } + /** + * Hands each request of {@code type} on {@code attempt} to what handles it now; one before the handshake completed + * fails the session, whatever its type. + */ + @SuppressWarnings("unchecked") + private void listenAuthenticated(Connection attempt, Class type) { + attempt.client.getMessageBus().listenAlways(type, message -> { + if (!attempt.authenticated()) { + failSession(attempt, "Companion sent " + type.getSimpleName() + " before authentication", null); + return; + } + if (message instanceof RunScriptMessage run && !compiledForThisInventory(run)) return; + ((ObjIntConsumer) this.requests.get(type)).accept(message, attempt.number); + }); + } + + /** Whether {@code run} was compiled against the runtime inventory the game has; answers it where it was not. */ + private boolean compiledForThisInventory(RunScriptMessage run) { + var inventory = this.runtimeInventoryState; + if (inventory.state() == PreparedFilePayload.State.READY && inventory.inventoryId().equals(run.inventoryId())) return true; + sendExecutionResult(run.scriptId(), ExecutionResult.fromStatus(ExecutionStatus.COMPILATION_FAILED, + "The script was compiled against a different runtime inventory. Wait for Companion to load the current index.")); + return false; + } + private void registerProtocol(Connection attempt) { var transport = attempt.client; transport.setMessageBus(new DefaultMessageBus() { @@ -389,7 +403,6 @@ private void registerProtocol(Connection attempt) { transport.getMessageProcessor().setMaxStringLength(DefaultMessageProcessor.DEFAULT_MAX_STRING_LENGTH); ProtocolBindings.registerMod(transport.getMessageProcessor()); transport.getMessageBus().listenAlways(ServerHelloMessage.class, message -> handleServerHello(attempt, message)); - transport.getMessageBus().listenAlways(RetryRuntimeInventoryMessage.class, message -> startRuntimeInventoryPreparation(true)); transport.getMessageBus().listenAlways(ReadyMessage.class, message -> { if (!attempt.authenticated()) { failSession(attempt, "Companion sent Ready before the session handshake completed", null); @@ -397,47 +410,7 @@ private void registerProtocol(Connection attempt) { } attempt.ready.complete(null); }); - transport.getMessageBus().listenAlways(ToServerMessage.class, message -> { - if (!attempt.authenticated()) { - failSession(attempt, "Companion sent a message for the server before authentication", null); - return; - } - this.toServerHandler.accept(message.payload(), attempt.number); - }); - transport.getMessageBus().listenAlways(RunScriptMessage.class, message -> { - if (!attempt.authenticated()) { - failSession(attempt, "Companion sent a script request before authentication", null); - return; - } - var inventory = this.runtimeInventoryState; - if (inventory.state() != PreparedFilePayload.State.READY || !inventory.inventoryId().equals(message.inventoryId())) { - sendExecutionResult(message.scriptId(), ExecutionResult.fromStatus(ExecutionStatus.COMPILATION_FAILED, - "The script was compiled against a different runtime inventory. Wait for Companion to load the current index.")); - return; - } - this.scriptRequestHandler.accept(message); - }); - transport.getMessageBus().listenAlways(ChangeMessage.class, message -> { - if (!attempt.authenticated()) { - failSession(attempt, "Companion sent a change before authentication", null); - return; - } - this.changeHandler.accept(message, attempt.number); - }); - transport.getMessageBus().listenAlways(ReloadMessage.class, message -> { - if (!attempt.authenticated()) { - failSession(attempt, "Companion sent a reload request before authentication", null); - return; - } - this.reloadHandler.accept(message, attempt.number); - }); - transport.getMessageBus().listenAlways(StopScriptMessage.class, message -> { - if (!attempt.authenticated()) { - failSession(attempt, "Companion sent a stop-script request before authentication", null); - return; - } - this.stopScriptHandler.accept(message.scriptId()); - }); + for (Class type : this.requests.keySet()) listenAuthenticated(attempt, type); transport.addConnectionListener(new IConnectionListener() { @Override public void onConnected() { if (closing || connection != attempt) { transport.close(); return; } diff --git a/mod/src/test/java/com/github/minecraft_ta/totaldebug/client/companion/CompanionAutomaticConnectionTest.java b/mod/src/test/java/com/github/minecraft_ta/totaldebug/client/companion/CompanionAutomaticConnectionTest.java index eab74620e..5c2fefcee 100644 --- a/mod/src/test/java/com/github/minecraft_ta/totaldebug/client/companion/CompanionAutomaticConnectionTest.java +++ b/mod/src/test/java/com/github/minecraft_ta/totaldebug/client/companion/CompanionAutomaticConnectionTest.java @@ -1,5 +1,6 @@ package com.github.minecraft_ta.totaldebug.client.companion; +import com.github.minecraft_ta.totaldebug.protocol.scnet.RetryRuntimeInventoryMessage; import com.github.minecraft_ta.totaldebug.protocol.CompanionProtocol; import com.github.minecraft_ta.totaldebug.protocol.ProjectSelectionRequest; import com.github.minecraft_ta.totaldebug.protocol.message.ChangePayload; @@ -287,7 +288,7 @@ void temporarilyUnreadableDescriptorRecoversWithoutAnotherPublication() throws E @Test void anAnswerGoesOnlyToTheConnectionThatAsked() throws Exception { var asked = new CompletableFuture(); try (var endpoint = new Endpoint(); var client = client()) { - client.setChangeHandler((message, companion) -> asked.complete(companion)); + client.on(ChangeMessage.class, (message, companion) -> asked.complete(companion)); endpoint.publish(profile()); client.startDiscovery(() -> true); await(client::isConnected); @@ -306,6 +307,30 @@ void temporarilyUnreadableDescriptorRecoversWithoutAnotherPublication() throws E } } + @Test void aRequestBeforeTheHandshakeFailsItsSessionWhateverItsType() throws Exception { + var handled = new AtomicInteger(); + try (var endpoint = new Endpoint(); var client = client()) { + client.on(ChangeMessage.class, (message, companion) -> handled.incrementAndGet()); + client.on(RetryRuntimeInventoryMessage.class, (message, companion) -> handled.incrementAndGet()); + endpoint.hold = true; + endpoint.publish(profile()); + client.startDiscovery(() -> true); + await(() -> endpoint.hellos.get() == 1); + + // Before Companion answered the hello, whatever it sends is refused, and the session starts again. + endpoint.server.getMessageProcessor().enqueueMessage(new RetryRuntimeInventoryMessage()); + await(() -> endpoint.hellos.get() == 2); + endpoint.server.getMessageProcessor().enqueueMessage(new ChangeMessage(new ChangePayload(1, + List.of(new ChangePayload.Edit("keyBinding", "key.jump", "key.keyboard.space", "key.keyboard.g"))))); + await(() -> endpoint.hellos.get() == 3); + assertEquals(0, handled.get(), "no request before the handshake reaches what handles it"); + + endpoint.hold = false; + endpoint.accept(); + await(client::isConnected); + } + } + private String profile() { return InstancePaths.profileId(game); } private static void await(BooleanSupplier condition) throws Exception {