diff --git a/packages/runtime-host/src/__tests__/client-capability-admission-integration.test.ts b/packages/runtime-host/src/__tests__/client-capability-admission-integration.test.ts index a68e1d6d71..4191e67b69 100644 --- a/packages/runtime-host/src/__tests__/client-capability-admission-integration.test.ts +++ b/packages/runtime-host/src/__tests__/client-capability-admission-integration.test.ts @@ -355,7 +355,8 @@ function createInteractionCoordinator( preflightSessionSnapshot: () => true, refreshCanonicalContinuity: async () => undefined, onPoison: () => undefined, - onSandboxBoundarySettled: async () => undefined, + resolveSandboxBoundaryRootSession: async () => undefined, + onSandboxBoundaryGraphWake: async () => undefined, }; return new HostInteractionCoordinator(options); } diff --git a/packages/runtime-host/src/__tests__/goal-root-authority.test.ts b/packages/runtime-host/src/__tests__/goal-root-authority.test.ts index c6e89e8376..98acd01f34 100644 --- a/packages/runtime-host/src/__tests__/goal-root-authority.test.ts +++ b/packages/runtime-host/src/__tests__/goal-root-authority.test.ts @@ -617,7 +617,8 @@ async function createFixture(options: { recoverAdmissions?: boolean } = {}): Pro onPoison: () => { requestedDrain = true; }, - onSandboxBoundarySettled: async () => {}, + resolveSandboxBoundaryRootSession: async () => undefined, + onSandboxBoundaryGraphWake: async () => {}, }); const backends = new BackendRegistry(); backends.register('ai-sdk', (context) => new FakeBackend(context)); diff --git a/packages/runtime-host/src/__tests__/interaction-coordinator.test.ts b/packages/runtime-host/src/__tests__/interaction-coordinator.test.ts index 52ebb394de..eac022b0b7 100644 --- a/packages/runtime-host/src/__tests__/interaction-coordinator.test.ts +++ b/packages/runtime-host/src/__tests__/interaction-coordinator.test.ts @@ -17,7 +17,7 @@ * under the License. */ -import { deferred } from '@maka/core/test-only/async-primitives'; +import { deferred, withTimeout } from '@maka/core/test-only/async-primitives'; import assert from 'node:assert/strict'; import { mkdir, mkdtemp, rm } from 'node:fs/promises'; import { tmpdir } from 'node:os'; @@ -294,7 +294,8 @@ describe('HostInteractionCoordinator', () => { return true; }, refreshCanonicalContinuity: async () => {}, - onSandboxBoundarySettled: async (sessionId) => { + resolveSandboxBoundaryRootSession: async (sessionId) => sessionId, + onSandboxBoundaryGraphWake: async (sessionId) => { assert.equal(sessionId, session.id); graphWakes += 1; }, @@ -388,6 +389,198 @@ describe('HostInteractionCoordinator', () => { }); }); + test('does not hold Session admission while graph wake reconciliation waits', async () => { + await withStore(async ({ owner, store, stores }) => { + const workspace = join(owner.capability.canonicalPath, 'wake-workspace'); + await mkdir(workspace); + const session = await stores.sessionStore.create({ + cwd: workspace, + llmConnectionId: 'cccccccc-cccc-4ccc-8ccc-cccccccccccc', + llmConnectionSlug: 'fake', + model: 'fake-model', + permissionMode: 'ask', + }); + const identity = { ...RUN, sessionId: session.id }; + const wakeStarted = deferred(); + const releaseWake = deferred(); + const wakeFinished = deferred(); + const resolverStarted = deferred(); + const releaseResolver = deferred(); + let resolvedRootSessionId: string | undefined; + let wakeNotificationStarted = false; + const coordinator = new HostInteractionCoordinator({ + store, + sandboxBoundaries: stores.sessionStore, + sessionAdmission: new SessionAdmissionGate(), + sessions: stores.sessionStore, + preflightSessionSnapshot: () => true, + refreshCanonicalContinuity: async () => {}, + resolveSandboxBoundaryRootSession: async (sessionId) => { + resolvedRootSessionId = sessionId; + resolverStarted.resolve(); + await releaseResolver.promise; + return session.id; + }, + onSandboxBoundaryGraphWake: async (rootSessionId) => { + assert.equal(rootSessionId, session.id); + wakeNotificationStarted = true; + wakeStarted.resolve(); + await releaseWake.promise; + wakeFinished.resolve(); + }, + onPoison: () => {}, + }); + const binding = coordinator.bindRun(identity); + const request = sandboxBoundaryEvent({ + sessionId: session.id, + requestId: 'boundary_wake_wait', + status: 'pending', + baseRevision: 0, + turnId: identity.turnId, + runId: identity.runId, + expansion: { network: { enabled: true } }, + justification: 'Connect to the requested service.', + createdAt: 1, + }); + await binding.acceptSandboxBoundaryRequest({ + request, + continuation: sandboxBoundaryContinuation(identity, request.requestId), + }); + + let answerSettled = false; + let answerResult: + | Awaited> + | undefined; + const answer = coordinator.handlers['interaction.answer']( + { + sessionId: session.id, + interactionId: request.requestId, + answer: { kind: 'sandbox_boundary', decision: 'allow' }, + }, + connection(), + ); + void answer.then( + (result) => { + answerResult = result; + answerSettled = true; + }, + () => { + answerSettled = true; + }, + ); + let querySettled = false; + let query: ReturnType<(typeof coordinator.handlers)['interaction.query']> | undefined; + try { + await withTimeout( + resolverStarted.promise, + 5_000, + 'sandbox boundary root-session resolver did not start', + ); + query = coordinator.handlers['interaction.query']( + { sessionId: session.id, interactionId: request.requestId }, + connection(), + ); + void query.then( + () => { + querySettled = true; + }, + () => { + querySettled = true; + }, + ); + await new Promise((resolve) => setImmediate(resolve)); + assert.equal( + querySettled, + false, + 'interaction query bypassed the resolver admission lease', + ); + releaseResolver.resolve(); + await answer; + assert.ok(query); + const queryResult = await query; + assert.equal(queryResult.ok, true); + assert.equal(querySettled, true); + await withTimeout(wakeStarted.promise, 5_000, 'sandbox boundary graph wake did not start'); + assert.equal( + answerSettled, + true, + 'interaction answer waited for graph wake reconciliation', + ); + assert.equal(answerResult?.ok, true); + if (answerResult?.ok) assert.equal(answerResult.result.status, 'answered'); + assert.equal(resolvedRootSessionId, session.id); + } finally { + releaseResolver.resolve(); + releaseWake.resolve(); + if (wakeNotificationStarted) await wakeFinished.promise; + await answer.catch(() => undefined); + } + await binding.close('turn_terminal'); + binding.release(); + await coordinator.close(); + }); + }); + + test('poisons when detached graph wake notification rejects', async () => { + await withStore(async ({ owner, store, stores }) => { + const workspace = join(owner.capability.canonicalPath, 'wake-rejection-workspace'); + await mkdir(workspace); + const session = await stores.sessionStore.create({ + cwd: workspace, + llmConnectionId: 'cccccccc-cccc-4ccc-8ccc-cccccccccccc', + llmConnectionSlug: 'fake', + model: 'fake-model', + permissionMode: 'ask', + }); + const identity = { ...RUN, sessionId: session.id }; + const poison: RuntimeInteractionFailStopError[] = []; + const coordinator = new HostInteractionCoordinator({ + store, + sandboxBoundaries: stores.sessionStore, + sessionAdmission: new SessionAdmissionGate(), + sessions: stores.sessionStore, + preflightSessionSnapshot: () => true, + refreshCanonicalContinuity: async () => {}, + resolveSandboxBoundaryRootSession: async () => session.id, + onSandboxBoundaryGraphWake: async () => { + throw new Error('graph wake notification failed'); + }, + onPoison: (error) => poison.push(error), + }); + const binding = coordinator.bindRun(identity); + const request = sandboxBoundaryEvent({ + sessionId: session.id, + requestId: 'boundary_wake_rejection', + status: 'pending', + baseRevision: 0, + turnId: identity.turnId, + runId: identity.runId, + expansion: { network: { enabled: true } }, + justification: 'Connect to the requested service.', + createdAt: 1, + }); + await binding.acceptSandboxBoundaryRequest({ + request, + continuation: sandboxBoundaryContinuation(identity, request.requestId), + }); + + const answerResult = await coordinator.handlers['interaction.answer']( + { + sessionId: session.id, + interactionId: request.requestId, + answer: { kind: 'sandbox_boundary', decision: 'allow' }, + }, + connection(), + ); + assert.equal(answerResult.ok, true); + await new Promise((resolve) => setImmediate(resolve)); + assert.equal(poison.length, 1); + assert.equal(coordinator.isPoisoned(), true); + await assert.rejects(binding.close('turn_terminal'), poison[0]); + await assert.rejects(coordinator.close(), poison[0]); + }); + }); + test('a queued stop waits for sandbox boundary publication before closing its Run', async () => { await withStore(async ({ owner, store, stores }) => { const workspace = join(owner.capability.canonicalPath, 'publication-workspace'); @@ -417,7 +610,8 @@ describe('HostInteractionCoordinator', () => { await releaseAdmissionRefresh.promise; }, onPoison: () => {}, - onSandboxBoundarySettled: async () => {}, + resolveSandboxBoundaryRootSession: async () => undefined, + onSandboxBoundaryGraphWake: async () => {}, }); const binding = await bindRuntimeInteractionRun(coordinator, identity); const request = sandboxBoundaryEvent({ @@ -501,7 +695,8 @@ describe('HostInteractionCoordinator', () => { preflightSessionSnapshot: () => false, refreshCanonicalContinuity: async () => {}, onPoison: () => {}, - onSandboxBoundarySettled: async () => {}, + resolveSandboxBoundaryRootSession: async () => undefined, + onSandboxBoundaryGraphWake: async () => {}, }); const ownerRun = coordinator.bindRun(identity); @@ -845,7 +1040,8 @@ function createCoordinator( preflightSessionSnapshot: () => true, refreshCanonicalContinuity: async () => {}, onPoison: () => {}, - onSandboxBoundarySettled: async () => {}, + resolveSandboxBoundaryRootSession: async () => undefined, + onSandboxBoundaryGraphWake: async () => {}, ...overrides, }); } diff --git a/packages/runtime-host/src/__tests__/root-turn-coordinator.test.ts b/packages/runtime-host/src/__tests__/root-turn-coordinator.test.ts index 553d5d1d4b..535a25837c 100644 --- a/packages/runtime-host/src/__tests__/root-turn-coordinator.test.ts +++ b/packages/runtime-host/src/__tests__/root-turn-coordinator.test.ts @@ -2717,7 +2717,8 @@ test('hosted linked child roots share admission, message, terminal, and stop aut onPoison: () => { drainRequested = true; }, - onSandboxBoundarySettled: async () => {}, + resolveSandboxBoundaryRootSession: async () => undefined, + onSandboxBoundaryGraphWake: async () => {}, }); const interactionAuthority: RuntimeInteractionAuthority = { bindRun: (identity) => { @@ -5286,7 +5287,8 @@ async function createFailureFixture(options: { refreshCanonicalContinuity: (sessionId, admission) => requireContinuity(continuity).refreshCanonical(sessionId, admission), onPoison: requestDrain, - onSandboxBoundarySettled: async () => {}, + resolveSandboxBoundaryRootSession: async () => undefined, + onSandboxBoundaryGraphWake: async () => {}, }) : undefined; const backends = new BackendRegistry(); diff --git a/packages/runtime-host/src/__tests__/sandbox-boundary-graph-wake.test.ts b/packages/runtime-host/src/__tests__/sandbox-boundary-graph-wake.test.ts index 20b56508e0..61269a9c14 100644 --- a/packages/runtime-host/src/__tests__/sandbox-boundary-graph-wake.test.ts +++ b/packages/runtime-host/src/__tests__/sandbox-boundary-graph-wake.test.ts @@ -21,7 +21,7 @@ import assert from 'node:assert/strict'; import { test } from 'node:test'; import { agentGraphIdForRootSession } from '@maka/runtime/stream-graph-coordinator'; import { - notifySandboxBoundaryGraphWake, + resolveSandboxBoundaryRootSession, sandboxBoundaryGraphWakeRoot, } from '../server/sandbox-boundary-graph-wake.js'; @@ -78,9 +78,8 @@ test('rejects graph operator lineage that is not owned by its parent Session', a ); }); -test('reads durable operator lineage before notifying only its root graph', async () => { +test('resolves durable operator lineage for only its root graph', async () => { const reads: string[] = []; - const wakes: string[] = []; const headers = new Map([ [ 'graph-operator', @@ -106,16 +105,17 @@ test('reads durable operator lineage before notifying only its root graph', asyn return header; }, }; - const notify = async (sessionId: string) => { - wakes.push(sessionId); - }; - const graphIds = idsFor('root-session'); - await notifySandboxBoundaryGraphWake('graph-operator', reader, graphIds, notify); - await notifySandboxBoundaryGraphWake('ordinary-child', reader, graphIds, notify); + assert.equal( + await resolveSandboxBoundaryRootSession('graph-operator', reader, graphIds), + 'root-session', + ); + assert.equal( + await resolveSandboxBoundaryRootSession('ordinary-child', reader, graphIds), + undefined, + ); assert.deepEqual(reads, ['graph-operator', 'ordinary-child']); - assert.deepEqual(wakes, ['root-session']); }); function idsFor(rootSessionId: string) { diff --git a/packages/runtime-host/src/server/execution-composition.ts b/packages/runtime-host/src/server/execution-composition.ts index 2e38c38f8d..6d83056a07 100644 --- a/packages/runtime-host/src/server/execution-composition.ts +++ b/packages/runtime-host/src/server/execution-composition.ts @@ -164,7 +164,7 @@ import { HostPluginPlatform } from './plugin-platform.js'; import { RootAdmissionOwner } from './root-admission-owner.js'; import { RootTurnCoordinator } from './root-turn-coordinator.js'; import { RuntimePolicyActivationGate } from './runtime-policy-activation-gate.js'; -import { notifySandboxBoundaryGraphWake } from './sandbox-boundary-graph-wake.js'; +import { resolveSandboxBoundaryRootSession } from './sandbox-boundary-graph-wake.js'; import { HostRuntimePolicyCoordinator } from './runtime-policy-coordinator.js'; import { startHostModelMetadataRefresh } from './model-metadata-refresh.js'; import { HostRuntimeResourceCoordinator } from './runtime-resource-coordinator.js'; @@ -680,17 +680,19 @@ export async function createExecutionRuntimeHostComposition( beginDrain(); context.requestDrain(); }, - onSandboxBoundarySettled: (sessionId) => - notifySandboxBoundaryGraphWake( - sessionId, - stores.sessionStore, - { + resolveSandboxBoundaryRootSession: async (sessionId) => { + try { + return await resolveSandboxBoundaryRootSession(sessionId, stores.sessionStore, { listGraphIds: (rootSessionId) => requireGraphCoordinator(graphCoordinator).listGraphIds(rootSessionId), - }, - (rootSessionId) => - requireGraphSupervisorWake(graphSupervisorWake).notifyPermissionResponse(rootSessionId), - ), + }); + } catch (error) { + if (isSessionNotFoundError(error)) return undefined; + throw error; + } + }, + onSandboxBoundaryGraphWake: (rootSessionId) => + requireGraphSupervisorWake(graphSupervisorWake).notifyPermissionResponse(rootSessionId), }); memory = new HostMemoryCoordinator({ store: memoryStore, diff --git a/packages/runtime-host/src/server/interaction-coordinator.ts b/packages/runtime-host/src/server/interaction-coordinator.ts index be899c983f..e3f4451745 100644 --- a/packages/runtime-host/src/server/interaction-coordinator.ts +++ b/packages/runtime-host/src/server/interaction-coordinator.ts @@ -108,7 +108,12 @@ export interface HostInteractionCoordinatorOptions { admission: SessionAdmissionLease, ) => Promise; readonly onPoison: (error: RuntimeInteractionFailStopError) => void; - readonly onSandboxBoundarySettled: (sessionId: string) => Promise | void; + /** Resolve the root Session while the settled Session still holds admission. */ + readonly resolveSandboxBoundaryRootSession: ( + sessionId: string, + ) => Promise | string | undefined; + /** Notify the root graph supervisor after the answer releases admission. */ + readonly onSandboxBoundaryGraphWake: (rootSessionId: string) => Promise | void; } interface RunClosure { @@ -205,7 +210,8 @@ export class HostInteractionCoordinator implements RuntimeInteractionAuthority { readonly #preflightSessionSnapshot: HostInteractionCoordinatorOptions['preflightSessionSnapshot']; readonly #refreshCanonicalContinuity: HostInteractionCoordinatorOptions['refreshCanonicalContinuity']; readonly #onPoison: HostInteractionCoordinatorOptions['onPoison']; - readonly #onSandboxBoundarySettled: HostInteractionCoordinatorOptions['onSandboxBoundarySettled']; + readonly #resolveSandboxBoundaryRootSession: HostInteractionCoordinatorOptions['resolveSandboxBoundaryRootSession']; + readonly #onSandboxBoundaryGraphWake: HostInteractionCoordinatorOptions['onSandboxBoundaryGraphWake']; readonly #runs = new Map(); readonly #live = new Map(); #accepting = true; @@ -220,7 +226,8 @@ export class HostInteractionCoordinator implements RuntimeInteractionAuthority { this.#preflightSessionSnapshot = options.preflightSessionSnapshot; this.#refreshCanonicalContinuity = options.refreshCanonicalContinuity; this.#onPoison = options.onPoison; - this.#onSandboxBoundarySettled = options.onSandboxBoundarySettled; + this.#resolveSandboxBoundaryRootSession = options.resolveSandboxBoundaryRootSession; + this.#onSandboxBoundaryGraphWake = options.onSandboxBoundaryGraphWake; } bindRun(identity: RuntimeInteractionRunIdentity): RuntimeInteractionRunOwner { @@ -886,7 +893,19 @@ export class HostInteractionCoordinator implements RuntimeInteractionAuthority { await this.#refreshCanonicalContinuity(request.sessionId, admission); this.#throwIfPoisoned(); await this.#applySandboxBoundaryDecisionAndDelete(entry, settlement); - await this.#onSandboxBoundarySettled(request.sessionId); + // The answer owns Session admission. Resolve graph-wake lineage while that + // admission is held, then detach only the notification which may need to + // acquire the activity lease held by the wake turn parked on this answer. + // Awaiting that notification here deadlocks the Session (#3328, #3866). + const resolvedRootSessionId = await this.#resolveSandboxBoundaryRootSession(request.sessionId); + void Promise.resolve() + .then(() => { + if (!resolvedRootSessionId) return; + return this.#onSandboxBoundaryGraphWake(resolvedRootSessionId); + }) + .catch((error: unknown) => { + this.#poison(error); + }); const result = projectSandboxBoundaryInteraction(settlement.request); if (result.status !== 'answered') { throw this.#poison( diff --git a/packages/runtime-host/src/server/sandbox-boundary-graph-wake.ts b/packages/runtime-host/src/server/sandbox-boundary-graph-wake.ts index 03505b0c75..99ec41e8d8 100644 --- a/packages/runtime-host/src/server/sandbox-boundary-graph-wake.ts +++ b/packages/runtime-host/src/server/sandbox-boundary-graph-wake.ts @@ -39,14 +39,12 @@ export async function sandboxBoundaryGraphWakeRoot( return parent.parentSessionId; } -/** Resolve durable lineage before notifying the root graph supervisor. */ -export async function notifySandboxBoundaryGraphWake( +/** Resolve the durable lineage for a settled sandbox boundary. */ +export async function resolveSandboxBoundaryRootSession( sessionId: string, sessions: SandboxBoundaryGraphWakeHeaderReader, graphIds: { listGraphIds(rootSessionId: string): Promise }, - notifyPermissionResponse: (rootSessionId: string) => Promise | void, -): Promise { +): Promise { const header = await sessions.readHeaderSnapshot(sessionId); - const rootSessionId = await sandboxBoundaryGraphWakeRoot(header, graphIds); - if (rootSessionId) await notifyPermissionResponse(rootSessionId); + return sandboxBoundaryGraphWakeRoot(header, graphIds); }