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 ( 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) } // 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 заполняются. func (s *Store) CreateDownloadIfNoActive(ctx context.Context, d *Download, hashes []string) (*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 err := tx.Commit(); err != nil { return nil, fmt.Errorf("create download: commit: %w", err) } return nil, 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) } // setState выполняет UPDATE состояния. reviveOK=true — вызов из гарда // (ActivateIfNoOtherActive), которому переход терминал→активное разрешён; // иначе предикат в UPDATE не даёт молча оживить терминальную задачу. func setState(ctx context.Context, e sqlx.ExecerContext, id string, state State, errCode, errMsg string, reviveOK bool) error { q := ` UPDATE download SET state = ?, error_code = ?, error_msg = ?, updated_at = ? WHERE id = ?` args := []any{string(state), nullArg(errCode), nullArg(errMsg), FormatTime(Now()), id} 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 fmt.Errorf("set download %s state %q: not found or terminal (revive requires ActivateIfNoOtherActive)", id, state) } return nil } // 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 }