Skip to content
Closed
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
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- Fixed `prime-agent status` terminating active IPython forkservers by excluding internal control sockets from daemon discovery and preserving the established forkserver connection.
22 changes: 21 additions & 1 deletion packages/coding-agent/src/cli/daemon-ps.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import { spawnSync } from "node:child_process";
import { existsSync, lstatSync, readdirSync, readFileSync, rmSync, unlinkSync } from "node:fs";
import { tmpdir } from "node:os";
import { basename, dirname, join, resolve } from "node:path";
import chalk from "chalk";
import { APP_NAME, getAgentDir, VERSION } from "../config.js";
Expand Down Expand Up @@ -99,7 +100,11 @@ export function parseSsListeners(stdout: string, appName: string): DiscoveredDae
if (!owner || !processNameMatches(owner[1]!, appName)) {
continue;
}
daemons.push({ pid: Number.parseInt(owner[2]!, 10), socketPath: normalizeSocketPath(socketPath) });
const normalized = normalizeSocketPath(socketPath);
if (isKernelForkServerSocketPath(normalized)) {
continue;
}
daemons.push({ pid: Number.parseInt(owner[2]!, 10), socketPath: normalized });
}
return daemons;
}
Expand All @@ -116,6 +121,9 @@ export function parseLsofListeners(stdout: string): DiscoveredDaemonProcess[] {
pid = Number.parseInt(value, 10);
} else if (field === "n" && pid !== undefined && value.startsWith("/")) {
const socketPath = normalizeSocketPath(value);
if (isKernelForkServerSocketPath(socketPath)) {
continue;
}
const key = `${pid}:${socketPath}`;
if (!seen.has(key)) {
seen.add(key);
Expand Down Expand Up @@ -881,6 +889,18 @@ export function isWorkerSocketPath(socketPath: string): boolean {
);
}

/** Internal kernel forkserver listeners are not daemon control sockets. */
export function isKernelForkServerSocketPath(socketPath: string): boolean {
if (process.platform === "win32") return false;
const normalized = normalizeSocketPath(socketPath);
const socketDirectory = dirname(normalized);
return (
basename(normalized) === "control.sock" &&
basename(socketDirectory).startsWith("prime-agent-forkserver-") &&
resolve(dirname(socketDirectory)) === resolve(tmpdir())
);
}

async function stopBackgroundService(
socketPath: string,
pid: number | undefined,
Expand Down
6 changes: 6 additions & 0 deletions packages/coding-agent/src/core/kernel/fork-server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -178,6 +178,12 @@ export class ForkServer {
};

server.on("connection", (socket) => {
// The Python forkserver owns the first connection for its lifetime. A
// discovery probe must not replace that control channel when it connects.
if (this.conn) {
socket.destroy();
return;
}
this.conn = socket;
socket.setEncoding("utf8");
socket.on("data", (chunk: string) => this.onData(chunk));
Expand Down
24 changes: 23 additions & 1 deletion packages/coding-agent/test/daemon-ps.test.ts
Original file line number Diff line number Diff line change
@@ -1,8 +1,10 @@
import { tmpdir } from "node:os";
import { join } from "node:path";
import { describe, expect, it } from "vitest";
import {
type DaemonInfo,
evaluateShutdownQuietPeriod,
isKernelForkServerSocketPath,
isWorkerSocketPath,
mergeDiscoveredDaemonProcesses,
parseLsofListeners,
Expand All @@ -26,11 +28,20 @@ describe("worker socket classification", () => {
});
});

describe("forkserver socket classification", () => {
it.runIf(process.platform !== "win32")("recognizes only internal forkserver control sockets", () => {
expect(isKernelForkServerSocketPath(join(tmpdir(), "prime-agent-forkserver-abc123", "control.sock"))).toBe(true);
expect(isKernelForkServerSocketPath(join(tmpdir(), "prime-agent-forkserver-abc123", "daemon.sock"))).toBe(false);
expect(isKernelForkServerSocketPath(join(tmpdir(), "custom", "control.sock"))).toBe(false);
});
});

describe("parseSsListeners", () => {
const stdout = [
"Netid State Recv-Q Send-Q Local Address:Port Peer Address:Port",
'u_str LISTEN 0 511 /tmp/custom.sock 10147608 * 0 users:(("prime-agent",pid=1234,fd=22))',
'u_str LISTEN 0 511 /tmp/prime-agent-1000/daemon.sock 79453846 * 0 users:(("prime-agent",pid=5678,fd=24))',
'u_str LISTEN 0 511 /tmp/prime-agent-forkserver-probe/control.sock 79453847 * 0 users:(("prime-agent",pid=2468,fd=25))',
'u_str LISTEN 0 4096 /run/dbus/system_bus_socket 123 * 0 users:(("dbus-daemon",pid=900,fd=3))',
'u_str ESTAB 0 0 /tmp/other.sock 456 * 0 users:(("prime-agent",pid=4321,fd=9))',
"",
Expand All @@ -48,6 +59,7 @@ describe("parseSsListeners", () => {
const daemons = parseSsListeners(stdout, "prime-agent");
expect(daemons.some((daemon) => daemon.socketPath.includes("dbus"))).toBe(false);
expect(daemons.some((daemon) => daemon.pid === 4321)).toBe(false);
expect(daemons.some((daemon) => daemon.pid === 2468)).toBe(false);
});

it("honors a different app name", () => {
Expand All @@ -57,7 +69,17 @@ describe("parseSsListeners", () => {

describe("parseLsofListeners", () => {
it("pairs each pid with its listening unix socket paths", () => {
const stdout = ["p1234", "fu", "n/tmp/a.sock", "p5678", "n/tmp/b.sock", "n0x0 (not a path)", ""].join("\n");
const stdout = [
"p1234",
"fu",
"n/tmp/a.sock",
"p2468",
"n/tmp/prime-agent-forkserver-probe/control.sock",
"p5678",
"n/tmp/b.sock",
"n0x0 (not a path)",
"",
].join("\n");
expect(parseLsofListeners(stdout)).toEqual([
{ pid: 1234, socketPath: "/tmp/a.sock" },
{ pid: 5678, socketPath: "/tmp/b.sock" },
Expand Down
18 changes: 18 additions & 0 deletions packages/coding-agent/test/kernel-fork-server-protocol.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import { spawnSync } from "node:child_process";
import { existsSync, mkdirSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from "node:fs";
import { createConnection, type Server } from "node:net";
import { homedir, tmpdir } from "node:os";
import { join } from "node:path";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
Expand Down Expand Up @@ -231,6 +232,23 @@ describeIf("forkserver kill/liveness protocol (stub python)", () => {
await expect(handle.isAlive()).rejects.toBeInstanceOf(ForkServerUnavailable);
}, 15_000);

it("keeps the established control channel when another client probes the socket", async () => {
const handle = await spawnStubKernel();
const listener = (server as unknown as { server?: Server }).server;
const address = listener?.address();
expect(typeof address).toBe("string");

await new Promise<void>((resolve, reject) => {
const probe = createConnection(address as string);
probe.once("connect", () => probe.end());
probe.once("close", () => resolve());
probe.once("error", reject);
});

expect(server!.isDead).toBe(false);
expect(await handle.isAlive()).toBe(true);
}, 15_000);

describe("forkserver orphan journal", () => {
let journalPath = "";
const savedJournal = process.env[ORPHAN_PROCESS_JOURNAL_ENV];
Expand Down