Files
jellybit/internal/ingest/ingest.go
T
avandClaude Fable 5 37f2f6481a Идентичность на ULID: download_infohash, guarded-дедуп, миграция (ulid-identity)
Все сущности переехали с INTEGER AUTOINCREMENT на TEXT ULID (lowercase,
internal/ident — единая точка генерации и разбора; oklog/ulid). Инфохэши
загрузки — множество (download_infohash, v1/v2 гибридных торрентов): дедуп
и сопоставление в поллинге по любому из хешей, magnet-парсер отдаёт оба
хеша гибридной ссылки, усечённый v2-хеш v2-only раздач не хранится.

Инвариант «не более одной активной загрузки на infohash» вместо снятого
unique-индекса держат guarded-методы store в одной write-транзакции
(_txlock=immediate): CreateDownloadIfNoActive (приём/adopt, с доносом
недостающих хешей), ActivateIfNoOtherActive (retry/recovery/relink, отказ
до побочных эффектов), guarded AddInfohashes; SetDownloadState отклоняет
терминал→активное как механический бэкстоп.

Миграция 0006 — первая Go-миграция goose: пересоздание таблиц при
включённых FK, backfill ULID с timestamp из created_at (хронология id
сохранена), разнос infohash, удаление idempotency_key. BREAKING: формат id
в URL/логах/Telegram, REST-поля id (string) и infohashes (список).

Новая конвенция docs/conventions/database.md (без числовых PK), корреляция
в логах grep'ом по голому ULID, ER-схема обновлена. Спеки: новая capability
identity, MODIFIED в state-reconciliation; change заархивирован. Пройдены
ревью дизайна и кода (по 8 углов), все находки исправлены с
регрессионными тестами.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-02 21:25:00 +03:00

207 lines
10 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). Принимает источник + контекст, отдаёт источник в
// 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 ещё нет». Имя от строки БД не зависит.
var rename string
if s.namer != nil {
rename = s.namer.DeriveName(ctx, req.Context, info.DisplayName)
}
d := &store.Download{
SourceType: store.SourceMagnet,
SourceRef: source,
DisplayName: rename, // то же имя, что уходит в qBittorrent (rename); заголовок в веб-UI
Context: req.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
}
// 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,
}
}