LibreChat/api/server/routes/convos.js
Danny Avila 84fa6aa820
🧹 feat: Eager HITL Checkpoint Cleanup (Expiry + Deletion) & Full-Wiring E2E (#14123)
* 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.
2026-07-05 11:29:30 -04:00

387 lines
12 KiB
JavaScript

const multer = require('multer');
const express = require('express');
const { sleep } = require('@librechat/agents');
const {
isEnabled,
deleteAgentCheckpoints,
resolveImportMaxFileSize,
restoreTenantContextFromReq,
deleteAllSharedLinksWithCleanup,
deleteConvoSharedLinksWithCleanup,
} = require('@librechat/api');
const { logger } = require('@librechat/data-schemas');
const { CacheKeys, EModelEndpoint } = require('librechat-data-provider');
const {
createImportLimiters,
validateConvoAccess,
createForkLimiters,
configMiddleware,
} = require('~/server/middleware');
const { forkConversation, duplicateConversation } = require('~/server/utils/import/fork');
const { storage, importFileFilter } = require('~/server/routes/files/multer');
const requireJwtAuth = require('~/server/middleware/requireJwtAuth');
const { importConversations } = require('~/server/utils/import');
const getLogStores = require('~/cache/getLogStores');
const db = require('~/models');
const assistantClients = {
[EModelEndpoint.azureAssistants]: require('~/server/services/Endpoints/azureAssistants'),
[EModelEndpoint.assistants]: require('~/server/services/Endpoints/assistants'),
};
const router = express.Router();
router.use(requireJwtAuth);
const isValidProjectFilter = (projectId) =>
!projectId || projectId === 'unassigned' || /^[a-f\d]{24}$/i.test(projectId);
router.get('/', async (req, res) => {
const limit = parseInt(req.query.limit, 10) || 25;
const cursor = req.query.cursor;
const isArchived = isEnabled(req.query.isArchived);
const search = req.query.search ? decodeURIComponent(req.query.search) : undefined;
const sortBy = req.query.sortBy || 'updatedAt';
const sortDirection = req.query.sortDirection || 'desc';
const projectId = Array.isArray(req.query.projectId)
? req.query.projectId[0]
: req.query.projectId;
if (!isValidProjectFilter(projectId)) {
return res.status(400).json({ error: 'projectId must be a valid project id or unassigned' });
}
let tags;
if (req.query.tags) {
tags = Array.isArray(req.query.tags) ? req.query.tags : [req.query.tags];
}
try {
const result = await db.getConvosByCursor(req.user.id, {
cursor,
limit,
isArchived,
tags,
search,
sortBy,
sortDirection,
projectId,
});
res.status(200).json(result);
} catch (error) {
logger.error('Error fetching conversations', error);
res.status(500).json({ error: 'Error fetching conversations' });
}
});
router.get('/:conversationId', async (req, res) => {
const { conversationId } = req.params;
const convo = await db.getConvo(req.user.id, conversationId);
if (convo) {
res.status(200).json(convo);
} else {
res.status(404).end();
}
});
router.get('/gen_title/:conversationId', async (req, res) => {
const { conversationId } = req.params;
const titleCache = getLogStores(CacheKeys.GEN_TITLE);
const key = `${req.user.id}-${conversationId}`;
let title = await titleCache.get(key);
if (!title) {
// Exponential backoff: 500ms, 1s, 2s, 4s, 8s (total ~15.5s max wait)
const delays = [500, 1000, 2000, 4000, 8000];
for (const delay of delays) {
await sleep(delay);
title = await titleCache.get(key);
if (title) {
break;
}
}
}
if (title) {
await titleCache.delete(key);
res.status(200).json({ title });
} else {
res.status(404).json({
message: "Title not found or method not implemented for the conversation's endpoint",
});
}
});
router.delete('/', configMiddleware, async (req, res) => {
let filter = {};
const { conversationId, source, thread_id, endpoint } = req.body?.arg ?? {};
// Prevent deletion of all conversations
if (!conversationId && !source && !thread_id && !endpoint) {
return res.status(400).json({
error: 'no parameters provided',
});
}
if (conversationId) {
filter = { conversationId };
} else if (source === 'button') {
return res.status(200).send('No conversationId provided');
}
if (
typeof endpoint !== 'undefined' &&
Object.prototype.propertyIsEnumerable.call(assistantClients, endpoint)
) {
/** @type {{ openai: OpenAI }} */
const { openai } = await assistantClients[endpoint].initializeClient({ req, res });
try {
const response = await openai.beta.threads.delete(thread_id);
logger.debug('Deleted OpenAI thread:', response);
} catch (error) {
logger.error('Error deleting OpenAI thread:', error);
}
}
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);
}
res.status(201).json(dbResponse);
} catch (error) {
logger.error('Error clearing conversations', error);
res.status(500).send('Error clearing conversations');
}
});
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);
} catch (error) {
logger.error('Error clearing conversations', error);
res.status(500).send('Error clearing conversations');
}
});
/**
* Archives or unarchives a conversation.
* @route POST /archive
* @param {string} req.body.arg.conversationId - The conversation ID to archive/unarchive.
* @param {boolean} req.body.arg.isArchived - Whether to archive (true) or unarchive (false).
* @returns {object} 200 - The updated conversation object.
*/
router.post('/archive', validateConvoAccess, async (req, res) => {
const { conversationId, isArchived } = req.body?.arg ?? {};
if (!conversationId) {
return res.status(400).json({ error: 'conversationId is required' });
}
if (typeof isArchived !== 'boolean') {
return res.status(400).json({ error: 'isArchived must be a boolean' });
}
try {
const dbResponse = await db.saveConvo(
{
userId: req?.user?.id,
isTemporary: req?.body?.isTemporary,
interfaceConfig: req?.config?.interfaceConfig,
},
{ conversationId, isArchived },
{ context: `POST /api/convos/archive ${conversationId}` },
);
res.status(200).json(dbResponse);
} catch (error) {
logger.error('Error archiving conversation', error);
res.status(500).send('Error archiving conversation');
}
});
router.post('/pin', validateConvoAccess, async (req, res) => {
const { conversationId, pinned } = req.body?.arg ?? {};
if (!conversationId) {
return res.status(400).json({ error: 'conversationId is required' });
}
if (pinned === undefined) {
return res.status(400).json({ error: 'pinned is required' });
}
if (typeof pinned !== 'boolean') {
return res.status(400).json({ error: 'pinned must be a boolean' });
}
try {
const dbResponse = await db.saveConvo(
{ userId: req.user.id },
{ conversationId, pinned },
{ context: `POST /api/convos/pin ${conversationId}` },
);
res.status(200).json(dbResponse);
} catch (error) {
logger.error('Error pinning conversation', error);
res.status(500).send('Error pinning conversation');
}
});
/** Maximum allowed length for conversation titles */
const MAX_CONVO_TITLE_LENGTH = 1024;
/**
* Updates a conversation's title.
* @route POST /update
* @param {string} req.body.arg.conversationId - The conversation ID to update.
* @param {string} req.body.arg.title - The new title for the conversation.
* @returns {object} 201 - The updated conversation object.
*/
router.post('/update', validateConvoAccess, async (req, res) => {
const { conversationId, title } = req.body?.arg ?? {};
if (!conversationId) {
return res.status(400).json({ error: 'conversationId is required' });
}
if (title === undefined) {
return res.status(400).json({ error: 'title is required' });
}
if (typeof title !== 'string') {
return res.status(400).json({ error: 'title must be a string' });
}
const sanitizedTitle = title.trim().slice(0, MAX_CONVO_TITLE_LENGTH);
try {
const dbResponse = await db.saveConvo(
{
userId: req?.user?.id,
isTemporary: req?.body?.isTemporary,
interfaceConfig: req?.config?.interfaceConfig,
},
{ conversationId, title: sanitizedTitle },
{ context: `POST /api/convos/update ${conversationId}` },
);
res.status(201).json(dbResponse);
} catch (error) {
logger.error('Error updating conversation', error);
res.status(500).send('Error updating conversation');
}
});
const { importIpLimiter, importUserLimiter } = createImportLimiters();
/** Fork and duplicate share one rate-limit budget (same "clone" operation class) */
const { forkIpLimiter, forkUserLimiter } = createForkLimiters();
const importMaxFileSize = resolveImportMaxFileSize();
const upload = multer({
storage,
fileFilter: importFileFilter,
limits: { fileSize: importMaxFileSize },
});
const uploadSingle = upload.single('file');
function handleUpload(req, res, next) {
uploadSingle(req, res, (err) => {
if (err && err.code === 'LIMIT_FILE_SIZE') {
return res.status(413).json({ message: 'File exceeds the maximum allowed size' });
}
if (err) {
return next(err);
}
next();
});
}
/**
* Imports a conversation from a JSON file and saves it to the database.
* @route POST /import
* @param {Express.Multer.File} req.file - The JSON file to import.
* @returns {object} 201 - success response - application/json
*/
router.post(
'/import',
importIpLimiter,
importUserLimiter,
configMiddleware,
handleUpload,
restoreTenantContextFromReq,
async (req, res) => {
try {
/* TODO: optimize to return imported conversations and add manually */
await importConversations({
filepath: req.file.path,
requestUserId: req.user.id,
userRole: req.user.role,
interfaceConfig: req.config?.interfaceConfig,
});
res.status(201).json({ message: 'Conversation(s) imported successfully' });
} catch (error) {
logger.error('Error processing file', error);
res.status(500).send('Error processing file');
}
},
);
/**
* POST /fork
* This route handles forking a conversation based on the TForkConvoRequest and responds with TForkConvoResponse.
* @route POST /fork
* @param {express.Request<{}, TForkConvoResponse, TForkConvoRequest>} req - Express request object.
* @param {express.Response<TForkConvoResponse>} res - Express response object.
* @returns {Promise<void>} - The response after forking the conversation.
*/
router.post('/fork', forkIpLimiter, forkUserLimiter, async (req, res) => {
try {
/** @type {TForkConvoRequest} */
const { conversationId, messageId, option, splitAtTarget, latestMessageId } = req.body;
const result = await forkConversation({
requestUserId: req.user.id,
originalConvoId: conversationId,
targetMessageId: messageId,
latestMessageId,
records: true,
splitAtTarget,
option,
});
res.json(result);
} catch (error) {
logger.error('Error forking conversation:', error);
res.status(500).send('Error forking conversation');
}
});
router.post('/duplicate', forkIpLimiter, forkUserLimiter, async (req, res) => {
const { conversationId, title } = req.body;
try {
const result = await duplicateConversation({
userId: req.user.id,
conversationId,
title,
});
res.status(201).json(result);
} catch (error) {
logger.error('Error duplicating conversation:', error);
res.status(500).send('Error duplicating conversation');
}
});
module.exports = router;