Live map: plugin/mod chunk+player sync, admin panel Live Map section, launcher transport
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01DjMbLQujBHunCCu5GpsHaT
This commit is contained in:
14 files changed
+769
-8
No files matched your search
@@ -0,0 +1,363 @@
|
||||
package net.scopenet.integration;
|
||||
|
||||
import com.google.gson.*;
|
||||
import java.io.ByteArrayOutputStream;
|
||||
import java.io.IOException;
|
||||
import java.net.URI;
|
||||
import java.net.http.WebSocket;
|
||||
import java.nio.ByteBuffer;
|
||||
import java.nio.channels.FileChannel;
|
||||
import java.nio.file.*;
|
||||
import java.time.Duration;
|
||||
import java.util.*;
|
||||
import java.util.concurrent.*;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
/**
|
||||
* Mirrors the world to the panel's central Live Map: changed chunks (raw,
|
||||
* already-compressed Anvil payloads read straight from the region files) over
|
||||
* HTTPS, and live player positions over a WebSocket (HTTPS as fallback).
|
||||
* Nothing here touches the server thread except {@link #offerPlayers}.
|
||||
*/
|
||||
public final class LiveMapSync implements AutoCloseable {
|
||||
public record Player(String uuid, String name, String dimension, double x, double y, double z, float yaw, float pitch) {}
|
||||
|
||||
private static final String[] DIMS = {"overworld", "the_nether", "the_end"};
|
||||
private static final int SECTOR = 4096;
|
||||
private static final int BATCH_BYTES = 6 * 1024 * 1024;
|
||||
private static final int SCAN_BUDGET_BYTES = 48 * 1024 * 1024;
|
||||
|
||||
private final PanelClient client;
|
||||
private final Path root;
|
||||
private final Consumer<String> log;
|
||||
private final ScheduledExecutorService worker = Executors.newSingleThreadScheduledExecutor(r -> {
|
||||
Thread t = new Thread(r, "scopenet-livemap");
|
||||
t.setDaemon(true);
|
||||
return t;
|
||||
});
|
||||
/** Per region, the timestamp the panel already holds for each of its 1024 chunks. */
|
||||
private final Map<String, int[]> known = new HashMap<>();
|
||||
private final Set<String> manifestLoaded = new HashSet<>();
|
||||
private final AtomicBoolean scanQueued = new AtomicBoolean();
|
||||
private volatile boolean panelEnabled;
|
||||
private volatile boolean closed;
|
||||
private volatile List<Player> players = List.of();
|
||||
private volatile boolean playersChanged;
|
||||
private volatile WebSocket socket;
|
||||
private volatile long socketRetryAt;
|
||||
private long nextScan;
|
||||
private long configCheckAt;
|
||||
private long levelModified = -1;
|
||||
private long lastPlayersSent;
|
||||
private long chunkIntervalSecs = 20;
|
||||
private boolean warnedMissing;
|
||||
|
||||
public LiveMapSync(PanelClient client, Path root, Consumer<String> log) {
|
||||
this.client = client;
|
||||
this.root = root;
|
||||
this.log = log;
|
||||
}
|
||||
|
||||
public void start() {
|
||||
worker.scheduleWithFixedDelay(this::scanSafely, 10, 5, TimeUnit.SECONDS);
|
||||
worker.scheduleWithFixedDelay(this::sendPlayersSafely, 5, 1, TimeUnit.SECONDS);
|
||||
}
|
||||
|
||||
public void offerPlayers(List<Player> next) {
|
||||
if (!next.equals(players)) {
|
||||
players = List.copyOf(next);
|
||||
playersChanged = true;
|
||||
}
|
||||
}
|
||||
|
||||
/** A save just happened; scan soon instead of waiting for the next interval. */
|
||||
public void requestScan() {
|
||||
if (closed || !scanQueued.compareAndSet(false, true)) return;
|
||||
try {
|
||||
worker.schedule(() -> {
|
||||
scanQueued.set(false);
|
||||
nextScan = 0;
|
||||
scanSafely();
|
||||
}, 3, TimeUnit.SECONDS);
|
||||
} catch (RejectedExecutionException e) {
|
||||
scanQueued.set(false);
|
||||
}
|
||||
}
|
||||
|
||||
@Override public void close() {
|
||||
closed = true;
|
||||
WebSocket ws = socket;
|
||||
if (ws != null) ws.abort();
|
||||
worker.shutdownNow();
|
||||
}
|
||||
|
||||
// ---- chunks ----
|
||||
|
||||
private void scanSafely() {
|
||||
if (closed || System.nanoTime() < nextScan) return;
|
||||
try {
|
||||
if (!refreshEnabled()) return;
|
||||
long uploaded = scan();
|
||||
nextScan = System.nanoTime() + TimeUnit.SECONDS.toNanos(uploaded >= SCAN_BUDGET_BYTES ? 1 : chunkIntervalSecs);
|
||||
} catch (PanelClient.HttpFailure e) {
|
||||
if (e.status == 403) panelEnabled = false;
|
||||
else log.accept("SCOPENET live map: the panel refused an upload (" + e.getMessage() + ")");
|
||||
nextScan = System.nanoTime() + TimeUnit.SECONDS.toNanos(30);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
} catch (Exception e) {
|
||||
log.accept("SCOPENET live map sync failed: " + e.getClass().getSimpleName() + ": " + e.getMessage());
|
||||
nextScan = System.nanoTime() + TimeUnit.SECONDS.toNanos(60);
|
||||
}
|
||||
}
|
||||
|
||||
/** True when the panel wants this server's map. Rechecks once a minute while it doesn't. */
|
||||
private boolean refreshEnabled() throws IOException, InterruptedException {
|
||||
if (panelEnabled) return true;
|
||||
long now = System.nanoTime();
|
||||
if (now < configCheckAt) return false;
|
||||
configCheckAt = now + TimeUnit.SECONDS.toNanos(60);
|
||||
try {
|
||||
JsonObject config = client.get("livemap/config");
|
||||
panelEnabled = config.has("enabled") && config.get("enabled").getAsBoolean();
|
||||
if (config.has("chunk_interval_secs")) chunkIntervalSecs = Math.max(5, config.get("chunk_interval_secs").getAsLong());
|
||||
if (panelEnabled) log.accept("SCOPENET live map is on: sending this world to the panel.");
|
||||
} catch (PanelClient.HttpFailure e) {
|
||||
panelEnabled = false; // an older panel without live map support
|
||||
}
|
||||
return panelEnabled;
|
||||
}
|
||||
|
||||
/** Returns how many bytes of chunk data were uploaded. */
|
||||
private long scan() throws IOException, InterruptedException {
|
||||
if (!Files.isDirectory(root)) {
|
||||
if (!warnedMissing) {
|
||||
warnedMissing = true;
|
||||
log.accept("SCOPENET live map: world folder " + root + " not found");
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
uploadLevelIfChanged();
|
||||
long total = 0;
|
||||
for (String dim : DIMS) {
|
||||
Path dir = regionDir(dim);
|
||||
if (dir == null) continue;
|
||||
loadManifest(dim);
|
||||
Batch batch = new Batch(dim);
|
||||
List<Path> files = new ArrayList<>();
|
||||
try (DirectoryStream<Path> stream = Files.newDirectoryStream(dir, "r.*.*.mca")) {
|
||||
stream.forEach(files::add);
|
||||
}
|
||||
files.sort(Comparator.comparing(Path::toString));
|
||||
for (Path file : files) {
|
||||
if (closed || total + batch.size >= SCAN_BUDGET_BYTES) break;
|
||||
total += readRegion(dim, file, batch);
|
||||
}
|
||||
total += batch.flush();
|
||||
if (total >= SCAN_BUDGET_BYTES) break;
|
||||
}
|
||||
return total;
|
||||
}
|
||||
|
||||
private Path regionDir(String dim) {
|
||||
String name = root.getFileName() == null ? "world" : root.getFileName().toString();
|
||||
Path parent = root.toAbsolutePath().getParent();
|
||||
List<Path> candidates = switch (dim) {
|
||||
case "the_nether" -> Arrays.asList(root.resolve("DIM-1/region"), sibling(parent, name + "_nether", "DIM-1/region"),
|
||||
root.resolve("dimensions/minecraft/the_nether/region"));
|
||||
case "the_end" -> Arrays.asList(root.resolve("DIM1/region"), sibling(parent, name + "_the_end", "DIM1/region"),
|
||||
root.resolve("dimensions/minecraft/the_end/region"));
|
||||
default -> List.of(root.resolve("region"), root.resolve("dimensions/minecraft/overworld/region"));
|
||||
};
|
||||
for (Path p : candidates) if (p != null && Files.isDirectory(p)) return p;
|
||||
return null;
|
||||
}
|
||||
|
||||
private static Path sibling(Path parent, String name, String rel) {
|
||||
return parent == null ? null : parent.resolve(name).resolve(rel);
|
||||
}
|
||||
|
||||
private void loadManifest(String dim) throws IOException, InterruptedException {
|
||||
if (manifestLoaded.contains(dim)) return;
|
||||
JsonObject manifest = client.get("livemap/manifest?dim=" + dim);
|
||||
for (JsonElement el : manifest.getAsJsonArray("regions")) {
|
||||
JsonObject region = el.getAsJsonObject();
|
||||
byte[] raw = Base64.getDecoder().decode(region.get("timestamps").getAsString());
|
||||
int[] ts = new int[1024];
|
||||
ByteBuffer.wrap(raw).asIntBuffer().get(ts, 0, Math.min(1024, raw.length / 4));
|
||||
known.put(key(dim, region.get("x").getAsInt(), region.get("z").getAsInt()), ts);
|
||||
}
|
||||
manifestLoaded.add(dim);
|
||||
}
|
||||
|
||||
private static String key(String dim, int rx, int rz) {
|
||||
return dim + ":" + rx + ":" + rz;
|
||||
}
|
||||
|
||||
/** Adds every newer chunk of one region file to the batch; returns bytes uploaded by flushes. */
|
||||
private long readRegion(String dim, Path file, Batch batch) throws IOException, InterruptedException {
|
||||
String[] parts = file.getFileName().toString().split("\\.");
|
||||
int rx, rz;
|
||||
try {
|
||||
rx = Integer.parseInt(parts[1]);
|
||||
rz = Integer.parseInt(parts[2]);
|
||||
} catch (RuntimeException e) {
|
||||
return 0;
|
||||
}
|
||||
long flushed = 0;
|
||||
int[] have = known.getOrDefault(key(dim, rx, rz), new int[1024]);
|
||||
try (FileChannel ch = FileChannel.open(file, StandardOpenOption.READ)) {
|
||||
long size = ch.size();
|
||||
if (size < 2L * SECTOR) return 0;
|
||||
ByteBuffer header = ByteBuffer.allocate(2 * SECTOR);
|
||||
while (header.hasRemaining()) if (ch.read(header, header.position()) < 0) return 0;
|
||||
for (int i = 0; i < 1024; i++) {
|
||||
int loc = header.getInt(i * 4);
|
||||
int ts = header.getInt(SECTOR + i * 4);
|
||||
int count = loc & 0xff;
|
||||
long offset = (loc >>> 8) & 0xffffffL;
|
||||
if (count == 0 || Integer.toUnsignedLong(ts) <= Integer.toUnsignedLong(have[i])) continue;
|
||||
if (offset < 2 || (offset + count) * SECTOR > size) continue;
|
||||
ByteBuffer prefix = ByteBuffer.allocate(5);
|
||||
while (prefix.hasRemaining()) if (ch.read(prefix, offset * SECTOR + prefix.position()) < 0) break;
|
||||
if (prefix.hasRemaining()) continue;
|
||||
int length = prefix.getInt(0);
|
||||
int compression = prefix.get(4) & 0xff;
|
||||
// Oversized chunks live in external .mcc files; skip those and anything malformed.
|
||||
if (length < 2 || length - 1 > count * SECTOR - 5 || (compression & 0x80) != 0 || compression == 0 || compression > 4) continue;
|
||||
ByteBuffer data = ByteBuffer.allocate(length - 1);
|
||||
while (data.hasRemaining()) if (ch.read(data, offset * SECTOR + 5 + data.position()) < 0) break;
|
||||
if (data.hasRemaining()) continue;
|
||||
batch.add(rx, rz, i, rx * 32 + (i & 31), rz * 32 + (i >> 5), ts, compression, data.array());
|
||||
if (batch.size >= BATCH_BYTES) flushed += batch.flush();
|
||||
}
|
||||
}
|
||||
return flushed;
|
||||
}
|
||||
|
||||
/** Chunk records waiting to be uploaded; a chunk counts as sent only once the panel accepted it. */
|
||||
private final class Batch {
|
||||
final String dim;
|
||||
final ByteArrayOutputStream out = new ByteArrayOutputStream();
|
||||
final List<int[]> marks = new ArrayList<>();
|
||||
int size;
|
||||
|
||||
Batch(String dim) {
|
||||
this.dim = dim;
|
||||
}
|
||||
|
||||
void add(int rx, int rz, int index, int x, int z, int ts, int compression, byte[] data) {
|
||||
ByteBuffer rec = ByteBuffer.allocate(17 + data.length);
|
||||
rec.putInt(x).putInt(z).putInt(ts).put((byte) compression).putInt(data.length).put(data);
|
||||
out.write(rec.array(), 0, rec.capacity());
|
||||
marks.add(new int[]{rx, rz, index, ts});
|
||||
size += rec.capacity();
|
||||
}
|
||||
|
||||
long flush() throws IOException, InterruptedException {
|
||||
if (marks.isEmpty()) return 0;
|
||||
long sent = size;
|
||||
try {
|
||||
client.postBytes("livemap/chunks?dim=" + dim, out.toByteArray());
|
||||
for (int[] m : marks) known.computeIfAbsent(key(dim, m[0], m[1]), k -> new int[1024])[m[2]] = m[3];
|
||||
} finally {
|
||||
out.reset();
|
||||
marks.clear();
|
||||
size = 0;
|
||||
}
|
||||
return sent;
|
||||
}
|
||||
}
|
||||
|
||||
private void uploadLevelIfChanged() throws IOException, InterruptedException {
|
||||
Path level = root.resolve("level.dat");
|
||||
if (!Files.isRegularFile(level)) return;
|
||||
long modified = Files.getLastModifiedTime(level).toMillis();
|
||||
if (modified == levelModified) return;
|
||||
byte[] bytes = Files.readAllBytes(level);
|
||||
if (bytes.length > 2 * 1024 * 1024) return;
|
||||
client.postBytes("livemap/level", bytes);
|
||||
levelModified = modified;
|
||||
}
|
||||
|
||||
// ---- players ----
|
||||
|
||||
private void sendPlayersSafely() {
|
||||
if (closed || !panelEnabled) return;
|
||||
long now = System.nanoTime();
|
||||
// Resend unchanged rosters now and then so the panel knows we're alive.
|
||||
boolean heartbeat = now - lastPlayersSent > TimeUnit.SECONDS.toNanos(10);
|
||||
if (!playersChanged && !heartbeat) return;
|
||||
playersChanged = false;
|
||||
lastPlayersSent = now;
|
||||
JsonArray list = new JsonArray();
|
||||
for (Player p : players) {
|
||||
JsonObject o = new JsonObject();
|
||||
o.addProperty("uuid", p.uuid());
|
||||
o.addProperty("name", p.name());
|
||||
o.addProperty("dimension", p.dimension());
|
||||
o.addProperty("x", p.x());
|
||||
o.addProperty("y", p.y());
|
||||
o.addProperty("z", p.z());
|
||||
o.addProperty("yaw", p.yaw());
|
||||
o.addProperty("pitch", p.pitch());
|
||||
list.add(o);
|
||||
}
|
||||
JsonObject body = new JsonObject();
|
||||
body.addProperty("type", "players");
|
||||
body.add("players", list);
|
||||
try {
|
||||
WebSocket ws = socket;
|
||||
if (ws != null && !ws.isOutputClosed()) {
|
||||
ws.sendText(body.toString(), true).get(5, TimeUnit.SECONDS);
|
||||
return;
|
||||
}
|
||||
connectSocket();
|
||||
client.post("livemap/players", body);
|
||||
} catch (PanelClient.HttpFailure e) {
|
||||
if (e.status == 403) panelEnabled = false;
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
} catch (Exception e) {
|
||||
socket = null; // reconnect later; HTTPS covers the gap
|
||||
}
|
||||
}
|
||||
|
||||
private void connectSocket() {
|
||||
long now = System.nanoTime();
|
||||
if (now < socketRetryAt) return;
|
||||
socketRetryAt = now + TimeUnit.SECONDS.toNanos(30);
|
||||
try {
|
||||
Settings settings = client.settings();
|
||||
URI panel = settings.panel();
|
||||
String scheme = "https".equals(panel.getScheme()) ? "wss" : "ws";
|
||||
String path = panel.getRawPath() == null ? "" : panel.getRawPath();
|
||||
URI uri = URI.create(scheme + "://" + panel.getRawAuthority() + path + "/api/server/v1/livemap/ws");
|
||||
client.http().newWebSocketBuilder()
|
||||
.header("Authorization", "Bearer " + settings.token())
|
||||
.connectTimeout(Duration.ofSeconds(5))
|
||||
.buildAsync(uri, new WebSocket.Listener() {
|
||||
@Override public CompletionStage<?> onText(WebSocket ws, CharSequence data, boolean last) {
|
||||
if (data.toString().contains("\"disabled\"")) panelEnabled = false;
|
||||
ws.request(1);
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override public CompletionStage<?> onClose(WebSocket ws, int code, String reason) {
|
||||
if (socket == ws) socket = null;
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override public void onError(WebSocket ws, Throwable error) {
|
||||
if (socket == ws) socket = null;
|
||||
}
|
||||
}).thenAccept(ws -> {
|
||||
if (closed) ws.abort();
|
||||
else socket = ws;
|
||||
});
|
||||
} catch (RuntimeException ignored) {
|
||||
// The HTTPS fallback keeps working.
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user