Пять независимых bugfix'ов из ревью приёма (docs/backlog/review-f7-f10-ingest-ui-fixes.md): - F7: oversized .torrent через веб отдавал 500. Введён sentinel ingest.ErrTorrentTooLarge, classifyErr транслирует его в 400. - F8: гонка fast-path attach с cancel. Пред-рид FindReingestBlockingByInfohash больше не короткозамыкает активную запись — авторитетное дедуп-решение принимает CreateDownloadIfNoActive под BEGIN IMMEDIATE; короткозамыкание оставлено только для терминальных desync-записей (target_missing/orphaned). F6-апгрейд сохранён. - F9: magnet — регистронезависимый URN-префикс xt (RFC 2141); tgbot.ParseMessage срезает хвостовую пунктуацию, приклеенную жадным matchем. - F10: cap контекста до 16 KiB в ingest.Ingest (единственное место слияния — покрывает все транспорты), рунобезопасная обрезка + маркер. - N2: httpapi.shorten режет по рунам, не байтам — кириллица не рвётся в U+FFFD. Добавлены юнит-тесты на каждое исправленное поведение. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
287 lines
17 KiB
Go
287 lines
17 KiB
Go
// 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 — итог приёма.
|
||
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) блокируют повторный приём. Здесь короткозамыкаем
|
||
// ТОЛЬКО 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.
|
||
// Фолбек на имя файла, если у раздачи нет содержательного имени
|
||
// (пустое или 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
|
||
}
|
||
|
||
// 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,
|
||
}
|
||
}
|