Files
jellybit/internal/ingest/ingest.go
T
avandClaude Opus 4.8 4475fbd548 Приём: гард дедуп-дозаписи хешей (F1) и апгрейд catched-magnet до torrent (F6)
Два дефекта дедуп-веток приёма (ревью Fable 2026-07-08), оба про инвариант
«≤1 активная загрузка на infohash» и сохранность источника.

F1: дедуп-ветка CreateDownloadIfNoActive дописывала все хеши входящего
источника в найденную активную задачу без пер-хеш гарда владения (в отличие
от AddInfohashes). Гибрид {v1,v2}, дедупнувшись на задачу B (владелец v2),
крал v1 у активной A → две активные владели v1. Теперь дозапись под тем же
гардом: хеш, которым владеет другая активная задача, не дописывается.

F6: при дедупе .torrent-байт на пойманную magnet-задачу (catched) байты
выбрасывались, source_type оставался magnet → worker добавлял по magnet-URL →
вечный metaDL → failed (magnet закрытого трекера без DHT метаданные не
докачает). Новый guarded-метод UpgradeCatchedMagnetToTorrent атомарно
сохраняет байты и меняет source_type magnet→torrent, но только пока задача в
catched (worker источник ещё не отдал). Ingest зовёт апгрейд на обоих
дедуп-путях. Это целевое исключение из правила спеки «при дедупе байты не
сохраняем» — оформлено MODIFIED-дельтой ingest.

Схема БД не меняется (download_torrent и source_type уже есть).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-08 17:18:28 +03:00

