mirror of
https://github.com/danny-avila/LibreChat.git
synced 2026-10-08 22:41:38 +00:00
♻️ refactor: Compute Context Gauge Client-Side, Drop Projection Endpoint (#13953)
* ♻️ refactor: Compute Context Gauge Client-Side, Drop Projection Endpoint The /api/endpoints/context-projection endpoint re-fetched a conversation's messages from Mongo and re-tokenized them to project the context gauge for snapshot-less branches. The browser already holds those messages and their per-message tokenCounts, so this duplicated work on the request path (an unbounded read + server-side BPE tokenization until it was later capped). Move the snapshot-less estimate fully client-side, from the in-memory index: - sumBranch accumulates an uncalibrated char/4 estimate (estTokens) for count-less messages (imports / pre-feature) under the same summary cutoff - useTokenUsage folds estTokens (calibrated via the existing calibrationFamily ratio) into the existing fallback; known per-message counts render unchanged - delete the endpoint, controller, rate limiter, route, the getMessageTextStats data-schemas method, and the data-provider surface (endpoint/key/type/service/query) No DB read, no server tokenization, no rate-limit knobs; the gauge recomputes reactively from the index. Net -793 lines. * 🩹 fix: Count quotes and object-form content in client context estimate Address Codex review on the client-side context estimate: - messageChars now reads object-form content text (part.text.value), not only string text/think, so imported / pre-feature messages whose body lives in content parts are no longer estimated as zero. - Count-less user messages include their merged quote excerpts in the estimate, mirroring what the send path prepends into the prompt. * 🩹 fix: Cap over-window estimate and surface estimated tokens in breakdown Address remaining Codex review on the client-side context estimate: - Clamp the snapshot-less estimate's displayed usedTokens to maxTokens. The send path prunes an over-window branch before calling the model, so the gauge never actually exceeds the window; this avoids impossible values (e.g. 50k / 8k) without re-introducing client-side pruning. - Surface the calibrated count-less estimate as its own "Estimated" row in the breakdown popover, so a branch of only count-less imported / pre-feature messages is no longer shown as Input 0 / Output 0 under a non-zero header. * 🩹 fix: Refine client context estimate per Codex re-review - Drop calibration from the snapshot-less estimate. The removed projection never actually calibrated (the client never sent a ratio), and a ratio inflated by provider-injected context over-estimates visible imported text. - Exclude reasoning (think) / error parts from the estimate; the send path strips them, so they are not part of the next call's context. - Fold quote text into the estimate even when a tokenCount is present, since the edit route recounts tokenCount from text only and drops the merged quote. * 🩹 fix: Recount quoted user turns instead of topping up the stored count The previous round added quote chars on top of a quoted message's stored tokenCount, which double-counts the common (unedited) case where the count already includes the merged quote prompt. Match the removed projection instead: for quoted user turns, ignore the stored count and estimate the full merged text. This both avoids the double-count and still corrects the stale text-only count an edit leaves behind. * 🩹 fix: Trust stored counts for quoted turns; count tool-call parts - Quoted user turns: revert to trusting a present tokenCount. The send path's stored count already includes the merged quote (and any calibration), and the client's char/4 path is coarser, so recounting regressed normal turns. Only count-less messages estimate quotes from text. - Count tool-call name/args/output for count-less assistant messages; the formatter sends them back as context, so omitting them under-reported imported branches with tool history. * 🩹 fix: Exclude in-flight tail from estimate to avoid resume double-count On resume the live path seeds liveTokens from the partial response and also writes that content into the messages cache, where the count-less response is estimated into estTokens too — double-counting the in-flight output on the snapshot-less estimate path. sumBranch now exposes the tail message's own estimate (tailEstTokens); the estimate path drops it while a stream is live, so the in-flight response is counted once (via liveTokens). The breakdown's Estimated row uses the same in-flight-adjusted value. * 🩹 fix: Recount quoted user turns in context estimate (match send path) A quoted user turn's stored tokenCount is unreliable for the gauge: a text-only Save edit recomputes it from text alone, and the send path (needsCanonicalTokenCount in agents/client.js) recounts the quote-merged prompt every turn regardless of the stored value. Mirror that on the client — estimate quoted turns from the merged text+quotes and ignore the stored count — so snapshot-less branches don't under-report by the quote block. Reverts the earlier "trust the count" assumption, which the server disproves. * 🧹 chore: Route useResumableSSE diagnostics through the frontend logger Convert the [ResumableSSE]/[Debug] console.log and console.error diagnostics to the gated frontend `logger` (client/src/utils/logger), splitting the tag from the message so object arguments are passed through as real args (logged expandably, not stringified) and the logs stay tag-filterable and off the production console unless explicitly enabled. All log statements preserved; nothing removed. * 🩹 fix: Prefer content over text when estimating count-less messages A stopped agent response is saved with both a `text` field and a structured `content` array, and the send path formats from content. messageChars early-returned on `text`, dropping the content array (and the tool-call tokens it carries) from the snapshot-less estimate — also making the tool_call handling dead for such messages. Prefer content when present, fall back to text.
This commit is contained in:
parent
e26ce4713a
commit
376370d610
20 changed files with 331 additions and 918 deletions
|
|
@ -1,35 +0,0 @@
|
|||
const { logger } = require('@librechat/data-schemas');
|
||||
const { resolveContextProjection } = require('@librechat/api');
|
||||
const db = require('~/models');
|
||||
|
||||
/**
|
||||
* Returns a server-side context-usage projection for the viewed branch + config
|
||||
* (agents SDK, no model call) — powers the gauge for snapshot-less branches and
|
||||
* after a model/window switch. Resolution lives in `@librechat/api`; this
|
||||
* controller only injects request-scoped model accessors.
|
||||
* @param {ServerRequest} req
|
||||
* @param {ServerResponse} res
|
||||
*/
|
||||
async function contextProjectionController(req, res) {
|
||||
try {
|
||||
const params = req.body ?? {};
|
||||
if (!params.conversationId || !params.messageId) {
|
||||
res.json(null);
|
||||
return;
|
||||
}
|
||||
const projection = await resolveContextProjection(
|
||||
{
|
||||
userId: req.user?.id,
|
||||
getMessages: db.getMessages,
|
||||
getMessageTextStats: db.getMessageTextStats,
|
||||
},
|
||||
params,
|
||||
);
|
||||
res.json(projection ?? null);
|
||||
} catch (error) {
|
||||
logger.error('[contextProjectionController]', error);
|
||||
res.status(500).json({ error: 'Failed to resolve context projection' });
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = contextProjectionController;
|
||||
|
|
@ -1,19 +0,0 @@
|
|||
const rateLimit = require('express-rate-limit');
|
||||
const { limiterCache } = require('@librechat/api');
|
||||
|
||||
const { CONTEXT_PROJECTION_WINDOW = 1, CONTEXT_PROJECTION_MAX = 20 } = process.env;
|
||||
|
||||
const windowMs = (parseInt(CONTEXT_PROJECTION_WINDOW, 10) || 1) * 60 * 1000;
|
||||
const max = parseInt(CONTEXT_PROJECTION_MAX, 10) || 20;
|
||||
|
||||
const contextProjectionLimiter = rateLimit({
|
||||
windowMs,
|
||||
max,
|
||||
handler: (_req, res) => {
|
||||
res.status(429).json({ message: 'Too many context projection requests. Try again later' });
|
||||
},
|
||||
keyGenerator: (req) => req.user?.id,
|
||||
store: limiterCache('context_projection_limiter'),
|
||||
});
|
||||
|
||||
module.exports = contextProjectionLimiter;
|
||||
|
|
@ -9,7 +9,6 @@ const registerLimiter = require('./registerLimiter');
|
|||
const toolCallLimiter = require('./toolCallLimiter');
|
||||
const messageLimiters = require('./messageLimiters');
|
||||
const promptUsageLimiter = require('./promptUsageLimiter');
|
||||
const contextProjectionLimiter = require('./contextProjectionLimiter');
|
||||
const verifyEmailLimiter = require('./verifyEmailLimiter');
|
||||
const resetPasswordLimiter = require('./resetPasswordLimiter');
|
||||
const twoFactorTempLimiter = require('./twoFactorTempLimiter');
|
||||
|
|
@ -23,7 +22,6 @@ module.exports = {
|
|||
loginLimiter,
|
||||
registerLimiter,
|
||||
toolCallLimiter,
|
||||
contextProjectionLimiter,
|
||||
createTTSLimiters,
|
||||
createSTTLimiters,
|
||||
verifyEmailLimiter,
|
||||
|
|
|
|||
|
|
@ -3,19 +3,10 @@ const requireJwtAuth = require('~/server/middleware/requireJwtAuth');
|
|||
const configMiddleware = require('~/server/middleware/config/app');
|
||||
const endpointController = require('~/server/controllers/EndpointController');
|
||||
const tokenConfigController = require('~/server/controllers/TokenConfigController');
|
||||
const contextProjectionController = require('~/server/controllers/ContextProjectionController');
|
||||
const { contextProjectionLimiter } = require('~/server/middleware/limiters');
|
||||
|
||||
const router = express.Router();
|
||||
/** Auth required for role/tenant-scoped endpoint config resolution. */
|
||||
router.get('/', requireJwtAuth, endpointController);
|
||||
router.get('/token-config', requireJwtAuth, configMiddleware, tokenConfigController);
|
||||
router.post(
|
||||
'/context-projection',
|
||||
requireJwtAuth,
|
||||
contextProjectionLimiter,
|
||||
configMiddleware,
|
||||
contextProjectionController,
|
||||
);
|
||||
|
||||
module.exports = router;
|
||||
|
|
|
|||
|
|
@ -150,6 +150,9 @@ export default function Breakdown({ view, showCost, currency }: BreakdownProps)
|
|||
label={localize('com_ui_output')}
|
||||
value={view.branchTotals.output + view.liveTokens}
|
||||
/>
|
||||
{view.estimatedTokens > 0 && (
|
||||
<Row label={localize('com_ui_context_estimated')} value={view.estimatedTokens} />
|
||||
)}
|
||||
{maxTokens == null && (
|
||||
<p className="text-xs text-text-secondary">{localize('com_ui_context_unknown')}</p>
|
||||
)}
|
||||
|
|
|
|||
|
|
@ -41,45 +41,6 @@ export const useTokenConfigQuery = (
|
|||
});
|
||||
};
|
||||
|
||||
/**
|
||||
* Server-side context-usage projection for the viewed branch + resolved config
|
||||
* (agents SDK, no model call). Keyed on the fields that change what the next
|
||||
* call would send — branch tail, endpoint/model/agent, window — so a branch or
|
||||
* model/window switch refetches. Disabled until a branch tail is known.
|
||||
*/
|
||||
export const useContextProjectionQuery = (
|
||||
params: t.TContextProjectionRequest | null,
|
||||
config?: UseQueryOptions<t.TContextUsageEvent | null>,
|
||||
): QueryObserverResult<t.TContextUsageEvent | null> => {
|
||||
const queriesEnabled = useRecoilValue<boolean>(store.queriesEnabled);
|
||||
return useQuery<t.TContextUsageEvent | null>(
|
||||
[
|
||||
QueryKeys.contextProjection,
|
||||
params?.conversationId,
|
||||
params?.messageId,
|
||||
params?.endpoint,
|
||||
params?.model,
|
||||
params?.agentId,
|
||||
params?.maxContextTokens,
|
||||
params?.revision,
|
||||
],
|
||||
() => dataService.getContextProjection(params as t.TContextProjectionRequest),
|
||||
{
|
||||
staleTime: Infinity,
|
||||
refetchOnWindowFocus: false,
|
||||
refetchOnReconnect: false,
|
||||
refetchOnMount: false,
|
||||
...config,
|
||||
enabled:
|
||||
(config?.enabled ?? true) === true &&
|
||||
queriesEnabled &&
|
||||
params != null &&
|
||||
(params.conversationId?.length ?? 0) > 0 &&
|
||||
(params.messageId?.length ?? 0) > 0,
|
||||
},
|
||||
);
|
||||
};
|
||||
|
||||
/**
|
||||
* Auth-aware query key so unauthenticated (login page) and authenticated
|
||||
* (chat page) configs are cached independently, preventing stale
|
||||
|
|
|
|||
|
|
@ -2,13 +2,7 @@ import { useEffect, useMemo, useRef } from 'react';
|
|||
import { useAtomValue, useSetAtom } from 'jotai';
|
||||
import { useQueryClient } from '@tanstack/react-query';
|
||||
import { Constants, QueryKeys } from 'librechat-data-provider';
|
||||
import type {
|
||||
TMessage,
|
||||
TConversation,
|
||||
TModelTokenomics,
|
||||
TContextUsageEvent,
|
||||
TContextProjectionRequest,
|
||||
} from 'librechat-data-provider';
|
||||
import type { TMessage, TConversation, TModelTokenomics } from 'librechat-data-provider';
|
||||
import type { BranchTotals, BranchUsage } from '~/utils/tokens';
|
||||
import type { ContextSnapshot } from '~/store/usage';
|
||||
import {
|
||||
|
|
@ -30,7 +24,6 @@ import {
|
|||
findBranchSnapshotAnchor,
|
||||
} from '~/utils';
|
||||
import { useLatestMessageId } from '~/hooks/Messages/useLatestMessage';
|
||||
import { useContextProjectionQuery } from '~/data-provider';
|
||||
import useTokenLimits from './useTokenLimits';
|
||||
|
||||
export interface TokenUsageParams {
|
||||
|
|
@ -60,6 +53,9 @@ export interface TokenUsageView {
|
|||
/** Authoritative cost across all branches (shown when it differs from branch) */
|
||||
totalCost: number;
|
||||
liveTokens: number;
|
||||
/** Estimated tokens for count-less messages (in-flight tail excluded while
|
||||
* streaming); 0 on snapshots. Rendered as its own breakdown row. */
|
||||
estimatedTokens: number;
|
||||
rates?: TModelTokenomics;
|
||||
}
|
||||
|
||||
|
|
@ -101,40 +97,6 @@ export default function useTokenUsage({
|
|||
return anchor != null ? (snapshotsByAnchor.get(anchor) ?? null) : null;
|
||||
}, [conversationKey, branchTotals.tailId, snapshotsByAnchor]);
|
||||
|
||||
const resolvedMax = limits.maxContextTokens;
|
||||
|
||||
/** Project the branch (agents SDK, no model call) ONLY when no persisted/live
|
||||
* snapshot covers it — snapshot-less branches (G2: pre-feature history,
|
||||
* imports, never-generated branches). A present snapshot stays authoritative;
|
||||
* reliable window-switch (G1) detection needs the snapshot to carry its
|
||||
* model/window (deferred to the fidelity follow-up), and the SDK window
|
||||
* (reserve-derived) doesn't equal the client-resolved raw window, so we must
|
||||
* NOT mis-flag a valid snapshot as stale here. Cached + refetched by branch/
|
||||
* endpoint/model/window/revision. */
|
||||
const projectionParams: TContextProjectionRequest | null =
|
||||
!isSubmitting &&
|
||||
branchSnapshot == null &&
|
||||
conversation?.conversationId != null &&
|
||||
conversation.conversationId !== Constants.NEW_CONVO &&
|
||||
branchTotals.tailId != null &&
|
||||
conversation.endpoint != null
|
||||
? {
|
||||
conversationId: conversation.conversationId,
|
||||
messageId: branchTotals.tailId,
|
||||
/** Resolved provider/model (e.g. an agent's actual provider, not the
|
||||
* `agents` endpoint) so the server picks the right tokenizer. */
|
||||
endpoint: limits.endpoint || conversation.endpoint,
|
||||
model: limits.model || conversation.model || undefined,
|
||||
agentId: conversation.agent_id ?? undefined,
|
||||
spec: conversation.spec ?? undefined,
|
||||
maxContextTokens: resolvedMax,
|
||||
/** Content revision so an in-place message edit (same tail id) refetches. */
|
||||
revision: branchTotals.input + branchTotals.output,
|
||||
}
|
||||
: null;
|
||||
const { data: projectionData } = useContextProjectionQuery(projectionParams);
|
||||
const projection = projectionData ?? null;
|
||||
|
||||
/** Branch/total provider usage is index-derived; the in-flight response is
|
||||
* the only live add (the pending holder), counted into both — it sits on the
|
||||
* active branch tail and inside the conversation. The backend prices each
|
||||
|
|
@ -242,25 +204,11 @@ export default function useTokenUsage({
|
|||
snapshot != null &&
|
||||
(isSubmitting || (snapshot.anchorMessageId != null && branchTotals.containsAnchor));
|
||||
|
||||
/** Precedence: live/active snapshot → persisted branch snapshot → server
|
||||
* projection (snapshot-less branches, G2) → per-message estimate. The first
|
||||
* two preserve the pre-projection behavior exactly; the projection only
|
||||
* slots in ahead of the estimate when no snapshot exists. Snapshot and
|
||||
* projection share the render-relevant fields, so they render uniformly. */
|
||||
let effective: ContextSnapshot | TContextUsageEvent | null = null;
|
||||
/** A server projection is the SDK's windowing but, in this first cut, omits
|
||||
* instruction/tool overhead — so it's surfaced as an ESTIMATE (a better one
|
||||
* than sumBranch), never a false-authoritative number. Real snapshots stay
|
||||
* authoritative. */
|
||||
let projected = false;
|
||||
if (currentActive) {
|
||||
effective = snapshot;
|
||||
} else if (branchSnapshot != null) {
|
||||
effective = branchSnapshot;
|
||||
} else if (projection != null) {
|
||||
effective = projection;
|
||||
projected = true;
|
||||
}
|
||||
/** Precedence: live/active snapshot → persisted branch snapshot →
|
||||
* per-message estimate. The first two are authoritative (real runs with the
|
||||
* feature on); the estimate covers snapshot-less branches (pre-feature
|
||||
* history, imports, never-generated branches) entirely client-side. */
|
||||
const effective: ContextSnapshot | null = currentActive ? snapshot : branchSnapshot;
|
||||
|
||||
if (effective != null) {
|
||||
const breakdown = effective.breakdown;
|
||||
|
|
@ -270,18 +218,18 @@ export default function useTokenUsage({
|
|||
effective.remainingContextTokens != null
|
||||
? maxTokens - effective.remainingContextTokens
|
||||
: instructionTokens + breakdown.messageTokens;
|
||||
/** The snapshot/projection is pre-invoke: in-flight output rides on
|
||||
* `liveTokens` (0 unless streaming this branch), the last call's finalized
|
||||
* output on `completedOutputTokens` (absent on a projection → 0). */
|
||||
/** The snapshot is pre-invoke: in-flight output rides on `liveTokens` (0
|
||||
* unless streaming this branch), the last call's finalized output on
|
||||
* `completedOutputTokens`. */
|
||||
const usedTokens =
|
||||
Math.max(0, baseUsed) + liveTokens + (effective.completedOutputTokens ?? 0);
|
||||
return {
|
||||
usedTokens,
|
||||
maxTokens,
|
||||
percent: maxTokens > 0 ? Math.min((usedTokens / maxTokens) * 100, 100) : 0,
|
||||
isEstimate: projected,
|
||||
snapshot: projected ? null : (effective as ContextSnapshot),
|
||||
snapshotActive: !projected,
|
||||
isEstimate: false,
|
||||
snapshot: effective,
|
||||
snapshotActive: true,
|
||||
branchTotals,
|
||||
branchUsage,
|
||||
totalUsage,
|
||||
|
|
@ -289,18 +237,38 @@ export default function useTokenUsage({
|
|||
branchCost: branchUsage.cost,
|
||||
totalCost: totalUsage.cost,
|
||||
liveTokens,
|
||||
estimatedTokens: 0,
|
||||
rates: limits.rates,
|
||||
};
|
||||
}
|
||||
|
||||
/** `summaryBaseline` is the compacted-context size from the deepest
|
||||
* summarized response on the branch (0 if none). The branch walk stops
|
||||
* there, so input/output are post-summary only — adding the baseline keeps
|
||||
* the estimate from re-summing the discarded pre-summary history (which
|
||||
* otherwise pins the gauge at 100% forever after a compaction). */
|
||||
const usedTokens =
|
||||
branchTotals.input + branchTotals.output + branchTotals.summaryBaseline + liveTokens;
|
||||
/** Snapshot-less estimate, computed from the in-memory message index — no
|
||||
* server round-trip. All terms are local per-message counts / char estimates
|
||||
* (uncalibrated): the learned calibration ratio reconciles provider-injected
|
||||
* context that isn't present in this visible text, so applying it here would
|
||||
* over-inflate. `summaryBaseline` is the compacted-context size from the
|
||||
* deepest summarized response on the branch (0 if none); the walk stops
|
||||
* there, so input/output are post-summary only — adding it keeps the estimate
|
||||
* from re-summing the discarded pre-summary history (which otherwise pins the
|
||||
* gauge at 100% after a compaction). */
|
||||
const maxTokens = limits.maxContextTokens;
|
||||
/** When a stream is live the tail is the in-flight response, already counted
|
||||
* by `liveTokens`; drop its static estimate so a resumed/partial response
|
||||
* isn't double-counted on the estimate path. */
|
||||
const estimatedTokens = Math.max(
|
||||
0,
|
||||
branchTotals.estTokens - (liveTokens > 0 ? branchTotals.tailEstTokens : 0),
|
||||
);
|
||||
const rawUsed =
|
||||
branchTotals.input +
|
||||
branchTotals.output +
|
||||
estimatedTokens +
|
||||
branchTotals.summaryBaseline +
|
||||
liveTokens;
|
||||
/** The send path prunes an over-window branch before calling the model, so the
|
||||
* live gauge never actually exceeds the window; clamp the display to the
|
||||
* window rather than show impossible values (e.g. 50k / 8k). */
|
||||
const usedTokens = maxTokens != null && maxTokens > 0 ? Math.min(rawUsed, maxTokens) : rawUsed;
|
||||
return {
|
||||
usedTokens,
|
||||
maxTokens,
|
||||
|
|
@ -316,6 +284,7 @@ export default function useTokenUsage({
|
|||
branchCost: branchUsage.cost,
|
||||
totalCost: totalUsage.cost,
|
||||
liveTokens,
|
||||
estimatedTokens,
|
||||
rates: limits.rates,
|
||||
};
|
||||
}, [
|
||||
|
|
@ -328,6 +297,5 @@ export default function useTokenUsage({
|
|||
liveTokens,
|
||||
limits,
|
||||
branchSnapshot,
|
||||
projection,
|
||||
]);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -27,6 +27,7 @@ import type { EventHandlerParams } from './useEventHandlers';
|
|||
import type { ActiveJobsResponse } from '~/data-provider';
|
||||
import type { TResData } from '~/common';
|
||||
import {
|
||||
logger,
|
||||
clearAllDrafts,
|
||||
removeConvoFromAllQueries,
|
||||
upsertConvoInAllQueries,
|
||||
|
|
@ -510,7 +511,7 @@ export default function useResumableSSE(
|
|||
|
||||
const baseUrl = `${apiBaseUrl()}/api/agents/chat/stream/${encodeURIComponent(currentStreamId)}`;
|
||||
const url = isResume ? `${baseUrl}?resume=true` : baseUrl;
|
||||
console.log('[ResumableSSE] Subscribing to stream:', url, { isResume });
|
||||
logger.log('ResumableSSE', 'Subscribing to stream:', url, { isResume });
|
||||
|
||||
const sse = new SSE(url, {
|
||||
headers: { Authorization: `Bearer ${token}` },
|
||||
|
|
@ -519,7 +520,7 @@ export default function useResumableSSE(
|
|||
sseRef.current = sse;
|
||||
|
||||
sse.addEventListener('open', () => {
|
||||
console.log('[ResumableSSE] Stream connected');
|
||||
logger.log('ResumableSSE', 'Stream connected');
|
||||
setAbortScroll(false);
|
||||
// Restore UI state on successful connection (including reconnection)
|
||||
setIsSubmitting(true);
|
||||
|
|
@ -532,7 +533,7 @@ export default function useResumableSSE(
|
|||
const data = JSON.parse(e.data);
|
||||
|
||||
if (data.final != null) {
|
||||
console.log('[ResumableSSE] Received FINAL event', {
|
||||
logger.log('ResumableSSE', 'Received FINAL event', {
|
||||
aborted: data.aborted,
|
||||
conversationId: data.conversation?.conversationId,
|
||||
hasResponseMessage: !!data.responseMessage,
|
||||
|
|
@ -545,7 +546,7 @@ export default function useResumableSSE(
|
|||
finalHandler(data, currentSubmission as EventSubmission);
|
||||
finalizeUsage(data, { ...currentSubmission, userMessage });
|
||||
} catch (error) {
|
||||
console.error('[ResumableSSE] Error in finalHandler:', error);
|
||||
logger.error('ResumableSSE', 'Error in finalHandler:', error);
|
||||
setIsSubmitting(false);
|
||||
setShowStopButton(false);
|
||||
}
|
||||
|
|
@ -562,7 +563,7 @@ export default function useResumableSSE(
|
|||
}
|
||||
|
||||
if (data.created != null) {
|
||||
console.log('[ResumableSSE] Received CREATED event', {
|
||||
logger.log('ResumableSSE', 'Received CREATED event', {
|
||||
messageId: data.message?.messageId,
|
||||
conversationId: data.message?.conversationId,
|
||||
});
|
||||
|
|
@ -634,7 +635,7 @@ export default function useResumableSSE(
|
|||
}
|
||||
|
||||
if (data.sync != null) {
|
||||
console.log('[ResumableSSE] SYNC received', {
|
||||
logger.log('ResumableSSE', 'SYNC received', {
|
||||
runSteps: data.resumeState?.runSteps?.length ?? 0,
|
||||
pendingEvents: data.pendingEvents?.length ?? 0,
|
||||
});
|
||||
|
|
@ -700,7 +701,7 @@ export default function useResumableSSE(
|
|||
);
|
||||
}
|
||||
|
||||
console.log('[ResumableSSE] SYNC update', {
|
||||
logger.log('ResumableSSE', 'SYNC update', {
|
||||
userMsgId,
|
||||
serverResponseId,
|
||||
responseIdx,
|
||||
|
|
@ -721,7 +722,7 @@ export default function useResumableSSE(
|
|||
model: preferDefinedString(messages[responseIdx]?.model, data.resumeState.model),
|
||||
} as TMessage;
|
||||
const updated = mergeResumeMessages(messages, userMessage, responseMessage);
|
||||
console.log('[ResumableSSE] SYNC updating message', {
|
||||
logger.log('ResumableSSE', 'SYNC updating message', {
|
||||
messageId: responseMessage.messageId,
|
||||
oldContentLength: Array.isArray(oldContent) ? oldContent.length : 0,
|
||||
newContentLength: data.resumeState.aggregatedContent?.length,
|
||||
|
|
@ -729,7 +730,7 @@ export default function useResumableSSE(
|
|||
setMessages(updated);
|
||||
resetContentHandler();
|
||||
syncStepMessage(responseMessage);
|
||||
console.log('[ResumableSSE] SYNC complete, handlers synced');
|
||||
logger.log('ResumableSSE', 'SYNC complete, handlers synced');
|
||||
} else {
|
||||
const responseId = serverResponseId ?? `${userMsgId}_`;
|
||||
const newMessage = {
|
||||
|
|
@ -753,8 +754,9 @@ export default function useResumableSSE(
|
|||
}
|
||||
|
||||
if (data.resumeState?.replayEvents?.length > 0) {
|
||||
console.log(
|
||||
`[ResumableSSE] Replaying ${data.resumeState.replayEvents.length} resume events`,
|
||||
logger.log(
|
||||
'ResumableSSE',
|
||||
`Replaying ${data.resumeState.replayEvents.length} resume events`,
|
||||
);
|
||||
for (const replayEvent of data.resumeState.replayEvents) {
|
||||
if (replayEvent.event === UsageEvents.ON_CONTEXT_USAGE) {
|
||||
|
|
@ -774,7 +776,7 @@ export default function useResumableSSE(
|
|||
}
|
||||
|
||||
if (data.pendingEvents?.length > 0) {
|
||||
console.log(`[ResumableSSE] Replaying ${data.pendingEvents.length} pending events`);
|
||||
logger.log('ResumableSSE', `Replaying ${data.pendingEvents.length} pending events`);
|
||||
for (const pendingEvent of data.pendingEvents) {
|
||||
if (pendingEvent.event === 'title') {
|
||||
titleHandler(pendingEvent);
|
||||
|
|
@ -827,7 +829,7 @@ export default function useResumableSSE(
|
|||
messageHandler(text, { ...currentSubmission, userMessage, initialResponse });
|
||||
}
|
||||
} catch (error) {
|
||||
console.error('[ResumableSSE] Error processing message:', error);
|
||||
logger.error('ResumableSSE', 'Error processing message:', error);
|
||||
}
|
||||
});
|
||||
|
||||
|
|
@ -848,7 +850,7 @@ export default function useResumableSSE(
|
|||
// Invalidate cache once so react-query refetches instead of showing an error.
|
||||
if (responseCode === 404) {
|
||||
const convoId = currentSubmission.conversation?.conversationId;
|
||||
console.log('[ResumableSSE] Stream 404, invalidating messages for:', convoId);
|
||||
logger.log('ResumableSSE', 'Stream 404, invalidating messages for:', convoId);
|
||||
sse.close();
|
||||
removeActiveJob(currentStreamId);
|
||||
/** Terminal: drop any in-flight live estimate so the gauge doesn't
|
||||
|
|
@ -893,7 +895,7 @@ export default function useResumableSSE(
|
|||
sse.stream();
|
||||
return;
|
||||
} catch (error) {
|
||||
console.log('[ResumableSSE] Token refresh failed:', error);
|
||||
logger.log('ResumableSSE', 'Token refresh failed:', error);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -905,7 +907,7 @@ export default function useResumableSSE(
|
|||
* not a server-sent error payload. Use `== null` to only match undefined/null (no HTTP status).
|
||||
*/
|
||||
if (responseCode == null && e.data) {
|
||||
console.log('[ResumableSSE] Server-sent error event received:', e.data);
|
||||
logger.log('ResumableSSE', 'Server-sent error event received:', e.data);
|
||||
sse.close();
|
||||
removeActiveJob(currentStreamId);
|
||||
resetLive({ ...currentSubmission, userMessage });
|
||||
|
|
@ -936,7 +938,7 @@ export default function useResumableSSE(
|
|||
// Not JSON or parsing failed - treat as generic error
|
||||
}
|
||||
|
||||
console.log('[ResumableSSE] Error type check:', { isKnownError, errorString });
|
||||
logger.log('ResumableSSE', 'Error type check:', { isKnownError, errorString });
|
||||
|
||||
// Display the error to user via errorHandler
|
||||
errorHandler({
|
||||
|
|
@ -944,7 +946,7 @@ export default function useResumableSSE(
|
|||
submission: currentSubmission as EventSubmission,
|
||||
});
|
||||
} catch (parseError) {
|
||||
console.error('[ResumableSSE] Failed to parse server error:', parseError);
|
||||
logger.error('ResumableSSE', 'Failed to parse server error:', parseError);
|
||||
errorHandler({
|
||||
data: { text: e.data } as unknown as Parameters<typeof errorHandler>[0]['data'],
|
||||
submission: currentSubmission as EventSubmission,
|
||||
|
|
@ -961,7 +963,7 @@ export default function useResumableSSE(
|
|||
}
|
||||
|
||||
// Network failure or unknown HTTP error - attempt reconnection with backoff
|
||||
console.log('[ResumableSSE] Stream error (network failure) - will attempt reconnect', {
|
||||
logger.log('ResumableSSE', 'Stream error (network failure) - will attempt reconnect', {
|
||||
responseCode,
|
||||
hasData: !!e.data,
|
||||
});
|
||||
|
|
@ -971,8 +973,9 @@ export default function useResumableSSE(
|
|||
reconnectAttemptRef.current++;
|
||||
const delay = Math.min(1000 * Math.pow(2, reconnectAttemptRef.current - 1), 30000);
|
||||
|
||||
console.log(
|
||||
`[ResumableSSE] Reconnecting in ${delay}ms (attempt ${reconnectAttemptRef.current}/${MAX_RETRIES})`,
|
||||
logger.log(
|
||||
'ResumableSSE',
|
||||
`Reconnecting in ${delay}ms (attempt ${reconnectAttemptRef.current}/${MAX_RETRIES})`,
|
||||
);
|
||||
|
||||
sse.close();
|
||||
|
|
@ -989,7 +992,7 @@ export default function useResumableSSE(
|
|||
setIsSubmitting(true);
|
||||
setShowStopButton(true);
|
||||
} else {
|
||||
console.error('[ResumableSSE] Max reconnect attempts reached');
|
||||
logger.error('ResumableSSE', 'Max reconnect attempts reached');
|
||||
sse.close();
|
||||
errorHandler({ data: undefined, submission: currentSubmission as EventSubmission });
|
||||
/** Terminal: clear the in-flight live estimate like the other
|
||||
|
|
@ -1020,11 +1023,11 @@ export default function useResumableSSE(
|
|||
// If we're in a reconnection cycle, don't reset state
|
||||
// (error handler will set up the reconnect timeout)
|
||||
if (reconnectAttemptRef.current > 0) {
|
||||
console.log('[ResumableSSE] Stream closed for reconnect - preserving state');
|
||||
logger.log('ResumableSSE', 'Stream closed for reconnect - preserving state');
|
||||
return;
|
||||
}
|
||||
|
||||
console.log('[ResumableSSE] Stream aborted (intentional close) - no reconnect');
|
||||
logger.log('ResumableSSE', 'Stream aborted (intentional close) - no reconnect');
|
||||
// Clear any pending reconnect attempts
|
||||
if (reconnectTimeoutRef.current) {
|
||||
clearTimeout(reconnectTimeoutRef.current);
|
||||
|
|
@ -1055,14 +1058,14 @@ export default function useResumableSSE(
|
|||
|
||||
/** Simulate network drop - triggers error event → reconnection */
|
||||
debugWindow.__killNetwork = () => {
|
||||
console.log('[Debug] Simulating network drop...');
|
||||
logger.log('Debug', 'Simulating network drop...');
|
||||
// @ts-ignore - sse.js types are incorrect, dispatchEvent actually takes Event
|
||||
sse.dispatchEvent(new Event('error'));
|
||||
};
|
||||
|
||||
/** Simulate clean close (navigation away) - triggers abort event → no reconnection */
|
||||
debugWindow.__closeClean = () => {
|
||||
console.log('[Debug] Simulating clean close (navigation away)...');
|
||||
logger.log('Debug', 'Simulating clean close (navigation away)...');
|
||||
sse.close();
|
||||
};
|
||||
}
|
||||
|
|
@ -1131,7 +1134,7 @@ export default function useResumableSSE(
|
|||
if (signal?.aborted) {
|
||||
return null;
|
||||
}
|
||||
console.log('[ResumableSSE] Generation started:', { streamId: data.streamId });
|
||||
logger.log('ResumableSSE', 'Generation started:', { streamId: data.streamId });
|
||||
return data.streamId;
|
||||
} catch (error) {
|
||||
if (signal?.aborted) {
|
||||
|
|
@ -1158,8 +1161,9 @@ export default function useResumableSSE(
|
|||
const limit = isServerNotReady
|
||||
? `${Math.ceil(START_GENERATION_READINESS_TIMEOUT_MS / 1000)}s readiness window`
|
||||
: `${START_GENERATION_NETWORK_RETRIES}`;
|
||||
console.log(
|
||||
`[ResumableSSE] ${reason} starting generation, retrying in ${retryDelay}ms (attempt ${attempt}/${limit})`,
|
||||
logger.log(
|
||||
'ResumableSSE',
|
||||
`${reason} starting generation, retrying in ${retryDelay}ms (attempt ${attempt}/${limit})`,
|
||||
);
|
||||
const shouldContinue = await waitForRetryDelay(retryDelay, signal);
|
||||
if (!shouldContinue) {
|
||||
|
|
@ -1177,7 +1181,7 @@ export default function useResumableSSE(
|
|||
return null;
|
||||
}
|
||||
|
||||
console.error('[ResumableSSE] Error starting generation:', lastError);
|
||||
logger.error('ResumableSSE', 'Error starting generation:', lastError);
|
||||
|
||||
const axiosError = lastError as { response?: { data?: Record<string, unknown> } };
|
||||
const errorData = axiosError?.response?.data;
|
||||
|
|
@ -1195,7 +1199,7 @@ export default function useResumableSSE(
|
|||
|
||||
useEffect(() => {
|
||||
if (!submission || Object.keys(submission).length === 0) {
|
||||
console.log('[ResumableSSE] No submission, cleaning up');
|
||||
logger.log('ResumableSSE', 'No submission, cleaning up');
|
||||
// Clear reconnect timeout if submission is cleared
|
||||
if (reconnectTimeoutRef.current) {
|
||||
clearTimeout(reconnectTimeoutRef.current);
|
||||
|
|
@ -1213,7 +1217,7 @@ export default function useResumableSSE(
|
|||
}
|
||||
|
||||
const resumeStreamId = (submission as TSubmission & { resumeStreamId?: string }).resumeStreamId;
|
||||
console.log('[ResumableSSE] Effect triggered', {
|
||||
logger.log('ResumableSSE', 'Effect triggered', {
|
||||
conversationId: submission.conversation?.conversationId,
|
||||
hasResumeStreamId: !!resumeStreamId,
|
||||
resumeStreamId,
|
||||
|
|
@ -1237,14 +1241,14 @@ export default function useResumableSSE(
|
|||
return;
|
||||
}
|
||||
// Resume: just subscribe to existing stream, don't start new generation
|
||||
console.log('[ResumableSSE] Resuming existing stream:', resumeStreamId);
|
||||
logger.log('ResumableSSE', 'Resuming existing stream:', resumeStreamId);
|
||||
setStreamId(resumeStreamId);
|
||||
// Optimistically add to active jobs (in case it's not already there)
|
||||
addActiveJob(resumeStreamId);
|
||||
subscribeToStream(resumeStreamId, submission, true); // isResume=true
|
||||
} else {
|
||||
// New generation: start and then subscribe
|
||||
console.log('[ResumableSSE] Starting NEW generation');
|
||||
logger.log('ResumableSSE', 'Starting NEW generation');
|
||||
const newStreamId = await startGeneration(submission, signal);
|
||||
if (signal.aborted) {
|
||||
return;
|
||||
|
|
@ -1268,7 +1272,7 @@ export default function useResumableSSE(
|
|||
submissionRef.current = streamSubmission;
|
||||
subscribeToStream(newStreamId, streamSubmission);
|
||||
} else {
|
||||
console.error('[ResumableSSE] Failed to get streamId from startGeneration');
|
||||
logger.error('ResumableSSE', 'Failed to get streamId from startGeneration');
|
||||
}
|
||||
}
|
||||
};
|
||||
|
|
@ -1276,7 +1280,7 @@ export default function useResumableSSE(
|
|||
initStream();
|
||||
|
||||
return () => {
|
||||
console.log('[ResumableSSE] Cleanup - closing SSE, resetting UI state');
|
||||
logger.log('ResumableSSE', 'Cleanup - closing SSE, resetting UI state');
|
||||
startController.abort();
|
||||
// Cleanup on unmount/navigation - close connection but DO NOT abort backend
|
||||
// Reset UI state so it doesn't leak to other conversations
|
||||
|
|
|
|||
|
|
@ -949,6 +949,7 @@
|
|||
"com_ui_context_cost": "Cost",
|
||||
"com_ui_context_cost_branch": "Cost (this branch)",
|
||||
"com_ui_context_cost_total": "All branches",
|
||||
"com_ui_context_estimated": "Estimated",
|
||||
"com_ui_context_filter_sort": "Filter and Sort by Context",
|
||||
"com_ui_context_free": "Free space",
|
||||
"com_ui_context_messages": "Messages",
|
||||
|
|
|
|||
|
|
@ -88,6 +88,136 @@ describe('token index', () => {
|
|||
expect(altTotals.output).toBe(1019);
|
||||
});
|
||||
|
||||
it('estimates count-less messages by text length without inflating counted totals', () => {
|
||||
buildIndex(CONVO, [
|
||||
msg('u1', Constants.NO_PARENT, true, 12),
|
||||
/** Imported message with no `tokenCount`: 40 chars of text → ~10 est tokens. */
|
||||
{
|
||||
messageId: 'a1',
|
||||
parentMessageId: 'u1',
|
||||
isCreatedByUser: false,
|
||||
conversationId: CONVO,
|
||||
text: 'x'.repeat(40),
|
||||
} as TMessage,
|
||||
]);
|
||||
|
||||
const totals = sumBranch(CONVO, 'a1');
|
||||
/** Known counts feed input/output/counted; the count-less message stays out
|
||||
* of those and lands in the separate (uncalibrated) estimate bucket. */
|
||||
expect(totals.input).toBe(12);
|
||||
expect(totals.output).toBe(0);
|
||||
expect(totals.counted).toBe(1);
|
||||
expect(totals.total).toBe(2);
|
||||
expect(totals.estTokens).toBe(10);
|
||||
});
|
||||
|
||||
it('estimates object-form content text and merged quote excerpts', () => {
|
||||
buildIndex(CONVO, [
|
||||
/** Assistant body lives only in object-form content (`text.value`). */
|
||||
{
|
||||
messageId: 'a1',
|
||||
parentMessageId: Constants.NO_PARENT,
|
||||
isCreatedByUser: false,
|
||||
conversationId: CONVO,
|
||||
content: [{ type: 'text', text: { value: 'y'.repeat(20) } }],
|
||||
} as unknown as TMessage,
|
||||
/** User turn whose quotes are merged into the prompt at send time. */
|
||||
{
|
||||
messageId: 'u1',
|
||||
parentMessageId: 'a1',
|
||||
isCreatedByUser: true,
|
||||
conversationId: CONVO,
|
||||
text: 'z'.repeat(16),
|
||||
quotes: ['q'.repeat(8)],
|
||||
} as TMessage,
|
||||
]);
|
||||
|
||||
/** a1: 20 content chars / 4 = 5; u1: (16 text + 8 quote) / 4 = 6. */
|
||||
const totals = sumBranch(CONVO, 'u1');
|
||||
expect(totals.counted).toBe(0);
|
||||
expect(totals.estTokens).toBe(11);
|
||||
});
|
||||
|
||||
it('recounts quoted user turns (ignoring stale counts), counts tool calls, skips reasoning', () => {
|
||||
buildIndex(CONVO, [
|
||||
/** Quoted user turn with a stale text-only stored count: the send path
|
||||
* recounts the merged prompt every turn, so the estimate ignores the count
|
||||
* and recounts from text+quotes. */
|
||||
{
|
||||
messageId: 'u1',
|
||||
parentMessageId: Constants.NO_PARENT,
|
||||
isCreatedByUser: true,
|
||||
conversationId: CONVO,
|
||||
tokenCount: 999,
|
||||
text: 'hi',
|
||||
quotes: ['q'.repeat(38)],
|
||||
} as TMessage,
|
||||
/** Count-less assistant turn: tool-call name/args/output count toward the
|
||||
* estimate (sent back as context); reasoning does not. */
|
||||
{
|
||||
messageId: 'a1',
|
||||
parentMessageId: 'u1',
|
||||
isCreatedByUser: false,
|
||||
conversationId: CONVO,
|
||||
content: [
|
||||
{ type: 'think', think: 'r'.repeat(40) },
|
||||
{ type: 'tool_call', tool_call: { name: 'sub', args: 'aa', output: 'o'.repeat(11) } },
|
||||
],
|
||||
} as unknown as TMessage,
|
||||
]);
|
||||
|
||||
const totals = sumBranch(CONVO, 'a1');
|
||||
/** u1 quoted: stored 999 ignored; (2 text + 38 quote) / 4 = 10. a1 tool_call
|
||||
* name 3 + args 2 + output 11 = 16 / 4 = 4 (think skipped). */
|
||||
expect(totals.input).toBe(0);
|
||||
expect(totals.counted).toBe(0);
|
||||
expect(totals.estTokens).toBe(14);
|
||||
});
|
||||
|
||||
it('prefers content over text for count-less messages carrying both', () => {
|
||||
buildIndex(CONVO, [
|
||||
msg('u1', Constants.NO_PARENT, true, 8),
|
||||
/** Stopped agent response: saved with both a short `text` and structured
|
||||
* `content` (a tool call). The send path formats from content, so the
|
||||
* estimate must use content (tool tokens), not the shorter text. */
|
||||
{
|
||||
messageId: 'a1',
|
||||
parentMessageId: 'u1',
|
||||
isCreatedByUser: false,
|
||||
conversationId: CONVO,
|
||||
text: 'hi',
|
||||
content: [
|
||||
{ type: 'tool_call', tool_call: { name: 'run', args: 'aa', output: 'o'.repeat(13) } },
|
||||
],
|
||||
} as unknown as TMessage,
|
||||
]);
|
||||
|
||||
const totals = sumBranch(CONVO, 'a1');
|
||||
/** a1 uses content (name 3 + args 2 + output 13 = 18 / 4 = 5), not text 'hi'. */
|
||||
expect(totals.input).toBe(8);
|
||||
expect(totals.estTokens).toBe(5);
|
||||
});
|
||||
|
||||
it('exposes the count-less tail estimate so live output is not double-counted', () => {
|
||||
buildIndex(CONVO, [
|
||||
msg('u1', Constants.NO_PARENT, true, 12),
|
||||
/** In-flight / resumed response: count-less, so it lands in estTokens; it is
|
||||
* also covered by liveTokens, so the estimate path drops tailEstTokens. */
|
||||
{
|
||||
messageId: 'a1',
|
||||
parentMessageId: 'u1',
|
||||
isCreatedByUser: false,
|
||||
conversationId: CONVO,
|
||||
text: 'o'.repeat(20),
|
||||
} as TMessage,
|
||||
]);
|
||||
|
||||
const totals = sumBranch(CONVO, 'a1');
|
||||
/** a1 is the tail: 20 / 4 = 5, surfaced both in estTokens and tailEstTokens. */
|
||||
expect(totals.estTokens).toBe(5);
|
||||
expect(totals.tailEstTokens).toBe(5);
|
||||
});
|
||||
|
||||
it('caps the branch at a summary marker instead of re-summing compacted history', () => {
|
||||
const summarized = {
|
||||
messageId: 'a2',
|
||||
|
|
|
|||
|
|
@ -27,6 +27,11 @@ export const EMPTY_USAGE: BranchUsage = {
|
|||
|
||||
export interface TokenEntry {
|
||||
tokenCount: number;
|
||||
/** Char/4 token estimate used in place of `tokenCount`: count-less (imported /
|
||||
* pre-feature) message bodies, plus quoted user turns (recounted from merged
|
||||
* text+quotes, since the send path ignores their stored count). Includes
|
||||
* tool-call name/args/output. Mutually exclusive with a counted `tokenCount`. */
|
||||
estTokens: number;
|
||||
isCreatedByUser: boolean;
|
||||
parentMessageId: string | null;
|
||||
/** Per-response provider usage from `metadata.usage` (response messages only) */
|
||||
|
|
@ -46,6 +51,13 @@ export interface BranchTotals {
|
|||
counted: number;
|
||||
/** Total messages on the branch */
|
||||
total: number;
|
||||
/** Uncalibrated estimate sum for count-less branch messages (imports /
|
||||
* pre-feature). Kept separate so known counts aren't re-estimated. */
|
||||
estTokens: number;
|
||||
/** The tail (latest) message's own `estTokens`. When a stream is live the tail
|
||||
* is the in-flight response, already covered by `liveTokens`, so the estimate
|
||||
* path excludes this to avoid double-counting a resumed/partial response. */
|
||||
tailEstTokens: number;
|
||||
tailId: string | null;
|
||||
/** Whether the latest run's anchor message is on this branch */
|
||||
containsAnchor: boolean;
|
||||
|
|
@ -63,6 +75,8 @@ export const EMPTY_BRANCH: BranchTotals = {
|
|||
output: 0,
|
||||
counted: 0,
|
||||
total: 0,
|
||||
estTokens: 0,
|
||||
tailEstTokens: 0,
|
||||
tailId: null,
|
||||
containsAnchor: false,
|
||||
usage: EMPTY_USAGE,
|
||||
|
|
@ -137,11 +151,101 @@ function addUsage(target: BranchUsage, usage?: BranchUsage): void {
|
|||
}
|
||||
}
|
||||
|
||||
/** Chars of a content part's text, handling both the string and `{ value }` forms.
|
||||
* Reasoning (`think`) and error parts are excluded — the send path strips them
|
||||
* before counting, so they aren't part of the next call's context. */
|
||||
function partTextChars(part: unknown): number {
|
||||
if (part == null || typeof part !== 'object') {
|
||||
return 0;
|
||||
}
|
||||
const type = (part as { type?: unknown }).type;
|
||||
if (type === 'think' || type === 'error') {
|
||||
return 0;
|
||||
}
|
||||
if (type === 'tool_call') {
|
||||
const call = (part as { tool_call?: { name?: unknown; args?: unknown; output?: unknown } })
|
||||
.tool_call;
|
||||
if (call == null) {
|
||||
return 0;
|
||||
}
|
||||
let chars = typeof call.name === 'string' ? call.name.length : 0;
|
||||
if (typeof call.args === 'string') {
|
||||
chars += call.args.length;
|
||||
} else if (call.args != null) {
|
||||
chars += JSON.stringify(call.args).length;
|
||||
}
|
||||
if (typeof call.output === 'string') {
|
||||
chars += call.output.length;
|
||||
}
|
||||
return chars;
|
||||
}
|
||||
const text = (part as { text?: unknown }).text;
|
||||
if (typeof text === 'string') {
|
||||
return text.length;
|
||||
}
|
||||
if (
|
||||
text != null &&
|
||||
typeof text === 'object' &&
|
||||
typeof (text as { value?: unknown }).value === 'string'
|
||||
) {
|
||||
return (text as { value: string }).value.length;
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
/** Char length of a message's rendered text, for estimating count-less messages.
|
||||
* Prefer structured `content` when present — the send path formats from it (incl.
|
||||
* tool calls), so a message carrying both `text` and `content` (e.g. a stopped
|
||||
* agent response) would otherwise drop its content/tool-call tokens. */
|
||||
function messageChars(message: Partial<TMessage>): number {
|
||||
if (Array.isArray(message.content) && message.content.length > 0) {
|
||||
let chars = 0;
|
||||
for (const part of message.content) {
|
||||
chars += partTextChars(part);
|
||||
}
|
||||
return chars;
|
||||
}
|
||||
if (typeof message.text === 'string') {
|
||||
return message.text.length;
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
/** Quoted excerpts the send path merges into a user message's prompt. */
|
||||
function quoteChars(message: Partial<TMessage>): number {
|
||||
if (!Array.isArray(message.quotes)) {
|
||||
return 0;
|
||||
}
|
||||
let chars = 0;
|
||||
for (const quote of message.quotes) {
|
||||
if (typeof quote === 'string') {
|
||||
chars += quote.length;
|
||||
}
|
||||
}
|
||||
return chars;
|
||||
}
|
||||
|
||||
function toEntry(message: Partial<TMessage>): TokenEntry {
|
||||
const summaryUsedTokens = message.metadata?.summaryUsedTokens;
|
||||
const tokenCount = typeof message.tokenCount === 'number' ? message.tokenCount : 0;
|
||||
const isCreatedByUser = message.isCreatedByUser === true;
|
||||
const quoted = isCreatedByUser && Array.isArray(message.quotes) && message.quotes.length > 0;
|
||||
/** A quoted user turn's stored `tokenCount` is unreliable: a text-only Save edit
|
||||
* recomputes it from `text` alone, and the send path recounts the quote-merged
|
||||
* prompt every turn regardless (`needsCanonicalTokenCount` in agents/client.js).
|
||||
* So mirror the server — estimate quoted turns from the merged text+quotes and
|
||||
* ignore the stored count. Other count-less imports / pre-feature messages
|
||||
* estimate from text. */
|
||||
let estTokens = 0;
|
||||
if (quoted) {
|
||||
estTokens = Math.round((messageChars(message) + quoteChars(message)) / 4);
|
||||
} else if (tokenCount === 0) {
|
||||
estTokens = Math.round(messageChars(message) / 4);
|
||||
}
|
||||
return {
|
||||
tokenCount: typeof message.tokenCount === 'number' ? message.tokenCount : 0,
|
||||
isCreatedByUser: message.isCreatedByUser === true,
|
||||
tokenCount: quoted ? 0 : tokenCount,
|
||||
estTokens,
|
||||
isCreatedByUser,
|
||||
parentMessageId: message.parentMessageId ?? null,
|
||||
usage: readPersistedUsage(message),
|
||||
summaryUsedTokens:
|
||||
|
|
@ -231,7 +335,10 @@ export function sumBranch(
|
|||
return EMPTY_BRANCH;
|
||||
}
|
||||
|
||||
const totals = { input: 0, output: 0, counted: 0, total: 0, containsAnchor: false };
|
||||
const totals = { input: 0, output: 0, counted: 0, total: 0, estTokens: 0, containsAnchor: false };
|
||||
/** The in-flight response, when streaming, is the branch tail and is covered by
|
||||
* `liveTokens`; expose its estimate so the estimate path can drop it. */
|
||||
const tailEstTokens = index.get(tailId)?.estTokens ?? 0;
|
||||
const usage: BranchUsage = { ...EMPTY_USAGE };
|
||||
let summaryBaseline = 0;
|
||||
/** Once a summary marker is crossed, older turns are out of the CONTEXT WINDOW
|
||||
|
|
@ -262,6 +369,8 @@ export function sumBranch(
|
|||
} else {
|
||||
totals.output += entry.tokenCount;
|
||||
}
|
||||
} else if (!contextCapped && entry.estTokens > 0) {
|
||||
totals.estTokens += entry.estTokens;
|
||||
}
|
||||
/** Cost/usage is cumulative spend — never truncated at the summary boundary. */
|
||||
addUsage(usage, entry.usage);
|
||||
|
|
@ -275,7 +384,7 @@ export function sumBranch(
|
|||
currentId = entry.parentMessageId;
|
||||
}
|
||||
|
||||
return { ...totals, tailId, usage, summaryBaseline };
|
||||
return { ...totals, tailEstTokens, tailId, usage, summaryBaseline };
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
|
|||
|
|
@ -6,5 +6,4 @@ export * from './google';
|
|||
export * from './models';
|
||||
export * from './openai';
|
||||
export * from './pricing';
|
||||
export * from './projection';
|
||||
export * from './tokenConfig';
|
||||
|
|
|
|||
|
|
@ -1,206 +0,0 @@
|
|||
import { resolveContextProjection } from './projection';
|
||||
import { QUOTE_MAX_COUNT } from '~/utils/quotes';
|
||||
|
||||
jest.mock('@librechat/agents', () => ({
|
||||
Providers: { OPENAI: 'openai' },
|
||||
createTokenCounter: jest.fn(async () => jest.fn(() => 1)),
|
||||
projectAgentContextUsage: jest.fn(() => ({ tokenCount: 1, maxContextTokens: 1000 })),
|
||||
}));
|
||||
|
||||
const GRAPH_SELECT = 'messageId parentMessageId metadata.summaryUsedTokens';
|
||||
const BODY_SELECT = 'messageId parentMessageId tokenCount isCreatedByUser text quotes';
|
||||
|
||||
function textStats(messageId: string, textBytes = 5) {
|
||||
return {
|
||||
messageId,
|
||||
textBytes,
|
||||
quoteCount: 0,
|
||||
quoteBytes: 0,
|
||||
quoteLineCount: 0,
|
||||
nonStringQuoteCount: 0,
|
||||
};
|
||||
}
|
||||
|
||||
describe('resolveContextProjection', () => {
|
||||
const baseParams = {
|
||||
conversationId: 'conversation-1',
|
||||
messageId: 'message-1',
|
||||
endpoint: 'openai',
|
||||
maxContextTokens: 1000,
|
||||
model: 'gpt-4o',
|
||||
};
|
||||
|
||||
beforeEach(() => {
|
||||
jest.clearAllMocks();
|
||||
});
|
||||
|
||||
it('returns null before tokenization when the conversation is too large', async () => {
|
||||
const { createTokenCounter } = jest.requireMock('@librechat/agents');
|
||||
const messages = Array.from({ length: 513 }, (_, index) => ({
|
||||
messageId: `message-${index}`,
|
||||
parentMessageId: index === 0 ? null : `message-${index - 1}`,
|
||||
isCreatedByUser: true,
|
||||
text: 'hello',
|
||||
}));
|
||||
const getMessages = jest.fn(async () => messages);
|
||||
const getMessageTextStats = jest.fn();
|
||||
|
||||
const result = await resolveContextProjection(
|
||||
{ userId: 'user-1', getMessages, getMessageTextStats },
|
||||
{ ...baseParams, messageId: 'message-512' },
|
||||
);
|
||||
|
||||
expect(result).toBeNull();
|
||||
expect(getMessages).toHaveBeenCalledTimes(1);
|
||||
expect(getMessages).toHaveBeenCalledWith(
|
||||
{ conversationId: 'conversation-1', user: 'user-1' },
|
||||
GRAPH_SELECT,
|
||||
{ limit: 513, sort: false },
|
||||
);
|
||||
expect(getMessageTextStats).not.toHaveBeenCalled();
|
||||
expect(createTokenCounter).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('returns null before tokenization when the branch is too long', async () => {
|
||||
const { createTokenCounter } = jest.requireMock('@librechat/agents');
|
||||
const messages = Array.from({ length: 257 }, (_, index) => ({
|
||||
messageId: `message-${index}`,
|
||||
parentMessageId: index === 0 ? null : `message-${index - 1}`,
|
||||
isCreatedByUser: true,
|
||||
text: 'hello',
|
||||
}));
|
||||
const getMessages = jest.fn(async () => messages);
|
||||
const getMessageTextStats = jest.fn();
|
||||
|
||||
const result = await resolveContextProjection(
|
||||
{ userId: 'user-1', getMessages, getMessageTextStats },
|
||||
{ ...baseParams, messageId: 'message-256' },
|
||||
);
|
||||
|
||||
expect(result).toBeNull();
|
||||
expect(getMessages).toHaveBeenCalledTimes(1);
|
||||
expect(getMessageTextStats).not.toHaveBeenCalled();
|
||||
expect(createTokenCounter).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('returns null before loading bodies when the branch text is too large', async () => {
|
||||
const { createTokenCounter } = jest.requireMock('@librechat/agents');
|
||||
const getMessages = jest.fn(async () => [
|
||||
{
|
||||
messageId: 'message-1',
|
||||
parentMessageId: null,
|
||||
},
|
||||
]);
|
||||
const getMessageTextStats = jest.fn(async () => [textStats('message-1', 512 * 1024 + 1)]);
|
||||
const result = await resolveContextProjection(
|
||||
{
|
||||
userId: 'user-1',
|
||||
getMessages,
|
||||
getMessageTextStats,
|
||||
},
|
||||
baseParams,
|
||||
);
|
||||
|
||||
expect(result).toBeNull();
|
||||
expect(getMessages).toHaveBeenCalledTimes(1);
|
||||
expect(getMessageTextStats).toHaveBeenCalledWith(
|
||||
{
|
||||
conversationId: 'conversation-1',
|
||||
user: 'user-1',
|
||||
messageId: { $in: ['message-1'] },
|
||||
},
|
||||
{ limit: 1 },
|
||||
);
|
||||
expect(createTokenCounter).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('loads only branch message bodies after resolving the graph', async () => {
|
||||
const graph = [
|
||||
{ messageId: 'message-1', parentMessageId: null },
|
||||
{ messageId: 'message-2', parentMessageId: 'message-1' },
|
||||
{ messageId: 'off-branch', parentMessageId: null },
|
||||
];
|
||||
const bodies = [
|
||||
{
|
||||
messageId: 'message-1',
|
||||
parentMessageId: null,
|
||||
isCreatedByUser: true,
|
||||
text: 'first',
|
||||
tokenCount: 5,
|
||||
},
|
||||
{
|
||||
messageId: 'message-2',
|
||||
parentMessageId: 'message-1',
|
||||
isCreatedByUser: false,
|
||||
text: 'second',
|
||||
tokenCount: 6,
|
||||
},
|
||||
];
|
||||
const getMessages = jest.fn(async (_filter: object, select?: string) =>
|
||||
select === GRAPH_SELECT ? graph : bodies,
|
||||
);
|
||||
const getMessageTextStats = jest.fn(async () => [
|
||||
textStats('message-1', 5),
|
||||
textStats('message-2', 6),
|
||||
]);
|
||||
|
||||
const result = await resolveContextProjection(
|
||||
{ userId: 'user-1', getMessages, getMessageTextStats },
|
||||
{ ...baseParams, messageId: 'message-2' },
|
||||
);
|
||||
|
||||
expect(result).toEqual({ tokenCount: 1, maxContextTokens: 1000 });
|
||||
expect(getMessages).toHaveBeenNthCalledWith(
|
||||
1,
|
||||
{ conversationId: 'conversation-1', user: 'user-1' },
|
||||
GRAPH_SELECT,
|
||||
{ limit: 513, sort: false },
|
||||
);
|
||||
expect(getMessageTextStats).toHaveBeenCalledWith(
|
||||
{
|
||||
conversationId: 'conversation-1',
|
||||
user: 'user-1',
|
||||
messageId: { $in: ['message-1', 'message-2'] },
|
||||
},
|
||||
{ limit: 2 },
|
||||
);
|
||||
expect(getMessages).toHaveBeenNthCalledWith(
|
||||
2,
|
||||
{
|
||||
conversationId: 'conversation-1',
|
||||
user: 'user-1',
|
||||
messageId: { $in: ['message-1', 'message-2'] },
|
||||
},
|
||||
BODY_SELECT,
|
||||
{ limit: 2, sort: false },
|
||||
);
|
||||
});
|
||||
|
||||
it('returns null before loading bodies when a branch message has too many quotes', async () => {
|
||||
const { createTokenCounter } = jest.requireMock('@librechat/agents');
|
||||
const getMessages = jest.fn(async () => [
|
||||
{
|
||||
messageId: 'message-1',
|
||||
parentMessageId: null,
|
||||
},
|
||||
]);
|
||||
const getMessageTextStats = jest.fn(async () => [
|
||||
{
|
||||
...textStats('message-1'),
|
||||
quoteCount: QUOTE_MAX_COUNT + 1,
|
||||
quoteBytes: 10,
|
||||
quoteLineCount: QUOTE_MAX_COUNT + 1,
|
||||
},
|
||||
]);
|
||||
|
||||
const result = await resolveContextProjection(
|
||||
{ userId: 'user-1', getMessages, getMessageTextStats },
|
||||
baseParams,
|
||||
);
|
||||
|
||||
expect(result).toBeNull();
|
||||
expect(getMessages).toHaveBeenCalledTimes(1);
|
||||
expect(getMessageTextStats).toHaveBeenCalledTimes(1);
|
||||
expect(createTokenCounter).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
|
|
@ -1,326 +0,0 @@
|
|||
import { HumanMessage, AIMessage } from '@langchain/core/messages';
|
||||
import { Providers, createTokenCounter, projectAgentContextUsage } from '@librechat/agents';
|
||||
import type { TContextProjectionRequest, TContextUsageEvent } from 'librechat-data-provider';
|
||||
import type { BaseMessage } from '@langchain/core/messages';
|
||||
import { QUOTE_MAX_COUNT, mergeQuotedText } from '~/utils/quotes';
|
||||
|
||||
const MAX_PROJECTION_MESSAGES = 512;
|
||||
const MAX_PROJECTION_BRANCH_MESSAGES = 256;
|
||||
const MAX_PROJECTION_BRANCH_TEXT_BYTES = 512 * 1024;
|
||||
const PROJECTION_GRAPH_SELECT = 'messageId parentMessageId metadata.summaryUsedTokens';
|
||||
const PROJECTION_BODY_SELECT = 'messageId parentMessageId tokenCount isCreatedByUser text quotes';
|
||||
|
||||
interface ProjectionMessage {
|
||||
messageId: string;
|
||||
parentMessageId?: string | null;
|
||||
tokenCount?: number;
|
||||
isCreatedByUser?: boolean;
|
||||
text?: string;
|
||||
/** Quoted excerpts merged into the model-facing text by the live path; must be
|
||||
* included here so the context gauge counts the same prompt the model sees. */
|
||||
quotes?: string[];
|
||||
/** Compaction marker written by the live path (`agents/usage.ts`); its
|
||||
* presence means the next call sends the summary + tail, not this raw chain. */
|
||||
metadata?: { summaryUsedTokens?: number };
|
||||
}
|
||||
|
||||
interface ProjectionMessageFilter {
|
||||
conversationId: string;
|
||||
user?: string;
|
||||
messageId?: string | { $in: string[] };
|
||||
}
|
||||
|
||||
interface ProjectionMessageQueryOptions {
|
||||
limit?: number;
|
||||
sort?: false;
|
||||
}
|
||||
|
||||
interface ProjectionMessageTextStats {
|
||||
messageId: string;
|
||||
textBytes: number;
|
||||
quoteCount: number;
|
||||
quoteBytes: number;
|
||||
quoteLineCount: number;
|
||||
nonStringQuoteCount: number;
|
||||
}
|
||||
|
||||
interface ProjectionMessageTextStatsOptions {
|
||||
limit?: number;
|
||||
}
|
||||
|
||||
export interface ContextProjectionDeps {
|
||||
/** Authenticated requester — branch lookups are scoped to this user. */
|
||||
userId?: string;
|
||||
getMessages: (
|
||||
filter: ProjectionMessageFilter,
|
||||
select?: string,
|
||||
options?: ProjectionMessageQueryOptions,
|
||||
) => Promise<ProjectionMessage[]>;
|
||||
getMessageTextStats: (
|
||||
filter: ProjectionMessageFilter,
|
||||
options?: ProjectionMessageTextStatsOptions,
|
||||
) => Promise<ProjectionMessageTextStats[]>;
|
||||
}
|
||||
|
||||
/**
|
||||
* Walks the parent chain from `tailId` to root and returns the branch messages
|
||||
* oldest→newest. The visited set guards against cycles / self-referential links.
|
||||
*/
|
||||
function resolveBranch(messages: ProjectionMessage[], tailId: string): ProjectionMessage[] {
|
||||
const byId = new Map<string, ProjectionMessage>();
|
||||
for (const message of messages) {
|
||||
byId.set(message.messageId, message);
|
||||
}
|
||||
const branch: ProjectionMessage[] = [];
|
||||
const seen = new Set<string>();
|
||||
let currentId: string | null | undefined = tailId;
|
||||
while (currentId != null && !seen.has(currentId)) {
|
||||
const message = byId.get(currentId);
|
||||
if (message == null) {
|
||||
break;
|
||||
}
|
||||
seen.add(currentId);
|
||||
branch.push(message);
|
||||
currentId = message.parentMessageId;
|
||||
}
|
||||
return branch.reverse();
|
||||
}
|
||||
|
||||
function hasValidProjectionIds(params: TContextProjectionRequest): boolean {
|
||||
return typeof params.conversationId === 'string' && typeof params.messageId === 'string';
|
||||
}
|
||||
|
||||
function getProjectionText(message: ProjectionMessage): string | null {
|
||||
const hasQuotes =
|
||||
message.isCreatedByUser === true && Array.isArray(message.quotes) && message.quotes.length > 0;
|
||||
if (!hasQuotes) {
|
||||
return message.text ?? '';
|
||||
}
|
||||
if (message.quotes == null || message.quotes.length > QUOTE_MAX_COUNT) {
|
||||
return null;
|
||||
}
|
||||
for (const quote of message.quotes) {
|
||||
if (typeof quote !== 'string') {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
return mergeQuotedText(message.text ?? '', message.quotes);
|
||||
}
|
||||
|
||||
function hasExceededBranchTextLimit(branch: ProjectionMessage[]): boolean {
|
||||
let bytes = 0;
|
||||
for (const message of branch) {
|
||||
const text = getProjectionText(message);
|
||||
if (text == null) {
|
||||
return true;
|
||||
}
|
||||
bytes += Buffer.byteLength(text, 'utf8');
|
||||
if (bytes > MAX_PROJECTION_BRANCH_TEXT_BYTES) {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
function getEstimatedMergedTextBytes(stats: ProjectionMessageTextStats): number | null {
|
||||
if (
|
||||
stats.nonStringQuoteCount > 0 ||
|
||||
stats.quoteCount > QUOTE_MAX_COUNT ||
|
||||
stats.quoteLineCount < stats.quoteCount
|
||||
) {
|
||||
return null;
|
||||
}
|
||||
if (stats.quoteCount === 0) {
|
||||
return stats.textBytes;
|
||||
}
|
||||
|
||||
const quotePrefixBytes = stats.quoteLineCount * 2;
|
||||
const quoteLineBreakBytes = stats.quoteLineCount - stats.quoteCount;
|
||||
const quoteSeparatorBytes = (stats.quoteCount - 1) * 2;
|
||||
const bodySeparatorBytes = stats.textBytes > 0 ? 2 : 0;
|
||||
return (
|
||||
stats.textBytes +
|
||||
stats.quoteBytes +
|
||||
quotePrefixBytes +
|
||||
quoteLineBreakBytes +
|
||||
quoteSeparatorBytes +
|
||||
bodySeparatorBytes
|
||||
);
|
||||
}
|
||||
|
||||
function hasExceededBranchTextStatsLimit(stats: ProjectionMessageTextStats[]): boolean {
|
||||
let bytes = 0;
|
||||
for (const messageStats of stats) {
|
||||
const messageBytes = getEstimatedMergedTextBytes(messageStats);
|
||||
if (messageBytes == null) {
|
||||
return true;
|
||||
}
|
||||
bytes += messageBytes;
|
||||
if (bytes > MAX_PROJECTION_BRANCH_TEXT_BYTES) {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
/** Maps an endpoint/provider string to the agents `Providers` enum. */
|
||||
function resolveProvider(value?: string): Providers {
|
||||
if (value == null || value === '') {
|
||||
return Providers.OPENAI;
|
||||
}
|
||||
const lower = value.toLowerCase();
|
||||
for (const provider of Object.values(Providers)) {
|
||||
if (provider.toLowerCase() === lower) {
|
||||
return provider;
|
||||
}
|
||||
}
|
||||
if (lower.includes('anthropic') || lower.includes('claude')) {
|
||||
return Providers.ANTHROPIC;
|
||||
}
|
||||
if (lower.includes('google') || lower.includes('gemini') || lower.includes('vertex')) {
|
||||
return Providers.GOOGLE;
|
||||
}
|
||||
if (lower.includes('bedrock')) {
|
||||
return Providers.BEDROCK;
|
||||
}
|
||||
return Providers.OPENAI;
|
||||
}
|
||||
|
||||
async function getBranchMessages(
|
||||
deps: ContextProjectionDeps,
|
||||
baseFilter: ProjectionMessageFilter,
|
||||
branch: ProjectionMessage[],
|
||||
): Promise<ProjectionMessage[] | null> {
|
||||
const branchIds = branch.map((message) => message.messageId);
|
||||
const stats = await deps.getMessageTextStats(
|
||||
{ ...baseFilter, messageId: { $in: branchIds } },
|
||||
{ limit: branchIds.length },
|
||||
);
|
||||
if (stats.length !== branchIds.length || hasExceededBranchTextStatsLimit(stats)) {
|
||||
return null;
|
||||
}
|
||||
|
||||
const stored = await deps.getMessages(
|
||||
{ ...baseFilter, messageId: { $in: branchIds } },
|
||||
PROJECTION_BODY_SELECT,
|
||||
{ limit: branchIds.length, sort: false },
|
||||
);
|
||||
if (stored.length !== branchIds.length) {
|
||||
return null;
|
||||
}
|
||||
const byId = new Map<string, ProjectionMessage>();
|
||||
for (const message of stored) {
|
||||
byId.set(message.messageId, message);
|
||||
}
|
||||
const ordered: ProjectionMessage[] = [];
|
||||
for (const messageId of branchIds) {
|
||||
const message = byId.get(messageId);
|
||||
if (message == null) {
|
||||
return null;
|
||||
}
|
||||
ordered.push(message);
|
||||
}
|
||||
return ordered;
|
||||
}
|
||||
|
||||
/**
|
||||
* Server-side context-usage projection: reconstructs the viewed branch and asks
|
||||
* the agents SDK what the next call's context would be, WITHOUT invoking the
|
||||
* model. Provider/model/window come from the (client-resolved) request — no
|
||||
* agent or model-spec config is loaded here, so there is no cross-user config
|
||||
* exposure. Reuses LibreChat's already-calibrated per-message `tokenCount`s (no
|
||||
* re-tokenizing). Returns null when there is no resolvable context window.
|
||||
* NOTE: this first cut targets message-windowing accuracy — instruction and
|
||||
* tool-schema tokens (agent instructions, `promptPrefix`, model-spec presets,
|
||||
* tool schemas) are NOT yet included; a follow-up will reuse the full
|
||||
* `initializeAgent`/send path for exact overhead and proper access control.
|
||||
*/
|
||||
export async function resolveContextProjection(
|
||||
deps: ContextProjectionDeps,
|
||||
params: TContextProjectionRequest,
|
||||
): Promise<TContextUsageEvent | null> {
|
||||
if (!hasValidProjectionIds(params)) {
|
||||
return null;
|
||||
}
|
||||
|
||||
const maxContextTokens = params.maxContextTokens;
|
||||
if (maxContextTokens == null || maxContextTokens <= 0) {
|
||||
return null;
|
||||
}
|
||||
|
||||
const baseFilter = { conversationId: params.conversationId, user: deps.userId };
|
||||
const stored = await deps.getMessages(baseFilter, PROJECTION_GRAPH_SELECT, {
|
||||
limit: MAX_PROJECTION_MESSAGES + 1,
|
||||
sort: false,
|
||||
});
|
||||
if (stored.length > MAX_PROJECTION_MESSAGES) {
|
||||
return null;
|
||||
}
|
||||
|
||||
const branch = resolveBranch(stored, params.messageId);
|
||||
if (branch.length === 0) {
|
||||
return null;
|
||||
}
|
||||
if (branch.length > MAX_PROJECTION_BRANCH_MESSAGES) {
|
||||
return null;
|
||||
}
|
||||
|
||||
/** A summarized/compacted branch's next call sends the saved summary + the
|
||||
* post-summary tail, NOT this raw parent chain — projecting from the full
|
||||
* history would prune/count the wrong context and omit the summary. Detect it
|
||||
* via the live path's `metadata.summaryUsedTokens` marker and fall back (null)
|
||||
* so the client's summary-baseline-aware estimate handles these branches until
|
||||
* a follow-up replays the summary boundary. */
|
||||
if (branch.some((message) => (message.metadata?.summaryUsedTokens ?? 0) > 0)) {
|
||||
return null;
|
||||
}
|
||||
|
||||
const bodyBranch = await getBranchMessages(deps, baseFilter, branch);
|
||||
if (bodyBranch == null || hasExceededBranchTextLimit(bodyBranch)) {
|
||||
return null;
|
||||
}
|
||||
|
||||
const model = params.model;
|
||||
const encoding = (model ?? '').toLowerCase().includes('claude') ? 'claude' : 'o200k_base';
|
||||
const tokenCounter = await createTokenCounter(encoding);
|
||||
|
||||
const messages: BaseMessage[] = [];
|
||||
const indexTokenCountMap: Record<string, number> = {};
|
||||
for (let i = 0; i < bodyBranch.length; i++) {
|
||||
const message = bodyBranch[i];
|
||||
/** Mirror the live path: prepend quoted excerpts into the user text the model
|
||||
* receives so the gauge counts the same prompt. */
|
||||
const hasQuotes =
|
||||
message.isCreatedByUser === true &&
|
||||
Array.isArray(message.quotes) &&
|
||||
message.quotes.length > 0;
|
||||
const text = getProjectionText(message);
|
||||
if (text == null) {
|
||||
return null;
|
||||
}
|
||||
const lcMessage =
|
||||
message.isCreatedByUser === true ? new HumanMessage(text) : new AIMessage(text);
|
||||
messages.push(lcMessage);
|
||||
/** Recount messages with no stored count (imported / pre-feature) rather
|
||||
* than charging 0 — a real 0 and "unknown" must not collapse, or the
|
||||
* snapshot-less histories this endpoint targets would under-report. Also
|
||||
* recount quoted messages: a text-only Save edit leaves a stale text-only
|
||||
* `tokenCount` that omits the quote block, so trust the merged recount. */
|
||||
indexTokenCountMap[String(i)] =
|
||||
!hasQuotes && message.tokenCount != null && message.tokenCount > 0
|
||||
? message.tokenCount
|
||||
: tokenCounter(lcMessage);
|
||||
}
|
||||
|
||||
return projectAgentContextUsage({
|
||||
agent: {
|
||||
agentId: params.agentId ?? 'projection',
|
||||
provider: resolveProvider(params.endpoint),
|
||||
maxContextTokens,
|
||||
},
|
||||
messages,
|
||||
tokenCounter,
|
||||
indexTokenCountMap,
|
||||
calibrationRatio: params.calibrationRatio,
|
||||
});
|
||||
}
|
||||
|
|
@ -160,8 +160,6 @@ export const aiEndpoints = () => `${BASE_URL}/api/endpoints`;
|
|||
|
||||
export const tokenConfig = () => `${BASE_URL}/api/endpoints/token-config`;
|
||||
|
||||
export const contextProjection = () => `${BASE_URL}/api/endpoints/context-projection`;
|
||||
|
||||
export const models = () => `${BASE_URL}/api/models`;
|
||||
|
||||
export const tokenizer = () => `${BASE_URL}/api/tokenizer`;
|
||||
|
|
|
|||
|
|
@ -1,5 +1,4 @@
|
|||
import type { AxiosResponse } from 'axios';
|
||||
import type { TContextProjectionRequest, TContextUsageEvent } from './types/runs';
|
||||
import type { TFileConfig } from './file-config';
|
||||
import type * as t from './types';
|
||||
import * as permissions from './accessPermissions';
|
||||
|
|
@ -259,12 +258,6 @@ export const getTokenConfig = (): Promise<t.TTokenConfigMap> => {
|
|||
return request.get(endpoints.tokenConfig());
|
||||
};
|
||||
|
||||
export const getContextProjection = (
|
||||
payload: TContextProjectionRequest,
|
||||
): Promise<TContextUsageEvent | null> => {
|
||||
return request.post(endpoints.contextProjection(), payload);
|
||||
};
|
||||
|
||||
export const getModels = async (): Promise<t.TModelsConfig> => {
|
||||
return request.get(endpoints.models());
|
||||
};
|
||||
|
|
|
|||
|
|
@ -14,7 +14,6 @@ export enum QueryKeys {
|
|||
balance = 'balance',
|
||||
endpoints = 'endpoints',
|
||||
tokenConfig = 'tokenConfig',
|
||||
contextProjection = 'contextProjection',
|
||||
presets = 'presets',
|
||||
searchResults = 'searchResults',
|
||||
tokenCount = 'tokenCount',
|
||||
|
|
|
|||
|
|
@ -88,29 +88,6 @@ export type TContextUsageEvent = {
|
|||
completedOutputTokens?: number;
|
||||
};
|
||||
|
||||
/**
|
||||
* Request payload for a server-side context-usage projection: "what context
|
||||
* would the next call send for this branch under this config", computed by the
|
||||
* agents SDK without invoking the model. Powers the gauge in states the live
|
||||
* snapshot can't cover (page load of a snapshot-less branch, window/model
|
||||
* switch). `messageId` is the viewed branch's tail; the server walks its parent
|
||||
* chain.
|
||||
*/
|
||||
export type TContextProjectionRequest = {
|
||||
conversationId: string;
|
||||
messageId: string;
|
||||
endpoint: string;
|
||||
model?: string;
|
||||
agentId?: string;
|
||||
spec?: string;
|
||||
maxContextTokens?: number;
|
||||
/** Provider-calibrated ratio from a prior snapshot, applied as a static seed. */
|
||||
calibrationRatio?: number;
|
||||
/** Client-only cache-bust: a branch content revision so a message edit
|
||||
* (which keeps the same tail id) refetches. The server ignores it. */
|
||||
revision?: number;
|
||||
};
|
||||
|
||||
/**
|
||||
* Per-response usage rollup persisted on `responseMessage.metadata.usage`, in
|
||||
* display units (input excludes cache; output includes repaired completion).
|
||||
|
|
|
|||
|
|
@ -21,7 +21,6 @@ let mongoServer: InstanceType<typeof MongoMemoryServer>;
|
|||
let Message: mongoose.Model<IMessage>;
|
||||
let saveMessage: ReturnType<typeof createMessageMethods>['saveMessage'];
|
||||
let getMessages: ReturnType<typeof createMessageMethods>['getMessages'];
|
||||
let getMessageTextStats: ReturnType<typeof createMessageMethods>['getMessageTextStats'];
|
||||
let updateMessage: ReturnType<typeof createMessageMethods>['updateMessage'];
|
||||
let deleteMessages: ReturnType<typeof createMessageMethods>['deleteMessages'];
|
||||
let bulkSaveMessages: ReturnType<typeof createMessageMethods>['bulkSaveMessages'];
|
||||
|
|
@ -40,7 +39,6 @@ beforeAll(async () => {
|
|||
const methods = createMessageMethods(mongoose);
|
||||
saveMessage = methods.saveMessage;
|
||||
getMessages = methods.getMessages;
|
||||
getMessageTextStats = methods.getMessageTextStats;
|
||||
updateMessage = methods.updateMessage;
|
||||
deleteMessages = methods.deleteMessages;
|
||||
bulkSaveMessages = methods.bulkSaveMessages;
|
||||
|
|
@ -273,31 +271,6 @@ describe('Message Operations', () => {
|
|||
expect(messages[0].text).toBe('First message');
|
||||
expect(messages[1].text).toBe('Second message');
|
||||
});
|
||||
|
||||
it('should retrieve message text stats without returning message bodies', async () => {
|
||||
const conversationId = uuidv4();
|
||||
|
||||
await saveMessage(mockCtx, {
|
||||
messageId: 'msg1',
|
||||
conversationId,
|
||||
text: 'hello',
|
||||
quotes: ['a\nb', ''],
|
||||
user: 'user123',
|
||||
});
|
||||
|
||||
const stats = await getMessageTextStats({ conversationId, user: 'user123' }, { limit: 1 });
|
||||
|
||||
expect(stats).toEqual([
|
||||
{
|
||||
messageId: 'msg1',
|
||||
textBytes: 5,
|
||||
quoteCount: 2,
|
||||
quoteBytes: 3,
|
||||
quoteLineCount: 3,
|
||||
nonStringQuoteCount: 0,
|
||||
},
|
||||
]);
|
||||
});
|
||||
});
|
||||
|
||||
describe('deleteMessages', () => {
|
||||
|
|
|
|||
|
|
@ -1,5 +1,5 @@
|
|||
import { RetentionMode } from 'librechat-data-provider';
|
||||
import type { DeleteResult, FilterQuery, Model, PipelineStage } from 'mongoose';
|
||||
import type { DeleteResult, FilterQuery, Model } from 'mongoose';
|
||||
import type { AppConfig, IMessage } from '~/types';
|
||||
import { createTempChatExpirationDate } from '~/utils/tempChatRetention';
|
||||
import { createFallbackRetentionDate } from '~/utils/retention';
|
||||
|
|
@ -14,19 +14,6 @@ interface MessageQueryOptions {
|
|||
sort?: Record<string, 1 | -1> | false;
|
||||
}
|
||||
|
||||
interface MessageTextStatsOptions {
|
||||
limit?: number;
|
||||
}
|
||||
|
||||
export interface MessageTextStats {
|
||||
messageId: string;
|
||||
textBytes: number;
|
||||
quoteCount: number;
|
||||
quoteBytes: number;
|
||||
quoteLineCount: number;
|
||||
nonStringQuoteCount: number;
|
||||
}
|
||||
|
||||
export interface MessageMethods {
|
||||
saveMessage(
|
||||
ctx: { userId: string; isTemporary?: boolean; interfaceConfig?: AppConfig['interfaceConfig'] },
|
||||
|
|
@ -60,10 +47,6 @@ export interface MessageMethods {
|
|||
select?: string,
|
||||
options?: MessageQueryOptions,
|
||||
): Promise<IMessage[]>;
|
||||
getMessageTextStats(
|
||||
filter: FilterQuery<IMessage>,
|
||||
options?: MessageTextStatsOptions,
|
||||
): Promise<MessageTextStats[]>;
|
||||
getMessage(params: { user: string; messageId: string }): Promise<IMessage | null>;
|
||||
getMessagesByCursor(
|
||||
filter: FilterQuery<IMessage>,
|
||||
|
|
@ -374,93 +357,6 @@ export function createMessageMethods(mongoose: typeof import('mongoose')): Messa
|
|||
}
|
||||
}
|
||||
|
||||
async function getMessageTextStats(
|
||||
filter: FilterQuery<IMessage>,
|
||||
options: MessageTextStatsOptions = {},
|
||||
) {
|
||||
try {
|
||||
const Message = mongoose.models.Message as Model<IMessage>;
|
||||
const pipeline: PipelineStage[] = [{ $match: filter }];
|
||||
if (options.limit != null && options.limit > 0) {
|
||||
pipeline.push({ $limit: options.limit });
|
||||
}
|
||||
pipeline.push({
|
||||
$project: {
|
||||
_id: 0,
|
||||
messageId: 1,
|
||||
textBytes: {
|
||||
$cond: [{ $eq: [{ $type: '$text' }, 'string'] }, { $strLenBytes: '$text' }, 0],
|
||||
},
|
||||
quoteCount: {
|
||||
$cond: [{ $isArray: '$quotes' }, { $size: '$quotes' }, 0],
|
||||
},
|
||||
quoteBytes: {
|
||||
$cond: [
|
||||
{ $isArray: '$quotes' },
|
||||
{
|
||||
$sum: {
|
||||
$map: {
|
||||
input: '$quotes',
|
||||
as: 'quote',
|
||||
in: {
|
||||
$cond: [
|
||||
{ $eq: [{ $type: '$$quote' }, 'string'] },
|
||||
{ $strLenBytes: '$$quote' },
|
||||
0,
|
||||
],
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
0,
|
||||
],
|
||||
},
|
||||
quoteLineCount: {
|
||||
$cond: [
|
||||
{ $isArray: '$quotes' },
|
||||
{
|
||||
$sum: {
|
||||
$map: {
|
||||
input: '$quotes',
|
||||
as: 'quote',
|
||||
in: {
|
||||
$cond: [
|
||||
{ $eq: [{ $type: '$$quote' }, 'string'] },
|
||||
{ $size: { $split: ['$$quote', '\n'] } },
|
||||
0,
|
||||
],
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
0,
|
||||
],
|
||||
},
|
||||
nonStringQuoteCount: {
|
||||
$cond: [
|
||||
{ $isArray: '$quotes' },
|
||||
{
|
||||
$size: {
|
||||
$filter: {
|
||||
input: '$quotes',
|
||||
as: 'quote',
|
||||
cond: { $ne: [{ $type: '$$quote' }, 'string'] },
|
||||
},
|
||||
},
|
||||
},
|
||||
0,
|
||||
],
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
return await Message.aggregate<MessageTextStats>(pipeline);
|
||||
} catch (err) {
|
||||
logger.error('Error getting message text stats:', err);
|
||||
throw err;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Retrieves a single message from the database.
|
||||
*/
|
||||
|
|
@ -547,7 +443,6 @@ export function createMessageMethods(mongoose: typeof import('mongoose')): Messa
|
|||
updateMessage,
|
||||
deleteMessagesSince,
|
||||
getMessages,
|
||||
getMessageTextStats,
|
||||
getMessage,
|
||||
getMessagesByCursor,
|
||||
searchMessages,
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue