Menu
popagent
publicLatest change bea7c647f451d774dcf4753d7ca5ca5e01d0cf54 - Rename supervisor agent ID to orchistrator by Ólafur Búi Ólafsson
import { z } from "zod";
import type { AgentRuntimeSettingsStore } from "./agent-runtime-settings";
import type { AgentWorkspaceStore } from "./agent-workspaces";
import type { ChannelSettingsInput } from "./api-types";
import { ChannelStore } from "./channels";
import { channelEvents, type ChannelEventBus } from "./channel-events";
import { isKnownModel } from "./models";
import type { AgentTaskRuntime, AgentTaskStore } from "./tasks";
const channelSettingsSchema: z.ZodType<ChannelSettingsInput> = z.object({
enabled: z.boolean(),
channelName: z.string().trim().regex(/^#[a-z0-9][a-z0-9_-]{0,63}$/),
dispatchMode: z.enum(["every-message", "mentions"]),
contextMessages: z.number().int().min(0).max(100),
streaming: z.boolean(),
toolDisplay: z.enum(["compact", "timeline"]),
}).strict();
const channelPostSchema = z.object({
content: z.string().trim().min(1).max(8000),
workspaceId: z.string().trim().min(1).max(256),
model: z.string().trim().min(1).max(200).optional(),
}).strict();
const MAX_BODY_BYTES = 32_000;
async function readJson(request: Request): Promise<unknown> {
const text = await request.text();
if (text.length > MAX_BODY_BYTES) throw new RangeError("body too large");
return JSON.parse(text);
}
function addressesOrchestrator(content: string): boolean {
return /@(orchistrator|orchestrator)\b/i.test(content);
}
function contextPrompt(channelName: string, history: Awaited<ReturnType<ChannelStore["listMessages"]>>, content: string) {
const context = history.map((message) => {
const result = message.task?.status === "completed" && message.task.output
? `\n[Orchestrator]: ${message.task.output}`
: "";
return `[${message.authorName}]: ${message.content}${result}`;
}).join("\n");
return [
`Respond to a task dispatch from the shared ${channelName} channel.`,
context ? `Recent channel context:\n${context}` : undefined,
`Current dispatch:\n${content}`,
"Complete the requested work using the selected repository workspace and return a concise channel-ready result.",
].filter(Boolean).join("\n\n");
}
type RouteRequest = Request & { params: { id: string; taskId: string } };
type RouteHandler = (request: RouteRequest) => Response | Promise<Response>;
export type ChannelRouteDefinition = RouteHandler | Partial<Record<"GET" | "POST" | "PATCH", RouteHandler>>;
export function createChannelRoutes(dependencies: {
channelStore: ChannelStore;
workspaceStore: Pick<AgentWorkspaceStore, "get">;
runtimeSettings: Pick<AgentRuntimeSettingsStore, "get">;
taskStore: Pick<AgentTaskStore, "init">;
taskRunner: Pick<AgentTaskRuntime, "enqueue" | "cancel">;
eventBus?: ChannelEventBus;
}): Record<string, ChannelRouteDefinition> {
const { channelStore, workspaceStore, runtimeSettings, taskStore, taskRunner } = dependencies;
const eventBus = dependencies.eventBus ?? channelEvents;
return {
"/api/settings/channels": {
GET: async () => Response.json(await channelStore.getSettings()),
PATCH: async (request) => {
const input = channelSettingsSchema.safeParse(await readJson(request).catch(() => undefined));
if (!input.success) {
return Response.json({ error: "Invalid channel settings", issues: input.error.issues }, { status: 400 });
}
try {
const settings = await channelStore.updateSettings(input.data);
eventBus.publish("general", { type: "settings", messageId: null, taskId: null, taskStatus: null });
return Response.json(settings);
} catch (error) {
return Response.json({ error: error instanceof Error ? error.message : String(error) }, { status: 400 });
}
},
},
"/api/channels": {
GET: async () => Response.json({ channels: await channelStore.listChannels() }),
},
"/api/channels/:id/messages": {
GET: async (request) => {
const channel = (await channelStore.listChannels()).find((item) => item.id === request.params.id);
if (!channel) return Response.json({ error: "Channel not found" }, { status: 404 });
const url = new URL(request.url);
const rawLimit = url.searchParams.get("limit");
const limit = rawLimit === null ? 100 : Number(rawLimit);
const before = url.searchParams.get("before") ?? undefined;
try {
const messages = await channelStore.listMessages(request.params.id, { limit, before });
return Response.json({ messages, nextCursor: messages.length === limit ? messages[0]?.id ?? null : null });
} catch (error) {
return Response.json({ error: error instanceof Error ? error.message : String(error) }, { status: 400 });
}
},
POST: async (request) => {
const input = channelPostSchema.safeParse(await readJson(request).catch(() => undefined));
if (!input.success) {
return Response.json({ error: "Invalid channel message", issues: input.error.issues }, { status: 400 });
}
const settings = await channelStore.getSettings();
if (!settings.enabled) return Response.json({ error: "Channels are disabled" }, { status: 409 });
if (!(await workspaceStore.get(input.data.workspaceId))) {
return Response.json({ error: "Workspace not found" }, { status: 404 });
}
const model = input.data.model ?? (await runtimeSettings.get()).defaultModel;
if (!(await isKnownModel(model))) {
return Response.json({ error: `Unknown model: ${model}` }, { status: 400 });
}
const shouldDispatch = settings.dispatchMode === "every-message" || addressesOrchestrator(input.data.content);
try {
if (!shouldDispatch) {
const message = await channelStore.createMessage({
channelId: request.params.id,
workspaceId: input.data.workspaceId,
content: input.data.content,
authorId: "popagent-user",
authorName: "You",
});
eventBus.message(request.params.id, message.id, null);
return Response.json({ message, task: null }, { status: 201 });
}
await taskStore.init();
const history = settings.contextMessages
? await channelStore.listMessages(request.params.id, { limit: settings.contextMessages })
: [];
const result = await channelStore.createDispatch({
channelId: request.params.id,
workspaceId: input.data.workspaceId,
content: input.data.content,
model,
authorId: "popagent-user",
authorName: "You",
taskPrompt: contextPrompt(settings.channelName, history, input.data.content),
});
taskRunner.enqueue(result.task);
eventBus.message(request.params.id, result.message.id, result.task.id);
return Response.json(result, { status: 201 });
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
return Response.json({ error: message }, { status: message === "Channel not found" ? 404 : 400 });
}
},
},
"/api/channels/:id/events": {
GET: async (request) => {
const channel = (await channelStore.listChannels()).find((item) => item.id === request.params.id);
return channel
? eventBus.stream(request.params.id, request.signal)
: Response.json({ error: "Channel not found" }, { status: 404 });
},
},
"/api/channels/:id/tasks/:taskId/cancel": {
POST: async (request) => {
const message = await channelStore.findMessageByTask(request.params.id, request.params.taskId);
if (!message) return Response.json({ error: "Channel task not found" }, { status: 404 });
if (!await taskRunner.cancel(request.params.taskId)) {
return Response.json({ error: "Task cannot be cancelled" }, { status: 409 });
}
eventBus.task(request.params.id, request.params.taskId, "cancelling");
return Response.json({ status: "cancelling" });
},
},
};
}