typescript · events · pub/sub

Building a tiny pub/sub system in TypeScript

One Map, one Set, and async iterators are enough to run GraphQL-style subscriptions without pulling in a library. This walks through the whole class, including the queue-versus-resolve branch that makes for await behave.

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

fn 1Everything below uses only 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.

PUBLISHER publish() topic chat.123 handler (sync) for await consumer awaiting next() late subscriber buffered
published4
resolved on arrival1
buffered for later3
dependencieszero
every message is either handed to a waiting iterator immediately or parked in its queue until the consumer asks
Fig. 1. Four publishes against three consumers. The sync handler takes everything. The 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;

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;
      }
    }),
fn 2GraphQL subscription resolvers are expected to return an AsyncIterator exactly like this one. The {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;
}
Lesson

An 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


7References

  1. MDN: iteration protocols: the async iterator contract used above.
  2. GraphQL specification: subscriptions consume exactly such iterators.
  3. Local source: bin/blogs/a-simple-pub-sub-written-in-ts.md, the original markdown note this page was rewritten from.