From 18244dbaf9615611d380fa0196600829fa661c2a Mon Sep 17 00:00:00 2001 From: Szymon Chmal Date: Tue, 6 Oct 2026 11:26:07 +0200 Subject: [PATCH 1/6] =?UTF-8?q?test:=20reproduce=20#395=20=E2=80=94=20a=20?= =?UTF-8?q?gateway=20keeps=20a=20worker=20that=20joined=20while=20starting?= =?UTF-8?q?=20as=20starting,=20and=20routes=20nothing=20to=20it?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/daemon/gateway-fleet.e2e.test.ts | 73 +++++++++++++++++++++++++++- 1 file changed, 71 insertions(+), 2 deletions(-) diff --git a/src/daemon/gateway-fleet.e2e.test.ts b/src/daemon/gateway-fleet.e2e.test.ts index 9844a5fb..0468fd3e 100644 --- a/src/daemon/gateway-fleet.e2e.test.ts +++ b/src/daemon/gateway-fleet.e2e.test.ts @@ -29,8 +29,14 @@ import { join } from "node:path"; import { afterEach, describe, expect, it, vi } from "vitest"; import { FakeDriver } from "../core/testing.js"; -import type { Filesystem } from "../ports/index.js"; -import { FakeHostInfo, MemoryFilesystem, NoopLogger, SystemClock } from "../ports/index.js"; +import type { Filesystem, IpcConnector, IpcListenerFactory } from "../ports/index.js"; +import { + FakeHostInfo, + MemoryFilesystem, + NodeIpcTransport, + NoopLogger, + SystemClock, +} from "../ports/index.js"; import type { DispatchSession } from "./dispatch.js"; import { startDaemon } from "./main.js"; import type { DaemonServer } from "./server.js"; @@ -176,6 +182,8 @@ describe("gateway fleet smoke (ADR 0005 §35)", () => { * directory and filesystem below. Omit for a worker that starts once and is never restarted. */ readonly existing?: { readonly filesystem: Filesystem; readonly directory: string }; + /** The worker's own socket transport; omitted, the real one. */ + readonly ipc?: IpcConnector & IpcListenerFactory; }): Promise<{ daemon: DaemonServer; filesystem: Filesystem; directory: string }> { const directory = options.existing?.directory ?? @@ -212,6 +220,7 @@ describe("gateway fleet smoke (ADR 0005 §35)", () => { }), ], filesystem, + ...(options.ipc === undefined ? {} : { ipc: options.ipc }), logger: new NoopLogger(), statePath: join(directory, "state.json"), version: "1.0.0-e2e", @@ -303,6 +312,66 @@ describe("gateway fleet smoke (ADR 0005 §35)", () => { await gateway.dispatch("lease.release", { leaseId: grant.lease.id }, agentSession()); }, 30_000); + /** + * #395: a worker whose own socket takes a while to claim. The gateway must still see it as + * running once it has started, and route to it, rather than keep the view it read while the + * worker was starting until its next periodic refresh. + */ + it("reports a worker running and grants a no-wait lease on it, when the worker's socket claim is slow", async () => { + const gateway = await startGateway(); + const { secret } = await gateway.dispatch( + "token.create", + { role: "worker", label: "worker-a" }, + adminSession(), + ); + const transport = new NodeIpcTransport(); + let claimSocket = (): void => undefined; + const socketClaimHeld = new Promise((resolve) => { + claimSocket = resolve; + }); + const slowClaim: IpcConnector & IpcListenerFactory = { + connect: (endpoint) => transport.connect(endpoint), + listen: async (endpoint, accept) => { + await socketClaimHeld; + return transport.listen(endpoint, accept); + }, + }; + + const started = startWorker({ + ipc: slowClaim, + label: "worker-a", + model: "Pixel-A", + stdout: "hello-from-a", + token: secret, + }); + // The claim is held long enough for an uplink that dials at once to join and be read. + await vi + .waitFor( + async () => { + const { workers } = await gateway.dispatch("worker.list", {}, adminSession()); + expect(workerServes(workers, "Pixel-A")).toBe(true); + }, + { timeout: 2_000 }, + ) + .catch(() => undefined); + claimSocket(); + await started; + await vi.waitFor(async () => { + const { workers } = await gateway.dispatch("worker.list", {}, adminSession()); + expect(workerServes(workers, "Pixel-A")).toBe(true); + }); + + const { workers } = await gateway.dispatch("worker.list", {}, adminSession()); + expect(workers.map((worker) => worker.health)).toEqual(["running"]); + const grant = await gateway.dispatch( + "lease.request", + { model: "Pixel-A", platform: "android", noWait: true }, + agentSession(), + ); + expect(grant.lease.worker?.label).toBe("worker-a"); + await gateway.dispatch("lease.release", { leaseId: grant.lease.id }, agentSession()); + }, 30_000); + /** * #119's own flagship: drain semantics, `WORKER_UNREACHABLE`-shaped exclusion, and reconnect * rebuild, proved together over real processes and a real WebSocket restart -- not the From 4942b4f4a5c5cb92f7845f7b131cfc1e23203a4f Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 20:40:48 +0000 Subject: [PATCH 2/6] test: add a failed-claim uplink test beside the #395 reproduction (#395) Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01JfqoY9qKUAGytNiMFDAE45 --- src/daemon/main.test.ts | 56 +++++++++++++++++++++++++++++++++++++++++ 1 file changed, 56 insertions(+) diff --git a/src/daemon/main.test.ts b/src/daemon/main.test.ts index 7affdbaf..80889004 100644 --- a/src/daemon/main.test.ts +++ b/src/daemon/main.test.ts @@ -1,6 +1,7 @@ import { mkdir, mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; +import { createServer as createHttpServer } from "node:http"; import { connect, createServer, Server } from "node:net"; import { afterEach, describe, expect, it, vi } from "vitest"; @@ -1194,6 +1195,61 @@ describe("startDaemon socket race with HTTP enabled", () => { }); }); +describe("startDaemon gateway uplink when the socket claim fails", () => { + // #395: the uplink is dialled from `onSocketClaimed`, so a daemon that loses the claim never + // dials its gateway -- it would otherwise join as a worker it is about to stop being. + it("never dials the gateway when the worker's socket is already claimed", async () => { + const directory = await mkdtemp(join(tmpdir(), "simlock-main-uplink-claim-")); + temporaryDirectories.push(directory); + const filesystem = new MemoryFilesystem(); + const dials: string[] = []; + const gateway = createHttpServer(); + gateway.on("upgrade", (request, socket) => { + dials.push(request.url ?? ""); + socket.destroy(); + }); + await new Promise((resolve) => gateway.listen(0, "127.0.0.1", resolve)); + const { port } = gateway.address() as { readonly port: number }; + const options = (gatewayUrl?: string): StartDaemonOptions => + ({ + clock: new FakeClock(1_000), + ...(gatewayUrl === undefined + ? {} + : { configOverrides: { gateway: { token: "join-token-395", url: gatewayUrl } } }), + dataDirectory: directory, + drivers: [ + new FakeDriver({ + availableOsVersions: ["26.5"], + clock: new FakeClock(1_000), + platform: "ios", + }), + ], + filesystem, + logger: new JsonLinesLogger({ + clock: new FakeClock(1_000), + level: "debug", + sink: new MemoryLogSink(), + }), + statePath: join(directory, "state.json"), + version: "1.2.3", + }) as StartDaemonOptions; + + const first = await startDaemon(options()); + try { + const outcome = await startDaemon(options(`ws://127.0.0.1:${port}`)).then( + () => "resolved" as const, + () => "rejected" as const, + ); + // The dial, were it made, reaches a loopback listener within milliseconds. + await new Promise((resolve) => setTimeout(resolve, 500)); + expect({ dials, outcome }).toEqual({ dials: [], outcome: "rejected" }); + } finally { + await first.stop("test-cleanup").catch(() => undefined); + await new Promise((resolve) => gateway.close(() => resolve())); + } + }); +}); + describe("startDaemon HTTP gateway bind failure", () => { // Review finding B6: before this fix, an HTTP bind failure (occupied port) logged and // stopped the daemon from inside `onSocketClaimed`'s handler without `startDaemon()` itself From e4cae848708914cdd6ea02251e4aa4abaa2648e9 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 20:41:09 +0000 Subject: [PATCH 3/6] fix: dial the gateway uplink from onSocketClaimed (#395) 2 failing -> 0 failing Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01JfqoY9qKUAGytNiMFDAE45 --- src/daemon/main.ts | 38 ++++++++++++++++++-------------------- 1 file changed, 18 insertions(+), 20 deletions(-) diff --git a/src/daemon/main.ts b/src/daemon/main.ts index 3b8c594e..8aa89673 100644 --- a/src/daemon/main.ts +++ b/src/daemon/main.ts @@ -502,28 +502,26 @@ export async function startDaemon(options: StartDaemonOptions = {}): Promise { - socketClaimed = true; - void startHttpGateway().then( - () => resolveGatewayStarted?.(), - (error: unknown) => { - logger.error("HTTP frontend failed to start", { message: errorMessage(error) }); - rejectGatewayStarted?.(error); - }, - ); - }, - } - : {}), + onSocketClaimed: () => { + socketClaimed = true; + // Dialled here, once the socket is claimed and `#readyPromise` is set, so the gateway's + // first `status.get` parks on startup readiness like any other request instead of + // reading health `starting` during convergence and keeping it. A daemon that fails its + // claim never reaches this line, so its uplink is never dialled. Nothing awaits the + // uplink: a worker whose gateway is down must still come up and serve its local agents + // (`GatewayUplink` retries on its own backoff). + gatewayUplink?.start(); + if (!config.http.enabled) return; + void startHttpGateway().then( + () => resolveGatewayStarted?.(), + (error: unknown) => { + logger.error("HTTP frontend failed to start", { message: errorMessage(error) }); + rejectGatewayStarted?.(error); + }, + ); + }, }); - // Dialled once the socket is claimed and the dispatcher can answer -- the gateway's first - // `status.get` then parks on startup readiness exactly like any other request, instead of - // racing convergence. Nothing awaits it: a worker whose gateway is down must still come up - // and serve its local agents (`GatewayUplink` retries on its own backoff). - gatewayUplink?.start(); - async function startHttpGateway(): Promise { const httpLogger = logger.child("http"); const app = createHttpApp({ From f175498367a640732dc8ed9e60c5d6045a0f9e39 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 20:42:33 +0000 Subject: [PATCH 4/6] test: pin that the claim callback starts no HTTP gateway with HTTP off (#395) Kills the ConditionalExpression mutant on the http.enabled guard. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01JfqoY9qKUAGytNiMFDAE45 --- src/daemon/main.test.ts | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/src/daemon/main.test.ts b/src/daemon/main.test.ts index 80889004..1dbf00e8 100644 --- a/src/daemon/main.test.ts +++ b/src/daemon/main.test.ts @@ -1250,6 +1250,16 @@ describe("startDaemon gateway uplink when the socket claim fails", () => { }); }); +describe("startDaemon with HTTP off", () => { + // #395: the socket-claim callback now runs with HTTP off too (it dials the uplink), so it + // must still leave the HTTP frontend unstarted. + it("starts no HTTP gateway when http.enabled is false", async () => { + const { sink } = await start({ configOverrides: { http: { enabled: false } } }); + await new Promise((resolve) => setTimeout(resolve, 200)); + expect(sink.records.some((record) => record.message === "HTTP gateway listening")).toBe(false); + }); +}); + describe("startDaemon HTTP gateway bind failure", () => { // Review finding B6: before this fix, an HTTP bind failure (occupied port) logged and // stopped the daemon from inside `onSocketClaimed`'s handler without `startDaemon()` itself From 0459840cc5a2e7c597d31e571e4fb213ef6449ae Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 20:43:58 +0000 Subject: [PATCH 5/6] test: count a refused HTTP bind as a started frontend too (#395) Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01JfqoY9qKUAGytNiMFDAE45 --- src/daemon/main.test.ts | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/src/daemon/main.test.ts b/src/daemon/main.test.ts index 1dbf00e8..d32b3a75 100644 --- a/src/daemon/main.test.ts +++ b/src/daemon/main.test.ts @@ -1256,7 +1256,14 @@ describe("startDaemon with HTTP off", () => { it("starts no HTTP gateway when http.enabled is false", async () => { const { sink } = await start({ configOverrides: { http: { enabled: false } } }); await new Promise((resolve) => setTimeout(resolve, 200)); - expect(sink.records.some((record) => record.message === "HTTP gateway listening")).toBe(false); + // Either outcome of a started frontend counts: bound ("listening") or refused ("failed"). + expect( + sink.records + .map((record) => record.message) + .filter( + (message) => message.startsWith("HTTP gateway") || message.startsWith("HTTP frontend"), + ), + ).toEqual([]); }); }); From 9b49b7a15ac33bb624a4e576aa8048cf22c2bfd2 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 20:49:51 +0000 Subject: [PATCH 6/6] docs: correct onSocketClaimed uplink comment (#395) status.get never parks; events.subscribe does. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01JfqoY9qKUAGytNiMFDAE45 --- src/daemon/main.ts | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/src/daemon/main.ts b/src/daemon/main.ts index 8aa89673..f98ecde0 100644 --- a/src/daemon/main.ts +++ b/src/daemon/main.ts @@ -504,9 +504,10 @@ export async function startDaemon(options: StartDaemonOptions = {}): Promise { socketClaimed = true; - // Dialled here, once the socket is claimed and `#readyPromise` is set, so the gateway's - // first `status.get` parks on startup readiness like any other request instead of - // reading health `starting` during convergence and keeping it. A daemon that fails its + // Dialled here, once the socket is claimed and `#readyPromise` is set. The gateway's first + // `status.get` never parks (it still answers `starting` during convergence), but its + // `events.subscribe` does park on startup readiness, so the refresh that follows the + // subscribe reads `running` instead of keeping `starting`. A daemon that fails its // claim never reaches this line, so its uplink is never dialled. Nothing awaits the // uplink: a worker whose gateway is down must still come up and serve its local agents // (`GatewayUplink` retries on its own backoff).