LibreChat/api/server/controllers/agents/request.js
Danny Avila 8c14f03432
🗂️ feat: Scope Scheduled Chats to Chat Projects (#15056)
* 🗂️ feat: Scope Scheduled Chats to Chat Projects

Adds an optional chat-project destination to a schedule, plus the operator
config to require one — or to pin every scheduled run to a specific project.

Feature:
- `chatProjectId` on the schedule row, accepted on create/update, projected on
  the wire, and carried into the run's conversation through the durable trigger
  envelope's `run` context.
- `interface.schedules.requireProject` refuses schedules that are not filed
  under a project; `interface.schedules.projectId` pins every run to one
  project and implies the requirement.
- Dialog gains a project picker (required when configured, a read-only row when
  pinned); the card shows the destination and the new disabled reasons.

Invariants:
- ONE resolver (`resolveScheduleProjectId`) decides the destination for the
  write handler, the fire path, and the wire projection alike, and an operator
  pin OUTRANKS the stored id in all three. Tightening the config therefore
  redirects — or stops — existing schedules instead of grandfathering where
  their runs land. A pin implies `requireProject` for the same reason: without
  it, a row created before the pin would keep firing with no project at all.
- Create/fire precheck symmetry, mirroring `resolveAgentFireAccess`: a write
  this handler accepts is one the next fire also accepts. Any edit leaving a
  schedule ENABLED re-validates its EFFECTIVE (possibly stored) project, like
  the existing stored-agent and cadence-floor rechecks. A DISABLING edit skips
  the requirement, or a schedule auto-disabled for `project_required` could
  never be turned off.
- Fire-time enforcement auto-disables rather than filing runs loose, matching
  agent_deleted: new `project_required` (requirement raised after creation) and
  `project_deleted` (gone, or pinned to a project this owner does not have)
  reasons, both refused BEFORE a billed generation is dispatched, and both
  advancing so a schedule can never wedge on the occurrence.
- `computeCreateDigest` appends the field only when present, so a payload
  without a project digests byte-identically to one from before this change —
  a create retried across the upgrade still matches its own row instead of
  reading as key reuse.
- Project reads are scoped to the owner, so ownership and existence are the
  same lookup; a read ERROR propagates instead of failing closed, so a Mongo
  blip retries the fire rather than auto-disabling the schedule.

The trigger idempotency key hashes principal/event/target and never
`envelope.run`, so the added run field cannot destabilize delivery identity.

* 🗜️ fix: Keep the Schedule Dialog Inside Its Height Budget

The project picker landed as a new ROW in the schedule dialog, which broke the
e2e edit spec: `md:overflow-visible` turns off the template's scrolling from
`md` up, so the dialog's content must fit the viewport. The extra row pushed the
footer's Save button below a 720x1280 window, where Playwright reported a
visible, enabled button it could never click — 226 scroll-into-view retries and
a 2-minute timeout, on all three attempts.

Measured against `dev` at 1280x720 (Save button's bottom edge, viewport 720):
  dev              688   (32px slack)
  project row      ~790  (off-screen, CI failure)
  3-column row     704   (16px slack — half the budget spent)
  this commit      686   (34px slack, 2px better than dev)

The identity row is now three columns — name, agent, project — and its caption
moved out of the agent cell to sit full width beneath the row: at a third of the
dialog that sentence wraps an extra line, and the row is the tallest thing
competing for the budget. The caption is grouped with the row rather than left
to the form's own 4-unit rhythm, which spent more height on the gap than the
caption occupies.

The e2e spec now asserts the button is in the viewport before clicking it, so
the next field that overflows this dialog says so in one line instead of a
two-minute timeout on a visible element. FOLLOWUPS.md records what the planned
dialog controls (multi-day weekly, timezone, attachments) need first: give
`ControlCombobox` the `portalElement` prop `Dropdown` already has, portal the
popovers into the dialog content, and let the form scroll again.

* 🧹 fix: Address Codex Review on Scheduled Chat Project Scope

Four P2 findings, all real.

Unreachable clearing path (ScheduleDialog). The picker only held live projects, so
`com_ui_schedule_project_none` was a PLACEHOLDER — nothing selectable. Once a
schedule had a project the owner could never take it away, leaving the server's
`chatProjectId: null` path reachable only by API. The picker now carries a real
"No project" option whenever a project is optional, and omits it when one is
required, where there is nothing valid to select.

Placeholder shown for a real project (ScheduleDialog). A stored or pinned project
outside the first loaded page had no name in the paged map, and the combobox
renders its placeholder for an empty display value — telling the owner a scoped
schedule had no project. That one project is now read by id, with the raw id as a
last resort: a poor label, but an honest one.

Project policy skipped at the resume boundary (service.ts). `claimScheduleResume`
re-applied the schedules gate, the revision fence, the kill switch and
SCHEDULES:USE, but not the project policy this PR added. Approving a paused run
whose project was deleted — or whose owner now sits under a requirement or a pin
it no longer satisfies — billed a continuation the very next scheduled fire would
refuse and auto-disable the schedule for. The effective-project resolution now
runs there too, refused before the lease and the capacity slot so a policy refusal
costs nothing and leaves no state to unwind. NOTE: agent access and balance are
still not rechecked on resume; that gap predates this PR and is left alone.

Per-card project derivation (ScheduleCard). Every card ran the projects hook and
rebuilt the full option array, name map, and one icon element per project, to use
a single name — O(schedules x projects) per render and per project-list refresh.
The hook is split: `useChatProjectNames` (map only, for the panel, which resolves
every card's name once and passes it down) and `useChatProjectPicker` (options and
pagination, for the dialog's one combobox). The panel skips the query entirely
until some schedule actually has a scope.

Tests: four at the resume boundary (verified failing without the gate) and four in
the dialog spec. The picker selection in one existing test now goes through the
search field — the popover's VIRTUALIZED renderer sizes its window from a scroll
height jsdom always reports as 0, so with three options it materialized only two.
Full schedules e2e re-run green against the rebuilt client.

* 🎯 fix: Settle Project-Policy Refusals and Keep Create Retries Idempotent

Second Codex round, four P2s. Three were consequences of the resume gate added in
the previous commit, which was half-built: it admitted where it should not and
stranded the run where it refused.

Project policy moves from `claimScheduleResume` into `isScheduleLive`'s `policy`
branch. Both entry points consult that branch FIRST, and both already route its
refusal through abort-and-settle — so a policy stop now settles the occurrence
instead of answering a bare 409 while the job stays `requires_action`, the card
keeps reading "Needs approval", and every retry repeats the same 409 until expiry.
No change to resume.js: the existing branch does the work.

The rule is deliberately NARROW. It refuses only where no valid destination is
left — the requirement is on with nothing satisfying it, or the schedule's own
project is gone (which also unset it on the conversation). It does NOT refuse
because an operator's pin moved: the paused conversation cannot be rebound
(`chatProjectId` is excluded from the resume context and the continuation reuses
the same conversationId), so refusing would strand a pending approval over a
destination it can never reach, for a pin that governs only where the NEXT run
lands — which the fire path already redirects.

Create retries are idempotent again. Project policy had been applied BEFORE the
`clientRequestId` replay lookup, so a raised requirement, a deleted project, or a
moved pin could answer 400 for a create that already committed — pushing the
client to rotate its key and create a DUPLICATE schedule, the exact failure the
key exists to prevent. Policy now applies only to a genuinely new insert, and the
digest is computed from the CLIENT's payload rather than the resolved destination,
so today's policy can no longer re-digest a genuine retry into a mismatch.

An explicit `chatProjectId: null` under a pin is refused rather than silently
resolved to the pin. The payload contract defines `null` as clearing the scope;
answering 201 while filing under the pin reported success for the opposite of what
was asked. Only an OMITTED field takes the pin silently.

Tests: five on the policy branch (both refusals verified failing without it, plus
guards that a moved pin and a live project still admit) and two on the handlers
(the pinned explicit clear, and a committed create recovered by retry after the
policy tightened — verified answering 400 without the reordering).

* 🧭 fix: Converge the Stored Project on the Destination a Fire Resolved

Third Codex round, four P2s.

The root confusion behind the resume findings: an operator pin outranks the stored
id at fire time, `fireSchedule` sends the pin in the trigger envelope, and the row
keeps its old value. The row therefore LIED about where that occurrence's
conversation went, and every later re-validation — the resume boundary above all —
checked a project the conversation was never filed under. A schedule storing A,
pinned to B, with B later deleted and A still live, was admitted for resume into a
conversation that had just been unscoped.

Fixed at the source rather than at each reader: a fire that resolves a destination
different from the stored one writes it back, claim-token fenced like every other
worker-side write and deliberately WITHOUT a configRevision bump — this is the
server reconciling itself to policy, not an owner edit, and a bump would fence an
in-flight occurrence off its own run. Written only AFTER the destination validates,
so an unusable pin never lands in the row, and best-effort: the envelope already
carries the right destination, so a failed write costs accuracy on a later recheck,
never the run itself. The wire projection already reported the pin, so this also
stops the row and the UI disagreeing.

An explicit `chatProjectId: null` under a pin is now refused on the DISABLING edit
path too. The pin check and the requirement are independent rules, and folding them
together let `{enabled: false, chatProjectId: null}` skip the pin check entirely,
unset the row, and answer with a wire projection still naming the pin. Only the
requirement is waived for a disabling edit.

The dialog no longer requires a project for an edit that leaves a schedule DISABLED.
The server waives the requirement there precisely so a row auto-disabled for
`project_required` can still be renamed or tidied up; requiring it in the form made
that unreachable, and an owner with no projects could not edit the stopped schedule
at all.

FOLLOWUPS.md records the two residual gaps with their exact triggers: the sub-second
deletion race inside the resume claim window (which needs the effective project
persisted per OCCURRENCE plus a distinct policy conflict routed through
abort-and-settle), and the fact that convergence happens only when a schedule fires.

Tests: four on convergence (pin written, no write when unchanged, no write for a
destination that failed validation, fire survives a failed write), two on the
disabling-edit pin rules, two in the dialog. Schedules e2e re-run green.

* 🔑 fix: Keep an Explicit Project Clear Out of an Omitted Field's Digest

`computeCreateDigest` appended `chatProjectId` on `!= null`, so an OMITTED field and
an explicit `null` produced the same digest. Because the replay lookup and
`matchesCreateIntent` deliberately run before project policy, a request could reuse a
pinned create's `clientRequestId` while explicitly sending `chatProjectId: null` and
receive 201 describing the pinned row — success reported for the opposite of what it
asked, and the pinned-clear refusal the normal create path applies never reached.

Now `!== undefined`: an omitted field still digests byte-identically to a payload from
before project scope existed, so a create in flight across the upgrade still matches
its own row, while an explicit clear is a distinct intent and digests differently. A
pre-scope client never sent the field at all, so nothing legacy can carry an explicit
null.

* 📍 fix: Validate a Paused Run Against the Project Its Own Occurrence Used

The schedule-wide convergence from the previous commit was not enough, and the
reason is the single-active run index: it covers `status: 'started'` only, so a
PAUSED run does not block the next occurrence. While run 1 sat paused in project A,
a pin move plus one later fire rewrote the schedule row to B — and the resume
policy then validated B while run 1's conversation was still filed under A. Delete
A and the continuation was admitted into a conversation that had just been unscoped.
That window lasts as long as the pause, not the sub-second race the previous commit
documented.

The reservation now records the destination THIS occurrence used, and
`isScheduleLive` validates that record when given the occurrence's `scheduledFor` —
which `resume.js` already reads two lines above the call. No new conflict type and
no settlement-path surgery: the refusal rides the abort-and-settle branch that check
already has.

An ABSENT record falls back to the schedule-level resolution rather than reading as
unscoped. A pre-scope occurrence, or one whose row is gone, must never be treated as
evidence to stop a run — the fallback keeps legacy paused runs behaving exactly as
they do today.

Schedule-wide convergence stays: it keeps the row honest for the UI and for every
check that has no occurrence in hand.

FOLLOWUPS.md now describes the one remaining gap accurately — a deletion inside the
claim window, which needs a distinct policy conflict routed through abort-and-settle
rather than the bare 409 an `inactive` conflict produces.

Tests: three on occurrence-vs-row precedence (the decisive one verified failing
without the lookup), two on what the reservation records, and the resume controller
spec now pins `scheduledFor` in the policy call. Schedules e2e re-run green.

* 🏷️ fix: Tell a Deliberately Unscoped Occurrence From an Unrecorded One

The occurrence fallback added in the previous commit conflated two different
absences. A post-upgrade run that deliberately went unscoped omitted the field
exactly like a row written before the field existed, so both took the fallback —
and a paused unscoped run was then validated against the schedule's CURRENT
project. Under a requirement or a pin added while it sat paused, that admitted a
billed continuation into a conversation satisfying no present policy.

The reservation now ALWAYS records its decision, `null` for unscoped, and the read
reports `recorded` from key PRESENCE rather than truthiness. Only an unknown record
— a pre-scope row, or no row at all — falls back to the schedule-level resolution;
a recorded null is the genuinely unscoped occurrence it says it is, and is refused
once a project becomes required.

The distinction rests entirely on a stored `null` surviving as a present key while a
never-written field stays absent, so that is asserted against real Mongo rather than
assumed: if it ever stopped holding, unscoped runs would silently start being
validated against the schedule's current project again.

Tests: two against mongodb-memory-server (recorded null vs never-written vs missing
row), plus refusal of a recorded-unscoped occurrence under a new requirement and the
preserved fallback for a pre-scope one. Schedules e2e green.

* 🔁 fix: Validate an Initial Scheduled Start Against Its Own Occurrence

The initial-start policy check in `request.js` called `isScheduleLive(..., { policy:
true })` without `scheduledFor`, so it fell back to the schedule-level resolution
even though the run row already carries the occurrence's recorded scope. An
occurrence reserved unscoped, with a pin introduced while its loopback request sat
queued, was therefore admitted against the new pin — producing a billed unscoped
conversation under a requirement it does not satisfy, from an envelope already built
without a project.

`scheduledFor` was already in scope there. Passing it makes the initial start and the
resume validate the same way: against the destination the occurrence itself recorded.
2026-08-21 11:09:04 -04:00

2267 lines
89 KiB
JavaScript

const { logger, tenantStorage } = require('@librechat/data-schemas');
const { v5: uuidv5 } = require('uuid');
const {
Constants,
EModelEndpoint,
ErrorTypes,
ViolationTypes,
isEphemeralAgentId,
} = require('librechat-data-provider');
const {
toPendingSteer,
getViolationInfo,
buildMessageFiles,
getReferencedQuotes,
resolveTitleTiming,
GenerationJobManager,
filterPersistableAbortContent,
decrementPendingRequest,
sanitizeMessageForTransmit,
checkAndIncrementPendingRequest,
exemptFromConcurrencyLimiter,
isScheduleFireRequest,
isUnpersistedPreliminaryParent,
resolveConversationAnchor,
getAgentStartupTelemetry,
acceptAgentStartupTelemetry,
isSteerPreemptSupported,
buildRecoveredSteerPayload,
deleteAgentCheckpoint,
getAttachmentTitleText,
} = require('@librechat/api');
const { disposeClient } = require('~/server/cleanup');
const {
getMCPRequestContext,
cleanupMCPRequestContextForReq,
} = require('~/server/services/MCPRequestContext');
const { logViolation } = require('~/cache');
const { recordScheduleOutcome, isScheduleLive } = require('~/server/services/Schedules');
const { saveMessage, getMessages, getConvo, isAgentTriggerPrincipalActive } = require('~/models');
const {
GENERATION_PROTOCOL_HEADER,
GENERATION_PROTOCOL_V2,
negotiateNewGenerationProtocol,
negotiateExistingGenerationProtocol,
} = require('./protocol');
function sendGenerationJson(res, status, body, generationProtocolVersion) {
if (typeof res.set === 'function') {
res.set(GENERATION_PROTOCOL_HEADER, String(generationProtocolVersion));
} else if (typeof res.setHeader === 'function') {
res.setHeader(GENERATION_PROTOCOL_HEADER, String(generationProtocolVersion));
}
return res.status(status).json({ ...body, generationProtocolVersion });
}
function getInitializationFailure(error) {
if (error?.code === ErrorTypes.RESOURCE_RECOVERY_REQUIRED) {
return {
status: 409,
code: ErrorTypes.RESOURCE_RECOVERY_REQUIRED,
error: error.message || 'Attached resources must be restored before retrying.',
};
}
const candidateStatus = error?.status ?? error?.statusCode;
if (!Number.isInteger(candidateStatus) || candidateStatus < 400 || candidateStatus >= 600) {
return null;
}
return {
status: candidateStatus,
...(typeof error?.code === 'string' ? { code: error.code } : {}),
error: error?.message || 'Failed to start generation',
};
}
function resolveConversationCreatedAt({ userId, conversationId, isNewConvo, conversation }) {
return resolveConversationAnchor({
isNewConversation: isNewConvo,
loadConversation: () =>
conversation !== undefined ? Promise.resolve(conversation) : getConvo(userId, conversationId),
onLoadError: (error) => {
logger.warn('[AgentController] Failed to resolve conversation timestamp anchor', {
conversationId,
error: error.message,
});
},
});
}
async function attachConversationCreatedAt(req, conversationId, conversationAnchorPromise) {
req.body.conversationId = conversationId;
const resolved = await conversationAnchorPromise;
req.conversationCreatedAt = resolved.createdAt;
if (resolved.conversation !== undefined) {
req.resolvedConversation = resolved.conversation ?? null;
}
}
function getPreliminaryResponseMessageId({ messageId, responseMessageId }) {
if (typeof responseMessageId === 'string' && responseMessageId.length > 0) {
return responseMessageId;
}
if (typeof messageId !== 'string' || messageId.length === 0) {
return null;
}
return `${messageId.replace(/_+$/, '')}_`;
}
function getPreliminaryUserMessage(
{ messageId, parentMessageId, text, quotes, files, manualSkills, alwaysAppliedSkills },
conversationId,
) {
if (typeof messageId !== 'string' || messageId.length === 0) {
return null;
}
/**
* Seed normalized quotes here too: if the user aborts before `sendMessage`
* reaches `onStart` (during init/tool loading), `abortMiddleware` falls back
* to this preliminary metadata, which must carry the excerpts so the stopped
* turn keeps its `MessageQuotes`.
*/
const referencedQuotes = getReferencedQuotes(quotes);
return {
messageId,
parentMessageId,
conversationId,
text,
...(referencedQuotes != null && { quotes: referencedQuotes }),
// Persist the turn's uploaded files on this AWAITED preliminary write so they land on
// job.metadata.userMessage BEFORE the run can reach its first interrupt. onStart's
// later writes are fire-and-forget, so a fast approval could otherwise read the job
// and resume an approved code/read-file tool without the paused turn's uploads.
...(Array.isArray(files) && files.length > 0 && { files }),
// Carry skill selections so a HITL-resumed turn's reconstructed `requestMessage`
// keeps its skill pills — the client's final handler replaces the user bubble from
// this object, and they'd otherwise vanish until a full reload refetches the row.
...(Array.isArray(manualSkills) && manualSkills.length > 0 && { manualSkills }),
...(Array.isArray(alwaysAppliedSkills) &&
alwaysAppliedSkills.length > 0 && { alwaysAppliedSkills }),
};
}
function getRequestModelSpec(req, endpointOption) {
const spec = endpointOption?.spec ?? req.body?.spec;
if (typeof spec !== 'string' || spec.length === 0) {
return;
}
const list = req.config?.modelSpecs?.list;
if (!Array.isArray(list)) {
return;
}
return list.find((modelSpec) => modelSpec?.name === spec);
}
function getModelSpecIconURL(modelSpec) {
return modelSpec?.iconURL ?? modelSpec?.preset?.iconURL ?? modelSpec?.preset?.endpoint ?? '';
}
function getEndpointIconURL(req, endpointOption) {
const iconURL =
endpointOption?.iconURL ?? getModelSpecIconURL(getRequestModelSpec(req, endpointOption));
return iconURL || undefined;
}
function getEndpointResponseModel(endpointOption) {
return endpointOption?.modelOptions?.model || endpointOption?.model_parameters?.model;
}
function getAgentResponseModel(req, endpointOption) {
const agentId = endpointOption?.agent_id || req.body?.agent_id;
if (typeof agentId === 'string' && agentId.length > 0 && !isEphemeralAgentId(agentId)) {
return agentId;
}
return getEndpointResponseModel(endpointOption);
}
async function finishResumableRequest(req, userId) {
try {
await cleanupMCPRequestContextForReq(req);
} finally {
if (req._scheduleConcurrencyExempt !== true) {
await decrementPendingRequest(userId);
}
}
}
function classifyScheduledFailure(error, aborted = false) {
if (aborted || error?.code === 'SCHEDULE_NO_LONGER_ACTIVE') {
return { status: 'interrupted', error: error?.message };
}
if (error?.message?.includes(ViolationTypes.TOKEN_BALANCE)) {
return { status: 'skipped_balance' };
}
return { status: 'error', error: error?.message || 'Generation failed' };
}
const JOB_RECORD_WAIT_ATTEMPTS = 5;
const JOB_RECORD_WAIT_DELAY_MS = 60;
// A winner writes its job record within a few ms of claiming; if a losing duplicate still
// sees no job within this window of the claim, the winner is still starting (retry rather
// than hand back a stream that would 404). Past it, a missing job means the original
// already completed and was cleaned up (attach and let the client refetch).
const IDEMPOTENCY_STARTUP_GRACE_MS = 5000;
const CLIENT_REQUEST_ID_PATTERN = /^[A-Za-z0-9:_-]{1,128}$/;
/** New-chat retries do not carry a conversation id, so derive the stream id
* from their stable per-submission id. This keeps both the dedupe key and the
* Redis hash slot identical across a lost-response retry. */
const NEW_CONVERSATION_IDEMPOTENCY_NAMESPACE = 'd7f2518c-94b8-4fe8-97ad-2d4bdb2c9f43';
function isValidGenerationClaim(value, streamId, conversationId, requireStarted = false) {
return (
value != null &&
typeof value === 'object' &&
value.streamId === streamId &&
value.conversationId === conversationId &&
Number.isSafeInteger(value.claimedAt) &&
value.claimedAt >= 0 &&
typeof value.claimToken === 'string' &&
value.claimToken.length > 0 &&
value.claimToken.length <= 128 &&
(value.generationProtocolVersion == null ||
value.generationProtocolVersion === 1 ||
value.generationProtocolVersion === GENERATION_PROTOCOL_V2) &&
(value.startedAt == null || (Number.isSafeInteger(value.startedAt) && value.startedAt >= 0)) &&
(!requireStarted || value.startedAt != null)
);
}
/** Pre-bridge servers wrote the legacy global key without a claim token and,
* for a new conversation, chose a random stream before claiming it. Accept
* only that tightly bounded legacy shape: existing conversations must still
* match the requested stream exactly; new-chat claims may point to the old
* random stream only when streamId === conversationId. Ownership is verified
* against the live job before attachment. */
function isValidLegacyGenerationClaim(value, streamId, isNewConvo) {
return (
value != null &&
typeof value === 'object' &&
typeof value.streamId === 'string' &&
value.streamId.length > 0 &&
value.streamId.length <= 512 &&
value.conversationId === value.streamId &&
(isNewConvo || value.streamId === streamId) &&
Number.isSafeInteger(value.claimedAt) &&
value.claimedAt >= 0 &&
value.claimToken == null &&
value.startedAt == null &&
(value.generationProtocolVersion == null || value.generationProtocolVersion === 1)
);
}
/** Store corruption must not turn a user-scoped idempotency claim into a
* pointer to another user's/tenant's live stream. Missing tenant metadata is
* kept as the explicit legacy case, but missing ownership never authorizes. */
function liveJobBelongsToRequester(job, user) {
return (
job?.metadata?.userId === user.id &&
(job.metadata?.tenantId == null || job.metadata.tenantId === user.tenantId)
);
}
/**
* Poll briefly for a job record to appear. A deduped retry that loses the idempotency
* claim must not be handed the winner's stream until its job exists, or the client's
* subscribe 404s terminally. The winner writes the record a few ms after claiming.
*/
async function waitForJobRecord(streamId) {
for (let attempt = 0; attempt < JOB_RECORD_WAIT_ATTEMPTS; attempt++) {
const job = await GenerationJobManager.getJob(streamId);
if (job) {
return job;
}
await new Promise((resolve) => setTimeout(resolve, JOB_RECORD_WAIT_DELAY_MS));
}
return GenerationJobManager.getJob(streamId);
}
/** The claimed generation already reached durable/terminal history, but its
* conversation stream id now belongs to no job or to a newer submission. A
* success shape with that streamId would attach the stale submission to the
* replacement, so tell the client to refetch without opening SSE. */
function sendSettledGeneration(
res,
streamId,
conversationId,
startupTelemetry,
generationProtocolVersion,
) {
startupTelemetry?.end('deduplicated');
if (generationProtocolVersion < GENERATION_PROTOCOL_V2) {
return sendGenerationJson(
res,
200,
{ streamId, conversationId, status: 'resumed' },
generationProtocolVersion,
);
}
return sendGenerationJson(
res,
200,
{ conversationId, status: 'settled' },
generationProtocolVersion,
);
}
function rejectPreliminaryParentMessageId(res, generationProtocolVersion) {
return sendGenerationJson(
res,
409,
{
code: 'PARENT_NOT_READY',
error:
'Cannot submit a follow-up while the selected parent response is still being saved. Please wait and try again.',
},
generationProtocolVersion,
);
}
function rejectMissingTriggerParentMessageId(res, generationProtocolVersion) {
return sendGenerationJson(
res,
404,
{
code: 'PARENT_NOT_FOUND',
error: 'The selected parent response is no longer available.',
},
generationProtocolVersion,
);
}
/**
* Resumable Agent Controller - Generation runs independently of HTTP connection.
* Returns streamId immediately, client subscribes separately via SSE.
*/
const ResumableAgentController = async (req, res, next, initializeClient, addTitle) => {
const startupTelemetry = getAgentStartupTelemetry(req);
let generationProtocolVersion = negotiateNewGenerationProtocol(req, GenerationJobManager);
const {
text,
isRegenerate,
endpointOption,
conversationId: reqConversationId,
isContinued = false,
editedContent = null,
parentMessageId = null,
overrideParentMessageId = null,
responseMessageId: editedResponseMessageId = null,
scheduleId: bodyScheduleId = null,
scheduledFor: bodyScheduledFor = null,
scheduleConfigRevision: bodyScheduleConfigRevision = null,
} = req.body;
const isScheduledFire = isScheduleFireRequest(req);
const scheduleId = isScheduledFire ? bodyScheduleId : null;
const scheduledFor = isScheduledFire ? bodyScheduledFor : null;
const scheduleConfigRevision = isScheduledFire ? bodyScheduleConfigRevision : undefined;
const userId = req.user.id;
const tenantId = req.user.tenantId;
const rawClientRequestId = req.body?.clientRequestId;
if (
rawClientRequestId != null &&
(typeof rawClientRequestId !== 'string' || !CLIENT_REQUEST_ID_PATTERN.test(rawClientRequestId))
) {
startupTelemetry?.end('rejected');
return sendGenerationJson(
res,
400,
{
code: 'INVALID_CLIENT_REQUEST_ID',
error: 'clientRequestId must be a 1-128 character identifier.',
},
generationProtocolVersion,
);
}
const clientRequestId = rawClientRequestId;
const rawExpectedPredecessorCreatedAt = req.body?.expectedPredecessorCreatedAt;
if (
rawExpectedPredecessorCreatedAt != null &&
(!Number.isSafeInteger(rawExpectedPredecessorCreatedAt) || rawExpectedPredecessorCreatedAt < 0)
) {
startupTelemetry?.end('rejected');
return sendGenerationJson(
res,
400,
{
code: 'INVALID_GENERATION_PREDECESSOR',
error: 'expectedPredecessorCreatedAt must be a non-negative safe integer.',
},
generationProtocolVersion,
);
}
const expectedPredecessorCreatedAt = rawExpectedPredecessorCreatedAt;
const legacyRecoveredSteerId =
clientRequestId?.startsWith('steer-recovery:') === true
? clientRequestId.slice('steer-recovery:'.length)
: undefined;
const explicitRecoveredSteerId = req.body?.recoverySteerId;
const invalidExplicitRecoveryId =
explicitRecoveredSteerId != null &&
(typeof explicitRecoveredSteerId !== 'string' ||
!CLIENT_REQUEST_ID_PATTERN.test(explicitRecoveredSteerId));
const mismatchedRecoveryIds =
explicitRecoveredSteerId != null &&
legacyRecoveredSteerId != null &&
explicitRecoveredSteerId !== legacyRecoveredSteerId;
if (invalidExplicitRecoveryId || mismatchedRecoveryIds) {
startupTelemetry?.end('rejected');
return sendGenerationJson(
res,
400,
{
code: 'INVALID_RECOVERY_REQUEST',
error: 'recoverySteerId must identify exactly one parked steer source.',
},
generationProtocolVersion,
);
}
const recoveredSteerId = explicitRecoveredSteerId ?? legacyRecoveredSteerId;
const isRecoveredSteerRequest = recoveredSteerId != null;
const recoveryUserMessageId = req.body?.overrideUserMessageId;
const recoveredSteerPayload = isRecoveredSteerRequest
? buildRecoveredSteerPayload(text, req.body?.files)
: undefined;
/** A recovered steer is handed off as a new ordinary user turn. Edit,
* regenerate, continue, and arbitrary override-id shapes can reuse an
* existing user row (or deliberately skip its save); consuming the parked
* source from one of those shapes would therefore erase the only durable
* copy of the recovered words without proving that a new user row contains
* them. The source steer id itself is the one permitted user-row override:
* retries intentionally upsert that stable recovery row while each
* generation attempt uses a fresh clientRequestId. */
if (
isRecoveredSteerRequest &&
(!clientRequestId ||
!recoveredSteerId ||
!recoveredSteerPayload ||
!!isRegenerate ||
!!isContinued ||
editedContent != null ||
overrideParentMessageId != null ||
editedResponseMessageId != null ||
(recoveryUserMessageId != null && recoveryUserMessageId !== recoveredSteerId) ||
!!req.body?.overrideConvoId)
) {
startupTelemetry?.end('rejected');
return sendGenerationJson(
res,
400,
{
code: 'INVALID_RECOVERY_REQUEST',
error: 'A recovered steer must be submitted as a new user turn.',
},
generationProtocolVersion,
);
}
if (isRecoveredSteerRequest && recoveryUserMessageId === recoveredSteerId) {
/** BaseClient treats a bare override id as an already-persisted row and
* skips its save. Recovery instead needs an idempotent upsert: preserve
* the source-derived row id while explicitly selecting save index zero. */
req.body.overrideUserMessageId = `${recoveredSteerId}${Constants.COMMON_DIVIDER}0`;
}
const isNewConvo = !reqConversationId || reqConversationId === 'new';
const scheduledNewConversationId =
isScheduledFire && typeof req.body?.newConversationId === 'string'
? req.body.newConversationId
: null;
let conversationId = reqConversationId;
if (isNewConvo) {
conversationId =
scheduledNewConversationId ??
(typeof clientRequestId === 'string' && clientRequestId.length > 0
? uuidv5(`${userId}:${clientRequestId}`, NEW_CONVERSATION_IDEMPOTENCY_NAMESPACE)
: crypto.randomUUID());
}
const conversationAnchorPromise = resolveConversationCreatedAt({
userId,
conversationId,
isNewConvo,
conversation: Object.prototype.hasOwnProperty.call(req, 'resolvedConversation')
? req.resolvedConversation
: undefined,
});
const isTriggerContinuation =
req._isAgentTrigger === true && !isNewConvo && parentMessageId !== Constants.NO_PARENT;
if (
await isUnpersistedPreliminaryParent({
userId,
conversationId: reqConversationId,
parentMessageId,
getMessages,
})
) {
if (isTriggerContinuation) {
let parentJob;
try {
parentJob = await GenerationJobManager.getJob(conversationId);
} catch (error) {
logger.warn('[ResumableAgentController] Trigger parent lookup failed', error);
res.set('Retry-After', '1');
startupTelemetry?.end('rejected');
return sendGenerationJson(
res,
503,
{ code: 'PARENT_STATE_UNAVAILABLE', error: 'Parent generation state is unavailable.' },
generationProtocolVersion,
);
}
if (
parentJob != null &&
liveJobBelongsToRequester(parentJob, req.user) &&
(parentJob.status === 'running' ||
parentJob.status === 'requires_action' ||
parentJob.metadata?.terminalPersistencePending === true) &&
!(
typeof clientRequestId === 'string' &&
parentJob.metadata?.idempotencyClientRequestId === clientRequestId
)
) {
startupTelemetry?.end('rejected');
return rejectPreliminaryParentMessageId(res, generationProtocolVersion);
}
startupTelemetry?.end('rejected');
return rejectMissingTriggerParentMessageId(res, generationProtocolVersion);
}
startupTelemetry?.end('rejected');
return rejectPreliminaryParentMessageId(res, generationProtocolVersion);
}
/** When to generate the conversation title. `immediate` (default) fires title
* generation in parallel with the response, from the user's first message;
* `final` defers it until the full response completes (legacy behavior).
* Resolved from the agent's actual endpoint once the client is initialized. */
let titleTiming = 'immediate';
// Generate conversationId upfront if not provided - streamId === conversationId always
// Treat "new" as a placeholder that needs a real UUID (frontend may send "new" for new convos)
const streamId = conversationId;
req.body.conversationId = conversationId;
/** A durable continuation trigger appends below a completed parent response. If
* that response belongs to a still-running or paused generation, admitting
* another generation on the same conversation stream would replace it.
* Defer without claiming the continuation idempotency key so the delivery engine
* can retry after the parent reaches a terminal state. */
if (isTriggerContinuation) {
let parentJob;
try {
parentJob = await GenerationJobManager.getJob(streamId);
} catch (error) {
logger.warn('[ResumableAgentController] Trigger continuation parent lookup failed', error);
res.set('Retry-After', '1');
startupTelemetry?.end('rejected');
return sendGenerationJson(
res,
503,
{
code: 'PARENT_STATE_UNAVAILABLE',
error: 'Parent generation state is temporarily unavailable.',
},
generationProtocolVersion,
);
}
if (
parentJob != null &&
liveJobBelongsToRequester(parentJob, req.user) &&
(parentJob.status === 'running' ||
parentJob.status === 'requires_action' ||
parentJob.metadata?.terminalPersistencePending === true) &&
!(
typeof clientRequestId === 'string' &&
parentJob.metadata?.idempotencyClientRequestId === clientRequestId
)
) {
res.set('Retry-After', '1');
startupTelemetry?.end('rejected');
return sendGenerationJson(
res,
409,
{ code: 'PARENT_NOT_READY', error: 'The parent generation has not settled yet.' },
generationProtocolVersion,
);
}
}
// Idempotency: a lost/reset start-generation response makes the client re-POST the
// identical payload, which would otherwise start a second fully-billed generation.
// Claim the submission's clientRequestId before creating the job so a retry attaches
// to the original stream instead of spawning a duplicate. Runs before the concurrency
// check so a deduped retry is never counted against the limiter. Once a
// stable id is present, an ambiguous store outcome must fail closed.
let ownedIdempotencyClaim = null;
if (clientRequestId) {
let claim = null;
try {
claim = await GenerationJobManager.claimGeneration(
userId,
clientRequestId,
streamId,
conversationId,
generationProtocolVersion,
);
} catch (err) {
logger.error(
'[ResumableAgentController] Idempotency claim outcome is unknown; asking the client to retry',
err,
);
res.set('Retry-After', '1');
startupTelemetry?.end('deduplicated');
return sendGenerationJson(
res,
503,
{
code: 'SERVER_NOT_READY',
error: 'Generation ownership could not be confirmed. Please retry shortly.',
},
generationProtocolVersion,
);
}
if (claim?.existing != null) {
generationProtocolVersion = Math.min(
generationProtocolVersion,
claim.existing.generationProtocolVersion === GENERATION_PROTOCOL_V2
? GENERATION_PROTOCOL_V2
: 1,
);
}
const isLegacyTokenlessClaim =
claim?.source === 'legacy' && claim?.existing != null && claim.existing.claimToken == null;
const validClaim = isLegacyTokenlessClaim
? isValidLegacyGenerationClaim(claim.existing, streamId, isNewConvo)
: isValidGenerationClaim(claim?.existing, streamId, conversationId);
if (claim?.existing != null && !validClaim) {
logger.error('[ResumableAgentController] Invalid or miscorrelated idempotency claim');
res.set('Retry-After', '1');
startupTelemetry?.end('deduplicated');
return sendGenerationJson(
res,
503,
{
code: 'SERVER_NOT_READY',
error: 'Generation ownership could not be confirmed. Please retry shortly.',
},
generationProtocolVersion,
);
}
if (claim?.claimed && claim.existing?.claimToken) {
ownedIdempotencyClaim = claim.existing;
try {
const existingLiveGeneration = await GenerationJobManager.resumeClaimedGeneration(
userId,
clientRequestId,
streamId,
ownedIdempotencyClaim,
);
if (
existingLiveGeneration &&
isValidGenerationClaim(existingLiveGeneration, streamId, conversationId, true)
) {
// A fresh lease may have been negotiated under a different rollout
// cap than the still-live job it was atomically rebound to. The
// job's immutable protocol wins; echoing the fresh request's marker
// would make the client use v2-only recovery against a v1 run (or
// unnecessarily downgrade a v2 run).
generationProtocolVersion = Math.min(
generationProtocolVersion,
existingLiveGeneration.generationProtocolVersion === GENERATION_PROTOCOL_V2
? GENERATION_PROTOCOL_V2
: 1,
);
startupTelemetry?.end('deduplicated');
return sendGenerationJson(
res,
200,
{
streamId: existingLiveGeneration.streamId,
conversationId: existingLiveGeneration.conversationId,
generationCreatedAt: existingLiveGeneration.startedAt,
status: 'resumed',
},
generationProtocolVersion,
);
} else if (existingLiveGeneration) {
throw new Error('Live generation idempotency adoption returned invalid ownership');
}
} catch (err) {
logger.error('[ResumableAgentController] Live generation idempotency adoption failed', err);
res.set('Retry-After', '1');
startupTelemetry?.end('deduplicated');
return sendGenerationJson(
res,
503,
{
code: 'SERVER_NOT_READY',
error: 'Generation ownership changed. Please retry shortly.',
},
generationProtocolVersion,
);
}
} else if (claim?.existing) {
// A duplicate is confirmed. Attach to the original stream — and never fall through to
// a second generation, even if the job lookup hiccups.
const existingStreamId = claim.existing.streamId;
let liveJob;
try {
// Wait briefly for the winner to write the job record (it does so a few ms after
// claiming) so a still-live stream isn't handed back before its job exists.
liveJob = await waitForJobRecord(existingStreamId);
} catch (err) {
// Store hiccup while checking the job: ask the client to retry rather than starting
// a second generation for a request we know is a duplicate.
logger.error(
'[ResumableAgentController] Job lookup failed for an existing claim; asking the client to retry',
err,
);
res.set('Retry-After', '1');
startupTelemetry?.end('deduplicated');
return sendGenerationJson(
res,
503,
{
code: 'SERVER_NOT_READY',
error: 'Generation is still starting. Please retry shortly.',
},
generationProtocolVersion,
);
}
const claimAgeMs = Date.now() - (claim.existing.claimedAt ?? 0);
if (!liveJob && claim.existing.startedAt != null) {
// createJob marked this claim in the same transaction that installed
// the job. A now-missing record therefore represents an already-owned
// generation (usually fast completion + cleanup), never an abandoned
// pre-create lease that may be taken over and billed again. There is no
// attachable stream; the settled response refetches persisted history.
return sendSettledGeneration(
res,
existingStreamId,
claim.existing.conversationId,
startupTelemetry,
generationProtocolVersion,
);
}
if (!liveJob && isLegacyTokenlessClaim && claimAgeMs >= IDEMPOTENCY_STARTUP_GRACE_MS) {
/** A legacy owner cannot be fenced (its value has no token), so it is
* never safe to take over. Return its original stream on the legacy
* attach/refetch path: this covers fast completion without starting a
* second billed generation, while an abandoned pre-create claim ages
* out under the old server's bounded TTL. */
return sendSettledGeneration(
res,
existingStreamId,
claim.existing.conversationId,
startupTelemetry,
generationProtocolVersion,
);
}
if (!liveJob && claimAgeMs < IDEMPOTENCY_STARTUP_GRACE_MS) {
// The winner claimed but has not written the job yet (still between claim and
// createJob). Handing back the stream now would 404 and tear down the client while
// the winner goes on to generate and bill with no UI attached — ask the client to
// retry via the readiness path instead.
res.set('Retry-After', '1');
startupTelemetry?.end('deduplicated');
return sendGenerationJson(
res,
503,
{
code: 'SERVER_NOT_READY',
error: 'Generation is still starting. Please retry shortly.',
},
generationProtocolVersion,
);
}
if (liveJob) {
generationProtocolVersion = negotiateExistingGenerationProtocol(req, liveJob);
if (!liveJobBelongsToRequester(liveJob, req.user)) {
logger.error(
'[ResumableAgentController] Existing idempotency claim resolved to a foreign generation',
);
res.set('Retry-After', '1');
startupTelemetry?.end('deduplicated');
return sendGenerationJson(
res,
503,
{
code: 'SERVER_NOT_READY',
error: 'Generation ownership could not be confirmed. Please retry shortly.',
},
generationProtocolVersion,
);
}
const liveClientRequestId = liveJob.metadata?.idempotencyClientRequestId;
const startedAt = claim.existing.startedAt;
if (liveJob.metadata?.terminalPersistencePending === true) {
/** The terminal owner has claimed the outcome but has not yet
* finished the required persistence hook. Do not let a duplicate
* start refetch history until that single-winner publication is
* finalized (or stale-pending recovery publishes failure). */
res.set('Retry-After', '1');
startupTelemetry?.end('deduplicated');
return sendGenerationJson(
res,
503,
{
code: 'SERVER_NOT_READY',
error: 'Generation is finalizing. Please retry shortly.',
},
generationProtocolVersion,
);
}
const terminalWithoutPayload =
['complete', 'error', 'aborted'].includes(liveJob.status) &&
!liveJob.finalEvent &&
!liveJob.error;
if (terminalWithoutPayload) {
/** A terminal CAS can precede its required DB save and durable FINAL
* by a narrow window. Returning an attachable/settled success here
* lets the retry refetch before persistence is complete. Keep the
* duplicate on the readiness path until the owner publishes its
* terminal payload (or cleanup makes the job disappear). */
res.set('Retry-After', '1');
startupTelemetry?.end('deduplicated');
return sendGenerationJson(
res,
503,
{
code: 'SERVER_NOT_READY',
error: 'Generation is finalizing. Please retry shortly.',
},
generationProtocolVersion,
);
}
const replacedGeneration =
(startedAt != null && liveJob.createdAt !== startedAt) ||
(liveClientRequestId != null && liveClientRequestId !== clientRequestId);
if (replacedGeneration) {
// streamId === conversationId, so a later turn reuses the same route.
// Never pair this stale POST's optimistic submission with that newer
// job's SSE snapshot. If the replacement is still active, distinguish
// it from an ordinary settled retry so the client hands off to the
// authoritative B submission instead of going idle and starting C.
if (liveJob.status === 'running' || liveJob.status === 'requires_action') {
startupTelemetry?.end('deduplicated');
if (generationProtocolVersion < GENERATION_PROTOCOL_V2) {
return sendGenerationJson(
res,
409,
{ code: 'RUN_REPLACED' },
generationProtocolVersion,
);
}
return sendGenerationJson(
res,
200,
{
streamId: existingStreamId,
conversationId: claim.existing.conversationId,
generationCreatedAt: liveJob.createdAt,
status: 'replaced',
},
generationProtocolVersion,
);
}
return sendSettledGeneration(
res,
existingStreamId,
claim.existing.conversationId,
startupTelemetry,
generationProtocolVersion,
);
}
if (liveClientRequestId == null && !isLegacyTokenlessClaim) {
// A syntactically valid claim plus an uncorrelated live job is
// outcome-ambiguous (legacy/corrupt/partially written state). Attaching
// risks cross-wiring two submissions; starting risks double billing.
res.set('Retry-After', '1');
startupTelemetry?.end('deduplicated');
return sendGenerationJson(
res,
503,
{
code: 'SERVER_NOT_READY',
error: 'Generation ownership could not be confirmed. Please retry shortly.',
},
generationProtocolVersion,
);
}
logger.debug('[ResumableAgentController] Deduped retried start-generation request', {
userId,
clientRequestId,
streamId: existingStreamId,
});
startupTelemetry?.end('deduplicated');
return sendGenerationJson(
res,
200,
{
streamId: existingStreamId,
conversationId: claim.existing.conversationId,
generationCreatedAt: liveJob.createdAt,
status: 'resumed',
},
generationProtocolVersion,
);
}
// The creator held the claim beyond the startup grace but never made a
// job. Atomically take over its lease; createJob verifies this token in
// the same Redis transaction as job creation, so the abandoned winner
// can no longer wake up and start a second generation.
const takeover = await GenerationJobManager.takeoverGeneration(
userId,
clientRequestId,
existingStreamId,
claim.existing,
).catch((err) => {
logger.error('[ResumableAgentController] Stale idempotency takeover failed', err);
return null;
});
if (
!takeover?.claimed ||
!isValidGenerationClaim(takeover.existing, streamId, conversationId)
) {
res.set('Retry-After', '1');
startupTelemetry?.end('deduplicated');
return sendGenerationJson(
res,
503,
{
code: 'SERVER_NOT_READY',
error: 'Generation ownership changed. Please retry shortly.',
},
generationProtocolVersion,
);
}
ownedIdempotencyClaim = takeover.existing;
} else {
// A malformed/unreadable existing claim is outcome-ambiguous. Starting
// anyway would turn a store parsing failure into duplicate generation.
res.set('Retry-After', '1');
startupTelemetry?.end('deduplicated');
return sendGenerationJson(
res,
503,
{
code: 'SERVER_NOT_READY',
error: 'Generation ownership could not be confirmed. Please retry shortly.',
},
generationProtocolVersion,
);
}
}
const scheduleConcurrencyExempt = exemptFromConcurrencyLimiter(req);
req._scheduleConcurrencyExempt = scheduleConcurrencyExempt;
if (!scheduleConcurrencyExempt) {
const { allowed, pendingRequests, limit } = await checkAndIncrementPendingRequest(userId);
if (!allowed) {
if (ownedIdempotencyClaim) {
await GenerationJobManager.releaseGeneration(
userId,
clientRequestId,
streamId,
ownedIdempotencyClaim,
).catch(() => {});
}
const violationInfo = getViolationInfo(pendingRequests, limit);
await logViolation(req, res, ViolationTypes.CONCURRENT, violationInfo, violationInfo.score);
startupTelemetry?.end('rejected');
return sendGenerationJson(res, 429, violationInfo, generationProtocolVersion);
}
}
startupTelemetry?.mark('request_admitted');
let client = null;
let jobCreatedAt;
let providerExecutionId;
let scheduleTerminalOutcomeRecorded = false;
const settleScheduledRun = async ({ status, error, clearConversationId = false }) => {
if (!scheduleId) {
return true;
}
if (status !== 'requires_action' && scheduleTerminalOutcomeRecorded) {
return true;
}
const recorded = await recordScheduleOutcome({
scheduleId,
scheduledFor,
streamId,
jobCreatedAt,
status,
conversationId,
clearConversationId,
error,
});
if (recorded && status !== 'requires_action') {
scheduleTerminalOutcomeRecorded = true;
}
return recorded;
};
try {
logger.debug(`[ResumableAgentController] Creating job`, {
streamId,
conversationId,
reqConversationId,
userId,
});
const endpointIconURL = getEndpointIconURL(req, endpointOption);
const responseModel = getAgentResponseModel(req, endpointOption);
const preliminaryUserMessage = getPreliminaryUserMessage(req.body, conversationId);
const preliminaryResponseMessageId = getPreliminaryResponseMessageId(req.body);
const job = await GenerationJobManager.createJob(streamId, userId, conversationId, {
startupTelemetry,
...(recoveredSteerId && { recoveredSteerId }),
...(recoveredSteerPayload && { recoveredSteerPayload }),
...(expectedPredecessorCreatedAt != null && { expectedPredecessorCreatedAt }),
...(isTriggerContinuation && { rejectActivePredecessor: true }),
...(ownedIdempotencyClaim?.claimToken && {
idempotencyClientRequestId: clientRequestId,
idempotencyClaimToken: ownedIdempotencyClaim.claimToken,
}),
initialMetadata: {
conversationId,
generationProtocolVersion,
endpoint: endpointOption.endpoint,
iconURL: endpointIconURL,
model: responseModel,
// Recorded HERE because this process owns the generation: the steer
// route may land on a different replica whose own SDK probe would
// answer for the wrong process during a rolling deploy.
preemptCapable: isSteerPreemptSupported(),
// Persist the originating agent so a HITL resume can refuse to rebuild this
// paused run on a different agent (see resume.js).
agent_id: endpointOption.agent_id ?? req.body?.agent_id,
// Persist temporary-chat state so a HITL resume keeps the resumed response
// non-persisted instead of trusting the resume request to re-send the flag.
isTemporary: req.body?.isTemporary,
...(scheduleId
? {
scheduleId,
scheduledFor,
preserveForScheduleReconcile: true,
...(Number.isSafeInteger(scheduleConfigRevision) && {
scheduleConfigRevision,
}),
...(req._isManualScheduledFire === true && { scheduleManual: true }),
}
: {}),
responseMessageId: preliminaryResponseMessageId,
userMessage: preliminaryUserMessage,
},
});
startupTelemetry?.mark('job_created');
generationProtocolVersion = negotiateExistingGenerationProtocol(req, job);
jobCreatedAt = job.createdAt; // Capture creation time to detect job replacement
providerExecutionId = job.metadata?.providerExecutionId;
/** Authentication can precede a slow admission path. Recheck the durable
* account-deletion fence after the job is committed but before execution
* starts. This ordering closes both sides of the race for ordinary and
* trigger-scoped sessions: a fence that wins first rejects this run; a
* fence that starts after this read must observe the already-created job
* in account deletion's active-generation drain. */
if (!(await isAgentTriggerPrincipalActive(userId))) {
throw Object.assign(new Error('Account deletion is in progress'), {
code: 'ACCOUNT_DELETION_IN_PROGRESS',
status: 409,
});
}
if (
scheduleId &&
!(await isScheduleLive(scheduleId, scheduleConfigRevision, {
automatic: req._isManualScheduledFire !== true,
policy: true,
// The occurrence's OWN recorded scope, exactly as the resume path passes it.
// The run row is reserved before this loopback request is dispatched, so a pin
// introduced while the request sat queued must not be validated in place of the
// destination this occurrence's envelope was already built with.
scheduledFor,
}))
) {
throw Object.assign(new Error('This scheduled occurrence is no longer active'), {
code: 'SCHEDULE_NO_LONGER_ACTIVE',
status: 409,
});
}
if (
providerExecutionId &&
!(await GenerationJobManager.beginProviderExecution(
streamId,
jobCreatedAt,
providerExecutionId,
))
) {
throw Object.assign(new Error('Generation stopped before provider startup'), {
code: 'RUN_REPLACED',
status: 409,
});
}
acceptAgentStartupTelemetry(req, streamId);
startupTelemetry?.mark('metadata_persisted');
req._resumableStreamId = streamId;
getMCPRequestContext(req, undefined, { cleanupOnResponse: false });
let recoveredSteerCommitted = false;
const commitRecoveredSteer = async () => {
if (!recoveredSteerId || recoveredSteerCommitted) {
return;
}
if (client?.skipSaveUserMessage) {
throw new Error('Recovered steer cannot skip user message persistence');
}
const committed = await GenerationJobManager.steering.consumeRecovered(
streamId,
recoveredSteerId,
{ userId, tenantId: req.user?.tenantId },
jobCreatedAt,
);
if (!committed) {
throw new Error('Recovered steer could not be committed after message persistence');
}
recoveredSteerCommitted = true;
};
// Send JSON response IMMEDIATELY so client can connect to SSE stream
// This is critical: tool loading (MCP OAuth) may emit events that the client needs to receive
sendGenerationJson(
res,
200,
{ streamId, conversationId, generationCreatedAt: jobCreatedAt, status: 'started' },
generationProtocolVersion,
);
await attachConversationCreatedAt(req, conversationId, conversationAnchorPromise).then(() =>
startupTelemetry?.mark('conversation_resolved'),
);
// Note: We no longer use res.on('close') to abort since we send JSON immediately.
// The response closes normally after res.json(), which is not an abort condition.
// Abort handling is done through GenerationJobManager via the SSE stream connection.
// Track if partial response was already saved to avoid duplicates
let partialResponseSaved = false;
/**
* Listen for all subscribers leaving to save partial response.
* This ensures the response is saved to DB even if all clients disconnect
* while generation continues.
*
* Note: The messageId used here falls back to `${userMessage.messageId}_` if the
* actual response messageId isn't available yet. The final response save will
* overwrite this with the complete response using the same messageId pattern.
*/
job.emitter.on('allSubscribersLeft', async (aggregatedContent) => {
if (partialResponseSaved || !aggregatedContent || aggregatedContent.length === 0) {
return;
}
const persistableContent = filterPersistableAbortContent(aggregatedContent);
if (persistableContent.length === 0) {
logger.debug('[ResumableAgentController] No persistable content to save partial response');
return;
}
const resumeState = await GenerationJobManager.getResumeState(streamId, jobCreatedAt);
if (!resumeState?.userMessage) {
logger.debug('[ResumableAgentController] No user message to save partial response for');
return;
}
partialResponseSaved = true;
const responseConversationId = resumeState.conversationId || conversationId;
try {
const partialMessage = {
messageId: resumeState.responseMessageId || `${resumeState.userMessage.messageId}_`,
conversationId: responseConversationId,
parentMessageId: resumeState.userMessage.messageId,
sender: client?.sender ?? 'AI',
content: persistableContent,
unfinished: true,
error: false,
isCreatedByUser: false,
user: userId,
endpoint: endpointOption.endpoint,
iconURL: resumeState.iconURL || endpointIconURL,
model: resumeState.model || responseModel,
};
if (req.body?.agent_id) {
partialMessage.agent_id = req.body.agent_id;
}
const savePartialMessage = () =>
saveMessage(
{
userId,
isTemporary: req?.body?.isTemporary,
interfaceConfig: req?.config?.interfaceConfig,
},
partialMessage,
{
context: 'api/server/controllers/agents/request.js - partial response on disconnect',
},
);
const savedPartialMessage = tenantId
? await tenantStorage.run({ tenantId, userId }, savePartialMessage)
: await savePartialMessage();
if (!savedPartialMessage) {
throw new Error('Partial response could not be persisted after disconnect');
}
logger.debug(
`[ResumableAgentController] Saved partial response for ${streamId}, content parts: ${persistableContent.length}`,
);
} catch (error) {
logger.error('[ResumableAgentController] Error saving partial response:', error);
// Reset flag so we can try again if subscribers reconnect and leave again
partialResponseSaved = false;
}
});
/** @type {{ client: TAgentClient; userMCPAuthMap?: Record<string, Record<string, string>> }} */
const result = await initializeClient({
req,
res,
endpointOption,
// Use the job's abort controller signal - allows abort via GenerationJobManager.abortJob()
signal: job.abortController.signal,
jobCreatedAt,
checkpointNamespace: job.metadata?.checkpointNamespace,
});
startupTelemetry?.mark('client_initialized');
client = result.client;
/** Request-shape validation rejects every known edit/regenerate path, but
* the client owns the final persistence decision. Fail closed if a future
* or provider-specific path still derives skip-save for a recovered turn;
* consuming its parked source would otherwise erase the only durable copy
* of the user's words. Re-checked inside commitRecoveredSteer in case a
* client mutates the flag while sending. */
if (recoveredSteerId && client?.skipSaveUserMessage) {
throw new Error('Recovered steer cannot skip user message persistence');
}
if (job.abortController.signal.aborted) {
await GenerationJobManager.completeJob(
streamId,
'Request aborted during initialization',
jobCreatedAt,
).catch((completeErr) => {
logger.warn(
'[ResumableAgentController] completeJob failed after initialization abort',
completeErr,
);
});
await settleScheduledRun({
status: 'interrupted',
error: 'Request aborted during initialization',
clearConversationId: job.createdEventEmitted !== true,
});
startupTelemetry?.end('aborted');
try {
await finishResumableRequest(req, userId);
} finally {
if (client) {
disposeClient(client);
}
client = null;
if (providerExecutionId) {
await GenerationJobManager.markProviderExecutionDrained?.(
streamId,
jobCreatedAt,
providerExecutionId,
).catch((drainError) => {
logger.warn(
'[ResumableAgentController] Failed to record initialization-abort provider drain',
drainError,
);
});
}
}
return;
}
// Tag the client with THIS generation's identity so HITL terminal side-effects
// (pause CAS, checkpoint prune) can tell whether a newer request has since replaced
// this job on the same conversationId before acting on it.
client.jobCreatedAt = jobCreatedAt;
// Resolve title timing from the public agents endpoint first, then fall
// back to the agent's actual backing provider/custom endpoint.
titleTiming = resolveTitleTiming({
appConfig: req.config,
endpoint: [endpointOption?.endpoint, client?.options?.agent?.endpoint],
});
if (client?.sender) {
void GenerationJobManager.updateMetadata(
streamId,
{ sender: client.sender },
jobCreatedAt,
).catch((err) => {
logger.warn('[ResumableAgentController] Failed to persist response sender', err);
});
}
// Store reference to client's contentParts - graph will be set when run is created
if (client?.contentParts) {
GenerationJobManager.setContentParts(streamId, client.contentParts, jobCreatedAt);
}
let userMessage;
const getReqData = (data = {}) => {
if (data.userMessage) {
userMessage = data.userMessage;
}
// conversationId is pre-generated, no need to update from callback
};
let immediateTitlePromise = null;
let trailingWritePromise = null;
let backgroundClientCleanupScheduled = false;
let terminalClaim = null;
let terminalClaimFinished = false;
let terminalPersistenceChecked = false;
let terminalWasAborted = false;
let preemptIncomplete = false;
/** A pause-row write failure is terminalized through the exact action/epoch
* barrier. Once that path starts, neither generic background error handler
* may call completeJob: the pause may already have been replaced by a newer
* action or generation by the time the persistence failure is observed. */
let pausePersistenceFailed = false;
let pausePersistenceFailureFinalized = false;
const finishOwnedTerminalClaim = async () => {
if (!terminalClaim || terminalClaimFinished) {
return;
}
try {
await GenerationJobManager.finishTerminalJob(terminalClaim);
} finally {
terminalClaimFinished = true;
}
};
/** Runs inside BaseClient immediately before it can start the completed
* response write. A lost claim returns false, and BaseClient skips that
* stale `unfinished:false` write entirely. The fallback invocation below
* supports test/custom clients that do not derive from BaseClient. */
const claimBeforeResponsePersistence = async () => {
if (terminalPersistenceChecked) {
return terminalClaim != null;
}
terminalPersistenceChecked = true;
if (client?.pendingApproval) {
// AgentClient installed a durable pause-persistence barrier in the
// running→requires_action CAS. BaseClient must not start its ordinary
// `unfinished:false` response write; the HITL branch persists the
// partial row as unfinished before releasing that barrier.
return false;
}
terminalWasAborted = job.abortController.signal.aborted;
const preemptStats = client?.run?.getPreemptStats?.();
preemptIncomplete =
(preemptStats?.emptyBoundaries ?? 0) > 0 ||
client?.run?.getHaltReason?.() === 'preempt_incomplete';
terminalClaim = await GenerationJobManager.claimTerminalJob(
streamId,
terminalWasAborted ? 'aborted' : 'complete',
undefined,
jobCreatedAt,
{ persistencePending: true },
);
return terminalClaim != null;
};
const disposeBackgroundClient = () => {
if (backgroundClientCleanupScheduled) {
return;
}
backgroundClientCleanupScheduled = true;
if (immediateTitlePromise) {
immediateTitlePromise.finally(() => {
if (client) {
disposeClient(client);
}
});
} else if (client) {
disposeClient(client);
}
};
// Start background generation immediately. The stream layer buffers and persists events
// until an SSE subscriber attaches, so generation no longer waits on subscriber readiness.
const startGeneration = async () => {
/** Immediate-mode title generation runs in parallel with the response, so
* the conversation row may not exist when the title resolves. `convoReady`
* resolves once the response (and thus the conversation) has been saved,
* gating the title's `saveConvo`. Declared here so both the success tail
* and the catch block can settle it and gate `disposeClient` on the title. */
let titleEventPromise = null;
let acceptsTitleEvents = true;
let resolveConvoReady;
const convoReady = new Promise((resolve) => {
resolveConvoReady = resolve;
});
/** Dedicated controller so a user Stop (or a replaced stream) cancels the
* in-flight title — kept separate from `job.abortController`, which
* `completeJob` also aborts on *successful* completion and would otherwise
* cancel a title that is merely slower than a short response. */
const titleAbortController = new AbortController();
/** Separate from `titleAbortController`: a user Stop cancels the in-flight
* title model call but keeps a title that already finished generating.
* Only a superseded/failed stream aborts this to discard such a title so it
* cannot clobber the conversation now owned by the newer run. */
const titleDiscardController = new AbortController();
const abortTitleOnJobAbort = () => titleAbortController.abort();
if (job.abortController.signal.aborted) {
titleAbortController.abort();
} else {
job.abortController.signal.addEventListener('abort', abortTitleOnJobAbort, { once: true });
}
const titleEligible =
addTitle && parentMessageId === Constants.NO_PARENT && isNewConvo && !req.body?.isTemporary;
const emitTitleEvent = ({ conversationId: titleConversationId, title }) => {
titleEventPromise = (async () => {
if (!acceptsTitleEvents || titleAbortController.signal.aborted) {
return;
}
const currentJob = await GenerationJobManager.getJob(streamId);
if (!currentJob || currentJob.createdAt !== jobCreatedAt) {
return;
}
if (titleAbortController.signal.aborted) {
return;
}
await GenerationJobManager.emitChunk(
streamId,
{
event: 'title',
data: {
conversationId: titleConversationId,
title,
},
},
{ expectedCreatedAt: jobCreatedAt },
);
})().catch((err) => {
logger.error('[ResumableAgentController] Error emitting title event', err);
});
return titleEventPromise;
};
try {
const onStart = (userMsg, respMsgId, _isNewConvo) => {
userMessage = userMsg;
// Store userMessage and responseMessageId upfront for resume capability
GenerationJobManager.updateMetadata(
streamId,
{
responseMessageId: respMsgId,
userMessage: {
messageId: userMsg.messageId,
parentMessageId: userMsg.parentMessageId,
conversationId: userMsg.conversationId,
text: userMsg.text,
quotes: userMsg.quotes,
// Persist the turn's uploaded files here (authoritative job metadata) so a
// HITL resume sources them from the job, not the user DB row — which the
// approval prompt can race (the row save may still be in flight when a fast
// /resume reads it). Without this an approved tool run can rebuild without the
// paused turn's files.
...(Array.isArray(req.body?.files) &&
req.body.files.length > 0 && { files: req.body.files }),
// Skill selections aren't on `userMsg` yet at onStart (BaseClient adds them
// later), so source them from the request — otherwise this update overwrites
// the preliminary metadata and a HITL-resumed turn loses its skill pills.
...(Array.isArray(req.body?.manualSkills) &&
req.body.manualSkills.length > 0 && { manualSkills: req.body.manualSkills }),
...(Array.isArray(req.body?.alwaysAppliedSkills) &&
req.body.alwaysAppliedSkills.length > 0 && {
alwaysAppliedSkills: req.body.alwaysAppliedSkills,
}),
},
},
jobCreatedAt,
).catch((err) => {
logger.error('[ResumableAgentController] Failed to persist start metadata', err);
});
GenerationJobManager.emitChunk(
streamId,
{
created: true,
// Skill selections aren't on `userMessage` yet at onStart (BaseClient adds
// them later), so attach them from the request — this is the message
// `trackUserMessage` persists as the authoritative job.metadata.userMessage,
// and it's what the live client renders the user bubble from.
message: {
...userMessage,
// Carry files so trackUserMessage (the authoritative writer) persists them on
// job.metadata.userMessage for a HITL resume (see the updateMetadata above).
...(Array.isArray(req.body?.files) &&
req.body.files.length > 0 && { files: req.body.files }),
...(Array.isArray(req.body?.manualSkills) &&
req.body.manualSkills.length > 0 && { manualSkills: req.body.manualSkills }),
...(Array.isArray(req.body?.alwaysAppliedSkills) &&
req.body.alwaysAppliedSkills.length > 0 && {
alwaysAppliedSkills: req.body.alwaysAppliedSkills,
}),
},
streamId,
},
{ expectedCreatedAt: jobCreatedAt },
).catch((err) => {
logger.error('[ResumableAgentController] Failed to queue created event', err);
});
};
const messageOptions = {
user: userId,
onStart,
getReqData,
isContinued,
isRegenerate,
editedContent,
conversationId,
parentMessageId,
abortController: job.abortController,
overrideParentMessageId,
isEdited: !!editedContent,
beforeResponsePersistence: claimBeforeResponsePersistence,
userMCPAuthMap: result.userMCPAuthMap,
responseMessageId: editedResponseMessageId,
progressOptions: {
res: {
write: () => true,
end: () => {},
headersSent: false,
writableEnded: false,
},
},
};
const sendPromise = client.sendMessage(text, messageOptions);
if (titleEligible && titleTiming === 'immediate') {
immediateTitlePromise = addTitle(req, {
text: text || getAttachmentTitleText(req.body.files),
conversationId,
client,
immediate: true,
convoReady,
signal: titleAbortController.signal,
discardSignal: titleDiscardController.signal,
onTitleGenerated: emitTitleEvent,
}).catch((err) => {
logger.error('[ResumableAgentController] Error in immediate title generation', err);
});
}
const response = await sendPromise;
// HITL: the turn paused for human review (see AgentClient.handleRunInterrupt).
// The job is already `requires_action` with the pending action persisted and
// emitted to the client; the resume route owns finishing this turn. Settle and
// verify the required unfinished history, then tear down without publishing a
// terminal event or completing a successfully persisted paused job.
if (client?.pendingApproval) {
if (response?.databasePromise) {
try {
await response.databasePromise;
} catch (dbErr) {
logger.error(
'[ResumableAgentController] Error settling databasePromise on HITL pause',
dbErr,
);
}
delete response.databasePromise;
}
const pauseActionId = client.pendingApproval.actionId;
const pauseCreatedAt = client.jobCreatedAt ?? jobCreatedAt;
const ownsPausePersistence = await GenerationJobManager.approvals.ownsPausePersistence(
streamId,
pauseActionId,
pauseCreatedAt,
);
if (ownsPausePersistence) {
try {
/** BaseClient awaits its first user/conversation write before the
* pause hook, but deliberately swallows a failed/falsy user save
* and may still record the id locally. Re-save idempotently for
* every ordinary user turn before exposing the approval. */
if (!client?.skipSaveUserMessage) {
if (!userMessage) {
throw new Error('User message was unavailable before HITL pause');
}
if (
typeof client.saveMessageToDatabase === 'function' &&
typeof client.getSaveOptions === 'function'
) {
/** Retry through BaseClient so a failure before its original
* saveConvo is repaired along with the message row. Direct
* saveMessage alone cannot recreate that conversation. */
const savedUserTurn = await client.saveMessageToDatabase(
userMessage,
client.getSaveOptions(),
userId,
);
if (!savedUserTurn?.message) {
throw new Error('User message could not be persisted before HITL pause');
}
if (!client.skipSaveConvo && !savedUserTurn.conversation) {
throw new Error('Conversation could not be persisted before HITL pause');
}
} else {
// Custom clients used by integrations/tests may not inherit BaseClient.
const savedUserMessage = await saveMessage(
{
userId,
isTemporary: req?.body?.isTemporary,
interfaceConfig: req?.config?.interfaceConfig,
},
userMessage,
{
context:
'api/server/controllers/agents/request.js - user message before HITL pause',
},
);
if (!savedUserMessage) {
throw new Error('User message could not be persisted before HITL pause');
}
}
}
if (!response?.messageId) {
throw new Error('Response message was unavailable before HITL pause');
}
const savedResponseMessage = await saveMessage(
{
userId,
isTemporary: req?.body?.isTemporary,
interfaceConfig: req?.config?.interfaceConfig,
},
{
...response,
endpoint: endpointOption.endpoint,
unfinished: true,
user: userId,
},
{
context:
'api/server/controllers/agents/request.js - HITL pause (persist unfinished)',
},
);
if (!savedResponseMessage) {
throw new Error('Paused response could not be persisted as unfinished');
}
await commitRecoveredSteer();
} catch (pausePersistenceError) {
pausePersistenceFailed = true;
try {
pausePersistenceFailureFinalized =
(await GenerationJobManager.failPausePersistence(
streamId,
pauseActionId,
pausePersistenceError?.message ?? 'Pause persistence failed',
pauseCreatedAt,
)) === true;
} catch (failError) {
logger.error(
`[ResumableAgentController] Failed to terminalize pause persistence error for ${streamId}`,
failError,
);
}
if (pausePersistenceFailureFinalized) {
/** Namespaced checkpoints belong exclusively to this epoch,
* so the exact pause-failure CAS winner can safely remove the
* now-unresumable graph state. Legacy shared namespaces are
* left to their guarded/TTL cleanup path. */
const checkpointNamespace = job.metadata?.checkpointNamespace;
if (typeof checkpointNamespace === 'string' && checkpointNamespace !== '') {
try {
await deleteAgentCheckpoint(
conversationId,
req.config?.endpoints?.[EModelEndpoint.agents]?.checkpointer,
undefined,
{ checkpointNamespace },
);
} catch (checkpointError) {
logger.error(
`[ResumableAgentController] Failed to prune checkpoint after pause persistence error for ${streamId}`,
checkpointError,
);
}
}
} else if (pausePersistenceFailureFinalized === false) {
logger.warn(
`[ResumableAgentController] Skipping stale pause persistence failure — ${streamId} no longer owns its barrier`,
);
}
throw pausePersistenceError;
}
const released = await GenerationJobManager.approvals.finishPausePersistence(
streamId,
pauseActionId,
pauseCreatedAt,
);
if (!released) {
logger.warn(
`[ResumableAgentController] Pause persistence barrier changed before release: ${streamId}`,
);
}
// The pause projection is what moves the run row off `started` and frees its
// GLOBAL capacity slot. recordScheduleOutcome already retried it; a `false`
// here means every attempt failed, leaving the row `started` while the job
// sits `requires_action`. Surface it — the armed engine's reconciler replays
// this state, and the clustered sweep now converges it too, but a silent drop
// gave neither a reason to look.
if (!(await settleScheduledRun({ status: 'requires_action' }))) {
logger.error(
`[ResumableAgentController] Failed to project the scheduled pause for ${streamId}; run stays active until reconciliation replays it`,
);
}
} else {
logger.debug(
`[ResumableAgentController] Skipping stale pause persistence — ${streamId} no longer owns its barrier`,
);
}
titleAbortController.abort();
acceptsTitleEvents = false;
resolveConvoReady();
// handleRunInterrupt already released the concurrency slot the moment it paused
// (so a fast /resume isn't 429'd); only release here if that didn't happen.
// Always run the MCP request-context cleanup.
await cleanupMCPRequestContextForReq(req);
if (!client?.pendingRequestReleased && req._scheduleConcurrencyExempt !== true) {
await decrementPendingRequest(userId);
}
if (client) {
disposeClient(client);
}
logger.debug(
`[ResumableAgentController] Turn paused for approval; awaiting resume: ${streamId}`,
);
startupTelemetry?.end('paused');
return;
}
// BaseClient invokes this before starting its response write. Custom
// clients/tests may return a database promise directly, so keep the
// controller-side fallback before awaiting that promise.
await claimBeforeResponsePersistence();
const endpoint = endpointOption.endpoint;
response.endpoint = endpoint;
const databasePromise = response.databasePromise;
delete response.databasePromise;
const { conversation: convoData = {} } = await databasePromise;
const conversation = { ...convoData };
conversation.title =
conversation && !conversation.title ? null : conversation?.title || 'New Chat';
if (!terminalClaim) {
/** Stop/replacement won before the response persistence hook. The
* BaseClient contract skipped its completed response write; cancel
* title work and leave terminal publication/persistence to the
* actual winner. */
titleAbortController.abort();
titleDiscardController.abort();
job.abortController.signal.removeEventListener('abort', abortTitleOnJobAbort);
acceptsTitleEvents = false;
resolveConvoReady();
await finishResumableRequest(req, userId);
disposeBackgroundClient();
startupTelemetry?.end(job.abortController.signal.aborted ? 'aborted' : 'replaced');
return;
}
if (req.body.files && Array.isArray(client.options.attachments)) {
const files = buildMessageFiles(req.body.files, client.options.attachments);
if (files.length > 0) {
userMessage.files = files;
}
delete userMessage.image_urls;
}
const shouldGenerateTitle =
addTitle &&
parentMessageId === Constants.NO_PARENT &&
isNewConvo &&
!terminalWasAborted &&
!preemptIncomplete;
// Save user message BEFORE sending final event to avoid race condition
// where client refetch happens before database is updated
const reqCtx = {
userId: req?.user?.id,
isTemporary: req?.body?.isTemporary,
interfaceConfig: req?.config?.interfaceConfig,
};
if (!client.skipSaveUserMessage) {
if (!userMessage) {
throw new Error('User message was unavailable before terminal persistence');
}
const savedUserMessage = await saveMessage(reqCtx, userMessage, {
context: 'api/server/controllers/agents/request.js - resumable user message',
});
if (!savedUserMessage) {
throw new Error('User message could not be persisted before terminal publication');
}
}
// Only consume the parked recovery source after the explicit user-row
// write above succeeds. `response.databasePromise` alone is insufficient:
// BaseClient intentionally swallows a failed first user-message save.
await commitRecoveredSteer();
// CRITICAL: Save response message BEFORE emitting final event.
// This prevents race conditions where the client sends a follow-up message
// before the response is saved to the database, causing orphaned parentMessageIds.
/** BaseClient can add the id to savedMessageIds even when its model-layer
* save resolved falsy. Re-save the terminal row idempotently and require
* the returned durable row before publishing the normal FINAL. */
const responseIsUnfinished = terminalWasAborted || preemptIncomplete;
const savedResponseMessage = await saveMessage(
reqCtx,
{
...response,
user: userId,
unfinished: responseIsUnfinished,
},
{
context: responseIsUnfinished
? 'api/server/controllers/agents/request.js - terminal response unfinished'
: 'api/server/controllers/agents/request.js - resumable response end',
},
);
if (!savedResponseMessage) {
throw new Error(
responseIsUnfinished
? 'Terminal response could not be persisted as unfinished'
: 'Response message could not be persisted before terminal publication',
);
}
// If the user stopped this turn — or an empty preempt boundary truncated
// it, which persists under the same honest `unfinished` contract — cancel
// the title BEFORE unblocking its persistence wait; otherwise resolving
// `convoReady` lets the title task resume and save before the later abort runs.
if (terminalWasAborted || preemptIncomplete) {
titleAbortController.abort();
} else {
job.abortController.signal.removeEventListener('abort', abortTitleOnJobAbort);
}
// The conversation row now exists and this stream is authoritative; allow
// any in-flight immediate title generation to persist (saveConvo uses noUpsert).
resolveConvoReady();
acceptsTitleEvents = false;
if (titleEventPromise) {
await titleEventPromise;
}
let scheduleCompletionError;
if (terminalWasAborted) {
scheduleCompletionError = 'Scheduled run was stopped';
} else if (preemptIncomplete) {
scheduleCompletionError = 'Scheduled run was interrupted before completion';
}
await settleScheduledRun({
status: terminalWasAborted || preemptIncomplete ? 'interrupted' : 'success',
...(scheduleCompletionError != null && { error: scheduleCompletionError }),
});
let terminalPublicationStarted = false;
try {
const pendingSteers = terminalClaim.drainedSteers.map(toPendingSteer);
const finalEvent = {
final: true,
conversation,
title: conversation.title,
requestMessage: sanitizeMessageForTransmit(userMessage),
responseMessage: {
...response,
...((terminalWasAborted || preemptIncomplete) && { unfinished: true }),
},
...(pendingSteers.length > 0 && { pendingSteers }),
};
logger.debug(
terminalWasAborted
? `[ResumableAgentController] Emitting ABORTED FINAL event`
: `[ResumableAgentController] Emitting FINAL event`,
{
streamId,
wasAbortedBeforeComplete: terminalWasAborted,
userMessageId: userMessage?.messageId,
responseMessageId: response?.messageId,
conversationId: conversation?.conversationId,
},
);
terminalPublicationStarted = true;
const publication = await GenerationJobManager.publishTerminalClaim(
terminalClaim,
finalEvent,
);
let terminalOutcome = 'completed_without_delta';
if (publication.persistenceFailed) {
terminalOutcome = 'error';
} else if (terminalWasAborted) {
terminalOutcome = 'aborted';
}
startupTelemetry?.end(terminalOutcome);
} catch (terminalError) {
/** A failure while constructing the payload happened after this
* controller's terminal CAS but before the manager could durably
* settle it. Publish conservative reconciliation immediately. Once
* publication starts, the manager either stores the payload or owns
* its bounded recovery marker, so retrying with a different payload
* here would only risk duplicate delivery. */
if (!terminalPublicationStarted) {
try {
await GenerationJobManager.publishTerminalClaim(terminalClaim, null);
} catch (reconcileError) {
logger.warn(
'[ResumableAgentController] Failed to publish terminal persistence reconciliation',
reconcileError,
);
}
}
throw terminalError;
} finally {
// Pair every successful claim even when final-event construction or
// transport publication throws. Cleanup is epoch/runtime guarded.
await finishOwnedTerminalClaim();
}
await finishResumableRequest(req, userId);
if (titleTiming === 'immediate') {
// Title was fired in parallel above (if eligible); a stopped turn already
// aborted it before `resolveConvoReady`. Defer disposal until it settles
// so the run/req aren't torn down mid-generation.
if (immediateTitlePromise) {
immediateTitlePromise.finally(() => {
if (client) {
disposeClient(client);
}
});
} else if (client) {
disposeClient(client);
}
} else if (shouldGenerateTitle) {
trailingWritePromise = addTitle(req, {
text: text || getAttachmentTitleText(req.body.files),
response: { ...response },
client,
})
.catch((err) => {
logger.error('[ResumableAgentController] Error in title generation', err);
})
.finally(() => {
if (client) {
disposeClient(client);
}
});
} else {
if (client) {
disposeClient(client);
}
}
} catch (error) {
// Any failure (user Stop, or a preflight/quota failure before the run is
// even created) must cancel the title and unblock its waits: the title's
// `_waitForRun` would otherwise never resolve, deferring client disposal
// until the 45s title timeout, and no title should persist for a failed turn.
titleAbortController.abort();
titleDiscardController.abort();
job.abortController.signal.removeEventListener('abort', abortTitleOnJobAbort);
acceptsTitleEvents = false;
resolveConvoReady();
// Once this controller owns terminal persistence, no competing error
// transition can win. Settle its pending marker with conservative
// reconciliation on any required-write/final-construction failure,
// then release exactly that claim.
let ownsScheduledFailure = false;
if (terminalClaim && !terminalClaimFinished) {
ownsScheduledFailure = true;
try {
await GenerationJobManager.publishTerminalClaim(terminalClaim, null);
} catch (publishError) {
logger.warn(
'[ResumableAgentController] Failed to publish terminal persistence reconciliation',
publishError,
);
} finally {
await finishOwnedTerminalClaim().catch((finishError) => {
logger.warn(
'[ResumableAgentController] Failed to finish terminal persistence claim',
finishError,
);
});
}
logger.error(
`[ResumableAgentController] Terminal persistence failed for ${streamId}:`,
error,
);
startupTelemetry?.end('error', error);
} else if (pausePersistenceFailed) {
ownsScheduledFailure = pausePersistenceFailureFinalized;
// failPausePersistence owns the only legal requires_action -> error
// transition for this exact action/epoch. Never fall through to
// completeJob, which could race a newer action or replacement job.
logger.error(
`[ResumableAgentController] Pause persistence failed for ${streamId}:`,
error,
);
startupTelemetry?.end('error', error);
} else if (job.abortController.signal.aborted || error.message?.includes('abort')) {
ownsScheduledFailure = true;
logger.debug(`[ResumableAgentController] Generation aborted for ${streamId}`);
startupTelemetry?.end('aborted');
// abortJob already handled emitDone and completeJob
} else {
logger.error(`[ResumableAgentController] Generation error for ${streamId}:`, error);
const generationError = error.message || 'Generation failed';
try {
// completeJob first wins running -> error and atomically parks
// steers, then publishes. A competing abort/pause emits nothing.
ownsScheduledFailure =
(await GenerationJobManager.completeJob(streamId, generationError, jobCreatedAt)) ===
true;
} catch (completeErr) {
logger.warn(
'[ResumableAgentController] completeJob failed during generation-error cleanup',
completeErr,
);
} finally {
startupTelemetry?.end('error', error);
}
}
if (ownsScheduledFailure && !scheduleTerminalOutcomeRecorded) {
const scheduledFailure = classifyScheduledFailure(
error,
job.abortController.signal.aborted,
);
await settleScheduledRun(scheduledFailure);
}
try {
await finishResumableRequest(req, userId);
} finally {
disposeBackgroundClient();
}
// Don't continue to title generation after error/abort
return;
}
};
// Start generation and handle any unhandled errors
void startGeneration()
.catch(async (err) => {
logger.error(
`[ResumableAgentController] Unhandled error in background generation: ${err.message}`,
);
startupTelemetry?.end('error', err);
let errorFinalized = false;
if (!pausePersistenceFailed) {
errorFinalized =
(await GenerationJobManager.completeJob(streamId, err.message, jobCreatedAt).catch(
(completeErr) => {
logger.warn(
'[ResumableAgentController] completeJob failed during background-error cleanup',
completeErr,
);
return false;
},
)) === true;
}
if (
(errorFinalized || (pausePersistenceFailed && pausePersistenceFailureFinalized)) &&
!scheduleTerminalOutcomeRecorded
) {
await settleScheduledRun(classifyScheduledFailure(err));
}
try {
await finishResumableRequest(req, userId);
} finally {
disposeBackgroundClient();
}
})
.finally(async () => {
await Promise.allSettled([immediateTitlePromise, trailingWritePromise].filter(Boolean));
if (providerExecutionId) {
await GenerationJobManager.markProviderExecutionDrained?.(
streamId,
jobCreatedAt,
providerExecutionId,
);
}
})
.catch((drainError) => {
logger.warn(
'[ResumableAgentController] Failed to record completed provider drain',
drainError,
);
});
} catch (error) {
logger.error('[ResumableAgentController] Initialization error:', error);
const initializationFailure = getInitializationFailure(error);
try {
if (!res.headersSent) {
if (error?.code === 'GENERATION_PREDECESSOR_MISMATCH') {
const currentJob = error.currentJob;
const currentStatus = currentJob?.status;
if (isTriggerContinuation && currentJob?.active === true) {
res.set('Retry-After', '1');
sendGenerationJson(
res,
409,
{
code: 'PARENT_NOT_READY',
error: 'Another generation became active before the continuation could start.',
},
generationProtocolVersion,
);
} else {
const predecessorVerified =
currentJob != null &&
Number.isSafeInteger(currentJob.createdAt) &&
currentJob.createdAt >= 0 &&
currentJob.verified !== false;
sendGenerationJson(
res,
409,
{
status: 'predecessor_mismatch',
code: 'GENERATION_PREDECESSOR_MISMATCH',
error: predecessorVerified
? 'A newer generation became current before this request could start.'
: 'The prior generation could not be verified. Please retry.',
streamId,
conversationId: currentJob?.conversationId ?? conversationId,
generationCreatedAt: currentJob?.createdAt,
predecessorVerified,
active:
typeof currentJob?.active === 'boolean'
? currentJob.active
: currentStatus === 'running' || currentStatus === 'requires_action',
},
generationProtocolVersion,
);
}
} else if (error?.code === 'RECOVERY_PAYLOAD_MISMATCH') {
sendGenerationJson(
res,
409,
{
code: 'RECOVERY_PAYLOAD_MISMATCH',
error: 'The queued message changed before it could be recovered. Please retry.',
},
generationProtocolVersion,
);
} else if (initializationFailure) {
sendGenerationJson(
res,
initializationFailure.status,
initializationFailure,
generationProtocolVersion,
);
} else {
sendGenerationJson(
res,
500,
{ error: error.message || 'Failed to start generation' },
generationProtocolVersion,
);
}
}
} catch (notificationError) {
logger.warn(
'[ResumableAgentController] Failed to send initialization error response',
notificationError,
);
} finally {
startupTelemetry?.end(
error?.code === 'GENERATION_PREDECESSOR_MISMATCH' ? 'deduplicated' : 'error',
error,
);
}
// Finalize THIS failed job before releasing the idempotency claim. Releasing first would
// let the client's retry win the same key and createJob() the same streamId while we are
// still here. The generation guard is defense-in-depth around that ordering. A
// completeJob() rejection (store hiccup) must NOT skip the
// release + pending-request decrement below, or the retry stays wedged behind the claim
// and the concurrency slot leaks — so swallow its error. (A failed completeJob did not
// finalize anything, so releasing afterward can't let it abort a later replacement.)
let initializationFinalized = jobCreatedAt == null;
if (jobCreatedAt != null) {
const initializationError = initializationFailure
? JSON.stringify(initializationFailure)
: error.message || 'Failed to start generation';
initializationFinalized =
(await GenerationJobManager.completeJob(streamId, initializationError, jobCreatedAt).catch(
(completeErr) => {
logger.warn(
'[ResumableAgentController] completeJob failed during init-error cleanup',
completeErr,
);
return false;
},
)) === true;
}
if (initializationFinalized && !scheduleTerminalOutcomeRecorded) {
await settleScheduledRun(classifyScheduledFailure(error));
}
if (ownedIdempotencyClaim) {
await GenerationJobManager.releaseGeneration(
userId,
clientRequestId,
streamId,
ownedIdempotencyClaim,
).catch(() => {});
}
await finishResumableRequest(req, userId);
if (client) {
disposeClient(client);
}
if (jobCreatedAt != null && providerExecutionId) {
await GenerationJobManager.markProviderExecutionDrained?.(
streamId,
jobCreatedAt,
providerExecutionId,
).catch((drainError) => {
logger.warn(
'[ResumableAgentController] Failed to record initialization-error provider drain',
drainError,
);
});
}
}
};
module.exports = ResumableAgentController;