diff --git a/package.json b/package.json index e906894..fa73e36 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "ai-pipelines", - "module": "src/index.ts", + "module": "src/main.ts", "devDependencies": { "@biomejs/biome": "^1.9.4", "@types/bun": "latest", @@ -16,7 +16,7 @@ "build": "npx tsc", "build:go": "npx tsgo", "prebuild": "npx rimraf dist", - "start": "node dist/index.js", + "start": "node dist/main.js", "________________ FORMAT AND LINT ________________": "", "check": "npx run-p ts:check lint", "lint": "biome check", diff --git a/src/index.ts b/src/index.ts deleted file mode 100644 index 53d06d9..0000000 --- a/src/index.ts +++ /dev/null @@ -1,347 +0,0 @@ -import { - AIMessage, - HumanMessage, - SystemMessage, - ToolMessage, -} from '@langchain/core/messages'; -import { Runnable } from '@langchain/core/runnables'; -import { DynamicStructuredTool } from '@langchain/core/tools'; -import { ChatOpenAI } from '@langchain/openai'; -import dotenv from 'dotenv'; -import Handlebars from 'handlebars'; -import { z } from 'zod'; - -import { Client } from '@modelcontextprotocol/sdk/client/index.js'; -import { SSEClientTransport } from '@modelcontextprotocol/sdk/client/sse.js'; -import { Tool as McpTool, TextContent } from '@modelcontextprotocol/sdk/types'; - -dotenv.config(); - -interface PipelineStep { - name: string; - prompt: string; -} - -interface Pipeline { - description: string; - systemPrompt: string; - steps: PipelineStep[]; -} - -class McpToolsService { - private mcpClient: Client | null = null; - private transport: SSEClientTransport | null = null; - private readonly serverUrl: string; - - constructor(serverUrl: string) { - this.serverUrl = serverUrl; - } - - async connect(): Promise { - if (this.mcpClient && this.transport) { - return; - } - try { - this.transport = new SSEClientTransport(new URL(this.serverUrl)); - this.mcpClient = new Client({ - name: 'my-langchain-mcp-client', - version: '1.0.0', - }); - await this.mcpClient.connect(this.transport); - console.log(`[MCP Service] Подключено к MCP серверу: ${this.serverUrl}`); - } catch (error) { - console.error( - `[MCP Service] Ошибка подключения к MCP серверу ${this.serverUrl}:`, - error, - ); - throw error; - } - } - - async disconnect(): Promise { - if (this.mcpClient) { - await this.mcpClient.close(); - this.mcpClient = null; - this.transport = null; - console.log(`[MCP Service] Отключено от MCP сервера: ${this.serverUrl}`); - } - } - - async getLangchainTools(): Promise { - if (!this.mcpClient) { - throw new Error('MCP клиент не подключен. Сначала вызовите .connect()'); - } - - const mcpTools: McpTool[] = (await this.mcpClient.listTools()).tools; - const langchainTools: DynamicStructuredTool[] = []; - - for (const mcpTool of mcpTools) { - const properties: Record = {}; - for (const key in mcpTool.inputSchema.properties) { - // biome-ignore lint/suspicious/noExplicitAny: - const prop = mcpTool.inputSchema.properties[key] as any; - let schemaType: z.ZodTypeAny; - switch (prop.type) { - case 'string': - schemaType = z.string(); - break; - case 'number': - schemaType = z.number(); - break; - case 'boolean': - schemaType = z.boolean(); - break; - case 'array': - schemaType = z.array(z.any()); - break; - case 'object': - schemaType = z.object({}); - break; - default: - schemaType = z.any(); - } - if (!(mcpTool.inputSchema.required || []).includes(key)) { - schemaType = schemaType.optional(); - } - properties[key] = schemaType.describe(prop.description || ''); - } - const zodSchema = z.object(properties); - - const langchainTool = new DynamicStructuredTool({ - name: mcpTool.name, - description: mcpTool.description || '', - schema: zodSchema, - func: async (args: Record) => { - console.log( - `[MCP Tool Call] Вызов инструмента MCP: ${mcpTool.name} с аргументами:`, - args, - ); - const result = await this.mcpClient!.callTool({ - name: mcpTool.name, - arguments: args, - }); - - if (result.isError) { - // biome-ignore lint/suspicious/noExplicitAny: - const errorContent = (result.content as any[]) - .map((c) => (c as TextContent).text || '') - .join('\n'); - throw new Error( - `Ошибка выполнения инструмента ${mcpTool.name}: ${errorContent}`, - ); - } - - if (result.structuredContent) { - return result.structuredContent; - } - // biome-ignore lint/suspicious/noExplicitAny: - return (result.content as any[]) - .map((c) => { - if (c.type === 'text') return (c as TextContent).text; - if (c.type === 'image') return `[Изображение: ${c.mimeType}]`; - if (c.type === 'audio') return `[Аудио: ${c.mimeType}]`; - if (c.type === 'resource') return `[Ресурс: ${c.uri}]`; - return JSON.stringify(c); - }) - .join('\n'); - }, - }); - langchainTools.push(langchainTool); - } - return langchainTools; - } -} - -class PipelineExecutor { - readonly mcpToolsService: McpToolsService; - private cachedTools: DynamicStructuredTool[] | null = null; - - constructor(mcpServerUrl: string) { - this.mcpToolsService = new McpToolsService(mcpServerUrl); - } - - private async _getOrLoadTools(): Promise { - if (!this.cachedTools) { - await this.mcpToolsService.connect(); - this.cachedTools = await this.mcpToolsService.getLangchainTools(); - } - return this.cachedTools; - } - - private async initChain(openAIApiKey: string): Promise { - const tools = await this._getOrLoadTools(); - const llm = new ChatOpenAI({ - model: 'gpt-4o', - temperature: 0.2, - openAIApiKey: openAIApiKey, - }); - return llm.bindTools(tools); - } - - private buildMessages( - systemPrompts: string[], - variables: Record = {}, - ): (HumanMessage | SystemMessage)[] { - return [ - ...systemPrompts.map( - (prompt) => new SystemMessage(Handlebars.compile(prompt)(variables)), - ), - new SystemMessage(`Текущая дата и время: ${new Date().toISOString()}`), - ]; - } - - private renderPrompt( - prompt: string, - variables: Record = {}, - ): string { - return Handlebars.compile(prompt)(variables); - } - - async executeChain(options: { - pipeline: Pipeline; - variables?: Record; - systemPrompts?: string[]; - openAIApiKey: string; - }): Promise { - const variables = options.variables ?? {}; - const systemPrompts = options.systemPrompts ?? []; - const chain = await this.initChain(options.openAIApiKey); - // Инициализируем массив сообщений, который будет накапливаться - const messages: (HumanMessage | SystemMessage | AIMessage | ToolMessage)[] = - this.buildMessages( - [options.pipeline.systemPrompt, ...systemPrompts], - variables, - ); - const results: string[] = []; - - for (const { prompt } of options.pipeline.steps) { - const renderedPrompt = this.renderPrompt(prompt, variables); - // Добавляем текущий промпт пользователя в общий массив сообщений - messages.push(new HumanMessage(renderedPrompt)); - - const finalResponse = await this._resolveToolCalls( - chain, - messages, // Передаем накопительный массив сообщений - ); - // Добавляем ответ ИИ в общий массив сообщений для сохранения контекста - messages.push(finalResponse); - results.push(finalResponse?.content.toString() ?? ''); - } - return results; - } - - private async _resolveToolCalls( - chain: Runnable, - messages: (HumanMessage | SystemMessage | AIMessage | ToolMessage)[], - ): Promise { - let response: AIMessage = await chain.invoke(messages); - - while (response.tool_calls && response.tool_calls.length > 0) { - messages.push(response); - - for (const toolCall of response.tool_calls) { - try { - const tools = await this._getOrLoadTools(); - const tool = tools.find((t) => t.name === toolCall.name); - - if (!tool) { - console.error( - `Инструмент ${toolCall.name} не найден в списке LangChain инструментов.`, - ); - messages.push( - new ToolMessage({ - tool_call_id: toolCall.id!, - content: `Ошибка: Инструмент ${toolCall.name} не найден.`, - }), - ); - continue; - } - - const toolResult = await tool.func(toolCall.args); - messages.push( - new ToolMessage({ - tool_call_id: toolCall.id!, - content: JSON.stringify(toolResult), - }), - ); - } catch (error) { - console.error( - `Ошибка выполнения инструмента ${toolCall.name}:`, - error, - ); - messages.push( - new ToolMessage({ - tool_call_id: toolCall.id!, - content: `Ошибка: ${error.message}`, - }), - ); - } - } - response = await chain.invoke(messages); - } - return response; - } -} - -async function main() { - const openAIApiKey = process.env.OPENAI_API_KEY; - const mcpServerUrl = - process.env.MCP_SERVER_URL || 'https://santiment-mcp.dev.mind-dev.com/sse'; - - if (!openAIApiKey) { - console.error( - 'Ошибка: Переменная окружения OPENAI_API_KEY не установлена.', - ); - process.exit(1); - } - - const pipelineExecutor = new PipelineExecutor(mcpServerUrl); - - const testPipeline: Pipeline = { - description: 'Тестовый пайплайн', - systemPrompt: - 'Ты полезный ассистент, который всегда отвечает на русском языке.', - steps: [ - { - name: 'Приветствие', - prompt: 'Скажи привет пользователю, его имя - {{username}}.', - }, - { - name: 'Использование инструмента', - prompt: - 'Напиши какие инструменты тебе доступны? (Очень кратко напиши суть и входящие параметры).', - }, - { - name: 'Использование инструмента', - prompt: 'Вызови **top_gainers** с pageSize=2 и покажи что получилось.', - }, - { - name: 'Итог', - prompt: - 'Отлично, теперь давай подытожим всё что мы сделали, соберём информацию вместе и оформим её как telegram пост.', - }, - ], - }; - - const variables = { - username: 'Вася', - }; - - const systemPrompts = ['Отвечай кратко и по существу.']; - - try { - const results = await pipelineExecutor.executeChain({ - pipeline: testPipeline, - variables: variables, - systemPrompts: systemPrompts, - openAIApiKey: openAIApiKey, - }); - console.log('Результаты выполнения пайплайна:', results); - } catch (error) { - console.error('Ошибка при выполнении пайплайна:', error); - } finally { - await pipelineExecutor.mcpToolsService.disconnect(); - } -} - -main().catch(console.error); diff --git a/src/main.ts b/src/main.ts new file mode 100644 index 0000000..51540bd --- /dev/null +++ b/src/main.ts @@ -0,0 +1,68 @@ +import dotenv from 'dotenv'; +import { Pipeline } from './models/pipeline'; +import { PipelineExecutor } from './pipelines/pipeline-executor'; + +dotenv.config(); + +async function main() { + const openAIApiKey = process.env.OPENAI_API_KEY; + const mcpServerUrl = + process.env.MCP_SERVER_URL || 'https://santiment-mcp.dev.mind-dev.com/sse'; + + if (!openAIApiKey) { + console.error( + 'Ошибка: Переменная окружения OPENAI_API_KEY не установлена.', + ); + process.exit(1); + } + + const pipelineExecutor = new PipelineExecutor(mcpServerUrl); + + const testPipeline: Pipeline = { + description: 'Тестовый пайплайн', + systemPrompt: + 'Ты полезный ассистент, который всегда отвечает на русском языке.', + steps: [ + { + name: 'Приветствие', + prompt: 'Скажи привет пользователю, его имя - {{username}}.', + }, + { + name: 'Использование инструмента', + prompt: + 'Напиши какие инструменты тебе доступны? (Очень кратко напиши суть и входящие параметры).', + }, + { + name: 'Использование инструмента', + prompt: 'Вызови **top_gainers** с pageSize=2 и покажи что получилось.', + }, + { + name: 'Итог', + prompt: + 'Отлично, теперь давай подытожим всё что мы сделали, соберём информацию вместе и оформим её как telegram пост.', + }, + ], + }; + + const variables = { + username: 'Вася', + }; + + const systemPrompts = ['Отвечай кратко и по существу.']; + + try { + const results = await pipelineExecutor.executeChain({ + pipeline: testPipeline, + variables: variables, + systemPrompts: systemPrompts, + openAIApiKey: openAIApiKey, + }); + console.log('Результаты выполнения пайплайна:', results); + } catch (error) { + console.error('Ошибка при выполнении пайплайна:', error); + } finally { + await pipelineExecutor.mcpToolsService.disconnect(); + } +} + +main().catch(console.error); diff --git a/src/models/pipeline.ts b/src/models/pipeline.ts new file mode 100644 index 0000000..096c7ec --- /dev/null +++ b/src/models/pipeline.ts @@ -0,0 +1,10 @@ +export interface PipelineStep { + name: string; + prompt: string; +} + +export interface Pipeline { + description: string; + systemPrompt: string; + steps: PipelineStep[]; +} diff --git a/src/pipelines/pipeline-executor.ts b/src/pipelines/pipeline-executor.ts new file mode 100644 index 0000000..6064199 --- /dev/null +++ b/src/pipelines/pipeline-executor.ts @@ -0,0 +1,138 @@ +import { + AIMessage, + HumanMessage, + SystemMessage, + ToolMessage, +} from '@langchain/core/messages'; +import { Runnable } from '@langchain/core/runnables'; +import { DynamicStructuredTool } from '@langchain/core/tools'; +import { ChatOpenAI } from '@langchain/openai'; +import Handlebars from 'handlebars'; + +import { Pipeline } from '../models/pipeline'; +import { McpToolsService } from '../services/mcp-tools.service'; + +export class PipelineExecutor { + readonly mcpToolsService: McpToolsService; + private cachedTools: DynamicStructuredTool[] | null = null; + + constructor(mcpServerUrl: string) { + this.mcpToolsService = new McpToolsService(mcpServerUrl); + } + + private async _getOrLoadTools(): Promise { + if (!this.cachedTools) { + await this.mcpToolsService.connect(); + this.cachedTools = await this.mcpToolsService.getLangchainTools(); + } + return this.cachedTools; + } + + private async initChain(openAIApiKey: string): Promise { + const tools = await this._getOrLoadTools(); + const llm = new ChatOpenAI({ + model: 'gpt-4o', + temperature: 0.2, + openAIApiKey: openAIApiKey, + }); + return llm.bindTools(tools); + } + + private buildMessages( + systemPrompts: string[], + variables: Record = {}, + ): (HumanMessage | SystemMessage)[] { + return [ + ...systemPrompts.map( + (prompt) => new SystemMessage(Handlebars.compile(prompt)(variables)), + ), + new SystemMessage(`Текущая дата и время: ${new Date().toISOString()}`), + ]; + } + + private renderPrompt( + prompt: string, + variables: Record = {}, + ): string { + return Handlebars.compile(prompt)(variables); + } + + async executeChain(options: { + pipeline: Pipeline; + variables?: Record; + systemPrompts?: string[]; + openAIApiKey: string; + }): Promise { + const variables = options.variables ?? {}; + const systemPrompts = options.systemPrompts ?? []; + const chain = await this.initChain(options.openAIApiKey); + const messages: (HumanMessage | SystemMessage | AIMessage | ToolMessage)[] = + this.buildMessages( + [options.pipeline.systemPrompt, ...systemPrompts], + variables, + ); + const results: string[] = []; + + for (const { prompt } of options.pipeline.steps) { + const renderedPrompt = this.renderPrompt(prompt, variables); + messages.push(new HumanMessage(renderedPrompt)); + + const finalResponse = await this._resolveToolCalls(chain, messages); + messages.push(finalResponse); + results.push(finalResponse?.content.toString() ?? ''); + } + return results; + } + + private async _resolveToolCalls( + chain: Runnable, + messages: (HumanMessage | SystemMessage | AIMessage | ToolMessage)[], + ): Promise { + let response: AIMessage = await chain.invoke(messages); + + while (response.tool_calls && response.tool_calls.length > 0) { + messages.push(response); + + for (const toolCall of response.tool_calls) { + try { + const tools = await this._getOrLoadTools(); + const tool = tools.find((t) => t.name === toolCall.name); + + if (!tool) { + console.error( + `Инструмент ${toolCall.name} не найден в списке LangChain инструментов.`, + ); + messages.push( + new ToolMessage({ + tool_call_id: toolCall.id!, + content: `Ошибка: Инструмент ${toolCall.name} не найден.`, + }), + ); + continue; + } + + const toolResult = await tool.func(toolCall.args); + messages.push( + new ToolMessage({ + tool_call_id: toolCall.id!, + content: JSON.stringify(toolResult), + }), + ); + } catch (error) { + console.error( + `Ошибка выполнения инструмента ${toolCall.name}:`, + error, + ); + messages.push( + new ToolMessage({ + tool_call_id: toolCall.id!, + content: `Ошибка: ${error.message}`, + }), + ); + } + } + response = await chain.invoke(messages); + } + return response; + } +} diff --git a/src/services/mcp-tools.service.ts b/src/services/mcp-tools.service.ts new file mode 100644 index 0000000..f8cb8ea --- /dev/null +++ b/src/services/mcp-tools.service.ts @@ -0,0 +1,129 @@ +import { DynamicStructuredTool } from '@langchain/core/tools'; +import { Client } from '@modelcontextprotocol/sdk/client/index.js'; +import { SSEClientTransport } from '@modelcontextprotocol/sdk/client/sse.js'; +import { Tool as McpTool, TextContent } from '@modelcontextprotocol/sdk/types'; +import { z } from 'zod'; + +export class McpToolsService { + private mcpClient: Client | null = null; + private transport: SSEClientTransport | null = null; + private readonly serverUrl: string; + + constructor(serverUrl: string) { + this.serverUrl = serverUrl; + } + + async connect(): Promise { + if (this.mcpClient && this.transport) { + return; + } + try { + this.transport = new SSEClientTransport(new URL(this.serverUrl)); + this.mcpClient = new Client({ + name: 'my-langchain-mcp-client', + version: '1.0.0', + }); + await this.mcpClient.connect(this.transport); + console.log(`[MCP Service] Подключено к MCP серверу: ${this.serverUrl}`); + } catch (error) { + console.error( + `[MCP Service] Ошибка подключения к MCP серверу ${this.serverUrl}:`, + error, + ); + throw error; + } + } + + async disconnect(): Promise { + if (this.mcpClient) { + await this.mcpClient.close(); + this.mcpClient = null; + this.transport = null; + console.log(`[MCP Service] Отключено от MCP сервера: ${this.serverUrl}`); + } + } + + async getLangchainTools(): Promise { + if (!this.mcpClient) { + throw new Error('MCP клиент не подключен. Сначала вызовите .connect()'); + } + + const mcpTools: McpTool[] = (await this.mcpClient.listTools()).tools; + const langchainTools: DynamicStructuredTool[] = []; + + for (const mcpTool of mcpTools) { + const properties: Record = {}; + for (const key in mcpTool.inputSchema.properties) { + // biome-ignore lint/suspicious/noExplicitAny: + const prop = mcpTool.inputSchema.properties[key] as any; + let schemaType: z.ZodTypeAny; + switch (prop.type) { + case 'string': + schemaType = z.string(); + break; + case 'number': + schemaType = z.number(); + break; + case 'boolean': + schemaType = z.boolean(); + break; + case 'array': + schemaType = z.array(z.any()); + break; + case 'object': + schemaType = z.object({}); + break; + default: + schemaType = z.any(); + } + if (!(mcpTool.inputSchema.required || []).includes(key)) { + schemaType = schemaType.optional(); + } + properties[key] = schemaType.describe(prop.description || ''); + } + const zodSchema = z.object(properties); + + const langchainTool = new DynamicStructuredTool({ + name: mcpTool.name, + description: mcpTool.description || '', + schema: zodSchema, + func: async (args: Record) => { + console.log( + `[MCP Tool Call] Вызов инструмента MCP: ${mcpTool.name} с аргументами:`, + args, + ); + const result = await this.mcpClient!.callTool({ + name: mcpTool.name, + arguments: args, + }); + + if (result.isError) { + // biome-ignore lint/suspicious/noExplicitAny: + const errorContent = (result.content as any[]) + .map((c) => (c as TextContent).text || '') + .join('\n'); + throw new Error( + `Ошибка выполнения инструмента ${mcpTool.name}: ${errorContent}`, + ); + } + + if (result.structuredContent) { + return result.structuredContent; + } + // biome-ignore lint/suspicious/noExplicitAny: + return (result.content as any[]) + .map((c) => { + if (c.type === 'text') return (c as TextContent).text; + if (c.type === 'image') return `[Изображение: ${c.mimeType}]`; + if (c.type === 'audio') return `[Аудио: ${c.mimeType}]`; + if (c.type === 'resource') return `[Ресурс: ${c.uri}]`; + return JSON.stringify(c); + }) + .join('\n'); + }, + }); + langchainTools.push(langchainTool); + } + return langchainTools; + } +} diff --git a/tsconfig.json b/tsconfig.json index a138ff4..f905176 100644 --- a/tsconfig.json +++ b/tsconfig.json @@ -20,9 +20,6 @@ "target": "ESNext", "sourceMap": true, "baseUrl": "./", - "paths": { - "@/*": ["src/*"] - }, "outDir": "dist", "noImplicitThis": true, "incremental": true,