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) } }