Synsema docsENES

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).

websocket.syn
-- 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.

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)

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"])

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}})
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 de ws_select/ws_recv/ws_status — por eso el patrón es un while con ws_select. Un ws_send suelto 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.