diff --git a/docs/cli/gateway/restart-and-supervision.md b/docs/cli/gateway/restart-and-supervision.md index 472ad587196b..4f6e8effbd48 100644 --- a/docs/cli/gateway/restart-and-supervision.md +++ b/docs/cli/gateway/restart-and-supervision.md @@ -68,7 +68,15 @@ the refusal, and deep status reports it instead of an unavailable shutdown recor Foreground/manual Gateways, in-process restarts selected by `OPENCLAW_NO_RESPAWN=1`, and other supervisors retain exit status `1` when cleanup cannot finish before the shutdown deadline. -`--force` skips the active-work drain and requests cancellation of active cron runs before cleanup. The normal shutdown path still joins accepted work; existing shutdown deadlines still apply. Plain `restart` normally uses the service-manager restart path. +`--force` begins restart immediately and closes new admissions. The current CLI supplies the normal drain budget, capped by the native service shutdown deadline. Only work remaining at that deadline is canceled before cleanup. A safe restart whose deferral budget has already expired does not get a second drain budget. Plain `restart` normally uses the service-manager restart path. + +Forced requests from older callers that supply no drain budget get at most +45 seconds to drain. Their 60-second replacement window reserves 10 seconds for +cleanup and 5 seconds for replacement. This applies to interactive and update +callers alike. If work remains when the drain expires, the Gateway records a +warning with the remaining work categories in its restart history and log, then +cancels that work through normal terminal recovery. Callers that supply a budget +retain that budget, subject to native service deadlines. During an upgrade, restart records its reason and drain options in the existing Gateway state without starting a schema migration while the old Gateway is still diff --git a/docs/gateway/restart-recovery.md b/docs/gateway/restart-recovery.md index cfecb8ac7a51..9e13eeed2237 100644 --- a/docs/gateway/restart-recovery.md +++ b/docs/gateway/restart-recovery.md @@ -190,9 +190,10 @@ shutdown deadline. A shorter supervisor timeout also caps requested restart wait The drained work, ordering, and interruption behavior stay the same. Service-child cleanup uses the remaining Gateway shutdown budget, leaving time -for final exit bookkeeping. A forced restart handed to a supervisor skips active-work -drain but retains the 10-second cleanup reserve; it does not start a fresh -85-second wait. A restart without a supervisor handoff uses the existing shutdown +for final exit bookkeeping. A forced restart drains admitted work within the same +budget. When the restart scheduler has already exhausted its deferral budget, +cleanup retains the 10-second reserve without starting a second drain. +A restart without a supervisor handoff uses the existing shutdown deadline for cleanup. This includes foreground Gateways inside another service's cgroup, restarts with `OPENCLAW_NO_RESPAWN=1`, and standalone updates that must launch their own replacement. Cgroup membership alone does not provide a supervisor @@ -434,6 +435,13 @@ Copying one does not grant permission to restart a service. ## How interrupted work is detected +Startup reconciles older subagent session rows that still say `running` but have +no live run, task, admission, or recovery owner. It records a diagnostic transcript +receipt and marks the row `interrupted` in one transaction. The end timestamp records when +startup observed the interruption, rather than an inferred execution finish time; +the original activity timestamps remain intact. A failed receipt write becomes a +warning and leaves the row eligible for a later repair. + Three complementary mechanisms mark sessions whose turn did not finish: - **At turn admission:** for an ordinary text turn on an existing main session, diff --git a/gateway-shutdown-budget.d.mts b/gateway-shutdown-budget.d.mts index 03637034718f..013678346a31 100644 --- a/gateway-shutdown-budget.d.mts +++ b/gateway-shutdown-budget.d.mts @@ -3,3 +3,4 @@ export const GATEWAY_SUPERVISOR_EXIT_MARGIN_MS: number; export const GATEWAY_SHUTDOWN_TIMEOUT_MS: number; export const GATEWAY_SERVICE_STOP_TIMEOUT_MS: number; export const LAUNCH_AGENT_EXIT_TIMEOUT_SECONDS: 20; +export const GATEWAY_RESTART_REPLACEMENT_TIMEOUT_MS: number; diff --git a/gateway-shutdown-budget.mjs b/gateway-shutdown-budget.mjs index b8aec66c937f..c6469692105c 100644 --- a/gateway-shutdown-budget.mjs +++ b/gateway-shutdown-budget.mjs @@ -1,4 +1,5 @@ -// Shared Gateway stop policy used by the run loop and native service definitions. +// Shared stop policy; v2026.9.5 restart-health.constants.ts fixes the replacement window. +export const GATEWAY_RESTART_REPLACEMENT_TIMEOUT_MS = 60_000; const GATEWAY_SHUTDOWN_DRAIN_TIMEOUT_MS = 315_000; export const GATEWAY_SHUTDOWN_RESERVE_MS = 10_000; export const GATEWAY_SUPERVISOR_EXIT_MARGIN_MS = 5_000; diff --git a/src/agents/subagents/registry/subagent-orphan-recovery.restart-integration.test.ts b/src/agents/subagents/registry/subagent-orphan-recovery.restart-integration.test.ts index 9255b6522376..2b7356939f42 100644 --- a/src/agents/subagents/registry/subagent-orphan-recovery.restart-integration.test.ts +++ b/src/agents/subagents/registry/subagent-orphan-recovery.restart-integration.test.ts @@ -1,3 +1,4 @@ +import { isRecord } from "@openclaw/normalization-core/record-coerce"; // Restart-path proof against the real registry sweeper and SQLite session store. import { describe, expect, it, vi } from "vitest"; import { createDeferred } from "../../../../test/helpers/promise.js"; @@ -444,6 +445,21 @@ describe("subagent orphan recovery — faithful restart path", () => { endedAt: expect.any(Number), }); expect(persistedSession?.abortedLastRun).toBeUndefined(); + expect( + ( + await loadTranscriptEvents({ + agentId: "main", + storePath, + sessionKey: childSessionKey, + sessionId: "sess-stale-aborted", + }) + ).filter((event) => isRecord(event) && event.customType === "run-failed-before-reply"), + ).toMatchObject([ + { + display: true, + details: { runId, error: expect.stringContaining("Gateway restart") }, + }, + ]); }); it.each([60_000, 3 * TWO_HOURS_MS])( diff --git a/src/cli/cli-entrypoint.test-support.ts b/src/cli/cli-entrypoint.test-support.ts index 601d5ca2b14f..a8431b28aab7 100644 --- a/src/cli/cli-entrypoint.test-support.ts +++ b/src/cli/cli-entrypoint.test-support.ts @@ -67,6 +67,11 @@ export const updateFinalizationOutputEntrypoint = { // Direct-stop children use the invocation's prepared graph before readiness starts. export const gatewayDirectStopEntrypoints = { + startupOrphanFixture: { + currentModuleUrl: import.meta.url, + sourceWorkerName: "../gateway/startup-orphan-process.test-support", + distWorkerPath: "gateway/startup-orphan-process.test-support.js", + }, forcedCronFixture: { currentModuleUrl: import.meta.url, sourceWorkerName: "gateway-cli/run-loop.forced-cron.test-support", @@ -121,6 +126,46 @@ export const gatewayDirectStopEntrypoints = { // Extra update roots share the native fixture generation. export const updateExecutorEntrypoints = { + sealedRegistry: { + currentModuleUrl: import.meta.url, + sourceWorkerName: "../infra/sealed-runtime-registry", + distWorkerPath: "infra/sealed-runtime-registry.js", + }, + ledger: { + currentModuleUrl: import.meta.url, + sourceWorkerName: "../infra/update-run-ledger", + distWorkerPath: "infra/update-run-ledger.js", + }, + handoff: { + currentModuleUrl: import.meta.url, + sourceWorkerName: "../infra/update-managed-service-handoff", + distWorkerPath: "infra/update-managed-service-handoff.js", + }, + sentinel: { + currentModuleUrl: import.meta.url, + sourceWorkerName: "../infra/update-control-plane-sentinel", + distWorkerPath: "infra/update-control-plane-sentinel.js", + }, + packageSteps: { + currentModuleUrl: import.meta.url, + sourceWorkerName: "../infra/package-update-steps", + distWorkerPath: "infra/package-update-steps.js", + }, + packageFixture: { + currentModuleUrl: import.meta.url, + sourceWorkerName: "../infra/package-update-steps.test-support", + distWorkerPath: "infra/package-update-steps.test-support.js", + }, + inventory: { + currentModuleUrl: import.meta.url, + sourceWorkerName: "../../scripts/lib/package-dist-inventory", + distWorkerPath: "scripts/lib/package-dist-inventory.js", + }, + exec: { + currentModuleUrl: import.meta.url, + sourceWorkerName: "../process/exec", + distWorkerPath: "process/exec.js", + }, lease: { currentModuleUrl: import.meta.url, sourceWorkerName: "../infra/update-managed-service-handoff-lease", diff --git a/src/cli/daemon-cli/lifecycle-safe-restart.ts b/src/cli/daemon-cli/lifecycle-safe-restart.ts index 9b5acb58ec27..ca4fdbb69e46 100644 --- a/src/cli/daemon-cli/lifecycle-safe-restart.ts +++ b/src/cli/daemon-cli/lifecycle-safe-restart.ts @@ -4,6 +4,7 @@ import { refreshLegacySystemdServiceMetadata } from "../../daemon/systemd.js"; import { callGatewayCli } from "../../gateway/call.js"; import { formatErrorMessage } from "../../infra/errors.js"; import { resolveGatewayServiceMutationError } from "../../infra/gateway-supervision.js"; +import { resolveGatewayRestartDeferralTimeoutMs } from "../../infra/restart-budget.js"; import type { SafeGatewayRestartRequestResult } from "../../infra/restart-coordinator.js"; import type { GatewayRestartIntent } from "../../infra/restart-intent.js"; import { defaultRuntime, writeRuntimeJson } from "../../runtime.js"; @@ -24,7 +25,7 @@ export function resolveGatewayRestartIntentOptions( throw new Error("--force cannot be combined with --wait"); } if (opts.force) { - return { force: true }; + return { force: true, waitMs: resolveGatewayRestartDeferralTimeoutMs() }; } return opts.wait === undefined ? undefined : { waitMs: parseDurationMs(opts.wait) }; } @@ -37,7 +38,9 @@ export async function runSafeGatewayRestart( target?: SafeRestartTarget, ): Promise { if (opts.force) { - throw new Error("--safe cannot be combined with --force; omit --safe to force restart now"); + throw new Error( + "--safe cannot be combined with --force; omit --safe to begin a forced restart", + ); } if (opts.wait !== undefined) { throw new Error("--safe cannot be combined with --wait; safe restart uses gateway deferral"); diff --git a/src/cli/daemon-cli/lifecycle-unmanaged.ts b/src/cli/daemon-cli/lifecycle-unmanaged.ts index 3193fb824d3e..8a7137703c22 100644 --- a/src/cli/daemon-cli/lifecycle-unmanaged.ts +++ b/src/cli/daemon-cli/lifecycle-unmanaged.ts @@ -67,6 +67,9 @@ export async function signalGatewayRestart( env?: NodeJS.ProcessEnv; }, ) { + const restartIntent = params.restartIntent?.force + ? { force: true, drainBudgetMs: params.restartIntent.waitMs } + : params.restartIntent; if (params.enforceRestartConfig) { await assertUnmanagedGatewayRestartEnabled(port); } @@ -144,7 +147,7 @@ export async function signalGatewayRestart( ownerId: previousLockIdentity.ownerId, port, }, - ...(params.restartIntent ? { restartIntent: params.restartIntent } : {}), + ...(restartIntent ? { restartIntent } : {}), }, localPortOverride: port, ignoreEnvUrlOverride: true, diff --git a/src/cli/daemon-cli/lifecycle.external-supervision.test.ts b/src/cli/daemon-cli/lifecycle.external-supervision.test.ts index 18b4edbd8fdd..05c5b6ebc4d1 100644 --- a/src/cli/daemon-cli/lifecycle.external-supervision.test.ts +++ b/src/cli/daemon-cli/lifecycle.external-supervision.test.ts @@ -220,7 +220,7 @@ describe("external gateway supervision lifecycle", () => { ownerId: "gateway-owner-old", port: 19_455, }, - restartIntent: { force: true }, + restartIntent: { force: true, drainBudgetMs: 300_000 }, }, localPortOverride: 19_455, ignoreEnvUrlOverride: true, @@ -237,7 +237,8 @@ describe("external gateway supervision lifecycle", () => { }); expect(waitForGatewayHealthyListener).toHaveBeenCalledWith({ port: 19_455, - attempts: 120, + // Allow the five-minute drain budget before the one-minute readiness window. + attempts: 720, delayMs: 500, previousLockIdentity: lockIdentity, waitIndefinitelyForPreviousOwner: false, diff --git a/src/cli/daemon-cli/lifecycle.test.ts b/src/cli/daemon-cli/lifecycle.test.ts index 73c7c1460b7e..31a277856d10 100644 --- a/src/cli/daemon-cli/lifecycle.test.ts +++ b/src/cli/daemon-cli/lifecycle.test.ts @@ -800,20 +800,25 @@ describe("runDaemonRestart health checks", () => { expect(signalVerifiedGatewayPidSync).not.toHaveBeenCalled(); }); - it.each(["win32", "linux", "darwin"] as const)( - "uses targeted RPC for an unmanaged %s gateway restart", - async (platform) => { + it.each( + (["win32", "linux", "darwin"] as const).flatMap((platform) => + [false, true].map((force) => ({ platform, force })), + ), + )( + "uses targeted RPC for an unmanaged $platform gateway restart (force=$force)", + async ({ platform, force }) => { vi.spyOn(process, "platform", "get").mockReturnValue(platform); callGatewayCli.mockResolvedValueOnce({ ok: true, status: "emitted", pid: 4200 }); findVerifiedGatewayListenerPidsOnPortSync.mockReturnValue([4200]); mockUnmanagedRestart({ runPostRestartCheck: true }); - await runDaemonRestart({ json: true }); + await runDaemonRestart({ json: true, force }); expect(callGatewayCli).toHaveBeenCalledWith({ method: "gateway.restart.request", params: { reason: "gateway.restart", + ...(force ? { restartIntent: { force: true, drainBudgetMs: 300_000 } } : {}), target: { pid: 4200, ownerId: "gateway-owner-old", diff --git a/src/cli/daemon-cli/lifecycle.ts b/src/cli/daemon-cli/lifecycle.ts index 7a3ae9543b67..85bd09d4a7aa 100644 --- a/src/cli/daemon-cli/lifecycle.ts +++ b/src/cli/daemon-cli/lifecycle.ts @@ -28,8 +28,8 @@ import { resolveGatewayServiceMutationError, } from "../../infra/gateway-supervision.js"; import { probePortUsage } from "../../infra/ports-probe.js"; +import { resolveGatewayRestartDrainTimeoutMs } from "../../infra/restart-budget.js"; import type { GatewayRestartIntent } from "../../infra/restart-intent.js"; -import { resolveGatewayRestartDeferralTimeoutMs } from "../../infra/restart.js"; import { defaultRuntime } from "../../runtime.js"; import { formatCliCommand } from "../command-format.js"; import { @@ -159,29 +159,13 @@ async function stopGatewayWithoutServiceManager( }; } -async function resolveRestartListenerHealthWait(restartIntent: GatewayRestartIntent | undefined) { - let drainTimeoutMs: number | undefined; - if (restartIntent?.force) { - drainTimeoutMs = 0; - } else if (typeof restartIntent?.waitMs === "number" && Number.isFinite(restartIntent.waitMs)) { - drainTimeoutMs = restartIntent.waitMs > 0 ? Math.floor(restartIntent.waitMs) : undefined; - } else { - drainTimeoutMs = resolveGatewayRestartDeferralTimeoutMs(); - } - - const replacementHealthAttempts = postRestartHealthAttempts(); - if (drainTimeoutMs === undefined) { - return { - attempts: replacementHealthAttempts, - waitIndefinitelyForPreviousOwner: true, - timeoutSeconds: Math.round((replacementHealthAttempts * POST_RESTART_HEALTH_DELAY_MS) / 1000), - }; - } +function resolveRestartListenerHealthWait(restartIntent: GatewayRestartIntent | undefined) { + const drainTimeoutMs = resolveGatewayRestartDrainTimeoutMs(restartIntent); const attempts = - replacementHealthAttempts + Math.ceil(drainTimeoutMs / POST_RESTART_HEALTH_DELAY_MS); + postRestartHealthAttempts() + Math.ceil((drainTimeoutMs ?? 0) / POST_RESTART_HEALTH_DELAY_MS); return { attempts, - waitIndefinitelyForPreviousOwner: false, + waitIndefinitelyForPreviousOwner: drainTimeoutMs === undefined, timeoutSeconds: Math.round((attempts * POST_RESTART_HEALTH_DELAY_MS) / 1000), }; } @@ -240,7 +224,7 @@ async function runExternalSupervisorRestart(opts: DaemonLifecycleOptions): Promi return false; } - const healthWait = await resolveRestartListenerHealthWait(restartIntent); + const healthWait = resolveRestartListenerHealthWait(restartIntent); const health = await waitForGatewayHealthyListener({ port: lockIdentity.port, attempts: healthWait.attempts, @@ -409,9 +393,7 @@ export async function runDaemonRestart(opts: DaemonLifecycleOptions = {}): Promi const restartHealthAttempts = postRestartHealthAttempts(); const restartWaitMs = restartHealthAttempts * POST_RESTART_HEALTH_DELAY_MS; const restartWaitSeconds = Math.round(restartWaitMs / 1000); - let unmanagedRestartHealthAttempts = restartHealthAttempts; - let unmanagedRestartWaitIndefinitely = false; - let unmanagedRestartWaitSeconds = restartWaitSeconds; + let unmanagedRestartWait: ReturnType | undefined; return await runServiceRestart({ serviceNoun: "Gateway", @@ -464,10 +446,7 @@ export async function runDaemonRestart(opts: DaemonLifecycleOptions = {}): Promi } restartedWithoutServiceManager = true; unmanagedPreviousLockIdentity = handled.previousLockIdentity; - const healthWait = await resolveRestartListenerHealthWait(restartIntent); - unmanagedRestartHealthAttempts = healthWait.attempts; - unmanagedRestartWaitIndefinitely = healthWait.waitIndefinitelyForPreviousOwner; - unmanagedRestartWaitSeconds = healthWait.timeoutSeconds; + unmanagedRestartWait = resolveRestartListenerHealthWait(restartIntent); return handled; }, repairLoadedService: preserveDefinition @@ -510,10 +489,7 @@ export async function runDaemonRestart(opts: DaemonLifecycleOptions = {}): Promi restartedWithoutServiceManager = true; if (isGatewaySignalRestartResult(handled) && handled.previousLockIdentity) { unmanagedPreviousLockIdentity = handled.previousLockIdentity; - const healthWait = await resolveRestartListenerHealthWait(restartIntent); - unmanagedRestartHealthAttempts = healthWait.attempts; - unmanagedRestartWaitIndefinitely = healthWait.waitIndefinitelyForPreviousOwner; - unmanagedRestartWaitSeconds = healthWait.timeoutSeconds; + unmanagedRestartWait = resolveRestartListenerHealthWait(restartIntent); } return handled; } @@ -530,12 +506,13 @@ export async function runDaemonRestart(opts: DaemonLifecycleOptions = {}): Promi const health = await waitForGatewayHealthyListener({ port: unmanagedPort, env: process.env, - attempts: unmanagedRestartHealthAttempts, + attempts: unmanagedRestartWait?.attempts ?? restartHealthAttempts, delayMs: POST_RESTART_HEALTH_DELAY_MS, ...(unmanagedPreviousLockIdentity ? { previousLockIdentity: unmanagedPreviousLockIdentity, - waitIndefinitelyForPreviousOwner: unmanagedRestartWaitIndefinitely, + waitIndefinitelyForPreviousOwner: + unmanagedRestartWait?.waitIndefinitelyForPreviousOwner ?? false, } : {}), }); @@ -544,7 +521,8 @@ export async function runDaemonRestart(opts: DaemonLifecycleOptions = {}): Promi } const diagnostics = renderGatewayPortHealthDiagnostics(health); - const timeoutLine = `Timed out after ${unmanagedRestartWaitSeconds}s waiting for gateway port ${unmanagedPort} to become healthy.`; + const waitSeconds = unmanagedRestartWait?.timeoutSeconds ?? restartWaitSeconds; + const timeoutLine = `Timed out after ${waitSeconds}s waiting for gateway port ${unmanagedPort} to become healthy.`; if (!jsonOutput) { defaultRuntime.log(theme.warn(timeoutLine)); for (const line of diagnostics) { @@ -556,7 +534,7 @@ export async function runDaemonRestart(opts: DaemonLifecycleOptions = {}): Promi } fail( - `Gateway restart timed out after ${unmanagedRestartWaitSeconds}s waiting for health checks.`, + `Gateway restart timed out after ${waitSeconds}s waiting for health checks.`, [formatCliCommand("openclaw gateway status --deep"), formatCliCommand("openclaw doctor")], activationAccepted ? "restart-health-failed" : undefined, ); diff --git a/src/cli/daemon-cli/register-service-commands.ts b/src/cli/daemon-cli/register-service-commands.ts index 75d5e011961b..1adcc54057d0 100644 --- a/src/cli/daemon-cli/register-service-commands.ts +++ b/src/cli/daemon-cli/register-service-commands.ts @@ -174,7 +174,7 @@ export function addGatewayServiceCommands(parent: Command, opts?: { statusDescri ) .description("Restart the Gateway service (launchd/systemd/schtasks)") .option("--preserve-definition", "Keep the native service definition", false) - .option("--force", "Restart immediately without waiting for active gateway work", false) + .option("--force", "Begin restart now; drain admitted work within the shutdown budget", false) .option( "--safe", "Request an OpenClaw-aware restart after active work drains " + diff --git a/src/cli/daemon-cli/restart-health.constants.ts b/src/cli/daemon-cli/restart-health.constants.ts index a75d386bee3f..f9de3ba87d6f 100644 --- a/src/cli/daemon-cli/restart-health.constants.ts +++ b/src/cli/daemon-cli/restart-health.constants.ts @@ -1,4 +1,6 @@ -export const DEFAULT_RESTART_HEALTH_TIMEOUT_MS = 60_000; +import { GATEWAY_RESTART_REPLACEMENT_TIMEOUT_MS } from "../../infra/gateway-shutdown-budget.js"; + +export const DEFAULT_RESTART_HEALTH_TIMEOUT_MS = GATEWAY_RESTART_REPLACEMENT_TIMEOUT_MS; export const DEFAULT_RESTART_HEALTH_DELAY_MS = 500; export const DEFAULT_RESTART_HEALTH_ATTEMPTS = Math.ceil( diff --git a/src/cli/gateway-cli/lifecycle-import-boundary.test.ts b/src/cli/gateway-cli/lifecycle-import-boundary.test.ts index 8fde2a5a499e..3d20d58fa881 100644 --- a/src/cli/gateway-cli/lifecycle-import-boundary.test.ts +++ b/src/cli/gateway-cli/lifecycle-import-boundary.test.ts @@ -113,7 +113,7 @@ describe("gateway lifecycle hub import boundaries", () => { markGatewayRestartHandled: vi.fn(), abortPendingChannelReloads: vi.fn(), markGatewayDraining: vi.fn(), - resolveGatewayRestartDeferralTimeoutMs: () => 300_000, + resolveGatewayRestartDrainTimeoutMs: () => 300_000, createGatewayActiveWorkSnapshot: () => idle, waitForGatewayActiveWork: vi.fn(async () => ({ drained: true, snapshot: idle })), stopGatewayManagedProviderLocalServices: vi.fn(async () => {}), diff --git a/src/cli/gateway-cli/lifecycle.runtime.ts b/src/cli/gateway-cli/lifecycle.runtime.ts index 5f14702bbf87..a4cf5c851da3 100644 --- a/src/cli/gateway-cli/lifecycle.runtime.ts +++ b/src/cli/gateway-cli/lifecycle.runtime.ts @@ -10,8 +10,8 @@ export { respawnGatewayProcessForUpdate, restartGatewayProcessWithFreshPid, } from "../../infra/process-respawn.js"; +export { resolveGatewayRestartDrainTimeoutMs } from "../../infra/restart-budget.js"; export { - resolveGatewayRestartDeferralTimeoutMs, consumeGatewayRestartIntent, consumeGatewayRestartAuthorization, isGatewayRestartExternallyAllowed, diff --git a/src/cli/gateway-cli/run-loop-drain.ts b/src/cli/gateway-cli/run-loop-drain.ts index 26cf0ae1d179..5d07a55739a1 100644 --- a/src/cli/gateway-cli/run-loop-drain.ts +++ b/src/cli/gateway-cli/run-loop-drain.ts @@ -11,44 +11,26 @@ import { formatDrainCounts, formatShutdownReason } from "./run-loop-shutdown-for const RESTART_DRAIN_STILL_PENDING_WARN_MS = 30_000; -export function resolveRestartDrainTimeoutMs( - restartIntent: GatewayRunSignalRequest["restartIntent"], - runtime: Pick, -): number | undefined { - if (restartIntent?.force) { - return 0; - } - if (typeof restartIntent?.waitMs === "number" && Number.isFinite(restartIntent.waitMs)) { - return restartIntent.waitMs > 0 ? Math.floor(restartIntent.waitMs) : undefined; - } - try { - return runtime.resolveGatewayRestartDeferralTimeoutMs(); - } catch { - return 300_000; - } -} - export async function drainGatewayActiveWork({ request, - restartIntent, runtime, - loadRuntime, drainTimeoutMs, restartDrainDeadlineAt, markDraining, recordCounts, + recordWarning, logger, }: { request: GatewayRunSignalRequest; - restartIntent: GatewayRunSignalRequest["restartIntent"]; runtime: typeof import("./lifecycle.runtime.js"); - loadRuntime: () => Promise; drainTimeoutMs: number | undefined; restartDrainDeadlineAt: number | undefined; markDraining: (reason: GatewayDrainReason) => void; recordCounts: (counts: string) => void; + recordWarning: (warning: string) => void; logger: Pick; }) { + const { restartIntent } = request; const reportDrainSnapshot = createGatewayDrainReporter( request.action, drainTimeoutMs, @@ -66,7 +48,7 @@ export async function drainGatewayActiveWork({ "restart.drain", async () => { const { abortEmbeddedAgentRun, createGatewayActiveWorkSnapshot, waitForGatewayActiveWork } = - await loadRuntime(); + runtime; // Reject new enqueues immediately during the drain window so // sessions get an explicit restart error instead of silent task loss. markDraining(formatShutdownReason(request)); @@ -78,27 +60,23 @@ export async function drainGatewayActiveWork({ } reportDrainSnapshot(initialSnapshot); - if (restartIntent?.force) { - logger.warn("forced restart requested; skipping active work drain"); - } else { - const remainingDrainTimeoutMs = - restartDrainDeadlineAt === undefined - ? undefined - : Math.max(0, restartDrainDeadlineAt - Date.now()); - const drain = await waitForGatewayActiveWork(remainingDrainTimeoutMs, { - onSnapshot: reportDrainSnapshot, - }); - if (drain.drained) { - if (!initialSnapshot.idle) { - logger.info("all active work drained"); - } - return; + const remainingDrainTimeoutMs = + restartDrainDeadlineAt === undefined + ? undefined + : Math.max(0, restartDrainDeadlineAt - Date.now()); + const drain = await waitForGatewayActiveWork(remainingDrainTimeoutMs, { + onSnapshot: reportDrainSnapshot, + }); + if (drain.drained) { + if (!initialSnapshot.idle) { + logger.info("all active work drained"); } - drainTimedOut = true; - logger.warn( - `active-work drain timeout reached; proceeding with restart: ${formatDrainCounts(drain.snapshot)}`, - ); + return; } + drainTimedOut = true; + const warning = `restart drain budget ${drainTimeoutMs}ms exhausted; cutting short ${formatDrainCounts(drain.snapshot)}`; + recordWarning(warning); + logger.warn(warning); // Connection work can retain cron cleanup; cancel before close joins it. runtime.abortActiveCronTaskRuns("Gateway restarting."); }, diff --git a/src/cli/gateway-cli/run-loop-force.test-support.ts b/src/cli/gateway-cli/run-loop-force.test-support.ts new file mode 100644 index 000000000000..ad3643d9de5a --- /dev/null +++ b/src/cli/gateway-cli/run-loop-force.test-support.ts @@ -0,0 +1,201 @@ +import { performance } from "node:perf_hooks"; +import { expect, it, vi } from "vitest"; +import type { GatewayActiveWorkSnapshot } from "../../infra/gateway-active-work.js"; +import { createDeferredCore } from "../../shared/deferred.js"; +import type { RequestFixtures } from "./run-loop-request-fixtures.test-support.js"; +import { + createActiveWorkSnapshot, + createCloseMock, + createRuntimeWithExitSignal, + createSignaledStart, + expectRestartCloseCall, + waitForStart, + withIsolatedSignals, +} from "./run-loop.test-support.js"; + +export function registerGatewayForcedRestartTests({ + createSignaledLoopHarness, + createGatewayActiveWorkSnapshot, + abortActiveCronTaskRuns, + runLoopWithStart, + waitForGatewayActiveWork, + consumeGatewayRestartIntent, + consumeGatewayRestartIntentPayloadSync, + isGatewayWorkAdmissionClosed, + gatewayLog, + readCgroup, + systemctl, +}: Pick< + RequestFixtures, + | "createSignaledLoopHarness" + | "createGatewayActiveWorkSnapshot" + | "abortActiveCronTaskRuns" + | "runLoopWithStart" + | "waitForGatewayActiveWork" + | "consumeGatewayRestartIntent" + | "consumeGatewayRestartIntentPayloadSync" + | "isGatewayWorkAdmissionClosed" + | "gatewayLog" + | "readCgroup" + | "systemctl" +>): void { + const idleActiveWorkSnapshot = createActiveWorkSnapshot(); + it.each( + (["SIGTERM", "SIGUSR2"] as const).flatMap((signal) => + [undefined, 180_000].map((waitMs) => ({ signal, waitMs, budget: waitMs ?? 45_000 })), + ), + )( + "drains admitted work before a forced $signal restart (budget=$budget)", + async ({ signal, waitMs, budget }) => { + (signal === "SIGTERM" + ? consumeGatewayRestartIntentPayloadSync + : consumeGatewayRestartIntent + ).mockReturnValueOnce({ force: true, ...(waitMs === undefined ? {} : { waitMs }) }); + createGatewayActiveWorkSnapshot.mockReturnValueOnce( + createActiveWorkSnapshot({ activeTasks: 1, embeddedRuns: 1 }, [ + { + kind: "task", + count: 1, + message: "taskId=task-force runId=run-force status=running runtime=cron label=forced", + }, + { kind: "embedded-run", count: 1, message: "1 active embedded run(s)" }, + ]), + ); + const drain = createDeferredCore<{ drained: boolean; snapshot: GatewayActiveWorkSnapshot }>(); + waitForGatewayActiveWork.mockImplementationOnce(() => drain.promise); + await withIsolatedSignals(async ({ captureSignal }) => { + const { close, start, exited } = await createSignaledLoopHarness(); + const sigint = captureSignal("SIGINT"); + vi.useFakeTimers(); + const clock = vi.spyOn(performance, "now").mockImplementation(() => Date.now()); + try { + captureSignal(signal)(); + await vi.advanceTimersByTimeAsync(0); + expect(isGatewayWorkAdmissionClosed()).toBe(true); + expect(close).not.toHaveBeenCalled(); + expect(waitForGatewayActiveWork).toHaveBeenCalledWith(budget, expect.any(Object)); + expect(abortActiveCronTaskRuns).not.toHaveBeenCalled(); + expect(gatewayLog.info.mock.calls.flat().join("\n")).not.toContain("task-force"); + drain.resolve({ drained: true, snapshot: idleActiveWorkSnapshot }); + await vi.advanceTimersByTimeAsync(0); + expectRestartCloseCall(close, budget); + expect(start).toHaveBeenCalledTimes(signal === "SIGTERM" ? 1 : 2); + } finally { + drain.resolve({ drained: true, snapshot: idleActiveWorkSnapshot }); + await vi.advanceTimersByTimeAsync(0); + sigint(); + await vi.advanceTimersByTimeAsync(0); + await expect(exited).resolves.toBe(0); + clock.mockRestore(); + vi.useRealTimers(); + } + }); + }, + ); + + it.each([ + { waitMs: undefined, refreshMs: 0, stallClose: false }, + { waitMs: undefined, refreshMs: 10_000, stallClose: false }, + { waitMs: 0, refreshMs: 0, stallClose: false }, + { waitMs: 180_000, refreshMs: 0, stallClose: false }, + { waitMs: undefined, refreshMs: 0, stallClose: true }, + { waitMs: undefined, refreshMs: 10_000, stallClose: true }, + { waitMs: 180_000, refreshMs: 0, stallClose: true }, + ])( + "records cut work only when the forced caller drain budget expires (waitMs=$waitMs, refresh=$refreshMs, stalled close=$stallClose)", + async ({ waitMs, refreshMs, stallClose }) => { + const budget = waitMs ?? 45_000; + const active = createActiveWorkSnapshot({ activeTasks: 1, cronRuns: 1 }); + const drain = createDeferredCore<{ drained: boolean; snapshot: GatewayActiveWorkSnapshot }>(); + let deadline: ReturnType | undefined; + const nativeReply = { code: 0, stdout: "LoadState=loaded\nTimeoutStopUSec=330s", stderr: "" }; + if (refreshMs) { + readCgroup.mockResolvedValue("0::/system.slice/openclaw-gateway.service\n"); + systemctl.mockResolvedValue(nativeReply); + } + consumeGatewayRestartIntent.mockReturnValueOnce({ + force: true, + ...(waitMs === undefined ? {} : { waitMs }), + }); + createGatewayActiveWorkSnapshot.mockReturnValueOnce(active); + waitForGatewayActiveWork.mockImplementationOnce((timeoutMs) => { + if (timeoutMs !== undefined) { + deadline = setTimeout( + () => drain.resolve({ drained: false, snapshot: active }), + timeoutMs, + ); + } + return drain.promise; + }); + await withIsolatedSignals(async ({ captureSignal }) => { + const closing = createDeferredCore(); + const close = createCloseMock(); + if (stallClose) { + close.mockImplementationOnce(() => closing.promise); + } + const { start, started } = createSignaledStart(close); + const { runtime, exited } = createRuntimeWithExitSignal(); + const completeBoot = vi.fn(); + await runLoopWithStart({ start, runtime, completeBoot }); + await waitForStart(started); + vi.useFakeTimers(); + const clock = vi.spyOn(performance, "now").mockImplementation(() => Date.now()); + if (refreshMs) { + systemctl.mockImplementationOnce( + () => + new Promise((resolve) => { + setTimeout(() => resolve(nativeReply), refreshMs); + }), + ); + } + try { + captureSignal("SIGUSR2")(); + if (budget > 0) { + await vi.advanceTimersByTimeAsync(budget - 1); + expect(close).not.toHaveBeenCalled(); + expect(abortActiveCronTaskRuns).not.toHaveBeenCalled(); + expect(completeBoot).not.toHaveBeenCalled(); + } + await vi.advanceTimersByTimeAsync(budget > 0 ? 1 : 0); + expect(abortActiveCronTaskRuns).toHaveBeenCalledWith("Gateway restarting."); + expectRestartCloseCall(close, 0); + const warning = `restart drain budget ${budget - refreshMs}ms exhausted; cutting short cronRuns=1 activeTasks=1`; + expect(gatewayLog.warn).toHaveBeenCalledWith(warning); + if (stallClose) { + expect(start).toHaveBeenCalledOnce(); + expect(completeBoot).not.toHaveBeenCalled(); + expect(runtime.exit).not.toHaveBeenCalled(); + await vi.advanceTimersByTimeAsync(9_999); + expect(completeBoot).not.toHaveBeenCalled(); + expect(runtime.exit).not.toHaveBeenCalled(); + await vi.advanceTimersByTimeAsync(1); + expect(runtime.exit).toHaveBeenCalledExactlyOnceWith(1); + expect(completeBoot).toHaveBeenCalledExactlyOnceWith({ + outcome: "forced_stop", + reason: `${warning}; gateway.restart_shutdown_timeout`, + }); + expect(start).toHaveBeenCalledOnce(); + } else { + expect(start).toHaveBeenCalledTimes(2); + expect(completeBoot).toHaveBeenCalledExactlyOnceWith({ + outcome: "planned_restart", + reason: `${warning}; restart (SIGUSR2)`, + }); + } + } finally { + clearTimeout(deadline); + drain.resolve({ drained: true, snapshot: idleActiveWorkSnapshot }); + closing.resolve(); + await vi.advanceTimersByTimeAsync(0); + if (runtime.exit.mock.calls.length === 0) { + captureSignal("SIGINT")(); + await vi.advanceTimersByTimeAsync(0); + } + await exited; + clock.mockRestore(); + vi.useRealTimers(); + } + }); + }, + ); +} diff --git a/src/cli/gateway-cli/run-loop-package-helper.test-support.ts b/src/cli/gateway-cli/run-loop-package-helper.test-support.ts index a17d1ff62031..e7fb2d7ff2fb 100644 --- a/src/cli/gateway-cli/run-loop-package-helper.test-support.ts +++ b/src/cli/gateway-cli/run-loop-package-helper.test-support.ts @@ -9,11 +9,14 @@ import { createDeferred } from "../../../test/helpers/promise.js"; import type { HostedGatewayStop } from "../../daemon/hosted-stop.js"; import type { GatewayServer } from "../../gateway/server-public.js"; import { withTimeout } from "../../infra/fs-safe.js"; +import { resolveRuntimeWorkerUrl } from "../../infra/runtime-worker-url.js"; import { createManagedServiceBoundaryCleanup } from "../../infra/update-managed-service-handoff-process.test-support.js"; +import { updateExecutorEntrypoints } from "../cli-entrypoint.test-support.js"; import type { UpdateRespawnFixtures } from "./run-loop.test-support.js"; const removeFixturePath = fs.rm; -const sourceUrl = (file: string) => JSON.stringify(new URL(`../../${file}`, import.meta.url).href); +const sourceUrl = (key: keyof typeof updateExecutorEntrypoints) => + JSON.stringify(resolveRuntimeWorkerUrl(updateExecutorEntrypoints[key]).href); export async function startPackageLifecycleStopFixture(params: { fixtures: UpdateRespawnFixtures; @@ -251,17 +254,19 @@ export async function writePackageLifecycleFixture(root: string, control: string const bootstrap = ` import fs from "node:fs/promises"; import path from "node:path"; - const { register } = await import(${JSON.stringify(pathToFileURL(createRequire(import.meta.url).resolve("tsx/esm/api")).href)}); - register({ tsconfig: ${JSON.stringify(path.resolve("tsconfig.json"))} }); - const { registerSealedRuntime } = await import(${sourceUrl("infra/sealed-runtime-registry.ts")}); + if (${sourceUrl("sealedRegistry")}.endsWith(".ts")) { + const { register } = await import(${JSON.stringify(pathToFileURL(createRequire(import.meta.url).resolve("tsx/esm/api")).href)}); + register({ tsconfig: ${JSON.stringify(path.resolve("tsconfig.json"))} }); + } + const { registerSealedRuntime } = await import(${sourceUrl("sealedRegistry")}); registerSealedRuntime({ json5: undefined, resolveSecureTempRoot: () => ${JSON.stringify(control)} }); `; await fs.writeFile( path.join(root, "dist", "cli", "daemon-cli.js"), `${bootstrap} - const ledger = await import(${sourceUrl("infra/update-run-ledger.ts")}); + const ledger = await import(${sourceUrl("ledger")}); export const { adoptUpdateRun, finishUpdateRun, getUpdateRun, recordUpdateRunStep, recordUpdateRunVerification } = ledger; - const handoff = await import(${sourceUrl("infra/update-managed-service-handoff.ts")}); + const handoff = await import(${sourceUrl("handoff")}); export const { assertForegroundUpdateOrigin } = handoff; `, ); @@ -273,12 +278,12 @@ export async function writePackageLifecycleFixture(root: string, control: string if (process.argv[2] === "triage") { process.stdout.write(JSON.stringify({ diagnostic: "isolated lifecycle fixture" })); } else { - const handoff = await import(${sourceUrl("infra/update-managed-service-handoff.ts")}); - const { readControlPlaneUpdateSentinelMeta } = await import(${sourceUrl("infra/update-control-plane-sentinel.ts")}); - const { runGlobalPackageUpdateSteps } = await import(${sourceUrl("infra/package-update-steps.ts")}); - const { createNpmTarget, createRootRunner } = await import(${sourceUrl("infra/package-update-steps.test-support.ts")}); - const { writePackageDistInventory } = await import(${JSON.stringify(new URL("../../../scripts/lib/package-dist-inventory.ts", import.meta.url).href)}); - const { runCommandWithTimeout } = await import(${sourceUrl("process/exec.ts")}); + const handoff = await import(${sourceUrl("handoff")}); + const { readControlPlaneUpdateSentinelMeta } = await import(${sourceUrl("sentinel")}); + const { runGlobalPackageUpdateSteps } = await import(${sourceUrl("packageSteps")}); + const { createNpmTarget, createRootRunner } = await import(${sourceUrl("packageFixture")}); + const { writePackageDistInventory } = await import(${sourceUrl("inventory")}); + const { runCommandWithTimeout } = await import(${sourceUrl("exec")}); const meta = await readControlPlaneUpdateSentinelMeta(); const run = { runId: meta.runId, env: process.env }; let outcome; diff --git a/src/cli/gateway-cli/run-loop-request-fixtures.test-support.ts b/src/cli/gateway-cli/run-loop-request-fixtures.test-support.ts new file mode 100644 index 000000000000..9863b8408e43 --- /dev/null +++ b/src/cli/gateway-cli/run-loop-request-fixtures.test-support.ts @@ -0,0 +1,49 @@ +import type { Mock } from "vitest"; +import type { GatewayActiveWorkSnapshot } from "../../infra/gateway-active-work.js"; +import type { GatewayBootLifecycleCompletion } from "../../infra/gateway-boot-lifecycle.js"; +import type { GatewayRestartIntent } from "../../infra/restart-intent.js"; +import type { + createRuntimeWithExitSignal, + createSignaledStart, + UpdateRespawnFixtures, +} from "./run-loop.test-support.js"; + +export type RequestFixtures = { + createSignaledLoopHarness: UpdateRespawnFixtures["createSignaledLoopHarness"]; + createGatewayActiveWorkSnapshot: Mock<() => GatewayActiveWorkSnapshot>; + abortActiveCronTaskRuns: Mock<(_reason?: string) => number>; + acquireGatewayLock: Mock< + (opts?: { port?: number }) => Promise<{ release: Mock<() => Promise> }> + >; + reloadTaskRuntimeStateFromStore: Mock<() => Promise>; + runLoopWithStart: (params: { + start: ReturnType["start"]; + runtime: ReturnType["runtime"]; + ownsProcessLifecycle?: boolean; + beginBoot?: (startedAtMs: number) => void | Promise; + completeBoot?: (completion: GatewayBootLifecycleCompletion) => void; + }) => Promise; + waitForGatewayActiveWork: Mock< + typeof import("../../infra/gateway-active-work.js").waitForGatewayActiveWork + >; + restartGatewayProcessWithFreshPid: Mock< + typeof import("../../infra/process-respawn.js").restartGatewayProcessWithFreshPid + >; + respawnGatewayProcessForUpdate: UpdateRespawnFixtures["respawnGatewayProcessForUpdate"]; + captureForegroundUpdateHandoffStop: UpdateRespawnFixtures["captureForegroundUpdateHandoffStop"]; + readCgroup: Mock; + systemctl: Mock; + armShutdownHardExitWatchdog: Mock; + cancelShutdownHardExitWatchdog: Mock; + consumeGatewayRestartIntent: Mock<() => GatewayRestartIntent | null>; + consumeGatewayRestartIntentPayloadSync: Mock< + () => Pick | null + >; + peekGatewayRestartReason: Mock<() => string | undefined>; + managedUpdateSuccessorOwner: NonNullable; + commitManagedServiceUpdateHandoff: Mock< + typeof import("../../infra/update-managed-service-handoff.js").commitManagedServiceUpdateHandoff + >; + isGatewayWorkAdmissionClosed: () => boolean; + gatewayLog: { info: Mock; warn: Mock; error: Mock }; +}; diff --git a/src/cli/gateway-cli/run-loop-request.test-support.ts b/src/cli/gateway-cli/run-loop-request.test-support.ts index 88d542852f05..07fac8a4928b 100644 --- a/src/cli/gateway-cli/run-loop-request.test-support.ts +++ b/src/cli/gateway-cli/run-loop-request.test-support.ts @@ -1,11 +1,11 @@ /** Shutdown request reasons and installation-replacement handoff cases share the run-loop fixture. */ import { fileURLToPath, pathToFileURL } from "node:url"; -import { expect, it, vi, type Mock } from "vitest"; +import { expect, it, vi } from "vitest"; import { withTimeout } from "../../infra/fs-safe.js"; import type { GatewayActiveWorkSnapshot } from "../../infra/gateway-active-work.js"; -import type { GatewayBootLifecycleCompletion } from "../../infra/gateway-boot-lifecycle.js"; -import type { GatewayRestartIntent } from "../../infra/restart-intent.js"; import { createDeferredCore } from "../../shared/deferred.js"; +import { registerGatewayForcedRestartTests } from "./run-loop-force.test-support.js"; +import type { RequestFixtures } from "./run-loop-request-fixtures.test-support.js"; import { createActiveWorkSnapshot, createCloseMock, @@ -14,47 +14,12 @@ import { waitForStart, waitForLoopCondition, withIsolatedSignals, - type UpdateRespawnFixtures, } from "./run-loop.test-support.js"; -type RequestFixtures = { - acquireGatewayLock: Mock< - (opts?: { port?: number }) => Promise<{ release: Mock<() => Promise> }> - >; - reloadTaskRuntimeStateFromStore: Mock<() => Promise>; - runLoopWithStart: (params: { - start: ReturnType["start"]; - runtime: ReturnType["runtime"]; - ownsProcessLifecycle?: boolean; - beginBoot?: (startedAtMs: number) => void | Promise; - completeBoot?: (completion: GatewayBootLifecycleCompletion) => void; - }) => Promise; - waitForGatewayActiveWork: Mock< - typeof import("../../infra/gateway-active-work.js").waitForGatewayActiveWork - >; - restartGatewayProcessWithFreshPid: Mock< - typeof import("../../infra/process-respawn.js").restartGatewayProcessWithFreshPid - >; - respawnGatewayProcessForUpdate: UpdateRespawnFixtures["respawnGatewayProcessForUpdate"]; - captureForegroundUpdateHandoffStop: UpdateRespawnFixtures["captureForegroundUpdateHandoffStop"]; - readCgroup: Mock; - systemctl: Mock; - armShutdownHardExitWatchdog: Mock; - cancelShutdownHardExitWatchdog: Mock; - consumeGatewayRestartIntent: Mock<() => GatewayRestartIntent | null>; - consumeGatewayRestartIntentPayloadSync: Mock< - () => Pick | null - >; - peekGatewayRestartReason: Mock<() => string | undefined>; - managedUpdateSuccessorOwner: NonNullable; - commitManagedServiceUpdateHandoff: Mock< - typeof import("../../infra/update-managed-service-handoff.js").commitManagedServiceUpdateHandoff - >; - isGatewayWorkAdmissionClosed: () => boolean; - gatewayLog: { info: Mock; error: Mock }; -}; - export function registerGatewayRequestTests({ + createSignaledLoopHarness, + createGatewayActiveWorkSnapshot, + abortActiveCronTaskRuns, acquireGatewayLock, reloadTaskRuntimeStateFromStore, runLoopWithStart, @@ -75,6 +40,20 @@ export function registerGatewayRequestTests({ gatewayLog, }: RequestFixtures): void { const idleActiveWorkSnapshot = createActiveWorkSnapshot(); + registerGatewayForcedRestartTests({ + createSignaledLoopHarness, + createGatewayActiveWorkSnapshot, + abortActiveCronTaskRuns, + runLoopWithStart, + waitForGatewayActiveWork, + consumeGatewayRestartIntent, + consumeGatewayRestartIntentPayloadSync, + isGatewayWorkAdmissionClosed, + gatewayLog, + readCgroup, + systemctl, + }); + it("keeps a captured pre-park Stop ahead of native budget refresh and drain completion", async () => { const nativeReply = { code: 0, diff --git a/src/cli/gateway-cli/run-loop-shutdown-budget.ts b/src/cli/gateway-cli/run-loop-shutdown-budget.ts index 64471c15bf35..2727d461daae 100644 --- a/src/cli/gateway-cli/run-loop-shutdown-budget.ts +++ b/src/cli/gateway-cli/run-loop-shutdown-budget.ts @@ -72,27 +72,39 @@ export async function resolveGatewayShutdownBudget( export function resolveGatewayShutdownDrainBudget(params: { budget: { nativeStopBudget: boolean; timeoutMs: number; reserveMs: number }; isRestart: boolean; + forceRestart: boolean; restartWithoutSupervisor: boolean; + acceptedAtMs: number; requestedRestartDrainTimeoutMs?: number; }) { const { budget, isRestart } = params; const requested = params.requestedRestartDrainTimeoutMs; + const elapsedMs = performance.now() - params.acceptedAtMs; + const remaining = requested === undefined ? undefined : Math.max(0, requested - elapsedMs); const restartDrainTimeoutMs = budget.nativeStopBudget - ? Math.min(requested ?? Infinity, Math.max(0, budget.timeoutMs - budget.reserveMs)) - : requested; + ? Math.min(remaining ?? Infinity, Math.max(0, budget.timeoutMs - budget.reserveMs)) + : remaining; const restartDrainDeadlineAt = isRestart && restartDrainTimeoutMs !== undefined ? Date.now() + restartDrainTimeoutMs : undefined; - const restartTimeoutMs = (drainTimeoutMs: number) => + const forcedRestartDeadlineAt = + params.forceRestart && restartDrainDeadlineAt !== undefined + ? restartDrainDeadlineAt + budget.reserveMs + : undefined; + const restartTimeoutMs = (drainTimeoutMs: number) => { + if (forcedRestartDeadlineAt !== undefined) { + return Math.max(0, forcedRestartDeadlineAt - Date.now()); + } // A containing service can bound an in-process restart without replacing it. - budget.nativeStopBudget && params.restartWithoutSupervisor + return budget.nativeStopBudget && params.restartWithoutSupervisor ? budget.timeoutMs : drainTimeoutMs + (budget.nativeStopBudget ? budget.reserveMs : GATEWAY_SHUTDOWN_TIMEOUT_MS); + }; return { restartDrainDeadlineAt, restartTimeoutMs: () => - budget.nativeStopBudget + budget.nativeStopBudget || params.forceRestart ? restartTimeoutMs(Math.max(0, (restartDrainDeadlineAt ?? Date.now()) - Date.now())) : GATEWAY_SHUTDOWN_TIMEOUT_MS, closeDrainTimeoutMs: () => diff --git a/src/cli/gateway-cli/run-loop-shutdown-completion.test-support.ts b/src/cli/gateway-cli/run-loop-shutdown-completion.test-support.ts index 8c19675040fb..f22107c4157d 100644 --- a/src/cli/gateway-cli/run-loop-shutdown-completion.test-support.ts +++ b/src/cli/gateway-cli/run-loop-shutdown-completion.test-support.ts @@ -124,13 +124,14 @@ export function registerShutdownCompletionTests({ ); it.each([true, false])( - "bounds abandoned cleanup after managed parking (restore commit=%s)", + "bounds abandoned cleanup after an exhausted deferral and managed parking (restore commit=%s)", async (restoreCommitted) => { vi.clearAllMocks(); process.env.OPENCLAW_SYSTEMD_UNIT = "openclaw-gateway.service"; setPlatform("linux"); consumeGatewayRestartIntent.mockReturnValueOnce({ force: true, + drainBudgetExhausted: true, reason: "update.run", successorOwner: managedUpdateSuccessorOwner, }); @@ -172,14 +173,16 @@ export function registerShutdownCompletionTests({ it("retains external supervisor recovery when timeout prevents a restart handoff", async () => { vi.clearAllMocks(); process.env.OPENCLAW_SUPERVISOR_MODE = "external"; - consumeGatewayRestartIntent.mockReturnValueOnce({ force: true }); + consumeGatewayRestartIntent.mockReturnValueOnce({ force: true, drainBudgetExhausted: true }); await withIsolatedSignals(async ({ captureSignal }) => { const { close, runtime } = await createSignaledLoopHarness(); close.mockReturnValue(new Promise(() => {})); vi.useFakeTimers(); try { captureSignal("SIGUSR2")(); - await vi.advanceTimersByTimeAsync(325_000); + await vi.advanceTimersByTimeAsync(9_999); + expect(runtime.exit).not.toHaveBeenCalled(); + await vi.advanceTimersByTimeAsync(1); expect(writeGatewayRestartHandoffSync).not.toHaveBeenCalled(); expect(runtime.exit).toHaveBeenCalledExactlyOnceWith(1); } finally { @@ -221,7 +224,7 @@ export function registerShutdownCompletionTests({ it("retains the restart deadline when a managed update arrives after final cleanup fails", async () => { vi.clearAllMocks(); process.env.OPENCLAW_SUPERVISOR_MODE = "external"; - consumeGatewayRestartIntent.mockReturnValueOnce({ force: true }); + consumeGatewayRestartIntent.mockReturnValueOnce({ force: true, drainBudgetExhausted: true }); restartGatewayProcessWithFreshPid.mockReturnValueOnce({ mode: "supervised" }); await withIsolatedSignals(async ({ captureSignal }) => { const { close, start, runtime, exited } = await createSignaledLoopHarness(undefined, true); @@ -230,7 +233,7 @@ export function registerShutdownCompletionTests({ vi.useFakeTimers(); try { restartSignal(); - await vi.advanceTimersByTimeAsync(10_000); + await vi.advanceTimersByTimeAsync(0); expect(close).toHaveBeenCalledOnce(); expect(gatewayLog.error).toHaveBeenCalledWith( "gateway lifecycle completion failed: shutdown cleanup failed", @@ -241,7 +244,7 @@ export function registerShutdownCompletionTests({ successorOwner: managedUpdateSuccessorOwner, }); restartSignal(); - await vi.advanceTimersByTimeAsync(314_999); + await vi.advanceTimersByTimeAsync(9_999); expect(runtime.exit).not.toHaveBeenCalled(); await vi.advanceTimersByTimeAsync(1); expect(runtime.exit).toHaveBeenCalledExactlyOnceWith(1); diff --git a/src/cli/gateway-cli/run-loop-shutdown-format.ts b/src/cli/gateway-cli/run-loop-shutdown-format.ts index d776de6e2d94..ffcbae883c64 100644 --- a/src/cli/gateway-cli/run-loop-shutdown-format.ts +++ b/src/cli/gateway-cli/run-loop-shutdown-format.ts @@ -1,10 +1,30 @@ import { truncateUtf16Safe } from "@openclaw/normalization-core/utf16-slice"; import type { GatewayActiveWorkSnapshot } from "../../infra/gateway-active-work.js"; +import { + GATEWAY_BOOT_REASON_MAX_UTF16_CODE_UNITS, + type GatewayBootLifecycleCompletion, +} from "../../infra/gateway-boot-lifecycle.js"; import type { GatewayDrainReason, GatewayShutdownTrigger, } from "../../process/gateway-work-admission.js"; +export function formatBootCompletionContext( + completion: GatewayBootLifecycleCompletion, + ...reasons: (string | undefined)[] +): GatewayBootLifecycleCompletion { + const context = reasons.filter(Boolean).join("; "); + return context + ? { + ...completion, + reason: truncateUtf16Safe( + `${context}; ${completion.reason ?? completion.outcome}`, + GATEWAY_BOOT_REASON_MAX_UTF16_CODE_UNITS, + ), + } + : completion; +} + export function formatShutdownReason(request: { action: "stop" | "restart" | "external-restart"; signal: GatewayShutdownTrigger; diff --git a/src/cli/gateway-cli/run-loop.forced-cron.process.test.ts b/src/cli/gateway-cli/run-loop.forced-cron.process.test.ts index c8fd3aaad826..27f839ed367f 100644 --- a/src/cli/gateway-cli/run-loop.forced-cron.process.test.ts +++ b/src/cli/gateway-cli/run-loop.forced-cron.process.test.ts @@ -1,8 +1,8 @@ import { spawn, type ChildProcess } from "node:child_process"; -import { once } from "node:events"; +import { EventEmitter, once } from "node:events"; import fs from "node:fs"; import path from "node:path"; -import { afterEach, expect, it, vi } from "vitest"; +import { afterEach, expect, it } from "vitest"; import { withTestTimeout } from "../../../test/helpers/promise.js"; import { useAutoCleanupTempDirTracker } from "../../../test/helpers/temp-dir.js"; import { @@ -35,7 +35,7 @@ it.skipIf(process.platform === "win32").each([ { signal: "SIGTERM", mode: "force" }, { signal: "SIGUSR2", mode: "timeout" }, ] as const)( - "cancels active cron work and joins cleanup before $signal $mode restart", + "settles admitted cron work before $signal $mode restart", async ({ signal, mode }) => { const root = tempDirs.make("openclaw-forced-cron-"); const home = path.join(root, "home"); @@ -56,23 +56,42 @@ it.skipIf(process.platform === "win32").each([ children.set(child, closed); void closed.catch(() => {}); let output = ""; - child.stdout?.on("data", (chunk: Buffer) => { + const changes = new EventEmitter(); + const recordOutput = (chunk: Buffer) => { output += chunk.toString(); - }); - child.stderr?.on("data", (chunk: Buffer) => { - output += chunk.toString(); - }); - const waitForOutput = (text: string, timeout = 5_000) => - vi.waitFor( - () => { - expect(output).toContain(text); - }, - { timeout, interval: 25 }, - ); + changes.emit("output"); + }; + child.stdout?.on("data", recordOutput); + child.stderr?.on("data", recordOutput); + const waitForOutput = async (text: string, timeout = 5_000) => { + let inspect: () => void = () => {}; + try { + await withTestTimeout( + new Promise((resolve) => { + inspect = () => { + if (output.includes(text)) { + resolve(); + } + }; + changes.on("output", inspect); + inspect(); + }), + timeout, + `Missing ${text}: ${output}`, + ); + } finally { + changes.off("output", inspect); + } + }; await waitForOutput("process proof: ready:1", 45_000); expect(child.kill(signal)).toBe(true); - await waitForOutput("process proof: cron-cancelled:Gateway restarting."); - await waitForOutput("process proof: close-entered"); + if (mode === "timeout") { + await waitForOutput("process proof: cron-cancelled:Gateway restarting."); + await waitForOutput("process proof: close-entered"); + } else { + await waitForOutput("draining active work before"); + expect(output).not.toContain("process proof: close-entered"); + } child.send("inspect"); await waitForOutput("process proof: held:starts=1:pending=true"); expect(fs.existsSync(path.join(root, "cleanup.txt"))).toBe(false); diff --git a/src/cli/gateway-cli/run-loop.forced-cron.test-support.ts b/src/cli/gateway-cli/run-loop.forced-cron.test-support.ts index 63884b9faec5..f1c513c0c4e8 100644 --- a/src/cli/gateway-cli/run-loop.forced-cron.test-support.ts +++ b/src/cli/gateway-cli/run-loop.forced-cron.test-support.ts @@ -43,9 +43,14 @@ const cron = new CronService({ onExecutionStarted?.(); coreStarted.resolve(); trace("cron-started"); - await cancelled.promise; - trace(`cron-cancelled:${String(abortSignal.reason)}`); + if (!force) { + await cancelled.promise; + trace(`cron-cancelled:${String(abortSignal.reason)}`); + } await cleanupMayFinish.promise; + if (force) { + assert(!abortSignal.aborted, "force restart cancelled admitted work before its budget"); + } await fs.writeFile(path.join(root, "cleanup.txt"), "settled\n"); trace("cron-cleanup-settled"); return { status: "ok" as const, summary: "settled" }; diff --git a/src/cli/gateway-cli/run-loop.test-support.ts b/src/cli/gateway-cli/run-loop.test-support.ts index a9972f598e67..9399e8778ff3 100644 --- a/src/cli/gateway-cli/run-loop.test-support.ts +++ b/src/cli/gateway-cli/run-loop.test-support.ts @@ -557,9 +557,9 @@ export function registerGatewayRestartOwnershipTests({ await vi.advanceTimersByTimeAsync(5_000); expect(close).toHaveBeenCalledOnce(); expect(runtime.exit).not.toHaveBeenCalled(); - await vi.advanceTimersByTimeAsync(outcome === "completed" ? 1_000 : 5_001); + await vi.advanceTimersByTimeAsync(outcome === "completed" ? 1_000 : 80_001); await expect(exited).resolves.toBe(outcome === "completed" ? 0 : 1); - expect(cleanupDeadline).toBe(10_000); + expect(cleanupDeadline).toBe(55_000); expect(start).toHaveBeenCalledOnce(); } finally { clock.mockRestore(); @@ -581,6 +581,7 @@ export function registerGatewayRestartOwnershipTests({ consumeGatewayRestartIntent.mockReturnValueOnce({ reason: "update.run", force: true, + waitMs: 300_000, successorOwner: managedUpdateSuccessorOwner, }); isForegroundUpdateHandoff.mockReturnValue(true); diff --git a/src/cli/gateway-cli/run-loop.test.ts b/src/cli/gateway-cli/run-loop.test.ts index ca1e82057600..15f78bdd9a62 100644 --- a/src/cli/gateway-cli/run-loop.test.ts +++ b/src/cli/gateway-cli/run-loop.test.ts @@ -542,6 +542,9 @@ afterEach(() => { describe("runGatewayLoop", () => { registerGatewayRequestTests({ + createSignaledLoopHarness, + createGatewayActiveWorkSnapshot, + abortActiveCronTaskRuns, acquireGatewayLock, reloadTaskRuntimeStateFromStore, runLoopWithStart, @@ -626,8 +629,8 @@ describe("runGatewayLoop", () => { ); sigterm(); expect(host?.externalRestart?.isCurrent()).toBe(false); - expectRestartCloseCall(close, 0); - expect(waitForGatewayActiveWork).not.toHaveBeenCalled(); + expectRestartCloseCall(close, DEFAULT_RESTART_DEFERRAL_TIMEOUT_MS); + expect(waitForGatewayActiveWork).toHaveBeenCalledOnce(); expect(runtime.exit).not.toHaveBeenCalled(); } finally { joined.resolve(); @@ -1749,6 +1752,7 @@ describe("runGatewayLoop", () => { it("hands timed-out active work to server close", async () => { vi.clearAllMocks(); + const clock = vi.spyOn(performance, "now").mockReturnValue(0); consumeGatewayRestartIntentPayloadSync.mockReturnValueOnce({}); const timedOutSnapshot = createActiveWorkSnapshot({ activeTasks: 1, embeddedRuns: 1 }, [ { kind: "task", count: 1, message: "1 active background task run(s)" }, @@ -1777,19 +1781,19 @@ describe("runGatewayLoop", () => { DEFAULT_RESTART_DEFERRAL_TIMEOUT_MS, ); expect(gatewayLog.warn).toHaveBeenCalledWith( - "active-work drain timeout reached; proceeding with restart: embeddedRuns=1 activeTasks=1", + "restart drain budget 300000ms exhausted; cutting short embeddedRuns=1 activeTasks=1", ); expectRestartCloseCall(close, DEFAULT_RESTART_DEFERRAL_TIMEOUT_MS); expect(start).toHaveBeenCalledOnce(); await expect(exited).resolves.toBe(0); - }); + }).finally(() => clock.mockRestore()); }); it("skips a second active-work drain after a SIGUSR2 deferral timeout intent", async () => { - vi.clearAllMocks(); consumeGatewayRestartIntent.mockReturnValueOnce({ force: true, + drainBudgetExhausted: true, reason: "config reload forced restart", }); createGatewayActiveWorkSnapshot.mockReturnValue( @@ -1812,7 +1816,7 @@ describe("runGatewayLoop", () => { setImmediate(resolve); }); - expect(waitForGatewayActiveWork).not.toHaveBeenCalled(); + expect(waitForGatewayActiveWork).toHaveBeenCalledWith(0, expect.any(Object)); expect(markGatewayRestartHandled).toHaveBeenCalledOnce(); expectRestartCloseCall(close, 0); expect(start).toHaveBeenCalledTimes(2); @@ -1822,46 +1826,6 @@ describe("runGatewayLoop", () => { }); }); - it("forces SIGTERM restarts without waiting for active task drain", async () => { - vi.clearAllMocks(); - consumeGatewayRestartIntentPayloadSync.mockReturnValueOnce({ force: true }); - createGatewayActiveWorkSnapshot.mockReturnValue( - createActiveWorkSnapshot({ activeTasks: 1, embeddedRuns: 1 }, [ - { - kind: "task", - count: 1, - message: "taskId=task-force runId=run-force status=running runtime=cron label=forced", - }, - { kind: "embedded-run", count: 1, message: "1 active embedded run(s)" }, - ]), - ); - await withIsolatedSignals(async ({ captureSignal }) => { - const { close, start, exited } = await createSignaledLoopHarness(); - const sigterm = captureSignal("SIGTERM"); - - sigterm(); - await new Promise((resolve) => { - setImmediate(resolve); - }); - await new Promise((resolve) => { - setImmediate(resolve); - }); - - expect(waitForGatewayActiveWork).not.toHaveBeenCalled(); - expect(gatewayLog.info).toHaveBeenCalledWith( - expect.stringContaining("embeddedRuns=1 activeTasks=1"), - ); - expect(gatewayLog.info.mock.calls.flat().join("\n")).not.toContain("task-force"); - expect(gatewayLog.warn).toHaveBeenCalledWith( - "forced restart requested; skipping active work drain", - ); - expectRestartCloseCall(close, 0); - expect(start).toHaveBeenCalledOnce(); - - await expect(exited).resolves.toBe(0); - }); - }); - registerGatewayRestartOwnershipTests({ consumeGatewayRestartIntentPayloadSync, readCgroup, @@ -1883,6 +1847,7 @@ describe("runGatewayLoop", () => { it("restarts after SIGUSR2 even when drain times out, and resets runtime state for the new iteration", async () => { vi.clearAllMocks(); + const clock = vi.spyOn(performance, "now").mockReturnValue(0); peekGatewayRestartReason.mockReturnValue(undefined); respawnGatewayProcessForUpdate.mockReturnValue({ mode: "disabled", @@ -1993,7 +1958,7 @@ describe("runGatewayLoop", () => { DEFAULT_RESTART_DEFERRAL_TIMEOUT_MS, ); expect(gatewayLog.warn).toHaveBeenCalledWith( - "active-work drain timeout reached; proceeding with restart: embeddedRuns=1 activeTasks=2", + "restart drain budget 300000ms exhausted; cutting short embeddedRuns=1 activeTasks=2", ); expectRestartCloseCall(closeFirst, DEFAULT_RESTART_DEFERRAL_TIMEOUT_MS); await startedThird; @@ -2039,7 +2004,7 @@ describe("runGatewayLoop", () => { reason: "gateway stopping", restartExpectedMs: null, }); - }); + }).finally(() => clock.mockRestore()); }); it("advances stale cron active markers after bounded restart cron-run drain", async () => { diff --git a/src/cli/gateway-cli/run-loop.ts b/src/cli/gateway-cli/run-loop.ts index 47c509f1551e..2105a5d0ebed 100644 --- a/src/cli/gateway-cli/run-loop.ts +++ b/src/cli/gateway-cli/run-loop.ts @@ -37,7 +37,7 @@ import type { RuntimeEnv } from "../../runtime.js"; import { createLazyImportLoader } from "../../shared/lazy-promise.js"; import { formatCliCommand } from "../command-format.js"; import { createGatewayHostLifecycle } from "./host-lifecycle.js"; -import { drainGatewayActiveWork, resolveRestartDrainTimeoutMs } from "./run-loop-drain.js"; +import { drainGatewayActiveWork } from "./run-loop-drain.js"; import * as loopLogs from "./run-loop-log-flush.js"; import { isUpdateProcessRestartReason, @@ -50,7 +50,7 @@ import { resolveGatewayShutdownDrainBudget, resolveGatewayShutdownBudget, } from "./run-loop-shutdown-budget.js"; -import { formatShutdownReason } from "./run-loop-shutdown-format.js"; +import { formatBootCompletionContext, formatShutdownReason } from "./run-loop-shutdown-format.js"; import { createGatewayStartupOperations, prepareGatewayRestartIteration, @@ -134,19 +134,13 @@ export async function runGatewayLoop(params: { let pendingStartupForceExitTimer: ReturnType | null = null; let installationReplacement: GatewayInstallationReplacement | undefined; let pendingRestartCompletion: GatewayBootLifecycleCompletion | undefined; + let restartDrainWarning: string | undefined; const completeBoot = (completion: GatewayBootLifecycleCompletion) => { pendingRestartCompletion = undefined; - params.completeBoot?.({ - ...completion, - ...(installationReplacement - ? { - reason: truncateUtf16Safe( - `${installationReplacement.reason}; ${completion.reason ?? completion.outcome}`, - GATEWAY_BOOT_REASON_MAX_UTF16_CODE_UNITS, - ), - } - : {}), - }); + params.completeBoot?.( + formatBootCompletionContext(completion, installationReplacement?.reason, restartDrainWarning), + ); + restartDrainWarning = undefined; }; let restartDrainingMarked = false; const observeSignal = loopLogs.createGatewaySignalObserver(gatewayLog); @@ -793,9 +787,11 @@ export async function runGatewayLoop(params: { const drainBudget = resolveGatewayShutdownDrainBudget({ budget, isRestart, + forceRestart: Boolean(restartIntent?.force || restartIntent?.drainBudgetExhausted), restartWithoutSupervisor, + acceptedAtMs: acceptedRequest.acceptedAtMs, requestedRestartDrainTimeoutMs: isRestart - ? resolveRestartDrainTimeoutMs(restartIntent, eagerLifecycleRuntime) + ? eagerLifecycleRuntime.resolveGatewayRestartDrainTimeoutMs(restartIntent) : 0, }); // Managed helpers must reach native parking before either exit watchdog can arm. @@ -812,15 +808,16 @@ export async function runGatewayLoop(params: { shutdownStep = "active-work-drain"; await drainGatewayActiveWork({ request: acceptedRequest, - restartIntent, runtime: eagerLifecycleRuntime, - loadRuntime: () => gatewayLifecycleRuntimeLoader.load(), drainTimeoutMs: drainBudget.drainTimeoutMs, restartDrainDeadlineAt: drainBudget.restartDrainDeadlineAt, markDraining: markRestartDraining, recordCounts: (counts) => { lastDrainCounts = counts; }, + recordWarning: (warning) => { + restartDrainWarning = warning; + }, logger: gatewayLog, }); diff --git a/src/cli/update-cli/update-command-handoff.ts b/src/cli/update-cli/update-command-handoff.ts index a67472d1b25d..622af3cf6e2a 100644 --- a/src/cli/update-cli/update-command-handoff.ts +++ b/src/cli/update-cli/update-command-handoff.ts @@ -7,8 +7,8 @@ import { } from "../../daemon/service-types.js"; import { resolveSystemdServiceName } from "../../daemon/systemd-service-files.js"; import { resolveInstallationTarget } from "../../infra/installation-target-context.js"; +import { resolveGatewayRestartDeferralTimeoutMs } from "../../infra/restart-budget.js"; import { getSelfAndAncestorPidsSync } from "../../infra/restart-stale-pids.js"; -import { resolveGatewayRestartDeferralTimeoutMs } from "../../infra/restart.js"; import { detectRespawnSupervisor } from "../../infra/supervisor-markers.js"; import { normalizeUpdateChannel } from "../../infra/update-channels.js"; import { diff --git a/src/commands/doctor-maintenance-foreground.ts b/src/commands/doctor-maintenance-foreground.ts index 0a18ff6f09b3..02a32ca441e4 100644 --- a/src/commands/doctor-maintenance-foreground.ts +++ b/src/commands/doctor-maintenance-foreground.ts @@ -5,7 +5,7 @@ import { GATEWAY_SERVICE_STOP_TIMEOUT_MS, GATEWAY_SHUTDOWN_RESERVE_MS, } from "../infra/gateway-shutdown-budget.js"; -import { resolveGatewayRestartDeferralTimeoutMs } from "../infra/restart.js"; +import { resolveGatewayRestartDeferralTimeoutMs } from "../infra/restart-budget.js"; import { acquireGatewayMaintenanceCoordinator, StateDatabaseCoordinatorContentionError, diff --git a/src/commands/doctor-maintenance.foreground-update.test.ts b/src/commands/doctor-maintenance.foreground-update.test.ts index 46e33c27ff32..a70e196ca60f 100644 --- a/src/commands/doctor-maintenance.foreground-update.test.ts +++ b/src/commands/doctor-maintenance.foreground-update.test.ts @@ -9,7 +9,7 @@ import { GATEWAY_SHUTDOWN_RESERVE_MS, GATEWAY_SHUTDOWN_TIMEOUT_MS, } from "../infra/gateway-shutdown-budget.js"; -import { resolveGatewayRestartDeferralTimeoutMs } from "../infra/restart.js"; +import { resolveGatewayRestartDeferralTimeoutMs } from "../infra/restart-budget.js"; import { tryAcquireExclusiveSqliteCoordinator } from "../infra/sqlite-coordinator.js"; import { acquireGatewayLifecycleCoordinator } from "../infra/state-database-coordinator.js"; import { createDeferredCore } from "../shared/deferred.js"; diff --git a/src/gateway/server-methods/restart-request.ts b/src/gateway/server-methods/restart-request.ts index 5f3287904ff9..fefb8b94d9dc 100644 --- a/src/gateway/server-methods/restart-request.ts +++ b/src/gateway/server-methods/restart-request.ts @@ -107,19 +107,22 @@ export function parseTargetedGatewayRestartIntent( if (value !== undefined && (!value || typeof value !== "object" || Array.isArray(value))) { return null; } - const raw = (value ?? {}) as { force?: unknown; waitMs?: unknown }; + const raw = (value ?? {}) as { force?: unknown; waitMs?: unknown; drainBudgetMs?: unknown }; const force = raw.force === true; + // Older Gateways ignore this optional field instead of rejecting force + waitMs. + const budget = force ? raw.drainBudgetMs : raw.waitMs; const waitMs = - typeof raw.waitMs === "number" && - Number.isSafeInteger(raw.waitMs) && - raw.waitMs >= 0 && - raw.waitMs <= MAX_TIMER_TIMEOUT_MS - ? raw.waitMs + typeof budget === "number" && + Number.isSafeInteger(budget) && + budget >= 0 && + budget <= MAX_TIMER_TIMEOUT_MS + ? budget : undefined; if ( (raw.force !== undefined && typeof raw.force !== "boolean") || - (raw.waitMs !== undefined && waitMs === undefined) || - (force && waitMs !== undefined) + (budget !== undefined && waitMs === undefined) || + (force && raw.waitMs !== undefined) || + (!force && raw.drainBudgetMs !== undefined) ) { return null; } diff --git a/src/gateway/server-methods/restart.test.ts b/src/gateway/server-methods/restart.test.ts index a8c95e5bf8ad..c894919174ee 100644 --- a/src/gateway/server-methods/restart.test.ts +++ b/src/gateway/server-methods/restart.test.ts @@ -196,28 +196,40 @@ describe("gateway restart handlers", () => { expectRestartRequest(false); }); - it("delivers a targeted restart only to the matching lock owner", async () => { - const respond = await invokeRestartRequest({ - reason: "operator", - target: { - pid: process.pid, - ownerId: "gateway-owner", - port: 18_789, - }, - restartIntent: { waitMs: 30_000 }, - }); + it.each([ + { restartIntent: { waitMs: 30_000 }, expected: { waitMs: 30_000 } }, + { restartIntent: { waitMs: 0 }, expected: { waitMs: 0 } }, + { restartIntent: { force: true }, expected: { force: true } }, + { restartIntent: { force: true, drainBudgetMs: 0 }, expected: { force: true, waitMs: 0 } }, + { + restartIntent: { force: true, drainBudgetMs: 180_000 }, + expected: { force: true, waitMs: 180_000 }, + }, + ])( + "delivers a targeted restart only to the matching lock owner ($restartIntent)", + async ({ restartIntent, expected }) => { + const respond = await invokeRestartRequest({ + reason: "operator", + target: { + pid: process.pid, + ownerId: "gateway-owner", + port: 18_789, + }, + restartIntent, + }); - expect(scheduleSafeGatewayRestart).not.toHaveBeenCalled(); - expect(requestGatewayRestartWithSignalAdmission).toHaveBeenCalledWith("operator", { - reason: "operator", - waitMs: 30_000, - }); - expect(respond).toHaveBeenCalledWith(true, { - ok: true, - status: "emitted", - pid: process.pid, - }); - }); + expect(scheduleSafeGatewayRestart).not.toHaveBeenCalled(); + expect(requestGatewayRestartWithSignalAdmission).toHaveBeenCalledWith("operator", { + reason: "operator", + ...expected, + }); + expect(respond).toHaveBeenCalledWith(true, { + ok: true, + status: "emitted", + pid: process.pid, + }); + }, + ); it("schedules a safe restart only after matching the target lock owner", async () => { mockScheduledRestart({ safe: false, summary: "restart deferred" }); @@ -390,26 +402,28 @@ describe("gateway restart handlers", () => { }); }); - it.each([0.5, -1, MAX_TIMER_TIMEOUT_MS + 1, Number.MAX_SAFE_INTEGER])( - "rejects an invalid targeted restart wait of %s ms", - async (waitMs) => { - const respond = await invokeRestartRequest({ - reason: "operator", - target: { - pid: process.pid, - ownerId: "gateway-owner", - port: 18_789, - }, - restartIntent: { waitMs }, - }); + it.each( + [0.5, -1, MAX_TIMER_TIMEOUT_MS + 1, Number.MAX_SAFE_INTEGER].flatMap((budget) => [ + { waitMs: budget }, + { force: true, drainBudgetMs: budget }, + ]), + )("rejects an invalid targeted restart budget: %j", async (restartIntent) => { + const respond = await invokeRestartRequest({ + reason: "operator", + target: { + pid: process.pid, + ownerId: "gateway-owner", + port: 18_789, + }, + restartIntent, + }); - expect(requestGatewayRestartWithSignalAdmission).not.toHaveBeenCalled(); - expect(respond).toHaveBeenCalledWith(false, undefined, { - code: "INVALID_REQUEST", - message: "invalid targeted gateway restart intent", - }); - }, - ); + expect(requestGatewayRestartWithSignalAdmission).not.toHaveBeenCalled(); + expect(respond).toHaveBeenCalledWith(false, undefined, { + code: "INVALID_REQUEST", + message: "invalid targeted gateway restart intent", + }); + }); it("accepts the maximum timer-safe targeted restart wait", async () => { const respond = await invokeRestartRequest({ diff --git a/src/gateway/server-methods/update.ts b/src/gateway/server-methods/update.ts index 682f0bc99f98..3ba538581072 100644 --- a/src/gateway/server-methods/update.ts +++ b/src/gateway/server-methods/update.ts @@ -24,17 +24,14 @@ import { isGatewayExternallySupervised, } from "../../infra/gateway-supervision.js"; import { readPackageVersion } from "../../infra/package-json.js"; +import { resolveGatewayRestartDeferralTimeoutMs } from "../../infra/restart-budget.js"; import type { GatewayRestartIntent } from "../../infra/restart-intent.js"; import { type RestartSentinelPayload, writeRestartSentinel, formatDoctorNonInteractiveHint, } from "../../infra/restart-sentinel.js"; -import { - normalizeGatewayRestartDelayMs, - resolveGatewayRestartDeferralTimeoutMs, - scheduleGatewayRestart, -} from "../../infra/restart.js"; +import { normalizeGatewayRestartDelayMs, scheduleGatewayRestart } from "../../infra/restart.js"; import { detectRespawnSupervisor } from "../../infra/supervisor-markers.js"; import { gatewayUpdateCampaign } from "../../infra/update-campaign.js"; import { diff --git a/src/gateway/server-reload-active-work.ts b/src/gateway/server-reload-active-work.ts index 75044c831237..3ba69903f7aa 100644 --- a/src/gateway/server-reload-active-work.ts +++ b/src/gateway/server-reload-active-work.ts @@ -1,7 +1,7 @@ import { getActiveBackgroundExecSessionCount } from "../agents/bash-process-registry.js"; import { getActiveEmbeddedRunCount } from "../agents/embedded-agent-runner/active-run-projections.js"; import { getTotalPendingReplies } from "../auto-reply/reply/dispatcher-registry.js"; -import { resolveGatewayRestartDeferralTimeoutMs } from "../infra/restart.js"; +import { resolveGatewayRestartDeferralTimeoutMs } from "../infra/restart-budget.js"; import { getTotalQueueSize } from "../process/command-queue.js"; import { getActiveGatewayRootWorkCount } from "../process/gateway-work-admission.js"; import { getInspectableActiveTaskRestartBlockers } from "../tasks/task-registry.maintenance.js"; diff --git a/src/gateway/server-reload-handlers.test.ts b/src/gateway/server-reload-handlers.test.ts index e84e07c16a1b..1ebc5216c3d3 100644 --- a/src/gateway/server-reload-handlers.test.ts +++ b/src/gateway/server-reload-handlers.test.ts @@ -4461,6 +4461,14 @@ describe("gateway restart deferral preflight", () => { title: "refresh all accounts", }), ); + const initialWarnings = [ + [ + "config change requires gateway restart (gateway.port) — deferring until 1 background task run(s) complete", + ], + [ + "restart blocked by active background task run(s): taskId=task-nightly runId=run-nightly status=running runtime=cron label=nightly sync title=refresh all accounts", + ], + ]; const signalSpy = vi.fn(); process.once("SIGUSR2", signalSpy); vi.useFakeTimers(); @@ -4470,16 +4478,7 @@ describe("gateway restart deferral preflight", () => { gateway: { reload: {} }, }); - expect(logReload.warn.mock.calls).toEqual( - expect.arrayContaining([ - [ - "config change requires gateway restart (gateway.port) — deferring until 1 background task run(s) complete", - ], - [ - "restart blocked by active background task run(s): taskId=task-nightly runId=run-nightly status=running runtime=cron label=nightly sync title=refresh all accounts", - ], - ]), - ); + expect(logReload.warn.mock.calls).toEqual(expect.arrayContaining(initialWarnings)); await vi.advanceTimersByTimeAsync(300_000); await Promise.resolve(); @@ -4487,17 +4486,13 @@ describe("gateway restart deferral preflight", () => { expect(signalSpy).toHaveBeenCalledTimes(1); expect(consumeGatewayRestartIntent()).toEqual({ force: true, + drainBudgetExhausted: true, reason: "config reload forced restart", }); expect(hoisted.markRestartAbortedMainSessions).not.toHaveBeenCalled(); expect(logReload.warn.mock.calls).toEqual( expect.arrayContaining([ - [ - "config change requires gateway restart (gateway.port) — deferring until 1 background task run(s) complete", - ], - [ - "restart blocked by active background task run(s): taskId=task-nightly runId=run-nightly status=running runtime=cron label=nightly sync title=refresh all accounts", - ], + ...initialWarnings, [ "restart timeout after 300000ms with 1 background task run(s) still active (taskId=task-nightly runId=run-nightly status=running runtime=cron label=nightly sync title=refresh all accounts); forcing restart", ], diff --git a/src/gateway/server-reload-restart.test.ts b/src/gateway/server-reload-restart.test.ts index a685c030ad62..0c868a798a78 100644 --- a/src/gateway/server-reload-restart.test.ts +++ b/src/gateway/server-reload-restart.test.ts @@ -157,7 +157,7 @@ describe("gateway restart readiness preflight", () => { expect(requestRecoveryRestart).toHaveBeenCalledExactlyOnceWith( "config reload: gateway.port", - { force: true, reason: "config reload forced restart" }, + { force: true, drainBudgetExhausted: true, reason: "config reload forced restart" }, ); expect(logReload.warn).toHaveBeenCalledWith( expect.stringContaining( diff --git a/src/gateway/server-reload-restart.ts b/src/gateway/server-reload-restart.ts index af15f2ca6059..3f60afdab568 100644 --- a/src/gateway/server-reload-restart.ts +++ b/src/gateway/server-reload-restart.ts @@ -3,11 +3,11 @@ import { isRestartEnabled } from "../config/commands.flags.js"; import { getConfigValueAtPath } from "../config/config-paths.js"; import { setRuntimeConfigAppliedHash } from "../config/config.js"; import type { OpenClawConfig } from "../config/types.openclaw.js"; +import { resolveGatewayRestartDeferralTimeoutMs } from "../infra/restart-budget.js"; import type { GatewayRestartIntent } from "../infra/restart-intent.js"; import { deferGatewayRestartUntilIdle, type RestartDeferralHandle, - resolveGatewayRestartDeferralTimeoutMs, setGatewayRestartPolicy, } from "../infra/restart.js"; import { runWithGatewayIndependentRootWorkAdmission } from "../process/gateway-work-admission.js"; diff --git a/src/gateway/server-startup-session-migration.ts b/src/gateway/server-startup-session-migration.ts index c5b68b3b6b7f..a7fdb246ca16 100644 --- a/src/gateway/server-startup-session-migration.ts +++ b/src/gateway/server-startup-session-migration.ts @@ -1,5 +1,7 @@ +import { buildAgentRunTerminalOutcome } from "../agents/agent-run-terminal-outcome.js"; import { hasSubagentSessionRecoveryOwner } from "../agents/subagents/registry/subagent-session-reconciliation.js"; -import { patchSessionEntryCore } from "../config/sessions/session-accessor.js"; +import { replaceSessionEntrySync } from "../config/sessions/session-accessor.js"; +import { readSessionEntryRow } from "../config/sessions/session-accessor.sqlite-entry-store.js"; import { hasSessionEntriesByStatus, readSessionEntriesByStatus, @@ -18,6 +20,8 @@ import { isIncognitoSessionKey, resolveAgentIdFromSessionKey, } from "../routing/session-key.js"; +import { isSessionWorkAdmissionActive } from "../sessions/session-lifecycle-admission.js"; +import { recordGatewaySessionRunFailure } from "../sessions/session-run-error.js"; import { withOpenClawAgentDatabaseReadOnly } from "../state/openclaw-agent-db-readonly.js"; import { openOpenClawAgentDatabase, @@ -105,40 +109,64 @@ async function reconcileStartupOrphans( continue; } const identity = { sessionKey, sessionId: entry.sessionId, env }; + const target = { + agentId: resolveAgentIdFromSessionKey(sessionKey), + env, + sessionKey, + storePath: connection.path, + }; + const matchesPredecessor = (current: InternalSessionEntry | undefined) => + current !== undefined && + current.sessionId === entry.sessionId && + current.lifecycleRevision === entry.lifecycleRevision && + current.lifecycleRunId === entry.lifecycleRunId && + current.updatedAt === entry.updatedAt && + current.startedAt === entry.startedAt && + isUnsettledPredecessor(current); const assertOwnerless = () => { assertGatewayOwner(); - if (hasSubagentSessionRecoveryOwner(identity)) { + if ( + hasSubagentSessionRecoveryOwner(identity) || + isSessionWorkAdmissionActive(connection.path, [sessionKey, entry.sessionId]) + ) { throw new Error("a current or retained run/task owns this session"); } }; try { assertOwnerless(); - const updated = await patchSessionEntryCore( - { - agentId: resolveAgentIdFromSessionKey(sessionKey), - env, - sessionKey, - storePath: connection.path, + const outcome = buildAgentRunTerminalOutcome({ + status: "error", + error: "subagent run was interrupted before a terminal lifecycle event was persisted", + startedAt: entry.startedAt, + // This is the repair observation, not a reconstructed execution finish time. + endedAt: Date.now(), + }); + await recordGatewaySessionRunFailure({ + target: { + ...target, + sessionId: entry.sessionId, + expectedLifecycleRevision: entry.lifecycleRevision, }, - (current) => - current.sessionId === entry.sessionId && - current.lifecycleRevision === entry.lifecycleRevision && - current.lifecycleRunId === entry.lifecycleRunId && - current.updatedAt === entry.updatedAt && - current.startedAt === entry.startedAt && - isUnsettledPredecessor(current) - ? { - status: "interrupted", - abortedLastRun: true, - lastRunError: - "subagent run was interrupted before a terminal lifecycle event was persisted", - } - : null, - { preserveActivity: true, skipMaintenance: true, assertCommitAllowed: assertOwnerless }, - ); - if (updated?.status === "interrupted") { - count++; - } + // A recovery-only receipt identity must not suppress the notice after partial output. + runId: `startup-orphan:${entry.sessionId}:${entry.lifecycleRunId ?? entry.startedAt}`, + error: outcome.error, + assertCommitAllowed: assertOwnerless, + settleSession: () => { + const current = readSessionEntryRow(connection, sessionKey)?.entry; + if (!current || !matchesPredecessor(current)) { + throw new Error("startup subagent session changed before interruption receipt"); + } + // The receipt owner holds the outer transaction; either both writes commit or neither does. + replaceSessionEntrySync(target, { + ...current, + status: "interrupted", + abortedLastRun: true, + endedAt: outcome.endedAt, + lastRunError: outcome.error, + }); + }, + }); + count++; } catch (error) { log.warn(`session: retained startup subagent ${sessionKey}: ${String(error)}`); } diff --git a/src/gateway/session-startup-orphan-races.test.ts b/src/gateway/session-startup-orphan-races.test.ts index 5788fe3392aa..1a2eb3e5a3a2 100644 --- a/src/gateway/session-startup-orphan-races.test.ts +++ b/src/gateway/session-startup-orphan-races.test.ts @@ -1,13 +1,18 @@ import fs from "node:fs"; import path from "node:path"; import { DatabaseSync } from "node:sqlite"; -import { afterEach, expect, it, vi } from "vitest"; +import { isRecord } from "@openclaw/normalization-core/record-coerce"; +import { afterEach, assert, expect, it, vi } from "vitest"; import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js"; import { resolveCompletionFromSessionEntry } from "../agents/subagents/registry/subagent-session-reconciliation.js"; import * as accessor from "../config/sessions/session-accessor.js"; +import * as entryStore from "../config/sessions/session-accessor.sqlite-entry-store.js"; import { readSessionEntriesByStatus } from "../config/sessions/session-accessor.sqlite-status.js"; +import * as transcriptStore from "../config/sessions/session-accessor.sqlite-transcript-store.js"; import { registerAgentRunContext, clearAgentRunContext } from "../infra/agent-run-registry.js"; import { acquireGatewayLock } from "../infra/gateway-lock.js"; +import { beginSessionWorkAdmission } from "../sessions/session-lifecycle-admission.js"; +import * as sessionRunError from "../sessions/session-run-error.js"; import { closeOpenClawAgentDatabasesForTest, openOpenClawAgentDatabase, @@ -32,134 +37,246 @@ it.each([ "session", "generation", "local-owner", + "session-admission", "durable-owner", "snapshot-owner", "snapshot-lease-release", "snapshot-lease-replace", -] as const)("retains the row when %s changes after startup prepares its patch", async (race) => { - const stateDir = fs.realpathSync.native(roots.make("startup-orphan-race-")); - await withEnvAsync( - { OPENCLAW_STATE_DIR: stateDir, OPENCLAW_CONFIG_PATH: path.join(stateDir, "openclaw.json") }, - async () => { - const lock = await acquireGatewayLock({ - allowInTests: true, - port: 24120, - listenerMode: "foreground", - }); - if (!lock) { - throw new Error("expected isolated Gateway ownership"); - } - try { - await lock.run(async () => { - const scope = { agentId: "main", sessionKey: "agent:main:subagent:race" }; - await accessor.replaceSessionEntry(scope, { - sessionId: "predecessor", - lifecycleRevision: "generation-1", - lifecycleRunId: "run-1", - status: "running", - startedAt: Math.floor(performance.timeOrigin) - 100, - updatedAt: Math.floor(performance.timeOrigin) - 100, - }); - const database = openOpenClawAgentDatabase({ agentId: "main" }); - const original = accessor.loadSessionEntryReadOnly(scope); - const patch = accessor.patchSessionEntryCore; - let prepared = false; - vi.spyOn(accessor, "patchSessionEntryCore").mockImplementation( - (target, update, options) => - patch( - target, - async (current, context) => { - const next = await update(current, context); - if (next?.status === "interrupted") { - prepared = true; - if (race === "session") { - accessor.replaceSessionEntrySync(scope, { - ...current, - sessionId: "successor", - }); - } else if (race === "generation") { - const other = new DatabaseSync(database.path); - try { - other - .prepare( - "UPDATE session_nodes SET entry_json = json_set(entry_json, ?, ?) WHERE session_key = ?", - ) - .run("$.lifecycleRevision", "successor", scope.sessionKey); - } finally { - other.close(); - } - } else if (race === "local-owner") { - registerAgentRunContext("startup-race-owner", { - sessionKey: scope.sessionKey, - sessionId: "predecessor", - projectSessionActive: false, - }); - } else if (race === "snapshot-lease-release") { - openOpenClawStateDatabase() - .db.prepare("DELETE FROM state_leases WHERE scope = ? AND lease_key = ?") - .run("gateway-owner", "global"); - } else if (race === "snapshot-lease-replace") { - openOpenClawStateDatabase() - .db.prepare( - "UPDATE state_leases SET owner = ? WHERE scope = ? AND lease_key = ?", - ) - .run("replacement-owner", "gateway-owner", "global"); - } else { - openOpenClawStateDatabase() - .db.prepare( - "INSERT INTO subagent_runs(run_id,child_session_key,requester_session_key,created_at,payload_json) VALUES(?,?,?,?,?)", - ) - .run( - "startup-race-owner", - scope.sessionKey, - "agent:main:main", - Date.now(), - "{}", - ); - } - } - return next; - }, - options, - ), - ); - const log = { info: vi.fn(), warn: vi.fn() }; - const runStartup = () => - runStartupSessionMigration({ cfg: { agents: { entries: { main: {} } } }, log }); - if ( - race === "snapshot-owner" || - race === "snapshot-lease-release" || - race === "snapshot-lease-replace" - ) { - openOpenClawStateDatabase(); - await withOpenClawStateDatabaseReadSnapshot(runStartup); - } else { - await runStartup(); - } - expect( - prepared, - JSON.stringify({ - warnings: log.warn.mock.calls, - info: log.info.mock.calls, - current: accessor.loadSessionEntryReadOnly(scope), - }), - ).toBe(true); - expect(log.warn).toHaveBeenCalled(); - const current = accessor.loadSessionEntryReadOnly(scope); - expect(current).toEqual({ - ...original, - ...(race === "session" ? { sessionId: "successor" } : {}), - ...(race === "generation" ? { lifecycleRevision: "successor" } : {}), - }); +] as const)( + "retains the row and emits no receipt when %s changes after repair preparation", + async (race) => { + const stateDir = fs.realpathSync.native(roots.make("startup-orphan-race-")); + await withEnvAsync( + { OPENCLAW_STATE_DIR: stateDir, OPENCLAW_CONFIG_PATH: path.join(stateDir, "openclaw.json") }, + async () => { + const lock = await acquireGatewayLock({ + allowInTests: true, + port: 24120, + listenerMode: "foreground", }); - } finally { - closeOpenClawAgentDatabasesForTest(); - await lock.release(); - closeOpenClawStateDatabaseForTest(); - } - }, - ); -}); + if (!lock) { + throw new Error("expected isolated Gateway ownership"); + } + let admission: Awaited> | undefined; + try { + await lock.run(async () => { + const scope = { agentId: "main", sessionKey: "agent:main:subagent:race" }; + await accessor.replaceSessionEntry(scope, { + sessionId: "predecessor", + lifecycleRevision: "generation-1", + lifecycleRunId: "run-1", + status: "running", + startedAt: Math.floor(performance.timeOrigin) - 100, + updatedAt: Math.floor(performance.timeOrigin) - 100, + }); + const database = openOpenClawAgentDatabase({ agentId: "main" }); + const original = accessor.loadSessionEntryReadOnly(scope); + assert(original); + const writeReceipt = sessionRunError.recordGatewaySessionRunFailure; + let prepared = false; + vi.spyOn(sessionRunError, "recordGatewaySessionRunFailure").mockImplementationOnce( + async (params) => { + await Promise.resolve(); + prepared = true; + if (race === "session") { + accessor.replaceSessionEntrySync(scope, { + ...original, + sessionId: "successor", + }); + } else if (race === "generation") { + const other = new DatabaseSync(database.path); + try { + other + .prepare( + "UPDATE session_nodes SET entry_json = json_set(entry_json, ?, ?) WHERE session_key = ?", + ) + .run("$.lifecycleRevision", "successor", scope.sessionKey); + } finally { + other.close(); + } + } else if (race === "local-owner") { + registerAgentRunContext("startup-race-owner", { + sessionKey: scope.sessionKey, + sessionId: "predecessor", + projectSessionActive: false, + }); + } else if (race === "session-admission") { + admission = await beginSessionWorkAdmission({ + scope: database.path, + identities: [scope.sessionKey, "predecessor"], + assertAllowed: () => {}, + }); + } else if (race === "snapshot-lease-release") { + openOpenClawStateDatabase() + .db.prepare("DELETE FROM state_leases WHERE scope = ? AND lease_key = ?") + .run("gateway-owner", "global"); + } else if (race === "snapshot-lease-replace") { + openOpenClawStateDatabase() + .db.prepare( + "UPDATE state_leases SET owner = ? WHERE scope = ? AND lease_key = ?", + ) + .run("replacement-owner", "gateway-owner", "global"); + } else { + openOpenClawStateDatabase() + .db.prepare( + "INSERT INTO subagent_runs(run_id,child_session_key,requester_session_key,created_at,payload_json) VALUES(?,?,?,?,?)", + ) + .run( + "startup-race-owner", + scope.sessionKey, + "agent:main:main", + Date.now(), + "{}", + ); + } + await writeReceipt(params); + }, + ); + const log = { info: vi.fn(), warn: vi.fn() }; + const runStartup = () => + runStartupSessionMigration({ cfg: { agents: { entries: { main: {} } } }, log }); + if ( + race === "snapshot-owner" || + race === "snapshot-lease-release" || + race === "snapshot-lease-replace" + ) { + openOpenClawStateDatabase(); + await withOpenClawStateDatabaseReadSnapshot(runStartup); + } else { + await runStartup(); + } + expect( + prepared, + JSON.stringify({ + warnings: log.warn.mock.calls, + info: log.info.mock.calls, + current: accessor.loadSessionEntryReadOnly(scope), + }), + ).toBe(true); + expect(log.warn).toHaveBeenCalled(); + expect( + (await accessor.loadTranscriptEvents({ ...scope, sessionId: "predecessor" })).filter( + (event) => isRecord(event) && event.customType === "run-failed-before-reply", + ), + ).toEqual([]); + const current = accessor.loadSessionEntryReadOnly(scope); + expect(current).toEqual({ + ...original, + ...(race === "session" ? { sessionId: "successor" } : {}), + ...(race === "generation" ? { lifecycleRevision: "successor" } : {}), + }); + }); + } finally { + admission?.release(); + closeOpenClawAgentDatabasesForTest(); + await lock.release(); + closeOpenClawStateDatabaseForTest(); + } + }, + ); + }, +); + +it.each(["owner", "settlement", "receipt"] as const)( + "settles the orphan and receipt atomically after a %s failure", + async (failure) => { + const stateDir = fs.realpathSync.native(roots.make("startup-orphan-receipt-")); + await withEnvAsync( + { OPENCLAW_STATE_DIR: stateDir, OPENCLAW_CONFIG_PATH: path.join(stateDir, "openclaw.json") }, + async () => { + const lock = await acquireGatewayLock({ + allowInTests: true, + port: 24120, + listenerMode: "foreground", + }); + if (!lock) { + throw new Error("expected isolated Gateway ownership"); + } + try { + await lock.run(async () => { + const target = { + agentId: "main", + sessionKey: "agent:main:subagent:receipt", + sessionId: "predecessor-receipt", + }; + await accessor.replaceSessionEntry(target, { + sessionId: target.sessionId, + lifecycleRevision: "predecessor-generation", + status: "running", + startedAt: Math.floor(performance.timeOrigin) - 100, + updatedAt: Math.floor(performance.timeOrigin) - 100, + }); + const original = accessor.loadSessionEntryReadOnly(target); + const log = { info: vi.fn(), warn: vi.fn() }; + const runStartup = () => + runStartupSessionMigration({ cfg: { agents: { entries: { main: {} } } }, log }); + if (failure === "receipt") { + vi.spyOn( + transcriptStore, + "appendTranscriptEventInTransaction", + ).mockImplementationOnce(() => { + throw new Error("synthetic repair receipt write failure"); + }); + } else { + const write = entryStore.writeSessionEntry; + vi.spyOn(entryStore, "writeSessionEntry").mockImplementationOnce((...args) => { + if (failure === "settlement") { + throw new Error("synthetic settlement failure"); + } + registerAgentRunContext("startup-race-owner", { + sessionKey: target.sessionKey, + sessionId: target.sessionId, + projectSessionActive: false, + }); + return write(...args); + }); + } + await runStartup(); + expect(accessor.loadSessionEntryReadOnly(target)).toEqual(original); + expect(log.warn).toHaveBeenCalled(); + expect( + (await accessor.loadTranscriptEvents(target)).filter( + (event) => isRecord(event) && event.customType === "run-failed-before-reply", + ), + ).toEqual([]); + clearAgentRunContext("startup-race-owner"); + + const repairObservedAt = Date.now(); + await runStartup(); + const repaired = accessor.loadSessionEntryReadOnly(target); + expect(repaired).toMatchObject({ + status: "interrupted", + abortedLastRun: true, + startedAt: original?.startedAt, + updatedAt: original?.updatedAt, + }); + expect(repaired?.endedAt).toBeGreaterThanOrEqual(repairObservedAt); + expect(repaired?.runtimeMs).toBeUndefined(); + const receipts = async () => + (await accessor.loadTranscriptEvents(target)).filter( + (event) => isRecord(event) && event.customType === "run-failed-before-reply", + ); + expect(await receipts()).toMatchObject([ + { + display: true, + details: { + error: expect.stringContaining("interrupted before a terminal lifecycle event"), + }, + }, + ]); + await runStartup(); + expect(accessor.loadSessionEntryReadOnly(target)).toEqual(repaired); + expect(await receipts()).toHaveLength(1); + }); + } finally { + closeOpenClawAgentDatabasesForTest(); + await lock.release(); + closeOpenClawStateDatabaseForTest(); + } + }, + ); + }, +); it("keeps interrupted status distinct in canonical reads without inventing registry completion", async () => { const env = { diff --git a/src/gateway/session-startup-orphans.test.ts b/src/gateway/session-startup-orphans.test.ts index 267eeaa2b892..959c125985ec 100644 --- a/src/gateway/session-startup-orphans.test.ts +++ b/src/gateway/session-startup-orphans.test.ts @@ -1,35 +1,41 @@ import { execFile } from "node:child_process"; import fs from "node:fs"; import path from "node:path"; -import { fileURLToPath } from "node:url"; import { promisify } from "node:util"; -import { afterEach, expect, it } from "vitest"; +import { afterAll, beforeAll, expect, it } from "vitest"; import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js"; +import { gatewayDirectStopEntrypoints } from "../cli/cli-entrypoint.test-support.js"; +import { resolveRuntimeWorkerArgv, resolveRuntimeWorkerUrl } from "../infra/runtime-worker-url.js"; + +const tempDirs = useAutoCleanupTempDirTracker(afterAll); +const fixture = resolveRuntimeWorkerUrl(gatewayDirectStopEntrypoints.startupOrphanFixture); +let stateRoot: string; + +beforeAll(async () => { + stateRoot = fs.realpathSync.native(tempDirs.make("openclaw-startup-orphan-")); + const env = { + ...process.env, + OPENCLAW_STATE_DIR: stateRoot, + OPENCLAW_CONFIG_PATH: path.join(stateRoot, "openclaw.json"), + }; + for (const mode of ["predecessor", "successor"]) { + await promisify(execFile)(process.execPath, [...resolveRuntimeWorkerArgv(fixture), mode], { + env, + timeout: 60000, + }); + } +}, 120000); -const tempDirs = useAutoCleanupTempDirTracker(afterEach); -const fixturePath = fileURLToPath( - new URL("./startup-orphan-process.test-support.ts", import.meta.url), -); it.each(["default", "shared", "embedded"])( "recovers only ownerless predecessor rows in %s state", - async (layout) => { - const stateDir = fs.realpathSync.native(tempDirs.make("openclaw-startup-orphan-")); - const env = { - ...process.env, - OPENCLAW_STATE_DIR: stateDir, - OPENCLAW_CONFIG_PATH: path.join(stateDir, "openclaw.json"), - }; - const run = async (mode: string) => { - await promisify(execFile)(process.execPath, ["--import", "tsx", fixturePath, mode, layout], { - env, - timeout: 60000, - }); - return JSON.parse(fs.readFileSync(path.join(stateDir, mode + ".json"), "utf8")); - }; - const predecessor = await run("predecessor"); + (layout) => { + const stateDir = path.join(stateRoot, layout); + const read = (mode: string) => + JSON.parse(fs.readFileSync(path.join(stateDir, mode + ".json"), "utf8")); + const predecessor = read("predecessor"); expect(() => process.kill(predecessor.pid, 0)).toThrow(); - const result = await run(layout === "embedded" ? "embedded" : "successor"); - const before = JSON.parse(fs.readFileSync(path.join(stateDir, "before-startup.json"), "utf8")); + const result = read(layout === "embedded" ? "embedded" : "successor"); + const before = read("before-startup"); expect(result.pid).not.toBe(predecessor.pid); const { running: mainOrphan, "ops-running": opsOrphan, ...controls } = result.rows; const { running: mainOriginal, "ops-running": opsOriginal, ...originalControls } = before.rows; @@ -47,7 +53,7 @@ it.each(["default", "shared", "embedded"])( expect(orphan.lifecycleRevision).toBe(original.lifecycleRevision); expect(orphan.updatedAt).toBe(original.updatedAt); expect(orphan.startedAt).toBe(original.startedAt); - expect(orphan.endedAt).toBeUndefined(); + expect(orphan.endedAt).toBeGreaterThan(original.startedAt); expect(orphan.runtimeMs).toBeUndefined(); } } diff --git a/src/gateway/startup-orphan-process.test-support.ts b/src/gateway/startup-orphan-process.test-support.ts index b305ec37c8f4..176079b4ad7e 100644 --- a/src/gateway/startup-orphan-process.test-support.ts +++ b/src/gateway/startup-orphan-process.test-support.ts @@ -9,7 +9,7 @@ import { } from "../config/sessions/session-accessor.js"; import type { InternalSessionEntry } from "../config/sessions/types.js"; import type { OpenClawConfig } from "../config/types.openclaw.js"; -import { registerAgentRunContext } from "../infra/agent-run-registry.js"; +import { clearAgentRunContext, registerAgentRunContext } from "../infra/agent-run-registry.js"; import { acquireGatewayLock } from "../infra/gateway-lock.js"; import { closeOpenClawAgentDatabasesForTest } from "../state/openclaw-agent-db.js"; import { @@ -19,167 +19,187 @@ import { import { upsertTaskWithDeliveryStateToSqlite } from "../tasks/task-registry.store.sqlite.js"; import { runStartupSessionMigration } from "./server-startup-session-migration.js"; -const stateDir = process.env.OPENCLAW_STATE_DIR!; -const mode = process.argv[2]; -const layout = process.argv[3]; -const storePath = layout === "shared" ? path.join(stateDir, "shared.sqlite") : undefined; -const cfg: OpenClawConfig = { - agents: { entries: { main: {}, ops: {} } }, - ...(storePath ? { session: { store: storePath } } : {}), -}; -const kinds = [ - "running", - "done", - "live", - "yielded", - "queued", - "recovering", - "retained-task", - "registry-queued", - "registry-recovering", - "registry-completion", - "malformed-owner", - "rebound", - "ops-running", - "incognito-control", -] as const; -const key = (kind: string) => - "agent:" + (kind === "ops-running" ? "ops" : "main") + ":subagent:" + kind; -const scope = (kind: string) => ({ - agentId: kind === "ops-running" ? "ops" : "main", - sessionKey: key(kind), - storePath, -}); -const rows = () => - Object.fromEntries(kinds.map((kind) => [kind, loadSessionEntryReadOnly(scope(kind))])); -const durableOwners = () => ({ - runs: openOpenClawStateDatabase().db.prepare("SELECT * FROM subagent_runs ORDER BY run_id").all(), - tasks: openOpenClawStateDatabase().db.prepare("SELECT * FROM task_runs ORDER BY task_id").all(), -}); -const lock = await acquireGatewayLock({ - allowInTests: true, - port: 24119, - ...(mode === "embedded" - ? { role: "agent-embedded" as const } - : { listenerMode: "foreground" as const }), -}); -if (!lock) { - throw new Error("proof requires actual process ownership"); +const stateRoot = process.env.OPENCLAW_STATE_DIR!; +const generation = process.argv[2]; +if (generation !== "predecessor" && generation !== "successor") { + throw new Error("unexpected fixture generation"); } -try { - await lock.run(async () => { - if (mode === "predecessor") { - for (const kind of kinds) { - if (kind === "incognito-control") { - continue; - } - const now = Date.now(); - const entry: InternalSessionEntry = { - sessionId: "predecessor-" + kind, - lifecycleRevision: "predecessor-" + kind, - startedAt: now, - updatedAt: now, - status: kind === "done" ? "done" : kind === "queued" ? "queued" : "running", - ...(kind === "done" ? { endedAt: now, runtimeMs: 0 } : {}), - ...(kind === "recovering" - ? { abortedLastRun: true, restartRecoveryForceSafeTools: true } - : {}), - }; - await upsertSessionEntryCore(scope(kind), entry); - } - const registry = new Map(); - for (const kind of ["registry-queued", "registry-recovering", "registry-completion"]) { - registry.set(kind, { - runId: kind, - childSessionKey: key(kind), - requesterSessionKey: "agent:main:main", - requesterDisplayKey: "main", - task: "retained control", - cleanup: "keep", - createdAt: Date.now(), - generation: 7, - execution: { - status: - kind === "registry-queued" - ? "queued" - : kind === "registry-recovering" - ? "interrupted" - : "terminal", - }, - completion: { required: true }, - delivery: { status: "pending" }, - ...(kind === "registry-recovering" - ? { terminalOwner: "interrupted-recovery" as const } - : {}), - }); - } - saveSubagentRegistryToSqlite(registry); - // A malformed retained claim is unresolved ownership, never permission to settle its session. - openOpenClawStateDatabase() - .db.prepare( - "INSERT INTO subagent_runs(run_id,child_session_key,requester_session_key,created_at,payload_json) VALUES(?,?,?,?,?)", - ) - .run("malformed-owner", key("malformed-owner"), "agent:main:main", Date.now(), "{}"); - upsertTaskWithDeliveryStateToSqlite({ - task: { - taskId: "retained-task", - runtime: "subagent", - requesterSessionKey: "agent:main:main", - ownerKey: "agent:main:main", - childSessionKey: key("retained-task"), - scopeKind: "session", - task: "retained completion", - status: "succeeded", - deliveryStatus: "pending", - notifyPolicy: "done_only", - createdAt: Date.now(), - }, - }); - } else if (mode === "successor" || mode === "embedded") { - await replaceSessionEntry(scope("incognito-control"), { - sessionId: "incognito", - incognito: true, - status: "running", - startedAt: Math.floor(performance.timeOrigin) - 100, - updatedAt: Math.floor(performance.timeOrigin) - 100, - }); - registerAgentRunContext("live-owner", { - sessionKey: key("live"), - sessionId: "predecessor-live", - projectSessionActive: true, - }); - registerAgentRunContext("yielded-owner", { - sessionKey: key("yielded"), - sessionId: "predecessor-yielded", - projectSessionActive: false, - }); - const rebound = loadSessionEntryReadOnly(scope("rebound"))!; - await upsertSessionEntryCore(scope("rebound"), { - ...rebound, - sessionId: "successor-rebound", - lifecycleRevision: "successor-rebound", - startedAt: Date.now(), - updatedAt: Date.now(), - }); - fs.writeFileSync( - path.join(stateDir, "before-startup.json"), - JSON.stringify({ rows: rows(), owners: durableOwners() }), - ); - await runStartupSessionMigration({ - cfg, - env: process.env, - log: { info: console.error, warn: console.error }, - }); - } else { - throw new Error("unexpected fixture mode"); - } - fs.writeFileSync( - path.join(stateDir, mode + ".json"), - JSON.stringify({ pid: process.pid, rows: rows(), owners: durableOwners() }), - ); +for (const layout of ["default", "shared", "embedded"]) { + const stateDir = path.join(stateRoot, layout); + fs.mkdirSync(stateDir, { recursive: true }); + process.env.OPENCLAW_STATE_DIR = stateDir; + process.env.OPENCLAW_CONFIG_PATH = path.join(stateDir, "openclaw.json"); + await runLayout( + stateDir, + layout, + generation === "predecessor" ? generation : layout === "embedded" ? "embedded" : generation, + ); +} + +async function runLayout(stateDir: string, layout: string, mode: string) { + const storePath = layout === "shared" ? path.join(stateDir, "shared.sqlite") : undefined; + const cfg: OpenClawConfig = { + agents: { entries: { main: {}, ops: {} } }, + ...(storePath ? { session: { store: storePath } } : {}), + }; + const kinds = [ + "running", + "done", + "live", + "yielded", + "queued", + "recovering", + "retained-task", + "registry-queued", + "registry-recovering", + "registry-completion", + "malformed-owner", + "rebound", + "ops-running", + "incognito-control", + ] as const; + const key = (kind: string) => + "agent:" + (kind === "ops-running" ? "ops" : "main") + ":subagent:" + kind; + const scope = (kind: string) => ({ + agentId: kind === "ops-running" ? "ops" : "main", + sessionKey: key(kind), + storePath, }); -} finally { - closeOpenClawAgentDatabasesForTest(); - await lock.release(); - closeOpenClawStateDatabaseForTest(); + const rows = () => + Object.fromEntries(kinds.map((kind) => [kind, loadSessionEntryReadOnly(scope(kind))])); + const durableOwners = () => ({ + runs: openOpenClawStateDatabase() + .db.prepare("SELECT * FROM subagent_runs ORDER BY run_id") + .all(), + tasks: openOpenClawStateDatabase().db.prepare("SELECT * FROM task_runs ORDER BY task_id").all(), + }); + const lock = await acquireGatewayLock({ + allowInTests: true, + port: 24119, + ...(mode === "embedded" + ? { role: "agent-embedded" as const } + : { listenerMode: "foreground" as const }), + }); + if (!lock) { + throw new Error("proof requires actual process ownership"); + } + try { + await lock.run(async () => { + if (mode === "predecessor") { + for (const kind of kinds) { + if (kind === "incognito-control") { + continue; + } + const now = Date.now(); + const entry: InternalSessionEntry = { + sessionId: "predecessor-" + kind, + lifecycleRevision: "predecessor-" + kind, + startedAt: now, + updatedAt: now, + status: kind === "done" ? "done" : kind === "queued" ? "queued" : "running", + ...(kind === "done" ? { endedAt: now, runtimeMs: 0 } : {}), + ...(kind === "recovering" + ? { abortedLastRun: true, restartRecoveryForceSafeTools: true } + : {}), + }; + await upsertSessionEntryCore(scope(kind), entry); + } + const registry = new Map(); + for (const kind of ["registry-queued", "registry-recovering", "registry-completion"]) { + registry.set(kind, { + runId: kind, + childSessionKey: key(kind), + requesterSessionKey: "agent:main:main", + requesterDisplayKey: "main", + task: "retained control", + cleanup: "keep", + createdAt: Date.now(), + generation: 7, + execution: { + status: + kind === "registry-queued" + ? "queued" + : kind === "registry-recovering" + ? "interrupted" + : "terminal", + }, + completion: { required: true }, + delivery: { status: "pending" }, + ...(kind === "registry-recovering" + ? { terminalOwner: "interrupted-recovery" as const } + : {}), + }); + } + saveSubagentRegistryToSqlite(registry); + // A malformed retained claim is unresolved ownership, never permission to settle its session. + openOpenClawStateDatabase() + .db.prepare( + "INSERT INTO subagent_runs(run_id,child_session_key,requester_session_key,created_at,payload_json) VALUES(?,?,?,?,?)", + ) + .run("malformed-owner", key("malformed-owner"), "agent:main:main", Date.now(), "{}"); + upsertTaskWithDeliveryStateToSqlite({ + task: { + taskId: "retained-task", + runtime: "subagent", + requesterSessionKey: "agent:main:main", + ownerKey: "agent:main:main", + childSessionKey: key("retained-task"), + scopeKind: "session", + task: "retained completion", + status: "succeeded", + deliveryStatus: "pending", + notifyPolicy: "done_only", + createdAt: Date.now(), + }, + }); + } else if (mode === "successor" || mode === "embedded") { + await replaceSessionEntry(scope("incognito-control"), { + sessionId: "incognito", + incognito: true, + status: "running", + startedAt: Math.floor(performance.timeOrigin) - 100, + updatedAt: Math.floor(performance.timeOrigin) - 100, + }); + registerAgentRunContext("live-owner", { + sessionKey: key("live"), + sessionId: "predecessor-live", + projectSessionActive: true, + }); + registerAgentRunContext("yielded-owner", { + sessionKey: key("yielded"), + sessionId: "predecessor-yielded", + projectSessionActive: false, + }); + const rebound = loadSessionEntryReadOnly(scope("rebound"))!; + await upsertSessionEntryCore(scope("rebound"), { + ...rebound, + sessionId: "successor-rebound", + lifecycleRevision: "successor-rebound", + startedAt: Date.now(), + updatedAt: Date.now(), + }); + fs.writeFileSync( + path.join(stateDir, "before-startup.json"), + JSON.stringify({ rows: rows(), owners: durableOwners() }), + ); + await runStartupSessionMigration({ + cfg, + env: process.env, + log: { info: console.error, warn: console.error }, + }); + } else { + throw new Error("unexpected fixture mode"); + } + fs.writeFileSync( + path.join(stateDir, mode + ".json"), + JSON.stringify({ pid: process.pid, rows: rows(), owners: durableOwners() }), + ); + }); + } finally { + clearAgentRunContext("live-owner"); + clearAgentRunContext("yielded-owner"); + closeOpenClawAgentDatabasesForTest(); + await lock.release(); + closeOpenClawStateDatabaseForTest(); + } } diff --git a/src/infra/gateway-shutdown-budget.ts b/src/infra/gateway-shutdown-budget.ts index 0cde0821380a..c07f05d3aadd 100644 --- a/src/infra/gateway-shutdown-budget.ts +++ b/src/infra/gateway-shutdown-budget.ts @@ -1,4 +1,5 @@ export { + GATEWAY_RESTART_REPLACEMENT_TIMEOUT_MS, GATEWAY_SHUTDOWN_RESERVE_MS, GATEWAY_SUPERVISOR_EXIT_MARGIN_MS, GATEWAY_SHUTDOWN_TIMEOUT_MS, diff --git a/src/infra/restart-budget.ts b/src/infra/restart-budget.ts new file mode 100644 index 000000000000..6d983a624898 --- /dev/null +++ b/src/infra/restart-budget.ts @@ -0,0 +1,26 @@ +import { + GATEWAY_RESTART_REPLACEMENT_TIMEOUT_MS, + GATEWAY_SHUTDOWN_RESERVE_MS, + GATEWAY_SUPERVISOR_EXIT_MARGIN_MS, +} from "./gateway-shutdown-budget.js"; +import type { GatewayRestartIntent } from "./restart-intent.js"; + +const DEFAULT_RESTART_DEFERRAL_TIMEOUT_MS = 300_000; + +export function resolveGatewayRestartDeferralTimeoutMs(): number; +export function resolveGatewayRestartDeferralTimeoutMs(timeoutMs: unknown): number | undefined; +export function resolveGatewayRestartDeferralTimeoutMs(timeoutMs?: unknown): number | undefined { + if (typeof timeoutMs !== "number" || !Number.isFinite(timeoutMs)) { + return DEFAULT_RESTART_DEFERRAL_TIMEOUT_MS; + } + return timeoutMs > 0 ? Math.floor(timeoutMs) : undefined; +} + +export function resolveGatewayRestartDrainTimeoutMs(intent?: GatewayRestartIntent) { + const reserveMs = GATEWAY_SHUTDOWN_RESERVE_MS + GATEWAY_SUPERVISOR_EXIT_MARGIN_MS; + const waitMs = + intent?.waitMs ?? + (intent?.force ? GATEWAY_RESTART_REPLACEMENT_TIMEOUT_MS - reserveMs : undefined); + const exhausted = intent?.drainBudgetExhausted || (intent?.force && waitMs === 0); + return exhausted ? 0 : resolveGatewayRestartDeferralTimeoutMs(waitMs); +} diff --git a/src/infra/restart-intent.ts b/src/infra/restart-intent.ts index 3f5782683298..685a41e5f5ff 100644 --- a/src/infra/restart-intent.ts +++ b/src/infra/restart-intent.ts @@ -56,6 +56,8 @@ export type GatewayRestartIntent = { reason?: string; force?: boolean; waitMs?: number; + // Only the in-process deferral owner can attest that the drain budget was spent. + drainBudgetExhausted?: true; // Process-local only: persisted restart requests cannot delegate successor ownership. successorOwner?: { kind: "managed-update-handoff"; diff --git a/src/infra/restart.deferral-timeout.test.ts b/src/infra/restart.deferral-timeout.test.ts index 21409c95e96e..d935dbf87fc6 100644 --- a/src/infra/restart.deferral-timeout.test.ts +++ b/src/infra/restart.deferral-timeout.test.ts @@ -113,6 +113,7 @@ describe("deferGatewayRestartUntilIdle timeout", () => { expect(hooks.onTimeout).toHaveBeenCalledOnce(); expect(consumeGatewayRestartIntent()).toEqual({ force: true, + drainBudgetExhausted: true, reason: "gateway.restart.deferral-timeout", }); }); @@ -366,6 +367,7 @@ describe("deferGatewayRestartUntilIdle timeout", () => { expect(hooks.onTimeout).toHaveBeenCalledOnce(); expect(consumeGatewayRestartIntent()).toEqual({ force: true, + drainBudgetExhausted: true, reason: "gateway.restart.deferral-timeout", }); }); @@ -549,6 +551,6 @@ describe("deferGatewayRestartUntilIdle timeout", () => { expect(emit).not.toHaveBeenCalledWith("SIGUSR2"); await vi.advanceTimersByTimeAsync(300_000); expect(emit.mock.calls.filter(([event]) => event === "SIGUSR2")).toHaveLength(1); - expect(consumeGatewayRestartIntent()).toEqual({ force: true }); + expect(consumeGatewayRestartIntent()).toEqual({ force: true, drainBudgetExhausted: true }); }); }); diff --git a/src/infra/restart.ts b/src/infra/restart.ts index cacec7708388..649ffebe6a51 100644 --- a/src/infra/restart.ts +++ b/src/infra/restart.ts @@ -11,6 +11,7 @@ import { type GatewayRestartSignalAdmissionLease, } from "../process/gateway-work-admission.js"; import { formatErrorMessage } from "./errors.js"; +import { resolveGatewayRestartDeferralTimeoutMs } from "./restart-budget.js"; import { type GatewayRestartIntent, normalizeRestartIntentReason } from "./restart-intent.js"; import { restartGatewayViaSupervisor } from "./restart-supervisor.js"; import type { RestartAttempt } from "./restart.types.js"; @@ -20,7 +21,6 @@ export { normalizeSystemdUnit } from "./restart-supervisor.js"; const RESTART_AUTH_GRACE_MS = 5000; const DEFAULT_DEFERRAL_POLL_MS = 500; const DEFAULT_DEFERRAL_STILL_PENDING_WARN_MS = 30_000; -const DEFAULT_RESTART_DEFERRAL_TIMEOUT_MS = 300_000; const RESTART_COOLDOWN_MS = 30_000; const restartLog = createSubsystemLogger("restart"); @@ -384,15 +384,6 @@ type GatewayRestartEmitResult = | { status: "coalesced" } | { status: "failed" }; -export function resolveGatewayRestartDeferralTimeoutMs(): number; -export function resolveGatewayRestartDeferralTimeoutMs(timeoutMs: unknown): number | undefined; -export function resolveGatewayRestartDeferralTimeoutMs(timeoutMs?: unknown): number | undefined { - if (typeof timeoutMs !== "number" || !Number.isFinite(timeoutMs)) { - return DEFAULT_RESTART_DEFERRAL_TIMEOUT_MS; - } - return timeoutMs > 0 ? Math.floor(timeoutMs) : undefined; -} - function canReplacePendingRestartEmitHooks( hooks: RestartEmitHooks | undefined, sessionKey: string | undefined, @@ -753,7 +744,7 @@ export function deferGatewayRestartUntilIdle(opts: { void emitPreparedGatewayRestart( opts.emitHooks, opts.reason, - timedOut ? opts.timeoutIntent : undefined, + timedOut ? { ...opts.timeoutIntent, drainBudgetExhausted: true } : undefined, timedOut ? undefined : () => { diff --git a/src/infra/update-startup-auto-run.ts b/src/infra/update-startup-auto-run.ts index cce3e4c7bc96..3ff8ea33a3b5 100644 --- a/src/infra/update-startup-auto-run.ts +++ b/src/infra/update-startup-auto-run.ts @@ -6,11 +6,11 @@ import { EXTERNAL_SUPERVISOR_UPDATE_REQUIRED_REASON, isGatewayExternallySupervised, } from "./gateway-supervision.js"; +import { resolveGatewayRestartDeferralTimeoutMs } from "./restart-budget.js"; import { readRestartSentinelSnapshot, writeRestartSentinelIfUnchanged, } from "./restart-sentinel.js"; -import { resolveGatewayRestartDeferralTimeoutMs } from "./restart.js"; import { detectRespawnSupervisor } from "./supervisor-markers.js"; import type { UpdateCampaignController } from "./update-campaign.js"; import { isPendingControlPlaneUpdateRestartSentinel } from "./update-control-plane-sentinel.js"; diff --git a/src/sessions/session-run-error.ts b/src/sessions/session-run-error.ts index 7da80e098a17..ac8de938d1f9 100644 --- a/src/sessions/session-run-error.ts +++ b/src/sessions/session-run-error.ts @@ -16,12 +16,13 @@ function sanitizeSessionRunError(error: unknown): string { return redactSensitiveText(text, { mode: "tools" }); } -/** Shared transcript outcome for owners that already committed a failed run. */ +/** Shared failure receipt; optional settlement joins the receipt's synchronous transaction. */ export async function recordGatewaySessionRunFailure(params: { target: SessionTranscriptWriteScope & { sessionId: string }; runId: string; error: unknown; assertCommitAllowed?: () => void; + settleSession?: () => undefined; }): Promise { const { runId } = params; const error = truncateUtf16Safe(sanitizeSessionRunError(params.error), 512) || "unknown error"; @@ -30,6 +31,8 @@ export async function recordGatewaySessionRunFailure(params: { customTypes: [RUN_FAILED_BEFORE_REPLY_TRANSCRIPT_TYPE], suppressWhenAssistantRun: runId, selectReport: (latest) => { + params.assertCommitAllowed?.(); + params.settleSession?.(); params.assertCommitAllowed?.(); if (isRecord(latest?.details) && latest.details.runId === runId) { return undefined;