const { logger } = require('@librechat/data-schemas'); const { Constants, ViolationTypes, isEphemeralAgentId } = require('librechat-data-provider'); const { sendEvent, toPendingSteer, getViolationInfo, buildMessageFiles, getReferencedQuotes, resolveTitleTiming, GenerationJobManager, filterPersistableAbortContent, decrementPendingRequest, sanitizeMessageForTransmit, checkAndIncrementPendingRequest, isUnpersistedPreliminaryParent, resolveConversationAnchor, getAgentStartupTelemetry, acceptAgentStartupTelemetry, } = require('@librechat/api'); const { disposeClient, clientRegistry, requestDataMap } = require('~/server/cleanup'); const { getMCPRequestContext, cleanupMCPRequestContextForReq, } = require('~/server/services/MCPRequestContext'); const { handleAbortError } = require('~/server/middleware'); const { logViolation } = require('~/cache'); const { saveMessage, getMessages, getConvo } = require('~/models'); function createCloseHandler(abortController) { return function (manual) { if (!manual) { logger.debug('[AgentController] Request closed'); } if (!abortController) { return; } else if (abortController.signal.aborted) { return; } else if (abortController.requestCompleted) { return; } abortController.abort(); logger.debug('[AgentController] Request aborted on close'); }; } function resolveConversationCreatedAt({ userId, conversationId, isNewConvo }) { return resolveConversationAnchor({ isNewConversation: isNewConvo, loadConversation: () => getConvo(userId, conversationId), onLoadError: (error) => { logger.warn('[AgentController] Failed to resolve conversation timestamp anchor', { conversationId, error: error.message, }); }, }); } async function attachConversationCreatedAt(req, conversationId, conversationAnchorPromise) { req.body.conversationId = conversationId; const resolved = await conversationAnchorPromise; req.conversationCreatedAt = resolved.createdAt; if (resolved.conversation !== undefined) { req.resolvedConversation = resolved.conversation ?? null; } } function getPreliminaryResponseMessageId({ messageId, responseMessageId }) { if (typeof responseMessageId === 'string' && responseMessageId.length > 0) { return responseMessageId; } if (typeof messageId !== 'string' || messageId.length === 0) { return null; } return `${messageId.replace(/_+$/, '')}_`; } function getPreliminaryUserMessage( { messageId, parentMessageId, text, quotes, files, manualSkills, alwaysAppliedSkills }, conversationId, ) { if (typeof messageId !== 'string' || messageId.length === 0) { return null; } /** * Seed normalized quotes here too: if the user aborts before `sendMessage` * reaches `onStart` (during init/tool loading), `abortMiddleware` falls back * to this preliminary metadata, which must carry the excerpts so the stopped * turn keeps its `MessageQuotes`. */ const referencedQuotes = getReferencedQuotes(quotes); return { messageId, parentMessageId, conversationId, text, ...(referencedQuotes != null && { quotes: referencedQuotes }), // Persist the turn's uploaded files on this AWAITED preliminary write so they land on // job.metadata.userMessage BEFORE the run can reach its first interrupt. onStart's // later writes are fire-and-forget, so a fast approval could otherwise read the job // and resume an approved code/read-file tool without the paused turn's uploads. ...(Array.isArray(files) && files.length > 0 && { files }), // Carry skill selections so a HITL-resumed turn's reconstructed `requestMessage` // keeps its skill pills — the client's final handler replaces the user bubble from // this object, and they'd otherwise vanish until a full reload refetches the row. ...(Array.isArray(manualSkills) && manualSkills.length > 0 && { manualSkills }), ...(Array.isArray(alwaysAppliedSkills) && alwaysAppliedSkills.length > 0 && { alwaysAppliedSkills }), }; } function getRequestModelSpec(req, endpointOption) { const spec = endpointOption?.spec ?? req.body?.spec; if (typeof spec !== 'string' || spec.length === 0) { return; } const list = req.config?.modelSpecs?.list; if (!Array.isArray(list)) { return; } return list.find((modelSpec) => modelSpec?.name === spec); } function getModelSpecIconURL(modelSpec) { return modelSpec?.iconURL ?? modelSpec?.preset?.iconURL ?? modelSpec?.preset?.endpoint ?? ''; } function getEndpointIconURL(req, endpointOption) { const iconURL = endpointOption?.iconURL ?? getModelSpecIconURL(getRequestModelSpec(req, endpointOption)); return iconURL || undefined; } function getEndpointResponseModel(endpointOption) { return endpointOption?.modelOptions?.model || endpointOption?.model_parameters?.model; } function getAgentResponseModel(req, endpointOption) { const agentId = endpointOption?.agent_id || req.body?.agent_id; if (typeof agentId === 'string' && agentId.length > 0 && !isEphemeralAgentId(agentId)) { return agentId; } return getEndpointResponseModel(endpointOption); } async function finishResumableRequest(req, userId) { try { await cleanupMCPRequestContextForReq(req); } finally { await decrementPendingRequest(userId); } } 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; /** * 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++) { if (await GenerationJobManager.hasJob(streamId)) { return true; } await new Promise((resolve) => setTimeout(resolve, JOB_RECORD_WAIT_DELAY_MS)); } return GenerationJobManager.hasJob(streamId); } function rejectPreliminaryParentMessageId(res) { return res.status(409).json({ error: 'Cannot submit a follow-up while the selected parent response is still being saved. Please wait and try again.', }); } /** * 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); const { text, isRegenerate, endpointOption, conversationId: reqConversationId, isContinued = false, editedContent = null, parentMessageId = null, overrideParentMessageId = null, responseMessageId: editedResponseMessageId = null, } = req.body; const userId = req.user.id; const isNewConvo = !reqConversationId || reqConversationId === 'new'; const conversationId = isNewConvo ? crypto.randomUUID() : reqConversationId; const conversationAnchorPromise = resolveConversationCreatedAt({ userId, conversationId, isNewConvo, }); if ( await isUnpersistedPreliminaryParent({ userId, conversationId: reqConversationId, parentMessageId, getMessages, }) ) { startupTelemetry?.end('rejected'); return rejectPreliminaryParentMessageId(res); } /** 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; // 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. Fail-open on errors. const clientRequestId = req.body?.clientRequestId; let ownsIdempotencyClaim = false; if (clientRequestId) { let claim = null; try { claim = await GenerationJobManager.claimGeneration( userId, clientRequestId, streamId, conversationId, ); } catch (err) { // The claim itself could not be determined (store unavailable): fail open and proceed // as a fresh request rather than blocking the send. This is the ONLY fail-open path — // once a duplicate is confirmed below, an error must never fall through to a second // billed generation. logger.error( '[ResumableAgentController] Idempotency claim failed; proceeding without dedup', err, ); } if (claim?.claimed) { ownsIdempotencyClaim = true; } 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 jobExists = false; 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. jobExists = 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 res.status(503).json({ code: 'SERVER_NOT_READY', error: 'Generation is still starting. Please retry shortly.', }); } const claimAgeMs = Date.now() - (claim.existing.claimedAt ?? 0); if (!jobExists && 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 res.status(503).json({ code: 'SERVER_NOT_READY', error: 'Generation is still starting. Please retry shortly.', }); } // Job exists (live), or the grace elapsed with none (the original already completed // and was cleaned up, or the winner died): attach. A then-missing job recovers via // the client's subscribe 404 handler (refetch persisted messages) rather than an // indefinite readiness loop. logger.debug('[ResumableAgentController] Deduped retried start-generation request', { userId, clientRequestId, streamId: existingStreamId, }); startupTelemetry?.end('deduplicated'); return res.json({ streamId: existingStreamId, conversationId: claim.existing.conversationId, status: 'resumed', }); } } const { allowed, pendingRequests, limit } = await checkAndIncrementPendingRequest(userId); if (!allowed) { if (ownsIdempotencyClaim) { await GenerationJobManager.releaseGeneration(userId, clientRequestId).catch(() => {}); } const violationInfo = getViolationInfo(pendingRequests, limit); await logViolation(req, res, ViolationTypes.CONCURRENT, violationInfo, violationInfo.score); startupTelemetry?.end('rejected'); return res.status(429).json(violationInfo); } startupTelemetry?.mark('request_admitted'); let client = null; let jobCreatedAt; try { logger.debug(`[ResumableAgentController] Creating job`, { streamId, conversationId, reqConversationId, userId, }); const endpointIconURL = getEndpointIconURL(req, endpointOption); const responseModel = getAgentResponseModel(req, endpointOption); const preliminaryUserMessage = getPreliminaryUserMessage(req.body, conversationId); const preliminaryResponseMessageId = getPreliminaryResponseMessageId(req.body); const job = await GenerationJobManager.createJob(streamId, userId, conversationId, { startupTelemetry, initialMetadata: { conversationId, endpoint: endpointOption.endpoint, iconURL: endpointIconURL, model: responseModel, // Persist the originating agent so a HITL resume can refuse to rebuild this // paused run on a different agent (see resume.js). agent_id: endpointOption.agent_id ?? req.body?.agent_id, // Persist temporary-chat state so a HITL resume keeps the resumed response // non-persisted instead of trusting the resume request to re-send the flag. isTemporary: req.body?.isTemporary, responseMessageId: preliminaryResponseMessageId, userMessage: preliminaryUserMessage, }, }); startupTelemetry?.mark('job_created'); acceptAgentStartupTelemetry(req, streamId); startupTelemetry?.mark('metadata_persisted'); jobCreatedAt = job.createdAt; // Capture creation time to detect job replacement req._resumableStreamId = streamId; getMCPRequestContext(req, undefined, { cleanupOnResponse: false }); // 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 res.json({ streamId, conversationId, status: 'started' }); 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); 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; } await saveMessage( { userId: req?.user?.id, isTemporary: req?.body?.isTemporary, interfaceConfig: req?.config?.interfaceConfig, }, partialMessage, { context: 'api/server/controllers/agents/request.js - partial response on 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> }} */ const result = await initializeClient({ req, res, endpointOption, // Use the job's abort controller signal - allows abort via GenerationJobManager.abortJob() signal: job.abortController.signal, jobCreatedAt, }); startupTelemetry?.mark('client_initialized'); client = result.client; 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, ); }); startupTelemetry?.end('aborted'); try { await finishResumableRequest(req, userId); } finally { if (client) { disposeClient(client); } client = null; } 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 backgroundClientCleanupScheduled = false; 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, userMCPAuthMap: result.userMCPAuthMap, responseMessageId: editedResponseMessageId, progressOptions: { res: { write: () => true, end: () => {}, headersSent: false, writableEnded: false, }, }, }; const sendPromise = client.sendMessage(text, messageOptions); if (titleEligible && titleTiming === 'immediate') { immediateTitlePromise = addTitle(req, { text, 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 the // in-flight user-message / conversation save, then tear down WITHOUT saving a // partial response, emitting a terminal event, or completing the 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; } // BaseClient saved the response as completed (unfinished:false), but the turn // is paused awaiting a decision. Re-mark it unfinished so an expired / never- // resumed approval doesn't leave a "finished" response in history; the resume // path overwrites it with the full completed message on success. if (response?.messageId) { // Guard against a fast /resume: the user can approve the instant the // pending-action SSE lands, and resume.js can then claim + finalize — saving // the COMPLETED response — while we're still awaiting `response.databasePromise` // above. Marking the row unfinished now would clobber that completed content // with this stale pre-pause response. Only mark unfinished while the job is // STILL paused on THIS generation's action: a claim transitions it out of // `requires_action`, and a replacement bumps `createdAt`. Fail open on a read // error so a genuinely never-resumed approval isn't left looking "finished". let stillPaused = true; try { const liveJob = await GenerationJobManager.getJob(streamId); stillPaused = !!liveJob && liveJob.status === 'requires_action' && (client?.jobCreatedAt == null || liveJob.createdAt === client.jobCreatedAt); } catch (readErr) { logger.warn( '[ResumableAgentController] Pause unfinished-save liveness check failed; proceeding', readErr?.message ?? readErr, ); } if (!stillPaused) { logger.debug( `[ResumableAgentController] Skipping pause unfinished-save — ${streamId} already resumed/replaced`, ); } else { try { await saveMessage( { userId, isTemporary: req?.body?.isTemporary, interfaceConfig: req?.config?.interfaceConfig, }, { ...response, endpoint: endpointOption.endpoint, unfinished: true, user: userId, }, { context: 'api/server/controllers/agents/request.js - HITL pause (mark unfinished)', }, ); } catch (saveErr) { logger.error( '[ResumableAgentController] Failed to mark paused response unfinished', saveErr, ); } } } 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) { await decrementPendingRequest(userId); } if (client) { disposeClient(client); } logger.debug( `[ResumableAgentController] Turn paused for approval; awaiting resume: ${streamId}`, ); startupTelemetry?.end('paused'); return; } const messageId = response.messageId; 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 (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; } // Check abort state BEFORE calling completeJob (which triggers abort signal for cleanup) const wasAbortedBeforeComplete = job.abortController.signal.aborted; const shouldGenerateTitle = addTitle && parentMessageId === Constants.NO_PARENT && isNewConvo && !wasAbortedBeforeComplete; // Save user message BEFORE sending final event to avoid race condition // where client refetch happens before database is updated const reqCtx = { userId: req?.user?.id, isTemporary: req?.body?.isTemporary, interfaceConfig: req?.config?.interfaceConfig, }; if (!client.skipSaveUserMessage && userMessage) { await saveMessage(reqCtx, userMessage, { context: 'api/server/controllers/agents/request.js - resumable user message', }); } // 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. if (client.savedMessageIds && !client.savedMessageIds.has(messageId)) { await saveMessage( reqCtx, { ...response, user: userId, unfinished: wasAbortedBeforeComplete }, { context: 'api/server/controllers/agents/request.js - resumable response end' }, ); } // Check if our job was replaced by a new request before emitting // This prevents stale requests from emitting events to newer jobs const currentJob = await GenerationJobManager.getJob(streamId); const jobWasReplaced = !currentJob || currentJob.createdAt !== jobCreatedAt; if (jobWasReplaced) { logger.debug(`[ResumableAgentController] Skipping FINAL emit - job was replaced`, { streamId, originalCreatedAt: jobCreatedAt, currentCreatedAt: currentJob?.createdAt, }); // Discard the stale title from this replaced stream: cancel it and // unblock its persistence wait without letting it save (the newer job // owns the conversation now). titleAbortController.abort(); titleDiscardController.abort(); job.abortController.signal.removeEventListener('abort', abortTitleOnJobAbort); acceptsTitleEvents = false; resolveConvoReady(); // Still decrement pending request since we incremented at start await finishResumableRequest(req, userId); startupTelemetry?.end('replaced'); if (immediateTitlePromise) { immediateTitlePromise.finally(() => { if (client) { disposeClient(client); } }); } else if (client) { disposeClient(client); } return; } // If the user stopped this turn, cancel the title BEFORE unblocking its // persistence wait — otherwise resolving `convoReady` lets the title task // resume and save before the later abort runs. if (wasAbortedBeforeComplete) { 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; } // Steers that never reached an injection boundary (queued after the last // tool batch, or the run had none). The close-and-drain atomically stops // new enqueues first — a steer POST racing this finalization gets 404 // (client sends it as a normal message) instead of a 202 whose payload // completeJob would then silently clear. Reported on the final event so // the client converts them to queued follow-up messages. let pendingSteers; try { const leftoverSteers = await GenerationJobManager.steering.closeAndDrain( streamId, jobCreatedAt, ); if (leftoverSteers.length > 0) { pendingSteers = leftoverSteers.map(toPendingSteer); // Parked BEFORE the final event: a client with no live subscriber // recovers these via /chat/status (claim-on-read) within the // recovery TTL — the SSE copy alone is transient. await GenerationJobManager.steering.park( streamId, pendingSteers, { userId, tenantId: req.user?.tenantId, }, jobCreatedAt, ); } } catch (err) { logger.warn(`[ResumableAgentController] Failed to drain leftover steers`, err); } if (!wasAbortedBeforeComplete) { const finalEvent = { final: true, conversation, title: conversation.title, requestMessage: sanitizeMessageForTransmit(userMessage), responseMessage: { ...response }, ...(pendingSteers && { pendingSteers }), }; logger.debug(`[ResumableAgentController] Emitting FINAL event`, { streamId, wasAbortedBeforeComplete, userMessageId: userMessage?.messageId, responseMessageId: response?.messageId, conversationId: conversation?.conversationId, }); await GenerationJobManager.emitDone(streamId, finalEvent, jobCreatedAt); startupTelemetry?.end('completed_without_delta'); void GenerationJobManager.completeJob(streamId, undefined, jobCreatedAt).catch((err) => { logger.warn('[ResumableAgentController] Failed to finalize completed job', err); }); await finishResumableRequest(req, userId); } else { const finalEvent = { final: true, conversation, title: conversation.title, requestMessage: sanitizeMessageForTransmit(userMessage), responseMessage: { ...response, unfinished: true }, ...(pendingSteers && { pendingSteers }), }; logger.debug(`[ResumableAgentController] Emitting ABORTED FINAL event`, { streamId, wasAbortedBeforeComplete, userMessageId: userMessage?.messageId, responseMessageId: response?.messageId, conversationId: conversation?.conversationId, }); await GenerationJobManager.emitDone(streamId, finalEvent, jobCreatedAt); startupTelemetry?.end('aborted'); void GenerationJobManager.completeJob(streamId, 'Request aborted', jobCreatedAt).catch( (err) => { logger.warn('[ResumableAgentController] Failed to finalize aborted job', err); }, ); 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) { addTitle(req, { text, 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(); // Check if this was an abort (not a real error) const wasAborted = job.abortController.signal.aborted || error.message?.includes('abort'); if (wasAborted) { 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); // Close the steer queue BEFORE the error event reaches clients: a // steer POST racing this failure gets 404 (client queues or sends it) // instead of a 202 whose payload would vanish with the job. Text // recovery is client-side — acknowledged chips convert to queued. try { const erroredLeftovers = await GenerationJobManager.steering.closeAndDrain( streamId, jobCreatedAt, ); if (erroredLeftovers.length > 0) { // The error event is a bare string — park the acknowledged // steers so a reloaded/disconnected client can still recover // them via /chat/status instead of losing them with the queue. await GenerationJobManager.steering.park( streamId, erroredLeftovers.map(toPendingSteer), { userId, tenantId: req.user?.tenantId }, jobCreatedAt, ); } } catch (drainErr) { logger.warn( `[ResumableAgentController] Failed to close steer queue on error`, drainErr, ); } try { await GenerationJobManager.emitError( streamId, error.message || 'Generation failed', jobCreatedAt, ); } catch (notificationError) { logger.warn( '[ResumableAgentController] Failed to notify client of generation error', notificationError, ); } finally { startupTelemetry?.end('error', error); } await GenerationJobManager.completeJob(streamId, error.message, jobCreatedAt).catch( (completeErr) => { logger.warn( '[ResumableAgentController] completeJob failed during generation-error cleanup', completeErr, ); }, ); } try { await finishResumableRequest(req, userId); } finally { disposeBackgroundClient(); } // Don't continue to title generation after error/abort return; } }; // Start generation and handle any unhandled errors startGeneration().catch(async (err) => { logger.error( `[ResumableAgentController] Unhandled error in background generation: ${err.message}`, ); startupTelemetry?.end('error', err); await GenerationJobManager.completeJob(streamId, err.message, jobCreatedAt).catch( (completeErr) => { logger.warn( '[ResumableAgentController] completeJob failed during background-error cleanup', completeErr, ); }, ); try { await finishResumableRequest(req, userId); } finally { disposeBackgroundClient(); } }); } catch (error) { logger.error('[ResumableAgentController] Initialization error:', error); try { if (!res.headersSent) { res.status(500).json({ error: error.message || 'Failed to start generation' }); } else if (jobCreatedAt != null) { // JSON already sent, emit error to stream so client can receive it await GenerationJobManager.emitError( streamId, error.message || 'Failed to start generation', jobCreatedAt, ); } } catch (notificationError) { logger.warn( '[ResumableAgentController] Failed to notify client of initialization error', notificationError, ); } finally { startupTelemetry?.end('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.) if (jobCreatedAt != null) { await GenerationJobManager.completeJob(streamId, error.message, jobCreatedAt).catch( (completeErr) => { logger.warn( '[ResumableAgentController] completeJob failed during init-error cleanup', completeErr, ); }, ); } if (ownsIdempotencyClaim) { await GenerationJobManager.releaseGeneration(userId, clientRequestId).catch(() => {}); } await finishResumableRequest(req, userId); if (client) { disposeClient(client); } } }; /** * Agent Controller - Routes to ResumableAgentController for all requests. * The legacy non-resumable path is kept below but no longer used by default. */ const AgentController = async (req, res, next, initializeClient, addTitle) => { return ResumableAgentController(req, res, next, initializeClient, addTitle); }; /** * Legacy Non-resumable Agent Controller - Uses GenerationJobManager for abort handling. * Response is streamed directly to client via res, but abort state is managed centrally. * @deprecated Use ResumableAgentController instead */ const _LegacyAgentController = async (req, res, next, initializeClient, addTitle) => { const { text, isRegenerate, endpointOption, conversationId: reqConversationId, isContinued = false, editedContent = null, parentMessageId = null, overrideParentMessageId = null, responseMessageId: editedResponseMessageId = null, } = req.body; // 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 isNewConvo = !reqConversationId || reqConversationId === 'new'; const conversationId = isNewConvo ? crypto.randomUUID() : reqConversationId; const streamId = conversationId; let userMessage; let userMessageId; let responseMessageId; let client = null; let jobCreatedAt; let cleanupHandlers = []; // Match the same logic used for conversationId generation above const userId = req.user.id; if ( await isUnpersistedPreliminaryParent({ userId, conversationId: reqConversationId, parentMessageId, getMessages, }) ) { return rejectPreliminaryParentMessageId(res); } await attachConversationCreatedAt(req, { userId, conversationId, isNewConvo }); // Create handler to avoid capturing the entire parent scope let getReqData = (data = {}) => { for (let key in data) { if (key === 'userMessage') { userMessage = data[key]; userMessageId = data[key].messageId; } else if (key === 'responseMessageId') { responseMessageId = data[key]; } else if (key === 'promptTokens') { // Update job metadata with prompt tokens for abort handling GenerationJobManager.updateMetadata(streamId, { promptTokens: data[key] }, jobCreatedAt); } else if (key === 'sender') { GenerationJobManager.updateMetadata(streamId, { sender: data[key] }, jobCreatedAt); } // conversationId is pre-generated, no need to update from callback } }; // Create a function to handle final cleanup const performCleanup = async () => { logger.debug('[AgentController] Performing cleanup'); if (Array.isArray(cleanupHandlers)) { for (const handler of cleanupHandlers) { try { if (typeof handler === 'function') { handler(); } } catch (e) { logger.error('[AgentController] Error in cleanup handler', e); } } } // Complete the job in GenerationJobManager if (jobCreatedAt != null) { logger.debug('[AgentController] Completing job in GenerationJobManager'); await GenerationJobManager.completeJob(streamId, undefined, jobCreatedAt); } // Dispose client properly if (client) { disposeClient(client); } // Clear all references client = null; getReqData = null; userMessage = null; cleanupHandlers = null; // Clear request data map if (requestDataMap.has(req)) { requestDataMap.delete(req); } logger.debug('[AgentController] Cleanup completed'); }; try { let prelimAbortController = new AbortController(); const prelimCloseHandler = createCloseHandler(prelimAbortController); res.on('close', prelimCloseHandler); const removePrelimHandler = (manual) => { try { prelimCloseHandler(manual); res.removeListener('close', prelimCloseHandler); } catch (e) { logger.error('[AgentController] Error removing close listener', e); } }; cleanupHandlers.push(removePrelimHandler); /** @type {{ client: TAgentClient; userMCPAuthMap?: Record> }} */ const result = await initializeClient({ req, res, endpointOption, signal: prelimAbortController.signal, }); if (prelimAbortController.signal?.aborted) { prelimAbortController = null; throw new Error('Request was aborted before initialization could complete'); } else { prelimAbortController = null; removePrelimHandler(true); cleanupHandlers.pop(); } client = result.client; // Register client with finalization registry if available if (clientRegistry) { clientRegistry.register(client, { userId }, client); } // Store request data in WeakMap keyed by req object requestDataMap.set(req, { client }); // Create job in GenerationJobManager for abort handling // streamId === conversationId (pre-generated above) const job = await GenerationJobManager.createJob(streamId, userId, conversationId); jobCreatedAt = job.createdAt; client.jobCreatedAt = jobCreatedAt; // Store endpoint metadata for abort handling GenerationJobManager.updateMetadata( streamId, { endpoint: endpointOption.endpoint, iconURL: getEndpointIconURL(req, endpointOption), model: getAgentResponseModel(req, endpointOption), sender: client?.sender, }, jobCreatedAt, ); // Store content parts reference for abort if (client?.contentParts) { GenerationJobManager.setContentParts(streamId, client.contentParts, jobCreatedAt); } const closeHandler = createCloseHandler(job.abortController); res.on('close', closeHandler); cleanupHandlers.push(() => { try { res.removeListener('close', closeHandler); } catch (e) { logger.error('[AgentController] Error removing close listener', e); } }); /** * onStart callback - stores user message and response ID for abort handling */ const onStart = (userMsg, respMsgId, _isNewConvo) => { sendEvent(res, { message: userMsg, created: true }); userMessage = userMsg; userMessageId = userMsg.messageId; responseMessageId = respMsgId; // Store metadata for abort handling (conversationId is pre-generated) GenerationJobManager.updateMetadata( streamId, { responseMessageId: respMsgId, userMessage: { messageId: userMsg.messageId, parentMessageId: userMsg.parentMessageId, conversationId, text: userMsg.text, quotes: userMsg.quotes, }, }, jobCreatedAt, ); }; const messageOptions = { user: userId, onStart, getReqData, isContinued, isRegenerate, editedContent, conversationId, parentMessageId, abortController: job.abortController, overrideParentMessageId, isEdited: !!editedContent, userMCPAuthMap: result.userMCPAuthMap, responseMessageId: editedResponseMessageId, progressOptions: { res, }, }; let response = await client.sendMessage(text, messageOptions); // Extract what we need and immediately break reference const messageId = response.messageId; const endpoint = endpointOption.endpoint; response.endpoint = endpoint; // Store database promise locally const databasePromise = response.databasePromise; delete response.databasePromise; // Resolve database-related data const { conversation: convoData = {} } = await databasePromise; const conversation = { ...convoData }; conversation.title = conversation && !conversation.title ? null : conversation?.title || 'New Chat'; 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; } // Only send if not aborted if (!job.abortController.signal.aborted) { // Create a new response object with minimal copies const finalResponse = { ...response }; sendEvent(res, { final: true, conversation, title: conversation.title, requestMessage: sanitizeMessageForTransmit(userMessage), responseMessage: finalResponse, }); res.end(); // Save the message if needed if (client.savedMessageIds && !client.savedMessageIds.has(messageId)) { await saveMessage( { userId: req?.user?.id, isTemporary: req?.body?.isTemporary, interfaceConfig: req?.config?.interfaceConfig, }, { ...finalResponse, user: userId }, { context: 'api/server/controllers/agents/request.js - response end' }, ); } } // Edge case: sendMessage completed but abort happened during sendCompletion // We need to ensure a final event is sent else if (!res.headersSent && !res.finished) { logger.debug( '[AgentController] Handling edge case: `sendMessage` completed but aborted during `sendCompletion`', ); const finalResponse = { ...response }; finalResponse.error = true; sendEvent(res, { final: true, conversation, title: conversation.title, requestMessage: sanitizeMessageForTransmit(userMessage), responseMessage: finalResponse, error: { message: 'Request was aborted during completion' }, }); res.end(); } // Save user message if needed if (!client.skipSaveUserMessage) { await saveMessage( { userId: req?.user?.id, isTemporary: req?.body?.isTemporary, interfaceConfig: req?.config?.interfaceConfig, }, userMessage, { context: "api/server/controllers/agents/request.js - don't skip saving user message" }, ); } // Add title if needed - extract minimal data if (addTitle && parentMessageId === Constants.NO_PARENT && isNewConvo) { addTitle(req, { text, response: { ...response }, client, }) .then(() => { logger.debug('[AgentController] Title generation started'); }) .catch((err) => { logger.error('[AgentController] Error in title generation', err); }) .finally(() => { logger.debug('[AgentController] Title generation completed'); performCleanup(); }); } else { performCleanup(); } } catch (error) { // Handle error without capturing much scope handleAbortError(res, req, error, { conversationId, sender: client?.sender, messageId: responseMessageId, parentMessageId: overrideParentMessageId ?? userMessageId ?? parentMessageId, userMessageId, }) .catch((err) => { logger.error('[api/server/controllers/agents/request] Error in `handleAbortError`', err); }) .finally(() => { performCleanup(); }); } }; module.exports = AgentController;