import { HTTP_STATUS_CODE_NOT_ACCEPTED } from '../constants.ts'; import type { BaseStream, StreamEvent } from './stream/base-stream.ts'; import type { Logger } from '../utils/logger.ts'; import { ApplicationError } from '../errors/index.ts'; /** Request-scoped view from which the broadcaster obtains a stable connection. */ export interface BroadcastStream { /** Connection identity shared by every request on the same transport session. */ readonly connection: BaseStream; /** Whether the connection can remain open to receive published events. */ readonly streaming: boolean; } /** One pending subscribe call and the topics whose removal will resolve it. */ interface SubscriptionWaiter { /** Only topics newly introduced by this particular subscribe call. */ readonly remainingTopics: Set; /** Completes the promise returned to the subscribing route. */ readonly resolve: () => void; } /** Topic delivery contract consumed by domain routes. */ export abstract class BaseBroadcaster { /** * Subscribe a connection and wait until the topics added by this call are removed. * * Fully duplicate subscriptions resolve immediately. * * @param stream - The stream to subscribe. * @param topics - The topics to subscribe to. */ abstract subscribe(stream: BroadcastStream, topics: string[]): Promise; /** * Unsubscribes a stream from a list of topics. * @param stream - The stream to unsubscribe. * @param topics - The topics to unsubscribe from. */ abstract unsubscribe(stream: BroadcastStream, topics?: string[]): Promise; /** * Publishes an event to a topic. * @param topic - The topic to publish to. * @param event - The event to publish. * @returns The published event. */ abstract publish(topic: string, event: Omit): Promise; /** * Sends an event to a stream. * @param stream - The stream to send the event to. * @param event - The event to send. */ abstract sendEvent(stream: BaseStream, event: StreamEvent): Promise; } /** * In-memory topic index with reverse lookup for deterministic stream cleanup. * * A stream is stored strongly only while it has topics. The WeakSet records that * its close observer has already been installed without extending its lifetime. */ export class Broadcaster extends BaseBroadcaster { /** Namespaced diagnostic logger for subscription and publication activity. */ private readonly debug: Logger; /** Forward index: topic name to the connections receiving that topic. */ private readonly topicStreams = new Map>(); /** Reverse index: connection to all topics currently attached to it. */ private readonly streamTopics = new Map>(); /** Pending subscribe calls grouped by their stable connection identity. */ private readonly subscriptionWaiters = new Map>(); /** Connections which already have the single required close observer. */ private readonly observedStreams = new WeakSet(); /** * Creates a new Broadcaster. * @param debug - The debug logger. */ constructor(debug: Logger) { super(); // Extend the debug logger to include the broadcaster namespace. this.debug = debug.extend('broadcaster'); } /** * Subscribe a connection and return a promise for this call's additions. * * Registration is synchronous. The returned promise resolves after all * topics newly added by this call are removed, or when the connection closes. * If every requested topic already exists, it resolves immediately. * * @param stream - The stream to subscribe. * @param topics - The topics to subscribe to. */ subscribe(stream: BroadcastStream, topics: string[]): Promise { // Normal HTTP cannot receive later publications, so fail before // mutating either subscription index. if (!stream.streaming) { throw new ApplicationError(HTTP_STATUS_CODE_NOT_ACCEPTED, 'This route requires a stream-capable connection'); } // ApplicationRouteStream is request-scoped, but subscriptions must survive // across requests. Always index by the shared underlying connection. const connection = stream.connection; // Reuse the connection's reverse-index entry when it already has topics. // A new Set is not stored until this call actually introduces a topic. const trackedTopics = this.streamTopics.get(connection) ?? new Set(); // Deduplicate the request itself const deduplicatedTopics = Array.from(new Set(topics)); // Filter topics that are already subscribed to by the connection. const topicsToAdd = deduplicatedTopics.filter((topic) => !trackedTopics.has(topic)); // A fully duplicate (or empty) subscription adds no lifetime to track. if (topicsToAdd.length === 0) { return Promise.resolve(); } // Store the reverse index before installing the close observer. An // already-closed connection invokes onClose immediately and must be able // to remove the topics registered by this call. this.streamTopics.set(connection, trackedTopics); for (const topic of topicsToAdd) { // Find or create the forward-index set for this topic. let streams = this.topicStreams.get(topic); if (!streams) { streams = new Set(); this.topicStreams.set(topic, streams); } // Update both indexes together: publication uses the forward index, // while unsubscribe and connection cleanup use the reverse index. streams.add(connection); trackedTopics.add(topic); } // Create the lifecycle promise returned to the route. Its waiter owns only // the topics added above, not duplicate topics owned by earlier calls. const removed = new Promise((resolve) => { // Several non-overlapping subscription requests can remain active on // one WebSocket connection, so each connection stores a set of waiters. const waiters = this.subscriptionWaiters.get(connection) ?? new Set(); waiters.add({ remainingTopics: new Set(topicsToAdd), resolve, }); this.subscriptionWaiters.set(connection, waiters); }); // Register the waiter before observing closure: onClose invokes its // callback immediately when registration races with an already-closed stream. if (!this.observedStreams.has(connection)) { // WeakSet prevents repeated subscribe requests from adding duplicate // close callbacks without retaining an otherwise unused connection. this.observedStreams.add(connection); // Remote disconnect, local close, and shutdown all use the same cleanup // path, which also resolves every affected subscription promise. connection.onClose(() => this.removeSubscriptions(connection)); } // Log only the topics introduced by this call; duplicates were no-ops. this.debug('subscribed stream to topics %o', topicsToAdd); // Keep the route dispatch pending for exactly this subscription's // lifetime. Duplicate calls return an already-resolved promise above. return removed; } /** * Unsubscribes a stream from a list of topics. * @param stream - The stream to unsubscribe. * @param topics - The topics to unsubscribe from. */ async unsubscribe(stream: BroadcastStream, topics?: string[]): Promise { // Resolve request-scoped facades to the same stable connection key used // during subscribe, allowing a later WebSocket request to unsubscribe. const topicsToRemove = this.removeSubscriptions(stream.connection, topics); // Logging remains useful even for idempotent removal of missing topics. this.debug('unsubscribed stream from topics %o', topicsToRemove); } /** * Remove topics from a connection and resolve affected subscription calls. * * @param connection - Stable connection stored in the topic index. * @param topics - Specific topics to remove, or all current topics. * @returns The topics considered for removal. */ private removeSubscriptions(connection: BaseStream, topics?: string[]): string[] { // Missing connections are valid: unsubscribe is deliberately idempotent. const trackedTopics = this.streamTopics.get(connection); // Omitting topics means connection cleanup, so remove every tracked topic. const topicsToRemove = topics ?? Array.from(trackedTopics ?? []); // Waiters should advance only for topics that were genuinely active. const removedTopics = new Set(); for (const topic of topicsToRemove) { // Deleting from the reverse index reports whether this call actually // removed an active connection/topic relationship. if (trackedTopics?.delete(topic)) { removedTopics.add(topic); } // Remove the same relationship from the publication index. this.topicStreams.get(topic)?.delete(connection); // Empty topic sets have no value and would unnecessarily retain maps. if (this.topicStreams.get(topic)?.size === 0) { this.topicStreams.delete(topic); } } // Stop strongly retaining connections after their final topic is removed. if (!trackedTopics || trackedTopics.size === 0) { this.streamTopics.delete(connection); } // Resolve subscribe calls whose newly added topics have all disappeared. const waiters = this.subscriptionWaiters.get(connection); if (waiters) { for (const waiter of waiters) { // Partial unsubscribe removes only the affected portion of each // waiter's outstanding topic set. for (const topic of removedTopics) { waiter.remainingTopics.delete(topic); } // The route completes once every topic introduced by its call has // been removed, even if the connection still has other topics. if (waiter.remainingTopics.size === 0) { waiters.delete(waiter); waiter.resolve(); } } // Avoid retaining an empty waiter collection after all routes settle. if (waiters.size === 0) { this.subscriptionWaiters.delete(connection); } } // Return the requested removal list for consistent unsubscribe logging. return topicsToRemove; } /** * Publishes an event to a topic. * @param topic - The topic to publish to. * @param event - The event to publish. * @returns The published event. */ async publish(topic: string, event: Omit): Promise { // Get the current timestamp. const timestamp = Date.now(); // Add an ID to the event. const eventWithId: StreamEvent = { ...event, id: String(timestamp), }; // Copy the current subscriber set and start every send immediately. // Promise.all provides concurrent fan-out while still allowing publish to // wait until every local delivery attempt has settled. await Promise.all(Array.from(this.topicStreams.get(topic) ?? [], (stream) => this.sendEvent(stream, eventWithId))); // Log the published event. this.debug('published %s to topic %s', event.type, topic); return eventWithId; } /** * Sends an event to a stream. * @param stream - The stream to send the event to. * @param event - The event to send. */ async sendEvent(stream: BaseStream, event: StreamEvent): Promise { try { // Broadcaster messages bypass the request facade because pushed events // are connection-level and must not inherit a request correlation ID. await stream.send(event); } catch (error) { // Log the error. this.debug('failed to send event to stream: %O', error); // A failed connection cannot receive future publications. Closing it // triggers the normal observer cleanup; the explicit removal also // makes this path safe for unusual stream implementations. stream.close(); this.removeSubscriptions(stream); } } }