Nyx con Ejemplos

nyx-kv Pub/Sub

PUBLISH difunde un mensaje a todos los clientes que hicieron SUBSCRIBE a un canal. Pub/Sub desacopla a los productores de los consumidores — ideal para notificaciones en tiempo real.

Código

// nyx-kv Pub/Sub -- SUBSCRIBE and PUBLISH for real-time messaging

fn resp_cmd(parts: Array) -> String {
    var sb: StringBuilder = StringBuilder.new()
    sb.append("*")
    sb.append(int_to_string(parts.length()))
    sb.append("\r\n")
    var i: int = 0
    while i < parts.length() {
        let p: String = parts[i]
        sb.append("$")
        sb.append(int_to_string(p.length()))
        sb.append("\r\n")
        sb.append(p)
        sb.append("\r\n")
        i = i + 1
    }
    return sb.to_string()
}

fn main() -> int {
    // Publisher connection
    let pub_fd: int = tcp_connect("127.0.0.1", 6380)
    if pub_fd < 0 {
        print("publisher connection failed")
        return 1
    }

    // Publish a message to a channel
    tcp_write(pub_fd, resp_cmd(["PUBLISH", "notifications", "user signed up"]))
    let pub_reply: String = tcp_read_line(pub_fd)
    print("PUBLISH -> " + pub_reply.trim() + " subscribers received")

    tcp_close(pub_fd)

    // Subscriber connection (usually runs in a separate thread)
    let sub_fd: int = tcp_connect("127.0.0.1", 6380)
    if sub_fd < 0 { return 1 }

    // SUBSCRIBE -- the connection becomes a message stream
    tcp_write(sub_fd, resp_cmd(["SUBSCRIBE", "notifications"]))
    let sub_hdr: String = tcp_read_line(sub_fd)
    print("SUBSCRIBE header -> " + sub_hdr.trim())

    // In a real app, loop reading messages:
    //   let msg_type: String = tcp_read_line(sub_fd)
    //   let channel: String = tcp_read_line(sub_fd)
    //   let payload: String = tcp_read_line(sub_fd)

    // UNSUBSCRIBE when done
    tcp_write(sub_fd, resp_cmd(["UNSUBSCRIBE", "notifications"]))
    tcp_close(sub_fd)
    return 0
}

Salida

PUBLISH -> :0 subscribers received
SUBSCRIBE header -> *3

Explicación

Pub/Sub es un modelo de difusión "fire-and-forget" (sin confirmación de entrega). El publicador envía PUBLISH channel payload y el servidor devuelve la cantidad de clientes suscritos actualmente que recibieron el mensaje. Los mensajes no se persisten — los suscriptores que se unen más tarde nunca ven mensajes pasados. Si se necesita reproducir mensajes anteriores, hay que usar una lista o un stream.

Una conexión suscriptora entra en un modo especial: deja de aceptar comandos regulares y empieza a recibir una secuencia de arreglos de tres elementos ["message", channel, payload] por cada publicación que coincida. Como el suscriptor se bloquea en las lecturas, normalmente corre en su propio thread o goroutine; el publicador puede vivir en cualquier otra conexión.

Pub/Sub sirve para dashboards en vivo, invalidación de caché en cascada, salas de chat y buses de eventos entre servicios. Para garantizar la entrega, se puede combinar con nyx-queue.

← Anterior Siguiente →

Source: examples/by-example/75-kv-pubsub.nx