fix: use bounded runtime paths for Unix sockets

This commit is contained in:
Christian Klotz
2026-08-14 17:02:03 +03:00
parent bf68c783a9
commit 9fcaac9b1a
14 changed files with 126 additions and 77 deletions
+5 -4
View File
@@ -63,15 +63,16 @@ const client = new PiClient({
await client.connect();
```
Unix discovery scans `~/.pi/server/*.sock`, derives each expected service ID from its filename, and verifies it through the existing handshake:
Unix discovery scans an explicit physical-route directory, derives each expected service ID from its filename, and verifies it through the existing handshake:
```ts
import { discoverUnixServices } from "@earendil-works/pi-client/unix";
const routes = await discoverUnixServices();
// [{ serviceId: "...", path: "/home/me/.pi/server/<serviceId>.sock" }]
const routes = await discoverUnixServices({ directory: "/run/user/1000/pi" });
// [{ serviceId: "...", path: "/run/user/1000/pi/<serviceId>.sock" }]
```
Malformed entries, non-sockets, stale or unresponsive endpoints, and service-ID mismatches are ignored. Discovery is read-only and probes at most 16 sockets concurrently. Unexpected filesystem and socket errors reject discovery. Pass `directory` or `timeoutMs` to override the defaults.
Malformed entries, non-sockets, stale or unresponsive endpoints, and service-ID mismatches are ignored. Discovery is read-only and probes at most 16 sockets concurrently. Unexpected filesystem and socket errors reject discovery. The caller must choose a short, private directory because Unix socket path limits are substantially lower than normal filesystem path limits.
Pass `timeoutMs` to override the default probe timeout.
`PiClientOptions.maxFrameLength` bounds protocol payloads. `maxPendingBytes` bounds queued Unix transport output. Configure matching limits on both peers.
+4 -5
View File
@@ -1,6 +1,5 @@
import { lstat, readdir } from "node:fs/promises";
import { createConnection, type Socket } from "node:net";
import { homedir } from "node:os";
import { join } from "node:path";
import {
DEFAULT_MAX_FRAME_LENGTH,
@@ -29,16 +28,16 @@ export interface UnixServiceRoute {
}
export interface DiscoverUnixServicesOptions {
/** Defaults to ~/.pi/server. */
directory?: string;
/** Directory containing service-addressed Unix sockets. */
directory: string;
/** Maximum time for each connection and handshake. Defaults to 1,000 ms. */
timeoutMs?: number;
}
/** Discover reachable local Pi services by probing service-addressed Unix sockets. */
export async function discoverUnixServices(options: DiscoverUnixServicesOptions = {}): Promise<UnixServiceRoute[]> {
export async function discoverUnixServices(options: DiscoverUnixServicesOptions): Promise<UnixServiceRoute[]> {
if (process.platform === "win32") throw new Error("Unix transport is not supported on Windows");
const directory = options.directory ?? join(homedir(), ".pi", "server");
const directory = options.directory;
const timeoutMs = options.timeoutMs ?? DEFAULT_DISCOVERY_TIMEOUT_MS;
if (!Number.isSafeInteger(timeoutMs) || timeoutMs <= 0 || timeoutMs > MAX_TIMER_DELAY_MS) {
throw new TypeError(`Unix discovery timeoutMs must be an integer between 1 and ${MAX_TIMER_DELAY_MS}`);
+1 -2
View File
@@ -1,6 +1,5 @@
import { mkdtemp, rm } from "node:fs/promises";
import { createServer, type Server, type Socket } from "node:net";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { ClientMessageDecoder, encodeServerMessage, PROTOCOL_VERSION } from "@earendil-works/pi-protocol";
import { afterEach, describe, expect, test } from "vitest";
@@ -13,7 +12,7 @@ const servers = new Set<Server>();
const sockets = new Set<Socket>();
async function makeSocketPath(): Promise<string> {
const directory = await mkdtemp(join(tmpdir(), "pi-client-transport-"));
const directory = await mkdtemp(join("/tmp", "pi-client-transport-"));
tempDirectories.add(directory);
return join(directory, "pi.sock");
}
+1 -2
View File
@@ -2,7 +2,6 @@ import { type ChildProcess, fork } from "node:child_process";
import { once } from "node:events";
import { lstat, mkdir, mkdtemp, rm, writeFile } from "node:fs/promises";
import { createServer, type Server, type Socket } from "node:net";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { afterEach, describe, expect, test } from "vitest";
import { PiServer } from "../../server/src/server.ts";
@@ -16,7 +15,7 @@ const rawSockets = new Set<Socket>();
const children = new Set<ChildProcess>();
async function makeDirectory(): Promise<string> {
const directory = await mkdtemp(join(tmpdir(), "pc-"));
const directory = await mkdtemp(join("/tmp", "pc-"));
tempDirectories.add(directory);
return directory;
}
@@ -1,4 +1,4 @@
import { lstat } from "node:fs/promises";
import { chmod, lstat, mkdir } from "node:fs/promises";
import { homedir } from "node:os";
import { basename, join } from "node:path";
import { setTimeout as delay } from "node:timers/promises";
@@ -20,6 +20,7 @@ import { startExperimentalSessionWorker } from "./session-worker.ts";
const SOCKET_RELEASE_TIMEOUT_MS = 10_000;
const SOCKET_RELEASE_POLL_MS = 10;
const EXPERIMENTAL_SOCKET_ROOT = "/tmp";
export interface ExperimentalMemoryServer {
readonly serviceId: string;
@@ -38,19 +39,21 @@ export type ExperimentalClientResult =
| { readonly kind: "attached"; readonly serviceId: string; readonly sessionId: string };
export interface StartExperimentalMemoryServerOptions {
/** Directory for service-addressed Unix sockets. Defaults to ~/.pi/server. */
/** Directory for service-addressed Unix sockets. Defaults to a short, private per-user runtime directory. */
readonly directory?: string;
readonly path?: string;
readonly serviceId?: string;
}
export interface StartExperimentalServerGenerationOptions {
/** One directory represents one experimental local service profile. */
/** Persistent profile directory. Defaults to ~/.pi/server. */
readonly directory?: string;
/** Physical socket directory. Defaults to a short, private per-user runtime directory. */
readonly socketDirectory?: string;
}
export interface RunExperimentalClientOptions {
/** Directory searched when --connect is omitted. Defaults to ~/.pi/server. */
/** Directory searched when --connect is omitted. Defaults to the experimental per-user runtime directory. */
readonly directory?: string;
}
@@ -95,7 +98,12 @@ export async function startExperimentalMemoryServer(
};
},
};
const socketPath = options.path ?? getUnixSocketPath(serviceId, options.directory);
let socketPath = options.path;
if (socketPath === undefined) {
const socketDirectory = options.directory ?? getExperimentalSocketDirectory();
await ensurePrivateSocketDirectory(socketDirectory);
socketPath = getUnixSocketPath(serviceId, socketDirectory);
}
const server = createUnixServer(host, { serviceId, path: socketPath, mode: 0o600 });
try {
await server.start();
@@ -142,10 +150,12 @@ export async function startExperimentalServerGeneration(
options: StartExperimentalServerGenerationOptions = {},
): Promise<ExperimentalMemoryServer> {
const directory = options.directory ?? join(homedir(), ".pi", "server");
const socketDirectory = options.socketDirectory ?? getExperimentalSocketDirectory();
const { serviceId, release } = await acquireExperimentalServiceProfile(directory);
let runtime: ExperimentalMemoryServer;
try {
const socketPath = getUnixSocketPath(serviceId, directory);
await ensurePrivateSocketDirectory(socketDirectory);
const socketPath = getUnixSocketPath(serviceId, socketDirectory);
await drainExistingGeneration(serviceId, socketPath);
runtime = await startExperimentalMemoryServer({
directory,
@@ -195,7 +205,9 @@ export async function runExperimentalClient(
options: RunExperimentalClientOptions = {},
): Promise<ExperimentalClientResult> {
if (command.auth !== undefined) throw new Error("Authentication is not supported by the local demo server");
const routes = command.connect ? [routeFromExplicitPath(command.connect.path)] : await discoverUnixServices(options);
const routes = command.connect
? [routeFromExplicitPath(command.connect.path)]
: await discoverUnixServices({ directory: options.directory ?? getExperimentalSocketDirectory() });
const discovered: { route: UnixServiceRoute; sessionIds: string[] }[] = [];
for (const route of routes) {
@@ -265,6 +277,23 @@ async function waitForSocketRelease(path: string): Promise<void> {
}
}
function getExperimentalSocketDirectory(): string {
if (process.platform === "win32" || typeof process.getuid !== "function") {
throw new Error("Experimental Unix server transport requires a POSIX user ID");
}
return join(EXPERIMENTAL_SOCKET_ROOT, `pi-server-${process.getuid()}`);
}
async function ensurePrivateSocketDirectory(directory: string): Promise<void> {
if (typeof process.getuid !== "function") throw new Error("Unix socket directory requires a POSIX user ID");
await mkdir(directory, { recursive: true, mode: 0o700 });
const stats = await lstat(directory);
if (!stats.isDirectory()) throw new Error(`Unix socket directory is not a directory: ${directory}`);
if (stats.uid !== process.getuid())
throw new Error(`Unix socket directory is not owned by the current user: ${directory}`);
await chmod(directory, 0o700);
}
function hasErrorCode(error: unknown, code: string): boolean {
let current = error;
const seen = new Set<unknown>();
@@ -1,4 +1,4 @@
import { lstat, mkdtemp, rm } from "node:fs/promises";
import { chmod, lstat, mkdtemp, rm } from "node:fs/promises";
import { join } from "node:path";
import { afterEach, describe, expect, test } from "vitest";
import {
@@ -27,6 +27,20 @@ afterEach(async () => {
});
describe("experimental memory server composition", () => {
test("does not change permissions on an explicit socket-path parent", async () => {
const directory = await mkdtemp(join("/tmp", "pep-"));
directories.add(directory);
await chmod(directory, 0o750);
const serviceId = "00000000000000000000000000000001";
const runtime = await startExperimentalMemoryServer({
path: join(directory, `${serviceId}.sock`),
serviceId,
});
servers.add(runtime);
expect((await lstat(directory)).mode & 0o777).toBe(0o750);
});
test("discovers and lists seeded sessions without hosting either session", async () => {
const { directory, runtime } = await makeServer();
@@ -99,8 +113,8 @@ describe("experimental memory server composition", () => {
const directory = await mkdtemp(join("/tmp", "pel-"));
directories.add(directory);
const generations = await Promise.all([
startExperimentalServerGeneration({ directory }),
startExperimentalServerGeneration({ directory }),
startExperimentalServerGeneration({ directory, socketDirectory: directory }),
startExperimentalServerGeneration({ directory, socketDirectory: directory }),
]);
for (const generation of generations) servers.add(generation);
@@ -125,8 +139,14 @@ describe("experimental memory server composition", () => {
const otherDirectory = await mkdtemp(join("/tmp", "per-"));
directories.add(firstDirectory);
directories.add(otherDirectory);
const first = await startExperimentalServerGeneration({ directory: firstDirectory });
const other = await startExperimentalServerGeneration({ directory: otherDirectory });
const first = await startExperimentalServerGeneration({
directory: firstDirectory,
socketDirectory: firstDirectory,
});
const other = await startExperimentalServerGeneration({
directory: otherDirectory,
socketDirectory: otherDirectory,
});
servers.add(first);
servers.add(other);
await runExperimentalClient({ command: "client", sessionId: "demo-1" }, { directory: firstDirectory });
@@ -136,7 +156,10 @@ describe("experimental memory server composition", () => {
expect(firstWorkerPid).toEqual(expect.any(Number));
expect(otherWorkerPid).toEqual(expect.any(Number));
const replacement = await startExperimentalServerGeneration({ directory: firstDirectory });
const replacement = await startExperimentalServerGeneration({
directory: firstDirectory,
socketDirectory: firstDirectory,
});
servers.add(replacement);
await first.closed;
@@ -1,6 +1,7 @@
import { type ChildProcess, spawn } from "node:child_process";
import { once } from "node:events";
import { mkdtemp, rm } from "node:fs/promises";
import { mkdir, mkdtemp, rm } from "node:fs/promises";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { fileURLToPath } from "node:url";
import { afterEach, describe, expect, test } from "vitest";
@@ -91,22 +92,35 @@ afterEach(async () => {
describe.skipIf(process.platform === "win32")("experimental CLI server replacement", () => {
test("a second CLI generation starts clean and accepts explicit reattachment", async () => {
const home = await mkdtemp(join("/tmp", "pcs-"));
directories.add(home);
const directory = join(home, ".pi", "server");
const root = await mkdtemp(join(tmpdir(), "pcs-"));
directories.add(root);
const home = join(root, "long-home-segment".repeat(8));
await mkdir(home);
const first = await startServer(home);
const firstIdentity = await waitForOutput(first, /Service: ([0-9a-f]{32})/);
await waitForOutput(first, /Socket: .+\.sock/);
await runExperimentalClient({ command: "client", sessionId: "demo-1" }, { directory });
const firstSocket = await waitForOutput(first, /Socket: (.+\.sock)/);
expect(firstSocket[1]).not.toContain(home);
expect(Buffer.byteLength(firstSocket[1]!)).toBeLessThanOrEqual(103);
await runExperimentalClient({
command: "client",
sessionId: "demo-1",
connect: { transport: "unix", path: firstSocket[1]! },
});
const replacement = await startServer(home);
const replacementIdentity = await waitForOutput(replacement, /Service: ([0-9a-f]{32})/);
await waitForOutput(replacement, /Socket: .+\.sock/);
const replacementSocket = await waitForOutput(replacement, /Socket: (.+\.sock)/);
await waitForExit(first.child);
expect(replacementIdentity[1]).toBe(firstIdentity[1]);
expect(replacementSocket[1]).toBe(firstSocket[1]);
expect(replacement.child.pid).not.toBe(first.child.pid);
await expect(runExperimentalClient({ command: "client" }, { directory })).resolves.toMatchObject({
await expect(
runExperimentalClient({
command: "client",
connect: { transport: "unix", path: replacementSocket[1]! },
}),
).resolves.toMatchObject({
kind: "list",
sessions: [
{ serviceId: firstIdentity[1], sessionId: "demo-1" },
@@ -114,7 +128,11 @@ describe.skipIf(process.platform === "win32")("experimental CLI server replaceme
],
});
await expect(
runExperimentalClient({ command: "client", sessionId: "demo-1" }, { directory }),
runExperimentalClient({
command: "client",
sessionId: "demo-1",
connect: { transport: "unix", path: replacementSocket[1]! },
}),
).resolves.toMatchObject({ kind: "attached", sessionId: "demo-1" });
});
});
-25
View File
@@ -1,4 +1,3 @@
import { fileURLToPath } from "node:url";
import { defineConfig, mergeConfig } from "vitest/config";
import baseConfig, { workspaceSourcePaths } from "../../vitest.base.ts";
@@ -22,30 +21,6 @@ export default mergeConfig(
},
resolve: {
alias: [
{
find: /^@earendil-works\/pi-client\/control$/,
replacement: fileURLToPath(new URL("../client/src/control.ts", import.meta.url)),
},
{
find: /^@earendil-works\/pi-client\/unix$/,
replacement: fileURLToPath(new URL("../client/src/unix.ts", import.meta.url)),
},
{
find: /^@earendil-works\/pi-server\/unix$/,
replacement: fileURLToPath(new URL("../server/src/transports/unix/index.ts", import.meta.url)),
},
{
find: /^@earendil-works\/pi-server$/,
replacement: fileURLToPath(new URL("../server/src/index.ts", import.meta.url)),
},
{
find: /^@earendil-works\/pi-client$/,
replacement: fileURLToPath(new URL("../client/src/index.ts", import.meta.url)),
},
{
find: /^@earendil-works\/pi-protocol$/,
replacement: fileURLToPath(new URL("../protocol/src/index.ts", import.meta.url)),
},
{ find: /^@earendil-works\/pi-ai$/, replacement: workspaceSourcePaths.aiIndex },
{ find: /^@earendil-works\/pi-agent-core$/, replacement: workspaceSourcePaths.agentIndex },
{ find: /^@mariozechner\/pi-ai$/, replacement: workspaceSourcePaths.aiIndex },
+5 -3
View File
@@ -14,7 +14,7 @@ Concurrent attachments to one session reuse one hosted Harness. Attachment is a
```ts
import { MemorySessionRepo } from "@earendil-works/pi-agent-core";
import { generateServiceId, type PiServerHost } from "@earendil-works/pi-server";
import { createUnixServer } from "@earendil-works/pi-server/unix";
import { createUnixServer, getUnixSocketPath } from "@earendil-works/pi-server/unix";
const sessions = new MemorySessionRepo();
const host: PiServerHost = {
@@ -24,13 +24,15 @@ const host: PiServerHost = {
},
};
const serviceId = generateServiceId();
const server = createUnixServer(host, {
serviceId: generateServiceId(),
serviceId,
path: getUnixSocketPath(serviceId, "/run/user/1000/pi"),
});
await server.start();
```
Applications supply the repository and Harness factory. `serviceId` is a logical identity supplied by the launcher, not a socket address. `generateServiceId()` creates an in-memory 128-bit identity. The Unix preset defaults to `~/.pi/server/<serviceId>.sock`; pass `path` to override it. A long-lived launcher can reuse the same ID and path when replacing a server process.
Applications supply the repository and Harness factory. `serviceId` is a logical identity supplied by the launcher, not a socket address. `generateServiceId()` creates an in-memory 128-bit identity. The Unix preset requires an explicit physical `path`; `getUnixSocketPath()` derives one from a caller-selected directory. Choose a short, private runtime directory rather than deriving the route from an unbounded home-directory path. A long-lived launcher can reuse the same ID and path when replacing a server process.
`PiServer` composes authenticated transports through `PiServerListener`. The Unix submodule provides `createUnixListener()` and `createUnixServer()`. Low-level CBOR framing and validation come from `@earendil-works/pi-protocol`.
@@ -1,9 +1,8 @@
import { homedir } from "node:os";
import { join } from "node:path";
import { validateUnixSocketPath } from "./listener.ts";
/** Derive the local Unix socket path for one logical service identity. */
export function getUnixSocketPath(serviceId: string, serverDirectory = join(homedir(), ".pi", "server")): string {
export function getUnixSocketPath(serviceId: string, serverDirectory: string): string {
if (!/^[0-9a-f]{32}$/.test(serviceId)) {
throw new TypeError("Unix serviceId must be 32 lowercase hexadecimal characters");
}
@@ -1,13 +1,12 @@
import { PiServer } from "../../server.ts";
import type { PiServerHost } from "../../types.ts";
import { getUnixSocketPath } from "./address.ts";
import { createUnixListener } from "./listener.ts";
import type { UnixServerOptions } from "./types.ts";
/** Compose PiServer with one Unix-domain socket listener. */
export function createUnixServer(host: PiServerHost, options: UnixServerOptions): PiServer {
const listener = createUnixListener({
path: options.path ?? getUnixSocketPath(options.serviceId),
path: options.path,
mode: options.mode,
maxFrameLength: options.maxFrameLength,
maxPendingBytes: options.maxPendingBytes,
+1 -4
View File
@@ -12,7 +12,4 @@ export interface UnixListenerOptions {
onError?: (error: Error) => void;
}
export interface UnixServerOptions extends Omit<PiServerOptions, "listeners">, Omit<UnixListenerOptions, "path"> {
/** Defaults to ~/.pi/server/<serviceId>.sock. */
path?: string;
}
export interface UnixServerOptions extends Omit<PiServerOptions, "listeners">, UnixListenerOptions {}
+3 -6
View File
@@ -1,7 +1,6 @@
import { type ChildProcess, fork } from "node:child_process";
import { once } from "node:events";
import { lstat, mkdtemp, readFile, rm, unlink, writeFile } from "node:fs/promises";
import { homedir, tmpdir } from "node:os";
import { join } from "node:path";
import { afterEach, describe, expect, test } from "vitest";
import { generateServiceId, type PiServer } from "../src/index.ts";
@@ -14,7 +13,7 @@ const children = new Set<ChildProcess>();
const tempDirectories = new Set<string>();
async function makeSocketPath(nested = false): Promise<string> {
const directory = await mkdtemp(join(tmpdir(), "ps-"));
const directory = await mkdtemp(join("/tmp", "ps-"));
tempDirectories.add(directory);
return nested ? join(directory, "p", "n", "server.sock") : join(directory, "server.sock");
}
@@ -48,16 +47,14 @@ test("generates unique service IDs", () => {
expect(second).not.toBe(first);
});
test("creates an in-memory service ID and derives its Unix socket path", async () => {
const directory = await mkdtemp(join(tmpdir(), "pi-server-"));
test("creates an in-memory service ID and derives its explicit Unix socket path", async () => {
const directory = await mkdtemp(join("/tmp", "pi-server-"));
tempDirectories.add(directory);
const serviceId = generateServiceId();
const path = getUnixSocketPath(serviceId, directory);
expect(serviceId).toMatch(/^[0-9a-f]{32}$/);
expect(path).toBe(join(directory, `${serviceId}.sock`));
expect(getUnixSocketPath(serviceId)).toBe(join(homedir(), ".pi", "server", `${serviceId}.sock`));
const first = createUnixServer(new TestServerHost(), { serviceId, path });
servers.add(first);
await first.start();
+12
View File
@@ -9,6 +9,12 @@ export const workspaceSourcePaths = {
aiOAuth: fileURLToPath(new URL("./packages/ai/src/oauth.ts", import.meta.url)),
aiProviders: fileURLToPath(new URL("./packages/ai/src/providers", import.meta.url)),
agentIndex: fileURLToPath(new URL("./packages/agent/src/index.ts", import.meta.url)),
protocolIndex: fileURLToPath(new URL("./packages/protocol/src/index.ts", import.meta.url)),
clientIndex: fileURLToPath(new URL("./packages/client/src/index.ts", import.meta.url)),
clientControl: fileURLToPath(new URL("./packages/client/src/control.ts", import.meta.url)),
clientUnix: fileURLToPath(new URL("./packages/client/src/unix.ts", import.meta.url)),
serverIndex: fileURLToPath(new URL("./packages/server/src/index.ts", import.meta.url)),
serverUnix: fileURLToPath(new URL("./packages/server/src/transports/unix/index.ts", import.meta.url)),
codingAgentIndex: fileURLToPath(new URL("./packages/coding-agent/src/index.ts", import.meta.url)),
tuiIndex: fileURLToPath(new URL("./packages/tui/src/index.ts", import.meta.url)),
} as const;
@@ -26,6 +32,12 @@ export default defineConfig({
replacement: `${workspaceSourcePaths.aiProviders}/$1.ts`,
},
{ find: /^@earendil-works\/pi-agent-core$/, replacement: workspaceSourcePaths.agentIndex },
{ find: /^@earendil-works\/pi-protocol$/, replacement: workspaceSourcePaths.protocolIndex },
{ find: /^@earendil-works\/pi-client$/, replacement: workspaceSourcePaths.clientIndex },
{ find: /^@earendil-works\/pi-client\/control$/, replacement: workspaceSourcePaths.clientControl },
{ find: /^@earendil-works\/pi-client\/unix$/, replacement: workspaceSourcePaths.clientUnix },
{ find: /^@earendil-works\/pi-server$/, replacement: workspaceSourcePaths.serverIndex },
{ find: /^@earendil-works\/pi-server\/unix$/, replacement: workspaceSourcePaths.serverUnix },
{ find: /^@earendil-works\/pi-tui$/, replacement: workspaceSourcePaths.tuiIndex },
],
},