diff --git a/src/config/io.write-config.test.ts b/src/config/io.write-config.test.ts index 43c6cd59ef07..c1dafcbb7fe0 100644 --- a/src/config/io.write-config.test.ts +++ b/src/config/io.write-config.test.ts @@ -3108,7 +3108,7 @@ gateway: { mode: "local", port: 18789 } initialPluginInstallRecords: {}, readPluginInstallRecords: async () => ({}), onNoopConfigCommit: async () => {}, - onHotReload: async () => {}, + onHotReload: async () => "applied" as const, onRestart: async () => {}, log: { info: () => {}, ...silentLogger }, watchPath: configPath, @@ -3192,7 +3192,7 @@ gateway: { mode: "local", port: 18789 } initialPluginInstallRecords: {}, readPluginInstallRecords: async () => ({}), onNoopConfigCommit: async () => {}, - onHotReload: async () => {}, + onHotReload: async () => "applied" as const, onRestart: async () => {}, log: { info: () => {}, ...silentLogger }, watchPath: configPathB, diff --git a/src/config/runtime-write-application.ts b/src/config/runtime-write-application.ts index c9dbbb91a31d..9f3a155c12be 100644 --- a/src/config/runtime-write-application.ts +++ b/src/config/runtime-write-application.ts @@ -2,6 +2,7 @@ import { createDeferredCore } from "../shared/deferred.js"; export type RuntimeConfigWriteApplicationStatus = | "applied" + | "restart-required" | "superseded" | "failed" | "stopped" diff --git a/src/gateway/config-reload-status.types.ts b/src/gateway/config-reload-status.types.ts index bfe2cc3a6b0b..a38d160c518c 100644 --- a/src/gateway/config-reload-status.types.ts +++ b/src/gateway/config-reload-status.types.ts @@ -7,3 +7,5 @@ // after watcher re-creation fails past the retry budget, so operators/callers // can detect silent degradation instead of assuming reloads still fire. export type GatewayHotReloadStatus = "active" | "disabled"; + +export type GatewayHotReloadApplicationStatus = "applied" | "restart-required"; diff --git a/src/gateway/config-reload.test.ts b/src/gateway/config-reload.test.ts index 7d82f3f8ab44..d3e6e799abee 100644 --- a/src/gateway/config-reload.test.ts +++ b/src/gateway/config-reload.test.ts @@ -870,7 +870,7 @@ function createReloaderHarness( nextConfig: OpenClawConfig, ownership: GatewayConfigReloadTransactionOwnership, sourceConfig: OpenClawConfig, - ) => Promise; + ) => Promise<"applied" | "restart-required">; onRestart?: ( plan: GatewayReloadPlan, nextConfig: OpenClawConfig, @@ -907,7 +907,7 @@ function createReloaderHarness( _plan: GatewayReloadPlan, _nextConfig: OpenClawConfig, _ownership: GatewayConfigReloadTransactionOwnership, - ) => {}), + ) => "applied" as const), ); const onRestart = vi.fn( options.onRestart ?? ((_plan: GatewayReloadPlan, _nextConfig: OpenClawConfig) => {}), @@ -1043,7 +1043,7 @@ describe("startGatewayConfigReloader include files", () => { logger: { error: vi.fn(), warn: vi.fn() }, }); const initialSnapshot = await configIo.readConfigFileSnapshot(); - const onHotReload = vi.fn(async () => {}); + const onHotReload = vi.fn(async () => "applied" as const); let signalWatcherReady!: () => void; const watcherReady = new Promise((resolve) => { signalWatcherReady = resolve; @@ -1650,7 +1650,10 @@ describe("startGatewayConfigReloader", () => { const harness = createReloaderHarness(vi.fn(), { initialSnapshotRawHash: null, initialAuthoredConfig: {}, - onHotReload: async () => await hotReloadGate, + onHotReload: async () => { + await hotReloadGate; + return "applied" as const; + }, }); let settled = false; void application.result.then(() => { @@ -1707,6 +1710,7 @@ describe("startGatewayConfigReloader", () => { markHotReloadStarted(); await hotReloadGate; ownership.markRuntimeCommitted(runtimeConfig, plan); + return "applied" as const; }, ); const log = { info: vi.fn(), warn: vi.fn(), error: vi.fn() }; @@ -1800,6 +1804,26 @@ describe("startGatewayConfigReloader", () => { await harness.reloader.stop(); }); + it("reports when a committed hot reload requires recovery restart", async () => { + const application = createRuntimeConfigWriteApplication(); + const harness = createReloaderHarness(vi.fn(), { + initialSnapshotRawHash: null, + initialAuthoredConfig: {}, + onHotReload: async () => "restart-required", + }); + + harness.emitWrite( + attachRuntimeConfigWriteApplication( + makeZeroDebounceHookWrite("application-restart-required"), + application, + ), + ); + await vi.runAllTimersAsync(); + + await expect(application.result).resolves.toBe("restart-required"); + await harness.reloader.stop(); + }); + it("settles a coalesced in-process write as superseded", async () => { const supersededApplication = createRuntimeConfigWriteApplication(); const appliedApplication = createRuntimeConfigWriteApplication(); @@ -2095,6 +2119,7 @@ describe("startGatewayConfigReloader", () => { writtenAtMs: Date.now(), afterWrite: { mode: "none", reason: "baseline-only acceptance" }, }); + return "applied" as const; }; const harness = createReloaderHarness(vi.fn(), { initialConfig, @@ -2155,6 +2180,7 @@ describe("startGatewayConfigReloader", () => { terminalPolicy.prepareConfig(nextConfig, { restartPending: false }); ownership.markRuntimeCommitted(nextConfig, plan); harness.watcher.emit("change"); + return "applied" as const; }; const harness = createReloaderHarness( vi.fn(async () => rejectedSnapshot), @@ -2220,6 +2246,7 @@ describe("startGatewayConfigReloader", () => { sourceFingerprint: "source-restart-b", writtenAtMs: Date.now(), }); + return "applied" as const; }; const harness = createReloaderHarness(vi.fn(), { initialConfig, @@ -2331,9 +2358,13 @@ describe("startGatewayConfigReloader", () => { await blocked; expect(ownership.isCurrent()).toBe(false); }; + const hotReloadA = async (...args: Parameters) => { + await publishA(...args); + return "applied" as const; + }; const harness = createReloaderHarness(readSnapshot, { initialConfig, - ...(kind === "noop" ? { onNoopConfigCommit: publishA } : { onHotReload: publishA }), + ...(kind === "noop" ? { onNoopConfigCommit: publishA } : { onHotReload: hotReloadA }), }); harness.watcher.emit("change"); @@ -2396,6 +2427,7 @@ describe("startGatewayConfigReloader", () => { recordCommitted?.(); await tailGate; } + return "applied" as const; }, ); const promoteSnapshot = vi.fn(async (_snapshot: ConfigFileSnapshot, _reason: string) => true); @@ -2468,6 +2500,7 @@ describe("startGatewayConfigReloader", () => { if (nextConfig === configA) { emitWrite(configB, "env-b", 2); } + return "applied"; }, }); const emitWrite = (config: OpenClawConfig, hash: string, revision: number) => { @@ -2661,6 +2694,7 @@ describe("startGatewayConfigReloader", () => { if (!ownership.isCurrent()) { throw new Error("unlinked config A was superseded"); } + return "applied"; }, }); @@ -2772,6 +2806,7 @@ describe("startGatewayConfigReloader", () => { harness.onHotReload.mockImplementationOnce(async () => { markReloadStarted?.(); await reloadBlocked; + return "applied"; }); harness.watcher.emit("change"); @@ -5044,7 +5079,7 @@ describe("startGatewayConfigReloader watcher error recovery", () => { initialPluginInstallRecords: {}, readPluginInstallRecords: async () => ({}), onNoopConfigCommit: vi.fn(async () => {}), - onHotReload: vi.fn(async () => {}), + onHotReload: vi.fn(async () => "applied" as const), onRestart: vi.fn(), log, watchPath: "/tmp/openclaw.json", diff --git a/src/gateway/config-reload.ts b/src/gateway/config-reload.ts index 8a50ae8b7ca9..7877f0096089 100644 --- a/src/gateway/config-reload.ts +++ b/src/gateway/config-reload.ts @@ -45,7 +45,10 @@ import { type GatewayReloadPlan, } from "./config-reload-plan.js"; import { resolveGatewayReloadSettings } from "./config-reload-settings.js"; -import type { GatewayHotReloadStatus } from "./config-reload-status.types.js"; +import type { + GatewayHotReloadApplicationStatus, + GatewayHotReloadStatus, +} from "./config-reload-status.types.js"; export type { GatewayReloadPlan } from "./config-reload-plan.js"; const MISSING_CONFIG_RETRY_DELAY_MS = 150; @@ -220,7 +223,7 @@ export function startGatewayConfigReloader(opts: { nextConfig: OpenClawConfig, ownership: GatewayConfigReloadTransactionOwnership, sourceConfig: OpenClawConfig, - ) => Promise; + ) => Promise; onRestart: ( plan: GatewayReloadPlan, nextConfig: OpenClawConfig, @@ -775,8 +778,9 @@ export function startGatewayConfigReloader(opts: { } await opts.onConfigChange?.(plan, nextConfig); + let applicationStatus: GatewayHotReloadApplicationStatus; try { - await opts.onHotReload(plan, nextConfig, ownership, nextSourceConfig); + applicationStatus = await opts.onHotReload(plan, nextConfig, ownership, nextSourceConfig); } catch (error) { ownership.rollbackRuntimeEnv(); throw error; @@ -784,7 +788,7 @@ export function startGatewayConfigReloader(opts: { assertCurrent(); await appliedRevision.apply(plan, nextConfig, nextConfigRevisionHash); await commitReloadBaseline(); - application?.settle("applied"); + application?.settle(applicationStatus); if (plan.reloadPlugins) { // The committed reload republished the metadata snapshot generation. markPluginMetadataRefreshApplied(); diff --git a/src/gateway/server-methods/config.test.ts b/src/gateway/server-methods/config.test.ts index 6c4d1d16ce00..cd58fb240c81 100644 --- a/src/gateway/server-methods/config.test.ts +++ b/src/gateway/server-methods/config.test.ts @@ -233,7 +233,7 @@ describe("config.patch application settlement", () => { ); }); - it.each(["superseded", "failed", "unclaimed"] as const)( + it.each(["restart-required", "superseded", "failed", "unclaimed"] as const)( "reports a persisted write whose runtime application was %s", async (outcome) => { const queueFollowUp = vi.fn(); @@ -259,6 +259,15 @@ describe("config.patch application settlement", () => { message: expect.stringContaining("persisted but was not applied"), }), ); + expect(harness.respond).toHaveBeenCalledWith( + false, + undefined, + expect.objectContaining({ + message: expect.stringContaining( + outcome === "restart-required" ? "Gateway is restarting" : "use config.apply", + ), + }), + ); expect(queueFollowUp).toHaveBeenCalledOnce(); }, ); diff --git a/src/gateway/server-methods/config.ts b/src/gateway/server-methods/config.ts index efc9651a1643..276e5cc21b2e 100644 --- a/src/gateway/server-methods/config.ts +++ b/src/gateway/server-methods/config.ts @@ -695,12 +695,16 @@ async function respondWithConfigRestartWrite(params: { if (params.writeResult.application) { const outcome = await params.writeResult.application; if (outcome !== "applied") { + const recovery = + outcome === "restart-required" + ? "the Gateway is restarting to finish applying it" + : "run config.get, then use config.apply to reapply the saved config or restart the Gateway"; params.respond( false, undefined, errorShape( ErrorCodes.UNAVAILABLE, - `${params.mode} persisted but was not applied to the active Gateway (${outcome}); run config.get, then retry or restart the Gateway`, + `${params.mode} persisted but was not applied to the active Gateway (${outcome}); ${recovery}`, ), ); params.writeResult.queueFollowUp(); diff --git a/src/gateway/server-reload-handlers.test.ts b/src/gateway/server-reload-handlers.test.ts index ca95bfde5179..fba2156e293f 100644 --- a/src/gateway/server-reload-handlers.test.ts +++ b/src/gateway/server-reload-handlers.test.ts @@ -82,6 +82,7 @@ import { type GatewayReloadPlan, } from "./config-reload-plan.js"; import { shouldRewarmProviderAuthState } from "./config-reload-recovery.js"; +import type { GatewayHotReloadApplicationStatus } from "./config-reload-status.types.js"; import { applyHookMappings } from "./hooks-mapping.js"; import { commitHooksConfigReload } from "./hooks.js"; import { createLazyGatewayCronState } from "./server-cron-lazy.js"; @@ -1360,7 +1361,7 @@ describe("gateway hot reload model state", () => { expect(abortSignal.aborted).toBe(false); releaseReconciliation.resolve(); if (publishes) { - await expect(reload).resolves.toBeUndefined(); + await expect(reload).resolves.toBe("applied"); } else { await expect(reload).rejects.toThrow("cron monitor"); } @@ -1675,7 +1676,7 @@ describe("gateway hot reload model state", () => { expect(setState).not.toHaveBeenCalled(); releaseReconciliation(); - await expect(reload).resolves.toBeUndefined(); + await expect(reload).resolves.toBe("applied"); expect(heartbeatRunner.updateConfig).toHaveBeenCalledWith(nextConfig); expect(setState).toHaveBeenCalledOnce(); }); @@ -1716,7 +1717,7 @@ describe("gateway hot reload model state", () => { isCurrent: () => true, publish: async (commit) => await commit(), }), - ).resolves.toBeUndefined(); + ).resolves.toBe("applied"); await waitForFast(() => expect(reconcileHeartbeatJobs).toHaveBeenCalledWith(nextConfig)); expect(heartbeatRunner.updateConfig).not.toHaveBeenCalled(); @@ -1786,7 +1787,7 @@ describe("gateway hot reload model state", () => { await expect( applyHotReload(createCronRestartPlan(), { cron: { enabled: true } }), - ).resolves.toBeUndefined(); + ).resolves.toBe("restart-required"); expect(setState).toHaveBeenCalledOnce(); await waitForFast(() => expect(signalSpy).toHaveBeenCalledOnce()); @@ -1864,7 +1865,7 @@ describe("gateway hot reload model state", () => { isCurrent: () => true, }, ), - ).resolves.toBeUndefined(); + ).resolves.toBe("restart-required"); expect(publish).toHaveBeenCalledOnce(); expect(setState).toHaveBeenCalledOnce(); @@ -2352,7 +2353,7 @@ describe("gateway hot reload superseded tail recovery", () => { await entered.promise; pendingConfig = invalidConfigB; release.resolve(); - await expect(reloadA).resolves.toBeUndefined(); + await expect(reloadA).resolves.toBe("applied"); expect(requestRecoveryRestart).not.toHaveBeenCalled(); expect(logReload.warn).toHaveBeenCalledWith( @@ -3378,10 +3379,10 @@ describe("gateway channel hot reload handlers", () => { }); } - async function withChannelReloadsEnabled(run: () => Promise) { + async function withChannelReloadsEnabled(run: () => Promise): Promise { const restoreChannelReloadEnv = enableChannelReloadsForTest(); try { - await run(); + return await run(); } finally { restoreChannelReloadEnv(); } @@ -3445,7 +3446,7 @@ describe("gateway channel hot reload handlers", () => { const { applyHotReload } = createReloadHandlersForTest(undefined, channels); const root = tryBeginGatewayRootWorkAdmission(); expect(root).not.toBeNull(); - let reload: Promise | undefined; + let reload: Promise | undefined; try { await root?.run(async () => { @@ -3669,7 +3670,7 @@ describe("gateway channel hot reload handlers", () => { const logReload = { info: vi.fn(), warn: vi.fn() }; const { applyHotReload } = createReloadHandlersForTest(logReload, channels, reloadPlugins); vi.useFakeTimers(); - let reload: Promise | undefined; + let reload: Promise | undefined; try { await withChannelReloadsEnabled(async () => { @@ -3811,7 +3812,7 @@ describe("gateway channel hot reload handlers", () => { await withChannelReloadsEnabled(async () => { await expect( applyHotReload(createChannelReloadPlan(["telegram", "discord"]), {}), - ).resolves.toBeUndefined(); + ).resolves.toBe("restart-required"); }); expect(signalSpy).toHaveBeenCalledOnce(); }); @@ -3876,7 +3877,7 @@ describe("gateway Gmail hot reload handlers", () => { await expect( applyHotReload(createGmailReloadPlan(), createGmailConfig("next@example.com")), - ).resolves.toBeUndefined(); + ).resolves.toBe("restart-required"); expect(stopPostReadySidecars).toHaveBeenCalledOnce(); expect(setState).toHaveBeenCalledOnce(); @@ -5822,7 +5823,7 @@ describe("gateway plugin hot reload handlers", () => { isCurrent: () => true, }, ), - ).resolves.toBeUndefined(); + ).resolves.toBe("restart-required"); expect(handlers.setState).toHaveBeenCalledTimes(1); expect(logReload.warn).toHaveBeenCalledWith( @@ -5853,7 +5854,7 @@ describe("gateway plugin hot reload handlers", () => { { plugins: { enabled: true } }, { sourceConfig: { plugins: { enabled: true } }, publish, isCurrent: () => true }, ), - ).resolves.toBeUndefined(); + ).resolves.toBe("restart-required"); expect(publish).toHaveBeenCalledOnce(); expect(handlers.setState).toHaveBeenCalledTimes(1); @@ -6526,7 +6527,7 @@ describe("deferred channel reload abort generation", () => { // Drain active work → should proceed to stop/start channels normally hoisted.activeTaskBlockers.length = 0; await vi.advanceTimersByTimeAsync(500); // wake up, see active=0, drain complete - await expect(reloadPromise).resolves.toBeUndefined(); + await expect(reloadPromise).resolves.toBe("applied"); expect(channels.stop).toHaveBeenCalledWith("whatsapp", undefined, { manual: false }); expect(channels.start).toHaveBeenCalledWith("whatsapp"); diff --git a/src/gateway/server-reload-hot.ts b/src/gateway/server-reload-hot.ts index 9631ca317ed2..2eff3afbe104 100644 --- a/src/gateway/server-reload-hot.ts +++ b/src/gateway/server-reload-hot.ts @@ -96,7 +96,7 @@ export function createGatewayReloadHandlers(params: GatewayReloadHandlerParams) plan: GatewayReloadPlan, nextConfig: OpenClawConfig, publication?: GatewayHotReloadPublication, - ): Promise => { + ) => { assertIrreversibleReloadPlanHasRecoveryOwner(plan, restartRecoveryAvailable); const isTransactionCurrent = () => !isRestartRetryStopped() && (publication?.isCurrent?.() ?? true); @@ -166,8 +166,7 @@ export function createGatewayReloadHandlers(params: GatewayReloadHandlerParams) let recoveryRestartScheduled = false; const laneConcurrency = resolveGatewayLaneConcurrency(nextConfig); const candidateEnv = publication?.runtimeEnv ?? process.env; - // Planning happens before candidate env publication, while channel starts - // happen after it. Use one candidate snapshot across both phases. + // Use one candidate env snapshot before publication and through later channel starts. const shouldSkipChannelRestart = isTruthyEnvValue(candidateEnv.OPENCLAW_SKIP_CHANNELS) || isTruthyEnvValue(candidateEnv.OPENCLAW_SKIP_PROVIDERS); @@ -335,8 +334,7 @@ export function createGatewayReloadHandlers(params: GatewayReloadHandlerParams) return; } try { - // Reuse the config-restart path: it excludes this reload root while - // draining other work and fences signal delivery until restart takes over. + // Reuse the config-restart path to drain other work and fence restart delivery. const restartTransaction = requestGatewayRestart( recoveryPlan, nextConfig, @@ -351,8 +349,7 @@ export function createGatewayReloadHandlers(params: GatewayReloadHandlerParams) }, ); settleRecoveryRestart(restartTransaction, surface); - // Immediate emission failure already owns a lifecycle retry. The runtime - // is committed, so keep this transaction accepted while that retry runs. + // Keep the committed transaction accepted while emission recovery retries. } catch (restartError) { params.logReload.warn( `failed to schedule post-commit gateway restart: ${formatErrorMessage(restartError)}`, @@ -553,7 +550,7 @@ export function createGatewayReloadHandlers(params: GatewayReloadHandlerParams) throw err; } scheduleRecoveryRestart("plugin runtime reload", err); - return; + return "restart-required"; } if (pluginReloadResult.cancelled) { pluginReloadAborted = true; @@ -609,7 +606,7 @@ export function createGatewayReloadHandlers(params: GatewayReloadHandlerParams) throw err; } scheduleRecoveryRestart("runtime commit", err); - return; + return "restart-required"; } try { @@ -621,7 +618,7 @@ export function createGatewayReloadHandlers(params: GatewayReloadHandlerParams) }); } catch (err) { scheduleRecoveryRestart("prepared model runtime reload", err); - return; + return "restart-required"; } if (plan.restartHealthMonitor) { @@ -720,6 +717,7 @@ export function createGatewayReloadHandlers(params: GatewayReloadHandlerParams) } else if (plan.noopPaths.length > 0) { params.logReload.info(`config change applied (dynamic reads: ${plan.noopPaths.join(", ")})`); } + return recoveryRestartScheduled ? "restart-required" : "applied"; }; return { diff --git a/src/gateway/server-reload-managed-secrets.ts b/src/gateway/server-reload-managed-secrets.ts index 74476a5833d6..7331cfa879da 100644 --- a/src/gateway/server-reload-managed-secrets.ts +++ b/src/gateway/server-reload-managed-secrets.ts @@ -374,8 +374,9 @@ export function createManagedReloadSecretHandlers(options: { let publishedSharedGatewaySessionGeneration: SharedGatewaySessionGenerationOwnership | null = null; let terminalConfigReconciled = false; + let applicationStatus: Awaited>; try { - await applyHotReload(plan, prepared.config, { + applicationStatus = await applyHotReload(plan, prepared.config, { isCurrent: transactionOwnership.isCurrent, ...(transactionOwnership.runtimeEnv ? { runtimeEnv: transactionOwnership.runtimeEnv.env } @@ -605,7 +606,7 @@ export function createManagedReloadSecretHandlers(options: { publishedSharedGatewaySessionGeneration, ); } - return; + return applicationStatus; } };