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
33 changes: 33 additions & 0 deletions packages/api/src/agents/background.shutdown.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -194,6 +194,39 @@ describe('background task shutdown', () => {
expect(summary).toEqual({ tracked: 1, interrupted: 1, flushed: 1, unsettled: 1 });
});

it('does not retain an unconfirmed result after its task expires', async () => {
const registry = new BackgroundTaskRegistryClass();
const task = createTask(registry, 'unconfirmed-expired');
const controlled = controlledHandle();
registry.trackShutdown(task, controlled.handle);
registry.fail('shutdown-owner', 'shutdown-conversation', task.id, 'no durable write');
const now = jest.spyOn(Date, 'now').mockReturnValue(task.updatedAt + 60 * 60 * 1000 + 1);
try {
expect(registry.get('shutdown-owner', 'shutdown-conversation', task.id)).toBeUndefined();
} finally {
now.mockRestore();
}

expect((await registry.drainForShutdown(drainOptions())).tracked).toBe(0);
});

it('does not retain an unconfirmed result after capacity pressure evicts its task', async () => {
const registry = new BackgroundTaskRegistryClass();
const task = createTask(registry, 'unconfirmed-evicted');
const controlled = controlledHandle();
registry.trackShutdown(task, controlled.handle);
registry.fail('shutdown-owner', 'shutdown-conversation', task.id, 'no durable write');
task.updatedAt = Date.now() - 1_000;
for (let i = 0; i < 199; i++) {
const settled = createTask(registry, `completed-${i}`);
registry.complete('shutdown-owner', 'shutdown-conversation', settled.id, { content: 'ok' });
}
createTask(registry, 'trigger-eviction');

expect(registry.get('shutdown-owner', 'shutdown-conversation', task.id)).toBeUndefined();
expect((await registry.drainForShutdown(drainOptions())).tracked).toBe(0);
});

