From 6d09a6ccee737bbb5721c80803f362dc176a4cb8 Mon Sep 17 00:00:00 2001 From: Danny Avila Date: Fri, 21 Aug 2026 03:35:46 -0400 Subject: [PATCH] =?UTF-8?q?=F0=9F=A7=BE=20feat:=20Sibling=20Task=20Manifes?= =?UTF-8?q?t=20for=20Resumed=20Parent=20Runs=20(#15063)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * feat: add bounded subagent orchestration snapshots * fix: harden orchestration snapshot selection * fix: close snapshot settlement race * fix: preserve retry lease uncertainty * fix: classify bounded sibling leases * fix: classify captured terminal leases * fix: enforce snapshot byte budget * fix: retain terminal lease evidence --- .../agents/subagentCompletionWakeup.spec.ts | 1006 ++++++++++++++++- .../src/agents/subagentCompletionWakeup.ts | 418 ++++++- packages/data-schemas/src/schema/message.ts | 9 + .../data-schemas/src/schema/subagent.spec.ts | 21 + 4 files changed, 1452 insertions(+), 2 deletions(-) create mode 100644 packages/data-schemas/src/schema/subagent.spec.ts diff --git a/packages/api/src/agents/subagentCompletionWakeup.spec.ts b/packages/api/src/agents/subagentCompletionWakeup.spec.ts index 9bd0f2df18..d2fd1fc95e 100644 --- a/packages/api/src/agents/subagentCompletionWakeup.spec.ts +++ b/packages/api/src/agents/subagentCompletionWakeup.spec.ts @@ -14,6 +14,15 @@ import { const NOW = 1_775_000_000_000; +interface TestMessageFilter { + conversationId?: string; + messageId?: { $in: string[] }; + user?: string; + 'subagentTask.parentRunId'?: string; + 'subagentTask.attemptKey'?: string; + 'subagentTask.status'?: string | { $in: string[] }; +} + function enqueueMock(): jest.MockedFunction { return jest.fn, Parameters>(async () => ({ id: 'delivery-1', @@ -207,12 +216,48 @@ function resolverMethods() { terminal, ], ), + listActiveSubagentThreadLeases: jest.fn( + async (): Promise< + Array<{ conversationId: string; parentConversationId: string; taskId: string }> + > => [], + ), claimSubagentTaskResult: jest.fn(async () => ({ status: 'acquired', message: terminal })), releaseSubagentTaskResultClaim: jest.fn(async () => true), }; return { methods, terminal }; } +function orchestrationSnapshot( + prepared: Awaited>>, +) { + if (prepared?.status !== 'ready') { + throw new Error('Expected a ready continuation.'); + } + const marker = 'Host-authored bounded orchestration snapshot:\n'; + const start = prepared.input.indexOf(marker); + if (start < 0) { + throw new Error('Expected an orchestration snapshot.'); + } + return { + rendered: prepared.input.slice(start + marker.length), + value: JSON.parse(prepared.input.slice(start + marker.length)) as { + parent_message_id?: string; + parent_message_id_truncated?: boolean; + completeness: 'complete' | 'bounded' | 'uncertain'; + known_children: Array<{ + background_task_id: string; + subagent_thread_id: string; + subagent_type: string; + status: string; + result_state: string; + current_completion: boolean; + }>; + additional_children_may_exist: boolean; + note: string; + }, + }; +} + describe('createSubagentCompletionWakeupResolver', () => { it('defers without claiming while the parent generation is active', async () => { const { methods } = resolverMethods(); @@ -262,6 +307,960 @@ describe('createSubagentCompletionWakeupResolver', () => { expect(prepared?.status === 'ready' && prepared.input.length).toBeLessThan(110_000); }); + it('keeps delayed sibling completions in deterministic host-authored order', async () => { + const { methods, terminal } = resolverMethods(); + terminal.subagentTask = { + ...terminal.subagentTask!, + parentRunId: 'response-1', + resultClaim: { kind: 'wakeup', claimId: 'trigger_claim_1', claimedAt: new Date(NOW) }, + }; + const delayedSibling = { + ...terminal, + messageId: 'task-2:assistant', + conversationId: 'thread-2', + sender: 'analyst', + text: 'private sibling transcript text', + subagentTranscript: { + taskId: 'task-2', + mode: 'replace' as const, + messagesJson: '[{"role":"assistant","content":"hidden reasoning"}]', + }, + createdAt: new Date(NOW - 500), + updatedAt: new Date(NOW - 400), + subagentTask: { + attemptKey: 'attempt-2', + parentRunId: 'response-1', + status: 'completed' as const, + resultClaim: { + kind: 'manual' as const, + claimId: 'older-poll', + claimedAt: new Date(NOW - 300), + }, + }, + }; + methods.getConvo.mockImplementation(async (_userId: string, conversationId: string) => { + if (conversationId === 'conversation-1') { + return { conversationId, tenantId: 'tenant-1' }; + } + return { + conversationId, + tenantId: 'tenant-1', + subagentThread: { + parentConversationId: 'conversation-1', + parentMessageId: 'response-1', + parentAgentId: 'agent_parent_1', + subagentType: conversationId === 'thread-2' ? 'analyst' : 'researcher', + }, + }; + }); + methods.getMessages.mockImplementation(async (filter: TestMessageFilter) => { + if (filter.conversationId === 'conversation-1') { + return [ + { + messageId: 'response-1', + parentMessageId: 'user-1', + isCreatedByUser: false, + createdAt: new Date(NOW - 30), + }, + ]; + } + if (filter['subagentTask.parentRunId'] === 'response-1') { + return [terminal, delayedSibling]; + } + return [ + { + messageId: 'task-1:user', + conversationId: 'thread-1', + isCreatedByUser: true, + }, + terminal, + ]; + }); + methods.claimSubagentTaskResult.mockResolvedValueOnce({ + status: 'acquired', + message: terminal, + }); + const resolve = createSubagentCompletionWakeupResolver({ + methods: methods as never, + getGenerationJob: async () => null, + }); + + const snapshot = orchestrationSnapshot( + await resolve(wakeupEnvelope(), { idempotencyKey: 'trigger_claim_1' } as never), + ); + + expect(snapshot.value.known_children).toEqual([ + expect.objectContaining({ + background_task_id: 'task-1', + status: 'completed', + result_state: 'claimed', + current_completion: true, + }), + expect.objectContaining({ + background_task_id: 'task-2', + subagent_type: 'analyst', + status: 'completed', + result_state: 'claimed', + current_completion: false, + }), + ]); + expect(snapshot.rendered).not.toContain('private sibling transcript text'); + expect(snapshot.rendered).not.toContain('hidden reasoning'); + }); + + it('bounds sibling count and rendered snapshot bytes', async () => { + const { methods, terminal } = resolverMethods(); + const long = '界'.repeat(240); + const siblingMessages = Array.from({ length: 40 }, (_, index) => ({ + ...terminal, + messageId: `task-${index}-${long}:assistant`, + conversationId: `thread-${index}-${long}`, + sender: `agent-${index}-${long}`, + createdAt: new Date(NOW - index), + updatedAt: new Date(NOW - index), + subagentTask: { + attemptKey: `attempt-${index}`, + parentRunId: 'response-1', + status: 'completed' as const, + }, + })); + methods.getMessages.mockImplementation(async (filter: TestMessageFilter) => { + if (filter.conversationId === 'conversation-1') { + return [ + { + messageId: 'response-1', + parentMessageId: 'user-1', + isCreatedByUser: false, + createdAt: new Date(NOW - 30), + }, + ]; + } + if (filter['subagentTask.parentRunId'] === 'response-1') { + return siblingMessages.slice(0, 12); + } + return [ + { + messageId: 'task-1:user', + conversationId: 'thread-1', + isCreatedByUser: true, + }, + terminal, + ]; + }); + methods.getConvo.mockImplementation(async (_userId: string, conversationId: string) => { + if (conversationId === 'conversation-1') { + return { conversationId, tenantId: 'tenant-1' }; + } + const match = /^thread-(\d+)-/.exec(conversationId); + return { + conversationId, + tenantId: 'tenant-1', + subagentThread: { + parentConversationId: 'conversation-1', + parentMessageId: 'response-1', + parentAgentId: 'agent_parent_1', + subagentType: match == null ? 'researcher' : `agent-${match[1]}-${long}`, + }, + }; + }); + const resolve = createSubagentCompletionWakeupResolver({ + methods: methods as never, + getGenerationJob: async () => null, + }); + + const snapshot = orchestrationSnapshot( + await resolve(wakeupEnvelope(), { idempotencyKey: 'trigger_claim_1' } as never), + ); + + expect(snapshot.value.known_children.length).toBeLessThanOrEqual(16); + expect(snapshot.value.known_children.length).toBeLessThan(13); + expect(Buffer.byteLength(snapshot.rendered, 'utf8')).toBeLessThanOrEqual(8 * 1_024); + expect(snapshot.value.completeness).toBe('bounded'); + expect(snapshot.value.additional_children_may_exist).toBe(true); + }); + + it('bounds an oversized parent identity before rendering the UTF-8 snapshot', async () => { + const { methods, terminal } = resolverMethods(); + const oversizedParentMessageId = '界'.repeat(10_000); + terminal.subagentTask = { + ...terminal.subagentTask!, + parentRunId: oversizedParentMessageId, + }; + methods.getMessages.mockImplementation(async (filter: TestMessageFilter) => { + if (filter.conversationId === 'conversation-1') { + return [ + { + messageId: oversizedParentMessageId, + parentMessageId: 'user-1', + isCreatedByUser: false, + createdAt: new Date(NOW - 30), + }, + ]; + } + if (filter.conversationId === 'thread-1') { + return [ + { + messageId: 'task-1:user', + conversationId: 'thread-1', + isCreatedByUser: true, + }, + terminal, + ]; + } + if (filter['subagentTask.parentRunId'] === oversizedParentMessageId) { + return [terminal]; + } + return []; + }); + methods.claimSubagentTaskResult.mockResolvedValueOnce({ + status: 'acquired', + message: terminal, + }); + const envelope = wakeupEnvelope(); + envelope.target.parentMessageId = oversizedParentMessageId; + const resolve = createSubagentCompletionWakeupResolver({ + methods: methods as never, + getGenerationJob: async () => null, + }); + + const snapshot = orchestrationSnapshot( + await resolve(envelope, { idempotencyKey: 'trigger_claim_1' } as never), + ); + + expect(Buffer.byteLength(snapshot.rendered, 'utf8')).toBeLessThanOrEqual(8 * 1_024); + expect(snapshot.value.parent_message_id).toBe('界'.repeat(256)); + expect(snapshot.value.parent_message_id_truncated).toBe(true); + expect(snapshot.value.known_children).toEqual([ + expect.objectContaining({ background_task_id: 'task-1', current_completion: true }), + ]); + }); + + it('keeps an older actively leased child ahead of newer settled siblings', async () => { + const { methods, terminal } = resolverMethods(); + const settled = Array.from({ length: 20 }, (_, index) => ({ + ...terminal, + messageId: `settled-${index}:assistant`, + conversationId: `thread-settled-${index}`, + sender: `settled-agent-${index}`, + createdAt: new Date(NOW - index), + updatedAt: new Date(NOW - index), + subagentTask: { + attemptKey: `settled-attempt-${index}`, + parentRunId: 'response-1', + status: 'completed' as const, + }, + })); + const running = { + ...terminal, + messageId: 'running-task:user', + conversationId: 'thread-running', + sender: 'User', + createdAt: new Date(NOW - 10_000), + updatedAt: new Date(NOW - 10_000), + subagentTask: { + attemptKey: 'running-attempt', + parentRunId: 'response-1', + status: 'running' as const, + }, + }; + methods.listActiveSubagentThreadLeases.mockResolvedValueOnce([ + { + conversationId: 'thread-running', + parentConversationId: 'conversation-1', + taskId: 'running-task', + }, + ]); + methods.getMessages.mockImplementation(async (filter: TestMessageFilter) => { + if (filter.conversationId === 'conversation-1') { + return [ + { + messageId: 'response-1', + parentMessageId: 'user-1', + isCreatedByUser: false, + createdAt: new Date(NOW - 30), + }, + ]; + } + if (filter.conversationId === 'thread-1') { + return [ + { + messageId: 'task-1:user', + conversationId: 'thread-1', + isCreatedByUser: true, + }, + terminal, + ]; + } + if (filter.messageId != null) { + return [running]; + } + if (filter['subagentTask.parentRunId'] === 'response-1') { + return [terminal, ...settled]; + } + return [ + { + messageId: 'task-1:user', + conversationId: 'thread-1', + isCreatedByUser: true, + }, + terminal, + ]; + }); + methods.getConvo.mockImplementation(async (_userId: string, conversationId: string) => { + if (conversationId === 'conversation-1') { + return { conversationId, tenantId: 'tenant-1' }; + } + const settledIndex = /^thread-settled-(\d+)$/.exec(conversationId)?.[1]; + let subagentType = 'researcher'; + if (conversationId === 'thread-running') { + subagentType = 'running-agent'; + } else if (settledIndex != null) { + subagentType = `settled-agent-${settledIndex}`; + } + return { + conversationId, + tenantId: 'tenant-1', + subagentThread: { + parentConversationId: 'conversation-1', + parentMessageId: 'response-1', + parentAgentId: 'agent_parent_1', + subagentType, + }, + }; + }); + const resolve = createSubagentCompletionWakeupResolver({ + methods: methods as never, + getGenerationJob: async () => null, + }); + + const snapshot = orchestrationSnapshot( + await resolve(wakeupEnvelope(), { idempotencyKey: 'trigger_claim_1' } as never), + ); + + expect(snapshot.value.known_children.slice(0, 2)).toEqual([ + expect.objectContaining({ background_task_id: 'task-1', current_completion: true }), + expect.objectContaining({ + background_task_id: 'running-task', + status: 'running', + result_state: 'pending', + }), + ]); + }); + + it('preserves uncertainty for a retry lease without a same-task seed or terminal', async () => { + const { methods, terminal } = resolverMethods(); + methods.listActiveSubagentThreadLeases.mockResolvedValueOnce([ + { + conversationId: 'thread-2', + parentConversationId: 'conversation-1', + taskId: 'retry-task', + }, + ]); + methods.getMessages.mockImplementation(async (filter: TestMessageFilter) => { + if (filter.conversationId === 'conversation-1') { + return [ + { + messageId: 'response-1', + parentMessageId: 'user-1', + isCreatedByUser: false, + createdAt: new Date(NOW - 30), + }, + ]; + } + if (filter.conversationId === 'thread-1') { + return [ + { + messageId: 'task-1:user', + conversationId: 'thread-1', + isCreatedByUser: true, + }, + terminal, + ]; + } + if (filter.messageId != null || filter['subagentTask.parentRunId'] === 'response-1') { + return []; + } + return []; + }); + const resolve = createSubagentCompletionWakeupResolver({ + methods: methods as never, + getGenerationJob: async () => null, + }); + + const snapshot = orchestrationSnapshot( + await resolve(wakeupEnvelope(), { idempotencyKey: 'trigger_claim_1' } as never), + ); + + expect(snapshot.value.known_children).toEqual([ + expect.objectContaining({ background_task_id: 'task-1', current_completion: true }), + ]); + expect(snapshot.value.completeness).toBe('uncertain'); + expect(snapshot.value.additional_children_may_exist).toBe(true); + expect(snapshot.value.note).toContain('Do not infer that no other children ran'); + }); + + it('checks an unmatched retry lease even when ordinary active seeds exceed the task cap', async () => { + const { methods, terminal } = resolverMethods(); + const activeSeeds = Array.from({ length: 17 }, (_, index) => ({ + messageId: `active-${index}:user`, + conversationId: `thread-active-${index}`, + sender: 'User', + isCreatedByUser: true, + createdAt: new Date(NOW - index), + updatedAt: new Date(NOW - index), + subagentTask: { + attemptKey: `active-attempt-${index}`, + parentRunId: 'response-1', + status: 'running' as const, + }, + })); + methods.listActiveSubagentThreadLeases.mockResolvedValueOnce([ + ...activeSeeds.map((message, index) => ({ + conversationId: message.conversationId, + parentConversationId: 'conversation-1', + taskId: `active-${index}`, + })), + { + conversationId: 'thread-retry', + parentConversationId: 'conversation-1', + taskId: 'retry-task', + }, + ]); + methods.getMessages.mockImplementation(async (filter: TestMessageFilter) => { + if (filter.conversationId === 'conversation-1') { + return [ + { + messageId: 'response-1', + parentMessageId: 'user-1', + isCreatedByUser: false, + createdAt: new Date(NOW - 30), + }, + ]; + } + if (filter.conversationId === 'thread-1') { + return [ + { + messageId: 'task-1:user', + conversationId: 'thread-1', + isCreatedByUser: true, + }, + terminal, + ]; + } + if (filter.messageId != null) { + return activeSeeds; + } + if (filter['subagentTask.parentRunId'] === 'response-1') { + return [terminal]; + } + return []; + }); + methods.getConvo.mockImplementation(async (_userId: string, conversationId: string) => ({ + conversationId, + tenantId: 'tenant-1', + ...(conversationId === 'conversation-1' + ? {} + : { + subagentThread: { + parentConversationId: 'conversation-1', + parentMessageId: 'response-1', + parentAgentId: 'agent_parent_1', + subagentType: + conversationId === 'thread-1' + ? 'researcher' + : `active-agent-${conversationId.slice('thread-active-'.length)}`, + }, + }), + })); + const resolve = createSubagentCompletionWakeupResolver({ + methods: methods as never, + getGenerationJob: async () => null, + }); + + const snapshot = orchestrationSnapshot( + await resolve(wakeupEnvelope(), { idempotencyKey: 'trigger_claim_1' } as never), + ); + + expect(snapshot.value.known_children).toHaveLength(16); + expect(snapshot.value.completeness).toBe('uncertain'); + expect(snapshot.value.additional_children_may_exist).toBe(true); + }); + + it('excludes a captured lease whose durable seed belongs to another parent run', async () => { + const { methods, terminal } = resolverMethods(); + methods.listActiveSubagentThreadLeases.mockResolvedValueOnce([ + { + conversationId: 'thread-other-run', + parentConversationId: 'conversation-1', + taskId: 'other-task', + }, + ]); + methods.getMessages.mockImplementation(async (filter: TestMessageFilter) => { + if (filter.conversationId === 'conversation-1') { + return [ + { + messageId: 'response-1', + parentMessageId: 'user-1', + isCreatedByUser: false, + createdAt: new Date(NOW - 30), + }, + ]; + } + if (filter.conversationId === 'thread-1') { + return [ + { + messageId: 'task-1:user', + conversationId: 'thread-1', + isCreatedByUser: true, + }, + terminal, + ]; + } + if (filter.messageId != null) { + return [ + { + messageId: 'other-task:user', + conversationId: 'thread-other-run', + sender: 'User', + isCreatedByUser: true, + subagentTask: { + attemptKey: 'other-attempt', + parentRunId: 'response-other', + status: 'running' as const, + }, + }, + ]; + } + if (filter['subagentTask.parentRunId'] === 'response-1') { + return [terminal]; + } + return []; + }); + const resolve = createSubagentCompletionWakeupResolver({ + methods: methods as never, + getGenerationJob: async () => null, + }); + + const snapshot = orchestrationSnapshot( + await resolve(wakeupEnvelope(), { idempotencyKey: 'trigger_claim_1' } as never), + ); + + expect(snapshot.value.known_children).toEqual([ + expect.objectContaining({ background_task_id: 'task-1', current_completion: true }), + ]); + expect(snapshot.value.completeness).toBe('complete'); + expect(snapshot.value.additional_children_may_exist).toBe(false); + }); + + it('excludes a captured lease whose same-task terminal belongs to another parent run', async () => { + const { methods, terminal } = resolverMethods(); + methods.listActiveSubagentThreadLeases.mockResolvedValueOnce([ + { + conversationId: 'thread-other-run', + parentConversationId: 'conversation-1', + taskId: 'other-retry', + }, + ]); + methods.getMessages.mockImplementation(async (filter: TestMessageFilter) => { + if (filter.conversationId === 'conversation-1') { + return [ + { + messageId: 'response-1', + parentMessageId: 'user-1', + isCreatedByUser: false, + createdAt: new Date(NOW - 30), + }, + ]; + } + if (filter.conversationId === 'thread-1') { + return [ + { + messageId: 'task-1:user', + conversationId: 'thread-1', + isCreatedByUser: true, + }, + terminal, + ]; + } + if (filter.messageId != null) { + return [ + { + messageId: 'other-retry:assistant', + conversationId: 'thread-other-run', + sender: 'other-agent', + isCreatedByUser: false, + subagentTask: { + attemptKey: 'other-attempt', + parentRunId: 'response-other', + status: 'error' as const, + }, + }, + ]; + } + if (filter['subagentTask.parentRunId'] === 'response-1') { + return [terminal]; + } + return []; + }); + const resolve = createSubagentCompletionWakeupResolver({ + methods: methods as never, + getGenerationJob: async () => null, + }); + + const snapshot = orchestrationSnapshot( + await resolve(wakeupEnvelope(), { idempotencyKey: 'trigger_claim_1' } as never), + ); + + expect(snapshot.value.known_children).toEqual([ + expect.objectContaining({ background_task_id: 'task-1', current_completion: true }), + ]); + expect(snapshot.value.completeness).toBe('complete'); + expect(snapshot.value.additional_children_may_exist).toBe(false); + }); + + it('retains a same-branch retry terminal seen only by lease evidence', async () => { + const { methods, terminal } = resolverMethods(); + const retryTerminal = { + ...terminal, + messageId: 'retry-task:assistant', + conversationId: 'thread-retry', + sender: 'reviewer', + createdAt: new Date(NOW - 20), + updatedAt: new Date(NOW - 10), + subagentTask: { + attemptKey: 'retry-attempt', + parentRunId: 'response-1', + status: 'error' as const, + }, + }; + methods.listActiveSubagentThreadLeases.mockResolvedValueOnce([ + { + conversationId: 'thread-retry', + parentConversationId: 'conversation-1', + taskId: 'retry-task', + }, + ]); + methods.getMessages.mockImplementation(async (filter: TestMessageFilter) => { + if (filter.conversationId === 'conversation-1') { + return [ + { + messageId: 'response-1', + parentMessageId: 'user-1', + isCreatedByUser: false, + createdAt: new Date(NOW - 30), + }, + ]; + } + if (filter.conversationId === 'thread-1') { + return [ + { + messageId: 'task-1:user', + conversationId: 'thread-1', + isCreatedByUser: true, + }, + terminal, + ]; + } + if (filter.messageId != null) { + return [retryTerminal]; + } + if (filter['subagentTask.parentRunId'] === 'response-1') { + return [terminal]; + } + return []; + }); + methods.getConvo.mockImplementation(async (_userId: string, conversationId: string) => ({ + conversationId, + tenantId: 'tenant-1', + ...(conversationId === 'conversation-1' + ? {} + : { + subagentThread: { + parentConversationId: 'conversation-1', + parentMessageId: 'response-1', + parentAgentId: 'agent_parent_1', + subagentType: conversationId === 'thread-retry' ? 'reviewer' : 'researcher', + }, + }), + })); + const resolve = createSubagentCompletionWakeupResolver({ + methods: methods as never, + getGenerationJob: async () => null, + }); + + const snapshot = orchestrationSnapshot( + await resolve(wakeupEnvelope(), { idempotencyKey: 'trigger_claim_1' } as never), + ); + + expect(snapshot.value.known_children).toEqual([ + expect.objectContaining({ background_task_id: 'task-1', current_completion: true }), + expect.objectContaining({ background_task_id: 'retry-task', status: 'error' }), + ]); + expect(snapshot.value.completeness).toBe('complete'); + expect(snapshot.value.additional_children_may_exist).toBe(false); + }); + + it('lets a lease-evidence terminal supersede the same attempt running seed', async () => { + const { methods, terminal } = resolverMethods(); + const runningSeed = { + messageId: 'settling-task:user', + conversationId: 'thread-settling', + sender: 'User', + isCreatedByUser: true, + createdAt: new Date(NOW - 30), + updatedAt: new Date(NOW - 30), + subagentTask: { + attemptKey: 'settling-attempt', + parentRunId: 'response-1', + status: 'running' as const, + }, + }; + const settledTerminal = { + ...terminal, + messageId: 'settling-task:assistant', + conversationId: 'thread-settling', + sender: 'reviewer', + createdAt: new Date(NOW - 20), + updatedAt: new Date(NOW - 10), + subagentTask: { + attemptKey: 'settling-attempt', + parentRunId: 'response-1', + status: 'completed' as const, + }, + }; + methods.listActiveSubagentThreadLeases.mockResolvedValueOnce([ + { + conversationId: 'thread-settling', + parentConversationId: 'conversation-1', + taskId: 'settling-task', + }, + ]); + methods.getMessages.mockImplementation(async (filter: TestMessageFilter) => { + if (filter.conversationId === 'conversation-1') { + return [ + { + messageId: 'response-1', + parentMessageId: 'user-1', + isCreatedByUser: false, + createdAt: new Date(NOW - 30), + }, + ]; + } + if (filter.conversationId === 'thread-1') { + return [ + { + messageId: 'task-1:user', + conversationId: 'thread-1', + isCreatedByUser: true, + }, + terminal, + ]; + } + if (filter.messageId != null) { + return [runningSeed, settledTerminal]; + } + if (filter['subagentTask.parentRunId'] === 'response-1') { + return [terminal]; + } + return []; + }); + methods.getConvo.mockImplementation(async (_userId: string, conversationId: string) => ({ + conversationId, + tenantId: 'tenant-1', + ...(conversationId === 'conversation-1' + ? {} + : { + subagentThread: { + parentConversationId: 'conversation-1', + parentMessageId: 'response-1', + parentAgentId: 'agent_parent_1', + subagentType: conversationId === 'thread-settling' ? 'reviewer' : 'researcher', + }, + }), + })); + const resolve = createSubagentCompletionWakeupResolver({ + methods: methods as never, + getGenerationJob: async () => null, + }); + + const snapshot = orchestrationSnapshot( + await resolve(wakeupEnvelope(), { idempotencyKey: 'trigger_claim_1' } as never), + ); + + expect(snapshot.value.known_children).toEqual([ + expect.objectContaining({ background_task_id: 'task-1', current_completion: true }), + expect.objectContaining({ + background_task_id: 'settling-task', + status: 'completed', + result_state: 'available', + }), + ]); + expect(snapshot.value.completeness).toBe('complete'); + }); + + it('omits sibling task records outside the authorized parent lineage', async () => { + const { methods, terminal } = resolverMethods(); + const sibling = (taskId: string, threadId: string, sender: string) => ({ + ...terminal, + messageId: `${taskId}:assistant`, + conversationId: threadId, + sender, + text: `secret-${taskId}`, + subagentTask: { + attemptKey: `attempt-${taskId}`, + parentRunId: 'response-1', + status: 'completed' as const, + }, + }); + methods.getMessages.mockImplementation(async (filter: TestMessageFilter) => { + if (filter.conversationId === 'conversation-1') { + return [ + { + messageId: 'response-1', + parentMessageId: 'user-1', + isCreatedByUser: false, + createdAt: new Date(NOW - 30), + }, + ]; + } + if (filter['subagentTask.parentRunId'] === 'response-1') { + expect(filter.user).toBe('user-1'); + return [ + terminal, + sibling('valid', 'thread-valid', 'valid-agent'), + sibling('wrong-tenant', 'thread-wrong-tenant', 'tenant-agent'), + sibling('wrong-parent', 'thread-wrong-parent', 'parent-agent'), + sibling('wrong-agent', 'thread-wrong-agent', 'agent-agent'), + ]; + } + return [ + { + messageId: 'task-1:user', + conversationId: 'thread-1', + isCreatedByUser: true, + }, + terminal, + ]; + }); + methods.getConvo.mockImplementation(async (_userId: string, conversationId: string) => { + if (conversationId === 'conversation-1') { + return { conversationId, tenantId: 'tenant-1' }; + } + const variants: Record< + string, + { + tenantId: string; + parentConversationId: string; + parentAgentId: string; + subagentType: string; + } + > = { + 'thread-1': { + tenantId: 'tenant-1', + parentConversationId: 'conversation-1', + parentAgentId: 'agent_parent_1', + subagentType: 'researcher', + }, + 'thread-valid': { + tenantId: 'tenant-1', + parentConversationId: 'conversation-1', + parentAgentId: 'agent_parent_1', + subagentType: 'valid-agent', + }, + 'thread-wrong-tenant': { + tenantId: 'tenant-2', + parentConversationId: 'conversation-1', + parentAgentId: 'agent_parent_1', + subagentType: 'tenant-agent', + }, + 'thread-wrong-parent': { + tenantId: 'tenant-1', + parentConversationId: 'conversation-2', + parentAgentId: 'agent_parent_1', + subagentType: 'parent-agent', + }, + 'thread-wrong-agent': { + tenantId: 'tenant-1', + parentConversationId: 'conversation-1', + parentAgentId: 'agent_parent_2', + subagentType: 'agent-agent', + }, + }; + const variant = variants[conversationId]; + return { + conversationId, + tenantId: variant.tenantId, + subagentThread: { + parentConversationId: variant.parentConversationId, + parentMessageId: 'response-1', + parentAgentId: variant.parentAgentId, + subagentType: variant.subagentType, + }, + }; + }); + const resolve = createSubagentCompletionWakeupResolver({ + methods: methods as never, + getGenerationJob: async () => null, + }); + + const snapshot = orchestrationSnapshot( + await resolve(wakeupEnvelope(), { idempotencyKey: 'trigger_claim_1' } as never), + ); + + expect( + snapshot.value.known_children.map(({ background_task_id }) => background_task_id), + ).toEqual(['task-1', 'valid']); + expect(snapshot.value.completeness).toBe('uncertain'); + expect(snapshot.value.additional_children_may_exist).toBe(true); + expect(snapshot.value.note).toContain('Do not infer that no other children ran'); + expect(snapshot.rendered).not.toContain('wrong-tenant'); + expect(snapshot.rendered).not.toContain('wrong-parent'); + expect(snapshot.rendered).not.toContain('wrong-agent'); + expect(snapshot.rendered).not.toContain('secret-'); + }); + + it('states uncertainty when the bounded sibling read is unavailable', async () => { + const { methods, terminal } = resolverMethods(); + methods.getMessages.mockImplementation(async (filter: TestMessageFilter) => { + if (filter['subagentTask.parentRunId'] === 'response-1') { + throw new Error('temporary sibling read failure'); + } + if (filter.conversationId === 'conversation-1') { + return [ + { + messageId: 'response-1', + parentMessageId: 'user-1', + isCreatedByUser: false, + createdAt: new Date(NOW - 30), + }, + ]; + } + return [ + { + messageId: 'task-1:user', + conversationId: 'thread-1', + isCreatedByUser: true, + }, + terminal, + ]; + }); + const resolve = createSubagentCompletionWakeupResolver({ + methods: methods as never, + getGenerationJob: async () => null, + }); + + const snapshot = orchestrationSnapshot( + await resolve(wakeupEnvelope(), { idempotencyKey: 'trigger_claim_1' } as never), + ); + + expect(snapshot.value.known_children).toEqual([ + expect.objectContaining({ background_task_id: 'task-1', current_completion: true }), + ]); + expect(snapshot.value.completeness).toBe('uncertain'); + expect(snapshot.value.additional_children_may_exist).toBe(true); + expect(snapshot.value.note).toContain('Do not infer that no other children ran'); + }); + it('dead-letters a child whose process disappeared after the task timeout grace', async () => { const { methods } = resolverMethods(); methods.getMessages.mockImplementation(async (filter: { conversationId: string }) => @@ -318,7 +1317,7 @@ describe('createSubagentCompletionWakeupResolver', () => { messageId: 'task-2:assistant', parentMessageId: 'task-1:user', }; - methods.getMessages.mockImplementation(async (filter: Record) => { + methods.getMessages.mockImplementation(async (filter: TestMessageFilter) => { if (filter.conversationId === 'conversation-1') { return [ { @@ -362,6 +1361,11 @@ describe('createSubagentCompletionWakeupResolver', () => { status: 'ready', input: expect.stringContaining('"background_task_id":"task-2"'), }); + expect( + orchestrationSnapshot(prepared).value.known_children.map( + ({ background_task_id }) => background_task_id, + ), + ).toEqual(['task-2']); expect(methods.claimSubagentTaskResult).toHaveBeenCalledWith({ userId: 'user-1', conversationId: 'thread-1', diff --git a/packages/api/src/agents/subagentCompletionWakeup.ts b/packages/api/src/agents/subagentCompletionWakeup.ts index dd6329f37e..915351cc30 100644 --- a/packages/api/src/agents/subagentCompletionWakeup.ts +++ b/packages/api/src/agents/subagentCompletionWakeup.ts @@ -21,6 +21,13 @@ const EVENT_TYPE = 'subagent.completion'; const MESSAGE_SELECT = 'messageId parentMessageId isCreatedByUser createdAt'; const TASK_SELECT = 'messageId conversationId parentMessageId sender text error createdAt updatedAt +subagentTask'; +const ORCHESTRATION_TASK_SELECT = + 'messageId conversationId sender isCreatedByUser createdAt updatedAt +subagentTask'; +const MAX_ORCHESTRATION_TASKS = 16; +const MAX_ORCHESTRATION_CANDIDATES = MAX_ORCHESTRATION_TASKS * 2 + 1; +const MAX_ORCHESTRATION_ACTIVE_LEASES = 200; +const MAX_ORCHESTRATION_SNAPSHOT_BYTES = 8 * 1_024; +const MAX_ORCHESTRATION_SCALAR_CHARS = 256; export type EnqueueAgentTrigger = ( envelope: unknown, @@ -28,6 +35,7 @@ export type EnqueueAgentTrigger = ( ) => Promise; type WakeupMethods = Pick & + Pick & Pick< MessageMethods, 'claimSubagentTaskResult' | 'getMessages' | 'releaseSubagentTaskResultClaim' @@ -41,6 +49,34 @@ interface GenerationState { }; } +type SubagentTaskStatus = NonNullable['status']; + +interface OrchestrationTaskCandidate { + attemptKey: string; + taskId: string; + threadId: string; + status: SubagentTaskStatus; + updatedAt: number; + resultClaimed: boolean; + sender?: string; +} + +interface OrchestrationTaskSnapshot { + background_task_id: string; + subagent_thread_id: string; + subagent_type: string; + status: SubagentTaskStatus; + result_state: 'pending' | 'available' | 'claimed'; + current_completion: boolean; +} + +interface OrchestrationSnapshotResolution { + tasks: OrchestrationTaskSnapshot[]; + candidateLimitReached: boolean; + lineageUncertain: boolean; + readUncertain: boolean; +} + export interface SubagentCompletionWakeupResolverDeps { methods: WakeupMethods; getGenerationJob: (conversationId: string) => Promise; @@ -116,6 +152,370 @@ function timestamp(message: Pick): number { return Number.isFinite(parsed) ? parsed : 0; } +function updatedTimestamp(message: Pick): number { + const value = message.updatedAt; + if (value instanceof Date) { + return value.getTime(); + } + const parsed = value == null ? Number.NaN : new Date(value).getTime(); + return Number.isFinite(parsed) ? parsed : timestamp(message); +} + +function taskIdFromMessage(message: Pick): string | undefined { + let suffix: ':assistant' | ':user'; + if (message.messageId.endsWith(':assistant')) { + suffix = ':assistant'; + } else if (message.messageId.endsWith(':user')) { + suffix = ':user'; + } else { + return; + } + const taskId = message.messageId.slice(0, -suffix.length); + return taskId.length > 0 && taskId.length <= 256 ? taskId : undefined; +} + +function candidateFromMessage(message: IMessage): OrchestrationTaskCandidate | undefined { + const taskId = taskIdFromMessage(message); + const threadId = message.conversationId; + const status = message.subagentTask?.status; + const attemptKey = message.subagentTask?.attemptKey; + if ( + taskId == null || + typeof threadId !== 'string' || + threadId.length === 0 || + threadId.length > 256 || + typeof attemptKey !== 'string' || + attemptKey.length === 0 || + attemptKey.length > 256 || + status == null + ) { + return; + } + const terminal = status !== 'running'; + if ( + (terminal && !message.messageId.endsWith(':assistant')) || + (!terminal && !message.messageId.endsWith(':user')) + ) { + return; + } + return { + attemptKey, + taskId, + threadId, + status, + updatedAt: updatedTimestamp(message), + resultClaimed: terminal && message.subagentTask?.resultClaim != null, + ...(typeof message.sender === 'string' && message.sender.length > 0 + ? { sender: message.sender } + : {}), + }; +} + +function preferCandidate( + current: OrchestrationTaskCandidate | undefined, + candidate: OrchestrationTaskCandidate, +): OrchestrationTaskCandidate { + if (current == null || (current.status === 'running' && candidate.status !== 'running')) { + return candidate; + } + return candidate.updatedAt > current.updatedAt ? candidate : current; +} + +function resultState( + candidate: OrchestrationTaskCandidate, +): OrchestrationTaskSnapshot['result_state'] { + if (candidate.status === 'running') { + return 'pending'; + } + return candidate.resultClaimed ? 'claimed' : 'available'; +} + +async function resolveOrchestrationSnapshot( + methods: WakeupMethods, + input: { + userId: string; + tenantId?: string; + parentConversationId: string; + parentMessageId: string; + parentAgentId: string; + currentThread: NonNullable>>; + currentTaskId: string; + currentTerminal: IMessage; + }, +): Promise { + const currentCandidate = candidateFromMessage(input.currentTerminal); + if (currentCandidate == null) { + return { + tasks: [], + candidateLimitReached: false, + lineageUncertain: true, + readUncertain: true, + }; + } + + let readUncertain = false; + let activeLeases: Awaited> = []; + /** Snapshot leases before terminal rows: a child settling between these reads is + * then visible either through its earlier lease or through its later terminal. */ + try { + activeLeases = ( + await methods.listActiveSubagentThreadLeases({ + user: input.userId, + now: new Date(), + ...(input.tenantId == null ? {} : { tenantId: input.tenantId }), + }) + ).filter((lease) => lease.parentConversationId === input.parentConversationId); + } catch { + readUncertain = true; + } + const boundedActiveLeases = activeLeases.slice(0, MAX_ORCHESTRATION_ACTIVE_LEASES); + const leaseEvidenceRead = + boundedActiveLeases.length === 0 + ? Promise.resolve([]) + : methods.getMessages( + { + user: input.userId, + messageId: { + $in: boundedActiveLeases.flatMap(({ taskId }) => [ + `${taskId}:user`, + `${taskId}:assistant`, + ]), + }, + 'subagentTask.status': { $in: ['running', 'completed', 'error', 'cancelled'] }, + }, + ORCHESTRATION_TASK_SELECT, + { sort: false, limit: MAX_ORCHESTRATION_ACTIVE_LEASES * 2 }, + ); + const [terminalResult, leaseEvidenceResult] = await Promise.allSettled([ + methods.getMessages( + { + user: input.userId, + 'subagentTask.parentRunId': input.parentMessageId, + 'subagentTask.status': { $in: ['completed', 'error', 'cancelled'] }, + }, + ORCHESTRATION_TASK_SELECT, + { sort: { updatedAt: -1, _id: -1 }, limit: MAX_ORCHESTRATION_CANDIDATES }, + ), + leaseEvidenceRead, + ]); + readUncertain ||= + terminalResult.status === 'rejected' || leaseEvidenceResult.status === 'rejected'; + const terminalMessages = terminalResult.status === 'fulfilled' ? terminalResult.value : []; + const leaseEvidenceMessages = + leaseEvidenceResult.status === 'fulfilled' ? leaseEvidenceResult.value : []; + const validLeaseEvidence = leaseEvidenceMessages.flatMap((message) => { + const candidate = candidateFromMessage(message); + return candidate == null ? [] : [{ message, candidate }]; + }); + const activeMessages = validLeaseEvidence + .filter( + ({ message, candidate }) => + candidate.status === 'running' && + message.subagentTask?.parentRunId === input.parentMessageId, + ) + .map(({ message }) => message); + const leaseTerminalMessages = validLeaseEvidence + .filter( + ({ message, candidate }) => + candidate.status !== 'running' && + message.subagentTask?.parentRunId === input.parentMessageId, + ) + .map(({ message }) => message); + if (leaseEvidenceResult.status === 'fulfilled') { + const resolvedLeaseTaskIds = new Set( + validLeaseEvidence + .filter(({ message }) => typeof message.subagentTask?.parentRunId === 'string') + .map(({ candidate }) => candidate.taskId), + ); + /** A retry can acquire a replacement task lease before persisting its terminal + * assistant row, while retaining only the abandoned attempt's seed. A valid + * seed from another parent run excludes that lease from this branch, and a + * visible same-task terminal resolves the lease during post-settlement cleanup. + * Anything left unmatched is an identity gap, not permission to invent one. */ + readUncertain ||= boundedActiveLeases.some(({ taskId }) => !resolvedLeaseTaskIds.has(taskId)); + } + if ( + terminalMessages.length === 0 && + activeMessages.length === 0 && + leaseTerminalMessages.length === 0 && + readUncertain + ) { + const lineage = input.currentThread.subagentThread; + if (lineage == null) { + return { + tasks: [], + candidateLimitReached: false, + lineageUncertain: true, + readUncertain: true, + }; + } + return { + tasks: [ + { + background_task_id: input.currentTaskId, + subagent_thread_id: input.currentThread.conversationId, + subagent_type: lineage.subagentType, + status: currentCandidate.status, + result_state: resultState(currentCandidate), + current_completion: true, + }, + ], + candidateLimitReached: false, + lineageUncertain: false, + readUncertain: true, + }; + } + + const byAttemptKey = new Map(); + for (const message of [...activeMessages, ...terminalMessages, ...leaseTerminalMessages]) { + const candidate = candidateFromMessage(message); + if (candidate == null) { + continue; + } + byAttemptKey.set( + candidate.attemptKey, + preferCandidate(byAttemptKey.get(candidate.attemptKey), candidate), + ); + } + byAttemptKey.set(currentCandidate.attemptKey, currentCandidate); + + const candidates = [...byAttemptKey.values()].sort((left, right) => { + if (left.taskId === input.currentTaskId) { + return -1; + } + if (right.taskId === input.currentTaskId) { + return 1; + } + if (left.status === 'running' && right.status !== 'running') { + return -1; + } + if (right.status === 'running' && left.status !== 'running') { + return 1; + } + const time = right.updatedAt - left.updatedAt; + return time === 0 ? left.taskId.localeCompare(right.taskId) : time; + }); + const selected = candidates.slice(0, MAX_ORCHESTRATION_TASKS); + const siblingThreadIds = [ + ...new Set( + selected + .filter((candidate) => candidate.threadId !== input.currentThread.conversationId) + .map((candidate) => candidate.threadId), + ), + ]; + const siblingThreads = await Promise.all( + siblingThreadIds.map(async (threadId) => { + try { + return await methods.getConvo(input.userId, threadId); + } catch { + return null; + } + }), + ); + const threads = new Map([ + [input.currentThread.conversationId, input.currentThread], + ...siblingThreads + .filter((thread): thread is NonNullable => thread != null) + .map((thread) => [thread.conversationId, thread] as const), + ]); + let lineageUncertain = siblingThreads.some((thread) => thread == null); + const tasks: OrchestrationTaskSnapshot[] = []; + for (const candidate of selected) { + const conversation = threads.get(candidate.threadId); + const lineage = conversation?.subagentThread; + if ( + conversation == null || + lineage == null || + !sameTenant(conversation.tenantId, input.tenantId) || + lineage.parentConversationId !== input.parentConversationId || + lineage.parentAgentId !== input.parentAgentId || + (candidate.status !== 'running' && candidate.sender !== lineage.subagentType) + ) { + lineageUncertain = true; + continue; + } + tasks.push({ + background_task_id: candidate.taskId, + subagent_thread_id: candidate.threadId, + subagent_type: lineage.subagentType, + status: candidate.status, + result_state: resultState(candidate), + current_completion: candidate.taskId === input.currentTaskId, + }); + } + return { + tasks, + candidateLimitReached: + terminalMessages.length === MAX_ORCHESTRATION_CANDIDATES || + activeLeases.length > MAX_ORCHESTRATION_ACTIVE_LEASES || + activeMessages.length > MAX_ORCHESTRATION_TASKS || + candidates.length > MAX_ORCHESTRATION_TASKS, + lineageUncertain, + readUncertain, + }; +} + +function renderOrchestrationSnapshot( + parentMessageId: string, + resolution: OrchestrationSnapshotResolution, +): string { + const knownChildren = resolution.tasks.slice(0, MAX_ORCHESTRATION_TASKS); + let omitted = resolution.tasks.length - knownChildren.length; + const boundedParentMessageId = parentMessageId.slice(0, MAX_ORCHESTRATION_SCALAR_CHARS); + const parentMessageIdTruncated = boundedParentMessageId !== parentMessageId; + const completeness = (): 'complete' | 'bounded' | 'uncertain' => { + if (resolution.readUncertain || resolution.lineageUncertain) { + return 'uncertain'; + } + return resolution.candidateLimitReached || omitted > 0 ? 'bounded' : 'complete'; + }; + const note = (): string => { + if (completeness() === 'uncertain') { + return 'Some sibling state could not be read or verified. Do not infer that no other children ran.'; + } + if (completeness() === 'bounded') { + return 'Additional durable child tasks may exist outside this bounded snapshot.'; + } + return 'This lists the known durable child tasks for this exact parent run.'; + }; + const serialize = () => + JSON.stringify({ + scope: 'current_parent_branch', + parent_message_id: boundedParentMessageId, + ...(parentMessageIdTruncated ? { parent_message_id_truncated: true } : {}), + completeness: completeness(), + known_children: knownChildren, + omitted_known_children: omitted, + additional_children_may_exist: + resolution.candidateLimitReached || + resolution.readUncertain || + resolution.lineageUncertain || + omitted > 0, + note: note(), + }); + let rendered = serialize(); + while ( + Buffer.byteLength(rendered, 'utf8') > MAX_ORCHESTRATION_SNAPSHOT_BYTES && + knownChildren.length > 0 + ) { + knownChildren.pop(); + omitted += 1; + rendered = serialize(); + } + if (Buffer.byteLength(rendered, 'utf8') > MAX_ORCHESTRATION_SNAPSHOT_BYTES) { + return JSON.stringify({ + scope: 'current_parent_branch', + completeness: 'uncertain', + known_children: [], + omitted_known_children: resolution.tasks.length, + additional_children_may_exist: true, + current_completion_in_preceding_result: true, + note: 'Snapshot metadata exceeded its byte budget. Do not infer that no other children ran.', + }); + } + return rendered; +} + /** Selects the newest persisted assistant on the branch below the original * parent. Re-resolving for every ordered delivery serializes sibling child * completions onto the branch produced by the preceding wakeup. */ @@ -155,6 +555,7 @@ function renderWakeupInput( registration: Pick, resultTaskId: string, terminal: IMessage, + orchestrationSnapshot: string, ): string { const status = terminal.subagentTask?.status ?? 'error'; return [ @@ -166,6 +567,8 @@ function renderWakeupInput( status, result: boundedSubagentTaskResult(terminal.text ?? ''), }), + 'Host-authored bounded orchestration snapshot:', + orchestrationSnapshot, ].join('\n'); } @@ -359,10 +762,23 @@ export function createSubagentCompletionWakeupResolver({ } return { status: 'settled' }; } + const orchestrationSnapshot = renderOrchestrationSnapshot( + envelope.target.parentMessageId, + await resolveOrchestrationSnapshot(methods, { + userId, + tenantId, + parentConversationId: envelope.target.conversationId, + parentMessageId: envelope.target.parentMessageId, + parentAgentId: envelope.target.agentId, + currentThread: child, + currentTaskId: resultTaskId, + currentTerminal: claim.message, + }), + ); return { status: 'ready', parentMessageId, - input: renderWakeupInput(registration, resultTaskId, claim.message), + input: renderWakeupInput(registration, resultTaskId, claim.message, orchestrationSnapshot), releaseOnDefiniteFailure: async () => { await methods.releaseSubagentTaskResultClaim({ userId, diff --git a/packages/data-schemas/src/schema/message.ts b/packages/data-schemas/src/schema/message.ts index ad1f59db9e..289b28b12d 100644 --- a/packages/data-schemas/src/schema/message.ts +++ b/packages/data-schemas/src/schema/message.ts @@ -252,6 +252,15 @@ messageSchema.index({ */ messageSchema.index({ conversationId: 1, user: 1, createdAt: 1 }); +/** Bounds parent-run completion snapshots without scanning a user's message history. */ +messageSchema.index( + { user: 1, 'subagentTask.parentRunId': 1, 'subagentTask.status': 1, updatedAt: -1, _id: -1 }, + { + name: 'subagent_parent_run_status_updated', + partialFilterExpression: { 'subagentTask.parentRunId': { $exists: true } }, + }, +); + // index for MeiliSearch sync operations messageSchema.index({ _meiliIndex: 1, isTemporary: 1, expiredAt: 1 }); diff --git a/packages/data-schemas/src/schema/subagent.spec.ts b/packages/data-schemas/src/schema/subagent.spec.ts new file mode 100644 index 0000000000..801c7e608a --- /dev/null +++ b/packages/data-schemas/src/schema/subagent.spec.ts @@ -0,0 +1,21 @@ +import type { IndexDefinition, IndexOptions } from 'mongoose'; +import messageSchema from './message'; + +describe('Subagent task indexes', () => { + it('indexes bounded terminal lookups by parent run and recency', () => { + const indexes = messageSchema.indexes() as Array<[IndexDefinition, IndexOptions]>; + expect(indexes).toContainEqual([ + { + user: 1, + 'subagentTask.parentRunId': 1, + 'subagentTask.status': 1, + updatedAt: -1, + _id: -1, + }, + expect.objectContaining({ + name: 'subagent_parent_run_status_updated', + partialFilterExpression: { 'subagentTask.parentRunId': { $exists: true } }, + }), + ]); + }); +});