Files
jellybit/internal/ingest/ingest.go
T
avandClaude Opus 4.8 b8017d65eb Приём/UI: пачка фиксов границ и парсинга (F7–F10, N2)
Пять независимых 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>
2026-07-17 20:35:14 +03:00

287 lines
17 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 — итог приёма.
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,
}
}