This commit is contained in:
kento2 2026-06-03 12:13:18 +02:00
parent feafe2af62
commit 72b4edf48b
8 changed files with 83 additions and 29 deletions

View file

@ -8,7 +8,7 @@ plugins {
} }
group = "de.kentoj.scrow" group = "de.kentoj.scrow"
version = "0.3" version = "0.4"
repositories { repositories {
mavenCentral() mavenCentral()

View file

@ -4,14 +4,17 @@ import io.nats.client.Message;
import lombok.AccessLevel; import lombok.AccessLevel;
import lombok.NoArgsConstructor; import lombok.NoArgsConstructor;
import org.bson.Document; import org.bson.Document;
import org.jetbrains.annotations.ApiStatus;
@NoArgsConstructor(access = AccessLevel.PRIVATE) @NoArgsConstructor(access = AccessLevel.PRIVATE)
public class DocumentRepository { public class DocumentRepository {
@ApiStatus.Obsolete
public static Document fromMessage(Message message) { public static Document fromMessage(Message message) {
return Document.parse(new String(message.getData())); return Document.parse(new String(message.getData()));
} }
@ApiStatus.Obsolete
public static byte[] toBytes(Document document) { public static byte[] toBytes(Document document) {
return document.toJson().getBytes(); return document.toJson().getBytes();
} }

View file

@ -32,11 +32,11 @@ public class NatsRepository {
return Mono.fromRunnable(() -> nats.publish(subject, data)); 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 -> { return Flux.create(sink -> {
var handler = dispatcher.subscribe(subject, msg -> { var handler = dispatcher.subscribe(subject, msg -> {
try { try {
sink.next(msg); sink.next(new SubscriptionContext(msg));
} catch (Throwable e) { } catch (Throwable e) {
sink.error(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))); 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)); return Mono.fromFuture(nats.request(subject, data));
} }
} }

View file

@ -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());
}
}

View file

@ -1,7 +1,6 @@
package de.kentoj.scrowlib.servermanager; package de.kentoj.scrowlib.servermanager;
import com.google.common.base.Preconditions; import com.google.common.base.Preconditions;
import de.kentoj.scrowlib.messaging.DocumentRepository;
import de.kentoj.scrowlib.messaging.NatsRepository; import de.kentoj.scrowlib.messaging.NatsRepository;
import de.kentoj.scrowlib.messaging.respose.SafeResult; import de.kentoj.scrowlib.messaging.respose.SafeResult;
import io.nats.client.Connection; import io.nats.client.Connection;
@ -22,9 +21,8 @@ public class ServerManagerImpl implements ServerManager {
@Override @Override
public Mono<@NotNull ServerInstance> deployServer(String template) { public Mono<@NotNull ServerInstance> deployServer(String template) {
var payload = DocumentRepository.toBytes(new Document("template", template)); return natsRepo.requestRaw("crown.deploy-server", new Document("template", template))
return natsRepo.request("crown.deploy-server", payload) .map(SafeResult::unwrap)
.map(DocumentRepository::fromMessage)
.map(doc -> new ServerInstance( .map(doc -> new ServerInstance(
template, template,
doc.getString("handle"), doc.getString("handle"),
@ -34,9 +32,8 @@ public class ServerManagerImpl implements ServerManager {
@Override @Override
public Flux<@NotNull ServerInstance> getRunningServers() { public Flux<@NotNull ServerInstance> getRunningServers() {
var payload = DocumentRepository.toBytes(new Document()); return natsRepo.requestRaw("crown.list-servers", new Document())
return natsRepo.request("crown.list-servers", payload) .map(SafeResult::unwrap)
.map(DocumentRepository::fromMessage)
.map(doc -> .map(doc ->
doc.getList("instances", Document.class) doc.getList("instances", Document.class)
) )
@ -51,16 +48,15 @@ public class ServerManagerImpl implements ServerManager {
@Override @Override
public Mono<@NotNull Void> destroyServer(String handle) { public Mono<@NotNull Void> destroyServer(String handle) {
var payload = DocumentRepository.toBytes(new Document("handle", handle)); return natsRepo.requestRaw("crown.destroy-server", new Document("handle", handle))
return natsRepo.request("crown.destroy-server", payload) .map(SafeResult::unwrap)
.then(); .then();
} }
@Override @Override
public Mono<@NotNull ServerInstance> getInfo(String handle) { public Mono<@NotNull ServerInstance> getInfo(String handle) {
var payload = DocumentRepository.toBytes(new Document("handle", handle)); return natsRepo.requestRaw("crown.get-info", new Document("handle", handle))
return natsRepo.request("crown.get-info", payload) .map(SafeResult::unwrap)
.map(DocumentRepository::fromMessage)
.map(doc -> new ServerInstance( .map(doc -> new ServerInstance(
doc.getString("templateName"), doc.getString("templateName"),
handle, handle,
@ -72,19 +68,20 @@ public class ServerManagerImpl implements ServerManager {
@Override @Override
public Mono<@NotNull Void> setState(String handle, String state) { public Mono<@NotNull Void> setState(String handle, String state) {
Preconditions.checkArgument(!state.contains(" "), "state may not contain spaces"); Preconditions.checkArgument(!state.contains(" "), "state may not contain spaces");
var payload = DocumentRepository.toBytes(new Document() var payload = new Document()
.append("handle", handle) .append("handle", handle)
.append("state", state)); .append("state", state);
return natsRepo.request("crown.set-state", payload) return natsRepo.requestRaw("crown.set-state", payload)
.map(SafeResult::unwrap)
.then(); .then();
} }
@Override @Override
public Mono<@NotNull Void> sendPlayer(UUID playerId, String serverHandle) { public Mono<@NotNull Void> sendPlayer(UUID playerId, String serverHandle) {
var payload = DocumentRepository.toBytes(new Document() var payload = new Document()
.append("playerId", playerId.toString()) .append("playerId", playerId.toString())
.append("handle", serverHandle)); .append("handle", serverHandle);
return natsRepo.request("proxy.send-player", payload) return natsRepo.requestRaw("proxy.send-player", payload)
.map(SafeResult::unwrap) .map(SafeResult::unwrap)
.then(); .then();
} }

View file

@ -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<MavenPublication>("ScrowVelocityApi") {
groupId = project.group.toString()
artifactId = project.name
version = project.version.toString()
from(components["java"])
}
}
}

View file

@ -1,14 +1,10 @@
package de.kentoj.scrow.corevelocity; package de.kentoj.scrow.corevelocity;
import com.velocitypowered.api.proxy.ProxyServer; 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.NatsRepository;
import de.kentoj.scrowlib.messaging.respose.SafeResult; import de.kentoj.scrowlib.messaging.respose.SafeResult;
import de.kentoj.scrowlib.messaging.respose.SafeSuccessResult;
import io.nats.client.Connection; import io.nats.client.Connection;
import org.bson.Document; import org.bson.Document;
import org.checkerframework.checker.optional.qual.Present;
import org.jetbrains.annotations.NotNull; import org.jetbrains.annotations.NotNull;
import reactor.core.publisher.Mono; import reactor.core.publisher.Mono;
@ -27,7 +23,7 @@ public class SendPlayerWatchdog {
public void listen() { public void listen() {
natsRepo.subscribe("proxy.send-player") natsRepo.subscribe("proxy.send-player")
.flatMap(msg -> .flatMap(msg ->
handleSendPlayer(DocumentRepository.fromMessage(msg)) handleSendPlayer(SafeResult.unwrap(msg))
.flatMap(result -> natsRepo.publishSafeResult(msg.getReplyTo(), SafeResult.wrap(result))) .flatMap(result -> natsRepo.publishSafeResult(msg.getReplyTo(), SafeResult.wrap(result)))
.onErrorResume(err -> natsRepo.publishSafeResult(msg.getReplyTo(), SafeResult.wrap(err))) .onErrorResume(err -> natsRepo.publishSafeResult(msg.getReplyTo(), SafeResult.wrap(err)))
) )

View file

@ -4,6 +4,7 @@ import com.velocitypowered.api.proxy.ProxyServer;
import com.velocitypowered.api.proxy.server.ServerInfo; import com.velocitypowered.api.proxy.server.ServerInfo;
import de.kentoj.scrowlib.messaging.DocumentRepository; import de.kentoj.scrowlib.messaging.DocumentRepository;
import de.kentoj.scrowlib.messaging.NatsRepository; import de.kentoj.scrowlib.messaging.NatsRepository;
import de.kentoj.scrowlib.messaging.respose.SafeResult;
import de.kentoj.scrowlib.servermanager.ServerManager; import de.kentoj.scrowlib.servermanager.ServerManager;
import io.nats.client.Connection; import io.nats.client.Connection;
import lombok.extern.java.Log; import lombok.extern.java.Log;
@ -27,7 +28,7 @@ public class ServerRegisterWatchdog {
public void listen() { public void listen() {
natsRepo.subscribe("crown.set-state") natsRepo.subscribe("crown.set-state")
.map(DocumentRepository::fromMessage) .map(SafeResult::unwrap)
.concatMap(doc -> { .concatMap(doc -> {
var state = doc.getString("state"); var state = doc.getString("state");
var handle = doc.getString("handle"); var handle = doc.getString("handle");