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.
1 POST /chat ----+
2 | publish(topic, event)
3 worker ------+----------------------> [ pubSub (EventEmitter) ]
4 |
5 | listen(topic) one per open connection
6 v
7 GET /chat/sse (async generator)
8 |
9 v yield sse({ data })
10 EventSource in React
11 |
12 v
13 TanStack Query cache -> UI
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.
1export class PubSub {
2 readonly #emitter = new EventEmitter().setMaxListeners(100);
3
4 publish(topic: PubSubTopic, message: unknown): void {
5 this.#emitter.emit(topic, message);
6 }
7
8 async *listen<T = unknown>(topic: PubSubTopic, options?: { signal?: AbortSignal }) {
9 for await (const [message] of on(this.#emitter, topic, options)) {
10 yield message as T;
11 }
12 }
13}
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:
1export const pubSub: PubSub = (() => {
2 const g = globalThis as GlobalWithPubSub;
3 g[GLOBAL_KEY] ??= new PubSub();
4 return g[GLOBAL_KEY];
5})();
Topics are a const map, so a typo is a type error (topics.ts):
1export const PUB_SUB_TOPICS = {
2 HELLO_MESSAGE: "hello-message",
3 CHAT_MESSAGE: "chat-message",
4 REPO_EMBED_PROGRESS: "repo-embed-progress",
5 USER_REPO_EMBED_PROGRESS: "user-repo-embed-progress",
6} as const;
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.
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:
1export type ChatRow = AwaitedData<ReturnType<ElysiaTreaty["chat"]["get"]>>[number];
2
3export type ChatSseEvent =
4 | { type: "created"; row: ChatRow }
5 | { type: "deleted"; id: number };
src/elysia/routes/chat/index.ts. Mutations write to PGlite first and publish only what actually happened:
1export const chatRoute = new Elysia({ prefix: "/chat" })
2 .get("/", () => db.query.chat.findMany({ orderBy: (row, { desc }) => [desc(row.createdAt)] }))
3 .post(
4 "/",
5 async ({ body }) => {
6 const [row] = await db.insert(chat).values({ message: body.message }).returning();
7 if (row) pubSub.publish(PUB_SUB_TOPICS.CHAT_MESSAGE, { type: "created", row } satisfies ChatSseEvent);
8 return row;
9 },
10 { body: t.Object({ message: t.String({ minLength: 1 }) }) },
11 )
12 .delete(
13 "/:id",
14 async ({ params }) => {
15 await db.delete(chat).where(eq(chat.id, params.id));
16 pubSub.publish(PUB_SUB_TOPICS.CHAT_MESSAGE, { type: "deleted", id: params.id } satisfies ChatSseEvent);
17 return { data: { message: "Chat deleted" }, error: null };
18 },
19 { params: t.Object({ id: t.Numeric() }) },
20 )
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.
1 .get("/sse", async function* ({ request }) {
2 for await (const event of pubSub.listen<ChatSseEvent>(PUB_SUB_TOPICS.CHAT_MESSAGE, {
3 signal: request.signal,
4 })) {
5 yield sse({ data: event });
6 }
7 });
That's the entire server side of live updates: no socket bookkeeping and no subscriber list.
A tiny generic helper parses JSON frames and returns an unsubscribe function (use-embedding-sse.ts):
1export function subscribeSseJson<T>(url: string, handlers: { onMessage: (data: T) => void }): () => void {
2 const source = new EventSource(url);
3 source.onmessage = (event) => handlers.onMessage(JSON.parse(event.data));
4 source.onerror = () => {
5 if (source.readyState !== EventSource.CONNECTING) source.close();
6 };
7 return () => source.close();
8}
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.
1function applyChatSseEvent(event: ChatSseEvent) {
2 getQueryClient().setQueryData<ChatRow[]>(chatQueryKey, (prev = []) => {
3 if (event.type === "created") {
4 return prev.some((row) => row.id === event.row.id) ? prev : [event.row, ...prev];
5 }
6 return prev.filter((row) => row.id !== event.id);
7 });
8}
9
10export function useChatSse() {
11 useEffect(() => {
12 const path = getElysiaTreaty().chat.sse["~path"];
13 return subscribeSseJson<ChatSseEvent>(path, { onMessage: applyChatSseEvent });
14 }, []);
15}
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.
The initial list comes from a normal query, and live changes come from the hook (Scratchpad.tsx):
1export function Scratchpad() {
2 useChatSse();
3 const { data, isLoading, error } = useQuery({ queryKey: chatQueryKey, queryFn: listChats });
4
5}
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 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):
1.get("/activity/events", async function* ({ request }) {
2 yield sse({ data: { status: getEmbedActivityStatus(), row: null } });
3
4 for await (const payload of pubSub.listen<EmbedActivitySsePayload>(
5 PUB_SUB_TOPICS.REPO_EMBED_PROGRESS,
6 { signal: request.signal },
7 )) {****
8 yield sse({ data: payload });
9 }****
10})
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:
1useEffect(() => {
2 return subscribeSseJson<EmbedActivitySsePayload>("/api/elysia/enrich/starred/activity/events", {
3 onMessage: (payload) => {
4 setStatus(payload.status);
5 if (payload.row) upsertEnrichedRepo(payload.row);
6 },
7 });
8}, []);****
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.
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.