Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
66 changes: 40 additions & 26 deletions plugins/pstack/pi/agents.ts
Original file line number Diff line number Diff line change
Expand Up @@ -78,7 +78,10 @@ interface Run {

// A remote agent runs under the live pi process that restore left it to.
type LocalAgent = { kind: "local"; record: RunningRecord; run: Run };
type AgentState = LocalAgent | { kind: "remote"; record: RunningRecord } | { kind: "ended"; record: EndedRecord };
type AgentState = LocalAgent
| { kind: "remote"; record: RunningRecord }
| { kind: "interrupted"; record: RunningRecord }
| { kind: "ended"; record: EndedRecord };

// How a run ended, from what its process left behind. A run that settled on a
// reply completed, whatever the exit code of the shutdown after it. Stderr
Expand Down Expand Up @@ -296,10 +299,10 @@ export class AgentRunner {

private find(to: string): AgentState {
const byId = this.agents.get(to);
if (byId) return byId;
if (byId) return this.refreshRestored(byId);
const states = [...this.agents.values()];
const matches = states.filter((state) => state.record.agent.description === to);
if (matches.length === 1) return matches[0];
if (matches.length === 1) return this.refreshRestored(matches[0]);
if (matches.length > 1) {
throw new Error(`"${to}" matches several agents (${matches.map((state) => state.record.agent.id).join(", ")}); pass an agentId.`);
}
Expand All @@ -309,11 +312,14 @@ export class AgentRunner {

// An agent this process can act on. One running under another pi process is
// not: only that process holds its stdin and can message or stop it.
private owned(to: string): Exclude<AgentState, { kind: "remote" }> {
private owned(to: string): Exclude<AgentState, { kind: "remote" | "interrupted" }> {
const state = this.find(to);
if (state.kind === "remote") {
throw new Error(`Agent ${state.record.agent.id} is running under another pi process (pid ${state.record.parentPid}); only that process can message or stop it.`);
}
if (state.kind === "interrupted") {
throw new Error(`Agent ${state.record.agent.id}'s previous process (pid ${state.record.pid}) has not exited; it cannot resume or be stopped here until its exit is confirmed.`);
}
return state;
}

Expand Down Expand Up @@ -360,40 +366,48 @@ export class AgentRunner {
}

list(): AgentRecord[] {
return [...this.agents.values()].map((state) => state.record);
return [...this.agents.values()].map((state) => this.refreshRestored(state).record);
}

private persist(record: AgentRecord): void {
this.pi.appendEntry(ENTRY_TYPE, record);
}

// Folds persisted snapshots. A snapshot still marked running is left to the
// live pi process that launched it while its process is still that one's
// child: a live parent pid alone may be a reused one. Any other running
// snapshot ends here, and its process is stopped when it is known to be the
// agent's.
private refreshRestored(state: AgentState): AgentState {
if (state.kind === "local" || state.kind === "ended") return state;
const { record } = state;
const fate = fateOf(record);
if (fate === "kept" && state.kind === "remote") return state;
if (fate !== "gone") {
if (state.kind === "interrupted") return state;
const interrupted: AgentState = { kind: "interrupted", record };
this.agents.set(record.agent.id, interrupted);
reapOrphan(record, this.settings.killGraceMs);
return interrupted;
}
const stopped: EndedRecord = {
agent: record.agent,
status: "stopped",
pid: record.pid,
exitCode: null,
endedAt: now(),
finalText: "(interrupted: the session that started it ended)",
};
const ended: AgentState = { kind: "ended", record: stopped };
this.agents.set(record.agent.id, ended);
this.persist(stopped);
return ended;
}

// A restored process stays unavailable until its exit is confirmed. Only a
// process whose identity matches the record can be signalled.
restore(entries: readonly SessionEntry[]): void {
for (const entry of entries) {
if (entry.type !== "custom" || entry.customType !== ENTRY_TYPE || !Value.Check(recordSchema, entry.data)) continue;
const { id } = entry.data.agent;
if (this.agents.get(id)?.kind === "local") continue;
this.agents.set(id, entry.data.status === "running" ? { kind: "remote", record: entry.data } : { kind: "ended", record: entry.data });
}
for (const [id, state] of this.agents) {
if (state.kind !== "remote") continue;
const { record } = state;
if (fateOf(record) === "kept") continue;
const stopped: EndedRecord = {
agent: record.agent,
status: "stopped",
pid: record.pid,
exitCode: null,
endedAt: now(),
finalText: "(interrupted: the session that started it ended)",
};
this.agents.set(id, { kind: "ended", record: stopped });
this.persist(stopped);
reapOrphan(record, this.settings.killGraceMs);
}
for (const state of this.agents.values()) this.refreshRestored(state);
}
}
2 changes: 1 addition & 1 deletion plugins/pstack/skills/poteto-mode/references/pi-tools.md
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ poteto-mode's Subagents section applies on Pi through the `agent` tool:
- `run_in_background: true` returns the agent id at once. The completion joins the conversation after your current tool calls finish, as on Claude Code, or starts a turn when you are idle. In print and JSON mode (`pi -p`), where Pi exits once the run settles, the extension holds the settle while a background agent runs, so each completion still arrives as a turn and the process ends after the last one. A child agent holds its settle the same way while its own background agents run, because its parent closes its stdin once it settles, and a message from the parent ends the hold. Interactive and RPC main sessions settle as usual and take the completion when it arrives.
- While a background agent runs, the extension blocks a `bash` command line in which the word `sleep` is followed by a literal duration of two seconds or more, anywhere in the line, poll loops included: `sleep 30`, `npm test; sleep 30`, `until gh pr checks 1; do sleep 30; done`, `sleep infinity`. A shorter sleep, a sleep backgrounded as `sleep 30 &`, and `sleep` inside a quoted string, a comment, or a heredoc all run. The rule errs toward blocking, so `timeout 5 sleep 30` and a sleep inside a backgrounded group are blocked too. The completion notice arrives on its own, so continue other work or end your turn.
- An agent's status follows its process. `completed` means the child exited, and `stop_agent` reports `stopped` only once the child has exited, so the Claude Code caveat about a `completed` agent that keeps running does not apply to the agent itself. Stopping an agent, like the session shutdown that stops every agent, ends the child pi and the `bash` command it still has running, because Pi kills that command's process tree on SIGTERM. A process that a finished `bash` command left in the background (`cmd &`, `nohup cmd`) runs in its own process group and outlives the agent; stop it yourself. A stopped agent's notice joins the conversation without starting a turn, since `stop_agent` already returned.
- Agents belong to the session that started them. Quitting, reloading, and starting, resuming, or forking a session all stop every running agent, and a parent process that exits signals its agents to stop. When a session starts, an agent still recorded as running is marked `stopped`, and its process is killed, only if the `pi` process that launched it is gone. An agent whose launching process is still alive is left running.
- Agents belong to the session that started them. Quitting, reloading, and starting, resuming, or forking a session all stop every running agent, and a parent process that exits signals its agents to stop. On restore, an agent still running under its launching Pi process is left alone. An orphan whose process identity matches the saved record is signalled to stop, and remains unavailable to `send_message` and `stop_agent` until its process exits. Without enough identity data to signal it safely, its process is left alone and continuation is refused until its exit is confirmed. Only then is the restored agent marked `stopped` and allowed to resume.
- `send_message` to a running agent is an RPC `steer` on the child's stdin. The agent reads it after its current tool calls, as on Claude Code, carries on in the same run, and sends one completion notice. A message the child rejects returns an error, and the agent keeps running. A message to a finished agent resumes its session in the background with the context of its earlier runs. A message that arrives after the agent settled but before its process exited resumes the agent once the process has exited.
- A role value's `@<level>` picks the same effort agent as on Claude Code, and the extension passes its level to the child as `--thinking`. `session`, or an agent with no effort, runs the child at the parent's current thinking level.
- Keep the rest of the policy unchanged. Pass file pointers not inlined context, give each worker its own worktree when they write, review every subagent's diff yourself.
Expand Down
12 changes: 9 additions & 3 deletions plugins/pstack/skills/poteto-mode/scripts/worktree-audit.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@
// Every probe yields a Fact, { known: true, value } or { known: false }. A hold
// bucket needs only its own fact; `safe` needs every fact known.
import { execFileSync } from "node:child_process";
import { existsSync, readdirSync, readFileSync, realpathSync, statSync } from "node:fs";
import { readdirSync, readFileSync, realpathSync, statSync } from "node:fs";
import { homedir } from "node:os";
import { basename, dirname, join } from "node:path";
import process from "node:process";
Expand Down Expand Up @@ -74,13 +74,19 @@ export function parseWorktrees(output) {
return worktrees;
}

export function defaultTranscriptRoots({ env = process.env, home = homedir(), exists = existsSync } = {}) {
export function defaultTranscriptRoots({ env = process.env, home = homedir(), stat = statSync } = {}) {
const claude = join(env.CLAUDE_CONFIG_DIR || join(home, ".claude"), "projects");
const codex = env.CODEX_HOME || join(home, ".codex");
const piAgent = env.PI_CODING_AGENT_DIR || join(home, ".pi", "agent");
const copilot = join(env.COPILOT_HOME || join(home, ".copilot"), "session-state");
const found = [claude, join(codex, "sessions"), join(codex, "archived_sessions"),
join(piAgent, "sessions"), join(piAgent, "pstack"), copilot].filter((root) => exists(root));
join(piAgent, "sessions"), join(piAgent, "pstack"), copilot].filter((root) => {
try {
return stat(root, { throwIfNoEntry: false }) !== undefined;
} catch (error) {
return !["ENOENT", "ENOTDIR"].includes(error.code);
}
});
return found.length ? found : [claude];
}

Expand Down
59 changes: 57 additions & 2 deletions tests/pi/agents-session.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -89,8 +89,8 @@ describe("registry", () => {

// The other pi: a second process that starts a background agent and prints
// the entries a pi process opening the same session would read.
async function otherPi() {
const { w, ctx } = setup({ script: { default: [{ spawn: "running" }, { sleep: 30000 }] } });
async function otherPi(options = {}) {
const { w, ctx } = setup({ script: { default: [{ spawn: "running" }, { sleep: 30000 }] }, ...options });
const host = w.spawn(process.execPath, [join(import.meta.dir, "host.mjs"), JSON.stringify(w.settings), "stay"], {
stdio: ["ignore", "pipe", "inherit"],
});
Expand All @@ -101,6 +101,61 @@ describe("registry", () => {
return { w, ctx, host, entries: JSON.parse(printed) };
}

test("an orphan cannot resume its session while its process awaits SIGKILL", async () => {
const { w, ctx, host, entries } = await otherPi({ script: { byPrompt: {
x: [{ ignoreSigterm: true }, { spawn: "running" }, { sleep: 30000 }],
again: [{ reply: "resumed" }],
} } });
const [first] = w.invocations();
const exited = exitOf(host);
host.kill("SIGKILL");
await exited;
const resumed = await restore(w, entries);
const id = entries[0].data.agent.id;

expect(alive(first.pid)).toBe(true);
expect((await listAgents(resumed, ctx))[0].status).toBe("running");
await expect(resumed.call("send_message", { to: id, message: "again" }, ctx)).rejects.toThrow(/process.*has not exited/);
expect(w.invocations()).toHaveLength(1);
await waitFor(() => !alive(first.pid));
expect((await listAgents(resumed, ctx))[0].status).toBe("stopped");

await resumed.call("send_message", { to: id, message: "again" }, ctx);
await waitFor(() => resumed.messages.length === 1);
const second = w.invocations()[1];
expect(flag(second, "--session-id")).toBe(flag(first, "--session-id"));
expect(second.cwd).toBe(first.cwd);
expect(alive(first.pid)).toBe(false);
});

test("an orphan without recorded process identity cannot resume or be signalled until it exits", async () => {
const { w, ctx, host, entries } = await otherPi();
const [first] = w.invocations();
delete entries[0].data.pidStart;
const exited = exitOf(host);
host.kill("SIGKILL");
await exited;
const resumed = await restore(w, entries);
const id = entries[0].data.agent.id;

expect((await listAgents(resumed, ctx))[0].status).toBe("running");
for (const [tool, params] of [
["send_message", { to: id, message: "again" }],
["stop_agent", { id }],
]) await expect(resumed.call(tool, params, ctx)).rejects.toThrow(/process.*has not exited/);
await sleep(2 * w.settings.killGraceMs);
expect(alive(first.pid)).toBe(true);
expect(w.invocations()).toHaveLength(1);

process.kill(-first.pid, "SIGTERM");
await waitFor(() => !alive(first.pid));
expect((await listAgents(resumed, ctx))[0].status).toBe("stopped");
await resumed.call("send_message", { to: id, message: "again" }, ctx);
await w.until("invocation", 2);
expect(flag(w.invocations()[1], "--session-id")).toBe(flag(first, "--session-id"));
await resumed.emit("session_shutdown", {}, ctx);
});

for (const [what, edit, env, status] of [
["is left running, untouched", (data) => data, {}, "running"],
["is left running on a record an earlier version wrote, which has no start time", ({ pidStart, ...data }) => data, {}, "running"],
Expand Down
38 changes: 36 additions & 2 deletions tests/worktree-audit.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -335,6 +335,24 @@ for (const dir of ["sessions", "archived_sessions"]) {
});
}

test.skipIf(noChmod)("an inaccessible default transcript root keeps a recently used worktree out of safe", () => {
const fixture = createFixture();
const chatted = addWorktree(fixture, "codex-inaccessible");
mkdirSync(join(fixture.root, ".claude", "projects"), { recursive: true });
const codex = join(fixture.root, ".codex");
const session = join(codex, "sessions", "recent.jsonl");
mkdirSync(dirname(session), { recursive: true });
writeFileSync(session, `${JSON.stringify({ type: "session_meta", payload: { cwd: chatted } })}\n`);
chmodSync(codex, 0o000);
locked.push(codex);

const transcripts = defaultTranscriptRoots({ env: {}, home: fixture.root });
const { rows, warnings } = runAudit(fixture, { transcripts });
expect(rowFor(rows, chatted).slice(6, 8)).toEqual(["-", "review"]);
expect(warnings).toHaveLength(1);
expect(warnings[0]).toMatch(/transcript scan failed.*EACCES/);
});

describe("lastChats matches a path as JSONL spells it, never a sibling's prefix", () => {
const scan = (path, cwd, spellings = [path]) => {
const root = realpathSync(mkdtempSync(join(tmpdir(), "worktree-audit-chats-")));
Expand Down Expand Up @@ -550,7 +568,7 @@ describe("default transcripts roots", () => {
const claude = "/home/u/.claude/projects";
// These paths are POSIX literals, which path.join spells with backslashes on Windows.
const roots = ({ exists, ...options }) =>
defaultTranscriptRoots({ ...options, home, exists: (path) => exists(gitPath(path)) }).map(gitPath);
defaultTranscriptRoots({ ...options, home, stat: (path) => exists(gitPath(path)) ? {} : undefined }).map(gitPath);

test("every runtime directory that exists, Pi's under PI_CODING_AGENT_DIR when set", () => {
const present = new Set([claude, "/home/u/.codex/sessions", "/pi/sessions", "/pi/pstack", "/home/u/.pi/agent/sessions"]);
Expand Down Expand Up @@ -587,12 +605,28 @@ describe("default transcripts roots", () => {
expect(roots({ env: {}, exists: () => false })).toEqual([claude]);
});

test.each(["EACCES", "EPERM", "EIO", "ELOOP"])("a %s while discovering a root retains it for the audit", (code) => {
const inaccessible = join(home, ".codex", "sessions");
const found = defaultTranscriptRoots({ env: {}, home, stat: (path) => {
if (path === inaccessible) throw Object.assign(new Error("cannot inspect root"), { code });
return gitPath(path) === claude ? {} : undefined;
} });
expect(found.map(gitPath)).toEqual([claude, gitPath(inaccessible)]);
});

test.each(["ENOENT", "ENOTDIR"])("a %s while discovering a root omits it", (code) => {
expect(defaultTranscriptRoots({ env: {}, home, stat: (path) => {
if (gitPath(path) === claude) return {};
throw Object.assign(new Error("root is absent"), { code });
} }).map(gitPath)).toEqual([claude]);
});

test.each([
["PI_CODING_AGENT_DIR", { PI_CODING_AGENT_DIR: "/pi" }],
["the default agent directory", {}],
])("the Pi roots are where the extension keeps sessions and agent state under %s", (_, env) => {
const { agentDir } = defaultSettings(() => 0, env);
const roots = defaultTranscriptRoots({ env, home: homedir(), exists: () => true });
const roots = defaultTranscriptRoots({ env, home: homedir(), stat: () => ({}) });
expect(roots).toEqual(expect.arrayContaining([join(agentDir, "sessions"), join(agentDir, PSTACK_STATE_DIR)]));
});
});
Expand Down
Loading