Skip to content
Open
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
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Fixed

- Fixed SQLite job list pagination skipping or repeating jobs by formatting cursor timestamps consistently with stored timestamps. [PR #1374](https://github.com/riverqueue/river/pull/1374).
- Improved PostgreSQL job listing performance when filtering by one finalized state (`completed`, `cancelled`, or `discarded`) and sorting by finalized time, including in River UI. [PR #1374](https://github.com/riverqueue/river/pull/1374).

## [0.47.0] - 2026-09-01

### Added
Expand Down
20 changes: 20 additions & 0 deletions delete_many_params_test.go
Original file line number Diff line number Diff line change
@@ -1,10 +1,13 @@
package river

import (
"context"
"testing"

"github.com/stretchr/testify/require"

"github.com/riverqueue/river/internal/dblist"
"github.com/riverqueue/river/riverdriver/riverpgxv5"
"github.com/riverqueue/river/rivertype"
)

Expand All @@ -29,3 +32,20 @@ func TestJobDeleteManyParams_UnsafeAll(t *testing.T) {
NewJobDeleteManyParams().IDs(123).UnsafeAll()
})
}

func TestJobDeleteManyParams_toDBParams(t *testing.T) {
t.Parallel()

for _, state := range []rivertype.JobState{rivertype.JobStateAvailable, rivertype.JobStateCancelled, rivertype.JobStateCompleted, rivertype.JobStateDiscarded} {
t.Run(string(state), func(t *testing.T) {
t.Parallel()

params := NewJobDeleteManyParams().States(state).toDBParams()
require.Empty(t, params.Where)
driverParams, err := dblist.JobMakeDriverParams(context.Background(), params, riverpgxv5.New(nil))
require.NoError(t, err)
require.Equal(t, "state = any(@state)", driverParams.WhereClause)
require.Equal(t, []string{string(state)}, driverParams.NamedArgs["state"])
})
}
}
43 changes: 37 additions & 6 deletions job_list_params.go
Original file line number Diff line number Diff line change
Expand Up @@ -270,20 +270,51 @@ func (p *JobListParams) toDBParams() (*dblist.JobListParams, error) {

orderBy = append(orderBy, dblist.JobListOrderBy{Expr: "id", Order: sortOrder})

// Preserve custom SQL and its argument types without trying to parse it.
// In particular, an ungrouped OR may bypass the typed state filter, and
// custom SQL can reference the existing @state array argument. Metadata
// predicates also live in p.where; conservatively keep that path unchanged.
states := p.states

// Copy conditions so reusing params does not accumulate generated cursor
// predicates or mix them into the caller's custom SQL.
where := append([]dblist.WherePredicate(nil), p.where...)
if len(p.where) == 0 && len(states) == 1 {
// Equality lets Postgres use the timestamp ordering of an index on
// (state, finalized_at). ANY does not establish that state is fixed.
where = append(where, dblist.WherePredicate{
NamedArgs: map[string]any{"state": string(states[0])},
SQL: "state = @state",
})

// Supported schemas enforce non-null finalized_at for these states.
// Make that explicit so Postgres can use the existing partial index.
if timeField == "finalized_at" {
switch states[0] {
case rivertype.JobStateCancelled, rivertype.JobStateCompleted, rivertype.JobStateDiscarded:
where = append(where, dblist.WherePredicate{SQL: "finalized_at IS NOT NULL"})
case rivertype.JobStateAvailable, rivertype.JobStatePending, rivertype.JobStateRetryable, rivertype.JobStateRunning, rivertype.JobStateScheduled:
}
}

// The state filter is already represented in where.
states = nil
}

if p.after != nil {
namedArgs := map[string]any{"after_id": p.after.id}
if p.after.time.IsZero() { // order by ID only
if sortOrder == dblist.SortOrderAsc {
p.where = append(p.where, dblist.WherePredicate{NamedArgs: namedArgs, SQL: "(id > @after_id)"})
where = append(where, dblist.WherePredicate{NamedArgs: namedArgs, SQL: "(id > @after_id)"})
} else {
p.where = append(p.where, dblist.WherePredicate{NamedArgs: namedArgs, SQL: "(id < @after_id)"})
where = append(where, dblist.WherePredicate{NamedArgs: namedArgs, SQL: "(id < @after_id)"})
}
} else {
namedArgs["cursor_time"] = p.after.time
if sortOrder == dblist.SortOrderAsc {
p.where = append(p.where, dblist.WherePredicate{NamedArgs: namedArgs, SQL: fmt.Sprintf(`("%s" > @cursor_time OR ("%s" = @cursor_time AND "id" > @after_id))`, timeField, timeField)})
where = append(where, dblist.WherePredicate{NamedArgs: namedArgs, SQL: fmt.Sprintf(`("%s" > @cursor_time OR ("%s" = @cursor_time AND "id" > @after_id))`, timeField, timeField)})
} else {
p.where = append(p.where, dblist.WherePredicate{NamedArgs: namedArgs, SQL: fmt.Sprintf(`("%s" < @cursor_time OR ("%s" = @cursor_time AND "id" < @after_id))`, timeField, timeField)})
where = append(where, dblist.WherePredicate{NamedArgs: namedArgs, SQL: fmt.Sprintf(`("%s" < @cursor_time OR ("%s" = @cursor_time AND "id" < @after_id))`, timeField, timeField)})
}
}
}
Expand All @@ -296,10 +327,10 @@ func (p *JobListParams) toDBParams() (*dblist.JobListParams, error) {
Priorities: p.priorities,
Queues: p.queues,
Schema: p.schema,
States: p.states,
States: states,
TagsAll: p.tagsAll,
TagsAny: p.tagsAny,
Where: p.where,
Where: where,
}, nil
}

