diff --git a/core-lib/build.gradle.kts b/core-lib/build.gradle.kts index 1bbcd9c..c1b9e43 100644 --- a/core-lib/build.gradle.kts +++ b/core-lib/build.gradle.kts @@ -8,7 +8,7 @@ plugins { } group = "de.kentoj.scrow" -version = "0.3" +version = "0.4" repositories { mavenCentral() diff --git a/core-lib/src/main/java/de/kentoj/scrowlib/messaging/DocumentRepository.java b/core-lib/src/main/java/de/kentoj/scrowlib/messaging/DocumentRepository.java index ff5cdfa..213458d 100644 --- a/core-lib/src/main/java/de/kentoj/scrowlib/messaging/DocumentRepository.java +++ b/core-lib/src/main/java/de/kentoj/scrowlib/messaging/DocumentRepository.java @@ -4,14 +4,17 @@ import io.nats.client.Message; import lombok.AccessLevel; import lombok.NoArgsConstructor; import org.bson.Document; +import org.jetbrains.annotations.ApiStatus; @NoArgsConstructor(access = AccessLevel.PRIVATE) public class DocumentRepository { + @ApiStatus.Obsolete public static Document fromMessage(Message message) { return Document.parse(new String(message.getData())); } + @ApiStatus.Obsolete public static byte[] toBytes(Document document) { return document.toJson().getBytes(); } 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 65d402f..359374c 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 @@ -32,11 +32,11 @@ public class NatsRepository { return Mono.fromRunnable(() -> nats.publish(subject, data)); } - public Flux<@NotNull Message> subscribe(String subject) { + public Flux<@NotNull SubscriptionContext> subscribe(String subject) { return Flux.create(sink -> { var handler = dispatcher.subscribe(subject, msg -> { try { - sink.next(msg); + sink.next(new SubscriptionContext(msg)); } catch (Throwable e) { sink.error(e); } @@ -46,11 +46,11 @@ public class NatsRepository { }); } - public Mono<@NotNull Message> request(String subject, Document document) { + public Mono<@NotNull Message> requestRaw(String subject, Document document) { return Mono.fromFuture(nats.request(subject, DocumentRepository.toBytes(document))); } - public Mono<@NotNull Message> request(String subject, byte[] data) { + public Mono<@NotNull Message> requestRaw(String subject, byte[] data) { return Mono.fromFuture(nats.request(subject, data)); } } diff --git a/core-lib/src/main/java/de/kentoj/scrowlib/messaging/SubscriptionContext.java b/core-lib/src/main/java/de/kentoj/scrowlib/messaging/SubscriptionContext.java new file mode 100644 index 0000000..5382e93 --- /dev/null +++ b/core-lib/src/main/java/de/kentoj/scrowlib/messaging/SubscriptionContext.java @@ -0,0 +1,18 @@ +package de.kentoj.scrowlib.messaging; + +import de.kentoj.scrowlib.messaging.respose.SafeResult; +import io.nats.client.Message; +import lombok.AllArgsConstructor; +import lombok.Value; +import org.bson.Document; + +@Value +@AllArgsConstructor +public class SubscriptionContext { + Message msg; + Document document; + + public SubscriptionContext(Message msg) { + this(msg, SafeResult.fromMessage(msg).unwrap()); + } +} diff --git a/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/ServerManagerImpl.java b/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/ServerManagerImpl.java index 37be894..7d71452 100644 --- a/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/ServerManagerImpl.java +++ b/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/ServerManagerImpl.java @@ -1,7 +1,6 @@ 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; @@ -22,9 +21,8 @@ public class ServerManagerImpl implements ServerManager { @Override public Mono<@NotNull ServerInstance> deployServer(String template) { - var payload = DocumentRepository.toBytes(new Document("template", template)); - return natsRepo.request("crown.deploy-server", payload) - .map(DocumentRepository::fromMessage) + return natsRepo.requestRaw("crown.deploy-server", new Document("template", template)) + .map(SafeResult::unwrap) .map(doc -> new ServerInstance( template, doc.getString("handle"), @@ -34,9 +32,8 @@ public class ServerManagerImpl implements ServerManager { @Override public Flux<@NotNull ServerInstance> getRunningServers() { - var payload = DocumentRepository.toBytes(new Document()); - return natsRepo.request("crown.list-servers", payload) - .map(DocumentRepository::fromMessage) + return natsRepo.requestRaw("crown.list-servers", new Document()) + .map(SafeResult::unwrap) .map(doc -> doc.getList("instances", Document.class) ) @@ -51,16 +48,15 @@ public class ServerManagerImpl implements ServerManager { @Override public Mono<@NotNull Void> destroyServer(String handle) { - var payload = DocumentRepository.toBytes(new Document("handle", handle)); - return natsRepo.request("crown.destroy-server", payload) + return natsRepo.requestRaw("crown.destroy-server", new Document("handle", handle)) + .map(SafeResult::unwrap) .then(); } @Override public Mono<@NotNull ServerInstance> getInfo(String handle) { - var payload = DocumentRepository.toBytes(new Document("handle", handle)); - return natsRepo.request("crown.get-info", payload) - .map(DocumentRepository::fromMessage) + return natsRepo.requestRaw("crown.get-info", new Document("handle", handle)) + .map(SafeResult::unwrap) .map(doc -> new ServerInstance( doc.getString("templateName"), handle, @@ -72,19 +68,20 @@ public class ServerManagerImpl implements ServerManager { @Override public Mono<@NotNull Void> setState(String handle, String state) { Preconditions.checkArgument(!state.contains(" "), "state may not contain spaces"); - var payload = DocumentRepository.toBytes(new Document() + var payload = new Document() .append("handle", handle) - .append("state", state)); - return natsRepo.request("crown.set-state", payload) + .append("state", state); + return natsRepo.requestRaw("crown.set-state", payload) + .map(SafeResult::unwrap) .then(); } @Override public Mono<@NotNull Void> sendPlayer(UUID playerId, String serverHandle) { - var payload = DocumentRepository.toBytes(new Document() + var payload = new Document() .append("playerId", playerId.toString()) - .append("handle", serverHandle)); - return natsRepo.request("proxy.send-player", payload) + .append("handle", serverHandle); + return natsRepo.requestRaw("proxy.send-player", payload) .map(SafeResult::unwrap) .then(); } diff --git a/core-velocity-api/build.gradle.kts b/core-velocity-api/build.gradle.kts new file mode 100644 index 0000000..5dc4e6d --- /dev/null +++ b/core-velocity-api/build.gradle.kts @@ -0,0 +1,39 @@ +plugins { + id("java") + id("java-library") + id("maven-publish") +} + +group = "de.kentoj.scrow" +version = "0.1" + +repositories { + mavenCentral() + maven { + name = "papermc" + url = uri("https://repo.papermc.io/repository/maven-public/") + } +} + +dependencies { + api(project(":core-lib")) + + // https://docs.papermc.io/velocity/dev/creating-your-first-plugin/ + api("com.velocitypowered:velocity-api:3.5.0-SNAPSHOT") + + // http://forge.kentoj.de/scrow/-/packages/maven/de.kentoj.scrow:kencommandapi-core + api("de.kentoj.scrow:kencommandapi-core:0.8") + // http://forge.kentoj.de/scrow/-/packages/maven/de.kentoj.scrow:kencommandapi-velocity + api("de.kentoj.scrow:kencommandapi-velocity:0.8") +} + +publishing { + publications { + create("ScrowVelocityApi") { + groupId = project.group.toString() + artifactId = project.name + version = project.version.toString() + from(components["java"]) + } + } +} diff --git a/core-velocity-impl/src/main/java/de/kentoj/scrow/corevelocity/SendPlayerWatchdog.java b/core-velocity-impl/src/main/java/de/kentoj/scrow/corevelocity/SendPlayerWatchdog.java index 08a8c4a..88ca3f4 100644 --- a/core-velocity-impl/src/main/java/de/kentoj/scrow/corevelocity/SendPlayerWatchdog.java +++ b/core-velocity-impl/src/main/java/de/kentoj/scrow/corevelocity/SendPlayerWatchdog.java @@ -1,14 +1,10 @@ package de.kentoj.scrow.corevelocity; import com.velocitypowered.api.proxy.ProxyServer; -import com.velocitypowered.api.proxy.server.ServerInfo; -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; @@ -27,7 +23,7 @@ public class SendPlayerWatchdog { public void listen() { natsRepo.subscribe("proxy.send-player") .flatMap(msg -> - handleSendPlayer(DocumentRepository.fromMessage(msg)) + handleSendPlayer(SafeResult.unwrap(msg)) .flatMap(result -> natsRepo.publishSafeResult(msg.getReplyTo(), SafeResult.wrap(result))) .onErrorResume(err -> natsRepo.publishSafeResult(msg.getReplyTo(), SafeResult.wrap(err))) ) diff --git a/core-velocity-impl/src/main/java/de/kentoj/scrow/corevelocity/ServerRegisterWatchdog.java b/core-velocity-impl/src/main/java/de/kentoj/scrow/corevelocity/ServerRegisterWatchdog.java index 77b89b2..576332a 100644 --- a/core-velocity-impl/src/main/java/de/kentoj/scrow/corevelocity/ServerRegisterWatchdog.java +++ b/core-velocity-impl/src/main/java/de/kentoj/scrow/corevelocity/ServerRegisterWatchdog.java @@ -4,6 +4,7 @@ import com.velocitypowered.api.proxy.ProxyServer; import com.velocitypowered.api.proxy.server.ServerInfo; import de.kentoj.scrowlib.messaging.DocumentRepository; import de.kentoj.scrowlib.messaging.NatsRepository; +import de.kentoj.scrowlib.messaging.respose.SafeResult; import de.kentoj.scrowlib.servermanager.ServerManager; import io.nats.client.Connection; import lombok.extern.java.Log; @@ -27,7 +28,7 @@ public class ServerRegisterWatchdog { public void listen() { natsRepo.subscribe("crown.set-state") - .map(DocumentRepository::fromMessage) + .map(SafeResult::unwrap) .concatMap(doc -> { var state = doc.getString("state"); var handle = doc.getString("handle");