From c81519808ec8447645795605c254fe1325370950 Mon Sep 17 00:00:00 2001 From: Brandur Date: Sun, 7 Jul 2024 17:52:45 -0700 Subject: [PATCH] Add migration that brings `line` column to the migration table Here, add a new River migration (version 005) that brings a `line` column to the `river_migration` table, allowing non-main lines to be supported. We also teach the migrator how to use it. This one is a little trickier than it sounds (partly thanks to sqlc) because the migrator needs to behave somewhat differently depending on whether the `line` column exits yet. So for example, when migrating to 004 it needs to upsert migration records that _do not_ include a `line` field, but then when migrating to 005 or beyond, a `line` value is included. The same consideration goes for performing a down migration, or checking the existing state of migrations in the database. This is implemented by pairing driver functions related to migrations into one for pre-`line` and another for post-`line`. e.g. // MigrationDeleteAssumingMainMany deletes many migrations assuming // everything is on the main line. This is suitable for use in databases on // a version before the `line` column exists. MigrationDeleteAssumingMainMany(ctx context.Context, versions []int) ([]*Migration, error) // MigrationDeleteByLineAndVersionMany deletes many migration versions on a // particular line. MigrationDeleteByLineAndVersionMany(ctx context.Context, line string, versions []int) ([]*Migration, error) We don't yet pull line support into the CLI, but this should make alternate lines fully supported up to that point. The CLI needs a little more thought because it might involve building a separate binary. --- CHANGELOG.md | 8 + cmd/river/go.mod | 22 +- cmd/river/go.sum | 26 +-- go.mod | 2 +- .../riverdrivertest/riverdrivertest.go | 119 +++++++++- .../testfactory/test_factory.go | 8 +- riverdriver/river_driver_interface.go | 35 ++- riverdriver/riverdatabasesql/go.mod | 2 +- .../internal/dbsqlc/models.go | 1 + .../internal/dbsqlc/river_migration.sql.go | 220 ++++++++++++++++-- .../005_river_migration_add_line.down.sql | 6 + .../main/005_river_migration_add_line.up.sql | 12 + .../riverdatabasesql/river_database_sql.go | 75 +++++- riverdriver/riverpgxv5/go.mod | 2 +- .../riverpgxv5/internal/dbsqlc/models.go | 1 + .../internal/dbsqlc/river_migration.sql | 51 +++- .../internal/dbsqlc/river_migration.sql.go | 211 +++++++++++++++-- .../005_river_migration_add_line.down.sql | 6 + .../main/005_river_migration_add_line.up.sql | 12 + riverdriver/riverpgxv5/river_pgx_v5_driver.go | 75 +++++- .../migration/alternate/001_premier.down.sql | 1 + .../migration/alternate/001_premier.up.sql | 1 + .../migration/alternate/002_deuxieme.down.sql | 1 + .../migration/alternate/002_deuxieme.up.sql | 1 + .../alternate/003_troisieme.down.sql | 1 + .../migration/alternate/003_troisieme.up.sql | 1 + .../alternate/004_quatrieme.down.sql | 1 + .../migration/alternate/004_quatrieme.up.sql | 1 + .../alternate/005_cinquieme.down.sql | 1 + .../migration/alternate/005_cinquieme.up.sql | 1 + .../migration/alternate/006_sixieme.down.sql | 1 + .../migration/alternate/006_sixieme.up.sql | 1 + .../migration/main/001_first.down.sql | 1 + rivermigrate/migration/main/001_first.up.sql | 1 + .../migration/main/002_second.down.sql | 1 + rivermigrate/migration/main/002_second.up.sql | 1 + rivermigrate/river_migrate.go | 75 +++++- rivermigrate/river_migrate_test.go | 212 +++++++++++++---- 38 files changed, 1052 insertions(+), 145 deletions(-) create mode 100644 riverdriver/riverdatabasesql/migration/main/005_river_migration_add_line.down.sql create mode 100644 riverdriver/riverdatabasesql/migration/main/005_river_migration_add_line.up.sql create mode 100644 riverdriver/riverpgxv5/migration/main/005_river_migration_add_line.down.sql create mode 100644 riverdriver/riverpgxv5/migration/main/005_river_migration_add_line.up.sql create mode 100644 rivermigrate/migration/alternate/003_troisieme.down.sql create mode 100644 rivermigrate/migration/alternate/003_troisieme.up.sql create mode 100644 rivermigrate/migration/alternate/004_quatrieme.down.sql create mode 100644 rivermigrate/migration/alternate/004_quatrieme.up.sql create mode 100644 rivermigrate/migration/alternate/005_cinquieme.down.sql create mode 100644 rivermigrate/migration/alternate/005_cinquieme.up.sql create mode 100644 rivermigrate/migration/alternate/006_sixieme.down.sql create mode 100644 rivermigrate/migration/alternate/006_sixieme.up.sql diff --git a/CHANGELOG.md b/CHANGELOG.md index bbc8170f..33d8e128 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,10 +7,18 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +⚠️ Version 0.10.0 contains a new database migration, version 5. See [documentation on running River migrations](https://riverqueue.com/docs/migrations). If migrating with the CLI, make sure to update it to its latest version: + +```shell +go install github.com/riverqueue/river/cmd/river@latest +river migrate-up --database-url "$DATABASE_URL" +``` + ### Added - Fully functional driver for `database/sql` for use with packages like Bun and GORM. [PR #351](https://github.com/riverqueue/river/pull/351). - Queues can be added after a client is initialized using `client.Queues().Add(queueName string, queueConfig QueueConfig)`. [PR #410](https://github.com/riverqueue/river/pull/410). +- Migration that adds a `line` column to the `river_migration` table so that it can support multiple migration lines. [PR #435](https://github.com/riverqueue/river/pull/435). ### Changed diff --git a/cmd/river/go.mod b/cmd/river/go.mod index 19f38ec0..d20620a6 100644 --- a/cmd/river/go.mod +++ b/cmd/river/go.mod @@ -1,22 +1,22 @@ module github.com/riverqueue/river/cmd/river -go 1.21.4 +go 1.22.5 -// replace github.com/riverqueue/river => ../.. +replace github.com/riverqueue/river => ../.. -// replace github.com/riverqueue/river/riverdriver => ../../riverdriver +replace github.com/riverqueue/river/riverdriver => ../../riverdriver -// replace github.com/riverqueue/river/riverdriver/riverdatabasesql => ../../riverdriver/riverdatabasesql +replace github.com/riverqueue/river/riverdriver/riverdatabasesql => ../../riverdriver/riverdatabasesql -// replace github.com/riverqueue/river/riverdriver/riverpgxv5 => ../../riverdriver/riverpgxv5 +replace github.com/riverqueue/river/riverdriver/riverpgxv5 => ../../riverdriver/riverpgxv5 require ( - github.com/jackc/pgx/v5 v5.5.5 + github.com/jackc/pgx/v5 v5.6.0 github.com/lmittmann/tint v1.0.4 github.com/riverqueue/river v0.6.1 - github.com/riverqueue/river/riverdriver v0.6.1 - github.com/riverqueue/river/riverdriver/riverpgxv5 v0.6.1 - github.com/riverqueue/river/rivertype v0.6.1 + github.com/riverqueue/river/riverdriver v0.9.0 + github.com/riverqueue/river/riverdriver/riverpgxv5 v0.9.0 + github.com/riverqueue/river/rivertype v0.9.0 github.com/spf13/cobra v1.8.0 github.com/stretchr/testify v1.9.0 ) @@ -28,9 +28,11 @@ require ( github.com/jackc/pgservicefile v0.0.0-20231201235250-de7065d80cb9 // indirect github.com/jackc/puddle/v2 v2.2.1 // indirect github.com/pmezard/go-difflib v1.0.0 // indirect + github.com/riverqueue/river/rivershared v0.0.0-20240707210043-f9063791ecb1 // indirect github.com/spf13/pflag v1.0.5 // indirect + go.uber.org/goleak v1.3.0 // indirect golang.org/x/crypto v0.23.0 // indirect golang.org/x/sync v0.7.0 // indirect - golang.org/x/text v0.15.0 // indirect + golang.org/x/text v0.16.0 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/cmd/river/go.sum b/cmd/river/go.sum index bb88fd1e..360d634d 100644 --- a/cmd/river/go.sum +++ b/cmd/river/go.sum @@ -10,8 +10,8 @@ github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsI github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= github.com/jackc/pgservicefile v0.0.0-20231201235250-de7065d80cb9 h1:L0QtFUgDarD7Fpv9jeVMgy/+Ec0mtnmYuImjTz6dtDA= github.com/jackc/pgservicefile v0.0.0-20231201235250-de7065d80cb9/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM= -github.com/jackc/pgx/v5 v5.5.5 h1:amBjrZVmksIdNjxGW/IiIMzxMKZFelXbUoPNb+8sjQw= -github.com/jackc/pgx/v5 v5.5.5/go.mod h1:ez9gk+OAat140fv9ErkZDYFWmXLfV+++K0uAOiwgm1A= +github.com/jackc/pgx/v5 v5.6.0 h1:SWJzexBzPL5jb0GEsrPMLIsi/3jOo7RHlzTjcAeDrPY= +github.com/jackc/pgx/v5 v5.6.0/go.mod h1:DNZ/vlrUnhWCoFGxHAG8U2ljioxukquj7utPDgtQdTw= github.com/jackc/puddle/v2 v2.2.1 h1:RhxXJtFG022u4ibrCSMSiu5aOq1i77R3OHKNJj77OAk= github.com/jackc/puddle/v2 v2.2.1/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= github.com/kr/pretty v0.3.0 h1:WgNl7dwNpEZ6jJ9k1snq4pZsg7DOEN8hP9Xw0Tsjwk0= @@ -24,20 +24,14 @@ github.com/lmittmann/tint v1.0.4 h1:LeYihpJ9hyGvE0w+K2okPTGUdVLfng1+nDNVR4vWISc= github.com/lmittmann/tint v1.0.4/go.mod h1:HIS3gSy7qNwGCj+5oRjAutErFBl4BzdQP6cJZ0NfMwE= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= -github.com/riverqueue/river v0.6.1 h1:D0A139oRJh3EvATSuMPwPVscoujdezxf32Osqf/ym9Q= -github.com/riverqueue/river v0.6.1/go.mod h1:saKYj0h5bwon4377dMF5Lnq0yxtYLerjssuhJESY/nI= -github.com/riverqueue/river/riverdriver v0.6.1 h1:xIgMgqC4WPpFT3yddvfpeeO7O2xJogBugOZVkDbVpjM= -github.com/riverqueue/river/riverdriver v0.6.1/go.mod h1:l1QlqYXt+CFnVuS9ovNIvntnTeUeHuuDvEb/utvWO9k= -github.com/riverqueue/river/riverdriver/riverdatabasesql v0.6.1 h1:ahJ9gpWEBxHazwRGiwxVbG1kVd3ceeMzq4eiZvW1C1E= -github.com/riverqueue/river/riverdriver/riverdatabasesql v0.6.1/go.mod h1:58F2aScxSs4V9uxjL02kdzGeTr+ZviYdrSM9bLBC+sc= -github.com/riverqueue/river/riverdriver/riverpgxv5 v0.6.1 h1:iHy7r+AHjSMBNRTVox/8K442zxhwfQhRqXPxqNpmJsY= -github.com/riverqueue/river/riverdriver/riverpgxv5 v0.6.1/go.mod h1:XBxscxZzQCkempAXvBUwkXsh+DinBRIWnPBUNTUYd14= -github.com/riverqueue/river/rivertype v0.6.1 h1:vW64T/sJN//5XkI2tNhNwK9nDfM4lobCm3d2AZyrQ70= -github.com/riverqueue/river/rivertype v0.6.1/go.mod h1:nDd50b/mIdxR/ezQzGS/JiAhBPERA7tUIne21GdfspQ= +github.com/riverqueue/river/rivershared v0.0.0-20240707210043-f9063791ecb1 h1:wCAWAmchE27xz40FBxmcoHxRWXzhiC94GecbguFk2S4= +github.com/riverqueue/river/rivershared v0.0.0-20240707210043-f9063791ecb1/go.mod h1:2egnQ7czNcW8IXKXMRjko0aEMrQzF4V3k3jddmYiihE= +github.com/riverqueue/river/rivertype v0.9.0 h1:xr2ktQ55lqqKgXIm0Z7GJDtGuKk9BUD9kbchoUL69Lg= +github.com/riverqueue/river/rivertype v0.9.0/go.mod h1:nDd50b/mIdxR/ezQzGS/JiAhBPERA7tUIne21GdfspQ= github.com/robfig/cron/v3 v3.0.1 h1:WdRxkvbJztn8LMz/QEvLN5sBU+xKpSqwwUO1Pjr4qDs= github.com/robfig/cron/v3 v3.0.1/go.mod h1:eQICP3HwyT7UooqI/z+Ov+PtYAWygg1TEWWzGIFLtro= -github.com/rogpeppe/go-internal v1.11.0 h1:cWPaGQEPrBb5/AsnsZesgZZ9yb1OQ+GOISoDNXVBh4M= -github.com/rogpeppe/go-internal v1.11.0/go.mod h1:ddIwULY96R17DhadqLgMfk9H9tvdUzkipdSkR5nkCZA= +github.com/rogpeppe/go-internal v1.12.0 h1:exVL4IDcn6na9z1rAb56Vxr+CgyK3nn3O+epU5NdKM8= +github.com/rogpeppe/go-internal v1.12.0/go.mod h1:E+RYuTGaKKdloAfM02xzb0FW3Paa99yedzYV+kq4uf4= github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM= github.com/spf13/cobra v1.8.0 h1:7aJaZx1B85qltLMc546zn58BxxfZdR/W22ej9CFoEf0= github.com/spf13/cobra v1.8.0/go.mod h1:WXLWApfZ71AjXPya3WOlMsY9yMs7YeiHhFVlvLyhcho= @@ -54,8 +48,8 @@ golang.org/x/crypto v0.23.0 h1:dIJU/v2J8Mdglj/8rJ6UUOM3Zc9zLZxVZwwxMooUSAI= golang.org/x/crypto v0.23.0/go.mod h1:CKFgDieR+mRhux2Lsu27y0fO304Db0wZe70UKqHu0v8= golang.org/x/sync v0.7.0 h1:YsImfSBoP9QPYL0xyKJPq0gcaJdG3rInoqxTWbfQu9M= golang.org/x/sync v0.7.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk= -golang.org/x/text v0.15.0 h1:h1V/4gjBv8v9cjcR6+AR5+/cIYK5N/WAgiv4xlsEtAk= -golang.org/x/text v0.15.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU= +golang.org/x/text v0.16.0 h1:a94ExnEXNtEwYLGJSIUxnWoxoRz/ZcCsV63ROupILh4= +golang.org/x/text v0.16.0/go.mod h1:GhwF1Be+LQoKShO3cGOHzqOgRrGaYc9AvblQOmPVHnI= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= diff --git a/go.mod b/go.mod index 495150dc..fa0e20b9 100644 --- a/go.mod +++ b/go.mod @@ -19,7 +19,7 @@ require ( github.com/riverqueue/river/riverdriver v0.9.0 github.com/riverqueue/river/riverdriver/riverdatabasesql v0.9.0 github.com/riverqueue/river/riverdriver/riverpgxv5 v0.9.0 - github.com/riverqueue/river/rivershared v0.0.0-20240707170519-d0685f5e0a5d + github.com/riverqueue/river/rivershared v0.0.0-20240707210043-f9063791ecb1 github.com/riverqueue/river/rivertype v0.9.0 github.com/robfig/cron/v3 v3.0.1 github.com/stretchr/testify v1.9.0 diff --git a/internal/riverinternaltest/riverdrivertest/riverdrivertest.go b/internal/riverinternaltest/riverdrivertest/riverdrivertest.go index 9a6fea9c..1739f8a5 100644 --- a/internal/riverinternaltest/riverdrivertest/riverdrivertest.go +++ b/internal/riverinternaltest/riverdrivertest/riverdrivertest.go @@ -180,6 +180,20 @@ func Exercise[TTx any](ctx context.Context, t *testing.T, }) }) + t.Run("ColumnExists", func(t *testing.T) { + t.Parallel() + + exec, _ := setup(ctx, t) + + exists, err := exec.ColumnExists(ctx, "river_job", "id") + require.NoError(t, err) + require.True(t, exists) + + exists, err = exec.ColumnExists(ctx, "river_job", "does_not_exist") + require.NoError(t, err) + require.False(t, exists) + }) + t.Run("Exec", func(t *testing.T) { t.Parallel() @@ -1877,7 +1891,7 @@ func Exercise[TTx any](ctx context.Context, t *testing.T, require.NoError(t, err) } - t.Run("MigrationDeleteByVersionMany", func(t *testing.T) { + t.Run("MigrationDeleteAssumingMainMany", func(t *testing.T) { t.Parallel() exec, _ := setup(ctx, t) @@ -1887,18 +1901,53 @@ func Exercise[TTx any](ctx context.Context, t *testing.T, migration1 := testfactory.Migration(ctx, t, exec, &testfactory.MigrationOpts{}) migration2 := testfactory.Migration(ctx, t, exec, &testfactory.MigrationOpts{}) - migrations, err := exec.MigrationDeleteByVersionMany(ctx, []int{ + // This query is designed to work before the `line` column was added to + // the `river_migration` table. These tests will be operating on a fully + // migrated database, so drop the column in this transaction to make + // sure we are really checking that this operation works as expected. + _, err := exec.Exec(ctx, "ALTER TABLE river_migration DROP COLUMN line") + require.NoError(t, err) + + migrations, err := exec.MigrationDeleteAssumingMainMany(ctx, []int{ + migration1.Version, + migration2.Version, + }) + require.NoError(t, err) + require.Len(t, migrations, 2) + slices.SortFunc(migrations, func(a, b *riverdriver.Migration) int { return a.Version - b.Version }) + require.Equal(t, riverdriver.MigrationLineMain, migrations[0].Line) + require.Equal(t, migration1.Version, migrations[0].Version) + require.Equal(t, riverdriver.MigrationLineMain, migrations[1].Line) + require.Equal(t, migration2.Version, migrations[1].Version) + }) + + t.Run("MigrationDeleteByLineAndVersionMany", func(t *testing.T) { + t.Parallel() + + exec, _ := setup(ctx, t) + + truncateMigrations(ctx, t, exec) + + // not touched + _ = testfactory.Migration(ctx, t, exec, &testfactory.MigrationOpts{}) + + migration1 := testfactory.Migration(ctx, t, exec, &testfactory.MigrationOpts{Line: ptrutil.Ptr("alternate")}) + migration2 := testfactory.Migration(ctx, t, exec, &testfactory.MigrationOpts{Line: ptrutil.Ptr("alternate")}) + + migrations, err := exec.MigrationDeleteByLineAndVersionMany(ctx, "alternate", []int{ migration1.Version, migration2.Version, }) require.NoError(t, err) require.Len(t, migrations, 2) slices.SortFunc(migrations, func(a, b *riverdriver.Migration) int { return a.Version - b.Version }) + require.Equal(t, "alternate", migrations[0].Line) require.Equal(t, migration1.Version, migrations[0].Version) + require.Equal(t, "alternate", migrations[1].Line) require.Equal(t, migration2.Version, migrations[1].Version) }) - t.Run("MigrationGetAll", func(t *testing.T) { + t.Run("MigrationGetAllAssumingMain", func(t *testing.T) { t.Parallel() exec, _ := setup(ctx, t) @@ -1908,7 +1957,14 @@ func Exercise[TTx any](ctx context.Context, t *testing.T, migration1 := testfactory.Migration(ctx, t, exec, &testfactory.MigrationOpts{}) migration2 := testfactory.Migration(ctx, t, exec, &testfactory.MigrationOpts{}) - migrations, err := exec.MigrationGetAll(ctx) + // This query is designed to work before the `line` column was added to + // the `river_migration` table. These tests will be operating on a fully + // migrated database, so drop the column in this transaction to make + // sure we are really checking that this operation works as expected. + _, err := exec.Exec(ctx, "ALTER TABLE river_migration DROP COLUMN line") + require.NoError(t, err) + + migrations, err := exec.MigrationGetAllAssumingMain(ctx) require.NoError(t, err) require.Len(t, migrations, 2) require.Equal(t, migration1.Version, migrations[0].Version) @@ -1918,6 +1974,34 @@ func Exercise[TTx any](ctx context.Context, t *testing.T, migration1Fetched := migrations[0] require.Equal(t, migration1.ID, migration1Fetched.ID) requireEqualTime(t, migration1.CreatedAt, migration1Fetched.CreatedAt) + require.Equal(t, riverdriver.MigrationLineMain, migration1Fetched.Line) + require.Equal(t, migration1.Version, migration1Fetched.Version) + }) + + t.Run("MigrationGetByLine", func(t *testing.T) { + t.Parallel() + + exec, _ := setup(ctx, t) + + truncateMigrations(ctx, t, exec) + + // not returned + _ = testfactory.Migration(ctx, t, exec, &testfactory.MigrationOpts{}) + + migration1 := testfactory.Migration(ctx, t, exec, &testfactory.MigrationOpts{Line: ptrutil.Ptr("alternate")}) + migration2 := testfactory.Migration(ctx, t, exec, &testfactory.MigrationOpts{Line: ptrutil.Ptr("alternate")}) + + migrations, err := exec.MigrationGetByLine(ctx, "alternate") + require.NoError(t, err) + require.Len(t, migrations, 2) + require.Equal(t, migration1.Version, migrations[0].Version) + require.Equal(t, migration2.Version, migrations[1].Version) + + // Check the full properties of one of the migrations. + migration1Fetched := migrations[0] + require.Equal(t, migration1.ID, migration1Fetched.ID) + requireEqualTime(t, migration1.CreatedAt, migration1Fetched.CreatedAt) + require.Equal(t, "alternate", migration1Fetched.Line) require.Equal(t, migration1.Version, migration1Fetched.Version) }) @@ -1928,10 +2012,35 @@ func Exercise[TTx any](ctx context.Context, t *testing.T, truncateMigrations(ctx, t, exec) - migrations, err := exec.MigrationInsertMany(ctx, []int{1, 2}) + migrations, err := exec.MigrationInsertMany(ctx, "alternate", []int{1, 2}) + require.NoError(t, err) + require.Len(t, migrations, 2) + require.Equal(t, "alternate", migrations[0].Line) + require.Equal(t, 1, migrations[0].Version) + require.Equal(t, "alternate", migrations[1].Line) + require.Equal(t, 2, migrations[1].Version) + }) + + t.Run("MigrationInsertManyAssumingMain", func(t *testing.T) { + t.Parallel() + + exec, _ := setup(ctx, t) + + truncateMigrations(ctx, t, exec) + + // This query is designed to work before the `line` column was added to + // the `river_migration` table. These tests will be operating on a fully + // migrated database, so drop the column in this transaction to make + // sure we are really checking that this operation works as expected. + _, err := exec.Exec(ctx, "ALTER TABLE river_migration DROP COLUMN line") + require.NoError(t, err) + + migrations, err := exec.MigrationInsertManyAssumingMain(ctx, []int{1, 2}) require.NoError(t, err) require.Len(t, migrations, 2) + require.Equal(t, riverdriver.MigrationLineMain, migrations[0].Line) require.Equal(t, 1, migrations[0].Version) + require.Equal(t, riverdriver.MigrationLineMain, migrations[1].Line) require.Equal(t, 2, migrations[1].Version) }) diff --git a/internal/riverinternaltest/testfactory/test_factory.go b/internal/riverinternaltest/testfactory/test_factory.go index 361f885c..faba3a5c 100644 --- a/internal/riverinternaltest/testfactory/test_factory.go +++ b/internal/riverinternaltest/testfactory/test_factory.go @@ -99,15 +99,17 @@ func Leader(ctx context.Context, tb testing.TB, exec riverdriver.Executor, opts } type MigrationOpts struct { + Line *string Version *int } func Migration(ctx context.Context, tb testing.TB, exec riverdriver.Executor, opts *MigrationOpts) *riverdriver.Migration { tb.Helper() - migration, err := exec.MigrationInsertMany(ctx, []int{ - ptrutil.ValOrDefaultFunc(opts.Version, nextSeq), - }) + migration, err := exec.MigrationInsertMany(ctx, + ptrutil.ValOrDefault(opts.Line, riverdriver.MigrationLineMain), + []int{ptrutil.ValOrDefaultFunc(opts.Version, nextSeq)}, + ) require.NoError(tb, err) return migration[0] } diff --git a/riverdriver/river_driver_interface.go b/riverdriver/river_driver_interface.go index 04ab7613..2883c6ed 100644 --- a/riverdriver/river_driver_interface.go +++ b/riverdriver/river_driver_interface.go @@ -98,6 +98,10 @@ type Executor interface { // subtransactions (like riverdriver/riverdatabasesql for database/sql). Begin(ctx context.Context) (ExecutorTx, error) + // ColumnExists checks whether a column for a particular table exists for + // the schema in the current search schema. + ColumnExists(ctx context.Context, tableName, columnName string) (bool, error) + // Exec executes raw SQL. Used for migrations. Exec(ctx context.Context, sql string) (struct{}, error) @@ -129,14 +133,30 @@ type Executor interface { LeaderInsert(ctx context.Context, params *LeaderInsertParams) (*Leader, error) LeaderResign(ctx context.Context, params *LeaderResignParams) (bool, error) - // MigrationDeleteByVersionMany deletes many migration versions. - MigrationDeleteByVersionMany(ctx context.Context, versions []int) ([]*Migration, error) + // MigrationDeleteAssumingMainMany deletes many migrations assuming + // everything is on the main line. This is suitable for use in databases on + // a version before the `line` column exists. + MigrationDeleteAssumingMainMany(ctx context.Context, versions []int) ([]*Migration, error) + + // MigrationDeleteByLineAndVersionMany deletes many migration versions on a + // particular line. + MigrationDeleteByLineAndVersionMany(ctx context.Context, line string, versions []int) ([]*Migration, error) - // MigrationGetAll gets all currently applied migrations. - MigrationGetAll(ctx context.Context) ([]*Migration, error) + // MigrationGetAllAssumingMain gets all migrations assuming everything is on + // the main line. This is suitable for use in databases on a version before + // the `line` column exists. + MigrationGetAllAssumingMain(ctx context.Context) ([]*Migration, error) + + // MigrationGetByLine gets all currently applied migrations. + MigrationGetByLine(ctx context.Context, line string) ([]*Migration, error) // MigrationInsertMany inserts many migration versions. - MigrationInsertMany(ctx context.Context, versions []int) ([]*Migration, error) + MigrationInsertMany(ctx context.Context, line string, versions []int) ([]*Migration, error) + + // MigrationInsertManyAssumingMain inserts many migration, assuming they're + // on the main line. This operation is necessary for compatibility before + // the `line` column was added to the migrations table. + MigrationInsertManyAssumingMain(ctx context.Context, versions []int) ([]*Migration, error) NotifyMany(ctx context.Context, params *NotifyManyParams) error PGAdvisoryXactLock(ctx context.Context, key int64) (*struct{}, error) @@ -375,6 +395,11 @@ type Migration struct { // API is not stable. DO NOT USE. CreatedAt time.Time + // Line is the migration line that the migration belongs to. + // + // API is not stable. DO NOT USE. + Line string + // Version is the version of the migration. // // API is not stable. DO NOT USE. diff --git a/riverdriver/riverdatabasesql/go.mod b/riverdriver/riverdatabasesql/go.mod index 6deb8eea..7eb7deb9 100644 --- a/riverdriver/riverdatabasesql/go.mod +++ b/riverdriver/riverdatabasesql/go.mod @@ -11,7 +11,7 @@ replace github.com/riverqueue/river/rivertype => ../../rivertype require ( github.com/lib/pq v1.10.9 github.com/riverqueue/river/riverdriver v0.9.0 - github.com/riverqueue/river/rivershared v0.0.0-20240707170519-d0685f5e0a5d + github.com/riverqueue/river/rivershared v0.0.0-20240707210043-f9063791ecb1 github.com/riverqueue/river/rivertype v0.9.0 github.com/stretchr/testify v1.9.0 ) diff --git a/riverdriver/riverdatabasesql/internal/dbsqlc/models.go b/riverdriver/riverdatabasesql/internal/dbsqlc/models.go index 1ff48d73..e9f50ae9 100644 --- a/riverdriver/riverdatabasesql/internal/dbsqlc/models.go +++ b/riverdriver/riverdatabasesql/internal/dbsqlc/models.go @@ -87,6 +87,7 @@ type RiverLeader struct { type RiverMigration struct { ID int64 CreatedAt time.Time + Line string Version int64 } diff --git a/riverdriver/riverdatabasesql/internal/dbsqlc/river_migration.sql.go b/riverdriver/riverdatabasesql/internal/dbsqlc/river_migration.sql.go index a827ba9e..c8561392 100644 --- a/riverdriver/riverdatabasesql/internal/dbsqlc/river_migration.sql.go +++ b/riverdriver/riverdatabasesql/internal/dbsqlc/river_migration.sql.go @@ -7,18 +7,83 @@ package dbsqlc import ( "context" + "time" "github.com/lib/pq" ) -const riverMigrationDeleteByVersionMany = `-- name: RiverMigrationDeleteByVersionMany :many +const columnExists = `-- name: ColumnExists :one +SELECT EXISTS ( + SELECT column_name + FROM information_schema.columns + WHERE table_name = $1 and column_name = $2 +) +` + +type ColumnExistsParams struct { + TableName interface{} + ColumnName interface{} +} + +func (q *Queries) ColumnExists(ctx context.Context, db DBTX, arg *ColumnExistsParams) (bool, error) { + row := db.QueryRowContext(ctx, columnExists, arg.TableName, arg.ColumnName) + var exists bool + err := row.Scan(&exists) + return exists, err +} + +const riverMigrationDeleteAssumingMainMany = `-- name: RiverMigrationDeleteAssumingMainMany :many DELETE FROM river_migration WHERE version = any($1::bigint[]) -RETURNING id, created_at, version +RETURNING + id, + created_at, + version ` -func (q *Queries) RiverMigrationDeleteByVersionMany(ctx context.Context, db DBTX, version []int64) ([]*RiverMigration, error) { - rows, err := db.QueryContext(ctx, riverMigrationDeleteByVersionMany, pq.Array(version)) +type RiverMigrationDeleteAssumingMainManyRow struct { + ID int64 + CreatedAt time.Time + Version int64 +} + +func (q *Queries) RiverMigrationDeleteAssumingMainMany(ctx context.Context, db DBTX, version []int64) ([]*RiverMigrationDeleteAssumingMainManyRow, error) { + rows, err := db.QueryContext(ctx, riverMigrationDeleteAssumingMainMany, pq.Array(version)) + if err != nil { + return nil, err + } + defer rows.Close() + var items []*RiverMigrationDeleteAssumingMainManyRow + for rows.Next() { + var i RiverMigrationDeleteAssumingMainManyRow + if err := rows.Scan(&i.ID, &i.CreatedAt, &i.Version); err != nil { + return nil, err + } + items = append(items, &i) + } + if err := rows.Close(); err != nil { + return nil, err + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} + +const riverMigrationDeleteByLineAndVersionMany = `-- name: RiverMigrationDeleteByLineAndVersionMany :many +DELETE FROM river_migration +WHERE line = $1 + AND version = any($2::bigint[]) +RETURNING id, created_at, line, version +` + +type RiverMigrationDeleteByLineAndVersionManyParams struct { + Line string + Version []int64 +} + +func (q *Queries) RiverMigrationDeleteByLineAndVersionMany(ctx context.Context, db DBTX, arg *RiverMigrationDeleteByLineAndVersionManyParams) ([]*RiverMigration, error) { + rows, err := db.QueryContext(ctx, riverMigrationDeleteByLineAndVersionMany, arg.Line, pq.Array(arg.Version)) if err != nil { return nil, err } @@ -26,6 +91,54 @@ func (q *Queries) RiverMigrationDeleteByVersionMany(ctx context.Context, db DBTX var items []*RiverMigration for rows.Next() { var i RiverMigration + if err := rows.Scan( + &i.ID, + &i.CreatedAt, + &i.Line, + &i.Version, + ); err != nil { + return nil, err + } + items = append(items, &i) + } + if err := rows.Close(); err != nil { + return nil, err + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} + +const riverMigrationGetAllAssumingMain = `-- name: RiverMigrationGetAllAssumingMain :many +SELECT + id, + created_at, + version +FROM river_migration +ORDER BY version +` + +type RiverMigrationGetAllAssumingMainRow struct { + ID int64 + CreatedAt time.Time + Version int64 +} + +// This is a compatibility query for getting existing migrations before the +// `line` column was added to the table in version 005. We need to make sure to +// only select non-line properties so the query doesn't error on older schemas. +// (Even if we use `SELECT *` below, sqlc materializes it to a list of column +// names in the generated query.) +func (q *Queries) RiverMigrationGetAllAssumingMain(ctx context.Context, db DBTX) ([]*RiverMigrationGetAllAssumingMainRow, error) { + rows, err := db.QueryContext(ctx, riverMigrationGetAllAssumingMain) + if err != nil { + return nil, err + } + defer rows.Close() + var items []*RiverMigrationGetAllAssumingMainRow + for rows.Next() { + var i RiverMigrationGetAllAssumingMainRow if err := rows.Scan(&i.ID, &i.CreatedAt, &i.Version); err != nil { return nil, err } @@ -40,14 +153,15 @@ func (q *Queries) RiverMigrationDeleteByVersionMany(ctx context.Context, db DBTX return items, nil } -const riverMigrationGetAll = `-- name: RiverMigrationGetAll :many -SELECT id, created_at, version +const riverMigrationGetByLine = `-- name: RiverMigrationGetByLine :many +SELECT id, created_at, line, version FROM river_migration +WHERE line = $1 ORDER BY version ` -func (q *Queries) RiverMigrationGetAll(ctx context.Context, db DBTX) ([]*RiverMigration, error) { - rows, err := db.QueryContext(ctx, riverMigrationGetAll) +func (q *Queries) RiverMigrationGetByLine(ctx context.Context, db DBTX, line string) ([]*RiverMigration, error) { + rows, err := db.QueryContext(ctx, riverMigrationGetByLine, line) if err != nil { return nil, err } @@ -55,7 +169,12 @@ func (q *Queries) RiverMigrationGetAll(ctx context.Context, db DBTX) ([]*RiverMi var items []*RiverMigration for rows.Next() { var i RiverMigration - if err := rows.Scan(&i.ID, &i.CreatedAt, &i.Version); err != nil { + if err := rows.Scan( + &i.ID, + &i.CreatedAt, + &i.Line, + &i.Version, + ); err != nil { return nil, err } items = append(items, &i) @@ -71,30 +190,49 @@ func (q *Queries) RiverMigrationGetAll(ctx context.Context, db DBTX) ([]*RiverMi const riverMigrationInsert = `-- name: RiverMigrationInsert :one INSERT INTO river_migration ( + line, version ) VALUES ( - $1 -) RETURNING id, created_at, version + $1, + $2 +) RETURNING id, created_at, line, version ` -func (q *Queries) RiverMigrationInsert(ctx context.Context, db DBTX, version int64) (*RiverMigration, error) { - row := db.QueryRowContext(ctx, riverMigrationInsert, version) +type RiverMigrationInsertParams struct { + Line string + Version int64 +} + +func (q *Queries) RiverMigrationInsert(ctx context.Context, db DBTX, arg *RiverMigrationInsertParams) (*RiverMigration, error) { + row := db.QueryRowContext(ctx, riverMigrationInsert, arg.Line, arg.Version) var i RiverMigration - err := row.Scan(&i.ID, &i.CreatedAt, &i.Version) + err := row.Scan( + &i.ID, + &i.CreatedAt, + &i.Line, + &i.Version, + ) return &i, err } const riverMigrationInsertMany = `-- name: RiverMigrationInsertMany :many INSERT INTO river_migration ( + line, version ) SELECT - unnest($1::bigint[]) -RETURNING id, created_at, version + $1, + unnest($2::bigint[]) +RETURNING id, created_at, line, version ` -func (q *Queries) RiverMigrationInsertMany(ctx context.Context, db DBTX, version []int64) ([]*RiverMigration, error) { - rows, err := db.QueryContext(ctx, riverMigrationInsertMany, pq.Array(version)) +type RiverMigrationInsertManyParams struct { + Line string + Version []int64 +} + +func (q *Queries) RiverMigrationInsertMany(ctx context.Context, db DBTX, arg *RiverMigrationInsertManyParams) ([]*RiverMigration, error) { + rows, err := db.QueryContext(ctx, riverMigrationInsertMany, arg.Line, pq.Array(arg.Version)) if err != nil { return nil, err } @@ -102,6 +240,52 @@ func (q *Queries) RiverMigrationInsertMany(ctx context.Context, db DBTX, version var items []*RiverMigration for rows.Next() { var i RiverMigration + if err := rows.Scan( + &i.ID, + &i.CreatedAt, + &i.Line, + &i.Version, + ); err != nil { + return nil, err + } + items = append(items, &i) + } + if err := rows.Close(); err != nil { + return nil, err + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} + +const riverMigrationInsertManyAssumingMain = `-- name: RiverMigrationInsertManyAssumingMain :many +INSERT INTO river_migration ( + version +) +SELECT + unnest($1::bigint[]) +RETURNING + id, + created_at, + version +` + +type RiverMigrationInsertManyAssumingMainRow struct { + ID int64 + CreatedAt time.Time + Version int64 +} + +func (q *Queries) RiverMigrationInsertManyAssumingMain(ctx context.Context, db DBTX, version []int64) ([]*RiverMigrationInsertManyAssumingMainRow, error) { + rows, err := db.QueryContext(ctx, riverMigrationInsertManyAssumingMain, pq.Array(version)) + if err != nil { + return nil, err + } + defer rows.Close() + var items []*RiverMigrationInsertManyAssumingMainRow + for rows.Next() { + var i RiverMigrationInsertManyAssumingMainRow if err := rows.Scan(&i.ID, &i.CreatedAt, &i.Version); err != nil { return nil, err } diff --git a/riverdriver/riverdatabasesql/migration/main/005_river_migration_add_line.down.sql b/riverdriver/riverdatabasesql/migration/main/005_river_migration_add_line.down.sql new file mode 100644 index 00000000..1d554601 --- /dev/null +++ b/riverdriver/riverdatabasesql/migration/main/005_river_migration_add_line.down.sql @@ -0,0 +1,6 @@ +DROP INDEX river_migration_line_version_idx; +CREATE UNIQUE INDEX river_migration_version_idx ON river_migration USING btree(version); + +ALTER TABLE river_migration + DROP CONSTRAINT line_length, + DROP COLUMN line; \ No newline at end of file diff --git a/riverdriver/riverdatabasesql/migration/main/005_river_migration_add_line.up.sql b/riverdriver/riverdatabasesql/migration/main/005_river_migration_add_line.up.sql new file mode 100644 index 00000000..aa5115a5 --- /dev/null +++ b/riverdriver/riverdatabasesql/migration/main/005_river_migration_add_line.up.sql @@ -0,0 +1,12 @@ +ALTER TABLE river_migration + ADD COLUMN line text; + +UPDATE river_migration +SET line = 'main'; + +ALTER TABLE river_migration + ALTER COLUMN line SET NOT NULL, + ADD CONSTRAINT line_length CHECK (char_length(line) > 0 AND char_length(line) < 128); + +CREATE UNIQUE INDEX river_migration_line_version_idx ON river_migration USING btree(line, version); +DROP INDEX river_migration_version_idx; \ No newline at end of file diff --git a/riverdriver/riverdatabasesql/river_database_sql.go b/riverdriver/riverdatabasesql/river_database_sql.go index d69d8637..bbdfb812 100644 --- a/riverdriver/riverdatabasesql/river_database_sql.go +++ b/riverdriver/riverdatabasesql/river_database_sql.go @@ -81,6 +81,14 @@ func (e *Executor) Begin(ctx context.Context) (riverdriver.ExecutorTx, error) { return &ExecutorTx{Executor: Executor{nil, tx, e.queries}, tx: tx}, nil } +func (e *Executor) ColumnExists(ctx context.Context, tableName, columnName string) (bool, error) { + exists, err := e.queries.ColumnExists(ctx, e.dbtx, &dbsqlc.ColumnExistsParams{ + ColumnName: columnName, + TableName: tableName, + }) + return exists, interpretError(err) +} + func (e *Executor) Exec(ctx context.Context, sql string) (struct{}, error) { _, err := e.dbtx.ExecContext(ctx, sql) return struct{}{}, interpretError(err) @@ -487,32 +495,84 @@ func (e *Executor) LeaderResign(ctx context.Context, params *riverdriver.LeaderR return numResigned > 0, nil } -func (e *Executor) MigrationDeleteByVersionMany(ctx context.Context, versions []int) ([]*riverdriver.Migration, error) { - migrations, err := e.queries.RiverMigrationDeleteByVersionMany(ctx, e.dbtx, +func (e *Executor) MigrationDeleteAssumingMainMany(ctx context.Context, versions []int) ([]*riverdriver.Migration, error) { + migrations, err := e.queries.RiverMigrationDeleteAssumingMainMany(ctx, e.dbtx, sliceutil.Map(versions, func(v int) int64 { return int64(v) })) if err != nil { return nil, interpretError(err) } + return sliceutil.Map(migrations, func(internal *dbsqlc.RiverMigrationDeleteAssumingMainManyRow) *riverdriver.Migration { + return &riverdriver.Migration{ + ID: int(internal.ID), + CreatedAt: internal.CreatedAt.UTC(), + Line: riverdriver.MigrationLineMain, + Version: int(internal.Version), + } + }), nil +} + +func (e *Executor) MigrationDeleteByLineAndVersionMany(ctx context.Context, line string, versions []int) ([]*riverdriver.Migration, error) { + migrations, err := e.queries.RiverMigrationDeleteByLineAndVersionMany(ctx, e.dbtx, &dbsqlc.RiverMigrationDeleteByLineAndVersionManyParams{ + Line: line, + Version: sliceutil.Map(versions, func(v int) int64 { return int64(v) }), + }) + if err != nil { + return nil, interpretError(err) + } return sliceutil.Map(migrations, migrationFromInternal), nil } -func (e *Executor) MigrationGetAll(ctx context.Context) ([]*riverdriver.Migration, error) { - migrations, err := e.queries.RiverMigrationGetAll(ctx, e.dbtx) +func (e *Executor) MigrationGetAllAssumingMain(ctx context.Context) ([]*riverdriver.Migration, error) { + migrations, err := e.queries.RiverMigrationGetAllAssumingMain(ctx, e.dbtx) + if err != nil { + return nil, interpretError(err) + } + return sliceutil.Map(migrations, func(internal *dbsqlc.RiverMigrationGetAllAssumingMainRow) *riverdriver.Migration { + return &riverdriver.Migration{ + ID: int(internal.ID), + CreatedAt: internal.CreatedAt.UTC(), + Line: riverdriver.MigrationLineMain, + Version: int(internal.Version), + } + }), nil +} + +func (e *Executor) MigrationGetByLine(ctx context.Context, line string) ([]*riverdriver.Migration, error) { + migrations, err := e.queries.RiverMigrationGetByLine(ctx, e.dbtx, line) if err != nil { return nil, interpretError(err) } return sliceutil.Map(migrations, migrationFromInternal), nil } -func (e *Executor) MigrationInsertMany(ctx context.Context, versions []int) ([]*riverdriver.Migration, error) { - migrations, err := e.queries.RiverMigrationInsertMany(ctx, e.dbtx, - sliceutil.Map(versions, func(v int) int64 { return int64(v) })) +func (e *Executor) MigrationInsertMany(ctx context.Context, line string, versions []int) ([]*riverdriver.Migration, error) { + migrations, err := e.queries.RiverMigrationInsertMany(ctx, e.dbtx, &dbsqlc.RiverMigrationInsertManyParams{ + Line: line, + Version: sliceutil.Map(versions, func(v int) int64 { return int64(v) }), + }) if err != nil { return nil, interpretError(err) } return sliceutil.Map(migrations, migrationFromInternal), nil } +func (e *Executor) MigrationInsertManyAssumingMain(ctx context.Context, versions []int) ([]*riverdriver.Migration, error) { + migrations, err := e.queries.RiverMigrationInsertManyAssumingMain(ctx, e.dbtx, + sliceutil.Map(versions, func(v int) int64 { return int64(v) }), + ) + if err != nil { + return nil, interpretError(err) + } + return sliceutil.Map(migrations, func(internal *dbsqlc.RiverMigrationInsertManyAssumingMainRow) *riverdriver.Migration { + return &riverdriver.Migration{ + ID: int(internal.ID), + CreatedAt: internal.CreatedAt.UTC(), + Line: riverdriver.MigrationLineMain, + Version: int(internal.Version), + } + }), nil +} + func (e *Executor) NotifyMany(ctx context.Context, params *riverdriver.NotifyManyParams) error { return e.queries.PGNotifyMany(ctx, e.dbtx, &dbsqlc.PGNotifyManyParams{ Payload: params.Payload, @@ -785,6 +845,7 @@ func migrationFromInternal(internal *dbsqlc.RiverMigration) *riverdriver.Migrati return &riverdriver.Migration{ ID: int(internal.ID), CreatedAt: internal.CreatedAt.UTC(), + Line: internal.Line, Version: int(internal.Version), } } diff --git a/riverdriver/riverpgxv5/go.mod b/riverdriver/riverpgxv5/go.mod index d82520e5..4ab3137c 100644 --- a/riverdriver/riverpgxv5/go.mod +++ b/riverdriver/riverpgxv5/go.mod @@ -12,7 +12,7 @@ require ( github.com/jackc/pgx/v5 v5.5.0 github.com/jackc/puddle/v2 v2.2.1 github.com/riverqueue/river/riverdriver v0.9.0 - github.com/riverqueue/river/rivershared v0.0.0-20240707170519-d0685f5e0a5d + github.com/riverqueue/river/rivershared v0.0.0-20240707210043-f9063791ecb1 github.com/riverqueue/river/rivertype v0.9.0 github.com/stretchr/testify v1.9.0 ) diff --git a/riverdriver/riverpgxv5/internal/dbsqlc/models.go b/riverdriver/riverpgxv5/internal/dbsqlc/models.go index 737322cb..623a57e9 100644 --- a/riverdriver/riverpgxv5/internal/dbsqlc/models.go +++ b/riverdriver/riverpgxv5/internal/dbsqlc/models.go @@ -87,6 +87,7 @@ type RiverLeader struct { type RiverMigration struct { ID int64 CreatedAt time.Time + Line string Version int64 } diff --git a/riverdriver/riverpgxv5/internal/dbsqlc/river_migration.sql b/riverdriver/riverpgxv5/internal/dbsqlc/river_migration.sql index c2701835..80f6f5aa 100644 --- a/riverdriver/riverpgxv5/internal/dbsqlc/river_migration.sql +++ b/riverdriver/riverpgxv5/internal/dbsqlc/river_migration.sql @@ -1,35 +1,82 @@ CREATE TABLE river_migration( id bigserial PRIMARY KEY, created_at timestamptz NOT NULL DEFAULT NOW(), + line TEXT NOT NULL, version bigint NOT NULL, + CONSTRAINT line_length CHECK (char_length(line) > 0 AND char_length(line) < 128), CONSTRAINT version CHECK (version >= 1) ); --- name: RiverMigrationDeleteByVersionMany :many +-- name: RiverMigrationDeleteAssumingMainMany :many DELETE FROM river_migration WHERE version = any(@version::bigint[]) +RETURNING + id, + created_at, + version; + +-- name: RiverMigrationDeleteByLineAndVersionMany :many +DELETE FROM river_migration +WHERE line = @line + AND version = any(@version::bigint[]) RETURNING *; --- name: RiverMigrationGetAll :many +-- This is a compatibility query for getting existing migrations before the +-- `line` column was added to the table in version 005. We need to make sure to +-- only select non-line properties so the query doesn't error on older schemas. +-- (Even if we use `SELECT *` below, sqlc materializes it to a list of column +-- names in the generated query.) +-- name: RiverMigrationGetAllAssumingMain :many +SELECT + id, + created_at, + version +FROM river_migration +ORDER BY version; + +-- name: RiverMigrationGetByLine :many SELECT * FROM river_migration +WHERE line = @line ORDER BY version; -- name: RiverMigrationInsert :one INSERT INTO river_migration ( + line, version ) VALUES ( + @line, @version ) RETURNING *; -- name: RiverMigrationInsertMany :many INSERT INTO river_migration ( + line, version ) SELECT + @line, unnest(@version::bigint[]) RETURNING *; +-- name: RiverMigrationInsertManyAssumingMain :many +INSERT INTO river_migration ( + version +) +SELECT + unnest(@version::bigint[]) +RETURNING + id, + created_at, + version; + +-- name: ColumnExists :one +SELECT EXISTS ( + SELECT column_name + FROM information_schema.columns + WHERE table_name = @table_name and column_name = @column_name +); + -- name: TableExists :one SELECT CASE WHEN to_regclass(@table_name) IS NULL THEN false ELSE true END; \ No newline at end of file diff --git a/riverdriver/riverpgxv5/internal/dbsqlc/river_migration.sql.go b/riverdriver/riverpgxv5/internal/dbsqlc/river_migration.sql.go index 5775591d..70577358 100644 --- a/riverdriver/riverpgxv5/internal/dbsqlc/river_migration.sql.go +++ b/riverdriver/riverpgxv5/internal/dbsqlc/river_migration.sql.go @@ -7,16 +7,78 @@ package dbsqlc import ( "context" + "time" ) -const riverMigrationDeleteByVersionMany = `-- name: RiverMigrationDeleteByVersionMany :many +const columnExists = `-- name: ColumnExists :one +SELECT EXISTS ( + SELECT column_name + FROM information_schema.columns + WHERE table_name = $1 and column_name = $2 +) +` + +type ColumnExistsParams struct { + TableName interface{} + ColumnName interface{} +} + +func (q *Queries) ColumnExists(ctx context.Context, db DBTX, arg *ColumnExistsParams) (bool, error) { + row := db.QueryRow(ctx, columnExists, arg.TableName, arg.ColumnName) + var exists bool + err := row.Scan(&exists) + return exists, err +} + +const riverMigrationDeleteAssumingMainMany = `-- name: RiverMigrationDeleteAssumingMainMany :many DELETE FROM river_migration WHERE version = any($1::bigint[]) -RETURNING id, created_at, version +RETURNING + id, + created_at, + version ` -func (q *Queries) RiverMigrationDeleteByVersionMany(ctx context.Context, db DBTX, version []int64) ([]*RiverMigration, error) { - rows, err := db.Query(ctx, riverMigrationDeleteByVersionMany, version) +type RiverMigrationDeleteAssumingMainManyRow struct { + ID int64 + CreatedAt time.Time + Version int64 +} + +func (q *Queries) RiverMigrationDeleteAssumingMainMany(ctx context.Context, db DBTX, version []int64) ([]*RiverMigrationDeleteAssumingMainManyRow, error) { + rows, err := db.Query(ctx, riverMigrationDeleteAssumingMainMany, version) + if err != nil { + return nil, err + } + defer rows.Close() + var items []*RiverMigrationDeleteAssumingMainManyRow + for rows.Next() { + var i RiverMigrationDeleteAssumingMainManyRow + if err := rows.Scan(&i.ID, &i.CreatedAt, &i.Version); err != nil { + return nil, err + } + items = append(items, &i) + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} + +const riverMigrationDeleteByLineAndVersionMany = `-- name: RiverMigrationDeleteByLineAndVersionMany :many +DELETE FROM river_migration +WHERE line = $1 + AND version = any($2::bigint[]) +RETURNING id, created_at, line, version +` + +type RiverMigrationDeleteByLineAndVersionManyParams struct { + Line string + Version []int64 +} + +func (q *Queries) RiverMigrationDeleteByLineAndVersionMany(ctx context.Context, db DBTX, arg *RiverMigrationDeleteByLineAndVersionManyParams) ([]*RiverMigration, error) { + rows, err := db.Query(ctx, riverMigrationDeleteByLineAndVersionMany, arg.Line, arg.Version) if err != nil { return nil, err } @@ -24,6 +86,51 @@ func (q *Queries) RiverMigrationDeleteByVersionMany(ctx context.Context, db DBTX var items []*RiverMigration for rows.Next() { var i RiverMigration + if err := rows.Scan( + &i.ID, + &i.CreatedAt, + &i.Line, + &i.Version, + ); err != nil { + return nil, err + } + items = append(items, &i) + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} + +const riverMigrationGetAllAssumingMain = `-- name: RiverMigrationGetAllAssumingMain :many +SELECT + id, + created_at, + version +FROM river_migration +ORDER BY version +` + +type RiverMigrationGetAllAssumingMainRow struct { + ID int64 + CreatedAt time.Time + Version int64 +} + +// This is a compatibility query for getting existing migrations before the +// `line` column was added to the table in version 005. We need to make sure to +// only select non-line properties so the query doesn't error on older schemas. +// (Even if we use `SELECT *` below, sqlc materializes it to a list of column +// names in the generated query.) +func (q *Queries) RiverMigrationGetAllAssumingMain(ctx context.Context, db DBTX) ([]*RiverMigrationGetAllAssumingMainRow, error) { + rows, err := db.Query(ctx, riverMigrationGetAllAssumingMain) + if err != nil { + return nil, err + } + defer rows.Close() + var items []*RiverMigrationGetAllAssumingMainRow + for rows.Next() { + var i RiverMigrationGetAllAssumingMainRow if err := rows.Scan(&i.ID, &i.CreatedAt, &i.Version); err != nil { return nil, err } @@ -35,14 +142,15 @@ func (q *Queries) RiverMigrationDeleteByVersionMany(ctx context.Context, db DBTX return items, nil } -const riverMigrationGetAll = `-- name: RiverMigrationGetAll :many -SELECT id, created_at, version +const riverMigrationGetByLine = `-- name: RiverMigrationGetByLine :many +SELECT id, created_at, line, version FROM river_migration +WHERE line = $1 ORDER BY version ` -func (q *Queries) RiverMigrationGetAll(ctx context.Context, db DBTX) ([]*RiverMigration, error) { - rows, err := db.Query(ctx, riverMigrationGetAll) +func (q *Queries) RiverMigrationGetByLine(ctx context.Context, db DBTX, line string) ([]*RiverMigration, error) { + rows, err := db.Query(ctx, riverMigrationGetByLine, line) if err != nil { return nil, err } @@ -50,7 +158,12 @@ func (q *Queries) RiverMigrationGetAll(ctx context.Context, db DBTX) ([]*RiverMi var items []*RiverMigration for rows.Next() { var i RiverMigration - if err := rows.Scan(&i.ID, &i.CreatedAt, &i.Version); err != nil { + if err := rows.Scan( + &i.ID, + &i.CreatedAt, + &i.Line, + &i.Version, + ); err != nil { return nil, err } items = append(items, &i) @@ -63,30 +176,49 @@ func (q *Queries) RiverMigrationGetAll(ctx context.Context, db DBTX) ([]*RiverMi const riverMigrationInsert = `-- name: RiverMigrationInsert :one INSERT INTO river_migration ( + line, version ) VALUES ( - $1 -) RETURNING id, created_at, version + $1, + $2 +) RETURNING id, created_at, line, version ` -func (q *Queries) RiverMigrationInsert(ctx context.Context, db DBTX, version int64) (*RiverMigration, error) { - row := db.QueryRow(ctx, riverMigrationInsert, version) +type RiverMigrationInsertParams struct { + Line string + Version int64 +} + +func (q *Queries) RiverMigrationInsert(ctx context.Context, db DBTX, arg *RiverMigrationInsertParams) (*RiverMigration, error) { + row := db.QueryRow(ctx, riverMigrationInsert, arg.Line, arg.Version) var i RiverMigration - err := row.Scan(&i.ID, &i.CreatedAt, &i.Version) + err := row.Scan( + &i.ID, + &i.CreatedAt, + &i.Line, + &i.Version, + ) return &i, err } const riverMigrationInsertMany = `-- name: RiverMigrationInsertMany :many INSERT INTO river_migration ( + line, version ) SELECT - unnest($1::bigint[]) -RETURNING id, created_at, version + $1, + unnest($2::bigint[]) +RETURNING id, created_at, line, version ` -func (q *Queries) RiverMigrationInsertMany(ctx context.Context, db DBTX, version []int64) ([]*RiverMigration, error) { - rows, err := db.Query(ctx, riverMigrationInsertMany, version) +type RiverMigrationInsertManyParams struct { + Line string + Version []int64 +} + +func (q *Queries) RiverMigrationInsertMany(ctx context.Context, db DBTX, arg *RiverMigrationInsertManyParams) ([]*RiverMigration, error) { + rows, err := db.Query(ctx, riverMigrationInsertMany, arg.Line, arg.Version) if err != nil { return nil, err } @@ -94,6 +226,49 @@ func (q *Queries) RiverMigrationInsertMany(ctx context.Context, db DBTX, version var items []*RiverMigration for rows.Next() { var i RiverMigration + if err := rows.Scan( + &i.ID, + &i.CreatedAt, + &i.Line, + &i.Version, + ); err != nil { + return nil, err + } + items = append(items, &i) + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} + +const riverMigrationInsertManyAssumingMain = `-- name: RiverMigrationInsertManyAssumingMain :many +INSERT INTO river_migration ( + version +) +SELECT + unnest($1::bigint[]) +RETURNING + id, + created_at, + version +` + +type RiverMigrationInsertManyAssumingMainRow struct { + ID int64 + CreatedAt time.Time + Version int64 +} + +func (q *Queries) RiverMigrationInsertManyAssumingMain(ctx context.Context, db DBTX, version []int64) ([]*RiverMigrationInsertManyAssumingMainRow, error) { + rows, err := db.Query(ctx, riverMigrationInsertManyAssumingMain, version) + if err != nil { + return nil, err + } + defer rows.Close() + var items []*RiverMigrationInsertManyAssumingMainRow + for rows.Next() { + var i RiverMigrationInsertManyAssumingMainRow if err := rows.Scan(&i.ID, &i.CreatedAt, &i.Version); err != nil { return nil, err } diff --git a/riverdriver/riverpgxv5/migration/main/005_river_migration_add_line.down.sql b/riverdriver/riverpgxv5/migration/main/005_river_migration_add_line.down.sql new file mode 100644 index 00000000..1d554601 --- /dev/null +++ b/riverdriver/riverpgxv5/migration/main/005_river_migration_add_line.down.sql @@ -0,0 +1,6 @@ +DROP INDEX river_migration_line_version_idx; +CREATE UNIQUE INDEX river_migration_version_idx ON river_migration USING btree(version); + +ALTER TABLE river_migration + DROP CONSTRAINT line_length, + DROP COLUMN line; \ No newline at end of file diff --git a/riverdriver/riverpgxv5/migration/main/005_river_migration_add_line.up.sql b/riverdriver/riverpgxv5/migration/main/005_river_migration_add_line.up.sql new file mode 100644 index 00000000..aa5115a5 --- /dev/null +++ b/riverdriver/riverpgxv5/migration/main/005_river_migration_add_line.up.sql @@ -0,0 +1,12 @@ +ALTER TABLE river_migration + ADD COLUMN line text; + +UPDATE river_migration +SET line = 'main'; + +ALTER TABLE river_migration + ALTER COLUMN line SET NOT NULL, + ADD CONSTRAINT line_length CHECK (char_length(line) > 0 AND char_length(line) < 128); + +CREATE UNIQUE INDEX river_migration_line_version_idx ON river_migration USING btree(line, version); +DROP INDEX river_migration_version_idx; \ No newline at end of file diff --git a/riverdriver/riverpgxv5/river_pgx_v5_driver.go b/riverdriver/riverpgxv5/river_pgx_v5_driver.go index 47e4bae5..ec60a912 100644 --- a/riverdriver/riverpgxv5/river_pgx_v5_driver.go +++ b/riverdriver/riverpgxv5/river_pgx_v5_driver.go @@ -83,6 +83,14 @@ func (e *Executor) Begin(ctx context.Context) (riverdriver.ExecutorTx, error) { return &ExecutorTx{Executor: Executor{tx, e.queries}, tx: tx}, nil } +func (e *Executor) ColumnExists(ctx context.Context, tableName, columnName string) (bool, error) { + exists, err := e.queries.ColumnExists(ctx, e.dbtx, &dbsqlc.ColumnExistsParams{ + ColumnName: columnName, + TableName: tableName, + }) + return exists, interpretError(err) +} + func (e *Executor) Exec(ctx context.Context, sql string) (struct{}, error) { _, err := e.dbtx.Exec(ctx, sql) return struct{}{}, interpretError(err) @@ -459,32 +467,84 @@ func (e *Executor) LeaderResign(ctx context.Context, params *riverdriver.LeaderR return numResigned > 0, nil } -func (e *Executor) MigrationDeleteByVersionMany(ctx context.Context, versions []int) ([]*riverdriver.Migration, error) { - migrations, err := e.queries.RiverMigrationDeleteByVersionMany(ctx, e.dbtx, +func (e *Executor) MigrationDeleteAssumingMainMany(ctx context.Context, versions []int) ([]*riverdriver.Migration, error) { + migrations, err := e.queries.RiverMigrationDeleteAssumingMainMany(ctx, e.dbtx, sliceutil.Map(versions, func(v int) int64 { return int64(v) })) if err != nil { return nil, interpretError(err) } + return sliceutil.Map(migrations, func(internal *dbsqlc.RiverMigrationDeleteAssumingMainManyRow) *riverdriver.Migration { + return &riverdriver.Migration{ + ID: int(internal.ID), + CreatedAt: internal.CreatedAt.UTC(), + Line: riverdriver.MigrationLineMain, + Version: int(internal.Version), + } + }), nil +} + +func (e *Executor) MigrationDeleteByLineAndVersionMany(ctx context.Context, line string, versions []int) ([]*riverdriver.Migration, error) { + migrations, err := e.queries.RiverMigrationDeleteByLineAndVersionMany(ctx, e.dbtx, &dbsqlc.RiverMigrationDeleteByLineAndVersionManyParams{ + Line: line, + Version: sliceutil.Map(versions, func(v int) int64 { return int64(v) }), + }) + if err != nil { + return nil, interpretError(err) + } return sliceutil.Map(migrations, migrationFromInternal), nil } -func (e *Executor) MigrationGetAll(ctx context.Context) ([]*riverdriver.Migration, error) { - migrations, err := e.queries.RiverMigrationGetAll(ctx, e.dbtx) +func (e *Executor) MigrationGetAllAssumingMain(ctx context.Context) ([]*riverdriver.Migration, error) { + migrations, err := e.queries.RiverMigrationGetAllAssumingMain(ctx, e.dbtx) + if err != nil { + return nil, interpretError(err) + } + return sliceutil.Map(migrations, func(internal *dbsqlc.RiverMigrationGetAllAssumingMainRow) *riverdriver.Migration { + return &riverdriver.Migration{ + ID: int(internal.ID), + CreatedAt: internal.CreatedAt.UTC(), + Line: riverdriver.MigrationLineMain, + Version: int(internal.Version), + } + }), nil +} + +func (e *Executor) MigrationGetByLine(ctx context.Context, line string) ([]*riverdriver.Migration, error) { + migrations, err := e.queries.RiverMigrationGetByLine(ctx, e.dbtx, line) if err != nil { return nil, interpretError(err) } return sliceutil.Map(migrations, migrationFromInternal), nil } -func (e *Executor) MigrationInsertMany(ctx context.Context, versions []int) ([]*riverdriver.Migration, error) { - migrations, err := e.queries.RiverMigrationInsertMany(ctx, e.dbtx, - sliceutil.Map(versions, func(v int) int64 { return int64(v) })) +func (e *Executor) MigrationInsertMany(ctx context.Context, line string, versions []int) ([]*riverdriver.Migration, error) { + migrations, err := e.queries.RiverMigrationInsertMany(ctx, e.dbtx, &dbsqlc.RiverMigrationInsertManyParams{ + Line: line, + Version: sliceutil.Map(versions, func(v int) int64 { return int64(v) }), + }) if err != nil { return nil, interpretError(err) } return sliceutil.Map(migrations, migrationFromInternal), nil } +func (e *Executor) MigrationInsertManyAssumingMain(ctx context.Context, versions []int) ([]*riverdriver.Migration, error) { + migrations, err := e.queries.RiverMigrationInsertManyAssumingMain(ctx, e.dbtx, + sliceutil.Map(versions, func(v int) int64 { return int64(v) }), + ) + if err != nil { + return nil, interpretError(err) + } + return sliceutil.Map(migrations, func(internal *dbsqlc.RiverMigrationInsertManyAssumingMainRow) *riverdriver.Migration { + return &riverdriver.Migration{ + ID: int(internal.ID), + CreatedAt: internal.CreatedAt.UTC(), + Line: riverdriver.MigrationLineMain, + Version: int(internal.Version), + } + }), nil +} + func (e *Executor) NotifyMany(ctx context.Context, params *riverdriver.NotifyManyParams) error { return e.queries.PGNotifyMany(ctx, e.dbtx, &dbsqlc.PGNotifyManyParams{ Payload: params.Payload, @@ -759,6 +819,7 @@ func migrationFromInternal(internal *dbsqlc.RiverMigration) *riverdriver.Migrati return &riverdriver.Migration{ ID: int(internal.ID), CreatedAt: internal.CreatedAt.UTC(), + Line: internal.Line, Version: int(internal.Version), } } diff --git a/rivermigrate/migration/alternate/001_premier.down.sql b/rivermigrate/migration/alternate/001_premier.down.sql index e69de29b..e0ac49d1 100644 --- a/rivermigrate/migration/alternate/001_premier.down.sql +++ b/rivermigrate/migration/alternate/001_premier.down.sql @@ -0,0 +1 @@ +SELECT 1; diff --git a/rivermigrate/migration/alternate/001_premier.up.sql b/rivermigrate/migration/alternate/001_premier.up.sql index e69de29b..e0ac49d1 100644 --- a/rivermigrate/migration/alternate/001_premier.up.sql +++ b/rivermigrate/migration/alternate/001_premier.up.sql @@ -0,0 +1 @@ +SELECT 1; diff --git a/rivermigrate/migration/alternate/002_deuxieme.down.sql b/rivermigrate/migration/alternate/002_deuxieme.down.sql index e69de29b..e0ac49d1 100644 --- a/rivermigrate/migration/alternate/002_deuxieme.down.sql +++ b/rivermigrate/migration/alternate/002_deuxieme.down.sql @@ -0,0 +1 @@ +SELECT 1; diff --git a/rivermigrate/migration/alternate/002_deuxieme.up.sql b/rivermigrate/migration/alternate/002_deuxieme.up.sql index e69de29b..e0ac49d1 100644 --- a/rivermigrate/migration/alternate/002_deuxieme.up.sql +++ b/rivermigrate/migration/alternate/002_deuxieme.up.sql @@ -0,0 +1 @@ +SELECT 1; diff --git a/rivermigrate/migration/alternate/003_troisieme.down.sql b/rivermigrate/migration/alternate/003_troisieme.down.sql new file mode 100644 index 00000000..e0ac49d1 --- /dev/null +++ b/rivermigrate/migration/alternate/003_troisieme.down.sql @@ -0,0 +1 @@ +SELECT 1; diff --git a/rivermigrate/migration/alternate/003_troisieme.up.sql b/rivermigrate/migration/alternate/003_troisieme.up.sql new file mode 100644 index 00000000..e0ac49d1 --- /dev/null +++ b/rivermigrate/migration/alternate/003_troisieme.up.sql @@ -0,0 +1 @@ +SELECT 1; diff --git a/rivermigrate/migration/alternate/004_quatrieme.down.sql b/rivermigrate/migration/alternate/004_quatrieme.down.sql new file mode 100644 index 00000000..e0ac49d1 --- /dev/null +++ b/rivermigrate/migration/alternate/004_quatrieme.down.sql @@ -0,0 +1 @@ +SELECT 1; diff --git a/rivermigrate/migration/alternate/004_quatrieme.up.sql b/rivermigrate/migration/alternate/004_quatrieme.up.sql new file mode 100644 index 00000000..e0ac49d1 --- /dev/null +++ b/rivermigrate/migration/alternate/004_quatrieme.up.sql @@ -0,0 +1 @@ +SELECT 1; diff --git a/rivermigrate/migration/alternate/005_cinquieme.down.sql b/rivermigrate/migration/alternate/005_cinquieme.down.sql new file mode 100644 index 00000000..e0ac49d1 --- /dev/null +++ b/rivermigrate/migration/alternate/005_cinquieme.down.sql @@ -0,0 +1 @@ +SELECT 1; diff --git a/rivermigrate/migration/alternate/005_cinquieme.up.sql b/rivermigrate/migration/alternate/005_cinquieme.up.sql new file mode 100644 index 00000000..e0ac49d1 --- /dev/null +++ b/rivermigrate/migration/alternate/005_cinquieme.up.sql @@ -0,0 +1 @@ +SELECT 1; diff --git a/rivermigrate/migration/alternate/006_sixieme.down.sql b/rivermigrate/migration/alternate/006_sixieme.down.sql new file mode 100644 index 00000000..e0ac49d1 --- /dev/null +++ b/rivermigrate/migration/alternate/006_sixieme.down.sql @@ -0,0 +1 @@ +SELECT 1; diff --git a/rivermigrate/migration/alternate/006_sixieme.up.sql b/rivermigrate/migration/alternate/006_sixieme.up.sql new file mode 100644 index 00000000..e0ac49d1 --- /dev/null +++ b/rivermigrate/migration/alternate/006_sixieme.up.sql @@ -0,0 +1 @@ +SELECT 1; diff --git a/rivermigrate/migration/main/001_first.down.sql b/rivermigrate/migration/main/001_first.down.sql index e69de29b..e0ac49d1 100644 --- a/rivermigrate/migration/main/001_first.down.sql +++ b/rivermigrate/migration/main/001_first.down.sql @@ -0,0 +1 @@ +SELECT 1; diff --git a/rivermigrate/migration/main/001_first.up.sql b/rivermigrate/migration/main/001_first.up.sql index e69de29b..e0ac49d1 100644 --- a/rivermigrate/migration/main/001_first.up.sql +++ b/rivermigrate/migration/main/001_first.up.sql @@ -0,0 +1 @@ +SELECT 1; diff --git a/rivermigrate/migration/main/002_second.down.sql b/rivermigrate/migration/main/002_second.down.sql index e69de29b..e0ac49d1 100644 --- a/rivermigrate/migration/main/002_second.down.sql +++ b/rivermigrate/migration/main/002_second.down.sql @@ -0,0 +1 @@ +SELECT 1; diff --git a/rivermigrate/migration/main/002_second.up.sql b/rivermigrate/migration/main/002_second.up.sql index e69de29b..e0ac49d1 100644 --- a/rivermigrate/migration/main/002_second.up.sql +++ b/rivermigrate/migration/main/002_second.up.sql @@ -0,0 +1 @@ +SELECT 1; diff --git a/rivermigrate/river_migrate.go b/rivermigrate/river_migrate.go index 94c9fc3a..6f2eb814 100644 --- a/rivermigrate/river_migrate.go +++ b/rivermigrate/river_migrate.go @@ -4,6 +4,7 @@ package rivermigrate import ( "context" + "errors" "fmt" "io" "io/fs" @@ -25,6 +26,11 @@ import ( "github.com/riverqueue/river/rivershared/util/valutil" ) +// The migrate version where the `line` column was added. Meaningful in that the +// migrator has to behave a little differently depending on whether it's working +// with versions before or after this boundary. +const migrateVersionLineColumnAdded = 5 + // Migration is a bundled migration containing a version (e.g. 1, 2, 3), and SQL // for up and down directions. type Migration struct { @@ -324,9 +330,21 @@ func (m *Migrator[TTx]) migrateDown(ctx context.Context, exec riverdriver.Execut return res, nil } - if !opts.DryRun { - if _, err := exec.MigrationDeleteByVersionMany(ctx, sliceutil.Map(res.Versions, migrateVersionToInt)); err != nil { - return nil, fmt.Errorf("error deleting migration rows for versions %+v: %w", res.Versions, err) + if !opts.DryRun && len(res.Versions) > 0 { + versions := sliceutil.Map(res.Versions, migrateVersionToInt) + + // Version 005 is hard-coded here because that's the version in which + // the migration `line` comes in. If migration to a point equal or above + // 005, we can remove migrations with a line included, but otherwise we + // must omit the `line` column from queries because it doesn't exist. + if m.line == riverdriver.MigrationLineMain && slices.Min(versions) <= migrateVersionLineColumnAdded { + if _, err := exec.MigrationDeleteAssumingMainMany(ctx, versions); err != nil { + return nil, fmt.Errorf("error inserting migration rows for versions %+v assuming main: %w", res.Versions, err) + } + } else { + if _, err := exec.MigrationDeleteByLineAndVersionMany(ctx, m.line, versions); err != nil { + return nil, fmt.Errorf("error deleting migration rows for versions %+v on line %q: %w", res.Versions, m.line, err) + } } } @@ -353,9 +371,24 @@ func (m *Migrator[TTx]) migrateUp(ctx context.Context, exec riverdriver.Executor return nil, err } - if opts == nil || !opts.DryRun { - if _, err := exec.MigrationInsertMany(ctx, sliceutil.Map(res.Versions, migrateVersionToInt)); err != nil { - return nil, fmt.Errorf("error inserting migration rows for versions %+v: %w", res.Versions, err) + if (opts == nil || !opts.DryRun) && len(res.Versions) > 0 { + versions := sliceutil.Map(res.Versions, migrateVersionToInt) + + // Version 005 is hard-coded here because that's the version in which + // the migration `line` comes in. If migration to a point equal or above + // 005, we can insert migrations with a line included, but otherwise we + // must omit the `line` column from queries because it doesn't exist. + if m.line == riverdriver.MigrationLineMain && slices.Max(versions) < migrateVersionLineColumnAdded { + if _, err := exec.MigrationInsertManyAssumingMain(ctx, versions); err != nil { + return nil, fmt.Errorf("error inserting migration rows for versions %+v assuming main: %w", res.Versions, err) + } + } else { + if _, err := exec.MigrationInsertMany(ctx, + m.line, + versions, + ); err != nil { + return nil, fmt.Errorf("error inserting migration rows for versions %+v on line %q: %w", res.Versions, m.line, err) + } } } @@ -480,17 +513,39 @@ func (m *Migrator[TTx]) applyMigrations(ctx context.Context, exec riverdriver.Ex // because otherwise the existing transaction would become aborted on an // unsuccessful `river_migration` check.) func (m *Migrator[TTx]) existingMigrations(ctx context.Context, exec riverdriver.Executor) ([]*riverdriver.Migration, error) { - exists, err := exec.TableExists(ctx, "river_migration") + migrateTableExists, err := exec.TableExists(ctx, "river_migration") if err != nil { return nil, fmt.Errorf("error checking if `%s` exists: %w", "river_migration", err) } - if !exists { + if !migrateTableExists { + if m.line != riverdriver.MigrationLineMain { + return nil, errors.New("can't add a non-main migration line until `river_migration` is raised; fully migrate the main migration line and try again") + } + return nil, nil } - migrations, err := exec.MigrationGetAll(ctx) + lineColumnExists, err := exec.ColumnExists(ctx, "river_migration", "line") + if err != nil { + return nil, fmt.Errorf("error checking if `%s.%s` exists: %w", "river_migration", "line", err) + } + + if !lineColumnExists { + if m.line != riverdriver.MigrationLineMain { + return nil, errors.New("can't add a non-main migration line until `river_migration.line` is raised; fully migrate the main migration line and try again") + } + + migrations, err := exec.MigrationGetAllAssumingMain(ctx) + if err != nil { + return nil, fmt.Errorf("error getting existing migrations: %w", err) + } + + return migrations, nil + } + + migrations, err := exec.MigrationGetByLine(ctx, m.line) if err != nil { - return nil, fmt.Errorf("error getting existing migrations: %w", err) + return nil, fmt.Errorf("error getting existing migrations for line %q: %w", m.line, err) } return migrations, nil diff --git a/rivermigrate/river_migrate_test.go b/rivermigrate/river_migrate_test.go index 17298138..5c15b89c 100644 --- a/rivermigrate/river_migrate_test.go +++ b/rivermigrate/river_migrate_test.go @@ -5,6 +5,7 @@ import ( "database/sql" "embed" "fmt" + "io/fs" "log/slog" "os" "slices" @@ -26,6 +27,35 @@ import ( "github.com/riverqueue/river/rivershared/util/sliceutil" ) +const ( + // The name of an actual migration line embedded in our test data below. + migrationLineAlternate = "alternate" + migrationLineAlternateMaxVersion = 6 +) + +//go:embed migration/*/*.sql +var migrationFS embed.FS + +// A test driver with the same migrations as the standard Pgx driver, but which +// includes an alternate line so we can test that those work. +type driverWithAlternateLine struct { + *riverpgxv5.Driver +} + +func (d *driverWithAlternateLine) GetMigrationFS(line string) fs.FS { + switch line { + case riverdriver.MigrationLineMain: + return d.Driver.GetMigrationFS(line) + case migrationLineAlternate: + return migrationFS + } + panic("migration line does not exist: " + line) +} + +func (d *driverWithAlternateLine) GetMigrationLines() []string { + return append(d.Driver.GetMigrationLines(), migrationLineAlternate) +} + func TestMigrator(t *testing.T) { t.Parallel() @@ -35,7 +65,7 @@ func TestMigrator(t *testing.T) { type testBundle struct { dbPool *pgxpool.Pool - driver *riverpgxv5.Driver + driver *driverWithAlternateLine logger *slog.Logger tx pgx.Tx } @@ -59,7 +89,7 @@ func TestMigrator(t *testing.T) { bundle := &testBundle{ dbPool: dbPool, - driver: riverpgxv5.New(dbPool), + driver: &driverWithAlternateLine{Driver: riverpgxv5.New(dbPool)}, logger: riversharedtest.Logger(t), tx: tx, } @@ -102,13 +132,16 @@ func TestMigrator(t *testing.T) { migrator, bundle := setup(t) + // The migration version in which `river_job` comes in. + const migrateVersionIncludingRiverJob = 2 + // Run two initial times to get to the version before river_job is dropped. // Defaults to only running one step when moving in the down direction. - for i := 0; i < 2; i++ { + for i := migrationsBundle.MaxVersion; i > migrateVersionIncludingRiverJob; i-- { res, err := migrator.MigrateTx(ctx, bundle.tx, DirectionDown, &MigrateOpts{}) require.NoError(t, err) require.Equal(t, DirectionDown, res.Direction) - require.Equal(t, []int{4 - i}, sliceutil.Map(res.Versions, migrateVersionToInt)) + require.Equal(t, []int{i}, sliceutil.Map(res.Versions, migrateVersionToInt)) err = dbExecError(ctx, bundle.driver.UnwrapExecutor(bundle.tx), "SELECT * FROM river_job") require.NoError(t, err) @@ -119,7 +152,7 @@ func TestMigrator(t *testing.T) { res, err := migrator.MigrateTx(ctx, bundle.tx, DirectionDown, &MigrateOpts{}) require.NoError(t, err) require.Equal(t, DirectionDown, res.Direction) - require.Equal(t, []int{2}, sliceutil.Map(res.Versions, migrateVersionToInt)) + require.Equal(t, []int{migrateVersionIncludingRiverJob}, sliceutil.Map(res.Versions, migrateVersionToInt)) err = dbExecError(ctx, bundle.driver.UnwrapExecutor(bundle.tx), "SELECT * FROM river_job") require.Error(t, err) @@ -152,7 +185,7 @@ func TestMigrator(t *testing.T) { require.Equal(t, []int{migrationsBundle.WithTestVersionsMaxVersion, migrationsBundle.WithTestVersionsMaxVersion - 1}, sliceutil.Map(res.Versions, migrateVersionToInt)) - migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetAll(ctx) + migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetByLine(ctx, riverdriver.MigrationLineMain) require.NoError(t, err) require.Equal(t, seqOneTo(migrationsBundle.WithTestVersionsMaxVersion-2), sliceutil.Map(migrations, driverMigrationToInt)) @@ -173,9 +206,9 @@ func TestMigrator(t *testing.T) { require.NoError(t, err) require.Equal(t, []int{}, sliceutil.Map(res.Versions, migrateVersionToInt)) - migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetAll(ctx) + migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetByLine(ctx, riverdriver.MigrationLineMain) require.NoError(t, err) - require.Equal(t, seqOneTo(4), + require.Equal(t, seqOneTo(migrationsBundle.MaxVersion), sliceutil.Map(migrations, driverMigrationToInt)) }) @@ -187,11 +220,11 @@ func TestMigrator(t *testing.T) { res, err := migrator.MigrateTx(ctx, tx, DirectionDown, &MigrateOpts{MaxSteps: 1}) require.NoError(t, err) - require.Equal(t, []int{4}, sliceutil.Map(res.Versions, migrateVersionToInt)) + require.Equal(t, []int{migrationsBundle.MaxVersion}, sliceutil.Map(res.Versions, migrateVersionToInt)) - migrations, err := migrator.driver.UnwrapExecutor(tx).MigrationGetAll(ctx) + migrations, err := migrator.driver.UnwrapExecutor(tx).MigrationGetAllAssumingMain(ctx) require.NoError(t, err) - require.Equal(t, seqOneTo(3), + require.Equal(t, seqOneTo(migrationsBundle.MaxVersion-1), sliceutil.Map(migrations, driverMigrationToInt)) }) @@ -205,10 +238,10 @@ func TestMigrator(t *testing.T) { res, err := migrator.MigrateTx(ctx, bundle.tx, DirectionDown, &MigrateOpts{TargetVersion: 4}) require.NoError(t, err) - require.Equal(t, []int{6, 5}, + require.Equal(t, seqDownTo(migrationsBundle.WithTestVersionsMaxVersion, 5), sliceutil.Map(res.Versions, migrateVersionToInt)) - migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetAll(ctx) + migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetAllAssumingMain(ctx) require.NoError(t, err) require.Equal(t, seqOneTo(4), sliceutil.Map(migrations, driverMigrationToInt)) @@ -227,7 +260,7 @@ func TestMigrator(t *testing.T) { res, err := migrator.MigrateTx(ctx, bundle.tx, DirectionDown, &MigrateOpts{TargetVersion: -1}) require.NoError(t, err) - require.Equal(t, seqToOne(6), + require.Equal(t, seqDownTo(migrationsBundle.WithTestVersionsMaxVersion, 1), sliceutil.Map(res.Versions, migrateVersionToInt)) err = dbExecError(ctx, bundle.driver.UnwrapExecutor(bundle.tx), "SELECT name FROM river_migrate") @@ -241,14 +274,14 @@ func TestMigrator(t *testing.T) { // migration doesn't exist { - _, err := migrator.MigrateTx(ctx, bundle.tx, DirectionDown, &MigrateOpts{TargetVersion: 77}) - require.EqualError(t, err, "version 77 is not a valid River migration version") + _, err := migrator.MigrateTx(ctx, bundle.tx, DirectionDown, &MigrateOpts{TargetVersion: migrationsBundle.MaxVersion + 77}) + require.EqualError(t, err, fmt.Sprintf("version %d is not a valid River migration version", migrationsBundle.MaxVersion+77)) } // migration exists but not one that's applied { - _, err := migrator.MigrateTx(ctx, bundle.tx, DirectionDown, &MigrateOpts{TargetVersion: 5}) - require.EqualError(t, err, "version 5 is not in target list of valid migrations to apply") + _, err := migrator.MigrateTx(ctx, bundle.tx, DirectionDown, &MigrateOpts{TargetVersion: migrationsBundle.MaxVersion + 1}) + require.EqualError(t, err, fmt.Sprintf("version %d is not in target list of valid migrations to apply", migrationsBundle.MaxVersion+1)) } }) @@ -267,7 +300,7 @@ func TestMigrator(t *testing.T) { // Migrate down returned a result above for a migration that was // removed, but because we're in a dry run, the database still shows // this version. - migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetAll(ctx) + migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetByLine(ctx, riverdriver.MigrationLineMain) require.NoError(t, err) require.Equal(t, seqOneTo(migrationsBundle.WithTestVersionsMaxVersion), sliceutil.Map(migrations, driverMigrationToInt)) @@ -298,7 +331,7 @@ func TestMigrator(t *testing.T) { res, err := migrator.MigrateTx(ctx, bundle.tx, DirectionUp, nil) require.NoError(t, err) - require.Equal(t, []int{5, 6}, sliceutil.Map(res.Versions, migrateVersionToInt)) + require.Equal(t, []int{migrationsBundle.MaxVersion + 1, migrationsBundle.MaxVersion + 2}, sliceutil.Map(res.Versions, migrateVersionToInt)) }) t.Run("MigrateUpDefault", func(t *testing.T) { @@ -314,7 +347,7 @@ func TestMigrator(t *testing.T) { require.Equal(t, []int{migrationsBundle.WithTestVersionsMaxVersion - 1, migrationsBundle.WithTestVersionsMaxVersion}, sliceutil.Map(res.Versions, migrateVersionToInt)) - migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetAll(ctx) + migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetByLine(ctx, riverdriver.MigrationLineMain) require.NoError(t, err) require.Equal(t, seqOneTo(migrationsBundle.WithTestVersionsMaxVersion), sliceutil.Map(migrations, driverMigrationToInt)) @@ -330,7 +363,7 @@ func TestMigrator(t *testing.T) { require.Equal(t, DirectionUp, res.Direction) require.Equal(t, []int{}, sliceutil.Map(res.Versions, migrateVersionToInt)) - migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetAll(ctx) + migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetByLine(ctx, riverdriver.MigrationLineMain) require.NoError(t, err) require.Equal(t, seqOneTo(migrationsBundle.WithTestVersionsMaxVersion), sliceutil.Map(migrations, driverMigrationToInt)) @@ -350,7 +383,7 @@ func TestMigrator(t *testing.T) { require.Equal(t, []int{migrationsBundle.WithTestVersionsMaxVersion - 1}, sliceutil.Map(res.Versions, migrateVersionToInt)) - migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetAll(ctx) + migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetByLine(ctx, riverdriver.MigrationLineMain) require.NoError(t, err) require.Equal(t, seqOneTo(migrationsBundle.WithTestVersionsMaxVersion-1), sliceutil.Map(migrations, driverMigrationToInt)) @@ -376,9 +409,9 @@ func TestMigrator(t *testing.T) { require.NoError(t, err) require.Equal(t, []int{}, sliceutil.Map(res.Versions, migrateVersionToInt)) - migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetAll(ctx) + migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetByLine(ctx, riverdriver.MigrationLineMain) require.NoError(t, err) - require.Equal(t, seqOneTo(4), + require.Equal(t, seqOneTo(migrationsBundle.MaxVersion), sliceutil.Map(migrations, driverMigrationToInt)) }) @@ -392,7 +425,7 @@ func TestMigrator(t *testing.T) { require.NoError(t, err) require.Equal(t, []int{migrationsBundle.MaxVersion + 1}, sliceutil.Map(res.Versions, migrateVersionToInt)) - migrations, err := migrator.driver.UnwrapExecutor(tx).MigrationGetAll(ctx) + migrations, err := migrator.driver.UnwrapExecutor(tx).MigrationGetByLine(ctx, riverdriver.MigrationLineMain) require.NoError(t, err) require.Equal(t, seqOneTo(migrationsBundle.MaxVersion+1), sliceutil.Map(migrations, driverMigrationToInt)) @@ -403,14 +436,14 @@ func TestMigrator(t *testing.T) { migrator, bundle := setup(t) - res, err := migrator.MigrateTx(ctx, bundle.tx, DirectionUp, &MigrateOpts{TargetVersion: 6}) + res, err := migrator.MigrateTx(ctx, bundle.tx, DirectionUp, &MigrateOpts{TargetVersion: migrationsBundle.MaxVersion + 2}) require.NoError(t, err) - require.Equal(t, []int{5, 6}, + require.Equal(t, []int{migrationsBundle.MaxVersion + 1, migrationsBundle.MaxVersion + 2}, sliceutil.Map(res.Versions, migrateVersionToInt)) - migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetAll(ctx) + migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetByLine(ctx, riverdriver.MigrationLineMain) require.NoError(t, err) - require.Equal(t, seqOneTo(6), sliceutil.Map(migrations, driverMigrationToInt)) + require.Equal(t, seqOneTo(migrationsBundle.MaxVersion+2), sliceutil.Map(migrations, driverMigrationToInt)) }) t.Run("MigrateUpWithTargetVersionInvalid", func(t *testing.T) { @@ -445,7 +478,7 @@ func TestMigrator(t *testing.T) { // Migrate up returned a result above for migrations that were applied, // but because we're in a dry run, the database still shows the test // migration versions not applied. - migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetAll(ctx) + migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetByLine(ctx, riverdriver.MigrationLineMain) require.NoError(t, err) require.Equal(t, seqOneTo(migrationsBundle.MaxVersion), sliceutil.Map(migrations, driverMigrationToInt)) @@ -476,10 +509,99 @@ func TestMigrator(t *testing.T) { Messages: []string{fmt.Sprintf("Unapplied migrations: [%d %d]", migrationsBundle.MaxVersion+1, migrationsBundle.MaxVersion+2)}, }, res) }) -} -//go:embed migration/*/*.sql -var migrationFS embed.FS + t.Run("MigrateDownToZeroAndBackUp", func(t *testing.T) { + t.Parallel() + + migrator, bundle := setup(t) + + requireMigrationTableExists := func(expectedExists bool) { + migrationExists, err := bundle.driver.UnwrapExecutor(bundle.tx).TableExists(ctx, "river_migration") + require.NoError(t, err) + require.Equal(t, expectedExists, migrationExists) + } + + requireMigrationTableExists(true) + + res, err := migrator.MigrateTx(ctx, bundle.tx, DirectionDown, &MigrateOpts{TargetVersion: -1}) + require.NoError(t, err) + require.Equal(t, seqDownTo(migrationsBundle.MaxVersion, 1), + sliceutil.Map(res.Versions, migrateVersionToInt)) + + requireMigrationTableExists(false) + + res, err = migrator.MigrateTx(ctx, bundle.tx, DirectionUp, &MigrateOpts{}) + require.NoError(t, err) + require.Equal(t, seqOneTo(migrationsBundle.WithTestVersionsMaxVersion), + sliceutil.Map(res.Versions, migrateVersionToInt)) + + migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetByLine(ctx, riverdriver.MigrationLineMain) + require.NoError(t, err) + require.Equal(t, seqOneTo(migrationsBundle.WithTestVersionsMaxVersion), + sliceutil.Map(migrations, driverMigrationToInt)) + }) + + t.Run("AlternateLineUpAndDown", func(t *testing.T) { + t.Parallel() + + _, bundle := setup(t) + + // We have to reinitialize the alternateMigrator because the migrations bundle is + // set in the constructor. + alternateMigrator := New(bundle.driver, &Config{ + Line: migrationLineAlternate, + Logger: bundle.logger, + }) + + res, err := alternateMigrator.MigrateTx(ctx, bundle.tx, DirectionUp, &MigrateOpts{}) + require.NoError(t, err) + require.Equal(t, seqOneTo(migrationLineAlternateMaxVersion), + sliceutil.Map(res.Versions, migrateVersionToInt)) + + migrations, err := bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetByLine(ctx, migrationLineAlternate) + require.NoError(t, err) + require.Equal(t, seqOneTo(migrationLineAlternateMaxVersion), + sliceutil.Map(migrations, driverMigrationToInt)) + + res, err = alternateMigrator.MigrateTx(ctx, bundle.tx, DirectionDown, &MigrateOpts{TargetVersion: -1}) + require.NoError(t, err) + require.Equal(t, seqDownTo(migrationLineAlternateMaxVersion, 1), + sliceutil.Map(res.Versions, migrateVersionToInt)) + + // The main migration line should not have been touched. + migrations, err = bundle.driver.UnwrapExecutor(bundle.tx).MigrationGetByLine(ctx, riverdriver.MigrationLineMain) + require.NoError(t, err) + require.Equal(t, seqOneTo(migrationsBundle.MaxVersion), + sliceutil.Map(migrations, driverMigrationToInt)) + }) + + t.Run("AlternateLineBeforeLineColumn", func(t *testing.T) { + t.Parallel() + + migrator, bundle := setup(t) + + // Main line to just before the `line` column was added. + _, err := migrator.MigrateTx(ctx, bundle.tx, DirectionDown, &MigrateOpts{TargetVersion: 4}) + require.NoError(t, err) + + alternateMigrator := New(bundle.driver, &Config{ + Line: migrationLineAlternate, + Logger: bundle.logger, + }) + + // Alternate line not allowed because `river_job.line` doesn't exist. + _, err = alternateMigrator.MigrateTx(ctx, bundle.tx, DirectionUp, &MigrateOpts{}) + require.EqualError(t, err, "can't add a non-main migration line until `river_migration.line` is raised; fully migrate the main migration line and try again") + + // Main line to zero. + _, err = migrator.MigrateTx(ctx, bundle.tx, DirectionDown, &MigrateOpts{TargetVersion: -1}) + require.NoError(t, err) + + // Alternate line not allowed because `river_job` doesn't exist. + _, err = alternateMigrator.MigrateTx(ctx, bundle.tx, DirectionUp, &MigrateOpts{}) + require.EqualError(t, err, "can't add a non-main migration line until `river_migration` is raised; fully migrate the main migration line and try again") + }) +} // This test uses a custom set of test-only migration files on the file system // in `rivermigrate/migrate/*`. @@ -489,6 +611,8 @@ func TestMigrationsFromFS(t *testing.T) { t.Run("Main", func(t *testing.T) { t.Parallel() + // This is not the actual main line, but rather one embedded in this + // package's test data (see `rivermigrate/migrate/*`). migrations, err := migrationsFromFS(migrationFS, "main") require.NoError(t, err) require.Equal(t, []int{1, 2}, sliceutil.Map(migrations, migrationToInt)) @@ -497,9 +621,9 @@ func TestMigrationsFromFS(t *testing.T) { t.Run("Alternate", func(t *testing.T) { t.Parallel() - migrations, err := migrationsFromFS(migrationFS, "alternate") + migrations, err := migrationsFromFS(migrationFS, migrationLineAlternate) require.NoError(t, err) - require.Equal(t, []int{1, 2}, sliceutil.Map(migrations, migrationToInt)) + require.Equal(t, seqOneTo(migrationLineAlternateMaxVersion), sliceutil.Map(migrations, migrationToInt)) }) t.Run("DoesNotExist", func(t *testing.T) { @@ -566,18 +690,24 @@ func dbExecError(ctx context.Context, exec riverdriver.Executor, sql string) err func driverMigrationToInt(r *riverdriver.Migration) int { return r.Version } func migrationToInt(migration Migration) int { return migration.Version } +// Produces a sequence down to one. Max is included. func seqOneTo(max int) []int { - seq := make([]int, max) + seq := make([]int, 0, max) - for i := 0; i < max; i++ { - seq[i] = i + 1 + for i := 1; i <= max; i++ { + seq = append(seq, i) } return seq } -func seqToOne(max int) []int { - seq := seqOneTo(max) +func seqDownTo(max, min int) []int { + seq := make([]int, 0, max-min+1) + + for i := min; i <= max; i++ { + seq = append(seq, i) + } + slices.Reverse(seq) return seq }