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 отклоняется. // - deferred — легальная цель из КАЖДОГО не-терминального состояния (Defer // проверяет лишь IsTerminal); инвариант закреплён тестом, а не ручной // аккуратностью. // - cancelled — легальная цель из ЛЮБОГО состояния, кроме deleted: помимо // Cancel из не-терминальных её даёт универсальный стоп-кран Dismiss, доступный // и из терминальных (done/failed/reverted/target_missing/orphaned) — только // смена статуса, файлы/раздачу не трогает (см. state-reconciliation «Ручное // закрытие загрузки»). deleted строго терминален и цель cancelled не получает. // // Правка воркера, вводящая новое ребро, ОБЯЗАНА отразить его здесь — иначе // 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, StateCancelled}, StateDeferred: {StateLinking, StateRecognizing, StateCancelled, StateOrphaned, StateDeleted}, StateStuck: {StateDownloading, StateCompleted, StateCancelled, StateDeferred}, StateFailed: {StateDownloading, StateCompleted, StateCancelled}, StateReverted: {StateRecognizing, StateOrphaned, StateDeleted, StateCancelled}, StateCancelled: {StateRecognizing, StateOrphaned, StateDeleted}, StateTargetMissing: {StateRecognizing, StateDone, StateOrphaned, StateDeleted, StateCancelled}, StateOrphaned: {StateDone, StateTargetMissing, StateDeleted, StateCancelled}, 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"` // RetriedAt — время последнего ручного retry (RFC 3339 UTC, суффикс Z), NULL // пока задачу не повторяли. Приподнимает базис отсчёта таймаутов, чтобы // возврат в downloading не ронял задачу снова на ближайшем тике (см. // state-reconciliation «Ручной повтор»). RetriedAt sql.NullString `db:"retried_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"` // RecMediaType — тип текущей попытки распознавания (movie/series; LEFT JOIN // recognition). Как и RecTitle, заполняется только листингом // ListDownloadsPage — для значка типа в строке списка; в прочих выборках пуст. RecMediaType sql.NullString `db:"rec_media_type"` } // 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) } // RetriedTime возвращает время последнего ручного retry (UTC) и ok=false, если // задачу ещё не повторяли (retried_at NULL) или метку не разобрать. func (d Download) RetriedTime() (time.Time, bool) { if !d.RetriedAt.Valid { return time.Time{}, false } t, err := ParseTime(d.RetriedAt.String) if err != nil { return time.Time{}, false } return t, true } // 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): дописываем // недостающие, иначе второй хеш молча теряется и последующий приём по // нему создал бы вторую активную задачу. Дозапись — под тем же пер-хеш // гардом владения, что и AddInfohashes: хеш, которым владеет ДРУГАЯ // активная задача (split-identity/крафт-магнет), не дописываем, иначе // две активные владели бы одним инфохэшем (нарушение инварианта). for _, h := range norm { owner, err := findActiveByInfohash(ctx, tx, []string{h}, existing.ID) if err != nil { return nil, fmt.Errorf("create download: %w", err) } if owner != nil { continue // чужой активный владелец — не крадём хеш } 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 } // UpgradeCatchedMagnetToTorrent апгрейдит пойманную magnet-задачу до torrent // при дедупе входящих байтов `.torrent`: в одной write-транзакции сохраняет // байты и меняет source_type magnet→torrent, но ТОЛЬКО пока задача в `catched` // (воркер источник ещё не отдал в qBittorrent) и её source_type всё ещё // `magnet`. Это целевое исключение из правила «при дедупе байты не сохраняем»: // magnet на закрытом трекере без DHT метаданные не докачает и застрянет в // metaDL, а поданный пользователем `.torrent` их несёт. Возвращает true, если // апгрейд применён; (false, nil) — если задача не подходит (уже добавлена, // отменена или не magnet) либо байты пусты (защитный no-op). Гард через // RowsAffected UPDATE — та же ре-валидация состояния, что у PromoteCatched. func (s *Store) UpgradeCatchedMagnetToTorrent(ctx context.Context, downloadID string, torrentBlob []byte) (bool, error) { if len(torrentBlob) == 0 { return false, nil } tx, err := s.DB.BeginTxx(ctx, nil) if err != nil { return false, fmt.Errorf("upgrade %s to torrent: begin tx: %w", downloadID, err) } defer func() { _ = tx.Rollback() }() res, err := tx.ExecContext(ctx, ` UPDATE download SET source_type = ?, updated_at = ? WHERE id = ? AND source_type = ? AND state = ?`, string(SourceTorrent), FormatTime(Now()), downloadID, string(SourceMagnet), string(StateCatched)) if err != nil { return false, fmt.Errorf("upgrade %s to torrent: %w", downloadID, err) } n, err := res.RowsAffected() if err != nil { return false, fmt.Errorf("upgrade %s to torrent: %w", downloadID, err) } if n == 0 { // Не catched-magnet (уже добавлена/отменена/torrent) — апгрейд не нужен. if err := tx.Commit(); err != nil { return false, fmt.Errorf("upgrade %s to torrent: commit: %w", downloadID, err) } return false, nil } // У magnet-задачи блоба не было; OR REPLACE — страховка идемпотентности. if _, err := tx.ExecContext(ctx, `INSERT OR REPLACE INTO download_torrent (download_id, data) VALUES (?, ?)`, downloadID, torrentBlob); err != nil { return false, fmt.Errorf("upgrade %s to torrent: store blob: %w", downloadID, err) } if err := tx.Commit(); err != nil { return false, fmt.Errorf("upgrade %s to torrent: commit: %w", downloadID, err) } return true, 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 } // SetRetriedAt проставляет время ручного retry задачи (сброс базиса отсчёта // таймаутов, см. state-reconciliation «Ручной повтор»). В отличие от // SetSourceAddedAt перезаписывает значение: retry можно повторять, и базис // должен смещаться на каждый. func (s *Store) SetRetriedAt(ctx context.Context, id string, t time.Time) error { res, err := s.DB.ExecContext(ctx, `UPDATE download SET retried_at = ? WHERE id = ?`, FormatTime(t), id) if err != nil { return fmt.Errorf("set download %s retried_at: %w", id, err) } if n, _ := res.RowsAffected(); n == 0 { return fmt.Errorf("set download %s retried_at: not found", id) } 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 } // reingestHoldingStates — desync-состояния, которые удерживают источник ради // незакрытого намерения и потому БЛОКИРУЮТ повторный приём наравне с активными: // target_missing (источник жив, ждёт relink) и orphaned (источник пропал, запись // держит претензию на последнюю копию). Прочие терминальные (done/cancelled/ // failed/reverted/deleted) повторный приём НЕ блокируют — это осознанная свежая // попытка. «Блокирующие» = активные (не-терминальные) ∪ reingestHoldingStates. var reingestHoldingStates = []State{StateTargetMissing, StateOrphaned} // FindReingestBlockingByInfohash возвращает задачу, блокирующую повторный приём // любого из hashes: активную (строго не-терминальную) ЛИБО удерживающую источник // desync-запись (target_missing/orphaned), приоритет — активной. Либо (nil, nil). // Читающая основа расширенного дедупа приёма (пред-рид ДО создания); инвариант // «≤1 активной на infohash» держит отдельный active-гард CreateDownloadIfNoActive, // в который desync-состояния НЕ заводятся. func (s *Store) FindReingestBlockingByInfohash(ctx context.Context, hashes ...string) (*Download, error) { norm := normalizeHashes(hashes) // Активная имеет приоритет: если по хешу есть и активная, и desync-запись // (инвариант это допускает), присоединяемся к активной. d, err := findActiveByInfohash(ctx, s.DB, norm, "") if err != nil { return nil, err } if d == nil { d, err = findByInfohashInStates(ctx, s.DB, norm, reingestHoldingStates) if err != nil { return nil, err } } if d != nil { if err := attachInfohashesOne(ctx, s.DB, d); err != nil { return nil, err } } return d, nil } // findByInfohashInStates — выборка «задача по любому из хешей в одном из states» // (позитивный фильтр `state IN (...)`, в отличие от findActiveByInfohash с // `NOT IN terminalStates`). hashes уже нормализованы; хеши найденной загрузки НЕ // подгружаются. Пустые hashes/states → (nil, nil). func findByInfohashInStates(ctx context.Context, q sqlx.QueryerContext, hashes []string, states []State) (*Download, error) { if len(hashes) == 0 || len(states) == 0 { return nil, nil } var args []any hashPh := placeholders(&args, hashes) statePh := placeholders(&args, states) query := `SELECT download.* FROM download JOIN download_infohash dh ON dh.download_id = download.id WHERE dh.infohash IN (` + hashPh + `) AND download.state IN (` + statePh + `) 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 by infohash in states: %w", 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 } // SetDisplayName обновляет отображаемое имя загрузки постфактум — перелив // канонического имени после распознавания (см. capability ingest). В отличие от // PromoteCatched не двигает FSM и не завязан на состояние (гарда state нет): // обновление имени валидно и в review, и в терминальных done/orphaned. Имя — // косметика (заголовок в UI + ярлык раздачи), пути на диске не затрагивает. func (s *Store) SetDisplayName(ctx context.Context, id, name string) error { if _, err := s.DB.ExecContext(ctx, ` UPDATE download SET display_name = ?, updated_at = ? WHERE id = ?`, name, FormatTime(Now()), id); err != nil { return fmt.Errorf("set display name %s: %w", id, err) } 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 }