feat: Add Telegram bot progress updates with step callbacks
This commit is contained in:
1 parent
f974dc8168
commit
1c585251a7
2 files changed
+94
-4
No files matched your search
+60
-2
@@ -5,7 +5,16 @@ import { PipelineExecutor } from '../pipelines/pipeline-executor';
|
|||||||
export class TelegramBot {
|
export class TelegramBot {
|
||||||
private bot: Telegraf;
|
private bot: Telegraf;
|
||||||
private pipelineExecutor: PipelineExecutor;
|
private pipelineExecutor: PipelineExecutor;
|
||||||
private activePipelines: Map<number, Pipeline> = new Map();
|
private activePipelines: Map<number, { pipeline: Pipeline; chatId: number }> =
|
||||||
|
new Map();
|
||||||
|
private stepCallbacks: Map<
|
||||||
|
string,
|
||||||
|
(
|
||||||
|
stepName: string,
|
||||||
|
status: 'started' | 'completed',
|
||||||
|
result?: string,
|
||||||
|
) => Promise<void>
|
||||||
|
> = new Map();
|
||||||
|
|
||||||
constructor(token: string, mcpServerUrl: string) {
|
constructor(token: string, mcpServerUrl: string) {
|
||||||
this.bot = new Telegraf(token);
|
this.bot = new Telegraf(token);
|
||||||
@@ -13,6 +22,39 @@ export class TelegramBot {
|
|||||||
this.setupHandlers();
|
this.setupHandlers();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private async sendMessage(chatId: number, message: string): Promise<void> {
|
||||||
|
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<void> {
|
||||||
|
return async (
|
||||||
|
stepName: string,
|
||||||
|
status: 'started' | 'completed',
|
||||||
|
result?: string,
|
||||||
|
) => {
|
||||||
|
const emoji = status === 'started' ? '🔄' : '✅';
|
||||||
|
const message =
|
||||||
|
status === 'started'
|
||||||
|
? `<b>${stepName}</b> ${emoji}`
|
||||||
|
: `<b>${stepName}</b> ${emoji}\n\n<pre>${result?.substring(0, 3000) || 'Готово'}</pre>`;
|
||||||
|
|
||||||
|
await this.sendMessage(chatId, message);
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
private setupHandlers(): void {
|
private setupHandlers(): void {
|
||||||
this.bot.start((ctx) => {
|
this.bot.start((ctx) => {
|
||||||
ctx.reply(
|
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.sendChatAction('typing');
|
||||||
|
await ctx.reply('🚀 Запускаю анализ...');
|
||||||
|
|
||||||
const results = await this.pipelineExecutor.executeChain({
|
const results = await this.pipelineExecutor.executeChain({
|
||||||
pipeline: userPipeline,
|
pipeline: userPipeline,
|
||||||
variables: {},
|
variables: {},
|
||||||
systemPrompts: [],
|
systemPrompts: [],
|
||||||
openAIApiKey: process.env.OPENAI_API_KEY || '',
|
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];
|
const lastResult = results[results.length - 1];
|
||||||
if (lastResult) {
|
if (lastResult) {
|
||||||
// remove ``` in start and ``` in end
|
// remove ``` in start and ``` in end
|
||||||
|
|||||||
@@ -76,6 +76,12 @@ export class PipelineExecutor {
|
|||||||
variables?: Record<string, string>;
|
variables?: Record<string, string>;
|
||||||
systemPrompts?: string[];
|
systemPrompts?: string[];
|
||||||
openAIApiKey: string;
|
openAIApiKey: string;
|
||||||
|
onStepStart?: (stepName: string, isLastStep: boolean) => Promise<void>;
|
||||||
|
onStepComplete?: (
|
||||||
|
stepName: string,
|
||||||
|
result: string,
|
||||||
|
isLastStep: boolean,
|
||||||
|
) => Promise<void>;
|
||||||
}): Promise<string[]> {
|
}): Promise<string[]> {
|
||||||
const variables = options.variables ?? {};
|
const variables = options.variables ?? {};
|
||||||
const systemPrompts = options.systemPrompts ?? [];
|
const systemPrompts = options.systemPrompts ?? [];
|
||||||
@@ -87,13 +93,39 @@ export class PipelineExecutor {
|
|||||||
);
|
);
|
||||||
const results: string[] = [];
|
const results: string[] = [];
|
||||||
|
|
||||||
for (const { prompt } of options.pipeline.steps) {
|
const steps = options.pipeline.steps;
|
||||||
|
const totalSteps = steps.length;
|
||||||
|
|
||||||
|
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);
|
const renderedPrompt = this.renderPrompt(prompt, variables);
|
||||||
messages.push(new HumanMessage(renderedPrompt));
|
messages.push(new HumanMessage(renderedPrompt));
|
||||||
|
|
||||||
const finalResponse = await this._resolveToolCalls(chain, messages);
|
const finalResponse = await this._resolveToolCalls(chain, messages);
|
||||||
|
const result = finalResponse?.content?.toString() || '<empty>';
|
||||||
|
|
||||||
messages.push(finalResponse);
|
messages.push(finalResponse);
|
||||||
results.push(finalResponse?.content.toString() ?? '<empty>');
|
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;
|
return results;
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in new issue
Block a user