From 72c8bf9946ee18ca9eaf22f3ba7180ec83789a8c Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Thu, 20 Aug 2026 23:39:12 -0700 Subject: [PATCH] fix(gateway): retain Bonjour cleanup on startup failure (#127062) --- src/gateway/server-core-runtime.ts | 6 +-- src/gateway/server-kernel.test.ts | 19 ++++++++- src/gateway/server-lifecycle.ts | 4 +- src/gateway/server-startup-early.test.ts | 50 ++++++++++++++++++++++++ src/gateway/server-startup-early.ts | 9 +++-- 5 files changed, 77 insertions(+), 11 deletions(-) diff --git a/src/gateway/server-core-runtime.ts b/src/gateway/server-core-runtime.ts index 15ec2b188c74..341d2538c95f 100644 --- a/src/gateway/server-core-runtime.ts +++ b/src/gateway/server-core-runtime.ts @@ -186,6 +186,7 @@ export async function startGatewayCoreRuntime(input: { log, logDiscovery, nodeRegistry, + swapBonjourStop: kernel.swapBonjourStop, pluginRegistry: pluginRuntime.registry, broadcast, nodeSendToAllSubscribed, @@ -226,10 +227,7 @@ export async function startGatewayCoreRuntime(input: { const discoveryResident = residentRegistry.register({ name: "bonjour-discovery", start: startEarlyRuntime, - stop: async () => { - const earlyRuntime = await startEarlyRuntime(); - await earlyRuntime.bonjourStop?.(); - }, + stop: async () => await kernel.swapBonjourStop(null)?.(), }); const taskAndSkillsResident = residentRegistry.register({ name: "task-and-skills-runtime", diff --git a/src/gateway/server-kernel.test.ts b/src/gateway/server-kernel.test.ts index eaa97aa58cc8..7b6591d33679 100644 --- a/src/gateway/server-kernel.test.ts +++ b/src/gateway/server-kernel.test.ts @@ -20,7 +20,7 @@ import { createSyntheticPluginRuntimeClient } from "./server-plugin-runtime-clie describe("createGatewayKernel", () => { it("reports startup and readiness as draining during a direct close", async () => { - const port = await getFreePort(); + const port = 19_789; const state = await createOpenClawTestState({ label: "gateway-kernel-direct-close-readiness", layout: "home", @@ -56,6 +56,20 @@ describe("createGatewayKernel", () => { expect(getStartup()).toMatchObject({ ok: true, status: "started" }); expect(getReadiness()).toMatchObject({ ready: true, failing: [] }); + const discoveryResident = kernel.residentRegistry + .list() + .find((resident) => resident.name === "bonjour-discovery"); + if (!discoveryResident) { + throw new Error("Expected the Gateway discovery resident"); + } + const residentFirstStop = vi.fn(async () => {}); + kernel.kernel.swapBonjourStop(residentFirstStop); + await discoveryResident.stop(); + expect(residentFirstStop).toHaveBeenCalledOnce(); + expect(kernel.runtimeState.bonjourStop).toBeNull(); + + const closeFirstStop = vi.fn(async () => {}); + kernel.kernel.swapBonjourStop(closeFirstStop); const configReloaderStop = createDeferred(); vi.spyOn(kernel.runtimeState.configReloader, "stop").mockReturnValue( configReloaderStop.promise, @@ -66,6 +80,9 @@ describe("createGatewayKernel", () => { expect(getReadiness()).toMatchObject({ ready: false, failing: ["gateway-draining"] }); configReloaderStop.resolve(); await closing; + await discoveryResident.stop(); + expect(closeFirstStop).toHaveBeenCalledOnce(); + expect(kernel.runtimeState.bonjourStop).toBeNull(); } finally { try { await kernel?.closeOnStartupFailure(); diff --git a/src/gateway/server-lifecycle.ts b/src/gateway/server-lifecycle.ts index 67064a1b02f4..288255b882ba 100644 --- a/src/gateway/server-lifecycle.ts +++ b/src/gateway/server-lifecycle.ts @@ -250,11 +250,9 @@ export async function prepareGatewayLifecycle(params: { runtimeState.gatewayMethods.splice(0, runtimeState.gatewayMethods.length, ...methods); }, setEarlyRuntimeHandles: (handles: { - bonjourStop: typeof runtimeState.bonjourStop; getActiveTaskCount: () => number; skillsChangeUnsub: typeof runtimeState.skillsChangeUnsub; }) => { - runtimeState.bonjourStop = handles.bonjourStop; activeTaskCount.get = handles.getActiveTaskCount; runtimeState.skillsChangeUnsub = handles.skillsChangeUnsub; }, @@ -515,7 +513,7 @@ export async function prepareGatewayLifecycle(params: { const transport = transportBridge.current(); await transport?.portalService.closeAll(); await shutdownRuntime.createGatewayCloseHandler({ - bonjourStop: runtimeState.bonjourStop, + bonjourStop: kernel.swapBonjourStop(null), tailscaleCleanup: runtimeState.tailscaleCleanup, clearSecretsRuntimeSnapshot: clearSecretsRuntimeSnapshotState, channelIds, diff --git a/src/gateway/server-startup-early.test.ts b/src/gateway/server-startup-early.test.ts index 052576602526..1f022bef6c5b 100644 --- a/src/gateway/server-startup-early.test.ts +++ b/src/gateway/server-startup-early.test.ts @@ -2,6 +2,7 @@ * Early gateway startup helper tests. */ import { beforeEach, describe, expect, it, vi } from "vitest"; +import { runGatewayShutdownSteps } from "./server-shutdown.js"; import { createGatewayMaintenanceStateForTest } from "./test-helpers.maintenance-state.js"; type StartGatewayDiscovery = typeof import("./server-discovery-runtime.js").startGatewayDiscovery; @@ -82,6 +83,7 @@ function earlyRuntimeInput( log, logDiscovery: log, nodeRegistry: {} as never, + swapBonjourStop: () => null, ...maintenanceState, skillsRefreshDelayMs: 30_000, getSkillsRefreshTimer: () => null, @@ -150,6 +152,48 @@ describe("startGatewayEarlyRuntime", () => { expect(mocks.closeSkillsWatchers).toHaveBeenCalledTimes(1); }); + it.each([false, true])( + "stops acquired discovery exactly once after later startup failure (cleanup rejects: %s)", + async (cleanupRejects) => { + const startupError = new Error("remote skills registry failed"); + const cleanupError = new Error("discovery cleanup failed"); + const stopDiscovery = vi.fn(async () => { + if (cleanupRejects) { + throw cleanupError; + } + }); + const owner: { current: (() => Promise) | null } = { current: null }; + const swapBonjourStop = (next: typeof owner.current) => { + const previous = owner.current; + owner.current = next; + return previous; + }; + mocks.startGatewayDiscovery.mockResolvedValueOnce({ bonjourStop: stopDiscovery }); + mocks.setSkillsRemoteRegistry.mockImplementationOnce(() => { + throw startupError; + }); + const onCleanupError = vi.fn(); + + const startup = startGatewayEarlyRuntime( + earlyRuntimeInput({ minimalTestGateway: false, swapBonjourStop }), + ).catch(async (error: unknown) => { + await runGatewayShutdownSteps({ + steps: [ + { name: "discovery resident", run: async () => await swapBonjourStop(null)?.() }, + { name: "gateway close", run: async () => await swapBonjourStop(null)?.() }, + ], + onError: onCleanupError, + }); + throw error; + }); + + await expect(startup).rejects.toBe(startupError); + expect(stopDiscovery).toHaveBeenCalledOnce(); + expect(owner.current).toBeNull(); + expect(onCleanupError).toHaveBeenCalledTimes(cleanupRejects ? 1 : 0); + }, + ); + it("broadcasts remote-node skill invalidations to operator clients", async () => { const broadcast = vi.fn(); @@ -212,6 +256,9 @@ describe("startGatewayEarlyRuntime", () => { }); it("fails before discovery and task maintenance when task state cannot restore", async () => { + const stopDiscovery = vi.fn(async () => {}); + const swapBonjourStop = vi.fn(() => null); + mocks.startGatewayDiscovery.mockResolvedValue({ bonjourStop: stopDiscovery }); mocks.ensureTaskRuntimeStateReady.mockImplementationOnce(() => { throw new Error("task-flow registry restore failed"); }); @@ -220,11 +267,14 @@ describe("startGatewayEarlyRuntime", () => { startGatewayEarlyRuntime( earlyRuntimeInput({ minimalTestGateway: false, + swapBonjourStop, }), ), ).rejects.toThrow("task-flow registry restore failed"); expect(mocks.startGatewayDiscovery).not.toHaveBeenCalled(); + expect(swapBonjourStop).not.toHaveBeenCalled(); + expect(stopDiscovery).not.toHaveBeenCalled(); expect(mocks.configureTaskRegistryMaintenance).not.toHaveBeenCalled(); expect(mocks.startTaskRegistryMaintenance).not.toHaveBeenCalled(); }); diff --git a/src/gateway/server-startup-early.ts b/src/gateway/server-startup-early.ts index 7a0c4c02b107..086fa065eb30 100644 --- a/src/gateway/server-startup-early.ts +++ b/src/gateway/server-startup-early.ts @@ -71,6 +71,7 @@ export async function startGatewayEarlyRuntime(params: { warn: (msg: string) => void; }; nodeRegistry: Parameters[0]; + swapBonjourStop: (next: (() => Promise) | null) => (() => Promise) | null; pluginRegistry?: PluginRegistry; broadcast: GatewayMaintenanceParams["broadcast"]; nodeSendToAllSubscribed: Parameters[0]["nodeSendToAllSubscribed"]; @@ -102,8 +103,11 @@ export async function startGatewayEarlyRuntime(params: { ensureTaskRuntimeStateReady(); }); } - const bonjourStop = await measureStartup(params.startupTrace, "runtime.early.discovery", () => - startGatewayPluginDiscovery(params), + // Startup failure can occur immediately after discovery; publish its owner first. + params.swapBonjourStop( + await measureStartup(params.startupTrace, "runtime.early.discovery", () => + startGatewayPluginDiscovery(params), + ), ); let getActiveTaskCount = () => 0; @@ -205,7 +209,6 @@ export async function startGatewayEarlyRuntime(params: { }; return { - bonjourStop, getActiveTaskCount, skillsChangeUnsub, startMaintenance,