it('stops tracking a task once its result is durable', async () => {
const registry = new BackgroundTaskRegistryClass();
const controlled = controlledHandle();
Expand Down
11 changes: 11 additions & 0 deletions packages/api/src/agents/background.ts
Original file line number Diff line number Diff line change
Expand Up @@ -887,6 +887,12 @@ export class BackgroundTaskRegistryClass {
return `${userId}::${conversationId}`;
}

private releaseEvictedShutdownHandle(task: BackgroundTask): void {
if (this.shutdownHandles.delete(task)) {
logger.warn(`[background] Unconfirmed shutdown result for evicted task ${task.id}.`);
}
}

private sweepBucketTasks(bucket: TaskBucket, now: number): void {
for (const [taskId, task] of bucket.tasks) {
if (
Expand All @@ -895,6 +901,7 @@ export class BackgroundTaskRegistryClass {
now - task.updatedAt > COMPLETED_TASK_TTL_MS
) {
bucket.tasks.delete(taskId);
this.releaseEvictedShutdownHandle(task);
}
}
/** Drop dedupe mappings whose task was evicted (keys are
Expand Down Expand Up @@ -925,6 +932,9 @@ export class BackgroundTaskRegistryClass {
(task) => task.status === 'running' || task.completionPersistencePending === true,
)
) {
for (const task of bucket.tasks.values()) {
this.releaseEvictedShutdownHandle(task);
}
this.buckets.delete(bucketKey);
continue;
}
Expand Down Expand Up @@ -1079,6 +1089,7 @@ export class BackgroundTaskRegistryClass {
const touched = new Set<TaskBucket>();
for (const [task, bucket] of selected) {
bucket.tasks.delete(task.id);
this.releaseEvictedShutdownHandle(task);
touched.add(bucket);
}
const now = Date.now();
Expand Down
215 changes: 215 additions & 0 deletions packages/api/src/agents/handlers.shutdown.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -234,6 +234,221 @@ describe('createToolExecuteHandler — background tasks at shutdown', () => {
);
});

it('waits for a legacy pre-admitted completion projection before declaring shutdown durable', async () => {
const { createToolExecuteHandler, registry } = loadModules();
let notifyFlush: () => void = () => undefined;
const flushStarted = new Promise<void>((resolve) => {
notifyFlush = resolve;
});
const trackShutdown = registry.trackShutdown.bind(registry);
jest.spyOn(registry, 'trackShutdown').mockImplementation((task, handle) =>
trackShutdown(task, {
...handle,
flush: (reason) => {
notifyFlush();
return handle.flush(reason);
},
}),
);
let resolveProjection: (persisted: boolean) => void = () => undefined;
const persist = jest.fn(
() =>
new Promise<boolean>((resolve) => {
resolveProjection = resolve;
}),
);
const retire = jest.fn(async () => true);
const tool = {
name: 'search_mcp_docs',
description: 'search docs',
schema: z.object({ q: z.string() }),
invoke: jest.fn(async () => ({ content: 'legacy projection result' })),
} as unknown as StructuredToolInterface;
const handler = createToolExecuteHandler({
loadTools: async () => ({ loadedTools: [tool] }),
backgroundToolCompletion: {
preregister: jest.fn(async () => ({ renew: async () => true, retire })),
persist,
claim: jest.fn(async () => ({ status: 'acquired' as const, results: [] })),
},
});

await runBatch(handler, {
toolCalls: [
{
id: 'call-legacy-projection',
name: tool.name,
args: { q: 'legacy', run_in_background: true },
stepId: 'step-legacy',
},
],
configurable: buildConfig([tool.name]),
metadata: { thread_id: 'shutdown_convo', run_id: 'response-legacy-projection' },
});
await flushMicrotasks();
expect(persist).toHaveBeenCalledTimes(1);

const draining = registry.drainForShutdown({
...drainOptions(),
deadlineAt: Date.now() + 1_750,
flushReserveMs: 1_400,
});
let drained = false;
void draining.then(() => {
drained = true;
});
await flushStarted;
await flushMicrotasks();
expect(drained).toBe(false);
resolveProjection(true);
expect(await draining).toEqual({ tracked: 1, interrupted: 0, flushed: 1, unsettled: 0 });
expect(retire).not.toHaveBeenCalled();
});

it('reports a pre-admitted legacy completion with no durable projection as unsettled', async () => {
const { createToolExecuteHandler, registry } = loadModules();
const persist = jest.fn(async () => false);
const tool = {
name: 'search_mcp_docs',
description: 'search docs',
schema: z.object({ q: z.string() }),
invoke: jest.fn(async () => ({ content: 'missing legacy result' })),
} as unknown as StructuredToolInterface;
const handler = createToolExecuteHandler({
loadTools: async () => ({ loadedTools: [tool] }),
backgroundToolCompletion: {
preregister: jest.fn(async () => ({ renew: async () => true, retire: async () => true })),
persist,
claim: jest.fn(async () => ({ status: 'acquired' as const, results: [] })),
},
});

await runBatch(handler, {
toolCalls: [
{
id: 'call-legacy-failed',
name: tool.name,
args: { q: 'failed', run_in_background: true },
stepId: 'step-legacy-failed',
},
],
configurable: buildConfig([tool.name]),
metadata: { thread_id: 'shutdown_convo', run_id: 'response-legacy-failed' },
});
await flushMicrotasks();
expect(persist).toHaveBeenCalledTimes(1);
expect(
await registry.drainForShutdown({
...drainOptions(),
deadlineAt: Date.now() + 190,
interruptGraceMs: 75,
flushReserveMs: 40,
}),
).toEqual({
tracked: 1,
interrupted: 0,
flushed: 1,
unsettled: 1,
});
});

it('projects an interrupted result for an uncooperative legacy completion without a receipt', async () => {
const { createToolExecuteHandler, registry } = loadModules();
const persist = jest.fn(async () => true);
const controlled = controlledTool({ honorAbort: false });
const handler = createToolExecuteHandler({
loadTools: async () => ({ loadedTools: [controlled.tool] }),
backgroundToolCompletion: {
preregister: jest.fn(async () => ({ renew: async () => true, retire: async () => true })),
persist,
claim: jest.fn(async () => ({ status: 'acquired' as const, results: [] })),
},
});

await runBatch(handler, {
toolCalls: [
{
id: 'call-legacy-stubborn',
name: controlled.tool.name,
args: { q: 'stubborn', run_in_background: true },
stepId: 'step-legacy-stubborn',
},
],
configurable: buildConfig([controlled.tool.name]),
metadata: { thread_id: 'shutdown_convo', run_id: 'response-legacy-stubborn' },
});

expect(
await registry.drainForShutdown({
...drainOptions(),
deadlineAt: Date.now() + 280,
interruptGraceMs: 90,
flushReserveMs: 70,
}),
).toEqual({
tracked: 1,
interrupted: 1,
flushed: 1,
unsettled: 0,
});
expect(persist).toHaveBeenCalledWith(
expect.objectContaining({
output: interruptedOutput,
backgroundTask: expect.objectContaining({ status: 'error' }),
}),
);

controlled.resolve('late success');
await flushMicrotasks();
await flushMicrotasks();
expect(persist).toHaveBeenCalledTimes(1);
});

it('reports a failed legacy interruption projection as unconfirmed', async () => {
const { createToolExecuteHandler, registry } = loadModules();
const persist = jest.fn(async () => false);
const controlled = controlledTool({ honorAbort: false });
const handler = createToolExecuteHandler({
loadTools: async () => ({ loadedTools: [controlled.tool] }),
backgroundToolCompletion: {
preregister: jest.fn(async () => ({ renew: async () => true, retire: async () => true })),
persist,
claim: jest.fn(async () => ({ status: 'acquired' as const, results: [] })),
},
});

await runBatch(handler, {
toolCalls: [
{
id: 'call-legacy-interruption-failed',
name: controlled.tool.name,
args: { q: 'stubborn', run_in_background: true },
stepId: 'step-legacy-interruption-failed',
},
],
configurable: buildConfig([controlled.tool.name]),
metadata: { thread_id: 'shutdown_convo', run_id: 'response-legacy-interruption-failed' },
});

expect(
await registry.drainForShutdown({
...drainOptions(),
deadlineAt: Date.now() + 280,
interruptGraceMs: 90,
flushReserveMs: 70,
}),
).toEqual({
tracked: 1,
interrupted: 1,
flushed: 1,
unsettled: 1,
});
expect(persist).toHaveBeenCalledTimes(1);
controlled.resolve('late success');
await flushMicrotasks();
await flushMicrotasks();
});

it('waits for the projected result when an independent receipt write returns false', async () => {
const { createToolExecuteHandler, registry } = loadModules();
const adapter = completionAdapter();
Expand Down
45 changes: 31 additions & 14 deletions packages/api/src/agents/handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5963,6 +5963,7 @@ export function createToolExecuteHandler(options: ToolExecuteOptions): EventHand
let durableReceiptWrite: Promise<boolean> | undefined;
let durableReceiptAmbiguous = false;
let durableResultConfirmed = false;
let forcedShutdownResult = false;
let resolveDurableReceipt: () => void = () => undefined;
const durableReceiptSettled = new Promise<void>((resolve) => {
resolveDurableReceipt = resolve;
Expand Down Expand Up @@ -6306,11 +6307,17 @@ export function createToolExecuteHandler(options: ToolExecuteOptions): EventHand
);
}
};
const persistSettledBackgroundResult = async (params: {
output?: string;
artifact?: unknown;
status: 'completed' | 'error' | 'cancelled';
}): Promise<void> => {
const persistSettledBackgroundResult = async (
params: {
output?: string;
artifact?: unknown;
status: 'completed' | 'error' | 'cancelled';
},
forced = false,
): Promise<void> => {
if (forcedShutdownResult && !forced) {
return;
}
settledReceipt = { status: params.status, output: params.output };
/** Held for the whole persist, including a code harvest that waits for
* a long dispatch turn, so retention pressure cannot evict the task. */
Expand Down Expand Up @@ -6671,7 +6678,7 @@ export function createToolExecuteHandler(options: ToolExecuteOptions): EventHand
})();
void settlement.then(
() => {
if (completionAdmission?.persistResult == null || !completionPreregistered) {
if (completionAdmission == null || !completionPreregistered) {
resolveDurableReceipt();
}
},
Expand All @@ -6693,6 +6700,7 @@ export function createToolExecuteHandler(options: ToolExecuteOptions): EventHand
} else if (settledReceipt != null) {
await writeDurableReceipt({ ...settledReceipt, settledAt: new Date() });
} else {
forcedShutdownResult = true;
const failure = toBackgroundToolFailure(tc.name, reason);
backgroundTaskRegistry.fail(
backgroundUserId,
Expand All @@ -6709,22 +6717,31 @@ export function createToolExecuteHandler(options: ToolExecuteOptions): EventHand
detachedError,
);
}
await writeDurableReceipt({
status: 'error',
output: failure,
settledAt: new Date(),
});
if (completionAdmission?.persistResult != null) {
await writeDurableReceipt({
status: 'error',
output: failure,
settledAt: new Date(),
});
} else if (completionAdmission != null) {
await persistSettledBackgroundResult(
{ status: 'error', output: failure },
true,
);
}
}
if (durableResultConfirmed || completionAdmission?.persistResult == null) {
if (durableResultConfirmed || completionAdmission == null) {
return;
}
if (settledReceipt != null) {
if (settledReceipt != null && !forcedShutdownResult) {
await settlement;
if (durableResultConfirmed) {
return;
}
}
throw new Error(`Background task ${task.id} has no durable shutdown result`);
if (completionPreregistered) {
throw new Error(`Background task ${task.id} has no durable shutdown result`);
}
},
});
if (
Expand Down
Loading