mirror of
https://github.com/danny-avila/LibreChat.git
synced 2026-09-09 16:10:04 +00:00
* 🧭 fix: Carry Quoted Excerpts Through Mid-Run Steering "Add to chat" quote chips were dropped by every during-run steer path: the steer POST had no quotes concept, so a composer-origin steer left the chip staged (gluing onto the NEXT send) and a queued item steered into the live run lost its quotes silently. Quotes now ride the steer protocol end to end: - POST + admission: `quotes` on the steer body, normalized like the chat route's (getReferencedQuotes caps), part of the idempotency fingerprint only when present so pre-existing receipts still replay. - Injection: merged into the model-bound turn as Markdown blockquotes at both boundaries (text-only and media paths), mirroring prependQuotes. - Persistence + replay: the STEER content part stores `quotes` separately from the typed text; stampSteerPartMedia re-merges them per turn (even with resendFiles off) via the SDK's transient media stamp, with the quote block folded into the token budget. - UI: composer steers/interrupt-steers drain the chips (skill picks stay staged — they configure a NEW turn's run); SteerPart and the in-flight bubble render the same MessageQuotes reference blocks as user bubbles; queued/failed rows show a quote count; reconnect reseeds fall back to the server item's quotes when no local chip survives. - buildMessages keeps its zero-await path to the parallel context kickoff via a synchronous stamp-target probe. * 🧭 fix: Keep Quotes in the Client-Safe Steer Projection toPendingSteer is the projection behind resume-state pendingSteers, abort responses, and terminal leftover claims — dropping quotes there would lose them on exactly the recovery paths the reconnect reseed's server fallback relies on. * 🧪 test: In-Flight Steer Bubble Renders Carried Quotes * 🔁 fix: Re-Stage Quotes When a Pre-Quotes Replica Accepts the Steer Codex flagged the rolling-deploy window: an old replica 202s a quoted steer while dropping the excerpts, so the client cleared the chips for context the model never received. The 202 (fresh and receipt replay) now echoes quotesAccepted from the DURABLE item; a missing echo on a quote-bearing composer-origin steer re-stages the excerpts as composer chips — the pre-steer behavior, so they ride the next send instead of vanishing — and strips them from the surviving chip so a later terminal conversion cannot duplicate them. Queued-origin steers keep quotes on the item, whose restore paths already return it intact. The residual cross-version lost-ACK retry stays fail-closed as a 409 idempotency conflict (failed chip with retry controls). * 🔁 fix: Close the Remaining Cross-Version Quote-Loss Windows Codex round 2: - Send now of a quoted queued item against a pre-quotes replica now re-stages the excerpts too (the row is consumed and the words inject bare, so the composer is their only remaining home); the strip clears the chip's captured origin copy so reclaims and terminal conversions cannot duplicate them. - A quoted retry whose lost first ACK was accepted by a pre-quotes replica now REPLAYS that legacy receipt instead of 409ing: the stored fingerprint matching the quote-less hash of the same words proves the cross-version case, and the replayed 202's missing echo drives the re-stage. Different quotes against a quote-bearing receipt still conflict. - TSteerAppliedEvent.part gains the quotes field (typed SSE consumers). * 🧪 test: Drop the Stale Narrow SteerDrainOutput Alias The spec's local intersection re-declared injectedMessages with content: string, predating the SDK pin that declares the field natively (content: string | MessageContentComplex[]). Under CI's clean install the hook's BaseHookOutput is no longer assignable to that narrower alias; the plain PostToolBatchHookOutput is the correct type for every drain/boundary assertion. Verified against the published 3.6.16 dist and the local one. * 🔁 fix: Honor the Generation Owner's Quote Capability End to End Codex round 4: - steerQuotesCapable rides job metadata (createJob + HITL resume rewrite), mirroring preemptCapable's owner-recorded pattern: an upgraded admission replica no longer stores quotes — or claims them accepted — for a generation whose older owning drain would silently drop them at injection. The missing echo drives the client re-stage, and a later capable handover cannot double-deliver restored context. - Applied events reconcile dropped quotes: when a quote-less applied part settles a quote-bearing chip (the lost-202 ordering the ACK-echo path cannot see), resolveSteerChip and both reconnect settle paths re-stage the chip's excerpts before removing their only copy. mergeRestagedQuotes dedupe keeps every trigger idempotent for the same excerpts. * 🔁 fix: Re-Read Quote Capability at the Last Moment and Cap Restaged Chips Codex round 5: - A HITL resume rewrites steerQuotesCapable without changing the generation's createdAt, so the enqueue fence cannot see a capable-to-legacy handover landing during admission's awaits. Re-read the owner's flag immediately before item construction (paid for only by quote-bearing requests); the residual between re-read and enqueue commit matches preemptCapable's documented race. - mergeRestagedQuotes now respects the 10-quote contract with the staged chips winning: a restored tail that cannot ride the next send is dropped explicitly instead of rendering as a chip the submission would silently discard. MAX_QUOTE_COUNT moves to utils/steer as the single client source; QuoteButton imports it. * 🔁 fix: Steer Quote Coverage for Preflights, Memory, and Single-Scan Stamping Codex round 6: - Stored-message policy inspection now extracts steer-part quotes as quote fragments (path /content/N/quotes/M), so conversation import and shared link preflights inspect the newly persisted field exactly like top-level message.quotes. - The memory copy gets its own quote-merge stamp (text only, resendFiles false): formatAgentMessages ignores part.quotes, so without it a steer whose substance lives in its excerpt reached the chat model but never memory extraction. - collectSteerStampTargets replaces the boolean probe: buildMessages collects once and hands the targets to stampSteerPartMedia, keeping the zero-await fast path without scanning the history twice. * 🔁 fix: Redis Quote Plumbing, Conversion-Race Guard, and Quote-Bound Recovery Proof Codex round 7: - RedisJobStore.deserializeJob now restores steerQuotesCapable (the explicit mapper otherwise dropped it on every read, leaving quote steering inert in Redis deployments), with the round-trip spec extended. - Both Lua parked-steer projections (terminal close + generation replacement) forward item.quotes, matching toPendingSteer — a lost final no longer strips excerpts from durable recovery in Redis mode. - The no-echo restage reads the SURVIVING chip (reclaimRejectedChipQuotes): a terminal conversion that beat the delayed 202 already moved the quotes onto the queued follow-up, and re-staging them again double-delivered. Regression-tested with the conversion-before-ACK ordering. - RecoveredSteerPayload binds normalized, order-significant quotes (builder, validator, TS matcher, and the Lua decode+matcher): a stale client presenting the same recoverySteerId with altered or missing quotes cannot consume the parked source. Quote-less sources keep matching quote-less recoveries. * 🔁 fix: Execution-Bound Quote Capability with an Atomic Enqueue Predicate Codex round 8: - steerQuotesCapable becomes a transient assertion translated (at createJob and in ApprovalLifecycle.resolve) into steerQuotesExecutionId, valid only while it equals the LIVE providerExecutionId. A legacy replica winning a HITL resume rewrites the execution id without knowing the marker, so its stale assertion self-invalidates — a bare boolean could not be cleared by code that predates it. - The fenced enqueue evaluates that equality atomically (all three Redis scripts decode-and-strip like the existing preemptCapable normalization; both InMemory sites mirror it) and returns the persisted item, so the quotesAccepted echo reflects exactly what was stored even when a handover lands between admission's read and the commit. The last-moment re-read is gone — the transaction is the authority. - Tests: capable-resume re-binding, legacy-resume omit-not-clear invalidation, the admission-vs-handover race (capability read true, then execution rewritten before enqueue), and the Redis round-trip of the marker. * 🔁 fix: Full Redis Parking Coverage and Loss-Moment Quote Restaging Codex round 9: - The two remaining Redis parking projections (terminal status CAS and stale-running cleanup) forward item.quotes — every field-picked steer projection now carries them (audited: 2 Lua 'projected' + 2 Lua 'clientItem' + toPendingSteer). - The ordinary no-echo ACK no longer re-stages: the steer has not injected yet, so the quotes stay carried on the pending chip. A quote-less applied event re-stages them at the actual loss; a terminal leftover conversion carries them onto the recovered row, whose normal send delivers quotes on any server — re-staging at the ACK let that leftover auto-send bare text while the excerpts glued onto an unrelated draft. Only the settled receipt replay (already injected, no future event) reclaims immediately. * 🔁 fix: Legacy-Replayable Receipts with Separate Quote Identity Codex round 10: an upgraded-first receipt stored a quote-inclusive fingerprint no pre-quotes replica could recompute, so a lost-ACK retry routed through one 409'd already-accepted words with duplicate-send controls. The durable fingerprint reverts to the quote-independent 3-field hash — the one shape EVERY deployed version computes, replayable across a rolling deploy in both directions — and quote identity moves beside it as requestedQuotesFingerprint (of the REQUESTED quotes, pre any capability strip, so an incapable-owner acceptance still replays its own retries). Absent records (legacy-written or quote-less) accept any same-words retry, preserving the round-5 rule; present records must match exactly, keeping different-quotes clientSteerId reuse a 409 on quote-aware readers. Under the keep-on-chip client contract a legacy replay's missing echo is harmless — the excerpts stay carried on the pending chip. * 🧪 chore: Re-Trigger CI After Dropped Workflow Events
2399 lines
94 KiB
JavaScript
2399 lines
94 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,
|
|
} = require('@librechat/api');
|
|
const { disposeClient } = require('~/server/cleanup');
|
|
const {
|
|
getMCPRequestContext,
|
|
cleanupMCPRequestContextForReq,
|
|
} = require('~/server/services/MCPRequestContext');
|
|
const { logViolation } = require('~/cache');
|
|
const { recordScheduleOutcome, isScheduleLive } = require('~/server/services/Schedules');
|
|
const {
|
|
saveMessage,
|
|
getMessages,
|
|
getConvo,
|
|
isAgentTriggerPrincipalActive,
|
|
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 getPreliminaryUserMessage(
|
|
{ messageId, parentMessageId, text, quotes, files, manualSkills, alwaysAppliedSkills },
|
|
conversationId,
|
|
) {
|
|
if (typeof messageId !== 'string' || messageId.length === 0) {
|
|
return null;
|
|
}
|
|
|
|
/**
|
|
* Seed normalized quotes here too: if the user aborts before `sendMessage`
|
|
* reaches `onStart` (during init/tool loading), `abortMiddleware` falls back
|
|
* to this preliminary metadata, which must carry the excerpts so the stopped
|
|
* turn keeps its `MessageQuotes`.
|
|
*/
|
|
const referencedQuotes = getReferencedQuotes(quotes);
|
|
|
|
return {
|
|
messageId,
|
|
parentMessageId,
|
|
conversationId,
|
|
text,
|
|
...(referencedQuotes != null && { quotes: referencedQuotes }),
|
|
// Persist the turn's uploaded files on this AWAITED preliminary write so they land on
|
|
// job.metadata.userMessage BEFORE the run can reach its first interrupt. onStart's
|
|
// later writes are fire-and-forget, so a fast approval could otherwise read the job
|
|
// and resume an approved code/read-file tool without the paused turn's uploads.
|
|
...(Array.isArray(files) && files.length > 0 && { files }),
|
|
// Carry skill selections so a HITL-resumed turn's reconstructed `requestMessage`
|
|
// keeps its skill pills — the client's final handler replaces the user bubble from
|
|
// this object, and they'd otherwise vanish until a full reload refetches the row.
|
|
...(Array.isArray(manualSkills) && manualSkills.length > 0 && { manualSkills }),
|
|
...(Array.isArray(alwaysAppliedSkills) &&
|
|
alwaysAppliedSkills.length > 0 && { alwaysAppliedSkills }),
|
|
};
|
|
}
|
|
|
|
function getRequestModelSpec(req, endpointOption) {
|
|
const spec = endpointOption?.spec ?? req.body?.spec;
|
|
if (typeof spec !== 'string' || spec.length === 0) {
|
|
return;
|
|
}
|
|
|
|
const list = req.config?.modelSpecs?.list;
|
|
if (!Array.isArray(list)) {
|
|
return;
|
|
}
|
|
|
|
return list.find((modelSpec) => modelSpec?.name === spec);
|
|
}
|
|
|
|
function getModelSpecIconURL(modelSpec) {
|
|
return modelSpec?.iconURL ?? modelSpec?.preset?.iconURL ?? modelSpec?.preset?.endpoint ?? '';
|
|
}
|
|
|
|
function getEndpointIconURL(req, endpointOption) {
|
|
const iconURL =
|
|
endpointOption?.iconURL ?? getModelSpecIconURL(getRequestModelSpec(req, endpointOption));
|
|
return iconURL || undefined;
|
|
}
|
|
|
|
function getEndpointResponseModel(endpointOption) {
|
|
return endpointOption?.modelOptions?.model || endpointOption?.model_parameters?.model;
|
|
}
|
|
|
|
function getAgentResponseModel(req, endpointOption) {
|
|
const agentId = endpointOption?.agent_id || req.body?.agent_id;
|
|
if (typeof agentId === 'string' && agentId.length > 0 && !isEphemeralAgentId(agentId)) {
|
|
return agentId;
|
|
}
|
|
|
|
return getEndpointResponseModel(endpointOption);
|
|
}
|
|
|
|
async function finishResumableRequest(req, userId) {
|
|
try {
|
|
await cleanupMCPRequestContextForReq(req);
|
|
} finally {
|
|
if (req._scheduleConcurrencyExempt !== true) {
|
|
await decrementPendingRequest(userId);
|
|
}
|
|
}
|
|
}
|
|
|
|
function classifyScheduledFailure(error, aborted = false) {
|
|
if (aborted || error?.code === 'SCHEDULE_NO_LONGER_ACTIVE') {
|
|
return { status: 'interrupted', error: error?.message };
|
|
}
|
|
if (error?.message?.includes(ViolationTypes.TOKEN_BALANCE)) {
|
|
return { status: 'skipped_balance' };
|
|
}
|
|
return { status: 'error', error: error?.message || 'Generation failed' };
|
|
}
|
|
|
|
const JOB_RECORD_WAIT_ATTEMPTS = 5;
|
|
const JOB_RECORD_WAIT_DELAY_MS = 60;
|
|
|
|
// A winner writes its job record within a few ms of claiming; if a losing duplicate still
|
|
// sees no job within this window of the claim, the winner is still starting (retry rather
|
|
// than hand back a stream that would 404). Past it, a missing job means the original
|
|
// already completed and was cleaned up (attach and let the client refetch).
|
|
const IDEMPOTENCY_STARTUP_GRACE_MS = 5000;
|
|
const CLIENT_REQUEST_ID_PATTERN = /^[A-Za-z0-9:_-]{1,128}$/;
|
|
/** New-chat retries do not carry a conversation id, so derive the stream id
|
|
* from their stable per-submission id. This keeps both the dedupe key and the
|
|
* Redis hash slot identical across a lost-response retry. */
|
|
const NEW_CONVERSATION_IDEMPOTENCY_NAMESPACE = 'd7f2518c-94b8-4fe8-97ad-2d4bdb2c9f43';
|
|
|
|
function isValidGenerationClaim(value, streamId, conversationId, requireStarted = false) {
|
|
return (
|
|
value != null &&
|
|
typeof value === 'object' &&
|
|
value.streamId === streamId &&
|
|
value.conversationId === conversationId &&
|
|
Number.isSafeInteger(value.claimedAt) &&
|
|
value.claimedAt >= 0 &&
|
|
typeof value.claimToken === 'string' &&
|
|
value.claimToken.length > 0 &&
|
|
value.claimToken.length <= 128 &&
|
|
(value.generationProtocolVersion == null ||
|
|
value.generationProtocolVersion === 1 ||
|
|
value.generationProtocolVersion === GENERATION_PROTOCOL_V2) &&
|
|
(value.startedAt == null || (Number.isSafeInteger(value.startedAt) && value.startedAt >= 0)) &&
|
|
(!requireStarted || value.startedAt != null)
|
|
);
|
|
}
|
|
|
|
/** Pre-bridge servers wrote the legacy global key without a claim token and,
|
|
* for a new conversation, chose a random stream before claiming it. Accept
|
|
* only that tightly bounded legacy shape: existing conversations must still
|
|
* match the requested stream exactly; new-chat claims may point to the old
|
|
* random stream only when streamId === conversationId. Ownership is verified
|
|
* against the live job before attachment. */
|
|
function isValidLegacyGenerationClaim(value, streamId, isNewConvo) {
|
|
return (
|
|
value != null &&
|
|
typeof value === 'object' &&
|
|
typeof value.streamId === 'string' &&
|
|
value.streamId.length > 0 &&
|
|
value.streamId.length <= 512 &&
|
|
value.conversationId === value.streamId &&
|
|
(isNewConvo || value.streamId === streamId) &&
|
|
Number.isSafeInteger(value.claimedAt) &&
|
|
value.claimedAt >= 0 &&
|
|
value.claimToken == null &&
|
|
value.startedAt == null &&
|
|
(value.generationProtocolVersion == null || value.generationProtocolVersion === 1)
|
|
);
|
|
}
|
|
|
|
/** Store corruption must not turn a user-scoped idempotency claim into a
|
|
* pointer to another user's/tenant's live stream. Missing tenant metadata is
|
|
* kept as the explicit legacy case, but missing ownership never authorizes. */
|
|
function liveJobBelongsToRequester(job, user) {
|
|
return (
|
|
job?.metadata?.userId === user.id &&
|
|
(job.metadata?.tenantId == null || job.metadata.tenantId === user.tenantId)
|
|
);
|
|
}
|
|
|
|
/**
|
|
* Poll briefly for a job record to appear. A deduped retry that loses the idempotency
|
|
* claim must not be handed the winner's stream until its job exists, or the client's
|
|
* subscribe 404s terminally. The winner writes the record a few ms after claiming.
|
|
*/
|
|
async function waitForJobRecord(streamId) {
|
|
for (let attempt = 0; attempt < JOB_RECORD_WAIT_ATTEMPTS; attempt++) {
|
|
const job = await GenerationJobManager.getJob(streamId);
|
|
if (job) {
|
|
return job;
|
|
}
|
|
await new Promise((resolve) => setTimeout(resolve, JOB_RECORD_WAIT_DELAY_MS));
|
|
}
|
|
return GenerationJobManager.getJob(streamId);
|
|
}
|
|
|
|
/** The claimed generation already reached durable/terminal history, but its
|
|
* conversation stream id now belongs to no job or to a newer submission. A
|
|
* success shape with that streamId would attach the stale submission to the
|
|
* replacement, so tell the client to refetch without opening SSE. */
|
|
function sendSettledGeneration(
|
|
res,
|
|
streamId,
|
|
conversationId,
|
|
startupTelemetry,
|
|
generationProtocolVersion,
|
|
) {
|
|
startupTelemetry?.end('deduplicated');
|
|
if (generationProtocolVersion < GENERATION_PROTOCOL_V2) {
|
|
return sendGenerationJson(
|
|
res,
|
|
200,
|
|
{ streamId, conversationId, status: 'resumed' },
|
|
generationProtocolVersion,
|
|
);
|
|
}
|
|
return sendGenerationJson(
|
|
res,
|
|
200,
|
|
{ conversationId, status: 'settled' },
|
|
generationProtocolVersion,
|
|
);
|
|
}
|
|
|
|
function rejectPreliminaryParentMessageId(res, generationProtocolVersion) {
|
|
return sendGenerationJson(
|
|
res,
|
|
409,
|
|
{
|
|
code: 'PARENT_NOT_READY',
|
|
error:
|
|
'Cannot submit a follow-up while the selected parent response is still being saved. Please wait and try again.',
|
|
},
|
|
generationProtocolVersion,
|
|
);
|
|
}
|
|
|
|
function rejectMissingTriggerParentMessageId(res, generationProtocolVersion) {
|
|
return sendGenerationJson(
|
|
res,
|
|
404,
|
|
{
|
|
code: 'PARENT_NOT_FOUND',
|
|
error: 'The selected parent response is no longer available.',
|
|
},
|
|
generationProtocolVersion,
|
|
);
|
|
}
|
|
|
|
/**
|
|
* Resumable Agent Controller - Generation runs independently of HTTP connection.
|
|
* Returns streamId immediately, client subscribes separately via SSE.
|
|
*/
|
|
const ResumableAgentController = async (req, res, next, initializeClient, addTitle) => {
|
|
const startupTelemetry = getAgentStartupTelemetry(req);
|
|
let generationProtocolVersion = negotiateNewGenerationProtocol(req, GenerationJobManager);
|
|
const {
|
|
text,
|
|
isRegenerate,
|
|
endpointOption,
|
|
conversationId: reqConversationId,
|
|
isContinued = false,
|
|
editedContent = null,
|
|
parentMessageId = null,
|
|
overrideParentMessageId = null,
|
|
responseMessageId: editedResponseMessageId = null,
|
|
scheduleId: bodyScheduleId = null,
|
|
scheduledFor: bodyScheduledFor = null,
|
|
scheduleConfigRevision: bodyScheduleConfigRevision = null,
|
|
} = req.body;
|
|
|
|
const isScheduledFire = isScheduleFireRequest(req);
|
|
const scheduleId = isScheduledFire ? bodyScheduleId : null;
|
|
const scheduledFor = isScheduledFire ? bodyScheduledFor : null;
|
|
const scheduleConfigRevision = isScheduledFire ? bodyScheduleConfigRevision : undefined;
|
|
|
|
const userId = req.user.id;
|
|
const tenantId = req.user.tenantId;
|
|
const rawClientRequestId = req.body?.clientRequestId;
|
|
if (
|
|
rawClientRequestId != null &&
|
|
(typeof rawClientRequestId !== 'string' || !CLIENT_REQUEST_ID_PATTERN.test(rawClientRequestId))
|
|
) {
|
|
startupTelemetry?.end('rejected');
|
|
return sendGenerationJson(
|
|
res,
|
|
400,
|
|
{
|
|
code: 'INVALID_CLIENT_REQUEST_ID',
|
|
error: 'clientRequestId must be a 1-128 character identifier.',
|
|
},
|
|
generationProtocolVersion,
|
|
);
|
|
}
|
|
const clientRequestId = rawClientRequestId;
|
|
const 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;
|
|
};
|
|
|
|
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,
|
|
...(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;
|
|
|
|
const getReqData = (data = {}) => {
|
|
if (data.userMessage) {
|
|
userMessage = data.userMessage;
|
|
}
|
|
// conversationId is pre-generated, no need to update from callback
|
|
};
|
|
|
|
let immediateTitlePromise = null;
|
|
let trailingWritePromise = null;
|
|
let backgroundClientCleanupScheduled = false;
|
|
let terminalClaim = null;
|
|
let terminalClaimFinished = false;
|
|
let terminalPersistenceChecked = false;
|
|
let terminalWasAborted = false;
|
|
let preemptIncomplete = false;
|
|
/** A pause-row write failure is terminalized through the exact action/epoch
|
|
* barrier. Once that path starts, neither generic background error handler
|
|
* may call completeJob: the pause may already have been replaced by a newer
|
|
* action or generation by the time the persistence failure is observed. */
|
|
let pausePersistenceFailed = false;
|
|
let pausePersistenceFailureFinalized = false;
|
|
const finishOwnedTerminalClaim = async () => {
|
|
if (!terminalClaim || terminalClaimFinished) {
|
|
return;
|
|
}
|
|
try {
|
|
await GenerationJobManager.finishTerminalJob(terminalClaim);
|
|
} finally {
|
|
terminalClaimFinished = true;
|
|
}
|
|
};
|
|
/** Runs inside BaseClient immediately before it can start the completed
|
|
* response write. A lost claim returns false, and BaseClient skips that
|
|
* stale `unfinished:false` write entirely. The fallback invocation below
|
|
* supports test/custom clients that do not derive from BaseClient. */
|
|
const claimBeforeResponsePersistence = async () => {
|
|
if (terminalPersistenceChecked) {
|
|
return terminalClaim != null;
|
|
}
|
|
terminalPersistenceChecked = true;
|
|
if (client?.pendingApproval) {
|
|
// AgentClient installed a durable pause-persistence barrier in the
|
|
// running→requires_action CAS. BaseClient must not start its ordinary
|
|
// `unfinished:false` response write; the HITL branch persists the
|
|
// partial row as unfinished before releasing that barrier.
|
|
return false;
|
|
}
|
|
terminalWasAborted = job.abortController.signal.aborted;
|
|
const preemptStats = client?.run?.getPreemptStats?.();
|
|
preemptIncomplete =
|
|
(preemptStats?.emptyBoundaries ?? 0) > 0 ||
|
|
client?.run?.getHaltReason?.() === 'preempt_incomplete';
|
|
terminalClaim = await GenerationJobManager.claimTerminalJob(
|
|
streamId,
|
|
terminalWasAborted ? 'aborted' : 'complete',
|
|
undefined,
|
|
jobCreatedAt,
|
|
{ persistencePending: true },
|
|
);
|
|
return terminalClaim != null;
|
|
};
|
|
const disposeBackgroundClient = () => {
|
|
if (backgroundClientCleanupScheduled) {
|
|
return;
|
|
}
|
|
backgroundClientCleanupScheduled = true;
|
|
|
|
if (immediateTitlePromise) {
|
|
immediateTitlePromise.finally(() => {
|
|
if (client) {
|
|
disposeClient(client);
|
|
}
|
|
});
|
|
} else if (client) {
|
|
disposeClient(client);
|
|
}
|
|
};
|
|
|
|
// Start background generation immediately. The stream layer buffers and persists events
|
|
// until an SSE subscriber attaches, so generation no longer waits on subscriber readiness.
|
|
const startGeneration = async () => {
|
|
/** Immediate-mode title generation runs in parallel with the response, so
|
|
* the conversation row may not exist when the title resolves. `convoReady`
|
|
* resolves once the response (and thus the conversation) has been saved,
|
|
* gating the title's `saveConvo`. Declared here so both the success tail
|
|
* and the catch block can settle it and gate `disposeClient` on the title. */
|
|
let titleEventPromise = null;
|
|
let acceptsTitleEvents = true;
|
|
let resolveConvoReady;
|
|
const convoReady = new Promise((resolve) => {
|
|
resolveConvoReady = resolve;
|
|
});
|
|
/** Dedicated controller so a user Stop (or a replaced stream) cancels the
|
|
* in-flight title — kept separate from `job.abortController`, which
|
|
* `completeJob` also aborts on *successful* completion and would otherwise
|
|
* cancel a title that is merely slower than a short response. */
|
|
const titleAbortController = new AbortController();
|
|
/** Separate from `titleAbortController`: a user Stop cancels the in-flight
|
|
* title model call but keeps a title that already finished generating.
|
|
* Only a superseded/failed stream aborts this to discard such a title so it
|
|
* cannot clobber the conversation now owned by the newer run. */
|
|
const titleDiscardController = new AbortController();
|
|
const abortTitleOnJobAbort = () => titleAbortController.abort();
|
|
if (job.abortController.signal.aborted) {
|
|
titleAbortController.abort();
|
|
} else {
|
|
job.abortController.signal.addEventListener('abort', abortTitleOnJobAbort, { once: true });
|
|
}
|
|
const titleEligible =
|
|
addTitle && parentMessageId === Constants.NO_PARENT && isNewConvo && !req.body?.isTemporary;
|
|
const emitTitleEvent = ({ conversationId: titleConversationId, title }) => {
|
|
titleEventPromise = (async () => {
|
|
if (!acceptsTitleEvents || titleAbortController.signal.aborted) {
|
|
return;
|
|
}
|
|
const currentJob = await GenerationJobManager.getJob(streamId);
|
|
if (!currentJob || currentJob.createdAt !== jobCreatedAt) {
|
|
return;
|
|
}
|
|
if (titleAbortController.signal.aborted) {
|
|
return;
|
|
}
|
|
await GenerationJobManager.emitChunk(
|
|
streamId,
|
|
{
|
|
event: 'title',
|
|
data: {
|
|
conversationId: titleConversationId,
|
|
title,
|
|
},
|
|
},
|
|
{ expectedCreatedAt: jobCreatedAt },
|
|
);
|
|
})().catch((err) => {
|
|
logger.error('[ResumableAgentController] Error emitting title event', err);
|
|
});
|
|
return titleEventPromise;
|
|
};
|
|
|
|
try {
|
|
const onStart = (userMsg, respMsgId, _isNewConvo) => {
|
|
userMessage = userMsg;
|
|
|
|
// Store userMessage and responseMessageId upfront for resume capability
|
|
GenerationJobManager.updateMetadata(
|
|
streamId,
|
|
{
|
|
responseMessageId: respMsgId,
|
|
userMessage: {
|
|
messageId: userMsg.messageId,
|
|
parentMessageId: userMsg.parentMessageId,
|
|
conversationId: userMsg.conversationId,
|
|
text: userMsg.text,
|
|
quotes: userMsg.quotes,
|
|
// Persist the turn's uploaded files here (authoritative job metadata) so a
|
|
// HITL resume sources them from the job, not the user DB row — which the
|
|
// approval prompt can race (the row save may still be in flight when a fast
|
|
// /resume reads it). Without this an approved tool run can rebuild without the
|
|
// paused turn's files.
|
|
...(Array.isArray(req.body?.files) &&
|
|
req.body.files.length > 0 && { files: req.body.files }),
|
|
// Skill selections aren't on `userMsg` yet at onStart (BaseClient adds them
|
|
// later), so source them from the request — otherwise this update overwrites
|
|
// the preliminary metadata and a HITL-resumed turn loses its skill pills.
|
|
...(Array.isArray(req.body?.manualSkills) &&
|
|
req.body.manualSkills.length > 0 && { manualSkills: req.body.manualSkills }),
|
|
...(Array.isArray(req.body?.alwaysAppliedSkills) &&
|
|
req.body.alwaysAppliedSkills.length > 0 && {
|
|
alwaysAppliedSkills: req.body.alwaysAppliedSkills,
|
|
}),
|
|
},
|
|
},
|
|
jobCreatedAt,
|
|
).catch((err) => {
|
|
logger.error('[ResumableAgentController] Failed to persist start metadata', err);
|
|
});
|
|
|
|
GenerationJobManager.emitChunk(
|
|
streamId,
|
|
{
|
|
created: true,
|
|
// Skill selections aren't on `userMessage` yet at onStart (BaseClient adds
|
|
// them later), so attach them from the request — this is the message
|
|
// `trackUserMessage` persists as the authoritative job.metadata.userMessage,
|
|
// and it's what the live client renders the user bubble from.
|
|
message: {
|
|
...userMessage,
|
|
// Carry files so trackUserMessage (the authoritative writer) persists them on
|
|
// job.metadata.userMessage for a HITL resume (see the updateMetadata above).
|
|
...(Array.isArray(req.body?.files) &&
|
|
req.body.files.length > 0 && { files: req.body.files }),
|
|
...(Array.isArray(req.body?.manualSkills) &&
|
|
req.body.manualSkills.length > 0 && { manualSkills: req.body.manualSkills }),
|
|
...(Array.isArray(req.body?.alwaysAppliedSkills) &&
|
|
req.body.alwaysAppliedSkills.length > 0 && {
|
|
alwaysAppliedSkills: req.body.alwaysAppliedSkills,
|
|
}),
|
|
},
|
|
streamId,
|
|
},
|
|
{ expectedCreatedAt: jobCreatedAt },
|
|
).catch((err) => {
|
|
logger.error('[ResumableAgentController] Failed to queue created event', err);
|
|
});
|
|
};
|
|
|
|
const messageOptions = {
|
|
user: userId,
|
|
onStart,
|
|
getReqData,
|
|
isContinued,
|
|
isRegenerate,
|
|
editedContent,
|
|
conversationId,
|
|
parentMessageId,
|
|
abortController: job.abortController,
|
|
overrideParentMessageId,
|
|
isEdited: !!editedContent,
|
|
beforeResponsePersistence: claimBeforeResponsePersistence,
|
|
userMCPAuthMap: result.userMCPAuthMap,
|
|
responseMessageId: editedResponseMessageId,
|
|
preallocatedUserMessageId,
|
|
preallocatedResponseMessageId,
|
|
progressOptions: {
|
|
res: {
|
|
write: () => true,
|
|
end: () => {},
|
|
headersSent: false,
|
|
writableEnded: false,
|
|
},
|
|
},
|
|
};
|
|
|
|
const sendPromise = client.sendMessage(text, messageOptions);
|
|
|
|
if (titleEligible && titleTiming === 'immediate') {
|
|
immediateTitlePromise = addTitle(req, {
|
|
text: text || getAttachmentTitleText(req.body.files),
|
|
conversationId,
|
|
client,
|
|
immediate: true,
|
|
convoReady,
|
|
signal: titleAbortController.signal,
|
|
discardSignal: titleDiscardController.signal,
|
|
onTitleGenerated: emitTitleEvent,
|
|
}).catch((err) => {
|
|
logger.error('[ResumableAgentController] Error in immediate title generation', err);
|
|
});
|
|
}
|
|
|
|
const response = await sendPromise;
|
|
|
|
// HITL: the turn paused for human review (see AgentClient.handleRunInterrupt).
|
|
// The job is already `requires_action` with the pending action persisted and
|
|
// emitted to the client; the resume route owns finishing this turn. Settle and
|
|
// verify the required unfinished history, then tear down without publishing a
|
|
// terminal event or completing a successfully persisted paused job.
|
|
if (client?.pendingApproval) {
|
|
if (response?.databasePromise) {
|
|
try {
|
|
await response.databasePromise;
|
|
} catch (dbErr) {
|
|
logger.error(
|
|
'[ResumableAgentController] Error settling databasePromise on HITL pause',
|
|
dbErr,
|
|
);
|
|
}
|
|
delete response.databasePromise;
|
|
}
|
|
const pauseActionId = client.pendingApproval.actionId;
|
|
const pauseCreatedAt = client.jobCreatedAt ?? jobCreatedAt;
|
|
const ownsPausePersistence = await GenerationJobManager.approvals.ownsPausePersistence(
|
|
streamId,
|
|
pauseActionId,
|
|
pauseCreatedAt,
|
|
);
|
|
if (ownsPausePersistence) {
|
|
try {
|
|
/** BaseClient awaits its first user/conversation write before the
|
|
* pause hook, but deliberately swallows a failed/falsy user save
|
|
* and may still record the id locally. Re-save idempotently for
|
|
* every ordinary user turn before exposing the approval. */
|
|
if (!client?.skipSaveUserMessage) {
|
|
if (!userMessage) {
|
|
throw new Error('User message was unavailable before HITL pause');
|
|
}
|
|
if (
|
|
typeof client.saveMessageToDatabase === 'function' &&
|
|
typeof client.getSaveOptions === 'function'
|
|
) {
|
|
/** Retry through BaseClient so a failure before its original
|
|
* saveConvo is repaired along with the message row. Direct
|
|
* saveMessage alone cannot recreate that conversation. */
|
|
const savedUserTurn = await client.saveMessageToDatabase(
|
|
userMessage,
|
|
client.getSaveOptions(),
|
|
userId,
|
|
);
|
|
if (!savedUserTurn?.message) {
|
|
throw new Error('User message could not be persisted before HITL pause');
|
|
}
|
|
if (!client.skipSaveConvo && !savedUserTurn.conversation) {
|
|
throw new Error('Conversation could not be persisted before HITL pause');
|
|
}
|
|
} else {
|
|
// Custom clients used by integrations/tests may not inherit BaseClient.
|
|
const savedUserMessage = await saveMessage(
|
|
{
|
|
userId,
|
|
isTemporary:
|
|
req?._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();
|
|
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 the user stopped this turn — or an empty preempt boundary truncated
|
|
// it, which persists under the same honest `unfinished` contract — cancel
|
|
// the title BEFORE unblocking its persistence wait; otherwise resolving
|
|
// `convoReady` lets the title task resume and save before the later abort runs.
|
|
if (terminalWasAborted || preemptIncomplete) {
|
|
titleAbortController.abort();
|
|
} else {
|
|
job.abortController.signal.removeEventListener('abort', abortTitleOnJobAbort);
|
|
}
|
|
|
|
// The conversation row now exists and this stream is authoritative; allow
|
|
// any in-flight immediate title generation to persist (saveConvo uses noUpsert).
|
|
resolveConvoReady();
|
|
acceptsTitleEvents = false;
|
|
|
|
if (titleEventPromise) {
|
|
await titleEventPromise;
|
|
}
|
|
|
|
let scheduleCompletionError;
|
|
if (terminalWasAborted) {
|
|
scheduleCompletionError = 'Scheduled run was stopped';
|
|
} else if (preemptIncomplete) {
|
|
scheduleCompletionError = 'Scheduled run was interrupted before completion';
|
|
}
|
|
await settleScheduledRun({
|
|
status: terminalWasAborted || preemptIncomplete ? 'interrupted' : 'success',
|
|
...(scheduleCompletionError != null && { error: scheduleCompletionError }),
|
|
});
|
|
|
|
let terminalPublicationStarted = false;
|
|
try {
|
|
const pendingSteers = terminalClaim.drainedSteers.map(toPendingSteer);
|
|
const finalEvent = {
|
|
final: true,
|
|
conversation,
|
|
title: conversation.title,
|
|
requestMessage: sanitizeMessageForTransmit(userMessage),
|
|
responseMessage: {
|
|
...response,
|
|
...((terminalWasAborted || preemptIncomplete) && { unfinished: true }),
|
|
},
|
|
...(pendingSteers.length > 0 && { pendingSteers }),
|
|
};
|
|
|
|
logger.debug(
|
|
terminalWasAborted
|
|
? `[ResumableAgentController] Emitting ABORTED FINAL event`
|
|
: `[ResumableAgentController] Emitting FINAL event`,
|
|
{
|
|
streamId,
|
|
wasAbortedBeforeComplete: terminalWasAborted,
|
|
userMessageId: userMessage?.messageId,
|
|
responseMessageId: response?.messageId,
|
|
conversationId: conversation?.conversationId,
|
|
},
|
|
);
|
|
|
|
terminalPublicationStarted = true;
|
|
const publication = await GenerationJobManager.publishTerminalClaim(
|
|
terminalClaim,
|
|
finalEvent,
|
|
);
|
|
let terminalOutcome = 'completed_without_delta';
|
|
if (publication.persistenceFailed) {
|
|
terminalOutcome = 'error';
|
|
} else if (terminalWasAborted) {
|
|
terminalOutcome = 'aborted';
|
|
}
|
|
startupTelemetry?.end(terminalOutcome);
|
|
} catch (terminalError) {
|
|
/** A failure while constructing the payload happened after this
|
|
* controller's terminal CAS but before the manager could durably
|
|
* settle it. Publish conservative reconciliation immediately. Once
|
|
* publication starts, the manager either stores the payload or owns
|
|
* its bounded recovery marker, so retrying with a different payload
|
|
* here would only risk duplicate delivery. */
|
|
if (!terminalPublicationStarted) {
|
|
try {
|
|
await GenerationJobManager.publishTerminalClaim(terminalClaim, null);
|
|
} catch (reconcileError) {
|
|
logger.warn(
|
|
'[ResumableAgentController] Failed to publish terminal persistence reconciliation',
|
|
reconcileError,
|
|
);
|
|
}
|
|
}
|
|
throw terminalError;
|
|
} finally {
|
|
// Pair every successful claim even when final-event construction or
|
|
// transport publication throws. Cleanup is epoch/runtime guarded.
|
|
await finishOwnedTerminalClaim();
|
|
}
|
|
await finishResumableRequest(req, userId);
|
|
|
|
if (titleTiming === 'immediate') {
|
|
// Title was fired in parallel above (if eligible); a stopped turn already
|
|
// aborted it before `resolveConvoReady`. Defer disposal until it settles
|
|
// so the run/req aren't torn down mid-generation.
|
|
if (immediateTitlePromise) {
|
|
immediateTitlePromise.finally(() => {
|
|
if (client) {
|
|
disposeClient(client);
|
|
}
|
|
});
|
|
} else if (client) {
|
|
disposeClient(client);
|
|
}
|
|
} else if (shouldGenerateTitle) {
|
|
trailingWritePromise = addTitle(req, {
|
|
text: text || getAttachmentTitleText(req.body.files),
|
|
response: { ...response },
|
|
client,
|
|
})
|
|
.catch((err) => {
|
|
logger.error('[ResumableAgentController] Error in title generation', err);
|
|
})
|
|
.finally(() => {
|
|
if (client) {
|
|
disposeClient(client);
|
|
}
|
|
});
|
|
} else {
|
|
if (client) {
|
|
disposeClient(client);
|
|
}
|
|
}
|
|
} catch (error) {
|
|
// Any failure (user Stop, or a preflight/quota failure before the run is
|
|
// even created) must cancel the title and unblock its waits: the title's
|
|
// `_waitForRun` would otherwise never resolve, deferring client disposal
|
|
// until the 45s title timeout, and no title should persist for a failed turn.
|
|
titleAbortController.abort();
|
|
titleDiscardController.abort();
|
|
job.abortController.signal.removeEventListener('abort', abortTitleOnJobAbort);
|
|
acceptsTitleEvents = false;
|
|
resolveConvoReady();
|
|
|
|
// Once this controller owns terminal persistence, no competing error
|
|
// transition can win. Settle its pending marker with conservative
|
|
// reconciliation on any required-write/final-construction failure,
|
|
// then release exactly that claim.
|
|
let ownsScheduledFailure = false;
|
|
if (terminalClaim && !terminalClaimFinished) {
|
|
ownsScheduledFailure = true;
|
|
try {
|
|
await GenerationJobManager.publishTerminalClaim(terminalClaim, null);
|
|
} catch (publishError) {
|
|
logger.warn(
|
|
'[ResumableAgentController] Failed to publish terminal persistence reconciliation',
|
|
publishError,
|
|
);
|
|
} finally {
|
|
await finishOwnedTerminalClaim().catch((finishError) => {
|
|
logger.warn(
|
|
'[ResumableAgentController] Failed to finish terminal persistence claim',
|
|
finishError,
|
|
);
|
|
});
|
|
}
|
|
logger.error(
|
|
`[ResumableAgentController] Terminal persistence failed for ${streamId}:`,
|
|
error,
|
|
);
|
|
startupTelemetry?.end('error', error);
|
|
} else if (pausePersistenceFailed) {
|
|
ownsScheduledFailure = pausePersistenceFailureFinalized;
|
|
// failPausePersistence owns the only legal requires_action -> error
|
|
// transition for this exact action/epoch. Never fall through to
|
|
// completeJob, which could race a newer action or replacement job.
|
|
logger.error(
|
|
`[ResumableAgentController] Pause persistence failed for ${streamId}:`,
|
|
error,
|
|
);
|
|
startupTelemetry?.end('error', error);
|
|
} else if (job.abortController.signal.aborted || error.message?.includes('abort')) {
|
|
ownsScheduledFailure = true;
|
|
logger.debug(`[ResumableAgentController] Generation aborted for ${streamId}`);
|
|
startupTelemetry?.end('aborted');
|
|
// abortJob already handled emitDone and completeJob
|
|
} else {
|
|
logger.error(`[ResumableAgentController] Generation error for ${streamId}:`, error);
|
|
const generationError = error.message || 'Generation failed';
|
|
try {
|
|
// completeJob first wins running -> error and atomically parks
|
|
// steers, then publishes. A competing abort/pause emits nothing.
|
|
ownsScheduledFailure =
|
|
(await GenerationJobManager.completeJob(streamId, generationError, jobCreatedAt)) ===
|
|
true;
|
|
} catch (completeErr) {
|
|
logger.warn(
|
|
'[ResumableAgentController] completeJob failed during generation-error cleanup',
|
|
completeErr,
|
|
);
|
|
} finally {
|
|
startupTelemetry?.end('error', error);
|
|
}
|
|
}
|
|
|
|
if (ownsScheduledFailure && !scheduleTerminalOutcomeRecorded) {
|
|
const scheduledFailure = classifyScheduledFailure(
|
|
error,
|
|
job.abortController.signal.aborted,
|
|
);
|
|
await settleScheduledRun(scheduledFailure);
|
|
}
|
|
|
|
try {
|
|
await finishResumableRequest(req, userId);
|
|
} finally {
|
|
disposeBackgroundClient();
|
|
}
|
|
|
|
// Don't continue to title generation after error/abort
|
|
return;
|
|
}
|
|
};
|
|
|
|
// Start generation and handle any unhandled errors
|
|
void startGeneration()
|
|
.catch(async (err) => {
|
|
logger.error(
|
|
`[ResumableAgentController] Unhandled error in background generation: ${err.message}`,
|
|
);
|
|
startupTelemetry?.end('error', err);
|
|
let errorFinalized = false;
|
|
if (!pausePersistenceFailed) {
|
|
errorFinalized =
|
|
(await GenerationJobManager.completeJob(streamId, err.message, jobCreatedAt).catch(
|
|
(completeErr) => {
|
|
logger.warn(
|
|
'[ResumableAgentController] completeJob failed during background-error cleanup',
|
|
completeErr,
|
|
);
|
|
return false;
|
|
},
|
|
)) === true;
|
|
}
|
|
if (
|
|
(errorFinalized || (pausePersistenceFailed && pausePersistenceFailureFinalized)) &&
|
|
!scheduleTerminalOutcomeRecorded
|
|
) {
|
|
await settleScheduledRun(classifyScheduledFailure(err));
|
|
}
|
|
try {
|
|
await finishResumableRequest(req, userId);
|
|
} finally {
|
|
disposeBackgroundClient();
|
|
}
|
|
})
|
|
.finally(async () => {
|
|
await Promise.allSettled([immediateTitlePromise, trailingWritePromise].filter(Boolean));
|
|
if (providerExecutionId) {
|
|
await GenerationJobManager.markProviderExecutionDrained?.(
|
|
streamId,
|
|
jobCreatedAt,
|
|
providerExecutionId,
|
|
);
|
|
}
|
|
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);
|
|
try {
|
|
if (!res.headersSent) {
|
|
if (error?.code === 'GENERATION_PREDECESSOR_MISMATCH') {
|
|
const currentJob = error.currentJob;
|
|
const currentStatus = currentJob?.status;
|
|
if (isTriggerContinuation && currentJob?.active === true) {
|
|
res.set('Retry-After', '1');
|
|
sendGenerationJson(
|
|
res,
|
|
409,
|
|
{
|
|
code: 'PARENT_NOT_READY',
|
|
error: 'Another generation became active before the continuation could start.',
|
|
},
|
|
generationProtocolVersion,
|
|
);
|
|
} else {
|
|
const predecessorVerified =
|
|
currentJob != null &&
|
|
Number.isSafeInteger(currentJob.createdAt) &&
|
|
currentJob.createdAt >= 0 &&
|
|
currentJob.verified !== false;
|
|
sendGenerationJson(
|
|
res,
|
|
409,
|
|
{
|
|
status: 'predecessor_mismatch',
|
|
code: 'GENERATION_PREDECESSOR_MISMATCH',
|
|
error: predecessorVerified
|
|
? 'A newer generation became current before this request could start.'
|
|
: 'The prior generation could not be verified. Please retry.',
|
|
streamId,
|
|
conversationId: currentJob?.conversationId ?? conversationId,
|
|
generationCreatedAt: currentJob?.createdAt,
|
|
predecessorVerified,
|
|
active:
|
|
typeof currentJob?.active === 'boolean'
|
|
? currentJob.active
|
|
: currentStatus === 'running' || currentStatus === 'requires_action',
|
|
},
|
|
generationProtocolVersion,
|
|
);
|
|
}
|
|
} else if (error?.code === 'RECOVERY_PAYLOAD_MISMATCH') {
|
|
sendGenerationJson(
|
|
res,
|
|
409,
|
|
{
|
|
code: 'RECOVERY_PAYLOAD_MISMATCH',
|
|
error: 'The queued message changed before it could be recovered. Please retry.',
|
|
},
|
|
generationProtocolVersion,
|
|
);
|
|
} else if (initializationFailure) {
|
|
sendGenerationJson(
|
|
res,
|
|
initializationFailure.status,
|
|
initializationFailure,
|
|
generationProtocolVersion,
|
|
);
|
|
} else {
|
|
sendGenerationJson(
|
|
res,
|
|
500,
|
|
{ error: error.message || 'Failed to start generation' },
|
|
generationProtocolVersion,
|
|
);
|
|
}
|
|
}
|
|
} catch (notificationError) {
|
|
logger.warn(
|
|
'[ResumableAgentController] Failed to send initialization error response',
|
|
notificationError,
|
|
);
|
|
} finally {
|
|
startupTelemetry?.end(
|
|
error?.code === 'GENERATION_PREDECESSOR_MISMATCH' ? 'deduplicated' : 'error',
|
|
error,
|
|
);
|
|
}
|
|
// Finalize THIS failed job before releasing the idempotency claim. Releasing first would
|
|
// let the client's retry win the same key and createJob() the same streamId while we are
|
|
// still here. The generation guard is defense-in-depth around that ordering. A
|
|
// completeJob() rejection (store hiccup) must NOT skip the
|
|
// release + pending-request decrement below, or the retry stays wedged behind the claim
|
|
// and the concurrency slot leaks — so swallow its error. (A failed completeJob did not
|
|
// finalize anything, so releasing afterward can't let it abort a later replacement.)
|
|
let initializationFinalized = jobCreatedAt == null;
|
|
if (jobCreatedAt != null) {
|
|
const initializationError = initializationFailure
|
|
? JSON.stringify(initializationFailure)
|
|
: error.message || 'Failed to start generation';
|
|
initializationFinalized =
|
|
(await GenerationJobManager.completeJob(streamId, initializationError, jobCreatedAt).catch(
|
|
(completeErr) => {
|
|
logger.warn(
|
|
'[ResumableAgentController] completeJob failed during init-error cleanup',
|
|
completeErr,
|
|
);
|
|
return false;
|
|
},
|
|
)) === true;
|
|
}
|
|
if (initializationFinalized && !scheduleTerminalOutcomeRecorded) {
|
|
await settleScheduledRun(classifyScheduledFailure(error));
|
|
}
|
|
if (ownedIdempotencyClaim) {
|
|
await GenerationJobManager.releaseGeneration(
|
|
userId,
|
|
clientRequestId,
|
|
streamId,
|
|
ownedIdempotencyClaim,
|
|
).catch(() => {});
|
|
}
|
|
await finishResumableRequest(req, userId);
|
|
if (client) {
|
|
disposeClient(client);
|
|
}
|
|
if (jobCreatedAt != null && providerExecutionId) {
|
|
await GenerationJobManager.markProviderExecutionDrained?.(
|
|
streamId,
|
|
jobCreatedAt,
|
|
providerExecutionId,
|
|
).catch((drainError) => {
|
|
logger.warn(
|
|
'[ResumableAgentController] Failed to record initialization-error provider drain',
|
|
drainError,
|
|
);
|
|
});
|
|
}
|
|
await releaseEventChildLease?.();
|
|
}
|
|
};
|
|
|
|
module.exports = ResumableAgentController;
|