Files
jellybit/internal/ingest/ingest.go
T
avandClaude Opus 4.8 1369a9cabe Приём: дедуп по target_missing/orphaned + стоп-кран «Закрыть»
Два дубля-близнеца на один инфохэш рождались, когда повторный приём
попадал на запись в target_missing: дедуп искал только активную задачу,
а target_missing терминален → заводилась новая загрузка, воркер усыновлял
уже присутствующий торрент и раскладывал его.

- Приём: критерий дедупа расширен до «блокирующей повторный приём» =
  активные ∪ {target_missing, orphaned}. Повторный приём такого инфохэша
  привязывается к существующей записи (спящей, без обращения к qBittorrent),
  а не плодит близнеца. Прочие терминальные (done/cancelled/failed/reverted/
  deleted) повторный приём не блокируют — осознанная свежая попытка. Новый
  read-метод FindReingestBlockingByInfohash (приоритет активной над desync);
  общий active-гард не тронут.
- Команда «Закрыть» (Dismiss) — универсальный стоп-кран из любого состояния,
  кроме deleted → cancelled (error_code=user_dismiss). Только меняет статус:
  файлы (в т.ч. хардлинки done/orphaned) и раздачу qBittorrent не трогает,
  в отличие от «Удалить». Веб — danger-зона внизу страницы; Telegram —
  кнопка с подтверждением; из cancelled — идемпотентный no-op.
- Транспорты при дедупе на desync-запись сообщают адресно (target_missing —
  привязать заново/закрыть; orphaned — закрыть и добавить заново); веб при
  дедупе ведёт на страницу существующей записи.

Спеки: ingest (дедуп), state-reconciliation (стоп-кран); граф переходов
допополнен рёбрами <терминал>→cancelled. OpenSpec change
dedup-target-missing-and-dismiss заархивирован.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-10 20:15:37 +03:00

243 lines
14 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 {
// FindReingestBlockingByInfohash — быстрый читающий дедуп-чек: активная задача
// ЛИБО удерживающая источник desync-запись (target_missing/orphaned). Активный
// инвариант «≤1 активной» авторитетно держит CreateDownloadIfNoActive; desync —
// устойчивый пред-рид, коротко замыкающий приём на возврат существующей записи.
FindReingestBlockingByInfohash(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)
// Быстрый дедуп-чек по ЛЮБОМУ из хешей источника (гибридный несёт и v1, и v2):
// активная задача ЛИБО удерживающая источник desync-запись
// (target_missing/orphaned) блокируют повторный приём. Для активной
// авторитетная (атомарная) проверка — внутри CreateDownloadIfNoActive ниже;
// desync-ветка сюда и завершается (в active-гард desync не заводим, чтобы не
// размыть инвариант «≤1 активной»).
if existing, err := s.store.FindReingestBlockingByInfohash(ctx, src.infohashes...); err != nil {
// Инфраструктурный сбой (БД) — операция приёма не выполнена: ERROR.
log.Error("ingest failed", "stage", "lookup-blocking", "error", err)
return Result{}, fmt.Errorf("ingest: lookup blocking: %w", err)
} else if existing != nil {
log.Info("download attached", "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 — итог дедупа на быстром чеке: присоединились к блокирующей записи
// (активной ЛИБО удерживающей источник desync — target_missing/orphaned) и
// доносим ей недостающие хеши источника (гибридный magnet мог принести хеш,
// которого задача ещё не знает; guarded-путь через CreateDownloadIfNoActive сюда
// не доходит). Для desync-записи состояние не меняем (возвращаем «спящей» — relink
// или закрытие делает пользователь). Донос — 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,
}
}