- контекст проложен от воркера и обоих входов до внешних вызовов: ffmpeg и ffprobe заводятся через exec.CommandContext, SpeechKit и Object Storage принимают ctx вместо context.Background, скачивание записи идёт запросом с контекстом. Прежде остановка сервиса не доходила до чужой работы вовсе - прерванный шаг приговора не выносит: убитый по контексту ffmpeg отдаёт «signal: killed», от настоящего отказа неотличимо ни типом, ни errors.Is, и различает их только ctx.Err(). Задача остаётся на повтор, попытку не тратит и отправителю о несуществующем сбое не сообщает; воркер не считает остановку отказом, а задача не забирается вовсе, если нас уже остановили - клиента Bot API заводит единая точка internal/adapter/telegram: токен стоит в пути каждого обращения, а http.Client кладёт адрес в *url.Error целиком. Чистка на месте употребления закрывала один вызов из пяти — теперь свой Do чистит отказ, подменённый логгер вычищает токен из строк самой библиотеки, а транспорт бота токена не получает вовсе - принятие операции распознавания защищено от отмены своим пределом: SpeechKit мог её принять и начать считать деньги, а потерянный идентификатор заставил бы повтор оплатить ту же запись второй раз - приём по HTTP доводит запись до задачи независимо от отправителя: на контексте запроса один обрыв соединения терял полностью загруженную запись - ответ Telegram с не-2xx кодом больше не становится записью: прежде тело отказа доезжало до хранилища и умирало на ffprobe, уводя диагностику
337 lines
13 KiB
Go
337 lines
13 KiB
Go
package tg
|
||
|
||
import (
|
||
"context"
|
||
"errors"
|
||
"fmt"
|
||
"io"
|
||
"log/slog"
|
||
"net/http"
|
||
"slices"
|
||
"strings"
|
||
|
||
// Транспорт знает адаптер Telegram ровно ради единой точки чистки отказа:
|
||
// второй экземпляр той же функции здесь был бы вторым способом делать одно
|
||
// и то же, а секрет в журнале — необратим. Направление «транспорт не знает
|
||
// адаптера» правилом не держится и уже нарушено HTTP-поверхностью
|
||
// (docs/conventions/go-linters.md, «Что остаётся прозой»).
|
||
"git.vakhrushev.me/av/transcriber/internal/adapter/telegram"
|
||
"git.vakhrushev.me/av/transcriber/internal/contract"
|
||
"git.vakhrushev.me/av/transcriber/internal/service"
|
||
tgbotapi "github.com/go-telegram-bot-api/telegram-bot-api/v5"
|
||
)
|
||
|
||
type TelegramController struct {
|
||
// deps
|
||
transcribeService *service.TranscribeService
|
||
jobRepo contract.TranscriptJobRepository
|
||
logger *slog.Logger
|
||
// params
|
||
bot *tgbotapi.BotAPI
|
||
userWhiteList []string
|
||
updateTimeout int
|
||
}
|
||
|
||
type TelegramConfig struct {
|
||
UpdateTimeout int
|
||
UserWhiteList []string
|
||
}
|
||
|
||
// NewTelegramController принимает готового клиента, а не токен: клиента заводит
|
||
// единая точка `internal/adapter/telegram`, и только её отказ не несёт секрета.
|
||
// Токен сюда не приезжает вовсе — значит, и утечь отсюда ему неоткуда.
|
||
func NewTelegramController(
|
||
config TelegramConfig,
|
||
bot *tgbotapi.BotAPI,
|
||
transcribeService *service.TranscribeService,
|
||
jobRepo contract.TranscriptJobRepository,
|
||
logger *slog.Logger,
|
||
) (*TelegramController, error) {
|
||
if bot == nil {
|
||
return nil, errors.New("telegram bot is not created")
|
||
}
|
||
|
||
controller := &TelegramController{
|
||
bot: bot,
|
||
transcribeService: transcribeService,
|
||
jobRepo: jobRepo,
|
||
logger: logger,
|
||
updateTimeout: config.UpdateTimeout,
|
||
userWhiteList: config.UserWhiteList,
|
||
}
|
||
|
||
return controller, nil
|
||
}
|
||
|
||
// Start принимает контекст жизни процесса и отдаёт его каждому обработчику:
|
||
// скачивание записи и разбор её метаданных — работа с внешним собеседником, и
|
||
// остановка сервиса обязана до неё доходить. Приём обновлений контекстом не
|
||
// правится: его прекращает Stop.
|
||
func (c *TelegramController) Start(ctx context.Context) {
|
||
c.logger.Info("Telegram bot started", "username", c.bot.Self.UserName)
|
||
|
||
u := tgbotapi.NewUpdate(0)
|
||
u.Timeout = c.updateTimeout
|
||
|
||
updates := c.bot.GetUpdatesChan(u)
|
||
|
||
for update := range updates {
|
||
if update.Message == nil { // ignore any non-Message updates
|
||
continue
|
||
}
|
||
|
||
author := update.Message.From.String()
|
||
c.logger.Info("New incoming message", "author", author)
|
||
|
||
if !slices.Contains(c.userWhiteList, author) {
|
||
c.logger.Info("User is not in white list, reject", "author", author)
|
||
c.handleForbiddenUser(update.Message)
|
||
continue
|
||
}
|
||
|
||
// Handle commands
|
||
if update.Message.IsCommand() {
|
||
// Extract the command from the Message
|
||
switch update.Message.Command() {
|
||
case "start":
|
||
c.handleStartCommand(update.Message)
|
||
case "help":
|
||
c.handleHelpCommand(update.Message)
|
||
}
|
||
continue
|
||
}
|
||
|
||
// Handle audio messages and files
|
||
if update.Message.Audio != nil {
|
||
c.handleAudioMessage(ctx, update.Message)
|
||
} else if update.Message.Voice != nil {
|
||
c.handleVoiceMessage(ctx, update.Message)
|
||
} else if update.Message.Document != nil {
|
||
c.handleDocumentMessage(ctx, update.Message)
|
||
}
|
||
}
|
||
}
|
||
|
||
func (c *TelegramController) Stop() {
|
||
c.bot.StopReceivingUpdates()
|
||
}
|
||
|
||
func (c *TelegramController) send(chattable tgbotapi.Chattable) (tgbotapi.Message, error) {
|
||
msg, err := c.bot.Send(chattable)
|
||
if err != nil {
|
||
c.logger.Error("Failed to send message to tg bot", "error", err)
|
||
}
|
||
return msg, err
|
||
}
|
||
|
||
func (c *TelegramController) handleStartCommand(message *tgbotapi.Message) {
|
||
msg := tgbotapi.NewMessage(message.Chat.ID, "Привет! Я бот для расшифровки аудиосообщений. Отправь мне голосовое сообщение или аудиофайл, и я пришлю тебе текст.")
|
||
msg.ReplyToMessageID = message.MessageID
|
||
|
||
c.send(msg)
|
||
}
|
||
|
||
func (c *TelegramController) handleForbiddenUser(message *tgbotapi.Message) {
|
||
msg := tgbotapi.NewMessage(message.Chat.ID, "Извини, тебе нельзя пользоваться этим ботом. Обратись к владельцу бота.")
|
||
msg.ReplyToMessageID = message.MessageID
|
||
|
||
c.send(msg)
|
||
}
|
||
|
||
func (c *TelegramController) handleHelpCommand(message *tgbotapi.Message) {
|
||
helpText := `Я бот для расшифровки аудиосообщений и аудиофайлов.
|
||
|
||
Просто отправь мне:
|
||
- Голосовое сообщение
|
||
- Аудиофайл (mp3, wav, ogg и др.)
|
||
|
||
Я пришлю тебе текст расшифровки.
|
||
|
||
Команды:
|
||
/start - Начало работы с ботом
|
||
/help - Показать эту справку`
|
||
|
||
msg := tgbotapi.NewMessage(message.Chat.ID, helpText)
|
||
msg.ReplyToMessageID = message.MessageID
|
||
|
||
c.send(msg)
|
||
}
|
||
|
||
func (c *TelegramController) handleAudioMessage(ctx context.Context, message *tgbotapi.Message) {
|
||
// Отправляем сообщение о начале обработки
|
||
progressMsg := tgbotapi.NewMessage(message.Chat.ID, "Обрабатываю аудиофайл...")
|
||
progressMsg.ReplyToMessageID = message.MessageID
|
||
sentProgressMsg, err := c.send(progressMsg)
|
||
if err != nil {
|
||
c.logger.Error("Failed to send progress message", "error", err)
|
||
return
|
||
}
|
||
|
||
// Скачиваем файл
|
||
fileReader, fileName, err := c.downloadAudioFile(ctx, message.Audio.FileID)
|
||
if err != nil {
|
||
c.logger.Error("Failed to download audio file", "error", err)
|
||
errorMsg := tgbotapi.NewMessage(message.Chat.ID, "Ошибка при скачивании аудиофайла. Попробуйте еще раз.")
|
||
c.send(errorMsg)
|
||
return
|
||
}
|
||
defer fileReader.Close()
|
||
|
||
// Обрабатываем файл
|
||
job, err := c.transcribeService.CreateJobFromTelegram(ctx, fileReader, fileName, message.Chat.ID, sentProgressMsg.MessageID)
|
||
if err != nil {
|
||
c.logger.Error("Failed to create transcribe job", "error", err)
|
||
errorMsg := tgbotapi.NewMessage(message.Chat.ID, "Ошибка при создании задачи на расшифровку. Попробуйте еще раз.")
|
||
c.send(errorMsg)
|
||
return
|
||
}
|
||
|
||
// Отправляем сообщение об успешном создании задачи
|
||
successMsg := tgbotapi.NewMessage(message.Chat.ID, fmt.Sprintf("Задача на расшифровку создана. ID задачи: %s", job.Id))
|
||
successMsg.ReplyToMessageID = message.MessageID
|
||
c.send(successMsg)
|
||
}
|
||
|
||
func (c *TelegramController) handleVoiceMessage(ctx context.Context, message *tgbotapi.Message) {
|
||
// Отправляем сообщение о начале обработки
|
||
progressMsg := tgbotapi.NewMessage(message.Chat.ID, "Обрабатываю голосовое сообщение...")
|
||
progressMsg.ReplyToMessageID = message.MessageID
|
||
sentProgressMsg, err := c.send(progressMsg)
|
||
if err != nil {
|
||
c.logger.Error("Failed to send progress message", "error", err)
|
||
return
|
||
}
|
||
|
||
// Скачиваем файл
|
||
fileReader, fileName, err := c.downloadAudioFile(ctx, message.Voice.FileID)
|
||
if err != nil {
|
||
c.logger.Error("Failed to download voice file", "error", err)
|
||
errorMsg := tgbotapi.NewMessage(message.Chat.ID, "Ошибка при скачивании голосового сообщения. Попробуйте еще раз.")
|
||
c.send(errorMsg)
|
||
return
|
||
}
|
||
defer fileReader.Close()
|
||
|
||
// Обрабатываем файл
|
||
job, err := c.transcribeService.CreateJobFromTelegram(ctx, fileReader, fileName, message.Chat.ID, sentProgressMsg.MessageID)
|
||
if err != nil {
|
||
c.logger.Error("Failed to create transcribe job", "error", err)
|
||
errorMsg := tgbotapi.NewMessage(message.Chat.ID, "Ошибка при создании задачи на расшифровку. Попробуйте еще раз.")
|
||
c.send(errorMsg)
|
||
return
|
||
}
|
||
|
||
// Отправляем сообщение об успешном создании задачи
|
||
successMsg := tgbotapi.NewMessage(message.Chat.ID, fmt.Sprintf("Задача на расшифровку создана. ID задачи: %s", job.Id))
|
||
successMsg.ReplyToMessageID = message.MessageID
|
||
c.send(successMsg)
|
||
}
|
||
|
||
func (c *TelegramController) handleDocumentMessage(ctx context.Context, message *tgbotapi.Message) {
|
||
// Проверяем, является ли документ аудиофайлом
|
||
if !c.isAudioDocument(message.Document) {
|
||
return
|
||
}
|
||
|
||
// Отправляем сообщение о начале обработки
|
||
progressMsg := tgbotapi.NewMessage(message.Chat.ID, "Обрабатываю аудиофайл...")
|
||
progressMsg.ReplyToMessageID = message.MessageID
|
||
sentProgressMsg, err := c.send(progressMsg)
|
||
if err != nil {
|
||
c.logger.Error("Failed to send progress message", "error", err)
|
||
return
|
||
}
|
||
|
||
// Скачиваем файл
|
||
fileReader, fileName, err := c.downloadAudioFile(ctx, message.Document.FileID)
|
||
if err != nil {
|
||
c.logger.Error("Failed to download document file", "error", err)
|
||
errorMsg := tgbotapi.NewMessage(message.Chat.ID, "Ошибка при скачивании аудиофайла. Попробуйте еще раз.")
|
||
c.send(errorMsg)
|
||
return
|
||
}
|
||
defer fileReader.Close()
|
||
|
||
// Обрабатываем файл
|
||
job, err := c.transcribeService.CreateJobFromTelegram(ctx, fileReader, fileName, message.Chat.ID, sentProgressMsg.MessageID)
|
||
if err != nil {
|
||
c.logger.Error("Failed to create transcribe job", "error", err)
|
||
errorMsg := tgbotapi.NewMessage(message.Chat.ID, "Ошибка при создании задачи на расшифровку. Попробуйте еще раз.")
|
||
c.send(errorMsg)
|
||
return
|
||
}
|
||
|
||
// Отправляем сообщение об успешном создании задачи
|
||
successMsg := tgbotapi.NewMessage(message.Chat.ID, fmt.Sprintf("Задача на расшифровку создана. ID задачи: %s", job.Id))
|
||
successMsg.ReplyToMessageID = message.MessageID
|
||
c.send(successMsg)
|
||
}
|
||
|
||
func (c *TelegramController) downloadAudioFile(ctx context.Context, fileID string) (io.ReadCloser, string, error) {
|
||
// Получаем информацию о файле
|
||
file, err := c.bot.GetFile(tgbotapi.FileConfig{FileID: fileID})
|
||
if err != nil {
|
||
return nil, "", fmt.Errorf("failed to get file info: %w", err)
|
||
}
|
||
|
||
// Скачиваем файл. Запрос заводится с контекстом: скачивание шестичасовой
|
||
// записи иначе продолжается и после остановки сервиса, а ссылка на файл
|
||
// несёт токен бота — держать её живой дольше нужного незачем.
|
||
//
|
||
// Клиент берётся у бота, а не `http.DefaultClient`: у бота он свой, и его
|
||
// отказ уже не несёт адреса (`internal/adapter/telegram`, единая точка).
|
||
fileURL := file.Link(c.bot.Token)
|
||
request, err := http.NewRequestWithContext(ctx, http.MethodGet, fileURL, nil)
|
||
if err != nil {
|
||
return nil, "", fmt.Errorf("failed to build download request: %w", telegram.WithoutURL(err))
|
||
}
|
||
|
||
resp, err := c.bot.Client.Do(request)
|
||
if err != nil {
|
||
return nil, "", fmt.Errorf("failed to download file: %w", err)
|
||
}
|
||
|
||
// Отказ выдачи файла — это не запись. Без проверки телом «записи» станет
|
||
// JSON вида `{"ok":false,…}`: он доедет до хранилища, ляжет рабочей копией
|
||
// и умрёт на `ffprobe`, а отправитель получит жалобу на свой файл вместо
|
||
// правды о протухшей ссылке.
|
||
if resp.StatusCode != http.StatusOK {
|
||
if err := resp.Body.Close(); err != nil {
|
||
c.logger.Error("Failed to close download response", "error", err)
|
||
}
|
||
return nil, "", fmt.Errorf("failed to download file: unexpected status %d", resp.StatusCode)
|
||
}
|
||
|
||
// Получаем имя файла из URL
|
||
fileName := file.FilePath
|
||
if fileName == "" {
|
||
fileName = "audio.ogg"
|
||
}
|
||
|
||
return resp.Body, fileName, nil
|
||
}
|
||
|
||
func (c *TelegramController) isAudioDocument(document *tgbotapi.Document) bool {
|
||
// Проверяем MIME-тип документа
|
||
if document.MimeType != "" {
|
||
return strings.HasPrefix(document.MimeType, "audio/") || strings.HasPrefix(document.MimeType, "video/")
|
||
}
|
||
|
||
// Проверяем расширение файла
|
||
audioExtensions := []string{".mp3", ".wav", ".ogg", ".flac", ".m4a", ".aac", ".wma"}
|
||
filename := document.FileName
|
||
for _, ext := range audioExtensions {
|
||
if len(filename) >= len(ext) && strings.ToLower(filename[len(filename)-len(ext):]) == ext {
|
||
return true
|
||
}
|
||
}
|
||
|
||
return false
|
||
}
|
||
|
||
type EmptyBotTokenError struct{}
|
||
|
||
func (e *EmptyBotTokenError) Error() string {
|
||
return "telegram bot token is empty"
|
||
}
|