474 lines
15 KiB
JavaScript
474 lines
15 KiB
JavaScript
import assert from "node:assert/strict";
|
|
import {
|
|
lstat,
|
|
mkdir,
|
|
mkdtemp,
|
|
readFile,
|
|
rename,
|
|
rm,
|
|
writeFile,
|
|
} from "node:fs/promises";
|
|
import { tmpdir } from "node:os";
|
|
import { join } from "node:path";
|
|
import test from "node:test";
|
|
import {
|
|
createAgentRegistry,
|
|
promoteQuickSessionDirectory,
|
|
sessionDirectoryPath,
|
|
} from "../src/bridge/agent-registry.js";
|
|
import { createWorkspaceStore } from "../src/bridge/workspace-store.js";
|
|
|
|
function adapters() {
|
|
const calls = [];
|
|
return {
|
|
calls,
|
|
startAdapter(options) {
|
|
const adapter = {
|
|
sent: [],
|
|
stopped: false,
|
|
state: {},
|
|
async send(command) {
|
|
this.sent.push(command);
|
|
return {
|
|
type: "response",
|
|
command: command.type,
|
|
success: true,
|
|
data: command.type === "get_state" ? this.state : {},
|
|
};
|
|
},
|
|
respondToExtension() {},
|
|
async stop() {
|
|
this.stopped = true;
|
|
},
|
|
};
|
|
calls.push({ options, adapter });
|
|
return adapter;
|
|
},
|
|
};
|
|
}
|
|
|
|
test("promotes an active quick runtime in place and makes stale quick cleanup harmless", async () => {
|
|
const root = await mkdtemp(join(tmpdir(), "pi-quick-promote-"));
|
|
const home = join(root, "home");
|
|
await mkdir(home);
|
|
const fixture = adapters();
|
|
const registry = createAgentRegistry({
|
|
homeWorktree: home,
|
|
sessionRoot: join(root, "sessions"),
|
|
startAdapter: fixture.startAdapter,
|
|
});
|
|
await registry.start();
|
|
const quick = await registry.createQuickRuntime(home);
|
|
const quickAdapter = fixture.calls.at(-1).adapter;
|
|
await registry.route(quick.agentId, "prompt", { message: "keep working" });
|
|
const events = [];
|
|
const subscription = registry.subscribeWorkspace(0, (event) =>
|
|
events.push(event),
|
|
);
|
|
const [promoted, repeated] = await Promise.all([
|
|
registry.promoteQuickRuntime(quick.runtimeId),
|
|
registry.promoteQuickRuntime(quick.runtimeId),
|
|
]);
|
|
assert.equal(promoted.runtimeId, quick.runtimeId);
|
|
assert.equal(repeated.agentId, quick.agentId);
|
|
assert.equal(fixture.calls.at(-1).adapter, quickAdapter);
|
|
assert.equal(quickAdapter.stopped, false);
|
|
assert.equal(promoted.sessionPath, undefined);
|
|
assert.ok(events.some((event) => event.type === "runtime_promoted"));
|
|
assert.equal(
|
|
registry
|
|
.getWorkspace()
|
|
.directories[0].runtimes.some(
|
|
(runtime) => runtime.runtimeId === quick.runtimeId,
|
|
),
|
|
true,
|
|
);
|
|
assert.match(
|
|
await readFile(join(root, "sessions", "bridge-workspace-v2.json"), "utf8"),
|
|
new RegExp(quick.runtimeId),
|
|
);
|
|
assert.match(
|
|
await readFile(
|
|
join(
|
|
sessionDirectoryPath(join(root, "sessions"), home),
|
|
"bridge-agent.json",
|
|
),
|
|
"utf8",
|
|
),
|
|
/"worktreePath"/,
|
|
);
|
|
await registry.closeQuickRuntime(quick.runtimeId);
|
|
assert.equal(quickAdapter.stopped, false);
|
|
await registry.route(quick.agentId, "abort");
|
|
subscription.unsubscribe();
|
|
await registry.stop();
|
|
});
|
|
|
|
test("promotion refreshes active JSONL identity without replacing runtime state", async () => {
|
|
const root = await mkdtemp(join(tmpdir(), "pi-quick-identity-"));
|
|
const home = join(root, "home");
|
|
await mkdir(home);
|
|
const fixture = adapters();
|
|
const registry = createAgentRegistry({
|
|
homeWorktree: home,
|
|
sessionRoot: join(root, "sessions"),
|
|
startAdapter: fixture.startAdapter,
|
|
});
|
|
await registry.start();
|
|
const quick = await registry.createQuickRuntime(home);
|
|
const call = fixture.calls.at(-1);
|
|
const quickPath = join(call.options.sessionDir, "active.jsonl");
|
|
await writeFile(quickPath, '{"type":"session"}\n');
|
|
call.adapter.state = { sessionFile: quickPath, sessionId: "session-active" };
|
|
call.options.onEvent({
|
|
type: "queue",
|
|
data: { event: { pendingMessageCount: 2 } },
|
|
});
|
|
call.options.onEvent({
|
|
type: "extension_ui_request",
|
|
data: { event: { id: "extension-1", method: "confirm" } },
|
|
});
|
|
await registry.route(quick.agentId, "set_model", {
|
|
provider: "test",
|
|
modelId: "model-1",
|
|
});
|
|
await registry.route(quick.agentId, "set_thinking_level", { level: "high" });
|
|
|
|
const promoted = await registry.promoteQuickRuntime(quick.runtimeId);
|
|
const snapshot = await registry.getSessionRuntimeSnapshot(quick.runtimeId);
|
|
assert.equal(promoted.runtimeId, quick.runtimeId);
|
|
assert.equal(promoted.agentId, quick.agentId);
|
|
assert.equal(promoted.sessionId, "session-active");
|
|
assert.match(promoted.sessionPath, /active\.jsonl$/);
|
|
assert.equal(snapshot.runtime.queueCount, 2);
|
|
assert.equal(snapshot.extensions[0].id, "extension-1");
|
|
assert.deepEqual(
|
|
call.adapter.sent.filter((command) =>
|
|
["set_model", "set_thinking_level"].includes(command.type),
|
|
),
|
|
[
|
|
{ type: "set_model", provider: "test", modelId: "model-1" },
|
|
{ type: "set_thinking_level", level: "high" },
|
|
],
|
|
);
|
|
assert.equal(fixture.calls.at(-1).adapter, call.adapter);
|
|
await registry.stop();
|
|
});
|
|
|
|
test("session-directory migration rolls back move and symlink failures for retry", async () => {
|
|
const root = await mkdtemp(join(tmpdir(), "pi-quick-directory-rollback-"));
|
|
const quickDir = join(root, "quick");
|
|
const sessionDir = join(root, "session");
|
|
await mkdir(quickDir);
|
|
await writeFile(join(quickDir, "active.jsonl"), "history");
|
|
let renameCalls = 0;
|
|
await assert.rejects(
|
|
promoteQuickSessionDirectory(quickDir, sessionDir, {
|
|
renameFile: async (...args) => {
|
|
renameCalls += 1;
|
|
if (renameCalls === 2) throw new Error("move failed");
|
|
return rename(...args);
|
|
},
|
|
}),
|
|
);
|
|
assert.equal(
|
|
await readFile(join(quickDir, "active.jsonl"), "utf8"),
|
|
"history",
|
|
);
|
|
await assert.rejects(
|
|
promoteQuickSessionDirectory(quickDir, sessionDir, {
|
|
makeSymlink: async () => {
|
|
throw new Error("symlink failed");
|
|
},
|
|
}),
|
|
);
|
|
assert.equal(
|
|
await readFile(join(quickDir, "active.jsonl"), "utf8"),
|
|
"history",
|
|
);
|
|
const transaction = await promoteQuickSessionDirectory(quickDir, sessionDir);
|
|
await transaction.finalize();
|
|
assert.equal(
|
|
await readFile(join(sessionDir, "active.jsonl"), "utf8"),
|
|
"history",
|
|
);
|
|
});
|
|
|
|
test("promotion queues abort and close until identity persistence completes", async () => {
|
|
const root = await mkdtemp(join(tmpdir(), "pi-quick-promote-race-"));
|
|
const home = join(root, "home");
|
|
await mkdir(home);
|
|
const fixture = adapters();
|
|
const registry = createAgentRegistry({
|
|
homeWorktree: home,
|
|
sessionRoot: join(root, "sessions"),
|
|
startAdapter: fixture.startAdapter,
|
|
});
|
|
await registry.start();
|
|
const quick = await registry.createQuickRuntime(home);
|
|
const call = fixture.calls.at(-1);
|
|
const quickPath = join(call.options.sessionDir, "active.jsonl");
|
|
await writeFile(quickPath, "history");
|
|
call.adapter.state = { sessionFile: quickPath, sessionId: "session-race" };
|
|
const send = call.adapter.send.bind(call.adapter);
|
|
let releaseIdentity;
|
|
let blockIdentity = true;
|
|
call.adapter.send = async (command) => {
|
|
if (command.type !== "get_state" || !blockIdentity) return send(command);
|
|
return new Promise((resolve) => {
|
|
releaseIdentity = () => {
|
|
blockIdentity = false;
|
|
resolve({
|
|
type: "response",
|
|
command: "get_state",
|
|
success: true,
|
|
data: call.adapter.state,
|
|
});
|
|
};
|
|
});
|
|
};
|
|
const promotion = registry.promoteQuickRuntime(quick.runtimeId);
|
|
while (!releaseIdentity) await new Promise((resolve) => setImmediate(resolve));
|
|
const abort = registry.route(quick.agentId, "abort");
|
|
const close = registry.closeSessionRuntime(quick.runtimeId);
|
|
assert.notEqual(call.adapter.sent.at(-1).type, "abort");
|
|
releaseIdentity();
|
|
await Promise.all([promotion, abort, close]);
|
|
assert.ok(call.adapter.sent.some((command) => command.type === "abort"));
|
|
assert.equal(call.adapter.stopped, true);
|
|
await registry.stop();
|
|
});
|
|
|
|
test("promotion reserves moved JSONL before blocking get_state and releases on rollback", async () => {
|
|
const root = await mkdtemp(
|
|
join(tmpdir(), "pi-quick-promote-get-state-lease-"),
|
|
);
|
|
const home = join(root, "home");
|
|
const sessionRoot = join(root, "sessions");
|
|
await mkdir(home);
|
|
const fixture = adapters();
|
|
const startAdapter = (options) => {
|
|
const adapter = fixture.startAdapter(options);
|
|
if (options.sessionPath) adapter.state = { sessionFile: options.sessionPath };
|
|
return adapter;
|
|
};
|
|
const registry = createAgentRegistry({
|
|
homeWorktree: home,
|
|
sessionRoot,
|
|
startAdapter,
|
|
});
|
|
await registry.start();
|
|
const quick = await registry.createQuickRuntime(home);
|
|
const call = fixture.calls.at(-1);
|
|
const quickPath = join(call.options.sessionDir, "active.jsonl");
|
|
const promotedPath = join(
|
|
sessionDirectoryPath(sessionRoot, home),
|
|
"active.jsonl",
|
|
);
|
|
await writeFile(quickPath, "history");
|
|
call.adapter.state = {
|
|
sessionFile: quickPath,
|
|
sessionId: "session-get-state-lease",
|
|
};
|
|
call.options.onEvent({ type: "agent_state", data: { state: "idle" } });
|
|
while (
|
|
!(await registry.getSessionRuntimeSnapshot(quick.runtimeId)).runtime
|
|
.sessionPath
|
|
)
|
|
await new Promise((resolve) => setImmediate(resolve));
|
|
|
|
const send = call.adapter.send.bind(call.adapter);
|
|
let rejectState;
|
|
call.adapter.send = async (command) => {
|
|
if (command.type !== "get_state") return send(command);
|
|
return new Promise((_, reject) => {
|
|
rejectState = () => reject(new Error("get_state failed"));
|
|
});
|
|
};
|
|
const promotion = registry.promoteQuickRuntime(quick.runtimeId);
|
|
while (!rejectState) await new Promise((resolve) => setImmediate(resolve));
|
|
assert.equal((await lstat(call.options.sessionDir)).isSymbolicLink(), true);
|
|
assert.equal(await readFile(promotedPath, "utf8"), "history");
|
|
const childrenBeforeOpen = fixture.calls.length;
|
|
await assert.rejects(
|
|
registry.openSessionRuntime(home, promotedPath),
|
|
/already open/,
|
|
);
|
|
assert.equal(fixture.calls.length, childrenBeforeOpen);
|
|
|
|
rejectState();
|
|
await assert.rejects(promotion, /get_state failed/);
|
|
await writeFile(promotedPath, "history");
|
|
await registry.openSessionRuntime(home, promotedPath);
|
|
assert.equal(fixture.calls.length, childrenBeforeOpen + 1);
|
|
await registry.stop();
|
|
});
|
|
|
|
test("promotion reserves its canonical JSONL before reference and manifest I/O", async () => {
|
|
const root = await mkdtemp(join(tmpdir(), "pi-quick-promote-lease-"));
|
|
const home = join(root, "home");
|
|
const sessionRoot = join(root, "sessions");
|
|
await mkdir(home);
|
|
const fixture = adapters();
|
|
let quick;
|
|
let releasePromotionSave;
|
|
const registry = createAgentRegistry({
|
|
homeWorktree: home,
|
|
sessionRoot,
|
|
startAdapter: fixture.startAdapter,
|
|
workspaceStoreFactory(rootPath) {
|
|
const store = createWorkspaceStore(rootPath);
|
|
return {
|
|
...store,
|
|
save(next) {
|
|
if (
|
|
!next.runtimes.some((runtime) => runtime.runtimeId === quick?.runtimeId)
|
|
)
|
|
return store.save(next);
|
|
return new Promise((resolve, reject) => {
|
|
releasePromotionSave = () => store.save(next).then(resolve, reject);
|
|
});
|
|
},
|
|
};
|
|
},
|
|
});
|
|
await registry.start();
|
|
quick = await registry.createQuickRuntime(home);
|
|
const call = fixture.calls.at(-1);
|
|
const quickPath = join(call.options.sessionDir, "active.jsonl");
|
|
await writeFile(quickPath, "history");
|
|
call.adapter.state = { sessionFile: quickPath, sessionId: "session-lease" };
|
|
const promotedPath = join(
|
|
sessionDirectoryPath(sessionRoot, home),
|
|
"active.jsonl",
|
|
);
|
|
const promotion = registry.promoteQuickRuntime(quick.runtimeId);
|
|
while (!releasePromotionSave)
|
|
await new Promise((resolve) => setImmediate(resolve));
|
|
await assert.rejects(
|
|
registry.openSessionRuntime(home, promotedPath),
|
|
/already open/,
|
|
);
|
|
releasePromotionSave();
|
|
await promotion;
|
|
await registry.stop();
|
|
});
|
|
|
|
test("concurrent manifest persistence retains promoted runtime for restart", async () => {
|
|
const root = await mkdtemp(join(tmpdir(), "pi-quick-promote-persist-"));
|
|
const home = join(root, "home");
|
|
const sessionRoot = join(root, "sessions");
|
|
await mkdir(home);
|
|
const fixture = adapters();
|
|
let releasePromotionSave;
|
|
let holdPromotionSave = true;
|
|
let quick;
|
|
const registry = createAgentRegistry({
|
|
homeWorktree: home,
|
|
sessionRoot,
|
|
startAdapter: fixture.startAdapter,
|
|
workspaceStoreFactory(rootPath) {
|
|
const store = createWorkspaceStore(rootPath);
|
|
return {
|
|
...store,
|
|
save(next) {
|
|
if (
|
|
!holdPromotionSave ||
|
|
!next.runtimes.some((runtime) => runtime.runtimeId === quick?.runtimeId)
|
|
)
|
|
return store.save(next);
|
|
holdPromotionSave = false;
|
|
return new Promise((resolve, reject) => {
|
|
releasePromotionSave = () => store.save(next).then(resolve, reject);
|
|
});
|
|
},
|
|
};
|
|
},
|
|
});
|
|
await registry.start();
|
|
quick = await registry.createQuickRuntime(home);
|
|
const call = fixture.calls.at(-1);
|
|
const quickPath = join(call.options.sessionDir, "active.jsonl");
|
|
await writeFile(quickPath, "history");
|
|
call.adapter.state = { sessionFile: quickPath, sessionId: "session-persist" };
|
|
const promotion = registry.promoteQuickRuntime(quick.runtimeId);
|
|
while (!releasePromotionSave)
|
|
await new Promise((resolve) => setImmediate(resolve));
|
|
const concurrentOpen = registry.createSessionRuntime(home);
|
|
releasePromotionSave();
|
|
await Promise.all([promotion, concurrentOpen]);
|
|
const manifest = JSON.parse(
|
|
await readFile(join(sessionRoot, "bridge-workspace-v2.json"), "utf8"),
|
|
);
|
|
assert.ok(
|
|
manifest.runtimes.some((runtime) => runtime.runtimeId === quick.runtimeId),
|
|
);
|
|
await registry.stop();
|
|
|
|
const restored = adapters();
|
|
const originalStart = restored.startAdapter;
|
|
restored.startAdapter = (options) => {
|
|
const adapter = originalStart(options);
|
|
if (options.sessionPath) adapter.state = { sessionFile: options.sessionPath };
|
|
return adapter;
|
|
};
|
|
const restarted = createAgentRegistry({
|
|
homeWorktree: home,
|
|
sessionRoot,
|
|
startAdapter: restored.startAdapter,
|
|
});
|
|
await restarted.start();
|
|
assert.ok(
|
|
restarted
|
|
.getWorkspace()
|
|
.directories.flatMap((directory) => directory.runtimes)
|
|
.some(
|
|
(runtime) =>
|
|
runtime.runtimeId === quick.runtimeId && runtime.state !== "failed",
|
|
),
|
|
);
|
|
await restarted.stop();
|
|
});
|
|
|
|
test("promotion reference and manifest failures roll back quick runtime for retry", async () => {
|
|
const root = await mkdtemp(join(tmpdir(), "pi-quick-persist-rollback-"));
|
|
const home = join(root, "home");
|
|
await mkdir(home);
|
|
const fixture = adapters();
|
|
const sessionRoot = join(root, "sessions");
|
|
const registry = createAgentRegistry({
|
|
homeWorktree: home,
|
|
sessionRoot,
|
|
startAdapter: fixture.startAdapter,
|
|
});
|
|
await registry.start();
|
|
const quick = await registry.createQuickRuntime(home);
|
|
const quickDir = fixture.calls.at(-1).options.sessionDir;
|
|
await writeFile(join(quickDir, "active.jsonl"), "history");
|
|
const sessionDir = sessionDirectoryPath(sessionRoot, home);
|
|
const referencePath = join(sessionDir, "bridge-agent.json");
|
|
await mkdir(referencePath);
|
|
await assert.rejects(registry.promoteQuickRuntime(quick.runtimeId));
|
|
assert.equal(
|
|
await readFile(join(quickDir, "active.jsonl"), "utf8"),
|
|
"history",
|
|
);
|
|
assert.equal((await lstat(quickDir)).isDirectory(), true);
|
|
await rm(referencePath, { recursive: true });
|
|
const workspacePath = join(sessionRoot, "bridge-workspace-v2.json");
|
|
await rm(workspacePath);
|
|
await mkdir(workspacePath);
|
|
await assert.rejects(registry.promoteQuickRuntime(quick.runtimeId));
|
|
assert.equal(
|
|
await readFile(join(quickDir, "active.jsonl"), "utf8"),
|
|
"history",
|
|
);
|
|
assert.equal((await lstat(quickDir)).isDirectory(), true);
|
|
await rm(workspacePath, { recursive: true });
|
|
assert.equal(
|
|
(await registry.promoteQuickRuntime(quick.runtimeId)).runtimeId,
|
|
quick.runtimeId,
|
|
);
|
|
await registry.stop();
|
|
});
|