diff --git a/.idea/gradle.xml b/.idea/gradle.xml index 7d5b9ed..6f0ae42 100644 --- a/.idea/gradle.xml +++ b/.idea/gradle.xml @@ -12,6 +12,7 @@ diff --git a/README.md b/README.md index 01d6d11..d43c654 100644 --- a/README.md +++ b/README.md @@ -1,5 +1,8 @@ TODO - timeout/error handling with reactive +- unregister server on shutdown +- infrastructure scripts/glue code w/ daemontools +- 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 310808a..45182e2 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 @@ -10,6 +10,8 @@ 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; @@ -60,6 +62,14 @@ public class CrownServerManager implements ServerManager { .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")); + } + @Override public Mono<@NotNull Void> setState(String handle, String state) { Preconditions.checkArgument(!state.contains(" "), "state may not contain spaces"); 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 cccba3a..7440a62 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 @@ -23,5 +23,7 @@ public interface ServerManager { Mono<@NotNull String> getState(String handle); + Mono<@NotNull Integer> getPort(String handle); + Mono<@NotNull Void> setState(String handle, String state); } 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 new file mode 100644 index 0000000..8519921 --- /dev/null +++ b/core-velocity/src/main/java/de/kentoj/scrow/corevelocity/CoreVelocity.java @@ -0,0 +1,49 @@ +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.ServerManager; +import io.nats.client.Connection; +import io.nats.client.Nats; +import io.nats.client.Options; +import lombok.extern.java.Log; + +import java.io.IOException; + +@Plugin( + id = "corevelocity", + name = "CoreVelocity", + version = "0.1" +) +@Log +public class CoreVelocity { + + private final ProxyServer server; + private final Connection nats; + private final ServerManager serverManager; + private final ServerRegisterWatchdog registerWatchdog; + + @Inject + public CoreVelocity(ProxyServer server) { + this.server = server; + + try { + var natsHost = System.getenv("HOST_NATS"); + nats = Nats.connect(Options.builder() + .pedantic() + .server(natsHost) + .build()); + serverManager = new CrownServerManager(nats); + registerWatchdog = new ServerRegisterWatchdog(nats, serverManager, server); + + log.info("listening on " + natsHost + " for server registrations"); + registerWatchdog.listen(); + } catch (IOException | InterruptedException e) { + throw new RuntimeException(e); + } + } +} \ No newline at end of file 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 new file mode 100644 index 0000000..2b00e9e --- /dev/null +++ b/core-velocity/src/main/java/de/kentoj/scrow/corevelocity/ServerRegisterWatchdog.java @@ -0,0 +1,43 @@ +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.servermanager.ServerManager; +import io.nats.client.Connection; +import lombok.extern.java.Log; +import reactor.core.publisher.Mono; + +import java.net.InetSocketAddress; + +@Log +public class ServerRegisterWatchdog { + + private final NatsRepository natsRepo; + private final ServerManager serverManager; + private final ProxyServer server; + + public ServerRegisterWatchdog(Connection nats, ServerManager serverManager, ProxyServer server) { + this.natsRepo = new NatsRepository(nats); + this.serverManager = serverManager; + this.server = server; + } + + public void listen() { + natsRepo.subscribe("crown.set-state") + .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))); + }); + }) + .subscribe(); + } +} diff --git a/settings.gradle.kts b/settings.gradle.kts index 710883e..049cb9b 100644 --- a/settings.gradle.kts +++ b/settings.gradle.kts @@ -9,4 +9,5 @@ plugins { rootProject.name = "core" include("core-bukkit-api") include("core-bukkit-impl") -include("core-lib") \ No newline at end of file +include("core-lib") +include("core-velocity") \ No newline at end of file