const { logger } = require('@librechat/data-schemas'); const { createContentAggregator, GraphNodeKeys } = require('@librechat/agents'); const { checkAccess, loadSkillStates, initializeAgent, isMemoryEnabled, primeInvokedSkillsForProfiles, validateAgentModel, extractManualSkills, GenerationJobManager, getCustomEndpointConfig, getProviderConfig, discoverConnectedAgents, resolveAgentTokenConfig, resolveAgentScopedSkillIds, resolveModelSpecSkillIds, getAgentStartupTelemetry, buildAgentContextAttachmentsByAgentId, collectCodeExecutionProfileRoutes, getLazySubagentConfigId, } = require('@librechat/api'); const { Permissions, ResourceType, EModelEndpoint, PermissionBits, PermissionTypes, MAX_SUBAGENT_DEPTH, isAgentsEndpoint, getResponseSender, AgentCapabilities, Tools, MAX_SUBAGENT_GRAPH_NODES, MAX_SUBAGENT_RUN_CONFIGS, isEphemeralAgentId, } = require('librechat-data-provider'); const { createToolEndCallback, createAttachmentEmitter, createBackgroundCodeResultHandler, getDefaultHandlers, } = require('~/server/controllers/agents/callbacks'); const { loadAgentTools, loadToolsForExecution, getAccessibleMcpServerNames, isFatalAgentInitializationError, } = require('~/server/services/ToolService'); const { filterFilesByAgentAccess } = require('~/server/services/Files/permissions'); const { getSkillToolDeps, getSkillDbMethods, canAuthorSkillFiles, withDeploymentSkillIds, buildAgentToolContext, enrichLoadedToolsWithAgentContext, } = require('./skillDeps'); const { getModelsConfig } = require('~/server/controllers/ModelController'); const { checkPermission, findAccessibleResources } = require('~/server/services/PermissionService'); const AgentClient = require('~/server/controllers/agents/client'); const { processAddedConvo } = require('./addedConvo'); const { logViolation } = require('~/cache'); const db = require('~/models'); /** * Creates a tool loader function for the agent. * @param {AbortSignal} signal - The abort signal * @param {string | null} [streamId] - The stream ID for resumable mode * @param {boolean} [definitionsOnly=false] - When true, returns only serializable * tool definitions without creating full tool instances (for event-driven mode) * @param {number} [jobCreatedAt] - The generation epoch that owns emitted tool events */ function createToolLoader(signal, streamId = null, definitionsOnly = false, jobCreatedAt) { /** * @param {object} params * @param {ServerRequest} params.req * @param {ServerResponse} params.res * @param {string} params.agentId * @param {string[]} params.tools * @param {string} params.provider * @param {string} params.model * @param {AgentToolResources} params.tool_resources * @returns {Promise<{ * tools?: StructuredTool[], * toolContextMap: Record, * toolDefinitions?: import('@librechat/agents').LCTool[], * userMCPAuthMap?: Record>, * toolRegistry?: import('@librechat/agents').LCToolRegistry * } | undefined>} */ return async function loadTools({ req, res, tools, model, agentId, provider, tool_options, tool_resources, codeExecutionContext, accessibleMcpServerNames, }) { const agent = { id: agentId, tools, provider, model, tool_options }; try { return await loadAgentTools({ req, res, agent, signal, streamId, jobCreatedAt, tool_resources, codeExecutionContext, definitionsOnly, accessibleMcpServerNames, }); } catch (error) { if (isFatalAgentInitializationError(error)) { throw error; } logger.error('Error loading tools for agent ' + agentId, error); } }; } /** * Initializes the AgentClient for a given request/response cycle. * @param {Object} params * @param {Express.Request} params.req * @param {Express.Response} params.res * @param {AbortSignal} params.signal * @param {Object} params.endpointOption * @param {number} [params.jobCreatedAt] * @param {string} [params.checkpointNamespace] Immutable saver-level generation scope */ const initializeClient = async ({ req, res, signal, endpointOption, jobCreatedAt, checkpointNamespace, }) => { if (!endpointOption) { throw new Error('Endpoint option not provided'); } const appConfig = req.config; const startupTelemetry = getAgentStartupTelemetry(req); /** @type {string | null} */ const streamId = req._resumableStreamId || null; /** @type {Array} */ const collectedUsage = []; /** * Vertex Gemini 3 thought signatures captured from `chat_model_end` events, * keyed by `tool_call_id`. Persisted on * `responseMessage.metadata.thoughtSignatures` so subsequent conversation * turns can restore each signature onto the right reconstructed AIMessage's * `additional_kwargs.signatures` and avoid 400s when resuming after a tool * round-trip without a final text reply. Always allocated; capture path * is a no-op for providers that don't emit signatures (OpenAI, Anthropic, * Bedrock, etc.). * @type {Record} */ const collectedThoughtSignatures = {}; /** @type {ArtifactPromises} */ const artifactPromises = []; /** @type {Map} */ const toolInputValidationErrors = new Map(); const { contentParts, aggregateContent, stepMap } = createContentAggregator(); const artifactToolEndCallback = createToolEndCallback({ req, res, artifactPromises, streamId, jobCreatedAt, }); /** Query accessible skill IDs once per run (shared across all agents). * Skills activate under strict opt-in semantics — see * `resolveAgentScopedSkillIds` for the per-agent activation predicate: * - Ephemeral agent → model-spec `skills` config first, otherwise the * per-conversation skills badge toggle (full catalog). * - Persisted agent → `agent.skills_enabled === true`. Optional * `agent.skills` allowlist narrows the catalog; empty/undefined * allowlist with the toggle on = full accessible catalog. */ const enabledCapabilities = new Set(appConfig?.endpoints?.[EModelEndpoint.agents]?.capabilities); const skillsCapabilityEnabled = enabledCapabilities.has(AgentCapabilities.skills); const codeEnvAvailable = enabledCapabilities.has(AgentCapabilities.execute_code); const backgroundToolsAvailable = enabledCapabilities.has(AgentCapabilities.run_in_background); const toolIntentsAvailable = enabledCapabilities.has(AgentCapabilities.tool_intents); const statefulSessionsAvailable = enabledCapabilities.has( AgentCapabilities.stateful_code_sessions, ); const ephemeralSkillsToggle = req.body?.ephemeralAgent?.skills === true; const skillDbMethods = getSkillDbMethods(); if (!endpointOption.agent) { throw new Error('No agent promise provided'); } /** Run-level gate for inline memory tools: the `memory` capability must be * enabled, memory must be configured, and the user must not have opted out. * Requires the memory WRITE permissions (CREATE + UPDATE) — both inline tools * mutate memory — so the tools aren't registered (and shown to the model) for * read-only-memory roles that the runtime loader would then refuse to build. * Agents (or the ephemeral memory badge) opt in per-agent via the `memory` * marker on `tools`. */ const memoryAvailablePromise = enabledCapabilities.has(AgentCapabilities.memory) && isMemoryEnabled(appConfig?.memory) && req.user?.personalization?.memories !== false && checkAccess({ user: req.user, permissionType: PermissionTypes.MEMORIES, permissions: [Permissions.USE, Permissions.CREATE, Permissions.UPDATE], getRoleByName: db.getRoleByName, }); const accessibleSkillIdsPromise = skillsCapabilityEnabled ? findAccessibleResources({ userId: req.user.id, role: req.user.role, resourceType: ResourceType.SKILL, requiredPermissions: PermissionBits.VIEW, }).then(withDeploymentSkillIds) : Promise.resolve([]); const editableSkillIdsPromise = skillsCapabilityEnabled ? findAccessibleResources({ userId: req.user.id, role: req.user.role, resourceType: ResourceType.SKILL, requiredPermissions: PermissionBits.EDIT, }) : Promise.resolve([]); const skillCreateAllowedPromise = skillsCapabilityEnabled ? getSkillToolDeps().canCreateSkill({ req }) : Promise.resolve(false); const skillStatesPromise = accessibleSkillIdsPromise.then((accessibleSkillIds) => loadSkillStates({ userId: req.user.id, appConfig, getUserById: db.getUserById, accessibleSkillIds, }), ); const primaryAgentPromise = endpointOption.agent; const modelsConfigPromise = getModelsConfig(req); const validatedPrimaryAgentPromise = Promise.all([primaryAgentPromise, modelsConfigPromise]).then( async ([primaryAgent, modelsConfig]) => { if (!primaryAgent) { throw new Error('Agent not found'); } const validationResult = await validateAgentModel({ req, res, modelsConfig, logViolation, agent: primaryAgent, }); if (!validationResult.isValid) { throw new Error(validationResult.error?.message); } return { primaryAgent, modelsConfig }; }, ); /** * Agent context store - populated after initialization, accessed by callback via closure. * Maps agentId -> { userMCPAuthMap, agent, tool_resources, toolRegistry, openAIApiKey } * @type {Map>, * agent?: object, * tool_resources?: object, * toolRegistry?: import('@librechat/agents').LCToolRegistry, * requestScopedConnections?: import('@librechat/api').RequestScopedMCPConnectionStore, * openAIApiKey?: string * }>} */ const agentToolContexts = new Map(); /** Attach only the host-resolved route for the actually executing agent. * Runnable metadata is transport data and may contain caller-controlled * keys, so discard any incoming route context before resolving from the * server-owned per-agent map. This covers both traditional TOOL_END events * and event-driven ON_TOOL_EXECUTE callbacks. */ const toolEndCallback = async (data, metadata = {}) => { const node = typeof metadata.langgraph_node === 'string' ? metadata.langgraph_node : ''; const nodeAgentId = node.startsWith(GraphNodeKeys.TOOLS) ? node.slice(GraphNodeKeys.TOOLS.length) : undefined; const executingAgentId = metadata.executingAgentId ?? metadata.agentId ?? metadata.agent_id ?? nodeAgentId; const soleContext = agentToolContexts.size === 1 ? agentToolContexts.values().next().value : null; const trustedContext = (typeof executingAgentId === 'string' ? agentToolContexts.get(executingAgentId) : null) ?? soleContext; const callbackMetadata = { ...metadata }; delete callbackMetadata.codeExecutionContext; if (trustedContext?.codeExecutionContext) { callbackMetadata.codeExecutionContext = trustedContext.codeExecutionContext; } return artifactToolEndCallback(data, callbackMetadata); }; /** @type {Map} */ const endpointTokenConfigByAgentId = new Map(); const toolExecuteOptions = { loadTools: async (toolNames, agentId) => { const ctx = agentToolContexts.get(agentId) ?? {}; logger.debug(`[ON_TOOL_EXECUTE] ctx found: ${!!ctx.userMCPAuthMap}, agent: ${ctx.agent?.id}`); logger.debug(`[ON_TOOL_EXECUTE] toolRegistry size: ${ctx.toolRegistry?.size ?? 'undefined'}`); const result = await loadToolsForExecution({ req, res, signal, streamId, conversationId, toolNames, agent: ctx.agent, toolRegistry: ctx.toolRegistry, backgroundToolNames: ctx.backgroundToolNames, intentToolNames: ctx.intentToolNames, mcpAvailableTools: ctx.mcpAvailableTools, requestScopedConnections: ctx.requestScopedConnections, userMCPAuthMap: ctx.userMCPAuthMap, tool_resources: ctx.tool_resources, actionsEnabled: ctx.actionsEnabled, accessibleMcpServerNames: ctx.accessibleMcpServerNames, jobCreatedAt, }); logger.debug(`[ON_TOOL_EXECUTE] loaded ${result.loadedTools?.length ?? 0} tools`); /** Per-agent narrowed flag (admin capability AND agent.tools * includes execute_code), captured in `agentToolContexts` when * the agent initialized. Falls back to `false` on any stray * ctx miss so a skills-only agent never gains sandbox access * even if capability lookup somehow skips. */ return enrichLoadedToolsWithAgentContext({ result, req, ctx, }); }, toolEndCallback, persistBackgroundCodeResult: createBackgroundCodeResultHandler({ req, updateToolCallResult: db.updateToolCallResult, }), emitAttachment: createAttachmentEmitter({ res, streamId, jobCreatedAt }), ...getSkillToolDeps(), }; const summarizationOptions = appConfig?.summarization?.enabled === false ? { enabled: false } : { enabled: true }; /** * Per-request map of per-subagent `createContentAggregator` instances * keyed by the parent's `tool_call_id`. The handler in `callbacks.js` * lazily creates an aggregator for each distinct `parentToolCallId` * and folds every `ON_SUBAGENT_UPDATE` event into it as they stream * in. `AgentClient` pulls each aggregator's `contentParts` at message * save time and attaches them to the matching `subagent` tool_call so * the child's reasoning / tool calls / final text survive a page * refresh — the client-side Recoil atom is best-effort live-only. */ const subagentAggregatorsByToolCallId = new Map(); /** Backend prices each model call authoritatively (premium tiers, cache * rates) and emits the cost on on_token_usage when contextCost is on, so * the gauge sums real costs instead of re-deriving from base rates. * `endpointTokenConfig` is filled in once `primaryConfig` resolves below so * custom-endpoint agents price with their configured rates, not defaults. */ const usageCost = { enabled: appConfig?.interfaceConfig?.contextCost === true, pricing: { getMultiplier: db.getMultiplier, getCacheMultiplier: db.getCacheMultiplier }, }; /** Latest visible context snapshot + every emitted usage payload for this * response, captured by the handlers and persisted on the response message's * metadata so the breakdown and branch/total cost survive a reload. * @type {{ latest: import('librechat-data-provider').TContextUsageEvent | null, count: number }} */ const contextUsageSink = { latest: null, count: 0 }; /** @type {Array} */ const usageEmitSink = []; const eventHandlers = getDefaultHandlers({ res, contentParts, stepMap, toolInputValidationErrors, toolExecuteOptions, summarizationOptions, aggregateContent, toolEndCallback, collectedUsage, collectedThoughtSignatures, streamId, jobCreatedAt, subagentAggregatorsByToolCallId, usageCost, contextUsageSink, usageEmitSink, }); const [ memoryAvailable, accessibleSkillIds, editableSkillIds, skillCreateAllowed, { skillStates, defaultActiveOnShare }, { primaryAgent, modelsConfig }, ] = await Promise.all([ memoryAvailablePromise, accessibleSkillIdsPromise, editableSkillIdsPromise, skillCreateAllowedPromise, skillStatesPromise, validatedPrimaryAgentPromise, ]); delete endpointOption.agent; const agentConfigs = new Map(); const allowedProviders = new Set(appConfig?.endpoints?.[EModelEndpoint.agents]?.allowedProviders); /** Event-driven mode: only load tool definitions, not full instances */ const loadTools = createToolLoader(signal, streamId, true, jobCreatedAt); /** @type {Array} */ const requestFiles = req.body.files ?? []; /** @type {string} */ const conversationId = req.body.conversationId; /** @type {string | undefined} */ const parentMessageId = req.body.parentMessageId; /** * Skill names the user invoked via the `$` popover for this turn. Only flows * to the primary agent — handoff agents are follow-up turns that don't see * the user's per-submission `$` selections. `extractManualSkills` also * drops non-string / empty elements so a crafted payload can't reach the * `getSkillByName` DB query with nonsense values. * @type {string[] | undefined} */ const manualSkills = extractManualSkills(req.body); const selectedModelSpec = endpointOption.spec && Array.isArray(appConfig?.modelSpecs?.list) ? appConfig.modelSpecs.list.find((modelSpec) => modelSpec.name === endpointOption.spec) : null; if ( primaryAgent && isEphemeralAgentId(primaryAgent.id) && selectedModelSpec && Object.hasOwn(selectedModelSpec, 'skills') ) { if (selectedModelSpec.skills === true) { primaryAgent.skills_enabled = true; delete primaryAgent.skills; } else if (selectedModelSpec.skills === false) { primaryAgent.skills_enabled = false; primaryAgent.skills = []; } else if (Array.isArray(selectedModelSpec.skills)) { const resolvedSkillIds = await resolveModelSpecSkillIds({ names: selectedModelSpec.skills, accessibleSkillIds, getSkillByName: skillDbMethods.getSkillByName, }); primaryAgent.skills_enabled = true; primaryAgent.skills = resolvedSkillIds.map((id) => id.toString()); } } const primaryScopedSkillIds = resolveAgentScopedSkillIds({ agent: primaryAgent, accessibleSkillIds, skillsCapabilityEnabled, ephemeralSkillsToggle, }); const primaryScopedEditableSkillIds = resolveAgentScopedSkillIds({ agent: primaryAgent, accessibleSkillIds: editableSkillIds, skillsCapabilityEnabled, ephemeralSkillsToggle, }); const primarySkillAuthoringAvailable = canAuthorSkillFiles({ agent: primaryAgent, scopedEditableSkillIds: primaryScopedEditableSkillIds, skillCreateAllowed, skillsCapabilityEnabled, ephemeralSkillsToggle, }); const primaryConfig = await initializeAgent( { req, res, loadTools, requestFiles, conversationId, parentMessageId, agent: primaryAgent, endpointOption, allowedProviders, isInitialAgent: true, accessibleSkillIds: primaryScopedSkillIds, skillAuthoringAvailable: primarySkillAuthoringAvailable, codeEnvAvailable, backgroundToolsAvailable, toolIntentsAvailable, statefulSessionsAvailable, memoryAvailable, skillStates, defaultActiveOnShare, manualSkills, }, { getFiles: db.getFiles, getUserKey: db.getUserKey, getMessages: db.getMessages, getConvoFiles: db.getConvoFiles, getAccessibleMcpServerNames, updateFilesUsage: db.updateFilesUsage, getUserKeyValues: db.getUserKeyValues, getUserCodeFiles: db.getUserCodeFiles, getToolFilesByIds: db.getToolFilesByIds, getCodeGeneratedFiles: db.getCodeGeneratedFiles, filterFilesByAgentAccess, listSkillsByAccess: skillDbMethods.listSkillsByAccess, listAlwaysApplySkills: skillDbMethods.listAlwaysApplySkills, getSkillByName: skillDbMethods.getSkillByName, }, ); /** Price emitted usage with the primary agent's resolved endpoint config so * custom-endpoint agents reflect configured rates (mirrors the AgentClient * spending path, which reads the same config). */ usageCost.endpointTokenConfig = primaryConfig.endpointTokenConfig; logger.debug( `[initializeClient] Storing tool context for ${primaryConfig.id}: ${primaryConfig.toolDefinitions?.length ?? 0} tools, registry size: ${primaryConfig.toolRegistry?.size ?? '0'}`, ); agentToolContexts.set( primaryConfig.id, buildAgentToolContext({ agent: primaryAgent, config: primaryConfig }), ); const { agentConfigs: discoveredConfigs, edges: discoveredEdges, userMCPAuthMap: discoveredMCPAuthMap, skippedAgentIds: discoveredSkippedIds, } = await discoverConnectedAgents( { req, res, primaryConfig, agent_ids: primaryConfig.agent_ids, endpointOption, allowedProviders, modelsConfig, loadTools, requestFiles, conversationId, parentMessageId, computeAccessibleSkillIds: (agent) => resolveAgentScopedSkillIds({ agent, accessibleSkillIds, skillsCapabilityEnabled, ephemeralSkillsToggle, }), computeSkillAuthoringAvailable: (agent) => canAuthorSkillFiles({ agent, scopedEditableSkillIds: resolveAgentScopedSkillIds({ agent, accessibleSkillIds: editableSkillIds, skillsCapabilityEnabled, ephemeralSkillsToggle, }), skillCreateAllowed, skillsCapabilityEnabled, ephemeralSkillsToggle, }), skillStates, defaultActiveOnShare, codeEnvAvailable, backgroundToolsAvailable, toolIntentsAvailable, statefulSessionsAvailable, memoryAvailable, }, { getAgent: db.getAgent, checkPermission, logViolation, db: { getFiles: db.getFiles, getUserKey: db.getUserKey, getMessages: db.getMessages, getConvoFiles: db.getConvoFiles, getAccessibleMcpServerNames, updateFilesUsage: db.updateFilesUsage, getUserKeyValues: db.getUserKeyValues, getUserCodeFiles: db.getUserCodeFiles, getToolFilesByIds: db.getToolFilesByIds, getCodeGeneratedFiles: db.getCodeGeneratedFiles, filterFilesByAgentAccess, listSkillsByAccess: skillDbMethods.listSkillsByAccess, listAlwaysApplySkills: skillDbMethods.listAlwaysApplySkills, getSkillByName: skillDbMethods.getSkillByName, }, // The callback fires during BFS, before the helper prunes agents // whose edges end up filtered. Don't populate `agentConfigs` here — // `discoveredConfigs` (returned below) is the authoritative pruned // set. The per-agent tool context map is OK to keep populated even // for pruned ids: it's only read by closure in ON_TOOL_EXECUTE, // stale entries are unreachable at runtime. onAgentInitialized: (agentId, agent, config) => { agentToolContexts.set(agentId, buildAgentToolContext({ agent, config })); }, // Pass through the `@librechat/api` exports so that tests which // `jest.mock('@librechat/api')` can override the initializer/validator. initializeAgent, validateAgentModel, }, ); // Copy the pruned discovery result into the outer map. Anything the // helper dropped (skipped or unreachable after edge filtering) is // intentionally absent. `processAddedConvo` below may still add more // entries for parallel multi-convo execution. for (const [agentId, config] of discoveredConfigs) { agentConfigs.set(agentId, config); } let userMCPAuthMap = discoveredMCPAuthMap; let edges = discoveredEdges; /** Multi-Convo: Process addedConvo for parallel agent execution */ const { userMCPAuthMap: updatedMCPAuthMap } = await processAddedConvo({ req, res, loadTools, logViolation, modelsConfig, requestFiles, agentConfigs, primaryAgent, endpointOption, userMCPAuthMap, conversationId, parentMessageId, allowedProviders, primaryAgentId: primaryConfig.id, accessibleSkillIds, editableSkillIds, skillsCapabilityEnabled, ephemeralSkillsToggle, skillCreateAllowed, skillStates, defaultActiveOnShare, codeEnvAvailable, backgroundToolsAvailable, toolIntentsAvailable, statefulSessionsAvailable, memoryAvailable, }); if (updatedMCPAuthMap) { userMCPAuthMap = updatedMCPAuthMap; } for (const [agentId, config] of agentConfigs) { if (agentToolContexts.has(agentId)) { continue; } agentToolContexts.set(agentId, buildAgentToolContext({ agent: config, config })); } // `discoverConnectedAgents` always returns a concrete array, so no // further normalization is needed before handing this to `createRun`. primaryConfig.edges = edges; // Subagents run in isolated context windows and are invoked via a dedicated // spawn tool, not handoff edges. Explicit children are advertised as inert, // VIEW-checked descriptors; model, tool, MCP, file, and skill initialization // happens only when the SDK selects one. const subagentsCapabilityEnabled = enabledCapabilities.has(AgentCapabilities.subagents); /** Track skipped ids locally so repeated failures short-circuit within * the subagent loading loop. Seeded from the discovery helper's skip * list so agents that already failed handoff loading don't get retried. */ const skippedAgentIds = new Set(discoveredSkippedIds ?? []); const lazyMetadataByAgentId = new Map(); const subagentGraphIds = new Set(); const expandedSubagentDescriptorState = { configCount: 0, rootAgentIds: [] }; const assertSubagentGraphRoom = (agentId) => { if (subagentGraphIds.has(agentId)) { return; } if (subagentGraphIds.size >= MAX_SUBAGENT_GRAPH_NODES) { logger.warn('[initializeClient] Subagent graph node limit exceeded', { agentId, primaryAgentId: primaryConfig.id, loadedSubagentCount: subagentGraphIds.size, maxSubagentGraphNodes: MAX_SUBAGENT_GRAPH_NODES, }); throw new Error( `Subagent graph exceeds the maximum of ${MAX_SUBAGENT_GRAPH_NODES} unique agents.`, ); } }; const countExpandedSubagentDescriptor = (agentId) => { expandedSubagentDescriptorState.configCount += 1; if (expandedSubagentDescriptorState.configCount <= MAX_SUBAGENT_RUN_CONFIGS) { return; } logger.warn('[initializeClient] Subagent run configuration limit exceeded', { agentId, expandedConfigCount: expandedSubagentDescriptorState.configCount, maxSubagentRunConfigs: MAX_SUBAGENT_RUN_CONFIGS, rootAgentIds: expandedSubagentDescriptorState.rootAgentIds, }); throw new Error( `Subagent run configuration exceeds the maximum of ${MAX_SUBAGENT_RUN_CONFIGS} expanded entries.`, ); }; const userId = req.user?.id; const userRole = req.user?.role; const throwIfAborted = (abortSignal) => { if (!abortSignal?.aborted) return; throw abortSignal.reason instanceof Error ? abortSignal.reason : new Error('Subagent resolution was aborted.'); }; const waitForAbort = (promise, abortSignal) => { throwIfAborted(abortSignal); if (!abortSignal) return promise; return new Promise((resolve, reject) => { const onAbort = () => { reject( abortSignal.reason instanceof Error ? abortSignal.reason : new Error('Subagent resolution was aborted.'), ); }; abortSignal.addEventListener('abort', onAbort, { once: true }); if (abortSignal.aborted) { onAbort(); } promise.then(resolve, reject).finally(() => { abortSignal.removeEventListener('abort', onAbort); }); }); }; const hasSubagentViewAccess = async (agent, agentId, abortSignal) => { throwIfAborted(abortSignal); if (!userId) return false; const hasAccess = await waitForAbort( checkPermission({ userId, role: userRole, resourceType: ResourceType.AGENT, resourceId: agent._id, requiredPermission: PermissionBits.VIEW, }), abortSignal, ); throwIfAborted(abortSignal); if (!hasAccess) { logger.warn( `[processAgent] User ${userId} lacks VIEW access to subagent ${agentId}, skipping`, ); } return hasAccess; }; const getIncludeReasoningHistory = (agent) => { if (!agent.provider) return undefined; try { return getProviderConfig({ provider: agent.provider, appConfig }).customEndpointConfig ?.customParams?.includeReasoningHistory; } catch { return undefined; } }; const getExplicitSubagentIds = (agent) => Array.from( new Set( Array.isArray(agent.subagents?.agent_ids) ? agent.subagents.agent_ids.filter( (id) => typeof id === 'string' && id && id !== agent.id, ) : [], ), ); const toLazySubagentMetadata = (agent) => ({ id: agent.id, name: agent.name, description: agent.description, provider: agent.provider, model: agent.model, model_parameters: { model: agent.model_parameters?.model }, recursion_limit: agent.recursion_limit, subagents: agent.subagents, configId: getLazySubagentConfigId(agent), codeEnvAvailable: codeEnvAvailable === true && agent.tools?.includes(Tools.execute_code) === true, statefulCodeSessions: statefulSessionsAvailable === true && codeEnvAvailable === true && agent.stateful_code_sessions === true && agent.tools?.includes(Tools.execute_code) === true, statefulCodeEnvironment: agent.stateful_code_environment, includeReasoningHistory: getIncludeReasoningHistory(agent), }); const loadSubagentMetadata = async (agentId) => { if (skippedAgentIds.has(agentId)) return null; const cached = lazyMetadataByAgentId.get(agentId); if (cached) return cached; try { const agent = await db.getAgentWithVersionCount({ id: agentId }); if (!agent || !(await hasSubagentViewAccess(agent, agentId))) { skippedAgentIds.add(agentId); return null; } const metadata = toLazySubagentMetadata(agent); lazyMetadataByAgentId.set(agentId, metadata); return metadata; } catch (error) { if (isFatalAgentInitializationError(error)) { throw error; } logger.error(`[initializeClient] Error loading subagent metadata ${agentId}:`, error); skippedAgentIds.add(agentId); return null; } }; /** * Resolves the selected descriptor inside the foreground request. The * legacy initializer requires request/response objects for tool and MCP * setup, so this intentionally remains request-scoped until AI-1597 gives * child execution a durable runtime context. */ const initializeLazySubagent = async ({ agentId, configId, context, lazyChildren }) => { throwIfAborted(context.signal); const agent = await waitForAbort(db.getAgentWithVersionCount({ id: agentId }), context.signal); throwIfAborted(context.signal); if (!agent || getLazySubagentConfigId(agent) !== configId) { throw new Error(`Subagent ${agentId} changed before it could be initialized.`); } if (!(await hasSubagentViewAccess(agent, agentId, context.signal))) { throw new Error(`You no longer have access to subagent ${agentId}.`); } const validation = await waitForAbort( validateAgentModel({ req, res, agent, modelsConfig, logViolation }), context.signal, ); throwIfAborted(context.signal); if (!validation.isValid) { throw new Error(validation.error?.message ?? `Subagent ${agentId} failed model validation.`); } const scopedSkillIds = resolveAgentScopedSkillIds({ agent, accessibleSkillIds, skillsCapabilityEnabled, ephemeralSkillsToggle, }); const scopedEditableSkillIds = resolveAgentScopedSkillIds({ agent, accessibleSkillIds: editableSkillIds, skillsCapabilityEnabled, ephemeralSkillsToggle, }); const config = await waitForAbort( initializeAgent( { req, res, agent, loadTools: createToolLoader(context.signal, streamId, true, jobCreatedAt), requestFiles, conversationId, parentMessageId, endpointOption: { ...endpointOption, endpoint: EModelEndpoint.agents }, allowedProviders, accessibleSkillIds: scopedSkillIds, skillAuthoringAvailable: canAuthorSkillFiles({ agent, scopedEditableSkillIds, skillCreateAllowed, skillsCapabilityEnabled, ephemeralSkillsToggle, }), codeEnvAvailable, statefulSessionsAvailable, memoryAvailable, skillStates, defaultActiveOnShare, }, { getFiles: db.getFiles, getUserKey: db.getUserKey, getMessages: db.getMessages, getConvoFiles: db.getConvoFiles, getAccessibleMcpServerNames, updateFilesUsage: db.updateFilesUsage, getUserKeyValues: db.getUserKeyValues, getUserCodeFiles: db.getUserCodeFiles, getToolFilesByIds: db.getToolFilesByIds, getCodeGeneratedFiles: db.getCodeGeneratedFiles, filterFilesByAgentAccess, listSkillsByAccess: skillDbMethods.listSkillsByAccess, listAlwaysApplySkills: skillDbMethods.listAlwaysApplySkills, getSkillByName: skillDbMethods.getSkillByName, }, ), context.signal, ); throwIfAborted(context.signal); config.lazySubagentConfigs = lazyChildren; agentToolContexts.set(agentId, buildAgentToolContext({ agent, config })); endpointTokenConfigByAgentId.set(agentId, config.endpointTokenConfig); return config; }; const buildLazySubagentDescriptors = async (agent, depth = 0, ancestors = new Set()) => { if (!subagentsCapabilityEnabled || !agent.subagents?.enabled) { return []; } if (agent.subagents.allowSelf !== false) { countExpandedSubagentDescriptor(agent.id); } const subagentIds = getExplicitSubagentIds(agent); if (subagentIds.length > 0 && depth >= MAX_SUBAGENT_DEPTH) { logger.warn('[initializeClient] Subagent graph depth limit exceeded', { agentId: agent.id, primaryAgentId: primaryConfig.id, depth, maxSubagentDepth: MAX_SUBAGENT_DEPTH, }); throw new Error( `Subagent graph exceeds the maximum depth of ${MAX_SUBAGENT_DEPTH} at agent ${agent.id}.`, ); } const nextAncestors = new Set(ancestors); nextAncestors.add(agent.id); const descriptors = []; for (const subagentId of subagentIds) { if (skippedAgentIds.has(subagentId) || nextAncestors.has(subagentId)) continue; if (subagentId !== primaryConfig.id) { assertSubagentGraphRoom(subagentId); } const existing = subagentId === primaryConfig.id ? primaryConfig : agentConfigs.get(subagentId); if (existing) { countExpandedSubagentDescriptor(subagentId); if (subagentId !== primaryConfig.id) { subagentGraphIds.add(subagentId); } const existingChildren = await buildLazySubagentDescriptors( existing, depth + 1, nextAncestors, ); existing.lazySubagentConfigs = existingChildren.filter((child) => child.configId); existing.subagentAgentConfigs = existingChildren.filter((child) => !child.configId); descriptors.push(existing); continue; } const metadata = await loadSubagentMetadata(subagentId); if (!metadata) continue; countExpandedSubagentDescriptor(subagentId); subagentGraphIds.add(subagentId); const childDescriptors = await buildLazySubagentDescriptors( metadata, depth + 1, nextAncestors, ); const lazyChildren = childDescriptors.filter((child) => child.configId); const eagerChildren = childDescriptors.filter((child) => !child.configId); descriptors.push({ id: metadata.id, name: metadata.name, description: metadata.description, provider: metadata.provider, model: metadata.model, model_parameters: metadata.model_parameters, recursion_limit: metadata.recursion_limit, subagents: metadata.subagents, configId: metadata.configId, codeEnvAvailable: metadata.codeEnvAvailable, statefulCodeSessions: metadata.statefulCodeSessions, statefulCodeEnvironment: metadata.statefulCodeEnvironment, includeReasoningHistory: metadata.includeReasoningHistory, lazySubagentConfigs: lazyChildren, subagentAgentConfigs: eagerChildren, resolve: (context) => initializeLazySubagent({ agentId: metadata.id, configId: metadata.configId, context, lazyChildren, }).then((config) => { config.subagentAgentConfigs = eagerChildren; return config; }), }); } return descriptors; }; const resolveSubagentTrees = async (rootConfigs) => { expandedSubagentDescriptorState.rootAgentIds = rootConfigs .filter((config) => config?.id) .map((config) => config.id); for (const config of rootConfigs) { if (!config?.id) continue; const descriptors = await buildLazySubagentDescriptors(config); config.lazySubagentConfigs = descriptors.filter((child) => child.configId); config.subagentAgentConfigs = descriptors.filter((child) => !child.configId); } }; await resolveSubagentTrees([primaryConfig, ...agentConfigs.values()]); primaryConfig.subagents = subagentsCapabilityEnabled ? primaryConfig.subagents : undefined; /** If the capability is off at the endpoint level, strip `subagents` on * every loaded config — not just the primary. `run.ts` calls * `buildSubagentConfigs` for every agent in the array, so a handoff * agent with `subagents.enabled: true` persisted on its document would * otherwise still expose self-spawn at runtime even though the admin * has disabled the capability globally. */ if (!subagentsCapabilityEnabled) { primaryConfig.lazySubagentConfigs = undefined; for (const config of agentConfigs.values()) { config.subagents = undefined; config.subagentAgentConfigs = undefined; config.lazySubagentConfigs = undefined; } } const agentContextAttachmentsByAgentId = buildAgentContextAttachmentsByAgentId([ primaryConfig, ...agentConfigs.values(), ]); let endpointConfig = appConfig.endpoints?.[primaryConfig.endpoint]; if (!isAgentsEndpoint(primaryConfig.endpoint) && !endpointConfig) { try { endpointConfig = getCustomEndpointConfig({ endpoint: primaryConfig.endpoint, appConfig, }); } catch (err) { logger.error( '[api/server/controllers/agents/client.js #titleConvo] Error getting custom endpoint config', err, ); } } const sender = primaryAgent.name ?? getResponseSender({ ...endpointOption, model: endpointOption.model_parameters.model, modelDisplayLabel: endpointConfig?.modelDisplayLabel, modelLabel: endpointOption.model_parameters.modelLabel, }); /** History priming uses the user's full ACL-accessible skill set (not * per-agent scoped) because prior turns may reference skills no longer * in any active agent's scope; the ACL check is the security gate. Each * selected Code API deployment receives its own upload, and only session * partitions routed to that deployment receive those storage pointers. */ const codeExecutionProfiles = collectCodeExecutionProfileRoutes( [primaryConfig, ...agentConfigs.values()], { userId: req.user.id, conversationId, }, ); const handlePrimeInvokedSkills = skillsCapabilityEnabled ? (payload) => primeInvokedSkillsForProfiles({ req, payload, accessibleSkillIds, executionProfiles: codeExecutionProfiles, ...getSkillToolDeps(), }) : undefined; /** Per-agent resolved endpoint token config, keyed by agent id. Built from * `agentToolContexts` (the one map holding every agent, including pure * subagents pruned from `agentConfigs`) so usage billed/emitted for a * connected or subagent on a different custom endpoint is priced with THAT * agent's configured rates instead of the primary's. Every known agent is * recorded — even with an `undefined` config — so the resolver can tell a * known non-custom agent (built-in pricing) from an untagged/unknown one * (primary fallback). * @type {Map} */ for (const [agentId, ctx] of agentToolContexts) { endpointTokenConfigByAgentId.set(agentId, ctx?.endpointTokenConfig); } /** Price emitted usage per producing agent too, so the streamed/persisted * `metadata.usage.cost` matches the per-agent balance transaction. */ usageCost.resolveEndpointTokenConfig = (usage) => resolveAgentTokenConfig({ agentId: usage?.agentId, byAgentId: endpointTokenConfigByAgentId, fallback: usageCost.endpointTokenConfig, }); const client = new AgentClient({ req, res, sender, contentParts, stepMap, agentConfigs, eventHandlers, collectedUsage, collectedThoughtSignatures, aggregateContent, artifactPromises, primeInvokedSkills: handlePrimeInvokedSkills, agent: primaryConfig, spec: endpointOption.spec, iconURL: endpointOption.iconURL, chatProjectId: endpointOption.chatProjectId, attachments: primaryConfig.requestAttachments ?? primaryConfig.attachments, agentContextAttachmentsByAgentId, endpointType: endpointOption.endpointType, resendFiles: primaryConfig.resendFiles ?? true, maxContextTokens: primaryConfig.maxContextTokens, endpoint: isEphemeralAgentId(primaryConfig.id) ? primaryConfig.endpoint : EModelEndpoint.agents, subagentAggregatorsByToolCallId, /** Resolved endpoint token/pricing config so spending and cost reflect * configured rates for custom-endpoint agents instead of defaults. */ endpointTokenConfig: primaryConfig.endpointTokenConfig, /** Per-agent override of the above for multi-endpoint graphs (connected * agents + subagents); falls back to the primary config when an agent * isn't present or has no configured rates. */ endpointTokenConfigByAgentId, /** Capture sinks the handlers fill during the run; `sendCompletion` reads * them to persist the breakdown + usage rollup on the response message. */ contextUsageSink, usageEmitSink, startupTelemetry, toolInputValidationErrors, jobCreatedAt, checkpointNamespace, }); if (streamId) { GenerationJobManager.setCollectedUsage(streamId, collectedUsage, jobCreatedAt); } return { client, userMCPAuthMap }; }; module.exports = { initializeClient };