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 async sendMessage(chatId: number, message: string): Promise { try { await this.bot.telegram.sendMessage(chatId, message, { parse_mode: 'HTML', }); } catch (error) { console.error('Error sending Telegram message:', error); } } private setupStepCallback( _pipelineId: string, chatId: number, ): ( stepName: string, status: 'started' | 'completed', result?: string, ) => Promise { return async ( stepName: string, status: 'started' | 'completed', result?: string, ) => { const emoji = status === 'started' ? '🔄' : '✅'; const message = status === 'started' ? `${stepName} ${emoji}` : `${stepName} ${emoji}\n\n
${result?.substring(0, 3000) || 'Done'}
`; await this.sendMessage(chatId, message); }; } private setupHandlers(): void { this.bot.start((ctx) => { ctx.reply( 'Hi there! 👋 I’m your personal Aptos ecosystem assistant. Just share your wallet address, and I’ll give you a quick overview of your portfolio along with some tips for improving your distribution.', ); }); this.bot.on('text', async (ctx) => { const chatId = ctx.chat.id; const userMessage = ctx.message.text; console.log(`Received message from chat ${chatId}: ${userMessage}`); try { const userPipeline: Pipeline = { description: 'Crypto wallet analyzer for Liquidswap', systemPrompt: `You are a professional crypto analyst and portfolio manager. Analyze the cryptocurrency market using up-to-date web data. Conduct a portfolio analysis, adapting to either a broad (market trends, news) or limited (portfolio only) dataset. If portfolio data is missing, request it. Cover the following aspects if possible: - Technical analysis (trends, volatility, indicators). - Fundamental analysis (news, projects). - Risks (correlations, liquidity, regulatory factors). - Diversification and optimization. Provide recommendations for portfolio improvement. Response format: a report with sections "Current Market Situation", "Portfolio Analysis", "Risks", "Recommendations". Maintain a professional tone without complex jargon. If data is insufficient, ask clarifying questions. Ensure accuracy ≥ 95%, avoiding unsubstantiated assumptions.`, steps: [ { name: 'Request Processing', prompt: userMessage, }, { name: 'Analysis', prompt: `I want to properly allocate my cryptocurrency portfolio using data from the Liquidswap DeFi protocol. Please provide asset allocation recommendations based on data from the following tools: \`get_tokens\`, \`get_pools\`, \`get_pools_historical_aprs\` (last 5 days), \`get_pools_historical_tvls\` (last 5 days), \`get_balances_by_address\` (my Aptos address, request it if needed) and \`multi_tool_use.parallel\`. Then provide the available information about the specified address.`, }, { name: 'Response Formatting', prompt: `Great, now let's create a Telegram post from this data using emojis. First, briefly explain in simple terms what the user should do, then provide detailed information. The post should end with the text: "${process.env.PROMO_TEXT}"`, }, ], }; this.activePipelines.set(chatId, { pipeline: userPipeline, chatId }); await ctx.sendChatAction('typing'); await ctx.reply('🚀 Starting analysis...'); const results = await this.pipelineExecutor.executeChain({ pipeline: userPipeline, variables: {}, systemPrompts: [], openAIApiKey: process.env.OPENAI_API_KEY || '', onStepStart: async (_stepName) => { await ctx.sendChatAction('typing'); }, onStepComplete: async (_stepName, _result, isLastStep) => { if (!isLastStep) { await ctx.sendChatAction('typing'); } }, }); const lastResult = results[results.length - 1]; if (lastResult) { // remove ``` in start and ``` in end const trimmedResult = lastResult .replace(/^```html/g, '') .replace(/^```HTML/g, '') .replace(/```$/g, ''); console.log(`Sending response to chat ${chatId}:`, trimmedResult); await ctx.replyWithHTML(trimmedResult); } else { await ctx.reply('Failed to process the request. Please try again.'); } } catch (error) { console.error('Error processing message:', error); await ctx.reply('An error occurred while processing your request.'); } }); 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'); process.once('SIGINT', () => { this.bot.stop('SIGINT'); this.cleanup(); }); process.once('SIGTERM', () => { this.bot.stop('SIGTERM'); this.cleanup(); }); } private async cleanup(): Promise { await this.pipelineExecutor.mcpToolsService.disconnect(); } }