mirror of
https://github.com/danny-avila/LibreChat.git
synced 2026-08-04 14:57:42 +00:00
Address all 6 findings from Codex review of 912d75bff:
- Exclude the resumed run from capacity checks (P2): reserveScheduledResume discounts
this occurrence's own started row (isOccurrenceStarted) so a self-active row (from a
transient pause-bookkeeping failure) doesn't push the global count over the cap and
block resuming the same occurrence.
- Preserve the lease marker while deleting (P2): markScheduleDeleting no longer unsets
leaseUntil/leaseBy, so a fire that already reserved a started row can prove
(holdsLease) it still owns the lease and roll back its own unposted row instead of
leaving a ghost.
- Treat a live lease as undrained (P2): eraseScheduleIfDrained also refuses to erase
while a LIVE lease is held (a worker claimed but hasn't reserved yet), so it can't
erase out from under a worker about to insert a reservation.
- Release own lease after owner-edit supersedes fire (P2): the superseded rollback path
now releaseLeaseByHolder(leaseBy) since advance() is token-fenced and no-ops after an
edit — otherwise the edited schedule / Run now reads 'already in progress' until TTL.
- Preserve early-abort jobs as aborted (P2): the createJob->liveness early abort uses
abortJob (stored 'aborted' -> reconcile 'interrupted') instead of completeJob, whose
preserved 'complete' the reconciler would map to 'success'.
- Show default weekly schedules as weekly (P3, client): describeCadence renders a
weekly cadence without daysOfWeek on the server's default weekly day instead of
falling through to the daily label.
1575 lines
62 KiB
JavaScript
1575 lines
62 KiB
JavaScript
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,
|
|
isScheduleFireRequest,
|
|
} = 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 { recordScheduleOutcome, isScheduleLive } = require('~/server/services/Schedules');
|
|
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 toValidISOString(value) {
|
|
if (value == null) {
|
|
return null;
|
|
}
|
|
|
|
const date = value instanceof Date ? value : new Date(value);
|
|
return Number.isNaN(date.getTime()) ? null : date.toISOString();
|
|
}
|
|
|
|
async function resolveConversationCreatedAt({ userId, conversationId, isNewConvo }) {
|
|
if (isNewConvo) {
|
|
return { createdAt: new Date().toISOString(), conversation: undefined };
|
|
}
|
|
|
|
try {
|
|
const conversation = await getConvo(userId, conversationId);
|
|
return {
|
|
conversation,
|
|
createdAt: toValidISOString(conversation?.createdAt) ?? new Date().toISOString(),
|
|
};
|
|
} catch (error) {
|
|
logger.warn('[AgentController] Failed to resolve conversation timestamp anchor', {
|
|
conversationId,
|
|
error: error?.message ?? error,
|
|
});
|
|
return { createdAt: new Date().toISOString(), conversation: undefined };
|
|
}
|
|
}
|
|
|
|
async function attachConversationCreatedAt(req, { userId, conversationId, isNewConvo }) {
|
|
req.body.conversationId = conversationId;
|
|
const resolved = await resolveConversationCreatedAt({
|
|
userId,
|
|
conversationId,
|
|
isNewConvo,
|
|
});
|
|
req.conversationCreatedAt = resolved.createdAt;
|
|
if (!isNewConvo && 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 {
|
|
// Scheduled fires never incremented the concurrent-request counter (they are
|
|
// governed by the scheduler's own caps), so they must not decrement it. Use
|
|
// the decision captured at request start, not a re-check of the expiring token.
|
|
if (!req._isScheduledFire) {
|
|
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 {
|
|
text,
|
|
isRegenerate,
|
|
endpointOption,
|
|
conversationId: reqConversationId,
|
|
isContinued = false,
|
|
editedContent = null,
|
|
parentMessageId = null,
|
|
overrideParentMessageId = null,
|
|
responseMessageId: editedResponseMessageId = null,
|
|
scheduleId: bodyScheduleId = null,
|
|
scheduledFor: bodyScheduledFor = null,
|
|
} = req.body;
|
|
|
|
// Only honor schedule bookkeeping fields on a server-minted scheduled fire
|
|
// (scope-claim verified). A normal chat could otherwise pass an arbitrary
|
|
// scheduleId and corrupt another schedule's lastRun/counters via the hook.
|
|
// Prefer the decision the chat route captured right after auth: the short-lived
|
|
// fire token can expire during the slower chat middleware, so re-verifying here
|
|
// would wrongly demote a valid fire to an ordinary chat and orphan its run.
|
|
const isScheduledFire =
|
|
typeof req._isScheduledFire === 'boolean' ? req._isScheduledFire : isScheduleFireRequest(req);
|
|
// Capture the decision on req: the short-lived fire token expires mid-run, so
|
|
// cleanup must not re-verify it (an expired token would read as non-scheduled
|
|
// and wrongly decrement the concurrent-request counter it never incremented).
|
|
req._isScheduledFire = isScheduledFire;
|
|
const scheduleId = isScheduledFire ? bodyScheduleId : null;
|
|
const scheduledFor = isScheduledFire ? bodyScheduledFor : null;
|
|
|
|
const userId = req.user.id;
|
|
|
|
if (
|
|
await isUnpersistedPreliminaryParent({
|
|
userId,
|
|
conversationId: reqConversationId,
|
|
parentMessageId,
|
|
getMessages,
|
|
})
|
|
) {
|
|
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 isNewConvo = !reqConversationId || reqConversationId === 'new';
|
|
// A verified scheduled fire may supply its own new-conversation id so its run row
|
|
// can record the id BEFORE the request runs (making the occurrence's job findable
|
|
// by reconciliation even if the post-accept detail write fails). Gated on the fire
|
|
// flag so an ordinary chat can't pin an arbitrary conversation id.
|
|
const scheduledNewConversationId =
|
|
req._isScheduledFire && typeof req.body?.newConversationId === 'string'
|
|
? req.body.newConversationId
|
|
: null;
|
|
const conversationId = isNewConvo
|
|
? (scheduledNewConversationId ?? crypto.randomUUID())
|
|
: reqConversationId;
|
|
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');
|
|
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');
|
|
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,
|
|
});
|
|
return res.json({
|
|
streamId: existingStreamId,
|
|
conversationId: claim.existing.conversationId,
|
|
status: 'resumed',
|
|
});
|
|
}
|
|
}
|
|
|
|
// Scheduled fires bypass the interactive concurrent-request limiter (like the
|
|
// message-rate limiters): a user with active chats or several schedules due at
|
|
// once would otherwise get 429s recorded as schedule errors, walking valid
|
|
// schedules toward auto-disable. Scheduler caps govern their concurrency.
|
|
if (!isScheduledFire) {
|
|
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);
|
|
return res.status(429).json(violationInfo);
|
|
}
|
|
}
|
|
|
|
let client = null;
|
|
|
|
try {
|
|
logger.debug(`[ResumableAgentController] Creating job`, {
|
|
streamId,
|
|
conversationId,
|
|
reqConversationId,
|
|
userId,
|
|
});
|
|
|
|
const job = await GenerationJobManager.createJob(streamId, userId, conversationId);
|
|
const jobCreatedAt = job.createdAt; // Capture creation time to detect job replacement
|
|
|
|
// Re-fence a scheduled fire against a delete/quiesce that landed in the claim ->
|
|
// POST window (the reservation row existed but this job did not yet, so the
|
|
// deletion's identity-guarded abort could not have seen it). Attach the schedule
|
|
// identity to the job FIRST — so a concurrent quiesce can now identity-match and
|
|
// abort THIS job — then re-check liveness: if a delete landed in the tiny
|
|
// createJob -> metadata window (before its abort could match), terminalize the
|
|
// run and tear the job down BEFORE any messages are persisted. Scoped to
|
|
// scheduled fires (scheduleId is null for interactive chat), so this never
|
|
// affects it.
|
|
if (scheduleId) {
|
|
await GenerationJobManager.updateMetadata(streamId, { scheduleId, scheduledFor });
|
|
if (!(await isScheduleLive(scheduleId))) {
|
|
logger.info(
|
|
`[AgentController] Scheduled fire aborted before start; schedule ${scheduleId} no longer active`,
|
|
);
|
|
const outcomeRecorded = await recordScheduleOutcome({
|
|
scheduleId,
|
|
scheduledFor,
|
|
status: 'interrupted',
|
|
conversationId: streamId,
|
|
});
|
|
// Terminalize as ABORTED, not complete: if the outcome write failed a
|
|
// preserved `complete` job would be mapped to `success` by the schedules
|
|
// reconciler, mislabeling this pre-start abort as a successful run. abortJob
|
|
// stores it as `aborted` (reconcile -> interrupted), and preserves it for
|
|
// reconcile only when the outcome write failed so the evidence survives.
|
|
await GenerationJobManager.abortJob(streamId, {
|
|
preserveForReconcile: !outcomeRecorded,
|
|
}).catch(() => undefined);
|
|
return res.json({ streamId, conversationId, status: 'aborted' });
|
|
}
|
|
}
|
|
|
|
req._resumableStreamId = streamId;
|
|
getMCPRequestContext(req, undefined, { cleanupOnResponse: false });
|
|
|
|
// The schedule identifiers were already persisted into the job metadata above
|
|
// (before the liveness re-check); the bulk updateMetadata below re-affirms them.
|
|
|
|
// 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, { userId, conversationId, isNewConvo });
|
|
|
|
const endpointIconURL = getEndpointIconURL(req, endpointOption);
|
|
const responseModel = getAgentResponseModel(req, endpointOption);
|
|
const preliminaryUserMessage = getPreliminaryUserMessage(req.body, conversationId);
|
|
const preliminaryResponseMessageId = getPreliminaryResponseMessageId(req.body);
|
|
await GenerationJobManager.updateMetadata(streamId, {
|
|
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,
|
|
// Persist schedule bookkeeping so a HITL resume can record the outcome
|
|
// (the resume request carries no scheduleId). Only set for verified fires.
|
|
...(scheduleId ? { scheduleId, scheduledFor } : {}),
|
|
responseMessageId: preliminaryResponseMessageId,
|
|
userMessage: preliminaryUserMessage,
|
|
});
|
|
|
|
// 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<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,
|
|
});
|
|
|
|
if (job.abortController.signal.aborted) {
|
|
GenerationJobManager.completeJob(streamId, 'Request aborted during initialization');
|
|
await finishResumableRequest(req, userId);
|
|
return;
|
|
}
|
|
|
|
client = result.client;
|
|
// 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) {
|
|
GenerationJobManager.updateMetadata(streamId, { sender: client.sender });
|
|
}
|
|
|
|
// Store reference to client's contentParts - graph will be set when run is created
|
|
if (client?.contentParts) {
|
|
GenerationJobManager.setContentParts(streamId, client.contentParts);
|
|
}
|
|
|
|
let userMessage;
|
|
|
|
const getReqData = (data = {}) => {
|
|
if (data.userMessage) {
|
|
userMessage = data.userMessage;
|
|
}
|
|
// conversationId is pre-generated, no need to update from callback
|
|
};
|
|
|
|
// Start background generation - readyPromise resolves immediately now
|
|
// (sync mechanism handles late subscribers)
|
|
const startGeneration = async () => {
|
|
try {
|
|
// Short timeout as safety net - promise should already be resolved
|
|
await Promise.race([job.readyPromise, new Promise((resolve) => setTimeout(resolve, 100))]);
|
|
} catch (waitError) {
|
|
logger.warn(
|
|
`[ResumableAgentController] Error waiting for subscriber: ${waitError.message}`,
|
|
);
|
|
}
|
|
|
|
/** 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 immediateTitlePromise = null;
|
|
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,
|
|
},
|
|
});
|
|
})().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,
|
|
}),
|
|
},
|
|
});
|
|
|
|
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,
|
|
});
|
|
};
|
|
|
|
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 (!isScheduledFire && !client?.pendingRequestReleased) {
|
|
await decrementPendingRequest(userId);
|
|
}
|
|
if (client) {
|
|
disposeClient(client);
|
|
}
|
|
// Record a scheduled fire's pause as `requires_action` NOW (not on the
|
|
// next reconcile sweep): overlap/capacity checks key on `started`, so a
|
|
// run left `started` while paused would wrongly block a run-now or the
|
|
// next occurrence. A failed write just falls back to reconcile marking it.
|
|
if (scheduleId) {
|
|
await recordScheduleOutcome({
|
|
scheduleId,
|
|
scheduledFor,
|
|
status: 'requires_action',
|
|
conversationId: streamId,
|
|
});
|
|
}
|
|
logger.debug(
|
|
`[ResumableAgentController] Turn paused for approval; awaiting resume: ${streamId}`,
|
|
);
|
|
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);
|
|
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,
|
|
});
|
|
}
|
|
} 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);
|
|
// Record the schedule success BEFORE completeJob: the default job
|
|
// manager deletes completed jobs immediately, so if this write ran
|
|
// after and the process died in between, reconciliation would see no
|
|
// job for the still-`started` run and mislabel a success as interrupted.
|
|
// If the write itself fails (transient Mongo outage across its retries),
|
|
// preserve the completed job so the reconciler can finalize it as
|
|
// success from the retained job status instead of deleting the evidence.
|
|
let scheduleOutcomeRecorded = true;
|
|
if (scheduleId) {
|
|
scheduleOutcomeRecorded = await recordScheduleOutcome({
|
|
scheduleId,
|
|
scheduledFor,
|
|
status: 'success',
|
|
conversationId: conversation?.conversationId,
|
|
});
|
|
}
|
|
GenerationJobManager.completeJob(streamId, undefined, {
|
|
preserveForReconcile: !scheduleOutcomeRecorded,
|
|
});
|
|
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);
|
|
// Record the abort BEFORE completeJob so the run doesn't linger as
|
|
// `started` (blocking run-now/overlap until the 30-minute orphan cutoff).
|
|
let abortOutcomeRecorded = true;
|
|
if (scheduleId) {
|
|
abortOutcomeRecorded = await recordScheduleOutcome({
|
|
scheduleId,
|
|
scheduledFor,
|
|
status: 'interrupted',
|
|
conversationId: conversation?.conversationId,
|
|
});
|
|
}
|
|
// Only finalize/clean the job when the outcome is recorded. When it isn't
|
|
// (Mongo down), leave any reconcile evidence the abort route preserved
|
|
// (an `aborted` job) intact — overwriting it here as an `error` job would
|
|
// make reconcile count a user stop toward failure auto-disable.
|
|
if (abortOutcomeRecorded) {
|
|
GenerationJobManager.completeJob(streamId, 'Request aborted');
|
|
}
|
|
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');
|
|
|
|
let errorScheduleOutcomeRecorded = true;
|
|
if (scheduleId) {
|
|
errorScheduleOutcomeRecorded = await recordScheduleOutcome({
|
|
scheduleId,
|
|
scheduledFor,
|
|
status: wasAborted ? 'interrupted' : 'error',
|
|
conversationId: streamId,
|
|
error: wasAborted ? undefined : error.message,
|
|
});
|
|
}
|
|
|
|
if (wasAborted) {
|
|
logger.debug(`[ResumableAgentController] Generation aborted for ${streamId}`);
|
|
// 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 },
|
|
);
|
|
}
|
|
} catch (drainErr) {
|
|
logger.warn(
|
|
`[ResumableAgentController] Failed to close steer queue on error`,
|
|
drainErr,
|
|
);
|
|
}
|
|
await GenerationJobManager.emitError(streamId, error.message || 'Generation failed');
|
|
GenerationJobManager.completeJob(streamId, error.message, {
|
|
preserveForReconcile: Boolean(scheduleId) && !errorScheduleOutcomeRecorded,
|
|
});
|
|
}
|
|
|
|
await finishResumableRequest(req, userId);
|
|
|
|
// Defer disposal until any immediate title settles (it holds the run/req).
|
|
if (immediateTitlePromise) {
|
|
immediateTitlePromise.finally(() => {
|
|
if (client) {
|
|
disposeClient(client);
|
|
}
|
|
});
|
|
} else if (client) {
|
|
disposeClient(client);
|
|
}
|
|
|
|
// 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}`,
|
|
);
|
|
GenerationJobManager.completeJob(streamId, err.message);
|
|
await finishResumableRequest(req, userId);
|
|
});
|
|
} catch (error) {
|
|
logger.error('[ResumableAgentController] Initialization error:', error);
|
|
if (!res.headersSent) {
|
|
res.status(500).json({ error: error.message || 'Failed to start generation' });
|
|
} else {
|
|
// JSON already sent, emit error to stream so client can receive it
|
|
await GenerationJobManager.emitError(streamId, error.message || 'Failed to start generation');
|
|
}
|
|
// A scheduled fire whose run row is already `started` (inserted before the
|
|
// early 200) must be terminalized here — an init failure after acceptance
|
|
// would otherwise leave it running until orphan reconciliation. Before
|
|
// completeJob deletes the job (same ordering rationale as the success path).
|
|
// If the outcome write fails (Mongo down across its retries), preserve the
|
|
// completed job so reconcile records the failure instead of losing the count.
|
|
let initScheduleOutcomeRecorded = true;
|
|
if (scheduleId) {
|
|
initScheduleOutcomeRecorded = await recordScheduleOutcome({
|
|
scheduleId,
|
|
scheduledFor,
|
|
status: 'error',
|
|
error: error.message,
|
|
conversationId: streamId,
|
|
});
|
|
}
|
|
// 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 — and completeJob() is not guarded by the original createdAt, so it would
|
|
// abort/error that replacement. 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.)
|
|
await GenerationJobManager.completeJob(streamId, error.message, {
|
|
preserveForReconcile: Boolean(scheduleId) && !initScheduleOutcomeRecorded,
|
|
}).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 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] });
|
|
} else if (key === 'sender') {
|
|
GenerationJobManager.updateMetadata(streamId, { sender: data[key] });
|
|
}
|
|
// 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 (streamId) {
|
|
logger.debug('[AgentController] Completing job in GenerationJobManager');
|
|
await GenerationJobManager.completeJob(streamId);
|
|
}
|
|
|
|
// 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<string, Record<string, string>> }} */
|
|
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);
|
|
|
|
// Store endpoint metadata for abort handling
|
|
GenerationJobManager.updateMetadata(streamId, {
|
|
endpoint: endpointOption.endpoint,
|
|
iconURL: getEndpointIconURL(req, endpointOption),
|
|
model: getAgentResponseModel(req, endpointOption),
|
|
sender: client?.sender,
|
|
});
|
|
|
|
// Store content parts reference for abort
|
|
if (client?.contentParts) {
|
|
GenerationJobManager.setContentParts(streamId, client.contentParts);
|
|
}
|
|
|
|
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,
|
|
},
|
|
});
|
|
};
|
|
|
|
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;
|