Files
sync-server-v2/source/services/broadcaster.ts
T
2026-07-27 10:19:13 +00:00

314 lines
13 KiB
TypeScript

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<string>;
/** 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<void>;
/**
* 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<void>;
/**
* 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<StreamEvent, 'id'>): Promise<StreamEvent>;
/**
* 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<void>;
}
/**
* 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<string, Set<BaseStream>>();
/** Reverse index: connection to all topics currently attached to it. */
private readonly streamTopics = new Map<BaseStream, Set<string>>();
/** Pending subscribe calls grouped by their stable connection identity. */
private readonly subscriptionWaiters = new Map<BaseStream, Set<SubscriptionWaiter>>();
/** Connections which already have the single required close observer. */
private readonly observedStreams = new WeakSet<BaseStream>();
/**
* 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<void> {
// 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<string>();
// 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<void>((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<SubscriptionWaiter>();
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<void> {
// 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<string>();
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<StreamEvent, 'id'>): Promise<StreamEvent> {
// 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<void> {
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);
}
}
}