Синтезируем контекст из полей самой magnet-ссылки (dn, xl, tr/xs, kt) без сети и дополняем им текст от транспорта: пользовательский текст первым, при пустом — синтез единственный. Обогащённый контекст идёт только в download.Context (его читают recognition и веб-UI); вход namer и отображаемое имя не меняются — строки-факты (Размер:/Трекер:) в display_name не текут. - magnet.Info: поля ExactLength/Sources/Keywords + Info.Context() (синтез, отсев dn-заглушек *-topic-<id>, домен трекера, человекочитаемый размер) - Parse устойчив к ссылке в процент-кодировке (разовый QueryUnescape) - ingest: mergeContext → download.Context, namer на сыром req.Context - Влита дельта capability ingest в openspec/specs, change заархивирован Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
226 lines
11 KiB
Go
226 lines
11 KiB
Go
// 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 ещё нет». Имя от строки БД не зависит.
|
||
// Namer получает СЫРОЙ req.Context (+ dn-hint), не обогащённый: строки-факты
|
||
// синтеза (Размер:/Трекер:) не должны становиться отображаемым именем.
|
||
var rename string
|
||
if s.namer != nil {
|
||
rename = s.namer.DeriveName(ctx, req.Context, info.DisplayName)
|
||
}
|
||
|
||
// Контекст распознавания дополняем фактами из полей самой magnet-ссылки
|
||
// (dn/xl/tr/xs/kt) — без сети. Пользовательский текст идёт первым. Результат
|
||
// уходит только в download.Context (его читают recognition и веб-UI).
|
||
d := &store.Download{
|
||
SourceType: store.SourceMagnet,
|
||
SourceRef: source,
|
||
DisplayName: rename, // то же имя, что уходит в qBittorrent (rename); заголовок в веб-UI
|
||
Context: mergeContext(req.Context, info.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
|
||
}
|
||
|
||
// 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, 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,
|
||
}
|
||
}
|