mirror of
https://github.com/danny-avila/LibreChat.git
synced 2026-09-04 05:28:30 +00:00
* style: Unify message row layout and edit surfaces Route chat, share, and search messages through a shared MessageRow so user turns render as right-aligned bubbles and assistant turns keep a visible identity column. Replace per-part text editors with one edit surface that keeps tools, errors, and artifacts visible. Preserve non-text fields when saving content parts, copy the full serialized message, and hide hover actions that do not apply during streaming or errors. * style: Align edit footer and lighten editor field in dark mode Drop the divider above the user edit footer so both edit surfaces share the same footer treatment. Move the editor fields to surface-tertiary-alt. Light mode is unchanged at #fff, while dark mode lifts from #0d0d0d to #2f2f2f so the field sits above the #212121 panel instead of sinking into near-black. * style: Drop focus border and ring from message editors The editor fields changed border color and added a ring on focus. Keep the border static and rely on the app-level focus handling instead. * fix: Keep a triggered message action visible when the row is not hovered Hover actions fade out on non-last rows, and mobile.css only restored display and visibility for an active button, never opacity. Opening the fork popover therefore left it anchored to an invisible trigger once the pointer left the row. Skip the fade entirely while a button is active. Extract the recipe the three toolbars repeated so the rule has one home. Rework the streaming guard to the contract the toolbar now implements: edit and fork are omitted from a streaming response rather than rendered disabled, and the settled turn above keeps its own actions. It asserted the removed disabled-and-transparent behaviour and its opacity check only held because the growing response shifted the row out from under the pointer. * style: Trim message edit chrome and stabilize the status row The edit surface was a titled card sitting inside the conversation: a bordered panel with an "Edit message" heading wrapping bordered fields, which read as a settings dialog rather than an inline editor. Drop the card background, border and heading, and take the footer buttons down to the small size so the editor reads as a field in the message flow. The captured row goes from 253px to 187px. Move "Unsaved changes" into the footer and merge the rerun hint into the same slot. Both previously added their own row, so typing pushed the rest of the conversation down. The slot is clamped to two lines, which stays under the 36px button row, so the footer height holds at 36px regardless of which message is showing. * test: Cover message edit layout stability Add a mock e2e spec that measures the edit footer and section boxes and asserts they hold steady as the status text appears, for both the single-part user editor and a multi-part response. The multi-part case needs an assistant message with two editable parts, so add an E2E_THINK_REPLY marker to the fake model. Its think tags are parsed downstream by the agents stream pipeline, which yields a reasoning part followed by a text part. * fix: Read the fork popover open state from its store Fork mirrored the popover state into its own useState and reset it from an onClose prop. Ariakit 0.4 has no onClose, and React's DOM types accept the name on any element, so it type-checked, landed on a div and never fired. Closing by Escape or an outside click therefore left the button reading as active until the trigger was clicked again. Read the state from the store instead so every close path clears it. * fix: Keep the whole toolbar visible while an action is open Only the triggered button escaped the hover fade, so opening the editor or the fork popover left the row as a single floating button once the pointer moved away. Mark the active button and have every action in the toolbar key off it, so the group stays opaque for as long as a surface is open. The marker is a dedicated class rather than the existing `active`, which HoverButtons pins to the edit button of every assistant message and would hold those toolbars open permanently. The existing guard pressed Escape to close the editor while focus sat on the body, so the editor never closed and its assertion only held because the sibling faded regardless. Close the editor through its own control, and drop focus before measuring the fade now that Escape returns it to the trigger. * fix: Withhold copy while a response is still streaming Text-to-speech, fork and feedback were all withheld from a message that is still generating, but copy was rendered throughout, so the button offered to put half a sentence on the clipboard. Gate it on the same condition. That empties the toolbar for the duration, and SubRow collapses an empty row, so a streaming response now carries no actions at all until it settles. Both guards encoded the old contract: the unit test asserted copy was present and counted a single button, and the browser guard used copy as its proof that the toolbar had mounted. The settled turn above takes over that role. * fix: Move retry navigation to the outer edge of a user turn A user turn is right-aligned, but its sibling navigation rendered ahead of the actions, so the retry counter sat inboard of the icons instead of under the edge of the bubble it belongs to. Order it last on user turns. * fix: Ride the stream instead of chasing it Following a generating answer went through a helper throttled at 145ms, so the thread caught up in visible jerks rather than flowing. It now writes the scroll position directly on each frame, which is what an answer arriving a few pixels at a time actually needs, and glides only for the one long trip a turn makes, when sending has to travel from wherever the reader was down to the newest word. Whether to follow at all is now answered by where the reader is and which way they were going, rather than by the abort flag. `useMessageProcess` raises that flag on any wheel at all, downward ones included, through a throttle whose trailing call lands after the gesture has ended, so nothing timed to the gesture could outlive it. Scrolling down to the newest word could therefore never resume the ride, while the scroll-to-bottom button, which touches no wheel, always could. Arrival is judged on the scroll it produces rather than the wheel tick that started it, because wheel scrolling is animated and at tick time the thread is still far short of where the tick is taking it. Arriving also counts from further out than leaving does: while an answer streams the end recedes between the last tick and the frame that measures it, so judging arrival as tightly as departure leaves a reader unable to catch it at all. * fix: Reveal retry navigation on hover while an answer generates Copy, edit, fork and read-aloud are all withheld from a response that is still generating, which left the retry counter as the only thing rendering under a half-written answer. It now reveals on hover there, like the actions it sits with, and stays put on a settled turn. * fix: Keep a refused rerun from discarding the edit While a response is streaming, the edit action stays available on every earlier row, and those editors see a per-message submitting flag that is false, so Update and rerun is enabled. The send itself is still refused: ask() returns false for the duration of the active submission. Both editors ignored that and closed anyway, so the draft went with them and no rerun ever started. Both rerun paths now check the result and leave the editor untouched when the send is refused, so the work survives until the thread is free. * fix: Let an upward gesture beat the pending send glide Sending arms a smooth glide down to the newest word, and the landing re-pins the thread to the bottom. The landing was scheduled two ways, on scrollend and on a 700ms fallback, and neither was ever cancelled. A reader who changed their mind and headed up mid-flight was pinned again regardless, then dragged back by the next streaming resize. The fallback fires for the whole window, so this held even after the glide had visibly settled. The gesture now marks the glide interrupted, wherever it lets go of the bottom, and the landing stands down when it sees that. A glide the reader leaves alone still re-affirms the ride. * fix: Fade retry navigation on every streaming response format Every other action is withheld from the row that is still generating, so the retry counter is the only thing left under a half-written answer. The plain text row already faded it to hover-only there; the structured rows did not, and left it sitting on its own. Both structured paths now apply the same condition, and the class string the three of them share moves next to the hover action styles it belongs with. * i18n: Correct the copy the edit surface rewrite left behind The multi-part hint told the reader to save first and then rerun, but a save closes the editor and reopening seeds the drafts from what was just saved, so there is nothing left to rerun and the button stays disabled. Rerunning carries a single edited section by design, so the hint now states that limit rather than pointing at a step that is not there. Drop com_ui_save_submit as well: the per-part editor that used it is gone. * test: Make the message visual baselines opt-in The suite asserts sixteen screenshots and the repository tracks none, so Playwright's default treats every one as a miss and the mock e2e job fails on Linux. Baselines only compare cleanly against the machine that produced them, and nothing here can generate ones that match the runner image. The flows keep running and asserting their structure, which is where their value was; only the pixel comparison is now gated behind E2E_VISUAL_SNAPSHOTS. * style: Restore import order in the reworked message files The repository sorter and CI disagreed with what these files were left holding after the edit surface rework. No behavior change. * test: Follow the reworded rerun hint in the edit layout spec The multi-part hint was restated in the previous commit; this assertion still expected the old wording and would have failed the mock e2e suite. * fix: Leave the send glide alone while the answer streams in Every delta of an answer reruns the scroll effect, and the plain follow writes scrollTop outright, which cancels an animation on its first frame. So the glide a send starts was killed by the first token to arrive and the reader was snapped down instead of carried. The follow now stands down while a glide is travelling, which is what the hook already documented but only enforced on the resize path. * fix: Write a saved edit onto the thread as it stands An earlier turn stays editable while the newest answer streams, and the save captured the thread before the request but wrote it back after. Every delta that landed during the round trip was overwritten. Most of the time the next delta re-merged and the damage showed as a one-frame truncation, but a save that resolved after the stream's final write left the cache wrong for the rest of the session. The thread is now read once the request has resolved, which is what the content part editor already did. The editor actions in this file also wrap again rather than hold one unbreakable row, for the reason given in the following commit. * fix: Let the editor actions wrap on a narrow row At 320px an assistant turn gives the editor about 252px once page padding, the identity column and the row gap are taken out, and Cancel, Save and Update & rerun need more than that in English alone. The group was pinned with shrink-0, so it ran past the edge of the row instead of wrapping. A longer translated label makes it worse, and the user turn had no margin left either. Both editors wrap again, which is what the footer did before the status row was folded into it. * fix: Catch up to the new bottom when the glide lands Following stands down for the length of the glide, so an answer that arrives while it travels moves the bottom past the target the glide aimed at. A short response that finished before the glide reported landing left the thread a few lines short of its own end, with nothing left to correct it. Landing now closes whatever gap opened, unless the reader took over on the way. * test: Follow the renamed rerun button in the edit flow specs The button became 'Update & rerun' when the edit surfaces were unified, but two edit-flow specs still located 'Save & Submit' and would have waited for it until they timed out. A type comment named the old button too. * fix: Judge the first thread scroll against a real position The direction check seeded its last-position ref at 0, so the first scroll event on an opened thread, which arrives carrying a large positive scrollTop, read as a jump downward. Near the end that cleared the abort flag and re-pinned a reader to the stream they were scrolling away from. Take the first event as a baseline and judge direction from the next. * fix: Hold the content part editor to what it replaced EditContentParts took over from EditTextPart and left two of its behaviors behind. An emptied box now blocks Save and rerun instead of persisting a blank part. EditTextPart refused the same edit through its form's required rule and the sibling EditMessage still does, so both editors hold one line. The keyboard shortcuts reach the save paths directly, so they are guarded there too, and the footer says why the buttons are down. The editor also follows the chat direction again, taking dir and text alignment from the same setting EditMessage reads. * fix: Hold the footer height while a response streams Every action is withheld from the row that is still generating, and a lone sibling counter renders nothing, so the footer measured zero until the answer landed and then sprang to the height of the buttons. The transcript stepped upward under the reader at the moment a response completed. The placeholder that used to reserve this space went when the footer became unconditional, so hold the height on the row itself instead. * fix: Remember where the thread was put before judging a gesture Direction is judged against the last sample, and the thread is placed at its end without the reader touching it. With no record of where it was put, their first gesture was spent taking the baseline instead of being obeyed: a single PageUp cleared no flag of its own, so the next streamed resize rode the reader straight back to the end they were leaving. Every programmatic move now records the position it left the thread at, so the sentinel stands only until something has actually placed it. * fix: Spend the start of a turn only once it can be honored A reader who scrolls away during one answer leaves the abort flag raised, and nothing lowers it until the next connection opens, which is after this effect has already seen the send. Marking the turn as started on that first pass spent it against a closed gate: by the time the flag cleared there was no start left to honor, the reader was still detached, and the answer they had just asked for streamed on offscreen. Record the turn as started only on the pass that acts on it. * fix: Show the part edits that survived a refused save The editor saves every changed part through one button, but the endpoint takes a single part per call and nothing rolls a write back. A part the server refused therefore left the earlier ones stored while the editor reported that the message could not be saved, so cancelling from there walked away from edits that were already live. Record the writes that landed and reconcile the transcript with them whichever way the save ended. The refused parts are the only ones left holding a draft, so a retry no longer rewrites what already arrived. * fix: Stop a shared transcript from calling the sharer the reader The share row reused the chat view's user label, which reads "You". It is the screen-reader heading for the user turn, so anyone opening a share link heard every prompt the sharer wrote credited to themselves. Use the neutral "User" label on this surface. It keeps the localization the row gained, unlike the untranslated string it replaced. * fix: Let go of the stream when an interaction settles over several resizes Expanding a tool result mid-answer renders the container first and fills it once its contents arrive, so one gesture produces more than one resize. Only the first was credited to the interaction. The second read the reader as still riding the stream and put them back on the bottom they had just left. The suppressed resize now settles the ride as well as the near-bottom measure, using the position the interaction actually left the reader at, so an interaction that kept them on the end still streams. * fix: Edit inside a structured text part instead of flattening it A text content part holds either a string or a { value, annotations } object. The Assistants thread sync persists the structured form with its file citations intact, and the editor reads the part through the same union, so saving an edit wrote a bare string over the whole object and took every citation with it. The same object was handed to the tokenizer, which measures length, so a part that had been edited this way also stored a NaN token count. Write the edit into value, keep the rest of the part, and count the text itself. * fix: Keep a saved part's citations in the transcript it is written back to A text or think part holds either a bare string or a { value, annotations } object, and the editor already read both through getPartText. Writing the draft back into the local message cache put the string over the whole value, so a response carrying file citations lost them the moment it was edited and did not get them back until a refetch. Reading and writing now go through the same accessor, so an edit lands in the shape it was read from and the rest of the part survives. * fix: Let the message editor follow the chosen font size Editing a message dropped the draft to a fixed 14px regardless of the Font Size setting. On dev the textarea carried the markdown class, so it read --markdown-font-size like the rendered message does; restyling it into a bordered box replaced that with text-sm, and the new per-part editor was written the same way. Anyone on Extra Small, Large or Extra Large saw the text jump the moment they entered edit mode. Share the .message-content typography with the editors through a message-editor-text class so a draft is sized like the message it replaces and keeps tracking the setting.
1630 lines
54 KiB
JavaScript
1630 lines
54 KiB
JavaScript
const crypto = require('crypto');
|
|
const fetch = require('node-fetch');
|
|
const { logger } = require('@librechat/data-schemas');
|
|
const {
|
|
countTokens,
|
|
checkBalance,
|
|
getBalanceConfig,
|
|
buildMessageFiles,
|
|
sanitizeFileForTransmit,
|
|
extractFileContext,
|
|
getReferencedQuotes,
|
|
encodeAndFormatAudios,
|
|
encodeAndFormatVideos,
|
|
getTransactionsConfig,
|
|
encodeAndFormatDocuments,
|
|
getLangfuseTraceDestinationIds,
|
|
isLangfuseTraceSampled,
|
|
traceIdForMessage,
|
|
} = require('@librechat/api');
|
|
const {
|
|
Constants,
|
|
FileSources,
|
|
Tools,
|
|
ContentTypes,
|
|
excludedKeys,
|
|
EModelEndpoint,
|
|
mergeFileConfig,
|
|
isParamEndpoint,
|
|
isAgentsEndpoint,
|
|
isEphemeralAgentId,
|
|
supportsBalanceCheck,
|
|
isBedrockDocumentType,
|
|
getEndpointFileConfig,
|
|
} = require('librechat-data-provider');
|
|
const { getStrategyFunctions } = require('~/server/services/Files/strategies');
|
|
const { logViolation } = require('~/cache');
|
|
const TextStream = require('./TextStream');
|
|
const db = require('~/models');
|
|
|
|
const collectHistoricalFileRefs = (message) => {
|
|
const refs = [];
|
|
if (Array.isArray(message.files)) {
|
|
refs.push(...message.files);
|
|
}
|
|
if (Array.isArray(message.attachments)) {
|
|
refs.push(...message.attachments);
|
|
}
|
|
/** Steer parts carry their own attachment refs inside assistant content;
|
|
* collecting them here folds the steer replay stamp's lookup into this
|
|
* single per-turn query (see `stampSteerPartMedia`). */
|
|
if (Array.isArray(message.content)) {
|
|
for (const part of message.content) {
|
|
if (part?.type === ContentTypes.STEER && Array.isArray(part.files)) {
|
|
refs.push(...part.files);
|
|
}
|
|
}
|
|
}
|
|
return refs;
|
|
};
|
|
|
|
const collectHistoricalFileIds = (messages) => {
|
|
const fileIds = new Set();
|
|
for (const message of messages) {
|
|
for (const ref of collectHistoricalFileRefs(message)) {
|
|
if (ref?.file_id) {
|
|
fileIds.add(ref.file_id);
|
|
}
|
|
}
|
|
}
|
|
return Array.from(fileIds);
|
|
};
|
|
|
|
const buildOwnerFileFilter = (fileIds, user) => {
|
|
if (!user?.id || fileIds.length === 0) {
|
|
return null;
|
|
}
|
|
|
|
const filter = {
|
|
file_id: { $in: fileIds },
|
|
user: user.id,
|
|
};
|
|
if (user.tenantId) {
|
|
filter.tenantId = user.tenantId;
|
|
}
|
|
return filter;
|
|
};
|
|
|
|
const TOOL_ATTACHMENT_KEYS = [
|
|
Tools.file_search,
|
|
Tools.web_search,
|
|
Tools.ui_resources,
|
|
Tools.memory,
|
|
];
|
|
const DISPLAY_ATTACHMENT_FIELDS = [
|
|
'filename',
|
|
'filepath',
|
|
'expiresAt',
|
|
'conversationId',
|
|
'messageId',
|
|
'toolCallId',
|
|
'name',
|
|
];
|
|
const PER_MESSAGE_FILE_ATTACHMENT_FIELDS = ['messageId', 'toolCallId'];
|
|
|
|
const pickFields = (source, fields) => {
|
|
const picked = {};
|
|
for (const field of fields) {
|
|
if (source?.[field] !== undefined) {
|
|
picked[field] = source[field];
|
|
}
|
|
}
|
|
return picked;
|
|
};
|
|
|
|
const sanitizeDisplayOnlyAttachment = (ref) => {
|
|
if (!ref || ref.file_id) {
|
|
return undefined;
|
|
}
|
|
|
|
const attachment = pickFields(ref, DISPLAY_ATTACHMENT_FIELDS);
|
|
if (TOOL_ATTACHMENT_KEYS.includes(ref.type)) {
|
|
attachment.type = ref.type;
|
|
}
|
|
for (const key of TOOL_ATTACHMENT_KEYS) {
|
|
if (ref[key] !== undefined) {
|
|
attachment[key] = ref[key];
|
|
}
|
|
}
|
|
|
|
return Object.keys(attachment).length > 0 ? attachment : undefined;
|
|
};
|
|
|
|
const rehydrateMessageFileRefs = (refs, filesById, { preserveDisplayOnly = false } = {}) => {
|
|
if (!Array.isArray(refs)) {
|
|
return undefined;
|
|
}
|
|
|
|
const files = [];
|
|
for (const ref of refs) {
|
|
const file = filesById.get(ref?.file_id);
|
|
if (file) {
|
|
files.push({
|
|
...sanitizeFileForTransmit(file),
|
|
...pickFields(ref, PER_MESSAGE_FILE_ATTACHMENT_FIELDS),
|
|
});
|
|
continue;
|
|
}
|
|
|
|
if (preserveDisplayOnly) {
|
|
const displayOnlyAttachment = sanitizeDisplayOnlyAttachment(ref);
|
|
if (displayOnlyAttachment) {
|
|
files.push(displayOnlyAttachment);
|
|
}
|
|
}
|
|
}
|
|
return files.length > 0 ? files : undefined;
|
|
};
|
|
|
|
class BaseClient {
|
|
constructor(apiKey, options = {}) {
|
|
this.apiKey = apiKey;
|
|
this.sender = options.sender ?? 'AI';
|
|
this.currentDateString = new Date().toLocaleDateString('en-us', {
|
|
year: 'numeric',
|
|
month: 'long',
|
|
day: 'numeric',
|
|
});
|
|
/** @type {boolean} */
|
|
this.skipSaveConvo = false;
|
|
/** @type {boolean} */
|
|
this.skipSaveUserMessage = false;
|
|
/** @type {string} */
|
|
this.user;
|
|
/** @type {string} */
|
|
this.conversationId;
|
|
/** @type {string} */
|
|
this.responseMessageId;
|
|
/** @type {string} */
|
|
this.parentMessageId;
|
|
/** @type {TAttachment[]} */
|
|
this.attachments;
|
|
/** The key for the usage object's input tokens
|
|
* @type {string} */
|
|
this.inputTokensKey = 'prompt_tokens';
|
|
/** The key for the usage object's output tokens
|
|
* @type {string} */
|
|
this.outputTokensKey = 'completion_tokens';
|
|
/** @type {Set<string>} */
|
|
this.savedMessageIds = new Set();
|
|
/**
|
|
* Flag to determine if the client re-submitted the latest assistant message.
|
|
* @type {boolean | undefined} */
|
|
this.continued;
|
|
/**
|
|
* Flag to determine if the client has already fetched the conversation while saving new messages.
|
|
* @type {boolean | undefined} */
|
|
this.fetchedConvo;
|
|
/** @type {TMessage[]} */
|
|
this.currentMessages = [];
|
|
/** @type {import('librechat-data-provider').VisionModes | undefined} */
|
|
this.visionMode;
|
|
/** @type {import('librechat-data-provider').FileConfig | undefined} */
|
|
this._mergedFileConfig;
|
|
/** @type {import('librechat-data-provider').EndpointFileConfig | undefined} */
|
|
this._endpointFileConfig;
|
|
}
|
|
|
|
setOptions() {
|
|
throw new Error("Method 'setOptions' must be implemented.");
|
|
}
|
|
|
|
async getCompletion() {
|
|
throw new Error("Method 'getCompletion' must be implemented.");
|
|
}
|
|
|
|
/** @type {sendCompletion} */
|
|
async sendCompletion() {
|
|
throw new Error("Method 'sendCompletion' must be implemented.");
|
|
}
|
|
|
|
getSaveOptions() {
|
|
throw new Error('Subclasses must implement getSaveOptions');
|
|
}
|
|
|
|
async buildMessages() {
|
|
throw new Error('Subclasses must implement buildMessages');
|
|
}
|
|
|
|
async summarizeMessages() {
|
|
throw new Error('Subclasses attempted to call summarizeMessages without implementing it');
|
|
}
|
|
|
|
/**
|
|
* @returns {string}
|
|
*/
|
|
getResponseModel() {
|
|
if (isAgentsEndpoint(this.options.endpoint) && this.options.agent && this.options.agent.id) {
|
|
return this.options.agent.id;
|
|
}
|
|
|
|
return this.modelOptions?.model ?? this.model;
|
|
}
|
|
|
|
/**
|
|
* Abstract method to get the token count for a message. Subclasses must implement this method.
|
|
* @param {TMessage} responseMessage
|
|
* @returns {number}
|
|
*/
|
|
getTokenCountForResponse(responseMessage) {
|
|
logger.debug('[BaseClient] `recordTokenUsage` not implemented.', {
|
|
messageId: responseMessage?.messageId,
|
|
});
|
|
}
|
|
|
|
/**
|
|
* Abstract method to record token usage. Subclasses must implement this method.
|
|
* If a correction to the token usage is needed, the method should return an object with the corrected token counts.
|
|
* Should only be used if `recordCollectedUsage` was not used instead.
|
|
* @param {string} [model]
|
|
* @param {AppConfig['balance']} [balance]
|
|
* @param {number} promptTokens
|
|
* @param {number} completionTokens
|
|
* @param {string} [messageId]
|
|
* @returns {Promise<void>}
|
|
*/
|
|
async recordTokenUsage({
|
|
model,
|
|
balance,
|
|
messageId,
|
|
transactions,
|
|
promptTokens,
|
|
completionTokens,
|
|
}) {
|
|
logger.debug('[BaseClient] `recordTokenUsage` not implemented.', {
|
|
model,
|
|
balance,
|
|
messageId,
|
|
transactions,
|
|
promptTokens,
|
|
completionTokens,
|
|
});
|
|
}
|
|
|
|
/**
|
|
* Makes an HTTP request and logs the process.
|
|
*
|
|
* @param {RequestInfo} url - The URL to make the request to. Can be a string or a Request object.
|
|
* @param {RequestInit} [init] - Optional init options for the request.
|
|
* @returns {Promise<Response>} - A promise that resolves to the response of the fetch request.
|
|
*/
|
|
async fetch(_url, init) {
|
|
let url = _url;
|
|
if (this.options.directEndpoint) {
|
|
url = this.options.reverseProxyUrl;
|
|
}
|
|
logger.debug(`Making request to ${url}`);
|
|
if (typeof Bun !== 'undefined') {
|
|
return await fetch(url, init);
|
|
}
|
|
return await fetch(url, init);
|
|
}
|
|
|
|
getBuildMessagesOptions() {
|
|
throw new Error('Subclasses must implement getBuildMessagesOptions');
|
|
}
|
|
|
|
async generateTextStream(text, onProgress, options = {}) {
|
|
const stream = new TextStream(text, options);
|
|
await stream.processTextStream(onProgress);
|
|
}
|
|
|
|
/**
|
|
* @returns {[string|undefined, string|undefined]}
|
|
*/
|
|
processOverideIds() {
|
|
/** @type {Record<string, string | undefined>} */
|
|
let { overrideConvoId, overrideUserMessageId } = this.options?.req?.body ?? {};
|
|
if (overrideConvoId) {
|
|
const [conversationId, index] = overrideConvoId.split(Constants.COMMON_DIVIDER);
|
|
overrideConvoId = conversationId;
|
|
if (index !== '0') {
|
|
this.skipSaveConvo = true;
|
|
}
|
|
}
|
|
if (overrideUserMessageId) {
|
|
const [userMessageId, index] = overrideUserMessageId.split(Constants.COMMON_DIVIDER);
|
|
overrideUserMessageId = userMessageId;
|
|
if (index !== '0') {
|
|
this.skipSaveUserMessage = true;
|
|
}
|
|
}
|
|
|
|
return [overrideConvoId, overrideUserMessageId];
|
|
}
|
|
|
|
async setMessageOptions(opts = {}) {
|
|
if (opts && opts.replaceOptions) {
|
|
this.setOptions(opts);
|
|
}
|
|
|
|
const [overrideConvoId, overrideUserMessageId] = this.processOverideIds();
|
|
const { isEdited, isContinued } = opts;
|
|
const user = opts.user ?? null;
|
|
this.user = user;
|
|
const saveOptions = this.getSaveOptions();
|
|
this.abortController = opts.abortController ?? new AbortController();
|
|
const requestConvoId = overrideConvoId ?? opts.conversationId;
|
|
const conversationId = requestConvoId ?? crypto.randomUUID();
|
|
const parentMessageId = opts.parentMessageId ?? Constants.NO_PARENT;
|
|
const userMessageId =
|
|
overrideUserMessageId ?? opts.overrideParentMessageId ?? crypto.randomUUID();
|
|
let responseMessageId = opts.responseMessageId ?? crypto.randomUUID();
|
|
let head = isEdited ? responseMessageId : parentMessageId;
|
|
this.currentMessages = (await this.loadHistory(conversationId, head)) ?? [];
|
|
this.conversationId = conversationId;
|
|
|
|
if (isEdited && !isContinued) {
|
|
responseMessageId = crypto.randomUUID();
|
|
head = responseMessageId;
|
|
this.currentMessages[this.currentMessages.length - 1].messageId = head;
|
|
}
|
|
|
|
if (opts.isRegenerate && responseMessageId.endsWith('_')) {
|
|
responseMessageId = crypto.randomUUID();
|
|
}
|
|
|
|
this.responseMessageId = responseMessageId;
|
|
|
|
return {
|
|
...opts,
|
|
user,
|
|
head,
|
|
saveOptions,
|
|
userMessageId,
|
|
requestConvoId,
|
|
conversationId,
|
|
parentMessageId,
|
|
responseMessageId,
|
|
};
|
|
}
|
|
|
|
createUserMessage({ messageId, parentMessageId, conversationId, text }) {
|
|
return {
|
|
messageId,
|
|
parentMessageId,
|
|
conversationId,
|
|
sender: 'User',
|
|
text,
|
|
isCreatedByUser: true,
|
|
};
|
|
}
|
|
|
|
async handleStartMethods(message, opts) {
|
|
const {
|
|
user,
|
|
head,
|
|
saveOptions,
|
|
userMessageId,
|
|
requestConvoId,
|
|
conversationId,
|
|
parentMessageId,
|
|
responseMessageId,
|
|
} = await this.setMessageOptions(opts);
|
|
this.options.startupTelemetry?.mark('history_loaded');
|
|
|
|
const userMessage = opts.isEdited
|
|
? this.currentMessages[this.currentMessages.length - 2]
|
|
: this.createUserMessage({
|
|
messageId: userMessageId,
|
|
parentMessageId,
|
|
conversationId,
|
|
text: message,
|
|
});
|
|
|
|
/**
|
|
* Attach quoted excerpts (the "Add to chat" selections from `req.body.quotes`)
|
|
* before `getReqData`/`onStart` fire, so the optimistic bubble, resumable job
|
|
* metadata, and the saved row all carry them. Only on fresh turns — edits
|
|
* replay an existing message that already has its quotes. The excerpts are
|
|
* merged into the model-facing text later, per message, in `buildMessages`,
|
|
* keeping the stored `text` clean while the count stays consistent.
|
|
*/
|
|
if (!opts.isEdited) {
|
|
const referencedQuotes = getReferencedQuotes(this.options.req?.body?.quotes);
|
|
if (referencedQuotes != null) {
|
|
userMessage.quotes = referencedQuotes;
|
|
}
|
|
}
|
|
|
|
if (typeof opts?.getReqData === 'function') {
|
|
opts.getReqData({
|
|
userMessage,
|
|
conversationId,
|
|
responseMessageId,
|
|
sender: this.sender,
|
|
});
|
|
}
|
|
|
|
if (typeof opts?.onStart === 'function') {
|
|
const isNewConvo = !requestConvoId && parentMessageId === Constants.NO_PARENT;
|
|
opts.onStart(userMessage, responseMessageId, isNewConvo);
|
|
}
|
|
|
|
return {
|
|
...opts,
|
|
user,
|
|
head,
|
|
conversationId,
|
|
responseMessageId,
|
|
saveOptions,
|
|
userMessage,
|
|
};
|
|
}
|
|
|
|
/**
|
|
* Adds instructions to the messages array. If the instructions object is empty or undefined,
|
|
* the original messages array is returned. Otherwise, the instructions are added to the messages
|
|
* array either at the beginning (default) or preserving the last message at the end.
|
|
*
|
|
* @param {Array} messages - An array of messages.
|
|
* @param {Object} instructions - An object containing instructions to be added to the messages.
|
|
* @param {boolean} [beforeLast=false] - If true, adds instructions before the last message; if false, adds at the beginning.
|
|
* @returns {Array} An array containing messages and instructions, or the original messages if instructions are empty.
|
|
*/
|
|
addInstructions(messages, instructions, beforeLast = false) {
|
|
if (!instructions || Object.keys(instructions).length === 0) {
|
|
return messages;
|
|
}
|
|
|
|
if (!beforeLast) {
|
|
return [instructions, ...messages];
|
|
}
|
|
|
|
// Legacy behavior: add instructions before the last message
|
|
const payload = [];
|
|
if (messages.length > 1) {
|
|
payload.push(...messages.slice(0, -1));
|
|
}
|
|
|
|
payload.push(instructions);
|
|
|
|
if (messages.length > 0) {
|
|
payload.push(messages[messages.length - 1]);
|
|
}
|
|
|
|
return payload;
|
|
}
|
|
|
|
concatenateMessages(messages) {
|
|
return messages.reduce((acc, message) => {
|
|
const nameOrRole = message.name ?? message.role;
|
|
return acc + `${nameOrRole}:\n${message.content}\n\n`;
|
|
}, '');
|
|
}
|
|
|
|
/**
|
|
* This method processes an array of messages and returns a context of messages that fit within a specified token limit.
|
|
* It iterates over the messages from newest to oldest, adding them to the context until the token limit is reached.
|
|
* If the token limit would be exceeded by adding a message, that message is not added to the context and remains in the original array.
|
|
* The method uses `push` and `pop` operations for efficient array manipulation, and reverses the context array at the end to maintain the original order of the messages.
|
|
*
|
|
* @param {Object} params
|
|
* @param {TMessage[]} params.messages - An array of messages, each with a `tokenCount` property. The messages should be ordered from oldest to newest.
|
|
* @param {number} [params.maxContextTokens] - The max number of tokens allowed in the context. If not provided, defaults to `this.maxContextTokens`.
|
|
* @param {{ role: 'system', content: text, tokenCount: number }} [params.instructions] - Instructions already added to the context at index 0.
|
|
* @returns {Promise<{
|
|
* context: TMessage[],
|
|
* remainingContextTokens: number,
|
|
* messagesToRefine: TMessage[],
|
|
* }>} An object with three properties: `context`, `remainingContextTokens`, and `messagesToRefine`.
|
|
* `context` is an array of messages that fit within the token limit.
|
|
* `remainingContextTokens` is the number of tokens remaining within the limit after adding the messages to the context.
|
|
* `messagesToRefine` is an array of messages that were not added to the context because they would have exceeded the token limit.
|
|
*/
|
|
async getMessagesWithinTokenLimit({ messages: _messages, maxContextTokens, instructions }) {
|
|
// Every reply is primed with <|start|>assistant<|message|>, so we
|
|
// start with 3 tokens for the label after all messages have been counted.
|
|
let currentTokenCount = 3;
|
|
const instructionsTokenCount = instructions?.tokenCount ?? 0;
|
|
let remainingContextTokens =
|
|
(maxContextTokens ?? this.maxContextTokens) - instructionsTokenCount;
|
|
const messages = [..._messages];
|
|
|
|
const context = [];
|
|
|
|
if (currentTokenCount < remainingContextTokens) {
|
|
while (messages.length > 0 && currentTokenCount < remainingContextTokens) {
|
|
if (messages.length === 1 && instructions) {
|
|
break;
|
|
}
|
|
const poppedMessage = messages.pop();
|
|
const { tokenCount } = poppedMessage;
|
|
|
|
if (poppedMessage && currentTokenCount + tokenCount <= remainingContextTokens) {
|
|
context.push(poppedMessage);
|
|
currentTokenCount += tokenCount;
|
|
} else {
|
|
messages.push(poppedMessage);
|
|
break;
|
|
}
|
|
}
|
|
}
|
|
|
|
if (instructions) {
|
|
context.push(_messages[0]);
|
|
messages.shift();
|
|
}
|
|
|
|
const prunedMemory = messages;
|
|
remainingContextTokens -= currentTokenCount;
|
|
|
|
return {
|
|
context: context.reverse(),
|
|
remainingContextTokens,
|
|
messagesToRefine: prunedMemory,
|
|
};
|
|
}
|
|
|
|
async sendMessage(message, opts = {}) {
|
|
const appConfig = this.options.req?.config;
|
|
/** @type {Promise<TMessage>} */
|
|
let userMessagePromise;
|
|
const { user, head, isEdited, conversationId, responseMessageId, saveOptions, userMessage } =
|
|
await this.handleStartMethods(message, opts);
|
|
|
|
if (opts.progressCallback) {
|
|
opts.onProgress = opts.progressCallback.call(null, {
|
|
...(opts.progressOptions ?? {}),
|
|
parentMessageId: userMessage.messageId,
|
|
messageId: responseMessageId,
|
|
});
|
|
}
|
|
|
|
const { editedContent } = opts;
|
|
|
|
// It's not necessary to push to currentMessages
|
|
// depending on subclass implementation of handling messages
|
|
// When this is an edit, all messages are already in currentMessages, both user and response
|
|
if (isEdited) {
|
|
let latestMessage = this.currentMessages[this.currentMessages.length - 1];
|
|
if (!latestMessage) {
|
|
latestMessage = {
|
|
messageId: responseMessageId,
|
|
conversationId,
|
|
parentMessageId: userMessage.messageId,
|
|
isCreatedByUser: false,
|
|
model: this.modelOptions?.model ?? this.model,
|
|
sender: this.sender,
|
|
};
|
|
this.currentMessages.push(userMessage, latestMessage);
|
|
} else if (editedContent != null) {
|
|
// Handle editedContent for content parts
|
|
if (editedContent && latestMessage.content && Array.isArray(latestMessage.content)) {
|
|
const { index, type } = editedContent;
|
|
const text = editedContent[type];
|
|
if (index >= 0 && index < latestMessage.content.length) {
|
|
const contentPart = latestMessage.content[index];
|
|
if (type === ContentTypes.THINK && contentPart.type === ContentTypes.THINK) {
|
|
contentPart[ContentTypes.THINK] = text;
|
|
} else if (type === ContentTypes.TEXT && contentPart.type === ContentTypes.TEXT) {
|
|
contentPart[ContentTypes.TEXT] = text;
|
|
}
|
|
}
|
|
}
|
|
}
|
|
this.continued = true;
|
|
} else {
|
|
this.currentMessages.push(userMessage);
|
|
}
|
|
|
|
/**
|
|
* When the userMessage is pushed to currentMessages, the parentMessage is the userMessageId.
|
|
* this only matters when buildMessages is utilizing the parentMessageId, and may vary on implementation
|
|
*/
|
|
const parentMessageId = isEdited ? head : userMessage.messageId;
|
|
this.parentMessageId = parentMessageId;
|
|
let {
|
|
prompt: payload,
|
|
tokenCountMap,
|
|
promptTokens,
|
|
} = await this.buildMessages(
|
|
this.currentMessages,
|
|
parentMessageId,
|
|
this.getBuildMessagesOptions(opts),
|
|
opts,
|
|
);
|
|
this.options.startupTelemetry?.mark('messages_built');
|
|
|
|
if (tokenCountMap && tokenCountMap[userMessage.messageId]) {
|
|
userMessage.tokenCount = tokenCountMap[userMessage.messageId];
|
|
logger.debug('[BaseClient] userMessage', {
|
|
messageId: userMessage.messageId,
|
|
tokenCount: userMessage.tokenCount,
|
|
conversationId: userMessage.conversationId,
|
|
});
|
|
}
|
|
|
|
if (!isEdited && !this.skipSaveUserMessage) {
|
|
const reqFiles = this.options.req?.body?.files;
|
|
if (reqFiles && Array.isArray(this.options.attachments)) {
|
|
const files = buildMessageFiles(reqFiles, this.options.attachments);
|
|
if (files.length > 0) {
|
|
userMessage.files = files;
|
|
}
|
|
delete userMessage.image_urls;
|
|
}
|
|
/**
|
|
* Persist the user's manual skill picks onto the user message so the
|
|
* frontend `SkillPills` component can render them in history
|
|
* after reload. UI-only metadata — the runtime skill resolution
|
|
* pipeline reads the top-level `req.body.manualSkills` separately.
|
|
* Filter is defense-in-depth on top of Mongoose schema validation:
|
|
* keeps the DB row free of empty/non-string entries even if a
|
|
* crafted payload slips past schema checks upstream.
|
|
*/
|
|
const rawManualSkills = this.options.req?.body?.manualSkills;
|
|
if (Array.isArray(rawManualSkills) && rawManualSkills.length > 0) {
|
|
const skills = rawManualSkills.filter((s) => typeof s === 'string' && s.length > 0);
|
|
if (skills.length > 0) {
|
|
userMessage.manualSkills = skills;
|
|
}
|
|
}
|
|
/**
|
|
* Persist the names of skills auto-primed this turn via `always-apply`
|
|
* frontmatter so `SkillPills` can render pinned-variant badges
|
|
* on the user bubble that survive reload and history render. Frozen
|
|
* at turn time (not reconstructed from `Skill.alwaysApply` at render
|
|
* time) because the flag is mutable — historical turns must keep
|
|
* their audit trail even if an admin flips `alwaysApply` off later.
|
|
*/
|
|
const alwaysApplySkillPrimes = this.options.agent?.alwaysApplySkillPrimes;
|
|
if (Array.isArray(alwaysApplySkillPrimes) && alwaysApplySkillPrimes.length > 0) {
|
|
const names = alwaysApplySkillPrimes
|
|
.map((p) => p?.name)
|
|
.filter((n) => typeof n === 'string' && n.length > 0);
|
|
if (names.length > 0) {
|
|
userMessage.alwaysAppliedSkills = names;
|
|
}
|
|
}
|
|
userMessagePromise = this.saveMessageToDatabase(userMessage, saveOptions, user).catch(
|
|
(err) => {
|
|
logger.error('[BaseClient] Failed to save user message:', err);
|
|
return {};
|
|
},
|
|
);
|
|
this.savedMessageIds.add(userMessage.messageId);
|
|
if (typeof opts?.getReqData === 'function') {
|
|
opts.getReqData({
|
|
userMessagePromise,
|
|
});
|
|
}
|
|
}
|
|
|
|
const balanceConfig = getBalanceConfig(appConfig);
|
|
const transactionsConfig = getTransactionsConfig(appConfig);
|
|
if (
|
|
balanceConfig?.enabled &&
|
|
supportsBalanceCheck[this.options.endpointType ?? this.options.endpoint]
|
|
) {
|
|
await checkBalance(
|
|
{
|
|
req: this.options.req,
|
|
res: this.options.res,
|
|
txData: {
|
|
user: this.user,
|
|
tokenType: 'prompt',
|
|
amount: promptTokens,
|
|
endpoint: this.options.endpoint,
|
|
model: this.modelOptions?.model ?? this.model,
|
|
endpointTokenConfig: this.options.endpointTokenConfig,
|
|
},
|
|
},
|
|
{
|
|
logViolation,
|
|
getMultiplier: db.getMultiplier,
|
|
findBalanceByUser: db.findBalanceByUser,
|
|
createAutoRefillTransaction: db.createAutoRefillTransaction,
|
|
balanceConfig,
|
|
upsertBalanceFields: db.upsertBalanceFields,
|
|
},
|
|
);
|
|
}
|
|
|
|
const { completion, metadata } = await this.sendCompletion(payload, opts);
|
|
if (this.abortController) {
|
|
this.abortController.requestCompleted = true;
|
|
}
|
|
|
|
const isAgentResponse = isAgentsEndpoint(this.options.endpoint);
|
|
const langfuseTraceId = isAgentResponse ? traceIdForMessage(responseMessageId) : undefined;
|
|
const langfuseSampled =
|
|
langfuseTraceId != null ? isLangfuseTraceSampled(langfuseTraceId) : undefined;
|
|
|
|
/** @type {TMessage} */
|
|
const responseMessage = {
|
|
messageId: responseMessageId,
|
|
conversationId,
|
|
parentMessageId: userMessage.messageId,
|
|
isCreatedByUser: false,
|
|
...(isAgentResponse && {
|
|
langfuseSampled,
|
|
langfuseDestinationIds: await getLangfuseTraceDestinationIds(
|
|
appConfig,
|
|
langfuseTraceId,
|
|
langfuseSampled,
|
|
),
|
|
}),
|
|
isEdited,
|
|
model: this.getResponseModel(),
|
|
sender: this.sender,
|
|
promptTokens,
|
|
iconURL: this.options.iconURL,
|
|
endpoint: this.options.endpoint,
|
|
...(this.metadata ?? {}),
|
|
metadata: Object.keys(metadata ?? {}).length > 0 ? metadata : undefined,
|
|
};
|
|
|
|
if (typeof completion === 'string') {
|
|
responseMessage.text = completion;
|
|
} else if (
|
|
Array.isArray(completion) &&
|
|
(this.clientName === EModelEndpoint.agents ||
|
|
isParamEndpoint(this.options.endpoint, this.options.endpointType))
|
|
) {
|
|
responseMessage.text = '';
|
|
|
|
if (!opts.editedContent || this.currentMessages.length === 0) {
|
|
responseMessage.content = completion;
|
|
} else {
|
|
const latestMessage = this.currentMessages[this.currentMessages.length - 1];
|
|
if (!latestMessage?.content) {
|
|
responseMessage.content = completion;
|
|
} else {
|
|
const existingContent = [...latestMessage.content];
|
|
const { type: editedType } = opts.editedContent;
|
|
responseMessage.content = this.mergeEditedContent(
|
|
existingContent,
|
|
completion,
|
|
editedType,
|
|
);
|
|
}
|
|
}
|
|
} else if (Array.isArray(completion)) {
|
|
responseMessage.text = completion.join('');
|
|
}
|
|
|
|
if (tokenCountMap && this.recordTokenUsage && this.getTokenCountForResponse) {
|
|
let completionTokens;
|
|
|
|
/**
|
|
* Metadata about input/output costs for the current message. The client
|
|
* should provide a function to get the current stream usage metadata; if not,
|
|
* use the legacy token estimations.
|
|
* @type {StreamUsage | null} */
|
|
const usage = this.getStreamUsage != null ? this.getStreamUsage() : null;
|
|
|
|
if (usage != null && Number(usage[this.outputTokensKey]) > 0) {
|
|
responseMessage.tokenCount = usage[this.outputTokensKey];
|
|
completionTokens = responseMessage.tokenCount;
|
|
} else {
|
|
responseMessage.tokenCount = this.getTokenCountForResponse(responseMessage);
|
|
completionTokens = responseMessage.tokenCount;
|
|
await this.recordTokenUsage({
|
|
usage,
|
|
promptTokens,
|
|
completionTokens,
|
|
balance: balanceConfig,
|
|
transactions: transactionsConfig,
|
|
/** Note: When using agents, responseMessage.model is the agent ID, not the model */
|
|
model: this.model,
|
|
messageId: this.responseMessageId,
|
|
});
|
|
}
|
|
|
|
logger.debug('[BaseClient] Response token usage', {
|
|
messageId: responseMessage.messageId,
|
|
model: responseMessage.model,
|
|
promptTokens,
|
|
completionTokens,
|
|
});
|
|
}
|
|
|
|
if (userMessagePromise) {
|
|
await userMessagePromise;
|
|
}
|
|
|
|
if (
|
|
this.contextMeta?.calibrationRatio > 0 &&
|
|
this.contextMeta.calibrationRatio !== 1 &&
|
|
userMessage.tokenCount > 0
|
|
) {
|
|
const calibrated = Math.round(userMessage.tokenCount * this.contextMeta.calibrationRatio);
|
|
if (calibrated !== userMessage.tokenCount) {
|
|
logger.debug('[BaseClient] Calibrated user message tokenCount', {
|
|
messageId: userMessage.messageId,
|
|
raw: userMessage.tokenCount,
|
|
calibrated,
|
|
ratio: this.contextMeta.calibrationRatio,
|
|
});
|
|
userMessage.tokenCount = calibrated;
|
|
await this.updateMessageInDatabase({
|
|
messageId: userMessage.messageId,
|
|
tokenCount: calibrated,
|
|
});
|
|
}
|
|
}
|
|
|
|
if (this.artifactPromises) {
|
|
responseMessage.attachments = (await Promise.all(this.artifactPromises)).filter((a) => a);
|
|
}
|
|
|
|
if (this.options.attachments) {
|
|
try {
|
|
saveOptions.files = this.options.attachments.map((attachments) => attachments.file_id);
|
|
} catch (error) {
|
|
logger.error('[BaseClient] Error mapping attachments for conversation', error);
|
|
}
|
|
}
|
|
|
|
if (this.contextMeta) {
|
|
responseMessage.contextMeta = this.contextMeta;
|
|
}
|
|
|
|
/** Resumable generation controllers must win the generation's terminal
|
|
* CAS before this outcome-defining `unfinished:false` write can begin.
|
|
* The hook is deliberately narrow: ordinary clients omit it, and `false`
|
|
* means another terminal owner (for example Stop) already won, so this
|
|
* stale completion must return without writing the response row. */
|
|
if (typeof opts.beforeResponsePersistence === 'function') {
|
|
const ownsTerminalPersistence = await opts.beforeResponsePersistence(responseMessage);
|
|
if (ownsTerminalPersistence === false) {
|
|
responseMessage.databasePromise = Promise.resolve({ persistenceSkipped: true });
|
|
return responseMessage;
|
|
}
|
|
}
|
|
|
|
responseMessage.databasePromise = this.saveMessageToDatabase(
|
|
responseMessage,
|
|
saveOptions,
|
|
user,
|
|
);
|
|
this.savedMessageIds.add(responseMessage.messageId);
|
|
return responseMessage;
|
|
}
|
|
|
|
async loadHistory(conversationId, parentMessageId = null) {
|
|
logger.debug('[BaseClient] Loading history:', { conversationId, parentMessageId });
|
|
|
|
const messages = (await db.getMessages({ conversationId, user: this.user })) ?? [];
|
|
|
|
if (messages.length === 0) {
|
|
return [];
|
|
}
|
|
|
|
let mapMethod = null;
|
|
if (this.getMessageMapMethod) {
|
|
mapMethod = this.getMessageMapMethod();
|
|
}
|
|
|
|
let _messages = this.constructor.getMessagesForConversation({
|
|
messages,
|
|
parentMessageId,
|
|
mapMethod,
|
|
});
|
|
|
|
_messages = await this.addPreviousAttachments(_messages);
|
|
|
|
if (!this.shouldSummarize) {
|
|
return _messages;
|
|
}
|
|
|
|
for (let i = _messages.length - 1; i >= 0; i--) {
|
|
const msg = _messages[i];
|
|
if (!msg) {
|
|
continue;
|
|
}
|
|
|
|
const summaryBlock = BaseClient.findSummaryContentBlock(msg);
|
|
if (summaryBlock) {
|
|
this.previous_summary = {
|
|
...msg,
|
|
summary: BaseClient.getSummaryText(summaryBlock),
|
|
summaryTokenCount: summaryBlock.tokenCount,
|
|
};
|
|
break;
|
|
}
|
|
|
|
if (msg.summary) {
|
|
this.previous_summary = msg;
|
|
break;
|
|
}
|
|
}
|
|
|
|
if (this.previous_summary) {
|
|
const { messageId, summary, tokenCount, summaryTokenCount } = this.previous_summary;
|
|
logger.debug('[BaseClient] Previous summary:', {
|
|
messageId,
|
|
summary,
|
|
tokenCount,
|
|
summaryTokenCount,
|
|
});
|
|
}
|
|
|
|
return _messages;
|
|
}
|
|
|
|
/**
|
|
* Save a message to the database.
|
|
* @param {TMessage} message
|
|
* @param {Partial<TConversation>} endpointOptions
|
|
* @param {string | null} user
|
|
*/
|
|
async saveMessageToDatabase(message, endpointOptions, user = null) {
|
|
// Snapshot options before any await; disposeClient may set client.options = null
|
|
// while this method is suspended at an I/O boundary, but the local reference
|
|
// remains valid (disposeClient nulls the property, not the object itself).
|
|
const options = this.options;
|
|
if (!options) {
|
|
logger.error('[BaseClient] saveMessageToDatabase: client disposed before save, skipping');
|
|
return {};
|
|
}
|
|
|
|
if (this.user && user !== this.user) {
|
|
throw new Error('User mismatch.');
|
|
}
|
|
|
|
const hasAddedConvo = options?.req?.body?.addedConvo != null;
|
|
const reqCtx = {
|
|
userId: options?.req?.user?.id,
|
|
isTemporary: options?.req?.body?.isTemporary,
|
|
interfaceConfig: options?.req?.config?.interfaceConfig,
|
|
};
|
|
const savedMessage = await db.saveMessage(
|
|
reqCtx,
|
|
{
|
|
...message,
|
|
endpoint: options.endpoint,
|
|
unfinished: false,
|
|
user,
|
|
...(hasAddedConvo && { addedConvo: true }),
|
|
},
|
|
{ context: 'api/app/clients/BaseClient.js - saveMessageToDatabase #saveMessage' },
|
|
);
|
|
|
|
if (this.skipSaveConvo) {
|
|
return { message: savedMessage };
|
|
}
|
|
|
|
const fieldsToKeep = {
|
|
conversationId: message.conversationId,
|
|
endpoint: options.endpoint,
|
|
endpointType: options.endpointType,
|
|
...endpointOptions,
|
|
};
|
|
const conversationCreatedAt = options?.req?.conversationCreatedAt;
|
|
const createdAtOnInsert =
|
|
conversationCreatedAt != null ? new Date(conversationCreatedAt) : undefined;
|
|
const validCreatedAtOnInsert =
|
|
createdAtOnInsert && !Number.isNaN(createdAtOnInsert.getTime())
|
|
? createdAtOnInsert
|
|
: undefined;
|
|
|
|
const req = options?.req;
|
|
const skippedExistingConvoLookup = this.fetchedConvo === true;
|
|
const hasResolvedConversation =
|
|
req != null && Object.prototype.hasOwnProperty.call(req, 'resolvedConversation');
|
|
let existingConvo = null;
|
|
if (!skippedExistingConvoLookup && hasResolvedConversation) {
|
|
existingConvo = req.resolvedConversation;
|
|
} else if (!skippedExistingConvoLookup) {
|
|
existingConvo = await db.getConvo(req?.user?.id, message.conversationId);
|
|
}
|
|
if (hasResolvedConversation) {
|
|
delete req.resolvedConversation;
|
|
}
|
|
const shouldSetCreatedAtOnInsert = !skippedExistingConvoLookup && existingConvo == null;
|
|
|
|
const unsetFields = {};
|
|
const exceptions = new Set(['spec', 'iconURL']);
|
|
const hasNonEphemeralAgent =
|
|
isAgentsEndpoint(options.endpoint) &&
|
|
endpointOptions?.agent_id &&
|
|
!isEphemeralAgentId(endpointOptions.agent_id);
|
|
if (hasNonEphemeralAgent) {
|
|
exceptions.add('model');
|
|
}
|
|
if (existingConvo != null) {
|
|
this.fetchedConvo = true;
|
|
for (const key in existingConvo) {
|
|
if (!key) {
|
|
continue;
|
|
}
|
|
if (excludedKeys.has(key) && !exceptions.has(key)) {
|
|
continue;
|
|
}
|
|
|
|
if (endpointOptions?.[key] === undefined) {
|
|
unsetFields[key] = 1;
|
|
}
|
|
}
|
|
}
|
|
|
|
const conversation = await db.saveConvo(reqCtx, fieldsToKeep, {
|
|
context: 'api/app/clients/BaseClient.js - saveMessageToDatabase #saveConvo',
|
|
unsetFields,
|
|
createdAtOnInsert: shouldSetCreatedAtOnInsert ? validCreatedAtOnInsert : undefined,
|
|
});
|
|
|
|
return { message: savedMessage, conversation };
|
|
}
|
|
|
|
/**
|
|
* Update a message in the database.
|
|
* @param {Partial<TMessage>} message
|
|
*/
|
|
async updateMessageInDatabase(message) {
|
|
await db.updateMessage(this.options?.req?.user?.id, message);
|
|
}
|
|
|
|
/** Extracts text from a summary block (handles both legacy `text` field and new `content` array format). */
|
|
static getSummaryText(summaryBlock) {
|
|
if (Array.isArray(summaryBlock.content)) {
|
|
return summaryBlock.content.map((b) => b.text ?? '').join('');
|
|
}
|
|
if (typeof summaryBlock.content === 'string') {
|
|
return summaryBlock.content;
|
|
}
|
|
return summaryBlock.text ?? '';
|
|
}
|
|
|
|
/** Finds the last summary content block in a message's content array (last-summary-wins). */
|
|
static findSummaryContentBlock(message) {
|
|
if (!Array.isArray(message?.content)) {
|
|
return null;
|
|
}
|
|
let lastSummary = null;
|
|
for (const part of message.content) {
|
|
if (
|
|
part?.type === ContentTypes.SUMMARY &&
|
|
BaseClient.getSummaryText(part).trim().length > 0
|
|
) {
|
|
lastSummary = part;
|
|
}
|
|
}
|
|
return lastSummary;
|
|
}
|
|
|
|
/**
|
|
* Iterate through messages, building an array based on the parentMessageId.
|
|
*
|
|
* This function constructs a conversation thread by traversing messages from a given parentMessageId up to the root message.
|
|
* It handles cyclic references by ensuring that a message is not processed more than once.
|
|
* If the 'summary' option is set to true and a message has a 'summary' property:
|
|
* - The message's 'role' is set to 'system'.
|
|
* - The message's 'text' is set to its 'summary'.
|
|
* - If the message has a 'summaryTokenCount', the message's 'tokenCount' is set to 'summaryTokenCount'.
|
|
* The traversal stops at the message with the 'summary' property.
|
|
*
|
|
* Each message object should have an 'id' or 'messageId' property and may have a 'parentMessageId' property.
|
|
* The 'parentMessageId' is the ID of the message that the current message is a reply to.
|
|
* If 'parentMessageId' is not present, null, or is Constants.NO_PARENT,
|
|
* the message is considered a root message.
|
|
*
|
|
* @param {Object} options - The options for the function.
|
|
* @param {TMessage[]} options.messages - An array of message objects. Each object should have either an 'id' or 'messageId' property, and may have a 'parentMessageId' property.
|
|
* @param {string} options.parentMessageId - The ID of the parent message to start the traversal from.
|
|
* @param {Function} [options.mapMethod] - An optional function to map over the ordered messages. Applied conditionally based on mapCondition.
|
|
* @param {(message: TMessage) => boolean} [options.mapCondition] - An optional function to determine whether mapMethod should be applied to a given message. If not provided and mapMethod is set, mapMethod applies to all messages.
|
|
* @param {boolean} [options.summary=false] - If set to true, the traversal modifies messages with 'summary' and 'summaryTokenCount' properties and stops at the message with a 'summary' property.
|
|
* @returns {TMessage[]} An array containing the messages in the order they should be displayed, starting with the most recent message with a 'summary' property if the 'summary' option is true, and ending with the message identified by 'parentMessageId'.
|
|
*/
|
|
static getMessagesForConversation({
|
|
messages,
|
|
parentMessageId,
|
|
mapMethod = null,
|
|
mapCondition = null,
|
|
summary = false,
|
|
}) {
|
|
if (!messages || messages.length === 0) {
|
|
return [];
|
|
}
|
|
|
|
const orderedMessages = [];
|
|
let currentMessageId = parentMessageId;
|
|
const visitedMessageIds = new Set();
|
|
|
|
while (currentMessageId) {
|
|
if (visitedMessageIds.has(currentMessageId)) {
|
|
break;
|
|
}
|
|
const message = messages.find((msg) => {
|
|
const messageId = msg.messageId ?? msg.id;
|
|
return messageId === currentMessageId;
|
|
});
|
|
|
|
visitedMessageIds.add(currentMessageId);
|
|
|
|
if (!message) {
|
|
break;
|
|
}
|
|
|
|
let resolved = message;
|
|
let hasSummary = false;
|
|
if (summary) {
|
|
const summaryBlock = BaseClient.findSummaryContentBlock(message);
|
|
if (summaryBlock) {
|
|
const summaryText = BaseClient.getSummaryText(summaryBlock);
|
|
resolved = {
|
|
...message,
|
|
role: 'system',
|
|
content: [{ type: ContentTypes.TEXT, text: summaryText }],
|
|
tokenCount: summaryBlock.tokenCount,
|
|
};
|
|
hasSummary = true;
|
|
} else if (message.summary) {
|
|
resolved = {
|
|
...message,
|
|
role: 'system',
|
|
content: [{ type: ContentTypes.TEXT, text: message.summary }],
|
|
tokenCount: message.summaryTokenCount ?? message.tokenCount,
|
|
};
|
|
hasSummary = true;
|
|
}
|
|
}
|
|
|
|
const shouldMap = mapMethod != null && (mapCondition != null ? mapCondition(resolved) : true);
|
|
const processedMessage = shouldMap ? mapMethod(resolved) : resolved;
|
|
orderedMessages.push(processedMessage);
|
|
|
|
if (hasSummary) {
|
|
break;
|
|
}
|
|
|
|
currentMessageId =
|
|
message.parentMessageId === Constants.NO_PARENT ? null : message.parentMessageId;
|
|
}
|
|
|
|
orderedMessages.reverse();
|
|
return orderedMessages;
|
|
}
|
|
|
|
/**
|
|
* Algorithm adapted from "6. Counting tokens for chat API calls" of
|
|
* https://github.com/openai/openai-cookbook/blob/main/examples/How_to_count_tokens_with_tiktoken.ipynb
|
|
*
|
|
* An additional 3 tokens need to be added for assistant label priming after all messages have been counted.
|
|
* In our implementation, this is accounted for in the getMessagesWithinTokenLimit method.
|
|
*
|
|
* The content parts example was adapted from the following example:
|
|
* https://github.com/openai/openai-cookbook/pull/881/files
|
|
*
|
|
* Note: image token calculation is to be done elsewhere where we have access to the image metadata
|
|
*
|
|
* @param {Object} message
|
|
*/
|
|
getTokenCountForMessage(message) {
|
|
// Note: gpt-3.5-turbo and gpt-4 may update over time. Use default for these as well as for unknown models
|
|
let tokensPerMessage = 3;
|
|
let tokensPerName = 1;
|
|
const model = this.modelOptions?.model ?? this.model;
|
|
|
|
if (model === 'gpt-3.5-turbo-0301') {
|
|
tokensPerMessage = 4;
|
|
tokensPerName = -1;
|
|
}
|
|
|
|
const processValue = (value) => {
|
|
if (Array.isArray(value)) {
|
|
for (let item of value) {
|
|
if (
|
|
!item ||
|
|
!item.type ||
|
|
item.type === ContentTypes.THINK ||
|
|
item.type === ContentTypes.ERROR ||
|
|
// UI-only progress headers — never model input, never billed output
|
|
item.type === ContentTypes.ACTIVITY_LABEL ||
|
|
item.type === ContentTypes.IMAGE_URL
|
|
) {
|
|
continue;
|
|
}
|
|
|
|
if (item.type === ContentTypes.TOOL_CALL && item.tool_call != null) {
|
|
const toolName = item.tool_call?.name || '';
|
|
if (toolName != null && toolName && typeof toolName === 'string') {
|
|
numTokens += this.getTokenCount(toolName);
|
|
}
|
|
|
|
const args = item.tool_call?.args || '';
|
|
if (args != null && args && typeof args === 'string') {
|
|
numTokens += this.getTokenCount(args);
|
|
}
|
|
|
|
const output = item.tool_call?.output || '';
|
|
if (output != null && output && typeof output === 'string') {
|
|
numTokens += this.getTokenCount(output);
|
|
}
|
|
continue;
|
|
}
|
|
|
|
const nestedValue = item[item.type];
|
|
|
|
if (!nestedValue) {
|
|
continue;
|
|
}
|
|
|
|
processValue(nestedValue);
|
|
}
|
|
} else if (typeof value === 'string') {
|
|
numTokens += this.getTokenCount(value);
|
|
} else if (typeof value === 'number') {
|
|
numTokens += this.getTokenCount(value.toString());
|
|
} else if (typeof value === 'boolean') {
|
|
numTokens += this.getTokenCount(value.toString());
|
|
}
|
|
};
|
|
|
|
let numTokens = tokensPerMessage;
|
|
for (let [key, value] of Object.entries(message)) {
|
|
processValue(value);
|
|
|
|
if (key === 'name') {
|
|
numTokens += tokensPerName;
|
|
}
|
|
}
|
|
return numTokens;
|
|
}
|
|
|
|
/**
|
|
* Merges completion content with existing content when editing TEXT or THINK types
|
|
* @param {Array} existingContent - The existing content array
|
|
* @param {Array} newCompletion - The new completion content
|
|
* @param {string} editedType - The type of content being edited
|
|
* @returns {Array} The merged content array
|
|
*/
|
|
mergeEditedContent(existingContent, newCompletion, editedType) {
|
|
if (!newCompletion.length) {
|
|
return existingContent.concat(newCompletion);
|
|
}
|
|
|
|
const lastIndex = existingContent.length - 1;
|
|
const lastExisting = existingContent[lastIndex];
|
|
const firstNew = newCompletion[0];
|
|
/** Phased and legacy/unphased text are distinct semantic streams. Merging
|
|
* either direction would stamp retained text with the wrong phase. */
|
|
const textPhaseCompatible =
|
|
editedType !== ContentTypes.TEXT ||
|
|
(lastExisting?.phase ?? null) === (firstNew?.phase ?? null);
|
|
const mergesFirstPart =
|
|
(editedType === ContentTypes.TEXT || editedType === ContentTypes.THINK) &&
|
|
lastExisting?.type === firstNew?.type &&
|
|
firstNew?.type === editedType &&
|
|
textPhaseCompatible;
|
|
/** Phase bounds are completion-local while the run streams. Persist them
|
|
* in the same absolute index space as the edited response assembled
|
|
* here. When the first new text/think part merges into the retained tail,
|
|
* every completion index shifts by prefixLength - 1; otherwise it shifts
|
|
* by the full retained prefix. */
|
|
const phaseIndexOffset = mergesFirstPart ? lastIndex : existingContent.length;
|
|
const adjustedCompletion = newCompletion.map((part) => {
|
|
if (
|
|
part?.type !== ContentTypes.ACTIVITY_LABEL ||
|
|
part.activity_label_type !== 'phase' ||
|
|
typeof part.activity_start_index !== 'number'
|
|
) {
|
|
return part;
|
|
}
|
|
return {
|
|
...part,
|
|
activity_start_index: part.activity_start_index + phaseIndexOffset,
|
|
...(typeof part.activity_end_index === 'number' && {
|
|
activity_end_index: part.activity_end_index + phaseIndexOffset,
|
|
}),
|
|
};
|
|
});
|
|
|
|
if (editedType !== ContentTypes.TEXT && editedType !== ContentTypes.THINK) {
|
|
return existingContent.concat(adjustedCompletion);
|
|
}
|
|
|
|
if (!mergesFirstPart) {
|
|
return existingContent.concat(adjustedCompletion);
|
|
}
|
|
|
|
const mergedContent = [...existingContent];
|
|
if (editedType === ContentTypes.TEXT) {
|
|
mergedContent[lastIndex] = {
|
|
...mergedContent[lastIndex],
|
|
...(firstNew.phase != null && { phase: firstNew.phase }),
|
|
[ContentTypes.TEXT]:
|
|
(mergedContent[lastIndex][ContentTypes.TEXT] || '') +
|
|
(adjustedCompletion[0][ContentTypes.TEXT] || ''),
|
|
};
|
|
} else {
|
|
mergedContent[lastIndex] = {
|
|
...mergedContent[lastIndex],
|
|
[ContentTypes.THINK]:
|
|
(mergedContent[lastIndex][ContentTypes.THINK] || '') +
|
|
(adjustedCompletion[0][ContentTypes.THINK] || ''),
|
|
};
|
|
}
|
|
|
|
// Add remaining completion items
|
|
return mergedContent.concat(adjustedCompletion.slice(1));
|
|
}
|
|
|
|
async sendPayload(payload, opts = {}) {
|
|
if (opts && typeof opts === 'object') {
|
|
this.setOptions(opts);
|
|
}
|
|
|
|
return await this.sendCompletion(payload, opts);
|
|
}
|
|
|
|
async addDocuments(message, attachments) {
|
|
const documentResult = await encodeAndFormatDocuments(
|
|
this.options.req,
|
|
attachments,
|
|
{
|
|
provider: this.options.agent?.provider ?? this.options.endpoint,
|
|
endpoint: this.options.agent?.endpoint ?? this.options.endpoint,
|
|
useResponsesApi: this.options.agent?.model_parameters?.useResponsesApi,
|
|
model: this.modelOptions?.model ?? this.model,
|
|
},
|
|
getStrategyFunctions,
|
|
);
|
|
message.documents =
|
|
documentResult.documents && documentResult.documents.length
|
|
? documentResult.documents
|
|
: undefined;
|
|
return documentResult.files;
|
|
}
|
|
|
|
async addVideos(message, attachments) {
|
|
const videoResult = await encodeAndFormatVideos(
|
|
this.options.req,
|
|
attachments,
|
|
{
|
|
provider: this.options.agent?.provider ?? this.options.endpoint,
|
|
endpoint: this.options.agent?.endpoint ?? this.options.endpoint,
|
|
},
|
|
getStrategyFunctions,
|
|
);
|
|
message.videos =
|
|
videoResult.videos && videoResult.videos.length ? videoResult.videos : undefined;
|
|
return videoResult.files;
|
|
}
|
|
|
|
async addAudios(message, attachments) {
|
|
const audioResult = await encodeAndFormatAudios(
|
|
this.options.req,
|
|
attachments,
|
|
{
|
|
provider: this.options.agent?.provider ?? this.options.endpoint,
|
|
endpoint: this.options.agent?.endpoint ?? this.options.endpoint,
|
|
},
|
|
getStrategyFunctions,
|
|
);
|
|
message.audios =
|
|
audioResult.audios && audioResult.audios.length ? audioResult.audios : undefined;
|
|
return audioResult.files;
|
|
}
|
|
|
|
/**
|
|
* Extracts text context from attachments and sets it on the message.
|
|
* This handles text that was already extracted from files (OCR, transcriptions, document text, etc.)
|
|
* @param {TMessage} message - The message to add context to
|
|
* @param {MongoFile[]} attachments - Array of file attachments
|
|
* @returns {Promise<void>}
|
|
*/
|
|
async addFileContextToMessage(message, attachments) {
|
|
const fileContext = await extractFileContext({
|
|
attachments,
|
|
req: this.options?.req,
|
|
tokenCountFn: (text) => countTokens(text),
|
|
});
|
|
|
|
if (fileContext) {
|
|
message.fileContext = fileContext;
|
|
}
|
|
}
|
|
|
|
async processAttachments(message, attachments) {
|
|
const categorizedAttachments = {
|
|
images: [],
|
|
videos: [],
|
|
audios: [],
|
|
documents: [],
|
|
};
|
|
|
|
const allFiles = [];
|
|
|
|
const provider = this.options.agent?.provider ?? this.options.endpoint;
|
|
const isBedrock = provider === EModelEndpoint.bedrock;
|
|
|
|
if (!this._mergedFileConfig) {
|
|
this._mergedFileConfig = mergeFileConfig(this.options.req?.config?.fileConfig);
|
|
const endpoint = this.options.agent?.endpoint ?? this.options.endpoint;
|
|
this._endpointFileConfig = getEndpointFileConfig({
|
|
fileConfig: this._mergedFileConfig,
|
|
endpoint,
|
|
endpointType: this.options.endpointType,
|
|
});
|
|
}
|
|
|
|
for (const file of attachments) {
|
|
/** @type {FileSources} */
|
|
const source = file.source ?? FileSources.local;
|
|
if (source === FileSources.text) {
|
|
allFiles.push(file);
|
|
continue;
|
|
}
|
|
if (
|
|
file.embedded === true ||
|
|
file.metadata?.codeEnvRef != null ||
|
|
file.metadata?.fileIdentifier != null
|
|
) {
|
|
allFiles.push(file);
|
|
continue;
|
|
}
|
|
|
|
if (file.type.startsWith('image/')) {
|
|
categorizedAttachments.images.push(file);
|
|
} else if (file.type === 'application/pdf') {
|
|
categorizedAttachments.documents.push(file);
|
|
allFiles.push(file);
|
|
} else if (isBedrock && isBedrockDocumentType(file.type)) {
|
|
categorizedAttachments.documents.push(file);
|
|
allFiles.push(file);
|
|
} else if (file.type.startsWith('video/')) {
|
|
categorizedAttachments.videos.push(file);
|
|
allFiles.push(file);
|
|
} else if (file.type.startsWith('audio/')) {
|
|
categorizedAttachments.audios.push(file);
|
|
allFiles.push(file);
|
|
} else if (
|
|
file.type &&
|
|
this._mergedFileConfig &&
|
|
this._endpointFileConfig?.supportedMimeTypes &&
|
|
this._mergedFileConfig.checkType(file.type, this._endpointFileConfig.supportedMimeTypes)
|
|
) {
|
|
categorizedAttachments.documents.push(file);
|
|
allFiles.push(file);
|
|
}
|
|
}
|
|
|
|
const [imageFiles] = await Promise.all([
|
|
categorizedAttachments.images.length > 0
|
|
? this.addImageURLs(message, categorizedAttachments.images)
|
|
: Promise.resolve([]),
|
|
categorizedAttachments.documents.length > 0
|
|
? this.addDocuments(message, categorizedAttachments.documents)
|
|
: Promise.resolve([]),
|
|
categorizedAttachments.videos.length > 0
|
|
? this.addVideos(message, categorizedAttachments.videos)
|
|
: Promise.resolve([]),
|
|
categorizedAttachments.audios.length > 0
|
|
? this.addAudios(message, categorizedAttachments.audios)
|
|
: Promise.resolve([]),
|
|
]);
|
|
|
|
allFiles.push(...imageFiles);
|
|
|
|
const seenFileIds = new Set();
|
|
const uniqueFiles = [];
|
|
|
|
for (const file of allFiles) {
|
|
if (file.file_id && !seenFileIds.has(file.file_id)) {
|
|
seenFileIds.add(file.file_id);
|
|
uniqueFiles.push(file);
|
|
} else if (!file.file_id) {
|
|
uniqueFiles.push(file);
|
|
}
|
|
}
|
|
|
|
return uniqueFiles;
|
|
}
|
|
|
|
/**
|
|
* @param {TMessage[]} _messages
|
|
* @returns {Promise<TMessage[]>}
|
|
*/
|
|
async addPreviousAttachments(_messages) {
|
|
if (!this.options.resendFiles) {
|
|
return _messages;
|
|
}
|
|
|
|
const contextSeen = new Set();
|
|
const attachmentsProcessed =
|
|
this.options.attachments && !(this.options.attachments instanceof Promise);
|
|
if (attachmentsProcessed) {
|
|
for (const attachment of this.options.attachments) {
|
|
if (attachment?.file_id) {
|
|
contextSeen.add(attachment.file_id);
|
|
}
|
|
}
|
|
}
|
|
|
|
const historicalFileIds = collectHistoricalFileIds(_messages);
|
|
const fileFilter = buildOwnerFileFilter(historicalFileIds, this.options.req?.user);
|
|
const authorizedFilesById = new Map();
|
|
if (fileFilter) {
|
|
const files = (await db.getFiles(fileFilter, {}, {})) ?? [];
|
|
for (const file of files) {
|
|
if (file?.file_id) {
|
|
authorizedFilesById.set(file.file_id, file);
|
|
}
|
|
}
|
|
}
|
|
/** Owner-scoped docs for THIS turn, including steer-part refs — the steer
|
|
* replay stamp consumes this instead of issuing a second query. */
|
|
this.authorizedHistoricalFiles = authorizedFilesById;
|
|
|
|
/**
|
|
*
|
|
* @param {TMessage} message
|
|
*/
|
|
const processMessage = async (message) => {
|
|
if (!this.message_file_map) {
|
|
/** @type {Record<string, MongoFile[]> */
|
|
this.message_file_map = {};
|
|
}
|
|
|
|
delete message.fileContext;
|
|
|
|
const contextFiles = [];
|
|
if (Array.isArray(message.files)) {
|
|
for (const file of message.files) {
|
|
if (!file?.file_id || contextSeen.has(file.file_id)) {
|
|
continue;
|
|
}
|
|
const authorizedFile = authorizedFilesById.get(file.file_id);
|
|
if (authorizedFile) {
|
|
contextFiles.push(authorizedFile);
|
|
contextSeen.add(file.file_id);
|
|
}
|
|
}
|
|
}
|
|
|
|
const rehydratedFiles = rehydrateMessageFileRefs(message.files, authorizedFilesById);
|
|
if (rehydratedFiles) {
|
|
message.files = rehydratedFiles;
|
|
} else {
|
|
delete message.files;
|
|
}
|
|
|
|
const rehydratedAttachments = rehydrateMessageFileRefs(
|
|
message.attachments,
|
|
authorizedFilesById,
|
|
{
|
|
preserveDisplayOnly: true,
|
|
},
|
|
);
|
|
if (rehydratedAttachments) {
|
|
message.attachments = rehydratedAttachments;
|
|
} else {
|
|
delete message.attachments;
|
|
}
|
|
|
|
if (contextFiles.length === 0) {
|
|
return message;
|
|
}
|
|
|
|
await Promise.all([
|
|
this.addFileContextToMessage(message, contextFiles),
|
|
this.processAttachments(message, contextFiles),
|
|
]);
|
|
|
|
this.message_file_map[message.messageId] = contextFiles;
|
|
return message;
|
|
};
|
|
|
|
const promises = [];
|
|
|
|
for (const message of _messages) {
|
|
if (!message.files && !message.attachments) {
|
|
promises.push(message);
|
|
continue;
|
|
}
|
|
|
|
promises.push(processMessage(message));
|
|
}
|
|
|
|
const messages = await Promise.all(promises);
|
|
|
|
this.checkVisionRequest(Object.values(this.message_file_map ?? {}).flat());
|
|
return messages;
|
|
}
|
|
}
|
|
|
|
module.exports = BaseClient;
|