DEV Community

Cover image for Live activity bus (Elysia SSE + pub/sub powered by nodejs event emitter)
Dennis kinuthia
Dennis kinuthia

Posted on AI-assisted

Live activity bus (Elysia SSE + pub/sub powered by nodejs event emitter)

Why

The workers we use run for minutes. The UI should show "embedding facebook/react", "42 done, 1 failed", and new rows appearing in the list as they land, without polling.

The shape is plain pub/sub: something publishes an event, and every open browser connection listens over Server-Sent Events. A small in-process bus sits in the middle.

 POST /chat  ----+
                 |  publish(topic, event)
 worker    ------+----------------------> [ pubSub (EventEmitter) ]
                                                 |
                                                 |  listen(topic)  one per open connection
                                                 v
                                      GET /chat/sse  (async generator)
                                                 |
                                                 v  yield sse({ data })
                                          EventSource in React
                                                 |
                                                 v
                                        TanStack Query cache -> UI
Enter fullscreen mode Exit fullscreen mode

The bus

src/lib/pub-sub/client.ts wraps a Node EventEmitter. The useful part is listen(): Node's events.on() turns an emitter into an async iterable, and passing an AbortSignal ends the loop and removes the listener.

export class PubSub {
  readonly #emitter = new EventEmitter().setMaxListeners(100);

  publish(topic: PubSubTopic, message: unknown): void {
    this.#emitter.emit(topic, message);
  }

  async *listen<T = unknown>(topic: PubSubTopic, options?: { signal?: AbortSignal }) {
    for await (const [message] of on(this.#emitter, topic, options)) {
      yield message as T;
    }
  }
}
Enter fullscreen mode Exit fullscreen mode

It's a process-wide singleton pinned on globalThis, so Vite HMR re-evaluating the module can't split routes and workers onto two different emitters:

export const pubSub: PubSub = (() => {
  const g = globalThis as GlobalWithPubSub;
  g[GLOBAL_KEY] ??= new PubSub();
  return g[GLOBAL_KEY];
})();
Enter fullscreen mode Exit fullscreen mode

Topics are a const map, so a typo is a type error (topics.ts):

export const PUB_SUB_TOPICS = {
  HELLO_MESSAGE: "hello-message",
  CHAT_MESSAGE: "chat-message",
  REPO_EMBED_PROGRESS: "repo-embed-progress",
  USER_REPO_EMBED_PROGRESS: "user-repo-embed-progress",
} as const;
Enter fullscreen mode Exit fullscreen mode

Worked example: chat

The scratchpad chat is the smallest complete version of the pattern: a list you can add to and delete from, synced live across every open window.

1. The event type

One discriminated union describes everything the stream can send (src/data-access-layer/chat/chat.ts). The row type is inferred from the route through Eden, so it can't drift:

export type ChatRow = AwaitedData<ReturnType<ElysiaTreaty["chat"]["get"]>>[number];

export type ChatSseEvent =
  | { type: "created"; row: ChatRow }
  | { type: "deleted"; id: number };
Enter fullscreen mode Exit fullscreen mode

2. The route: write, then publish

src/elysia/routes/chat/index.ts. Mutations write to PGlite first and publish only what actually happened:

export const chatRoute = new Elysia({ prefix: "/chat" })
  .get("/", () => db.query.chat.findMany({ orderBy: (row, { desc }) => [desc(row.createdAt)] }))
  .post(
    "/",
    async ({ body }) => {
      const [row] = await db.insert(chat).values({ message: body.message }).returning();
      if (row) pubSub.publish(PUB_SUB_TOPICS.CHAT_MESSAGE, { type: "created", row } satisfies ChatSseEvent);
      return row;
    },
    { body: t.Object({ message: t.String({ minLength: 1 }) }) },
  )
  .delete(
    "/:id",
    async ({ params }) => {
      await db.delete(chat).where(eq(chat.id, params.id));
      pubSub.publish(PUB_SUB_TOPICS.CHAT_MESSAGE, { type: "deleted", id: params.id } satisfies ChatSseEvent);
      return { data: { message: "Chat deleted" }, error: null };
    },
    { params: t.Object({ id: t.Numeric() }) },
  )
Enter fullscreen mode Exit fullscreen mode

3. The stream: an async generator

With Elysia, an SSE endpoint is a generator that yields sse(...) frames. Elysia sets the text/event-stream headers and stops the generator when the client disconnects. Passing request.signal to listen() also removes the emitter listener at that moment.

