From cc402a71e768b8d7d016cd5207bddd44d7b0e40f Mon Sep 17 00:00:00 2001 From: kento2 Date: Mon, 1 Jun 2026 09:57:58 +0200 Subject: [PATCH] . --- README.md | 2 +- .../servermanager/CrownServerManager.java | 21 +++++---------- .../scrowlib/servermanager/ServerManager.java | 12 ++++----- .../corevelocity/ServerRegisterWatchdog.java | 27 ++++++++++++++----- 4 files changed, 33 insertions(+), 29 deletions(-) diff --git a/README.md b/README.md index d43c654..d88995e 100644 --- a/README.md +++ b/README.md @@ -2,7 +2,7 @@ TODO - timeout/error handling with reactive - unregister server on shutdown - infrastructure scripts/glue code w/ daemontools -- djb's multilogd for logging +- djb's multilogd for logging(?) environment ----------- diff --git a/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/CrownServerManager.java b/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/CrownServerManager.java index 45182e2..fb0f150 100644 --- a/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/CrownServerManager.java +++ b/core-lib/src/main/java/de/kentoj/scrowlib/servermanager/CrownServerManager.java @@ -1,7 +1,6 @@ package de.kentoj.scrowlib.servermanager; import com.google.common.base.Preconditions; -import com.google.common.io.Files; import de.kentoj.scrowlib.messaging.DocumentRepository; import de.kentoj.scrowlib.messaging.NatsRepository; import io.nats.client.Connection; @@ -10,8 +9,6 @@ import org.jetbrains.annotations.NotNull; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; -import javax.print.Doc; - public class CrownServerManager implements ServerManager { private final NatsRepository natsRepo; @@ -55,19 +52,15 @@ public class CrownServerManager implements ServerManager { } @Override - public Mono<@NotNull String> getState(String handle) { + public Mono<@NotNull ServerInstance> getInfo(String handle) { var payload = DocumentRepository.toBytes(new Document("handle", handle)); - return natsRepo.request("crown.get-state", payload) + return natsRepo.request("crown.get-info", payload) .map(DocumentRepository::fromMessage) - .map(doc -> doc.getString("state")); - } - - @Override - public Mono<@NotNull Integer> getPort(String handle) { - var payload = DocumentRepository.toBytes(new Document("handle", handle)); - return natsRepo.request("crown.get-port", payload) - .map(DocumentRepository::fromMessage) - .map(doc -> doc.getInteger("port")); + .map(doc -> new ServerInstance( + handle, + doc.getInteger("port"), + doc.getString("state") + )); } @Override 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 7440a62..6271544 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 @@ -11,19 +11,17 @@ public interface ServerManager { */ Mono<@NotNull ServerInstance> deployServer(String template); - /** - * @return Flux with handles of running servers - */ - Flux<@NotNull ServerInstance> getRunningServers(); - /** * Attempts a clean shutdown, if that doesn't succeed, kills the server. */ Mono<@NotNull Void> destroyServer(String handle); - Mono<@NotNull String> getState(String handle); + /** + * @return Flux with handles of running servers + */ + Flux<@NotNull ServerInstance> getRunningServers(); - Mono<@NotNull Integer> getPort(String handle); + Mono<@NotNull ServerInstance> getInfo(String handle); Mono<@NotNull Void> setState(String handle, String state); } diff --git a/core-velocity/src/main/java/de/kentoj/scrow/corevelocity/ServerRegisterWatchdog.java b/core-velocity/src/main/java/de/kentoj/scrow/corevelocity/ServerRegisterWatchdog.java index 2b00e9e..5c441fa 100644 --- a/core-velocity/src/main/java/de/kentoj/scrow/corevelocity/ServerRegisterWatchdog.java +++ b/core-velocity/src/main/java/de/kentoj/scrow/corevelocity/ServerRegisterWatchdog.java @@ -7,6 +7,7 @@ import de.kentoj.scrowlib.messaging.NatsRepository; import de.kentoj.scrowlib.servermanager.ServerManager; import io.nats.client.Connection; import lombok.extern.java.Log; +import org.jetbrains.annotations.NotNull; import reactor.core.publisher.Mono; import java.net.InetSocketAddress; @@ -29,15 +30,27 @@ public class ServerRegisterWatchdog { .map(DocumentRepository::fromMessage) .flatMap(doc -> { var state = doc.getString("state"); - if (!"UP".equals(state)) return Mono.empty(); - var handle = doc.getString("handle"); - return serverManager.getPort(handle) - .map(port -> { - log.info("registering " + handle); - return server.registerServer(new ServerInfo(handle, new InetSocketAddress(port))); - }); + return handleStateChange(state, handle); }) .subscribe(); } + + private @NotNull Mono<@NotNull Void> handleStateChange(String state, String handle) { + return switch (state) { + case "UP" -> serverManager.getInfo(handle) + .doOnNext(instance -> { + log.info("registering " + handle); + server.registerServer(new ServerInfo(handle, new InetSocketAddress(instance.getPort()))); + }) + .then(); + case "DOWN" -> serverManager.getInfo(handle) + .doOnNext(instance -> { + log.info("unregistering " + handle); + server.unregisterServer(new ServerInfo(handle, new InetSocketAddress(instance.getPort()))); + }) + .then(); + default -> Mono.empty(); + }; + } }