feat: simplify AI chat architecture and add record links (#16463)
## Summary This PR significantly simplifies the AI chat architecture by removing complex routing/planning mechanisms and introduces clickable record links in AI responses. ## Changes ### AI Chat Architecture Simplification - **Removed** the entire `ai-chat-router` module (~850 lines) including: - Strategy decider service - Plan generator service - Complex routing logic - **Removed** agent execution planning services (~700 lines): - `agent-execution.service.ts` - `agent-plan-executor.service.ts` - `agent-tool-generator.service.ts` - **Added** centralized `ToolRegistryService` for tool management: - Builds searchable tool index (database, action, workflow tools) - Provides tool lookup by name - Supports agent search for loading expertise - **Added** `ChatExecutionService` as simple replacement: - Includes full tool catalog in system prompt - Pre-loads common tools (find/create/update for company, person, opportunity, task, note) - Uses `load_tools` mechanism for dynamic tool activation - Enables native web search by default ### Record References in AI Responses - Added `recordReferences` field to tool outputs for create, find, and update operations - Implemented `[[record:objectName:recordId:displayName]]` syntax for AI to reference records - Created `RecordLink` component that renders clickable chips with object icons - Integrated record link parsing into the markdown renderer - Users can now click directly on created/found records in AI responses ### Workflow Agent Fixes - Fixed cache invalidation issue when creating agents in workflows - Added default prompt for workflow-created agents to prevent validation errors - Relaxed agent validation to only check properties being updated (not all required properties) ### Code Quality Improvements - Extracted `getRecordDisplayName` utility that mirrors frontend's `getLabelIdentifierFieldValue` logic - Uses object metadata to determine the correct label identifier field - Handles `FULL_NAME` composite type for person/workspaceMember objects - Shared across create, find, and update record services ## Net Impact - **~1,200 lines deleted** (complex routing/planning code) - **~500 lines added** (simpler tool registry + record links) - Significantly reduced code complexity - Better tool discovery through full catalog in system prompt - Improved UX with clickable record references ## Testing - Typecheck passes - Lint passes - Manual testing of AI chat with record creation and linking
This commit is contained in:
+1
-13
@@ -18,9 +18,6 @@ import { AgentMessageEntity } from './entities/agent-message.entity';
|
||||
import { AgentTurnEntity } from './entities/agent-turn.entity';
|
||||
import { AgentActorContextService } from './services/agent-actor-context.service';
|
||||
import { AgentAsyncExecutorService } from './services/agent-async-executor.service';
|
||||
import { AgentExecutionService } from './services/agent-execution.service';
|
||||
import { AgentPlanExecutorService } from './services/agent-plan-executor.service';
|
||||
import { AgentToolGeneratorService } from './services/agent-tool-generator.service';
|
||||
|
||||
@Module({
|
||||
imports: [
|
||||
@@ -41,18 +38,9 @@ import { AgentToolGeneratorService } from './services/agent-tool-generator.servi
|
||||
RoleTargetEntity,
|
||||
]),
|
||||
],
|
||||
providers: [
|
||||
AgentAsyncExecutorService,
|
||||
AgentExecutionService,
|
||||
AgentToolGeneratorService,
|
||||
AgentActorContextService,
|
||||
AgentPlanExecutorService,
|
||||
],
|
||||
providers: [AgentAsyncExecutorService, AgentActorContextService],
|
||||
exports: [
|
||||
AgentAsyncExecutorService,
|
||||
AgentExecutionService,
|
||||
AgentPlanExecutorService,
|
||||
AgentToolGeneratorService,
|
||||
AgentActorContextService,
|
||||
TypeOrmModule.forFeature([
|
||||
AgentMessageEntity,
|
||||
|
||||
+3
-3
@@ -19,7 +19,7 @@ import {
|
||||
AgentExceptionCode,
|
||||
} from 'src/engine/metadata-modules/ai/ai-agent/agent.exception';
|
||||
import { AGENT_CONFIG } from 'src/engine/metadata-modules/ai/ai-agent/constants/agent-config.const';
|
||||
import { AGENT_SYSTEM_PROMPTS } from 'src/engine/metadata-modules/ai/ai-agent/constants/agent-system-prompts.const';
|
||||
import { WORKFLOW_SYSTEM_PROMPTS } from 'src/engine/metadata-modules/ai/ai-agent/constants/agent-system-prompts.const';
|
||||
import { AgentEntity } from 'src/engine/metadata-modules/ai/ai-agent/entities/agent.entity';
|
||||
import { repairToolCall } from 'src/engine/metadata-modules/ai/ai-agent/utils/repair-tool-call.util';
|
||||
import { AI_TELEMETRY_CONFIG } from 'src/engine/metadata-modules/ai/ai-models/constants/ai-telemetry.const';
|
||||
@@ -141,7 +141,7 @@ export class AgentAsyncExecutorService {
|
||||
this.logger.log(`Generated ${Object.keys(tools).length} tools for agent`);
|
||||
|
||||
const textResponse = await generateText({
|
||||
system: `${AGENT_SYSTEM_PROMPTS.BASE}\n${AGENT_SYSTEM_PROMPTS.WORKFLOW_ADDITIONS}\n\n${agent ? agent.prompt : ''}`,
|
||||
system: `${WORKFLOW_SYSTEM_PROMPTS.BASE}\n\n${agent ? agent.prompt : ''}`,
|
||||
tools,
|
||||
model: registeredModel.model,
|
||||
prompt: userPrompt,
|
||||
@@ -177,7 +177,7 @@ export class AgentAsyncExecutorService {
|
||||
}
|
||||
|
||||
const output = await generateObject({
|
||||
system: AGENT_SYSTEM_PROMPTS.OUTPUT_GENERATOR,
|
||||
system: WORKFLOW_SYSTEM_PROMPTS.OUTPUT_GENERATOR,
|
||||
model: registeredModel.model,
|
||||
prompt: `Based on the following execution results, generate the structured output according to the schema:
|
||||
|
||||
|
||||
-411
@@ -1,411 +0,0 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
|
||||
import {
|
||||
convertToModelMessages,
|
||||
stepCountIs,
|
||||
streamText,
|
||||
ToolSet,
|
||||
UIDataTypes,
|
||||
UIMessage,
|
||||
UITools,
|
||||
} from 'ai';
|
||||
import { AppPath, type ActorMetadata } from 'twenty-shared/types';
|
||||
import { getAppPath } from 'twenty-shared/utils';
|
||||
import { In } from 'typeorm';
|
||||
|
||||
import { getAllSelectableColumnNames } from 'src/engine/api/utils/get-all-selectable-column-names.utils';
|
||||
import { WorkspaceDomainsService } from 'src/engine/core-modules/domain/workspace-domains/services/workspace-domains.service';
|
||||
import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
|
||||
import {
|
||||
AgentException,
|
||||
AgentExceptionCode,
|
||||
} from 'src/engine/metadata-modules/ai/ai-agent/agent.exception';
|
||||
import { AgentService } from 'src/engine/metadata-modules/ai/ai-agent/agent.service';
|
||||
import { AGENT_CONFIG } from 'src/engine/metadata-modules/ai/ai-agent/constants/agent-config.const';
|
||||
import { AGENT_SYSTEM_PROMPTS } from 'src/engine/metadata-modules/ai/ai-agent/constants/agent-system-prompts.const';
|
||||
import { RecordIdsByObjectMetadataNameSingularType } from 'src/engine/metadata-modules/ai/ai-agent/types/recordIdsByObjectMetadataNameSingular.type';
|
||||
import { repairToolCall } from 'src/engine/metadata-modules/ai/ai-agent/utils/repair-tool-call.util';
|
||||
import { AIBillingService } from 'src/engine/metadata-modules/ai/ai-billing/services/ai-billing.service';
|
||||
import { ToolHints } from 'src/engine/metadata-modules/ai/ai-chat-router/types/tool-hints.interface';
|
||||
import { AI_TELEMETRY_CONFIG } from 'src/engine/metadata-modules/ai/ai-models/constants/ai-telemetry.const';
|
||||
import { AiModelRegistryService } from 'src/engine/metadata-modules/ai/ai-models/services/ai-model-registry.service';
|
||||
import { FlatAgentWithRoleId } from 'src/engine/metadata-modules/flat-agent/types/flat-agent.type';
|
||||
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
|
||||
import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service';
|
||||
import { AgentModelConfigService } from 'src/engine/metadata-modules/ai/ai-models/services/agent-model-config.service';
|
||||
|
||||
import { AgentActorContextService } from './agent-actor-context.service';
|
||||
import { AgentToolGeneratorService } from './agent-tool-generator.service';
|
||||
|
||||
// Re-export for backward compatibility
|
||||
export { type AgentExecutionResult } from 'src/engine/metadata-modules/ai/ai-agent-execution/types/agent-execution-result.type';
|
||||
|
||||
export interface StreamChatResponseResult {
|
||||
stream: ReturnType<typeof streamText>;
|
||||
timings: {
|
||||
contextBuildTimeMs: number;
|
||||
toolGenerationTimeMs: number;
|
||||
aiRequestPrepTimeMs: number;
|
||||
toolCount: number;
|
||||
};
|
||||
}
|
||||
|
||||
@Injectable()
|
||||
export class AgentExecutionService {
|
||||
private readonly logger = new Logger(AgentExecutionService.name);
|
||||
|
||||
constructor(
|
||||
private readonly workspaceDomainsService: WorkspaceDomainsService,
|
||||
private readonly twentyORMGlobalManager: TwentyORMGlobalManager,
|
||||
private readonly aiModelRegistryService: AiModelRegistryService,
|
||||
private readonly agentToolGeneratorService: AgentToolGeneratorService,
|
||||
private readonly agentModelConfigService: AgentModelConfigService,
|
||||
private readonly aiBillingService: AIBillingService,
|
||||
private readonly agentActorContextService: AgentActorContextService,
|
||||
private readonly agentService: AgentService,
|
||||
private readonly workspaceCacheService: WorkspaceCacheService,
|
||||
) {}
|
||||
|
||||
async prepareAIRequestConfig({
|
||||
messages,
|
||||
system,
|
||||
agent,
|
||||
actorContext,
|
||||
roleIds,
|
||||
toolHints,
|
||||
additionalTools,
|
||||
}: {
|
||||
system: string;
|
||||
agent: FlatAgentWithRoleId | null;
|
||||
messages: UIMessage<unknown, UIDataTypes, UITools>[];
|
||||
actorContext?: ActorMetadata;
|
||||
roleIds?: string[];
|
||||
toolHints?: ToolHints;
|
||||
additionalTools?: ToolSet;
|
||||
}) {
|
||||
try {
|
||||
if (agent) {
|
||||
this.logger.log(
|
||||
`Preparing AI request config for agent ${agent.id} with model ${agent.modelId}`,
|
||||
);
|
||||
}
|
||||
|
||||
const registeredModel =
|
||||
await this.aiModelRegistryService.resolveModelForAgent(agent);
|
||||
|
||||
let tools: ToolSet = {};
|
||||
let providerOptions;
|
||||
|
||||
if (agent) {
|
||||
const baseTools =
|
||||
await this.agentToolGeneratorService.generateToolsForAgent(
|
||||
agent.id,
|
||||
agent.workspaceId,
|
||||
actorContext,
|
||||
roleIds,
|
||||
toolHints,
|
||||
);
|
||||
|
||||
const nativeModelTools =
|
||||
this.agentModelConfigService.getNativeModelTools(
|
||||
registeredModel,
|
||||
agent,
|
||||
);
|
||||
|
||||
tools = {
|
||||
...baseTools,
|
||||
...nativeModelTools,
|
||||
...(additionalTools || {}),
|
||||
};
|
||||
|
||||
providerOptions = this.agentModelConfigService.getProviderOptions(
|
||||
registeredModel,
|
||||
agent,
|
||||
);
|
||||
}
|
||||
|
||||
this.logger.log(
|
||||
`Generated ${Object.keys(tools).length} tools for agent (including ${Object.keys(additionalTools || {}).length} additional tools)`,
|
||||
);
|
||||
|
||||
return {
|
||||
system,
|
||||
tools,
|
||||
model: registeredModel.model,
|
||||
messages: convertToModelMessages(messages),
|
||||
stopWhen: stepCountIs(AGENT_CONFIG.MAX_STEPS),
|
||||
providerOptions,
|
||||
experimental_telemetry: AI_TELEMETRY_CONFIG,
|
||||
experimental_repairToolCall: async ({
|
||||
toolCall,
|
||||
tools: toolsForRepair,
|
||||
inputSchema,
|
||||
error,
|
||||
}: {
|
||||
toolCall: {
|
||||
type: 'tool-call';
|
||||
toolCallId: string;
|
||||
toolName: string;
|
||||
input: string;
|
||||
};
|
||||
tools: Record<string, unknown>;
|
||||
inputSchema: (toolCall: { toolName: string }) => unknown;
|
||||
error: Error;
|
||||
}) => {
|
||||
return repairToolCall({
|
||||
toolCall,
|
||||
tools: toolsForRepair,
|
||||
inputSchema,
|
||||
error,
|
||||
model: registeredModel.model,
|
||||
});
|
||||
},
|
||||
};
|
||||
} catch (error) {
|
||||
this.logger.error(
|
||||
`Failed to prepare AI request config for agent ${agent?.id ?? 'no agent'}`,
|
||||
error instanceof Error ? error.stack : error,
|
||||
);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async getContextForSystemPrompt(
|
||||
workspace: WorkspaceEntity,
|
||||
recordIdsByObjectMetadataNameSingular: RecordIdsByObjectMetadataNameSingularType,
|
||||
userWorkspaceId: string,
|
||||
) {
|
||||
const { userWorkspaceRoleMap } =
|
||||
await this.workspaceCacheService.getOrRecompute(workspace.id, [
|
||||
'userWorkspaceRoleMap',
|
||||
]);
|
||||
|
||||
const roleId = userWorkspaceRoleMap[userWorkspaceId];
|
||||
|
||||
if (!roleId) {
|
||||
throw new AgentException(
|
||||
'Failed to retrieve user role.',
|
||||
AgentExceptionCode.ROLE_NOT_FOUND,
|
||||
);
|
||||
}
|
||||
|
||||
const workspaceDataSource =
|
||||
await this.twentyORMGlobalManager.getDataSourceForWorkspace({
|
||||
workspaceId: workspace.id,
|
||||
});
|
||||
|
||||
const flatObjectMetadataMaps =
|
||||
workspaceDataSource.internalContext.flatObjectMetadataMaps;
|
||||
const flatFieldMetadataMaps =
|
||||
workspaceDataSource.internalContext.flatFieldMetadataMaps;
|
||||
const objectIdByNameSingular =
|
||||
workspaceDataSource.internalContext.objectIdByNameSingular;
|
||||
const objectMetadataPermissions = workspaceDataSource.permissionsPerRoleId;
|
||||
|
||||
const contextObject = (
|
||||
await Promise.all(
|
||||
recordIdsByObjectMetadataNameSingular.map(
|
||||
async (recordsWithObjectMetadataNameSingular) => {
|
||||
if (recordsWithObjectMetadataNameSingular.recordIds.length === 0) {
|
||||
return [];
|
||||
}
|
||||
|
||||
const objectMetadataId =
|
||||
objectIdByNameSingular[
|
||||
recordsWithObjectMetadataNameSingular.objectMetadataNameSingular
|
||||
];
|
||||
const objectMetadataMapItem = objectMetadataId
|
||||
? flatObjectMetadataMaps.byId[objectMetadataId]
|
||||
: undefined;
|
||||
|
||||
if (!objectMetadataMapItem) {
|
||||
this.logger.warn(
|
||||
`Object metadata not found for ${recordsWithObjectMetadataNameSingular.objectMetadataNameSingular}`,
|
||||
);
|
||||
|
||||
return [];
|
||||
}
|
||||
|
||||
const repository = workspaceDataSource.getRepository(
|
||||
recordsWithObjectMetadataNameSingular.objectMetadataNameSingular,
|
||||
{ unionOf: [roleId] },
|
||||
);
|
||||
|
||||
const restrictedFields =
|
||||
objectMetadataPermissions?.[roleId]?.[objectMetadataMapItem.id]
|
||||
?.restrictedFields ?? {};
|
||||
|
||||
const hasRestrictedFields = Object.values(restrictedFields).some(
|
||||
(field) => field.canRead === false,
|
||||
);
|
||||
|
||||
const selectOptions = hasRestrictedFields
|
||||
? getAllSelectableColumnNames({
|
||||
restrictedFields,
|
||||
objectMetadata: {
|
||||
objectMetadataMapItem,
|
||||
flatFieldMetadataMaps,
|
||||
},
|
||||
})
|
||||
: undefined;
|
||||
|
||||
return (
|
||||
await repository.find({
|
||||
...(selectOptions && { select: selectOptions }),
|
||||
where: {
|
||||
id: In(recordsWithObjectMetadataNameSingular.recordIds),
|
||||
},
|
||||
})
|
||||
).map((record) => {
|
||||
return {
|
||||
...record,
|
||||
resourceUrl: this.workspaceDomainsService.buildWorkspaceURL({
|
||||
workspace,
|
||||
pathname: getAppPath(AppPath.RecordShowPage, {
|
||||
objectNameSingular:
|
||||
recordsWithObjectMetadataNameSingular.objectMetadataNameSingular,
|
||||
objectRecordId: record.id,
|
||||
}),
|
||||
}),
|
||||
};
|
||||
});
|
||||
},
|
||||
),
|
||||
)
|
||||
).flat(2);
|
||||
|
||||
return JSON.stringify(contextObject);
|
||||
}
|
||||
|
||||
async streamChatResponse({
|
||||
workspace,
|
||||
userWorkspaceId,
|
||||
agentId,
|
||||
messages,
|
||||
recordIdsByObjectMetadataNameSingular,
|
||||
toolHints,
|
||||
additionalTools,
|
||||
}: {
|
||||
workspace: WorkspaceEntity;
|
||||
userWorkspaceId: string;
|
||||
agentId: string;
|
||||
messages: UIMessage<unknown, UIDataTypes, UITools>[];
|
||||
recordIdsByObjectMetadataNameSingular: RecordIdsByObjectMetadataNameSingularType;
|
||||
toolHints?: ToolHints;
|
||||
additionalTools?: ToolSet;
|
||||
}): Promise<{
|
||||
stream: ReturnType<typeof streamText>;
|
||||
timings: {
|
||||
contextBuildTimeMs: number;
|
||||
toolGenerationTimeMs: number;
|
||||
aiRequestPrepTimeMs: number;
|
||||
toolCount: number;
|
||||
};
|
||||
contextInfo: {
|
||||
contextString: string;
|
||||
contextRecordCount: number;
|
||||
contextSizeBytes: number;
|
||||
};
|
||||
}> {
|
||||
try {
|
||||
const agent = await this.agentService.findOneAgentById({
|
||||
workspaceId: workspace.id,
|
||||
id: agentId,
|
||||
});
|
||||
|
||||
const contextBuildStart = Date.now();
|
||||
let contextPart = '';
|
||||
let contextRecordCount = 0;
|
||||
|
||||
if (recordIdsByObjectMetadataNameSingular.length > 0) {
|
||||
contextPart = await this.getContextForSystemPrompt(
|
||||
workspace,
|
||||
recordIdsByObjectMetadataNameSingular,
|
||||
userWorkspaceId,
|
||||
);
|
||||
|
||||
try {
|
||||
const contextData = JSON.parse(contextPart);
|
||||
|
||||
contextRecordCount = Array.isArray(contextData)
|
||||
? contextData.length
|
||||
: 0;
|
||||
} catch (error) {
|
||||
this.logger.warn('Failed to parse context for record count:', error);
|
||||
}
|
||||
}
|
||||
|
||||
const contextString = contextPart ? `\n\nCONTEXT:\n${contextPart}` : '';
|
||||
const contextBuildTime = Date.now() - contextBuildStart;
|
||||
|
||||
const { actorContext, roleId } =
|
||||
await this.agentActorContextService.buildUserAndAgentActorContext(
|
||||
userWorkspaceId,
|
||||
workspace.id,
|
||||
);
|
||||
|
||||
const aiRequestPrepStart = Date.now();
|
||||
|
||||
const aiRequestConfig = await this.prepareAIRequestConfig({
|
||||
system: `${AGENT_SYSTEM_PROMPTS.BASE}\n${AGENT_SYSTEM_PROMPTS.CHAT_ADDITIONS}\n\n${agent.prompt}${contextString}`,
|
||||
agent,
|
||||
messages,
|
||||
actorContext,
|
||||
roleIds: [roleId, ...(agent?.roleId ? [agent?.roleId] : [])],
|
||||
toolHints,
|
||||
additionalTools,
|
||||
});
|
||||
|
||||
const aiRequestPrepTime = Date.now() - aiRequestPrepStart;
|
||||
const toolCount = Object.keys(aiRequestConfig.tools || {}).length;
|
||||
const toolGenerationTime = aiRequestPrepTime;
|
||||
|
||||
this.logger.log(
|
||||
`Sending request to AI model with ${messages.length} messages and ${toolCount} tools`,
|
||||
);
|
||||
|
||||
const model =
|
||||
await this.aiModelRegistryService.resolveModelForAgent(agent);
|
||||
|
||||
const stream = streamText(aiRequestConfig);
|
||||
|
||||
stream.usage
|
||||
.then((usage) => {
|
||||
this.aiBillingService.calculateAndBillUsage(
|
||||
model.modelId,
|
||||
usage,
|
||||
workspace.id,
|
||||
agent.id,
|
||||
);
|
||||
})
|
||||
.catch((usageError) => {
|
||||
this.logger.error('Failed to get usage information:', usageError);
|
||||
});
|
||||
|
||||
return {
|
||||
stream,
|
||||
timings: {
|
||||
contextBuildTimeMs: contextBuildTime,
|
||||
toolGenerationTimeMs: toolGenerationTime,
|
||||
aiRequestPrepTimeMs: aiRequestPrepTime,
|
||||
toolCount,
|
||||
},
|
||||
contextInfo: {
|
||||
contextString: contextPart,
|
||||
contextRecordCount,
|
||||
contextSizeBytes: contextPart
|
||||
? Buffer.byteLength(contextPart, 'utf8')
|
||||
: 0,
|
||||
},
|
||||
};
|
||||
} catch (error) {
|
||||
this.logger.error('Error in streamChatResponse:', error);
|
||||
throw new AgentException(
|
||||
error instanceof Error
|
||||
? error.message
|
||||
: 'Failed to stream chat response',
|
||||
AgentExceptionCode.AGENT_EXECUTION_FAILED,
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
-251
@@ -1,251 +0,0 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
|
||||
import { type WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
|
||||
import { AgentService } from 'src/engine/metadata-modules/ai/ai-agent/agent.service';
|
||||
import { type RecordIdsByObjectMetadataNameSingularType } from 'src/engine/metadata-modules/ai/ai-agent/types/recordIdsByObjectMetadataNameSingular.type';
|
||||
import { type PlanStep } from 'src/engine/metadata-modules/ai/ai-chat-router/types/router-result.interface';
|
||||
import { STANDARD_AGENT_DEFINITIONS } from 'src/engine/workspace-manager/workspace-sync-metadata/standard-agents/standard-agent-definitions';
|
||||
|
||||
import { AgentExecutionService } from './agent-execution.service';
|
||||
|
||||
export type PlanExecutionProgress = {
|
||||
type: 'plan-generated' | 'step-started' | 'step-completed';
|
||||
stepNumber?: number;
|
||||
agentName?: string;
|
||||
task?: string;
|
||||
output?: string;
|
||||
totalSteps?: number;
|
||||
reasoning?: string;
|
||||
};
|
||||
|
||||
export type StepResult = {
|
||||
stepNumber: number;
|
||||
agentName: string;
|
||||
output: string;
|
||||
};
|
||||
|
||||
export type PlanExecutionResult = {
|
||||
finalOutput: string;
|
||||
stepResults: StepResult[];
|
||||
};
|
||||
|
||||
@Injectable()
|
||||
export class AgentPlanExecutorService {
|
||||
private readonly logger = new Logger(AgentPlanExecutorService.name);
|
||||
|
||||
constructor(
|
||||
private readonly agentExecutionService: AgentExecutionService,
|
||||
private readonly agentService: AgentService,
|
||||
) {}
|
||||
|
||||
async executePlan({
|
||||
steps,
|
||||
reasoning,
|
||||
workspace,
|
||||
userWorkspaceId,
|
||||
recordIdsByObjectMetadataNameSingular,
|
||||
onProgress,
|
||||
writer,
|
||||
}: {
|
||||
steps: PlanStep[];
|
||||
reasoning: string;
|
||||
workspace: WorkspaceEntity;
|
||||
userWorkspaceId: string;
|
||||
recordIdsByObjectMetadataNameSingular: RecordIdsByObjectMetadataNameSingularType;
|
||||
onProgress?: (progress: PlanExecutionProgress) => void;
|
||||
writer?: {
|
||||
write: (chunk: unknown) => void;
|
||||
merge: (stream: unknown) => void;
|
||||
};
|
||||
}): Promise<PlanExecutionResult> {
|
||||
this.logger.log(`Executing plan with ${steps.length} steps`);
|
||||
|
||||
onProgress?.({
|
||||
type: 'plan-generated',
|
||||
totalSteps: steps.length,
|
||||
reasoning,
|
||||
});
|
||||
|
||||
const stepResults: StepResult[] = [];
|
||||
|
||||
for (const step of steps) {
|
||||
try {
|
||||
this.logger.log(
|
||||
`[PLAN EXECUTION] Step ${step.stepNumber}: Looking up agent "${step.agentName}"`,
|
||||
);
|
||||
|
||||
const agent = await this.agentService.findOneAgentByName({
|
||||
name: step.agentName,
|
||||
workspaceId: workspace.id,
|
||||
});
|
||||
|
||||
this.logger.log(
|
||||
`[PLAN EXECUTION] Step ${step.stepNumber}: Found agent "${agent.label}" (${agent.id})`,
|
||||
);
|
||||
|
||||
onProgress?.({
|
||||
type: 'step-started',
|
||||
stepNumber: step.stepNumber,
|
||||
agentName: step.agentName,
|
||||
task: step.task,
|
||||
});
|
||||
|
||||
const dependencyOutputs = this.gatherDependencyOutputs(
|
||||
step,
|
||||
stepResults,
|
||||
);
|
||||
|
||||
const promptWithContext = this.buildStepPrompt(step, dependencyOutputs);
|
||||
|
||||
const { stream: stepStream } =
|
||||
await this.agentExecutionService.streamChatResponse({
|
||||
workspace,
|
||||
agentId: agent.id,
|
||||
userWorkspaceId,
|
||||
messages: [
|
||||
{
|
||||
id: `step-${step.stepNumber}`,
|
||||
role: 'user' as const,
|
||||
parts: [{ type: 'text' as const, text: promptWithContext }],
|
||||
},
|
||||
],
|
||||
recordIdsByObjectMetadataNameSingular,
|
||||
});
|
||||
|
||||
let stepOutput = '';
|
||||
|
||||
if (writer) {
|
||||
writer.merge(
|
||||
stepStream.toUIMessageStream({
|
||||
onError: (error) => {
|
||||
return error instanceof Error ? error.message : String(error);
|
||||
},
|
||||
sendStart: false,
|
||||
onFinish: async ({ responseMessage }) => {
|
||||
stepOutput = responseMessage.parts
|
||||
.filter((part) => part.type === 'text')
|
||||
.map((part) => {
|
||||
if (part.type === 'text') {
|
||||
return part.text;
|
||||
}
|
||||
|
||||
return '';
|
||||
})
|
||||
.join('');
|
||||
},
|
||||
}),
|
||||
);
|
||||
|
||||
await stepStream.text;
|
||||
} else {
|
||||
stepOutput = await stepStream.text;
|
||||
}
|
||||
|
||||
stepResults.push({
|
||||
stepNumber: step.stepNumber,
|
||||
agentName: step.agentName,
|
||||
output: stepOutput,
|
||||
});
|
||||
|
||||
onProgress?.({
|
||||
type: 'step-completed',
|
||||
stepNumber: step.stepNumber,
|
||||
agentName: step.agentName,
|
||||
output: stepOutput,
|
||||
});
|
||||
|
||||
this.logger.log(
|
||||
`Completed step ${step.stepNumber}: ${step.task.substring(0, 50)}...`,
|
||||
);
|
||||
} catch (error) {
|
||||
this.logger.error(
|
||||
`Failed to execute step ${step.stepNumber}: ${step.task}`,
|
||||
error,
|
||||
);
|
||||
|
||||
throw new Error(
|
||||
`Plan execution failed at step ${step.stepNumber}: ${error.message}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
const finalOutput = this.synthesizeResults(stepResults, steps);
|
||||
|
||||
return {
|
||||
finalOutput,
|
||||
stepResults,
|
||||
};
|
||||
}
|
||||
|
||||
private gatherDependencyOutputs(
|
||||
step: PlanStep,
|
||||
previousResults: StepResult[],
|
||||
): string {
|
||||
if (!step.dependsOn || step.dependsOn.length === 0) {
|
||||
return '';
|
||||
}
|
||||
|
||||
const dependencyOutputs = step.dependsOn
|
||||
.map((depStepNum) => {
|
||||
const depResult = previousResults.find(
|
||||
(result) => result.stepNumber === depStepNum,
|
||||
);
|
||||
|
||||
if (!depResult) {
|
||||
throw new Error(
|
||||
`Dependency step ${depStepNum} not found for step ${step.stepNumber}`,
|
||||
);
|
||||
}
|
||||
|
||||
return `Step ${depStepNum} output:\n${depResult.output}`;
|
||||
})
|
||||
.join('\n\n');
|
||||
|
||||
return dependencyOutputs;
|
||||
}
|
||||
|
||||
private buildStepPrompt(step: PlanStep, dependencyOutputs: string): string {
|
||||
let prompt = `Task: ${step.task}\n\nExpected output: ${step.expectedOutput}`;
|
||||
|
||||
if (dependencyOutputs) {
|
||||
prompt += `\n\nPrevious step results:\n${dependencyOutputs}`;
|
||||
}
|
||||
|
||||
return prompt;
|
||||
}
|
||||
|
||||
private synthesizeResults(
|
||||
stepResults: StepResult[],
|
||||
steps: PlanStep[],
|
||||
): string {
|
||||
const lastStep = stepResults[stepResults.length - 1];
|
||||
|
||||
if (!lastStep) {
|
||||
return 'No results produced';
|
||||
}
|
||||
|
||||
const lastStepDefinition = steps.find(
|
||||
(s) => s.stepNumber === lastStep.stepNumber,
|
||||
);
|
||||
|
||||
if (lastStepDefinition) {
|
||||
const agentDefinition = STANDARD_AGENT_DEFINITIONS.find(
|
||||
(def) => def.name === lastStepDefinition.agentName,
|
||||
);
|
||||
|
||||
if (agentDefinition?.outputStrategy === 'direct') {
|
||||
return lastStep.output;
|
||||
}
|
||||
}
|
||||
|
||||
const summary = stepResults
|
||||
.map((result) => {
|
||||
const step = steps.find((s) => s.stepNumber === result.stepNumber);
|
||||
|
||||
return `**Step ${result.stepNumber}: ${step?.task || 'Unknown task'}**\n${result.output}`;
|
||||
})
|
||||
.join('\n\n---\n\n');
|
||||
|
||||
return summary;
|
||||
}
|
||||
}
|
||||
-47
@@ -1,47 +0,0 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
|
||||
import { type ToolSet } from 'ai';
|
||||
|
||||
import type { ActorMetadata } from 'twenty-shared/types';
|
||||
|
||||
import { ToolCategory } from 'src/engine/core-modules/tool-provider/enums/tool-category.enum';
|
||||
import { ToolProviderService } from 'src/engine/core-modules/tool-provider/services/tool-provider.service';
|
||||
import type { ToolHints } from 'src/engine/metadata-modules/ai/ai-chat-router/types/tool-hints.interface';
|
||||
|
||||
@Injectable()
|
||||
export class AgentToolGeneratorService {
|
||||
private readonly logger = new Logger(AgentToolGeneratorService.name);
|
||||
|
||||
constructor(private readonly toolProvider: ToolProviderService) {}
|
||||
|
||||
// Generates base tools for chat context (DATABASE_CRUD and ACTION)
|
||||
// Additional tools (WORKFLOW, METADATA) are provided via additionalTools
|
||||
// from ChatToolsProviderService to avoid circular dependencies
|
||||
async generateToolsForAgent(
|
||||
agentId: string,
|
||||
workspaceId: string,
|
||||
actorContext?: ActorMetadata,
|
||||
roleIds?: string[],
|
||||
toolHints?: ToolHints,
|
||||
): Promise<ToolSet> {
|
||||
try {
|
||||
return await this.toolProvider.getTools({
|
||||
workspaceId,
|
||||
categories: [ToolCategory.DATABASE_CRUD, ToolCategory.ACTION],
|
||||
rolePermissionConfig: roleIds ? { intersectionOf: roleIds } : undefined,
|
||||
actorContext,
|
||||
toolHints,
|
||||
wrapWithErrorContext: true,
|
||||
});
|
||||
} catch (toolError) {
|
||||
const errorMessage =
|
||||
toolError instanceof Error ? toolError.message : 'Unknown error';
|
||||
|
||||
this.logger.warn(
|
||||
`Failed to generate tools for agent ${agentId}: ${errorMessage}. Proceeding without tools.`,
|
||||
);
|
||||
|
||||
return {};
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -2,7 +2,7 @@ import { Injectable } from '@nestjs/common';
|
||||
import { InjectRepository } from '@nestjs/typeorm';
|
||||
|
||||
import { isDefined } from 'twenty-shared/utils';
|
||||
import { Repository } from 'typeorm';
|
||||
import { ILike, IsNull, Repository } from 'typeorm';
|
||||
|
||||
import { ApplicationService } from 'src/engine/core-modules/application/application.service';
|
||||
import { type CreateAgentInput } from 'src/engine/metadata-modules/ai/ai-agent/dtos/create-agent.input';
|
||||
@@ -300,4 +300,26 @@ export class AgentService {
|
||||
roleId: roleId ?? null,
|
||||
};
|
||||
}
|
||||
|
||||
async searchAgents(
|
||||
query: string,
|
||||
workspaceId: string,
|
||||
options: { limit: number } = { limit: 2 },
|
||||
): Promise<AgentEntity[]> {
|
||||
const queryLower = query.toLowerCase();
|
||||
|
||||
return this.agentRepository.find({
|
||||
where: [
|
||||
{ workspaceId, deletedAt: IsNull(), name: ILike(`%${queryLower}%`) },
|
||||
{
|
||||
workspaceId,
|
||||
deletedAt: IsNull(),
|
||||
description: ILike(`%${queryLower}%`),
|
||||
},
|
||||
{ workspaceId, deletedAt: IsNull(), label: ILike(`%${queryLower}%`) },
|
||||
],
|
||||
take: options.limit,
|
||||
order: { name: 'ASC' },
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
+14
-96
@@ -1,107 +1,25 @@
|
||||
export const AGENT_SYSTEM_PROMPTS = {
|
||||
BASE: `Tool usage strategy:
|
||||
// System prompts for Workflow Agents (automated execution only)
|
||||
// NOTE: For user-facing chat, use CHAT_SYSTEM_PROMPTS from ai-chat/constants
|
||||
|
||||
export const WORKFLOW_SYSTEM_PROMPTS = {
|
||||
// Core workflow execution behavior
|
||||
BASE: `You are executing as part of a workflow automation in Twenty CRM.
|
||||
|
||||
Tool usage strategy:
|
||||
- Chain multiple tools to solve complex tasks
|
||||
- If a tool fails, try alternative approaches
|
||||
- Use results from one tool to inform the next
|
||||
- Don't give up after first failure - be persistent
|
||||
- Validate assumptions before making changes
|
||||
|
||||
Error recovery:
|
||||
- Analyze error messages to understand what went wrong
|
||||
- Adjust parameters or try different tools
|
||||
- Only give up after exhausting reasonable alternatives
|
||||
Context:
|
||||
- Your output may be used by downstream workflow nodes
|
||||
- Be thorough and include all relevant data
|
||||
- Focus on completing the task efficiently
|
||||
|
||||
Permissions:
|
||||
- Only perform actions your role allows
|
||||
- Explain limitations if you lack permissions`,
|
||||
|
||||
CHAT_ADDITIONS: `
|
||||
Format responses with markdown for clarity (headings, lists, code blocks, tables).`,
|
||||
|
||||
WORKFLOW_ADDITIONS: `
|
||||
Context:
|
||||
- You are executing as part of a workflow automation
|
||||
- Your output may be used by downstream nodes
|
||||
- Be thorough and include all relevant data`,
|
||||
|
||||
ROUTER: (
|
||||
agentDescriptions: string,
|
||||
) => `You are an AI router that decides how to handle user messages.
|
||||
|
||||
Available agents:
|
||||
${agentDescriptions}
|
||||
|
||||
Decision process:
|
||||
1. Can ONE agent handle this entirely? → Use "simple" strategy
|
||||
2. Does it require MULTIPLE agents working together? → Use "planned" strategy
|
||||
|
||||
Agent selection rules (CRITICAL):
|
||||
- **metadata-builder**: For modifying the DATA MODEL/SCHEMA - creating new object types, adding fields to objects, managing object structure. Use when user wants to define NEW TYPES of entities or add properties to existing types.
|
||||
- **data-manipulator**: For CRUD operations on existing RECORDS/DATA - creating company records, finding people, updating opportunities. Use when user wants to work with actual data entries.
|
||||
- **helper**: ONLY for questions about HOW TO USE Twenty (features, setup, documentation)
|
||||
- **researcher**: For finding external information from the web
|
||||
- **workflow-builder**: For creating automation workflows
|
||||
- **dashboard-builder**: For creating and managing dashboards and visualizations
|
||||
|
||||
CRITICAL DISTINCTION:
|
||||
- "Create an object called Project" → metadata-builder (creating a new object TYPE in the schema)
|
||||
- "Create a company called Acme" → data-manipulator (creating a company RECORD)
|
||||
- "Add a field to Company" → metadata-builder (modifying schema)
|
||||
- "Update the company's phone number" → data-manipulator (modifying data)
|
||||
|
||||
Use "planned" strategy when:
|
||||
- Request needs custom code AND context from data/research
|
||||
- Code generation requires knowing schemas, APIs, or external data
|
||||
- Multiple specialized capabilities must combine (code + data + research)
|
||||
|
||||
Use "simple" strategy for:
|
||||
- Single-agent tasks (data operations, research, documentation lookup)
|
||||
- Standard workflow creation (no custom code needed)
|
||||
|
||||
Examples:
|
||||
|
||||
Simple: "Show me all companies with >100 employees"
|
||||
→ { strategy: "simple", agentName: "data-manipulator", toolHints: { relevantObjects: ["company"], operations: ["find"] } }
|
||||
|
||||
Simple: "Create 30 companies in the automobile industry with 2 people each"
|
||||
→ { strategy: "simple", agentName: "data-manipulator", toolHints: { relevantObjects: ["company", "person"], operations: ["create"] } }
|
||||
|
||||
Simple: "Update all opportunities in stage 'Qualified' to 'Proposal'"
|
||||
→ { strategy: "simple", agentName: "data-manipulator", toolHints: { relevantObjects: ["opportunity"], operations: ["find", "update"] } }
|
||||
|
||||
Simple: "What's the latest news about AI trends?"
|
||||
→ { strategy: "simple", agentName: "researcher" }
|
||||
|
||||
Simple: "How do I set up email sync in Twenty?"
|
||||
→ { strategy: "simple", agentName: "helper" }
|
||||
|
||||
Simple: "Create a workflow that emails customers when deals close"
|
||||
→ { strategy: "simple", agentName: "workflow-builder" }
|
||||
|
||||
Simple: "Create an object called Project" or "Create a new custom object for tracking invoices"
|
||||
→ { strategy: "simple", agentName: "metadata-builder" }
|
||||
|
||||
Simple: "Add a budget field to the Project object" or "Add a phone field to Company"
|
||||
→ { strategy: "simple", agentName: "metadata-builder" }
|
||||
|
||||
Planned: "Research information about Meta and update the company record"
|
||||
→ {
|
||||
strategy: "planned",
|
||||
plan: {
|
||||
steps: [
|
||||
{ stepNumber: 1, agentName: "researcher", task: "Look up current information about Meta (employee count, headquarters, revenue, etc.)", expectedOutput: "Company facts and data" },
|
||||
{ stepNumber: 2, agentName: "data-manipulator", task: "Update the Meta company record with the researched information", expectedOutput: "Updated company record", dependsOn: [1] }
|
||||
],
|
||||
reasoning: "Requires web research followed by database update"
|
||||
}
|
||||
}
|
||||
|
||||
For simple strategy toolHints:
|
||||
- relevantObjects: Extract object names (e.g., ["company", "person"])
|
||||
- operations: ["find", "create", "update", "delete"]
|
||||
|
||||
Keep plans minimal and only use planning when truly necessary.`,
|
||||
- Only perform actions your role allows`,
|
||||
|
||||
// Structured output generation for workflow data passing
|
||||
OUTPUT_GENERATOR: `You are a structured output generator for a workflow system. Your role is to convert the provided execution results into a structured format according to a specific schema.
|
||||
|
||||
Context: Before this call, the system executed generateText with tools to perform any required actions and gather information. The execution results you receive include both the AI agent's analysis and any tool outputs from database operations, HTTP requests, data retrieval, or other actions.
|
||||
|
||||
-27
@@ -1,27 +0,0 @@
|
||||
import { Module } from '@nestjs/common';
|
||||
import { TypeOrmModule } from '@nestjs/typeorm';
|
||||
|
||||
import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
|
||||
import { AgentEntity } from 'src/engine/metadata-modules/ai/ai-agent/entities/agent.entity';
|
||||
import { AiModelsModule } from 'src/engine/metadata-modules/ai/ai-models/ai-models.module';
|
||||
import { ObjectMetadataModule } from 'src/engine/metadata-modules/object-metadata/object-metadata.module';
|
||||
|
||||
import { AiChatRouterService } from './ai-chat-router.service';
|
||||
|
||||
import { AiChatRouterPlanGeneratorService } from './services/ai-chat-router-plan-generator.service';
|
||||
import { AiChatRouterStrategyDeciderService } from './services/ai-chat-router-strategy-decider.service';
|
||||
|
||||
@Module({
|
||||
imports: [
|
||||
TypeOrmModule.forFeature([AgentEntity, WorkspaceEntity]),
|
||||
AiModelsModule,
|
||||
ObjectMetadataModule,
|
||||
],
|
||||
providers: [
|
||||
AiChatRouterService,
|
||||
AiChatRouterStrategyDeciderService,
|
||||
AiChatRouterPlanGeneratorService,
|
||||
],
|
||||
exports: [AiChatRouterService],
|
||||
})
|
||||
export class AiChatRouterModule {}
|
||||
-391
@@ -1,391 +0,0 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
import { InjectRepository } from '@nestjs/typeorm';
|
||||
|
||||
import { type UIDataTypes, type UIMessage, type UITools } from 'ai';
|
||||
import { IsNull, type Repository } from 'typeorm';
|
||||
|
||||
import { AgentEntity } from 'src/engine/metadata-modules/ai/ai-agent/entities/agent.entity';
|
||||
import { isWorkflowRelatedObject } from 'src/engine/metadata-modules/ai/ai-agent/utils/is-workflow-related-object.util';
|
||||
import { type ModelId } from 'src/engine/metadata-modules/ai/ai-models/constants/ai-models.const';
|
||||
import { ObjectMetadataService } from 'src/engine/metadata-modules/object-metadata/object-metadata.service';
|
||||
import { HELPER_AGENT } from 'src/engine/workspace-manager/workspace-sync-metadata/standard-agents/agents/helper-agent';
|
||||
|
||||
import { AiChatRouterPlanGeneratorService } from './services/ai-chat-router-plan-generator.service';
|
||||
import {
|
||||
AiChatRouterStrategyDeciderService,
|
||||
type StrategyDecision,
|
||||
} from './services/ai-chat-router-strategy-decider.service';
|
||||
import {
|
||||
type RouterDebugInfo,
|
||||
type UnifiedRouterResult,
|
||||
} from './types/router-result.interface';
|
||||
import { type ToolHints } from './types/tool-hints.interface';
|
||||
|
||||
export interface AiChatRouterContext {
|
||||
messages: UIMessage<unknown, UIDataTypes, UITools>[];
|
||||
workspaceId: string;
|
||||
fastModel: ModelId;
|
||||
smartModel: ModelId;
|
||||
}
|
||||
|
||||
@Injectable()
|
||||
export class AiChatRouterService {
|
||||
private readonly logger = new Logger(AiChatRouterService.name);
|
||||
|
||||
constructor(
|
||||
@InjectRepository(AgentEntity)
|
||||
private readonly agentRepository: Repository<AgentEntity>,
|
||||
private readonly strategyDecider: AiChatRouterStrategyDeciderService,
|
||||
private readonly planGenerator: AiChatRouterPlanGeneratorService,
|
||||
private readonly objectMetadataService: ObjectMetadataService,
|
||||
) {}
|
||||
|
||||
async routeMessage(
|
||||
context: AiChatRouterContext,
|
||||
includeDebugInfo = false,
|
||||
): Promise<UnifiedRouterResult> {
|
||||
try {
|
||||
const { messages, workspaceId, fastModel, smartModel } = context;
|
||||
const availableAgents = await this.getAvailableAgents(workspaceId);
|
||||
|
||||
this.logger.log(
|
||||
`[ROUTER] Available agents (${availableAgents.length}): ${availableAgents.map((a) => `${a.label} (${a.name})`).join(', ')}`,
|
||||
);
|
||||
|
||||
if (availableAgents.length === 0) {
|
||||
return await this.handleNoAgentsAvailable(workspaceId);
|
||||
}
|
||||
|
||||
const debugInfo = this.createDebugInfo(
|
||||
includeDebugInfo,
|
||||
availableAgents,
|
||||
smartModel,
|
||||
fastModel,
|
||||
);
|
||||
|
||||
if (availableAgents.length === 1) {
|
||||
return this.createSimpleResult(availableAgents[0], debugInfo);
|
||||
}
|
||||
|
||||
return await this.routeToMultipleAgents({
|
||||
messages,
|
||||
workspaceId,
|
||||
availableAgents,
|
||||
fastModel,
|
||||
smartModel,
|
||||
debugInfo,
|
||||
});
|
||||
} catch (error) {
|
||||
return await this.handleRoutingError(error, context.workspaceId);
|
||||
}
|
||||
}
|
||||
|
||||
private async handleNoAgentsAvailable(
|
||||
workspaceId: string,
|
||||
): Promise<UnifiedRouterResult> {
|
||||
this.logger.warn('No agents available for routing');
|
||||
|
||||
const helperAgent = await this.getHelperAgent(workspaceId);
|
||||
|
||||
if (!helperAgent) {
|
||||
throw new Error('No helper agent available');
|
||||
}
|
||||
|
||||
return {
|
||||
strategy: 'simple',
|
||||
agent: helperAgent,
|
||||
};
|
||||
}
|
||||
|
||||
private createDebugInfo(
|
||||
includeDebugInfo: boolean,
|
||||
availableAgents: AgentEntity[],
|
||||
smartModel: ModelId,
|
||||
fastModel: ModelId,
|
||||
): RouterDebugInfo | undefined {
|
||||
if (!includeDebugInfo) {
|
||||
return undefined;
|
||||
}
|
||||
|
||||
return {
|
||||
availableAgents: availableAgents.map((agent) => ({
|
||||
id: agent.id,
|
||||
label: agent.label,
|
||||
})),
|
||||
routerModel: String(smartModel ?? fastModel),
|
||||
promptTokens: 0,
|
||||
completionTokens: 0,
|
||||
totalTokens: 0,
|
||||
};
|
||||
}
|
||||
|
||||
private createSimpleResult(
|
||||
agent: AgentEntity,
|
||||
debugInfo?: RouterDebugInfo,
|
||||
toolHints?: ToolHints,
|
||||
): UnifiedRouterResult {
|
||||
return {
|
||||
strategy: 'simple',
|
||||
agent,
|
||||
toolHints,
|
||||
debugInfo,
|
||||
};
|
||||
}
|
||||
|
||||
private async routeToMultipleAgents(params: {
|
||||
messages: UIMessage<unknown, UIDataTypes, UITools>[];
|
||||
workspaceId: string;
|
||||
availableAgents: AgentEntity[];
|
||||
fastModel: ModelId;
|
||||
smartModel: ModelId;
|
||||
debugInfo?: RouterDebugInfo;
|
||||
}): Promise<UnifiedRouterResult> {
|
||||
const {
|
||||
messages,
|
||||
workspaceId,
|
||||
availableAgents,
|
||||
fastModel,
|
||||
smartModel,
|
||||
debugInfo,
|
||||
} = params;
|
||||
|
||||
const workspaceObjectsList =
|
||||
await this.buildWorkspaceObjectsList(workspaceId);
|
||||
const agentDescriptions = this.buildAgentDescriptions(
|
||||
availableAgents,
|
||||
workspaceObjectsList,
|
||||
);
|
||||
|
||||
this.logRoutingContext(messages, agentDescriptions);
|
||||
|
||||
const strategyDecision = await this.strategyDecider.decideStrategy({
|
||||
messages,
|
||||
availableAgents,
|
||||
agentDescriptions,
|
||||
fastModel,
|
||||
});
|
||||
|
||||
if (strategyDecision.strategy === 'simple') {
|
||||
return this.handleSimpleStrategy(
|
||||
strategyDecision,
|
||||
availableAgents,
|
||||
debugInfo,
|
||||
);
|
||||
}
|
||||
|
||||
return await this.handlePlannedStrategy({
|
||||
messages,
|
||||
availableAgents,
|
||||
agentDescriptions,
|
||||
smartModel,
|
||||
debugInfo,
|
||||
});
|
||||
}
|
||||
|
||||
private logRoutingContext(
|
||||
messages: UIMessage<unknown, UIDataTypes, UITools>[],
|
||||
agentDescriptions: string,
|
||||
) {
|
||||
this.logger.log(`[ROUTER] Agent descriptions:\n${agentDescriptions}`);
|
||||
|
||||
const currentMessage =
|
||||
messages[messages.length - 1]?.parts.find((part) => part.type === 'text')
|
||||
?.text || '';
|
||||
|
||||
this.logger.log(
|
||||
`[ROUTER] User message: "${currentMessage.substring(0, 100)}..."`,
|
||||
);
|
||||
}
|
||||
|
||||
private handleSimpleStrategy(
|
||||
strategyDecision: StrategyDecision,
|
||||
availableAgents: AgentEntity[],
|
||||
debugInfo?: RouterDebugInfo,
|
||||
): UnifiedRouterResult {
|
||||
if (!strategyDecision.agentName) {
|
||||
throw new Error('agentName is required for simple strategy');
|
||||
}
|
||||
|
||||
const selectedAgent = this.findAgentByName(
|
||||
strategyDecision.agentName,
|
||||
availableAgents,
|
||||
);
|
||||
|
||||
this.logger.log(
|
||||
`[ROUTER] Routing to ${selectedAgent.label} (${selectedAgent.name})`,
|
||||
);
|
||||
|
||||
return this.createSimpleResult(
|
||||
selectedAgent,
|
||||
debugInfo,
|
||||
strategyDecision.toolHints,
|
||||
);
|
||||
}
|
||||
|
||||
private async handlePlannedStrategy(params: {
|
||||
messages: UIMessage<unknown, UIDataTypes, UITools>[];
|
||||
availableAgents: AgentEntity[];
|
||||
agentDescriptions: string;
|
||||
smartModel: ModelId;
|
||||
debugInfo?: RouterDebugInfo;
|
||||
}): Promise<UnifiedRouterResult> {
|
||||
const {
|
||||
messages,
|
||||
availableAgents,
|
||||
agentDescriptions,
|
||||
smartModel,
|
||||
debugInfo,
|
||||
} = params;
|
||||
|
||||
const plan = await this.planGenerator.generatePlan({
|
||||
messages,
|
||||
availableAgents,
|
||||
agentDescriptions,
|
||||
smartModel,
|
||||
});
|
||||
|
||||
if (plan.steps.length === 1) {
|
||||
return this.convertSingleStepPlanToSimple(
|
||||
plan.steps[0],
|
||||
availableAgents,
|
||||
debugInfo,
|
||||
);
|
||||
}
|
||||
|
||||
this.logger.log(
|
||||
`[ROUTER] Executing planned strategy with ${plan.steps.length} steps`,
|
||||
);
|
||||
|
||||
return {
|
||||
strategy: 'planned',
|
||||
plan,
|
||||
debugInfo,
|
||||
};
|
||||
}
|
||||
|
||||
private convertSingleStepPlanToSimple(
|
||||
step: { agentName: string },
|
||||
availableAgents: AgentEntity[],
|
||||
debugInfo?: RouterDebugInfo,
|
||||
): UnifiedRouterResult {
|
||||
this.logger.log(
|
||||
`[ROUTER] Plan has only 1 step, converting to simple strategy`,
|
||||
);
|
||||
|
||||
const selectedAgent = this.findAgentByName(step.agentName, availableAgents);
|
||||
|
||||
return this.createSimpleResult(selectedAgent, debugInfo);
|
||||
}
|
||||
|
||||
private findAgentByName(
|
||||
agentName: string,
|
||||
availableAgents: AgentEntity[],
|
||||
): AgentEntity {
|
||||
const selectedAgent = availableAgents.find(
|
||||
(agent) => agent.name === agentName,
|
||||
);
|
||||
|
||||
if (!selectedAgent) {
|
||||
this.logger.error(
|
||||
`[ROUTER] Agent "${agentName}" not found in available agents: ${availableAgents.map((a) => a.name).join(', ')}`,
|
||||
);
|
||||
throw new Error(`Selected agent ${agentName} not found`);
|
||||
}
|
||||
|
||||
return selectedAgent;
|
||||
}
|
||||
|
||||
private async handleRoutingError(
|
||||
error: unknown,
|
||||
workspaceId: string,
|
||||
): Promise<UnifiedRouterResult> {
|
||||
this.logger.error(
|
||||
'Routing with planning failed, falling back to Helper agent:',
|
||||
error,
|
||||
);
|
||||
|
||||
const helperAgent = await this.getHelperAgent(workspaceId);
|
||||
|
||||
if (!helperAgent) {
|
||||
throw new Error('No helper agent available for fallback');
|
||||
}
|
||||
|
||||
return {
|
||||
strategy: 'simple',
|
||||
agent: helperAgent,
|
||||
};
|
||||
}
|
||||
|
||||
private async getAvailableAgents(
|
||||
workspaceId: string,
|
||||
): Promise<AgentEntity[]> {
|
||||
const agents = await this.agentRepository.find({
|
||||
where: { workspaceId, deletedAt: IsNull() },
|
||||
order: { createdAt: 'ASC' },
|
||||
});
|
||||
|
||||
return agents.filter(
|
||||
(agent) => !agent.name.includes('workflow-service-agent'),
|
||||
);
|
||||
}
|
||||
|
||||
private async getHelperAgent(workspaceId: string) {
|
||||
const helperAgent = await this.agentRepository.findOne({
|
||||
where: {
|
||||
workspaceId,
|
||||
standardId: HELPER_AGENT.standardId,
|
||||
},
|
||||
});
|
||||
|
||||
return helperAgent;
|
||||
}
|
||||
|
||||
private async buildWorkspaceObjectsList(
|
||||
workspaceId: string,
|
||||
): Promise<string> {
|
||||
try {
|
||||
const objects = await this.objectMetadataService.findManyWithinWorkspace(
|
||||
workspaceId,
|
||||
{
|
||||
where: { isActive: true, isSystem: false },
|
||||
},
|
||||
);
|
||||
|
||||
const filteredObjects = objects.filter(
|
||||
(obj) => !isWorkflowRelatedObject(obj),
|
||||
);
|
||||
|
||||
if (filteredObjects.length === 0) {
|
||||
return '';
|
||||
}
|
||||
|
||||
return filteredObjects
|
||||
.map((obj) => `- ${obj.labelSingular} (${obj.nameSingular})`)
|
||||
.join('\n');
|
||||
} catch (error) {
|
||||
this.logger.warn('Failed to build workspace objects list:', error);
|
||||
|
||||
return '';
|
||||
}
|
||||
}
|
||||
|
||||
private buildAgentDescriptions(
|
||||
agents: AgentEntity[],
|
||||
workspaceObjectsList: string,
|
||||
): string {
|
||||
const agentDescriptions = agents
|
||||
.map((agent) => {
|
||||
return `- ${agent.label} (${agent.name}): ${agent.description}`;
|
||||
})
|
||||
.join('\n');
|
||||
|
||||
if (workspaceObjectsList) {
|
||||
return `${agentDescriptions}
|
||||
|
||||
Available workspace objects for data-manipulator:
|
||||
${workspaceObjectsList}`;
|
||||
}
|
||||
|
||||
return agentDescriptions;
|
||||
}
|
||||
}
|
||||
-192
@@ -1,192 +0,0 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
|
||||
import {
|
||||
generateObject,
|
||||
type UIDataTypes,
|
||||
type UIMessage,
|
||||
type UITools,
|
||||
} from 'ai';
|
||||
import { z } from 'zod';
|
||||
|
||||
import { type AgentEntity } from 'src/engine/metadata-modules/ai/ai-agent/entities/agent.entity';
|
||||
import { type ExecutionPlan } from 'src/engine/metadata-modules/ai/ai-chat-router/types/router-result.interface';
|
||||
import {
|
||||
DEFAULT_SMART_MODEL,
|
||||
type ModelId,
|
||||
} from 'src/engine/metadata-modules/ai/ai-models/constants/ai-models.const';
|
||||
import { AI_TELEMETRY_CONFIG } from 'src/engine/metadata-modules/ai/ai-models/constants/ai-telemetry.const';
|
||||
import { AiModelRegistryService } from 'src/engine/metadata-modules/ai/ai-models/services/ai-model-registry.service';
|
||||
|
||||
@Injectable()
|
||||
export class AiChatRouterPlanGeneratorService {
|
||||
private readonly logger = new Logger(AiChatRouterPlanGeneratorService.name);
|
||||
|
||||
constructor(
|
||||
private readonly aiModelRegistryService: AiModelRegistryService,
|
||||
) {}
|
||||
|
||||
async generatePlan({
|
||||
messages,
|
||||
availableAgents,
|
||||
agentDescriptions,
|
||||
smartModel,
|
||||
}: {
|
||||
messages: UIMessage<unknown, UIDataTypes, UITools>[];
|
||||
availableAgents: AgentEntity[];
|
||||
agentDescriptions: string;
|
||||
smartModel: ModelId;
|
||||
}): Promise<ExecutionPlan> {
|
||||
const model = this.getSmartModel(smartModel);
|
||||
const agentNames = availableAgents.map((agent) => agent.name);
|
||||
|
||||
const conversationHistory = messages
|
||||
.slice(0, -1)
|
||||
.map((msg) => {
|
||||
const textContent =
|
||||
msg.parts.find((part) => part.type === 'text')?.text || '';
|
||||
|
||||
return `${msg.role}: ${textContent}`;
|
||||
})
|
||||
.join('\n');
|
||||
|
||||
const currentMessage =
|
||||
messages[messages.length - 1]?.parts.find((part) => part.type === 'text')
|
||||
?.text || '';
|
||||
|
||||
const systemPrompt = `You are an AI planner that creates execution plans for multi-agent tasks.
|
||||
|
||||
Available agents:
|
||||
${agentDescriptions}
|
||||
|
||||
Create a step-by-step execution plan. Each step should:
|
||||
- Assign to the most appropriate agent
|
||||
- Have a clear, specific task
|
||||
- Specify expected output
|
||||
- List dependencies on previous steps (if any)
|
||||
|
||||
Keep plans focused and efficient.`;
|
||||
|
||||
const userPrompt = `${conversationHistory ? `Conversation history:\n${conversationHistory}\n\n` : ''}Current request:\n${currentMessage}\n\nCreate a detailed execution plan with specific steps.`;
|
||||
|
||||
const planStepSchema = z.object({
|
||||
stepNumber: z.number().describe('Step number in execution order'),
|
||||
agentName: z
|
||||
.enum([agentNames[0], ...agentNames.slice(1)])
|
||||
.describe('Agent name to execute this step'),
|
||||
task: z.string().describe('Specific task for this agent'),
|
||||
expectedOutput: z.string().describe('Expected output from this step'),
|
||||
dependsOn: z
|
||||
.array(z.number())
|
||||
.optional()
|
||||
.describe('Step numbers this step depends on'),
|
||||
});
|
||||
|
||||
const planSchema = z.object({
|
||||
steps: z.array(planStepSchema).describe('Execution steps in order'),
|
||||
reasoning: z.string().describe('Why multi-agent planning is needed'),
|
||||
});
|
||||
|
||||
const PLANNER_TEMPERATURE = 0.1;
|
||||
|
||||
const result = await generateObject({
|
||||
model,
|
||||
system: systemPrompt,
|
||||
prompt: userPrompt,
|
||||
schema: planSchema,
|
||||
temperature: PLANNER_TEMPERATURE,
|
||||
experimental_telemetry: AI_TELEMETRY_CONFIG,
|
||||
});
|
||||
|
||||
this.logger.log(
|
||||
`[PLANNER] Generated plan with ${result.object.steps.length} steps`,
|
||||
);
|
||||
|
||||
this.validatePlan(result.object);
|
||||
|
||||
return result.object as ExecutionPlan;
|
||||
}
|
||||
|
||||
private validatePlan(plan: ExecutionPlan): void {
|
||||
const stepNumbers = new Set(plan.steps.map((s) => s.stepNumber));
|
||||
|
||||
for (const step of plan.steps) {
|
||||
this.validateStepDependencies(step, stepNumbers);
|
||||
}
|
||||
|
||||
this.logger.log(`[PLANNER] Plan validation passed`);
|
||||
}
|
||||
|
||||
private validateStepDependencies(
|
||||
step: { stepNumber: number; dependsOn?: number[] },
|
||||
validStepNumbers: Set<number>,
|
||||
): void {
|
||||
if (!step.dependsOn) {
|
||||
return;
|
||||
}
|
||||
|
||||
this.checkForSelfDependency(step);
|
||||
|
||||
for (const dependency of step.dependsOn) {
|
||||
this.validateDependencyExists(dependency, validStepNumbers);
|
||||
this.validateDependencyOrder(step.stepNumber, dependency);
|
||||
}
|
||||
}
|
||||
|
||||
private checkForSelfDependency(step: {
|
||||
stepNumber: number;
|
||||
dependsOn?: number[];
|
||||
}): void {
|
||||
if (step.dependsOn?.includes(step.stepNumber)) {
|
||||
throw new Error(`Step ${step.stepNumber} cannot depend on itself`);
|
||||
}
|
||||
}
|
||||
|
||||
private validateDependencyExists(
|
||||
dependency: number,
|
||||
validStepNumbers: Set<number>,
|
||||
): void {
|
||||
if (!validStepNumbers.has(dependency)) {
|
||||
throw new Error(`Invalid dependency: step ${dependency} not found`);
|
||||
}
|
||||
}
|
||||
|
||||
private validateDependencyOrder(
|
||||
currentStepNumber: number,
|
||||
dependency: number,
|
||||
): void {
|
||||
if (dependency >= currentStepNumber) {
|
||||
throw new Error(
|
||||
`Step ${currentStepNumber} depends on future step ${dependency}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
private getSmartModel(modelId: ModelId) {
|
||||
if (modelId === DEFAULT_SMART_MODEL) {
|
||||
return this.getDefaultSmartModel();
|
||||
}
|
||||
|
||||
return this.getSpecificSmartModel(modelId);
|
||||
}
|
||||
|
||||
private getDefaultSmartModel() {
|
||||
const registeredModel =
|
||||
this.aiModelRegistryService.getDefaultPerformanceModel();
|
||||
|
||||
if (!registeredModel) {
|
||||
throw new Error('No smart model available');
|
||||
}
|
||||
|
||||
return registeredModel.model;
|
||||
}
|
||||
|
||||
private getSpecificSmartModel(modelId: ModelId) {
|
||||
const registeredModel = this.aiModelRegistryService.getModel(modelId);
|
||||
|
||||
if (!registeredModel) {
|
||||
throw new Error(`Smart model "${modelId}" not available`);
|
||||
}
|
||||
|
||||
return registeredModel.model;
|
||||
}
|
||||
}
|
||||
-198
@@ -1,198 +0,0 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
|
||||
import {
|
||||
generateObject,
|
||||
type LanguageModel,
|
||||
type UIDataTypes,
|
||||
type UIMessage,
|
||||
type UITools,
|
||||
} from 'ai';
|
||||
import { z } from 'zod';
|
||||
|
||||
import { AGENT_SYSTEM_PROMPTS } from 'src/engine/metadata-modules/ai/ai-agent/constants/agent-system-prompts.const';
|
||||
import { type AgentEntity } from 'src/engine/metadata-modules/ai/ai-agent/entities/agent.entity';
|
||||
import {
|
||||
DEFAULT_FAST_MODEL,
|
||||
type ModelId,
|
||||
} from 'src/engine/metadata-modules/ai/ai-models/constants/ai-models.const';
|
||||
import { AI_TELEMETRY_CONFIG } from 'src/engine/metadata-modules/ai/ai-models/constants/ai-telemetry.const';
|
||||
import { AiModelRegistryService } from 'src/engine/metadata-modules/ai/ai-models/services/ai-model-registry.service';
|
||||
|
||||
export type StrategyDecision = {
|
||||
strategy: 'simple' | 'planned';
|
||||
agentName?: string;
|
||||
toolHints?: {
|
||||
relevantObjects?: string[];
|
||||
operations?: Array<'find' | 'create' | 'update' | 'delete'>;
|
||||
};
|
||||
};
|
||||
|
||||
@Injectable()
|
||||
export class AiChatRouterStrategyDeciderService {
|
||||
private readonly logger = new Logger(AiChatRouterStrategyDeciderService.name);
|
||||
|
||||
constructor(
|
||||
private readonly aiModelRegistryService: AiModelRegistryService,
|
||||
) {}
|
||||
|
||||
async decideStrategy({
|
||||
messages,
|
||||
availableAgents,
|
||||
agentDescriptions,
|
||||
fastModel,
|
||||
}: {
|
||||
messages: UIMessage<unknown, UIDataTypes, UITools>[];
|
||||
availableAgents: AgentEntity[];
|
||||
agentDescriptions: string;
|
||||
fastModel: ModelId;
|
||||
}): Promise<StrategyDecision> {
|
||||
if (availableAgents.length === 1) {
|
||||
return this.createSingleAgentDecision(availableAgents[0]);
|
||||
}
|
||||
|
||||
const model = this.getFastModel(fastModel);
|
||||
const agentNames = availableAgents.map((agent) => agent.name);
|
||||
const conversationHistory = this.buildConversationHistory(messages);
|
||||
const currentMessage = this.extractCurrentMessage(messages);
|
||||
|
||||
const systemPrompt = AGENT_SYSTEM_PROMPTS.ROUTER(agentDescriptions);
|
||||
const userPrompt = this.buildUserPrompt(
|
||||
conversationHistory,
|
||||
currentMessage,
|
||||
);
|
||||
|
||||
const strategySchema = this.buildStrategySchema(agentNames);
|
||||
const decision = await this.generateStrategyDecision(
|
||||
model,
|
||||
systemPrompt,
|
||||
userPrompt,
|
||||
strategySchema,
|
||||
);
|
||||
|
||||
this.logger.log(
|
||||
`[STRATEGY] Decision: ${JSON.stringify(decision, null, 2)}`,
|
||||
);
|
||||
|
||||
return decision;
|
||||
}
|
||||
|
||||
private createSingleAgentDecision(agent: AgentEntity): StrategyDecision {
|
||||
return {
|
||||
strategy: 'simple',
|
||||
agentName: agent.name,
|
||||
};
|
||||
}
|
||||
|
||||
private buildConversationHistory(
|
||||
messages: UIMessage<unknown, UIDataTypes, UITools>[],
|
||||
): string {
|
||||
return messages
|
||||
.slice(0, -1)
|
||||
.map((message) => {
|
||||
const textContent =
|
||||
message.parts.find((part) => part.type === 'text')?.text || '';
|
||||
|
||||
return `${message.role}: ${textContent}`;
|
||||
})
|
||||
.join('\n');
|
||||
}
|
||||
|
||||
private extractCurrentMessage(
|
||||
messages: UIMessage<unknown, UIDataTypes, UITools>[],
|
||||
): string {
|
||||
return (
|
||||
messages[messages.length - 1]?.parts.find((part) => part.type === 'text')
|
||||
?.text || ''
|
||||
);
|
||||
}
|
||||
|
||||
private buildStrategySchema(agentNames: string[]) {
|
||||
return z.object({
|
||||
strategy: z
|
||||
.enum(['simple', 'planned'])
|
||||
.describe(
|
||||
'Routing strategy: "simple" for single agent, "planned" for multi-agent coordination',
|
||||
),
|
||||
agentName: z
|
||||
.enum([agentNames[0], ...agentNames.slice(1)])
|
||||
.optional()
|
||||
.describe(
|
||||
'Agent name (REQUIRED if strategy is "simple", omit if "planned")',
|
||||
),
|
||||
toolHints: z
|
||||
.object({
|
||||
relevantObjects: z
|
||||
.array(z.string())
|
||||
.optional()
|
||||
.describe('Names of objects mentioned (e.g., "company", "person")'),
|
||||
operations: z
|
||||
.array(z.enum(['find', 'create', 'update', 'delete']))
|
||||
.optional()
|
||||
.describe('Required database operations'),
|
||||
})
|
||||
.optional()
|
||||
.describe('Tool hints for simple strategy (optional)'),
|
||||
});
|
||||
}
|
||||
|
||||
private async generateStrategyDecision(
|
||||
model: LanguageModel,
|
||||
systemPrompt: string,
|
||||
userPrompt: string,
|
||||
schema: z.ZodTypeAny,
|
||||
): Promise<StrategyDecision> {
|
||||
const ROUTER_TEMPERATURE = 0.1;
|
||||
|
||||
const result = await generateObject({
|
||||
model,
|
||||
system: systemPrompt,
|
||||
prompt: userPrompt,
|
||||
schema,
|
||||
temperature: ROUTER_TEMPERATURE,
|
||||
experimental_telemetry: AI_TELEMETRY_CONFIG,
|
||||
});
|
||||
|
||||
return result.object as StrategyDecision;
|
||||
}
|
||||
|
||||
private getFastModel(modelId: ModelId) {
|
||||
if (modelId === DEFAULT_FAST_MODEL) {
|
||||
return this.getDefaultFastModel();
|
||||
}
|
||||
|
||||
return this.getSpecificFastModel(modelId);
|
||||
}
|
||||
|
||||
private getDefaultFastModel() {
|
||||
const registeredModel = this.aiModelRegistryService.getDefaultSpeedModel();
|
||||
|
||||
if (!registeredModel) {
|
||||
throw new Error('No fast model available');
|
||||
}
|
||||
|
||||
return registeredModel.model;
|
||||
}
|
||||
|
||||
private getSpecificFastModel(modelId: ModelId) {
|
||||
const registeredModel = this.aiModelRegistryService.getModel(modelId);
|
||||
|
||||
if (!registeredModel) {
|
||||
throw new Error(`Fast model "${modelId}" not available`);
|
||||
}
|
||||
|
||||
return registeredModel.model;
|
||||
}
|
||||
|
||||
private buildUserPrompt(
|
||||
conversationHistory: string,
|
||||
currentMessage: string,
|
||||
): string {
|
||||
return `Conversation history:
|
||||
${conversationHistory || 'No previous conversation'}
|
||||
|
||||
Current user message:
|
||||
${currentMessage}
|
||||
|
||||
Which agent should handle this message?`;
|
||||
}
|
||||
}
|
||||
-39
@@ -1,39 +0,0 @@
|
||||
import { type AgentEntity } from 'src/engine/metadata-modules/ai/ai-agent/entities/agent.entity';
|
||||
|
||||
import { type ToolHints } from './tool-hints.interface';
|
||||
|
||||
export type PlanStep = {
|
||||
stepNumber: number;
|
||||
agentName: string;
|
||||
task: string;
|
||||
expectedOutput: string;
|
||||
dependsOn?: number[];
|
||||
};
|
||||
|
||||
export type ExecutionPlan = {
|
||||
steps: PlanStep[];
|
||||
reasoning: string;
|
||||
};
|
||||
|
||||
export type RouterDebugInfo = {
|
||||
availableAgents: Array<{ id: string; label: string }>;
|
||||
routerModel: string;
|
||||
promptTokens?: number;
|
||||
completionTokens?: number;
|
||||
totalTokens?: number;
|
||||
};
|
||||
|
||||
export type SimpleRouterResult = {
|
||||
strategy: 'simple';
|
||||
agent: AgentEntity;
|
||||
toolHints?: ToolHints;
|
||||
debugInfo?: RouterDebugInfo;
|
||||
};
|
||||
|
||||
export type PlannedRouterResult = {
|
||||
strategy: 'planned';
|
||||
plan: ExecutionPlan;
|
||||
debugInfo?: RouterDebugInfo;
|
||||
};
|
||||
|
||||
export type UnifiedRouterResult = SimpleRouterResult | PlannedRouterResult;
|
||||
-8
@@ -1,8 +0,0 @@
|
||||
export type ToolOperation = 'find' | 'create' | 'update' | 'delete';
|
||||
|
||||
export interface ToolHints {
|
||||
// Object names (singular or plural) that are relevant to the query
|
||||
relevantObjects?: string[];
|
||||
// Specific CRUD operations needed for the query
|
||||
operations?: ToolOperation[];
|
||||
}
|
||||
@@ -2,32 +2,30 @@ import { Module } from '@nestjs/common';
|
||||
import { TypeOrmModule } from '@nestjs/typeorm';
|
||||
|
||||
import { TokenModule } from 'src/engine/core-modules/auth/token/token.module';
|
||||
import { WorkspaceDomainsModule } from 'src/engine/core-modules/domain/workspace-domains/workspace-domains.module';
|
||||
import { FeatureFlagModule } from 'src/engine/core-modules/feature-flag/feature-flag.module';
|
||||
import { FileEntity } from 'src/engine/core-modules/file/entities/file.entity';
|
||||
import { FileUploadModule } from 'src/engine/core-modules/file/file-upload/file-upload.module';
|
||||
import { FileModule } from 'src/engine/core-modules/file/file.module';
|
||||
import { ThrottlerModule } from 'src/engine/core-modules/throttler/throttler.module';
|
||||
import { WORKFLOW_TOOL_SERVICE_TOKEN } from 'src/engine/core-modules/tool-provider/constants/workflow-tool-service.token';
|
||||
import { ToolProviderModule } from 'src/engine/core-modules/tool-provider/tool-provider.module';
|
||||
import { UserWorkspaceEntity } from 'src/engine/core-modules/user-workspace/user-workspace.entity';
|
||||
import { UserWorkspaceModule } from 'src/engine/core-modules/user-workspace/user-workspace.module';
|
||||
import { AiAgentExecutionModule } from 'src/engine/metadata-modules/ai/ai-agent-execution/ai-agent-execution.module';
|
||||
import { AiAgentModule } from 'src/engine/metadata-modules/ai/ai-agent/ai-agent.module';
|
||||
import { AiBillingModule } from 'src/engine/metadata-modules/ai/ai-billing/ai-billing.module';
|
||||
import { AiChatRouterModule } from 'src/engine/metadata-modules/ai/ai-chat-router/ai-chat-router.module';
|
||||
import { PermissionsModule } from 'src/engine/metadata-modules/permissions/permissions.module';
|
||||
import { TwentyORMModule } from 'src/engine/twenty-orm/twenty-orm.module';
|
||||
import { WorkspaceCacheStorageModule } from 'src/engine/workspace-cache-storage/workspace-cache-storage.module';
|
||||
import { WorkflowToolWorkspaceService } from 'src/modules/workflow/workflow-tools/services/workflow-tool.workspace-service';
|
||||
import { WorkflowToolsModule } from 'src/modules/workflow/workflow-tools/workflow-tools.module';
|
||||
import { WorkspaceCacheModule } from 'src/engine/workspace-cache/workspace-cache.module';
|
||||
|
||||
import { AgentChatController } from './controllers/agent-chat.controller';
|
||||
import { AgentChatThreadEntity } from './entities/agent-chat-thread.entity';
|
||||
import { AgentChatResolver } from './resolvers/agent-chat.resolver';
|
||||
import { AgentChatRoutingService } from './services/agent-chat-routing.service';
|
||||
import { AgentChatStreamingService } from './services/agent-chat-streaming.service';
|
||||
import { AgentChatService } from './services/agent-chat.service';
|
||||
import { AgentTitleGenerationService } from './services/agent-title-generation.service';
|
||||
import { ChatToolsProviderService } from './services/chat-tools-provider.service';
|
||||
import { ChatExecutionService } from './services/chat-execution.service';
|
||||
|
||||
@Module({
|
||||
imports: [
|
||||
@@ -38,33 +36,27 @@ import { ChatToolsProviderService } from './services/chat-tools-provider.service
|
||||
]),
|
||||
AiAgentModule,
|
||||
AiAgentExecutionModule,
|
||||
AiChatRouterModule,
|
||||
ThrottlerModule,
|
||||
FeatureFlagModule,
|
||||
FileUploadModule,
|
||||
FileModule,
|
||||
PermissionsModule,
|
||||
WorkspaceCacheStorageModule,
|
||||
WorkspaceCacheModule,
|
||||
WorkspaceDomainsModule,
|
||||
TwentyORMModule,
|
||||
TokenModule,
|
||||
UserWorkspaceModule,
|
||||
AiBillingModule,
|
||||
ToolProviderModule,
|
||||
// WorkflowToolsModule provides workflow tools for chat context
|
||||
WorkflowToolsModule,
|
||||
],
|
||||
controllers: [AgentChatController],
|
||||
providers: [
|
||||
AgentChatResolver,
|
||||
AgentChatService,
|
||||
AgentChatStreamingService,
|
||||
AgentChatRoutingService,
|
||||
AgentTitleGenerationService,
|
||||
ChatToolsProviderService,
|
||||
// Provide WorkflowToolWorkspaceService via token for ToolProviderService
|
||||
{
|
||||
provide: WORKFLOW_TOOL_SERVICE_TOKEN,
|
||||
useExisting: WorkflowToolWorkspaceService,
|
||||
},
|
||||
ChatExecutionService,
|
||||
],
|
||||
exports: [
|
||||
AgentChatService,
|
||||
|
||||
+34
@@ -0,0 +1,34 @@
|
||||
// System prompts for AI Chat (user-facing conversational interface)
|
||||
export const CHAT_SYSTEM_PROMPTS = {
|
||||
// Core chat behavior and tool strategy
|
||||
BASE: `You are a helpful AI assistant integrated into Twenty CRM.
|
||||
|
||||
Tool usage strategy:
|
||||
- Chain multiple tools to solve complex tasks
|
||||
- If a tool fails, try alternative approaches
|
||||
- Use results from one tool to inform the next
|
||||
- Don't give up after first failure - be persistent
|
||||
- Validate assumptions before making changes
|
||||
|
||||
Error recovery:
|
||||
- Analyze error messages to understand what went wrong
|
||||
- Adjust parameters or try different tools
|
||||
- Only give up after exhausting reasonable alternatives
|
||||
|
||||
Permissions:
|
||||
- Only perform actions your role allows
|
||||
- Explain limitations if you lack permissions`,
|
||||
|
||||
// Response formatting and record references
|
||||
RESPONSE_FORMAT: `
|
||||
Format responses with markdown for clarity (headings, lists, code blocks, tables).
|
||||
|
||||
Record References - IMPORTANT:
|
||||
- Tool responses include a "recordReferences" array with clickable links
|
||||
- ONLY use record references that are returned by tools - NEVER make up IDs
|
||||
- Copy the exact format from the tool response: [[record:objectName:recordId:displayName]]
|
||||
- The recordId MUST be a real UUID (like "abc12345-1234-5678-abcd-123456789012")
|
||||
- DO NOT create record references before calling the tool
|
||||
- DO NOT use placeholder IDs like "rec-snowflake" or "rec-person-1"
|
||||
- If a tool hasn't been called yet, don't reference records that don't exist`,
|
||||
};
|
||||
+3
-3
@@ -12,7 +12,7 @@ import { type ExtendedUIMessage } from 'twenty-shared/ai';
|
||||
import { PermissionFlagType } from 'twenty-shared/constants';
|
||||
|
||||
import { RestApiExceptionFilter } from 'src/engine/api/rest/rest-api-exception.filter';
|
||||
import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
|
||||
import { type WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
|
||||
import { AuthUserWorkspaceId } from 'src/engine/decorators/auth/auth-user-workspace-id.decorator';
|
||||
import { AuthWorkspace } from 'src/engine/decorators/auth/auth-workspace.decorator';
|
||||
import { JwtAuthGuard } from 'src/engine/guards/jwt-auth.guard';
|
||||
@@ -45,10 +45,10 @@ export class AgentChatController {
|
||||
this.agentStreamingService.streamAgentChat({
|
||||
threadId: body.threadId,
|
||||
messages: body.messages,
|
||||
recordIdsByObjectMetadataNameSingular:
|
||||
body.recordIdsByObjectMetadataNameSingular ?? [],
|
||||
userWorkspaceId,
|
||||
workspace,
|
||||
recordIdsByObjectMetadataNameSingular:
|
||||
body.recordIdsByObjectMetadataNameSingular || [],
|
||||
response,
|
||||
});
|
||||
}
|
||||
|
||||
-449
@@ -1,449 +0,0 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
|
||||
import { createUIMessageStream, pipeUIMessageStreamToResponse } from 'ai';
|
||||
import { type Response } from 'express';
|
||||
import { type ExtendedUIMessage } from 'twenty-shared/ai';
|
||||
|
||||
import { type WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
|
||||
import { AgentMessageRole } from 'src/engine/metadata-modules/ai/ai-agent-execution/entities/agent-message.entity';
|
||||
import { AgentActorContextService } from 'src/engine/metadata-modules/ai/ai-agent-execution/services/agent-actor-context.service';
|
||||
import { AgentExecutionService } from 'src/engine/metadata-modules/ai/ai-agent-execution/services/agent-execution.service';
|
||||
import { AgentPlanExecutorService } from 'src/engine/metadata-modules/ai/ai-agent-execution/services/agent-plan-executor.service';
|
||||
import { type RecordIdsByObjectMetadataNameSingularType } from 'src/engine/metadata-modules/ai/ai-agent/types/recordIdsByObjectMetadataNameSingular.type';
|
||||
import { AIBillingService } from 'src/engine/metadata-modules/ai/ai-billing/services/ai-billing.service';
|
||||
import { convertCentsToBillingCredits } from 'src/engine/metadata-modules/ai/ai-billing/utils/convert-cents-to-billing-credits.util';
|
||||
import { AiChatRouterService } from 'src/engine/metadata-modules/ai/ai-chat-router/ai-chat-router.service';
|
||||
import { type ModelId } from 'src/engine/metadata-modules/ai/ai-models/constants/ai-models.const';
|
||||
|
||||
import { ChatToolsProviderService } from './chat-tools-provider.service';
|
||||
|
||||
export type TokenUsage = {
|
||||
promptTokens: number;
|
||||
completionTokens: number;
|
||||
totalTokens: number;
|
||||
};
|
||||
|
||||
export type MessagePersistenceCallbacks = {
|
||||
saveSystemMessage?: (message: Omit<ExtendedUIMessage, 'id'>) => Promise<{
|
||||
turnId: string;
|
||||
}>;
|
||||
saveUserMessage: (
|
||||
message: Omit<ExtendedUIMessage, 'id'>,
|
||||
turnId?: string,
|
||||
) => Promise<{ turnId: string }>;
|
||||
saveAssistantMessage: (
|
||||
message: Omit<ExtendedUIMessage, 'id'>,
|
||||
turnId: string,
|
||||
agentId: string,
|
||||
) => Promise<void>;
|
||||
};
|
||||
|
||||
export type StreamAgentExecutionOptions = {
|
||||
userWorkspaceId: string;
|
||||
workspace: WorkspaceEntity;
|
||||
recordIdsByObjectMetadataNameSingular: RecordIdsByObjectMetadataNameSingularType;
|
||||
response: Response;
|
||||
messages: ExtendedUIMessage[];
|
||||
persistenceCallbacks: MessagePersistenceCallbacks;
|
||||
};
|
||||
|
||||
@Injectable()
|
||||
export class AgentChatRoutingService {
|
||||
private readonly logger = new Logger(AgentChatRoutingService.name);
|
||||
|
||||
constructor(
|
||||
private readonly agentExecutionService: AgentExecutionService,
|
||||
private readonly agentPlanExecutorService: AgentPlanExecutorService,
|
||||
private readonly aiChatRouterService: AiChatRouterService,
|
||||
private readonly aiBillingService: AIBillingService,
|
||||
private readonly chatToolsProviderService: ChatToolsProviderService,
|
||||
private readonly agentActorContextService: AgentActorContextService,
|
||||
) {}
|
||||
|
||||
async streamAgentExecution({
|
||||
userWorkspaceId,
|
||||
workspace,
|
||||
messages,
|
||||
recordIdsByObjectMetadataNameSingular,
|
||||
response,
|
||||
persistenceCallbacks,
|
||||
}: StreamAgentExecutionOptions) {
|
||||
try {
|
||||
const stream = createUIMessageStream<ExtendedUIMessage>({
|
||||
execute: async ({ writer }) => {
|
||||
const startTime = Date.now();
|
||||
|
||||
writer.write({
|
||||
type: 'data-routing-status' as const,
|
||||
id: 'routing-status',
|
||||
data: {
|
||||
text: 'Finding the best agent for your request...',
|
||||
state: 'loading',
|
||||
},
|
||||
});
|
||||
|
||||
const routingStart = Date.now();
|
||||
const includeDebugInfo = true;
|
||||
const routeResult = await this.aiChatRouterService.routeMessage(
|
||||
{
|
||||
messages,
|
||||
workspaceId: workspace.id,
|
||||
fastModel: workspace.fastModel,
|
||||
smartModel: workspace.smartModel,
|
||||
},
|
||||
includeDebugInfo,
|
||||
);
|
||||
|
||||
const routingTime = Date.now() - routingStart;
|
||||
|
||||
const { debugInfo } = routeResult;
|
||||
|
||||
let routingCostInCredits: number | undefined;
|
||||
|
||||
if (
|
||||
debugInfo?.routerModel &&
|
||||
debugInfo?.promptTokens !== undefined &&
|
||||
debugInfo?.completionTokens !== undefined
|
||||
) {
|
||||
try {
|
||||
const routingCostInCents =
|
||||
await this.aiBillingService.calculateCost(
|
||||
debugInfo.routerModel as ModelId,
|
||||
{
|
||||
inputTokens: debugInfo.promptTokens,
|
||||
outputTokens: debugInfo.completionTokens,
|
||||
totalTokens: debugInfo.totalTokens || 0,
|
||||
},
|
||||
);
|
||||
|
||||
routingCostInCredits = Math.round(
|
||||
convertCentsToBillingCredits(routingCostInCents),
|
||||
);
|
||||
} catch (error) {
|
||||
this.logger.warn('Failed to calculate routing cost:', error);
|
||||
}
|
||||
}
|
||||
|
||||
if (routeResult.strategy === 'planned') {
|
||||
this.logger.log(
|
||||
`Executing planned strategy with ${routeResult.plan.steps.length} steps`,
|
||||
);
|
||||
this.logger.log(
|
||||
`Plan steps: ${routeResult.plan.steps.map((s) => `${s.stepNumber}. ${s.agentName}: ${s.task}`).join('; ')}`,
|
||||
);
|
||||
|
||||
writer.write({
|
||||
type: 'data-routing-status' as const,
|
||||
id: 'routing-status',
|
||||
data: {
|
||||
text: `Executing ${routeResult.plan.steps.length}-step plan`,
|
||||
state: 'routed',
|
||||
debug: {
|
||||
routingTimeMs: routingTime,
|
||||
planReasoning: routeResult.plan.reasoning,
|
||||
totalSteps: routeResult.plan.steps.length,
|
||||
steps: routeResult.plan.steps.map((s) => ({
|
||||
stepNumber: s.stepNumber,
|
||||
agent: s.agentName,
|
||||
task: s.task,
|
||||
})),
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
const planResult = await this.agentPlanExecutorService.executePlan({
|
||||
steps: routeResult.plan.steps,
|
||||
reasoning: routeResult.plan.reasoning,
|
||||
workspace,
|
||||
userWorkspaceId,
|
||||
recordIdsByObjectMetadataNameSingular,
|
||||
writer,
|
||||
onProgress: (progress) => {
|
||||
if (progress.type === 'step-started') {
|
||||
this.logger.log(
|
||||
`Starting step ${progress.stepNumber}: ${progress.agentName} - ${progress.task}`,
|
||||
);
|
||||
writer.write({
|
||||
type: 'data-routing-status' as const,
|
||||
id: `step-${progress.stepNumber}`,
|
||||
data: {
|
||||
text: `Step ${progress.stepNumber}/${routeResult.plan.steps.length}: ${progress.agentName} → ${progress.task}`,
|
||||
state: 'loading',
|
||||
},
|
||||
});
|
||||
} else if (progress.type === 'step-completed') {
|
||||
this.logger.log(
|
||||
`Completed step ${progress.stepNumber}: ${progress.agentName}`,
|
||||
);
|
||||
writer.write({
|
||||
type: 'data-routing-status' as const,
|
||||
id: `step-${progress.stepNumber}`,
|
||||
data: {
|
||||
text: `Step ${progress.stepNumber}/${routeResult.plan.steps.length}: ✓ ${progress.agentName} completed`,
|
||||
state: 'routed',
|
||||
},
|
||||
});
|
||||
}
|
||||
},
|
||||
});
|
||||
|
||||
const systemMessage = messages.find((msg) => msg.role === 'system');
|
||||
let turnId: string | undefined;
|
||||
|
||||
if (systemMessage && persistenceCallbacks.saveSystemMessage) {
|
||||
const savedSystemMessage =
|
||||
await persistenceCallbacks.saveSystemMessage({
|
||||
role: AgentMessageRole.SYSTEM,
|
||||
parts: systemMessage.parts,
|
||||
});
|
||||
|
||||
turnId = savedSystemMessage.turnId;
|
||||
}
|
||||
|
||||
const userMessage = await persistenceCallbacks.saveUserMessage(
|
||||
{
|
||||
role: AgentMessageRole.USER,
|
||||
parts: [
|
||||
{
|
||||
type: 'text',
|
||||
text:
|
||||
messages[messages.length - 1].parts.find(
|
||||
(part) => part.type === 'text',
|
||||
)?.text ?? '',
|
||||
},
|
||||
],
|
||||
},
|
||||
turnId,
|
||||
);
|
||||
|
||||
await persistenceCallbacks.saveAssistantMessage(
|
||||
{
|
||||
role: AgentMessageRole.ASSISTANT,
|
||||
parts: [
|
||||
{
|
||||
type: 'text',
|
||||
text: planResult.finalOutput,
|
||||
},
|
||||
],
|
||||
},
|
||||
userMessage.turnId,
|
||||
'',
|
||||
);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
const { agent, toolHints } = routeResult;
|
||||
|
||||
this.logger.log(`Using agent ${agent.id} for message routing`);
|
||||
|
||||
const agentExecutionStart = Date.now();
|
||||
|
||||
// Get permission-based tools for chat context (workflow, metadata, etc.)
|
||||
// These tools are NOT available in workflow executor to prevent circular dependencies
|
||||
const { roleId } =
|
||||
await this.agentActorContextService.buildUserAndAgentActorContext(
|
||||
userWorkspaceId,
|
||||
workspace.id,
|
||||
);
|
||||
|
||||
const roleIds = [roleId];
|
||||
|
||||
const chatTools = await this.chatToolsProviderService.getChatTools(
|
||||
workspace.id,
|
||||
roleIds,
|
||||
toolHints,
|
||||
);
|
||||
|
||||
const {
|
||||
stream: result,
|
||||
timings,
|
||||
contextInfo,
|
||||
} = await this.agentExecutionService.streamChatResponse({
|
||||
workspace,
|
||||
agentId: agent.id,
|
||||
userWorkspaceId,
|
||||
messages,
|
||||
recordIdsByObjectMetadataNameSingular,
|
||||
toolHints,
|
||||
additionalTools: chatTools,
|
||||
});
|
||||
|
||||
const routedStatusPart = {
|
||||
type: 'data-routing-status' as const,
|
||||
id: 'routing-status',
|
||||
data: {
|
||||
text: `Routed to ${agent.label} agent`,
|
||||
state: 'routed',
|
||||
debug: {
|
||||
routingTimeMs: routingTime,
|
||||
contextBuildTimeMs: timings.contextBuildTimeMs,
|
||||
agentExecutionStartTimeMs: Date.now() - startTime,
|
||||
selectedAgentId: agent.id,
|
||||
selectedAgentLabel: agent.label,
|
||||
availableAgents: debugInfo?.availableAgents,
|
||||
routerModel: debugInfo?.routerModel,
|
||||
agentModel: agent.modelId,
|
||||
context: contextInfo.contextString || undefined,
|
||||
contextRecordCount: contextInfo.contextRecordCount,
|
||||
contextSizeBytes: contextInfo.contextSizeBytes,
|
||||
routingPromptTokens: debugInfo?.promptTokens,
|
||||
routingCompletionTokens: debugInfo?.completionTokens,
|
||||
routingTotalTokens: debugInfo?.totalTokens,
|
||||
routingCostInCredits,
|
||||
},
|
||||
},
|
||||
};
|
||||
|
||||
writer.write(routedStatusPart);
|
||||
|
||||
writer.merge(
|
||||
result.toUIMessageStream({
|
||||
onError: (error) => {
|
||||
return error instanceof Error ? error.message : String(error);
|
||||
},
|
||||
sendStart: false,
|
||||
onFinish: async ({ responseMessage }) => {
|
||||
if (responseMessage.parts.length === 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
const toolCallCount = responseMessage.parts.filter((part) =>
|
||||
part.type.startsWith('tool-'),
|
||||
).length;
|
||||
|
||||
const tokenUsage = await this.extractTokenUsage(result.usage);
|
||||
|
||||
const agentExecutionTime = Date.now() - agentExecutionStart;
|
||||
|
||||
let agentCostInCredits: number | undefined;
|
||||
let totalCostInCredits: number | undefined;
|
||||
|
||||
if (
|
||||
agent.modelId &&
|
||||
tokenUsage &&
|
||||
tokenUsage.promptTokens > 0 &&
|
||||
tokenUsage.completionTokens > 0
|
||||
) {
|
||||
try {
|
||||
const agentCostInCents =
|
||||
await this.aiBillingService.calculateCost(
|
||||
agent.modelId as ModelId,
|
||||
{
|
||||
inputTokens: tokenUsage.promptTokens,
|
||||
outputTokens: tokenUsage.completionTokens,
|
||||
totalTokens: tokenUsage.totalTokens,
|
||||
},
|
||||
);
|
||||
|
||||
agentCostInCredits = Math.round(
|
||||
convertCentsToBillingCredits(agentCostInCents),
|
||||
);
|
||||
|
||||
totalCostInCredits =
|
||||
(routingCostInCredits || 0) + agentCostInCredits;
|
||||
} catch (error) {
|
||||
this.logger.warn('Failed to calculate agent cost:', error);
|
||||
}
|
||||
}
|
||||
|
||||
const updatedRoutedStatusPart = {
|
||||
...routedStatusPart,
|
||||
data: {
|
||||
...routedStatusPart.data,
|
||||
debug: {
|
||||
...routedStatusPart.data.debug,
|
||||
agentExecutionTimeMs: agentExecutionTime,
|
||||
toolCallCount,
|
||||
toolCount: timings.toolCount,
|
||||
agentContextBuildTimeMs: timings.contextBuildTimeMs,
|
||||
toolGenerationTimeMs: timings.toolGenerationTimeMs,
|
||||
aiRequestPrepTimeMs: timings.aiRequestPrepTimeMs,
|
||||
...(tokenUsage && {
|
||||
agentPromptTokens: tokenUsage.promptTokens,
|
||||
agentCompletionTokens: tokenUsage.completionTokens,
|
||||
agentTotalTokens: tokenUsage.totalTokens,
|
||||
}),
|
||||
agentCostInCredits,
|
||||
totalCostInCredits,
|
||||
},
|
||||
},
|
||||
};
|
||||
|
||||
writer.write(updatedRoutedStatusPart);
|
||||
|
||||
const userMessage = await persistenceCallbacks.saveUserMessage(
|
||||
{
|
||||
role: AgentMessageRole.USER,
|
||||
parts: [
|
||||
{
|
||||
type: 'text',
|
||||
text:
|
||||
messages[messages.length - 1].parts.find(
|
||||
(part) => part.type === 'text',
|
||||
)?.text ?? '',
|
||||
},
|
||||
],
|
||||
},
|
||||
undefined,
|
||||
);
|
||||
|
||||
await persistenceCallbacks.saveAssistantMessage(
|
||||
{
|
||||
...responseMessage,
|
||||
parts: [updatedRoutedStatusPart, ...responseMessage.parts],
|
||||
},
|
||||
userMessage.turnId,
|
||||
agent.id,
|
||||
);
|
||||
},
|
||||
sendReasoning: true,
|
||||
}),
|
||||
);
|
||||
},
|
||||
});
|
||||
|
||||
pipeUIMessageStreamToResponse({ stream, response });
|
||||
} catch (error) {
|
||||
this.logger.error(
|
||||
'Failed to stream agent execution:',
|
||||
error instanceof Error ? error.message : String(error),
|
||||
);
|
||||
response.end();
|
||||
}
|
||||
}
|
||||
|
||||
private async extractTokenUsage(
|
||||
usagePromise: Promise<unknown>,
|
||||
): Promise<TokenUsage | null> {
|
||||
try {
|
||||
const usage = await usagePromise;
|
||||
|
||||
const usageWithTokens = usage as {
|
||||
inputTokens?: number;
|
||||
outputTokens?: number;
|
||||
promptTokens?: number;
|
||||
completionTokens?: number;
|
||||
totalTokens?: number;
|
||||
};
|
||||
|
||||
const tokenUsage = {
|
||||
promptTokens:
|
||||
usageWithTokens.inputTokens ?? usageWithTokens.promptTokens ?? 0,
|
||||
completionTokens:
|
||||
usageWithTokens.outputTokens ?? usageWithTokens.completionTokens ?? 0,
|
||||
totalTokens: usageWithTokens.totalTokens ?? 0,
|
||||
};
|
||||
|
||||
this.logger.log(
|
||||
`Agent execution usage: ${tokenUsage.promptTokens} prompt + ${tokenUsage.completionTokens} completion = ${tokenUsage.totalTokens} total tokens`,
|
||||
);
|
||||
|
||||
return tokenUsage;
|
||||
} catch (error) {
|
||||
this.logger.warn('Failed to get token usage:', error);
|
||||
|
||||
return null;
|
||||
}
|
||||
}
|
||||
}
|
||||
+104
-35
@@ -1,36 +1,41 @@
|
||||
import { Injectable } from '@nestjs/common';
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
import { InjectRepository } from '@nestjs/typeorm';
|
||||
|
||||
import { createUIMessageStream, pipeUIMessageStreamToResponse } from 'ai';
|
||||
import { type Response } from 'express';
|
||||
import { type ExtendedUIMessage } from 'twenty-shared/ai';
|
||||
import { type Repository } from 'typeorm';
|
||||
|
||||
import { type WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
|
||||
import { AgentMessageRole } from 'src/engine/metadata-modules/ai/ai-agent-execution/entities/agent-message.entity';
|
||||
import {
|
||||
AgentException,
|
||||
AgentExceptionCode,
|
||||
} from 'src/engine/metadata-modules/ai/ai-agent/agent.exception';
|
||||
import { type RecordIdsByObjectMetadataNameSingularType } from 'src/engine/metadata-modules/ai/ai-agent/types/recordIdsByObjectMetadataNameSingular.type';
|
||||
import { AgentChatThreadEntity } from 'src/engine/metadata-modules/ai/ai-chat/entities/agent-chat-thread.entity';
|
||||
import { AgentChatRoutingService } from 'src/engine/metadata-modules/ai/ai-chat/services/agent-chat-routing.service';
|
||||
import { AgentChatService } from 'src/engine/metadata-modules/ai/ai-chat/services/agent-chat.service';
|
||||
|
||||
import { AgentChatService } from './agent-chat.service';
|
||||
import { ChatExecutionService } from './chat-execution.service';
|
||||
|
||||
export type StreamAgentChatOptions = {
|
||||
threadId: string;
|
||||
userWorkspaceId: string;
|
||||
workspace: WorkspaceEntity;
|
||||
recordIdsByObjectMetadataNameSingular: RecordIdsByObjectMetadataNameSingularType;
|
||||
response: Response;
|
||||
messages: ExtendedUIMessage[];
|
||||
recordIdsByObjectMetadataNameSingular: RecordIdsByObjectMetadataNameSingularType;
|
||||
};
|
||||
|
||||
@Injectable()
|
||||
export class AgentChatStreamingService {
|
||||
private readonly logger = new Logger(AgentChatStreamingService.name);
|
||||
|
||||
constructor(
|
||||
@InjectRepository(AgentChatThreadEntity)
|
||||
private readonly threadRepository: Repository<AgentChatThreadEntity>,
|
||||
private readonly agentChatService: AgentChatService,
|
||||
private readonly agentChatRoutingService: AgentChatRoutingService,
|
||||
private readonly chatExecutionService: ChatExecutionService,
|
||||
) {}
|
||||
|
||||
async streamAgentChat({
|
||||
@@ -46,7 +51,6 @@ export class AgentChatStreamingService {
|
||||
id: threadId,
|
||||
userWorkspaceId,
|
||||
},
|
||||
relations: ['messages'],
|
||||
});
|
||||
|
||||
if (!thread) {
|
||||
@@ -56,39 +60,104 @@ export class AgentChatStreamingService {
|
||||
);
|
||||
}
|
||||
|
||||
await this.agentChatRoutingService.streamAgentExecution({
|
||||
userWorkspaceId,
|
||||
workspace,
|
||||
messages,
|
||||
recordIdsByObjectMetadataNameSingular,
|
||||
response,
|
||||
persistenceCallbacks: {
|
||||
saveSystemMessage: async (message) => {
|
||||
const savedMessage = await this.agentChatService.addMessage({
|
||||
threadId,
|
||||
uiMessage: message,
|
||||
try {
|
||||
const uiStream = createUIMessageStream<ExtendedUIMessage>({
|
||||
execute: async ({ writer }) => {
|
||||
const { stream } = await this.chatExecutionService.streamChat({
|
||||
workspace,
|
||||
userWorkspaceId,
|
||||
messages,
|
||||
recordIdsByObjectMetadataNameSingular,
|
||||
});
|
||||
|
||||
return { turnId: savedMessage.turnId };
|
||||
},
|
||||
saveUserMessage: async (message, turnId) => {
|
||||
const savedMessage = await this.agentChatService.addMessage({
|
||||
threadId,
|
||||
uiMessage: message,
|
||||
turnId,
|
||||
// Write initial status
|
||||
writer.write({
|
||||
type: 'data-routing-status' as const,
|
||||
id: 'execution-status',
|
||||
data: {
|
||||
text: 'Processing your request...',
|
||||
state: 'loading',
|
||||
},
|
||||
});
|
||||
|
||||
return { turnId: savedMessage.turnId };
|
||||
// Merge the AI stream
|
||||
writer.merge(
|
||||
stream.toUIMessageStream({
|
||||
onError: (error) => {
|
||||
this.logger.error('Stream error:', error);
|
||||
|
||||
return error instanceof Error ? error.message : String(error);
|
||||
},
|
||||
sendStart: false,
|
||||
onFinish: async ({ responseMessage }) => {
|
||||
if (responseMessage.parts.length === 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
// Update status to completed
|
||||
writer.write({
|
||||
type: 'data-routing-status' as const,
|
||||
id: 'execution-status',
|
||||
data: {
|
||||
text: 'Completed',
|
||||
state: 'routed',
|
||||
},
|
||||
});
|
||||
|
||||
// Save messages to database
|
||||
// Use thread.id from the validated thread object to ensure it's not null
|
||||
const validThreadId = thread.id;
|
||||
|
||||
if (!validThreadId) {
|
||||
this.logger.error('Thread ID is unexpectedly null/undefined');
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
const userMessage = await this.agentChatService.addMessage({
|
||||
threadId: validThreadId,
|
||||
uiMessage: {
|
||||
role: AgentMessageRole.USER,
|
||||
parts: [
|
||||
{
|
||||
type: 'text',
|
||||
text:
|
||||
messages[messages.length - 1].parts.find(
|
||||
(part) => part.type === 'text',
|
||||
)?.text ?? '',
|
||||
},
|
||||
],
|
||||
},
|
||||
});
|
||||
|
||||
await this.agentChatService.addMessage({
|
||||
threadId: validThreadId,
|
||||
uiMessage: responseMessage,
|
||||
turnId: userMessage.turnId,
|
||||
});
|
||||
} catch (saveError) {
|
||||
this.logger.error(
|
||||
'Failed to save messages:',
|
||||
saveError instanceof Error
|
||||
? saveError.message
|
||||
: String(saveError),
|
||||
);
|
||||
}
|
||||
},
|
||||
sendReasoning: true,
|
||||
}),
|
||||
);
|
||||
},
|
||||
saveAssistantMessage: async (message, turnId, agentId) => {
|
||||
await this.agentChatService.addMessage({
|
||||
threadId,
|
||||
uiMessage: message,
|
||||
agentId,
|
||||
turnId,
|
||||
});
|
||||
},
|
||||
},
|
||||
});
|
||||
});
|
||||
|
||||
pipeUIMessageStreamToResponse({ stream: uiStream, response });
|
||||
} catch (error) {
|
||||
this.logger.error(
|
||||
'Failed to stream chat:',
|
||||
error instanceof Error ? error.message : String(error),
|
||||
);
|
||||
response.end();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+453
@@ -0,0 +1,453 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
|
||||
import { anthropic } from '@ai-sdk/anthropic';
|
||||
import { openai } from '@ai-sdk/openai';
|
||||
import {
|
||||
convertToModelMessages,
|
||||
stepCountIs,
|
||||
streamText,
|
||||
type ToolSet,
|
||||
type UIDataTypes,
|
||||
type UIMessage,
|
||||
type UITools,
|
||||
} from 'ai';
|
||||
import { AppPath } from 'twenty-shared/types';
|
||||
import { getAppPath } from 'twenty-shared/utils';
|
||||
import { In } from 'typeorm';
|
||||
|
||||
import { getAllSelectableColumnNames } from 'src/engine/api/utils/get-all-selectable-column-names.utils';
|
||||
import { WorkspaceDomainsService } from 'src/engine/core-modules/domain/workspace-domains/services/workspace-domains.service';
|
||||
import {
|
||||
type ToolIndexEntry,
|
||||
ToolRegistryService,
|
||||
} from 'src/engine/core-modules/tool-provider/services/tool-registry.service';
|
||||
import {
|
||||
AGENT_SEARCH_TOOL_NAME,
|
||||
createAgentSearchTool,
|
||||
createLoadToolsTool,
|
||||
type DynamicToolStore,
|
||||
LOAD_TOOLS_TOOL_NAME,
|
||||
} from 'src/engine/core-modules/tool-provider/tools';
|
||||
import { type WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
|
||||
import { AgentActorContextService } from 'src/engine/metadata-modules/ai/ai-agent-execution/services/agent-actor-context.service';
|
||||
import {
|
||||
AgentException,
|
||||
AgentExceptionCode,
|
||||
} from 'src/engine/metadata-modules/ai/ai-agent/agent.exception';
|
||||
import { AgentService } from 'src/engine/metadata-modules/ai/ai-agent/agent.service';
|
||||
import { AGENT_CONFIG } from 'src/engine/metadata-modules/ai/ai-agent/constants/agent-config.const';
|
||||
import { type AgentEntity } from 'src/engine/metadata-modules/ai/ai-agent/entities/agent.entity';
|
||||
import { type RecordIdsByObjectMetadataNameSingularType } from 'src/engine/metadata-modules/ai/ai-agent/types/recordIdsByObjectMetadataNameSingular.type';
|
||||
import { repairToolCall } from 'src/engine/metadata-modules/ai/ai-agent/utils/repair-tool-call.util';
|
||||
import { AIBillingService } from 'src/engine/metadata-modules/ai/ai-billing/services/ai-billing.service';
|
||||
import { CHAT_SYSTEM_PROMPTS } from 'src/engine/metadata-modules/ai/ai-chat/constants/chat-system-prompts.const';
|
||||
import { ModelProvider } from 'src/engine/metadata-modules/ai/ai-models/constants/ai-models.const';
|
||||
import { AI_TELEMETRY_CONFIG } from 'src/engine/metadata-modules/ai/ai-models/constants/ai-telemetry.const';
|
||||
import { AiModelRegistryService } from 'src/engine/metadata-modules/ai/ai-models/services/ai-model-registry.service';
|
||||
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
|
||||
import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service';
|
||||
|
||||
export type ChatExecutionOptions = {
|
||||
workspace: WorkspaceEntity;
|
||||
userWorkspaceId: string;
|
||||
messages: UIMessage<unknown, UIDataTypes, UITools>[];
|
||||
recordIdsByObjectMetadataNameSingular: RecordIdsByObjectMetadataNameSingularType;
|
||||
};
|
||||
|
||||
export type ChatExecutionResult = {
|
||||
stream: ReturnType<typeof streamText>;
|
||||
preloadedTools: string[];
|
||||
initialAgents: string[];
|
||||
};
|
||||
|
||||
const INITIAL_AGENTS_LIMIT = 2;
|
||||
|
||||
// Common tools to pre-load for quick access
|
||||
const COMMON_PRELOAD_TOOLS = ['http_request', 'search_articles'];
|
||||
|
||||
@Injectable()
|
||||
export class ChatExecutionService {
|
||||
private readonly logger = new Logger(ChatExecutionService.name);
|
||||
|
||||
constructor(
|
||||
private readonly toolRegistry: ToolRegistryService,
|
||||
private readonly agentService: AgentService,
|
||||
private readonly aiModelRegistryService: AiModelRegistryService,
|
||||
private readonly aiBillingService: AIBillingService,
|
||||
private readonly agentActorContextService: AgentActorContextService,
|
||||
private readonly twentyORMGlobalManager: TwentyORMGlobalManager,
|
||||
private readonly workspaceCacheService: WorkspaceCacheService,
|
||||
private readonly workspaceDomainsService: WorkspaceDomainsService,
|
||||
) {}
|
||||
|
||||
async streamChat({
|
||||
workspace,
|
||||
userWorkspaceId,
|
||||
messages,
|
||||
recordIdsByObjectMetadataNameSingular,
|
||||
}: ChatExecutionOptions): Promise<ChatExecutionResult> {
|
||||
const { actorContext, roleId } =
|
||||
await this.agentActorContextService.buildUserAndAgentActorContext(
|
||||
userWorkspaceId,
|
||||
workspace.id,
|
||||
);
|
||||
|
||||
const toolContext = { workspaceId: workspace.id, roleId, actorContext };
|
||||
|
||||
const lastUserMessage = this.getLastUserMessage(messages);
|
||||
|
||||
let recordContext: string | undefined;
|
||||
|
||||
if (recordIdsByObjectMetadataNameSingular.length > 0) {
|
||||
recordContext = await this.buildContextFromRecords(
|
||||
workspace,
|
||||
recordIdsByObjectMetadataNameSingular,
|
||||
userWorkspaceId,
|
||||
);
|
||||
}
|
||||
|
||||
const [toolCatalog, initialAgents] = await Promise.all([
|
||||
this.toolRegistry.buildToolIndex(workspace.id, roleId),
|
||||
this.agentService.searchAgents(lastUserMessage, workspace.id, {
|
||||
limit: INITIAL_AGENTS_LIMIT,
|
||||
}),
|
||||
]);
|
||||
|
||||
this.logger.log(
|
||||
`Built tool catalog with ${toolCatalog.length} tools, ${initialAgents.length} agents`,
|
||||
);
|
||||
|
||||
const preloadedTools = await this.toolRegistry.getToolsByName(
|
||||
COMMON_PRELOAD_TOOLS,
|
||||
toolContext,
|
||||
);
|
||||
|
||||
const preloadedToolNames = Object.keys(preloadedTools);
|
||||
|
||||
const dynamicToolStore: DynamicToolStore = {
|
||||
loadedTools: new Set(preloadedToolNames),
|
||||
};
|
||||
|
||||
const registeredModel =
|
||||
this.aiModelRegistryService.getDefaultPerformanceModel();
|
||||
|
||||
const activeTools: ToolSet = {
|
||||
...preloadedTools,
|
||||
...this.getNativeWebSearchTool(registeredModel.provider),
|
||||
[LOAD_TOOLS_TOOL_NAME]: createLoadToolsTool(
|
||||
this.toolRegistry,
|
||||
toolContext,
|
||||
dynamicToolStore,
|
||||
async (toolNames) => {
|
||||
const newTools = await this.toolRegistry.getToolsByName(
|
||||
toolNames,
|
||||
toolContext,
|
||||
);
|
||||
|
||||
Object.assign(activeTools, newTools);
|
||||
this.logger.log(`Dynamically loaded tools: ${toolNames.join(', ')}`);
|
||||
},
|
||||
),
|
||||
[AGENT_SEARCH_TOOL_NAME]: createAgentSearchTool((query, options) =>
|
||||
this.agentService.searchAgents(query, workspace.id, options),
|
||||
),
|
||||
};
|
||||
|
||||
const systemPrompt = this.buildSystemPrompt(
|
||||
toolCatalog,
|
||||
initialAgents,
|
||||
preloadedToolNames,
|
||||
recordContext,
|
||||
);
|
||||
|
||||
this.logger.log(
|
||||
`Starting chat execution with model ${registeredModel.modelId}, ${Object.keys(activeTools).length} active tools`,
|
||||
);
|
||||
|
||||
const stream = streamText({
|
||||
model: registeredModel.model,
|
||||
system: systemPrompt,
|
||||
messages: convertToModelMessages(messages),
|
||||
tools: activeTools,
|
||||
stopWhen: stepCountIs(AGENT_CONFIG.MAX_STEPS),
|
||||
experimental_telemetry: AI_TELEMETRY_CONFIG,
|
||||
experimental_repairToolCall: async ({
|
||||
toolCall,
|
||||
tools: toolsForRepair,
|
||||
inputSchema,
|
||||
error,
|
||||
}) => {
|
||||
return repairToolCall({
|
||||
toolCall,
|
||||
tools: toolsForRepair,
|
||||
inputSchema,
|
||||
error,
|
||||
model: registeredModel.model,
|
||||
});
|
||||
},
|
||||
});
|
||||
|
||||
stream.usage
|
||||
.then((usage) => {
|
||||
this.aiBillingService.calculateAndBillUsage(
|
||||
registeredModel.modelId,
|
||||
usage,
|
||||
workspace.id,
|
||||
null,
|
||||
);
|
||||
})
|
||||
.catch((error) => {
|
||||
this.logger.error('Failed to bill usage:', error);
|
||||
});
|
||||
|
||||
return {
|
||||
stream,
|
||||
preloadedTools: preloadedToolNames,
|
||||
initialAgents: initialAgents.map((a) => a.name),
|
||||
};
|
||||
}
|
||||
|
||||
private async buildContextFromRecords(
|
||||
workspace: WorkspaceEntity,
|
||||
recordIdsByObjectMetadataNameSingular: RecordIdsByObjectMetadataNameSingularType,
|
||||
userWorkspaceId: string,
|
||||
): Promise<string> {
|
||||
const { userWorkspaceRoleMap } =
|
||||
await this.workspaceCacheService.getOrRecompute(workspace.id, [
|
||||
'userWorkspaceRoleMap',
|
||||
]);
|
||||
|
||||
const roleId = userWorkspaceRoleMap[userWorkspaceId];
|
||||
|
||||
if (!roleId) {
|
||||
throw new AgentException(
|
||||
'Failed to retrieve user role.',
|
||||
AgentExceptionCode.ROLE_NOT_FOUND,
|
||||
);
|
||||
}
|
||||
|
||||
const workspaceDataSource =
|
||||
await this.twentyORMGlobalManager.getDataSourceForWorkspace({
|
||||
workspaceId: workspace.id,
|
||||
});
|
||||
|
||||
const flatObjectMetadataMaps =
|
||||
workspaceDataSource.internalContext.flatObjectMetadataMaps;
|
||||
const flatFieldMetadataMaps =
|
||||
workspaceDataSource.internalContext.flatFieldMetadataMaps;
|
||||
const objectIdByNameSingular =
|
||||
workspaceDataSource.internalContext.objectIdByNameSingular;
|
||||
const objectMetadataPermissions = workspaceDataSource.permissionsPerRoleId;
|
||||
|
||||
const contextObject = (
|
||||
await Promise.all(
|
||||
recordIdsByObjectMetadataNameSingular.map(
|
||||
async (recordsWithObjectMetadataNameSingular) => {
|
||||
if (recordsWithObjectMetadataNameSingular.recordIds.length === 0) {
|
||||
return [];
|
||||
}
|
||||
|
||||
const objectMetadataId =
|
||||
objectIdByNameSingular[
|
||||
recordsWithObjectMetadataNameSingular.objectMetadataNameSingular
|
||||
];
|
||||
const objectMetadataMapItem = objectMetadataId
|
||||
? flatObjectMetadataMaps.byId[objectMetadataId]
|
||||
: undefined;
|
||||
|
||||
if (!objectMetadataMapItem) {
|
||||
this.logger.warn(
|
||||
`Object metadata not found for ${recordsWithObjectMetadataNameSingular.objectMetadataNameSingular}`,
|
||||
);
|
||||
|
||||
return [];
|
||||
}
|
||||
|
||||
const repository = workspaceDataSource.getRepository(
|
||||
recordsWithObjectMetadataNameSingular.objectMetadataNameSingular,
|
||||
{ unionOf: [roleId] },
|
||||
);
|
||||
|
||||
const restrictedFields =
|
||||
objectMetadataPermissions?.[roleId]?.[objectMetadataMapItem.id]
|
||||
?.restrictedFields ?? {};
|
||||
|
||||
const hasRestrictedFields = Object.values(restrictedFields).some(
|
||||
(field) => field.canRead === false,
|
||||
);
|
||||
|
||||
const selectOptions = hasRestrictedFields
|
||||
? getAllSelectableColumnNames({
|
||||
restrictedFields,
|
||||
objectMetadata: {
|
||||
objectMetadataMapItem,
|
||||
flatFieldMetadataMaps,
|
||||
},
|
||||
})
|
||||
: undefined;
|
||||
|
||||
return (
|
||||
await repository.find({
|
||||
...(selectOptions && { select: selectOptions }),
|
||||
where: {
|
||||
id: In(recordsWithObjectMetadataNameSingular.recordIds),
|
||||
},
|
||||
})
|
||||
).map((record) => {
|
||||
return {
|
||||
...record,
|
||||
resourceUrl: this.workspaceDomainsService.buildWorkspaceURL({
|
||||
workspace,
|
||||
pathname: getAppPath(AppPath.RecordShowPage, {
|
||||
objectNameSingular:
|
||||
recordsWithObjectMetadataNameSingular.objectMetadataNameSingular,
|
||||
objectRecordId: record.id,
|
||||
}),
|
||||
}),
|
||||
};
|
||||
});
|
||||
},
|
||||
),
|
||||
)
|
||||
).flat(2);
|
||||
|
||||
return JSON.stringify(contextObject);
|
||||
}
|
||||
|
||||
private getLastUserMessage(
|
||||
messages: UIMessage<unknown, UIDataTypes, UITools>[],
|
||||
): string {
|
||||
for (let i = messages.length - 1; i >= 0; i--) {
|
||||
const message = messages[i];
|
||||
|
||||
if (message.role === 'user') {
|
||||
const textPart = message.parts.find((part) => part.type === 'text');
|
||||
|
||||
if (textPart && 'text' in textPart) {
|
||||
return textPart.text;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return '';
|
||||
}
|
||||
|
||||
private buildSystemPrompt(
|
||||
toolCatalog: ToolIndexEntry[],
|
||||
agents: AgentEntity[],
|
||||
preloadedTools: string[],
|
||||
recordContext?: string,
|
||||
): string {
|
||||
const parts: string[] = [
|
||||
CHAT_SYSTEM_PROMPTS.BASE,
|
||||
CHAT_SYSTEM_PROMPTS.RESPONSE_FORMAT,
|
||||
];
|
||||
|
||||
if (agents.length > 0) {
|
||||
const skillsSection = agents
|
||||
.map((agent) => `## ${agent.label} Expertise\n${agent.prompt}`)
|
||||
.join('\n\n');
|
||||
|
||||
parts.push(`\nYou have the following expertise:\n\n${skillsSection}`);
|
||||
}
|
||||
|
||||
parts.push(this.buildToolCatalogSection(toolCatalog, preloadedTools));
|
||||
|
||||
if (recordContext) {
|
||||
parts.push(
|
||||
`\nCONTEXT (records the user is currently viewing):\n${recordContext}`,
|
||||
);
|
||||
}
|
||||
|
||||
return parts.join('\n');
|
||||
}
|
||||
|
||||
private buildToolCatalogSection(
|
||||
toolCatalog: ToolIndexEntry[],
|
||||
preloadedTools: string[],
|
||||
): string {
|
||||
const preloadedSet = new Set(preloadedTools);
|
||||
|
||||
const toolsByCategory = new Map<string, ToolIndexEntry[]>();
|
||||
|
||||
for (const tool of toolCatalog) {
|
||||
const category = tool.category;
|
||||
const existing = toolsByCategory.get(category) ?? [];
|
||||
|
||||
existing.push(tool);
|
||||
toolsByCategory.set(category, existing);
|
||||
}
|
||||
|
||||
const sections: string[] = [];
|
||||
|
||||
sections.push(`
|
||||
## Available Tools
|
||||
|
||||
You have access to ${toolCatalog.length} tools plus native web search. Some are pre-loaded and ready to use immediately.
|
||||
To use a tool that isn't pre-loaded, call \`${LOAD_TOOLS_TOOL_NAME}\` with the exact tool name(s) first.
|
||||
|
||||
### Pre-loaded Tools (ready to use now)
|
||||
- \`web_search\` ✓: Search the web for real-time information (ALWAYS use this for current data, news, research)
|
||||
${preloadedTools.length > 0 ? preloadedTools.map((t) => `- \`${t}\` ✓`).join('\n') : ''}
|
||||
|
||||
### Tool Catalog by Category`);
|
||||
|
||||
const categoryOrder = ['database', 'action', 'workflow', 'metadata'];
|
||||
|
||||
for (const category of categoryOrder) {
|
||||
const tools = toolsByCategory.get(category);
|
||||
|
||||
if (!tools || tools.length === 0) {
|
||||
continue;
|
||||
}
|
||||
|
||||
const categoryLabel = this.getCategoryLabel(category);
|
||||
|
||||
sections.push(`
|
||||
#### ${categoryLabel} (${tools.length} tools)
|
||||
${tools
|
||||
.map((t) => {
|
||||
const status = preloadedSet.has(t.name) ? ' ✓' : '';
|
||||
|
||||
return `- \`${t.name}\`${status}: ${t.description}`;
|
||||
})
|
||||
.join('\n')}`);
|
||||
}
|
||||
|
||||
sections.push(`
|
||||
### How to Use Tools
|
||||
1. **Web search** (\`web_search\`): Use for ANY request requiring current/real-time information from the internet
|
||||
2. **Pre-loaded tools** (marked with ✓): Use directly
|
||||
3. **Other tools**: First call \`${LOAD_TOOLS_TOOL_NAME}({toolNames: ["tool_name"]})\`, then use the tool
|
||||
4. **Agent expertise**: Call \`${AGENT_SEARCH_TOOL_NAME}\` to load specialized knowledge for workflows, etc.`);
|
||||
|
||||
return sections.join('\n');
|
||||
}
|
||||
|
||||
private getCategoryLabel(category: string): string {
|
||||
switch (category) {
|
||||
case 'database':
|
||||
return 'Database Tools (CRUD operations)';
|
||||
case 'action':
|
||||
return 'Action Tools (HTTP, Email, etc.)';
|
||||
case 'workflow':
|
||||
return 'Workflow Tools (create/manage workflows)';
|
||||
case 'metadata':
|
||||
return 'Metadata Tools (schema management)';
|
||||
default:
|
||||
return category;
|
||||
}
|
||||
}
|
||||
|
||||
private getNativeWebSearchTool(provider: ModelProvider): ToolSet {
|
||||
switch (provider) {
|
||||
case ModelProvider.ANTHROPIC:
|
||||
return { web_search: anthropic.tools.webSearch_20250305() };
|
||||
case ModelProvider.OPENAI:
|
||||
return { web_search: openai.tools.webSearch() };
|
||||
default:
|
||||
// Other providers don't have native web search
|
||||
return {};
|
||||
}
|
||||
}
|
||||
}
|
||||
-37
@@ -1,37 +0,0 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
|
||||
import { type ToolSet } from 'ai';
|
||||
|
||||
import { ToolCategory } from 'src/engine/core-modules/tool-provider/enums/tool-category.enum';
|
||||
import { ToolProviderService } from 'src/engine/core-modules/tool-provider/services/tool-provider.service';
|
||||
import { type ToolHints } from 'src/engine/metadata-modules/ai/ai-chat-router/types/tool-hints.interface';
|
||||
|
||||
@Injectable()
|
||||
export class ChatToolsProviderService {
|
||||
private readonly logger = new Logger(ChatToolsProviderService.name);
|
||||
|
||||
constructor(private readonly toolProvider: ToolProviderService) {}
|
||||
|
||||
// Provides additional tools for the chat context (WORKFLOW and METADATA)
|
||||
// These tools are NOT available in the workflow executor context to prevent circular dependencies
|
||||
// Base tools (DATABASE_CRUD, ACTION) are provided by AgentToolGeneratorService
|
||||
async getChatTools(
|
||||
workspaceId: string,
|
||||
roleIds: string[],
|
||||
toolHints?: ToolHints,
|
||||
): Promise<ToolSet> {
|
||||
const tools = await this.toolProvider.getTools({
|
||||
workspaceId,
|
||||
categories: [ToolCategory.WORKFLOW, ToolCategory.METADATA],
|
||||
rolePermissionConfig: { intersectionOf: roleIds },
|
||||
toolHints,
|
||||
wrapWithErrorContext: false,
|
||||
});
|
||||
|
||||
this.logger.log(
|
||||
`Generated ${Object.keys(tools).length} additional chat tools (workflow + metadata)`,
|
||||
);
|
||||
|
||||
return tools;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user