diff --git a/.env.example b/.env.example index b310ffe2cce..fa196f5e3cb 100644 --- a/.env.example +++ b/.env.example @@ -1232,11 +1232,13 @@ HELP_AND_FAQ_URL=https://librechat.ai # SCHEDULES_DISABLED=true # Coalesce streamed model/tool-argument deltas into windowed Redis publications (ms). -# Unset or 0 (default) publishes per delta. 25 is recommended: it batches the publish -# EVAL and the durable append across the window (fewer Redis round trips and lower -# Redis CPU at high token rates) at the cost of up to one window of added delivery -# latency. Enable only after EVERY replica runs a build with batch-frame support; -# older subscribers drop coalesced frames. Values are capped at 1000. +# Defaults to 25ms when unset; explicit 0 publishes per delta. Batches both publish +# and durable append operations (fewer Redis round trips and repeated guard/TTL work at high +# token rates), adding up to one window of buffering latency and crash-loss exposure +# for unflushed deltas. In-memory streams are unchanged. Values are capped at 1000. +# Batched publications retain individual chunk frames for older subscribers; Pub/Sub +# message count is unchanged. Incoming chunk_batch frames remain supported for existing +# opt-in producers. See UPGRADING.md for compatibility guidance. # Keep the window <= the stream-smoothing cadence (`streamRate`, default 25ms): # each smoothing tick emits its pieces in one burst, so a tick-sized window # captures exactly one batch per tick; a larger window re-batches the paced diff --git a/.github/pull_request_template.md b/.github/pull_request_template.md index 088c0ac33c1..82997ea8bf5 100644 --- a/.github/pull_request_template.md +++ b/.github/pull_request_template.md @@ -1,39 +1,40 @@ -# Pull Request Template +# Pull Request -⚠️ Before Submitting a PR, Please Review: -- Please ensure that you have thoroughly read and understood the [Contributing Docs](https://github.com/danny-avila/LibreChat/blob/main/.github/CONTRIBUTING.md) before submitting your Pull Request. - -⚠️ Documentation Updates Notice: -- Kindly note that documentation updates are managed in this repository: [librechat.ai](https://github.com/LibreChat-AI/librechat.ai) +> Before submitting, please review the [Contributing Guide](https://github.com/danny-avila/LibreChat/blob/main/.github/CONTRIBUTING.md). +> +> Documentation changes belong in the [LibreChat documentation repository](https://github.com/LibreChat-AI/librechat.ai). ## Summary ## How it works -## Change Type +## Type of change -Please delete any irrelevant options. + -- [ ] Bug fix (non-breaking change which fixes an issue) -- [ ] New feature (non-breaking change which adds functionality) -- [ ] Breaking change (fix or feature that would cause existing functionality to not work as expected) -- [ ] This change requires a documentation update -- [ ] Translation update +* [ ] Bug fix +* [ ] Feature +* [ ] Refactor +* [ ] Performance improvement +* [ ] Breaking change +* [ ] Documentation +* [ ] Translation +* [ ] Tests / tooling / CI ## Testing -Please describe your test process and include instructions so that we can reproduce your test. If there are any important variables for your testing configuration, list them here. + -### **Test Configuration**: +**Tested environments/configuration:** + + + +**Automated tests:** + + + +## Screenshots / recordings + + + +## Risk / compatibility + + ## Checklist -Please delete any irrelevant options. - -- [ ] My code adheres to this project's style guidelines -- [ ] I have performed a self-review of my own code -- [ ] I have commented in any complex areas of my code -- [ ] I have made pertinent documentation changes -- [ ] My changes do not introduce new warnings -- [ ] I have written tests demonstrating that my changes are effective or that my feature works -- [ ] Local unit tests pass with my changes -- [ ] Any changes dependent on mine have been merged and published in downstream modules. -- [ ] A pull request for updating the documentation has been submitted. +* [ ] I reviewed my own changes +* [ ] Relevant tests have been added or updated +* [ ] Existing relevant tests pass +* [ ] The change does not introduce new warnings or errors +* [ ] User-facing or complex behavior is documented where necessary +* [ ] Required dependency changes have been merged/published +* [ ] Required documentation PR: diff --git a/.github/workflows/backend-review.yml b/.github/workflows/backend-review.yml index 6271abc91f2..82d0088cdc9 100644 --- a/.github/workflows/backend-review.yml +++ b/.github/workflows/backend-review.yml @@ -433,7 +433,9 @@ jobs: - name: Run unit tests (shard ${{ matrix.shard }}/3) env: SELECTED: ${{ needs.codegraph_select.outputs.api_files }} + JEST_JSON: --json --outputFile=${{ github.workspace }}/jest-results/jest-results-api-${{ matrix.shard }}.json run: | + mkdir -p "$GITHUB_WORKSPACE/jest-results" cd api # A selected path can be stale in exactly two ways at this checkout (Codex P2, #15145 r6): # deleted on the branch — dropped, which matches full CI (the file runs nowhere) — or @@ -448,14 +450,26 @@ jobs: KEEP="${KEEP# }" if [ -z "$KEEP" ]; then echo "no selected test file exists at HEAD (stale selection); running FULL" - npm run test:ci -- --shard=${{ matrix.shard }}/3 + npm run test:ci -- --shard=${{ matrix.shard }}/3 $JEST_JSON else echo "codegraph: $(echo $KEEP | wc -w) selected test files (safe mode)" - npm run test:ci -- --shard=${{ matrix.shard }}/3 --passWithNoTests --runTestsByPath $KEEP + npm run test:ci -- --shard=${{ matrix.shard }}/3 --passWithNoTests --runTestsByPath $KEEP $JEST_JSON fi else - npm run test:ci -- --shard=${{ matrix.shard }}/3 + npm run test:ci -- --shard=${{ matrix.shard }}/3 $JEST_JSON fi + # Per-test results keyed by head SHA (run.head_sha), for the codegraph test-evidence feed. + # Never part of the gate: it cannot fail the job, and a re-run overwrites its own artifact. + - name: Upload Jest results + if: ${{ !cancelled() }} + continue-on-error: true + uses: actions/upload-artifact@v6 + with: + name: jest-results-api-${{ matrix.shard }} + path: jest-results/ + retention-days: 7 + if-no-files-found: ignore + overwrite: true test-data-provider: name: 'Tests: data-provider' @@ -509,7 +523,9 @@ jobs: - name: Run unit tests env: SELECTED: ${{ needs.codegraph_select.outputs.dataprovider_files }} + JEST_JSON: --json --outputFile=${{ github.workspace }}/jest-results/jest-results-data-provider.json run: | + mkdir -p "$GITHUB_WORKSPACE/jest-results" cd packages/data-provider if [ -n "$SELECTED" ]; then KEEP="" @@ -519,14 +535,26 @@ jobs: KEEP="${KEEP# }" if [ -z "$KEEP" ]; then echo "no selected test file exists at HEAD (stale selection); running FULL" - npm run test:ci + npm run test:ci -- $JEST_JSON else echo "codegraph: $(echo $KEEP | wc -w) selected test files (safe mode)" - npm run test:ci -- --passWithNoTests --runTestsByPath $KEEP + npm run test:ci -- --passWithNoTests --runTestsByPath $KEEP $JEST_JSON fi else - npm run test:ci + npm run test:ci -- $JEST_JSON fi + # Per-test results keyed by head SHA (run.head_sha), for the codegraph test-evidence feed. + # Never part of the gate: it cannot fail the job, and a re-run overwrites its own artifact. + - name: Upload Jest results + if: ${{ !cancelled() }} + continue-on-error: true + uses: actions/upload-artifact@v6 + with: + name: jest-results-data-provider + path: jest-results/ + retention-days: 7 + if-no-files-found: ignore + overwrite: true test-data-schemas: name: 'Tests: data-schemas' @@ -586,7 +614,9 @@ jobs: - name: Run unit tests env: SELECTED: ${{ needs.codegraph_select.outputs.dataschemas_files }} + JEST_JSON: --json --outputFile=${{ github.workspace }}/jest-results/jest-results-data-schemas.json run: | + mkdir -p "$GITHUB_WORKSPACE/jest-results" cd packages/data-schemas if [ -n "$SELECTED" ]; then KEEP="" @@ -596,14 +626,26 @@ jobs: KEEP="${KEEP# }" if [ -z "$KEEP" ]; then echo "no selected test file exists at HEAD (stale selection); running FULL" - npm run test:ci + npm run test:ci -- $JEST_JSON else echo "codegraph: $(echo $KEEP | wc -w) selected test files (safe mode)" - npm run test:ci -- --passWithNoTests --runTestsByPath $KEEP + npm run test:ci -- --passWithNoTests --runTestsByPath $KEEP $JEST_JSON fi else - npm run test:ci + npm run test:ci -- $JEST_JSON fi + # Per-test results keyed by head SHA (run.head_sha), for the codegraph test-evidence feed. + # Never part of the gate: it cannot fail the job, and a re-run overwrites its own artifact. + - name: Upload Jest results + if: ${{ !cancelled() }} + continue-on-error: true + uses: actions/upload-artifact@v6 + with: + name: jest-results-data-schemas + path: jest-results/ + retention-days: 7 + if-no-files-found: ignore + overwrite: true test-packages-api: name: 'Tests: @librechat/api (shard ${{ matrix.shard }}/4)' @@ -677,7 +719,9 @@ jobs: - name: Run unit tests (shard ${{ matrix.shard }}/4) env: SELECTED: ${{ needs.codegraph_select.outputs.pkgapi_files }} + JEST_JSON: --json --outputFile=${{ github.workspace }}/jest-results/jest-results-packages-api-${{ matrix.shard }}.json run: | + mkdir -p "$GITHUB_WORKSPACE/jest-results" cd packages/api if [ -n "$SELECTED" ]; then KEEP="" @@ -687,11 +731,23 @@ jobs: KEEP="${KEEP# }" if [ -z "$KEEP" ]; then echo "no selected test file exists at HEAD (stale selection); running FULL" - npm run test:ci -- --shard=${{ matrix.shard }}/4 + npm run test:ci -- --shard=${{ matrix.shard }}/4 $JEST_JSON else echo "codegraph: $(echo $KEEP | wc -w) selected test files (safe mode)" - npm run test:ci -- --shard=${{ matrix.shard }}/4 --passWithNoTests --runTestsByPath $KEEP + npm run test:ci -- --shard=${{ matrix.shard }}/4 --passWithNoTests --runTestsByPath $KEEP $JEST_JSON fi else - npm run test:ci -- --shard=${{ matrix.shard }}/4 + npm run test:ci -- --shard=${{ matrix.shard }}/4 $JEST_JSON fi + # Per-test results keyed by head SHA (run.head_sha), for the codegraph test-evidence feed. + # Never part of the gate: it cannot fail the job, and a re-run overwrites its own artifact. + - name: Upload Jest results + if: ${{ !cancelled() }} + continue-on-error: true + uses: actions/upload-artifact@v6 + with: + name: jest-results-packages-api-${{ matrix.shard }} + path: jest-results/ + retention-days: 7 + if-no-files-found: ignore + overwrite: true diff --git a/.github/workflows/frontend-review.yml b/.github/workflows/frontend-review.yml index 8fec3bac4c1..40ebe5b7221 100644 --- a/.github/workflows/frontend-review.yml +++ b/.github/workflows/frontend-review.yml @@ -304,7 +304,9 @@ jobs: - name: Run unit tests env: SELECTED: ${{ needs.codegraph_select.outputs.clientpkg_files }} + JEST_JSON: --json --outputFile=${{ github.workspace }}/jest-results/jest-results-packages-client.json run: | + mkdir -p "$GITHUB_WORKSPACE/jest-results" # A selected path can be stale in exactly two ways at this checkout (Codex P2, #15145 r6): # deleted on the branch — dropped, which matches full CI (the file runs nowhere) — or # renamed, where the NEW path is a changed test file and is selected independently. If @@ -318,15 +320,27 @@ jobs: KEEP="${KEEP# }" if [ -z "$KEEP" ]; then echo "no selected test file exists at HEAD (stale selection); running FULL" - npm run test:ci + npm run test:ci -- $JEST_JSON else echo "codegraph: $(echo $KEEP | wc -w) selected test files (safe mode)" - npm run test:ci -- --passWithNoTests --runTestsByPath $KEEP + npm run test:ci -- --passWithNoTests --runTestsByPath $KEEP $JEST_JSON fi else - npm run test:ci + npm run test:ci -- $JEST_JSON fi working-directory: packages/client + # Per-test results keyed by head SHA (run.head_sha), for the codegraph test-evidence feed. + # Never part of the gate: it cannot fail the job, and a re-run overwrites its own artifact. + - name: Upload Jest results + if: ${{ !cancelled() }} + continue-on-error: true + uses: actions/upload-artifact@v6 + with: + name: jest-results-packages-client + path: jest-results/ + retention-days: 7 + if-no-files-found: ignore + overwrite: true test-ubuntu: name: 'Tests: Ubuntu (shard ${{ matrix.shard }}/2)' @@ -378,7 +392,9 @@ jobs: - name: Run unit tests (shard ${{ matrix.shard }}/2) env: SELECTED: ${{ needs.codegraph_select.outputs.client_files }} + JEST_JSON: --json --outputFile=${{ github.workspace }}/jest-results/jest-results-client-${{ matrix.shard }}.json run: | + mkdir -p "$GITHUB_WORKSPACE/jest-results" if [ -n "$SELECTED" ]; then KEEP="" for f in $SELECTED; do @@ -387,15 +403,27 @@ jobs: KEEP="${KEEP# }" if [ -z "$KEEP" ]; then echo "no selected test file exists at HEAD (stale selection); running FULL" - npm run test:ci -- --shard=${{ matrix.shard }}/2 + npm run test:ci -- --shard=${{ matrix.shard }}/2 $JEST_JSON else echo "codegraph: $(echo $KEEP | wc -w) selected test files (safe mode)" - npm run test:ci -- --shard=${{ matrix.shard }}/2 --passWithNoTests --runTestsByPath $KEEP + npm run test:ci -- --shard=${{ matrix.shard }}/2 --passWithNoTests --runTestsByPath $KEEP $JEST_JSON fi else - npm run test:ci -- --shard=${{ matrix.shard }}/2 + npm run test:ci -- --shard=${{ matrix.shard }}/2 $JEST_JSON fi working-directory: client + # Per-test results keyed by head SHA (run.head_sha), for the codegraph test-evidence feed. + # Never part of the gate: it cannot fail the job, and a re-run overwrites its own artifact. + - name: Upload Jest results + if: ${{ !cancelled() }} + continue-on-error: true + uses: actions/upload-artifact@v6 + with: + name: jest-results-client-${{ matrix.shard }} + path: jest-results/ + retention-days: 7 + if-no-files-found: ignore + overwrite: true build-verify: name: Vite build verification diff --git a/UPGRADING.md b/UPGRADING.md index 51eab7cc678..aecd00630c6 100644 --- a/UPGRADING.md +++ b/UPGRADING.md @@ -1,5 +1,29 @@ # Upgrading LibreChat +## Redis streaming now coalesces deltas by default + +Redis-backed streams now batch model and tool-argument deltas in a **25 ms** +window when `STREAM_DELTA_COALESCE_MS` is unset. Explicit `0` keeps per-delta +publication; existing explicit values retain their behavior. In-memory streams +are unchanged, including deployments using `USE_REDIS_STREAMS=false`. + +The same window batches durable appends and publications. It reduces Redis +round trips and repeated guard/TTL work, but not Pub/Sub message count at the cost of up to one window of buffering latency and possible loss +of unflushed deltas on process crash. Terminal and non-coalescable barriers still +flush pending batches before proceeding. + +Coalescing batches Redis requests, not the subscriber wire format: each event +still publishes as an individually sequenced `chunk` frame. Subscribers from +before batch-frame support can read these publications without a preparatory +configuration change. New subscribers also retain `chunk_batch` decoding for +interoperation with existing opt-in batching producers. + +This removes the new default's batch-frame compatibility hazard; it does not +promise compatibility across unrelated generation-protocol changes. If an +existing producer already emits `chunk_batch` frames through explicit opt-in, +keep those producers away from subscribers that predate batch-frame support, +or disable coalescing on those producers before introducing older subscribers. + ## Tenant index migration (v0.8.7 and earlier databases) Upgrading an existing database can log `Index build failed` for User, Role, diff --git a/api/app/clients/BaseClient.js b/api/app/clients/BaseClient.js index a685ee3091d..4a4f6975bda 100644 --- a/api/app/clients/BaseClient.js +++ b/api/app/clients/BaseClient.js @@ -21,9 +21,15 @@ const { collectModelBoundHistoricalFileIdState, projectModelBoundSourceFiles, isModelBoundAttachmentFile, + isToolOwnedAttachment, withBalanceReservations, findCheckpointSummaryPart, getSummaryPartText, + runAfterSeed, + saveTurnConversation, + seedTurnConversation, + needsRetentionConversation, + getConversationWriteContext, } = require('@librechat/api'); const { Constants, @@ -32,11 +38,9 @@ const { ErrorTypes, ContentTypes, isCompactedLeaf, - excludedKeys, EModelEndpoint, isParamEndpoint, isAgentsEndpoint, - isEphemeralAgentId, supportsBalanceCheck, isBedrockDocumentType, HITL_MESSAGE_FILTER_FIELDS, @@ -282,6 +286,12 @@ class BaseClient { return false; } + /** Whether a deferred parent write may still create a new conversation's row up front, so + * the conversation lists can return it while the run is in flight. */ + shouldSeedDeferredConversation() { + return false; + } + /** Returns the request-scoped deferred parent-write controller, when any. */ getModelBoundUserMessagePersistence() { return this.modelBoundUserMessagePersistence; @@ -909,6 +919,18 @@ class BaseClient { if (this.shouldDeferUserMessagePersistence()) { let state = 'pending'; let startPersistence = startUserMessagePersistence; + if (!this.skipSaveConvo && this.shouldSeedDeferredConversation()) { + const seed = seedTurnConversation( + db, + this.getTurnConversationFields( + this.options, + userMessage.conversationId, + saveOptions, + 'api/app/clients/BaseClient.js - sendMessage #seedConversation', + ), + ); + startPersistence = runAfterSeed(seed, startUserMessagePersistence); + } let resolvePersistence; let removeAbortListener = () => {}; const persistencePromise = new Promise((resolve) => { @@ -1323,25 +1345,10 @@ class BaseClient { const hasAddedConvo = options?.req?.body?.addedConvo != null; const req = options?.req; - if ( - req?.config?.interfaceConfig?.retentionMode === 'all' && - req?.config?.interfaceConfig?.generalChatRetention !== undefined && - !Object.prototype.hasOwnProperty.call(req, 'resolvedConversation') - ) { + if (needsRetentionConversation(req)) { req.resolvedConversation = await db.getConvo(req.user.id, message.conversationId); } - const hasResolvedConversation = - req != null && Object.prototype.hasOwnProperty.call(req, 'resolvedConversation'); - const resolvedRetention = hasResolvedConversation ? req.resolvedConversation : null; - const reqCtx = { - userId: req?.user?.id, - isTemporary: - req?._agentEventBindingRetention?.isTemporary ?? - resolvedRetention?.isTemporary ?? - req?.body?.isTemporary, - expiredAt: req?._agentEventBindingRetention?.expiredAt ?? resolvedRetention?.expiredAt, - interfaceConfig: req?.config?.interfaceConfig, - }; + const reqCtx = getConversationWriteContext(req); const savedMessage = await db.saveMessage( reqCtx, { @@ -1358,72 +1365,43 @@ class BaseClient { 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 skippedExistingConvoLookup = this.fetchedConvo === true; - let existingConvo = null; - if (!skippedExistingConvoLookup && hasResolvedConversation) { - existingConvo = req.resolvedConversation; - } else if (!skippedExistingConvoLookup) { - existingConvo = await db.getConvo(req?.user?.id, message.conversationId); - } - // Keep the authenticated conversation available for response, abort, and retry saves. - // fetchedConvo already prevents repeating the conversation initialization work. - 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, - noUpsert: req?._agentEventBindingParentConversationId != null, - initialAgentId: hasNonEphemeralAgent ? options.agent?.id : null, - createdAtOnInsert: shouldSetCreatedAtOnInsert ? validCreatedAtOnInsert : undefined, - ...(savedMessage?._id != null ? { appendMessageIds: [savedMessage._id] } : {}), + const { conversation, initialized } = await saveTurnConversation(db, { + ...this.getTurnConversationFields( + options, + message.conversationId, + endpointOptions, + 'api/app/clients/BaseClient.js - saveMessageToDatabase #saveConvo', + ), + ctx: reqCtx, + initialized: this.fetchedConvo === true, + savedMessageId: savedMessage?._id, }); - - if (req != null && conversation != null) { - req.resolvedConversation = conversation; + if (initialized) { + this.fetchedConvo = true; } return { message: savedMessage, conversation }; } + /** + * The conversation fields a turn's writes share. + * @param {Object} options - The client options snapshot. + * @param {string} conversationId + * @param {Partial} endpointOptions + * @param {string} context - Names the write in the save log. + */ + getTurnConversationFields(options, conversationId, endpointOptions, context) { + return { + req: options.req, + conversationId, + endpoint: options.endpoint, + endpointType: options.endpointType, + endpointOptions, + agentId: options.agent?.id, + context, + }; + } + /** * Update a message in the database. * @param {Partial} message @@ -1844,13 +1822,7 @@ class BaseClient { /* An explicit `provider` path is authoritative: lazy provisioning stamps * `embedded`/`codeEnvRef` on files that are still meant for the model, so the * legacy tool-provisioning exclusion only applies to records without one. */ - if ( - deliveryPath !== 'provider' && - (file.embedded === true || - file.metadata?.codeEnvRef != null || - file.metadata?.codeEnvRefs != null || - file.metadata?.fileIdentifier != null) - ) { + if (deliveryPath !== 'provider' && isToolOwnedAttachment(file)) { allFiles.push(file); continue; } diff --git a/api/app/clients/specs/BaseClient.test.js b/api/app/clients/specs/BaseClient.test.js index 7277e634879..2c675aa5dcd 100644 --- a/api/app/clients/specs/BaseClient.test.js +++ b/api/app/clients/specs/BaseClient.test.js @@ -839,6 +839,124 @@ describe('BaseClient', () => { ); }); + describe('seeding the conversation for a deferred user message', () => { + const seedContext = 'api/app/clients/BaseClient.js - sendMessage #seedConversation'; + const flush = () => new Promise((resolve) => setImmediate(resolve)); + const savedUserMessage = () => + saveMessage.mock.calls.some(([, message]) => message.isCreatedByUser === true); + + beforeEach(() => { + saveMessage.mockReset(); + saveConvo.mockReset().mockImplementation(async (_ctx, fields) => ({ ...fields })); + getConvo.mockReset().mockResolvedValue(null); + /** A fresh options object: the suite-level one is shared by reference across clients. */ + TestClient.options = { + ...TestClient.options, + req: { user: { id: 'seed-user' }, body: {} }, + }; + TestClient.shouldDeferUserMessagePersistence = jest.fn(() => true); + TestClient.shouldSeedDeferredConversation = jest.fn(() => true); + }); + + afterEach(() => { + saveMessage.mockReset(); + saveConvo.mockReset(); + getConvo.mockReset(); + }); + + test('creates a new conversation row while the message itself waits for admission', async () => { + TestClient.sendCompletion.mockImplementation(async () => { + await flush(); + expect(savedUserMessage()).toBe(false); + expect(saveConvo).toHaveBeenCalledTimes(1); + const [, fields, seedOptions] = saveConvo.mock.calls[0]; + expect(fields).toEqual(expect.objectContaining({ conversationId: expect.any(String) })); + expect(seedOptions).toEqual(expect.objectContaining({ context: seedContext })); + /** An empty append set spares `saveConvo` the read of a message list the seed lacks. */ + expect(seedOptions).toEqual(expect.objectContaining({ appendMessageIds: [] })); + return { completion: 'Safe response', metadata: undefined }; + }); + + await TestClient.sendMessage('Message with an attachment'); + + expect(savedUserMessage()).toBe(true); + /** The message save reuses the seeded row instead of looking the conversation up again. */ + expect(getConvo).toHaveBeenCalledTimes(1); + }); + + test('leaves an existing conversation to the deferred message write', async () => { + getConvo.mockResolvedValue({ conversationId: 'existing-convo', endpoint: 'openAI' }); + TestClient.sendCompletion.mockImplementation(async () => { + await flush(); + expect(saveConvo).not.toHaveBeenCalled(); + return { completion: 'Safe response', metadata: undefined }; + }); + + await TestClient.sendMessage('Message with an attachment'); + + expect(savedUserMessage()).toBe(true); + expect(saveConvo.mock.calls.every(([, , opts]) => opts.context !== seedContext)).toBe(true); + }); + + test('holds the deferred message write until an in-flight seed lands', async () => { + const seedWrite = deferred(); + saveConvo.mockImplementationOnce(() => seedWrite.promise); + const completionStarted = deferred(); + const completionResult = deferred(); + const abortController = new AbortController(); + TestClient.sendCompletion.mockImplementation(() => { + completionStarted.resolve(); + return completionResult.promise; + }); + + const sendPromise = TestClient.sendMessage('Message with an attachment', { + abortController, + }); + await completionStarted.promise; + await flush(); + expect(saveConvo).toHaveBeenCalledTimes(1); + + abortController.abort(); + await flush(); + expect(savedUserMessage()).toBe(false); + + seedWrite.resolve({ conversationId: 'seeded-convo' }); + await flush(); + expect(savedUserMessage()).toBe(true); + + completionResult.resolve({ completion: 'Partial response', metadata: undefined }); + await sendPromise; + }); + + test('keeps the deferral cancellable after seeding when the model boundary rejects content', async () => { + const policyError = new ContentFilterError({ source: 'message', field: 'text' }); + TestClient.sendCompletion.mockImplementation(async () => { + await flush(); + throw policyError; + }); + + await expect(TestClient.sendMessage('Message with an attachment')).rejects.toBe( + policyError, + ); + + expect(savedUserMessage()).toBe(false); + expect(saveConvo).toHaveBeenCalledTimes(1); + }); + + test('does not seed when the client holds back every write', async () => { + TestClient.shouldSeedDeferredConversation = jest.fn(() => false); + TestClient.sendCompletion.mockImplementation(async () => { + await flush(); + expect(saveConvo).not.toHaveBeenCalled(); + return { completion: 'Safe response', metadata: undefined }; + }); + + await TestClient.sendMessage('Message with an attachment'); + + expect(savedUserMessage()).toBe(true); + }); + }); + test('preserves eager user-message persistence for non-policy provider failures', async () => { saveMessage.mockClear(); saveConvo.mockClear(); @@ -3584,13 +3702,13 @@ describe('BaseClient', () => { expect(TestClient.getTextContextAttachments([file])).toEqual([]); }); - test('delivers the text when the tool that would read the file has yet to receive it', () => { - /* An upload that named no destination is filed under no tool, so the enabled tool alone - * cannot serve it and withholding the text left it readable by nothing. */ + test('delivers the text when File Search has yet to receive the file', () => { + /* An upload that named no destination is filed under no tool, so an enabled search tool + * alone cannot serve it and withholding the text left it readable by nothing. */ routeCsvToTools(); TestClient.options.agent = { provider: EModelEndpoint.openAI, - fileConsumers: { executeCode: true, fileSearch: true }, + fileConsumers: { executeCode: false, fileSearch: true }, }; TestClient.options.agent.deliveryRouting = resolveTurnDeliveryRouting({ agent: TestClient.options.agent, @@ -3614,6 +3732,32 @@ describe('BaseClient', () => { ]); }); + test('leaves a file Run Code can read with Run Code before the sandbox holds it', () => { + /* Run Code uploads the file on its first call, so the text stays off the prompt, where it + * would otherwise count toward the history limits until that call. */ + routeCsvToTools(); + TestClient.options.agent = { + provider: EModelEndpoint.openAI, + fileConsumers: { executeCode: true, fileSearch: true }, + }; + TestClient.options.agent.deliveryRouting = resolveTurnDeliveryRouting({ + agent: TestClient.options.agent, + config: TestClient.options.req?.config, + }); + const file = { + file_id: 'unprovisioned-csv', + filename: 'sales.csv', + type: 'text/csv', + text: 'region,total', + llmDeliveryPath: 'none', + metadata: { destinationChosen: false }, + }; + + expect(TestClient.getAttachmentDeliveryPath(file)).toBe('none'); + expect(TestClient.getTextContextAttachments([file])).toEqual([]); + expect(TestClient.resolveTurnAttachments([file])).toEqual([file]); + }); + test('does not fall back on an endpoint that has not enabled it', () => { routeCsvToTools({ textFallbackWithoutTools: false }); TestClient.options.agent = { @@ -3931,6 +4075,29 @@ describe('BaseClient', () => { expect(message.image_urls).toEqual(['encoded-image']); }); + test('keeps a code output out of the prompt once its expired sandbox reference is cleared', async () => { + /* Priming clears a dead sandbox reference on the turn's copy of the record so the file + * is re-provisioned. The output still belongs to the sandbox that wrote it. */ + const message = {}; + const file = { + user: 'user1', + file_id: 'code-output-chart', + filename: 'chart.png', + filepath: '/uploads/chart.png', + type: 'image/png', + bytes: 100, + source: 'local', + context: 'execute_code', + metadata: {}, + }; + + const result = await TestClient.processAttachments(message, [file]); + + expect(result).toEqual([file]); + expect(TestClient.addImageURLs).not.toHaveBeenCalled(); + expect(message.image_urls).toBeUndefined(); + }); + test('keeps excluding embedded legacy files that have no delivery path', async () => { const message = {}; const file = { diff --git a/api/cache/logViolation.js b/api/cache/logViolation.js index 1ff65c6ccdd..fd2549c58dd 100644 --- a/api/cache/logViolation.js +++ b/api/cache/logViolation.js @@ -13,7 +13,7 @@ const banViolation = require('./banViolation'); * @param {number | string} [score=1] - The severity of the violation. Defaults to 1 */ const logViolation = async (req, res, type, errorMessage, score = 1) => { - const userId = req.user?.id ?? req.user?._id; + const userId = req.user?.id ?? req.user?._id?.toString(); if (!userId) { return; } diff --git a/api/package.json b/api/package.json index f226bf0f558..6c925bfad76 100644 --- a/api/package.json +++ b/api/package.json @@ -46,7 +46,7 @@ "@azure/storage-blob": "^12.30.0", "@google/genai": "^2.8.0", "@keyv/redis": "5.1.6", - "@librechat/agents": "^3.8.7", + "@librechat/agents": "^3.8.8", "@librechat/api": "*", "@librechat/data-schemas": "*", "@microsoft/microsoft-graph-client": "^3.0.7", diff --git a/api/server/controllers/agents/__tests__/openai.spec.js b/api/server/controllers/agents/__tests__/openai.spec.js index 6fac81d9768..b02664babc4 100644 --- a/api/server/controllers/agents/__tests__/openai.spec.js +++ b/api/server/controllers/agents/__tests__/openai.spec.js @@ -1355,7 +1355,10 @@ describe('OpenAIChatCompletionController', () => { const { loadAgentTools, loadToolsForExecution } = require('~/server/services/ToolService'); const { filterFilesByAgentAccess } = require('~/server/services/Files/permissions'); - req.config.endpoints.agents.backgroundTasks = { ordinaryToolCancellation: true }; + req.config.endpoints.agents.backgroundTasks = { + ordinaryToolCancellation: true, + completionResultMaxChars: 4096, + }; await OpenAIChatCompletionController(req, res); const [initializeParams, dbMethods] = initializeAgent.mock.calls.at(-1); @@ -1384,6 +1387,7 @@ describe('OpenAIChatCompletionController', () => { const toolExecuteOptions = createToolExecuteHandler.mock.calls.at(-1)[0]; expect(toolExecuteOptions.ordinaryToolCancellation).toBe(true); + expect(toolExecuteOptions.backgroundCompletionResultMaxChars).toBe(4096); expect(toolExecuteOptions.runSignal).toBe(mockExecution.signal); expect(toolExecuteOptions.foregroundRunId).toBe(initializeParams.requestBody.messageId); const effectiveSignal = new AbortController().signal; diff --git a/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js b/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js index 90b9c89a330..e4524db39ab 100644 --- a/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js +++ b/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js @@ -67,6 +67,8 @@ jest.mock('@librechat/api', () => ({ resolveRunCodeWorkspaces: jest.requireActual('@librechat/api').resolveRunCodeWorkspaces, shouldPersistCodeWorkspaceInitializationError: jest.requireActual('@librechat/api').shouldPersistCodeWorkspaceInitializationError, + resolvePersistableCodeEnvironmentDecision: (...args) => + jest.requireActual('@librechat/api').resolvePersistableCodeEnvironmentDecision(...args), getSafeErrorMetadata: jest.requireActual('@librechat/api').getSafeErrorMetadata, getSafeErrorText: jest.requireActual('@librechat/api').getSafeErrorText, GenerationJobManager: mockGenerationJobManager, diff --git a/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js b/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js index d778f6b2e58..eebe5f35186 100644 --- a/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js +++ b/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js @@ -281,6 +281,8 @@ jest.mock('@librechat/api', () => ({ jest.requireActual('@librechat/api').getCodeWorkspaceSelectionErrorDetails, shouldPersistCodeWorkspaceInitializationError: jest.requireActual('@librechat/api').shouldPersistCodeWorkspaceInitializationError, + resolvePersistableCodeEnvironmentDecision: (...args) => + jest.requireActual('@librechat/api').resolvePersistableCodeEnvironmentDecision(...args), getSafeErrorMetadata: jest.requireActual('@librechat/api').getSafeErrorMetadata, getSafeErrorText: jest.requireActual('@librechat/api').getSafeErrorText, resolveFailedTurnContent: jest.requireActual('@librechat/api').resolveFailedTurnContent, @@ -3999,6 +4001,68 @@ describe('ResumableAgentController resume metadata', () => { expect.objectContaining({ initialAgentId: null }), ); }); + + it('records the decision a failed turn of a saved chat ran under', async () => { + const codeWorkspaces = [{ environmentId: 'personal-vm', workspaceId: 'project-a' }]; + const initializeClient = jest.fn().mockImplementation(async ({ req: request }) => { + request._codeEnvironmentDecision = { mode: 'attached', codeWorkspaces }; + throw new Error('model unavailable'); + }); + + await AgentController( + createFailedRequest(), + createResumableResponse(), + jest.fn(), + initializeClient, + null, + ); + + expect(mockSaveConvo).toHaveBeenCalledWith( + expect.objectContaining({ userId: 'user-123' }), + { conversationId, codeEnvironmentMode: 'attached', codeWorkspaces }, + expect.objectContaining({ noUpsert: true }), + ); + }); + + it('does not rewrite the decision a chat already holds', async () => { + const req = createFailedRequest(); + req.resolvedConversation = { + conversationId, + codeEnvironmentMode: 'attached', + codeWorkspaces: [{ environmentId: 'personal-vm', workspaceId: 'project-a' }], + }; + const initializeClient = jest.fn().mockImplementation(async ({ req: request }) => { + request._codeEnvironmentDecision = { + mode: 'attached', + codeWorkspaces: [{ environmentId: 'personal-vm', workspaceId: 'project-b' }], + }; + throw new Error('model unavailable'); + }); + + await AgentController(req, createResumableResponse(), jest.fn(), initializeClient, null); + + expect(mockSaveConvo).toHaveBeenCalledWith( + expect.objectContaining({ userId: 'user-123' }), + { conversationId }, + expect.objectContaining({ noUpsert: true }), + ); + }); + + it('leaves a saved chat undecided when the turn resolved no decision', async () => { + await AgentController( + createFailedRequest(), + createResumableResponse(), + jest.fn(), + jest.fn().mockRejectedValue(new Error('model unavailable')), + null, + ); + + expect(mockSaveConvo).toHaveBeenCalledWith( + expect.objectContaining({ userId: 'user-123' }), + { conversationId }, + expect.objectContaining({ noUpsert: true }), + ); + }); }); it('finalizes the failed job before releasing the idempotency claim', async () => { diff --git a/api/server/controllers/agents/__tests__/responses.unit.spec.js b/api/server/controllers/agents/__tests__/responses.unit.spec.js index 45d00b56842..70f0e937545 100644 --- a/api/server/controllers/agents/__tests__/responses.unit.spec.js +++ b/api/server/controllers/agents/__tests__/responses.unit.spec.js @@ -2065,7 +2065,10 @@ describe('createResponse controller', () => { const { loadAgentTools, loadToolsForExecution } = require('~/server/services/ToolService'); const { filterFilesByAgentAccess } = require('~/server/services/Files/permissions'); - req.config.endpoints.agents.backgroundTasks = { ordinaryToolCancellation: true }; + req.config.endpoints.agents.backgroundTasks = { + ordinaryToolCancellation: true, + completionResultMaxChars: 4096, + }; req.body.stream = stream; await createResponse(req, res); @@ -2095,6 +2098,7 @@ describe('createResponse controller', () => { const toolExecuteOptions = createToolExecuteHandler.mock.calls.at(-1)[0]; expect(toolExecuteOptions.ordinaryToolCancellation).toBe(true); + expect(toolExecuteOptions.backgroundCompletionResultMaxChars).toBe(4096); expect(toolExecuteOptions.runSignal).toBe(mockExecution.signal); expect(toolExecuteOptions.foregroundRunId).toBe(initializeParams.requestBody.messageId); const effectiveSignal = new AbortController().signal; diff --git a/api/server/controllers/agents/client.js b/api/server/controllers/agents/client.js index d11da051856..451fcf83fcd 100644 --- a/api/server/controllers/agents/client.js +++ b/api/server/controllers/agents/client.js @@ -1989,6 +1989,15 @@ class AgentClient extends BaseClient { ); } + /** Attachments alone defer only the message, so a new conversation still gets its row when + * the run starts, as it did before that deferral. A content policy holds back every write. */ + shouldSeedDeferredConversation() { + return !hasModelBoundContentProtection( + this.options.req?.config?.filters, + this.options.req?.config?.messageFilter?.pii, + ); + } + /** Legacy `messageFilter.pii` historically covered the restored branch * before model-input construction and persistence. Retain that contract * without scanning new source-aware filters before SDK pruning. */ diff --git a/api/server/controllers/agents/client.test.js b/api/server/controllers/agents/client.test.js index 82ccb23afd7..c646df0436d 100644 --- a/api/server/controllers/agents/client.test.js +++ b/api/server/controllers/agents/client.test.js @@ -6314,6 +6314,26 @@ describe('AgentClient - titleConvo', () => { expect(client.shouldDeferUserMessagePersistence()).toBe(true); }); + it('still seeds the conversation row when only attachments defer the message', () => { + client.modelBoundCurrentFiles = [makeTextFile('pending', 'pending.txt', 'context')]; + + expect(client.shouldDeferUserMessagePersistence()).toBe(true); + expect(client.shouldSeedDeferredConversation()).toBe(true); + }); + + it('holds back the conversation row while a content policy defers every write', () => { + client.modelBoundCurrentFiles = [makeTextFile('pending', 'pending.txt', 'context')]; + mockReq.config.messageFilter = { + pii: { + starterPatterns: [], + customPatterns: [{ id: 'secret', label: 'secret', regex: 'SECRET-[A-Z]+' }], + }, + }; + + expect(client.shouldDeferUserMessagePersistence()).toBe(true); + expect(client.shouldSeedDeferredConversation()).toBe(false); + }); + it('keeps repeated lazy scoped-text admission cumulative across resolutions', () => { mockReq.config.fileConfig = { fileContextCharLimit: 1_000_000 }; const repeated = makeTextFile('lazy-context', 'lazy.txt', 'x'.repeat(600_000)); diff --git a/api/server/controllers/agents/openai.js b/api/server/controllers/agents/openai.js index 7bcbe4b9fae..65e22305943 100644 --- a/api/server/controllers/agents/openai.js +++ b/api/server/controllers/agents/openai.js @@ -490,6 +490,8 @@ const executeOpenAIChatCompletion = async (envelope, { req, res }) => { const allowedProviders = new Set(agentsEConfig?.allowedProviders); const ordinaryToolCancellationEnabled = agentsEConfig?.backgroundTasks?.ordinaryToolCancellation === true; + const backgroundCompletionResultMaxChars = + agentsEConfig?.backgroundTasks?.completionResultMaxChars; // Create tool loader const loadTools = createToolLoader({ req, res, signal: execution.signal }); @@ -840,6 +842,7 @@ const executeOpenAIChatCompletion = async (envelope, { req, res }) => { runSignal: execution.signal, foregroundRunId: responseId, ordinaryToolCancellation: ordinaryToolCancellationEnabled, + backgroundCompletionResultMaxChars, provisionFiles: createProvisionFilesCallback({ req, agentToolContexts, diff --git a/api/server/controllers/agents/request.js b/api/server/controllers/agents/request.js index 9f9b5935d74..c1a4cfb3894 100644 --- a/api/server/controllers/agents/request.js +++ b/api/server/controllers/agents/request.js @@ -45,6 +45,7 @@ const { logAgentMemorySnapshot, getCodeWorkspaceSelectionErrorDetails, shouldPersistCodeWorkspaceInitializationError, + resolvePersistableCodeEnvironmentDecision, getFailedTurnTraceFields, resolveFailedTurnContent, } = require('@librechat/api'); @@ -502,7 +503,16 @@ async function saveErrorTurn( const agentId = endpointOption?.agent_id ?? req.body?.agent_id; const chatProjectId = endpointOption?.chatProjectId ?? req.body?.chatProjectId; const seedConvo = isNewConvo || req.resolvedConversation === null; - const codeEnvironmentDecision = req._codeEnvironmentDecision; + /** A stored turn seals the decision it ran under, on a saved chat as much as on a new one: the + * error turn below enters the conversation, so leaving its validated decision out would let a + * retry choose a different workspace than the failure already recorded. The resolver the + * streaming saves already use decides what this turn may write, so an error turn and a + * completed one record a decision under one rule. */ + const decisionFields = resolvePersistableCodeEnvironmentDecision({ + conversationId, + decision: req._codeEnvironmentDecision, + conversation: req.resolvedConversation, + }); const convoFields = seedConvo ? { ...(endpoint != null && { endpoint }), @@ -514,14 +524,9 @@ async function saveErrorTurn( ...(endpointOption?.spec != null && { spec: endpointOption.spec }), ...(agentId != null && { agent_id: agentId }), ...(typeof chatProjectId === 'string' && chatProjectId.length > 0 && { chatProjectId }), - ...(codeEnvironmentDecision?.mode != null && { - codeEnvironmentMode: codeEnvironmentDecision.mode, - ...(codeEnvironmentDecision.codeWorkspaces != null && { - codeWorkspaces: codeEnvironmentDecision.codeWorkspaces, - }), - }), + ...decisionFields, } - : {}; + : decisionFields; await saveConvo( reqCtx, { conversationId, ...convoFields }, diff --git a/api/server/controllers/agents/responses.js b/api/server/controllers/agents/responses.js index 24e668576b1..1b4772ffd07 100644 --- a/api/server/controllers/agents/responses.js +++ b/api/server/controllers/agents/responses.js @@ -732,6 +732,8 @@ const executeResponse = async (envelope, { req, res }) => { const agentsEConfig = appConfig?.endpoints?.[EModelEndpoint.agents]; const ordinaryToolCancellationEnabled = agentsEConfig?.backgroundTasks?.ordinaryToolCancellation === true; + const backgroundCompletionResultMaxChars = + agentsEConfig?.backgroundTasks?.completionResultMaxChars; const previousMessages = request.previous_response_id ? await loadPreviousMessages(request.previous_response_id, principal.userId) : []; @@ -1194,6 +1196,7 @@ const executeResponse = async (envelope, { req, res }) => { runSignal: execution.signal, foregroundRunId: responseId, ordinaryToolCancellation: ordinaryToolCancellationEnabled, + backgroundCompletionResultMaxChars, provisionFiles: createProvisionFilesCallback({ req, agentToolContexts, @@ -1424,6 +1427,7 @@ const executeResponse = async (envelope, { req, res }) => { runSignal: execution.signal, foregroundRunId: responseId, ordinaryToolCancellation: ordinaryToolCancellationEnabled, + backgroundCompletionResultMaxChars, provisionFiles: createProvisionFilesCallback({ req, agentToolContexts, diff --git a/api/server/experimental.js b/api/server/experimental.js index ef3de373cea..d6b610a5956 100644 --- a/api/server/experimental.js +++ b/api/server/experimental.js @@ -705,7 +705,11 @@ if (cluster.isMaster) { await initializeMCPs(); await initializeOAuthReconnectManager(); await checkMigrations(); - await initializeAgentTriggerService({ address: server.address() }); + await initializeAgentTriggerService({ + address: server.address(), + completionResultBatchSize: + baseAppConfig?.endpoints?.agents?.backgroundTasks?.completionResultBatchSize, + }); } catch (initErr) { logger.error(`Worker ${process.pid} post-listen initialization failed:`, initErr); process.exit(1); diff --git a/api/server/index.js b/api/server/index.js index 610d5c2c9a2..b4555e7b262 100644 --- a/api/server/index.js +++ b/api/server/index.js @@ -498,7 +498,11 @@ const startServer = async () => { if (inspectFlags || isEnabled(process.env.MEM_DIAG)) { memoryDiagnostics.start(); } - await initializeAgentTriggerService({ address: server.address() }); + await initializeAgentTriggerService({ + address: server.address(), + completionResultBatchSize: + appConfig?.endpoints?.agents?.backgroundTasks?.completionResultBatchSize, + }); const scheduleEngineArmed = (await initializeScheduleEngine()) != null; scheduleEngineState = scheduleEngineArmed ? 'armed' : 'unavailable'; if (!scheduleEngineArmed) { diff --git a/api/server/middleware/checkBan.js b/api/server/middleware/checkBan.js index 7194a555b0e..0871e51c047 100644 --- a/api/server/middleware/checkBan.js +++ b/api/server/middleware/checkBan.js @@ -88,7 +88,7 @@ const checkBan = async (req, res, next = () => {}) => { } req.ip = removePorts(req); - let userId = req.user?.id ?? req.user?._id ?? null; + let userId = req.user?.id ?? req.user?._id?.toString() ?? null; if (!userId && req?.body?.email) { const user = await findUser({ email: req.body.email }, '_id'); diff --git a/api/server/middleware/checkBan.spec.js b/api/server/middleware/checkBan.spec.js new file mode 100644 index 00000000000..ee83f0bef8f --- /dev/null +++ b/api/server/middleware/checkBan.spec.js @@ -0,0 +1,96 @@ +const mongoose = require('mongoose'); +const { MongoMemoryServer } = require('mongodb-memory-server'); +const { ErrorTypes, ViolationTypes } = require('librechat-data-provider'); + +/** Violation logs are file-backed in production; keep the real namespaced Keyv, swap only the store. */ +jest.mock('~/cache/getLogStores', () => { + const { Keyv } = jest.requireActual('keyv'); + const { ViolationTypes } = jest.requireActual('librechat-data-provider'); + const getLogStores = jest.requireActual('~/cache/getLogStores'); + const violationLogs = new Map(); + return (type) => { + if (type === ViolationTypes.BAN) { + return getLogStores(type); + } + if (!violationLogs.has(type)) { + const namespace = type === ViolationTypes.GENERAL ? 'violations' : `violations:${type}`; + violationLogs.set(type, new Keyv({ store: new Map(), namespace })); + } + return violationLogs.get(type); + }; +}); + +jest.mock('~/models', () => ({ + ...jest.requireActual('~/models'), + deleteAllUserSessions: jest.fn().mockResolvedValue(true), +})); + +process.env.BAN_VIOLATIONS = 'true'; +process.env.BAN_INTERVAL = '20'; +delete process.env.USE_REDIS; + +const logViolation = require('~/cache/logViolation'); +const checkBan = require('./checkBan'); + +/** Passport's social strategies hand the callback a lean user: an ObjectId `_id` and no `id`. */ +const createOAuthCallbackReq = (userId, ip) => ({ + ip, + user: { _id: userId }, + method: 'GET', + headers: { 'user-agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) Chrome/140.0.0.0' }, + body: {}, + baseUrl: '/oauth', + originalUrl: '/oauth/google/callback', + isOAuthNavigation: true, +}); + +const createRes = () => ({ + status: jest.fn().mockReturnThis(), + json: jest.fn().mockReturnThis(), + redirect: jest.fn().mockReturnThis(), + clearCookie: jest.fn(), +}); + +describe('checkBan with namespaced Keyv stores and ObjectId user ids', () => { + let mongoServer; + + beforeAll(async () => { + mongoServer = await MongoMemoryServer.create(); + await mongoose.connect(mongoServer.getUri()); + }); + + afterAll(async () => { + await mongoose.disconnect(); + await mongoServer.stop(); + }); + + it('lets an unbanned OAuth user through without Redis', async () => { + const next = jest.fn(); + const req = createOAuthCallbackReq(new mongoose.Types.ObjectId(), '10.0.0.1'); + + await checkBan(req, createRes(), next); + + expect(next).toHaveBeenCalledTimes(1); + expect(next).toHaveBeenCalledWith(); + expect(req.banned).toBeUndefined(); + }); + + it('enforces a ban recorded for the same user from another address', async () => { + const userId = new mongoose.Types.ObjectId(); + const violationReq = createOAuthCallbackReq(userId, '10.0.0.2'); + const errorMessage = { type: ViolationTypes.LOGINS }; + + await logViolation(violationReq, createRes(), ViolationTypes.LOGINS, errorMessage, 20); + expect(errorMessage.ban).toBe(true); + + const next = jest.fn(); + const res = createRes(); + const req = createOAuthCallbackReq(new mongoose.Types.ObjectId(userId.toString()), '10.0.0.3'); + + await checkBan(req, res, next); + + expect(next).not.toHaveBeenCalled(); + expect(req.banned).toBe(true); + expect(res.redirect).toHaveBeenCalledWith(expect.stringContaining(ErrorTypes.AUTH_BANNED)); + }); +}); diff --git a/api/server/routes/files/files.js b/api/server/routes/files/files.js index de84c67c950..f6b97e587f4 100644 --- a/api/server/routes/files/files.js +++ b/api/server/routes/files/files.js @@ -8,6 +8,7 @@ const { refreshS3FileUrls, handleFilesUsageRequest, buildDeleteFilesResponse, + deleteAgentResourceFiles, shouldUseUploadSse, startUploadSseStream, sendUploadPolicyError, @@ -245,20 +246,36 @@ router.delete('/', async (req, res) => { }); } - const toolResourceFiles = agent.tool_resources?.[req.body.tool_resource]?.file_ids ?? []; - const agentFiles = files - .filter((f) => toolResourceFiles.includes(f.file_id)) - .map((file) => ({ tool_resource: req.body.tool_resource, file_id: file.file_id })); - if (agentFiles.length === 0) { + const agentDeletion = await deleteAgentResourceFiles( + { + agentId: req.body.agent_id, + agentObjectId: agent._id.toString(), + toolResource: req.body.tool_resource, + requestedFileIds: fileIds, + attachedFileIds: agent.tool_resources?.[req.body.tool_resource]?.file_ids ?? [], + files: dbFiles.map((file) => ({ + file_id: file.file_id, + owner: file.user?.toString() ?? null, + file, + })), + userId: req.user.id.toString(), + }, + { + getSharedResourceFileIds: db.getSharedResourceFileIds, + removeAgentResourceFiles: db.removeAgentResourceFiles, + deleteFiles: (agentFiles) => processDeleteRequest({ req, files: agentFiles }), + }, + ); + + if (agentDeletion.outcome == null) { res.status(200).json({ message: 'File associations removed successfully from agent' }); return; } - await db.removeAgentResourceFiles({ - agent_id: req.body.agent_id, - files: agentFiles, - }); - res.status(200).json({ message: 'File associations removed successfully from agent' }); + logger.debug( + `[/files] Agent files deleted successfully: ${agentDeletion.destroyedFileIds.join(', ')}`, + ); + sendDeleteResult(agentDeletion.outcome, 'Files deleted successfully'); return; } diff --git a/api/server/routes/files/files.test.js b/api/server/routes/files/files.test.js index 2ad58d350dd..9df2acb0623 100644 --- a/api/server/routes/files/files.test.js +++ b/api/server/routes/files/files.test.js @@ -327,6 +327,245 @@ describe('File Routes - Delete with Agent Access', () => { expect(updatedAgent.tool_resources.file_search.file_ids).toEqual([]); }); + it('deletes storage and embeddings for an attached file the caller owns', async () => { + const ownedFileId = uuidv4(); + await createFile({ + user: otherUserId, + file_id: ownedFileId, + filename: 'owned-knowledge.txt', + filepath: '/uploads/owned-knowledge.txt', + bytes: 100, + type: 'text/plain', + source: FileSources.vectordb, + embedded: true, + }); + + const agent = await createAgent({ + id: uuidv4(), + name: 'Test Agent', + provider: 'openai', + model: 'gpt-4', + author: otherUserId, + tool_resources: { + file_search: { + file_ids: [ownedFileId], + }, + }, + }); + + const response = await request(app) + .delete('/files') + .send({ + agent_id: agent.id, + tool_resource: 'file_search', + files: [{ file_id: ownedFileId, filepath: '/uploads/owned-knowledge.txt' }], + }); + + expect(response.status).toBe(200); + expect(response.body.message).toBe('Files deleted successfully'); + expect(processDeleteRequest).toHaveBeenCalledTimes(1); + + const [{ req, files: deletedFiles }] = processDeleteRequest.mock.calls[0]; + expect(deletedFiles.map((file) => file.file_id)).toEqual([ownedFileId]); + expect(deletedFiles[0].source).toBe(FileSources.vectordb); + expect(req.body.agent_id).toBe(agent.id); + expect(req.body.tool_resource).toBe('file_search'); + }); + + it('unlinks another user’s attached file while deleting the caller’s own', async () => { + const ownedFileId = uuidv4(); + await createFile({ + user: otherUserId, + file_id: ownedFileId, + filename: 'owned-knowledge.txt', + filepath: '/uploads/owned-knowledge.txt', + bytes: 100, + type: 'text/plain', + source: FileSources.vectordb, + embedded: true, + }); + + const agent = await createAgent({ + id: uuidv4(), + name: 'Test Agent', + provider: 'openai', + model: 'gpt-4', + author: otherUserId, + tool_resources: { + file_search: { + file_ids: [ownedFileId, fileId], + }, + }, + }); + + const response = await request(app) + .delete('/files') + .send({ + agent_id: agent.id, + tool_resource: 'file_search', + files: [ + { file_id: ownedFileId, filepath: '/uploads/owned-knowledge.txt' }, + { file_id: fileId, filepath: '/uploads/test.txt' }, + ], + }); + + expect(response.status).toBe(200); + expect(response.body.message).toBe('Files deleted successfully'); + + const [{ files: deletedFiles }] = processDeleteRequest.mock.calls[0]; + expect(deletedFiles.map((file) => file.file_id)).toEqual([ownedFileId]); + + const updatedAgent = await Agent.findOne({ id: agent.id }).lean(); + expect(updatedAgent.tool_resources.file_search.file_ids).toEqual([ownedFileId]); + + const retainedFile = await File.findOne({ file_id: fileId }).lean(); + expect(retainedFile).toBeTruthy(); + }); + + it('keeps a file the same agent holds under another tool resource', async () => { + const sharedFileId = uuidv4(); + await createFile({ + user: otherUserId, + file_id: sharedFileId, + filename: 'dual-purpose.txt', + filepath: '/uploads/dual-purpose.txt', + bytes: 100, + type: 'text/plain', + source: FileSources.vectordb, + embedded: true, + }); + + /* One agent can hold the same file under two resources, so the reference being removed is the + `(agent, tool_resource)` pair rather than the agent. */ + const agent = await createAgent({ + id: uuidv4(), + name: 'Test Agent', + provider: 'openai', + model: 'gpt-4', + author: otherUserId, + tool_resources: { + file_search: { file_ids: [sharedFileId] }, + context: { file_ids: [sharedFileId] }, + }, + }); + + const response = await request(app) + .delete('/files') + .send({ + agent_id: agent.id, + tool_resource: 'file_search', + files: [{ file_id: sharedFileId, filepath: '/uploads/dual-purpose.txt' }], + }); + + expect(response.status).toBe(200); + expect(response.body.message).toBe('File associations removed successfully from agent'); + expect(processDeleteRequest).not.toHaveBeenCalled(); + + const updatedAgent = await Agent.findOne({ id: agent.id }).lean(); + expect(updatedAgent.tool_resources.file_search.file_ids).toEqual([]); + expect(updatedAgent.tool_resources.context.file_ids).toEqual([sharedFileId]); + + const retainedFile = await File.findOne({ file_id: sharedFileId }).lean(); + expect(retainedFile).toBeTruthy(); + }); + + it('keeps a file a duplicated agent still references, unlinking it here only', async () => { + const sharedFileId = uuidv4(); + await createFile({ + user: otherUserId, + file_id: sharedFileId, + filename: 'shared-knowledge.txt', + filepath: '/uploads/shared-knowledge.txt', + bytes: 100, + type: 'text/plain', + source: FileSources.vectordb, + embedded: true, + }); + + const agent = await createAgent({ + id: uuidv4(), + name: 'Test Agent', + provider: 'openai', + model: 'gpt-4', + author: otherUserId, + tool_resources: { file_search: { file_ids: [sharedFileId] } }, + }); + + /* Duplicating an agent copies file_ids rather than the files behind them, and lands them + under `context`, so the second holder is found across tool resources. */ + const duplicate = await createAgent({ + id: uuidv4(), + name: 'Test Agent (copy)', + provider: 'openai', + model: 'gpt-4', + author: otherUserId, + tool_resources: { context: { file_ids: [sharedFileId] } }, + }); + + const response = await request(app) + .delete('/files') + .send({ + agent_id: agent.id, + tool_resource: 'file_search', + files: [{ file_id: sharedFileId, filepath: '/uploads/shared-knowledge.txt' }], + }); + + expect(response.status).toBe(200); + expect(response.body.message).toBe('File associations removed successfully from agent'); + expect(processDeleteRequest).not.toHaveBeenCalled(); + + const updatedAgent = await Agent.findOne({ id: agent.id }).lean(); + expect(updatedAgent.tool_resources.file_search.file_ids).toEqual([]); + + const untouchedDuplicate = await Agent.findOne({ id: duplicate.id }).lean(); + expect(untouchedDuplicate.tool_resources.context.file_ids).toEqual([sharedFileId]); + + const retainedFile = await File.findOne({ file_id: sharedFileId }).lean(); + expect(retainedFile).toBeTruthy(); + }); + + it('leaves an owned file alone when the tool resource does not hold it', async () => { + const ownedFileId = uuidv4(); + await createFile({ + user: otherUserId, + file_id: ownedFileId, + filename: 'detached-knowledge.txt', + filepath: '/uploads/detached-knowledge.txt', + bytes: 100, + type: 'text/plain', + source: FileSources.vectordb, + embedded: true, + }); + + const agent = await createAgent({ + id: uuidv4(), + name: 'Test Agent', + provider: 'openai', + model: 'gpt-4', + author: otherUserId, + tool_resources: { + file_search: { + file_ids: [fileId], + }, + }, + }); + + const response = await request(app) + .delete('/files') + .send({ + agent_id: agent.id, + tool_resource: 'file_search', + files: [{ file_id: ownedFileId, filepath: '/uploads/detached-knowledge.txt' }], + }); + + expect(response.status).toBe(200); + expect(response.body.message).toBe('File associations removed successfully from agent'); + expect(processDeleteRequest).not.toHaveBeenCalled(); + + const updatedAgent = await Agent.findOne({ id: agent.id }).lean(); + expect(updatedAgent.tool_resources.file_search.file_ids).toEqual([fileId]); + }); + it('rejects invalid agent tool_resource values before unlinking', async () => { const agent = await createAgent({ id: uuidv4(), diff --git a/api/server/routes/skills.test.js b/api/server/routes/skills.test.js index 3508cb512dd..992cb2d8497 100644 --- a/api/server/routes/skills.test.js +++ b/api/server/routes/skills.test.js @@ -553,6 +553,54 @@ describe('Skill routes', () => { ); }); + it('rolls back the skill, files, ACL, and stored blob when a bundled file fails', async () => { + const saveBuffer = jest.fn().mockResolvedValue('/uploads/test/queries.sql'); + const deleteFile = jest.fn().mockResolvedValue(undefined); + const { getStrategyFunctions } = require('~/server/services/Files/strategies'); + /** Both the write and the rollback delete resolve a strategy, so the + * override outlives a single call and is restored below. */ + const defaultStrategies = getStrategyFunctions(); + getStrategyFunctions.mockReturnValue({ saveBuffer, deleteFile }); + + const zip = new JSZip(); + zip.file( + 'SKILL.md', + [ + '---', + 'name: partial-import', + 'description: Imported skill that loses one bundled file.', + '---', + '# Partial Import', + ].join('\n'), + ); + zip.file('queries.sql', 'select 1;'); + /** A space is outside the stored path charset, so this entry can never + * persist — the import must fail instead of dropping it silently. */ + zip.file('references/region mapping.md', '# regions'); + const buffer = await zip.generateAsync({ type: 'nodebuffer' }); + + const res = await request(app).post('/api/skills/import').attach('file', buffer, { + filename: 'partial-import.skill', + contentType: 'application/zip', + }); + getStrategyFunctions.mockReturnValue(defaultStrategies); + + expect(res.status).toBe(422); + expect(res.body).toEqual( + expect.objectContaining({ + error: 'skill_import_incomplete', + failedFiles: [{ path: 'references/region mapping.md', reason: 'invalid_path' }], + }), + ); + expect(await Skill.countDocuments({ name: 'partial-import' })).toBe(0); + expect(await SkillFile.countDocuments()).toBe(0); + expect(await AclEntry.countDocuments({ resourceType: ResourceType.SKILL })).toBe(0); + expect(deleteFile).toHaveBeenCalledWith( + expect.anything(), + expect.objectContaining({ filepath: '/uploads/test/queries.sql' }), + ); + }); + it('blocks filtered Markdown before creating the imported skill', async () => { mockFilters = { files: { diff --git a/api/server/services/Agents/triggers.js b/api/server/services/Agents/triggers.js index 66926c2991a..0f2d81218e7 100644 --- a/api/server/services/Agents/triggers.js +++ b/api/server/services/Agents/triggers.js @@ -28,6 +28,7 @@ const subagentCompletionAdapter = createSubagentCompletionWakeupResolver({ const backgroundToolCompletionAdapter = createBackgroundToolCompletionWakeupResolver({ methods, getGenerationJob: (conversationId) => GenerationJobManager.getJob(conversationId), + getResultBatchSize: () => service.getBackgroundCompletionResultBatchSize(), }); const eventActorAdapter = createAgentEventContinueResolver({ methods, @@ -86,6 +87,9 @@ module.exports = { requeueAgentTrigger: service.requeue, retireAgentTrigger: service.retire, renewAgentTriggerProducerLease: service.renewProducerLease, + persistAgentBackgroundToolResult: service.persistBackgroundToolResult, + getAgentBackgroundToolResultClaim: service.getBackgroundToolResultClaim, + releaseAgentBackgroundToolResultClaims: service.releaseBackgroundToolResultClaims, drainAgentTriggerDeliveriesForUser: service.drainUser, prepareAgentTriggerUserPurge: service.prepareUserPurge, cancelAgentTriggerUserPurge: service.cancelUserPurge, diff --git a/api/server/services/Agents/triggers.spec.js b/api/server/services/Agents/triggers.spec.js index 16965b5abe8..4c636f46928 100644 --- a/api/server/services/Agents/triggers.spec.js +++ b/api/server/services/Agents/triggers.spec.js @@ -1,4 +1,5 @@ const mockCreateAgentTriggerService = jest.fn(); +const mockCreateBackgroundToolCompletionWakeupResolver = jest.fn(() => jest.fn()); const mockGenerationJobManager = { supportsDetachedAgentEventActions: true, getJob: jest.fn(), @@ -21,7 +22,8 @@ jest.mock('@librechat/api', () => ({ createAgentTriggerService: (...args) => mockCreateAgentTriggerService(...args), createAgentContinuationResolver: jest.fn(() => jest.fn()), createAgentEventContinueResolver: jest.fn(() => jest.fn()), - createBackgroundToolCompletionWakeupResolver: jest.fn(() => jest.fn()), + createBackgroundToolCompletionWakeupResolver: (...args) => + mockCreateBackgroundToolCompletionWakeupResolver(...args), createSubagentCompletionWakeupResolver: jest.fn(() => jest.fn()), createAgentQueuedTurnLifecycle: jest.fn(() => mockQueuedTurnLifecycle), BACKGROUND_TOOL_COMPLETION_SOURCE: 'background-tool-completion', @@ -39,8 +41,11 @@ describe('agent trigger service composition', () => { jest.resetModules(); jest.clearAllMocks(); mockGenerationJobManager.supportsDetachedAgentEventActions = true; + let completionResultBatchSize = 8; mockCreateAgentTriggerService.mockReturnValue({ - initialize: jest.fn(), + initialize: jest.fn(async (options) => { + completionResultBatchSize = options.completionResultBatchSize ?? 8; + }), stop: jest.fn(), dispatch: jest.fn(), enqueue: jest.fn(), @@ -52,6 +57,7 @@ describe('agent trigger service composition', () => { prepareUserPurge: jest.fn(), cancelUserPurge: jest.fn(), purgeUser: jest.fn(), + getBackgroundCompletionResultBatchSize: () => completionResultBatchSize, }); }); @@ -73,4 +79,16 @@ describe('agent trigger service composition', () => { mockGenerationJobManager.supportsDetachedAgentEventActions = false; expect(supportsDetachedActionCompletion()).toBe(false); }); + + it('injects the configured background completion batch size', async () => { + const { initializeAgentTriggerService } = require('./triggers'); + await initializeAgentTriggerService({ address: 'local', completionResultBatchSize: 12 }); + + const resolverDeps = mockCreateBackgroundToolCompletionWakeupResolver.mock.calls[0][0]; + expect(resolverDeps.getResultBatchSize()).toBe(12); + expect(mockCreateAgentTriggerService.mock.results[0].value.initialize).toHaveBeenCalledWith({ + address: 'local', + completionResultBatchSize: 12, + }); + }); }); diff --git a/api/server/services/Endpoints/agents/backgroundCompletion.js b/api/server/services/Endpoints/agents/backgroundCompletion.js index 6fc0996e623..8ca1cae6efc 100644 --- a/api/server/services/Endpoints/agents/backgroundCompletion.js +++ b/api/server/services/Endpoints/agents/backgroundCompletion.js @@ -2,9 +2,13 @@ const { createBackgroundToolCompletionWakeupHandler, createBackgroundToolDeadClaimRecovery, createBackgroundToolResultHandler, + claimBackgroundToolResult: claimResult, } = require('@librechat/api'); const { enqueueAgentTrigger, + persistAgentBackgroundToolResult, + getAgentBackgroundToolResultClaim, + releaseAgentBackgroundToolResultClaims, renewAgentTriggerProducerLease, retireAgentTrigger, } = require('../../Agents/triggers'); @@ -13,12 +17,17 @@ const preregisterBackgroundToolCompletion = createBackgroundToolCompletionWakeup enqueueAgentTrigger, retireAgentTrigger, renewAgentTriggerProducerLease, + (deliveryKey, sourceId, result) => + persistAgentBackgroundToolResult({ deliveryKey, sourceId, result }), ); function createBackgroundToolResultPersistence({ req, updateToolCallResult }) { return createBackgroundToolResultHandler({ req, updateToolCallResult }); } +const claimBackgroundToolResult = (db, input) => + claimResult(db, getAgentBackgroundToolResultClaim, input); + function createDeadBackgroundToolClaimRecovery( releaseBackgroundToolResultClaims, getGenerationJob, @@ -29,11 +38,13 @@ function createDeadBackgroundToolClaimRecovery( releaseBackgroundToolResultClaims, getGenerationJob, fenceGenerationClaim, + releaseAgentBackgroundToolResultClaims, ); } module.exports = { preregisterBackgroundToolCompletion, createBackgroundToolResultPersistence, + claimBackgroundToolResult, createDeadBackgroundToolClaimRecovery, }; diff --git a/api/server/services/Endpoints/agents/initialize.js b/api/server/services/Endpoints/agents/initialize.js index ee8ebb207d8..875d085951b 100644 --- a/api/server/services/Endpoints/agents/initialize.js +++ b/api/server/services/Endpoints/agents/initialize.js @@ -95,6 +95,7 @@ const subagentThreadTaskStore = require('./subagentThreadStore'); const { preregisterBackgroundToolCompletion, createBackgroundToolResultPersistence, + claimBackgroundToolResult, createDeadBackgroundToolClaimRecovery, } = require('./backgroundCompletion'); const { logViolation } = require('~/cache'); @@ -215,6 +216,8 @@ const initializeClientWithProvider = async ({ const ordinaryToolCancellationEnabled = appConfig?.endpoints?.[EModelEndpoint.agents]?.backgroundTasks?.ordinaryToolCancellation === true; + const backgroundCompletionResultMaxChars = + appConfig?.endpoints?.[EModelEndpoint.agents]?.backgroundTasks?.completionResultMaxChars; /** The normal controller resolves this once for timestamp anchoring. Reuse * that trusted document for child-thread execution policy; resume and direct * callers fall back to the same owner-scoped lookup. */ @@ -440,6 +443,7 @@ const initializeClientWithProvider = async ({ runSignal: signal, foregroundRunId, ordinaryToolCancellation: ordinaryToolCancellationEnabled, + backgroundCompletionResultMaxChars, loadTools: async ( toolNames, agentId, @@ -515,7 +519,7 @@ const initializeClientWithProvider = async ({ req, updateToolCallResult: db.updateToolCallResult, }), - claim: db.claimBackgroundToolResults, + claim: (input) => claimBackgroundToolResult(db, input), recoverDeadClaim: createDeadBackgroundToolClaimRecovery( db.releaseBackgroundToolResultClaims, (conversationId) => GenerationJobManager.getJob(conversationId), diff --git a/api/server/services/Endpoints/agents/initialize.spec.js b/api/server/services/Endpoints/agents/initialize.spec.js index b01af987604..fe65c2493ab 100644 --- a/api/server/services/Endpoints/agents/initialize.spec.js +++ b/api/server/services/Endpoints/agents/initialize.spec.js @@ -408,10 +408,11 @@ describe('initializeClient — processAgent ACL gate', () => { endpointOption: makeEndpointOption(), }); expect(capturedToolExecuteOptions.ordinaryToolCancellation).toBe(false); + expect(capturedToolExecuteOptions.backgroundCompletionResultMaxChars).toBeUndefined(); const enabledReq = makeReq(); enabledReq.config.endpoints.agents = { - backgroundTasks: { ordinaryToolCancellation: true }, + backgroundTasks: { ordinaryToolCancellation: true, completionResultMaxChars: 4096 }, }; await initializeClient({ req: enabledReq, @@ -420,6 +421,7 @@ describe('initializeClient — processAgent ACL gate', () => { endpointOption: makeEndpointOption(), }); expect(capturedToolExecuteOptions.ordinaryToolCancellation).toBe(true); + expect(capturedToolExecuteOptions.backgroundCompletionResultMaxChars).toBe(4096); }); it('propagates an expected-MCP-tools failure from the runtime tool loader', async () => { diff --git a/api/server/services/Files/Code/process.js b/api/server/services/Files/Code/process.js index ed800aaf37a..73efab77009 100644 --- a/api/server/services/Files/Code/process.js +++ b/api/server/services/Files/Code/process.js @@ -1112,6 +1112,7 @@ async function readSandboxFile({ * @param {Object} params * @param {string} params.file_path * @param {string} params.workspace_id + * @param {string} [params.workspace_instance_id] * @param {number} params.start_line * @param {number} params.max_lines * @param {string} params.codeApiBaseUrl @@ -1123,6 +1124,7 @@ async function readSandboxFile({ async function readWorkspaceFile({ file_path, workspace_id, + workspace_instance_id, start_line, max_lines, codeApiBaseUrl, @@ -1130,18 +1132,21 @@ async function readWorkspaceFile({ bridgeWorkerId, req, signal, + maxQueueWaitMs, }) { - const authHeaders = await getCodeApiAuthHeaders(req, bridgeWorkerId); return executeWorkspaceTool({ baseURL: codeApiBaseUrl, - authHeaders: { - ...authHeaders, + maxQueueWaitMs, + /** Minted per admission attempt: a queued call outlives one token TTL. */ + authHeaders: async () => ({ + ...(await getCodeApiAuthHeaders(req, bridgeWorkerId)), ...codeExecutionHeaders({ executionProfile, bridgeWorkerId }), - }, + }), request: { protocolVersion: 1, operation: 'read_file', workspaceId: workspace_id, + ...(workspace_instance_id ? { workspaceInstanceId: workspace_instance_id } : {}), path: file_path, startLine: start_line, maxLines: max_lines, @@ -1156,6 +1161,7 @@ async function readWorkspaceFile({ * @param {Object} params * @param {string} params.query * @param {string} params.workspace_id + * @param {string} [params.workspace_instance_id] * @param {string} [params.path] * @param {number} params.max_results * @param {string} params.codeApiBaseUrl @@ -1167,6 +1173,7 @@ async function readWorkspaceFile({ async function searchWorkspace({ query, workspace_id, + workspace_instance_id, path, max_results, codeApiBaseUrl, @@ -1174,18 +1181,21 @@ async function searchWorkspace({ bridgeWorkerId, req, signal, + maxQueueWaitMs, }) { - const authHeaders = await getCodeApiAuthHeaders(req, bridgeWorkerId); return executeWorkspaceTool({ baseURL: codeApiBaseUrl, - authHeaders: { - ...authHeaders, + maxQueueWaitMs, + /** Minted per admission attempt: a queued call outlives one token TTL. */ + authHeaders: async () => ({ + ...(await getCodeApiAuthHeaders(req, bridgeWorkerId)), ...codeExecutionHeaders({ executionProfile, bridgeWorkerId }), - }, + }), request: { protocolVersion: 1, operation: 'search_text', workspaceId: workspace_id, + ...(workspace_instance_id ? { workspaceInstanceId: workspace_instance_id } : {}), query, ...(path ? { path } : {}), maxResults: max_results, @@ -1199,6 +1209,7 @@ async function searchWorkspace({ * * @param {Object} params * @param {string} params.workspace_id + * @param {string} [params.workspace_instance_id] * @param {string} [params.path] * @param {string} [params.after_path] * @param {number} params.max_results @@ -1210,6 +1221,7 @@ async function searchWorkspace({ */ async function listWorkspaceFiles({ workspace_id, + workspace_instance_id, path, after_path, max_results, @@ -1218,18 +1230,21 @@ async function listWorkspaceFiles({ bridgeWorkerId, req, signal, + maxQueueWaitMs, }) { - const authHeaders = await getCodeApiAuthHeaders(req, bridgeWorkerId); return executeWorkspaceTool({ baseURL: codeApiBaseUrl, - authHeaders: { - ...authHeaders, + maxQueueWaitMs, + /** Minted per admission attempt: a queued call outlives one token TTL. */ + authHeaders: async () => ({ + ...(await getCodeApiAuthHeaders(req, bridgeWorkerId)), ...codeExecutionHeaders({ executionProfile, bridgeWorkerId }), - }, + }), request: { protocolVersion: 1, operation: 'list_files', workspaceId: workspace_id, + ...(workspace_instance_id ? { workspaceInstanceId: workspace_instance_id } : {}), ...(path ? { path } : {}), ...(after_path ? { afterPath: after_path } : {}), maxResults: max_results, @@ -1244,23 +1259,27 @@ async function writeWorkspaceFile({ content, overwrite, workspace_id, + workspace_instance_id, codeApiBaseUrl, executionProfile, bridgeWorkerId, req, signal, + maxQueueWaitMs, }) { - const authHeaders = await getCodeApiAuthHeaders(req, bridgeWorkerId); return executeWorkspaceTool({ baseURL: codeApiBaseUrl, - authHeaders: { - ...authHeaders, + maxQueueWaitMs, + /** Minted per admission attempt: a queued call outlives one token TTL. */ + authHeaders: async () => ({ + ...(await getCodeApiAuthHeaders(req, bridgeWorkerId)), ...codeExecutionHeaders({ executionProfile, bridgeWorkerId }), - }, + }), request: { protocolVersion: 1, operation: 'write_file', workspaceId: workspace_id, + ...(workspace_instance_id ? { workspaceInstanceId: workspace_instance_id } : {}), path: file_path, content, overwrite, @@ -1275,23 +1294,27 @@ async function editWorkspaceFile({ edits, expected_base_sha256, workspace_id, + workspace_instance_id, codeApiBaseUrl, executionProfile, bridgeWorkerId, req, signal, + maxQueueWaitMs, }) { - const authHeaders = await getCodeApiAuthHeaders(req, bridgeWorkerId); return executeWorkspaceTool({ baseURL: codeApiBaseUrl, - authHeaders: { - ...authHeaders, + maxQueueWaitMs, + /** Minted per admission attempt: a queued call outlives one token TTL. */ + authHeaders: async () => ({ + ...(await getCodeApiAuthHeaders(req, bridgeWorkerId)), ...codeExecutionHeaders({ executionProfile, bridgeWorkerId }), - }, + }), request: { protocolVersion: 1, operation: 'edit_file', workspaceId: workspace_id, + ...(workspace_instance_id ? { workspaceInstanceId: workspace_instance_id } : {}), path: file_path, edits, ...(expected_base_sha256 ? { expectedBaseSha256: expected_base_sha256 } : {}), @@ -1305,23 +1328,27 @@ async function previewWorkspaceEdit({ file_path, edits, workspace_id, + workspace_instance_id, codeApiBaseUrl, executionProfile, bridgeWorkerId, req, signal, + maxQueueWaitMs, }) { - const authHeaders = await getCodeApiAuthHeaders(req, bridgeWorkerId); return executeWorkspaceTool({ baseURL: codeApiBaseUrl, - authHeaders: { - ...authHeaders, + maxQueueWaitMs, + /** Minted per admission attempt: a queued call outlives one token TTL. */ + authHeaders: async () => ({ + ...(await getCodeApiAuthHeaders(req, bridgeWorkerId)), ...codeExecutionHeaders({ executionProfile, bridgeWorkerId }), - }, + }), request: { protocolVersion: 1, operation: 'preview_edit', workspaceId: workspace_id, + ...(workspace_instance_id ? { workspaceInstanceId: workspace_instance_id } : {}), path: file_path, edits, }, diff --git a/api/server/services/Files/Code/process.spec.js b/api/server/services/Files/Code/process.spec.js index 0a87df47bcc..176bf2229b6 100644 --- a/api/server/services/Files/Code/process.spec.js +++ b/api/server/services/Files/Code/process.spec.js @@ -2050,6 +2050,7 @@ describe('Code Process', () => { readWorkspaceFile({ file_path: 'src/app.ts', workspace_id: 'primary', + workspace_instance_id: 'a'.repeat(64), start_line: 1, max_lines: 200, codeApiBaseUrl: 'https://attached-code.example.com/v1', @@ -2057,21 +2058,32 @@ describe('Code Process', () => { bridgeWorkerId: 'worker-user-1', req: mockReq, signal: controller.signal, + maxQueueWaitMs: 0, }), ).resolves.toBe(result); - expect(getCodeApiAuthHeaders).toHaveBeenCalledWith(mockReq, 'worker-user-1'); + expect(getCodeApiAuthHeaders).not.toHaveBeenCalled(); + const { authHeaders } = mockExecuteWorkspaceTool.mock.calls[0][0]; + await expect(authHeaders()).resolves.toEqual({ + Authorization: 'Bearer workspace-token', + 'X-CodeAPI-Expected-Profile': 'stateful', + 'X-LibreChat-Code-Worker-ID': 'worker-user-1', + }); + getCodeApiAuthHeaders.mockResolvedValueOnce({ Authorization: 'Bearer refreshed-token' }); + await expect(authHeaders()).resolves.toMatchObject({ + Authorization: 'Bearer refreshed-token', + 'X-LibreChat-Code-Worker-ID': 'worker-user-1', + }); + expect(getCodeApiAuthHeaders).toHaveBeenNthCalledWith(2, mockReq, 'worker-user-1'); expect(mockExecuteWorkspaceTool).toHaveBeenCalledWith({ baseURL: 'https://attached-code.example.com/v1', - authHeaders: { - Authorization: 'Bearer workspace-token', - 'X-CodeAPI-Expected-Profile': 'stateful', - 'X-LibreChat-Code-Worker-ID': 'worker-user-1', - }, + authHeaders: expect.any(Function), + maxQueueWaitMs: 0, request: { protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', + workspaceInstanceId: 'a'.repeat(64), path: 'src/app.ts', startLine: 1, maxLines: 200, @@ -2105,16 +2117,27 @@ describe('Code Process', () => { bridgeWorkerId: 'worker-user-1', req: mockReq, signal: controller.signal, + maxQueueWaitMs: 0, }), ).resolves.toBe(result); + expect(getCodeApiAuthHeaders).not.toHaveBeenCalled(); + const { authHeaders } = mockExecuteWorkspaceTool.mock.calls[0][0]; + await expect(authHeaders()).resolves.toEqual({ + Authorization: 'Bearer workspace-token', + 'X-CodeAPI-Expected-Profile': 'stateful', + 'X-LibreChat-Code-Worker-ID': 'worker-user-1', + }); + getCodeApiAuthHeaders.mockResolvedValueOnce({ Authorization: 'Bearer refreshed-token' }); + await expect(authHeaders()).resolves.toMatchObject({ + Authorization: 'Bearer refreshed-token', + 'X-LibreChat-Code-Worker-ID': 'worker-user-1', + }); + expect(getCodeApiAuthHeaders).toHaveBeenNthCalledWith(2, mockReq, 'worker-user-1'); expect(mockExecuteWorkspaceTool).toHaveBeenCalledWith({ baseURL: 'https://attached-code.example.com/v1', - authHeaders: { - Authorization: 'Bearer workspace-token', - 'X-CodeAPI-Expected-Profile': 'stateful', - 'X-LibreChat-Code-Worker-ID': 'worker-user-1', - }, + authHeaders: expect.any(Function), + maxQueueWaitMs: 0, request: { protocolVersion: 1, operation: 'search_text', @@ -2152,16 +2175,27 @@ describe('Code Process', () => { bridgeWorkerId: 'worker-user-1', req: mockReq, signal: controller.signal, + maxQueueWaitMs: 0, }), ).resolves.toBe(result); + expect(getCodeApiAuthHeaders).not.toHaveBeenCalled(); + const { authHeaders } = mockExecuteWorkspaceTool.mock.calls[0][0]; + await expect(authHeaders()).resolves.toEqual({ + Authorization: 'Bearer workspace-token', + 'X-CodeAPI-Expected-Profile': 'stateful', + 'X-LibreChat-Code-Worker-ID': 'worker-user-1', + }); + getCodeApiAuthHeaders.mockResolvedValueOnce({ Authorization: 'Bearer refreshed-token' }); + await expect(authHeaders()).resolves.toMatchObject({ + Authorization: 'Bearer refreshed-token', + 'X-LibreChat-Code-Worker-ID': 'worker-user-1', + }); + expect(getCodeApiAuthHeaders).toHaveBeenNthCalledWith(2, mockReq, 'worker-user-1'); expect(mockExecuteWorkspaceTool).toHaveBeenCalledWith({ baseURL: 'https://attached-code.example.com/v1', - authHeaders: { - Authorization: 'Bearer workspace-token', - 'X-CodeAPI-Expected-Profile': 'stateful', - 'X-LibreChat-Code-Worker-ID': 'worker-user-1', - }, + authHeaders: expect.any(Function), + maxQueueWaitMs: 0, request: { protocolVersion: 1, operation: 'list_files', @@ -2195,25 +2229,38 @@ describe('Code Process', () => { content: 'ready', overwrite: false, workspace_id: 'primary', + workspace_instance_id: 'b'.repeat(64), codeApiBaseUrl: 'https://attached-code.example.com/v1', executionProfile: 'stateful', bridgeWorkerId: 'worker-user-1', req: mockReq, signal: controller.signal, + maxQueueWaitMs: 0, }), ).resolves.toBe(result); + expect(getCodeApiAuthHeaders).not.toHaveBeenCalled(); + const { authHeaders } = mockExecuteWorkspaceTool.mock.calls[0][0]; + await expect(authHeaders()).resolves.toEqual({ + Authorization: 'Bearer workspace-token', + 'X-CodeAPI-Expected-Profile': 'stateful', + 'X-LibreChat-Code-Worker-ID': 'worker-user-1', + }); + getCodeApiAuthHeaders.mockResolvedValueOnce({ Authorization: 'Bearer refreshed-token' }); + await expect(authHeaders()).resolves.toMatchObject({ + Authorization: 'Bearer refreshed-token', + 'X-LibreChat-Code-Worker-ID': 'worker-user-1', + }); + expect(getCodeApiAuthHeaders).toHaveBeenNthCalledWith(2, mockReq, 'worker-user-1'); expect(mockExecuteWorkspaceTool).toHaveBeenCalledWith({ baseURL: 'https://attached-code.example.com/v1', - authHeaders: { - Authorization: 'Bearer workspace-token', - 'X-CodeAPI-Expected-Profile': 'stateful', - 'X-LibreChat-Code-Worker-ID': 'worker-user-1', - }, + authHeaders: expect.any(Function), + maxQueueWaitMs: 0, request: { protocolVersion: 1, operation: 'write_file', workspaceId: 'primary', + workspaceInstanceId: 'b'.repeat(64), path: 'src/new.ts', content: 'ready', overwrite: false, @@ -2859,6 +2906,7 @@ describe('Code Process', () => { tool_resources: { execute_code: { file_ids: [dbFile.file_id], files: [] } }, agentId: 'agent-id', signal: controller.signal, + maxQueueWaitMs: 0, }), ).rejects.toMatchObject({ name: 'AbortError' }); expect(handleFileUpload).toHaveBeenCalledTimes(1); diff --git a/api/server/services/Files/Local/__tests__/crud-delete.spec.js b/api/server/services/Files/Local/__tests__/crud-delete.spec.js new file mode 100644 index 00000000000..7a7a30340ab --- /dev/null +++ b/api/server/services/Files/Local/__tests__/crud-delete.spec.js @@ -0,0 +1,70 @@ +/** Exactly what `deleteLocalFile` reaches for, so the double does not depend on a built package. */ +jest.mock('@librechat/api', () => ({ + deleteRagFile: jest.fn().mockResolvedValue(undefined), + stripCacheBust: (filepath) => String(filepath).split('?')[0], +})); +jest.mock('@librechat/data-schemas', () => ({ + logger: { warn: jest.fn(), error: jest.fn() }, +})); + +const fs = require('fs'); +const os = require('os'); +const path = require('path'); +const { deleteLocalFile } = require('../crud'); + +/* The resolved promise of a delete is what `processDeleteRequest` reads to decide that a record may + lose its metadata and its agent references, so this adapter may only resolve once the bytes are + actually gone. Storage that was already missing is the one benign case. */ +describe('deleteLocalFile failure reporting', () => { + const userId = 'user-1'; + let tmpBase; + let req; + + beforeEach(() => { + jest.restoreAllMocks(); + tmpBase = fs.mkdtempSync(path.join(os.tmpdir(), 'crud-delete-')); + fs.mkdirSync(path.join(tmpBase, 'uploads', userId), { recursive: true }); + req = { + user: { id: userId }, + config: { + paths: { + publicPath: path.join(tmpBase, 'public'), + uploads: path.join(tmpBase, 'uploads'), + }, + }, + }; + }); + + afterEach(() => { + fs.rmSync(tmpBase, { recursive: true, force: true }); + }); + + const uploadedFile = (filename) => { + const filepath = path.join(tmpBase, 'uploads', userId, filename); + fs.writeFileSync(filepath, 'contents'); + return { file_id: 'file-1', filepath: `/uploads/${userId}/${filename}` }; + }; + + it('removes the file and resolves', async () => { + const file = uploadedFile('knowledge.txt'); + + await expect(deleteLocalFile(req, file)).resolves.toBeUndefined(); + expect(fs.existsSync(path.join(tmpBase, 'uploads', userId, 'knowledge.txt'))).toBe(false); + }); + + it('resolves when the file is already gone', async () => { + await expect( + deleteLocalFile(req, { file_id: 'file-1', filepath: `/uploads/${userId}/missing.txt` }), + ).resolves.toBeUndefined(); + }); + + it('rejects when the bytes survive the delete', async () => { + const file = uploadedFile('locked.txt'); + jest + .spyOn(fs.promises, 'unlink') + .mockRejectedValue(Object.assign(new Error('permission denied'), { code: 'EACCES' })); + + await expect(deleteLocalFile(req, file)).rejects.toThrow('permission denied'); + expect(fs.existsSync(path.join(tmpBase, 'uploads', userId, 'locked.txt'))).toBe(true); + }); +}); diff --git a/api/server/services/Files/Local/crud.js b/api/server/services/Files/Local/crud.js index fe524550986..9555df49c03 100644 --- a/api/server/services/Files/Local/crud.js +++ b/api/server/services/Files/Local/crud.js @@ -208,11 +208,21 @@ const isValidPath = (req, base, subfolder, filepath) => { /** * @param {string} filepath */ +/** + * A file whose bytes are still on disk must not be reported as deleted: callers use the resolved + * promise to decide that a record may lose its metadata and its agent references. Storage that was + * already gone is the one benign case, and `processDeleteRequest` treats it as deleted by design. + */ const unlinkFile = async (filepath) => { try { await fs.promises.unlink(filepath); } catch (error) { + if (error?.code === 'ENOENT') { + logger.warn('Local file was already missing during delete:', error); + return; + } logger.error('Error deleting file:', error); + throw error; } }; diff --git a/api/server/services/Files/process.js b/api/server/services/Files/process.js index 05b02ae6df6..35117ab248e 100644 --- a/api/server/services/Files/process.js +++ b/api/server/services/Files/process.js @@ -265,16 +265,8 @@ const processDeleteRequest = async ({ req, files }) => { await initializeClients(); } - const agentFiles = []; - for (const file of files) { const source = file.source ?? FileSources.local; - if (req.body.agent_id && req.body.tool_resource) { - agentFiles.push({ - tool_resource: req.body.tool_resource, - file_id: file.file_id, - }); - } if (source === FileSources.text) { resolvedFileIds.add(file.file_id); @@ -313,15 +305,6 @@ const processDeleteRequest = async ({ req, files }) => { }); } - if (agentFiles.length > 0) { - promises.push( - db.removeAgentResourceFiles({ - agent_id: req.body.agent_id, - files: agentFiles, - }), - ); - } - await Promise.allSettled(promises); const deletedFileIds = [...resolvedFileIds]; let metadataDeletedFileIds = deletedFileIds; @@ -334,6 +317,9 @@ const processDeleteRequest = async ({ req, files }) => { metadataDeletedFileIds = []; throw error; } + /* The only place a delete removes agent references, and it runs after the metadata delete + succeeded: a file that kept its storage, its chunks or its record keeps its references too, + so the agent it was removed from can be asked again (see issue #12776). */ if (metadataDeletedFileIds.length > 0) { try { await db.removeAgentResourceFilesFromAllAgents({ file_ids: metadataDeletedFileIds }); diff --git a/api/server/services/Files/process.spec.js b/api/server/services/Files/process.spec.js index efbb8a1b008..33ec48b7861 100644 --- a/api/server/services/Files/process.spec.js +++ b/api/server/services/Files/process.spec.js @@ -2610,6 +2610,89 @@ describe('processDeleteRequest', () => { expect(result).toEqual({ deletedFileIds: [], failedFileIds: ['embedded-file'] }); }); + it('keeps a failed agent file attached so the delete can be retried', async () => { + const deleteFile = jest.fn().mockRejectedValue(new Error('rag unavailable')); + getStrategyFunctions.mockReturnValue({ deleteFile }); + const req = { + body: { agent_id: 'agent_1', tool_resource: 'file_search' }, + config: {}, + user: { id: 'user-123', tenantId: 'tenant-a' }, + }; + + const result = await processDeleteRequest({ + req, + files: [ + { + file_id: 'knowledge-file', + filepath: '/uploads/knowledge.txt', + source: FileSources.local, + }, + ], + }); + + expect(result).toEqual({ deletedFileIds: [], failedFileIds: ['knowledge-file'] }); + expect(db.deleteFiles).not.toHaveBeenCalled(); + expect(db.removeAgentResourceFiles).not.toHaveBeenCalled(); + expect(db.removeAgentResourceFilesFromAllAgents).not.toHaveBeenCalled(); + }); + + it('strips agent references only for the files it deleted', async () => { + getStrategyFunctions.mockReturnValue({ + deleteFile: jest + .fn() + .mockImplementation((_req, file) => + file.file_id === 'kept-file' + ? Promise.reject(new Error('rag unavailable')) + : Promise.resolve(undefined), + ), + }); + db.deleteFiles.mockResolvedValue({ deletedCount: 1 }); + const req = { + body: { agent_id: 'agent_1', tool_resource: 'file_search' }, + config: {}, + user: { id: 'user-123', tenantId: 'tenant-a' }, + }; + + const result = await processDeleteRequest({ + req, + files: [ + { file_id: 'gone-file', filepath: '/uploads/gone.txt', source: FileSources.local }, + { file_id: 'kept-file', filepath: '/uploads/kept.txt', source: FileSources.local }, + ], + }); + + expect(result).toEqual({ deletedFileIds: ['gone-file'], failedFileIds: ['kept-file'] }); + expect(db.removeAgentResourceFilesFromAllAgents).toHaveBeenCalledWith({ + file_ids: ['gone-file'], + }); + }); + + it('keeps agent references when the metadata delete fails', async () => { + getStrategyFunctions.mockReturnValue({ deleteFile: jest.fn().mockResolvedValue(undefined) }); + db.deleteFiles.mockRejectedValue(new Error('mongo unavailable')); + const req = { + body: { agent_id: 'agent_1', tool_resource: 'file_search' }, + config: {}, + user: { id: 'user-123', tenantId: 'tenant-a' }, + }; + + await expect( + processDeleteRequest({ + req, + files: [ + { + file_id: 'knowledge-file', + filepath: '/uploads/knowledge.txt', + source: FileSources.local, + }, + ], + }), + ).rejects.toThrow('mongo unavailable'); + + expect(db.removeAgentResourceFiles).not.toHaveBeenCalled(); + expect(db.removeAgentResourceFilesFromAllAgents).not.toHaveBeenCalled(); + }); + it('does not delete vector storage when primary embedded file deletion fails', async () => { const primaryDelete = jest.fn().mockRejectedValue(new Error('permission denied')); const vectorDelete = jest.fn().mockResolvedValue(undefined); diff --git a/api/server/services/ToolService.js b/api/server/services/ToolService.js index 7e4ceca5451..05126536df4 100644 --- a/api/server/services/ToolService.js +++ b/api/server/services/ToolService.js @@ -51,6 +51,7 @@ const { createRepositoryInstructionSource, createRepositoryInstructionLoader, resolveAttachedWorkspaceCommandTimeoutMax, + resolveAttachedWorkspaceQueueWaitMs, createContextProgrammaticBashTool, resolveCodeExecutionContext, resolveCodeExecutionWorkspaceContext, @@ -2295,10 +2296,15 @@ async function loadToolsForExecution({ authHeaders, baseUrl: codeExecutionContext.baseUrl, workspaceId: codeExecutionContext.codeWorkspace.workspaceId, + workspaceInstanceId: codeExecutionContext.codeWorkspace.workspaceInstanceId, environment: codeExecutionContext.codeWorkspace.environment, gitIdentity: agent?.git_identity, maxTimeoutMs: resolveAttachedWorkspaceCommandTimeoutMax( codeExecutionContext.codeEnvironmentConfigSchema, + codeExecutionContext.codeWorkspace?.maxCommandTimeoutMs, + ), + maxQueueWaitMs: resolveAttachedWorkspaceQueueWaitMs( + codeExecutionContext.codeEnvironmentConfigSchema, ), }) : createBashExecutionTool({ diff --git a/api/server/services/__tests__/ToolService.spec.js b/api/server/services/__tests__/ToolService.spec.js index 05641b3a671..558be1d5e0f 100644 --- a/api/server/services/__tests__/ToolService.spec.js +++ b/api/server/services/__tests__/ToolService.spec.js @@ -2898,7 +2898,7 @@ describe('ToolService - Action Capability Gating', () => { environmentType: 'attached', environmentId: 'personal-machine', bridgeWorkerId: 'worker-abc', - codeEnvironmentConfigSchema: { limits: { maxCommandTimeoutMs: 120000 } }, + codeEnvironmentConfigSchema: { limits: { maxCommandTimeoutMs: 120000, maxQueueWaitMs: 0 } }, }); const toolRegistry = new Map([ [AgentConstants.BASH_TOOL, { name: AgentConstants.BASH_TOOL }], @@ -2925,6 +2925,7 @@ describe('ToolService - Action Capability Gating', () => { workspaceId: 'project-a', gitIdentity: { name: 'LibreChat Agent', email: 'agent@example.com' }, maxTimeoutMs: 120000, + maxQueueWaitMs: 0, }); expect(mockResolveCodeExecutionWorkspaceContext).toHaveBeenCalledWith( expect.objectContaining({ requestedSelections: req.body.codeWorkspaces }), diff --git a/api/test/server/middleware/checkBan.test.js b/api/test/server/middleware/checkBan.test.js index 237282d9d0b..39775c389a4 100644 --- a/api/test/server/middleware/checkBan.test.js +++ b/api/test/server/middleware/checkBan.test.js @@ -401,6 +401,55 @@ describe('checkBan middleware', () => { }); }); + describe('non-string cache keys (#16025)', () => { + const objectIdHex = '507f1f77bcf86cd799439011'; + const objectId = { + toString() { + return objectIdHex; + }, + }; + + it('stringifies ObjectId user keys when Redis is off', async () => { + const next = jest.fn(); + const req = createReq({ user: { _id: objectId } }); + + await checkBan(req, createRes(), next); + + expect(next).toHaveBeenCalledWith(); + expect(mockBanCacheGet).toHaveBeenCalledWith(objectIdHex); + expect(mockBanLogsGet).toHaveBeenCalledWith(objectIdHex); + for (const [key] of mockBanCacheGet.mock.calls) { + expect(typeof key).toBe('string'); + } + }); + + it('stringifies numeric user keys when Redis is off', async () => { + await checkBan(createReq({ user: { _id: 12345 } }), createRes(), jest.fn()); + + expect(mockBanCacheGet).toHaveBeenCalledWith('12345'); + expect(mockBanLogsGet).toHaveBeenCalledWith('12345'); + }); + + it('stringifies ObjectId user keys in Redis-prefixed cache lookups', async () => { + process.env.USE_REDIS = 'true'; + + await checkBan(createReq({ user: { _id: objectId } }), createRes(), jest.fn()); + + expect(mockBanCacheGet).toHaveBeenCalledWith(`ban_cache:user:${objectIdHex}`); + expect(mockBanLogsGet).toHaveBeenCalledWith(objectIdHex); + }); + + it('stringifies ObjectId user keys from email lookup', async () => { + findUser.mockResolvedValueOnce({ _id: objectId }); + const req = createReq({ user: null, body: { email: 'oauth@example.com' } }); + + await checkBan(req, createRes(), jest.fn()); + + expect(mockBanCacheGet).toHaveBeenCalledWith(objectIdHex); + expect(mockBanLogsGet).toHaveBeenCalledWith(objectIdHex); + }); + }); + describe('Redis key paths (Finding 2 regression)', () => { beforeEach(() => { process.env.USE_REDIS = 'true'; diff --git a/client/src/Providers/AuthorContext.tsx b/client/src/Providers/AuthorContext.tsx new file mode 100644 index 00000000000..e8cec934465 --- /dev/null +++ b/client/src/Providers/AuthorContext.tsx @@ -0,0 +1,17 @@ +import { createContext, useContext } from 'react'; +import type { ReactNode } from 'react'; + +/** The author a message restates wherever its content resumes after a steer. */ +export type TMessageAuthor = { + icon: ReactNode; + label: string; +}; + +/** + * Carries the message author to the headers inside its content. The author can + * resolve after the message paints, when an agent's name and avatar arrive with the + * agents list, and a context change reaches only the headers that read it rather + * than every part of the message. + */ +export const AuthorContext = createContext(null); +export const useAuthorContext = () => useContext(AuthorContext); diff --git a/client/src/Providers/index.ts b/client/src/Providers/index.ts index 76aad5b4560..3e499921fca 100644 --- a/client/src/Providers/index.ts +++ b/client/src/Providers/index.ts @@ -10,6 +10,7 @@ export * from './EditorContext'; export * from './ChatFormContext'; export * from './BookmarkContext'; export * from './MessageContext'; +export * from './AuthorContext'; export * from './AssistantsContext'; export * from './AgentsContext'; export * from './AssistantsMapContext'; diff --git a/client/src/components/Chat/Messages/Content/Markdown.tsx b/client/src/components/Chat/Messages/Content/Markdown.tsx index 24b3cc26249..222c3d0a0d3 100644 --- a/client/src/components/Chat/Messages/Content/Markdown.tsx +++ b/client/src/components/Chat/Messages/Content/Markdown.tsx @@ -19,7 +19,8 @@ const Markdown = memo(function Markdown({ content = '', isLatestMessage }: TCont const LaTeXParsing = useRecoilValue(store.LaTeXParsing); const isInitializing = content === ''; - const animate = smoothStreaming && isLatestMessage && isSubmitting; + const streaming = isLatestMessage && isSubmitting; + const animate = smoothStreaming && streaming; // Hydration signal for the fade: substantial content already present at the // render where `animate` flips on means resumed/switched-to/follow-up @@ -52,6 +53,7 @@ const Markdown = memo(function Markdown({ content = '', isLatestMessage }: TCont ( ); +/** The per-block renderer a streaming message uses, without the fade. */ const NewMarkdown = ({ content }: { content: string }) => ( - + +); + +/** A finished message as a conversation opens with it. */ +const SettledMarkdown = ({ content }: { content: string }) => ( + ); const streamingContext = { @@ -186,4 +200,53 @@ describe('Markdown streaming benchmark (OLD whole-message vs NEW per-block)', () // Memoization should cut total code-block renders by a wide margin. expect(newRenders).toBeLessThan(oldRenders * 0.5); }); + + it('reports mount cost for a message that was already finished', () => { + const iterations = 5; + const rows = [12, 40].map((sections) => { + const prefixes = [buildMessage(sections)]; + measure(OldMarkdown, prefixes); + measure(NewMarkdown, prefixes); + measure(SettledMarkdown, prefixes); + + const old: Array<{ totalMs: number; codeBlockRenders: number }> = []; + const perBlock: Array<{ totalMs: number; codeBlockRenders: number }> = []; + const settled: Array<{ totalMs: number; codeBlockRenders: number }> = []; + for (let i = 0; i < iterations; i += 1) { + old.push(measure(OldMarkdown, prefixes)); + perBlock.push(measure(NewMarkdown, prefixes)); + settled.push(measure(SettledMarkdown, prefixes)); + } + const minMs = (rs: Array<{ totalMs: number }>) => Math.min(...rs.map((r) => r.totalMs)); + return { + chars: prefixes[0].length, + oldMs: minMs(old), + perBlockMs: minMs(perBlock), + settledMs: minMs(settled), + oldRenders: old[0].codeBlockRenders, + settledRenders: settled[0].codeBlockRenders, + }; + }); + + console.log( + [ + '', + '================ Markdown finished-message mount benchmark ================', + `min of ${iterations} mounts, summed Profiler actualDuration; jsdom`, + ...rows.map( + (r) => + ` ${r.chars} chars: OLD ${r.oldMs.toFixed(1)} ms | PER-BLOCK ${r.perBlockMs.toFixed(1)} ms | ` + + `SETTLED ${r.settledMs.toFixed(1)} ms (${(r.perBlockMs / r.settledMs).toFixed(2)}x faster than per-block)`, + ), + '==========================================================================', + '', + ].join('\n'), + ); + + // A finished message mounts through the whole-message pipeline: every code + // block renders exactly once, as it did before per-block rendering existed. + for (const r of rows) { + expect(r.settledRenders).toBe(r.oldRenders); + } + }); }); diff --git a/client/src/components/Chat/Messages/Content/MarkdownBlocks.tsx b/client/src/components/Chat/Messages/Content/MarkdownBlocks.tsx index d5b76f3ae62..256bc9cc753 100644 --- a/client/src/components/Chat/Messages/Content/MarkdownBlocks.tsx +++ b/client/src/components/Chat/Messages/Content/MarkdownBlocks.tsx @@ -1,4 +1,4 @@ -import React, { memo, useMemo, useLayoutEffect } from 'react'; +import React, { memo, useMemo, useState, useLayoutEffect } from 'react'; import ReactMarkdown from 'react-markdown'; import type { PluggableList } from 'unified'; import type { ElementType } from 'react'; @@ -86,46 +86,103 @@ MarkdownBlock.displayName = 'MarkdownBlock'; type MarkdownBlocksProps = SharedProps & { content: string; + /** Whether this message is the one generating right now. */ + streaming: boolean; +}; + +type BlockEntry = { + /** + * Code and artifact blocks capture their index in a ref when they mount, so a + * block has to remount whenever an index it holds could shift, and must not + * remount otherwise. + */ + key: string; + raw: string; + codeBaseIndex: number; + artifactBaseIndex: number; + mermaidBaseIndex: number; +}; + +/** + * Each top-level block, seeded with the code, artifact and Mermaid indices of the + * blocks before it. The key carries those bases, so an in-place edit that inserts + * a block before existing code or artifacts remounts the blocks it shifted, while + * append-only streaming keeps completed blocks mounted. + */ +const toBlockEntries = (content: string): BlockEntry[] => { + let codeBaseIndex = 0; + let artifactBaseIndex = 0; + let mermaidBaseIndex = 0; + return splitMarkdownIntoBlocks(content).map((block, index) => { + const entry = { + key: `${index}-${codeBaseIndex}-${artifactBaseIndex}-${mermaidBaseIndex}`, + raw: block.raw, + codeBaseIndex, + artifactBaseIndex, + mermaidBaseIndex, + }; + codeBaseIndex += block.codeBlockCount; + artifactBaseIndex += block.artifactCount; + mermaidBaseIndex += block.mermaidCount; + return entry; + }); }; /** - * Splits a message into top-level blocks and renders each independently so - * that, during streaming, only the last block re-parses while earlier blocks - * (tables, code, etc.) stay memoized. Each block's executable code and artifact - * indices are preserved in document order via per-block providers seeded with - * prefix-summed base indices. + * The whole message as one block, which numbers its code and artifacts from zero + * itself. It only ever renders the message exactly as it mounted, so no index it + * assigns can go stale and its key never has to change. + */ +const toWholeMessage = (content: string): BlockEntry[] => + content + ? [{ key: 'whole', raw: content, codeBaseIndex: 0, artifactBaseIndex: 0, mermaidBaseIndex: 0 }] + : []; + +/** + * Renders a message's markdown. + * + * While the message streams, each top-level block renders and memoizes on its + * own, so only the last, still-growing block re-parses on each token. Each + * block's executable code and artifact indices stay in document order through + * per-block providers seeded with prefix-summed base indices. + * + * A message that mounts finished renders as one pipeline instead, for as long as + * it stays exactly as it mounted. Splitting it would buy memoization nothing + * uses, at the price of a whole extra parse to find block boundaries and one + * pipeline per block, which is most of the cost of opening a long conversation. + * + * Once the message generates or its content changes, it moves to the split for + * good. Code and artifact blocks capture their index at mount, and per block only + * the blocks a change touches re-render, so finishing an answer never remounts its + * blocks and neither do repeated edits such as artifact saves. The move itself + * remounts the message once. */ const MarkdownBlocks = memo(function MarkdownBlocks({ content, + streaming, remarkPlugins, rehypePlugins, components, animate, hydrated, }: MarkdownBlocksProps) { - const blocks = useMemo(() => { - let codeBaseIndex = 0; - let artifactBaseIndex = 0; - let mermaidBaseIndex = 0; - return splitMarkdownIntoBlocks(content).map((block) => { - const entry = { raw: block.raw, codeBaseIndex, artifactBaseIndex, mermaidBaseIndex }; - codeBaseIndex += block.codeBlockCount; - artifactBaseIndex += block.artifactCount; - mermaidBaseIndex += block.mermaidCount; - return entry; - }); - }, [content]); + const [mountedContent] = useState(content); + const [hasChanged, setHasChanged] = useState(streaming); + const changing = streaming || content !== mountedContent; + if (changing && !hasChanged) { + setHasChanged(true); + } + const perBlock = hasChanged || changing; + const blocks = useMemo( + () => (perBlock ? toBlockEntries(content) : toWholeMessage(content)), + [content, perBlock], + ); return ( <> - {blocks.map((block, index) => ( - // Key includes the base indices so that an in-place edit which inserts a - // block before existing code/artifact blocks (shifting their base) forces - // a remount, refreshing the index each code/artifact block captures in a - // ref. During append-only streaming these stay constant, so completed - // blocks keep a stable key and are not remounted. + {blocks.map((block) => ( ; +}); + +export default ResumeAuthorHeader; diff --git a/client/src/components/Chat/Messages/Content/Parts/index.ts b/client/src/components/Chat/Messages/Content/Parts/index.ts index 229db2fe36a..56ef2380ee0 100644 --- a/client/src/components/Chat/Messages/Content/Parts/index.ts +++ b/client/src/components/Chat/Messages/Content/Parts/index.ts @@ -18,3 +18,4 @@ export { default as PtcToolTrace } from './PtcToolTrace'; export { default as SubagentCall } from './SubagentCall'; export { default as SteerPart } from './SteerPart'; export { default as AuthorHeader } from './AuthorHeader'; +export { default as ResumeAuthorHeader } from './ResumeAuthorHeader'; diff --git a/client/src/components/Chat/Messages/Content/ToolCall.tsx b/client/src/components/Chat/Messages/Content/ToolCall.tsx index ee14b6971e1..6bd9e522300 100644 --- a/client/src/components/Chat/Messages/Content/ToolCall.tsx +++ b/client/src/components/Chat/Messages/Content/ToolCall.tsx @@ -10,10 +10,10 @@ import { } from 'librechat-data-provider'; import type { TAttachment, PartMetadata } from 'librechat-data-provider'; import { useLocalize, useProgress, useExpandCollapse, useLazyCollapseBody } from '~/hooks'; +import { cn, getToolDisplayLabel, logger, openInNewTab } from '~/utils'; import { ToolIcon, getToolIconType, isError } from './ToolOutput'; import { useMCPIconMap, useMCPServerNames } from '~/hooks/MCP'; import { resolveToolCallPhase } from '~/utils/toolCallPhase'; -import { cn, getToolDisplayLabel, logger } from '~/utils'; import { toolPanelSpacingClassName } from './disclosure'; import { useToolCallIntent } from './Parts/intent'; import { AttachmentGroup } from './Parts'; @@ -54,6 +54,7 @@ export default function ToolCall({ }) { const localize = useLocalize(); const [oauthError, setOAuthError] = useState(null); + const [oauthBinding, setOAuthBinding] = useState<'pending' | 'bound' | 'failed'>('pending'); const autoExpand = useRecoilValue(store.autoExpandTools); const hasOutput = (output?.length ?? 0) > 0; const [showInfo, setShowInfo] = useState(() => autoExpand && hasOutput); @@ -139,24 +140,23 @@ export default function ToolCall({ return match?.[1] || ''; }, [parsedAuthUrl, isMCPToolCall]); - const handleOAuthClick = useCallback(async () => { - if (!auth) { - return; + /** + * Sets the CSRF cookie the OAuth callback checks when the provider redirects back, or returns + * null when this prompt has nothing to bind. + */ + const bindOAuth = useCallback((): Promise | null => { + const bindsMCP = isMCPToolCall && mcpServerName.length > 0; + if (!bindsMCP && !actionId) { + return null; } - setOAuthError(null); - try { - if (isMCPToolCall && mcpServerName) { + return (async () => { + if (bindsMCP) { await dataService.bindMCPOAuth(mcpServerName); - } else if (actionId) { + } else { await dataService.bindActionOAuth(actionId); } - } catch (e) { - logger.error('Failed to bind OAuth CSRF cookie', e); - setOAuthError(localize('com_ui_oauth_error_generic')); - return; - } - window.open(auth, '_blank', 'noopener,noreferrer'); - }, [auth, isMCPToolCall, mcpServerName, actionId, localize]); + })(); + }, [isMCPToolCall, mcpServerName, actionId]); const hasError = (typeof output === 'string' && isError(output)) || runStepStatus === 'failed'; /** @@ -216,6 +216,72 @@ export default function ToolCall({ isSubmitting, hasError, }); + const showOAuth = Boolean(auth) && phase === 'running'; + + /** + * Binds when the sign-in prompt appears instead of on tap, so the tap opens the provider + * synchronously: an iOS home-screen app drops a tab opened after an awaited request. The button + * stays disabled until the bind lands, so the provider cannot redirect back before its cookie. + */ + useEffect(() => { + if (!showOAuth) { + return; + } + const binding = bindOAuth(); + if (binding == null) { + setOAuthBinding('bound'); + return; + } + let active = true; + setOAuthBinding('pending'); + binding.then( + () => { + if (active) { + setOAuthBinding('bound'); + } + }, + (error: unknown) => { + logger.error('Failed to bind OAuth CSRF cookie', error); + if (active) { + setOAuthBinding('failed'); + } + }, + ); + return () => { + active = false; + }; + }, [showOAuth, auth, bindOAuth]); + + const handleOAuthClick = useCallback(() => { + if (!auth || oauthBinding === 'pending') { + return; + } + if (oauthBinding === 'bound') { + setOAuthError(null); + openInNewTab(auth); + /** + * Live prompts share one CSRF cookie per callback path, so the last prompt to bind owns it. + * The tapped prompt claims it again after opening; the provider cannot redirect back before + * the user signs in, and the session cookie from the earlier bind covers a faster callback. + */ + bindOAuth()?.catch((error: unknown) => { + logger.error('Failed to bind OAuth CSRF cookie', error); + }); + return; + } + setOAuthError(localize('com_ui_oauth_error_generic')); + setOAuthBinding('pending'); + (bindOAuth() ?? Promise.resolve()).then( + () => { + setOAuthBinding('bound'); + setOAuthError(null); + }, + (error: unknown) => { + logger.error('Failed to bind OAuth CSRF cookie', error); + setOAuthBinding('failed'); + }, + ); + }, [auth, oauthBinding, bindOAuth, localize]); const handleToggleInfo = useCallback(() => { mountBody(); @@ -329,13 +395,15 @@ export default function ToolCall({ )} - {auth != null && auth && phase === 'running' && ( + {showOAuth && (
+ )} + {readable && !raw && ( + + )} + {(!readable || raw) && ( + + )}
); } +function RawContent({ + input, + output, + metadata, +}: Pick) { + return ( +
+ + + +
+ ); +} + +const COPIED_MS = 2000; +const CALL_CONTENT_LENGTH = 4000; + +function toContent(value?: string): TTraceContent | undefined { + if (value == null) { + return undefined; + } + const truncated = value.length > CALL_CONTENT_LENGTH; + return { value: truncated ? value.slice(0, CALL_CONTENT_LENGTH) : value, truncated }; +} + +/** The saved agent a record ran: who it is, its id to copy, and a chat with it one click away. */ +function AgentCard({ agent, agentId }: { agent: RecordPresentation['agent']; agentId: string }) { + const localize = useLocalize(); + const [copied, setCopied] = useState(false); + const copyId = () => { + copy(agentId); + setCopied(true); + setTimeout(() => setCopied(false), COPIED_MS); + }; + return ( +
+ {agent?.description != null && agent.description !== '' && ( +

{agent.description}

+ )} + {agent == null && ( +

{localize('com_ui_trace_agent_unavailable')}

+ )} +
+ + {agentId} + + +
+ {agent != null && ( + +
+ ); +} + +/** A tool round's calls as the chat's own tool cards hold them: no trace read, no content gate. */ +function ToolCalls({ + calls, + mcpIconMap, +}: { + calls: ToolCallView[]; + mcpIconMap: Map; +}) { + const localize = useLocalize(); + return ( +
+ {calls.map((call, index) => ( +
+

+ + {call.title} + {call.caption != null && ( + + {call.caption} + + )} +

+ + +
+ ))} +

{localize('com_ui_trace_from_conversation')}

+
+ ); +} + /** Details for the selected record; input and output load only when the deployment allows them. */ function Inspector({ node, + presentation, + mcpIconMap, + toolFor, turnStart, sourceId, conversationId, @@ -122,6 +266,9 @@ function Inspector({ onClose, }: { node: TraceNode; + presentation: RecordPresentation; + mcpIconMap: Map; + toolFor: ToolFor; turnStart: number; /** The page source that listed this record. */ sourceId?: string; @@ -135,7 +282,7 @@ function Inspector({ const format = useTraceFormat(); const headingId = useId(); const { record } = node; - const appearance = KIND_APPEARANCE[record.kind]; + const appearance = appearanceOf(record); const Icon = appearance.icon; const { usage } = record; @@ -179,6 +326,25 @@ function Inspector({ }); } + const preview = presentation.calls == null && presentation.preview != null && ( +

+ {presentation.preview} +

+ ); + /** A model call is read as a conversation, which is what the panel is opened for, so it leads. */ + const leadsWithContent = showContent && record.kind === 'generation'; + const content = ( + + ); + return (