rpcWebSocketConnections function
- HttpServer server, {
- Duration? pingInterval = const Duration(seconds: 30),
- dynamic protocolSelector()?,
- CompressionOptions compression = CompressionOptions.compressionOff,
- Set<
String> ? allowedOrigins, - bool allowUpgrade(
- 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);
});
}