diff --git a/docs/gateway/configuration.md b/docs/gateway/configuration.md index f38e0ed5b174..497e6844560a 100644 --- a/docs/gateway/configuration.md +++ b/docs/gateway/configuration.md @@ -630,6 +630,11 @@ through replay; persistence alone is not an application acknowledgment. Shutdown supersession by different content, or failed application returns `UNAVAILABLE` with recovery guidance. `config.set` acknowledges persistence only. +Once a reload has committed, it finishes its model and channel work before a +newer config is applied. If that work needs restart recovery, the RPC returns +`UNAVAILABLE`; wait for the Gateway to restart, then use `config.get` to verify +the active revision. + `config.patch` also accepts `replacePaths`, an array of config paths whose array replacement is intentional. If a patch would replace or delete an existing array with fewer entries, the Gateway rejects the write unless that exact path appears diff --git a/qa/scenarios/memory/memory-dreaming-sweep.yaml b/qa/scenarios/memory/memory-dreaming-sweep.yaml index 43895fb619f3..c20e652acdb2 100644 --- a/qa/scenarios/memory/memory-dreaming-sweep.yaml +++ b/qa/scenarios/memory/memory-dreaming-sweep.yaml @@ -167,6 +167,8 @@ flow: - set: transcriptPath value: expr: "path.join(sessionsDir, `${config.transcriptId}.jsonl`)" + - set: recallEvidence + value: [] - try: actions: - call: fs.mkdir @@ -228,6 +230,9 @@ flow: expr: "JSON.stringify(payload.results ?? []).includes(config.expectedNeedle)" message: expr: "`memory search missed dreaming canary for query: ${query}`" + - set: recallEvidence + value: + expr: "[...recallEvidence, { query, results: (payload.results ?? []).map(({ path, snippet }) => ({ path, snippet })) }]" - set: cronRunStartedAt value: expr: "Date.now()" @@ -269,6 +274,7 @@ flow: - 1000 finally: - call: patchConfig + saveAs: restoreResult args: - env: ref: env @@ -286,4 +292,30 @@ flow: args: - ref: env - 60000 - detailsExpr: "JSON.stringify({ promotedTotal: promoted.status.dreaming?.promotedTotal ?? 0, shortTermCount: promoted.status.dreaming?.shortTermCount ?? 0, phaseSignalCount: promoted.status.dreaming?.phaseSignalCount ?? 0, lightSleep: promoted.lightReport.includes('# Light Sleep'), remSleep: promoted.remReport.includes('# REM Sleep'), promotionMarkerPresent: promoted.promotedMemory.includes(promoted.promotionMarker) })" + detailsExpr: "JSON.stringify({ promotedTotal: promoted.status.dreaming?.promotedTotal ?? 0, shortTermCount: promoted.status.dreaming?.shortTermCount ?? 0, phaseSignalCount: promoted.status.dreaming?.phaseSignalCount ?? 0, lightSleep: promoted.lightReport.includes('# Light Sleep'), remSleep: promoted.remReport.includes('# REM Sleep'), promotionMarkerPresent: promoted.promotedMemory.includes(promoted.promotionMarker), recallEvidence, cronRun: { id: cronId, status: finishedRun.status, ts: finishedRun.ts } })" + + - name: reports restored configuration and runtime readiness + actions: + - call: readConfigSnapshot + saveAs: restored + args: + - ref: env + - call: fetchJson + saveAs: readyz + args: + - expr: "`${env.gateway.baseUrl}/readyz`" + - call: env.gateway.call + saveAs: restoredChannels + args: + - channels.status + - probe: false + timeoutMs: 10000 + - timeoutMs: 15000 + detailsExpr: >- + JSON.stringify({ dreamingBefore: dreamingOriginal ?? null, + dreamingAfter: restored.config.plugins?.entries?.['memory-core']?.config?.dreaming ?? null, + restoreReceipt: { ok: restoreResult?.ok ?? null, noop: restoreResult?.noop ?? false, + restarted: restoreResult?.restarted ?? false }, readyz, + channels: Object.fromEntries(Object.entries(restoredChannels.channelAccounts ?? {}).map( + ([channel, accounts]) => [channel, accounts.map(({ accountId, running, connected, + lifecycle, restartPending }) => ({ accountId, running, connected, lifecycle, restartPending }))])) }) diff --git a/src/gateway/config-reload.ts b/src/gateway/config-reload.ts index e1a5aed229e0..735d773d99d5 100644 --- a/src/gateway/config-reload.ts +++ b/src/gateway/config-reload.ts @@ -171,6 +171,8 @@ export function startGatewayConfigReloader(opts: { onRuntimeConfigCommitted?: (plan: GatewayReloadPlan, nextConfig: OpenClawConfig) => void; /** Publishes the resolved source-config revision accepted by the active runtime. */ onConfigRevisionApplied?: (hash: string) => void; + /** Reads the same restart owner that fences publication of the applied revision. */ + hasOutstandingGatewayRestart?: () => boolean; /** Retires rejected lifecycle work after any newer config transaction is accepted. */ onConfigAccepted?: ( nextConfig: OpenClawConfig, @@ -464,6 +466,12 @@ export function startGatewayConfigReloader(opts: { authoredConfig?: unknown, application?: RuntimeConfigWriteApplicationClaim, ) => { + const settleRuntimeApplication = (status: GatewayHotReloadApplicationStatus = "applied") => { + // A watcher replay must not turn recovery-owned runtime work into a success receipt. + application?.settle( + opts.hasOutstandingGatewayRestart?.() ? "applied-restart-required" : status, + ); + }; // Reprepare against the current accepted env owner. A managed write can // finish preflight while another watcher transaction accepts first. const preparedCandidate = @@ -690,7 +698,7 @@ export function startGatewayConfigReloader(opts: { publishedSource?.commit?.(); } opts.onConfigRevisionApplied?.(nextConfigRevisionHash); - application?.settle("applied"); + settleRuntimeApplication(); return; } @@ -741,7 +749,7 @@ export function startGatewayConfigReloader(opts: { assertCurrent(); await appliedRevision.apply(plan, nextConfig, nextConfigRevisionHash); await commitReloadBaseline(); - application?.settle("applied"); + settleRuntimeApplication(); return; } if (followUp.requiresRestart) { @@ -778,7 +786,7 @@ export function startGatewayConfigReloader(opts: { assertCurrent(); await appliedRevision.apply(plan, nextConfig, nextConfigRevisionHash); await commitReloadBaseline(); - application?.settle(applicationStatus); + settleRuntimeApplication(applicationStatus); if (plan.reloadPlugins) { // The committed reload republished the metadata snapshot generation. markPluginMetadataRefreshApplied(); diff --git a/src/gateway/server-reload-channel-restart.ts b/src/gateway/server-reload-channel-restart.ts index a92f309d6d87..158e1b400c30 100644 --- a/src/gateway/server-reload-channel-restart.ts +++ b/src/gateway/server-reload-channel-restart.ts @@ -4,12 +4,11 @@ import { formatErrorMessage } from "../infra/errors.js"; import { requireActivePluginChannelRegistry } from "../plugins/runtime.js"; import { withPluginRuntimeRegistryScope } from "../plugins/runtime/gateway-request-scope.js"; import { runOutsideGatewayRootWorkAdmission } from "../process/gateway-work-admission.js"; -import type { ChannelKind } from "./config-reload-plan.js"; -import type { GatewayReloadPlan } from "./config-reload.js"; +import type { ChannelKind, GatewayReloadPlan } from "./config-reload-plan.js"; import type { GatewayReloadHandlerParams } from "./server-reload-contracts.js"; import { collectChannelOperationFailures } from "./server-reload-utils.js"; -export function startGatewayChannelFromActiveRegistry( +function startGatewayChannelFromActiveRegistry( params: Pick, channel: ChannelKind, accountId?: string, @@ -23,6 +22,47 @@ export function startGatewayChannelFromActiveRegistry( ); } +export async function rollbackStoppedGatewayChannels( + params: Pick, + channels: Set, + accounts: Map>, + reason: string, +): Promise { + const failures: string[] = []; + for (const [channel, accountIds] of accounts) { + for (const accountId of accountIds) { + try { + params.logChannels.info(`restarting ${channel} account ${accountId} after ${reason}`); + await startGatewayChannelFromActiveRegistry(params, channel, accountId); + accountIds.delete(accountId); + } catch (err) { + failures.push(`${channel}[${accountId}]`); + params.logChannels.error( + `failed to restart ${channel} account ${accountId} after ${reason}: ${formatErrorMessage(err)}`, + ); + } + } + if (accountIds.size === 0) { + accounts.delete(channel); + } + } + return failures.concat( + await collectChannelOperationFailures({ + channels: [...channels], + run: async (channel) => { + params.logChannels.info(`restarting ${channel} channel after ${reason}`); + await startGatewayChannelFromActiveRegistry(params, channel); + channels.delete(channel); + }, + onFailure: (channel, err) => { + params.logChannels.error( + `failed to restart ${channel} channel after ${reason}: ${formatErrorMessage(err)}`, + ); + }, + }), + ); +} + export async function restartGatewayChannels(options: { params: GatewayReloadHandlerParams; plan: GatewayReloadPlan; @@ -34,7 +74,6 @@ export async function restartGatewayChannels(options: { accountsStoppedBeforePluginReload: ReadonlyMap>; shouldSkipChannelRestart: boolean; skipChannelRestartLogMessage: string; - pluginReloadAborted: boolean; isLifecycleReloadAborted: () => boolean; getChannelAutostartSuppression: () => unknown; channelReloadTargets: () => Set; @@ -52,7 +91,6 @@ export async function restartGatewayChannels(options: { accountsStoppedBeforePluginReload, shouldSkipChannelRestart, skipChannelRestartLogMessage, - pluginReloadAborted, isLifecycleReloadAborted, getChannelAutostartSuppression, channelReloadTargets, @@ -102,108 +140,67 @@ export async function restartGatewayChannels(options: { return targets; }; - if (channelsToRestart.size > 0 || restartChannelAccounts.size > 0) { - if (shouldSkipChannelRestart) { - params.logChannels.info(skipChannelRestartLogMessage); - } else if (getChannelAutostartSuppression()) { - const cancelledByRestart = pluginReloadAborted; - if (cancelledByRestart) { - params.logChannels.info("channel restart cancelled by in-process restart"); - } else { - const accountStops = collectChannelAccountTargets(); - const accountStopFailures: string[] = []; - for (const [channel, accountId] of accountStops) { - try { - params.logChannels.info( - `stopping ${channel} account ${accountId} before suppressed hot reload`, - ); - if (!wasStoppedBeforePluginReload(channel, accountId)) { - await params.stopChannel(channel, accountId, { manual: false }); - } - } catch (err) { - accountStopFailures.push(`${channel}[${accountId}]`); - params.logChannels.error( - `failed to stop ${channel} account ${accountId} during suppressed hot reload: ${formatErrorMessage(err)}`, - ); - } - } - const stopFailures = await collectChannelOperationFailures({ - channels: channelsToRestart, - run: async (channel) => { - if (plan.reloadPlugins && activePluginChannelsAfterReload?.has(channel) === false) { - return; - } - if (channelsStoppedBeforePluginReload.has(channel)) { - return; - } - params.logChannels.info(`stopping ${channel} channel before suppressed hot reload`); - await params.stopChannel(channel, undefined, { manual: false }); - }, - onFailure: (channel, err) => { - params.logChannels.error( - `failed to stop ${channel} channel during suppressed hot reload: ${formatErrorMessage( - err, - )}`, - ); - }, - }); - const allStopFailures = [...accountStopFailures, ...stopFailures]; - if (allStopFailures.length > 0) { - scheduleRecoveryRestart(`channel stop (${allStopFailures.join(", ")})`); - } - logSuppressedChannelRestart(channelReloadTargets(), "channel restart during hot reload"); + if (channelsToRestart.size === 0 && restartChannelAccounts.size === 0) { + return; + } + if (shouldSkipChannelRestart) { + params.logChannels.info(skipChannelRestartLogMessage); + return; + } + const suppressed = Boolean(getChannelAutostartSuppression()); + const operation = suppressed ? "stop" : "restart"; + const phase = suppressed ? "suppressed hot reload" : "hot reload"; + const accountTargets = collectChannelAccountTargets(); + const accountFailures: string[] = []; + for (const [channel, accountId] of accountTargets) { + try { + params.logChannels.info( + suppressed + ? `stopping ${channel} account ${accountId} before suppressed hot reload` + : `restarting ${channel} account ${accountId}`, + ); + if (!wasStoppedBeforePluginReload(channel, accountId)) { + await params.stopChannel(channel, accountId, { manual: false }); } - } else { - const cancelledByRestart = pluginReloadAborted; - if (cancelledByRestart) { - params.logChannels.info("channel restart cancelled by in-process restart"); - } else { - const accountRestarts = collectChannelAccountTargets(); - const accountRestartFailures: string[] = []; - for (const [channel, accountId] of accountRestarts) { - try { - params.logChannels.info(`restarting ${channel} account ${accountId}`); - if (!wasStoppedBeforePluginReload(channel, accountId)) { - await params.stopChannel(channel, accountId, { manual: false }); - } - if (isLifecycleReloadAborted()) { - continue; - } - await startGatewayChannelFromActiveRegistry(params, channel, accountId); - } catch (err) { - accountRestartFailures.push(`${channel}[${accountId}]`); - params.logChannels.error( - `failed to restart ${channel} account ${accountId} during hot reload: ${formatErrorMessage(err)}`, - ); - } - } - const restartChannel = async (name: ChannelKind) => { - if (plan.reloadPlugins && activePluginChannelsAfterReload?.has(name) === false) { - return; - } - params.logChannels.info(`restarting ${name} channel`); - if (!channelsStoppedBeforePluginReload.has(name)) { - await params.stopChannel(name, undefined, { manual: false }); - } - if (isLifecycleReloadAborted()) { - return; - } - await startGatewayChannelFromActiveRegistry(params, name); - }; - const restartFailures = await collectChannelOperationFailures({ - channels: channelsToRestart, - run: restartChannel, - onFailure: (channel, err) => { - params.logChannels.error( - `failed to restart ${channel} channel during hot reload: ${formatErrorMessage(err)}`, - ); - }, - }); - const allRestartFailures = [...accountRestartFailures, ...restartFailures]; - if (allRestartFailures.length > 0) { - scheduleRecoveryRestart(`channel restart (${allRestartFailures.join(", ")})`); - } + if (!suppressed && !isLifecycleReloadAborted()) { + await startGatewayChannelFromActiveRegistry(params, channel, accountId); } + } catch (err) { + accountFailures.push(`${channel}[${accountId}]`); + params.logChannels.error( + `failed to ${operation} ${channel} account ${accountId} during ${phase}: ${formatErrorMessage(err)}`, + ); } } + const channelFailures = await collectChannelOperationFailures({ + channels: channelsToRestart, + run: async (channel) => { + if (plan.reloadPlugins && activePluginChannelsAfterReload?.has(channel) === false) { + return; + } + params.logChannels.info( + suppressed + ? `stopping ${channel} channel before suppressed hot reload` + : `restarting ${channel} channel`, + ); + if (!channelsStoppedBeforePluginReload.has(channel)) { + await params.stopChannel(channel, undefined, { manual: false }); + } + if (!suppressed && !isLifecycleReloadAborted()) { + await startGatewayChannelFromActiveRegistry(params, channel); + } + }, + onFailure: (channel, err) => { + params.logChannels.error( + `failed to ${operation} ${channel} channel during ${phase}: ${formatErrorMessage(err)}`, + ); + }, + }); + const failures = [...accountFailures, ...channelFailures]; + if (failures.length > 0) { + scheduleRecoveryRestart(`channel ${operation} (${failures.join(", ")})`); + } + if (suppressed) { + logSuppressedChannelRestart(channelReloadTargets(), "channel restart during hot reload"); + } } diff --git a/src/gateway/server-reload-handlers.test.ts b/src/gateway/server-reload-handlers.test.ts index 41684bb3411a..01e152fed7f6 100644 --- a/src/gateway/server-reload-handlers.test.ts +++ b/src/gateway/server-reload-handlers.test.ts @@ -29,6 +29,7 @@ import { import { attachRuntimeConfigWriteApplication, createRuntimeConfigWriteApplication, + type RuntimeConfigWriteApplicationStatus, } from "../config/runtime-write-application.js"; import type { OpenClawConfig } from "../config/types.openclaw.js"; import { CronService } from "../cron/service.js"; @@ -60,6 +61,7 @@ import { setCommandLaneConcurrency, } from "../process/command-queue.js"; import { + captureGatewayRootWorkAdmissionContinuationScope, getActiveGatewayRootWorkCount, isGatewayWorkAdmissionClosed, resetGatewayWorkAdmission, @@ -6433,9 +6435,14 @@ describe("deferred channel reload abort generation", () => { ); }); - it.each([false, true])( - "cancels superseded active-work deferral (runtime committed: %s)", - async (committed) => { + it.each([ + [false, "config"], + [true, "config"], + [false, "lifecycle"], + [true, "lifecycle"], + ] as const)( + "settles active-work deferral (runtime committed: %s, cancellation: %s)", + async (committed, cancellationKind) => { const logChannels = { info: vi.fn(), error: vi.fn() }; const channels = { start: vi.fn(async () => {}), @@ -6443,7 +6450,10 @@ describe("deferred channel reload abort generation", () => { }; const reloadPlugins: ReloadHandlerParams["reloadPlugins"] = async (params) => { await params.commitRuntime(); - return makePluginReloadResult({ restartChannels: new Set(["whatsapp"]) }); + return makePluginReloadResult({ + restartChannels: new Set(["whatsapp"]), + activeChannels: new Set(["whatsapp"]), + }); }; const { applyHotReload } = createTestHandlers(logChannels, channels, { reloadPlugins }); hoisted.activeTaskBlockers.push( @@ -6462,18 +6472,31 @@ describe("deferred channel reload abort generation", () => { publish: async (commit) => await commit(), }, ); - const cancellation = committed - ? GatewayHotReloadCancelledError - : GatewayConfigReloadSupersededError; + const cancellation = + committed || cancellationKind === "lifecycle" + ? GatewayHotReloadCancelledError + : GatewayConfigReloadSupersededError; const reloadError = reloadPromise.then( () => null, (error: unknown) => error, ); await vi.advanceTimersByTimeAsync(10); - transactionCurrent = false; + if (cancellationKind === "lifecycle") { + abortPendingChannelReloads(); + } else { + transactionCurrent = false; + } await vi.advanceTimersByTimeAsync(500); const error = await reloadError; + if (committed && cancellationKind === "config") { + expect(error).toBeNull(); + expect(channels.stop).toHaveBeenCalledOnce(); + expect(channels.start).toHaveBeenCalledOnce(); + expect(hoisted.refreshPreparedModelRuntimeSnapshots).toHaveBeenCalledOnce(); + expect(hoisted.rejectPendingPreparedModelRuntimeReplacement).not.toHaveBeenCalled(); + return; + } expect(error).toBeInstanceOf(cancellation); expect(channels.stop).not.toHaveBeenCalled(); @@ -6495,14 +6518,13 @@ describe("deferred channel reload abort generation", () => { }, ); - it.each([ - ...(["channel", "plugin"] as const).flatMap((surface) => + it.each( + (["channel", "plugin"] as const).flatMap((surface) => (["same write", "newer content", "lifecycle stop", "publication failure"] as const).map( (outcome) => ({ surface, outcome }), ), ), - { surface: "committed plugin", outcome: "same write" } as const, - ])( + )( "settles a deferred $surface receipt after watcher handoff: $outcome", async ({ surface, outcome }) => { const initialConfig = { @@ -6542,13 +6564,6 @@ describe("deferred channel reload abort generation", () => { ); }); const reloadPlugins = vi.fn(async (params) => { - if (surface === "committed plugin") { - await params.commitRuntime(); - return makePluginReloadResult({ - restartChannels: new Set(["whatsapp"]), - activeChannels: new Set(["whatsapp"]), - }); - } await params.beforeReplace(new Set(["whatsapp"])); if (params.isAborted?.()) { return makePluginReloadResult({ cancelled: true }); @@ -6609,17 +6624,6 @@ describe("deferred channel reload abort generation", () => { await vi.advanceTimersByTimeAsync(500); await vi.advanceTimersByTimeAsync(300); await snapshotStarted.promise; - if (surface === "committed plugin") { - // A no-op replay cannot finish the interrupted runtime tail after commit. - expect(settled).toHaveBeenCalledExactlyOnceWith("failed"); - expect(setState).toHaveBeenCalledOnce(); - expect(startChannel).not.toHaveBeenCalled(); - snapshotGate.resolve(); - unrelatedRequest.release(); - await vi.advanceTimersByTimeAsync(500); - await expect(application.result).resolves.toBe("failed"); - return; - } expect(settled).not.toHaveBeenCalled(); expect(setState).not.toHaveBeenCalled(); expect(stopChannel).not.toHaveBeenCalled(); @@ -6677,6 +6681,146 @@ describe("deferred channel reload abort generation", () => { }, ); + it.each(["same watcher echo", "newer admitted write", "prepared refresh failure"] as const)( + "finishes committed runtime work before settling its receipt: %s", + async (successor) => { + const initialConfig: OpenClawConfig = { + gateway: { reload: {} }, + channels: { whatsapp: { enabled: true, selfChatMode: false } }, + }; + const nextConfig: OpenClawConfig = { + ...initialConfig, + plugins: { entries: { fixture: { enabled: true } } }, + }; + activateSecretsRuntimeSnapshot(makePreparedSecretsSnapshot(initialConfig)); + setActivePluginRegistry( + createTestRegistry([ + { + pluginId: "whatsapp", + source: "test", + plugin: { + ...createChannelTestPluginBase({ id: "whatsapp" }), + reload: { + configPrefixes: ["channels.whatsapp.selfChatMode"], + noopPrefixes: ["channels.whatsapp"], + }, + }, + }, + ]), + ); + const watcher = new chokidar.FSWatcher(); + const watch = vi.spyOn(chokidar, "watch").mockReturnValue(watcher); + const writeListenerRef = createConfigWriteListenerRef(); + const channels = { start: vi.fn(async () => {}), stop: vi.fn(async () => {}) }; + const commitTerminalConfig = vi.fn(); + const logReload = { info: vi.fn(), warn: vi.fn(), error: vi.fn() }; + const watchedConfig = successor === "newer admitted write" ? initialConfig : nextConfig; + const continuePlugin = createDeferred(); + const requestRecoveryRestart = vi.fn(() => ({ status: "emitted" as const })); + let blocker: ReturnType = null; + let successorRequest: Promise | undefined; + const submitWrite = (config: OpenClawConfig, hash: string, revision: number) => + runWithGatewayIndependentRootWorkAdmission(async () => { + const application = createRuntimeConfigWriteApplication( + captureGatewayRootWorkAdmissionContinuationScope()?.run, + ); + writeListenerRef.current!( + attachRuntimeConfigWriteApplication( + createConfigWriteNotification(config, hash, revision, "runtime", "source"), + application, + ), + ); + return await application.result; + }); + const reloader = startManagedGatewayConfigReloader({ + initialConfig, + readSnapshot: vi.fn(async () => createValidConfigSnapshot(watchedConfig, "same-write")), + subscribeToWrites: captureConfigWriteListener(writeListenerRef), + startChannel: channels.start, + stopChannel: channels.stop, + reloadPlugins: async (params) => { + await params.beforeReplace(new Set(), new Map([["whatsapp", new Set(["default"])]])); + if (params.isAborted?.()) { + return makePluginReloadResult({ cancelled: true }); + } + await params.commitRuntime(); + if (params.sourceConfig === nextConfig) { + await continuePlugin.promise; + } + return makePluginReloadResult({ activeChannels: new Set(["whatsapp"]) }); + }, + commitTerminalConfig, + logReload, + requestRecoveryRestart, + }); + if (successor === "prepared refresh failure") { + hoisted.refreshPreparedModelRuntimeSnapshots.mockRejectedValueOnce( + new Error("prepared runtime refresh failed"), + ); + } + vi.useFakeTimers(); + const request = submitWrite(nextConfig, "same-write", 1); + + try { + await vi.advanceTimersByTimeAsync(10); + expect(channels.stop).toHaveBeenCalledOnce(); + blocker = tryBeginGatewayRootWorkAdmission(); + if (!blocker) { + throw new Error("Expected unrelated gateway request admission"); + } + continuePlugin.resolve(); + await vi.advanceTimersByTimeAsync(10); + expect(logReload.warn).toHaveBeenCalledWith(expect.stringContaining("deferring until")); + // A newer RPC remains admitted while waiting for the queued successor reload. + // The committed tail must not wait for that request to finish first. + if (successor === "newer admitted write") { + successorRequest = submitWrite(initialConfig, "newer-write", 2); + } else { + watcher.emit("change", "/tmp/openclaw.json"); + } + await vi.advanceTimersByTimeAsync(500); + blocker.release(); + await vi.advanceTimersByTimeAsync(1_000); + await expect(request).resolves.toBe( + successor === "prepared refresh failure" + ? "applied-restart-required" + : successor === "newer admitted write" + ? "superseded" + : "applied", + ); + if (successor === "prepared refresh failure") { + expect(channels.stop).toHaveBeenCalledOnce(); + expect(channels.start).not.toHaveBeenCalled(); + expect(requestRecoveryRestart).toHaveBeenCalledOnce(); + expect(hoisted.rejectPendingPreparedModelRuntimeReplacement).toHaveBeenCalledOnce(); + } else { + const reloadCount = successorRequest ? 2 : 1; + expect(channels.stop).toHaveBeenCalledTimes(reloadCount); + expect(channels.start).toHaveBeenCalledTimes(reloadCount); + expect(hoisted.refreshPreparedModelRuntimeSnapshots).toHaveBeenCalledWith( + nextConfig, + expect.anything(), + ); + expect(hoisted.rejectPendingPreparedModelRuntimeReplacement).not.toHaveBeenCalled(); + expect(commitTerminalConfig).toHaveBeenCalledWith(nextConfig); + if (successorRequest) { + await expect(successorRequest).resolves.toBe("applied"); + expect(commitTerminalConfig).toHaveBeenLastCalledWith(initialConfig); + } + } + expect(getActiveSecretsRuntimeSnapshot()?.sourceConfig).toEqual(watchedConfig); + expect(logReload.error).not.toHaveBeenCalled(); + } finally { + continuePlugin.resolve(); + blocker?.release(); + await reloader.stop(); + await request; + await successorRequest; + watch.mockRestore(); + } + }, + ); + it("new reload lifecycle is not affected by a previous lifecycle abort", async () => { const logChannels = { info: vi.fn(), error: vi.fn() }; const channels = { diff --git a/src/gateway/server-reload-hot.ts b/src/gateway/server-reload-hot.ts index 128a9f40905b..e85bf2fd3b92 100644 --- a/src/gateway/server-reload-hot.ts +++ b/src/gateway/server-reload-hot.ts @@ -12,19 +12,18 @@ import { isTruthyEnvValue } from "../infra/env.js"; import { formatErrorMessage } from "../infra/errors.js"; import { resetDirectoryCache } from "../infra/outbound/target-resolver.js"; import { setGatewaySigusr1RestartPolicy } from "../infra/restart.js"; -import type { ChannelKind } from "./config-reload-plan.js"; +import type { ChannelKind, GatewayReloadPlan } from "./config-reload-plan.js"; import { shouldRefreshContextWindowCache, shouldRewarmProviderAuthState, } from "./config-reload-recovery.js"; -import type { GatewayReloadPlan } from "./config-reload.js"; import { commitHooksConfigReload, resolveHooksConfig } from "./hooks.js"; import { buildGatewayCronService, type GatewayCronExitWatcherHandoff } from "./server-cron.js"; import { applyGatewayLaneConcurrency, resolveGatewayLaneConcurrency } from "./server-lanes.js"; import { createGatewayActiveWorkTracker } from "./server-reload-active-work.js"; import { restartGatewayChannels, - startGatewayChannelFromActiveRegistry, + rollbackStoppedGatewayChannels, } from "./server-reload-channel-restart.js"; import { assertReloadPublicationCurrent, @@ -365,45 +364,13 @@ export function createGatewayReloadHandlers(params: GatewayReloadHandlerParams) } }; if (plan.reloadPlugins) { - const restartStoppedPluginAccounts = async (reason: string): Promise => { - const failures: string[] = []; - for (const [channel, accountIds] of accountsStoppedBeforePluginReload) { - for (const accountId of accountIds) { - try { - params.logChannels.info(`restarting ${channel} account ${accountId} after ${reason}`); - await startGatewayChannelFromActiveRegistry(params, channel, accountId); - accountIds.delete(accountId); - } catch (err) { - failures.push(`${channel}[${accountId}]`); - params.logChannels.error( - `failed to restart ${channel} account ${accountId} after ${reason}: ${formatErrorMessage(err)}`, - ); - } - } - if (accountIds.size === 0) { - accountsStoppedBeforePluginReload.delete(channel); - } - } - return failures; - }; - const restartStoppedPluginChannels = async (reason: string) => - await collectChannelOperationFailures({ - channels: [...channelsStoppedBeforePluginReload], - run: async (channel) => { - params.logChannels.info(`restarting ${channel} channel after ${reason}`); - await startGatewayChannelFromActiveRegistry(params, channel); - channelsStoppedBeforePluginReload.delete(channel); - }, - onFailure: (channel, err) => { - params.logChannels.error( - `failed to restart ${channel} channel after ${reason}: ${formatErrorMessage(err)}`, - ); - }, - }); - const rollbackStoppedPluginTargets = async (reason: string): Promise => [ - ...(await restartStoppedPluginAccounts(reason)), - ...(await restartStoppedPluginChannels(reason)), - ]; + const rollbackStoppedPluginTargets = (reason: string) => + rollbackStoppedGatewayChannels( + params, + channelsStoppedBeforePluginReload, + accountsStoppedBeforePluginReload, + reason, + ); const failPluginChannelRollback = (reason: string, failures: string[]): never => { const error = new Error( `plugin reload cancellation rollback failed for: ${failures.join(", ")}`, @@ -590,7 +557,12 @@ export function createGatewayReloadHandlers(params: GatewayReloadHandlerParams) // Plugin replacement can admit new agent work while an account monitor stays live. // Recheck that work here; durable ingress replay remains owned by the fresh monitor drain. if (!pluginReloadAborted && hasLiveChannelTargets && !shouldSkipChannelRestart) { - pluginReloadAborted = await waitForActiveWorkBeforeChannelReload(channelTargets, isCurrent); + const waitCancelled = await waitForActiveWorkBeforeChannelReload(channelTargets, isCurrent); + // A committed owner must finish its model/channel tail before the next config runs. + // Supersession ends this wait: a newer writer may itself be awaiting that next reload. + pluginReloadAborted = + waitCancelled && + (!runtimeCommitted || isRestartRetryStopped() || isLifecycleReloadAborted()); } if (pluginReloadAborted) { // Only an uncommitted reload can transfer its receipt to the watcher. After @@ -690,7 +662,6 @@ export function createGatewayReloadHandlers(params: GatewayReloadHandlerParams) shouldSkipChannelRestart, skipChannelRestartLogMessage: "skipping channel reload (OPENCLAW_SKIP_CHANNELS=1 or OPENCLAW_SKIP_PROVIDERS=1)", - pluginReloadAborted, isLifecycleReloadAborted, getChannelAutostartSuppression, channelReloadTargets, diff --git a/src/gateway/server-reload-managed.ts b/src/gateway/server-reload-managed.ts index ac96c8198756..64fadc0bcb33 100644 --- a/src/gateway/server-reload-managed.ts +++ b/src/gateway/server-reload-managed.ts @@ -133,6 +133,7 @@ export function startManagedGatewayConfigReloader( applyHotReload, acceptRestartConfig, beginGatewayRestartLifecycle, + hasOutstandingGatewayRestart, pauseGatewayRestartForConfigCandidate, publishAppliedConfigHash, publishAcceptedRestartTarget, @@ -489,6 +490,7 @@ export function startManagedGatewayConfigReloader( params.commitTerminalConfig(nextConfig); }, onConfigRevisionApplied: publishAppliedConfigHash, + hasOutstandingGatewayRestart, onEffectiveConfigUnchanged, onNoopConfigCommit: async (plan, nextConfig, ownership, sourceConfig) => { // Cleared per transaction so a rebuild can never inherit a config committed