Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,20 @@
-- target state rather than no-op'ing over a wrong column type (project rule
-- feedback_migration_full_restore).

-- ------------------------------------------------------------------
-- H3 (prerequisite): drop the materialized views FIRST.
-- monthly_savings_summary / daily_savings_trend / provider_savings_summary
-- all reference savings_snapshots.account_id, so the H4 ALTER COLUMN TYPE
-- below fails with "cannot alter type of a column used by a view or rule"
-- while they still exist. They are unconditionally recreated further down,
-- so dropping them up front is safe and makes this migration applyable on a
-- fresh DB (where an earlier migration already created them). DROP ... IF
-- EXISTS keeps the step idempotent across re-runs / auto-heal.
-- ------------------------------------------------------------------
DROP MATERIALIZED VIEW IF EXISTS provider_savings_summary CASCADE;
DROP MATERIALIZED VIEW IF EXISTS daily_savings_trend CASCADE;
DROP MATERIALIZED VIEW IF EXISTS monthly_savings_summary CASCADE;

-- ------------------------------------------------------------------
-- H4: widen account_id on the partitioned parent.
-- Postgres propagates the type change to all existing partitions.
Expand Down Expand Up @@ -71,12 +85,9 @@ ALTER TABLE savings_snapshots ALTER COLUMN coverage_percentage DROP DEFAULT;
-- DROP + CREATE because a materialized view's column list / GROUP BY
-- cannot be altered in place. Unique indexes are recreated for the
-- CONCURRENTLY refresh path. AVG(coverage_percentage) now skips NULLs
-- automatically (H2).
-- automatically (H2). The DROPs ran above (before the ALTER COLUMN) so
-- the CREATEs below land on a clean slate.
-- ------------------------------------------------------------------
DROP MATERIALIZED VIEW IF EXISTS provider_savings_summary CASCADE;
DROP MATERIALIZED VIEW IF EXISTS daily_savings_trend CASCADE;
DROP MATERIALIZED VIEW IF EXISTS monthly_savings_summary CASCADE;

