From 4d75e6b75c9a2df309ea50cee6a3db8d29aa4760 Mon Sep 17 00:00:00 2001 From: Joshua Castle <26531652+Kas-tle@users.noreply.github.com> Date: Mon, 16 Feb 2026 18:27:05 -0800 Subject: [PATCH 1/4] Add support for JSON-RPC signaling used by realms Signed-off-by: Joshua Castle <26531652+Kas-tle@users.noreply.github.com> --- .../channel/nethernet/NetherNetConstants.java | 12 + .../AbstractNetherNetXboxSignaling.java | 265 +++++++++++++ .../signaling/NetherNetXboxRpcSignaling.java | 212 +++++++++++ .../signaling/NetherNetXboxSignaling.java | 351 ++---------------- 4 files changed, 524 insertions(+), 316 deletions(-) create mode 100644 transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/signaling/AbstractNetherNetXboxSignaling.java create mode 100644 transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/signaling/NetherNetXboxRpcSignaling.java diff --git a/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/NetherNetConstants.java b/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/NetherNetConstants.java index abf5899..dcc6a0f 100644 --- a/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/NetherNetConstants.java +++ b/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/NetherNetConstants.java @@ -27,6 +27,9 @@ public class NetherNetConstants { public static final String RTC_NEGOTIATION_CANDIDATE_ADD = "CANDIDATEADD"; public static final String RTC_NEGOTIATION_CONNECT_ERROR = "CONNECTERROR"; + // Signaling User Agent String + public static final String SIGNALING_USER_AGENT = "libHttpClient/1.0.0.0"; + // Xbox Signaling Message Types public static final int XBOX_SIGNAL_NOT_FOUND = 0; public static final int XBOX_SIGNAL_SIGNAL = 1; @@ -34,6 +37,15 @@ public class NetherNetConstants { public static final int XBOX_SIGNAL_ACCEPTED = 3; public static final int XBOX_SIGNAL_ACK = 4; + // Xbox JSON-RPC Signaling Method Names + public static final String XBOX_RPC_METHOD_TURN_AUTH = "Signaling_TurnAuth_v1_0"; + public static final String XBOX_RPC_METHOD_SEND_MESSAGE = "Signaling_SendClientMessage_v1_0"; + public static final String XBOX_RPC_METHOD_RECEIVE_MESSAGE = "Signaling_ReceiveMessage_v1_0"; + public static final String XBOX_RPC_METHOD_PING = "System_Ping_v1_0"; + public static final String XBOX_RPC_METHOD_PONG = "System_Pong_v1_0"; + public static final String XBOX_RPC_INNER_METHOD_WEBRTC = "Signaling_WebRtc_v1_0"; + public static final String XBOX_RPC_INNER_METHOD_DELIVERY = "Signaling_DeliveryNotification_V1_0"; + // SCTP Constants public static final int MAX_SCTP_MESSAGE_SIZE = 10000; public static final String RELIABLE_CHANNEL_LABEL = "ReliableDataChannel"; diff --git a/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/signaling/AbstractNetherNetXboxSignaling.java b/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/signaling/AbstractNetherNetXboxSignaling.java new file mode 100644 index 0000000..782d66c --- /dev/null +++ b/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/signaling/AbstractNetherNetXboxSignaling.java @@ -0,0 +1,265 @@ +package dev.kastle.netty.channel.nethernet.signaling; + +import com.google.gson.JsonArray; +import com.google.gson.JsonElement; +import com.google.gson.JsonObject; +import dev.kastle.netty.channel.nethernet.NetherNetConstants; +import io.netty.bootstrap.Bootstrap; +import io.netty.channel.socket.SocketChannel; +import io.netty.channel.socket.nio.NioSocketChannel; +import io.netty.channel.Channel; +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.ChannelInitializer; +import io.netty.channel.ChannelPipeline; +import io.netty.channel.EventLoopGroup; +import io.netty.channel.SimpleChannelInboundHandler; +import io.netty.channel.nio.NioEventLoopGroup; +import io.netty.handler.codec.http.DefaultHttpHeaders; +import io.netty.handler.codec.http.HttpClientCodec; +import io.netty.handler.codec.http.HttpObjectAggregator; +import io.netty.handler.codec.http.websocketx.TextWebSocketFrame; +import io.netty.handler.codec.http.websocketx.WebSocketClientHandshaker; +import io.netty.handler.codec.http.websocketx.WebSocketClientHandshakerFactory; +import io.netty.handler.codec.http.websocketx.WebSocketClientProtocolHandler; +import io.netty.handler.codec.http.websocketx.WebSocketVersion; +import io.netty.handler.ssl.SslContext; +import io.netty.handler.ssl.SslContextBuilder; +import io.netty.util.internal.logging.InternalLogger; +import io.netty.util.internal.logging.InternalLoggerFactory; + +import java.net.ConnectException; +import java.net.SocketAddress; +import java.net.URI; +import java.nio.channels.ClosedChannelException; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentHashMap; + +public abstract class AbstractNetherNetXboxSignaling extends SimpleChannelInboundHandler + implements NetherNetClientSignaling, NetherNetServerSignaling { + + protected final InternalLogger log = InternalLoggerFactory.getInstance(getClass()); + + protected final String xboxToken; + protected final String localNetworkId; + protected final URI uri; + protected final EventLoopGroup eventLoopGroup; + + protected Channel channel; + protected CompletableFuture> connectFuture; + protected volatile List iceServers = new ArrayList<>(); + + protected final Map handlers = new ConcurrentHashMap<>(); + protected NetherNetServerSignaling.NewConnectionHandler newConnectionHandler; + protected volatile NetherNetClientSignaling.NotFoundHandler notFoundHandler; + + protected AbstractNetherNetXboxSignaling(String localNetworkId, String xboxToken, URI uri) { + this.localNetworkId = localNetworkId; + this.xboxToken = xboxToken; + this.uri = uri; + this.eventLoopGroup = new NioEventLoopGroup(1); + } + + @Override + public String getLocalNetworkId() { + return this.localNetworkId; + } + + @Override + public synchronized CompletableFuture> connect(SocketAddress remoteAddress) { + return connectInternal(); + } + + @Override + public void bind(SocketAddress localAddress) throws ConnectException { + try { + connectInternal().join(); + } catch (Exception e) { + Throwable cause = e.getCause() != null ? e.getCause() : e; + close(); + if (cause instanceof ConnectException) throw (ConnectException) cause; + ConnectException ce = new ConnectException("Failed to connect to Xbox Signaling: " + cause.getMessage()); + ce.initCause(cause); + throw ce; + } + } + + protected synchronized CompletableFuture> connectInternal() { + if (connectFuture != null) return connectFuture; + + connectFuture = new CompletableFuture<>(); + connectFuture.thenAccept(servers -> this.iceServers = servers); + + try { + SslContext sslCtx = SslContextBuilder.forClient().build(); + WebSocketClientHandshaker handshaker = WebSocketClientHandshakerFactory.newHandshaker( + uri, WebSocketVersion.V13, null, false, + new DefaultHttpHeaders() + .add("Authorization", xboxToken) + .add("User-Agent", NetherNetConstants.SIGNALING_USER_AGENT) + .add("session-id", UUID.randomUUID().toString()) + .add("request-id", UUID.randomUUID().toString()) + ); + + Bootstrap b = new Bootstrap(); + b.group(eventLoopGroup) + .channel(NioSocketChannel.class) + .handler(new ChannelInitializer() { + @Override + protected void initChannel(SocketChannel ch) { + ChannelPipeline p = ch.pipeline(); + p.addLast(sslCtx.newHandler(ch.alloc(), uri.getHost(), 443)); + p.addLast(new HttpClientCodec(), new HttpObjectAggregator(8192)); + p.addLast("ws-handshake", new WebSocketClientProtocolHandler(handshaker)); + p.addLast("handler", AbstractNetherNetXboxSignaling.this); + } + }); + + this.channel = b.connect(uri.getHost(), 443).sync().channel(); + } catch (Exception e) { + Throwable cause = e.getCause() != null ? e.getCause() : e; + if (connectFuture != null) connectFuture.completeExceptionally(cause); + } + return connectFuture; + } + + @Override + public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception { + if (evt == WebSocketClientProtocolHandler.ClientHandshakeStateEvent.HANDSHAKE_COMPLETE) { + log.debug("{} WebSocket Connected", getClass().getSimpleName()); + onConnected(ctx); + } else { + super.userEventTriggered(ctx, evt); + } + } + + /** + * Called when the WebSocket handshake is complete. + */ + protected abstract void onConnected(ChannelHandlerContext ctx); + + @Override + public List getIceServers() { + return this.iceServers; + } + + @Override + public void setNewConnectionHandler(NetherNetServerSignaling.NewConnectionHandler handler) { + this.newConnectionHandler = handler; + } + + @Override + public void setNotFoundHandler(NotFoundHandler handler) { + this.notFoundHandler = handler; + } + + @Override + public void setSignalHandler(long connectionId, SignalHandler handler) { + this.handlers.put(connectionId, handler); + } + + @Override + public void removeSignalHandler(long connectionId) { + this.handlers.remove(connectionId); + } + + @Override + public void setAdvertisementData(PongData pongData) { + // No-op for Xbox Signaling. + } + + @Override + public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception { + if (connectFuture != null && !connectFuture.isDone()) { + connectFuture.completeExceptionally(cause); + } + log.error("Signaling Exception: {}", cause.getMessage(), cause); + ctx.close(); + } + + @Override + public void channelInactive(ChannelHandlerContext ctx) throws Exception { + synchronized (this) { + if (connectFuture != null && !connectFuture.isDone()) { + connectFuture.completeExceptionally(new ClosedChannelException()); + } + connectFuture = null; + this.channel = null; + } + super.channelInactive(ctx); + } + + @Override + public void close() { + if (channel != null) channel.close(); + eventLoopGroup.shutdownGracefully(); + } + + protected void dispatchSignalToPipeline(String sender, String rawMsg) { + try { + // Signal Format: + String[] parts = rawMsg.split(" ", 3); + if (parts.length < 2) return; + + long connectionId = Long.parseUnsignedLong(parts[1]); + + SignalHandler handler = handlers.get(connectionId); + if (handler != null) { + handler.onSignal(rawMsg); + return; + } + + if (NetherNetConstants.RTC_NEGOTIATION_CONNECT_REQUEST.equals(parts[0]) && newConnectionHandler != null) { + String payload = parts.length > 2 ? parts[2] : ""; + newConnectionHandler.onConnect(connectionId, sender, payload); + } else { + log.debug("No handler found for connection ID: {} (Type: {})", connectionId, parts[0]); + } + } catch (Exception e) { + log.error("Failed to dispatch signal: {}", rawMsg, e); + } + } + + protected List parseTurnServers(JsonObject json) { + List result = new ArrayList<>(); + try { + JsonArray servers = null; + if (json.has("TurnAuthServers")) servers = json.getAsJsonArray("TurnAuthServers"); + else if (json.has("turnAuthServers")) servers = json.getAsJsonArray("turnAuthServers"); + + if (servers != null) { + for (JsonElement el : servers) { + JsonObject server = el.getAsJsonObject(); + List urls = new ArrayList<>(); + + JsonArray urlsArray = null; + if (server.has("Urls")) urlsArray = server.getAsJsonArray("Urls"); + else if (server.has("urls")) urlsArray = server.getAsJsonArray("urls"); + + if (urlsArray != null) { + urlsArray.forEach(u -> urls.add(u.getAsString())); + + IceServerInfo.Builder info = new IceServerInfo.Builder().setUrls(urls); + + if (server.has("Username")) info.setUsername(server.get("Username").getAsString()); + else if (server.has("username")) info.setUsername(server.get("username").getAsString()); + + if (server.has("Password")) info.setPassword(server.get("Password").getAsString()); + else if (server.has("password")) info.setPassword(server.get("password").getAsString()); + else if (server.has("Credential")) info.setPassword(server.get("Credential").getAsString()); + else if (server.has("credential")) info.setPassword(server.get("credential").getAsString()); + + result.add(info.build()); + } + } + } + } catch (Exception e) { + log.error("Failed to parse TURN servers", e); + } + log.debug("Successfully parsed {} ICE servers.", result.size()); + return result; + } +} diff --git a/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/signaling/NetherNetXboxRpcSignaling.java b/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/signaling/NetherNetXboxRpcSignaling.java new file mode 100644 index 0000000..010a14e --- /dev/null +++ b/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/signaling/NetherNetXboxRpcSignaling.java @@ -0,0 +1,212 @@ +package dev.kastle.netty.channel.nethernet.signaling; + +import com.google.gson.Gson; +import com.google.gson.GsonBuilder; +import com.google.gson.JsonArray; +import com.google.gson.JsonElement; +import com.google.gson.JsonObject; +import com.google.gson.JsonParser; +import dev.kastle.netty.channel.nethernet.NetherNetConstants; +import io.netty.channel.ChannelHandlerContext; +import io.netty.handler.codec.http.websocketx.TextWebSocketFrame; + +import java.net.URI; +import java.nio.channels.ClosedChannelException; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ThreadLocalRandom; +import java.util.concurrent.TimeUnit; + +public class NetherNetXboxRpcSignaling extends AbstractNetherNetXboxSignaling { + private static final Gson gson = new GsonBuilder().serializeNulls().create(); + private final Map> pendingRequests = new ConcurrentHashMap<>(); + + /** + * Creates a NetherNetXboxRpcSignaling instance. + * + * @param networkId The Network ID to use. + * @param xboxToken The Minecraft Bedrock Session authorization header ('MCToken ***'). + */ + public NetherNetXboxRpcSignaling(String localNetworkId, String xboxToken) { + super(localNetworkId, xboxToken, URI.create("wss://signal.franchise.minecraft-services.net/ws/v1.0/messaging/connect")); + } + + /** + * Creates a NetherNetXboxRpcSignaling instance. + * + * @param localNetworkId The local Network ID to use. + * @param xboxToken The Minecraft Bedrock Session authorization header ('MCToken ***'). + */ + public NetherNetXboxRpcSignaling(long localNetworkId, String xboxToken) { + this(Long.toUnsignedString(localNetworkId), xboxToken); + } + + /** + * Creates a NetherNetXboxRpcSignaling instance with a random local Network ID. + * + * @param xboxToken The Minecraft Bedrock Session authorization header ('MCToken ***'). + */ + public NetherNetXboxRpcSignaling(String xboxToken) { + this(Long.toUnsignedString(ThreadLocalRandom.current().nextLong(1, Long.MAX_VALUE)), xboxToken); + } + + @Override + protected void onConnected(ChannelHandlerContext ctx) { + ctx.executor().scheduleAtFixedRate(() -> { + if (channel != null && channel.isActive()) { + sendJsonRpcRequest(NetherNetConstants.XBOX_RPC_METHOD_PING, new JsonObject()); + } + }, 30, 50, TimeUnit.SECONDS); + + sendJsonRpcRequest(NetherNetConstants.XBOX_RPC_METHOD_TURN_AUTH, new JsonObject()) + .thenAccept(response -> { + List servers = parseTurnServers(response); + if (connectFuture != null && !connectFuture.isDone()) connectFuture.complete(servers); + }) + .exceptionally(t -> { + log.error("Failed to fetch TURN credentials", t); + if (connectFuture != null && !connectFuture.isDone()) connectFuture.completeExceptionally(t); + return null; + }); + } + + @Override + protected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame frame) { + String text = frame.text(); + try { + JsonObject json = JsonParser.parseString(text).getAsJsonObject(); + + if (json.has("result") || (json.has("error") && json.has("id"))) { + handleResponse(json); + } else if (json.has("method")) { + handleRequest(json); + } + } catch (Exception e) { + log.error("Error processing signaling frame: " + text, e); + } + } + + private void handleResponse(JsonObject json) { + if (!json.has("id") || json.get("id").isJsonNull()) return; + String id = json.get("id").getAsString(); + CompletableFuture future = pendingRequests.remove(id); + + if (future != null) { + if (json.has("error") && !json.get("error").isJsonNull()) { + JsonObject error = json.getAsJsonObject("error"); + String msg = error.has("message") ? error.get("message").getAsString() : error.toString(); + + boolean isNotFound = msg.contains("Player not registered"); + if (!isNotFound && error.has("data") && error.get("data").isJsonObject()) { + JsonObject data = error.getAsJsonObject("data"); + if (data.has("Code") && "MissingOrExpiredIdentity".equals(data.get("Code").getAsString())) { + isNotFound = true; + } + } + + if (isNotFound && notFoundHandler != null) { + notFoundHandler.onNotFound(msg); + } + future.completeExceptionally(new RuntimeException(msg)); + } else { + future.complete(json.has("result") && !json.get("result").isJsonNull() ? json.getAsJsonObject("result") : new JsonObject()); + } + } + } + + private void handleRequest(JsonObject json) { + String method = json.get("method").getAsString(); + JsonElement id = json.get("id"); + + switch (method) { + case NetherNetConstants.XBOX_RPC_METHOD_RECEIVE_MESSAGE -> { + if (id != null) sendJsonRpcResult(id, null); + JsonArray params = json.getAsJsonArray("params"); + if (params != null) { + for (JsonElement el : params) processIncomingMessage(el.getAsJsonObject()); + } + } + case NetherNetConstants.XBOX_RPC_METHOD_PONG, NetherNetConstants.XBOX_RPC_METHOD_PING -> { + if (id != null) sendJsonRpcResult(id, null); + } + } + } + + private void processIncomingMessage(JsonObject msgObj) { + String from = msgObj.get("From").getAsString(); + String rawInner = msgObj.get("Message").getAsString(); + String msgId = msgObj.has("Id") ? msgObj.get("Id").getAsString() : UUID.randomUUID().toString(); + + JsonObject innerParams = new JsonObject(); + innerParams.addProperty("messageId", msgId); + JsonObject innerMsg = new JsonObject(); + innerMsg.add("params", innerParams); + innerMsg.addProperty("jsonrpc", "2.0"); + innerMsg.addProperty("method", NetherNetConstants.XBOX_RPC_INNER_METHOD_DELIVERY); + sendJsonRpcRequest(NetherNetConstants.XBOX_RPC_METHOD_SEND_MESSAGE, createSendParams(from, innerMsg.toString())); + + try { + JsonObject innerJson = JsonParser.parseString(rawInner).getAsJsonObject(); + if (innerJson.has("method") && NetherNetConstants.XBOX_RPC_INNER_METHOD_WEBRTC.equals(innerJson.get("method").getAsString())) { + String payload = innerJson.getAsJsonObject("params").get("message").getAsString(); + dispatchSignalToPipeline(from, payload); + } + } catch (Exception e) { + log.error("Failed to parse inner signaling message from " + from, e); + } + } + + @Override + public void sendSignal(String targetNetworkId, String data) { + if (channel == null || !channel.isActive()) throw new IllegalStateException("Signaling channel is not active"); + + JsonObject innerParams = new JsonObject(); + innerParams.addProperty("netherNetId", localNetworkId); + innerParams.addProperty("message", data); + + JsonObject innerMsg = new JsonObject(); + innerMsg.add("params", innerParams); + innerMsg.addProperty("jsonrpc", "2.0"); + innerMsg.addProperty("method", NetherNetConstants.XBOX_RPC_INNER_METHOD_WEBRTC); + + sendJsonRpcRequest(NetherNetConstants.XBOX_RPC_METHOD_SEND_MESSAGE, createSendParams(targetNetworkId, innerMsg.toString())); + } + + private JsonObject createSendParams(String toPlayerId, String message) { + JsonObject params = new JsonObject(); + params.addProperty("toPlayerId", toPlayerId); + params.addProperty("messageId", UUID.randomUUID().toString()); + params.addProperty("message", message); + return params; + } + + private CompletableFuture sendJsonRpcRequest(String method, JsonObject params) { + String id = UUID.randomUUID().toString(); + JsonObject rpc = new JsonObject(); + rpc.add("params", params); + rpc.addProperty("jsonrpc", "2.0"); + rpc.addProperty("method", method); + rpc.addProperty("id", id); + + CompletableFuture future = new CompletableFuture<>(); + pendingRequests.put(id, future); + + if (channel != null && channel.isActive()) { + channel.writeAndFlush(new TextWebSocketFrame(gson.toJson(rpc))); + } else { + future.completeExceptionally(new ClosedChannelException()); + } + return future; + } + + private void sendJsonRpcResult(JsonElement id, JsonElement result) { + JsonObject response = new JsonObject(); + response.add("id", id); + response.add("result", result); + response.addProperty("jsonrpc", "2.0"); + if (channel != null && channel.isActive()) channel.writeAndFlush(new TextWebSocketFrame(gson.toJson(response))); + } +} \ No newline at end of file diff --git a/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/signaling/NetherNetXboxSignaling.java b/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/signaling/NetherNetXboxSignaling.java index 419722e..8e00d65 100644 --- a/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/signaling/NetherNetXboxSignaling.java +++ b/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/signaling/NetherNetXboxSignaling.java @@ -1,66 +1,21 @@ package dev.kastle.netty.channel.nethernet.signaling; import com.google.gson.Gson; -import com.google.gson.JsonArray; -import com.google.gson.JsonElement; import com.google.gson.JsonObject; +import com.google.gson.JsonParser; import dev.kastle.netty.channel.nethernet.NetherNetConstants; -import io.netty.bootstrap.Bootstrap; -import io.netty.channel.Channel; -import io.netty.channel.ChannelHandlerContext; -import io.netty.channel.ChannelInitializer; -import io.netty.channel.ChannelPipeline; -import io.netty.channel.EventLoopGroup; -import io.netty.channel.SimpleChannelInboundHandler; import io.netty.channel.ChannelHandler.Sharable; -import io.netty.channel.nio.NioEventLoopGroup; -import io.netty.channel.socket.SocketChannel; -import io.netty.channel.socket.nio.NioSocketChannel; -import io.netty.handler.codec.http.DefaultHttpHeaders; -import io.netty.handler.codec.http.HttpClientCodec; -import io.netty.handler.codec.http.HttpObjectAggregator; +import io.netty.channel.ChannelHandlerContext; import io.netty.handler.codec.http.websocketx.TextWebSocketFrame; -import io.netty.handler.codec.http.websocketx.WebSocketClientHandshaker; -import io.netty.handler.codec.http.websocketx.WebSocketClientHandshakerFactory; -import io.netty.handler.codec.http.websocketx.WebSocketClientProtocolHandler; -import io.netty.handler.codec.http.websocketx.WebSocketVersion; -import io.netty.handler.ssl.SslContext; -import io.netty.handler.ssl.SslContextBuilder; -import io.netty.util.internal.logging.InternalLogger; -import io.netty.util.internal.logging.InternalLoggerFactory; -import java.net.ConnectException; -import java.net.SocketAddress; import java.net.URI; -import java.nio.channels.ClosedChannelException; -import java.util.ArrayList; -import java.util.List; -import java.util.Map; -import java.util.UUID; -import java.util.concurrent.CompletableFuture; -import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.TimeUnit; @Sharable -public class NetherNetXboxSignaling extends SimpleChannelInboundHandler implements NetherNetClientSignaling, NetherNetServerSignaling { - private static final InternalLogger log = InternalLoggerFactory.getInstance(NetherNetXboxSignaling.class); +public class NetherNetXboxSignaling extends AbstractNetherNetXboxSignaling { private static final Gson gson = new Gson(); - private final String xboxToken; - private final String localNetworkId; - private final URI uri; - private final EventLoopGroup eventLoopGroup; - - private Channel channel; - private CompletableFuture> connectFuture; - - private final Map handlers = new ConcurrentHashMap<>(); - private NetherNetServerSignaling.NewConnectionHandler newConnectionHandler; - private volatile NetherNetClientSignaling.NotFoundHandler notFoundHandler; - - private volatile List iceServers = new ArrayList<>(); - /** * Creates a NetherNetXboxSignaling instance. * @@ -68,10 +23,7 @@ public class NetherNetXboxSignaling extends SimpleChannelInboundHandler> connect(SocketAddress remoteAddress) { - // SocketAddress is ignored for Xbox Signaling Service connection - return connectInternal(); - } - - @Override - public void bind(SocketAddress localAddress) throws ConnectException { - try { - connectInternal().join(); - } catch (Exception e) { - Throwable cause = e.getCause() != null ? e.getCause() : e; - close(); - - if (cause instanceof ConnectException) { - throw (ConnectException) cause; - } - - ConnectException ce = new ConnectException("Failed to connect to Xbox Signaling: " + cause.getMessage()); - ce.initCause(cause); - throw ce; - } - } - - private synchronized CompletableFuture> connectInternal() { - if (connectFuture != null) { - return connectFuture; - } - - connectFuture = new CompletableFuture<>(); - connectFuture.thenAccept(servers -> this.iceServers = servers); - - try { - SslContext sslCtx = SslContextBuilder.forClient().build(); - WebSocketClientHandshaker handshaker = WebSocketClientHandshakerFactory.newHandshaker( - uri, WebSocketVersion.V13, null, false, - new DefaultHttpHeaders() - .add("Authorization", xboxToken) - .add("Session-Id", UUID.randomUUID().toString()) - .add("Request-Id", UUID.randomUUID().toString()) - ); - - Bootstrap b = new Bootstrap(); - b.group(eventLoopGroup) - .channel(NioSocketChannel.class) - .handler(new ChannelInitializer() { - @Override - protected void initChannel(SocketChannel ch) { - ChannelPipeline p = ch.pipeline(); - p.addLast(sslCtx.newHandler(ch.alloc(), uri.getHost(), 443)); - p.addLast(new HttpClientCodec(), new HttpObjectAggregator(8192)); - p.addLast("ws-handshake", new WebSocketClientProtocolHandler(handshaker)); - p.addLast("handler", NetherNetXboxSignaling.this); - } - }); - - this.channel = b.connect(uri.getHost(), 443).sync().channel(); - } catch (Exception e) { - Throwable cause = e.getCause() != null ? e.getCause() : e; - if (cause instanceof ConnectException) { - connectFuture.completeExceptionally(cause); - } else { - ConnectException ce = new ConnectException("Failed to connect to Xbox Signaling: " + cause.getMessage()); - ce.initCause(cause); - connectFuture.completeExceptionally(ce); - } - } - return connectFuture; - } - - public List getIceServers() { - return this.iceServers; - } - - @Override - public void setNewConnectionHandler(NetherNetServerSignaling.NewConnectionHandler handler) { - this.newConnectionHandler = handler; - } - - @Override - public void setNotFoundHandler(NotFoundHandler handler) { - this.notFoundHandler = handler; - } - - @Override - public void setAdvertisementData(PongData pongData) { - // No-op for Xbox Signaling. - // Advertisement is handled via the Session Directory service (PUT /session/...). - } - - @Override - public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception { - if (evt == WebSocketClientProtocolHandler.ClientHandshakeStateEvent.HANDSHAKE_COMPLETE) { - log.debug("NetherNet Signaling WebSocket Connected"); - startPingLoop(ctx); - } else { - super.userEventTriggered(ctx, evt); - } + protected void onConnected(ChannelHandlerContext ctx) { + ctx.executor().scheduleAtFixedRate(() -> { + JsonObject ping = new JsonObject(); + ping.addProperty("Type", 0); + ctx.writeAndFlush(new TextWebSocketFrame(gson.toJson(ping))); + }, 5, 5, TimeUnit.SECONDS); } @Override @@ -203,18 +59,33 @@ protected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame frame) String text = frame.text(); try { JsonObject json = gson.fromJson(text, JsonObject.class); - if (!json.has("Type")) { - log.debug("Ignored message without Type: {}", text); - return; - } + if (!json.has("Type")) return; int type = json.get("Type").getAsInt(); switch (type) { - case NetherNetConstants.XBOX_SIGNAL_NOT_FOUND -> handleNotFound(json, text); - case NetherNetConstants.XBOX_SIGNAL_SIGNAL -> handleSignal(json); - case NetherNetConstants.XBOX_SIGNAL_CREDENTIALS -> handleCredentials(json, text); - case NetherNetConstants.XBOX_SIGNAL_ACCEPTED -> log.trace("Signal Accepted: {}", text); - case NetherNetConstants.XBOX_SIGNAL_ACK -> log.trace("Delivery Ack: {}", text); + case NetherNetConstants.XBOX_SIGNAL_NOT_FOUND -> { + log.debug("Peer Not Found: {}", text); + if (notFoundHandler != null) { + String reason = json.has("Message") ? json.get("Message").getAsString() : text; + notFoundHandler.onNotFound(reason); + } + } + case NetherNetConstants.XBOX_SIGNAL_SIGNAL -> { + String sender = json.has("From") ? json.get("From").getAsString() : "0"; + if (json.has("Message")) { + dispatchSignalToPipeline(sender, json.get("Message").getAsString()); + } + } + case NetherNetConstants.XBOX_SIGNAL_CREDENTIALS -> { + log.trace("Received Credentials"); + if (json.has("Message") && connectFuture != null && !connectFuture.isDone()) { + String rawMsg = json.get("Message").getAsString(); + JsonObject credentials = JsonParser.parseString(rawMsg).getAsJsonObject(); + + connectFuture.complete(parseTurnServers(credentials)); + } + } + case NetherNetConstants.XBOX_SIGNAL_ACCEPTED, NetherNetConstants.XBOX_SIGNAL_ACK -> log.trace("Signal Ack: {}", text); default -> log.debug("Unknown message type {}: {}", type, text); } } catch (Exception e) { @@ -222,92 +93,6 @@ protected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame frame) } } - private void handleNotFound(JsonObject json, String rawText) { - log.debug("Peer Not Found. Payload: {}", rawText); - if (notFoundHandler != null) { - String reason = json.has("Message") ? json.get("Message").getAsString() : rawText; - notFoundHandler.onNotFound(reason); - } - } - - private void handleSignal(JsonObject json) { - log.trace("Received Signal: {}", json.toString()); - String sender = json.has("From") ? json.get("From").getAsString() : "0"; - if (!json.has("Message")) { - log.warn("Received SIGNAL (1) without Message payload."); - return; - } - - String rawMsg = json.get("Message").getAsString(); - dispatchSignalToPipeline(sender, rawMsg); - } - - private void handleCredentials(JsonObject json, String rawText) { - log.trace("Received Credentials: {}", rawText); - if (json.has("Message") && connectFuture != null && !connectFuture.isDone()) { - connectFuture.complete(parseTurnServers(json.get("Message").getAsString())); - } - } - - @Override - public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception { - if (connectFuture != null && !connectFuture.isDone()) { - connectFuture.completeExceptionally(cause); - } - log.error("Signaling Exception: {}", cause.getMessage(), cause); - ctx.close(); - } - - private void dispatchSignalToPipeline(String sender, String rawMsg) { - try { - // Signal Format: - String[] parts = rawMsg.split(" ", 3); - if (parts.length < 2) return; - - long connectionId = Long.parseUnsignedLong(parts[1]); - - // Try specific connection handlers (Existing Connections) - SignalHandler handler = handlers.get(connectionId); - if (handler != null) { - handler.onSignal(rawMsg); - return; - } - - // Try New Connection Handler (Server Mode) - if (NetherNetConstants.RTC_NEGOTIATION_CONNECT_REQUEST.equals(parts[0]) && newConnectionHandler != null) { - String payload = parts.length > 2 ? parts[2] : ""; - newConnectionHandler.onConnect(connectionId, sender, payload); - } - - } catch (NumberFormatException e) { - log.debug("Malformed Connection ID in signal: {}", rawMsg); - } catch (Exception e) { - log.error("Failed to dispatch signal: {}", rawMsg, e); - } - } - - @Override - public void channelInactive(ChannelHandlerContext ctx) throws Exception { - synchronized (this) { - if (connectFuture != null) { - if (!connectFuture.isDone()) { - connectFuture.completeExceptionally(new ClosedChannelException()); - } - connectFuture = null; - } - this.channel = null; - } - super.channelInactive(ctx); - } - - private void startPingLoop(ChannelHandlerContext ctx) { - ctx.executor().scheduleAtFixedRate(() -> { - JsonObject ping = new JsonObject(); - ping.addProperty("Type", 0); - ctx.writeAndFlush(new TextWebSocketFrame(gson.toJson(ping))); - }, 5, 5, TimeUnit.SECONDS); - } - @Override public void sendSignal(String targetNetworkId, String data) { if (channel != null && channel.isActive()) { @@ -317,73 +102,7 @@ public void sendSignal(String targetNetworkId, String data) { msg.addProperty("Message", data); channel.writeAndFlush(new TextWebSocketFrame(gson.toJson(msg))); } else { - throw new IllegalStateException("Attempted to send signal to " + targetNetworkId + " but WebSocket is closed or null!"); - } - } - - @Override - public void setSignalHandler(long connectionId, SignalHandler handler) { - this.handlers.put(connectionId, handler); - } - - @Override - public void removeSignalHandler(long connectionId) { - this.handlers.remove(connectionId); - } - - private List parseTurnServers(String jsonString) { - List result = new ArrayList<>(); - try { - JsonObject root = gson.fromJson(jsonString, JsonObject.class); - - JsonArray servers = null; - if (root.has("TurnAuthServers")) { - servers = root.getAsJsonArray("TurnAuthServers"); - } else if (root.has("turnAuthServers")) { - servers = root.getAsJsonArray("turnAuthServers"); - } - - if (servers != null) { - for (JsonElement el : servers) { - JsonObject server = el.getAsJsonObject(); - List urls = new ArrayList<>(); - - JsonArray urlsArray = null; - if (server.has("Urls")) { - urlsArray = server.getAsJsonArray("Urls"); - } else if (server.has("urls")) { - urlsArray = server.getAsJsonArray("urls"); - } - - if (urlsArray != null) { - urlsArray.forEach(u -> urls.add(u.getAsString())); - - IceServerInfo.Builder info = new IceServerInfo.Builder(); - info.setUrls(urls); - - if (server.has("Username")) info.setUsername(server.get("Username").getAsString()); - else if (server.has("username")) info.setUsername(server.get("username").getAsString()); - - if (server.has("Password")) info.setPassword(server.get("Password").getAsString()); - else if (server.has("password")) info.setPassword(server.get("password").getAsString()); - else if (server.has("Credential")) info.setPassword(server.get("Credential").getAsString()); - else if (server.has("credential")) info.setPassword(server.get("credential").getAsString()); - - result.add(info.build()); - } - } - } - } catch (Exception e) { - log.error("Failed to parse TURN servers", e); + throw new IllegalStateException("Attempted to send signal to " + targetNetworkId + " but WebSocket is closed!"); } - - log.debug("Successfully parsed " + result.size() + " ICE servers."); - return result; - } - - @Override - public void close() { - if (channel != null) channel.close(); - eventLoopGroup.shutdownGracefully(); } } \ No newline at end of file From cc730f548eea3923da443d50b713174e1ac53035 Mon Sep 17 00:00:00 2001 From: Joshua Castle <26531652+Kas-tle@users.noreply.github.com> Date: Mon, 16 Feb 2026 18:34:51 -0800 Subject: [PATCH 2/4] Fix javadoc Signed-off-by: Joshua Castle <26531652+Kas-tle@users.noreply.github.com> --- .../nethernet/signaling/NetherNetXboxRpcSignaling.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/signaling/NetherNetXboxRpcSignaling.java b/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/signaling/NetherNetXboxRpcSignaling.java index 010a14e..4345a77 100644 --- a/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/signaling/NetherNetXboxRpcSignaling.java +++ b/transport-nethernet/src/main/java/dev/kastle/netty/channel/nethernet/signaling/NetherNetXboxRpcSignaling.java @@ -30,8 +30,8 @@ public class NetherNetXboxRpcSignaling extends AbstractNetherNetXboxSignaling { * @param networkId The Network ID to use. * @param xboxToken The Minecraft Bedrock Session authorization header ('MCToken ***'). */ - public NetherNetXboxRpcSignaling(String localNetworkId, String xboxToken) { - super(localNetworkId, xboxToken, URI.create("wss://signal.franchise.minecraft-services.net/ws/v1.0/messaging/connect")); + public NetherNetXboxRpcSignaling(String networkId, String xboxToken) { + super(networkId, xboxToken, URI.create("wss://signal.franchise.minecraft-services.net/ws/v1.0/messaging/connect")); } /** From 6e558f24ec07dacbfc7aec413045cdd3757c866e Mon Sep 17 00:00:00 2001 From: Kas-tle <26531652+Kas-tle@users.noreply.github.com> Date: Tue, 17 Feb 2026 15:50:57 -0800 Subject: [PATCH 3/4] Bump version to 1.7.0 --- gradle.properties | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/gradle.properties b/gradle.properties index cc78677..7262f46 100644 --- a/gradle.properties +++ b/gradle.properties @@ -1,2 +1,2 @@ # Only update version on publishing to Maven Central -version=1.6.1 +version=1.7.0 From d13d1a2aaa852d794ad709ecd65d87d57e888307 Mon Sep 17 00:00:00 2001 From: Kas-tle <26531652+Kas-tle@users.noreply.github.com> Date: Tue, 17 Feb 2026 15:53:39 -0800 Subject: [PATCH 4/4] Include transport-nethernet in GitHub release --- .github/workflows/publish.yml | 1 + 1 file changed, 1 insertion(+) diff --git a/.github/workflows/publish.yml b/.github/workflows/publish.yml index 83f6422..33b9461 100644 --- a/.github/workflows/publish.yml +++ b/.github/workflows/publish.yml @@ -32,6 +32,7 @@ jobs: with: files: | transport-raknet/build/libs/*.jar + transport-nethernet/build/libs/*.jar appID: ${{ secrets.RELEASE_APP_ID }} appPrivateKey: ${{ secrets.RELEASE_APP_PK }} discordWebhook: ${{ secrets.DISCORD_WEBHOOK }}