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

146 lines
6.0 KiB
Go

package yandex
import (
"context"
"fmt"
"io"
"time"
"git.vakhrushev.me/av/transcriber/internal/entity"
)
// ProviderName — имя провайдера, под которым сохраняется попытка распознавания.
// По нему видно, чем считана запись, когда провайдеров станет больше одного.
const ProviderName = "yandex-speechkit"
type YandexAudioRecognizerConfig struct {
// s3
Region string
AccessKey string
SecretKey string
BucketName string
Endpoint string
// speech kit
ApiKey string
FolderID string
}
type YandexAudioRecognizerService struct {
s3Sevice *yandexS3Service
sttService *speechKitService
}
func NewYandexAudioRecognizerService(cfg YandexAudioRecognizerConfig) (*YandexAudioRecognizerService, error) {
s3, err := newYandexS3Service(s3Config{
Region: cfg.Region,
AccessKey: cfg.AccessKey,
SecretKey: cfg.SecretKey,
BucketName: cfg.BucketName,
Endpoint: cfg.Endpoint,
})
if err != nil {
return nil, err
}
stt, err := newSpeechKitService(speechKitConfig{
ApiKey: cfg.ApiKey,
FolderID: cfg.FolderID,
})
if err != nil {
return nil, err
}
return &YandexAudioRecognizerService{
s3Sevice: s3,
sttService: stt,
}, nil
}
func (s *YandexAudioRecognizerService) Close() error {
return s.sttService.Close()
}
func (s *YandexAudioRecognizerService) Provider() string { return ProviderName }
func (s *YandexAudioRecognizerService) Model() string { return RecognitionModel }
// startRecognitionTimeout — сколько ждём принятия операции, когда нас уже
// остановили. Число меньше жёсткого предела остановки: иначе процесс убьют
// прежде, чем ответ дойдёт, и защита ничего не даст.
const startRecognitionTimeout = 10 * time.Second
// Upload кладёт аудио туда, откуда провайдер его прочитает.
//
// Отменяется штатно: заливка дорога по времени, а повтор её бесплатен — объект
// ложится под тем же ключом.
func (s *YandexAudioRecognizerService) Upload(ctx context.Context, file io.Reader, objectKey string) (string, error) {
if err := s.s3Sevice.uploadFile(ctx, file, objectKey); err != nil {
return "", err
}
return s.s3Sevice.fileUrl(objectKey), nil
}
// ObjectExists отвечает, лежит ли объект нужного размера.
//
// Сверка идёт по присутствию и длине, а не по отпечатку содержимого: признак
// целостности у составного объекта не равен отпечатку, и сверка хешем дала бы
// расхождение на всякой большой записи.
func (s *YandexAudioRecognizerService) ObjectExists(ctx context.Context, objectKey string, size int64) (bool, error) {
return s.s3Sevice.objectExists(ctx, objectKey, size)
}
// Submit заводит операцию распознавания. Оплачивается наружу, поэтому от отмены
// защищён: окно короткое и дорогое — SpeechKit может операцию принять и начать
// считать деньги, а ответ до нас не доедет, и повтор оплатит ту же запись второй
// раз. Свой предел вызову оставлен, чтобы остановка не ждала вечно.
func (s *YandexAudioRecognizerService) Submit(ctx context.Context, sourceURI string) (string, error) {
startCtx, cancel := protectFromCancel(ctx, startRecognitionTimeout)
defer cancel()
return s.sttService.recognizeFileFromS3(startCtx, sourceURI)
}
// protectFromCancel отвязывает вызов от отмены родителя, оставляя ему значения
// родителя и собственный предел по времени. Употребляется там, где обрыв стоит
// дороже ожидания: у платной операции, чей результат нельзя переспросить.
func protectFromCancel(ctx context.Context, timeout time.Duration) (context.Context, context.CancelFunc) {
return context.WithTimeout(context.WithoutCancel(ctx), timeout)
}
// Fetch забирает готовый результат и отдаёт его доменным: реплики со временем,
// плоский текст и байты ответа на хранение. Формата провайдера наружу не выходит
// ничего — ни один шаг конвейера не знает, каким потоком тот отвечает.
func (s *YandexAudioRecognizerService) Fetch(ctx context.Context, operationID string) (*entity.RecognitionOutcome, error) {
return s.sttService.fetchRecognition(ctx, operationID)
}
// Parse строит доменный результат из **сохранённого** ответа, не обращаясь к
// провайдеру. По нему архив пересчитывается без единого рубля.
func (s *YandexAudioRecognizerService) Parse(raw []byte) (*entity.RecognitionOutcome, error) {
responses, err := decodeResponses(raw)
if err != nil {
return nil, err
}
outcome := outcomeFromResponses(responses)
outcome.Raw = raw
return outcome, nil
}
func (s *YandexAudioRecognizerService) CheckStatus(ctx context.Context, operationID string) (*entity.RecognitionResult, error) {
operation, err := s.sttService.checkOperationStatus(ctx, operationID)
if err != nil {
return nil, err
}
if !operation.Done {
return entity.NewInProgressResult(), nil
}
if opErr := operation.GetError(); opErr != nil {
errorText := fmt.Sprintf("operation failed: code %d, message: %s", opErr.Code, opErr.Message)
return entity.NewFailedResult(errorText), nil
}
return entity.NewCompletedResult(), nil
}