diff --git a/api/server/controllers/agents/__tests__/resume.spec.js b/api/server/controllers/agents/__tests__/resume.spec.js index 866d09024db..33ff36bc5df 100644 --- a/api/server/controllers/agents/__tests__/resume.spec.js +++ b/api/server/controllers/agents/__tests__/resume.spec.js @@ -1275,6 +1275,13 @@ describe('ResumeAgentController (POST /agents/chat/resume)', () => { await settled; expect(capturedInit.isScheduledFire).toBe(true); + expect(mockInitializeClient.mock.calls[0][0].scheduledTokenContext).toEqual({ + scheduleId: 'schedule-1', + ownerId: USER_ID, + tenantId: TENANT_ID, + agentId: AGENT_ID, + invocationMode: 'delegated', + }); expect(mockClaimScheduleResume).toHaveBeenCalledWith('schedule-1', scheduledFor, { expectedConfigRevision: 4, diff --git a/api/server/controllers/agents/client.codeDecision.spec.js b/api/server/controllers/agents/client.codeDecision.spec.js new file mode 100644 index 00000000000..fc38cea553d --- /dev/null +++ b/api/server/controllers/agents/client.codeDecision.spec.js @@ -0,0 +1,50 @@ +const AgentClient = require('./client'); + +describe('AgentClient code environment save options', () => { + const mac = { environmentId: 'code-mac', workspaceId: 'primary' }; + const vm = { environmentId: 'code-vm', workspaceId: 'primary' }; + const conversationId = 'conversation-1'; + + const buildClient = (resolvedConversation) => { + const client = Object.create(AgentClient.prototype); + client.options = { + req: { + body: { conversationId }, + config: {}, + resolvedConversation, + _codeEnvironmentDecision: { mode: 'attached', codeWorkspaces: [mac] }, + }, + endpoint: 'agents', + agent: { id: 'agent_1', provider: 'openai' }, + }; + return client; + }; + + it('recognizes an existing conversation before the client initializes its ID', () => { + const client = buildClient({ + conversationId, + codeEnvironmentMode: 'attached', + codeWorkspaces: [vm], + }); + + expect(client.conversationId).toBeUndefined(); + const initial = client.getSaveOptions(); + expect(initial).not.toHaveProperty('codeEnvironmentMode'); + expect(initial).not.toHaveProperty('codeWorkspaces'); + + client.conversationId = conversationId; + expect(client.getSaveOptions()).toEqual(initial); + }); + + it('seeds a new conversation and omits its decision after the row becomes available', () => { + const client = buildClient(null); + const initial = client.getSaveOptions(); + expect(initial).toMatchObject({ codeEnvironmentMode: 'attached', codeWorkspaces: [mac] }); + + client.conversationId = conversationId; + client.options.req.resolvedConversation = { conversationId, ...initial }; + const paused = client.getSaveOptions(); + expect(paused).not.toHaveProperty('codeEnvironmentMode'); + expect(paused).not.toHaveProperty('codeWorkspaces'); + }); +}); diff --git a/api/server/controllers/agents/client.js b/api/server/controllers/agents/client.js index c98ec0c583a..9398181489c 100644 --- a/api/server/controllers/agents/client.js +++ b/api/server/controllers/agents/client.js @@ -1943,7 +1943,7 @@ class AgentClient extends BaseClient { agentsEConfig?.toolApproval?.enabled !== false, ); const persistedCodeEnvironmentDecision = resolvePersistableCodeEnvironmentDecision({ - conversationId: this.conversationId, + conversationId: this.options.req.body.conversationId, decision: this.options.req._codeEnvironmentDecision, conversation: this.options.req.resolvedConversation, requested: this.options.req.body, diff --git a/api/server/controllers/agents/client.test.js b/api/server/controllers/agents/client.test.js index 852487c2de9..c084201298a 100644 --- a/api/server/controllers/agents/client.test.js +++ b/api/server/controllers/agents/client.test.js @@ -107,7 +107,7 @@ describe('AgentClient code approval persistence', () => { endpoint: EModelEndpoint.agents, agent: { id: 'attached-agent' }, req: { - body: {}, + body: { conversationId: 'convo-1' }, _codeEnvironmentDecision: { mode: 'attached', codeWorkspaces: [{ environmentId: 'mac', workspaceId: 'primary' }], @@ -135,7 +135,7 @@ describe('AgentClient code approval persistence', () => { endpoint: EModelEndpoint.agents, agent: { id: 'attached-agent' }, req: { - body: {}, + body: { conversationId: 'convo-1' }, _codeEnvironmentDecision: { mode: 'attached', codeWorkspaces }, resolvedConversation: { conversationId: 'convo-1', codeWorkspaces }, config: { endpoints: { [EModelEndpoint.agents]: {} } }, diff --git a/api/server/controllers/agents/resume.js b/api/server/controllers/agents/resume.js index 96e481c7f47..2cc9a112d73 100644 --- a/api/server/controllers/agents/resume.js +++ b/api/server/controllers/agents/resume.js @@ -47,6 +47,7 @@ const { findAgentEventAppliedAction, assertCodeExecutionApprovalBinding, collectReachableAgents, + restoreScheduledTokenContext, } = require('@librechat/api'); const { disposeClient } = require('~/server/cleanup'); const { decryptMetadata } = require('~/server/services/ActionService'); @@ -1812,6 +1813,7 @@ const ResumeAgentController = async (req, res, next, initializeClient, addTitle) parentMessageId: job.metadata.userMessage?.messageId ?? Constants.NO_PARENT, }); const result = await initializeClient({ + scheduledTokenContext: restoreScheduledTokenContext(req, job.metadata), req, res, endpointOption: req.body.endpointOption, diff --git a/api/server/services/Endpoints/agents/initialize.js b/api/server/services/Endpoints/agents/initialize.js index 6c826220c50..ee8ebb207d8 100644 --- a/api/server/services/Endpoints/agents/initialize.js +++ b/api/server/services/Endpoints/agents/initialize.js @@ -1873,7 +1873,7 @@ const initializeClientWithProvider = async ({ * token material so refresh remains owned by the host integration. * * @param {object} [dependencies] - * @param {(user: import('@librechat/data-schemas').IUser, options: { signal?: AbortSignal }) => import('@librechat/api').UpstreamTokenProvider | undefined | Promise} [dependencies.resolveUpstreamTokenProvider] + * @param {import('@librechat/api').HostUpstreamTokenProviderResolver} [dependencies.resolveUpstreamTokenProvider] */ function createInitializeClient(dependencies = {}) { return async (params) => { @@ -1881,6 +1881,7 @@ function createInitializeClient(dependencies = {}) { params.req, dependencies.resolveUpstreamTokenProvider, params.signal, + params.scheduledTokenContext, ); return initializeClientWithProvider({ ...params, upstreamTokenProviderResolver }); }; diff --git a/api/server/services/Endpoints/agents/initialize.spec.js b/api/server/services/Endpoints/agents/initialize.spec.js index e76bdbb8c1a..b01af987604 100644 --- a/api/server/services/Endpoints/agents/initialize.spec.js +++ b/api/server/services/Endpoints/agents/initialize.spec.js @@ -182,12 +182,22 @@ describe('initializeClient — processAgent ACL gate', () => { }, }); - it('defers host credential resolution during scheduled agent initialization', async () => { + it.each([false, true])('defers host credential resolution with restored=%s', async (restored) => { const upstreamTokenProvider = jest.fn(); const resolveUpstreamTokenProvider = jest.fn().mockResolvedValue(upstreamTokenProvider); const hostInitializeClient = createInitializeClient({ resolveUpstreamTokenProvider }); const req = makeReq(); req._isScheduledFire = true; + req._isAgentTrigger = !restored; + req.body.agent_id = PRIMARY_ID; + req.body.agentTrigger = { + version: 1, + event: { + type: 'schedule.occurrence', + occurredAt: 0, + source: { type: 'schedule', id: 'sched-1' }, + }, + }; const signal = new AbortController().signal; mockInitializeAgent.mockImplementationOnce(async ({ loadTools, agent }) => { await loadTools({ @@ -203,6 +213,14 @@ describe('initializeClient — processAgent ACL gate', () => { req, res: {}, signal, + scheduledTokenContext: restored + ? { + scheduleId: 'sched-1', + ownerId: req.user.id, + agentId: PRIMARY_ID, + invocationMode: 'delegated', + } + : undefined, endpointOption: makeEndpointOption(), }); @@ -211,7 +229,15 @@ describe('initializeClient — processAgent ACL gate', () => { expect(toolLoadParams.upstreamTokenProvider).toBeUndefined(); const resolver = toolLoadParams.upstreamTokenProviderResolver; await expect(resolver({ signal })).resolves.toBe(upstreamTokenProvider); - expect(resolveUpstreamTokenProvider).toHaveBeenCalledWith(req.user, { signal }); + expect(resolveUpstreamTokenProvider).toHaveBeenCalledWith(req.user, { + signal, + context: { + scheduleId: 'sched-1', + ownerId: req.user.id, + agentId: PRIMARY_ID, + invocationMode: 'delegated', + }, + }); }); it('keeps interactive agent initialization independent of the host resolver', async () => { diff --git a/api/server/services/Schedules/mcp.js b/api/server/services/Schedules/mcp.js index e87a8909f97..ae358c08f39 100644 --- a/api/server/services/Schedules/mcp.js +++ b/api/server/services/Schedules/mcp.js @@ -15,7 +15,7 @@ const methods = require('~/models'); * credential source and therefore remains fail-closed for unattended OBO. * * @param {object} [options] - * @param {(user: import('@librechat/data-schemas').IUser, options: { signal?: AbortSignal }) => import('@librechat/api').UpstreamTokenProvider | undefined | Promise} [options.resolveUpstreamTokenProvider] + * @param {import('@librechat/api').HostUpstreamTokenProviderResolver} [options.resolveUpstreamTokenProvider] */ function createMCPPreflight(options = {}) { return createScheduleMCPPreflight({ diff --git a/client/src/components/Chat/Input/CodeApprovalMenu.tsx b/client/src/components/Chat/Input/CodeApprovalMenu.tsx index 85446bf57ee..13e856a15ff 100644 --- a/client/src/components/Chat/Input/CodeApprovalMenu.tsx +++ b/client/src/components/Chat/Input/CodeApprovalMenu.tsx @@ -7,6 +7,7 @@ import type { CodeApprovalMode, TConversation } from 'librechat-data-provider'; import type { LucideIcon } from 'lucide-react'; import type { SetterOrUpdater } from 'recoil'; import type { TranslationKeys } from '~/hooks'; +import { useCodeApprovalModePreference } from '~/hooks/Agents/codeApprovalPreference'; import { useCodeApprovalMode, useLocalize } from '~/hooks'; import { cn } from '~/utils'; @@ -51,6 +52,7 @@ export default function CodeApprovalMenu({ const localize = useLocalize(); const queryClient = useQueryClient(); const { available, modes, selected } = useCodeApprovalMode(conversation, addedConversation); + const preference = useCodeApprovalModePreference(); const menuStore = Ariakit.useMenuStore({ focusLoop: true, placement: 'top-start' }); const isOpen = menuStore.useState('open'); @@ -62,11 +64,14 @@ export default function CodeApprovalMenu({ * cache, so the pick lands there too, seeding the record from the live * conversation when none exists yet (the resumable transport seeds the same * key optimistically). A chat that has no id yet keeps the pick in - * conversation state alone until the run assigns one. */ + * conversation state alone until the run assigns one. The pick is also this + * browser's remembered default, so the next chat opens on it rather than back + * at `ask`; policy is re-checked before it is ever shown or submitted. */ const selectMode = (mode: CodeApprovalMode) => { if (!modes.includes(mode)) { return; } + preference.remember(mode); setConversation((current) => current == null ? current : { ...current, codeApprovalMode: mode }, ); diff --git a/client/src/components/Chat/Input/__tests__/CodeApprovalMenu.spec.tsx b/client/src/components/Chat/Input/__tests__/CodeApprovalMenu.spec.tsx index 4d48c8981e1..8a584b561d4 100644 --- a/client/src/components/Chat/Input/__tests__/CodeApprovalMenu.spec.tsx +++ b/client/src/components/Chat/Input/__tests__/CodeApprovalMenu.spec.tsx @@ -1,7 +1,7 @@ import userEvent from '@testing-library/user-event'; -import { QueryKeys, Constants } from 'librechat-data-provider'; import { render, screen, within } from '@testing-library/react'; import { QueryClient, QueryClientProvider } from '@tanstack/react-query'; +import { QueryKeys, Constants, LocalStorageKeys } from 'librechat-data-provider'; import type { TConversation } from 'librechat-data-provider'; import CodeApprovalMenu from '../CodeApprovalMenu'; @@ -163,3 +163,28 @@ describe('CodeApprovalMenu', () => { }); }); }); + +test("remembers the pick as this browser's default for the next chat", async () => { + localStorage.clear(); + mockUseCodeApprovalMode.mockReturnValue({ + available: true, + modes: ['ask', 'acceptEdits'], + selected: 'ask', + }); + render( + + + , + ); + + await userEvent.click(screen.getByTestId('code-approval-mode')); + await userEvent.click(await screen.findByText('com_ui_code_approval_accept_edits')); + + expect(localStorage.getItem(LocalStorageKeys.LAST_CODE_APPROVAL_MODE)).toBe( + JSON.stringify('acceptEdits'), + ); +}); diff --git a/client/src/hooks/Agents/__tests__/useCodeApprovalMode.test.ts b/client/src/hooks/Agents/__tests__/useCodeApprovalMode.test.ts index 35d58e3f7fd..688bbf037ab 100644 --- a/client/src/hooks/Agents/__tests__/useCodeApprovalMode.test.ts +++ b/client/src/hooks/Agents/__tests__/useCodeApprovalMode.test.ts @@ -1,4 +1,6 @@ +import { Provider } from 'jotai'; import { renderHook } from '@testing-library/react'; +import { LocalStorageKeys } from 'librechat-data-provider'; import type { TConversation } from 'librechat-data-provider'; import useCodeApprovalMode from '../useCodeApprovalMode'; @@ -371,3 +373,70 @@ describe('useCodeApprovalMode', () => { expect(result.current.selected).toBe('ask'); }); }); + +describe('useCodeApprovalMode remembered pick', () => { + const withoutMode = { ...conversation, codeApprovalMode: undefined } as TConversation; + + beforeEach(() => { + jest.clearAllMocks(); + localStorage.clear(); + mockUseAgentsMapContext.mockReturnValue({}); + mockUseAgentToolPermissions.mockReturnValue({ + agent: { + id: 'agent_1', + tools: ['execute_code'], + stateful_code_sessions: true, + code_environment_id: 'mac', + }, + }); + mockUseGetAgentsConfig.mockReturnValue({ + agentsConfig: { + statefulCodeSessions: { + approvalsEnabled: true, + approvalModes: ['ask', 'acceptEdits'], + environments: [ + { + id: 'mac', + name: 'Mac', + type: 'attached', + configSchema: { + permissions: { + fileWrite: { allowed: ['ask', 'allow'], default: 'ask' }, + commandExecution: { allowed: ['ask'], default: 'ask' }, + }, + }, + }, + ], + }, + }, + }); + }); + + test('opens a conversation with no stored mode on the last pick', () => { + localStorage.setItem(LocalStorageKeys.LAST_CODE_APPROVAL_MODE, JSON.stringify('acceptEdits')); + const { result } = renderHook(() => useCodeApprovalMode(withoutMode), { wrapper: Provider }); + expect(result.current.selected).toBe('acceptEdits'); + }); + + test('a stored conversation mode still wins over the remembered pick', () => { + localStorage.setItem(LocalStorageKeys.LAST_CODE_APPROVAL_MODE, JSON.stringify('acceptEdits')); + const { result } = renderHook( + () => useCodeApprovalMode({ ...conversation, codeApprovalMode: 'ask' } as TConversation), + { wrapper: Provider }, + ); + expect(result.current.selected).toBe('ask'); + }); + + test('falls back to ask when policy no longer allows the remembered pick', () => { + localStorage.setItem(LocalStorageKeys.LAST_CODE_APPROVAL_MODE, JSON.stringify('fullAccess')); + const { result } = renderHook(() => useCodeApprovalMode(withoutMode), { wrapper: Provider }); + expect(result.current.modes).not.toContain('fullAccess'); + expect(result.current.selected).toBe('ask'); + }); + + test('ignores a corrupted remembered value', () => { + localStorage.setItem(LocalStorageKeys.LAST_CODE_APPROVAL_MODE, JSON.stringify('nonsense')); + const { result } = renderHook(() => useCodeApprovalMode(withoutMode), { wrapper: Provider }); + expect(result.current.selected).toBe('ask'); + }); +}); diff --git a/client/src/hooks/Agents/codeApprovalPreference.ts b/client/src/hooks/Agents/codeApprovalPreference.ts new file mode 100644 index 00000000000..415d9d79367 --- /dev/null +++ b/client/src/hooks/Agents/codeApprovalPreference.ts @@ -0,0 +1,33 @@ +import { useAtom } from 'jotai'; +import { CODE_APPROVAL_MODES, LocalStorageKeys } from 'librechat-data-provider'; +import type { CodeApprovalMode } from 'librechat-data-provider'; +import { createStorageAtom } from '~/store/jotai-utils'; + +const preference = createStorageAtom(LocalStorageKeys.LAST_CODE_APPROVAL_MODE, null); + +function validMode(saved: string | null): CodeApprovalMode | undefined { + return (CODE_APPROVAL_MODES as readonly string[]).includes(saved ?? '') + ? (saved as CodeApprovalMode) + : undefined; +} + +/** + * The mode the reader last picked in this browser, so a new chat opens the way + * they left the last one instead of falling back to `ask` every time. A + * disposable hint and never authorization: the caller re-checks it against the + * modes current policy allows before showing or submitting it, and logout + * clears it with the rest of this browser's conversation state. + */ +export function useCodeApprovalModePreference() { + const [saved, setSaved] = useAtom(preference); + return { + get: () => validMode(saved), + remember: (mode: CodeApprovalMode) => { + try { + setSaved(mode); + } catch { + // Disabled or full browser storage must not prevent an explicit pick. + } + }, + }; +} diff --git a/client/src/hooks/Agents/useCodeApprovalMode.ts b/client/src/hooks/Agents/useCodeApprovalMode.ts index 3ee381cd145..1836784291f 100644 --- a/client/src/hooks/Agents/useCodeApprovalMode.ts +++ b/client/src/hooks/Agents/useCodeApprovalMode.ts @@ -8,6 +8,7 @@ import { } from 'librechat-data-provider'; import type { Agent, TAgentsMap, TConfig, TPublicCodeEnvironment } from 'librechat-data-provider'; import type { CodeApprovalMode, TConversation } from 'librechat-data-provider'; +import { useCodeApprovalModePreference } from './codeApprovalPreference'; import useAgentToolPermissions from './useAgentToolPermissions'; import useGetAgentsConfig from './useGetAgentsConfig'; import { useAgentsMapContext } from '~/Providers'; @@ -22,6 +23,7 @@ export default function useCodeApprovalMode( } { const { agentsConfig } = useGetAgentsConfig(); const agentsMap = useAgentsMapContext(); + const preference = useCodeApprovalModePreference(); const { agent: primaryAgent } = useAgentToolPermissions(conversation?.agent_id); const { agent: addedAgent } = useAgentToolPermissions(addedConversation?.agent_id); const statefulCodeSessions = agentsConfig?.statefulCodeSessions as @@ -81,7 +83,12 @@ export default function useCodeApprovalMode( if (fullAccessAllowed) allowed.add('fullAccess'); return CODE_APPROVAL_MODES.filter((mode) => allowed.has(mode)); }, [attachedEnvironments, available, codeEnvironments, endpointModes, reachable.complete]); - const requested = conversation?.codeApprovalMode ?? 'ask'; + /** A conversation that carries no mode of its own opens on the reader's last + * pick in this browser, so choosing `acceptEdits` or `fullAccess` survives a + * new chat and a reload instead of being re-picked every time. The remembered + * value is a preference, not a grant: it passes the same policy gate below as + * a stored one, so a mode current policy no longer allows falls back to `ask`. */ + const requested = conversation?.codeApprovalMode ?? preference.get() ?? 'ask'; /** * Fail closed while agent/environment metadata is incomplete. An affirmative * server capability means `ask` is safe to submit even before an attached diff --git a/client/src/hooks/SSE/__tests__/useEventHandlers.spec.ts b/client/src/hooks/SSE/__tests__/useEventHandlers.spec.ts index 524b6966386..d38e5b0dbb1 100644 --- a/client/src/hooks/SSE/__tests__/useEventHandlers.spec.ts +++ b/client/src/hooks/SSE/__tests__/useEventHandlers.spec.ts @@ -1,5 +1,11 @@ -import { Constants } from 'librechat-data-provider'; -import type { EventSubmission, TMessage, TConversation } from 'librechat-data-provider'; +import { Constants, ContentTypes } from 'librechat-data-provider'; +import type { + TMessageContentParts, + EventSubmission, + TConversation, + TMessage, +} from 'librechat-data-provider'; +import type { TResData } from '~/common'; import { buildCreatedInitialResponse, getExistingConversationAbortMessages, @@ -8,8 +14,10 @@ import { buildRecoveryPreset, mergeErrorMessages, mergeRegenerateFinalMessages, + resolveErrorTurn, startedAsNewConversation, } from '~/hooks/SSE/useEventHandlers'; +import { stripStreamedIndexStamps, getPartKeyIndex } from '~/utils'; describe('buildCreatedInitialResponse', () => { const userMessage = { @@ -352,3 +360,258 @@ describe('buildRecoveryPreset', () => { expect(buildRecoveryPreset(sent, undefined, '_fresh').codeApprovalMode).toBe('ask'); }); }); + +describe('resolveErrorTurn', () => { + const userMessage = { + messageId: 'user-1', + conversationId: 'conversation-1', + parentMessageId: Constants.NO_PARENT, + isCreatedByUser: true, + text: 'Look up the issue', + sender: 'User', + } as TMessage; + const initialResponse = { + messageId: 'user-1_', + parentMessageId: 'user-1', + conversationId: 'conversation-1', + isCreatedByUser: false, + text: '', + sender: 'Lia', + endpoint: 'agents', + model: 'agent_1', + } as TMessage; + const submission = { + messages: [], + userMessage, + initialResponse, + conversation: { conversationId: 'conversation-1' }, + } as unknown as EventSubmission; + const streamedParts = [ + { type: ContentTypes.THINK, think: 'Checking the issue' }, + { type: ContentTypes.TEXT, text: 'Let me try the GitHub CLI from the workspace.' }, + { + type: ContentTypes.TOOL_CALL, + tool_call: { id: 'call-1', name: 'execute_code', args: '{}', progress: 1 }, + }, + ] as TMessageContentParts[]; + const streamedResponse = { ...initialResponse, content: streamedParts }; + const startFailureText = JSON.stringify({ code: 'code_workspace_unavailable', reason: 'locked' }); + const startFailure = { + text: startFailureText, + metadata: { streamStartFailed: true }, + } as unknown as TResData; + + it('keeps what the run streamed and takes the failure as one more part', () => { + const { conversationId, errorResponse, recover } = resolveErrorTurn({ + data: startFailure, + submission, + getMessages: () => [userMessage, streamedResponse], + isNewConversationRoute: false, + }); + + expect(conversationId).toBe('conversation-1'); + expect(recover).toBe(false); + expect(errorResponse.content).toEqual([ + ...streamedParts, + { type: ContentTypes.ERROR, error: startFailureText }, + ]); + expect(errorResponse.error).toBeUndefined(); + expect(errorResponse.text).toBe(''); + expect(errorResponse.messageId).toBe('user-1_'); + expect(errorResponse.parentMessageId).toBe('user-1'); + expect(errorResponse.metadata).toEqual({ streamStartFailed: true }); + }); + + it('drops holes, keeps empty slots, and keeps the identity every part streamed under', () => { + const openedThink = { type: ContentTypes.THINK, think: '' }; + const openedText = { type: ContentTypes.TEXT, text: '' }; + const { errorResponse } = resolveErrorTurn({ + data: startFailure, + submission, + getMessages: () => [ + userMessage, + { + ...initialResponse, + content: [ + openedThink, + streamedParts[0], + undefined, + streamedParts[1], + streamedParts[2], + openedText, + ], + } as TMessage, + ], + isNewConversationRoute: false, + }); + + const content = errorResponse.content ?? []; + expect(stripStreamedIndexStamps(content)).toEqual([ + openedThink, + ...streamedParts, + openedText, + { type: ContentTypes.ERROR, error: startFailureText }, + ]); + expect(content.map((part, idx) => getPartKeyIndex(part, idx))).toEqual([0, 1, 3, 4, 5, 6]); + }); + + it('keeps the comparison lanes when one side streamed before the failure', () => { + const lanes = [ + { type: '', agentId: 'agent_a', groupId: 1 }, + { type: ContentTypes.TEXT, text: 'From the added agent', agentId: 'agent_b', groupId: 1 }, + ]; + const { errorResponse } = resolveErrorTurn({ + data: startFailure, + submission, + getMessages: () => [ + userMessage, + { ...initialResponse, content: lanes } as unknown as TMessage, + ], + isNewConversationRoute: false, + }); + + expect(errorResponse.content).toEqual([ + ...lanes, + { type: ContentTypes.ERROR, error: startFailureText }, + ]); + }); + + it('is the whole row when a comparison run failed before either lane streamed', () => { + const { errorResponse } = resolveErrorTurn({ + data: startFailure, + submission, + getMessages: () => [ + userMessage, + { + ...initialResponse, + content: [ + { type: '', agentId: 'agent_a', groupId: 1 }, + { type: '', agentId: 'agent_b', groupId: 1 }, + ], + } as unknown as TMessage, + ], + isNewConversationRoute: false, + }); + + expect(errorResponse.content).toBeUndefined(); + expect(errorResponse.error).toBe(true); + expect(errorResponse.text).toBe(startFailureText); + }); + + it('is the whole row when nothing streamed', () => { + const { errorResponse } = resolveErrorTurn({ + data: startFailure, + submission, + getMessages: () => [userMessage, initialResponse], + isNewConversationRoute: false, + }); + + expect(errorResponse.content).toBeUndefined(); + expect(errorResponse.error).toBe(true); + expect(errorResponse.text).toBe(startFailureText); + expect(errorResponse.messageId).toBe('user-1_'); + expect(errorResponse.parentMessageId).toBe('user-1'); + }); + + it('never takes a user row at the tail for the failed response', () => { + const { errorResponse } = resolveErrorTurn({ + data: startFailure, + submission, + getMessages: () => [ + { ...userMessage, content: [{ type: ContentTypes.TEXT, text: 'Look up the issue' }] }, + ], + isNewConversationRoute: false, + }); + + expect(errorResponse.content).toBeUndefined(); + expect(errorResponse.error).toBe(true); + }); + + it('records a lost connection as a part of the streamed response and rebuilds the chat', () => { + const { conversationId, errorResponse, recover } = resolveErrorTurn({ + data: undefined, + submission, + getMessages: () => [userMessage, streamedResponse], + isNewConversationRoute: false, + }); + + expect(conversationId).toBe('conversation-1'); + expect(recover).toBe(true); + expect(errorResponse.content).toEqual([ + ...streamedParts, + { + type: ContentTypes.ERROR, + error: 'Error connecting to server, try refreshing the page.', + }, + ]); + }); + + it("keeps the streamed row's envelope when a server failure carries its own", () => { + const streamedAt = '2026-09-15T13:00:00.000Z'; + const { errorResponse } = resolveErrorTurn({ + data: { + conversationId: 'conversation-1', + messageId: 'user-1_', + isCreatedByUser: false, + sender: 'System', + model: null, + iconURL: null, + createdAt: '2026-09-15T13:05:00.000Z', + text: startFailureText, + metadata: { streamStartFailed: true }, + } as unknown as TResData, + submission, + getMessages: () => [ + userMessage, + { ...streamedResponse, iconURL: 'lia.png', createdAt: streamedAt, metadata: { seed: 1 } }, + ], + isNewConversationRoute: false, + }); + + expect(errorResponse).toEqual( + expect.objectContaining({ + sender: 'Lia', + model: 'agent_1', + iconURL: 'lia.png', + createdAt: streamedAt, + metadata: { seed: 1, streamStartFailed: true }, + }), + ); + expect(errorResponse.content?.at(-1)).toEqual({ + type: ContentTypes.ERROR, + error: startFailureText, + }); + }); + + it('keeps the streamed parts under a failure the server addressed to the conversation', () => { + const serverText = JSON.stringify({ type: 'invalid_request' }); + const data = { + conversationId: 'conversation-1', + messageId: 'user-1_', + isCreatedByUser: false, + sender: 'Lia', + text: serverText, + } as unknown as TResData; + + const fromChat = resolveErrorTurn({ + data, + submission, + getMessages: () => [userMessage, streamedResponse], + isNewConversationRoute: false, + }); + const fromNewChat = resolveErrorTurn({ + data, + submission, + getMessages: () => [userMessage, streamedResponse], + isNewConversationRoute: true, + }); + + expect(fromChat.recover).toBe(false); + expect(fromNewChat.recover).toBe(true); + expect(fromChat.errorResponse.content).toEqual([ + ...streamedParts, + { type: ContentTypes.ERROR, error: serverText }, + ]); + expect(fromChat.errorResponse.parentMessageId).toBe('user-1'); + }); +}); diff --git a/client/src/hooks/SSE/useEventHandlers.ts b/client/src/hooks/SSE/useEventHandlers.ts index dbc285664c1..91aa6103667 100644 --- a/client/src/hooks/SSE/useEventHandlers.ts +++ b/client/src/hooks/SSE/useEventHandlers.ts @@ -17,8 +17,9 @@ import type { TPreset, TMessage, TConversation, - EventSubmission, TStartupConfig, + EventSubmission, + TMessageContentParts, } from 'librechat-data-provider'; import type { InfiniteData } from '@tanstack/react-query'; import type { SetterOrUpdater } from 'recoil'; @@ -39,6 +40,8 @@ import { removeConvoFromAllQueries, findConversationInInfinite, preserveStreamedContentIdentity, + isEmptyContentPart, + getPartKeyIndex, } from '~/utils'; import { startupConfigKey, @@ -229,6 +232,50 @@ export type EventHandlerParams = { setShowStopButton: SetterOrUpdater; }; +const CONNECTION_ERROR_TEXT = 'Error connecting to server, try refreshing the page.'; + +/** + * The parts the in-flight response has streamed so far: the transcript's tail (a user row there + * means no response was placed), without the holes an interrupted stream leaves. Whether anything + * streamed is judged without the slots that never received content — a comparison run's + * `type: ''` placeholders, a text or think part opened before its first delta — but once something + * did, those slots stay: a placeholder is what keeps a comparison lane's layout and attribution + * when only the other lane produced output. The parts carry the render identity they streamed + * under, as the final path stamps it, so the settled row does not remount. + */ +const getStreamedContent = (message?: TMessage): TMessageContentParts[] => { + if (message == null || message.isCreatedByUser === true) { + return []; + } + const streamed = message.content ?? []; + const parts = streamed.filter((part): part is TMessageContentParts => part != null); + if (!parts.some((part) => !isEmptyContentPart(part))) { + return []; + } + return preserveStreamedContentIdentity(streamed, parts) ?? parts; +}; + +/** Keys the appended failure past every key the kept parts render under, physical or stamped. */ +const appendErrorPart = ( + content: TMessageContentParts[], + errorText: string, +): TMessageContentParts[] => { + const keyIndex = content.reduce( + (next, part, idx) => Math.max(next, getPartKeyIndex(part, idx) + 1), + 0, + ); + const errorPart: TMessageContentParts = { type: ContentTypes.ERROR, error: errorText }; + return [ + ...content, + keyIndex === content.length ? errorPart : { ...errorPart, streamedIndex: keyIndex }, + ]; +}; + +/** + * A failure that lands after the run has streamed keeps what streamed and takes the failure as + * one more part, the way the server records a failure inside a run; one that lands before + * anything streamed is the whole row. + */ const createErrorMessage = ({ errorMetadata, getMessages, @@ -242,41 +289,24 @@ const createErrorMessage = ({ }): TMessage => { const currentMessages = getMessages(); const latestMessage = currentMessages?.[currentMessages.length - 1]; - let errorMessage: TMessage; const text = submission.initialResponse.text.length > 45 ? submission.initialResponse.text : ''; const errorText = (errorMetadata?.text || text || (error as Error | undefined)?.message) ?? 'Error cancelling request'; - const latestContent = latestMessage?.content ?? []; - let isValidContentPart = false; - if (latestContent.length > 0) { - const latestContentPart = latestContent[latestContent.length - 1]; - if (latestContentPart != null) { - const latestPartValue = latestContentPart[latestContentPart.type ?? '']; - isValidContentPart = - latestContentPart.type !== ContentTypes.TEXT || - (latestContentPart.type === ContentTypes.TEXT && typeof latestPartValue === 'string') - ? true - : latestPartValue?.value !== ''; - } - } - if ( - latestMessage?.conversationId && - latestMessage?.messageId && - latestContent && - isValidContentPart - ) { - const content = [...latestContent]; - content.push({ - type: ContentTypes.ERROR, - error: errorText, - }); - errorMessage = { + const streamedContent = getStreamedContent(latestMessage); + if (latestMessage?.conversationId && latestMessage.messageId && streamedContent.length > 0) { + /** The row keeps the envelope it streamed under — author, model, icon, creation time — and + * takes from the failure only what describes the failure: a server payload names `System` + * as its sender and the schema dates a fresh envelope now, neither of which applies to a + * response that is merely gaining a part. */ + const errorMessage: TMessage = { ...latestMessage, - ...errorMetadata, + ...(errorMetadata?.metadata != null + ? { metadata: { ...latestMessage.metadata, ...errorMetadata.metadata } } + : {}), error: undefined, text: '', - content, + content: appendErrorPart(streamedContent, errorText), }; if ( submission.userMessage.messageId && @@ -285,18 +315,94 @@ const createErrorMessage = ({ errorMessage.parentMessageId = submission.userMessage.messageId; } return errorMessage; - } else if (errorMetadata) { + } + if (errorMetadata) { return errorMetadata as TMessage; - } else { - errorMessage = { - ...submission, - ...submission.initialResponse, - text: errorText, - unfinished: !!text.length, + } + return tMessageSchema.parse({ + ...submission, + ...submission.initialResponse, + text: errorText, + unfinished: !!text.length, + error: true, + }) as TMessage; +}; + +export interface ErrorTurn { + conversationId: string; + errorResponse: TMessage; + /** The conversation must be rebuilt: it never learned its id, the transport failed, or the + * server named a conversation the new-chat route has not navigated to. */ + recover: boolean; +} + +/** Builds the row a failed turn leaves in the transcript and names the conversation it belongs to. */ +export const resolveErrorTurn = ({ + data, + submission, + getMessages, + isNewConversationRoute, +}: { + data?: TResData; + submission: EventSubmission; + getMessages: () => TMessage[] | undefined; + isNewConversationRoute: boolean; +}): ErrorTurn => { + const { userMessage, initialResponse } = submission; + const conversationId = + userMessage.conversationId ?? submission.conversation?.conversationId ?? ''; + + const parseErrorResponse = (payload: TResData | Partial): TMessage => { + const metadata = payload['responseMessage'] ?? payload; + const errorMessage: Partial = { + ...initialResponse, + ...metadata, error: true, + parentMessageId: userMessage.messageId, }; + + if (errorMessage.messageId === undefined || errorMessage.messageId === '') { + errorMessage.messageId = v4(); + } + + return tMessageSchema.parse(errorMessage) as TMessage; + }; + const build = (errorMetadata: TMessage): TMessage => + createErrorMessage({ errorMetadata, getMessages, submission }); + + if (!data) { + const convoId = conversationId || `_${v4()}`; + return { + conversationId: convoId, + recover: true, + errorResponse: build( + parseErrorResponse({ text: CONNECTION_ERROR_TEXT, ...submission, conversationId: convoId }), + ), + }; + } + + const receivedConvoId = data.conversationId ?? ''; + if (!conversationId && !receivedConvoId) { + return { + conversationId: `_${v4()}`, + recover: true, + errorResponse: build(parseErrorResponse(data)), + }; + } + if (!receivedConvoId) { + return { conversationId, recover: false, errorResponse: build(parseErrorResponse(data)) }; } - return tMessageSchema.parse(errorMessage) as TMessage; + return { + conversationId: receivedConvoId, + recover: isNewConversationRoute, + errorResponse: build( + tMessageSchema.parse({ + ...data, + error: true, + parentMessageId: userMessage.messageId, + }) as TMessage, + ), + }; }; /** @@ -1056,81 +1162,22 @@ export default function useEventHandlers({ const errorHandler = useCallback( ({ data, submission }: { data?: TResData; submission: EventSubmission }) => { - const { userMessage, initialResponse } = submission; - setCompleted((prev) => new Set(prev.add(initialResponse.messageId))); + setCompleted((prev) => new Set(prev.add(submission.initialResponse.messageId))); setSubmissionStart(null); - const conversationId = - userMessage.conversationId ?? submission.conversation?.conversationId ?? ''; - - const setErrorMessages = (convoId: string, errorMessage: TMessage) => { - const finalMessages = mergeErrorMessages({ ...submission, errorMessage }); - setMessages(finalMessages); - queryClient.setQueryData([QueryKeys.messages, convoId], finalMessages); - }; - - const parseErrorResponse = (data: TResData | Partial): TMessage => { - const metadata = data['responseMessage'] ?? data; - const errorMessage: Partial = { - ...initialResponse, - ...metadata, - error: true, - parentMessageId: userMessage.messageId, - }; - - if (errorMessage.messageId === undefined || errorMessage.messageId === '') { - errorMessage.messageId = v4(); - } - - return tMessageSchema.parse(errorMessage) as TMessage; - }; - - if (!data) { - const convoId = conversationId || `_${v4()}`; - const errorMetadata = parseErrorResponse({ - text: 'Error connecting to server, try refreshing the page.', - ...submission, - conversationId: convoId, - }); - const errorResponse = createErrorMessage({ - errorMetadata, - getMessages, - submission, - }); - setErrorMessages(convoId, errorResponse); - recoverConversation(convoId, submission); - setIsSubmitting(false); - return; - } - - const receivedConvoId = data.conversationId ?? ''; - if (!conversationId && !receivedConvoId) { - const convoId = `_${v4()}`; - const errorResponse = parseErrorResponse(data); - setErrorMessages(convoId, errorResponse); - recoverConversation(convoId, submission); - setIsSubmitting(false); - return; - } else if (!receivedConvoId) { - const errorResponse = parseErrorResponse(data); - setErrorMessages(conversationId, errorResponse); - setIsSubmitting(false); - return; - } - - const errorResponse = tMessageSchema.parse({ - ...data, - error: true, - parentMessageId: userMessage.messageId, - }) as TMessage; - - setErrorMessages(receivedConvoId, errorResponse); - if (receivedConvoId && paramId === Constants.NEW_CONVO) { - recoverConversation(receivedConvoId, submission); + const { conversationId, errorResponse, recover } = resolveErrorTurn({ + data, + submission, + getMessages, + isNewConversationRoute: paramId === Constants.NEW_CONVO, + }); + const finalMessages = mergeErrorMessages({ ...submission, errorMessage: errorResponse }); + setMessages(finalMessages); + queryClient.setQueryData([QueryKeys.messages, conversationId], finalMessages); + if (recover) { + recoverConversation(conversationId, submission); } - setIsSubmitting(false); - return; }, [ setCompleted, diff --git a/client/src/utils/localStorage.ts b/client/src/utils/localStorage.ts index 17900689219..1cf590b8ee3 100644 --- a/client/src/utils/localStorage.ts +++ b/client/src/utils/localStorage.ts @@ -60,7 +60,10 @@ export function clearLocalStorage(skipFirst?: boolean) { key === LocalStorageKeys.LAST_SPEC || key === LocalStorageKeys.LAST_TOOLS || key === LocalStorageKeys.LAST_MODEL || - key === LocalStorageKeys.FILES_TO_DELETE + key === LocalStorageKeys.FILES_TO_DELETE || + /** A permissive code approval default belongs to the account that chose it, not to + * whoever signs in next on a shared browser. */ + key === LocalStorageKeys.LAST_CODE_APPROVAL_MODE ) { localStorage.removeItem(key); } diff --git a/client/src/utils/messages.ts b/client/src/utils/messages.ts index 532add8de91..f279360eb79 100644 --- a/client/src/utils/messages.ts +++ b/client/src/utils/messages.ts @@ -204,7 +204,7 @@ const getPartToolCall = (part: TMessageContentParts): Agents.ToolCall | undefine /** Slots the persistence compaction leaves nothing behind for: the * dual-message `type: ''` placeholders, text/think parts that never received a * delta, and tool calls missing their `tool_call` payload. */ -const isEmptyContentPart = (part: TMessageContentParts): boolean => { +export const isEmptyContentPart = (part: TMessageContentParts): boolean => { if (!part.type) { return true; } diff --git a/packages/api/src/agents/subagentThreads.spec.ts b/packages/api/src/agents/subagentThreads.spec.ts index 1d44c2f3cf8..32daa78c447 100644 --- a/packages/api/src/agents/subagentThreads.spec.ts +++ b/packages/api/src/agents/subagentThreads.spec.ts @@ -1413,9 +1413,12 @@ describe('SubagentThreadTaskStore', () => { let ownerActive = true; const options = { isOwnerActive: async () => ownerActive, - leaseTtlMs: 60, + // Exercise owner cancellation, not lease expiry. Keep the lease beyond the + // drain deadline so slow CI database operations cannot bypass child startup + // or make the drain succeed without the worker releasing its lease. + leaseTtlMs: 30_000, leaseHeartbeatMs: 10, - ownerDrainTimeoutMs: 1_000, + ownerDrainTimeoutMs: 5_000, ownerDrainPollMs: 5, }; const workerStore = new SubagentThreadTaskStore(methods, options); diff --git a/packages/api/src/mcp/MCPConnectionFactory.ts b/packages/api/src/mcp/MCPConnectionFactory.ts index 66b3368ea18..6c416ae65ab 100644 --- a/packages/api/src/mcp/MCPConnectionFactory.ts +++ b/packages/api/src/mcp/MCPConnectionFactory.ts @@ -654,6 +654,7 @@ export class MCPConnectionFactory { this.upstreamTokenProvider = createLazyOboUpstreamTokenProvider( this.upstreamTokenProviderResolver, this.signal, + { mcpServer: this.serverName, scopes: oboConfig.scopes }, ); } if (!this.upstreamTokenProvider) { diff --git a/packages/api/src/mcp/MCPManager.ts b/packages/api/src/mcp/MCPManager.ts index 5e9a7825cf4..bf5196c3828 100644 --- a/packages/api/src/mcp/MCPManager.ts +++ b/packages/api/src/mcp/MCPManager.ts @@ -1302,6 +1302,7 @@ Please follow these instructions when using tools from the respective MCP server oboUpstreamTokenProvider = createLazyOboUpstreamTokenProvider( upstreamTokenProviderResolver, options?.signal, + { mcpServer: serverName, scopes: oboConfig.scopes }, ); } if (!oboUpstreamTokenProvider) { diff --git a/packages/api/src/mcp/__tests__/MCPConnectionFactory.test.ts b/packages/api/src/mcp/__tests__/MCPConnectionFactory.test.ts index f370a8d12af..ffe06f27734 100644 --- a/packages/api/src/mcp/__tests__/MCPConnectionFactory.test.ts +++ b/packages/api/src/mcp/__tests__/MCPConnectionFactory.test.ts @@ -5586,6 +5586,11 @@ describe('MCPConnectionFactory', () => { ); expect(upstreamTokenProviderResolver).toHaveBeenCalledTimes(1); + expect(upstreamTokenProviderResolver).toHaveBeenCalledWith( + expect.objectContaining({ + target: { mcpServer: 'obo-srv', scopes: oboServerConfig.obo?.scopes }, + }), + ); expect(resolveOboToken).toHaveBeenCalledWith( mockUser, oboServerConfig.obo, diff --git a/packages/api/src/mcp/__tests__/MCPManager.test.ts b/packages/api/src/mcp/__tests__/MCPManager.test.ts index ac8c638cd18..32b75090e3c 100644 --- a/packages/api/src/mcp/__tests__/MCPManager.test.ts +++ b/packages/api/src/mcp/__tests__/MCPManager.test.ts @@ -2914,6 +2914,11 @@ describe('MCPManager', () => { }); expect(upstreamTokenProviderResolver).toHaveBeenCalledTimes(1); + expect(upstreamTokenProviderResolver).toHaveBeenCalledWith( + expect.objectContaining({ + target: { mcpServer: serverName, scopes: serverConfig.obo?.scopes }, + }), + ); expect(mockResolveOboToken).toHaveBeenCalledWith( mockUser, serverConfig.obo, diff --git a/packages/api/src/mcp/oauth/obo.ts b/packages/api/src/mcp/oauth/obo.ts index 878dfee4acc..32fca35f2e9 100644 --- a/packages/api/src/mcp/oauth/obo.ts +++ b/packages/api/src/mcp/oauth/obo.ts @@ -45,9 +45,16 @@ export type UpstreamTokenProvider = (options?: { signal?: AbortSignal; }) => Promise; +/** Target resolved from server configuration after the OBO trust check. Scopes are not an audience. */ +export interface UpstreamTokenTarget { + readonly mcpServer: string; + readonly scopes: string; +} + /** Lazily supplies a renewable upstream-token provider when an OBO server actually needs one. */ export type UpstreamTokenProviderResolver = (options?: { signal?: AbortSignal; + target?: UpstreamTokenTarget; }) => UpstreamTokenProvider | undefined | Promise; /** Scheduled OBO credentials must not replace the browser's direct-bearer source. */ @@ -90,6 +97,7 @@ export async function awaitOboOperation( export function createLazyOboUpstreamTokenProvider( resolver: UpstreamTokenProviderResolver, signal?: AbortSignal, + target?: UpstreamTokenTarget, ): UpstreamTokenProvider { let pending: Promise | undefined; return async (options) => { @@ -99,7 +107,7 @@ export function createLazyOboUpstreamTokenProvider( pending ??= Promise.resolve() .then(() => { effectiveSignal?.throwIfAborted(); - return resolver({ signal: effectiveSignal }); + return resolver({ signal: effectiveSignal, ...(target ? { target } : {}) }); }) .catch((error) => { pending = undefined; diff --git a/packages/api/src/schedules/README.md b/packages/api/src/schedules/README.md index 408641fa6d7..e155a4b4024 100644 --- a/packages/api/src/schedules/README.md +++ b/packages/api/src/schedules/README.md @@ -38,3 +38,45 @@ messages or OAuth URLs. Successful dispatch records also retain the server outco The schedule card shows failed servers and links to the selected agent for recovery. A pure pause remains available even when MCP validation fails. This change does not re-enable existing schedules automatically or repair credentials on the user's behalf. + +## Host token-provider context + +`createMCPPreflight` and `createInitializeClient` accept a +`HostUpstreamTokenProviderResolver`. The default application does not install one. +The host receives the persisted/authenticated user and these optional fields: + +```ts +resolveUpstreamTokenProvider(user, { + signal, + context: { scheduleId, ownerId, tenantId, agentId, invocationMode: 'delegated' }, + target: { mcpServer, scopes }, +}); +``` + +Admission allocates a proposed schedule ID before preflight and persists that same ID +only if creation succeeds. Preflight is validation, not a provisioning/consent hook: +hosts must not create durable grants keyed by this proposed ID. Failed admission or a +concurrent idempotency-key winner can discard it. Edits +and dispatch use the existing schedule ID. Execution derives the context from the +verified schedule trigger and authenticated owner, then captures it before tool loading. +`agentId` identifies the root scheduled agent throughout child execution and handoffs. +After an approval pause, the resume host restores the context from the saved job after +validating ownership, tenancy, agent identity, and schedule liveness. An explicit initializer +argument carries this restored context; resume body fields cannot replace it. +The existing resolver closure is passed through tool discovery, execution, and reconnects; +it never goes into tool arguments or durable job payloads. + +Each OBO consumer supplies its server name and configured scopes after its existing trust +check. A run shares in-flight lookups and successful providers only for identical server +and scope pairs. Failed or empty lookups may retry; cancellation belongs to the owning run, +so cancelling one child does not cancel a sibling's lookup. + +`tenantId` is absent in deployments without tenancy. `context` is absent for legacy callers +that supply no schedule ID or verified trigger, and `target` is absent for legacy consumers. +Existing callbacks may ignore the new fields. Hosts needing either field must reject its +absence. Context describes execution; it is not a consent grant or permission to mint. +Current schedules execute with their owner's authority, hence `delegated`. Dedicated agent +authorization requires a separate implementation. Scopes are not an STS audience; audience +mapping and authorization remain the host's responsibility. + +This interface adds no token store, STS exchange, consent API, or new MCP credential mode. diff --git a/packages/api/src/schedules/context.spec.ts b/packages/api/src/schedules/context.spec.ts new file mode 100644 index 00000000000..494922feea5 --- /dev/null +++ b/packages/api/src/schedules/context.spec.ts @@ -0,0 +1,158 @@ +import type { IUser } from '@librechat/data-schemas'; +import type { ScheduledTokenContext } from './context'; +import { + bindUpstreamTokenProviderResolver, + createScheduleUpstreamTokenProviderResolver, +} from './mcp'; +import { createLazyOboUpstreamTokenProvider } from '../mcp/oauth/obo'; +import { restoreScheduledTokenContext } from './context'; + +const user = { id: 'owner', tenantId: 'tenant' } as IUser; +const context: ScheduledTokenContext = { + scheduleId: 'schedule', + ownerId: 'owner', + tenantId: 'tenant', + agentId: 'root-agent', + invocationMode: 'delegated', +}; +const target = { mcpServer: 'warehouse', scopes: 'api://warehouse/.default' }; + +function request(manual = false) { + return { + user, + _isAgentTrigger: true, + body: { + agent_id: 'root-agent', + scheduleId: 'spoofed', + agentTrigger: { + version: 1, + event: { + type: 'schedule.occurrence', + occurredAt: 0, + source: { type: 'schedule', id: 'schedule' }, + }, + metadata: { manual }, + }, + }, + }; +} + +it.each([false, true])('captures verified root identity for manual=%s', async (manual) => { + const req = request(manual); + const provider = jest.fn().mockResolvedValue({ access_token: 'token' }); + const resolve = jest.fn().mockResolvedValue(provider); + const signal = new AbortController().signal; + const bound = createScheduleUpstreamTokenProviderResolver(req, resolve, signal)!; + req.body.agent_id = 'child-agent'; + req.body.agentTrigger.event.source.id = 'changed'; + await createLazyOboUpstreamTokenProvider(bound, signal, target)(); + expect(resolve).toHaveBeenCalledWith(user, { signal, context, target }); + expect(Object.isFrozen(resolve.mock.calls[0][1].context)).toBe(true); +}); + +it('ignores copied schedule metadata on an ordinary interactive request', () => { + const resolve = jest.fn(); + expect( + createScheduleUpstreamTokenProviderResolver({ ...request(), _isAgentTrigger: false }, resolve), + ).toBeUndefined(); + expect(resolve).not.toHaveBeenCalled(); +}); + +it('leaves context absent for legacy schedule classification without verified metadata', async () => { + const resolve = jest.fn().mockResolvedValue(jest.fn()); + const bound = createScheduleUpstreamTokenProviderResolver( + { ...request(), _isAgentTrigger: false, _isScheduledFire: true }, + resolve, + )!; + await bound(); + expect(resolve).toHaveBeenCalledWith(user, { signal: undefined }); +}); + +it('shares lookup for sibling consumers and reconnects, but isolates server and scope changes', async () => { + const resolve = jest.fn(async () => jest.fn(async () => ({ access_token: 'token' }))); + const bound = bindUpstreamTokenProviderResolver(user, resolve, undefined, context)!; + const targets = [ + target, + target, + { ...target, mcpServer: 'second' }, + { ...target, scopes: 'read' }, + ]; + await Promise.all( + targets.map((item) => createLazyOboUpstreamTokenProvider(bound, undefined, item)()), + ); + expect(resolve).toHaveBeenCalledTimes(3); + await createLazyOboUpstreamTokenProvider(bound, undefined, target)({ forceRefresh: true }); + expect(resolve).toHaveBeenCalledTimes(3); +}); + +it('isolates provider caches between schedules and retries only the failed target', async () => { + const resolve = jest + .fn() + .mockRejectedValueOnce(new Error('temporary')) + .mockResolvedValue(jest.fn()); + const first = bindUpstreamTokenProviderResolver(user, resolve, undefined, context)!; + const second = bindUpstreamTokenProviderResolver(user, resolve, undefined, { + ...context, + scheduleId: 'second', + })!; + await expect(first({ target })).rejects.toThrow('temporary'); + await second({ target }); + await first({ target }); + await first({ target }); + expect(resolve).toHaveBeenCalledTimes(3); + expect(resolve.mock.calls[1][1].context.scheduleId).toBe('second'); + expect(resolve.mock.calls[2][1].context.scheduleId).toBe('schedule'); +}); + +it('restores a paused run from job identity without trusting resume body fields', async () => { + const req = { + user, + _isScheduledFire: true, + body: { scheduleId: 'spoofed', agent_id: 'spoofed' }, + }; + const restored = restoreScheduledTokenContext(req, { + userId: 'owner', + tenantId: 'tenant', + scheduleId: 'schedule', + agent_id: 'root-agent', + }); + const resolve = jest.fn().mockResolvedValue(jest.fn()); + await createScheduleUpstreamTokenProviderResolver(req, resolve, undefined, restored)!({ target }); + expect(resolve).toHaveBeenCalledWith(user, { signal: undefined, context, target }); + expect(JSON.stringify(req)).not.toContain('root-agent'); +}); + +it.each([{ userId: 'different' }, { tenantId: 'different' }])( + 'rejects invalid restored identity %j', + (override) => { + expect(() => + restoreScheduledTokenContext( + { user }, + { + userId: 'owner', + tenantId: 'tenant', + scheduleId: 'schedule', + agent_id: 'root-agent', + ...override, + }, + ), + ).toThrow('Scheduled job identity'); + }, +); + +it.each([{ agent_id: undefined }, { tenantId: undefined }])( + 'leaves legacy job context absent for %j', + async (override) => { + const req = { user, _isScheduledFire: true }; + const restored = restoreScheduledTokenContext(req, { + userId: 'owner', + tenantId: 'tenant', + scheduleId: 'schedule', + agent_id: 'root-agent', + ...override, + }); + const resolve = jest.fn().mockResolvedValue(jest.fn()); + await createScheduleUpstreamTokenProviderResolver(req, resolve, undefined, restored)!(); + expect(resolve).toHaveBeenCalledWith(user, { signal: undefined }); + }, +); diff --git a/packages/api/src/schedules/context.ts b/packages/api/src/schedules/context.ts new file mode 100644 index 00000000000..0c2393bbc7b --- /dev/null +++ b/packages/api/src/schedules/context.ts @@ -0,0 +1,38 @@ +/** Root actor for an owner-authorized schedule, preserved across subagents and handoffs. */ +export interface ScheduledTokenContext { + readonly scheduleId: string; + readonly ownerId: string; + readonly tenantId?: string; + readonly agentId: string; + readonly invocationMode: 'delegated'; +} + +interface ScheduleJobIdentity { + userId?: string; + tenantId?: string; + scheduleId?: string; + agent_id?: string; +} + +/** Called by the resume host after job ownership, tenant, agent, and schedule checks. */ +export function restoreScheduledTokenContext( + req: { user: { id: string; tenantId?: string } }, + metadata?: ScheduleJobIdentity, +): ScheduledTokenContext | undefined { + if (!metadata?.scheduleId) return; + if ( + metadata.userId !== req.user.id || + (metadata.tenantId != null && metadata.tenantId !== req.user.tenantId) + ) { + throw new Error('Scheduled job identity does not match the authenticated owner.'); + } + /** Legacy jobs remain resumable, but cannot claim a complete minting context. */ + if (!metadata.agent_id || (req.user.tenantId && metadata.tenantId == null)) return; + return Object.freeze({ + scheduleId: metadata.scheduleId, + ownerId: metadata.userId, + ...(metadata.tenantId ? { tenantId: metadata.tenantId } : {}), + agentId: metadata.agent_id, + invocationMode: 'delegated', + }); +} diff --git a/packages/api/src/schedules/fire.spec.ts b/packages/api/src/schedules/fire.spec.ts index c01138a2b2c..a2073f5f50a 100644 --- a/packages/api/src/schedules/fire.spec.ts +++ b/packages/api/src/schedules/fire.spec.ts @@ -949,7 +949,11 @@ it('bounds MCP preflight by the claim lease and the stricter concurrency config' expect(preflightMCP).toHaveBeenCalledWith( 'agent-1', OWNER, - expect.objectContaining({ concurrency: 2, deadlineMs: leaseUntil.getTime() }), + expect.objectContaining({ + scheduleId: 'sched-1', + concurrency: 2, + deadlineMs: leaseUntil.getTime(), + }), ); }); diff --git a/packages/api/src/schedules/fire.ts b/packages/api/src/schedules/fire.ts index a8cea114115..64f63a8bdb6 100644 --- a/packages/api/src/schedules/fire.ts +++ b/packages/api/src/schedules/fire.ts @@ -448,6 +448,7 @@ export async function fireSchedule( Date.now() + Math.min(ownerLimits.mcpPreflightTimeoutMs, deploymentLimits.mcpPreflightTimeoutMs); mcp = await deps.preflightMCP(schedule.agent_id, user, { + scheduleId: schedule.id, signal: options?.signal, concurrency: Math.min( ownerLimits.mcpPreflightConcurrency, diff --git a/packages/api/src/schedules/handlers.spec.ts b/packages/api/src/schedules/handlers.spec.ts index 7afc4dcdf15..787b19fbf2e 100644 --- a/packages/api/src/schedules/handlers.spec.ts +++ b/packages/api/src/schedules/handlers.spec.ts @@ -264,6 +264,11 @@ describe('createSchedule late-create compensation', () => { const inserted = (deps.methods.createScheduleWithSlot as jest.Mock).mock.calls[0][0]; expect(inserted.nextRunAt).toBeUndefined(); + expect(deps.preflightMCP).toHaveBeenCalledWith( + CREATE_BODY.agent_id, + expect.objectContaining({ id: 'user-1' }), + expect.objectContaining({ scheduleId: inserted.id }), + ); // And it is never armed afterwards, because the barrier refused the create. expect(deps.methods.updateScheduleById).not.toHaveBeenCalled(); }); @@ -1185,6 +1190,11 @@ describe('updateSchedule re-enable attachment revalidation', () => { await createSchedulesHandlers(deps).updateSchedule(makeReEnableReq(), res); expect(markFilesUsed).toHaveBeenCalledWith(['file-a', 'file-b'], 'user-1'); + expect(deps.preflightMCP).toHaveBeenCalledWith( + expect.any(String), + expect.objectContaining({ id: 'user-1', tenantId: 't1' }), + expect.objectContaining({ scheduleId: 'sched-1' }), + ); expect(captured.status ?? 200).toBe(200); expect(deps.methods.updateScheduleById).toHaveBeenCalled(); }); diff --git a/packages/api/src/schedules/handlers.ts b/packages/api/src/schedules/handlers.ts index 2468ca4ed9a..1d30f50ca2e 100644 --- a/packages/api/src/schedules/handlers.ts +++ b/packages/api/src/schedules/handlers.ts @@ -387,9 +387,11 @@ export function createSchedulesHandlers(deps: SchedulesHandlersDeps): SchedulesH res: Response, signal: AbortSignal, limits: ScheduleLimits, + scheduleId: string, ): Promise { try { await deps.preflightMCP(agentId, requestUser(req), { + scheduleId, signal, concurrency: limits.mcpPreflightConcurrency, deadlineMs: Date.now() + limits.mcpPreflightTimeoutMs, @@ -705,9 +707,10 @@ export function createSchedulesHandlers(deps: SchedulesHandlersDeps): SchedulesH await respondToReplay(replayed); return; } + const id = `sched_${randomUUID()}`; if ( parsed.data.enabled && - !(await validateMCP(parsed.data.agent_id, req, res, mcpSignal, limits)) + !(await validateMCP(parsed.data.agent_id, req, res, mcpSignal, limits, id)) ) return; // Project policy applies to a NEW insert only, and is therefore resolved AFTER every @@ -751,7 +754,6 @@ export function createSchedulesHandlers(deps: SchedulesHandlersDeps): SchedulesH if (!withinIntervalFloor(res, parsed.data.cadence, parsed.data.timezone, limits)) { return; } - const id = `sched_${randomUUID()}`; const nextRunAt = parsed.data.enabled ? computeNextRunAt({ cadence: parsed.data.cadence, @@ -972,7 +974,14 @@ export function createSchedulesHandlers(deps: SchedulesHandlersDeps): SchedulesH } if ( enabled && - !(await validateMCP(parsed.data.agent_id ?? existing.agent_id, req, res, mcpSignal, limits)) + !(await validateMCP( + parsed.data.agent_id ?? existing.agent_id, + req, + res, + mcpSignal, + limits, + existing.id, + )) ) return; // The destination is re-resolved on every edit that leaves the schedule ENABLED, diff --git a/packages/api/src/schedules/index.ts b/packages/api/src/schedules/index.ts index bd702113d71..709c9a6cfd4 100644 --- a/packages/api/src/schedules/index.ts +++ b/packages/api/src/schedules/index.ts @@ -1,5 +1,6 @@ export * from './access'; export * from './cadence'; +export * from './context'; export * from './engine'; export * from './erasure'; export * from './fire'; diff --git a/packages/api/src/schedules/mcp.spec.ts b/packages/api/src/schedules/mcp.spec.ts index 4def3a3d6cd..35c20b47c14 100644 --- a/packages/api/src/schedules/mcp.spec.ts +++ b/packages/api/src/schedules/mcp.spec.ts @@ -93,8 +93,13 @@ function setup(tools = ['search_mcp_docs']) { disconnect, check: ( agentId: string, - user: typeof principal, - options?: { concurrency?: number; signal?: AbortSignal; deadlineMs?: number }, + user: typeof principal & { tenantId?: string }, + options?: { + concurrency?: number; + signal?: AbortSignal; + deadlineMs?: number; + scheduleId?: string; + }, ) => preflight(agentId, user, { concurrency: 3, ...options }), }; } @@ -157,6 +162,36 @@ it('does not resolve upstream credentials for non-OBO servers', async () => { expect(deps.resolveUpstreamTokenProvider).not.toHaveBeenCalled(); }); +it('passes admission identity to the host and rejects a changed tenant before connection', async () => { + const { check, deps } = setup(); + const user = { ...principal, tenantId: 'tenant' }; + deps.getUser = jest.fn(async () => user as IUser); + deps.resolveUpstreamTokenProvider = jest.fn(async () => + jest.fn(async () => ({ access_token: 'token' })), + ); + const connect = deps.connect; + deps.connect = jest.fn(async (options) => { + await options.upstreamTokenProviderResolver?.(); + return connect(options); + }); + await check('agent', user, { scheduleId: 'schedule' }); + expect(deps.resolveUpstreamTokenProvider).toHaveBeenCalledWith(user, { + signal: undefined, + context: { + scheduleId: 'schedule', + ownerId: 'owner', + tenantId: 'tenant', + agentId: 'agent', + invocationMode: 'delegated', + }, + }); + jest.mocked(deps.connect).mockClear(); + await expect( + check('agent', { ...user, tenantId: 'different' }, { scheduleId: 'schedule' }), + ).rejects.toBeInstanceOf(ScheduleMCPError); + expect(deps.connect).not.toHaveBeenCalled(); +}); + it('does not expose an OBO provider to a sibling direct-bearer server', async () => { const { check, deps } = setup(['search_mcp_obo', 'search_mcp_direct']); deps.getServerConfigs = jest.fn(async () => ({ diff --git a/packages/api/src/schedules/mcp.ts b/packages/api/src/schedules/mcp.ts index 63a7e4f8db3..37d3a344894 100644 --- a/packages/api/src/schedules/mcp.ts +++ b/packages/api/src/schedules/mcp.ts @@ -22,11 +22,16 @@ import type { AgentGraphAccessContext, } from '@librechat/data-schemas'; import type { TModelsConfig, ScheduleMCPStatus, ScheduleMCPOutcome } from 'librechat-data-provider'; -import type { UpstreamTokenProvider, UpstreamTokenProviderResolver } from '../mcp/oauth/obo'; +import type { + UpstreamTokenProvider, + UpstreamTokenProviderResolver, + UpstreamTokenTarget, +} from '../mcp/oauth/obo'; import type { ParsedServerConfig, UserMCPConnectionOptions } from '../mcp/types'; import type { CheckAccessParams } from '../middleware/access'; import type { MCPToolsSnapshot } from '../mcp/connection'; import type { GetAppConfigOptions } from '../app/service'; +import type { ScheduledTokenContext } from './context'; import type { ScheduleMCPPreflight } from './types'; import { MCPAuthenticationRejectedError, @@ -41,6 +46,7 @@ import { } from '../mcp/utils'; import { MCPConfigInitializationCanceledError } from '../mcp/registry/MCPServersRegistry'; import { createMCPRequestContext, cleanupMCPRequestContext } from '../mcp/request'; +import { isScheduleFireRequest, readScheduleFireContext } from './trigger'; import { getAppConfigOptionsFromUser } from '../app/service'; import { createConcurrencyLimiter } from '../utils/promise'; import { OboTokenResolutionError } from '../mcp/oauth/obo'; @@ -50,48 +56,85 @@ import { formatMCPServerTools } from '../mcp/tools'; import { checkAccess } from '../middleware/access'; import { detachOnAbort } from '../utils/promises'; import { getPluginAuthMap } from '../agents/auth'; -import { isScheduleFireRequest } from './trigger'; -type HostUpstreamTokenProviderResolver = ( - user: IUser, - options: { signal?: AbortSignal }, +export interface ScheduledTokenIdentity { + readonly id: string; + readonly tenantId?: string; + readonly role?: string; + readonly provider?: string; + readonly openidId?: string; + readonly openidIssuer?: string; +} + +export type HostUpstreamTokenProviderResolver = ( + user: ScheduledTokenIdentity, + options: { + signal?: AbortSignal; + context?: ScheduledTokenContext; + target?: UpstreamTokenTarget; + }, ) => ReturnType; /** Bind a credential lookup to the trusted principal and owning run's cancellation. */ export function bindUpstreamTokenProviderResolver( - user: IUser, + user: ScheduledTokenIdentity, resolve: HostUpstreamTokenProviderResolver | undefined, signal?: AbortSignal, + context?: ScheduledTokenContext, ): UpstreamTokenProviderResolver | undefined { if (!resolve) return undefined; - let pending: Promise | undefined; - return () => { + const capturedContext = context && Object.freeze({ ...context }); + const pending = new Map>(); + return (options) => { signal?.throwIfAborted(); - pending ??= Promise.resolve() + const target = options?.target && Object.freeze({ ...options.target }); + const key = JSON.stringify([target?.mcpServer, target?.scopes]); + const cached = pending.get(key); + if (cached) return cached; + const lookup = Promise.resolve() .then(() => { signal?.throwIfAborted(); - return resolve(user, { signal }); + return resolve(user, { + signal, + ...(capturedContext ? { context: capturedContext } : {}), + ...(target ? { target } : {}), + }); }) .then((provider) => { - if (!provider) pending = undefined; + if (!provider) pending.delete(key); return provider; }) .catch((error) => { - pending = undefined; + pending.delete(key); throw error; }); - return pending; + pending.set(key, lookup); + return lookup; }; } export function createScheduleUpstreamTokenProviderResolver( - req: Parameters[0] & { user: IUser }, + req: Parameters[0] & { user: ScheduledTokenIdentity }, resolve: HostUpstreamTokenProviderResolver | undefined, signal?: AbortSignal, + restoredContext?: ScheduledTokenContext, ): UpstreamTokenProviderResolver | undefined { - return isScheduleFireRequest(req) - ? bindUpstreamTokenProviderResolver(req.user, resolve, signal) - : undefined; + if (!isScheduleFireRequest(req)) return undefined; + if (restoredContext) + return bindUpstreamTokenProviderResolver(req.user, resolve, signal, restoredContext); + const fire = readScheduleFireContext(req); + const agentId = req.body?.agent_id; + const context: ScheduledTokenContext | undefined = + fire && typeof agentId === 'string' && agentId.trim().length > 0 + ? { + scheduleId: fire.scheduleId, + ownerId: req.user.id, + ...(req.user.tenantId ? { tenantId: req.user.tenantId } : {}), + agentId, + invocationMode: 'delegated', + } + : undefined; + return bindUpstreamTokenProviderResolver(req.user, resolve, signal, context); } // The public schedule schema caps mcpPreflightConcurrency at 10. Keep the same @@ -145,10 +188,7 @@ interface ScheduleMCPDeps { role?: string, ) => Promise>; findPluginAuthsByKeys: PluginAuthMethods['findPluginAuthsByKeys']; - resolveUpstreamTokenProvider?: ( - user: IUser, - options: { signal?: AbortSignal }, - ) => UpstreamTokenProvider | undefined | Promise; + resolveUpstreamTokenProvider?: HostUpstreamTokenProviderResolver; connect: (options: UserMCPConnectionOptions) => Promise<{ fetchToolsSnapshot: (deadlineMs?: number, signal?: AbortSignal) => Promise; }>; @@ -165,7 +205,7 @@ export function createScheduleMCPPreflight(deps: ScheduleMCPDeps): ScheduleMCPPr throwIfAborted(); const user = await deps.getUser(principal.id); throwIfAborted(); - if (!user) throw new ScheduleMCPError([]); + if (!user || user.tenantId !== principal.tenantId) throw new ScheduleMCPError([]); user.id = principal.id; let appConfig: AppConfig | undefined; const loadAppConfig = async (): Promise => { @@ -624,6 +664,15 @@ export function createScheduleMCPPreflight(deps: ScheduleMCPDeps): ScheduleMCPPr user, deps.resolveUpstreamTokenProvider, options.signal, + options.scheduleId + ? { + scheduleId: options.scheduleId, + ownerId: principal.id, + ...(user.tenantId ? { tenantId: user.tenantId } : {}), + agentId, + invocationMode: 'delegated', + } + : undefined, ); throwIfAborted(); const requestBody = { diff --git a/packages/api/src/schedules/types.ts b/packages/api/src/schedules/types.ts index 6f3f29ba1ca..1e4c659a3f2 100644 --- a/packages/api/src/schedules/types.ts +++ b/packages/api/src/schedules/types.ts @@ -329,5 +329,5 @@ export type FireableSchedule = ISchedule; export type ScheduleMCPPreflight = ( agentId: string, user: ScheduleUserContext, - options: { concurrency: number; signal?: AbortSignal; deadlineMs?: number }, + options: { concurrency: number; signal?: AbortSignal; deadlineMs?: number; scheduleId?: string }, ) => Promise; diff --git a/packages/data-provider/specs/parsers.spec.ts b/packages/data-provider/specs/parsers.spec.ts index 36d2ff12703..c1cf132f5dc 100644 --- a/packages/data-provider/specs/parsers.spec.ts +++ b/packages/data-provider/specs/parsers.spec.ts @@ -12,7 +12,7 @@ import { import { specialVariables } from '../src/config'; import { EModelEndpoint, Providers } from '../src/schemas'; import { ContentTypes } from '../src/types/runs'; -import type { TMessageContentParts } from '../src/types/assistants'; +import type { TMessageContentParts } from '../src/types/content'; import type { TUser, TConversation } from '../src/types'; // Mock dayjs module with consistent date/time values regardless of environment diff --git a/packages/data-provider/specs/stateful-code.spec.ts b/packages/data-provider/specs/stateful-code.spec.ts index 866a853ea7d..a2ed3648e40 100644 --- a/packages/data-provider/specs/stateful-code.spec.ts +++ b/packages/data-provider/specs/stateful-code.spec.ts @@ -2,7 +2,7 @@ import { STATEFUL_CODE_ENVIRONMENTS, resolveStatefulCodeEnvironment, resolveAllowedStatefulCodeEnvironments, -} from '../src/types/assistants'; +} from '../src/types/agents'; describe('stateful code environment policy', () => { it('allows every environment when deployment configuration is omitted', () => { diff --git a/packages/data-provider/src/actions.ts b/packages/data-provider/src/actions.ts index dd9a6744441..639ec3ca7fb 100644 --- a/packages/data-provider/src/actions.ts +++ b/packages/data-provider/src/actions.ts @@ -5,9 +5,9 @@ import crypto from 'crypto'; import { load } from 'js-yaml'; import type { OpenAPIV3 } from 'openapi-types'; import type { ActionMetadata, ActionMetadataRuntime } from './types/agents'; -import type { FunctionTool, Schema, Reference } from './types/assistants'; +import type { FunctionTool, Schema, Reference } from './types/tools'; import { AuthTypeEnum, AuthorizationTypeEnum } from './types/agents'; -import { Tools } from './types/assistants'; +import { Tools } from './types/tools'; export type ParametersSchema = { type: string; diff --git a/packages/data-provider/src/agentToolOptions.spec.ts b/packages/data-provider/src/agentToolOptions.spec.ts index be52b21b98b..70ac60238f9 100644 --- a/packages/data-provider/src/agentToolOptions.spec.ts +++ b/packages/data-provider/src/agentToolOptions.spec.ts @@ -1,4 +1,4 @@ -import type { AgentToolOptions } from './types/assistants'; +import type { AgentToolOptions } from './types/tools'; import { normalizeActionToolName, removeCodeExecutionCaller } from './agentToolOptions'; describe('normalizeActionToolName', () => { diff --git a/packages/data-provider/src/agentToolOptions.ts b/packages/data-provider/src/agentToolOptions.ts index 269e27a7b7b..8de99fd47bb 100644 --- a/packages/data-provider/src/agentToolOptions.ts +++ b/packages/data-provider/src/agentToolOptions.ts @@ -4,7 +4,7 @@ import { isActionTool, type AgentToolOptions, type AllowedCaller, -} from './types/assistants'; +} from './types/tools'; const actionDomainSeparatorRegex = new RegExp(actionDomainSeparator, 'g'); diff --git a/packages/data-provider/src/config.ts b/packages/data-provider/src/config.ts index 73d64858adc..cafd0e11454 100644 --- a/packages/data-provider/src/config.ts +++ b/packages/data-provider/src/config.ts @@ -31,8 +31,8 @@ import { CODE_ENVIRONMENT_DECISION_VERSION, CODE_ENVIRONMENT_MOVE_VERSION } from import { ComponentTypes, SettingTypes, OptionTypes } from './generate'; import { STATEFUL_CODE_ENVIRONMENTS } from './stateful-code'; import { specsConfigSchema, TSpecsConfig } from './models'; -import { isActionTool } from './types/assistants'; import { fileConfigSchema } from './file-config'; +import { isActionTool } from './types/tools'; import { apiBaseUrl } from './api-endpoints'; import { FileSources } from './types/files'; import { MCPServersSchema } from './mcp'; @@ -4270,6 +4270,8 @@ export enum LocalStorageKeys { PIN_WEB_SEARCH_ = 'PIN_WEB_SEARCH_', /** Pin state for Code Interpreter per conversation ID */ PIN_CODE_INTERPRETER_ = 'PIN_CODE_INTERPRETER_', + /** Key for the last selected code approval mode */ + LAST_CODE_APPROVAL_MODE = 'lastCodeApprovalMode', } export enum ForkOptions { diff --git a/packages/data-provider/src/data-service.ts b/packages/data-provider/src/data-service.ts index 60e2832b52f..cf5ade19a1c 100644 --- a/packages/data-provider/src/data-service.ts +++ b/packages/data-provider/src/data-service.ts @@ -8,6 +8,7 @@ import type { } from './types/traces'; import type { TInsightsAccessResponse, TInsightsParams, TInsightsResponse } from './types/insights'; import type { TFileConfig } from './file-config'; +import type * as tl from './types/tools'; import type * as t from './types'; import * as permissions from './accessPermissions'; import * as endpoints from './api-endpoints'; @@ -667,11 +668,11 @@ export const deleteAction = async ({ * Agents */ -export const createAgent = ({ ...data }: a.AgentCreateParams): Promise => { +export const createAgent = ({ ...data }: ag.AgentCreateParams): Promise => { return request.post(endpoints.agents({}), data); }; -export const getAgentById = ({ agent_id }: { agent_id: string }): Promise => { +export const getAgentById = ({ agent_id }: { agent_id: string }): Promise => { return request.get( endpoints.agents({ path: agent_id, @@ -679,7 +680,7 @@ export const getAgentById = ({ agent_id }: { agent_id: string }): Promise => { +export const getExpandedAgentById = ({ agent_id }: { agent_id: string }): Promise => { return request.get( endpoints.agents({ path: `${agent_id}/expanded`, @@ -687,7 +688,7 @@ export const getExpandedAgentById = ({ agent_id }: { agent_id: string }): Promis ); }; -export const getAgentVersions = ({ agent_id }: { agent_id: string }): Promise => { +export const getAgentVersions = ({ agent_id }: { agent_id: string }): Promise => { return request.get( endpoints.agents({ path: `${agent_id}/versions`, @@ -700,8 +701,8 @@ export const updateAgent = ({ data, }: { agent_id: string; - data: a.AgentUpdateParams; -}): Promise => { + data: ag.AgentUpdateParams; +}): Promise => { return request.patch( endpoints.agents({ path: agent_id, @@ -712,7 +713,7 @@ export const updateAgent = ({ export const duplicateAgent = ({ agent_id, -}: m.DuplicateAgentBody): Promise<{ agent: a.Agent; actions: ag.Action[] }> => { +}: m.DuplicateAgentBody): Promise<{ agent: ag.Agent; actions: ag.Action[] }> => { return request.post( endpoints.agents({ path: `${agent_id}/duplicate`, @@ -728,7 +729,7 @@ export const deleteAgent = ({ agent_id }: m.DeleteAgentBody): Promise => { ); }; -export const listAgents = (params: a.AgentListParams): Promise => { +export const listAgents = (params: ag.AgentListParams): Promise => { return request.get( endpoints.agents({ options: params, @@ -742,7 +743,7 @@ export const revertAgentVersion = ({ }: { agent_id: string; version_index: number; -}): Promise => request.post(endpoints.revertAgentVersion(agent_id), { version_index }); +}): Promise => request.post(endpoints.revertAgentVersion(agent_id), { version_index }); /* Marketplace */ @@ -763,7 +764,7 @@ export const getMarketplaceAgents = (params: { limit?: number; cursor?: string; promoted?: 0 | 1; -}): Promise => { +}): Promise => { return request.get( endpoints.agents({ // path: 'marketplace', @@ -877,7 +878,7 @@ export const uploadAssistantAvatar = (data: m.AssistantAvatarVariables): Promise ); }; -export const uploadAgentAvatar = (data: m.AgentAvatarVariables): Promise => { +export const uploadAgentAvatar = (data: m.AgentAvatarVariables): Promise => { return request.postMultiPart( `${endpoints.images()}/agents/${data.agent_id}/avatar`, data.formData, @@ -926,7 +927,7 @@ export const deleteFiles = async (payload: { files: f.BatchFile[]; agent_id?: string; assistant_id?: string; - tool_resource?: a.EToolResources; + tool_resource?: tl.EToolResources; }): Promise => request.deleteWithOptions(endpoints.files(), { data: payload, diff --git a/packages/data-provider/src/index.ts b/packages/data-provider/src/index.ts index 27305337887..bcc234688b2 100644 --- a/packages/data-provider/src/index.ts +++ b/packages/data-provider/src/index.ts @@ -29,6 +29,8 @@ export * from './roles'; export * from './types'; export * from './types/agents'; export * from './types/assistants'; +export * from './types/content'; +export * from './types/tools'; export * from './types/files'; export * from './types/mcpServers'; export * from './types/mutations'; diff --git a/packages/data-provider/src/messages.spec.ts b/packages/data-provider/src/messages.spec.ts index 600478fdc0c..d453f6b8562 100644 --- a/packages/data-provider/src/messages.spec.ts +++ b/packages/data-provider/src/messages.spec.ts @@ -1,4 +1,4 @@ -import type { SummaryContentPart } from './types/assistants'; +import type { SummaryContentPart } from './types/content'; import type { ParentMessage } from './messages'; import type { TFile } from './types/files'; import type { TMessage } from './types'; diff --git a/packages/data-provider/src/messages.ts b/packages/data-provider/src/messages.ts index df952ba58a1..b424f648907 100644 --- a/packages/data-provider/src/messages.ts +++ b/packages/data-provider/src/messages.ts @@ -1,4 +1,4 @@ -import type { TMessageContentParts } from './types/assistants'; +import type { TMessageContentParts } from './types/content'; import type { TFile } from './types/files'; import type { TMessage } from './types'; import { ContentTypes } from './types/runs'; diff --git a/packages/data-provider/src/models.ts b/packages/data-provider/src/models.ts index de7fb5d3337..4d14e25773a 100644 --- a/packages/data-provider/src/models.ts +++ b/packages/data-provider/src/models.ts @@ -1,5 +1,5 @@ import { z } from 'zod'; -import type { AgentSubagentsConfig } from './types/assistants'; +import type { AgentSubagentsConfig } from './types/agents'; import type { TModelSpecPreset } from './schemas'; import { EModelEndpoint, diff --git a/packages/data-provider/src/parsers.ts b/packages/data-provider/src/parsers.ts index 3de6365db7a..9f49538d268 100644 --- a/packages/data-provider/src/parsers.ts +++ b/packages/data-provider/src/parsers.ts @@ -2,7 +2,7 @@ import dayjs from 'dayjs'; import utc from 'dayjs/plugin/utc.js'; import timezonePlugin from 'dayjs/plugin/timezone.js'; import type { ZodIssue } from 'zod'; -import type * as a from './types/assistants'; +import type * as a from './types/content'; import type * as s from './schemas'; import type * as t from './types'; import { diff --git a/packages/data-provider/src/resolve-llm-delivery-path.ts b/packages/data-provider/src/resolve-llm-delivery-path.ts index 09726f74c2a..25fe9ab48fe 100644 --- a/packages/data-provider/src/resolve-llm-delivery-path.ts +++ b/packages/data-provider/src/resolve-llm-delivery-path.ts @@ -15,8 +15,8 @@ import { isMediaSupportedProvider, isDocumentSupportedProvider, } from './schemas'; -import { EToolResources } from './types/assistants'; import { normalizeEndpointName } from './utils'; +import { EToolResources } from './types/tools'; /** * The native provider a custom endpoint declares, when it declares one. A custom endpoint diff --git a/packages/data-provider/src/schemas.ts b/packages/data-provider/src/schemas.ts index 7b131053639..527f26ba95f 100644 --- a/packages/data-provider/src/schemas.ts +++ b/packages/data-provider/src/schemas.ts @@ -1,12 +1,14 @@ import { z } from 'zod'; -import type { TMessageContentParts, AgentSubagentGraph, FunctionTool } from './types/assistants'; +import type { TMessageContentParts } from './types/content'; +import type { AgentSubagentGraph } from './types/agents'; import type { SearchResultData } from './types/web'; +import type { FunctionTool } from './types/tools'; import type { TFile } from './types/files'; import { CODE_ENVIRONMENT_MODES, CODE_WORKSPACE_ID_PATTERN } from './code/workspace'; import { userSubmittedMessageFieldPathSchema } from './filters'; import { TFeedback, feedbackSchema } from './feedback'; import { CODE_APPROVAL_MODES } from './code/approval'; -import { Tools } from './types/assistants'; +import { Tools } from './types/tools'; export const isUUID = z.string().uuid(); diff --git a/packages/data-provider/src/types.ts b/packages/data-provider/src/types.ts index e619f933c0f..aa94eba897d 100644 --- a/packages/data-provider/src/types.ts +++ b/packages/data-provider/src/types.ts @@ -21,13 +21,15 @@ import type { CodeEnvironmentUserSettings, TAgentsEndpoint, } from './config'; -import type { Agent, EToolResources, StatefulCodeEnvironment } from './types/assistants'; +import type { StatefulCodeEnvironment } from './stateful-code'; import type { CodeApprovalMode } from './code/approval'; +import type { EToolResources } from './types/tools'; import type { RefillIntervalUnit } from './balance'; import type { SettingDefinition } from './generate'; import type { TMinimalFeedback } from './feedback'; import type { ContentTypes } from './types/runs'; import type { ProviderId } from './providers'; +import type { Agent } from './types/agents'; export * from './schemas'; export * from './types/subagents'; diff --git a/packages/data-provider/src/types/agents.ts b/packages/data-provider/src/types/agents.ts index e95baff7434..65c39d9f8bd 100644 --- a/packages/data-provider/src/types/agents.ts +++ b/packages/data-provider/src/types/agents.ts @@ -1,8 +1,20 @@ /* eslint-disable @typescript-eslint/no-namespace */ +import { z } from 'zod'; +import type { TAttachment, TPlugin, AgentProvider, MemoryScope, SkillsScope } from 'src/schemas'; import type { TTokenUsageEvent, TContextUsageEvent, TPendingSteer } from './runs'; -import type { TAttachment, TPlugin } from 'src/schemas'; -import type { SummaryContentPart } from './assistants'; +import type { FunctionTool, ToolResources, AgentToolOptions } from './tools'; +import type { StatefulCodeEnvironment } from '../stateful-code'; +import type { SummaryContentPart } from './content'; +import type { TFile } from './files'; import { StepTypes, ContentTypes, ToolCallTypes } from './runs'; +import { ArtifactModes } from 'src/artifacts'; +import { EToolResources } from './tools'; +export { + STATEFUL_CODE_ENVIRONMENTS, + resolveStatefulCodeEnvironment, + resolveAllowedStatefulCodeEnvironments, +} from '../stateful-code'; +export type { StatefulCodeEnvironment } from '../stateful-code'; export namespace Agents { export type MessageType = 'human' | 'ai' | 'generic' | 'system' | 'function' | 'tool' | 'remove'; @@ -835,3 +847,306 @@ export type GraphEdge = { */ promptKey?: string; }; + +/* Agent types */ + +export type AgentAvatar = { + filepath: string; + source: string; +}; + +export type AgentParameterValue = number | string | null; + +export type AgentModelParameters = { + model?: string; + temperature: AgentParameterValue; + maxContextTokens: AgentParameterValue; + max_context_tokens: AgentParameterValue; + max_output_tokens: AgentParameterValue; + top_p: AgentParameterValue; + frequency_penalty: AgentParameterValue; + presence_penalty: AgentParameterValue; + useResponsesApi?: boolean; +}; + +export interface AgentBaseResource { + /** + * A list of file IDs made available to the tool. + */ + file_ids?: Array; + /** + * A list of files already fetched. + */ + files?: Array; +} + +export interface AgentToolResources { + [EToolResources.image_edit]?: AgentBaseResource; + [EToolResources.execute_code]?: ExecuteCodeResource; + [EToolResources.file_search]?: AgentFileResource; + [EToolResources.context]?: AgentBaseResource; + /** @deprecated Use context instead */ + [EToolResources.ocr]?: AgentBaseResource; +} +/** + * A resource for the execute_code tool. + * Contains file IDs made available to the tool (max 20 files) and already fetched files. + */ +export type ExecuteCodeResource = AgentBaseResource; + +export interface AgentFileResource extends AgentBaseResource { + /** + * The ID of the vector store attached to this agent. There + * can be a maximum of 1 vector store attached to the agent. + */ + vector_store_ids?: Array; +} +export type SupportContact = { + name?: string; + email?: string; +}; + +export type AgentOwnerContact = { + name?: string; +}; + +/** + * Configuration for spawning subagents (isolated-context child agents) from an agent. + * When `enabled` is true, the agent gets a subagent-spawn tool that can delegate work + * to itself, listed single-agent targets, and/or explicit saved-agent teams. + */ +export type AgentSubagentGraphEdge = Omit< + GraphEdge, + 'edgeType' | 'condition' | 'prompt' | 'promptKey' +> & { + edgeType: 'direct'; + condition?: never; + prompt?: string; + promptKey?: never; +}; + +/** A bounded saved-agent team that can be spawned as one isolated child graph. */ +export type AgentSubagentGraph = { + /** Stable spawn-tool enum value for the team. */ + type: string; + name: string; + description: string; + /** Member IDs. In create/update payloads, an empty ID refers to the current agent. */ + agent_ids: string[]; + edges: AgentSubagentGraphEdge[]; + /** Entry member ID. In create/update payloads, an empty ID refers to the current agent. */ + entry_agent_id: string; + /** Result member ID. In create/update payloads, an empty ID refers to the current agent. */ + result_agent_id: string; +}; + +export type AgentSubagentsConfig = { + enabled?: boolean; + /** When true (default), the agent may spawn itself in an isolated context. */ + allowSelf?: boolean; + /** Share current-turn files with authorized descendants. Off unless explicitly enabled. */ + shareFiles?: boolean; + /** Specific agents that may be spawned as subagents. */ + agent_ids?: string[]; + /** Explicit saved-agent teams that may be spawned as bounded child graphs. */ + graphs?: AgentSubagentGraph[]; +}; + +export type AgentGitIdentity = { + /** Commit author and committer display name. */ + name: string; + /** Commit author and committer email address. */ + email: string; +}; + +export const agentGitIdentitySchema: z.ZodType = z + .object({ + name: z + .string() + .trim() + .min(1) + .max(128) + .refine((value) => !/[\0\r\n]/.test(value)), + email: z + .string() + .trim() + .email() + .max(254) + .refine((value) => !/[\0\r\n]/.test(value)), + }) + .optional(); + +export type Agent = { + _id?: string; + id: string; + name: string | null; + author?: string | null; + /** The original custom endpoint name, lowercased */ + endpoint?: string | null; + authorName?: string | null; + description: string | null; + created_at: number; + avatar: AgentAvatar | null; + instructions?: string | null; + additional_instructions?: string | null; + tools?: string[]; + tool_kwargs?: Record; + metadata?: Record; + provider: AgentProvider; + model: string | null; + model_parameters: AgentModelParameters; + conversation_starters?: string[]; + tool_resources?: AgentToolResources; + /** @deprecated Use edges instead */ + agent_ids?: string[]; + edges?: GraphEdge[]; + end_after_tools?: boolean; + hide_sequential_outputs?: boolean; + /** Per-agent opt-in for stateful code sessions (requires the app-level capability). */ + stateful_code_sessions?: boolean; + /** Stateful workspace sharing scope. Defaults to one workspace per user. */ + stateful_code_environment?: StatefulCodeEnvironment; + /** Operator-configured managed or attached stateful execution environment. */ + code_environment_id?: string | null; + /** Default attached workspace for new chats; empty means no agent default. */ + code_workspace_id?: string; + /** Non-secret Git authorship injected into this agent's sandboxed commands. */ + git_identity?: AgentGitIdentity | null; + artifacts?: ArtifactModes; + recursion_limit?: number; + isPublic?: boolean; + /** + * Whether the requesting user holds EDIT on this agent, so a single VIEW-scoped fetch can + * serve consumers that only need the editable subset instead of issuing a second full + * paginated walk under an EDIT-scoped cache key. + * + * Set by the list endpoint only; single-agent responses omit it. Treat absence as unknown + * and fail open (`isEditable !== false`), never as `false`, since a client on an older + * server would otherwise see an empty list rather than too many rows. + * + * Reflects the caller's ACL grant. The `MANAGE_AGENTS` capability bypasses ACL on write, + * so a capability holder can edit agents this flag reports as not editable. + */ + isEditable?: boolean; + version?: number; + category?: string; + support_contact?: SupportContact; + owner_contact?: AgentOwnerContact; + /** Per-tool configuration options (deferred loading, allowed callers, etc.) */ + tool_options?: AgentToolOptions; + /** Attached action registrations, each `${encodedDomain}${actionDelimiter}${action_id}` */ + actions?: string[]; + /** Optional allowlist of skill ObjectIds. Only applies when `skills_enabled`. */ + skills?: string[]; + /** Master toggle for skill use on this agent. `true` = active (full catalog unless + * `skills` narrows it). `false`/undefined = inactive (no skills available). */ + skills_enabled?: boolean; + /** Enables runtime skill creation without exposing an existing skill catalog. */ + skill_authoring_enabled?: boolean; + /** Explicit catalog exposure while skills are enabled. Missing preserves legacy semantics. */ + skills_scope?: SkillsScope; + /** Subagent spawning configuration — isolated-context child agents. */ + subagents?: AgentSubagentsConfig; + /** Memory partition: `agent` isolates memories per (user, agent); default shared pool */ + memory_scope?: MemoryScope; +}; + +export type TAgentsMap = Record; + +export type AgentCreateParams = { + git_identity?: AgentGitIdentity; + name?: string | null; + description?: string | null; + avatar?: AgentAvatar | null; + file_ids?: string[]; + instructions?: string | null; + tools?: Array; + provider: AgentProvider; + model: string | null; + model_parameters: AgentModelParameters; +} & Pick< + Agent, + | 'agent_ids' + | 'edges' + | 'end_after_tools' + | 'hide_sequential_outputs' + | 'stateful_code_sessions' + | 'stateful_code_environment' + | 'code_environment_id' + | 'code_workspace_id' + | 'artifacts' + | 'recursion_limit' + | 'category' + | 'support_contact' + | 'tool_options' + | 'skills' + | 'skills_enabled' + | 'skill_authoring_enabled' + | 'skills_scope' + | 'subagents' + | 'memory_scope' +>; + +export type AgentUpdateParams = { + name?: string | null; + description?: string | null; + avatar?: AgentAvatar | null; + file_ids?: string[]; + instructions?: string | null; + tools?: Array; + tool_resources?: ToolResources; + provider?: AgentProvider; + model?: string | null; + model_parameters?: AgentModelParameters; +} & Pick< + Agent, + | 'agent_ids' + | 'edges' + | 'end_after_tools' + | 'hide_sequential_outputs' + | 'stateful_code_sessions' + | 'stateful_code_environment' + | 'code_environment_id' + | 'git_identity' + | 'code_workspace_id' + | 'artifacts' + | 'recursion_limit' + | 'category' + | 'support_contact' + | 'tool_options' + | 'skills' + | 'skills_enabled' + | 'skill_authoring_enabled' + | 'skills_scope' + | 'subagents' + | 'memory_scope' +>; + +export type AgentListParams = { + limit?: number; + requiredPermission: number; + category?: string; + search?: string; + cursor?: string; + promoted?: 0 | 1; +}; + +export type AgentListResponse = { + object: string; + data: Agent[]; + first_id: string; + last_id: string; + has_more: boolean; + after?: string; +}; + +export type AgentFile = { + file_id: string; + id?: string; + temp_file_id?: string; + bytes: number; + created_at: number; + filename: string; + object: string; + purpose: 'fine-tune' | 'fine-tune-results' | 'agents' | 'agents_output'; +}; diff --git a/packages/data-provider/src/types/assistants.ts b/packages/data-provider/src/types/assistants.ts index 38b60f4e0cd..8b33fb8fb79 100644 --- a/packages/data-provider/src/types/assistants.ts +++ b/packages/data-provider/src/types/assistants.ts @@ -1,20 +1,5 @@ -import { z } from 'zod'; -import type { OpenAPIV3 } from 'openapi-types'; -import type { AssistantsEndpoint, AgentProvider, MemoryScope, SkillsScope } from 'src/schemas'; -import type { StatefulCodeEnvironment } from '../stateful-code'; -import type { ContentTypes, SubagentIdentity } from './runs'; -import type { Agents, GraphEdge } from './agents'; -import type { TFile } from './files'; -import { ArtifactModes } from 'src/artifacts'; -export { - STATEFUL_CODE_ENVIRONMENTS, - resolveStatefulCodeEnvironment, - resolveAllowedStatefulCodeEnvironments, -} from '../stateful-code'; -export type { StatefulCodeEnvironment } from '../stateful-code'; - -export type Schema = OpenAPIV3.SchemaObject & { description?: string }; -export type Reference = OpenAPIV3.ReferenceObject & { description?: string }; +import type { FunctionTool, ToolResources } from './tools'; +import type { AssistantsEndpoint } from 'src/schemas'; export type Metadata = { avatar?: string; @@ -23,73 +8,6 @@ export type Metadata = { [key: string]: unknown; }; -export enum Tools { - execute_code = 'execute_code', - code_interpreter = 'code_interpreter', - file_search = 'file_search', - web_search = 'web_search', - retrieval = 'retrieval', - function = 'function', - memory = 'memory', - ui_resources = 'ui_resources', - skill = 'skill', - read_file = 'read_file', - bash_tool = 'bash_tool', -} - -export enum EToolResources { - code_interpreter = 'code_interpreter', - execute_code = 'execute_code', - file_search = 'file_search', - image_edit = 'image_edit', - context = 'context', - ocr = 'ocr', -} - -export type Tool = { - [type: string]: Tools; -}; - -export type FunctionTool = { - type: Tools; - function?: { - description: string; - name: string; - parameters: Record; - strict?: boolean; - additionalProperties?: boolean; // must be false if strict is true https://platform.openai.com/docs/guides/structured-outputs/some-type-specific-keywords-are-not-yet-supported - }; -}; - -/** - * A set of resources that are used by the assistant's tools. The resources are - * specific to the type of tool. For example, the `code_interpreter` tool requires - * a list of file IDs, while the `file_search` tool requires a list of vector store - * IDs. - */ -export interface ToolResources { - code_interpreter?: CodeInterpreterResource; - file_search?: FileSearchResource; -} -export interface CodeInterpreterResource { - /** - * A list of [file](https://platform.openai.com/docs/api-reference/files) IDs made - * available to the `code_interpreter`` tool. There can be a maximum of 20 files - * associated with the tool. - */ - file_ids?: Array; -} - -export interface FileSearchResource { - /** - * The ID of the - * [vector store](https://platform.openai.com/docs/api-reference/vector-stores/object) - * attached to this assistant. There can be a maximum of 1 vector store attached to - * the assistant. - */ - vector_store_ids?: Array; -} - /* Assistant types */ export type Assistant = { @@ -164,476 +82,6 @@ export type File = { purpose: 'fine-tune' | 'fine-tune-results' | 'assistants' | 'assistants_output'; }; -/* Agent types */ - -export type AgentParameterValue = number | string | null; - -export type AgentModelParameters = { - model?: string; - temperature: AgentParameterValue; - maxContextTokens: AgentParameterValue; - max_context_tokens: AgentParameterValue; - max_output_tokens: AgentParameterValue; - top_p: AgentParameterValue; - frequency_penalty: AgentParameterValue; - presence_penalty: AgentParameterValue; - useResponsesApi?: boolean; -}; - -export interface AgentBaseResource { - /** - * A list of file IDs made available to the tool. - */ - file_ids?: Array; - /** - * A list of files already fetched. - */ - files?: Array; -} - -export interface AgentToolResources { - [EToolResources.image_edit]?: AgentBaseResource; - [EToolResources.execute_code]?: ExecuteCodeResource; - [EToolResources.file_search]?: AgentFileResource; - [EToolResources.context]?: AgentBaseResource; - /** @deprecated Use context instead */ - [EToolResources.ocr]?: AgentBaseResource; -} -/** - * A resource for the execute_code tool. - * Contains file IDs made available to the tool (max 20 files) and already fetched files. - */ -export type ExecuteCodeResource = AgentBaseResource; - -export interface AgentFileResource extends AgentBaseResource { - /** - * The ID of the vector store attached to this agent. There - * can be a maximum of 1 vector store attached to the agent. - */ - vector_store_ids?: Array; -} -export type SupportContact = { - name?: string; - email?: string; -}; - -export type AgentOwnerContact = { - name?: string; -}; - -/** - * Specifies who can invoke a tool. - * - 'direct': LLM can call directly - * - 'code_execution': Only callable via programmatic tool calling (PTC) - */ -export type AllowedCaller = 'direct' | 'code_execution'; - -/** - * Per-tool configuration options stored at the agent level. - * Keyed by tool_id (e.g., "search_mcp_github"). - */ -export type ToolOptions = { - /** - * If true, the tool uses deferred loading (discoverable via tool search). - * @default false - */ - defer_loading?: boolean; - /** - * Specifies who can invoke this tool. - * - 'direct': LLM can call directly (default behavior) - * - 'code_execution': Only callable via PTC sandbox - * @default ['direct'] - */ - allowed_callers?: AllowedCaller[]; - /** - * If true (and the `run_in_background` capability is enabled), the tool's - * schema gains a `run_in_background` boolean so the model can dispatch the - * call detached and poll its result via `check_background_task`. - * @default false - */ - run_in_background?: boolean; - /** - * If true (and the `tool_intents` capability is enabled), the tool's schema - * gains an `intent` string as its FIRST property — one model-authored - * sentence per call, rendered as the call's live status label. Native host - * tools default on while the capability is enabled; `false` opts one out. - * @default false - */ - describe_intent?: boolean; -}; - -/** - * Map of tool_id to its configuration options. - * Used to customize tool behavior per agent. - */ -export type AgentToolOptions = Record; - -/** - * Configuration for spawning subagents (isolated-context child agents) from an agent. - * When `enabled` is true, the agent gets a subagent-spawn tool that can delegate work - * to itself, listed single-agent targets, and/or explicit saved-agent teams. - */ -export type AgentSubagentGraphEdge = Omit< - GraphEdge, - 'edgeType' | 'condition' | 'prompt' | 'promptKey' -> & { - edgeType: 'direct'; - condition?: never; - prompt?: string; - promptKey?: never; -}; - -/** A bounded saved-agent team that can be spawned as one isolated child graph. */ -export type AgentSubagentGraph = { - /** Stable spawn-tool enum value for the team. */ - type: string; - name: string; - description: string; - /** Member IDs. In create/update payloads, an empty ID refers to the current agent. */ - agent_ids: string[]; - edges: AgentSubagentGraphEdge[]; - /** Entry member ID. In create/update payloads, an empty ID refers to the current agent. */ - entry_agent_id: string; - /** Result member ID. In create/update payloads, an empty ID refers to the current agent. */ - result_agent_id: string; -}; - -export type AgentSubagentsConfig = { - enabled?: boolean; - /** When true (default), the agent may spawn itself in an isolated context. */ - allowSelf?: boolean; - /** Share current-turn files with authorized descendants. Off unless explicitly enabled. */ - shareFiles?: boolean; - /** Specific agents that may be spawned as subagents. */ - agent_ids?: string[]; - /** Explicit saved-agent teams that may be spawned as bounded child graphs. */ - graphs?: AgentSubagentGraph[]; -}; - -export type AgentGitIdentity = { - /** Commit author and committer display name. */ - name: string; - /** Commit author and committer email address. */ - email: string; -}; - -export const agentGitIdentitySchema: z.ZodType = z - .object({ - name: z - .string() - .trim() - .min(1) - .max(128) - .refine((value) => !/[\0\r\n]/.test(value)), - email: z - .string() - .trim() - .email() - .max(254) - .refine((value) => !/[\0\r\n]/.test(value)), - }) - .optional(); - -export type Agent = { - _id?: string; - id: string; - name: string | null; - author?: string | null; - /** The original custom endpoint name, lowercased */ - endpoint?: string | null; - authorName?: string | null; - description: string | null; - created_at: number; - avatar: AgentAvatar | null; - instructions?: string | null; - additional_instructions?: string | null; - tools?: string[]; - tool_kwargs?: Record; - metadata?: Record; - provider: AgentProvider; - model: string | null; - model_parameters: AgentModelParameters; - conversation_starters?: string[]; - tool_resources?: AgentToolResources; - /** @deprecated Use edges instead */ - agent_ids?: string[]; - edges?: GraphEdge[]; - end_after_tools?: boolean; - hide_sequential_outputs?: boolean; - /** Per-agent opt-in for stateful code sessions (requires the app-level capability). */ - stateful_code_sessions?: boolean; - /** Stateful workspace sharing scope. Defaults to one workspace per user. */ - stateful_code_environment?: StatefulCodeEnvironment; - /** Operator-configured managed or attached stateful execution environment. */ - code_environment_id?: string | null; - /** Default attached workspace for new chats; empty means no agent default. */ - code_workspace_id?: string; - /** Non-secret Git authorship injected into this agent's sandboxed commands. */ - git_identity?: AgentGitIdentity | null; - artifacts?: ArtifactModes; - recursion_limit?: number; - isPublic?: boolean; - /** - * Whether the requesting user holds EDIT on this agent, so a single VIEW-scoped fetch can - * serve consumers that only need the editable subset instead of issuing a second full - * paginated walk under an EDIT-scoped cache key. - * - * Set by the list endpoint only; single-agent responses omit it. Treat absence as unknown - * and fail open (`isEditable !== false`), never as `false`, since a client on an older - * server would otherwise see an empty list rather than too many rows. - * - * Reflects the caller's ACL grant. The `MANAGE_AGENTS` capability bypasses ACL on write, - * so a capability holder can edit agents this flag reports as not editable. - */ - isEditable?: boolean; - version?: number; - category?: string; - support_contact?: SupportContact; - owner_contact?: AgentOwnerContact; - /** Per-tool configuration options (deferred loading, allowed callers, etc.) */ - tool_options?: AgentToolOptions; - /** Attached action registrations, each `${encodedDomain}${actionDelimiter}${action_id}` */ - actions?: string[]; - /** Optional allowlist of skill ObjectIds. Only applies when `skills_enabled`. */ - skills?: string[]; - /** Master toggle for skill use on this agent. `true` = active (full catalog unless - * `skills` narrows it). `false`/undefined = inactive (no skills available). */ - skills_enabled?: boolean; - /** Enables runtime skill creation without exposing an existing skill catalog. */ - skill_authoring_enabled?: boolean; - /** Explicit catalog exposure while skills are enabled. Missing preserves legacy semantics. */ - skills_scope?: SkillsScope; - /** Subagent spawning configuration — isolated-context child agents. */ - subagents?: AgentSubagentsConfig; - /** Memory partition: `agent` isolates memories per (user, agent); default shared pool */ - memory_scope?: MemoryScope; -}; - -export type TAgentsMap = Record; - -export type AgentCreateParams = { - git_identity?: AgentGitIdentity; - name?: string | null; - description?: string | null; - avatar?: AgentAvatar | null; - file_ids?: string[]; - instructions?: string | null; - tools?: Array; - provider: AgentProvider; - model: string | null; - model_parameters: AgentModelParameters; -} & Pick< - Agent, - | 'agent_ids' - | 'edges' - | 'end_after_tools' - | 'hide_sequential_outputs' - | 'stateful_code_sessions' - | 'stateful_code_environment' - | 'code_environment_id' - | 'code_workspace_id' - | 'artifacts' - | 'recursion_limit' - | 'category' - | 'support_contact' - | 'tool_options' - | 'skills' - | 'skills_enabled' - | 'skill_authoring_enabled' - | 'skills_scope' - | 'subagents' - | 'memory_scope' ->; - -export type AgentUpdateParams = { - name?: string | null; - description?: string | null; - avatar?: AgentAvatar | null; - file_ids?: string[]; - instructions?: string | null; - tools?: Array; - tool_resources?: ToolResources; - provider?: AgentProvider; - model?: string | null; - model_parameters?: AgentModelParameters; -} & Pick< - Agent, - | 'agent_ids' - | 'edges' - | 'end_after_tools' - | 'hide_sequential_outputs' - | 'stateful_code_sessions' - | 'stateful_code_environment' - | 'code_environment_id' - | 'git_identity' - | 'code_workspace_id' - | 'artifacts' - | 'recursion_limit' - | 'category' - | 'support_contact' - | 'tool_options' - | 'skills' - | 'skills_enabled' - | 'skill_authoring_enabled' - | 'skills_scope' - | 'subagents' - | 'memory_scope' ->; - -export type AgentListParams = { - limit?: number; - requiredPermission: number; - category?: string; - search?: string; - cursor?: string; - promoted?: 0 | 1; -}; - -export type AgentListResponse = { - object: string; - data: Agent[]; - first_id: string; - last_id: string; - has_more: boolean; - after?: string; -}; - -export type AgentFile = { - file_id: string; - id?: string; - temp_file_id?: string; - bytes: number; - created_at: number; - filename: string; - object: string; - purpose: 'fine-tune' | 'fine-tune-results' | 'agents' | 'agents_output'; -}; - -/** - * Details of the Code Interpreter tool call the run step was involved in. - * Includes the tool call ID, the code interpreter definition, and the type of tool call. - */ -export type CodeToolCall = { - id: string; // The ID of the tool call. - code_interpreter: { - input: string; // The input to the Code Interpreter tool call. - outputs: Array>; // The outputs from the Code Interpreter tool call. - }; - type: 'code_interpreter'; // The type of tool call, always 'code_interpreter'. -}; - -/** - * Details of a Function tool call the run step was involved in. - * Includes the tool call ID, the function definition, and the type of tool call. - */ -export type FunctionToolCall = { - id: string; // The ID of the tool call object. - function: { - arguments: string; // The arguments passed to the function. - name: string; // The name of the function. - output: string | null; // The output of the function, null if not submitted. - }; - type: 'function'; // The type of tool call, always 'function'. -}; - -/** - * Details of a Retrieval tool call the run step was involved in. - * Includes the tool call ID and the type of tool call. - */ -export type RetrievalToolCall = { - id: string; // The ID of the tool call object. - retrieval: unknown; // An empty object for now. - type: 'retrieval'; // The type of tool call, always 'retrieval'. -}; - -/** - * Details of a Retrieval tool call the run step was involved in. - * Includes the tool call ID and the type of tool call. - */ -export type FileSearchToolCall = { - id: string; // The ID of the tool call object. - file_search: unknown; // An empty object for now. - type: 'file_search'; // The type of tool call, always 'retrieval'. -}; - -/** - * Details of the tool calls involved in a run step. - * Can be associated with one of three types of tools: `code_interpreter`, `retrieval`, or `function`. - */ -export type ToolCallsStepDetails = { - tool_calls: Array; // An array of tool calls the run step was involved in. - type: 'tool_calls'; // Always 'tool_calls'. -}; - -export type ImageFile = TFile & { - /** - * The [File](https://platform.openai.com/docs/api-reference/files) ID of the image - * in the message content. - */ - file_id: string; - filename: string; - filepath: string; - height: number; - width: number; - /** - * Prompt used to generate the image if applicable. - */ - prompt?: string; - /** - * Additional metadata used to generate or about the image/tool_call. - */ - metadata?: Record; -}; - -// FileCitation.ts -export type FileCitation = { - end_index: number; - file_citation: FileCitationDetails; - start_index: number; - text: string; - type: 'file_citation'; -}; - -export type FileCitationDetails = { - file_id: string; - quote: string; -}; - -export type FilePath = { - end_index: number; - file_path: FilePathDetails; - start_index: number; - text: string; - type: 'file_path'; -}; - -export type FilePathDetails = { - file_id: string; -}; - -export type Text = { - annotations?: Array; - value: string; -}; - -export enum AnnotationTypes { - FILE_CITATION = 'file_citation', - FILE_PATH = 'file_path', -} - -export enum StepStatus { - IN_PROGRESS = 'in_progress', - CANCELLED = 'cancelled', - FAILED = 'failed', - COMPLETED = 'completed', - EXPIRED = 'expired', -} - -export enum MessageContentTypes { - TEXT = 'text', - IMAGE_FILE = 'image_file', -} - //enum for RunStatus // The status of the run: queued, in_progress, requires_action, cancelling, cancelled, failed, completed, or expired. export enum RunStatus { @@ -647,242 +95,6 @@ export enum RunStatus { EXPIRED = 'expired', } -export type PartMetadata = { - /** Host-resolved execution identity for a saved subagent invocation. */ - subagentIdentity?: SubagentIdentity; - progress?: number; - asset_pointer?: string; - status?: string; - action?: boolean; - auth?: string; - expires_at?: number; - /** Index indicating parallel sibling content (same stepIndex in multi-agent runs) */ - siblingIndex?: number; - /** Agent ID for parallel agent rendering - identifies which agent produced this content */ - agentId?: string; - /** Group ID for parallel content - parts with same groupId are displayed in columns */ - groupId?: number; - /** - * Terminal lifecycle status of the run step that produced this part, from - * `on_run_step_closed`. Distinct from `status`, which is already claimed by - * activity-label and question-form parts. Absent on parts predating the - * event or from endpoints that do not emit it, in which case renderers fall - * back to inferring "stopped" from `progress` and `isSubmitting`. - */ - runStepStatus?: Agents.RunStepClosedStatus; - /** - * Wall-clock milliseconds the run step took, derived from the same - * `on_run_step_closed` event as {@link runStepStatus} via - * `getRunStepDurationMs`. Only written when the event carried both - * timestamps and they agree in order — so its absence means "not - * derivable", never "instant". The raw value is persisted unfiltered; - * whether it is worth showing (`isReportableRunStepDuration`) is decided - * at render time. - */ - runStepDurationMs?: number; - /** - * Stamped by the background harvester when a detached task's final output - * replaces the dispatch handle in `tool_call.output`. The handle JSON and - * the live status-marker attachment are both transient, so after the patch - * (or a reload) this is the only signal that the call ran in the - * background — renderers use it to keep treating {@link runStepDurationMs} - * as dispatch time rather than the task's runtime. - */ - backgrounded?: boolean; - /** - * Content index this part occupied while its run streamed. The aggregator - * writes parts at provider-source indexes, so the streamed array is sparse; - * persistence compacts it and every part after a hole shifts down. The - * client's final handler stamps the streamed position onto the compacted - * parts it adopts, so index-derived render identity survives the swap - * instead of remounting the settled message. Client-only and absent - * everywhere else — persisted content never carries it. - */ - streamedIndex?: number; -}; - -/** Metadata for parallel content rendering - subset of PartMetadata */ -export type ContentMetadata = Pick; - -export type ContentPart = ( - | CodeToolCall - | RetrievalToolCall - | FileSearchToolCall - | FunctionToolCall - | Agents.AgentToolCall - | ImageFile - | Text -) & - PartMetadata; - -export type TextData = (Text & PartMetadata) | undefined; - -export type SummaryContentPart = { - type: ContentTypes.SUMMARY; - content?: Array<{ type: ContentTypes.TEXT; text: string }>; - tokenCount?: number; - summarizing?: boolean; - /** A summarize round that ended in error. Partial deltas already streamed - * into this slot are kept, so the renderer needs this to avoid presenting - * truncated text under the "Conversation summarized" label. */ - failed?: boolean; - /** Set when the user compacted the context manually rather than the - * automatic detour firing on context pressure. */ - initiatedBy?: 'user'; - summaryVersion?: number; - model?: string; - provider?: string; - createdAt?: string; - boundary?: { - messageId: string; - contentIndex: number; - }; -}; - -/** - * A user steering message injected mid-run at a tool-batch boundary. - * Persisted inline in the response message's content array (keyed by the - * type name like `text`/`think` so token counting reads it for free); - * replayed as a user message on subsequent turns by `formatAgentMessages`. - */ -export type SteerContentPart = { - type: ContentTypes.STEER; - steer: string; - steerId?: string; - /** Stable optimistic-client id used to settle a POST whose response was lost. */ - clientSteerId?: string; - createdAt?: number; - /** Attachments steered with the message; re-encoded per turn on replay - * like any other user-message media (refs only, never encoded data). */ - files?: Partial[]; - /** Quoted excerpts steered with the message, persisted separately from the - * typed text (mirroring `TMessage.quotes`) so the UI renders them as - * reference blocks; merged into the model-bound user turn on every replay. */ - quotes?: string[]; -}; - -export type TMessageContentParts = - | ({ - type: ContentTypes.ERROR; - text?: string | TextData; - error?: string; - /** Set when this failure is what a manual compaction produced instead of a - * summary. The turn has no other record of having been one, so the rerun - * controls read it the same way they read a summary's marker. */ - initiatedBy?: 'user'; - } & ContentMetadata) - | ({ - type: ContentTypes.THINK; - think?: string | TextData; - /** Generated orientation for this user-visible reasoning step. */ - reasoning_label?: string; - /** Stable SDK run-step identity used to correlate live revisions. */ - reasoning_label_step_id?: string; - /** Durable provider-call count used to enforce the per-run cost cap across resumes. */ - reasoning_label_attempts?: number; - /** Visible reasoning length included in this step's latest provider call. */ - reasoning_label_submitted_chars?: number; - /** Monotonic provider-call revision; gaps are allowed after unsuccessful attempts. */ - reasoning_label_revision?: number; - /** Whether the reasoning step can still produce a newer label. */ - reasoning_label_status?: 'streaming' | 'complete'; - /** The reasoning happened but its text is not available to this view - * (e.g. detached subagent projections retain only a marker). */ - reasoning_unavailable?: boolean; - } & ContentMetadata) - | (SteerContentPart & ContentMetadata) - | ({ - type: ContentTypes.TEXT; - text?: string | TextData; - tool_call_ids?: string[]; - /** Open Responses semantic channel for assistant text. */ - phase?: 'commentary' | 'final_answer'; - } & ContentMetadata) - | ({ - type: ContentTypes.TOOL_CALL; - tool_call: ( - | CodeToolCall - | RetrievalToolCall - | FileSearchToolCall - | FunctionToolCall - | Agents.AgentToolCall - ) & - PartMetadata; - } & ContentMetadata) - | ({ type: ContentTypes.IMAGE_FILE; image_file: ImageFile & PartMetadata } & ContentMetadata) - | (SummaryContentPart & ContentMetadata) - | ({ - /** One-line LLM-generated note describing a completed tool batch. UI-only: - * never sent to the model (stripped before payload formatting). */ - type: ContentTypes.ACTIVITY_LABEL; - activity_label?: string; - /** Missing means the legacy/per-batch activity label. */ - activity_label_type?: 'phase'; - tool_call_ids?: string[]; - /** Parent phase bounds and telemetry. */ - activity_start_index?: number; - /** Exclusive end of the grouped content; may precede the marker itself. */ - activity_end_index?: number; - activity_count?: number; - agent_ids?: string[]; - /** ok = all tools succeeded, failed = all failed, partial = mixed. */ - status?: 'ok' | 'partial' | 'failed'; - pending?: boolean; - } & ContentMetadata) - | (Agents.AgentUpdate & ContentMetadata) - | (Agents.MessageContentImageUrl & ContentMetadata) - | (Agents.MessageContentVideoUrl & ContentMetadata) - | (Agents.MessageContentInputAudio & ContentMetadata); - -export type StreamContentData = TMessageContentParts & { - /** The index of the current content part */ - index: number; - /** The current text content was already served but edited to replace elements therein */ - edited?: boolean; -}; - -export type TContentData = StreamContentData & { - messageId: string; - conversationId: string; - userMessageId: string; - thread_id: string; - stream?: boolean; -}; - -export const actionDelimiter = '_action_'; -export const actionDomainSeparator = '---'; -/** Mirrors `Constants.mcp_delimiter`; duplicated here to avoid a circular import from `config.ts`. */ -const mcpDelimiter = '_mcp_'; - -/** - * Checks whether a tool name is an OpenAPI action tool. - * - * Action format: `operationId_action_normalizedDomain` - * MCP format: `toolName_mcp_serverName` - * - * Cross-delimiter collision: an MCP tool like `get_action_mcp_srv` contains - * `_action_` as a false positive. Guarded by checking whether `_mcp_` appears - * after `_action_`. In the collision case the `_mcp_` suffix always follows - * `_action_`; in a valid action tool whose operationId contains `_mcp_`, the - * `_mcp_` precedes `_action_`. - * - * Theoretical limitation: a non-RFC-compliant domain containing literal - * underscores that form `_mcp_` (e.g. `api_mcp_internal.com`) would produce - * a false negative. RFC 952/1123 prohibit underscores in hostnames, so this - * is not expected in practice. - */ -export function isActionTool(toolName: string): boolean { - const actionIdx = toolName.indexOf(actionDelimiter); - if (actionIdx < 0) { - return false; - } - const mcpIdx = toolName.indexOf(mcpDelimiter); - return mcpIdx < 0 || mcpIdx < actionIdx; -} - -export const hostImageIdSuffix = '_host_copy'; -export const hostImageNamePrefix = 'host_copy_'; - export type AssistantAvatar = { filepath: string; source: string; @@ -901,13 +113,6 @@ export type AssistantDocument = { append_current_datetime?: boolean; }; -/* Agent types */ - -export type AgentAvatar = { - filepath: string; - source: string; -}; - export enum FilePurpose { Vision = 'vision', FineTune = 'fine-tune', diff --git a/packages/data-provider/src/types/content.ts b/packages/data-provider/src/types/content.ts new file mode 100644 index 00000000000..bdeeac8c7e7 --- /dev/null +++ b/packages/data-provider/src/types/content.ts @@ -0,0 +1,333 @@ +import type { SubagentIdentity, ContentTypes } from './runs'; +import type { Agents } from './agents'; +import type { TFile } from './files'; + +/** + * Details of the Code Interpreter tool call the run step was involved in. + * Includes the tool call ID, the code interpreter definition, and the type of tool call. + */ +export type CodeToolCall = { + id: string; // The ID of the tool call. + code_interpreter: { + input: string; // The input to the Code Interpreter tool call. + outputs: Array>; // The outputs from the Code Interpreter tool call. + }; + type: 'code_interpreter'; // The type of tool call, always 'code_interpreter'. +}; + +/** + * Details of a Function tool call the run step was involved in. + * Includes the tool call ID, the function definition, and the type of tool call. + */ +export type FunctionToolCall = { + id: string; // The ID of the tool call object. + function: { + arguments: string; // The arguments passed to the function. + name: string; // The name of the function. + output: string | null; // The output of the function, null if not submitted. + }; + type: 'function'; // The type of tool call, always 'function'. +}; + +/** + * Details of a Retrieval tool call the run step was involved in. + * Includes the tool call ID and the type of tool call. + */ +export type RetrievalToolCall = { + id: string; // The ID of the tool call object. + retrieval: unknown; // An empty object for now. + type: 'retrieval'; // The type of tool call, always 'retrieval'. +}; + +/** + * Details of a Retrieval tool call the run step was involved in. + * Includes the tool call ID and the type of tool call. + */ +export type FileSearchToolCall = { + id: string; // The ID of the tool call object. + file_search: unknown; // An empty object for now. + type: 'file_search'; // The type of tool call, always 'retrieval'. +}; + +/** + * Details of the tool calls involved in a run step. + * Can be associated with one of three types of tools: `code_interpreter`, `retrieval`, or `function`. + */ +export type ToolCallsStepDetails = { + tool_calls: Array; // An array of tool calls the run step was involved in. + type: 'tool_calls'; // Always 'tool_calls'. +}; + +export type ImageFile = TFile & { + /** + * The [File](https://platform.openai.com/docs/api-reference/files) ID of the image + * in the message content. + */ + file_id: string; + filename: string; + filepath: string; + height: number; + width: number; + /** + * Prompt used to generate the image if applicable. + */ + prompt?: string; + /** + * Additional metadata used to generate or about the image/tool_call. + */ + metadata?: Record; +}; + +// FileCitation.ts +export type FileCitation = { + end_index: number; + file_citation: FileCitationDetails; + start_index: number; + text: string; + type: 'file_citation'; +}; + +export type FileCitationDetails = { + file_id: string; + quote: string; +}; + +export type FilePath = { + end_index: number; + file_path: FilePathDetails; + start_index: number; + text: string; + type: 'file_path'; +}; + +export type FilePathDetails = { + file_id: string; +}; + +export type Text = { + annotations?: Array; + value: string; +}; + +export enum AnnotationTypes { + FILE_CITATION = 'file_citation', + FILE_PATH = 'file_path', +} + +export enum StepStatus { + IN_PROGRESS = 'in_progress', + CANCELLED = 'cancelled', + FAILED = 'failed', + COMPLETED = 'completed', + EXPIRED = 'expired', +} + +export enum MessageContentTypes { + TEXT = 'text', + IMAGE_FILE = 'image_file', +} + +export type PartMetadata = { + /** Host-resolved execution identity for a saved subagent invocation. */ + subagentIdentity?: SubagentIdentity; + progress?: number; + asset_pointer?: string; + status?: string; + action?: boolean; + auth?: string; + expires_at?: number; + /** Index indicating parallel sibling content (same stepIndex in multi-agent runs) */ + siblingIndex?: number; + /** Agent ID for parallel agent rendering - identifies which agent produced this content */ + agentId?: string; + /** Group ID for parallel content - parts with same groupId are displayed in columns */ + groupId?: number; + /** + * Terminal lifecycle status of the run step that produced this part, from + * `on_run_step_closed`. Distinct from `status`, which is already claimed by + * activity-label and question-form parts. Absent on parts predating the + * event or from endpoints that do not emit it, in which case renderers fall + * back to inferring "stopped" from `progress` and `isSubmitting`. + */ + runStepStatus?: Agents.RunStepClosedStatus; + /** + * Wall-clock milliseconds the run step took, derived from the same + * `on_run_step_closed` event as {@link runStepStatus} via + * `getRunStepDurationMs`. Only written when the event carried both + * timestamps and they agree in order — so its absence means "not + * derivable", never "instant". The raw value is persisted unfiltered; + * whether it is worth showing (`isReportableRunStepDuration`) is decided + * at render time. + */ + runStepDurationMs?: number; + /** + * Stamped by the background harvester when a detached task's final output + * replaces the dispatch handle in `tool_call.output`. The handle JSON and + * the live status-marker attachment are both transient, so after the patch + * (or a reload) this is the only signal that the call ran in the + * background — renderers use it to keep treating {@link runStepDurationMs} + * as dispatch time rather than the task's runtime. + */ + backgrounded?: boolean; + /** + * Content index this part occupied while its run streamed. The aggregator + * writes parts at provider-source indexes, so the streamed array is sparse; + * persistence compacts it and every part after a hole shifts down. The + * client's final handler stamps the streamed position onto the compacted + * parts it adopts, so index-derived render identity survives the swap + * instead of remounting the settled message. Client-only and absent + * everywhere else — persisted content never carries it. + */ + streamedIndex?: number; +}; + +/** Metadata for parallel content rendering - subset of PartMetadata */ +export type ContentMetadata = Pick; + +export type ContentPart = ( + | CodeToolCall + | RetrievalToolCall + | FileSearchToolCall + | FunctionToolCall + | Agents.AgentToolCall + | ImageFile + | Text +) & + PartMetadata; + +export type TextData = (Text & PartMetadata) | undefined; + +export type SummaryContentPart = { + type: ContentTypes.SUMMARY; + content?: Array<{ type: ContentTypes.TEXT; text: string }>; + tokenCount?: number; + summarizing?: boolean; + /** A summarize round that ended in error. Partial deltas already streamed + * into this slot are kept, so the renderer needs this to avoid presenting + * truncated text under the "Conversation summarized" label. */ + failed?: boolean; + /** Set when the user compacted the context manually rather than the + * automatic detour firing on context pressure. */ + initiatedBy?: 'user'; + summaryVersion?: number; + model?: string; + provider?: string; + createdAt?: string; + boundary?: { + messageId: string; + contentIndex: number; + }; +}; + +/** + * A user steering message injected mid-run at a tool-batch boundary. + * Persisted inline in the response message's content array (keyed by the + * type name like `text`/`think` so token counting reads it for free); + * replayed as a user message on subsequent turns by `formatAgentMessages`. + */ +export type SteerContentPart = { + type: ContentTypes.STEER; + steer: string; + steerId?: string; + /** Stable optimistic-client id used to settle a POST whose response was lost. */ + clientSteerId?: string; + createdAt?: number; + /** Attachments steered with the message; re-encoded per turn on replay + * like any other user-message media (refs only, never encoded data). */ + files?: Partial[]; + /** Quoted excerpts steered with the message, persisted separately from the + * typed text (mirroring `TMessage.quotes`) so the UI renders them as + * reference blocks; merged into the model-bound user turn on every replay. */ + quotes?: string[]; +}; + +export type TMessageContentParts = + | ({ + type: ContentTypes.ERROR; + text?: string | TextData; + error?: string; + /** Set when this failure is what a manual compaction produced instead of a + * summary. The turn has no other record of having been one, so the rerun + * controls read it the same way they read a summary's marker. */ + initiatedBy?: 'user'; + } & ContentMetadata) + | ({ + type: ContentTypes.THINK; + think?: string | TextData; + /** Generated orientation for this user-visible reasoning step. */ + reasoning_label?: string; + /** Stable SDK run-step identity used to correlate live revisions. */ + reasoning_label_step_id?: string; + /** Durable provider-call count used to enforce the per-run cost cap across resumes. */ + reasoning_label_attempts?: number; + /** Visible reasoning length included in this step's latest provider call. */ + reasoning_label_submitted_chars?: number; + /** Monotonic provider-call revision; gaps are allowed after unsuccessful attempts. */ + reasoning_label_revision?: number; + /** Whether the reasoning step can still produce a newer label. */ + reasoning_label_status?: 'streaming' | 'complete'; + /** The reasoning happened but its text is not available to this view + * (e.g. detached subagent projections retain only a marker). */ + reasoning_unavailable?: boolean; + } & ContentMetadata) + | (SteerContentPart & ContentMetadata) + | ({ + type: ContentTypes.TEXT; + text?: string | TextData; + tool_call_ids?: string[]; + /** Open Responses semantic channel for assistant text. */ + phase?: 'commentary' | 'final_answer'; + } & ContentMetadata) + | ({ + type: ContentTypes.TOOL_CALL; + tool_call: ( + | CodeToolCall + | RetrievalToolCall + | FileSearchToolCall + | FunctionToolCall + | Agents.AgentToolCall + ) & + PartMetadata; + } & ContentMetadata) + | ({ type: ContentTypes.IMAGE_FILE; image_file: ImageFile & PartMetadata } & ContentMetadata) + | (SummaryContentPart & ContentMetadata) + | ({ + /** One-line LLM-generated note describing a completed tool batch. UI-only: + * never sent to the model (stripped before payload formatting). */ + type: ContentTypes.ACTIVITY_LABEL; + activity_label?: string; + /** Missing means the legacy/per-batch activity label. */ + activity_label_type?: 'phase'; + tool_call_ids?: string[]; + /** Parent phase bounds and telemetry. */ + activity_start_index?: number; + /** Exclusive end of the grouped content; may precede the marker itself. */ + activity_end_index?: number; + activity_count?: number; + agent_ids?: string[]; + /** ok = all tools succeeded, failed = all failed, partial = mixed. */ + status?: 'ok' | 'partial' | 'failed'; + pending?: boolean; + } & ContentMetadata) + | (Agents.AgentUpdate & ContentMetadata) + | (Agents.MessageContentImageUrl & ContentMetadata) + | (Agents.MessageContentVideoUrl & ContentMetadata) + | (Agents.MessageContentInputAudio & ContentMetadata); + +export type StreamContentData = TMessageContentParts & { + /** The index of the current content part */ + index: number; + /** The current text content was already served but edited to replace elements therein */ + edited?: boolean; +}; + +export type TContentData = StreamContentData & { + messageId: string; + conversationId: string; + userMessageId: string; + thread_id: string; + stream?: boolean; +}; + +export const hostImageIdSuffix = '_host_copy'; +export const hostImageNamePrefix = 'host_copy_'; diff --git a/packages/data-provider/src/types/files.ts b/packages/data-provider/src/types/files.ts index e7a80d81d6c..2651cc6ddb2 100644 --- a/packages/data-provider/src/types/files.ts +++ b/packages/data-provider/src/types/files.ts @@ -1,6 +1,6 @@ import type { TDefaultLLMDeliveryPathConfig } from '../file-config'; import type { CodeEnvRef, CodeEnvRefMap } from '../codeEnvRef'; -import { EToolResources } from './assistants'; +import { EToolResources } from './tools'; export enum FileSources { local = 'local', diff --git a/packages/data-provider/src/types/mutations.ts b/packages/data-provider/src/types/mutations.ts index 414d479b5c1..a99e5a45e79 100644 --- a/packages/data-provider/src/types/mutations.ts +++ b/packages/data-provider/src/types/mutations.ts @@ -12,17 +12,13 @@ import type { TSkillListResponse, } from './skills'; import { - Tools, Assistant, AssistantCreateParams, AssistantUpdateParams, - FunctionTool, AssistantDocument, - Agent, - AgentCreateParams, - AgentUpdateParams, } from './assistants'; -import { Action, ActionMetadata } from './agents'; +import { Action, ActionMetadata, Agent, AgentCreateParams, AgentUpdateParams } from './agents'; +import { Tools, FunctionTool } from './tools'; import * as p from '../permissions'; import * as types from '../types'; import * as r from '../roles'; diff --git a/packages/data-provider/src/types/tools.ts b/packages/data-provider/src/types/tools.ts new file mode 100644 index 00000000000..26db147d12e --- /dev/null +++ b/packages/data-provider/src/types/tools.ts @@ -0,0 +1,149 @@ +import type { OpenAPIV3 } from 'openapi-types'; + +export type Schema = OpenAPIV3.SchemaObject & { description?: string }; +export type Reference = OpenAPIV3.ReferenceObject & { description?: string }; + +export enum Tools { + execute_code = 'execute_code', + code_interpreter = 'code_interpreter', + file_search = 'file_search', + web_search = 'web_search', + retrieval = 'retrieval', + function = 'function', + memory = 'memory', + ui_resources = 'ui_resources', + skill = 'skill', + read_file = 'read_file', + bash_tool = 'bash_tool', +} + +export enum EToolResources { + code_interpreter = 'code_interpreter', + execute_code = 'execute_code', + file_search = 'file_search', + image_edit = 'image_edit', + context = 'context', + ocr = 'ocr', +} + +export type Tool = { + [type: string]: Tools; +}; + +export type FunctionTool = { + type: Tools; + function?: { + description: string; + name: string; + parameters: Record; + strict?: boolean; + additionalProperties?: boolean; // must be false if strict is true https://platform.openai.com/docs/guides/structured-outputs/some-type-specific-keywords-are-not-yet-supported + }; +}; + +/** + * A set of resources that are used by the assistant's tools. The resources are + * specific to the type of tool. For example, the `code_interpreter` tool requires + * a list of file IDs, while the `file_search` tool requires a list of vector store + * IDs. + */ +export interface ToolResources { + code_interpreter?: CodeInterpreterResource; + file_search?: FileSearchResource; +} +export interface CodeInterpreterResource { + /** + * A list of [file](https://platform.openai.com/docs/api-reference/files) IDs made + * available to the `code_interpreter`` tool. There can be a maximum of 20 files + * associated with the tool. + */ + file_ids?: Array; +} + +export interface FileSearchResource { + /** + * The ID of the + * [vector store](https://platform.openai.com/docs/api-reference/vector-stores/object) + * attached to this assistant. There can be a maximum of 1 vector store attached to + * the assistant. + */ + vector_store_ids?: Array; +} + +/** + * Specifies who can invoke a tool. + * - 'direct': LLM can call directly + * - 'code_execution': Only callable via programmatic tool calling (PTC) + */ +export type AllowedCaller = 'direct' | 'code_execution'; + +/** + * Per-tool configuration options stored at the agent level. + * Keyed by tool_id (e.g., "search_mcp_github"). + */ +export type ToolOptions = { + /** + * If true, the tool uses deferred loading (discoverable via tool search). + * @default false + */ + defer_loading?: boolean; + /** + * Specifies who can invoke this tool. + * - 'direct': LLM can call directly (default behavior) + * - 'code_execution': Only callable via PTC sandbox + * @default ['direct'] + */ + allowed_callers?: AllowedCaller[]; + /** + * If true (and the `run_in_background` capability is enabled), the tool's + * schema gains a `run_in_background` boolean so the model can dispatch the + * call detached and poll its result via `check_background_task`. + * @default false + */ + run_in_background?: boolean; + /** + * If true (and the `tool_intents` capability is enabled), the tool's schema + * gains an `intent` string as its FIRST property — one model-authored + * sentence per call, rendered as the call's live status label. Native host + * tools default on while the capability is enabled; `false` opts one out. + * @default false + */ + describe_intent?: boolean; +}; + +/** + * Map of tool_id to its configuration options. + * Used to customize tool behavior per agent. + */ +export type AgentToolOptions = Record; + +export const actionDelimiter = '_action_'; +export const actionDomainSeparator = '---'; +/** Mirrors `Constants.mcp_delimiter`; duplicated here to avoid a circular import from `config.ts`. */ +const mcpDelimiter = '_mcp_'; + +/** + * Checks whether a tool name is an OpenAPI action tool. + * + * Action format: `operationId_action_normalizedDomain` + * MCP format: `toolName_mcp_serverName` + * + * Cross-delimiter collision: an MCP tool like `get_action_mcp_srv` contains + * `_action_` as a false positive. Guarded by checking whether `_mcp_` appears + * after `_action_`. In the collision case the `_mcp_` suffix always follows + * `_action_`; in a valid action tool whose operationId contains `_mcp_`, the + * `_mcp_` precedes `_action_`. + * + * Theoretical limitation: a non-RFC-compliant domain containing literal + * underscores that form `_mcp_` (e.g. `api_mcp_internal.com`) would produce + * a false negative. RFC 952/1123 prohibit underscores in hostnames, so this + * is not expected in practice. + */ +export function isActionTool(toolName: string): boolean { + const actionIdx = toolName.indexOf(actionDelimiter); + if (actionIdx < 0) { + return false; + } + const mcpIdx = toolName.indexOf(mcpDelimiter); + return mcpIdx < 0 || mcpIdx < actionIdx; +} diff --git a/packages/data-provider/src/upload.ts b/packages/data-provider/src/upload.ts index 0edb86d5df6..797a04b027c 100644 --- a/packages/data-provider/src/upload.ts +++ b/packages/data-provider/src/upload.ts @@ -1,4 +1,4 @@ -import type { EToolResources } from './types/assistants'; +import type { EToolResources } from './types/tools'; import type { TFileUpload } from './types/files'; import request from './request'; diff --git a/packages/data-schemas/src/methods/conversation.spec.ts b/packages/data-schemas/src/methods/conversation.spec.ts index 00e3dae652b..5fa9ae689f9 100644 --- a/packages/data-schemas/src/methods/conversation.spec.ts +++ b/packages/data-schemas/src/methods/conversation.spec.ts @@ -1079,6 +1079,88 @@ describe('Conversation Operations', () => { }); }); + describe('code environment persistence during ordinary saves', () => { + const mac = { environmentId: 'code-mac', workspaceId: 'primary' }; + const vm = { environmentId: 'code-vm', workspaceId: 'primary' }; + const userId = 'user123'; + const unsetFields = { codeEnvironmentMode: 1, codeWorkspaces: 1 }; + + it.each([ + { codeEnvironmentMode: 'attached', codeWorkspaces: [mac] }, + { codeEnvironmentMode: 'without_attached' }, + ])( + 'preserves $codeEnvironmentMode through an approval pause and later saves', + async (decision) => { + const conversationId = uuidv4(); + await saveConvo({ userId }, { conversationId, ...decision }); + + for (const title of ['Approval required', 'Resumed', 'Completed']) { + await saveConvo({ userId }, { conversationId, title }, { unsetFields }); + const stored = await getConvo(userId, conversationId); + expect(stored?.codeEnvironmentMode).toBe(decision.codeEnvironmentMode); + expect(stored?.codeWorkspaces).toEqual(decision.codeWorkspaces); + expect(stored?.title).toBe(title); + } + }, + ); + + it('cannot overwrite an explicit move with a stale turn-start snapshot', async () => { + const conversationId = uuidv4(); + const decision = { codeEnvironmentMode: 'attached' as const, codeWorkspaces: [mac] }; + await saveConvo({ userId }, { conversationId, ...decision }); + await methods.replaceConvoCodeEnvironmentDecision({ + user: userId, + conversationId, + expected: decision, + codeWorkspaces: [vm], + }); + + await saveConvo({ userId }, { conversationId, ...decision }); + expect((await getConvo(userId, conversationId))?.codeWorkspaces).toEqual([vm]); + await saveConvo( + { userId }, + { conversationId, codeEnvironmentMode: null, codeWorkspaces: [] }, + ); + const stored = await getConvo(userId, conversationId); + expect(stored?.codeEnvironmentMode).toBe('attached'); + expect(stored?.codeWorkspaces).toEqual([vm]); + }); + + it('seeds imported decisions without letting repeated imports undo a move', async () => { + const conversationId = uuidv4(); + const decision = { codeEnvironmentMode: 'attached' as const, codeWorkspaces: [mac] }; + const imported = { conversationId, user: userId, ...decision }; + await methods.bulkSaveConvos([imported]); + expect((await getConvo(userId, conversationId))?.codeWorkspaces).toEqual([mac]); + await methods.replaceConvoCodeEnvironmentDecision({ + user: userId, + conversationId, + expected: decision, + codeWorkspaces: [vm], + }); + await methods.bulkSaveConvos([imported]); + const stored = await getConvo(userId, conversationId); + expect(stored?.codeEnvironmentMode).toBe('attached'); + expect(stored?.codeWorkspaces).toEqual([vm]); + }); + + it.each([{}, { codeWorkspaces: [mac] }])( + 'preserves a legacy decision: %j', + async (decision) => { + const conversationId = uuidv4(); + await Conversation.collection.insertOne({ conversationId, user: userId, ...decision }); + await saveConvo( + { userId }, + { conversationId, codeEnvironmentMode: 'attached', codeWorkspaces: [vm] }, + { unsetFields }, + ); + const stored = await getConvo(userId, conversationId); + expect(stored?.codeEnvironmentMode).toBeUndefined(); + expect(stored?.codeWorkspaces).toEqual(decision.codeWorkspaces); + }, + ); + }); + describe('replaceConvoCodeEnvironmentDecision', () => { const anchor = new Date('2026-09-12T11:22:22.976Z'); const mac = { environmentId: 'code-mac', workspaceId: 'primary' }; diff --git a/packages/data-schemas/src/methods/conversation.ts b/packages/data-schemas/src/methods/conversation.ts index 10ad260b8d4..36b7f365f2e 100644 --- a/packages/data-schemas/src/methods/conversation.ts +++ b/packages/data-schemas/src/methods/conversation.ts @@ -2181,6 +2181,15 @@ export function createConversationMethods( delete update.isTemporary; delete update.expiredAt; delete update.initial_agent_id; + /** Ordinary saves may seed a decision, but only an explicit move may replace it. */ + const decisionOnInsert = { + ...(convo.codeEnvironmentMode != null && { + codeEnvironmentMode: convo.codeEnvironmentMode, + }), + ...(convo.codeWorkspaces != null && { codeWorkspaces: convo.codeWorkspaces }), + }; + delete update.codeEnvironmentMode; + delete update.codeWorkspaces; stripActorCheckpointFields(update); if (appendMessageIds == null) { update.messages = await getMessages({ conversationId, user: userId }, '_id'); @@ -2189,6 +2198,8 @@ export function createConversationMethods( } const unsetFields: Record = { ...(metadata?.unsetFields ?? {}) }; delete unsetFields.initial_agent_id; + delete unsetFields.codeEnvironmentMode; + delete unsetFields.codeWorkspaces; stripActorCheckpointFields(unsetFields); if (Object.prototype.hasOwnProperty.call(update, 'chatProjectId') && update.chatProjectId) { @@ -2318,6 +2329,7 @@ export function createConversationMethods( : createdAtOnInsert; operation.$setOnInsert = { initial_agent_id: initialAgentId, + ...decisionOnInsert, ...retentionOnInsert, ...(createdAtForInsert ? { createdAt: createdAtForInsert } : {}), }; @@ -2642,7 +2654,7 @@ export function createConversationMethods( const affectedProjectStats = new Map(); const bulkOps = conversations.map((convo) => { - const sanitized = { ...convo }; + const { codeEnvironmentMode, codeWorkspaces, ...sanitized } = convo; delete sanitized.initial_agent_id; stripActorCheckpointFields(sanitized); if (typeof sanitized.user === 'string' && typeof sanitized.chatProjectId === 'string') { @@ -2676,7 +2688,11 @@ export function createConversationMethods( }, update: { $set: sanitized, - $setOnInsert: { initial_agent_id: null }, + $setOnInsert: { + initial_agent_id: null, + ...(codeEnvironmentMode != null && { codeEnvironmentMode }), + ...(codeWorkspaces != null && { codeWorkspaces }), + }, }, upsert: true, timestamps: false,