158 lines
4.5 KiB
JavaScript
158 lines
4.5 KiB
JavaScript
import path from "node:path";
|
|
import {
|
|
createAgentRegistry,
|
|
UnknownAgentError,
|
|
UnknownRuntimeError,
|
|
} from "./agent-registry.js";
|
|
import { startPiRpcAdapter } from "./pi-rpc-adapter.js";
|
|
import { ProtocolError } from "../protocol/index.js";
|
|
import { createRuntimePaths, startBridgeServer } from "./runtime.js";
|
|
|
|
export async function startBridgeService({
|
|
homeWorktree,
|
|
runtimeDir,
|
|
sessionRoot,
|
|
startAdapter = startPiRpcAdapter,
|
|
processAlive,
|
|
maxFrameBytes,
|
|
}) {
|
|
const paths = await createRuntimePaths(runtimeDir);
|
|
const registry = createAgentRegistry({
|
|
homeWorktree,
|
|
sessionRoot: sessionRoot ?? path.join(paths.directory, "sessions"),
|
|
startAdapter,
|
|
});
|
|
|
|
const dispatch = async (request) => {
|
|
try {
|
|
switch (request.op) {
|
|
case "list_agents":
|
|
return { agents: registry.listAgents() };
|
|
case "list_directories":
|
|
return { directories: await registry.listDirectories() };
|
|
case "select_agent":
|
|
return {
|
|
agent: await registry.selectWorktree(request.payload.worktreePath),
|
|
};
|
|
case "forget_directory":
|
|
return registry.forgetDirectory(request.payload.worktreePath);
|
|
case "get_workspace":
|
|
return registry.getWorkspace();
|
|
case "get_workspace_summary":
|
|
return registry.getWorkspaceSummary();
|
|
case "get_model_catalog":
|
|
return { models: await registry.getModelCatalog() };
|
|
case "create_session_runtime":
|
|
return {
|
|
runtime: await registry.createSessionRuntime(
|
|
request.payload.worktreePath,
|
|
),
|
|
};
|
|
case "create_quick_runtime":
|
|
return {
|
|
runtime: await registry.createQuickRuntime(
|
|
request.payload.worktreePath,
|
|
),
|
|
};
|
|
case "promote_quick_runtime":
|
|
return {
|
|
runtime: await registry.promoteQuickRuntime(request.payload.runtimeId),
|
|
};
|
|
case "open_session_runtime":
|
|
return {
|
|
runtime: await registry.openSessionRuntime(
|
|
request.payload.worktreePath,
|
|
request.payload.sessionPath,
|
|
),
|
|
};
|
|
case "close_session_runtime":
|
|
return registry.closeSessionRuntime(request.payload.runtimeId);
|
|
case "close_quick_runtime":
|
|
return registry.closeQuickRuntime(request.payload.runtimeId);
|
|
case "list_directory_sessions":
|
|
return {
|
|
sessions: await registry.listDirectorySessions(
|
|
request.payload.worktreePath,
|
|
),
|
|
};
|
|
case "get_session_runtime_snapshot":
|
|
return registry.getSessionRuntimeSnapshot(request.payload.runtimeId);
|
|
case "subscribe_workspace":
|
|
return registry.workspaceEventsAfter(request.payload.cursor ?? 0);
|
|
case "subscribe":
|
|
return {
|
|
events: registry.eventsAfter(
|
|
request.agentId,
|
|
request.payload.cursor ?? 0,
|
|
),
|
|
};
|
|
case "retry":
|
|
return { agent: await registry.retry(request.agentId) };
|
|
case "restart":
|
|
return { agent: await registry.restart(request.agentId) };
|
|
case "renew_session":
|
|
return { agent: await registry.renewSession(request.agentId) };
|
|
case "list_sessions":
|
|
return { sessions: await registry.listSessions(request.agentId) };
|
|
default:
|
|
return registry.route(request.agentId, request.op, request.payload);
|
|
}
|
|
} catch (error) {
|
|
if (error instanceof UnknownAgentError)
|
|
throw new ProtocolError("unknown_agent", error.message);
|
|
if (error instanceof UnknownRuntimeError)
|
|
throw new ProtocolError("unknown_runtime", error.message);
|
|
throw error;
|
|
}
|
|
};
|
|
|
|
const server = await startBridgeServer({
|
|
runtimeDir,
|
|
handleRequest: dispatch,
|
|
subscribe: (request, notify) => {
|
|
if (request.op === "subscribe_workspace") {
|
|
const subscription = registry.subscribeWorkspace(
|
|
request.payload.cursor ?? 0,
|
|
notify,
|
|
);
|
|
return {
|
|
result: {
|
|
bridgeInstanceId: subscription.bridgeInstanceId,
|
|
firstAvailableSeq: subscription.firstAvailableSeq,
|
|
latestSeq: subscription.latestSeq,
|
|
truncated: subscription.truncated,
|
|
events: subscription.events,
|
|
},
|
|
unsubscribe: subscription.unsubscribe,
|
|
};
|
|
}
|
|
const subscription = registry.subscribe(
|
|
request.agentId,
|
|
request.payload.cursor ?? 0,
|
|
notify,
|
|
);
|
|
return {
|
|
result: { events: subscription.events },
|
|
unsubscribe: subscription.unsubscribe,
|
|
};
|
|
},
|
|
...(processAlive ? { processAlive } : {}),
|
|
...(maxFrameBytes ? { maxFrameBytes } : {}),
|
|
});
|
|
try {
|
|
await registry.start();
|
|
} catch (error) {
|
|
await server.close();
|
|
throw error;
|
|
}
|
|
|
|
return {
|
|
socketPath: server.socketPath,
|
|
listAgents: () => registry.listAgents(),
|
|
async close() {
|
|
await registry.stop();
|
|
await server.close();
|
|
},
|
|
};
|
|
}
|