feat(status-ui): improve concurrent session workflows
Add in-place session renewal, searchable model selection, clearer progress feedback, tool-result previews, and per-session presentation state. Keep commands, errors, loading, and lifecycle indicators attached to their owning runtime so parallel tabs cannot interfere.
This commit is contained in:
@@ -51,8 +51,7 @@ async function readSessionReference(sessionDir) {
|
||||
);
|
||||
return {
|
||||
sessionPath:
|
||||
typeof value?.sessionPath === "string" &&
|
||||
path.isAbsolute(value.sessionPath)
|
||||
typeof value?.sessionPath === "string" && path.isAbsolute(value.sessionPath)
|
||||
? value.sessionPath
|
||||
: undefined,
|
||||
worktreePath:
|
||||
@@ -110,11 +109,7 @@ async function readSessionSummary(sessionPath, currentPath) {
|
||||
firstMessage = sessionPreview(entry.message.content);
|
||||
}
|
||||
}
|
||||
if (
|
||||
!header ||
|
||||
typeof header.id !== "string" ||
|
||||
typeof header.cwd !== "string"
|
||||
)
|
||||
if (!header || typeof header.id !== "string" || typeof header.cwd !== "string")
|
||||
return undefined;
|
||||
return {
|
||||
path: sessionPath,
|
||||
@@ -380,8 +375,7 @@ export function createAgentRegistry({
|
||||
}
|
||||
|
||||
function enqueue(runtime, operation) {
|
||||
if (stopping)
|
||||
return Promise.reject(new Error("agent registry is stopping"));
|
||||
if (stopping) return Promise.reject(new Error("agent registry is stopping"));
|
||||
if (
|
||||
runtime.state === "closing" ||
|
||||
runtime.state === "stopped" ||
|
||||
@@ -420,10 +414,7 @@ export function createAgentRegistry({
|
||||
} catch (error) {
|
||||
if (error?.code !== "ENOENT" || !allowMissing) throw error;
|
||||
const resolvedParent = await realpath(path.dirname(sessionPath));
|
||||
resolvedSessionPath = path.join(
|
||||
resolvedParent,
|
||||
path.basename(sessionPath),
|
||||
);
|
||||
resolvedSessionPath = path.join(resolvedParent, path.basename(sessionPath));
|
||||
}
|
||||
const relative = path.relative(resolvedSessionDir, resolvedSessionPath);
|
||||
if (
|
||||
@@ -461,9 +452,7 @@ export function createAgentRegistry({
|
||||
resolved = candidate;
|
||||
} catch (error) {
|
||||
if (error?.code !== "ENOENT")
|
||||
throw new Error(
|
||||
"Pi reported a session outside its managed directory",
|
||||
);
|
||||
throw new Error("Pi reported a session outside its managed directory");
|
||||
}
|
||||
}
|
||||
if (expectedSessionPath && resolved !== expectedSessionPath)
|
||||
@@ -483,8 +472,7 @@ export function createAgentRegistry({
|
||||
throw new Error("session runtime is closing");
|
||||
if (typeof state.sessionName === "string")
|
||||
runtime.sessionName = state.sessionName;
|
||||
if (typeof state.sessionId === "string")
|
||||
runtime.sessionId = state.sessionId;
|
||||
if (typeof state.sessionId === "string") runtime.sessionId = state.sessionId;
|
||||
if (resolved) {
|
||||
if (
|
||||
runtime.sessionPath &&
|
||||
@@ -532,12 +520,8 @@ export function createAgentRegistry({
|
||||
}
|
||||
}
|
||||
if (type === "queue") {
|
||||
const followUp = Array.isArray(event.followUp)
|
||||
? event.followUp.length
|
||||
: 0;
|
||||
const steering = Array.isArray(event.steering)
|
||||
? event.steering.length
|
||||
: 0;
|
||||
const followUp = Array.isArray(event.followUp) ? event.followUp.length : 0;
|
||||
const steering = Array.isArray(event.steering) ? event.steering.length : 0;
|
||||
runtime.queueCount =
|
||||
typeof event.pendingMessageCount === "number"
|
||||
? event.pendingMessageCount
|
||||
@@ -553,8 +537,7 @@ export function createAgentRegistry({
|
||||
const method = event.method;
|
||||
if (["select", "confirm", "input", "editor"].includes(method)) {
|
||||
runtime.attention = true;
|
||||
if (typeof event.id === "string")
|
||||
runtime.extensions.set(event.id, event);
|
||||
if (typeof event.id === "string") runtime.extensions.set(event.id, event);
|
||||
}
|
||||
}
|
||||
const message = event.message;
|
||||
@@ -624,8 +607,7 @@ export function createAgentRegistry({
|
||||
const state = await adapter.send({ type: "get_state" });
|
||||
await updateRuntimeIdentity(runtime, state, { adapter });
|
||||
}).catch((error) => {
|
||||
if (runtime.stopped || runtime.state === "closing" || stopping)
|
||||
return;
|
||||
if (runtime.stopped || runtime.state === "closing" || stopping) return;
|
||||
runtime.error = {
|
||||
code: "identity_refresh",
|
||||
message: error.message,
|
||||
@@ -829,8 +811,7 @@ export function createAgentRegistry({
|
||||
}
|
||||
|
||||
function openRuntime(record, options) {
|
||||
if (stopping)
|
||||
return Promise.reject(new Error("agent registry is stopping"));
|
||||
if (stopping) return Promise.reject(new Error("agent registry is stopping"));
|
||||
const operation = openRuntimeInternal(record, options);
|
||||
pendingOpenOperations.add(operation);
|
||||
void operation
|
||||
@@ -867,14 +848,8 @@ export function createAgentRegistry({
|
||||
let canonicalSession;
|
||||
if (record.sessionPath) {
|
||||
const canonicalWorktree = await realpath(record.worktreePath);
|
||||
const sessionDir = sessionDirectoryFor(
|
||||
sessionRoot,
|
||||
canonicalWorktree,
|
||||
);
|
||||
canonicalSession = await ownedSessionPath(
|
||||
sessionDir,
|
||||
record.sessionPath,
|
||||
);
|
||||
const sessionDir = sessionDirectoryFor(sessionRoot, canonicalWorktree);
|
||||
canonicalSession = await ownedSessionPath(sessionDir, record.sessionPath);
|
||||
}
|
||||
prepared.push({ record, canonicalSession });
|
||||
} catch (error) {
|
||||
@@ -1057,6 +1032,13 @@ export function createAgentRegistry({
|
||||
runtime.sessionId = undefined;
|
||||
runtime.sessionName = undefined;
|
||||
runtime.firstMessage = undefined;
|
||||
runtime.extensions.clear();
|
||||
runtime.queueCount = 0;
|
||||
runtime.activeTool = undefined;
|
||||
runtime.error = undefined;
|
||||
runtime.errorAttention = false;
|
||||
runtime.attention = false;
|
||||
runtime.lastActivity = new Date().toISOString();
|
||||
if (runtime.legacyDefault)
|
||||
await persistSessionReference(
|
||||
runtime.sessionDir,
|
||||
@@ -1066,6 +1048,12 @@ export function createAgentRegistry({
|
||||
await persistWorkspace();
|
||||
const state = await runtime.adapter.send({ type: "get_state" });
|
||||
await updateRuntimeIdentity(runtime, state, { adapter: runtime.adapter });
|
||||
runtime.state = state?.data?.isStreaming ? "streaming" : "idle";
|
||||
runtime.stateVersion += 1;
|
||||
publishAgent(runtime, "agent_state", { state: runtime.state });
|
||||
publishWorkspace("runtime_renewed", runtime, {
|
||||
runtime: publicRuntime(runtime),
|
||||
});
|
||||
return response;
|
||||
}
|
||||
|
||||
@@ -1105,17 +1093,14 @@ export function createAgentRegistry({
|
||||
workingCount: runtimes.filter((entry) => entry.state === "streaming")
|
||||
.length,
|
||||
attentionCount: runtimes.filter((entry) => entry.attention).length,
|
||||
recoveringCount: runtimes.filter(
|
||||
(entry) => entry.state === "recovering",
|
||||
).length,
|
||||
recoveringCount: runtimes.filter((entry) => entry.state === "recovering")
|
||||
.length,
|
||||
errorCount: runtimes.filter((entry) =>
|
||||
["error", "failed"].includes(entry.state),
|
||||
).length,
|
||||
runtimes,
|
||||
}))
|
||||
.sort((left, right) =>
|
||||
left.worktreePath.localeCompare(right.worktreePath),
|
||||
),
|
||||
.sort((left, right) => left.worktreePath.localeCompare(right.worktreePath)),
|
||||
...(manifestIssue ? { issue: manifestIssue } : {}),
|
||||
};
|
||||
}
|
||||
@@ -1130,10 +1115,7 @@ export function createAgentRegistry({
|
||||
const loaded = await workspaceStore.load();
|
||||
manifestIssue = loaded.issue;
|
||||
if (loaded.migrating) {
|
||||
const sessionDir = sessionDirectoryFor(
|
||||
sessionRoot,
|
||||
canonicalHomeWorktree,
|
||||
);
|
||||
const sessionDir = sessionDirectoryFor(sessionRoot, canonicalHomeWorktree);
|
||||
await mkdir(sessionDir, {
|
||||
recursive: true,
|
||||
mode: SESSION_DIRECTORY_MODE,
|
||||
@@ -1143,9 +1125,7 @@ export function createAgentRegistry({
|
||||
runtimeId: randomUUID(),
|
||||
worktreePath: canonicalHomeWorktree,
|
||||
legacyDefault: true,
|
||||
...(reference.sessionPath
|
||||
? { sessionPath: reference.sessionPath }
|
||||
: {}),
|
||||
...(reference.sessionPath ? { sessionPath: reference.sessionPath } : {}),
|
||||
};
|
||||
if (reference.sessionPath) await restoreRecords([migrationRecord]);
|
||||
else await openRuntime(migrationRecord);
|
||||
@@ -1179,8 +1159,7 @@ export function createAgentRegistry({
|
||||
...(runtimesByWorktreePath.get(canonicalPath) ?? []),
|
||||
].filter((runtime) => runtime.agentId);
|
||||
const existing =
|
||||
liveRuntimes.find((runtime) => runtime.legacyDefault) ??
|
||||
liveRuntimes[0];
|
||||
liveRuntimes.find((runtime) => runtime.legacyDefault) ?? liveRuntimes[0];
|
||||
if (existing) {
|
||||
if (!existing.legacyDefault) {
|
||||
await setLegacyDefault(existing);
|
||||
@@ -1197,9 +1176,7 @@ export function createAgentRegistry({
|
||||
const opened = await openRuntime({
|
||||
worktreePath: canonicalPath,
|
||||
legacyDefault: true,
|
||||
...(reference.sessionPath
|
||||
? { sessionPath: reference.sessionPath }
|
||||
: {}),
|
||||
...(reference.sessionPath ? { sessionPath: reference.sessionPath } : {}),
|
||||
});
|
||||
return publicAgent(getRuntime(opened.runtimeId));
|
||||
},
|
||||
@@ -1213,9 +1190,7 @@ export function createAgentRegistry({
|
||||
listAgents() {
|
||||
return [...runtimesByAgentId.values()]
|
||||
.map(publicAgent)
|
||||
.sort((left, right) =>
|
||||
left.worktreePath.localeCompare(right.worktreePath),
|
||||
);
|
||||
.sort((left, right) => left.worktreePath.localeCompare(right.worktreePath));
|
||||
},
|
||||
async listDirectories() {
|
||||
let entries = [];
|
||||
@@ -1240,9 +1215,9 @@ export function createAgentRegistry({
|
||||
]);
|
||||
return [...paths]
|
||||
.map((worktreePath) => {
|
||||
const runtime = [
|
||||
...(runtimesByWorktreePath.get(worktreePath) ?? []),
|
||||
].find((entry) => entry.agentId);
|
||||
const runtime = [...(runtimesByWorktreePath.get(worktreePath) ?? [])].find(
|
||||
(entry) => entry.agentId,
|
||||
);
|
||||
return {
|
||||
worktreePath,
|
||||
state: runtime?.state ?? "inactive",
|
||||
@@ -1250,9 +1225,7 @@ export function createAgentRegistry({
|
||||
...(runtime ? { agentId: runtime.agentId } : {}),
|
||||
};
|
||||
})
|
||||
.sort((left, right) =>
|
||||
left.worktreePath.localeCompare(right.worktreePath),
|
||||
);
|
||||
.sort((left, right) => left.worktreePath.localeCompare(right.worktreePath));
|
||||
},
|
||||
async forgetDirectory(worktreePath) {
|
||||
let canonicalPath = path.normalize(worktreePath);
|
||||
@@ -1297,9 +1270,8 @@ export function createAgentRegistry({
|
||||
bridgeInstanceId,
|
||||
latestSeq: workspaceSequence,
|
||||
openCount: runtimes.length,
|
||||
workingCount: runtimes.filter(
|
||||
(runtime) => runtime.state === "streaming",
|
||||
).length,
|
||||
workingCount: runtimes.filter((runtime) => runtime.state === "streaming")
|
||||
.length,
|
||||
attentionCount: runtimes.filter((runtime) => runtime.attention).length,
|
||||
recoveringCount: runtimes.filter(
|
||||
(runtime) => runtime.state === "recovering",
|
||||
@@ -1323,11 +1295,11 @@ export function createAgentRegistry({
|
||||
};
|
||||
return enqueue(runtime, async () => {
|
||||
const adapter = runtime.adapter;
|
||||
const safe = async (command, fallback) => {
|
||||
const safe = async (command, fallback, { includeError = false } = {}) => {
|
||||
try {
|
||||
return await adapter.send({ type: command });
|
||||
} catch {
|
||||
return fallback;
|
||||
} catch (error) {
|
||||
return includeError ? { ...fallback, error: String(error) } : fallback;
|
||||
}
|
||||
};
|
||||
const state = await safe("get_state", { data: {} });
|
||||
@@ -1336,9 +1308,11 @@ export function createAgentRegistry({
|
||||
data: { messages: [] },
|
||||
});
|
||||
const commands = await safe("get_commands", { data: { commands: [] } });
|
||||
const models = await safe("get_available_models", {
|
||||
data: { models: [] },
|
||||
});
|
||||
const models = await safe(
|
||||
"get_available_models",
|
||||
{ data: { models: [] } },
|
||||
{ includeError: true },
|
||||
);
|
||||
await updateRuntimeIdentity(runtime, state, { adapter });
|
||||
return {
|
||||
bridgeInstanceId,
|
||||
@@ -1392,10 +1366,23 @@ export function createAgentRegistry({
|
||||
return publicAgent(runtime);
|
||||
});
|
||||
},
|
||||
async renewSession(agentId) {
|
||||
const runtime = getAgent(agentId);
|
||||
return enqueue(runtime, async () => {
|
||||
if (!runtime.adapter) throw new Error("Pi is not currently running");
|
||||
if (runtime.state === "streaming" || runtime.queueCount > 0)
|
||||
await runtime.adapter.send({ type: "abort" });
|
||||
const response = await beginLegacyNewSession(runtime);
|
||||
if (response?.data?.cancelled)
|
||||
throw new Error("Pi cancelled the session renewal");
|
||||
return publicAgent(runtime);
|
||||
});
|
||||
},
|
||||
async restart(agentId) {
|
||||
const runtime = getAgent(agentId);
|
||||
return enqueue(runtime, async () => {
|
||||
const replacedAdapter = runtime.adapter;
|
||||
if (!replacedAdapter) throw new Error("Pi is not currently running");
|
||||
// Detach the old child before stopping it so late events cannot mutate
|
||||
// presentation state that belongs to the replacement generation.
|
||||
runtime.adapter = undefined;
|
||||
@@ -1406,7 +1393,29 @@ export function createAgentRegistry({
|
||||
runtime.attention = runtime.errorAttention;
|
||||
runtime.supervisor.markHealthy();
|
||||
runtime.recoveryPromise = runtime.supervisor.handleUnexpectedExit();
|
||||
await runtime.recoveryPromise;
|
||||
const nextAdapter = await runtime.recoveryPromise;
|
||||
if (!nextAdapter) {
|
||||
runtime.state = "failed";
|
||||
runtime.stateVersion += 1;
|
||||
runtime.error = {
|
||||
code: "restart_failed",
|
||||
message: "Pi did not restart after the recovery attempts were exhausted",
|
||||
};
|
||||
runtime.errorAttention = true;
|
||||
runtime.attention = true;
|
||||
publishAgent(runtime, "agent_state", {
|
||||
state: runtime.state,
|
||||
error: runtime.error,
|
||||
});
|
||||
throw new Error(runtime.error.message);
|
||||
}
|
||||
runtime.error = undefined;
|
||||
runtime.errorAttention = false;
|
||||
runtime.attention = runtime.extensions.size > 0;
|
||||
publishAgent(runtime, "agent_state", { state: runtime.state });
|
||||
publishWorkspace("runtime_recovered", runtime, {
|
||||
runtime: publicRuntime(runtime),
|
||||
});
|
||||
return publicAgent(runtime);
|
||||
});
|
||||
},
|
||||
@@ -1458,9 +1467,7 @@ export function createAgentRegistry({
|
||||
operation === "extension_response" &&
|
||||
!runtime.extensions.has(payload.requestId)
|
||||
)
|
||||
throw new Error(
|
||||
"extension request is not actionable for this runtime",
|
||||
);
|
||||
throw new Error("extension request is not actionable for this runtime");
|
||||
const response = await routeCommand(
|
||||
runtime.adapter,
|
||||
operation,
|
||||
@@ -1473,8 +1480,7 @@ export function createAgentRegistry({
|
||||
}
|
||||
if (operation === "extension_response") {
|
||||
runtime.extensions.delete(payload.requestId);
|
||||
runtime.attention =
|
||||
runtime.errorAttention || runtime.extensions.size > 0;
|
||||
runtime.attention = runtime.errorAttention || runtime.extensions.size > 0;
|
||||
}
|
||||
return delivery ? { ...response, delivery } : response;
|
||||
} catch (error) {
|
||||
|
||||
@@ -76,6 +76,8 @@ export async function startBridgeService({
|
||||
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:
|
||||
|
||||
Reference in New Issue
Block a user