- audiorecords вместо transcribe_jobs: приложения (texts, structures, recognitions, record_events, topics) живут своими коллекциями, ссылки на исходник и на приведённую копию перестали переставляться - рубеж называет достигнутое, отказ стал признаком остановки с причиной, а сторожей стало двое: число отказов и время в рубеже - воркеры потеряли специализацию, их число задаётся [pipeline] workers, шаг выбирается по рубежу, а захват отдаёт идентификатор и признак захвата
198 lines
7.6 KiB
Go
198 lines
7.6 KiB
Go
package service
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"io"
|
|
"sync/atomic"
|
|
"testing"
|
|
|
|
"github.com/google/uuid"
|
|
"github.com/stretchr/testify/assert"
|
|
"github.com/stretchr/testify/require"
|
|
|
|
"git.vakhrushev.me/av/transcriber/internal/entity"
|
|
)
|
|
|
|
// countingRecognizer считает обращения наружу: за них платят по факту, и повтор
|
|
// оплаченного шага — самый дорогой класс дефекта в этом сервисе.
|
|
type countingRecognizer struct {
|
|
uploads atomic.Int64
|
|
submits atomic.Int64
|
|
fetches atomic.Int64
|
|
parses atomic.Int64
|
|
objectHere atomic.Bool
|
|
inProgress atomic.Int64
|
|
}
|
|
|
|
func (r *countingRecognizer) Provider() string { return "counting" }
|
|
|
|
func (r *countingRecognizer) Model() string { return "counting" }
|
|
|
|
func (r *countingRecognizer) Upload(_ context.Context, _ io.Reader, objectKey string) (string, error) {
|
|
r.uploads.Add(1)
|
|
r.objectHere.Store(true)
|
|
return "counting://" + objectKey, nil
|
|
}
|
|
|
|
func (r *countingRecognizer) ObjectExists(context.Context, string, int64) (bool, error) {
|
|
return r.objectHere.Load(), nil
|
|
}
|
|
|
|
func (r *countingRecognizer) SourceURI(objectKey string) string {
|
|
return "counting://" + objectKey
|
|
}
|
|
|
|
func (r *countingRecognizer) Submit(context.Context, string) (string, error) {
|
|
r.submits.Add(1)
|
|
return uuid.NewString(), nil
|
|
}
|
|
|
|
func (r *countingRecognizer) CheckStatus(context.Context, string) (*entity.RecognitionResult, error) {
|
|
if r.inProgress.Load() > 0 {
|
|
r.inProgress.Add(-1)
|
|
return entity.NewInProgressResult(), nil
|
|
}
|
|
return entity.NewCompletedResult(), nil
|
|
}
|
|
|
|
func (r *countingRecognizer) Fetch(context.Context, string) (*entity.RecognitionOutcome, error) {
|
|
r.fetches.Add(1)
|
|
return r.Parse(countingPayload(t0Replicas))
|
|
}
|
|
|
|
func (r *countingRecognizer) Parse(raw []byte) (*entity.RecognitionOutcome, error) {
|
|
r.parses.Add(1)
|
|
|
|
var replicas []entity.Replica
|
|
if err := json.Unmarshal(raw, &replicas); err != nil {
|
|
return nil, errors.New("не разобрать сохранённый ответ")
|
|
}
|
|
|
|
var plain []byte
|
|
for _, replica := range replicas {
|
|
if len(plain) > 0 {
|
|
plain = append(plain, ' ')
|
|
}
|
|
plain = append(plain, replica.Text...)
|
|
}
|
|
|
|
return &entity.RecognitionOutcome{Replicas: replicas, PlainText: string(plain), Raw: raw}, nil
|
|
}
|
|
|
|
var t0Replicas = []entity.Replica{
|
|
{StartMs: 0, EndMs: 900, Text: "Первая реплика."},
|
|
{StartMs: 900, EndMs: 1800, Text: "Вторая реплика."},
|
|
}
|
|
|
|
func countingPayload(replicas []entity.Replica) []byte {
|
|
raw, err := json.Marshal(replicas)
|
|
if err != nil {
|
|
panic(err)
|
|
}
|
|
return raw
|
|
}
|
|
|
|
// Критерий приёмки 4. Структура реплик строится из сохранённого ответа
|
|
// провайдера без единого обращения к нему: результат операции не
|
|
// переспрашивается, и архив пересчитывается из сохранённого без рубля.
|
|
func TestStructureIsBuiltFromStoredPayload(t *testing.T) {
|
|
rec := &countingRecognizer{}
|
|
env := newPipelineEnvWith(t, &okMetaViewer{}, &okConverter{}, rec)
|
|
|
|
record := newTelegramRecord(t, env)
|
|
drain(t, env, record.Id)
|
|
|
|
after := readRecord(t, env, record.Id)
|
|
require.Equal(t, entity.StateDone, after.State)
|
|
require.NotNil(t, after.RecognitionID)
|
|
|
|
// Сохранённый ответ на месте и читается целиком.
|
|
raw, err := env.repos.Recognitions.ReadRaw(*after.RecognitionID)
|
|
require.NoError(t, err)
|
|
require.NotEmpty(t, raw, "сырой ответ провайдера сохранён")
|
|
|
|
// А теперь — построение структуры из сохранённого, без обращений наружу.
|
|
fetchesBefore := rec.fetches.Load()
|
|
submitsBefore := rec.submits.Load()
|
|
|
|
outcome, err := env.service.recognizer.Parse(raw)
|
|
require.NoError(t, err)
|
|
require.Len(t, outcome.Replicas, len(t0Replicas), "структура собрана")
|
|
assert.Equal(t, t0Replicas[0].Text, outcome.Replicas[0].Text)
|
|
assert.Equal(t, t0Replicas[0].StartMs, outcome.Replicas[0].StartMs, "время реплики сохранено")
|
|
|
|
assert.Equal(t, fetchesBefore, rec.fetches.Load(), "к провайдеру за результатом не ходили")
|
|
assert.Equal(t, submitsBefore, rec.submits.Load(), "и новой операции не заводили")
|
|
|
|
// Структура записи собрана из того же ответа и лежит своей строкой.
|
|
require.NotNil(t, after.StructureID)
|
|
structure, err := env.repos.Structures.GetByID(*after.StructureID)
|
|
require.NoError(t, err)
|
|
assert.Equal(t, entity.StructureVersion, structure.Version)
|
|
require.Len(t, structure.Replicas, len(t0Replicas))
|
|
assert.Equal(t, t0Replicas[1].Text, structure.Replicas[1].Text)
|
|
}
|
|
|
|
// Шаг с внешней оплатой проверяет сделанное: объект нужного размера на месте —
|
|
// заливку не повторяем; операция заведена — не платим второй раз.
|
|
func TestPaidWorkIsNotRepeated(t *testing.T) {
|
|
rec := &countingRecognizer{}
|
|
env := newPipelineEnvWith(t, &okMetaViewer{}, &okConverter{}, rec)
|
|
|
|
record := newTelegramRecord(t, env)
|
|
|
|
// Приведение.
|
|
require.NoError(t, env.service.RunStep(t.Context()))
|
|
require.Equal(t, entity.StateNormalized, readRecord(t, env, record.Id).State)
|
|
|
|
// Отправка.
|
|
require.NoError(t, env.service.RunStep(t.Context()))
|
|
require.Equal(t, int64(1), rec.uploads.Load(), "залили один раз")
|
|
require.Equal(t, int64(1), rec.submits.Load(), "и заплатили один раз")
|
|
|
|
// Возвращаем запись на рубеж отправки — так это выглядит, когда шаг оборвался
|
|
// после оплаты, а захват протух.
|
|
after := readRecord(t, env, record.Id)
|
|
after.MoveToState(entity.StateNormalized)
|
|
require.NoError(t, env.recordRepo.Save(after, ""))
|
|
|
|
require.NoError(t, env.service.RunStep(t.Context()))
|
|
|
|
assert.Equal(t, int64(1), rec.uploads.Load(), "объект на месте — заливка не повторилась")
|
|
assert.Equal(t, int64(1), rec.submits.Load(), "операция заведена — второй раз не платим")
|
|
}
|
|
|
|
// Ожидание чужой операции откладывается своей задержкой, отказов не тратит и
|
|
// рубежа не двигает.
|
|
func TestPollingPostponesWithoutSpendingAttempts(t *testing.T) {
|
|
rec := &countingRecognizer{}
|
|
rec.inProgress.Store(3)
|
|
env := newPipelineEnvWith(t, &okMetaViewer{}, &okConverter{}, rec)
|
|
|
|
record := newTelegramRecord(t, env)
|
|
|
|
require.NoError(t, env.service.RunStep(t.Context())) // приведение
|
|
require.NoError(t, env.service.RunStep(t.Context())) // отправка
|
|
require.Equal(t, entity.StateSubmitted, readRecord(t, env, record.Id).State)
|
|
|
|
entered := readRecord(t, env, record.Id).StateEnteredAt
|
|
|
|
var delays []int64
|
|
for range 3 {
|
|
clearDelay(t, env, record.Id)
|
|
require.NoError(t, env.service.RunStep(t.Context()))
|
|
|
|
after := readRecord(t, env, record.Id)
|
|
require.Equal(t, entity.StateSubmitted, after.State, "рубеж не сдвинут откладыванием")
|
|
assert.Equal(t, 0, after.Attempts, "ожидание чужой операции отказа не тратит")
|
|
assert.WithinDuration(t, entered, after.StateEnteredAt, 0, "и отсчёт застревания не двигает")
|
|
|
|
require.NotNil(t, after.DelayTime)
|
|
delays = append(delays, after.DelayTime.Unix())
|
|
}
|
|
|
|
assert.Len(t, delays, 3, "все три прогона отложили работу")
|
|
}
|