-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathtelegram_whisper_bot.py
More file actions
483 lines (404 loc) · 21.7 KB
/
Copy pathtelegram_whisper_bot.py
File metadata and controls
483 lines (404 loc) · 21.7 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
import os
import asyncio
import logging
import queue
import threading
import time
from pathlib import Path
from typing import Optional, Dict, Any
import tempfile
import hashlib
from datetime import datetime
import whisper
from telegram import Update, Message
from telegram.ext import Application, MessageHandler, filters, ContextTypes
from dotenv import load_dotenv
# Загружаем переменные окружения
load_dotenv()
# Конфигурация
class Config:
def __init__(self):
self.TG_BOT_TOKEN = os.getenv('TG_BOT_TOKEN')
self.WHISPER_MODEL = os.getenv('WHISPER_MODEL', 'base')
self.WHISPER_DEVICE = os.getenv('WHISPER_DEVICE', 'cpu') # cpu/cuda/auto
self.WORKER_COUNT = int(os.getenv('WORKER_COUNT', '2'))
self.QUEUE_MAXSIZE = int(os.getenv('QUEUE_MAXSIZE', '200'))
self.SAVE_VOICES = os.getenv('SAVE_VOICES', 'True').lower() == 'true'
self.VOICES_DIR = os.getenv('VOICES_DIR', 'voices')
self.LOG_LEVEL = os.getenv('LOG_LEVEL', 'INFO')
# Валидация обязательных параметров
if not self.TG_BOT_TOKEN:
raise ValueError("TG_BOT_TOKEN не установлен в .env файле")
# Создаем директорию для голосовых, если нужно сохранять
if self.SAVE_VOICES:
Path(self.VOICES_DIR).mkdir(exist_ok=True)
config = Config()
# Настройка логирования
def setup_logging():
if config.LOG_LEVEL.upper() == 'NONE':
logging.disable(logging.CRITICAL)
return
level = getattr(logging, config.LOG_LEVEL.upper(), logging.INFO)
logging.basicConfig(
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
level=level
)
setup_logging()
logger = logging.getLogger(__name__)
class VoiceProcessor:
"""Обработчик голосовых сообщений с использованием Whisper"""
def __init__(self):
# Определяем устройство для обработки
self.device = self._get_device()
logger.info(f"Загружаем модель Whisper: {config.WHISPER_MODEL} на устройство: {self.device}")
# Загружаем модель с явным указанием устройства
self.model = whisper.load_model(config.WHISPER_MODEL, device=self.device)
logger.info("Модель Whisper загружена успешно")
def _get_device(self):
"""Определяет устройство для обработки"""
import torch
if config.WHISPER_DEVICE.lower() == 'auto':
if torch.cuda.is_available():
device = 'cuda'
logger.info("CUDA доступна, используем GPU")
else:
device = 'cpu'
logger.info("CUDA недоступна, используем CPU")
elif config.WHISPER_DEVICE.lower() == 'cuda':
if torch.cuda.is_available():
device = 'cuda'
logger.info("Принудительно используем CUDA/GPU")
else:
logger.warning("CUDA запрошена, но недоступна. Используем CPU")
device = 'cpu'
else: # cpu или любое другое значение
device = 'cpu'
logger.info("Принудительно используем CPU")
return device
def transcribe_audio(self, audio_path: str) -> Dict[str, Any]:
"""Транскрибирует аудио файл"""
try:
logger.info(f"Начинаем транскрипцию: {audio_path}")
start_time = time.time()
result = self.model.transcribe(
audio_path,
# language='ru', # Можно сделать автоопределение, убрав этот параметр
task='transcribe'
)
processing_time = time.time() - start_time
logger.info(f"Транскрипция завершена за {processing_time:.2f} сек")
return {
'success': True,
'text': result['text'].strip(),
'language': result.get('language', 'unknown'),
'processing_time': processing_time
}
except Exception as e:
logger.error(f"Ошибка при транскрипции: {e}")
return {
'success': False,
'error': str(e)
}
class QueueManager:
"""Менеджер очереди обработки голосовых сообщений"""
def __init__(self):
self.processing_queue = queue.Queue(maxsize=config.QUEUE_MAXSIZE)
self.result_queue = queue.Queue() # Очередь для результатов
self.active_tasks = {} # chat_id -> task_info
self.processor = VoiceProcessor()
self.workers = []
self.running = True
self.event_loop = None # Ссылка на основной event loop
# Запускаем worker'ов
for i in range(config.WORKER_COUNT):
worker = threading.Thread(target=self._worker, args=(i,), daemon=True)
worker.start()
self.workers.append(worker)
logger.info(f"Запущен worker #{i}")
def add_task(self, chat_id: int, message: Message, audio_path: str) -> bool:
"""Добавляет задачу в очередь"""
try:
task = {
'chat_id': chat_id,
'message': message,
'audio_path': audio_path,
'timestamp': datetime.now()
}
self.processing_queue.put_nowait(task)
self.active_tasks[chat_id] = {
'status': 'queued',
'position': self.processing_queue.qsize()
}
logger.info(f"Задача добавлена в очередь для chat_id={chat_id}, позиция: {self.processing_queue.qsize()}")
return True
except queue.Full:
logger.warning(f"Очередь переполнена, отклоняем задачу для chat_id={chat_id}")
return False
def get_queue_status(self, chat_id: int) -> Optional[Dict[str, Any]]:
"""Возвращает статус задачи в очереди"""
return self.active_tasks.get(chat_id)
def set_event_loop(self, loop):
"""Устанавливает ссылку на основной event loop"""
self.event_loop = loop
def _schedule_coroutine(self, coro):
"""Планирует выполнение корутины в основном event loop"""
if self.event_loop and not self.event_loop.is_closed():
try:
asyncio.run_coroutine_threadsafe(coro, self.event_loop)
except Exception as e:
logger.error(f"Ошибка при планировании корутины: {e}")
else:
logger.error("Event loop недоступен для планирования корутины")
def _worker(self, worker_id: int):
"""Worker для обработки задач из очереди"""
logger.info(f"Worker #{worker_id} запущен")
while self.running:
try:
# Получаем задачу из очереди
task = self.processing_queue.get(timeout=1)
chat_id = task['chat_id']
message = task['message']
audio_path = task['audio_path']
logger.info(f"Worker #{worker_id} обрабатывает задачу для chat_id={chat_id}")
# Обновляем статус
self.active_tasks[chat_id] = {'status': 'processing'}
try:
# Обрабатываем аудио
result = self.processor.transcribe_audio(audio_path)
# Планируем отправку результата в основном event loop
self._schedule_coroutine(self._send_result(message, result, audio_path))
except Exception as e:
logger.error(f"Ошибка при обработке задачи: {e}")
self._schedule_coroutine(self._send_error(message, str(e), audio_path))
finally:
# Удаляем из активных задач
self.active_tasks.pop(chat_id, None)
self.processing_queue.task_done()
except queue.Empty:
continue
except Exception as e:
logger.error(f"Критическая ошибка в worker #{worker_id}: {e}")
async def _send_result(self, original_message: Message, result: Dict[str, Any], audio_path: str):
"""Отправляет результат транскрипции"""
try:
if result['success']:
text = result['text']
if not text:
response = "🤔 Не удалось распознать речь в голосовом сообщении"
else:
processing_time = result.get('processing_time', 0)
response = f"📝 *Текст голосового сообщения:*\n\n{text}\n\n_Обработано за {processing_time:.1f} сек_"
await original_message.reply_text(
response,
parse_mode='Markdown',
reply_to_message_id=original_message.message_id
)
else:
await self._send_error(original_message, result.get('error', 'Неизвестная ошибка'), audio_path)
except Exception as e:
logger.error(f"Ошибка при отправке результата: {e}")
finally:
# Удаляем временный файл, если не нужно сохранять
if not config.SAVE_VOICES and os.path.exists(audio_path):
try:
os.remove(audio_path)
logger.info(f"Временный файл удален: {audio_path}")
except Exception as e:
logger.error(f"Не удалось удалить временный файл: {e}")
async def _send_error(self, original_message: Message, error: str, audio_path: str):
"""Отправляет сообщение об ошибке"""
try:
await original_message.reply_text(
f"❌ Произошла ошибка при обработке голосового сообщения:\n{error}",
reply_to_message_id=original_message.message_id
)
except Exception as e:
logger.error(f"Ошибка при отправке сообщения об ошибке: {e}")
finally:
# Удаляем файл при ошибке в любом случае
if os.path.exists(audio_path):
try:
os.remove(audio_path)
logger.info(f"Файл удален после ошибки: {audio_path}")
except Exception as e:
logger.error(f"Не удалось удалить файл после ошибки: {e}")
def shutdown(self):
"""Корректное завершение работы"""
logger.info("Завершаем работу QueueManager...")
self.running = False
# Ждем завершения всех workers
for worker in self.workers:
worker.join(timeout=5)
# Глобальный менеджер очереди
queue_manager = QueueManager()
async def handle_voice(update: Update, context: ContextTypes.DEFAULT_TYPE):
"""Обработчик голосовых сообщений"""
chat_id = update.effective_chat.id
message = update.message
# Проверяем, есть ли уже активная задача для этого чата
if queue_manager.get_queue_status(chat_id):
await message.reply_text(
"⏳ У вас уже есть голосовое сообщение в обработке. Пожалуйста, дождитесь завершения.",
reply_to_message_id=message.message_id
)
return
try:
# Получаем файл голосового сообщения
voice = message.voice
file = await context.bot.get_file(voice.file_id)
# Генерируем имя файла
timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
file_hash = hashlib.md5(f"{chat_id}_{voice.file_id}".encode()).hexdigest()[:8]
if config.SAVE_VOICES:
# Сохраняем в постоянную директорию
filename = f"voice_{timestamp}_{file_hash}.ogg"
audio_path = os.path.join(config.VOICES_DIR, filename)
else:
# Используем временный файл
temp_file = tempfile.NamedTemporaryFile(
suffix='.ogg',
delete=False,
prefix=f'voice_{timestamp}_{file_hash}_'
)
audio_path = temp_file.name
temp_file.close()
# Скачиваем файл
await file.download_to_drive(audio_path)
logger.info(f"Голосовое сообщение сохранено: {audio_path}")
# Добавляем в очередь
if queue_manager.add_task(chat_id, message, audio_path):
queue_position = queue_manager.processing_queue.qsize()
if queue_position > 1:
await message.reply_text(
f"📥 Голосовое сообщение добавлено в очередь на обработку.\n"
f"Позиция в очереди: {queue_position}\n"
f"⏱ Примерное время ожидания: {queue_position * 10-30} сек",
reply_to_message_id=message.message_id
)
else:
await message.reply_text(
"🔄 Начинаю обработку голосового сообщения...",
reply_to_message_id=message.message_id
)
else:
# Очередь переполнена
await message.reply_text(
"😔 Извините, сервер перегружен. Попробуйте позже.\n"
f"Максимальная длина очереди: {config.QUEUE_MAXSIZE}",
reply_to_message_id=message.message_id
)
# Удаляем файл, если очередь переполнена
if os.path.exists(audio_path):
os.remove(audio_path)
except Exception as e:
logger.error(f"Ошибка при обработке голосового сообщения: {e}")
await message.reply_text(
f"❌ Произошла ошибка при получении голосового сообщения: {e}",
reply_to_message_id=message.message_id
)
async def handle_video_note(update: Update, context: ContextTypes.DEFAULT_TYPE):
"""Обработчик кружочков (video note)"""
# Аналогично голосовым сообщениям, но для видео-кружочков
chat_id = update.effective_chat.id
message = update.message
if queue_manager.get_queue_status(chat_id):
await message.reply_text(
"⏳ У вас уже есть сообщение в обработке. Пожалуйста, дождитесь завершения.",
reply_to_message_id=message.message_id
)
return
try:
video_note = message.video_note
file = await context.bot.get_file(video_note.file_id)
timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
file_hash = hashlib.md5(f"{chat_id}_{video_note.file_id}".encode()).hexdigest()[:8]
if config.SAVE_VOICES:
filename = f"videonote_{timestamp}_{file_hash}.mp4"
audio_path = os.path.join(config.VOICES_DIR, filename)
else:
temp_file = tempfile.NamedTemporaryFile(
suffix='.mp4',
delete=False,
prefix=f'videonote_{timestamp}_{file_hash}_'
)
audio_path = temp_file.name
temp_file.close()
await file.download_to_drive(audio_path)
logger.info(f"Видео-кружочек сохранен: {audio_path}")
if queue_manager.add_task(chat_id, message, audio_path):
await message.reply_text(
"🔄 Начинаю извлечение аудио из видео-кружочка...",
reply_to_message_id=message.message_id
)
else:
await message.reply_text(
"😔 Извините, сервер перегружен. Попробуйте позже.",
reply_to_message_id=message.message_id
)
if os.path.exists(audio_path):
os.remove(audio_path)
except Exception as e:
logger.error(f"Ошибка при обработке видео-кружочка: {e}")
await message.reply_text(
f"❌ Произошла ошибка: {e}",
reply_to_message_id=message.message_id
)
async def start_command(update: Update, context: ContextTypes.DEFAULT_TYPE):
"""Команда /start"""
welcome_text = (
"🎤 *Привет! Я бот для распознавания речи*\n\n"
"Отправь мне голосовое сообщение или видео-кружочек, "
"и я преобразую речь в текст с помощью OpenAI Whisper.\n\n"
"📊 *Информация о сервере:*\n"
f"• Модель Whisper: `{config.WHISPER_MODEL}`\n"
f"• Устройство обработки: `{config.WHISPER_DEVICE}`\n"
f"• Количество обработчиков: {config.WORKER_COUNT}\n"
f"• Максимальный размер очереди: {config.QUEUE_MAXSIZE}\n"
f"• Сохранение файлов: {'✅' if config.SAVE_VOICES else '❌'}\n\n"
"Просто отправь голосовое сообщение! 🚀"
)
await update.message.reply_text(welcome_text, parse_mode='Markdown')
async def status_command(update: Update, context: ContextTypes.DEFAULT_TYPE):
"""Команда /status - показывает статус очереди"""
chat_id = update.effective_chat.id
queue_size = queue_manager.processing_queue.qsize()
status_text = f"📊 *Статус обработки:*\n\n"
status_text += f"• Задач в очереди: {queue_size}\n"
status_text += f"• Активных обработчиков: {config.WORKER_COUNT}\n"
user_status = queue_manager.get_queue_status(chat_id)
if user_status:
if user_status['status'] == 'queued':
status_text += f"• Ваша позиция в очереди: {user_status['position']}\n"
elif user_status['status'] == 'processing':
status_text += "• Ваше сообщение обрабатывается прямо сейчас ⚡\n"
else:
status_text += "• У вас нет активных задач\n"
await update.message.reply_text(status_text, parse_mode='Markdown')
def main():
"""Главная функция"""
logger.info("Запускаем Telegram бота...")
# Создаем приложение
application = Application.builder().token(config.TG_BOT_TOKEN).build()
# Регистрируем обработчики
application.add_handler(MessageHandler(filters.VOICE, handle_voice))
application.add_handler(MessageHandler(filters.VIDEO_NOTE, handle_video_note))
application.add_handler(MessageHandler(filters.COMMAND & filters.Regex(r'^/start'), start_command))
application.add_handler(MessageHandler(filters.COMMAND & filters.Regex(r'^/status'), status_command))
# Устанавливаем event loop после инициализации приложения
async def post_init(application):
loop = asyncio.get_running_loop()
queue_manager.set_event_loop(loop)
logger.info("Event loop установлен для QueueManager")
application.post_init = post_init
try:
logger.info("Бот запущен и готов к работе!")
# Запускаем бота (это создаст и запустит свой event loop)
application.run_polling(allowed_updates=Update.ALL_TYPES)
except KeyboardInterrupt:
logger.info("Получен сигнал завершения...")
finally:
# Корректное завершение
logger.info("Завершаем работу...")
queue_manager.shutdown()
if __name__ == '__main__':
main()