Menu
popagent
publicLatest change 7f0ff66d6d9fb6468416c58bee46bd3d08169501 - Checkpoint browser channels and memory work by AkurAI Build
import type { AgentTaskStatus, ChannelEvent } from "./api-types";
type Listener = (event: ChannelEvent) => void;
export class ChannelEventBus {
private readonly listeners = new Map<string, Set<Listener>>();
publish(channelId: string, event: Omit<ChannelEvent, "id" | "channelId" | "createdAt">): ChannelEvent {
const published: ChannelEvent = {
...event,
id: crypto.randomUUID(),
channelId,
createdAt: new Date().toISOString(),
};
for (const listener of this.listeners.get(channelId) ?? []) listener(published);
return published;
}
message(channelId: string, messageId: string, taskId: string | null) {
return this.publish(channelId, { type: "message", messageId, taskId, taskStatus: null });
}
task(channelId: string, taskId: string, taskStatus: AgentTaskStatus) {
return this.publish(channelId, { type: "task", messageId: null, taskId, taskStatus });
}
subscribe(channelId: string, listener: Listener): () => void {
const channelListeners = this.listeners.get(channelId) ?? new Set<Listener>();
channelListeners.add(listener);
this.listeners.set(channelId, channelListeners);
return () => {
channelListeners.delete(listener);
if (!channelListeners.size) this.listeners.delete(channelId);
};
}
stream(channelId: string, signal?: AbortSignal): Response {
const encoder = new TextEncoder();
let unsubscribe: () => void = () => undefined;
let heartbeat: Timer | undefined;
const stream = new ReadableStream<Uint8Array>({
start: (controller) => {
let closed = false;
const send = (event: ChannelEvent) => {
if (!closed) controller.enqueue(encoder.encode(`id: ${event.id}\nevent: channel\ndata: ${JSON.stringify(event)}\n\n`));
};
unsubscribe = this.subscribe(channelId, send);
controller.enqueue(encoder.encode(": connected\n\n"));
heartbeat = setInterval(() => {
if (!closed) controller.enqueue(encoder.encode(": heartbeat\n\n"));
}, 15_000);
signal?.addEventListener("abort", () => {
if (closed) return;
closed = true;
clearInterval(heartbeat);
unsubscribe();
controller.close();
}, { once: true });
},
cancel: () => {
clearInterval(heartbeat);
unsubscribe();
},
});
return new Response(stream, {
headers: {
"content-type": "text/event-stream; charset=utf-8",
"cache-control": "no-cache, no-transform",
"x-accel-buffering": "no",
},
});
}
}
export const channelEvents = new ChannelEventBus();