diff --git a/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/CoreImplPlugin.java b/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/CoreImplPlugin.java index 62ed252..3d52edc 100644 --- a/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/CoreImplPlugin.java +++ b/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/CoreImplPlugin.java @@ -3,6 +3,7 @@ package de.kentoj.scrow.bukkit; import de.kentoj.scrow.bukkit.crown.CrownCommand; import de.kentoj.scrow.bukkit.economy.CoinsCommand; import de.kentoj.scrow.bukkit.friends.command.FriendCommand; +import de.kentoj.scrow.bukkit.server.PlayCommand; import org.bukkit.Bukkit; import org.bukkit.plugin.java.JavaPlugin; @@ -18,6 +19,7 @@ public class CoreImplPlugin extends JavaPlugin { ScrowAPI.getCommandManager().register(new FriendCommand().getRootNode()); ScrowAPI.getCommandManager().register(new CoinsCommand().getRootNode()); ScrowAPI.getCommandManager().register(new CrownCommand().getRootNode()); + ScrowAPI.getCommandManager().register(new PlayCommand().getRootNode()); } registerServer(); diff --git a/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/ScrowAPISurface.java b/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/ScrowAPISurface.java index fa3ce75..432c465 100644 --- a/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/ScrowAPISurface.java +++ b/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/ScrowAPISurface.java @@ -11,7 +11,7 @@ import de.kentoj.scrow.bukkit.friends.friendship.InMemoryFriendshipService; import de.kentoj.scrow.bukkit.friends.friendship.MongoFriendshipService; import de.kentoj.scrow.bukkit.friends.request.InMemoryFriendRequestService; import de.kentoj.scrow.bukkit.friends.request.MongoFriendRequestService; -import de.kentoj.scrowlib.servermanager.CrownServerManager; +import de.kentoj.scrowlib.servermanager.ServerManagerImpl; import io.nats.client.Nats; import io.nats.client.Options; import lombok.AccessLevel; @@ -64,7 +64,7 @@ public class ScrowAPISurface { ScrowAPI.setDatabase(mongoClient.getDatabase("scrow")); } - ScrowAPI.setServerManager(new CrownServerManager(ScrowAPI.getNats())); + ScrowAPI.setServerManager(new ServerManagerImpl(ScrowAPI.getNats())); ScrowAPI.setEconomyService(ScrowAPI.getDatabase() != null ? new MongoEconomyService(ScrowAPI.getDatabase()) : new InMemoryEconomyService()); diff --git a/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/crown/CrownDeployServerLiteral.java b/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/crown/CrownDeployServerLiteral.java index 571318b..60e3d15 100644 --- a/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/crown/CrownDeployServerLiteral.java +++ b/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/crown/CrownDeployServerLiteral.java @@ -28,8 +28,8 @@ public class CrownDeployServerLiteral implements CommandExecutor { - ctx.getSender().sendMessage("Deployed instance of " + template + " with handle " + server.getHandle()); - }); + .subscribe(server -> + ctx.getSender().sendMessage("Deployed instance of " + template + " with handle " + server.getHandle()) + ); } } diff --git a/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/crown/CrownDestroyServerLiteral.java b/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/crown/CrownDestroyServerLiteral.java index 0e96e75..f14c9a0 100644 --- a/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/crown/CrownDestroyServerLiteral.java +++ b/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/crown/CrownDestroyServerLiteral.java @@ -38,8 +38,8 @@ public class CrownDestroyServerLiteral implements CommandExecutor { - ctx.getSender().sendMessage("Destroyed server with handle " + handle); - }); + .subscribe(__ -> + ctx.getSender().sendMessage("Destroyed server with handle " + handle) + ); } } diff --git a/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/crown/CrownListServersLiteral.java b/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/crown/CrownListServersLiteral.java index 487a47a..feed9f9 100644 --- a/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/crown/CrownListServersLiteral.java +++ b/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/crown/CrownListServersLiteral.java @@ -28,8 +28,8 @@ public class CrownListServersLiteral implements CommandExecutor { - ctx.getPlayerSender().sendMessage(msg.toString()); - }); + .subscribe(msg -> + ctx.getPlayerSender().sendMessage(msg.toString()) + ); } } diff --git a/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/server/PlayCommand.java b/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/server/PlayCommand.java new file mode 100644 index 0000000..db73b27 --- /dev/null +++ b/core-bukkit-impl/src/main/java/de/kentoj/scrow/bukkit/server/PlayCommand.java @@ -0,0 +1,41 @@ +package de.kentoj.scrow.bukkit.server; + +import de.kentoj.kencommandapi.BukkitCommandContext; +import de.kentoj.kencommandapi.api.processing.CommandExecutor; +import de.kentoj.kencommandapi.api.structure.argument.CommandArgument; +import de.kentoj.kencommandapi.api.structure.argument.types.StringArgumentType; +import de.kentoj.kencommandapi.api.structure.node.CommandNode; +import de.kentoj.scrow.bukkit.ScrowAPI; +import de.kentoj.scrowlib.servermanager.ServerInstance; +import lombok.Getter; +import org.bukkit.command.CommandSender; + +public class PlayCommand implements CommandExecutor { + + @Getter + private final CommandNode rootNode; + private final CommandArgument gamemodeArg; + + public PlayCommand() { + this.rootNode = ScrowAPI.getCommandManager().createNode("play", "queue"); + this.gamemodeArg = ScrowAPI.getCommandManager().createArgument("gamemode", new StringArgumentType<>()); + this.rootNode.addArgument(this.gamemodeArg); + this.rootNode.setExecutor(this); + } + + @Override + public void execute(BukkitCommandContext ctx) { + var gamemode = ctx.getArg(gamemodeArg); + ScrowAPI.getServerManager().getRunningServers().filter( + inst -> inst.getTemplateName().equals(gamemode) && inst.getState().equals("UP") + ) + .single() + .map(ServerInstance::getHandle) + .flatMap(handle -> + ScrowAPI.getServerManager().sendPlayer(ctx.getPlayerSender().getUniqueId(), handle) + ) + .subscribe(null, ex -> { + ctx.getPlayerSender().sendMessage("failed sending you to " + gamemode + ": " + ex.getMessage()); + }); + } +} diff --git a/core-lib/build.gradle.kts b/core-lib/build.gradle.kts index f5d0a3f..1bbcd9c 100644 --- a/core-lib/build.gradle.kts +++ b/core-lib/build.gradle.kts @@ -8,7 +8,7 @@ plugins { } group = "de.kentoj.scrow" -version = "0.2" +version = "0.3" repositories { mavenCentral() diff --git a/core-lib/src/main/java/de/kentoj/scrowlib/messaging/NatsRepository.java b/core-lib/src/main/java/de/kentoj/scrowlib/messaging/NatsRepository.java index 7b33187..65d402f 100644 --- a/core-lib/src/main/java/de/kentoj/scrowlib/messaging/NatsRepository.java +++ b/core-lib/src/main/java/de/kentoj/scrowlib/messaging/NatsRepository.java @@ -1,8 +1,10 @@ package de.kentoj.scrowlib.messaging; +import de.kentoj.scrowlib.messaging.respose.SafeResult; import io.nats.client.Connection; import io.nats.client.Dispatcher; import io.nats.client.Message; +import org.bson.Document; import org.jetbrains.annotations.NotNull; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @@ -18,23 +20,34 @@ public class NatsRepository { this.dispatcher = nats.createDispatcher(); } - public Mono<@NotNull Void> publish(String subject, byte[] data) { + public Mono<@NotNull Void> publishSafeResult(String subject, SafeResult safeResult) { + return this.publish(subject, safeResult.toDocument()); + } + + public Mono<@NotNull Void> publish(String subject, Document document) { + return publishRaw(subject, DocumentRepository.toBytes(document)); + } + + public Mono<@NotNull Void> publishRaw(String subject, byte[] data) { return Mono.fromRunnable(() -> nats.publish(subject, data)); } public Flux<@NotNull Message> subscribe(String subject) { return Flux.create(sink -> { - var handler = dispatcher.subscribe(subject, msg -> { - try { - sink.next(msg); - } catch (Throwable e) { - sink.error(e); - } - }); + var handler = dispatcher.subscribe(subject, msg -> { + try { + sink.next(msg); + } catch (Throwable e) { + sink.error(e); + } + }); + sink.onCancel(handler::unsubscribe); + sink.onDispose(handler::unsubscribe); + }); + } - sink.onCancel(handler::unsubscribe); - sink.onDispose(handler::unsubscribe); - }); + public Mono<@NotNull Message> request(String subject, Document document) { + return Mono.fromFuture(nats.request(subject, DocumentRepository.toBytes(document))); } public Mono<@NotNull Message> request(String subject, byte[] data) { diff --git a/core-lib/src/main/java/de/kentoj/scrowlib/messaging/respose/SafeErrorResult.java b/core-lib/src/main/java/de/kentoj/scrowlib/messaging/respose/SafeErrorResult.java new file mode 100644 index 0000000..2de9027 --- /dev/null +++ b/core-lib/src/main/java/de/kentoj/scrowlib/messaging/respose/SafeErrorResult.java @@ -0,0 +1,34 @@ +package de.kentoj.scrowlib.messaging.respose; + +import lombok.AccessLevel; +import lombok.AllArgsConstructor; +import lombok.Getter; +import org.bson.Document; + +@Getter +@AllArgsConstructor(access = AccessLevel.PACKAGE) +public final class SafeErrorResult extends SafeResult { + private final String type; + private final String message; + + public SafeErrorResult(Throwable throwable) { + this.type = throwable.getClass().getCanonicalName(); + this.message = throwable.getMessage(); + } + + @Override + public Document unwrap() { + throw this.toException(); + } + + public RuntimeException toException() { + return new RuntimeException(message); + } + + @Override + public Document toDocument() { + return new Document("error", new Document() + .append("type", this.type) + .append("message", this.getMessage())); + } +} diff --git a/core-lib/src/main/java/de/kentoj/scrowlib/messaging/respose/SafeResult.java b/core-lib/src/main/java/de/kentoj/scrowlib/messaging/respose/SafeResult.java new file mode 100644 index 0000000..6988f62 --- /dev/null +++ b/core-lib/src/main/java/de/kentoj/scrowlib/messaging/respose/SafeResult.java @@ -0,0 +1,64 @@ +package de.kentoj.scrowlib.messaging.respose; + +import de.kentoj.scrowlib.messaging.DocumentRepository; +import io.nats.client.Message; +import org.bson.Document; + +public abstract sealed class SafeResult permits SafeErrorResult, SafeSuccessResult { + + /** + * @return the SafeResult as a Document for transit + */ + public abstract Document toDocument(); + + /** + * @return the data stored in the SafeResult + * @throws RuntimeException if the SafeResult represents an error + */ + public abstract Document unwrap(); + + public boolean isError() { + return this instanceof SafeErrorResult; + } + + public boolean isSuccess() { + return this instanceof SafeSuccessResult; + } + + public SafeErrorResult asError() { + return (SafeErrorResult) this; + } + + public SafeSuccessResult asSuccess() { + return (SafeSuccessResult) this; + } + + public static Document unwrap(Message message) { + return fromMessage(message).unwrap(); + } + + public static SafeResult wrap(Document document) { + return new SafeSuccessResult(document); + } + + public static SafeResult wrap(Throwable throwable) { + return new SafeErrorResult(throwable); + } + + public static SafeResult fromMessage(Message message) { + return fromDocument(DocumentRepository.fromMessage(message)); + } + + public static SafeResult fromDocument(Document document) { + if (document.containsKey("error")) { + var doc = document.get("error", Document.class); + return new SafeErrorResult( + doc.getString("type"), + doc.getString("message") + ); + } else if (document.containsKey("data")) { + return new SafeSuccessResult(document.get("data", Document.class)); + } + throw new IllegalArgumentException("document does not conform to SafeDocument format"); + } +} diff --git a/core-lib/src/main/java/de/kentoj/scrowlib/messaging/respose/SafeSuccessResult.java b/core-lib/src/main/java/de/kentoj/scrowlib/messaging/respose/SafeSuccessResult.java new file mode 100644 index 0000000..69242d1 --- /dev/null +++ b/core-lib/src/main/java/de/kentoj/scrowlib/messaging/respose/SafeSuccessResult.java @@ -0,0 +1,22 @@ +package de.kentoj.scrowlib.messaging.respose; + +import lombok.AccessLevel; +import lombok.AllArgsConstructor; +import lombok.Getter; +import org.bson.Document; + +@Getter +@AllArgsConstructor(access = AccessLevel.PACKAGE) +public final class SafeSuccessResult extends SafeResult { + private final Document data; + + @Override + public Document unwrap() { + return this.getData(); + } + + @Override + public Document toDocument() { + return new Document("data", this.getData()); + } +} diff --git a/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/ServerInstance.java b/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/ServerInstance.java index 2ee727f..0750e1b 100644 --- a/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/ServerInstance.java +++ b/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/ServerInstance.java @@ -1,16 +1,26 @@ package de.kentoj.scrowlib.servermanager; -import lombok.AllArgsConstructor; +import com.google.common.base.Preconditions; import lombok.Value; @Value -@AllArgsConstructor public class ServerInstance { + String templateName; String handle; int port; String state; - public ServerInstance(String handle, int port) { - this(handle, port, "STARTING"); + public ServerInstance(String templateName, String handle, int port, String state) { + Preconditions.checkNotNull(templateName); + Preconditions.checkNotNull(handle); + Preconditions.checkNotNull(state); + this.templateName = templateName; + this.handle = handle; + this.port = port; + this.state = state; + } + + public ServerInstance(String templateName, String handle, int port) { + this(templateName, handle, port, "STARTING"); } } diff --git a/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/ServerManager.java b/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/ServerManager.java index 6271544..f092bc4 100644 --- a/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/ServerManager.java +++ b/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/ServerManager.java @@ -4,6 +4,8 @@ import org.jetbrains.annotations.NotNull; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; +import java.util.UUID; + public interface ServerManager { /** @@ -24,4 +26,6 @@ public interface ServerManager { Mono<@NotNull ServerInstance> getInfo(String handle); Mono<@NotNull Void> setState(String handle, String state); + + Mono<@NotNull Void> sendPlayer(UUID playerId, String serverHandle); } diff --git a/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/CrownServerManager.java b/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/ServerManagerImpl.java similarity index 78% rename from core-lib/src/main/java/de/kentoj/scrowlib/servermanager/CrownServerManager.java rename to core-lib/src/main/java/de/kentoj/scrowlib/servermanager/ServerManagerImpl.java index fb0f150..37be894 100644 --- a/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/CrownServerManager.java +++ b/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/ServerManagerImpl.java @@ -3,17 +3,20 @@ package de.kentoj.scrowlib.servermanager; import com.google.common.base.Preconditions; import de.kentoj.scrowlib.messaging.DocumentRepository; import de.kentoj.scrowlib.messaging.NatsRepository; +import de.kentoj.scrowlib.messaging.respose.SafeResult; import io.nats.client.Connection; import org.bson.Document; import org.jetbrains.annotations.NotNull; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; -public class CrownServerManager implements ServerManager { +import java.util.UUID; + +public class ServerManagerImpl implements ServerManager { private final NatsRepository natsRepo; - public CrownServerManager(Connection nats) { + public ServerManagerImpl(Connection nats) { this.natsRepo = new NatsRepository(nats); } @@ -23,6 +26,7 @@ public class CrownServerManager implements ServerManager { return natsRepo.request("crown.deploy-server", payload) .map(DocumentRepository::fromMessage) .map(doc -> new ServerInstance( + template, doc.getString("handle"), doc.getInteger("port") )); @@ -38,6 +42,7 @@ public class CrownServerManager implements ServerManager { ) .flatMapIterable(l -> l) .map(doc -> new ServerInstance( + doc.getString("templateName"), doc.getString("handle"), doc.getInteger("port"), doc.getString("state") @@ -57,6 +62,7 @@ public class CrownServerManager implements ServerManager { return natsRepo.request("crown.get-info", payload) .map(DocumentRepository::fromMessage) .map(doc -> new ServerInstance( + doc.getString("templateName"), handle, doc.getInteger("port"), doc.getString("state") @@ -72,4 +78,14 @@ public class CrownServerManager implements ServerManager { return natsRepo.request("crown.set-state", payload) .then(); } + + @Override + public Mono<@NotNull Void> sendPlayer(UUID playerId, String serverHandle) { + var payload = DocumentRepository.toBytes(new Document() + .append("playerId", playerId.toString()) + .append("handle", serverHandle)); + return natsRepo.request("proxy.send-player", payload) + .map(SafeResult::unwrap) + .then(); + } } diff --git a/core-velocity/src/main/java/de/kentoj/scrow/corevelocity/CoreVelocity.java b/core-velocity/src/main/java/de/kentoj/scrow/corevelocity/CoreVelocity.java index 8519921..cae6960 100644 --- a/core-velocity/src/main/java/de/kentoj/scrow/corevelocity/CoreVelocity.java +++ b/core-velocity/src/main/java/de/kentoj/scrow/corevelocity/CoreVelocity.java @@ -1,11 +1,9 @@ package de.kentoj.scrow.corevelocity; import com.google.inject.Inject; -import com.velocitypowered.api.event.Subscribe; -import com.velocitypowered.api.event.proxy.ProxyInitializeEvent; import com.velocitypowered.api.plugin.Plugin; import com.velocitypowered.api.proxy.ProxyServer; -import de.kentoj.scrowlib.servermanager.CrownServerManager; +import de.kentoj.scrowlib.servermanager.ServerManagerImpl; import de.kentoj.scrowlib.servermanager.ServerManager; import io.nats.client.Connection; import io.nats.client.Nats; @@ -26,6 +24,7 @@ public class CoreVelocity { private final Connection nats; private final ServerManager serverManager; private final ServerRegisterWatchdog registerWatchdog; + private final SendPlayerWatchdog sendPlayerWatchdog; @Inject public CoreVelocity(ProxyServer server) { @@ -37,11 +36,13 @@ public class CoreVelocity { .pedantic() .server(natsHost) .build()); - serverManager = new CrownServerManager(nats); + serverManager = new ServerManagerImpl(nats); registerWatchdog = new ServerRegisterWatchdog(nats, serverManager, server); + sendPlayerWatchdog = new SendPlayerWatchdog(nats, server); log.info("listening on " + natsHost + " for server registrations"); registerWatchdog.listen(); + sendPlayerWatchdog.listen(); } catch (IOException | InterruptedException e) { throw new RuntimeException(e); } diff --git a/core-velocity/src/main/java/de/kentoj/scrow/corevelocity/SendPlayerWatchdog.java b/core-velocity/src/main/java/de/kentoj/scrow/corevelocity/SendPlayerWatchdog.java new file mode 100644 index 0000000..27ecdac --- /dev/null +++ b/core-velocity/src/main/java/de/kentoj/scrow/corevelocity/SendPlayerWatchdog.java @@ -0,0 +1,48 @@ +package de.kentoj.scrow.corevelocity; + +import com.velocitypowered.api.proxy.ProxyServer; +import de.kentoj.scrowlib.messaging.DocumentRepository; +import de.kentoj.scrowlib.messaging.NatsRepository; +import de.kentoj.scrowlib.messaging.respose.SafeResult; +import de.kentoj.scrowlib.messaging.respose.SafeSuccessResult; +import io.nats.client.Connection; +import org.bson.Document; +import org.checkerframework.checker.optional.qual.Present; +import org.jetbrains.annotations.NotNull; +import reactor.core.publisher.Mono; + +import java.util.UUID; + +public class SendPlayerWatchdog { + + private final NatsRepository natsRepo; + private final ProxyServer server; + + public SendPlayerWatchdog(Connection nats, ProxyServer server) { + this.natsRepo = new NatsRepository(nats); + this.server = server; + } + + public void listen() { + natsRepo.subscribe("proxy.send-player") + .flatMap(msg -> + handleSendPlayer(DocumentRepository.fromMessage(msg)) + .flatMap(result -> natsRepo.publishSafeResult(msg.getReplyTo(), SafeResult.wrap(result))) + .onErrorResume(err -> natsRepo.publishSafeResult(msg.getReplyTo(), SafeResult.wrap(err))) + ) + .subscribe(); + } + + private Mono<@NotNull Document> handleSendPlayer(Document doc) { + return Mono.fromCallable(() -> { + var player = server.getPlayer(UUID.fromString(doc.getString("playerId"))).orElseThrow(); + var targetServer = server.getServer(doc.getString("handle")).orElseThrow(); + return player.createConnectionRequest(targetServer).connect(); + }) + .flatMap(Mono::fromFuture) + .flatMap(res -> res.isSuccessful() + ? Mono.just(new Document()) + : Mono.error(new RuntimeException(res.getStatus().name().toLowerCase().replace('_', ' '))) + ); + } +}