From bd96945139fb8c64826752cf1d6f2dc3d34fa666 Mon Sep 17 00:00:00 2001 From: Danny Avila Date: Thu, 24 Sep 2026 18:50:30 -0400 Subject: [PATCH 1/5] =?UTF-8?q?=F0=9F=93=AE=20fix:=20Count=20Undelivered?= =?UTF-8?q?=20Background=20Results=20as=20Outstanding?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../Endpoints/agents/backgroundCompletion.js | 8 + .../services/Endpoints/agents/initialize.js | 8 +- packages/api/src/agents/background.spec.ts | 220 ++++++++++++++++++ packages/api/src/agents/background.ts | 132 ++++++++++- .../api/src/agents/backgroundCompletion.ts | 37 +++ .../agents/backgroundCompletionWakeup.spec.ts | 88 +++++++ .../src/agents/backgroundCompletionWakeup.ts | 53 +++++ packages/api/src/agents/handlers.ts | 15 +- .../src/methods/triggerDelivery.spec.ts | 119 ++++++++++ .../src/methods/triggerDelivery.ts | 102 ++++++++ 10 files changed, 764 insertions(+), 18 deletions(-) diff --git a/api/server/services/Endpoints/agents/backgroundCompletion.js b/api/server/services/Endpoints/agents/backgroundCompletion.js index 8ca1cae6efc..e048d3aa7b7 100644 --- a/api/server/services/Endpoints/agents/backgroundCompletion.js +++ b/api/server/services/Endpoints/agents/backgroundCompletion.js @@ -1,9 +1,11 @@ const { createBackgroundToolCompletionWakeupHandler, createBackgroundToolDeadClaimRecovery, + createPendingBackgroundCompletions, createBackgroundToolResultHandler, claimBackgroundToolResult: claimResult, } = require('@librechat/api'); +const { listPendingAgentBackgroundToolCompletions } = require('~/models'); const { enqueueAgentTrigger, persistAgentBackgroundToolResult, @@ -21,6 +23,11 @@ const preregisterBackgroundToolCompletion = createBackgroundToolCompletionWakeup persistAgentBackgroundToolResult({ deliveryKey, sourceId, result }), ); +const pendingBackgroundToolCompletions = createPendingBackgroundCompletions({ + list: listPendingAgentBackgroundToolCompletions, + retire: retireAgentTrigger, +}); + function createBackgroundToolResultPersistence({ req, updateToolCallResult }) { return createBackgroundToolResultHandler({ req, updateToolCallResult }); } @@ -44,6 +51,7 @@ function createDeadBackgroundToolClaimRecovery( module.exports = { preregisterBackgroundToolCompletion, + pendingBackgroundToolCompletions, createBackgroundToolResultPersistence, claimBackgroundToolResult, createDeadBackgroundToolClaimRecovery, diff --git a/api/server/services/Endpoints/agents/initialize.js b/api/server/services/Endpoints/agents/initialize.js index 875d085951b..d759c98a1a2 100644 --- a/api/server/services/Endpoints/agents/initialize.js +++ b/api/server/services/Endpoints/agents/initialize.js @@ -94,6 +94,7 @@ const { processAddedConvo } = require('./addedConvo'); const subagentThreadTaskStore = require('./subagentThreadStore'); const { preregisterBackgroundToolCompletion, + pendingBackgroundToolCompletions, createBackgroundToolResultPersistence, claimBackgroundToolResult, createDeadBackgroundToolClaimRecovery, @@ -514,7 +515,12 @@ const initializeClientWithProvider = async ({ updateToolCallResult: db.updateToolCallResult, }), backgroundToolCompletion: { - ...(completionWakeupsEnabled ? { preregister: preregisterBackgroundToolCompletion } : {}), + ...(completionWakeupsEnabled + ? { + preregister: preregisterBackgroundToolCompletion, + pending: pendingBackgroundToolCompletions, + } + : {}), persist: createBackgroundToolResultPersistence({ req, updateToolCallResult: db.updateToolCallResult, diff --git a/packages/api/src/agents/background.spec.ts b/packages/api/src/agents/background.spec.ts index 5418fdeb3f1..08d03f4a3ab 100644 --- a/packages/api/src/agents/background.spec.ts +++ b/packages/api/src/agents/background.spec.ts @@ -433,6 +433,12 @@ describe('registerBackgroundTaskTool', () => { expect(automatic.toolDefinitions[0].description).toContain( 'Ordinary tool execution remains process-local', ); + expect(automatic.toolDefinitions[0].description).toContain( + 'A task is outstanding until its result is delivered', + ); + expect(automatic.toolDefinitions[0].description).toContain( + 'Polling or cancelling a finished task retires its pending delivery', + ); }); }); @@ -3411,3 +3417,217 @@ describe('toolOptionsSchema', () => { expect(parsed).toEqual({ run_in_background: true }); }); }); + +describe('runCheckBackgroundTask delivery semantics', () => { + const pendingControls = ( + overrides: { + list?: () => Promise; + discard?: () => Promise; + } = {}, + ) => + ({ + list: jest.fn(overrides.list ?? (async () => [])), + discard: jest.fn(overrides.discard ?? (async () => 'not_pending')), + }) as never; + + function completedWithWakeup(userId: string, conversationId: string, toolCallId: string) { + const created = backgroundTaskRegistry.create({ + userId, + conversationId, + toolCallId, + toolName: 'bash_tool', + messageId: `${toolCallId}-message`, + }); + if ('atCapacity' in created) { + throw new Error('unexpected capacity'); + } + backgroundTaskRegistry.markCompletionWakeup(userId, conversationId, created.task.id, { + renew: jest.fn(async () => true), + retire: jest.fn(async () => true), + }); + backgroundTaskRegistry.complete(userId, conversationId, created.task.id, { + content: 'finished output', + }); + return created.task.id; + } + + it('counts a finished task as outstanding until its result is delivered', async () => { + const taskId = completedWithWakeup('outstanding-user', 'outstanding-convo', 'outstanding-call'); + + const before = JSON.parse( + await runCheckBackgroundTask({ + userId: 'outstanding-user', + conversationId: 'outstanding-convo', + args: {}, + }), + ); + expect(before.tasks[0]).toEqual( + expect.objectContaining({ + background_task_id: taskId, + status: 'completed', + delivery: 'pending', + }), + ); + expect(before.outstanding).toBe(1); + expect(before.message).toContain('have not been delivered yet'); + + expect( + backgroundTaskRegistry.claimResult('outstanding-user', 'outstanding-convo', taskId, { + kind: 'wakeup', + claimId: 'automatic-delivery', + }), + ).toBe('acquired'); + const after = JSON.parse( + await runCheckBackgroundTask({ + userId: 'outstanding-user', + conversationId: 'outstanding-convo', + args: {}, + }), + ); + expect(after.tasks[0].delivery).toBe('delivered'); + expect(after.outstanding).toBe(0); + expect(after.message).toBeUndefined(); + }); + + it('counts running work as outstanding without claiming undelivered results', async () => { + const created = backgroundTaskRegistry.create({ + userId: 'running-user', + conversationId: 'running-convo', + toolCallId: 'running-call', + toolName: 'bash_tool', + }); + if ('atCapacity' in created) { + throw new Error('unexpected capacity'); + } + + const listed = JSON.parse( + await runCheckBackgroundTask({ + userId: 'running-user', + conversationId: 'running-convo', + args: {}, + }), + ); + expect(listed.tasks[0]).toEqual(expect.objectContaining({ status: 'running' })); + expect(listed.tasks[0].delivery).toBeUndefined(); + expect(listed.outstanding).toBe(1); + expect(listed.message).toBeUndefined(); + }); + + it('lists undelivered results the process-local registry no longer holds', async () => { + const localTaskId = completedWithWakeup('durable-user', 'durable-convo', 'durable-local'); + const pendingCompletions = pendingControls({ + list: async () => [ + { + taskId: localTaskId, + toolName: 'bash_tool', + dispatchedAt: new Date('2026-09-24T12:00:00Z'), + result: { status: 'completed', settledAt: new Date('2026-09-24T12:01:00Z') }, + claimedByWakeup: false, + }, + { + taskId: 'earlier-turn-task', + toolName: 'slow_task', + dispatchedAt: new Date('2026-09-24T11:00:00Z'), + result: { status: 'error', settledAt: new Date('2026-09-24T11:05:00Z') }, + claimedByWakeup: false, + }, + { + taskId: 'other-replica-task', + toolName: 'slow_task', + dispatchedAt: new Date('2026-09-24T11:30:00Z'), + claimedByWakeup: false, + }, + ], + }); + + const listed = JSON.parse( + await runCheckBackgroundTask({ + userId: 'durable-user', + conversationId: 'durable-convo', + args: {}, + pendingCompletions, + }), + ); + + expect( + listed.tasks.map((task: { background_task_id: string }) => task.background_task_id), + ).toEqual([localTaskId, 'earlier-turn-task', 'other-replica-task']); + expect(listed.tasks[1]).toEqual( + expect.objectContaining({ + status: 'error', + delivery: 'pending', + started_at: '2026-09-24T11:00:00.000Z', + settled_at: '2026-09-24T11:05:00.000Z', + }), + ); + expect(listed.tasks[2]).toEqual( + expect.objectContaining({ status: 'running', delivery: 'pending', progress: 0 }), + ); + expect(listed.tasks[1].result).toBeUndefined(); + expect(listed.outstanding).toBe(3); + }); + + it('keeps listing local work when the durable view is unavailable', async () => { + const taskId = completedWithWakeup('degraded-user', 'degraded-convo', 'degraded-call'); + const pendingCompletions = pendingControls({ + list: async () => { + throw new Error('delivery store unavailable'); + }, + }); + + const listed = JSON.parse( + await runCheckBackgroundTask({ + userId: 'degraded-user', + conversationId: 'degraded-convo', + args: {}, + pendingCompletions, + }), + ); + + expect( + listed.tasks.map((task: { background_task_id: string }) => task.background_task_id), + ).toEqual([taskId]); + expect(listed.partial).toBe(true); + expect(listed.warning).toContain('Undelivered results from earlier turns could not be listed'); + }); + + it.each([ + ['discarded', 'cancelled', 'will not arrive as a new turn'], + ['running', 'unavailable', 'cannot be stopped from here'], + ['delivering', 'delivery_scheduled', 'already being delivered'], + ])( + 'reports a %s undelivered completion this process does not hold', + async (outcome, status, message) => { + const pendingCompletions = pendingControls({ discard: async () => outcome }); + + const cancelled = JSON.parse( + await runCheckBackgroundTask({ + userId: 'discard-user', + conversationId: 'discard-convo', + args: { background_task_id: 'earlier-turn-task', action: 'cancel' }, + pendingCompletions, + }), + ); + + expect(cancelled).toEqual( + expect.objectContaining({ status, background_task_id: 'earlier-turn-task' }), + ); + expect(cancelled.message).toContain(message); + }, + ); + + it('falls through to the ordinary lookup when nothing is pending for the task', async () => { + const pendingCompletions = pendingControls({ discard: async () => 'not_pending' }); + + const cancelled = JSON.parse( + await runCheckBackgroundTask({ + userId: 'discard-user', + conversationId: 'discard-convo', + args: { background_task_id: 'unknown-task', action: 'cancel' }, + pendingCompletions, + }), + ); + + expect(cancelled).toEqual(expect.objectContaining({ status: 'not_found' })); + }); +}); diff --git a/packages/api/src/agents/background.ts b/packages/api/src/agents/background.ts index 2c9b744ee47..e4bb29f1f0e 100644 --- a/packages/api/src/agents/background.ts +++ b/packages/api/src/agents/background.ts @@ -57,8 +57,10 @@ import type { BackgroundToolResultState } from './harvest'; import type { CapabilityToolNames } from './selection'; import { BACKGROUND_TASK_TIMEOUT_MS, + type PendingBackgroundCompletion, type BackgroundToolDeadClaimRecovery, type BackgroundToolWakeupAdmission, + type PendingBackgroundCompletionControls, } from './backgroundCompletion'; import { CREATE_FILE_TOOL_NAME, @@ -358,7 +360,7 @@ Provide a background_task_id to poll one task; omit it to list every background const CHECK_BACKGROUND_TASK_WAKEUP_DESCRIPTION = `Check, control, and retrieve tool or subagent tasks previously dispatched in the background (with run_in_background: true). -Provide a background_task_id to inspect one task; omit it to list every background task in this thread. Background tools and detached subagents use automatic completion delivery: continue independent work or end the turn instead of repeatedly polling an unchanged running task, and the host will resume you when one finishes. Use this tool for explicit status, steer, queue, interrupt, cancel, or cancel_message actions, or as a fallback if automatic delivery is unavailable. Ordinary tool execution remains process-local and does not survive restart; once its result is persisted, completion delivery may continue on another replica. Live subagent controls route across API replicas but do not survive a restart of the process that owns the executor. A completed subagent thread may be continued later through the subagent tool's durable thread id.`; +Provide a background_task_id to inspect one task; omit it to list every background task in this thread. Background tools and detached subagents use automatic completion delivery: continue independent work or end the turn instead of repeatedly polling an unchanged running task, and the host will resume you when one finishes. Use this tool for explicit status, steer, queue, interrupt, cancel, or cancel_message actions, or as a fallback if automatic delivery is unavailable. A task is outstanding until its result is delivered, not merely until it stops running: a finished task whose delivery is "pending" will still arrive as a new turn, so never report it as done or cancelled on the strength of its status alone. Polling or cancelling a finished task retires its pending delivery so it never arrives as a new turn: a poll returns the result now, and cancelling a result this turn can no longer poll discards it. Ordinary tool execution remains process-local and does not survive restart; once its result is persisted, completion delivery may continue on another replica. Live subagent controls route across API replicas but do not survive a restart of the process that owns the executor. A completed subagent thread may be continued later through the subagent tool's durable thread id.`; function checkBackgroundTaskDescription(subagentCompletionWakeups: boolean): string { return subagentCompletionWakeups @@ -1840,10 +1842,17 @@ interface SerializedBackgroundTask { result?: string; result_available?: boolean; result_chars?: number; + /** Whether the result has reached the conversation. `pending` results still + * arrive as a new turn unless polled or cancelled first. Absent when the task + * has no automatic delivery, so only a poll ever surfaces its result. */ + delivery?: 'pending' | 'delivered'; note?: string; error?: string; } +const PENDING_DELIVERY_GUIDANCE = + 'Some finished tasks have not been delivered yet (delivery: "pending"); each will arrive as a new turn. Poll one to collect its result now, or cancel it so it does not arrive. Do not report these tasks as finished or cancelled until then.'; + /** * Model-facing task timings. The registry keeps epoch milliseconds; everything the * app serializes carries ISO-8601 (`toISOString`), so the poll payload does too. @@ -1900,6 +1909,37 @@ function taskNote(task: BackgroundTask): Pick return {}; } +function taskDelivery(task: BackgroundTask): Pick { + if (task.completionWakeup !== true || task.completionPersistenceFailed === true) { + return {}; + } + if (task.resultClaim != null || task.completionWakeupRetired === true) { + return { delivery: 'delivered' }; + } + return { delivery: 'pending' }; +} + +/** A completion known only to the durable delivery store: dispatched in an earlier + * turn, on another replica, or before a restart, and not delivered yet. */ +function serializePendingCompletion( + completion: PendingBackgroundCompletion, +): SerializedBackgroundTask { + const settled = completion.result; + return { + background_task_id: completion.taskId, + tool: completion.toolName, + status: settled?.status ?? 'running', + progress: settled == null ? 0 : 1, + started_at: completion.dispatchedAt.toISOString(), + ...(settled != null && { settled_at: settled.settledAt.toISOString() }), + delivery: 'pending', + note: + settled == null + ? 'Still running outside this turn; its result will arrive as a new turn when it finishes.' + : 'Finished, but its result has not been delivered; it will arrive as a new turn unless you poll or cancel it.', + }; +} + function serializeTask( task: BackgroundTask, { includeResult }: { includeResult: boolean }, @@ -1914,6 +1954,7 @@ function serializeTask( : {}), ...taskTimings(task), ...resultFields(task, includeResult), + ...taskDelivery(task), ...taskNote(task), ...(task.error !== undefined ? { error: task.error } : {}), }; @@ -2113,6 +2154,8 @@ export async function runCheckBackgroundTask(params: { recoverDeadBackgroundToolClaim?: BackgroundToolDeadClaimRecovery; /** Trusted deployment policy. Defaults false for backward compatibility. */ ordinaryToolCancellation?: boolean; + /** Durable view of this conversation's undelivered background completions. */ + pendingCompletions?: PendingBackgroundCompletionControls; }): Promise { const { userId, conversationId } = params; const args = coerceArgsObject(params.args) ?? {}; @@ -2399,6 +2442,44 @@ export async function runCheckBackgroundTask(params: { return JSON.stringify(serializeTask(task, { includeResult: true })); } + if (action === 'cancel' && params.pendingCompletions != null) { + let outcome: Awaited>; + try { + outcome = await params.pendingCompletions.discard({ userId, conversationId, taskId }); + } catch (error) { + logger.warn(`[background] Failed to discard pending completion ${taskId}:`, error); + return JSON.stringify({ + status: 'unavailable', + background_task_id: taskId, + message: + 'The pending result could not be discarded right now. It may still arrive as a new turn; retry the cancel shortly.', + }); + } + if (outcome === 'discarded') { + return JSON.stringify({ + status: 'cancelled', + background_task_id: taskId, + message: 'The finished result was discarded and will not arrive as a new turn.', + }); + } + if (outcome === 'running') { + return JSON.stringify({ + status: 'unavailable', + background_task_id: taskId, + delivery: 'pending', + message: + 'This task is still running outside this turn and cannot be stopped from here. Its result will arrive as a new turn when it finishes; cancel it then to discard the result.', + }); + } + if (outcome === 'delivering') { + return JSON.stringify({ + status: 'delivery_scheduled', + background_task_id: taskId, + message: 'This result is already being delivered as a new turn.', + }); + } + } + const subagentTasks = params.subagentTasks; let subagentPollChecked = false; let subagentPollError: unknown; @@ -2586,7 +2667,21 @@ export async function runCheckBackgroundTask(params: { const tasks = backgroundTaskRegistry.list(userId, conversationId); let subagentTasks: SerializedSubagentTask[] = []; - let listWarning: string | undefined; + const listWarnings: string[] = []; + let pendingCompletions: PendingBackgroundCompletion[] = []; + if (params.pendingCompletions != null) { + try { + const localTaskIds = new Set(tasks.map((task) => task.id)); + pendingCompletions = ( + await params.pendingCompletions.list({ userId, conversationId }) + ).filter((completion) => !localTaskIds.has(completion.taskId)); + } catch (error) { + logger.warn('[background] Failed to list undelivered background completions:', error); + listWarnings.push( + 'Undelivered results from earlier turns could not be listed; some may still arrive as new turns.', + ); + } + } const completionWakeups = agentUsesSubagentCompletionWakeups( params.subagentTasks, params.agentId, @@ -2607,24 +2702,37 @@ export async function runCheckBackgroundTask(params: { subagentTasks = params.subagentTasks.store .list(params.subagentTasks.scopeId) .map((task) => serializeSubagentSnapshot(task)); - listWarning = `Cross-replica subagent tasks could not be listed: ${error.message}`; + listWarnings.push(`Cross-replica subagent tasks could not be listed: ${error.message}`); } else { throw error; } } } + const ordinaryTasks = [ + ...tasks.map((task) => serializeTask(task, { includeResult: false })), + ...pendingCompletions.map(serializePendingCompletion), + ]; + /** Work is outstanding until its result reaches the conversation: a finished + * task with a pending delivery is still going to resume the agent. */ + const outstanding = + ordinaryTasks.filter((task) => task.status === 'running' || task.delivery === 'pending') + .length + subagentTasks.filter((task) => task.status === 'running').length; + const guidance = [ + ...(ordinaryTasks.some((task) => task.status !== 'running' && task.delivery === 'pending') + ? [PENDING_DELIVERY_GUIDANCE] + : []), + ...(completionWakeups && subagentTasks.some((task) => task.status === 'running') + ? [SUBAGENT_WAKEUP_GUIDANCE] + : []), + ]; logger.debug( - `[background] check_background_task listed ${tasks.length + subagentTasks.length} task(s)`, + `[background] check_background_task listed ${ordinaryTasks.length + subagentTasks.length} task(s), ${outstanding} outstanding`, ); return JSON.stringify({ - tasks: [ - ...tasks.map((task) => serializeTask(task, { includeResult: false })), - ...subagentTasks, - ], - ...(completionWakeups && subagentTasks.some((task) => task.status === 'running') - ? { message: SUBAGENT_WAKEUP_GUIDANCE } - : {}), - ...(listWarning != null && { partial: true, warning: listWarning }), + tasks: [...ordinaryTasks, ...subagentTasks], + outstanding, + ...(guidance.length > 0 && { message: guidance.join(' ') }), + ...(listWarnings.length > 0 && { partial: true, warning: listWarnings.join(' ') }), }); } diff --git a/packages/api/src/agents/backgroundCompletion.ts b/packages/api/src/agents/backgroundCompletion.ts index 8ef1d74cca5..c27697997e2 100644 --- a/packages/api/src/agents/backgroundCompletion.ts +++ b/packages/api/src/agents/backgroundCompletion.ts @@ -62,3 +62,40 @@ export interface BackgroundToolDeadClaimRecoveryInput { export type BackgroundToolDeadClaimRecovery = ( input: BackgroundToolDeadClaimRecoveryInput, ) => Promise; + +/** A background tool completion whose result has not reached its conversation yet, + * read from the durable delivery store rather than a process-local registry. */ +export interface PendingBackgroundCompletion { + taskId: string; + toolName: string; + dispatchedAt: Date; + /** The tool's terminal outcome once it settled; absent while it still runs. */ + result?: { status: 'completed' | 'error' | 'cancelled'; settledAt: Date }; + /** An automatic delivery holds the result and is starting its turn. */ + claimedByWakeup: boolean; +} + +/** + * What cancelling an undelivered completion did: `discarded` retired its delivery, + * so the result never arrives; `running` found the tool still executing where this + * process cannot stop it; `delivering` found the result already being delivered; + * `not_pending` found no undelivered completion for the task. + */ +export type BackgroundCompletionDiscardOutcome = + | 'discarded' + | 'running' + | 'delivering' + | 'not_pending'; + +/** Durable view and control of one principal's undelivered background completions. */ +export interface PendingBackgroundCompletionControls { + list: (input: { + userId: string; + conversationId: string; + }) => Promise; + discard: (input: { + userId: string; + conversationId: string; + taskId: string; + }) => Promise; +} diff --git a/packages/api/src/agents/backgroundCompletionWakeup.spec.ts b/packages/api/src/agents/backgroundCompletionWakeup.spec.ts index 1e87d6a2f55..ae60ba4eb39 100644 --- a/packages/api/src/agents/backgroundCompletionWakeup.spec.ts +++ b/packages/api/src/agents/backgroundCompletionWakeup.spec.ts @@ -4,6 +4,7 @@ import type { CodeApprovalMode } from 'librechat-data-provider'; import type { EnqueueBackgroundToolCompletion } from './backgroundCompletionWakeup'; import { BACKGROUND_TOOL_WAKEUP_INPUT_MAX_CHARS, + createPendingBackgroundCompletions, createBackgroundToolCompletionWakeupHandler, createBackgroundToolCompletionWakeupResolver, createBackgroundToolDeadClaimRecovery, @@ -790,3 +791,90 @@ describe('background tool completion wakeups', () => { ); }); }); + +describe('pending background completions', () => { + const dispatchedAt = new Date(NOW - 60_000); + const row = (overrides = {}) => ({ + deliveryKey: 'delivery-key-1', + taskId: 'task-1', + toolCallId: 'call-1', + toolName: 'slow_tool', + dispatchedAt, + claimedByWakeup: false, + ...overrides, + }); + + it('lists the durable view for the owner without delivery internals', async () => { + const list = jest.fn(async () => [ + row({ result: { status: 'completed' as const, settledAt: new Date(NOW) } }), + ]); + const pending = createPendingBackgroundCompletions({ list, retire: jest.fn() }); + + await expect( + pending.list({ userId: 'user-1', conversationId: 'conversation-1' }), + ).resolves.toEqual([ + { + taskId: 'task-1', + toolName: 'slow_tool', + dispatchedAt, + result: { status: 'completed', settledAt: new Date(NOW) }, + claimedByWakeup: false, + }, + ]); + expect(list).toHaveBeenCalledWith({ + user: 'user-1', + conversationId: 'conversation-1', + sourceId: 'background-tool-completion', + }); + }); + + it('discards a settled, unclaimed result by retiring only an unclaimed delivery', async () => { + const retire = jest.fn(async () => true); + const pending = createPendingBackgroundCompletions({ + list: async () => [ + row({ result: { status: 'completed' as const, settledAt: new Date(NOW) } }), + ], + retire, + }); + + await expect( + pending.discard({ userId: 'user-1', conversationId: 'conversation-1', taskId: 'task-1' }), + ).resolves.toBe('discarded'); + expect(retire).toHaveBeenCalledWith( + 'delivery-key-1', + 'background-tool-completion', + 'background result discarded by its owner', + { onlyIfUnclaimed: true }, + ); + }); + + it.each([ + ['not_pending', [], true], + ['running', [row()], true], + [ + 'delivering', + [ + row({ + result: { status: 'completed' as const, settledAt: new Date(NOW) }, + claimedByWakeup: true, + }), + ], + true, + ], + [ + 'delivering', + [row({ result: { status: 'completed' as const, settledAt: new Date(NOW) } })], + false, + ], + ])('reports %s without discarding what it cannot', async (outcome, rows, retired) => { + const retire = jest.fn(async () => retired); + const pending = createPendingBackgroundCompletions({ list: async () => rows, retire }); + + await expect( + pending.discard({ userId: 'user-1', conversationId: 'conversation-1', taskId: 'task-1' }), + ).resolves.toBe(outcome); + if (outcome !== 'delivering' || rows[0]?.claimedByWakeup === true) { + expect(retire).not.toHaveBeenCalled(); + } + }); +}); diff --git a/packages/api/src/agents/backgroundCompletionWakeup.ts b/packages/api/src/agents/backgroundCompletionWakeup.ts index bfdb9f21d3b..dcc2e3243b4 100644 --- a/packages/api/src/agents/backgroundCompletionWakeup.ts +++ b/packages/api/src/agents/backgroundCompletionWakeup.ts @@ -10,9 +10,11 @@ import type { } from '@librechat/data-schemas'; import type { BackgroundToolDeadClaimRecovery, + PendingBackgroundCompletion, BackgroundToolWakeupAdmission, BackgroundToolWakeupRegistration, BackgroundToolWakeupRetireOptions, + PendingBackgroundCompletionControls, } from './backgroundCompletion'; import type { AgentTriggerContinuePreparation, @@ -532,6 +534,57 @@ export function createBackgroundToolCompletionWakeupResolver({ }; } +/** Lists and discards a conversation's undelivered background completions from the + * durable delivery store, which outlives the process-local task registry: a result + * dispatched in an earlier turn, on another replica, or before a restart is still + * going to arrive, and the owner must be able to see and stop that. */ +export function createPendingBackgroundCompletions(deps: { + list: (input: { + user: string; + conversationId: string; + sourceId: string; + }) => Promise>; + retire: RetireBackgroundToolCompletion; +}): PendingBackgroundCompletionControls { + const read = (input: { userId: string; conversationId: string }) => + deps.list({ + user: input.userId, + conversationId: input.conversationId, + sourceId: BACKGROUND_TOOL_COMPLETION_SOURCE, + }); + return { + list: async (input) => + (await read(input)).map(({ taskId, toolName, dispatchedAt, result, claimedByWakeup }) => ({ + taskId, + toolName, + dispatchedAt, + ...(result != null && { result }), + claimedByWakeup, + })), + discard: async (input) => { + const completion = (await read(input)).find(({ taskId }) => taskId === input.taskId); + if (completion == null) { + return 'not_pending'; + } + if (completion.result == null) { + return 'running'; + } + if (completion.claimedByWakeup) { + return 'delivering'; + } + /** Unclaimed-only: once a resolver owns the delivery its continuation can no + * longer be withdrawn, so that race reports as already delivering. */ + const retired = await deps.retire( + completion.deliveryKey, + BACKGROUND_TOOL_COMPLETION_SOURCE, + 'background result discarded by its owner', + { onlyIfUnclaimed: true }, + ); + return retired ? 'discarded' : 'delivering'; + }, + }; +} + /** Pre-registers the ordered completion delivery before external tool work starts. */ export function createBackgroundToolCompletionWakeupHandler( enqueue: EnqueueBackgroundToolCompletion, diff --git a/packages/api/src/agents/handlers.ts b/packages/api/src/agents/handlers.ts index f991418273d..9aff965f292 100644 --- a/packages/api/src/agents/handlers.ts +++ b/packages/api/src/agents/handlers.ts @@ -35,6 +35,12 @@ import type { import type { CodeEnvRef, CodeWorkspaceOperation, PtcToolCallEvent } from 'librechat-data-provider'; import type { StructuredToolInterface } from '@librechat/agents/langchain/tools'; import type { CodeEnvFile, CodeSessionContext } from '@librechat/agents'; +import type { + BackgroundToolDeadClaimRecovery, + BackgroundToolWakeupAdmission, + BackgroundToolWakeupRegistration, + PendingBackgroundCompletionControls, +} from './backgroundCompletion'; import type { WorkspaceEditResult, WorkspacePreviewEditResult, @@ -43,11 +49,6 @@ import type { WorkspaceSearchResult, WorkspaceWriteResult, } from '~/code/workspace'; -import type { - BackgroundToolDeadClaimRecovery, - BackgroundToolWakeupAdmission, - BackgroundToolWakeupRegistration, -} from './backgroundCompletion'; import type { SkillFileRecord, PrimeSkillFilesResult } from './skillFiles'; import type { ArtifactDeliveryFailure } from '~/files/code'; import type { BackgroundToolResultState } from './harvest'; @@ -345,6 +346,9 @@ export interface ToolExecuteOptions { allowUnfinished?: boolean; }) => Promise; recoverDeadClaim?: BackgroundToolDeadClaimRecovery; + /** Durable view of undelivered completions, so a status check counts results + * dispatched in earlier turns, on other replicas, or before a restart. */ + pending?: PendingBackgroundCompletionControls; }; /** Emits an `attachment` SSE event on the current request's live stream. */ emitAttachment?: (attachment: unknown) => void; @@ -6669,6 +6673,7 @@ export function createToolExecuteHandler(options: ToolExecuteOptions): EventHand subagentTasks, claimBackgroundToolResult: backgroundToolCompletion?.claim, recoverDeadBackgroundToolClaim: backgroundToolCompletion?.recoverDeadClaim, + pendingCompletions: backgroundToolCompletion?.pending, ordinaryToolCancellation, }); const taskSnapshot = getBackgroundTaskSnapshot({ diff --git a/packages/data-schemas/src/methods/triggerDelivery.spec.ts b/packages/data-schemas/src/methods/triggerDelivery.spec.ts index c2b61109308..e6dc2f5101a 100644 --- a/packages/data-schemas/src/methods/triggerDelivery.spec.ts +++ b/packages/data-schemas/src/methods/triggerDelivery.spec.ts @@ -453,6 +453,125 @@ describe('agent trigger delivery methods', () => { ).resolves.toEqual({ status: 'live', leaseUntil: renewedUntil }); }); + describe('listPendingAgentBackgroundToolCompletions', () => { + const background = { id: 'background-tool-completion', type: 'internal' }; + const completion = ( + user: mongoose.Types.ObjectId, + taskId: string, + overrides: Partial[0]> = {}, + ) => + methods.enqueueAgentTriggerDelivery( + enqueueInput({ + user, + orderingKey: `background-lane-${taskId}`, + envelope: { + event: { + source: background, + payload: { taskId, toolCallId: `call-${taskId}`, toolName: 'slow_task' }, + }, + target: { conversationId: 'conversation-1' }, + }, + requiredWorkerCapability: + AGENT_TRIGGER_WORKER_CAPABILITY_BACKGROUND_COMPLETION_RECEIPT_V2, + ...overrides, + }), + ); + + it('lists what is still going to arrive, running or settled, without result content', async () => { + const user = new mongoose.Types.ObjectId(); + const running = await completion(user, 'task-running'); + const settled = await completion(user, 'task-settled'); + await methods.persistAgentBackgroundToolResult({ + deliveryKey: settled.delivery.deliveryKey, + sourceId: background.id, + result: { status: 'completed', output: 'secret output', settledAt: START }, + }); + + const pending = await methods.listPendingAgentBackgroundToolCompletions({ + user, + conversationId: 'conversation-1', + sourceId: background.id, + }); + + expect(pending).toEqual([ + { + deliveryKey: running.delivery.deliveryKey, + taskId: 'task-running', + toolCallId: 'call-task-running', + toolName: 'slow_task', + dispatchedAt: expect.any(Date), + claimedByWakeup: false, + }, + { + deliveryKey: settled.delivery.deliveryKey, + taskId: 'task-settled', + toolCallId: 'call-task-settled', + toolName: 'slow_task', + dispatchedAt: expect.any(Date), + result: { status: 'completed', settledAt: START }, + claimedByWakeup: false, + }, + ]); + expect(JSON.stringify(pending)).not.toContain('secret output'); + }); + + it('excludes delivered rows and everything outside the conversation, user, and source', async () => { + const user = new mongoose.Types.ObjectId(); + const delivered = await completion(user, 'task-delivered'); + await Delivery.updateOne({ _id: delivered.delivery.id }, { $set: { status: 'succeeded' } }); + await completion(new mongoose.Types.ObjectId(), 'task-other-user'); + await completion(user, 'task-other-conversation', { + envelope: { + event: { + source: background, + payload: { + taskId: 'task-other-conversation', + toolCallId: 'call', + toolName: 'slow_task', + }, + }, + target: { conversationId: 'conversation-2' }, + }, + }); + await completion(user, 'task-other-source', { + envelope: { + event: { + source: { id: 'agent-queued-turn', type: 'internal' }, + payload: { taskId: 'task-other-source', toolCallId: 'call', toolName: 'slow_task' }, + }, + target: { conversationId: 'conversation-1' }, + }, + }); + const waiting = await completion(user, 'task-waiting'); + + const pending = await methods.listPendingAgentBackgroundToolCompletions({ + user, + conversationId: 'conversation-1', + sourceId: background.id, + }); + + expect(pending.map((entry) => entry.deliveryKey)).toEqual([waiting.delivery.deliveryKey]); + }); + + it('refuses a malformed lookup', async () => { + await expect( + methods.listPendingAgentBackgroundToolCompletions({ + user: new mongoose.Types.ObjectId(), + conversationId: '', + sourceId: background.id, + }), + ).rejects.toThrow(TypeError); + await expect( + methods.listPendingAgentBackgroundToolCompletions({ + user: new mongoose.Types.ObjectId(), + conversationId: 'conversation-1', + sourceId: background.id, + limit: 0, + }), + ).rejects.toThrow(TypeError); + }); + }); + it('persists one private background result receipt independently of message rows', async () => { const source = { id: 'background-tool-completion', type: 'internal' }; const queued = await methods.enqueueAgentTriggerDelivery( diff --git a/packages/data-schemas/src/methods/triggerDelivery.ts b/packages/data-schemas/src/methods/triggerDelivery.ts index e3fe4c2d0fa..fd87e713e14 100644 --- a/packages/data-schemas/src/methods/triggerDelivery.ts +++ b/packages/data-schemas/src/methods/triggerDelivery.ts @@ -44,6 +44,18 @@ export const CLAIM_CAS_MAX_ATTEMPTS = 16; /** Candidates fetched per claim read; losing claimers advance through the * batch instead of re-reading the same head-of-queue row. */ const CLAIM_CANDIDATE_BATCH = 8; +/** A conversation's undelivered completions are few; this bounds a pathological listing. */ +const MAX_PENDING_BACKGROUND_COMPLETIONS = 50; +/** Every status before a delivery settles, i.e. whose result has not reached its conversation. */ +const UNDELIVERED_STATUSES: IAgentTriggerDelivery['status'][] = [ + 'staging', + 'capability_staging', + 'batched', + 'pending', + 'capability_pending', + 'leased', + 'capability_leased', +]; /** Capability work is inert to legacy claimers while preserving their lane * behavior: publishing is `staging`; queued work is `leased` without a lease * owner/deadline; execution adds a private lease; dead work is terminal. */ @@ -246,6 +258,19 @@ export interface AgentEventActorReceiptStorageMetrics { deadDeliveries: number; } +/** A background tool completion that has not reached its conversation yet. */ +export interface PendingAgentBackgroundToolCompletion { + deliveryKey: string; + taskId: string; + toolCallId: string; + toolName: string; + dispatchedAt: Date; + /** The tool's terminal outcome once it settled; absent while it still runs. */ + result?: { status: AgentBackgroundToolResultReceipt['status']; settledAt: Date }; + /** An automatic delivery holds the result and is starting its turn. */ + claimedByWakeup: boolean; +} + export interface AgentTriggerDeliveryMethods { ensureAgentTriggerDeliveryIndexes: () => Promise; enqueueAgentTriggerDelivery: ( @@ -307,6 +332,12 @@ export interface AgentTriggerDeliveryMethods { sourceId: string; now: Date; }) => Promise; + listPendingAgentBackgroundToolCompletions: (input: { + user: string | Types.ObjectId; + conversationId: string; + sourceId: string; + limit?: number; + }) => Promise; persistAgentBackgroundToolResult: ( input: PersistAgentBackgroundToolResultInput, ) => Promise; @@ -2287,6 +2318,76 @@ export function createAgentTriggerDeliveryMethods( : { status: 'expired', leaseUntil: delivery.producerLeaseUntil }; } + /** Every background completion of one conversation that has not been + * delivered yet: still running, or settled and waiting for a wake-up. The + * durable delivery row outlives the process-local task registry (another + * replica, a restart, or the registry's retention), so it is the record of + * what is still going to arrive. Result content is never returned here. */ + async function listPendingAgentBackgroundToolCompletions(input: { + user: string | Types.ObjectId; + conversationId: string; + sourceId: string; + limit?: number; + }): Promise { + const limit = input.limit ?? MAX_PENDING_BACKGROUND_COMPLETIONS; + if ( + input.conversationId.length === 0 || + input.conversationId.length > 256 || + input.sourceId.length === 0 || + input.sourceId.length > 256 || + !Number.isSafeInteger(limit) || + limit <= 0 + ) { + throw new TypeError('Invalid pending background completion lookup'); + } + const rows = await Delivery() + .find({ + user: input.user, + 'envelope.event.source.type': 'internal', + 'envelope.event.source.id': input.sourceId, + 'envelope.target.conversationId': input.conversationId, + status: { $in: UNDELIVERED_STATUSES }, + }) + .select('+backgroundToolResult deliveryKey createdAt envelope.event.payload') + .sort({ createdAt: 1, _id: 1 }) + .limit(Math.min(limit, MAX_PENDING_BACKGROUND_COMPLETIONS)) + .lean< + Array< + Pick & { + envelope?: { event?: { payload?: Record } }; + } + > + >(); + return rows.flatMap((row) => { + const payload = row.envelope?.event?.payload; + const taskId = payload?.taskId; + const toolCallId = payload?.toolCallId; + const toolName = payload?.toolName; + if ( + row.createdAt == null || + typeof taskId !== 'string' || + typeof toolCallId !== 'string' || + typeof toolName !== 'string' + ) { + return []; + } + const receipt = row.backgroundToolResult; + return [ + { + deliveryKey: row.deliveryKey, + taskId, + toolCallId, + toolName, + dispatchedAt: row.createdAt, + ...(receipt != null && { + result: { status: receipt.status, settledAt: receipt.settledAt }, + }), + claimedByWakeup: receipt?.resultClaim != null, + }, + ]; + }); + } + /** Stores terminal output on the pre-admitted delivery before attempting the * parent-message projection. The first terminal receipt wins; exact retries * are idempotent and conflicting rewrites fail closed. */ @@ -4049,6 +4150,7 @@ export function createAgentTriggerDeliveryMethods( retireAgentTriggerDelivery, renewAgentTriggerDeliveryProducerLease, getAgentTriggerDeliveryProducerLease, + listPendingAgentBackgroundToolCompletions, persistAgentBackgroundToolResult, getAgentBackgroundToolResult, getAgentBackgroundToolResultClaim, From 352e8cede237adc4833c1843f718125d02c3b996 Mon Sep 17 00:00:00 2001 From: Danny Avila Date: Thu, 24 Sep 2026 19:13:59 -0400 Subject: [PATCH 2/5] fix: Reconcile delivered tasks, discard exactly, and count pending subagent results --- packages/api/src/agents/background.spec.ts | 125 +++++++++++++++++- packages/api/src/agents/background.ts | 62 +++++++-- .../api/src/agents/backgroundCompletion.ts | 5 +- .../agents/backgroundCompletionWakeup.spec.ts | 65 ++++----- .../src/agents/backgroundCompletionWakeup.ts | 39 ++++-- packages/api/src/agents/triggers/service.ts | 3 +- .../src/methods/triggerDelivery.spec.ts | 80 ++++++++++- .../src/methods/triggerDelivery.ts | 46 +++++-- 8 files changed, 355 insertions(+), 70 deletions(-) diff --git a/packages/api/src/agents/background.spec.ts b/packages/api/src/agents/background.spec.ts index 08d03f4a3ab..ff784c4ccd5 100644 --- a/packages/api/src/agents/background.spec.ts +++ b/packages/api/src/agents/background.spec.ts @@ -3063,6 +3063,52 @@ describe('runCheckBackgroundTask (singleton)', () => { expect(claimBackgroundToolResult).toHaveBeenCalledTimes(1); }); + it('counts a finished subagent as outstanding until its result is delivered', async () => { + const store = new InMemorySubagentTaskStore(); + const subagentTasks: HostSubagentTaskConfig = { + store, + scopeId: 'owner:settled-subagent-parent', + completionDelivery: SUBAGENT_COMPLETION_DELIVERY, + }; + const started = store.start({ + scopeId: subagentTasks.scopeId, + idempotencyKey: 'parent-run:parent-agent:call-settled', + parentRunId: 'parent-run', + parentAgentId: 'parent-agent', + parentToolCallId: 'call-settled', + input: 'Research this.', + subagentKind: 'agent', + subagentType: 'researcher', + run: async () => ({ content: 'research done' }), + }); + if (!started.accepted) { + throw new Error('Expected subagent task to start.'); + } + for ( + let i = 0; + i < 50 && store.get(subagentTasks.scopeId, started.task.taskId)?.status === 'running'; + i++ + ) { + await new Promise((resolve) => setImmediate(resolve)); + } + + const listed = JSON.parse( + await runCheckBackgroundTask({ + userId: 'owner', + conversationId: 'settled-subagent-parent', + agentId: 'agent_parent', + args: {}, + subagentTasks, + }), + ); + + expect(listed.tasks[0]).toEqual( + expect.objectContaining({ status: 'completed', result_available: true, delivery: 'pending' }), + ); + expect(listed.outstanding).toBe(1); + expect(listed.message).toContain('have not been delivered yet'); + }); + it('tells a wakeup-enabled parent to yield on an unchanged running subagent', async () => { const store = new InMemorySubagentTaskStore(); const subagentTasks: HostSubagentTaskConfig = { @@ -3422,11 +3468,15 @@ describe('runCheckBackgroundTask delivery semantics', () => { const pendingControls = ( overrides: { list?: () => Promise; + complete?: boolean; discard?: () => Promise; } = {}, ) => ({ - list: jest.fn(overrides.list ?? (async () => [])), + list: jest.fn(async () => ({ + completions: await (overrides.list ?? (async () => []))(), + complete: overrides.complete ?? true, + })), discard: jest.fn(overrides.discard ?? (async () => 'not_pending')), }) as never; @@ -3616,6 +3666,79 @@ describe('runCheckBackgroundTask delivery semantics', () => { }, ); + it('reports a local task delivered once the durable store no longer holds its delivery', async () => { + const delivered = completedWithWakeup('reconcile-user', 'reconcile-convo', 'reconcile-done'); + const waiting = completedWithWakeup('reconcile-user', 'reconcile-convo', 'reconcile-waiting'); + const listed = JSON.parse( + await runCheckBackgroundTask({ + userId: 'reconcile-user', + conversationId: 'reconcile-convo', + args: {}, + pendingCompletions: pendingControls({ + list: async () => [ + { + taskId: waiting, + toolName: 'bash_tool', + dispatchedAt: new Date('2026-09-24T12:00:00Z'), + result: { status: 'completed', settledAt: new Date('2026-09-24T12:01:00Z') }, + claimedByWakeup: false, + }, + ], + }), + }), + ); + + const byId = new Map( + listed.tasks.map((task: { background_task_id: string; delivery?: string }) => [ + task.background_task_id, + task.delivery, + ]), + ); + expect(byId.get(delivered)).toBe('delivered'); + expect(byId.get(waiting)).toBe('pending'); + expect(listed.outstanding).toBe(1); + }); + + it('keeps the local view and warns when the durable listing is incomplete', async () => { + const taskId = completedWithWakeup('truncated-user', 'truncated-convo', 'truncated-call'); + const listed = JSON.parse( + await runCheckBackgroundTask({ + userId: 'truncated-user', + conversationId: 'truncated-convo', + args: {}, + pendingCompletions: pendingControls({ complete: false }), + }), + ); + + expect(listed.tasks[0]).toEqual( + expect.objectContaining({ background_task_id: taskId, delivery: 'pending' }), + ); + expect(listed.outstanding).toBe(1); + expect(listed.partial).toBe(true); + expect(listed.warning).toContain('More undelivered results exist'); + }); + + it('lets a finished local task be cancelled without the live-cancellation policy', async () => { + const taskId = completedWithWakeup( + 'settled-cancel-user', + 'settled-cancel-convo', + 'settled-cancel', + ); + + const cancelled = JSON.parse( + await runCheckBackgroundTask({ + userId: 'settled-cancel-user', + conversationId: 'settled-cancel-convo', + args: { background_task_id: taskId, action: 'cancel' }, + }), + ); + + expect(cancelled.status).not.toBe('invalid'); + expect(cancelled).toEqual( + expect.objectContaining({ background_task_id: taskId, status: 'completed' }), + ); + }); + it('falls through to the ordinary lookup when nothing is pending for the task', async () => { const pendingCompletions = pendingControls({ discard: async () => 'not_pending' }); diff --git a/packages/api/src/agents/background.ts b/packages/api/src/agents/background.ts index e4bb29f1f0e..0075e04ed36 100644 --- a/packages/api/src/agents/background.ts +++ b/packages/api/src/agents/background.ts @@ -1940,6 +1940,24 @@ function serializePendingCompletion( }; } +/** A local task whose durable delivery settled elsewhere (an automatic wake-up + * on any replica) no longer holds a local claim; the complete durable listing is + * the evidence. An incomplete listing proves nothing, so the local view stands. */ +function reconcileDelivery( + task: SerializedBackgroundTask, + durablePendingTaskIds: ReadonlySet | undefined, +): SerializedBackgroundTask { + if ( + task.delivery !== 'pending' || + task.status === 'running' || + durablePendingTaskIds == null || + durablePendingTaskIds.has(task.background_task_id) + ) { + return task; + } + return { ...task, delivery: 'delivered' }; +} + function serializeTask( task: BackgroundTask, { includeResult }: { includeResult: boolean }, @@ -1994,6 +2012,8 @@ interface SerializedSubagentTask { error?: string; control_id?: string; message?: string; + /** A finished subagent whose result will still resume the parent turn. */ + delivery?: 'pending'; } function serializeSubagentSnapshot( @@ -2175,7 +2195,10 @@ export async function runCheckBackgroundTask(params: { if (task != null) { if (action !== 'poll') { if (action === 'cancel') { - if (params.ordinaryToolCancellation !== true) { + /** A finished task has no execution to stop, so the live-cancellation + * policy does not apply: cancelling it falls through to the poll path, + * which retires its pending delivery. */ + if (params.ordinaryToolCancellation !== true && task.status === 'running') { return JSON.stringify({ status: 'invalid', background_task_id: taskId, @@ -2669,12 +2692,23 @@ export async function runCheckBackgroundTask(params: { let subagentTasks: SerializedSubagentTask[] = []; const listWarnings: string[] = []; let pendingCompletions: PendingBackgroundCompletion[] = []; + /** Undelivered task ids from the durable store, when the listing was complete: + * a local task absent from it was delivered on this or another replica. */ + let durablePendingTaskIds: ReadonlySet | undefined; if (params.pendingCompletions != null) { try { const localTaskIds = new Set(tasks.map((task) => task.id)); - pendingCompletions = ( - await params.pendingCompletions.list({ userId, conversationId }) - ).filter((completion) => !localTaskIds.has(completion.taskId)); + const durable = await params.pendingCompletions.list({ userId, conversationId }); + pendingCompletions = durable.completions.filter( + (completion) => !localTaskIds.has(completion.taskId), + ); + if (durable.complete) { + durablePendingTaskIds = new Set(durable.completions.map(({ taskId }) => taskId)); + } else { + listWarnings.push( + 'More undelivered results exist than could be listed; some not shown may still arrive as new turns.', + ); + } } catch (error) { logger.warn('[background] Failed to list undelivered background completions:', error); listWarnings.push( @@ -2709,16 +2743,28 @@ export async function runCheckBackgroundTask(params: { } } const ordinaryTasks = [ - ...tasks.map((task) => serializeTask(task, { includeResult: false })), + ...tasks.map((task) => + reconcileDelivery(serializeTask(task, { includeResult: false }), durablePendingTaskIds), + ), ...pendingCompletions.map(serializePendingCompletion), ]; + if (completionWakeups) { + subagentTasks = subagentTasks.map((task) => + task.status !== 'running' && task.result_available === true && task.result_claimed !== true + ? { ...task, delivery: 'pending' as const } + : task, + ); + } /** Work is outstanding until its result reaches the conversation: a finished * task with a pending delivery is still going to resume the agent. */ + const isOutstanding = (task: { status: string; delivery?: string }): boolean => + task.status === 'running' || task.delivery === 'pending'; const outstanding = - ordinaryTasks.filter((task) => task.status === 'running' || task.delivery === 'pending') - .length + subagentTasks.filter((task) => task.status === 'running').length; + ordinaryTasks.filter(isOutstanding).length + subagentTasks.filter(isOutstanding).length; const guidance = [ - ...(ordinaryTasks.some((task) => task.status !== 'running' && task.delivery === 'pending') + ...([...ordinaryTasks, ...subagentTasks].some( + (task) => task.status !== 'running' && task.delivery === 'pending', + ) ? [PENDING_DELIVERY_GUIDANCE] : []), ...(completionWakeups && subagentTasks.some((task) => task.status === 'running') diff --git a/packages/api/src/agents/backgroundCompletion.ts b/packages/api/src/agents/backgroundCompletion.ts index c27697997e2..54be1beb803 100644 --- a/packages/api/src/agents/backgroundCompletion.ts +++ b/packages/api/src/agents/backgroundCompletion.ts @@ -25,6 +25,8 @@ export interface BackgroundToolWakeupRetireOptions { onlyIfUnclaimed?: boolean; /** Reconcile only after the delivery is irreversibly dead-lettered. */ onlyIfDead?: boolean; + /** Report success only when this call retired it, not when it was already delivered. */ + requireTransition?: boolean; } /** Process-local handle for the durable delivery admitted before launch. */ @@ -89,10 +91,11 @@ export type BackgroundCompletionDiscardOutcome = /** Durable view and control of one principal's undelivered background completions. */ export interface PendingBackgroundCompletionControls { + /** `complete` is false when more undelivered completions exist than were listed. */ list: (input: { userId: string; conversationId: string; - }) => Promise; + }) => Promise<{ completions: PendingBackgroundCompletion[]; complete: boolean }>; discard: (input: { userId: string; conversationId: string; diff --git a/packages/api/src/agents/backgroundCompletionWakeup.spec.ts b/packages/api/src/agents/backgroundCompletionWakeup.spec.ts index ae60ba4eb39..ecde99a153d 100644 --- a/packages/api/src/agents/backgroundCompletionWakeup.spec.ts +++ b/packages/api/src/agents/backgroundCompletionWakeup.spec.ts @@ -794,6 +794,7 @@ describe('background tool completion wakeups', () => { describe('pending background completions', () => { const dispatchedAt = new Date(NOW - 60_000); + const settled = { status: 'completed' as const, settledAt: new Date(NOW) }; const row = (overrides = {}) => ({ deliveryKey: 'delivery-key-1', taskId: 'task-1', @@ -803,24 +804,27 @@ describe('pending background completions', () => { claimedByWakeup: false, ...overrides, }); + const listing = (completions: Array>, truncated = false) => + jest.fn(async () => ({ completions, truncated })); it('lists the durable view for the owner without delivery internals', async () => { - const list = jest.fn(async () => [ - row({ result: { status: 'completed' as const, settledAt: new Date(NOW) } }), - ]); + const list = listing([row({ result: settled })], true); const pending = createPendingBackgroundCompletions({ list, retire: jest.fn() }); await expect( pending.list({ userId: 'user-1', conversationId: 'conversation-1' }), - ).resolves.toEqual([ - { - taskId: 'task-1', - toolName: 'slow_tool', - dispatchedAt, - result: { status: 'completed', settledAt: new Date(NOW) }, - claimedByWakeup: false, - }, - ]); + ).resolves.toEqual({ + completions: [ + { + taskId: 'task-1', + toolName: 'slow_tool', + dispatchedAt, + result: settled, + claimedByWakeup: false, + }, + ], + complete: false, + }); expect(list).toHaveBeenCalledWith({ user: 'user-1', conversationId: 'conversation-1', @@ -828,47 +832,36 @@ describe('pending background completions', () => { }); }); - it('discards a settled, unclaimed result by retiring only an unclaimed delivery', async () => { + it('discards a settled, unclaimed result by looking up that task and retiring it exactly', async () => { const retire = jest.fn(async () => true); - const pending = createPendingBackgroundCompletions({ - list: async () => [ - row({ result: { status: 'completed' as const, settledAt: new Date(NOW) } }), - ], - retire, - }); + const list = listing([row({ result: settled })]); + const pending = createPendingBackgroundCompletions({ list, retire }); await expect( pending.discard({ userId: 'user-1', conversationId: 'conversation-1', taskId: 'task-1' }), ).resolves.toBe('discarded'); + expect(list).toHaveBeenCalledWith({ + user: 'user-1', + conversationId: 'conversation-1', + sourceId: 'background-tool-completion', + taskId: 'task-1', + }); expect(retire).toHaveBeenCalledWith( 'delivery-key-1', 'background-tool-completion', 'background result discarded by its owner', - { onlyIfUnclaimed: true }, + { onlyIfUnclaimed: true, requireTransition: true }, ); }); it.each([ ['not_pending', [], true], ['running', [row()], true], - [ - 'delivering', - [ - row({ - result: { status: 'completed' as const, settledAt: new Date(NOW) }, - claimedByWakeup: true, - }), - ], - true, - ], - [ - 'delivering', - [row({ result: { status: 'completed' as const, settledAt: new Date(NOW) } })], - false, - ], + ['delivering', [row({ result: settled, claimedByWakeup: true })], true], + ['delivering', [row({ result: settled })], false], ])('reports %s without discarding what it cannot', async (outcome, rows, retired) => { const retire = jest.fn(async () => retired); - const pending = createPendingBackgroundCompletions({ list: async () => rows, retire }); + const pending = createPendingBackgroundCompletions({ list: listing(rows), retire }); await expect( pending.discard({ userId: 'user-1', conversationId: 'conversation-1', taskId: 'task-1' }), diff --git a/packages/api/src/agents/backgroundCompletionWakeup.ts b/packages/api/src/agents/backgroundCompletionWakeup.ts index dcc2e3243b4..9c6381bd95d 100644 --- a/packages/api/src/agents/backgroundCompletionWakeup.ts +++ b/packages/api/src/agents/backgroundCompletionWakeup.ts @@ -543,26 +543,38 @@ export function createPendingBackgroundCompletions(deps: { user: string; conversationId: string; sourceId: string; - }) => Promise>; + taskId?: string; + }) => Promise<{ + completions: Array; + truncated: boolean; + }>; retire: RetireBackgroundToolCompletion; }): PendingBackgroundCompletionControls { - const read = (input: { userId: string; conversationId: string }) => + const read = (input: { userId: string; conversationId: string; taskId?: string }) => deps.list({ user: input.userId, conversationId: input.conversationId, sourceId: BACKGROUND_TOOL_COMPLETION_SOURCE, + ...(input.taskId != null && { taskId: input.taskId }), }); return { - list: async (input) => - (await read(input)).map(({ taskId, toolName, dispatchedAt, result, claimedByWakeup }) => ({ - taskId, - toolName, - dispatchedAt, - ...(result != null && { result }), - claimedByWakeup, - })), + list: async (input) => { + const { completions, truncated } = await read(input); + return { + completions: completions.map( + ({ taskId, toolName, dispatchedAt, result, claimedByWakeup }) => ({ + taskId, + toolName, + dispatchedAt, + ...(result != null && { result }), + claimedByWakeup, + }), + ), + complete: !truncated, + }; + }, discard: async (input) => { - const completion = (await read(input)).find(({ taskId }) => taskId === input.taskId); + const [completion] = (await read(input)).completions; if (completion == null) { return 'not_pending'; } @@ -573,12 +585,13 @@ export function createPendingBackgroundCompletions(deps: { return 'delivering'; } /** Unclaimed-only: once a resolver owns the delivery its continuation can no - * longer be withdrawn, so that race reports as already delivering. */ + * longer be withdrawn, so that race, including one it already finished, + * reports as delivering rather than discarded. */ const retired = await deps.retire( completion.deliveryKey, BACKGROUND_TOOL_COMPLETION_SOURCE, 'background result discarded by its owner', - { onlyIfUnclaimed: true }, + { onlyIfUnclaimed: true, requireTransition: true }, ); return retired ? 'discarded' : 'delivering'; }, diff --git a/packages/api/src/agents/triggers/service.ts b/packages/api/src/agents/triggers/service.ts index a37e1821f35..0e41408261c 100644 --- a/packages/api/src/agents/triggers/service.ts +++ b/packages/api/src/agents/triggers/service.ts @@ -190,7 +190,7 @@ export interface AgentTriggerService { deliveryKey: string, sourceId: string, reason: string, - options?: { onlyIfUnclaimed?: boolean; onlyIfDead?: boolean }, + options?: { onlyIfUnclaimed?: boolean; onlyIfDead?: boolean; requireTransition?: boolean }, ) => Promise; renewProducerLease: (deliveryKey: string, sourceId: string, leaseUntil: Date) => Promise; persistBackgroundToolResult: (input: { @@ -699,6 +699,7 @@ export function createAgentTriggerService(deps: AgentTriggerServiceDeps = {}): A settledAt: new Date(), ...(options?.onlyIfUnclaimed === true ? { onlyIfUnclaimed: true } : {}), ...(options?.onlyIfDead === true ? { onlyIfDead: true } : {}), + ...(options?.requireTransition === true ? { requireTransition: true } : {}), }, recovery, ), diff --git a/packages/data-schemas/src/methods/triggerDelivery.spec.ts b/packages/data-schemas/src/methods/triggerDelivery.spec.ts index e6dc2f5101a..a2189da3dea 100644 --- a/packages/data-schemas/src/methods/triggerDelivery.spec.ts +++ b/packages/data-schemas/src/methods/triggerDelivery.spec.ts @@ -493,7 +493,8 @@ describe('agent trigger delivery methods', () => { sourceId: background.id, }); - expect(pending).toEqual([ + expect(pending.truncated).toBe(false); + expect(pending.completions).toEqual([ { deliveryKey: running.delivery.deliveryKey, taskId: 'task-running', @@ -550,7 +551,82 @@ describe('agent trigger delivery methods', () => { sourceId: background.id, }); - expect(pending.map((entry) => entry.deliveryKey)).toEqual([waiting.delivery.deliveryKey]); + expect(pending.completions.map((entry) => entry.deliveryKey)).toEqual([ + waiting.delivery.deliveryKey, + ]); + }); + + it('looks up one task and reports a truncated listing', async () => { + const user = new mongoose.Types.ObjectId(); + await completion(user, 'task-first'); + const second = await completion(user, 'task-second'); + + const one = await methods.listPendingAgentBackgroundToolCompletions({ + user, + conversationId: 'conversation-1', + sourceId: background.id, + taskId: 'task-second', + }); + expect(one).toEqual({ + completions: [expect.objectContaining({ deliveryKey: second.delivery.deliveryKey })], + truncated: false, + }); + + const page = await methods.listPendingAgentBackgroundToolCompletions({ + user, + conversationId: 'conversation-1', + sourceId: background.id, + limit: 1, + }); + expect(page.completions.map((entry) => entry.taskId)).toEqual(['task-first']); + expect(page.truncated).toBe(true); + }); + + it('omits capability-dead rows, which no worker will deliver', async () => { + const user = new mongoose.Types.ObjectId(); + const dead = await completion(user, 'task-dead'); + await Delivery.updateOne( + { _id: dead.delivery.id }, + { $set: { status: 'leased', capabilityStatus: 'dead' } }, + ); + + const pending = await methods.listPendingAgentBackgroundToolCompletions({ + user, + conversationId: 'conversation-1', + sourceId: background.id, + }); + + expect(pending.completions).toEqual([]); + }); + + it('distinguishes retiring a completion from finding it already delivered', async () => { + const user = new mongoose.Types.ObjectId(); + const delivered = await completion(user, 'task-already-delivered'); + await Delivery.updateOne({ _id: delivered.delivery.id }, { $set: { status: 'succeeded' } }); + const retire = (requireTransition?: true) => + methods.retireAgentTriggerDelivery({ + deliveryKey: delivered.delivery.deliveryKey, + sourceId: background.id, + settledAt: START, + reason: 'background result discarded by its owner', + onlyIfUnclaimed: true, + ...(requireTransition != null && { requireTransition }), + }); + + await expect(retire()).resolves.toBe(true); + await expect(retire(true)).resolves.toBe(false); + + const waiting = await completion(user, 'task-still-waiting'); + await expect( + methods.retireAgentTriggerDelivery({ + deliveryKey: waiting.delivery.deliveryKey, + sourceId: background.id, + settledAt: START, + reason: 'background result discarded by its owner', + onlyIfUnclaimed: true, + requireTransition: true, + }), + ).resolves.toBe(true); }); it('refuses a malformed lookup', async () => { diff --git a/packages/data-schemas/src/methods/triggerDelivery.ts b/packages/data-schemas/src/methods/triggerDelivery.ts index fd87e713e14..9f2d59219f5 100644 --- a/packages/data-schemas/src/methods/triggerDelivery.ts +++ b/packages/data-schemas/src/methods/triggerDelivery.ts @@ -44,8 +44,9 @@ export const CLAIM_CAS_MAX_ATTEMPTS = 16; /** Candidates fetched per claim read; losing claimers advance through the * batch instead of re-reading the same head-of-queue row. */ const CLAIM_CANDIDATE_BATCH = 8; -/** A conversation's undelivered completions are few; this bounds a pathological listing. */ -const MAX_PENDING_BACKGROUND_COMPLETIONS = 50; +/** Matches the per-conversation background task limit, so a listing is complete + * unless durable rows outlived that limit across restarts; it then says so. */ +const MAX_PENDING_BACKGROUND_COMPLETIONS = 200; /** Every status before a delivery settles, i.e. whose result has not reached its conversation. */ const UNDELIVERED_STATUSES: IAgentTriggerDelivery['status'][] = [ 'staging', @@ -271,6 +272,12 @@ export interface PendingAgentBackgroundToolCompletion { claimedByWakeup: boolean; } +export interface PendingAgentBackgroundToolCompletions { + completions: PendingAgentBackgroundToolCompletion[]; + /** More undelivered completions exist than were returned. */ + truncated: boolean; +} + export interface AgentTriggerDeliveryMethods { ensureAgentTriggerDeliveryIndexes: () => Promise; enqueueAgentTriggerDelivery: ( @@ -319,6 +326,9 @@ export interface AgentTriggerDeliveryMethods { /** Accept transport success without a terminal handling receipt, unless the * delivery explicitly keeps its lane open for terminal handling. */ allowSucceeded?: boolean; + /** True only when this call retired the delivery, not when it had already + * succeeded, e.g. delivered by a resolver that won the race. */ + requireTransition?: boolean; }, recovery?: { required: boolean }, ) => Promise; @@ -336,8 +346,10 @@ export interface AgentTriggerDeliveryMethods { user: string | Types.ObjectId; conversationId: string; sourceId: string; + /** One task's completion, e.g. to discard it. */ + taskId?: string; limit?: number; - }) => Promise; + }) => Promise; persistAgentBackgroundToolResult: ( input: PersistAgentBackgroundToolResultInput, ) => Promise; @@ -2132,6 +2144,9 @@ export function createAgentTriggerDeliveryMethods( /** Accept transport success without a terminal handling receipt, unless the * delivery explicitly keeps its lane open for terminal handling. */ allowSucceeded?: boolean; + /** True only when this call retired the delivery, not when it had already + * succeeded, e.g. delivered by a resolver that won the race. */ + requireTransition?: boolean; }, recovery?: { required: boolean }, ): Promise { @@ -2219,6 +2234,9 @@ export function createAgentTriggerDeliveryMethods( } return true; } + if (input.requireTransition === true) { + return false; + } return ( (await Delivery().exists({ deliveryKey: input.deliveryKey, @@ -2327,14 +2345,19 @@ export function createAgentTriggerDeliveryMethods( user: string | Types.ObjectId; conversationId: string; sourceId: string; + taskId?: string; limit?: number; - }): Promise { - const limit = input.limit ?? MAX_PENDING_BACKGROUND_COMPLETIONS; + }): Promise { + const limit = Math.min( + input.limit ?? MAX_PENDING_BACKGROUND_COMPLETIONS, + MAX_PENDING_BACKGROUND_COMPLETIONS, + ); if ( input.conversationId.length === 0 || input.conversationId.length > 256 || input.sourceId.length === 0 || input.sourceId.length > 256 || + (input.taskId != null && (input.taskId.length === 0 || input.taskId.length > 256)) || !Number.isSafeInteger(limit) || limit <= 0 ) { @@ -2346,11 +2369,17 @@ export function createAgentTriggerDeliveryMethods( 'envelope.event.source.type': 'internal', 'envelope.event.source.id': input.sourceId, 'envelope.target.conversationId': input.conversationId, + ...(input.taskId != null && { 'envelope.event.payload.taskId': input.taskId }), status: { $in: UNDELIVERED_STATUSES }, + /** Capability-dead rows are dead letters to every worker version. */ + capabilityStatus: { $ne: 'dead' }, }) - .select('+backgroundToolResult deliveryKey createdAt envelope.event.payload') + .select( + 'deliveryKey createdAt envelope.event.payload ' + + 'backgroundToolResult.status backgroundToolResult.settledAt backgroundToolResult.resultClaim', + ) .sort({ createdAt: 1, _id: 1 }) - .limit(Math.min(limit, MAX_PENDING_BACKGROUND_COMPLETIONS)) + .limit(limit + 1) .lean< Array< Pick & { @@ -2358,7 +2387,7 @@ export function createAgentTriggerDeliveryMethods( } > >(); - return rows.flatMap((row) => { + const completions = rows.slice(0, limit).flatMap((row) => { const payload = row.envelope?.event?.payload; const taskId = payload?.taskId; const toolCallId = payload?.toolCallId; @@ -2386,6 +2415,7 @@ export function createAgentTriggerDeliveryMethods( }, ]; }); + return { completions, truncated: rows.length > limit }; } /** Stores terminal output on the pre-admitted delivery before attempting the From ecb5d3e59438624ead4c275084e9c132bc499181 Mon Sep 17 00:00:00 2001 From: Danny Avila Date: Thu, 24 Sep 2026 19:31:51 -0400 Subject: [PATCH 3/5] fix: Settle manually claimed remote deliveries and keep pending controls while wake-ups drain --- .../services/Endpoints/agents/initialize.js | 9 +-- packages/api/src/agents/background.spec.ts | 59 ++++++++++++++++++- packages/api/src/agents/background.ts | 41 +++++++++---- .../api/src/agents/backgroundCompletion.ts | 6 ++ .../agents/backgroundCompletionWakeup.spec.ts | 29 +++++++++ .../src/agents/backgroundCompletionWakeup.ts | 12 ++++ .../src/methods/triggerDelivery.spec.ts | 15 +++++ .../src/methods/triggerDelivery.ts | 3 + 8 files changed, 155 insertions(+), 19 deletions(-) diff --git a/api/server/services/Endpoints/agents/initialize.js b/api/server/services/Endpoints/agents/initialize.js index d759c98a1a2..a3ba9cb531b 100644 --- a/api/server/services/Endpoints/agents/initialize.js +++ b/api/server/services/Endpoints/agents/initialize.js @@ -515,12 +515,9 @@ const initializeClientWithProvider = async ({ updateToolCallResult: db.updateToolCallResult, }), backgroundToolCompletion: { - ...(completionWakeupsEnabled - ? { - preregister: preregisterBackgroundToolCompletion, - pending: pendingBackgroundToolCompletions, - } - : {}), + ...(completionWakeupsEnabled ? { preregister: preregisterBackgroundToolCompletion } : {}), + /** Deliveries admitted before wake-ups were disabled still drain and still count. */ + pending: pendingBackgroundToolCompletions, persist: createBackgroundToolResultPersistence({ req, updateToolCallResult: db.updateToolCallResult, diff --git a/packages/api/src/agents/background.spec.ts b/packages/api/src/agents/background.spec.ts index ff784c4ccd5..ca6cca14c8d 100644 --- a/packages/api/src/agents/background.spec.ts +++ b/packages/api/src/agents/background.spec.ts @@ -3106,7 +3106,8 @@ describe('runCheckBackgroundTask (singleton)', () => { expect.objectContaining({ status: 'completed', result_available: true, delivery: 'pending' }), ); expect(listed.outstanding).toBe(1); - expect(listed.message).toContain('have not been delivered yet'); + expect(listed.message).toContain('Some finished subagents have not been delivered yet'); + expect(listed.message).not.toContain('cancel'); }); it('tells a wakeup-enabled parent to yield on an unchanged running subagent', async () => { @@ -3478,6 +3479,7 @@ describe('runCheckBackgroundTask delivery semantics', () => { complete: overrides.complete ?? true, })), discard: jest.fn(overrides.discard ?? (async () => 'not_pending')), + settleClaimed: jest.fn(async () => true), }) as never; function completedWithWakeup(userId: string, conversationId: string, toolCallId: string) { @@ -3739,6 +3741,61 @@ describe('runCheckBackgroundTask delivery semantics', () => { ); }); + it('retires the pending delivery of a remote task a manual poll just claimed', async () => { + const pendingCompletions = pendingControls(); + const claimBackgroundToolResult = jest.fn(async () => ({ + status: 'acquired' as const, + results: [ + { + taskId: 'remote-task', + toolName: 'slow_task', + status: 'completed' as const, + output: 'remote result', + settledAt: new Date('2026-09-24T12:00:00Z'), + }, + ], + })); + + const polled = JSON.parse( + await runCheckBackgroundTask({ + userId: 'remote-user', + conversationId: 'remote-convo', + args: { background_task_id: 'remote-task' }, + claimBackgroundToolResult: claimBackgroundToolResult as never, + pendingCompletions, + }), + ); + + expect(polled).toEqual( + expect.objectContaining({ status: 'completed', result: 'remote result' }), + ); + expect( + (pendingCompletions as unknown as { settleClaimed: jest.Mock }).settleClaimed, + ).toHaveBeenCalledWith({ + userId: 'remote-user', + conversationId: 'remote-convo', + taskId: 'remote-task', + }); + }); + + it('reports a failed discard lookup only when nothing else claims the task', async () => { + const cancelled = JSON.parse( + await runCheckBackgroundTask({ + userId: 'discard-user', + conversationId: 'discard-convo', + args: { background_task_id: 'unknown-task', action: 'cancel' }, + pendingCompletions: pendingControls({ + discard: async () => { + throw new Error('delivery store unavailable'); + }, + }), + }), + ); + + expect(cancelled).toEqual(expect.objectContaining({ status: 'unavailable' })); + expect(cancelled.message).toContain('could not be discarded right now'); + }); + it('falls through to the ordinary lookup when nothing is pending for the task', async () => { const pendingCompletions = pendingControls({ discard: async () => 'not_pending' }); diff --git a/packages/api/src/agents/background.ts b/packages/api/src/agents/background.ts index 0075e04ed36..e4c2900c6d2 100644 --- a/packages/api/src/agents/background.ts +++ b/packages/api/src/agents/background.ts @@ -1850,6 +1850,9 @@ interface SerializedBackgroundTask { error?: string; } +const SUBAGENT_PENDING_DELIVERY_GUIDANCE = + 'Some finished subagents have not been delivered yet (delivery: "pending"); each will arrive as a new turn. Poll one to collect its result now. Do not report them as finished until then.'; + const PENDING_DELIVERY_GUIDANCE = 'Some finished tasks have not been delivered yet (delivery: "pending"); each will arrive as a new turn. Poll one to collect its result now, or cancel it so it does not arrive. Do not report these tasks as finished or cancelled until then.'; @@ -2465,18 +2468,16 @@ export async function runCheckBackgroundTask(params: { return JSON.stringify(serializeTask(task, { includeResult: true })); } + /** A failed lookup must not mask a subagent the controls below can still reach. */ + let discardFailed = false; if (action === 'cancel' && params.pendingCompletions != null) { - let outcome: Awaited>; + let outcome: Awaited> = + 'not_pending'; try { outcome = await params.pendingCompletions.discard({ userId, conversationId, taskId }); } catch (error) { logger.warn(`[background] Failed to discard pending completion ${taskId}:`, error); - return JSON.stringify({ - status: 'unavailable', - background_task_id: taskId, - message: - 'The pending result could not be discarded right now. It may still arrive as a new turn; retry the cancel shortly.', - }); + discardFailed = true; } if (outcome === 'discarded') { return JSON.stringify({ @@ -2584,6 +2585,15 @@ export async function runCheckBackgroundTask(params: { if (durableClaim.status === 'acquired') { const durableTask = durableClaim.results.find((result) => result.taskId === taskId); if (durableTask != null) { + /** The result reaches the agent here, so its automatic delivery is redundant. */ + await params.pendingCompletions + ?.settleClaimed({ userId, conversationId, taskId }) + .catch((error: unknown) => + logger.warn( + `[background] Failed to retire the delivery of manually claimed task ${taskId}:`, + error, + ), + ); return JSON.stringify(serializeDurableTask(durableTask)); } return JSON.stringify({ @@ -2674,6 +2684,14 @@ export async function runCheckBackgroundTask(params: { } } + if (discardFailed) { + return JSON.stringify({ + status: 'unavailable', + background_task_id: taskId, + message: + 'The pending result could not be discarded right now. It may still arrive as a new turn; retry the cancel shortly.', + }); + } return JSON.stringify({ status: 'not_found', background_task_id: taskId, @@ -2761,12 +2779,11 @@ export async function runCheckBackgroundTask(params: { task.status === 'running' || task.delivery === 'pending'; const outstanding = ordinaryTasks.filter(isOutstanding).length + subagentTasks.filter(isOutstanding).length; + const isFinishedPending = (task: { status: string; delivery?: string }): boolean => + task.status !== 'running' && task.delivery === 'pending'; const guidance = [ - ...([...ordinaryTasks, ...subagentTasks].some( - (task) => task.status !== 'running' && task.delivery === 'pending', - ) - ? [PENDING_DELIVERY_GUIDANCE] - : []), + ...(ordinaryTasks.some(isFinishedPending) ? [PENDING_DELIVERY_GUIDANCE] : []), + ...(subagentTasks.some(isFinishedPending) ? [SUBAGENT_PENDING_DELIVERY_GUIDANCE] : []), ...(completionWakeups && subagentTasks.some((task) => task.status === 'running') ? [SUBAGENT_WAKEUP_GUIDANCE] : []), diff --git a/packages/api/src/agents/backgroundCompletion.ts b/packages/api/src/agents/backgroundCompletion.ts index 54be1beb803..36702b34d14 100644 --- a/packages/api/src/agents/backgroundCompletion.ts +++ b/packages/api/src/agents/backgroundCompletion.ts @@ -101,4 +101,10 @@ export interface PendingBackgroundCompletionControls { conversationId: string; taskId: string; }) => Promise; + /** Retires a task's pending delivery after a manual poll claimed its result. */ + settleClaimed: (input: { + userId: string; + conversationId: string; + taskId: string; + }) => Promise; } diff --git a/packages/api/src/agents/backgroundCompletionWakeup.spec.ts b/packages/api/src/agents/backgroundCompletionWakeup.spec.ts index ecde99a153d..7f114a90dda 100644 --- a/packages/api/src/agents/backgroundCompletionWakeup.spec.ts +++ b/packages/api/src/agents/backgroundCompletionWakeup.spec.ts @@ -854,6 +854,35 @@ describe('pending background completions', () => { ); }); + it("retires a manually claimed task's delivery only while no wake-up holds it", async () => { + const retire = jest.fn(async () => true); + const pending = createPendingBackgroundCompletions({ + list: listing([row({ result: settled })]), + retire, + }); + + await expect( + pending.settleClaimed({ + userId: 'user-1', + conversationId: 'conversation-1', + taskId: 'task-1', + }), + ).resolves.toBe(true); + expect(retire).toHaveBeenCalledWith( + 'delivery-key-1', + 'background-tool-completion', + 'completion claimed by manual poll', + { onlyIfUnclaimed: true }, + ); + await expect( + createPendingBackgroundCompletions({ list: listing([]), retire }).settleClaimed({ + userId: 'user-1', + conversationId: 'conversation-1', + taskId: 'task-1', + }), + ).resolves.toBe(false); + }); + it.each([ ['not_pending', [], true], ['running', [row()], true], diff --git a/packages/api/src/agents/backgroundCompletionWakeup.ts b/packages/api/src/agents/backgroundCompletionWakeup.ts index 9c6381bd95d..fbba7b0e425 100644 --- a/packages/api/src/agents/backgroundCompletionWakeup.ts +++ b/packages/api/src/agents/backgroundCompletionWakeup.ts @@ -595,6 +595,18 @@ export function createPendingBackgroundCompletions(deps: { ); return retired ? 'discarded' : 'delivering'; }, + settleClaimed: async (input) => { + const [completion] = (await read(input)).completions; + if (completion == null) { + return false; + } + return deps.retire( + completion.deliveryKey, + BACKGROUND_TOOL_COMPLETION_SOURCE, + 'completion claimed by manual poll', + { onlyIfUnclaimed: true }, + ); + }, }; } diff --git a/packages/data-schemas/src/methods/triggerDelivery.spec.ts b/packages/data-schemas/src/methods/triggerDelivery.spec.ts index a2189da3dea..4cd6034aa52 100644 --- a/packages/data-schemas/src/methods/triggerDelivery.spec.ts +++ b/packages/data-schemas/src/methods/triggerDelivery.spec.ts @@ -582,6 +582,21 @@ describe('agent trigger delivery methods', () => { expect(page.truncated).toBe(true); }); + it('omits legacy rows whose results live only on the parent message', async () => { + const user = new mongoose.Types.ObjectId(); + await completion(user, 'task-legacy', { + requiredWorkerCapability: AGENT_TRIGGER_WORKER_CAPABILITY_BACKGROUND_COMPLETION_V1, + }); + + const pending = await methods.listPendingAgentBackgroundToolCompletions({ + user, + conversationId: 'conversation-1', + sourceId: background.id, + }); + + expect(pending.completions).toEqual([]); + }); + it('omits capability-dead rows, which no worker will deliver', async () => { const user = new mongoose.Types.ObjectId(); const dead = await completion(user, 'task-dead'); diff --git a/packages/data-schemas/src/methods/triggerDelivery.ts b/packages/data-schemas/src/methods/triggerDelivery.ts index 9f2d59219f5..2d05863a7d9 100644 --- a/packages/data-schemas/src/methods/triggerDelivery.ts +++ b/packages/data-schemas/src/methods/triggerDelivery.ts @@ -2370,6 +2370,9 @@ export function createAgentTriggerDeliveryMethods( 'envelope.event.source.id': input.sourceId, 'envelope.target.conversationId': input.conversationId, ...(input.taskId != null && { 'envelope.event.payload.taskId': input.taskId }), + /** Legacy rows keep results only on the parent message, so a missing + * receipt cannot tell running from finished; they drain on their own path. */ + requiredWorkerCapability: AGENT_TRIGGER_WORKER_CAPABILITY_BACKGROUND_COMPLETION_RECEIPT_V2, status: { $in: UNDELIVERED_STATUSES }, /** Capability-dead rows are dead letters to every worker version. */ capabilityStatus: { $ne: 'dead' }, From 8887c467ec9bc364d4d5d9ccda330cc5b39b00ff Mon Sep 17 00:00:00 2001 From: Danny Avila Date: Thu, 24 Sep 2026 19:53:05 -0400 Subject: [PATCH 4/5] fix: Report dead-lettered deliveries as failed and read subagent wake-ups from the durable store --- .../Endpoints/agents/backgroundCompletion.js | 6 +- packages/api/src/agents/background.spec.ts | 59 ++++++++++++--- packages/api/src/agents/background.ts | 61 ++++++++++++---- .../api/src/agents/backgroundCompletion.ts | 11 ++- .../agents/backgroundCompletionWakeup.spec.ts | 34 +++++++-- .../src/agents/backgroundCompletionWakeup.ts | 18 ++++- .../src/methods/triggerDelivery.spec.ts | 32 +++++++- .../src/methods/triggerDelivery.ts | 73 +++++++++++++++++-- 8 files changed, 254 insertions(+), 40 deletions(-) diff --git a/api/server/services/Endpoints/agents/backgroundCompletion.js b/api/server/services/Endpoints/agents/backgroundCompletion.js index e048d3aa7b7..730a8467d7f 100644 --- a/api/server/services/Endpoints/agents/backgroundCompletion.js +++ b/api/server/services/Endpoints/agents/backgroundCompletion.js @@ -5,7 +5,10 @@ const { createBackgroundToolResultHandler, claimBackgroundToolResult: claimResult, } = require('@librechat/api'); -const { listPendingAgentBackgroundToolCompletions } = require('~/models'); +const { + listPendingAgentBackgroundToolCompletions, + listUndeliveredAgentTriggerTaskIds, +} = require('~/models'); const { enqueueAgentTrigger, persistAgentBackgroundToolResult, @@ -25,6 +28,7 @@ const preregisterBackgroundToolCompletion = createBackgroundToolCompletionWakeup const pendingBackgroundToolCompletions = createPendingBackgroundCompletions({ list: listPendingAgentBackgroundToolCompletions, + listTaskIds: listUndeliveredAgentTriggerTaskIds, retire: retireAgentTrigger, }); diff --git a/packages/api/src/agents/background.spec.ts b/packages/api/src/agents/background.spec.ts index ca6cca14c8d..6ad16653659 100644 --- a/packages/api/src/agents/background.spec.ts +++ b/packages/api/src/agents/background.spec.ts @@ -3092,16 +3092,32 @@ describe('runCheckBackgroundTask (singleton)', () => { await new Promise((resolve) => setImmediate(resolve)); } - const listed = JSON.parse( - await runCheckBackgroundTask({ - userId: 'owner', - conversationId: 'settled-subagent-parent', - agentId: 'agent_parent', - args: {}, - subagentTasks, - }), - ); + const listWithWakeups = async (subagentWakeups: string[]) => + JSON.parse( + await runCheckBackgroundTask({ + userId: 'owner', + conversationId: 'settled-subagent-parent', + agentId: 'agent_parent', + args: {}, + subagentTasks, + pendingCompletions: { + list: jest.fn(async () => ({ completions: [], deadTaskIds: [], complete: true })), + listSubagentWakeups: jest.fn(async () => ({ + taskIds: subagentWakeups, + complete: true, + })), + discard: jest.fn(async () => 'not_pending' as const), + settleClaimed: jest.fn(async () => false), + }, + }), + ); + /** No durable wake-up (admitted poll-only): nothing will arrive, so nothing is pending. */ + const pollOnly = await listWithWakeups([]); + expect(pollOnly.tasks[0].delivery).toBeUndefined(); + expect(pollOnly.outstanding).toBe(0); + + const listed = await listWithWakeups([started.task.taskId]); expect(listed.tasks[0]).toEqual( expect.objectContaining({ status: 'completed', result_available: true, delivery: 'pending' }), ); @@ -3470,14 +3486,21 @@ describe('runCheckBackgroundTask delivery semantics', () => { overrides: { list?: () => Promise; complete?: boolean; + deadTaskIds?: string[]; + subagentWakeups?: string[]; discard?: () => Promise; } = {}, ) => ({ list: jest.fn(async () => ({ completions: await (overrides.list ?? (async () => []))(), + deadTaskIds: overrides.deadTaskIds ?? [], complete: overrides.complete ?? true, })), + listSubagentWakeups: jest.fn(async () => ({ + taskIds: overrides.subagentWakeups ?? [], + complete: true, + })), discard: jest.fn(overrides.discard ?? (async () => 'not_pending')), settleClaimed: jest.fn(async () => true), }) as never; @@ -3701,6 +3724,24 @@ describe('runCheckBackgroundTask delivery semantics', () => { expect(listed.outstanding).toBe(1); }); + it('reports a local task whose automatic delivery dead-lettered as failed, not delivered', async () => { + const taskId = completedWithWakeup('dead-user', 'dead-convo', 'dead-call'); + const listed = JSON.parse( + await runCheckBackgroundTask({ + userId: 'dead-user', + conversationId: 'dead-convo', + args: {}, + pendingCompletions: pendingControls({ deadTaskIds: [taskId] }), + }), + ); + + expect(listed.tasks[0]).toEqual( + expect.objectContaining({ background_task_id: taskId, delivery: 'failed' }), + ); + expect(listed.outstanding).toBe(1); + expect(listed.message).toContain('Automatic delivery failed'); + }); + it('keeps the local view and warns when the durable listing is incomplete', async () => { const taskId = completedWithWakeup('truncated-user', 'truncated-convo', 'truncated-call'); const listed = JSON.parse( diff --git a/packages/api/src/agents/background.ts b/packages/api/src/agents/background.ts index e4c2900c6d2..4d591b28978 100644 --- a/packages/api/src/agents/background.ts +++ b/packages/api/src/agents/background.ts @@ -1845,11 +1845,14 @@ interface SerializedBackgroundTask { /** Whether the result has reached the conversation. `pending` results still * arrive as a new turn unless polled or cancelled first. Absent when the task * has no automatic delivery, so only a poll ever surfaces its result. */ - delivery?: 'pending' | 'delivered'; + delivery?: 'pending' | 'delivered' | 'failed'; note?: string; error?: string; } +const FAILED_DELIVERY_GUIDANCE = + 'Automatic delivery failed for some finished tasks (delivery: "failed"); they will not arrive as a new turn. Poll each to collect its result.'; + const SUBAGENT_PENDING_DELIVERY_GUIDANCE = 'Some finished subagents have not been delivered yet (delivery: "pending"); each will arrive as a new turn. Poll one to collect its result now. Do not report them as finished until then.'; @@ -1949,13 +1952,15 @@ function serializePendingCompletion( function reconcileDelivery( task: SerializedBackgroundTask, durablePendingTaskIds: ReadonlySet | undefined, + deadTaskIds: ReadonlySet, ): SerializedBackgroundTask { - if ( - task.delivery !== 'pending' || - task.status === 'running' || - durablePendingTaskIds == null || - durablePendingTaskIds.has(task.background_task_id) - ) { + if (task.delivery !== 'pending' || task.status === 'running') { + return task; + } + if (deadTaskIds.has(task.background_task_id)) { + return { ...task, delivery: 'failed' }; + } + if (durablePendingTaskIds == null || durablePendingTaskIds.has(task.background_task_id)) { return task; } return { ...task, delivery: 'delivered' }; @@ -2713,10 +2718,12 @@ export async function runCheckBackgroundTask(params: { /** Undelivered task ids from the durable store, when the listing was complete: * a local task absent from it was delivered on this or another replica. */ let durablePendingTaskIds: ReadonlySet | undefined; + let deadTaskIds: ReadonlySet = new Set(); if (params.pendingCompletions != null) { try { const localTaskIds = new Set(tasks.map((task) => task.id)); const durable = await params.pendingCompletions.list({ userId, conversationId }); + deadTaskIds = new Set(durable.deadTaskIds); pendingCompletions = durable.completions.filter( (completion) => !localTaskIds.has(completion.taskId), ); @@ -2762,27 +2769,51 @@ export async function runCheckBackgroundTask(params: { } const ordinaryTasks = [ ...tasks.map((task) => - reconcileDelivery(serializeTask(task, { includeResult: false }), durablePendingTaskIds), + reconcileDelivery( + serializeTask(task, { includeResult: false }), + durablePendingTaskIds, + deadTaskIds, + ), ), ...pendingCompletions.map(serializePendingCompletion), ]; - if (completionWakeups) { - subagentTasks = subagentTasks.map((task) => - task.status !== 'running' && task.result_available === true && task.result_claimed !== true - ? { ...task, delivery: 'pending' as const } - : task, - ); + /** Pending only with durable evidence: whether a subagent's wake-up exists depends on + * the policy when it was admitted, not on this request's configuration. */ + const finishedSubagents = subagentTasks.filter( + (task) => task.status !== 'running' && task.result_claimed !== true, + ); + if (finishedSubagents.length > 0 && params.pendingCompletions != null) { + try { + const wakeups = await params.pendingCompletions.listSubagentWakeups({ + userId, + conversationId, + }); + const waiting = new Set(wakeups.taskIds); + subagentTasks = subagentTasks.map((task) => + task.status !== 'running' && + task.result_claimed !== true && + waiting.has(task.background_task_id) + ? { ...task, delivery: 'pending' as const } + : task, + ); + } catch (error) { + logger.warn('[background] Failed to list undelivered subagent completions:', error); + listWarnings.push( + 'Undelivered subagent results could not be checked; some may still arrive as new turns.', + ); + } } /** Work is outstanding until its result reaches the conversation: a finished * task with a pending delivery is still going to resume the agent. */ const isOutstanding = (task: { status: string; delivery?: string }): boolean => - task.status === 'running' || task.delivery === 'pending'; + task.status === 'running' || task.delivery === 'pending' || task.delivery === 'failed'; const outstanding = ordinaryTasks.filter(isOutstanding).length + subagentTasks.filter(isOutstanding).length; const isFinishedPending = (task: { status: string; delivery?: string }): boolean => task.status !== 'running' && task.delivery === 'pending'; const guidance = [ ...(ordinaryTasks.some(isFinishedPending) ? [PENDING_DELIVERY_GUIDANCE] : []), + ...(ordinaryTasks.some((task) => task.delivery === 'failed') ? [FAILED_DELIVERY_GUIDANCE] : []), ...(subagentTasks.some(isFinishedPending) ? [SUBAGENT_PENDING_DELIVERY_GUIDANCE] : []), ...(completionWakeups && subagentTasks.some((task) => task.status === 'running') ? [SUBAGENT_WAKEUP_GUIDANCE] diff --git a/packages/api/src/agents/backgroundCompletion.ts b/packages/api/src/agents/backgroundCompletion.ts index 36702b34d14..2957048afeb 100644 --- a/packages/api/src/agents/backgroundCompletion.ts +++ b/packages/api/src/agents/backgroundCompletion.ts @@ -92,10 +92,17 @@ export type BackgroundCompletionDiscardOutcome = /** Durable view and control of one principal's undelivered background completions. */ export interface PendingBackgroundCompletionControls { /** `complete` is false when more undelivered completions exist than were listed. */ - list: (input: { + list: (input: { userId: string; conversationId: string }) => Promise<{ + completions: PendingBackgroundCompletion[]; + /** Tasks whose automatic delivery dead-lettered; only a poll recovers them. */ + deadTaskIds: string[]; + complete: boolean; + }>; + /** Subagent tasks whose completion wake-up has not been delivered yet. */ + listSubagentWakeups: (input: { userId: string; conversationId: string; - }) => Promise<{ completions: PendingBackgroundCompletion[]; complete: boolean }>; + }) => Promise<{ taskIds: string[]; complete: boolean }>; discard: (input: { userId: string; conversationId: string; diff --git a/packages/api/src/agents/backgroundCompletionWakeup.spec.ts b/packages/api/src/agents/backgroundCompletionWakeup.spec.ts index 7f114a90dda..7a84ac953d6 100644 --- a/packages/api/src/agents/backgroundCompletionWakeup.spec.ts +++ b/packages/api/src/agents/backgroundCompletionWakeup.spec.ts @@ -805,15 +805,17 @@ describe('pending background completions', () => { ...overrides, }); const listing = (completions: Array>, truncated = false) => - jest.fn(async () => ({ completions, truncated })); + jest.fn(async () => ({ completions, deadTaskIds: ['task-dead'], truncated })); + const listTaskIds = jest.fn(async () => ({ taskIds: ['child-1'], truncated: false })); it('lists the durable view for the owner without delivery internals', async () => { const list = listing([row({ result: settled })], true); - const pending = createPendingBackgroundCompletions({ list, retire: jest.fn() }); + const pending = createPendingBackgroundCompletions({ list, listTaskIds, retire: jest.fn() }); await expect( pending.list({ userId: 'user-1', conversationId: 'conversation-1' }), ).resolves.toEqual({ + deadTaskIds: ['task-dead'], completions: [ { taskId: 'task-1', @@ -835,7 +837,7 @@ describe('pending background completions', () => { it('discards a settled, unclaimed result by looking up that task and retiring it exactly', async () => { const retire = jest.fn(async () => true); const list = listing([row({ result: settled })]); - const pending = createPendingBackgroundCompletions({ list, retire }); + const pending = createPendingBackgroundCompletions({ list, listTaskIds, retire }); await expect( pending.discard({ userId: 'user-1', conversationId: 'conversation-1', taskId: 'task-1' }), @@ -858,6 +860,7 @@ describe('pending background completions', () => { const retire = jest.fn(async () => true); const pending = createPendingBackgroundCompletions({ list: listing([row({ result: settled })]), + listTaskIds, retire, }); @@ -875,7 +878,7 @@ describe('pending background completions', () => { { onlyIfUnclaimed: true }, ); await expect( - createPendingBackgroundCompletions({ list: listing([]), retire }).settleClaimed({ + createPendingBackgroundCompletions({ list: listing([]), listTaskIds, retire }).settleClaimed({ userId: 'user-1', conversationId: 'conversation-1', taskId: 'task-1', @@ -883,6 +886,23 @@ describe('pending background completions', () => { ).resolves.toBe(false); }); + it('lists undelivered subagent wake-ups from their own source', async () => { + const pending = createPendingBackgroundCompletions({ + list: listing([]), + listTaskIds, + retire: jest.fn(), + }); + + await expect( + pending.listSubagentWakeups({ userId: 'user-1', conversationId: 'conversation-1' }), + ).resolves.toEqual({ taskIds: ['child-1'], complete: true }); + expect(listTaskIds).toHaveBeenCalledWith({ + user: 'user-1', + conversationId: 'conversation-1', + sourceId: 'subagent-completion', + }); + }); + it.each([ ['not_pending', [], true], ['running', [row()], true], @@ -890,7 +910,11 @@ describe('pending background completions', () => { ['delivering', [row({ result: settled })], false], ])('reports %s without discarding what it cannot', async (outcome, rows, retired) => { const retire = jest.fn(async () => retired); - const pending = createPendingBackgroundCompletions({ list: listing(rows), retire }); + const pending = createPendingBackgroundCompletions({ + list: listing(rows), + listTaskIds, + retire, + }); await expect( pending.discard({ userId: 'user-1', conversationId: 'conversation-1', taskId: 'task-1' }), diff --git a/packages/api/src/agents/backgroundCompletionWakeup.ts b/packages/api/src/agents/backgroundCompletionWakeup.ts index fbba7b0e425..21ecb5b4dc2 100644 --- a/packages/api/src/agents/backgroundCompletionWakeup.ts +++ b/packages/api/src/agents/backgroundCompletionWakeup.ts @@ -24,6 +24,7 @@ import type { AgentContinueTriggerEnvelope } from './triggers/envelope'; import type { AgentTriggerDispatchContext } from './triggers/dispatch'; import type { AgentTriggerEnqueueOptions } from './triggers/delivery'; import { BACKGROUND_TOOL_PRODUCER_LEASE_MS } from './backgroundCompletion'; +import { SUBAGENT_COMPLETION_SOURCE } from './subagentCompletionWakeup'; import { createAgentTriggerEnvelope } from './triggers/envelope'; import { AgentTriggerExecutionError } from './triggers/host'; import { truncateMiddle } from '~/utils'; @@ -546,8 +547,14 @@ export function createPendingBackgroundCompletions(deps: { taskId?: string; }) => Promise<{ completions: Array; + deadTaskIds: string[]; truncated: boolean; }>; + listTaskIds: (input: { + user: string; + conversationId: string; + sourceId: string; + }) => Promise<{ taskIds: string[]; truncated: boolean }>; retire: RetireBackgroundToolCompletion; }): PendingBackgroundCompletionControls { const read = (input: { userId: string; conversationId: string; taskId?: string }) => @@ -559,8 +566,9 @@ export function createPendingBackgroundCompletions(deps: { }); return { list: async (input) => { - const { completions, truncated } = await read(input); + const { completions, deadTaskIds, truncated } = await read(input); return { + deadTaskIds, completions: completions.map( ({ taskId, toolName, dispatchedAt, result, claimedByWakeup }) => ({ taskId, @@ -595,6 +603,14 @@ export function createPendingBackgroundCompletions(deps: { ); return retired ? 'discarded' : 'delivering'; }, + listSubagentWakeups: async (input) => { + const { taskIds, truncated } = await deps.listTaskIds({ + user: input.userId, + conversationId: input.conversationId, + sourceId: SUBAGENT_COMPLETION_SOURCE, + }); + return { taskIds, complete: !truncated }; + }, settleClaimed: async (input) => { const [completion] = (await read(input)).completions; if (completion == null) { diff --git a/packages/data-schemas/src/methods/triggerDelivery.spec.ts b/packages/data-schemas/src/methods/triggerDelivery.spec.ts index 4cd6034aa52..ec403659693 100644 --- a/packages/data-schemas/src/methods/triggerDelivery.spec.ts +++ b/packages/data-schemas/src/methods/triggerDelivery.spec.ts @@ -569,6 +569,7 @@ describe('agent trigger delivery methods', () => { }); expect(one).toEqual({ completions: [expect.objectContaining({ deliveryKey: second.delivery.deliveryKey })], + deadTaskIds: [], truncated: false, }); @@ -597,13 +598,15 @@ describe('agent trigger delivery methods', () => { expect(pending.completions).toEqual([]); }); - it('omits capability-dead rows, which no worker will deliver', async () => { + it('reports dead-lettered tasks apart from pending ones', async () => { const user = new mongoose.Types.ObjectId(); const dead = await completion(user, 'task-dead'); await Delivery.updateOne( { _id: dead.delivery.id }, { $set: { status: 'leased', capabilityStatus: 'dead' } }, ); + const deadLetter = await completion(user, 'task-dead-letter'); + await Delivery.updateOne({ _id: deadLetter.delivery.id }, { $set: { status: 'dead' } }); const pending = await methods.listPendingAgentBackgroundToolCompletions({ user, @@ -612,6 +615,33 @@ describe('agent trigger delivery methods', () => { }); expect(pending.completions).toEqual([]); + expect(pending.deadTaskIds.sort()).toEqual(['task-dead', 'task-dead-letter']); + }); + + it("lists a conversation's undelivered task ids for one source", async () => { + const user = new mongoose.Types.ObjectId(); + const subagent = { id: 'subagent-completion', type: 'internal' }; + await completion(user, 'child-waiting', { + envelope: { + event: { source: subagent, payload: { taskId: 'child-waiting' } }, + target: { conversationId: 'conversation-1' }, + }, + }); + const delivered = await completion(user, 'child-delivered', { + envelope: { + event: { source: subagent, payload: { taskId: 'child-delivered' } }, + target: { conversationId: 'conversation-1' }, + }, + }); + await Delivery.updateOne({ _id: delivered.delivery.id }, { $set: { status: 'succeeded' } }); + + await expect( + methods.listUndeliveredAgentTriggerTaskIds({ + user, + conversationId: 'conversation-1', + sourceId: subagent.id, + }), + ).resolves.toEqual({ taskIds: ['child-waiting'], truncated: false }); }); it('distinguishes retiring a completion from finding it already delivered', async () => { diff --git a/packages/data-schemas/src/methods/triggerDelivery.ts b/packages/data-schemas/src/methods/triggerDelivery.ts index 2d05863a7d9..4e8ef9870ad 100644 --- a/packages/data-schemas/src/methods/triggerDelivery.ts +++ b/packages/data-schemas/src/methods/triggerDelivery.ts @@ -48,6 +48,11 @@ const CLAIM_CANDIDATE_BATCH = 8; * unless durable rows outlived that limit across restarts; it then says so. */ const MAX_PENDING_BACKGROUND_COMPLETIONS = 200; /** Every status before a delivery settles, i.e. whose result has not reached its conversation. */ +const DEAD_STATUSES: IAgentTriggerDelivery['status'][] = ['dead', 'capability_dead']; +/** Dead to every worker version, including a capability row a legacy worker still sees leased. */ +function isDeadDelivery(row: Pick): boolean { + return DEAD_STATUSES.includes(row.status) || row.capabilityStatus === 'dead'; +} const UNDELIVERED_STATUSES: IAgentTriggerDelivery['status'][] = [ 'staging', 'capability_staging', @@ -274,10 +279,17 @@ export interface PendingAgentBackgroundToolCompletion { export interface PendingAgentBackgroundToolCompletions { completions: PendingAgentBackgroundToolCompletion[]; + /** Tasks whose delivery dead-lettered: never delivered, recoverable only by a poll. */ + deadTaskIds: string[]; /** More undelivered completions exist than were returned. */ truncated: boolean; } +export interface UndeliveredAgentTriggerTaskIds { + taskIds: string[]; + truncated: boolean; +} + export interface AgentTriggerDeliveryMethods { ensureAgentTriggerDeliveryIndexes: () => Promise; enqueueAgentTriggerDelivery: ( @@ -350,6 +362,12 @@ export interface AgentTriggerDeliveryMethods { taskId?: string; limit?: number; }) => Promise; + /** Task ids of one conversation's undelivered internal deliveries from one source. */ + listUndeliveredAgentTriggerTaskIds: (input: { + user: string | Types.ObjectId; + conversationId: string; + sourceId: string; + }) => Promise; persistAgentBackgroundToolResult: ( input: PersistAgentBackgroundToolResultInput, ) => Promise; @@ -2373,25 +2391,34 @@ export function createAgentTriggerDeliveryMethods( /** Legacy rows keep results only on the parent message, so a missing * receipt cannot tell running from finished; they drain on their own path. */ requiredWorkerCapability: AGENT_TRIGGER_WORKER_CAPABILITY_BACKGROUND_COMPLETION_RECEIPT_V2, - status: { $in: UNDELIVERED_STATUSES }, - /** Capability-dead rows are dead letters to every worker version. */ - capabilityStatus: { $ne: 'dead' }, + /** Dead letters are read too, so a caller can tell "delivered" from "failed". */ + status: { $in: [...UNDELIVERED_STATUSES, ...DEAD_STATUSES] }, }) .select( - 'deliveryKey createdAt envelope.event.payload ' + + 'deliveryKey createdAt status capabilityStatus envelope.event.payload ' + 'backgroundToolResult.status backgroundToolResult.settledAt backgroundToolResult.resultClaim', ) .sort({ createdAt: 1, _id: 1 }) .limit(limit + 1) .lean< Array< - Pick & { + Pick< + IAgentTriggerDelivery, + 'deliveryKey' | 'createdAt' | 'backgroundToolResult' | 'status' | 'capabilityStatus' + > & { envelope?: { event?: { payload?: Record } }; } > >(); + const deadTaskIds: string[] = []; const completions = rows.slice(0, limit).flatMap((row) => { const payload = row.envelope?.event?.payload; + if (isDeadDelivery(row)) { + if (typeof payload?.taskId === 'string') { + deadTaskIds.push(payload.taskId); + } + return []; + } const taskId = payload?.taskId; const toolCallId = payload?.toolCallId; const toolName = payload?.toolName; @@ -2418,7 +2445,40 @@ export function createAgentTriggerDeliveryMethods( }, ]; }); - return { completions, truncated: rows.length > limit }; + return { completions, deadTaskIds, truncated: rows.length > limit }; + } + + async function listUndeliveredAgentTriggerTaskIds(input: { + user: string | Types.ObjectId; + conversationId: string; + sourceId: string; + }): Promise { + if ( + input.conversationId.length === 0 || + input.conversationId.length > 256 || + input.sourceId.length === 0 || + input.sourceId.length > 256 + ) { + throw new TypeError('Invalid undelivered task lookup'); + } + const rows = await Delivery() + .find({ + user: input.user, + 'envelope.event.source.type': 'internal', + 'envelope.event.source.id': input.sourceId, + 'envelope.target.conversationId': input.conversationId, + status: { $in: UNDELIVERED_STATUSES }, + capabilityStatus: { $ne: 'dead' }, + }) + .select('envelope.event.payload.taskId') + .sort({ createdAt: 1, _id: 1 }) + .limit(MAX_PENDING_BACKGROUND_COMPLETIONS + 1) + .lean>(); + const taskIds = rows + .slice(0, MAX_PENDING_BACKGROUND_COMPLETIONS) + .map((row) => row.envelope?.event?.payload?.taskId) + .filter((taskId): taskId is string => typeof taskId === 'string'); + return { taskIds, truncated: rows.length > MAX_PENDING_BACKGROUND_COMPLETIONS }; } /** Stores terminal output on the pre-admitted delivery before attempting the @@ -4184,6 +4244,7 @@ export function createAgentTriggerDeliveryMethods( renewAgentTriggerDeliveryProducerLease, getAgentTriggerDeliveryProducerLease, listPendingAgentBackgroundToolCompletions, + listUndeliveredAgentTriggerTaskIds, persistAgentBackgroundToolResult, getAgentBackgroundToolResult, getAgentBackgroundToolResultClaim, From bb0382355ddc4aa8e184b85137cf7d95b05dc4fd Mon Sep 17 00:00:00 2001 From: Danny Avila Date: Thu, 24 Sep 2026 20:06:59 -0400 Subject: [PATCH 5/5] fix: List restored dead letters and retire deliveries on local durable claims --- packages/api/src/agents/background.spec.ts | 86 +++++++++++++++++-- packages/api/src/agents/background.ts | 35 +++++++- .../api/src/agents/backgroundCompletion.ts | 4 +- .../agents/backgroundCompletionWakeup.spec.ts | 11 ++- .../src/agents/backgroundCompletionWakeup.ts | 29 ++++--- .../src/methods/triggerDelivery.spec.ts | 10 ++- .../src/methods/triggerDelivery.ts | 22 ++--- 7 files changed, 162 insertions(+), 35 deletions(-) diff --git a/packages/api/src/agents/background.spec.ts b/packages/api/src/agents/background.spec.ts index 6ad16653659..c6f169c4761 100644 --- a/packages/api/src/agents/background.spec.ts +++ b/packages/api/src/agents/background.spec.ts @@ -3101,7 +3101,7 @@ describe('runCheckBackgroundTask (singleton)', () => { args: {}, subagentTasks, pendingCompletions: { - list: jest.fn(async () => ({ completions: [], deadTaskIds: [], complete: true })), + list: jest.fn(async () => ({ completions: [], dead: [], complete: true })), listSubagentWakeups: jest.fn(async () => ({ taskIds: subagentWakeups, complete: true, @@ -3486,7 +3486,7 @@ describe('runCheckBackgroundTask delivery semantics', () => { overrides: { list?: () => Promise; complete?: boolean; - deadTaskIds?: string[]; + dead?: unknown[]; subagentWakeups?: string[]; discard?: () => Promise; } = {}, @@ -3494,7 +3494,7 @@ describe('runCheckBackgroundTask delivery semantics', () => { ({ list: jest.fn(async () => ({ completions: await (overrides.list ?? (async () => []))(), - deadTaskIds: overrides.deadTaskIds ?? [], + dead: overrides.dead ?? [], complete: overrides.complete ?? true, })), listSubagentWakeups: jest.fn(async () => ({ @@ -3731,17 +3731,93 @@ describe('runCheckBackgroundTask delivery semantics', () => { userId: 'dead-user', conversationId: 'dead-convo', args: {}, - pendingCompletions: pendingControls({ deadTaskIds: [taskId] }), + pendingCompletions: pendingControls({ + dead: [ + { + taskId, + toolName: 'bash_tool', + dispatchedAt: new Date('2026-09-24T12:00:00Z'), + claimedByWakeup: false, + }, + { + taskId: 'restored-dead-task', + toolName: 'slow_task', + dispatchedAt: new Date('2026-09-24T11:00:00Z'), + result: { status: 'completed', settledAt: new Date('2026-09-24T11:01:00Z') }, + claimedByWakeup: false, + }, + ], + }), }), ); expect(listed.tasks[0]).toEqual( expect.objectContaining({ background_task_id: taskId, delivery: 'failed' }), ); - expect(listed.outstanding).toBe(1); + /** A dead letter this process no longer holds is listed from the durable store. */ + expect(listed.tasks[1]).toEqual( + expect.objectContaining({ + background_task_id: 'restored-dead-task', + status: 'completed', + delivery: 'failed', + }), + ); + expect(listed.outstanding).toBe(2); expect(listed.message).toContain('Automatic delivery failed'); }); + it('retires the pending delivery when a local poll claims the durable result', async () => { + const created = backgroundTaskRegistry.create({ + userId: 'local-claim-user', + conversationId: 'local-claim-convo', + toolCallId: 'local-claim-call', + toolName: 'bash_tool', + messageId: 'local-claim-message', + }); + if ('atCapacity' in created) { + throw new Error('unexpected capacity'); + } + const retire = jest.fn(async () => true); + backgroundTaskRegistry.markCompletionWakeup( + 'local-claim-user', + 'local-claim-convo', + created.task.id, + { + renew: jest.fn(async () => true), + retire, + }, + ); + backgroundTaskRegistry.complete('local-claim-user', 'local-claim-convo', created.task.id, { + content: 'finished output', + }); + + const polled = JSON.parse( + await runCheckBackgroundTask({ + userId: 'local-claim-user', + conversationId: 'local-claim-convo', + args: { background_task_id: created.task.id }, + claimBackgroundToolResult: jest.fn(async () => ({ + status: 'acquired' as const, + results: [], + })) as never, + }), + ); + + expect(polled.status).toBe('completed'); + expect(retire).toHaveBeenCalledWith('completion claimed by manual poll', { + onlyIfUnclaimed: true, + }); + const listed = JSON.parse( + await runCheckBackgroundTask({ + userId: 'local-claim-user', + conversationId: 'local-claim-convo', + args: {}, + }), + ); + expect(listed.tasks[0].delivery).toBe('delivered'); + expect(listed.outstanding).toBe(0); + }); + it('keeps the local view and warns when the durable listing is incomplete', async () => { const taskId = completedWithWakeup('truncated-user', 'truncated-convo', 'truncated-call'); const listed = JSON.parse( diff --git a/packages/api/src/agents/background.ts b/packages/api/src/agents/background.ts index 4d591b28978..f0533e21983 100644 --- a/packages/api/src/agents/background.ts +++ b/packages/api/src/agents/background.ts @@ -1966,6 +1966,18 @@ function reconcileDelivery( return { ...task, delivery: 'delivered' }; } +/** A completion whose automatic delivery dead-lettered and that this process no + * longer holds: never delivered, so only a poll can still collect its result. */ +function serializeDeadCompletion( + completion: PendingBackgroundCompletion, +): SerializedBackgroundTask { + return { + ...serializePendingCompletion(completion), + delivery: 'failed', + note: 'Automatic delivery failed; this result will not arrive as a new turn. Poll it to collect the result.', + }; +} + function serializeTask( task: BackgroundTask, { includeResult }: { includeResult: boolean }, @@ -2325,6 +2337,23 @@ export async function runCheckBackgroundTask(params: { }); } } + if (durableClaim.status === 'acquired' && task.completionWakeupRetired !== true) { + /** The result reaches the agent here, so its automatic delivery is redundant. */ + await backgroundTaskRegistry + .retireCompletionWakeup( + userId, + conversationId, + taskId, + 'completion claimed by manual poll', + { onlyIfUnclaimed: true }, + ) + .catch((error: unknown) => + logger.warn( + `[background] Failed to retire the delivery of manually claimed task ${taskId}:`, + error, + ), + ); + } if (durableClaim.status === 'not_found' || durableClaim.status === 'not_ready') { const localReplay = task.resultClaim?.kind === 'manual' && task.resultClaim.claimId === invocationId; @@ -2719,11 +2748,14 @@ export async function runCheckBackgroundTask(params: { * a local task absent from it was delivered on this or another replica. */ let durablePendingTaskIds: ReadonlySet | undefined; let deadTaskIds: ReadonlySet = new Set(); + let deadCompletions: PendingBackgroundCompletion[] = []; if (params.pendingCompletions != null) { try { const localTaskIds = new Set(tasks.map((task) => task.id)); const durable = await params.pendingCompletions.list({ userId, conversationId }); - deadTaskIds = new Set(durable.deadTaskIds); + deadTaskIds = new Set(durable.dead.map(({ taskId }) => taskId)); + /** A dead letter this process no longer holds is still recoverable by a poll. */ + deadCompletions = durable.dead.filter((completion) => !localTaskIds.has(completion.taskId)); pendingCompletions = durable.completions.filter( (completion) => !localTaskIds.has(completion.taskId), ); @@ -2776,6 +2808,7 @@ export async function runCheckBackgroundTask(params: { ), ), ...pendingCompletions.map(serializePendingCompletion), + ...deadCompletions.map(serializeDeadCompletion), ]; /** Pending only with durable evidence: whether a subagent's wake-up exists depends on * the policy when it was admitted, not on this request's configuration. */ diff --git a/packages/api/src/agents/backgroundCompletion.ts b/packages/api/src/agents/backgroundCompletion.ts index 2957048afeb..fe85447463f 100644 --- a/packages/api/src/agents/backgroundCompletion.ts +++ b/packages/api/src/agents/backgroundCompletion.ts @@ -94,8 +94,8 @@ export interface PendingBackgroundCompletionControls { /** `complete` is false when more undelivered completions exist than were listed. */ list: (input: { userId: string; conversationId: string }) => Promise<{ completions: PendingBackgroundCompletion[]; - /** Tasks whose automatic delivery dead-lettered; only a poll recovers them. */ - deadTaskIds: string[]; + /** Completions whose automatic delivery dead-lettered; only a poll recovers them. */ + dead: PendingBackgroundCompletion[]; complete: boolean; }>; /** Subagent tasks whose completion wake-up has not been delivered yet. */ diff --git a/packages/api/src/agents/backgroundCompletionWakeup.spec.ts b/packages/api/src/agents/backgroundCompletionWakeup.spec.ts index 7a84ac953d6..9f452ea5b40 100644 --- a/packages/api/src/agents/backgroundCompletionWakeup.spec.ts +++ b/packages/api/src/agents/backgroundCompletionWakeup.spec.ts @@ -805,7 +805,7 @@ describe('pending background completions', () => { ...overrides, }); const listing = (completions: Array>, truncated = false) => - jest.fn(async () => ({ completions, deadTaskIds: ['task-dead'], truncated })); + jest.fn(async () => ({ completions, dead: [row({ taskId: 'task-dead' })], truncated })); const listTaskIds = jest.fn(async () => ({ taskIds: ['child-1'], truncated: false })); it('lists the durable view for the owner without delivery internals', async () => { @@ -815,7 +815,14 @@ describe('pending background completions', () => { await expect( pending.list({ userId: 'user-1', conversationId: 'conversation-1' }), ).resolves.toEqual({ - deadTaskIds: ['task-dead'], + dead: [ + { + taskId: 'task-dead', + toolName: 'slow_tool', + dispatchedAt, + claimedByWakeup: false, + }, + ], completions: [ { taskId: 'task-1', diff --git a/packages/api/src/agents/backgroundCompletionWakeup.ts b/packages/api/src/agents/backgroundCompletionWakeup.ts index 21ecb5b4dc2..f28b2aaa06c 100644 --- a/packages/api/src/agents/backgroundCompletionWakeup.ts +++ b/packages/api/src/agents/backgroundCompletionWakeup.ts @@ -547,7 +547,7 @@ export function createPendingBackgroundCompletions(deps: { taskId?: string; }) => Promise<{ completions: Array; - deadTaskIds: string[]; + dead: Array; truncated: boolean; }>; listTaskIds: (input: { @@ -566,18 +566,23 @@ export function createPendingBackgroundCompletions(deps: { }); return { list: async (input) => { - const { completions, deadTaskIds, truncated } = await read(input); + const { completions, dead, truncated } = await read(input); + const project = ({ + taskId, + toolName, + dispatchedAt, + result, + claimedByWakeup, + }: PendingBackgroundCompletion): PendingBackgroundCompletion => ({ + taskId, + toolName, + dispatchedAt, + ...(result != null && { result }), + claimedByWakeup, + }); return { - deadTaskIds, - completions: completions.map( - ({ taskId, toolName, dispatchedAt, result, claimedByWakeup }) => ({ - taskId, - toolName, - dispatchedAt, - ...(result != null && { result }), - claimedByWakeup, - }), - ), + completions: completions.map(project), + dead: dead.map(project), complete: !truncated, }; }, diff --git a/packages/data-schemas/src/methods/triggerDelivery.spec.ts b/packages/data-schemas/src/methods/triggerDelivery.spec.ts index ec403659693..3c151fd9441 100644 --- a/packages/data-schemas/src/methods/triggerDelivery.spec.ts +++ b/packages/data-schemas/src/methods/triggerDelivery.spec.ts @@ -569,7 +569,7 @@ describe('agent trigger delivery methods', () => { }); expect(one).toEqual({ completions: [expect.objectContaining({ deliveryKey: second.delivery.deliveryKey })], - deadTaskIds: [], + dead: [], truncated: false, }); @@ -615,7 +615,13 @@ describe('agent trigger delivery methods', () => { }); expect(pending.completions).toEqual([]); - expect(pending.deadTaskIds.sort()).toEqual(['task-dead', 'task-dead-letter']); + expect(pending.dead.map(({ taskId }) => taskId).sort()).toEqual([ + 'task-dead', + 'task-dead-letter', + ]); + expect(pending.dead[0]).toEqual( + expect.objectContaining({ toolName: 'slow_task', dispatchedAt: expect.any(Date) }), + ); }); it("lists a conversation's undelivered task ids for one source", async () => { diff --git a/packages/data-schemas/src/methods/triggerDelivery.ts b/packages/data-schemas/src/methods/triggerDelivery.ts index 4e8ef9870ad..df3bdf4859a 100644 --- a/packages/data-schemas/src/methods/triggerDelivery.ts +++ b/packages/data-schemas/src/methods/triggerDelivery.ts @@ -279,8 +279,8 @@ export interface PendingAgentBackgroundToolCompletion { export interface PendingAgentBackgroundToolCompletions { completions: PendingAgentBackgroundToolCompletion[]; - /** Tasks whose delivery dead-lettered: never delivered, recoverable only by a poll. */ - deadTaskIds: string[]; + /** Completions whose delivery dead-lettered: never delivered, recoverable only by a poll. */ + dead: PendingAgentBackgroundToolCompletion[]; /** More undelivered completions exist than were returned. */ truncated: boolean; } @@ -2410,15 +2410,9 @@ export function createAgentTriggerDeliveryMethods( } > >(); - const deadTaskIds: string[] = []; + const dead: PendingAgentBackgroundToolCompletion[] = []; const completions = rows.slice(0, limit).flatMap((row) => { const payload = row.envelope?.event?.payload; - if (isDeadDelivery(row)) { - if (typeof payload?.taskId === 'string') { - deadTaskIds.push(payload.taskId); - } - return []; - } const taskId = payload?.taskId; const toolCallId = payload?.toolCallId; const toolName = payload?.toolName; @@ -2443,9 +2437,15 @@ export function createAgentTriggerDeliveryMethods( }), claimedByWakeup: receipt?.resultClaim != null, }, - ]; + ].filter((completion) => { + if (!isDeadDelivery(row)) { + return true; + } + dead.push(completion); + return false; + }); }); - return { completions, deadTaskIds, truncated: rows.length > limit }; + return { completions, dead, truncated: rows.length > limit }; } async function listUndeliveredAgentTriggerTaskIds(input: {