Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -51,3 +51,9 @@ GEMINI.md
# JFR
/.profileconfig.json


# Web
*.db
*.db-shm
*.db-wal
/assets
9 changes: 9 additions & 0 deletions demo/build.gradle.kts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ plugins {

dependencies {
implementation(rootProject)
implementation(project(":web"))

runtimeOnly(libs.bundles.logback)
}
Expand All @@ -13,4 +14,12 @@ application {
mainModule.set("net.minestom.demo")

applicationDefaultJvmArgs += "-ea"

// Javalin / its Jetty deps are automatic modules with no explicit requires from anyone;
// ALL-MODULE-PATH makes the JVM resolve every module on the module path so they load.
applicationDefaultJvmArgs = listOf("--add-modules", "ALL-MODULE-PATH")
}

tasks.named<JavaExec>("run") {
jvmArgs("--add-modules", "ALL-MODULE-PATH")
}
3 changes: 3 additions & 0 deletions demo/src/main/java/module-info.java
Original file line number Diff line number Diff line change
@@ -1,3 +1,6 @@
module net.minestom.demo {
requires net.minestom.server;
requires net.minestom.web;
requires java.management;
requires jdk.management;
}
6 changes: 5 additions & 1 deletion demo/src/main/java/net/minestom/demo/Main.java
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@
public class Main {

static void main(String[] args) {
System.setProperty("minestom.registry.unsafe-ops", "true"); // TEMP for proxy
System.setProperty("minestom.new-socket-write-lock", "true");
System.setProperty("minestom.registry.unsafe-ops", "true");
MinecraftServer.setCompressionThreshold(0);
Expand Down Expand Up @@ -174,7 +175,10 @@ static void main(String[] args) {
// useful for testing - we don't need to worry about event calls so just set this to a long time
OpenToLAN.open(new OpenToLANConfig().eventCallDelay(Duration.of(1, TimeUnit.DAY)));

minecraftServer.start("0.0.0.0", 25565);
// Optional web dashboard. When enabled the proxy holds the public port and forwards to
// the server below; when disabled the server binds the public port directly.
WebInterface.register();
minecraftServer.start(WebInterface.bindHost(), WebInterface.bindPort());
// minecraftServer.start(java.net.UnixDomainSocketAddress.of("minestom-demo.sock"));
//Runtime.getRuntime().addShutdownHook(new Thread(MinecraftServer::stopCleanly));
}
Expand Down
208 changes: 208 additions & 0 deletions demo/src/main/java/net/minestom/demo/WebInterface.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,208 @@
package net.minestom.demo;

import com.sun.management.OperatingSystemMXBean;
import net.kyori.adventure.nbt.BinaryTagIO;
import net.kyori.adventure.nbt.CompoundBinaryTag;
import net.kyori.adventure.text.Component;
import net.minestom.server.MinecraftServer;
import net.minestom.server.adventure.audience.Audiences;
import net.minestom.server.event.player.PlayerSpawnEvent;
import net.minestom.server.event.server.ServerTickMonitorEvent;
import net.minestom.server.timer.TaskSchedule;
import net.minestom.web.ControlBridge;
import net.minestom.web.ControlPacket;
import net.minestom.web.ProxyConfig;
import net.minestom.web.ProxyServer;

import java.io.ByteArrayOutputStream;
import java.io.OutputStream;
import java.io.PrintStream;
import java.lang.management.ManagementFactory;
import java.net.InetSocketAddress;
import java.nio.charset.StandardCharsets;
import java.util.UUID;
import java.util.concurrent.atomic.AtomicLong;

/// Demo wiring for the web interface. The proxy holds the public Minecraft port and forwards
/// to an upstream Minestom server bound to a loopback port; the dashboard binds to
/// `MINESTOM_WEB_DASHBOARD_PORT` (default 8080). The control bridge carries console lines,
/// 1 Hz JVM/tick metrics, and global NBT into the dashboard.
public final class WebInterface {

private static final boolean ENABLED = Boolean.parseBoolean(
System.getenv().getOrDefault("MINESTOM_WEB_INTERFACE", "true"));

public static String bindHost() {
return ENABLED ? "127.0.0.1" : "0.0.0.0";
}

public static int bindPort() {
return ENABLED ? env("MINESTOM_WEB_UPSTREAM_PORT", 25566) : 25565;
}

public static void register() {
if (!ENABLED) return;
final int proxy = env("MINESTOM_WEB_PROXY_PORT", 25565);
final int dashboard = env("MINESTOM_WEB_DASHBOARD_PORT", 8080);

final ProxyServer web = ProxyServer.builder()
.bindProxy(new InetSocketAddress("0.0.0.0", proxy))
.defaultBackend(new InetSocketAddress("127.0.0.1", bindPort()))
.bindDashboard(new InetSocketAddress("127.0.0.1", dashboard))
.token(System.getenv("MINESTOM_WEB_TOKEN"))
.build();
web.start();
Runtime.getRuntime().addShutdownHook(new Thread(web::close, "Minestom-Web-Shutdown"));

final ControlBridge bridge = web.control();
bridge.setOnOutbound(WebInterface::handleOutbound);
teeConsole(bridge);
schedulePumps(bridge);

System.out.printf("[web] proxy on 0.0.0.0:%d → 127.0.0.1:%d · dashboard http://127.0.0.1:%d/%n",
proxy, bindPort(), dashboard);
}

/// Run dashboard-initiated packets on the tick thread so handlers see the same threading
/// guarantees as a player-typed command.
private static void handleOutbound(ControlPacket packet) {
MinecraftServer.getSchedulerManager().scheduleNextTick(() -> {
final var cm = MinecraftServer.getConnectionManager();
switch (packet) {
case ControlPacket.Command(String c) -> {
final var commands = MinecraftServer.getCommandManager();
commands.execute(commands.getConsoleSender(), c.startsWith("/") ? c.substring(1) : c);
}
case ControlPacket.Broadcast(Component m) -> Audiences.players().sendMessage(m);
case ControlPacket.Kick(UUID id, String reason) -> {
final var p = cm.getOnlinePlayerByUuid(id);
if (p != null) p.kick(Component.text(reason));
}
default -> {
}
}
});
}

/// Tee stdout/stderr into ConsoleLine packets so dashboard subscribers see SLF4J output.
/// Logback resolves the underlying stream per-write, so swapping in after init still works.
private static void teeConsole(ControlBridge bridge) {
System.setOut(linePump(System.out, bridge, "INFO"));
System.setErr(linePump(System.err, bridge, "ERROR"));
}

private static PrintStream linePump(PrintStream original, ControlBridge bridge, String level) {
final ByteArrayOutputStream buf = new ByteArrayOutputStream(256);
final ThreadLocal<Boolean> reentrant = ThreadLocal.withInitial(() -> false);
return new PrintStream(new OutputStream() {
@Override
public synchronized void write(int b) {
original.write(b);
capture(b);
}

@Override
public synchronized void write(byte[] b, int off, int len) {
original.write(b, off, len);
for (int i = 0; i < len; i++) capture(b[off + i] & 0xFF);
}

@Override
public void flush() {
original.flush();
}

private void capture(int b) {
if (b == '\n') flushLine();
else if (b != '\r') buf.write(b);
}

private void flushLine() {
if (buf.size() == 0 || reentrant.get()) {
buf.reset();
return;
}
final String msg = buf.toString(StandardCharsets.UTF_8);
buf.reset();
reentrant.set(true);
try {
bridge.receive(new ControlPacket.ConsoleLine(System.currentTimeMillis(), level, msg));
} catch (Throwable ignored) {
} finally {
reentrant.set(false);
}
}
}, true);
}

private static void schedulePumps(ControlBridge bridge) {
final var os = (OperatingSystemMXBean) ManagementFactory.getOperatingSystemMXBean();
final var runtime = ManagementFactory.getRuntimeMXBean();
final var threads = ManagementFactory.getThreadMXBean();
final var heap = ManagementFactory.getMemoryMXBean();
final long maxMem = Runtime.getRuntime().maxMemory();
final var scheduler = MinecraftServer.getSchedulerManager();
final var connections = MinecraftServer.getConnectionManager();

final AtomicLong msptNanos = new AtomicLong();
MinecraftServer.getGlobalEventHandler().addListener(ServerTickMonitorEvent.class,
e -> msptNanos.set((long) (e.getTickMonitor().getTickTime() * 1_000_000.0)));

scheduler.submitTask(() -> {
final double mspt = msptNanos.get() / 1_000_000.0;
final double tps = mspt > 0 ? Math.min(MinecraftServer.TICK_PER_SECOND, 1000.0 / mspt) : MinecraftServer.TICK_PER_SECOND;
bridge.receive(new ControlPacket.Metrics(
System.currentTimeMillis(),
Math.max(0.0, os.getCpuLoad()),
heap.getHeapMemoryUsage().getUsed(), maxMem,
threads.getThreadCount(), runtime.getUptime(),
mspt, tps,
connections.getOnlinePlayers().size()));
return TaskSchedule.seconds(1);
});

scheduler.submitTask(() -> {
bridge.receive(new ControlPacket.ServerData(CompoundBinaryTag.builder()
.putString("event", "winter_celebration")
.putInt("season", 2)
.putInt("onlinePlayers", connections.getOnlinePlayers().size())
.putLong("epochMs", System.currentTimeMillis())
.build()));
return TaskSchedule.seconds(2);
});

// Per-player NBT on the reserved minestom:web/data channel — proxy intercepts it, the
// client never sees the packet but the dashboard sees the decoded NBT.
MinecraftServer.getGlobalEventHandler().addListener(PlayerSpawnEvent.class, event -> {
final var player = event.getPlayer();
player.scheduler().submitTask(() -> {
if (!player.isOnline()) return TaskSchedule.stop();
final CompoundBinaryTag data = CompoundBinaryTag.builder()
.putString("rank", (player.getUuid().hashCode() & 0xF) == 0 ? "vip" : "member")
.putInt("kills", (int) ((System.currentTimeMillis() / 1000) % 50))
.putString("partyId", UUID.nameUUIDFromBytes(player.getUuid().toString().getBytes()).toString())
.putLong("lastSeenMs", System.currentTimeMillis())
.build();
player.sendPluginMessage(ProxyConfig.DEFAULT_DATA_CHANNEL, encode(data));
return TaskSchedule.seconds(1);
});
});
}

private static byte[] encode(CompoundBinaryTag tag) {
final ByteArrayOutputStream out = new ByteArrayOutputStream();
try {
BinaryTagIO.writer().write(tag, out, BinaryTagIO.Compression.NONE);
} catch (java.io.IOException e) {
throw new RuntimeException(e);
}
return out.toByteArray();
}

private static int env(String name, int def) {
return Integer.parseInt(System.getenv().getOrDefault(name, Integer.toString(def)));
}

private WebInterface() {
}
}
14 changes: 14 additions & 0 deletions demo/src/main/resources/logback.xml
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
<configuration>
<appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>%d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} -- %msg%n</pattern>
</encoder>
</appender>

<root level="INFO">
<appender-ref ref="STDOUT"/>
</root>

<logger name="org.eclipse.jetty" level="WARN"/>
<logger name="io.javalin" level="INFO"/>
</configuration>
1 change: 1 addition & 0 deletions settings.gradle.kts
Original file line number Diff line number Diff line change
Expand Up @@ -9,3 +9,4 @@ include("jmh-benchmarks")
include("jcstress-tests")

include("demo")
include("web")
4 changes: 4 additions & 0 deletions src/main/java/net/minestom/server/item/ItemStack.java
Original file line number Diff line number Diff line change
Expand Up @@ -321,6 +321,10 @@ static Hash of(ItemStack itemStack) {
return ItemStackHashImpl.of(new RegistryTranscoder<>(Transcoder.CRC32_HASH, MinecraftServer.process()), itemStack);
}

default ItemStack asItemStack() {
return ItemStack.AIR;
}

NetworkBuffer.Type<Hash> NETWORK_TYPE = ItemStackHashImpl.NETWORK_TYPE;
}

Expand Down
5 changes: 5 additions & 0 deletions src/main/java/net/minestom/server/item/ItemStackHashImpl.java
Original file line number Diff line number Diff line change
Expand Up @@ -75,5 +75,10 @@ record Item(
addedComponents = Map.copyOf(addedComponents);
removedComponents = Set.copyOf(removedComponents);
}

@Override
public ItemStack asItemStack() {
return ItemStack.of(material, amount);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
import net.minestom.server.gamedata.DataPack;
import net.minestom.server.item.enchant.Enchantment;
import net.minestom.server.network.packet.server.SendablePacket;
import net.minestom.server.network.packet.server.configuration.RegistryDataPacket;
import org.jetbrains.annotations.ApiStatus;
import org.jetbrains.annotations.Nullable;

Expand Down Expand Up @@ -180,4 +181,7 @@ default RegistryKey<T> register(String id, T object, DataPack pack) {
@ApiStatus.Internal
SendablePacket registryDataPacket(Registries registries, boolean excludeVanilla);

@ApiStatus.Internal
void applyRegistryDataPacket(Registries registries, RegistryDataPacket packet);

}
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,7 @@ final class DynamicRegistryImpl<T> implements DynamicRegistry<T> {
private final Map<TagKey<T>, RegistryTagImpl.Backed<T>> tags;

private final Key key;
private final Codec<T> codec;
private final @Nullable Codec<T> codec;

DynamicRegistryImpl(Key key, @Nullable Codec<T> codec) {
this.key = key;
Expand Down Expand Up @@ -115,9 +115,8 @@ public Key key() {

@Override
public @Nullable RegistryKey<T> getKey(Key key) {
if (!keyToValue.containsKey(key))
return null;
return new RegistryKeyImpl<>(key);
final RegistryKey<T> registryKey = new RegistryKeyImpl<>(key);
return keyToId.containsKey(registryKey) ? registryKey : null;
}

@Override
Expand Down Expand Up @@ -246,6 +245,51 @@ public SendablePacket registryDataPacket(Registries registries, boolean excludeV
return createRegistryDataPacket(registries, false);
}

@Override
public void applyRegistryDataPacket(Registries registries, RegistryDataPacket packet) {
Check.argCondition(!key.asString().equals(packet.registryId()),
"Registry data packet {0} cannot be applied to registry {1}", packet.registryId(), key);
final Transcoder<BinaryTag> transcoder = codec != null ? new RegistryTranscoder<>(Transcoder.NBT, registries) : null;
synchronized (REGISTRY_LOCK) {
final Map<Key, T> previousValues = new HashMap<>(keyToValue);
final Map<RegistryKey<T>, DataPack> previousPacks = new HashMap<>(packById.size() * 2);
for (int i = 0; i < idToKey.size(); i++) {
previousPacks.put(idToKey.get(i), packById.get(i));
}

idToValue.clear();
idToKey.clear();
keyToId.clear();
keyToValue.clear();
valueToKey.clear();
packById.clear();

final List<RegistryDataPacket.Entry> entries = packet.entries();
for (int id = 0; id < entries.size(); id++) {
final RegistryDataPacket.Entry entry = entries.get(id);
final RegistryKey<T> registryKey = new RegistryKeyImpl<>(Key.key(entry.id()));
final T value = decodeRegistryDataValue(transcoder, entry, previousValues.get(registryKey.key()));

idToKey.add(registryKey);
idToValue.add(value);
keyToId.put(registryKey, id);
if (value != null) {
keyToValue.put(registryKey.key(), value);
valueToKey.put(value, registryKey);
}
packById.add(previousPacks.get(registryKey));
}
vanillaRegistryDataPacket.invalidate();
}
}

private @Nullable T decodeRegistryDataValue(@Nullable Transcoder<BinaryTag> transcoder,
RegistryDataPacket.Entry entry, @Nullable T fallback) {
if (transcoder == null || entry.data() == null) return fallback;
final Result<T> result = codec.decode(transcoder, entry.data());
return result instanceof Result.Ok(T value) ? value : fallback;
}

@Override
public TagsPacket.Registry tagRegistry() {
final List<TagsPacket.Tag> tagList = new ArrayList<>(tags.size());
Expand Down
Loading