/** * @wrnexus/pubsub — topic-based publish/subscribe with a pluggable driver. * The default is in-process; swap in a Redis/NATS driver for cross-instance * messaging (it also backs @wrnexus/core's realtime bridge). * * const bus = createPubSub(); * const off = bus.subscribe("order:*", (msg, topic) => {...}); * await bus.publish("order:created", { id: 7 }); * * Subscriptions match exact topics, "ns:*" prefixes, and "*" (everything). */ export type Handler = (message: T, topic: string) => void | Promise; export interface PubSubDriver { publish(topic: string, message: unknown): void | Promise; subscribe(pattern: string, handler: Handler): () => void; close?(): void | Promise; } export interface PubSub { publish(topic: string, message: T): Promise; subscribe(pattern: string, handler: Handler): () => void; /** Stop new work, remove subscriptions, and close the backing driver. */ close(): Promise; } function patternMatches(pattern: string, topic: string): boolean { if (pattern === "*" || pattern === topic) return true; if (pattern.endsWith(":*")) return topic.startsWith(pattern.slice(0, -1)); // "post:" prefix return false; } /** In-process pub/sub driver (default). */ export function memoryDriver(): PubSubDriver { const subs = new Map>(); let closed = false; return { async publish(topic, message) { if (closed) throw new Error("WRN-PUBSUB-CLOSED: pubsub driver is closed"); const pending: Promise[] = []; for (const [pattern, handlers] of subs) { if (!patternMatches(pattern, topic)) continue; for (const handler of handlers) pending.push(Promise.resolve(handler(message, topic))); } await Promise.all(pending); }, subscribe(pattern, handler) { if (closed) throw new Error("WRN-PUBSUB-CLOSED: pubsub driver is closed"); let set = subs.get(pattern); if (!set) subs.set(pattern, (set = new Set())); set.add(handler); return () => { set!.delete(handler); if (!set!.size) subs.delete(pattern); }; }, close() { closed = true; subs.clear(); }, }; } /** Create a pub/sub bus over a driver (in-memory by default). */ export function createPubSub(driver: PubSubDriver = memoryDriver()): PubSub { let closed = false; return { async publish(topic, message) { if (closed) throw new Error("WRN-PUBSUB-CLOSED: pubsub is closed"); if (!topic.trim()) throw new TypeError("pubsub topic cannot be empty"); await driver.publish(topic, message); }, subscribe(pattern, handler) { if (closed) throw new Error("WRN-PUBSUB-CLOSED: pubsub is closed"); if (!pattern.trim()) throw new TypeError("pubsub pattern cannot be empty"); return driver.subscribe(pattern, handler as Handler); }, async close() { if (closed) return; closed = true; await driver.close?.(); }, }; } export { createResilientPubSub, PresenceChannel } from "./resilient.ts"; export type { MessageEnvelope, ResilientPubSubOptions, PresenceMember } from "./resilient.ts"; export { natsDriver, kafkaDriver } from "./brokers.ts"; export type { NatsClient, KafkaClient } from "./brokers.ts"; export { subjectPubSub } from "./subject.ts"; export type { SubjectPubSub } from "./subject.ts";