Files
jellybit/internal/ingest/ingest.go
T
av d081ef1d30 ingest: закрыты мелочи приёма — вырожденное имя, контракт Result, корреляция add
- имя раздачи нормализуется на границе разбора: вырожденное `-`
  (metainfo.NoName) даёт пустое имя, пробельное схлопывается — сентинел больше
  не доходит ни до контекста распознавания, ни до source_ref, ни до подсказки
  вывода имени
- контракт «на любом пути ошибки приёма результат нулевой» объявлен в ingest и
  удерживается структурно; три транспорта перестали обещать идентификатор,
  которого нет, и коррелируют отказ по request_id
- scoped-логгер загрузки ставится до вызова внешнего сервиса в семи командах
  воркера — записи об отказе qBittorrent и метабаз получили download_id
  и infohash; граница разбора bencode записана в docs/research
2026-08-06 18:20:12 +03:00

306 lines
18 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"
"errors"
"fmt"
"log/slog"
"strings"
"unicode/utf8"
"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 — итог приёма. При ненулевой ошибке Ingest возвращает НУЛЕВОЙ Result:
// идентификатор загрузки, хеши, состояние и признак дедупликации не
// публикуются (см. Ingest).
type Result struct {
DownloadID string
Infohashes []string // все хеши источника (гибридный magnet: v1 и v2, v1 первым)
State store.State
Deduplicated bool // присоединились к уже активной задаче, нового добавления не было
}
// Ingest быстро принимает источник: извлекает infohash, синтезирует контекст из
// полей ссылки, дедуплицирует по активной задаче, иначе сохраняет загрузку в
// `catched` и сразу возвращает результат. Добавление в qBittorrent и вывод
// имени выполняет worker (см. download-tracking).
//
// Контракт: на ЛЮБОМ пути ошибки возвращается нулевой Result. Приём не создаёт
// наблюдаемых последствий раньше, чем способен вернуть успех, а всё, что может
// отказать после заведения загрузки, делает worker. Транспорты на это
// опираются и не обещают идентификатора, которого нет: HTTP коррелирует отказ
// по request_id, Telegram — ключа не даёт (см. ingest-спеку, требование
// «Результат приёма при ошибке пуст»).
//
// Гарантия структурная — обнуление в одном defer, а не аккуратность каждой
// ветки возврата: перечень веток растёт, и именно расхождение перечня с
// комментариями транспортов породило исходный дефект.
func (s *Service) Ingest(ctx context.Context, req Request) (res Result, err error) {
defer func() {
if err != nil {
res = Result{}
}
}()
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) блокируют повторный приём. Здесь короткозамыкаем
// ТОЛЬКО desync-запись (терминальную): её CreateDownloadIfNoActive не увидит
// (тот проверяет лишь активных), а состояния она не меняет. Активную же НЕ
// короткозамыкаем — этот чек без транзакции, и в гонке с параллельным cancel
// вернул бы stale «уже в работе» при пустом активном множестве. Авторитетное
// дедуп-решение по активной примет CreateDownloadIfNoActive под BEGIN IMMEDIATE.
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 && existing.State.IsTerminal() {
// FindReingestBlockingByInfohash отдаёт терминальную запись только из
// удерживающих desync-состояний (target_missing/orphaned) — присоединяемся.
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: capContext(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
// ErrTorrentTooLarge — принятый `.torrent` превышает MaxTorrentSize. Это промах
// ввода пользователя (норма, не сбой сервера), поэтому транспорт транслирует его
// в 400, а не 500 (веб MaxBytesReader пропускает файлы чуть больше лимита —
// отсекает уже приём). Проверяется через errors.Is.
var ErrTorrentTooLarge = errors.New("torrent too large")
// 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: %w", len(req.TorrentData), MaxTorrentSize, ErrTorrentTooLarge)
}
info, err := torrent.Parse(req.TorrentData)
if err != nil {
return parsedSource{}, fmt.Errorf("ingest: parse torrent: %w", err)
}
// SourceRef — человекочитаемый референс (имя раздачи), НЕ адрес
// добавления: torrent добавляется байтами (см. worker), не по SourceRef.
// Фолбек на имя файла, если у раздачи нет содержательного имени:
// вырожденное значение отбросил разборщик, здесь остаётся пустота.
ref := info.DisplayName
if 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
}
// MaxContextSize — предел размера контекста распознавания (пользовательский текст
// + синтез из полей источника). Кап здесь, на единственном месте слияния,
// покрывает все транспорты: REST ограничен телом (64 KiB), Telegram — лимитом
// подписи, но веб-форма (multipart-бюджет на всё тело) иначе пропустила бы
// мегабайты в поле context → в БД, рендер карточки и LLM-промпты.
const MaxContextSize = 16 << 10 // 16 KiB
// contextTruncMarker дописывается к усечённому контексту как явный маркер.
const contextTruncMarker = "\n…[контекст усечён]"
// capContext ограничивает контекст MaxContextSize байтами, обрезая по границе
// руны (кириллица — 2 байта/руна; обрезка посреди руны дала бы U+FFFD) и добавляя
// маркер усечения. Пустой/короткий контекст возвращается как есть.
func capContext(s string) string {
if len(s) <= MaxContextSize {
return s
}
return trimToRune(s[:MaxContextSize]) + contextTruncMarker
}
// trimToRune отбрасывает незавершённую многобайтовую руну на конце строки
// (результат обрезки по фиксированному числу байт), не трогая корректный хвост.
func trimToRune(s string) string {
for len(s) > 0 {
if r, size := utf8.DecodeLastRuneInString(s); r != utf8.RuneError || size > 1 {
break
}
s = s[:len(s)-1]
}
return s
}
// 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,
}
}