From 373fb350d8ee689f2e4e7de5f8fc1945900241d7 Mon Sep 17 00:00:00 2001 From: Dax Raad Date: Thu, 6 Aug 2026 00:37:07 -0400 Subject: [PATCH] fix(client): avoid duplicate service contenders --- packages/client/src/effect/service.ts | 21 ++++++++--------- packages/client/src/promise/service.ts | 17 +++++++------- packages/client/test/fixture/service.ts | 14 +++--------- packages/client/test/promise-service.test.ts | 20 +--------------- packages/client/test/service.test.ts | 24 +++----------------- 5 files changed, 25 insertions(+), 71 deletions(-) diff --git a/packages/client/src/effect/service.ts b/packages/client/src/effect/service.ts index f27b273620..06f9ad2080 100644 --- a/packages/client/src/effect/service.ts +++ b/packages/client/src/effect/service.ts @@ -48,11 +48,11 @@ const discoverLocal = Effect.fnUntraced(function* (options: DiscoverOptions) { }) // Idempotent ensure-running: reuses a healthy compatible server, replaces a -// version-mismatched one, and otherwise spawns small contenders until a server -// becomes discoverable. A contender is never killed merely for slow startup. +// version-mismatched one, and otherwise spawns one contender until a server +// becomes discoverable. The contender is never killed merely for slow startup. /** Ensure a healthy, compatible local service is running. */ export const ensure = Effect.fn("service.ensure")(function* (options: EnsureOptions = {}) { - const contenders = new Set() + let contender: Contender | undefined let announced = false let lastSpawn = 0 let spawnDelay = 5_000 @@ -95,17 +95,16 @@ export const ensure = Effect.fn("service.ensure")(function* (options: EnsureOpti return Option.none() } else if (lastSpawn === 0 && info !== undefined) lastSpawn = Date.now() - const finished = [...contenders].filter(contenderFinished) - const failure = finished.map(contenderFailure).find((error): error is Error => error !== undefined) - if (finished.some((item) => item.child.exitCode === 0)) { + const finished = contender !== undefined && contenderFinished(contender) ? contender : undefined + const failure = finished === undefined ? undefined : contenderFailure(finished) + if (finished?.child.exitCode === 0) { spawnDelay = Math.min(spawnDelay * 2, 30_000) } - finished.forEach((item) => contenders.delete(item)) - if (failure !== undefined && contenders.size === 0) return yield* Effect.fail(failure) - // Keep one candidate plus one lock probe so a pre-lock stall cannot block recovery. - if (contenders.size < 2 && Date.now() - lastSpawn >= spawnDelay) { + if (finished !== undefined) contender = undefined + if (failure !== undefined) return yield* Effect.fail(failure) + if (contender === undefined && Date.now() - lastSpawn >= spawnDelay) { yield* announce("missing") - contenders.add(yield* spawnContender) + contender = yield* spawnContender lastSpawn = Date.now() } return Option.none() diff --git a/packages/client/src/promise/service.ts b/packages/client/src/promise/service.ts index c03f4fe551..2bde897be3 100644 --- a/packages/client/src/promise/service.ts +++ b/packages/client/src/promise/service.ts @@ -33,7 +33,7 @@ async function discoverLocal(options: DiscoverOptions) { /** Ensure a healthy, compatible local service is running. */ export async function ensure(options: EnsureOptions = {}): Promise { const deadline = Date.now() + 120_000 - const contenders = new Set() + let contender: Contender | undefined let announced = false let lastSpawn = 0 let spawnDelay = 5_000 @@ -76,17 +76,16 @@ export async function ensure(options: EnsureOptions = {}): Promise { } } else { if (lastSpawn === 0 && registration.info !== undefined) lastSpawn = Date.now() - const finished = [...contenders].filter(contenderFinished) - const failure = finished.map(contenderFailure).find((error) => error !== undefined) - if (finished.some((item) => item.child.exitCode === 0)) { + const finished = contender !== undefined && contenderFinished(contender) ? contender : undefined + const failure = finished === undefined ? undefined : contenderFailure(finished) + if (finished?.child.exitCode === 0) { spawnDelay = Math.min(spawnDelay * 2, 30_000) } - finished.forEach((item) => contenders.delete(item)) - if (failure !== undefined && contenders.size === 0) throw failure - // Keep one candidate plus one lock probe so a pre-lock stall cannot block recovery. - if (contenders.size < 2 && Date.now() - lastSpawn >= spawnDelay) { + if (finished !== undefined) contender = undefined + if (failure !== undefined) throw failure + if (contender === undefined && Date.now() - lastSpawn >= spawnDelay) { announce("missing") - contenders.add(spawnContender()) + contender = spawnContender() lastSpawn = Date.now() } } diff --git a/packages/client/test/fixture/service.ts b/packages/client/test/fixture/service.ts index 7cbbec20ee..766ca00c96 100644 --- a/packages/client/test/fixture/service.ts +++ b/packages/client/test/fixture/service.ts @@ -9,21 +9,13 @@ if (mode === "record-start") { } if (mode === "signal") process.kill(process.pid, process.platform === "win32" ? "SIGTERM" : "SIGKILL") -if ( - mode === "delayed" || - mode === "delayed-failed" || - mode === "coordinated" || - mode === "coordinated-failed-loser" -) { +if (mode === "delayed" || mode === "delayed-failed") { await appendFile(registration + ".starts", process.pid + "\n") const owner = await writeFile(registration + ".owner", String(process.pid), { flag: "wx" }) .then(() => true) .catch(() => false) - if (!owner) process.exit(mode === "coordinated-failed-loser" ? 1 : 0) - if (mode === "coordinated" || mode === "coordinated-failed-loser") { - while ((await Bun.file(registration + ".starts").text()).trim().split("\n").length < 2) await Bun.sleep(10) - if (mode === "coordinated-failed-loser") await Bun.sleep(1_500) - } else await Bun.sleep(Number(delay)) + if (!owner) process.exit(0) + await Bun.sleep(Number(delay)) if (mode === "delayed-failed") process.exit(1) } diff --git a/packages/client/test/promise-service.test.ts b/packages/client/test/promise-service.test.ts index 51744cf50e..f2e04710ad 100644 --- a/packages/client/test/promise-service.test.ts +++ b/packages/client/test/promise-service.test.ts @@ -31,7 +31,7 @@ test("ensures a missing service with native promises", async () => { const endpoint = await Service.ensure({ file: registration, version: "test", - command: [process.execPath, fixture, registration, "coordinated"], + command: [process.execPath, fixture, registration, "delayed", "100"], onStart: (reason) => starts.push(reason), }) const info = await Bun.file(registration).json() @@ -44,24 +44,6 @@ test("ensures a missing service with native promises", async () => { } }, 15_000) -test("waits for a live contender when another native contender fails", async () => { - const directory = await temp() - const registration = join(directory, "service.json") - - const endpoint = await Service.ensure({ - file: registration, - version: "test", - command: [process.execPath, fixture, registration, "coordinated-failed-loser"], - }) - const info = await Bun.file(registration).json() - try { - expect(endpoint.url).toBe(info.url) - } finally { - process.kill(info.pid, "SIGTERM") - await waitForExit(info.pid) - } -}, 15_000) - test("reports a failed registered service", async () => { const registration = await setup("failed-owner") diff --git a/packages/client/test/service.test.ts b/packages/client/test/service.test.ts index 519c2a46df..0f534108c8 100644 --- a/packages/client/test/service.test.ts +++ b/packages/client/test/service.test.ts @@ -121,39 +121,21 @@ test("a legacy health response is still replaced", async () => { await existing.exited }, 10_000) -test("waits for a slow winner while bounding lock probes", async () => { +test("does not spawn another contender while the first is starting", async () => { const directory = await temp() const registration = join(directory, "service.json") const endpoint = await run( Service.ensure({ file: registration, version: "test", - command: [process.execPath, fixture, registration, "coordinated"], + command: [process.execPath, fixture, registration, "delayed", "6000"], }), ) const info = await Bun.file(registration).json() try { expect(endpoint.url).toBe(info.url) expect(await health(endpoint.url)).toEqual({ healthy: true, version: "test", pid: info.pid }) - expect((await Bun.file(registration + ".starts").text()).trim().split("\n")).toHaveLength(2) - } finally { - process.kill(info.pid, "SIGTERM") - } -}, 15_000) - -test("waits for a live contender when another contender fails", async () => { - const directory = await temp() - const registration = join(directory, "service.json") - const endpoint = await run( - Service.ensure({ - file: registration, - version: "test", - command: [process.execPath, fixture, registration, "coordinated-failed-loser"], - }), - ) - const info = await Bun.file(registration).json() - try { - expect(endpoint.url).toBe(info.url) + expect((await Bun.file(registration + ".starts").text()).trim().split("\n")).toHaveLength(1) } finally { process.kill(info.pid, "SIGTERM") }