This commit is contained in:
kento2 2026-05-31 23:06:21 +02:00
parent d3e20606d9
commit 5a255aeac0
16 changed files with 128 additions and 74 deletions

View file

@ -1,3 +1,6 @@
TODO
- timeout/error handling with reactive
environment
-----------

View file

@ -36,8 +36,9 @@ dependencies {
api("io.nats:jnats:2.25.2")
api("de.kentoj.scrow:kencommandapi-core:0.4")
implementation("de.kentoj.scrow:kencommandapi-bukkit:0.4")
api("de.kentoj.scrow:kencommandapi-bukkit:0.4")
api(project(":core-lib"))
}
publishing {

View file

@ -4,6 +4,7 @@ import com.mongodb.reactivestreams.client.MongoDatabase;
import de.kentoj.kencommandapi.BukkitKenCommandApi;
import de.kentoj.scrow.bukkit.friends.FriendRequestService;
import de.kentoj.scrow.bukkit.friends.FriendshipService;
import de.kentoj.scrowlib.servermanager.ServerManager;
import io.nats.client.Connection;
import lombok.AccessLevel;
import lombok.Getter;

View file

@ -22,6 +22,8 @@ dependencies {
implementation(project(":core-bukkit-api"))
implementation("org.spigotmc:spigot-api:1.21.11-R0.1-SNAPSHOT")
implementation("de.kentoj.scrow:kencommandapi-bukkit:0.4")
implementation(project(":core-lib"))
}
tasks {

View file

@ -1,13 +1,10 @@
package de.kentoj.scrow.bukkit;
import de.kentoj.scrow.bukkit.crown.command.CrownCommand;
import de.kentoj.scrow.bukkit.crown.CrownCommand;
import de.kentoj.scrow.bukkit.economy.CoinsCommand;
import de.kentoj.scrow.bukkit.friends.command.FriendCommand;
import org.bukkit.Bukkit;
import org.bukkit.plugin.java.JavaPlugin;
import reactor.core.publisher.Mono;
import java.util.UUID;
public class CoreImplPlugin extends JavaPlugin {
@ -27,7 +24,7 @@ public class CoreImplPlugin extends JavaPlugin {
}
private void registerServer() {
ScrowAPI.getServerManager().updateState(System.getenv("SERVER_HANDLE"), "UP")
ScrowAPI.getServerManager().setState(System.getenv("SERVER_HANDLE"), "UP")
.subscribe(__ -> Bukkit.getLogger().info("registered self"));
}

View file

@ -5,13 +5,13 @@ import com.mongodb.MongoClientSettings;
import com.mongodb.reactivestreams.client.MongoClient;
import com.mongodb.reactivestreams.client.MongoClients;
import de.kentoj.kencommandapi.BukkitKenCommandApi;
import de.kentoj.scrow.bukkit.crown.CrownServerManager;
import de.kentoj.scrow.bukkit.economy.InMemoryEconomyService;
import de.kentoj.scrow.bukkit.economy.MongoEconomyService;
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 io.nats.client.Nats;
import io.nats.client.Options;
import lombok.AccessLevel;

View file

@ -1,4 +1,4 @@
package de.kentoj.scrow.bukkit.crown.command;
package de.kentoj.scrow.bukkit.crown;
import de.kentoj.kencommandapi.BukkitCommandContext;
import de.kentoj.kencommandapi.api.structure.node.CommandNode;

View file

@ -1,4 +1,4 @@
package de.kentoj.scrow.bukkit.crown.command;
package de.kentoj.scrow.bukkit.crown;
import de.kentoj.kencommandapi.BukkitCommandContext;
import de.kentoj.kencommandapi.api.processing.CommandExecutor;
@ -28,8 +28,8 @@ public class CrownDeployServerLiteral implements CommandExecutor<CommandSender,
public void execute(BukkitCommandContext ctx) {
var template = ctx.getArg(templateArg);
ScrowAPI.getServerManager().deployServer(template)
.subscribe(handle -> {
ctx.getSender().sendMessage("Deployed instance of " + template + " with handle " + handle);
.subscribe(server -> {
ctx.getSender().sendMessage("Deployed instance of " + template + " with handle " + server.getHandle());
});
}
}

View file

@ -1,4 +1,4 @@
package de.kentoj.scrow.bukkit.crown.command;
package de.kentoj.scrow.bukkit.crown;
import de.kentoj.kencommandapi.BukkitCommandContext;
import de.kentoj.kencommandapi.api.processing.CommandExecutor;
@ -6,6 +6,7 @@ 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;
@ -24,7 +25,10 @@ public class CrownDestroyServerLiteral implements CommandExecutor<CommandSender,
this.handleArg = ScrowAPI.getCommandManager().createArgument("template", new StringArgumentType<>());
this.handleArg.setSuggestionProvider(ctx -> {
//noinspection CodeBlock2Expr
return ScrowAPI.getServerManager().getRunningServers().collect(Collectors.toSet()).toFuture();
return ScrowAPI.getServerManager().getRunningServers()
.map(ServerInstance::getHandle)
.collect(Collectors.toSet())
.toFuture();
});
this.node.addArgument(handleArg);
}
@ -35,7 +39,7 @@ public class CrownDestroyServerLiteral implements CommandExecutor<CommandSender,
var handle = ctx.getArg(handleArg);
ScrowAPI.getServerManager().destroyServer(handle)
.subscribe(__ -> {
ctx.getSender().sendMessage("Destroyed server with handle " + handle);
ctx.getSender().sendMessage("Destroyed server with handle " + handle);
});
}
}

View file

@ -1,4 +1,4 @@
package de.kentoj.scrow.bukkit.crown.command;
package de.kentoj.scrow.bukkit.crown;
import de.kentoj.kencommandapi.BukkitCommandContext;
import de.kentoj.kencommandapi.api.processing.CommandExecutor;
@ -20,11 +20,15 @@ public class CrownListServersLiteral implements CommandExecutor<CommandSender, B
@Override
public void execute(BukkitCommandContext ctx) {
ScrowAPI.getServerManager().getRunningServers()
.collectList()
.subscribe(uids -> {
var msg = new StringBuilder().append('\n');
uids.forEach(uid -> msg.append(" > ").append(uid));
msg.append('\n');
.reduce(new StringBuilder("\n"), (acc, cur) -> acc.append("> ")
.append(cur.getHandle())
.append(" on port ")
.append(cur.getPort())
.append(": ")
.append(cur.getState())
.append('\n')
)
.subscribe(msg -> {
ctx.getPlayerSender().sendMessage(msg.toString());
});
}

View file

@ -1,48 +0,0 @@
package de.kentoj.scrow.bukkit.crown;
import com.google.common.base.Preconditions;
import de.kentoj.scrow.bukkit.ScrowAPI;
import de.kentoj.scrow.bukkit.ServerManager;
import io.nats.client.Connection;
import lombok.RequiredArgsConstructor;
import org.bson.Document;
import org.jetbrains.annotations.NotNull;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import java.util.Arrays;
public class CrownServerManager implements ServerManager {
private final
@Override
public Mono<@NotNull String> deployServer(String template) {
return Mono.fromFuture(nats.request("crown.deploy-server", new Document()
.append("template", template))
.map(msg -> new String(msg.getData()));
}
@Override
public Flux<@NotNull String> getRunningServers() {
return Mono.fromFuture(nats.request("crown.deploy-server", new byte[0]))
.flatMapMany(msg -> {
var ids = new String(msg.getData()).split("\u001F");
return Flux.fromStream(Arrays.stream(ids));
});
}
@Override
public Mono<@NotNull Void> destroyServer(String handle) {
return Mono.fromFuture(nats.request("crown.destroy-server",handle.getBytes()))
.then();
}
@Override
public Mono<@NotNull Void> updateState(String handle, String state) {
Preconditions.checkArgument(!state.contains(" "), "state may not contain spaces");
var payload = (handle + " " + state).getBytes();
return Mono.fromFuture(nats.request("crown.update-state", payload))
.then();
}
}

View file

@ -1,4 +1,4 @@
package de.kentoj.scrowlib;
package de.kentoj.scrowlib.messaging;
import io.nats.client.Message;
import lombok.AccessLevel;

View file

@ -1,4 +1,4 @@
package de.kentoj.scrowlib;
package de.kentoj.scrowlib.messaging;
import io.nats.client.Connection;
import io.nats.client.Dispatcher;

View file

@ -0,0 +1,72 @@
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;
import org.bson.Document;
import org.jetbrains.annotations.NotNull;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
public class CrownServerManager implements ServerManager {
private final NatsRepository natsRepo;
public CrownServerManager(Connection nats) {
this.natsRepo = new NatsRepository(nats);
}
@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)
.map(doc -> new ServerInstance(
doc.getString("handle"),
doc.getInteger("port")
));
}
@Override
public Flux<@NotNull ServerInstance> getRunningServers() {
var payload = DocumentRepository.toBytes(new Document());
return natsRepo.request("crown.list-servers", payload)
.map(DocumentRepository::fromMessage)
.map(doc ->
doc.getList("instances", Document.class)
)
.flatMapIterable(l -> l)
.map(doc -> new ServerInstance(
doc.getString("handle"),
doc.getInteger("port"),
doc.getString("state")
));
}
@Override
public Mono<@NotNull Void> destroyServer(String handle) {
var payload = DocumentRepository.toBytes(new Document("handle", handle));
return natsRepo.request("crown.destroy-server", payload)
.then();
}
@Override
public Mono<@NotNull String> getState(String handle) {
var payload = DocumentRepository.toBytes(new Document("handle", handle));
return natsRepo.request("crown.get-state", payload)
.map(DocumentRepository::fromMessage)
.map(doc -> doc.getString("state"));
}
@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()
.append("handle", handle)
.append("state", state));
return natsRepo.request("crown.set-state", payload)
.then();
}
}

View file

@ -0,0 +1,16 @@
package de.kentoj.scrowlib.servermanager;
import lombok.AllArgsConstructor;
import lombok.Value;
@Value
@AllArgsConstructor
public class ServerInstance {
String handle;
int port;
String state;
public ServerInstance(String handle, int port) {
this(handle, port, "STARTING");
}
}

View file

@ -1,4 +1,4 @@
package de.kentoj.scrow.bukkit;
package de.kentoj.scrowlib.servermanager;
import org.jetbrains.annotations.NotNull;
import reactor.core.publisher.Flux;
@ -9,17 +9,19 @@ public interface ServerManager {
/**
* @return Mono with handle of deployed server
*/
Mono<@NotNull String> deployServer(String template);
Mono<@NotNull ServerInstance> deployServer(String template);
/**
* @return Flux with handles of running servers
*/
Flux<@NotNull String> getRunningServers();
Flux<@NotNull ServerInstance> getRunningServers();
/**
* Attempts a clean shutdown, if that doesn't succeed, kills the server.
*/
Mono<@NotNull Void> destroyServer(String handle);
Mono<@NotNull Void> updateState(String handle, String state);
Mono<@NotNull String> getState(String handle);
Mono<@NotNull Void> setState(String handle, String state);
}