// 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 { // Невалидный источник — норма (адресат не команда, а пользователь, и он // получит отказ на транспорте): DEBUG, чтобы не шуметь в аудите. s.log.Debug("ingest source rejected", "capability", capIngest, "error", err) 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 { // Инфраструктурный сбой (БД) — операция приёма не выполнена: ERROR. log.Error("ingest failed", "stage", "lookup-active", "error", err) 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 { // Инфраструктурный сбой (БД) — операция приёма не выполнена: ERROR. log.Error("ingest failed", "stage", "create-download", "error", err) 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, } }