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
2 changes: 1 addition & 1 deletion services/cloud-agent-next/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -127,7 +127,7 @@ This pattern blocks API endpoints from running for external contributors who don
- New pending-message writes use a versioned record containing one nested immutable `SessionMessageIntent`, delivery retry state, and callback snapshot; flat pending rows are decoded only inside `src/session/pending-messages.ts`.
- `SessionMessageState` owns lifecycle/outbox status, terminal effect accounting, and a named immutable `admissionSnapshot` only for post-pending replay validation and recovery; predecessor records normalize into partial `legacyAdmissionConstraints` and never fabricate missing immutable input. Terminal and accepted/sent effects are repairable from pending/alarm replay and events use deterministic uniqueness.
- Wrapper handoff is currently at-least-once under ambiguous delivery failures: the wrapper forwards prompt/command submissions directly to Kilo and does not query Kilo to suppress or recover duplicate `messageId` submissions. Duplicate prompt/command processing is an accepted edge-case trade-off until Kilo provides an atomic submit-or-return-existing contract.
- When accepted work has no pending residue and its fenced wrapper runtime/socket is gone, disconnect or liveness expiry first reconciles each accepted message against the DO's stored kilocode events (`getAssistantMessageForUserMessage`): positive terminal evidence (assistant `time.completed` or a terminal assistant error) settles the message as `idle_reconciliation`; anything else terminalizes as wrapper failure without redispatch. There is still no live authoritative Kilo terminal query for redispatch; adding one remains separate lifecycle capability work.
- When accepted work has no pending residue and its fenced wrapper runtime/socket is gone, disconnect or liveness expiry first reconciles each accepted message against the DO's stored kilocode events (`getAssistantMessageForUserMessage`): positive terminal evidence (assistant `time.completed` or a terminal assistant error) settles the message as `idle_reconciliation`; anything else terminalizes as wrapper failure. A detection of no output on an accepted turn is recoverable once: the accepted-message inactivity path here and the wrapper no-output watchdog in the legacy plane each re-dispatch the same durable intent on a fresh runtime, and only a second identical detection terminalizes as wrapper failure with the attempt count recorded in the failure payload (`attempts`). There is still no live authoritative Kilo terminal query; adding one remains separate lifecycle capability work.
- Legacy wrapper-supervisor physical cleanup exhaustion (`WRAPPER_STOP_MAX_ATTEMPTS` reached) is fenced but recoverable: the lease re-observes the sandbox on a slow cadence (`WRAPPER_CLEANUP_EXHAUSTED_RECHECK_MS`) and releases to `none` only after a confirmed `absent` observation, so a wedged sandbox reaped later by the container runtime does not brick the session. Recovery is observation-only (`observeWrappersWithoutWaking`) and never issues another stop: the attempt budget and its rollback fence still hold, and the probe must not wake a stopped container to ask about a process that cannot outlive it. Background rechecks stop after `WRAPPER_CLEANUP_EXHAUSTED_RECHECK_WINDOW_MS` so an unrecoverable exhaustion stops re-arming the DO alarm; explicit sends still force one probe afterwards.
- A legacy pending-message flush blocked on exhaustion forces one out-of-cadence recheck (`recoverExhaustedDeliveryBlock` → `recheckExhaustedCleanup`) because a user is actively waiting, then retries on the `WRAPPER_CLEANUP_EXHAUSTED` budget before failing closed. The retry budget is what keeps the two halves consistent — recovery takes minutes, so terminalizing on the first blocked attempt would discard messages a later probe would have delivered — and failing closed at the end of it is what keeps a message from sitting `queued` with no terminal signal. The flush failure code must stay authoritative: `INTERNAL` is treated as non-authoritative by `recordPendingFlushFailure` and would terminalize the message under whatever earlier cause it carried.
- Callback delivery retry policy is paired with `wrangler.jsonc`: `CALLBACK_DELIVERY_MAX_ATTEMPTS` includes the initial attempt, and each Cloud Agent Next callback queue consumer must configure `max_retries` for the remaining redeliveries.
Expand Down
113 changes: 108 additions & 5 deletions services/cloud-agent-next/src/sandbox-session/SandboxSession.ts
Original file line number Diff line number Diff line change
Expand Up @@ -236,9 +236,11 @@ import {
markSessionOperationRejection,
matchesSessionMessageReplay,
nextQueuedMessageId,
noOutputRecoveryAllowed,
providerOwnershipOf,
queuedAtOf,
recordAcceptedMessageActivity,
redispatchAcceptedMessage,
releaseCompletedRetryableAttach,
releaseUnadmittedWaitingMessages,
replacePreparationAttemptId,
Expand All @@ -256,6 +258,11 @@ import { bindingForAttachment } from './session-binding.js';
import { decideSession, terminalMessageState } from '../sandbox-state/session/reduce.js';
import type { Binding } from '../sandbox-state/model/session.js';
import { createMessageCallbacks, type MessageCallbacks } from './message-callbacks.js';
import {
applyStoredAssistantSettlement,
projectStoredAssistantSettlement,
type AcceptedStoredSettlement,
} from './accepted-stored-settlement.js';
import {
createReportOutbox,
readReportAnchor,
Expand Down Expand Up @@ -3236,6 +3243,11 @@ export class SandboxSession extends DurableObject<Env> {
report('inactivity');
return;
}
if (!this.acceptedTurnIsDrivable(accepted.messageId)) {
diagnostic.reason = 'accepted_message_changed';
report('superseded');
return;
}
diagnostic.healthy = true;
report('healthy');
await this.scheduleAcceptedRecheck(epoch, accepted.messageId);
Expand Down Expand Up @@ -3268,14 +3280,22 @@ export class SandboxSession extends DurableObject<Env> {
report('inactivity');
return;
}
if (!this.acceptedTurnIsDrivable(accepted.messageId)) {
diagnostic.reason = 'accepted_message_changed';
report('superseded');
return;
}
if (!waiting) {
diagnostic.reason = 'inactive_snapshot';
throw new Error('Accepted execution is no longer active');
}
report('healthy');
await this.scheduleAcceptedRecheck(epoch, accepted.messageId);
} catch {
if (this.isCurrentAcceptedMessage(accepted, epoch)) {
if (
this.isCurrentAcceptedMessage(accepted, epoch) &&
this.acceptedTurnIsDrivable(accepted.messageId)
) {
diagnostic.reason ??= 'sync_failed';
report('runtime_unhealthy');
await this.failDelivery(
Expand All @@ -3299,10 +3319,37 @@ export class SandboxSession extends DurableObject<Env> {
}

/**
* The inactivity fail path. Re-reads the accepted row and re-evaluates the
* bound synchronously before persisting, so real progress that landed during
* the preceding observe/sync keeps the turn alive. Message-only: no
* quarantine and no native-runtime retirement.
* Positive terminal evidence for an accepted turn that already answered over
* the ingest channel, read from the DO's stored kilocode events. The record
* stays `accepted` because settlement here comes from the prompt operation
* result, so recovery must settle (not re-run) a turn whose answer already
* arrived. Mirrors the legacy wrapper-death reconciliation.
*/
private storedAssistantSettlement(current: MessageRecord): AcceptedStoredSettlement | undefined {
const metadata = this.terminalLifecycle.getStoredMetadata();
const kiloSessionId = metadata?.auth.kiloSessionId;
if (!metadata || !kiloSessionId) return undefined;
return projectStoredAssistantSettlement(
this.eventQueries.getAssistantMessageForUserMessage(
metadata.identity.sessionId,
kiloSessionId,
current.messageId
)
);
}

/**
* The inactivity path. Re-reads the accepted row and re-evaluates the bound
* synchronously before persisting, so real progress that landed during the
* preceding observe/sync keeps the turn alive.
*
* The first no-output detection is recoverable: the accepted turn is
* re-queued under its durable intent, so the alarm re-dispatches it once on a
* replacement runtime (the wedged runtime is left to the health machine's own
* recovery ladder). A second identical detection keeps the terminal path
* (message-only: no quarantine and no native-runtime retirement) and records
* the attempt count. The follow-up abort is best-effort and fenced with the
* captured wrapper instance.
*/
private async failOverdueAcceptedMessage(
accepted: MessageRecord,
Expand All @@ -3315,13 +3362,54 @@ export class SandboxSession extends DurableObject<Env> {
if (
!current ||
!currentState ||
// A Stop admitted while this alarm awaited its observation keeps the
// record `accepted` and only adds `cancellation`. Recovering it would
// re-queue the turn the user just stopped.
current.cancellation !== undefined ||
activityAt === undefined ||
!acceptedInactivityDue(activityAt, Date.now())
)
return false;
diagnostic.stage = 'inactivity';
diagnostic.reason = 'inactivity_due';
diagnostic.lastActivityAt = activityAt;
// Reconcile before recovering: settlement here normally comes from the
// prompt operation result, so a runtime that went silent after streaming
// its answer still leaves the record `accepted`. Re-dispatching it would run
// the turn and its side effects a second time.
const settlement = this.storedAssistantSettlement(current);
if (settlement) {
const settled = applyStoredAssistantSettlement(
this.loadMessages(),
current.messageId,
settlement,
Date.now()
);
if (!settled || !this.saveMessages(settled, epoch)) return false;
diagnostic.recovery = 'reconciled';
logControlDiagnostic('accepted_no_output_reconciled', {
...diagnostic,
producer: 'accepted_inactivity',
result: settlement.state,
});
return true;
}
const recoveryAttempts = currentState.recoveryAttempts ?? 0;
if (noOutputRecoveryAllowed(recoveryAttempts)) {
Comment thread
iscekic marked this conversation as resolved.
const next = redispatchAcceptedMessage(this.loadMessages(), current.messageId);
Comment thread
iscekic marked this conversation as resolved.
if (!next || !this.saveMessages(next, epoch)) return false;
diagnostic.recovery = 'redispatched';
logControlDiagnostic('accepted_no_output_recovery', {
...diagnostic,
producer: 'accepted_inactivity',
recoveryAttempts: recoveryAttempts + 1,
});
const turn = getSessionMessageTurn(current);
this.broadcastQueuedMessage(current.messageId, turn ? renderExecutionTurnContent(turn) : '');
await this.armQueueRetry();
return true;
}
diagnostic.attempts = recoveryAttempts + 1;
const failed = failAcceptedMessage(
this.sessionAggregate(this.loadMessages()),
current.messageId,
Expand All @@ -3333,6 +3421,20 @@ export class SandboxSession extends DurableObject<Env> {
return true;
}

/**
* Whether the accepted turn is still the watchdog's to drive. A Stop admitted
* during an awaited observation keeps the record `accepted` and only adds
* `cancellation`; the stop lifecycle settles the turn from there, so the
* watchdog must release it instead of re-arming it or failing it as a runtime
* failure. `failOverdueAcceptedMessage` already refuses to recover such a
* record, which is why the callers re-check here before treating "not overdue"
* as healthy.
*/
private acceptedTurnIsDrivable(messageId: string): boolean {
const current = this.loadMessages().find(item => item.messageId === messageId);
return current?.state.kind === 'accepted' && current.cancellation === undefined;
}

/**
* Schedule the next non-failing check from the freshly read clock:
* `min(now + acceptedAlarmCap, activityAt + kiloInactivity)`. Never rearm at
Expand Down Expand Up @@ -3366,6 +3468,7 @@ export class SandboxSession extends DurableObject<Env> {
if (!this.terminalLifecycle.isCurrent(epoch)) return;
if (this.isCurrentAcceptedMessage(accepted, epoch)) {
if (await this.failOverdueAcceptedMessage(accepted, epoch, diagnostic)) return;
if (!this.acceptedTurnIsDrivable(accepted.messageId)) return;
await this.scheduleAcceptedRecheck(epoch, accepted.messageId);
return;
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,123 @@
import { describe, expect, it } from 'vitest';
import type { LatestAssistantMessage } from '../session/types.js';
import type { SessionMessage, SessionMessageIntent } from '../sandbox-state/model/session.js';
import {
applyStoredAssistantSettlement,
projectStoredAssistantSettlement,
} from './accepted-stored-settlement.js';

function assistant(
info: Record<string, unknown>,
parts: LatestAssistantMessage['parts'] = []
): LatestAssistantMessage {
return {
eventId: 1,
timestamp: 2,
info: { id: 'ase_1', role: 'assistant', ...info } as LatestAssistantMessage['info'],
parts,
};
}

const intent = (messageId: string): SessionMessageIntent => ({
turn: { type: 'prompt', messageId, prompt: 'hello' },
agent: { mode: 'build' },
});

function accepted(messageId = 'msg_1'): SessionMessage {
return {
messageId,
state: { kind: 'accepted', intent: intent(messageId), acceptedAt: 5, executionDeadlineAt: 500 },
};
}

describe('projectStoredAssistantSettlement', () => {
it('settles a stored assistant answer with a completion marker', () => {
expect(
projectStoredAssistantSettlement(assistant({ time: { created: 1, completed: 2 } }))
).toEqual({ state: 'completed' });
});

it('settles a terminal assistant error with its bounded classification', () => {
expect(
projectStoredAssistantSettlement(assistant({ error: 'Provider request timed out' }))
).toMatchObject({
state: 'failed',
failedReason: 'assistant_error',
assistantReason: 'timeout',
});
});

it('settles an interrupted assistant error as a cancellation', () => {
expect(
projectStoredAssistantSettlement(
assistant({ error: { name: 'MessageAbortedError', message: 'interrupted' } })
)
).toEqual({ state: 'cancelled', failedReason: 'interrupted' });
});

it('does not settle a partial answer with no terminal evidence', () => {
expect(projectStoredAssistantSettlement(assistant({ time: { created: 1 } }))).toBeUndefined();
expect(projectStoredAssistantSettlement(null)).toBeUndefined();
});
});

describe('applyStoredAssistantSettlement', () => {
it('only settles a record that is still accepted', () => {
const messages = [accepted()];

expect(
applyStoredAssistantSettlement(messages, 'msg_1', { state: 'completed' }, 9)?.[0]
).toMatchObject({ state: { kind: 'completed', at: 9, source: 'coordinator' } });
expect(
applyStoredAssistantSettlement(messages, 'msg_1', { state: 'completed' }, 9)
).toHaveLength(1);
expect(
applyStoredAssistantSettlement(
[
{
messageId: 'msg_1',
state: {
kind: 'queued',
intent: intent('msg_1'),
deliveryStep: 'waiting',
deadlineAt: null,
attachFailures: 0,
promptFailures: 0,
},
},
],
'msg_1',
{ state: 'completed' },
9
)
).toBeUndefined();
expect(applyStoredAssistantSettlement([], 'msg_1', { state: 'completed' }, 9)).toBeUndefined();
});

it('carries the assistant failure facts onto the terminal record', () => {
const [settled] =
applyStoredAssistantSettlement(
[accepted()],
'msg_1',
{
state: 'failed',
failedReason: 'assistant_error',
failedDetail: 'Assistant request timed out',
assistantReason: 'timeout',
providerOwnership: 'unknown',
},
9
) ?? [];

expect(settled).toMatchObject({
state: {
kind: 'failed',
reason: 'assistant_error',
detail: 'Assistant request timed out',
assistantReason: 'timeout',
providerOwnership: 'unknown',
source: 'coordinator',
},
});
});
});
Loading
Loading