Expand Down
233 changes: 232 additions & 1 deletion job_list_params_test.go
Original file line number Diff line number Diff line change
@@ -1,13 +1,17 @@
package river

import (
"context"
"encoding/json"
"fmt"
"testing"
"time"

"github.com/stretchr/testify/require"

"github.com/riverqueue/river/internal/dblist"
"github.com/riverqueue/river/riverdriver"
"github.com/riverqueue/river/riverdriver/riverpgxv5"
"github.com/riverqueue/river/rivertype"
)

Expand Down Expand Up @@ -210,7 +214,10 @@ func Test_JobListParams_toDBParams(t *testing.T) {
OrderBy(JobListOrderByFinalizedAt, SortOrderDesc).
toDBParams()
require.NoError(t, err)
require.Equal(t, []rivertype.JobState{rivertype.JobStateCompleted}, dbParams.States)
driverParams, err := dblist.JobMakeDriverParams(context.Background(), dbParams, riverpgxv5.New(nil))
require.NoError(t, err)
require.Equal(t, "completed", driverParams.NamedArgs["state"])
require.Equal(t, "state = @state\n AND finalized_at IS NOT NULL", driverParams.WhereClause)
})

t.Run("FinalizedAtWithMixedStates", func(t *testing.T) {
Expand Down Expand Up @@ -283,3 +290,227 @@ func Test_JobListParams_toDBParams(t *testing.T) {
require.Empty(t, dbParams.TagsAny)
})
}

func Test_JobListParams_toDBParams_CustomConditions(t *testing.T) {
t.Parallel()

type testBundle struct {
driver *riverpgxv5.Driver
params *JobListParams
}

setup := func(t *testing.T) *testBundle {
t.Helper()

return &testBundle{
driver: riverpgxv5.New(nil),
params: NewJobListParams().States(rivertype.JobStateCompleted).
OrderBy(JobListOrderByTime, SortOrderDesc),
}
}

driverParamsFunc := func(t *testing.T, bundle *testBundle) *riverdriver.JobListParams {
t.Helper()

dbParams, err := bundle.params.toDBParams()
require.NoError(t, err)
driverParams, err := dblist.JobMakeDriverParams(context.Background(), dbParams, bundle.driver)
require.NoError(t, err)
return driverParams
}

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

bundle := setup(t)

// A condition that excludes completed jobs is intentional; don't repair it.
bundle.params = bundle.params.Where("finalized_at IS NULL")
driverParams := driverParamsFunc(t, bundle)
require.Equal(t, "state = any(@state)\n AND finalized_at IS NULL", driverParams.WhereClause)
require.Equal(t, map[string]any{"state": []string{"completed"}}, driverParams.NamedArgs)
})

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

