доменные ошибки сравниваются через errors.As, отказ Close не теряется

- признаки «работы нет» и «задача не найдена» узнаются по смыслу, а не
  приведением типа: обёртка `%w` на пути больше не превращает пустой прогон
  воркера в отказ раз в секунду
- отказ закрытия соединения с распознавателем доходит до вызывающего
  (`errors.Join`) либо до журнала; у `errcheck` включён `check-blank`, иначе
  критерий принимал реализацию, выбрасывающую отказ в пустоту
- заведены первые тесты пакета worker и capability `pipeline`; долг из четырёх
  замечаний линтера закрыт, гейт зелёный целиком
This commit is contained in:
av
2026-08-11 18:04:25 +03:00
parent b7dd060ba0
commit 2559d09fc8
24 changed files with 1554 additions and 48 deletions
@@ -2,6 +2,7 @@ package yandex
import (
"context"
"errors"
"fmt"
"strings"
@@ -52,8 +53,16 @@ func newSpeechKitService(cfg speechKitConfig) (*speechKitService, error) {
// Создаем защищенное соединение для Operations API
opConn, err := grpc.NewClient(OperationEndpoint, grpc.WithTransportCredentials(creds))
if err != nil {
sttConn.Close()
return nil, fmt.Errorf("failed to connect to Operations API: %w", err)
// Отказы независимы, и второй не теряется. На сегодняшнем клиенте он
// почти наверняка не наступит: grpc.NewClient ленив, соединение к этому
// моменту не открыто, и Close вернёт отказ только при повторном
// закрытии — то есть сообщит о нашей ошибке, а не о Yandex. Сборка
// оставлена как защита от смены реализации клиента; nil от закрытия
// errors.Join отбрасывает, и форма ошибки в обычном случае не меняется.
return nil, errors.Join(
fmt.Errorf("failed to connect to Operations API: %w", err),
sttConn.Close(),
)
}
sttClient := stt.NewAsyncRecognizerClient(sttConn)
@@ -77,10 +86,10 @@ func (s *speechKitService) Close() error {
if s.opConn != nil {
err2 = s.opConn.Close()
}
if err1 != nil {
return err1
}
return err2
// Отказы двух соединений независимы, и вернуть только первый — значит
// потерять половину причины: журнал пишется при остановке процесса, и
// восстановить утраченное будет уже негде.
return errors.Join(err1, err2)
}
// recognizeFileFromS3 запускает асинхронное распознавание файла из S3
+6 -1
View File
@@ -2,6 +2,7 @@ package worker
import (
"context"
"errors"
"log/slog"
"strconv"
"time"
@@ -48,7 +49,11 @@ func (w *CallbackWorker) Start(ctx context.Context) {
return
default:
err := w.f()
_, isNoop := err.(*contract.NoopJobError)
// Признак узнаётся по смыслу, а не по точной форме значения:
// приведение типа видело только вершину цепочки и сломалось бы от
// первой же обёртки `%w`, которая в проекте — умолчание.
var noop *contract.NoopJobError
isNoop := errors.As(err, &noop)
if !isNoop {
metrics.WorkerJobCounter.WithLabelValues(w.Name(), strconv.FormatBool(err != nil)).Inc()
}
+204
View File
@@ -0,0 +1,204 @@
package worker
import (
"context"
"errors"
"fmt"
"log/slog"
"strings"
"sync"
"testing"
"time"
"git.vakhrushev.me/av/transcriber/internal/contract"
"github.com/prometheus/client_golang/prometheus"
)
// Проверки этого файла судят одну развилку воркера: пустой прогон против
// отказа. Инвариант проекта — «NoopJobError не ошибка» — стоит ровно на ней, а
// цена срабатывания отложенная: три воркера опрашивают базу раз в секунду, и
// пустой прогон, принятый за отказ, даёт три записи в секунду и столько же
// засчитанных сбоев, которых не было.
// journalBuffer собирает журнал прогона. Пишут в него из горутины воркера, а
// читает проверка — отсюда мьютекс.
type journalBuffer struct {
mu sync.Mutex
text strings.Builder
}
func (b *journalBuffer) Write(p []byte) (int, error) {
b.mu.Lock()
defer b.mu.Unlock()
return b.text.Write(p)
}
func (b *journalBuffer) String() string {
b.mu.Lock()
defer b.mu.Unlock()
return b.text.String()
}
// runOnce прогоняет воркер ровно один раз и возвращает журнал этого прогона.
// Цикл воркера бесконечен и спит секунду между прогонами, поэтому контекст
// отменяется сразу после первого вызова работы: ждать второго прогона нечего, а
// секунда сна на проверку — цена ни за что.
func runOnce(t *testing.T, name string, work func() error) string {
t.Helper()
journal := &journalBuffer{}
logger := slog.New(slog.NewTextHandler(journal, nil))
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
var once sync.Once
done := make(chan struct{})
w := NewCallbackWorker(name, func() error {
err := work()
once.Do(func() {
cancel()
close(done)
})
return err
}, logger)
finished := make(chan struct{})
go func() {
w.Start(ctx)
close(finished)
}()
select {
case <-done:
case <-time.After(5 * time.Second):
t.Fatal("работа воркера не была вызвана")
}
select {
case <-finished:
case <-time.After(5 * time.Second):
t.Fatal("воркер не остановился по отмене контекста")
}
return journal.String()
}
// runRecords оставляет от журнала только записи об исходе прогона. Жизненный
// цикл самого воркера — старт и остановка — по конвенции идёт на INFO и к
// прогону не относится; требование говорит о том, что воркер пишет про свой
// прогон, а не о том, что он молчит вообще.
func runRecords(journal string) string {
var kept []string
for _, line := range strings.Split(strings.TrimSpace(journal), "\n") {
if line == "" {
continue
}
if strings.Contains(line, "msg=\"Worker started\"") ||
strings.Contains(line, "msg=\"Worker received shutdown signal") {
continue
}
kept = append(kept, line)
}
return strings.Join(kept, "\n")
}
// jobCount читает счётчик работы воркера из общего реестра процесса. Судит
// реестр, а не переменную пакета: метка, потерянная в точке употребления,
// переменную не ломает, а на странице метрик видна.
func jobCount(t *testing.T, worker, errLabel string) float64 {
t.Helper()
families, err := prometheus.DefaultGatherer.Gather()
if err != nil {
t.Fatalf("не удалось собрать метрики: %v", err)
}
for _, mf := range families {
if mf.GetName() != "transcriber_worker_job_count" {
continue
}
for _, m := range mf.GetMetric() {
var gotWorker, gotErr string
for _, label := range m.GetLabel() {
switch label.GetName() {
case "name":
gotWorker = label.GetValue()
case "error":
gotErr = label.GetValue()
}
}
if gotWorker == worker && gotErr == errLabel {
return m.GetCounter().GetValue()
}
}
}
return 0
}
// Обёртка `%w` объявлена конвенцией проекта умолчанием, и до этой задачи первая
// же обёртка на пути сломала бы распознавание молча. Оракул держит именно
// обёрнутое значение: на голом признак узнавался и приведением типа, то есть
// проверка прошла бы и на починенном, и на сломанном коде.
func TestWrappedNoopIsNotAFailure(t *testing.T) {
const name = "wrapped_noop_worker"
before := jobCount(t, name, "false")
beforeErr := jobCount(t, name, "true")
journal := runOnce(t, name, func() error {
return fmt.Errorf("find and acquire job: %w", &contract.NoopJobError{State: "created"})
})
// Записи о старте и остановке воркера законны и к прогону не относятся —
// проверяется отсутствие записи об исходе прогона.
if got := runRecords(journal); got != "" {
t.Errorf("пустой прогон попал в журнал: %q", got)
}
if got := jobCount(t, name, "false"); got != before {
t.Errorf("счётчик успешных прогонов вырос на пустом прогоне: было %v, стало %v", before, got)
}
if got := jobCount(t, name, "true"); got != beforeErr {
t.Errorf("пустой прогон засчитан отказом: было %v, стало %v", beforeErr, got)
}
}
// Без этой проверки оракул был бы зелен и на коде, который не считает отказом
// вообще ничего.
func TestFailureIsLoggedAndCounted(t *testing.T) {
const name = "failing_worker"
before := jobCount(t, name, "true")
journal := runOnce(t, name, func() error {
return errors.New("database is gone")
})
if !strings.Contains(journal, "database is gone") {
t.Errorf("отказ не виден владельцу: журнал %q", journal)
}
if got := jobCount(t, name, "true"); got != before+1 {
t.Errorf("отказ не засчитан: было %v, стало %v", before, got)
}
}
// Счёт успешных прогонов — знаменатель доли отказов. Реализация, снявшая его,
// проходит обе проверки выше, а владелец теряет способность отличить «три
// прогона в секунду, все отказали» от «три отказа среди тысячи прогонов».
func TestSuccessIsCounted(t *testing.T) {
const name = "successful_worker"
before := jobCount(t, name, "false")
journal := runOnce(t, name, func() error {
return nil
})
if got := jobCount(t, name, "false"); got != before+1 {
t.Errorf("успешный прогон не засчитан: было %v, стало %v", before, got)
}
if strings.Contains(journal, "Worker error") {
t.Errorf("успешный прогон записан отказом: журнал %q", journal)
}
}
+78
View File
@@ -0,0 +1,78 @@
package service
import (
"errors"
"fmt"
"io"
"log/slog"
"testing"
"time"
"git.vakhrushev.me/av/transcriber/internal/contract"
"git.vakhrushev.me/av/transcriber/internal/entity"
)
// Путь признака «работы нет» состоит из двух звеньев: репозиторий рождает
// «подходящей задачи не нашлось», сервис переводит это в «работы нет», и уже
// его читает воркер. Проверки воркера подменяют работу целиком и второе звено
// не видят — без этого файла правку в сервисе принимал бы только линтер, а он
// судит форму записи, а не то, узнаётся ли признак на самом деле.
// stubJobRepo отдаёт заданную ошибку на запрос задачи. Прочих методов запроса
// задачи проверки этого файла не зовут.
type stubJobRepo struct {
err error
}
func (r *stubJobRepo) Create(*entity.TranscribeJob) error { return nil }
func (r *stubJobRepo) Save(*entity.TranscribeJob) error { return nil }
func (r *stubJobRepo) GetByID(string) (*entity.TranscribeJob, error) {
return nil, errors.New("не зовётся этими проверками")
}
func (r *stubJobRepo) FindAndAcquire(string, string, time.Time) (*entity.TranscribeJob, error) {
return nil, r.err
}
func serviceWithRepo(repo contract.TranscriptJobRepository) *TranscribeService {
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
return NewTranscribeService(repo, nil, nil, nil, nil, nil, "", logger)
}
// Репозиторий вправе добавить своему отказу пояснение — соседние ветки того же
// метода уже оборачивают ошибки `%w` подряд. Пока признак узнавался приведением
// типа, первая такая обёртка превратила бы пустой прогон в отказ: воркер начал
// бы писать в журнал раз в секунду на каждом из трёх воркеров.
func TestFindJobTranslatesWrappedNotFoundToNoop(t *testing.T) {
svc := serviceWithRepo(&stubJobRepo{
err: fmt.Errorf("find and acquire job: %w",
&contract.JobNotFoundError{State: "created", Message: "appropriate job not found"}),
})
_, err := svc.findJob("created", time.Minute)
var noop *contract.NoopJobError
if !errors.As(err, &noop) {
t.Fatalf("обёрнутое «задачи нет» не переведено в пустой прогон: получено %v", err)
}
if noop.State != "created" {
t.Errorf("состояние потеряно при переводе: %q", noop.State)
}
}
// Оборотная сторона: настоящий отказ хранилища пустым прогоном считаться не
// должен, иначе задача молча крутилась бы в цикле без единой записи.
func TestFindJobKeepsRealFailure(t *testing.T) {
svc := serviceWithRepo(&stubJobRepo{err: errors.New("database is gone")})
_, err := svc.findJob("created", time.Minute)
var noop *contract.NoopJobError
if errors.As(err, &noop) {
t.Fatalf("отказ хранилища зачтён пустым прогоном: %v", err)
}
if err == nil {
t.Fatal("отказ хранилища потерян")
}
}
+4 -1
View File
@@ -389,7 +389,10 @@ func (s *TranscribeService) findJob(state string, expiration time.Duration) (job
job, err = s.jobRepo.FindAndAcquire(state, acquisitionId, rottingTime)
if err != nil {
if _, ok := err.(*contract.JobNotFoundError); ok {
// Признак узнаётся по смыслу: репозиторий вправе обернуть свой отказ
// пояснением, и приведение типа от этого сломалось бы молча.
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)