From 2bd485825215d91a8de2db0ff7e5cdbb34a1e82e Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Mon, 27 Jul 2026 05:46:00 -0400 Subject: [PATCH] refactor(gateway): make reload restart transactional (#114368) --- src/gateway/server-reload-active-work.ts | 6 + src/gateway/server-reload-hot.ts | 2 + src/gateway/server-reload-restart.ts | 625 ++++++++++++----------- 3 files changed, 341 insertions(+), 292 deletions(-) diff --git a/src/gateway/server-reload-active-work.ts b/src/gateway/server-reload-active-work.ts index 53dee69690d7..67306c9a9a1e 100644 --- a/src/gateway/server-reload-active-work.ts +++ b/src/gateway/server-reload-active-work.ts @@ -74,6 +74,11 @@ export function createGatewayActiveWorkTracker(options: { const omitted = blockers.length - shown.length; return omitted > 0 ? `${shown.join("; ")}; +${omitted} more` : shown.join("; "); }; + const formatDeferredWorkStatus = (status: "active" | "still active") => { + const details = formatActiveDetails(getActiveCounts()).join(", "); + const taskBlockers = formatTaskBlockers(); + return `${details} ${status}${taskBlockers ? ` (${taskBlockers})` : ""}`; + }; const waitForActiveWorkBeforeChannelReload = async ( channels: Iterable, isTransactionCurrent: () => boolean, @@ -134,6 +139,7 @@ export function createGatewayActiveWorkTracker(options: { return { formatActiveDetails, + formatDeferredWorkStatus, formatTaskBlockers, getActiveCounts, waitForActiveWorkBeforeChannelReload, diff --git a/src/gateway/server-reload-hot.ts b/src/gateway/server-reload-hot.ts index d84d1b803d03..969732eeb2bc 100644 --- a/src/gateway/server-reload-hot.ts +++ b/src/gateway/server-reload-hot.ts @@ -52,6 +52,7 @@ export function createGatewayReloadHandlers(params: GatewayReloadHandlerParams) const { formatActiveDetails, + formatDeferredWorkStatus, formatTaskBlockers, getActiveCounts, waitForActiveWorkBeforeChannelReload, @@ -81,6 +82,7 @@ export function createGatewayReloadHandlers(params: GatewayReloadHandlerParams) restartRecoveryAvailable, getActiveCounts, formatActiveDetails, + formatDeferredWorkStatus, formatTaskBlockers, }); diff --git a/src/gateway/server-reload-restart.ts b/src/gateway/server-reload-restart.ts index 2c2339595333..39d8d19a4187 100644 --- a/src/gateway/server-reload-restart.ts +++ b/src/gateway/server-reload-restart.ts @@ -36,229 +36,143 @@ type GatewayActiveCounts = { totalActive: number; }; -export function createGatewayRestartCoordinator(coordinatorOptions: { +type RestartRequestDetails = { + plan: GatewayReloadPlan; + nextConfig: OpenClawConfig; + restartOwnedPaths: string[]; + retainDebtAcrossConfigChanges: boolean; +}; + +type GatewayRestartOperation = + | { kind: "idle" } + | { + kind: "lifecycle"; + transaction: { state: GatewayRestartTransactionState }; + } + | { + kind: "request"; + transaction: { state: GatewayRestartTransactionState }; + details: RestartRequestDetails; + emissionSettled: boolean; + }; + +type AcceptedRestartTargetState = + | { kind: "empty"; generation: number } + | { + kind: "candidate-pending"; + generation: number; + previousTarget: AcceptedRestartTarget | undefined; + } + | { + kind: "accepted"; + generation: number; + target: AcceptedRestartTarget; + }; + +type GatewayRestartCoordinatorOptions = { params: GatewayReloadHandlerParams; myGeneration: number; restartRecoveryAvailable: boolean; getActiveCounts: () => GatewayActiveCounts; formatActiveDetails: (counts: GatewayActiveCounts) => string[]; + formatDeferredWorkStatus: (status: "active" | "still active") => string; formatTaskBlockers: () => string | null; -}) { - const { - params, - myGeneration, - restartRecoveryAvailable, - getActiveCounts, - formatActiveDetails, - formatTaskBlockers, - } = coordinatorOptions; - let restartPending = false; - let restartRetryStopped = false; - let restartRetryTimer: ReturnType | null = null; - let restartDeferral: RestartDeferralHandle | null = null; - let restartRequestGeneration = 0; - let restartRequestTransaction: { state: GatewayRestartTransactionState } | null = null; +}; + +class GatewayRestartTransaction { + private restartPending = false; + private retryStopped = false; + private retryTimer: ReturnType | null = null; + private restartDeferral: RestartDeferralHandle | null = null; + private requestGeneration = 0; // onReady/onTimeout precede async restart preparation. Keep committed details // debt-eligible until the emitter confirms this generation won. - let restartEmissionSettled = false; - type RestartRequestDetails = { - plan: GatewayReloadPlan; - nextConfig: OpenClawConfig; - restartOwnedPaths: string[]; - retainDebtAcrossConfigChanges: boolean; - }; - let restartRequestDetails: RestartRequestDetails | null = null; - let pausedRestartDebt: RestartRequestDetails | null = null; + private operation: GatewayRestartOperation = { kind: "idle" }; + private pausedDebt: RestartRequestDetails | null = null; // Post-commit recovery is satisfied only by an accepted restart emission. // Keep it separate from config-owned debt that later baselines may retire. - let conservativeRestartDebt: RestartRequestDetails | null = null; - let latestAcceptedRestartTarget: AcceptedRestartTarget | null = null; - let acceptedRestartTargetGeneration = 0; - let configCandidatePending = false; + private conservativeDebt: RestartRequestDetails | null = null; + private acceptedTargetState: AcceptedRestartTargetState = { kind: "empty", generation: 0 }; - const recordAcceptedRestartTarget = (target: AcceptedRestartTarget) => { - const generation = ++acceptedRestartTargetGeneration; + readonly appliedConfigHashPublisher = createAppliedConfigHashPublisher({ + hasPendingRestart: () => + this.operation.kind === "request" || + this.pausedDebt !== null || + this.conservativeDebt !== null, + publish: setRuntimeConfigAppliedHash, + }); + + constructor(private readonly options: GatewayRestartCoordinatorOptions) {} + + readonly isStopped = () => this.retryStopped; + readonly hasPendingConfigCandidate = () => this.acceptedTargetState.kind === "candidate-pending"; + readonly hasOperation = () => this.operation.kind !== "idle"; + readonly getAcceptedTarget = (): AcceptedRestartTarget | null => + this.acceptedTargetState.kind === "accepted" ? this.acceptedTargetState.target : null; + + recordAcceptedTarget(target: AcceptedRestartTarget): AcceptedRestartTargetOwnership { + const generation = this.acceptedTargetState.generation + 1; const acceptedTarget: AcceptedRestartTarget = { ...target, prepareRuntimeConfig: async () => { - if ( - configCandidatePending || - generation !== acceptedRestartTargetGeneration || - latestAcceptedRestartTarget !== acceptedTarget - ) { + if (this.acceptedTargetState !== acceptedState) { throw new GatewayConfigReloadSupersededError(); } const prepared = await target.prepareRuntimeConfig(); - if ( - configCandidatePending || - generation !== acceptedRestartTargetGeneration || - latestAcceptedRestartTarget !== acceptedTarget - ) { + if (this.acceptedTargetState !== acceptedState) { throw new GatewayConfigReloadSupersededError(); } return prepared; }, }; - latestAcceptedRestartTarget = acceptedTarget; - configCandidatePending = false; + const acceptedState = { kind: "accepted", generation, target: acceptedTarget } as const; + this.acceptedTargetState = acceptedState; return { reject: () => { - if (latestAcceptedRestartTarget !== acceptedTarget) { + const state = this.acceptedTargetState; + const ownsAcceptedTarget = + (state.kind === "accepted" && state.target === acceptedTarget) || + (state.kind === "candidate-pending" && state.previousTarget === acceptedTarget); + if (!ownsAcceptedTarget) { return; } - acceptedRestartTargetGeneration += 1; - latestAcceptedRestartTarget = null; - configCandidatePending = true; + this.acceptedTargetState = { + kind: "candidate-pending", + generation: generation + 1, + previousTarget: undefined, + }; }, - } satisfies AcceptedRestartTargetOwnership; - }; - - const createRestartRequestDetails = ( - plan: GatewayReloadPlan, - nextConfig: OpenClawConfig, - options?: GatewayRestartRequestOptions, - ): RestartRequestDetails => { - const explicitRestartPaths = plan.restartReasons.filter((path) => - plan.changedPaths.includes(path), - ); - return { - plan, - nextConfig: options?.debtConfig ?? nextConfig, - restartOwnedPaths: - explicitRestartPaths.length > 0 ? explicitRestartPaths : [...plan.changedPaths], - retainDebtAcrossConfigChanges: options?.retainDebtAcrossConfigChanges === true, }; - }; + } - const deferGatewayRestartDebt = ( + publishAcceptedTarget(target: AcceptedRestartTarget) { + return { + ownership: this.recordAcceptedTarget(target), + conservativeDebt: this.takeConservativeDebt(), + }; + } + + restoreConservativeDebt(debt: RestartRequestDetails): void { + this.conservativeDebt ??= debt; + } + + deferDebt( plan: GatewayReloadPlan, nextConfig: OpenClawConfig, options?: GatewayRestartRequestOptions, - ) => { - const details = createRestartRequestDetails(plan, nextConfig, options); - if (details.retainDebtAcrossConfigChanges) { - conservativeRestartDebt = details; - } else { - pausedRestartDebt = details; - } - }; + ): void { + this.preserveDebt(this.createRequestDetails(plan, nextConfig, options)); + } - const preserveRestartDebt = (details: RestartRequestDetails) => { - if (details.retainDebtAcrossConfigChanges) { - conservativeRestartDebt = details; - } else { - pausedRestartDebt = details; - } - }; - - const takeConservativeRestartDebt = (): RestartRequestDetails | null => { - const debt = conservativeRestartDebt; - conservativeRestartDebt = null; - return debt; - }; - - const restoreConservativeRestartDebt = (debt: RestartRequestDetails) => { - conservativeRestartDebt ??= debt; - }; - - const publishAcceptedRestartTarget = (target: AcceptedRestartTarget) => ({ - ownership: recordAcceptedRestartTarget(target), - conservativeDebt: takeConservativeRestartDebt(), - }); - - const markRestartEmissionSettled = () => { - restartEmissionSettled = true; - conservativeRestartDebt = null; - }; - - const isCurrentRestartRetry = (retry: { requestGeneration: number }) => - !restartRetryStopped && - retry.requestGeneration === restartRequestGeneration && - isCurrentGatewayReloadGeneration(myGeneration); - - const supersedeRestartRequest = () => { - restartRequestGeneration += 1; - restartPending = false; - restartDeferral?.cancel(); - restartDeferral = null; - if (restartRetryTimer) { - clearTimeout(restartRetryTimer); - restartRetryTimer = null; - } - restartRequestTransaction = null; - restartRequestDetails = null; - restartEmissionSettled = false; - }; - - const stopRestartRetries = () => { - restartRetryStopped = true; - pausedRestartDebt = null; - conservativeRestartDebt = null; - supersedeRestartRequest(); - }; - - const appliedConfigHashPublisher = createAppliedConfigHashPublisher({ - hasPendingRestart: () => - restartRequestDetails !== null || - pausedRestartDebt !== null || - conservativeRestartDebt !== null, - publish: setRuntimeConfigAppliedHash, - }); - - const scheduleRestartEmissionRetry = (retry: { - reason: string; - intent?: GatewayRestartIntent; - requestGeneration: number; - prepareForEmit?: () => Promise; - }) => { - if (restartRetryTimer || !isCurrentRestartRetry(retry)) { - return; - } - // Retry the exact failed emission. Re-entering request planning would start - // a fresh idle deferral and discard a timeout's force/deadline decision. - restartPending = true; - restartRetryTimer = setTimeout(() => { - restartRetryTimer = null; - if (!isCurrentRestartRetry(retry)) { - return; - } - // Timer callbacks outlive the config transaction root. Re-enter process - // admission so prepared host suspension cannot race signal delivery. - void runWithGatewayIndependentRootWorkAdmission(async () => { - if (!isCurrentRestartRetry(retry)) { - return; - } - restartPending = false; - if (retry.prepareForEmit && !(await retry.prepareForEmit())) { - scheduleRestartEmissionRetry(retry); - return; - } - const emitResult = params.requestRecoveryRestart?.(retry.reason, retry.intent); - if (emitResult && emitResult.status !== "failed") { - markRestartEmissionSettled(); - } - if (!emitResult || emitResult.status === "failed") { - scheduleRestartEmissionRetry(retry); - } - }).catch((err: unknown) => { - if (isCurrentRestartRetry(retry)) { - params.logReload.warn(`gateway restart recovery retry stopped: ${String(err)}`); - } - }); - }, RESTART_EMISSION_RETRY_MS); - restartRetryTimer.unref?.(); - }; - - const acceptRestartConfig = (acceptedConfig?: OpenClawConfig) => { - if (restartRequestTransaction?.state !== "rejected") { + acceptConfig(acceptedConfig?: OpenClawConfig) { + if (this.operation.kind === "idle" || this.operation.transaction.state !== "rejected") { return { retireRejectedRestart: false }; } - const rejectedDebt = !restartEmissionSettled ? restartRequestDetails : null; - if (rejectedDebt) { - preserveRestartDebt(rejectedDebt); + if (this.operation.kind === "request" && !this.operation.emissionSettled) { + this.preserveDebt(this.operation.details); } - supersedeRestartRequest(); - const configDebt = pausedRestartDebt; + this.supersedeRequest(); + const configDebt = this.pausedDebt; const retainsConfigDebt = configDebt && acceptedConfig && @@ -275,61 +189,213 @@ export function createGatewayRestartCoordinator(coordinatorOptions: { ), ); if (!retainsConfigDebt) { - pausedRestartDebt = null; + this.pausedDebt = null; } - const debt = (retainsConfigDebt ? configDebt : null) ?? conservativeRestartDebt; - if (debt) { - return { retireRejectedRestart: false, debt }; - } - return { retireRejectedRestart: true }; - }; - const retireRejectedRestartRequest = () => acceptRestartConfig().retireRejectedRestart; + const debt = (retainsConfigDebt ? configDebt : null) ?? this.conservativeDebt; + return debt ? { retireRejectedRestart: false, debt } : { retireRejectedRestart: true }; + } - const beginGatewayRestartLifecycle = () => { + retireRejectedRequest(): boolean { + return this.acceptConfig().retireRejectedRestart; + } + + beginLifecycle() { // A newer restart candidate owns the disk config now. Cancel any older // emission before async preflight so it cannot restart into stale secrets. if ( - !restartEmissionSettled && - restartRequestTransaction?.state !== "pending" && - restartRequestDetails + this.operation.kind === "request" && + !this.operation.emissionSettled && + this.operation.transaction.state !== "pending" ) { - preserveRestartDebt(restartRequestDetails); + this.preserveDebt(this.operation.details); } - supersedeRestartRequest(); + this.supersedeRequest(); const transaction = { state: "pending" as GatewayRestartTransactionState }; - restartRequestTransaction = transaction; + this.operation = { kind: "lifecycle", transaction }; return { settle: (state: Exclude) => { if (transaction.state === "pending") { transaction.state = state; if (state === "committed") { - pausedRestartDebt = null; + this.pausedDebt = null; } } }, }; - }; + } - const pauseGatewayRestartForConfigCandidate = () => { - configCandidatePending = true; - const lifecycle = beginGatewayRestartLifecycle(); + pauseForConfigCandidate(): void { + const state = this.acceptedTargetState; + const previousTarget = + state.kind === "accepted" + ? state.target + : state.kind === "candidate-pending" + ? state.previousTarget + : undefined; + this.acceptedTargetState = { + kind: "candidate-pending", + generation: state.generation, + previousTarget, + }; // Candidate acceptance owns debt rearm. Until then, invalid/failed config // must leave the prior committed restart paused. - lifecycle.settle("rejected"); - }; + this.beginLifecycle().settle("rejected"); + } - const requestGatewayRestartForGeneration = ( + request( + plan: GatewayReloadPlan, + nextConfig: OpenClawConfig, + options?: GatewayRestartRequestOptions, + ): GatewayRestartTransactionResult { + if (this.retryStopped) { + return { status: "recovery-pending", settle: () => {} }; + } + // Only another restart requirement supersedes accepted restart work. A + // duplicate, hot-only, or failed config transaction must preserve it. + this.supersedeRequest(); + const transaction = { state: "pending" as GatewayRestartTransactionState }; + this.operation = { + kind: "request", + transaction, + details: this.createRequestDetails(plan, nextConfig, options), + emissionSettled: false, + }; + const requestGeneration = this.requestGeneration; + const accepted = this.requestForGeneration(plan, nextConfig, requestGeneration, options); + return { + status: accepted ? "accepted" : "recovery-pending", + settle: (state) => { + if (transaction.state === "pending") { + transaction.state = state; + } + }, + }; + } + + stop(): void { + this.retryStopped = true; + this.pausedDebt = null; + this.conservativeDebt = null; + this.supersedeRequest(); + } + + private createRequestDetails( + plan: GatewayReloadPlan, + nextConfig: OpenClawConfig, + options?: GatewayRestartRequestOptions, + ): RestartRequestDetails { + const explicitRestartPaths = plan.restartReasons.filter((path) => + plan.changedPaths.includes(path), + ); + return { + plan, + nextConfig: options?.debtConfig ?? nextConfig, + restartOwnedPaths: + explicitRestartPaths.length > 0 ? explicitRestartPaths : [...plan.changedPaths], + retainDebtAcrossConfigChanges: options?.retainDebtAcrossConfigChanges === true, + }; + } + + private preserveDebt(details: RestartRequestDetails): void { + if (details.retainDebtAcrossConfigChanges) { + this.conservativeDebt = details; + } else { + this.pausedDebt = details; + } + } + + private takeConservativeDebt(): RestartRequestDetails | null { + const debt = this.conservativeDebt; + this.conservativeDebt = null; + return debt; + } + + private markEmissionSettled(): void { + if (this.operation.kind === "request") { + this.operation.emissionSettled = true; + } + this.conservativeDebt = null; + } + + private isCurrentRequest(requestGeneration: number): boolean { + return ( + !this.retryStopped && + requestGeneration === this.requestGeneration && + isCurrentGatewayReloadGeneration(this.options.myGeneration) + ); + } + + private supersedeRequest(): void { + this.requestGeneration += 1; + this.restartPending = false; + this.restartDeferral?.cancel(); + this.restartDeferral = null; + if (this.retryTimer) { + clearTimeout(this.retryTimer); + this.retryTimer = null; + } + this.operation = { kind: "idle" }; + } + + private scheduleEmissionRetry(retry: { + reason: string; + intent?: GatewayRestartIntent; + requestGeneration: number; + prepareForEmit?: () => Promise; + }): void { + if (this.retryTimer || !this.isCurrentRequest(retry.requestGeneration)) { + return; + } + // Retry the exact failed emission. Re-entering request planning would start + // a fresh idle deferral and discard a timeout's force/deadline decision. + this.restartPending = true; + this.retryTimer = setTimeout(() => { + this.retryTimer = null; + if (!this.isCurrentRequest(retry.requestGeneration)) { + return; + } + // Timer callbacks outlive the config transaction root. Re-enter process + // admission so prepared host suspension cannot race signal delivery. + void runWithGatewayIndependentRootWorkAdmission(async () => { + if (!this.isCurrentRequest(retry.requestGeneration)) { + return; + } + this.restartPending = false; + if (retry.prepareForEmit && !(await retry.prepareForEmit())) { + this.scheduleEmissionRetry(retry); + return; + } + const emitResult = this.options.params.requestRecoveryRestart?.(retry.reason, retry.intent); + if (emitResult && emitResult.status !== "failed") { + this.markEmissionSettled(); + } + if (!emitResult || emitResult.status === "failed") { + this.scheduleEmissionRetry(retry); + } + }).catch((err: unknown) => { + if (this.isCurrentRequest(retry.requestGeneration)) { + this.options.params.logReload.warn( + `gateway restart recovery retry stopped: ${String(err)}`, + ); + } + }); + }, RESTART_EMISSION_RETRY_MS); + this.retryTimer.unref?.(); + } + + private requestForGeneration( plan: GatewayReloadPlan, nextConfig: OpenClawConfig, requestGeneration: number, options?: GatewayRestartRequestOptions, - ): boolean => { + ): boolean { + const { params } = this.options; const reasons = plan.restartReasons.length ? plan.restartReasons.join(", ") : plan.changedPaths.join(", "); const restartReason = `config reload: ${reasons}`; - if (!restartRecoveryAvailable) { + if (!this.options.restartRecoveryAvailable) { params.logReload.warn( "gateway restart recovery unavailable; restart-required reload rejected", ); @@ -346,12 +412,12 @@ export function createGatewayRestartCoordinator(coordinatorOptions: { const preparedConfig = options?.prepareRuntimeConfig ? await options.prepareRuntimeConfig() : nextConfig; - if (requestGeneration !== restartRequestGeneration) { + if (!this.isCurrentRequest(requestGeneration)) { return false; } emissionPrepared = true; setGatewaySigusr1RestartPolicy({ allowExternal: isRestartEnabled(preparedConfig) }); - return requestGeneration === restartRequestGeneration; + return this.isCurrentRequest(requestGeneration); } catch (err) { emissionPrepared = false; params.logReload.warn(`gateway restart secrets preflight failed: ${String(err)}`); @@ -359,23 +425,23 @@ export function createGatewayRestartCoordinator(coordinatorOptions: { } }; - const active = getActiveCounts(); + const active = this.options.getActiveCounts(); if (active.totalActive > 0 || options?.prepareRuntimeConfig) { // Avoid spinning up duplicate polling loops from repeated config changes. - if (restartPending) { + if (this.restartPending) { params.logReload.info( `config change requires gateway restart (${reasons}) — already waiting for operations to complete`, ); return true; } - restartPending = true; + this.restartPending = true; if (active.totalActive > 0) { - const initialDetails = formatActiveDetails(active); + const initialDetails = this.options.formatActiveDetails(active); params.logReload.warn( `config change requires gateway restart (${reasons}) — deferring until ${initialDetails.join(", ")} complete`, ); - const taskBlockers = formatTaskBlockers(); + const taskBlockers = this.options.formatTaskBlockers(); if (taskBlockers) { params.logReload.warn( `restart blocked by active background task run(s): ${taskBlockers}`, @@ -386,8 +452,8 @@ export function createGatewayRestartCoordinator(coordinatorOptions: { } let failedEmission: { reason: string; intent?: GatewayRestartIntent } | undefined; - restartDeferral = deferGatewayRestartUntilIdle({ - getPendingCount: () => getActiveCounts().totalActive, + this.restartDeferral = deferGatewayRestartUntilIdle({ + getPendingCount: () => this.options.getActiveCounts().totalActive, maxWaitMs: resolveGatewayRestartDeferralTimeoutMs(undefined), timeoutIntent: { force: true, reason: "config reload forced restart" }, reason: restartReason, @@ -396,7 +462,7 @@ export function createGatewayRestartCoordinator(coordinatorOptions: { emissionPrepared = await prepareForEmit(); }, emitRestart: (reason, intent) => { - if (requestGeneration !== restartRequestGeneration) { + if (!this.isCurrentRequest(requestGeneration)) { return { status: "coalesced" }; } const resolvedReason = reason ?? restartReason; @@ -406,22 +472,22 @@ export function createGatewayRestartCoordinator(coordinatorOptions: { } const emitResult = requestRecoveryRestart(resolvedReason, intent); if (emitResult.status !== "failed") { - markRestartEmissionSettled(); + this.markEmissionSettled(); } failedEmission = emitResult.status === "failed" ? { reason: resolvedReason, intent } : undefined; return emitResult; }, afterEmitFailed: async () => { - if (requestGeneration !== restartRequestGeneration || !failedEmission) { + if (!this.isCurrentRequest(requestGeneration) || !failedEmission) { return; } - if (!restartRecoveryAvailable) { + if (!this.options.restartRecoveryAvailable) { params.logReload.warn("gateway restart recovery unavailable; retry skipped"); return; } params.logReload.warn("gateway restart recovery emission failed; retrying"); - scheduleRestartEmissionRetry({ + this.scheduleEmissionRetry({ ...failedEmission, requestGeneration, prepareForEmit, @@ -430,33 +496,25 @@ export function createGatewayRestartCoordinator(coordinatorOptions: { }, hooks: { onReady: () => { - restartPending = false; - restartDeferral = null; + this.restartPending = false; + this.restartDeferral = null; params.logReload.info("all operations and replies completed; restarting gateway now"); }, onStillPending: (_pending, elapsedMs) => { - const remaining = formatActiveDetails(getActiveCounts()); - const taskBlockersValue = formatTaskBlockers(); params.logReload.warn( - `restart still deferred after ${elapsedMs}ms with ${remaining.join(", ")} active${ - taskBlockersValue ? ` (${taskBlockersValue})` : "" - }`, + `restart still deferred after ${elapsedMs}ms with ${this.options.formatDeferredWorkStatus("active")}`, ); }, onTimeout: (_pending, elapsedMs) => { - const remaining = formatActiveDetails(getActiveCounts()); - const taskBlockersLocal = formatTaskBlockers(); - restartPending = false; - restartDeferral = null; + this.restartPending = false; + this.restartDeferral = null; params.logReload.warn( - `restart timeout after ${elapsedMs}ms with ${remaining.join(", ")} still active${ - taskBlockersLocal ? ` (${taskBlockersLocal})` : "" - }; forcing restart`, + `restart timeout after ${elapsedMs}ms with ${this.options.formatDeferredWorkStatus("still active")}; forcing restart`, ); }, onCheckError: (err) => { - restartPending = false; - restartDeferral = null; + this.restartPending = false; + this.restartDeferral = null; params.logReload.warn( `restart deferral check failed (${String(err)}); restarting gateway now`, ); @@ -473,12 +531,12 @@ export function createGatewayRestartCoordinator(coordinatorOptions: { // atomically promotes it to one-way restart drain. const emitResult = requestRecoveryRestart(restartReason); if (emitResult.status !== "failed") { - markRestartEmissionSettled(); + this.markEmissionSettled(); } if (emitResult.status === "failed") { params.logReload.warn("gateway restart recovery emission failed"); - if (restartRecoveryAvailable) { - scheduleRestartEmissionRetry({ + if (this.options.restartRecoveryAvailable) { + this.scheduleEmissionRetry({ reason: restartReason, requestGeneration, prepareForEmit, @@ -491,54 +549,37 @@ export function createGatewayRestartCoordinator(coordinatorOptions: { } setGatewaySigusr1RestartPolicy({ allowExternal: isRestartEnabled(nextConfig) }); return true; - }; - - const requestGatewayRestart = ( - plan: GatewayReloadPlan, - nextConfig: OpenClawConfig, - options?: GatewayRestartRequestOptions, - ): GatewayRestartTransactionResult => { - if (restartRetryStopped) { - return { status: "recovery-pending", settle: () => {} }; - } - // Only another restart requirement supersedes accepted restart work. A - // duplicate, hot-only, or failed config transaction must preserve it. - supersedeRestartRequest(); - const transaction = { state: "pending" as GatewayRestartTransactionState }; - restartRequestTransaction = transaction; - restartEmissionSettled = false; - restartRequestDetails = createRestartRequestDetails(plan, nextConfig, options); - const accepted = requestGatewayRestartForGeneration( - plan, - nextConfig, - restartRequestGeneration, - options, - ); - return { - status: accepted ? "accepted" : "recovery-pending", - settle: (state) => { - if (transaction.state === "pending") { - transaction.state = state; - } - }, - }; - }; + } +} +export function createGatewayRestartCoordinator(options: GatewayRestartCoordinatorOptions) { + const transaction = new GatewayRestartTransaction(options); return { - acceptRestartConfig, - ...appliedConfigHashPublisher, - beginGatewayRestartLifecycle, - pauseGatewayRestartForConfigCandidate, - publishAcceptedRestartTarget, - recordAcceptedRestartTarget, - requestGatewayRestart, - restoreConservativeRestartDebt, - retireRejectedRestartRequest, - stopRestartRetries, - deferGatewayRestartDebt, - getLatestAcceptedRestartTarget: () => latestAcceptedRestartTarget, - hasConfigCandidatePending: () => configCandidatePending, - hasRestartRequestTransaction: () => restartRequestTransaction !== null, - isRestartRetryStopped: () => restartRetryStopped, + acceptRestartConfig: (config?: OpenClawConfig) => transaction.acceptConfig(config), + ...transaction.appliedConfigHashPublisher, + beginGatewayRestartLifecycle: () => transaction.beginLifecycle(), + pauseGatewayRestartForConfigCandidate: () => transaction.pauseForConfigCandidate(), + publishAcceptedRestartTarget: (target: AcceptedRestartTarget) => + transaction.publishAcceptedTarget(target), + recordAcceptedRestartTarget: (target: AcceptedRestartTarget) => + transaction.recordAcceptedTarget(target), + requestGatewayRestart: ( + plan: GatewayReloadPlan, + nextConfig: OpenClawConfig, + requestOptions?: GatewayRestartRequestOptions, + ) => transaction.request(plan, nextConfig, requestOptions), + restoreConservativeRestartDebt: (debt: RestartRequestDetails) => + transaction.restoreConservativeDebt(debt), + retireRejectedRestartRequest: () => transaction.retireRejectedRequest(), + stopRestartRetries: () => transaction.stop(), + deferGatewayRestartDebt: ( + plan: GatewayReloadPlan, + nextConfig: OpenClawConfig, + requestOptions?: GatewayRestartRequestOptions, + ) => transaction.deferDebt(plan, nextConfig, requestOptions), + getLatestAcceptedRestartTarget: transaction.getAcceptedTarget, + hasConfigCandidatePending: transaction.hasPendingConfigCandidate, + hasRestartRequestTransaction: transaction.hasOperation, + isRestartRetryStopped: transaction.isStopped, }; }