// Package ingest — use-case приёма загрузки, общий для всех транспортов // (HTTP, Telegram, CLI). Принимает источник + контекст, отдаёт источник в // qBittorrent и заводит/находит задачу в БД. 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/qbt" "git.vakhrushev.me/av/jellybit/internal/store" ) // capIngest — стадия приёма для поля capability в логах. const capIngest = "ingest" // errCodeQbitAdd — error_code задачи, упавшей на добавлении источника в // qBittorrent (раздачи в qBittorrent нет, восстановлению не подлежит). const errCodeQbitAdd = "qbit_add" // Store — нужная ingest часть хранилища. type Store interface { // FindActiveByInfohash — быстрый читающий дедуп-чек (до вызова LLM-namer); // авторитетная проверка — внутри CreateDownloadIfNoActive. FindActiveByInfohash(ctx context.Context, hashes ...string) (*store.Download, error) // CreateDownloadIfNoActive атомарно проверяет инвариант «одна активная // загрузка на infohash» и заводит задачу; вернувшаяся existing ≠ nil — // дедуп на активную задачу (недостающие хеши вызова метод доносит сам). CreateDownloadIfNoActive(ctx context.Context, d *store.Download, hashes []string) (*store.Download, error) // AddInfohashes доносит задаче недостающие хеши (guarded). Нужен на // быстром дедуп-пути, который не доходит до CreateDownloadIfNoActive. AddInfohashes(ctx context.Context, downloadID string, hashes []string) error SetDownloadState(ctx context.Context, id string, state store.State, errCode, errMsg string) error } // QBittorrent — нужная ingest часть клиента qBittorrent. type QBittorrent interface { Add(ctx context.Context, ar qbt.AddRequest) error } // Namer выводит человекочитаемое отображаемое имя торрента из контекста. // Пустой результат → имя в qBittorrent не задаём. nil → шаг пропускается. type Namer interface { DeriveName(ctx context.Context, contextText, hint string) string } // Config — параметры добавления в qBittorrent. type Config struct { Category string SavePath string } // Service — реализация приёма. type Service struct { store Store qbt QBittorrent namer Namer cfg Config log *slog.Logger // notifyFailed — опц. пинг автору о падении приёма (добавление в qBittorrent // не удалось). Closure, а не worker.Notifier: приёмное падение в qBit не // попадает в поллинг-цикл worker (раздачи нет), поэтому уведомляет ingest // сам; closure избавляет ядро приёма от зависимости на пакет worker. notifyFailed func(downloadID string) } // New собирает сервис приёма. namer опционален (nil → отображаемое имя не // выводится; qBittorrent оставит своё). func New(st Store, qb QBittorrent, namer Namer, cfg Config, log *slog.Logger) *Service { return &Service{store: st, qbt: qb, namer: namer, cfg: cfg, log: log} } // SetFailureNotifier подключает пинг о падении приёма (до начала работы). func (s *Service) SetFailureNotifier(fn func(downloadID string)) { s.notifyFailed = fn } // Request — входной запрос приёма. type Request struct { Source string // пока — magnet-ссылка Context string // подсказка для распознавания (опц.) } // Result — итог приёма. type Result struct { DownloadID string Infohashes []string // все хеши источника (гибридный magnet: v1 и v2, v1 первым) State store.State Deduplicated bool // присоединились к уже активной задаче, нового добавления не было } // Ingest принимает источник: извлекает infohash, дедуплицирует по активной // задаче, иначе заводит задачу и отдаёт источник в qBittorrent. func (s *Service) Ingest(ctx context.Context, req Request) (Result, error) { source := strings.TrimSpace(req.Source) info, err := magnet.Parse(source) if err != nil { // Ф1: поддержан только magnet. .torrent/url — следующий заход. return Result{}, fmt.Errorf("ingest: parse source: %w", err) } // Scoped-логгер стадии приёма: download_id допишется после CreateDownload. // Кладём в ctx, чтобы внешние клиенты (qBittorrent, LLM-namer) дописывали // ключи корреляции к своим ext.*-записям сами. log := s.log.With("capability", capIngest, "infohash", info.Infohash) ctx = logctx.With(ctx, log) // Быстрый дедуп-чек до дорогого LLM-namer; авторитетная (атомарная) // проверка — внутри CreateDownloadIfNoActive ниже. Дедуп — по ЛЮБОМУ из // хешей источника: гибридный magnet несёт и v1, и v2. if existing, err := s.store.FindActiveByInfohash(ctx, info.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, info, existing), nil } // Отображаемое имя для списка qBit — best-effort: не валит приём. // Выводится синхронно (param rename действует только при добавлении) и // ДО CreateDownload, чтобы возможный медленный вызов LLM не расширял окно // «строка в БД есть, в qBittorrent ещё нет». Имя от строки БД не зависит. var rename string if s.namer != nil { rename = s.namer.DeriveName(ctx, req.Context, info.DisplayName) } d := &store.Download{ SourceType: store.SourceMagnet, SourceRef: source, DisplayName: rename, // то же имя, что уходит в qBittorrent (rename); заголовок в веб-UI Context: req.Context, State: store.StateDownloading, } // Все хеши из magnet (гибридный несёт v1 и v2); kind store выведет по длине. existing, err := s.store.CreateDownloadIfNoActive(ctx, d, info.Infohashes) if err != nil { return Result{}, fmt.Errorf("ingest: create download: %w", err) } if existing != nil { // Гонка с параллельным приёмом/discover: активная задача появилась // после быстрого чека — присоединяемся к ней (хеши донёс сам // CreateDownloadIfNoActive). log.Info("download attached to active", "download_id", existing.ID, "state", existing.State) return Result{ DownloadID: existing.ID, Infohashes: info.Infohashes, State: existing.State, Deduplicated: true, }, nil } id := d.ID log = log.With("download_id", id) ctx = logctx.With(ctx, log) addErr := s.qbt.Add(ctx, qbt.AddRequest{ URLs: []string{source}, Category: s.cfg.Category, SavePath: s.cfg.SavePath, Rename: rename, }) if addErr != nil { // Граница доменной операции приёма: логируем исход один раз (ERROR). // Поведение самого вызова qBittorrent уже залогировал клиент (ext.*) — // это разные факты, не дубль. log.Error("download accept failed", "error", addErr) // Задача уже в БД — помечаем failed, чтобы worker её не подхватил. if setErr := s.store.SetDownloadState(ctx, id, store.StateFailed, errCodeQbitAdd, addErr.Error()); setErr != nil { log.Error("mark download failed after qbit error failed", "error", setErr) } else if s.notifyFailed != nil { // Это падение минует worker.transition (раздачи в qBit нет) — уведомляем // сами, чтобы приёмные провалы тоже доходили до автора. go s.notifyFailed(id) } return Result{DownloadID: id, Infohashes: info.Infohashes, State: store.StateFailed}, fmt.Errorf("ingest: add to qbittorrent: %w", addErr) } log.Info("download accepted", "category", s.cfg.Category) return Result{ DownloadID: id, Infohashes: info.Infohashes, State: store.StateDownloading, }, nil } // attached — итог дедупа на быстром чеке: присоединились к уже активной // задаче и доносим ей недостающие хеши источника (гибридный magnet мог // принести хеш, которого задача ещё не знает; guarded-путь через // CreateDownloadIfNoActive сюда не доходит). Донос — best-effort: конфликт // хеша с другой активной задачей логируется, приём не валится. func (s *Service) attached(ctx context.Context, info magnet.Info, existing *store.Download) Result { if len(info.Infohashes) > len(existing.Infohashes) { if err := s.store.AddInfohashes(ctx, existing.ID, info.Infohashes); err != nil { logctx.FromOr(ctx, s.log).Warn("ingest top-up infohashes failed", "error", err) } } return Result{ DownloadID: existing.ID, Infohashes: info.Infohashes, State: existing.State, Deduplicated: true, } }