Files
transcriber/internal/service/transcribe.go
T
av 8af8ec2e54 у записи появился владелец: чужую больше не отдают
- колонка `owner` связью с `users` в обеих коллекциях новым шагом схемы
  `202608140001`; чтение задачи сужено владельцем, и чужая, ничья и
  несуществующая дают один ответ; правило просмотра файлов сужено им же
- приём по HTTP берёт владельца из сессии, а предъявителя без учётной записи
  пользователя отвергает до чтения тела: позже пришлось бы убирать уложенный
  файл, а уборки файлов сервис не умеет. Выборка воркера владельцем не сужается
- удаление учётной записи с записями отвергается стражем, и вешает его сама
  сборка хранилища: сборка, забывшая его позвать, теряла защиту молча
2026-08-14 12:18:11 +03:00

663 lines
29 KiB
Go

package service
import (
"context"
"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"
"git.vakhrushev.me/av/transcriber/internal/clock"
)
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(ctx context.Context, 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(ctx, job, file, fileName)
}
// CreateJobFromApi заводит задачу от имени вошедшего. Владелец обязателен:
// пустой отвергается здесь, потому что колонка владельца допускает пустое
// значение ради записей бота, и приём по HTTP — то место, где обязательность
// держится.
//
// Отказ этот — последний рубеж, а не первый: предъявителя без учётной записи
// пользователя транспорт отвергает раньше, до чтения тела. Здесь он остаётся на
// случай нового вызывающего, который такой проверки не поставит.
func (s *TranscribeService) CreateJobFromApi(ctx context.Context, file io.Reader, fileName, ownerID string) (*entity.TranscribeJob, error) {
if ownerID == "" {
s.logger.Error("Refusing to create job without owner")
return nil, contract.ErrOwnerRequired
}
job := &entity.TranscribeJob{
State: entity.StateCreated,
Source: entity.SourceApi,
OwnerID: &ownerID,
}
return s.createTranscribeJob(ctx, job, file, fileName)
}
func (s *TranscribeService) createTranscribeJob(ctx context.Context, 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(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
}
fileRecord, err := s.fileRepo.CreateLocal(storageFileName, work, ownerOf(job))
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(ctx context.Context) error {
return s.runStep(ctx, entity.StateCreated, conversionAcquireTimeout, s.convertJob)
}
func (s *TranscribeService) FindAndRunTranscribeJob(ctx context.Context) error {
return s.runStep(ctx, entity.StateConverted, transcribeAcquireTimeout, s.transcribeJob)
}
func (s *TranscribeService) FindAndRunTranscribeCheckJob(ctx context.Context) error {
return s.runStep(ctx, entity.StateTranscribe, checkAcquireTimeout, s.checkTranscribeJob)
}
// runStep забирает задачу и отдаёт её шагу. Отказ шага не оставляет задачу
// захваченной до конца срока: захват снимается, и задача ждёт нарастающую паузу
// — иначе повтор наступал бы через восемь часов, а не через секунду.
//
// Контекст доходит до шага, а через него — до внешнего собеседника: остановка
// сервиса убивает `ffmpeg` и обрывает запрос к распознаванию. Прерванный шаг
// приговора не выносит: задача остаётся пригодной к повтору, попытки не тратит
// и отправителю о несуществующем сбое не сообщает — исход остановки отличается
// от исхода отказа на каждом шаге.
func (s *TranscribeService) runStep(ctx context.Context, state string, expiration time.Duration, step func(ctx context.Context, job *entity.TranscribeJob, holder string) error) error {
// Нас уже остановили — задачу не забираем: захват стоил бы ей попытки, а
// работы всё равно не будет. Исход «шаг не сделал ничего» — это `NoopJobError`
// по смыслу, и он же не поднимает уровень и не считается в метрику.
if ctx.Err() != nil {
return &contract.NoopJobError{State: state}
}
job, holder, err := s.findJob(state, expiration)
if err != nil {
return err
}
if err := step(ctx, job, holder); err != nil {
s.scheduleRetry(job, holder, err)
return err
}
return nil
}
func (s *TranscribeService) convertJob(ctx context.Context, 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 := clock.Start()
err = s.converter.Convert(ctx, src.Path(), dest.Path())
conversionDuration := time.Since(startTime)
// Записываем метрику времени конвертации
metrics.ObserveConversionDuration(srcExt, "ogg", err != nil, conversionDuration.Seconds())
if err != nil {
// Остановка сервиса — не приговор записи. Убитый по контексту `ffmpeg`
// отдаёт `signal: killed`, и от настоящего отказа конвертации
// (`exit status N`) эта ошибка неотличима ни типом, ни `errors.Is`:
// различает их только контекст шага. Без этой развилки каждый деплой
// хоронил бы конвертируемую запись в `failed` — состояние терминальное,
// и вернуть её оттуда может только владелец правкой в панели, — да ещё
// и сообщал бы отправителю о сбое, которого не было.
if ctxErr := ctx.Err(); ctxErr != nil {
s.logger.Info("File conversion interrupted by shutdown",
"job_id", job.Id,
"duration", conversionDuration)
return fmt.Errorf("conversion interrupted: %w", ctxErr)
}
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, ownerOf(job))
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(ctx context.Context, 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(ctx, content, fileRecord.FileName)
if err != nil {
if ctxErr := ctx.Err(); ctxErr != nil {
s.logger.Info("Recognition interrupted by shutdown", "job_id", job.Id)
return fmt.Errorf("recognition interrupted: %w", ctxErr)
}
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, ownerOf(job))
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 := clock.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(ctx context.Context, 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("recognition 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(ctx, opId)
if err != nil {
if ctxErr := ctx.Err(); ctxErr != nil {
s.logger.Info("Status check interrupted by shutdown", "job_id", job.Id)
return fmt.Errorf("status check interrupted: %w", ctxErr)
}
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 := clock.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(ctx, opId)
if err != nil {
if ctxErr := ctx.Err(); ctxErr != nil {
s.logger.Info("Text fetch interrupted by shutdown", "job_id", job.Id)
return fmt.Errorf("text fetch interrupted: %w", ctxErr)
}
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 := clock.Now().Add(-1 * expiration)
job, err := s.jobRepo.FindAndAcquire(state, acquisitionId, rottingTime)
if err != nil {
// Признак узнаётся по смыслу: репозиторий вправе обернуть свой отказ
// пояснением, и приведение типа от этого сломалось бы молча.
var notFound *contract.JobNotFoundError
if errors.As(err, &notFound) {
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
}
// Остановка попытки не тратит: задача не виновата в том, что нас
// перезапустили. Счётчик растёт при захвате, поэтому здесь его возвращают
// назад — иначе пять выкладок подряд уводят живую запись в «мертва» с
// приговором «попытки исчерпаны».
if errors.Is(stepErr, context.Canceled) || errors.Is(stepErr, context.DeadlineExceeded) {
if job.Attempts > 0 {
job.Attempts--
}
}
job.RetryAfter(clock.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
}
// Адресата у задачи нет: отвечать некуда, и повторять нечего. Уровень здесь
// выше, чем у неподнятого канала, и это не педантизм: пустой чат у задачи
// из Telegram — симптом порчи записи, а самый коварный её источник назван
// инвариантом «колонки очереди правятся в четырёх местах». Утони этот
// сигнал в одном ряду со штатным «бот не настроен» — и обнуление колонки
// заметит только отправитель, переставший получать ответы.
if job.TgChatId == nil {
s.undelivered(job, slog.LevelError, "chat is not specified")
return nil
}
if err := s.tgSender.Send(text, *job.TgChatId, job.TgReplyMessageId); err != nil {
// Канал не поднят: сервис работает без этого входа, и это объявленный
// режим, а не поломка.
if errors.Is(err, contract.ErrDeliveryChannelDown) {
s.undelivered(job, slog.LevelWarn, "delivery channel is down")
return 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
}
// undelivered записывает недоставленный ответ и считает его в метрику. Уровень
// приходит от причины: объявленный режим — «может стать проблемой», порча
// записи — событие для разбора.
//
// Идентификатор задачи обязателен, иначе владелец видит, что ответ не ушёл, но
// не может найти, чей; текста ответа в записи нет — он содержимое чужой записи.
//
// Счётчик нужен потому, что журнал контейнера живёт до ротации, а вопрос «кому
// не ответили за последние сутки» задают позже.
func (s *TranscribeService) undelivered(job *entity.TranscribeJob, level slog.Level, reason string) {
metrics.UndeliveredReplyCounter.WithLabelValues(reason).Inc()
// Уровень выбирается ветвлением, а не передачей контекста: контекст здесь
// брать неоткуда — ответ идёт после сохранения состояния, — а выдуманный
// `context.Background()` соврал бы про отмену и цеплялся бы правилами.
switch level {
case slog.LevelError:
s.logger.Error(undeliveredMessage, "job_id", job.Id, "reason", reason)
default:
s.logger.Warn(undeliveredMessage, "job_id", job.Id, "reason", reason)
}
}
const undeliveredMessage = "Reply was not delivered"
// 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)
}
}
// ownerOf — владелец задачи строкой; пустая значит «владельца нет», и таковы
// записи, принятые ботом. Файл наследует владельца своей задачи: правило
// просмотра коллекции файлов сужено этой колонкой, и файл, заведённый шагом
// конвейера без неё, перестал бы доставаться собственному владельцу.
func ownerOf(job *entity.TranscribeJob) string {
if job.OwnerID == nil {
return ""
}
return *job.OwnerID
}