  .get("/sse", async function* ({ request }) {
    for await (const event of pubSub.listen<ChatSseEvent>(PUB_SUB_TOPICS.CHAT_MESSAGE, {
      signal: request.signal,
    })) {
      yield sse({ data: event });
    }
  });
Enter fullscreen mode Exit fullscreen mode

That's the entire server side of live updates: no socket bookkeeping and no subscriber list.

4. The client: EventSource → query cache

A tiny generic helper parses JSON frames and returns an unsubscribe function (use-embedding-sse.ts):

export function subscribeSseJson<T>(url: string, handlers: { onMessage: (data: T) => void }): () => void {
  const source = new EventSource(url);
  source.onmessage = (event) => handlers.onMessage(JSON.parse(event.data));
  source.onerror = () => {
    if (source.readyState !== EventSource.CONNECTING) source.close();
  };
  return () => source.close();
}
Enter fullscreen mode Exit fullscreen mode

The chat hook (use-chat-sse.ts) applies each event straight to the TanStack Query cache, so the list re-renders without a refetch. The route path comes from Eden (~path), not a hardcoded string.

function applyChatSseEvent(event: ChatSseEvent) {
  getQueryClient().setQueryData<ChatRow[]>(chatQueryKey, (prev = []) => {
    if (event.type === "created") {
      return prev.some((row) => row.id === event.row.id) ? prev : [event.row, ...prev];
    }
    return prev.filter((row) => row.id !== event.id);
  });
}

export function useChatSse() {
  useEffect(() => {
    const path = getElysiaTreaty().chat.sse["~path"];
    return subscribeSseJson<ChatSseEvent>(path, { onMessage: applyChatSseEvent });
  }, []);
}
Enter fullscreen mode Exit fullscreen mode

The some(...) check matters: the window that sent the message also receives its own created event, and this stops it from showing the row twice.

5. The component

The initial list comes from a normal query, and live changes come from the hook (Scratchpad.tsx):

export function Scratchpad() {
  useChatSse();
  const { data, isLoading, error } = useQuery({ queryKey: chatQueryKey, queryFn: listChats });
  // ...render ChatInput + ChatList
}
Enter fullscreen mode Exit fullscreen mode

The input's mutation just calls createChat(message). It doesn't touch the cache, because the SSE event does that for every window, including this one.

The same pattern for embedding progress

The chapter 3 worker publishes { status, row } on REPO_EMBED_PROGRESS from patchEmbedActivity(). The stream route adds one thing: it sends the current snapshot first, so a window that opens mid-run starts from the real state instead of waiting for the next event (enrich/starred/index.ts):

.get("/activity/events", async function* ({ request }) {
  yield sse({ data: { status: getEmbedActivityStatus(), row: null } });

  for await (const payload of pubSub.listen<EmbedActivitySsePayload>(
    PUB_SUB_TOPICS.REPO_EMBED_PROGRESS,
    { signal: request.signal },
  )) {****
    yield sse({ data: payload });
  }****
})
Enter fullscreen mode Exit fullscreen mode

On the client (use-embed-activity-sse.ts), each frame updates the status pill, and any row gets upserted into the enriched-repos collection, so newly embedded repos appear in the list while the run continues:

useEffect(() => {
  return subscribeSseJson<EmbedActivitySsePayload>("/api/elysia/enrich/starred/activity/events", {
    onMessage: (payload) => {
      setStatus(payload.status);
      if (payload.row) upsertEnrichedRepo(payload.row);
    },
  });
}, []);****
Enter fullscreen mode Exit fullscreen mode

As a safety net, while a run is live the hook also polls GET /activity every 2 s. That way a frame lost during HMR can't leave the spinner stuck.

Limitations

This isn't RabbitMQ or Redis pub/sub. It's an EventEmitter in one Deno process, and for a desktop app with one embedded server and a few windows, that's all it needs to be.

  • Single process only. Publishers and listeners must live in the same process. That's always true here.
  • No replay. A client that connects late misses earlier events. The initial fetch, plus the snapshot-first frame on the activity stream, covers that.
  • Not durable. Events are for live UI only. The real state lives in PGlite and the job queue, so a restart loses progress messages, not data.
  • No backpressure. Each connection queues events until its generator pulls them, which is fine at the rate a paced worker publishes.

If it ever needed to span processes, only PubSub would change: Postgres LISTEN/NOTIFY or Redis could sit behind the same publish/listen pair.

Top comments (0)