diff --git a/.gitignore b/.gitignore index d28d18530e..89d6c2e99a 100644 --- a/.gitignore +++ b/.gitignore @@ -52,6 +52,7 @@ __blobstorage__/**/* # Deployed apps should consider commenting these lines out: # see https://npmjs.org/doc/faq.html#Should-I-check-my-node_modules-folder-into-git node_modules/ +.node_modules-* meili_data/ api/node_modules/ client/node_modules/ diff --git a/api/server/middleware/limiters/uploadLimiters.js b/api/server/middleware/limiters/uploadLimiters.js index 8c878cfa86..ab138b8679 100644 --- a/api/server/middleware/limiters/uploadLimiters.js +++ b/api/server/middleware/limiters/uploadLimiters.js @@ -80,6 +80,39 @@ const createFileLimiters = () => { return { fileUploadIpLimiter, fileUploadUserLimiter }; }; +/** + * Per-user limiter for the `/files/usage` TTL hold. Deliberately separate from + * the upload limiters: a metadata touch must not consume upload quota, but it + * still writes to the DB and so cannot go unmetered. Sized well above the + * enqueue-driven call rate a real client produces. + */ +const createFileUsageLimiter = () => { + const windowMinutes = parseInt(process.env.FILE_USAGE_USER_WINDOW) || 15; + const max = parseInt(process.env.FILE_USAGE_USER_MAX) || 120; + const windowMs = windowMinutes * 60 * 1000; + + return rateLimit({ + windowMs, + max, + handler: async (req, res) => { + const type = ViolationTypes.FILE_UPLOAD_LIMIT; + await logViolation( + req, + res, + type, + { type, max, limiter: 'user', windowInMinutes: windowMinutes }, + process.env.FILE_UPLOAD_VIOLATION_SCORE, + ); + res.status(429).json({ message: 'Too many file usage requests. Try again later' }); + }, + keyGenerator: function (req) { + return req.user?.id; + }, + store: limiterCache('file_usage_user_limiter'), + }); +}; + module.exports = { createFileLimiters, + createFileUsageLimiter, }; diff --git a/api/server/routes/files/files.js b/api/server/routes/files/files.js index d98ed3b6f8..458d9ec1d0 100644 --- a/api/server/routes/files/files.js +++ b/api/server/routes/files/files.js @@ -3,6 +3,7 @@ const express = require('express'); const { logger, SystemCapabilities } = require('@librechat/data-schemas'); const { logAxiosError, + getApprovalTtlMs, refreshS3FileUrls, handleFilesUsageRequest, shouldUseUploadSse, @@ -144,15 +145,19 @@ router.get('/config', async (req, res) => { /** * POST /files/usage * - * Owner-scoped TTL touch for uploads held in a client-side queue (mid-run + * Owner-scoped TTL hold for uploads sitting in a client-side queue (mid-run * queued messages), so the upload-window TTL cannot reap them before drain. - * Thin wrapper: validation, cap, and best-effort semantics live in - * `@librechat/api` (`handleFilesUsageRequest`). + * Extends the deadline rather than clearing it; the real release happens at + * send. The approval window is passed through so a queue waiting on a paused + * run outlives that pause. Thin wrapper: validation, cap, hold window, and + * best-effort semantics live in `@librechat/api` (`handleFilesUsageRequest`). */ router.post('/usage', async (req, res) => { try { + const checkpointerCfg = req.config?.endpoints?.[EModelEndpoint.agents]?.checkpointer; const { status, body } = await handleFilesUsageRequest(req.user ?? {}, req.body ?? {}, { - updateFilesUsage: db.updateFilesUsage, + extendFilesTTL: db.extendFilesTTL, + approvalTtlMs: getApprovalTtlMs(checkpointerCfg), }); return res.status(status).json(body); } catch (error) { diff --git a/api/server/routes/files/files.test.js b/api/server/routes/files/files.test.js index 6cd60d43d6..d087fc508d 100644 --- a/api/server/routes/files/files.test.js +++ b/api/server/routes/files/files.test.js @@ -938,7 +938,7 @@ describe('File Routes - Delete with Agent Access', () => { }); describe('POST /files/usage', () => { - it('marks owned files used and clears the upload TTL', async () => { + const createQueuedFile = async (expiresAt) => { const ownFileId = uuidv4(); await createFile({ user: otherUserId, @@ -948,31 +948,110 @@ describe('File Routes - Delete with Agent Access', () => { bytes: 10, type: 'image/png', }); - await File.updateOne({ file_id: ownFileId }, { $set: { expiresAt: new Date() } }); + await File.updateOne({ file_id: ownFileId }, { $set: { expiresAt } }); + return ownFileId; + }; + + it('extends the upload TTL of owned files without clearing it', async () => { + const soon = new Date(Date.now() + 60 * 1000); + const ownFileId = await createQueuedFile(soon); const response = await request(app) .post('/files/usage') .send({ file_ids: [ownFileId] }); expect(response.status).toBe(200); - expect(response.body).toEqual({ marked: 1 }); - const marked = await File.findOne({ file_id: ownFileId }).lean(); - expect(marked.usage).toBe(1); - expect(marked.expiresAt).toBeUndefined(); + expect(response.body).toEqual({ held: 1 }); + const held = await File.findOne({ file_id: ownFileId }).lean(); + /* The hold must remain a hold: still reapable, just later. */ + expect(held.expiresAt).toBeDefined(); + expect(held.expiresAt.getTime()).toBeGreaterThan(soon.getTime()); + /* The 24h baseline plus the default 24h approval window, so a queue + * waiting on a paused run outlives that pause. Renewed from now, but + * never past the ceiling measured from upload time. */ + const HOUR = 60 * 60 * 1000; + expect(held.expiresAt.getTime()).toBeGreaterThan(Date.now() + 47 * HOUR); + expect(held.expiresAt.getTime()).toBeLessThanOrEqual( + held.createdAt.getTime() + 24 * HOUR + 8 * 24 * HOUR, + ); + /* A queue touch is not a send, so it must not inflate usage. */ + expect(held.usage).toBe(0); + }); + + it('cannot be replayed to preserve a file indefinitely', async () => { + const ownFileId = await createQueuedFile(new Date(Date.now() + 60 * 1000)); + + const first = await request(app) + .post('/files/usage') + .send({ file_ids: [ownFileId] }); + expect(first.body).toEqual({ held: 1 }); + + for (let i = 0; i < 5; i++) { + const repeat = await request(app) + .post('/files/usage') + .send({ file_ids: [ownFileId] }); + expect(repeat.status).toBe(200); + } + + /* Every renewal is clamped to the ceiling measured from upload time, so + * replay converges there instead of advancing a window per call. */ + const HOUR = 60 * 60 * 1000; + const held = await File.findOne({ file_id: ownFileId }).lean(); + expect(held.expiresAt).toBeDefined(); + expect(held.expiresAt.getTime()).toBeLessThanOrEqual( + held.createdAt.getTime() + 24 * HOUR + 8 * 24 * HOUR, + ); + }); + + it('never re-adds a TTL to a file that was already sent', async () => { + const sentFileId = uuidv4(); + await createFile({ + user: otherUserId, + file_id: sentFileId, + filename: 'sent.png', + filepath: '/uploads/sent.png', + bytes: 10, + type: 'image/png', + }); + await File.updateOne({ file_id: sentFileId }, { $unset: { expiresAt: '' } }); + + const response = await request(app) + .post('/files/usage') + .send({ file_ids: [sentFileId] }); + + expect(response.status).toBe(200); + expect(response.body).toEqual({ held: 0 }); + const permanent = await File.findOne({ file_id: sentFileId }).lean(); + expect(permanent.expiresAt).toBeUndefined(); + }); + + it('never shortens an existing hold', async () => { + const farOut = new Date(Date.now() + 90 * 24 * 60 * 60 * 1000); + const ownFileId = await createQueuedFile(farOut); + + const response = await request(app) + .post('/files/usage') + .send({ file_ids: [ownFileId] }); + + expect(response.status).toBe(200); + expect(response.body).toEqual({ held: 0 }); + const untouched = await File.findOne({ file_id: ownFileId }).lean(); + expect(untouched.expiresAt.getTime()).toBe(farOut.getTime()); }); it("is owner-scoped: another user's file stays untouched (best-effort 200)", async () => { - await File.updateOne({ file_id: fileId }, { $set: { expiresAt: new Date() } }); + const soon = new Date(Date.now() + 60 * 1000); + await File.updateOne({ file_id: fileId }, { $set: { expiresAt: soon } }); const response = await request(app) .post('/files/usage') .send({ file_ids: [fileId] }); expect(response.status).toBe(200); - expect(response.body).toEqual({ marked: 0 }); + expect(response.body).toEqual({ held: 0 }); const untouched = await File.findOne({ file_id: fileId }).lean(); expect(untouched.usage).toBe(0); - expect(untouched.expiresAt).toBeDefined(); + expect(untouched.expiresAt.getTime()).toBe(soon.getTime()); }); it('rejects a list over the cap', async () => { diff --git a/api/server/routes/files/index.js b/api/server/routes/files/index.js index f7e2428c4c..e192ac02ca 100644 --- a/api/server/routes/files/index.js +++ b/api/server/routes/files/index.js @@ -1,5 +1,6 @@ const express = require('express'); const { + createFileUsageLimiter, createFileLimiters, configMiddleware, requireJwtAuth, @@ -30,19 +31,29 @@ const initialize = async () => { router.use('/speech', speech); const { fileUploadIpLimiter, fileUploadUserLimiter } = createFileLimiters(); + const fileUsageLimiter = createFileUsageLimiter(); + + /** Non-strict routing means `/usage/` reaches the same handler, so match the + * route the way Express does. An exact comparison would push a + * trailing-slash request onto the upload quota instead. */ + const isUsagePath = (path) => path.replace(/\/+$/, '') === '/usage'; /** Apply rate limiters to all POST routes (excluding /speech which is handled - * above, and /usage — a metadata touch that must not consume upload quota) */ + * above). `/usage` is a metadata touch, so it gets its own limiter rather + * than consuming upload quota, but it is never left unmetered. */ router.use((req, res, next) => { - if (req.method === 'POST' && !req.path.startsWith('/speech') && req.path !== '/usage') { - return fileUploadIpLimiter(req, res, (err) => { - if (err) { - return next(err); - } - return fileUploadUserLimiter(req, res, next); - }); + if (req.method !== 'POST' || req.path.startsWith('/speech')) { + return next(); } - next(); + if (isUsagePath(req.path)) { + return fileUsageLimiter(req, res, next); + } + return fileUploadIpLimiter(req, res, (err) => { + if (err) { + return next(err); + } + return fileUploadUserLimiter(req, res, next); + }); }); router.post('/', upload.single('file'), restoreTenantContextFromReq); diff --git a/api/server/routes/files/index.limiters.test.js b/api/server/routes/files/index.limiters.test.js new file mode 100644 index 0000000000..0e7d1cf0ec --- /dev/null +++ b/api/server/routes/files/index.limiters.test.js @@ -0,0 +1,129 @@ +const express = require('express'); +const request = require('supertest'); + +const hits = { uploadIp: 0, uploadUser: 0, usage: 0 }; + +jest.mock('~/server/middleware', () => ({ + createFileLimiters: jest.fn(() => ({ + fileUploadIpLimiter: (req, res, next) => { + hits.uploadIp++; + next(); + }, + fileUploadUserLimiter: (req, res, next) => { + hits.uploadUser++; + next(); + }, + })), + createFileUsageLimiter: jest.fn(() => (req, res, next) => { + hits.usage++; + next(); + }), + configMiddleware: (req, res, next) => { + req.config = { fileStrategy: 'local', paths: { uploads: '/tmp', images: '/tmp' } }; + next(); + }, + requireJwtAuth: (req, res, next) => { + req.user = { id: 'user-1', role: 'USER' }; + next(); + }, + uaParser: (req, res, next) => next(), + checkBan: (req, res, next) => next(), +})); + +jest.mock('./multer', () => ({ + createMulterInstance: jest.fn(async () => ({ + single: jest.fn(() => (req, res, next) => next()), + })), +})); + +const okRouter = (paths) => { + const router = express.Router(); + for (const path of paths) { + router.post(path, (req, res) => res.status(200).json({ ok: true })); + } + return router; +}; + +jest.mock('./files', () => { + const express = require('express'); + const router = express.Router(); + router.post('/', (req, res) => res.status(200).json({ ok: true })); + router.post('/usage', (req, res) => res.status(200).json({ ok: true })); + return router; +}); + +jest.mock('./images', () => okRouter(['/'])); +jest.mock('./avatar', () => okRouter(['/'])); +jest.mock('./speech', () => okRouter(['/stt'])); + +jest.mock('~/server/routes/agents/v1', () => ({ + avatar: okRouter(['/:agent_id/avatar/']), +})); +jest.mock('~/server/routes/assistants/v1', () => ({ + avatar: okRouter(['/:assistant_id/avatar/']), +})); + +describe('file route limiter wiring', () => { + let app; + + beforeAll(async () => { + const { initialize } = require('./index'); + app = express(); + app.use(express.json()); + app.use('/api/files', await initialize()); + }); + + beforeEach(() => { + hits.uploadIp = 0; + hits.uploadUser = 0; + hits.usage = 0; + }); + + it('meters POST /usage with the usage limiter, never leaving it unlimited', async () => { + const response = await request(app) + .post('/api/files/usage') + .send({ file_ids: ['f1'] }); + + expect(response.status).toBe(200); + expect(hits.usage).toBe(1); + }); + + it('keeps POST /usage off the upload quota', async () => { + await request(app) + .post('/api/files/usage') + .send({ file_ids: ['f1'] }); + + expect(hits.uploadIp).toBe(0); + expect(hits.uploadUser).toBe(0); + }); + + /* Non-strict routing sends `/usage/` to the same handler, so an exact path + * comparison would bill a trailing-slash client's heartbeats to the upload + * quota and hand them file-upload violations. */ + it('routes a trailing-slash /usage/ through the usage limiter too', async () => { + const response = await request(app) + .post('/api/files/usage/') + .send({ file_ids: ['f1'] }); + + expect(response.status).toBe(200); + expect(hits.usage).toBe(1); + expect(hits.uploadIp).toBe(0); + expect(hits.uploadUser).toBe(0); + }); + + it('still applies the upload limiters to real uploads', async () => { + await request(app).post('/api/files').send({}); + + expect(hits.uploadIp).toBe(1); + expect(hits.uploadUser).toBe(1); + expect(hits.usage).toBe(0); + }); + + it('leaves /speech exempt from both limiters', async () => { + await request(app).post('/api/files/speech/stt').send({}); + + expect(hits.uploadIp).toBe(0); + expect(hits.uploadUser).toBe(0); + expect(hits.usage).toBe(0); + }); +}); diff --git a/api/server/routes/files/index.tenant.test.js b/api/server/routes/files/index.tenant.test.js index 9b703d7524..b2071f06d7 100644 --- a/api/server/routes/files/index.tenant.test.js +++ b/api/server/routes/files/index.tenant.test.js @@ -21,6 +21,7 @@ jest.mock('~/server/middleware', () => ({ fileUploadIpLimiter: (req, res, next) => next(), fileUploadUserLimiter: (req, res, next) => next(), })), + createFileUsageLimiter: jest.fn(() => (req, res, next) => next()), configMiddleware: (req, res, next) => { req.config = { fileStrategy: 'local', diff --git a/client/src/hooks/Chat/__tests__/useQueueDrain.spec.tsx b/client/src/hooks/Chat/__tests__/useQueueDrain.spec.tsx index fce29c79a2..c33b3f979e 100644 --- a/client/src/hooks/Chat/__tests__/useQueueDrain.spec.tsx +++ b/client/src/hooks/Chat/__tests__/useQueueDrain.spec.tsx @@ -1,11 +1,18 @@ import React from 'react'; import { Constants } from 'librechat-data-provider'; import { act, renderHook, waitFor } from '@testing-library/react'; +import { QueryClient, QueryClientProvider } from '@tanstack/react-query'; import { RecoilRoot, useRecoilValue, useSetRecoilState, type MutableSnapshot } from 'recoil'; import type { RunEnd, QueuedMessage } from '~/store/families'; import useQueueDrain from '../useQueueDrain'; import store from '~/store'; +const mockMarkFilesUsage = jest.fn(); +jest.mock('~/data-provider', () => ({ + ...jest.requireActual('~/data-provider'), + useMarkFilesUsageMutation: () => ({ mutate: mockMarkFilesUsage }), +})); + const INDEX = 0; const CONVO_ID = 'convo-drain'; @@ -35,11 +42,18 @@ function setup( return null; } + /* The drain renews the TTL hold on what stays queued, which goes through + * react-query, so the harness needs a client. */ + const queryClient = new QueryClient({ + defaultOptions: { mutations: { retry: false } }, + }); const wrapper = ({ children }: { children: React.ReactNode }) => ( - - - {children} - + + + + {children} + + ); renderHook(() => null, { wrapper }); return { ask, setters }; @@ -57,6 +71,10 @@ const runEnd = (overrides: Partial = {}): RunEnd => ({ }); describe('useQueueDrain', () => { + beforeEach(() => { + mockMarkFilesUsage.mockClear(); + }); + it('drains exactly one queued message on clean completion', async () => { const { ask, setters } = setup(({ set }) => { set(store.queuedMessagesByConvoId(CONVO_ID), [ @@ -117,6 +135,172 @@ describe('useQueueDrain', () => { ); }); + /** Items behind the drained one wait another full run, which may itself + * pause for approval, so a hold taken once at enqueue would lapse before a + * deep queue finishes draining. */ + it('renews the TTL hold on attachments that stay queued', async () => { + const { setters } = setup(({ set }) => { + set(store.queuedMessagesByConvoId(CONVO_ID), [ + { ...queuedMessage('q1', 'first'), files: [{ file_id: 'sent', type: 'image/png' }] }, + { + ...queuedMessage('q2', 'second'), + files: [{ file_id: 'still-queued', type: 'image/png' }], + }, + { ...queuedMessage('q3', 'third'), files: [{ file_id: 'also-queued', type: 'image/png' }] }, + ]); + }); + + mockMarkFilesUsage.mockClear(); + act(() => { + setters.setRunEnd!(runEnd()); + }); + + await waitFor(() => expect(mockMarkFilesUsage).toHaveBeenCalledTimes(1)); + /* Only what remains: the drained item's file is released at send. */ + expect(mockMarkFilesUsage).toHaveBeenCalledWith({ + file_ids: ['still-queued', 'also-queued'], + }); + }); + + it('does not renew when nothing with attachments stays queued', async () => { + const { setters } = setup(({ set }) => { + set(store.queuedMessagesByConvoId(CONVO_ID), [ + { ...queuedMessage('q1', 'only one'), files: [{ file_id: 'sent', type: 'image/png' }] }, + ]); + }); + + mockMarkFilesUsage.mockClear(); + act(() => { + setters.setRunEnd!(runEnd()); + }); + + await waitFor(() => expect(mockMarkFilesUsage).not.toHaveBeenCalled()); + }); + + /** The server caps ids per request, so a queue holding more than one batch + * must renew in several; truncating would leave later messages on their + * enqueue-time hold. */ + it('renews every queued attachment across multiple capped batches', async () => { + const manyFiles = (prefix: string, n: number) => + Array.from({ length: n }, (_, i) => ({ file_id: `${prefix}-${i}`, type: 'image/png' })); + const { setters } = setup(({ set }) => { + set(store.queuedMessagesByConvoId(CONVO_ID), [ + { ...queuedMessage('q1', 'first'), files: manyFiles('sent', 2) }, + { ...queuedMessage('q2', 'second'), files: manyFiles('a', 10) }, + { ...queuedMessage('q3', 'third'), files: manyFiles('b', 4) }, + ]); + }); + + mockMarkFilesUsage.mockClear(); + act(() => { + setters.setRunEnd!(runEnd()); + }); + + await waitFor(() => expect(mockMarkFilesUsage).toHaveBeenCalledTimes(2)); + const sent = mockMarkFilesUsage.mock.calls.flatMap((call) => call[0].file_ids); + expect(sent).toHaveLength(14); + expect(sent).toContain('b-3'); + expect(sent).not.toContain('sent-0'); + for (const call of mockMarkFilesUsage.mock.calls) { + expect(call[0].file_ids.length).toBeLessThanOrEqual(10); + } + }); + + /** A refused send puts the item back with its run-end signal already + * consumed, so nothing else would touch it before the next drain. */ + it('renews a restored item when the send is refused', async () => { + const { ask, setters } = setup(({ set }) => { + set(store.queuedMessagesByConvoId(CONVO_ID), [ + { ...queuedMessage('q1', 'refused'), files: [{ file_id: 'restored', type: 'image/png' }] }, + ]); + }); + ask.mockReturnValue(false); + + mockMarkFilesUsage.mockClear(); + act(() => { + setters.setRunEnd!(runEnd()); + }); + + await waitFor(() => expect(mockMarkFilesUsage).toHaveBeenCalledTimes(1)); + expect(mockMarkFilesUsage).toHaveBeenCalledWith({ file_ids: ['restored'] }); + }); + + /** A single run can pause for approval more than once, so renewal cannot + * depend on catching drain transitions alone. */ + it('renews on a heartbeat while items stay queued', async () => { + jest.useFakeTimers(); + try { + setup(({ set }) => { + set(store.queuedMessagesByConvoId(CONVO_ID), [ + { ...queuedMessage('q1', 'waiting'), files: [{ file_id: 'held', type: 'image/png' }] }, + ]); + }); + + /* Immediate, then on each interval. */ + expect(mockMarkFilesUsage).toHaveBeenCalledTimes(1); + expect(mockMarkFilesUsage).toHaveBeenCalledWith({ file_ids: ['held'] }); + + act(() => { + jest.advanceTimersByTime(30 * 60 * 1000); + }); + expect(mockMarkFilesUsage).toHaveBeenCalledTimes(2); + + act(() => { + jest.advanceTimersByTime(30 * 60 * 1000); + }); + expect(mockMarkFilesUsage).toHaveBeenCalledTimes(3); + } finally { + jest.useRealTimers(); + } + }); + + it('emits no heartbeat when the queue holds no attachments', async () => { + jest.useFakeTimers(); + try { + setup(({ set }) => { + set(store.queuedMessagesByConvoId(CONVO_ID), [queuedMessage('q1', 'no files')]); + }); + + act(() => { + jest.advanceTimersByTime(2 * 60 * 60 * 1000); + }); + expect(mockMarkFilesUsage).not.toHaveBeenCalled(); + } finally { + jest.useRealTimers(); + } + }); + + /** Items queued during the first turn stay keyed under NEW_CONVO until that + * run ends, so renewing only the active id would skip them for its whole + * duration. */ + it('renews the pre-migration NEW_CONVO queue alongside the active one', async () => { + setup(({ set }) => { + set(store.queuedMessagesByConvoId(Constants.NEW_CONVO), [ + { ...queuedMessage('n1', 'queued pre-migration'), files: [{ file_id: 'pending-migrate' }] }, + ]); + set(store.queuedMessagesByConvoId(CONVO_ID), [ + { ...queuedMessage('q1', 'queued after'), files: [{ file_id: 'already-migrated' }] }, + ]); + }); + + await waitFor(() => expect(mockMarkFilesUsage).toHaveBeenCalled()); + const sent = mockMarkFilesUsage.mock.calls.flatMap((call) => call[0].file_ids); + expect(sent).toContain('pending-migrate'); + expect(sent).toContain('already-migrated'); + }); + + it('does not double-count the queue before migration', async () => { + setup(({ set }) => { + set(store.queuedMessagesByConvoId(Constants.NEW_CONVO), [ + { ...queuedMessage('n1', 'new convo'), files: [{ file_id: 'only-once' }] }, + ]); + }, Constants.NEW_CONVO as string); + + await waitFor(() => expect(mockMarkFilesUsage).toHaveBeenCalled()); + const sent = mockMarkFilesUsage.mock.calls.flatMap((call) => call[0].file_ids); + expect(sent).toEqual(['only-once']); + }); + it('passes carried quotes + manual skills through as overrides', async () => { const { ask, setters } = setup(({ set }) => { set(store.queuedMessagesByConvoId(CONVO_ID), [ @@ -235,11 +419,16 @@ describe('useQueueDrain', () => { return null; } + const queryClient = new QueryClient({ + defaultOptions: { mutations: { retry: false } }, + }); const wrapper = ({ children }: { children: React.ReactNode }) => ( - - - {children} - + + + + {children} + + ); renderHook(() => null, { wrapper }); return { ask, setters, state }; diff --git a/client/src/hooks/Chat/useQueueDrain.ts b/client/src/hooks/Chat/useQueueDrain.ts index 588038012e..e4092304a9 100644 --- a/client/src/hooks/Chat/useQueueDrain.ts +++ b/client/src/hooks/Chat/useQueueDrain.ts @@ -1,10 +1,41 @@ -import { useEffect } from 'react'; +import { useEffect, useMemo } from 'react'; import { Constants } from 'librechat-data-provider'; import { useRecoilValue, useRecoilCallback } from 'recoil'; import type { QueuedMessage } from '~/store/families'; import type { TAskFunction } from '~/common'; +import { useMarkFilesUsageMutation } from '~/data-provider'; import store from '~/store'; +/** Mirrors the server's per-request cap on a usage touch. */ +const QUEUE_USAGE_MAX_FILES = 10; + +/** Well under the server's smallest hold (24h), so no gap between renewals + * can outlive one, however long a run pauses. */ +const QUEUE_USAGE_RENEW_INTERVAL_MS = 30 * 60 * 1000; + +const collectQueuedFileIds = (items: QueuedMessage[]): string[] => { + const fileIds: string[] = []; + for (const item of items) { + for (const file of item.files ?? []) { + if (typeof file.file_id === 'string' && file.file_id.length > 0) { + fileIds.push(file.file_id); + } + } + } + return fileIds; +}; + +/** The server caps ids per request, so a queue holding more than one batch + * has to renew in several. Truncating instead would leave everything after + * the first batch on its enqueue-time hold. */ +const batchFileIds = (fileIds: string[]): string[][] => { + const batches: string[][] = []; + for (let i = 0; i < fileIds.length; i += QUEUE_USAGE_MAX_FILES) { + batches.push(fileIds.slice(i, i + QUEUE_USAGE_MAX_FILES)); + } + return batches; +}; + /** * Auto-sends queued follow-up messages when a run finishes. * @@ -30,6 +61,50 @@ export default function useQueueDrain( store.pendingRunEndByConvoId(activeConversationId ?? Constants.NEW_CONVO), ); const isSubmitting = useRecoilValue(store.isSubmittingFamily(index)); + const { mutate: markFilesUsage } = useMarkFilesUsageMutation(); + const ownQueue = useRecoilValue( + store.queuedMessagesByConvoId(activeConversationId ?? Constants.NEW_CONVO), + ); + /** `drainNext` merges this in, and it outlives the URL update: items queued + * during the first turn stay keyed here until that run ends. Renewing only + * the active id would skip them for the whole of that run. */ + const newConvoQueue = useRecoilValue(store.queuedMessagesByConvoId(Constants.NEW_CONVO)); + + /* Deduped because the two subscriptions are the same atom before migration. + * Keyed by id list so the effect re-runs when the held set changes, not + * whenever Recoil hands back a new array for the same contents. */ + const renewKey = [ + ...new Set([...collectQueuedFileIds(ownQueue), ...collectQueuedFileIds(newConvoQueue)]), + ].join(','); + const queuedFileIds = useMemo(() => (renewKey ? renewKey.split(',') : []), [renewKey]); + + /** + * Heartbeat renewal while anything is queued. + * + * Renewing only at drain transitions ties an attachment's survival to + * catching every state change, and a single run can stretch well past one + * hold: it may interrupt for approval more than once, and each pause can + * run to the configured window. Rather than hook every transition, renew on + * a cadence far shorter than the hold itself, so no single gap can outlive + * it. Bounded regardless: the server clamps each renewal against the file's + * upload time. A queue nobody has open stops emitting these and lapses + * normally. + */ + useEffect(() => { + if (queuedFileIds.length === 0) { + return; + } + const renew = () => { + for (const file_ids of batchFileIds(queuedFileIds)) { + markFilesUsage({ file_ids }); + } + }; + /* Immediately, not one interval later: returning to a conversation whose + * hold is nearly up would otherwise wait out a full period first. */ + renew(); + const timer = setInterval(renew, QUEUE_USAGE_RENEW_INTERVAL_MS); + return () => clearInterval(timer); + }, [queuedFileIds, markFilesUsage]); // Fully synchronous reads (getLoadable): a useRecoilCallback snapshot is // only guaranteed valid for the callback's synchronous execution, so no @@ -153,9 +228,25 @@ export default function useQueueDrain( if (accepted === false) { // `ask` refused without sending (e.g. the conversation history is not // in the query cache yet, right after navigating back). Restore the - // item so the user's text is never silently dropped — the chip stays + // item so the user's text is never silently dropped, the chip stays // available for manual send. restoreQueued(conversationId, next); + /** Popping and restoring leaves the held set identical, so the renewal + * effect sees no change and will not re-run. Draining normally does + * change the set, and renews itself. Fire-and-forget: send-time + * marking is the backstop. */ + for (const file_ids of batchFileIds(collectQueuedFileIds([next]))) { + markFilesUsage({ file_ids }); + } } - }, [runEnd, parkedRunEnd, isSubmitting, activeConversationId, drainNext, restoreQueued, ask]); + }, [ + runEnd, + parkedRunEnd, + isSubmitting, + activeConversationId, + drainNext, + restoreQueued, + markFilesUsage, + ask, + ]); } diff --git a/client/src/hooks/Chat/useSteering.ts b/client/src/hooks/Chat/useSteering.ts index 2e5467464b..a0ad175c5a 100644 --- a/client/src/hooks/Chat/useSteering.ts +++ b/client/src/hooks/Chat/useSteering.ts @@ -251,10 +251,11 @@ export default function useSteering({ [], ); - /** Fire-and-forget TTL touch for uploads entering the client queue: a + /** Fire-and-forget TTL hold for uploads entering the client queue: a * queued message can outlive the upload window (long run, approval pause) - * and send-time marking only happens at drain. Failure is tolerated — - * the send-time marking remains the backstop. */ + * and send-time marking only happens at drain. The server extends the + * deadline rather than clearing it, so a queue this tab never drains still + * gets reaped. Failure is tolerated; send-time marking is the backstop. */ const markQueuedFilesUsage = useCallback( (files?: TMessage['files']) => { if (files == null || files.length === 0) { @@ -285,8 +286,8 @@ export default function useSteering({ files?: TMessage['files']; quotes?: string[]; manualSkills?: string[]; - /** Set when the files were ALREADY queued/steered — their usage was - * marked when they first entered the queue (or at the steer 202). */ + /** Set when the files were ALREADY queued/steered: their TTL was + * held when they first entered the queue (or at the steer 202). */ skipUsageMark?: boolean; }, ) => { @@ -694,7 +695,7 @@ export default function useSteering({ } return; } - // Files were already marked used when the item first entered the queue. + // Files were already TTL-held when the item first entered the queue. enqueue(taken.text, { front: true, files: taken.files, diff --git a/packages/api/src/files/usage.spec.ts b/packages/api/src/files/usage.spec.ts index 59de9cbffd..102f235cfb 100644 --- a/packages/api/src/files/usage.spec.ts +++ b/packages/api/src/files/usage.spec.ts @@ -1,17 +1,23 @@ -import { handleFilesUsageRequest, FILES_USAGE_MAX_IDS } from './usage'; +import { + handleFilesUsageRequest, + resolveFilesUsageHold, + FILES_USAGE_QUEUED_RUN_ALLOWANCE, + FILES_USAGE_BASE_HOLD_MS, + FILES_USAGE_MAX_IDS, +} from './usage'; describe('handleFilesUsageRequest', () => { const user = { id: 'user-1', tenantId: 'tenant-1' }; - const createDeps = (marked: unknown[] = []) => ({ - updateFilesUsage: jest.fn().mockResolvedValue(marked), + const createDeps = (held = 0) => ({ + extendFilesTTL: jest.fn().mockResolvedValue(held), }); it('rejects unauthenticated requests without touching the DB', async () => { const deps = createDeps(); const result = await handleFilesUsageRequest({}, { file_ids: ['f1'] }, deps); expect(result).toEqual({ status: 401, body: { code: 'UNAUTHORIZED' } }); - expect(deps.updateFilesUsage).not.toHaveBeenCalled(); + expect(deps.extendFilesTTL).not.toHaveBeenCalled(); }); it.each([ @@ -25,7 +31,7 @@ describe('handleFilesUsageRequest', () => { const result = await handleFilesUsageRequest(user, body, deps); expect(result.status).toBe(400); expect(result.body).toEqual({ code: 'INVALID_FILE_IDS' }); - expect(deps.updateFilesUsage).not.toHaveBeenCalled(); + expect(deps.extendFilesTTL).not.toHaveBeenCalled(); }); it('caps the list at FILES_USAGE_MAX_IDS', async () => { @@ -34,24 +40,80 @@ describe('handleFilesUsageRequest', () => { const result = await handleFilesUsageRequest(user, { file_ids }, deps); expect(result.status).toBe(400); expect(result.body).toEqual({ code: 'TOO_MANY_FILES', max: FILES_USAGE_MAX_IDS }); - expect(deps.updateFilesUsage).not.toHaveBeenCalled(); + expect(deps.extendFilesTTL).not.toHaveBeenCalled(); }); - it('marks usage owner-scoped and returns best-effort 200', async () => { - const deps = createDeps([{ file_id: 'f1' }]); + it('requests an owner-scoped hold of the resolved window', async () => { + const deps = createDeps(1); const result = await handleFilesUsageRequest(user, { file_ids: ['f1', 'f2'] }, deps); - expect(deps.updateFilesUsage).toHaveBeenCalledTimes(1); - expect(deps.updateFilesUsage).toHaveBeenCalledWith( - [{ file_id: 'f1' }, { file_id: 'f2' }], - undefined, + expect(deps.extendFilesTTL).toHaveBeenCalledTimes(1); + expect(deps.extendFilesTTL).toHaveBeenCalledWith( + ['f1', 'f2'], + { renewMs: FILES_USAGE_BASE_HOLD_MS, maxLifetimeMs: FILES_USAGE_BASE_HOLD_MS }, { user: 'user-1', tenantId: 'tenant-1' }, ); - expect(result).toEqual({ status: 200, body: { marked: 1 } }); + expect(result).toEqual({ status: 200, body: { held: 1 } }); }); - it('returns 200 with zero marked when no id resolves to an owned file', async () => { - const deps = createDeps([]); + /** The ceiling must be a constant the data layer clamps against the upload + * time, not a deadline derived here: a request-clock deadline with no + * ceiling would let a caller walk the file's lifetime forward per call. */ + it('passes a constant ceiling, never a request-derived deadline', async () => { + const deps = createDeps(1); + await handleFilesUsageRequest(user, { file_ids: ['f1'] }, deps); + await handleFilesUsageRequest(user, { file_ids: ['f1'] }, deps); + + const [, firstHold] = deps.extendFilesTTL.mock.calls[0]; + const [, secondHold] = deps.extendFilesTTL.mock.calls[1]; + expect(secondHold).toEqual(firstHold); + expect(typeof firstHold.maxLifetimeMs).toBe('number'); + }); + + /** A paused run's approval window has no configured upper bound, so a fixed + * window shorter than it would reap an attachment whose approval is still + * live and leave the drain sending a missing file. */ + it('stretches the hold to outlast a long configured approval window', async () => { + const sevenDaysMs = 7 * 24 * 60 * 60 * 1000; + const deps = { ...createDeps(1), approvalTtlMs: sevenDaysMs }; + await handleFilesUsageRequest(user, { file_ids: ['f1'] }, deps); + + const [, hold] = deps.extendFilesTTL.mock.calls[0]; + expect(hold.renewMs).toBe(FILES_USAGE_BASE_HOLD_MS + sevenDaysMs); + expect(hold.renewMs).toBeGreaterThan(sevenDaysMs); + }); + + describe('resolveFilesUsageHold', () => { + it('covers one approval window per renewal', () => { + expect(resolveFilesUsageHold(1000).renewMs).toBe(FILES_USAGE_BASE_HOLD_MS + 1000); + }); + + /** The drain sends one queued item per run completion, so a deep queue + * waits behind several approval windows; the ceiling has to span them. */ + it('scales the ceiling to the queued run chain', () => { + const hold = resolveFilesUsageHold(1000); + expect(hold.maxLifetimeMs).toBe( + FILES_USAGE_BASE_HOLD_MS + 1000 * FILES_USAGE_QUEUED_RUN_ALLOWANCE, + ); + expect(hold.maxLifetimeMs).toBeGreaterThan(hold.renewMs); + }); + + it.each([ + ['undefined', undefined], + ['zero', 0], + ['negative', -1], + ['NaN', Number.NaN], + ['Infinity', Number.POSITIVE_INFINITY], + ])('falls back to the baseline for %s', (_label, value) => { + expect(resolveFilesUsageHold(value)).toEqual({ + renewMs: FILES_USAGE_BASE_HOLD_MS, + maxLifetimeMs: FILES_USAGE_BASE_HOLD_MS, + }); + }); + }); + + it('returns 200 with zero held when no id resolves to an owned file', async () => { + const deps = createDeps(0); const result = await handleFilesUsageRequest(user, { file_ids: ['not-owned'] }, deps); - expect(result).toEqual({ status: 200, body: { marked: 0 } }); + expect(result).toEqual({ status: 200, body: { held: 0 } }); }); }); diff --git a/packages/api/src/files/usage.ts b/packages/api/src/files/usage.ts index fdb20127ed..149e6e1ddb 100644 --- a/packages/api/src/files/usage.ts +++ b/packages/api/src/files/usage.ts @@ -1,6 +1,53 @@ /** Cap per usage touch, mirroring the composer's practical attachment limit. */ export const FILES_USAGE_MAX_IDS: number = 10; +/** + * Baseline lifetime a held upload gets, measured from its upload time. Covers + * the span a queued attachment spends before a run can even pause: upload, + * enqueue, and the run itself reaching an approval point. + */ +export const FILES_USAGE_BASE_HOLD_MS: number = 24 * 60 * 60 * 1000; + +/** + * Queued runs a held attachment is assumed to wait behind. The queue is + * unbounded, so no multiple is provably sufficient; this is the depth the + * ceiling covers before a still-queued attachment can lapse. + */ +export const FILES_USAGE_QUEUED_RUN_ALLOWANCE: number = 8; + +export interface FilesUsageHold { + /** Granted from now, so an actively draining queue keeps renewing. */ + renewMs: number; + /** Hard ceiling from upload time, so renewals converge instead of walking. */ + maxLifetimeMs: number; +} + +/** + * Resolves the hold window for one touch. + * + * A queued message legitimately outlives the upload window when a run pauses + * for approval, and that pause is bounded by the deployment's configured + * approval window (`endpoints.agents.checkpointer.ttl`), which has no upper + * limit. Deeper queues wait behind several such runs, since the drain sends + * one item per run completion. + * + * Hence two numbers rather than one. `renewMs` covers a single run's wait and + * is granted from now, so a queue that is still draining re-asserts it at each + * transition. `maxLifetimeMs` is measured from the immutable upload time and + * caps every renewal, so repeated touches converge on a fixed ceiling instead + * of advancing a window per call. An abandoned queue therefore lapses one + * `renewMs` after its last touch rather than surviving to the ceiling. + * + * @param approvalTtlMs - Configured approval window; falsy means none configured + */ +export function resolveFilesUsageHold(approvalTtlMs?: number): FilesUsageHold { + const approval = Number.isFinite(approvalTtlMs) && approvalTtlMs! > 0 ? approvalTtlMs! : 0; + return { + renewMs: FILES_USAGE_BASE_HOLD_MS + approval, + maxLifetimeMs: FILES_USAGE_BASE_HOLD_MS + approval * FILES_USAGE_QUEUED_RUN_ALLOWANCE, + }; +} + export interface FilesUsageUser { id?: string; tenantId?: string; @@ -17,20 +64,35 @@ export interface FilesUsageResult { } export interface FilesUsageDeps { - /** Owner-scoped usage marker (`db.updateFilesUsage`-shaped). */ - updateFilesUsage: ( - files: Array<{ file_id: string }>, - fileIds?: string[], - options?: { user?: string; tenantId?: string | null }, - ) => Promise; + /** Owner-scoped TTL hold (`db.extendFilesTTL`-shaped). */ + extendFilesTTL: ( + fileIds: string[], + hold: FilesUsageHold, + owner: { user: string; tenantId?: string | null }, + ) => Promise; + /** Configured approval window, so the hold outlasts a paused run. */ + approvalTtlMs?: number; } /** - * Owner-scoped usage touch for attachments entering a client-side queue: a - * queued message can outlive the upload window (long run, approval pause), so - * marking at queue time stops the TTL from reaping files the drain will send. - * Best-effort 200 — ids that do not resolve to an owned file are not errors - * (send-time marking remains the backstop). + * Owner-scoped TTL hold for attachments entering a client-side queue: a + * queued message can outlive the upload window (long run, approval pause), + * so holding at queue time stops the TTL from reaping files the drain will + * send. + * + * A hold, not a release. The client queue is ephemeral browser state, so a + * closed tab or a cleared queue leaves nothing referencing these files, and + * clearing the TTL outright would strand them in storage permanently. The + * hold only ever widens, and every renewal is capped against the file's + * upload time, so replaying it converges on a ceiling instead of advancing + * indefinitely; the real release happens at send, where `updateFilesUsage` + * marks the files used against an actual message. + * + * The window tracks the deployment's configured approval TTL so a queue + * waiting on a paused run cannot be reaped while that approval is live. + * + * Best-effort 200: ids that do not resolve to a held file are not errors + * (they may be already-sent files, or not owned). */ export async function handleFilesUsageRequest( user: FilesUsageUser, @@ -47,16 +109,16 @@ export async function handleFilesUsageRequest( if (raw.length > FILES_USAGE_MAX_IDS) { return { status: 400, body: { code: 'TOO_MANY_FILES', max: FILES_USAGE_MAX_IDS } }; } - const files: Array<{ file_id: string }> = []; + const fileIds: string[] = []; for (const value of raw) { if (typeof value !== 'string' || value.length === 0) { return { status: 400, body: { code: 'INVALID_FILE_IDS' } }; } - files.push({ file_id: value }); + fileIds.push(value); } - const marked = await deps.updateFilesUsage(files, undefined, { + const held = await deps.extendFilesTTL(fileIds, resolveFilesUsageHold(deps.approvalTtlMs), { user: user.id, tenantId: user.tenantId, }); - return { status: 200, body: { marked: marked.length } }; + return { status: 200, body: { held } }; } diff --git a/packages/data-provider/src/types/files.ts b/packages/data-provider/src/types/files.ts index 4d1210cd7d..b889ca6f17 100644 --- a/packages/data-provider/src/types/files.ts +++ b/packages/data-provider/src/types/files.ts @@ -242,7 +242,8 @@ export type TFilesUsageBody = { }; export type TFilesUsageResponse = { - marked: number; + /** Count of queued uploads whose TTL hold was extended. */ + held: number; }; export type DeleteFilesResponse = { diff --git a/packages/data-schemas/src/methods/file.spec.ts b/packages/data-schemas/src/methods/file.spec.ts index 71b12ff5aa..4aa424301b 100644 --- a/packages/data-schemas/src/methods/file.spec.ts +++ b/packages/data-schemas/src/methods/file.spec.ts @@ -1123,6 +1123,167 @@ describe('File Methods', () => { }); }); + describe('extendFilesTTL', () => { + const HOUR = 3_600_000; + const HOLD = { renewMs: 24 * HOUR, maxLifetimeMs: 48 * HOUR }; + + const seedTempFile = async ( + userId: mongoose.Types.ObjectId, + expiresAt: Date, + createdAt?: Date, + ) => { + const fileId = uuidv4(); + await fileMethods.createFile({ + file_id: fileId, + user: userId, + filename: `${fileId}.txt`, + filepath: `/uploads/${fileId}.txt`, + type: 'text/plain', + bytes: 100, + }); + await mongoose.models.File.updateOne( + { file_id: fileId }, + { $set: { expiresAt, ...(createdAt ? { createdAt } : {}) } }, + { timestamps: false }, + ); + return fileId; + }; + + const readFile = async (fileId: string) => + (await mongoose.models.File.findOne({ file_id: fileId }) + .lean<{ createdAt: Date; expiresAt?: Date }>() + .exec())!; + + it('widens the TTL toward the renewal window without unsetting it', async () => { + const userId = new mongoose.Types.ObjectId(); + const fileId = await seedTempFile(userId, new Date(Date.now() + 60_000)); + + const count = await fileMethods.extendFilesTTL([fileId], HOLD, { user: String(userId) }); + + expect(count).toBe(1); + const file = await readFile(fileId); + expect(file.expiresAt).toBeDefined(); + expect(file.expiresAt!.getTime()).toBeGreaterThan(Date.now() + 23 * HOUR); + expect(file.expiresAt!.getTime()).toBeLessThanOrEqual( + file.createdAt.getTime() + HOLD.maxLifetimeMs, + ); + }); + + /** The bound that makes the endpoint safe to expose: every renewal is + * clamped to createdAt + maxLifetimeMs, so replaying the call converges + * on a ceiling instead of advancing a window at a time. */ + it('clamps every renewal to the ceiling, however often it is replayed', async () => { + const userId = new mongoose.Types.ObjectId(); + const fileId = await seedTempFile(userId, new Date(Date.now() + 60_000)); + const ceilingHold = { renewMs: 24 * HOUR, maxLifetimeMs: 10 * 60_000 }; + const { createdAt } = await readFile(fileId); + + const first = await fileMethods.extendFilesTTL([fileId], ceilingHold, { + user: String(userId), + }); + expect(first).toBe(1); + + for (let i = 0; i < 5; i++) { + expect( + await fileMethods.extendFilesTTL([fileId], ceilingHold, { user: String(userId) }), + ).toBe(0); + } + + const file = await readFile(fileId); + expect(file.expiresAt!.getTime()).toBe(createdAt.getTime() + ceilingHold.maxLifetimeMs); + /* The renewal window alone would have granted a full day. */ + expect(file.expiresAt!.getTime()).toBeLessThan(Date.now() + HOUR); + }); + + /** Renewal is measured from now, so a queue still draining across + * successive runs keeps its attachments past the first hold. */ + it('renews from now while the ceiling still allows it', async () => { + const userId = new mongoose.Types.ObjectId(); + const twoHoursAgo = new Date(Date.now() - 2 * HOUR); + const fileId = await seedTempFile(userId, new Date(Date.now() + 60_000), twoHoursAgo); + + const count = await fileMethods.extendFilesTTL( + [fileId], + { renewMs: HOUR, maxLifetimeMs: 24 * HOUR }, + { user: String(userId) }, + ); + + expect(count).toBe(1); + const file = await readFile(fileId); + /* Anchoring to createdAt alone would have expired this an hour ago. */ + expect(file.expiresAt!.getTime()).toBeGreaterThan(Date.now() + 50 * 60_000); + }); + + it('does not resurrect a TTL on an already-released file', async () => { + const userId = new mongoose.Types.ObjectId(); + const fileId = await seedTempFile(userId, new Date(Date.now() + 60_000)); + await fileMethods.updateFileUsage({ file_id: fileId, user: String(userId) }); + + const count = await fileMethods.extendFilesTTL([fileId], HOLD, { user: String(userId) }); + + expect(count).toBe(0); + const file = await readFile(fileId); + expect(file.expiresAt).toBeUndefined(); + }); + + it('never moves an expiry earlier', async () => { + const userId = new mongoose.Types.ObjectId(); + const farOut = new Date(Date.now() + 7 * 24 * HOUR); + const fileId = await seedTempFile(userId, farOut); + + const count = await fileMethods.extendFilesTTL([fileId], HOLD, { user: String(userId) }); + + expect(count).toBe(0); + const file = await readFile(fileId); + expect(file.expiresAt!.getTime()).toBe(farOut.getTime()); + }); + + it("leaves another user's file untouched", async () => { + const ownerId = new mongoose.Types.ObjectId(); + const attackerId = new mongoose.Types.ObjectId(); + const soon = new Date(Date.now() + 60_000); + const fileId = await seedTempFile(ownerId, soon); + + const count = await fileMethods.extendFilesTTL([fileId], HOLD, { + user: String(attackerId), + }); + + expect(count).toBe(0); + const file = await readFile(fileId); + expect(file.expiresAt!.getTime()).toBe(soon.getTime()); + }); + + it('is a no-op without an owner scope rather than a cross-user update', async () => { + const userId = new mongoose.Types.ObjectId(); + const soon = new Date(Date.now() + 60_000); + const fileId = await seedTempFile(userId, soon); + + const count = await fileMethods.extendFilesTTL([fileId], HOLD, { user: '' } as { + user: string; + }); + + expect(count).toBe(0); + const file = await readFile(fileId); + expect(file.expiresAt!.getTime()).toBe(soon.getTime()); + }); + + it('is a no-op for a non-positive hold', async () => { + const userId = new mongoose.Types.ObjectId(); + const soon = new Date(Date.now() + 60_000); + const fileId = await seedTempFile(userId, soon); + const owner = { user: String(userId) }; + + expect( + await fileMethods.extendFilesTTL([fileId], { renewMs: 0, maxLifetimeMs: 0 }, owner), + ).toBe(0); + expect( + await fileMethods.extendFilesTTL([fileId], { renewMs: -1, maxLifetimeMs: -1 }, owner), + ).toBe(0); + const file = await readFile(fileId); + expect(file.expiresAt!.getTime()).toBe(soon.getTime()); + }); + }); + describe('deleteFile', () => { it('should delete a file by file_id', async () => { const fileId = uuidv4(); diff --git a/packages/data-schemas/src/methods/file.ts b/packages/data-schemas/src/methods/file.ts index 2e1d9159a5..ac428bdce0 100644 --- a/packages/data-schemas/src/methods/file.ts +++ b/packages/data-schemas/src/methods/file.ts @@ -82,6 +82,11 @@ export function createFileMethods(mongoose: typeof import('mongoose')): { fileIds?: string[], options?: { user?: string; tenantId?: string | null }, ) => Promise; + extendFilesTTL: ( + fileIds: string[], + hold: { renewMs: number; maxLifetimeMs: number }, + owner: { user: string; tenantId?: string | null }, + ) => Promise; sweepOrphanedPreviews: (maxAgeMs?: number) => Promise; } { /** @@ -546,6 +551,78 @@ export function createFileMethods(mongoose: typeof import('mongoose')): { return results.filter((result): result is IMongoFile => result != null); } + /** + * Widens the upload-window TTL of owned, still-temporary files to + * `min(now + renewMs, createdAt + maxLifetimeMs)`. + * + * A renewable hold, not a release: unlike `updateFileUsage` this never + * unsets `expiresAt`, so a file that is held but never actually sent is + * still reaped once the hold lapses. Four properties hold by + * construction, which is what makes the write safe to drive from a + * client-supplied id list: + * - `$min` against `createdAt + maxLifetimeMs` caps every renewal against + * an immutable anchor, so repeated calls converge on a fixed ceiling + * instead of walking a file's lifetime forward a window at a time; + * - renewing from `now` up to that ceiling lets a queue that is still + * draining keep its attachments alive across successive runs, while an + * abandoned queue lapses a single `renewMs` after its last touch rather + * than surviving to the ceiling; + * - `$max` against the current value means a hold only ever widens; + * - `expiresAt: { $exists: true }` means a file whose TTL was already + * cleared by a real send stays permanent. Re-adding `expiresAt` there + * would schedule a live file for deletion. + * + * `createdAt` is required rather than defaulted: without the anchor there + * is no ceiling to enforce, so such a file is skipped instead of held. + * + * The owner scope is required, not optional: an unscoped call would hold + * every user's matching file. A missing owner is a no-op, not a wide + * update. + * + * @param fileIds - File IDs to hold + * @param hold - `renewMs` granted from now, capped at `maxLifetimeMs` from upload + * @param owner - Owner scope; mismatches leave the TTL unchanged + * @returns Number of files whose hold was widened + */ + async function extendFilesTTL( + fileIds: string[], + hold: { renewMs: number; maxLifetimeMs: number }, + owner: { user: string; tenantId?: string | null }, + ): Promise { + const renewMs = hold?.renewMs; + const maxLifetimeMs = hold?.maxLifetimeMs; + if (fileIds.length === 0 || !owner?.user || !(renewMs > 0) || !(maxLifetimeMs > 0)) { + return 0; + } + const File = mongoose.models.File as Model; + const filter = withOwnerScope( + { + file_id: { $in: [...new Set(fileIds)] }, + expiresAt: { $exists: true }, + createdAt: { $exists: true }, + }, + { userId: owner.user, tenantId: owner.tenantId }, + ); + const renewUntil = new Date(Date.now() + renewMs); + const result = await File.updateMany( + filter, + [ + { + $set: { + expiresAt: { + $max: ['$expiresAt', { $min: [renewUntil, { $add: ['$createdAt', maxLifetimeMs] }] }], + }, + }, + }, + ], + /** `timestamps: false`: a hold is TTL bookkeeping, not a content write. + * Bumping `updatedAt` would also make every re-touch count as a + * modification, hiding whether the deadline actually moved. */ + { timestamps: false }, + ); + return result.modifiedCount ?? 0; + } + /** * Mark stale `status: 'pending'` file records as `'failed'` with * `previewError: 'orphaned'`. Recovers from the one case the @@ -594,6 +671,7 @@ export function createFileMethods(mongoose: typeof import('mongoose')): { deleteFileByFilter, batchUpdateFiles, updateFilesUsage, + extendFilesTTL, sweepOrphanedPreviews, }; }