ingest: закрыты мелочи приёма — вырожденное имя, контракт Result, корреляция add

- имя раздачи нормализуется на границе разбора: вырожденное `-`
  (metainfo.NoName) даёт пустое имя, пробельное схлопывается — сентинел больше
  не доходит ни до контекста распознавания, ни до source_ref, ни до подсказки
  вывода имени
- контракт «на любом пути ошибки приёма результат нулевой» объявлен в ingest и
  удерживается структурно; три транспорта перестали обещать идентификатор,
  которого нет, и коррелируют отказ по request_id
- scoped-логгер загрузки ставится до вызова внешнего сервиса в семи командах
  воркера — записи об отказе qBittorrent и метабаз получили download_id
  и infohash; граница разбора bencode записана в docs/research
This commit is contained in:
av
2026-08-06 18:20:12 +03:00
parent 52615e4e49
commit d081ef1d30
30 changed files with 2190 additions and 63 deletions
+3 -1
View File
@@ -18,13 +18,15 @@ type fakeNamer struct {
name string
fields naming.Fields // извлечённая структура (ok=true, если Title непуст)
gotContext string
gotHint string // подсказка имени — третий потребитель torrent.Info.DisplayName
calls int
onCall func()
}
func (f *fakeNamer) Derive(_ context.Context, contextText, _ string) (string, naming.Fields, bool) {
func (f *fakeNamer) Derive(_ context.Context, contextText, hint string) (string, naming.Fields, bool) {
f.calls++
f.gotContext = contextText
f.gotHint = hint
if f.onCall != nil {
f.onCall()
}
+26 -5
View File
@@ -404,6 +404,10 @@ func (w *Worker) Relink(ctx context.Context, id string) (err error) {
if d.State != store.StateReverted && d.State != store.StateCancelled && d.State != store.StateTargetMissing {
return fmt.Errorf("relink: download %s is in state %s (expected reverted/cancelled/target_missing): %w", id, d.State, ErrConflict)
}
// Scoped-логгер загрузки — ДО первого вызова внешнего сервиса (обязанность
// вызывающего, capability identity): иначе запись клиента об отказе уходит
// без download_id/infohash.
ctx = w.scoped(ctx, capReview, id, d.PrimaryInfohash())
// Источник нужен для распознавания и должен быть докачан — проверяем
// синхронно (без дебаунса); отсутствие приводит состояние к реальности
// (orphaned/deleted), недокачанный — ErrNotReady.
@@ -414,7 +418,6 @@ func (w *Worker) Relink(ctx context.Context, id string) (err error) {
if err := w.store.SetOverride(ctx, id, ovrForceReview, "1"); err != nil {
return fmt.Errorf("relink: %w", err)
}
ctx = w.scoped(ctx, capReview, id, d.PrimaryInfohash())
// Возврат в активную обработку — только через атомарный гард инварианта
// «не более одной активной задачи на infohash» (см. design ulid-identity, D4).
if err := w.store.ActivateIfNoOtherActive(ctx, id, store.StateRecognizing, "", ""); err != nil {
@@ -439,10 +442,13 @@ func (w *Worker) Rerecognize(ctx context.Context, id string) (err error) {
if err != nil {
return err
}
// Scoped-логгер загрузки — ДО первого вызова внешнего сервиса (обязанность
// вызывающего, capability identity): иначе запись клиента об отказе уходит
// без download_id/infohash.
ctx = w.scoped(ctx, capReview, id, d.PrimaryInfohash())
if err := w.ensureSourceReady(ctx, d, "rerecognize"); err != nil {
return err
}
ctx = w.scoped(ctx, capReview, id, d.PrimaryInfohash())
logctx.From(ctx).Info("review re-recognizing without hint")
w.transition(ctx, *d, store.StateRecognizing, "", "")
return nil
@@ -462,10 +468,13 @@ func (w *Worker) Refine(ctx context.Context, id string, hint string) (err error)
if err != nil {
return err
}
// Scoped-логгер загрузки — ДО первого вызова внешнего сервиса (обязанность
// вызывающего, capability identity): иначе запись клиента об отказе уходит
// без download_id/infohash.
ctx = w.scoped(ctx, capReview, id, d.PrimaryInfohash())
if err := w.ensureSourceReady(ctx, d, "refine"); err != nil {
return err
}
ctx = w.scoped(ctx, capReview, id, d.PrimaryInfohash())
if err := w.store.AddHint(ctx, id, hint); err != nil {
return fmt.Errorf("refine: %w", err)
}
@@ -688,6 +697,10 @@ func (w *Worker) ChooseCandidate(ctx context.Context, id, candidateID string) (e
if err != nil {
return err
}
// Scoped-логгер загрузки — ДО первого вызова внешнего сервиса (обязанность
// вызывающего, capability identity): иначе запись клиента об отказе уходит
// без download_id/infohash.
ctx = w.scoped(ctx, capReview, id, d.PrimaryInfohash())
rec, err := w.store.GetCurrentRecognition(ctx, id)
if err != nil {
return fmt.Errorf("choose candidate: %w", err)
@@ -724,6 +737,10 @@ func (w *Worker) AddManualSource(ctx context.Context, id, provider, providerID s
if err != nil {
return err
}
// Scoped-логгер загрузки — ДО первого вызова внешнего сервиса (обязанность
// вызывающего, capability identity): иначе запись клиента об отказе уходит
// без download_id/infohash.
ctx = w.scoped(ctx, capReview, id, d.PrimaryInfohash())
rec, err := w.store.GetCurrentRecognition(ctx, id)
if err != nil {
return fmt.Errorf("add source: %w", err)
@@ -809,7 +826,7 @@ func (w *Worker) chooseCandidateLocked(ctx context.Context, id string, d *store.
// раздачи (best-effort, косметика). Сбой обновления имени не должен ронять
// выбор кандидата: логируем и продолжаем.
if err := w.refreshDisplayNameLocked(ctx, id); err != nil {
logctx.From(w.scoped(ctx, capReview, id, d.PrimaryInfohash())).
logctx.From(ctx).
Warn("display name refresh after candidate choice failed", "error", err)
}
return nil
@@ -835,6 +852,10 @@ func (w *Worker) SetProviderID(ctx context.Context, id string, provider, provide
if err != nil {
return err
}
// Scoped-логгер загрузки — ДО первого вызова внешнего сервиса (обязанность
// вызывающего, capability identity): иначе запись клиента об отказе уходит
// без download_id/infohash.
ctx = w.scoped(ctx, capReview, id, d.PrimaryInfohash())
// Режиссёр вручную заданного источника из метабазы (credits) — best-effort
// косметика: сбой чтения рекогниции не валит смену источника (тип по умолчанию
// movie, как трактует candMediaType(nil)).
@@ -842,7 +863,7 @@ func (w *Worker) SetProviderID(ctx context.Context, id string, provider, provide
if w.recognizer != nil {
rec, rerr := w.store.GetCurrentRecognition(ctx, id)
if rerr != nil {
logctx.From(w.scoped(ctx, capReview, id, d.PrimaryInfohash())).
logctx.From(ctx).
Warn("set provider: recognition lookup for director failed", "error", rerr)
rec = nil
}
+162
View File
@@ -5,11 +5,14 @@ import (
"context"
"crypto/sha1"
"encoding/hex"
"errors"
"log/slog"
"testing"
"github.com/anacrolix/torrent/bencode"
"github.com/anacrolix/torrent/metainfo"
"git.vakhrushev.me/av/jellybit/internal/logctx"
"git.vakhrushev.me/av/jellybit/internal/store"
)
@@ -118,3 +121,162 @@ func TestRetryTorrentMissingBytesRollsBack(t *testing.T) {
t.Errorf("state = %q, want failed (откат активации)", st.downloads["1"].State)
}
}
// Третий потребитель torrent.Info.DisplayName — подсказка вывода отображаемого
// имени, то есть вход LLM. Раздача, объявившая вырожденное имя `-`, не должна
// отдавать его подсказкой: нормализация стоит на границе разбора, и здесь
// проверяется, что она туда доехала (собственный оракул, а не транзитивный).
func TestProcessCatchedNoNameGivesEmptyHint(t *testing.T) {
data, ih := buildTorrent(t, metainfo.NoName)
st := catchedTorrentStore("1", data, ih)
qb := &fakeQbt{}
w := newTestWorker(st, qb)
namer := &fakeNamer{name: "Дюна (2024)"}
w.SetNamer(namer)
w.processCatched(context.Background())
if namer.calls == 0 {
t.Fatal("namer не вызывался — тест ничего не проверил")
}
if namer.gotHint != "" {
t.Errorf("подсказка имени = %q, want пусто (вырожденное имя раздачи)", namer.gotHint)
}
}
// probeHandler захватывает атрибуты, доклеенные логгером. Свой WithAttrs
// обязателен: делегированный вернул бы чужой хендлер, и scoped-логгер (он
// собирается через With) перестал бы захватываться.
type probeHandler struct {
attrs *[]slog.Attr
}
func (h probeHandler) Enabled(context.Context, slog.Level) bool { return true }
func (h probeHandler) Handle(context.Context, slog.Record) error { return nil }
func (h probeHandler) WithAttrs(as []slog.Attr) slog.Handler {
*h.attrs = append(*h.attrs, as...)
return h
}
func (h probeHandler) WithGroup(string) slog.Handler { return h }
// Обязанность вызывающего (capability identity): команда, делающая вызов
// внешнего сервиса в контексте загрузки, кладёт scoped-логгер в ctx ДО этого
// вызова. Без этого запись клиента qBittorrent об отказе добавления уходит без
// download_id/infohash, а причины отказа qBittorrent не сообщает («Fails.») —
// корреляция там единственное, что делает запись пригодной для разбора.
func TestRetryPassesScopedContextToAdd(t *testing.T) {
data, ih := buildTorrent(t, "Fargo.mkv")
st := catchedTorrentStore("1", data, ih)
st.downloads["1"].State = store.StateFailed
qb := &fakeQbt{}
w := newTestWorker(st, qb)
var attrs []slog.Attr
w.log = slog.New(probeHandler{attrs: &attrs})
if err := w.Retry(context.Background(), "1"); err != nil {
t.Fatalf("Retry: %v", err)
}
if len(qb.addCtx) != 1 {
t.Fatalf("вызовов Add = %d, want 1", len(qb.addCtx))
}
if logctx.FromOr(qb.addCtx[0], nil) == nil {
t.Fatal("в ctx вызова Add нет scoped-логгера загрузки")
}
got := map[string]string{}
for _, a := range attrs {
got[a.Key] = a.Value.String()
}
if got["download_id"] != "1" {
t.Errorf("download_id = %q, want 1", got["download_id"])
}
if got["infohash"] != ih {
t.Errorf("infohash = %q, want %q", got["infohash"], ih)
}
if got["capability"] == "" {
t.Errorf("нет capability в scoped-логгере: %v", got)
}
}
// Best-effort ветки Retry: откат активации при провале Add и сбросы базиса
// таймаутов/счётчика пропусков. Ветки переписаны этим изменением (ручные
// w.log.… схлопнуты в logctx.From(ctx)), поэтому проверяем не только исход, но
// и что аварийная запись несёт корреляцию из scoped-контекста — ради неё
// scoped-логгер сюда и заводился.
func TestRetryBestEffortBranchesCarryScope(t *testing.T) {
t.Run("откат активации при провале Add", func(t *testing.T) {
data, ih := buildTorrent(t, "Fargo.mkv")
st := catchedTorrentStore("1", data, ih)
st.downloads["1"].State = store.StateFailed
qb := &fakeQbt{addErr: errors.New("qbit down")}
w := newTestWorker(st, qb)
st.setStateErr = errors.New("db down") // и откат тоже не проходит
var attrs []slog.Attr
w.log = slog.New(probeHandler{attrs: &attrs})
if err := w.Retry(context.Background(), "1"); err == nil {
t.Fatal("ожидалась ошибка добавления")
}
assertScopeAttrs(t, attrs, ih)
})
t.Run("сброс базиса таймаутов не прошёл", func(t *testing.T) {
data, ih := buildTorrent(t, "Fargo.mkv")
st := catchedTorrentStore("1", data, ih)
st.downloads["1"].State = store.StateFailed
st.retriedErr = errors.New("db down")
qb := &fakeQbt{}
w := newTestWorker(st, qb)
var attrs []slog.Attr
w.log = slog.New(probeHandler{attrs: &attrs})
// Best-effort: сам retry состоялся, несмотря на отказ сброса.
if err := w.Retry(context.Background(), "1"); err != nil {
t.Fatalf("Retry: %v", err)
}
if st.downloads["1"].State != store.StateDownloading {
t.Errorf("state = %q, want downloading", st.downloads["1"].State)
}
assertScopeAttrs(t, attrs, ih)
})
t.Run("сброс счётчика пропусков не прошёл", func(t *testing.T) {
data, ih := buildTorrent(t, "Fargo.mkv")
st := catchedTorrentStore("1", data, ih)
st.downloads["1"].State = store.StateFailed
st.downloads["1"].SourceMissCount = 3
st.missCntErr = errors.New("db down")
qb := &fakeQbt{}
w := newTestWorker(st, qb)
var attrs []slog.Attr
w.log = slog.New(probeHandler{attrs: &attrs})
if err := w.Retry(context.Background(), "1"); err != nil {
t.Fatalf("Retry: %v", err)
}
assertScopeAttrs(t, attrs, ih)
})
}
// assertScopeAttrs проверяет, что scoped-логгер загрузки собран и несёт
// корреляционные поля.
func assertScopeAttrs(t *testing.T, attrs []slog.Attr, infohash string) {
t.Helper()
got := map[string]string{}
for _, a := range attrs {
got[a.Key] = a.Value.String()
}
if got["download_id"] != "1" {
t.Errorf("download_id = %q, want 1", got["download_id"])
}
if got["infohash"] != infohash {
t.Errorf("infohash = %q, want %q", got["infohash"], infohash)
}
}
+11 -9
View File
@@ -1035,6 +1035,12 @@ func (w *Worker) Retry(ctx context.Context, id string) (err error) {
if d.State != store.StateFailed && d.State != store.StateStuck {
return fmt.Errorf("retry: download %s is %s, only failed/stuck are retriable: %w", id, d.State, ErrConflict)
}
// Scoped-логгер загрузки — ДО первого вызова внешнего сервиса (обязанность
// вызывающего, capability identity). Иначе запись клиента qBittorrent об
// отказе добавления уходит без download_id/infohash, а причины отказа
// qBittorrent не сообщает («Fails.») — корреляция там единственное, что
// делает запись пригодной для разбора. Форма та же, что у Delete.
ctx = w.scoped(ctx, capReview, id, d.PrimaryInfohash())
// Если раздача уже жива и ЗДОРОВА в qBittorrent — перецепляемся к ней,
// повторный Add не нужен (и вреден: вслепую дублировал бы торрент). Add — когда
// источника в qBittorrent нет. Базис таймаута сбрасывается ниже через
@@ -1076,8 +1082,7 @@ func (w *Worker) Retry(ctx context.Context, id string) (err error) {
// Байты torrent недоступны — откатываем активацию, задача не должна
// «качаться» без раздачи в qBittorrent.
if rbErr := w.store.SetDownloadState(ctx, id, d.State, d.ErrorCode.String, d.ErrorMsg.String); rbErr != nil {
w.log.Error("retry rollback failed",
"capability", capReview, "download_id", id, "error", rbErr)
logctx.From(ctx).Error("retry rollback failed", "error", rbErr)
}
return fmt.Errorf("retry: prepare add: %w", prepErr)
}
@@ -1085,8 +1090,7 @@ func (w *Worker) Retry(ctx context.Context, id string) (err error) {
// Активация уже прошла — откатываем задачу в прежнее состояние,
// чтобы не оставить «качающуюся» задачу без раздачи в qBittorrent.
if rbErr := w.store.SetDownloadState(ctx, id, d.State, d.ErrorCode.String, d.ErrorMsg.String); rbErr != nil {
w.log.Error("retry rollback failed",
"capability", capReview, "download_id", id, "error", rbErr)
logctx.From(ctx).Error("retry rollback failed", "error", rbErr)
}
return fmt.Errorf("retry: add to qbittorrent: %w", err)
}
@@ -1096,8 +1100,7 @@ func (w *Worker) Retry(ctx context.Context, id string) (err error) {
// stuck_after на ближайшем тике. Best-effort: сбой лишь лишает свежего окна
// (WARN), сам retry уже состоялся.
if err := w.store.SetRetriedAt(ctx, id, w.now()); err != nil {
w.log.Warn("retry basis reset failed",
"capability", capReview, "download_id", id, "error", err)
logctx.From(ctx).Warn("retry basis reset failed", "error", err)
}
// Сброс счётчика пропусков источника: retried source_gone-задача иначе вошла бы
// в downloading с source_miss_count == threshold и упала бы снова на ближайшем
@@ -1105,11 +1108,10 @@ func (w *Worker) Retry(ctx context.Context, id string) (err error) {
// обещанного грейс-окна (MAJOR-3). Best-effort: сбой лишь лишает свежего окна.
if d.SourceMissCount != 0 {
if err := w.store.SetSourceMissCount(ctx, id, 0); err != nil {
w.log.Warn("retry miss count reset failed",
"capability", capReview, "download_id", id, "error", err)
logctx.From(ctx).Warn("retry miss count reset failed", "error", err)
}
}
logctx.From(w.scoped(ctx, capReview, id, d.PrimaryInfohash())).Info("state transition",
logctx.From(ctx).Info("state transition",
"from", d.State, "to", store.StateDownloading)
return nil
}
+15 -1
View File
@@ -24,6 +24,9 @@ type fakeStore struct {
transitions []transition
torrents map[string][]byte // download_id → байты .torrent
promoteErr error // если задан — PromoteCatched возвращает его, НЕ меняя state (симуляция транзиентного сбоя БД)
setStateErr error // если задан — SetDownloadState отказывает (ветка отката Retry)
retriedErr error // если задан — SetRetriedAt отказывает (best-effort сброс базиса)
missCntErr error // если задан — SetSourceMissCount отказывает (best-effort сброс счётчика)
}
type transition struct {
@@ -166,6 +169,9 @@ func (f *fakeStore) AddInfohashes(_ context.Context, id string, hashes []string)
}
func (f *fakeStore) SetDownloadState(_ context.Context, id string, st store.State, code, msg string) error {
if f.setStateErr != nil {
return f.setStateErr
}
d, ok := f.downloads[id]
if !ok {
return fmt.Errorf("download %s not found", id)
@@ -213,6 +219,9 @@ func (f *fakeStore) SetParsedContext(_ context.Context, id, jsonStr string) erro
}
func (f *fakeStore) SetSourceMissCount(_ context.Context, id string, n int) error {
if f.missCntErr != nil {
return f.missCntErr
}
d, ok := f.downloads[id]
if !ok {
return fmt.Errorf("download %s not found", id)
@@ -233,6 +242,9 @@ func (f *fakeStore) SetSourceAddedAt(_ context.Context, id string, t time.Time)
}
func (f *fakeStore) SetRetriedAt(_ context.Context, id string, t time.Time) error {
if f.retriedErr != nil {
return f.retriedErr
}
d, ok := f.downloads[id]
if !ok {
return fmt.Errorf("download %s not found", id)
@@ -283,6 +295,7 @@ type fakeQbt struct {
torrentsErr error
onTorrents func() // вклинивается в момент листинга (симуляция гонки между снимком и re-read)
added []qbt.AddRequest
addCtx []context.Context // ctx каждого вызова Add (проверка scoped-логгера)
addErr error
onAdd func() // вклинивается в момент Add (симуляция отмены в окне после add)
files []qbt.File
@@ -321,7 +334,8 @@ func (f *fakeQbt) Torrents(_ context.Context, category string) ([]qbt.Torrent, e
return out, nil
}
func (f *fakeQbt) Add(_ context.Context, ar qbt.AddRequest) error {
func (f *fakeQbt) Add(ctx context.Context, ar qbt.AddRequest) error {
f.addCtx = append(f.addCtx, ctx)
if f.onAdd != nil {
f.onAdd()
}