Files
transcriber/internal/service/transcribe.go
T
av 1576d06735 внутренняя модель перестроена вокруг аудиозаписи
- audiorecords вместо transcribe_jobs: приложения (texts, structures,
  recognitions, record_events, topics) живут своими коллекциями, ссылки на
  исходник и на приведённую копию перестали переставляться
- рубеж называет достигнутое, отказ стал признаком остановки с причиной, а
  сторожей стало двое: число отказов и время в рубеже
- воркеры потеряли специализацию, их число задаётся [pipeline] workers, шаг
  выбирается по рубежу, а захват отдаёт идентификатор и признак захвата
2026-08-14 20:20:33 +03:00

984 lines
45 KiB
Go

package service
import (
"context"
"errors"
"fmt"
"io"
"log/slog"
"math"
"path/filepath"
"strings"
"time"
"github.com/google/uuid"
"git.vakhrushev.me/av/transcriber/internal/clock"
"git.vakhrushev.me/av/transcriber/internal/contract"
"git.vakhrushev.me/av/transcriber/internal/entity"
"git.vakhrushev.me/av/transcriber/internal/metrics"
)
const (
defaultAudioExt = "audio"
// Предел отказов. Число обратимо и живёт здесь одним местом; счётчик растёт
// при захвате и обнуляется на шаге, завершившемся без отказа либо отложившем
// работу.
maxAttempts = 5
// Пауза перед повтором отказавшей записи растёт с числом отказов до потолка.
// Ожидание чужой операции этой паузой не выражается — у него своя задержка
// числом, и отказов оно не тратит.
retryDelayBase = time.Second
retryDelayCap = 5 * time.Minute
// Задержки опроса операции распознавания. Числа, а не функция числа отказов:
// счётчик на ожидании обнулён, и выведенная из него пауза выродилась бы в
// своё наименьшее значение, учащая опрос платного сервиса.
firstCheckDelay = 10 * time.Second
nextCheckDelay = 5 * time.Second
)
// Имена шагов. Идут меткой в журнал событий записи и в счётчик работы воркера
// вместе с рубежом, с которого запись взята.
const (
stepNormalize = "normalize"
stepSubmit = "submit"
stepPoll = "poll"
stepFinish = "finish"
)
// Repositories — хранилища, с которыми работает конвейер. Собраны структурой, а
// не восемью доводами: перечень растёт с моделью, а порядок восьми одинаковых
// указателей в вызове перепутать нечем, кроме внимательности.
type Repositories struct {
Records contract.AudioRecordRepository
Files contract.FileRepository
Texts contract.TextRepository
Structures contract.StructureRepository
Recognitions contract.RecognitionRepository
Events contract.RecordEventRepository
}
type TranscribeService struct {
repos Repositories
metaviewer contract.AudioMetaViewer
converter contract.AudioFileConverter
recognizer contract.AudioRecognizer
tgSender contract.TelegramMessageSender
limits entity.StuckLimits
logger *slog.Logger
}
func NewTranscribeService(
repos Repositories,
metaviewer contract.AudioMetaViewer,
converter contract.AudioFileConverter,
recognizer contract.AudioRecognizer,
tgSender contract.TelegramMessageSender,
limits entity.StuckLimits,
logger *slog.Logger,
) *TranscribeService {
if logger == nil {
logger = slog.Default()
}
return &TranscribeService{
repos: repos,
metaviewer: metaviewer,
converter: converter,
recognizer: recognizer,
tgSender: tgSender,
limits: limits,
logger: logger,
}
}
// stepOutcome — чем кончился шаг.
//
// Исход называется **явно**, а не выводится из «шаг не вернул ошибку»: приговор
// шага и откладывание опроса оба возвращают отсутствие отказа, и по одному лишь
// `nil` они неотличимы от сделанной работы. Пока их различал `nil`, остановленная
// запись писала в журнал `done` вслед за `halted` и растила счётчик успехов, а
// часовое ожидание чужой операции оставляло в журнале сотни строк ни о чём.
type stepOutcome int
const (
// outcomeDone — шаг сделал работу и двинул рубеж.
outcomeDone stepOutcome = iota
// outcomePostponed — шаг отложил работу: рубеж прежний, писать нечего.
outcomePostponed
// outcomeHalted — шаг вынес приговор и остановил запись. Строку журнала и
// счётчик отказа ставит сама остановка.
outcomeHalted
)
// step — шаг конвейера. Держит захват и записывает результат только по нему.
type step func(ctx context.Context, r *entity.AudioRecord, holder string) (stepOutcome, error)
// stepFor выбирает шаг по рубежу записи — таблицей, а не тем, какой воркер
// пришёл. Воркеры одинаковы и не привязаны к шагу, поэтому решение принимается
// здесь и по одному признаку.
func (s *TranscribeService) stepFor(state string) (string, step, bool) {
switch state {
case entity.StateUploaded:
return stepNormalize, s.normalize, true
case entity.StateNormalized:
return stepSubmit, s.submit, true
case entity.StateSubmitted:
return stepPoll, s.poll, true
case entity.StateTranscribed:
return stepFinish, s.finish, true
default:
return "", nil, false
}
}
func (s *TranscribeService) CreateJobFromTelegram(ctx context.Context, file io.Reader, fileName string, chatId int64, replyMsgId int) (*entity.AudioRecord, error) {
record := &entity.AudioRecord{
State: entity.StateUploaded,
Source: entity.SourceTelegram,
TgChatId: &chatId,
TgReplyMessageId: &replyMsgId,
}
return s.createRecord(ctx, record, file, fileName)
}
// CreateJobFromApi заводит запись от имени вошедшего. Владелец обязателен:
// пустой отвергается здесь, потому что колонка владельца допускает пустое
// значение ради записей бота, и приём по HTTP — то место, где обязательность
// держится.
//
// Отказ этот — последний рубеж, а не первый: предъявителя без учётной записи
// пользователя транспорт отвергает раньше, до чтения тела. Здесь он остаётся на
// случай нового вызывающего, который такой проверки не поставит.
func (s *TranscribeService) CreateJobFromApi(ctx context.Context, file io.Reader, fileName, ownerID string) (*entity.AudioRecord, error) {
if ownerID == "" {
s.logger.Error("Refusing to create record without owner")
return nil, contract.ErrOwnerRequired
}
record := &entity.AudioRecord{
State: entity.StateUploaded,
Source: entity.SourceApi,
OwnerID: &ownerID,
}
return s.createRecord(ctx, record, file, fileName)
}
func (s *TranscribeService) createRecord(ctx context.Context, r *entity.AudioRecord, file io.Reader, fileName string) (*entity.AudioRecord, 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.repos.Files.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 audio record", "file_ext", ext)
info, err := s.metaviewer.GetInfo(ctx, 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
}
meta := contract.FileMeta{
Format: formatOf(ext),
DurationMs: int64(info.Seconds) * 1000,
}
fileRecord, err := s.repos.Files.Create(storageFileName, work, meta, ownerOf(r))
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)
// Ссылка на принятую копию ставится один раз и больше не переставляется:
// исходник остаётся доступным после того, как запись прошла конвейер.
r.OriginalFileID = &fileRecord.Id
r.StateEnteredAt = clock.Now()
if err := s.repos.Records.Create(r); err != nil {
s.logger.Error("Failed to create audio record", "error", err, "file_id", fileRecord.Id)
return nil, err
}
s.logger.Info("Audio record created successfully", "record_id", r.Id, "file_id", fileRecord.Id)
return r, nil
}
// RunStep забирает любую пригодную к работе запись и отдаёт её шагу, выбранному
// по рубежу.
//
// Воркер к шагу не привязан: он не знает заранее, что вытянет, и потому срок
// протухания захвата приезжает с рубежом, а не с ним.
//
// Отказ шага не оставляет запись захваченной до конца срока: захват снимается, и
// запись ждёт нарастающую паузу — иначе повтор наступал бы через часы, а не
// через секунду.
//
// Контекст доходит до шага, а через него — до внешнего собеседника: остановка
// сервиса убивает `ffmpeg` и обрывает запрос к распознаванию. Прерванный шаг
// приговора не выносит: запись остаётся пригодной к повтору, отказа не тратит и
// отправителю о несуществующем сбое не сообщает.
func (s *TranscribeService) RunStep(ctx context.Context) error {
// Нас уже остановили — запись не забираем: захват стоил бы ей отказа, а
// работы всё равно не будет. Исход «шаг не сделал ничего» — это
// `NoopJobError` по смыслу, и он же не поднимает уровень и не считается в
// метрику.
if ctx.Err() != nil {
return &contract.NoopJobError{}
}
record, holder, err := s.acquire()
if err != nil {
return err
}
name, run, ok := s.stepFor(record.State)
if !ok {
// Рубеж, для которого шага нет, — это порча записи либо забытая правка
// таблицы. Молчать нельзя: запись выпала бы из работы без единого следа.
s.logger.Error("No step declared for state", "record_id", record.Id, "state", record.State)
s.halt(record, holder, record.State, entity.HaltReasonStepFailed,
fmt.Sprintf("no step for state %s", record.State),
"сервис не знает, что делать с этой записью")
return &contract.NoopJobError{State: record.State}
}
// Рубеж запоминается **до** шага: шаг двигает его на свой результат, и метка,
// снятая после, называла бы уже следующий рубеж — «падает приведение» и
// «падает распознавание» перестали бы различаться ровно там, где владелец
// сервиса на это смотрит.
stage := record.State
started := clock.Start()
outcome, stepErr := run(ctx, record, holder)
duration := time.Since(started)
stopped := stepErr != nil && ctx.Err() != nil
if stepErr != nil {
if !stopped {
metrics.WorkerJobCounter.WithLabelValues(stage, labelFailed).Inc()
}
s.scheduleRetry(record, holder, stepErr)
return stepErr
}
switch outcome {
case outcomeDone:
metrics.WorkerJobCounter.WithLabelValues(stage, labelOk).Inc()
s.appendEvent(record, entity.EventOriginPipeline, name, entity.EventOutcomeDone, "", duration)
case outcomePostponed:
// Откладывание работой не является: рубеж прежний, и строка журнала о
// нём была бы строкой ни о чём — часовое ожидание чужой операции дало бы
// их сотни. Счётчик растёт: прогон состоялся, и его исход — успех.
metrics.WorkerJobCounter.WithLabelValues(stage, labelOk).Inc()
case outcomeHalted:
// Строку журнала и счётчик отказа поставила сама остановка: она знает
// причину, а `RunStep` — только то, что отказа не было.
}
return nil
}
// acquire забирает запись, читает её колонки и проверяет обоих сторожей.
//
// Захват отдаёт идентификатор и признак **этого** захвата; колонки читаются
// отдельным чтением. Прежде захват возвращал перечень колонок, и всякая новая
// колонка записи попадала под инвариант проекта о колонках очереди — теперь
// перечня нет вовсе.
func (s *TranscribeService) acquire() (*entity.AudioRecord, string, error) {
acquired, err := s.repos.Records.FindAndAcquire(entity.WorkingStages())
if err != nil {
// Признак узнаётся по смыслу: репозиторий вправе обернуть свой отказ
// пояснением, и приведение типа от этого сломалось бы молча.
var notFound *contract.JobNotFoundError
if errors.As(err, &notFound) {
return nil, "", &contract.NoopJobError{}
}
s.logger.Error("Failed to find and acquire a record", "error", err)
return nil, "", fmt.Errorf("failed to find and acquire a record: %w", err)
}
record, err := s.repos.Records.Get(acquired.ID)
if err != nil {
s.logger.Error("Failed to read acquired record", "error", err, "record_id", acquired.ID)
return nil, "", fmt.Errorf("failed to read acquired record: %w", err)
}
// Сторож отказов: повторы внутри шага.
if record.Attempts > maxAttempts {
s.logger.Error("Record exhausted its attempts",
"record_id", record.Id, "state", record.State, "attempts", record.Attempts)
s.halt(record, acquired.Holder, record.State, entity.HaltReasonAttempts,
fmt.Sprintf("attempts exhausted: %d", record.Attempts),
"попытки исчерпаны")
return nil, "", &contract.NoopJobError{State: record.State}
}
// Сторож времени: застревание в рубеже. Считает от входа в рубеж, и
// откладывание опроса его не двигает — иначе запись, чью чужую операцию
// опрашивают раз в несколько секунд, не достигла бы предела никогда.
if s.isStuck(record) {
s.logger.Error("Record is stuck in its state",
"record_id", record.Id, "state", record.State,
"state_entered_at", record.StateEnteredAt)
s.halt(record, acquired.Holder, record.State, entity.HaltReasonStuck,
fmt.Sprintf("stuck in %s", record.State),
"обработка застряла")
return nil, "", &contract.NoopJobError{State: record.State}
}
return record, acquired.Holder, nil
}
// isStuck — простояла ли запись в рубеже дольше предела. У конечного рубежа
// предела нет: стоять в нём запись будет вечно по построению.
func (s *TranscribeService) isStuck(r *entity.AudioRecord) bool {
stage, ok := entity.StageByName(r.State)
if !ok {
return false
}
limit, hasLimit := stage.Limit(s.limits)
if !hasLimit || limit <= 0 || r.StateEnteredAt.IsZero() {
return false
}
return clock.Now().Sub(r.StateEnteredAt) > limit
}
// normalize приводит принятую копию к рабочему формату.
//
// Ссылка на исходник не переставляется: у записи две ссылки, и обе живут до
// конца.
func (s *TranscribeService) normalize(ctx context.Context, r *entity.AudioRecord, holder string) (stepOutcome, error) {
s.logger.Info("Starting normalize step", "record_id", r.Id)
if r.OriginalFileID == nil {
s.logger.Error("Record has no original file", "record_id", r.Id)
return s.failStep(r, holder, stepNormalize, errors.New("record has no original file"), "у записи нет файла")
}
srcFile, err := s.repos.Files.GetByID(*r.OriginalFileID)
if err != nil {
s.logger.Error("Failed to get original file", "error", err, "file_id", *r.OriginalFileID)
return outcomeDone, err
}
srcFormat := srcFile.Format
if srcFormat == "" {
srcFormat = defaultAudioExt
}
src, err := s.repos.Files.Localize(*r.OriginalFileID)
if err != nil {
s.logger.Error("Failed to localize original file", "error", err, "file_id", *r.OriginalFileID)
return outcomeDone, err
}
defer s.closeWork(src)
dest, err := s.repos.Files.StageEmpty(".ogg")
if err != nil {
s.logger.Error("Failed to stage normalized file", "error", err, "record_id", r.Id)
return outcomeDone, err
}
defer s.closeWork(dest)
s.logger.Info("Converting file", "record_id", r.Id, "src_format", srcFormat)
startTime := clock.Start()
err = s.converter.Convert(ctx, src.Path(), dest.Path())
conversionDuration := time.Since(startTime)
metrics.ObserveConversionDuration(srcFormat, "ogg", err != nil, conversionDuration.Seconds())
if err != nil {
// Остановка сервиса — не приговор записи. Убитый по контексту `ffmpeg`
// отдаёт `signal: killed`, и от настоящего отказа конвертации
// (`exit status N`) эта ошибка неотличима ни типом, ни `errors.Is`:
// различает их только контекст шага. Без этой развилки каждый деплой
// останавливал бы конвертируемую запись.
if ctxErr := ctx.Err(); ctxErr != nil {
s.logger.Info("File conversion interrupted by shutdown",
"record_id", r.Id, "duration", conversionDuration)
return outcomeDone, fmt.Errorf("conversion interrupted: %w", ctxErr)
}
s.logger.Error("File conversion failed",
"error", err, "record_id", r.Id, "duration", conversionDuration)
return s.failStep(r, holder, stepNormalize, err, "сбой конвертации файла")
}
destSize, err := dest.Size()
if err != nil {
s.logger.Error("Failed to measure normalized file", "error", err, "record_id", r.Id)
return outcomeDone, err
}
s.logger.Info("File conversion completed",
"record_id", r.Id, "duration", conversionDuration, "output_size", destSize)
metrics.OutputFileSizeHistogram.WithLabelValues("ogg").Observe(float64(destSize))
destFileName := fmt.Sprintf("%s%s", uuid.NewString(), ".ogg")
destMeta := contract.FileMeta{Format: "ogg", DurationMs: srcFile.DurationMs}
destFileRecord, err := s.repos.Files.Create(destFileName, dest, destMeta, ownerOf(r))
if err != nil {
s.logger.Error("Failed to create normalized file record", "error", err, "record_id", r.Id)
return outcomeDone, err
}
// Ссылка ставится только после того, как запись о файле есть: иначе повтор
// оставил бы запись указывающей на файл, которого нет.
r.NormalizedFileID = &destFileRecord.Id
r.MoveToState(entity.StateNormalized)
if err := s.repos.Records.Save(r, holder); err != nil {
s.logger.Error("Failed to save record", "error", err, "record_id", r.Id)
return outcomeDone, err
}
s.logger.Info("Normalize step completed successfully", "record_id", r.Id)
return outcomeDone, nil
}
// submit кладёт аудио туда, откуда провайдер его прочитает, и заводит операцию
// распознавания.
//
// Заливка и отправка разделены: повтор заливки бесплатен, повтор отправки
// оплачивается наружу. Строка попытки заводится **до** обращения к провайдеру —
// окно между его ответом и записью идентификатора это то место, где теряется
// оплаченное.
func (s *TranscribeService) submit(ctx context.Context, r *entity.AudioRecord, holder string) (stepOutcome, error) {
s.logger.Info("Starting submit step", "record_id", r.Id)
if r.NormalizedFileID == nil {
s.logger.Error("Record has no normalized file", "record_id", r.Id)
return s.failStep(r, holder, stepSubmit, errors.New("record has no normalized file"), "у записи нет приведённого файла")
}
fileRecord, err := s.repos.Files.GetByID(*r.NormalizedFileID)
if err != nil {
s.logger.Error("Failed to get normalized file", "error", err, "file_id", *r.NormalizedFileID)
return outcomeDone, err
}
attempt, err := s.recognitionAttempt(r, holder)
if err != nil {
return outcomeDone, err
}
// Работа уже сделана — второй раз наружу не платим.
if attempt.ExternalID != "" {
s.logger.Info("Recognition operation already submitted", "record_id", r.Id)
return s.awaitOperation(r, holder)
}
sourceURI, err := s.uploadSource(ctx, r, attempt, fileRecord)
if err != nil {
return outcomeDone, err
}
operationID, err := s.recognizer.Submit(ctx, sourceURI)
if err != nil {
if ctxErr := ctx.Err(); ctxErr != nil {
s.logger.Info("Recognition submit interrupted by shutdown", "record_id", r.Id)
return outcomeDone, fmt.Errorf("recognition submit interrupted: %w", ctxErr)
}
s.logger.Error("Failed to submit recognition", "error", err, "record_id", r.Id)
return outcomeDone, err
}
if err := s.repos.Recognitions.Submitted(attempt.Id, sourceURI, operationID); err != nil {
// Идентификатор операции получен, а записать его не вышло: следующая
// попытка оплатит ту же запись второй раз. Уровень здесь владельцу
// сервиса, а не отправителю.
s.logger.Error("Failed to store operation id", "error", err, "record_id", r.Id)
return outcomeDone, err
}
s.logger.Info("Recognition submitted", "record_id", r.Id, "operation_id", operationID)
return s.awaitOperation(r, holder)
}
// awaitOperation двигает запись на рубеж ожидания и назначает первую задержку
// опроса.
func (s *TranscribeService) awaitOperation(r *entity.AudioRecord, holder string) (stepOutcome, error) {
r.MoveToState(entity.StateSubmitted)
r.DelayTime = ptr(clock.Now().Add(firstCheckDelay))
if err := s.repos.Records.Save(r, holder); err != nil {
s.logger.Error("Failed to save record", "error", err, "record_id", r.Id)
return outcomeDone, err
}
return outcomeDone, nil
}
// recognitionAttempt возвращает строку попытки распознавания, заводя её при
// первом проходе.
//
// Строка заводится **до** обращения к провайдеру и сохраняется на записи тем же
// движением: без этого прерванный между заведением и сохранением шаг завёл бы на
// повторе вторую попытку, и первая осталась бы сиротой — а по ней он и узнаёт,
// что за запись уже заплачено.
func (s *TranscribeService) recognitionAttempt(r *entity.AudioRecord, holder string) (*entity.Recognition, error) {
if r.RecognitionID != nil {
attempt, err := s.repos.Recognitions.GetByID(*r.RecognitionID)
if err != nil {
s.logger.Error("Failed to read recognition attempt", "error", err, "record_id", r.Id)
return nil, err
}
return attempt, nil
}
attempt := &entity.Recognition{
RecordID: r.Id,
Provider: s.recognizer.Provider(),
Model: s.recognizer.Model(),
}
if err := s.repos.Recognitions.Create(attempt); err != nil {
s.logger.Error("Failed to create recognition attempt", "error", err, "record_id", r.Id)
return nil, err
}
r.RecognitionID = &attempt.Id
if err := s.repos.Records.Save(r, holder); err != nil {
s.logger.Error("Failed to link recognition attempt", "error", err, "record_id", r.Id)
return nil, err
}
return attempt, nil
}
// uploadSource кладёт приведённую копию туда, откуда провайдер её прочитает, —
// если её там ещё нет, — и отдаёт адрес.
//
// Повтор заливки бесплатен, но дорог по времени на многочасовой записи, поэтому
// шаг сперва смотрит на наблюдаемый признак сделанного.
//
// Адрес уже уложенного берётся из **строки попытки**, а не пересчитывается: он
// сохранён тем шагом, который заливал, и пересчёт разошёлся бы с сохранённым
// молча при первой же смене правила именования объекта или адреса хранилища.
func (s *TranscribeService) uploadSource(
ctx context.Context,
r *entity.AudioRecord,
attempt *entity.Recognition,
file *entity.File,
) (string, error) {
objectKey := file.FileName
exists, err := s.recognizer.ObjectExists(ctx, objectKey, file.Size)
if err != nil {
s.logger.Error("Failed to check uploaded object", "error", err, "record_id", r.Id)
return "", err
}
if exists && attempt.SourceURI != "" {
s.logger.Info("Audio is already uploaded, skipping", "record_id", r.Id)
return attempt.SourceURI, nil
}
content, err := s.repos.Files.Open(*r.NormalizedFileID)
if err != nil {
s.logger.Error("Failed to open normalized file", "error", err, "file_id", *r.NormalizedFileID)
return "", err
}
defer func() {
if err := content.Close(); err != nil {
s.logger.Error("Failed to close file", "error", err, "file_id", *r.NormalizedFileID)
}
}()
sourceURI, err := s.recognizer.Upload(ctx, content, objectKey)
if err != nil {
if ctxErr := ctx.Err(); ctxErr != nil {
s.logger.Info("Upload interrupted by shutdown", "record_id", r.Id)
return "", fmt.Errorf("upload interrupted: %w", ctxErr)
}
s.logger.Error("Failed to upload audio for recognition", "error", err, "record_id", r.Id)
return "", err
}
return sourceURI, nil
}
// poll опрашивает операцию у провайдера.
//
// Операция ещё идёт — работа **откладывается**, а не переводится в тот же рубеж:
// переходом это никогда не было, и именно мнимость перехода прежде обнуляла
// сторожа времени.
func (s *TranscribeService) poll(ctx context.Context, r *entity.AudioRecord, holder string) (stepOutcome, error) {
if r.RecognitionID == nil {
s.logger.Error("Record has no recognition attempt", "record_id", r.Id)
return s.failStep(r, holder, stepPoll, errors.New("record has no recognition attempt"), "сведений о распознавании нет")
}
attempt, err := s.repos.Recognitions.GetByID(*r.RecognitionID)
if err != nil {
s.logger.Error("Failed to read recognition attempt", "error", err, "record_id", r.Id)
return outcomeDone, err
}
if attempt.ExternalID == "" {
s.logger.Error("Recognition attempt has no operation id", "record_id", r.Id)
return s.failStep(r, holder, stepPoll, errors.New("recognition attempt has no operation id"), "сведений о распознавании нет")
}
// Опрос идёт раз в несколько секунд всё время распознавания: часовая запись
// дала бы полторы тысячи строк. Конвенция журнала относит проверку готовности
// операции к отладочному уровню поимённо.
s.logger.Debug("Checking operation status", "record_id", r.Id, "operation_id", attempt.ExternalID)
result, err := s.recognizer.CheckStatus(ctx, attempt.ExternalID)
if err != nil {
if ctxErr := ctx.Err(); ctxErr != nil {
s.logger.Info("Status check interrupted by shutdown", "record_id", r.Id)
return outcomeDone, fmt.Errorf("status check interrupted: %w", ctxErr)
}
s.logger.Error("Failed to check recognition status", "error", err, "operation_id", attempt.ExternalID)
return outcomeDone, err
}
if result.IsInProgress() {
// Шаг отработал без отказа: ожидание чужой операции отказом не является.
// Задержка своя, числом; рубеж и время входа в него не трогаются.
s.logger.Debug("Operation in progress", "record_id", r.Id, "operation_id", attempt.ExternalID)
r.Postpone(clock.Now().Add(nextCheckDelay))
if err := s.repos.Records.Save(r, holder); err != nil {
s.logger.Error("Failed to save record", "error", err, "record_id", r.Id)
return outcomeDone, err
}
return outcomePostponed, nil
}
if result.IsFailed() {
errorText := result.GetError()
s.logger.Error("Operation failed",
"record_id", r.Id, "operation_id", attempt.ExternalID, "error_message", errorText)
return s.failStep(r, holder, stepPoll, errors.New(errorText), "сбой при распознавании файла")
}
outcome, err := s.recognizer.Fetch(ctx, attempt.ExternalID)
if err != nil {
if ctxErr := ctx.Err(); ctxErr != nil {
s.logger.Info("Result fetch interrupted by shutdown", "record_id", r.Id)
return outcomeDone, fmt.Errorf("result fetch interrupted: %w", ctxErr)
}
s.logger.Error("Failed to fetch recognition result", "error", err, "operation_id", attempt.ExternalID)
return outcomeDone, err
}
s.logger.Info("Recognition completed",
"record_id", r.Id,
"operation_id", attempt.ExternalID,
"text_length", len(outcome.PlainText),
"replicas", len(outcome.Replicas))
// Сырой ответ сохраняется целиком: результат операции у провайдера не
// переспрашивается, и когда мы научимся размечать говорящих, архив
// пересчитается из сохранённого без единого рубля.
if err := s.repos.Recognitions.Finish(attempt.Id, outcome.Raw); err != nil {
s.logger.Error("Failed to store provider payload", "error", err, "record_id", r.Id)
return outcomeDone, err
}
if err := s.storeOutcome(r, outcome); err != nil {
return outcomeDone, err
}
r.MoveToState(entity.StateTranscribed)
if err := s.repos.Records.Save(r, holder); err != nil {
s.logger.Error("Failed to save record", "error", err, "record_id", r.Id)
return outcomeDone, err
}
return outcomeDone, nil
}
// storeOutcome кладёт расшифровку и структуру реплик. Обе строки уникальны по
// своей паре, поэтому повтор прерванного шага второго комплекта не заводит.
func (s *TranscribeService) storeOutcome(r *entity.AudioRecord, outcome *entity.RecognitionOutcome) error {
text, err := s.repos.Texts.Put(r.Id, entity.TextKindTranscript, outcome.PlainText)
if err != nil {
s.logger.Error("Failed to store transcript", "error", err, "record_id", r.Id)
return err
}
r.TranscriptTextID = &text.Id
structure, err := s.repos.Structures.Put(r.Id, entity.StructureVersion, outcome.Replicas)
if err != nil {
s.logger.Error("Failed to store structure", "error", err, "record_id", r.Id)
return err
}
r.StructureID = &structure.Id
return nil
}
// finish отвечает отправителю и доводит запись до конечного рубежа. Доставка —
// хвост последнего шага, а не отдельный узел конвейера.
func (s *TranscribeService) finish(ctx context.Context, r *entity.AudioRecord, holder string) (stepOutcome, error) {
text := "Ой, кажется, на аудиозаписи нет текста."
if r.TranscriptTextID != nil {
stored, err := s.repos.Texts.GetByID(*r.TranscriptTextID)
if err != nil {
s.logger.Error("Failed to read transcript", "error", err, "record_id", r.Id)
return outcomeDone, err
}
if stored.Contents != "" {
text = stored.Contents
}
}
r.MoveToState(entity.StateDone)
if err := s.repos.Records.Save(r, holder); err != nil {
s.logger.Error("Failed to save record", "error", err, "record_id", r.Id)
return outcomeDone, err
}
return outcomeDone, s.send(r, text)
}
// failStep останавливает запись приговором шага и сообщает отправителю.
//
// Исход возвращается **явно**: остановка отказом шага не является — шаг
// рассудил об этой записи окончательно, — но и работой она не была, и
// засчитывать её успехом нельзя.
func (s *TranscribeService) failStep(r *entity.AudioRecord, holder, step string, stepErr error, humanText string) (stepOutcome, error) {
s.halt(r, holder, step, entity.HaltReasonStepFailed, stepErr.Error(),
fmt.Sprintf("При обработке записи произошла ошибка: %s", humanText))
return outcomeHalted, nil
}
// halt ставит признак остановки, считает её отказом, пишет строку журнала
// событий и сообщает отправителю.
//
// Сообщение уходит при **любой** причине остановки: инвариант проекта «Принятая
// запись не теряется молча» допускает два исхода — запись пригодна к повтору
// либо об отказе сказано, — а остановленная запись захвату не выдаётся, значит
// первый исход исключён.
//
// Счётчик растит **сама остановка**, а не воркер, и это не стилистика.
// Остановка по сторожам наступает в захвате, до всякого шага, и воркер о ней
// узнаёт признаком «работы нет» — а считать его в метрику запрещено инвариантом
// «`NoopJobError` — не ошибка». Значит единственное место, где известны и факт
// остановки, и рубеж, — здесь.
func (s *TranscribeService) halt(r *entity.AudioRecord, holder, step, reason, errText, humanText string) {
// Рубеж читается до остановки: она его не двигает, но читать состояние после
// мутации — привычка, из-за которой метка счётчика уже однажды разъехалась.
stage := r.State
r.Halt(reason, errText)
if err := s.repos.Records.Save(r, holder); err != nil {
var lost *contract.LostAcquisitionError
if errors.As(err, &lost) {
return
}
s.logger.Error("Failed to halt record", "error", err, "record_id", r.Id)
return
}
metrics.WorkerJobCounter.WithLabelValues(stage, labelFailed).Inc()
// Приговор шага и приговор сторожа — разные события, и журнал их различает:
// первый вынес шаг, рассудивший об этой записи, второй — счётчик, у которого
// кончилось терпение.
outcome := entity.EventOutcomeHalted
if reason == entity.HaltReasonStepFailed {
outcome = entity.EventOutcomeFailed
}
s.appendEvent(r, entity.EventOriginPipeline, step, outcome, reason, 0)
s.notify(r, humanText+"\nПожалуйста, попробуйте еще раз.")
}
// scheduleRetry снимает захват с отказавшей записи и ставит нарастающую паузу.
// Захват, оставленный до конца срока, отложил бы повтор на часы.
func (s *TranscribeService) scheduleRetry(r *entity.AudioRecord, holder string, stepErr error) {
// Шаг, потерявший захват, запись уже не трогает: ею занят другой.
var lost *contract.LostAcquisitionError
if errors.As(stepErr, &lost) {
return
}
// Остановка сервиса отказа не тратит: запись не виновата в том, что нас
// перезапустили. Счётчик растёт при захвате, поэтому здесь его возвращают
// назад — иначе пять выкладок подряд останавливают живую запись с приговором
// «попытки исчерпаны».
if errors.Is(stepErr, context.Canceled) || errors.Is(stepErr, context.DeadlineExceeded) {
if r.Attempts > 0 {
r.Attempts--
}
}
r.RetryAfter(clock.Now().Add(retryDelay(r.Attempts)))
if err := s.repos.Records.Save(r, holder); err != nil {
var lostOnSave *contract.LostAcquisitionError
if errors.As(err, &lostOnSave) {
return
}
s.logger.Error("Failed to schedule record retry", "error", err, "record_id", r.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
}
// appendEvent пишет строку журнала событий записи. Отказ записи журнала шаг не
// роняет: работа сделана, а журнал никем не читается ради решения.
func (s *TranscribeService) appendEvent(r *entity.AudioRecord, origin, step, outcome, outcomeText string, duration time.Duration) {
event := &entity.RecordEvent{
RecordID: r.Id,
Origin: origin,
Step: step,
Outcome: outcome,
OutcomeText: outcomeText,
DurationMs: duration.Milliseconds(),
}
if err := s.repos.Events.Append(event); err != nil {
s.logger.Error("Failed to append record event", "error", err, "record_id", r.Id)
}
}
// send отвечает отправителю там, откуда пришла запись. Отказ отправки поднимает
// вверх: он принадлежит шагу.
//
// Кроме недоставки — её шаг записывает и завершается без отказа. Ответ уходит
// после того, как достигнутый рубеж сохранён: работа к этой минуте сделана, и
// объявленный отказ засчитался бы воркеру сбоем и лёг бы владельцу записью
// отказа. Повтор делу не помогает — ни бот, ни адресат от ожидания не появятся,
// — поэтому причина недоставки живёт в журнале, а не в рубеже записи.
func (s *TranscribeService) send(r *entity.AudioRecord, text string) error {
if r.Source != entity.SourceTelegram {
return nil
}
// Адресата у записи нет: отвечать некуда, и повторять нечего. Уровень здесь
// выше, чем у неподнятого канала, и это не педантизм: пустой чат у записи
// из Telegram — симптом порчи записи.
if r.TgChatId == nil {
s.undelivered(r, slog.LevelError, "chat is not specified")
return nil
}
if err := s.tgSender.Send(text, *r.TgChatId, r.TgReplyMessageId); err != nil {
// Канал не поднят: сервис работает без этого входа, и это объявленный
// режим, а не поломка.
if errors.Is(err, contract.ErrDeliveryChannelDown) {
s.undelivered(r, slog.LevelWarn, "delivery channel is down")
return nil
}
s.logger.Error("Failed to sent message to client", "record_id", r.Id)
return fmt.Errorf("failed to sent message to client, record id: %s, err: %w", r.Id, err)
}
return nil
}
// undelivered записывает недоставленный ответ и считает его в метрику. Уровень
// приходит от причины: объявленный режим — «может стать проблемой», порча
// записи — событие для разбора.
//
// Идентификатор записи обязателен, иначе владелец видит, что ответ не ушёл, но
// не может найти, чей; текста ответа в записи нет — он содержимое чужой записи.
func (s *TranscribeService) undelivered(r *entity.AudioRecord, level slog.Level, reason string) {
metrics.UndeliveredReplyCounter.WithLabelValues(reason).Inc()
switch level {
case slog.LevelError:
s.logger.Error(undeliveredMessage, "record_id", r.Id, "reason", reason)
default:
s.logger.Warn(undeliveredMessage, "record_id", r.Id, "reason", reason)
}
}
const undeliveredMessage = "Reply was not delivered"
// notify отвечает отправителю там, где поднимать отказ некуда: запись уже
// доведена до своего исхода, и отказ отправки остаётся записью в журнале.
func (s *TranscribeService) notify(r *entity.AudioRecord, text string) {
if err := s.send(r, text); err != nil {
s.logger.Error("Failed to notify sender", "error", err, "record_id", r.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)
}
}
// ownerOf — владелец записи строкой; пустая значит «владельца нет», и таковы
// записи, принятые ботом. Файл наследует владельца своей записи: правило
// просмотра коллекции файлов сужено этой колонкой.
func ownerOf(r *entity.AudioRecord) string {
if r.OwnerID == nil {
return ""
}
return *r.OwnerID
}
// formatOf приводит расширение к виду колонки формата: без точки, в нижнем
// регистре. Наружу оно выходит только приведённым к перечню известных форматов —
// это делает метка метрики.
func formatOf(ext string) string {
return strings.ToLower(strings.TrimPrefix(ext, "."))
}
// Метки счётчика работы. Строками, а не приведением булева: значение метки —
// часть наблюдаемой поверхности, и опечатка в ней разводит один ряд на два.
const (
labelOk = "false"
labelFailed = "true"
)
func ptr[T any](v T) *T { return &v }