diff --git a/services/cloud-agent-next/src/persistence/SandboxControl.ts b/services/cloud-agent-next/src/persistence/SandboxControl.ts index a08bdb2397..cf061548c8 100644 --- a/services/cloud-agent-next/src/persistence/SandboxControl.ts +++ b/services/cloud-agent-next/src/persistence/SandboxControl.ts @@ -506,6 +506,7 @@ export type SandboxControlStatus = StatusProjection & { allocationIncarnation?: string; operationResults?: true; runtimeRecovery?: true; + runtimeReplacementInFlight?: true; }; export type ControlRuntimeCredentialProxyFence = { @@ -1982,7 +1983,7 @@ export class SandboxControl extends DurableObject { ) { throw new SandboxAcquisitionLostError(); } - return this.statusForAllocation(current, this.allocationIncarnationOf(current)); + return this.statusForAllocation(current, this.allocationIncarnationOf(current), sessionId); }); } @@ -2853,21 +2854,23 @@ export class SandboxControl extends DurableObject { ); } - async getStatus(): Promise { + async getStatus(input?: { sessionId?: string }): Promise { await this.ensureOperationalInitialized(); const record = await this.readCanonicalAllocation(); - return this.statusForAllocation(record, this.allocationIncarnationOf(record)); + return this.statusForAllocation(record, this.allocationIncarnationOf(record), input?.sessionId); } private async statusForAllocation( record: AllocationRecord, - allocationIncarnation?: string + allocationIncarnation?: string, + sessionId?: string ): Promise { const connection = this.connectionState(); const work = await this.workState(); const runtime = this.readyWrapperRuntime(); const physical = legacyPhysicalState(record); const projection = projectStatus({ allocation: record, ownerPresent: true, now: Date.now() }); + const runtimeReplacementInFlight = this.replacementInFlight(record, sessionId); return { ...projection, physical, @@ -2882,6 +2885,7 @@ export class SandboxControl extends DurableObject { ? { operationResults: true as const } : {}), ...(runtime?.runtimeRecovery ? { runtimeRecovery: true as const } : {}), + ...(runtimeReplacementInFlight ? { runtimeReplacementInFlight: true as const } : {}), }; } @@ -2900,6 +2904,21 @@ export class SandboxControl extends DurableObject { return { ...rest, connection: 'connected' }; } + /** + * True while a runtime replacement is in flight for this workspace: the + * canonical allocation is still creating its runtime, so the workspace has no + * bound runtime and one is on the way. `stopped` and `unknown` allocations are + * not a replacement, so a runtime that never comes back still reaches the + * caller's terminal preparation path. + * + * The canonical aggregate is per workspace, not per session, so the probe is + * not scoped to one session; `sessionId` is accepted for the RPC contract. + */ + private replacementInFlight(record: AllocationRecord, sessionId?: string): boolean { + void sessionId; + return record.state.kind === 'creating'; + } + async getSandboxStatus(input: { ownerId: string; provider: AgentSandboxProvider; diff --git a/services/cloud-agent-next/src/sandbox-session/SandboxSession.ts b/services/cloud-agent-next/src/sandbox-session/SandboxSession.ts index d811f5ae41..dcc4193587 100644 --- a/services/cloud-agent-next/src/sandbox-session/SandboxSession.ts +++ b/services/cloud-agent-next/src/sandbox-session/SandboxSession.ts @@ -246,6 +246,8 @@ import { replacePreparationAttemptId, rotateLostPreparationAttempt, resolveSessionMessageIntent, + retireAttachProof, + RUNTIME_REPLACEMENT_WAIT_LIMIT, streamCloudStatus, streamQueuedSnapshots, terminalAtOf, @@ -2274,6 +2276,7 @@ export class SandboxSession extends DurableObject { authorization: authorization.data, }) ); + this.clearRuntimeReplacementWaits(epoch); logControlDiagnostic('native_fence_transition', { sessionId: this.sessionId, sandboxId: input.sandboxId, @@ -3749,12 +3752,12 @@ export class SandboxSession extends DurableObject { const deadlineAt = queuedState?.deadlineAt ?? undefined; if (!queued || !queuedState || deadlineAt === undefined) return; if (Date.now() >= deadlineAt && !queued.proofs?.prompt?.dispatched) { - await this.failDelivery( + await this.awaitRuntimeReplacementOrFail({ messageId, - 'preparation_timeout', - activeWrapperInstanceId(queued), - queuedState.deliveryRetryScope - ); + sandboxId, + wrapperInstanceId: activeWrapperInstanceId(queued), + scope: queuedState.deliveryRetryScope, + }); return; } if (queuedState.retryNotBefore !== undefined && queuedState.retryNotBefore > Date.now()) { @@ -3804,6 +3807,16 @@ export class SandboxSession extends DurableObject { message.state.wrapperInstanceId === runtime.data ? message.state.unresolvedDispatch : undefined, + // A different wrapper incarnation means the replacement the + // head was waiting on has rebound, so the deferral chain that + // spent the previous budget is over. The next replacement + // starts with a fresh budget instead of inheriting the spent + // one. A directory-native replacement keeps the wrapper + // incarnation and is reset where the native runtime is + // recorded instead (`recordNativeRuntime`). + ...(message.state.wrapperInstanceId === runtime.data + ? {} + : { replacementWaits: undefined }), }, } : message @@ -3964,12 +3977,12 @@ export class SandboxSession extends DurableObject { await this.armQueueRetry(Math.min(deadlineAt, Date.now() + QUEUE_RETRY_MS)); if (!isCurrent()) return; if (Date.now() >= deadlineAt && !queued.proofs?.prompt?.dispatched) { - await this.failDelivery( + await this.awaitRuntimeReplacementOrFail({ messageId, - 'preparation_timeout', + sandboxId, wrapperInstanceId, - queuedState.deliveryRetryScope - ); + scope: queuedState.deliveryRetryScope, + }); return; } const intent = queuedState.intent; @@ -4133,7 +4146,7 @@ export class SandboxSession extends DurableObject { await this.armQueueRetry(Math.min(deadlineAt, Date.now() + QUEUE_RETRY_MS)); return; } - await this.failDelivery(messageId, 'preparation_timeout', wrapperInstanceId); + await this.awaitRuntimeReplacementOrFail({ messageId, sandboxId, wrapperInstanceId }); return; } const provision = provisionPreparingStep(observed.physical, allowCreate); @@ -4485,12 +4498,12 @@ export class SandboxSession extends DurableObject { if (!message) return; if (input.hadAcquisition && isSandboxAcquisitionLostError(error)) { if (Date.now() >= deadlineAt) { - await this.failDelivery( + await this.awaitRuntimeReplacementOrFail({ messageId, - 'preparation_timeout', + sandboxId: this.terminalLifecycle.getStoredMetadata()?.workspace?.sandboxId, wrapperInstanceId, - message.state.kind === 'queued' ? message.state.deliveryRetryScope : undefined - ); + scope: message.state.kind === 'queued' ? message.state.deliveryRetryScope : undefined, + }); return; } const retryNotBefore = Math.min(deadlineAt, Date.now() + QUEUE_RETRY_MS); @@ -4535,7 +4548,13 @@ export class SandboxSession extends DurableObject { : retryableRejection && !unresolvedDispatch; const scope: 'message' | 'runtime' = messageRetryScope ? 'message' : 'runtime'; if (Date.now() >= deadlineAt) { - await this.failDelivery(messageId, 'preparation_timeout', wrapperInstanceId, scope, detail); + await this.awaitRuntimeReplacementOrFail({ + messageId, + sandboxId: this.terminalLifecycle.getStoredMetadata()?.workspace?.sandboxId, + wrapperInstanceId, + scope, + detail, + }); return; } const busy = retryableRejection && error.code === 'session_busy'; @@ -4580,6 +4599,134 @@ export class SandboxSession extends DurableObject { ); } + /** + * True while the control plane still reports a runtime replacement in flight + * for this workspace: the canonical allocation has not bound a runtime yet and + * a replacement is being created. The probe is best-effort: an absent session + * id or a transport failure reports "no replacement in flight", so it adds no + * failure mode of its own and the existing terminal path runs unchanged. + */ + private async runtimeReplacementInFlight(sandboxId: string): Promise { + if (this.sessionId === undefined) return false; + try { + const status = await sandboxControlRpc(this.env, sandboxId).getStatus({ + sessionId: this.sessionId, + }); + return status.runtimeReplacementInFlight === true; + } catch { + return false; + } + } + + /** + * End the deferral chain once the session binds a native runtime, so the next + * replacement gets a full budget instead of inheriting a spent one. Keyed on + * the native runtime identity because a directory-native replacement recreates + * the native runtime in place inside the same wrapper incarnation, which the + * wrapper-identity reset in `recordRuntime` cannot observe. + */ + private clearRuntimeReplacementWaits(epoch: number): void { + const messages = this.loadMessages(); + if ( + !messages.some( + message => message.state.kind === 'queued' && message.state.replacementWaits !== undefined + ) + ) + return; + this.saveMessages( + messages.map(message => + message.state.kind !== 'queued' || message.state.replacementWaits === undefined + ? message + : { ...message, state: { ...message.state, replacementWaits: undefined } } + ), + epoch + ); + } + + /** + * A preparation deadline that lands while the control plane reports a + * runtime replacement in flight does not terminalize the head. It re-arms the + * head's existing durable delivery deadline and the existing 5 s queue-retry + * alarm, so the delivery resumes once the replacement rebinds. + * + * The deferral is re-validated after the probe: `runtimeReplacementInFlight` + * is a cross-DO RPC, so another event may deliver the head while it is + * outstanding, and `commitSavedMessages` treats an accepted row as mutable. + * Only a still-queued head is rewritten. + * + * The wait is bounded: each deferral spends one unit of + * `RUNTIME_REPLACEMENT_WAIT_LIMIT` for the current replacement cycle. A + * replacement that never completes exhausts the budget and the existing + * terminal path fails the head exactly as before, so a runtime that never + * returns still reaches `preparation_timeout`. Binding a replacement runtime + * ends the cycle and resets the budget, so a later, unrelated replacement does + * not inherit a partially spent one. + * + * Every re-arm mints a fresh preparation attempt. A preparation attempt is + * the acquisition request identity, and `SandboxControl` binds it to its + * original deadline and rejects a changed deadline for the same id, so the + * re-armed window must be a new acquisition rather than a mutated one. A live + * attach proof bound the old attempt identity, so it is retired into the slot + * a late result is matched against and the replacement attach can carry the + * new identity. + */ + private async awaitRuntimeReplacementOrFail(input: { + messageId: string; + sandboxId?: string; + wrapperInstanceId?: string; + scope?: 'message' | 'runtime'; + detail?: string; + }): Promise { + if (input.sandboxId !== undefined && (await this.runtimeReplacementInFlight(input.sandboxId))) { + const epoch = this.terminalLifecycle.captureEpoch(); + if (epoch === null) return; + const current = this.loadMessages().find(message => message.messageId === input.messageId); + if (!current || current.state.kind !== 'queued') return; + const waits = current.state.replacementWaits ?? 0; + if (waits >= RUNTIME_REPLACEMENT_WAIT_LIMIT) { + await this.failDelivery( + input.messageId, + 'preparation_timeout', + input.wrapperInstanceId, + input.scope, + input.detail + ); + return; + } + const now = Date.now(); + const extended = this.loadMessages().map(message => { + if (message.messageId !== input.messageId || message.state.kind !== 'queued') + return message; + const proofs = message.proofs ? { ...message.proofs } : undefined; + if (proofs?.attach) { + proofs.retiredAttach = retireAttachProof(proofs.attach, proofs.retiredAttach); + delete proofs.attach; + } + return { + ...message, + ...(proofs ? { proofs } : {}), + state: { + ...message.state, + preparationAttemptId: crypto.randomUUID(), + preparationWait: undefined, + deadlineAt: now + SESSION_DELIVERY_TIMEOUT_MS, + replacementWaits: waits + 1, + }, + }; + }); + if (!this.saveMessages(extended, epoch)) return; + await this.armQueueRetry(now + QUEUE_RETRY_MS); + return; + } + await this.failDelivery( + input.messageId, + 'preparation_timeout', + input.wrapperInstanceId, + input.scope, + input.detail + ); + } + private async failDelivery( messageId: string, reason: string, diff --git a/services/cloud-agent-next/src/sandbox-session/control-rpc.ts b/services/cloud-agent-next/src/sandbox-session/control-rpc.ts index dc7b9708c8..a6a6e623a1 100644 --- a/services/cloud-agent-next/src/sandbox-session/control-rpc.ts +++ b/services/cloud-agent-next/src/sandbox-session/control-rpc.ts @@ -50,15 +50,17 @@ type SandboxControlRpc = { allocationIncarnation?: string; operationResults?: true; runtimeRecovery?: true; + runtimeReplacementInFlight?: true; attachment?: SessionAttachPayload; }>; - getStatus(): Promise<{ + getStatus(input?: { sessionId?: string }): Promise<{ connection: ConnectionState; physical: PhysicalState; wrapperInstanceId?: string; allocationIncarnation?: string; operationResults?: true; runtimeRecovery?: true; + runtimeReplacementInFlight?: true; }>; getRuntimeCredentialProxyFence(input: { ownerId: string; @@ -110,7 +112,8 @@ export function sandboxControlRpc( 'prepareSessionCredentials' ), ensureReady: input => stub().ensureReady(input), - getStatus: () => withDORetry(stub, control => control.getStatus(), 'getStatus', config()), + getStatus: input => + withDORetry(stub, control => control.getStatus(input), 'getStatus', config()), getRuntimeCredentialProxyFence: input => withDORetry( stub, diff --git a/services/cloud-agent-next/src/sandbox-session/recovery/awaiting-runtime-replacement-delivers-queue.test.ts b/services/cloud-agent-next/src/sandbox-session/recovery/awaiting-runtime-replacement-delivers-queue.test.ts new file mode 100644 index 0000000000..af6552d446 --- /dev/null +++ b/services/cloud-agent-next/src/sandbox-session/recovery/awaiting-runtime-replacement-delivers-queue.test.ts @@ -0,0 +1,540 @@ +/** + * The production shape from Pylon 28572: a queued message whose workspace has no + * runtime while a runtime replacement is in flight. The preparation deadline must + * not terminalize the head with `preparation_timeout`; the head waits for the + * replacement and is delivered once it rebinds. A replacement that never binds + * still reaches the terminal path once the deferral budget is spent. + */ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { + DIRECTORY, + NEXT_RUNTIME_ID, + RUNTIME_ID, + SANDBOX_ID, + SESSION_ID, + controlFailure, + controlResponse, + createSessionFixture, + delegateRequest, +} from '../session-fixture.test-helpers.js'; +import { RUNTIME_REPLACEMENT_WAIT_LIMIT } from '../session-message-queue.js'; +import type { SessionEnvelope, SessionMessage } from '../../sandbox-state/model/session.js'; +import { readSessionValueSync, writeSessionMessages } from '../../sandbox-state/persist/access.js'; +import { readRawSessionMessages } from '../../sandbox-state/persist/load.js'; +import type { SessionOperationAuthorization } from '../../shared/sandbox-control-protocol.js'; + +const orchestrationMocks = vi.hoisted(() => ({ + eventQueries: vi.fn(), + signedAttachments: vi.fn(), + broadcast: vi.fn(), +})); + +vi.mock('cloudflare:workers', () => ({ + DurableObject: class { + constructor( + protected ctx: unknown, + protected env: unknown + ) {} + }, +})); +vi.mock('@cloudflare/sandbox', () => ({ getSandbox: vi.fn() })); +vi.mock('drizzle-orm/durable-sqlite', () => ({ drizzle: vi.fn() })); +vi.mock('drizzle-orm/durable-sqlite/migrator', () => ({ migrate: vi.fn(async () => undefined) })); +vi.mock('../../../drizzle/migrations', () => ({ default: {} })); +vi.mock('../../session/queries/index.js', () => ({ + createEventQueries: orchestrationMocks.eventQueries, +})); +vi.mock('../../model-validation.js', () => ({ + assertKiloModelAvailable: vi.fn(async () => undefined), +})); +vi.mock('../../execution/attachment-prompt-parts.js', () => ({ + buildSignedPromptAttachments: orchestrationMocks.signedAttachments, +})); +vi.mock('../../websocket/stream.js', () => ({ + createStreamHandler: ( + _state: unknown, + _queries: unknown, + _sessionId: string, + options?: { + deriveCloudStatus?: () => Promise; + deriveQueuedMessages?: () => Promise; + readPendingInteractions?: () => unknown; + deriveSessionStatus?: () => Promise; + getPreparationSnapshots?: () => Promise; + } + ) => ({ + broadcastEvent: orchestrationMocks.broadcast, + handleStreamRequest: async () => + Response.json({ + cloudStatus: await options?.deriveCloudStatus?.(), + queuedMessages: await options?.deriveQueuedMessages?.(), + pendingInteractions: options?.readPendingInteractions?.(), + sessionStatus: await options?.deriveSessionStatus?.(), + preparationSnapshots: await options?.getPreparationSnapshots?.(), + }), + }), +})); + +const fixtureDeps = { + eventQueries: orchestrationMocks.eventQueries, + signedAttachments: orchestrationMocks.signedAttachments, +}; + +type Fixture = ReturnType; + +/** The head's durable queued state, or undefined when it is not queued. */ +function queuedStateOf(fixture: Fixture, messageId: string) { + const record = fixture.record(messageId); + return record?.state.kind === 'queued' ? record.state : undefined; +} + +function queuedDeadline(fixture: Fixture, messageId: string): number { + const deadlineAt = queuedStateOf(fixture, messageId)?.deadlineAt; + if (deadlineAt === undefined || deadlineAt === null) throw new Error('missing delivery deadline'); + return deadlineAt; +} + +/** Rewrites the stored head through the same writer the DO uses. */ +function rewriteHead( + fixture: Fixture, + messageId: string, + rewrite: (message: SessionMessage) => SessionMessage, + binding?: SessionEnvelope['binding'] +): void { + const envelope = readSessionValueSync(fixture.storage.kv); + if (!envelope) throw new Error('missing session envelope'); + writeSessionMessages( + fixture.storage.kv, + binding ?? envelope.binding, + readRawSessionMessages(fixture.storage.kv).map(message => + message.messageId === messageId ? rewrite(message) : message + ) + ); +} + +/** Forces the head's durable delivery deadline, as an expired window would. */ +function setDeliveryDeadline(fixture: Fixture, messageId: string, deadlineAt: number): void { + rewriteHead(fixture, messageId, message => + message.state.kind === 'queued' + ? { ...message, state: { ...message.state, deadlineAt } } + : message + ); +} + +describe('awaiting a runtime replacement', () => { + beforeEach(() => { + vi.useFakeTimers(); + vi.setSystemTime(1_000_000); + orchestrationMocks.broadcast.mockClear(); + }); + afterEach(() => { + vi.useRealTimers(); + }); + + it('defers the head past its deadline and delivers once the replacement rebinds', async () => { + const fixture = createSessionFixture(fixtureDeps); + fixture.setStatus({ + physical: 'running', + connection: 'connected', + wrapperInstanceId: RUNTIME_ID, + allocationIncarnation: 'incarnation_1', + runtimeReplacementInFlight: true, + }); + await fixture.admit('a'); + await fixture.flush(); + + const deadlineAt = Date.now() + 20_000; + setDeliveryDeadline(fixture, 'a', deadlineAt); + + // The deadline lands while the workspace has no runtime and a replacement is + // in flight: the head stays queued on a fresh window instead of failing. + vi.setSystemTime(deadlineAt); + await fixture.fireAlarm(); + await fixture.flush(); + expect(fixture.record('a')?.state.kind).toBe('queued'); + expect(queuedDeadline(fixture, 'a')).toBeGreaterThan(deadlineAt); + + // The replacement rebinds inside the window: the head is delivered. + fixture.setStatus({ + physical: 'running', + connection: 'ready', + wrapperInstanceId: NEXT_RUNTIME_ID, + allocationIncarnation: 'incarnation_1', + }); + await fixture.fireAlarm(); + await fixture.flush(); + expect(fixture.record('a')?.state.kind).toBe('accepted'); + }); + + it('re-arms a head that already bound an acquisition as a new acquisition', async () => { + const fixture = createSessionFixture(fixtureDeps); + const originalEnsureReady = fixture.control.ensureReady.getMockImplementation(); + if (!originalEnsureReady) throw new Error('Missing ensureReady fixture'); + // The control plane binds an acquisition id to its original deadline and + // rejects a changed deadline for the same id, so a re-armed window must be a + // new acquisition rather than a mutated one. + const boundDeadlines = new Map(); + fixture.control.ensureReady.mockImplementation(async input => { + const acquisition = input.acquisition; + if (acquisition) { + const bound = boundDeadlines.get(acquisition.id); + if (bound !== undefined && bound !== acquisition.deadlineAt) + throw new Error('Sandbox acquisition deadline changed'); + boundDeadlines.set(acquisition.id, acquisition.deadlineAt); + } + return originalEnsureReady(input); + }); + fixture.setStatus({ + physical: 'running', + connection: 'connected', + wrapperInstanceId: RUNTIME_ID, + allocationIncarnation: 'incarnation_1', + }); + await fixture.admit('a'); + await fixture.flush(); + expect(queuedStateOf(fixture, 'a')?.preparationAttemptId).toBeDefined(); + expect(boundDeadlines.size).toBe(1); + + const deadlineAt = Date.now() + 20_000; + setDeliveryDeadline(fixture, 'a', deadlineAt); + fixture.setStatus({ + physical: 'running', + connection: 'connected', + wrapperInstanceId: RUNTIME_ID, + allocationIncarnation: 'incarnation_1', + runtimeReplacementInFlight: true, + }); + vi.setSystemTime(deadlineAt); + await fixture.fireAlarm(); + await fixture.flush(); + expect(fixture.record('a')?.state.kind).toBe('queued'); + expect(queuedDeadline(fixture, 'a')).toBeGreaterThan(deadlineAt); + + fixture.setStatus({ + physical: 'running', + connection: 'ready', + wrapperInstanceId: NEXT_RUNTIME_ID, + allocationIncarnation: 'incarnation_1', + }); + await fixture.fireAlarm(); + await fixture.flush(); + expect(fixture.record('a')?.state.kind).toBe('accepted'); + expect(boundDeadlines.size).toBe(2); + }); + + it('retires a bound attach proof into retiredAttach while deferring', async () => { + const fixture = createSessionFixture(fixtureDeps); + delegateRequest(fixture, 'session.attach', async () => controlFailure(true, 'not_ready')); + fixture.setStatus({ + physical: 'running', + connection: 'ready', + wrapperInstanceId: RUNTIME_ID, + allocationIncarnation: 'incarnation_1', + operationResults: true, + runtimeReplacementInFlight: true, + }); + await fixture.admit('a'); + await fixture.flush(); + + const attach = fixture.record('a')?.proofs?.attach; + if (!attach?.dispatched) throw new Error('Missing dispatched attach proof'); + + const deadlineAt = Date.now() + 20_000; + setDeliveryDeadline(fixture, 'a', deadlineAt); + vi.setSystemTime(deadlineAt); + await fixture.fireAlarm(); + await fixture.flush(); + + const record = fixture.record('a'); + expect(record?.state.kind).toBe('queued'); + expect(record?.proofs?.attach).toBeUndefined(); + expect(record?.proofs?.retiredAttach).toEqual(attach); + }); + + it('does not rewrite a head another event delivered during the probe', async () => { + const fixture = createSessionFixture(fixtureDeps); + fixture.setStatus({ + physical: 'running', + connection: 'connected', + wrapperInstanceId: RUNTIME_ID, + allocationIncarnation: 'incarnation_1', + }); + await fixture.admit('a'); + await fixture.flush(); + const attemptId = queuedStateOf(fixture, 'a')?.preparationAttemptId; + if (!attemptId) throw new Error('Missing preparation attempt'); + + const deadlineAt = Date.now() + 20_000; + setDeliveryDeadline(fixture, 'a', deadlineAt); + + // `runtimeReplacementInFlight` is a cross-DO probe: another event may deliver + // the head while it is outstanding, so only a still-queued head is rewritten. + const originalGetStatus = fixture.control.getStatus.getMockImplementation(); + if (!originalGetStatus) throw new Error('Missing getStatus fixture'); + fixture.control.getStatus.mockImplementation(async () => { + // The head was delivered by the concurrent event; acceptance requires a + // bound attachment, so the stored envelope carries the bound handle. + rewriteHead( + fixture, + 'a', + message => + message.state.kind === 'queued' + ? { + ...message, + state: { + kind: 'accepted', + intent: message.state.intent, + acceptedAt: Date.now(), + executionDeadlineAt: Date.now() + 60_000, + ...(message.state.legacy === undefined ? {} : { legacy: message.state.legacy }), + ...(message.state.legacyInvalidIntent === undefined + ? {} + : { legacyInvalidIntent: message.state.legacyInvalidIntent }), + ...(message.state.queuedAt === undefined + ? {} + : { queuedAt: message.state.queuedAt }), + ...(message.state.wrapperInstanceId === undefined + ? {} + : { wrapperInstanceId: message.state.wrapperInstanceId }), + ...(message.state.preparationAttemptId === undefined + ? {} + : { preparationAttemptId: message.state.preparationAttemptId }), + }, + } + : message, + { + kind: 'bound', + handle: { incarnation: 'incarnation_1', wrapper: RUNTIME_ID, epoch: 0 }, + } + ); + return { ...(await originalGetStatus()), runtimeReplacementInFlight: true as const }; + }); + + vi.setSystemTime(deadlineAt); + await fixture.fireAlarm(); + await fixture.flush(); + + expect(fixture.record('a')).toMatchObject({ + state: { kind: 'accepted', preparationAttemptId: attemptId }, + }); + }); + + it('fails the head once the deferral budget is spent', async () => { + const fixture = createSessionFixture(fixtureDeps); + fixture.setStatus({ + physical: 'running', + connection: 'connected', + wrapperInstanceId: RUNTIME_ID, + allocationIncarnation: 'incarnation_1', + runtimeReplacementInFlight: true, + }); + await fixture.admit('a'); + await fixture.flush(); + setDeliveryDeadline(fixture, 'a', Date.now() + 20_000); + + // Each deferral spends one unit of the budget and grants a fresh window. + for (let wait = 0; wait < RUNTIME_REPLACEMENT_WAIT_LIMIT; wait += 1) { + const deadlineAt = queuedDeadline(fixture, 'a'); + vi.setSystemTime(deadlineAt); + await fixture.fireAlarm(); + await fixture.flush(); + expect(fixture.record('a')?.state.kind).toBe('queued'); + expect(queuedDeadline(fixture, 'a')).toBeGreaterThan(deadlineAt); + expect(queuedStateOf(fixture, 'a')?.replacementWaits).toBe(wait + 1); + } + + // A replacement that never completes exhausts the budget and the existing + // terminal path fails the head exactly as before. + vi.setSystemTime(queuedDeadline(fixture, 'a')); + await fixture.fireAlarm(); + await fixture.flush(); + expect(fixture.record('a')?.state).toMatchObject({ + kind: 'failed', + reason: 'preparation_timeout', + }); + }); + + it('resets the deferral budget when a native runtime rebinds in place', async () => { + const fixture = createSessionFixture(fixtureDeps); + const retiredNativeRuntimeId = '11111111-1111-4111-8111-111111111111'; + const replacementNativeRuntimeId = '44444444-4444-4444-8444-444444444444'; + + fixture.setStatus({ + physical: 'running', + connection: 'connected', + wrapperInstanceId: RUNTIME_ID, + allocationIncarnation: 'incarnation_1', + runtimeReplacementInFlight: true, + }); + await fixture.admit('a'); + await fixture.flush(); + + // The head binds a native runtime before the retirement, so the fence holds + // its attach epoch and the in-place rebind must advance past it. + const authorization: SessionOperationAuthorization = { + operation: 'session.attach', + operationId: '11111111-1111-4111-8111-111111111111', + messageId: 'a', + session: { sessionId: SESSION_ID, kiloSessionId: 'kilo_root', directory: DIRECTORY }, + wrapperInstanceId: RUNTIME_ID, + dispatchDeadlineAt: Date.now() + 60_000, + }; + rewriteHead(fixture, 'a', message => ({ + ...message, + proofs: { + attach: { + authorization, + dispatched: true, + completedAt: Date.now(), + attachmentEpoch: 1, + }, + }, + })); + await fixture.session.recordNativeRuntime({ + sandboxId: SANDBOX_ID, + wrapperInstanceId: RUNTIME_ID, + nativeRuntimeId: retiredNativeRuntimeId, + authorization, + }); + expect(fixture.values.get('native_runtime_fence')).toMatchObject({ + nativeRuntimeId: retiredNativeRuntimeId, + attachmentEpoch: 1, + }); + + // Fence A: the head defers once and spends one unit of the budget. + const fenceDeadline = Date.now() + 20_000; + setDeliveryDeadline(fixture, 'a', fenceDeadline); + vi.setSystemTime(fenceDeadline); + await fixture.fireAlarm(); + await fixture.flush(); + expect(queuedStateOf(fixture, 'a')?.replacementWaits).toBe(1); + expect(fixture.record('a')?.proofs?.retiredAttach).toMatchObject({ attachmentEpoch: 1 }); + + // Fence A clears by recreating only the native runtime in place: the wrapper + // incarnation stays RUNTIME_ID, so only the attach result's native runtime + // identity changes. The attach binds it while the prompt stays retryable, so + // the head is still queued when the replacement has bound. + delegateRequest(fixture, 'session.attach', async () => + controlResponse({ attached: true, nativeRuntimeId: replacementNativeRuntimeId }) + ); + delegateRequest(fixture, 'session.prompt', async () => controlFailure(true, 'not_ready')); + fixture.setStatus({ + physical: 'running', + connection: 'ready', + wrapperInstanceId: RUNTIME_ID, + allocationIncarnation: 'incarnation_1', + operationResults: true, + }); + await fixture.fireAlarm(); + await fixture.flush(); + + expect(fixture.values.get('native_runtime_fence')).toMatchObject({ + nativeRuntimeId: replacementNativeRuntimeId, + }); + expect(fixture.record('a')).toMatchObject({ + state: { kind: 'queued', wrapperInstanceId: RUNTIME_ID, replacementWaits: undefined }, + }); + }); + + it('keeps the attach-epoch pool across a second deferral so the rebind still binds', async () => { + const fixture = createSessionFixture(fixtureDeps); + const retiredNativeRuntimeId = '11111111-1111-4111-8111-111111111111'; + const replacementNativeRuntimeId = '44444444-4444-4444-8444-444444444444'; + + fixture.setStatus({ + physical: 'running', + connection: 'connected', + wrapperInstanceId: RUNTIME_ID, + allocationIncarnation: 'incarnation_1', + runtimeReplacementInFlight: true, + }); + await fixture.admit('a'); + await fixture.flush(); + + // The head binds a native runtime at attach epoch 1 before the retirement, so + // the fence holds 1 and the in-place rebind must advance past it. + const boundAuthorization: SessionOperationAuthorization = { + operation: 'session.attach', + operationId: '11111111-1111-4111-8111-111111111111', + messageId: 'a', + session: { sessionId: SESSION_ID, kiloSessionId: 'kilo_root', directory: DIRECTORY }, + wrapperInstanceId: RUNTIME_ID, + dispatchDeadlineAt: Date.now() + 60_000, + }; + rewriteHead(fixture, 'a', message => ({ + ...message, + proofs: { + attach: { + authorization: boundAuthorization, + dispatched: true, + completedAt: Date.now(), + attachmentEpoch: 1, + }, + }, + })); + await fixture.session.recordNativeRuntime({ + sandboxId: SANDBOX_ID, + wrapperInstanceId: RUNTIME_ID, + nativeRuntimeId: retiredNativeRuntimeId, + authorization: boundAuthorization, + }); + expect(fixture.values.get('native_runtime_fence')).toMatchObject({ attachmentEpoch: 1 }); + + // Deferral one retires the completed attach, whose epoch the fence holds. + const firstDeadline = Date.now() + 20_000; + setDeliveryDeadline(fixture, 'a', firstDeadline); + vi.setSystemTime(firstDeadline); + await fixture.fireAlarm(); + await fixture.flush(); + expect(fixture.record('a')?.proofs?.retiredAttach).toMatchObject({ attachmentEpoch: 1 }); + + // The redelivery dispatches a fresh attach, and its result has not arrived + // when the next deadline lands: the second deferral retires an epoch-less + // proof over the completed one. + const pendingAuthorization: SessionOperationAuthorization = { + operation: 'session.attach', + operationId: '22222222-2222-4222-8222-222222222222', + messageId: 'a', + session: { sessionId: SESSION_ID, kiloSessionId: 'kilo_root', directory: DIRECTORY }, + wrapperInstanceId: RUNTIME_ID, + dispatchDeadlineAt: Date.now() + 60_000, + }; + rewriteHead(fixture, 'a', message => ({ + ...message, + proofs: { + ...message.proofs, + attach: { authorization: pendingAuthorization, dispatched: true }, + }, + })); + const secondDeadline = Date.now() + 20_000; + setDeliveryDeadline(fixture, 'a', secondDeadline); + vi.setSystemTime(secondDeadline); + await fixture.fireAlarm(); + await fixture.flush(); + + expect(queuedStateOf(fixture, 'a')?.replacementWaits).toBe(2); + // Dropping the epoch here makes the next attach mint at or below the fence's + // epoch, so the replacement looks like a stale result and never binds. + expect(fixture.record('a')?.proofs?.retiredAttach).toMatchObject({ attachmentEpoch: 1 }); + + // The replacement binds in place: the attach result carries the replacement + // native runtime identity, and the fence must advance to it. + delegateRequest(fixture, 'session.attach', async () => + controlResponse({ attached: true, nativeRuntimeId: replacementNativeRuntimeId }) + ); + delegateRequest(fixture, 'session.prompt', async () => controlFailure(true, 'not_ready')); + fixture.setStatus({ + physical: 'running', + connection: 'ready', + wrapperInstanceId: RUNTIME_ID, + allocationIncarnation: 'incarnation_1', + operationResults: true, + }); + await fixture.fireAlarm(); + await fixture.flush(); + + expect(fixture.values.get('native_runtime_fence')).toMatchObject({ + nativeRuntimeId: replacementNativeRuntimeId, + }); + }); +}); diff --git a/services/cloud-agent-next/src/sandbox-session/session-message-queue.test.ts b/services/cloud-agent-next/src/sandbox-session/session-message-queue.test.ts index 85d9133946..c994e02c1e 100644 --- a/services/cloud-agent-next/src/sandbox-session/session-message-queue.test.ts +++ b/services/cloud-agent-next/src/sandbox-session/session-message-queue.test.ts @@ -972,6 +972,32 @@ describe('releaseUnconfirmedAttach', () => { expect(released?.[0]?.proofs?.attach).toBeUndefined(); }); + it('keeps the retired attach epoch when the replacement proof has none', () => { + const message = queuedMessage({ authorization, dispatched: true }); + if (!message.proofs) throw new Error('Missing attach proof'); + const released = releaseUnconfirmedAttach( + [ + { + ...message, + proofs: { + ...message.proofs, + retiredAttach: { + authorization, + dispatched: true, + completedAt: 400, + attachmentEpoch: 3, + }, + }, + }, + ], + authorization + ); + + // The dispatched attach carries no epoch until its result arrives, so the + // overwritten slot would drop the watermark the fence holds. + expect(released?.[0]?.proofs?.retiredAttach).toMatchObject({ attachmentEpoch: 3 }); + }); + it('refuses a message id that is not in the messages array', () => { expect(releaseUnconfirmedAttach([], authorization)).toBeUndefined(); }); @@ -5622,7 +5648,10 @@ describe('SandboxSession orchestration', () => { expect(input.acquisition).toEqual(acquisition); expect(input.allowCreate).toBeUndefined(); } - expect(fixture.control.getStatus).not.toHaveBeenCalled(); + // The deadline check probes the control plane for an in-flight runtime + // replacement before terminalizing. Absent one, it terminalizes without + // dispatching or quarantining. + expect(fixture.control.getStatus).toHaveBeenCalledWith({ sessionId: SESSION_ID }); expect(fixture.control.request).not.toHaveBeenCalled(); }); diff --git a/services/cloud-agent-next/src/sandbox-session/session-message-queue.ts b/services/cloud-agent-next/src/sandbox-session/session-message-queue.ts index 84eb047789..b78123f6c1 100644 --- a/services/cloud-agent-next/src/sandbox-session/session-message-queue.ts +++ b/services/cloud-agent-next/src/sandbox-session/session-message-queue.ts @@ -74,6 +74,15 @@ export function noOutputRecoveryAllowed(recoveryAttempts: number): boolean { return recoveryAttempts < NO_OUTPUT_RECOVERY_LIMIT; } +/** + * Deferral budget for a head whose preparation deadline lands while the control + * plane reports a runtime replacement in flight. Each deferral grants a fresh + * delivery window, so the budget caps the wait inside one replacement cycle and + * leaves the existing terminal path to run when a replacement never completes. + * Binding the replacement runtime ends the cycle and starts the next budget. + */ +export const RUNTIME_REPLACEMENT_WAIT_LIMIT = 6; + // --------------------------------------------------------------------------- // Canonical field access. The wire model nests per-state fields under `state`; // these helpers keep the one mapping in one place. @@ -470,7 +479,10 @@ export function releaseUnadmittedWaitingMessages( ]); // Preserve intent and deadline: the head keeps its original preparation // bound. A released attach proof is retained for late results. - const proofs = releaseAttach && attach ? { retiredAttach: attach } : undefined; + const proofs = + releaseAttach && attach + ? { retiredAttach: retireAttachProof(attach, message.proofs?.retiredAttach) } + : undefined; return withProofs({ ...message, state: cleared }, proofs); }), releasedIds, @@ -553,7 +565,10 @@ export function releaseCompletedRetryableAttach( return messages.map(message => { const attach = message.messageId === messageId ? message.proofs?.attach : undefined; if (!attach?.dispatched || attach.result?.ok !== false) return message; - const proofs: MessageProofs = { ...message.proofs, retiredAttach: attach }; + const proofs: MessageProofs = { + ...message.proofs, + retiredAttach: retireAttachProof(attach, message.proofs?.retiredAttach), + }; delete proofs.attach; return withProofs( { @@ -569,6 +584,24 @@ export function releaseCompletedRetryableAttach( }); } +/** + * Retire an attach proof into the single `retiredAttach` slot without letting + * the slot's attach epoch drop. A dispatched attach carries no `attachmentEpoch` + * until its result arrives, and every other retired proof's epoch was already + * counted by `nextAttachmentEpoch`. Replacing the slot with an epoch-less proof + * would drop the pool to the epochs of the proofs that remain, so the next attach + * is minted at or below the epoch the native-runtime fence holds and looks like a + * stale result — refusing the in-place rebind the head is waiting for. Carry the + * highest epoch either proof has carried so the pool never falls below the fence. + */ +export function retireAttachProof( + attach: SessionOperationProof, + previous: SessionOperationProof | undefined +): SessionOperationProof { + const attachmentEpoch = Math.max(attach.attachmentEpoch ?? 0, previous?.attachmentEpoch ?? 0); + return { ...attach, ...(attachmentEpoch > 0 ? { attachmentEpoch } : {}) }; +} + /** * Retire a dispatched attach with no result so this delivery can record a * fresh attach against the same runtime, attempt, and deadline. Late results @@ -588,7 +621,10 @@ export function releaseUnconfirmedAttach( !sameSessionOperation(attach.authorization, authorization) ) return undefined; - const proofs: MessageProofs = { ...message.proofs, retiredAttach: attach }; + const proofs: MessageProofs = { + ...message.proofs, + retiredAttach: retireAttachProof(attach, message.proofs?.retiredAttach), + }; delete proofs.attach; return messages.map(item => item.messageId !== message.messageId @@ -876,6 +912,24 @@ export function streamCloudStatus( return messages.length > 0 ? { type: 'ready' } : null; } +/** + * Mint the next attach epoch. A retired proof keeps its epoch in the pool: the + * native-runtime fence rejects an in-place rebind recorded at an epoch it + * already holds, and the retired proof's epoch is the one that fence carries, so + * re-minting it would make the replacement attach look like a stale result. + */ +function nextAttachmentEpoch(messages: readonly SessionMessage[]): number { + return ( + Math.max( + 0, + ...messages.flatMap(message => [ + message.proofs?.attach?.attachmentEpoch ?? 0, + message.proofs?.retiredAttach?.attachmentEpoch ?? 0, + ]) + ) + 1 + ); +} + export function applySessionOperationResult( aggregate: SessionAggregate, delivery: SessionOperationDelivery, @@ -960,10 +1014,7 @@ export function applySessionOperationResult( : delivery.completedAt, }; const attachmentEpoch = - kind === 'attach' - ? (proof.attachmentEpoch ?? - Math.max(0, ...messages.map(item => item.proofs?.attach?.attachmentEpoch ?? 0)) + 1) - : undefined; + kind === 'attach' ? (proof.attachmentEpoch ?? nextAttachmentEpoch(messages)) : undefined; return { messages: applied.map(item => item.messageId === message.messageId @@ -1113,9 +1164,7 @@ export function completeSessionOperationAttachment( nextQueuedMessageId(messages) !== message.messageId ) return undefined; - const attachmentEpoch = - proof.attachmentEpoch ?? - Math.max(0, ...messages.map(item => item.proofs?.attach?.attachmentEpoch ?? 0)) + 1; + const attachmentEpoch = proof.attachmentEpoch ?? nextAttachmentEpoch(messages); return messages.map(item => item.messageId === message.messageId ? { diff --git a/services/cloud-agent-next/src/sandbox-state/model/session.ts b/services/cloud-agent-next/src/sandbox-state/model/session.ts index 8b73f74d1a..b5e5f3af8d 100644 --- a/services/cloud-agent-next/src/sandbox-state/model/session.ts +++ b/services/cloud-agent-next/src/sandbox-state/model/session.ts @@ -281,6 +281,14 @@ export const queuedMessageStateSchema = z * original scope cannot be recovered from it. */ deliveryRetryScope: z.enum(['message', 'runtime']).optional(), + /** + * Consecutive delivery-deadline deferrals granted while the control plane + * reported a runtime replacement in flight. Bounds a replacement that never + * completes so the head still reaches its terminal preparation timeout. The + * count resets when the head binds a replacement runtime: that closes the + * chain, so a later replacement gets its own budget. + */ + replacementWaits: z.number().int().nonnegative().optional(), }) .strict();