Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -151,7 +151,9 @@ public CompanionApplication(CompanionLaunchConfiguration configuration, String t
try {
JDTHacks.init(configuration.paths().jdtCache());
runtimeIndexService = new RuntimeIndexService(lifecycleLock, this::installRuntimeSnapshot);
runtimeIndexService.addStatusListener(this::updateRuntimeIndexUi);
RuntimeIndexService index = runtimeIndexService;
index.statusChanged().subscribe(() -> updateRuntimeIndexUi(index.status()));
updateRuntimeIndexUi(index.status());
debuggerController = createDebuggerController();
debuggerController.addListener(new DebuggerSessionController.Listener() {
private Throwable lastFailure;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -56,12 +56,9 @@ public ConfigChanges(GameLocation location, ChangeRecord record) {
this.location = Objects.requireNonNull(location, "location");
this.workspace = location.workspace();
this.record = Objects.requireNonNull(record, "record");
location.addListener(change -> {
switch (change) {
case PROCESS -> gameProcess(location.process());
case DISCONNECTED -> gameDisconnected();
default -> { }
}
location.processChanged().subscribe(() -> gameProcess(location.process()));
location.connectionChanged().subscribe(() -> {
if (location.connection() == null) gameDisconnected();
});
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,12 @@ public WorldReading(GameLocation location, GamePacks packs) {
() -> Workers.later(SETTLE_MILLIS, this.strand, this::follow));
// Companion changed the datapacks, also in the level.dat of a world the game does not hold.
this.stopFollowingDatapacks = packs.changed(ChangeRecord.PackSide.DATA).subscribe(() -> this.strand.execute(this::follow));
this.stopFollowingGame = location.addListener(change -> this.strand.execute(this::follow));
Runnable stopConnection = location.connectionChanged().subscribe(() -> this.strand.execute(this::follow));
Runnable stopPlaying = location.playingChanged().subscribe(() -> this.strand.execute(this::follow));
this.stopFollowingGame = () -> {
stopConnection.run();
stopPlaying.run();
};
// Followed after the triggers are, so none is missed in between; on the strand, as every later look.
CompletableFuture.runAsync(this::follow, this.strand).join();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -74,8 +74,8 @@ public ChangePipeline(GameLocation location, ChangeRecord record, Executor write
this.record = Objects.requireNonNull(record, "record");
this.writes = Objects.requireNonNull(writes, "writes");
this.reloads = new Reloads(location);
location.addListener(change -> {
if (change == GameLocation.Change.DISCONNECTED) gameDisconnected();
location.connectionChanged().subscribe(() -> {
if (location.connection() == null) gameDisconnected();
});
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -119,10 +119,10 @@ void sendIfIdle() {

public Reloads(GameLocation location) {
this.location = Objects.requireNonNull(location, "location");
location.addListener(change -> {
if (change == GameLocation.Change.DISCONNECTED) gameDisconnected();
else if (change == GameLocation.Change.PLAYING) leftWorld();
location.connectionChanged().subscribe(() -> {
if (location.connection() == null) gameDisconnected();
});
location.playingChanged().subscribe(this::leftWorld);
}

/** A write is queued whose reload the next reloads wait for, so writes in quick succession take one reload. */
Expand Down
Original file line number Diff line number Diff line change
@@ -1,16 +1,14 @@
package com.github.minecraft_ta.totalDebugCompanion.game;

import com.github.minecraft_ta.totalDebugCompanion.util.Signal;
import com.github.minecraft_ta.totalDebugCompanion.catalog.Worlds;
import com.github.minecraft_ta.totaldebug.protocol.message.PlayingPayload;
import com.github.minecraft_ta.totaldebug.storage.GameLock;
import com.github.minecraft_ta.totaldebug.storage.InstancePaths;
import com.github.tth05.scnet.message.AbstractMessage;

import java.nio.file.Path;
import java.util.List;
import java.util.Objects;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.function.Consumer;

/**
* Where a project's game is, as {@code docs/GAME_LOCATION.md} describes: closed, running without a connection, or
Expand All @@ -25,9 +23,6 @@ public interface Connection {
boolean send(AbstractMessage message);
}

/** What changed, told to listeners on the thread that changed it. */
public enum Change { CONNECTED, DISCONNECTED, PROCESS, PLAYING }

/** The instance's files that tell whether a game runs and whether a world is held. Blocking. */
public interface Files {
/** Whether a game holds the instance's game lock; a lock that cannot be checked counts as held. */
Expand Down Expand Up @@ -60,7 +55,9 @@ record Link(Connection connection, long process, PlayingPayload playing) {

private final Path workspace;
private final Files files;
private final List<Consumer<Change>> listeners = new CopyOnWriteArrayList<>();
private final Signal connectionChanged = new Signal();
private final Signal processChanged = new Signal();
private final Signal playingChanged = new Signal();
private volatile Link link;
/**
* What the game told on the current connection, kept from its first message: Companion may take the connection as
Expand All @@ -85,8 +82,8 @@ public Path workspace() {
}

/**
* The game's connection is established, with what the game has told on it so far; listeners hear of that as well, as
* if it was told now.
* The game's connection is established, with what the game has told on it so far; its followers hear of that as well,
* as if it was told now.
*/
public void connected(Connection connection) {
Objects.requireNonNull(connection, "connection");
Expand All @@ -95,24 +92,24 @@ public void connected(Connection connection) {
established = new Link(connection, this.toldProcess, this.toldPlaying);
this.link = established;
}
tell(Change.CONNECTED);
if (established.process() != 0) tell(Change.PROCESS);
if (established.playing() != null) tell(Change.PLAYING);
this.connectionChanged.fire();
if (established.process() != 0) this.processChanged.fire();
if (established.playing() != null) this.playingChanged.fire();
}

/** The game's process, as it announced it on the current connection; listeners hear of it when it changed. */
/** The game's process, as it announced it on the current connection; followers hear of it when it changed. */
public void process(long process) {
synchronized (this) {
this.toldProcess = process;
Link current = this.link;
if (current == null || current.process() == process) return;
this.link = new Link(current.connection(), process, current.playing());
}
tell(Change.PROCESS);
this.processChanged.fire();
}

/**
* What the game plays, as it told on the current connection; listeners hear of it when it changed. The game tells it
* What the game plays, as it told on the current connection; followers hear of it when it changed. The game tells it
* again after each handshake, which may repeat what {@link #connected} already announced.
*/
public void playing(PlayingPayload playing) {
Expand All @@ -123,7 +120,7 @@ public void playing(PlayingPayload playing) {
if (current == null || playing.equals(current.playing())) return;
this.link = new Link(current.connection(), current.process(), playing);
}
tell(Change.PLAYING);
this.playingChanged.fire();
}

/** The game's connection ended; what it told on it no longer holds. */
Expand All @@ -134,7 +131,7 @@ public void disconnected() {
if (this.link == null) return;
this.link = null;
}
tell(Change.DISCONNECTED);
this.connectionChanged.fire();
}

/** The connection to the game, or null without one. Not blocking. */
Expand Down Expand Up @@ -165,13 +162,21 @@ public GameState read() {
return new GameState(this.workspace, this.files, new GameState.Game.Unconnected());
}

/** Tells {@code listener} of each change, on the thread that made it; returns what removes it. */
public Runnable addListener(Consumer<Change> listener) {
this.listeners.add(Objects.requireNonNull(listener, "listener"));
return () -> this.listeners.remove(listener);
/**
* Fires when the game connected or disconnected, on the thread that saw it; {@link #connection()} tells which. The
* followers of one change hear of it before anything else changes the location.
*/
public Signal connectionChanged() {
return this.connectionChanged;
}

/** Fires when the connected game's process became known or changed, on the thread that saw it. */
public Signal processChanged() {
return this.processChanged;
}

private void tell(Change change) {
this.listeners.forEach(listener -> listener.accept(change));
/** Fires when what the connected game plays changed, also when a connection is established, on the thread that saw it. */
public Signal playingChanged() {
return this.playingChanged;
}
}
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
package com.github.minecraft_ta.totalDebugCompanion.inspection;

import com.github.minecraft_ta.totalDebugCompanion.util.Signal;
import com.github.minecraft_ta.totalDebugCompanion.catalog.CatalogIndex;
import com.github.minecraft_ta.totalDebugCompanion.itemrender.ItemModelId;
import com.github.minecraft_ta.totalDebugCompanion.itemrender.ItemRenderBackend;
Expand All @@ -21,7 +22,6 @@
import java.util.Objects;
import java.util.Optional;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.Executor;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
Expand Down Expand Up @@ -54,7 +54,9 @@ private record Key(String model, Map<Integer, Integer> tints, int size) {
.daemon()
.name("Companion item icons")
.unstarted(task));
private final List<Runnable> listeners = new CopyOnWriteArrayList<>();
private final Signal changed = new Signal();
/** Held while an adoption decides whether it is the newest. */
private final Object adoption = new Object();
private final Map<Key, Optional<BufferedImage>> cache = new LinkedHashMap<>(64, 0.75f, true) {
@Override
protected boolean removeEldestEntry(Map.Entry<Key, Optional<BufferedImage>> eldest) {
Expand Down Expand Up @@ -97,20 +99,20 @@ CompletableFuture<Void> restore(Path directory, Executor reader) {

/** Starts an adoption, superseding every earlier one that has not finished. */
private long start() {
synchronized (this.listeners) {
synchronized (this.adoption) {
return ++this.adoptions;
}
}

/** Replaces the snapshot, unless a newer adoption started since {@code generation}. */
private void adopt(long generation, Snapshot next) {
synchronized (this.listeners) {
synchronized (this.adoption) {
if (this.closed || generation != this.adoptions || Objects.equals(next, this.snapshot)) {
return;
}
this.snapshot = next;
}
SwingUtilities.invokeLater(() -> this.listeners.forEach(Runnable::run));
SwingUtilities.invokeLater(this.changed::fire);
}

static Snapshot newestSnapshot(Path directory) {
Expand Down Expand Up @@ -165,10 +167,9 @@ Snapshot snapshot() {
return this.snapshot;
}

/** Registers a listener for new snapshots and returns its removal. */
public Runnable addListener(Runnable listener) {
this.listeners.add(Objects.requireNonNull(listener, "listener"));
return () -> this.listeners.remove(listener);
/** Fires on the Swing thread after a newer snapshot was adopted, so views draw their icons again. */
public Signal changed() {
return this.changed;
}

/** Renders {@code model} (for example {@code minecraft:item/furnace}); empty when it cannot be drawn. */
Expand Down Expand Up @@ -312,7 +313,6 @@ public void close() {
return;
}
this.closed = true;
this.listeners.clear();
this.worker.execute(this::closeBackend);
this.worker.shutdown();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -54,20 +54,20 @@ ChangeRecord.PackSide.RESOURCES, new Signal(),
public GamePacks(GameLocation location) {
this.location = Objects.requireNonNull(location, "location");
this.workspace = location.workspace();
location.addListener(change -> {
if (change == GameLocation.Change.DISCONNECTED) gameDisconnected();
else if (change == GameLocation.Change.PLAYING) {
// The datapacks the server named belong to the world it played; its next world's server is asked.
synchronized (this) {
if (!Objects.equals(this.datapacksFor, location.playing())) {
this.datapacks = null;
this.worldRefusal = "";
}
location.connectionChanged().subscribe(() -> {
if (location.connection() == null) gameDisconnected();
});
location.playingChanged().subscribe(() -> {
// The datapacks the server named belong to the world it played; its next world's server is asked.
synchronized (this) {
if (!Objects.equals(this.datapacksFor, location.playing())) {
this.datapacks = null;
this.worldRefusal = "";
}
askForDatapacks();
// Without datapacks named yet, the world's own files stand for them, and they are another world's now.
tell(ChangeRecord.PackSide.DATA);
}
askForDatapacks();
// Without datapacks named yet, the world's own files stand for them, and they are another world's now.
tell(ChangeRecord.PackSide.DATA);
});
}

Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
package com.github.minecraft_ta.totalDebugCompanion.runtime;

import com.github.minecraft_ta.totalDebugCompanion.util.Signal;
import com.github.minecraft_ta.totaldebug.storage.RuntimeInventory;
import com.github.minecraft_ta.totaldebug.storage.InstancePaths;
import com.github.minecraft_ta.totaldebug.storage.AtomicFiles;
Expand All @@ -24,7 +25,6 @@
import java.util.Objects;
import java.util.Set;
import java.util.zip.ZipFile;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.CancellationException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
Expand Down Expand Up @@ -114,7 +114,7 @@ record PreparedInput(IndexSource indexSource, RuntimeSnapshotBytecodeSource.Sour
private final Object lifecycleLock;
private final Consumer<ReadySnapshot> readyHandler;
private final Function<String, ClassIndex> indexLoader;
private final CopyOnWriteArrayList<Consumer<Status>> listeners = new CopyOnWriteArrayList<>();
private final Signal statusChanged = new Signal();
private volatile Status status = new Status(Phase.WAITING, "Waiting for runtime inventory", null);
private String activeInventoryId;
private Metrics activeMetrics;
Expand Down Expand Up @@ -166,10 +166,9 @@ public Status status() {
return this.status;
}

public void addStatusListener(Consumer<Status> listener) {
Consumer<Status> checked = Objects.requireNonNull(listener, "listener");
this.listeners.add(checked);
checked.accept(this.status);
/** Fires after {@link #status()} changed, on the thread that changed it, under the lifecycle lock. */
public Signal statusChanged() {
return this.statusChanged;
}

public void waiting(String detail) {
Expand All @@ -185,10 +184,6 @@ public void waiting(String detail) {
}
}

public void removeStatusListener(Consumer<Status> listener) {
this.listeners.remove(listener);
}

public void restore(Path dataDirectory) {
restore(dataDirectory, null);
}
Expand Down Expand Up @@ -608,9 +603,10 @@ private void update(Work work, Status status) {
}

private void update(Status replacement) {
this.status = replacement;
for (Consumer<Status> listener : this.listeners) {
listener.accept(replacement);
// Under the lock, so what the followers read of the status is this one.
synchronized (this.lifecycleLock) {
this.status = replacement;
this.statusChanged.fire();
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ public CatalogIcons(ItemIconService service, int size) {
this.size = size;
this.icons = new IconLoader<>(4_096, 8, item -> service.render(item.model(), item.tints(), size)
.thenApply(image -> image.<Icon>map(ImageIcon::new)));
this.removeListener = service.addListener(this.icons::clear);
this.removeListener = service.changed().subscribe(this.icons::clear);
}

public int size() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -112,12 +112,10 @@ final class ConfigPanel extends JPanel {
failure -> showFailure("Could not list the worlds' copies: " + failure.getMessage()))
// A server configuration is shown from the world the game has open first; it follows the game to another
// world, unless its text has unsaved changes, which stay with the file they were made in.
.follow(listener -> this.location.addListener(change -> {
if (change == GameLocation.Change.PLAYING) SwingUtilities.invokeLater(() -> {
PackCatalog.ConfigFile file = selectedFile();
if (file != null && file.type() == PackCatalog.ConfigType.SERVER && !this.textEditor.modified()) listener.run();
});
}));
.follow(listener -> this.location.playingChanged().subscribe(() -> SwingUtilities.invokeLater(() -> {
PackCatalog.ConfigFile file = selectedFile();
if (file != null && file.type() == PackCatalog.ConfigType.SERVER && !this.textEditor.modified()) listener.run();
})));
configureFiles();
this.content.add(toolbar(), BorderLayout.NORTH);
JPanel settings = new JPanel(new BorderLayout());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -108,7 +108,7 @@ public DefinitionDetails(SubjectRef.Definition subject, Services services, JComp
this.appearance = found.appearance();
this.appearancePreviews = found.previews();
showExtras();
}, failure -> { }).waitsWhileHidden(page).follow(services.icons()::addListener);
}, failure -> { }).waitsWhileHidden(page).follow(services.icons().changed()::subscribe);
this.resourceLoader = new PageLoader<>(this::prepareResources, list -> {
this.owned = list;
this.matched = matching(this.owned, this.subject.namespace(), resourceName());
Expand Down
Loading