diff --git a/api/server/services/Endpoints/agents/backgroundCompletion.js b/api/server/services/Endpoints/agents/backgroundCompletion.js index da93d982424..94d1f52eec9 100644 --- a/api/server/services/Endpoints/agents/backgroundCompletion.js +++ b/api/server/services/Endpoints/agents/backgroundCompletion.js @@ -1,9 +1,14 @@ const { createBackgroundToolCompletionWakeupHandler, createBackgroundToolDeadClaimRecovery, + createPendingBackgroundCompletions, createBackgroundToolResultHandler, claimBackgroundToolResult: claimResult, } = require('@librechat/api'); +const { + listPendingAgentBackgroundToolCompletions, + listUndeliveredAgentTriggerTaskIds, +} = require('~/models'); const { enqueueAgentTrigger, persistAgentBackgroundToolResult, @@ -23,6 +28,12 @@ const preregisterBackgroundToolCompletion = createBackgroundToolCompletionWakeup (deliveryKey) => expediteCompletionWakeups({ deliveryKeys: [deliveryKey] }), ); +const pendingBackgroundToolCompletions = createPendingBackgroundCompletions({ + list: listPendingAgentBackgroundToolCompletions, + listTaskIds: listUndeliveredAgentTriggerTaskIds, + retire: retireAgentTrigger, +}); + function createBackgroundToolResultPersistence({ req, updateToolCallResult }) { return createBackgroundToolResultHandler({ req, updateToolCallResult }); } @@ -46,6 +57,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..a3ba9cb531b 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, @@ -515,6 +516,8 @@ const initializeClientWithProvider = async ({ }), backgroundToolCompletion: { ...(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 5418fdeb3f1..c6f169c4761 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', + ); }); }); @@ -3057,6 +3063,69 @@ 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 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: [], dead: [], 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' }), + ); + expect(listed.outstanding).toBe(1); + 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 () => { const store = new InMemorySubagentTaskStore(); const subagentTasks: HostSubagentTaskConfig = { @@ -3411,3 +3480,451 @@ describe('toolOptionsSchema', () => { expect(parsed).toEqual({ run_in_background: true }); }); }); + +describe('runCheckBackgroundTask delivery semantics', () => { + const pendingControls = ( + overrides: { + list?: () => Promise; + complete?: boolean; + dead?: unknown[]; + subagentWakeups?: string[]; + discard?: () => Promise; + } = {}, + ) => + ({ + list: jest.fn(async () => ({ + completions: await (overrides.list ?? (async () => []))(), + dead: overrides.dead ?? [], + 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; + + 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('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('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({ + 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' }), + ); + /** 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( + 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('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' }); + + 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..f0533e21983 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,23 @@ 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' | '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.'; + +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 +1915,69 @@ 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.', + }; +} + +/** 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, + deadTaskIds: ReadonlySet, +): SerializedBackgroundTask { + 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' }; +} + +/** 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 }, @@ -1914,6 +1992,7 @@ function serializeTask( : {}), ...taskTimings(task), ...resultFields(task, includeResult), + ...taskDelivery(task), ...taskNote(task), ...(task.error !== undefined ? { error: task.error } : {}), }; @@ -1953,6 +2032,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( @@ -2113,6 +2194,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) ?? {}; @@ -2132,7 +2215,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, @@ -2251,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; @@ -2399,6 +2502,42 @@ 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> = + 'not_pending'; + try { + outcome = await params.pendingCompletions.discard({ userId, conversationId, taskId }); + } catch (error) { + logger.warn(`[background] Failed to discard pending completion ${taskId}:`, error); + discardFailed = true; + } + 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; @@ -2480,6 +2619,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({ @@ -2570,6 +2718,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, @@ -2586,7 +2742,37 @@ export async function runCheckBackgroundTask(params: { const tasks = backgroundTaskRegistry.list(userId, conversationId); let subagentTasks: SerializedSubagentTask[] = []; - let listWarning: string | undefined; + 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; + 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.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), + ); + 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( + '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 +2793,73 @@ 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) => + reconcileDelivery( + serializeTask(task, { includeResult: false }), + durablePendingTaskIds, + deadTaskIds, + ), + ), + ...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. */ + 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.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] + : []), + ]; 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 02d57cb7d8e..c6574030963 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. */ @@ -65,3 +67,54 @@ 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 { + /** `complete` is false when more undelivered completions exist than were listed. */ + list: (input: { userId: string; conversationId: string }) => Promise<{ + completions: PendingBackgroundCompletion[]; + /** 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. */ + listSubagentWakeups: (input: { + userId: string; + conversationId: string; + }) => Promise<{ taskIds: string[]; complete: boolean }>; + discard: (input: { + userId: string; + 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 a30d5a4b499..652ba4df253 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, @@ -891,3 +892,143 @@ 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', + toolCallId: 'call-1', + toolName: 'slow_tool', + dispatchedAt, + claimedByWakeup: false, + ...overrides, + }); + const listing = (completions: Array>, truncated = false) => + 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 () => { + const list = listing([row({ result: settled })], true); + const pending = createPendingBackgroundCompletions({ list, listTaskIds, retire: jest.fn() }); + + await expect( + pending.list({ userId: 'user-1', conversationId: 'conversation-1' }), + ).resolves.toEqual({ + dead: [ + { + taskId: 'task-dead', + toolName: 'slow_tool', + dispatchedAt, + claimedByWakeup: false, + }, + ], + completions: [ + { + taskId: 'task-1', + toolName: 'slow_tool', + dispatchedAt, + result: settled, + claimedByWakeup: false, + }, + ], + complete: false, + }); + expect(list).toHaveBeenCalledWith({ + user: 'user-1', + conversationId: 'conversation-1', + sourceId: 'background-tool-completion', + }); + }); + + 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, listTaskIds, 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, requireTransition: true }, + ); + }); + + 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 })]), + listTaskIds, + 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([]), listTaskIds, retire }).settleClaimed({ + userId: 'user-1', + conversationId: 'conversation-1', + taskId: 'task-1', + }), + ).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], + ['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: listing(rows), + listTaskIds, + 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 e5ef7ca6312..63e6f45c89e 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, @@ -23,6 +25,7 @@ import type { AgentTriggerDispatchContext } from './triggers/dispatch'; import type { AgentTriggerEnqueueOptions } from './triggers/delivery'; import { WAITING_RETRY_CAP_MS, waitingRetryAfter } from './triggers/backoff'; 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'; @@ -544,6 +547,102 @@ 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; + taskId?: string; + }) => Promise<{ + completions: Array; + dead: Array; + 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 }) => + deps.list({ + user: input.userId, + conversationId: input.conversationId, + sourceId: BACKGROUND_TOOL_COMPLETION_SOURCE, + ...(input.taskId != null && { taskId: input.taskId }), + }); + return { + list: async (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 { + completions: completions.map(project), + dead: dead.map(project), + complete: !truncated, + }; + }, + discard: async (input) => { + const [completion] = (await read(input)).completions; + 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, 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, requireTransition: true }, + ); + 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) { + return false; + } + return deps.retire( + completion.deliveryKey, + BACKGROUND_TOOL_COMPLETION_SOURCE, + 'completion claimed by manual poll', + { onlyIfUnclaimed: true }, + ); + }, + }; +} + /** 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 edbab15caa5..c0505905303 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; @@ -6671,6 +6675,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/api/src/agents/triggers/service.ts b/packages/api/src/agents/triggers/service.ts index fc2125bef36..63aa707b22c 100644 --- a/packages/api/src/agents/triggers/service.ts +++ b/packages/api/src/agents/triggers/service.ts @@ -213,7 +213,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: { @@ -770,6 +770,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 7498b932c49..a6ad23cb447 100644 --- a/packages/data-schemas/src/methods/triggerDelivery.spec.ts +++ b/packages/data-schemas/src/methods/triggerDelivery.spec.ts @@ -795,6 +795,252 @@ describe('agent trigger delivery methods', () => { }); }); + 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.truncated).toBe(false); + expect(pending.completions).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.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 })], + dead: [], + 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 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('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, + conversationId: 'conversation-1', + sourceId: background.id, + }); + + expect(pending.completions).toEqual([]); + 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 () => { + 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 () => { + 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 () => { + 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 1e2b6ead6fb..441977f8186 100644 --- a/packages/data-schemas/src/methods/triggerDelivery.ts +++ b/packages/data-schemas/src/methods/triggerDelivery.ts @@ -44,6 +44,24 @@ 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; +/** 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 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', + '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 +264,32 @@ 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 PendingAgentBackgroundToolCompletions { + completions: PendingAgentBackgroundToolCompletion[]; + /** Completions whose delivery dead-lettered: never delivered, recoverable only by a poll. */ + dead: PendingAgentBackgroundToolCompletion[]; + /** More undelivered completions exist than were returned. */ + truncated: boolean; +} + +export interface UndeliveredAgentTriggerTaskIds { + taskIds: string[]; + truncated: boolean; +} + /** Selects deferred internal deliveries whose readiness condition just changed. */ export interface ExpediteAgentTriggerDeliveriesInput { sourceIds: readonly string[]; @@ -316,6 +360,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; @@ -329,6 +376,20 @@ export interface AgentTriggerDeliveryMethods { sourceId: string; now: Date; }) => Promise; + listPendingAgentBackgroundToolCompletions: (input: { + user: string | Types.ObjectId; + conversationId: string; + sourceId: string; + /** One task's completion, e.g. to discard it. */ + 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; expediteAgentTriggerDeliveries: ( input: ExpediteAgentTriggerDeliveriesInput, ) => Promise; @@ -2077,6 +2138,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 { @@ -2164,6 +2228,9 @@ export function createAgentTriggerDeliveryMethods( } return true; } + if (input.requireTransition === true) { + return false; + } return ( (await Delivery().exists({ deliveryKey: input.deliveryKey, @@ -2263,6 +2330,133 @@ 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; + taskId?: string; + limit?: number; + }): 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 + ) { + 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, + ...(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, + /** Dead letters are read too, so a caller can tell "delivered" from "failed". */ + status: { $in: [...UNDELIVERED_STATUSES, ...DEAD_STATUSES] }, + }) + .select( + 'deliveryKey createdAt status capabilityStatus envelope.event.payload ' + + 'backgroundToolResult.status backgroundToolResult.settledAt backgroundToolResult.resultClaim', + ) + .sort({ createdAt: 1, _id: 1 }) + .limit(limit + 1) + .lean< + Array< + Pick< + IAgentTriggerDelivery, + 'deliveryKey' | 'createdAt' | 'backgroundToolResult' | 'status' | 'capabilityStatus' + > & { + envelope?: { event?: { payload?: Record } }; + } + > + >(); + const dead: PendingAgentBackgroundToolCompletion[] = []; + const completions = rows.slice(0, limit).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, + }, + ].filter((completion) => { + if (!isDeadDelivery(row)) { + return true; + } + dead.push(completion); + return false; + }); + }); + return { completions, dead, 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 }; + } + /** Pulls deferred deliveries back to `now` when the condition they were * waiting on has changed, so a waiting delivery can back off without delaying * the moment it becomes deliverable. Unclaimed rows move now; held rows retain @@ -4094,6 +4288,8 @@ export function createAgentTriggerDeliveryMethods( retireAgentTriggerDelivery, renewAgentTriggerDeliveryProducerLease, getAgentTriggerDeliveryProducerLease, + listPendingAgentBackgroundToolCompletions, + listUndeliveredAgentTriggerTaskIds, expediteAgentTriggerDeliveries, persistAgentBackgroundToolResult, getAgentBackgroundToolResult,