diff --git a/.gitignore b/.gitignore index 5128661..f974daf 100644 --- a/.gitignore +++ b/.gitignore @@ -19,6 +19,7 @@ data/ # Temp *.log .claude/ +.worktrees/ CLAUDE.md .serena/ .playwright-mcp/ diff --git a/docs/plans/2026-02-22-realtime-websocket-editor-design.md b/docs/plans/2026-02-22-realtime-websocket-editor-design.md new file mode 100644 index 0000000..cb50edb --- /dev/null +++ b/docs/plans/2026-02-22-realtime-websocket-editor-design.md @@ -0,0 +1,233 @@ +# Realtime WebSocket Editor Design + +**Date:** 2026-02-22 +**Status:** Approved +**Scope:** HyperPerms plugin, HyperPermsWeb editor, new CF Workers relay + +## Problem + +The web editor at hyperperms.com uses a polling/REST model: the plugin POSTs session +data to the API, the user edits in the browser, then runs `/hp apply ` to pull +changes back. This is slow, manual, and one-directional. + +## Goal + +Two-way realtime sync between the browser editor and the Hytale plugin via WebSocket, +so edits flow in both directions without manual commands. + +## Decisions + +| Decision | Choice | Rationale | +|---|---|---| +| Direction | Two-way full sync | Editor pushes to server, server pushes to editor | +| Transport | WebSocket via Cloudflare Durable Objects | Already on CF Workers, no extra vendor, ~$5/month | +| Relay model | API-side (CF DO) | Both plugin and browser connect as WS clients to CF | +| Default apply mode | Batch with confirm | Safety: changes accumulate, user clicks Apply | +| Live mode | Toggle in editor UI | Opt-in immediate apply for power users | +| Undo | Server-side undo history (last 20 ops) | Reversible batch applies | +| Message format | Incremental diffs (operations) | Small payloads, natural undo, efficient | +| Conflict handling | Auto-merge non-conflicting, UI for conflicts | Server changes to different entities merge silently | + +## Architecture + +``` +Browser (Next.js) CF Durable Object Hytale Plugin (Java) + | (1 DO per session) | + |--- WS connect -------->| | + | |<-------- WS connect ---------------| + | | | + |-- editor.change ------->| (stores op in pending list) | + | | | + |-- batch.apply --------->| | + | |--- batch.apply ------------------->| + | | [applies to live server] + | |<--- apply.result ------------------| + |<-- apply.result --------| | + | | | + | |<--- server.change -----------------| + |<-- server.change -------| [in-game command ran] +``` + +### Session Lifecycle + +1. Player runs `/hp editor` in-game +2. Plugin creates session via existing REST API (POST /api/session/create) +3. API returns session ID + editor URL (unchanged) +4. Plugin opens WebSocket to `wss://ws.hyperperms.com/session/` as `plugin` role +5. Player opens editor in browser, editor opens WebSocket to same URL as `editor` role +6. DO has both connections -- relay begins +7. On disconnect/timeout, DO hibernates (cost-efficient) +8. Session expires after 24h (same as current TTL) + +### Why CF Durable Objects + +- CF Workers already in use (no new vendor) +- DOs can hold persistent WebSocket connections (unlike Vercel serverless) +- Built-in hibernation API saves costs when idle +- Native WebSocket support with `acceptWebSocket()` API +- Cost: ~$0.15/M requests + $0.50/GB-month -- effectively free at our volume +- Stays within Cloudflare ecosystem + +## Message Protocol + +### Operation Types + +All messages are JSON with a `type` field. Operations are the atomic unit of change. + +```typescript +// === Editor -> DO -> Plugin === + +// Individual changes (accumulated in batch mode, applied immediately in live mode) +{ type: "op", op: "permission.add", target: "group:admin", node: "server.kick", value: true } +{ type: "op", op: "permission.remove", target: "group:admin", node: "server.kick" } +{ type: "op", op: "permission.set", target: "user:", node: "chat.color", value: false } +{ type: "op", op: "group.create", name: "moderator", weight: 50, parents: ["default"] } +{ type: "op", op: "group.delete", name: "moderator" } +{ type: "op", op: "group.setMeta", target: "admin", key: "prefix", value: "&c[Admin] " } +{ type: "op", op: "group.setWeight", target: "admin", weight: 100 } +{ type: "op", op: "group.addParent", target: "moderator", parent: "default" } +{ type: "op", op: "group.removeParent", target: "moderator", parent: "default" } +{ type: "op", op: "user.addGroup", target: "", group: "admin" } +{ type: "op", op: "user.removeGroup", target: "", group: "admin" } +{ type: "op", op: "track.create", name: "staff", groups: ["default", "helper", "mod", "admin"] } +{ type: "op", op: "track.delete", name: "staff" } +{ type: "op", op: "track.addGroup", target: "staff", group: "moderator", position: 2 } +{ type: "op", op: "track.removeGroup", target: "staff", group: "moderator" } + +// Batch apply (batch mode only) +{ type: "batch.apply" } + +// Mode switch +{ type: "session.mode", mode: "live" | "batch" } + +// === Plugin -> DO -> Editor === + +// Server-side change (in-game command, API, another plugin) +{ type: "server.change", op: "permission.add", target: "group:admin", node: "fly.use", + value: true, source: "console", timestamp: 1708617600000 } + +// Apply result +{ type: "apply.result", success: true, applied: 5, failed: 0, + errors: [], undoId: "abc123" } + +// Undo result +{ type: "undo.result", success: true, undoId: "abc123", reverted: 5 } + +// === Control messages (both directions) === + +{ type: "session.sync", data: { groups: [...], users: [...], tracks: [...] } } +{ type: "session.ping" } +{ type: "session.pong" } +{ type: "session.error", code: "CONFLICT", message: "...", details: {...} } +``` + +### Operation Inversion (for undo) + +Each operation has a natural inverse stored in the undo history: + +| Operation | Inverse | +|---|---| +| `permission.add` | `permission.remove` (same target/node) | +| `permission.remove` | `permission.add` (with original value) | +| `group.create` | `group.delete` | +| `group.delete` | `group.create` (with full group data snapshot) | +| `group.setWeight` | `group.setWeight` (with previous weight) | +| `user.addGroup` | `user.removeGroup` | +| `user.removeGroup` | `user.addGroup` | +| `track.addGroup` | `track.removeGroup` | + +## Conflict Detection + +### Non-conflicting (auto-merge) + +Server change to entity A while editor has pending changes to entity B: +- Apply server change to editor state silently +- Show subtle indicator that server state updated (e.g. pulse animation on the changed entity) + +### Conflicting + +Server change to entity A while editor has pending changes to entity A: +- Mark the entity with a conflict indicator in the UI +- Show conflict banner: "Server changed [group: admin] -- your local changes may conflict" +- Options: "Keep mine", "Accept server", "View diff" +- Conflict state clears when user resolves + +### Detection Logic + +Conflict = same `target` value in both a pending editor op and an incoming server change. +Tracked per-entity, not per-field (simple and predictable). + +## Undo History + +- Stored in the Durable Object's transactional storage +- Last 20 batch applies, each with its inverse operations +- Each undo entry: `{ undoId, timestamp, ops: Operation[], inverseOps: Operation[] }` +- Undo request: `{ type: "undo", undoId: "abc123" }` -> sends inverse ops to plugin +- History clears on session expiry + +## Three Codebases + +### 1. CF Workers (new Durable Object) + +**Location:** New worker in existing CF project or standalone +**Responsibilities:** +- Accept WebSocket connections from editor and plugin +- Route messages between them based on type +- Store pending operations (batch mode) +- Maintain undo history +- Handle hibernation for cost efficiency +- Session authentication (validate session ID against Upstash Redis) + +**Key files:** +- `src/session-do.ts` -- Durable Object class with WebSocket handlers +- `src/index.ts` -- Worker entry point, routes `/session/:id` to DO +- `wrangler.toml` -- DO binding configuration + +### 2. HyperPermsWeb (Next.js editor) + +**Changes:** +- New WebSocket hook: `useRealtimeSession(sessionId)` -- manages WS connection, reconnection, message handling +- Editor state: track pending operations, conflict state +- UI: "Live" toggle switch, Apply button (batch mode), conflict banners, undo button +- Connection status indicator (connected/reconnecting/disconnected) +- Replace current `PUT /api/session` save flow with WebSocket ops + +**Key files:** +- `src/hooks/useRealtimeSession.ts` -- WebSocket connection + state management +- `src/components/editor/RealtimeStatus.tsx` -- Connection indicator +- `src/components/editor/ConflictBanner.tsx` -- Conflict resolution UI +- `src/components/editor/LiveModeToggle.tsx` -- Live/batch toggle +- Modifications to existing editor components to dispatch ops instead of direct state mutation + +### 3. HyperPerms (Java plugin) + +**Changes:** +- New `WebSocketClient` class using Java 11+ `HttpClient` WebSocket API +- Change listener: hook into existing EventBus to detect in-game permission changes +- Apply handler: receive ops from editor, apply them through existing manager layer +- Undo support: store pre-apply snapshots for revert +- Auto-reconnect on connection loss +- New config options: `webEditor.realtimeEnabled`, `webEditor.wsUrl` + +**Key files:** +- `src/main/java/com/hyperperms/web/RealtimeClient.java` -- WebSocket client +- `src/main/java/com/hyperperms/web/RealtimeMessageHandler.java` -- Message routing +- `src/main/java/com/hyperperms/web/OperationApplier.java` -- Applies ops to live state +- `src/main/java/com/hyperperms/web/ChangeListener.java` -- Detects in-game changes, sends to WS +- Modifications to `WebEditorService.java` -- Start WS after session creation + +## Security + +- Session ID serves as auth token (same as current model) +- DO validates session exists in Redis before accepting WS upgrade +- Plugin identifies as `role: plugin` on connect, editor as `role: editor` +- Only one plugin connection per session (reject duplicates) +- Rate limiting on operations (prevent abuse from editor) +- All traffic over WSS (TLS) + +## Migration / Backwards Compatibility + +- Existing REST flow continues to work (no breaking changes) +- WebSocket is additive -- if WS connection fails, editor falls back to REST save/load +- Plugin config: `webEditor.realtimeEnabled: true` (default true for new installs) +- `/hp apply` command still works as fallback diff --git a/src/main/java/com/hyperperms/HyperPerms.java b/src/main/java/com/hyperperms/HyperPerms.java index ad1aa33..504ffad 100644 --- a/src/main/java/com/hyperperms/HyperPerms.java +++ b/src/main/java/com/hyperperms/HyperPerms.java @@ -217,7 +217,7 @@ public void enable() { // Initialize managers with event bus groupManager = new GroupManagerImpl(storage, cacheInvalidator, eventBus); - trackManager = new TrackManagerImpl(storage); + trackManager = new TrackManagerImpl(storage, eventBus); userManager = new UserManagerImpl(storage, cache, eventBus, config.getDefaultGroup()); // Load data @@ -428,6 +428,11 @@ public void disable() { placeholderApiIntegration.unregister(); } + // Disconnect WebSocket before shutting down scheduler + if (webEditorService != null) { + webEditorService.disconnectWebSocket(); + } + // Stop scheduled tasks FIRST to prevent new storage executor submissions if (expiryTask != null) { expiryTask.cancel(true); diff --git a/src/main/java/com/hyperperms/api/events/HyperPermsEvent.java b/src/main/java/com/hyperperms/api/events/HyperPermsEvent.java index aef17b3..a097263 100644 --- a/src/main/java/com/hyperperms/api/events/HyperPermsEvent.java +++ b/src/main/java/com/hyperperms/api/events/HyperPermsEvent.java @@ -71,6 +71,21 @@ enum EventType { */ DATA_RELOAD, + /** + * Fired when a track is created. + */ + TRACK_CREATE, + + /** + * Fired when a track is deleted. + */ + TRACK_DELETE, + + /** + * Fired when a track is modified. + */ + TRACK_MODIFY, + /** * Fired when a user is promoted along a track. */ diff --git a/src/main/java/com/hyperperms/api/events/TrackCreateEvent.java b/src/main/java/com/hyperperms/api/events/TrackCreateEvent.java new file mode 100644 index 0000000..3d8cb7b --- /dev/null +++ b/src/main/java/com/hyperperms/api/events/TrackCreateEvent.java @@ -0,0 +1,112 @@ +package com.hyperperms.api.events; + +import com.hyperperms.model.Track; +import org.jetbrains.annotations.NotNull; +import org.jetbrains.annotations.Nullable; + +/** + * Event fired when a track is created. + *

