mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-27 12:56:01 -06:00
refactor(gateway): make reload restart transactional (#114368)
This commit is contained in:
committed by
GitHub
parent
ba499528d4
commit
2bd4858252
@@ -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<ChannelKind>,
|
||||
isTransactionCurrent: () => boolean,
|
||||
@@ -134,6 +139,7 @@ export function createGatewayActiveWorkTracker(options: {
|
||||
|
||||
return {
|
||||
formatActiveDetails,
|
||||
formatDeferredWorkStatus,
|
||||
formatTaskBlockers,
|
||||
getActiveCounts,
|
||||
waitForActiveWorkBeforeChannelReload,
|
||||
|
||||
@@ -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,
|
||||
});
|
||||
|
||||
|
||||
@@ -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<typeof setTimeout> | 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<typeof setTimeout> | 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<boolean>;
|
||||
}) => {
|
||||
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<GatewayRestartTransactionState, "pending">) => {
|
||||
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<boolean>;
|
||||
}): 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,
|
||||
};
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user