diff --git a/src/agents/subagent-orphan-recovery.restart-integration.test.ts b/src/agents/subagent-orphan-recovery.restart-integration.test.ts index 82945bdba84b..74e9ad1bb233 100644 --- a/src/agents/subagent-orphan-recovery.restart-integration.test.ts +++ b/src/agents/subagent-orphan-recovery.restart-integration.test.ts @@ -38,6 +38,7 @@ import { import { loadSubagentRegistryFromSqlite } from "./subagents/registry/subagent-registry.store.sqlite.js"; import { addSubagentRunForTests, + activateSubagentRegistry, getSubagentRunByChildSessionKey, initSubagentRegistry, listSubagentRunsForRequester, @@ -74,9 +75,13 @@ async function acceptRecoveryDispatch(payload: Record) { const dispatchAgent = vi.fn(acceptRecoveryDispatch); const gatewayRuntime: GatewayRecoveryRuntime = { dispatchAgent: dispatchAgent as GatewayRecoveryRuntime["dispatchAgent"], - waitForAgent: vi.fn(), + waitForAgent: vi.fn(async () => ({ + status: "pending", + })) as GatewayRecoveryRuntime["waitForAgent"], sendRecoveryNotice: vi.fn(), }; +const activateGatewayRuntime = () => + activateSubagentRegistry(() => ({ recoveryRuntime: gatewayRuntime }) as never); vi.mock("../gateway/session-utils.fs.js", () => ({ readSessionMessagesAsync: vi.fn(async () => []), @@ -114,10 +119,10 @@ describe("subagent orphan recovery — faithful restart path", () => { // external side effects) are recorded so completeSubagentRun runs in-process. testing.setDepsForTest({ ...createSubagentRegistryTestDeps(), - getGatewayRecoveryRuntime: () => gatewayRuntime, runSubagentAnnounceFlow: vi.fn(async () => "delivered" as const), onAgentEvent: vi.fn(() => () => undefined), }); + activateGatewayRuntime(); dispatchAgent.mockReset(); dispatchAgent.mockImplementation(acceptRecoveryDispatch); }); @@ -255,7 +260,6 @@ describe("subagent orphan recovery — faithful restart path", () => { let strictWriteCount = 0; testing.setDepsForTest({ ...createSubagentRegistryTestDeps(), - getGatewayRecoveryRuntime: () => gatewayRuntime, runSubagentAnnounceFlow: vi.fn(async () => "delivered" as const), onAgentEvent: vi.fn(() => () => undefined), persistSubagentRunsToDiskOrThrow: (runs, changedRunIds) => { @@ -353,6 +357,7 @@ describe("subagent orphan recovery — faithful restart path", () => { acceptedAdmission?.release(); rotateAgentEventLifecycleGeneration(); initSubagentRegistry(); + activateGatewayRuntime(); const restored = subagentRuns.get(runId); expect(restored?.execution.restartRecovery).toMatchObject({ sessionMarker: `sess-lost-acceptance:${now}`, @@ -407,7 +412,6 @@ describe("subagent orphan recovery — faithful restart path", () => { let strictWriteCount = 0; testing.setDepsForTest({ ...createSubagentRegistryTestDeps(), - getGatewayRecoveryRuntime: () => gatewayRuntime, runSubagentAnnounceFlow: vi.fn(async () => "delivered" as const), onAgentEvent: vi.fn(() => () => undefined), persistSubagentRunsToDiskOrThrow: (runs, changedRunIds) => { @@ -455,11 +459,11 @@ describe("subagent orphan recovery — faithful restart path", () => { testing.setDepsForTest({ ...createSubagentRegistryTestDeps(), callGateway, - getGatewayRecoveryRuntime: () => gatewayRuntime, runSubagentAnnounceFlow: vi.fn(async () => "delivered" as const), onAgentEvent: vi.fn(() => () => undefined), }); initSubagentRegistry(); + activateGatewayRuntime(); await Promise.resolve(); expect( callGatewayRequests.mock.calls.some( @@ -527,11 +531,11 @@ describe("subagent orphan recovery — faithful restart path", () => { rotateAgentEventLifecycleGeneration(); testing.setDepsForTest({ ...createSubagentRegistryTestDeps(), - getGatewayRecoveryRuntime: () => gatewayRuntime, runSubagentAnnounceFlow: vi.fn(async () => "delivered" as const), onAgentEvent: vi.fn(() => () => undefined), }); initSubagentRegistry(); + activateGatewayRuntime(); await Promise.resolve(); await testing.sweepOnceForTests(); @@ -560,6 +564,7 @@ describe("subagent orphan recovery — faithful restart path", () => { resetSubagentRegistryForTests({ persist: false }); rotateAgentEventLifecycleGeneration(); initSubagentRegistry(); + activateGatewayRuntime(); await Promise.resolve(); await testing.sweepOnceForTests(); @@ -602,7 +607,6 @@ describe("subagent orphan recovery — faithful restart path", () => { let strictWriteCount = 0; testing.setDepsForTest({ ...createSubagentRegistryTestDeps(), - getGatewayRecoveryRuntime: () => gatewayRuntime, runSubagentAnnounceFlow: vi.fn(async () => "delivered" as const), onAgentEvent: vi.fn(() => () => undefined), persistSubagentRunsToDiskOrThrow: (runs, changedRunIds) => { diff --git a/src/agents/subagents/registry/subagent-registry-deps.ts b/src/agents/subagents/registry/subagent-registry-deps.ts index 19d487b6d108..7724ee82c14f 100644 --- a/src/agents/subagents/registry/subagent-registry-deps.ts +++ b/src/agents/subagents/registry/subagent-registry-deps.ts @@ -4,8 +4,6 @@ import type { OpenClawConfig } from "../../../config/types.openclaw.js"; import type { ResolveContextEngineOptions } from "../../../context-engine/registry.js"; import type { ContextEngine } from "../../../context-engine/types.js"; import { callGateway } from "../../../gateway/call.js"; -import type { GatewayRecoveryRuntime } from "../../../gateway/server-instance-runtime.types.js"; -import { getGatewayRecoveryRuntime } from "../../../gateway/server-recovery-runtime-context.js"; import { onAgentEvent, type AgentEventPayload } from "../../../infra/agent-events.js"; import type { PluginRegistry } from "../../../plugins/registry-types.js"; import { createLazyImportLoader, createLazyPromiseLoader } from "../../../shared/lazy-promise.js"; @@ -32,7 +30,6 @@ type BrowserCleanupModule = Pick< export type SubagentRegistryDeps = { callGateway: typeof callGateway; - getGatewayRecoveryRuntime: () => GatewayRecoveryRuntime | undefined; captureSubagentCompletionReply: SubagentAnnounceModule["captureSubagentCompletionReply"]; cleanupBrowserSessionsForLifecycleEnd: typeof cleanupBrowserSessionsForLifecycleEnd; getRuntimeConfig: typeof getRuntimeConfig; @@ -74,7 +71,6 @@ async function loadCleanupBrowserSessionsForLifecycleEnd(): Promise< const defaultSubagentRegistryDeps: SubagentRegistryDeps = { callGateway, - getGatewayRecoveryRuntime, captureSubagentCompletionReply: async (sessionKey, options) => (await loadSubagentAnnounceModule()).captureSubagentCompletionReply(sessionKey, options), cleanupBrowserSessionsForLifecycleEnd: async (params) => diff --git a/src/agents/subagents/registry/subagent-registry-restore.ts b/src/agents/subagents/registry/subagent-registry-restore.ts index 6f493f2c4401..b53aa3c9d65d 100644 --- a/src/agents/subagents/registry/subagent-registry-restore.ts +++ b/src/agents/subagents/registry/subagent-registry-restore.ts @@ -1,4 +1,5 @@ import { ADMIN_SCOPE } from "../../../gateway/method-scopes.js"; +import type { GatewayContextResolver } from "../../../gateway/server-methods/types.js"; import { getAgentEventLifecycleGeneration, isAgentEventLifecycleGenerationCurrent, @@ -55,6 +56,7 @@ export function createSubagentRegistryRestorer(config: { runs: Map; resumedRuns: Set; deps: () => SubagentRegistryDeps; + getGatewayContextResolver: () => GatewayContextResolver | undefined; persist: (...runIds: string[]) => void; persistOrThrow: (...runIds: string[]) => void; settleRequesterTurn: SubagentLifecycleController["settleRequesterTurnAfterSessionSpawns"]; @@ -84,13 +86,13 @@ export function createSubagentRegistryRestorer(config: { ) => Promise; settleFailedQueuedSubagentLaunch: (runId: string, error: string) => boolean; completeCollectorLaunchCleanup: (runId: string) => void; - scheduleSweep: (params?: { delayMs?: number }) => void; warn: (message: string, meta?: Record) => void; }) { const { runs, resumedRuns, deps, + getGatewayContextResolver, persist, persistOrThrow, settleRequesterTurn, @@ -103,10 +105,11 @@ export function createSubagentRegistryRestorer(config: { cleanupCollectorLaunchResources, settleFailedQueuedSubagentLaunch, completeCollectorLaunchCleanup, - scheduleSweep, warn, } = config; let restoreState: "idle" | "in-progress" | "succeeded" = "idle"; + let activationRequested = false; + let activated = false; // A dependency can merge rows before throwing. Keep their reconciliation // pending because mergeOnly correctly reports them as existing on retry. let restoredRowsPending = false; @@ -138,6 +141,170 @@ export function createSubagentRegistryRestorer(config: { restoredRowsPending = false; restoreState = "succeeded"; clearRestoreRetryTimer(); + if (activationRequested) { + activateRestoredRuns(); + } + } + + function activateRestoredRuns() { + activationRequested = true; + if (restoreState !== "succeeded" || activated) { + return; + } + const cfg = deps().getRuntimeConfig(); + const requesterTurns = new Map>(); + const resolveRequesterAgentId = (entry: SubagentRunRecord) => + resolveSubagentRequesterAgentId(cfg, entry); + for (const entry of runs.values()) { + const requesterTurnRunId = entry.requesterTurnRunId?.trim(); + if (!requesterTurnRunId) { + continue; + } + const requesterIdentity = `${resolveRequesterAgentId(entry) ?? "unknown"}\0${entry.requesterSessionKey}`; + let turns = requesterTurns.get(requesterIdentity); + if (!turns) { + turns = new Map(); + requesterTurns.set(requesterIdentity, turns); + } + const entries = turns.get(requesterTurnRunId) ?? []; + entries.push(entry); + turns.set(requesterTurnRunId, entries); + } + for (const [, turns] of requesterTurns) { + for (const [requesterTurnRunId, entries] of turns) { + const firstEntry = entries[0]; + if (!firstEntry) { + continue; + } + settleRequesterTurn({ + requesterSessionKey: firstEntry.requesterSessionKey, + requesterAgentId: resolveRequesterAgentId(firstEntry), + requesterTurnRunId, + requesterYielded: entries.every((entry) => entry.requesterTurnYielded === true), + acceptedSessionSpawns: entries.map((entry) => ({ + runId: entry.taskRunId ?? entry.runId, + childSessionKey: entry.childSessionKey, + })), + }); + } + } + if (runs.size === 0) { + activated = true; + return; + } + + ensureListener(); + // Session-mode runs have no archive deadline but still need TTL cleanup. + startSweeper(); + const restoredSessionCache: SubagentSessionStoreCache = new Map(); + for (const [runId, entry] of runs) { + // Restart recovery exclusively owns receipt-bearing source rows until it + // remaps or terminalizes them. Generic resume would wait on an obsolete run. + if (entry.execution.restartRecovery || entry.killIntent || entry.killReconciliation) { + continue; + } + if (entry.collect && entry.execution.status === "queued") { + const cleanupSessionEntry = loadSubagentSessionEntry({ + childSessionKey: entry.childSessionKey, + storeCache: restoredSessionCache, + }); + const launch = entry.queuedLaunch; + if (!launch) { + const cleanupLifecycleGeneration = getAgentEventLifecycleGeneration(); + void failAndCleanupRestoredQueuedRun( + runId, + entry, + "queued collector launch state was unavailable after restart", + false, + cleanupLifecycleGeneration, + cleanupSessionEntry?.sessionId, + cleanupSessionEntry?.lifecycleRevision, + ); + continue; + } + const groupRuns = listSwarmRunsForGroup( + entry.groupId ?? "", + entry.swarmRequesterSessionKey ?? entry.requesterSessionKey, + entry.requesterAgentId, + ); + const currentSwarmConfig = resolveSwarmConfig(cfg, entry.requesterAgentId); + let launchTerminationConfirmed = false; + let launchLifecycleGeneration: string | undefined; + enqueueSwarmRun({ + groupId: launch.schedulerGroupKey, + runId, + maxConcurrent: currentSwarmConfig.maxConcurrent, + activeRunIds: groupRuns + .filter((candidate) => candidate.execution.status === "running") + .map((candidate) => candidate.schedulerSlotId ?? candidate.runId), + start: async () => { + await runWithGatewayIndependentRootWorkAdmission(async () => { + launchLifecycleGeneration = getAgentEventLifecycleGeneration(); + const request = { + params: applySubagentLaunchAuthorization(launch.request, launch.authorization), + timeoutMs: launch.timeoutMs, + }; + const gatewayRuntime = getGatewayContextResolver()?.()?.recoveryRuntime; + if (!gatewayRuntime) { + throw new GatewayDrainingError(); + } + const response = await gatewayRuntime.dispatchAgent( + request.params as Parameters[0], + request.timeoutMs, + launch.authorization + ? { allowModelOverride: true, scopes: [ADMIN_SCOPE] } + : undefined, + ); + const gatewayRunId = readGatewayRunId(response) ?? runId; + try { + if (!startQueuedSubagentRun(runId, gatewayRunId, launchLifecycleGeneration)) { + throw new Error( + "collector registry row could not transition from queued to running", + ); + } + } catch (error) { + await terminateAcceptedRestoredCollectorRun({ + entry, + gatewayRunId, + timeoutMs: launch.timeoutMs, + expectedSessionId: cleanupSessionEntry?.sessionId, + expectedLifecycleRevision: cleanupSessionEntry?.lifecycleRevision, + }); + launchTerminationConfirmed = true; + throw error; + } + }); + }, + onStartFailure: (error) => { + if (error instanceof GatewayDrainingError) { + return false; + } + return failAndCleanupRestoredQueuedRun( + runId, + entry, + error instanceof Error ? error.message : String(error), + launchTerminationConfirmed, + launchLifecycleGeneration ?? getAgentEventLifecycleGeneration(), + cleanupSessionEntry?.sessionId, + cleanupSessionEntry?.lifecycleRevision, + ); + }, + }); + continue; + } + // An aborted persisted session belongs to orphan recovery. Waiting on its + // pre-restart run can terminalize it before the replacement turn starts. + if ( + loadSubagentSessionEntry({ + childSessionKey: entry.childSessionKey, + storeCache: restoredSessionCache, + })?.abortedLastRun === true + ) { + continue; + } + resumeRun(runId); + } + activated = true; } function restoreSubagentRunsOnce(retryDelayMs = RESTORE_RETRY_DELAY_MS) { @@ -172,168 +339,6 @@ export function createSubagentRegistryRestorer(config: { if (restoredStateChanged) { persist(); } - const requesterTurns = new Map>(); - const resolveRequesterAgentId = (entry: SubagentRunRecord) => - resolveSubagentRequesterAgentId(cfg, entry); - for (const entry of runs.values()) { - const requesterTurnRunId = entry.requesterTurnRunId?.trim(); - if (!requesterTurnRunId) { - continue; - } - const requesterIdentity = `${resolveRequesterAgentId(entry) ?? "unknown"}\0${entry.requesterSessionKey}`; - let turns = requesterTurns.get(requesterIdentity); - if (!turns) { - turns = new Map(); - requesterTurns.set(requesterIdentity, turns); - } - const entries = turns.get(requesterTurnRunId) ?? []; - entries.push(entry); - turns.set(requesterTurnRunId, entries); - } - for (const [, turns] of requesterTurns) { - for (const [requesterTurnRunId, entries] of turns) { - const firstEntry = entries[0]; - if (!firstEntry) { - continue; - } - settleRequesterTurn({ - requesterSessionKey: firstEntry.requesterSessionKey, - requesterAgentId: resolveRequesterAgentId(firstEntry), - requesterTurnRunId, - requesterYielded: entries.every((entry) => entry.requesterTurnYielded === true), - acceptedSessionSpawns: entries.map((entry) => ({ - runId: entry.taskRunId ?? entry.runId, - childSessionKey: entry.childSessionKey, - })), - }); - } - } - if (runs.size === 0) { - completeRestore(); - return; - } - // Resume pending work. - ensureListener(); - // Always start sweeper — session-mode runs (no archiveAtMs) also need TTL cleanup. - startSweeper(); - const restoredSessionCache: SubagentSessionStoreCache = new Map(); - for (const [runId, entry] of runs) { - // Restart recovery exclusively owns receipt-bearing source rows until it - // remaps or terminalizes them. Generic resume would wait on an obsolete run. - if (entry.execution.restartRecovery || entry.killIntent || entry.killReconciliation) { - continue; - } - if (entry.collect && entry.execution.status === "queued") { - const cleanupSessionEntry = loadSubagentSessionEntry({ - childSessionKey: entry.childSessionKey, - storeCache: restoredSessionCache, - }); - const launch = entry.queuedLaunch; - if (!launch) { - const cleanupLifecycleGeneration = getAgentEventLifecycleGeneration(); - void failAndCleanupRestoredQueuedRun( - runId, - entry, - "queued collector launch state was unavailable after restart", - false, - cleanupLifecycleGeneration, - cleanupSessionEntry?.sessionId, - cleanupSessionEntry?.lifecycleRevision, - ); - continue; - } - const groupRuns = listSwarmRunsForGroup( - entry.groupId ?? "", - entry.swarmRequesterSessionKey ?? entry.requesterSessionKey, - entry.requesterAgentId, - ); - const currentSwarmConfig = resolveSwarmConfig( - deps().getRuntimeConfig(), - entry.requesterAgentId, - ); - let launchTerminationConfirmed = false; - let launchLifecycleGeneration: string | undefined; - enqueueSwarmRun({ - groupId: launch.schedulerGroupKey, - runId, - maxConcurrent: currentSwarmConfig.maxConcurrent, - activeRunIds: groupRuns - .filter((candidate) => candidate.execution.status === "running") - .map((candidate) => candidate.schedulerSlotId ?? candidate.runId), - start: async () => { - await runWithGatewayIndependentRootWorkAdmission(async () => { - launchLifecycleGeneration = getAgentEventLifecycleGeneration(); - const request = { - method: "agent", - params: applySubagentLaunchAuthorization(launch.request, launch.authorization), - // Restart replay must restore the trusted launch capability; otherwise - // the queued child silently falls back to its session/default route. - ...(launch.authorization ? { scopes: [ADMIN_SCOPE] } : {}), - timeoutMs: launch.timeoutMs, - }; - const gatewayRuntime = deps().getGatewayRecoveryRuntime(); - const response = gatewayRuntime - ? await gatewayRuntime.dispatchAgent( - request.params as Parameters[0], - request.timeoutMs, - launch.authorization - ? { allowModelOverride: true, scopes: [ADMIN_SCOPE] } - : undefined, - ) - : await deps().callGateway(request); - const gatewayRunId = readGatewayRunId(response) ?? runId; - try { - if (!startQueuedSubagentRun(runId, gatewayRunId, launchLifecycleGeneration)) { - throw new Error( - "collector registry row could not transition from queued to running", - ); - } - } catch (error) { - await terminateAcceptedRestoredCollectorRun({ - entry, - gatewayRunId, - timeoutMs: launch.timeoutMs, - expectedSessionId: cleanupSessionEntry?.sessionId, - expectedLifecycleRevision: cleanupSessionEntry?.lifecycleRevision, - }); - launchTerminationConfirmed = true; - throw error; - } - }); - }, - onStartFailure: (error) => { - if (error instanceof GatewayDrainingError) { - return false; - } - return failAndCleanupRestoredQueuedRun( - runId, - entry, - error instanceof Error ? error.message : String(error), - launchTerminationConfirmed, - launchLifecycleGeneration ?? getAgentEventLifecycleGeneration(), - cleanupSessionEntry?.sessionId, - cleanupSessionEntry?.lifecycleRevision, - ); - }, - }); - continue; - } - // An aborted persisted session belongs to orphan recovery. Waiting on its - // pre-restart run can terminalize it before the replacement turn starts. - if ( - loadSubagentSessionEntry({ - childSessionKey: entry.childSessionKey, - storeCache: restoredSessionCache, - })?.abortedLastRun === true - ) { - continue; - } - resumeRun(runId); - } - - // Cold-start restore can precede instance-runtime registration. The post-attach - // startup pass retries this seam once the lifecycle-bound principal exists. - scheduleSweep(); completeRestore(); } catch (err) { restoredRowsPending ||= runs.size > runCountBeforeRestore; @@ -517,10 +522,13 @@ export function createSubagentRegistryRestorer(config: { return { restoreOnce: restoreSubagentRunsOnce, + activate: activateRestoredRuns, reset: () => { clearRestoreRetryTimer(); restoreState = "idle"; restoredRowsPending = false; + activationRequested = false; + activated = false; }, }; } diff --git a/src/agents/subagents/registry/subagent-registry.announce-loop-guard.test.ts b/src/agents/subagents/registry/subagent-registry.announce-loop-guard.test.ts index 59acf7aa1977..fbfc759d88d3 100644 --- a/src/agents/subagents/registry/subagent-registry.announce-loop-guard.test.ts +++ b/src/agents/subagents/registry/subagent-registry.announce-loop-guard.test.ts @@ -79,6 +79,16 @@ vi.mock("../../timeout.js", () => ({ describe("announce loop guard (#18264)", () => { let registry: typeof import("./subagent-registry.test-helpers.js"); + function hydrateAndActivateRegistry() { + registry.initSubagentRegistry(); + const recoveryRuntime = { + dispatchAgent: vi.fn(), + waitForAgent: vi.fn(async () => ({ status: "pending" })), + sendRecoveryNotice: vi.fn(), + }; + registry.activateSubagentRegistry(() => ({ recoveryRuntime }) as never); + } + function requireRunById(runs: SubagentRunRecord[], runId: string): SubagentRunRecord { const entry = runs.find((run) => run.runId === runId); if (!entry) { @@ -198,7 +208,7 @@ describe("announce loop guard (#18264)", () => { // Initialization finalizes expired pending rows without another recipient-visible attempt. const beforeInit = Date.now(); - registry.initSubagentRegistry(); + hydrateAndActivateRegistry(); await flushAsync(); expect(mocks.runSubagentAnnounceFlow).not.toHaveBeenCalled(); @@ -231,7 +241,7 @@ describe("announce loop guard (#18264)", () => { }; mocks.loadSubagentRegistryFromSqlite.mockReturnValue(new Map([[entry.runId, entry]])); - registry.initSubagentRegistry(); + hydrateAndActivateRegistry(); const resumed = await waitForRun( entry.runId, (run) => run.delivery?.attemptCount === 4 && typeof run.delivery.nextAttemptAt === "number", @@ -279,7 +289,7 @@ describe("announce loop guard (#18264)", () => { ]), ); - registry.initSubagentRegistry(); + hydrateAndActivateRegistry(); await flushAsync(); expect(mocks.runSubagentAnnounceFlow).toHaveBeenCalledTimes(1); @@ -315,7 +325,7 @@ describe("announce loop guard (#18264)", () => { ]), ); - registry.initSubagentRegistry(); + hydrateAndActivateRegistry(); await flushAsync(); const stored = await waitForRun( diff --git a/src/agents/subagents/registry/subagent-registry.persistence.resume.test.ts b/src/agents/subagents/registry/subagent-registry.persistence.resume.test.ts index 0e5664bad7d8..162921eb0572 100644 --- a/src/agents/subagents/registry/subagent-registry.persistence.resume.test.ts +++ b/src/agents/subagents/registry/subagent-registry.persistence.resume.test.ts @@ -28,6 +28,17 @@ let callGatewayModule: typeof import("../../../gateway/call.js"); let agentEventsModule: typeof import("../../../infra/agent-events.js"); let registryStateDbModule: typeof import("../../../state/openclaw-state-db.js"); +function activateRegistry() { + const recoveryRuntime = { + dispatchAgent: (params: Record, timeoutMs?: number) => + callGatewayModule.callGateway({ method: "agent", params, timeoutMs }), + waitForAgent: (params: Record, timeoutMs?: number) => + callGatewayModule.callGateway({ method: "agent.wait", params, timeoutMs }), + sendRecoveryNotice: vi.fn(), + }; + mod.activateSubagentRegistry(() => ({ recoveryRuntime }) as never); +} + describe("subagent registry persistence resume", () => { let tempStateDir: string | null = null; @@ -97,6 +108,7 @@ describe("subagent registry persistence resume", () => { }); mod.initSubagentRegistry(); + activateRegistry(); await vi.waitFor(() => expect(announceSpy).toHaveBeenCalled(), { timeout: 1_000, @@ -169,6 +181,7 @@ describe("subagent registry persistence resume", () => { }); mod.initSubagentRegistry(); + activateRegistry(); await vi.waitFor(() => expect(announceSpy).toHaveBeenCalled(), { timeout: 1_000, @@ -180,6 +193,170 @@ describe("subagent registry persistence resume", () => { }); }); + it("keeps restored recovery dormant until the Gateway lifecycle activates it", async () => { + tempStateDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-subagent-")); + const stateDir = tempStateDir; + const wakeRequester = vi.fn(async () => false); + mod.testing.setDepsForTest({ + ...createSubagentRegistryTestDeps({ + callGateway: vi.mocked(callGatewayModule.callGateway), + maybeWakeRequesterAfterAllChildrenSettled: wakeRequester, + }), + }); + + await withEnvAsync({ OPENCLAW_STATE_DIR: stateDir }, async () => { + const endedAt = Date.now(); + const yieldedRun: SubagentRunRecord = { + runId: "run-hydrated-yield", + taskRunId: "run-hydrated-yield", + requesterTurnRunId: "run-requester", + requesterTurnYielded: true, + childSessionKey: "agent:main:subagent:hydrated-yield", + requesterSessionKey: "agent:main:main", + requesterDisplayKey: "main", + task: "wake only after lifecycle activation", + cleanup: "keep", + createdAt: endedAt - 1_000, + endedReason: "subagent-complete", + execution: { + status: "terminal", + startedAt: endedAt - 500, + endedAt, + outcome: { status: "ok" }, + }, + expectsCompletionMessage: true, + completion: { required: true, resultText: "done", capturedAt: endedAt }, + delivery: { status: "delivered", deliveredAt: endedAt }, + cleanupHandled: true, + cleanupCompletedAt: endedAt, + }; + const queuedCollector: SubagentRunRecord = { + runId: "run-hydrated-collector", + childSessionKey: "agent:main:subagent:hydrated-collector", + requesterSessionKey: "agent:main:main", + requesterDisplayKey: "main", + task: "clean only after lifecycle activation", + cleanup: "keep", + createdAt: endedAt - 500, + collect: true, + swarmRequesterSessionKey: "agent:main:main", + groupId: "hydrated-group", + archiveAtMs: endedAt - 1, + execution: { + status: "terminal", + startedAt: endedAt - 400, + endedAt, + outcome: { status: "error", error: "launch failed" }, + }, + completion: { required: true }, + delivery: { status: "pending" }, + collectorCompletion: { status: "failed" }, + collectorLaunchCleanupPending: true, + }; + const runningRun: SubagentRunRecord = { + runId: "run-hydrated-running", + childSessionKey: "agent:main:subagent:hydrated-running", + requesterSessionKey: "agent:main:main", + requesterDisplayKey: "main", + task: "wait through the activated instance", + cleanup: "keep", + createdAt: endedAt, + execution: { status: "running", startedAt: endedAt }, + completion: { required: false }, + delivery: { status: "not_required" }, + }; + saveSubagentRegistryToSqlite( + new Map([ + [yieldedRun.runId, yieldedRun], + [queuedCollector.runId, queuedCollector], + [runningRun.runId, runningRun], + ]), + ); + await writeSubagentSessionEntry({ + stateDir, + agentId: "main", + sessionKey: yieldedRun.childSessionKey, + sessionId: "sess-hydrated-yield", + defaultSessionId: "sess-hydrated-yield", + }); + await writeSubagentSessionEntry({ + stateDir, + agentId: "main", + sessionKey: queuedCollector.childSessionKey, + sessionId: "sess-hydrated-collector", + defaultSessionId: "sess-hydrated-collector", + lifecycleRevision: "revision-hydrated-collector", + }); + await writeSubagentSessionEntry({ + stateDir, + agentId: "main", + sessionKey: runningRun.childSessionKey, + sessionId: "sess-hydrated-running", + defaultSessionId: "sess-hydrated-running", + }); + + mod.initSubagentRegistry(); + await new Promise((resolve) => { + setImmediate(resolve); + }); + + expect(mod.getSubagentRunByRunId(yieldedRun.runId)).toBeDefined(); + expect(mod.getSubagentRunByRunId(queuedCollector.runId)).toBeDefined(); + expect(mod.getSubagentRunByRunId(runningRun.runId)).toBeDefined(); + expect(wakeRequester).not.toHaveBeenCalled(); + expect(callGatewayModule.callGateway).not.toHaveBeenCalledWith( + expect.objectContaining({ method: "sessions.delete" }), + ); + + const recoveryRuntime = { + dispatchAgent: vi.fn(), + waitForAgent: vi.fn(async () => ({ status: "pending" })), + sendRecoveryNotice: vi.fn(), + }; + let firstLifecycleOpen = true; + const resolveGatewayContext = vi.fn(() => + firstLifecycleOpen ? ({ recoveryRuntime } as never) : undefined, + ); + mod.activateSubagentRegistry(resolveGatewayContext); + mod.activateSubagentRegistry(resolveGatewayContext); + + await vi.waitFor(() => { + expect(wakeRequester).toHaveBeenCalledOnce(); + expect(recoveryRuntime.waitForAgent).toHaveBeenCalledOnce(); + }); + expect(recoveryRuntime.dispatchAgent).not.toHaveBeenCalled(); + expect(callGatewayModule.callGateway).not.toHaveBeenCalledWith( + expect.objectContaining({ method: "agent.wait" }), + ); + + firstLifecycleOpen = false; + expect(resolveGatewayContext()).toBeUndefined(); + const replacementRuntime = { + dispatchAgent: vi.fn(), + waitForAgent: vi.fn(async () => ({ status: "pending" })), + sendRecoveryNotice: vi.fn(), + }; + const resolveReplacementContext = () => ({ recoveryRuntime: replacementRuntime }) as never; + mod.activateSubagentRegistry(resolveReplacementContext); + mod.activateSubagentRegistry(resolveReplacementContext); + expect(wakeRequester).toHaveBeenCalledOnce(); + expect(recoveryRuntime.waitForAgent).toHaveBeenCalledOnce(); + expect(replacementRuntime.waitForAgent).not.toHaveBeenCalled(); + + await mod.testing.runSweeperTickForTests(); + expect(callGatewayModule.callGateway).toHaveBeenCalledTimes(1); + expect(callGatewayModule.callGateway).toHaveBeenCalledWith( + expect.objectContaining({ + method: "sessions.delete", + params: expect.objectContaining({ + expectedSessionId: "sess-hydrated-collector", + expectedLifecycleRevision: "revision-hydrated-collector", + }), + }), + ); + }); + }); + it("keeps dismissed terminal delivery dormant and TTL-eligible after restore", async () => { tempStateDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-subagent-")); const stateDir = tempStateDir; @@ -269,6 +446,7 @@ describe("subagent registry persistence resume", () => { }); mod.initSubagentRegistry(); + activateRegistry(); const restored = mod.getSubagentRunByRunId(run.runId); expect(restored).toMatchObject({ runId: run.runId, taskRunId: run.taskRunId }); diff --git a/src/agents/subagents/registry/subagent-registry.persistence.test-support.ts b/src/agents/subagents/registry/subagent-registry.persistence.test-support.ts index 43a6706b7ec3..b836800b358a 100644 --- a/src/agents/subagents/registry/subagent-registry.persistence.test-support.ts +++ b/src/agents/subagents/registry/subagent-registry.persistence.test-support.ts @@ -69,6 +69,7 @@ export async function writeSubagentSessionEntry(params: { sessionId?: string; updatedAt?: number; abortedLastRun?: boolean; + lifecycleRevision?: string; agentId: string; defaultSessionId: string; }): Promise { @@ -81,6 +82,7 @@ export async function writeSubagentSessionEntry(params: { ...(typeof params.abortedLastRun === "boolean" ? { abortedLastRun: params.abortedLastRun } : {}), + ...(params.lifecycleRevision ? { lifecycleRevision: params.lifecycleRevision } : {}), }; await replaceSessionEntry({ storePath, sessionKey: params.sessionKey }, entry); return storePath; @@ -111,11 +113,6 @@ export function createSubagentRegistryTestDeps( ensureContextEnginesInitialized: vi.fn(), loadAgentRuntimePluginRegistryHandle: vi.fn(), getRuntimeConfig: vi.fn(() => ({})), - getGatewayRecoveryRuntime: vi.fn(() => ({ - dispatchAgent: vi.fn(), - waitForAgent: vi.fn(), - sendRecoveryNotice: vi.fn(), - })), resolveAgentTimeoutMs: vi.fn(() => 100), resolveContextEngine: vi.fn(async () => ({ info: { id: "test", name: "Test", version: "0.0.1" }, diff --git a/src/agents/subagents/registry/subagent-registry.persistence.test.ts b/src/agents/subagents/registry/subagent-registry.persistence.test.ts index c4cfa53e6dea..6f2450077f1e 100644 --- a/src/agents/subagents/registry/subagent-registry.persistence.test.ts +++ b/src/agents/subagents/registry/subagent-registry.persistence.test.ts @@ -31,6 +31,7 @@ import { } from "./subagent-registry.store.sqlite.js"; import { testing, + activateSubagentRegistry, addSubagentRunForTests, clearSubagentRunSteerRestart, getSubagentRunByChildSessionKey, @@ -187,6 +188,14 @@ describe("subagent registry persistence", () => { const restartRegistry = () => { resetSubagentRegistryForTests({ persist: false }); initSubagentRegistry(); + const recoveryRuntime = { + dispatchAgent: (params: Record, timeoutMs?: number) => + callGateway({ method: "agent", params, timeoutMs }), + waitForAgent: (params: Record, timeoutMs?: number) => + callGateway({ method: "agent.wait", params, timeoutMs }), + sendRecoveryNotice: vi.fn(), + }; + activateSubagentRegistry(() => ({ recoveryRuntime }) as never); }; const fastPersistSubagentRunsToDisk = (runs: Map) => diff --git a/src/agents/subagents/registry/subagent-registry.test-helpers.ts b/src/agents/subagents/registry/subagent-registry.test-helpers.ts index d14dfa616aff..4aac662d8667 100644 --- a/src/agents/subagents/registry/subagent-registry.test-helpers.ts +++ b/src/agents/subagents/registry/subagent-registry.test-helpers.ts @@ -50,9 +50,6 @@ type RegistryTestApi = { type RegistryDeps = { callGateway: typeof import("../../../gateway/call.js").callGateway; - getGatewayRecoveryRuntime: () => - | import("../../../gateway/server-instance-runtime.types.js").GatewayRecoveryRuntime - | undefined; captureSubagentCompletionReply: typeof import("../announce/subagent-announce.js").captureSubagentCompletionReply; cleanupBrowserSessionsForLifecycleEnd: typeof import("../../../browser-lifecycle-cleanup.js").cleanupBrowserSessionsForLifecycleEnd; getRuntimeConfig: typeof import("../../../config/config.js").getRuntimeConfig; diff --git a/src/agents/subagents/registry/subagent-registry.test.ts b/src/agents/subagents/registry/subagent-registry.test.ts index f3145974992e..f0eae216e48b 100644 --- a/src/agents/subagents/registry/subagent-registry.test.ts +++ b/src/agents/subagents/registry/subagent-registry.test.ts @@ -423,6 +423,21 @@ describe("subagent registry seam flow", () => { }, }); let mod: RegistryHarness; + const recoveryRuntime: GatewayRecoveryRuntime = { + dispatchAgent: mocks.dispatchRecoveryAgent as GatewayRecoveryRuntime["dispatchAgent"], + waitForAgent: (params, timeoutMs) => + mocks.callGateway({ + method: "agent.wait", + params: params as unknown as Record, + timeoutMs, + }) as never, + sendRecoveryNotice: vi.fn(), + }; + const activateRegistry = () => mod.activateSubagentRegistry(() => ({ recoveryRuntime }) as never); + const hydrateAndActivateRegistry = () => { + mod.initSubagentRegistry(); + activateRegistry(); + }; const findRequesterRun = (runId: string) => mod.listSubagentRunsForRequester("agent:main:main").find((entry) => entry.runId === runId); const mockPendingAgentWait = () => @@ -530,7 +545,6 @@ describe("subagent registry seam flow", () => { callGateway: mocks.callGateway as typeof import("../../../gateway/call.js").callGateway, captureSubagentCompletionReply: mocks.captureSubagentCompletionReply, cleanupBrowserSessionsForLifecycleEnd: mocks.cleanupBrowserSessionsForLifecycleEnd, - getGatewayRecoveryRuntime: mocks.getGatewayRecoveryRuntime, onAgentEvent: mocks.onAgentEvent, persistSubagentRunsToDisk: mocks.persistSubagentRunsToDisk, persistSubagentRunsToDiskOrThrow: mocks.persistSubagentRunsToDiskOrThrow, @@ -1251,7 +1265,7 @@ describe("subagent registry seam flow", () => { }) as never) .mockReturnValue(0); - mod.initSubagentRegistry(); + hydrateAndActivateRegistry(); expect(mocks.restoreSubagentRunsFromDisk).toHaveBeenCalledOnce(); expect(mocks.onAgentEvent).not.toHaveBeenCalled(); @@ -1314,7 +1328,7 @@ describe("subagent registry seam flow", () => { return 1; }) as never); - mod.initSubagentRegistry(); + hydrateAndActivateRegistry(); await waitForFast(() => { expect(mocks.maybeWakeRequesterAfterAllChildrenSettled).toHaveBeenCalledTimes(1); @@ -1353,7 +1367,7 @@ describe("subagent registry seam flow", () => { return 1; }) as never); - mod.initSubagentRegistry(); + hydrateAndActivateRegistry(); await Promise.resolve(); await Promise.resolve(); @@ -1402,7 +1416,7 @@ describe("subagent registry seam flow", () => { const suspension = tryBeginGatewaySuspendAdmission(() => {}); expect(suspension?.commit()).toBe(true); - mod.initSubagentRegistry(); + hydrateAndActivateRegistry(); await Promise.resolve(); expect(mocks.callGateway.mock.calls.filter(([request]) => request.method === "agent")).toEqual( [], @@ -1619,7 +1633,7 @@ describe("subagent registry seam flow", () => { return request.method === "agent.wait" ? { status: "pending" } : {}; }); - mod.initSubagentRegistry(); + hydrateAndActivateRegistry(); await waitForFast(() => expect(releaseAbort).toBeTypeOf("function")); expect(agentCalls).toBe(1); @@ -1700,7 +1714,7 @@ describe("subagent registry seam flow", () => { return request.method === "agent.wait" ? { status: "pending" } : {}; }); - mod.initSubagentRegistry(); + hydrateAndActivateRegistry(); await waitForFast(() => expect(releaseDelete).toBeTypeOf("function")); expect(agentCalls).toBe(1); @@ -1767,7 +1781,7 @@ describe("subagent registry seam flow", () => { return request.method === "agent.wait" ? { status: "pending" } : {}; }); - mod.initSubagentRegistry(); + hydrateAndActivateRegistry(); await waitForFast(() => expect(persistenceCalls).toBeGreaterThanOrEqual(3)); await concurrentSweep; @@ -1818,7 +1832,7 @@ describe("subagent registry seam flow", () => { agent: new Error("launch failed"), }); - mod.initSubagentRegistry(); + hydrateAndActivateRegistry(); await waitForFast(() => expect(mod.getSubagentRunByRunId("run-queued-failure")).toMatchObject({ @@ -1905,7 +1919,7 @@ describe("subagent registry seam flow", () => { return {}; }); - mod.initSubagentRegistry(); + hydrateAndActivateRegistry(); await waitForFast(() => expect(mod.getSubagentRunByRunId("run-queued-cleanup-retry")).toMatchObject({ @@ -2079,9 +2093,7 @@ describe("subagent registry seam flow", () => { }); it("does not fall back to network recovery without an instance-bound runtime", async () => { - mod.testing.setDepsForTest({ - getGatewayRecoveryRuntime: () => undefined, - }); + mod.activateSubagentRegistry(() => undefined); mod.scheduleSubagentRegistrySweep({ delayMs: 1 }); await vi.advanceTimersByTimeAsync(1); @@ -2789,7 +2801,7 @@ describe("subagent registry seam flow", () => { }, }); - mod.initSubagentRegistry(); + hydrateAndActivateRegistry(); await waitForFast(() => { const completedRun = findRequesterRun(runId); @@ -3088,7 +3100,7 @@ describe("subagent registry seam flow", () => { }, ); - mod.initSubagentRegistry(); + hydrateAndActivateRegistry(); await waitForFast(() => { expect(waitTimeouts).toEqual([1_000]); @@ -5774,7 +5786,7 @@ describe("subagent registry seam flow", () => { return 1; }) as never); - mod.initSubagentRegistry(); + hydrateAndActivateRegistry(); await Promise.resolve(); await Promise.resolve(); @@ -5817,7 +5829,7 @@ describe("subagent registry seam flow", () => { return 1; }) as never); - mod.initSubagentRegistry(); + hydrateAndActivateRegistry(); await Promise.resolve(); await Promise.resolve(); diff --git a/src/agents/subagents/registry/subagent-registry.ts b/src/agents/subagents/registry/subagent-registry.ts index 5fe64ad4eabd..ba508619a51b 100644 --- a/src/agents/subagents/registry/subagent-registry.ts +++ b/src/agents/subagents/registry/subagent-registry.ts @@ -3,8 +3,8 @@ import type { AgentWaitParams } from "../../../../packages/gateway-protocol/src/ import type { OpenClawConfig } from "../../../config/types.openclaw.js"; import { callGateway } from "../../../gateway/call.js"; import type { GatewayContextResolver } from "../../../gateway/server-methods/types.js"; -import { getGatewayRecoveryRuntime } from "../../../gateway/server-recovery-runtime-context.js"; import { createSubsystemLogger } from "../../../logging/subsystem.js"; +import { bindGatewayContextResolver } from "../../../plugins/runtime/gateway-request-scope.js"; import { isGatewayRestartDraining, runWithGatewayIndependentRootWorkAdmission, @@ -63,6 +63,7 @@ const subagentRegistryBootstrapState: { } = {}; const resumeRetryTimers = new Set>(); +let activeGatewayContextResolver: GatewayContextResolver | undefined; const SUBAGENT_ANNOUNCE_TIMEOUT_MS = 120_000; const GATEWAY_ADMISSION_RETRY_DELAY_MS = 1_000; /** Admission pressure for recoverable completion deliveries; rows are never pruned for capacity. */ @@ -322,6 +323,7 @@ const subagentRestorer = createSubagentRegistryRestorer({ runs: subagentRuns, resumedRuns, deps: () => subagentRegistryDeps, + getGatewayContextResolver: () => activeGatewayContextResolver, persist: persistSubagentRuns, persistOrThrow: persistSubagentRunsOrThrow, settleRequesterTurn: settleRequesterTurnAfterSessionSpawns, @@ -351,7 +353,6 @@ const subagentRestorer = createSubagentRegistryRestorer({ settleFailedQueuedSubagentLaunch: (runId, error) => subagentRunManager.settleFailedQueuedSubagentLaunch(runId, error), completeCollectorLaunchCleanup: (runId) => publicApi.completeCollectorLaunchCleanup(runId), - scheduleSweep: scheduleSubagentRegistrySweep, warn: (message, meta) => log.warn(message, meta), }); @@ -380,7 +381,7 @@ const subagentSweeper = createSubagentRegistrySweeper({ clearPendingLifecycleTimeout, sweepPendingLifecycle: (now) => pendingLifecycle.sweepExpired(now), completeSubagentRunWithRecovery: completionRuntime.completeSubagentRunWithRecovery, - getGatewayRecoveryRuntime: () => subagentRegistryDeps.getGatewayRecoveryRuntime(), + getGatewayRecoveryRuntime: () => activeGatewayContextResolver?.()?.recoveryRuntime, abandonSubagentRestartRecoveryLaunch: (params) => subagentRunManager.abandonSubagentRestartRecoveryLaunch(params), clearAcceptedSubagentRestartRecovery: (params) => @@ -433,7 +434,7 @@ const subagentRunManager = createSubagentRunManager({ persistOrThrow: persistSubagentRunsOrThrow, callGateway: async (request: Parameters[0]) => { if (request.method === "agent.wait") { - const gatewayRuntime = getGatewayRecoveryRuntime(); + const gatewayRuntime = activeGatewayContextResolver?.()?.recoveryRuntime; if (gatewayRuntime) { // Registry waits are Gateway-owned lifecycle work. Keep them on the // owning instance when one exists; standalone processes authenticate normally. @@ -544,6 +545,7 @@ function resetSubagentRegistryForTests(opts?: { persist?: boolean }) { clearSubagentRunsReadCacheForTest(); subagentSweeper.reset(); subagentRestorer.reset(); + activeGatewayContextResolver = undefined; subagentListener.reset(); if (opts?.persist !== false) { persistSubagentRuns(); @@ -599,6 +601,15 @@ export function initSubagentRegistry() { } state.restorer.restoreOnce(); } +export function activateSubagentRegistry(resolveGatewayContext: GatewayContextResolver) { + activeGatewayContextResolver = resolveGatewayContext; + for (const entry of subagentRuns.values()) { + bindGatewayContextResolver(entry, resolveGatewayContext); + } + subagentRestorer.activate(); + // Post-ready only: collector cleanup retains the canonical sessions.delete RPC owner. + scheduleSubagentRegistrySweep(); +} export const settleRequesterAfterSessionSpawns = publicApi.settleRequesterAfterSessionSpawns; export const markRequesterTurnYielded = publicApi.markRequesterTurnYielded; diff --git a/src/agents/tools/swarm-tools.integration.test.ts b/src/agents/tools/swarm-tools.integration.test.ts index e13183dfbc03..7d17b983735b 100644 --- a/src/agents/tools/swarm-tools.integration.test.ts +++ b/src/agents/tools/swarm-tools.integration.test.ts @@ -121,7 +121,6 @@ describe("swarm tools integration", () => { return resultTextBySession.get(sessionKey) ?? ""; }) as never, cleanupBrowserSessionsForLifecycleEnd: vi.fn(async () => undefined), - getGatewayRecoveryRuntime: () => undefined, getRuntimeConfig: () => config, maybeWakeRequesterAfterAllChildrenSettled: vi.fn(async () => false), onAgentEvent: vi.fn(() => () => undefined) as never, diff --git a/src/gateway/server-startup-finish.ts b/src/gateway/server-startup-finish.ts index aaa815f36618..006314df77d6 100644 --- a/src/gateway/server-startup-finish.ts +++ b/src/gateway/server-startup-finish.ts @@ -266,6 +266,7 @@ export async function finishGatewayStartup(params: { deps, startChannels, recoveryRuntime: gatewayInstanceRuntime.recovery, + resolveGatewayContext: gatewayRequestContext.resolveGatewayContext!, logHooks, logChannels, unlockStartupMethods: kernel.unlockStartupMethods, diff --git a/src/gateway/server-startup-plugins.test.ts b/src/gateway/server-startup-plugins.test.ts index f00da3ced48b..7d09a50626a9 100644 --- a/src/gateway/server-startup-plugins.test.ts +++ b/src/gateway/server-startup-plugins.test.ts @@ -326,6 +326,15 @@ describe("prepareGatewayPluginBootstrap startup plugins", () => { expect(migrateLegacyNodePairingStore).not.toHaveBeenCalled(); }); + it("hydrates the subagent registry before plugin bootstrap", async () => { + await prepareBootstrapWithRuntimeConfig({}); + + expect(initSubagentRegistry).toHaveBeenCalledOnce(); + expect(initSubagentRegistry.mock.invocationCallOrder[0]).toBeLessThan( + loadPluginLookUpTable.mock.invocationCallOrder[0]!, + ); + }); + it("derives startup activation from source config instead of runtime plugin defaults", async () => { const sourceConfig = { channels: { diff --git a/src/gateway/server-startup-post-attach.test.ts b/src/gateway/server-startup-post-attach.test.ts index 406669cb9119..f166c69d20bf 100644 --- a/src/gateway/server-startup-post-attach.test.ts +++ b/src/gateway/server-startup-post-attach.test.ts @@ -39,7 +39,7 @@ const hoisted = vi.hoisted(() => { })); const scheduleGatewayUpdateCheck = vi.fn(() => () => {}); const logGatewayStartup = vi.fn(); - const scheduleSubagentRegistrySweep = vi.fn(); + const activateSubagentRegistry = vi.fn(); const markStartupOrphanedMainSessionsForRecovery = vi.fn(async () => ({ marked: 0, skipped: 0, @@ -102,7 +102,7 @@ const hoisted = vi.hoisted(() => { initializeGatewayUpdateStatus, scheduleGatewayUpdateCheck, logGatewayStartup, - scheduleSubagentRegistrySweep, + activateSubagentRegistry, markStartupOrphanedMainSessionsForRecovery, scheduleRestartAbortedMainSessionRecovery, scheduleRestartSentinelWake, @@ -133,7 +133,7 @@ vi.mock("../agents/session-dirs.js", () => ({ })); vi.mock("../agents/subagents/registry/subagent-registry.js", () => ({ - scheduleSubagentRegistrySweep: hoisted.scheduleSubagentRegistrySweep, + activateSubagentRegistry: hoisted.activateSubagentRegistry, })); vi.mock("../agents/main-session-recovery/main-session-restart-recovery-marking.js", () => ({ @@ -479,7 +479,7 @@ describe("startGatewayPostAttachRuntime", () => { hoisted.initializeGatewayUpdateStatus.mockClear(); hoisted.scheduleGatewayUpdateCheck.mockClear(); hoisted.logGatewayStartup.mockClear(); - hoisted.scheduleSubagentRegistrySweep.mockClear(); + hoisted.activateSubagentRegistry.mockClear(); hoisted.markStartupOrphanedMainSessionsForRecovery.mockReset(); hoisted.markStartupOrphanedMainSessionsForRecovery.mockResolvedValue({ marked: 0, @@ -569,6 +569,7 @@ describe("startGatewayPostAttachRuntime", () => { it("re-enables startup-gated methods after post-attach sidecars start", async () => { const unavailableGatewayMethods = new Set(["chat.history", "models.list"]); + const startupOrder: string[] = []; const methodsAtRecoveryRegistration: string[][] = []; const currentConfig = { agents: { list: [{ id: "main" }, { id: "work" }] } }; hoisted.scheduleRestartAbortedMainSessionRecovery.mockImplementationOnce( @@ -577,14 +578,20 @@ describe("startGatewayPostAttachRuntime", () => { expect(params.getConfig()).toBe(currentConfig); }, ); - const onSidecarsReady = vi.fn(); + const onSidecarsReady = vi.fn(() => startupOrder.push("ready")); + hoisted.activateSubagentRegistry.mockImplementationOnce(() => { + startupOrder.push("registry"); + }); const log = { info: vi.fn(), warn: vi.fn() }; await startGatewayPostAttachRuntime({ ...createPostAttachParams(), getConfig: () => currentConfig, log, - unlockStartupMethods: createStartupMethodUnlocker(unavailableGatewayMethods), + unlockStartupMethods: () => { + startupOrder.push("unlock"); + createStartupMethodUnlocker(unavailableGatewayMethods)(); + }, onSidecarsReady, }); @@ -610,7 +617,8 @@ describe("startGatewayPostAttachRuntime", () => { waitForStart: undefined, gatewayRuntime: expect.any(Object), }); - expect(hoisted.scheduleSubagentRegistrySweep).toHaveBeenCalledWith(); + expect(hoisted.activateSubagentRegistry).toHaveBeenCalledWith(expect.any(Function)); + expect(startupOrder).toEqual(["unlock", "ready", "registry"]); expect(methodsAtRecoveryRegistration).toStrictEqual([["chat.history", "models.list"]]); }); @@ -3052,7 +3060,7 @@ describe("startGatewayPostAttachRuntime", () => { expect(startWorkerEnvironmentRuntime).not.toHaveBeenCalled(); }); - it("keeps startup methods fenced when close begins during late recovery loading", async () => { + it("does not activate restored recovery when close begins during activation loading", async () => { let closeStarted = false; let releaseRecoveryLoad: (() => void) | undefined; const recoveryLoadReady = new Promise((resolve) => { @@ -3067,7 +3075,7 @@ describe("startGatewayPostAttachRuntime", () => { const workerSidecar = { stop: vi.fn(async () => {}) }; let ownedWorkerSidecar: typeof workerSidecar | undefined; const unlockStartupMethods = vi.fn(); - const scheduleSubagentRegistrySweep = vi.fn(); + const activateSubagentRegistry = vi.fn(); const onPluginServices = vi.fn(); const onGatewayLifetimeSidecars = vi.fn(); const runtime = await startGatewayPostAttachRuntime( @@ -3091,10 +3099,10 @@ describe("startGatewayPostAttachRuntime", () => { return { pluginServices, postReadySidecars: [postReadySidecar] }; }, ), - loadSubagentRegistrySweep: vi.fn(async () => { + loadSubagentRegistryActivation: vi.fn(async () => { markRecoveryLoadStarted?.(); await recoveryLoadReady; - return scheduleSubagentRegistrySweep; + return activateSubagentRegistry; }), }), ); @@ -3104,8 +3112,8 @@ describe("startGatewayPostAttachRuntime", () => { releaseRecoveryLoad?.(); await expect(runtime.startupSettled).resolves.toBeUndefined(); - expect(scheduleSubagentRegistrySweep).not.toHaveBeenCalled(); - expect(unlockStartupMethods).not.toHaveBeenCalled(); + expect(activateSubagentRegistry).not.toHaveBeenCalled(); + expect(unlockStartupMethods).toHaveBeenCalledOnce(); expect(workerSidecar.stop).not.toHaveBeenCalled(); expect(pluginServices.stop).toHaveBeenCalledOnce(); expect(postReadySidecar.stop).toHaveBeenCalledOnce(); @@ -3457,7 +3465,7 @@ function createPostAttachRuntimeDeps( scheduleGatewayUpdateCheck: hoisted.scheduleGatewayUpdateCheck, startGatewaySidecars: vi.fn(async () => ({ pluginServices: null, postReadySidecars: [] })), warmSystemCa: vi.fn(async () => {}), - loadSubagentRegistrySweep: vi.fn(async () => hoisted.scheduleSubagentRegistrySweep), + loadSubagentRegistryActivation: vi.fn(async () => hoisted.activateSubagentRegistry), ...overrides, }; } @@ -3496,6 +3504,7 @@ function createPostAttachParams(overrides: Partial = {}): Post waitForAgent: vi.fn(), sendRecoveryNotice: vi.fn(), }, + resolveGatewayContext: vi.fn(() => ({ recoveryRuntime: {} }) as never), logHooks: { info: vi.fn(), warn: vi.fn(), diff --git a/src/gateway/server-startup-post-attach.ts b/src/gateway/server-startup-post-attach.ts index 5df2e9a6b6f2..d281d7c9c03a 100644 --- a/src/gateway/server-startup-post-attach.ts +++ b/src/gateway/server-startup-post-attach.ts @@ -34,7 +34,7 @@ import { import type { GatewayBroadcastToConnIdsFn } from "./server-broadcast-types.js"; import type { GatewayControlUiRootLifecycle } from "./server-control-ui-root.js"; import type { GatewayRecoveryRuntime } from "./server-instance-runtime.types.js"; -import type { GatewayClient } from "./server-methods/shared-types.js"; +import type { GatewayClient, GatewayContextResolver } from "./server-methods/shared-types.js"; import type { GatewayResidentRegistry } from "./server-resident-registry.js"; import type { refreshLatestUpdateRestartSentinel } from "./server-restart-sentinel.js"; import type { GatewaySidecarStartupMode } from "./server-sidecar-startup-mode.js"; @@ -943,7 +943,9 @@ type GatewayPostAttachRuntimeDeps = { ) => Awaitable>; startGatewaySidecars: typeof startGatewaySidecars; warmSystemCa: typeof warmMacOSSystemCaOffMainThread; - loadSubagentRegistrySweep: () => Awaitable<() => void>; + loadSubagentRegistryActivation: () => Awaitable< + (resolveGatewayContext: GatewayContextResolver) => void + >; }; const defaultGatewayPostAttachRuntimeDeps: GatewayPostAttachRuntimeDeps = { @@ -958,9 +960,8 @@ const defaultGatewayPostAttachRuntimeDeps: GatewayPostAttachRuntimeDeps = { (await import("../infra/update-startup.js")).scheduleGatewayUpdateCheck(...args), startGatewaySidecars, warmSystemCa: warmMacOSSystemCaOffMainThread, - loadSubagentRegistrySweep: async () => - (await import("../agents/subagents/registry/subagent-registry.js")) - .scheduleSubagentRegistrySweep, + loadSubagentRegistryActivation: async () => + (await import("../agents/subagents/registry/subagent-registry.js")).activateSubagentRegistry, }; function createDeferredGatewayUpdateCheck(params: { @@ -1116,6 +1117,7 @@ export async function startGatewayPostAttachRuntime( startChannels: () => Promise; refreshChatMetadata?: () => Promise; recoveryRuntime: GatewayRecoveryRuntime; + resolveGatewayContext: GatewayContextResolver; logHooks: { info: (msg: string) => void; warn: (msg: string) => void; @@ -1434,17 +1436,6 @@ export async function startGatewayPostAttachRuntime( ], ]); let mainSessionRecoverySidecar: GatewayPostReadySidecarHandle | undefined; - try { - const scheduleSubagentRegistrySweep = await runtimeDeps.loadSubagentRegistrySweep(); - if (params.isClosing?.() !== true) { - scheduleSubagentRegistrySweep(); - } - } catch (err) { - params.log.warn(`subagent restart recovery failed to schedule: ${String(err)}`); - } - if (params.isClosing?.()) { - return await stopStartupSidecars(mainSessionRecoverySidecar); - } try { await startupLog; } catch (error) { @@ -1518,6 +1509,17 @@ export async function startGatewayPostAttachRuntime( ]; params.log.info(formatGatewayStartupOutcomes(startupOutcomes.snapshot())); params.onSidecarsReady?.(); + try { + const activateSubagentRegistry = await runtimeDeps.loadSubagentRegistryActivation(); + if (params.isClosing?.() !== true) { + activateSubagentRegistry(params.resolveGatewayContext); + } + } catch (err) { + params.log.warn(`subagent restart recovery failed to activate: ${String(err)}`); + } + if (params.isClosing?.()) { + return await stopStartupSidecars(mainSessionRecoverySidecar); + } params.startupTrace?.detail("sidecars.ready", [ [ "loadedPluginCount",