diff --git a/.env.example b/.env.example index 7091d3f9ad0..22b42658820 100644 --- a/.env.example +++ b/.env.example @@ -1059,6 +1059,11 @@ HELP_AND_FAQ_URL=https://librechat.ai # Enable Redis for resumable LLM streams (defaults to USE_REDIS value if not set) # Set to false to use in-memory storage for streams while keeping Redis for other caches # USE_REDIS_STREAMS=true +# Scheduled chats require shared Redis streams in multi-replica deployments. +# Set this only when the deployment truly runs one LibreChat process without Redis. +# SCHEDULES_SINGLE_PROCESS=true +# Emergency global stop for both automatic and manual scheduled runs. +# SCHEDULES_DISABLED=true # Generation stream wire/state protocol. Redis-backed deployments default to the # rolling-upgrade-safe v1 protocol when this is unset; in-memory deployments use v2. diff --git a/api/package.json b/api/package.json index c625b46fca2..58810c14875 100644 --- a/api/package.json +++ b/api/package.json @@ -46,7 +46,7 @@ "@azure/storage-blob": "^12.30.0", "@google/genai": "^2.8.0", "@keyv/redis": "^4.3.3", - "@librechat/agents": "^3.6.8", + "@librechat/agents": "^3.6.9", "@librechat/api": "*", "@librechat/data-schemas": "*", "@microsoft/microsoft-graph-client": "^3.0.7", diff --git a/api/server/controllers/UserController.js b/api/server/controllers/UserController.js index 9821f47db77..040b2ae21b2 100644 --- a/api/server/controllers/UserController.js +++ b/api/server/controllers/UserController.js @@ -33,6 +33,11 @@ const { purgeAgentTriggerDeliveriesForUser, } = require('~/server/services/Agents/triggers'); const { getAppConfig } = require('~/server/services/Config'); +const { randomUUID } = require('node:crypto'); +const { + quiesceUserSchedules, + restoreUserSchedulesFromDeletion, +} = require('~/server/services/Schedules'); const { getLogStores } = require('~/cache'); const db = require('~/models'); @@ -360,6 +365,7 @@ const updateUserPluginsController = async (req, res) => { const deleteUserController = async (req, res) => { const { user } = req; let triggerDeletionFence; + let scheduleSuspensionToken; let userDeleted = false; try { @@ -395,6 +401,14 @@ const deleteUserController = async (req, res) => { } await drainAgentTriggerDeliveriesForUser(user.id); await subagentThreadTaskStore.cancelAndDrainForOwner(user.id, user.tenantId); + // Reversibly suspend the user's schedules under a per-attempt token BEFORE draining. + // A later cascade step (or this drain) can still fail and cancel the deletion, and the + // catch below restores exactly this attempt's rows — so a failed deletion never leaves + // a live user with silently disabled, erasure-eligible schedules. + scheduleSuspensionToken = randomUUID(); + if (!(await quiesceUserSchedules(user.id, scheduleSuspensionToken))) { + throw new Error('Scheduled executions could not be confirmed stopped'); + } const activeAgentRuns = await GenerationJobManager.getCleanupBlockingJobIdsForUser( user.id, user.tenantId, @@ -446,6 +460,7 @@ const deleteUserController = async (req, res) => { await db.deleteTokens({ userId: user.id }); await db.removeUserFromAllGroups(user.id); await db.deleteAclEntries({ principalId: user._id }); + await db.deleteSchedulesByUser(user.id); const deleteResult = await db.deleteUserById(user.id); if (deleteResult.deletedCount !== 1) { throw new Error('User disappeared before account deletion could commit'); @@ -455,6 +470,33 @@ const deleteUserController = async (req, res) => { logger.info(`User deleted account. Email: ${user.email} ID: ${user.id}`); res.status(200).send({ message: 'User deleted' }); } catch (err) { + // The account survives this failed attempt, so its schedules must too: restore the + // exact rows this attempt suspended (re-enabled/re-armed from their snapshot). Fenced + // to the token, so a schedule the owner deleted meanwhile is not resurrected. A + // successful deletion never reaches here (userDeleted short-circuits it). + // + // RESTORE BEFORE RELEASING THE DELETION FENCE. That fence is what refuses new schedule + // writes/claims for this user; releasing it first opens a window where an owner PATCH + // could edit a still-suspended row and then have its enabled/next-run state overwritten + // by this older snapshot, and where a second deletion attempt could re-suspend these + // rows under a new token — making this restore a no-op and stranding the disabled + // snapshot permanently. + if (scheduleSuspensionToken != null && !userDeleted) { + try { + await restoreUserSchedulesFromDeletion(user.id, scheduleSuspensionToken); + } catch (restoreError) { + // Every retry is exhausted at this point. The fence is still released below on + // purpose: retaining it would refuse this live account's schedule writes AND make + // `beginAgentTriggerUserDeletion` report `in_progress` forever, blocking the retry + // that is the convergence path — a later attempt re-suspends by ADOPTING this + // snapshot, so its cancel restores these exact rows. Log the token so the state is + // recoverable directly if that never happens. + logger.error( + `[deleteUserController] Failed to restore suspended schedules after a cancelled deletion; they remain disabled for user ${user.id} under suspension token ${scheduleSuspensionToken}`, + restoreError, + ); + } + } if (triggerDeletionFence != null && !userDeleted) { try { await cancelAgentTriggerUserPurge(user.id, triggerDeletionFence); diff --git a/api/server/controllers/UserController.spec.js b/api/server/controllers/UserController.spec.js index c9a3f6eb447..89af705a563 100644 --- a/api/server/controllers/UserController.spec.js +++ b/api/server/controllers/UserController.spec.js @@ -8,6 +8,8 @@ const mockPrepareAgentTriggerUserPurge = jest.fn().mockResolvedValue(undefined); const mockCancelAgentTriggerUserPurge = jest.fn().mockResolvedValue(true); const mockPurgeAgentTriggerDeliveriesForUser = jest.fn().mockResolvedValue(undefined); const mockCancelAndDrainSubagentThreads = jest.fn().mockResolvedValue(undefined); +const mockQuiesceUserSchedules = jest.fn().mockResolvedValue(true); +const mockRestoreUserSchedules = jest.fn().mockResolvedValue(undefined); jest.mock('@librechat/data-schemas', () => { const actual = jest.requireActual('@librechat/data-schemas'); @@ -30,6 +32,7 @@ jest.mock('~/models', () => { deleteAllAgentApiKeys: jest.fn().mockResolvedValue(undefined), deleteConversationTags: jest.fn().mockResolvedValue(undefined), deleteAllUserMemories: jest.fn().mockResolvedValue(undefined), + deleteSchedulesByUser: jest.fn().mockResolvedValue(undefined), deleteTransactions: jest.fn().mockResolvedValue(undefined), deleteAclEntries: jest.fn().mockResolvedValue(undefined), updateUserPlugins: jest.fn(), @@ -100,6 +103,11 @@ jest.mock('~/server/services/Endpoints/agents/subagentThreadStore', () => ({ cancelAndDrainForOwner: (...args) => mockCancelAndDrainSubagentThreads(...args), })); +jest.mock('~/server/services/Schedules', () => ({ + quiesceUserSchedules: (...args) => mockQuiesceUserSchedules(...args), + restoreUserSchedulesFromDeletion: (...args) => mockRestoreUserSchedules(...args), +})); + jest.mock('~/server/services/Files/process', () => ({ processDeleteRequest: jest.fn().mockResolvedValue({ deletedFileIds: [], failedFileIds: [] }), })); @@ -160,6 +168,7 @@ describe('verifyEmailController', () => { beforeEach(() => { jest.clearAllMocks(); + mockQuiesceUserSchedules.mockResolvedValue(true); }); it('returns the generic verification error message from service failures', async () => { @@ -343,6 +352,7 @@ describe('deleteUserController', () => { beforeEach(() => { jest.clearAllMocks(); + mockQuiesceUserSchedules.mockResolvedValue(true); }); it('should return 200 on successful deletion', async () => { @@ -361,6 +371,7 @@ describe('deleteUserController', () => { ); expect(mockDrainAgentTriggerDeliveriesForUser).toHaveBeenCalledWith(userId.toString()); expect(mockCancelAndDrainSubagentThreads).toHaveBeenCalledWith(userId.toString(), undefined); + expect(mockQuiesceUserSchedules).toHaveBeenCalledWith(userId.toString(), expect.any(String)); expect(beginAgentTriggerUserDeletion.mock.invocationCallOrder[0]).toBeLessThan( mockPrepareAgentTriggerUserPurge.mock.invocationCallOrder[0], ); @@ -371,6 +382,9 @@ describe('deleteUserController', () => { mockCancelAndDrainSubagentThreads.mock.invocationCallOrder[0], ); expect(mockCancelAndDrainSubagentThreads.mock.invocationCallOrder[0]).toBeLessThan( + mockQuiesceUserSchedules.mock.invocationCallOrder[0], + ); + expect(mockQuiesceUserSchedules.mock.invocationCallOrder[0]).toBeLessThan( deleteMessages.mock.invocationCallOrder[0], ); expect(deleteMessages.mock.invocationCallOrder[0]).toBeLessThan( @@ -382,6 +396,8 @@ describe('deleteUserController', () => { expect(mockPurgeAgentTriggerDeliveriesForUser).toHaveBeenCalledWith(userId.toString()); expect(cancelAgentTriggerUserDeletion).not.toHaveBeenCalled(); expect(mockCancelAgentTriggerUserPurge).not.toHaveBeenCalled(); + // A successful deletion hard-deletes the schedules; it must never restore them. + expect(mockRestoreUserSchedules).not.toHaveBeenCalled(); }); it('aborts generations admitted before the deletion fence before erasing messages', async () => { @@ -428,6 +444,17 @@ describe('deleteUserController', () => { expect(mockCancelAgentTriggerUserPurge).toHaveBeenCalledWith(userIdString, deletionFence); expect(cancelAgentTriggerUserDeletion).toHaveBeenCalledWith(userIdString, deletionFence); expect(deleteUserById).not.toHaveBeenCalled(); + // Account survives -> its suspended schedules are restored under the quiesce token. + expect(mockRestoreUserSchedules).toHaveBeenCalledWith( + userIdString, + mockQuiesceUserSchedules.mock.calls[0][1], + ); + // BEFORE the deletion fence is released: that fence is what refuses new schedule writes, + // so restoring after it would let an owner PATCH — or a second deletion attempt + // re-suspending under a new token — race the restore and strand the disabled snapshot. + expect(mockRestoreUserSchedules.mock.invocationCallOrder[0]).toBeLessThan( + cancelAgentTriggerUserDeletion.mock.invocationCallOrder[0], + ); }); it('fails closed before data cleanup when detached subagents do not drain', async () => { @@ -453,6 +480,41 @@ describe('deleteUserController', () => { expect(deleteUserById).not.toHaveBeenCalled(); }); + it('fails closed and releases deletion fences when schedules cannot be quiesced', async () => { + const userId = new mongoose.Types.ObjectId(); + const userIdString = userId.toString(); + mockQuiesceUserSchedules.mockResolvedValueOnce(false); + const req = { + user: { + id: userIdString, + _id: userId, + email: 'scheduled@test.com', + tenantId: 'tenant-1', + }, + }; + + await deleteUserController(req, mockRes); + + expect(mockRes.status).toHaveBeenCalledWith(500); + expect(deleteMessages).not.toHaveBeenCalled(); + expect(mockGetActiveJobIdsForUser).not.toHaveBeenCalled(); + const deletionFence = beginAgentTriggerUserDeletion.mock.calls[0][1]; + expect(mockCancelAgentTriggerUserPurge).toHaveBeenCalledWith(userIdString, deletionFence); + expect(cancelAgentTriggerUserDeletion).toHaveBeenCalledWith(userIdString, deletionFence); + expect(deleteUserById).not.toHaveBeenCalled(); + // Account survives -> its suspended schedules are restored under the quiesce token. + expect(mockRestoreUserSchedules).toHaveBeenCalledWith( + userIdString, + mockQuiesceUserSchedules.mock.calls[0][1], + ); + // BEFORE the deletion fence is released: that fence is what refuses new schedule writes, + // so restoring after it would let an owner PATCH — or a second deletion attempt + // re-suspending under a new token — race the restore and strand the disabled snapshot. + expect(mockRestoreUserSchedules.mock.invocationCallOrder[0]).toBeLessThan( + cancelAgentTriggerUserDeletion.mock.invocationCallOrder[0], + ); + }); + it('should remove the user from all groups via $pullAll', async () => { const userId = new mongoose.Types.ObjectId(); const userIdStr = userId.toString(); diff --git a/api/server/controllers/__tests__/deleteUser.spec.js b/api/server/controllers/__tests__/deleteUser.spec.js index f6777172a6e..c64a40de8d1 100644 --- a/api/server/controllers/__tests__/deleteUser.spec.js +++ b/api/server/controllers/__tests__/deleteUser.spec.js @@ -28,6 +28,8 @@ const mockPurgeAgentTriggerDeliveriesForUser = jest.fn(); const mockBeginAgentTriggerUserDeletion = jest.fn(); const mockCancelAgentTriggerUserDeletion = jest.fn(); const mockCancelAndDrainSubagentThreads = jest.fn(); +const mockQuiesceUserSchedules = jest.fn(); +const mockDeleteSchedulesByUser = jest.fn(); jest.mock('@librechat/data-schemas', () => ({ logger: { error: jest.fn(), info: jest.fn() }, @@ -85,6 +87,7 @@ jest.mock('~/models', () => ({ deleteTokens: jest.fn(), removeUserFromAllGroups: jest.fn(), deleteAclEntries: jest.fn(), + deleteSchedulesByUser: (...args) => mockDeleteSchedulesByUser(...args), getSoleOwnedResourceIds: jest.fn().mockResolvedValue([]), })); @@ -127,6 +130,10 @@ jest.mock('~/server/services/Endpoints/agents/subagentThreadStore', () => ({ cancelAndDrainForOwner: (...args) => mockCancelAndDrainSubagentThreads(...args), })); +jest.mock('~/server/services/Schedules', () => ({ + quiesceUserSchedules: (...args) => mockQuiesceUserSchedules(...args), +})); + jest.mock('~/server/services/Config', () => ({ getAppConfig: jest.fn(), })); @@ -171,6 +178,8 @@ function stubDeletionMocks() { mockBeginAgentTriggerUserDeletion.mockResolvedValue('acquired'); mockCancelAgentTriggerUserDeletion.mockResolvedValue(true); mockCancelAndDrainSubagentThreads.mockResolvedValue(); + mockQuiesceUserSchedules.mockResolvedValue(true); + mockDeleteSchedulesByUser.mockResolvedValue(); } beforeEach(() => { diff --git a/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js b/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js index c2f11732577..7428c2287f1 100644 --- a/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js +++ b/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js @@ -47,6 +47,8 @@ jest.mock('@librechat/data-schemas', () => ({ jest.mock('@librechat/api', () => ({ sendEvent: jest.fn(), + isScheduleFireRequest: jest.fn(() => false), + exemptFromConcurrencyLimiter: jest.fn(() => false), toPendingSteer: jest.fn((item) => item), isSteerPreemptSupported: jest.fn(() => true), buildRecoveredSteerPayload: jest.fn(() => null), diff --git a/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js b/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js index e6592429772..5b6fb7dab6c 100644 --- a/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js +++ b/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js @@ -64,6 +64,10 @@ const mockGetConvo = jest.fn(); const mockGetMessages = jest.fn(); const mockSaveMessage = jest.fn(); const mockIsAgentTriggerPrincipalActive = jest.fn(); +const mockIsScheduleFireRequest = jest.fn(); +const mockExemptFromConcurrencyLimiter = jest.fn(); +const mockRecordScheduleOutcome = jest.fn(); +const mockIsScheduleLive = jest.fn(); const mockDeleteAgentCheckpoint = jest.fn(); const mockStartupTelemetry = { mark: jest.fn(), @@ -132,6 +136,8 @@ jest.mock('@librechat/data-schemas', () => ({ jest.mock('@librechat/api', () => ({ sendEvent: jest.fn(), + isScheduleFireRequest: (...args) => mockIsScheduleFireRequest(...args), + exemptFromConcurrencyLimiter: (...args) => mockExemptFromConcurrencyLimiter(...args), toPendingSteer: jest.fn((item) => item), /** Recorded onto the job so the steer route can honour the OWNING replica's * seal capability rather than its own probe. */ @@ -214,6 +220,11 @@ jest.mock('~/models', () => ({ isAgentTriggerPrincipalActive: (...args) => mockIsAgentTriggerPrincipalActive(...args), })); +jest.mock('~/server/services/Schedules', () => ({ + recordScheduleOutcome: (...args) => mockRecordScheduleOutcome(...args), + isScheduleLive: (...args) => mockIsScheduleLive(...args), +})); + const AgentController = require('../request'); const { ErrorTypes } = require('librechat-data-provider'); const { disposeClient: mockDisposeClient } = require('~/server/cleanup'); @@ -250,6 +261,12 @@ describe('ResumableAgentController resume metadata', () => { mockGetConvo.mockResolvedValue({ createdAt: '2026-06-07T00:00:00.000Z' }); mockGetMessages.mockResolvedValue([]); mockIsAgentTriggerPrincipalActive.mockResolvedValue(true); + mockIsScheduleFireRequest.mockImplementation((req) => req?._isScheduledFire === true); + mockExemptFromConcurrencyLimiter.mockImplementation( + (req) => req?._isScheduledFire === true && req?._isManualScheduledFire !== true, + ); + mockRecordScheduleOutcome.mockResolvedValue(true); + mockIsScheduleLive.mockResolvedValue(true); mockGenerationJobManager.createJob.mockResolvedValue({ createdAt: 1000, metadata: { @@ -286,7 +303,7 @@ describe('ResumableAgentController resume metadata', () => { }), ); mockGenerationJobManager.finishTerminalJob.mockResolvedValue(undefined); - mockGenerationJobManager.completeJob.mockResolvedValue(undefined); + mockGenerationJobManager.completeJob.mockResolvedValue(true); mockGenerationJobManager.beginProviderExecution.mockResolvedValue(true); mockGenerationJobManager.markProviderExecutionDrained.mockResolvedValue(true); mockGenerationJobManager.failPausePersistence.mockResolvedValue(true); @@ -797,6 +814,73 @@ describe('ResumableAgentController resume metadata', () => { expect(mockDecrementPendingRequest).toHaveBeenCalledWith('user-123'); }); + it('rejects a superseded automatic occurrence after durable job creation and preserves its outcome', async () => { + const conversationId = 'scheduled-conversation-123'; + mockGenerationJobManager.claimGeneration.mockResolvedValue( + wonGenerationClaim({ streamId: conversationId, conversationId }), + ); + mockIsScheduleLive.mockResolvedValue(false); + const req = { + _isScheduledFire: true, + _isManualScheduledFire: false, + user: { id: 'user-123' }, + body: { + text: 'Run the scheduled digest.', + messageId: 'scheduled-user-message', + clientRequestId: 'sched:schedule-1:2026-08-17T12:00:00-000Z', + conversationId: 'new', + newConversationId: conversationId, + scheduleId: 'schedule-1', + scheduledFor: '2026-08-17T12:00:00.000Z', + scheduleConfigRevision: 7, + endpointOption: { endpoint: 'agents', agent_id: 'agent-1' }, + }, + config: {}, + }; + const res = createResumableResponse(); + const initializeClient = jest.fn(); + + await AgentController(req, res, jest.fn(), initializeClient, null); + + expect(mockCheckAndIncrementPendingRequest).not.toHaveBeenCalled(); + expect(mockGenerationJobManager.createJob).toHaveBeenCalledWith( + conversationId, + 'user-123', + conversationId, + expect.objectContaining({ + initialMetadata: expect.objectContaining({ + scheduleId: 'schedule-1', + scheduledFor: '2026-08-17T12:00:00.000Z', + scheduleConfigRevision: 7, + preserveForScheduleReconcile: true, + }), + }), + ); + expect(mockIsScheduleLive).toHaveBeenCalledWith('schedule-1', 7, { + automatic: true, + policy: true, + }); + expect(initializeClient).not.toHaveBeenCalled(); + expect(res.status).toHaveBeenCalledWith(409); + expect(res.json).toHaveBeenCalledWith({ + status: 409, + code: 'SCHEDULE_NO_LONGER_ACTIVE', + error: 'This scheduled occurrence is no longer active', + generationProtocolVersion: 1, + }); + expect(mockRecordScheduleOutcome).toHaveBeenCalledWith({ + scheduleId: 'schedule-1', + scheduledFor: '2026-08-17T12:00:00.000Z', + streamId: conversationId, + jobCreatedAt: 1000, + status: 'interrupted', + conversationId, + clearConversationId: false, + error: 'This scheduled occurrence is no longer active', + }); + expect(mockDecrementPendingRequest).not.toHaveBeenCalled(); + }); + it('does not start a provider when account deletion or replacement wins the startup CAS', async () => { mockGenerationJobManager.beginProviderExecution.mockResolvedValue(false); const req = { diff --git a/api/server/controllers/agents/__tests__/resume.spec.js b/api/server/controllers/agents/__tests__/resume.spec.js index eb308080540..034c53189ee 100644 --- a/api/server/controllers/agents/__tests__/resume.spec.js +++ b/api/server/controllers/agents/__tests__/resume.spec.js @@ -60,6 +60,7 @@ const mockGenerationJobManager = { publishTerminalClaim: jest.fn(), finishTerminalJob: jest.fn(), completeJob: jest.fn(), + abortJob: jest.fn(), beginProviderExecution: jest.fn(), markProviderExecutionDrained: jest.fn(), failPausePersistence: jest.fn(), @@ -82,6 +83,12 @@ const mockGetMessages = jest.fn(); const mockDisposeClient = jest.fn(); const mockGetMCPRequestContext = jest.fn(); const mockCleanupMCPRequestContextForReq = jest.fn(); +const mockRecordScheduleOutcome = jest.fn(); +const mockIsScheduleLive = jest.fn(); +const mockClaimScheduleResume = jest.fn(); +const mockReleaseScheduleResumeClaim = jest.fn(); +const mockFinalizeScheduleResumeClaim = jest.fn(); +const mockReleaseScheduleResumeFence = jest.fn(); jest.mock('@librechat/data-schemas', () => ({ ...jest.requireActual('@librechat/data-schemas'), @@ -104,6 +111,15 @@ jest.mock('~/models', () => ({ getMessages: (...args) => mockGetMessages(...args), })); +jest.mock('~/server/services/Schedules', () => ({ + recordScheduleOutcome: (...args) => mockRecordScheduleOutcome(...args), + claimScheduleResume: (...args) => mockClaimScheduleResume(...args), + releaseScheduleResumeClaim: (...args) => mockReleaseScheduleResumeClaim(...args), + finalizeScheduleResumeClaim: (...args) => mockFinalizeScheduleResumeClaim(...args), + releaseScheduleResumeFence: (...args) => mockReleaseScheduleResumeFence(...args), + isScheduleLive: (...args) => mockIsScheduleLive(...args), +})); + jest.mock('~/server/cleanup', () => ({ disposeClient: (...args) => mockDisposeClient(...args), })); @@ -248,12 +264,23 @@ describe('ResumeAgentController (POST /agents/chat/resume)', () => { ); mockGenerationJobManager.finishTerminalJob.mockResolvedValue(undefined); mockGenerationJobManager.completeJob.mockResolvedValue(true); + mockGenerationJobManager.abortJob.mockResolvedValue({ success: true }); mockGenerationJobManager.beginProviderExecution.mockResolvedValue(true); mockGenerationJobManager.markProviderExecutionDrained.mockResolvedValue(true); mockGenerationJobManager.failPausePersistence.mockResolvedValue(true); mockGenerationJobManager.approvals.resolve.mockResolvedValue(true); mockGenerationJobManager.approvals.ownsPausePersistence.mockResolvedValue(true); mockGenerationJobManager.approvals.finishPausePersistence.mockResolvedValue(true); + mockRecordScheduleOutcome.mockResolvedValue(true); + mockIsScheduleLive.mockResolvedValue(true); + mockClaimScheduleResume.mockResolvedValue({ + capacitySlot: 0, + claimToken: 'resume-token', + leaseBy: 'resume:resume-token', + }); + mockReleaseScheduleResumeClaim.mockResolvedValue(true); + mockFinalizeScheduleResumeClaim.mockResolvedValue(true); + mockReleaseScheduleResumeFence.mockResolvedValue(undefined); // `decrementPendingRequest` runs in the controller's `finally` on every // post-ACK path, so resolving on it signals the async continuation is done. @@ -304,6 +331,222 @@ describe('ResumeAgentController (POST /agents/chat/resume)', () => { ...extra, }); + describe('scheduled occurrence lifecycle', () => { + const scheduledFor = '2026-08-17T12:00:00.000Z'; + const makeScheduledJob = () => + makeToolApprovalJob({ + metadata: { + scheduleId: 'schedule-1', + scheduledFor, + scheduleConfigRevision: 4, + checkpointNamespace: '1000', + }, + }); + + it('stops and settles an occurrence that became inactive while awaiting approval', async () => { + mockGenerationJobManager.getJob.mockResolvedValue(makeScheduledJob()); + mockIsScheduleLive.mockResolvedValue(false); + + const res = await post(approveBody()); + + expect(res.status).toBe(409); + expect(res.body).toMatchObject({ code: 'SCHEDULE_NO_LONGER_ACTIVE' }); + expect(mockIsScheduleLive).toHaveBeenCalledWith('schedule-1', 4, { + automatic: true, + policy: true, + }); + expect(mockGenerationJobManager.abortJob).toHaveBeenCalledWith(CONVO_ID, { + expectedCreatedAt: 1000, + awaitProviderDrain: true, + }); + expect(mockRecordScheduleOutcome).toHaveBeenCalledWith({ + scheduleId: 'schedule-1', + scheduledFor, + streamId: CONVO_ID, + jobCreatedAt: 1000, + status: 'interrupted', + conversationId: CONVO_ID, + error: 'Schedule was disabled, changed, or deleted before approval', + }); + expect(mockDeleteAgentCheckpoint).toHaveBeenCalledWith( + CONVO_ID, + { type: 'mongo' }, + undefined, + { checkpointNamespace: '1000' }, + ); + expect(mockGenerationJobManager.approvals.resolve).not.toHaveBeenCalled(); + }); + + it('fails closed without settling or pruning when provider drain cannot be confirmed', async () => { + mockGenerationJobManager.getJob.mockResolvedValue(makeScheduledJob()); + mockIsScheduleLive.mockResolvedValue(false); + mockGenerationJobManager.abortJob.mockRejectedValue(new Error('drain timed out')); + + const res = await post(approveBody()); + + expect(res.status).toBe(503); + expect(res.headers['retry-after']).toBe('1'); + expect(res.body).toMatchObject({ code: 'SCHEDULE_STOP_UNCONFIRMED' }); + expect(mockRecordScheduleOutcome).not.toHaveBeenCalled(); + expect(mockDeleteAgentCheckpoint).not.toHaveBeenCalled(); + expect(mockGenerationJobManager.approvals.resolve).not.toHaveBeenCalled(); + }); + + it('records success after resumed persistence and before terminal publication', async () => { + mockGenerationJobManager.getJob.mockResolvedValue(makeScheduledJob()); + + const res = await post(approveBody()); + expect(res.status).toBe(200); + await settled; + + expect(mockClaimScheduleResume).toHaveBeenCalledWith('schedule-1', scheduledFor, { + expectedConfigRevision: 4, + automatic: true, + }); + expect(mockFinalizeScheduleResumeClaim).toHaveBeenCalledWith( + 'schedule-1', + 'resume-token', + 'resume:resume-token', + { expectedConfigRevision: 4, automatic: true }, + ); + + expect(mockRecordScheduleOutcome).toHaveBeenCalledWith({ + scheduleId: 'schedule-1', + scheduledFor, + streamId: CONVO_ID, + jobCreatedAt: 1000, + status: 'success', + conversationId: CONVO_ID, + }); + expect(mockRecordScheduleOutcome.mock.invocationCallOrder[0]).toBeLessThan( + mockGenerationJobManager.publishTerminalClaim.mock.invocationCallOrder[0], + ); + }); + + it('keeps the approval paused when global scheduled-run capacity is full', async () => { + mockGenerationJobManager.getJob.mockResolvedValue(makeScheduledJob()); + mockClaimScheduleResume.mockResolvedValue({ conflict: 'capacity' }); + + const res = await post(approveBody()); + + expect(res.status).toBe(429); + expect(res.headers['retry-after']).toBe('1'); + expect(res.body).toMatchObject({ code: 'SCHEDULE_CAPACITY' }); + expect(mockGenerationJobManager.approvals.resolve).not.toHaveBeenCalled(); + expect(mockDecrementPendingRequest).toHaveBeenCalledWith(USER_ID); + }); + + it('rolls back scheduled capacity when the approval CAS does not consume the action', async () => { + const job = makeScheduledJob(); + mockGenerationJobManager.getJob.mockResolvedValue(job); + mockGenerationJobManager.approvals.resolve.mockResolvedValue(false); + + const res = await post(approveBody()); + + expect(res.status).toBe(409); + expect(mockReleaseScheduleResumeClaim).toHaveBeenCalledWith('schedule-1', scheduledFor, 0); + expect(mockReleaseScheduleResumeFence).toHaveBeenCalledWith( + 'schedule-1', + 'resume:resume-token', + ); + }); + + it('stops before provider execution when an edit wins the final resume handoff', async () => { + mockGenerationJobManager.getJob.mockResolvedValue(makeScheduledJob()); + mockFinalizeScheduleResumeClaim.mockResolvedValue(false); + + const res = await post(approveBody()); + + expect(res.status).toBe(409); + expect(res.body).toMatchObject({ code: 'SCHEDULE_NO_LONGER_ACTIVE' }); + expect(mockGenerationJobManager.approvals.resolve).toHaveBeenCalled(); + expect(mockGenerationJobManager.abortJob).toHaveBeenCalledWith(CONVO_ID, { + expectedCreatedAt: 1000, + awaitProviderDrain: true, + }); + expect(mockRecordScheduleOutcome).toHaveBeenCalledWith({ + scheduleId: 'schedule-1', + scheduledFor, + streamId: CONVO_ID, + jobCreatedAt: 1000, + status: 'interrupted', + conversationId: CONVO_ID, + error: 'Schedule was disabled, changed, or deleted before approval', + }); + expect(mockInitializeClient).not.toHaveBeenCalled(); + }); + + it('settles a scheduled continuation stopped during its resumed segment', async () => { + const job = makeScheduledJob(); + mockGenerationJobManager.getJob.mockResolvedValue(job); + mockInitializeClient.mockImplementation(async () => { + job.abortController.abort(); + return { client: makeClient(), userMCPAuthMap: {} }; + }); + + const res = await post(approveBody()); + expect(res.status).toBe(200); + await settled; + await flush(); + + expect(mockRecordScheduleOutcome).toHaveBeenCalledWith({ + scheduleId: 'schedule-1', + scheduledFor, + streamId: CONVO_ID, + jobCreatedAt: 1000, + status: 'interrupted', + conversationId: CONVO_ID, + error: 'Scheduled run was stopped', + }); + }); + + it('records an empty-preempt resumed segment as interrupted, not successful', async () => { + mockGenerationJobManager.getJob.mockResolvedValue(makeScheduledJob()); + mockInitializeClient.mockResolvedValue({ + client: makeClient({ + run: { + getPreemptStats: () => ({ emptyBoundaries: 1 }), + getHaltReason: () => 'preempt_incomplete', + }, + }), + userMCPAuthMap: {}, + }); + + const res = await post(approveBody()); + expect(res.status).toBe(200); + await settled; + await flush(); + + expect(mockSaveMessage).toHaveBeenCalledWith( + expect.anything(), + expect.objectContaining({ unfinished: true }), + expect.anything(), + ); + expect(mockRecordScheduleOutcome).toHaveBeenCalledWith({ + scheduleId: 'schedule-1', + scheduledFor, + streamId: CONVO_ID, + jobCreatedAt: 1000, + status: 'interrupted', + conversationId: CONVO_ID, + error: 'Scheduled run was interrupted before completion', + }); + }); + + it('does not overwrite the terminal winner when failed-resume finalization loses its CAS', async () => { + mockGenerationJobManager.getJob.mockResolvedValue(makeScheduledJob()); + mockGenerationJobManager.completeJob.mockResolvedValue(false); + mockInitializeClient.mockRejectedValue(new Error('resume reconstruction failed')); + + const res = await post(approveBody()); + expect(res.status).toBe(200); + await settled; + + expect(mockGenerationJobManager.completeJob).toHaveBeenCalled(); + expect(mockRecordScheduleOutcome).not.toHaveBeenCalled(); + }); + }); + describe('temporal context restore', () => { it('restores req.conversationCreatedAt from the convo before initializeClient', async () => { // Temporal prompt vars must resolve against the paused anchor, not resume wall-clock. diff --git a/api/server/controllers/agents/client.js b/api/server/controllers/agents/client.js index 3da58904e82..bb9201ab899 100644 --- a/api/server/controllers/agents/client.js +++ b/api/server/controllers/agents/client.js @@ -2884,7 +2884,9 @@ class AgentClient extends BaseClient { // the flag; if it fails here, the teardown still releases (it checks the flag). if (!this.pendingRequestReleased) { try { - await decrementPendingRequest(this.options.req?.user?.id); + if (this.options.req?._scheduleConcurrencyExempt !== true) { + await decrementPendingRequest(this.options.req?.user?.id); + } this.pendingRequestReleased = true; } catch (err) { logger.error(`[AgentClient] Failed to release request slot on pause ${streamId}`, err); diff --git a/api/server/controllers/agents/request.js b/api/server/controllers/agents/request.js index 580b545f4cc..5c993831805 100644 --- a/api/server/controllers/agents/request.js +++ b/api/server/controllers/agents/request.js @@ -18,6 +18,8 @@ const { decrementPendingRequest, sanitizeMessageForTransmit, checkAndIncrementPendingRequest, + exemptFromConcurrencyLimiter, + isScheduleFireRequest, isUnpersistedPreliminaryParent, resolveConversationAnchor, getAgentStartupTelemetry, @@ -33,6 +35,7 @@ const { cleanupMCPRequestContextForReq, } = require('~/server/services/MCPRequestContext'); const { logViolation } = require('~/cache'); +const { recordScheduleOutcome, isScheduleLive } = require('~/server/services/Schedules'); const { saveMessage, getMessages, getConvo, isAgentTriggerPrincipalActive } = require('~/models'); const { GENERATION_PROTOCOL_HEADER, @@ -182,8 +185,20 @@ async function finishResumableRequest(req, userId) { try { await cleanupMCPRequestContextForReq(req); } finally { - await decrementPendingRequest(userId); + if (req._scheduleConcurrencyExempt !== true) { + await decrementPendingRequest(userId); + } + } +} + +function classifyScheduledFailure(error, aborted = false) { + if (aborted || error?.code === 'SCHEDULE_NO_LONGER_ACTIVE') { + return { status: 'interrupted', error: error?.message }; } + if (error?.message?.includes(ViolationTypes.TOKEN_BALANCE)) { + return { status: 'skipped_balance' }; + } + return { status: 'error', error: error?.message || 'Generation failed' }; } const JOB_RECORD_WAIT_ATTEMPTS = 5; @@ -338,8 +353,16 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit parentMessageId = null, overrideParentMessageId = null, responseMessageId: editedResponseMessageId = null, + scheduleId: bodyScheduleId = null, + scheduledFor: bodyScheduledFor = null, + scheduleConfigRevision: bodyScheduleConfigRevision = null, } = req.body; + const isScheduledFire = isScheduleFireRequest(req); + const scheduleId = isScheduledFire ? bodyScheduleId : null; + const scheduledFor = isScheduledFire ? bodyScheduledFor : null; + const scheduleConfigRevision = isScheduledFire ? bodyScheduleConfigRevision : undefined; + const userId = req.user.id; const tenantId = req.user.tenantId; const rawClientRequestId = req.body?.clientRequestId; @@ -446,12 +469,17 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit req.body.overrideUserMessageId = `${recoveredSteerId}${Constants.COMMON_DIVIDER}0`; } const isNewConvo = !reqConversationId || reqConversationId === 'new'; + const scheduledNewConversationId = + isScheduledFire && typeof req.body?.newConversationId === 'string' + ? req.body.newConversationId + : null; let conversationId = reqConversationId; if (isNewConvo) { conversationId = - typeof clientRequestId === 'string' && clientRequestId.length > 0 + scheduledNewConversationId ?? + (typeof clientRequestId === 'string' && clientRequestId.length > 0 ? uuidv5(`${userId}:${clientRequestId}`, NEW_CONVERSATION_IDEMPOTENCY_NAMESPACE) - : crypto.randomUUID(); + : crypto.randomUUID()); } const conversationAnchorPromise = resolveConversationCreatedAt({ userId, @@ -936,26 +964,53 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit } } - const { allowed, pendingRequests, limit } = await checkAndIncrementPendingRequest(userId); - if (!allowed) { - if (ownedIdempotencyClaim) { - await GenerationJobManager.releaseGeneration( - userId, - clientRequestId, - streamId, - ownedIdempotencyClaim, - ).catch(() => {}); + const scheduleConcurrencyExempt = exemptFromConcurrencyLimiter(req); + req._scheduleConcurrencyExempt = scheduleConcurrencyExempt; + if (!scheduleConcurrencyExempt) { + const { allowed, pendingRequests, limit } = await checkAndIncrementPendingRequest(userId); + if (!allowed) { + if (ownedIdempotencyClaim) { + await GenerationJobManager.releaseGeneration( + userId, + clientRequestId, + streamId, + ownedIdempotencyClaim, + ).catch(() => {}); + } + const violationInfo = getViolationInfo(pendingRequests, limit); + await logViolation(req, res, ViolationTypes.CONCURRENT, violationInfo, violationInfo.score); + startupTelemetry?.end('rejected'); + return sendGenerationJson(res, 429, violationInfo, generationProtocolVersion); } - const violationInfo = getViolationInfo(pendingRequests, limit); - await logViolation(req, res, ViolationTypes.CONCURRENT, violationInfo, violationInfo.score); - startupTelemetry?.end('rejected'); - return sendGenerationJson(res, 429, violationInfo, generationProtocolVersion); } startupTelemetry?.mark('request_admitted'); let client = null; let jobCreatedAt; let providerExecutionId; + let scheduleTerminalOutcomeRecorded = false; + const settleScheduledRun = async ({ status, error, clearConversationId = false }) => { + if (!scheduleId) { + return true; + } + if (status !== 'requires_action' && scheduleTerminalOutcomeRecorded) { + return true; + } + const recorded = await recordScheduleOutcome({ + scheduleId, + scheduledFor, + streamId, + jobCreatedAt, + status, + conversationId, + clearConversationId, + error, + }); + if (recorded && status !== 'requires_action') { + scheduleTerminalOutcomeRecorded = true; + } + return recorded; + }; try { logger.debug(`[ResumableAgentController] Creating job`, { @@ -995,6 +1050,17 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit // Persist temporary-chat state so a HITL resume keeps the resumed response // non-persisted instead of trusting the resume request to re-send the flag. isTemporary: req.body?.isTemporary, + ...(scheduleId + ? { + scheduleId, + scheduledFor, + preserveForScheduleReconcile: true, + ...(Number.isSafeInteger(scheduleConfigRevision) && { + scheduleConfigRevision, + }), + ...(req._isManualScheduledFire === true && { scheduleManual: true }), + } + : {}), responseMessageId: preliminaryResponseMessageId, userMessage: preliminaryUserMessage, }, @@ -1016,6 +1082,18 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit status: 409, }); } + if ( + scheduleId && + !(await isScheduleLive(scheduleId, scheduleConfigRevision, { + automatic: req._isManualScheduledFire !== true, + policy: true, + })) + ) { + throw Object.assign(new Error('This scheduled occurrence is no longer active'), { + code: 'SCHEDULE_NO_LONGER_ACTIVE', + status: 409, + }); + } if ( providerExecutionId && !(await GenerationJobManager.beginProviderExecution( @@ -1187,6 +1265,11 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit completeErr, ); }); + await settleScheduledRun({ + status: 'interrupted', + error: 'Request aborted during initialization', + clearConversationId: job.createdEventEmitted !== true, + }); startupTelemetry?.end('aborted'); try { await finishResumableRequest(req, userId); @@ -1260,6 +1343,7 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit * may call completeJob: the pause may already have been replaced by a newer * action or generation by the time the persistence failure is observed. */ let pausePersistenceFailed = false; + let pausePersistenceFailureFinalized = false; const finishOwnedTerminalClaim = async () => { if (!terminalClaim || terminalClaimFinished) { return; @@ -1586,21 +1670,21 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit await commitRecoveredSteer(); } catch (pausePersistenceError) { pausePersistenceFailed = true; - let failed; try { - failed = await GenerationJobManager.failPausePersistence( - streamId, - pauseActionId, - pausePersistenceError?.message ?? 'Pause persistence failed', - pauseCreatedAt, - ); + pausePersistenceFailureFinalized = + (await GenerationJobManager.failPausePersistence( + streamId, + pauseActionId, + pausePersistenceError?.message ?? 'Pause persistence failed', + pauseCreatedAt, + )) === true; } catch (failError) { logger.error( `[ResumableAgentController] Failed to terminalize pause persistence error for ${streamId}`, failError, ); } - if (failed === true) { + if (pausePersistenceFailureFinalized) { /** Namespaced checkpoints belong exclusively to this epoch, * so the exact pause-failure CAS winner can safely remove the * now-unresumable graph state. Legacy shared namespaces are @@ -1621,7 +1705,7 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit ); } } - } else if (failed === false) { + } else if (pausePersistenceFailureFinalized === false) { logger.warn( `[ResumableAgentController] Skipping stale pause persistence failure — ${streamId} no longer owns its barrier`, ); @@ -1638,6 +1722,17 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit `[ResumableAgentController] Pause persistence barrier changed before release: ${streamId}`, ); } + // The pause projection is what moves the run row off `started` and frees its + // GLOBAL capacity slot. recordScheduleOutcome already retried it; a `false` + // here means every attempt failed, leaving the row `started` while the job + // sits `requires_action`. Surface it — the armed engine's reconciler replays + // this state, and the clustered sweep now converges it too, but a silent drop + // gave neither a reason to look. + if (!(await settleScheduledRun({ status: 'requires_action' }))) { + logger.error( + `[ResumableAgentController] Failed to project the scheduled pause for ${streamId}; run stays active until reconciliation replays it`, + ); + } } else { logger.debug( `[ResumableAgentController] Skipping stale pause persistence — ${streamId} no longer owns its barrier`, @@ -1650,7 +1745,7 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit // (so a fast /resume isn't 429'd); only release here if that didn't happen. // Always run the MCP request-context cleanup. await cleanupMCPRequestContextForReq(req); - if (!client?.pendingRequestReleased) { + if (!client?.pendingRequestReleased && req._scheduleConcurrencyExempt !== true) { await decrementPendingRequest(userId); } if (client) { @@ -1781,6 +1876,17 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit await titleEventPromise; } + let scheduleCompletionError; + if (terminalWasAborted) { + scheduleCompletionError = 'Scheduled run was stopped'; + } else if (preemptIncomplete) { + scheduleCompletionError = 'Scheduled run was interrupted before completion'; + } + await settleScheduledRun({ + status: terminalWasAborted || preemptIncomplete ? 'interrupted' : 'success', + ...(scheduleCompletionError != null && { error: scheduleCompletionError }), + }); + let terminalPublicationStarted = false; try { const pendingSteers = terminalClaim.drainedSteers.map(toPendingSteer); @@ -1893,7 +1999,9 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit // transition can win. Settle its pending marker with conservative // reconciliation on any required-write/final-construction failure, // then release exactly that claim. + let ownsScheduledFailure = false; if (terminalClaim && !terminalClaimFinished) { + ownsScheduledFailure = true; try { await GenerationJobManager.publishTerminalClaim(terminalClaim, null); } catch (publishError) { @@ -1915,6 +2023,7 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit ); startupTelemetry?.end('error', error); } else if (pausePersistenceFailed) { + ownsScheduledFailure = pausePersistenceFailureFinalized; // failPausePersistence owns the only legal requires_action -> error // transition for this exact action/epoch. Never fall through to // completeJob, which could race a newer action or replacement job. @@ -1924,6 +2033,7 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit ); startupTelemetry?.end('error', error); } else if (job.abortController.signal.aborted || error.message?.includes('abort')) { + ownsScheduledFailure = true; logger.debug(`[ResumableAgentController] Generation aborted for ${streamId}`); startupTelemetry?.end('aborted'); // abortJob already handled emitDone and completeJob @@ -1933,7 +2043,9 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit try { // completeJob first wins running -> error and atomically parks // steers, then publishes. A competing abort/pause emits nothing. - await GenerationJobManager.completeJob(streamId, generationError, jobCreatedAt); + ownsScheduledFailure = + (await GenerationJobManager.completeJob(streamId, generationError, jobCreatedAt)) === + true; } catch (completeErr) { logger.warn( '[ResumableAgentController] completeJob failed during generation-error cleanup', @@ -1944,6 +2056,14 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit } } + if (ownsScheduledFailure && !scheduleTerminalOutcomeRecorded) { + const scheduledFailure = classifyScheduledFailure( + error, + job.abortController.signal.aborted, + ); + await settleScheduledRun(scheduledFailure); + } + try { await finishResumableRequest(req, userId); } finally { @@ -1962,15 +2082,24 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit `[ResumableAgentController] Unhandled error in background generation: ${err.message}`, ); startupTelemetry?.end('error', err); + let errorFinalized = false; if (!pausePersistenceFailed) { - await GenerationJobManager.completeJob(streamId, err.message, jobCreatedAt).catch( - (completeErr) => { - logger.warn( - '[ResumableAgentController] completeJob failed during background-error cleanup', - completeErr, - ); - }, - ); + errorFinalized = + (await GenerationJobManager.completeJob(streamId, err.message, jobCreatedAt).catch( + (completeErr) => { + logger.warn( + '[ResumableAgentController] completeJob failed during background-error cleanup', + completeErr, + ); + return false; + }, + )) === true; + } + if ( + (errorFinalized || (pausePersistenceFailed && pausePersistenceFailureFinalized)) && + !scheduleTerminalOutcomeRecorded + ) { + await settleScheduledRun(classifyScheduledFailure(err)); } try { await finishResumableRequest(req, userId); @@ -2084,18 +2213,24 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit // release + pending-request decrement below, or the retry stays wedged behind the claim // and the concurrency slot leaks — so swallow its error. (A failed completeJob did not // finalize anything, so releasing afterward can't let it abort a later replacement.) + let initializationFinalized = jobCreatedAt == null; if (jobCreatedAt != null) { const initializationError = initializationFailure ? JSON.stringify(initializationFailure) : error.message || 'Failed to start generation'; - await GenerationJobManager.completeJob(streamId, initializationError, jobCreatedAt).catch( - (completeErr) => { - logger.warn( - '[ResumableAgentController] completeJob failed during init-error cleanup', - completeErr, - ); - }, - ); + initializationFinalized = + (await GenerationJobManager.completeJob(streamId, initializationError, jobCreatedAt).catch( + (completeErr) => { + logger.warn( + '[ResumableAgentController] completeJob failed during init-error cleanup', + completeErr, + ); + return false; + }, + )) === true; + } + if (initializationFinalized && !scheduleTerminalOutcomeRecorded) { + await settleScheduledRun(classifyScheduledFailure(error)); } if (ownedIdempotencyClaim) { await GenerationJobManager.releaseGeneration( diff --git a/api/server/controllers/agents/resume.js b/api/server/controllers/agents/resume.js index cb3ea829084..e4eb76e2953 100644 --- a/api/server/controllers/agents/resume.js +++ b/api/server/controllers/agents/resume.js @@ -1,6 +1,6 @@ const { randomUUID } = require('crypto'); const { logger } = require('@librechat/data-schemas'); -const { Constants, EModelEndpoint } = require('librechat-data-provider'); +const { Constants, EModelEndpoint, ViolationTypes } = require('librechat-data-provider'); const { GenerationJobManager, isPendingActionStale, @@ -30,6 +30,14 @@ const { cleanupMCPRequestContextForReq, } = require('~/server/services/MCPRequestContext'); const { saveMessage, getConvo, getMessages } = require('~/models'); +const { + recordScheduleOutcome, + claimScheduleResume, + releaseScheduleResumeClaim, + finalizeScheduleResumeClaim, + releaseScheduleResumeFence, + isScheduleLive, +} = require('~/server/services/Schedules'); const { GENERATION_PROTOCOL_HEADER, negotiateNewGenerationProtocol, @@ -422,6 +430,20 @@ async function finalizeResumedTurn({ } conversation.title = conversation.title || 'New Chat'; + if (meta.scheduleId) { + await recordScheduleOutcome({ + scheduleId: meta.scheduleId, + scheduledFor: meta.scheduledFor, + streamId, + jobCreatedAt: job.createdAt, + status: preemptIncomplete ? 'interrupted' : 'success', + conversationId, + ...(preemptIncomplete && { + error: 'Scheduled run was interrupted before completion', + }), + }); + } + const pendingSteers = terminalClaim.drainedSteers.map(toPendingSteer); const finalEvent = { final: true, @@ -571,6 +593,65 @@ const ResumeAgentController = async (req, res, next, initializeClient, addTitle) ); } + const scheduleId = job.metadata?.scheduleId; + const scheduledFor = job.metadata?.scheduledFor; + if ( + scheduleId && + !(await isScheduleLive(scheduleId, job.metadata?.scheduleConfigRevision, { + automatic: job.metadata?.scheduleManual !== true, + policy: true, + })) + ) { + let stopped = false; + try { + const abortResult = await GenerationJobManager.abortJob(streamId, { + expectedCreatedAt: job.createdAt, + awaitProviderDrain: true, + }); + stopped = abortResult != null && abortResult.failureReason == null; + } catch (error) { + logger.warn('[ResumeAgentController] Failed to stop inactive scheduled run', error); + } + if (!stopped) { + res.set('Retry-After', '1'); + return sendGenerationJson( + res, + 503, + { + code: 'SCHEDULE_STOP_UNCONFIRMED', + error: 'The inactive scheduled run could not be confirmed stopped. Please retry.', + }, + generationProtocolVersion, + ); + } + await recordScheduleOutcome({ + scheduleId, + scheduledFor, + streamId, + jobCreatedAt: job.createdAt, + status: 'interrupted', + conversationId, + error: 'Schedule was disabled, changed, or deleted before approval', + }); + const checkpointNamespace = job.metadata?.checkpointNamespace; + if (typeof checkpointNamespace === 'string' && checkpointNamespace !== '') { + await deleteAgentCheckpoint( + conversationId, + req.config?.endpoints?.[EModelEndpoint.agents]?.checkpointer, + undefined, + { checkpointNamespace }, + ).catch((error) => { + logger.warn('[ResumeAgentController] Failed to prune inactive schedule checkpoint', error); + }); + } + return sendGenerationJson( + res, + 409, + { code: 'SCHEDULE_NO_LONGER_ACTIVE', error: 'This schedule can no longer be resumed' }, + generationProtocolVersion, + ); + } + const pendingAction = job.metadata?.pendingAction; if (job.status !== 'requires_action') { return sendGenerationJson( @@ -720,6 +801,109 @@ const ResumeAgentController = async (req, res, next, initializeClient, addTitle) ); } + // Finish the legacy checkpoint snapshot before claiming scheduled capacity. + // It is independent of the approval claim, and holding a deployment-wide slot + // while an indexed saver read stalls would unnecessarily block other schedules + // and lengthen the Mongo-claim -> approval-CAS hand-off window below. + const checkpointGeneration = await checkpointGenerationPromise; + + // A pause frees its scheduled-run capacity slot. Before consuming the approval, + // atomically promote the run row back to `started` and claim a fresh global slot. + // The database's partial unique indexes arbitrate both deployment capacity and a + // concurrent active occurrence of the same schedule. + let scheduleCapacitySlot; + let scheduleResumeClaimToken; + let scheduleResumeLeaseBy; + const scheduleResumeOptions = { + expectedConfigRevision: job.metadata?.scheduleConfigRevision, + automatic: job.metadata?.scheduleManual !== true, + }; + if (scheduleId) { + let scheduleClaim; + try { + scheduleClaim = await claimScheduleResume(scheduleId, scheduledFor, scheduleResumeOptions); + } catch (err) { + await decrementPendingRequest(userId); + logger.error('[ResumeAgentController] Failed to claim scheduled resume capacity', err); + return sendGenerationJson( + res, + 500, + { error: 'Failed to reserve scheduled-run capacity' }, + generationProtocolVersion, + ); + } + if ('conflict' in scheduleClaim) { + await decrementPendingRequest(userId); + if (scheduleClaim.conflict === 'capacity' || scheduleClaim.conflict === 'overlap') { + res.set('Retry-After', '1'); + return sendGenerationJson( + res, + 429, + { + code: + scheduleClaim.conflict === 'capacity' + ? 'SCHEDULE_CAPACITY' + : 'SCHEDULE_OCCURRENCE_ACTIVE', + error: + scheduleClaim.conflict === 'capacity' + ? 'Scheduled-run capacity is currently full. Please retry.' + : 'Another occurrence of this schedule is still running. Please retry.', + }, + generationProtocolVersion, + ); + } + return sendGenerationJson( + res, + 409, + { + code: + scheduleClaim.conflict === 'inactive' + ? 'SCHEDULE_NO_LONGER_ACTIVE' + : 'SCHEDULE_RUN_NOT_PAUSED', + error: 'This scheduled run can no longer be resumed', + }, + generationProtocolVersion, + ); + } + scheduleCapacitySlot = scheduleClaim.capacitySlot; + scheduleResumeClaimToken = scheduleClaim.claimToken; + scheduleResumeLeaseBy = scheduleClaim.leaseBy; + } + + const releaseScheduleFence = async () => { + if (scheduleId == null || scheduleResumeLeaseBy == null) { + return; + } + try { + await releaseScheduleResumeFence(scheduleId, scheduleResumeLeaseBy); + } catch (releaseError) { + logger.warn('[ResumeAgentController] Failed to release scheduled resume fence', releaseError); + } + }; + + /** Release only when the exact generation demonstrably remains paused. If the + * approval CAS reply is ambiguous and the job cannot be read, retaining the slot + * until reconciliation is the safe direction: releasing it could exceed the cap + * while a committed continuation is already running. */ + const rollbackUnconsumedScheduleClaim = async (currentJob) => { + if ( + scheduleId == null || + scheduleCapacitySlot == null || + currentJob?.createdAt !== job.createdAt || + currentJob?.status !== 'requires_action' + ) { + return; + } + try { + await releaseScheduleResumeClaim(scheduleId, scheduledFor, scheduleCapacitySlot); + } catch (rollbackError) { + logger.warn( + '[ResumeAgentController] Failed to release unconsumed scheduled resume capacity', + rollbackError, + ); + } + }; + // Atomically claim the resume. The single winner drives the run; a racing second // submit (double-click, two tabs) gets false and must not re-drive — that would // re-execute tools and double-bill. @@ -729,10 +913,8 @@ const ResumeAgentController = async (req, res, next, initializeClient, addTitle) // would leak the concurrency slot until the counter TTL expires — spuriously 429'ing // the user when they retry the still-paused approval. Release the slot on that path too. let claimed; - let checkpointGeneration; const providerExecutionId = randomUUID(); try { - checkpointGeneration = await checkpointGenerationPromise; /** The CAS that reopens steering must also publish THIS owner's seal * capability. A separate write after status=`running` leaves a window in * which steer/arm requests read the previous replica's capability. */ @@ -748,6 +930,9 @@ const ResumeAgentController = async (req, res, next, initializeClient, addTitle) job.createdAt, ); } catch (err) { + const currentJob = await GenerationJobManager.getJob(streamId).catch(() => null); + await rollbackUnconsumedScheduleClaim(currentJob); + await releaseScheduleFence(); await decrementPendingRequest(userId); logger.error('[ResumeAgentController] Failed to claim resume', err); return sendGenerationJson(res, 500, { error: 'Failed to resume' }, generationProtocolVersion); @@ -755,6 +940,8 @@ const ResumeAgentController = async (req, res, next, initializeClient, addTitle) if (!claimed) { await decrementPendingRequest(userId); const currentJob = await GenerationJobManager.getJob(streamId).catch(() => null); + await rollbackUnconsumedScheduleClaim(currentJob); + await releaseScheduleFence(); if (currentJob != null && currentJob.createdAt !== job.createdAt) { return sendGenerationJson(res, 409, { code: 'RUN_REPLACED' }, generationProtocolVersion); } @@ -766,6 +953,73 @@ const ResumeAgentController = async (req, res, next, initializeClient, addTitle) ); } + // Linearize the consumed approval against the schedule's live config. The schedule + // document fence was acquired only after all async policy reads, and this atomic + // consume checks its token/revision/enabled state immediately after the approval CAS. + // An edit/disable that won first makes this fail; one that lands afterward is ordered + // after the continuation has started. Never begin provider execution on a stale claim. + if (scheduleId) { + let scheduleClaimCurrent = false; + try { + scheduleClaimCurrent = await finalizeScheduleResumeClaim( + scheduleId, + scheduleResumeClaimToken, + scheduleResumeLeaseBy, + scheduleResumeOptions, + ); + } catch (error) { + logger.error('[ResumeAgentController] Failed to finalize scheduled resume fence', error); + await releaseScheduleFence(); + } + if (!scheduleClaimCurrent) { + await decrementPendingRequest(userId); + let stopped = false; + try { + const abortResult = await GenerationJobManager.abortJob(streamId, { + expectedCreatedAt: job.createdAt, + awaitProviderDrain: true, + }); + stopped = abortResult != null && abortResult.failureReason == null; + } catch (error) { + logger.warn('[ResumeAgentController] Failed to stop stale scheduled resume', error); + } + if (!stopped) { + res.set('Retry-After', '1'); + return sendGenerationJson( + res, + 503, + { + code: 'SCHEDULE_STOP_UNCONFIRMED', + error: 'The stale scheduled resume could not be confirmed stopped.', + }, + generationProtocolVersion, + ); + } + await recordScheduleOutcome({ + scheduleId, + scheduledFor, + streamId, + jobCreatedAt: job.createdAt, + status: 'interrupted', + conversationId, + error: 'Schedule was disabled, changed, or deleted before approval', + }); + if (checkpointNamespace !== '') { + await deleteAgentCheckpoint(conversationId, checkpointerCfg, undefined, { + checkpointNamespace, + }).catch((error) => { + logger.warn('[ResumeAgentController] Failed to prune stale schedule checkpoint', error); + }); + } + return sendGenerationJson( + res, + 409, + { code: 'SCHEDULE_NO_LONGER_ACTIVE', error: 'This schedule can no longer be resumed' }, + generationProtocolVersion, + ); + } + } + /** * An interrupt steer enqueued just before the pause survives durably with * its `preempt` flag, but the ARM lived only in the previous owner's @@ -1018,6 +1272,16 @@ const ResumeAgentController = async (req, res, next, initializeClient, addTitle) `[ResumeAgentController] Re-pause persistence barrier changed before release: ${streamId}`, ); } + if (scheduleId) { + await recordScheduleOutcome({ + scheduleId, + scheduledFor, + streamId, + jobCreatedAt: job.createdAt, + status: 'requires_action', + conversationId, + }); + } } else { logger.debug( `[ResumeAgentController] Skipping stale re-pause persistence — ${streamId} no longer owns its barrier`, @@ -1027,11 +1291,25 @@ const ResumeAgentController = async (req, res, next, initializeClient, addTitle) } // If the user aborted mid-resume, the abort route already emitted the terminal - // event and finalized the job — don't double-save / double-finalize here. + // event and finalized the job — don't double-save / double-finalize here. This + // continuation is nevertheless the scheduled-run owner, so it must settle the + // run row after observing its own abort; the generic Stop route deliberately + // delegates a running generation's settlement to that generation owner. if (job.abortController.signal.aborted) { logger.debug( `[ResumeAgentController] Aborted during resume; abort route finalizes: ${streamId}`, ); + if (scheduleId) { + await recordScheduleOutcome({ + scheduleId, + scheduledFor, + streamId, + jobCreatedAt: job.createdAt, + status: 'interrupted', + conversationId, + error: 'Scheduled run was stopped', + }); + } return; } @@ -1061,6 +1339,17 @@ const ResumeAgentController = async (req, res, next, initializeClient, addTitle) 're-pause persistence failure', ); } + if (scheduleId && pausePersistenceFailureFinalized) { + await recordScheduleOutcome({ + scheduleId, + scheduledFor, + streamId, + jobCreatedAt: job.createdAt, + status: 'error', + conversationId, + error: err?.message ?? 'Re-pause persistence failed', + }); + } return; } // Job-replacement guard (mirrors finalizeResumedTurn's success-path guard): if a @@ -1104,6 +1393,18 @@ const ResumeAgentController = async (req, res, next, initializeClient, addTitle) 'failed resume finalization', ); } + if (scheduleId && errorFinalized) { + const balanceRefusal = err?.message?.includes(ViolationTypes.TOKEN_BALANCE); + await recordScheduleOutcome({ + scheduleId, + scheduledFor, + streamId, + jobCreatedAt: job.createdAt, + status: balanceRefusal ? 'skipped_balance' : 'error', + conversationId, + ...(!balanceRefusal && { error: err?.message ?? 'Resume failed' }), + }); + } } } finally { try { diff --git a/api/server/experimental.js b/api/server/experimental.js index 698cb52ff26..8d03f8ae59e 100644 --- a/api/server/experimental.js +++ b/api/server/experimental.js @@ -28,6 +28,7 @@ const { setupGracefulShutdown, configureMessageFilterRegexValidator, configureFileConfigRegexEngine, + GenerationJobManager, waitForKeyvRedisClient, } = require('@librechat/api'); const { connectDb, indexSync } = require('~/db'); @@ -37,6 +38,10 @@ const createValidateImageRequest = require('./middleware/validateImageRequest'); const { startExpiredFileSweep } = require('./services/Files/process'); const { initializeGitHubSkillSync } = require('./services/Skills/sync'); const { initializeAgentTriggerService } = require('./services/Agents/triggers'); +const { + recordExpiredScheduleApproval, + initializeScheduleErasureSweep, +} = require('./services/Schedules'); const { configureSubagentTaskRouting } = require('./services/Endpoints/agents/subagentThreadStore'); const { jwtLogin, ldapLogin, passportLogin } = require('~/strategies'); const { updateInterfacePermissions: updateInterfacePerms } = require('@librechat/api'); @@ -276,6 +281,11 @@ if (cluster.isMaster) { * Each worker runs a full Express server instance */ const app = express(); + // The clustered entrypoint deliberately does not arm the v1 schedule engine, + // but an already-fired scheduled generation can still reach HITL here. Settle + // its durable run when the generic approval runtime expires it. + GenerationJobManager.setApprovalExpiredHandler(recordExpiredScheduleApproval); + GenerationJobManager.initialize(); /** * The master may assign the sweep worker before or after this worker has * loaded app config. These flags join the IPC assignment with config @@ -284,6 +294,17 @@ if (cluster.isMaster) { let shouldStartExpiredFileSweep = false; let expiredFileSweepOptions = null; let expiredFileSweepStarted = false; + const SCHEDULE_ENGINE_OPTIONAL_METHODS = new Set(['GET', 'HEAD', 'OPTIONS', 'DELETE']); + + const rejectScheduleWritesUntilReady = (req, res, next) => { + if (SCHEDULE_ENGINE_OPTIONAL_METHODS.has(req.method)) { + return next(); + } + return res.status(501).json({ + code: 'SCHEDULES_NOT_SUPPORTED', + error: 'Scheduled chats are not available in clustered mode.', + }); + }; const startExpiredFileSweepOnce = () => { if (!shouldStartExpiredFileSweep || expiredFileSweepStarted || !expiredFileSweepOptions) { @@ -322,6 +343,13 @@ if (cluster.isMaster) { logger.error(`[Worker ${process.pid}][indexSync] Background sync failed:`, err); }); + // This entrypoint deliberately does not arm the schedule engine, but DELETE stays + // open — so soft-deleted rows still accrue with no reconciler to erase them. Start + // the erasure-ONLY sweep (Mongo is up, GenerationJobManager was initialized above): + // it never claims, fires, advances, or infers owner death from a process-local + // missing job, and its idempotent guard makes this safe once per worker. + initializeScheduleErasureSweep(); + app.disable('x-powered-by'); app.set('trust proxy', trusted_proxy); @@ -488,6 +516,7 @@ if (cluster.isMaster) { app.use('/api/agents', routes.agents); app.use('/api/banner', routes.banner); app.use('/api/memories', routes.memories); + app.use('/api/schedules', rejectScheduleWritesUntilReady, routes.schedules); app.use('/api/permissions', routes.accessPermissions); app.use('/api/tags', routes.tags); app.use('/api/mcp', routes.mcp); diff --git a/api/server/experimental.spec.js b/api/server/experimental.spec.js index 4ff9d6b924a..58349d3091d 100644 --- a/api/server/experimental.spec.js +++ b/api/server/experimental.spec.js @@ -23,6 +23,34 @@ describe('Experimental server configuration', () => { expect(source).toMatch(/if \(shuttingDown\) \{[\s\S]*?return;[\s\S]*?Starting a new worker/); }); + it('starts approval expiry after installing the scheduled-run callback', () => { + const handlerIndex = source.indexOf( + 'GenerationJobManager.setApprovalExpiredHandler(recordExpiredScheduleApproval);', + ); + const initializeIndex = source.indexOf('GenerationJobManager.initialize();'); + + expect(handlerIndex).toBeGreaterThan(-1); + expect(initializeIndex).toBeGreaterThan(handlerIndex); + }); + + it('starts erasure-only schedule maintenance after connecting to Mongo, once per worker', () => { + const connectIndex = source.indexOf('await connectDb();'); + const sweepIndex = source.indexOf('initializeScheduleErasureSweep();'); + + expect(connectIndex).toBeGreaterThan(-1); + expect(sweepIndex).toBeGreaterThan(-1); + // Mongo must be up before the sweep reads soft-deleted rows. + expect(sweepIndex).toBeGreaterThan(connectIndex); + // Idempotent guard lives in the service; started exactly once from this entrypoint. + expect(source.match(/initializeScheduleErasureSweep\(\);/g)).toHaveLength(1); + }); + + it('never arms the full schedule engine in a clustered worker', () => { + // The clustered entrypoint runs erasure-only maintenance: arming the engine here + // would claim/fire/absence-reconcile runs whose peer generations it cannot see. + expect(source).not.toContain('initializeScheduleEngine('); + }); + it('runs cross-tenant startup work in the system context', () => { expect(source).toContain('await runAsSystem(seedDatabase);'); expect(source).toMatch( diff --git a/api/server/index.js b/api/server/index.js index 0768a27c50d..9203250d70e 100644 --- a/api/server/index.js +++ b/api/server/index.js @@ -57,6 +57,7 @@ const { capabilityContextMiddleware } = require('./middleware/roles/capabilities const createValidateImageRequest = require('./middleware/validateImageRequest'); const { initializeGitHubSkillSync } = require('./services/Skills/sync'); const { initializeAgentTriggerService } = require('./services/Agents/triggers'); +const { initializeScheduleEngine, recordExpiredScheduleApproval } = require('./services/Schedules'); const { jwtLogin, ldapLogin, passportLogin } = require('~/strategies'); const { startExpiredFileSweep } = require('./services/Files/process'); const { checkMigrations } = require('./services/start/migration'); @@ -85,9 +86,12 @@ const trusted_proxy = Number(TRUST_PROXY) || 1; /* trust first proxy by default const app = express(); let serverReady = false; +let schedulesReady = false; const SERVER_NOT_READY_CODE = 'SERVER_NOT_READY'; const CHAT_START_RETRY_AFTER_SECONDS = '1'; +const SCHEDULES_NOT_READY_CODE = 'SCHEDULES_NOT_READY'; +const SCHEDULE_ENGINE_OPTIONAL_METHODS = new Set(['GET', 'HEAD', 'OPTIONS', 'DELETE']); const rejectChatStartsUntilReady = (req, res, next) => { if (serverReady || req.method !== 'POST' || req.path === '/abort') { @@ -101,12 +105,24 @@ const rejectChatStartsUntilReady = (req, res, next) => { }); }; +const rejectScheduleWritesUntilReady = (req, res, next) => { + if (schedulesReady || SCHEDULE_ENGINE_OPTIONAL_METHODS.has(req.method)) { + return next(); + } + res.set('Retry-After', CHAT_START_RETRY_AFTER_SECONDS); + return res.status(503).json({ + code: SCHEDULES_NOT_READY_CODE, + error: 'Scheduler is still starting. Please retry shortly.', + }); +}; + const configureGenerationStreams = () => { const streamServices = createStreamServices(); GenerationJobManager.configure({ ...streamServices, cleanupOnComplete: !isEnabled(process.env.STREAM_KEEP_COMPLETED_JOBS), }); + GenerationJobManager.setApprovalExpiredHandler(recordExpiredScheduleApproval); GenerationJobManager.initialize(); // Stop active generations and close their SSE streams while the HTTP server drains. registerShutdownTask( @@ -345,6 +361,7 @@ const startServer = async () => { app.use('/api/agents', routes.agents); app.use('/api/banner', routes.banner); app.use('/api/memories', routes.memories); + app.use('/api/schedules', rejectScheduleWritesUntilReady, routes.schedules); app.use('/api/permissions', routes.accessPermissions); app.use('/api/tags', routes.tags); @@ -402,6 +419,10 @@ const startServer = async () => { memoryDiagnostics.start(); } await initializeAgentTriggerService({ address: server.address() }); + schedulesReady = (await initializeScheduleEngine()) != null; + if (!schedulesReady) { + logger.warn('[schedules] write routes remain unavailable because the engine did not arm.'); + } serverReady = true; logger.info('Server readiness checks passing.'); } catch (initErr) { diff --git a/api/server/index.metrics.spec.js b/api/server/index.metrics.spec.js index 6b95104c3de..80a386c6079 100644 --- a/api/server/index.metrics.spec.js +++ b/api/server/index.metrics.spec.js @@ -38,6 +38,10 @@ jest.mock('~/server/services/Agents/triggers', () => ({ initializeAgentTriggerService: jest.fn().mockResolvedValue(undefined), })); +jest.mock('~/server/services/Schedules', () => ({ + initializeScheduleEngine: jest.fn().mockResolvedValue(undefined), +})); + describe('Server metrics route', () => { jest.setTimeout(30_000); diff --git a/api/server/index.spec.js b/api/server/index.spec.js index 013ece204ef..73ad042865e 100644 --- a/api/server/index.spec.js +++ b/api/server/index.spec.js @@ -39,6 +39,10 @@ jest.mock('~/server/services/Agents/triggers', () => ({ initializeAgentTriggerService: jest.fn().mockResolvedValue(undefined), })); +jest.mock('~/server/services/Schedules', () => ({ + initializeScheduleEngine: jest.fn().mockResolvedValue(undefined), +})); + jest.mock( '@librechat/api/telemetry', () => ({ diff --git a/api/server/routes/__tests__/mcp.spec.js b/api/server/routes/__tests__/mcp.spec.js index 44dcfe26288..99ebff511a7 100644 --- a/api/server/routes/__tests__/mcp.spec.js +++ b/api/server/routes/__tests__/mcp.spec.js @@ -228,6 +228,12 @@ describe('MCP Routes', () => { currentUser = undefined; mockResolveAllMcpConfigs.mockResolvedValue({}); mockResolveMcpConfigNames.mockResolvedValue([]); + // `clearAllMocks` preserves queued `mockReturnValueOnce` entries. A callback + // test can legitimately leave one unconsumed, which then masks a later test's + // default implementation when this file shares a CI shard with other suites. + // Reset the cache accessor itself so every test starts from a deterministic + // default and can opt into a one-shot value deliberately. + require('~/cache').getLogStores.mockReset().mockReturnValue({}); const { MCPOAuthHandler } = require('@librechat/api'); const { getTenantId } = require('@librechat/data-schemas'); getTenantId.mockReturnValue(undefined); diff --git a/api/server/routes/agents/__tests__/abort.spec.js b/api/server/routes/agents/__tests__/abort.spec.js index f6a3ca559ff..43516fbb161 100644 --- a/api/server/routes/agents/__tests__/abort.spec.js +++ b/api/server/routes/agents/__tests__/abort.spec.js @@ -25,6 +25,10 @@ const mockGenerationJobManager = { const mockSaveMessage = jest.fn(); +const mockRecordScheduleOutcome = jest.fn(); +const mockBeginScheduledStop = jest.fn(); +const mockAcknowledgeScheduledStopPersistence = jest.fn(); + jest.mock('@librechat/data-schemas', () => ({ ...jest.requireActual('@librechat/data-schemas'), logger: mockLogger, @@ -33,6 +37,11 @@ jest.mock('@librechat/data-schemas', () => ({ jest.mock('@librechat/api', () => ({ ...jest.requireActual('@librechat/api'), isEnabled: jest.fn().mockReturnValue(false), + isAgentTriggerRequest: jest.fn(() => false), + captureScheduleFireContext: jest.fn((req) => { + req._isScheduledFire = false; + req._isManualScheduledFire = false; + }), GenerationJobManager: mockGenerationJobManager, })); @@ -40,6 +49,13 @@ jest.mock('~/models', () => ({ saveMessage: (...args) => mockSaveMessage(...args), })); +jest.mock('~/server/services/Schedules', () => ({ + recordScheduleOutcome: (...args) => mockRecordScheduleOutcome(...args), + beginScheduledStop: (...args) => mockBeginScheduledStop(...args), + acknowledgeScheduledStopPersistence: (...args) => + mockAcknowledgeScheduledStopPersistence(...args), +})); + jest.mock('~/server/middleware', () => ({ uaParser: (req, res, next) => next(), checkBan: (req, res, next) => next(), @@ -81,6 +97,12 @@ describe('Agent Abort Endpoint', () => { mockGenerationJobManager.getActiveJobIdsForUser.mockReset(); mockSaveMessage.mockReset(); mockSaveMessage.mockImplementation(async (_context, message) => message); + mockRecordScheduleOutcome.mockReset(); + mockRecordScheduleOutcome.mockResolvedValue(true); + mockBeginScheduledStop.mockReset(); + mockBeginScheduledStop.mockResolvedValue(true); + mockAcknowledgeScheduledStopPersistence.mockReset(); + mockAcknowledgeScheduledStopPersistence.mockResolvedValue(undefined); }); describe('POST /chat/abort', () => { @@ -841,5 +863,140 @@ describe('Agent Abort Endpoint', () => { }); }); }); + + describe('Scheduled Stop persistence protocol', () => { + const scheduledJob = { + status: 'running', + createdAt: 111, + metadata: { + userId: 'test-user-123', + generationProtocolVersion: 2, + scheduleId: 's1', + scheduledFor: '2026-01-01T00:00:00.000Z', + conversationId: 'conv-1', + }, + }; + + it('stamps the Stop before signalling abort, then acknowledges after persistence', async () => { + mockGenerationJobManager.getJob.mockResolvedValue(scheduledJob); + mockGenerationJobManager.abortJob.mockResolvedValue({ + success: true, + jobData: { conversationId: 'conv-1' }, + content: [], + text: '', + }); + + const response = await request(app) + .post('/api/agents/chat/abort') + .set('X-LibreChat-Generation-Protocol', '2') + .send({ conversationId: 'conv-1', generationProtocolVersion: 2 }); + + expect(response.status).toBe(200); + expect(mockBeginScheduledStop).toHaveBeenCalledWith({ + scheduleId: 's1', + scheduledFor: '2026-01-01T00:00:00.000Z', + }); + // Stamp BEFORE the abort signal; acknowledgement AFTER it (persistence done). + expect(mockBeginScheduledStop.mock.invocationCallOrder[0]).toBeLessThan( + mockGenerationJobManager.abortJob.mock.invocationCallOrder[0], + ); + // The ack also carries a terminal outcome to re-drive: the owner calls + // recordScheduleOutcome once, and if that call's Stop barrier deferred (slow + // beforePublish), nothing would settle the run where no reconciler is armed. + expect(mockAcknowledgeScheduledStopPersistence).toHaveBeenCalledWith( + expect.objectContaining({ + scheduleId: 's1', + scheduledFor: '2026-01-01T00:00:00.000Z', + settle: expect.objectContaining({ status: 'interrupted' }), + }), + ); + expect(mockGenerationJobManager.abortJob.mock.invocationCallOrder[0]).toBeLessThan( + mockAcknowledgeScheduledStopPersistence.mock.invocationCallOrder[0], + ); + }); + + it('returns STOP_IN_PROGRESS and never signals a second abort when a Stop is already live', async () => { + mockGenerationJobManager.getJob.mockResolvedValue(scheduledJob); + mockBeginScheduledStop.mockResolvedValue('in_progress'); + + const response = await request(app) + .post('/api/agents/chat/abort') + .set('X-LibreChat-Generation-Protocol', '2') + .send({ conversationId: 'conv-1', generationProtocolVersion: 2 }); + + expect(response.status).toBe(409); + expect(response.body.code).toBe('STOP_IN_PROGRESS'); + expect(mockGenerationJobManager.abortJob).not.toHaveBeenCalled(); + expect(mockAcknowledgeScheduledStopPersistence).not.toHaveBeenCalled(); + }); + + it('does NOT acknowledge when the abort persistence failed (run stays preserved)', async () => { + mockGenerationJobManager.getJob.mockResolvedValue(scheduledJob); + mockGenerationJobManager.abortJob.mockResolvedValue({ + success: true, + persistenceFailed: true, + jobData: { conversationId: 'conv-1' }, + content: [], + text: '', + }); + + const response = await request(app) + .post('/api/agents/chat/abort') + .set('X-LibreChat-Generation-Protocol', '2') + .send({ conversationId: 'conv-1', generationProtocolVersion: 2 }); + + expect(response.status).toBe(200); + expect(response.body.persistenceFailed).toBe(true); + expect(mockBeginScheduledStop).toHaveBeenCalled(); + expect(mockAcknowledgeScheduledStopPersistence).not.toHaveBeenCalled(); + }); + + it('does not re-drive settlement for a PAUSED job, which settles explicitly', async () => { + mockGenerationJobManager.getJob.mockResolvedValue({ + ...scheduledJob, + status: 'requires_action', + }); + mockGenerationJobManager.abortJob.mockResolvedValue({ + success: true, + jobData: { conversationId: 'conv-1' }, + content: [], + text: '', + }); + + const response = await request(app) + .post('/api/agents/chat/abort') + .set('X-LibreChat-Generation-Protocol', '2') + .send({ conversationId: 'conv-1', generationProtocolVersion: 2 }); + + expect(response.status).toBe(200); + const ack = mockAcknowledgeScheduledStopPersistence.mock.calls[0][0]; + expect(ack.settle).toBeUndefined(); + // The paused path records its own outcome explicitly. + expect(mockRecordScheduleOutcome).toHaveBeenCalled(); + }); + + it('leaves non-scheduled aborts untouched by the Stop protocol', async () => { + mockGenerationJobManager.getJob.mockResolvedValue({ + status: 'running', + createdAt: 1, + metadata: { userId: 'test-user-123', generationProtocolVersion: 2 }, + }); + mockGenerationJobManager.abortJob.mockResolvedValue({ + success: true, + jobData: { conversationId: 'conv-x' }, + content: [], + text: '', + }); + + const response = await request(app) + .post('/api/agents/chat/abort') + .set('X-LibreChat-Generation-Protocol', '2') + .send({ conversationId: 'conv-x', generationProtocolVersion: 2 }); + + expect(response.status).toBe(200); + expect(mockBeginScheduledStop).not.toHaveBeenCalled(); + expect(mockAcknowledgeScheduledStopPersistence).not.toHaveBeenCalled(); + }); + }); }); }); diff --git a/api/server/routes/agents/index.js b/api/server/routes/agents/index.js index ee3ed4e7bbb..1d73ea5d985 100644 --- a/api/server/routes/agents/index.js +++ b/api/server/routes/agents/index.js @@ -15,6 +15,8 @@ const { createMessageFilterPii, isAgentTriggerRequest, exemptAgentTriggerFromIpLimiter, + captureScheduleFireContext, + exemptFromUserLimiter: exemptScheduleFromUserLimiter, } = require('@librechat/api'); const { createSseStreamTelemetry } = require('@librechat/api/telemetry'); const { logger } = require('@librechat/data-schemas'); @@ -36,6 +38,11 @@ const { negotiateExistingGenerationProtocol, } = require('~/server/controllers/agents/protocol'); const { saveMessage } = require('~/models'); +const { + recordScheduleOutcome, + beginScheduledStop, + acknowledgeScheduledStopPersistence, +} = require('~/server/services/Schedules'); const responses = require('./responses'); const openai = require('./openai'); const { v1 } = require('./v1'); @@ -119,6 +126,7 @@ router.use(requireJwtAuth); // middleware reads this stable decision instead of re-verifying an expired token. router.use((req, _res, next) => { req._isAgentTrigger = isAgentTriggerRequest(req); + captureScheduleFireContext(req); next(); }); router.use(checkBan); @@ -669,6 +677,27 @@ router.post('/chat/abort', configMiddleware, async (req, res, next) => { throwOnError: true, }) : undefined; + // Stamp a scheduled run's Stop BEFORE signalling the abort. `abortJob` flips the job + // to `aborted` immediately, then runs the partial-message/checkpoint persistence in + // `beforePublish`. Without this stamp the owner settlement, reconciliation, and + // schedule/account deletion could observe `aborted` and terminalize/erase the run + // mid-write. The stamp is serialized: a fresh Stop already owning it means another + // request is persisting, so we must not signal a second abort. + const stopScheduleId = job.metadata?.scheduleId; + const stopScheduledFor = job.metadata?.scheduledFor; + const isScheduledStop = stopScheduleId != null && stopScheduledFor != null; + let scheduledStopStamped = false; + if (isScheduledStop) { + const stopStamp = await beginScheduledStop({ + scheduleId: stopScheduleId, + scheduledFor: stopScheduledFor, + }); + if (stopStamp === 'in_progress') { + res.set('Retry-After', '1'); + return res.status(409).json({ code: 'STOP_IN_PROGRESS', generationProtocolVersion }); + } + scheduledStopStamped = stopStamp === true; + } const abortResult = await GenerationJobManager.abortJob(jobStreamId, { expectedCreatedAt: job.createdAt, transformAbortContent: (content, abortJobData) => { @@ -807,6 +836,18 @@ router.post('/chat/abort', configMiddleware, async (req, res, next) => { } }, }); + // The abort did not land (replaced/still-active/already-settled), so no persistence + // is in flight: release the Stop barrier we armed rather than deferring this + // occurrence's settlement for the full stale-owner window. A retry or a replacement + // generation re-stamps its own; the acknowledgement is fenced to the occurrence, not + // a generation, so it cannot settle a successor through its predecessor. + if (scheduledStopStamped && !abortResult.success) { + await acknowledgeScheduledStopPersistence({ + scheduleId: stopScheduleId, + scheduledFor: stopScheduledFor, + }); + scheduledStopStamped = false; + } if (abortResult.failureReason === 'generation_replaced') { return res.status(409).json({ code: 'RUN_REPLACED', generationProtocolVersion }); } @@ -858,6 +899,53 @@ router.post('/chat/abort', configMiddleware, async (req, res, next) => { abortResultResponseMessageId: abortResult.jobData?.responseMessageId, }); + // `beforePublish` has run: its partial-message and checkpoint writes have either + // landed or failed. Acknowledge the Stop ONLY on success — that releases the owner's + // settlement barrier so the run can terminalize. On a persistence failure we leave + // the barrier unresolved: the run stays preserved (client retries; the stale-owner + // timeout is the bounded recovery) rather than settling over an incomplete write. + if (scheduledStopStamped && !abortResult.persistenceFailed) { + await acknowledgeScheduledStopPersistence({ + scheduleId: stopScheduleId, + scheduledFor: stopScheduledFor, + // Re-drive the terminal outcome from here for a RUNNING generation: its owner + // calls recordScheduleOutcome once, and if that call's Stop barrier deferred + // (slow beforePublish), nothing would settle the run where no schedule + // reconciler is armed. recordRunOutcome is match-guarded and idempotent, so an + // owner that already settled makes this a no-op. A paused job is settled + // explicitly below and needs no re-drive here. + ...(job.status !== 'requires_action' && { + settle: { + status: 'interrupted', + conversationId: job.metadata?.conversationId ?? jobStreamId, + error: 'Scheduled run was stopped', + }, + }), + }); + } + + // A paused generation has no live provider owner left to report the stop. + // Persist its scheduled occurrence here, after abortJob's required partial + // response/checkpoint work AND its acknowledgement above, while running generations + // continue to settle from their owning request/resume controller. A failed + // persistence skips settlement so the incomplete run is not terminalized. + if ( + job.status === 'requires_action' && + job.metadata?.scheduleId && + !abortResult.persistenceFailed + ) { + await recordScheduleOutcome({ + scheduleId: job.metadata.scheduleId, + scheduledFor: job.metadata.scheduledFor, + streamId: jobStreamId, + jobCreatedAt: job.createdAt, + status: 'interrupted', + conversationId: job.metadata.conversationId ?? jobStreamId, + clearConversationId: abortResult.jobData?.createdEventEmitted !== true, + error: 'Scheduled run was stopped while awaiting approval', + }); + } + if (abortResult.persistenceFailed && generationProtocolVersion < GENERATION_PROTOCOL_V2) { res.set('Retry-After', '1'); return res.status(409).json({ @@ -977,7 +1065,7 @@ if (isEnabled(LIMIT_MESSAGE_IP)) { } if (isEnabled(LIMIT_MESSAGE_USER)) { - chatRouter.use(messageUserLimiter); + chatRouter.use(unless(exemptScheduleFromUserLimiter, messageUserLimiter)); } chatRouter.use('/', chat); diff --git a/api/server/routes/index.js b/api/server/routes/index.js index 96baaf92b6d..6b305686b6b 100644 --- a/api/server/routes/index.js +++ b/api/server/routes/index.js @@ -17,6 +17,7 @@ const memories = require('./memories'); const presets = require('./presets'); const projects = require('./projects'); const prompts = require('./prompts'); +const schedules = require('./schedules'); const skills = require('./skills'); const balance = require('./balance'); const actions = require('./actions'); @@ -69,6 +70,7 @@ module.exports = { models, prompts, projects, + schedules, skills, actions, presets, diff --git a/api/server/routes/schedules.js b/api/server/routes/schedules.js new file mode 100644 index 00000000000..b426710990f --- /dev/null +++ b/api/server/routes/schedules.js @@ -0,0 +1,102 @@ +const express = require('express'); +const { Permissions, PermissionTypes } = require('librechat-data-provider'); +const { + isEnabled, + SCHEDULE_FILE_HOLD, + generateCheckAccess, + createSchedulesHandlers, +} = require('@librechat/api'); +const { requireJwtAuth, configMiddleware, messageIpLimiter } = require('~/server/middleware'); +const { + getLimits, + fireScheduleNow, + deleteScheduleForOwner, + isUserDeleting, +} = require('~/server/services/Schedules'); +const { resolveAgentFireAccess } = require('~/server/services/Schedules/access'); +const methods = require('~/models'); + +const { getRoleByName } = methods; + +const router = express.Router(); +router.use(requireJwtAuth); +router.use(configMiddleware); + +const checkSchedulesAccess = generateCheckAccess({ + permissionType: PermissionTypes.SCHEDULES, + permissions: [Permissions.USE], + getRoleByName, +}); +const checkSchedulesCreate = generateCheckAccess({ + permissionType: PermissionTypes.SCHEDULES, + permissions: [Permissions.USE, Permissions.CREATE], + getRoleByName, +}); + +const handlers = createSchedulesHandlers({ + methods, + getLimits, + // Full fire-equivalent access check (not mere existence): the body-based agent + // middleware is a no-op for schedule payloads (no `endpoint: 'agents'`), so + // enforce the same role AGENTS:USE + resource VIEW (with manage:agents bypass) + // the actual fire requires — otherwise a role without AGENTS:USE could schedule + // runs the chat route rejects and walks toward auto-disable. + canViewAgent: async (agentId, req) => (await resolveAgentFireAccess(agentId, req.user)) === 'ok', + filterOwnedFileIds: async (fileIds, userId) => { + const files = await methods.getFiles({ file_id: { $in: fileIds }, user: userId }, null, { + file_id: 1, + }); + return (files ?? []).map((file) => file.file_id); + }, + markFilesUsed: async (fileIds, userId) => { + // BOUNDED renewable hold (extendFilesTTL), not a permanent `$unset` of the upload + // TTL: permanence made a schedule deleted before its first run, an edit that + // replaced file_ids, or a failed creation leak the upload forever, since nothing + // ever restored an expiry. The first fire that actually SENDS the file clears its + // TTL through the ordinary consumption path; until then the hold is renewed at + // create/edit and each fire preflight, and lapses when the schedule stops touching + // it. Files already made permanent by a real send are skipped by construction. + // Idempotent, and never touches the usage counter (a retention is not a consumption). + // Then VERIFY every requested file still exists: one can be deleted between the + // ownership check and here, and a silent success would persist a schedule whose + // attachments the first fire drops. Existence is the check — not the hold's + // modified-count, which reads 0 for an already-permanent or already-held file. + const unique = [...new Set(fileIds)]; + await methods.extendFilesTTL(unique, SCHEDULE_FILE_HOLD, { user: userId }); + const files = await methods.getFiles({ file_id: { $in: unique }, user: userId }, null, { + file_id: 1, + }); + const present = (files ?? []).length; + if (present !== unique.length) { + throw new Error(`attachment retention incomplete: ${present}/${unique.length} files exist`); + } + }, + fireNow: fireScheduleNow, + // Quiesce-then-erase delete: stops new claims, settles provably job-less runs + // synchronously, aborts live ones, and reports honestly (see ScheduleDeleteResult); + // a delivered abort erases on the generation's own outcome write, in any topology. + deleteSchedule: deleteScheduleForOwner, + // Durable account-deletion barrier. A one-shot disable scan cannot close the + // create race, so every scheduling WRITE consults the user-level flag instead. + isUserDeleting, +}); + +router.get('/', checkSchedulesAccess, handlers.listSchedules); +router.get('/:id', checkSchedulesAccess, handlers.getSchedule); +router.post('/', checkSchedulesCreate, handlers.createSchedule); +router.patch('/:id', checkSchedulesCreate, handlers.updateSchedule); +router.delete('/:id', checkSchedulesCreate, handlers.deleteSchedule); +// Run-now mutates runtime state; gate it on CREATE like the UI does (not USE). +// This is also where LIMIT_MESSAGE_IP has to apply to a manual run: the fire itself is a +// loopback POST carrying the server's address, so limiting by IP there would pool every +// user into one bucket. Here the initiating client's address is still on the request. +// The USER limiter is NOT duplicated here; the fire token carries the authenticated id, +// so the chat router applies it to the loopback exactly once. +router.post( + '/:id/run', + ...(isEnabled(process.env.LIMIT_MESSAGE_IP) ? [messageIpLimiter] : []), + checkSchedulesCreate, + handlers.runScheduleNow, +); + +module.exports = router; diff --git a/api/server/services/Config/loadCustomConfig.js b/api/server/services/Config/loadCustomConfig.js index 2629ed1c8f0..45cec141607 100644 --- a/api/server/services/Config/loadCustomConfig.js +++ b/api/server/services/Config/loadCustomConfig.js @@ -8,7 +8,9 @@ const { logger } = require('@librechat/data-schemas'); const { configSchema, paramSettings, + EModelEndpoint, EImageOutputType, + setMaxSubagents, agentParamSettings, validateSettingDefinitions, } = require('librechat-data-provider'); @@ -109,6 +111,11 @@ async function loadCustomConfig(printConfig = true) { } } + // Applied before parsing so specs validated in the same pass (whose subagent + // presets share the cap) check against the configured limit. Invalid values + // are ignored here and rejected by the schema parse below. + setMaxSubagents(customConfig?.endpoints?.[EModelEndpoint.agents]?.maxSubagents); + const result = configSchema.strict().safeParse(customConfig); if (result?.error?.errors?.some((err) => err?.path && err.path?.includes('imageOutputType'))) { throw new Error( diff --git a/api/server/services/Schedules/access.js b/api/server/services/Schedules/access.js new file mode 100644 index 00000000000..a278b4e3398 --- /dev/null +++ b/api/server/services/Schedules/access.js @@ -0,0 +1,18 @@ +const mongoose = require('mongoose'); +const { createResolveAgentFireAccess } = require('@librechat/api'); +const { checkPermission } = require('~/server/services/PermissionService'); +const { hasCapability } = require('~/server/middleware/roles/capabilities'); +const { getRoleByName } = require('~/models'); + +// Thin wiring over the TypeScript implementation in packages/api: inject the +// api-layer lookups (agent id, role, capability, resource ACL); the authorization +// logic itself lives in @librechat/api so it stays type-checked and in-boundary. +const resolveAgentFireAccess = createResolveAgentFireAccess({ + findAgentObjectId: (agentId) => + mongoose.models.Agent.findOne({ id: agentId }).select('_id').lean(), + getRoleByName, + hasCapability, + checkPermission, +}); + +module.exports = { resolveAgentFireAccess }; diff --git a/api/server/services/Schedules/index.js b/api/server/services/Schedules/index.js new file mode 100644 index 00000000000..72f76fa7dc1 --- /dev/null +++ b/api/server/services/Schedules/index.js @@ -0,0 +1,112 @@ +let service; + +/** + * Build the schedule service on first use. Controllers import this facade in many + * focused test/runtime paths that never use scheduling; eagerly loading Config and + * the trigger worker made those paths initialize two unrelated service graphs just + * by requiring a controller. + */ +function getService() { + if (service != null) { + return service; + } + const mongoose = require('mongoose'); + const { createSchedulesService } = require('@librechat/api'); + const { getAppConfig } = require('~/server/services/Config/app'); + const { + enqueueAgentTrigger, + getAgentTriggerDelivery, + } = require('~/server/services/Agents/triggers'); + const { resolveAgentFireAccess } = require('./access'); + const methods = require('~/models'); + const isUserDeleting = async (userId) => !(await methods.isAgentTriggerPrincipalActive(userId)); + + service = createSchedulesService({ + methods, + getAppConfig, + findUserById: (userId) => + mongoose.models.User.findById(userId).select('_id tenantId role').lean(), + findBalance: (userId) => mongoose.models.Balance.findOne({ user: userId }).lean(), + upsertBalance: (userId, { set, setOnInsert }) => + mongoose.models.Balance.findOneAndUpdate( + { user: userId }, + { + ...(set && Object.keys(set).length > 0 ? { $set: set } : {}), + ...(setOnInsert && Object.keys(setOnInsert).length > 0 + ? { $setOnInsert: setOnInsert } + : {}), + }, + { upsert: true, new: true }, + ).lean(), + // Compare-and-set: only initialize an existing record while its credit is still null. + // No upsert — a CAS miss must re-read the winner, not insert a fresh balance. `null` + // in the filter also matches a legacy record whose `tokenCredits` field is absent. + initializeNullBalance: (userId, { tokenCredits, sync }) => + mongoose.models.Balance.findOneAndUpdate( + { user: userId, tokenCredits: null }, + { $set: { tokenCredits, ...(sync && Object.keys(sync).length > 0 ? sync : {}) } }, + { new: true }, + ).lean(), + enqueueAgentTrigger, + // Reconciliation reads the durable delivery to tell a still-live admission (a + // deferred Retry-After) or a dead-letter apart from a genuinely orphaned run. + getTriggerDelivery: getAgentTriggerDelivery, + resolveAgentFireAccess, + // Reuse the merged trigger/deletion admission fence. A schedule fire and every + // generic trigger now fail closed on the same durable principal state. + isUserDeleting, + }); + return service; +} + +const invoke = + (method) => + (...args) => + getService()[method](...args); + +/** Host adapter for the generic generation runtime's approval-expiry seam. The + * callback is intentionally idempotent: a replica relay or schedule reconciliation + * may re-drive the same retained terminal evidence after a transient failure. */ +async function recordExpiredScheduleApproval(streamId, job) { + if (!job?.scheduleId || !job?.scheduledFor) { + return; + } + const recorded = await getService().recordScheduleOutcome({ + scheduleId: job.scheduleId, + scheduledFor: job.scheduledFor, + streamId, + jobCreatedAt: job.createdAt, + status: 'interrupted', + conversationId: job.conversationId ?? streamId, + error: 'Approval expired before a decision was made', + }); + if (!recorded) { + throw new Error(`Failed to settle expired scheduled approval ${job.scheduleId}`); + } +} + +module.exports = { + getLimits: invoke('getLimits'), + fireScheduleNow: invoke('fireScheduleNow'), + recordScheduleOutcome: invoke('recordScheduleOutcome'), + beginScheduledStop: invoke('beginScheduledStop'), + acknowledgeScheduledStopPersistence: invoke('acknowledgeScheduledStopPersistence'), + claimScheduleResume: invoke('claimScheduleResume'), + releaseScheduleResumeClaim: invoke('releaseScheduleResumeClaim'), + finalizeScheduleResumeClaim: invoke('finalizeScheduleResumeClaim'), + releaseScheduleResumeFence: invoke('releaseScheduleResumeFence'), + isScheduleLive: invoke('isScheduleLive'), + deleteScheduleForOwner: invoke('deleteScheduleForOwner'), + quiesceUserSchedules: invoke('quiesceUserSchedules'), + restoreUserSchedulesFromDeletion: invoke('restoreUserSchedulesFromDeletion'), + initializeScheduleEngine: invoke('initializeScheduleEngine'), + initializeScheduleErasureSweep: invoke('initializeScheduleErasureSweep'), + recordExpiredScheduleApproval, + isUserDeleting: async (userId) => { + const methods = require('~/models'); + return !(await methods.isAgentTriggerPrincipalActive(userId)); + }, + // Identity-fenced delete for a retained job whose run has since been settled + // inline; reconciliation only reaches jobs whose run is still active. + clearScheduledJob: (...args) => getService().engineDeps.clearReconciledJob(...args), +}; diff --git a/api/server/services/Schedules/index.spec.js b/api/server/services/Schedules/index.spec.js new file mode 100644 index 00000000000..194af1ea684 --- /dev/null +++ b/api/server/services/Schedules/index.spec.js @@ -0,0 +1,60 @@ +const mockRecordScheduleOutcome = jest.fn(); +const mockCreateSchedulesService = jest.fn(() => ({ + recordScheduleOutcome: mockRecordScheduleOutcome, +})); + +jest.mock('@librechat/api', () => ({ + createSchedulesService: (...args) => mockCreateSchedulesService(...args), +})); +jest.mock('mongoose', () => ({ models: {} })); +jest.mock('~/server/services/Config/app', () => ({ getAppConfig: jest.fn() })); +jest.mock('~/server/services/Agents/triggers', () => ({ enqueueAgentTrigger: jest.fn() })); +jest.mock('./access', () => ({ resolveAgentFireAccess: jest.fn() })); +jest.mock('~/models', () => ({ isAgentTriggerPrincipalActive: jest.fn() })); + +const { recordExpiredScheduleApproval } = require('./index'); + +describe('recordExpiredScheduleApproval', () => { + beforeEach(() => { + mockRecordScheduleOutcome.mockReset().mockResolvedValue(true); + }); + + it('ignores ordinary generation approvals', async () => { + await expect( + recordExpiredScheduleApproval('conversation-1', { createdAt: 1000 }), + ).resolves.toBeUndefined(); + + expect(mockRecordScheduleOutcome).not.toHaveBeenCalled(); + }); + + it('settles the exact scheduled generation as interrupted', async () => { + await recordExpiredScheduleApproval('conversation-1', { + createdAt: 1000, + conversationId: 'conversation-1', + scheduleId: 'schedule-1', + scheduledFor: '2026-08-17T12:00:00.000Z', + }); + + expect(mockRecordScheduleOutcome).toHaveBeenCalledWith({ + scheduleId: 'schedule-1', + scheduledFor: '2026-08-17T12:00:00.000Z', + streamId: 'conversation-1', + jobCreatedAt: 1000, + status: 'interrupted', + conversationId: 'conversation-1', + error: 'Approval expired before a decision was made', + }); + }); + + it('surfaces a failed durable settlement instead of reporting false success', async () => { + mockRecordScheduleOutcome.mockResolvedValue(false); + + await expect( + recordExpiredScheduleApproval('conversation-1', { + createdAt: 1000, + scheduleId: 'schedule-1', + scheduledFor: '2026-08-17T12:00:00.000Z', + }), + ).rejects.toThrow('Failed to settle expired scheduled approval schedule-1'); + }); +}); diff --git a/client/src/components/Chat/Footer.tsx b/client/src/components/Chat/Footer.tsx index 90d78737b94..31cc0d48eee 100644 --- a/client/src/components/Chat/Footer.tsx +++ b/client/src/components/Chat/Footer.tsx @@ -63,7 +63,8 @@ function Footer({ className, startupConfig }: FooterProps) { {children} @@ -90,7 +91,6 @@ function Footer({ className, startupConfig }: FooterProps) { className ?? 'absolute bottom-0 left-0 right-0 hidden items-center justify-center gap-2 px-2 py-2 text-center text-xs text-text-primary sm:flex md:px-[60px]' } - role="contentinfo" > {footerElements.map((contentRender, index) => { const isLastElement = index === footerElements.length - 1; diff --git a/client/src/components/Chat/Input/ToolsDropdown.tsx b/client/src/components/Chat/Input/ToolsDropdown.tsx index 2365a4f128a..66d5684ddc4 100644 --- a/client/src/components/Chat/Input/ToolsDropdown.tsx +++ b/client/src/components/Chat/Input/ToolsDropdown.tsx @@ -27,6 +27,10 @@ interface ToolsDropdownProps { disabled?: boolean; } +/** Ariakit portals to document.body by default, which puts the menu outside every landmark. + * Returning null falls back to that default. */ +const getMainLandmark = () => document.querySelector('main'); + const ToolsDropdown = ({ disabled }: ToolsDropdownProps) => { const localize = useLocalize(); const { user } = useAuthContext(); @@ -400,7 +404,10 @@ const ToolsDropdown = ({ disabled }: ToolsDropdownProps) => { menuId="tools-dropdown-menu" isOpen={isPopoverActive} setIsOpen={setIsPopoverActive} - modal={true} + modal={false} + portal={true} + portalElement={getMainLandmark} + preserveTabOrder={false} unmountOnHide={true} trigger={menuTrigger} items={dropdownItems} diff --git a/client/src/components/Chat/__tests__/Footer.spec.tsx b/client/src/components/Chat/__tests__/Footer.spec.tsx new file mode 100644 index 00000000000..c855b746f08 --- /dev/null +++ b/client/src/components/Chat/__tests__/Footer.spec.tsx @@ -0,0 +1,56 @@ +import React from 'react'; +import { render, screen } from '@testing-library/react'; +import '@testing-library/jest-dom/extend-expect'; +import Footer from '../Footer'; + +jest.mock('react-gtm-module', () => ({ + __esModule: true, + default: { initialize: jest.fn() }, +})); + +jest.mock('~/data-provider', () => ({ + useGetStartupConfig: jest.fn(() => ({ data: undefined, isFetching: false, error: null })), +})); + +const mockTranslations: Record = { + com_ui_latest_footer: 'Every AI for Everyone.', + com_ui_privacy_policy: 'Privacy policy', + com_ui_terms_of_service: 'Terms of service', +}; + +jest.mock('~/hooks', () => ({ + useLocalize: () => (key: string) => mockTranslations[key] ?? key, +})); + +describe('Footer', () => { + test('opens the default LibreChat site link in a new tab', () => { + render(