- записи, метаданные и файлы съехались под один каталог данных; появилась панель владельца, а gin, goqu, goose и требование CGO ушли - захват задачи стал одним запросом с RETURNING; заведены число попыток, состояние dead и нарастающая пауза вместо признака is_error - имя файла в хранилище задаёт сервис и в журнал не идёт: вместе с идентификатором записи оно собирало бы ссылку на скачивание
537 lines
20 KiB
Go
537 lines
20 KiB
Go
package service
|
|
|
|
import (
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"log/slog"
|
|
"math"
|
|
"path/filepath"
|
|
"strings"
|
|
"time"
|
|
|
|
"git.vakhrushev.me/av/transcriber/internal/contract"
|
|
"git.vakhrushev.me/av/transcriber/internal/entity"
|
|
"git.vakhrushev.me/av/transcriber/internal/metrics"
|
|
"github.com/google/uuid"
|
|
)
|
|
|
|
const (
|
|
defaultAudioExt = "audio"
|
|
|
|
// Предел попыток. Число обратимо и живёт здесь одним местом; счётчик растёт
|
|
// при захвате и обнуляется на шаге, завершившемся без отказа.
|
|
maxAttempts = 5
|
|
|
|
// Пауза перед повтором отказавшей задачи растёт с числом попыток до
|
|
// потолка. Ожидание чужой операции этой паузой не выражается — у него своя
|
|
// задержка числом, и попытки оно не тратит.
|
|
retryDelayBase = time.Second
|
|
retryDelayCap = 5 * time.Minute
|
|
|
|
// Сроки захвата. Каждый не меньше того, что его шаг может занять на самом
|
|
// длинном допустимом входе: расчётный потолок записи — шесть часов, и
|
|
// конвертация такой записи идёт дольше часа по построению.
|
|
conversionAcquireTimeout = 8 * time.Hour
|
|
transcribeAcquireTimeout = 8 * time.Hour
|
|
checkAcquireTimeout = time.Hour
|
|
|
|
// Задержки опроса операции распознавания. Числа, а не функция числа попыток:
|
|
// счётчик на ожидании обнулён, и выведенная из него пауза выродилась бы в
|
|
// своё наименьшее значение, учащая опрос платного сервиса.
|
|
firstCheckDelay = 10 * time.Second
|
|
nextCheckDelay = 5 * time.Second
|
|
)
|
|
|
|
type TranscribeService struct {
|
|
jobRepo contract.TranscriptJobRepository
|
|
fileRepo contract.FileRepository
|
|
metaviewer contract.AudioMetaViewer
|
|
converter contract.AudioFileConverter
|
|
recognizer contract.AudioRecognizer
|
|
tgSender contract.TelegramMessageSender
|
|
logger *slog.Logger
|
|
}
|
|
|
|
func NewTranscribeService(
|
|
jobRepo contract.TranscriptJobRepository,
|
|
fileRepo contract.FileRepository,
|
|
metaviewer contract.AudioMetaViewer,
|
|
converter contract.AudioFileConverter,
|
|
recognizer contract.AudioRecognizer,
|
|
tgSender contract.TelegramMessageSender,
|
|
logger *slog.Logger,
|
|
) *TranscribeService {
|
|
return &TranscribeService{
|
|
jobRepo: jobRepo,
|
|
fileRepo: fileRepo,
|
|
metaviewer: metaviewer,
|
|
converter: converter,
|
|
recognizer: recognizer,
|
|
tgSender: tgSender,
|
|
logger: logger,
|
|
}
|
|
}
|
|
|
|
func (s *TranscribeService) CreateJobFromTelegram(file io.Reader, fileName string, chatId int64, replyMsgId int) (*entity.TranscribeJob, error) {
|
|
job := &entity.TranscribeJob{
|
|
State: entity.StateCreated,
|
|
Source: entity.SourceTelegram,
|
|
TgChatId: &chatId,
|
|
TgReplyMessageId: &replyMsgId,
|
|
}
|
|
|
|
return s.createTranscribeJob(job, file, fileName)
|
|
}
|
|
|
|
func (s *TranscribeService) CreateJobFromApi(file io.Reader, fileName string) (*entity.TranscribeJob, error) {
|
|
job := &entity.TranscribeJob{
|
|
State: entity.StateCreated,
|
|
Source: entity.SourceApi,
|
|
}
|
|
|
|
return s.createTranscribeJob(job, file, fileName)
|
|
}
|
|
|
|
func (s *TranscribeService) createTranscribeJob(job *entity.TranscribeJob, file io.Reader, fileName string) (*entity.TranscribeJob, error) {
|
|
// Определяем расширение файла
|
|
ext := filepath.Ext(fileName)
|
|
if ext == "" {
|
|
ext = fmt.Sprintf(".%s", defaultAudioExt) // fallback если расширение не определено
|
|
}
|
|
|
|
// Собственное имя записи: идентификатор с расширением. Имя, данное
|
|
// отправителем, в хранилище не попадает — от него взято только расширение.
|
|
fileId := uuid.NewString()
|
|
storageFileName := fmt.Sprintf("%s%s", fileId, ext)
|
|
|
|
// Содержимое ложится в рабочую копию потоком: в память запись целиком не
|
|
// читается, расчётный потолок — шесть часов.
|
|
work, err := s.fileRepo.Stage(ext, file)
|
|
if err != nil {
|
|
s.logger.Error("Failed to stage uploaded file", "error", err)
|
|
return nil, err
|
|
}
|
|
defer s.closeWork(work)
|
|
|
|
// В журнал идёт расширение, и только оно. Имя, данное отправителем, не
|
|
// пишется по инварианту приватности; имя, под которым файл ложится в
|
|
// хранилище, — потому что оно последняя часть ссылки на скачивание, и
|
|
// строка журнала вместе с идентификатором записи собрала бы её целиком.
|
|
s.logger.Info("Creating transcribe job", "file_ext", ext)
|
|
|
|
info, err := s.metaviewer.GetInfo(work.Path())
|
|
if err != nil {
|
|
s.logger.Error("Failed to get file info", "error", err, "file_ext", ext)
|
|
return nil, err
|
|
}
|
|
|
|
size, err := work.Size()
|
|
if err != nil {
|
|
s.logger.Error("Failed to measure uploaded file", "error", err)
|
|
return nil, err
|
|
}
|
|
|
|
fileRecord, err := s.fileRepo.CreateLocal(storageFileName, work)
|
|
if err != nil {
|
|
s.logger.Error("Failed to create file record", "error", err, "file_ext", ext)
|
|
return nil, err
|
|
}
|
|
|
|
s.logger.Info("File uploaded successfully",
|
|
"file_id", fileRecord.Id,
|
|
"size", size,
|
|
"duration_seconds", info.Seconds)
|
|
|
|
metrics.InputFileDurationHistogram.WithLabelValues().Observe(float64(info.Seconds))
|
|
metrics.ObserveInputFileSize(ext, size)
|
|
|
|
job.FileID = &fileRecord.Id
|
|
|
|
if err := s.jobRepo.Create(job); err != nil {
|
|
s.logger.Error("Failed to create job record", "error", err, "file_id", fileRecord.Id)
|
|
return nil, err
|
|
}
|
|
|
|
s.logger.Info("Transcribe job created successfully", "job_id", job.Id, "file_id", fileRecord.Id)
|
|
|
|
return job, nil
|
|
}
|
|
|
|
func (s *TranscribeService) FindAndRunConversionJob() error {
|
|
return s.runStep(entity.StateCreated, conversionAcquireTimeout, s.convertJob)
|
|
}
|
|
|
|
func (s *TranscribeService) FindAndRunTranscribeJob() error {
|
|
return s.runStep(entity.StateConverted, transcribeAcquireTimeout, s.transcribeJob)
|
|
}
|
|
|
|
func (s *TranscribeService) FindAndRunTranscribeCheckJob() error {
|
|
return s.runStep(entity.StateTranscribe, checkAcquireTimeout, s.checkTranscribeJob)
|
|
}
|
|
|
|
// runStep забирает задачу и отдаёт её шагу. Отказ шага не оставляет задачу
|
|
// захваченной до конца срока: захват снимается, и задача ждёт нарастающую паузу
|
|
// — иначе повтор наступал бы через восемь часов, а не через секунду.
|
|
func (s *TranscribeService) runStep(state string, expiration time.Duration, step func(job *entity.TranscribeJob, holder string) error) error {
|
|
job, holder, err := s.findJob(state, expiration)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
if err := step(job, holder); err != nil {
|
|
s.scheduleRetry(job, holder, err)
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (s *TranscribeService) convertJob(job *entity.TranscribeJob, holder string) error {
|
|
s.logger.Info("Starting conversion job", "job_id", job.Id)
|
|
|
|
if job.FileID == nil {
|
|
s.logger.Error("Job has no file", "job_id", job.Id)
|
|
return s.failJob(job, holder, errors.New("job has no file"), "у задачи нет записи")
|
|
}
|
|
|
|
srcFile, err := s.fileRepo.GetByID(*job.FileID)
|
|
if err != nil {
|
|
s.logger.Error("Failed to get source file", "error", err, "file_id", *job.FileID)
|
|
return err
|
|
}
|
|
|
|
// Получаем расширение исходного файла для метрики
|
|
srcExt := strings.TrimPrefix(filepath.Ext(srcFile.FileName), ".")
|
|
if srcExt == "" {
|
|
srcExt = defaultAudioExt
|
|
}
|
|
|
|
src, err := s.fileRepo.Localize(*job.FileID)
|
|
if err != nil {
|
|
s.logger.Error("Failed to localize source file", "error", err, "file_id", *job.FileID)
|
|
return err
|
|
}
|
|
defer s.closeWork(src)
|
|
|
|
dest, err := s.fileRepo.StageEmpty(".ogg")
|
|
if err != nil {
|
|
s.logger.Error("Failed to stage converted file", "error", err, "job_id", job.Id)
|
|
return err
|
|
}
|
|
defer s.closeWork(dest)
|
|
|
|
s.logger.Info("Converting file", "job_id", job.Id, "src_format", srcExt)
|
|
|
|
// Измеряем время конвертации
|
|
startTime := time.Now()
|
|
err = s.converter.Convert(src.Path(), dest.Path())
|
|
conversionDuration := time.Since(startTime)
|
|
|
|
// Записываем метрику времени конвертации
|
|
metrics.ObserveConversionDuration(srcExt, "ogg", err != nil, conversionDuration.Seconds())
|
|
|
|
if err != nil {
|
|
s.logger.Error("File conversion failed",
|
|
"error", err,
|
|
"job_id", job.Id,
|
|
"duration", conversionDuration)
|
|
return s.failJob(job, holder, err, "сбой конвертации файла")
|
|
}
|
|
|
|
destSize, err := dest.Size()
|
|
if err != nil {
|
|
s.logger.Error("Failed to measure converted file", "error", err, "job_id", job.Id)
|
|
return err
|
|
}
|
|
|
|
s.logger.Info("File conversion completed",
|
|
"job_id", job.Id,
|
|
"duration", conversionDuration,
|
|
"output_size", destSize)
|
|
|
|
// Записываем метрику размера выходного файла
|
|
metrics.OutputFileSizeHistogram.WithLabelValues("ogg").Observe(float64(destSize))
|
|
|
|
destFileName := fmt.Sprintf("%s%s", uuid.NewString(), ".ogg")
|
|
destFileRecord, err := s.fileRepo.CreateLocal(destFileName, dest)
|
|
if err != nil {
|
|
s.logger.Error("Failed to create converted file record", "error", err, "job_id", job.Id)
|
|
return err
|
|
}
|
|
|
|
// Ссылка переставляется только после того, как запись о новом файле есть:
|
|
// иначе повтор оставил бы задачу указывающей на файл, которого нет.
|
|
job.FileID = &destFileRecord.Id
|
|
job.MoveToState(entity.StateConverted)
|
|
|
|
if err := s.jobRepo.Save(job, holder); err != nil {
|
|
s.logger.Error("Failed to save job", "error", err, "job_id", job.Id)
|
|
return err
|
|
}
|
|
|
|
s.logger.Info("Conversion job completed successfully", "job_id", job.Id)
|
|
return nil
|
|
}
|
|
|
|
func (s *TranscribeService) transcribeJob(job *entity.TranscribeJob, holder string) error {
|
|
s.logger.Info("Starting transcribe job", "job_id", job.Id)
|
|
|
|
if job.FileID == nil {
|
|
s.logger.Error("Job has no file", "job_id", job.Id)
|
|
return s.failJob(job, holder, errors.New("job has no file"), "у задачи нет записи")
|
|
}
|
|
|
|
fileRecord, err := s.fileRepo.GetByID(*job.FileID)
|
|
if err != nil {
|
|
s.logger.Error("Failed to get file record", "error", err, "file_id", *job.FileID)
|
|
return err
|
|
}
|
|
|
|
content, err := s.fileRepo.Open(*job.FileID)
|
|
if err != nil {
|
|
s.logger.Error("Failed to open file", "error", err, "file_id", *job.FileID)
|
|
return err
|
|
}
|
|
defer func() {
|
|
if err := content.Close(); err != nil {
|
|
s.logger.Error("Failed to close file", "error", err, "file_id", *job.FileID)
|
|
}
|
|
}()
|
|
|
|
s.logger.Info("Starting recognition", "job_id", job.Id, "file_id", *job.FileID)
|
|
|
|
// Запускаем асинхронное распознавание
|
|
operationID, err := s.recognizer.Recognize(content, fileRecord.FileName)
|
|
if err != nil {
|
|
s.logger.Error("Failed to start recognition", "error", err, "job_id", job.Id)
|
|
return err
|
|
}
|
|
|
|
s.logger.Info("Recognition started",
|
|
"job_id", job.Id,
|
|
"operation_id", operationID)
|
|
|
|
destFileRecord, err := s.fileRepo.CreateRemote(fileRecord.FileName, fileRecord.Size)
|
|
if err != nil {
|
|
s.logger.Error("Failed to create S3 file record", "error", err, "job_id", job.Id)
|
|
return err
|
|
}
|
|
|
|
// Обновляем задачу с ID операции распознавания
|
|
job.FileID = &destFileRecord.Id
|
|
job.RecognitionOpID = &operationID
|
|
delayTime := time.Now().Add(firstCheckDelay)
|
|
job.MoveToStateAndDelay(entity.StateTranscribe, &delayTime)
|
|
|
|
if err := s.jobRepo.Save(job, holder); err != nil {
|
|
s.logger.Error("Failed to save job", "error", err, "job_id", job.Id)
|
|
return err
|
|
}
|
|
|
|
s.logger.Info("Transcribe job updated successfully", "job_id", job.Id)
|
|
return nil
|
|
}
|
|
|
|
func (s *TranscribeService) checkTranscribeJob(job *entity.TranscribeJob, holder string) error {
|
|
if job.RecognitionOpID == nil {
|
|
s.logger.Error("Recognition operation ID not found", "job_id", job.Id)
|
|
return fmt.Errorf("recogniton opId not found for job: %s", job.Id)
|
|
}
|
|
|
|
opId := *job.RecognitionOpID
|
|
|
|
// Проверяем статус операции
|
|
s.logger.Info("Checking operation status", "job_id", job.Id, "operation_id", opId)
|
|
recResult, err := s.recognizer.CheckRecognitionStatus(opId)
|
|
if err != nil {
|
|
s.logger.Error("Failed to check recognition status", "error", err, "operation_id", opId)
|
|
return err
|
|
}
|
|
|
|
if recResult.IsInProgress() {
|
|
// Операция ещё не завершена. Шаг отработал без отказа, поэтому задержка
|
|
// здесь своя, числом, а число попыток обнуляется переходом: ожидание
|
|
// чужой операции попытку не тратит.
|
|
s.logger.Info("Operation in progress", "job_id", job.Id, "operation_id", opId)
|
|
delayTime := time.Now().Add(nextCheckDelay)
|
|
job.MoveToStateAndDelay(entity.StateTranscribe, &delayTime)
|
|
if err := s.jobRepo.Save(job, holder); err != nil {
|
|
s.logger.Error("Failed to save job", "error", err, "job_id", job.Id)
|
|
return err
|
|
}
|
|
return nil
|
|
}
|
|
|
|
if recResult.IsFailed() {
|
|
errorText := recResult.GetError()
|
|
s.logger.Error("Operation failed",
|
|
"job_id", job.Id,
|
|
"operation_id", opId,
|
|
"error_message", errorText)
|
|
return s.failJob(job, holder, errors.New(errorText), "сбой при распознавании файла")
|
|
}
|
|
|
|
// Операция завершена, получаем результат
|
|
transcriptionText, err := s.recognizer.GetRecognitionText(opId)
|
|
if err != nil {
|
|
s.logger.Error("Failed to get recognition text", "error", err, "operation_id", opId)
|
|
return err
|
|
}
|
|
|
|
s.logger.Info("Transcribe operation completed successfully",
|
|
"job_id", job.Id,
|
|
"operation_id", opId,
|
|
"text_length", len(transcriptionText))
|
|
|
|
if len(transcriptionText) == 0 {
|
|
return s.completeJob(job, holder, "Ой, кажется, на аудиозаписи нет текста.")
|
|
}
|
|
|
|
// Завершаем задачу
|
|
return s.completeJob(job, holder, transcriptionText)
|
|
}
|
|
|
|
// findJob забирает задачу и отдаёт её вместе с признаком захвата, который шаг
|
|
// держит. Задача, захваченная сверх предела попыток, до шага не доходит: её
|
|
// переводят в «мертва» и сообщают об этом отправителю.
|
|
func (s *TranscribeService) findJob(state string, expiration time.Duration) (*entity.TranscribeJob, string, error) {
|
|
acquisitionId := uuid.NewString()
|
|
rottingTime := time.Now().Add(-1 * expiration)
|
|
|
|
job, err := s.jobRepo.FindAndAcquire(state, acquisitionId, rottingTime)
|
|
if err != nil {
|
|
// Признак узнаётся по смыслу: репозиторий вправе обернуть свой отказ
|
|
// пояснением, и приведение типа от этого сломалось бы молча.
|
|
var notFound *contract.JobNotFoundError
|
|
if errors.As(err, ¬Found) {
|
|
return nil, "", &contract.NoopJobError{State: state}
|
|
}
|
|
s.logger.Error("Failed to find and acquire job", "state", state, "error", err)
|
|
return nil, "", fmt.Errorf("failed find and acquire job: %s, %w", state, err)
|
|
}
|
|
|
|
if job.Attempts > maxAttempts {
|
|
s.killJob(job, acquisitionId)
|
|
return nil, "", &contract.NoopJobError{State: state}
|
|
}
|
|
|
|
return job, acquisitionId, nil
|
|
}
|
|
|
|
// killJob переводит исчерпавшую попытки задачу в «мертва» и сообщает об этом
|
|
// отправителю. Инвариант «Принятая запись не теряется молча» допускает два
|
|
// исхода — задача пригодна к повтору либо об отказе сказано, — и молчаливая
|
|
// смерть не подходит ни под один.
|
|
func (s *TranscribeService) killJob(job *entity.TranscribeJob, holder string) {
|
|
s.logger.Error("Job exhausted its attempts",
|
|
"job_id", job.Id,
|
|
"state", job.State,
|
|
"attempts", job.Attempts)
|
|
|
|
job.Die(fmt.Sprintf("attempts exhausted: %d", job.Attempts))
|
|
|
|
if err := s.jobRepo.Save(job, holder); err != nil {
|
|
s.logger.Error("Failed to save dead job", "error", err, "job_id", job.Id)
|
|
return
|
|
}
|
|
|
|
s.notify(job, "Не удалось обработать запись: попытки исчерпаны.\nПожалуйста, попробуйте еще раз.")
|
|
}
|
|
|
|
// scheduleRetry снимает захват с отказавшей задачи и ставит нарастающую паузу.
|
|
// Захват, оставленный до конца срока, отложил бы повтор на часы.
|
|
func (s *TranscribeService) scheduleRetry(job *entity.TranscribeJob, holder string, stepErr error) {
|
|
// Шаг, потерявший захват, задачу уже не трогает: ею занят другой.
|
|
var lost *contract.LostAcquisitionError
|
|
if errors.As(stepErr, &lost) {
|
|
return
|
|
}
|
|
|
|
job.RetryAfter(time.Now().Add(retryDelay(job.Attempts)))
|
|
|
|
if err := s.jobRepo.Save(job, holder); err != nil {
|
|
var lostOnSave *contract.LostAcquisitionError
|
|
if errors.As(err, &lostOnSave) {
|
|
return
|
|
}
|
|
s.logger.Error("Failed to schedule job retry", "error", err, "job_id", job.Id)
|
|
}
|
|
}
|
|
|
|
// retryDelay растит паузу с числом попыток до потолка.
|
|
func retryDelay(attempts int) time.Duration {
|
|
if attempts < 1 {
|
|
attempts = 1
|
|
}
|
|
delay := time.Duration(math.Pow(2, float64(attempts-1))) * retryDelayBase
|
|
if delay > retryDelayCap || delay <= 0 {
|
|
return retryDelayCap
|
|
}
|
|
return delay
|
|
}
|
|
|
|
func (s *TranscribeService) completeJob(job *entity.TranscribeJob, holder string, transcriptionText string) error {
|
|
// Обновляем задачу с результатом
|
|
job.Done(transcriptionText)
|
|
|
|
// Сохраняем задачу в базу
|
|
if err := s.jobRepo.Save(job, holder); err != nil {
|
|
s.logger.Error("Failed to save job", "error", err, "job_id", job.Id)
|
|
return fmt.Errorf("failed to save job: %w", err)
|
|
}
|
|
|
|
// Отправляем распознанный текст обратно пользователю
|
|
return s.send(job, transcriptionText)
|
|
}
|
|
|
|
func (s *TranscribeService) failJob(job *entity.TranscribeJob, holder string, jobErr error, humanErrorText string) error {
|
|
// Обновляем задачу с результатом
|
|
job.Fail(jobErr.Error())
|
|
|
|
// Сохраняем задачу в базу
|
|
if err := s.jobRepo.Save(job, holder); err != nil {
|
|
s.logger.Error("Failed to save job", "error", err, "job_id", job.Id)
|
|
return fmt.Errorf("failed to save job: %w", err)
|
|
}
|
|
|
|
errorMessage := fmt.Sprintf("При обработке задачи произошла ошибка: %s.\nПожалуйста, попробуйте еще раз.", humanErrorText)
|
|
return s.send(job, errorMessage)
|
|
}
|
|
|
|
// send отвечает отправителю там, откуда пришла запись, и отказ отправки
|
|
// поднимает вверх: он принадлежит шагу.
|
|
func (s *TranscribeService) send(job *entity.TranscribeJob, text string) error {
|
|
if job.Source != entity.SourceTelegram {
|
|
return nil
|
|
}
|
|
|
|
if job.TgChatId == nil {
|
|
s.logger.Error("Telegram chat not specified", "job_id", job.Id)
|
|
return fmt.Errorf("tg chat id not specified, job id: %s", job.Id)
|
|
}
|
|
|
|
if err := s.tgSender.Send(text, *job.TgChatId, job.TgReplyMessageId); err != nil {
|
|
s.logger.Error("Failed to sent message to client", "job_id", job.Id)
|
|
return fmt.Errorf("failed to sent message to client, job id: %s, err: %w", job.Id, err)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// notify отвечает отправителю там, где поднимать отказ некуда: задача уже
|
|
// доведена до конца, и отказ отправки остаётся записью в журнале владельца.
|
|
func (s *TranscribeService) notify(job *entity.TranscribeJob, text string) {
|
|
if err := s.send(job, text); err != nil {
|
|
s.logger.Error("Failed to notify sender", "error", err, "job_id", job.Id)
|
|
}
|
|
}
|
|
|
|
// closeWork убирает рабочую копию. Отказ уборки не роняет шаг, но и не
|
|
// проглатывается: забытая копия это шестичасовая запись во временном каталоге.
|
|
func (s *TranscribeService) closeWork(work contract.WorkFile) {
|
|
if err := work.Close(); err != nil {
|
|
s.logger.Error("Failed to remove work file", "error", err)
|
|
}
|
|
}
|