Все сущности переехали с INTEGER AUTOINCREMENT на TEXT ULID (lowercase, internal/ident — единая точка генерации и разбора; oklog/ulid). Инфохэши загрузки — множество (download_infohash, v1/v2 гибридных торрентов): дедуп и сопоставление в поллинге по любому из хешей, magnet-парсер отдаёт оба хеша гибридной ссылки, усечённый v2-хеш v2-only раздач не хранится. Инвариант «не более одной активной загрузки на infohash» вместо снятого unique-индекса держат guarded-методы store в одной write-транзакции (_txlock=immediate): CreateDownloadIfNoActive (приём/adopt, с доносом недостающих хешей), ActivateIfNoOtherActive (retry/recovery/relink, отказ до побочных эффектов), guarded AddInfohashes; SetDownloadState отклоняет терминал→активное как механический бэкстоп. Миграция 0006 — первая Go-миграция goose: пересоздание таблиц при включённых FK, backfill ULID с timestamp из created_at (хронология id сохранена), разнос infohash, удаление idempotency_key. BREAKING: формат id в URL/логах/Telegram, REST-поля id (string) и infohashes (список). Новая конвенция docs/conventions/database.md (без числовых PK), корреляция в логах grep'ом по голому ULID, ER-схема обновлена. Спеки: новая capability identity, MODIFIED в state-reconciliation; change заархивирован. Пройдены ревью дизайна и кода (по 8 углов), все находки исправлены с регрессионными тестами. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
440 lines
16 KiB
Go
440 lines
16 KiB
Go
// Package migrations содержит Go-миграции goose (SQL-миграции лежат рядом
|
|
// *.sql-файлами и прогоняются из embed FS пакета store). Регистрация — в
|
|
// init(); чтобы она сработала, пакет blank-импортируется из store.
|
|
package migrations
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"fmt"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/pressly/goose/v3"
|
|
|
|
"git.vakhrushev.me/av/jellybit/internal/ident"
|
|
)
|
|
|
|
func init() {
|
|
goose.AddMigrationContext(upUlidIdentity, downUlidIdentity)
|
|
}
|
|
|
|
// upUlidIdentity переводит все таблицы на ULID-идентификаторы (TEXT PK),
|
|
// разносит download.infohash в download_infohash и убирает idempotency_key
|
|
// (см. openspec/changes/ulid-identity/design.md, D6).
|
|
//
|
|
// Работает при включённых foreign_keys (PRAGMA внутри транзакции — no-op),
|
|
// поэтому порядок жёсткий: новые таблицы и данные — родители первыми, DROP
|
|
// старых — дети первыми, затем RENAME (SQLite ≥ 3.25 переписывает REFERENCES
|
|
// в ссылающихся таблицах).
|
|
func upUlidIdentity(ctx context.Context, tx *sql.Tx) error {
|
|
if err := createNewTables(ctx, tx); err != nil {
|
|
return err
|
|
}
|
|
|
|
downloadIDs, err := migrateDownloads(ctx, tx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
recognitionIDs, err := migrateRecognitions(ctx, tx, downloadIDs)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := migrateHints(ctx, tx, downloadIDs); err != nil {
|
|
return err
|
|
}
|
|
if err := migrateOverrides(ctx, tx, downloadIDs); err != nil {
|
|
return err
|
|
}
|
|
if err := migrateCandidates(ctx, tx, recognitionIDs); err != nil {
|
|
return err
|
|
}
|
|
if err := migrateFileLinks(ctx, tx, downloadIDs); err != nil {
|
|
return err
|
|
}
|
|
|
|
// Старые таблицы: дети первыми, родитель последним (FK включены).
|
|
for _, stmt := range []string{
|
|
`DROP TABLE file_link`,
|
|
`DROP TABLE metadata_candidate`,
|
|
`DROP TABLE override`,
|
|
`DROP TABLE hint`,
|
|
`DROP TABLE recognition`,
|
|
`DROP TABLE download`,
|
|
`ALTER TABLE download_new RENAME TO download`,
|
|
`ALTER TABLE recognition_new RENAME TO recognition`,
|
|
`ALTER TABLE hint_new RENAME TO hint`,
|
|
`ALTER TABLE override_new RENAME TO override`,
|
|
`ALTER TABLE metadata_candidate_new RENAME TO metadata_candidate`,
|
|
`ALTER TABLE file_link_new RENAME TO file_link`,
|
|
// Индексы — после переименований, с каноническими именами (старые
|
|
// одноимённые ушли вместе со старыми таблицами).
|
|
`CREATE INDEX idx_download_state ON download (state)`,
|
|
`CREATE INDEX idx_download_infohash_download ON download_infohash (download_id)`,
|
|
`CREATE INDEX idx_recognition_download ON recognition (download_id)`,
|
|
`CREATE INDEX idx_hint_download ON hint (download_id)`,
|
|
`CREATE INDEX idx_candidate_recognition ON metadata_candidate (recognition_id)`,
|
|
`CREATE INDEX idx_file_link_download ON file_link (download_id)`,
|
|
`CREATE INDEX idx_file_link_batch ON file_link (apply_batch_id)`,
|
|
} {
|
|
if _, err := tx.ExecContext(ctx, stmt); err != nil {
|
|
return fmt.Errorf("ulid migration: %q: %w", stmt, err)
|
|
}
|
|
}
|
|
|
|
return checkForeignKeys(ctx, tx)
|
|
}
|
|
|
|
// downUlidIdentity: обратной миграции нет — ULID → числовые id невосстановимы.
|
|
// Откат — восстановление файла БД из копии (см. design.md, Migration Plan).
|
|
func downUlidIdentity(context.Context, *sql.Tx) error {
|
|
return fmt.Errorf("ulid identity migration is irreversible; restore the database file from a backup")
|
|
}
|
|
|
|
func createNewTables(ctx context.Context, tx *sql.Tx) error {
|
|
for _, stmt := range []string{
|
|
`CREATE TABLE download_new (
|
|
id TEXT PRIMARY KEY,
|
|
source_type TEXT NOT NULL,
|
|
source_ref TEXT NOT NULL,
|
|
display_name TEXT NOT NULL DEFAULT '',
|
|
context TEXT NOT NULL DEFAULT '',
|
|
state TEXT NOT NULL,
|
|
error_code TEXT,
|
|
error_msg TEXT,
|
|
source_miss_count INTEGER NOT NULL DEFAULT 0,
|
|
source_added_at TEXT,
|
|
created_at TEXT NOT NULL DEFAULT (datetime('now')),
|
|
updated_at TEXT NOT NULL DEFAULT (datetime('now'))
|
|
)`,
|
|
`CREATE TABLE download_infohash (
|
|
download_id TEXT NOT NULL REFERENCES download_new (id) ON DELETE CASCADE,
|
|
infohash TEXT NOT NULL,
|
|
kind TEXT NOT NULL,
|
|
created_at TEXT NOT NULL DEFAULT (datetime('now')),
|
|
PRIMARY KEY (infohash, download_id)
|
|
)`,
|
|
`CREATE TABLE recognition_new (
|
|
id TEXT PRIMARY KEY,
|
|
download_id TEXT NOT NULL REFERENCES download_new (id) ON DELETE CASCADE,
|
|
attempt_no INTEGER NOT NULL DEFAULT 1,
|
|
is_current INTEGER NOT NULL DEFAULT 1,
|
|
media_type TEXT,
|
|
title TEXT,
|
|
original_title TEXT,
|
|
year INTEGER,
|
|
provider TEXT,
|
|
provider_id TEXT,
|
|
confidence REAL,
|
|
reasons TEXT NOT NULL DEFAULT '[]',
|
|
raw_llm TEXT,
|
|
plan TEXT,
|
|
created_at TEXT NOT NULL DEFAULT (datetime('now'))
|
|
)`,
|
|
`CREATE TABLE hint_new (
|
|
id TEXT PRIMARY KEY,
|
|
download_id TEXT NOT NULL REFERENCES download_new (id) ON DELETE CASCADE,
|
|
text TEXT NOT NULL,
|
|
created_at TEXT NOT NULL DEFAULT (datetime('now'))
|
|
)`,
|
|
`CREATE TABLE override_new (
|
|
id TEXT PRIMARY KEY,
|
|
download_id TEXT NOT NULL REFERENCES download_new (id) ON DELETE CASCADE,
|
|
field TEXT NOT NULL,
|
|
value TEXT NOT NULL,
|
|
created_at TEXT NOT NULL DEFAULT (datetime('now')),
|
|
UNIQUE (download_id, field)
|
|
)`,
|
|
`CREATE TABLE metadata_candidate_new (
|
|
id TEXT PRIMARY KEY,
|
|
recognition_id TEXT NOT NULL REFERENCES recognition_new (id) ON DELETE CASCADE,
|
|
provider TEXT NOT NULL,
|
|
provider_id TEXT NOT NULL,
|
|
title TEXT,
|
|
year INTEGER,
|
|
chosen INTEGER NOT NULL DEFAULT 0,
|
|
url TEXT,
|
|
created_at TEXT NOT NULL DEFAULT (datetime('now'))
|
|
)`,
|
|
`CREATE TABLE file_link_new (
|
|
id TEXT PRIMARY KEY,
|
|
download_id TEXT NOT NULL REFERENCES download_new (id) ON DELETE CASCADE,
|
|
apply_batch_id TEXT NOT NULL,
|
|
src_path TEXT NOT NULL,
|
|
dst_path TEXT NOT NULL,
|
|
kind TEXT NOT NULL,
|
|
status TEXT NOT NULL,
|
|
created_at TEXT NOT NULL DEFAULT (datetime('now'))
|
|
)`,
|
|
} {
|
|
if _, err := tx.ExecContext(ctx, stmt); err != nil {
|
|
return fmt.Errorf("ulid migration: create tables: %w", err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// idMap строит маппинг «старый int id → ULID» для таблицы: строки читаются в
|
|
// порядке старого id (хронология), timestamp-часть ULID — из created_at, так
|
|
// что лексикографический порядок новых id сохраняет исторический. Равные
|
|
// секунды created_at упорядочивает monotonic-энтропия по порядку обхода.
|
|
func idMap(ctx context.Context, tx *sql.Tx, table string) (map[int64]string, error) {
|
|
rows, err := tx.QueryContext(ctx,
|
|
`SELECT id, created_at FROM `+table+` ORDER BY id`) //nolint:gosec // имена таблиц — константы этого файла
|
|
if err != nil {
|
|
return nil, fmt.Errorf("ulid migration: read %s ids: %w", table, err)
|
|
}
|
|
defer func() { _ = rows.Close() }()
|
|
|
|
out := map[int64]string{}
|
|
for rows.Next() {
|
|
var id int64
|
|
var createdAt string
|
|
if err := rows.Scan(&id, &createdAt); err != nil {
|
|
return nil, fmt.Errorf("ulid migration: scan %s id: %w", table, err)
|
|
}
|
|
out[id] = ident.NewIDAt(parseCreatedAt(createdAt))
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// parseCreatedAt разбирает метку datetime('now') (UTC); непарсибельная метка
|
|
// → текущее время (порядок в пределах таблицы всё равно монотонен).
|
|
func parseCreatedAt(s string) time.Time {
|
|
t, err := time.ParseInLocation("2006-01-02 15:04:05", s, time.UTC)
|
|
if err != nil {
|
|
return time.Now()
|
|
}
|
|
return t
|
|
}
|
|
|
|
func migrateDownloads(ctx context.Context, tx *sql.Tx) (map[int64]string, error) {
|
|
ids, err := idMap(ctx, tx, "download")
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
// Сначала вычитываем всё и закрываем курсор, потом вставляем: Tx держит
|
|
// одно соединение, Exec при открытых Rows на нём невозможен.
|
|
type downloadRow struct {
|
|
id int64
|
|
sourceType, sourceRef, displayName string
|
|
contextText, state string
|
|
infohash, errorCode, errorMsg sql.NullString
|
|
sourceMissCount int
|
|
sourceAddedAt sql.NullString
|
|
createdAt, updatedAt string
|
|
}
|
|
var all []downloadRow
|
|
rows, err := tx.QueryContext(ctx, `
|
|
SELECT id, source_type, source_ref, display_name, context, infohash, state,
|
|
error_code, error_msg, source_miss_count, source_added_at,
|
|
created_at, updated_at
|
|
FROM download ORDER BY id`)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("ulid migration: read downloads: %w", err)
|
|
}
|
|
for rows.Next() {
|
|
var r downloadRow
|
|
if err := rows.Scan(&r.id, &r.sourceType, &r.sourceRef, &r.displayName,
|
|
&r.contextText, &r.infohash, &r.state, &r.errorCode, &r.errorMsg,
|
|
&r.sourceMissCount, &r.sourceAddedAt, &r.createdAt, &r.updatedAt); err != nil {
|
|
_ = rows.Close()
|
|
return nil, fmt.Errorf("ulid migration: scan download: %w", err)
|
|
}
|
|
all = append(all, r)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
_ = rows.Close()
|
|
return nil, fmt.Errorf("ulid migration: iterate downloads: %w", err)
|
|
}
|
|
_ = rows.Close()
|
|
|
|
for _, r := range all {
|
|
if _, err := tx.ExecContext(ctx, `
|
|
INSERT INTO download_new (id, source_type, source_ref, display_name, context,
|
|
state, error_code, error_msg, source_miss_count,
|
|
source_added_at, created_at, updated_at)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
|
|
ids[r.id], r.sourceType, r.sourceRef, r.displayName, r.contextText,
|
|
r.state, r.errorCode, r.errorMsg, r.sourceMissCount, r.sourceAddedAt,
|
|
r.createdAt, r.updatedAt); err != nil {
|
|
return nil, fmt.Errorf("ulid migration: insert download %d: %w", r.id, err)
|
|
}
|
|
if r.infohash.Valid && r.infohash.String != "" {
|
|
h := strings.ToLower(r.infohash.String)
|
|
if _, err := tx.ExecContext(ctx, `
|
|
INSERT INTO download_infohash (download_id, infohash, kind, created_at)
|
|
VALUES (?, ?, ?, ?)`,
|
|
ids[r.id], h, hashKind(h), r.createdAt); err != nil {
|
|
return nil, fmt.Errorf("ulid migration: insert infohash for %d: %w", r.id, err)
|
|
}
|
|
}
|
|
}
|
|
return ids, nil
|
|
}
|
|
|
|
// hashKind — вид инфохэша по длине hex: 40 — v1 (SHA-1), 64 — v2 (SHA-256).
|
|
func hashKind(h string) string {
|
|
if len(h) == 64 {
|
|
return "v2"
|
|
}
|
|
return "v1"
|
|
}
|
|
|
|
func migrateRecognitions(ctx context.Context, tx *sql.Tx, downloads map[int64]string) (map[int64]string, error) {
|
|
ids, err := idMap(ctx, tx, "recognition")
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if err := copyRows(ctx, tx, copySpec{
|
|
from: "recognition", to: "recognition_new",
|
|
cols: []string{"attempt_no", "is_current", "media_type", "title", "original_title", "year", "provider", "provider_id", "confidence", "reasons", "raw_llm", "plan", "created_at"},
|
|
ids: ids,
|
|
parent: parentRef{col: "download_id", ids: downloads},
|
|
}); err != nil {
|
|
return nil, err
|
|
}
|
|
return ids, nil
|
|
}
|
|
|
|
func migrateHints(ctx context.Context, tx *sql.Tx, downloads map[int64]string) error {
|
|
ids, err := idMap(ctx, tx, "hint")
|
|
if err != nil {
|
|
return err
|
|
}
|
|
return copyRows(ctx, tx, copySpec{
|
|
from: "hint", to: "hint_new",
|
|
cols: []string{"text", "created_at"},
|
|
ids: ids,
|
|
parent: parentRef{col: "download_id", ids: downloads},
|
|
})
|
|
}
|
|
|
|
func migrateOverrides(ctx context.Context, tx *sql.Tx, downloads map[int64]string) error {
|
|
ids, err := idMap(ctx, tx, "override")
|
|
if err != nil {
|
|
return err
|
|
}
|
|
return copyRows(ctx, tx, copySpec{
|
|
from: "override", to: "override_new",
|
|
cols: []string{"field", "value", "created_at"},
|
|
ids: ids,
|
|
parent: parentRef{col: "download_id", ids: downloads},
|
|
})
|
|
}
|
|
|
|
func migrateCandidates(ctx context.Context, tx *sql.Tx, recognitions map[int64]string) error {
|
|
ids, err := idMap(ctx, tx, "metadata_candidate")
|
|
if err != nil {
|
|
return err
|
|
}
|
|
return copyRows(ctx, tx, copySpec{
|
|
from: "metadata_candidate", to: "metadata_candidate_new",
|
|
cols: []string{"provider", "provider_id", "title", "year", "chosen", "url", "created_at"},
|
|
ids: ids,
|
|
parent: parentRef{col: "recognition_id", ids: recognitions},
|
|
})
|
|
}
|
|
|
|
func migrateFileLinks(ctx context.Context, tx *sql.Tx, downloads map[int64]string) error {
|
|
ids, err := idMap(ctx, tx, "file_link")
|
|
if err != nil {
|
|
return err
|
|
}
|
|
return copyRows(ctx, tx, copySpec{
|
|
from: "file_link", to: "file_link_new",
|
|
cols: []string{"apply_batch_id", "src_path", "dst_path", "kind", "status", "created_at"},
|
|
ids: ids,
|
|
parent: parentRef{col: "download_id", ids: downloads},
|
|
})
|
|
}
|
|
|
|
// copySpec описывает перенос таблицы: собственный маппинг id, FK-родитель и
|
|
// прочие столбцы, копируемые как есть.
|
|
type copySpec struct {
|
|
from, to string
|
|
cols []string
|
|
ids map[int64]string
|
|
parent parentRef
|
|
}
|
|
|
|
type parentRef struct {
|
|
col string
|
|
ids map[int64]string
|
|
}
|
|
|
|
func copyRows(ctx context.Context, tx *sql.Tx, spec copySpec) error {
|
|
colList := strings.Join(spec.cols, ", ")
|
|
// Сначала вычитываем всё и закрываем курсор, потом вставляем (Exec при
|
|
// открытых Rows на соединении транзакции невозможен).
|
|
type rowData struct {
|
|
oldID, parentID int64
|
|
rest []any
|
|
}
|
|
var all []rowData
|
|
//nolint:gosec // имена таблиц/столбцов — константы этого файла
|
|
rows, err := tx.QueryContext(ctx, fmt.Sprintf(
|
|
`SELECT id, %s, %s FROM %s ORDER BY id`, spec.parent.col, colList, spec.from))
|
|
if err != nil {
|
|
return fmt.Errorf("ulid migration: read %s: %w", spec.from, err)
|
|
}
|
|
for rows.Next() {
|
|
r := rowData{rest: make([]any, len(spec.cols))}
|
|
dest := append([]any{&r.oldID, &r.parentID}, scanPtrs(r.rest)...)
|
|
if err := rows.Scan(dest...); err != nil {
|
|
_ = rows.Close()
|
|
return fmt.Errorf("ulid migration: scan %s: %w", spec.from, err)
|
|
}
|
|
all = append(all, r)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
_ = rows.Close()
|
|
return fmt.Errorf("ulid migration: iterate %s: %w", spec.from, err)
|
|
}
|
|
_ = rows.Close()
|
|
|
|
ph := strings.TrimSuffix(strings.Repeat("?, ", len(spec.cols)), ", ")
|
|
//nolint:gosec // имена таблиц/столбцов — константы этого файла
|
|
insert := fmt.Sprintf(`INSERT INTO %s (id, %s, %s) VALUES (?, ?, %s)`,
|
|
spec.to, spec.parent.col, colList, ph)
|
|
|
|
for _, r := range all {
|
|
newParent, ok := spec.parent.ids[r.parentID]
|
|
if !ok {
|
|
return fmt.Errorf("ulid migration: %s row %d references unknown %s %d",
|
|
spec.from, r.oldID, spec.parent.col, r.parentID)
|
|
}
|
|
args := append([]any{spec.ids[r.oldID], newParent}, r.rest...)
|
|
if _, err := tx.ExecContext(ctx, insert, args...); err != nil {
|
|
return fmt.Errorf("ulid migration: insert %s row %d: %w", spec.to, r.oldID, err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// scanPtrs — указатели на элементы среза для rows.Scan (значения любого типа
|
|
// SQLite едут через any и вставляются обратно как есть).
|
|
func scanPtrs(vals []any) []any {
|
|
out := make([]any, len(vals))
|
|
for i := range vals {
|
|
out[i] = &vals[i]
|
|
}
|
|
return out
|
|
}
|
|
|
|
// checkForeignKeys — финальная самопроверка целостности после rebuild.
|
|
func checkForeignKeys(ctx context.Context, tx *sql.Tx) error {
|
|
rows, err := tx.QueryContext(ctx, `PRAGMA foreign_key_check`)
|
|
if err != nil {
|
|
return fmt.Errorf("ulid migration: foreign_key_check: %w", err)
|
|
}
|
|
defer func() { _ = rows.Close() }()
|
|
if rows.Next() {
|
|
var table string
|
|
var rowid, parent, fkid any
|
|
_ = rows.Scan(&table, &rowid, &parent, &fkid)
|
|
return fmt.Errorf("ulid migration: foreign key violation in %s after rebuild", table)
|
|
}
|
|
return rows.Err()
|
|
}
|