Приём (Ingest) стал быстрым: синхронно только парс magnet, синтез контекста из полей ссылки, атомарный дедуп и запись загрузки в новое состояние `catched` — ответ клиенту сразу. Медленный вывод имени (LLM) и добавление в qBittorrent вынесены в асинхронный шаг машины состояний, который двигает worker. - store: состояние `catched` (нетерминальное, активная группа); атомарный переход PromoteCatched (catched → downloading + display_name) с гардом state='catched' (ре-валидация после сетевых вызовов вне блокировки) - ingest: убраны namer/qbt из пути приёма; пишем `catched`, отвечаем сразу - worker.processCatched: вне w.mu выводит имя и qbt.Add, под w.mu — короткий переход; сбой add оставляет catched (ретрай тиком); предохранитель catch_timeout → failed(qbit_add)+notify; catched исключён из проверок пропажи - config: worker.catch_timeout (дефолт 10m) - веб-UI: бейдж catched, активная группа, самозавершающийся htmx-поллинг карточки/страницы до перехода в downloading; Telegram-текст без сырого catched - OpenSpec: дельты ingest/download-tracking/web-ui влиты в спеки, change заархивирован Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
156 lines
7.7 KiB
Go
156 lines
7.7 KiB
Go
// 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"
|
||
)
|
||
|
||
// capIngest — стадия приёма для поля capability в логах.
|
||
const capIngest = "ingest"
|
||
|
||
// Store — нужная ingest часть хранилища.
|
||
type Store interface {
|
||
// FindActiveByInfohash — быстрый читающий дедуп-чек; авторитетная проверка —
|
||
// внутри 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
|
||
}
|
||
|
||
// 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 — входной запрос приёма.
|
||
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, синтезирует контекст из
|
||
// полей ссылки, дедуплицирует по активной задаче, иначе сохраняет загрузку в
|
||
// `catched` и сразу возвращает результат. Добавление в qBittorrent и вывод
|
||
// имени выполняет worker (см. download-tracking).
|
||
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.
|
||
log := s.log.With("capability", capIngest, "infohash", info.Infohash)
|
||
ctx = logctx.With(ctx, log)
|
||
|
||
// Быстрый дедуп-чек; авторитетная (атомарная) проверка — внутри
|
||
// 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
|
||
}
|
||
|
||
// Контекст распознавания дополняем фактами из полей самой magnet-ссылки
|
||
// (dn/xl/tr/xs/kt) — без сети. Пользовательский текст идёт первым.
|
||
// DisplayName пуст: имя выведет worker на шаге добавления (rename действует
|
||
// только при добавлении, а тут медленный LLM в пути ответа недопустим).
|
||
d := &store.Download{
|
||
SourceType: store.SourceMagnet,
|
||
SourceRef: source,
|
||
Context: mergeContext(req.Context, info.Context()),
|
||
State: store.StateCatched,
|
||
}
|
||
// Все хеши из 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
|
||
}
|
||
|
||
log.Info("download catched", "download_id", d.ID)
|
||
return Result{
|
||
DownloadID: d.ID,
|
||
Infohashes: info.Infohashes,
|
||
State: store.StateCatched,
|
||
}, 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 — итог дедупа на быстром чеке: присоединились к уже активной
|
||
// задаче и доносим ей недостающие хеши источника (гибридный 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,
|
||
}
|
||
}
|