- audiorecords вместо transcribe_jobs: приложения (texts, structures, recognitions, record_events, topics) живут своими коллекциями, ссылки на исходник и на приведённую копию перестали переставляться - рубеж называет достигнутое, отказ стал признаком остановки с причиной, а сторожей стало двое: число отказов и время в рубеже - воркеры потеряли специализацию, их число задаётся [pipeline] workers, шаг выбирается по рубежу, а захват отдаёт идентификатор и признак захвата
146 lines
6.0 KiB
Go
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
|
|
}
|