CREATE MATERIALIZED VIEW monthly_savings_summary AS
SELECT
DATE_TRUNC('month', timestamp) as month,
Expand Down
156 changes: 138 additions & 18 deletions internal/database/postgres/migrations/migrate.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,27 +31,14 @@ const defaultAdminGroupID = "00000000-0000-5000-8000-000000000001"
// adminEmail is optional - if provided, admin user will be created after migrations complete
// adminPassword is optional - if provided, admin is created with hashed password and active=true
func RunMigrations(ctx context.Context, pool *pgxpool.Pool, migrationsPath string, adminEmail string, adminPassword string) error {
// Get database connection string from pool config (without admin email parameter - RDS Proxy doesn't support options)
dsn := buildMigrateDSN(pool.Config(), "")

// Create migrator
m, err := migrate.New(
fmt.Sprintf("file://%s", migrationsPath),
dsn,
)
// Create the migrator and run the pre-Up recovery hooks (operator force,
// then default-on dirty auto-heal). Kept in a helper so RunMigrations stays
// under the cyclomatic-complexity budget as recovery paths grow.
m, err := newMigratorWithRecovery(pool, migrationsPath)
if err != nil {
return fmt.Errorf("failed to create migrator: %w", err)
}
defer m.Close()

// One-shot operator recovery: if CUDLY_FORCE_MIGRATION_VERSION is set,
// call Force(N) before Up(). Clears the dirty flag and pins state to
// the given version. Used to recover from a partially-applied
// migration without direct DB access. Remove the env var after the
// next successful deploy.
if err := maybeForceMigrationVersion(m); err != nil {
return err
}
defer m.Close()

// Run migrations
if err := m.Up(); err != nil && err != migrate.ErrNoChange {
Expand Down Expand Up @@ -85,6 +72,52 @@ func RunMigrations(ctx context.Context, pool *pgxpool.Pool, migrationsPath strin
return nil
}

// newMigratorWithRecovery builds the golang-migrate migrator for migrationsPath
// and runs the pre-Up recovery hooks in order: the one-shot operator force
// (CUDLY_FORCE_MIGRATION_VERSION) first, then the default-on dirty auto-heal.
// On success it returns a migrator the caller must Close(); on any error it
// closes the migrator itself and returns the error so the caller never sees a
// half-initialized migrator.
//
// Ordering note: maybeAutoHealDirty runs AFTER maybeForceMigrationVersion so an
// explicit force always takes precedence (it pins+cleans first, leaving nothing
// dirty for auto-heal to act on). The auto-heal error propagates so a heal
// failure is recorded (and the app fail-opens in ensureDB) rather than being
// masked by the later dirty check.
func newMigratorWithRecovery(pool *pgxpool.Pool, migrationsPath string) (*migrate.Migrate, error) {
// Get database connection string from pool config (without admin email parameter - RDS Proxy doesn't support options)
dsn := buildMigrateDSN(pool.Config(), "")

m, err := migrate.New(
fmt.Sprintf("file://%s", migrationsPath),
dsn,
)
if err != nil {
return nil, fmt.Errorf("failed to create migrator: %w", err)
}

// One-shot operator recovery: if CUDLY_FORCE_MIGRATION_VERSION is set,
// call Force(N) before Up(). Clears the dirty flag and pins state to
// the given version. Used to recover from a partially-applied
// migration without direct DB access. Remove the env var after the
// next successful deploy.
if err := maybeForceMigrationVersion(m); err != nil {
m.Close()
return nil, err
}

// Default-on dirty auto-heal: when the schema_migrations row is dirty,
// clear the dirty flag at the CURRENT recorded version so the subsequent
// Up() re-applies any pending migrations, letting a cold start self-recover
// instead of staying broken until a manual force.
if err := maybeAutoHealDirty(m); err != nil {
m.Close()
return nil, err
}

return m, nil
}

// ensureAdminUser creates the admin user if it doesn't exist.
// When password is provided, the admin is created with a bcrypt-hashed password and active=true.
// When password is empty, the admin is created inactive and must use password reset to log in.
Expand Down Expand Up @@ -292,6 +325,93 @@ func maybeForceMigrationVersion(m *migrate.Migrate) error {
return nil
}

// maybeAutoHealDirty self-recovers a dirty schema_migrations row: when the DB
// is dirty it calls m.Force(currentRecordedVersion) to clear the dirty flag at
// the version golang-migrate last recorded, and the caller's subsequent m.Up()
// re-applies only the still-pending migrations. This runs ONCE per boot, so a
// cold start that finds a dirty DB self-heals instead of staying broken until
// an operator manually forces a version (the multi-hour outage this addresses).
//
// DEFAULT-ON (not opt-in). The whole outage was "the app started but stayed
// broken for hours needing a manual force", and TestMigrations_FullStackIdempotent
// proves re-running migrations is safe, so auto-heal runs by default. Set the
// escape hatch CUDLY_MIGRATION_AUTOHEAL=false (any strconv.ParseBool falsey
// value) to disable it in an environment whose migrations are not idempotent.
// When disabled, a dirty DB is left untouched here and surfaces as the usual
// "database is in dirty state" error after Up(); the caller (app.go ensureDB)
// still FAIL-OPENS on that error so the app always starts -- a broken schema
// surfaces via /health + the CloudWatch alarm, never via a crash-loop.
//
// CRITICAL -- Force to the CURRENT recorded version, NEVER a lower one. Up()
// then applies only the pending tail. Forcing BELOW already-applied migrations
// would re-run seed/data migrations, and some RAISE on a second run (e.g.
// 000059 errors with "a group named Purchaser already exists with a different
// id" because 000064 already relocated it). Force(current)+Up() is the only
// safe shape; Force(lower) would turn a recoverable dirty state into a hard
// failure. (Proven live during the incident.)
//
// IDEMPOTENCY INVARIANT: Force(version) clears the dirty marker WITHOUT rolling
// back the partial effects of the interrupted migration, so when Up() re-runs
// the pending tail those migrations must tolerate already-applied state
// (CREATE ... IF NOT EXISTS, DROP ... IF EXISTS, DO-blocks that no-op when the
// target already exists). The full-stack idempotency test guards this for the
// whole directory as migrations are added.
//
// CUDLY_FORCE_MIGRATION_VERSION still takes precedence: it runs earlier and
// leaves the row clean, so this is a no-op when both are set. If auto-heal
// itself fails (e.g. Force errors, or the re-applied Up() still fails), the
// error propagates to ensureDB, which fail-opens AND records the failure so
// the migration-failed alarm fires -- the app still starts.
//
// The ParseBool gate is duplicated from internal/server.getEnvBool rather than
// imported because the migrations package must not depend on internal/server,
// which would invert the dependency direction.
func maybeAutoHealDirty(m *migrate.Migrate) error {
if !autoHealEnabled() {
return nil
}

version, dirty, err := m.Version()
if err == migrate.ErrNilVersion {
// No migrations recorded yet -> nothing to heal.
return nil
}
if err != nil {
return fmt.Errorf("auto-heal: failed to read migration version: %w", err)
}
if !dirty {
return nil
}

// Force the CURRENT recorded version (never lower -- see the doc comment),
// then let the caller's Up() re-apply only the pending tail.
log.Printf("Database is DIRTY at version %d: auto-heal forcing the current version %d to clear the dirty flag, then re-applying pending migrations (set CUDLY_MIGRATION_AUTOHEAL=false to disable)", version, version)
if err := m.Force(int(version)); err != nil {
return fmt.Errorf("auto-heal: failed to force version %d to clear dirty flag: %w", version, err)
}
log.Printf("Auto-heal cleared dirty flag at version %d; proceeding to re-apply pending migrations", version)
return nil
}

// autoHealEnabled reports whether dirty auto-heal should run. DEFAULT-ON: it
// returns true unless CUDLY_MIGRATION_AUTOHEAL is explicitly set to a
// strconv.ParseBool falsey value (0/f/false/...). Unset, empty, or unparseable
// values keep auto-heal enabled. Kept local to the migrations package to avoid
// an inverted dependency on internal/server (see maybeAutoHealDirty).
func autoHealEnabled() bool {
v := os.Getenv("CUDLY_MIGRATION_AUTOHEAL")
if v == "" {
return true
}
b, err := strconv.ParseBool(v)
if err != nil {
// Unparseable value -> keep the safe default (enabled) rather than
// silently disabling self-recovery on a typo.
return true
}
return b
}

// RollbackMigrations rolls back N migrations
func RollbackMigrations(ctx context.Context, pool *pgxpool.Pool, migrationsPath string, steps int) error {
if steps <= 0 {
Expand Down
116 changes: 116 additions & 0 deletions internal/database/postgres/migrations/migrate_autoheal_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,116 @@
//go:build integration
// +build integration

package migrations_test

import (
"context"
"testing"

"github.com/LeanerCloud/CUDly/internal/database/postgres/migrations"
"github.com/LeanerCloud/CUDly/internal/database/postgres/testhelpers"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

// markDirty forces schema_migrations into the dirty state at the given version,
// reproducing what golang-migrate leaves behind when a migration is interrupted
// mid-run (Lambda timeout, ENI drop, etc.). golang-migrate keeps exactly one
// row in this table.
func markDirty(ctx context.Context, t *testing.T, container *testhelpers.PostgresContainer, version uint) {
t.Helper()
_, err := container.DB.Exec(ctx,
`UPDATE schema_migrations SET version = $1, dirty = true`, version)
require.NoError(t, err, "failed to mark schema_migrations dirty")
}

// TestMigrations_AutoHealDirty covers the DEFAULT-ON CUDLY_MIGRATION_AUTOHEAL
// behavior from migrate.go:
//
// - BY DEFAULT (flag unset), a dirty DB is healed: Force(current) clears the
// dirty flag and the subsequent Up() reaches head clean -- a cold start
// self-recovers without operator intervention.
// - With CUDLY_MIGRATION_AUTOHEAL=false, auto-heal is disabled and a dirty DB
// still returns an error with the dirty flag left intact (escape hatch).
//
// In both cases the Force target is the CURRENT recorded version (never lower),
// which is the only safe shape -- forcing below already-applied migrations
// would re-run guarded seed migrations that raise on a second run.
func TestMigrations_AutoHealDirty(t *testing.T) {
ctx := context.Background()
migrationsPath := getMigrationsPath()

t.Run("by default a dirty DB self-heals and head is reached", func(t *testing.T) {
container, err := testhelpers.SetupPostgresContainer(ctx, t)
require.NoError(t, err)
defer container.Cleanup(ctx)
pool := container.DB.Pool()

// Migrate to head, then simulate an interrupted migration by forcing
// the dirty flag at the head version.
require.NoError(t, migrations.RunMigrations(ctx, pool, migrationsPath, "", ""))
headVersion, _, err := migrations.GetMigrationVersion(ctx, pool, migrationsPath)
require.NoError(t, err)
require.Greater(t, headVersion, uint(0))
markDirty(ctx, t, container, headVersion)

// Sanity: confirm the DB is genuinely dirty before healing.
_, dirtyBefore, err := migrations.GetMigrationVersion(ctx, pool, migrationsPath)
require.NoError(t, err)
require.True(t, dirtyBefore, "precondition: DB must be dirty before the heal run")

// Flag explicitly unset: auto-heal is default-on, so the re-run must
// self-recover. t.Setenv auto-restores and forbids t.Parallel().
t.Setenv("CUDLY_MIGRATION_AUTOHEAL", "")
require.NoError(t, migrations.RunMigrations(ctx, pool, migrationsPath, "", ""),
"default-on auto-heal should clear the dirty flag and reach head without error")

versionAfter, dirtyAfter, err := migrations.GetMigrationVersion(ctx, pool, migrationsPath)
require.NoError(t, err)
assert.False(t, dirtyAfter, "auto-heal must clear the dirty flag")
assert.Equal(t, headVersion, versionAfter, "auto-heal must leave the DB at head")
})

t.Run("with CUDLY_MIGRATION_AUTOHEAL=true a dirty DB self-heals", func(t *testing.T) {
container, err := testhelpers.SetupPostgresContainer(ctx, t)
require.NoError(t, err)
defer container.Cleanup(ctx)
pool := container.DB.Pool()

require.NoError(t, migrations.RunMigrations(ctx, pool, migrationsPath, "", ""))
headVersion, _, err := migrations.GetMigrationVersion(ctx, pool, migrationsPath)
require.NoError(t, err)
markDirty(ctx, t, container, headVersion)

t.Setenv("CUDLY_MIGRATION_AUTOHEAL", "true")
require.NoError(t, migrations.RunMigrations(ctx, pool, migrationsPath, "", ""),
"explicit true should also clear the dirty flag and reach head")

versionAfter, dirtyAfter, err := migrations.GetMigrationVersion(ctx, pool, migrationsPath)
require.NoError(t, err)
assert.False(t, dirtyAfter, "auto-heal must clear the dirty flag")
assert.Equal(t, headVersion, versionAfter, "auto-heal must leave the DB at head")
})

t.Run("CUDLY_MIGRATION_AUTOHEAL=false disables auto-heal and a dirty DB errors", func(t *testing.T) {
container, err := testhelpers.SetupPostgresContainer(ctx, t)
require.NoError(t, err)
defer container.Cleanup(ctx)
pool := container.DB.Pool()

require.NoError(t, migrations.RunMigrations(ctx, pool, migrationsPath, "", ""))
headVersion, _, err := migrations.GetMigrationVersion(ctx, pool, migrationsPath)
require.NoError(t, err)
markDirty(ctx, t, container, headVersion)

// Escape hatch: the falsey value disables auto-heal, so the dirty error
// must surface (the caller fail-opens on it; the app still starts).
t.Setenv("CUDLY_MIGRATION_AUTOHEAL", "false")
err = migrations.RunMigrations(ctx, pool, migrationsPath, "", "")
require.Error(t, err, "CUDLY_MIGRATION_AUTOHEAL=false must disable auto-heal so a dirty DB errors")

_, dirtyAfter, verr := migrations.GetMigrationVersion(ctx, pool, migrationsPath)
require.NoError(t, verr)
assert.True(t, dirtyAfter, "with auto-heal disabled, the dirty flag must be left intact")
})
}
53 changes: 53 additions & 0 deletions internal/database/postgres/migrations/migrate_idempotency_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
//go:build integration
// +build integration

package migrations_test

import (
"context"
"testing"

"github.com/LeanerCloud/CUDly/internal/database/postgres/migrations"
"github.com/LeanerCloud/CUDly/internal/database/postgres/testhelpers"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

// TestMigrations_FullStackIdempotent proves the entire migration stack is
// idempotent: running it to head twice against a fresh DB leaves the second
// run a no-op and the schema_migrations row clean (not dirty).
//
// This is the invariant the opt-in CUDLY_MIGRATION_AUTOHEAL path relies on
// (maybeAutoHealDirty in migrate.go Force()s past a dirty version and lets
// Up() re-apply the pending tail). If a migration were NOT idempotent, a
// re-apply would error or corrupt data, so this test guards the whole
// directory as new migrations are added.
func TestMigrations_FullStackIdempotent(t *testing.T) {
ctx := context.Background()
migrationsPath := getMigrationsPath()

container, err := testhelpers.SetupPostgresContainer(ctx, t)
require.NoError(t, err)
defer container.Cleanup(ctx)
pool := container.DB.Pool()

// First run: migrate to head on a fresh DB.
require.NoError(t, migrations.RunMigrations(ctx, pool, migrationsPath, "", ""),
"first migration run to head should succeed on a fresh DB")

versionAfterFirst, dirtyAfterFirst, err := migrations.GetMigrationVersion(ctx, pool, migrationsPath)
require.NoError(t, err)
require.False(t, dirtyAfterFirst, "DB must not be dirty after the first run")
require.Greater(t, versionAfterFirst, uint(0), "head version should be > 0")

// Second run: must be a no-op (golang-migrate's m.Up() returns ErrNoChange,
// which RunMigrations swallows) and must not flip the dirty flag.
require.NoError(t, migrations.RunMigrations(ctx, pool, migrationsPath, "", ""),
"second migration run must be a clean no-op (ErrNoChange), proving idempotency")

versionAfterSecond, dirtyAfterSecond, err := migrations.GetMigrationVersion(ctx, pool, migrationsPath)
require.NoError(t, err)
assert.False(t, dirtyAfterSecond, "DB must not be dirty after the second (no-op) run")
assert.Equal(t, versionAfterFirst, versionAfterSecond,
"version must be unchanged after a no-op second run")
}
Loading
Loading