feat: add multi-session desktop workspace

This commit is contained in:
2026-08-13 21:34:04 +02:00
parent 12b9b83d51
commit 475f95ba70
56 changed files with 8998 additions and 3805 deletions
+1192 -346
View File
File diff suppressed because it is too large Load Diff
+9 -1
View File
@@ -4,6 +4,7 @@ import path from "node:path";
import { StringDecoder } from "node:string_decoder";
export const DEFAULT_COMMAND_TIMEOUT_MS = 30_000;
export const COMPACT_COMMAND_TIMEOUT_MS = 5 * 60_000;
export const MAX_PI_RPC_FRAME_BYTES = 1024 * 1024;
const supportedCommands = new Set([
@@ -18,6 +19,8 @@ const supportedCommands = new Set([
"get_messages",
"get_available_models",
"get_commands",
"set_session_name",
"compact",
"set_model",
"set_thinking_level",
]);
@@ -126,6 +129,7 @@ export function startPiRpcAdapter({
onEvent = () => {},
onError = () => {},
commandTimeoutMs = DEFAULT_COMMAND_TIMEOUT_MS,
compactCommandTimeoutMs = COMPACT_COMMAND_TIMEOUT_MS,
maxFrameBytes = MAX_PI_RPC_FRAME_BYTES,
}) {
assertAbsolutePath(cwd, "cwd");
@@ -316,6 +320,10 @@ export function startPiRpcAdapter({
}
const id = `bridge-${randomUUID()}`;
const timeoutMs =
commandInput.type === "compact"
? compactCommandTimeoutMs
: commandTimeoutMs;
return new Promise((resolve, reject) => {
const timeout = setTimeout(() => {
pending.delete(id);
@@ -325,7 +333,7 @@ export function startPiRpcAdapter({
`Pi RPC command timed out: ${commandInput.type}`,
),
);
}, commandTimeoutMs);
}, timeoutMs);
timeout.unref?.();
pending.set(id, { resolve, reject, timeout });
try {
+52 -1
View File
@@ -1,5 +1,9 @@
import path from "node:path";
import { createAgentRegistry, UnknownAgentError } from "./agent-registry.js";
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";
@@ -32,6 +36,35 @@ export async function startBridgeService({
};
case "forget_directory":
return registry.forgetDirectory(request.payload.worktreePath);
case "get_workspace":
return registry.getWorkspace();
case "get_workspace_summary":
return registry.getWorkspaceSummary();
case "create_session_runtime":
return {
runtime: await registry.createSessionRuntime(
request.payload.worktreePath,
),
};
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 "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(
@@ -51,6 +84,8 @@ export async function startBridgeService({
} 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;
}
};
@@ -59,6 +94,22 @@ export async function startBridgeService({
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,
+4 -1
View File
@@ -88,7 +88,10 @@ export async function startUnixSocketServer({
const normalized = line.endsWith("\r") ? line.slice(0, -1) : line;
try {
const request = parseRequestFrame(normalized, maxFrameBytes);
if (request.op === "subscribe" && subscribe) {
if (
(request.op === "subscribe" || request.op === "subscribe_workspace") &&
subscribe
) {
await handleSubscription(request);
return;
}
+153
View File
@@ -0,0 +1,153 @@
import {
chmod,
mkdir,
readFile,
rename,
rm,
writeFile,
} from "node:fs/promises";
import path from "node:path";
import { randomUUID } from "node:crypto";
const WORKSPACE_VERSION = 2;
const WORKSPACE_FILE = "bridge-workspace-v2.json";
const DIRECTORY_MODE = 0o700;
const FILE_MODE = 0o600;
function validRuntime(record) {
return (
record &&
typeof record === "object" &&
typeof record.runtimeId === "string" &&
record.runtimeId.length > 0 &&
typeof record.worktreePath === "string" &&
path.isAbsolute(record.worktreePath) &&
(record.sessionPath === undefined ||
(typeof record.sessionPath === "string" &&
path.isAbsolute(record.sessionPath))) &&
(record.sessionId === undefined || typeof record.sessionId === "string") &&
(record.legacyDefault === undefined ||
typeof record.legacyDefault === "boolean")
);
}
function normalizeManifest(value) {
if (
!value ||
typeof value !== "object" ||
value.version !== WORKSPACE_VERSION ||
!Array.isArray(value.runtimes) ||
!value.runtimes.every(validRuntime) ||
(value.directories !== undefined &&
(!Array.isArray(value.directories) ||
!value.directories.every(
(directory) =>
typeof directory === "string" && path.isAbsolute(directory),
)))
)
throw new Error("workspace manifest has an unsupported or invalid shape");
const runtimeIds = new Set();
const sessionPaths = new Set();
const directories = new Set(value.directories ?? []);
if (directories.size !== (value.directories ?? []).length)
throw new Error("workspace manifest repeats a directory");
for (const runtime of value.runtimes) directories.add(runtime.worktreePath);
for (const runtime of value.runtimes) {
if (runtimeIds.has(runtime.runtimeId))
throw new Error(
`workspace manifest repeats runtime ${runtime.runtimeId}`,
);
runtimeIds.add(runtime.runtimeId);
if (runtime.sessionPath) {
if (sessionPaths.has(runtime.sessionPath))
throw new Error(
`workspace manifest repeats session ${runtime.sessionPath}`,
);
sessionPaths.add(runtime.sessionPath);
}
}
return {
version: WORKSPACE_VERSION,
migrated: value.migrated === true,
directories: [...directories].sort(),
runtimes: value.runtimes.map((runtime) => ({ ...runtime })),
};
}
export function createWorkspaceStore(sessionRoot) {
if (typeof sessionRoot !== "string" || !path.isAbsolute(sessionRoot))
throw new TypeError("sessionRoot must be an absolute path");
const filePath = path.join(sessionRoot, WORKSPACE_FILE);
let writeQueue = Promise.resolve();
async function load() {
await mkdir(sessionRoot, { recursive: true, mode: DIRECTORY_MODE });
await chmod(sessionRoot, DIRECTORY_MODE);
try {
return {
manifest: normalizeManifest(
JSON.parse(await readFile(filePath, "utf8")),
),
migrating: false,
};
} catch (error) {
if (error?.code === "ENOENT") {
return {
manifest: {
version: WORKSPACE_VERSION,
migrated: false,
directories: [],
runtimes: [],
},
migrating: true,
};
}
return {
manifest: {
version: WORKSPACE_VERSION,
migrated: true,
directories: [],
runtimes: [],
},
migrating: false,
issue: {
code: "invalid_workspace_manifest",
message: error instanceof Error ? error.message : String(error),
},
};
}
}
function save(manifest) {
const normalized = normalizeManifest({ ...manifest, migrated: true });
const operation = writeQueue
.catch(() => {})
.then(async () => {
await mkdir(sessionRoot, { recursive: true, mode: DIRECTORY_MODE });
await chmod(sessionRoot, DIRECTORY_MODE);
const temporaryPath = `${filePath}.${process.pid}.${randomUUID()}.tmp`;
try {
await writeFile(
temporaryPath,
`${JSON.stringify(normalized, null, 2)}\n`,
{
encoding: "utf8",
mode: FILE_MODE,
},
);
await chmod(temporaryPath, FILE_MODE);
await rename(temporaryPath, filePath);
await chmod(filePath, FILE_MODE);
} finally {
await rm(temporaryPath, { force: true }).catch(() => {});
}
});
writeQueue = operation.catch(() => {});
return operation;
}
return { filePath, load, save };
}
export const workspaceManifestVersion = WORKSPACE_VERSION;
export const workspaceManifestFile = WORKSPACE_FILE;
+8 -21
View File
@@ -7,29 +7,16 @@ export async function publishNoctaliaState(
state,
{ execute = async (command, args) => execFileAsync(command, args) } = {},
) {
if (!state || typeof state !== "object")
throw new TypeError("state is required");
const { state: bridgeState, projectLabel, attentionCount, detail } = state;
if (!state || typeof state !== "object") throw new TypeError("state is required");
const { state: bridgeState, projectLabel, attentionCount, detail, openCount = 0, workingCount = 0, recoveringCount = 0, errorCount = 0 } = state;
if (
typeof bridgeState !== "string" ||
typeof projectLabel !== "string" ||
!Number.isSafeInteger(attentionCount) ||
typeof bridgeState !== "string" || typeof projectLabel !== "string" ||
![attentionCount, openCount, workingCount, recoveringCount, errorCount].every(Number.isSafeInteger) ||
typeof detail !== "string"
) {
throw new TypeError(
"state must contain string state/projectLabel/detail and integer attentionCount",
);
}
) throw new TypeError("state must contain presentation strings and integer counts");
await execute("qs", [
"-c",
"noctalia-shell",
"ipc",
"call",
"plugin:pi-status-bridge",
"setBridgeState",
bridgeState,
projectLabel,
String(attentionCount),
detail,
"-c", "noctalia-shell", "ipc", "call", "plugin:pi-status-bridge", "setBridgeState",
bridgeState, projectLabel, String(attentionCount), detail, String(openCount),
String(workingCount), String(recoveringCount), String(errorCount),
]);
}
+14 -18
View File
@@ -10,25 +10,23 @@ function flag(args, name) {
async function run() {
const args = process.argv.slice(2);
const socketPath =
flag(args, "--socket") ?? process.env.PI_STATUS_BRIDGE_SOCKET;
const socketPath = flag(args, "--socket") ?? process.env.PI_STATUS_BRIDGE_SOCKET;
const agentId = flag(args, "--agent");
if (!socketPath || !agentId)
throw new Error(
"Usage: pi-status-bridge-noctalia-relay --socket <path> --agent <id>",
);
if (!socketPath)
throw new Error("Usage: pi-status-bridge-noctalia-relay --socket <path> [--agent <id>]");
const client = await connectLocalClient({ socketPath });
const agents = await client.request("list_agents");
const agent = agents.agents.find((candidate) => candidate.id === agentId);
if (!agent) throw new Error(`unknown agent: ${agentId}`);
let agent;
if (agentId) {
const agents = await client.request("list_agents");
agent = agents.agents.find((candidate) => candidate.id === agentId);
if (!agent) throw new Error(`unknown agent: ${agentId}`);
}
const relay = createNoctaliaStateRelay({
client,
agent,
onState: (state) => {
void publishNoctaliaState(state).catch((error) =>
process.stderr.write(`${error.message}\n`),
);
},
...(agent ? { agent } : {}),
onState: (state) => void publishNoctaliaState(state).catch((error) =>
process.stderr.write(`${error.message}\n`),
),
});
await relay.start();
await new Promise((resolve) => {
@@ -41,8 +39,6 @@ async function run() {
}
run().catch((error) => {
process.stderr.write(
`${error instanceof Error ? error.message : "Noctalia relay failed"}\n`,
);
process.stderr.write(`${error instanceof Error ? error.message : "Noctalia relay failed"}\n`);
process.exitCode = 1;
});
+85 -41
View File
@@ -1,58 +1,102 @@
import path from "node:path";
function detailFor(state) {
function detailForSummary(summary) {
return [
`${summary.openCount} open`,
`${summary.workingCount} working`,
`${summary.attentionCount} attention`,
`${summary.recoveringCount} recovering`,
`${summary.errorCount} errors`,
].join(" · ");
}
function stateForSummary(summary) {
if (summary.errorCount) return "error";
if (summary.recoveringCount) return "recovering";
if (summary.workingCount) return "streaming";
return "idle";
}
function summaryState(summary) {
return {
state: stateForSummary(summary),
projectLabel: `${summary.openCount} session${summary.openCount === 1 ? "" : "s"}`,
attentionCount: summary.attentionCount,
detail: detailForSummary(summary),
openCount: summary.openCount,
workingCount: summary.workingCount,
recoveringCount: summary.recoveringCount,
errorCount: summary.errorCount,
};
}
function legacyDetail(state) {
switch (state) {
case "streaming":
return "Streaming";
case "recovering":
return "Recovering Pi session";
case "failed":
return "Recovery failed";
case "error":
return "Bridge error";
default:
return "Idle";
case "streaming": return "Streaming";
case "recovering": return "Recovering Pi session";
case "failed": return "Recovery failed";
case "error": return "Bridge error";
default: return "Idle";
}
}
export function createNoctaliaStateRelay({ client, agent, onState }) {
if (
!client ||
typeof client.request !== "function" ||
typeof client.subscribe !== "function"
)
throw new TypeError("client must support request and subscribe");
if (!agent?.id || !agent?.worktreePath)
/**
* Relays workspace-wide presentation state to Noctalia. Passing `agent` keeps
* the former agent-scoped behavior for external compatibility callers only.
*/
export function createNoctaliaStateRelay({ client, agent, onState, pollMs = 2_000 }) {
if (!client || typeof client.request !== "function")
throw new TypeError("client must support request");
if (agent && (!agent.id || !agent.worktreePath))
throw new TypeError("agent id and worktreePath are required");
if (typeof onState !== "function")
throw new TypeError("onState must be a function");
if (!agent && (!Number.isSafeInteger(pollMs) || pollMs < 250))
throw new TypeError("pollMs must be an integer of at least 250ms");
if (typeof onState !== "function") throw new TypeError("onState must be a function");
let state = "idle";
let attentionCount = 0;
let unsubscribe = () => {};
const projectLabel = path.basename(agent.worktreePath) || agent.worktreePath;
const publish = () =>
onState({ state, projectLabel, attentionCount, detail: detailFor(state) });
const handleEvent = (event) => {
if (event.type === "agent_state" && typeof event.data?.state === "string")
state = event.data.state;
if (event.type === "queue")
attentionCount =
(event.data?.event?.steering?.length ?? 0) +
(event.data?.event?.followUp?.length ?? 0);
if (event.type === "extension_ui_request")
attentionCount = Math.max(attentionCount, 1);
publish();
};
let timer;
let stopped = false;
const publishLegacy = (state, attentionCount) => onState({
state,
projectLabel: path.basename(agent.worktreePath) || agent.worktreePath,
attentionCount,
detail: legacyDetail(state),
openCount: 1,
workingCount: state === "streaming" ? 1 : 0,
recoveringCount: state === "recovering" ? 1 : 0,
errorCount: ["error", "failed"].includes(state) ? 1 : 0,
});
return {
async start() {
const response = await client.request("get_state", { agentId: agent.id });
state = response?.data?.isStreaming ? "streaming" : "idle";
publish();
unsubscribe = await client.subscribe(agent.id, 0, handleEvent);
if (agent) {
if (typeof client.subscribe !== "function")
throw new TypeError("legacy agent relay requires client.subscribe");
const response = await client.request("get_state", { agentId: agent.id });
let state = response?.data?.isStreaming ? "streaming" : "idle";
let attentionCount = 0;
publishLegacy(state, attentionCount);
unsubscribe = await client.subscribe(agent.id, 0, (event) => {
if (event.type === "agent_state" && typeof event.data?.state === "string") state = event.data.state;
if (event.type === "queue") attentionCount = (event.data?.event?.steering?.length ?? 0) + (event.data?.event?.followUp?.length ?? 0);
if (event.type === "extension_ui_request") attentionCount = Math.max(attentionCount, 1);
publishLegacy(state, attentionCount);
});
return;
}
const refresh = async () => {
try {
const summary = await client.request("get_workspace_summary");
if (!stopped) onState(summaryState(summary));
} finally {
if (!stopped) timer = setTimeout(() => void refresh(), pollMs);
}
};
await refresh();
},
stop() {
stopped = true;
clearTimeout(timer);
unsubscribe();
unsubscribe = () => {};
},
+74 -7
View File
@@ -8,6 +8,14 @@ const requestOperations = new Map([
["list_directories", { agent: false, payload: "none" }],
["select_agent", { agent: false, payload: "worktree" }],
["forget_directory", { agent: false, payload: "worktree" }],
["get_workspace", { agent: false, payload: "none" }],
["get_workspace_summary", { agent: false, payload: "none" }],
["create_session_runtime", { agent: false, payload: "worktree" }],
["open_session_runtime", { agent: false, payload: "worktreeSession" }],
["close_session_runtime", { agent: false, payload: "runtime" }],
["list_directory_sessions", { agent: false, payload: "worktree" }],
["get_session_runtime_snapshot", { agent: false, payload: "runtime" }],
["subscribe_workspace", { agent: false, payload: "cursor" }],
["get_state", { agent: true, payload: "none" }],
["get_session_stats", { agent: true, payload: "none" }],
["list_sessions", { agent: true, payload: "none" }],
@@ -25,6 +33,8 @@ const requestOperations = new Map([
["restart", { agent: true, payload: "none" }],
["get_available_models", { agent: true, payload: "none" }],
["get_commands", { agent: true, payload: "none" }],
["set_session_name", { agent: true, payload: "sessionName" }],
["compact", { agent: true, payload: "compact" }],
["set_model", { agent: true, payload: "model" }],
["set_thinking_level", { agent: true, payload: "thinking" }],
]);
@@ -100,10 +110,10 @@ function assertString(value, field, { maxLength = 4096, pattern } = {}) {
function assertOptionalCursor(value) {
if (value === undefined) return undefined;
if (Number.isSafeInteger(value) && value >= 0) return value;
return assertString(value, "payload.cursor", {
maxLength: 256,
pattern: idPattern,
});
throw new ProtocolError(
"invalid_message",
"payload.cursor must be a non-negative integer",
);
}
function validatePayload(kind, value) {
@@ -115,6 +125,7 @@ function validatePayload(kind, value) {
);
return undefined;
}
if (kind === "compact" && value === undefined) return undefined;
const payload = assertRecord(value, "payload");
switch (kind) {
@@ -137,9 +148,7 @@ function validatePayload(kind, value) {
const sessionPath = assertString(
payload.sessionPath,
"payload.sessionPath",
{
maxLength: 4096,
},
{ maxLength: 4096 },
);
if (!path.isAbsolute(sessionPath))
throw new ProtocolError(
@@ -148,6 +157,38 @@ function validatePayload(kind, value) {
);
return { sessionPath };
}
case "worktreeSession": {
assertAllowedKeys(
payload,
new Set(["worktreePath", "sessionPath"]),
"payload",
);
const worktreePath = assertString(
payload.worktreePath,
"payload.worktreePath",
{ maxLength: 4096 },
);
const sessionPath = assertString(
payload.sessionPath,
"payload.sessionPath",
{ maxLength: 4096 },
);
if (!path.isAbsolute(worktreePath) || !path.isAbsolute(sessionPath))
throw new ProtocolError(
"invalid_message",
"payload worktree and session paths must be absolute",
);
return { worktreePath, sessionPath };
}
case "runtime": {
assertAllowedKeys(payload, new Set(["runtimeId"]), "payload");
return {
runtimeId: assertString(payload.runtimeId, "payload.runtimeId", {
maxLength: 128,
pattern: idPattern,
}),
};
}
case "cursor": {
assertAllowedKeys(payload, new Set(["cursor"]), "payload");
return {
@@ -185,6 +226,32 @@ function validatePayload(kind, value) {
response,
};
}
case "sessionName": {
assertAllowedKeys(payload, new Set(["name"]), "payload");
const name = assertString(payload.name, "payload.name", {
maxLength: 4096,
}).trim();
if (!name)
throw new ProtocolError(
"invalid_message",
"payload.name must not be blank",
);
return { name };
}
case "compact": {
assertAllowedKeys(payload, new Set(["customInstructions"]), "payload");
const customInstructions = assertString(
payload.customInstructions,
"payload.customInstructions",
{ maxLength: 32 * 1024 },
).trim();
if (!customInstructions)
throw new ProtocolError(
"invalid_message",
"payload.customInstructions must not be blank",
);
return { customInstructions };
}
case "model": {
assertAllowedKeys(payload, new Set(["provider", "modelId"]), "payload");
return {