Files

566 lines
19 KiB
Go
Raw Permalink Normal View History

package database
import (
"context"
"database/sql"
"fmt"
"os"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/jackc/pgx/v5/stdlib"
"github.com/Silo-Server/silo-server/internal/jellycompat/displayprefs"
"github.com/Silo-Server/silo-server/migrations"
)
// TestPostgresDisplayPrefsMove runs the real goose provider — which registers
// the Go move migration — against a real database, then exercises the move
// directly over seeded legacy rows. The parsing rules are unit-tested in
// internal/jellycompat/displayprefs; this covers what only a live database
// shows: registration, the table's constraints, verbatim copy through real
// text columns, and that re-running or rolling back behaves.
func TestPostgresDisplayPrefsMove(t *testing.T) {
dsn := os.Getenv("SILO_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("SILO_TEST_DATABASE_URL is not set")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
t.Fatalf("connect test database: %v", err)
}
t.Cleanup(pool.Close)
// Migrate first, then seed legacy rows and run the move directly: the
// goose version gate has already consumed the registered migration, so
// calling the function is how the upgrade path is exercised against data.
if err := RunMigrations(ctx, pool, migrations.FS, "sql"); err != nil {
t.Fatalf("initial migration: %v", err)
}
userID := seedLegacyDisplayPrefsRows(ctx, t, pool)
sqlDB := stdlib.OpenDBFromPool(pool)
t.Cleanup(func() { _ = sqlDB.Close() })
runMove := func(fn func(context.Context, *sql.Tx) error, label string) {
t.Helper()
tx, err := sqlDB.BeginTx(ctx, nil)
if err != nil {
t.Fatalf("begin %s: %v", label, err)
}
if err := fn(ctx, tx); err != nil {
t.Fatalf("%s: %v", label, err)
}
if err := tx.Commit(); err != nil {
t.Fatalf("commit %s: %v", label, err)
}
}
runMove(moveDisplayPrefs, "moveDisplayPrefs")
countLegacy := func() int {
var count int
if err := pool.QueryRow(ctx, `
SELECT COUNT(*) FROM user_settings
WHERE user_id = $1 AND key LIKE 'jellycompat:%'`, userID).Scan(&count); err != nil {
t.Fatalf("counting legacy rows: %v", err)
}
return count
}
t.Run("blobs move verbatim", func(t *testing.T) {
var value string
if err := pool.QueryRow(ctx, `
SELECT value FROM jellycompat_displayprefs
WHERE user_id = $1 AND prefs_id = 'usersettings' AND client = 'emby'`, userID).
Scan(&value); err != nil {
t.Fatalf("reading moved blob: %v", err)
}
if value != `{"SortBy":"SortName", "CustomPrefs":{"b":"2","a":"1"}}` {
t.Errorf("blob = %q, want it byte-for-byte", value)
}
// The empty client is a real identity and must survive the key split.
if err := pool.QueryRow(ctx, `
SELECT value FROM jellycompat_displayprefs
WHERE user_id = $1 AND prefs_id = 'f137a2dd' AND client = ''`, userID).
Scan(&value); err != nil {
t.Fatalf("reading empty-client blob: %v", err)
}
})
t.Run("user_settings keeps no jellycompat tenants", func(t *testing.T) {
if count := countLegacy(); count != 0 {
t.Errorf("%d jellycompat rows still ride user_settings", count)
}
var theme string
if err := pool.QueryRow(ctx, `
SELECT value FROM user_settings WHERE user_id = $1 AND key = 'ui_theme'`, userID).
Scan(&theme); err != nil || theme != "cobalt-studio" {
t.Errorf("ui_theme = (%q, %v); the move touched a non-jellycompat row", theme, err)
}
})
t.Run("unparseable rows are recorded, not silently deleted", func(t *testing.T) {
var value, reason string
err := pool.QueryRow(ctx, `
SELECT value, reason FROM user_setting_migration_rejects
WHERE user_id = $1 AND source_table = 'user_settings' AND source_key = 'jellycompat:stray'`,
userID).Scan(&value, &reason)
if err != nil {
t.Fatalf("the stray row was dropped rather than recorded: %v", err)
}
if value != "not a displayprefs blob" || reason == "" {
t.Errorf("reject = (%q, %q); the original value and a reason must survive", value, reason)
}
})
t.Run("a second run is a no-op", func(t *testing.T) {
counts := func() (blobs, rejects int) {
t.Helper()
if err := pool.QueryRow(ctx,
`SELECT COUNT(*) FROM jellycompat_displayprefs WHERE user_id = $1`, userID).
Scan(&blobs); err != nil {
t.Fatalf("counting blobs: %v", err)
}
if err := pool.QueryRow(ctx, `
SELECT COUNT(*) FROM user_setting_migration_rejects
WHERE user_id = $1 AND source_key LIKE 'jellycompat:%'`, userID).Scan(&rejects); err != nil {
t.Fatalf("counting rejects: %v", err)
}
return blobs, rejects
}
blobsBefore, rejectsBefore := counts()
runMove(moveDisplayPrefs, "moveDisplayPrefs re-run")
blobsAfter, rejectsAfter := counts()
if blobsAfter != blobsBefore || rejectsAfter != rejectsBefore {
t.Errorf("re-run changed counts: blobs %d→%d, rejects %d→%d",
blobsBefore, blobsAfter, rejectsBefore, rejectsAfter)
}
})
t.Run("rollback restores the legacy rows", func(t *testing.T) {
runMove(unmoveDisplayPrefs, "unmoveDisplayPrefs")
if count := countLegacy(); count != 3 {
t.Errorf("rollback restored %d legacy rows, want 3", count)
}
var value string
if err := pool.QueryRow(ctx, `
SELECT value FROM user_settings
WHERE user_id = $1 AND key = 'jellycompat:displayprefs:usersettings:emby'`, userID).
Scan(&value); err != nil {
t.Fatalf("reading restored blob row: %v", err)
}
if value != `{"SortBy":"SortName", "CustomPrefs":{"b":"2","a":"1"}}` {
t.Errorf("restored blob = %q, want it byte-for-byte", value)
}
var blobs int
if err := pool.QueryRow(ctx,
`SELECT COUNT(*) FROM jellycompat_displayprefs WHERE user_id = $1`, userID).
Scan(&blobs); err != nil {
t.Fatalf("counting blobs after rollback: %v", err)
}
if blobs != 0 {
t.Errorf("rollback left %d rows in jellycompat_displayprefs", blobs)
}
})
}
// TestPostgresDisplayPrefsMoveDoesNotDeleteConcurrentWrites reproduces the
// rolling-deploy window: the goose lock only excludes other migrators, so an
// old-binary app instance can commit a jellycompat row into user_settings
// while the move transaction sits between its SELECT and its deletes. Under
// READ COMMITTED each statement snapshots independently, so a pattern-based
// DELETE would see — and destroy — a row the SELECT never copied. The move
// must instead delete only the exact rows it read, leaving the late row
// stranded in user_settings for a re-run to pick up.
//
// The stall is real: an uncommitted conflicting insert on
// jellycompat_displayprefs blocks the move's ON CONFLICT insert, the late
// legacy row commits during the stall, and releasing the blocker lets the
// move finish.
func TestPostgresDisplayPrefsMoveDoesNotDeleteConcurrentWrites(t *testing.T) {
dsn := os.Getenv("SILO_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("SILO_TEST_DATABASE_URL is not set")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
t.Fatalf("connect test database: %v", err)
}
t.Cleanup(pool.Close)
if err := RunMigrations(ctx, pool, migrations.FS, "sql"); err != nil {
t.Fatalf("initial migration: %v", err)
}
var userID int
err = pool.QueryRow(ctx, `
INSERT INTO users (username, email, password_hash, role)
VALUES ('displayprefs-racetest', 'displayprefs-racetest@example.com', 'x', 'user')
ON CONFLICT (username) DO UPDATE SET email = EXCLUDED.email
RETURNING id`).Scan(&userID)
if err != nil {
t.Fatalf("seeding user: %v", err)
}
t.Cleanup(func() {
_, _ = pool.Exec(context.Background(), `DELETE FROM users WHERE id = $1`, userID)
})
for _, stmt := range []string{
`DELETE FROM user_settings WHERE user_id = $1`,
`DELETE FROM jellycompat_displayprefs WHERE user_id = $1`,
`DELETE FROM user_setting_migration_rejects WHERE user_id = $1`,
} {
if _, err := pool.Exec(ctx, stmt, userID); err != nil {
t.Fatalf("clearing prior rows: %v", err)
}
}
if _, err := pool.Exec(ctx, `
INSERT INTO user_settings (user_id, key, value) VALUES ($1, $2, $3)`,
userID, "jellycompat:displayprefs:usersettings:emby", `{"SortBy":"SortName"}`); err != nil {
t.Fatalf("seeding legacy row: %v", err)
}
sqlDB := stdlib.OpenDBFromPool(pool)
t.Cleanup(func() { _ = sqlDB.Close() })
// The blocker: an uncommitted insert on the identity the seeded legacy row
// maps to, so the move's copy insert waits on this transaction.
blockerTx, err := sqlDB.BeginTx(ctx, nil)
if err != nil {
t.Fatalf("begin blocker: %v", err)
}
blockerReleased := false
defer func() {
if !blockerReleased {
_ = blockerTx.Rollback()
}
}()
if _, err := blockerTx.ExecContext(ctx, `
INSERT INTO jellycompat_displayprefs (user_id, prefs_id, client, value)
VALUES ($1, 'usersettings', 'emby', 'blocker')`, userID); err != nil {
t.Fatalf("blocker insert: %v", err)
}
moveDone := make(chan error, 1)
go func() {
tx, err := sqlDB.BeginTx(ctx, nil)
if err != nil {
moveDone <- fmt.Errorf("begin move: %w", err)
return
}
if err := moveDisplayPrefs(ctx, tx); err != nil {
_ = tx.Rollback()
moveDone <- fmt.Errorf("moveDisplayPrefs: %w", err)
return
}
moveDone <- tx.Commit()
}()
// Wait until the move transaction is provably parked on the blocker's
// lock: its SELECT over user_settings has happened, its deletes have not.
waitDeadline := time.Now().Add(10 * time.Second)
for {
var waiting int
if err := pool.QueryRow(ctx, `
SELECT COUNT(*) FROM pg_stat_activity
WHERE wait_event_type = 'Lock' AND query LIKE '%jellycompat_displayprefs%'`).
Scan(&waiting); err != nil {
t.Fatalf("polling pg_stat_activity: %v", err)
}
if waiting > 0 {
break
}
if time.Now().After(waitDeadline) {
t.Fatal("the move never blocked on the conflicting insert")
}
time.Sleep(20 * time.Millisecond)
}
// The old binary's handler commits a fresh DisplayPreferences row now —
// after the move's SELECT, before its deletes.
const lateKey = "jellycompat:displayprefs:late:acme"
if _, err := pool.Exec(ctx, `
INSERT INTO user_settings (user_id, key, value) VALUES ($1, $2, $3)`,
userID, lateKey, `{"SortBy":"DateCreated"}`); err != nil {
t.Fatalf("committing the late legacy row: %v", err)
}
// And it updates a row the move has already read: the delete predicate
// pins the value the SELECT saw, so this newer write must survive too.
const updatedKey = "jellycompat:displayprefs:usersettings:emby"
const updatedValue = `{"SortBy":"Runtime"}`
if _, err := pool.Exec(ctx, `
UPDATE user_settings SET value = $3 WHERE user_id = $1 AND key = $2`,
userID, updatedKey, updatedValue); err != nil {
t.Fatalf("committing the late update: %v", err)
}
blockerReleased = true
if err := blockerTx.Rollback(); err != nil {
t.Fatalf("releasing blocker: %v", err)
}
if err := <-moveDone; err != nil {
t.Fatalf("move under contention: %v", err)
}
// The late row must not have been destroyed: it was never copied, so it
// must still ride user_settings.
var lateRows int
if err := pool.QueryRow(ctx, `
SELECT COUNT(*) FROM user_settings WHERE user_id = $1 AND key = $2`,
userID, lateKey).Scan(&lateRows); err != nil {
t.Fatalf("counting the late row: %v", err)
}
if lateRows != 1 {
var copied, rejected int
_ = pool.QueryRow(ctx, `
SELECT COUNT(*) FROM jellycompat_displayprefs
WHERE user_id = $1 AND prefs_id = 'late'`, userID).Scan(&copied)
_ = pool.QueryRow(ctx, `
SELECT COUNT(*) FROM user_setting_migration_rejects
WHERE user_id = $1 AND source_key = $2`, userID, lateKey).Scan(&rejected)
t.Fatalf("late row gone from user_settings (copied=%d rejected=%d): "+
"a concurrently committed row was deleted without being moved",
copied, rejected)
}
// The concurrently updated row must also still be in user_settings, with
// the newer value: the move copied the old value but its delete named that
// old value, so it must not have matched the updated row.
var survivingValue string
if err := pool.QueryRow(ctx, `
SELECT value FROM user_settings WHERE user_id = $1 AND key = $2`,
userID, updatedKey).Scan(&survivingValue); err != nil {
t.Fatalf("updated legacy row gone from user_settings: %v — "+
"a concurrently updated row was deleted with only its old value moved", err)
}
if survivingValue != updatedValue {
t.Errorf("surviving legacy value = %q, want the late update %q",
survivingValue, updatedValue)
}
// A re-run — which the stranded row exists to allow — picks it up.
tx, err := sqlDB.BeginTx(ctx, nil)
if err != nil {
t.Fatalf("begin re-run: %v", err)
}
if err := moveDisplayPrefs(ctx, tx); err != nil {
t.Fatalf("re-run: %v", err)
}
if err := tx.Commit(); err != nil {
t.Fatalf("commit re-run: %v", err)
}
var copied int
if err := pool.QueryRow(ctx, `
SELECT COUNT(*) FROM jellycompat_displayprefs
WHERE user_id = $1 AND prefs_id = 'late' AND client = 'acme'`, userID).Scan(&copied); err != nil {
t.Fatalf("counting the re-run copy: %v", err)
}
if copied != 1 {
t.Errorf("re-run copied %d late rows, want 1", copied)
}
}
// TestPostgresDisplayPrefsRollbackDoesNotDeleteConcurrentUpdates reproduces
// the inverse rolling-deploy window: a new-binary app instance updates a blob
// after the rollback has read it but before the rollback deletes it. The
// rollback may restore its older snapshot to user_settings, but it must leave
// the newer canonical value in place rather than deleting a value it never
// restored.
func TestPostgresDisplayPrefsRollbackDoesNotDeleteConcurrentUpdates(t *testing.T) {
dsn := os.Getenv("SILO_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("SILO_TEST_DATABASE_URL is not set")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
t.Fatalf("connect test database: %v", err)
}
t.Cleanup(pool.Close)
if err := RunMigrations(ctx, pool, migrations.FS, "sql"); err != nil {
t.Fatalf("initial migration: %v", err)
}
var userID int
err = pool.QueryRow(ctx, `
INSERT INTO users (username, email, password_hash, role)
VALUES ('displayprefs-rollback-racetest', 'displayprefs-rollback-racetest@example.com', 'x', 'user')
ON CONFLICT (username) DO UPDATE SET email = EXCLUDED.email
RETURNING id`).Scan(&userID)
if err != nil {
t.Fatalf("seeding user: %v", err)
}
t.Cleanup(func() {
_, _ = pool.Exec(context.Background(), `DELETE FROM users WHERE id = $1`, userID)
})
for _, stmt := range []string{
`DELETE FROM user_settings WHERE user_id = $1`,
`DELETE FROM jellycompat_displayprefs WHERE user_id = $1`,
`DELETE FROM user_setting_migration_rejects WHERE user_id = $1`,
} {
if _, err := pool.Exec(ctx, stmt, userID); err != nil {
t.Fatalf("clearing prior rows: %v", err)
}
}
const (
prefsID = "usersettings"
client = "emby"
originalValue = `{"SortBy":"SortName"}`
updatedValue = `{"SortBy":"Runtime"}`
)
if _, err := pool.Exec(ctx, `
INSERT INTO jellycompat_displayprefs (user_id, prefs_id, client, value)
VALUES ($1, $2, $3, $4)`, userID, prefsID, client, originalValue); err != nil {
t.Fatalf("seeding moved row: %v", err)
}
sqlDB := stdlib.OpenDBFromPool(pool)
t.Cleanup(func() { _ = sqlDB.Close() })
// Hold the destination identity open so the rollback stalls after reading
// the canonical row but before its insert and delete.
blockerTx, err := sqlDB.BeginTx(ctx, nil)
if err != nil {
t.Fatalf("begin blocker: %v", err)
}
blockerReleased := false
defer func() {
if !blockerReleased {
_ = blockerTx.Rollback()
}
}()
legacyKey := displayprefs.LegacyKey(prefsID, client)
if _, err := blockerTx.ExecContext(ctx, `
INSERT INTO user_settings (user_id, key, value)
VALUES ($1, $2, 'blocker')`, userID, legacyKey); err != nil {
t.Fatalf("blocker insert: %v", err)
}
rollbackDone := make(chan error, 1)
go func() {
tx, err := sqlDB.BeginTx(ctx, nil)
if err != nil {
rollbackDone <- fmt.Errorf("begin rollback: %w", err)
return
}
if err := unmoveDisplayPrefs(ctx, tx); err != nil {
_ = tx.Rollback()
rollbackDone <- fmt.Errorf("unmoveDisplayPrefs: %w", err)
return
}
rollbackDone <- tx.Commit()
}()
waitDeadline := time.Now().Add(10 * time.Second)
for {
var waiting int
if err := pool.QueryRow(ctx, `
SELECT COUNT(*) FROM pg_stat_activity
WHERE wait_event_type = 'Lock' AND query LIKE '%INSERT INTO user_settings%'`).
Scan(&waiting); err != nil {
t.Fatalf("polling pg_stat_activity: %v", err)
}
if waiting > 0 {
break
}
if time.Now().After(waitDeadline) {
t.Fatal("the rollback never blocked on the conflicting insert")
}
time.Sleep(20 * time.Millisecond)
}
if _, err := pool.Exec(ctx, `
UPDATE jellycompat_displayprefs SET value = $4
WHERE user_id = $1 AND prefs_id = $2 AND client = $3`,
userID, prefsID, client, updatedValue); err != nil {
t.Fatalf("committing the concurrent canonical update: %v", err)
}
blockerReleased = true
if err := blockerTx.Rollback(); err != nil {
t.Fatalf("releasing blocker: %v", err)
}
if err := <-rollbackDone; err != nil {
t.Fatalf("rollback under contention: %v", err)
}
var survivingValue string
if err := pool.QueryRow(ctx, `
SELECT value FROM jellycompat_displayprefs
WHERE user_id = $1 AND prefs_id = $2 AND client = $3`,
userID, prefsID, client).Scan(&survivingValue); err != nil {
t.Fatalf("concurrently updated canonical row was deleted: %v", err)
}
if survivingValue != updatedValue {
t.Errorf("surviving canonical value = %q, want %q", survivingValue, updatedValue)
}
var restoredValue string
if err := pool.QueryRow(ctx, `
SELECT value FROM user_settings WHERE user_id = $1 AND key = $2`,
userID, legacyKey).Scan(&restoredValue); err != nil {
t.Fatalf("reading restored legacy snapshot: %v", err)
}
if restoredValue != originalValue {
t.Errorf("restored legacy value = %q, want the rollback snapshot %q",
restoredValue, originalValue)
}
}
// seedLegacyDisplayPrefsRows writes the pre-cutover user_settings rows: two
// handler-written DisplayPreferences blobs, one jellycompat row only the legacy
// settings API's removed unknown-key carve-out could have produced, and a real
// user setting that must not move.
func seedLegacyDisplayPrefsRows(ctx context.Context, t *testing.T, pool *pgxpool.Pool) int {
t.Helper()
var userID int
err := pool.QueryRow(ctx, `
INSERT INTO users (username, email, password_hash, role)
VALUES ('displayprefs-migtest', 'displayprefs-migtest@example.com', 'x', 'user')
ON CONFLICT (username) DO UPDATE SET email = EXCLUDED.email
RETURNING id`).Scan(&userID)
if err != nil {
t.Fatalf("seeding user: %v", err)
}
t.Cleanup(func() {
_, _ = pool.Exec(context.Background(), `DELETE FROM users WHERE id = $1`, userID)
})
// Clear anything a prior run left so the assertions see only this seed.
for _, stmt := range []string{
`DELETE FROM user_settings WHERE user_id = $1`,
`DELETE FROM jellycompat_displayprefs WHERE user_id = $1`,
`DELETE FROM user_setting_migration_rejects WHERE user_id = $1`,
} {
if _, err := pool.Exec(ctx, stmt, userID); err != nil {
t.Fatalf("clearing prior rows: %v", err)
}
}
for key, value := range map[string]string{
"jellycompat:displayprefs:usersettings:emby": `{"SortBy":"SortName", "CustomPrefs":{"b":"2","a":"1"}}`,
"jellycompat:displayprefs:f137a2dd:": `{"SortBy":"DateCreated"}`,
"jellycompat:stray": "not a displayprefs blob",
"ui_theme": "cobalt-studio",
} {
if _, err := pool.Exec(ctx, `
INSERT INTO user_settings (user_id, key, value) VALUES ($1, $2, $3)
ON CONFLICT (user_id, key) DO UPDATE SET value = EXCLUDED.value`,
userID, key, value); err != nil {
t.Fatalf("seeding user_settings %s: %v", key, err)
}
}
return userID
}