Cliente WebSocket
Cualquier cosa que streamea — una suscripción RPC, el feed de órdenes de un exchange, un gateway de Discord/Slack, la mempool — antes significaba cron + polling. El cliente WebSocket lo reemplaza con una conexión en vivo. Es un transporte general, no una feature de blockchain: vive en la biblioteca estándar al lado del cliente HTTP, y lo gatea la misma capability net(host) con el mismo scope.
No es un cliente mínimo de una-conexión-por-vez: multiplexa miles de feeds en un solo thread (readiness-driven vía epoll/kqueue/IOCP — CPU ~0 mientras está ocioso), reconecta solo con backoff acotado y resubscribe, detecta conexiones half-open por keepalive, y acota la memoria con backpressure. Los dos ejes de escala: ws_select (fan-in vertical, N feeds en 1 thread) y parallel_map (fan-out horizontal, N workers × 1 conexión).
-- Doc example: the WebSocket client is a GENERAL transport (RPC subscriptions,
-- exchange feeds, chat, mempool — anything that streams), gated by the SAME
-- net(host) capability as HTTP. A connection to an undeclared host is refused at
-- the capability check, before any socket opens. (Opening a live feed needs real
-- net, which the docs sandbox denies — so this example proves the gate; the live
-- round-trip and the WS→SSE bridge are shown in the page prose.)
intent: "doc example: WebSocket client capability gating"
require net("stream.allowed.com")
task connect_undeclared()
give ws_connect("wss://feed.evil.example.com/ws") -- host not declared → denied
task bad_scheme()
give ws_connect("https://stream.allowed.com/ws") -- WebSocket needs ws:// or wss://
test "a WebSocket to an undeclared host is blocked (same net gate as HTTP)"
assert_error(connect_undeclared)
test "the scheme must be ws:// or wss:// — an http(s) URL is a directed error"
assert_error(bad_scheme)
-- Batch 15: ws_select multiplexes many feeds in one thread (Go's `select {}` in
-- one call). Its opts are validated up front — a typo doesn't become a silent
-- misconfiguration. These are PURE checks (no live server; the multiplex, reconnect
-- and keepalive round-trips are exercised against a local mock in the test suite).
task bad_on_full()
-- on_full only accepts "block" (default), "drop_oldest", "error"
give ws_connect("wss://stream.allowed.com/ws", nothing, {"on_full": "explode"})
task unknown_opt()
give ws_connect("wss://stream.allowed.com/ws", nothing, {"reconect": {}}) -- typo
test "ws_select over an empty set is nothing (no feeds → nothing to wait for)"
assert_eq(ws_select([], 0), nothing)
test "ws_status of an unknown handle is \"closed\" (never lies, never panics)"
assert_eq(ws_status(999999), "closed")
test "connect opts are validated: a bad on_full / an unknown option error clearly"
assert_error(bad_on_full)
assert_error(unknown_opt)
El modelo de capabilities (sin puerta nueva)
ws_connect necesita net(host) — deny-by-default, con scope por hostname, exactamente como http_get. No hay un permiso específico de WebSocket: un WebSocket es transporte, gateado como cualquier otra salida de red.
- Host no declarado → denegado en el chequeo de capability, antes de abrir socket alguno.
- Dentro de un
sandbox(capabilities vaciadas) → denegado. wss://valida el certificado del servidor contra los root CAs del SO, igual quehttps://.
require net("stream.exchange.com")
let conn be ws_connect("wss://stream.exchange.com/ws") -- handle opaco
Enviar, recibir, cerrar
ws_send(conn, json_encode({"op": "subscribe", "channel": "trades"})) -- text o bytes
let msg be ws_recv(conn, 5) -- próximo mensaje, o `nothing` tras el timeout de 5s
-- msg es {"type": "text" | "binary" | "close", "data": …}
when msg != nothing and msg["type"] == "text"
print(msg["data"])
ws_close(conn) -- close frame limpio (idempotente)
ws_connect(url, headers?, opts?)→ un handle opaco (como una conexióndb_open).headerses un map text→text — un valorsecret(p.ej.bearer(...)) se materializa solo en el socket, igual que los headers de HTTP.opts={"timeout", "max_message_size", "subprotocols", "max_queue", "max_queue_bytes", "on_full", "reconnect", "keepalive"}(ver abajo).ws_recv(conn, timeout?)→ el map del mensaje, onothingal vencer el timeout (default 30s, el mismo criterio quewait_for). Nunca bloquea para siempre; los frames ping/pong se manejan de forma transparente.ws_send(conn, data)envía text (sidataes text) o binary (si es bytes). Unsecretse rechaza —reveal()-lo primero si de verdad hace falta.ws_close(conn)envía un close frame limpio; cerrar un handle ya cerrado es un no-op.
Multiplexar miles de feeds: ws_select (el event-loop)
Un ws_recv espera UNA conexión. Para mirar N feeds sin quemar un thread por cada uno, ws_select espera la primera que tenga datos — el select {} de Go sobre channels, en una sola llamada. Es readiness-driven (epoll/kqueue/IOCP vía mio): CPU ~0 mientras está ocioso (duerme en el kernel, no busy-spinea) y escala a miles.
require net("feed.example.com")
let feeds be {"trades": ws_connect("wss://feed.example.com/trades"),
"book": ws_connect("wss://feed.example.com/book")} -- map nombre→handle
let live be true
while live
let m be ws_select(feeds, 30) -- la primera lista de TODAS, o nothing al vencer
when m == nothing
set live to false -- 30s ociosos → cortar (ajustá a tu feed)
when m != nothing
-- m agrega "conn" (CUÁL handle disparó) y, si es un map, "name"
when m["type"] == "close"
print("se cayó el feed " + m["name"]) -- sabés exactamente cuál resubscribir
when m["type"] == "text"
print(m["name"] + ": " + m["data"])
ws_select(conns, timeout?)→{conn, type, data, name?}de la primera conexión lista, onothingal vencer.connses una lista de handles o un map nombre→handle (y el resultado traename). Una conexión que se cae sale como{type: "close", conn}y se retira — sabés CUÁL resubscribir. Un error fatal de protocolo sale como error atrapable (conconn).ws_select_all(conns, timeout?)→ una lista con TODOS los mensajes listos de ese tick (uno por conexión; procesamiento por lote).ws_broadcast(conns, data)→ manda el mismo mensaje a muchas conexiones de una → devuelve cuántas.
Resiliencia: reconexión y keepalive (opt-in)
Sin estas opciones, el comportamiento es el de siempre — nada silencioso. Con ellas, un agente de vida larga sobrevive caídas y detecta sockets muertos sin plumbing en userland.
let conn be ws_connect("wss://feed.example.com/ws", nothing, {
"reconnect": {"max_retries": 10, "backoff": 0.5, "on_reconnect": resubscribe},
"keepalive": {"interval": 20, "timeout": 10}})
reconnect={max_retries?, backoff?, backoff_max?, on_reconnect?}— reconexión transparente con backoff exponencial acotado.on_reconnectes una task que corre tras reconectar (recibe el handle) — resubscribí ahí para no perder el estado del feed. Cada reconexión re-chequeanet(host)(nunca escala el scope) ywssre-valida el cert. Sinreconnect, una caída sale comoclose.keepalive={interval, timeout?}— auto-ping cadainterval; si no hay pong entimeout, la conexión está muerta → reconecta (si está habilitado) o emiteclose. Detección de half-open que la mayoría de las libs no hacen.ws_status(conn)→"open" | "reconnecting" | "closed".ws_stats(conn)→{sent, received, reconnects, queued, queued_bytes, last_pong_ago, status, subprotocol}— para verse (alinea con observabilidad). Las stats no mienten:sentcuenta solo mensajes que salieron (o quedaron encolados para flushear);last_pong_agoson segundos desde el último pong (nothingsi nunca llegó). Cualquier tráfico inbound cuenta como prueba de vida — un pong demorado detrás de un feed ocupado jamás mata en falso una conexión viva.
La frontera del motor síncrono. No hay hilo de fondo por conexión (rompería el aislamiento CSP y la promesa de "sin thread por conexión"). El keepalive y la reconexión avanzan mientras el programa está dentro dews_select/ws_recv/ws_status— por eso el patrón es unwhileconws_select. Unws_sendsuelto sin un recv posterior no hace tickear los timers.
Fan-out con parallel_map: N workers × 1 conexión
ws_select es el eje vertical (muchos feeds, un thread). El horizontal es parallel_map: cada worker recibe su propio intérprete (hereda las caps) y su propio registro de WebSockets — abre su feed, lo procesa y devuelve. Los handles no cruzan workers (aislamiento CSP): nunca compartas un handle entre workers.
task watch(url)
let c be ws_connect(url)
let m be ws_recv(c, 30)
ws_close(c)
give m
let resultados be parallel_map(watch, miles_de_urls) -- miles de feeds repartidos en el pool
Nunca cuelga, no se lo puede inundar, no busy-spinea
Una conexión muerta jamás cuelga al agente: connect, el handshake y cada ws_recv/ws_select están acotados por el timeout — vencido → nothing, no un programa trabado. Memoria acotada bajo entrada hostil: cada mensaje tiene tope (default 16 MiB, techo 64 MiB) y la cola inbound por conexión está acotada en las dos dimensiones — mensajes (max_queue, default 1024) y bytes (max_queue_bytes, default 64 MiB) — así que una inundación de frames grandes no puede OOMearte más que una de frames chicos. La política por default al tocar cualquiera de los topes es backpressure TCP real: se deja de leer el socket, la ventana TCP frena al peer — sin pérdida de datos ni OOM. Las alternativas son explícitas: "drop_oldest" descarta los más viejos hasta que el nuevo quepa, y "error" drena lo ya encolado y después emerge un error atrapable que nombra el overflow — jamás un descarte silencioso. La entrega entre feeds es equitativa round-robin: una conexión parlanchina con backlog no puede hambrear a las demás. Y ws_select ocioso duerme en el kernel (una sola espera de readiness), no gira en vacío. Un tope blando de conexiones por intérprete (SYNSEMA_WS_MAX_CONNS, default 4096) evita que un loop abra 100k sockets.
Subprotocolos
opts.subprotocols (lista) negocia el header Sec-WebSocket-Protocol; el subprotocolo acordado queda en ws_stats(conn)["subprotocol"].
let c be ws_connect("wss://feed.example.com/ws", nothing, {"subprotocols": ["json", "cbor"]})
let acordado be ws_stats(c)["subprotocol"] -- "json" si el server lo eligió
Bajo serve: el puente WS→SSE
El patrón de servidor más común es un puente: una ruta abre un WebSocket saliente hacia un RPC o exchange, y reenvía lo que recibe al navegador por SSE. Esto es de las pocas cosas que no se pueden mostrar del todo en run/test (necesita un servidor corriendo y un upstream en vivo), así que acá va la forma:
require serve(8080)
require net("stream.exchange.com")
serve on 8080
route "GET /live"
let conn be ws_connect("wss://stream.exchange.com/ws")
ws_send(conn, json_encode({"op": "subscribe", "channel": "trades"}))
stream -- Server-Sent Events al navegador
let live be true
while live
let msg be ws_recv(conn, 30)
when msg == nothing or msg["type"] == "close"
set live to false
when msg != nothing and msg["type"] == "text"
send msg["data"] -- reenviar cada mensaje del upstream hacia abajo
ws_close(conn)
El navegador abre un EventSource("/live") común; el servidor sostiene el WebSocket upstream. Aceptar conexiones WebSocket entrantes (un servidor WS) es una feature aparte — para chat y notificaciones, POST + SSE ya sale (ver Serve).
Qué desbloquea
Como es una primitiva general, no es solo para blockchain. Cualquier programa Synsema puede ahora construir sobre un stream bidireccional en vivo: dashboards en tiempo real, un agente reaccionando a un gateway de webhooks, herramientas colaborativas, un grabador de market-data, un cliente de chat. Donde antes escribías un loop de polling, ahora sostenés una conexión abierta.
Un ejemplo concreto: la suscripción eth_subscribe de un nodo EVM (newHeads/logs) se compone hoy en userland — ws_connect al endpoint WS del nodo, ws_send del frame JSON-RPC de subscribe, ws_recv + json_decode de cada notificación. Las lecturas one-shot (nonce, fees, saldos, receipts) van por el lado de lectura tipado de blockchain.