diff --git a/.env.example b/.env.example index ad448e2..65ddfe3 100644 --- a/.env.example +++ b/.env.example @@ -1 +1,8 @@ -OPENAI_API_KEY=sk-********************************************* \ No newline at end of file +# OpenAI API Key +OPENAI_API_KEY=sk-********************************************* + +# Telegram Bot Token (get it from @BotFather) +TELEGRAM_BOT_TOKEN=your_telegram_bot_token_here + +# MCP Server URL (optional) +MCP_SERVER_URL=https://santiment-mcp.dev.mind-dev.com/sse \ No newline at end of file diff --git a/README.md b/README.md index 40a15cd..0522e0e 100644 --- a/README.md +++ b/README.md @@ -1,15 +1,75 @@ -# ai-pipelines +# AI Pipelines with Telegram Bot -To install dependencies: +This project provides a framework for creating and executing AI-powered pipelines, with a Telegram bot interface for easy interaction. + +## Features + +- Execute AI-powered pipelines using OpenAI's models +- Integration with Model Context Protocol (MCP) tools +- Telegram bot interface for easy interaction +- Support for dynamic prompts and variables + +## Prerequisites + +- Node.js and Bun installed +- OpenAI API key +- Telegram Bot Token (get it from [@BotFather](https://t.me/botfather)) + +## Setup + +1. Clone the repository +2. Install dependencies: ```bash bun install ``` -To run: +3. Copy `.env.example` to `.env` and fill in your credentials: ```bash -bun run index.ts +cp .env.example .env ``` -This project was created using `bun init` in bun v1.2.14. [Bun](https://bun.sh) is a fast all-in-one JavaScript runtime. +4. Edit the `.env` file and add your: + - `OPENAI_API_KEY` - Your OpenAI API key + - `TELEGRAM_BOT_TOKEN` - Your Telegram bot token from @BotFather + - `MCP_SERVER_URL` - (Optional) MCP server URL (defaults to Santiment's MCP server) + +## Running the Bot + +Start the bot with: + +```bash +bun start +``` + +Or in development mode with auto-reload: + +```bash +bun --watch src/main.ts +``` + +## Using the Bot + +1. Find your bot on Telegram using the username you set with @BotFather +2. Start a chat with the bot +3. Send any message to the bot, and it will process it through the AI pipeline +4. The bot will respond with the AI's analysis or action + +## Project Structure + +- `src/bot/` - Telegram bot implementation +- `src/models/` - Data models and interfaces +- `src/pipelines/` - Pipeline execution logic +- `src/services/` - External service integrations +- `src/main.ts` - Application entry point + +## Development + +- Run linter: `bun run lint` +- Run type checking: `bun run ts:check` +- Run both: `bun run check` + +## License + +MIT diff --git a/bun.lock b/bun.lock index 4e3be2e..ab70dcb 100644 --- a/bun.lock +++ b/bun.lock @@ -11,6 +11,7 @@ "dotenv": "^16.5.0", "handlebars": "^4.7.8", "rimraf": "^6.0.1", + "telegraf": "^4.16.3", "yaml": "^2.8.0", }, "devDependencies": { @@ -55,6 +56,8 @@ "@modelcontextprotocol/sdk": ["@modelcontextprotocol/sdk@1.12.0", "", { "dependencies": { "ajv": "^6.12.6", "content-type": "^1.0.5", "cors": "^2.8.5", "cross-spawn": "^7.0.5", "eventsource": "^3.0.2", "express": "^5.0.1", "express-rate-limit": "^7.5.0", "pkce-challenge": "^5.0.0", "raw-body": "^3.0.0", "zod": "^3.23.8", "zod-to-json-schema": "^3.24.1" } }, "sha512-m//7RlINx1F3sz3KqwY1WWzVgTcYX52HYk4bJ1hkBXV3zccAEth+jRvG8DBRrdaQuRsPAJOx2MH3zaHNCKL7Zg=="], + "@telegraf/types": ["@telegraf/types@7.1.0", "", {}, "sha512-kGevOIbpMcIlCDeorKGpwZmdH7kHbqlk/Yj6dEpJMKEQw5lk0KVQY0OLXaCswy8GqlIVLd5625OB+rAntP9xVw=="], + "@types/bun": ["@types/bun@1.2.14", "", { "dependencies": { "bun-types": "1.2.14" } }, "sha512-VsFZKs8oKHzI7zwvECiAJ5oSorWndIWEVhfbYqZd4HI/45kzW7PN2Rr5biAzvGvRuNmYLSANY+H59ubHq8xw7Q=="], "@types/node": ["@types/node@22.15.21", "", { "dependencies": { "undici-types": "~6.21.0" } }, "sha512-EV/37Td6c+MgKAbkcLG6vqZ2zEYHD7bvSrzqqs2RIhbA6w3x+Dqz8MZM3sP6kGTeLrdoOgKZe+Xja7tUB2DNkQ=="], @@ -111,6 +114,12 @@ "brace-expansion": ["brace-expansion@1.1.11", "", { "dependencies": { "balanced-match": "^1.0.0", "concat-map": "0.0.1" } }, "sha512-iCuPHDFgrHX7H2vEI/5xpz07zSHB00TpugqhmYtVmMO6518mCuRMoOYFldEBl0g187ufozdaHgWKcYFb61qGiA=="], + "buffer-alloc": ["buffer-alloc@1.2.0", "", { "dependencies": { "buffer-alloc-unsafe": "^1.1.0", "buffer-fill": "^1.0.0" } }, "sha512-CFsHQgjtW1UChdXgbyJGtnm+O/uLQeZdtbDo8mfUgYXCHSM1wgrVxXm6bSyrUuErEb+4sYVGCzASBRot7zyrow=="], + + "buffer-alloc-unsafe": ["buffer-alloc-unsafe@1.1.0", "", {}, "sha512-TEM2iMIEQdJ2yjPJoSIsldnleVaAk1oW3DBVUykyOLsEsFmEc9kn+SFFPz+gl54KQNxlDnAwCXosOS9Okx2xAg=="], + + "buffer-fill": ["buffer-fill@1.0.0", "", {}, "sha512-T7zexNBwiiaCOGDg9xNX9PBmjrubblRkENuptryuI64URkXDFum9il/JGL8Lm8wYfAXpredVXXZz7eMHilimiQ=="], + "bun-types": ["bun-types@1.2.14", "", { "dependencies": { "@types/node": "*" } }, "sha512-Kuh4Ub28ucMRWeiUUWMHsT9Wcbr4H3kLIO72RZZElSDxSu7vpetRvxIUDUaW6QtaIeixIpm7OXtNnZPf82EzwA=="], "bytes": ["bytes@3.1.2", "", {}, "sha512-/Nf7TyzTx6S3yRJObOAV7956r8cr2+Oj8AC5dt8wSP3BQAoeX58NoHyCU8P8zGkNXStjTSi6fzO6F0pBdcYbEg=="], @@ -387,6 +396,8 @@ "minipass": ["minipass@7.1.2", "", {}, "sha512-qOOzS1cBTWYF4BH8fVePDBOO9iptMnGUEZwNc/cMWnTV2nVLZ7VoNWEPHkYczZA0pdoA7dl6e7FL659nX9S2aw=="], + "mri": ["mri@1.2.0", "", {}, "sha512-tzzskb3bG8LvYGFF/mDTpq3jpI6Q9wc3LEmBaghu+DdCssd1FakN7Bc0hVNmEyGq1bq3RgfkCb3cmQLpNPOroA=="], + "ms": ["ms@2.1.3", "", {}, "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA=="], "mustache": ["mustache@4.2.0", "", { "bin": { "mustache": "bin/mustache" } }, "sha512-71ippSywq5Yb7/tVYyGbkBggbU8H3u5Rz56fH60jGFgr8uHwxs+aSKeqmluIVzM0m0kB7xQjKS6qPfd0b2ZoqQ=="], @@ -429,7 +440,7 @@ "p-retry": ["p-retry@4.6.2", "", { "dependencies": { "@types/retry": "0.12.0", "retry": "^0.13.1" } }, "sha512-312Id396EbJdvRONlngUx0NydfrIQ5lsYu0znKVUzVvArzEIt08V1qhtyESbGVd1FGX7UKtiFp5uwKZdM8wIuQ=="], - "p-timeout": ["p-timeout@3.2.0", "", { "dependencies": { "p-finally": "^1.0.0" } }, "sha512-rhIwUycgwwKcP9yTOOFK/AKsAopjjCakVqLHePO3CC6Mir1Z99xT+R63jZxAT5lFZLa2inS5h+ZS2GvR99/FBg=="], + "p-timeout": ["p-timeout@4.1.0", "", {}, "sha512-+/wmHtzJuWii1sXn3HCuH/FTwGhrp4tmJTxSKJbfS+vkipci6osxXM5mY0jUiRzWKMTgUT8l7HFbeSwZAynqHw=="], "package-json-from-dist": ["package-json-from-dist@1.0.1", "", {}, "sha512-UEZIS3/by4OC8vL3P2dTXRETpebLI2NiI5vIrjaD/5UtrkFX/tNbwjTSRAGC/+7CAo2pIcBaRgWmcBBHcsaCIw=="], @@ -483,12 +494,16 @@ "safe-buffer": ["safe-buffer@5.2.1", "", {}, "sha512-rp3So07KcdmmKbGvgaNxQSJr7bGVSVk5S9Eq1F+ppbRo70+YeaDxkw5Dd8NPN+GD6bjnYm2VuPuCXmpuYvmCXQ=="], + "safe-compare": ["safe-compare@1.1.4", "", { "dependencies": { "buffer-alloc": "^1.2.0" } }, "sha512-b9wZ986HHCo/HbKrRpBJb2kqXMK9CEWIE1egeEvZsYn69ay3kdfl9nG3RyOcR+jInTDf7a86WQ1d4VJX7goSSQ=="], + "safe-push-apply": ["safe-push-apply@1.0.0", "", { "dependencies": { "es-errors": "^1.3.0", "isarray": "^2.0.5" } }, "sha512-iKE9w/Z7xCzUMIZqdBsp6pEQvwuEebH4vdpjcDWnyzaI6yl6O9FHvVpmGelvEHNsoY6wGblkxR6Zty/h00WiSA=="], "safe-regex-test": ["safe-regex-test@1.1.0", "", { "dependencies": { "call-bound": "^1.0.2", "es-errors": "^1.3.0", "is-regex": "^1.2.1" } }, "sha512-x/+Cz4YrimQxQccJf5mKEbIa1NzeCRNI5Ecl/ekmlYaampdNLPalVyIcCZNNH3MvmqBugV5TMYZXv0ljslUlaw=="], "safer-buffer": ["safer-buffer@2.1.2", "", {}, "sha512-YZo3K82SD7Riyi0E1EQPojLz7kpepnSQI9IyPbHHg1XXXevb5dJI7tpyN2ADxGcQbHG7vcyRHk0cbwqcQriUtg=="], + "sandwich-stream": ["sandwich-stream@2.0.2", "", {}, "sha512-jLYV0DORrzY3xaz/S9ydJL6Iz7essZeAfnAavsJ+zsJGZ1MOnsS52yRjU3uF3pJa/lla7+wisp//fxOwOH8SKQ=="], + "semver": ["semver@5.7.2", "", { "bin": { "semver": "bin/semver" } }, "sha512-cBznnQ9KjJqU67B52RMC65CMarK2600WFnbkcaiwWq3xy/5haFJlshgnpjovMVJ+Hff49d8GEn0b87C5pDQ10g=="], "send": ["send@1.2.0", "", { "dependencies": { "debug": "^4.3.5", "encodeurl": "^2.0.0", "escape-html": "^1.0.3", "etag": "^1.8.1", "fresh": "^2.0.0", "http-errors": "^2.0.0", "mime-types": "^3.0.1", "ms": "^2.1.3", "on-finished": "^2.4.1", "range-parser": "^1.2.1", "statuses": "^2.0.1" } }, "sha512-uaW0WwXKpL9blXE2o0bRhoL2EGXIrZxQ2ZQ4mgcfoBxdFmQold+qWsD2jLrfZ0trjKL6vOw0j//eAwcALFjKSw=="], @@ -555,6 +570,8 @@ "supports-preserve-symlinks-flag": ["supports-preserve-symlinks-flag@1.0.0", "", {}, "sha512-ot0WnXS9fgdkgIcePe6RHNk1WA8+muPa6cSjeR3V8K27q9BB1rTE3R1p7Hv0z1ZyAc8s6Vvv8DIyWf681MAt0w=="], + "telegraf": ["telegraf@4.16.3", "", { "dependencies": { "@telegraf/types": "^7.1.0", "abort-controller": "^3.0.0", "debug": "^4.3.4", "mri": "^1.2.0", "node-fetch": "^2.7.0", "p-timeout": "^4.1.0", "safe-compare": "^1.1.4", "sandwich-stream": "^2.0.2" }, "bin": { "telegraf": "lib/cli.mjs" } }, "sha512-yjEu2NwkHlXu0OARWoNhJlIjX09dRktiMQFsM678BAH/PEPVwctzL67+tvXqLCRQQvm3SDtki2saGO9hLlz68w=="], + "toidentifier": ["toidentifier@1.0.1", "", {}, "sha512-o5sSPKEkg/DIQNmH43V0/uerLrpzVedkUh8tGNvaeXpfpuwjKenlSox/2O/BTlZUtEe+JG7s5YhEz608PlAHRA=="], "tr46": ["tr46@0.0.3", "", {}, "sha512-N3WMsuqV66lT30CrXNbEjx4GEwlow3v6rr4mCcv6prnfwhS01rkgyFdjPNBYd9br7LpXV1+Emh01fHnq2Gdgrw=="], @@ -633,6 +650,8 @@ "openai/@types/node": ["@types/node@18.19.103", "", { "dependencies": { "undici-types": "~5.26.4" } }, "sha512-hHTHp+sEz6SxFsp+SA+Tqrua3AbmlAw+Y//aEwdHrdZkYVRWdvWD3y5uPZ0flYOkgskaFWqZ/YGFm3FaFQ0pRw=="], + "p-queue/p-timeout": ["p-timeout@3.2.0", "", { "dependencies": { "p-finally": "^1.0.0" } }, "sha512-rhIwUycgwwKcP9yTOOFK/AKsAopjjCakVqLHePO3CC6Mir1Z99xT+R63jZxAT5lFZLa2inS5h+ZS2GvR99/FBg=="], + "string-width-cjs/emoji-regex": ["emoji-regex@8.0.0", "", {}, "sha512-MSjYzcWNOA0ewAHpz0MxpYFvwg6yjy1NG3xteoqz644VCo/RPgnr1/GGt+ic3iJTzQ8Eu3TdM14SawnVUmGE6A=="], "string-width-cjs/strip-ansi": ["strip-ansi@6.0.1", "", { "dependencies": { "ansi-regex": "^5.0.1" } }, "sha512-Y38VPSHcqkFrCpFnQ9vuSXmquuv5oXOKpGeT6aGrr3o3Gc9AlVa6JBfUSOCnbxGGZF+/0ooI7KrPuUSztUdU5A=="], diff --git a/package.json b/package.json index fa73e36..0065041 100644 --- a/package.json +++ b/package.json @@ -38,6 +38,7 @@ "dotenv": "^16.5.0", "handlebars": "^4.7.8", "rimraf": "^6.0.1", + "telegraf": "^4.16.3", "yaml": "^2.8.0" } } diff --git a/src/bot/telegram-bot.ts b/src/bot/telegram-bot.ts new file mode 100644 index 0000000..694bd00 --- /dev/null +++ b/src/bot/telegram-bot.ts @@ -0,0 +1,103 @@ +import { Telegraf } from 'telegraf'; +import { Pipeline } from '../models/pipeline'; +import { PipelineExecutor } from '../pipelines/pipeline-executor'; + +export class TelegramBot { + private bot: Telegraf; + private pipelineExecutor: PipelineExecutor; + private activePipelines: Map = new Map(); + + constructor(token: string, mcpServerUrl: string) { + this.bot = new Telegraf(token); + this.pipelineExecutor = new PipelineExecutor(mcpServerUrl); + this.setupHandlers(); + } + + private setupHandlers(): void { + // Start command + this.bot.start((ctx) => { + ctx.reply( + 'Привет! Я бот для работы с AI пайплайнами. Просто отправь мне сообщение, и я обработаю его через пайплайн.', + ); + }); + + // Handle text messages + this.bot.on('text', async (ctx) => { + const chatId = ctx.chat.id; + const userMessage = ctx.message.text; + + try { + // Create a simple pipeline with the user's message as the first step + const userPipeline: Pipeline = { + description: 'User message processing pipeline', + systemPrompt: + 'Ты полезный ассистент, который всегда отвечает на русском языке. Анализируй сообщения пользователя и давай развернутые ответы.', + steps: [ + { + name: 'Start', + prompt: 'Всегда начинай своё сообщение со слов: Мой господин,', + }, + { + name: 'Обработка сообщения', + prompt: userMessage, + }, + ], + }; + + // Store the pipeline for this chat + this.activePipelines.set(chatId, userPipeline); + + // Show typing action + await ctx.sendChatAction('typing'); + + // Execute the pipeline + const results = await this.pipelineExecutor.executeChain({ + pipeline: userPipeline, + variables: {}, + systemPrompts: [], + openAIApiKey: process.env.OPENAI_API_KEY || '', + }); + + // Send the result back to the user + const lastResult = results[results.length - 1]; + if (lastResult) { + await ctx.reply(lastResult); + } else { + await ctx.reply( + 'Не удалось обработать запрос. Пожалуйста, попробуйте еще раз.', + ); + } + } catch (error) { + console.error('Error processing message:', error); + await ctx.reply('Произошла ошибка при обработке вашего запроса.'); + } + }); + + // Error handling + this.bot.catch((error: unknown) => { + const errorMessage = + error instanceof Error ? error.message : 'Unknown error'; + console.error('Telegram bot error:', errorMessage); + }); + } + + public launch(): void { + this.bot.launch(); + console.log('Telegram bot started'); + + // Enable graceful stop + process.once('SIGINT', () => { + this.bot.stop('SIGINT'); + this.cleanup(); + }); + process.once('SIGTERM', () => { + this.bot.stop('SIGTERM'); + this.cleanup(); + }); + } + + private async cleanup(): Promise { + // Cleanup resources + await this.pipelineExecutor.mcpToolsService.disconnect(); + } +} diff --git a/src/main.ts b/src/main.ts index 51540bd..4b71fad 100644 --- a/src/main.ts +++ b/src/main.ts @@ -1,13 +1,19 @@ import dotenv from 'dotenv'; -import { Pipeline } from './models/pipeline'; -import { PipelineExecutor } from './pipelines/pipeline-executor'; +import { TelegramBot } from './bot/telegram-bot'; dotenv.config(); async function main() { + const telegramToken = process.env.TELEGRAM_BOT_TOKEN; const openAIApiKey = process.env.OPENAI_API_KEY; - const mcpServerUrl = - process.env.MCP_SERVER_URL || 'https://santiment-mcp.dev.mind-dev.com/sse'; + const mcpServerUrl = process.env.MCP_SERVER_URL; + + if (!telegramToken) { + console.error( + 'Ошибка: Переменная окружения TELEGRAM_BOT_TOKEN не установлена.', + ); + process.exit(1); + } if (!openAIApiKey) { console.error( @@ -16,53 +22,21 @@ async function main() { 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(); + if (!mcpServerUrl) { + console.error( + 'Ошибка: Переменная окружения MCP_SERVER_URL не установлена.', + ); + process.exit(1); } + + // Initialize and start the Telegram bot + const bot = new TelegramBot(telegramToken, mcpServerUrl); + bot.launch(); + + console.log('Бот запущен и готов к работе!'); } -main().catch(console.error); +main().catch((error) => { + console.error('Произошла ошибка:', error); + process.exit(1); +});