удалён вход Telegram, владелец записи стал обязателен в схеме
- убраны клиент бота, транспорт обновлений, отправитель сообщений, сборка входа при старте, секция настроек и зависимость go-telegram-bot-api; из конвейера ушла доставка ответа отправителю — исход виден опросом готовности. Колонки адресата и значение источника остались в схеме: применённые шаги не переписываются - шаг 202608140003 запрещает пустого владельца у аудиозаписи и у файла; существующие строки он не проверяет, и это принято сознательно — искать их надо запросом до выкладки - ревью нашло два пред-существующих дефекта, оба закрыты: пустой второй ответ распознавателя стирал сохранённую расшифровку, а пустая расшифровка перестала быть заметной вместе с убранной доставкой. Попутно поднят golang.org/x/image до v0.45.0 — красный шаг vulns, воспроизводился и на чистом master
This commit is contained in:
@@ -40,7 +40,7 @@ func (r *stubRecordRepo) FindAndAcquire([]entity.Stage) (*contract.AcquiredRecor
|
||||
func serviceWithRepo(repo contract.AudioRecordRepository) *TranscribeService {
|
||||
return NewTranscribeService(
|
||||
Repositories{Records: repo},
|
||||
nil, nil, nil, nil,
|
||||
nil, nil, nil,
|
||||
entity.StuckLimits{},
|
||||
slog.New(slog.DiscardHandler),
|
||||
)
|
||||
|
||||
@@ -55,7 +55,7 @@ func TestFailureIsCountedWithStageLabel(t *testing.T) {
|
||||
|
||||
beforeOk := stageCount(t, entity.StateUploaded, "false")
|
||||
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
require.NoError(t, env.service.RunStep(t.Context()))
|
||||
require.Equal(t, entity.StateNormalized, readRecord(t, env, record.Id).State)
|
||||
|
||||
@@ -66,7 +66,7 @@ func TestFailureIsCountedWithStageLabel(t *testing.T) {
|
||||
failing := newPipelineEnv(t, &okMetaViewer{}, &failingStepConverter{})
|
||||
beforeErr := stageCount(t, entity.StateUploaded, "true")
|
||||
|
||||
newTelegramRecord(t, failing)
|
||||
newRecord(t, failing)
|
||||
require.Error(t, failing.service.RunStep(t.Context()))
|
||||
|
||||
assert.InDelta(t, beforeErr+1, stageCount(t, entity.StateUploaded, "true"), 0,
|
||||
|
||||
@@ -12,9 +12,8 @@ import (
|
||||
)
|
||||
|
||||
// Выборка воркера владельцем не сужается: владелец решает, кому запись
|
||||
// показывать, а не кому её считать. Сужение остановило бы расшифровку записей
|
||||
// бота вовсе, а записи остальных поставило бы в зависимость от того, кто первым
|
||||
// завёл учётную запись.
|
||||
// показывать, а не кому её считать. Сужение поставило бы записи одних людей в
|
||||
// зависимость от того, кто первым завёл учётную запись.
|
||||
func TestWorkerTakesRecordsOfEveryOwner(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
|
||||
|
||||
@@ -26,8 +25,8 @@ func TestWorkerTakesRecordsOfEveryOwner(t *testing.T) {
|
||||
strings.NewReader("вторая"), "two.mp3", newOwner(t, env.app))
|
||||
require.NoError(t, err)
|
||||
|
||||
// Третья пришла ботом, и владельца у неё нет вовсе.
|
||||
third := newTelegramRecord(t, env)
|
||||
// Третья — ещё одного владельца: воркер не сужается ни одним из них.
|
||||
third := newRecord(t, env)
|
||||
|
||||
// Захваченная запись остаётся за держателем, и следующий вызов берёт
|
||||
// следующую, а не ту же самую.
|
||||
@@ -40,7 +39,7 @@ func TestWorkerTakesRecordsOfEveryOwner(t *testing.T) {
|
||||
|
||||
assert.True(t, taken[first.Id], "запись первого владельца досталась воркеру")
|
||||
assert.True(t, taken[second.Id], "и второго")
|
||||
assert.True(t, taken[third.Id], "и запись без владельца")
|
||||
assert.True(t, taken[third.Id], "и третьего")
|
||||
|
||||
_, err = env.recordRepo.FindAndAcquire(entity.WorkingStages())
|
||||
var missing *contract.JobNotFoundError
|
||||
@@ -65,8 +64,7 @@ func TestAcquireReturnsIdentifierAndHolder(t *testing.T) {
|
||||
|
||||
// Колонки читаются отдельным чтением, и владелец среди них.
|
||||
read := readRecord(t, env, record.Id)
|
||||
require.NotNil(t, read.OwnerID)
|
||||
assert.Equal(t, owner, *read.OwnerID)
|
||||
assert.Equal(t, owner, read.OwnerID)
|
||||
assert.Equal(t, acquired.Holder, *read.AcquisitionID, "признак захвата записан в саму запись")
|
||||
require.NotNil(t, read.AcquireExpiresAt, "срок протухания приехал с рубежом")
|
||||
}
|
||||
@@ -86,12 +84,12 @@ func TestPipelineStepKeepsOwner(t *testing.T) {
|
||||
|
||||
after := readRecord(t, env, record.Id)
|
||||
require.True(t, after.IsHalted(), "шаг записал свой приговор")
|
||||
require.NotNil(t, after.OwnerID, "владелец пережил шаг")
|
||||
assert.Equal(t, owner, *after.OwnerID)
|
||||
assert.Equal(t, owner, after.OwnerID, "владелец пережил шаг")
|
||||
}
|
||||
|
||||
// Приём из веба без владельца записи не заводит. Обязательность держит здесь
|
||||
// код, а не схема: колонка допускает пустое значение ради записей бота.
|
||||
// Приём без владельца записи не заводит. Отказ стоит здесь раньше схемы: он
|
||||
// отвечает понятной ошибкой до того, как запись ляжет в хранилище, а схема
|
||||
// отказала бы уже после укладки файла.
|
||||
func TestCreateJobFromApiRequiresOwner(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
|
||||
|
||||
@@ -110,17 +108,11 @@ func TestGetByIDHidesForeignRecords(t *testing.T) {
|
||||
record, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "one.mp3", owner)
|
||||
require.NoError(t, err)
|
||||
|
||||
// Запись без владельца — принятая ботом.
|
||||
orphan := newTelegramRecord(t, env)
|
||||
|
||||
var missing *contract.JobNotFoundError
|
||||
|
||||
_, err = env.recordRepo.GetByID(record.Id, stranger)
|
||||
require.ErrorAs(t, err, &missing, "чужая запись неотличима от несуществующей")
|
||||
|
||||
_, err = env.recordRepo.GetByID(orphan.Id, stranger)
|
||||
require.ErrorAs(t, err, &missing, "ничья запись не достаётся никому")
|
||||
|
||||
_, err = env.recordRepo.GetByID(record.Id, "")
|
||||
require.ErrorAs(t, err, &missing, "пустой владелец не совпадает ни с чем")
|
||||
|
||||
|
||||
@@ -60,32 +60,12 @@ func (m *failingMetaViewer) GetInfo(context.Context, string) (*contract.AudioInf
|
||||
return nil, errors.New("запись не читается")
|
||||
}
|
||||
|
||||
// recordingSender запоминает, что отправлено.
|
||||
type recordingSender struct {
|
||||
mu sync.Mutex
|
||||
messages []string
|
||||
}
|
||||
|
||||
func (s *recordingSender) Send(text string, chatId int64, replyMsgId *int) error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
s.messages = append(s.messages, text)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *recordingSender) sent() []string {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return append([]string(nil), s.messages...)
|
||||
}
|
||||
|
||||
type pipelineEnv struct {
|
||||
app core.App
|
||||
service *TranscribeService
|
||||
repos Repositories
|
||||
recordRepo *pbrepo.AudioRecordRepository
|
||||
fileRepo *pbrepo.FileRepository
|
||||
sender *recordingSender
|
||||
}
|
||||
|
||||
// testLimits — пределы простоя проверок. Числа боевые; проверка застревания
|
||||
@@ -104,6 +84,19 @@ func newPipelineEnvWith(
|
||||
rec contract.AudioRecognizer,
|
||||
) *pipelineEnv {
|
||||
t.Helper()
|
||||
return newPipelineEnvWithLogger(t, metaviewer, converter, rec, slog.New(slog.DiscardHandler))
|
||||
}
|
||||
|
||||
// newPipelineEnvWithLogger — то же с подменённым журналом: проверке, судящей
|
||||
// строку журнала, нужен свой, а не общий.
|
||||
func newPipelineEnvWithLogger(
|
||||
t *testing.T,
|
||||
metaviewer contract.AudioMetaViewer,
|
||||
converter contract.AudioFileConverter,
|
||||
rec contract.AudioRecognizer,
|
||||
logger *slog.Logger,
|
||||
) *pipelineEnv {
|
||||
t.Helper()
|
||||
|
||||
app, err := pbrepo.New(t.TempDir())
|
||||
require.NoError(t, err)
|
||||
@@ -128,9 +121,7 @@ func newPipelineEnvWith(
|
||||
Recognitions: pbrepo.NewRecognitionRepository(app),
|
||||
Events: pbrepo.NewRecordEventRepository(app),
|
||||
}
|
||||
sender := &recordingSender{}
|
||||
|
||||
svc := NewTranscribeService(repos, metaviewer, converter, rec, sender, testLimits, slog.New(slog.DiscardHandler))
|
||||
svc := NewTranscribeService(repos, metaviewer, converter, rec, testLimits, logger)
|
||||
|
||||
return &pipelineEnv{
|
||||
app: app,
|
||||
@@ -138,15 +129,16 @@ func newPipelineEnvWith(
|
||||
repos: repos,
|
||||
recordRepo: recordRepo,
|
||||
fileRepo: fileRepo,
|
||||
sender: sender,
|
||||
}
|
||||
}
|
||||
|
||||
// newTelegramRecord заводит запись — так, как её завёл бы приём из бота.
|
||||
func newTelegramRecord(t *testing.T, env *pipelineEnv) *entity.AudioRecord {
|
||||
// newRecord заводит запись — так, как её заводит приём по HTTP: от имени
|
||||
// вошедшего, потому что ничьей записи в хранилище не бывает.
|
||||
func newRecord(t *testing.T, env *pipelineEnv) *entity.AudioRecord {
|
||||
t.Helper()
|
||||
|
||||
record, err := env.service.CreateJobFromTelegram(t.Context(), strings.NewReader("запись"), "voice.ogg", 100, 1)
|
||||
record, err := env.service.CreateJobFromApi(
|
||||
t.Context(), strings.NewReader("запись"), "voice.ogg", newOwner(t, env.app))
|
||||
require.NoError(t, err)
|
||||
return record
|
||||
}
|
||||
@@ -175,6 +167,17 @@ func enteredStateAt(t *testing.T, env *pipelineEnv, recordID string, moment time
|
||||
require.NoError(t, env.app.Save(record))
|
||||
}
|
||||
|
||||
// setAttempts ставит записи число отказов: так она выглядит, когда шаг отказал
|
||||
// столько раз подряд, что предел исчерпан.
|
||||
func setAttempts(t *testing.T, env *pipelineEnv, recordID string, attempts int) {
|
||||
t.Helper()
|
||||
|
||||
record, err := env.app.FindRecordById(migrations.RecordsCollection, recordID)
|
||||
require.NoError(t, err)
|
||||
record.Set("attempts", attempts)
|
||||
require.NoError(t, env.app.Save(record))
|
||||
}
|
||||
|
||||
// drain крутит конвейер, пока он двигает записи. Паузы опроса снимаются: они
|
||||
// проверяются отдельно, а здесь мешают дойти до конца.
|
||||
func drain(t *testing.T, env *pipelineEnv, recordID string) {
|
||||
@@ -217,7 +220,7 @@ func readRecord(t *testing.T, env *pipelineEnv, id string) *entity.AudioRecord {
|
||||
func TestRecordHaltsAfterAttemptLimit(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
|
||||
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
|
||||
// Захват без выполнения шага — так это выглядит при гибели процесса: отказа
|
||||
// шаг объявить не успевает, а попытка засчитана.
|
||||
@@ -239,9 +242,6 @@ func TestRecordHaltsAfterAttemptLimit(t *testing.T) {
|
||||
assert.Equal(t, entity.StateUploaded, after.State, "рубеж пережил остановку")
|
||||
assert.Greater(t, after.Attempts, maxAttempts, "число отказов сохранено")
|
||||
|
||||
require.Len(t, env.sender.sent(), 1, "отправитель узнал о неудаче")
|
||||
assert.Contains(t, env.sender.sent()[0], "попытки исчерпаны")
|
||||
|
||||
// И из выборки она исчезла.
|
||||
_, err = env.recordRepo.FindAndAcquire(entity.WorkingStages())
|
||||
var missing *contract.JobNotFoundError
|
||||
@@ -265,7 +265,7 @@ func expireAcquisition(t *testing.T, env *pipelineEnv, recordID string) {
|
||||
func TestHaltedRecordResumesFromItsStage(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{})
|
||||
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
|
||||
// Доводим до рубежа приведения и останавливаем на нём.
|
||||
require.NoError(t, env.service.RunStep(t.Context()))
|
||||
@@ -302,7 +302,7 @@ func TestHaltedRecordResumesFromItsStage(t *testing.T) {
|
||||
func TestBothFileLinksSurvivePipeline(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{})
|
||||
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
drain(t, env, record.Id)
|
||||
|
||||
after := readRecord(t, env, record.Id)
|
||||
@@ -327,7 +327,7 @@ func TestBothFileLinksSurvivePipeline(t *testing.T) {
|
||||
func TestOutcomeDoesNotDependOnWorkerCount(t *testing.T) {
|
||||
for _, workers := range []int{1, 4} {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{})
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
|
||||
for range 20 {
|
||||
clearDelay(t, env, record.Id)
|
||||
@@ -361,7 +361,7 @@ func TestOutcomeDoesNotDependOnWorkerCount(t *testing.T) {
|
||||
// Записи принимаются и не двигаются, и это режим, а не поломка.
|
||||
func TestZeroWorkersLeaveRecordUntouched(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{})
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
|
||||
// Пул нулевого размера ни одного прогона не делает — потому запись остаётся
|
||||
// там, где её оставил приём.
|
||||
@@ -393,7 +393,7 @@ func TestPostponeKeepsStuckCountdown(t *testing.T) {
|
||||
func TestStuckRecordIsHalted(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{})
|
||||
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
enteredStateAt(t, env, record.Id, time.Now().Add(-2*testLimits.Own))
|
||||
|
||||
err := env.service.RunStep(t.Context())
|
||||
@@ -406,8 +406,6 @@ func TestStuckRecordIsHalted(t *testing.T) {
|
||||
assert.Equal(t, entity.HaltReasonStuck, *after.HaltReason, "причина названа")
|
||||
assert.Equal(t, entity.StateUploaded, after.State, "рубеж сохранён")
|
||||
|
||||
require.Len(t, env.sender.sent(), 1, "отправитель узнал о неудаче")
|
||||
assert.Contains(t, env.sender.sent()[0], "застряла")
|
||||
}
|
||||
|
||||
// Конечный рубеж под предел простоя не подпадает: стоять в нём запись будет
|
||||
@@ -415,7 +413,7 @@ func TestStuckRecordIsHalted(t *testing.T) {
|
||||
func TestDoneRecordIsNeverStuck(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{})
|
||||
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
drain(t, env, record.Id)
|
||||
enteredStateAt(t, env, record.Id, time.Now().Add(-10*24*time.Hour))
|
||||
|
||||
@@ -434,7 +432,7 @@ func TestDoneRecordIsNeverStuck(t *testing.T) {
|
||||
func TestOnlyHolderWritesResult(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{})
|
||||
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
|
||||
first, err := env.recordRepo.FindAndAcquire(entity.WorkingStages())
|
||||
require.NoError(t, err)
|
||||
@@ -457,12 +455,12 @@ func TestOnlyHolderWritesResult(t *testing.T) {
|
||||
|
||||
after := readRecord(t, env, record.Id)
|
||||
assert.Equal(t, entity.StateUploaded, after.State, "рубеж не сдвинут потерявшим захват")
|
||||
assert.Empty(t, env.sender.sent(), "и отправителю от него ничего не ушло")
|
||||
}
|
||||
|
||||
// Критерий приёмки 7. Всякий способ вывести запись из работы сообщает
|
||||
// отправителю: причин остановки больше одной, и обязанность у них общая.
|
||||
func TestEveryHaltReasonNotifiesSender(t *testing.T) {
|
||||
// Критерий приёмки 7. Всякий способ вывести запись из работы оставляет причину
|
||||
// остановки: причин больше одной, и обязанность у них общая. Отправитель узнаёт
|
||||
// исход опросом готовности, а владелец сервиса — журналом событий записи.
|
||||
func TestEveryHaltReasonRecordsItsCause(t *testing.T) {
|
||||
reasons := []struct {
|
||||
name string
|
||||
halt func(t *testing.T, env *pipelineEnv, recordID string)
|
||||
@@ -481,18 +479,40 @@ func TestEveryHaltReasonNotifiesSender(t *testing.T) {
|
||||
require.ErrorAs(t, env.service.RunStep(t.Context()), &noop)
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "исчерпанные отказы",
|
||||
halt: func(t *testing.T, env *pipelineEnv, recordID string) {
|
||||
setAttempts(t, env, recordID, maxAttempts+1)
|
||||
var noop *contract.NoopJobError
|
||||
require.ErrorAs(t, env.service.RunStep(t.Context()), &noop)
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
for _, reason := range reasons {
|
||||
t.Run(reason.name, func(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
|
||||
reason.halt(t, env, record.Id)
|
||||
|
||||
after := readRecord(t, env, record.Id)
|
||||
require.True(t, after.IsHalted(), "запись остановлена")
|
||||
require.Len(t, env.sender.sent(), 1, "отправитель узнал о неудаче")
|
||||
|
||||
// Признак остановки и причина ставятся одним движением, поэтому
|
||||
// вторым утверждением берётся **журнал событий**: он пишется
|
||||
// отдельной строкой, отдельным сохранением, и упасть может сам по
|
||||
// себе. Прежде эту роль играл счёт ответов отправителю; ответы ушли
|
||||
// вместе с входом Telegram, и без замены проверка вывелась бы из
|
||||
// собственной предыдущей строки.
|
||||
events := recordEvents(t, env, record.Id)
|
||||
require.Len(t, events, 1, "остановка оставила строку журнала событий")
|
||||
|
||||
outcome := events[0].GetString("outcome")
|
||||
assert.Contains(t,
|
||||
[]string{entity.EventOutcomeHalted, entity.EventOutcomeFailed}, outcome,
|
||||
"строка журнала называет исход остановкой")
|
||||
assert.NotEmpty(t, events[0].GetString("outcome_text"), "и несёт причину")
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -503,7 +523,7 @@ func TestEveryHaltReasonNotifiesSender(t *testing.T) {
|
||||
func TestRepeatedStoreKeepsSingleText(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{})
|
||||
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
|
||||
first, err := env.repos.Texts.Put(record.Id, entity.TextKindTranscript, "первый разбор")
|
||||
require.NoError(t, err)
|
||||
@@ -535,7 +555,7 @@ func TestRetryDelayGrowsAndCaps(t *testing.T) {
|
||||
func TestFailedStepSchedulesRetryWithGrowingDelay(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
|
||||
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
|
||||
// Ссылка переставляется на запись о файле без содержимого: шаг отказывает на
|
||||
// получении рабочей копии — то есть отказом, а не приговором записи.
|
||||
@@ -544,6 +564,8 @@ func TestFailedStepSchedulesRetryWithGrowingDelay(t *testing.T) {
|
||||
empty := core.NewRecord(files)
|
||||
empty.Set("location", entity.LocationLocal)
|
||||
empty.Set("size", 1)
|
||||
// Владелец обязателен и у файла: схема ничьих не принимает.
|
||||
empty.Set("owner", record.OwnerID)
|
||||
require.NoError(t, env.app.Save(empty))
|
||||
|
||||
stored, err := env.app.FindRecordById(migrations.RecordsCollection, record.Id)
|
||||
@@ -606,7 +628,7 @@ func TestWorkFilesRemovedAfterConversionFailure(t *testing.T) {
|
||||
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
|
||||
|
||||
newTelegramRecord(t, env)
|
||||
newRecord(t, env)
|
||||
|
||||
leftovers, err := filepath.Glob(filepath.Join(tempDir, "transcriber-*"))
|
||||
require.NoError(t, err)
|
||||
@@ -624,7 +646,7 @@ func TestWorkFilesRemovedAfterConversionFailure(t *testing.T) {
|
||||
func TestRecordNeverPointsToMissingFile(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
|
||||
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
|
||||
require.NoError(t, env.service.RunStep(t.Context()))
|
||||
|
||||
@@ -714,7 +736,7 @@ func TestHaltIsCountedAsFailureAndLoggedOnce(t *testing.T) {
|
||||
beforeOk := stageCount(t, entity.StateUploaded, "false")
|
||||
beforeErr := stageCount(t, entity.StateUploaded, "true")
|
||||
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
require.NoError(t, env.service.RunStep(t.Context()))
|
||||
|
||||
after := readRecord(t, env, record.Id)
|
||||
@@ -739,7 +761,7 @@ func TestPostponeWritesNoEvent(t *testing.T) {
|
||||
rec.inProgress.Store(5)
|
||||
env := newPipelineEnvWith(t, &okMetaViewer{}, &okConverter{}, rec)
|
||||
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(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)
|
||||
|
||||
@@ -1,10 +1,12 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"log/slog"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
|
||||
@@ -101,7 +103,7 @@ func TestStructureIsBuiltFromStoredPayload(t *testing.T) {
|
||||
rec := &countingRecognizer{}
|
||||
env := newPipelineEnvWith(t, &okMetaViewer{}, &okConverter{}, rec)
|
||||
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
drain(t, env, record.Id)
|
||||
|
||||
after := readRecord(t, env, record.Id)
|
||||
@@ -141,7 +143,7 @@ func TestPaidWorkIsNotRepeated(t *testing.T) {
|
||||
rec := &countingRecognizer{}
|
||||
env := newPipelineEnvWith(t, &okMetaViewer{}, &okConverter{}, rec)
|
||||
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
|
||||
// Приведение.
|
||||
require.NoError(t, env.service.RunStep(t.Context()))
|
||||
@@ -171,7 +173,7 @@ func TestPollingPostponesWithoutSpendingAttempts(t *testing.T) {
|
||||
rec.inProgress.Store(3)
|
||||
env := newPipelineEnvWith(t, &okMetaViewer{}, &okConverter{}, rec)
|
||||
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
|
||||
require.NoError(t, env.service.RunStep(t.Context())) // приведение
|
||||
require.NoError(t, env.service.RunStep(t.Context())) // отправка
|
||||
@@ -195,3 +197,88 @@ func TestPollingPostponesWithoutSpendingAttempts(t *testing.T) {
|
||||
|
||||
assert.Len(t, delays, 3, "все три прогона отложили работу")
|
||||
}
|
||||
|
||||
// emptyRecognizer отдаёт готовую операцию с пустым результатом: так выглядит
|
||||
// оплаченное распознавание, из которого ничего не вышло.
|
||||
type emptyRecognizer struct {
|
||||
countingRecognizer
|
||||
}
|
||||
|
||||
func (r *emptyRecognizer) Fetch(context.Context, string) (*entity.RecognitionOutcome, error) {
|
||||
r.fetches.Add(1)
|
||||
return &entity.RecognitionOutcome{}, nil
|
||||
}
|
||||
|
||||
// Пустой результат виден владельцу сервиса записью журнала «может стать
|
||||
// проблемой». Прежде этот случай был виден ответом отправителю — «на записи нет
|
||||
// текста», — и вместе с убранной доставкой он исчез бы вовсе: запись доходит до
|
||||
// конца молча и от успешной не отличается. Текста в записи нет по инварианту
|
||||
// приватности: пустой ему взяться неоткуда, а проверка судит уровень и
|
||||
// идентификатор.
|
||||
func TestEmptyRecognitionIsNamedInJournal(t *testing.T) {
|
||||
journal := &bytes.Buffer{}
|
||||
|
||||
env := newPipelineEnvWithLogger(t, &okMetaViewer{}, &okConverter{}, &emptyRecognizer{},
|
||||
slog.New(slog.NewTextHandler(journal, &slog.HandlerOptions{Level: slog.LevelDebug})))
|
||||
|
||||
record := newRecord(t, env)
|
||||
drain(t, env, record.Id)
|
||||
|
||||
written := journal.String()
|
||||
assert.Contains(t, written, "Recognition returned empty text", "пустой результат назван")
|
||||
assert.Contains(t, written, "level=WARN", "уровень — «может стать проблемой»")
|
||||
assert.Contains(t, written, record.Id, "запись названа идентификатором")
|
||||
}
|
||||
|
||||
// exhaustibleRecognizer отдаёт полный результат один раз, а на всяком следующем
|
||||
// обращении — пустой: так выглядит провайдер, чей поток закрылся на первом же
|
||||
// ответе. Отказом это не считается ни у него, ни у нас.
|
||||
type exhaustibleRecognizer struct {
|
||||
countingRecognizer
|
||||
served atomic.Bool
|
||||
}
|
||||
|
||||
func (r *exhaustibleRecognizer) Fetch(context.Context, string) (*entity.RecognitionOutcome, error) {
|
||||
r.fetches.Add(1)
|
||||
if r.served.Swap(true) {
|
||||
return &entity.RecognitionOutcome{}, nil
|
||||
}
|
||||
return r.Parse(countingPayload(t0Replicas))
|
||||
}
|
||||
|
||||
// Повторный опрос той же операции — обычное дело: держатель захвата умер,
|
||||
// сохранение рубежа отказало, человек снял остановку в панели. Пустой ответ
|
||||
// провайдера при этом не должен стирать уже сохранённую расшифровку: сервис
|
||||
// объявлен архивом, а восстановления у стёртого текста нет.
|
||||
func TestEmptySecondAnswerKeepsArchivedText(t *testing.T) {
|
||||
rec := &exhaustibleRecognizer{}
|
||||
env := newPipelineEnvWith(t, &okMetaViewer{}, &okConverter{}, rec)
|
||||
|
||||
record := newRecord(t, env)
|
||||
drain(t, env, record.Id)
|
||||
|
||||
stored := readRecord(t, env, record.Id)
|
||||
require.NotNil(t, stored.TranscriptTextID, "расшифровка сохранена первым ответом")
|
||||
|
||||
before, err := env.repos.Texts.GetByID(*stored.TranscriptTextID)
|
||||
require.NoError(t, err)
|
||||
require.NotEmpty(t, before.Contents)
|
||||
|
||||
// Второй опрос той же записи: рубеж возвращается на отправленный, и шаг
|
||||
// забирает результат заново — теперь пустой.
|
||||
back := readRecord(t, env, record.Id)
|
||||
back.MoveToState(entity.StateSubmitted)
|
||||
require.NoError(t, env.recordRepo.Save(back, ""))
|
||||
clearDelay(t, env, record.Id)
|
||||
require.NoError(t, env.service.RunStep(t.Context()))
|
||||
|
||||
// Проверка обязана дойти до второго ответа: иначе она зеленела бы, ничего не
|
||||
// проверив, — а класс «проверка, не способная упасть» в этом проекте ловили
|
||||
// уже трижды.
|
||||
require.EqualValues(t, 2, rec.fetches.Load(), "результат забирали дважды")
|
||||
|
||||
after, err := env.repos.Texts.GetByID(*stored.TranscriptTextID)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, before.Contents, after.Contents,
|
||||
"пустой ответ провайдера не стирает сохранённую расшифровку")
|
||||
}
|
||||
|
||||
@@ -37,7 +37,7 @@ func TestShutdownDuringConversionKeepsRecordRetryable(t *testing.T) {
|
||||
|
||||
converter := &killedConverter{cancel: cancel}
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, converter)
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
|
||||
err := env.service.RunStep(ctx)
|
||||
|
||||
@@ -51,14 +51,13 @@ func TestShutdownDuringConversionKeepsRecordRetryable(t *testing.T) {
|
||||
assert.Nil(t, after.AcquisitionID, "захват снят: запись возьмёт следующий прогон")
|
||||
assert.Equal(t, 0, after.Attempts, "остановка отказа не тратит")
|
||||
assert.Nil(t, after.ErrorText)
|
||||
assert.Empty(t, env.sender.sent(), "отправителю о несуществующем сбое не сообщают")
|
||||
}
|
||||
|
||||
// Запись, которую шаг не успел взять, потому что нас уже остановили, остаётся
|
||||
// нетронутой: захват не случился, отказ не потрачен.
|
||||
func TestShutdownBeforeStepLeavesRecordUntouched(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
|
||||
record := newTelegramRecord(t, env)
|
||||
record := newRecord(t, env)
|
||||
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
cancel()
|
||||
|
||||
+45
-139
@@ -66,7 +66,6 @@ type TranscribeService struct {
|
||||
metaviewer contract.AudioMetaViewer
|
||||
converter contract.AudioFileConverter
|
||||
recognizer contract.AudioRecognizer
|
||||
tgSender contract.TelegramMessageSender
|
||||
limits entity.StuckLimits
|
||||
logger *slog.Logger
|
||||
}
|
||||
@@ -76,7 +75,6 @@ func NewTranscribeService(
|
||||
metaviewer contract.AudioMetaViewer,
|
||||
converter contract.AudioFileConverter,
|
||||
recognizer contract.AudioRecognizer,
|
||||
tgSender contract.TelegramMessageSender,
|
||||
limits entity.StuckLimits,
|
||||
logger *slog.Logger,
|
||||
) *TranscribeService {
|
||||
@@ -88,7 +86,6 @@ func NewTranscribeService(
|
||||
metaviewer: metaviewer,
|
||||
converter: converter,
|
||||
recognizer: recognizer,
|
||||
tgSender: tgSender,
|
||||
limits: limits,
|
||||
logger: logger,
|
||||
}
|
||||
@@ -134,25 +131,14 @@ func (s *TranscribeService) stepFor(state string) (string, step, bool) {
|
||||
}
|
||||
}
|
||||
|
||||
func (s *TranscribeService) CreateJobFromTelegram(ctx context.Context, file io.Reader, fileName string, chatId int64, replyMsgId int) (*entity.AudioRecord, error) {
|
||||
record := &entity.AudioRecord{
|
||||
State: entity.StateUploaded,
|
||||
Source: entity.SourceTelegram,
|
||||
TgChatId: &chatId,
|
||||
TgReplyMessageId: &replyMsgId,
|
||||
}
|
||||
|
||||
return s.createRecord(ctx, record, file, fileName)
|
||||
}
|
||||
|
||||
// CreateJobFromApi заводит запись от имени вошедшего. Владелец обязателен:
|
||||
// пустой отвергается здесь, потому что колонка владельца допускает пустое
|
||||
// значение ради записей бота, и приём по HTTP — то место, где обязательность
|
||||
// держится.
|
||||
// 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")
|
||||
@@ -162,7 +148,7 @@ func (s *TranscribeService) CreateJobFromApi(ctx context.Context, file io.Reader
|
||||
record := &entity.AudioRecord{
|
||||
State: entity.StateUploaded,
|
||||
Source: entity.SourceApi,
|
||||
OwnerID: &ownerID,
|
||||
OwnerID: ownerID,
|
||||
}
|
||||
|
||||
return s.createRecord(ctx, record, file, fileName)
|
||||
@@ -212,7 +198,7 @@ func (s *TranscribeService) createRecord(ctx context.Context, r *entity.AudioRec
|
||||
DurationMs: int64(info.Seconds) * 1000,
|
||||
}
|
||||
|
||||
fileRecord, err := s.repos.Files.Create(storageFileName, work, meta, ownerOf(r))
|
||||
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
|
||||
@@ -253,8 +239,8 @@ func (s *TranscribeService) createRecord(ctx context.Context, r *entity.AudioRec
|
||||
//
|
||||
// Контекст доходит до шага, а через него — до внешнего собеседника: остановка
|
||||
// сервиса убивает `ffmpeg` и обрывает запрос к распознаванию. Прерванный шаг
|
||||
// приговора не выносит: запись остаётся пригодной к повтору, отказа не тратит и
|
||||
// отправителю о несуществующем сбое не сообщает.
|
||||
// приговора не выносит: запись остаётся пригодной к повтору и отказа не
|
||||
// тратит.
|
||||
func (s *TranscribeService) RunStep(ctx context.Context) error {
|
||||
// Нас уже остановили — запись не забираем: захват стоил бы ей отказа, а
|
||||
// работы всё равно не будет. Исход «шаг не сделал ничего» — это
|
||||
@@ -275,8 +261,7 @@ func (s *TranscribeService) RunStep(ctx context.Context) error {
|
||||
// таблицы. Молчать нельзя: запись выпала бы из работы без единого следа.
|
||||
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),
|
||||
"сервис не знает, что делать с этой записью")
|
||||
fmt.Sprintf("no step for state %s", record.State))
|
||||
return &contract.NoopJobError{State: record.State}
|
||||
}
|
||||
|
||||
@@ -347,8 +332,7 @@ func (s *TranscribeService) acquire() (*entity.AudioRecord, string, error) {
|
||||
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),
|
||||
"попытки исчерпаны")
|
||||
fmt.Sprintf("attempts exhausted: %d", record.Attempts))
|
||||
return nil, "", &contract.NoopJobError{State: record.State}
|
||||
}
|
||||
|
||||
@@ -360,8 +344,7 @@ func (s *TranscribeService) acquire() (*entity.AudioRecord, string, error) {
|
||||
"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),
|
||||
"обработка застряла")
|
||||
fmt.Sprintf("stuck in %s", record.State))
|
||||
return nil, "", &contract.NoopJobError{State: record.State}
|
||||
}
|
||||
|
||||
@@ -391,7 +374,7 @@ func (s *TranscribeService) normalize(ctx context.Context, r *entity.AudioRecord
|
||||
|
||||
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"), "у записи нет файла")
|
||||
return s.failStep(r, holder, stepNormalize, errors.New("record has no original file"))
|
||||
}
|
||||
|
||||
srcFile, err := s.repos.Files.GetByID(*r.OriginalFileID)
|
||||
@@ -441,7 +424,7 @@ func (s *TranscribeService) normalize(ctx context.Context, r *entity.AudioRecord
|
||||
|
||||
s.logger.Error("File conversion failed",
|
||||
"error", err, "record_id", r.Id, "duration", conversionDuration)
|
||||
return s.failStep(r, holder, stepNormalize, err, "сбой конвертации файла")
|
||||
return s.failStep(r, holder, stepNormalize, err)
|
||||
}
|
||||
|
||||
destSize, err := dest.Size()
|
||||
@@ -457,7 +440,7 @@ func (s *TranscribeService) normalize(ctx context.Context, r *entity.AudioRecord
|
||||
|
||||
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, ownerOf(r))
|
||||
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
|
||||
@@ -489,7 +472,7 @@ func (s *TranscribeService) submit(ctx context.Context, r *entity.AudioRecord, h
|
||||
|
||||
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"), "у записи нет приведённого файла")
|
||||
return s.failStep(r, holder, stepSubmit, errors.New("record has no normalized file"))
|
||||
}
|
||||
|
||||
fileRecord, err := s.repos.Files.GetByID(*r.NormalizedFileID)
|
||||
@@ -646,7 +629,7 @@ func (s *TranscribeService) uploadSource(
|
||||
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"), "сведений о распознавании нет")
|
||||
return s.failStep(r, holder, stepPoll, errors.New("record has no recognition attempt"))
|
||||
}
|
||||
|
||||
attempt, err := s.repos.Recognitions.GetByID(*r.RecognitionID)
|
||||
@@ -656,7 +639,7 @@ func (s *TranscribeService) poll(ctx context.Context, r *entity.AudioRecord, hol
|
||||
}
|
||||
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"), "сведений о распознавании нет")
|
||||
return s.failStep(r, holder, stepPoll, errors.New("recognition attempt has no operation id"))
|
||||
}
|
||||
|
||||
// Опрос идёт раз в несколько секунд всё время распознавания: часовая запись
|
||||
@@ -689,7 +672,7 @@ func (s *TranscribeService) poll(ctx context.Context, r *entity.AudioRecord, hol
|
||||
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), "сбой при распознавании файла")
|
||||
return s.failStep(r, holder, stepPoll, errors.New(errorText))
|
||||
}
|
||||
|
||||
outcome, err := s.recognizer.Fetch(ctx, attempt.ExternalID)
|
||||
@@ -708,6 +691,17 @@ func (s *TranscribeService) poll(ctx context.Context, r *entity.AudioRecord, hol
|
||||
"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)
|
||||
}
|
||||
|
||||
// Сырой ответ сохраняется целиком: результат операции у провайдера не
|
||||
// переспрашивается, и когда мы научимся размечать говорящих, архив
|
||||
// пересчитается из сохранённого без единого рубля.
|
||||
@@ -749,55 +743,42 @@ func (s *TranscribeService) storeOutcome(r *entity.AudioRecord, outcome *entity.
|
||||
return nil
|
||||
}
|
||||
|
||||
// finish отвечает отправителю и доводит запись до конечного рубежа. Доставка —
|
||||
// хвост последнего шага, а не отдельный узел конвейера.
|
||||
// finish доводит запись до конечного рубежа. Наружу шаг не обращается: доставки
|
||||
// ответа отправителю у сервиса нет, и свой исход отправитель узнаёт опросом
|
||||
// готовности.
|
||||
func (s *TranscribeService) finish(ctx context.Context, r *entity.AudioRecord, holder string) (stepOutcome, error) {
|
||||
text := "Ой, кажется, на аудиозаписи нет текста."
|
||||
if r.TranscriptTextID != nil {
|
||||
stored, err := s.repos.Texts.GetByID(*r.TranscriptTextID)
|
||||
if err != nil {
|
||||
s.logger.Error("Failed to read transcript", "error", err, "record_id", r.Id)
|
||||
return outcomeDone, err
|
||||
}
|
||||
if stored.Contents != "" {
|
||||
text = stored.Contents
|
||||
}
|
||||
}
|
||||
|
||||
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, s.send(r, text)
|
||||
return outcomeDone, nil
|
||||
}
|
||||
|
||||
// failStep останавливает запись приговором шага и сообщает отправителю.
|
||||
// failStep останавливает запись приговором шага.
|
||||
//
|
||||
// Исход возвращается **явно**: остановка отказом шага не является — шаг
|
||||
// рассудил об этой записи окончательно, — но и работой она не была, и
|
||||
// засчитывать её успехом нельзя.
|
||||
func (s *TranscribeService) failStep(r *entity.AudioRecord, holder, step string, stepErr error, humanText string) (stepOutcome, error) {
|
||||
s.halt(r, holder, step, entity.HaltReasonStepFailed, stepErr.Error(),
|
||||
fmt.Sprintf("При обработке записи произошла ошибка: %s", humanText))
|
||||
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 ставит признак остановки, считает её отказом, пишет строку журнала
|
||||
// событий и сообщает отправителю.
|
||||
// halt ставит признак остановки, считает её отказом и пишет строку журнала
|
||||
// событий.
|
||||
//
|
||||
// Сообщение уходит при **любой** причине остановки: инвариант проекта «Принятая
|
||||
// запись не теряется молча» допускает два исхода — запись пригодна к повтору
|
||||
// либо об отказе сказано, — а остановленная запись захвату не выдаётся, значит
|
||||
// первый исход исключён.
|
||||
// Отправителю отсюда ничего не уходит: инвариант проекта «Принятая запись не
|
||||
// теряется молча» держится теперь опросом готовности — остановка видна там
|
||||
// признаком — и журналом владельца, где у неё стоит причина.
|
||||
//
|
||||
// Счётчик растит **сама остановка**, а не воркер, и это не стилистика.
|
||||
// Остановка по сторожам наступает в захвате, до всякого шага, и воркер о ней
|
||||
// узнаёт признаком «работы нет» — а считать его в метрику запрещено инвариантом
|
||||
// «`NoopJobError` — не ошибка». Значит единственное место, где известны и факт
|
||||
// остановки, и рубеж, — здесь.
|
||||
func (s *TranscribeService) halt(r *entity.AudioRecord, holder, step, reason, errText, humanText string) {
|
||||
func (s *TranscribeService) halt(r *entity.AudioRecord, holder, step, reason, errText string) {
|
||||
// Рубеж читается до остановки: она его не двигает, но читать состояние после
|
||||
// мутации — привычка, из-за которой метка счётчика уже однажды разъехалась.
|
||||
stage := r.State
|
||||
@@ -823,8 +804,6 @@ func (s *TranscribeService) halt(r *entity.AudioRecord, holder, step, reason, er
|
||||
outcome = entity.EventOutcomeFailed
|
||||
}
|
||||
s.appendEvent(r, entity.EventOriginPipeline, step, outcome, reason, 0)
|
||||
|
||||
s.notify(r, humanText+"\nПожалуйста, попробуйте еще раз.")
|
||||
}
|
||||
|
||||
// scheduleRetry снимает захват с отказавшей записи и ставит нарастающую паузу.
|
||||
@@ -885,69 +864,6 @@ func (s *TranscribeService) appendEvent(r *entity.AudioRecord, origin, step, out
|
||||
}
|
||||
}
|
||||
|
||||
// send отвечает отправителю там, откуда пришла запись. Отказ отправки поднимает
|
||||
// вверх: он принадлежит шагу.
|
||||
//
|
||||
// Кроме недоставки — её шаг записывает и завершается без отказа. Ответ уходит
|
||||
// после того, как достигнутый рубеж сохранён: работа к этой минуте сделана, и
|
||||
// объявленный отказ засчитался бы воркеру сбоем и лёг бы владельцу записью
|
||||
// отказа. Повтор делу не помогает — ни бот, ни адресат от ожидания не появятся,
|
||||
// — поэтому причина недоставки живёт в журнале, а не в рубеже записи.
|
||||
func (s *TranscribeService) send(r *entity.AudioRecord, text string) error {
|
||||
if r.Source != entity.SourceTelegram {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Адресата у записи нет: отвечать некуда, и повторять нечего. Уровень здесь
|
||||
// выше, чем у неподнятого канала, и это не педантизм: пустой чат у записи
|
||||
// из Telegram — симптом порчи записи.
|
||||
if r.TgChatId == nil {
|
||||
s.undelivered(r, slog.LevelError, "chat is not specified")
|
||||
return nil
|
||||
}
|
||||
|
||||
if err := s.tgSender.Send(text, *r.TgChatId, r.TgReplyMessageId); err != nil {
|
||||
// Канал не поднят: сервис работает без этого входа, и это объявленный
|
||||
// режим, а не поломка.
|
||||
if errors.Is(err, contract.ErrDeliveryChannelDown) {
|
||||
s.undelivered(r, slog.LevelWarn, "delivery channel is down")
|
||||
return nil
|
||||
}
|
||||
|
||||
s.logger.Error("Failed to sent message to client", "record_id", r.Id)
|
||||
return fmt.Errorf("failed to sent message to client, record id: %s, err: %w", r.Id, err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// undelivered записывает недоставленный ответ и считает его в метрику. Уровень
|
||||
// приходит от причины: объявленный режим — «может стать проблемой», порча
|
||||
// записи — событие для разбора.
|
||||
//
|
||||
// Идентификатор записи обязателен, иначе владелец видит, что ответ не ушёл, но
|
||||
// не может найти, чей; текста ответа в записи нет — он содержимое чужой записи.
|
||||
func (s *TranscribeService) undelivered(r *entity.AudioRecord, level slog.Level, reason string) {
|
||||
metrics.UndeliveredReplyCounter.WithLabelValues(reason).Inc()
|
||||
|
||||
switch level {
|
||||
case slog.LevelError:
|
||||
s.logger.Error(undeliveredMessage, "record_id", r.Id, "reason", reason)
|
||||
default:
|
||||
s.logger.Warn(undeliveredMessage, "record_id", r.Id, "reason", reason)
|
||||
}
|
||||
}
|
||||
|
||||
const undeliveredMessage = "Reply was not delivered"
|
||||
|
||||
// notify отвечает отправителю там, где поднимать отказ некуда: запись уже
|
||||
// доведена до своего исхода, и отказ отправки остаётся записью в журнале.
|
||||
func (s *TranscribeService) notify(r *entity.AudioRecord, text string) {
|
||||
if err := s.send(r, text); err != nil {
|
||||
s.logger.Error("Failed to notify sender", "error", err, "record_id", r.Id)
|
||||
}
|
||||
}
|
||||
|
||||
// closeWork убирает рабочую копию. Отказ уборки не роняет шаг, но и не
|
||||
// проглатывается: забытая копия это шестичасовая запись во временном каталоге.
|
||||
func (s *TranscribeService) closeWork(work contract.WorkFile) {
|
||||
@@ -956,16 +872,6 @@ func (s *TranscribeService) closeWork(work contract.WorkFile) {
|
||||
}
|
||||
}
|
||||
|
||||
// ownerOf — владелец записи строкой; пустая значит «владельца нет», и таковы
|
||||
// записи, принятые ботом. Файл наследует владельца своей записи: правило
|
||||
// просмотра коллекции файлов сужено этой колонкой.
|
||||
func ownerOf(r *entity.AudioRecord) string {
|
||||
if r.OwnerID == nil {
|
||||
return ""
|
||||
}
|
||||
return *r.OwnerID
|
||||
}
|
||||
|
||||
// formatOf приводит расширение к виду колонки формата: без точки, в нижнем
|
||||
// регистре. Наружу оно выходит только приведённым к перечню известных форматов —
|
||||
// это делает метка метрики.
|
||||
|
||||
@@ -1,154 +0,0 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"log/slog"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"git.vakhrushev.me/av/transcriber/internal/adapter/repo/pocketbase/migrations"
|
||||
"git.vakhrushev.me/av/transcriber/internal/contract"
|
||||
"git.vakhrushev.me/av/transcriber/internal/entity"
|
||||
)
|
||||
|
||||
// Ответ отправителю уходит после того, как достигнутый рубеж сохранён. Значит,
|
||||
// недоставка не может быть отказом шага: объявленный отказ засчитался бы воркеру
|
||||
// сбоем, лёг бы владельцу записью отказа и переписал бы служебные поля
|
||||
// доведённой записи. Причин недоставки две, исход у них общий.
|
||||
|
||||
// downSender изображает неподнятый канал доставки: так ведёт себя заглушка,
|
||||
// которую ядро получает вместо отправителя Telegram.
|
||||
type downSender struct {
|
||||
calls int
|
||||
}
|
||||
|
||||
func (s *downSender) Send(string, int64, *int) error {
|
||||
s.calls++
|
||||
return contract.ErrDeliveryChannelDown
|
||||
}
|
||||
|
||||
// journalEnv пересобирает сервис с названным отправителем и своим журналом:
|
||||
// утверждения судят и состояние записи, и то, что увидел владелец.
|
||||
func journalEnv(
|
||||
t *testing.T,
|
||||
env *pipelineEnv,
|
||||
sender contract.TelegramMessageSender,
|
||||
) (*TranscribeService, *bytes.Buffer) {
|
||||
t.Helper()
|
||||
|
||||
journal := &bytes.Buffer{}
|
||||
svc := NewTranscribeService(
|
||||
env.repos,
|
||||
&okMetaViewer{},
|
||||
&okConverter{},
|
||||
env.service.recognizer,
|
||||
sender,
|
||||
testLimits,
|
||||
slog.New(slog.NewTextHandler(journal, &slog.HandlerOptions{Level: slog.LevelDebug})),
|
||||
)
|
||||
|
||||
return svc, journal
|
||||
}
|
||||
|
||||
// transcribedRecord доводит запись до рубежа, с которого уходит ответ.
|
||||
func transcribedRecord(t *testing.T, env *pipelineEnv, svc *TranscribeService) *entity.AudioRecord {
|
||||
t.Helper()
|
||||
|
||||
record := newTelegramRecord(t, env)
|
||||
for range 5 {
|
||||
clearDelay(t, env, record.Id)
|
||||
require.NoError(t, svc.RunStep(t.Context()))
|
||||
if readRecord(t, env, record.Id).State == entity.StateTranscribed {
|
||||
return readRecord(t, env, record.Id)
|
||||
}
|
||||
}
|
||||
t.Fatal("запись не дошла до рубежа расшифровки")
|
||||
return nil
|
||||
}
|
||||
|
||||
// Канал не поднят: запись доводится до конца, шаг отказа не объявляет, а
|
||||
// владелец узнаёт о недоставке из журнала.
|
||||
func TestUndeliveredOnDownChannelKeepsRecordDone(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{})
|
||||
|
||||
sender := &downSender{}
|
||||
svc, journal := journalEnv(t, env, sender)
|
||||
|
||||
record := transcribedRecord(t, env, svc)
|
||||
|
||||
// Шаг завершается без отказа — именно это воркер считает в свой счётчик.
|
||||
clearDelay(t, env, record.Id)
|
||||
require.NoError(t, svc.RunStep(t.Context()))
|
||||
|
||||
assert.Equal(t, 1, sender.calls, "ответ до отправителя доехал")
|
||||
|
||||
after := readRecord(t, env, record.Id)
|
||||
assert.Equal(t, entity.StateDone, after.State, "запись осталась на достигнутом рубеже")
|
||||
assert.False(t, after.IsHalted(), "отказ записи не приписан")
|
||||
require.NotNil(t, after.TranscriptTextID, "расшифровка сохранена")
|
||||
|
||||
text, err := env.repos.Texts.GetByID(*after.TranscriptTextID)
|
||||
require.NoError(t, err)
|
||||
require.NotEmpty(t, text.Contents)
|
||||
|
||||
written := journal.String()
|
||||
assert.Contains(t, written, "Reply was not delivered", "недоставка названа")
|
||||
assert.Contains(t, written, record.Id, "запись несёт идентификатор")
|
||||
assert.Contains(t, written, "level=WARN", "объявленный режим — «может стать проблемой»")
|
||||
assert.NotContains(t, written, text.Contents, "текста расшифровки в журнале нет")
|
||||
}
|
||||
|
||||
// Адресат у записи не назван: исход тот же. Прежде эта ветка объявляла отказ
|
||||
// шага на уже завершённой работе.
|
||||
func TestUndeliveredWithoutChatKeepsRecordDone(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{})
|
||||
|
||||
sender := &downSender{}
|
||||
svc, journal := journalEnv(t, env, sender)
|
||||
|
||||
record := transcribedRecord(t, env, svc)
|
||||
|
||||
// Запись из Telegram, у которой чат не назван: такую отдаёт правка в панели.
|
||||
// Колонка чистится мимо захвата — иначе подготовка унесла бы запись у шага.
|
||||
stored, err := env.app.FindRecordById(migrations.RecordsCollection, record.Id)
|
||||
require.NoError(t, err)
|
||||
stored.Set("tg_chat_id", nil)
|
||||
require.NoError(t, env.app.Save(stored))
|
||||
|
||||
clearDelay(t, env, record.Id)
|
||||
require.NoError(t, svc.RunStep(t.Context()))
|
||||
|
||||
assert.Equal(t, 0, sender.calls, "до отправителя дело не дошло: адресата нет")
|
||||
|
||||
after := readRecord(t, env, record.Id)
|
||||
assert.Equal(t, entity.StateDone, after.State)
|
||||
assert.False(t, after.IsHalted(), "отказ записи не приписан")
|
||||
|
||||
written := journal.String()
|
||||
assert.Contains(t, written, "Reply was not delivered")
|
||||
assert.Contains(t, written, record.Id)
|
||||
assert.Contains(t, written, "chat is not specified", "причина названа")
|
||||
assert.Contains(t, written, "level=ERROR",
|
||||
"порча записи громче штатного «бот не настроен»: иначе сигнал утонет")
|
||||
}
|
||||
|
||||
// Запись, принятая по HTTP, до отправителя не доходит вовсе: недоставки нет, и
|
||||
// записи о ней в журнале быть не должно — иначе журнал владельца заполнят
|
||||
// строки о записях основного входа.
|
||||
func TestApiRecordDoesNotReachSenderAndLogsNothing(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{})
|
||||
|
||||
record, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "voice.ogg", newOwner(t, env.app))
|
||||
require.NoError(t, err)
|
||||
|
||||
sender := &downSender{}
|
||||
svc, journal := journalEnv(t, env, sender)
|
||||
|
||||
require.NoError(t, svc.send(record, "расшифровка записи"))
|
||||
|
||||
assert.Equal(t, 0, sender.calls, "отправителя не звали")
|
||||
assert.NotContains(t, journal.String(), "Reply was not delivered", "недоставки не было")
|
||||
}
|
||||
Reference in New Issue
Block a user