feat(status-ui): improve bridge recovery and desktop controls
This commit is contained in:
+158
-14
@@ -5,6 +5,7 @@ import {
|
||||
readdir,
|
||||
readFile,
|
||||
realpath,
|
||||
rm,
|
||||
stat,
|
||||
writeFile,
|
||||
} from "node:fs/promises";
|
||||
@@ -141,7 +142,7 @@ async function listSessionSummaries(sessionDir, currentPath) {
|
||||
.sort((left, right) => right.modified.localeCompare(left.modified));
|
||||
}
|
||||
|
||||
async function listDirectoryCatalog(sessionRoot, agentsByPath) {
|
||||
async function listDirectoryCatalog(sessionRoot, agentsByPath, homeWorktree) {
|
||||
let entries;
|
||||
try {
|
||||
entries = await readdir(sessionRoot, { withFileTypes: true });
|
||||
@@ -174,6 +175,7 @@ async function listDirectoryCatalog(sessionRoot, agentsByPath) {
|
||||
return {
|
||||
worktreePath,
|
||||
state: agent?.state ?? "inactive",
|
||||
isHome: worktreePath === homeWorktree,
|
||||
...(agent ? { agentId: agent.id } : {}),
|
||||
};
|
||||
})
|
||||
@@ -249,11 +251,19 @@ export function createAgentRegistry({
|
||||
const agentsByPath = new Map();
|
||||
const agentsById = new Map();
|
||||
const inFlight = new Map();
|
||||
const forgettingPaths = new Map();
|
||||
const pendingForgetOperations = new Set();
|
||||
let canonicalHomeWorktree;
|
||||
let stopping = false;
|
||||
|
||||
async function ensureAgent(worktreePath) {
|
||||
if (stopping) throw new Error("agent registry is stopping");
|
||||
if (typeof worktreePath !== "string" || !path.isAbsolute(worktreePath))
|
||||
throw new TypeError("worktreePath must be an absolute path");
|
||||
const canonicalPath = await realpath(worktreePath);
|
||||
if (stopping) throw new Error("agent registry is stopping");
|
||||
if (forgettingPaths.has(canonicalPath))
|
||||
throw new Error("directory is being forgotten");
|
||||
const existing = agentsByPath.get(canonicalPath);
|
||||
if (existing) return publicAgent(existing);
|
||||
const creating = inFlight.get(canonicalPath);
|
||||
@@ -276,10 +286,14 @@ export function createAgentRegistry({
|
||||
events: [],
|
||||
listeners: new Set(),
|
||||
nextSequence: 1,
|
||||
stateVersion: 0,
|
||||
adapter: undefined,
|
||||
stopped: false,
|
||||
recoveryPromise: undefined,
|
||||
supervisor: undefined,
|
||||
forgetting: false,
|
||||
activeOperations: 0,
|
||||
operationWaiters: new Set(),
|
||||
};
|
||||
const publish = (type, data) => {
|
||||
const event = {
|
||||
@@ -304,12 +318,14 @@ export function createAgentRegistry({
|
||||
typeof event.data?.state === "string"
|
||||
) {
|
||||
agent.state = event.data.state;
|
||||
agent.stateVersion += 1;
|
||||
if (event.data.state === "idle") agent.supervisor.markHealthy();
|
||||
}
|
||||
publish(event.type, event.data);
|
||||
},
|
||||
onError: (error) => {
|
||||
agent.state = "error";
|
||||
agent.stateVersion += 1;
|
||||
publish("agent_state", {
|
||||
state: "error",
|
||||
error: { code: error.code ?? "unknown", message: error.message },
|
||||
@@ -323,6 +339,7 @@ export function createAgentRegistry({
|
||||
.then((adapterAfterRecovery) => {
|
||||
if (agent.stopped) return undefined;
|
||||
agent.state = adapterAfterRecovery ? "idle" : "failed";
|
||||
agent.stateVersion += 1;
|
||||
publish("agent_state", { state: agent.state });
|
||||
return adapterAfterRecovery;
|
||||
});
|
||||
@@ -332,6 +349,7 @@ export function createAgentRegistry({
|
||||
agent.adapter = adapter;
|
||||
const state = await adapter.send({ type: "get_state" });
|
||||
agent.state = state?.data?.isStreaming ? "streaming" : "idle";
|
||||
agent.stateVersion += 1;
|
||||
const sessionPath = state?.data?.sessionFile;
|
||||
if (typeof sessionPath === "string" && path.isAbsolute(sessionPath)) {
|
||||
agent.sessionPath = sessionPath;
|
||||
@@ -367,6 +385,25 @@ export function createAgentRegistry({
|
||||
return agent;
|
||||
}
|
||||
|
||||
async function withAgentOperation(agent, operation) {
|
||||
if (agent.forgetting) throw new Error("directory is being forgotten");
|
||||
agent.activeOperations += 1;
|
||||
try {
|
||||
return await operation();
|
||||
} finally {
|
||||
agent.activeOperations -= 1;
|
||||
if (agent.activeOperations === 0) {
|
||||
for (const resolve of agent.operationWaiters) resolve();
|
||||
agent.operationWaiters.clear();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function waitForAgentOperations(agent) {
|
||||
if (agent.activeOperations === 0) return Promise.resolve();
|
||||
return new Promise((resolve) => agent.operationWaiters.add(resolve));
|
||||
}
|
||||
|
||||
async function refreshAgentSessionReference(agent) {
|
||||
const state = await agent.adapter.send({ type: "get_state" });
|
||||
const sessionPath = state?.data?.sessionFile;
|
||||
@@ -417,9 +454,78 @@ export function createAgentRegistry({
|
||||
return response;
|
||||
}
|
||||
|
||||
async function forgetDirectory(worktreePath) {
|
||||
if (stopping) throw new Error("agent registry is stopping");
|
||||
if (typeof worktreePath !== "string" || !path.isAbsolute(worktreePath))
|
||||
throw new TypeError("worktreePath must be an absolute path");
|
||||
let finishPendingForget;
|
||||
const pendingForget = new Promise((resolve) => {
|
||||
finishPendingForget = resolve;
|
||||
});
|
||||
pendingForgetOperations.add(pendingForget);
|
||||
try {
|
||||
let canonicalPath = path.normalize(worktreePath);
|
||||
if (!agentsByPath.has(canonicalPath)) {
|
||||
try {
|
||||
canonicalPath = await realpath(worktreePath);
|
||||
} catch (error) {
|
||||
if (error?.code !== "ENOENT") throw error;
|
||||
}
|
||||
}
|
||||
const homePath = canonicalHomeWorktree ?? (await realpath(homeWorktree));
|
||||
if (canonicalPath === homePath)
|
||||
throw new Error("cannot forget the home directory");
|
||||
|
||||
const existingForget = forgettingPaths.get(canonicalPath);
|
||||
if (existingForget) return existingForget;
|
||||
const forgetting = (async () => {
|
||||
try {
|
||||
const creation = inFlight.get(canonicalPath);
|
||||
if (creation) await creation;
|
||||
const agent = agentsByPath.get(canonicalPath);
|
||||
if (agent) {
|
||||
agent.forgetting = true;
|
||||
await waitForAgentOperations(agent);
|
||||
if (agent.state === "streaming")
|
||||
throw new Error("cannot forget a directory while Pi is working");
|
||||
agent.stopped = true;
|
||||
agent.supervisor.stop();
|
||||
await agent.adapter.stop();
|
||||
agent.state = "stopped";
|
||||
agent.stateVersion += 1;
|
||||
agent.listeners.clear();
|
||||
agentsByPath.delete(canonicalPath);
|
||||
agentsById.delete(agent.id);
|
||||
}
|
||||
const sessionDir = sessionDirectoryFor(sessionRoot, canonicalPath);
|
||||
await rm(path.join(sessionDir, SESSION_REFERENCE_FILE), {
|
||||
force: true,
|
||||
});
|
||||
return { worktreePath: canonicalPath };
|
||||
} catch (error) {
|
||||
const agent = agentsByPath.get(canonicalPath);
|
||||
if (agent) agent.forgetting = false;
|
||||
throw error;
|
||||
}
|
||||
})();
|
||||
forgettingPaths.set(canonicalPath, forgetting);
|
||||
try {
|
||||
return await forgetting;
|
||||
} finally {
|
||||
if (forgettingPaths.get(canonicalPath) === forgetting)
|
||||
forgettingPaths.delete(canonicalPath);
|
||||
}
|
||||
} finally {
|
||||
pendingForgetOperations.delete(pendingForget);
|
||||
finishPendingForget();
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
async start() {
|
||||
return ensureAgent(homeWorktree);
|
||||
const agent = await ensureAgent(homeWorktree);
|
||||
canonicalHomeWorktree = agent.worktreePath;
|
||||
return agent;
|
||||
},
|
||||
async selectWorktree(worktreePath) {
|
||||
return ensureAgent(worktreePath);
|
||||
@@ -432,8 +538,13 @@ export function createAgentRegistry({
|
||||
);
|
||||
},
|
||||
async listDirectories() {
|
||||
return listDirectoryCatalog(sessionRoot, agentsByPath);
|
||||
return listDirectoryCatalog(
|
||||
sessionRoot,
|
||||
agentsByPath,
|
||||
canonicalHomeWorktree,
|
||||
);
|
||||
},
|
||||
forgetDirectory,
|
||||
async listSessions(agentId) {
|
||||
const agent = getAgent(agentId);
|
||||
return listSessionSummaries(agent.sessionDir, agent.sessionPath);
|
||||
@@ -459,30 +570,63 @@ export function createAgentRegistry({
|
||||
},
|
||||
async retry(agentId) {
|
||||
const agent = getAgent(agentId);
|
||||
agent.supervisor.markHealthy();
|
||||
agent.recoveryPromise = agent.supervisor.handleUnexpectedExit();
|
||||
await agent.recoveryPromise;
|
||||
return publicAgent(agent);
|
||||
return withAgentOperation(agent, async () => {
|
||||
agent.supervisor.markHealthy();
|
||||
agent.recoveryPromise = agent.supervisor.handleUnexpectedExit();
|
||||
await agent.recoveryPromise;
|
||||
return publicAgent(agent);
|
||||
});
|
||||
},
|
||||
async restart(agentId) {
|
||||
const agent = getAgent(agentId);
|
||||
await agent.adapter.stop();
|
||||
return this.retry(agentId);
|
||||
return withAgentOperation(agent, async () => {
|
||||
await agent.adapter.stop();
|
||||
agent.supervisor.markHealthy();
|
||||
agent.recoveryPromise = agent.supervisor.handleUnexpectedExit();
|
||||
await agent.recoveryPromise;
|
||||
return publicAgent(agent);
|
||||
});
|
||||
},
|
||||
async route(agentId, operation, payload) {
|
||||
const agent = getAgent(agentId);
|
||||
if (operation === "switch_session")
|
||||
return switchSession(agent, payload.sessionPath);
|
||||
if (operation === "new_session") return newSession(agent);
|
||||
return routeCommand(agent.adapter, operation, payload, agent.state);
|
||||
return withAgentOperation(agent, async () => {
|
||||
if (operation === "switch_session")
|
||||
return switchSession(agent, payload.sessionPath);
|
||||
if (operation === "new_session") return newSession(agent);
|
||||
const previousState = agent.state;
|
||||
const previousStateVersion = agent.stateVersion;
|
||||
const startsWork =
|
||||
operation === "prompt" ||
|
||||
(operation === "submit_prompt" && previousState !== "streaming");
|
||||
if (startsWork) agent.state = "streaming";
|
||||
try {
|
||||
return await routeCommand(
|
||||
agent.adapter,
|
||||
operation,
|
||||
payload,
|
||||
previousState,
|
||||
);
|
||||
} catch (error) {
|
||||
if (startsWork && agent.stateVersion === previousStateVersion)
|
||||
agent.state = previousState;
|
||||
throw error;
|
||||
}
|
||||
});
|
||||
},
|
||||
async stop() {
|
||||
stopping = true;
|
||||
await Promise.allSettled([...inFlight.values()]);
|
||||
await Promise.all([...pendingForgetOperations]);
|
||||
const agents = [...agentsById.values()];
|
||||
for (const agent of agents) agent.forgetting = true;
|
||||
await Promise.all(agents.map(waitForAgentOperations));
|
||||
await Promise.all(
|
||||
[...agentsById.values()].map(async (agent) => {
|
||||
agents.map(async (agent) => {
|
||||
agent.stopped = true;
|
||||
agent.supervisor.stop();
|
||||
await agent.adapter.stop();
|
||||
agent.state = "stopped";
|
||||
agent.stateVersion += 1;
|
||||
}),
|
||||
);
|
||||
},
|
||||
|
||||
Reference in New Issue
Block a user