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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
- Fixed maintenance startup failures leaving a client renewing leadership with maintenance stopped in poll-only mode. After exhausting startup retries, clients now request local resignation without depending on database notifications. [PR #1347](https://github.com/riverqueue/river/pull/1347).
- Fixed YugabyteDB clients relying on notifications when `LISTEN/NOTIFY` is unavailable or disabled. Clients now automatically poll for running job cancellations and queue pause, resume, and metadata changes with the default `PollOnly: false`, and skip unsupported notification broadcasts. Native notifications require YugabyteDB 2025.2.3 or later with `ysql_yb_enable_listen_notify=true` on both Masters and TServers. [PR #1347](https://github.com/riverqueue/river/pull/1347).
- Fixed the default retry policy scheduling a job's retry about 292 years in the past on amd64 once the job had errored 310 or more times, which made it run again immediately. The capped retry delay is now exactly the maximum duration on every architecture. [PR #1402](https://github.com/riverqueue/river/pull/1402).
- Fixed up migrations targeting an already-applied version to do nothing instead of applying later pending migrations. [PR #1403](https://github.com/riverqueue/river/pull/1403).
- Fixed job cancellations received during a fetch being lost before the fetched jobs started. Matching jobs now receive cancellation before work begins. [PR #1397](https://github.com/riverqueue/river/pull/1397).
- Fixed SQLite `JobCancel` and `JobCancelTx` notifying running workers through the shared control outbox, so their contexts are cancelled when the transaction commits. [PR #1398](https://github.com/riverqueue/river/pull/1398).
- Fixed SQLite `InsertMany` and `InsertManyTx` reporting multiple inserted jobs when a batch contains the same active unique key more than once. Such batches now fail atomically, matching PostgreSQL. [PR #1399](https://github.com/riverqueue/river/pull/1399).
Expand Down
41 changes: 41 additions & 0 deletions cmd/river/rivercli/river_cli_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import (
"maps"
"net/url"
"runtime/debug"
"strconv"
"strings"
"testing"
"time"
Expand All @@ -23,6 +24,7 @@ import (
"github.com/riverqueue/river/riverdriver/riverpgxv5"
"github.com/riverqueue/river/rivermigrate"
"github.com/riverqueue/river/rivershared/riversharedtest"
"github.com/riverqueue/river/rivershared/util/randutil"
)

type DriverProcurerStub struct {
Expand Down Expand Up @@ -171,6 +173,45 @@ func TestBaseCommandSetIntegration(t *testing.T) {
require.EqualError(t, cmd.Execute(), "unsupported database URL (`post://`); try one with a `postgres://`, `postgresql://`, or `sqlite://` scheme/prefix")
})

t.Run("MigrateUpWithAppliedTargetVersion", func(t *testing.T) {
t.Parallel()

ctx := context.Background()
cmd, bundle := setup(t)
dbPool := riversharedtest.DBPool(ctx, t)
schema := "river_cli_target_test_" + randutil.Hex(8)
_, err := dbPool.Exec(ctx, "CREATE SCHEMA "+schema)
require.NoError(t, err)
t.Cleanup(func() {
_, err := dbPool.Exec(ctx, "DROP SCHEMA "+schema+" CASCADE")
require.NoError(t, err)
})

migrator, err := rivermigrate.New(riverpgxv5.New(dbPool), &rivermigrate.Config{Schema: schema})
require.NoError(t, err)

// Leave the latest migration pending, then ask the CLI to migrate up
// to an older version that is already applied.
versions := migrator.AllVersions()
currentVersion := versions[len(versions)-2].Version
targetVersion := versions[len(versions)-3].Version

_, err = migrator.Migrate(ctx, rivermigrate.DirectionUp, &rivermigrate.MigrateOpts{TargetVersion: currentVersion})
require.NoError(t, err)

cmd.SetArgs([]string{
"migrate-up", "--database-url", riversharedtest.TestDatabaseURL(),
"--schema", schema, "--target-version", strconv.Itoa(targetVersion),
})
require.NoError(t, cmd.Execute())
require.Equal(t, "no migrations to apply\n", bundle.out.String())

existingVersions, err := migrator.ExistingVersions(ctx)
require.NoError(t, err)
require.Len(t, existingVersions, len(versions)-1)
require.Equal(t, currentVersion, existingVersions[len(existingVersions)-1].Version)
})

t.Run("MissingDatabaseURLAndPGEnv", func(t *testing.T) {
t.Parallel()

Expand Down
19 changes: 12 additions & 7 deletions rivermigrate/river_migrate.go
Original file line number Diff line number Diff line change
Expand Up @@ -226,9 +226,8 @@ type MigrateOpts struct {
MaxSteps int

// TargetVersion is a specific migration version to apply migrations to. The
// version must exist and it must be in the possible list of migrations to
// apply. e.g. If requesting an up migration with version 3, version 3 must
// not already be applied.
// version must exist. An up migration whose target is already applied does
// nothing, even if later migrations are pending.
//
// When applying migrations up, migrations are applied including the target
// version, so when starting at version 0 and requesting version 3, versions
Expand Down Expand Up @@ -525,6 +524,10 @@ func (m *Migrator[TTx]) validate(ctx context.Context, exec riverdriver.Executor,
// Common code shared between the up and down migration directions that walks
// through each target migration and applies it, logging appropriately.
func (m *Migrator[TTx]) applyMigrations(ctx context.Context, exec riverdriver.Executor, direction Direction, opts *MigrateOpts, inOuterTx bool, sortedTargetMigrations []Migration) (*MigrateResult, error) {
targetWasPending := slices.ContainsFunc(sortedTargetMigrations, func(migration Migration) bool {
return migration.Version == opts.TargetVersion
})

var maxSteps int
switch {
case opts.MaxSteps != 0:
Expand All @@ -547,12 +550,14 @@ func (m *Migrator[TTx]) applyMigrations(ctx context.Context, exec riverdriver.Ex

targetIndex := slices.IndexFunc(sortedTargetMigrations, func(b Migration) bool { return b.Version == opts.TargetVersion })
if targetIndex == -1 {
// Error, but only if the migration doesn't exist or was never
// applied on a down migration. Up migrations with TargetVersion
// that's already applied should fall through with a no-op.
if _, ok := m.migrations[opts.TargetVersion]; !ok || direction == DirectionDown {
if direction == DirectionDown {
return nil, fmt.Errorf("version %d is not in target list of valid migrations to apply", opts.TargetVersion)
}
// The target was already applied if it was absent before MaxSteps
// trimmed the list. Keep the trimmed list when the target is pending.
if !targetWasPending {
sortedTargetMigrations = []Migration{}
}
} else {
// Replace target list with list up to target index. Migrations are
// sorted according to the direction we're migrating in, so when down
Expand Down
60 changes: 60 additions & 0 deletions rivermigrate/river_migrate_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -540,6 +540,31 @@ func TestMigrator(t *testing.T) {
require.Equal(t, pgerrcode.UndefinedColumn, pgErr.Code)
})

t.Run("MigrateUpWithMaxStepsAndTargetVersion", func(t *testing.T) {
t.Parallel()

migrator, bundle := setup(t)

_, err := migrator.Migrate(ctx, DirectionUp, &MigrateOpts{TargetVersion: migrationsBundle.MaxVersion})
require.NoError(t, err)

// The target is pending but beyond the one-step limit, so only the
// first pending migration should be applied.
res, err := migrator.Migrate(ctx, DirectionUp, &MigrateOpts{
MaxSteps: 1,
TargetVersion: migrationsBundle.WithTestVersionsMaxVersion,
})
require.NoError(t, err)
require.Equal(t, []int{migrationsBundle.MaxVersion + 1}, sliceutil.Map(res.Versions, migrateVersionToInt))

migrations, err := bundle.driver.GetExecutor().MigrationGetByLine(ctx, &riverdriver.MigrationGetByLineParams{
Line: riverdriver.MigrationLineMain,
Schema: bundle.schema,
})
require.NoError(t, err)
require.Equal(t, seqOneTo(migrationsBundle.MaxVersion+1), sliceutil.Map(migrations, driverMigrationToInt))
})

t.Run("MigrateUpWithPool", func(t *testing.T) {
t.Parallel()

Expand Down Expand Up @@ -581,6 +606,41 @@ func TestMigrator(t *testing.T) {
require.Empty(t, res.Versions)
})

t.Run("MigrateUpWithTargetVersionAlreadyApplied", func(t *testing.T) {
t.Parallel()

for _, testCase := range []struct {
name string
targetVersion int
}{
{name: "TargetBelowCurrentVersion", targetVersion: migrationsBundle.MaxVersion - 1},
{name: "TargetEqualsCurrentVersion", targetVersion: migrationsBundle.MaxVersion},
} {
t.Run(testCase.name, func(t *testing.T) {
t.Parallel()

migrator, bundle := setup(t)

// Stop before the two test migrations, leaving newer versions
// pending for both already-applied target cases.
_, err := migrator.Migrate(ctx, DirectionUp, &MigrateOpts{TargetVersion: migrationsBundle.MaxVersion})
require.NoError(t, err)

res, err := migrator.Migrate(ctx, DirectionUp, &MigrateOpts{TargetVersion: testCase.targetVersion})
require.NoError(t, err)
require.Equal(t, DirectionUp, res.Direction)
require.Empty(t, res.Versions)

migrations, err := bundle.driver.GetExecutor().MigrationGetByLine(ctx, &riverdriver.MigrationGetByLineParams{
Line: riverdriver.MigrationLineMain,
Schema: bundle.schema,
})
require.NoError(t, err)
require.Equal(t, seqOneTo(migrationsBundle.MaxVersion), sliceutil.Map(migrations, driverMigrationToInt))
})
}
})

t.Run("MigrateUpWithTargetVersionInvalid", func(t *testing.T) {
t.Parallel()

Expand Down
Loading