bundle := setup(t)

// Keep both state conditions, even though no row can satisfy both.
bundle.params = bundle.params.Where("state = @other_state", NamedArgs{"other_state": "available"})
driverParams := driverParamsFunc(t, bundle)
require.Equal(t, "state = any(@state)\n AND state = @other_state", driverParams.WhereClause)
require.Equal(t, map[string]any{"other_state": "available", "state": []string{"completed"}}, driverParams.NamedArgs)
})

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

bundle := setup(t)

// The typed state filter owns @state. A custom argument with the same
// name must still report a conflict, rather than replacing the filter.
bundle.params = bundle.params.Where("state = @state", NamedArgs{"state": "available"})
dbParams, err := bundle.params.toDBParams()
require.NoError(t, err)
_, err = dblist.JobMakeDriverParams(context.Background(), dbParams, bundle.driver)
require.EqualError(t, err, "named argument @state already registered")
})

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

bundle := setup(t)

// Custom SQL can reuse @state, so its value must remain an array.
bundle.params = bundle.params.Where("state = ANY(@state)")
driverParams := driverParamsFunc(t, bundle)
require.Equal(t, "state = any(@state)\n AND state = ANY(@state)", driverParams.WhereClause)
require.Equal(t, map[string]any{"state": []string{"completed"}}, driverParams.NamedArgs)
})

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

bundle := setup(t)

bundle.params = bundle.params.Where("(state = @other_state OR finalized_at IS NULL)", NamedArgs{"other_state": "completed"})
driverParams := driverParamsFunc(t, bundle)
require.Equal(t, "state = any(@state)\n AND (state = @other_state OR finalized_at IS NULL)", driverParams.WhereClause)
require.Equal(t, map[string]any{"other_state": "completed", "state": []string{"completed"}}, driverParams.NamedArgs)
})

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

bundle := setup(t)

// Metadata shares Where's internal condition list, so it also retains
// the array filter and receives no inferred null condition.
bundle.params = bundle.params.Metadata(`{"selected":true}`)
driverParams := driverParamsFunc(t, bundle)
require.Equal(t, "state = any(@state)\n AND metadata @> @metadata_fragment::jsonb", driverParams.WhereClause)
require.Equal(t, map[string]any{"metadata_fragment": `{"selected":true}`, "state": []string{"completed"}}, driverParams.NamedArgs)
})

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

bundle := setup(t)

// Appending the cursor must preserve the custom OR's existing grouping
// and leave its named argument independent of the cursor arguments.
cursorTime := time.Date(2026, 9, 9, 12, 0, 0, 0, time.UTC)
bundle.params = bundle.params.
Where("state = @other_state OR finalized_at IS NULL", NamedArgs{"other_state": "completed"}).
After(&JobListCursor{id: 42, time: cursorTime})
driverParams := driverParamsFunc(t, bundle)
const cursorSQL = `("finalized_at" < @cursor_time OR ("finalized_at" = @cursor_time AND "id" < @after_id))`
require.Equal(t, "state = any(@state)\n AND state = @other_state OR finalized_at IS NULL\n AND "+cursorSQL, driverParams.WhereClause)
require.Equal(t, map[string]any{
"after_id": int64(42),
"cursor_time": cursorTime,
"other_state": "completed",
"state": []string{"completed"},
}, driverParams.NamedArgs)
})

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

bundle := setup(t)

bundle.params = bundle.params.Where("state = ANY(@state)").
After(&JobListCursor{id: 42, time: time.Date(2026, 9, 9, 12, 0, 0, 0, time.UTC)})

// Conversion must not store generated conditions in the caller's params
// or append another cursor each time the same params are used.
first := driverParamsFunc(t, bundle)
second := driverParamsFunc(t, bundle)
require.Equal(t, first, second)
require.Equal(t, []dblist.WherePredicate{{SQL: "state = ANY(@state)"}}, bundle.params.where)
})

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

bundle := setup(t)

// This OR can admit non-finalized rows despite the typed state filter.
// Adding parentheses or a non-null condition would change the results.
bundle.params = bundle.params.Where("false OR finalized_at IS NULL")
driverParams := driverParamsFunc(t, bundle)
require.Equal(t, "state = any(@state)\n AND false OR finalized_at IS NULL", driverParams.WhereClause)
require.Equal(t, map[string]any{"state": []string{"completed"}}, driverParams.NamedArgs)
})
}

