rpcWebSocketConnections function

Stream<WebSocketChannel> rpcWebSocketConnections(
  1. HttpServer server, {
  2. Duration? pingInterval = const Duration(seconds: 30),
  3. dynamic protocolSelector(
    1. List<String> protocols
    )?,
  4. CompressionOptions compression = CompressionOptions.compressionOff,
  5. Set<String>? allowedOrigins,
  6. bool allowUpgrade(
    1. HttpRequest request
    )?,
})

Turns an HttpServer into the Stream<WebSocketChannel> that RpcWebSocketServer consumes, applying server-side keepalive.

Exists because the dart:io WebSocket is only reachable between the upgrade and the wrap — IOWebSocketChannel hides it.

pingInterval defaults to 30s, and an idle connection is pinged and closed if no pong comes. Without it a peer that completes the upgrade and goes silent holds its endpoint, and the contracts on it, forever: HttpServer.idleTimeout does not reach an upgraded socket. Pass null to disable; raise it where radio wake-ups matter.

protocolSelector is forwarded to WebSocketTransformer for subprotocol negotiation.

compression defaults to OFF where dart:io's default is ON. dart:io inflates each message with no output limit before rpc_dart sees it, so RpcSecurityPolicy.maxMessageLengthBytes cannot bound it and a few hundred KiB on the wire can peak at hundreds of MiB of RSS. Turn it on only between peers you control.

allowedOrigins refuses cross-origin handshakes. WebSocket is not subject to the same-origin policy, so a browser page will open a socket anywhere and attach ambient credentials; checking Origin at the handshake is the only protocol-level defence. Compared case-insensitively against the whole origin (scheme://host[:port]).

A request with NO Origin is ALLOWED even when this is set — every non-browser client sends none, and the attack being stopped is a browser riding cookies it cannot read. More than one Origin header is refused.

allowUpgrade runs after allowedOrigins; both must accept. Synchronous on purpose — it is in the accept path, where an await serializes every handshake. If it throws the connection is refused and the server lives: the accept loop is the ROOT ZONE, where an escaping error kills the isolate.

A refused request is answered 403 and never upgraded.

final http = await HttpServer.bind(host, port);
final server = RpcWebSocketServer(
  connections: rpcWebSocketConnections(
    http,
    allowedOrigins: {'https://app.example.com'},
  ),
  onEndpointCreated: ...,
);

Implementation

Stream<WebSocketChannel> rpcWebSocketConnections(
  HttpServer server, {
  Duration? pingInterval = const Duration(seconds: 30),
  dynamic Function(List<String> protocols)? protocolSelector,
  CompressionOptions compression = CompressionOptions.compressionOff,
  Set<String>? allowedOrigins,
  bool Function(HttpRequest request)? allowUpgrade,
}) {
  // Filtered BEFORE the transformer rather than after: once WebSocketTransformer
  // has upgraded the request the response is already committed, and the only
  // thing left to do would be to close a socket the peer believes is open.
  final gated = allowedOrigins == null && allowUpgrade == null
      ? server
      : server.where((request) {
          // This runs inside the accept loop's event handler, which is the ROOT
          // ZONE: anything thrown here is an unhandled async error and kills
          // the isolate -- unauthenticated, in one request. The origin read is
          // safe now, but [allowUpgrade] is USER code and cannot be, so failing
          // closed keeps a throwing predicate to a refused connection instead
          // of a dead server.
          bool allowed;
          try {
            allowed = _upgradeAllowed(request, allowedOrigins, allowUpgrade);
          } catch (_) {
            allowed = false;
          }
          if (allowed) return true;
          _refuse(request);
          return false;
        });

  return gated
      .transform(
        WebSocketTransformer(
          protocolSelector: protocolSelector,
          compression: compression,
        ),
      )
      .map((socket) {
        // Set BEFORE wrapping: once inside IOWebSocketChannel the socket is no
        // longer reachable.
        socket.pingInterval = pingInterval;
        return IOWebSocketChannel(socket);
      });
}