Add SCOPENET developer API and plugin compatibility layer (LuckPerms, PlaceholderAPI, Vault, CoreProtect, WorldGuard, Spark)
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01DjMbLQujBHunCCu5GpsHaT
This commit is contained in:
38 files changed
+2185
-4
No files matched your search
@@ -19,6 +19,7 @@ public final class Integration implements AutoCloseable {
|
||||
new ArrayBlockingQueue<>(64), r -> daemon(r, "scopenet-login"), new ThreadPoolExecutor.AbortPolicy());
|
||||
private final AtomicBoolean syncing = new AtomicBoolean();
|
||||
private volatile boolean closed;
|
||||
private volatile Consumer<JsonArray> notificationHandler;
|
||||
private volatile LiveMapSync liveMap;
|
||||
private boolean registered;
|
||||
private JsonObject pending;
|
||||
@@ -106,6 +107,12 @@ public final class Integration implements AutoCloseable {
|
||||
if (!response.has("ok") || !response.get("ok").getAsBoolean() || !response.has("kick")
|
||||
|| !response.get("kick").isJsonArray()) throw new IllegalStateException("Invalid sync response");
|
||||
pending = null;
|
||||
Consumer<JsonArray> handler = notificationHandler;
|
||||
if (handler != null && response.has("notifications") && response.get("notifications").isJsonArray()
|
||||
&& response.getAsJsonArray("notifications").size() > 0) {
|
||||
try { handler.accept(response.getAsJsonArray("notifications")); }
|
||||
catch (RuntimeException e) { log.accept("SCOPENET could not process panel events: " + e.getMessage()); }
|
||||
}
|
||||
for (JsonElement element : response.getAsJsonArray("kick")) {
|
||||
try {
|
||||
JsonObject entry = element.getAsJsonObject();
|
||||
@@ -129,6 +136,14 @@ public final class Integration implements AutoCloseable {
|
||||
return client;
|
||||
}
|
||||
|
||||
/**
|
||||
* Receive events the panel queued for this server (level-ups, achievements,
|
||||
* guild changes) with each sync. Called on the sync worker thread.
|
||||
*/
|
||||
public void onNotifications(Consumer<JsonArray> handler) {
|
||||
this.notificationHandler = handler;
|
||||
}
|
||||
|
||||
/**
|
||||
* Start mirroring this server's world and player positions to the panel's
|
||||
* central live map. {@code worldRoot} is the level folder (e.g. {@code world}).
|
||||
|
||||
@@ -0,0 +1,138 @@
|
||||
package net.scopenet.integration;
|
||||
|
||||
import com.google.gson.JsonObject;
|
||||
import com.google.gson.JsonParser;
|
||||
import java.io.IOException;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.nio.file.StandardCopyOption;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.BiConsumer;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
/**
|
||||
* A write-ahead queue for panel requests that must not be lost or repeated,
|
||||
* like the money moved by Vault plugins. Each request gets an operation id and
|
||||
* is saved to disk before it is sent; it is retried until the panel answers
|
||||
* (the panel applies an operation id once), and only then forgotten. A panel
|
||||
* outage or a server crash therefore never drops or duplicates a payment.
|
||||
*/
|
||||
public final class OperationLedger implements AutoCloseable {
|
||||
private record Pending(String endpoint, JsonObject payload) {}
|
||||
|
||||
private final Path file;
|
||||
private final PanelClient client;
|
||||
private final Consumer<String> log;
|
||||
private final BiConsumer<String, JsonObject> done;
|
||||
private final Map<String, Pending> pending = new LinkedHashMap<>();
|
||||
private final ScheduledExecutorService worker = Executors.newSingleThreadScheduledExecutor(r -> {
|
||||
Thread t = new Thread(r, "scopenet-ledger");
|
||||
t.setDaemon(true);
|
||||
return t;
|
||||
});
|
||||
private volatile boolean closed;
|
||||
|
||||
/**
|
||||
* @param done called with (operation id, response) after the panel accepted an operation, or with
|
||||
* (operation id, {"error": "..."}) after it permanently refused it
|
||||
*/
|
||||
public OperationLedger(Path file, PanelClient client, Consumer<String> log, BiConsumer<String, JsonObject> done) {
|
||||
this.file = file;
|
||||
this.client = client;
|
||||
this.log = log;
|
||||
this.done = done;
|
||||
load();
|
||||
worker.scheduleWithFixedDelay(this::drain, 1, 2, TimeUnit.SECONDS);
|
||||
}
|
||||
|
||||
public synchronized int size() { return pending.size(); }
|
||||
|
||||
/** Save the operation and start sending it. Returns its operation id. */
|
||||
public String submit(String endpoint, JsonObject payload) throws IOException {
|
||||
String id = UUID.randomUUID().toString();
|
||||
payload.addProperty("operation_id", id);
|
||||
synchronized (this) {
|
||||
pending.put(id, new Pending(endpoint, payload));
|
||||
try { save(); } catch (IOException e) { pending.remove(id); throw e; }
|
||||
}
|
||||
try { worker.execute(this::drain); } catch (RuntimeException ignored) { /* closing */ }
|
||||
return id;
|
||||
}
|
||||
|
||||
private void drain() {
|
||||
if (closed) return;
|
||||
Map<String, Pending> batch;
|
||||
synchronized (this) { batch = new LinkedHashMap<>(pending); }
|
||||
for (Map.Entry<String, Pending> entry : batch.entrySet()) {
|
||||
if (closed) return;
|
||||
String id = entry.getKey();
|
||||
JsonObject result;
|
||||
try {
|
||||
result = client.post(entry.getValue().endpoint(), entry.getValue().payload().deepCopy());
|
||||
} catch (PanelClient.HttpFailure e) {
|
||||
if (e.status >= 400 && e.status < 500 && e.status != 408 && e.status != 429) {
|
||||
// The panel will never accept this one; report it and move on.
|
||||
JsonObject error = new JsonObject();
|
||||
error.addProperty("error", e.getMessage());
|
||||
finish(id, error);
|
||||
continue;
|
||||
}
|
||||
return; // panel trouble: try again later, keeping the order
|
||||
} catch (IOException e) {
|
||||
return; // unreachable: try again later
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
return;
|
||||
}
|
||||
finish(id, result);
|
||||
}
|
||||
}
|
||||
|
||||
private void finish(String id, JsonObject result) {
|
||||
synchronized (this) {
|
||||
pending.remove(id);
|
||||
try { save(); } catch (IOException e) { log.accept("SCOPENET could not update " + file.getFileName() + ": " + e.getMessage()); }
|
||||
}
|
||||
try { done.accept(id, result); } catch (RuntimeException e) { log.accept("SCOPENET operation callback failed: " + e.getMessage()); }
|
||||
}
|
||||
|
||||
private void load() {
|
||||
try {
|
||||
if (!Files.isRegularFile(file)) return;
|
||||
JsonObject root = JsonParser.parseString(Files.readString(file, StandardCharsets.UTF_8)).getAsJsonObject();
|
||||
for (Map.Entry<String, com.google.gson.JsonElement> e : root.entrySet()) {
|
||||
JsonObject op = e.getValue().getAsJsonObject();
|
||||
pending.put(e.getKey(), new Pending(op.get("endpoint").getAsString(), op.getAsJsonObject("payload")));
|
||||
}
|
||||
if (!pending.isEmpty()) log.accept("SCOPENET is resending " + pending.size() + " saved transaction(s)");
|
||||
} catch (Exception e) {
|
||||
log.accept("SCOPENET could not read " + file.getFileName() + ": " + e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
private void save() throws IOException {
|
||||
JsonObject root = new JsonObject();
|
||||
for (Map.Entry<String, Pending> e : pending.entrySet()) {
|
||||
JsonObject op = new JsonObject();
|
||||
op.addProperty("endpoint", e.getValue().endpoint());
|
||||
op.add("payload", e.getValue().payload());
|
||||
root.add(e.getKey(), op);
|
||||
}
|
||||
Path parent = file.toAbsolutePath().getParent();
|
||||
if (parent != null) Files.createDirectories(parent);
|
||||
Path temp = file.resolveSibling(file.getFileName() + ".tmp");
|
||||
Files.writeString(temp, root.toString(), StandardCharsets.UTF_8);
|
||||
Files.move(temp, file, StandardCopyOption.REPLACE_EXISTING);
|
||||
}
|
||||
|
||||
@Override public void close() {
|
||||
closed = true;
|
||||
worker.shutdownNow();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,106 @@
|
||||
package net.scopenet.integration;
|
||||
|
||||
import com.google.gson.JsonObject;
|
||||
import com.google.gson.JsonParser;
|
||||
import com.sun.net.httpserver.HttpServer;
|
||||
import org.junit.jupiter.api.*;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import java.net.InetSocketAddress;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.*;
|
||||
import java.util.concurrent.*;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
class OperationLedgerTest {
|
||||
private static final String TOKEN = "sn_" + "b".repeat(40);
|
||||
private HttpServer server;
|
||||
private PanelClient client;
|
||||
@TempDir Path dir;
|
||||
|
||||
@BeforeEach void start() throws Exception {
|
||||
server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0);
|
||||
server.start();
|
||||
client = new PanelClient(Settings.of("http://127.0.0.1:" + server.getAddress().getPort(), TOKEN));
|
||||
}
|
||||
|
||||
@AfterEach void stop() {
|
||||
client.close();
|
||||
server.stop(0);
|
||||
}
|
||||
|
||||
private void reply(String path, int[] statuses, AtomicInteger calls, List<String> seen) {
|
||||
server.createContext("/api/server/v1/" + path, exchange -> {
|
||||
String body = new String(exchange.getRequestBody().readAllBytes(), StandardCharsets.UTF_8);
|
||||
int n = calls.getAndIncrement();
|
||||
seen.add(JsonParser.parseString(body).getAsJsonObject().get("operation_id").getAsString());
|
||||
int status = statuses[Math.min(n, statuses.length - 1)];
|
||||
byte[] out = (status == 200 ? "{\"ok\":true,\"balance\":42}" : "{\"error\":\"nope\"}").getBytes(StandardCharsets.UTF_8);
|
||||
exchange.sendResponseHeaders(status, out.length);
|
||||
exchange.getResponseBody().write(out);
|
||||
exchange.close();
|
||||
});
|
||||
}
|
||||
|
||||
private static JsonObject payload() {
|
||||
JsonObject p = new JsonObject();
|
||||
p.addProperty("delta", 5);
|
||||
return p;
|
||||
}
|
||||
|
||||
@Test void retriesWithTheSameIdUntilThePanelAnswers() throws Exception {
|
||||
AtomicInteger calls = new AtomicInteger();
|
||||
List<String> ids = new CopyOnWriteArrayList<>();
|
||||
reply("economy/adjust", new int[]{503, 503, 200}, calls, ids);
|
||||
BlockingQueue<JsonObject> results = new LinkedBlockingQueue<>();
|
||||
try (OperationLedger ledger = new OperationLedger(dir.resolve("pending.json"), client, s -> {}, (id, r) -> results.add(r))) {
|
||||
String id = ledger.submit("economy/adjust", payload());
|
||||
JsonObject result = results.poll(15, TimeUnit.SECONDS);
|
||||
assertNotNull(result);
|
||||
assertEquals(42, result.get("balance").getAsInt());
|
||||
assertEquals(3, calls.get());
|
||||
assertEquals(Set.of(id), new HashSet<>(ids), "every retry carries the same operation id");
|
||||
assertEquals(0, ledger.size());
|
||||
}
|
||||
assertEquals("{}", Files.readString(dir.resolve("pending.json")).trim());
|
||||
}
|
||||
|
||||
@Test void permanentRefusalsAreReportedNotRetried() throws Exception {
|
||||
AtomicInteger calls = new AtomicInteger();
|
||||
reply("economy/adjust", new int[]{400}, calls, new CopyOnWriteArrayList<>());
|
||||
BlockingQueue<JsonObject> results = new LinkedBlockingQueue<>();
|
||||
try (OperationLedger ledger = new OperationLedger(dir.resolve("pending.json"), client, s -> {}, (id, r) -> results.add(r))) {
|
||||
ledger.submit("economy/adjust", payload());
|
||||
JsonObject result = results.poll(10, TimeUnit.SECONDS);
|
||||
assertNotNull(result);
|
||||
assertEquals("nope", result.get("error").getAsString());
|
||||
Thread.sleep(2500);
|
||||
assertEquals(1, calls.get());
|
||||
}
|
||||
}
|
||||
|
||||
@Test void pendingOperationsSurviveARestart() throws Exception {
|
||||
Path file = dir.resolve("pending.json");
|
||||
AtomicInteger calls = new AtomicInteger();
|
||||
List<String> ids = new CopyOnWriteArrayList<>();
|
||||
reply("economy/adjust", new int[]{503}, calls, ids);
|
||||
String id;
|
||||
try (OperationLedger first = new OperationLedger(file, client, s -> {}, (i, r) -> {})) {
|
||||
id = first.submit("economy/adjust", payload());
|
||||
assertTrue(Files.readString(file).contains(id), "saved before anything is sent");
|
||||
}
|
||||
server.removeContext("/api/server/v1/economy/adjust");
|
||||
AtomicInteger ok = new AtomicInteger();
|
||||
List<String> resent = new CopyOnWriteArrayList<>();
|
||||
reply("economy/adjust", new int[]{200}, ok, resent);
|
||||
BlockingQueue<JsonObject> results = new LinkedBlockingQueue<>();
|
||||
List<String> logs = new CopyOnWriteArrayList<>();
|
||||
try (OperationLedger second = new OperationLedger(file, client, logs::add, (i, r) -> results.add(r))) {
|
||||
assertNotNull(results.poll(10, TimeUnit.SECONDS));
|
||||
assertEquals(List.of(id), resent);
|
||||
assertTrue(logs.stream().anyMatch(l -> l.contains("resending 1")));
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user