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
20 changes: 13 additions & 7 deletions api/src/api/v2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import type { TFile } from '../job';
import { getLatestRuntimeMatchingLanguageVersion, getRuntimes } from '../runtime';
import { logger } from '../logger';
import { config } from '../config';
import { reconcileArtifactDelivery } from '../delivery';
import {
Job,
SessionWorkspaceDirtyError,
Expand Down Expand Up @@ -563,15 +564,20 @@ router.post('/execute', express.json({ limit: config.execute_body_limit }), asyn
return new Set<string>();
});

const generatedIds = new Set(job.getGeneratedFileIds());
const before = result.files.length;
result.files = result.files.filter(
f => !generatedIds.has(f.id) || uploaded.has(f.id),
const delivery = reconcileArtifactDelivery(
result.files,
job.getGeneratedFileIds(),
uploaded,
);
const dropped = before - result.files.length;
if (dropped > 0) {
result.files = delivery.files;
result.artifact_delivery = delivery.artifact_delivery;
if (delivery.artifact_delivery) {
logger.warn(
{ job: job.uuid, dropped, kept: result.files.length },
{
job: job.uuid,
dropped: delivery.artifact_delivery.failed,
kept: result.files.length,
},
'Pruned files from response because upload did not reach file_server',
);
}
Expand Down
73 changes: 73 additions & 0 deletions api/src/delivery.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
import { describe, expect, test } from 'bun:test';
import { reconcileArtifactDelivery } from './delivery';

describe('reconcileArtifactDelivery', () => {
test('leaves successful and inherited file references unchanged', () => {
const files = [
{ id: 'generated', name: 'result.txt' },
{ id: 'inherited', name: 'input.txt', inherited: true as const },
];

expect(
reconcileArtifactDelivery(
files,
['generated'],
new Set(['generated']),
),
).toEqual({ files });
});

test('reports a complete delivery failure without returning phantom references', () => {
const files = [
{ id: 'generated', name: 'result.txt' },
{ id: 'inherited', name: 'input.txt', inherited: true as const },
];

expect(
reconcileArtifactDelivery(files, ['generated'], new Set()),
).toEqual({
files: [{ id: 'inherited', name: 'input.txt', inherited: true }],
artifact_delivery: {
code: 'artifact_delivery_failed',
status: 'failed',
attempted: 1,
delivered: 0,
failed: 1,
},
});
});

test('reports partial delivery and counts only expected generated ids', () => {
const files = [
{ id: 'first', name: 'first.txt' },
{ id: 'second', name: 'second.txt' },
];

expect(
reconcileArtifactDelivery(
files,
['first', 'second'],
new Set(['first', 'unknown']),
),
).toEqual({
files: [{ id: 'first', name: 'first.txt' }],
artifact_delivery: {
code: 'artifact_delivery_failed',
status: 'partial',
attempted: 2,
delivered: 1,
failed: 1,
},
});
});

test('does not report a failure when there were no generated files', () => {
const files = [
{ id: 'inherited', name: 'input.txt', inherited: true as const },
];

expect(reconcileArtifactDelivery(files, [], new Set())).toEqual({
files,
});
});
});
43 changes: 43 additions & 0 deletions api/src/delivery.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
export interface ArtifactDeliveryFailure {
code: 'artifact_delivery_failed';
status: 'partial' | 'failed';
attempted: number;
delivered: number;
failed: number;
}

export interface ArtifactDeliveryResult<T> {
files: T[];
artifact_delivery?: ArtifactDeliveryFailure;
}

/** Removes unusable generated refs while preserving an explicit delivery failure for callers. */
export function reconcileArtifactDelivery<T extends { id: string }>(
files: T[],
generatedFileIds: Iterable<string>,
uploadedFileIds: ReadonlySet<string>,
): ArtifactDeliveryResult<T> {
const generatedIds = new Set(generatedFileIds);
if (generatedIds.size === 0) return { files };

let delivered = 0;
const retained = files.filter(file => {
if (!generatedIds.has(file.id)) return true;
if (!uploadedFileIds.has(file.id)) return false;
delivered++;
return true;
});
const failed = generatedIds.size - delivered;
if (failed === 0) return { files: retained };

return {
files: retained,
artifact_delivery: {
code: 'artifact_delivery_failed',
status: delivered === 0 ? 'failed' : 'partial',
attempted: generatedIds.size,
delivered,
failed,
},
};
}
2 changes: 2 additions & 0 deletions api/src/job.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import * as fsp from 'fs/promises';
import { pipeline } from 'stream/promises';
import { Readable, Transform } from 'stream';
import type { Logger } from 'pino';
import type { ArtifactDeliveryFailure } from './delivery';
import type { NsJailResult } from './nsjail';
import type { Runtime } from './runtime';
import { logger as rootLogger } from './logger';
Expand Down Expand Up @@ -676,6 +677,7 @@ interface ExecuteResult {
/** Top-level execution session id (one sandbox `/exec` invocation). */
session_id: string;
files: FileRef[];
artifact_delivery?: ArtifactDeliveryFailure;
}

const jobQueue: Array<() => void> = [];
Expand Down
17 changes: 16 additions & 1 deletion service/src/execution-log.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,14 @@ describe('execution log summaries', () => {
{ id: 'file_1', name: 'a.txt', inherited: true },
{ id: 'file_2', name: 'b.txt', modified_from: { id: 'file_1', storage_session_id: 'sess_old' } },
],
artifact_delivery: {
code: 'artifact_delivery_failed',
status: 'partial',
attempted: 3,
delivered: 2,
failed: 1,
detail: 'private storage failure',
},
run: {
code: 0,
stdout: 'top secret stdout',
Expand All @@ -33,9 +41,17 @@ describe('execution log summaries', () => {
expect(JSON.stringify(summary)).not.toContain('top secret stdout');
expect(JSON.stringify(summary)).not.toContain('sensitive stderr');
expect(JSON.stringify(summary)).not.toContain('combined output');
expect(JSON.stringify(summary)).not.toContain('private storage failure');
expect(summary).toMatchObject({
session_id: 'sess_123',
files: { count: 2, inheritedCount: 1, modifiedCount: 1 },
artifact_delivery: {
code: 'artifact_delivery_failed',
status: 'partial',
attempted: 3,
delivered: 2,
failed: 1,
},
run: {
stdout: { length: 17, present: true },
stderr: { length: 16, present: true },
Expand All @@ -56,4 +72,3 @@ describe('execution log summaries', () => {
expect(summary).toEqual({ count: 3, skillCount: 1, agentCount: 1, userCount: 1 });
});
});

21 changes: 20 additions & 1 deletion service/src/execution-log.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,9 +18,28 @@ type SandboxResponseLike = {
language?: unknown;
version?: unknown;
files?: unknown;
artifact_delivery?: unknown;
run?: RunLike;
};

function summarizeArtifactDelivery(value: unknown): Record<string, unknown> | undefined {
if (value == null || typeof value !== 'object' || Array.isArray(value)) return undefined;
const delivery = value as {
code?: unknown;
status?: unknown;
attempted?: unknown;
delivered?: unknown;
failed?: unknown;
};
return {
code: delivery.code,
status: delivery.status,
attempted: delivery.attempted,
delivered: delivery.delivered,
failed: delivery.failed,
};
}

export function summarizeText(value: unknown): { length: number; present: boolean } {
if (typeof value !== 'string') {
return { length: 0, present: false };
Expand Down Expand Up @@ -67,6 +86,7 @@ export function summarizeSandboxResponse(data: SandboxResponseLike): Record<stri
language: data.language,
version: data.version,
files: summarizeFiles(data.files),
artifact_delivery: summarizeArtifactDelivery(data.artifact_delivery),
run: run == null
? undefined
: {
Expand All @@ -83,4 +103,3 @@ export function summarizeSandboxResponse(data: SandboxResponseLike): Record<stri
},
};
}

66 changes: 66 additions & 0 deletions service/src/service/blocking-poll.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
import { describe, expect, test } from 'bun:test';
import { pollBlockingExecution, type BlockingPollDependencies } from './blocking-poll';
import type * as t from '../types';

const result: t.ExecuteResult = {
session_id: 'session', stdout: 'successful code', stderr: '', files: [],
artifact_delivery: {
code: 'artifact_delivery_failed', status: 'failed', attempted: 1, delivered: 0, failed: 1,
},
};

function fixture(): BlockingPollDependencies {
let tick = 0;
return {
getExecutionState: async () => ({ jobCompleted: tick >= 2 }),
getBlockingResult: async () => result,
getPending: async () => ({ status: 'completed' }),
isNotFound: (error: unknown) => error === 'missing',
sleep: async (): Promise<void> => { tick++; },
now: () => tick,
};
}

describe('blocking worker settlement', () => {
test('completed callback waits for delayed upload reconciliation', async () => {
const deps = fixture();
expect(await pollBlockingExecution('exec', 5, deps)).toEqual({
status: 'completed', stdout: result.stdout, stderr: '', files: [],
artifact_delivery: result.artifact_delivery,
});
expect(deps.now()).toBe(2);
});

test('missing callback session still receives the worker result', async () => {
const deps = fixture();
deps.getPending = async (): Promise<never> => { throw 'missing'; };
expect((await pollBlockingExecution('exec', 5, deps)).artifact_delivery).toEqual(result.artifact_delivery);
});

test('never reports success if uploads remain unsettled at timeout', async () => {
expect(await pollBlockingExecution('exec', 1, fixture())).toEqual({ status: 'error' });
});

test('accepts the legacy inline worker result during rolling deployments', async () => {
const deps = fixture();
expect(await pollBlockingExecution('exec', 5, {
...deps,
getExecutionState: async () => ({ jobCompleted: true, jobResult: result }),
getBlockingResult: async () => null,
})).toMatchObject({ status: 'completed', artifact_delivery: result.artifact_delivery });
});

test('preserves waiting calls and worker failures', async () => {
expect(await pollBlockingExecution('exec', 5, {
...fixture(),
getPending: async () => ({ status: 'waiting', pending_calls: [
{ call_id: 'call', tool_name: 'search', tool_input: { query: 'hello' } },
] }),
})).toEqual({ status: 'waiting', pending_calls: [
{ id: 'call', name: 'search', input: { query: 'hello' } },
] });
expect(await pollBlockingExecution('exec', 5, {
...fixture(), getExecutionState: async () => ({ jobError: 'worker failed' }),
})).toEqual({ status: 'error' });
});
});
76 changes: 76 additions & 0 deletions service/src/service/blocking-poll.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,76 @@
import type * as t from '../types';

export interface BlockingPendingState {
status: string;
pending_calls?: Array<{
call_id: string;
tool_name: string;
tool_input: Record<string, unknown>;
}>;
}

export interface BlockingPollDependencies {
getExecutionState(id: string): Promise<{
jobCompleted?: boolean;
jobResult?: t.ExecuteResult;
jobError?: string;
} | null>;
getBlockingResult(id: string): Promise<t.ExecuteResult | null>;
getPending(id: string): Promise<BlockingPendingState>;
isNotFound(error: unknown): boolean;
sleep(): Promise<void>;
now(): number;
}

/** Tool-call completion precedes upload reconciliation. Only a worker result is final. */
export async function pollBlockingExecution(
id: string,
timeout: number,
deps: BlockingPollDependencies,
): Promise<{
status: 'waiting' | 'completed' | 'error';
pending_calls?: t.ProgrammaticToolCall[];
stdout?: string;
stderr?: string;
files?: t.FileRefs;
artifact_delivery?: t.ArtifactDeliveryFailure;
}> {
const start = deps.now();
while (deps.now() - start < timeout) {
const execution = await deps.getExecutionState(id);
if (execution?.jobCompleted === true) {
// Preserve the inline result fallback for in-flight jobs from older binaries.
const result = (await deps.getBlockingResult(id)) ?? execution.jobResult;
if (result) {
return {
status: 'completed',
stdout: result.stdout,
stderr: result.stderr,
files: result.files,
artifact_delivery: result.artifact_delivery,
};
}
}
if (execution?.jobError != null) return { status: 'error' };

try {
const pending = await deps.getPending(id);
if (pending.status === 'waiting' && pending.pending_calls != null && pending.pending_calls.length > 0) {
return {
status: 'waiting',
pending_calls: pending.pending_calls.map(call => ({
id: call.call_id,
name: call.tool_name,
input: call.tool_input,
})),
};
}
if (pending.status === 'error') return { status: 'error' };
// Both completed and missing callback sessions must await the worker result.
} catch (error) {
if (!deps.isNotFound(error)) throw error;
}
await deps.sleep();
}
return { status: 'error' };
}
Loading