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
10 changes: 8 additions & 2 deletions services/platform/backend/core/automations/agent_retry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,8 +30,8 @@ export type WorkflowAgentFailureCode =
| 'turn_stalled'
/** The sandbox ran out of memory and the kernel's OOM killer ended the
* harness or its session. Re-kicked like any failure, but only after
* `resourceExhaustedRetryDelayMs` for the node's attempt: at once it
* would meet the same limit. */
* {@link RESOURCE_EXHAUSTED_REKICK_DELAY_MS}: at once it would meet the
* same limit. */
| 'resource_exhausted'
| 'session_gone'
| 'start_failed'
Expand Down Expand Up @@ -76,6 +76,12 @@ export const SANDBOX_ROOM_MAX_WAIT_MS = 2 * 60 * 60_000;
* past it, a node still waiting asks about once a minute on average. */
export const SANDBOX_ROOM_RETRY_CEILING_MS = 2 * 60_000;

/** How long the re-kick of a node whose sandbox ran out of memory is held:
* as long as the kick may hold a start (its op row must stay inside the
* stalled-turn sweep's window), so the 10 and 30 minutes a task run waits
* are not available here. */
export const RESOURCE_EXHAUSTED_REKICK_DELAY_MS = SANDBOX_ROOM_RETRY_CEILING_MS;

