From 84fa6aa8202a373e7f6b744b961ea7e20038885b Mon Sep 17 00:00:00 2001 From: Danny Avila Date: Sun, 5 Jul 2026 11:29:30 -0400 Subject: [PATCH] =?UTF-8?q?=F0=9F=A7=B9=20feat:=20Eager=20HITL=20Checkpoin?= =?UTF-8?q?t=20Cleanup=20(Expiry=20+=20Deletion)=20&=20Full-Wiring=20E2E?= =?UTF-8?q?=20(#14123)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * feat: eager HITL checkpoint cleanup on expiry + deletion, full-wiring e2e Follow-up to the lazy checkpointer (#14024): two paths still left a paused run's durable checkpoint to the 24h Mongo TTL, and no test exercised the whole HITL seam with real components. 1. Approval expiry: GenerationJobManager.setApprovalExpiredHandler(fn) — a non-destructive host hook fired after expireApproval's CAS succeeds (periodic sweeper AND stale-submit path), safe on startups that run constructor defaults. Both startups (index.js configureGenerationStreams, experimental.js) register a handler that prunes the checkpoint, resolving config lazily per expiry (streamId === conversationId === thread_id). 2. Conversation deletion: deleteConvos now returns the deleted conversationIds; the three deletion paths (DELETE /convos, DELETE /convos/all, account deletion) prune them via the new bulk deleteAgentCheckpoints (one $in deleteMany per collection). The delete routes gain configMiddleware for the checkpointer config. 3. Full-wiring e2e (hitlCheckpoint.e2e.spec.js): real SDK Run driven by FakeChatModel calling a gated tool, real PreToolUse/humanInTheLoop wiring, real LazyMongoSaver over mongodb-memory-server, real GenerationJobManager, real /resume controller via supertest. Asserts: clean turn persists nothing; error turn persists nothing; pause -> HTTP approve -> gated tool executes exactly once -> finalize prunes the checkpoint; expiry prunes the abandoned pause eagerly. Tests: 3 handler unit tests (pendingAction.spec), 2 bulk-prune integration tests (checkpointer.integration.spec), convos route + deleteUser specs updated, 4 e2e scenarios. 207 tests green across changed areas. * fix: tenant-scoped expiry prune, store-won expiry relay, resilient deleteConvos ids Codex round 1 on #14123 — all three valid: 1. The approval-expired handler now receives the expired JOB so both startups resolve config in the paused job's tenant/user scope (getAppConfig({userId, tenantId})) — a tenant checkpointer override no longer sends the prune to the base config's collections. expireApproval fetches the job best-effort. 2. Multi-replica: when RedisJobStore.cleanupRequiresActionIndex wins the expiry CAS on another replica, this replica's sweeper relay branch now runs the approval-expired cleanup too (prune is idempotent) — store-driven expiry no longer bypasses the hook. 3. deleteConvos: post-delete cleanup (deleteMessages, project stats refresh) is now best-effort — the conversations are already gone, so throwing hid the deletion and dropped the conversationIds the checkpoint prune needs, unrecoverable on retry. Updated the existing tag-decrement-on-failure test to the new contract (ids still returned). Tests: handler-receives-job, store-won relay path, ids-survive-cleanup-failure. 134 tests green across changed suites. * fix: relay cleanup independent of cached errorEvent; enter tenant ALS context Codex round 2 on #14123: 1. The sweeper's relay branch gated BOTH the terminal-error emit and the new checkpoint cleanup on !runtime.errorEvent — but a reconnect seeds errorEvent from the aborted job (runtime-state creation), which then suppressed the cleanup entirely. The emit stays gated; the idempotent cleanup now runs independent of the cached error, once per runtime lifetime (approvalCleanupRan flag — the aborted job is swept repeatedly). 2. Passing userId/tenantId to getAppConfig only keys the config cache; the Config query is ALS-scoped by the tenant-isolation plugin. Both startup handlers now ENTER the paused job's tenant context via tenantStorage.run before resolving config + pruning, so a tenant checkpointer override is honored in strict and non-strict modes. Tests: relay-cleanup-with-cached-error (reconnect simulation), repeated sweeps run the cleanup once. 27+4 green. * fix: dedup expiry cleanup across winner and relay paths Codex round 3 (P3): expireApproval ran the handler without marking the runtime's approvalCleanupRan flag, so the next sweep's relay branch (the aborted job outlives expiry for the completed-job TTL) ran the cleanup a second time. The dedup now lives inside runApprovalExpiredHandler — the single choke point both paths call — set-before-run, once per runtime lifetime. Test: local expiry followed by a sweep fires the handler once. --- api/server/controllers/UserController.js | 16 +- .../__tests__/hitlCheckpoint.e2e.spec.js | 347 ++++++++++++++++++ api/server/experimental.js | 20 +- api/server/index.js | 16 +- .../__test-utils__/convos-route-mocks.js | 1 + api/server/routes/__tests__/convos.spec.js | 15 + api/server/routes/convos.js | 16 +- .../agents/checkpointer.integration.spec.ts | 36 ++ packages/api/src/agents/checkpointer.ts | 44 +++ .../api/src/stream/GenerationJobManager.ts | 87 ++++- .../stream/__tests__/pendingAction.spec.ts | 105 ++++++ .../src/methods/conversation.spec.ts | 12 +- .../data-schemas/src/methods/conversation.ts | 40 +- 13 files changed, 726 insertions(+), 29 deletions(-) create mode 100644 api/server/controllers/agents/__tests__/hitlCheckpoint.e2e.spec.js diff --git a/api/server/controllers/UserController.js b/api/server/controllers/UserController.js index fe3a206458..5a64cd33a7 100644 --- a/api/server/controllers/UserController.js +++ b/api/server/controllers/UserController.js @@ -7,6 +7,7 @@ const { MCPTokenStorage, normalizeHttpError, extractWebSearchEnvVars, + deleteAgentCheckpoints, deleteAllSharedLinksWithCleanup, } = require('@librechat/api'); const { @@ -360,7 +361,20 @@ const deleteUserController = async (req, res) => { await db.deleteBalances({ user: user._id }); await db.deletePresets(user.id); try { - await db.deleteConvos(user.id); + const convoDeletion = await db.deleteConvos(user.id); + // HITL: prune the deleted conversations' durable checkpoints — a paused run's + // checkpoint would otherwise persist until the Mongo TTL. Never throws. + const appConfig = + req.config ?? + (await getAppConfig({ + role: req.user?.role, + userId: req.user?.id, + tenantId: req.user?.tenantId, + })); + await deleteAgentCheckpoints( + convoDeletion?.conversationIds, + appConfig?.endpoints?.agents?.checkpointer, + ); } catch (error) { logger.error('[deleteUserController] Error deleting user convos, likely no convos', error); } diff --git a/api/server/controllers/agents/__tests__/hitlCheckpoint.e2e.spec.js b/api/server/controllers/agents/__tests__/hitlCheckpoint.e2e.spec.js new file mode 100644 index 0000000000..29e594bd7a --- /dev/null +++ b/api/server/controllers/agents/__tests__/hitlCheckpoint.e2e.spec.js @@ -0,0 +1,347 @@ +/** + * Full-wiring HITL checkpoint lifecycle e2e. + * + * Every HITL-specific component here is REAL: the `@librechat/agents` Run (driven by the + * SDK's FakeChatModel scripted to call a gated tool), the PreToolUse approval hook + + * `humanInTheLoop` wiring, the LazyMongoSaver over mongodb-memory-server, the + * GenerationJobManager (in-memory services), and the `/agents/chat/resume` controller via + * supertest. Only LibreChat's persistence adapters (`~/models`), request cleanup, and the + * concurrency gate are mocked. This is the cross-layer seam none of the unit suites cover: + * pause → durable checkpoint → HTTP approval → rebuilt-run resume → finalize prune. + */ +const express = require('express'); +const request = require('supertest'); +const mongoose = require('mongoose'); +const { MongoMemoryServer } = require('mongodb-memory-server'); +const { z } = require('zod'); +const { tool } = require('@langchain/core/tools'); +const { HumanMessage } = require('@langchain/core/messages'); +const { Run, Providers, FakeChatModel } = require('@librechat/agents'); + +const mockLogger = { debug: jest.fn(), info: jest.fn(), warn: jest.fn(), error: jest.fn() }; + +jest.mock('@librechat/data-schemas', () => ({ + ...jest.requireActual('@librechat/data-schemas'), + logger: mockLogger, +})); + +jest.mock('@librechat/api', () => ({ + ...jest.requireActual('@librechat/api'), + checkAndIncrementPendingRequest: jest.fn(async () => ({ allowed: true })), + decrementPendingRequest: jest.fn(async () => {}), +})); + +jest.mock('~/models', () => ({ + saveMessage: jest.fn(async (req, message) => message), + getConvo: jest.fn(async () => null), + getMessages: jest.fn(async () => []), +})); + +jest.mock('~/server/cleanup', () => ({ + disposeClient: jest.fn(), +})); + +jest.mock('~/server/services/MCPRequestContext', () => ({ + getMCPRequestContext: jest.fn(() => null), + cleanupMCPRequestContextForReq: jest.fn(), +})); + +// Import after mocks — these are the REAL implementations. +const { + GenerationJobManager, + createStreamServices, + buildPendingAction, + getAgentCheckpointer, + deleteAgentCheckpoint, + buildHITLRunWiring, + resolveToolApprovalPolicy, + __resetCheckpointerForTests, +} = require('@librechat/api'); +const ResumeAgentController = require('~/server/controllers/agents/resume'); + +const USER_ID = 'hitl-e2e-user'; +const MONGO_CFG = { type: 'mongo', ttl: 3600 }; +const GATED_TOOL = 'guarded_echo'; + +/** Side-effect counter: proves the gated tool runs exactly once across pause+resume. */ +let toolExecutions = 0; +const guardedTool = tool(async ({ text }) => `echo:${text}`, { + name: GATED_TOOL, + description: 'Echoes text back, but requires human approval first.', + schema: z.object({ text: z.string() }), +}); +guardedTool.func = async ({ text }) => { + toolExecutions += 1; + return `echo:${text}`; +}; + +/** Build a REAL run with the HITL wiring + durable checkpointer attached (mirrors createRun). */ +async function buildHitlRun({ saver, conversationId, responses, toolCalls, runId }) { + const hitl = buildHITLRunWiring( + resolveToolApprovalPolicy({ endpoint: { enabled: true, ask: [GATED_TOOL] } }), + { userId: USER_ID, conversationId, appConfig: {} }, + ); + const run = await Run.create({ + runId, + graphConfig: { + type: 'standard', + llmConfig: { + provider: Providers.OPENAI, + model: 'gpt-4o-mini', + streaming: true, + streamUsage: false, + }, + instructions: 'You are a helpful assistant.', + tools: [guardedTool], + compileOptions: { checkpointer: saver }, + }, + returnContent: true, + customHandlers: {}, + tokenCounter: (text) => String(text ?? '').length, + indexTokenCountMap: {}, + ...(hitl && { humanInTheLoop: hitl.humanInTheLoop, hooks: hitl.hooks }), + }); + run.Graph.overrideModel = new FakeChatModel({ responses, toolCalls }); + return run; +} + +const runConfig = (conversationId) => ({ + runName: 'AgentRun', + configurable: { thread_id: conversationId, user_id: USER_ID }, + streamMode: 'values', + version: 'v2', +}); + +/** Poll until `predicate` returns true (the resume continuation is fire-and-forget). */ +async function waitFor(predicate, { timeoutMs = 10_000, intervalMs = 50 } = {}) { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (await predicate()) { + return; + } + await new Promise((resolve) => setTimeout(resolve, intervalMs)); + } + throw new Error('waitFor: condition not met within timeout'); +} + +async function checkpointCounts(conversationId) { + const db = mongoose.connection.db; + return { + checkpoints: await db + .collection('agent_checkpoints') + .countDocuments({ thread_id: conversationId }), + writes: await db + .collection('agent_checkpoint_writes') + .countDocuments({ thread_id: conversationId }), + }; +} + +let mongoServer; +let saver; + +beforeAll(async () => { + mongoServer = await MongoMemoryServer.create(); + await mongoose.connect(mongoServer.getUri()); + __resetCheckpointerForTests(); + saver = await getAgentCheckpointer(MONGO_CFG); + + GenerationJobManager.configure({ ...createStreamServices(), cleanupOnComplete: false }); + GenerationJobManager.initialize(); + // Mirrors api/server/index.js: expiry prunes the paused run's durable checkpoint. + GenerationJobManager.setApprovalExpiredHandler(async (conversationId) => { + await deleteAgentCheckpoint(conversationId, MONGO_CFG); + }); +}, 60000); + +afterAll(async () => { + GenerationJobManager.setApprovalExpiredHandler(null); + await GenerationJobManager.destroy(); + await mongoose.disconnect(); + await mongoServer.stop(); +}); + +beforeEach(() => { + toolExecutions = 0; + jest.clearAllMocks(); +}); + +describe('HITL checkpoint lifecycle (full wiring)', () => { + jest.setTimeout(30000); + + test('a clean turn (no tool gating triggered) persists NOTHING durable', async () => { + const conversationId = `e2e-clean-${Date.now()}`; + const run = await buildHitlRun({ + saver, + conversationId, + responses: ['Hello there!'], + runId: 'resp-clean', + }); + await run.processStream({ messages: [new HumanMessage('hi')] }, runConfig(conversationId)); + + expect(run.getInterrupt?.()).toBeFalsy(); + expect(await checkpointCounts(conversationId)).toEqual({ checkpoints: 0, writes: 0 }); + }); + + test('a turn that ERRORS before pausing persists NOTHING durable', async () => { + const conversationId = `e2e-error-${Date.now()}`; + const run = await buildHitlRun({ + saver, + conversationId, + responses: ['unused'], + runId: 'resp-error', + }); + class BoomModel extends FakeChatModel { + // eslint-disable-next-line require-yield + async *_streamResponseChunks() { + throw new Error('model boom'); + } + } + run.Graph.overrideModel = new BoomModel({ responses: ['unused'] }); + + await expect( + run.processStream({ messages: [new HumanMessage('hi')] }, runConfig(conversationId)), + ).rejects.toThrow('model boom'); + + expect(await checkpointCounts(conversationId)).toEqual({ checkpoints: 0, writes: 0 }); + }); + + test('pause → approve over the REAL /resume controller → tool runs once → checkpoint pruned', async () => { + const conversationId = `e2e-resume-${Date.now()}`; + const responseMessageId = 'resp-pause-1'; + + // --- Turn 1: the model calls the gated tool → PreToolUse 'ask' → interrupt. --- + const run = await buildHitlRun({ + saver, + conversationId, + responses: ['Let me run that.'], + toolCalls: [{ name: GATED_TOOL, args: { text: 'hi' }, id: 'tc_1', type: 'tool_call' }], + runId: responseMessageId, + }); + await run.processStream( + { messages: [new HumanMessage('run the guarded tool')] }, + runConfig(conversationId), + ); + + const interrupt = run.getInterrupt(); + expect(interrupt?.payload?.type).toBe('tool_approval'); + expect(toolExecutions).toBe(0); // gated — must NOT have run pre-approval + const paused = await checkpointCounts(conversationId); + expect(paused.checkpoints).toBeGreaterThan(0); // the interrupt checkpoint is durable + + // --- Pause bookkeeping (mirrors AgentClient.handleRunInterrupt). --- + const job = await GenerationJobManager.createJob(conversationId, USER_ID, conversationId); + await GenerationJobManager.updateMetadata(conversationId, { + endpoint: 'agents', + agent_id: 'agent-e2e', + responseMessageId, + }); + const pendingAction = buildPendingAction(interrupt.payload, { + streamId: conversationId, + conversationId, + runId: responseMessageId, + responseMessageId, + ttlMs: 60_000, + }); + expect(await GenerationJobManager.approvals.pause(conversationId, pendingAction)).toBe(true); + + // --- Turn 2: approve through the REAL controller; the thin client rebuilds a REAL run. --- + const thinClient = { + contentParts: [], + artifactPromises: [], + conversationId, + responseMessageId, + pendingApproval: null, + async resumeCompletion({ resumeValue, abortController }) { + const resumed = await buildHitlRun({ + saver, + conversationId, + responses: ['Done after approval.'], + runId: responseMessageId, + }); + await resumed.resume(resumeValue, { + ...runConfig(conversationId), + signal: (abortController ?? new AbortController()).signal, + }); + const reInterrupt = resumed.getInterrupt?.(); + if (reInterrupt?.payload) { + this.pendingApproval = reInterrupt.payload; + } + this.contentParts.push({ type: 'text', text: 'Done after approval.' }); + return resumed; + }, + }; + const initializeClient = jest.fn(async () => ({ client: thinClient })); + const addTitle = jest.fn(); + + const app = express(); + app.use(express.json()); + app.use((req, _res, next) => { + req.user = { id: USER_ID }; + req.config = { endpoints: { agents: { checkpointer: MONGO_CFG } }, interfaceConfig: {} }; + next(); + }); + app.post('/api/agents/chat/resume', (req, res, next) => + ResumeAgentController(req, res, next, initializeClient, addTitle), + ); + + const response = await request(app) + .post('/api/agents/chat/resume') + .send({ + conversationId, + actionId: pendingAction.actionId, + agent_id: 'agent-e2e', + endpoint: 'agents', + decisions: [{ tool_call_id: 'tc_1', decision: 'approve' }], + }); + + // The controller ACKs immediately ({ status: 'resuming' }) and drives the resumed run + // asynchronously — wait for the terminal side effects before asserting. + expect(response.status).toBe(200); + expect(response.body.status).toBe('resuming'); + await waitFor(async () => { + const liveJob = await GenerationJobManager.getJob(conversationId); + return liveJob?.status !== 'requires_action' && liveJob?.status !== 'running'; + }); + + expect(initializeClient).toHaveBeenCalledTimes(1); + expect(toolExecutions).toBe(1); // approved tool ran exactly ONCE across pause+resume + + // Terminal state: the checkpoint was pruned by the REAL finalize path. + await waitFor(async () => (await checkpointCounts(conversationId)).checkpoints === 0); + expect(await checkpointCounts(conversationId)).toEqual({ checkpoints: 0, writes: 0 }); + + expect(job).toBeDefined(); + }); + + test('an abandoned pause is pruned eagerly on approval EXPIRY (not left to the TTL)', async () => { + const conversationId = `e2e-expiry-${Date.now()}`; + const run = await buildHitlRun({ + saver, + conversationId, + responses: ['Let me run that.'], + toolCalls: [{ name: GATED_TOOL, args: { text: 'x' }, id: 'tc_exp', type: 'tool_call' }], + runId: 'resp-expire', + }); + await run.processStream({ messages: [new HumanMessage('run it')] }, runConfig(conversationId)); + const interrupt = run.getInterrupt(); + expect((await checkpointCounts(conversationId)).checkpoints).toBeGreaterThan(0); + + await GenerationJobManager.createJob(conversationId, USER_ID, conversationId); + const pendingAction = buildPendingAction(interrupt.payload, { + streamId: conversationId, + conversationId, + runId: 'resp-expire', + responseMessageId: 'resp-expire', + ttlMs: 60_000, + }); + await GenerationJobManager.approvals.pause(conversationId, pendingAction); + + // The sweeper/stale-submit path: expiry fires the registered checkpoint prune. + expect(await GenerationJobManager.expireApproval(conversationId, pendingAction.actionId)).toBe( + true, + ); + + expect(await GenerationJobManager.getJobStatus(conversationId)).toBe('aborted'); + expect(await checkpointCounts(conversationId)).toEqual({ checkpoints: 0, writes: 0 }); + }); +}); diff --git a/api/server/experimental.js b/api/server/experimental.js index dadb87dd0b..e670dd07e0 100644 --- a/api/server/experimental.js +++ b/api/server/experimental.js @@ -10,15 +10,17 @@ const express = require('express'); const passport = require('passport'); const compression = require('compression'); const cookieParser = require('cookie-parser'); -const { logger, runAsSystem } = require('@librechat/data-schemas'); +const { logger, runAsSystem, tenantStorage } = require('@librechat/data-schemas'); const mongoSanitize = require('express-mongo-sanitize'); const { isEnabled, apiNotFound, ErrorController, + GenerationJobManager, QUERY_DEVTOOLS_HEADER, performStartupChecks, handleJsonParseError, + deleteAgentCheckpoint, initializeFileStorage, loadToolApprovalHooks, maybeInjectQueryDevtoolsBootstrap, @@ -311,6 +313,22 @@ if (cluster.isMaster) { await loadToolApprovalHooks(toolApproval?.enabled ? toolApproval.hooks : undefined, { basePath: path.resolve(__dirname, '../..'), }); + // Prune the paused run's durable checkpoint when its approval EXPIRES (a stale submit — + // this startup never runs the periodic sweeper) instead of leaving it until the Mongo + // TTL. Mirrors api/server/index.js's configureGenerationStreams wiring; safe here even + // though this startup runs the manager on constructor defaults (the setter never resets + // services). streamId === conversationId === the LangGraph thread_id. + GenerationJobManager.setApprovalExpiredHandler(async (conversationId, job) => { + // Resolve config in the PAUSED JOB's tenant/user scope (mirrors index.js): enter the + // tenant ALS context — getAppConfig args alone only key the cache. + await tenantStorage.run({ tenantId: job?.tenantId, userId: job?.userId }, async () => { + const currentConfig = await getAppConfig({ + userId: job?.userId, + tenantId: job?.tenantId, + }); + await deleteAgentCheckpoint(conversationId, currentConfig?.endpoints?.agents?.checkpointer); + }); + }); expiredFileSweepOptions = { appConfig, loadAppConfig: getAppConfig }; startExpiredFileSweepOnce(); await performStartupChecks(appConfig); diff --git a/api/server/index.js b/api/server/index.js index bcc38ae3fd..f88b2fd258 100644 --- a/api/server/index.js +++ b/api/server/index.js @@ -9,7 +9,7 @@ const passport = require('passport'); const compression = require('compression'); const cookieParser = require('cookie-parser'); const mongoSanitize = require('express-mongo-sanitize'); -const { logger, runAsSystem } = require('@librechat/data-schemas'); +const { logger, runAsSystem, tenantStorage } = require('@librechat/data-schemas'); const { isEnabled, apiNotFound, @@ -21,6 +21,7 @@ const { GenerationJobManager, QUERY_DEVTOOLS_HEADER, createStreamServices, + deleteAgentCheckpoint, initializeFileStorage, initializeDeploymentSkills, loadToolApprovalHooks, @@ -84,6 +85,19 @@ const configureGenerationStreams = () => { cleanupOnComplete: !isEnabled(process.env.STREAM_KEEP_COMPLETED_JOBS), }); GenerationJobManager.initialize(); + // Prune the paused run's durable checkpoint when its approval EXPIRES (periodic sweeper + // or a stale submit) instead of leaving it until the Mongo TTL. streamId === conversationId + // === the LangGraph thread_id. Config is resolved lazily per expiry so the prune always + // targets the currently configured checkpoint collections. + GenerationJobManager.setApprovalExpiredHandler(async (conversationId, job) => { + // Resolve config in the PAUSED JOB's tenant/user scope — the expiry runs outside any + // request context. Passing ids to getAppConfig only keys the cache; the Config query + // itself is ALS-scoped by the tenant-isolation plugin, so ENTER the tenant context. + await tenantStorage.run({ tenantId: job?.tenantId, userId: job?.userId }, async () => { + const appConfig = await getAppConfig({ userId: job?.userId, tenantId: job?.tenantId }); + await deleteAgentCheckpoint(conversationId, appConfig?.endpoints?.agents?.checkpointer); + }); + }); }; const startServer = async () => { diff --git a/api/server/routes/__test-utils__/convos-route-mocks.js b/api/server/routes/__test-utils__/convos-route-mocks.js index a0eb6fe312..769d2b61d8 100644 --- a/api/server/routes/__test-utils__/convos-route-mocks.js +++ b/api/server/routes/__test-utils__/convos-route-mocks.js @@ -14,6 +14,7 @@ module.exports = { restoreTenantContextFromReq: jest.fn((req, res, next) => next()), deleteConvoSharedLinksWithCleanup: jest.fn(), deleteAllSharedLinksWithCleanup: jest.fn(), + deleteAgentCheckpoints: jest.fn(), ...overrides, }), diff --git a/api/server/routes/__tests__/convos.spec.js b/api/server/routes/__tests__/convos.spec.js index 9c760f5057..6c78287003 100644 --- a/api/server/routes/__tests__/convos.spec.js +++ b/api/server/routes/__tests__/convos.spec.js @@ -23,6 +23,7 @@ describe('Convos Routes', () => { let convosRouter; const { deleteToolCalls, deleteConvos, saveConvo } = require('~/models'); const { + deleteAgentCheckpoints, deleteAllSharedLinksWithCleanup, deleteConvoSharedLinksWithCleanup, } = require('@librechat/api'); @@ -47,6 +48,20 @@ describe('Convos Routes', () => { }); describe('DELETE /all', () => { + it('prunes the deleted conversations’ agent checkpoints (bulk, ids from deleteConvos)', async () => { + // HITL: a paused conversation's durable checkpoint must not outlive the conversation. + const conversationIds = ['conv-a', 'conv-b']; + deleteConvos.mockResolvedValue({ deletedCount: 2, conversationIds }); + deleteToolCalls.mockResolvedValue({ deletedCount: 0 }); + deleteAllSharedLinksWithCleanup.mockResolvedValue({ deletedCount: 0 }); + + const response = await request(app).delete('/api/convos/all'); + + expect(response.status).toBe(201); + expect(deleteAgentCheckpoints).toHaveBeenCalledTimes(1); + expect(deleteAgentCheckpoints.mock.calls[0][0]).toEqual(conversationIds); + }); + it('should delete all conversations, tool calls, and shared links for a user', async () => { const mockDbResponse = { deletedCount: 5, diff --git a/api/server/routes/convos.js b/api/server/routes/convos.js index 879b701d2e..b38e73b576 100644 --- a/api/server/routes/convos.js +++ b/api/server/routes/convos.js @@ -3,6 +3,7 @@ const express = require('express'); const { sleep } = require('@librechat/agents'); const { isEnabled, + deleteAgentCheckpoints, resolveImportMaxFileSize, restoreTenantContextFromReq, deleteAllSharedLinksWithCleanup, @@ -111,7 +112,7 @@ router.get('/gen_title/:conversationId', async (req, res) => { } }); -router.delete('/', async (req, res) => { +router.delete('/', configMiddleware, async (req, res) => { let filter = {}; const { conversationId, source, thread_id, endpoint } = req.body?.arg ?? {}; @@ -144,6 +145,12 @@ router.delete('/', async (req, res) => { try { const dbResponse = await db.deleteConvos(req.user.id, filter); + // HITL: prune the deleted conversations' durable checkpoints — a paused run's + // checkpoint would otherwise persist until the Mongo TTL. Never throws. + await deleteAgentCheckpoints( + dbResponse.conversationIds, + req.config?.endpoints?.[EModelEndpoint.agents]?.checkpointer, + ); if (filter.conversationId) { await db.deleteToolCalls(req.user.id, filter.conversationId); await deleteConvoSharedLinksWithCleanup(req.user.id, filter.conversationId); @@ -155,9 +162,14 @@ router.delete('/', async (req, res) => { } }); -router.delete('/all', async (req, res) => { +router.delete('/all', configMiddleware, async (req, res) => { try { const dbResponse = await db.deleteConvos(req.user.id, {}); + // HITL: prune ALL the deleted conversations' durable checkpoints in one bulk pass. + await deleteAgentCheckpoints( + dbResponse.conversationIds, + req.config?.endpoints?.[EModelEndpoint.agents]?.checkpointer, + ); await db.deleteToolCalls(req.user.id); await deleteAllSharedLinksWithCleanup(req.user.id); res.status(201).json(dbResponse); diff --git a/packages/api/src/agents/checkpointer.integration.spec.ts b/packages/api/src/agents/checkpointer.integration.spec.ts index e11451c54e..cc67b96817 100644 --- a/packages/api/src/agents/checkpointer.integration.spec.ts +++ b/packages/api/src/agents/checkpointer.integration.spec.ts @@ -5,6 +5,7 @@ import { emptyCheckpoint, ERROR, INTERRUPT } from '@langchain/langgraph-checkpoi import { getAgentCheckpointer, deleteAgentCheckpoint, + deleteAgentCheckpoints, __resetCheckpointerForTests, } from './checkpointer'; @@ -120,6 +121,41 @@ describe('checkpointer (mongodb-memory-server integration)', () => { it('deleteAgentCheckpoint is a no-op for an undefined threadId', async () => { await expect(deleteAgentCheckpoint(undefined, MONGO_CFG)).resolves.toBeUndefined(); }); + + it('deleteAgentCheckpoints bulk-prunes exactly the given threads (checkpoints AND writes)', async () => { + // The bulk path behind conversation deletion / delete-all / account deletion: + // one $in deleteMany per collection instead of two round-trips per thread. + const saver = await getAgentCheckpointer(MONGO_CFG); + const threadA = `convo-${new mongoose.Types.ObjectId().toString()}`; + const threadB = `convo-${new mongoose.Types.ObjectId().toString()}`; + const threadC = `convo-${new mongoose.Types.ObjectId().toString()}`; + + await seedInterruptCheckpoint(saver!, threadA); + await seedInterruptCheckpoint(saver!, threadB); + await seedInterruptCheckpoint(saver!, threadC); + + // Falsy entries are skipped rather than widening the delete. + await deleteAgentCheckpoints([threadA, undefined, threadB, null], MONGO_CFG); + + expect(await saver!.getTuple(readConfig(threadA))).toBeUndefined(); + expect(await saver!.getTuple(readConfig(threadB))).toBeUndefined(); + expect(await saver!.getTuple(readConfig(threadC))).toBeDefined(); + + const db = mongoose.connection.db!; + const writesFilter = { thread_id: { $in: [threadA, threadB] } }; + expect(await db.collection('agent_checkpoints').countDocuments(writesFilter)).toBe(0); + expect(await db.collection('agent_checkpoint_writes').countDocuments(writesFilter)).toBe(0); + // The untouched thread keeps its interrupt write row. + expect( + await db.collection('agent_checkpoint_writes').countDocuments({ thread_id: threadC }), + ).toBe(1); + }); + + it('deleteAgentCheckpoints is a no-op for an empty or all-falsy list', async () => { + await expect(deleteAgentCheckpoints([], MONGO_CFG)).resolves.toBeUndefined(); + await expect(deleteAgentCheckpoints([undefined, null], MONGO_CFG)).resolves.toBeUndefined(); + await expect(deleteAgentCheckpoints(undefined, MONGO_CFG)).resolves.toBeUndefined(); + }); }); describe('LazyMongoSaver (lazy persistence — mongodb-memory-server)', () => { diff --git a/packages/api/src/agents/checkpointer.ts b/packages/api/src/agents/checkpointer.ts index 2a24810e03..79b6954346 100644 --- a/packages/api/src/agents/checkpointer.ts +++ b/packages/api/src/agents/checkpointer.ts @@ -394,6 +394,50 @@ export async function deleteAgentCheckpoint( } } +/** + * Bulk variant of {@link deleteAgentCheckpoint} for terminal transitions that cover MANY + * threads at once — deleting conversations, "delete all", account deletion. One indexed + * `deleteMany` per collection instead of two round-trips per thread. Deletes through the + * same live mongoose connection the saver is built on, using the same resolved collection + * names; like the single-thread variant it no-ops in memory mode or before Mongo is + * connected, and never throws (the conversations are already gone — the Mongo TTL remains + * the backstop for anything this misses). + * + * @param threadIds - LangGraph `thread_id`s (LibreChat conversationIds); falsy entries skipped. + */ +export async function deleteAgentCheckpoints( + threadIds: Array | undefined, + cfg?: TCheckpointerConfig, +): Promise { + const ids = (threadIds ?? []).filter((id): id is string => Boolean(id)); + if (ids.length === 0) { + return; + } + // Reuse the saver gate: memory mode / no connection ⇒ nothing durable to delete. + const saver = await getAgentCheckpointer(cfg); + if (!saver) { + return; + } + const resolved = resolveCheckpointerConfig(cfg); + try { + const db = mongoose.connection.db; + if (!db) { + return; + } + await Promise.all([ + db.collection(resolved.checkpointCollectionName).deleteMany({ thread_id: { $in: ids } }), + db + .collection(resolved.checkpointWritesCollectionName) + .deleteMany({ thread_id: { $in: ids } }), + ]); + } catch (err) { + logger.warn( + `[checkpointer] Failed to bulk-delete checkpoints for ${ids.length} thread(s):`, + err, + ); + } +} + /** Test-only: drop the memoized saver so a fresh build is forced. */ export function __resetCheckpointerForTests(): void { saverPromise = undefined; diff --git a/packages/api/src/stream/GenerationJobManager.ts b/packages/api/src/stream/GenerationJobManager.ts index 1509335276..fa66f0a051 100644 --- a/packages/api/src/stream/GenerationJobManager.ts +++ b/packages/api/src/stream/GenerationJobManager.ts @@ -152,6 +152,8 @@ interface RuntimeJobState { resolveReady: () => void; finalEvent?: t.ServerSentEvent; errorEvent?: string; + /** Approval-expired host cleanup already ran for this runtime (relay path is swept repeatedly). */ + approvalCleanupRan?: boolean; syncSent: boolean; earlyEventBuffer: t.ServerSentEvent[]; hasSubscriber: boolean; @@ -207,6 +209,15 @@ class GenerationJobManagerClass { /** Whether to cleanup event transport immediately on job completion */ private _cleanupOnComplete = true; + /** + * Host cleanup fired after an approval EXPIRES (periodic sweeper or a stale submit) — + * e.g. prune the paused run's durable checkpoint eagerly instead of letting it sit + * until its store TTL. Best-effort: failures are logged, never break the expiry. + */ + private _onApprovalExpired: + | ((streamId: string, job?: SerializableJobData | null) => void | Promise) + | null = null; + constructor(options?: GenerationJobManagerOptions) { this.jobStore = options?.jobStore ?? new InMemoryJobStore({ ttlAfterComplete: 0, maxJobs: 1000 }); @@ -280,6 +291,19 @@ class GenerationJobManagerClass { ); } + /** + * Register a host callback fired after an approval EXPIRES — from the periodic sweeper or + * a stale submit — e.g. to prune the paused run's durable checkpoint eagerly instead of + * waiting out its TTL. Unlike {@link configure} this never resets services, so it is safe + * to call from any startup path (including ones that run on constructor defaults). The + * `streamId` argument equals the LangGraph `thread_id` (LibreChat's conversationId). + */ + setApprovalExpiredHandler( + handler: ((streamId: string, job?: SerializableJobData | null) => void | Promise) | null, + ): void { + this._onApprovalExpired = handler; + } + /** * Check if using Redis stores. */ @@ -1690,10 +1714,48 @@ class GenerationJobManagerClass { } catch (err) { logger.error(`[GenerationJobManager] Failed to notify expired approval ${streamId}`, err); } + await this.runApprovalExpiredHandler(streamId); this.runningJobs.delete(streamId); return true; } + /** + * Invoke the host approval-expired cleanup, passing the job so the host can resolve + * tenant/user-scoped config (the expiry runs outside any request context). Best-effort: + * the job read and the handler itself may fail without breaking the expiry. + */ + private async runApprovalExpiredHandler( + streamId: string, + job?: SerializableJobData | null, + ): Promise { + if (!this._onApprovalExpired) { + return; + } + // Dedup across the expiry paths: a locally expired approval (expireApproval) stays in + // the store/runtime for the completed-job TTL, so later sweeps re-enter the relay + // branch for the same aborted approval — run the cleanup once per runtime lifetime. + const runtime = this.runtimeState.get(streamId); + if (runtime?.approvalCleanupRan) { + return; + } + if (runtime) { + runtime.approvalCleanupRan = true; + } + let resolvedJob = job; + if (resolvedJob === undefined) { + try { + resolvedJob = await this.jobStore.getJob(streamId); + } catch { + resolvedJob = null; + } + } + try { + await this._onApprovalExpired(streamId, resolvedJob); + } catch (err) { + logger.warn(`[GenerationJobManager] Approval-expired cleanup failed for ${streamId}`, err); + } + } + private async expireStaleApprovals(): Promise { let changed = false; for (const streamId of this.runtimeState.keys()) { @@ -1715,16 +1777,23 @@ class GenerationJobManagerClass { // expiry* and we haven't emitted here, relay the terminal error to our subscriber. // The `errorEvent` flag (set by emitError) keeps this idempotent vs the win path. const runtime = this.runtimeState.get(streamId); - if ( - job?.status === 'aborted' && - job.error === APPROVAL_EXPIRED_ERROR && - !runtime?.errorEvent - ) { - try { - await this.emitError(streamId, APPROVAL_EXPIRED_ERROR); - } catch (err) { - logger.error(`[GenerationJobManager] Failed to relay expired approval ${streamId}`, err); + if (job?.status === 'aborted' && job.error === APPROVAL_EXPIRED_ERROR) { + if (!runtime?.errorEvent) { + try { + await this.emitError(streamId, APPROVAL_EXPIRED_ERROR); + } catch (err) { + logger.error( + `[GenerationJobManager] Failed to relay expired approval ${streamId}`, + err, + ); + } } + // The winning store cleanup (`cleanupRequiresActionIndex`) transitions status + // directly and can't run host cleanup — do it on relay. Deliberately NOT gated on + // `errorEvent`: a reconnect seeds that flag from the aborted job, which must not + // suppress the (idempotent) prune. The handler dedups per runtime lifetime, which + // also covers approvals expired LOCALLY via expireApproval. + await this.runApprovalExpiredHandler(streamId, job); changed = this.runningJobs.delete(streamId) || changed; continue; } diff --git a/packages/api/src/stream/__tests__/pendingAction.spec.ts b/packages/api/src/stream/__tests__/pendingAction.spec.ts index 4edfed9619..0098e9865e 100644 --- a/packages/api/src/stream/__tests__/pendingAction.spec.ts +++ b/packages/api/src/stream/__tests__/pendingAction.spec.ts @@ -199,6 +199,111 @@ describe('ApprovalLifecycle via GenerationJobManager.approvals (in-memory)', () }); }); + describe('expireApproval → approval-expired handler', () => { + // The host registers this to prune the paused run's durable checkpoint eagerly on + // expiry (sweeper or stale submit) instead of waiting out the checkpoint TTL. + test('fires the registered handler with the streamId after a successful expiry', async () => { + const streamId = 'stream-expire-handler'; + await manager.createJob(streamId, 'user-1'); + await manager.approvals.pause(streamId, buildAction(streamId, { actionId: 'action-A' })); + + const handler = jest.fn(); + manager.setApprovalExpiredHandler(handler); + + expect(await manager.expireApproval(streamId, 'action-A')).toBe(true); + expect(handler).toHaveBeenCalledTimes(1); + // The expired job rides along so the host can resolve tenant/user-scoped config. + expect(handler).toHaveBeenCalledWith(streamId, expect.objectContaining({ userId: 'user-1' })); + + // The aborted job outlives the expiry (completed-job TTL), so the next sweep enters + // the relay branch for the SAME approval — the cleanup must not run a second time. + await ( + manager as unknown as { expireStaleApprovals(): Promise } + ).expireStaleApprovals(); + expect(handler).toHaveBeenCalledTimes(1); + }); + + test('relays a store-won expiry through the handler (multi-replica path)', async () => { + const streamId = 'stream-expire-relay'; + await manager.createJob(streamId, 'user-1'); + await manager.approvals.pause(streamId, buildAction(streamId)); + + // Another replica's store cleanup wins the expiry CAS: the status flips via the + // lifecycle primitive with NO emit, NO handler, and no errorEvent on this replica. + expect(await manager.approvals.expire(streamId)).toBe(true); + + const handler = jest.fn(); + manager.setApprovalExpiredHandler(handler); + // This replica's sweep observes the already-aborted expiry and relays it. + await ( + manager as unknown as { expireStaleApprovals(): Promise } + ).expireStaleApprovals(); + + expect(handler).toHaveBeenCalledTimes(1); + expect(handler).toHaveBeenCalledWith( + streamId, + expect.objectContaining({ userId: 'user-1', status: 'aborted' }), + ); + + // Repeated sweeps must not re-run the (idempotent but not free) cleanup. + await ( + manager as unknown as { expireStaleApprovals(): Promise } + ).expireStaleApprovals(); + expect(handler).toHaveBeenCalledTimes(1); + }); + + test('relay cleanup still runs when the terminal error is already cached (reconnect)', async () => { + const streamId = 'stream-expire-relay-cached'; + await manager.createJob(streamId, 'user-1'); + await manager.approvals.pause(streamId, buildAction(streamId)); + expect(await manager.approvals.expire(streamId)).toBe(true); // store-won CAS + + // A reconnect seeds runtime.errorEvent from the aborted job BEFORE any sweep — + // that must gate the relay emit, not the checkpoint cleanup. + const internals = manager as unknown as { + runtimeState: Map; + expireStaleApprovals(): Promise; + }; + const runtime = internals.runtimeState.get(streamId); + expect(runtime).toBeDefined(); + runtime!.errorEvent = 'cached-terminal-error'; + + const handler = jest.fn(); + manager.setApprovalExpiredHandler(handler); + await internals.expireStaleApprovals(); + + expect(handler).toHaveBeenCalledTimes(1); + expect(handler).toHaveBeenCalledWith( + streamId, + expect.objectContaining({ status: 'aborted' }), + ); + }); + + test('does NOT fire when nothing was expired (failed CAS)', async () => { + const streamId = 'stream-expire-handler-noop'; + await manager.createJob(streamId, 'user-1'); // running — no pending action to expire + + const handler = jest.fn(); + manager.setApprovalExpiredHandler(handler); + + expect(await manager.expireApproval(streamId)).toBe(false); + expect(handler).not.toHaveBeenCalled(); + }); + + test('a throwing handler never breaks the expiry itself', async () => { + const streamId = 'stream-expire-handler-throws'; + await manager.createJob(streamId, 'user-1'); + await manager.approvals.pause(streamId, buildAction(streamId)); + + manager.setApprovalExpiredHandler(() => { + throw new Error('prune failed'); + }); + + expect(await manager.expireApproval(streamId)).toBe(true); + expect(await manager.getJobStatus(streamId)).toBe('aborted'); + }); + }); + describe('facade integration', () => { test('requires_action drops the running count but keeps the user-active set', async () => { const streamId = 'stream-counts'; diff --git a/packages/data-schemas/src/methods/conversation.spec.ts b/packages/data-schemas/src/methods/conversation.spec.ts index 1d9621cc0f..124522a9b0 100644 --- a/packages/data-schemas/src/methods/conversation.spec.ts +++ b/packages/data-schemas/src/methods/conversation.spec.ts @@ -1110,7 +1110,7 @@ describe('Conversation Operations', () => { expect(tag?.count).toBe(2); }); - it('should still decrement tag counts when message deletion fails after the delete', async () => { + it('still decrements tag counts AND returns the deleted ids when message deletion fails', async () => { await ConversationTag.create({ user: 'user123', tag: 'work', count: 2, position: 1 }); const convoId = uuidv4(); await Conversation.create({ @@ -1122,9 +1122,13 @@ describe('Conversation Operations', () => { deleteMessages.mockRejectedValueOnce(new Error('message cleanup failed')); - await expect(deleteConvos('user123', { conversationId: convoId })).rejects.toThrow( - 'message cleanup failed', - ); + // Post-delete cleanup is best-effort: the conversations are already gone, so the + // caller must still receive the deleted ids (downstream cleanup — e.g. agent + // checkpoint pruning — depends on them, and a retry would find nothing). + const result = await deleteConvos('user123', { conversationId: convoId }); + expect(result.deletedCount).toBe(1); + expect(result.conversationIds).toEqual([convoId]); + expect(result.messages.deletedCount).toBe(0); const tag = await ConversationTag.findOne({ user: 'user123', tag: 'work' }).lean(); expect(tag?.count).toBe(1); diff --git a/packages/data-schemas/src/methods/conversation.ts b/packages/data-schemas/src/methods/conversation.ts index cc2594fd26..723e464681 100644 --- a/packages/data-schemas/src/methods/conversation.ts +++ b/packages/data-schemas/src/methods/conversation.ts @@ -64,7 +64,7 @@ export interface ConversationMethods { deleteConvos( user: string, filter: FilterQuery, - ): Promise; + ): Promise; } export function createConversationMethods( @@ -793,10 +793,21 @@ export function createConversationMethods( await decrementTagCounts(mongoose, user, tagDecrements); } - const deleteMessagesResult = await deleteMessages({ - conversationId: { $in: conversationIds }, - user, - }); + /** + * Post-delete cleanup is best-effort: the conversations are already gone, so a + * thrown error here would hide the deletion from the caller — dropping the + * `conversationIds` that downstream cleanup (e.g. agent-checkpoint pruning) + * needs, with no way to recover them on retry (the query finds nothing). + */ + let deleteMessagesResult: DeleteResult = { acknowledged: false, deletedCount: 0 }; + try { + deleteMessagesResult = await deleteMessages({ + conversationId: { $in: conversationIds }, + user, + }); + } catch (error) { + logger.error('[deleteConvos] Conversations deleted but message cleanup failed', error); + } /** * Refresh project stats after message cleanup so a stats-refresh error cannot @@ -804,14 +815,21 @@ export function createConversationMethods( * conversations' messages. */ if (deleted && projectIds.size > 0) { - await Promise.all( - [...projectIds].map((projectId) => - refreshChatProjectStatsForUser(mongoose, user, projectId), - ), - ); + try { + await Promise.all( + [...projectIds].map((projectId) => + refreshChatProjectStatsForUser(mongoose, user, projectId), + ), + ); + } catch (error) { + logger.error('[deleteConvos] Conversations deleted but stats refresh failed', error); + } } - return { ...deleteConvoResult, messages: deleteMessagesResult }; + // conversationIds lets callers run sibling cleanup that lives in higher layers + // (e.g. pruning the conversations' durable agent checkpoints) without re-querying + // documents that no longer exist. + return { ...deleteConvoResult, messages: deleteMessagesResult, conversationIds }; } catch (error) { logger.error('[deleteConvos] Error deleting conversations and messages', error); throw error;