Refactor: Rename entry point to main.ts and update start script.

This commit is contained in:
uxname committed 2025-05-24 16:32:50 +03:00
1 parent 40e71274dd
commit e9320e438f
7 files changed
+347 -352

No files matched your search

-347
View File
@@ -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<void> {
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<void> {
if (this.mcpClient) {
await this.mcpClient.close();
this.mcpClient = null;
this.transport = null;
console.log(`[MCP Service] Отключено от MCP сервера: ${this.serverUrl}`);
}
}
async getLangchainTools(): Promise<DynamicStructuredTool[]> {
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<string, z.ZodTypeAny> = {};
for (const key in mcpTool.inputSchema.properties) {
// biome-ignore lint/suspicious/noExplicitAny: <explanation>
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<string, unknown>) => {
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: <explanation>
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: <explanation>
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<DynamicStructuredTool[]> {
if (!this.cachedTools) {
await this.mcpToolsService.connect();
this.cachedTools = await this.mcpToolsService.getLangchainTools();
}
return this.cachedTools;
}
private async initChain(openAIApiKey: string): Promise<Runnable> {
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<string, string> = {},
): (HumanMessage | SystemMessage)[] {
return [
...systemPrompts.map(
(prompt) => new SystemMessage(Handlebars.compile(prompt)(variables)),
),
new SystemMessage(`Текущая дата и время: ${new Date().toISOString()}`),
];
}
private renderPrompt(
prompt: string,
variables: Record<string, string> = {},
): string {
return Handlebars.compile(prompt)(variables);
}
async executeChain(options: {
pipeline: Pipeline;
variables?: Record<string, string>;
systemPrompts?: string[];
openAIApiKey: string;
}): Promise<string[]> {
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() ?? '<empty>');
}
return results;
}
private async _resolveToolCalls(
chain: Runnable,
messages: (HumanMessage | SystemMessage | AIMessage | ToolMessage)[],
): Promise<AIMessage> {
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);
+68
View File
@@ -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);
+10
View File
@@ -0,0 +1,10 @@
export interface PipelineStep {
name: string;
prompt: string;
}
export interface Pipeline {
description: string;
systemPrompt: string;
steps: PipelineStep[];
}
+138
View File
@@ -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<DynamicStructuredTool[]> {
if (!this.cachedTools) {
await this.mcpToolsService.connect();
this.cachedTools = await this.mcpToolsService.getLangchainTools();
}
return this.cachedTools;
}
private async initChain(openAIApiKey: string): Promise<Runnable> {
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<string, string> = {},
): (HumanMessage | SystemMessage)[] {
return [
...systemPrompts.map(
(prompt) => new SystemMessage(Handlebars.compile(prompt)(variables)),
),
new SystemMessage(`Текущая дата и время: ${new Date().toISOString()}`),
];
}
private renderPrompt(
prompt: string,
variables: Record<string, string> = {},
): string {
return Handlebars.compile(prompt)(variables);
}
async executeChain(options: {
pipeline: Pipeline;
variables?: Record<string, string>;
systemPrompts?: string[];
openAIApiKey: string;
}): Promise<string[]> {
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() ?? '<empty>');
}
return results;
}
private async _resolveToolCalls(
chain: Runnable,
messages: (HumanMessage | SystemMessage | AIMessage | ToolMessage)[],
): Promise<AIMessage> {
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;
}
}
+129
View File
@@ -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<void> {
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<void> {
if (this.mcpClient) {
await this.mcpClient.close();
this.mcpClient = null;
this.transport = null;
console.log(`[MCP Service] Отключено от MCP сервера: ${this.serverUrl}`);
}
}
async getLangchainTools(): Promise<DynamicStructuredTool[]> {
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<string, z.ZodTypeAny> = {};
for (const key in mcpTool.inputSchema.properties) {
// biome-ignore lint/suspicious/noExplicitAny: <explanation>
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<string, unknown>) => {
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: <explanation>
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: <explanation>
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;
}
}