From 4c68b51c1ba4ca4fce4ba870ae1cfe8d94f04e85 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Tue, 1 Sep 2026 21:32:36 +0000 Subject: [PATCH] fix(events): claim the IPC endpoint atomically before removing stale sockets (#209 follow-up) --- .changeset/atomic-event-endpoint.md | 5 + packages/agent-bundle/src/events/ipc.ts | 167 +++++++++++++++--- packages/agent-bundle/tests/event-ipc.test.ts | 88 +++++++++ 3 files changed, 239 insertions(+), 21 deletions(-) create mode 100644 .changeset/atomic-event-endpoint.md diff --git a/.changeset/atomic-event-endpoint.md b/.changeset/atomic-event-endpoint.md new file mode 100644 index 000000000..aa5196b48 --- /dev/null +++ b/.changeset/atomic-event-endpoint.md @@ -0,0 +1,5 @@ +--- +'agent-bundle': patch +--- + +Prevent concurrent event runtime startups from unlinking a newly claimed IPC socket. diff --git a/packages/agent-bundle/src/events/ipc.ts b/packages/agent-bundle/src/events/ipc.ts index 957e5c564..981e1968a 100644 --- a/packages/agent-bundle/src/events/ipc.ts +++ b/packages/agent-bundle/src/events/ipc.ts @@ -1,5 +1,5 @@ import { createHash } from 'node:crypto'; -import { chmod, mkdir, rm, stat } from 'node:fs/promises'; +import { chmod, mkdir, open, rm, stat, type FileHandle } from 'node:fs/promises'; import { createConnection, createServer, type Server, type Socket } from 'node:net'; import { dirname, join } from 'node:path'; import { StringDecoder } from 'node:string_decoder'; @@ -12,6 +12,8 @@ import { liftPromise } from '../effect/lift.ts'; const EVENT_RUNTIME_PROTOCOL_VERSION = 1 as const; const MAX_EVENT_MESSAGE_BYTES = 1024 * 1024; +const ENDPOINT_CLAIM_RETRY_COUNT = 100; +const ENDPOINT_CLAIM_RETRY_DELAY = Duration.millis(10); export type EventRuntimeTransportErrorCode = | 'epoch-mismatch' @@ -227,6 +229,16 @@ class EventSocketService extends Context.Service Promise; +} + +interface EndpointClaim { + readonly handle: FileHandle; + readonly identity: Readonly<{ readonly device: number; readonly inode: number }>; + readonly path: string; +} + const probeEndpoint = (endpoint: string): Effect.Effect => Effect.callback((resume) => { const socket = createConnection(endpoint); @@ -265,17 +277,54 @@ const probeEndpoint = (endpoint: string): Effect.Effect => Effect.gen(function*() { - const endpoint = eventRuntimeEndpoint(options.endpointId); - if (process.platform !== 'win32') { - yield* liftPromise(async () => { - await mkdir(dirname(endpoint), { mode: 0o700, recursive: true }); - await chmod(dirname(endpoint), 0o700); - }).pipe( - Effect.mapError((error) => transportError('runtime-failed', 'Unable to prepare the event runtime endpoint.', error)), - ); +const tryClaimEndpoint = ( + endpoint: string, +): Effect.Effect => + liftPromise(async () => { + const path = `${endpoint}.lock`; + let handle: FileHandle | undefined; + try { + handle = await open(path, 'wx', 0o600); + const lockStat = await handle.stat(); + return { + handle, + identity: { device: lockStat.dev, inode: lockStat.ino }, + path, + }; + } catch (error) { + await handle?.close(); + if ((error as NodeJS.ErrnoException).code === 'EEXIST') return undefined; + throw error; + } + }).pipe( + Effect.mapError((error) => transportError('runtime-failed', 'Unable to claim the event runtime endpoint.', error)), + ); + +const releaseEndpointClaim = (claim: EndpointClaim): Effect.Effect => + liftPromise(async () => { + await claim.handle.close(); + let current; + try { + current = await stat(claim.path); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return; + throw error; + } + if ( + current.dev === claim.identity.device + && current.ino === claim.identity.inode + ) { + await rm(claim.path, { force: true }); + } + }).pipe(Effect.ignore); + +const claimEndpoint = Effect.fnUntraced(function*( + endpoint: string, +): Effect.fn.Return { + for (let attempt = 0; attempt < ENDPOINT_CLAIM_RETRY_COUNT; attempt += 1) { + const claim = yield* tryClaimEndpoint(endpoint); + if (claim !== undefined) return claim; + const endpointState = yield* probeEndpoint(endpoint); if (endpointState === 'live') { return yield* Effect.fail(transportError( @@ -283,13 +332,25 @@ const openServer = ( 'Event runtime endpoint already has a live server.', )); } - if (endpointState === 'stale') { - yield* liftPromise(() => rm(endpoint)).pipe( - Effect.mapError((error) => transportError('runtime-failed', 'Unable to remove the stale event runtime endpoint.', error)), - ); + if (attempt + 1 < ENDPOINT_CLAIM_RETRY_COUNT) { + yield* Effect.sleep(ENDPOINT_CLAIM_RETRY_DELAY); } } - return yield* Effect.callback((resume) => { + + // A crashed process can leave the claim file behind. Retrying the endpoint + // probe is bounded rather than stealing an unverifiable claim and recreating + // the same unlink race that this lock prevents. + return yield* Effect.fail(transportError( + 'runtime-failed', + 'Event runtime endpoint claim did not clear after bounded retries.', + )); +}); + +const listenServer = ( + options: CreateEventRuntimeServerOptions, + endpoint: string, +): Effect.Effect => + Effect.callback((resume) => { const sockets = new Set(); const server = createServer({ allowHalfOpen: true }, (socket) => { sockets.add(socket); @@ -337,6 +398,57 @@ const openServer = ( server.close(); }); }); + +const openServer = ( + options: CreateEventRuntimeServerOptions, + testHooks?: EventRuntimeServerTestHooks, +): Effect.Effect => Effect.gen(function*() { + const endpoint = eventRuntimeEndpoint(options.endpointId); + if (process.platform === 'win32') return yield* listenServer(options, endpoint); + + yield* liftPromise(async () => { + await mkdir(dirname(endpoint), { mode: 0o700, recursive: true }); + await chmod(dirname(endpoint), 0o700); + }).pipe( + Effect.mapError((error) => transportError('runtime-failed', 'Unable to prepare the event runtime endpoint.', error)), + ); + const endpointState = yield* probeEndpoint(endpoint); + const afterEndpointProbe = testHooks?.afterEndpointProbe; + if (afterEndpointProbe !== undefined) { + yield* liftPromise(() => afterEndpointProbe(endpointState)).pipe( + Effect.mapError((error) => transportError('runtime-failed', 'Event runtime endpoint probe hook failed.', error)), + ); + } + if (endpointState === 'live') { + return yield* Effect.fail(transportError( + 'runtime-failed', + 'Event runtime endpoint already has a live server.', + )); + } + + return yield* Effect.acquireUseRelease( + claimEndpoint(endpoint), + () => Effect.gen(function*() { + const claimedEndpointState = yield* probeEndpoint(endpoint); + if (claimedEndpointState === 'live') { + return yield* Effect.fail(transportError( + 'runtime-failed', + 'Event runtime endpoint already has a live server.', + )); + } + if (claimedEndpointState === 'stale') { + yield* liftPromise(() => rm(endpoint)).pipe( + Effect.mapError((error) => transportError( + 'runtime-failed', + 'Unable to remove the stale event runtime endpoint.', + error, + )), + ); + } + return yield* listenServer(options, endpoint); + }), + releaseEndpointClaim, + ); }); const removeOwnedEndpoint = (service: EventSocketServiceShape): Effect.Effect => @@ -374,13 +486,17 @@ const closeServer = (service: EventSocketServiceShape): Effect.Effect => ), ); -const eventSocketLayer = (options: CreateEventRuntimeServerOptions): Layer.Layer => - Layer.effect(EventSocketService, Effect.acquireRelease(openServer(options), closeServer)); +const eventSocketLayer = ( + options: CreateEventRuntimeServerOptions, + testHooks?: EventRuntimeServerTestHooks, +): Layer.Layer => + Layer.effect(EventSocketService, Effect.acquireRelease(openServer(options, testHooks), closeServer)); -export const createEventRuntimeServer = async ( +const createEventRuntimeServerWithHooks = async ( options: CreateEventRuntimeServerOptions, + testHooks?: EventRuntimeServerTestHooks, ): Promise => { - const runtime: ScopedEffectRuntime = makeScopedEffectRuntime(eventSocketLayer(options)); + const runtime: ScopedEffectRuntime = makeScopedEffectRuntime(eventSocketLayer(options, testHooks)); const service = await runtime.run(EventSocketService); return Object.freeze({ close: () => runtime.close(), @@ -388,6 +504,15 @@ export const createEventRuntimeServer = async ( }); }; +export const createEventRuntimeServer = async ( + options: CreateEventRuntimeServerOptions, +): Promise => createEventRuntimeServerWithHooks(options); + +export const createEventRuntimeServerForTest = async ( + options: CreateEventRuntimeServerOptions, + testHooks: EventRuntimeServerTestHooks, +): Promise => createEventRuntimeServerWithHooks(options, testHooks); + const connect = (endpoint: string): Effect.Effect => Effect.callback((resume) => { const socket = createConnection(endpoint); diff --git a/packages/agent-bundle/tests/event-ipc.test.ts b/packages/agent-bundle/tests/event-ipc.test.ts index 25896d7d8..f641991be 100644 --- a/packages/agent-bundle/tests/event-ipc.test.ts +++ b/packages/agent-bundle/tests/event-ipc.test.ts @@ -6,6 +6,7 @@ import { expect, it } from 'effect-rstest'; import { createEventRuntimeServer, + createEventRuntimeServerForTest, EventRuntimeTransportError, requestEventRuntime, } from '../src/events/ipc.ts'; @@ -119,6 +120,93 @@ it.live('replaces a stale event runtime socket file', () => Effect.gen(function* expect(response).toEqual({ replaced: true }); })); +it.live('does not unlink a concurrent winner after both servers probe a stale endpoint', () => Effect.gen(function*() { + if (process.platform === 'win32') return; + const endpointId = `event-ipc-stale-race-${crypto.randomUUID()}`; + const endpoint = yield* Effect.acquireUseRelease( + Effect.promise(() => createEventRuntimeServer({ + artifactEpoch: 'epoch-1', + endpointId, + handle: async () => undefined, + })), + (server) => Effect.succeed(server.endpoint), + (server) => Effect.promise(() => server.close()), + ); + yield* Effect.promise(() => writeFile(endpoint, 'stale socket')); + + let markFirstProbed!: () => void; + let releaseFirst!: () => void; + const firstProbed = new Promise((resolve) => { markFirstProbed = resolve; }); + const firstCanContinue = new Promise((resolve) => { releaseFirst = resolve; }); + const firstServer = createEventRuntimeServerForTest({ + artifactEpoch: 'epoch-1', + endpointId, + handle: async () => ({ owner: 'first' }), + }, { + afterEndpointProbe: async (state) => { + expect(state).toBe('stale'); + markFirstProbed(); + await firstCanContinue; + }, + }); + yield* Effect.promise(() => firstProbed); + + let markSecondProbed!: () => void; + let releaseSecond!: () => void; + const secondProbed = new Promise((resolve) => { markSecondProbed = resolve; }); + const secondCanContinue = new Promise((resolve) => { releaseSecond = resolve; }); + const secondServer = createEventRuntimeServerForTest({ + artifactEpoch: 'epoch-1', + endpointId, + handle: async () => ({ owner: 'second' }), + }, { + afterEndpointProbe: async (state) => { + expect(state).toBe('stale'); + markSecondProbed(); + await secondCanContinue; + }, + }); + yield* Effect.promise(() => secondProbed); + + releaseFirst(); + const winner = yield* Effect.acquireRelease( + Effect.promise(() => firstServer), + (server) => Effect.promise(() => server.close()), + ); + expect(winner.endpoint).toBe(endpoint); + + releaseSecond(); + const loser = yield* Effect.promise(async () => { + try { + return { server: await secondServer, status: 'opened' as const }; + } catch (error) { + return { error, status: 'failed' as const }; + } + }); + if (loser.status === 'opened') { + yield* Effect.promise(() => loser.server.close()); + expect(loser.status).toBe('failed'); + return; + } + expect(loser.error).toMatchObject({ + code: 'runtime-failed', + message: expect.stringMatching(/already has a live server/u), + name: EventRuntimeTransportError.name, + }); + + const response = yield* Effect.promise(() => requestEventRuntime({ + artifactEpoch: 'epoch-1', + endpointId, + event: 'session/start', + hostContractRevision: '2.1.250', + native: { hook_event_name: 'SessionStart' }, + signal: new AbortController().signal, + target: 'claude', + timeoutMs: 1_000, + })); + expect(response).toEqual({ owner: 'first' }); +})); + it.live('interrupts an in-flight event handler when the client disconnects', () => Effect.gen(function*() { const endpointId = `event-ipc-disconnect-${crypto.randomUUID()}`; let markStarted!: () => void;