LibreChat/api/server/routes/agents/index.js
Danny Avila e7f1838515
Some checks are pending
Docker Dev Branch Images Build / build (Dockerfile, lc-dev, node) (push) Waiting to run
Docker Dev Branch Images Build / build (Dockerfile.multi, lc-dev-api, api-build) (push) Waiting to run
GitNexus Index / index (push) Waiting to run
GitNexus Index / post-index (push) Blocked by required conditions
feat: Reliable Interrupt & Steer Escalation and Recovery (#14558)
* feat: surface interrupt-steer escalation on waiting messages

The interrupt & steer feature shipped reachable only through the
composer chord, the send-button hovercard, and the composer button; a
message already waiting (queued for after the run, or steered and
parked at the next tool boundary) had no path to it. Both waiting
surfaces now carry one:

- Queued rows get an icon-only ZapOff escalation button beside the
  existing Steer primary. It routes through sendQueuedNow, which now
  takes a preempt option on its live-run path. The tooltip teaches the
  composer chord, derived through resolveComposerKeyDown so a rebound
  or yielded chord is never advertised.
- In-flight steer bubbles get an "Interrupt now" overflow entry with
  the same race rules as Edit: reclaim first, and only a `reclaimed`
  outcome resubmits (via retrySteer with preempt, swapping the chip
  for an interrupting one). `applied` and run-ended-mid-reclaim
  outcomes stop at the existing informational toasts, so the words can
  never land twice. Not offered on a steer already preempting.
- Every during-run overflow menu gains an "Always interrupt instead"
  toggle for steerInterruptsByDefault, next to the existing steer/queue
  default toggle. MenuEntry supports disabled for the new entries.

Only one interrupt can be unresolved at a time: while one preempt is
pending (or the run is paused on approval, where the server 409s),
every escalation control disables instead of racing the same seal.

Ten new tests across both surfaces; 381 green in the affected suites.

* fix: lock escalation across its reclaim window, keep the paused control visible, label as steer

Codex round 1, all three findings.

P2, escalation race. The single-interrupt invariant had a window between
clicking "Interrupt now" and the reclaim resolving, where no preempt
chip existed for the chip-derived gate to see: two bubbles escalated
back-to-back could both resubmit. A shared escalating flag (Jotai,
per-conversation) now covers the window and disables every escalation
control on both surfaces, and a fresh recheck before resubmitting
catches an interrupt armed elsewhere meanwhile (composer chord, queued
row); those words re-home to the queue with an informational toast
instead of breaking the invariant.

P2, unreachable paused state. canSteer is defined as
hasRealConvoId && !pausedOnApproval, so gating the button on canSteer
removed it exactly when it was meant to render disabled; the test only
passed on an impossible stub combination. The render gate is now
duringRunActive && (canSteer || pausedOnApproval), and the test uses the
real invariant.

P2, label semantics. "Interrupt & send now" borrowed the name of the
hard-abort action; this one preserves the partial answer and steers.
Renamed to "Interrupt & steer now" (com_ui_interrupt_steer_now).

Both behavior fixes counterfactually verified; 384 tests green across
the affected suites.

* fix: disable bubble escalation while the run cannot accept a steer

Codex round 2, one P2. Answer mode (ask_user_question) sets
duringRunActive false while pausedOnApproval stays false, since that
flag only detects approval-bearing tool calls. The bubble's escalation
entry stayed enabled there, so clicking it cancelled a healthy waiting
steer and the preempt resubmission bounced off RUN_PAUSED, degrading
the words to the queue. The entry now also disables on
!duringRunActive, matching the queued-row control's gate.

Counterfactually verified: reverting the gate fails the new
answer-mode test.

* fix: recheck live run state after the reclaim, not just at the click

Codex round 3, one P2, and it is the round-1 recheck principle applied
one level deeper: the entry-time disable cannot see a run that pauses
(tool approval, answer mode) while the reclaim round-trip is in flight,
and the .then closure held the render's stale steering controls, so the
resubmit would fire into a RUN_PAUSED rejection after the reclaim had
already surrendered the steer's boundary slot.

The escalation continuation now reads the LIVE controls through a
latest-ref: if the run can no longer accept a steer, the words re-home
to the queue with an informational toast instead of resubmitting, and
the resubmit itself also goes through the live controls.

Counterfactually verified: reading the stale closure instead of the ref
fails the new mid-reclaim pause test.

* refactor: make escalation one atomic server-side arm, in place

Codex round 4: four P2s, every one an interleaving of the same window —
escalation as reclaim-then-repost is a compound, non-atomic operation
whose continuation must revalidate the world (FIFO position lost, ref
assigned too late, no run fence, competing bubble actions). Rounds 1-3
patched that window with a lock and rechecks; round 4 shows the window
itself is the defect, so this removes it instead of guarding it again.

Escalation is now POST /chat/steer/arm: the server flips preempt on the
EXISTING queued item in one atomic store op (new IJobStore.armSteer; a
decode-patch-encode LSET Lua on Redis, an in-place mutation in memory),
fenced to the validated generation and refused once the queue closes.
The handler mirrors the steer POST's preempt contract exactly: durable
flag gated on the owner's recorded capability, volatile requestPreempt
fire-and-forget because the durable flag is the truth resume/handover
re-arm from.

By construction this resolves all four findings: FIFO survives (the
item never moves; the whole queue still drains in instruction order at
the seal), there is no continuation to hold stale controls, the store
op is fenced to the original run, and a competing Edit/Queue/Cancel
either beats the arm (armed:false, chip untouched) or operates on the
armed item, whose cancel already disarms.

The client escalation entry becomes one mutation: armed:true relabels
the chip in place (same steerId, same position), PREEMPT_UNSUPPORTED
and lost races toast honestly, and the round 1-3 machinery — the
escalating lock atom, the latest-ref, the post-reclaim rechecks and
their two toast strings — is deleted rather than extended.

Verified: 7 new handler tests on the real in-memory manager (including
FIFO preservation and the stale-generation fence), 2 Redis integration
tests against real Redis (in-place arm keeps order and every field;
missing/stale/closed all refuse), client suites 396 green.

* fix: decide capability inside the atomic arm, neutralize the lost-race toast

Codex round 5, both findings, both edges of the new arm design rather
than its mechanism.

P2, capability TOCTOU. A HITL resume on a rolling deploy rewrites
preemptCapable for the SAME generation, so the handler's read could go
stale between validation and the flag flip, arming a steer the live
owner cannot seal. armSteer now returns armed | missing | incapable,
with the owner's live capability part of the same atomic predicate as
the generation fence (HGET preemptCapable inside the Lua; the flat job
field, not a metadata blob — the in-memory store reads the same field).
The handler's pre-check is deleted rather than kept alongside; the
store predicate is the single source. New handler test rewrites the
capability after queueing and expects PREEMPT_UNSUPPORTED with the item
left unflagged; the Redis guards test now asserts the incapable refusal
against real Redis.

P2, ambiguous toast. armed:false covers injected, cancelled, re-homed,
and run-over alike, so telling the user the message "already reached
the agent" claimed one specific outcome. The lost-race branch now uses
a neutral message (com_ui_steer_arm_lost_race) and defers to the events
for what actually happened.

* fix: flip the escalation lock synchronously before the arm request

Codex round 6, one P2. Round 4 deleted the escalating flag along with
the reclaim continuation it guarded, but that left the one-interrupt
gate blind during the arm request's own round trip: the chip-derived
check cannot see an arm until its response relabels the chip, so on a
slow connection two bubbles could both arm before either response
landed. Double-arm is harmless server-side now (the run seals once and
drains the whole queue in order), but every escalation control
advertises "one interrupt at a time" by disabling, and the controls
must tell the truth.

The per-conversation escalating flag returns as a pure UX gate: set
synchronously at click, before the mutation, cleared on settlement, and
folded into interruptPending on both surfaces. Unlike its round 1-3
ancestor there is no continuation behind it to guard and no recheck to
pair with it.

Counterfactually verified: without the synchronous set, the two-bubble
race test arms twice. 207 tests green across the Chat Input suites.

* test(e2e): cover escalation of waiting messages through the real seal

Three mock-harness tests on E2E_SLOW_REPLY, a 160-chunk stream with no
tool boundary, so an in-thread steer part can ONLY come from a genuine
mid-stream seal — which makes each test a behavioral proof rather than
a UI check:

- Queued row escalation: the ZapOff button turns a waiting queued
  message into a preempt-armed steer (202 echoes preempt: true) that
  seals and injects, where the sibling steering.spec test proves the
  unescalated path waits for run end instead.
- Bubble in-place arm: an ordinary steer (202 with no preempt echo)
  waits as a bubble, POST /chat/steer/arm answers armed: true, the
  bubble relabels in place (same single bubble, same text, escalation
  no longer offered on reopen), and the armed steer seals mid-stream.
- Always-interrupt toggle: flipped from a waiting row's overflow menu,
  plain Enter now produces a preempt: true steer that seals in the SAME
  run, and the menu offers the way back. An afterEach clears the
  localStorage preference so a mid-test failure cannot leak
  preempt-by-default into the rest of the serial suite.

All three verified locally through the full harness (real backend, mock
LLM, seeded DB): 3 passed in 27s.

* feat: dedicated escalation arrow + shortcut, menu split into actions and preferences

The escalation was still half-hidden: the bubble only offered it inside
the overflow menu, and the tooltip taught the composer chord, which does
a different thing (interrupts with typed text, not this chip). Three
changes make it a first-class command:

- A shared EscalateNowButton (circular arrow, ghost-bordered like the
  composer's interrupt control) is always visible on BOTH surfaces:
  beside each queued row's Steer primary and on every waiting steer
  bubble next to its menu. It disappears once a steer is interrupting.
- A dedicated registry shortcut, escalateSteer (Cmd/Ctrl+Shift+.),
  editing-allowed and rebindable like every other action. Deliberately
  NOT an Enter chord: the composer owns every Enter chord, and the
  yield design rests on no default binding using Enter besides submit.
  Its handler clicks the newest enabled arrow control (bubbles beat
  queued rows), so the shortcut can never diverge from the button, and
  the arrow's tooltip teaches THIS command via the registry display.
- The overflow menus separate one-off actions from sticky behavior
  changes: Edit, Cancel, Queue, then a smaller "Preferences" section
  holding the queueing and always-interrupt toggles, each with the
  standard InfoHoverCard reusing the Settings panel's descriptions.
  "Interrupt & steer now" leaves the menu entirely.

386 client tests green, including a menu-structure test locking the
order and the absence of the escalation entry; bubble escalation tests
drive the visible arrow. The e2e spec's bubble test now clicks the
arrow, and a fourth test drives the dedicated shortcut end to end
through a real mid-stream seal.

* style: bind the escalation arrow to its message (variant A anatomy)

Two same-weight circles in a row read as one control group, leaving the
arrow's ownership ambiguous, and a floating arrow stops meaning anything
once several messages stack. The shared control now carries variant A's
anatomy: a thin divider binds a small SOLID arrow (filled, inverted) to
the message region on its left, and the menu ellipsis stays a bare
glyph, so the two affordances can no longer blur together — and the
divider+arrow pairing repeats cleanly per chip at N messages.

* chore: drop the unused within import CI lint caught

* fix: advertise the escalation shortcut only while the control is live

Codex on the e2e head, one P2: the tooltip appended the chord hint even
while the button was disabled, advertising a shortcut that does nothing
during an approval pause. The flagged control (InterruptNowButton) was
since replaced by the shared EscalateNowButton, which inherited the
pattern; the successor now omits the chord whenever the control is
disabled, matching the rule the during-run hovercard already follows.

* fix: harden steer escalation lifecycle and recovery

* test(e2e): disambiguate accessible steer preferences

* test: align abort persistence coverage with prerequisites

* chore(i18n): remove obsolete steer race message

* chore: normalize imports across steering changes

* test: exercise stream integration on Redis Cluster

* test: scope HITL checkpoints to generation

* test: fix cluster cleanup and locale policy

* fix: keep escalation visible during ask pauses

* fix: fence recovery downgrade and stale predecessors

* fix: require generation owner abort acknowledgement

* fix: validate delayed preempt arms

* test: align final escalation fixtures

* fix: preserve in-memory predecessor abort handoff

* fix: restore controls for recovered queued messages

* test: cover recovered queue controls

* fix: close final steering review gaps
2026-07-31 20:07:56 -04:00

942 lines
36 KiB
JavaScript

const express = require('express');
const {
isEnabled,
GenerationJobManager,
TERMINAL_PUBLICATION_RECONNECT_ERROR,
hasPersistableAbortContent,
buildAbortedResponseMetadata,
isPendingActionStale,
toClientPendingAction,
isHITLEnabled,
captureAgentCheckpointGeneration,
deleteAgentCheckpoint,
attachAskUserQuestionArgs,
createMessageFilterPii,
} = require('@librechat/api');
const { createSseStreamTelemetry } = require('@librechat/api/telemetry');
const { logger } = require('@librechat/data-schemas');
const {
uaParser,
checkBan,
moderateText,
requireJwtAuth,
messageIpLimiter,
configMiddleware,
messageUserLimiter,
} = require('~/server/middleware');
const SteerController = require('~/server/controllers/agents/steer');
const {
GENERATION_PROTOCOL_HEADER,
GENERATION_PROTOCOL_V2,
getRequestedGenerationProtocol,
getServerGenerationProtocol,
negotiateExistingGenerationProtocol,
} = require('~/server/controllers/agents/protocol');
const { saveMessage } = require('~/models');
const responses = require('./responses');
const openai = require('./openai');
const { v1 } = require('./v1');
const chat = require('./chat');
const { LIMIT_MESSAGE_IP, LIMIT_MESSAGE_USER } = process.env ?? {};
/** Untenanted jobs (pre-multi-tenancy) remain accessible if the userId check passes. */
function hasTenantMismatch(job, user) {
return job.metadata?.tenantId != null && job.metadata.tenantId !== user.tenantId;
}
/** Protocol selected before a job has been authorized/read. This is used for
* validation, not-found, and authorization envelopes; it never leaks an
* existing job's marker to an unauthorized caller. */
function negotiateRequestGenerationProtocol(req) {
return Math.min(
getRequestedGenerationProtocol(req),
getServerGenerationProtocol(GenerationJobManager),
);
}
/** Every generation-control JSON envelope carries the exact numeric protocol
* that governs it. The response header is useful to fetch/Axios callers, while
* the body survives auth-refresh adapters and is the client's fail-closed
* source of truth. */
function sendGenerationJson(res, status, body, generationProtocolVersion) {
res.set(GENERATION_PROTOCOL_HEADER, String(generationProtocolVersion));
return res.status(status).json({ ...body, generationProtocolVersion });
}
async function sendJoblessStatus(req, res, conversationId) {
// The default completeJob path deletes the job record immediately, so the
// jobless branch IS the common reload-after-terminal case — parked steers
// live under their own bounded-TTL key and authorize from their stored owner.
const requestedProtocolVersion = getRequestedGenerationProtocol(req);
const claimed = await GenerationJobManager.steering.claimDetailed(
conversationId,
{
userId: req.user.id,
tenantId: req.user.tenantId,
},
requestedProtocolVersion,
);
const generationProtocolVersion = Math.min(
requestedProtocolVersion,
claimed.steers.length > 0
? claimed.generationProtocolVersion
: getServerGenerationProtocol(GenerationJobManager),
);
res.set(GENERATION_PROTOCOL_HEADER, String(generationProtocolVersion));
return res.json({
active: false,
generationProtocolVersion,
...(claimed.steers.length > 0 && { unrecoveredSteers: claimed.steers }),
});
}
const router = express.Router();
/**
* Open Responses API routes (API key authentication handled in route file)
* Mounted at /agents/v1/responses (full path: /api/agents/v1/responses)
* NOTE: Must be mounted BEFORE /v1 to avoid being caught by the less specific route
* @see https://openresponses.org/specification
*/
router.use('/v1/responses', responses);
/**
* OpenAI-compatible API routes (API key authentication handled in route file)
* Mounted at /agents/v1 (full path: /api/agents/v1/chat/completions)
*/
router.use('/v1', openai);
router.use(requireJwtAuth);
router.use(checkBan);
router.use(uaParser);
/**
* Stream endpoints - mounted before chatRouter to bypass rate limiters
* These are GET requests and don't need message body validation or rate limiting
*/
/**
* @route GET /chat/stream/:streamId
* @desc Subscribe to an ongoing generation job's SSE stream with replay support
* @access Private
* @description Sends sync event with resume state, replays missed chunks, then streams live
* @query resume=true - Indicates this is a reconnection (sends sync event)
*/
router.get('/chat/stream/:streamId', async (req, res) => {
const { streamId } = req.params;
const isResume = req.query.resume === 'true';
const requestProtocolVersion = negotiateRequestGenerationProtocol(req);
const rawGenerationCreatedAt = req.query.generationCreatedAt;
let expectedGenerationCreatedAt;
if (rawGenerationCreatedAt != null) {
if (
typeof rawGenerationCreatedAt !== 'string' ||
!/^\d+$/.test(rawGenerationCreatedAt) ||
!Number.isSafeInteger(Number(rawGenerationCreatedAt))
) {
return sendGenerationJson(
res,
400,
{ error: 'Invalid generation identity' },
requestProtocolVersion,
);
}
expectedGenerationCreatedAt = Number(rawGenerationCreatedAt);
}
let result;
const attachmentAbortController = new AbortController();
req.on('close', () => {
logger.debug(`[AgentStream] Client disconnected from ${streamId}`);
attachmentAbortController.abort();
result?.unsubscribe();
});
const job = await GenerationJobManager.getJob(streamId);
if (attachmentAbortController.signal.aborted) {
return;
}
if (!job) {
return sendGenerationJson(
res,
404,
{
error: 'Stream not found',
message: 'The generation job does not exist or has expired.',
},
requestProtocolVersion,
);
}
// Every job has an owner at creation time. Treat a missing/corrupt owner as
// unauthorized instead of turning malformed store state into a public
// stream for anyone who knows the conversation id.
if (job.metadata?.userId !== req.user.id) {
return sendGenerationJson(res, 403, { error: 'Unauthorized' }, requestProtocolVersion);
}
if (hasTenantMismatch(job, req.user)) {
return sendGenerationJson(res, 403, { error: 'Unauthorized' }, requestProtocolVersion);
}
const generationProtocolVersion = negotiateExistingGenerationProtocol(req, job);
if (expectedGenerationCreatedAt != null && job.createdAt !== expectedGenerationCreatedAt) {
// streamId is conversation-scoped and may now belong to a newer turn. A
// stale start/reconnect gets a dedicated handoff signal instead of either
// following the ordinary terminal path or receiving replacement content.
return sendGenerationJson(
res,
409,
{
code: 'GENERATION_REPLACED',
error: 'Generation replaced',
message: 'The requested generation has completed or was replaced.',
},
generationProtocolVersion,
);
}
/** Pin even legacy (unfenced-query) subscribers to the exact job snapshot
* that passed the owner + tenant checks above. `streamId` is conversation-
* scoped, so a replacement can otherwise land between this authorization
* read and the manager attachment and expose the replacement generation
* without ever authorizing its owner. */
const authorizedGenerationCreatedAt = job.createdAt;
if (!Number.isSafeInteger(authorizedGenerationCreatedAt) || authorizedGenerationCreatedAt < 0) {
logger.warn(`[AgentStream] Refusing stream with invalid generation identity: ${streamId}`);
return sendGenerationJson(res, 403, { error: 'Unauthorized' }, requestProtocolVersion);
}
const streamTelemetry = createSseStreamTelemetry({ req, res, streamId, isResume });
res.setHeader('Content-Encoding', 'identity');
res.setHeader('Content-Type', 'text/event-stream');
res.setHeader('Cache-Control', 'no-cache, no-transform');
res.setHeader('Connection', 'keep-alive');
res.setHeader('X-Accel-Buffering', 'no');
res.setHeader(GENERATION_PROTOCOL_HEADER, String(generationProtocolVersion));
res.flushHeaders();
streamTelemetry.recordHeadersFlushed();
logger.debug(`[AgentStream] Client subscribed to ${streamId}, resume: ${isResume}`);
const writeEvent = (event, options = {}) => {
if (generationProtocolVersion < GENERATION_PROTOCOL_V2 && event?.event === 'on_steer_updated') {
return true;
}
if (!res.writableEnded) {
const eventName = options.eventName ?? 'message';
const payload = `event: ${eventName}\ndata: ${JSON.stringify(event)}\n\n`;
res.write(payload);
streamTelemetry.recordWrite(payload, { final: options.final });
if (typeof res.flush === 'function') {
res.flush();
}
return true;
}
return false;
};
const onDone = (event) => {
streamTelemetry.recordFinalEventEmitted();
if (event?.reconcile === true && generationProtocolVersion < GENERATION_PROTOCOL_V2) {
/** Legacy clients treat an ordinary `final: true` as the completion of
* their optimistic submission. A reconciliation frame has no response
* payload and may describe a replacement, so expose it only as a
* transport error; the v1 reconnect/status path will refetch safely. */
writeEvent(
{
error: 'Generation state changed; reconnect to load the saved response.',
generationProtocolVersion,
},
{ eventName: 'error', final: true },
);
res.end();
return;
}
writeEvent(
event != null && typeof event === 'object'
? { ...event, generationProtocolVersion }
: { final: true, generationProtocolVersion },
{ final: true },
);
res.end();
};
const onError = (error) => {
if (!res.writableEnded) {
streamTelemetry.recordErrorEventEmitted();
if (error === TERMINAL_PUBLICATION_RECONNECT_ERROR) {
/** A durable terminal payload exists, but cross-replica DONE publish
* failed. Tear down the HTTP stream without an application error frame:
* sse.js treats the transport close as reconnectable, and the retained
* terminal job then replays its authoritative final payload. */
res.destroy();
return;
}
writeEvent({ error, generationProtocolVersion }, { eventName: 'error' });
res.end();
}
};
if (isResume) {
const { subscription, resumeState, pendingEvents } =
await GenerationJobManager.subscribeWithResume(streamId, writeEvent, onDone, onError, {
signal: attachmentAbortController.signal,
expectedCreatedAt: authorizedGenerationCreatedAt,
});
if (subscription && !attachmentAbortController.signal.aborted && !res.writableEnded) {
if (resumeState) {
writeEvent({ sync: true, resumeState, pendingEvents });
GenerationJobManager.markSyncSent(streamId, authorizedGenerationCreatedAt);
logger.debug(
`[AgentStream] Sent sync event for ${streamId} with ${resumeState.runSteps.length} run steps, ${pendingEvents.length} pending events`,
);
} else if (pendingEvents.length > 0) {
for (const event of pendingEvents) {
writeEvent(event);
}
logger.warn(
`[AgentStream] Resume state null for ${streamId}, replayed ${pendingEvents.length} gap events directly`,
);
}
subscription.activate();
} else {
subscription?.unsubscribe();
}
result = subscription;
} else {
result = await GenerationJobManager.subscribe(streamId, writeEvent, onDone, onError, {
signal: attachmentAbortController.signal,
expectedCreatedAt: authorizedGenerationCreatedAt,
});
}
if (attachmentAbortController.signal.aborted) {
result?.unsubscribe();
return;
}
if (!result) {
streamTelemetry.recordSubscribeFailed();
{
let currentJob;
let currentJobReadSucceeded = false;
try {
currentJob = await GenerationJobManager.getJob(streamId);
currentJobReadSucceeded = true;
} catch (error) {
logger.warn(`[AgentStream] Failed to reconcile fenced subscription for ${streamId}`, error);
}
if (attachmentAbortController.signal.aborted || res.writableEnded) {
return;
}
const currentJobAuthorized =
currentJobReadSucceeded &&
(!currentJob ||
(currentJob.metadata?.userId === req.user.id &&
!hasTenantMismatch(currentJob, req.user)));
const generationReplaced =
currentJobAuthorized &&
currentJob != null &&
currentJob.createdAt !== authorizedGenerationCreatedAt;
const expectedGenerationTerminal =
currentJobAuthorized &&
currentJob?.createdAt === authorizedGenerationCreatedAt &&
['complete', 'error', 'aborted'].includes(currentJob.status);
/** The route already flushed SSE headers before the manager's final
* generation fence ran. A generic error here would misreport the common
* snapshot-to-attach race where the requested run terminalized or was
* replaced. Send the same control-only reconciliation frame used by the
* manager so the client refetches authoritative state instead. */
if (
currentJobReadSucceeded &&
(generationReplaced || !currentJob || expectedGenerationTerminal)
) {
onDone({
final: true,
reconcile: true,
reconcileReason: generationReplaced ? 'generation_replaced' : 'terminal_payload_missing',
...(expectedGenerationTerminal && { terminalStatus: currentJob.status }),
generationCreatedAt: authorizedGenerationCreatedAt,
conversation: {
conversationId: currentJob?.conversationId ?? job.conversationId ?? streamId,
},
});
return;
}
}
onError('Failed to subscribe to stream');
return;
}
});
/**
* @route GET /chat/active
* @desc Get all active generation job IDs for the current user
* @access Private
* @returns { activeJobIds: string[] }
*/
router.get('/chat/active', async (req, res) => {
const activeJobIds = await GenerationJobManager.getActiveJobIdsForUser(
req.user.id,
req.user.tenantId,
);
res.json({ activeJobIds });
});
/**
* @route GET /chat/status/:conversationId
* @desc Check if there's an active generation job for a conversation
* @access Private
* @returns { active, streamId, status, aggregatedContent, createdAt, resumeState }
*/
router.get('/chat/status/:conversationId', async (req, res) => {
const { conversationId } = req.params;
const requestProtocolVersion = negotiateRequestGenerationProtocol(req);
// streamId === conversationId, so we can use getJob directly
let job = await GenerationJobManager.getJob(conversationId);
if (!job) {
return sendJoblessStatus(req, res, conversationId);
}
let resumeState;
let snapshotVerified = false;
/** `getResumeState` begins with its own streamId lookup. A replacement can
* land after this route authorizes A but before that lookup and make it read
* B's content. Verify the epoch after each read and discard mismatched
* snapshots; every replacement snapshot is re-authorized before use. */
for (let attempt = 0; attempt < 3; attempt++) {
if (job.metadata?.userId !== req.user.id || hasTenantMismatch(job, req.user)) {
return sendGenerationJson(res, 403, { error: 'Unauthorized' }, requestProtocolVersion);
}
if (!Number.isSafeInteger(job.createdAt) || job.createdAt < 0) {
return sendGenerationJson(res, 403, { error: 'Unauthorized' }, requestProtocolVersion);
}
const authorizedCreatedAt = job.createdAt;
resumeState = await GenerationJobManager.getResumeState(conversationId, authorizedCreatedAt);
const verifiedJob = await GenerationJobManager.getJob(conversationId);
if (!verifiedJob) {
return sendJoblessStatus(req, res, conversationId);
}
if (verifiedJob.createdAt === authorizedCreatedAt) {
if (
verifiedJob.metadata?.userId !== req.user.id ||
hasTenantMismatch(verifiedJob, req.user)
) {
return sendGenerationJson(res, 403, { error: 'Unauthorized' }, requestProtocolVersion);
}
job = verifiedJob;
snapshotVerified = true;
break;
}
job = verifiedJob;
}
if (!snapshotVerified) {
res.set('Retry-After', '1');
return sendGenerationJson(res, 503, { code: 'SERVER_NOT_READY' }, requestProtocolVersion);
}
/** Abort has won terminal ownership, but its required message/checkpoint
* persistence has not finished yet. Reporting this snapshot as inactive
* would let a reloading client clear its live state and refetch history
* before the terminal owner has made that history authoritative. `getJob`
* recovers a stale pending marker; while the verified marker remains live,
* keep every status consumer on the same readiness path as duplicate starts. */
if (job.metadata?.terminalPersistencePending === true) {
res.set('Retry-After', '1');
return sendGenerationJson(
res,
503,
{ code: 'SERVER_NOT_READY' },
negotiateExistingGenerationProtocol(req, job),
);
}
let generationProtocolVersion = negotiateExistingGenerationProtocol(req, job);
res.set(GENERATION_PROTOCOL_HEADER, String(generationProtocolVersion));
// A job paused for human review is still active (consistent with /chat/active),
// so the client resumes/subscribes rather than treating it as finished — but
// only while it has a live, resolvable prompt: a missing/malformed or
// past-expiry pendingAction reads as inactive (cleanup/expiry will finalize it).
const pendingAction = job.metadata.pendingAction;
const pendingLive = job.status === 'requires_action' && !isPendingActionStale({ pendingAction });
const isActive = job.status === 'running' || pendingLive;
/** Acknowledged steers the terminal drains parked because no subscriber was
* live to receive the final/abort event. Reads are replayable; a recovery
* turn leases its exact source and removes it only after durable persistence. */
let unrecoveredSteers;
if (!isActive || job.metadata.steersClosed === true) {
const claimed = await GenerationJobManager.steering.claimDetailed(
conversationId,
{
userId: req.user.id,
tenantId: req.user.tenantId,
},
getRequestedGenerationProtocol(req),
);
if (claimed.steers.length > 0) {
generationProtocolVersion = Math.min(
generationProtocolVersion,
claimed.generationProtocolVersion,
);
res.set(GENERATION_PROTOCOL_HEADER, String(generationProtocolVersion));
unrecoveredSteers = claimed.steers;
}
}
res.json({
active: isActive,
generationProtocolVersion,
...(unrecoveredSteers && { unrecoveredSteers }),
streamId: conversationId,
status: job.status,
aggregatedContent: resumeState?.aggregatedContent ?? [],
createdAt: job.createdAt,
resumeState,
// Surface the live pending approval so a client rebuilding from /chat/status
// (reload / cross-replica) has the action id + payload to render and submit
// the prompt, not just the knowledge that the stream is paused. Client-safe
// projection only — resumeContext/requestFingerprint stay server-side.
pendingAction:
job.status === 'requires_action' && pendingLive
? toClientPendingAction(pendingAction)
: undefined,
});
});
/**
* @route POST /chat/abort
* @desc Abort an ongoing generation job
* @access Private
* @description Mounted before chatRouter to bypass buildEndpointOption middleware
*/
router.post('/chat/abort', configMiddleware, async (req, res, next) => {
logger.debug(`[AgentStream] ========== ABORT ENDPOINT HIT ==========`);
logger.debug(`[AgentStream] Method: ${req.method}, Path: ${req.path}`);
logger.debug(`[AgentStream] Body:`, req.body);
const requestProtocolVersion = negotiateRequestGenerationProtocol(req);
let responseProtocolVersion = requestProtocolVersion;
try {
if (req.body == null || typeof req.body !== 'object' || Array.isArray(req.body)) {
return sendGenerationJson(res, 400, { code: 'INVALID_ABORT_TARGET' }, requestProtocolVersion);
}
const { streamId, conversationId, abortKey, generationCreatedAt } = req.body;
for (const value of [streamId, conversationId, abortKey]) {
if (
value != null &&
(typeof value !== 'string' || value.length === 0 || value.length > 512)
) {
return sendGenerationJson(
res,
400,
{ code: 'INVALID_ABORT_TARGET' },
requestProtocolVersion,
);
}
}
const userId = req.user?.id;
if (
generationCreatedAt != null &&
(!Number.isSafeInteger(generationCreatedAt) || generationCreatedAt < 0)
) {
return sendGenerationJson(
res,
400,
{ code: 'INVALID_GENERATION_IDENTITY' },
requestProtocolVersion,
);
}
// streamId === conversationId, so try any of the provided IDs
// Skip "new" as it's a placeholder for new conversations, not an actual ID.
const streamCandidate = streamId && streamId !== 'new' ? streamId : null;
const conversationCandidate =
conversationId && conversationId !== 'new' ? conversationId : null;
const abortCandidate = abortKey?.split(':')[0];
const abortKeyCandidate = abortCandidate && abortCandidate !== 'new' ? abortCandidate : null;
let jobStreamId = streamCandidate || conversationCandidate || abortKeyCandidate || null;
let job = jobStreamId ? await GenerationJobManager.getJob(jobStreamId) : null;
/** Fallback only for the explicit new-conversation placeholder. An unknown
* concrete id (including a typo/stale tab) must never abort an unrelated
* active job. If several new-chat starts are active, the epoch selects the
* exact one; an unfenced legacy request is safe only when unambiguous. */
const canResolveNewPlaceholder =
!jobStreamId && (streamId === 'new' || conversationId === 'new') && userId;
if (!job && canResolveNewPlaceholder) {
logger.debug(`[AgentStream] Job not found by ID, checking active jobs for user: ${userId}`);
const activeJobIds = await GenerationJobManager.getActiveJobIdsForUser(
userId,
req.user.tenantId,
);
const candidates = [];
for (const activeJobId of activeJobIds) {
const activeJob = await GenerationJobManager.getJob(activeJobId);
if (
!activeJob ||
(activeJob.status !== 'running' && activeJob.status !== 'requires_action') ||
activeJob.metadata?.userId !== userId ||
hasTenantMismatch(activeJob, req.user) ||
(generationCreatedAt != null && activeJob.createdAt !== generationCreatedAt)
) {
continue;
}
candidates.push({ streamId: activeJobId, job: activeJob });
}
if (candidates.length > 1) {
return sendGenerationJson(
res,
409,
{ code: 'AMBIGUOUS_ACTIVE_RUN' },
requestProtocolVersion,
);
}
if (candidates.length === 1) {
jobStreamId = candidates[0].streamId;
job = candidates[0].job;
logger.debug(`[AgentStream] Found active job for user: ${jobStreamId}`);
}
}
logger.debug(`[AgentStream] Computed jobStreamId: ${jobStreamId}`);
if (job && jobStreamId) {
if (job.metadata?.userId !== userId) {
logger.warn(
`[AgentStream] Unauthorized abort attempt for ${jobStreamId} by user ${userId}`,
);
return sendGenerationJson(res, 403, { error: 'Unauthorized' }, requestProtocolVersion);
}
if (hasTenantMismatch(job, req.user)) {
return sendGenerationJson(res, 403, { error: 'Unauthorized' }, requestProtocolVersion);
}
const generationProtocolVersion = negotiateExistingGenerationProtocol(req, job);
responseProtocolVersion = generationProtocolVersion;
res.set(GENERATION_PROTOCOL_HEADER, String(generationProtocolVersion));
if (generationCreatedAt != null && job.createdAt !== generationCreatedAt) {
return res.status(409).json({ code: 'RUN_REPLACED', generationProtocolVersion });
}
logger.debug(`[AgentStream] Job found, aborting: ${jobStreamId}`);
// Re-attach a paused ask_user_question's args to the abort content BEFORE
// abortJob emits the final SSE. Redis reconstructs abort content from the
// chunk log, which never saw the pause-time stamp applied to the in-process
// contentParts — stamping inside abortJob (not after) means the LIVE client
// gets the question too, not just the saved message on reload.
const abortedAskPayload = job.metadata?.pendingAction?.payload;
const agentsCfg = req.config?.endpoints?.agents;
const shouldPruneCheckpoint =
isHITLEnabled(agentsCfg?.toolApproval) || job.metadata?.pendingAction != null;
const checkpointNamespace =
typeof job.metadata?.checkpointNamespace === 'string'
? job.metadata.checkpointNamespace
: '';
/** New jobs have an immutable saver-level namespace, so the terminal
* owner can delete that entire namespace (including a checkpoint written
* after this route's initial read) without touching a replacement. Legacy
* jobs share the root namespace and still need an id snapshot before CAS. */
const checkpointGeneration =
shouldPruneCheckpoint && checkpointNamespace === ''
? await captureAgentCheckpointGeneration(jobStreamId, agentsCfg?.checkpointer, {
throwOnError: true,
})
: undefined;
const abortResult = await GenerationJobManager.abortJob(jobStreamId, {
expectedCreatedAt: job.createdAt,
transformAbortContent: (content) =>
abortedAskPayload?.type === 'ask_user_question' && Array.isArray(content)
? attachAskUserQuestionArgs(content, abortedAskPayload.question)
: content,
/** Persist every parent-row prerequisite before publishing the ordinary
* abort FINAL. That frame can immediately drain a queued follow-up, whose
* parent must already exist and whose graph must not see a stale HITL
* checkpoint. Throwing makes the manager publish a conservative
* reconciliation frame instead of an unsafe normal FINAL. */
beforePublish: async (pendingAbortResult) => {
const persistenceErrors = [];
const { jobData, text, content } = pendingAbortResult;
/** `abortJob` treats a delivered `created` event as a real turn even
* when every streamed part is filtered out (for example, an
* interrupt before the model's first non-whitespace token). Its
* normal FINAL therefore carries an empty unfinished assistant.
* Persist that same row before publishing, including when its id is
* the underscore-suffixed preliminary id rendered by `created`.
* Otherwise interrupt-and-send immediately posts that unsaved id as
* its parent and the preliminary-parent fence correctly rejects it. */
const shouldPersistAbortedTurn =
hasPersistableAbortContent(content) || jobData?.createdEventEmitted === true;
if (
jobData?.userMessage?.messageId &&
jobData?.responseMessageId &&
shouldPersistAbortedTurn
) {
const messageContext = {
userId: req?.user?.id,
// Source from the job: the stop request does not carry the
// original temporary-chat flag.
isTemporary: jobData?.isTemporary ?? req?.body?.isTemporary,
interfaceConfig: req?.config?.interfaceConfig,
};
const requestMessage = {
...jobData.userMessage,
conversationId: jobData.conversationId,
sender: 'User',
endpoint: jobData.endpoint,
isCreatedByUser: true,
user: userId,
};
const responseMessage = {
messageId: jobData.responseMessageId,
parentMessageId: jobData.userMessage.messageId,
conversationId: jobData.conversationId,
content: content || [],
text: text || '',
sender: jobData.sender || 'AI',
endpoint: jobData.endpoint,
iconURL: jobData.iconURL,
model: jobData.model,
unfinished: true,
error: false,
isCreatedByUser: false,
user: userId,
};
const abortMetadata = buildAbortedResponseMetadata(jobData);
if (abortMetadata) {
responseMessage.metadata = abortMetadata;
}
/** `created` fires before BaseClient starts its asynchronous user
* write. A very early interrupt can therefore reach this barrier
* with neither row stored. Both writes are idempotent upserts;
* await the user prerequisite first, but still attempt the child
* write and checkpoint cleanup so every independently useful
* operation gets a chance to succeed. */
try {
const persistedRequest = await saveMessage(messageContext, requestMessage, {
context: 'api/server/routes/agents/index.js - abort user prerequisite',
});
if (!persistedRequest) {
throw new Error('Abort user prerequisite was not persisted');
}
} catch (error) {
persistenceErrors.push(error);
}
try {
const persistedResponse = await saveMessage(messageContext, responseMessage, {
context: 'api/server/routes/agents/index.js - abort endpoint',
});
if (!persistedResponse) {
throw new Error('Abort response was not persisted');
}
logger.debug(`[AgentStream] Saved partial response for: ${jobStreamId}`);
} catch (error) {
persistenceErrors.push(error);
}
}
/** Attempt checkpoint cleanup even when the message write failed, and
* attempt the message write even when cleanup will fail. Both are
* independently valuable; any failure still suppresses the normal
* FINAL after all required work has been attempted. */
if (shouldPruneCheckpoint) {
try {
await deleteAgentCheckpoint(
jobStreamId,
agentsCfg?.checkpointer,
checkpointGeneration,
checkpointNamespace !== ''
? { throwOnError: true, checkpointNamespace }
: { throwOnError: true },
);
} catch (error) {
persistenceErrors.push(error);
}
}
if (persistenceErrors.length === 1) {
throw persistenceErrors[0];
}
if (persistenceErrors.length > 1) {
const error = new Error('Abort message persistence and checkpoint cleanup failed');
error.causes = persistenceErrors;
throw error;
}
},
});
if (abortResult.failureReason === 'generation_replaced') {
return res.status(409).json({ code: 'RUN_REPLACED', generationProtocolVersion });
}
if (abortResult.failureReason === 'job_still_active') {
res.set('Retry-After', '1');
return res.status(409).json({ code: 'RUN_STILL_ACTIVE', generationProtocolVersion });
}
if (!abortResult.success) {
// The route authorized a live generation, but the manager can lose its
// terminal CAS to natural completion/error (or observe deletion before
// its own lookup). Never claim that Stop won when no abort FINAL exists.
if (!abortResult.jobData) {
return res.status(404).json({
success: false,
error: 'Job not found',
streamId: jobStreamId,
generationProtocolVersion,
});
}
const currentJob = await GenerationJobManager.getJob(jobStreamId);
if (currentJob && currentJob.createdAt !== job.createdAt) {
return res.status(409).json({ code: 'RUN_REPLACED', generationProtocolVersion });
}
if (currentJob?.status === 'running' || currentJob?.status === 'requires_action') {
res.set('Retry-After', '1');
return res.status(409).json({ code: 'RUN_STILL_ACTIVE', generationProtocolVersion });
}
if (generationProtocolVersion < GENERATION_PROTOCOL_V2) {
return res.json({
success: true,
aborted: jobStreamId,
generationProtocolVersion,
});
}
return res.json({
success: false,
settled: true,
code: 'RUN_ALREADY_SETTLED',
streamId: jobStreamId,
generationProtocolVersion,
...(currentJob?.status && { terminalStatus: currentJob.status }),
});
}
logger.debug(`[AgentStream] Job aborted successfully: ${jobStreamId}`, {
abortResultSuccess: abortResult.success,
abortResultUserMessageId: abortResult.jobData?.userMessage?.messageId,
abortResultResponseMessageId: abortResult.jobData?.responseMessageId,
});
if (abortResult.persistenceFailed && generationProtocolVersion < GENERATION_PROTOCOL_V2) {
res.set('Retry-After', '1');
return res.status(409).json({
code: 'ABORT_PERSISTENCE_FAILED',
generationProtocolVersion,
});
}
return res.json({
success: true,
aborted: jobStreamId,
generationProtocolVersion,
...(abortResult.persistenceFailed && { persistenceFailed: true }),
// Steers that never reached an injection boundary — restored client-side
// as queued chips so the user's words aren't dropped with the abort.
...(!abortResult.persistenceFailed &&
abortResult.pendingSteers?.length > 0 && { pendingSteers: abortResult.pendingSteers }),
});
}
logger.warn(`[AgentStream] Job not found for streamId: ${jobStreamId}`);
return sendGenerationJson(
res,
404,
{ error: 'Job not found', streamId: jobStreamId },
requestProtocolVersion,
);
} catch (error) {
logger.error('[AgentStream] Abort request failed', error);
if (res.headersSent) {
return next(error);
}
return sendGenerationJson(
res,
500,
{ code: 'ABORT_FAILED', error: 'Failed to abort generation' },
responseProtocolVersion,
);
}
});
/**
* @route POST /chat/steer
* @desc Queue a mid-run user message for injection at the next tool boundary
* @access Private
* @description Mounted before chatRouter to bypass buildEndpointOption middleware,
* but a steer is model-bound user text, so it carries the same guards as a normal
* message IN THE SAME ORDER as chat.js: the configured IP/user rate limiters,
* the PII filter FIRST (blocked sensitive text must never reach the external
* moderation endpoint), then `moderateText`.
*/
const steerLimiters = [];
if (isEnabled(LIMIT_MESSAGE_IP)) {
steerLimiters.push(messageIpLimiter);
}
if (isEnabled(LIMIT_MESSAGE_USER)) {
steerLimiters.push(messageUserLimiter);
}
router.post(
'/chat/steer',
configMiddleware,
...steerLimiters,
createMessageFilterPii({ getConfig: (req) => req.config?.messageFilter?.pii }),
moderateText,
SteerController,
);
/**
* @route POST /chat/steer/cancel
* @desc Remove a still-queued steer before injection (no model-bound content,
* so no PII/moderation pass — just the shared rate limiters)
* @access Private
*/
router.post(
'/chat/steer/cancel',
configMiddleware,
...steerLimiters,
SteerController.SteerCancelController,
);
/**
* @route POST /chat/steer/arm
* @desc Escalate a still-queued steer to an interrupt in place (no new
* model-bound content, so no PII/moderation pass — just the shared limiters)
* @access Private
*/
router.post(
'/chat/steer/arm',
configMiddleware,
...steerLimiters,
SteerController.SteerArmController,
);
router.use('/', v1);
const chatRouter = express.Router();
chatRouter.use(configMiddleware);
if (isEnabled(LIMIT_MESSAGE_IP)) {
chatRouter.use(messageIpLimiter);
}
if (isEnabled(LIMIT_MESSAGE_USER)) {
chatRouter.use(messageUserLimiter);
}
chatRouter.use('/', chat);
router.use('/chat', chatRouter);
module.exports = router;