110 / Redes

Eventos del servidor (SSE)

std/serve sirve Server-Sent Events con una ruta NORMAL: pasa por hooks, middlewares, mounts y wraps, así que un middleware de autenticación que responde 401 corta el canal antes de que se abra. sse_open(req, room) abre la conexión, sse_broadcast(room, evento, datos) difunde desde cualquier handler o thread, y el canal no ocupa un worker: el proceso escribe la cabecera, registra el fd en el room y suelta la conexión. Detrás de un proxy inverso que no bufferice el stream, el canal funciona sin cambios; uno que sí bufferiza recién suelta los eventos cuando el origen cierra la conexión, y hasta entonces el cliente no ve nada.

110-serve-sse.nxFuente →
// std/serve Server-Sent Events — a normal route opens the channel, rooms broadcast
// Eventos del servidor (SSE) con std/serve — una ruta normal abre el canal, los rooms difunden
//
// The browser side is browser_sse_fn(url, fn(evento, datos) {...}) from
// std/browser. This recipe plays the client with a raw TCP socket so it can
// run (and finish) on its own.
//
// Behind a gateway: nyx-proxy 0.4.4 tunnels SSE end to end, so this works
// unchanged behind it. Older proxies buffer the stream until the upstream
// closes (the client sees nothing), so with one of those, expose SSE only on
// a server the client reaches directly.

import "std/http"
import "std/web"
import "std/serve"

// A NORMAL route: it runs through hooks, middlewares, mounts and wraps, so
// an auth middleware that answers 401 stops it before any channel opens.
fn open_inventory(req: Request) -> Response {
    return sse_open(req, "inventario")
}

fn main() -> int {
    let app: App = app_new()
    let port: int = 18208
    app_get(app, "/eventos/inventario", open_inventory)

    // Real server in a detached goroutine; it dies with the process.
    spawn {
        serve_app(app, port, 2)
    }

    // Wait for the listener.
    var fd: int = 0 - 1
    var tries: int = 0
    while fd < 0 and tries < 50 {
        fd = tcp_connect("127.0.0.1", port)
        if fd < 0 {
            sleep(100)
        }
        tries = tries + 1
    }
    if fd < 0 {
        print("the server did not start")
        return 1
    }
    tcp_set_timeout(fd, 5)

    // Open the channel like a browser would.
    tcp_write(fd, "GET /eventos/inventario HTTP/1.1\r\nHost: 127.0.0.1\r\nAccept: text/event-stream\r\n\r\n")
    let head: String = tcp_read_partial(fd, 4096)
    if head.indexOf("text/event-stream") >= 0 {
        print("head: 200 text/event-stream, no Content-Length")
    } else {
        print("unexpected head: " + head)
        tcp_close(fd)
        return 1
    }

    // The channel is registered once the head is written; it does NOT hold
    // a worker, so the server keeps answering other requests meanwhile.
    var waited: int = 0
    while sse_count_room("inventario") == 0 and waited < 50 {
        sleep(20)
        waited = waited + 1
    }
    print("open channels in the room: " + int_to_string(sse_count_room("inventario")))

    // From any handler or thread: one event to everyone in the room. The
    // return value counts clients that received the WHOLE frame.
    let delivered: int = sse_broadcast("inventario", "stock", "{\"sku\":\"A-1\",\"cantidad\":7}")
    print("delivered to: " + int_to_string(delivered))

    // What the client reads: `event:` + one `data:` per line + a blank line.
    let frame: String = tcp_read_partial(fd, 4096)
    print(frame)

    // An event name with a line break is rejected (0), never sanitized.
    let rejected: int = sse_broadcast("inventario", "bad\nname", "x")
    print("rejected event name delivered to: " + int_to_string(rejected))

    tcp_close(fd)
    return 0
}
Salidastdout
nyx-serve starting on 127.0.0.1:18208 with 2 workers
head: 200 text/event-stream, no Content-Length
open channels in the room: 1
delivered to: 1
event: stock
data: {"sku":"A-1","cantidad":7}


rejected event name delivered to: 0

Cómo funciona

Esta receta juega el rol del cliente con un socket TCP crudo, así puede correr (y terminar) sola: el header de respuesta llega sin Content-Length, como corresponde a un stream que no sabe cuándo termina. Del lado del navegador el equivalente es browser_sse_fn(url, fn(evento, datos) {...}) de std/browser.

sse_broadcast devuelve cuántos clientes recibieron el frame COMPLETO — acá 1, el único canal abierto en el room "inventario". Lo que el cliente lee es exactamente event: más una línea data: por línea del payload, seguido de una línea en blanco que cierra el frame.

Un nombre de evento con \r o \n se RECHAZA (sse_broadcast devuelve 0), nunca se sanea en silencio — de ahí que la última línea muestre 0 entregas.