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" // maxExtLen — потолок длины расширения вместе с точкой. // // Сторож от патологии, а не перечень допустимого: расширения известных // форматов укладываются в пять знаков, и щедрый потолок ничего у отправителя // не отнимает. Он нужен против другого — имени `x.` с четырьмястами знаками // после точки: оно роняет заведение временного файла, и отправитель получает // `500` на входе, за который отвечает сам, а владелец сервиса — строку `ERROR`, // неотличимую от аварии хранилища. maxExtLen = 32 // Предел отказов. Число обратимо и живёт здесь одним местом; счётчик растёт // при захвате и обнуляется на шаге, завершившемся без отказа либо отложившем // работу. 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 limits entity.StuckLimits logger *slog.Logger } func NewTranscribeService( repos Repositories, metaviewer contract.AudioMetaViewer, converter contract.AudioFileConverter, recognizer contract.AudioRecognizer, limits entity.StuckLimits, logger *slog.Logger, ) *TranscribeService { if logger == nil { logger = slog.Default() } return &TranscribeService{ repos: repos, metaviewer: metaviewer, converter: converter, recognizer: recognizer, 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 } } // CreateJobFromApi заводит запись от имени вошедшего. Владелец обязателен, и // обязательность эту держит схема хранилища: колонка владельца пустого значения // не принимает. Отказ стоит и здесь, раньше схемы, потому что отвечает // отправителю понятной ошибкой до того, как запись ляжет в хранилище: схема // отказала бы уже после укладки файла, а уборки файлов сервис не умеет. // // Первый рубеж при этом ещё раньше: предъявителя без учётной записи пользователя // транспорт отвергает до чтения тела. 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) { // Расширение приходит из имени, которое дал отправитель, и потому может быть // чем угодно. Длину назначает сервис: `filepath.Ext` режет по последней точке // и всё, что после неё, берёт дословно, а имя `x.` плюс четыреста знаков // роняет заведение временного файла — отправитель получал бы `500` на входе, // за который отвечает сам, и клал бы в журнал владельца строку `ERROR`, // неотличимую от аварии хранилища. ext := filepath.Ext(fileName) if ext == "" || len(ext) > maxExtLen { ext = fmt.Sprintf(".%s", defaultAudioExt) } // Собственное имя записи: идентификатор с расширением. Имя, данное // отправителем, в хранилище не попадает — от него взято только расширение. 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) // Признак заводится здесь, а не остаётся голой ошибкой источника // метаданных: причина отказа — присланная запись, а не сбой сервиса, и // без признака ветвь по умолчанию отдала бы `500`. «Файл негоден» // читалось бы как «сломался сервер», и человек не понял бы, что делать. return nil, fmt.Errorf("%w: %w", contract.ErrRecordUnreadable, 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, r.OwnerID) 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 cleaned := entity.SanitizeOriginalFilename(fileName); cleaned != "" { r.OriginalFilename = &cleaned } // Величины принятого — снимок с этой минуты. Со строкой файла они намеренно не // сверяются: там лежат величины сегодняшней копии, и уточнение длительности // меняет их, не трогая эти. r.DurationMs = &meta.DurationMs r.SizeBytes = &size 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, ¬Found) { 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, r.OwnerID) 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 len(outcome.PlainText) == 0 && len(outcome.Replicas) == 0 { s.logger.Warn("Recognition returned empty text", "record_id", r.Id, "operation_id", attempt.ExternalID) } // Сырой ответ сохраняется целиком: результат операции у провайдера не // переспрашивается, и когда мы научимся размечать говорящих, архив // пересчитается из сохранённого без единого рубля. 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) { 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, nil } // failStep останавливает запись приговором шага. // // Исход возвращается **явно**: остановка отказом шага не является — шаг // рассудил об этой записи окончательно, — но и работой она не была, и // засчитывать её успехом нельзя. func (s *TranscribeService) failStep(r *entity.AudioRecord, holder, step string, stepErr error) (stepOutcome, error) { s.halt(r, holder, step, entity.HaltReasonStepFailed, stepErr.Error()) return outcomeHalted, nil } // halt ставит признак остановки, считает её отказом и пишет строку журнала // событий. // // Отправителю отсюда ничего не уходит: инвариант проекта «Принятая запись не // теряется молча» держится теперь опросом готовности — остановка видна там // признаком — и журналом владельца, где у неё стоит причина. // // Счётчик растит **сама остановка**, а не воркер, и это не стилистика. // Остановка по сторожам наступает в захвате, до всякого шага, и воркер о ней // узнаёт признаком «работы нет» — а считать его в метрику запрещено инвариантом // «`NoopJobError` — не ошибка». Значит единственное место, где известны и факт // остановки, и рубеж, — здесь. func (s *TranscribeService) halt(r *entity.AudioRecord, holder, step, reason, errText 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) } // 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) } } // closeWork убирает рабочую копию. Отказ уборки не роняет шаг, но и не // проглатывается: забытая копия это шестичасовая запись во временном каталоге. func (s *TranscribeService) closeWork(work contract.WorkFile) { if err := work.Close(); err != nil { s.logger.Error("Failed to remove work file", "error", err) } } // 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 }