/** The most a start that holds a place in the spawner's line comes back
* after its hint: enough to keep waiters refused together apart. */
const QUEUED_RETRY_JITTER_MS = 1_000;
Expand Down
8 changes: 4 additions & 4 deletions services/platform/backend/core/automations/stepper.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,13 +33,13 @@ import { harnessResumesConversations } from '../chat/external_turn_shared';
import type { ActionCtx } from '../lib/ctx';
import { internal } from '../lib/handler_names';
import type { Id } from '../lib/rows';
import { resourceExhaustedRetryDelayMs } from '../tasks/task_auto_retry';
import {
automationAgentHost,
type AutomationAgentHost,
type WorkflowAgentRequest,
} from './agent_host';
import {
RESOURCE_EXHAUSTED_REKICK_DELAY_MS,
SANDBOX_ROOM_MAX_WAIT_MS,
isWorkflowAgentRetryable,
planWorkflowAgentRetry,
Expand Down Expand Up @@ -1720,8 +1720,8 @@ async function stepAgentNode(args: AgentStepArgs): Promise<StepOutcome> {
// more with each refusal in a row; any other refusal with a hint (a
// broker pool cooling down) waits for exactly that.
const now = Date.now();
// One whose sandbox ran out of memory waits 2, 10, then 30 minutes:
// at once it would meet the same limit.
// One whose sandbox ran out of memory waits as long as a start may be
// held (two minutes): at once it would meet the same limit.
const notBefore = waitingForRoom
? sandboxRoomRetryAtMs({
now,
Expand All @@ -1732,7 +1732,7 @@ async function stepAgentNode(args: AgentStepArgs): Promise<StepOutcome> {
queued: settled.roomQueued === true,
})
: settled.failureCode === 'resource_exhausted'
? now + resourceExhaustedRetryDelayMs(parked.attempt)
? now + RESOURCE_EXHAUSTED_REKICK_DELAY_MS
: settled.retryAtMs;
const kicked = await run.agent.kick({
runId: run.runId,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1136,7 +1136,7 @@ const OUT_OF_MEMORY_CODES: ReadonlySet<string> = new Set([
/** The reason a turn settles failed with when its sandbox ran out of
* memory. */
export const OUT_OF_MEMORY_TURN_REASON =
"The agent's sandbox ran out of memory: the kernel's OOM killer ended the agent. A retry follows after a pause; if it keeps happening, the agent sessions need a larger memory limit (SANDBOX_AGENT_MEMORY).";
"The agent's sandbox ran out of memory: the kernel's OOM killer ended the agent. If it keeps happening, the agent sessions need a larger memory limit (SANDBOX_AGENT_MEMORY).";

/** The reason a turn settles failed with when the sandbox ended its harness
* as stalled. */
Expand Down
34 changes: 30 additions & 4 deletions services/platform/backend/core/tasks/task_input_mirrors.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -146,7 +146,7 @@ describe('the start’s pass over the worker’s copies of task inputs', () => {
);
});

it('looks at the oldest copies first, a bounded number per start', async () => {
it('removes the oldest stale copies first, a bounded number per start', async () => {
io.dirs['/agent/inputs'] = Array.from(
{ length: MAX_INPUT_MIRRORS_PER_PASS + 10 },
(_, n) => dir(`task-${n}`, 10_000 - n),
Expand All @@ -158,16 +158,42 @@ describe('the start’s pass over the worker’s copies of task inputs', () => {

await pruneStaleTaskInputMirrors(ctx, ARGS);

// Every copy is asked about, the least recently changed first.
const checked = asked[0]?.taskIds as string[];
expect(checked).toHaveLength(MAX_INPUT_MIRRORS_PER_PASS);
// The least recently changed first: task-59 down to task-10.
expect(checked).toHaveLength(MAX_INPUT_MIRRORS_PER_PASS + 10);
expect(checked[0]).toBe(`task-${MAX_INPUT_MIRRORS_PER_PASS + 9}`);
expect(checked).not.toContain('task-0');
// A bounded number goes: the oldest, task-59 down to task-10.
expect(io.deletes[0]).toHaveLength(MAX_INPUT_MIRRORS_PER_PASS);
expect(io.deletes[0]?.[0]).toBe(
`/agent/inputs/task-${MAX_INPUT_MIRRORS_PER_PASS + 9}`,
);
expect(io.deletes[0]).not.toContain('/agent/inputs/task-0');
// No reviews directory, so none is listed.
expect(io.listings).toEqual(['/agent/inputs']);
});

it('old copies that are still needed never hide the stale ones behind them', async () => {
// The oldest 60 copies belong to tasks still open; the 5 newer ones are
// stale.
io.dirs['/agent/inputs'] = Array.from({ length: 65 }, (_, n) =>
dir(`task-${n}`, n < 60 ? n : 10_000 + n),
);
const { ctx } = makeCtx(({ taskIds }) => ({
taskIds: taskIds.filter((id) => Number(id.slice('task-'.length)) >= 60),
reviewHashes: [],
}));

await pruneStaleTaskInputMirrors(ctx, ARGS);

expect(io.deletes[0]).toEqual([
'/agent/inputs/task-60',
'/agent/inputs/task-61',
'/agent/inputs/task-62',
'/agent/inputs/task-63',
'/agent/inputs/task-64',
]);
});

it('asks nothing and removes nothing when the worker holds no other copy', async () => {
io.dirs['/agent/inputs'] = [dir('task-current')];
const { ctx, asked } = makeCtx(() => ({ taskIds: [], reviewHashes: [] }));
Expand Down
21 changes: 16 additions & 5 deletions services/platform/backend/core/tasks/task_input_mirrors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,11 +31,17 @@ const REVIEWS_DIR_NAME = 'reviews';

const REVIEW_INPUTS_ROOT = `${TASK_INPUTS_ROOT}/${REVIEWS_DIR_NAME}`;

/** The most copies of each kind one pass looks at, the oldest first, so a
/** The most copies of each kind one pass removes, the oldest first, so a
* worker that gathered many is cleared over several starts rather than
* holding up one. */
export const MAX_INPUT_MIRRORS_PER_PASS = 50;

/** The most copies of each kind one pass asks about, the oldest first: far
* more than it removes, so copies that are still needed — the oldest are
* often long-open tasks — never hide the stale ones behind them. One
* indexed lookup answers them all. */
const MAX_INPUT_MIRROR_CANDIDATES = 1000;

/** A directory name a task id can be; anything else under the root is left
* alone, since nothing the platform stages is named so. */
const TASK_DIR_NAME_RE = /^[A-Za-z0-9_-]{1,128}$/;
Expand All @@ -52,7 +58,7 @@ export function reviewInputsDir(taskId: string): string {
return `${REVIEW_INPUTS_ROOT}/${hash}`;
}

/** The names of up to {@link MAX_INPUT_MIRRORS_PER_PASS} directories that
/** The names of up to {@link MAX_INPUT_MIRROR_CANDIDATES} directories that
* `accept` takes, the least recently changed first. */
function oldestDirs(
entries: readonly SessionFsEntry[],
Expand All @@ -61,7 +67,7 @@ function oldestDirs(
return entries
.filter((entry) => entry.type === 'dir' && accept(entry.name))
.sort((a, b) => a.mtimeMs - b.mtimeMs)
.slice(0, MAX_INPUT_MIRRORS_PER_PASS)
.slice(0, MAX_INPUT_MIRROR_CANDIDATES)
.map((entry) => entry.name);
}

Expand Down Expand Up @@ -119,9 +125,14 @@ export async function pruneStaleTaskInputMirrors(
taskIds,
reviewHashes,
});
// The answer keeps the candidates' order, oldest first.
const paths = [
...stale.taskIds.map((id) => `${TASK_INPUTS_ROOT}/${id}`),
...stale.reviewHashes.map((hash) => `${REVIEW_INPUTS_ROOT}/${hash}`),
...stale.taskIds
.slice(0, MAX_INPUT_MIRRORS_PER_PASS)
.map((id) => `${TASK_INPUTS_ROOT}/${id}`),
...stale.reviewHashes
.slice(0, MAX_INPUT_MIRRORS_PER_PASS)
.map((hash) => `${REVIEW_INPUTS_ROOT}/${hash}`),
];
if (paths.length === 0) return;
const removed = await sessionDeleteFiles(args.sessionId, paths);
Expand Down
42 changes: 42 additions & 0 deletions services/platform/backend/domains/tasks/agent-runs.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -315,6 +315,48 @@ describe('the turn host’s terminal marks write the provenance entry', () => {
);
});

it.each([
[undefined, 2 * 60_000],
[1, 10 * 60_000],
[2, 30 * 60_000],
[5, 30 * 60_000],
])(
'holds the retry decision of a run out of memory (attempt %p) for %p ms',
async (autoRetryAttempt, waitMs) => {
const { sql } = fakeSql((text) =>
text.startsWith('UPDATE app.project_agent_runs')
? [
{
organizationId: 'org-1',
taskId: 'task-1',
agentId: 'agent-1',
autoRetryAttempt: autoRetryAttempt ?? null,
},
]
: [],
);
await failAgentRunFromTurn(sql, {
runId: 'run-1',
execId: 'exec-1',
error: "the agent's sandbox ran out of memory",
failureCode: 'resource_exhausted',
});
// The job itself waits: no queued run sits out the wait for the
// stranded-queued-run sweep to start early.
expect(addJobInTx).toHaveBeenCalledExactlyOnceWith(
expect.anything(),
'task.agent_retry',
{
organizationId: 'org-1',
taskId: 'task-1',
agentId: 'agent-1',
expectedRunId: 'run-1',
},
{ startAfter: new Date(NOW + waitMs) },
);
},
);

it('starts the model-capacity floor after a terminal-update lock wait, not before it', async () => {
const { sql } = fakeSql((text) => {
if (!text.startsWith('UPDATE app.project_agent_runs')) return [];
Expand Down
26 changes: 19 additions & 7 deletions services/platform/backend/domains/tasks/agent-runs.ts
Original file line number Diff line number Diff line change
Expand Up @@ -706,18 +706,30 @@ export async function failAgentRunFromTurn(
const startAfterMs =
args.failureCode === 'model_capacity'
? Date.now() + MODEL_CAPACITY_RETRY_DELAY_MS
: args.failureCode === 'resource_exhausted'
? Date.now() + resourceExhaustedRetryDelayMs(run.autoRetryAttempt)
: args.retryAtMs !== undefined && args.retryAtMs > now
? Math.min(args.retryAtMs, now + BROKER_RATE_LIMIT_COOLDOWN_MS)
: undefined;
await addJobInTx(tx, 'task.agent_retry', {
: args.retryAtMs !== undefined && args.retryAtMs > now
? Math.min(args.retryAtMs, now + BROKER_RATE_LIMIT_COOLDOWN_MS)
: undefined;
// A run its sandbox's memory limit ended waits 2, 10, then 30
// minutes — longer than the stranded-queued-run sweep lets a queued
// run wait, so the retry DECISION is held instead: no queued run
// exists meanwhile, the failed run shows its retry pending, and every
// guard is re-derived when the wait ends.
const decideAfter =
args.failureCode === 'resource_exhausted'
? new Date(
Date.now() + resourceExhaustedRetryDelayMs(run.autoRetryAttempt),
)
: undefined;
const retry = {
organizationId: run.organizationId,
taskId: run.taskId,
agentId: run.agentId,
expectedRunId: args.runId,
...(startAfterMs !== undefined && { startAfterMs }),
});
};
await (decideAfter !== undefined
? addJobInTx(tx, 'task.agent_retry', retry, { startAfter: decideAfter })
: addJobInTx(tx, 'task.agent_retry', retry));
} else {
await announceAgentRunFailed(tx, {
organizationId: run.organizationId,
Expand Down
2 changes: 1 addition & 1 deletion services/platform/messages/de/tasks.yml
Original file line number Diff line number Diff line change
Expand Up @@ -121,7 +121,7 @@ agentRun:
model: Das KI-Modell hinter diesem Agenten ist ausgefallen, bevor die Arbeit erledigt war. Starte den Agenten erneut — schlägt er wieder fehl, bitte einen Admin, den KI-Anbieter zu prüfen.
start: 'Der Lauf konnte nicht starten. Versuche es erneut — schlägt er wieder fehl, zeig einem Admin, was der Lauf unter "Details" gemeldet hat.'
interrupted: Der Lauf wurde unterbrochen, bevor der Agent fertig war. Starte den Agenten erneut.
out_of_memory: Der Sandbox des Agenten ist der Arbeitsspeicher ausgegangen, deshalb wurde seine Arbeit gestoppt. Tale versucht es nach einer Pause erneut — passiert das wieder, bitte einen Admin, Agenten mehr Arbeitsspeicher zu geben.
out_of_memory: Der Sandbox des Agenten ist der Arbeitsspeicher ausgegangen, deshalb wurde seine Arbeit gestoppt. Starte den Agenten erneut — passiert das wieder, bitte einen Admin, Agenten mehr Arbeitsspeicher zu geben.
stalled: 'Der Agent hat nicht mehr reagiert: Er hat lange nichts ausgegeben und kaum gearbeitet, deshalb hat seine Sandbox ihn gestoppt. Starte den Agenten erneut — bleibt er wieder stehen, zeig einem Admin, was der Lauf unter "Details" gemeldet hat.'
unknown: 'Der Lauf ist fehlgeschlagen. Starte den Agenten erneut — schlägt er wieder fehl, zeig einem Admin, was der Lauf unter "Details" gemeldet hat.'
standardAgent:
Expand Down
2 changes: 1 addition & 1 deletion services/platform/messages/en/tasks.yml
Original file line number Diff line number Diff line change
Expand Up @@ -123,7 +123,7 @@ agentRun:
model: The AI model behind this agent failed before the work was done. Start the agent again — if it fails again, ask an Admin to check the AI provider.
start: The run couldn't start. Try again — if it fails again, show an Admin what the run reported under Details.
interrupted: The run was interrupted before the agent finished. Start the agent again.
out_of_memory: The agent's sandbox ran out of memory, so its work was stopped. Tale tries again after a pause — if it keeps happening, ask an Admin to give agents more memory.
out_of_memory: The agent's sandbox ran out of memory, so its work was stopped. Start the agent again — if it keeps happening, ask an Admin to give agents more memory.
stalled: The agent stopped responding — it printed nothing and did almost no work for a long time — so its sandbox stopped it. Start the agent again; if it stops again, show an Admin what the run reported under Details.
unknown: The run failed. Start the agent again — if it fails again, show an Admin what the run reported under Details.
standardAgent:
Expand Down
2 changes: 1 addition & 1 deletion services/platform/messages/fr/tasks.yml
Original file line number Diff line number Diff line change
Expand Up @@ -122,7 +122,7 @@ agentRun:
model: Le modèle d’IA derrière cet agent a échoué avant la fin du travail. Relance l’agent — s’il échoue encore, demande à un admin de vérifier le fournisseur d’IA.
start: L’exécution n’a pas pu démarrer. Réessaie — si elle échoue encore, montre à un admin ce que l’exécution a signalé dans « Détails ».
interrupted: L’exécution a été interrompue avant que l’agent ait terminé. Relance l’agent.
out_of_memory: La sandbox de l’agent a manqué de mémoire, son travail a donc été arrêté. Tale réessaie après une pause — si cela se reproduit, demande à un admin de donner plus de mémoire aux agents.
out_of_memory: La sandbox de l’agent a manqué de mémoire, son travail a donc été arrêté. Relance l’agent — si cela se reproduit, demande à un admin de donner plus de mémoire aux agents.
stalled: L’agent ne répondait plus — il n’a rien affiché et n’a presque rien fait pendant longtemps —, sa sandbox l’a donc arrêté. Relance l’agent — s’il s’arrête encore, montre à un admin ce que l’exécution a signalé dans « Détails ».
unknown: L’exécution a échoué. Relance l’agent — si elle échoue encore, montre à un admin ce que l’exécution a signalé dans « Détails ».
standardAgent:
Expand Down
Loading
Loading