- имя раздачи нормализуется на границе разбора: вырожденное `-` (metainfo.NoName) даёт пустое имя, пробельное схлопывается — сентинел больше не доходит ни до контекста распознавания, ни до source_ref, ни до подсказки вывода имени - контракт «на любом пути ошибки приёма результат нулевой» объявлен в ingest и удерживается структурно; три транспорта перестали обещать идентификатор, которого нет, и коррелируют отказ по request_id - scoped-логгер загрузки ставится до вызова внешнего сервиса в семи командах воркера — записи об отказе qBittorrent и метабаз получили download_id и infohash; граница разбора bencode записана в docs/research
283 lines
11 KiB
Go
283 lines
11 KiB
Go
package worker
|
||
|
||
import (
|
||
"bytes"
|
||
"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"
|
||
)
|
||
|
||
// buildTorrent — валидные байты .torrent и их v1-инфохэш.
|
||
func buildTorrent(t *testing.T, name string) (data []byte, infohash string) {
|
||
t.Helper()
|
||
info := metainfo.Info{Name: name, Length: 2048, PieceLength: 1024, Pieces: make([]byte, 40)}
|
||
infoBytes, err := bencode.Marshal(info)
|
||
if err != nil {
|
||
t.Fatalf("marshal info: %v", err)
|
||
}
|
||
sum := sha1.Sum(infoBytes)
|
||
mi := metainfo.MetaInfo{InfoBytes: infoBytes, Announce: "http://t/ann"}
|
||
var buf bytes.Buffer
|
||
if err := mi.Write(&buf); err != nil {
|
||
t.Fatalf("write metainfo: %v", err)
|
||
}
|
||
return buf.Bytes(), hex.EncodeToString(sum[:])
|
||
}
|
||
|
||
// catchedTorrentStore — пойманная torrent-загрузка с сохранёнными байтами.
|
||
func catchedTorrentStore(id string, data []byte, infohash string) *fakeStore {
|
||
return &fakeStore{
|
||
downloads: map[string]*store.Download{
|
||
id: {
|
||
ID: id,
|
||
State: store.StateCatched,
|
||
SourceType: store.SourceTorrent,
|
||
SourceRef: "Some.Release.Name", // человекочитаемый референс, не URL
|
||
Infohashes: hashesOf(id, infohash),
|
||
CreatedAt: nowStr,
|
||
},
|
||
},
|
||
torrents: map[string][]byte{id: data},
|
||
}
|
||
}
|
||
|
||
// Пойманная torrent-загрузка добавляется в qBittorrent ФАЙЛОМ (Torrents), не URL.
|
||
func TestProcessCatchedTorrentAddsAsFile(t *testing.T) {
|
||
data, ih := buildTorrent(t, "Dune.mkv")
|
||
st := catchedTorrentStore("1", data, ih)
|
||
qb := &fakeQbt{}
|
||
w := newTestWorker(st, qb)
|
||
w.SetNamer(&fakeNamer{name: "Дюна (2024)"})
|
||
|
||
w.processCatched(context.Background())
|
||
|
||
if len(qb.added) != 1 {
|
||
t.Fatalf("qbt.Add calls = %d, want 1", len(qb.added))
|
||
}
|
||
add := qb.added[0]
|
||
if len(add.URLs) != 0 {
|
||
t.Errorf("torrent не должен добавляться по URLs: %v", add.URLs)
|
||
}
|
||
if len(add.Torrents) != 1 || !bytes.Equal(add.Torrents[0], data) {
|
||
t.Errorf("байты .torrent не переданы в Add.Torrents")
|
||
}
|
||
if add.Rename != "Дюна (2024)" {
|
||
t.Errorf("rename = %q", add.Rename)
|
||
}
|
||
if st.downloads["1"].State != store.StateDownloading {
|
||
t.Errorf("state = %q, want downloading", st.downloads["1"].State)
|
||
}
|
||
}
|
||
|
||
// Retry torrent-загрузки без живой раздачи добавляет её ФАЙЛОМ (регрессия
|
||
// BLOCKER 1: раньше add гейтился на magnet и torrent активировался без раздачи).
|
||
func TestRetryTorrentAddsAsFile(t *testing.T) {
|
||
data, ih := buildTorrent(t, "Fargo.mkv")
|
||
st := catchedTorrentStore("1", data, ih)
|
||
st.downloads["1"].State = store.StateFailed
|
||
qb := &fakeQbt{} // раздачи в qBittorrent нет → должен добавить заново
|
||
w := newTestWorker(st, qb)
|
||
|
||
if err := w.Retry(context.Background(), "1"); err != nil {
|
||
t.Fatalf("Retry: %v", err)
|
||
}
|
||
if len(qb.added) != 1 {
|
||
t.Fatalf("ожидалось повторное добавление, got %d", len(qb.added))
|
||
}
|
||
if len(qb.added[0].Torrents) != 1 || len(qb.added[0].URLs) != 0 {
|
||
t.Errorf("retry torrent должен добавляться файлом, add = %+v", qb.added[0])
|
||
}
|
||
if st.downloads["1"].State != store.StateDownloading {
|
||
t.Errorf("state = %q, want downloading", st.downloads["1"].State)
|
||
}
|
||
}
|
||
|
||
// Retry torrent-загрузки без сохранённых байтов — откат активации (не оставляем
|
||
// «качающуюся» задачу без раздачи).
|
||
func TestRetryTorrentMissingBytesRollsBack(t *testing.T) {
|
||
_, ih := buildTorrent(t, "X.mkv")
|
||
st := catchedTorrentStore("1", nil, ih) // байтов нет
|
||
delete(st.torrents, "1")
|
||
st.downloads["1"].State = store.StateFailed
|
||
qb := &fakeQbt{}
|
||
w := newTestWorker(st, qb)
|
||
|
||
if err := w.Retry(context.Background(), "1"); err == nil {
|
||
t.Fatal("ожидалась ошибка (нет байтов torrent)")
|
||
}
|
||
if len(qb.added) != 0 {
|
||
t.Errorf("add не должен вызываться без байтов")
|
||
}
|
||
if st.downloads["1"].State != store.StateFailed {
|
||
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)
|
||
}
|
||
}
|