Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
155 changes: 155 additions & 0 deletions packages/runtime-host/src/__tests__/interaction-coordinator.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -388,6 +388,161 @@ 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();
let resolvedRootSessionId: string | undefined;
const coordinator = new HostInteractionCoordinator({
store,
sandboxBoundaries: stores.sessionStore,
sessionAdmission: new SessionAdmissionGate(),
sessions: stores.sessionStore,
preflightSessionSnapshot: () => true,
refreshCanonicalContinuity: async () => {},
resolveSandboxBoundaryGraphWake: async (sessionId) => {
resolvedRootSessionId = sessionId;
return session.id;
},
onSandboxBoundarySettled: async (rootSessionId) => {
assert.equal(rootSessionId, session.id);
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<ReturnType<(typeof coordinator.handlers)['interaction.answer']>>
| 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;
},
);
await wakeStarted.promise;
try {
await new Promise<void>((resolve) => setImmediate(resolve));
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 {
releaseWake.resolve();
await wakeFinished.promise;
await answer;
}
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 () => {},
resolveSandboxBoundaryGraphWake: async () => session.id,
onSandboxBoundarySettled: 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<void>((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');
Expand Down
22 changes: 12 additions & 10 deletions packages/runtime-host/src/server/execution-composition.ts
Original file line number Diff line number Diff line change
Expand Up @@ -157,7 +157,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 { resolveSandboxBoundaryGraphWake } 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';
Expand Down Expand Up @@ -672,17 +672,19 @@ export async function createExecutionRuntimeHostComposition(
beginDrain();
context.requestDrain();
},
onSandboxBoundarySettled: (sessionId) =>
notifySandboxBoundaryGraphWake(
sessionId,
stores.sessionStore,
{
resolveSandboxBoundaryGraphWake: async (sessionId) => {
try {
return await resolveSandboxBoundaryGraphWake(sessionId, stores.sessionStore, {
listGraphIds: (rootSessionId) =>
requireGraphCoordinator(graphCoordinator).listGraphIds(rootSessionId),
},
(rootSessionId) =>
requireGraphSupervisorWake(graphSupervisorWake).notifyPermissionResponse(rootSessionId),
),
});
} catch (error) {
if (isSessionNotFoundError(error)) return undefined;
throw error;
}
},
onSandboxBoundarySettled: (sessionId) =>
requireGraphSupervisorWake(graphSupervisorWake).notifyPermissionResponse(sessionId),
});
memory = new HostMemoryCoordinator({
store: memoryStore,
Expand Down
23 changes: 22 additions & 1 deletion packages/runtime-host/src/server/interaction-coordinator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,11 @@ export interface HostInteractionCoordinatorOptions {
admission: SessionAdmissionLease,
) => Promise<void>;
readonly onPoison: (error: RuntimeInteractionFailStopError) => void;
/** Resolve graph-wake lineage while the settled Session still holds admission. */
readonly resolveSandboxBoundaryGraphWake?: (
sessionId: string,
admission: SessionAdmissionLease,
) => Promise<string | undefined> | string | undefined;
readonly onSandboxBoundarySettled: (sessionId: string) => Promise<void> | void;
}

Expand Down Expand Up @@ -205,6 +210,7 @@ export class HostInteractionCoordinator implements RuntimeInteractionAuthority {
readonly #preflightSessionSnapshot: HostInteractionCoordinatorOptions['preflightSessionSnapshot'];
readonly #refreshCanonicalContinuity: HostInteractionCoordinatorOptions['refreshCanonicalContinuity'];
readonly #onPoison: HostInteractionCoordinatorOptions['onPoison'];
readonly #resolveSandboxBoundaryGraphWake: HostInteractionCoordinatorOptions['resolveSandboxBoundaryGraphWake'];
readonly #onSandboxBoundarySettled: HostInteractionCoordinatorOptions['onSandboxBoundarySettled'];
readonly #runs = new Map<string, BoundRun>();
readonly #live = new Map<string, LiveEntry>();
Expand All @@ -220,6 +226,7 @@ export class HostInteractionCoordinator implements RuntimeInteractionAuthority {
this.#preflightSessionSnapshot = options.preflightSessionSnapshot;
this.#refreshCanonicalContinuity = options.refreshCanonicalContinuity;
this.#onPoison = options.onPoison;
this.#resolveSandboxBoundaryGraphWake = options.resolveSandboxBoundaryGraphWake;
this.#onSandboxBoundarySettled = options.onSandboxBoundarySettled;
}

Expand Down Expand Up @@ -886,7 +893,21 @@ 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 = this.#resolveSandboxBoundaryGraphWake
? await this.#resolveSandboxBoundaryGraphWake(request.sessionId, admission)
: undefined;
void Promise.resolve()
.then(() => {
if (this.#resolveSandboxBoundaryGraphWake && !resolvedRootSessionId) return;
return this.#onSandboxBoundarySettled(resolvedRootSessionId ?? request.sessionId);
})
.catch((error: unknown) => {
this.#poison(error);
});
const result = projectSandboxBoundaryInteraction(settlement.request);
if (result.status !== 'answered') {
throw this.#poison(
Expand Down
13 changes: 11 additions & 2 deletions packages/runtime-host/src/server/sandbox-boundary-graph-wake.ts
Original file line number Diff line number Diff line change
Expand Up @@ -39,14 +39,23 @@ export async function sandboxBoundaryGraphWakeRoot(
return parent.parentSessionId;
}

/** Resolve the durable lineage for a settled sandbox boundary. */
export async function resolveSandboxBoundaryGraphWake(
sessionId: string,
sessions: SandboxBoundaryGraphWakeHeaderReader,
graphIds: { listGraphIds(rootSessionId: string): Promise<readonly string[]> },
): Promise<string | undefined> {
const header = await sessions.readHeaderSnapshot(sessionId);
return sandboxBoundaryGraphWakeRoot(header, graphIds);
}

/** Resolve durable lineage before notifying the root graph supervisor. */
export async function notifySandboxBoundaryGraphWake(
sessionId: string,
sessions: SandboxBoundaryGraphWakeHeaderReader,
graphIds: { listGraphIds(rootSessionId: string): Promise<readonly string[]> },
notifyPermissionResponse: (rootSessionId: string) => Promise<void> | void,
): Promise<void> {
const header = await sessions.readHeaderSnapshot(sessionId);
const rootSessionId = await sandboxBoundaryGraphWakeRoot(header, graphIds);
const rootSessionId = await resolveSandboxBoundaryGraphWake(sessionId, sessions, graphIds);
if (rootSessionId) await notifyPermissionResponse(rootSessionId);
}