func Test_JobListParams_toDBParamsFinalizedIndex(t *testing.T) {
t.Parallel()

for _, state := range []rivertype.JobState{rivertype.JobStateCancelled, rivertype.JobStateCompleted, rivertype.JobStateDiscarded} {
for _, field := range []JobListOrderByField{JobListOrderByFinalizedAt, JobListOrderByTime} {
for _, order := range []SortOrder{SortOrderAsc, SortOrderDesc} {
t.Run(fmt.Sprintf("%s/%s/%d", state, field, order), func(t *testing.T) {
t.Parallel()

params := NewJobListParams().States(state).OrderBy(field, order)
dbParams, err := params.toDBParams()
require.NoError(t, err)
driverParams, err := dblist.JobMakeDriverParams(context.Background(), dbParams, riverpgxv5.New(nil))
require.NoError(t, err)
require.Equal(t, "state = @state\n AND finalized_at IS NOT NULL", driverParams.WhereClause)
require.Equal(t, map[string]any{"state": string(state)}, driverParams.NamedArgs)
direction := "ASC"
if order == SortOrderDesc {
direction = "DESC"
}
require.Equal(t, "finalized_at "+direction+", id "+direction, driverParams.OrderByClause)
params = params.After(&JobListCursor{id: 42, time: time.Now().UTC()})
dbParams, err = params.toDBParams()
require.NoError(t, err)
driverParams, err = dblist.JobMakeDriverParams(context.Background(), dbParams, riverpgxv5.New(nil))
require.NoError(t, err)
require.Len(t, dbParams.Where, 3)
require.Equal(t, "state = @state\n AND finalized_at IS NOT NULL\n AND "+dbParams.Where[2].SQL, driverParams.WhereClause)
require.Equal(t, map[string]any{"after_id": int64(42), "cursor_time": params.after.time, "state": string(state)}, driverParams.NamedArgs)
require.Empty(t, params.where)
repeated, err := params.toDBParams()
require.NoError(t, err)
require.Equal(t, dbParams, repeated)
})
}
}
}
}

func Test_JobListParams_toDBParamsWithoutFinalizedIndex(t *testing.T) {
t.Parallel()

for _, tt := range []struct {
name string
params *JobListParams
where string
}{
{"Default", NewJobListParams(), "state = any(@state)"},
{"FinalizedByID", NewJobListParams().States(rivertype.JobStateCompleted), "state = @state"},
{"FinalizedByScheduledAt", NewJobListParams().States(rivertype.JobStateCompleted).OrderBy(JobListOrderByScheduledAt, SortOrderAsc), "state = @state"},
{"MixedStatesByTime", NewJobListParams().States(rivertype.JobStateCompleted, rivertype.JobStateAvailable).OrderBy(JobListOrderByTime, SortOrderDesc), "state = any(@state)"},
{"MultipleFinalizedStates", NewJobListParams().OrderBy(JobListOrderByFinalizedAt, SortOrderAsc), "state = any(@state)"},
{"MultipleNonFinalizedStates", NewJobListParams().States(rivertype.JobStateAvailable, rivertype.JobStateRunning).OrderBy(JobListOrderByTime, SortOrderDesc), "state = any(@state)"},
{"NonFinalized", NewJobListParams().States(rivertype.JobStateAvailable).OrderBy(JobListOrderByTime, SortOrderAsc), "state = @state"},
{"Unfiltered", NewJobListParams().States(), "true"},
{"Unknown", NewJobListParams().States("unknown").OrderBy(JobListOrderByFinalizedAt, SortOrderAsc), "state = @state"},
{"UnknownAndFinalized", NewJobListParams().States(rivertype.JobStateCompleted, "unknown").OrderBy(JobListOrderByTime, SortOrderDesc), "state = any(@state)"},
} {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()

dbParams, err := tt.params.toDBParams()
require.NoError(t, err)
driverParams, err := dblist.JobMakeDriverParams(context.Background(), dbParams, riverpgxv5.New(nil))
require.NoError(t, err)
require.Equal(t, tt.where, driverParams.WhereClause)
})
}
}
Loading
Loading