From 1c585251a70816937912cfca2c127fd4aee1d8ec Mon Sep 17 00:00:00 2001 From: uxname Date: Sat, 24 May 2025 23:39:23 +0300 Subject: [PATCH] feat: Add Telegram bot progress updates with step callbacks --- src/bot/telegram-bot.ts | 62 +++++++++++++++++++++++++++++- src/pipelines/pipeline-executor.ts | 44 ++++++++++++++++++--- 2 files changed, 98 insertions(+), 8 deletions(-) diff --git a/src/bot/telegram-bot.ts b/src/bot/telegram-bot.ts index 9a05e0b..38da6ea 100644 --- a/src/bot/telegram-bot.ts +++ b/src/bot/telegram-bot.ts @@ -5,7 +5,16 @@ import { PipelineExecutor } from '../pipelines/pipeline-executor'; export class TelegramBot { private bot: Telegraf; private pipelineExecutor: PipelineExecutor; - private activePipelines: Map = new Map(); + private activePipelines: Map = + new Map(); + private stepCallbacks: Map< + string, + ( + stepName: string, + status: 'started' | 'completed', + result?: string, + ) => Promise + > = new Map(); constructor(token: string, mcpServerUrl: string) { this.bot = new Telegraf(token); @@ -13,6 +22,39 @@ export class TelegramBot { 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) || 'Готово'}
`; + + await this.sendMessage(chatId, message); + }; + } + private setupHandlers(): void { this.bot.start((ctx) => { ctx.reply( @@ -67,17 +109,33 @@ export class TelegramBot { }, ], }; - this.activePipelines.set(chatId, userPipeline); + const pipelineId = `pipeline_${Date.now()}_${chatId}`; + const stepCallback = this.setupStepCallback(pipelineId, chatId); + this.stepCallbacks.set(pipelineId, stepCallback); + this.activePipelines.set(chatId, { pipeline: userPipeline, chatId }); await ctx.sendChatAction('typing'); + await ctx.reply('🚀 Запускаю анализ...'); const results = await this.pipelineExecutor.executeChain({ pipeline: userPipeline, variables: {}, systemPrompts: [], openAIApiKey: process.env.OPENAI_API_KEY || '', + onStepStart: async (stepName) => { + await stepCallback(stepName, 'started'); + await ctx.sendChatAction('typing'); + }, + onStepComplete: async (stepName, result, isLastStep) => { + // Skip the completion notification for the last step as it will be sent as the final response + if (!isLastStep) { + await stepCallback(stepName, 'completed', result); + } + }, }); + this.stepCallbacks.delete(pipelineId); + const lastResult = results[results.length - 1]; if (lastResult) { // remove ``` in start and ``` in end diff --git a/src/pipelines/pipeline-executor.ts b/src/pipelines/pipeline-executor.ts index e99bc4d..59a037c 100644 --- a/src/pipelines/pipeline-executor.ts +++ b/src/pipelines/pipeline-executor.ts @@ -76,6 +76,12 @@ export class PipelineExecutor { variables?: Record; systemPrompts?: string[]; openAIApiKey: string; + onStepStart?: (stepName: string, isLastStep: boolean) => Promise; + onStepComplete?: ( + stepName: string, + result: string, + isLastStep: boolean, + ) => Promise; }): Promise { const variables = options.variables ?? {}; const systemPrompts = options.systemPrompts ?? []; @@ -87,13 +93,39 @@ export class PipelineExecutor { ); const results: string[] = []; - for (const { prompt } of options.pipeline.steps) { - const renderedPrompt = this.renderPrompt(prompt, variables); - messages.push(new HumanMessage(renderedPrompt)); + const steps = options.pipeline.steps; + const totalSteps = steps.length; - const finalResponse = await this._resolveToolCalls(chain, messages); - messages.push(finalResponse); - results.push(finalResponse?.content.toString() ?? ''); + for (let i = 0; i < totalSteps; i++) { + const { name, prompt } = steps[i]!; + const stepName = name || 'Без названия'; + const isLastStep = i === totalSteps - 1; + + try { + if (options.onStepStart) { + await options.onStepStart(stepName, isLastStep); + } + + const renderedPrompt = this.renderPrompt(prompt, variables); + messages.push(new HumanMessage(renderedPrompt)); + + const finalResponse = await this._resolveToolCalls(chain, messages); + const result = finalResponse?.content?.toString() || ''; + + messages.push(finalResponse); + results.push(result); + + if (options.onStepComplete) { + await options.onStepComplete(stepName, result, isLastStep); + } + } catch (error) { + console.error(`Error in step "${stepName}":`, error); + const errorMessage = `❌ Ошибка на шаге "${stepName}": ${error instanceof Error ? error.message : String(error)}`; + if (options.onStepComplete) { + await options.onStepComplete(stepName, errorMessage, isLastStep); + } + throw error; + } } return results; }