package ingest import ( "context" "errors" "io" "log/slog" "strings" "testing" "git.vakhrushev.me/av/jellybit/internal/ident" "git.vakhrushev.me/av/jellybit/internal/store" ) const sampleMagnet = "magnet:?xt=urn:btih:541ADCFF3B6DD5DBA7088EA83317D9D6FAC331D6&dn=Dune" const sampleInfohash = "541adcff3b6dd5dba7088ea83317d9d6fac331d6" type fakeStore struct { active *store.Download created []store.Download hashes [][]string blobs [][]byte toppedUp []string upgradeID string // downloadID последнего вызова UpgradeCatchedMagnetToTorrent upgradeBlob []byte // байты, переданные в апгрейд upgradeUp bool // что вернуть из UpgradeCatchedMagnetToTorrent } func (f *fakeStore) FindReingestBlockingByInfohash(_ context.Context, _ ...string) (*store.Download, error) { return f.active, nil } func (f *fakeStore) CreateDownloadIfNoActive(_ context.Context, d *store.Download, hashes []string, torrentBlob []byte) (*store.Download, error) { if f.active != nil { return f.active, nil } d.ID = ident.NewID() f.created = append(f.created, *d) f.hashes = append(f.hashes, hashes) f.blobs = append(f.blobs, torrentBlob) return nil, nil } func (f *fakeStore) AddInfohashes(_ context.Context, id string, hashes []string) error { f.toppedUp = append(f.toppedUp, hashes...) _ = id return nil } func (f *fakeStore) UpgradeCatchedMagnetToTorrent(_ context.Context, id string, blob []byte) (bool, error) { f.upgradeID = id f.upgradeBlob = blob return f.upgradeUp, nil } // raceStore моделирует гонку F8: пред-рид FindReingestBlockingByInfohash видит // активную запись (blocking), но create-гард CreateDownloadIfNoActive её уже не // находит (в параллели отменена) и заводит свежую задачу. type raceStore struct { blocking *store.Download created []store.Download } func (r *raceStore) FindReingestBlockingByInfohash(_ context.Context, _ ...string) (*store.Download, error) { return r.blocking, nil } func (r *raceStore) CreateDownloadIfNoActive(_ context.Context, d *store.Download, _ []string, _ []byte) (*store.Download, error) { d.ID = ident.NewID() r.created = append(r.created, *d) return nil, nil // активной уже нет — создаём новую } func (r *raceStore) AddInfohashes(_ context.Context, _ string, _ []string) error { return nil } func (r *raceStore) UpgradeCatchedMagnetToTorrent(_ context.Context, _ string, _ []byte) (bool, error) { return false, nil } func newService(st Store) *Service { return New(st, slog.New(slog.NewTextHandler(io.Discard, nil))) } // Быстрый приём: сохраняем загрузку в catched и сразу отвечаем; qBittorrent и // вывод имени в пути приёма не участвуют (это делает worker). func TestIngestCatchesFast(t *testing.T) { fs := &fakeStore{} res, err := newService(fs).Ingest(context.Background(), Request{Source: sampleMagnet, Context: "Дюна 2"}) if err != nil { t.Fatalf("Ingest: %v", err) } if len(res.Infohashes) != 1 || res.Infohashes[0] != sampleInfohash { t.Errorf("infohashes = %v", res.Infohashes) } if res.State != store.StateCatched || res.Deduplicated { t.Errorf("res = %+v", res) } if len(fs.created) != 1 { t.Fatalf("создано задач: %d, want 1", len(fs.created)) } got := fs.created[0] if got.State != store.StateCatched { t.Errorf("state задачи = %q, want catched", got.State) } // Имя выводит worker на шаге добавления — при приёме display_name пуст. if got.DisplayName != "" { t.Errorf("display_name при приёме = %q, want пусто", got.DisplayName) } // download.Context = пользовательский текст + синтез из полей magnet // (dn=Dune). Текст пользователя идёт первым. if !strings.HasPrefix(got.Context, "Дюна 2") || !strings.Contains(got.Context, "Dune") { t.Errorf("сохранённый контекст = %q", got.Context) } if len(fs.hashes) != 1 || len(fs.hashes[0]) != 1 || fs.hashes[0][0] != sampleInfohash { t.Errorf("хеши задачи: %v", fs.hashes) } } // Голый magnet без текста: download.Context синтезируется из полей ссылки // (dn-имя + размер), приём проходит штатно. func TestIngestMagnetOnlySynthesizesContext(t *testing.T) { const raw = "magnet:?xt=urn:btih:541ADCFF3B6DD5DBA7088EA83317D9D6FAC331D6" + "&dn=Dune.Part.Two.2024.2160p&xl=2200000000" fs := &fakeStore{} if _, err := newService(fs).Ingest(context.Background(), Request{Source: raw}); err != nil { t.Fatalf("Ingest: %v", err) } if len(fs.created) != 1 { t.Fatalf("создано задач: %d, want 1", len(fs.created)) } ctx := fs.created[0].Context if !strings.Contains(ctx, "Dune.Part.Two.2024.2160p") || !strings.Contains(ctx, "Размер:") { t.Errorf("контекст не синтезирован из magnet: %q", ctx) } } // Заглушка-dn (rutracker-topic-*) как строка-название в контекст не попадает, // но домен трекера — попадает (сигнал для recognition). func TestIngestSynthDropsStubName(t *testing.T) { const raw = "magnet:?xt=urn:btih:541ADCFF3B6DD5DBA7088EA83317D9D6FAC331D6" + "&dn=rutracker-topic-6514485&tr=http%3A%2F%2Fbt.t-ru.org%2Fann%3Fmagnet" fs := &fakeStore{} if _, err := newService(fs).Ingest(context.Background(), Request{Source: raw}); err != nil { t.Fatalf("Ingest: %v", err) } got := fs.created[0].Context if !strings.Contains(got, "t-ru.org") { t.Errorf("download.Context не обогащён доменом трекера: %q", got) } if strings.Contains(got, "rutracker-topic") { t.Errorf("заглушка-dn просочилась в контекст как имя: %q", got) } } // Реальная рутрекер-ссылка без текста: download.Context = релиз-заголовок из // dn (раскодирован) + домен трекера. func TestIngestRealRutrackerMagnetOnly(t *testing.T) { const raw = "magnet:?xt=urn:btih:BACA24E18C7382A9E9A44132C8D7DB86C4D319C2" + "&tr=http%3A%2F%2Fbt4.t-ru.org%2Fann%3Fmagnet" + "&dn=%D0%91%D1%83%D1%85%D1%82%D0%B0%20%D0%B2%D0%B4%D0%BE%D0%B2%20%2F%20Widow's%20Bay%20%2F%20%D0%A1%D0%B5%D0%B7%D0%BE%D0%BD%3A%201%20%5B2026%2C%20%D0%A1%D0%A8%D0%90%2C%20WEB-DL%201080p%5D" fs := &fakeStore{} if _, err := newService(fs).Ingest(context.Background(), Request{Source: raw}); err != nil { t.Fatalf("Ingest: %v", err) } ctx := fs.created[0].Context if !strings.Contains(ctx, "Widow's Bay") || !strings.Contains(ctx, "Трекер: t-ru.org") { t.Errorf("download.Context не обогащён: %q", ctx) } } func TestIngestIdempotent(t *testing.T) { existing := &store.Download{ID: "01hzzzexisting000000000000", State: store.StateCatched} fs := &fakeStore{active: existing} res, err := newService(fs).Ingest(context.Background(), Request{Source: sampleMagnet}) if err != nil { t.Fatalf("Ingest: %v", err) } if !res.Deduplicated || res.DownloadID != existing.ID { t.Errorf("ожидалось присоединение к существующей задаче: %+v", res) } if len(fs.created) != 0 { t.Error("не должно создаваться новой задачи") } } // Повторный приём привязывается к удерживающей источник desync-записи // (target_missing/orphaned) вместо создания близнеца: возвращается существующая // «спящей» (её состояние не меняется, к qBittorrent не ходим), Deduplicated=true. func TestIngestAttachesToDesyncRecord(t *testing.T) { for _, s := range []store.State{store.StateTargetMissing, store.StateOrphaned} { t.Run(string(s), func(t *testing.T) { existing := &store.Download{ID: "01hzzzexisting000000000000", State: s} fs := &fakeStore{active: existing} res, err := newService(fs).Ingest(context.Background(), Request{Source: sampleMagnet}) if err != nil { t.Fatalf("Ingest: %v", err) } if !res.Deduplicated || res.DownloadID != existing.ID { t.Errorf("ожидалось присоединение к desync-записи: %+v", res) } if res.State != s { t.Errorf("состояние существующей записи должно вернуться как есть (%s), got %s", s, res.State) } if len(fs.created) != 0 { t.Error("не должно создаваться новой задачи (близнеца)") } }) } } // Быстрый дедуп-путь доносит существующей задаче недостающие хеши // гибридного magnet (иначе последующий приём по второму хешу создал бы // вторую активную задачу). func TestIngestDedupTopsUpHashes(t *testing.T) { const v2 = "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef" existing := &store.Download{ ID: "01hzzzexisting000000000000", State: store.StateCatched, Infohashes: []store.Infohash{{DownloadID: "01hzzzexisting000000000000", Infohash: sampleInfohash, Kind: store.HashV1}}, } fs := &fakeStore{active: existing} res, err := newService(fs).Ingest(context.Background(), Request{Source: sampleMagnet + "&xt=urn:btmh:1220" + v2}) if err != nil { t.Fatalf("Ingest: %v", err) } if !res.Deduplicated { t.Fatalf("ожидался дедуп: %+v", res) } if len(fs.toppedUp) != 2 { t.Errorf("хеши не донесены существующей задаче: %v", fs.toppedUp) } } // F8: пред-рид FindReingestBlockingByInfohash увидел активную задачу, но к моменту // создания она отменена (гонка с cancel). Активную запись пред-рид НЕ // короткозамыкает — авторитетное дедуп-решение принимает CreateDownloadIfNoActive // под BEGIN IMMEDIATE: активной больше нет → заводим свежую задачу, а не // возвращаем stale Deduplicated «уже в работе». func TestIngestActivePreReadNotShortCircuited(t *testing.T) { stale := &store.Download{ID: "01hzzzstale00000000000000000", State: store.StateCatched} fs := &raceStore{blocking: stale} // пред-рид видит активную; create-гард — уже нет res, err := newService(fs).Ingest(context.Background(), Request{Source: sampleMagnet}) if err != nil { t.Fatalf("Ingest: %v", err) } if res.Deduplicated { t.Errorf("активный пред-рид не должен коротко замыкать дедуп: %+v", res) } if res.DownloadID == stale.ID || res.State != store.StateCatched { t.Errorf("ожидалась свежая задача, а не stale: %+v", res) } if len(fs.created) != 1 { t.Errorf("должна быть создана новая задача, created=%d", len(fs.created)) } } // F7: oversized `.torrent` — доменная ошибка размера класса ErrTorrentTooLarge // (транспорт транслирует в 400, а не 500). Задача не заводится. func TestIngestRejectsOversizedTorrent(t *testing.T) { fs := &fakeStore{} big := make([]byte, MaxTorrentSize+1) _, err := newService(fs).Ingest(context.Background(), Request{TorrentData: big}) if !errors.Is(err, ErrTorrentTooLarge) { t.Fatalf("err = %v, want ErrTorrentTooLarge", err) } if len(fs.created) != 0 { t.Error("не должно быть записи задачи") } } // F10: контекст из веб-формы может быть огромным (multipart-бюджет на всё тело) — // Ingest режет его до MaxContextSize по границе руны (без U+FFFD) и метит маркером. func TestIngestCapsContext(t *testing.T) { // «Ё» — 2 байта; ASCII-префикс сдвигает границу MaxContextSize на нечётный // байт, чтобы обрезка s[:MaxContextSize] пришлась ВНУТРЬ двухбайтовой руны — // тогда trimToRune реально срабатывает (иначе граница попадёт между рунами). huge := "x" + strings.Repeat("Ё", MaxContextSize) fs := &fakeStore{} if _, err := newService(fs).Ingest(context.Background(), Request{Source: sampleMagnet, Context: huge}); err != nil { t.Fatalf("Ingest: %v", err) } got := fs.created[0].Context if len(got) > MaxContextSize+len(contextTruncMarker) { t.Errorf("контекст не ограничен: %d байт", len(got)) } if !strings.HasSuffix(got, contextTruncMarker) { t.Errorf("нет маркера усечения: …%q", got[max(0, len(got)-40):]) } if strings.ContainsRune(got, '�') { t.Error("обрезка порвала руну (U+FFFD)") } } func TestIngestRejectsNonMagnet(t *testing.T) { fs := &fakeStore{} if _, err := newService(fs).Ingest(context.Background(), Request{Source: "https://example.com/x.torrent"}); err == nil { t.Fatal("ожидалась ошибка для не-magnet источника") } if len(fs.created) != 0 { t.Error("не должно быть записи задачи") } }