mirror of
https://github.com/danny-avila/LibreChat.git
synced 2026-08-28 04:37:37 +00:00
💰 fix: Multi-Agent Token Spending & Prevent Double-Spend (#11433)
* fix: Token Spending Logic for Multi-Agents on Abort Scenarios * Implemented logic to skip token spending if a conversation is aborted, preventing double-spending. * Introduced `spendCollectedUsage` function to handle token spending for multiple models during aborts, ensuring accurate accounting for parallel agents. * Updated `GenerationJobManager` to store and retrieve collected usage data for improved abort handling. * Added comprehensive tests for the new functionality, covering various scenarios including cache token handling and parallel agent usage. * fix: Memory Context Handling for Multi-Agents * Refactored `buildMessages` method to pass memory context to parallel agents, ensuring they share the same user context. * Improved handling of memory context when no existing instructions are present for parallel agents. * Added comprehensive tests to verify memory context propagation and behavior under various scenarios, including cases with no memory available and empty agent configurations. * Enhanced logging for better traceability of memory context additions to agents. * chore: Memory Context Documentation for Parallel Agents * Updated documentation in the `AgentClient` class to clarify the in-place mutation of agentConfig objects when passing memory context to parallel agents. * Added notes on the implications of mutating objects directly to ensure all parallel agents receive the correct memory context before execution. * chore: UsageMetadata Interface docs for Token Spending * Expanded the UsageMetadata interface to support both OpenAI and Anthropic cache token formats. * Added detailed documentation for cache token properties, including mutually exclusive fields for different model types. * Improved clarity on how to access cache token details for accurate token spending tracking. * fix: Enhance Token Spending Logic in Abort Middleware * Refactored `spendCollectedUsage` function to utilize Promise.all for concurrent token spending, improving performance and ensuring all operations complete before clearing the collectedUsage array. * Added documentation to clarify the importance of clearing the collectedUsage array to prevent double-spending in abort scenarios. * Updated tests to verify the correct behavior of the spending logic and the clearing of the array after spending operations.
This commit is contained in:
parent
32e6f3b8e5
commit
36c5a88c4e
11 changed files with 1440 additions and 28 deletions
|
|
@ -1,7 +1,12 @@
|
|||
import { logger } from '@librechat/data-schemas';
|
||||
import type { StandardGraph } from '@librechat/agents';
|
||||
import type { Agents } from 'librechat-data-provider';
|
||||
import type { IJobStore, SerializableJobData, JobStatus } from '~/stream/interfaces/IJobStore';
|
||||
import type {
|
||||
SerializableJobData,
|
||||
UsageMetadata,
|
||||
IJobStore,
|
||||
JobStatus,
|
||||
} from '~/stream/interfaces/IJobStore';
|
||||
|
||||
/**
|
||||
* Content state for a job - volatile, in-memory only.
|
||||
|
|
@ -10,6 +15,7 @@ import type { IJobStore, SerializableJobData, JobStatus } from '~/stream/interfa
|
|||
interface ContentState {
|
||||
contentParts: Agents.MessageContentComplex[];
|
||||
graphRef: WeakRef<StandardGraph> | null;
|
||||
collectedUsage: UsageMetadata[];
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -240,6 +246,7 @@ export class InMemoryJobStore implements IJobStore {
|
|||
this.contentState.set(streamId, {
|
||||
contentParts: [],
|
||||
graphRef: new WeakRef(graph),
|
||||
collectedUsage: [],
|
||||
});
|
||||
}
|
||||
}
|
||||
|
|
@ -252,10 +259,30 @@ export class InMemoryJobStore implements IJobStore {
|
|||
if (existing) {
|
||||
existing.contentParts = contentParts;
|
||||
} else {
|
||||
this.contentState.set(streamId, { contentParts, graphRef: null });
|
||||
this.contentState.set(streamId, { contentParts, graphRef: null, collectedUsage: [] });
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Set collected usage reference for a job.
|
||||
*/
|
||||
setCollectedUsage(streamId: string, collectedUsage: UsageMetadata[]): void {
|
||||
const existing = this.contentState.get(streamId);
|
||||
if (existing) {
|
||||
existing.collectedUsage = collectedUsage;
|
||||
} else {
|
||||
this.contentState.set(streamId, { contentParts: [], graphRef: null, collectedUsage });
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Get collected usage for a job.
|
||||
*/
|
||||
getCollectedUsage(streamId: string): UsageMetadata[] {
|
||||
const state = this.contentState.get(streamId);
|
||||
return state?.collectedUsage ?? [];
|
||||
}
|
||||
|
||||
/**
|
||||
* Get content parts for a job.
|
||||
* Returns live content from stored reference.
|
||||
|
|
|
|||
|
|
@ -1,9 +1,14 @@
|
|||
import { logger } from '@librechat/data-schemas';
|
||||
import { createContentAggregator } from '@librechat/agents';
|
||||
import type { IJobStore, SerializableJobData, JobStatus } from '~/stream/interfaces/IJobStore';
|
||||
import type { StandardGraph } from '@librechat/agents';
|
||||
import type { Agents } from 'librechat-data-provider';
|
||||
import type { Redis, Cluster } from 'ioredis';
|
||||
import type {
|
||||
SerializableJobData,
|
||||
UsageMetadata,
|
||||
IJobStore,
|
||||
JobStatus,
|
||||
} from '~/stream/interfaces/IJobStore';
|
||||
|
||||
/**
|
||||
* Key prefixes for Redis storage.
|
||||
|
|
@ -90,6 +95,13 @@ export class RedisJobStore implements IJobStore {
|
|||
*/
|
||||
private localGraphCache = new Map<string, WeakRef<StandardGraph>>();
|
||||
|
||||
/**
|
||||
* Local cache for collectedUsage arrays.
|
||||
* Generation happens on a single instance, so collectedUsage is only available locally.
|
||||
* For cross-replica abort, the abort handler falls back to text-based token counting.
|
||||
*/
|
||||
private localCollectedUsageCache = new Map<string, UsageMetadata[]>();
|
||||
|
||||
/** Cleanup interval in ms (1 minute) */
|
||||
private cleanupIntervalMs = 60000;
|
||||
|
||||
|
|
@ -227,6 +239,7 @@ export class RedisJobStore implements IJobStore {
|
|||
async deleteJob(streamId: string): Promise<void> {
|
||||
// Clear local caches
|
||||
this.localGraphCache.delete(streamId);
|
||||
this.localCollectedUsageCache.delete(streamId);
|
||||
|
||||
// Note: userJobs cleanup is handled lazily via self-healing in getActiveJobIdsByUser
|
||||
// In cluster mode, separate runningJobs (global) from stream-specific keys (same slot)
|
||||
|
|
@ -290,6 +303,7 @@ export class RedisJobStore implements IJobStore {
|
|||
if (!job) {
|
||||
await this.redis.srem(KEYS.runningJobs, streamId);
|
||||
this.localGraphCache.delete(streamId);
|
||||
this.localCollectedUsageCache.delete(streamId);
|
||||
cleaned++;
|
||||
continue;
|
||||
}
|
||||
|
|
@ -298,6 +312,7 @@ export class RedisJobStore implements IJobStore {
|
|||
if (job.status !== 'running') {
|
||||
await this.redis.srem(KEYS.runningJobs, streamId);
|
||||
this.localGraphCache.delete(streamId);
|
||||
this.localCollectedUsageCache.delete(streamId);
|
||||
cleaned++;
|
||||
continue;
|
||||
}
|
||||
|
|
@ -382,6 +397,7 @@ export class RedisJobStore implements IJobStore {
|
|||
}
|
||||
// Clear local caches
|
||||
this.localGraphCache.clear();
|
||||
this.localCollectedUsageCache.clear();
|
||||
// Don't close the Redis connection - it's shared
|
||||
logger.info('[RedisJobStore] Destroyed');
|
||||
}
|
||||
|
|
@ -406,11 +422,28 @@ export class RedisJobStore implements IJobStore {
|
|||
* No-op for Redis - content parts are reconstructed from chunks.
|
||||
* Metadata (agentId, groupId) is embedded directly on content parts by the agent runtime.
|
||||
*/
|
||||
setContentParts(_streamId: string, _contentParts: Agents.MessageContentComplex[]): void {
|
||||
setContentParts(): void {
|
||||
// Content parts are reconstructed from chunks during getContentParts
|
||||
// No separate storage needed
|
||||
}
|
||||
|
||||
/**
|
||||
* Store collectedUsage reference in local cache.
|
||||
* This is used for abort handling to spend tokens for all models.
|
||||
* Note: Only available on the generating instance; cross-replica abort uses fallback.
|
||||
*/
|
||||
setCollectedUsage(streamId: string, collectedUsage: UsageMetadata[]): void {
|
||||
this.localCollectedUsageCache.set(streamId, collectedUsage);
|
||||
}
|
||||
|
||||
/**
|
||||
* Get collected usage for a job.
|
||||
* Only available if this is the generating instance.
|
||||
*/
|
||||
getCollectedUsage(streamId: string): UsageMetadata[] {
|
||||
return this.localCollectedUsageCache.get(streamId) ?? [];
|
||||
}
|
||||
|
||||
/**
|
||||
* Get aggregated content - tries local cache first, falls back to Redis reconstruction.
|
||||
*
|
||||
|
|
@ -528,6 +561,7 @@ export class RedisJobStore implements IJobStore {
|
|||
clearContentState(streamId: string): void {
|
||||
// Clear local caches immediately
|
||||
this.localGraphCache.delete(streamId);
|
||||
this.localCollectedUsageCache.delete(streamId);
|
||||
|
||||
// Fire and forget - async cleanup for Redis
|
||||
this.clearContentStateAsync(streamId).catch((err) => {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue