From 47e06b824924bc027c5b4c44116884ee80b44dfc Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Fri, 28 Aug 2026 05:03:09 +0800 Subject: [PATCH] fix(gateway): pin runtime generation at admission (#131120) --- ...arness-source-delivery.integration.test.ts | 41 +++-- .../agent-turn/agent-run-admission-phase.ts | 148 +++++++++++++----- .../agent-run-execution-phase.owner.test.ts | 142 +++++++++++++---- .../agent-turn/agent-run-execution-phase.ts | 38 ++--- .../server-methods/agent.create-event.test.ts | 14 ++ .../server-methods/agent.test-harness.ts | 6 + ....gateway-server-agent-auth-refresh.test.ts | 113 +++++++++++-- 7 files changed, 384 insertions(+), 118 deletions(-) diff --git a/src/agents/embedded-agent-runner/run.prepared-harness-source-delivery.integration.test.ts b/src/agents/embedded-agent-runner/run.prepared-harness-source-delivery.integration.test.ts index 49ddfac046a7..255511615874 100644 --- a/src/agents/embedded-agent-runner/run.prepared-harness-source-delivery.integration.test.ts +++ b/src/agents/embedded-agent-runner/run.prepared-harness-source-delivery.integration.test.ts @@ -32,7 +32,10 @@ import type { GetReplyOptions, ReplyPayload } from "../../auto-reply/types.js"; import { createEmptyPluginRegistry } from "../../plugins/registry-empty.js"; import type { FailoverReason } from "../failover/signal.js"; import { registerAgentHarness } from "../harness/registry.js"; -import { withPreparedModelRuntimePluginGenerationScope } from "../prepared-model-runtime-generation-scope.js"; +import { + getPreparedModelRuntimeBorrowedSnapshot, + withPreparedModelRuntimePluginGenerationScope, +} from "../prepared-model-runtime-generation-scope.js"; import type { PreparedModelRuntimePluginGeneration } from "../prepared-model-runtime.types.js"; import { makeAttemptResult } from "./run.overflow-compaction.fixture.js"; import { @@ -534,7 +537,7 @@ describe("prepared harness source delivery", () => { ); }); - it("completes an admitted turn on its generation after a plugin-runtime replacement", async () => { + it("completes an admitted turn on A after plugin-runtime generation B publishes", async () => { const { runEmbeddedAgent } = await loadRunOverflowCompactionHarness(); const config = {}; const workspaceDir = "/tmp/workspace"; @@ -560,28 +563,37 @@ describe("prepared harness source delivery", () => { policyHash: "replacement", workspaceDir, }; + const admittedSnapshot = { + ...baseLease.snapshot, + config, + workspaceDir, + pluginRegistry, + metadataSnapshot: admittedMetadataSnapshot, + } as NonNullable>; + let publishedMetadataSnapshot = admittedMetadataSnapshot; const release = vi.fn(); let servedMetadataSnapshot: unknown; + let publishedMetadataAtAcquire: unknown; mockedAcquireAgentRunPreparedModelRuntime.mockClear(); mockedAcquireAgentRunPreparedModelRuntime.mockImplementationOnce( async ( _input, options?: { - pluginGeneration?: { pluginMetadataSnapshot: typeof admittedMetadataSnapshot }; + pluginGeneration?: PreparedModelRuntimePluginGeneration; }, ) => { - const metadataSnapshot = - options?.pluginGeneration?.pluginMetadataSnapshot ?? replacementMetadataSnapshot; - servedMetadataSnapshot = metadataSnapshot; + const generation = options?.pluginGeneration; + const borrowed = generation + ? getPreparedModelRuntimeBorrowedSnapshot(generation) + : undefined; + if (!borrowed) { + throw new Error("prepared model runtime plugin generation was superseded"); + } + publishedMetadataAtAcquire = publishedMetadataSnapshot; + servedMetadataSnapshot = borrowed.metadataSnapshot; return { ...baseLease, - snapshot: { - ...baseLease.snapshot, - config, - workspaceDir, - pluginRegistry, - metadataSnapshot, - }, + snapshot: borrowed as typeof baseLease.snapshot, release, }; }, @@ -589,6 +601,7 @@ describe("prepared harness source delivery", () => { mockedBuildEmbeddedRunPayloads.mockReturnValue([{ text: "ok" }]); mockedRunEmbeddedAttempt.mockResolvedValueOnce(makeAttemptResult({ assistantTexts: ["ok"] })); useOpenAIPlatformAuthFixture(); + publishedMetadataSnapshot = replacementMetadataSnapshot; const result = await withPreparedModelRuntimePluginGenerationScope( admittedGeneration, @@ -601,12 +614,14 @@ describe("prepared harness source delivery", () => { runId: "admitted-generation-replacement", sessionKey: undefined, }), + () => admittedSnapshot, ); expect(mockedAcquireAgentRunPreparedModelRuntime).toHaveBeenCalledWith( expect.objectContaining({ config, workspaceDir }), expect.objectContaining({ pluginGeneration: admittedGeneration }), ); + expect(publishedMetadataAtAcquire).toBe(replacementMetadataSnapshot); expect(servedMetadataSnapshot).toBe(admittedGeneration.pluginMetadataSnapshot); expect(result.payloads).toEqual([{ text: "ok" }]); expect(release).toHaveBeenCalledOnce(); diff --git a/src/gateway/agent-turn/agent-run-admission-phase.ts b/src/gateway/agent-turn/agent-run-admission-phase.ts index 3c9530885493..a4925a03ef18 100644 --- a/src/gateway/agent-turn/agent-run-admission-phase.ts +++ b/src/gateway/agent-turn/agent-run-admission-phase.ts @@ -14,7 +14,14 @@ import { type MainSessionRecoveryPendingTarget, } from "../../agents/main-session-recovery/main-session-recovery-store.js"; import { resolvePersistedOverrideModelRef } from "../../agents/model-selection.js"; +import { + acquireAgentRunPreparedModelRuntime, + loadPublishedGatewayReplyDispatchRuntime, + type PreparedModelRuntimeLease, + type PreparedReplyDispatchRuntime, +} from "../../agents/prepared-model-runtime.js"; import { resolveProviderIdForAuth } from "../../agents/provider-auth-aliases.js"; +import { resolveIngressWorkspaceOverrideForSessionRun } from "../../agents/spawned-context.js"; import { resolveExactSubagentCompletionEvent, type TrustedSubagentCompletionHandoff, @@ -46,6 +53,7 @@ import { formatForLog } from "../ws-log.js"; import { isPreRegistrationAbortedAgentDedupeEntryForSession, readGatewayDedupeEntry, + setAbortedAgentDedupeEntries, setGatewayDedupeEntries, } from "./agent-dedupe.js"; import type { AgentDeliveryPhaseResult } from "./agent-delivery-phase.js"; @@ -71,8 +79,11 @@ export type PreparedAgentRunDispatch = { lifecycleStorePath: string; resolvedThreadId?: string | number; dispatchTaskTrackingMode: Exclude; + preparedModelRuntimeLease: PreparedModelRuntimeLease; + replyDispatchRuntime: PreparedReplyDispatchRuntime; unpersistedOffloadedRefs: OffloadedRef[]; userTurn: PreparedAgentRunUserTurn; + workspaceOverride?: string; restoreAdmittedRestartRecoveryInterrupted?: () => Promise< MainSessionRecoveryPendingTarget | undefined >; @@ -320,6 +331,84 @@ export async function prepareAgentRunDispatch(params: { } } + const workspaceOverride = resolveIngressWorkspaceOverrideForSessionRun({ + spawnedBy: params.sessionEntry?.spawnedBy, + workspaceDir: params.sessionEntry?.spawnedWorkspaceDir, + cwd: params.sessionEntry?.spawnedCwd, + }); + let preparedModelRuntimeLease: PreparedModelRuntimeLease | undefined; + const cleanupPreaccept = (admissionReleased = false) => { + preparedModelRuntimeLease?.release(); + preparedModelRuntimeLease = undefined; + activeRunAbort.cleanup({ force: true }); + if (!admissionReleased) { + activeGatewayWorkAdmission.release(); + } + }; + const rejectPreaccept = (error: ReturnType) => { + cleanupPreaccept(); + params.io.emitAcceptance([false, undefined, error]); + return undefined; + }; + const revalidateAdmission = () => { + if (activeRunAbort.controller.signal.aborted) { + setAbortedAgentDedupeEntries({ + dedupe: params.context.dedupe, + keys: params.agentDedupeKeys, + agentId: params.admissionAgentId(), + runId: params.runId, + stopReason: activeRunAbort.entry?.abortStopReason ?? "rpc", + }); + } + try { + params.assertGatewayWorkAdmissionAllowed(); + } catch (err) { + rejectPreaccept(errorShapeFromError(ErrorCodes.INVALID_REQUEST, err)); + return false; + } + if (!params.respondToGatewayAdmissionOutcome()) { + return true; + } + cleanupPreaccept(true); + return false; + }; + let replyDispatchRuntime: PreparedReplyDispatchRuntime; + try { + const publishedRuntime = await loadPublishedGatewayReplyDispatchRuntime({ + agentId: params.activeSessionAgentId, + abortSignal: activeRunAbort.controller.signal, + }); + if (!revalidateAdmission()) { + return undefined; + } + if (!publishedRuntime) { + throw new Error(`published reply runtime missing for ${params.activeSessionAgentId}`); + } + replyDispatchRuntime = publishedRuntime; + preparedModelRuntimeLease = await acquireAgentRunPreparedModelRuntime( + { + config: replyDispatchRuntime.config, + agentId: replyDispatchRuntime.agentId, + agentDir: replyDispatchRuntime.agentDir, + allowGatewaySubagentBinding: true, + workspaceDir: workspaceOverride ?? replyDispatchRuntime.workspaceDir, + }, + { + catalogMode: "static", + pluginGeneration: replyDispatchRuntime.pluginGeneration, + abortSignal: activeRunAbort.controller.signal, + }, + ); + if (!revalidateAdmission()) { + return undefined; + } + } catch (err) { + if (!revalidateAdmission()) { + return undefined; + } + return rejectPreaccept(errorShapeFromError(ErrorCodes.UNAVAILABLE, err)); + } + const resolvedThreadId = params.delivery.explicitThreadId ?? params.delivery.deliveryPlan.resolvedThreadId; const completionEvent = resolveExactSubagentCompletionEvent({ @@ -367,23 +456,21 @@ export async function prepareAgentRunDispatch(params: { pluginId: normalizeOptionalString(params.client?.internal?.pluginRuntimeOwnerId), gatewayContextResolver: params.context.resolveGatewayContext, }); + if (!revalidateAdmission()) { + return undefined; + } } catch (err) { params.context.logGateway.warn( `failed to register plugin subagent run ${params.runId}; rejecting untracked dispatch: ${formatForLog(err)}`, ); - activeRunAbort.cleanup({ force: true }); - activeGatewayWorkAdmission.release(); - params.io.emitAcceptance([ - false, - undefined, + return rejectPreaccept( errorShapeFromError( ErrorCodes.UNAVAILABLE, new Error("plugin subagent registry persistence failed; run was not started", { cause: err, }), ), - ]); - return undefined; + ); } } let restoreAdmittedRestartRecoveryInterrupted: @@ -392,14 +479,9 @@ export async function prepareAgentRunDispatch(params: { if (params.isRestartRecoveryResumeRun) { const recoverySessionKey = params.resolvedSessionKey; if (!recoverySessionKey) { - activeRunAbort.cleanup({ force: true }); - activeGatewayWorkAdmission.release(); - params.io.emitAcceptance([ - false, - undefined, + return rejectPreaccept( errorShape(ErrorCodes.UNAVAILABLE, "restart recovery session target is unavailable"), - ]); - return undefined; + ); } try { const recoveryAdmission = await commitMainSessionRecovery({ @@ -413,6 +495,9 @@ export async function prepareAgentRunDispatch(params: { requireWriteSuccess: true, target: { sessionKey: recoverySessionKey, storePath: lifecycleStorePath }, }); + if (!revalidateAdmission()) { + return undefined; + } if (recoveryAdmission.transition.kind !== "admitted_recovery") { throw new Error( `Session "${recoverySessionKey}" restart recovery reservation is stale; recovery was skipped.`, @@ -449,14 +534,7 @@ export async function prepareAgentRunDispatch(params: { : undefined; }; } catch (err) { - activeRunAbort.cleanup({ force: true }); - activeGatewayWorkAdmission.release(); - params.io.emitAcceptance([ - false, - undefined, - errorShape(ErrorCodes.UNAVAILABLE, formatForLog(err)), - ]); - return undefined; + return rejectPreaccept(errorShape(ErrorCodes.UNAVAILABLE, formatForLog(err))); } } let userTurn: PreparedAgentRunUserTurn; @@ -491,29 +569,10 @@ export async function prepareAgentRunDispatch(params: { params.onUserTurnMediaPersisted(); } } catch (err) { - activeRunAbort.cleanup({ force: true }); - activeGatewayWorkAdmission.release(); - params.io.emitAcceptance([false, undefined, errorShapeFromError(ErrorCodes.UNAVAILABLE, err)]); - return undefined; + return rejectPreaccept(errorShapeFromError(ErrorCodes.UNAVAILABLE, err)); } - try { - // Transcript persistence can yield. Revalidate the exact live admission - // before its durable turn is allowed to cross the acceptance boundary. - params.assertGatewayWorkAdmissionAllowed(); - } catch (err) { + if (!revalidateAdmission()) { releasePreparedAgentRunUserTurn(userTurn); - activeRunAbort.cleanup({ force: true }); - activeGatewayWorkAdmission.release(); - params.io.emitAcceptance([ - false, - undefined, - errorShapeFromError(ErrorCodes.INVALID_REQUEST, err), - ]); - return undefined; - } - if (params.respondToGatewayAdmissionOutcome()) { - releasePreparedAgentRunUserTurn(userTurn); - activeRunAbort.cleanup({ force: true }); return undefined; } const accepted = { @@ -566,8 +625,11 @@ export async function prepareAgentRunDispatch(params: { lifecycleStorePath, resolvedThreadId, dispatchTaskTrackingMode, + preparedModelRuntimeLease, + replyDispatchRuntime, unpersistedOffloadedRefs: userTurn.recorder ? [] : params.offloadedRefs, userTurn, + workspaceOverride, restoreAdmittedRestartRecoveryInterrupted, }; } diff --git a/src/gateway/agent-turn/agent-run-execution-phase.owner.test.ts b/src/gateway/agent-turn/agent-run-execution-phase.owner.test.ts index c77cc7fb142c..58cbd6b02150 100644 --- a/src/gateway/agent-turn/agent-run-execution-phase.owner.test.ts +++ b/src/gateway/agent-turn/agent-run-execution-phase.owner.test.ts @@ -1,47 +1,55 @@ -import { describe, expect, it, vi } from "vitest"; +import { beforeEach, describe, expect, it, vi } from "vitest"; +import { + getPreparedModelRuntimeBorrowedSnapshot, + getPreparedModelRuntimePluginGeneration, +} from "../../agents/prepared-model-runtime-generation-scope.js"; import { startAgentRunExecution } from "./agent-run-execution-phase.js"; const dispatchAgentRunFromGateway = vi.hoisted(() => vi.fn()); -vi.mock("../../agents/prepared-model-runtime.js", () => ({ - loadPublishedGatewayReplyDispatchRuntime: async () => ({ - config: {}, - pluginGeneration: "test", - }), -})); - vi.mock("./agent-run-dispatch.js", () => ({ dispatchAgentRunFromGateway, resolveAbortedAgentStopReason: () => "rpc", })); -describe("startAgentRunExecution Gateway ownership", () => { - it("rejects a retired owner after preparation and before final dispatch", async () => { - const cleanup = vi.fn(); - const release = vi.fn(); - let resolveFinal!: () => void; - const final = new Promise((resolve) => { - resolveFinal = resolve; - }); - - startAgentRunExecution({ - assertContextCurrent: () => { - throw new Error("Gateway owner retired"); - }, +function createExecution(options: { aborted?: boolean; assertContextCurrent?: () => void } = {}) { + const abortCleanup = vi.fn(); + const gatewayRelease = vi.fn(); + let resolveRuntimeReleased!: () => void; + const runtimeReleased = new Promise((resolve) => { + resolveRuntimeReleased = resolve; + }); + const runtimeRelease = vi.fn(resolveRuntimeReleased); + const controller = new AbortController(); + if (options.aborted) { + controller.abort(); + } + return { + abortCleanup, + gatewayRelease, + runtimeRelease, + runtimeReleased, + params: { + assertContextCurrent: options.assertContextCurrent, prepared: { activeGatewayWorkAdmission: { - release, + release: gatewayRelease, run: async (run: () => Promise) => await run(), }, activeRunAbort: { - cleanup, - controller: new AbortController(), + cleanup: abortCleanup, + controller, registered: false, }, dispatchTaskTrackingMode: "none", effectiveAllowModelOverride: false, lifecycleStorePath: "", operationalRunInstance: {}, + preparedModelRuntimeLease: { release: runtimeRelease, snapshot: {} }, + replyDispatchRuntime: { + config: { runtime: "A" }, + pluginGeneration: "generation-A", + }, unpersistedOffloadedRefs: [], userTurn: { execApprovalFollowupHandoffClaimId: "claim", @@ -49,6 +57,7 @@ describe("startAgentRunExecution Gateway ownership", () => { senderIsOwner: false, suppressPromptPersistence: false, }, + workspaceOverride: "/workspace/A", }, request: {}, cfg: {}, @@ -62,7 +71,7 @@ describe("startAgentRunExecution Gateway ownership", () => { images: [], imageOrder: [], media: [], - runId: "owner-retired", + runId: "owner-test", agentDedupeKeys: [], bestEffortDeliver: false, lifecycleGeneration: "test", @@ -77,14 +86,87 @@ describe("startAgentRunExecution Gateway ownership", () => { }, io: { emitAcceptance: vi.fn(), - emitFinal: () => resolveFinal(), + emitFinal: vi.fn(), }, releaseCronContinuationClaimWithRecovery: async () => true, - } as never); + } as unknown as Parameters[0], + }; +} - await final; - await vi.waitFor(() => expect(cleanup).toHaveBeenCalledOnce()); +describe("startAgentRunExecution Gateway ownership", () => { + beforeEach(() => dispatchAgentRunFromGateway.mockReset()); + + it("dispatches with the runtime generation frozen at admission", async () => { + const execution = createExecution(); + let resolveDispatched!: () => void; + const dispatched = new Promise((resolve) => { + resolveDispatched = resolve; + }); + let resolveCleanupObserved!: () => void; + const cleanupObserved = new Promise((resolve) => { + resolveCleanupObserved = resolve; + }); + let borrowedAfterCleanup: Promise | undefined; + let dispatchedGeneration: unknown; + let dispatchedSnapshot: unknown; + dispatchAgentRunFromGateway.mockImplementationOnce(() => { + const generation = execution.params.prepared.replyDispatchRuntime.pluginGeneration; + dispatchedGeneration = getPreparedModelRuntimePluginGeneration(); + dispatchedSnapshot = getPreparedModelRuntimeBorrowedSnapshot(generation); + borrowedAfterCleanup = (async () => { + await cleanupObserved; + return getPreparedModelRuntimeBorrowedSnapshot(generation); + })(); + resolveDispatched(); + }); + + startAgentRunExecution(execution.params); + + await dispatched; + expect(dispatchedGeneration).toBe( + execution.params.prepared.replyDispatchRuntime.pluginGeneration, + ); + expect(dispatchedSnapshot).toBe(execution.params.prepared.preparedModelRuntimeLease.snapshot); + const dispatch = dispatchAgentRunFromGateway.mock.calls[0]?.[0]; + expect(dispatch?.commandRuntimeContext).toEqual({ + config: { runtime: "A" }, + pluginGeneration: "generation-A", + }); + expect(dispatch?.ingressOpts.workspaceDir).toBe("/workspace/A"); + expect(execution.runtimeRelease).not.toHaveBeenCalled(); + + dispatch?.cleanupAbortController(); + dispatch?.cleanupAbortController(); + resolveCleanupObserved(); + await expect(borrowedAfterCleanup).resolves.toBeUndefined(); + expect(execution.runtimeRelease).toHaveBeenCalledOnce(); + }); + + it("releases the admitted runtime once when aborted before dispatch", async () => { + const execution = createExecution({ aborted: true }); + + startAgentRunExecution(execution.params); + + await execution.runtimeReleased; expect(dispatchAgentRunFromGateway).not.toHaveBeenCalled(); - expect(release).toHaveBeenCalledOnce(); + expect(execution.abortCleanup).toHaveBeenCalledOnce(); + expect(execution.gatewayRelease).toHaveBeenCalledOnce(); + expect(execution.runtimeRelease).toHaveBeenCalledOnce(); + }); + + it("releases the admitted runtime once when its owner retires before dispatch", async () => { + const execution = createExecution({ + assertContextCurrent: () => { + throw new Error("Gateway owner retired"); + }, + }); + + startAgentRunExecution(execution.params); + + await execution.runtimeReleased; + expect(dispatchAgentRunFromGateway).not.toHaveBeenCalled(); + expect(execution.abortCleanup).toHaveBeenCalledOnce(); + expect(execution.gatewayRelease).toHaveBeenCalledOnce(); + expect(execution.runtimeRelease).toHaveBeenCalledOnce(); }); }); diff --git a/src/gateway/agent-turn/agent-run-execution-phase.ts b/src/gateway/agent-turn/agent-run-execution-phase.ts index 99cb4098ac14..be81a6a8ed14 100644 --- a/src/gateway/agent-turn/agent-run-execution-phase.ts +++ b/src/gateway/agent-turn/agent-run-execution-phase.ts @@ -13,9 +13,8 @@ import { type MainSessionRecoveryPendingTarget, type MainSessionRecoveryOwnerLease, } from "../../agents/main-session-recovery/main-session-recovery-store.js"; -import { loadPublishedGatewayReplyDispatchRuntime } from "../../agents/prepared-model-runtime.js"; +import { withPreparedModelRuntimePluginGenerationScope } from "../../agents/prepared-model-runtime-generation-scope.js"; import { resolveScheduledToolPolicyContext } from "../../agents/scheduled-tool-policy.js"; -import { resolveIngressWorkspaceOverrideForSessionRun } from "../../agents/spawned-context.js"; import { isExecutionIdentityCollectionEnabled } from "../../audit/audit-config.js"; import { setChannelSourceTurnId, @@ -112,16 +111,27 @@ export function startAgentRunExecution(params: { }): void { const { prepared } = params; let unpersistedOffloadedRefs = prepared.unpersistedOffloadedRefs; + let preparedModelRuntimeLease: typeof prepared.preparedModelRuntimeLease | undefined = + prepared.preparedModelRuntimeLease; let releaseGatewayRootContinuation = retainGatewayRootWorkAdmissionContinuation() ?? undefined; const cleanupAdmittedRun: typeof prepared.activeRunAbort.cleanup = (options) => { const refsToDiscard = unpersistedOffloadedRefs; unpersistedOffloadedRefs = []; prepared.activeRunAbort.cleanup(options); prepared.activeGatewayWorkAdmission.release(); + const runtimeLease = preparedModelRuntimeLease; + preparedModelRuntimeLease = undefined; + runtimeLease?.release(); releaseGatewayRootContinuation?.(); releaseGatewayRootContinuation = undefined; void discardPreparedInboundMedia(refsToDiscard, params.context.logGateway); }; + const dispatchAdmittedAgentRun = (dispatch: Parameters[0]) => + withPreparedModelRuntimePluginGenerationScope( + prepared.replyDispatchRuntime.pluginGeneration, + () => dispatchAgentRunFromGateway(dispatch), + () => preparedModelRuntimeLease?.snapshot, + ); void prepared.activeGatewayWorkAdmission.run(async () => { await yieldAfterAgentAcceptedAck(); let dispatched = false; @@ -212,15 +222,6 @@ export function startAgentRunExecution(params: { const ingressAgentId = params.resolvedSessionKey ? params.activeSessionAgentId : params.agentId; - const replyDispatchRuntime = await loadPublishedGatewayReplyDispatchRuntime({ - agentId: params.activeSessionAgentId, - abortSignal: prepared.activeRunAbort.controller.signal, - }); - if (!replyDispatchRuntime?.pluginGeneration) { - throw new Error( - `prepared reply dispatch runtime was not published for ${params.activeSessionAgentId}`, - ); - } // Plugin-owned additive grants stay internal to the authenticated in-process run. // Public agent params cannot supply them, and normal tool policy still filters them. const runtimePluginToolGrant = @@ -280,16 +281,15 @@ export function startAgentRunExecution(params: { } else if (localUserIngress) { attachAgentCommandAdmissionFacts(runContext, localUserIngress.facts); } - // Routing and runtime publication await after admission. Retired owners - // must fail before the prepared user turn becomes an agent run. + // Awaited routing can retire this owner before final dispatch. params.assertContextCurrent?.(); finalizePreparedAgentRunUserTurn(prepared.userTurn); - dispatchAgentRunFromGateway( + dispatchAdmittedAgentRun( withAgentRunDispatchExecutionIdentity( { commandRuntimeContext: { - config: replyDispatchRuntime.config, - pluginGeneration: replyDispatchRuntime.pluginGeneration, + config: prepared.replyDispatchRuntime.config, + pluginGeneration: prepared.replyDispatchRuntime.pluginGeneration, }, cronCreatorAuthority: prepared.cronCreatorAuthority, ingressOpts: { @@ -425,11 +425,7 @@ export function startAgentRunExecution(params: { prepared.activeRunAbort.entry.sessionId = sessionId; } }, - workspaceDir: resolveIngressWorkspaceOverrideForSessionRun({ - spawnedBy: params.spawnedBy, - workspaceDir: params.sessionEntry?.spawnedWorkspaceDir, - cwd: params.sessionEntry?.spawnedCwd, - }), + workspaceDir: prepared.workspaceOverride, cwd: resolveSessionRuntimeCwd({ requestedCwd: params.request.cwd, sessionEntry: params.sessionEntry, diff --git a/src/gateway/server-methods/agent.create-event.test.ts b/src/gateway/server-methods/agent.create-event.test.ts index 25dc31f7f802..0f6609e92867 100644 --- a/src/gateway/server-methods/agent.create-event.test.ts +++ b/src/gateway/server-methods/agent.create-event.test.ts @@ -41,6 +41,20 @@ vi.mock("../../commands/agent.js", () => ({ agentCommandFromIngress: agentIngressMocks.agentCommandFromIngress, })); +vi.mock("../../agents/prepared-model-runtime.js", () => ({ + acquireAgentRunPreparedModelRuntime: vi.fn(async () => ({ + release: vi.fn(), + snapshot: {}, + })), + loadPublishedGatewayReplyDispatchRuntime: vi.fn(async ({ agentId }: { agentId: string }) => ({ + agentId, + agentDir: configMocks.workspaceDir, + config: configMocks.getRuntimeConfig(), + pluginGeneration: { pluginMetadataSnapshot: {} }, + workspaceDir: configMocks.workspaceDir, + })), +})); + vi.mock("../../runtime.js", () => ({ defaultRuntime: {}, })); diff --git a/src/gateway/server-methods/agent.test-harness.ts b/src/gateway/server-methods/agent.test-harness.ts index c86da2c4e232..fa193ff008a7 100644 --- a/src/gateway/server-methods/agent.test-harness.ts +++ b/src/gateway/server-methods/agent.test-harness.ts @@ -193,10 +193,16 @@ vi.mock("../../commands/agent.js", () => { vi.mock("../../agents/prepared-model-runtime.js", () => ({ // Direct handler tests bypass Gateway startup, so provide the lifecycle fact // that production publishes before admitting agent RPCs. + acquireAgentRunPreparedModelRuntime: vi.fn(async () => ({ + release: vi.fn(), + snapshot: {}, + })), loadPublishedGatewayReplyDispatchRuntime: async ({ agentId }: { agentId: string }) => ({ agentId, + agentDir: "/tmp/agent", config: resolveAgentTestConfig(), pluginGeneration: { pluginMetadataSnapshot: {} }, + workspaceDir: "/tmp/workspace", }), })); diff --git a/src/gateway/server.agent.gateway-server-agent-auth-refresh.test.ts b/src/gateway/server.agent.gateway-server-agent-auth-refresh.test.ts index 6a32a70efd3e..a8c9894ac3ee 100644 --- a/src/gateway/server.agent.gateway-server-agent-auth-refresh.test.ts +++ b/src/gateway/server.agent.gateway-server-agent-auth-refresh.test.ts @@ -70,6 +70,31 @@ function sendAgentRpc(socket: WebSocket, params: { agentId: string; runId: strin return { accepted, final }; } +function sendPreacceptAgentRpc(socket: WebSocket, params: { agentId: string; runId: string }) { + const response = onceMessage( + socket, + (frame) => frame.type === "res" && frame.id === params.runId, + ); + const final = onceMessage( + socket, + (frame) => + frame.type === "res" && frame.id === params.runId && frame.payload?.status !== "accepted", + ); + socket.send( + JSON.stringify({ + type: "req", + id: params.runId, + method: "agent", + params: { + agentId: params.agentId, + message: `dispatch ${params.runId}`, + idempotencyKey: params.runId, + }, + }), + ); + return { response, final }; +} + function agentCommandCallsFor(runId: string) { return vi .mocked(agentCommandMock) @@ -103,6 +128,71 @@ describe("gateway agent auth refresh dispatch", () => { testState.agentsConfig = undefined; }); + test("keeps an accepted run on its admitted runtime generation", async () => { + const affectedAgentId = "auth-pinned"; + const admittedRunId = "idem-agent-auth-admitted"; + const subsequentRunId = "idem-agent-auth-next"; + const before = await prepareAuthDispatchAgents(affectedAgentId); + const published = createDeferred(); + const unregister = registerPreparedModelRuntimePublicationListener((event) => { + if (event.phase === "published") { + published.resolve(); + } + }); + try { + const admitted = sendAgentRpc(gatewaySuite.ws, { + agentId: affectedAgentId, + runId: admittedRunId, + }); + await admitted.accepted; + expect(agentCommandCallsFor(admittedRunId)).toHaveLength(0); + + setRuntimeAuthProfileStoreSnapshot( + { + version: 1, + profiles: { + "anthropic:default": { + type: "api_key", + provider: "anthropic", + key: "next-generation-key", + }, + }, + }, + before.agentDir, + ); + await published.promise; + const after = await loadPublishedGatewayReplyDispatchRuntime({ + agentId: affectedAgentId, + }); + expect(after).not.toBe(before.runtime); + + await expect(admitted.final).resolves.toMatchObject({ + ok: true, + payload: { status: "ok" }, + }); + expect(agentCommandCallsFor(admittedRunId)[0]?.[4]).toMatchObject({ + config: before.runtime?.config, + pluginGeneration: before.runtime?.pluginGeneration, + }); + + const subsequent = sendAgentRpc(gatewaySuite.ws, { + agentId: affectedAgentId, + runId: subsequentRunId, + }); + await subsequent.accepted; + await expect(subsequent.final).resolves.toMatchObject({ + ok: true, + payload: { status: "ok" }, + }); + expect(agentCommandCallsFor(subsequentRunId)[0]?.[4]).toMatchObject({ + config: after?.config, + pluginGeneration: after?.pluginGeneration, + }); + } finally { + unregister(); + } + }); + test("aborts one affected waiter without cancelling shared auth publication", async () => { const affectedAgentId = "auth-wait"; const abortedRunId = "idem-agent-auth-aborted"; @@ -142,15 +232,14 @@ describe("gateway agent auth refresh dispatch", () => { before.agentDir, ); - const aborted = sendAgentRpc(gatewaySuite.ws, { + const aborted = sendPreacceptAgentRpc(gatewaySuite.ws, { agentId: affectedAgentId, runId: abortedRunId, }); - const waiting = sendAgentRpc(gatewaySuite.ws, { + const waiting = sendPreacceptAgentRpc(gatewaySuite.ws, { agentId: affectedAgentId, runId: waitingRunId, }); - await Promise.all([aborted.accepted, waiting.accepted]); const sibling = sendAgentRpc(gatewaySuite.ws, { agentId: "main", runId: siblingRunId }); await sibling.accepted; await expect(sibling.final).resolves.toMatchObject({ ok: true, payload: { status: "ok" } }); @@ -168,7 +257,7 @@ describe("gateway agent auth refresh dispatch", () => { payload: { aborted: true, runIds: [abortedRunId] }, }); await vi.waitFor(() => expect(getActiveGatewayRootWorkCount()).toBe(activeWorkBefore + 1)); - await expect(aborted.final).resolves.toMatchObject({ + await expect(aborted.response).resolves.toMatchObject({ ok: true, payload: { status: "timeout", @@ -178,13 +267,17 @@ describe("gateway agent auth refresh dispatch", () => { }, }); await expect( - Promise.race([waiting.final.then(() => "settled"), Promise.resolve("pending")]), + Promise.race([waiting.response.then(() => "settled"), Promise.resolve("pending")]), ).resolves.toBe("pending"); publicationGate.resolve({ agentDir: before.agentDir, wrote: false }); await published.promise; const after = await loadPublishedGatewayReplyDispatchRuntime({ agentId: affectedAgentId }); expect(after).not.toBe(before.runtime); + await expect(waiting.response).resolves.toMatchObject({ + ok: true, + payload: { status: "accepted" }, + }); await expect(waiting.final).resolves.toMatchObject({ ok: true, payload: { status: "ok" }, @@ -253,14 +346,12 @@ describe("gateway agent auth refresh dispatch", () => { `prepared reply dispatch runtime owner was not published for ${affectedAgentId}`, ); - const dispatched = sendAgentRpc(gatewaySuite.ws, { agentId: affectedAgentId, runId }); - await expect(dispatched.accepted).resolves.toMatchObject({ - ok: true, - payload: { status: "accepted" }, + const rejected = sendPreacceptAgentRpc(gatewaySuite.ws, { + agentId: affectedAgentId, + runId, }); - await expect(dispatched.final).resolves.toMatchObject({ + await expect(rejected.response).resolves.toMatchObject({ ok: false, - payload: { status: "error" }, error: { code: "UNAVAILABLE", message: expect.stringContaining(