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.