diff --git a/packages/client/src/effect/service.ts b/packages/client/src/effect/service.ts index b079c0dd5c..eb83145344 100644 --- a/packages/client/src/effect/service.ts +++ b/packages/client/src/effect/service.ts @@ -87,6 +87,10 @@ const discoverLocal = Effect.fnUntraced(function* (options: Options) { // version-mismatched one, and otherwise spawns small contenders until a server // becomes discoverable. A contender is never killed merely for slow startup. export const start = Effect.fn("service.start")(function* (options: StartOptions = {}) { + return yield* startInternal(options) +}) + +const startInternal = Effect.fnUntraced(function* (options: StartOptions, stale?: Info) { const contenders = new Set() let announced = false let reported: Status | undefined @@ -139,7 +143,10 @@ export const start = Effect.fn("service.start")(function* (options: StartOptions yield* kill(service, options, options.version).pipe(Effect.ignore) lastSpawn = 0 return Option.none() - } else if (lastSpawn === 0 && info !== undefined) lastSpawn = Date.now() + } else if (lastSpawn === 0 && info !== undefined) { + // A process observed exiting cannot still own its unchanged registration. + if (stale === undefined || !same(info, stale)) lastSpawn = Date.now() + } const failure = [...contenders].map(contenderFailure).find((error): error is Error => error !== undefined) if (failure !== undefined) return yield* Effect.fail(failure) @@ -203,16 +210,16 @@ export const stop = Effect.fn("service.stop")(function* (options: Options = {}, // registered process that no longer answers authenticated health checks. export const restart = Effect.fn("service.restart")(function* (options: StartOptions = {}) { const existing = yield* registered(options.file, true) - if (existing.service !== undefined) { - const result = yield* kill(existing.service, options, options.version) - if (result === "rejected") return yield* Effect.fail(new Error("Background service rejected restart")) - if (result === "changed") { - const current = yield* read(options.file) - if (current !== undefined && same(current, existing.info)) - yield* terminate(existing.info, read(options.file), true) - } - } else if (existing.info !== undefined) yield* terminate(existing.info, read(options.file), true) - return yield* start(options) + if (existing.info === undefined) return yield* startInternal(options) + + const requested = existing.service === undefined ? "unsupported" : yield* requestStop(existing.service, options.version) + if (requested === "rejected") return yield* Effect.fail(new Error("Background service rejected restart")) + const result = yield* terminate( + existing.info, + read(options.file), + requested === "unsupported" ? "signal" : "requested", + ) + return yield* startInternal(options, result === "stopped" ? existing.info : undefined) }) function fallback() { @@ -264,10 +271,10 @@ const probe = Effect.fnUntraced(function* (info: Info, allowLegacy = false) { ? undefined : { type: "basic" as const, username: "opencode", password: info.password }, } satisfies Endpoint - const response = yield* Effect.tryPromise(() => + const response = yield* Effect.tryPromise((signal) => fetch(new URL("/api/health", info.url), { headers: headers(endpoint), - signal: AbortSignal.timeout(2_000), + signal: AbortSignal.any([signal, AbortSignal.timeout(2_000)]), }), ).pipe(Effect.option, Effect.map(Option.getOrUndefined)) if (response === undefined) return undefined @@ -325,9 +332,9 @@ const stopped = Effect.fnUntraced(function* (pid: number) { const terminate = Effect.fnUntraced(function* ( info: Info, current: Effect.Effect, - graceful: boolean, + mode: "requested" | "signal", ) { - if (graceful) { + if (mode === "signal") { const owner = yield* current if (owner === undefined || !same(owner, info)) return "changed" as const yield* signal(info.pid, "SIGTERM") @@ -348,7 +355,7 @@ const kill = Effect.fnUntraced(function* (service: LocalService, options: Option return yield* terminate( service.info, registered(options.file, true).pipe(Effect.map((result) => result.service?.info)), - requested === "unsupported", + requested === "unsupported" ? "signal" : "requested", ) }) @@ -356,12 +363,12 @@ const decodeStopResponse = Schema.decodeUnknownOption(ServiceStatus.StopResponse const requestStop = Effect.fnUntraced(function* (service: LocalService, targetVersion?: string) { if (service.info.id === undefined || service.legacy) return "unsupported" as const - const response = yield* Effect.tryPromise(() => + const response = yield* Effect.tryPromise((signal) => fetch(new URL("/api/service/stop", service.info.url), { method: "POST", headers: { ...headers(service.endpoint), "content-type": "application/json" }, body: JSON.stringify({ instanceID: service.info.id, targetVersion }), - signal: AbortSignal.timeout(2_000), + signal: AbortSignal.any([signal, AbortSignal.timeout(2_000)]), }), ).pipe(Effect.option, Effect.map(Option.getOrUndefined)) if (response === undefined || response.status === 404 || response.status === 405) return "unsupported" as const diff --git a/packages/client/test/service.test.ts b/packages/client/test/service.test.ts index 9e2fcdc4b2..6c352cd97a 100644 --- a/packages/client/test/service.test.ts +++ b/packages/client/test/service.test.ts @@ -120,6 +120,35 @@ test("explicit restart replaces an unresponsive registered process", async () => } }, 15_000) +test("restart waits for accepted shutdown before starting a replacement", async () => { + const directory = await temp() + const registration = join(directory, "service.json") + const existing = spawn(registration, "graceful") + await waitForFile(registration) + const original = await Bun.file(registration).json() + + const endpoint = await run( + Service.restart({ + file: registration, + version: "test", + command: [process.execPath, fixture, registration, "ready"], + }), + ) + await existing.exited + const info = await Bun.file(registration).json() + + try { + expect(endpoint.url).toBe(info.url) + expect(info.pid).not.toBe(existing.pid) + expect(await Bun.file(registration + ".stop").json()).toEqual({ + instanceID: original.id, + targetVersion: "test", + }) + } finally { + await run(Service.stop({ file: registration })) + } +}) + test("restart recovers when a healthy owner stops responding during shutdown", async () => { const directory = await temp() const registration = join(directory, "service.json") diff --git a/packages/tui/src/context/client.tsx b/packages/tui/src/context/client.tsx index 58308d0302..6dc21e2cf2 100644 --- a/packages/tui/src/context/client.tsx +++ b/packages/tui/src/context/client.tsx @@ -141,6 +141,19 @@ export const { use: useClient, provider: ClientProvider } = createSimpleContext( events.clear() }) + const reloadService = props.reload + const reload = + reloadService === undefined + ? undefined + : async (signal?: AbortSignal) => { + stream?.abort() + try { + await reloadService(signal) + } finally { + if (!abort.signal.aborted) start() + } + } + return { get api() { return api @@ -168,7 +181,7 @@ export const { use: useClient, provider: ClientProvider } = createSimpleContext( }, }, }, - reload: props.reload, + reload, } }, }) diff --git a/packages/tui/test/cli/tui/use-event.test.tsx b/packages/tui/test/cli/tui/use-event.test.tsx index 542f33b3d1..ed47e8f5b0 100644 --- a/packages/tui/test/cli/tui/use-event.test.tsx +++ b/packages/tui/test/cli/tui/use-event.test.tsx @@ -56,6 +56,7 @@ function update(version: string): OpenCodeEvent { async function mount( reconnect?: (onStatus: (status: Service.Status) => void, signal: AbortSignal) => Promise<{ api: OpenCodeClient }>, log?: LogSink, + reload?: (signal?: AbortSignal) => Promise, ) { const events = createEventStream() const calls = createFetch(undefined, events) @@ -70,7 +71,7 @@ async function mount( const app = await testRender(() => ( - + { @@ -349,4 +350,38 @@ describe("useEvent", () => { app.renderer.destroy() await wait(() => aborted) }) + + test("cancels pending endpoint resolution before explicit restart", async () => { + let aborted = false + let restarts = 0 + const { app, events, client } = await mount( + (_onStatus, signal) => + new Promise((_, reject) => { + signal.addEventListener( + "abort", + () => { + aborted = true + reject(signal.reason) + }, + { once: true }, + ) + }), + undefined, + async () => { + restarts += 1 + }, + ) + + try { + await wait(() => client.connection.status() === "connected") + events.disconnect() + await wait(() => client.connection.status() === "reconnecting") + await client.reload?.() + await wait(() => aborted && client.connection.status() === "connected") + + expect(restarts).toBe(1) + } finally { + app.renderer.destroy() + } + }) }) diff --git a/packages/www/content/docs/build/client.mdx b/packages/www/content/docs/build/client.mdx index 0c251ca386..a29c49c7c9 100644 --- a/packages/www/content/docs/build/client.mdx +++ b/packages/www/content/docs/build/client.mdx @@ -115,6 +115,10 @@ Node application: a process. - `Service.start()` reuses a compatible service or starts one when needed. - `Service.stop()` stops the registered service. +- `Service.restart()` explicitly replaces the registered service, including + signaling an unchanged registered PID that no longer answers health checks. + Reserve it for deliberate user recovery; automatic reconnect should use + `Service.start()` to avoid the narrow stale-PID reuse risk. - `Service.headers(endpoint)` creates the authentication headers for a client. ```sh