Skip to content
Merged
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
27 changes: 23 additions & 4 deletions services/cloud-agent-next/src/persistence/SandboxControl.ts
Original file line number Diff line number Diff line change
Expand Up @@ -506,6 +506,7 @@ export type SandboxControlStatus = StatusProjection & {
allocationIncarnation?: string;
operationResults?: true;
runtimeRecovery?: true;
runtimeReplacementInFlight?: true;
};

export type ControlRuntimeCredentialProxyFence = {
Expand Down Expand Up @@ -1982,7 +1983,7 @@ export class SandboxControl extends DurableObject<Env> {
) {
throw new SandboxAcquisitionLostError();
}
return this.statusForAllocation(current, this.allocationIncarnationOf(current));
return this.statusForAllocation(current, this.allocationIncarnationOf(current), sessionId);
});
}

Expand Down Expand Up @@ -2853,21 +2854,23 @@ export class SandboxControl extends DurableObject<Env> {
);
}

async getStatus(): Promise<SandboxControlStatus> {
async getStatus(input?: { sessionId?: string }): Promise<SandboxControlStatus> {
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<SandboxControlStatus> {
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,
Expand All @@ -2882,6 +2885,7 @@ export class SandboxControl extends DurableObject<Env> {
? { operationResults: true as const }
: {}),
...(runtime?.runtimeRecovery ? { runtimeRecovery: true as const } : {}),
...(runtimeReplacementInFlight ? { runtimeReplacementInFlight: true as const } : {}),
};
}

Expand All @@ -2900,6 +2904,21 @@ export class SandboxControl extends DurableObject<Env> {
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;
Expand Down
177 changes: 162 additions & 15 deletions services/cloud-agent-next/src/sandbox-session/SandboxSession.ts
Original file line number Diff line number Diff line change
Expand Up @@ -246,6 +246,8 @@ import {
replacePreparationAttemptId,
rotateLostPreparationAttempt,
resolveSessionMessageIntent,
retireAttachProof,
RUNTIME_REPLACEMENT_WAIT_LIMIT,
streamCloudStatus,
streamQueuedSnapshots,
terminalAtOf,
Expand Down Expand Up @@ -2274,6 +2276,7 @@ export class SandboxSession extends DurableObject<Env> {
authorization: authorization.data,
})
);
this.clearRuntimeReplacementWaits(epoch);
logControlDiagnostic('native_fence_transition', {
sessionId: this.sessionId,
sandboxId: input.sandboxId,
Expand Down Expand Up @@ -3749,12 +3752,12 @@ export class SandboxSession extends DurableObject<Env> {
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()) {
Expand Down Expand Up @@ -3804,6 +3807,16 @@ export class SandboxSession extends DurableObject<Env> {
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
Expand Down Expand Up @@ -3964,12 +3977,12 @@ export class SandboxSession extends DurableObject<Env> {
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;
Expand Down Expand Up @@ -4133,7 +4146,7 @@ export class SandboxSession extends DurableObject<Env> {
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);
Expand Down Expand Up @@ -4485,12 +4498,12 @@ export class SandboxSession extends DurableObject<Env> {
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);
Expand Down Expand Up @@ -4535,7 +4548,13 @@ export class SandboxSession extends DurableObject<Env> {
: 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';
Expand Down Expand Up @@ -4580,6 +4599,134 @@ export class SandboxSession extends DurableObject<Env> {
);
}

/**
* 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<boolean> {
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<void> {
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 => {
Comment thread
iscekic marked this conversation as resolved.
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,
Expand Down
7 changes: 5 additions & 2 deletions services/cloud-agent-next/src/sandbox-session/control-rpc.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand Down
Loading