LibreChat/api/server/controllers/agents/request.js
Danny Avila 68fc46a055
🪢 feat: Resume Bound Event Actors from Checkpoint Forks (#15227)
* feat: resume event actors from checkpoint forks

* fix: fence event actor checkpoint uncertainty

* fix: satisfy event actor type contracts

* fix: make actor reconciliation recoverable

* fix: preserve event actor lifecycle transitions

* fix: fence event actor lifecycle outcomes

* fix: enforce event actor lifecycle ownership

* fix: retain event actor settlement proof

* 🔒 fix: Retain Event Actor Receipts Through Repair and Bound Their Journal

Repair and compensation deleted the reconciliation row they resolved, which
was the only durable proof that the invocation had already applied an external
action. A delayed duplicate owner could then reacquire the same invocation id
and repeat that action. Both resolutions now retire their receipt to `settled`
and record how it settled, so the same-id tombstone survives; `history_repaired`
and `action_compensated` still force a cold rebuild. Compensation undoes the
effect without re-authorizing the delivery, so a legitimate retry must arrive
under a new invocation id. A retried repair converges on its own receipt.

Bound the journal so a long-lived actor cannot grow its conversation document
without limit: a new fence is admitted only when no active lifecycle row
exists, so a capped push can evict nothing but the oldest settled receipts.

Stop shipping the unbounded source payload on every bound-child continuation.
It rode the delivery body regardless of the feature flag while the sibling
`fire` body deliberately sends event identity alone, so a large webhook payload
could push a previously working delivery past the chat route's body limit. The
actor binds an invocation from identity and never builds the prompt from it.

Blank the positional token map on warm continuations. It is derived from the
full DB history, while a warm run executes on checkpoint-restored state, so its
indices address different messages and the pruner never recounts them —
misattributing cached counts to the wrong messages in both directions.

Keep the replaced-claim exit on its cleanup path when preserving reconciliation
fails: the committing CAS already left a blocking row, so the failed status
upgrade costs provenance, not safety.

* 🧪 test: Pin the Warm Continuation's Map/Summary Asymmetry

Give the warm-continuation client test a populated token map and a real
cross-run summary so its assertions bite: the positional map must arrive
blank (checkpoint-restored state no longer matches DB-derived indices, and
the pruner never recounts a populated entry) while `initialSummary` must
pass through unchanged — it rides the system tail and summarizes
pre-boundary turns that were excluded from the very history the committed
checkpoint was built from, so blanking it would silently drop context no
warm run can recover.

* ⚖️ fix: Honor Compensation in Settlement and Age-Bound the Receipt Journal

A compensated receipt still tombstones its invocation id, but its external
effect was explicitly undone — the terminal handler nonetheless replayed
every settled lifecycle's stored action as authoritative and settled the
public outcome as applied, telling an action-aware source the operation
stands and suppressing the new-invocation retry compensation requires. The
handler now settles a compensated invocation as failed with an explicit
compensation error, overriding even fresh applied run evidence from a
replayed generation; verified and repaired receipts continue to replay
applied.

Receipt eviction is now primarily age-based: a stale same-id owner is
bounded by time, not by how many newer invocations settle, so the previous
count-only slice let a high-rate actor evict a tombstone while its delayed
duplicate owner could still wake and repeat the action. Admission prunes
only settled receipts older than a retention window that dwarfs every
generation, job, and delivery-retry lifetime, and the count cap is demoted
to a raised document-size backstop.

* 🔀 fix: Serialize Compensation Against Settlement and Never Evict Live Receipts

The terminal handler read its lifecycle snapshot, verified history, and then
settled the public outcome — so a compensation resolving the same receipt
during that window lost: the handler settled applied from its stale snapshot
and no retry could ever change the replay-identity-locked outcome. The
receipt's status CAS is now the serialization point: verification resolves
the receipt BEFORE settling, whichever transition wins determines the public
outcome, and a crash between resolve and settle converges through the
retained receipt's replay. The verified-replay probe requires the receipt's
own resolution, so a compensated receipt can never satisfy a verification
retry. This inverts the settle-before-receipt ordering deliberately: that
ordering guarded proof that resolution used to delete, and the receipt now
retains its full action proof through resolution.

The document-size cap is no longer an eviction quota. A receipt inside its
retention window is never discarded: when the journal holds a full cap of
unexpired receipts, new invocations are refused fail-closed until receipts
age out, making duplicate protection and document integrity simultaneous
invariants instead of a rate-dependent trade.

* 🎓 fix: Keep Skill-Bearing Event Actors on the Legacy Path

Skill primes are spliced into the message list directly ahead of the newest
message, and a warm continuation forwards only that newest message — so a
checkpoint-restored actor would keep serving the prime bodies baked in at
its last cold start and never observe an edited or newly attached skill.
Until the actor head carries a context fingerprint that forces a cold
rebuild when the agent's skill context changes, agents with always-apply or
manual skill primes stay on the legacy path, which re-primes fresh bodies
every turn: correct on every event, just never warm.

* 📜 fix: Gate Fork Mode on the Skills Capability, Not Just Request-Time Primes

History-derived re-priming was a third path into the same staleness class:
an actor that previously invoked a skill carries no request-time prime
arrays, yet primeInvokedSkills re-resolves that skill's current body from
history each turn and the warm slice drops the reconstruction — leaving the
checkpoint's old body active after edits. The fork gate now keys on the
priming hook itself (present exactly when the skills capability is enabled)
alongside the request-time arrays, so every skill-body path routes to the
legacy rebuild until #15235's context fingerprint restores warm
continuation for skill-bearing actors.

* 🧾 fix: Capture Applied-Action Proof at Tool Execution, Not After sendMessage

The executor read applied-action evidence from the run-step collection the
instant sendMessage resolved, but that collection is populated
asynchronously — an applied invocation could classify as actionless
(runSteps still empty while the tool result already streamed), discarding
its fork and stranding the actor cold while the terminal handler later
settled the same delivery as applied from the persisted evidence.

Authoritative proof is now recorded in graph context the moment the
expected tool executes: the request-owned recorder observes the tool-end
chain (which ToolNode dispatches synchronously with both input and output)
and applies the same fences as run-step evidence — exact tool name with the
MCP-suffixed form, the declared argument subset against the execution
input, an error-free result, and the background non-execution receipt
exclusion. readAppliedAction consults the receipt first; run-step
inspection remains the fallback for paths that bypass the tool-end chain.

Regression coverage reproduces the observed ordering: the real executor
commits the head from the receipt while run steps are empty, warm-continues
the next event, and never re-executes the action; recorder fences and the
receipt-first controller wiring are covered separately.

* 🎯 fix: Supply Execution Arguments to the Tool End Callback

The live Vertex + MCP canary exposed a contract mismatch the synthetic
fixtures hid: the ON_TOOL_EXECUTE execution path invoked its tool end
callback with output only, while the action recorder must verify the
declared argument subset against the execution input. The receipt never
qualified, every turn fell back to cold history rebuilds, and the
tournament advanced with zero actor heads and zero retained checkpoints
while looking successful.

The execution handler owns both halves at the same moment, so the fix is
at the source rather than a correlation store: ToolEndCallbackData gains
the executed call's input and every callback site passes tc.args. A
handler-level regression drives the real createToolExecuteHandler and
asserts the callback receives both fields; recorder regressions pin the
production shapes — an output-only tool end must starve an
argument-fenced receipt rather than trust an unfenced match, and still
qualifies a name-only expected action.

* 🕵️ fix: Mark Background Deliveries So They Cannot Impersonate Applied Actions

The background-claim callback reports the ORIGINAL tool's name for artifact
attribution on the poll turn that harvests a completed task. A name-only
expected action could therefore be impersonated by work some earlier turn
dispatched: the recorder would attribute that delivery to the current
invocation and commit a head whose state never contained the invocation's
own action. The run-step evidence path never had this hole — it sees the
poll tool's name — so the recorder must match its provenance discipline.

Delivery callbacks now carry an explicit backgroundDelivery marker set at
the one site that rewrites the name, and the recorder ignores marked
deliveries outright. Regressions pin both halves of the contract: the
delivery callback must carry the marker with the poll call's arguments,
and a marked delivery can never qualify even a name-only expected action.

* 🧿 fix: Version Invalidations, Keep Evidence Ahead of Output Policy, Gate Detachable Actions

Three closeout-round findings, each converted into an invariant.

Every legacy-path invalidation now advances a durable epoch — including for
headless and already cold-marked actors, where the marker alone leaves no
CAS-visible trace — and the actor-head CAS requires the epoch observed at
preparation. A concurrently prepared fork whose history predates an
intervening legacy turn can no longer commit past it; the commit reports an
ordinary conflict and journals.

Execution identity is now emitted before post-execution output policy: when
a side-effecting tool succeeds but its returned content is withheld by the
output filter, the callback delivers an outputFiltered receipt with blank
content — the recorder accepts it as proof (rejecting model-detached calls
it cannot distinguish through the blank shape), the artifact path never
sees it, and an applied action is no longer reclassified actionless and
re-executed on retry.

Background-capable expected actions stay off the fork path: dispatch
returns a launch handle every evidence fence correctly rejects, and the
completion is provenance-marked as another turn's work, so a fork would
settle actionless before the external effect lands with nothing to stop a
retry from dispatching it again. The gate mirrors the MCP-suffix name
matching of the evidence path.

* 🚧 fix: Seal the Whole Legacy Turn Behind a Second Epoch Advance

The epoch fenced only the legacy turn's start: a fork preparing after the
begin invalidation but before the turn's message persistence observed the
new epoch and cold marker, rebuilt from history that did not yet contain
the turn, and committed cleanly because nothing advanced the epoch again —
making the incomplete rebuild authoritative and clearing the marker.

Every legacy event turn now seals its invalidation at terminal persistence
with a second epoch advance, on the success, replaced-claim, and error
exits alike. Sealing deliberately carries no quiescence requirement — it
must succeed while a fork fence is active, because that is exactly the
mid-turn race it defeats — and a fork that already committed against the
begin epoch is healed the same way: the seal re-marks the head cold, so
the next event rebuilds with complete history. Seal failure never diverts
the turn's own exit; the begin bump still fences everything prepared
before the turn.

* 🔗 fix: Replace the Best-Effort Epoch Bump With a Durable Legacy-Turn Fence

The second epoch advance could not make a legacy turn atomic, and three
findings shared that root cause: two conditional updates left a headless
gap a fork could create the head inside; the error exit sealed before
saveErrorTurn made the error history durable; and any crash or failure
between persistence and sealing left an incomplete fork authoritative,
because the seal was best-effort and its failure was swallowed.

A legacy turn now carries one durable fence. A token is written before
execution by a single update-pipeline write — no two-write gap, and the
cold marker is applied only where a head exists via $cond/$$REMOVE. While
the token is present no fork may prepare (the adapter refuses) or commit
(the CAS requires its absence), because the turn's messages are not yet
durable. One atomic write clears the exact token and advances the epoch
once history is persisted — after saveErrorTurn on the error exit — and
success, replacement, and error exits all route through it.

Failure is now fail-closed rather than silent: a failed seal leaves the
token set, which keeps blocking forks and is logged as such, and a fence
abandoned by a crashed turn is reclaimed only once stale, advancing the
epoch and marking any head cold so the next event rebuilds from whatever
history actually survived.

* fix: serialize legacy event actor turns

* fix: close legacy actor fence ownership gaps

* fix: preserve legacy actor persistence fences
2026-08-26 08:39:00 -04:00

2881 lines
114 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,
createMCPRuntimeRequestBody,
isAgentEventRetentionActive,
executeAgentEventActor,
findAgentEventAppliedAction,
createAgentEventActionRecorder,
isHITLEnabled,
agentRequestsAskUserQuestion,
} = 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,
saveConvo,
getMessages,
getConvo,
getAgentEventActorSnapshot,
commitAgentEventActorState,
beginAgentEventActorLegacyTurn,
completeAgentEventActorLegacyTurn,
recordAgentEventActorReconciliation,
resolveAgentEventActorReconciliation,
isAgentTriggerPrincipalActive,
isSubagentOwnerAdmissible,
} = require('~/models');
const {
acquireEventChildGenerationLease,
} = require('~/server/services/Endpoints/agents/eventChildLease');
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);
}
}
}
async function saveErrorTurn(
req,
{
conversationId,
endpointOption,
isNewConvo,
errorText,
liveUserMessage,
liveResponseMessageId,
sender,
},
) {
try {
const { isContinued, isRegenerate, editedContent, responseMessageId, overrideParentMessageId } =
req.body ?? {};
if (
isContinued ||
editedContent != null ||
(responseMessageId && !isRegenerate) ||
req.body?.recoverySteerId != null ||
req.body?.clientRequestId?.startsWith?.('steer-recovery:') === true
) {
return;
}
let userMessage = null;
let errorMessageId = null;
let errorParentMessageId = null;
if (isRegenerate) {
errorMessageId =
typeof responseMessageId === 'string' && responseMessageId.length > 0
? responseMessageId
: null;
errorParentMessageId = liveUserMessage?.messageId ?? overrideParentMessageId ?? null;
} else {
userMessage =
liveUserMessage != null
? {
...liveUserMessage,
...(liveUserMessage.files == null &&
Array.isArray(req.body?.files) &&
req.body.files.length > 0 && { files: req.body.files }),
...(liveUserMessage.manualSkills == null &&
Array.isArray(req.body?.manualSkills) &&
req.body.manualSkills.length > 0 && { manualSkills: req.body.manualSkills }),
...(liveUserMessage.alwaysAppliedSkills == null &&
Array.isArray(req.body?.alwaysAppliedSkills) &&
req.body.alwaysAppliedSkills.length > 0 && {
alwaysAppliedSkills: req.body.alwaysAppliedSkills,
}),
}
: getPreliminaryUserMessage(req.body, conversationId);
if (!userMessage) {
return;
}
errorMessageId = getPreliminaryResponseMessageId(
liveUserMessage != null ? { messageId: liveUserMessage.messageId } : req.body,
);
errorParentMessageId = userMessage.messageId;
}
if (!errorMessageId || !errorParentMessageId) {
return;
}
const userId = req.user.id;
const existing = await getMessages(
{ user: userId, messageId: errorMessageId, conversationId },
'_id',
);
if (existing.length > 0) {
return;
}
if (liveResponseMessageId != null && liveResponseMessageId !== errorMessageId) {
const partial = await getMessages(
{ user: userId, messageId: liveResponseMessageId, conversationId },
'_id',
);
if (partial.length > 0) {
return;
}
}
const reqCtx = {
userId,
isTemporary: req?._agentEventBindingRetention?.isTemporary ?? req?.body?.isTemporary,
expiredAt: req?._agentEventBindingRetention?.expiredAt,
interfaceConfig: req?.config?.interfaceConfig,
};
const context = 'api/server/controllers/agents/request.js - failed turn';
const endpoint = endpointOption?.endpoint;
const model = getAgentResponseModel(req, endpointOption);
const iconURL = getEndpointIconURL(req, endpointOption);
if (userMessage) {
const savedUserMessage = await saveMessage(
reqCtx,
{
...userMessage,
user: userId,
sender: 'User',
isCreatedByUser: true,
error: false,
unfinished: false,
},
{ context },
);
if (!savedUserMessage) {
throw new Error('Failed user message could not be persisted');
}
}
const savedErrorMessage = await saveMessage(
reqCtx,
{
messageId: errorMessageId,
conversationId,
parentMessageId: errorParentMessageId,
sender: sender ?? 'AI',
...(endpoint != null && { endpoint }),
...(model != null && { model }),
...(iconURL != null && { iconURL }),
user: userId,
text: errorText,
error: true,
unfinished: false,
isCreatedByUser: false,
},
{ context },
);
if (!savedErrorMessage) {
throw new Error('Failed response message could not be persisted');
}
const agentId = endpointOption?.agent_id ?? req.body?.agent_id;
const chatProjectId = endpointOption?.chatProjectId ?? req.body?.chatProjectId;
const seedConvo = isNewConvo || req.resolvedConversation === null;
const convoFields = seedConvo
? {
...(endpoint != null && { endpoint }),
...(endpointOption?.endpointType != null && {
endpointType: endpointOption.endpointType,
}),
...(model != null && { model }),
...(iconURL != null && { iconURL }),
...(endpointOption?.spec != null && { spec: endpointOption.spec }),
...(agentId != null && { agent_id: agentId }),
...(typeof chatProjectId === 'string' && chatProjectId.length > 0 && { chatProjectId }),
}
: {};
await saveConvo(
reqCtx,
{ conversationId, ...convoFields },
seedConvo ? { context } : { context, noUpsert: true },
);
} catch (err) {
logger.error('[AgentController] Failed to persist error turn', err);
throw err;
}
}
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 rawOverrideUserMessageId = req.body?.overrideUserMessageId;
const rawOverrideConversationId = req.body?.overrideConvoId;
if (
(rawOverrideUserMessageId != null && typeof rawOverrideUserMessageId !== 'string') ||
(rawOverrideConversationId != null && typeof rawOverrideConversationId !== 'string')
) {
startupTelemetry?.end('rejected');
return sendGenerationJson(
res,
400,
{
code: 'INVALID_OVERRIDE_ID',
error: 'overrideUserMessageId and overrideConvoId must be strings.',
},
generationProtocolVersion,
);
}
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 = rawOverrideUserMessageId;
const recoveredSteerPayload = isRecoveredSteerRequest
? buildRecoveredSteerPayload(text, req.body?.files, req.body?.quotes)
: 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');
/** Allocate the turn identities before Agent initialization. Request-scoped
* MCP transports resolve BODY placeholders while tools are discovered, so
* discovery and graph execution must receive the same response-scoped body.
* BaseClient otherwise allocates these IDs later in `sendMessage`, after MCP
* connections already exist. */
const overrideUserMessageId = rawOverrideUserMessageId
? rawOverrideUserMessageId.split(Constants.COMMON_DIVIDER)[0]
: undefined;
/** Event deliveries already carry a stable, retry-safe idempotency key. Reuse
* it as the public child-task identity so the lease, persisted turn, and
* parent activity index continue to agree after the live lease is released. */
const eventTaskId =
req._agentEventBindingParentConversationId != null
? (clientRequestId ?? crypto.randomUUID())
: undefined;
if (eventTaskId != null) {
req._agentEventTaskId = eventTaskId;
}
const preallocatedUserMessageId =
eventTaskId == null
? (overrideUserMessageId ?? overrideParentMessageId ?? crypto.randomUUID())
: `${eventTaskId}:user`;
const overrideConversationId = rawOverrideConversationId
? rawOverrideConversationId.split(Constants.COMMON_DIVIDER)[0]
: undefined;
const effectiveConversationId = overrideConversationId ?? conversationId;
let preallocatedResponseMessageId =
eventTaskId == null
? (editedResponseMessageId ?? crypto.randomUUID())
: `${eventTaskId}:assistant`;
if (
(editedContent != null && !isContinued) ||
(isRegenerate && preallocatedResponseMessageId.endsWith('_'))
) {
preallocatedResponseMessageId = crypto.randomUUID();
}
const mcpRequestBody = createMCPRuntimeRequestBody({
messageId: preallocatedResponseMessageId,
conversationId: effectiveConversationId,
parentMessageId:
editedContent != null ? preallocatedResponseMessageId : preallocatedUserMessageId,
});
let client = null;
let jobCreatedAt;
let providerExecutionId;
let releaseEventChildLease;
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;
};
/** The loopback trigger host binds this lifecycle identity to the same
* idempotency key that owns generation admission. Ignore mismatched or
* direct-chat metadata rather than letting callers relabel another run. */
const rawAgentEventDelivery = req.body?.agentEventDelivery;
const agentEventDelivery =
isTriggerContinuation &&
rawAgentEventDelivery != null &&
typeof rawAgentEventDelivery === 'object' &&
rawAgentEventDelivery.deliveryKey === clientRequestId
? rawAgentEventDelivery
: undefined;
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, messageId: preallocatedUserMessageId },
conversationId,
);
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(),
// Same owner-recorded pattern: this build's drain merges queued steer
// quotes into the injected turn. Admission on another replica must
// not store/acknowledge quotes an older owner would drop.
steerQuotesCapable: true,
// 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._agentEventBindingRetention?.isTemporary ?? req.body?.isTemporary,
...(agentEventDelivery != null && {
agentEventDeliveryKey: agentEventDelivery.deliveryKey,
...(agentEventDelivery.expectedAction != null && {
agentEventExpectedAction: agentEventDelivery.expectedAction,
}),
}),
...(isRegenerate && { isRegenerate: true }),
...(scheduleId
? {
scheduleId,
scheduledFor,
preserveForScheduleReconcile: true,
...(Number.isSafeInteger(scheduleConfigRevision) && {
scheduleConfigRevision,
}),
...(req._isManualScheduledFire === true && { scheduleManual: true }),
}
: {}),
responseMessageId: preallocatedResponseMessageId,
mcpRequestBody,
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 (req._agentEventBindingParentConversationId != null) {
/** The generation job is the durable marker that a deletion on another replica
* can abort. Recheck only after that marker exists: either the deletion fence
* wins and this run stops here, or the deletion observes and drains this job. */
releaseEventChildLease = await acquireEventChildGenerationLease({
userId,
tenantId: req._agentEventBindingTenantId,
conversationId,
streamId,
taskId: eventTaskId,
jobCreatedAt,
retentionExpiresAt: req._agentEventBindingRetention?.expiredAt,
});
if (releaseEventChildLease == null) {
const bindingActive = isAgentEventRetentionActive(
req._agentEventBindingRetention?.expiredAt,
);
throw Object.assign(
new Error(
bindingActive
? 'The event actor is already handling another turn'
: 'The event binding parent is no longer available',
),
{
code: bindingActive ? 'EVENT_ACTOR_NOT_READY' : 'EVENT_BINDING_PARENT_ENDED',
status: 409,
},
);
}
const [eventParent, ownerAdmissible] = await Promise.all([
getConvo(userId, req._agentEventBindingParentConversationId),
isSubagentOwnerAdmissible(userId),
]);
if (!ownerAdmissible) {
throw Object.assign(new Error('The event actor is temporarily unavailable'), {
code: 'EVENT_ACTOR_NOT_READY',
status: 409,
});
}
if (
eventParent == null ||
eventParent.subagentThread != null ||
eventParent.agent_id !== req._agentEventBindingParentAgentId ||
(eventParent.tenantId ?? undefined) !== req._agentEventBindingTenantId ||
!isAgentEventRetentionActive(req._agentEventBindingRetention?.expiredAt) ||
!isAgentEventRetentionActive(eventParent.expiredAt)
) {
throw Object.assign(new Error('The event binding parent is no longer available'), {
code: 'EVENT_BINDING_PARENT_ENDED',
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?._agentEventBindingRetention?.isTemporary ?? req?.body?.isTemporary,
expiredAt: req?._agentEventBindingRetention?.expiredAt,
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,
requestBody: mcpRequestBody,
});
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;
let liveResponseMessageId = preallocatedResponseMessageId;
const getReqData = (data = {}) => {
if (data.userMessage) {
userMessage = data.userMessage;
}
if (data.responseMessageId) {
liveResponseMessageId = data.responseMessageId;
}
// 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;
};
const eventActorTenantId = req._agentEventBindingTenantId;
let appliedEventActor;
let eventActorPersistenceComplete = false;
let legacyEventActorTurnToken;
/** Closes the durable legacy-turn fence once this turn's history is
* persisted, clearing the token and advancing the epoch in one write. A
* failure here leaves the token SET, which keeps blocking forks until an
* explicit reconciliation proves the ambiguous turn safe to release. */
const sealLegacyEventActorTurn = async () => {
if (legacyEventActorTurnToken == null) {
return;
}
const token = legacyEventActorTurnToken;
legacyEventActorTurnToken = undefined;
try {
const sealed = await completeAgentEventActorLegacyTurn({
user: userId,
conversationId,
...(eventActorTenantId == null ? {} : { tenantId: eventActorTenantId }),
token,
});
if (!sealed) {
logger.error(
`[event-actor] Legacy turn fence ${token} was not sealed; forks stay blocked pending reconciliation`,
);
}
} catch (sealError) {
logger.error(
`[event-actor] Failed to seal legacy turn fence ${token}; forks stay blocked pending reconciliation`,
sealError,
);
}
};
const recordEventActorPersistenceFailure = async (error) => {
if (appliedEventActor == null || eventActorPersistenceComplete) {
return;
}
const recorded = await recordAgentEventActorReconciliation({
user: userId,
conversationId,
...(eventActorTenantId == null ? {} : { tenantId: eventActorTenantId }),
reconciliation: {
invocationId: appliedEventActor.invocationId,
status: 'persistence_failed',
checkpoint: appliedEventActor.checkpoint,
action: appliedEventActor.action,
error: String(error?.message ?? error).slice(0, 1024),
observedAt: new Date(),
},
});
if (!recorded) {
throw new Error('Failed to preserve applied event actor persistence reconciliation');
}
};
try {
const onStart = (userMsg, respMsgId, _isNewConvo) => {
userMessage = userMsg;
liveResponseMessageId = respMsgId;
// 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,
preallocatedUserMessageId,
preallocatedResponseMessageId,
progressOptions: {
res: {
write: () => true,
end: () => {},
headersSent: false,
writableEnded: false,
},
},
};
const checkpointForksEnabled =
req.config?.endpoints?.[EModelEndpoint.agents]?.eventDriven?.checkpointForks === true;
const agentsConfig = req.config?.endpoints?.[EModelEndpoint.agents];
const eventActorAgents = [
client?.options?.agent,
...(client?.agentConfigs?.values?.() ?? []),
].filter(Boolean);
const eventActorMayPause =
isHITLEnabled(agentsConfig?.toolApproval) ||
eventActorAgents.some(agentRequestsAskUserQuestion);
/** Skill bodies are spliced or reconstructed into the message list
* ahead of the newest message, and a warm continuation forwards only
* that newest message — a checkpoint-restored actor would keep serving
* the bodies baked in at its last cold start and never observe an
* edited or newly attached skill. That covers request-time primes AND
* history-derived re-priming: `primeInvokedSkills` re-resolves any
* skill the actor previously invoked, so its mere presence (the skills
* capability) keeps the actor on the legacy path, which rebuilds fresh
* bodies every turn. #15235 tracks the context fingerprint that will
* restore warm continuation for skill-bearing actors. */
const eventActorHasSkillPrimes =
client?.options?.primeInvokedSkills != null ||
(client?.options?.agent?.alwaysApplySkillPrimes?.length ?? 0) > 0 ||
(client?.options?.agent?.manualSkillPrimes?.length ?? 0) > 0;
/** A background-opted expected action can detach: dispatch returns a
* launch handle the evidence fences correctly reject, and the later
* completion is provenance-marked as another turn's work — so a fork
* would classify actionless and settle before the external effect
* lands, with no receipt to stop a retry from dispatching it again.
* The name comparison mirrors matchesExpectedAction's MCP suffix. */
const expectedActionToolName = agentEventDelivery?.expectedAction?.toolName;
const eventActorActionMayDetach =
typeof expectedActionToolName === 'string' &&
eventActorAgents.some((agent) =>
(agent.backgroundToolNames ?? []).some(
(name) =>
name === expectedActionToolName ||
name.startsWith(`${expectedActionToolName}_mcp_`),
),
);
const canUseEventActorFork =
checkpointForksEnabled &&
agentsConfig?.checkpointer?.type !== 'memory' &&
!eventActorMayPause &&
!eventActorHasSkillPrimes &&
!eventActorActionMayDetach &&
agentEventDelivery?.event != null &&
agentEventDelivery.expectedAction != null &&
typeof eventTaskId === 'string' &&
req._agentEventBindingParentConversationId != null;
/** Authoritative action proof is captured in graph context the moment
* the expected tool executes (see the observer tee in initialize.js);
* run-step inspection stays only as a fallback, because the run-step
* collection is populated asynchronously and can still be empty the
* instant sendMessage resolves — misreading an applied invocation as
* actionless would discard its fork and strand the actor cold. */
const eventActorActionRecorder = canUseEventActorFork
? createAgentEventActionRecorder(agentEventDelivery.expectedAction)
: undefined;
if (eventActorActionRecorder != null) {
req._agentEventActionObserver = eventActorActionRecorder.observeToolEnd;
}
const sendPromise = canUseEventActorFork
? executeAgentEventActor(
{
user: userId,
...(eventActorTenantId == null ? {} : { tenantId: eventActorTenantId }),
conversationId,
invocationId: eventTaskId,
event: agentEventDelivery.event,
expectedAction: agentEventDelivery.expectedAction,
signal: job.abortController.signal,
checkpointer: req.config?.endpoints?.[EModelEndpoint.agents]?.checkpointer,
invoke: async (actorContext) => {
client.checkpointNamespace = actorContext.checkpointNamespace;
client.eventActorCheckpointId = actorContext.checkpointId;
client.eventActorInvocationId = actorContext.invocationId;
client.eventActorContinuation = actorContext.continuation;
return client.sendMessage(text, messageOptions);
},
readAppliedAction: () =>
eventActorActionRecorder.read() ??
findAgentEventAppliedAction(
agentEventDelivery.expectedAction,
client?.run?.getRunSteps?.() ?? [],
client?.contentParts ?? [],
),
},
{
getSnapshot: getAgentEventActorSnapshot,
commitState: commitAgentEventActorState,
recordReconciliation: recordAgentEventActorReconciliation,
resolveReconciliation: resolveAgentEventActorReconciliation,
},
).then(({ value, execution }) => {
if (execution.status === 'applied') {
appliedEventActor = {
invocationId: eventTaskId,
checkpoint: execution.head.checkpoint,
action: execution.result.action,
};
}
logger.info('[event-actor] Bound child event completed', {
conversationId,
invocationId: eventTaskId,
status: execution.status,
continuation: execution.continuation,
});
return value;
})
: (async () => {
if (
agentEventDelivery?.event != null &&
req._agentEventBindingParentConversationId != null
) {
const token = crypto.randomUUID();
/** Open the fence BEFORE execution so a fork can neither run
* nor commit while this turn's messages are not yet durable. */
const legacyTurnOwner = {
user: userId,
conversationId,
...(eventActorTenantId == null ? {} : { tenantId: eventActorTenantId }),
};
const acquiredLegacyTurn = await beginAgentEventActorLegacyTurn({
...legacyTurnOwner,
token,
});
if (!acquiredLegacyTurn) {
throw Object.assign(new Error('The event actor is temporarily unavailable'), {
code: 'EVENT_ACTOR_NOT_READY',
status: 409,
});
}
/** The same exact token must survive a HITL pause. Retain local
* ownership before the Redis write so a metadata failure can
* still seal this token after its durable error turn is saved;
* provider execution starts only after metadata persists. */
legacyEventActorTurnToken = token;
await GenerationJobManager.updateMetadata(
streamId,
{ agentEventLegacyTurnToken: token },
jobCreatedAt,
);
}
return 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?._agentEventBindingRetention?.isTemporary ?? req?.body?.isTemporary,
expiredAt: req?._agentEventBindingRetention?.expiredAt,
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?._agentEventBindingRetention?.isTemporary ?? req?.body?.isTemporary,
expiredAt: req?._agentEventBindingRetention?.expiredAt,
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();
try {
await recordEventActorPersistenceFailure(
new Error('Event actor terminal persistence claim was replaced'),
);
} catch (reconciliationError) {
/** The committing CAS already left a non-settled row that blocks
* later actor turns, so a failed status upgrade costs provenance,
* not safety. Never divert this clean exit past its cleanup. */
logger.error(
'[event-actor] Failed to preserve replaced-claim reconciliation',
reconciliationError,
);
}
/** This controller lost terminal persistence ownership, so it cannot
* prove the winning Stop/replacement has written the unfinished
* response yet. Keep the conversation fence closed; a HITL resume
* carries the exact token, while every other orphan is handled by
* bounded stale reclaim. */
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?._agentEventBindingRetention?.isTemporary ?? req?.body?.isTemporary,
expiredAt: req?._agentEventBindingRetention?.expiredAt,
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 (appliedEventActor != null) {
const recorded = await recordAgentEventActorReconciliation({
user: userId,
conversationId,
...(eventActorTenantId == null ? {} : { tenantId: eventActorTenantId }),
reconciliation: {
invocationId: appliedEventActor.invocationId,
status: 'history_persisted',
checkpoint: appliedEventActor.checkpoint,
action: appliedEventActor.action,
observedAt: new Date(),
},
});
if (!recorded) {
throw new Error('Applied event actor history barrier could not be durably recorded');
}
}
await sealLegacyEventActorTurn();
eventActorPersistenceComplete = true;
// 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();
try {
await recordEventActorPersistenceFailure(error);
} catch (reconciliationError) {
logger.error(
'[event-actor] Failed to preserve terminal persistence reconciliation',
reconciliationError,
);
}
// 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;
let legacyEventActorErrorHistoryDurable = 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, {
beforeErrorPublication: () =>
saveErrorTurn(req, {
conversationId,
endpointOption,
isNewConvo,
errorText: generationError,
liveUserMessage: userMessage,
liveResponseMessageId,
sender: client?.sender,
}),
})) === true;
/** A true completion means this owner won the terminal CAS and
* the beforeErrorPublication barrier above finished. Only that
* combination proves the failed-turn rows are durable enough to
* let a checkpoint fork rebuild past this legacy turn. */
legacyEventActorErrorHistoryDurable = ownsScheduledFailure;
} catch (completeErr) {
logger.warn(
'[ResumableAgentController] completeJob failed during generation-error cleanup',
completeErr,
);
} finally {
startupTelemetry?.end('error', error);
}
}
/** Leave the fence set when terminal persistence loses ownership or
* fails. Time cannot prove whether an external action occurred, so an
* ambiguous fence remains fail-closed pending explicit reconciliation. */
if (legacyEventActorErrorHistoryDurable) {
await sealLegacyEventActorTurn();
}
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,
);
}
await releaseEventChildLease?.();
})
.catch((drainError) => {
logger.warn(
'[ResumableAgentController] Failed to record completed provider drain',
drainError,
);
});
} catch (error) {
logger.error('[ResumableAgentController] Initialization error:', error);
const initializationFailure = getInitializationFailure(error);
const streamStarted = res.headersSent;
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';
const completionPromise = streamStarted
? GenerationJobManager.completeJob(streamId, initializationError, jobCreatedAt, {
beforeErrorPublication: () =>
saveErrorTurn(req, {
conversationId,
endpointOption,
isNewConvo,
errorText: initializationError,
}),
})
: GenerationJobManager.completeJob(streamId, initializationError, jobCreatedAt);
initializationFinalized =
(await completionPromise.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,
);
});
}
await releaseEventChildLease?.();
}
};
module.exports = ResumableAgentController;