diff --git a/CONTEXT.md b/CONTEXT.md index ff6dd95e0f..9764a82526 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -11,4 +11,5 @@ - **Agent event expected action**: an optional source-declared tool name and bounded argument subset evaluated against host-observed completed run steps. It is evidence policy, not authorization and not a model-authored success claim. - **Event actor head**: the private, durable pointer on an event-bound child conversation to its latest committed LangGraph checkpoint, plus one previous checkpoint for safe cleanup. Only a qualifying applied action advances it through compare-and-swap; failed, cancelled, or no-action invocations leave it unchanged. A legacy-path event marks the head for a cold rebuild from durable message history before fork mode can resume. Every applied commit conflict, unverified commit, or post-commit persistence failure is retained in a private reconciliation journal that blocks later actor turns instead of continuing from stale state; an exact marker can be cleared only after its checkpoint is verified authoritative, its history is repaired, or its external action is explicitly compensated. - **Event actor invocation fork**: a delivery-owned checkpoint namespace copied from the event actor head. A warm invocation receives only the new trusted event, then commits its terminal checkpoint when the expected action is observed or deletes the fork otherwise. Pause-capable actors remain on the existing resumable path until forked HITL has an explicit contract. +- **Agent event actor mailbox**: the durable delivery-ordering lane for one authenticated source binding. When enabled after a fleet-wide rollout, it keeps later deliveries queued after transport admission until the current child turn records an authoritative terminal handling outcome. It serializes existing coalesced batches and individual events without becoming a second execution controller or actor checkpoint store. - **Theme definition**: a versioned, data-only description of LibreChat semantic colors and shared appearance roles, optionally specialized by light or dark mode. The theme module validates and resolves partial definitions against bundled defaults before adapters apply them. A theme definition does not contain arbitrary CSS, application behavior, or alternate feature layouts. diff --git a/librechat.example.yaml b/librechat.example.yaml index f167077acb..e5d63e3053 100644 --- a/librechat.example.yaml +++ b/librechat.example.yaml @@ -568,6 +568,10 @@ endpoints: # eventDriven: # childTurns: false # completionWakeups: false + # # Serialize each authenticated bound actor's event lane through the + # # authoritative terminal child-turn outcome. Enable only after every + # # API replica supports terminal-handling mailbox blockers. + # actorMailbox: false # # Reuse bound child state through isolated checkpoint forks. Enable only # # after every API replica supports this lifecycle. Requires the default # # Mongo checkpointer; it is incompatible with checkpointer.type: memory. diff --git a/packages/api/src/agents/triggers/README.md b/packages/api/src/agents/triggers/README.md index 5502d8af7e..732c9eed5f 100644 --- a/packages/api/src/agents/triggers/README.md +++ b/packages/api/src/agents/triggers/README.md @@ -173,6 +173,16 @@ immediately before dispatch, so queued events do not persist stale chat topology is its default ordering lane. A short-lived internal trigger token plus a second binding lookup is required to pass the child-thread write guard; possessing a binding id alone grants no access. +Set `endpoints.agents.eventDriven.actorMailbox: true` only after every API replica runs a release +that understands terminal-handling mailbox blockers. The flag defaults to false for rolling +deployment safety, and `ENABLE_AGENT_EVENT_ACTOR_MAILBOX` remains a compatibility fallback. Once +enabled, a bound actor's next delivery stays queued after the current delivery reaches transport +success and does not dispatch until that child generation records `applied`, +`completed_no_action`, `failed`, or `cancelled`. Different bindings remain independent and can run +in parallel. Existing coalesced batches occupy one mailbox position and retain each member's +individual receipt. An active mailbox record does not receive its normal success TTL; the 90-day +retention window begins only after terminal handling is recorded. + ### Coalescing observational child events Sources that can prove several bound `continue` events are interchangeable observations may opt diff --git a/packages/api/src/agents/triggers/delivery.ts b/packages/api/src/agents/triggers/delivery.ts index 9e77e912a9..33c0e2d345 100644 --- a/packages/api/src/agents/triggers/delivery.ts +++ b/packages/api/src/agents/triggers/delivery.ts @@ -35,6 +35,9 @@ export interface PreparedAgentTriggerDelivery { coalesceKey?: string; coalesceFrom?: Date; coalesceUntil?: Date; + /** Persisted rollout marker: keep this bound actor lane queued until the + * admitted child turn records an authoritative terminal outcome. */ + awaitTerminalHandling?: boolean; } export class AgentTriggerDeliveryError extends TypeError { diff --git a/packages/api/src/agents/triggers/engine.spec.ts b/packages/api/src/agents/triggers/engine.spec.ts index f30e89f9bf..e456fc692c 100644 --- a/packages/api/src/agents/triggers/engine.spec.ts +++ b/packages/api/src/agents/triggers/engine.spec.ts @@ -106,7 +106,9 @@ describe('createAgentTriggerDeliveryEngine', () => { }, input: 'Take the turn.', }); - const store = storeWith({ claimNext: jest.fn(async () => delivery({ envelope })) }); + const store = storeWith({ + claimNext: jest.fn(async () => delivery({ envelope, awaitTerminalHandling: true })), + }); const result: AgentTriggerExecutionResult = { mode: 'continue', status: 'started', @@ -123,6 +125,7 @@ describe('createAgentTriggerDeliveryEngine', () => { expect(store.complete).toHaveBeenCalledWith( expect.objectContaining({ + awaitTerminalHandling: true, handling: { status: 'started', conversationId: 'conversation-1', @@ -545,6 +548,28 @@ describe('createAgentTriggerDeliveryEngine', () => { ); }); + it('rechecks an active actor turn without a tight delivery-lease polling loop', async () => { + const store = storeWith({ + findEarlierUnsettled: jest.fn(async () => ({ + availableAt: START, + reason: 'active_handling' as const, + })), + }); + const dispatch = jest.fn(async () => successResult()); + const engine = createAgentTriggerDeliveryEngine( + { store, dispatch, now: () => START }, + { concurrency: 1 }, + ); + + await engine.runTick(); + + expect(store.release).toHaveBeenCalledWith( + expect.objectContaining({ availableAt: new Date(START.getTime() + 5_000) }), + ); + expect(store.beginAttempt).not.toHaveBeenCalled(); + expect(dispatch).not.toHaveBeenCalled(); + }); + it('starts independent deliveries up to the configured concurrency', async () => { let releaseDispatch: (() => void) | undefined; const gate = new Promise((resolve) => { diff --git a/packages/api/src/agents/triggers/engine.ts b/packages/api/src/agents/triggers/engine.ts index 77e46c515a..502702a6c9 100644 --- a/packages/api/src/agents/triggers/engine.ts +++ b/packages/api/src/agents/triggers/engine.ts @@ -13,6 +13,7 @@ const DEFAULT_RETRY_CAP_MS = 5 * 60_000; const DEFAULT_TICK_MS = 1_000; const DEFAULT_MAX_IDLE_TICK_MS = 15_000; const ORDERING_RECHECK_MS = 250; +const ACTIVE_HANDLING_RECHECK_MS = 5_000; const DEFAULT_DEFER_MS = 5_000; const MAX_RETRY_AFTER_MS = 24 * 60 * 60_000; @@ -97,6 +98,7 @@ export interface AgentTriggerDeliveryRecord { batchMemberIds?: Array<{ toString(): string } | string>; batchRootId?: { toString(): string } | string; batchMembersSettledAt?: Date; + awaitTerminalHandling?: boolean; leaseBy?: string; leaseUntil?: Date; lastError?: AgentTriggerDeliveryFailure; @@ -115,6 +117,7 @@ export interface AgentTriggerDeliveryRecord { export interface AgentTriggerOrderingBlock { availableAt: Date; leaseUntil?: Date; + reason?: 'active_handling'; } export interface AgentTriggerDeliveryStore { @@ -157,6 +160,7 @@ export interface AgentTriggerDeliveryStore { result: AgentTriggerExecutionResult; settledAt: Date; handling?: AgentTriggerDeliveryRecord['handling']; + awaitTerminalHandling?: true; }) => Promise; retry: (input: { id: string; @@ -338,7 +342,9 @@ export function createAgentTriggerDeliveryEngine( const block = await deps.store.findEarlierUnsettled(delivery); if (block != null) { - const recheckAt = now().getTime() + ORDERING_RECHECK_MS; + const recheckAt = + now().getTime() + + (block.reason === 'active_handling' ? ACTIVE_HANDLING_RECHECK_MS : ORDERING_RECHECK_MS); const nextCheck = block.leaseUntil == null ? Math.max(recheckAt, block.availableAt.getTime()) : recheckAt; noteEligibleAt(new Date(nextCheck)); @@ -506,6 +512,7 @@ export function createAgentTriggerDeliveryEngine( attempt, result, settledAt, + ...(delivery.awaitTerminalHandling === true && { awaitTerminalHandling: true }), ...(handling != null && { handling }), }); } catch (error) { diff --git a/packages/api/src/agents/triggers/service.delivery.spec.ts b/packages/api/src/agents/triggers/service.delivery.spec.ts index 26574fd92f..3eec41de66 100644 --- a/packages/api/src/agents/triggers/service.delivery.spec.ts +++ b/packages/api/src/agents/triggers/service.delivery.spec.ts @@ -31,6 +31,29 @@ const envelope = () => input: 'Handle the ready resource.', }); +const boundEnvelope = () => + createAgentTriggerEnvelope({ + mode: 'continue', + requestId: 'request-bound-1', + deliveryId: 'delivery-bound-1', + receivedAt: 20, + principal: { id: '507f1f77bcf86cd799439011', tenantId: 'tenant-1' }, + target: { + agentId: 'agent-1', + conversationId: 'child-conversation-1', + parentMessageId: 'parent-message-1', + bindingId: 'binding-1', + sourceKeyId: 'source-key-1', + }, + event: { + id: 'event-bound-1', + type: 'game.turn', + occurredAt: 10, + source: { id: 'source-key-1', type: 'remote_api_key' }, + }, + input: 'Make the next move.', + }); + function deliveryRecord(overrides: Partial = {}) { return { id: 'delivery-row-1', @@ -182,6 +205,46 @@ describe('durable agent trigger service', () => { await service.stop(); }); + it('persists terminal handling serialization only for opted-in bound continuations', async () => { + const methods = deliveryMethods(); + const service = createAgentTriggerService({ + methods, + actorMailboxEnabled: () => true, + deliveryOptions: { concurrency: 1, tickMs: 60_000 }, + }); + await service.initialize({ address: { address: '127.0.0.1', family: 'IPv4', port: 3080 } }); + + await service.enqueue(boundEnvelope()); + await service.enqueue(envelope()); + + expect(methods.enqueueAgentTriggerDelivery).toHaveBeenNthCalledWith( + 1, + expect.objectContaining({ awaitTerminalHandling: true }), + ); + expect(methods.enqueueAgentTriggerDelivery).toHaveBeenNthCalledWith( + 2, + expect.not.objectContaining({ awaitTerminalHandling: expect.anything() }), + ); + await service.stop(); + }); + + it('does not persist mailbox semantics before the rollout is enabled', async () => { + const methods = deliveryMethods(); + const service = createAgentTriggerService({ + methods, + actorMailboxEnabled: () => false, + deliveryOptions: { concurrency: 1, tickMs: 60_000 }, + }); + await service.initialize({ address: { address: '127.0.0.1', family: 'IPv4', port: 3080 } }); + + await service.enqueue(boundEnvelope()); + + expect(methods.enqueueAgentTriggerDelivery).toHaveBeenCalledWith( + expect.not.objectContaining({ awaitTerminalHandling: expect.anything() }), + ); + await service.stop(); + }); + it('rejects enqueue before persistence when the principal was deleted', async () => { const methods = deliveryMethods(); const service = createAgentTriggerService({ diff --git a/packages/api/src/agents/triggers/service.ts b/packages/api/src/agents/triggers/service.ts index 9fa6566b13..c891fec73f 100644 --- a/packages/api/src/agents/triggers/service.ts +++ b/packages/api/src/agents/triggers/service.ts @@ -47,6 +47,7 @@ export interface AgentTriggerServiceDeps { purgeRecoveryIntervalMs?: number; purgeRecoveryLimit?: number; coalescingEnabled?: () => boolean; + actorMailboxEnabled?: () => boolean; } export interface AgentTriggerDeliveryReceipt { @@ -195,6 +196,8 @@ export function createAgentTriggerService(deps: AgentTriggerServiceDeps = {}): A const purgeRecoveryLimit = deps.purgeRecoveryLimit ?? DEFAULT_PURGE_RECOVERY_LIMIT; const coalescingEnabled = deps.coalescingEnabled ?? (() => isEnabled(process.env.ENABLE_AGENT_EVENT_COALESCING)); + const actorMailboxEnabled = + deps.actorMailboxEnabled ?? (() => isEnabled(process.env.ENABLE_AGENT_EVENT_ACTOR_MAILBOX)); if (!Number.isSafeInteger(userDrainTimeoutMs) || userDrainTimeoutMs <= 0) { throw new TypeError('userDrainTimeoutMs must be a positive integer'); } @@ -406,8 +409,19 @@ export function createAgentTriggerService(deps: AgentTriggerServiceDeps = {}): A throw new AgentTriggerDeliveryError('Agent event coalescing is not enabled on this server'); } const prepared = prepareAgentTriggerDelivery(envelope, options); + const awaitTerminalHandling = + actorMailboxEnabled() && + prepared.envelope.mode === 'continue' && + prepared.envelope.target.bindingId != null && + prepared.envelope.target.sourceKeyId != null; + const durableDelivery: PreparedAgentTriggerDelivery = { + ...prepared, + ...(awaitTerminalHandling && { awaitTerminalHandling: true }), + }; await requireActivePrincipal(String(prepared.user)); - const queued = await runAsSystem(async () => methods.enqueueAgentTriggerDelivery(prepared)); + const queued = await runAsSystem(async () => + methods.enqueueAgentTriggerDelivery(durableDelivery), + ); try { await requireActivePrincipal(String(prepared.user)); } catch (error) { diff --git a/packages/api/src/app/agents.spec.ts b/packages/api/src/app/agents.spec.ts index 88fa43b29b..c15d695afb 100644 --- a/packages/api/src/app/agents.spec.ts +++ b/packages/api/src/app/agents.spec.ts @@ -8,6 +8,7 @@ describe('configureAgentEventRuntime', () => { delete process.env.ENABLE_AGENT_EVENT_CHILD_TURNS; delete process.env.ENABLE_SUBAGENT_COMPLETION_WAKEUPS; delete process.env.ENABLE_AGENT_EVENT_COALESCING; + delete process.env.ENABLE_AGENT_EVENT_ACTOR_MAILBOX; delete process.env.AGENT_TRIGGERS_SELF_URL; }); @@ -20,12 +21,14 @@ describe('configureAgentEventRuntime', () => { childTurns: true, completionWakeups: false, coalescing: true, + actorMailbox: true, selfUrl: 'https://triggers.internal', }); expect(process.env.ENABLE_AGENT_EVENT_CHILD_TURNS).toBe('true'); expect(process.env.ENABLE_SUBAGENT_COMPLETION_WAKEUPS).toBe('false'); expect(process.env.ENABLE_AGENT_EVENT_COALESCING).toBe('true'); + expect(process.env.ENABLE_AGENT_EVENT_ACTOR_MAILBOX).toBe('true'); expect(process.env.AGENT_TRIGGERS_SELF_URL).toBe('https://triggers.internal'); }); @@ -33,6 +36,7 @@ describe('configureAgentEventRuntime', () => { process.env.ENABLE_AGENT_EVENT_CHILD_TURNS = 'true'; process.env.ENABLE_SUBAGENT_COMPLETION_WAKEUPS = 'true'; process.env.ENABLE_AGENT_EVENT_COALESCING = 'true'; + process.env.ENABLE_AGENT_EVENT_ACTOR_MAILBOX = 'true'; process.env.AGENT_TRIGGERS_SELF_URL = 'https://legacy.internal'; configureAgentEventRuntime(undefined); @@ -40,6 +44,7 @@ describe('configureAgentEventRuntime', () => { expect(process.env.ENABLE_AGENT_EVENT_CHILD_TURNS).toBe('true'); expect(process.env.ENABLE_SUBAGENT_COMPLETION_WAKEUPS).toBe('true'); expect(process.env.ENABLE_AGENT_EVENT_COALESCING).toBe('true'); + expect(process.env.ENABLE_AGENT_EVENT_ACTOR_MAILBOX).toBe('true'); expect(process.env.AGENT_TRIGGERS_SELF_URL).toBe('https://legacy.internal'); }); }); diff --git a/packages/api/src/app/agents.ts b/packages/api/src/app/agents.ts index d5299b4dd2..c14fdfa4f2 100644 --- a/packages/api/src/app/agents.ts +++ b/packages/api/src/app/agents.ts @@ -13,6 +13,7 @@ export const configureAgentEventRuntime = (config?: AgentEventRuntimeConfig): vo setBooleanEnvironmentFallback('ENABLE_AGENT_EVENT_CHILD_TURNS', config?.childTurns); setBooleanEnvironmentFallback('ENABLE_SUBAGENT_COMPLETION_WAKEUPS', config?.completionWakeups); setBooleanEnvironmentFallback('ENABLE_AGENT_EVENT_COALESCING', config?.coalescing); + setBooleanEnvironmentFallback('ENABLE_AGENT_EVENT_ACTOR_MAILBOX', config?.actorMailbox); if (config?.selfUrl != null) { process.env.AGENT_TRIGGERS_SELF_URL = config.selfUrl; } diff --git a/packages/data-provider/src/config.spec.ts b/packages/data-provider/src/config.spec.ts index f635494e98..7fc011b59b 100644 --- a/packages/data-provider/src/config.spec.ts +++ b/packages/data-provider/src/config.spec.ts @@ -76,6 +76,7 @@ describe('agent event runtime config', () => { childTurns: true, completionWakeups: false, coalescing: true, + actorMailbox: true, checkpointForks: true, selfUrl: 'https://triggers.internal', }, @@ -94,6 +95,7 @@ describe('agent event runtime config', () => { childTurns: true, completionWakeups: false, coalescing: true, + actorMailbox: true, checkpointForks: true, selfUrl: 'https://triggers.internal', }); diff --git a/packages/data-provider/src/config.ts b/packages/data-provider/src/config.ts index a82c332215..cfa6927027 100644 --- a/packages/data-provider/src/config.ts +++ b/packages/data-provider/src/config.ts @@ -1051,6 +1051,9 @@ export const agentsEndpointSchema = baseEndpointSchema completionWakeups: z.boolean().optional(), /** Enable only after every API worker can consume coalesced deliveries. */ coalescing: z.boolean().optional(), + /** Keep each bound actor's durable delivery lane queued through the + * admitted child turn's authoritative terminal outcome. */ + actorMailbox: z.boolean().optional(), /** Reuse a bound event actor's committed checkpoint through isolated * per-invocation forks. Keep off until every API worker runs an SDK * and host adapter that understand the fork lifecycle. */ diff --git a/packages/data-schemas/src/methods/triggerDelivery.spec.ts b/packages/data-schemas/src/methods/triggerDelivery.spec.ts index 2203430d4e..68772251a6 100644 --- a/packages/data-schemas/src/methods/triggerDelivery.spec.ts +++ b/packages/data-schemas/src/methods/triggerDelivery.spec.ts @@ -103,6 +103,32 @@ describe('agent trigger delivery methods', () => { expect(await Delivery.countDocuments()).toBe(1); }); + it('keeps the original mailbox rollout semantics on idempotent replay', async () => { + const enabled = enqueueInput({ awaitTerminalHandling: true }); + const first = await methods.enqueueAgentTriggerDelivery(enabled); + const disabledReplay = await methods.enqueueAgentTriggerDelivery({ + ...enabled, + awaitTerminalHandling: undefined, + }); + + expect(first.delivery.awaitTerminalHandling).toBe(true); + expect(disabledReplay).toMatchObject({ + replayed: true, + delivery: { id: first.delivery.id, awaitTerminalHandling: true }, + }); + + const disabled = enqueueInput(); + const second = await methods.enqueueAgentTriggerDelivery(disabled); + const enabledReplay = await methods.enqueueAgentTriggerDelivery({ + ...disabled, + awaitTerminalHandling: true, + }); + + expect(second.delivery.awaitTerminalHandling).toBeUndefined(); + expect(enabledReplay).toMatchObject({ replayed: true, delivery: { id: second.delivery.id } }); + expect(enabledReplay.delivery.awaitTerminalHandling).toBeUndefined(); + }); + it('projects public status while enforcing API key, owner, and tenant in the query', async () => { const user = new mongoose.Types.ObjectId(); const queued = await methods.enqueueAgentTriggerDelivery( @@ -226,6 +252,7 @@ describe('agent trigger delivery methods', () => { coalesceFrom, coalesceUntil, availableAt: coalesceUntil, + awaitTerminalHandling: true, envelopeBytes: 128, envelope: { mode: 'continue', @@ -289,6 +316,7 @@ describe('agent trigger delivery methods', () => { attempt: attempt!, result, settledAt: coalesceUntil, + awaitTerminalHandling: true, handling: { status: 'started', conversationId: 'child-thread', @@ -312,6 +340,30 @@ describe('agent trigger delivery methods', () => { expect(receipts).toHaveLength(4); expect(receipts.every((receipt) => receipt?.status === 'succeeded')).toBe(true); expect(receipts.every((receipt) => receipt?.handling?.status === 'started')).toBe(true); + expect( + ( + await Delivery.find({ _id: { $in: queued.map(({ delivery }) => delivery.id) } }).lean() + ).every((row) => row.expiresAt == null), + ).toBe(true); + + const later = await methods.enqueueAgentTriggerDelivery( + enqueueInput({ + user, + orderingKey: 'commentary-lane', + awaitTerminalHandling: true, + availableAt: coalesceUntil, + }), + ); + const laterClaim = await methods.claimNextAgentTriggerDelivery({ + workerId: 'worker-later', + claimToken: 'claim-later', + now: coalesceUntil, + leaseUntil: new Date(coalesceUntil.getTime() + 60_000), + }); + expect(laterClaim?.id).toBe(later.delivery.id); + await expect(methods.findEarlierAgentTriggerDelivery(laterClaim!)).resolves.toMatchObject({ + reason: 'active_handling', + }); await expect( methods.settleAgentTriggerHandlingOutcome({ @@ -322,11 +374,24 @@ describe('agent trigger delivery methods', () => { settledAt: new Date(coalesceUntil.getTime() + 1_000), }), ).resolves.toBe(true); + await expect(methods.findEarlierAgentTriggerDelivery(laterClaim!)).resolves.toBeNull(); let terminal = await Delivery.find({ orderingKey: 'commentary-lane' }).lean(); - expect(terminal.every((row) => row.handling?.status === 'completed_no_action')).toBe(true); + const batchIds = new Set(queued.map(({ delivery }) => delivery.id)); + expect( + terminal + .filter((row) => batchIds.has(String(row._id))) + .every((row) => row.handling?.status === 'completed_no_action'), + ).toBe(true); - const recoveringMember = terminal.find((row) => String(row._id) !== String(root!._id)); + const recoveringMember = terminal.find( + (row) => batchIds.has(String(row._id)) && String(row._id) !== String(root!._id), + ); + await Delivery.updateOne( + { _id: recoveringMember!._id }, + { $set: { 'handling.status': 'started' } }, + ); + await expect(methods.findEarlierAgentTriggerDelivery(laterClaim!)).resolves.toBeNull(); await Delivery.updateOne({ _id: recoveringMember!._id }, { $unset: { handling: 1 } }); await expect( methods.settleAgentTriggerHandlingOutcome({ @@ -338,7 +403,44 @@ describe('agent trigger delivery methods', () => { }), ).resolves.toBe(true); terminal = await Delivery.find({ orderingKey: 'commentary-lane' }).lean(); - expect(terminal.every((row) => row.handling?.status === 'completed_no_action')).toBe(true); + expect( + terminal + .filter((row) => batchIds.has(String(row._id))) + .every((row) => row.handling?.status === 'completed_no_action'), + ).toBe(true); + expect( + terminal + .filter((row) => batchIds.has(String(row._id))) + .every( + (row) => + row.expiresAt?.getTime() === coalesceUntil.getTime() + 1_000 + 90 * 24 * 60 * 60_000, + ), + ).toBe(true); + }); + + it('promotes mailbox semantics when an enabled delivery joins an unmarked batch root', async () => { + const user = new mongoose.Types.ObjectId(); + const coalesceUntil = new Date(Date.now() + 60_000); + const shared = { + user, + orderingKey: 'mixed-rollout-batch', + coalesceKey: 'trigger_batch_mixed_rollout', + coalesceFrom: new Date(coalesceUntil.getTime() - 750), + coalesceUntil, + availableAt: coalesceUntil, + envelopeBytes: 128, + }; + const first = await methods.enqueueAgentTriggerDelivery(enqueueInput(shared)); + await methods.enqueueAgentTriggerDelivery( + enqueueInput({ ...shared, awaitTerminalHandling: true }), + ); + + const root = await Delivery.findById(first.delivery.id).lean(); + expect(root).toMatchObject({ + status: 'pending', + awaitTerminalHandling: true, + batchSize: 2, + }); }); it('recovers batch receipts and lane cleanup after root settlement was interrupted', async () => { @@ -1163,6 +1265,135 @@ describe('agent trigger delivery methods', () => { }); }); + it('keeps a binding lane queued until the earlier child turn settles', async () => { + const user = new mongoose.Types.ObjectId(); + const first = await methods.enqueueAgentTriggerDelivery( + enqueueInput({ user, orderingKey: 'binding-lane', awaitTerminalHandling: true }), + ); + const firstClaim = await methods.claimNextAgentTriggerDelivery({ + workerId: 'worker-1', + claimToken: 'claim-1', + now: START, + leaseUntil: new Date(START.getTime() + 60_000), + }); + const firstAttempt = await methods.beginAgentTriggerDeliveryAttempt({ + id: firstClaim!.id, + workerId: 'worker-1', + claimToken: 'claim-1', + now: START, + }); + const generationCreatedAt = START.getTime() + 1_000; + await methods.completeAgentTriggerDelivery({ + id: firstClaim!.id, + workerId: 'worker-1', + claimToken: 'claim-1', + attempt: firstAttempt!, + result: { + mode: 'continue', + status: 'started', + conversationId: 'actor-thread', + streamId: 'actor-thread', + generationCreatedAt, + }, + settledAt: new Date(generationCreatedAt), + awaitTerminalHandling: true, + handling: { + status: 'started', + conversationId: 'actor-thread', + streamId: 'actor-thread', + generationCreatedAt, + startedAt: new Date(generationCreatedAt), + }, + }); + const second = await methods.enqueueAgentTriggerDelivery( + enqueueInput({ user, orderingKey: 'binding-lane', awaitTerminalHandling: true }), + ); + const secondClaim = await methods.claimNextAgentTriggerDelivery({ + workerId: 'worker-2', + claimToken: 'claim-2', + now: new Date(generationCreatedAt), + leaseUntil: new Date(generationCreatedAt + 60_000), + }); + + expect(firstClaim?.id).toBe(first.delivery.id); + expect(second.delivery.laneSequence).toBe(first.delivery.laneSequence + 1); + await expect(methods.findEarlierAgentTriggerDelivery(secondClaim!)).resolves.toMatchObject({ + availableAt: START, + reason: 'active_handling', + }); + await expect(Delivery.findById(first.delivery.id).lean()).resolves.not.toHaveProperty( + 'expiresAt', + ); + + await expect( + methods.settleAgentTriggerHandlingOutcome({ + deliveryKey: first.delivery.deliveryKey, + conversationId: 'actor-thread', + generationCreatedAt, + status: 'applied', + settledAt: new Date(generationCreatedAt + 1_000), + action: { toolName: 'submit_action' }, + }), + ).resolves.toBe(true); + await expect(Delivery.findById(first.delivery.id).lean()).resolves.toMatchObject({ + expiresAt: new Date(generationCreatedAt + 1_000 + 90 * 24 * 60 * 60_000), + }); + await expect(methods.findEarlierAgentTriggerDelivery(secondClaim!)).resolves.toBeNull(); + }); + + it('reclaims a binding lane only after its admitted child turn settles', async () => { + const queued = await methods.enqueueAgentTriggerDelivery( + enqueueInput({ orderingKey: 'settling-binding-lane', awaitTerminalHandling: true }), + ); + const claim = await methods.claimNextAgentTriggerDelivery({ + workerId: 'worker-1', + claimToken: 'claim-1', + now: START, + leaseUntil: new Date(START.getTime() + 60_000), + }); + const attempt = await methods.beginAgentTriggerDeliveryAttempt({ + id: claim!.id, + workerId: 'worker-1', + claimToken: 'claim-1', + now: START, + }); + const generationCreatedAt = START.getTime() + 1_000; + await methods.completeAgentTriggerDelivery({ + id: claim!.id, + workerId: 'worker-1', + claimToken: 'claim-1', + attempt: attempt!, + result: { + mode: 'continue', + status: 'started', + conversationId: 'actor-thread', + streamId: 'actor-thread', + generationCreatedAt, + }, + settledAt: new Date(generationCreatedAt), + awaitTerminalHandling: true, + handling: { + status: 'started', + conversationId: 'actor-thread', + streamId: 'actor-thread', + generationCreatedAt, + startedAt: new Date(generationCreatedAt), + }, + }); + + await expect(LaneSequence.findById('settling-binding-lane')).resolves.not.toBeNull(); + await expect( + methods.settleAgentTriggerHandlingOutcome({ + deliveryKey: queued.delivery.deliveryKey, + conversationId: 'actor-thread', + generationCreatedAt, + status: 'completed_no_action', + settledAt: new Date(generationCreatedAt + 1_000), + }), + ).resolves.toBe(true); + await expect(LaneSequence.findById('settling-binding-lane')).resolves.toBeNull(); + }); + it('records retry and success outcomes and expires only successful rows', async () => { const queued = await methods.enqueueAgentTriggerDelivery(enqueueInput()); const first = await methods.claimNextAgentTriggerDelivery({ diff --git a/packages/data-schemas/src/methods/triggerDelivery.ts b/packages/data-schemas/src/methods/triggerDelivery.ts index 58cb662567..4161f24876 100644 --- a/packages/data-schemas/src/methods/triggerDelivery.ts +++ b/packages/data-schemas/src/methods/triggerDelivery.ts @@ -49,6 +49,7 @@ export interface EnqueueAgentTriggerDeliveryInput { coalesceKey?: string; coalesceFrom?: Date; coalesceUntil?: Date; + awaitTerminalHandling?: boolean; } export interface AgentTriggerDeliveryFence { @@ -99,6 +100,7 @@ export interface AgentTriggerDeliveryMethods { result: unknown; settledAt: Date; handling?: AgentTriggerHandlingState; + awaitTerminalHandling?: true; }, ) => Promise; settleAgentTriggerHandlingOutcome: ( @@ -196,6 +198,9 @@ function toRecord(delivery: IAgentTriggerDelivery): AgentTriggerDeliveryRecord { ...(delivery.batchMembersSettledAt != null && { batchMembersSettledAt: delivery.batchMembersSettledAt, }), + ...(delivery.awaitTerminalHandling != null && { + awaitTerminalHandling: delivery.awaitTerminalHandling, + }), ...(delivery.handling != null && { handling: delivery.handling }), }; } @@ -315,6 +320,9 @@ export function createAgentTriggerDeliveryMethods( availableAt: staged.coalesceUntil, }, $push: { batchMemberIds: staged._id }, + ...(staged.awaitTerminalHandling === true && { + $set: { awaitTerminalHandling: true }, + }), }, { new: true, sort: { laneSequence: -1, _id: -1 } }, ) @@ -325,6 +333,20 @@ export function createAgentTriggerDeliveryMethods( .lean(); } } + if ( + batchRoot?._id != null && + staged.awaitTerminalHandling === true && + batchRoot.awaitTerminalHandling !== true + ) { + const promoted = await Delivery().updateOne( + { _id: batchRoot._id, batchMemberIds: staged._id }, + { $set: { awaitTerminalHandling: true } }, + ); + if (promoted.matchedCount !== 1) { + throw new Error('Failed to promote terminal handling onto the trigger batch root'); + } + batchRoot.awaitTerminalHandling = true; + } } const published = await Delivery().updateOne( { _id: publisherDeliveryId, orderingKey: lane._id, status: 'staging' }, @@ -657,7 +679,15 @@ export function createAgentTriggerDeliveryMethods( } const stillRetained = await Delivery().exists({ orderingKey, - status: { $in: ['staging', 'batched', 'pending', 'leased', 'dead'] }, + $or: [ + { status: { $in: ['staging', 'batched', 'pending', 'leased', 'dead'] } }, + { + status: 'succeeded', + batchRootId: { $exists: false }, + awaitTerminalHandling: true, + 'handling.status': 'started', + }, + ], }); if (stillRetained != null) { await LaneSequence().updateOne( @@ -774,18 +804,28 @@ export function createAgentTriggerDeliveryMethods( const earlier = await Delivery() .findOne({ orderingKey: delivery.orderingKey, - status: { $in: ['pending', 'leased'] }, laneSequence: { $lt: delivery.laneSequence }, + $or: [ + { status: { $in: ['pending', 'leased'] } }, + { + status: 'succeeded', + batchRootId: { $exists: false }, + awaitTerminalHandling: true, + 'handling.status': 'started', + }, + ], }) .sort({ laneSequence: 1 }) - .select('availableAt leaseUntil') - .lean>(); + .select('availableAt leaseUntil status handling.status') + .lean>(); if (earlier == null) { return null; } return { availableAt: earlier.availableAt, ...(earlier.leaseUntil != null && { leaseUntil: earlier.leaseUntil }), + ...(earlier.status === 'succeeded' && + earlier.handling?.status === 'started' && { reason: 'active_handling' as const }), }; } @@ -831,7 +871,15 @@ export function createAgentTriggerDeliveryMethods( { 'handling.status': { $exists: false } }, ], }, - { $set: { handling: root.handling } }, + { + $set: { + handling: root.handling, + ...(status !== 'started' && + root.handling.settledAt != null && { + expiresAt: new Date(root.handling.settledAt.getTime() + SUCCESS_RETENTION_MS), + }), + }, + }, ); } @@ -843,6 +891,7 @@ export function createAgentTriggerDeliveryMethods( | 'batchMemberIds' | 'batchMembersSettledAt' | 'requeueCount' + | 'awaitTerminalHandling' | 'handling' >, input: { @@ -861,6 +910,8 @@ export function createAgentTriggerDeliveryMethods( if (memberIds.length > 0) { const error = input.error == null ? undefined : normalizeFailure(input.error); const batchRootRequeueCount = root.requeueCount ?? 0; + const awaitsTerminalHandling = + root.awaitTerminalHandling === true && root.handling?.status === 'started'; const settlement = input.status === 'succeeded' ? { @@ -868,7 +919,9 @@ export function createAgentTriggerDeliveryMethods( attempts: input.attempt, result: input.result, settledAt: input.settledAt, - expiresAt: new Date(input.settledAt.getTime() + SUCCESS_RETENTION_MS), + ...(!awaitsTerminalHandling && { + expiresAt: new Date(input.settledAt.getTime() + SUCCESS_RETENTION_MS), + }), batchRootId: root._id, ...(root.handling != null && { handling: root.handling }), } @@ -896,7 +949,9 @@ export function createAgentTriggerDeliveryMethods( leaseBy: 1, leaseUntil: 1, claimToken: 1, - ...(input.status === 'succeeded' ? { lastError: 1 } : { result: 1, expiresAt: 1 }), + ...(input.status === 'succeeded' + ? { lastError: 1, ...(awaitsTerminalHandling && { expiresAt: 1 }) } + : { result: 1, expiresAt: 1 }), }, $push: { history: { @@ -1033,21 +1088,37 @@ export function createAgentTriggerDeliveryMethods( result: unknown; settledAt: Date; handling?: AgentTriggerHandlingState; + awaitTerminalHandling?: true; }, ): Promise { + const awaitsTerminalHandling = + input.awaitTerminalHandling === true && input.handling?.status === 'started'; const completed = await Delivery() .findOneAndUpdate( - fence(input), + { + ...fence(input), + ...(input.awaitTerminalHandling === true + ? { awaitTerminalHandling: true } + : { awaitTerminalHandling: { $ne: true } }), + }, { $set: { status: 'succeeded', result: input.result, settledAt: input.settledAt, - expiresAt: new Date(input.settledAt.getTime() + SUCCESS_RETENTION_MS), + ...(!awaitsTerminalHandling && { + expiresAt: new Date(input.settledAt.getTime() + SUCCESS_RETENTION_MS), + }), laneCleanupPendingAt: input.settledAt, ...(input.handling != null && { handling: input.handling }), }, - $unset: { leaseBy: 1, leaseUntil: 1, claimToken: 1, lastError: 1 }, + $unset: { + leaseBy: 1, + leaseUntil: 1, + claimToken: 1, + lastError: 1, + ...(awaitsTerminalHandling && { expiresAt: 1 }), + }, $push: { history: { $each: [ @@ -1065,7 +1136,7 @@ export function createAgentTriggerDeliveryMethods( { new: true }, ) .select( - '_id orderingKey laneCleanupPendingAt batchMemberIds batchMembersSettledAt requeueCount handling', + '_id orderingKey laneCleanupPendingAt batchMemberIds batchMembersSettledAt requeueCount awaitTerminalHandling handling', ) .lean< Pick< @@ -1076,6 +1147,7 @@ export function createAgentTriggerDeliveryMethods( | 'batchMemberIds' | 'batchMembersSettledAt' | 'requeueCount' + | 'awaitTerminalHandling' | 'handling' > >(); @@ -1110,6 +1182,7 @@ export function createAgentTriggerDeliveryMethods( const terminalHandling = { 'handling.status': input.status, 'handling.settledAt': input.settledAt, + expiresAt: new Date(input.settledAt.getTime() + SUCCESS_RETENTION_MS), ...(error != null && { 'handling.error': error }), ...(input.action != null && { 'handling.action': input.action }), }; @@ -1130,8 +1203,13 @@ export function createAgentTriggerDeliveryMethods( }, { new: true }, ) - .select('_id orderingKey batchMemberIds handling') - .lean>(); + .select('_id orderingKey batchMemberIds awaitTerminalHandling handling') + .lean< + Pick< + IAgentTriggerDelivery, + '_id' | 'orderingKey' | 'batchMemberIds' | 'awaitTerminalHandling' | 'handling' + > + >(); let authoritative = terminal; if (authoritative == null) { @@ -1141,8 +1219,13 @@ export function createAgentTriggerDeliveryMethods( 'handling.conversationId': input.conversationId, 'handling.generationCreatedAt': input.generationCreatedAt, }) - .select('_id orderingKey batchMemberIds handling') - .lean>(); + .select('_id orderingKey batchMemberIds awaitTerminalHandling handling') + .lean< + Pick< + IAgentTriggerDelivery, + '_id' | 'orderingKey' | 'batchMemberIds' | 'awaitTerminalHandling' | 'handling' + > + >(); const replayed = existing?.handling?.status === input.status && existing.handling.error === error && @@ -1155,6 +1238,13 @@ export function createAgentTriggerDeliveryMethods( } await propagateBatchHandling(authoritative); + if (authoritative.awaitTerminalHandling === true) { + await LaneSequence().updateOne( + { _id: authoritative.orderingKey }, + { $set: { cleanupRequestedAt: input.settledAt } }, + ); + await reclaimLaneIfInactive(authoritative.orderingKey); + } return true; } diff --git a/packages/data-schemas/src/schema/triggerDelivery.ts b/packages/data-schemas/src/schema/triggerDelivery.ts index 7b1b1cdf9b..a7bbd27bea 100644 --- a/packages/data-schemas/src/schema/triggerDelivery.ts +++ b/packages/data-schemas/src/schema/triggerDelivery.ts @@ -82,6 +82,7 @@ const triggerDeliverySchema: Schema = new Schema( batchRootId: { type: Schema.Types.ObjectId, ref: 'AgentTriggerDelivery' }, batchRootRequeueCount: { type: Number, min: 0 }, batchMembersSettledAt: { type: Date }, + awaitTerminalHandling: { type: Boolean }, handling: { type: handlingSchema }, leaseBy: { type: String }, leaseUntil: { type: Date }, @@ -102,6 +103,14 @@ triggerDeliverySchema.index({ deliveryKey: 1 }, { unique: true }); triggerDeliverySchema.index({ status: 1, availableAt: 1, createdAt: 1 }); triggerDeliverySchema.index({ status: 1, leaseUntil: 1, createdAt: 1 }); triggerDeliverySchema.index({ orderingKey: 1, status: 1, laneSequence: 1 }); +triggerDeliverySchema.index({ + orderingKey: 1, + awaitTerminalHandling: 1, + status: 1, + 'handling.status': 1, + batchRootId: 1, + laneSequence: 1, +}); triggerDeliverySchema.index({ batchRootId: 1 }, { sparse: true }); triggerDeliverySchema.index( { orderingKey: 1, coalesceKey: 1, status: 1, coalesceUntil: 1 }, diff --git a/packages/data-schemas/src/types/triggerDelivery.ts b/packages/data-schemas/src/types/triggerDelivery.ts index fb6b3c378c..79a3f6501b 100644 --- a/packages/data-schemas/src/types/triggerDelivery.ts +++ b/packages/data-schemas/src/types/triggerDelivery.ts @@ -63,6 +63,9 @@ export interface IAgentTriggerDelivery { batchRootId?: Types.ObjectId; batchRootRequeueCount?: number; batchMembersSettledAt?: Date; + /** Keeps this binding lane serialized until its admitted child turn reaches + * an authoritative terminal handling outcome. */ + awaitTerminalHandling?: boolean; handling?: AgentTriggerHandlingState; leaseBy?: string; leaseUntil?: Date; @@ -147,4 +150,5 @@ export interface AgentTriggerDeliveryClaim extends AgentTriggerDeliveryRecord { export interface AgentTriggerOrderingBlock { availableAt: Date; leaseUntil?: Date; + reason?: 'active_handling'; }