229 lines
12 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// Package ingest — use-case быстрого приёма загрузки, общий для всех
// транспортов (HTTP, Telegram, CLI). Синхронно только парсит источник,
// синтезирует контекст из полей ссылки, дедуплицирует и сохраняет загрузку в
// состоянии `catched`, сразу возвращая ответ. Вывод отображаемого имени
// (медленный LLM) и добавление в qBittorrent — отдельный асинхронный шаг
// worker'а (см. download-tracking).
package ingest
import (
"context"
"fmt"
"log/slog"
"strings"
"git.vakhrushev.me/av/jellybit/internal/logctx"
"git.vakhrushev.me/av/jellybit/internal/magnet"
"git.vakhrushev.me/av/jellybit/internal/store"
"git.vakhrushev.me/av/jellybit/internal/torrent"
)
// capIngest — стадия приёма для поля capability в логах.
const capIngest = "ingest"
// Store — нужная ingest часть хранилища.
type Store interface {
// FindActiveByInfohash — быстрый читающий дедуп-чек; авторитетная проверка —
// внутри CreateDownloadIfNoActive.
FindActiveByInfohash(ctx context.Context, hashes ...string) (*store.Download, error)
// CreateDownloadIfNoActive атомарно проверяет инвариант «одна активная
// загрузка на infohash» и заводит задачу; вернувшаяся existing ≠ nil —
// дедуп на активную задачу (недостающие хеши вызова метод доносит сам).
// torrentBlob (для source_type=torrent) пишется в той же транзакции только
// на ветке создания; при дедупе не пишется. Для magnet — nil.
CreateDownloadIfNoActive(ctx context.Context, d *store.Download, hashes []string, torrentBlob []byte) (*store.Download, error)
// AddInfohashes доносит задаче недостающие хеши (guarded). Нужен на
// быстром дедуп-пути, который не доходит до CreateDownloadIfNoActive.
AddInfohashes(ctx context.Context, downloadID string, hashes []string) error
// UpgradeCatchedMagnetToTorrent при дедупе входящих байтов `.torrent` на
// пойманную (`catched`) magnet-задачу сохраняет байты и меняет source_type
// на torrent (атомарно). No-op, если задача уже добавлена/не magnet. Чинит
// magnet закрытого трекера, который иначе застрянет в metaDL.
UpgradeCatchedMagnetToTorrent(ctx context.Context, downloadID string, torrentBlob []byte) (bool, error)
}
// Service — реализация быстрого приёма.
type Service struct {
store Store
log *slog.Logger
}
// New собирает сервис приёма.
func New(st Store, log *slog.Logger) *Service {
return &Service{store: st, log: log}
}
// Request — входной запрос приёма. Источник задаётся либо строкой (magnet), либо
// байтами `.torrent` (TorrentData); при непустом TorrentData он в приоритете.
type Request struct {
Source string // magnet-ссылка (когда TorrentData пуст)
TorrentData []byte // байты .torrent-файла (опц.); при наличии — источник torrent
TorrentName string // имя загруженного файла (опц.); фолбек для source_ref
Context string // подсказка для распознавания (опц.)
}
// Result — итог приёма.
type Result struct {
DownloadID string
Infohashes []string // все хеши источника (гибридный magnet: v1 и v2, v1 первым)
State store.State
Deduplicated bool // присоединились к уже активной задаче, нового добавления не было
}
// Ingest быстро принимает источник: извлекает infohash, синтезирует контекст из
// полей ссылки, дедуплицирует по активной задаче, иначе сохраняет загрузку в
// `catched` и сразу возвращает результат. Добавление в qBittorrent и вывод
// имени выполняет worker (см. download-tracking).
func (s *Service) Ingest(ctx context.Context, req Request) (Result, error) {
src, err := s.parse(req)
if err != nil {
return Result{}, err
}
// Scoped-логгер стадии приёма: download_id допишется после CreateDownload.
log := s.log.With("capability", capIngest, "infohash", src.infohashes[0])
ctx = logctx.With(ctx, log)
// Быстрый дедуп-чек; авторитетная (атомарная) проверка — внутри
// CreateDownloadIfNoActive ниже. Дедуп — по ЛЮБОМУ из хешей источника:
// гибридный несёт и v1, и v2.
if existing, err := s.store.FindActiveByInfohash(ctx, src.infohashes...); err != nil {
return Result{}, fmt.Errorf("ingest: lookup active: %w", err)
} else if existing != nil {
log.Info("download attached to active", "download_id", existing.ID, "state", existing.State)
return s.attached(ctx, src, existing), nil
}
// Контекст распознавания дополняем фактами из полей источника (magnet-поля
// или дерево `.torrent`) — без сети. Пользовательский текст идёт первым.
// DisplayName пуст: имя выведет worker на шаге добавления (rename действует
// только при добавлении, а тут медленный LLM в пути ответа недопустим).
d := &store.Download{
SourceType: src.sourceType,
SourceRef: src.sourceRef,
Context: mergeContext(req.Context, src.synthContext),
State: store.StateCatched,
}
// Все хеши источника (гибрид несёт v1 и v2); kind store выведет по длине.
// torrentBlob непуст только для source_type=torrent — пишется в той же
// транзакции лишь на ветке создания.
existing, err := s.store.CreateDownloadIfNoActive(ctx, d, src.infohashes, src.torrentBlob)
if err != nil {
return Result{}, fmt.Errorf("ingest: create download: %w", err)
}
if existing != nil {
// Гонка с параллельным приёмом/discover: активная задача появилась
// после быстрого чека — присоединяемся к ней (хеши донёс сам
// CreateDownloadIfNoActive; existing уже с подгруженными хешами).
log.Info("download attached to active", "download_id", existing.ID, "state", existing.State)
return s.attached(ctx, src, existing), nil
}
log.Info("download catched", "download_id", d.ID)
return Result{
DownloadID: d.ID,
Infohashes: src.infohashes,
State: store.StateCatched,
}, nil
}
// MaxTorrentSize — предел размера принимаемого `.torrent` (защита от разбухания
// БД и oversized-загрузок). Реальные торрент-файлы много меньше; крупные (много
// файлов → много piece-хешей) отсекаются здесь.
const MaxTorrentSize = 8 << 20 // 8 MiB
// parsedSource — нормализованный источник приёма (magnet или .torrent).
type parsedSource struct {
sourceType store.SourceType
sourceRef string // референс для человека/логов (magnet-URI или имя torrent)
infohashes []string // все хеши источника; v1 раньше v2
synthContext string // синтез контекста из полей источника (без сети)
torrentBlob []byte // байты .torrent (для source_type=torrent); иначе nil
}
// parse разбирает источник запроса: при непустом TorrentData — как `.torrent`
// (в приоритете), иначе — как magnet. Синтез контекста и хеши берутся из полей
// источника, без сети.
func (s *Service) parse(req Request) (parsedSource, error) {
if len(req.TorrentData) > 0 {
if len(req.TorrentData) > MaxTorrentSize {
return parsedSource{}, fmt.Errorf("ingest: torrent too large: %d > %d bytes", len(req.TorrentData), MaxTorrentSize)
}
info, err := torrent.Parse(req.TorrentData)
if err != nil {
return parsedSource{}, fmt.Errorf("ingest: parse torrent: %w", err)
}
// SourceRef — человекочитаемый референс (имя раздачи), НЕ адрес
// добавления: torrent добавляется байтами (см. worker), не по SourceRef.
// Фолбек на имя файла, если у раздачи нет содержательного имени
// (пустое или NoName-сентинел "-").
ref := strings.TrimSpace(info.DisplayName)
if ref == "" || ref == "-" {
ref = strings.TrimSpace(req.TorrentName)
}
return parsedSource{
sourceType: store.SourceTorrent,
sourceRef: ref,
infohashes: info.Infohashes,
synthContext: info.Context(),
torrentBlob: req.TorrentData,
}, nil
}
source := strings.TrimSpace(req.Source)
info, err := magnet.Parse(source)
if err != nil {
return parsedSource{}, fmt.Errorf("ingest: parse source: %w", err)
}
return parsedSource{
sourceType: store.SourceMagnet,
sourceRef: source,
infohashes: info.Infohashes,
synthContext: info.Context(),
}, nil
}
// mergeContext склеивает контекст от транспорта с синтезом из полей magnet:
// пользовательский текст идёт первым, затем факты из ссылки. Пустые части
// опускаются; при пустых обеих — пустая строка (пустой контекст допустим).
func mergeContext(userText, synth string) string {
parts := make([]string, 0, 2)
if s := strings.TrimSpace(userText); s != "" {
parts = append(parts, s)
}
if s := strings.TrimSpace(synth); s != "" {
parts = append(parts, s)
}
return strings.Join(parts, "\n")
}
// attached — итог дедупа на быстром чеке: присоединились к уже активной
// задаче и доносим ей недостающие хеши источника (гибридный magnet мог
// принести хеш, которого задача ещё не знает; guarded-путь через
// CreateDownloadIfNoActive сюда не доходит). Донос — best-effort: конфликт
// хеша с другой активной задачей логируется, приём не валится.
func (s *Service) attached(ctx context.Context, src parsedSource, existing *store.Download) Result {
if len(src.infohashes) > len(existing.Infohashes) {
if err := s.store.AddInfohashes(ctx, existing.ID, src.infohashes); err != nil {
logctx.FromOr(ctx, s.log).Warn("ingest top-up infohashes failed", "error", err)
}
}
// F6: входящее — байты `.torrent`, а активная задача поймана как magnet и
// ещё не отдана в qBittorrent (catched) → сохраняем байты и переключаем
// источник на torrent, чтобы воркер добавил раздачу файлом (magnet
// закрытого трекера иначе застрянет в metaDL). Best-effort: не валит приём.
if src.sourceType == store.SourceTorrent && len(src.torrentBlob) > 0 {
if upgraded, err := s.store.UpgradeCatchedMagnetToTorrent(ctx, existing.ID, src.torrentBlob); err != nil {
logctx.FromOr(ctx, s.log).Warn("ingest torrent upgrade failed", "error", err)
} else if upgraded {
logctx.FromOr(ctx, s.log).Info("catched magnet upgraded to torrent", "download_id", existing.ID)
}
}
return Result{
DownloadID: existing.ID,
Infohashes: src.infohashes,
State: existing.State,
Deduplicated: true,
}
}