import assert from "node:assert/strict"; import { describe, it } from "node:test"; import * as fs from "node:fs"; import * as os from "node:os"; import * as path from "node:path"; import * as net from "node:net"; import { spawn } from "node:child_process"; import { createHash } from "node:crypto"; import type { ChildSession, ChildSessionFactory, ChildSessionLaunch } from "../../src/runs/shared/child-session.ts"; import { boundedHerdrReconnect, BridgeChannel, createPlacementAwareChildSessionFactory, HERDR_MAX_OWNED_PANES, HerdrPiSession, HerdrPlacedRunOwner, herdrPaneAllocationKey, ownsHerdrPane, parseHerdrAgent, parseHerdrCreated, parseHerdrSessionSnapshot, provisionHerdrPane, reconnectHerdrPiSession, serializeHerdrPiLaunch, withHerdrPaneAllocationLock, type HerdrRunIdentity } from "../../src/runs/shared/herdr-placed-run.ts"; import { HerdrPiFrameDecoder, encodeHerdrPiFrame, HERDR_PI_MAX_FRAME_BYTES, HERDR_PI_MODE_ENV, HERDR_PI_RUN_ENV, HERDR_PI_RUNTIME_DIR_ENV } from "../../src/runs/shared/herdr-pi-protocol.ts"; import { parseHerdrEndpoint } from "../../src/runs/shared/herdr-connection.ts"; import { formatHerdrMachineRunnerUnsupported, resolveHerdrMachinePlacement } from "../../src/runs/shared/herdr-machine.ts"; import { buildRunnerChildLaunch } from "../../src/runs/background/runner-child-launch.ts"; import { runChildSession } from "../../src/runs/background/run-child-session.ts"; import { getAgentDir } from "../../src/shared/utils.ts"; import registerHerdrPiBridge, { resolveRemoteHerdrResources } from "../../src/extension/herdr-pi-bridge.ts"; const machine = { provider: "herdr" as const, id: "m1", label: "workmac", target: "remote.example", cwd: "/remote/repo" }; const runtime = { runId: "run", agent: "worker", childIndex: 0, fanoutChild: false, depth: 1, maxDepth: 1, inheritProjectContext: true, inheritGlobalContext: true, inheritSkills: false, fast: false, waitTool: { enabled: true } }; const packageVersion = (JSON.parse(fs.readFileSync(new URL("../../package.json", import.meta.url), "utf8")) as { version: string }).version; function launch(remote: boolean): ChildSessionLaunch { return { cwd: remote ? machine.cwd : "/local/repo", ...(remote ? { machine, remoteResources: { agent: "worker" } } : {}), storage: { kind: "memory" }, extensionPaths: [], ambientExtensions: true, hooks: [], noSkills: true, noContextFiles: true, tools: ["read"], runtime }; } const child: ChildSession = { subscribe: () => () => {}, prompt: async () => {}, steer: async () => {}, followUp: async () => {}, abort: async () => {}, dispose: async () => {}, messages: [], sessionFile: undefined, sessionId: "native", modelId: "provider/model" }; describe("pane-native Herdr placement public contracts", { skip: process.platform === "win32" }, () => { it("preserves explicit default versus named endpoint identity", () => { assert.equal(parseHerdrEndpoint(JSON.stringify({ socket: "/tmp/herdr.sock", session: null, version: "0.9.0", protocol: 22, running: true, compatible: true })).session, null); assert.equal(parseHerdrEndpoint(JSON.stringify({ socket: "/tmp/herdr.sock", session: null, version: "0.9.0", protocol: 22, running: true, compatible: true }), "default").session, null); assert.equal(parseHerdrEndpoint(JSON.stringify({ socket: "/tmp/herdr.sock", session: "default", version: "0.9.0", protocol: 22, running: true, compatible: true })).session, null); assert.equal(parseHerdrEndpoint(JSON.stringify({ socket: "/tmp/named.sock", session: "team", version: "0.9.0", protocol: 22, running: true, compatible: true }), "team").session, "team"); assert.throws(() => parseHerdrEndpoint(JSON.stringify({ socket: "/tmp/wrong.sock", session: "wrong", version: "0.9.0", protocol: 22, running: true, compatible: true }), "team"), /identity mismatch/u); assert.throws(() => parseHerdrEndpoint(JSON.stringify({ socket: "/tmp/herdr.sock", session: null, version: "0.9.0", protocol: 22, running: false, compatible: true })), /stopped or incompatible/u); }); it("allows native Pi placement while retaining unsupported-runner and managed-worktree gates", () => { assert.equal(formatHerdrMachineRunnerUnsupported({ machine: "workmac", agentName: "worker" }), process.platform === "win32" ? "Herdr saved-machine pane transport requires hardened OpenSSH StreamLocal forwarding, which is not supported from a Windows host yet." : undefined); assert.match(formatHerdrMachineRunnerUnsupported({ machine: "workmac", agentName: "job", runnerType: "external-job" }) ?? "", /cannot use pane-native/u); assert.match(formatHerdrMachineRunnerUnsupported({ machine: "workmac", agentName: "worker", worktree: true }) ?? "", /managed worktrees/u); }); it("rejects local machine environment with an actionable remote-configuration message", () => { assert.throws(() => resolveHerdrMachinePlacement({ machine: "workmac", cwd: "/local/repo", catalogJson: JSON.stringify([{ id: "m1", label: "workmac", target: "remote.example", enabled: true }]), settings: { cwd: "/remote/repo", env: { API_TOKEN: "secret" } } }), /configured remotely.*remove machines\.workmac\.env/iu); }); it("routes only placed native launches through the remote factory", async () => { let local = 0; let remote = 0; let seen: ChildSessionLaunch | undefined; const localFactory: ChildSessionFactory = { create: async () => { local++; return child; }, dispose: async () => {} }; const factory = createPlacementAwareChildSessionFactory(localFactory, async (value) => { remote++; seen = value; return child; }); await factory.create(launch(false)); await factory.create(launch(true)); assert.equal(local, 1); assert.equal(remote, 1); assert.equal(seen?.machine?.cwd, "/remote/repo"); assert.equal(seen?.cwd, "/remote/repo"); }); it("rejects unrepresented required and MCP tool contracts at the public factory boundary before provisioning", async () => { let provisioned = 0; const factory = createPlacementAwareChildSessionFactory({ async create() { return child; }, async dispose() {} }, async () => { provisioned++; return child; }); await assert.rejects(factory.create({ ...launch(true), runtime: { ...runtime, requiredTools: ["required_missing"] } }), /required-tool/u); await assert.rejects(factory.create({ ...launch(true), runtime: { ...runtime, mcpDirectTools: ["mcp_tool"] } }), /MCP direct tools/u); assert.equal(provisioned, 0); }); it("propagates resolved placement into the detached native child launch", () => { const built = buildRunnerChildLaunch({ agent: "worker", task: "inspect", machine, cwd: machine.cwd, tools: ["read"], inheritProjectContext: true, inheritGlobalContext: true, inheritSkills: false }, { cwd: "/local/repo", id: "run-route", flatIndex: 0 }, { sessionEnabled: true, watchdogStatus() {} }); assert.deepEqual(built.session.machine, machine); assert.equal(built.session.cwd, machine.cwd); assert.deepEqual(built.session.remoteResources?.toolCeiling, ["read"]); }); it("projects native Git evidence through the detached child result", async () => { const built = buildRunnerChildLaunch({ agent: "worker", task: "inspect", machine, cwd: machine.cwd, tools: ["read"], inheritProjectContext: true, inheritGlobalContext: true, inheritSkills: false }, { cwd: "/local/repo", id: "run-native-result", flatIndex: 0 }, { sessionEnabled: true, watchdogStatus() {} }); const listeners = new Set<(event: any) => void>(); const placed: ChildSession = { ...child, machineEvidence: { machineId: "m1", initial: { head: "a", dirty: false }, final: { head: "b", dirty: true } }, subscribe(listener) { listeners.add(listener); return () => listeners.delete(listener); }, async prompt() { const message = { role: "assistant", content: [{ type: "text", text: "done" }], stopReason: "stop", usage: { input: 1, output: 1, cacheRead: 0, cacheWrite: 0, cost: { total: 0 } } }; for (const listener of listeners) listener({ type: "message_end", message }); for (const listener of listeners) listener({ type: "agent_settled" }); } }; const result = await runChildSession({ factory: { async create() { return placed; }, async dispose() {} }, launch: built, prompt: "Task: inspect", appendChildEvent() {}, writeOutputLine() {} }); assert.deepEqual(result.nativeMachine, { provider: "herdr", machineId: "m1", initialGit: { head: "a", dirty: false }, finalGit: { head: "b", dirty: true } }); }); it("keeps native launch data-only and rejects local resources, callbacks, fork files, and nested delegation", () => { const clean = { ...launch(true), systemPrompt: "LOCAL_SECRET_PROMPT /Users/local/private" }; const wire = serializeHerdrPiLaunch(clean); assert.equal("tools" in wire, false); assert.doesNotMatch(JSON.stringify(wire), /remote\.example|\/local\/repo|callback|token/iu); assert.deepEqual(serializeHerdrPiLaunch({ ...clean, remoteResources: { ...clean.remoteResources!, reads: ["remote-only.md"] } }).resources.reads, ["remote-only.md"]); assert.doesNotMatch(JSON.stringify(wire), /LOCAL_SECRET_PROMPT|Users\/local/u); assert.throws(() => serializeHerdrPiLaunch({ ...clean, storage: { kind: "file", sessionFile: "/local/session.jsonl" } }), /fork, resume, or revival/u); assert.throws(() => serializeHerdrPiLaunch({ ...clean, extensionPaths: ["/local/private.ts"] }), /cannot transfer local extension paths/u); assert.doesNotMatch(JSON.stringify(serializeHerdrPiLaunch({ ...clean, hooks: [{ name: "local", factory() {} }] })), /callback|factory|local/u); assert.throws(() => serializeHerdrPiLaunch({ ...clean, runtime: { ...runtime, nestedRoute: { rootRunId: "r", eventSink: "/local/events", controlInbox: "/local/control", capabilityToken: "SECRET" } } }), /Nested delegation/u); assert.throws(() => serializeHerdrPiLaunch({ ...clean, processEnv: { MCP_DIRECT_TOOLS: "local-server" } }), /bindings, MCP selections, or process environment/u); assert.throws(() => serializeHerdrPiLaunch({ ...clean, runtime: { ...runtime, permissions: { rules: { bash: "deny" } } as never } }), /permission rules/u); }); it("parses literal tagged Herdr 0.9 success envelopes", () => { assert.deepEqual(parseHerdrCreated({ type: "workspace_created", workspace: { workspace_id: "w" }, tab: { tab_id: "t" }, root_pane: { pane_id: "p" } }, "workspace_created"), { workspaceId: "w", tabId: "t", paneId: "p" }); assert.deepEqual(parseHerdrCreated({ type: "tab_created", tab: { tab_id: "t2" }, root_pane: { pane_id: "p2" } }, "tab_created"), { tabId: "t2", paneId: "p2" }); assert.equal(parseHerdrAgent({ type: "agent_started", agent: { terminal_id: "term", pane_id: "p" }, argv: [] }, "agent_started").terminal_id, "term"); assert.equal(parseHerdrAgent({ type: "agent_info", agent: { terminal_id: "term" } }, "agent_info").terminal_id, "term"); assert.deepEqual(parseHerdrSessionSnapshot({ type: "session_snapshot", snapshot: { workspaces: [] } }), { workspaces: [] }); }); it("subscribable provisioning uses session.snapshot and tagged tab envelopes", { skip: process.platform === "win32" }, async () => { const methods: string[] = []; const client = { async call(method: string) { methods.push(method); if (method === "session.snapshot") return { type: "session_snapshot", snapshot: { workspaces: [{ workspace_id: "w", label: "repo" }], panes: [{ workspace_id: "w", cwd: "/remote/repo" }], tabs: [], agents: [], layouts: [], version: "0.9.0", protocol: 1 } }; if (method === "tab.create") return { type: "tab_created", tab: { tab_id: "t" }, root_pane: { pane_id: "p" } }; throw new Error(method); }, async subscribe() { return () => {}; } }; assert.deepEqual(await provisionHerdrPane(client, "/remote/repo", "run_12345678", "/tmp/runtime"), { workspaceId: "w", tabId: "t", paneId: "p" }); assert.deepEqual(methods, ["session.snapshot", "tab.create"]); }); it("caps a selected workspace at twenty panes before allocation", async () => { let creates = 0; const panes = Array.from({ length: HERDR_MAX_OWNED_PANES }, (_, index) => ({ workspace_id: "w", pane_id: `p${index}`, cwd: "/remote/repo" })); const client = { async call(method: string) { if (method === "session.snapshot") return { type: "session_snapshot", snapshot: { workspaces: [{ workspace_id: "w", label: "repo" }], panes } }; creates++; return {}; } }; await assert.rejects(provisionHerdrPane(client, "/remote/repo", "cap", "/tmp/runtime", undefined, "cap-key"), /pane limit of 20/u); assert.equal(creates, 0); }); it("serializes concurrent reservations so only one twentieth pane allocates", async () => { const panes = Array.from({ length: HERDR_MAX_OWNED_PANES - 1 }, (_, index) => ({ workspace_id: "w", pane_id: `p${index}`, cwd: "/remote/repo" })); let creates = 0; const client = { async call(method: string) { if (method === "session.snapshot") return { type: "session_snapshot", snapshot: { workspaces: [{ workspace_id: "w", label: "repo" }], panes: [...panes] } }; if (method === "tab.create") { const paneId = `new-${++creates}`; panes.push({ workspace_id: "w", pane_id: paneId, cwd: "/remote/repo" }); await new Promise((resolve) => setTimeout(resolve, 5)); return { type: "tab_created", tab: { tab_id: `t-${creates}` }, root_pane: { pane_id: paneId } }; } throw new Error(method); } }; const results = await Promise.allSettled([provisionHerdrPane(client, "/remote/repo", "one", "/tmp/one", undefined, "shared-key"), provisionHerdrPane(client, "/remote/repo", "two", "/tmp/two", undefined, "shared-key")]); assert.equal(results.filter((result) => result.status === "fulfilled").length, 1); assert.equal(results.filter((result) => result.status === "rejected" && /pane limit of 20/u.test(String(result.reason))).length, 1); assert.equal(creates, 1); }); it("uses alias-equivalent process locks, never steals a live holder, and recovers a dead holder", { skip: process.platform === "win32" }, async () => { const keyA = herdrPaneAllocationKey("host.example", "default", "/remote/repo"), keyB = herdrPaneAllocationKey("host.example", "default", "/remote/repo"); assert.equal(keyA, keyB); const fixture = path.resolve("test/fixtures/herdr-allocation-lock-child.ts"), encodedKey = Buffer.from(keyA).toString("base64url"), launch = (hold: number, timeout: number) => spawn(process.execPath, ["--experimental-strip-types", fixture, encodedKey, String(hold), String(timeout)], { env: process.env, stdio: ["ignore", "pipe", "pipe"] }), waitFor = (child: ReturnType, stream: "stdout" | "stderr", match: RegExp) => new Promise((resolve, reject) => { let text = ""; child[stream].on("data", (chunk) => { text += chunk; if (match.test(text)) resolve(text); }); child.on("error", reject); child.on("exit", (code) => { if (!match.test(text)) reject(new Error(`child exited ${code}: ${text}`)); }); }); const holder = launch(30_000, 2_000); await waitFor(holder, "stdout", /locked/u); const publishedPath = path.join(getAgentDir(), "pi-subagents", "herdr-allocation-locks", `${createHash("sha256").update(keyA).digest("hex")}.lock`), publishedStat = fs.lstatSync(publishedPath), published = JSON.parse(fs.readFileSync(publishedPath, "utf8")) as { pid: number; token: string }; assert.equal(publishedStat.isFile() && !publishedStat.isSymbolicLink() && (publishedStat.mode & 0o777) === 0o600, true); assert.equal(published.pid, holder.pid); assert.match(published.token, /^[a-f0-9]{32}$/u); const blocked = launch(0, 100), blockedExit = new Promise((resolve) => blocked.once("exit", resolve)); await waitFor(blocked, "stderr", /Timed out waiting/u); assert.equal(await blockedExit, 2); holder.kill("SIGKILL"); await new Promise((resolve) => holder.once("exit", resolve)); let acquired = false; await withHerdrPaneAllocationLock(keyB, async () => { acquired = true; }, { timeoutMs: 2_000, pollMs: 10 }); assert.equal(acquired, true); }); it("publishes only validated owner files, refuses unsafe owners, and preserves action plus release failures", { skip: process.platform === "win32" }, async () => { const root = path.join(getAgentDir(), "pi-subagents", "herdr-allocation-locks"), lockPath = (key: string) => path.join(root, `${createHash("sha256").update(key).digest("hex")}.lock`); await withHerdrPaneAllocationLock("initialize-safety-root", async () => {}); for (const kind of ["mode", "symlink"] as const) { const key = `unsafe-${kind}-${Date.now()}`, canonical = lockPath(key), target = `${canonical}.target`; try { fs.writeFileSync(target, JSON.stringify({ pid: process.pid, token: "a".repeat(32) }), { mode: kind === "mode" ? 0o644 : 0o600 }); if (kind === "symlink") fs.symlinkSync(target, canonical); else fs.renameSync(target, canonical); await assert.rejects(withHerdrPaneAllocationLock(key, async () => {}, { timeoutMs: 50, pollMs: 5 }), /owner record is unsafe|ELOOP/u); } finally { fs.rmSync(canonical, { force: true }); fs.rmSync(target, { force: true }); } } const key = `release-replacement-${Date.now()}`, canonical = lockPath(key), actionFailure = new Error("action failed"); try { await assert.rejects(withHerdrPaneAllocationLock(key, async () => { fs.unlinkSync(canonical); fs.writeFileSync(canonical, JSON.stringify({ pid: process.pid, token: "b".repeat(32) }), { mode: 0o600 }); throw actionFailure; }), (error) => error instanceof AggregateError && error.cause === actionFailure && error.errors[0] === actionFailure && /ownership changed/u.test(String(error.errors[1]))); } finally { fs.rmSync(canonical, { force: true }); } }); it("preserves falsy action rejections with and without a release failure", async () => { let rejected = false; try { await withHerdrPaneAllocationLock(`falsy-${Date.now()}`, async () => Promise.reject(undefined)); } catch (error) { rejected = true; assert.equal(error, undefined); } assert.equal(rejected, true); const root = path.join(getAgentDir(), "pi-subagents", "herdr-allocation-locks"), key = `falsy-release-${Date.now()}`, canonical = path.join(root, `${createHash("sha256").update(key).digest("hex")}.lock`); try { await assert.rejects(withHerdrPaneAllocationLock(key, async () => { fs.unlinkSync(canonical); fs.writeFileSync(canonical, JSON.stringify({ pid: process.pid, token: "b".repeat(32) }), { mode: 0o600 }); return Promise.reject(undefined); }), (error) => error instanceof AggregateError && error.cause === undefined && error.errors[0] === undefined && /ownership changed/u.test(String(error.errors[1]))); } finally { fs.rmSync(canonical, { force: true }); } }); it("prevents two detached aliases from allocating past twenty panes", { skip: process.platform === "win32" }, async () => { const dir = fs.mkdtempSync(path.join(os.tmpdir(), "herdr-cap-")), statePath = path.join(dir, "state.json"), panes = Array.from({ length: HERDR_MAX_OWNED_PANES - 1 }, (_, index) => ({ workspace_id: "w", pane_id: `p${index}`, cwd: "/remote/repo" })); fs.writeFileSync(statePath, JSON.stringify({ workspaces: [{ workspace_id: "w", label: "repo" }], panes })); const fixture = path.resolve("test/fixtures/herdr-provision-child.ts"), key = Buffer.from(herdrPaneAllocationKey("same-target", "session", "/remote/repo")).toString("base64url"), launch = (id: string) => spawn(process.execPath, ["--experimental-strip-types", fixture, statePath, key, id], { env: process.env, stdio: ["ignore", "pipe", "pipe"] }), finish = (child: ReturnType) => new Promise<{ code: number | null; out: string; err: string }>((resolve, reject) => { let out = "", err = ""; child.stdout.on("data", (value) => { out += value; }); child.stderr.on("data", (value) => { err += value; }); child.on("error", reject); child.on("exit", (code) => resolve({ code, out, err })); }); try { const children = [launch("alias-a"), launch("alias-b")], results = await Promise.all(children.map(finish)); assert.equal(results.filter((result) => result.code === 0 && /allocated/u.test(result.out)).length, 1); assert.equal(results.filter((result) => result.code === 2 && /pane limit of 20/u.test(result.err)).length, 1); assert.equal((JSON.parse(fs.readFileSync(statePath, "utf8")) as { panes: unknown[] }).panes.length, 20); } finally { fs.rmSync(dir, { recursive: true, force: true }); } }); it("rejects duplicate deterministic workspaces and reuses one despite cwd projection drift", async () => { const label = `pi-subagents-${createHash("sha256").update("/remote/repo").digest("hex").slice(0, 10)}`; let creates = 0; const duplicate = { async call(method: string) { if (method === "session.snapshot") return { type: "session_snapshot", snapshot: { workspaces: [{ workspace_id: "a", label }, { workspace_id: "b", label }], panes: [] } }; creates++; return {}; } }; await assert.rejects(provisionHerdrPane(duplicate, "/remote/repo", "dup", "/tmp/dup", undefined, "dup-key"), /owned workspace identity is ambiguous/u); assert.equal(creates, 0); const reused = { async call(method: string, params: Record) { if (method === "session.snapshot") return { type: "session_snapshot", snapshot: { workspaces: [{ workspace_id: "owned", label }], panes: [{ workspace_id: "owned", cwd: "/changed" }] } }; assert.equal(method, "tab.create"); assert.equal(params.workspace_id, "owned"); return { type: "tab_created", tab: { tab_id: "t" }, root_pane: { pane_id: "p" } }; } }; assert.equal((await provisionHerdrPane(reused, "/remote/repo", "reuse", "/tmp/reuse", undefined, "reuse-key")).workspaceId, "owned"); }); it("creates one deterministic fallback when multiple unowned cwd workspaces exist", async () => { let workspaceCreates = 0; const label = `pi-subagents-${createHash("sha256").update("/remote/repo").digest("hex").slice(0, 10)}`; const workspaces: Record[] = [{ workspace_id: "a", label: "a" }, { workspace_id: "b", label: "b" }]; const panes: Record[] = [{ workspace_id: "a", cwd: "/remote/repo" }, { workspace_id: "b", cwd: "/remote/repo" }]; const client = { async call(method: string, params: Record) { if (method === "session.snapshot") return { type: "session_snapshot", snapshot: { workspaces: [...workspaces], panes: [...panes] } }; if (method === "workspace.create") { workspaceCreates++; workspaces.push({ workspace_id: "owned", label }); panes.push({ workspace_id: "owned", cwd: "/remote/repo" }); return { type: "workspace_created", workspace: { workspace_id: "owned" }, tab: { tab_id: "t" }, root_pane: { pane_id: "p" } }; } if (method === "tab.create") { assert.equal(params.workspace_id, "owned"); return { type: "tab_created", tab: { tab_id: "t2" }, root_pane: { pane_id: "p2" } }; } throw new Error(method); } }; await provisionHerdrPane(client, "/remote/repo", "first", "/tmp/first", undefined, "fallback-key"); await provisionHerdrPane(client, "/remote/repo", "second", "/tmp/second", undefined, "fallback-key"); assert.equal(workspaceCreates, 1); }); it("permits destructive cleanup only for the exact current terminal and pane", () => { assert.equal(ownsHerdrPane({ terminal_id: "term", pane_id: "pane" }, "term", "pane"), true); assert.equal(ownsHerdrPane({ terminal_id: "term", pane_id: "moved" }, "term", "pane"), false); assert.equal(ownsHerdrPane({ terminal_id: "reused", pane_id: "pane" }, "term", "pane"), false); assert.equal(ownsHerdrPane({}, "term", "pane"), false); }); it("starts immediately with one owner RPC and retains the exact requested argv", async () => { const methods: string[] = [], connection = { client: { async call(method: string) { methods.push(method); return { type: "agent_started", agent: { terminal_id: "term", pane_id: "pane" } }; } } }, owner = new (HerdrPlacedRunOwner as unknown as new (...args: unknown[]) => HerdrPlacedRunOwner)({ id: "m", target: "host", cwd: "/repo" }, "run", "/tmp/pi-subagents-herdr-owned", connection, () => {}); owner.owned = { workspaceId: "w", tabId: "t", paneId: "pane" }; const started = await owner.start("claude", "a", ["--restricted"]); assert.deepEqual(methods, ["agent.start"]); assert.deepEqual(started.argv, ["claude", "--restricted"]); assert.deepEqual(owner.startedArgv, started.argv); }); it("ignores the superseded connection's disconnect callback during replacement", async () => { let lost!: (error: Error) => void, disconnected = 0; const original = { client: { async subscribe(filters: unknown, _listener: unknown, onDisconnect: (error: Error) => void) { assert.deepEqual(filters, [{ type: "pane.agent_status_changed", pane_id: "p" }, { type: "pane.closed" }, { type: "pane.moved" }]); lost = onDisconnect; return () => {}; } }, async close() { lost(new Error("superseded connection closed")); } }, replacement = { client: {}, async close() {} }, owner = new (HerdrPlacedRunOwner as unknown as new (...args: unknown[]) => HerdrPlacedRunOwner)(machine, "stale-disconnect", "/tmp/pi-subagents-herdr-stale", original, () => {}); owner.owned = { workspaceId: "w", tabId: "t", paneId: "p" }; await owner.subscribe(() => {}, () => { disconnected++; }); await owner.replaceConnection({ connection: replacement as never, unsubscribe: () => {}, paneId: "p", state: "ready" }); assert.equal(disconnected, 0); assert.equal(owner.snapshot.connection, "connected"); }); it("rejects an initial subscription lost after acknowledgement but before adoption", async () => { let starts = 0, stops = 0; const connection = { client: { async subscribe(_filters: unknown, _listener: unknown, lost: (error: Error) => void) { lost(new Error("post-ack loss")); return () => { stops++; }; } } }, owner = new (HerdrPlacedRunOwner as unknown as new (...args: unknown[]) => HerdrPlacedRunOwner)(machine, "initial-race", "/tmp/pi-subagents-herdr-race", connection, () => {}); owner.owned = { workspaceId: "w", tabId: "t", paneId: "p" }; await assert.rejects(owner.subscribe(() => {}, () => { starts++; }), /post-ack loss/u); assert.equal(starts, 0); assert.equal(stops, 1); assert.equal(owner.snapshot.connection, "connected"); }); it("never adopts a Pi reconnect candidate lost before subscribe returns", async () => { const identity: HerdrRunIdentity = { runId: "run_12345678", machineId: "m1", target: "remote", session: null, workspaceId: "w", tabId: "t", paneId: "p", terminalId: "term", agentName: "pi-r", nativeSessionId: "native", cwd: "/remote", runtimeDir: "/tmp/run" }; let connects = 0, closes = 0, replacements = 0; const fakeSession = { identity, evidenceCursor: 0, placementSnapshot: { connection: "unknown" }, observeHerdr() {}, reconnect() {}, async replaceTransport() { replacements++; } } as never; const fakeBridge = () => ({ frames: [], failure: new Promise(() => {}), listeners: { add(listener: (frame: unknown) => void) { listener({ type: "ready", configured: true, runId: identity.runId, nativeSessionId: identity.nativeSessionId, packageVersion, protocol: 1, evidenceCursor: 0, evidenceFloor: 0 }); }, delete() {} }, close() {} }) as never; const connect = async () => { connects++; return { endpoint: { session: null, protocol: 22, version: "0.9.0" }, client: { async call() { return { type: "session_snapshot", snapshot: { agents: [{ terminal_id: "term", pane_id: "p", agent_status: "idle" }] } }; }, async subscribe(_filters: unknown, _listener: unknown, lost: (error: Error) => void) { lost(new Error("post-ack Pi loss")); return () => {}; } }, async forwardRemoteSocket() { return { socketPath: "/unused", async close() {} }; }, async close() { closes++; } }; }; await assert.rejects(reconnectHerdrPiSession(fakeSession, { ...launch(true), machine: { ...machine, cwd: "/remote" } }, { connect: connect as never, discoverManifest: (() => ({ socketPath: "/remote/bridge", nativeSessionId: "native", packageVersion, protocol: 1 })) as never, createBridge: fakeBridge }), /post-ack Pi loss/u); assert.equal(replacements, 0); assert.equal(connects, 3); assert.equal(closes, 3); }); it("closes each adopted transport resource once when disposal races replacement", async () => { const identity: HerdrRunIdentity = { runId: "run_12345678", machineId: "m1", target: "remote", session: null, workspaceId: "w", tabId: "t", paneId: "p", terminalId: "term", agentName: "pi-r", nativeSessionId: "native", cwd: "/remote", runtimeDir: "/tmp/run" }; let releaseOld!: () => void, oldCloseStarted!: () => void, connectionCloses = 0, bridgeForwardCloses = 0, bridgeCloses = 0, unsubscribes = 0, cleanupCalls = 0; const oldClosePending = new Promise((resolve) => { releaseOld = resolve; }), oldCloseEntered = new Promise((resolve) => { oldCloseStarted = resolve; }); const adoptedBridge = { frames: [], listeners: new Set(), close() { bridgeCloses++; } } as unknown as BridgeChannel; const connection = { async close() { connectionCloses++; } } as never; let cleanupPromise: Promise | undefined; const owner = { identity, connection: { async close() {} }, unsubscribe: () => {}, snapshot: { connection: "connected", state: "ready", identity, revision: 0 }, async replaceConnection(input: { connection: typeof connection; unsubscribe: () => void }) { const old = this.connection; this.connection = input.connection; this.unsubscribe = input.unsubscribe; await old.close(); }, cleanup() { return cleanupPromise ??= (async () => { cleanupCalls++; this.unsubscribe(); await this.connection.close(); })(); } } as unknown as HerdrPlacedRunOwner; const session = new HerdrPiSession(owner, { close: async () => { oldCloseStarted(); await oldClosePending; } }, { frames: [], listeners: new Set(), close() {} } as unknown as BridgeChannel, undefined, undefined, runtime); const replacing = session.replaceTransport({ connection, bridgeForward: { async close() { bridgeForwardCloses++; } }, bridge: adoptedBridge, unsubscribe() { unsubscribes++; }, paneId: "p2", cursor: 0, state: "working" }); await oldCloseEntered; await session.dispose(); releaseOld(); await replacing; assert.deepEqual({ connectionCloses, bridgeForwardCloses, bridgeCloses, unsubscribes, cleanupCalls }, { connectionCloses: 1, bridgeForwardCloses: 1, bridgeCloses: 1, unsubscribes: 1, cleanupCalls: 1 }); }); it("retries the identical start only after a transient rejection and fresh unoccupied pane proof", async () => { const methods: string[] = []; let starts = 0, now = 0; const connection = { client: { async call(method: string, params: unknown) { methods.push(method); if (method === "pane.get") return { type: "pane_info", pane: { pane_id: "pane", workspace_id: "w", tab_id: "t", agent: null, agent_session: null } }; if (++starts === 1) throw new Error("agent target pane pane is not an available shell"); assert.deepEqual(params, { name: "a", kind: "codex", pane_id: "pane", args: ["--sandbox", "read-only"], timeout_ms: 45_000 }); return { type: "agent_started", agent: { terminal_id: "term", pane_id: "pane" } }; } } }, owner = new (HerdrPlacedRunOwner as unknown as new (...args: unknown[]) => HerdrPlacedRunOwner)({ id: "m", target: "host", cwd: "/repo" }, "run", "/tmp/pi-subagents-herdr-owned", connection, () => {}); owner.owned = { workspaceId: "w", tabId: "t", paneId: "pane" }; const result = await owner.start("codex", "a", ["--sandbox", "read-only"], { timeoutMs: 20, pollMs: 10, clock: { now: () => now, async wait(ms) { now += ms; } } }); assert.deepEqual(methods, ["agent.start", "pane.get", "agent.start"]); assert.equal(starts, 2); assert.deepEqual(result, { terminalId: "term", paneId: "pane", argv: ["codex", "--sandbox", "read-only"] }); }); it("does not retry a non-transient start error", async () => { let calls = 0; const connection = { client: { async call() { calls++; throw new Error("permission denied"); } } }, owner = new (HerdrPlacedRunOwner as unknown as new (...args: unknown[]) => HerdrPlacedRunOwner)({ id: "m", target: "host", cwd: "/repo" }, "run", "/tmp/pi-subagents-herdr-owned", connection, () => {}); owner.owned = { workspaceId: "w", tabId: "t", paneId: "pane" }; await assert.rejects(owner.start("codex", "a", []), /permission denied/u); assert.equal(calls, 1); }); it("rejects pane drift or occupation after a transient start rejection", async () => { for (const pane of [{ pane_id: "other", workspace_id: "w", tab_id: "t", agent: null }, { pane_id: "pane", workspace_id: "w", tab_id: "t", agent: "other" }]) { let starts = 0; const connection = { client: { async call(method: string) { if (method === "pane.get") return { type: "pane_info", pane }; starts++; throw new Error("is not an available shell"); } } }, owner = new (HerdrPlacedRunOwner as unknown as new (...args: unknown[]) => HerdrPlacedRunOwner)({ id: "m", target: "host", cwd: "/repo" }, "run", "/tmp/pi-subagents-herdr-owned", connection, () => {}); owner.owned = { workspaceId: "w", tabId: "t", paneId: "pane" }; await assert.rejects(owner.start("codex", "a", []), /identity changed|became occupied/u); assert.equal(starts, 1); } }); it("deadline rejection leaves ownership unset and cleanup closes the unoccupied pane", async () => { let starts = 0, closes = 0, now = 0, removals = 0; const connection = { client: { async call(method: string) { if (method === "agent.start") { starts++; throw new Error("is not an available shell"); } if (method === "pane.get") return { type: "pane_info", pane: { pane_id: "pane", workspace_id: "w", tab_id: "t", agent: null, agent_session: null } }; if (method === "session.snapshot") return { type: "session_snapshot", snapshot: { agents: [] } }; if (method === "pane.close") { closes++; return {}; } throw new Error(method); } }, async close() {} }, owner = new (HerdrPlacedRunOwner as unknown as new (...args: unknown[]) => HerdrPlacedRunOwner)({ id: "m", target: "host", cwd: "/repo" }, "run", "/tmp/pi-subagents-herdr-owned", connection, () => { removals++; }); owner.owned = { workspaceId: "w", tabId: "t", paneId: "pane" }; await assert.rejects(owner.start("codex", "a", [], { timeoutMs: 20, pollMs: 10, clock: { now: () => now, async wait(ms) { now += ms; } } }), /startup deadline/u); assert.equal(starts, 3); assert.equal(owner.agentName, undefined); await owner.cleanup(true); assert.equal(closes, 1); assert.equal(removals, 1); }); it("resolves agent instructions, skills, memory, and refinements from a distinct remote fixture", () => { const remote = fs.mkdtempSync(path.join(os.tmpdir(), "pi-herdr-remote-resources-")); fs.mkdirSync(path.join(remote, ".git")); fs.mkdirSync(path.join(remote, ".selesai", "agents"), { recursive: true }); fs.writeFileSync(path.join(remote, ".selesai", "agents", "remote-worker.md"), "---\nname: remote-worker\ndescription: Remote\ntools: read\nskills:\n - remote-skill\nmemory: { scope: project, path: remote-worker }\ninheritProjectContext: true\ninheritGlobalContext: false\ninheritSkills: true\n---\nREMOTE_AGENT_INSTRUCTIONS\n"); fs.mkdirSync(path.join(remote, ".selesai", "skills", "remote-skill"), { recursive: true }); fs.writeFileSync(path.join(remote, ".selesai", "skills", "remote-skill", "SKILL.md"), "---\ndescription: REMOTE_SKILL_DESCRIPTION\n---\nREMOTE_SKILL_CONTENT\n"); fs.mkdirSync(path.join(remote, ".selesai", "agent-memory", "remote-worker"), { recursive: true }); fs.writeFileSync(path.join(remote, ".selesai", "agent-memory", "remote-worker", "MEMORY.md"), "REMOTE_MEMORY_CONTENT\n"); const refinement = path.join(remote, ".pi-subagents", "refinements", "remote-worker.md"); fs.mkdirSync(path.dirname(refinement), { recursive: true }); fs.writeFileSync(refinement, `\n\n\`\`\`pi-subagents-refinement-current\nREMOTE_REFINEMENT_CONTENT\n\`\`\`\n\n\`\`\`pi-subagents-refinement-snapshots-json\n[]\n\`\`\`\n`); try { const resolved = resolveRemoteHerdrResources(remote, { agent: "remote-worker" }); assert.match(resolved.systemPrompt, /REMOTE_AGENT_INSTRUCTIONS/u); assert.match(resolved.systemPrompt, /REMOTE_SKILL_DESCRIPTION/u); assert.match(resolved.systemPrompt, /REMOTE_MEMORY_CONTENT/u); assert.match(resolved.systemPrompt, /REMOTE_REFINEMENT_CONTENT/u); assert.doesNotMatch(resolved.systemPrompt, /LOCAL_/u); assert.throws(() => resolveRemoteHerdrResources(remote, { agent: "remote-worker", skills: ["missing"] }), /Remote Pi skills not found/u); } finally { fs.rmSync(remote, { recursive: true, force: true }); } }); it("uses remote defaults for omitted tools and intersects wider remote tools with the caller ceiling", () => { const remote = fs.mkdtempSync(path.join(os.tmpdir(), "pi-herdr-tools-")); fs.mkdirSync(path.join(remote, ".git")); fs.mkdirSync(path.join(remote, ".selesai", "agents"), { recursive: true }); fs.writeFileSync(path.join(remote, ".selesai", "agents", "defaults.md"), "---\nname: defaults\ndescription: defaults\ndefaultReads: remote-only.md\ninheritProjectContext: true\ninheritGlobalContext: true\ninheritSkills: true\n---\nDefaults\n"); fs.writeFileSync(path.join(remote, "remote-only.md"), "remote context\n"); fs.writeFileSync(path.join(remote, ".selesai", "agents", "wide.md"), "---\nname: wide\ndescription: wide\ntools: read, bash\ninheritProjectContext: true\ninheritGlobalContext: true\ninheritSkills: true\n---\nWide\n"); try { const defaults = resolveRemoteHerdrResources(remote, { agent: "defaults", toolCeiling: ["read"] }, ["read", "bash"]); assert.deepEqual(defaults.tools, ["read"]); assert.ok(defaults.systemPrompt.includes(`[Read from: ${path.join(remote, "remote-only.md")}]`)); assert.doesNotMatch(resolveRemoteHerdrResources(remote, { agent: "defaults", reads: false }, ["read"]).systemPrompt, /remote-only\.md/u); assert.deepEqual(resolveRemoteHerdrResources(remote, { agent: "wide", toolCeiling: ["read"] }, ["read", "bash"]).tools, ["read"]); } finally { fs.rmSync(remote, { recursive: true, force: true }); } }); it("preserves canonical contact_supervisor progress and structured-interview semantics", { skip: process.platform === "win32" }, async () => { const dir = fs.mkdtempSync(path.join(os.tmpdir(), "pi-herdr-supervisor-")); const old = { mode: process.env[HERDR_PI_MODE_ENV], run: process.env[HERDR_PI_RUN_ENV], dir: process.env[HERDR_PI_RUNTIME_DIR_ENV] }; process.env[HERDR_PI_MODE_ENV] = "1"; process.env[HERDR_PI_RUN_ENV] = "run_supervisor1"; process.env[HERDR_PI_RUNTIME_DIR_ENV] = dir; const handlers = new Map(); let tool: any; let command: any; let socket: net.Socket | undefined; let activeTools = ["read"]; const pi = { registerTool(value: unknown) { tool = value; }, registerCommand(_name: string, value: unknown) { command = value; }, on(name: string, handler: Function) { handlers.set(name, handler); }, setActiveTools(value: string[]) { activeTools = value; } }; try { registerHerdrPiBridge(pi as never); const ctx = { sessionManager: { getSessionId: () => "native", getSessionFile: () => "/remote/session" }, getActiveTools: () => activeTools, model: { provider: "p", id: "m" }, isIdle: () => true, hasPendingMessages: () => false }; handlers.get("session_start")?.({}, ctx); socket = net.createConnection(path.join(dir, "bridge.sock")); const decoder = new HerdrPiFrameDecoder(); const frames: any[] = []; socket.on("data", (chunk) => frames.push(...decoder.push(chunk))); await new Promise((resolve) => socket!.once("connect", resolve)); socket.write(encodeHerdrPiFrame({ protocol: 1, runId: "run_supervisor1", type: "configure", requestId: "cfg", resources: { agent: "worker", toolCeiling: ["read"] } })); await new Promise((resolve) => setTimeout(resolve, 30)); assert.ok(frames.some((frame) => frame.type === "configured")); const progressPromise = tool.execute("tc1", { reason: "progress_update", message: "halfway" }); await new Promise((resolve) => setTimeout(resolve, 10)); const progressRequest = frames.find((frame) => frame.type === "supervisor-request" && frame.reason === "progress_update"); assert.equal(progressRequest?.expectsReply, false); socket.write(encodeHerdrPiFrame({ protocol: 1, runId: "run_supervisor1", type: "supervisor-delivered", requestId: progressRequest.requestId })); const progress = await progressPromise; assert.equal(progress.details.delivered, true); const interviewPromise = tool.execute("tc2", { reason: "interview_request", interview: { title: "Choose", questions: [{ id: "x" }] } }); await new Promise((resolve) => setTimeout(resolve, 10)); const request = frames.find((frame) => frame.type === "supervisor-request" && frame.reason === "interview_request"); assert.deepEqual(request.interview, { title: "Choose", questions: [{ id: "x" }] }); socket.write(encodeHerdrPiFrame({ protocol: 1, runId: "run_supervisor1", type: "prepare", requestId: "reply_1", operation: "supervisor-reply", supervisorId: request.requestId, text: '{"x":"yes"}' })); await new Promise((resolve) => setTimeout(resolve, 10)); await command.handler("reply_1", ctx); const interview = await interviewPromise; assert.deepEqual(interview.details.structuredReply, { x: "yes" }); socket.destroy(); handlers.get("session_shutdown")?.(); } finally { socket?.destroy(); handlers.get("session_shutdown")?.(); old.mode === undefined ? delete process.env[HERDR_PI_MODE_ENV] : process.env[HERDR_PI_MODE_ENV] = old.mode; old.run === undefined ? delete process.env[HERDR_PI_RUN_ENV] : process.env[HERDR_PI_RUN_ENV] = old.run; old.dir === undefined ? delete process.env[HERDR_PI_RUNTIME_DIR_ENV] : process.env[HERDR_PI_RUNTIME_DIR_ENV] = old.dir; fs.rmSync(dir, { recursive: true, force: true }); } }); it("frames correlated bridge data and fails closed on malformed, partial, or oversized evidence", () => { const decoder = new HerdrPiFrameDecoder(); const encoded = encodeHerdrPiFrame({ protocol: 1, runId: "run_12345678", type: "settled", requestId: "q1", idle: true, pending: false }); assert.deepEqual(decoder.push(encoded)[0], { protocol: 1, runId: "run_12345678", type: "settled", requestId: "q1", idle: true, pending: false }); assert.throws(() => new HerdrPiFrameDecoder().push(Buffer.from("not-json\n")), /malformed JSON/u); const partial = new HerdrPiFrameDecoder(); partial.push(Buffer.from('{"protocol":1')); assert.throws(() => partial.end(), /partial frame/u); assert.throws(() => new HerdrPiFrameDecoder().push(Buffer.alloc(HERDR_PI_MAX_FRAME_BYTES + 1)), /exceeds/u); assert.throws(() => new HerdrPiFrameDecoder().push(Buffer.concat([Buffer.from('{"protocol":1,"runId":"run_12345678","type":"x","v":"'), Buffer.from([0xc3, 0x28]), Buffer.from('"}\n')])), /UTF-8/u); }); it("reconnects by exact identities and evidence cursor without redispatch", async () => { const identity: HerdrRunIdentity = { runId: "run_12345678", machineId: "m1", target: "remote", session: null, workspaceId: "w1", tabId: "t1", paneId: "p1", terminalId: "term1", agentName: "agent", nativeSessionId: "native1", cwd: "/remote", runtimeDir: "/tmp/run" }; let prompts = 0; const order: string[] = []; const result = await boundedHerdrReconnect(identity, 4, async () => { order.push("subscribe", "snapshot"); return { endpoint: { session: null, protocol: 22, version: "0.9.0" }, agents: [{ terminal_id: "term1", pane_id: "p2", agent_status: "working" }], bridge: { runId: identity.runId, nativeSessionId: "native1", packageVersion, protocol: 1, evidenceFloor: 3, evidenceCursor: 6 } }; }); assert.equal(result.ok, true); assert.deepEqual(order, ["subscribe", "snapshot"]); assert.equal(prompts, 0); if (result.ok) assert.equal(result.paneId, "p2"); }); it("rejects reconnect candidates from a different package version", async () => { const identity: HerdrRunIdentity = { runId: "run_12345678", machineId: "m1", target: "remote", session: null, workspaceId: "w1", tabId: "t1", paneId: "p1", terminalId: "term1", agentName: "agent", nativeSessionId: "native1", cwd: "/remote", runtimeDir: "/tmp/run" }; const candidate = { endpoint: { session: null, protocol: 22, version: "0.9.0" }, agents: [{ terminal_id: "term1", pane_id: "p1", agent_status: "idle" }], bridge: { runId: identity.runId, nativeSessionId: identity.nativeSessionId, packageVersion: `${packageVersion}-different`, protocol: 1, evidenceFloor: 1, evidenceCursor: 5 } }; assert.deepEqual(await boundedHerdrReconnect(identity, 5, async () => candidate, { attempts: 1 }), { ok: false, reason: "Bridge run, package, protocol, or native session identity changed." }); }); it("keeps an accepted public prompt pending across transport replacement and replay without redispatch", { skip: process.platform === "win32" }, async () => { const dir = fs.mkdtempSync(path.join(os.tmpdir(), "pi-herdr-replay-")); const one = path.join(dir, "one.sock"); const two = path.join(dir, "two.sock"); let firstSocket: net.Socket | undefined; let requestId = ""; let dispatches = 0; const server1 = net.createServer((socket) => { firstSocket = socket; let text = ""; socket.on("data", (chunk) => { text += chunk; const line = text.split("\n").find(Boolean); if (!line) return; const frame = JSON.parse(line); requestId = frame.requestId; socket.write(encodeHerdrPiFrame({ protocol: 1, runId: "run_12345678", type: "prepared", requestId, nativeSessionId: "native1", cursor: 1 })); }); }); await new Promise((resolve) => server1.listen(one, resolve)); const server2 = net.createServer((socket) => { const send = () => { socket.write(encodeHerdrPiFrame({ protocol: 1, runId: "run_12345678", type: "accepted", requestId, nativeSessionId: "native1", cursor: 2 })); socket.write(encodeHerdrPiFrame({ protocol: 1, runId: "run_12345678", type: "settled", requestId, nativeSessionId: "native1", cursor: 3, idle: true, pending: false })); }; requestId ? send() : setTimeout(send, 20); }); await new Promise((resolve) => server2.listen(two, resolve)); const fakeConnection = () => ({ endpoint: { socket: "/tmp/h", session: null, version: "0.9.0", protocol: 1, compatible: true as const, running: true as const }, socketPath: "/tmp/h", client: { async call() { dispatches++; firstSocket?.write(encodeHerdrPiFrame({ protocol: 1, runId: "run_12345678", type: "accepted", requestId, nativeSessionId: "native1", cursor: 2 })); setTimeout(() => firstSocket?.destroy(), 5); return {}; }, async subscribe() { return () => {}; } }, async forwardRemoteSocket() { throw new Error("unused"); }, async close() {} }); const identity: HerdrRunIdentity = { runId: "run_12345678", machineId: "m1", target: "remote", session: null, workspaceId: "w1", tabId: "t1", paneId: "p1", terminalId: "term1", agentName: "agent", nativeSessionId: "native1", cwd: "/remote", runtimeDir: "/tmp/run" }; const bridge = new BridgeChannel(one, identity.runId); const connection = fakeConnection(); const owner = { identity, connection, snapshot: { connection: "connected", state: "ready", identity, revision: 0 }, async replaceConnection(input: { connection: ReturnType; paneId: string; state: "ready" | "working" | "blocked" | "settled" | "unknown" }) { this.connection = input.connection; this.identity = { ...this.identity, paneId: input.paneId }; this.snapshot = { connection: "connected", state: input.state, identity: this.identity, revision: this.snapshot.revision + 1 }; }, async cleanup() { await this.connection.close(); }, statusSubscriptions: () => [], observeEvent() {}, markConnectionUnknown: () => true } as unknown as HerdrPlacedRunOwner; const session = new HerdrPiSession(owner, { async close() {} }, bridge, ["read"], "p/m", runtime); session.armReconnect(async () => { const next = new BridgeChannel(two, identity.runId); await new Promise((resolve) => setTimeout(resolve, 40)); await session.replaceTransport({ connection: fakeConnection(), bridgeForward: { async close() {} }, bridge: next, unsubscribe() {}, paneId: "p1", cursor: 3, state: "settled" }); }); try { await session.prompt("work"); assert.equal(dispatches, 1); assert.equal(session.placementSnapshot.connection, "connected"); } finally { await session.dispose(); await Promise.all([new Promise((resolve) => server1.close(() => resolve())), new Promise((resolve) => server2.close(() => resolve()))]); fs.rmSync(dir, { recursive: true, force: true }); } }); it("keeps reconnect unknown for terminal/native mismatch, missing pane, restart evidence gap, and disposal", async () => { const identity: HerdrRunIdentity = { runId: "run_12345678", machineId: "m1", target: "remote", session: null, workspaceId: "w1", tabId: "t1", paneId: "p1", terminalId: "term1", agentName: "agent", nativeSessionId: "native1", cwd: "/remote", runtimeDir: "/tmp/run" }; const base = { endpoint: { session: null, protocol: 22, version: "0.9.0" }, agents: [{ terminal_id: "term1", pane_id: "p1", agent_status: "idle" }], bridge: { runId: identity.runId, nativeSessionId: "native1", packageVersion, protocol: 1, evidenceFloor: 1, evidenceCursor: 5 } }; assert.match((await boundedHerdrReconnect(identity, 5, async () => ({ ...base, bridge: { ...base.bridge, nativeSessionId: "other" } }), { attempts: 1 })).reason, /identity changed/u); assert.match((await boundedHerdrReconnect(identity, 5, async () => ({ ...base, agents: [] }), { attempts: 1 })).reason, /terminal is missing/u); assert.match((await boundedHerdrReconnect(identity, 5, async () => ({ ...base, bridge: { ...base.bridge, evidenceFloor: 9, evidenceCursor: 9 } }), { attempts: 1 })).reason, /evidence has a gap/u); let called = 0; assert.match((await boundedHerdrReconnect(identity, 5, async () => { called++; return base; }, { disposed: () => true })).reason, /disposed/u); assert.equal(called, 0); }); it("cleans common-owner success, startup failure, mismatch, and runtime removal exactly once", async () => { const make = (mode: "exact" | "empty" | "mismatch" | "occupied", removalFails = false) => { const calls: string[] = []; const connection = { endpoint: { socket: "/h", session: null, version: "0.9.0", protocol: 1, compatible: true as const, running: true as const }, socketPath: "/h", client: { async call(method: string) { calls.push(method); if (method === "agent.get") return { type: "agent_info", agent: { terminal_id: mode === "exact" ? "term" : "other", pane_id: mode === "exact" ? "pane" : "elsewhere" } }; if (method === "session.snapshot") return { type: "session_snapshot", snapshot: { agents: mode === "occupied" ? [{ terminal_id: "foreign", pane_id: "pane" }] : [] } }; return {}; }, async subscribe() { return () => {}; } }, async forwardRemoteSocket() { throw new Error("unused"); }, async close() { calls.push("close"); } }; const owner = new (HerdrPlacedRunOwner as unknown as new (...args: unknown[]) => HerdrPlacedRunOwner)({ id: "m", target: "host", cwd: "/repo" }, "run", "/tmp/pi-subagents-herdr-owned", connection, () => { calls.push("remove"); if (removalFails) throw new Error("remove failed"); }); owner.owned = { workspaceId: "w", tabId: "t", paneId: "pane" }; if (mode !== "empty") { owner.terminalId = "term"; owner.agentName = "agent"; } return { owner, calls }; }; for (const mode of ["exact", "empty", "mismatch"] as const) { const { owner, calls } = make(mode); await owner.cleanup(true); await owner.cleanup(true); assert.equal(calls.filter((value) => value === "pane.close").length, 1); assert.equal(calls.filter((value) => value === "remove").length, 1); } const uncertain = make("occupied"); await assert.rejects(uncertain.owner.cleanup(true), /uncertain/u); assert.equal(uncertain.calls.includes("pane.close"), false); assert.equal(uncertain.calls.filter((value) => value === "remove").length, 1); const removal = make("exact", true); const concurrent = await Promise.allSettled([removal.owner.cleanup(true), removal.owner.cleanup(true), removal.owner.cleanup(true)]); assert.equal(concurrent.every((result) => result.status === "rejected" && /remove failed/u.test(String(result.reason))), true); await assert.rejects(removal.owner.cleanup(true), /remove failed/u); assert.equal(removal.calls.filter((value) => value === "pane.close").length, 1); assert.equal(removal.calls.filter((value) => value === "remove").length, 1); }); });