0One class, no dependencies
This started as a small piece of code that behaves like a minimal pub/sub (publish-subscribe) system. Nothing fancy: one Map of listener Sets plus async iterators on top. It works, it has zero dependencies, and once you see how it is wired up it is genuinely neat.1
Map, Set, promises, and the async iterator protocol. No event emitter library, no RxJS.1The Listener type and the SimplePubSub class
type Listener<T> = (payload: T) => void;
class SimplePubSub {
private topics = new Map<string, Set<Listener<any>>>();
A listener is a function that takes some data (the payload) and does something with it. That's it.
The class holds a Map called topics. Each topic name is a string; each value is a Set of listener functions. Why a Set? So duplicate listeners cannot float around: adding the same function twice is a no-op.
2Publishing: fan-out to whoever listens
publish<T>(topic: string, payload: T) {
const listeners = this.topics.get(topic);
if (!listeners) return;
for (const l of listeners) l(payload);
}
publish looks up the topic, finds every subscribed listener, and calls each one with the payload. If three listeners watch "chat.123", all three get the new message.
for await loop was mid-next() when message one landed, so it resolved instantly; messages two through four park in its queue (small dots). The late subscriber only sees what arrives after it joins.3Subscribing and the async iterator
Subscribing either finds the existing listener set for a topic or creates a fresh one:
subscribe<T>(topic: string) {
const listeners = this.topics.get(topic) ?? new Set<Listener<T>>();
this.topics.set(topic, listeners);
Then it sets up a little queue and some plumbing for async iteration:
const queue: T[] = [];
let resolveNext: ((v: IteratorResult<T>) => void) | null = null;
- If a message comes in before anyone asked for it, it gets stored in
queue. - If someone is already waiting for the next message, it resolves immediately.
That branch lives in one closure:
const onMessage = (payload: T) => {
if (resolveNext) {
resolveNext({ value: payload, done: false });
resolveNext = null;
} else {
queue.push(payload);
}
};
And the closure joins the set:
listeners.add(onMessage as Listener<T>);
The async iterator part
This is where it gets GraphQL-friendly:2
const asyncIterator: AsyncIterator<T> = {
next: () =>
new Promise<IteratorResult<T>>((resolve) => {
if (queue.length) {
const value = queue.shift()!;
resolve({ value, done: false });
} else {
resolveNext = resolve;
}
}),
- When GraphQL asks for the next value and something is already queued, return it.
- Otherwise hold onto
resolveNextand wait until a new payload comes in.
{value, done} envelope is the whole contract.4Cleanup: return and throw
Unsubscribing or erroring out removes the listener so topics do not leak:
return: () => {
listeners.delete(onMessage as Listener<T>);
return Promise.resolve({ value: undefined, done: true });
},
throw: (err) => {
listeners.delete(onMessage as Listener<T>);
return Promise.reject(err);
},
And finally, this line makes the object work nicely with for await loops:
[Symbol.asyncIterator]() {
return this;
}
LessonAn async iterator is two states in a loop: if the consumer is slow, buffer; if the consumer is waiting, resolve. Get that duality right and subscriptions feel instant without ever dropping a message.
5Topic helpers and the export
export const topicUserChatEvents = (userId: string) =>
`user.${userId}.chat-events`;
export const topicChatMessageEvents = (chatId: string) =>
`chat.${chatId}.message-events`;
Helpers keep topic naming consistent: instead of remembering to type "chat.123.message-events", call topicChatMessageEvents("123").
At the end, export one instance so any module can import it:
export const pubsub = new SimplePubSub();
6Little notes
- Super lightweight: no external libs, just
Map,Set, and async iterators. - Listeners are not "users" directly. They are functions, but usually each one maps to a user's active subscription.
- Perfect for simple GraphQL subscriptions or any little event system worth hacking together.
7References
- MDN: iteration protocols: the async iterator contract used above.
- GraphQL specification: subscriptions consume exactly such iterators.
- Local source:
bin/blogs/a-simple-pub-sub-written-in-ts.md, the original markdown note this page was rewritten from.