From cd3768ed1fc2491806689678cc5b64765e29cf22 Mon Sep 17 00:00:00 2001 From: Danny Avila Date: Sun, 30 Aug 2026 11:54:17 -0400 Subject: [PATCH] =?UTF-8?q?=F0=9F=8D=B5=20feat:=20Continue=20Late=20Steers?= =?UTF-8?q?=20in=20Warm=20Agent=20Runs=20(#15357)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * feat: continue late steers in warm agent runs * chore: bump agents sdk to v3.7.9 * fix: guard terminal steer admission --- CONTEXT.md | 1 + api/package.json | 2 +- .../__tests__/client.steerWiring.spec.js | 19 +- .../request.partialDisconnect.spec.js | 1 + .../__tests__/request.resumeMetadata.spec.js | 1 + .../agents/__tests__/resume.spec.js | 1 + api/server/controllers/agents/client.js | 11 +- package-lock.json | 10 +- packages/api/package.json | 2 +- .../src/agents/hooks/compatibility.spec.ts | 6 +- packages/api/src/agents/hooks/executor.ts | 1 + packages/api/src/agents/hooks/runtime.ts | 3 + packages/api/src/agents/run.ts | 18 +- .../agents/steering/__tests__/runtime.spec.ts | 167 ++++++++++++++++++ packages/api/src/agents/steering/index.ts | 9 +- packages/api/src/agents/steering/runtime.ts | 62 ++++++- packages/api/src/stream/SteeringLifecycle.ts | 10 ++ .../RedisJobStore.stream_integration.spec.ts | 86 +++++++++ .../api/src/stream/__tests__/steering.spec.ts | 105 +++++++++++ .../implementations/InMemoryJobStore.ts | 35 ++++ .../stream/implementations/RedisJobStore.ts | 67 +++++++ .../api/src/stream/interfaces/IJobStore.ts | 24 +++ .../api/src/stream/jobStoreCapabilities.ts | 1 + 23 files changed, 624 insertions(+), 18 deletions(-) diff --git a/CONTEXT.md b/CONTEXT.md index d084383e7a..b51a341892 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -10,6 +10,7 @@ - **Live subagent task owner**: the one API process holding a detached child execution, its abort controller, and its bounded control queue. Redis may route trusted poll/control envelopes to that owner, but it does not migrate or persist the executor; Mongo persists only the logical child thread and its continuation fence. - **Subagent completion wakeup**: a durable internal `continue` trigger pre-registered before detached child execution so a process crash cannot lose the wakeup. Delivery defers until the child's terminal transcript is persisted, targets the initiating agent and exact parent response branch, carries task metadata rather than child output, waits for the parent generation to settle, and starts the parent turn that collects the result through the existing task store. - **Agent continuation preparation**: the single source-dispatch seam that resolves a durable `continue` delivery immediately before admission. Bound Event Actor work selects its binding adapter; internal completion work selects an adapter by stable source identity. Preparation may resolve authoritative input and branch state or settle already-consumed work, but it does not own source result truth, delivery ordering, or generation execution. +- **Warm terminal steer continuation**: a queued steer accepted before a generation's terminal boundary may continue the same SDK `Run` without creating a replacement generation. After parallel Stop hooks fold, the serialized StopFinalize phase tells the job store whether another continuation is already planned or terminal progress is forbidden. The store atomically chooses among claiming the current protocol-v2 FIFO batch, keeping empty admission open for an already-planned segment, and sealing admission so every racing or later message becomes an ordinary follow-up. Claimed steer receipts remain the crash-recovery authority; protocol-v1 generations always seal because they cannot recover an ambiguous terminal claim. Tool-batch, preemption, and terminal boundaries share one durable apply-and-inject adapter, while the SDK owns the bounded Stop-continuation loop. - **Subagent activity stream**: an observational, task-scoped live projection of bounded child progress for the currently open private panel. It may cross API replicas through Redis, never carries hidden reasoning text, and never controls or settles execution. The durable child thread remains canonical and its existing polling view is the fallback for missed or unavailable live events. - **Agent event handling outcome**: the durable, generation-fenced result of a previously accepted event delivery. `started` proves generation admission; terminal states distinguish verified tool application, clean completion without action, failure, and cancellation. Transport success remains separate so an accepted event cannot masquerade as completed work. - **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. diff --git a/api/package.json b/api/package.json index d0e6c09161..7ed9fbcef4 100644 --- a/api/package.json +++ b/api/package.json @@ -46,7 +46,7 @@ "@azure/storage-blob": "^12.30.0", "@google/genai": "^2.8.0", "@keyv/redis": "5.1.6", - "@librechat/agents": "^3.7.8", + "@librechat/agents": "^3.7.9", "@librechat/api": "*", "@librechat/data-schemas": "*", "@microsoft/microsoft-graph-client": "^3.0.7", diff --git a/api/server/controllers/agents/__tests__/client.steerWiring.spec.js b/api/server/controllers/agents/__tests__/client.steerWiring.spec.js index 97efceb72f..f599d14b1e 100644 --- a/api/server/controllers/agents/__tests__/client.steerWiring.spec.js +++ b/api/server/controllers/agents/__tests__/client.steerWiring.spec.js @@ -1,14 +1,20 @@ const AgentClient = require('../client'); -const { isSteeringSupported, isSteerPreemptSupported } = require('@librechat/api'); +const { + isSteeringSupported, + isSteerPreemptSupported, + isSteerTerminalContinuationSupported, +} = require('@librechat/api'); jest.mock('@librechat/api', () => ({ ...jest.requireActual('@librechat/api'), isSteeringSupported: jest.fn(() => true), isSteerPreemptSupported: jest.fn(() => true), + isSteerTerminalContinuationSupported: jest.fn(() => true), })); const mockIsSteeringSupported = isSteeringSupported; const mockIsPreemptSupported = isSteerPreemptSupported; +const mockIsTerminalContinuationSupported = isSteerTerminalContinuationSupported; /** Minimal `this` for the wiring builder — it only reads these three. */ function buildWiring(streamId, { jobCreatedAt = 1700000000000 } = {}) { @@ -25,6 +31,7 @@ describe('AgentClient.buildSteerWiring — preempt capability gating', () => { jest.clearAllMocks(); mockIsSteeringSupported.mockReturnValue(true); mockIsPreemptSupported.mockReturnValue(true); + mockIsTerminalContinuationSupported.mockReturnValue(true); }); it('returns both boundary hooks and the poll when preempt is supported', () => { @@ -33,6 +40,16 @@ describe('AgentClient.buildSteerWiring — preempt capability gating', () => { expect(typeof wiring.hook).toBe('function'); expect(typeof wiring.preemptHook).toBe('function'); expect(typeof wiring.preemption?.shouldPreempt).toBe('function'); + expect(typeof wiring.terminalHook).toBe('function'); + }); + + it('omits only terminal continuation when the SDK lacks Stop continuation', () => { + mockIsTerminalContinuationSupported.mockReturnValue(false); + const wiring = buildWiring('stream-terminal-unsupported'); + + expect(typeof wiring.hook).toBe('function'); + expect(typeof wiring.preemptHook).toBe('function'); + expect(wiring.terminalHook).toBeUndefined(); }); /** diff --git a/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js b/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js index f8b30502a2..85caf062cb 100644 --- a/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js +++ b/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js @@ -51,6 +51,7 @@ jest.mock('@librechat/api', () => ({ exemptFromConcurrencyLimiter: jest.fn(() => false), toPendingSteer: jest.fn((item) => item), isSteerPreemptSupported: jest.fn(() => true), + isSteerTerminalContinuationSupported: jest.fn(() => false), buildRecoveredSteerPayload: jest.fn(() => null), deleteAgentCheckpoint: jest.fn(), getViolationInfo: jest.fn(() => ({ diff --git a/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js b/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js index d8bbb55b26..693aa6cb56 100644 --- a/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js +++ b/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js @@ -254,6 +254,7 @@ jest.mock('@librechat/api', () => ({ /** Recorded onto the job so the steer route can honour the OWNING replica's * seal capability rather than its own probe. */ isSteerPreemptSupported: jest.fn(() => true), + isSteerTerminalContinuationSupported: jest.fn(() => false), buildRecoveredSteerPayload: jest.fn((text, files) => { if (typeof text !== 'string' || (files != null && !Array.isArray(files))) { return null; diff --git a/api/server/controllers/agents/__tests__/resume.spec.js b/api/server/controllers/agents/__tests__/resume.spec.js index 0677f12222..8dee4b6aa4 100644 --- a/api/server/controllers/agents/__tests__/resume.spec.js +++ b/api/server/controllers/agents/__tests__/resume.spec.js @@ -132,6 +132,7 @@ jest.mock('@librechat/api', () => ({ decrementPendingRequest: (...args) => mockDecrementPendingRequest(...args), checkAndIncrementPendingRequest: (...args) => mockCheckAndIncrementPendingRequest(...args), isSteerPreemptSupported: jest.fn(() => true), + isSteerTerminalContinuationSupported: jest.fn(() => false), createMCPRuntimeRequestBody: ({ messageId, conversationId, parentMessageId }) => ({ messageId, conversationId, diff --git a/api/server/controllers/agents/client.js b/api/server/controllers/agents/client.js index 37bc93ea66..7ce1bde2e1 100644 --- a/api/server/controllers/agents/client.js +++ b/api/server/controllers/agents/client.js @@ -65,9 +65,11 @@ const { createSteerIndexOffsetHandlers, createSteerDrainHook, createSteerPreemptBoundaryHook, + createSteerTerminalContinuationHook, createSteerPreemptPoll, isSteeringSupported, isSteerPreemptSupported, + isSteerTerminalContinuationSupported, buildSteerMedia, collectSteerStampTargets, stampSteerPartMedia, @@ -557,9 +559,9 @@ class AgentClient extends BaseClient { /** * The `steering` fragment for `createRun`: the run-scoped PostToolBatch - * drain hook — plus, when the SDK can seal mid-stream, the PreemptBoundary - * twin and the preempt poll built from the SAME drain closures, so both - * boundaries inject byte-identical shapes. `undefined` when there is no + * drain hook — plus the capability-gated PreemptBoundary and terminal Stop + * twins built from the SAME drain closures, so every boundary injects + * byte-identical shapes. `undefined` when there is no * resumable job surface or the installed SDK cannot inject hook messages * (draining would drop them). * @@ -592,6 +594,9 @@ class AgentClient extends BaseClient { preemptHook: createSteerPreemptBoundaryHook(drainOptions), preemption: createSteerPreemptPoll(streamId), }), + ...(isSteerTerminalContinuationSupported() && { + terminalHook: createSteerTerminalContinuationHook(drainOptions), + }), }; } diff --git a/package-lock.json b/package-lock.json index 472ee46cec..acc0c71b66 100644 --- a/package-lock.json +++ b/package-lock.json @@ -63,7 +63,7 @@ "@azure/storage-blob": "^12.30.0", "@google/genai": "^2.8.0", "@keyv/redis": "5.1.6", - "@librechat/agents": "^3.7.8", + "@librechat/agents": "^3.7.9", "@librechat/api": "*", "@librechat/data-schemas": "*", "@microsoft/microsoft-graph-client": "^3.0.7", @@ -10633,9 +10633,9 @@ } }, "node_modules/@librechat/agents": { - "version": "3.7.8", - "resolved": "https://registry.npmjs.org/@librechat/agents/-/agents-3.7.8.tgz", - "integrity": "sha512-XLUnZaGdPovtwXAc1SxSfUmQi/D4zu6hNRpEoqfLIF4QkA8DOV8wj6ruztuPR4bbMOP1dBze8wqmHEUNVM4zdw==", + "version": "3.7.9", + "resolved": "https://registry.npmjs.org/@librechat/agents/-/agents-3.7.9.tgz", + "integrity": "sha512-PF7KM3W1K5PwwtsW84X9xBoCRS3Hr0HB8BZ2cuehBUHCpXDCXVathaqlzgxwBUHYGg/XOrnTwR/IvTZ1hkzUYg==", "license": "MIT", "dependencies": { "@anthropic-ai/sdk": "^0.115.0", @@ -42839,7 +42839,7 @@ "@azure/storage-blob": "^12.30.0", "@google/genai": "^2.8.0", "@keyv/redis": "5.1.6", - "@librechat/agents": "^3.7.8", + "@librechat/agents": "^3.7.9", "@librechat/data-schemas": "*", "@modelcontextprotocol/sdk": "^1.30.0", "@opentelemetry/api": "^1.9.0", diff --git a/packages/api/package.json b/packages/api/package.json index b601ec625c..8056d28b95 100644 --- a/packages/api/package.json +++ b/packages/api/package.json @@ -113,7 +113,7 @@ "@azure/storage-blob": "^12.30.0", "@google/genai": "^2.8.0", "@keyv/redis": "5.1.6", - "@librechat/agents": "^3.7.8", + "@librechat/agents": "^3.7.9", "@librechat/data-schemas": "*", "@modelcontextprotocol/sdk": "^1.30.0", "@opentelemetry/api": "^1.9.0", diff --git a/packages/api/src/agents/hooks/compatibility.spec.ts b/packages/api/src/agents/hooks/compatibility.spec.ts index 5c97c3d5e9..10472169b3 100644 --- a/packages/api/src/agents/hooks/compatibility.spec.ts +++ b/packages/api/src/agents/hooks/compatibility.spec.ts @@ -357,16 +357,20 @@ describe('planPluginHooks', () => { const plan = planPluginHooks( document({ UserPromptExpansion: [{ hooks: [{ type: 'command', command: 'banner' }] }], + StopFinalize: [{ hooks: [{ type: 'command', command: 'intercept-finalization' }] }], Stop: [{ hooks: [{ type: 'prompt', prompt: 'Verify completion' }] }], }), commandCapabilities, ); - expect(plan.summary).toEqual({ declared: 2, ready: 0, unsupported: 2 }); + expect(plan.summary).toEqual({ declared: 3, ready: 0, unsupported: 3 }); expect(plan.entries[0].issues).toEqual([ expect.objectContaining({ code: 'unsupported_event' }), ]); expect(plan.entries[1].issues).toEqual([ + expect.objectContaining({ code: 'unsupported_event' }), + ]); + expect(plan.entries[2].issues).toEqual([ expect.objectContaining({ code: 'unsupported_handler' }), ]); }); diff --git a/packages/api/src/agents/hooks/executor.ts b/packages/api/src/agents/hooks/executor.ts index e0b626fde4..ea2e6d7099 100644 --- a/packages/api/src/agents/hooks/executor.ts +++ b/packages/api/src/agents/hooks/executor.ts @@ -54,6 +54,7 @@ const EVENT_TRAITS = { SubagentStart: { toolMatcher: false, decisions: 'tool', exitTwo: 'deny' }, SubagentStop: { toolMatcher: false, decisions: 'tool', exitTwo: 'prevent' }, Stop: { toolMatcher: false, decisions: 'stop', exitTwo: 'block' }, + StopFinalize: { toolMatcher: false, decisions: 'stop', exitTwo: 'block' }, StopFailure: { toolMatcher: false, decisions: 'tool', exitTwo: 'prevent' }, PreCompact: { toolMatcher: false, decisions: 'tool', exitTwo: 'prevent' }, PostCompact: { toolMatcher: false, decisions: 'tool', exitTwo: 'prevent' }, diff --git a/packages/api/src/agents/hooks/runtime.ts b/packages/api/src/agents/hooks/runtime.ts index 0425835f53..984b2739a3 100644 --- a/packages/api/src/agents/hooks/runtime.ts +++ b/packages/api/src/agents/hooks/runtime.ts @@ -281,6 +281,9 @@ export function createPluginHookPayload( return { ...payload, agent_type: input.agentType }; case 'Stop': return { ...payload, stop_hook_active: input.stopHookActive }; + /** Internal terminal-admission phase; plugin declarations cannot target it. */ + case 'StopFinalize': + return payload; case 'StopFailure': { const lastAssistantMessage = getMessageText(input.lastAssistantMessage); return { diff --git a/packages/api/src/agents/run.ts b/packages/api/src/agents/run.ts index 0757b8439d..e49965a68c 100644 --- a/packages/api/src/agents/run.ts +++ b/packages/api/src/agents/run.ts @@ -48,6 +48,7 @@ import type { AppConfig, IUser } from '@librechat/data-schemas'; import type { ModelBoundChatModelCallback } from '~/middleware/modelBoundContent'; import type { ToolInputValidationError } from '~/agents/toolValidation'; import type { ResolvedToolApprovalHook } from '~/agents/hitl/hooks'; +import type { TerminalSteerHook } from '~/agents/steering/runtime'; import type { ResolvedAlwaysApplySkill } from '~/agents/skills'; import type { MCPToolAlias } from '~/tools/classification'; import type { SubagentUsageEvent } from '~/agents/usage'; @@ -63,6 +64,11 @@ import { agentUsesSubagentCompletionWakeups, usesSubagentCompletionWakeups, } from '~/agents/subagentDelivery'; +import { + isSteeringSupported, + isSteerPreemptSupported, + isSteerTerminalContinuationSupported, +} from '~/agents/steering/runtime'; import { resolveToolApprovalPolicy, healToolApprovalPolicy, @@ -74,7 +80,6 @@ import { } from '~/agents/hitl/askUserQuestionTool'; import { applyCustomHandoffPromptKeyCompatibility } from '~/agents/handoffPromptKeyCompatibility'; import { stripIntentFromToolRegistry, stripIntentFromToolDefinitions } from '~/agents/intent'; -import { isSteeringSupported, isSteerPreemptSupported } from '~/agents/steering/runtime'; import { extractDefaultParams, resolveReasoningParams } from '~/endpoints/openai/llm'; import { getLLMConfig as getAnthropicLLMConfig } from '~/endpoints/anthropic/llm'; import { resolveStreamLimits, resolveSubagentMaxTurns } from '~/agents/config'; @@ -1473,6 +1478,11 @@ export async function createRun({ * when the SDK seals a model stream mid-generation on a preempt request. */ preemptHook?: HookCallback<'PreemptBoundary'>; + /** + * Atomically claims queued steers at the SDK's terminal Stop boundary or + * seals admission so later messages become ordinary follow-up turns. + */ + terminalHook?: TerminalSteerHook; /** * Level-triggered O(1) poll over the job's armed preempt requests * (`createSteerPreemptPoll`). Threaded into `RunConfig.preemption`, which @@ -1974,6 +1984,12 @@ export async function createRun({ if (steering.preemptHook != null && isSteerPreemptSupported()) { hooks.register('PreemptBoundary', { hooks: [steering.preemptHook] }); } + if (steering.terminalHook != null && isSteerTerminalContinuationSupported()) { + const stopFinalizeRegistry = hooks as unknown as { + register: (event: 'StopFinalize', matcher: { hooks: TerminalSteerHook[] }) => () => void; + }; + stopFinalizeRegistry.register('StopFinalize', { hooks: [steering.terminalHook] }); + } } /** * Deployment-plugin hooks (Agent Plugins `ai.librechat/hooks/hooks.json`) diff --git a/packages/api/src/agents/steering/__tests__/runtime.spec.ts b/packages/api/src/agents/steering/__tests__/runtime.spec.ts index 883a5dbc49..ff37d594d4 100644 --- a/packages/api/src/agents/steering/__tests__/runtime.spec.ts +++ b/packages/api/src/agents/steering/__tests__/runtime.spec.ts @@ -5,6 +5,7 @@ import type { PreemptBoundaryHookInput, } from '@librechat/agents'; import type { SteerQueueItem } from '~/stream/interfaces/IJobStore'; +import { STEER_ENQUEUE_NOT_RUNNING } from '~/stream/interfaces/IJobStore'; /** The pinned SDK's hook output declares `injectedMessages` natively; a * narrower local re-declaration would no longer be assignable from it. */ @@ -15,10 +16,13 @@ import { GenerationJobManager } from '~/stream/GenerationJobManager'; import { createSteerDrainHook, createSteerPreemptBoundaryHook, + createSteerTerminalContinuationHook, createSteerPreemptPoll, isSteeringSupported, isSteerPreemptSupported, + isSteerTerminalContinuationSupported, } from '../runtime'; +import type { TerminalSteerHookInput } from '../runtime'; jest.spyOn(console, 'log').mockImplementation(); @@ -48,6 +52,22 @@ function boundaryInput( }; } +function stopInput( + continuationBudgetRemaining: number, + stopReason?: string, + overrides: Partial = {}, +): TerminalSteerHookInput { + return { + hook_event_name: 'StopFinalize', + runId: 'run-1', + continuationBudgetRemaining, + continuationPlanned: false, + continuationPrevented: stopReason != null, + ...(stopReason != null && { stopReason }), + ...overrides, + } as unknown as TerminalSteerHookInput; +} + describe('isSteeringSupported', () => { it('mirrors the installed SDK capability flag AND replay support', () => { // CI runs against the published SDK pin (possibly pre-injectedMessages); @@ -65,6 +85,153 @@ describe('isSteeringSupported', () => { sdk.HOOK_INJECTED_MESSAGES_CAPABLE === true && sdk.ContentTypes?.STEER === 'steer'; expect(isSteeringSupported()).toBe(capable); }); + + it('gates terminal continuation on its separate SDK capability', () => { + const sdk = agentsSdk as { HOOK_STOP_CONTINUATION_CAPABLE?: boolean }; + expect(isSteerTerminalContinuationSupported()).toBe( + isSteeringSupported() && sdk.HOOK_STOP_CONTINUATION_CAPABLE === true, + ); + }); +}); + +describe('createSteerTerminalContinuationHook', () => { + beforeEach(() => { + GenerationJobManager.configure({ + jobStore: new InMemoryJobStore({ ttlAfterComplete: 60000 }), + eventTransport: new InMemoryEventTransport(), + isRedis: false, + cleanupOnComplete: false, + }); + GenerationJobManager.initialize(); + }); + + afterEach(async () => { + await GenerationJobManager.destroy(); + }); + + it('claims queued steers and blocks Stop into the same warm Run', async () => { + const streamId = `terminal-claim-${Date.now()}`; + const job = await GenerationJobManager.createJob(streamId, 'user-1'); + await GenerationJobManager.steering.enqueue(streamId, buildSteer('s1', 'keep going')); + const applied = jest.fn(); + const hook = createSteerTerminalContinuationHook({ + streamId, + jobCreatedAt: job.createdAt, + applySteer: applied, + }); + + await expect(hook(stopInput(1), abortSignal)).resolves.toEqual({ + decision: 'block', + injectedMessages: [{ role: 'user', content: 'keep going', source: 'steer' }], + }); + expect(applied).toHaveBeenCalledWith(expect.objectContaining({ steerId: 's1' })); + await expect(GenerationJobManager.steering.peek(streamId, job.createdAt)).resolves.toEqual([]); + await expect( + GenerationJobManager.steering.enqueue( + streamId, + buildSteer('s2', 'next continuation'), + job.createdAt, + ), + ).resolves.toBe(1); + }); + + it('never admits terminal steers inside a subagent scope', async () => { + const streamId = `terminal-subagent-${Date.now()}`; + const job = await GenerationJobManager.createJob(streamId, 'user-1'); + const queued = buildSteer('s1', 'keep for the parent'); + await GenerationJobManager.steering.enqueue(streamId, queued); + const applySteer = jest.fn(); + const hook = createSteerTerminalContinuationHook({ + streamId, + jobCreatedAt: job.createdAt, + applySteer, + }); + + await expect( + hook(stopInput(1, undefined, { agentId: 'child-agent' }), abortSignal), + ).resolves.toEqual({ decision: 'continue' }); + expect(applySteer).not.toHaveBeenCalled(); + await expect(GenerationJobManager.steering.peek(streamId, job.createdAt)).resolves.toEqual([ + queued, + ]); + }); + + it('seals admission when the continuation budget is exhausted without losing the queue', async () => { + const streamId = `terminal-budget-${Date.now()}`; + const job = await GenerationJobManager.createJob(streamId, 'user-1'); + const queued = buildSteer('s1', 'ordinary follow-up'); + await GenerationJobManager.steering.enqueue(streamId, queued); + const hook = createSteerTerminalContinuationHook({ + streamId, + jobCreatedAt: job.createdAt, + applySteer: jest.fn(), + }); + + await expect(hook(stopInput(0), abortSignal)).resolves.toEqual({ decision: 'continue' }); + await expect( + GenerationJobManager.steering.enqueue(streamId, buildSteer('s2', 'too late'), job.createdAt), + ).resolves.toBe(STEER_ENQUEUE_NOT_RUNNING); + await expect( + GenerationJobManager.steering.closeAndDrain(streamId, job.createdAt), + ).resolves.toEqual([queued]); + }); + + it('seals an empty terminal boundary so a later steer becomes a new turn', async () => { + const streamId = `terminal-empty-${Date.now()}`; + const job = await GenerationJobManager.createJob(streamId, 'user-1'); + const hook = createSteerTerminalContinuationHook({ + streamId, + jobCreatedAt: job.createdAt, + applySteer: jest.fn(), + }); + + await expect(hook(stopInput(1), abortSignal)).resolves.toEqual({ decision: 'continue' }); + await expect( + GenerationJobManager.steering.enqueue(streamId, buildSteer('s1', 'new turn'), job.createdAt), + ).resolves.toBe(STEER_ENQUEUE_NOT_RUNNING); + }); + + it('keeps empty admission open when another Stop hook already planned a continuation', async () => { + const streamId = `terminal-other-hook-${Date.now()}`; + const job = await GenerationJobManager.createJob(streamId, 'user-1'); + const hook = createSteerTerminalContinuationHook({ + streamId, + jobCreatedAt: job.createdAt, + applySteer: jest.fn(), + }); + + await expect( + hook(stopInput(1, undefined, { continuationPlanned: true }), abortSignal), + ).resolves.toEqual({ decision: 'continue' }); + await expect( + GenerationJobManager.steering.enqueue( + streamId, + buildSteer('s1', 'join the planned continuation'), + job.createdAt, + ), + ).resolves.toBe(1); + }); + + it('seals instead of claiming when the graph has a terminal halt reason', async () => { + const streamId = `terminal-halt-${Date.now()}`; + const job = await GenerationJobManager.createJob(streamId, 'user-1'); + const queued = buildSteer('s1', 'retry in a fresh turn'); + await GenerationJobManager.steering.enqueue(streamId, queued, job.createdAt); + const applySteer = jest.fn(); + const hook = createSteerTerminalContinuationHook({ + streamId, + jobCreatedAt: job.createdAt, + applySteer, + }); + + await expect(hook(stopInput(1, 'preempt_incomplete'), abortSignal)).resolves.toEqual({ + decision: 'continue', + }); + expect(applySteer).not.toHaveBeenCalled(); + await expect( + GenerationJobManager.steering.closeAndDrain(streamId, job.createdAt), + ).resolves.toEqual([queued]); + }); }); describe('createSteerDrainHook', () => { diff --git a/packages/api/src/agents/steering/index.ts b/packages/api/src/agents/steering/index.ts index 4f1eab09c1..22af745578 100644 --- a/packages/api/src/agents/steering/index.ts +++ b/packages/api/src/agents/steering/index.ts @@ -1,11 +1,18 @@ export { createSteerDrainHook, createSteerPreemptBoundaryHook, + createSteerTerminalContinuationHook, createSteerPreemptPoll, isSteeringSupported, isSteerPreemptSupported, + isSteerTerminalContinuationSupported, +} from './runtime'; +export type { + SteerDrainHookOptions, + SteerMediaResult, + TerminalSteerHook, + TerminalSteerHookInput, } from './runtime'; -export type { SteerDrainHookOptions, SteerMediaResult } from './runtime'; export { handleSteerRequest, handleSteerCancel, diff --git a/packages/api/src/agents/steering/runtime.ts b/packages/api/src/agents/steering/runtime.ts index 502c8f9308..44a0a06b46 100644 --- a/packages/api/src/agents/steering/runtime.ts +++ b/packages/api/src/agents/steering/runtime.ts @@ -12,6 +12,16 @@ import { GenerationJobManager } from '~/stream/GenerationJobManager'; import { getReferencedQuotes, mergeQuotedText } from '~/utils'; type SteerDrainOutput = HookOutputByEvent['PostToolBatch']; +export type TerminalSteerHookInput = Omit & { + hook_event_name: 'StopFinalize'; + continuationBudgetRemaining: number; + continuationPlanned: boolean; + continuationPrevented: boolean; +}; +export type TerminalSteerHook = ( + input: TerminalSteerHookInput, + signal: AbortSignal, +) => HookOutputByEvent['Stop'] | Promise; /** * Whether the installed `@librechat/agents` supports the FULL steering @@ -73,8 +83,8 @@ export interface SteerDrainHookOptions { } /** - * Shared drain body for BOTH injection boundaries (PostToolBatch and - * PreemptBoundary) — the provider-safety argument rests on the two sites + * Shared drain body for all injection boundaries (PostToolBatch, + * PreemptBoundary, and terminal Stop) — the provider-safety argument rests on each site * emitting identical `InjectedMessage` shapes, so they must share one body. * * Every part is durably applied before ANY slow media encoding begins. A @@ -87,7 +97,11 @@ export interface SteerDrainHookOptions { * preempt request for them must clear even on the failure path, or a * satisfied preempt could seal a later, unrelated stretch of generation. */ -async function drainAndBuildInjections(opts: SteerDrainHookOptions): Promise { +async function drainAndBuildInjections( + opts: SteerDrainHookOptions, + claim: () => Promise = () => + GenerationJobManager.steering.drain(opts.streamId, opts.jobCreatedAt), +): Promise { const { streamId, jobCreatedAt, applySteer, buildMedia } = opts; // The replacement guard lives INSIDE the store's atomic drain: a separate // check-then-drain could still consume a replacement job's queue if @@ -99,7 +113,7 @@ async function drainAndBuildInjections(opts: SteerDrainHookOptions): Promise => { + if (input.agentId != null) { + return { decision: 'continue' }; + } + const allowClaim = + input.continuationBudgetRemaining > 0 && + input.stopReason == null && + !input.continuationPrevented; + const injectedMessages = await drainAndBuildInjections(opts, async () => { + const admission = await GenerationJobManager.steering.admitTerminal( + opts.streamId, + { + allowClaim, + keepOpenWhenEmpty: allowClaim && input.continuationPlanned, + }, + opts.jobCreatedAt, + ); + return admission.outcome === 'claimed' ? admission.items : []; + }); + if (injectedMessages.length === 0) { + return { decision: 'continue' }; + } + return { decision: 'block', injectedMessages }; + }; +} + /** * The run's `RunConfig.preemption` — a level-triggered O(1) poll over the * job's armed preempt requests, exactly as the SDK contract requires: it @@ -283,3 +332,8 @@ export function isSteerPreemptSupported(): boolean { const sdk = agentsSdk as { HOOK_PREEMPT_BOUNDARY_CAPABLE?: boolean }; return isSteeringSupported() && sdk.HOOK_PREEMPT_BOUNDARY_CAPABLE === true; } + +export function isSteerTerminalContinuationSupported(): boolean { + const sdk = agentsSdk as { HOOK_STOP_CONTINUATION_CAPABLE?: boolean }; + return isSteeringSupported() && sdk.HOOK_STOP_CONTINUATION_CAPABLE === true; +} diff --git a/packages/api/src/stream/SteeringLifecycle.ts b/packages/api/src/stream/SteeringLifecycle.ts index 7ccde362a0..849b4d5f7e 100644 --- a/packages/api/src/stream/SteeringLifecycle.ts +++ b/packages/api/src/stream/SteeringLifecycle.ts @@ -11,6 +11,8 @@ import type { SteerQueueItem, SteerReceipt, SteerReceiptInput, + TerminalSteerAdmissionPolicy, + TerminalSteerAdmissionResult, } from '~/stream/interfaces/IJobStore'; import type { ServerSentEvent } from '~/types'; @@ -175,6 +177,14 @@ export class SteeringLifecycle { return this.store.restoreClaimedSteers(streamId, items, expectedCreatedAt); } + admitTerminal( + streamId: string, + policy: TerminalSteerAdmissionPolicy, + expectedCreatedAt?: number, + ): Promise { + return this.store.admitTerminalSteers(streamId, policy, expectedCreatedAt); + } + /** * Terminal drain: atomically CLOSE the queue to new steers, then take all * queued items. Finalization paths use this so a steer POST racing the diff --git a/packages/api/src/stream/__tests__/RedisJobStore.stream_integration.spec.ts b/packages/api/src/stream/__tests__/RedisJobStore.stream_integration.spec.ts index 3259a910cb..6805b63db1 100644 --- a/packages/api/src/stream/__tests__/RedisJobStore.stream_integration.spec.ts +++ b/packages/api/src/stream/__tests__/RedisJobStore.stream_integration.spec.ts @@ -3109,6 +3109,92 @@ describe('RedisJobStore Integration Tests', () => { await store.destroy(); }); + test('terminal claim-or-seal atomically assigns the final steer race', async () => { + if (!ioredisClient) { + return; + } + + const { RedisJobStore } = await import('../implementations/RedisJobStore'); + const { STEER_ENQUEUE_NOT_RUNNING } = await import('../interfaces/IJobStore'); + const store = new RedisJobStore(ioredisClient); + await store.initialize(); + + const claimStream = `steer-terminal-claim-${Date.now()}`; + const claimJob = await store.createJob(claimStream, 'steer-user', claimStream); + const first = buildSteer('terminal-first', 'first'); + const second = buildSteer('terminal-second', 'second'); + await store.enqueueSteer(claimStream, first, claimJob.createdAt); + await store.enqueueSteer(claimStream, second, claimJob.createdAt); + + await expect( + store.admitTerminalSteers( + claimStream, + { allowClaim: true, keepOpenWhenEmpty: false }, + claimJob.createdAt, + ), + ).resolves.toEqual({ outcome: 'claimed', items: [first, second] }); + await expect(store.peekClaimedSteers(claimStream, claimJob.createdAt)).resolves.toEqual([ + first, + second, + ]); + await expect( + store.enqueueSteer(claimStream, buildSteer('terminal-later', 'later'), claimJob.createdAt), + ).resolves.toBe(1); + + const openStream = `steer-terminal-open-${Date.now()}`; + const openJob = await store.createJob(openStream, 'steer-user', openStream); + await expect( + store.admitTerminalSteers( + openStream, + { allowClaim: true, keepOpenWhenEmpty: true }, + openJob.createdAt, + ), + ).resolves.toEqual({ outcome: 'open' }); + await expect( + store.enqueueSteer( + openStream, + buildSteer('terminal-planned', 'planned'), + openJob.createdAt, + ), + ).resolves.toBe(1); + + const sealStream = `steer-terminal-seal-${Date.now()}`; + const sealJob = await store.createJob(sealStream, 'steer-user', sealStream); + const queued = buildSteer('terminal-queued', 'ordinary follow-up'); + await store.enqueueSteer(sealStream, queued, sealJob.createdAt); + await expect( + store.admitTerminalSteers( + sealStream, + { allowClaim: false, keepOpenWhenEmpty: false }, + sealJob.createdAt, + ), + ).resolves.toEqual({ outcome: 'sealed' }); + await expect( + store.enqueueSteer(sealStream, buildSteer('terminal-raced', 'raced'), sealJob.createdAt), + ).resolves.toBe(STEER_ENQUEUE_NOT_RUNNING); + await expect(store.closeAndDrainSteers(sealStream, sealJob.createdAt)).resolves.toEqual([ + queued, + ]); + + const replacement = await store.createJob(sealStream, 'steer-user', sealStream); + await expect( + store.admitTerminalSteers( + sealStream, + { allowClaim: true, keepOpenWhenEmpty: false }, + sealJob.createdAt, + ), + ).resolves.toEqual({ outcome: 'unavailable' }); + await expect( + store.enqueueSteer( + sealStream, + buildSteer('terminal-replacement', 'replacement'), + replacement.createdAt, + ), + ).resolves.toBe(1); + + await store.destroy(); + }); + test('terminal CAS atomically returns and parks claimed plus queued steers', async () => { if (!ioredisClient) { return; diff --git a/packages/api/src/stream/__tests__/steering.spec.ts b/packages/api/src/stream/__tests__/steering.spec.ts index b8355781f6..b444d0248e 100644 --- a/packages/api/src/stream/__tests__/steering.spec.ts +++ b/packages/api/src/stream/__tests__/steering.spec.ts @@ -203,6 +203,111 @@ describe('SteeringLifecycle via GenerationJobManager.steering (in-memory)', () = }); }); + describe('terminal claim-or-seal', () => { + test('claims a v2 FIFO batch without closing admission', async () => { + const streamId = 'steer-terminal-claim'; + const job = await manager.createJob(streamId, 'user-1'); + const first = buildSteer('first'); + const second = buildSteer('second'); + await manager.steering.enqueue(streamId, first, job.createdAt); + await manager.steering.enqueue(streamId, second, job.createdAt); + + await expect( + manager.steering.admitTerminal( + streamId, + { allowClaim: true, keepOpenWhenEmpty: false }, + job.createdAt, + ), + ).resolves.toEqual({ outcome: 'claimed', items: [first, second] }); + await expect(jobStore.peekClaimedSteers(streamId, job.createdAt)).resolves.toEqual([ + first, + second, + ]); + await expect( + manager.steering.enqueue(streamId, buildSteer('later'), job.createdAt), + ).resolves.toBe(1); + }); + + test.each([ + { name: 'empty queue', queued: false, allowClaim: true }, + { name: 'exhausted budget', queued: true, allowClaim: false }, + ])('seals admission for $name', async ({ queued, allowClaim }) => { + const streamId = `steer-terminal-seal-${queued}-${allowClaim}`; + const job = await manager.createJob(streamId, 'user-1'); + const existing = buildSteer('existing'); + if (queued) { + await manager.steering.enqueue(streamId, existing, job.createdAt); + } + + await expect( + manager.steering.admitTerminal( + streamId, + { allowClaim, keepOpenWhenEmpty: false }, + job.createdAt, + ), + ).resolves.toEqual({ outcome: 'sealed' }); + await expect( + manager.steering.enqueue(streamId, buildSteer('raced'), job.createdAt), + ).resolves.toBe(STEER_ENQUEUE_NOT_RUNNING); + await expect(manager.steering.closeAndDrain(streamId, job.createdAt)).resolves.toEqual( + queued ? [existing] : [], + ); + }); + + test('refuses a stale generation without sealing its replacement', async () => { + const streamId = 'steer-terminal-stale'; + const stale = await manager.createJob(streamId, 'user-1'); + const replacement = await manager.createJob(streamId, 'user-1'); + + await expect( + manager.steering.admitTerminal( + streamId, + { allowClaim: true, keepOpenWhenEmpty: false }, + stale.createdAt, + ), + ).resolves.toEqual({ outcome: 'unavailable' }); + await expect( + manager.steering.enqueue(streamId, buildSteer('replacement'), replacement.createdAt), + ).resolves.toBe(1); + }); + + test('seals protocol v1 instead of claiming an unrecoverable batch', async () => { + const streamId = 'steer-terminal-v1'; + const job = await manager.createJob(streamId, 'user-1', undefined, { + initialMetadata: { generationProtocolVersion: 1 }, + }); + const queued = buildSteer('legacy'); + await manager.steering.enqueue(streamId, queued, job.createdAt); + + await expect( + manager.steering.admitTerminal( + streamId, + { allowClaim: true, keepOpenWhenEmpty: false }, + job.createdAt, + ), + ).resolves.toEqual({ outcome: 'sealed' }); + await expect(manager.steering.closeAndDrain(streamId, job.createdAt)).resolves.toEqual([ + queued, + ]); + }); + + test('keeps an empty queue open while another continuation is already planned', async () => { + const streamId = 'steer-terminal-open'; + const job = await manager.createJob(streamId, 'user-1'); + + await expect( + manager.steering.admitTerminal( + streamId, + { allowClaim: true, keepOpenWhenEmpty: true }, + job.createdAt, + ), + ).resolves.toEqual({ outcome: 'open' }); + await expect( + manager.steering.enqueue(streamId, buildSteer('planned continuation'), job.createdAt), + ).resolves.toBe(1); + }); + }); + describe('cancel', () => { test('removes exactly the cancelled steer and preserves queue order', async () => { const streamId = 'steer-cancel'; diff --git a/packages/api/src/stream/implementations/InMemoryJobStore.ts b/packages/api/src/stream/implementations/InMemoryJobStore.ts index 888c708113..605291bd40 100644 --- a/packages/api/src/stream/implementations/InMemoryJobStore.ts +++ b/packages/api/src/stream/implementations/InMemoryJobStore.ts @@ -9,6 +9,8 @@ import type { SteerArmResult, SteerEnqueueReceiptResult, SteerEnqueueVersionedResult, + TerminalSteerAdmissionPolicy, + TerminalSteerAdmissionResult, SteerQueueItem, SteerReceipt, SteerReceiptInput, @@ -1901,6 +1903,39 @@ export class InMemoryJobStore implements IJobStoreV2 { return true; } + async admitTerminalSteers( + streamId: string, + policy: TerminalSteerAdmissionPolicy, + expectedCreatedAt?: number, + ): Promise { + const job = this.jobs.get(streamId); + if ( + job?.status !== 'running' || + this.closedSteerQueues.has(streamId) || + (expectedCreatedAt != null && job.createdAt !== expectedCreatedAt) + ) { + return { outcome: 'unavailable' }; + } + const queue = this.steerQueues.get(streamId); + if (!policy.allowClaim || job.generationProtocolVersion !== 2) { + this.closedSteerQueues.add(streamId); + return { outcome: 'sealed' }; + } + if (queue == null || queue.length === 0) { + if (policy.keepOpenWhenEmpty) { + return { outcome: 'open' }; + } + this.closedSteerQueues.add(streamId); + return { outcome: 'sealed' }; + } + const items = await this.drainSteers(streamId, expectedCreatedAt); + if (items.length === 0) { + this.closedSteerQueues.add(streamId); + return { outcome: 'sealed' }; + } + return { outcome: 'claimed', items }; + } + async closeAndDrainSteers( streamId: string, expectedCreatedAt?: number, diff --git a/packages/api/src/stream/implementations/RedisJobStore.ts b/packages/api/src/stream/implementations/RedisJobStore.ts index 7f2ee875d5..eef54ff58a 100644 --- a/packages/api/src/stream/implementations/RedisJobStore.ts +++ b/packages/api/src/stream/implementations/RedisJobStore.ts @@ -20,6 +20,8 @@ import type { SteerArmResult, SteerEnqueueReceiptResult, SteerEnqueueVersionedResult, + TerminalSteerAdmissionPolicy, + TerminalSteerAdmissionResult, SteerReceipt, SteerReceiptInput, ParkedSteerClaim, @@ -1115,6 +1117,42 @@ const STEER_DRAIN_LUA = 'redis.call("DEL", KEYS[2]) ' + 'return items'; +/** Atomic terminal choice: claim the current v2 FIFO batch, or close steer + * admission when the queue is empty / the continuation budget is exhausted. + * The first return entry is the outcome; claimed item JSON follows it. */ +const STEER_TERMINAL_ADMISSION_LUA = + 'if ARGV[1] ~= "" and redis.call("HGET", KEYS[1], "createdAt") ~= ARGV[1] then return { "unavailable" } end ' + + 'if redis.call("HGET", KEYS[1], "status") ~= "running" ' + + 'or redis.call("HGET", KEYS[1], "steersClosed") == "1" then return { "unavailable" } end ' + + 'local items = redis.call("LRANGE", KEYS[2], 0, -1) ' + + 'if ARGV[3] ~= "1" or redis.call("HGET", KEYS[1], "generationProtocolVersion") ~= "2" then ' + + 'redis.call("HSET", KEYS[1], "steersClosed", "1") return { "sealed" } end ' + + 'if #items == 0 and ARGV[4] == "1" then return { "open" } end ' + + 'if #items == 0 then redis.call("HSET", KEYS[1], "steersClosed", "1") return { "sealed" } end ' + + 'local currentCreatedAt = redis.call("HGET", KEYS[1], "createdAt") ' + + 'local decodedItems = {} local decodedReceipts = {} ' + + 'for i = 1, #items do local ok, item = pcall(cjson.decode, items[i]) ' + + 'if not ok or type(item) ~= "table" or not item.steerId then ' + + 'return redis.error_reply("invalid steer queue item") end decodedItems[i] = item ' + + 'if item.clientSteerId then local raw = redis.call("HGET", KEYS[4], item.clientSteerId) ' + + 'if not raw then return redis.error_reply("missing steer receipt") end ' + + 'local receiptOk, receipt = pcall(cjson.decode, raw) ' + + 'if not receiptOk or type(receipt) ~= "table" ' + + 'or receipt.clientSteerId ~= item.clientSteerId ' + + 'or tostring(receipt.generationCreatedAt or "") ~= currentCreatedAt ' + + 'or receipt.state ~= "queued" or not receipt.item ' + + 'or receipt.item.steerId ~= item.steerId ' + + 'or receipt.item.clientSteerId ~= item.clientSteerId then ' + + 'return redis.error_reply("invalid steer receipt") end decodedReceipts[i] = receipt end end ' + + 'redis.call("RPUSH", KEYS[3], unpack(items)) redis.call("EXPIRE", KEYS[3], ARGV[2]) ' + + 'for i = 1, #items do local item = decodedItems[i] local receipt = decodedReceipts[i] ' + + 'if receipt then receipt.item = item receipt.state = "claimed" ' + + 'redis.call("HSET", KEYS[4], item.clientSteerId, cjson.encode(receipt)) end end ' + + 'for i = 4, 5 do local ttl = redis.call("TTL", KEYS[i]) ' + + 'if ttl >= 0 and ttl < tonumber(ARGV[2]) then redis.call("EXPIRE", KEYS[i], ARGV[2]) end end ' + + 'redis.call("DEL", KEYS[2]) local result = { "claimed" } ' + + 'for i = 1, #items do result[#result + 1] = items[i] end return result'; + /** Roll back claimed items whose durable applied-part write failed. New * enqueues may have landed after the drain, so the failed accepted batch is * prepended in its original order rather than replacing the live queue. */ @@ -4144,6 +4182,35 @@ export class RedisJobStore implements IJobStoreV2 { return this.parseSteerItems(raw); } + async admitTerminalSteers( + streamId: string, + policy: TerminalSteerAdmissionPolicy, + expectedCreatedAt?: number, + ): Promise { + const raw = await this.redis.eval( + STEER_TERMINAL_ADMISSION_LUA, + 5, + KEYS.job(streamId), + KEYS.steers(streamId), + KEYS.claimedSteers(streamId), + KEYS.steerReceipts(streamId), + KEYS.steerReceiptOrder(streamId), + expectedCreatedAt != null ? String(expectedCreatedAt) : '', + String(this.runningStorageTtlSeconds()), + policy.allowClaim ? '1' : '0', + policy.keepOpenWhenEmpty ? '1' : '0', + ); + if (!Array.isArray(raw) || typeof raw[0] !== 'string') { + return { outcome: 'unavailable' }; + } + if (raw[0] !== 'claimed') { + return { + outcome: raw[0] === 'sealed' || raw[0] === 'open' ? raw[0] : 'unavailable', + }; + } + return { outcome: 'claimed', items: this.parseSteerItems(raw.slice(1)) }; + } + async restoreClaimedSteers( streamId: string, items: SteerQueueItem[], diff --git a/packages/api/src/stream/interfaces/IJobStore.ts b/packages/api/src/stream/interfaces/IJobStore.ts index b56a259d30..2013fcf9ff 100644 --- a/packages/api/src/stream/interfaces/IJobStore.ts +++ b/packages/api/src/stream/interfaces/IJobStore.ts @@ -562,6 +562,15 @@ export interface SteerEnqueueResult { position: number; } +export type TerminalSteerAdmissionResult = + | { outcome: 'claimed'; items: SteerQueueItem[] } + | { outcome: 'open' | 'sealed' | 'unavailable' }; + +export interface TerminalSteerAdmissionPolicy { + allowClaim: boolean; + keepOpenWhenEmpty: boolean; +} + export type SteerEnqueueVersionedResult = SteerEnqueueResult | number; /** @@ -1388,6 +1397,21 @@ export interface IJobStoreV2 extends IJobStore { expectedCreatedAt?: number, ): Promise; + /** + * Terminal admission fence. When `allowClaim` and queued work are both + * present, atomically claim the FIFO batch while leaving admission open for + * the continued run. Otherwise atomically close admission so a racing steer + * is rejected and remains an ordinary follow-up, unless + * `keepOpenWhenEmpty` proves another folded Stop hook already planned a + * continuation. V1 generations always seal because they lack + * crash-recoverable claimed-steer receipts. + */ + admitTerminalSteers( + streamId: string, + policy: TerminalSteerAdmissionPolicy, + expectedCreatedAt?: number, + ): Promise; + /** * Atomically CLOSE the queue to new steers, then take all queued items * FIFO. Used by the terminal paths (final event, abort) so a steer POST diff --git a/packages/api/src/stream/jobStoreCapabilities.ts b/packages/api/src/stream/jobStoreCapabilities.ts index 2df4cb6bed..92a6f915cb 100644 --- a/packages/api/src/stream/jobStoreCapabilities.ts +++ b/packages/api/src/stream/jobStoreCapabilities.ts @@ -22,6 +22,7 @@ export const JOB_STORE_V2_REQUIRED_METHODS = [ 'enqueueSteerWithReceipt', 'getSteerReceipt', 'restoreClaimedSteers', + 'admitTerminalSteers', 'peekClaimedSteers', 'armSteer', 'armSteerVersioned',