Принимаем .torrent как загруженные байты — через файл-пикер в веб-форме и Telegram-документ, наряду с magnet. Файл несёт полные метаданные: работает там, где magnet не резолвится (закрытые трекеры, без DHT), и даёт максимум контекста для распознавания без сети. - internal/torrent: парсер поверх anacrolix/torrent/metainfo — инфохэш(и) (v1 SHA1 исходных байтов info; v2 BEP52 при наличии) + Context() из имени, дерева файлов, размера, трекеров. Извлечение файлов панико-безопасно (недоверенный вход). - Персистентность байтов: таблица-спутник download_torrent (миграция 0009); пишется в транзакции создания загрузки, только на ветке создания (не при дедупе). Байты живут весь срок строки — нужны для повторного добавления при retry. - ingest: Request.TorrentData/TorrentName, диспетч парсера; source_ref — человекочитаемый референс (имя раздачи/файла), не адрес добавления. - worker: общий sourceAddParts ветвит по source_type в ОБОИХ add-путях — processCatched и Retry (torrent добавляется файлом, не magnet-хешем). - Транспорты: multipart-форма с файл-пикером (деградация без JS) и приём Telegram-документа (скачивание с редактированием токена из ошибок — секрет не в логи; обработка до ветки pending/текста). Разработка по OpenSpec (SDD): change torrent-file-ingest, два чекпоинта ревью (дизайн до кода, код до архива) сабагентами; дельты влиты в спеки, change архивирован. Ручная проверка на живом qBittorrent (7.3) — за деплоем. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
769 lines
36 KiB
Go
769 lines
36 KiB
Go
package store
|
||
|
||
import (
|
||
"context"
|
||
"database/sql"
|
||
"errors"
|
||
"fmt"
|
||
"slices"
|
||
"strings"
|
||
"time"
|
||
|
||
"github.com/jmoiron/sqlx"
|
||
|
||
"git.vakhrushev.me/av/jellybit/internal/ident"
|
||
)
|
||
|
||
// State — состояние загрузки в машине состояний (см. architecture.md).
|
||
// В Ф1 используется подмножество: downloading → completed, плюс stuck,
|
||
// failed, cancelled. Остальные состояния заведены под будущие фазы.
|
||
type State string
|
||
|
||
const (
|
||
StateCatched State = "catched" // поймано и сохранено; worker добавит в qBittorrent
|
||
StateDownloading State = "downloading"
|
||
StateCompleted State = "completed"
|
||
StateRecognizing State = "recognizing" // Ф2
|
||
StateReview State = "review" // Ф3
|
||
StateLinking State = "linking" // Ф3
|
||
StateDone State = "done" // Ф3
|
||
StateDeferred State = "deferred" // Ф3
|
||
StateStuck State = "stuck"
|
||
StateFailed State = "failed"
|
||
StateCancelled State = "cancelled"
|
||
StateReverted State = "reverted" // Ф3
|
||
|
||
// Состояния рассинхрона с реальностью (см. state-reconciliation).
|
||
StateTargetMissing State = "target_missing" // источник есть, цель удалена → relink
|
||
StateOrphaned State = "orphaned" // источник пропал, цель (последняя копия) есть
|
||
StateDeleted State = "deleted" // нет ни источника, ни цели
|
||
)
|
||
|
||
// terminalStates — единый список окончательно остановленных состояний:
|
||
// источник истины и для IsTerminal, и для выборки «активных» задач.
|
||
// Активность выводится ТОЛЬКО из state (отдельного ключа идемпотентности
|
||
// нет); любое новое терминальное состояние добавляется ТОЛЬКО сюда — иначе
|
||
// семантика «активности» разъедется.
|
||
//
|
||
// Состояния рассинхрона (target_missing/orphaned/deleted), а также failed —
|
||
// терминальны для активности, но не «мертвы»: дальше двигает либо человек
|
||
// (relink из target_missing, retry из failed), либо фоновая сверка (healing/
|
||
// прогрессия desync; авто-восстановление при оживлении источника, см.
|
||
// state-reconciliation) — через ActivateIfNoOtherActive, который атомарно
|
||
// проверяет инвариант «не более одной активной загрузки на infohash».
|
||
var terminalStates = []State{
|
||
StateDone, StateCancelled, StateFailed, StateReverted,
|
||
StateTargetMissing, StateOrphaned, StateDeleted,
|
||
}
|
||
|
||
// IsTerminal сообщает, завершена ли задача окончательно. Терминальная задача
|
||
// не «владеет» своими инфохэшами — тот же infohash можно завести заново
|
||
// новой задачей (см. architecture.md, «повторное добавление»). stuck
|
||
// терминальным не считается: задача восстановима (retry).
|
||
func (s State) IsTerminal() bool {
|
||
return slices.Contains(terminalStates, s)
|
||
}
|
||
|
||
// allowedTransitions — декларативный граф легальных переходов машины состояний
|
||
// (from → множество допустимых to). Единственный источник истины о легальности
|
||
// рёбер: покрывает все переходы, которые worker выполняет по всем capability
|
||
// (прямой путь download-tracking, state-reconciliation, review). Выведен построчно
|
||
// из кода воркера — таблица в change `state-transition-graph`/design.md.
|
||
//
|
||
// Правила:
|
||
// - Самопереход (from == to, идемпотентная переустановка того же состояния —
|
||
// напр. повторная запись error_msg, Defer на уже deferred) разрешён ВСЕГДА и
|
||
// здесь НЕ перечисляется (его добавляет invertTransitions).
|
||
// - Ребро может присутствовать здесь, но всё равно требовать revive-путь
|
||
// (ActivateIfNoOtherActive): гейт графа ортогонален гарду терминальности в
|
||
// setState — граф говорит «ребро есть», гард «но не мимо ActivateIfNoOtherActive».
|
||
// Так, failed → downloading объявлено, но обычным SetDownloadState отклоняется.
|
||
// - cancelled/deferred — легальная цель из КАЖДОГО не-терминального состояния
|
||
// (Cancel/Defer проверяют лишь IsTerminal); инвариант закреплён тестом, а не
|
||
// ручной аккуратностью.
|
||
//
|
||
// Правка воркера, вводящая новое ребро, ОБЯЗАНА отразить его здесь — иначе
|
||
// setState отклонит переход (0 строк UPDATE → ошибка).
|
||
var allowedTransitions = map[State][]State{
|
||
StateCatched: {StateDownloading, StateFailed, StateCancelled, StateDeferred},
|
||
StateDownloading: {StateCompleted, StateFailed, StateStuck, StateCancelled, StateDeferred},
|
||
StateCompleted: {StateRecognizing, StateCancelled, StateDeferred},
|
||
StateRecognizing: {StateLinking, StateReview, StateCancelled, StateDeferred},
|
||
StateReview: {StateLinking, StateRecognizing, StateCancelled, StateDeferred, StateOrphaned, StateDeleted},
|
||
StateLinking: {StateDone, StateReview, StateFailed, StateCancelled, StateDeferred},
|
||
StateDone: {StateReverted, StateTargetMissing, StateOrphaned, StateDeleted},
|
||
StateDeferred: {StateLinking, StateRecognizing, StateCancelled, StateOrphaned, StateDeleted},
|
||
StateStuck: {StateDownloading, StateCompleted, StateCancelled, StateDeferred},
|
||
StateFailed: {StateDownloading, StateCompleted},
|
||
StateReverted: {StateRecognizing, StateOrphaned, StateDeleted},
|
||
StateCancelled: {StateRecognizing, StateOrphaned, StateDeleted},
|
||
StateTargetMissing: {StateRecognizing, StateDone, StateOrphaned, StateDeleted},
|
||
StateOrphaned: {StateDone, StateTargetMissing, StateDeleted},
|
||
StateDeleted: nil, // окончательно терминально: сверка его не переоценивает
|
||
}
|
||
|
||
// transitionSources — обратное отображение (to → множество легальных from),
|
||
// построенное из allowedTransitions один раз при инициализации пакета. Сам to
|
||
// всегда входит в своё множество (самопереход). Основа предиката setState
|
||
// `state IN (...)`.
|
||
var transitionSources = invertTransitions(allowedTransitions)
|
||
|
||
// invertTransitions переворачивает граф from→to в to→from, добавляя каждому
|
||
// состоянию его самого (самопереход всегда легален).
|
||
func invertTransitions(fwd map[State][]State) map[State][]State {
|
||
into := make(map[State][]State, len(fwd))
|
||
ensureSelf := func(s State) {
|
||
if !slices.Contains(into[s], s) {
|
||
into[s] = append(into[s], s)
|
||
}
|
||
}
|
||
for from, tos := range fwd {
|
||
ensureSelf(from)
|
||
for _, to := range tos {
|
||
ensureSelf(to)
|
||
if !slices.Contains(into[to], from) {
|
||
into[to] = append(into[to], from)
|
||
}
|
||
}
|
||
}
|
||
return into
|
||
}
|
||
|
||
// SourceType — вид источника загрузки.
|
||
type SourceType string
|
||
|
||
const (
|
||
SourceMagnet SourceType = "magnet"
|
||
SourceTorrent SourceType = "torrent"
|
||
SourceURL SourceType = "url"
|
||
)
|
||
|
||
// Виды инфохэша (download_infohash.kind).
|
||
const (
|
||
HashV1 = "v1" // SHA-1, 40 hex
|
||
HashV2 = "v2" // SHA-256, 64 hex
|
||
)
|
||
|
||
// HashKind — вид инфохэша по длине hex-строки: 64 — v2, иначе v1.
|
||
func HashKind(h string) string {
|
||
if len(h) == 64 {
|
||
return HashV2
|
||
}
|
||
return HashV1
|
||
}
|
||
|
||
// NormalizeHash приводит инфохэш к канонической форме хранения (нижний hex —
|
||
// в этом же виде его отдаёт qBittorrent).
|
||
func NormalizeHash(h string) string {
|
||
return strings.ToLower(strings.TrimSpace(h))
|
||
}
|
||
|
||
// Infohash — строка таблицы download_infohash: один из хешей загрузки
|
||
// (у одной загрузки их несколько: v1/v2 гибридного торрента).
|
||
type Infohash struct {
|
||
DownloadID string `db:"download_id"`
|
||
Infohash string `db:"infohash"`
|
||
Kind string `db:"kind"`
|
||
}
|
||
|
||
// Download — строка таблицы download.
|
||
type Download struct {
|
||
ID string `db:"id"` // ULID (lowercase), публичный ключ домена
|
||
SourceType SourceType `db:"source_type"`
|
||
SourceRef string `db:"source_ref"`
|
||
DisplayName string `db:"display_name"` // имя раздачи (rename в qBittorrent), заголовок в веб-UI
|
||
Context string `db:"context"`
|
||
State State `db:"state"`
|
||
ErrorCode sql.NullString `db:"error_code"`
|
||
ErrorMsg sql.NullString `db:"error_msg"`
|
||
// SourceMissCount — счётчик подряд идущих тиков сверки без раздачи в
|
||
// qBittorrent (дебаунс пропажи источника, см. state-reconciliation).
|
||
SourceMissCount int `db:"source_miss_count"`
|
||
// SourceAddedAt — время добавления торрента в qBittorrent (added_on), базис
|
||
// сортировки списка. NULL, пока воркер не наблюдал раздачу. Хранится в
|
||
// формате RFC 3339 (UTC, суффикс Z), как created_at.
|
||
SourceAddedAt sql.NullString `db:"source_added_at"`
|
||
CreatedAt string `db:"created_at"`
|
||
UpdatedAt string `db:"updated_at"`
|
||
|
||
// Infohashes — хеши загрузки (download_infohash); подгружаются вместе с
|
||
// записью методами чтения store (v1 раньше v2 — порядок стабильный).
|
||
Infohashes []Infohash `db:"-"`
|
||
|
||
// RecTitle — распознанное название текущей попытки (LEFT JOIN recognition).
|
||
// Заполняется только листингом ListDownloadsPage для фолбека заголовка; в
|
||
// прочих выборках остаётся пустым.
|
||
RecTitle sql.NullString `db:"rec_title"`
|
||
}
|
||
|
||
// HashList — все хеши загрузки списком (для сопоставления с qBittorrent).
|
||
func (d Download) HashList() []string {
|
||
out := make([]string, len(d.Infohashes))
|
||
for i, h := range d.Infohashes {
|
||
out[i] = h.Infohash
|
||
}
|
||
return out
|
||
}
|
||
|
||
// PrimaryInfohash — первый известный хеш (v1 приоритетно) для показа и
|
||
// scoped-логгера; пустая строка, если хешей нет.
|
||
func (d Download) PrimaryInfohash() string {
|
||
if len(d.Infohashes) == 0 {
|
||
return ""
|
||
}
|
||
return d.Infohashes[0].Infohash
|
||
}
|
||
|
||
// ParseTime разбирает временную метку хранилища (RFC 3339, всегда UTC).
|
||
func ParseTime(s string) (time.Time, error) {
|
||
return time.Parse(time.RFC3339, s)
|
||
}
|
||
|
||
// FormatTime форматирует время в формат меток хранилища — RFC 3339 в UTC
|
||
// (суффикс Z), напр. «2006-01-02T15:04:05Z». Единый формат всех меток; DEFAULT
|
||
// в схеме нет, время всегда пишет приложение через Now (см. ниже). Фиксированная
|
||
// ширина сохраняет лексикографическое сравнение строк времени = хронологию
|
||
// (COALESCE(source_added_at, created_at) в сортировке списка).
|
||
func FormatTime(t time.Time) string {
|
||
return t.UTC().Format(time.RFC3339)
|
||
}
|
||
|
||
// Now — единая точка получения «сейчас» для меток времени store (UTC). Всё время
|
||
// в БД генерирует приложение через неё (аналогично ident.NewID для id), а не
|
||
// SQLite DEFAULT: один источник формата, тестируемая точка.
|
||
func Now() time.Time { return time.Now().UTC() }
|
||
|
||
// CreatedTime возвращает время создания загрузки как time.Time (UTC).
|
||
func (d Download) CreatedTime() (time.Time, error) { return ParseTime(d.CreatedAt) }
|
||
|
||
// NullString строит sql.NullString: пустая строка → NULL.
|
||
func NullString(s string) sql.NullString {
|
||
return sql.NullString{String: s, Valid: s != ""}
|
||
}
|
||
|
||
// CreateDownloadIfNoActive атомарно (одна write-транзакция, BEGIN IMMEDIATE
|
||
// через _txlock) проверяет инвариант «не более одной активной загрузки на
|
||
// infohash» и заводит загрузку: если активная задача с любым из hashes уже
|
||
// есть — возвращает её (дедуп, ничего не создавая); иначе вставляет d с новым
|
||
// ULID и его хешами и возвращает (nil, nil). d.ID и d.Infohashes заполняются.
|
||
//
|
||
// torrentBlob (если непуст) — исходные байты `.torrent`; пишутся в
|
||
// download_torrent В ТОЙ ЖЕ транзакции только на ветке создания (для
|
||
// source_type=torrent, чтобы воркер добавил раздачу файлом). При дедупе байты
|
||
// не пишутся. Для magnet/url — nil.
|
||
func (s *Store) CreateDownloadIfNoActive(ctx context.Context, d *Download, hashes []string, torrentBlob []byte) (*Download, error) {
|
||
norm := normalizeHashes(hashes)
|
||
if len(norm) == 0 {
|
||
return nil, fmt.Errorf("create download: no infohash")
|
||
}
|
||
now := FormatTime(Now())
|
||
|
||
tx, err := s.DB.BeginTxx(ctx, nil)
|
||
if err != nil {
|
||
return nil, fmt.Errorf("create download: begin tx: %w", err)
|
||
}
|
||
defer func() { _ = tx.Rollback() }()
|
||
|
||
existing, err := findActiveByInfohash(ctx, tx, norm, "")
|
||
if err != nil {
|
||
return nil, fmt.Errorf("create download: %w", err)
|
||
}
|
||
if existing != nil {
|
||
// Дедуп нашёл активного владельца по одному из хешей — остальные хеши
|
||
// norm принадлежат тому же торренту (гибридный magnet): дописываем
|
||
// недостающие, иначе второй хеш молча теряется и последующий приём по
|
||
// нему создал бы вторую активную задачу.
|
||
for _, h := range norm {
|
||
if _, err := tx.ExecContext(ctx,
|
||
`INSERT OR IGNORE INTO download_infohash (download_id, infohash, kind, created_at) VALUES (?, ?, ?, ?)`,
|
||
existing.ID, h, HashKind(h), now); err != nil {
|
||
return nil, fmt.Errorf("create download: top up infohash: %w", err)
|
||
}
|
||
}
|
||
if err := attachInfohashesOne(ctx, tx, existing); err != nil {
|
||
return nil, fmt.Errorf("create download: %w", err)
|
||
}
|
||
if err := tx.Commit(); err != nil {
|
||
return nil, fmt.Errorf("create download: commit dedup: %w", err)
|
||
}
|
||
return existing, nil
|
||
}
|
||
|
||
d.ID = ident.NewID()
|
||
const q = `
|
||
INSERT INTO download (id, source_type, source_ref, display_name, context, state, created_at, updated_at)
|
||
VALUES (?, ?, ?, ?, ?, ?, ?, ?)`
|
||
if _, err := tx.ExecContext(ctx, q,
|
||
d.ID, d.SourceType, d.SourceRef, d.DisplayName, d.Context, d.State, now, now); err != nil {
|
||
return nil, fmt.Errorf("insert download: %w", err)
|
||
}
|
||
d.Infohashes = d.Infohashes[:0]
|
||
for _, h := range norm {
|
||
if _, err := tx.ExecContext(ctx,
|
||
`INSERT INTO download_infohash (download_id, infohash, kind, created_at) VALUES (?, ?, ?, ?)`,
|
||
d.ID, h, HashKind(h), now); err != nil {
|
||
return nil, fmt.Errorf("insert download infohash: %w", err)
|
||
}
|
||
d.Infohashes = append(d.Infohashes, Infohash{DownloadID: d.ID, Infohash: h, Kind: HashKind(h)})
|
||
}
|
||
if len(torrentBlob) > 0 {
|
||
if _, err := tx.ExecContext(ctx,
|
||
`INSERT INTO download_torrent (download_id, data) VALUES (?, ?)`,
|
||
d.ID, torrentBlob); err != nil {
|
||
return nil, fmt.Errorf("insert download torrent: %w", err)
|
||
}
|
||
}
|
||
if err := tx.Commit(); err != nil {
|
||
return nil, fmt.Errorf("create download: commit: %w", err)
|
||
}
|
||
return nil, nil
|
||
}
|
||
|
||
// GetTorrentData возвращает сохранённые байты `.torrent` загрузки (для
|
||
// добавления раздачи файлом воркером). ErrNotFound — если байтов нет (не
|
||
// torrent-источник или запись отсутствует).
|
||
func (s *Store) GetTorrentData(ctx context.Context, downloadID string) ([]byte, error) {
|
||
var data []byte
|
||
err := s.DB.QueryRowxContext(ctx,
|
||
`SELECT data FROM download_torrent WHERE download_id = ?`, downloadID).Scan(&data)
|
||
if errors.Is(err, sql.ErrNoRows) {
|
||
return nil, fmt.Errorf("torrent data for %s: %w", downloadID, ErrNotFound)
|
||
}
|
||
if err != nil {
|
||
return nil, fmt.Errorf("get torrent data: %w", err)
|
||
}
|
||
return data, nil
|
||
}
|
||
|
||
// ActivateIfNoOtherActive атомарно возвращает загрузку в активное состояние
|
||
// (retry/восстановление сверкой/relink): в одной write-транзакции проверяет,
|
||
// что никакая ДРУГАЯ активная загрузка не владеет любым из хешей этой, и
|
||
// переводит состояние. При владении возвращает ErrInfohashTaken (обёрнутый
|
||
// с id владельца).
|
||
func (s *Store) ActivateIfNoOtherActive(ctx context.Context, id string, state State, errCode, errMsg string) error {
|
||
tx, err := s.DB.BeginTxx(ctx, nil)
|
||
if err != nil {
|
||
return fmt.Errorf("activate %s: begin tx: %w", id, err)
|
||
}
|
||
defer func() { _ = tx.Rollback() }()
|
||
|
||
var hashes []string
|
||
if err := tx.SelectContext(ctx, &hashes,
|
||
`SELECT infohash FROM download_infohash WHERE download_id = ?`, id); err != nil {
|
||
return fmt.Errorf("activate %s: read hashes: %w", id, err)
|
||
}
|
||
if len(hashes) > 0 {
|
||
// Исключаем саму задачу: при retry из stuck она сама активна и без
|
||
// исключения LIMIT 1 мог бы вернуть её, замаскировав другого владельца.
|
||
other, err := findActiveByInfohash(ctx, tx, hashes, id)
|
||
if err != nil {
|
||
return fmt.Errorf("activate %s: %w", id, err)
|
||
}
|
||
if other != nil {
|
||
return fmt.Errorf("activate %s: infohash owned by download %s: %w",
|
||
id, other.ID, ErrInfohashTaken)
|
||
}
|
||
}
|
||
if err := setState(ctx, tx, id, state, errCode, errMsg, true); err != nil {
|
||
return err
|
||
}
|
||
if err := tx.Commit(); err != nil {
|
||
return fmt.Errorf("activate %s: commit: %w", id, err)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// AddInfohashes дописывает загрузке недостающие хеши (qBittorrent раскрыл
|
||
// оба хеша гибридного торрента, а приём знал один). Это тоже мутация
|
||
// владения хешем, поэтому она под тем же гардом, что и create/activate:
|
||
// в одной write-транзакции каждый хеш проверяется на владение ДРУГОЙ
|
||
// активной задачей; конфликтные хеши не дописываются, метод возвращает
|
||
// ErrInfohashTaken (неконфликтные при этом дописаны — частичный успех).
|
||
func (s *Store) AddInfohashes(ctx context.Context, downloadID string, hashes []string) error {
|
||
norm := normalizeHashes(hashes)
|
||
if len(norm) == 0 {
|
||
return nil
|
||
}
|
||
now := FormatTime(Now())
|
||
tx, err := s.DB.BeginTxx(ctx, nil)
|
||
if err != nil {
|
||
return fmt.Errorf("add infohashes to %s: begin tx: %w", downloadID, err)
|
||
}
|
||
defer func() { _ = tx.Rollback() }()
|
||
|
||
var taken []string
|
||
for _, h := range norm {
|
||
other, err := findActiveByInfohash(ctx, tx, []string{h}, downloadID)
|
||
if err != nil {
|
||
return fmt.Errorf("add infohashes to %s: %w", downloadID, err)
|
||
}
|
||
if other != nil {
|
||
taken = append(taken, h)
|
||
continue
|
||
}
|
||
if _, err := tx.ExecContext(ctx,
|
||
`INSERT OR IGNORE INTO download_infohash (download_id, infohash, kind, created_at) VALUES (?, ?, ?, ?)`,
|
||
downloadID, h, HashKind(h), now); err != nil {
|
||
return fmt.Errorf("add infohash %s to %s: %w", h, downloadID, err)
|
||
}
|
||
}
|
||
if err := tx.Commit(); err != nil {
|
||
return fmt.Errorf("add infohashes to %s: commit: %w", downloadID, err)
|
||
}
|
||
if len(taken) > 0 {
|
||
return fmt.Errorf("add infohashes to %s: %v owned by another active download: %w",
|
||
downloadID, taken, ErrInfohashTaken)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// GetDownload возвращает загрузку по id (с хешами).
|
||
func (s *Store) GetDownload(ctx context.Context, id string) (*Download, error) {
|
||
var d Download
|
||
if err := s.DB.GetContext(ctx, &d, `SELECT * FROM download WHERE id = ?`, id); err != nil {
|
||
if errors.Is(err, sql.ErrNoRows) {
|
||
return nil, fmt.Errorf("get download %s: %w", id, ErrNotFound)
|
||
}
|
||
return nil, fmt.Errorf("get download %s: %w", id, err)
|
||
}
|
||
if err := attachInfohashesOne(ctx, s.DB, &d); err != nil {
|
||
return nil, fmt.Errorf("get download %s: %w", id, err)
|
||
}
|
||
return &d, nil
|
||
}
|
||
|
||
// SetSourceAddedAt однократно фиксирует время добавления торрента в источник
|
||
// (qBittorrent added_on). SQL-гард `source_added_at IS NULL` не даёт перезапи-
|
||
// сать значение на повторных наблюдениях: время добавления неизменно.
|
||
func (s *Store) SetSourceAddedAt(ctx context.Context, id string, t time.Time) error {
|
||
const q = `UPDATE download SET source_added_at = ? WHERE id = ? AND source_added_at IS NULL`
|
||
if _, err := s.DB.ExecContext(ctx, q, FormatTime(t), id); err != nil {
|
||
return fmt.Errorf("set source added at %s: %w", id, err)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// ListDownloads возвращает все загрузки, новые сверху (id — ULID, сортировка
|
||
// по нему хронологична).
|
||
func (s *Store) ListDownloads(ctx context.Context) ([]Download, error) {
|
||
var out []Download
|
||
if err := s.DB.SelectContext(ctx, &out, `SELECT * FROM download ORDER BY id DESC`); err != nil {
|
||
return nil, fmt.Errorf("list downloads: %w", err)
|
||
}
|
||
if err := s.attachInfohashes(ctx, out); err != nil {
|
||
return nil, fmt.Errorf("list downloads: %w", err)
|
||
}
|
||
return out, nil
|
||
}
|
||
|
||
// ListDownloadsByState возвращает загрузки в одном из указанных состояний.
|
||
func (s *Store) ListDownloadsByState(ctx context.Context, states ...State) ([]Download, error) {
|
||
if len(states) == 0 {
|
||
return nil, nil
|
||
}
|
||
ph := make([]string, len(states))
|
||
args := make([]any, len(states))
|
||
for i, st := range states {
|
||
ph[i] = "?"
|
||
args[i] = string(st)
|
||
}
|
||
q := `SELECT * FROM download WHERE state IN (` + strings.Join(ph, ",") + `) ORDER BY id DESC`
|
||
var out []Download
|
||
if err := s.DB.SelectContext(ctx, &out, q, args...); err != nil {
|
||
return nil, fmt.Errorf("list downloads by state: %w", err)
|
||
}
|
||
if err := s.attachInfohashes(ctx, out); err != nil {
|
||
return nil, fmt.Errorf("list downloads by state: %w", err)
|
||
}
|
||
return out, nil
|
||
}
|
||
|
||
// ListRecoverable возвращает задачи в failed/stuck с одним из переданных
|
||
// error_code — кандидатов на авто-восстановление (см. state-reconciliation).
|
||
// Фильтр по коду в SQL, чтобы не вычитывать на каждом тике поллинга все
|
||
// накопленные провалы (qbit_error и пр.), которые восстановлению не подлежат.
|
||
func (s *Store) ListRecoverable(ctx context.Context, codes ...string) ([]Download, error) {
|
||
if len(codes) == 0 {
|
||
return nil, nil
|
||
}
|
||
ph := make([]string, len(codes))
|
||
args := make([]any, 0, len(codes)+2)
|
||
for i, c := range codes {
|
||
ph[i] = "?"
|
||
args = append(args, c)
|
||
}
|
||
args = append(args, string(StateFailed), string(StateStuck))
|
||
q := `SELECT * FROM download WHERE error_code IN (` + strings.Join(ph, ",") +
|
||
`) AND state IN (?, ?) ORDER BY id DESC`
|
||
var out []Download
|
||
if err := s.DB.SelectContext(ctx, &out, q, args...); err != nil {
|
||
return nil, fmt.Errorf("list recoverable: %w", err)
|
||
}
|
||
if err := s.attachInfohashes(ctx, out); err != nil {
|
||
return nil, fmt.Errorf("list recoverable: %w", err)
|
||
}
|
||
return out, nil
|
||
}
|
||
|
||
// FindActiveByInfohash возвращает незавершённую задачу, владеющую любым из
|
||
// hashes, либо (nil, nil). Читающая основа дедупа; сам инвариант держат
|
||
// guarded-методы (CreateDownloadIfNoActive / ActivateIfNoOtherActive).
|
||
func (s *Store) FindActiveByInfohash(ctx context.Context, hashes ...string) (*Download, error) {
|
||
d, err := findActiveByInfohash(ctx, s.DB, normalizeHashes(hashes), "")
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if d != nil {
|
||
if err := attachInfohashesOne(ctx, s.DB, d); err != nil {
|
||
return nil, err
|
||
}
|
||
}
|
||
return d, nil
|
||
}
|
||
|
||
// findActiveByInfohash — общая выборка «активная задача по любому из хешей»
|
||
// (для guarded-методов — внутри их транзакции). hashes уже нормализованы;
|
||
// excludeID исключает саму проверяемую задачу (она может быть активной,
|
||
// например stuck при retry, — и не должна маскировать другого владельца);
|
||
// пустой excludeID — без исключения. Хеши найденной загрузки НЕ подгружаются.
|
||
func findActiveByInfohash(ctx context.Context, q sqlx.QueryerContext, hashes []string, excludeID string) (*Download, error) {
|
||
if len(hashes) == 0 {
|
||
return nil, nil
|
||
}
|
||
// «Активна» = не в терминальном состоянии. Список — единый с IsTerminal
|
||
// (terminalStates), иначе семантика активности разъедется.
|
||
var args []any
|
||
hashPh := placeholders(&args, hashes)
|
||
statePh := placeholders(&args, terminalStates)
|
||
query := `SELECT download.* FROM download
|
||
JOIN download_infohash dh ON dh.download_id = download.id
|
||
WHERE dh.infohash IN (` + hashPh + `) AND download.state NOT IN (` + statePh + `)`
|
||
if excludeID != "" {
|
||
query += ` AND download.id != ?`
|
||
args = append(args, excludeID)
|
||
}
|
||
query += ` ORDER BY download.id DESC LIMIT 1`
|
||
var d Download
|
||
err := sqlx.GetContext(ctx, q, &d, query, args...)
|
||
if errors.Is(err, sql.ErrNoRows) {
|
||
return nil, nil
|
||
}
|
||
if err != nil {
|
||
return nil, fmt.Errorf("find active by infohash: %w", err)
|
||
}
|
||
return &d, nil
|
||
}
|
||
|
||
// placeholders дописывает значения in в args и возвращает строку "?,?,…"
|
||
// той же длины — сборка IN-списков без ручного жонглирования срезами.
|
||
func placeholders[T ~string](args *[]any, in []T) string {
|
||
ph := make([]string, len(in))
|
||
for i, v := range in {
|
||
ph[i] = "?"
|
||
*args = append(*args, string(v))
|
||
}
|
||
return strings.Join(ph, ",")
|
||
}
|
||
|
||
// ExistsByInfohash сообщает, есть ли хоть одна загрузка (в любом состоянии)
|
||
// с любым из hashes. Discovery усыновляет раздачу только если её ещё не
|
||
// видели — так готовые задачи не переобрабатываются на каждом тике.
|
||
func (s *Store) ExistsByInfohash(ctx context.Context, hashes ...string) (bool, error) {
|
||
norm := normalizeHashes(hashes)
|
||
if len(norm) == 0 {
|
||
return false, nil
|
||
}
|
||
ph := make([]string, len(norm))
|
||
args := make([]any, len(norm))
|
||
for i, h := range norm {
|
||
ph[i] = "?"
|
||
args[i] = h
|
||
}
|
||
var n int
|
||
if err := s.DB.GetContext(ctx, &n,
|
||
`SELECT COUNT(1) FROM download_infohash WHERE infohash IN (`+strings.Join(ph, ",")+`)`,
|
||
args...); err != nil {
|
||
return false, fmt.Errorf("exists by infohash: %w", err)
|
||
}
|
||
return n > 0, nil
|
||
}
|
||
|
||
// SetDownloadState переводит загрузку в новое состояние. Механический
|
||
// бэкстоп инварианта «одна активная загрузка на infohash» (заменяет
|
||
// удалённый unique-индекс): переход из терминального состояния в активное
|
||
// этим методом отклоняется — возврат в активное идёт ТОЛЬКО через
|
||
// ActivateIfNoOtherActive, который проверяет владение хешами.
|
||
func (s *Store) SetDownloadState(ctx context.Context, id string, state State, errCode, errMsg string) error {
|
||
return setState(ctx, s.DB, id, state, errCode, errMsg, false)
|
||
}
|
||
|
||
// PromoteCatched переводит пойманную загрузку catched → downloading, попутно
|
||
// записывая выведенное отображаемое имя. Ребро catched → downloading объявлено в
|
||
// allowedTransitions; setState этот путь не проходит намеренно — собственный гард
|
||
// `state = 'catched'` жёстче (фиксирует ровно from=catched) и служит ре-валидацией:
|
||
// если загрузку успели отменить (catched → cancelled) во время вывода имени/
|
||
// добавления вне блокировки переходов, UPDATE не заденет ни строки и вернёт ошибку,
|
||
// а переход не применится. Пустое имя допустимо (rename не задавали) — тогда
|
||
// display_name так и остаётся пустым.
|
||
func (s *Store) PromoteCatched(ctx context.Context, id, displayName string) error {
|
||
res, err := s.DB.ExecContext(ctx, `
|
||
UPDATE download
|
||
SET state = ?, display_name = ?, updated_at = ?
|
||
WHERE id = ? AND state = ?`,
|
||
string(StateDownloading), displayName, FormatTime(Now()), id, string(StateCatched))
|
||
if err != nil {
|
||
return fmt.Errorf("promote catched %s: %w", id, err)
|
||
}
|
||
n, err := res.RowsAffected()
|
||
if err != nil {
|
||
return fmt.Errorf("promote catched %s: %w", id, err)
|
||
}
|
||
if n == 0 {
|
||
return fmt.Errorf("promote catched %s: not in catched (already added or cancelled)", id)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// setState выполняет UPDATE состояния. reviveOK=true — вызов из гарда
|
||
// (ActivateIfNoOtherActive), которому переход терминал→активное разрешён;
|
||
// иначе предикат в UPDATE не даёт молча оживить терминальную задачу.
|
||
func setState(ctx context.Context, e sqlx.ExecerContext, id string, state State, errCode, errMsg string, reviveOK bool) error {
|
||
sources := transitionSources[state]
|
||
if len(sources) == 0 {
|
||
// Целевое состояние не объявлено в графе переходов — fail-closed (это
|
||
// программная ошибка: новая State без ребра). Тест графа ловит на этапе CI.
|
||
return fmt.Errorf("set download %s state %q: target not declared in transition graph", id, state)
|
||
}
|
||
q := `
|
||
UPDATE download
|
||
SET state = ?,
|
||
error_code = ?,
|
||
error_msg = ?,
|
||
updated_at = ?
|
||
WHERE id = ?`
|
||
args := []any{string(state), nullArg(errCode), nullArg(errMsg), FormatTime(Now()), id}
|
||
// Гейт графа: переход применяется, только если текущее состояние — легальный
|
||
// источник для state (объявленное ребро или самопереход from == to). Аддитивен
|
||
// к гарду терминальности ниже и его НЕ ослабляет.
|
||
q += ` AND state IN (` + placeholders(&args, sources) + `)`
|
||
if !reviveOK && !state.IsTerminal() {
|
||
q += ` AND state NOT IN (` + placeholders(&args, terminalStates) + `)`
|
||
}
|
||
res, err := e.ExecContext(ctx, q, args...)
|
||
if err != nil {
|
||
return fmt.Errorf("set download %s state %q: %w", id, state, err)
|
||
}
|
||
n, err := res.RowsAffected()
|
||
if err != nil {
|
||
return fmt.Errorf("set download %s state %q: %w", id, state, err)
|
||
}
|
||
if n == 0 {
|
||
return setStateRejected(ctx, e, id, state)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// setStateRejected формирует точную ошибку отклонённого перехода (0 строк UPDATE):
|
||
// читает текущее состояние и различает «не найдено» / «нелегальное ребро» /
|
||
// «терминал без revive». Только путь ошибки (редкий), поэтому доп. чтение дёшево.
|
||
func setStateRejected(ctx context.Context, e sqlx.ExecerContext, id string, state State) error {
|
||
q, ok := e.(sqlx.QueryerContext)
|
||
if !ok {
|
||
return fmt.Errorf("set download %s state %q: rejected (not found, illegal transition, or terminal without revive)", id, state)
|
||
}
|
||
var cur State
|
||
err := sqlx.GetContext(ctx, q, &cur, `SELECT state FROM download WHERE id = ?`, id)
|
||
if errors.Is(err, sql.ErrNoRows) {
|
||
return fmt.Errorf("set download %s state %q: %w", id, state, ErrNotFound)
|
||
}
|
||
if err != nil {
|
||
return fmt.Errorf("set download %s state %q: rejected, current state unreadable: %w", id, state, err)
|
||
}
|
||
if slices.Contains(transitionSources[state], cur) {
|
||
// Ребро cur → state легально — значит зарубил гард терминальности.
|
||
return fmt.Errorf("set download %s: %s → %s rejected: terminal revive requires ActivateIfNoOtherActive",
|
||
id, cur, state)
|
||
}
|
||
return fmt.Errorf("set download %s: illegal transition %s → %s (not in transition graph)", id, cur, state)
|
||
}
|
||
|
||
// SetSourceMissCount записывает счётчик пропусков источника (дебаунс сверки).
|
||
// Состояние не трогает — это отдельная от перехода фоновая отметка.
|
||
func (s *Store) SetSourceMissCount(ctx context.Context, id string, n int) error {
|
||
res, err := s.DB.ExecContext(ctx,
|
||
`UPDATE download SET source_miss_count = ? WHERE id = ?`, n, id)
|
||
if err != nil {
|
||
return fmt.Errorf("set download %s source_miss_count: %w", id, err)
|
||
}
|
||
if affected, _ := res.RowsAffected(); affected == 0 {
|
||
return fmt.Errorf("set download %s source_miss_count: not found", id)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// attachInfohashes подгружает хеши для набора загрузок батч-запросами
|
||
// (порядок стабильный: v1 раньше v2). IN-список режется на чанки, чтобы
|
||
// безлимитные выборки (ListDownloads за всю историю) не упирались в
|
||
// SQLITE_MAX_VARIABLE_NUMBER.
|
||
func (s *Store) attachInfohashes(ctx context.Context, ds []Download) error {
|
||
if len(ds) == 0 {
|
||
return nil
|
||
}
|
||
const chunkSize = 500
|
||
byID := make(map[string][]Infohash, len(ds))
|
||
for start := 0; start < len(ds); start += chunkSize {
|
||
end := min(start+chunkSize, len(ds))
|
||
var args []any
|
||
ids := make([]string, 0, end-start)
|
||
for i := start; i < end; i++ {
|
||
ids = append(ids, ds[i].ID)
|
||
}
|
||
ph := placeholders(&args, ids)
|
||
var rows []Infohash
|
||
if err := s.DB.SelectContext(ctx, &rows,
|
||
`SELECT download_id, infohash, kind FROM download_infohash
|
||
WHERE download_id IN (`+ph+`) ORDER BY kind, infohash`,
|
||
args...); err != nil {
|
||
return fmt.Errorf("attach infohashes: %w", err)
|
||
}
|
||
for _, r := range rows {
|
||
byID[r.DownloadID] = append(byID[r.DownloadID], r)
|
||
}
|
||
}
|
||
for i := range ds {
|
||
ds[i].Infohashes = byID[ds[i].ID]
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// attachInfohashesOne подгружает хеши одной загрузки (в т.ч. внутри tx).
|
||
func attachInfohashesOne(ctx context.Context, q sqlx.QueryerContext, d *Download) error {
|
||
if err := sqlx.SelectContext(ctx, q, &d.Infohashes,
|
||
`SELECT download_id, infohash, kind FROM download_infohash
|
||
WHERE download_id = ? ORDER BY kind, infohash`, d.ID); err != nil {
|
||
return fmt.Errorf("attach infohashes: %w", err)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// normalizeHashes нормализует и дедуплицирует хеши, отбрасывая пустые.
|
||
func normalizeHashes(hashes []string) []string {
|
||
out := make([]string, 0, len(hashes))
|
||
for _, h := range hashes {
|
||
h = NormalizeHash(h)
|
||
if h == "" || slices.Contains(out, h) {
|
||
continue
|
||
}
|
||
out = append(out, h)
|
||
}
|
||
return out
|
||
}
|
||
|
||
// nullArg возвращает nil для пустой строки (чтобы писать NULL, не "").
|
||
func nullArg(s string) any {
|
||
if s == "" {
|
||
return nil
|
||
}
|
||
return s
|
||
}
|