Menu
popagent
publicLatest change d6a6ee8ba42f872e46ff85d4be7ee1a6c2acf5da - Render channel messages as Markdown and announce a task report once by AkurAI Build
import { handleChatStream } from "@mastra/ai-sdk";
import {
createUIMessageStream,
createUIMessageStreamResponse,
type UIMessage,
type UIMessageStreamWriter,
} from "ai";
import { Cron } from "croner";
import { z } from "zod";
import index from "./ui/index.html";
import {
AGENT_TOOL_NAMES,
DEFAULT_WORKSPACE_ID,
MAX_AGENT_INSTRUCTIONS_CHARACTERS,
MODEL_SOURCES,
runtimeModelId,
type AgentActivity,
type AgentRuntimeSettingsInput,
type AutonomySettingsInput,
type BrowserSettingsInput,
type AgentSettingsInput,
type AgentSecretInput,
type AgentSecretUpdateInput,
type AgentSkillInput,
type ChatMessage,
type MemorySettingsInput,
type SessionSummary,
} from "./api-types";
import { companyAgentProfiles } from "./company-roster";
import { agentActivity, AgentExecutionRuntime, createAgentTaskExecutor, createTraceId } from "./agent-execution";
import { autonomyActivity } from "./autonomy-runtime";
import { agent, mastra, taskParticipantAgents } from "./agent";
import { agentRuntimeSettings, type AgentRuntimeSettingsStore } from "./agent-runtime-settings";
import { agentSettings, type AgentSettingsStore } from "./agent-settings";
import { autonomySettings, type AutonomySettingsStore } from "./autonomy-settings";
import { agentSkills, type AgentSkillStore } from "./agent-skills";
import { validateAgentHealth } from "./agent-health";
import { browserRuntime, browserSettings, type BrowserRuntime, type BrowserSettingsStore } from "./browser-settings";
import { browserRecordings } from "./browser-recordings";
import { browserProfiles } from "./browser-profiles";
import { agentWorkspaceRuntime } from "./workspace";
import { latestUserText, replaceLatestUserText, sanitizeSensitiveToolMessages } from "./chat-messages";
import { companyAvatarIds, companyDepartments } from "./company-roster";
import { createChannelRoutes } from "./channel-routes";
import { createCodeIntelligenceRoutes } from "./code-intelligence-routes";
import { channelEvents } from "./channel-events";
import { buildEvolutionEvidence, evolutionStore, type EvolutionStore } from "./evolution-store";
import { evolutionRuntime } from "./evolution-runtime";
import { channels, taskUsesChannelTranscript } from "./channels";
import { hookAudit } from "./hook-audit";
import { createTaskLifecycleHooks } from "./hook-lifecycle";
import { createHookEvent, HookBlockedError, loadHookRuntime, type HookRuntime } from "./hooks";
import { longTermMemory, MAX_FACT_CHARACTERS, type LongTermMemoryStore } from "./long-term-memory";
import { memorySettings, type MemorySettingsStore } from "./memory-settings";
import { isKnownModel, listModelCatalog } from "./models";
import {
getObservabilityOverview,
getObservabilityTrace,
listObservabilityLogs,
listObservabilityTraces,
observabilityHealth,
parseObservabilityQuery,
} from "./observability-store";
import { RESOURCE_ID, SessionStore } from "./sessions";
import { AgentTaskRuntime, AgentTaskStore } from "./tasks";
import { SelfUpdateScheduler } from "./self-update-scheduler";
import { requireApiKey } from "./auth";
import { startMaintenance } from "./maintenance";
import { appLogger, observability } from "./observability";
import { agentWorkspaces, type AgentWorkspaceStore } from "./agent-workspaces";
import { repositoryBrief } from "./repo-brief";
import { documentationIndex, type DocumentationIndex } from "./documentation-index";
import { DocumentationConflictError, MAX_DOCUMENTATION_CHARACTERS } from "./documentation-files";
import {
agentSecrets,
MAX_SECRET_NAME_CHARACTERS,
MAX_SECRET_VALUE_CHARACTERS,
normalizeSecretOrigin,
SECRET_NAME_PATTERN,
} from "./secrets";
import { TaskAgentCommunicationSession } from "./task-agent-communication";
const observabilityFeedbackSchema = z.object({
traceId: z.string().regex(/^[a-f0-9]{32}$/),
value: z.union([z.literal(1), z.literal(-1)]),
comment: z.string().trim().max(2_000).optional(),
}).strict();
const skillInputSchema: z.ZodType<AgentSkillInput> = z.object({
name: z.string().regex(/^[a-z0-9]+(?:-[a-z0-9]+)*$/).max(64),
description: z.string().min(1).max(1024),
instructions: z.string().min(1),
references: z.record(z.string(), z.string()),
enabled: z.boolean(),
userInvocable: z.boolean(),
});
const browserHostSchema = z.string().trim().min(1).max(255).regex(/^(?:\*\.)?[a-z0-9.-]+(?::\d{1,5})?$/i);
const browserSettingsSchema: z.ZodType<BrowserSettingsInput> = z.object({
enabled: z.boolean(),
provider: z.enum(["agent-browser", "bifrost-navigator"]),
scope: z.enum(["thread", "shared"]),
viewportWidth: z.number().int().min(320).max(3840),
viewportHeight: z.number().int().min(240).max(2160),
timeoutMs: z.number().int().min(1_000).max(120_000),
maxSessions: z.number().int().min(1).max(16),
idleTimeoutMs: z.number().int().min(60_000).max(86_400_000),
screencastEnabled: z.boolean(),
screenshotsEnabled: z.boolean(),
multiTabEnabled: z.boolean(),
formsEnabled: z.boolean(),
dialogsEnabled: z.boolean(),
dragEnabled: z.boolean(),
evaluateEnabled: z.boolean(),
recordingEnabled: z.boolean(),
recordingRetentionDays: z.number().int().min(1).max(90),
recordingMaxFiles: z.number().int().min(1).max(100),
allowHosts: z.array(browserHostSchema).max(100),
denyHosts: z.array(browserHostSchema).max(100),
}).strict();
const browserProfileSchema = z.object({
name: z.string().trim().min(1).max(120),
enabled: z.boolean(),
state: z.unknown(),
}).strict();
const browserProfileUpdateSchema = browserProfileSchema.extend({ state: z.unknown().optional() });
const browserMouseEventSchema = z.object({
type: z.enum(["mousePressed", "mouseReleased", "mouseMoved", "mouseWheel"]),
x: z.number().finite().min(0).max(10_000),
y: z.number().finite().min(0).max(10_000),
button: z.enum(["left", "right", "middle", "none"]).optional(),
clickCount: z.number().int().min(0).max(3).optional(),
deltaX: z.number().finite().min(-10_000).max(10_000).optional(),
deltaY: z.number().finite().min(-10_000).max(10_000).optional(),
}).strict();
const browserKeyboardEventSchema = z.object({
type: z.enum(["keyDown", "keyUp", "char"]),
key: z.string().min(1).max(64),
code: z.string().max(64).optional(),
text: z.string().max(1000).optional(),
modifiers: z.number().int().min(0).max(15).optional(),
}).strict();
const browserInputSchema = z.discriminatedUnion("kind", [
z.object({ kind: z.literal("mouse"), event: browserMouseEventSchema }).strict(),
z.object({ kind: z.literal("keyboard"), event: browserKeyboardEventSchema }).strict(),
]);
const secretNameSchema = z.string()
.trim()
.min(1)
.max(MAX_SECRET_NAME_CHARACTERS)
.regex(SECRET_NAME_PATTERN);
const secretOriginSchema = z.string().trim().min(1).max(2_048).superRefine((value, context) => {
try {
normalizeSecretOrigin(value);
} catch (error) {
context.addIssue({ code: "custom", message: messageOf(error) });
}
});
const agentSecretInputSchema: z.ZodType<AgentSecretInput> = z.object({
name: secretNameSchema,
value: z.string().min(1).max(MAX_SECRET_VALUE_CHARACTERS),
allowedOrigin: secretOriginSchema,
}).strict();
const agentSecretUpdateSchema: z.ZodType<AgentSecretUpdateInput> = z.object({
value: z.string().min(1).max(MAX_SECRET_VALUE_CHARACTERS).optional(),
allowedOrigin: secretOriginSchema.optional(),
}).strict().refine(
(input) => input.value !== undefined || input.allowedOrigin !== undefined,
"Secret update must include a value or allowed origin",
);
const memorySettingsSchema: z.ZodType<MemorySettingsInput> = z.object({
autoCompact: z.boolean(),
observationTokens: z.number().int().min(4_000).max(200_000),
reflectionTokens: z.number().int().min(4_000).max(400_000),
recentMessagePercent: z.number().int().min(5).max(75),
asyncBuffering: z.boolean(),
bufferIntervalPercent: z.number().int().min(5).max(90),
bufferOnIdle: z.boolean(),
observationBlockPercent: z.number().int().min(101).max(199),
reflectionBufferPercent: z.number().int().min(10).max(90),
reflectionBlockPercent: z.number().int().min(101).max(199),
optimizeObserverContext: z.boolean(),
previousObserverTokens: z.number().int().min(0).max(100_000),
retrievalEnabled: z.boolean(),
retrievalScope: z.enum(["thread", "resource"]),
temporalMarkers: z.boolean(),
activateAfterIdle: z.enum(["off", "auto", "5m", "1hr", "24hr"]),
activateOnProviderChange: z.boolean(),
shareTokenBudget: z.boolean(),
observeAttachments: z.enum(["auto", "all", "none"]),
observationInstruction: z.string().max(8_000),
reflectionInstruction: z.string().max(8_000),
internalRecall: z.boolean(),
internalRetention: z.boolean(),
}).refine(
(settings) => !settings.shareTokenBudget || !settings.asyncBuffering,
{ message: "Shared token budget requires async buffering to be disabled", path: ["shareTokenBudget"] },
);
const agentSettingsInputSchema: z.ZodType<AgentSettingsInput> = z.object({
name: z.string().trim().min(1).max(120),
description: z.string().trim().min(1).max(1024),
instructions: z.string().min(1).max(MAX_AGENT_INSTRUCTIONS_CHARACTERS),
model: z.string().trim().min(1).max(200).nullable(),
workspaceAccess: z.enum(["read-write", "read-only", "none"]),
browserAccess: z.enum(["interactive", "read-only", "none"]),
delegationEnabled: z.boolean(),
tools: z.array(z.enum(AGENT_TOOL_NAMES)).max(AGENT_TOOL_NAMES.length)
.refine((tools) => new Set(tools).size === tools.length, "Tools must be unique"),
}).strict();
const agentRuntimeSettingsSchema: z.ZodType<AgentRuntimeSettingsInput> = z.object({
defaultModel: z.string().trim().min(1).max(200),
modelSource: z.enum(MODEL_SOURCES),
supervisorMaxSteps: z.number().int().min(1).max(128),
specialistMaxSteps: z.number().int().min(1).max(128),
toolConcurrency: z.number().int().min(1).max(16),
delegationContextMessages: z.number().int().min(1).max(100),
delegationResultCharacters: z.number().int().min(1_000).max(100_000),
maxProcessorRetries: z.number().int().min(0).max(10),
finalResponseFeedback: z.string().trim().min(1).max(4_000),
delegationFailureFeedback: z.string().trim().min(1).max(4_000),
delegationResultTruncationMarker: z.string().min(1).max(500),
taskConcurrency: z.number().int().min(1).max(16),
taskPollIntervalMs: z.number().int().min(250).max(300_000),
taskTimeoutMs: z.number().int().min(1_000).max(86_400_000),
taskStaleAfterMs: z.number().int().min(60_000).max(604_800_000),
}).strict();
const autonomySettingsSchema: z.ZodType<AutonomySettingsInput> = z.object({
enabled: z.boolean(),
reflectionIntervalMs: z.number().int().min(60_000).max(86_400_000),
batchSize: z.number().int().min(1).max(100),
maxAttempts: z.number().int().min(1).max(10),
autoApplyStrategies: z.boolean(),
autoCreateSkills: z.boolean(),
autoRetainFacts: z.boolean(),
selfUpdateEnabled: z.boolean(),
idleImprovementEnabled: z.boolean(),
idleDeploymentEnabled: z.boolean(),
idleWorkspaceIds: z.array(z.string().trim().min(1).max(256)).max(100)
.refine((ids) => new Set(ids).size === ids.length, "Workspace IDs must be unique"),
selfUpdateCron: z.string().trim().min(1).max(100).refine((value) => {
try {
new Cron(value).nextRun();
return true;
} catch {
return false;
}
}, "Invalid cron expression"),
}).strict();
const autonomyHistoryQuerySchema = z.object({
limit: z.coerce.number().int().min(1).max(200).default(100),
}).strict();
const autonomyRevisionIdSchema = z.string().uuid();
const autonomySignalIdSchema = z.string().uuid();
const memoryKeySchema = z.string().trim().min(1).max(128);
const memoryContentSchema = z.string().trim().min(1).max(MAX_FACT_CHARACTERS);
const memoryCreateSchema = z.object({
content: memoryContentSchema,
key: memoryKeySchema.optional(),
importance: z.number().min(0).max(1).default(0.8),
sessionId: z.string().trim().min(1).max(256).default("memory-manager"),
}).strict();
const memoryUpdateSchema = z.object({
content: memoryContentSchema,
key: memoryKeySchema.optional(),
importance: z.number().min(0).max(1),
}).strict();
const memoryRecallSchema = z.object({
query: z.string().trim().min(1).max(4_000),
limit: z.coerce.number().int().min(1).max(20).default(8),
}).strict();
const memoryEpisodeSchema = z.object({
sessionId: z.string().trim().min(1).max(256),
userText: z.string().max(4_000),
assistantText: z.string().max(8_000),
}).strict();
const workspaceInputSchema = z.object({
name: z.string().trim().min(1).max(120),
repositoryPath: z.string().trim().min(1).max(512),
}).strict();
const documentationSaveSchema = z.object({
path: z.string().trim().min(1).max(512),
content: z.string().max(MAX_DOCUMENTATION_CHARACTERS),
revision: z.string().regex(/^[a-f0-9]{64}$/).optional(),
}).strict();
const documentationFolderSchema = z.object({
path: z.string().trim().min(1).max(512),
}).strict();
const documentationDeleteSchema = z.object({
path: z.string().trim().min(1).max(512),
revision: z.string().regex(/^[a-f0-9]{64}$/),
}).strict();
const documentationMoveSchema = documentationDeleteSchema.extend({
nextPath: z.string().trim().min(1).max(512),
}).strict();
const taskResolveSchema = z.object({
output: z.string().trim().min(1).max(32_000),
}).strict();
const sessionPostSchema = z.object({
model: z.string().min(1).max(200).optional(),
workspaceId: z.string().min(1).max(256).default(DEFAULT_WORKSPACE_ID),
}).strict();
const messageSchema = z.object({
id: z.string().optional(),
role: z.enum(["user", "assistant", "system"]),
parts: z.array(z.unknown()),
}).passthrough();
const sessionPutSchema = z.object({
model: z.string().min(1).max(200),
messages: z.array(messageSchema).max(2000),
revision: z.number().int().min(0),
});
const chatPostSchema = z.object({
id: z.string().min(1).max(256).optional(),
model: z.string().min(1).max(200).optional(),
workspaceId: z.string().min(1).max(256).optional(),
messages: z.array(messageSchema).min(1).max(2000),
});
const taskPostSchema = z.object({
prompt: z.string().trim().min(1).max(8000),
model: z.string().min(1).max(200).optional(),
sessionId: z.string().min(1).max(256).optional(),
workspaceId: z.string().min(1).max(256).default(DEFAULT_WORKSPACE_ID),
workflow: z.boolean().default(false),
});
const scheduleSchema = z.object({
name: z.string().trim().min(1).max(120),
prompt: z.string().trim().min(1).max(8000),
model: z.string().min(1).max(200).optional(),
sessionId: z.string().min(1).max(256).optional(),
workspaceId: z.string().min(1).max(256).default(DEFAULT_WORKSPACE_ID),
cron: z.string().trim().min(1).max(120),
enabled: z.boolean().optional(),
});
const MAX_BODY_BYTES = 2 * 1024 * 1024;
async function readJson(req: Request): Promise<unknown> {
const text = await req.text();
if (text.length > MAX_BODY_BYTES) throw new RangeError("body too large");
return JSON.parse(text);
}
function guardApiRoutes<T extends Record<string, unknown>>(routes: T): T {
const guarded = { ...routes } as T;
for (const [path, route] of Object.entries(routes)) {
if (!path.startsWith("/api/")) continue;
if (typeof route === "function") {
(guarded as Record<string, unknown>)[path] = async (req: Request) => requireApiKey(req) ?? route(req);
} else if (route && typeof route === "object") {
const methods = { ...(route as Record<string, unknown>) };
for (const [method, handler] of Object.entries(methods)) {
if (typeof handler === "function") methods[method] = async (req: Request) => requireApiKey(req) ?? handler(req);
}
(guarded as Record<string, unknown>)[path] = methods;
}
}
return guarded;
}
export const sessions = new SessionStore();
const hooks = loadHookRuntime((run) => hookAudit.record(run));
export const taskStore = new AgentTaskStore();
const agentExecution = new AgentExecutionRuntime(
agent,
hooks,
longTermMemory,
browserRuntime,
memorySettings,
agentRuntimeSettings,
repositoryBrief,
);
const taskHookLifecycle = createTaskLifecycleHooks(hooks);
const taskRuntime = new AgentTaskRuntime(
taskStore,
createAgentTaskExecutor(agent, agentExecution, undefined, (task, signal, wake) =>
new TaskAgentCommunicationSession({
taskId: task.id,
workspaceId: task.workspaceId,
resourceId: RESOURCE_ID,
participants: companyAgentProfiles.map((profile) => profile.id),
signal,
wake,
persist: async (communication) => {
const message = await channels.recordAgentCommunication(communication);
try {
channelEvents.message(message.channelId, message.id, null);
} catch (error) {
appLogger().warn("task.communication_refresh_failed", {
taskId: task.id,
messageId: message.id,
error: error instanceof Error ? error.message : String(error),
});
}
},
}), taskParticipantAgents),
async (event) => {
if (event.status === "progress") {
if (!taskUsesChannelTranscript(event.task)) return;
const activity = event.progress?.activity;
// The final turn's text becomes the task outcome and is announced once
// by the completion event below; posting it here too printed every
// report twice in the channel.
if (activity?.isFinal && activity.text.trim()) return;
const content = activity?.text.trim()
|| (activity?.tools.length ? `Using ${activity.tools.join(", ")}` : event.progress?.progress)
|| "Task updated";
const announcement = await channels.announceTaskUpdate(
event.task,
activity?.agentName || "Orchistrator",
content,
);
channelEvents.message(announcement.channelId, announcement.id, null);
return;
}
await taskHookLifecycle(event);
if (event.status === "running") {
const announcement = await channels.announceTaskStart(event.task);
if (announcement) channelEvents.message(announcement.channelId, announcement.id, event.task.id);
} else if (taskUsesChannelTranscript(event.task)) {
const label = event.status === "completed"
? `Task completed\n\n${event.outcome || "Completed without output."}`
: event.status === "cancelled"
? "Task cancelled"
: `Task ${event.status}: ${event.error || "No error detail."}`;
const announcement = await channels.announceTaskUpdate(event.task, "Orchistrator", label);
channelEvents.message(announcement.channelId, announcement.id, null);
}
const evidence = {
sessionId: event.task.sessionId,
turnId: event.turnId,
traceId: event.traceId,
agentId: "orchistrator",
workspaceId: event.task.workspaceId,
};
if (event.status === "completed") {
await evolutionStore.enqueueCompletedTurn({
...evidence,
summary: buildEvolutionEvidence({
goal: event.task.prompt,
outcome: event.outcome || "Completed without task output.",
classification: "success",
}),
}).catch((error) => appLogger().warn("evolution.signal.failed", { error: messageOf(error) }));
} else if (event.status === "failed" || event.status === "dead-letter") {
await evolutionStore.enqueueFailedTurn({
...evidence,
summary: buildEvolutionEvidence({
goal: event.task.prompt,
failure: event.error || "Task failed without error detail.",
classification: event.errorClass ?? event.status,
}),
}).catch((error) => appLogger().warn("evolution.signal.failed", { error: messageOf(error) }));
}
const channelMessage = await channels.findMessageByTask("general", event.task.id);
if (channelMessage) {
const status = event.status === "interrupted" ? "failed" : event.status;
channelEvents.task(channelMessage.channelId, event.task.id, status);
}
if (!event.task.sessionId && event.status !== "running") {
await browserRuntime.sessions.closeThread(`task:${event.task.id}`);
}
},
);
const selfUpdateScheduler = new SelfUpdateScheduler({ tasks: taskStore });
type AgentActivityMessage = UIMessage<unknown, { agentActivity: AgentActivity }>;
export const SERVER_IDLE_TIMEOUT_SECONDS = 255;
export function startServer(
port = 3000,
sessionStore = sessions,
hookRuntime: HookRuntime = hooks,
memoryStore: LongTermMemoryStore = longTermMemory,
agentStore: AgentSettingsStore = agentSettings,
browserStore: BrowserSettingsStore = browserSettings,
activeBrowser: BrowserRuntime = browserRuntime,
skillStore: AgentSkillStore = agentSkills,
tasks: AgentTaskStore = taskStore,
taskRunner: AgentTaskRuntime = taskRuntime,
memoryConfigStore: MemorySettingsStore = memorySettings,
runtimeConfigStore: Pick<AgentRuntimeSettingsStore, "get" | "update"> = agentRuntimeSettings,
documentation: DocumentationIndex = documentationIndex,
secretStore: Pick<typeof agentSecrets, "list" | "create" | "update" | "delete"> = agentSecrets,
autonomyConfigStore: Pick<AutonomySettingsStore, "get" | "update"> = autonomySettings,
evolution: Pick<EvolutionStore, "listSignals" | "manageSignal" | "listRevisions" | "revertRevision" | "enqueueCompletedTurn" | "enqueueFailedTurn"> = evolutionStore,
workspaceStore: AgentWorkspaceStore = agentWorkspaces,
) {
const startedSessions = new Set<string>();
const executionRuntime = new AgentExecutionRuntime(
agent,
hookRuntime,
memoryStore,
activeBrowser,
memoryConfigStore,
runtimeConfigStore,
repositoryBrief,
);
const permanentlyDeleteSession = async (session: SessionSummary) => {
try {
await hookRuntime.dispatch(createHookEvent("SessionEnd", {
sessionId: session.id,
model: session.model,
detail: { source: "delete" },
}));
} catch (error) {
appLogger().warn("hook.session_end.failed", { error: messageOf(error) });
}
await activeBrowser.sessions.closeThread(session.id);
await hookAudit.deleteForSession(session.id);
await sessionStore.delete(session.id);
startedSessions.delete(session.id);
};
type ApiRequest = Request & { params: { id: string; skillId: string; name: string; action: string } };
type ApiHandler = (req: ApiRequest) => Response | Promise<Response>;
type RouteDef = typeof index | ApiHandler | Partial<Record<"GET" | "POST" | "PUT" | "PATCH" | "DELETE", ApiHandler>>;
const routes: Record<string, RouteDef> = {
"/": index,
...createChannelRoutes({
channelStore: channels,
workspaceStore,
runtimeSettings: runtimeConfigStore,
taskStore: tasks,
taskRunner,
}),
...createCodeIntelligenceRoutes({ workspaceStore }),
"/api/health": async () => {
let applicationDatabase: "ok" | "error" = "ok";
try {
await Promise.all([runtimeConfigStore.get(), workspaceStore.list()]);
} catch {
applicationDatabase = "error";
}
const observabilityDatabase = (await observabilityHealth()).status;
const status = applicationDatabase === "ok" && observabilityDatabase === "ok" ? "ok" : "error";
return Response.json(
{ status, applicationDatabase, observabilityDatabase },
{ status: status === "ok" ? 200 : 503 },
);
},
"/api/models": async () => {
const runtimeSettings = await runtimeConfigStore.get();
return Response.json({ ...(await listModelCatalog()), defaultModel: runtimeSettings.defaultModel });
},
"/api/workspaces": {
GET: async () => Response.json({ workspaces: await workspaceStore.list() }),
POST: async (req) => {
const input = workspaceInputSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) {
return Response.json({ error: "Invalid workspace", issues: input.error.issues }, { status: 400 });
}
try {
return Response.json(await workspaceStore.create(input.data), { status: 201 });
} catch (error) {
return Response.json({ error: messageOf(error) }, { status: 400 });
}
},
},
"/api/workspaces/:id/repo-brief": async (req) => {
if (!(await workspaceStore.get(req.params.id))) return Response.json({ error: "Workspace not found" }, { status: 404 });
const query = new URL(req.url).searchParams.get("q")?.trim() ?? "";
if (!query || query.length > 2_000) return Response.json({ error: "Invalid repository brief query" }, { status: 400 });
try {
return Response.json(await repositoryBrief.build(req.params.id, query));
} catch (error) {
return Response.json({ error: messageOf(error) }, { status: 400 });
}
},
"/api/docs/search": async (req) => {
const url = new URL(req.url);
const query = url.searchParams.get("q")?.trim() ?? "";
const workspaceId = url.searchParams.get("workspaceId")?.trim() || undefined;
if (!query || query.length > 2_000) return Response.json({ error: "Invalid documentation query" }, { status: 400 });
if (workspaceId && !(await workspaceStore.get(workspaceId))) {
return Response.json({ error: "Workspace not found" }, { status: 404 });
}
return Response.json({ results: await documentation.search(query, { workspaceId }) });
},
"/api/workspaces/:id/docs": {
GET: async (req) => {
if (!(await workspaceStore.get(req.params.id))) return Response.json({ error: "Workspace not found" }, { status: 404 });
return Response.json(await documentation.files.list(req.params.id));
},
POST: async (req) => {
if (!(await workspaceStore.get(req.params.id))) return Response.json({ error: "Workspace not found" }, { status: 404 });
const input = documentationSaveSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) return Response.json({ error: "Invalid documentation page", issues: input.error.issues }, { status: 400 });
try {
const page = await documentation.files.save(req.params.id, input.data);
let indexed = true;
await documentation.indexPage(req.params.id, page).catch((error) => {
indexed = false;
appLogger().warn("documentation.index.page.failed", { workspaceId: req.params.id, path: page.path, error });
});
return Response.json(page, {
status: input.data.revision ? 200 : 201,
headers: { "x-popagent-indexed": String(indexed) },
});
} catch (error) {
if (error instanceof DocumentationConflictError) {
return Response.json({ error: "revision conflict", page: error.current }, { status: 409 });
}
return Response.json({ error: messageOf(error) }, { status: 400 });
}
},
PATCH: async (req) => {
if (!(await workspaceStore.get(req.params.id))) return Response.json({ error: "Workspace not found" }, { status: 404 });
const input = documentationMoveSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) return Response.json({ error: "Invalid documentation move", issues: input.error.issues }, { status: 400 });
try {
const page = await documentation.files.move(
req.params.id,
input.data.path,
input.data.nextPath,
input.data.revision,
);
await documentation.removePage(req.params.id, input.data.path).catch((error) => {
appLogger().warn("documentation.index.move-delete.failed", { workspaceId: req.params.id, path: input.data.path, error });
});
let indexed = true;
await documentation.indexPage(req.params.id, page).catch((error) => {
indexed = false;
appLogger().warn("documentation.index.move.failed", { workspaceId: req.params.id, path: page.path, error });
});
return Response.json(page, { headers: { "x-popagent-indexed": String(indexed) } });
} catch (error) {
if (error instanceof DocumentationConflictError) {
return Response.json({ error: "revision conflict", page: error.current }, { status: 409 });
}
return Response.json({ error: messageOf(error) }, { status: 400 });
}
},
DELETE: async (req) => {
if (!(await workspaceStore.get(req.params.id))) return Response.json({ error: "Workspace not found" }, { status: 404 });
const input = documentationDeleteSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) return Response.json({ error: "Invalid documentation page", issues: input.error.issues }, { status: 400 });
try {
await documentation.files.remove(req.params.id, input.data.path, input.data.revision);
await documentation.removePage(req.params.id, input.data.path).catch((error) => {
appLogger().warn("documentation.index.delete.failed", { workspaceId: req.params.id, path: input.data.path, error });
});
return new Response(null, { status: 204 });
} catch (error) {
if (error instanceof DocumentationConflictError) {
return Response.json({ error: "revision conflict", page: error.current }, { status: 409 });
}
return Response.json({ error: messageOf(error) }, { status: 400 });
}
},
},
"/api/workspaces/:id/docs/folders": {
POST: async (req) => {
if (!(await workspaceStore.get(req.params.id))) return Response.json({ error: "Workspace not found" }, { status: 404 });
const input = documentationFolderSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) return Response.json({ error: "Invalid documentation folder", issues: input.error.issues }, { status: 400 });
try {
return Response.json(await documentation.files.createFolder(req.params.id, input.data.path), { status: 201 });
} catch (error) {
return Response.json({ error: messageOf(error) }, { status: 400 });
}
},
DELETE: async (req) => {
if (!(await workspaceStore.get(req.params.id))) return Response.json({ error: "Workspace not found" }, { status: 404 });
const input = documentationFolderSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) return Response.json({ error: "Invalid documentation folder", issues: input.error.issues }, { status: 400 });
try {
const paths = await documentation.files.removeFolder(req.params.id, input.data.path);
await Promise.all(paths.map((path) => documentation.removePage(req.params.id, path)));
return new Response(null, { status: 204 });
} catch (error) {
return Response.json({ error: messageOf(error) }, { status: 400 });
}
},
},
"/api/workspaces/:id/docs/content": async (req) => {
if (!(await workspaceStore.get(req.params.id))) return Response.json({ error: "Workspace not found" }, { status: 404 });
const path = new URL(req.url).searchParams.get("path");
if (!path) return Response.json({ error: "Documentation path is required" }, { status: 400 });
try {
return Response.json(await documentation.files.read(req.params.id, path));
} catch (error) {
return Response.json({ error: messageOf(error) }, { status: 404 });
}
},
"/api/workspaces/:id/docs/reindex": {
POST: async (req) => {
if (!(await workspaceStore.get(req.params.id))) return Response.json({ error: "Workspace not found" }, { status: 404 });
try {
return Response.json(await documentation.indexWorkspace(req.params.id));
} catch (error) {
appLogger().warn("documentation.index.workspace.failed", { workspaceId: req.params.id, error });
return Response.json({ error: "Unable to index documentation" }, { status: 503 });
}
},
},
"/api/workspaces/:id/docs/index-status": async (req) => {
if (!(await workspaceStore.get(req.params.id))) return Response.json({ error: "Workspace not found" }, { status: 404 });
return Response.json(await documentation.status(req.params.id));
},
"/api/workspaces/:id/sessions": {
DELETE: async (req) => {
if (!(await workspaceStore.get(req.params.id))) {
return Response.json({ error: "Workspace not found" }, { status: 404 });
}
const workspaceSessions = await sessionStore.listForWorkspace(req.params.id);
for (const session of workspaceSessions) await permanentlyDeleteSession(session);
return Response.json({ deleted: workspaceSessions.length });
},
},
"/api/workspaces/:id": {
PATCH: async (req) => {
const input = workspaceInputSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) {
return Response.json({ error: "Invalid workspace", issues: input.error.issues }, { status: 400 });
}
try {
const workspace = await workspaceStore.update(req.params.id, input.data);
return workspace
? Response.json(workspace)
: Response.json({ error: "Workspace not found" }, { status: 404 });
} catch (error) {
return Response.json({ error: messageOf(error) }, { status: 400 });
}
},
DELETE: async (req) => {
if (req.params.id === DEFAULT_WORKSPACE_ID) {
return Response.json({ error: "The default workspace cannot be deleted" }, { status: 409 });
}
if (!(await workspaceStore.get(req.params.id))) {
return Response.json({ error: "Workspace not found" }, { status: 404 });
}
const [active, archived, workspaceTasks, schedules] = await Promise.all([
sessionStore.list(false, req.params.id),
sessionStore.list(true, req.params.id),
tasks.listTasks(req.params.id),
tasks.listSchedules(req.params.id),
]);
if (active.length || archived.length || workspaceTasks.length || schedules.length || await channels.hasWorkspaceMessages(req.params.id)) {
return Response.json({ error: "Workspace still has sessions, tasks, schedules, or channel history" }, { status: 409 });
}
await documentation.removeWorkspace(req.params.id);
await workspaceStore.delete(req.params.id);
return new Response(null, { status: 204 });
},
},
"/api/company": {
GET: () => Response.json({ departments: companyDepartments }),
},
"/api/company/avatars/:id": {
GET: (req) => companyAvatarIds.has(req.params.id)
? new Response(Bun.file(new URL(`./company-avatars/${req.params.id}.webp`, import.meta.url)))
: Response.json({ error: "Company avatar not found" }, { status: 404 }),
},
"/api/agents": {
GET: async () => Response.json({ agents: await agentStore.list() }),
},
"/api/agents/:id": {
PATCH: async (req) => {
const input = agentSettingsInputSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) {
return Response.json({ error: "Invalid agent settings", issues: input.error.issues }, { status: 400 });
}
if (input.data.model && !(await isKnownModel(input.data.model))) {
return Response.json({ error: `Unknown model: ${input.data.model}` }, { status: 400 });
}
try {
const settings = await agentStore.update(req.params.id, input.data);
return settings
? Response.json(settings)
: Response.json({ error: "Agent not found" }, { status: 404 });
} catch (error) {
return Response.json({ error: messageOf(error) }, { status: 400 });
}
},
},
"/api/settings/browser": {
GET: async () => Response.json(await browserStore.get()),
PATCH: async (req) => {
const input = browserSettingsSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) {
return Response.json({ error: "Invalid browser settings", issues: input.error.issues }, { status: 400 });
}
const settings = await browserStore.update(input.data);
await browserRecordings.cleanup({ retentionDays: settings.recordingRetentionDays, maxFiles: settings.recordingMaxFiles });
await activeBrowser.apply(agent, settings);
return Response.json(settings);
},
},
"/api/browser/test": {
POST: async () => {
try {
await activeBrowser.apply(agent, await browserStore.get());
return Response.json(await activeBrowser.test());
} catch (error) {
return Response.json({ error: messageOf(error) }, { status: 503 });
}
},
},
"/api/browser/health": {
GET: async () => {
const settings = await browserStore.get();
const [activeSessions, profiles, health] = await Promise.all([
activeBrowser.sessions.list(),
browserProfiles.list(),
activeBrowser.health(settings),
]);
return Response.json({
...health,
enabled: settings.enabled,
scope: settings.scope,
activeSessions: activeSessions.length,
profileConfigured: profiles.some((profile) => profile.enabled),
screencastEnabled: settings.screencastEnabled,
recordingEnabled: settings.recordingEnabled,
});
},
},
"/api/browser/profiles": {
GET: async () => Response.json({ profiles: await browserProfiles.list() }),
POST: async (req) => {
const input = browserProfileSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) return Response.json({ error: "Invalid browser profile", issues: input.error.issues }, { status: 400 });
try {
const profile = await browserProfiles.create(input.data);
activeBrowser.invalidate();
await activeBrowser.apply(agent, await browserStore.get());
return Response.json(profile, { status: 201 });
} catch (error) {
return Response.json({ error: messageOf(error) }, { status: 400 });
}
},
},
"/api/browser/profiles/:id": {
PATCH: async (req) => {
const input = browserProfileUpdateSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) return Response.json({ error: "Invalid browser profile", issues: input.error.issues }, { status: 400 });
try {
const profile = await browserProfiles.update(req.params.id, input.data);
if (!profile) return Response.json({ error: "Browser profile not found" }, { status: 404 });
activeBrowser.invalidate();
await activeBrowser.apply(agent, await browserStore.get());
return Response.json(profile);
} catch (error) {
return Response.json({ error: messageOf(error) }, { status: 400 });
}
},
DELETE: async (req) => {
if (!await browserProfiles.delete(req.params.id)) return Response.json({ error: "Browser profile not found" }, { status: 404 });
activeBrowser.invalidate();
await activeBrowser.apply(agent, await browserStore.get());
return new Response(null, { status: 204 });
},
},
"/api/browser/recordings": {
GET: async () => {
const settings = await browserStore.get();
await browserRecordings.cleanup({ retentionDays: settings.recordingRetentionDays, maxFiles: settings.recordingMaxFiles });
return Response.json({ recordings: await browserRecordings.list() });
},
},
"/api/browser/recordings/:id": {
GET: async (req) => {
try {
const path = await browserRecordings.pathFor(req.params.id);
return new Response(Bun.file(path), { headers: { "content-type": "video/x-msvideo", "content-disposition": `attachment; filename="${req.params.id}"` } });
} catch {
return Response.json({ error: "Recording not found" }, { status: 404 });
}
},
DELETE: async (req) => await browserRecordings.delete(req.params.id)
? new Response(null, { status: 204 })
: Response.json({ error: "Recording not found" }, { status: 404 }),
},
"/api/browser/sessions": {
GET: async (req) => {
await activeBrowser.sessions.cleanupIdle();
const threadId = new URL(req.url).searchParams.get("threadId") || undefined;
return Response.json({ sessions: await activeBrowser.sessions.list(threadId) });
},
},
"/api/browser/sessions/:id/close": {
POST: async (req) => await activeBrowser.sessions.close(req.params.id)
? new Response(null, { status: 204 })
: Response.json({ error: "Browser session not found" }, { status: 404 }),
},
"/api/browser/sessions/:id/takeover": {
PATCH: async (req) => {
const body = await readJson(req).catch(() => undefined);
const input = z.object({ enabled: z.boolean() }).strict().safeParse(body);
if (!input.success) return Response.json({ error: "Invalid takeover setting" }, { status: 400 });
return activeBrowser.sessions.setTakeover(req.params.id, input.data.enabled)
? Response.json({ enabled: input.data.enabled })
: Response.json({ error: "Interactive browser session not found" }, { status: 404 });
},
},
"/api/browser/sessions/:id/input": {
POST: async (req) => {
const input = browserInputSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) return Response.json({ error: "Invalid browser input", issues: input.error.issues }, { status: 400 });
try {
if (input.data.kind === "mouse") await activeBrowser.sessions.injectMouse(req.params.id, input.data.event);
else await activeBrowser.sessions.injectKeyboard(req.params.id, input.data.event);
return new Response(null, { status: 204 });
} catch (error) {
const message = messageOf(error);
return Response.json({ error: message }, { status: message.includes("takeover") ? 409 : 404 });
}
},
},
"/api/browser/sessions/:id/stream": {
GET: async (req) => {
if (!(await browserStore.get()).screencastEnabled) {
return Response.json({ error: "Browser screencast is disabled" }, { status: 409 });
}
try {
const screencast = await activeBrowser.sessions.startScreencast(req.params.id);
const encoder = new TextEncoder();
let closed = false;
const stream = new ReadableStream<Uint8Array>({
start(controller) {
const send = (event: string, value: unknown) => {
if (!closed) controller.enqueue(encoder.encode(`event: ${event}\ndata: ${JSON.stringify(value)}\n\n`));
};
screencast.on("frame", (frame) => send("frame", frame));
screencast.on("url", (url) => send("url", { url }));
screencast.on("error", (error) => send("error", { error: error.message }));
screencast.on("stop", (reason) => {
if (closed) return;
send("stop", { reason });
closed = true;
controller.close();
});
},
async cancel() {
closed = true;
await screencast.stop();
},
});
return new Response(stream, {
headers: {
"content-type": "text/event-stream; charset=utf-8",
"cache-control": "no-cache, no-transform",
"x-accel-buffering": "no",
},
});
} catch (error) {
return Response.json({ error: messageOf(error) }, { status: 404 });
}
},
},
"/api/secrets": {
GET: async () => {
try {
return Response.json({ secrets: await secretStore.list(RESOURCE_ID) });
} catch (error) {
appLogger().error("secret.list.failed", { error: messageOf(error) });
return Response.json({ error: "Secrets could not be listed" }, { status: 500 });
}
},
POST: async (req) => {
const input = agentSecretInputSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) {
return Response.json({ error: "Invalid secret", issues: input.error.issues }, { status: 400 });
}
try {
const secret = await secretStore.create(RESOURCE_ID, input.data);
return secret
? Response.json(secret, { status: 201 })
: Response.json({ error: "Secret already exists" }, { status: 409 });
} catch (error) {
appLogger().error("secret.create.failed", { error: messageOf(error) });
return Response.json({ error: "Secret could not be stored" }, { status: 500 });
}
},
},
"/api/secrets/:name": {
PATCH: async (req) => {
const name = secretNameSchema.safeParse(req.params.name);
const input = agentSecretUpdateSchema.safeParse(await readJson(req).catch(() => undefined));
if (!name.success || !input.success) {
return Response.json({ error: "Invalid secret update" }, { status: 400 });
}
try {
const secret = await secretStore.update(RESOURCE_ID, name.data, input.data);
return secret
? Response.json(secret)
: Response.json({ error: "Secret not found" }, { status: 404 });
} catch (error) {
appLogger().error("secret.update.failed", { error: messageOf(error) });
return Response.json({ error: "Secret could not be updated" }, { status: 500 });
}
},
DELETE: async (req) => {
const name = secretNameSchema.safeParse(req.params.name);
if (!name.success) return Response.json({ error: "Invalid secret name" }, { status: 400 });
try {
return await secretStore.delete(RESOURCE_ID, name.data)
? new Response(null, { status: 204 })
: Response.json({ error: "Secret not found" }, { status: 404 });
} catch (error) {
appLogger().error("secret.delete.failed", { error: messageOf(error) });
return Response.json({ error: "Secret could not be deleted" }, { status: 500 });
}
},
},
"/api/settings/memory": {
GET: async () => Response.json(await memoryConfigStore.get()),
PATCH: async (req) => {
const input = memorySettingsSchema.safeParse(await req.json().catch(() => ({})));
if (!input.success) {
return Response.json({ error: "Invalid memory settings", issues: input.error.issues }, { status: 400 });
}
return Response.json(await memoryConfigStore.update(input.data));
},
},
"/api/settings/agent-runtime": {
GET: async () => Response.json(await runtimeConfigStore.get()),
PATCH: async (req) => {
const input = agentRuntimeSettingsSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) {
return Response.json({ error: "Invalid agent runtime settings", issues: input.error.issues }, { status: 400 });
}
const model = runtimeModelId(input.data.defaultModel, input.data.modelSource);
if (!(await isKnownModel(model))) {
return Response.json({ error: `Unknown model: ${model}` }, { status: 400 });
}
try {
const settings = await runtimeConfigStore.update(input.data);
taskRunner.configure(settings);
return Response.json(settings);
} catch (error) {
return Response.json({ error: messageOf(error) }, { status: 400 });
}
},
},
"/api/settings/autonomy": {
GET: async () => Response.json(await autonomyConfigStore.get()),
PATCH: async (req) => {
const input = autonomySettingsSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) {
return Response.json({ error: "Invalid autonomy settings", issues: input.error.issues }, { status: 400 });
}
return Response.json(await autonomyConfigStore.update(input.data));
},
},
"/api/autonomy/status": async () => Response.json(await tasks.automationStatus()),
"/api/autonomy/run": {
POST: async (req) => {
const input = z.object({ workspaceId: z.string().min(1).max(256) }).strict()
.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) return Response.json({ error: "Invalid automation run", issues: input.error.issues }, { status: 400 });
const task = await tasks.queueSelfUpdate(input.data.workspaceId);
if (!task) return Response.json({ error: "Workspace is not enabled or automation is busy" }, { status: 409 });
taskRunner.enqueue(task);
return Response.json(task, { status: 201 });
},
},
"/api/autonomy/reconcile": {
POST: async () => {
await selfUpdateScheduler.reconcile();
return Response.json({ reconciled: true });
},
},
"/api/autonomy/signals": async (req) => {
const query = autonomyHistoryQuerySchema.safeParse(Object.fromEntries(new URL(req.url).searchParams));
if (!query.success) {
return Response.json({ error: "Invalid autonomy history query", issues: query.error.issues }, { status: 400 });
}
return Response.json({ signals: await evolution.listSignals(query.data.limit) });
},
"/api/autonomy/signals/:id/:action": {
POST: async (req) => {
const id = autonomySignalIdSchema.safeParse(req.params.id);
const action = z.enum(["ignore", "retry"]).safeParse(req.params.action);
if (!id.success || !action.success) {
return Response.json({ error: "Invalid evolution signal operation" }, { status: 400 });
}
const signal = await evolution.manageSignal(id.data, action.data);
if (signal === "conflict") return Response.json({ error: `Signal cannot be ${action.data}d from its current state` }, { status: 409 });
return signal
? Response.json(signal)
: Response.json({ error: "Evolution signal not found" }, { status: 404 });
},
},
"/api/autonomy/revisions": async (req) => {
const query = autonomyHistoryQuerySchema.safeParse(Object.fromEntries(new URL(req.url).searchParams));
if (!query.success) {
return Response.json({ error: "Invalid autonomy history query", issues: query.error.issues }, { status: 400 });
}
return Response.json({ revisions: await evolution.listRevisions(query.data.limit) });
},
"/api/autonomy/revisions/:id/revert": {
POST: async (req) => {
const id = autonomyRevisionIdSchema.safeParse(req.params.id);
if (!id.success) {
return Response.json({ error: "Invalid evolution revision id", issues: id.error.issues }, { status: 400 });
}
const revision = await evolution.revertRevision(id.data);
return revision
? Response.json(revision)
: Response.json({ error: "Evolution revision not found" }, { status: 404 });
},
},
"/api/agents/health": async () => {
const agents = await agentStore.list();
const skills = (await Promise.all(agents.map(({ id }) => skillStore.list(id)))).flat();
return Response.json(validateAgentHealth(agents, skills));
},
"/api/agents/:id/skills": {
GET: async (req) => Response.json({ skills: await skillStore.list(req.params.id) }),
POST: async (req) => {
const input = skillInputSchema.safeParse(await req.json().catch(() => ({})));
if (!input.success) return Response.json({ error: "Invalid skill", issues: input.error.issues }, { status: 400 });
return Response.json(await skillStore.create(req.params.id, input.data), { status: 201 });
},
},
"/api/agents/:id/skills/:skillId": {
PATCH: async (req) => {
const input = skillInputSchema.safeParse(await req.json().catch(() => ({})));
if (!input.success) return Response.json({ error: "Invalid skill", issues: input.error.issues }, { status: 400 });
const skill = await skillStore.update(req.params.id, req.params.skillId, input.data);
return skill ? Response.json(skill) : Response.json({ error: "Skill not found" }, { status: 404 });
},
DELETE: async (req) => await skillStore.delete(req.params.id, req.params.skillId)
? new Response(null, { status: 204 })
: Response.json({ error: "Skill not found" }, { status: 404 }),
},
"/api/sessions": {
GET: async (req) => {
const params = new URL(req.url).searchParams;
const archived = params.get("archived") === "true";
const workspaceId = params.get("workspaceId") || undefined;
return Response.json({ sessions: await sessionStore.list(archived, workspaceId) });
},
POST: async (req) => {
const input = sessionPostSchema.safeParse(await readJson(req).catch(() => ({})));
if (!input.success) {
return Response.json({ error: "Invalid session", issues: input.error.issues }, { status: 400 });
}
const model = input.data.model ?? (await runtimeConfigStore.get()).defaultModel;
if (!(await isKnownModel(model))) {
return Response.json({ error: `Unknown model: ${model}` }, { status: 400 });
}
if (!(await workspaceStore.get(input.data.workspaceId))) {
return Response.json({ error: "Workspace not found" }, { status: 404 });
}
const session = await sessionStore.create(model, input.data.workspaceId);
try {
await hookRuntime.dispatch(createHookEvent("SessionStart", {
sessionId: session.id,
model,
detail: {
source: "create",
session: { title: session.title, model, workspaceId: session.workspaceId },
},
}));
startedSessions.add(session.id);
} catch (error) {
appLogger().warn("hook.session_start.failed", { error: messageOf(error) });
}
return Response.json(session, { status: 201 });
},
},
"/api/sessions/:id": {
GET: async (req) => {
const session = await sessionStore.get(req.params.id);
return session
? Response.json(session)
: Response.json({ error: "Session not found" }, { status: 404 });
},
PUT: async (req) => {
const input = sessionPutSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) {
return Response.json({ error: "Invalid session", issues: input.error.issues }, { status: 400 });
}
if (!(await isKnownModel(input.data.model))) {
return Response.json({ error: `Unknown model: ${input.data.model}` }, { status: 400 });
}
const session = await sessionStore.save(req.params.id, {
model: input.data.model,
messages: input.data.messages as ChatMessage[],
expectedRevision: input.data.revision,
});
if (session === "conflict") {
const current = await sessionStore.get(req.params.id);
return Response.json({ error: "revision conflict", session: current }, { status: 409 });
}
return session
? Response.json(session)
: Response.json({ error: "Session not found" }, { status: 404 });
},
PATCH: async (req) => {
const body = (await req.json().catch(() => ({}))) as {
title?: string;
archived?: boolean;
};
if (typeof body.title === "string") {
const title = body.title.trim();
if (!title) return Response.json({ error: "title is required" }, { status: 400 });
const session = await sessionStore.rename(req.params.id, title);
return session
? Response.json(session)
: Response.json({ error: "Session not found" }, { status: 404 });
}
if (typeof body.archived === "boolean") {
const session = await sessionStore.setArchived(req.params.id, body.archived);
return session
? Response.json(session)
: Response.json({ error: "Session not found" }, { status: 404 });
}
return Response.json({ error: "title or archived is required" }, { status: 400 });
},
DELETE: async (req) => {
const session = await sessionStore.get(req.params.id);
if (!session) return Response.json({ error: "Session not found" }, { status: 404 });
await permanentlyDeleteSession(session);
return new Response(null, { status: 204 });
},
},
"/api/sessions/:id/hooks": {
GET: async (req) => Response.json({ runs: await hookAudit.listForSession(req.params.id) }),
},
"/api/memories/status": {
GET: async () => Response.json(await memoryStore.status(RESOURCE_ID)),
},
"/api/memories": {
GET: async (req) => {
const params = new URL(req.url).searchParams;
const kind = params.get("kind");
if (kind !== null && kind !== "fact" && kind !== "episode") {
return Response.json({ error: "kind must be fact or episode" }, { status: 400 });
}
return Response.json({
memories: await memoryStore.list({
resourceId: RESOURCE_ID,
query: params.get("query") ?? undefined,
kind: kind ?? undefined,
}),
});
},
POST: async (req) => {
const input = memoryCreateSchema.safeParse(await req.json().catch(() => ({})));
if (!input.success) {
return Response.json({ error: "Invalid memory", issues: input.error.issues }, { status: 400 });
}
await memoryStore.rememberFact({
resourceId: RESOURCE_ID,
sessionId: input.data.sessionId,
key: input.data.key,
content: input.data.content,
importance: input.data.importance,
});
const memories = await memoryStore.list({
resourceId: RESOURCE_ID,
query: input.data.content,
kind: "fact",
limit: 1,
});
return Response.json(memories[0], { status: 201 });
},
},
"/api/memories/recall": {
GET: async (req) => {
const input = memoryRecallSchema.safeParse(Object.fromEntries(new URL(req.url).searchParams));
if (!input.success) {
return Response.json({ error: "Invalid memory recall", issues: input.error.issues }, { status: 400 });
}
return Response.json({
memories: await memoryStore.recall({ resourceId: RESOURCE_ID, ...input.data }),
});
},
},
"/api/memories/episodes": {
POST: async (req) => {
const input = memoryEpisodeSchema.safeParse(await req.json().catch(() => ({})));
if (!input.success) {
return Response.json({ error: "Invalid memory episode", issues: input.error.issues }, { status: 400 });
}
await memoryStore.retainEpisode({ resourceId: RESOURCE_ID, ...input.data });
return new Response(null, { status: 204 });
},
},
"/api/memories/:id": {
PATCH: async (req) => {
const input = memoryUpdateSchema.safeParse(await req.json().catch(() => ({})));
if (!input.success) {
return Response.json({ error: "Invalid memory", issues: input.error.issues }, { status: 400 });
}
const memory = await memoryStore.updateFact({
id: req.params.id,
resourceId: RESOURCE_ID,
key: input.data.key,
content: input.data.content,
importance: input.data.importance,
});
return memory
? Response.json(memory)
: Response.json({ error: "Editable fact not found" }, { status: 404 });
},
DELETE: async (req) => {
const deleted = await memoryStore.delete(req.params.id, RESOURCE_ID);
return deleted
? new Response(null, { status: 204 })
: Response.json({ error: "Memory not found" }, { status: 404 });
},
},
"/api/tasks": {
GET: async (req) => {
const query = new URL(req.url).searchParams;
const workspaceId = query.get("workspaceId") || undefined;
return Response.json({ tasks: await tasks.listTasks(workspaceId, query.get("workflow") === "true") });
},
POST: async (req) => {
const input = taskPostSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) return Response.json({ error: "Invalid task", issues: input.error.issues }, { status: 400 });
const model = input.data.model ?? (await runtimeConfigStore.get()).defaultModel;
if (!(await isKnownModel(model))) return Response.json({ error: `Unknown model: ${model}` }, { status: 400 });
if (!(await workspaceStore.get(input.data.workspaceId))) return Response.json({ error: "Workspace not found" }, { status: 404 });
const task = await (input.data.workflow ? tasks.createWorkflowTask({
sessionId: input.data.sessionId,
workspaceId: input.data.workspaceId,
prompt: input.data.prompt,
model,
}) : tasks.createTask({
sessionId: input.data.sessionId,
workspaceId: input.data.workspaceId,
prompt: input.data.prompt,
model,
}));
taskRunner.enqueue(task);
return Response.json(task, { status: 201 });
},
},
"/api/tasks/:id": {
GET: async (req) => {
const task = await tasks.getTask(req.params.id);
return task
? Response.json(task)
: Response.json({ error: "Task not found" }, { status: 404 });
},
DELETE: async (req) => {
const result = await tasks.deleteTerminalTask(req.params.id);
if (result === "deleted") return new Response(null, { status: 204 });
if (result === "active") return Response.json({ error: "Active tasks cannot be removed" }, { status: 409 });
return Response.json({ error: "Task not found" }, { status: 404 });
},
},
"/api/tasks/:id/workflow": {
GET: async (req) => {
const workflow = await tasks.getWorkflow(req.params.id);
if (!workflow) return Response.json({ error: "Self-update workflow not found" }, { status: 404 });
return Response.json(workflow);
},
},
"/api/tasks/:id/cancel": {
POST: async (req) => await taskRunner.cancel(req.params.id)
? new Response(null, { status: 204 })
: Response.json({ error: "Task cannot be cancelled" }, { status: 409 }),
},
"/api/tasks/:id/resolve": {
POST: async (req) => {
const input = taskResolveSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) return Response.json({ error: "Invalid task resolution", issues: input.error.issues }, { status: 400 });
const result = await tasks.resolve(req.params.id, input.data.output);
if (result === "missing") return Response.json({ error: "Task not found" }, { status: 404 });
if (result === "active") return Response.json({ error: "Only failed, cancelled, or dead-lettered tasks can be resolved" }, { status: 409 });
const task = await tasks.getTask(req.params.id);
return Response.json(task);
},
},
"/api/schedules": {
GET: async (req) => {
const workspaceId = new URL(req.url).searchParams.get("workspaceId") || undefined;
return Response.json({ schedules: await tasks.listSchedules(workspaceId) });
},
POST: async (req) => {
const input = scheduleSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) return Response.json({ error: "Invalid schedule", issues: input.error.issues }, { status: 400 });
const model = input.data.model ?? (await runtimeConfigStore.get()).defaultModel;
if (!(await isKnownModel(model))) return Response.json({ error: `Unknown model: ${model}` }, { status: 400 });
if (!(await workspaceStore.get(input.data.workspaceId))) return Response.json({ error: "Workspace not found" }, { status: 404 });
try {
return Response.json(await tasks.createSchedule({
sessionId: input.data.sessionId,
workspaceId: input.data.workspaceId,
name: input.data.name,
prompt: input.data.prompt,
model,
cron: input.data.cron,
enabled: input.data.enabled ?? true,
}), { status: 201 });
} catch (error) {
return Response.json({ error: messageOf(error) }, { status: 400 });
}
},
},
"/api/schedules/:id": {
PATCH: async (req) => {
const input = scheduleSchema.extend({ enabled: z.boolean() }).safeParse(await readJson(req).catch(() => undefined));
if (!input.success) return Response.json({ error: "Invalid schedule", issues: input.error.issues }, { status: 400 });
const model = input.data.model ?? (await runtimeConfigStore.get()).defaultModel;
if (!(await isKnownModel(model))) return Response.json({ error: `Unknown model: ${model}` }, { status: 400 });
if (!(await workspaceStore.get(input.data.workspaceId))) return Response.json({ error: "Workspace not found" }, { status: 404 });
try {
const schedule = await tasks.updateSchedule(req.params.id, {
workspaceId: input.data.workspaceId,
name: input.data.name,
prompt: input.data.prompt,
model,
cron: input.data.cron,
enabled: input.data.enabled,
});
return schedule ? Response.json(schedule) : Response.json({ error: "Schedule not found" }, { status: 404 });
} catch (error) {
return Response.json({ error: messageOf(error) }, { status: 400 });
}
},
DELETE: async (req) => await tasks.deleteSchedule(req.params.id)
? new Response(null, { status: 204 })
: Response.json({ error: "Schedule not found" }, { status: 404 }),
},
"/api/observability/overview": async (req) => {
try {
return Response.json(await getObservabilityOverview(parseObservabilityQuery(req.url)));
} catch (error) {
return Response.json({ error: messageOf(error) }, { status: 400 });
}
},
"/api/observability/traces": async (req) => {
try {
return Response.json(await listObservabilityTraces(parseObservabilityQuery(req.url)));
} catch (error) {
return Response.json({ error: messageOf(error) }, { status: 400 });
}
},
"/api/observability/traces/:id": async (req) => {
if (!/^[a-f0-9]{32}$/.test(req.params.id)) {
return Response.json({ error: "Invalid trace ID" }, { status: 400 });
}
const trace = await getObservabilityTrace(req.params.id);
return trace ? Response.json(trace) : Response.json({ error: "Trace not found" }, { status: 404 });
},
"/api/observability/logs": async (req) => {
try {
return Response.json(await listObservabilityLogs(parseObservabilityQuery(req.url)));
} catch (error) {
return Response.json({ error: messageOf(error) }, { status: 400 });
}
},
"/api/observability/health": async () => {
const health = await observabilityHealth();
return Response.json(health, { status: health.status === "ok" ? 200 : 503 });
},
"/api/observability/feedback": {
POST: async (req) => {
const input = observabilityFeedbackSchema.safeParse(await readJson(req).catch(() => undefined));
if (!input.success) {
return Response.json({ error: "Invalid feedback", issues: input.error.issues }, { status: 400 });
}
if (!await getObservabilityTrace(input.data.traceId)) {
return Response.json({ error: "Trace not found" }, { status: 404 });
}
await observability.addFeedback({
traceId: input.data.traceId,
feedback: {
feedbackSource: "user",
feedbackType: "thumbs",
value: input.data.value,
comment: input.data.comment,
},
});
return Response.json({ recorded: true });
},
},
"/api/chat": {
POST: async (req) => {
let evidenceContext: {
sessionId: string;
turnId: string;
traceId: string;
agentId: string;
workspaceId: string;
} | undefined;
let releaseInteractive: (() => void) | undefined;
let evidenceGoal: string | undefined;
let evidenceFailureSignaled = false;
let streamStarted = false;
try {
const parsed = chatPostSchema.safeParse(await readJson(req).catch(() => undefined));
if (!parsed.success) return Response.json({ error: "Invalid chat request", issues: parsed.error.issues }, { status: 400 });
const body = parsed.data;
const model = body.model ?? (await runtimeConfigStore.get()).defaultModel;
if (!(await isKnownModel(model))) return Response.json({ error: `Unknown model: ${model}` }, { status: 400 });
const sessionId = body.id ?? crypto.randomUUID();
const turnId = crypto.randomUUID();
const traceId = createTraceId();
const messages = sanitizeSensitiveToolMessages(body.messages as ChatMessage[]);
const userText = latestUserText(messages);
if (!userText) return Response.json({ error: "A user message is required" }, { status: 400 });
const existing = body.id ? await sessionStore.get(body.id) : undefined;
const workspaceId = existing?.workspaceId ?? body.workspaceId ?? DEFAULT_WORKSPACE_ID;
if (!(await workspaceStore.get(workspaceId))) {
return Response.json({ error: "Workspace not found" }, { status: 404 });
}
releaseInteractive = autonomyActivity.enterInteractive(turnId);
evidenceContext = { sessionId, turnId, traceId, agentId: "orchistrator", workspaceId };
evidenceGoal = userText;
if (!startedSessions.has(sessionId)) {
await hookRuntime.dispatch(createHookEvent("SessionStart", {
sessionId,
model,
detail: { source: existing ? "resume" : "create", workspaceId },
}));
startedSessions.add(sessionId);
}
let activityWriter: UIMessageStreamWriter<AgentActivityMessage> | undefined;
const pendingActivity: Array<Parameters<UIMessageStreamWriter<AgentActivityMessage>["write"]>[0]> = [];
const lifecycle = { sessionId, turnId, model };
const activeEvidenceContext = evidenceContext!;
let streamFailed = false;
const prepared = await executionRuntime.prepare({
...lifecycle,
traceId,
workspaceId,
executionSource: "chat",
prompt: userText,
onIterationComplete: async (iteration) => {
if (iteration.agentId === "orchistrator") return;
const chunk = {
type: "data-agentActivity" as const,
id: iteration.runId,
data: agentActivity(turnId, iteration),
};
if (activityWriter) activityWriter.write(chunk);
else pendingActivity.push(chunk);
},
});
if (prepared.promptWasReplaced) replaceLatestUserText(messages, prepared.prompt);
const agentStream = await handleChatStream({
mastra,
agentId: "orchistrator",
version: "v6",
params: {
id: sessionId,
messages: [messages.at(-1)!],
memory: { thread: sessionId, resource: RESOURCE_ID },
} as Parameters<typeof handleChatStream>[0]["params"],
defaultOptions: prepared.options as Parameters<typeof handleChatStream>[0]["defaultOptions"],
onError: (error) => {
streamFailed = true;
evidenceFailureSignaled = true;
void evolution.enqueueFailedTurn({
...activeEvidenceContext,
summary: buildEvolutionEvidence({
goal: userText,
failure: messageOf(error),
classification: error instanceof Error ? error.name : "unknown",
}),
}).catch((signalError) => appLogger().warn("evolution.signal.failed", { error: messageOf(signalError) }));
void hookRuntime.dispatch(createHookEvent("StopFailure", {
...lifecycle,
detail: { error: { message: messageOf(error) } },
})).catch(() => undefined);
return "The agent could not complete this turn.";
},
messageMetadata: ({ part }) => {
if (part.type !== "finish" || !("totalUsage" in part)) return undefined;
return { custom: { usage: part.totalUsage, turnId, traceId } };
},
});
const stream = createUIMessageStream<AgentActivityMessage>({
execute: ({ writer }) => {
activityWriter = writer;
for (const chunk of pendingActivity) writer.write(chunk);
pendingActivity.length = 0;
writer.merge(agentStream as unknown as Parameters<typeof writer.merge>[0]);
},
onEnd: async ({ responseMessage, isAborted }) => {
try {
if (streamFailed || isAborted) return;
const outcome = responseMessage.parts
.filter((part): part is Extract<typeof part, { type: "text" }> => part.type === "text")
.map((part) => part.text)
.join("\n");
await evolution.enqueueCompletedTurn({
...activeEvidenceContext,
summary: buildEvolutionEvidence({
goal: userText,
outcome: outcome || "Completed without assistant text.",
classification: "success",
}),
}).catch((signalError) => appLogger().warn("evolution.signal.failed", { error: messageOf(signalError) }));
} finally {
releaseInteractive?.();
releaseInteractive = undefined;
}
},
});
const response = createUIMessageStreamResponse({
stream: stream as unknown as Parameters<typeof createUIMessageStreamResponse>[0]["stream"],
});
streamStarted = true;
return response;
} catch (error) {
releaseInteractive?.();
releaseInteractive = undefined;
if (evidenceContext && evidenceGoal && !streamStarted && !evidenceFailureSignaled) {
await evolution.enqueueFailedTurn({
...evidenceContext,
summary: buildEvolutionEvidence({
goal: evidenceGoal,
failure: messageOf(error),
classification: error instanceof Error ? error.name : "unknown",
}),
}).catch((signalError) => appLogger().warn("evolution.signal.failed", { error: messageOf(signalError) }));
}
if (error instanceof HookBlockedError) return Response.json({ error: error.reason }, { status: 403 });
appLogger().error("chat.turn.failed", { error });
return Response.json({ error: "Unable to start the agent turn" }, { status: 500 });
}
},
},
"/api/*": () => Response.json({ error: "Not found" }, { status: 404 }),
"/*": index,
};
return Bun.serve({
port,
idleTimeout: SERVER_IDLE_TIMEOUT_SECONDS,
routes: guardApiRoutes(routes) as Bun.Serve.Routes<undefined, string>,
});
}
const messageOf = (error: unknown) => error instanceof Error ? error.message : String(error);
if (import.meta.main) {
const server = startServer(Number(process.env.PORT ?? 3000));
await selfUpdateScheduler.start();
await evolutionRuntime.start();
await taskRuntime.start();
const maintenance = startMaintenance();
appLogger().info("server.started", { url: String(server.url) });
// Warm the semantic code index for the default repository in the background
// so researchers can query it from their first step instead of paying a cold
// full-index build (minutes on the local embedder) mid-task. Other
// workspaces warm lazily on their first semantic lookup: indexing every
// registered repository at boot kept the CPU embedder saturated for hours.
void (async () => {
const { codeIntelligence } = await import("./code-intelligence");
try {
const { path } = await agentWorkspaces.resolveRepository();
const startedAt = Date.now();
await codeIntelligence.warm(path);
appLogger().info("code-intelligence.warm.ready", { workspaceId: "default", elapsedMs: Date.now() - startedAt });
} catch (error) {
appLogger().warn("code-intelligence.warm.failed", { workspaceId: "default", error: messageOf(error) });
}
})();
let stopping = false;
const shutdown = async (signal: string) => {
if (stopping) return;
stopping = true;
appLogger().info("server.stopping", { signal });
maintenance.stop();
server.stop(false);
await selfUpdateScheduler.stop();
await evolutionRuntime.stop();
await taskRuntime.stop();
await mastra.observability.flush();
await agentWorkspaceRuntime.shutdown();
await browserRuntime.shutdown();
await mastra.shutdown();
process.exit(0);
};
process.once("SIGINT", () => void shutdown("SIGINT"));
process.once("SIGTERM", () => void shutdown("SIGTERM"));
}