+ * This event can be cancelled to prevent the track from being created. + * The PRE state fires before creation, the POST state fires after. + */ +public final class TrackCreateEvent implements HyperPermsEvent, Cancellable { + + /** + * The state of the track creation event. + */ + public enum State { + /** + * Before the track is created. The event can be cancelled at this point. + */ + PRE, + + /** + * After the track has been created. Cancellation has no effect. + */ + POST + } + + private final String trackName; + private final Track track; + private final State state; + private boolean cancelled; + + /** + * Creates a PRE event for track creation. + * + * @param trackName the name of the track being created + */ + public TrackCreateEvent(@NotNull String trackName) { + this.trackName = trackName; + this.track = null; + this.state = State.PRE; + this.cancelled = false; + } + + /** + * Creates a POST event for track creation. + * + * @param track the created track + */ + public TrackCreateEvent(@NotNull Track track) { + this.trackName = track.getName(); + this.track = track; + this.state = State.POST; + this.cancelled = false; + } + + @Override + public EventType getType() { + return EventType.TRACK_CREATE; + } + + /** + * Gets the name of the track being created. + * + * @return the track name + */ + @NotNull + public String getTrackName() { + return trackName; + } + + /** + * Gets the created track. + *

+ * This is only available in the POST state. Returns null in PRE state. + * + * @return the track, or null if PRE state + */ + @Nullable + public Track getTrack() { + return track; + } + + /** + * Gets the state of this event. + * + * @return the state + */ + @NotNull + public State getState() { + return state; + } + + @Override + public boolean isCancelled() { + return cancelled; + } + + @Override + public void setCancelled(boolean cancelled) { + if (state == State.PRE) { + this.cancelled = cancelled; + } + } + + @Override + public String toString() { + return "TrackCreateEvent{trackName='" + trackName + "', state=" + state + ", cancelled=" + cancelled + "}"; + } +} diff --git a/src/main/java/com/hyperperms/api/events/TrackDeleteEvent.java b/src/main/java/com/hyperperms/api/events/TrackDeleteEvent.java new file mode 100644 index 0000000..d3ca80b --- /dev/null +++ b/src/main/java/com/hyperperms/api/events/TrackDeleteEvent.java @@ -0,0 +1,112 @@ +package com.hyperperms.api.events; + +import com.hyperperms.model.Track; +import org.jetbrains.annotations.NotNull; +import org.jetbrains.annotations.Nullable; + +/** + * Event fired when a track is deleted. + *

+ * This event can be cancelled to prevent the track from being deleted. + * The PRE state fires before deletion, the POST state fires after. + */ +public final class TrackDeleteEvent implements HyperPermsEvent, Cancellable { + + /** + * The state of the track deletion event. + */ + public enum State { + /** + * Before the track is deleted. The event can be cancelled at this point. + */ + PRE, + + /** + * After the track has been deleted. Cancellation has no effect. + */ + POST + } + + private final String trackName; + private final Track track; + private final State state; + private boolean cancelled; + + /** + * Creates a PRE event for track deletion. + * + * @param track the track being deleted + */ + public TrackDeleteEvent(@NotNull Track track) { + this.trackName = track.getName(); + this.track = track; + this.state = State.PRE; + this.cancelled = false; + } + + /** + * Creates a POST event for track deletion. + * + * @param trackName the name of the deleted track + */ + public TrackDeleteEvent(@NotNull String trackName) { + this.trackName = trackName; + this.track = null; + this.state = State.POST; + this.cancelled = false; + } + + @Override + public EventType getType() { + return EventType.TRACK_DELETE; + } + + /** + * Gets the name of the track being deleted. + * + * @return the track name + */ + @NotNull + public String getTrackName() { + return trackName; + } + + /** + * Gets the track being deleted. + *

+ * This is only available in the PRE state. Returns null in POST state. + * + * @return the track, or null if POST state + */ + @Nullable + public Track getTrack() { + return track; + } + + /** + * Gets the state of this event. + * + * @return the state + */ + @NotNull + public State getState() { + return state; + } + + @Override + public boolean isCancelled() { + return cancelled; + } + + @Override + public void setCancelled(boolean cancelled) { + if (state == State.PRE) { + this.cancelled = cancelled; + } + } + + @Override + public String toString() { + return "TrackDeleteEvent{trackName='" + trackName + "', state=" + state + ", cancelled=" + cancelled + "}"; + } +} diff --git a/src/main/java/com/hyperperms/api/events/TrackModifyEvent.java b/src/main/java/com/hyperperms/api/events/TrackModifyEvent.java new file mode 100644 index 0000000..cf5efad --- /dev/null +++ b/src/main/java/com/hyperperms/api/events/TrackModifyEvent.java @@ -0,0 +1,71 @@ +package com.hyperperms.api.events; + +import com.hyperperms.model.Track; +import org.jetbrains.annotations.NotNull; + +import java.util.List; + +/** + * Event fired when a track's groups are modified. + *

+ * This event is fired after a track is saved with updated groups. + */ +public final class TrackModifyEvent implements HyperPermsEvent { + + private final Track track; + private final List oldGroups; + private final List newGroups; + + /** + * Creates a new track modify event. + * + * @param track the modified track + * @param oldGroups the previous group list + * @param newGroups the new group list + */ + public TrackModifyEvent(@NotNull Track track, @NotNull List oldGroups, @NotNull List newGroups) { + this.track = track; + this.oldGroups = List.copyOf(oldGroups); + this.newGroups = List.copyOf(newGroups); + } + + @Override + public EventType getType() { + return EventType.TRACK_MODIFY; + } + + /** + * Gets the modified track. + * + * @return the track + */ + @NotNull + public Track getTrack() { + return track; + } + + /** + * Gets the previous group list. + * + * @return the old groups + */ + @NotNull + public List getOldGroups() { + return oldGroups; + } + + /** + * Gets the new group list. + * + * @return the new groups + */ + @NotNull + public List getNewGroups() { + return newGroups; + } + + @Override + public String toString() { + return "TrackModifyEvent{track='" + track.getName() + "', oldGroups=" + oldGroups + ", newGroups=" + newGroups + "}"; + } +} diff --git a/src/main/java/com/hyperperms/commands/ApplySubCommand.java b/src/main/java/com/hyperperms/commands/ApplySubCommand.java index 12b90a1..0ac1fa5 100644 --- a/src/main/java/com/hyperperms/commands/ApplySubCommand.java +++ b/src/main/java/com/hyperperms/commands/ApplySubCommand.java @@ -4,6 +4,7 @@ import com.hyperperms.util.Logger; import com.hyperperms.web.ChangeApplier; import com.hyperperms.web.WebEditorService; +import com.hyperperms.web.WebSocketSessionClient; import com.hyperperms.web.dto.Change; import com.hypixel.hytale.server.core.Message; import com.hypixel.hytale.server.core.command.system.AbstractCommand; @@ -44,9 +45,55 @@ protected CompletableFuture execute(CommandContext ctx) { return CompletableFuture.completedFuture(null); } + // Try WebSocket apply if connected, otherwise fall back to REST + WebSocketSessionClient wsClient = webEditorService.getActiveWsClient(); + if (wsClient != null && wsClient.isConnected()) { + ctx.sender().sendMessage(Message.raw("Fetching changes via live connection...")); + + webEditorService.fetchChanges(sessionId) + .thenCompose(result -> { + if (!result.isSuccess()) { + ctx.sender().sendMessage(Message.raw("Failed to fetch changes: " + result.getError())); + return CompletableFuture.completedFuture((Void) null); + } + + List changes = result.getChanges(); + if (changes == null || changes.isEmpty()) { + ctx.sender().sendMessage(Message.raw("No changes found for this session.")); + return CompletableFuture.completedFuture((Void) null); + } + + ctx.sender().sendMessage(Message.raw("Found " + changes.size() + " change(s). Applying via WebSocket...")); + + return wsClient.sendBatchApply(changes) + .thenAccept(batchResult -> { + ctx.sender().sendMessage(Message.raw("")); + ctx.sender().sendMessage(Message.raw("=== Apply Results (WebSocket) ===")); + ctx.sender().sendMessage(Message.raw("Successful: " + batchResult.getSuccessCount())); + ctx.sender().sendMessage(Message.raw("Failed: " + batchResult.getFailedCount())); + if (batchResult.getSuccessCount() > 0) { + ctx.sender().sendMessage(Message.raw("")); + ctx.sender().sendMessage(Message.raw("Changes applied and synced.")); + Logger.info("Applied " + batchResult.getSuccessCount() + " changes via WebSocket for session " + sessionId); + } + }); + }) + .exceptionally(e -> { + ctx.sender().sendMessage(Message.raw("WebSocket apply failed, falling back to REST: " + e.getMessage())); + Logger.warn("WebSocket apply failed: " + e.getMessage()); + applyViaRest(ctx, sessionId); + return null; + }); + } else { + applyViaRest(ctx, sessionId); + } + + return CompletableFuture.completedFuture(null); + } + + private void applyViaRest(CommandContext ctx, String sessionId) { ctx.sender().sendMessage(Message.raw("Fetching changes from web editor...")); - // Fetch changes from the API webEditorService.fetchChanges(sessionId) .thenAccept(result -> { if (!result.isSuccess()) { @@ -69,10 +116,8 @@ protected CompletableFuture execute(CommandContext ctx) { ctx.sender().sendMessage(Message.raw("Found " + changes.size() + " change(s). Applying...")); - // Apply changes ChangeApplier.ApplyResult applyResult = changeApplier.applyChanges(changes); - // Report results ctx.sender().sendMessage(Message.raw("")); ctx.sender().sendMessage(Message.raw("=== Apply Results ===")); ctx.sender().sendMessage(Message.raw("Successful: " + applyResult.getSuccessCount())); @@ -97,7 +142,5 @@ protected CompletableFuture execute(CommandContext ctx) { Logger.warn("Failed to apply changes: " + e.getMessage()); return null; }); - - return CompletableFuture.completedFuture(null); } } diff --git a/src/main/java/com/hyperperms/commands/EditorSubCommand.java b/src/main/java/com/hyperperms/commands/EditorSubCommand.java index 4dda727..39a902d 100644 --- a/src/main/java/com/hyperperms/commands/EditorSubCommand.java +++ b/src/main/java/com/hyperperms/commands/EditorSubCommand.java @@ -72,6 +72,13 @@ protected CompletableFuture execute(CommandContext ctx) { .insert(Message.raw(" and enter your session ID manually.")) ); + // Connect WebSocket if available and enabled + if (response.getWsUrl() != null && hyperPerms.getConfig().isWebEditorWebsocketEnabled()) { + webEditorService.connectWebSocket(response.getSessionId(), response.getWsUrl()); + ctx.sender().sendMessage(Message.raw("Live connection established - changes sync in realtime.") + .color(new java.awt.Color(0x55FFFF))); + } + // Log full URL to console for easy copying Logger.info("Web editor session created: " + response.getSessionId()); Logger.info("Editor URL: " + editorUrl); diff --git a/src/main/java/com/hyperperms/config/HyperPermsConfig.java b/src/main/java/com/hyperperms/config/HyperPermsConfig.java index 68e0e5c..fe5d380 100644 --- a/src/main/java/com/hyperperms/config/HyperPermsConfig.java +++ b/src/main/java/com/hyperperms/config/HyperPermsConfig.java @@ -297,6 +297,10 @@ private JsonObject createDefaultConfig() { webEditor.addProperty("url", "https://www.hyperperms.com"); webEditor.addProperty("apiUrl", ""); // Empty = use main URL for backward compatibility webEditor.addProperty("timeoutSeconds", 10); + webEditor.addProperty("websocketEnabled", true); + webEditor.addProperty("websocketReconnectMaxAttempts", 10); + webEditor.addProperty("websocketReconnectMaxDelaySeconds", 30); + webEditor.addProperty("websocketPingTimeoutSeconds", 90); root.add("webEditor", webEditor); // Tab list settings @@ -596,6 +600,43 @@ public int getWebEditorTimeoutSeconds() { return getNestedInt("webEditor", "timeoutSeconds", 10); } + /** + * Checks if WebSocket realtime sync is enabled. + * + * @return true if WebSocket is enabled + */ + public boolean isWebEditorWebsocketEnabled() { + return getNestedBoolean("webEditor", "websocketEnabled", true); + } + + /** + * Gets the maximum number of WebSocket reconnect attempts. + * + * @return the max reconnect attempts + */ + public int getWebEditorWebsocketReconnectMaxAttempts() { + return getNestedInt("webEditor", "websocketReconnectMaxAttempts", 10); + } + + /** + * Gets the maximum delay in seconds between WebSocket reconnect attempts. + * + * @return the max reconnect delay in seconds + */ + public int getWebEditorWebsocketReconnectMaxDelaySeconds() { + return getNestedInt("webEditor", "websocketReconnectMaxDelaySeconds", 30); + } + + /** + * Gets the WebSocket ping timeout in seconds. + * If no ping is received within this period, the connection is considered dead. + * + * @return the ping timeout in seconds + */ + public int getWebEditorWebsocketPingTimeoutSeconds() { + return getNestedInt("webEditor", "websocketPingTimeoutSeconds", 90); + } + // ==================== Faction Integration Settings ==================== /** diff --git a/src/main/java/com/hyperperms/config/WebEditorConfig.java b/src/main/java/com/hyperperms/config/WebEditorConfig.java index 2eef426..c664acd 100644 --- a/src/main/java/com/hyperperms/config/WebEditorConfig.java +++ b/src/main/java/com/hyperperms/config/WebEditorConfig.java @@ -15,6 +15,10 @@ public final class WebEditorConfig extends ConfigFile { private String url; private String apiUrl; private int timeoutSeconds; + private boolean websocketEnabled; + private int websocketReconnectMaxAttempts; + private int websocketReconnectMaxDelaySeconds; + private int websocketPingTimeoutSeconds; public WebEditorConfig(@NotNull Path dataDirectory) { super(dataDirectory.resolve("webeditor.json")); @@ -25,6 +29,10 @@ protected void createDefaults() { url = "https://www.hyperperms.com"; apiUrl = ""; timeoutSeconds = 10; + websocketEnabled = true; + websocketReconnectMaxAttempts = 10; + websocketReconnectMaxDelaySeconds = 30; + websocketPingTimeoutSeconds = 90; } @Override @@ -32,6 +40,10 @@ protected void loadFromJson(@NotNull JsonObject root) { url = getString(root, "url", "https://www.hyperperms.com"); apiUrl = getString(root, "apiUrl", ""); timeoutSeconds = getInt(root, "timeoutSeconds", 10); + websocketEnabled = getBool(root, "websocketEnabled", true); + websocketReconnectMaxAttempts = getInt(root, "websocketReconnectMaxAttempts", 10); + websocketReconnectMaxDelaySeconds = getInt(root, "websocketReconnectMaxDelaySeconds", 30); + websocketPingTimeoutSeconds = getInt(root, "websocketPingTimeoutSeconds", 90); } @Override @@ -41,6 +53,10 @@ protected JsonObject toJson() { root.addProperty("url", url); root.addProperty("apiUrl", apiUrl); root.addProperty("timeoutSeconds", timeoutSeconds); + root.addProperty("websocketEnabled", websocketEnabled); + root.addProperty("websocketReconnectMaxAttempts", websocketReconnectMaxAttempts); + root.addProperty("websocketReconnectMaxDelaySeconds", websocketReconnectMaxDelaySeconds); + root.addProperty("websocketPingTimeoutSeconds", websocketPingTimeoutSeconds); return root; } @@ -49,6 +65,9 @@ protected JsonObject toJson() { public ValidationResult validate() { ValidationResult result = new ValidationResult(); timeoutSeconds = validateRange(result, "timeoutSeconds", timeoutSeconds, 1, 120, 10); + websocketReconnectMaxAttempts = validateRange(result, "websocketReconnectMaxAttempts", websocketReconnectMaxAttempts, 1, 100, 10); + websocketReconnectMaxDelaySeconds = validateRange(result, "websocketReconnectMaxDelaySeconds", websocketReconnectMaxDelaySeconds, 1, 300, 30); + websocketPingTimeoutSeconds = validateRange(result, "websocketPingTimeoutSeconds", websocketPingTimeoutSeconds, 10, 600, 90); return result; } @@ -64,4 +83,12 @@ public String getApiUrl() { } public int getTimeoutSeconds() { return timeoutSeconds; } + + public boolean isWebsocketEnabled() { return websocketEnabled; } + + public int getWebsocketReconnectMaxAttempts() { return websocketReconnectMaxAttempts; } + + public int getWebsocketReconnectMaxDelaySeconds() { return websocketReconnectMaxDelaySeconds; } + + public int getWebsocketPingTimeoutSeconds() { return websocketPingTimeoutSeconds; } } diff --git a/src/main/java/com/hyperperms/manager/TrackManagerImpl.java b/src/main/java/com/hyperperms/manager/TrackManagerImpl.java index 03328a6..f43e474 100644 --- a/src/main/java/com/hyperperms/manager/TrackManagerImpl.java +++ b/src/main/java/com/hyperperms/manager/TrackManagerImpl.java @@ -1,6 +1,10 @@ package com.hyperperms.manager; import com.hyperperms.api.HyperPermsAPI.TrackManager; +import com.hyperperms.api.events.EventBus; +import com.hyperperms.api.events.TrackCreateEvent; +import com.hyperperms.api.events.TrackDeleteEvent; +import com.hyperperms.api.events.TrackModifyEvent; import com.hyperperms.model.Track; import com.hyperperms.storage.StorageProvider; import com.hyperperms.util.Logger; @@ -17,10 +21,12 @@ public final class TrackManagerImpl implements TrackManager { private final StorageProvider storage; + private final EventBus eventBus; private final Map loadedTracks = new ConcurrentHashMap<>(); - public TrackManagerImpl(@NotNull StorageProvider storage) { + public TrackManagerImpl(@NotNull StorageProvider storage, @NotNull EventBus eventBus) { this.storage = storage; + this.eventBus = eventBus; } @Override @@ -47,6 +53,14 @@ public CompletableFuture> loadTrack(@NotNull String name) { @NotNull public Track createTrack(@NotNull String name) { String lowerName = name.toLowerCase(); + + // Fire PRE event + TrackCreateEvent preEvent = new TrackCreateEvent(lowerName); + eventBus.fire(preEvent); + if (preEvent.isCancelled()) { + throw new IllegalStateException("Track creation cancelled by event handler: " + name); + } + Track track = new Track(lowerName); // putIfAbsent is atomic - prevents concurrent duplicate creation @@ -57,20 +71,48 @@ public Track createTrack(@NotNull String name) { storage.saveTrack(track); Logger.info("Created track: " + name); + + // Fire POST event + eventBus.fire(new TrackCreateEvent(track)); + return track; } @Override public CompletableFuture deleteTrack(@NotNull String name) { String lowerName = name.toLowerCase(); + Track track = loadedTracks.get(lowerName); + + // Fire PRE event if track exists + if (track != null) { + TrackDeleteEvent preEvent = new TrackDeleteEvent(track); + eventBus.fire(preEvent); + if (preEvent.isCancelled()) { + return CompletableFuture.completedFuture(null); + } + } + loadedTracks.remove(lowerName); - return storage.deleteTrack(lowerName); + return storage.deleteTrack(lowerName).thenRun(() -> { + // Fire POST event + eventBus.fire(new TrackDeleteEvent(lowerName)); + }); } @Override public CompletableFuture saveTrack(@NotNull Track track) { + // Capture old groups for modify event + Track existing = loadedTracks.get(track.getName()); + List oldGroups = existing != null ? List.copyOf(existing.getGroups()) : List.of(); + loadedTracks.put(track.getName(), track); - return storage.saveTrack(track); + return storage.saveTrack(track).thenRun(() -> { + // Fire modify event if groups changed + List newGroups = track.getGroups(); + if (!oldGroups.equals(newGroups)) { + eventBus.fire(new TrackModifyEvent(track, oldGroups, newGroups)); + } + }); } @Override diff --git a/src/main/java/com/hyperperms/web/ChangeApplier.java b/src/main/java/com/hyperperms/web/ChangeApplier.java index 060069e..be6ab3e 100644 --- a/src/main/java/com/hyperperms/web/ChangeApplier.java +++ b/src/main/java/com/hyperperms/web/ChangeApplier.java @@ -21,12 +21,41 @@ */ public final class ChangeApplier { + private static final ThreadLocal applyingFromWebSocket = ThreadLocal.withInitial(() -> false); + + /** + * Checks if the current thread is applying changes from a WebSocket message. + * EventBus listeners can use this to avoid echo loops. + * + * @return true if changes are being applied from a WebSocket message + */ + public static boolean isApplyingFromWebSocket() { + return applyingFromWebSocket.get(); + } + private final HyperPerms hyperPerms; public ChangeApplier(@NotNull HyperPerms hyperPerms) { this.hyperPerms = hyperPerms; } + /** + * Applies changes originating from a WebSocket message. + * Sets a thread-local flag so EventBus listeners can detect the source + * and avoid echo loops. + * + * @param changes The changes to apply + * @return Result with success/failure counts + */ + public ApplyResult applyChangesFromWebSocket(@NotNull List changes) { + applyingFromWebSocket.set(true); + try { + return applyChanges(changes); + } finally { + applyingFromWebSocket.set(false); + } + } + /** * Result of applying changes. */ diff --git a/src/main/java/com/hyperperms/web/WebEditorService.java b/src/main/java/com/hyperperms/web/WebEditorService.java index 161f1f0..b1684d0 100644 --- a/src/main/java/com/hyperperms/web/WebEditorService.java +++ b/src/main/java/com/hyperperms/web/WebEditorService.java @@ -11,6 +11,7 @@ import com.hyperperms.web.dto.Change; import com.hyperperms.web.dto.SessionCreateResponse; import org.jetbrains.annotations.NotNull; +import org.jetbrains.annotations.Nullable; import java.net.URI; import java.net.URLEncoder; @@ -26,6 +27,8 @@ import java.util.List; import java.util.Map; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; /** * Service for communicating with the web editor API. @@ -37,6 +40,8 @@ public final class WebEditorService { private final HyperPerms hyperPerms; private final HttpClient httpClient; + private final ScheduledExecutorService wsScheduler; + private volatile WebSocketSessionClient activeWsClient; public WebEditorService(@NotNull HyperPerms hyperPerms) { this.hyperPerms = hyperPerms; @@ -44,6 +49,46 @@ public WebEditorService(@NotNull HyperPerms hyperPerms) { .version(HttpClient.Version.HTTP_1_1) .connectTimeout(Duration.ofSeconds(hyperPerms.getConfig().getWebEditorTimeoutSeconds())) .build(); + this.wsScheduler = Executors.newSingleThreadScheduledExecutor(r -> { + Thread t = new Thread(r, "HyperPerms-WebSocket"); + t.setDaemon(true); + return t; + }); + } + + /** + * Connects the WebSocket client for realtime sync. + * + * @param sessionId the session ID + * @param wsUrl the WebSocket URL from the session response + */ + public void connectWebSocket(@NotNull String sessionId, @NotNull String wsUrl) { + disconnectWebSocket(); + WebSocketSessionClient client = new WebSocketSessionClient( + hyperPerms, httpClient, sessionId, wsUrl, wsScheduler); + activeWsClient = client; + client.connect(); + } + + /** + * Disconnects the active WebSocket client. + */ + public void disconnectWebSocket() { + WebSocketSessionClient client = activeWsClient; + if (client != null) { + client.disconnect(); + activeWsClient = null; + } + } + + /** + * Gets the active WebSocket client, if connected. + * + * @return the active client, or null if not connected + */ + @Nullable + public WebSocketSessionClient getActiveWsClient() { + return activeWsClient; } /** @@ -78,7 +123,9 @@ public CompletableFuture createSession(int playerCount) { String sessionId = obj.get("sessionId").getAsString(); String editorUrl = obj.get("url").getAsString(); String expiresAt = obj.has("expiresAt") ? obj.get("expiresAt").getAsString() : ""; - return new SessionCreateResponse(sessionId, editorUrl, expiresAt); + String wsUrl = obj.has("wsUrl") && !obj.get("wsUrl").isJsonNull() + ? obj.get("wsUrl").getAsString() : null; + return new SessionCreateResponse(sessionId, editorUrl, expiresAt, wsUrl); } catch (Exception e) { Logger.warn("Failed to parse session response: " + e.getMessage()); return SessionCreateResponse.error("Invalid response from server"); diff --git a/src/main/java/com/hyperperms/web/WebSocketSessionClient.java b/src/main/java/com/hyperperms/web/WebSocketSessionClient.java new file mode 100644 index 0000000..b365944 --- /dev/null +++ b/src/main/java/com/hyperperms/web/WebSocketSessionClient.java @@ -0,0 +1,724 @@ +package com.hyperperms.web; + +import com.google.gson.Gson; +import com.google.gson.GsonBuilder; +import com.google.gson.JsonArray; +import com.google.gson.JsonObject; +import com.google.gson.JsonParser; +import com.hyperperms.HyperPerms; +import com.hyperperms.api.events.*; +import com.hyperperms.model.Group; +import com.hyperperms.model.Node; +import com.hyperperms.model.User; +import com.hyperperms.util.Logger; +import com.hyperperms.web.dto.Change; +import org.jetbrains.annotations.NotNull; +import org.jetbrains.annotations.Nullable; + +import java.net.URI; +import java.net.http.HttpClient; +import java.net.http.WebSocket; +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionStage; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; + +/** + * WebSocket client for realtime bidirectional sync with the web editor + * via the Cloudflare Durable Object relay. + */ +public final class WebSocketSessionClient { + + private static final Gson GSON = new GsonBuilder().create(); + + private final HyperPerms hyperPerms; + private final HttpClient httpClient; + private final ChangeApplier changeApplier; + private final String sessionId; + private final String wsUrl; + private final ScheduledExecutorService scheduler; + + private final int maxReconnectAttempts; + private final int maxReconnectDelaySeconds; + private final int pingTimeoutSeconds; + + private final AtomicReference webSocket = new AtomicReference<>(); + private final AtomicBoolean connected = new AtomicBoolean(false); + private final AtomicBoolean intentionallyClosed = new AtomicBoolean(false); + private final AtomicInteger reconnectAttempts = new AtomicInteger(0); + + private volatile long lastPingReceived; + private volatile ScheduledFuture pingWatchdog; + + private final List eventSubscriptions = new ArrayList<>(); + + // For batch.apply results + private final Map> pendingBatchApplies = + Collections.synchronizedMap(new HashMap<>()); + private final AtomicInteger batchIdCounter = new AtomicInteger(0); + + // Buffer for partial WebSocket text messages + private final StringBuilder messageBuffer = new StringBuilder(); + + public WebSocketSessionClient( + @NotNull HyperPerms hyperPerms, + @NotNull HttpClient httpClient, + @NotNull String sessionId, + @NotNull String wsUrl, + @NotNull ScheduledExecutorService scheduler + ) { + this.hyperPerms = hyperPerms; + this.httpClient = httpClient; + this.changeApplier = new ChangeApplier(hyperPerms); + this.sessionId = sessionId; + this.wsUrl = wsUrl; + this.scheduler = scheduler; + + this.maxReconnectAttempts = hyperPerms.getConfig().getWebEditorWebsocketReconnectMaxAttempts(); + this.maxReconnectDelaySeconds = hyperPerms.getConfig().getWebEditorWebsocketReconnectMaxDelaySeconds(); + this.pingTimeoutSeconds = hyperPerms.getConfig().getWebEditorWebsocketPingTimeoutSeconds(); + } + + /** + * Connects to the WebSocket relay. + */ + public void connect() { + intentionallyClosed.set(false); + reconnectAttempts.set(0); + doConnect(); + } + + private void doConnect() { + if (intentionallyClosed.get()) return; + + String connectUrl = wsUrl + (wsUrl.contains("?") ? "&" : "?") + "source=plugin"; + Logger.info("[WebSocket] Connecting to: " + connectUrl); + + httpClient.newWebSocketBuilder() + .buildAsync(URI.create(connectUrl), new WsListener()) + .thenAccept(ws -> { + webSocket.set(ws); + connected.set(true); + reconnectAttempts.set(0); + lastPingReceived = System.currentTimeMillis(); + Logger.info("[WebSocket] Connected to session " + sessionId); + subscribeToEvents(); + startPingWatchdog(); + }) + .exceptionally(e -> { + Logger.warn("[WebSocket] Connection failed: " + e.getMessage()); + scheduleReconnect(); + return null; + }); + } + + /** + * Disconnects from the WebSocket relay. + */ + public void disconnect() { + intentionallyClosed.set(true); + connected.set(false); + unsubscribeFromEvents(); + stopPingWatchdog(); + + WebSocket ws = webSocket.getAndSet(null); + if (ws != null) { + try { + ws.sendClose(WebSocket.NORMAL_CLOSURE, "Plugin disconnect"); + } catch (Exception e) { + Logger.debug("[WebSocket] Error during close: " + e.getMessage()); + } + } + + // Fail any pending batch applies + pendingBatchApplies.values().forEach(f -> + f.completeExceptionally(new RuntimeException("WebSocket disconnected"))); + pendingBatchApplies.clear(); + + Logger.info("[WebSocket] Disconnected from session " + sessionId); + } + + public boolean isConnected() { + return connected.get(); + } + + // ==================== Outgoing: Send ops to editor ==================== + + private void sendJson(JsonObject msg) { + WebSocket ws = webSocket.get(); + if (ws != null && connected.get()) { + String json = GSON.toJson(msg); + ws.sendText(json, true); + } + } + + private void sendServerChange(String op, JsonObject data) { + JsonObject msg = new JsonObject(); + msg.addProperty("type", "server.change"); + msg.addProperty("op", op); + msg.add("data", data); + sendJson(msg); + } + + /** + * Sends a batch of changes to the editor for application via WebSocket. + * + * @param changes the changes to apply + * @return a future that completes with the result + */ + public CompletableFuture sendBatchApply(@NotNull List changes) { + String batchId = "batch-" + batchIdCounter.incrementAndGet(); + + JsonObject msg = new JsonObject(); + msg.addProperty("type", "batch.apply"); + msg.addProperty("batchId", batchId); + + JsonArray opsArray = new JsonArray(); + for (Change change : changes) { + opsArray.add(changeToJson(change)); + } + msg.add("ops", opsArray); + + CompletableFuture future = new CompletableFuture<>(); + pendingBatchApplies.put(batchId, future); + + // Timeout after 30s + scheduler.schedule(() -> { + CompletableFuture pending = pendingBatchApplies.remove(batchId); + if (pending != null && !pending.isDone()) { + pending.completeExceptionally(new RuntimeException("batch.apply timed out")); + } + }, 30, TimeUnit.SECONDS); + + sendJson(msg); + return future; + } + + // ==================== Incoming: Handle messages from editor ==================== + + private void handleMessage(String text) { + try { + JsonObject msg = JsonParser.parseString(text).getAsJsonObject(); + String type = msg.has("type") ? msg.get("type").getAsString() : ""; + + switch (type) { + case "ping" -> handlePing(); + case "batch.apply" -> handleBatchApply(msg); + case "session.sync" -> handleSessionSync(msg); + case "apply.result" -> handleApplyResult(msg); + default -> Logger.debug("[WebSocket] Unknown message type: " + type); + } + } catch (Exception e) { + Logger.warn("[WebSocket] Error handling message: " + e.getMessage()); + } + } + + private void handlePing() { + lastPingReceived = System.currentTimeMillis(); + JsonObject pong = new JsonObject(); + pong.addProperty("type", "pong"); + sendJson(pong); + } + + private void handleBatchApply(JsonObject msg) { + String batchId = msg.has("batchId") ? msg.get("batchId").getAsString() : null; + + try { + JsonArray ops = msg.has("ops") ? msg.getAsJsonArray("ops") : new JsonArray(); + List changes = new ArrayList<>(); + + for (var elem : ops) { + Change change = jsonToChange(elem.getAsJsonObject()); + if (change != null) { + changes.add(change); + } + } + + ChangeApplier.ApplyResult result = changeApplier.applyChangesFromWebSocket(changes); + + // Send result back + JsonObject response = new JsonObject(); + response.addProperty("type", "apply.result"); + if (batchId != null) response.addProperty("batchId", batchId); + response.addProperty("success", result.getSuccessCount()); + response.addProperty("failed", result.getFailureCount()); + if (result.hasErrors()) { + JsonArray errors = new JsonArray(); + result.getErrors().forEach(errors::add); + response.add("errors", errors); + } + sendJson(response); + + Logger.info("[WebSocket] Applied batch: " + result.getSuccessCount() + " ok, " + + result.getFailureCount() + " failed"); + } catch (Exception e) { + Logger.warn("[WebSocket] Error applying batch: " + e.getMessage()); + JsonObject response = new JsonObject(); + response.addProperty("type", "apply.result"); + if (batchId != null) response.addProperty("batchId", batchId); + response.addProperty("success", 0); + response.addProperty("failed", 1); + JsonArray errors = new JsonArray(); + errors.add("Internal error: " + e.getMessage()); + response.add("errors", errors); + sendJson(response); + } + } + + private void handleSessionSync(JsonObject msg) { + Logger.debug("[WebSocket] Received session.sync request"); + try { + int playerCount = hyperPerms.getUserManager().getLoadedUsers().size(); + SessionData data = SessionData.fromHyperPerms(hyperPerms, playerCount); + String json = GSON.toJson(data); + + JsonObject response = new JsonObject(); + response.addProperty("type", "session.sync"); + response.add("data", JsonParser.parseString(json)); + sendJson(response); + + Logger.debug("[WebSocket] Sent session.sync response"); + } catch (Exception e) { + Logger.warn("[WebSocket] Error serializing session sync: " + e.getMessage()); + } + } + + private void handleApplyResult(JsonObject msg) { + String batchId = msg.has("batchId") ? msg.get("batchId").getAsString() : null; + if (batchId == null) return; + + CompletableFuture future = pendingBatchApplies.remove(batchId); + if (future != null) { + int success = msg.has("success") ? msg.get("success").getAsInt() : 0; + int failed = msg.has("failed") ? msg.get("failed").getAsInt() : 0; + String undoId = msg.has("undoId") ? msg.get("undoId").getAsString() : null; + future.complete(new BatchApplyResult(success, failed, undoId)); + } + } + + // ==================== Event subscriptions (plugin -> editor) ==================== + + private void subscribeToEvents() { + EventBus eventBus = hyperPerms.getEventBus(); + + eventSubscriptions.add(eventBus.subscribe(PermissionChangeEvent.class, this::onPermissionChange)); + eventSubscriptions.add(eventBus.subscribe(GroupCreateEvent.class, this::onGroupCreate)); + eventSubscriptions.add(eventBus.subscribe(GroupDeleteEvent.class, this::onGroupDelete)); + eventSubscriptions.add(eventBus.subscribe(GroupModifyEvent.class, this::onGroupModify)); + eventSubscriptions.add(eventBus.subscribe(UserGroupChangeEvent.class, this::onUserGroupChange)); + eventSubscriptions.add(eventBus.subscribe(TrackCreateEvent.class, this::onTrackCreate)); + eventSubscriptions.add(eventBus.subscribe(TrackDeleteEvent.class, this::onTrackDelete)); + eventSubscriptions.add(eventBus.subscribe(TrackModifyEvent.class, this::onTrackModify)); + } + + private void unsubscribeFromEvents() { + eventSubscriptions.forEach(EventBus.Subscription::unsubscribe); + eventSubscriptions.clear(); + } + + private boolean shouldSkipEvent() { + return ChangeApplier.isApplyingFromWebSocket() || !connected.get(); + } + + private void onPermissionChange(PermissionChangeEvent event) { + if (shouldSkipEvent()) return; + + String op = switch (event.getChangeType()) { + case ADD -> "permission.add"; + case REMOVE -> "permission.remove"; + case UPDATE -> "permission.set"; + default -> null; + }; + if (op == null) return; + + JsonObject data = new JsonObject(); + data.addProperty("holder", event.getHolder().getIdentifier()); + data.addProperty("holderType", event.getHolder() instanceof Group ? "group" : "user"); + data.addProperty("node", event.getNode().getPermission()); + data.addProperty("value", event.getNode().getValue()); + + if (!event.getNode().getContexts().isEmpty()) { + JsonObject contexts = new JsonObject(); + for (var entry : event.getNode().getContexts().toSet()) { + contexts.addProperty(entry.key(), entry.value()); + } + data.add("contexts", contexts); + } + + if (event.getNode().getExpiry() != null) { + data.addProperty("expiry", event.getNode().getExpiry().toEpochMilli()); + } + + sendServerChange(op, data); + } + + private void onGroupCreate(GroupCreateEvent event) { + if (shouldSkipEvent()) return; + if (event.getState() != GroupCreateEvent.State.POST) return; + + JsonObject data = new JsonObject(); + data.addProperty("name", event.getGroupName()); + sendServerChange("group.create", data); + } + + private void onGroupDelete(GroupDeleteEvent event) { + if (shouldSkipEvent()) return; + if (event.getState() != GroupDeleteEvent.State.POST) return; + + JsonObject data = new JsonObject(); + data.addProperty("name", event.getGroupName()); + sendServerChange("group.delete", data); + } + + private void onGroupModify(GroupModifyEvent event) { + if (shouldSkipEvent()) return; + + switch (event.getModifyType()) { + case PARENT_ADD -> { + JsonObject data = new JsonObject(); + data.addProperty("group", event.getGroup().getName()); + data.addProperty("parent", String.valueOf(event.getNewValue())); + sendServerChange("group.addParent", data); + } + case PARENT_REMOVE -> { + JsonObject data = new JsonObject(); + data.addProperty("group", event.getGroup().getName()); + data.addProperty("parent", String.valueOf(event.getOldValue())); + sendServerChange("group.removeParent", data); + } + default -> { + // Meta changes (weight, prefix, suffix, etc.) - send generic modify + JsonObject data = new JsonObject(); + data.addProperty("group", event.getGroup().getName()); + data.addProperty("property", event.getModifyType().name().toLowerCase()); + if (event.getOldValue() != null) { + data.addProperty("oldValue", String.valueOf(event.getOldValue())); + } + if (event.getNewValue() != null) { + data.addProperty("newValue", String.valueOf(event.getNewValue())); + } + sendServerChange("group.modify", data); + } + } + } + + private void onUserGroupChange(UserGroupChangeEvent event) { + if (shouldSkipEvent()) return; + + String op = switch (event.getChangeType()) { + case ADD -> "user.addGroup"; + case REMOVE -> "user.removeGroup"; + default -> null; + }; + if (op == null) return; + + JsonObject data = new JsonObject(); + data.addProperty("uuid", event.getUuid().toString()); + data.addProperty("group", event.getGroupName()); + sendServerChange(op, data); + } + + private void onTrackCreate(TrackCreateEvent event) { + if (shouldSkipEvent()) return; + if (event.getState() != TrackCreateEvent.State.POST) return; + + JsonObject data = new JsonObject(); + data.addProperty("name", event.getTrackName()); + if (event.getTrack() != null) { + JsonArray groups = new JsonArray(); + event.getTrack().getGroups().forEach(groups::add); + data.add("groups", groups); + } + sendServerChange("track.create", data); + } + + private void onTrackDelete(TrackDeleteEvent event) { + if (shouldSkipEvent()) return; + if (event.getState() != TrackDeleteEvent.State.POST) return; + + JsonObject data = new JsonObject(); + data.addProperty("name", event.getTrackName()); + sendServerChange("track.delete", data); + } + + private void onTrackModify(TrackModifyEvent event) { + if (shouldSkipEvent()) return; + + JsonObject data = new JsonObject(); + data.addProperty("name", event.getTrack().getName()); + JsonArray groups = new JsonArray(); + event.getNewGroups().forEach(groups::add); + data.add("groups", groups); + sendServerChange("track.modify", data); + } + + // ==================== Change serialization ==================== + + private JsonObject changeToJson(Change change) { + JsonObject obj = new JsonObject(); + obj.addProperty("type", change.getType().name().toLowerCase()); + if (change.getTargetType() != null) obj.addProperty("targetType", change.getTargetType()); + if (change.getTarget() != null) obj.addProperty("target", change.getTarget()); + if (change.getNode() != null) obj.addProperty("node", change.getNode()); + if (change.getValue() != null) obj.addProperty("value", change.getValue()); + if (change.getGroupName() != null) obj.addProperty("groupName", change.getGroupName()); + if (change.getParent() != null) obj.addProperty("parent", change.getParent()); + if (change.getTrackName() != null) obj.addProperty("trackName", change.getTrackName()); + if (change.getNewGroups() != null) { + JsonArray arr = new JsonArray(); + change.getNewGroups().forEach(arr::add); + obj.add("newGroups", arr); + } + return obj; + } + + @Nullable + private Change jsonToChange(JsonObject obj) { + String typeStr = obj.has("type") ? obj.get("type").getAsString() : null; + if (typeStr == null) return null; + + // Map op-style types to Change.Type + Change.Type type = switch (typeStr) { + case "permission.add", "permission_added" -> Change.Type.PERMISSION_ADDED; + case "permission.remove", "permission_removed" -> Change.Type.PERMISSION_REMOVED; + case "permission.set", "permission_modified" -> Change.Type.PERMISSION_MODIFIED; + case "group.create", "group_created" -> Change.Type.GROUP_CREATED; + case "group.delete", "group_deleted" -> Change.Type.GROUP_DELETED; + case "group.addParent", "parent_added" -> Change.Type.PARENT_ADDED; + case "group.removeParent", "parent_removed" -> Change.Type.PARENT_REMOVED; + case "group.modify", "meta_changed" -> Change.Type.META_CHANGED; + case "weight_changed" -> Change.Type.WEIGHT_CHANGED; + case "track.create", "track_created" -> Change.Type.TRACK_CREATED; + case "track.delete", "track_deleted" -> Change.Type.TRACK_DELETED; + case "track.modify", "track_modified" -> Change.Type.TRACK_MODIFIED; + default -> null; + }; + if (type == null) { + Logger.warn("[WebSocket] Unknown change type: " + typeStr); + return null; + } + + Change.Builder builder = Change.builder(type); + + if (obj.has("targetType")) builder.targetType(obj.get("targetType").getAsString()); + if (obj.has("target")) builder.target(obj.get("target").getAsString()); + if (obj.has("node")) builder.node(obj.get("node").getAsString()); + if (obj.has("value") && !obj.get("value").isJsonNull()) builder.value(obj.get("value").getAsBoolean()); + if (obj.has("groupName")) builder.groupName(obj.get("groupName").getAsString()); + if (obj.has("parent")) builder.parent(obj.get("parent").getAsString()); + if (obj.has("trackName")) builder.trackName(obj.get("trackName").getAsString()); + if (obj.has("key")) builder.key(obj.get("key").getAsString()); + if (obj.has("newValue") && obj.get("newValue").isJsonPrimitive()) { + if (obj.get("newValue").getAsJsonPrimitive().isString()) { + builder.metaNewValue(obj.get("newValue").getAsString()); + } else if (obj.get("newValue").getAsJsonPrimitive().isBoolean()) { + builder.newValue(obj.get("newValue").getAsBoolean()); + } + } + if (obj.has("newGroups") && obj.get("newGroups").isJsonArray()) { + List groups = new ArrayList<>(); + obj.getAsJsonArray("newGroups").forEach(e -> groups.add(e.getAsString())); + builder.newGroups(groups); + } + if (obj.has("contexts") && obj.get("contexts").isJsonObject()) { + Map contexts = new HashMap<>(); + obj.getAsJsonObject("contexts").entrySet().forEach(e -> + contexts.put(e.getKey(), e.getValue().getAsString())); + builder.contexts(contexts); + } + if (obj.has("expiry") && !obj.get("expiry").isJsonNull()) { + builder.expiry(obj.get("expiry").getAsLong()); + } + if (obj.has("group") && obj.get("group").isJsonObject()) { + JsonObject groupObj = obj.getAsJsonObject("group"); + builder.group(parseGroupData(groupObj)); + } + if (obj.has("track") && obj.get("track").isJsonObject()) { + JsonObject trackObj = obj.getAsJsonObject("track"); + String name = trackObj.has("name") ? trackObj.get("name").getAsString() : ""; + List groups = new ArrayList<>(); + if (trackObj.has("groups") && trackObj.get("groups").isJsonArray()) { + trackObj.getAsJsonArray("groups").forEach(e -> groups.add(e.getAsString())); + } + builder.track(new Change.TrackData(name, groups)); + } + + return builder.build(); + } + + private Change.GroupData parseGroupData(JsonObject obj) { + String name = obj.has("name") ? obj.get("name").getAsString() : ""; + String displayName = obj.has("displayName") ? obj.get("displayName").getAsString() : null; + int weight = obj.has("weight") ? obj.get("weight").getAsInt() : 0; + String prefix = obj.has("prefix") ? obj.get("prefix").getAsString() : null; + String suffix = obj.has("suffix") ? obj.get("suffix").getAsString() : null; + int prefixPriority = obj.has("prefixPriority") ? obj.get("prefixPriority").getAsInt() : 0; + int suffixPriority = obj.has("suffixPriority") ? obj.get("suffixPriority").getAsInt() : 0; + + List perms = new ArrayList<>(); + if (obj.has("permissions") && obj.get("permissions").isJsonArray()) { + for (var elem : obj.getAsJsonArray("permissions")) { + JsonObject p = elem.getAsJsonObject(); + String node = p.has("node") ? p.get("node").getAsString() : null; + if (node == null) continue; + boolean value = !p.has("value") || p.get("value").getAsBoolean(); + Map ctx = Collections.emptyMap(); + if (p.has("contexts") && p.get("contexts").isJsonObject()) { + ctx = new HashMap<>(); + for (var e : p.getAsJsonObject("contexts").entrySet()) { + ctx.put(e.getKey(), e.getValue().getAsString()); + } + } + Long expiry = p.has("expiry") && !p.get("expiry").isJsonNull() + ? p.get("expiry").getAsLong() : null; + perms.add(new Change.PermissionNode(node, value, ctx, expiry)); + } + } + + List parents = new ArrayList<>(); + if (obj.has("parents") && obj.get("parents").isJsonArray()) { + obj.getAsJsonArray("parents").forEach(e -> parents.add(e.getAsString())); + } + + return new Change.GroupData(name, displayName, weight, prefix, suffix, + prefixPriority, suffixPriority, perms, parents); + } + + // ==================== Reconnection ==================== + + private void scheduleReconnect() { + if (intentionallyClosed.get()) return; + + int attempt = reconnectAttempts.incrementAndGet(); + if (attempt > maxReconnectAttempts) { + Logger.warn("[WebSocket] Max reconnect attempts (" + maxReconnectAttempts + ") reached. Giving up."); + connected.set(false); + return; + } + + // Exponential backoff: 1s, 2s, 4s, 8s... capped at maxReconnectDelaySeconds + long delaySec = Math.min((long) Math.pow(2, attempt - 1), maxReconnectDelaySeconds); + Logger.info("[WebSocket] Reconnecting in " + delaySec + "s (attempt " + attempt + "/" + maxReconnectAttempts + ")"); + + scheduler.schedule(this::doConnect, delaySec, TimeUnit.SECONDS); + } + + // ==================== Ping watchdog ==================== + + private void startPingWatchdog() { + stopPingWatchdog(); + pingWatchdog = scheduler.scheduleAtFixedRate(() -> { + if (!connected.get() || intentionallyClosed.get()) return; + + long elapsed = System.currentTimeMillis() - lastPingReceived; + if (elapsed > pingTimeoutSeconds * 1000L) { + Logger.warn("[WebSocket] No ping for " + (elapsed / 1000) + "s, reconnecting..."); + WebSocket ws = webSocket.getAndSet(null); + connected.set(false); + if (ws != null) { + try { + ws.sendClose(WebSocket.NORMAL_CLOSURE, "Ping timeout"); + } catch (Exception ignored) {} + } + unsubscribeFromEvents(); + scheduleReconnect(); + } + }, 30, 30, TimeUnit.SECONDS); + } + + private void stopPingWatchdog() { + if (pingWatchdog != null) { + pingWatchdog.cancel(false); + pingWatchdog = null; + } + } + + // ==================== WebSocket listener ==================== + + private class WsListener implements WebSocket.Listener { + + @Override + public void onOpen(WebSocket webSocket) { + Logger.debug("[WebSocket] Connection opened"); + webSocket.request(1); + } + + @Override + public CompletionStage onText(WebSocket webSocket, CharSequence data, boolean last) { + messageBuffer.append(data); + if (last) { + String fullMessage = messageBuffer.toString(); + messageBuffer.setLength(0); + handleMessage(fullMessage); + } + webSocket.request(1); + return null; + } + + @Override + public CompletionStage onPing(WebSocket webSocket, ByteBuffer message) { + lastPingReceived = System.currentTimeMillis(); + webSocket.request(1); + return null; + } + + @Override + public CompletionStage onClose(WebSocket webSocket, int statusCode, String reason) { + Logger.info("[WebSocket] Connection closed: " + statusCode + " " + reason); + connected.set(false); + WebSocketSessionClient.this.webSocket.set(null); + unsubscribeFromEvents(); + stopPingWatchdog(); + + if (!intentionallyClosed.get()) { + scheduleReconnect(); + } + return null; + } + + @Override + public void onError(WebSocket webSocket, Throwable error) { + Logger.warn("[WebSocket] Error: " + error.getMessage()); + connected.set(false); + WebSocketSessionClient.this.webSocket.set(null); + unsubscribeFromEvents(); + stopPingWatchdog(); + + if (!intentionallyClosed.get()) { + scheduleReconnect(); + } + } + } + + // ==================== Result types ==================== + + /** + * Result of a batch apply operation sent over WebSocket. + */ + public static final class BatchApplyResult { + private final int successCount; + private final int failedCount; + private final String undoId; + + public BatchApplyResult(int successCount, int failedCount, @Nullable String undoId) { + this.successCount = successCount; + this.failedCount = failedCount; + this.undoId = undoId; + } + + public int getSuccessCount() { return successCount; } + public int getFailedCount() { return failedCount; } + @Nullable public String getUndoId() { return undoId; } + public boolean hasErrors() { return failedCount > 0; } + } +} diff --git a/src/main/java/com/hyperperms/web/dto/SessionCreateResponse.java b/src/main/java/com/hyperperms/web/dto/SessionCreateResponse.java index a6443f5..f1974e4 100644 --- a/src/main/java/com/hyperperms/web/dto/SessionCreateResponse.java +++ b/src/main/java/com/hyperperms/web/dto/SessionCreateResponse.java @@ -11,16 +11,19 @@ public final class SessionCreateResponse { private final String sessionId; private final String url; private final String expiresAt; + private final String wsUrl; private final String error; public SessionCreateResponse( @NotNull String sessionId, @NotNull String url, - @NotNull String expiresAt + @NotNull String expiresAt, + @Nullable String wsUrl ) { this.sessionId = sessionId; this.url = url; this.expiresAt = expiresAt; + this.wsUrl = wsUrl; this.error = null; } @@ -28,6 +31,7 @@ private SessionCreateResponse(@NotNull String error) { this.sessionId = null; this.url = null; this.expiresAt = null; + this.wsUrl = null; this.error = error; } @@ -53,6 +57,16 @@ public String getExpiresAt() { return expiresAt; } + /** + * Gets the WebSocket URL for realtime sync. + * + * @return the WebSocket URL, or null if not available + */ + @Nullable + public String getWsUrl() { + return wsUrl; + } + @Nullable public String getError() { return error;