mirror of
https://github.com/danny-avila/LibreChat.git
synced 2026-09-23 08:35:32 +00:00
* feat: ask_user_question tool — agent-initiated questions with durable pause/resume The HITL runtime merged in #13942/#14024/#14025/#14123 already ships the full ask_user_question lifecycle (payload-agnostic handleRunInterrupt, resume validation via mapAskUserAnswer, reconnect rehydration, and the client question card) — but nothing ever raised the interrupt. This adds the producer: - packages/api/agents/hitl/askUserQuestionTool.ts: LLM-callable tool whose func calls the SDK askUserQuestion() helper (LangGraph interrupt() from the tool body); zod schema with length caps mirroring AskUserQuestionRequest, plus a JSON-schema twin for the schema-only registry - Registration: agentToolDefinitions, manifest.json (Tools dialog, admin filteredTools/includedTools kill switch), basicToolInstances, handleTools constructor branch - run.ts gating: checkpointer now attaches for hitlCapable runs whose agents carry the ask tool even with the tool-approval policy disabled (the interrupt needs only durability, not humanInTheLoop/hooks); the tool is stripped fail-closed from non-HITL callers (OpenAI-compat/Responses) and subagent child configs; excluded from eager event execution (interrupts must be raised inside the Pregel task frame) - resume.js: 16k length cap on the answer wire field - e2e (real Run + FakeChatModel + LazyMongoSaver + supertest resume): tool-body interrupt pauses durably with NO approval policy, answer round-trips as the ToolMessage content, tool body re-runs once on resume, sequential questions re-pause * fix: adversarial-review findings — in-graph execution, orphan prunes, endpoint scoping, real kill switch Pre-PR multi-agent review confirmed 5 defects in the initial commit; all fixed: 1. CRITICAL — the tool never paused on the real agents endpoint: production loads tools definitions-only, flipping the SDK ToolNode to event-driven dispatch, and the host ON_TOOL_EXECUTE handler runs outside the Pregel task frame (under runOutsideTracing), where interrupt() throws and becomes an error ToolMessage. Reworked: the ask tool never rides toolDefinitions/ toolRegistry — on HITL-capable top-level agents a real instance is supplied via AgentInputs.graphTools (agents#289, requires @librechat/agents > 3.2.57), the SDK's in-graph direct-tool seam; new production-shape e2e pins the event-driven mode end to end. 2. CRITICAL — ask-only runs left orphaned interrupted checkpoints (silent context duplication on every later turn): both orphan prunes were gated on toolApproval.enabled. The pre-turn prune now also fires for ask-capable agents (exported agentRequestsAskUserQuestion), and the abort-route prune fires when the aborted job carries a pendingAction. 3. MAJOR — self-spawned subagents bypassed the strip (self config resolves from the parent's _sourceInputs): fixed SDK-side (buildChildInputs clears graphTools) and the tool is now never present on child surfaces host-side. 4. MINOR — the manifest entry leaked into the Assistants tools dialog and the legacy plugins endpoint, where tools execute with no run to pause: new agentsOnly manifest flag, scoped out of both listings. 5. MINOR — filteredTools/includedTools only hid the tool from the dialog: now enforced at run build (strip + no checkpointer), making the admin filter a real kill switch for already-saved agents. * chore: update @librechat/agents dependency to version 3.2.58 in package-lock.json and package.json files * fix: reject agents-only tools at assistant create/update (Codex round 1) The tools-dialog scoping keeps ask_user_question out of the assistants LISTING, but the v1/v2 create/update handlers resolve arbitrary posted tool strings from the shared getCachedTools map — a REST client or stale saved payload could still attach it, and the assistants runtime executes tools with no run to pause, so every call would error. New isAgentsOnlyTool(tool) (manifest-driven, handles string and function-object shapes) drops such tools with a warn at all four resolution sites (v1+v2, create+update). * fix: offset resumed-run content indices past the pre-pause seed A resumed run rebuilds the graph from the checkpoint, and the fresh graph numbers content indices from its own empty contentData — starting at 0. The resume path seeds the (also fresh) content aggregator with the pre-pause parts at exactly those indices, so the resumed model turn collided with the seed: type-matching parts silently MERGED (post-resume text appended into a pre-pause text block), and type-mismatching parts (a reasoning/think part at index 0 — any Anthropic reasoning agent) dropped EVERY delta with 'Content type mismatch', losing the entire post-resume output from the live stream and the saved message. Latent since #13942 — tool-approval resumes corrupt content the same way (probe-verified); it surfaced now because ask_user_question makes pausing a first-class flow and reasoning models make the loss total. - createContentIndexOffsetHandlers(handlers, offset): wraps ON_RUN_STEP (the single point where a content index enters the pipeline — deltas resolve through the aggregator's stepMap) and ON_AGENT_UPDATE's inline index; every other handler passes through by reference. Probe-validated: resumed output now lands as a new part after the paused tool call. - resumeCompletion wires it with offset = seedContent.length. - logToolError: a GraphInterrupt unwinding out of a tool body is the HITL pause working as designed — no longer logged as a Tool Error. * fix: unblock live streaming of the resumed segment after an answer With resume indices now ABSOLUTE (server continues after the pre-pause parts), the synthetic ask-user-question card was squatting on exactly the index the resumed segment streams into: applyAskUserQuestion appends the card at the end of the message content, so on the answering device every incoming part at that index was blocked and nothing rendered between the answer submission and the finalize replacing the message. removeAskUserQuestionPart(message, actionId) strips the pause-scoped card on successful answer submission (useResumeSubmit onSuccess) — the durable record of the Q&A is the ask_user_question tool call itself. Pure helper + specs; same-reference no-op when nothing matches. * fix: displace the synthetic question card in the streaming content writer The store-level strip on answer submit wasn't enough: the SSE step handler keeps its own in-flight copy of the streaming message, so on the answering device the synthetic ask-user-question card still occupied the ABSOLUTE index the resumed segment streams into — every delta warned 'Content type mismatch' (existing ask_user_question vs incoming text) and nothing rendered between the pending_action and finalize. Displace the card inside updateContent when any real part claims its slot — the same displacement pattern as the OAuth prompt part directly above it. Covers the streaming handler's own copy, reconnecting tabs, and other devices; once real content streams, the pause is over by definition. Spec drives a runStep + text delta into the card's index and pins: no mismatch warn, card gone, text rendered. * feat: dedicated UI + durable data for completed ask_user_question calls The completed ask call rendered as a generic tool card labeled 'Cancelled' with raw (and empty) JSON args. Two layers fixed: Data: the saved tool_call part had args:'' and no output — streamed arg chunks carry no tool name so the aggregator drops them (normal tools recover via the completion event, which never fires for a tool that interrupts mid-execution and resumes on a rebuilt run with no step id). The resume controller now stamps the paused ask part with the pendingAction's authoritative question as args and the user's answer as output (attachAskUserQuestionAnswer — pure, targets the newest unanswered ask part, so sequential questions each keep their own answer). UI: Part.tsx routes ask_user_question tool calls to AskUserQuestionCall — a compact Q&A record ('Asked a question' header, question, description, 'You answered: <label>' preferring the picked option's label, or 'No answer was given' for an abandoned pause) instead of the generic card. New i18n keys; parseAskUserQuestionArgs degrades to null on malformed model args. * fix: single question UI per pause + immediate answer display Two live-turn issues with the new durable Q&A card: 1. Duplicate question on ask: during a live pause the message carries BOTH the ask tool_call part (now rendered by AskUserQuestionCall, showing a misleading 'No answer was given' while paused) and the synthetic interactive card. The durable card now defers while the turn is live and unanswered (isSubmitting) — the interactive card owns the question UI until it's answered; an abandoned pause still shows its no-answer state once the turn settles. 2. 'No answer was given' after answering: the server stamps the answer onto the part at resume seed, but the client only received that at finalize. No stream emission needed — the client knows the answer it just submitted: resolveAskUserQuestionPart (replacing the plain strip on submit success) removes the synthetic card AND stamps output/progress onto the newest unanswered ask tool_call, seeding args from the synthetic part's question when the streamed args were lost — mirroring the server-side attachAskUserQuestionAnswer, so the Q&A record shows the answer the moment the user submits. * fix: keep the Q&A record visible while the resumed segment streams The optimistic output stamp lives in the message store, but the SSE step handler evolves its own cached copy of the streaming message (created at turn start) — the first resumed event overwrites the store with that copy, wiping the stamp, so the Q&A card blinked out during streaming and only returned at finalize. Render-layer fallback instead of fighting the handler's copy: submitted answers are recorded by ask tool_call id when resolveAskUserQuestionPart stamps the part, and AskUserQuestionCall reads the recorded answer whenever the part's own output is missing — the record survives any message-copy churn until finalize delivers the server-stamped part. * feat: present Ask User as a native builtin in the tools dialog It ships with the app and pauses the run like a first-class feature, so it belongs with the builtins (Run Code, Web Search, Memory, ...) rather than in the third-party plugin list — while its mechanics stay exactly a plugin's: - BuiltinId += 'ask_user_question' (documented exception: a native TOOL, not a capability; selection reads agent.tools, the toggle emits tool-add/remove patches instead of a capability field) - buildCatalog surfaces it as a builtin gated on the same signals as before (tools capability on + the server lists the plugin, i.e. not admin-filtered) and skips it in the plugin loop so it never double-lists - On-theme icon: lucide MessageCircleQuestion in a teal chip via the builtin icon map, matching the other native entries; the bespoke purple SVG and the manifest icon field are gone - i18n'd name/description keys like the other builtins * feat: composer popover for answering questions (mentions-style) Answering moves to the composer, matching the existing mentions/prompts popover pattern: while an ask_user_question pause is live, a popover anchors above the textarea with the question as its header, numbered option rows (hover/click, or ↑/↓ + Enter from the empty composer), and an × to dismiss. The main textarea doubles as the free-form answer — its placeholder flips to 'Something else...' and form submit routes the text to the paused run as the answer instead of starting a new turn. Dismissing (× or Escape) restores normal sends; the inline transcript surfaces stay as before (interactive card while paused, durable Q&A record after) so the question remains visible in history. - findLiveAskUserQuestion (pure, spec'd): newest unanswered synthetic part across the conversation IS the popover signal — applied on on_pending_action, stripped on answer submit, so visibility tracks the pause lifecycle with no extra state - useLiveAskUserQuestion hook shared by the popover and ChatForm; dismissals in a recoil atom so both react - popover only mounts on the primary composer (index 0), mirroring QuoteButton * feat: number-key selection + return glyph in the question popover Pressing 1-9 in the empty composer picks the matching option directly, mirroring the numbered row chips; the highlighted row shows a return-key glyph as the Enter affordance. Same empty-composer guard as the arrow keys — typing a free-form answer is never intercepted. * refactor: first-class composer answer mode (useAskAnswerMode) Replaces the bolted-on integration (inline onSubmit interception + raw capture-phase keydown listeners on the textarea ref) with a single hook that owns the whole answer mode: live-question derivation, dismissal + highlighted option (shared recoil state), option selection, free-form submit routing (submitText returns whether it consumed the submission), and keyboard handling (handleKeyDown returns whether it consumed the key, composed ahead of the textarea's normal handler — no more addEventListener). The popover is now pure rendering off the hook; ChatForm wires placeholder, onKeyDown, and onSubmit through the same instance. Deliberately scoped to the composer rather than useSubmitMessage: starters/prompt-commands keep new-turn semantics (and the existing job-replacement behavior while paused). * fix: Codex round 2 — inline answer input, approval exemption, pause-time args F1 (composer submit unreachable while paused — isSubmitting keeps Stop shown and useTextarea eats Enter): redesigned around it, borrowing Claude Code's AskUserQuestion semantics. The popover now owns free-form input via an inline 'Other' row (numbered last, 'Something else…'), with select-then-confirm rows (click/arrows/digits highlight; Submit ↵, Enter, or double-click fires; Skip dismisses). The composer returns to being a plain composer — no placeholder swap, no submit interception; Stop keeps meaning stop. F2: ask_user_question is exempt from the tool-approval prompt unless the admin explicitly lists it (allow/ask/deny all win) — approving the right to ask a question was a pure double pause; the tool is side-effect-free. F3: the question is stamped onto the paused ask tool_call's args at PAUSE time (attachAskUserQuestionArgs in handleRunInterrupt), so abandoned/expired/ stopped turns persist with the question intact and the record card can render it — previously only the answer-resume path stamped args. * fix: fold model-supplied 'Other' options into the inline free-form row The model can generate its own catch-all option ('Other (type your own)', value 'other'), duplicating the popover's built-in free-form row — two other-ish rows, one pickable as a literal answer. Two layers: - Tool description now tells the model NOT to include catch-all options (the answer UI always offers free-form input on its own) - splitOtherOption (pure, spec'd) folds a catch-all option that arrives anyway out of the choice rows and uses its label as the inline input's placeholder — conservative match (value 'other', or a label reading as a free-form invitation), no false positives on real choices * fix: single question surface + clean free-form-only popover Two live-pause confusions: (1) the inline transcript card and the composer popover both rendered — the card now defers while the popover is up for its action, returning as the fallback surface when the user dismisses the popover (and in contexts without a ChatContext, where the popover can't exist); (2) an options-less question showed a pointless numbered '1 Something else…' row — free-form-only questions now render the inline input alone, with the 'Type your answer…' placeholder (a folded model 'Other' label still wins). * feat: the composer is the free-form answer box (like the main chat input) While a question pause is live, the main chat textarea composes the free-form answer — placeholder swaps to 'Something else…' (or a folded model 'Other' label), Enter with text submits the answer through answer-mode key handling (composed BEFORE useTextarea's submitting-lock, so the lock can't swallow it), and the Stop button swaps to Send (enabled despite isSubmitting) per the select-then-confirm design. The popover slims to the question header, numbered option rows, and Skip/Submit — its inline input is gone since the composer owns free-form now. Dismissing the popover restores normal composer semantics (Stop button, normal sends). * fix: Codex round 3 + real Skip semantics - Skip now ANSWERS instead of hiding UI (danny): it resumes the run with a decline notice ('The user chose not to answer this question.') so the model moves on — a client-side dismiss left the run paused until expiry, a hung turn. × / Escape remain pure dismiss (switch to the inline card surface). - P1 (resumed approval tool indices): resumed tool_calls steps whose tool_call id matches a seeded UNRESOLVED part now rebind to that seeded slot instead of offsetting — the original part resolves in place (output attaches) and no duplicate appears; message steps keep the offset, so the text-loss fix stands. createContentIndexOffsetHandlers now takes the seed array; resolved seeded calls are not rebind targets. - P2 (stale selection across questions): selection state resets when the live actionId changes; the vestigial inline-Other state ('other' selection + text atom) is gone — the composer owns free-form. - P2 (Redis abort path loses the args stamp): the abort route re-stamps the question onto the ask tool_call in the reconstructed abort content, so a Stop-abandoned question persists with its question intact. - P2 (malformed args crash): parseAskUserQuestionArgs normalizes untrusted shapes (options: {} / non-string entries) instead of throwing in render. * feat: free-form hint in the question popover footer Left-aligned in the footer row (opposite Skip/Submit): 'Or type your answer below' — points open-ended answering at the composer, whose placeholder already reads 'Something else…'. * feat: preserve composer drafts across the answer-mode swap The answer phase gets its own draft key (ask-answer:<actionId>), passed as a draftId override into useAutoSave — the key change itself drives the existing save/restore machinery, so the conversation draft (or mid-run PENDING draft) is stashed when a question pause takes the composer and restored once the user answers, skips, or dismisses. Ask keys are exempt from the PENDING migration branch, which would otherwise move-and-delete the stashed draft. A half-typed answer survives reload/navigation while its question stays live. Answer submission (option pick, free-form, skip) resets the composer via a new non-throwing useOptionalChatFormContext, so the swap-back restores into an empty box even outside ChatView-less render contexts (Share/search). * fix: rebind resumed steps for ALL seeded tool call ids The resume controller pre-stamps the user's answer onto the seeded ask_user_question part, so the unresolved-only rebind predicate treated it as settled and shifted the tool's re-run step to a fresh offset slot, leaving a duplicate ask record in streamed/saved content. Tool call ids are provider-minted per call: a resumed step bearing a seeded id can only be the interrupted batch re-executing, so rebinding every seeded id is always correct. * feat: popover UX round 4 — clickable hint, collapse, click-submit, multiSelect - Footer hint is a button that focuses the composer; reads 'Type your answer below' (no 'Or') when the question has no options. - Collapse (chevron) hides the popover WITHOUT closing the pause: answer mode stays live (placeholder, Enter routing, draft key), the chat card renders the question with a ChevronUp affordance to re-expand. x remains dismiss. - Single-select options submit on a single click; the Submit button renders only for multi-select. - multiSelect end-to-end: tool zod schema + JSON definition twin, wire type, client parse, popover check-chips, card toggles, record-card label mapping; answer = option values joined ', '; composer Enter and the multi Submit button both fold free-form text in with the checked values. - Hardening from adversarial review: in-flight status guard on every submit path (no duplicate resumes on double-click), popover locks while submitting, collapsed mode disarms invisible digit/arrow steering, the card shares the hook's checked state while the pause is live, the card folds catch-all 'Other' options, record mapping is all-or-nothing to avoid phantom labels, composer resets only when its text was consumed or the draft machinery will restore the stash. * feat: ask_user_question in model specs and ephemeral agents A librechat.yaml modelSpec can now equip the tool the same way it equips webSearch/executeCode/fileSearch/memory: modelSpecs: list: - name: my-spec askUserQuestion: true loadEphemeralAgent pushes the tool name when the spec flag (or the ephemeralAgent request flag, wired for parity) is set; everything downstream is the existing persisted-agent machinery — createRun's hitlCapable gating, graphTools injection, checkpointer attach, subagent strip, and the admin filteredTools/includedTools kill switch all apply unchanged. * feat: tense-aware Q&A record label (Asking / Asked) Shorten the record card header per feedback: 'Asking' while the question is still unanswered (abandoned/awaiting), 'Asked' once answered — replacing the single 'Asked a question' label. * fix: Codex round 4 — added-agent ask parity + preserve answer on failed resume F1 (added.ts): mirror loadEphemeralAgent's ask_user_question branch in the added-agent loader so a model spec's askUserQuestion flag (or the ephemeral request flag) equips added top-level agents too, matching execute_code / web_search / memory. Two load.spec cases added. F3 (composer): submitAskAnswer now takes an onSuccess callback and useAskAnswerMode defers clearing the selection/composer until the resume is accepted. A failed resume (16k answer-cap 400, expired action, network error) leaves status re-answerable, so wiping the composer up front lost the user's only copy of a free-form answer; now it survives for trim/retry. (F2 — a claimed Tools-capability bypass — was verified NOT reproducible: agentRequestsAskUserQuestion matches only loaded instances/toolDefinitions/ toolRegistry, all capability-filtered; a raw tools string has no .name and never triggers the install. Replied on-thread with the probe evidence.) * fix: Codex round 5 — expired question exits answer mode so its message shows An expired question (e.g. resume returns the stale-action 409) previously left the popover open with locked controls and no explanation, because the chat card — which carries the only 'this action expired' message — was suppressed by the popover-open guard. Treat 'expired' as no longer active: the popover closes, the composer reverts to normal, and the card becomes the sole surface and renders the expired message. 'error' stays active (retryable). * feat: group ask_user_question calls as their own category A homogeneous group of ask_user_question tool calls now reads 'Asked N questions' (present tense 'Asking N questions' while the turn streams) with a question glyph and no raw-name suffix — mirroring the subagent 'Ran N agents' category treatment, instead of 'Used N tools — ask_user_question'. Mixed groups keep 'Used N tools' but humanize the suffix to 'Question' and show a question icon for the ask entries (TOOL_FRIENDLY_NAME_KEYS + ToolIcon map). A group only forms at count >= 2, so the plural is always grammatical. Three ToolCallGroup.test cases cover homogeneous label/icon/suffix, present tense while streaming, and the mixed-group fallback. * fix: Codex round 6 — composer submit lock + abort stamp before emit F7 (composer status lock): the ask submit status lived on ApprovalContext, a React context mounted only around message content (ContentParts). The PRIMARY answer surface — the composer in ChatForm — renders outside it, so useApprovalContext returned the inert FALLBACK: status was always 'idle', setStatus a no-op. The in-flight double-submit guard (round 4) and the expired-exits-answer-mode fix (round 5) therefore never engaged for the composer. Move ask submit status to a global Recoil atom (useAskSubmitStatus) read/written by the composer, the popover, and the card alike, so a fast double-click/Enter is actually blocked and expired/error surfaces on every surface. Tool-approval status stays on the context (unchanged). F5 (abort stamp before emit): the abort route re-stamped a paused ask_user_question's args AFTER GenerationJobManager.abortJob had already emitted the final SSE from the unstamped content, so a Redis/cross-replica Stop left the live client showing an empty question until reload. abortJob now takes an optional transformAbortContent applied to the persistable content BEFORE the final event is built (and returned), so the live client and the saved message agree. New abort.spec case + updated call assertions. * feat: gate ask_user_question behind its own agent capability Add a first-class AgentCapabilities.ask_user_question (in defaultAgentCapabilities, on by default) so admins can enable/disable questions independently via endpoints.agents.capabilities, exactly like execute_code / web_search — not lumped under the generic tools capability. - ToolService: both filteredTools predicates (definitions-only and instance loaders) gate ask_user_question on checkCapability(ask_user_question) before the generic tools fallthrough. When off, the tool is dropped from toolDefinitions/toolRegistry, so run.ts's agentRequestsAskUserQuestion (which keys on the loaded surface) declines to install it and attach a checkpointer — the capability is enforced end-to-end at the loader, no run.ts change needed. - Tools dialog catalog: surface the ask builtin under its own capability rather than the generic tools one, so the UI matches the backend gate. - Tests: ToolService capability on/off filtering + defaults membership; catalog builtin visibility keyed on the dedicated capability. * style: sort imports in ToolCallGroup.test (CI import-order gate) * fix: Codex round 7 — surface ask-answer errors in the open popover A failed answer submission (16k reject, network error) sets the ask status to 'error', which — unlike 'expired' — deliberately keeps the question active and retryable. But the chat card that renders the error message is suppressed while the popover is open, so a composer/popover answer failed silently. Expose an 'errored' flag from useAskAnswerMode and render a warning line (com_ui_ask_answer_error) in the popover, so the user gets feedback and retry guidance without having to collapse/dismiss. It clears automatically on retry (status flips to 'submitting'). * fix: Codex round 8 — respect IME composition before submitting answers handleComposerKeyDown runs before useTextarea's composition guard, so with a CJK/IME keyboard the Enter that commits an in-progress composition was being intercepted and submitting the partial answer (and the composition buffer can leave value empty mid-compose, mis-triggering digit/arrow steering too). Bail at the top when composing — nativeEvent.isComposing, or key==='Process' / keyCode===229 for Safari's inconsistent reporting — mirroring the existing composer guard so the character commits normally. * chore: update `@librechat/agents` to v3.2.60 * 🔧 chore: Update @opentelemetry/core to version 2.9.0 and clean up package-lock.json * feat: digit shortcuts select options when the popover has focus Previously a number key (1..N) only selected an option from the empty composer (handleComposerKeyDown on the textarea) — if focus moved into the popover (a row/Skip/Submit button clicked or tabbed to), the number keys went dead. Add handlePopoverKeyDown, wired to the popover container's onKeyDown so it catches digits bubbling from the focused control: a digit activates its option exactly like a click (single-select submits, multi toggles). No highlight/Enter dance on this path — the options are buttons whose action is the click, and intercepting Enter would fight the focused button. Gated on active && !locked so it no-ops while a submit is in flight. * chore: update @librechat/agents to version 3.2.61 and @opentelemetry packages to latest versions
725 lines
32 KiB
JavaScript
725 lines
32 KiB
JavaScript
const { logger } = require('@librechat/data-schemas');
|
|
const { Constants, EModelEndpoint } = require('librechat-data-provider');
|
|
const {
|
|
GenerationJobManager,
|
|
isPendingActionStale,
|
|
mapToolApprovalResolutions,
|
|
mapAskUserAnswer,
|
|
attachAskUserQuestionAnswer,
|
|
findUndecidedToolCalls,
|
|
findDisallowedDecisions,
|
|
findIncompleteDecisions,
|
|
computeAgentRequestFingerprint,
|
|
deleteAgentCheckpoint,
|
|
buildAbortedResponseMetadata,
|
|
sanitizeMessageForTransmit,
|
|
filterMalformedContentParts,
|
|
decrementPendingRequest,
|
|
checkAndIncrementPendingRequest,
|
|
} = require('@librechat/api');
|
|
const { disposeClient } = require('~/server/cleanup');
|
|
const {
|
|
getMCPRequestContext,
|
|
cleanupMCPRequestContextForReq,
|
|
} = require('~/server/services/MCPRequestContext');
|
|
const { saveMessage, getConvo, getMessages } = require('~/models');
|
|
|
|
/**
|
|
* Upper bound on an `ask_user_question` answer (characters). Generous for any real
|
|
* reply typed into the question card while still bounding what a crafted POST can
|
|
* inject into the resumed run's ToolMessage.
|
|
*/
|
|
const MAX_ASK_ANSWER_LENGTH = 16_000;
|
|
|
|
/** De-duplicate a merged attachment list by a stable artifact identity. */
|
|
function mergeAttachments(existing, incoming) {
|
|
const seen = new Set();
|
|
const out = [];
|
|
for (const attachment of [...(existing ?? []), ...(incoming ?? [])]) {
|
|
if (!attachment) {
|
|
continue;
|
|
}
|
|
const key =
|
|
attachment.file_id ??
|
|
attachment.filepath ??
|
|
attachment.filename ??
|
|
JSON.stringify(attachment);
|
|
if (seen.has(key)) {
|
|
continue;
|
|
}
|
|
seen.add(key);
|
|
out.push(attachment);
|
|
}
|
|
return out;
|
|
}
|
|
|
|
/**
|
|
* Resolve the current segment's tool artifacts and merge them with any already
|
|
* persisted on the response row. A resumed turn can span multiple pause segments;
|
|
* each rebuilt client has its own `artifactPromises`, and the final finalize would
|
|
* otherwise OVERWRITE the row's attachments with only the last segment's. Reading
|
|
* the persisted row and merging keeps every segment's artifacts on the saved message.
|
|
*/
|
|
async function resolveAccumulatedAttachments({ client, conversationId, responseMessageId }) {
|
|
const promises = Array.isArray(client?.artifactPromises) ? client.artifactPromises : [];
|
|
const resolved = promises.length > 0 ? (await Promise.all(promises)).filter(Boolean) : [];
|
|
let existing = [];
|
|
if (responseMessageId) {
|
|
try {
|
|
const [row] = await getMessages(
|
|
{ conversationId, messageId: responseMessageId },
|
|
'attachments',
|
|
);
|
|
existing = Array.isArray(row?.attachments) ? row.attachments : [];
|
|
} catch (err) {
|
|
logger.warn(
|
|
'[ResumeAgentController] Failed to read prior attachments for merge',
|
|
err?.message ?? err,
|
|
);
|
|
}
|
|
}
|
|
return mergeAttachments(existing, resolved);
|
|
}
|
|
|
|
/** Resolve the segment's content for an unfinished save (mirrors finalize's source). */
|
|
async function resolveSegmentContent(client, streamId) {
|
|
const liveContent = Array.isArray(client?.contentParts) ? client.contentParts : [];
|
|
const rawContent =
|
|
liveContent.length > 0
|
|
? liveContent
|
|
: ((await GenerationJobManager.getResumeState(streamId))?.aggregatedContent ?? []);
|
|
return filterMalformedContentParts(rawContent);
|
|
}
|
|
|
|
/**
|
|
* A resumed segment that streamed content / produced artifacts and then paused AGAIN
|
|
* must persist that progress before returning. The next resume rebuilds a fresh client
|
|
* (empty `contentParts`/`artifactPromises`), so without this an approval that later
|
|
* expires or is reaped would leave only the EARLIER pause's content on the saved row —
|
|
* the user loses everything streamed during this segment. Saved as a partial (`$set`,
|
|
* still `unfinished`) so a subsequent successful resume overwrites it on finalize.
|
|
*/
|
|
async function persistRePauseProgress({ req, client, job, streamId, conversationId }) {
|
|
const userId = req.user.id;
|
|
const meta = job.metadata ?? {};
|
|
const responseMessageId = meta.responseMessageId ?? client.responseMessageId;
|
|
if (!responseMessageId) {
|
|
return;
|
|
}
|
|
const content = await resolveSegmentContent(client, streamId);
|
|
const attachments = await resolveAccumulatedAttachments({
|
|
client,
|
|
conversationId,
|
|
responseMessageId,
|
|
});
|
|
if (content.length === 0 && attachments.length === 0) {
|
|
return;
|
|
}
|
|
try {
|
|
await saveMessage(
|
|
{
|
|
userId,
|
|
isTemporary: meta.isTemporary ?? req.body?.isTemporary,
|
|
interfaceConfig: req?.config?.interfaceConfig,
|
|
},
|
|
{
|
|
messageId: responseMessageId,
|
|
conversationId,
|
|
...(content.length > 0 && { content }),
|
|
...(attachments.length > 0 && { attachments }),
|
|
unfinished: true,
|
|
user: userId,
|
|
},
|
|
{ context: 'api/server/controllers/agents/resume.js - re-pause progress persist' },
|
|
);
|
|
} catch (err) {
|
|
logger.error('[ResumeAgentController] Failed to persist re-pause progress', err);
|
|
}
|
|
}
|
|
|
|
/** Untenanted jobs (pre-multi-tenancy) remain accessible if the userId check passes. */
|
|
function hasTenantMismatch(job, user) {
|
|
return job.metadata?.tenantId != null && job.metadata.tenantId !== user.tenantId;
|
|
}
|
|
|
|
/**
|
|
* Build the SDK resume value from the wire decision payload, validating against the
|
|
* pending action. Returns `{ resumeValue }` on success or `{ error }` with an HTTP
|
|
* status for the route to surface.
|
|
*/
|
|
function resolveResumeValue(pendingAction, body) {
|
|
const payload = pendingAction.payload;
|
|
if (payload?.type === 'tool_approval') {
|
|
const resolutions = Array.isArray(body.decisions) ? body.decisions : [];
|
|
const undecided = findUndecidedToolCalls(payload, resolutions);
|
|
if (undecided.length > 0) {
|
|
return { status: 400, error: 'Every paused tool call must be decided', undecided };
|
|
}
|
|
// Enforce the policy's per-tool allowed_decisions — a crafted POST must not
|
|
// approve a tool the policy restricted to (e.g.) reject/respond.
|
|
const disallowed = findDisallowedDecisions(payload, resolutions);
|
|
if (disallowed.length > 0) {
|
|
return { status: 403, error: 'Decision not permitted for one or more tools', disallowed };
|
|
}
|
|
// `edit`/`respond` must carry their payload — otherwise toSdkDecision's defensive
|
|
// defaults ({} / '') would resume with an empty input/result the user didn't approve.
|
|
const incomplete = findIncompleteDecisions(resolutions);
|
|
if (incomplete.length > 0) {
|
|
return {
|
|
status: 400,
|
|
error: 'edit requires editedArguments and respond requires responseText',
|
|
incomplete,
|
|
};
|
|
}
|
|
return { resumeValue: mapToolApprovalResolutions(resolutions) };
|
|
}
|
|
if (payload?.type === 'ask_user_question') {
|
|
if (typeof body.answer !== 'string' || body.answer.length === 0) {
|
|
return { status: 400, error: 'An answer is required' };
|
|
}
|
|
// The answer becomes a ToolMessage the model must ingest — bound it like any
|
|
// other user-controlled wire field rather than trusting the client.
|
|
if (body.answer.length > MAX_ASK_ANSWER_LENGTH) {
|
|
return { status: 400, error: 'Answer exceeds the maximum length' };
|
|
}
|
|
return { resumeValue: mapAskUserAnswer({ answer: body.answer }) };
|
|
}
|
|
return { status: 400, error: 'Unsupported pending action type' };
|
|
}
|
|
|
|
/**
|
|
* Finalize a resumed turn that ran to completion: persist the (now complete)
|
|
* response message, emit the terminal event over the existing SSE, complete the
|
|
* job, and prune the checkpoint. Mirrors the abort route's save shape but for a
|
|
* successful finish. Best-effort title generation for a first-turn pause.
|
|
*/
|
|
async function finalizeResumedTurn({ req, client, job, streamId, conversationId, addTitle }) {
|
|
const userId = req.user.id;
|
|
const checkpointerCfg = req.config?.endpoints?.[EModelEndpoint.agents]?.checkpointer;
|
|
const meta = job.metadata ?? {};
|
|
const userMessage = meta.userMessage;
|
|
// The response hangs off the user message; the *user* message's own parent decides
|
|
// whether this is the first turn of the conversation (title eligibility).
|
|
const parentMessageId = userMessage?.messageId ?? Constants.NO_PARENT;
|
|
const isFirstTurn = (userMessage?.parentMessageId ?? Constants.NO_PARENT) === Constants.NO_PARENT;
|
|
const responseMessageId = meta.responseMessageId ?? `${userMessage?.messageId ?? 'resumed'}_`;
|
|
// Sourced from the paused job (persisted at creation), not the resume body — a
|
|
// temporary chat must stay temporary on resume so its messages aren't persisted.
|
|
const isTemporary = meta.isTemporary ?? req.body?.isTemporary;
|
|
|
|
// Read the raw job data BEFORE completeJob deletes it — its tracked token/context
|
|
// usage backs the response message's cost rollup (parity with normal completion).
|
|
const jobData = await GenerationJobManager.getJobStore().getJob(streamId);
|
|
|
|
// Job-replacement guard (mirrors the normal request path): jobs are keyed by streamId
|
|
// (== conversationId), so a new/concurrent request reusing this conversation overwrites
|
|
// the record with a fresh createdAt. If that happened while we were resuming, finalizing
|
|
// now would emit `done` to / complete / delete the NEWER turn's job. Skip all terminal
|
|
// side effects when the job we paused is no longer the live one; the caller's `finally`
|
|
// still disposes the client + releases the slot.
|
|
if (!jobData || jobData.createdAt !== job.createdAt) {
|
|
logger.warn(
|
|
`[ResumeAgentController] Skipping resumed finalization — job ${streamId} was replaced`,
|
|
);
|
|
return;
|
|
}
|
|
// Prefer the resumed run's live content: it's complete (seeded with the pre-pause
|
|
// content) and avoids a Redis re-read that can race appendChunk writes still in
|
|
// flight. Fall back to the aggregated store content only when the live array is empty.
|
|
const liveContent = Array.isArray(client?.contentParts) ? client.contentParts : [];
|
|
const rawContent =
|
|
liveContent.length > 0
|
|
? liveContent
|
|
: ((await GenerationJobManager.getResumeState(streamId))?.aggregatedContent ?? []);
|
|
// Parity with the normal agents path (AgentClient strips these before saving):
|
|
// drop empty/malformed tool_call parts so a resumed turn can't persist an invalid
|
|
// part that breaks reload/rendering.
|
|
const content = filterMalformedContentParts(rawContent);
|
|
|
|
const responseMessage = {
|
|
messageId: responseMessageId,
|
|
parentMessageId,
|
|
conversationId,
|
|
content,
|
|
sender: meta.sender ?? client?.sender ?? 'AI',
|
|
endpoint: meta.endpoint,
|
|
iconURL: meta.iconURL,
|
|
model: meta.model,
|
|
unfinished: false,
|
|
error: false,
|
|
isCreatedByUser: false,
|
|
user: userId,
|
|
};
|
|
if (meta.agent_id ?? req.body?.agent_id) {
|
|
responseMessage.agent_id = meta.agent_id ?? req.body.agent_id;
|
|
}
|
|
// Persist tool artifacts (code files, images, UI resources) the resumed continuation
|
|
// produced — BaseClient.sendMessage awaits these before saving, but the lean resume
|
|
// path bypasses it, so do it here or they vanish on reload / for late subscribers.
|
|
// MERGE with any already on the row (earlier pause segments) rather than overwrite —
|
|
// the final segment's client only holds its own segment's artifacts.
|
|
const attachments = await resolveAccumulatedAttachments({
|
|
client,
|
|
conversationId,
|
|
responseMessageId,
|
|
});
|
|
if (attachments.length > 0) {
|
|
responseMessage.attachments = attachments;
|
|
}
|
|
|
|
// Response metadata: the resume client only sees POST-resume usage, while the job's
|
|
// tracked tokenUsage is cumulative across the pause. Take the cumulative usage (+
|
|
// summary marker) from the job, and contextUsage / thoughtSignatures from the client
|
|
// (which the abort-only helper drops). Cumulative usage wins so cost isn't underreported.
|
|
const clientMeta = client?.buildResponseMetadata?.() ?? null;
|
|
const cumulativeMeta = jobData ? buildAbortedResponseMetadata(jobData) : null;
|
|
const responseMetadata = {
|
|
...(clientMeta ?? {}),
|
|
...(cumulativeMeta?.usage ? { usage: cumulativeMeta.usage } : {}),
|
|
...(cumulativeMeta?.summaryUsedTokens != null
|
|
? { summaryUsedTokens: cumulativeMeta.summaryUsedTokens }
|
|
: {}),
|
|
};
|
|
if (Object.keys(responseMetadata).length > 0) {
|
|
responseMessage.metadata = responseMetadata;
|
|
}
|
|
// Carry the resumed run's context-window calibration (BaseClient.sendMessage persists
|
|
// this on the response). Without it, the NEXT turn can't seed its pruner from this
|
|
// run and falls back to uncalibrated token accounting.
|
|
if (client?.contextMeta != null) {
|
|
responseMessage.contextMeta = client.contextMeta;
|
|
}
|
|
|
|
await saveMessage(
|
|
{ userId, isTemporary, interfaceConfig: req?.config?.interfaceConfig },
|
|
responseMessage,
|
|
{ context: 'api/server/controllers/agents/resume.js - resumed response end' },
|
|
);
|
|
|
|
const convo = await getConvo(userId, conversationId);
|
|
const conversation = { ...(convo ?? {}), conversationId };
|
|
|
|
// First-turn pause: the title was deferred when the turn paused. Generate it BEFORE
|
|
// completing the stream so the `title` event still reaches the live client (emitChunk
|
|
// no-ops once completeJob tears down the runtime) and the final event carries the real
|
|
// title instead of "New Chat". Best-effort — a failure must not fail the resumed turn.
|
|
if (
|
|
addTitle &&
|
|
isFirstTurn &&
|
|
!isTemporary &&
|
|
userMessage?.text &&
|
|
(!convo || !convo.title || convo.title === 'New Chat')
|
|
) {
|
|
try {
|
|
await addTitle(req, {
|
|
text: userMessage.text,
|
|
conversationId,
|
|
client,
|
|
onTitleGenerated: ({ conversationId: titleConvoId, title }) => {
|
|
conversation.title = title;
|
|
return GenerationJobManager.emitChunk(streamId, {
|
|
event: 'title',
|
|
data: { conversationId: titleConvoId, title },
|
|
});
|
|
},
|
|
});
|
|
} catch (err) {
|
|
logger.error('[ResumeAgentController] Title generation failed after resume', err);
|
|
}
|
|
}
|
|
conversation.title = conversation.title || 'New Chat';
|
|
|
|
// Re-check ownership immediately before the terminal writes. The start-of-function
|
|
// guard can go stale across the awaits above: saveMessage and (first-turn) title
|
|
// generation can take long enough for a new request to replace this job on the same
|
|
// conversationId (streamId == conversationId). Without this second read, emitDone /
|
|
// completeJob / prune below would emit `done` to and tear down the REPLACEMENT job —
|
|
// the same hazard the catch-path guard prevents on the failure path.
|
|
const liveJobBeforeFinalize = await GenerationJobManager.getJobStore().getJob(streamId);
|
|
if (!liveJobBeforeFinalize || liveJobBeforeFinalize.createdAt !== job.createdAt) {
|
|
logger.warn(
|
|
`[ResumeAgentController] Skipping resumed terminal writes — job ${streamId} was replaced mid-finalize`,
|
|
);
|
|
return;
|
|
}
|
|
|
|
const finalEvent = {
|
|
final: true,
|
|
conversation,
|
|
title: conversation.title,
|
|
requestMessage: userMessage
|
|
? sanitizeMessageForTransmit({
|
|
...userMessage,
|
|
conversationId,
|
|
isCreatedByUser: true,
|
|
// job.metadata.userMessage is persisted without files; carry the restored
|
|
// uploads (seeded onto req.body.files before reconstruction) so the final SSE
|
|
// doesn't blank the user bubble's attachments — matching the normal path.
|
|
...(Array.isArray(req.body?.files) && req.body.files.length > 0
|
|
? { files: req.body.files }
|
|
: {}),
|
|
})
|
|
: null,
|
|
responseMessage: { ...responseMessage },
|
|
};
|
|
|
|
await GenerationJobManager.emitDone(streamId, finalEvent);
|
|
// Awaited (not fire-and-forget) so the job's terminal write lands before the
|
|
// checkpoint prune, and so a failure here doesn't race the controller's error path.
|
|
try {
|
|
await GenerationJobManager.completeJob(streamId);
|
|
} catch (completeErr) {
|
|
logger.error('[ResumeAgentController] Failed to complete resumed turn', completeErr);
|
|
}
|
|
await deleteAgentCheckpoint(conversationId, checkpointerCfg);
|
|
}
|
|
|
|
/**
|
|
* Resume a generation that paused for human-in-the-loop review.
|
|
*
|
|
* The original run lives in a detached background task that exits when the run
|
|
* pauses, so this REBUILDS the run from the durable checkpoint (same `thread_id`)
|
|
* and continues it with the user's decision. The continuation streams over the
|
|
* client's existing SSE (events flow through the same `streamId`).
|
|
*
|
|
* Flow: authorize → map decisions → atomically claim the resume (single-winner) →
|
|
* ACK → reconstruct the client → `resumeCompletion` → finalize (or re-pause).
|
|
*
|
|
* Shares chat.js's middleware (auth, agent access, `buildEndpointOption`) so the
|
|
* agent/endpoint are reconstructed from the request exactly like a normal turn.
|
|
*
|
|
* @param {express.Request} req
|
|
* @param {express.Response} res
|
|
* @param {express.NextFunction} next
|
|
* @param {Function} initializeClient
|
|
* @param {Function} addTitle
|
|
*/
|
|
const ResumeAgentController = async (req, res, next, initializeClient, addTitle) => {
|
|
const userId = req.user.id;
|
|
const { conversationId, actionId } = req.body;
|
|
const streamId = conversationId;
|
|
|
|
if (!streamId || streamId === 'new') {
|
|
return res.status(400).json({ error: 'conversationId is required to resume' });
|
|
}
|
|
|
|
const job = await GenerationJobManager.getJob(streamId);
|
|
if (!job) {
|
|
return res.status(404).json({ error: 'No paused generation for this conversation' });
|
|
}
|
|
if (job.metadata?.userId && job.metadata.userId !== userId) {
|
|
return res.status(403).json({ error: 'Unauthorized' });
|
|
}
|
|
if (hasTenantMismatch(job, req.user)) {
|
|
return res.status(403).json({ error: 'Unauthorized' });
|
|
}
|
|
|
|
// The resume must rebuild the SAME agent/endpoint that paused. Require an EXACT
|
|
// agent_id match when the paused job had one — a request that omits agent_id (or
|
|
// claims an ephemeral / non-agents endpoint) must not rebuild the claimed checkpoint
|
|
// on a different graph. The conversation's agent is stable, so a correct client always
|
|
// sends the right one.
|
|
const originalAgentId = job.metadata?.agent_id;
|
|
if (originalAgentId && req.body.agent_id !== originalAgentId) {
|
|
return res.status(403).json({ error: 'Cannot resume with a different agent' });
|
|
}
|
|
// Require an EXACT endpoint match (like agent_id): a request that OMITS endpoint must
|
|
// not fall through — the shared chat middleware treats a missing/non-agents endpoint
|
|
// as the ephemeral agent, so omitting it could rebuild the claimed checkpoint on a
|
|
// different graph. A correct client always echoes the paused endpoint.
|
|
const originalEndpoint = job.metadata?.endpoint;
|
|
if (originalEndpoint && req.body.endpoint !== originalEndpoint) {
|
|
return res.status(403).json({ error: 'Cannot resume on a different endpoint' });
|
|
}
|
|
|
|
const pendingAction = job.metadata?.pendingAction;
|
|
if (job.status !== 'requires_action') {
|
|
return res.status(409).json({ error: 'No live pending action to resume' });
|
|
}
|
|
if (isPendingActionStale({ pendingAction })) {
|
|
// The action expired between the pending-action SSE and this submit. Drive the expiry
|
|
// NOW (expire CAS + terminal SSE) instead of waiting for the periodic sweeper —
|
|
// otherwise the job sits `requires_action` with a dead action and any attached SSE
|
|
// client never gets a terminal event, so the stream appears to hang even though the
|
|
// UI already reported the action as expired.
|
|
try {
|
|
await GenerationJobManager.expireApproval(streamId, pendingAction?.actionId);
|
|
} catch (err) {
|
|
logger.warn(
|
|
'[ResumeAgentController] Failed to expire stale action on submit',
|
|
err?.message ?? err,
|
|
);
|
|
}
|
|
return res.status(409).json({ error: 'No live pending action to resume' });
|
|
}
|
|
// Require the actionId the UI sends: without it, a stale/malformed client could
|
|
// resolve whatever action is currently pending (e.g. answer a different question).
|
|
if (!actionId) {
|
|
return res.status(400).json({ error: 'actionId is required to resume' });
|
|
}
|
|
if (pendingAction.actionId !== actionId) {
|
|
return res.status(409).json({ error: 'This decision targets a stale action' });
|
|
}
|
|
|
|
// Pin the graph identity: the resume must rebuild the SAME agent/graph + tool set the
|
|
// run paused on. The agent_id + endpoint guards above cover saved agents; the
|
|
// fingerprint additionally catches an ephemeral-agent config swap (its agent_id is
|
|
// undefined, so the id guard can't tell two ephemeral configs apart). Enforced only
|
|
// when the paused action carries a fingerprint (in-flight pauses from before this
|
|
// change won't), and recomputed from the resume body's graph-determining fields.
|
|
const pinnedFingerprint = pendingAction.requestFingerprint;
|
|
if (pinnedFingerprint && pinnedFingerprint !== computeAgentRequestFingerprint(req.body ?? {})) {
|
|
return res.status(403).json({ error: 'Cannot resume with a different agent configuration' });
|
|
}
|
|
|
|
const mapped = resolveResumeValue(pendingAction, req.body);
|
|
if (mapped.error) {
|
|
return res.status(mapped.status).json({
|
|
error: mapped.error,
|
|
...(mapped.undecided && { undecided: mapped.undecided }),
|
|
...(mapped.disallowed && { disallowed: mapped.disallowed }),
|
|
...(mapped.incomplete && { incomplete: mapped.incomplete }),
|
|
});
|
|
}
|
|
|
|
// Count the resume against the concurrency limit. The original turn released its slot
|
|
// when it paused, so resuming must re-acquire one — otherwise pausing several turns
|
|
// and resuming them at once would bypass LIMIT_CONCURRENT_MESSAGES.
|
|
const { allowed } = await checkAndIncrementPendingRequest(userId);
|
|
if (!allowed) {
|
|
return res.status(429).json({ error: 'Too many concurrent requests' });
|
|
}
|
|
|
|
// Atomically claim the resume. The single winner drives the run; a racing second
|
|
// submit (double-click, two tabs) gets false and must not re-drive — that would
|
|
// re-execute tools and double-bill.
|
|
//
|
|
// The claim runs AFTER the slot increment above but BEFORE the run's own try/finally
|
|
// that releases it, so a store/Redis error here (unlike the clean `!claimed` branch)
|
|
// would leak the concurrency slot until the counter TTL expires — spuriously 429'ing
|
|
// the user when they retry the still-paused approval. Release the slot on that path too.
|
|
let claimed;
|
|
try {
|
|
claimed = await GenerationJobManager.approvals.resolve(streamId, pendingAction.actionId);
|
|
} catch (err) {
|
|
await decrementPendingRequest(userId);
|
|
logger.error('[ResumeAgentController] Failed to claim resume', err);
|
|
return res.status(500).json({ error: 'Failed to resume' });
|
|
}
|
|
if (!claimed) {
|
|
await decrementPendingRequest(userId);
|
|
return res.status(409).json({ error: 'This action was already resolved or has expired' });
|
|
}
|
|
|
|
// Seed the run-scoped MCP request-context store BEFORE the ACK: once `res.json`
|
|
// finishes the response, a later `getMCPRequestContext(req, res)` (from tool loading)
|
|
// sees `res` as ended and returns undefined, leaving the resumed run without its MCP
|
|
// connection store — approved MCP / OAuth-overlay tools would then run without their
|
|
// request-scoped connections. Pre-seeding with a null `res` + `cleanupOnResponse:false`
|
|
// mirrors the normal stream path (request.js); torn down in the `finally` below.
|
|
req._resumableStreamId = streamId;
|
|
getMCPRequestContext(req, undefined, { cleanupOnResponse: false });
|
|
|
|
// ACK immediately; the continuation streams over the client's existing SSE.
|
|
res.json({ streamId, conversationId, status: 'resuming' });
|
|
|
|
// Seed the original thread parent BEFORE initializeClient: initializeAgent scopes
|
|
// thread files / code artifacts off `req.body.parentMessageId`, and the resume body
|
|
// doesn't carry it. This is the user message's parent (the thread position);
|
|
// `client.parentMessageId` below is a different value — the response's parent, i.e.
|
|
// the user message id.
|
|
req.body.parentMessageId = job.metadata.userMessage?.parentMessageId ?? Constants.NO_PARENT;
|
|
|
|
// Restore the paused user message's OWN uploaded files. initializeAgent rebuilds
|
|
// code/file sessions by walking the conversation from `parentMessageId`, but
|
|
// execute-code files are excluded from that lookup, so files uploaded on the paused
|
|
// turn would be dropped — an approved code/read-file tool would resume without them.
|
|
//
|
|
// SECURITY: ALWAYS source files from the paused job, never from the `/resume` body.
|
|
// `files` is not pinned by the resume fingerprint or replayed via resumeContext, so
|
|
// honoring a client-supplied `files` array would let a crafted/buggy client resume an
|
|
// approved code/read-file tool against a DIFFERENT file set than the one the user
|
|
// approved. A resume reconstructs the SAME paused turn, so there is no legitimate
|
|
// reason for the client to supply its own files. Prefer the files persisted on the JOB
|
|
// at onStart (race-free), fall back to the DB row for older jobs, and CLEAR otherwise
|
|
// so a client-supplied set can never leak through.
|
|
const metaFiles = job.metadata.userMessage?.files;
|
|
if (Array.isArray(metaFiles) && metaFiles.length > 0) {
|
|
req.body.files = metaFiles;
|
|
} else {
|
|
let restoredFiles = false;
|
|
const pausedUserMessageId = job.metadata.userMessage?.messageId;
|
|
if (pausedUserMessageId) {
|
|
try {
|
|
const [row] = await getMessages(
|
|
{ conversationId, messageId: pausedUserMessageId },
|
|
'files',
|
|
);
|
|
if (Array.isArray(row?.files) && row.files.length > 0) {
|
|
req.body.files = row.files;
|
|
restoredFiles = true;
|
|
}
|
|
} catch (err) {
|
|
logger.warn(
|
|
'[ResumeAgentController] Failed to restore paused user message files',
|
|
err?.message ?? err,
|
|
);
|
|
}
|
|
}
|
|
if (!restoredFiles) {
|
|
// No paused files (or the lookup failed): drop any client-supplied files so a
|
|
// crafted resume body can't inject a file set the paused turn never had.
|
|
req.body.files = [];
|
|
}
|
|
}
|
|
|
|
// Restore the conversation's createdAt so temporal prompt vars ({{current_datetime}},
|
|
// {{iso_datetime}}, ...) resolve against the SAME anchor the paused graph used rather
|
|
// than the resume wall-clock. initializeAgent reads `req.conversationCreatedAt`; the
|
|
// normal path sets it from the convo timestamp (resolveConversationCreatedAt), so mirror
|
|
// that here. (The original `timezone` is replayed onto req.body via RESUME_CONTEXT_KEYS.)
|
|
try {
|
|
const resumedConvo = await getConvo(userId, conversationId);
|
|
const createdAt = resumedConvo?.createdAt ? new Date(resumedConvo.createdAt) : null;
|
|
if (createdAt && !Number.isNaN(createdAt.getTime())) {
|
|
req.conversationCreatedAt = createdAt.toISOString();
|
|
}
|
|
} catch (err) {
|
|
logger.warn(
|
|
'[ResumeAgentController] Failed to restore conversation timestamp anchor',
|
|
err?.message ?? err,
|
|
);
|
|
}
|
|
|
|
let client = null;
|
|
try {
|
|
const result = await initializeClient({
|
|
req,
|
|
res,
|
|
endpointOption: req.body.endpointOption,
|
|
signal: job.abortController.signal,
|
|
});
|
|
client = result.client;
|
|
|
|
// Bind the rebuilt client to the in-flight turn's identity (no new user message).
|
|
client.conversationId = streamId;
|
|
// The resume operates on the SAME job (it moved it running again), so its identity is
|
|
// the paused job's createdAt — used by the re-pause CAS pre-check + checkpoint prune to
|
|
// avoid acting on a job a newer request has since replaced.
|
|
client.jobCreatedAt = job.createdAt;
|
|
client.responseMessageId = job.metadata.responseMessageId;
|
|
client.parentMessageId = job.metadata.userMessage?.messageId ?? Constants.NO_PARENT;
|
|
// Read the pre-pause content BEFORE swapping the store's content reference: the
|
|
// in-memory store's setContentParts REPLACES the stored array, so reading the
|
|
// resume state afterward would see the new (empty) client array and lose the seed.
|
|
const resumeState = await GenerationJobManager.getResumeState(streamId);
|
|
let seedContent = resumeState?.aggregatedContent ?? [];
|
|
// Stamp the answered question onto the paused ask_user_question tool-call part
|
|
// (args = the pendingAction's authoritative question, output = the user's answer):
|
|
// the streamed arg chunks carry no tool name so the aggregator dropped them, and
|
|
// no completion event ever fires for this tool — without this the saved part is
|
|
// an empty "cancelled-looking" tool call. See attachAskUserQuestionAnswer.
|
|
if (pendingAction.payload?.type === 'ask_user_question') {
|
|
seedContent = attachAskUserQuestionAnswer(
|
|
seedContent,
|
|
pendingAction.payload.question,
|
|
req.body.answer,
|
|
);
|
|
}
|
|
if (client.contentParts) {
|
|
GenerationJobManager.setContentParts(streamId, client.contentParts);
|
|
}
|
|
|
|
await client.resumeCompletion({
|
|
resumeValue: mapped.resumeValue,
|
|
seedContent,
|
|
abortController: job.abortController,
|
|
// Carry the user's MCP auth so approved MCP tools run with their credentials.
|
|
userMCPAuthMap: result.userMCPAuthMap,
|
|
// Replay deferred tools discovered before the pause (captured at pause). The rebuilt
|
|
// graph passes `messages: []`, so without these an approved deferred tool would be
|
|
// absent from the schema-only toolMap and resume would fail with "unknown tool".
|
|
discoveredToolNames: job.metadata?.discoveredTools,
|
|
});
|
|
|
|
// The model may pause AGAIN (another tool, or a follow-up question). The pending
|
|
// action is already persisted + emitted; leave the job `requires_action`.
|
|
if (client.pendingApproval) {
|
|
logger.debug(`[ResumeAgentController] Re-paused for approval: ${streamId}`);
|
|
// Persist this segment's content + artifacts before the fresh client (next
|
|
// resume) drops them, so an expiring re-pause doesn't lose them; finalize later
|
|
// overwrites content and merges attachments onto the saved message.
|
|
await persistRePauseProgress({ req, client, job, streamId, conversationId });
|
|
return;
|
|
}
|
|
|
|
// If the user aborted mid-resume, the abort route already emitted the terminal
|
|
// event and finalized the job — don't double-save / double-finalize here.
|
|
if (job.abortController.signal.aborted) {
|
|
logger.debug(
|
|
`[ResumeAgentController] Aborted during resume; abort route finalizes: ${streamId}`,
|
|
);
|
|
return;
|
|
}
|
|
|
|
await finalizeResumedTurn({ req, client, job, streamId, conversationId, addTitle });
|
|
} catch (err) {
|
|
logger.error('[ResumeAgentController] Resume failed', err);
|
|
// Job-replacement guard (mirrors finalizeResumedTurn's success-path guard): if a
|
|
// newer request reused this conversationId while the resume was failing, do NOT emit
|
|
// the error to / complete / prune the NEWER turn's job. The finally still releases
|
|
// the slot + disposes. Proceed with finalization if the replacement check itself fails.
|
|
let stillLive = true;
|
|
try {
|
|
const liveJob = await GenerationJobManager.getJobStore().getJob(streamId);
|
|
stillLive = !!liveJob && liveJob.createdAt === job.createdAt;
|
|
} catch (readErr) {
|
|
logger.warn('[ResumeAgentController] Replacement check failed; finalizing anyway', readErr);
|
|
}
|
|
if (!stillLive) {
|
|
logger.warn(
|
|
`[ResumeAgentController] Skipping failed-resume finalization — job ${streamId} was replaced`,
|
|
);
|
|
} else {
|
|
try {
|
|
await GenerationJobManager.emitError(streamId, err?.message ?? 'Resume failed');
|
|
} catch (emitErr) {
|
|
logger.error('[ResumeAgentController] Failed to emit resume error', emitErr);
|
|
}
|
|
try {
|
|
await GenerationJobManager.completeJob(streamId, err?.message ?? 'Resume failed');
|
|
} catch (completeErr) {
|
|
logger.error('[ResumeAgentController] Failed to finalize failed resume', completeErr);
|
|
// Last resort: force a terminal state so the job isn't orphaned in `running`.
|
|
await GenerationJobManager.getJobStore()
|
|
.updateJob(streamId, {
|
|
status: 'error',
|
|
completedAt: Date.now(),
|
|
error: 'Resume failed',
|
|
})
|
|
.catch((updErr) =>
|
|
logger.error('[ResumeAgentController] Fallback job finalize failed', updErr),
|
|
);
|
|
}
|
|
await deleteAgentCheckpoint(
|
|
conversationId,
|
|
req.config?.endpoints?.[EModelEndpoint.agents]?.checkpointer,
|
|
);
|
|
}
|
|
} finally {
|
|
// Tear down the MCP request-context store seeded before the ACK (parity with
|
|
// request.js's finishResumableRequest). No-op if it was never seeded.
|
|
await cleanupMCPRequestContextForReq(req);
|
|
// Release the concurrency slot taken above — UNLESS handleRunInterrupt already
|
|
// released it on a re-pause (so a fast /resume isn't 429'd). On a normal finish or
|
|
// error it didn't, so release here. A re-pause re-acquires its own slot next resume.
|
|
if (!client?.pendingRequestReleased) {
|
|
await decrementPendingRequest(userId);
|
|
}
|
|
if (client) {
|
|
disposeClient(client);
|
|
}
|
|
}
|
|
};
|
|
|
|
module.exports = ResumeAgentController;
|