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
2 changes: 1 addition & 1 deletion cmd/river/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@ require (
github.com/ncruces/go-strftime v1.0.0 // indirect
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect
github.com/spf13/pflag v1.0.9 // indirect
github.com/tidwall/gjson v1.19.0 // indirect
github.com/tidwall/gjson v1.20.0 // indirect
github.com/tidwall/match v1.2.0 // indirect
github.com/tidwall/pretty v1.2.1 // indirect
github.com/tidwall/sjson v1.2.5 // indirect
Expand Down
28 changes: 14 additions & 14 deletions cmd/river/go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -29,18 +29,18 @@ github.com/ncruces/go-strftime v1.0.0/go.mod h1:Fwc5htZGVVkseilnfgOVb9mKy6w1naJm
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec h1:W09IVJc94icq4NjY3clb7Lk8O1qJ8BdBEF8z0ibU0rE=
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo=
github.com/riverqueue/river v0.48.0 h1:SJwDBsNY/GuSdcW4d7xf2BeRrC1XThqd2Sns+i8ufQU=
github.com/riverqueue/river v0.48.0/go.mod h1:p0oR3A4EPF9IkyKhFiyPPRWgHRP4lF2zqtcPWs39Kp0=
github.com/riverqueue/river/riverdriver v0.48.0 h1:7oKPZb3tNjvNyQsi+JapWWREpLnOdUamb9YTnLGY3Lg=
github.com/riverqueue/river/riverdriver v0.48.0/go.mod h1:1MpM6Mf/VqlJp/iF4pzTcPrF/cMsLZ96ARw1tjfmewI=
github.com/riverqueue/river/riverdriver/riverpgxv5 v0.48.0 h1:IIS1ZfusT5zdTC5oUF/MnIthYKjUEyUGBZB3QmuqyOU=
github.com/riverqueue/river/riverdriver/riverpgxv5 v0.48.0/go.mod h1:CZe68VANQs/2fjvw96AJC0XXyptJW7KiGHSFT5XSUvo=
github.com/riverqueue/river/riverdriver/riversqlite v0.48.0 h1:6qdp9pVIlw7aPHYo3W96Vfk750j1CIKhzeUCLP0/U/U=
github.com/riverqueue/river/riverdriver/riversqlite v0.48.0/go.mod h1:g9uElXcqkIwfq7/ebVz3bL2vD7i9mLSlTltqK5sXD2I=
github.com/riverqueue/river/rivershared v0.48.0 h1:PGKa+ke7nqgBqchaSOXtQJ6Ghok5wSWJg3Su5m+m0PU=
github.com/riverqueue/river/rivershared v0.48.0/go.mod h1:FmqY+WQVCot+obBqsTWoPcIOHjYXdacHRMqZNw/oXRk=
github.com/riverqueue/river/rivertype v0.48.0 h1:t9giVes2Y2w9pGaeuRyUyGC7/QoJgTdr2vhY5W1PazU=
github.com/riverqueue/river/rivertype v0.48.0/go.mod h1:XKkcRQR6zm8RR/JQa1Q2ywpj8uXQu21quPa4Lpw1Xhw=
github.com/riverqueue/river v0.49.0 h1:JUCLFgregbX1Wu+bSTHCjF/WHGdrwGbtz/ddInHfeb0=
github.com/riverqueue/river v0.49.0/go.mod h1:USYb57gpMBXLQm2l7poskAalOPaf5uHjSZgp2AOoYDw=
github.com/riverqueue/river/riverdriver v0.49.0 h1:kSykNQJNB7AeG6kudjB0mThV29PvykoOBIFT8ipCHEc=
github.com/riverqueue/river/riverdriver v0.49.0/go.mod h1:rVUuX/fTF2kAiJjpPT+10Lft8n90kvLqjm6b+rFXaVE=
github.com/riverqueue/river/riverdriver/riverpgxv5 v0.49.0 h1:c7YA1plP/nrNS3SEHuKizf4wHx9OQ87Yxuhj3V0eQ50=
github.com/riverqueue/river/riverdriver/riverpgxv5 v0.49.0/go.mod h1:o7zkFstM+Fk+mnlOW2sx92mQW4JswSc6zvN191l8Isg=
github.com/riverqueue/river/riverdriver/riversqlite v0.49.0 h1:DfATYWIuQ4bKjhhqj/XYn3ZGpluXqG0VAZfk+voL558=
github.com/riverqueue/river/riverdriver/riversqlite v0.49.0/go.mod h1:sOSH7dNsGt2UWLFUk7Cy9qNRNkRTswFrTwYf+TJkpeg=
github.com/riverqueue/river/rivershared v0.49.0 h1:wnCYVwftMiu85kT1JrUPKEgujKkBIWoRtSNZJTUtotY=
github.com/riverqueue/river/rivershared v0.49.0/go.mod h1:E8UzQAdDutFT8rVL1wZeNbuphmIS81TgngU7/1F4U5E=
github.com/riverqueue/river/rivertype v0.49.0 h1:3up3P2DtOqnM2yEYSC1xeK9AjrljOtIfl10yZlaF4aM=
github.com/riverqueue/river/rivertype v0.49.0/go.mod h1:XKkcRQR6zm8RR/JQa1Q2ywpj8uXQu21quPa4Lpw1Xhw=
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/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM=
Expand All @@ -54,8 +54,8 @@ github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/
github.com/stretchr/testify v1.12.1 h1:EuwCh5fleGS7H32xRwO3wRGT7DxrDhLAT6FF8MpWDWE=
github.com/stretchr/testify v1.12.1/go.mod h1:MDEgiDPPsNp5cuIrHPPCyornHKgEVbtFUmoNlxoYthg=
github.com/tidwall/gjson v1.14.2/go.mod h1:/wbyibRr2FHMks5tjHJ5F8dMZh3AcwJEMf5vlfC0lxk=
github.com/tidwall/gjson v1.19.0 h1:xwxm7n691Uf3u5OFjzngavjGTh55KX5q/9w9xHW88JU=
github.com/tidwall/gjson v1.19.0/go.mod h1:V37/opeE/JbLUOfH0QTXiNez2l0RUjYUhpT4szFQAfc=
github.com/tidwall/gjson v1.20.0 h1:+agJ3rEzKcCXKDKo0ml26UROcpcAlnZyjq+TjbY5Fto=
github.com/tidwall/gjson v1.20.0/go.mod h1:V37/opeE/JbLUOfH0QTXiNez2l0RUjYUhpT4szFQAfc=
github.com/tidwall/match v1.1.1/go.mod h1:eRSPERbgtNPcGhD8UCthc6PmLEQXEWd3PRB5JTxsfmM=
github.com/tidwall/match v1.2.0 h1:0pt8FlkOwjN2fPt4bIl4BoNxb98gGHN2ObFEDkrfZnM=
github.com/tidwall/match v1.2.0/go.mod h1:eRSPERbgtNPcGhD8UCthc6PmLEQXEWd3PRB5JTxsfmM=
Expand Down
23 changes: 6 additions & 17 deletions conformance/cmd/generatefixtures/snooze_counters.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,9 +17,9 @@ type snoozeCounterCase struct {
Name string `json:"name"`
}

// makeSnoozeCounters records the `snoozes` count the job executor writes when
// a job with the given metadata snoozes, including how it coerces values that
// aren't integers.
// makeSnoozeCounters records increments of canonical non-negative integer
// counters and initialization when the counter is absent. Recovery from invalid
// counter values is implementation-specific, not part of the shared protocol.
func makeSnoozeCounters() snoozeCounters {
fixture := snoozeCounters{
Comment: generatedComment("jobexecutor.NextSnoozeCount"),
Expand All @@ -30,21 +30,10 @@ func makeSnoozeCounters() snoozeCounters {
name string
}{
{metadata: `{}`, name: "absent"},
{metadata: `{"snoozes":2}`, name: "integer"},
{metadata: `{"snoozes":2.9}`, name: "fraction_truncates"},
{metadata: `{"snoozes":-2.5}`, name: "negative_fraction_truncates_toward_zero"},
{metadata: `{"snoozes":1e3}`, name: "exponent"},
{metadata: `{"snoozes":9007199254740993}`, name: "beyond_float_precision"},
{metadata: `{"snoozes":"4"}`, name: "numeric_string"},
{metadata: `{"snoozes":"-7"}`, name: "negative_numeric_string"},
{metadata: `{"snoozes":"4.5"}`, name: "fractional_string_is_zero"},
{metadata: `{"snoozes":" 5"}`, name: "padded_string_is_zero"},
{metadata: `{"snoozes":"abc"}`, name: "non_numeric_string_is_zero"},
{metadata: `{"snoozes":true}`, name: "true_is_one"},
{metadata: `{"snoozes":false}`, name: "false_is_zero"},
{metadata: `{"snoozes":null}`, name: "null_is_zero"},
{metadata: `{"snoozes":[3]}`, name: "array_is_zero"},
{metadata: `{"snoozes":{"count":3}}`, name: "object_is_zero"},
{metadata: `{"snoozes":2}`, name: "integer"},
{metadata: `{"snoozes":9223372036854775806}`, name: "largest_increment"},
{metadata: `{"snoozes":0}`, name: "zero"},
} {
fixture.SnoozeCounters = append(fixture.SnoozeCounters, snoozeCounterCase{
ExpectedSnoozes: jobexecutor.NextSnoozeCount([]byte(testCase.metadata)),
Expand Down
3 changes: 2 additions & 1 deletion conformance/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -17,9 +17,10 @@ require (
)

require (
github.com/tidwall/gjson v1.19.0 // indirect
github.com/tidwall/gjson v1.20.0 // indirect
github.com/tidwall/match v1.2.0 // indirect
github.com/tidwall/pretty v1.2.1 // indirect
github.com/tidwall/sjson v1.2.5 // indirect
golang.org/x/sync v0.23.0 // indirect
golang.org/x/tools v0.50.0 // indirect
)
8 changes: 4 additions & 4 deletions conformance/go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -21,8 +21,8 @@ github.com/robfig/cron/v3 v3.0.1/go.mod h1:eQICP3HwyT7UooqI/z+Ov+PtYAWygg1TEWWzG
github.com/stretchr/testify v1.12.1 h1:EuwCh5fleGS7H32xRwO3wRGT7DxrDhLAT6FF8MpWDWE=
github.com/stretchr/testify v1.12.1/go.mod h1:MDEgiDPPsNp5cuIrHPPCyornHKgEVbtFUmoNlxoYthg=
github.com/tidwall/gjson v1.14.2/go.mod h1:/wbyibRr2FHMks5tjHJ5F8dMZh3AcwJEMf5vlfC0lxk=
github.com/tidwall/gjson v1.19.0 h1:xwxm7n691Uf3u5OFjzngavjGTh55KX5q/9w9xHW88JU=
github.com/tidwall/gjson v1.19.0/go.mod h1:V37/opeE/JbLUOfH0QTXiNez2l0RUjYUhpT4szFQAfc=
github.com/tidwall/gjson v1.20.0 h1:+agJ3rEzKcCXKDKo0ml26UROcpcAlnZyjq+TjbY5Fto=
github.com/tidwall/gjson v1.20.0/go.mod h1:V37/opeE/JbLUOfH0QTXiNez2l0RUjYUhpT4szFQAfc=
github.com/tidwall/match v1.1.1/go.mod h1:eRSPERbgtNPcGhD8UCthc6PmLEQXEWd3PRB5JTxsfmM=
github.com/tidwall/match v1.2.0 h1:0pt8FlkOwjN2fPt4bIl4BoNxb98gGHN2ObFEDkrfZnM=
github.com/tidwall/match v1.2.0/go.mod h1:eRSPERbgtNPcGhD8UCthc6PmLEQXEWd3PRB5JTxsfmM=
Expand All @@ -41,5 +41,5 @@ golang.org/x/sync v0.23.0 h1:KameEIfc1IkluZyXWLn39Wd4tURc6GbCiISGiZm2bQk=
golang.org/x/sync v0.23.0/go.mod h1:sUUOizhqBxiL6pEWpqNLUiaJn1ShEbZ6BBqskPbjZm0=
golang.org/x/text v0.42.0 h1:JbOZXgfeCPU9gacVtYliJqOhD+zhrEqK4LfdpmlUZqI=
golang.org/x/text v0.42.0/go.mod h1:ojzP1Z+2QtioaF8DTtO8K5q7JWVVYwZKenzujK0Zd0E=
golang.org/x/tools v0.49.0 h1:3NI7VXzL9+1WZD52Dx2ttoPwD5DWrFGpl9mFZDlmisI=
golang.org/x/tools v0.49.0/go.mod h1:SJNXV9DBKT0UbdttsQjbfJlAE/q+y36++zo3uL3N0Oo=
golang.org/x/tools v0.50.0 h1:c2ifzfcuY7L90lZ2aKd8S4K2NpASF08SZx9ZuJkHmSU=
golang.org/x/tools v0.50.0/go.mod h1:7ulVMw3831Mwi5EZD6RomGyffr4VFjuNYXf2BbCEAV0=
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ require (
github.com/riverqueue/river/rivertype v0.49.0
github.com/robfig/cron/v3 v3.0.1
github.com/stretchr/testify v1.12.1
github.com/tidwall/gjson v1.19.0
github.com/tidwall/gjson v1.20.0
github.com/tidwall/sjson v1.2.5
golang.org/x/sync v0.23.0
)
Expand Down
20 changes: 10 additions & 10 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -10,14 +10,14 @@ github.com/jackc/pgx/v5 v5.11.0/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QII
github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo=
github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/riverqueue/river/riverdriver v0.48.0 h1:7oKPZb3tNjvNyQsi+JapWWREpLnOdUamb9YTnLGY3Lg=
github.com/riverqueue/river/riverdriver v0.48.0/go.mod h1:1MpM6Mf/VqlJp/iF4pzTcPrF/cMsLZ96ARw1tjfmewI=
github.com/riverqueue/river/riverdriver/riverpgxv5 v0.48.0 h1:IIS1ZfusT5zdTC5oUF/MnIthYKjUEyUGBZB3QmuqyOU=
github.com/riverqueue/river/riverdriver/riverpgxv5 v0.48.0/go.mod h1:CZe68VANQs/2fjvw96AJC0XXyptJW7KiGHSFT5XSUvo=
github.com/riverqueue/river/rivershared v0.48.0 h1:PGKa+ke7nqgBqchaSOXtQJ6Ghok5wSWJg3Su5m+m0PU=
github.com/riverqueue/river/rivershared v0.48.0/go.mod h1:FmqY+WQVCot+obBqsTWoPcIOHjYXdacHRMqZNw/oXRk=
github.com/riverqueue/river/rivertype v0.48.0 h1:t9giVes2Y2w9pGaeuRyUyGC7/QoJgTdr2vhY5W1PazU=
github.com/riverqueue/river/rivertype v0.48.0/go.mod h1:XKkcRQR6zm8RR/JQa1Q2ywpj8uXQu21quPa4Lpw1Xhw=
github.com/riverqueue/river/riverdriver v0.49.0 h1:kSykNQJNB7AeG6kudjB0mThV29PvykoOBIFT8ipCHEc=
github.com/riverqueue/river/riverdriver v0.49.0/go.mod h1:rVUuX/fTF2kAiJjpPT+10Lft8n90kvLqjm6b+rFXaVE=
github.com/riverqueue/river/riverdriver/riverpgxv5 v0.49.0 h1:c7YA1plP/nrNS3SEHuKizf4wHx9OQ87Yxuhj3V0eQ50=
github.com/riverqueue/river/riverdriver/riverpgxv5 v0.49.0/go.mod h1:o7zkFstM+Fk+mnlOW2sx92mQW4JswSc6zvN191l8Isg=
github.com/riverqueue/river/rivershared v0.49.0 h1:wnCYVwftMiu85kT1JrUPKEgujKkBIWoRtSNZJTUtotY=
github.com/riverqueue/river/rivershared v0.49.0/go.mod h1:E8UzQAdDutFT8rVL1wZeNbuphmIS81TgngU7/1F4U5E=
github.com/riverqueue/river/rivertype v0.49.0 h1:3up3P2DtOqnM2yEYSC1xeK9AjrljOtIfl10yZlaF4aM=
github.com/riverqueue/river/rivertype v0.49.0/go.mod h1:XKkcRQR6zm8RR/JQa1Q2ywpj8uXQu21quPa4Lpw1Xhw=
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/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
Expand All @@ -26,8 +26,8 @@ github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/
github.com/stretchr/testify v1.12.1 h1:EuwCh5fleGS7H32xRwO3wRGT7DxrDhLAT6FF8MpWDWE=
github.com/stretchr/testify v1.12.1/go.mod h1:MDEgiDPPsNp5cuIrHPPCyornHKgEVbtFUmoNlxoYthg=
github.com/tidwall/gjson v1.14.2/go.mod h1:/wbyibRr2FHMks5tjHJ5F8dMZh3AcwJEMf5vlfC0lxk=
github.com/tidwall/gjson v1.19.0 h1:xwxm7n691Uf3u5OFjzngavjGTh55KX5q/9w9xHW88JU=
github.com/tidwall/gjson v1.19.0/go.mod h1:V37/opeE/JbLUOfH0QTXiNez2l0RUjYUhpT4szFQAfc=
github.com/tidwall/gjson v1.20.0 h1:+agJ3rEzKcCXKDKo0ml26UROcpcAlnZyjq+TjbY5Fto=
github.com/tidwall/gjson v1.20.0/go.mod h1:V37/opeE/JbLUOfH0QTXiNez2l0RUjYUhpT4szFQAfc=
github.com/tidwall/match v1.1.1/go.mod h1:eRSPERbgtNPcGhD8UCthc6PmLEQXEWd3PRB5JTxsfmM=
github.com/tidwall/match v1.2.0 h1:0pt8FlkOwjN2fPt4bIl4BoNxb98gGHN2ObFEDkrfZnM=
github.com/tidwall/match v1.2.0/go.mod h1:eRSPERbgtNPcGhD8UCthc6PmLEQXEWd3PRB5JTxsfmM=
Expand Down
6 changes: 3 additions & 3 deletions internal/jobexecutor/job_executor.go
Original file line number Diff line number Diff line change
Expand Up @@ -81,9 +81,9 @@ func MetadataUpdatesFromWorkContext(ctx context.Context) (map[string]any, bool)
}

// NextSnoozeCount returns the snooze count recorded on a job that's snoozed
// again: its metadata's current `snoozes` value plus one. The current value is
// read leniently, so a missing, non-numeric, or malformed count reads as zero
// and a fractional one truncates toward zero.
// again: its metadata's current integer `snoozes` value plus one, starting at
// one when absent. Invalid values are read leniently; their exact coercion is
// not part of the cross-language protocol.
func NextSnoozeCount(metadata []byte) int64 {
return gjson.GetBytes(metadata, "snoozes").Int() + 1
}
Expand Down
34 changes: 34 additions & 0 deletions internal/jobexecutor/job_executor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1480,6 +1480,40 @@ func TestJobExecutor_Execute(t *testing.T) {
})
}

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

for _, testCase := range []struct {
name string
value string
}{
{name: "Array", value: `[3]`},
{name: "Boolean", value: `true`},
{name: "ExcessiveInteger", value: `9223372036854775808`},
{name: "ExcessiveNumber", value: `1e400`},
{name: "ExponentString", value: `"1e3"`},
{name: "FractionalNumber", value: `2.9`},
{name: "FractionalString", value: `"4.5"`},
{name: "MaxInteger", value: `9223372036854775807`},
{name: "NegativeInteger", value: `-2`},
{name: "Null", value: `null`},
{name: "NumericString", value: `"4"`},
{name: "Object", value: `{"count":3}`},
{name: "String", value: `"abc"`},
} {
t.Run(testCase.name, func(t *testing.T) {
t.Parallel()

metadata := []byte(`{"snoozes":` + testCase.value + `}`)
require.NotPanics(t, func() {
// The return type guarantees an integer. Invalid-value coercion
// is intentionally not pinned to a particular result.
_ = NextSnoozeCount(metadata)
})
})
}
}

//
// *Func types are copied from the top level River package because they can't be
// accessed from here.
Expand Down
2 changes: 1 addition & 1 deletion js/src/runtime.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3425,7 +3425,7 @@ describe("Client runtime", () => {
driver.claim = [
{
...fakeJob(),
metadata: { snoozes: "2", user: true },
metadata: { snoozes: 2, user: true },
},
];
const definition = defineJob({ kind: "test" });
Expand Down
56 changes: 42 additions & 14 deletions js/src/runtime/completion-command.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -50,20 +50,10 @@ describe("completionCommand", () => {
readonly snooze_counters: readonly SnoozeCounterCase[];
};
expect(golden.snooze_counters.length).toBeGreaterThan(0);
const now = Temporal.Instant.from("2026-09-01T00:00:00Z");

const results = golden.snooze_counters.map(({ metadata, name }) => {
const command = completionCommand(
job(metadata),
"client",
{ outcome: snooze({ seconds: 30 }), status: "succeeded" },
now,
now,
now,
5_000
);
return [name, numberText(command.metadata?.snoozes)];
});
const results = golden.snooze_counters.map(({ metadata, name }) => [
name,
numberText(snoozeCommand(metadata).metadata?.snoozes),
]);

expect(results).toEqual(
golden.snooze_counters.map(({ expected_snoozes, name }) => [
Expand All @@ -72,6 +62,31 @@ describe("completionCommand", () => {
])
);
});

it.each([
"2.9",
"-2",
"1e400",
"9223372036854775807",
"9223372036854775808",
'"4"',
'"4.5"',
'"1e3"',
'" 5"',
'"abc"',
"true",
"false",
"null",
"[3]",
'{"count":3}',
])("recovers safely from invalid snooze counter %s", (raw) => {
const metadata = parseJson(`{"snoozes":${raw}}`) as JsonObject;
const command = snoozeCommand(metadata);

expect(command.kind).toBe("snooze");
expect(command.scheduledAt).not.toBeNull();
expect(BigInt(numberText(command.metadata?.snoozes))).toBeGreaterThan(0n);
});
});

function job(metadata: JsonObject): JobRow {
Expand Down Expand Up @@ -104,3 +119,16 @@ function numberText(value: JsonValue | undefined): string {
if (isExactJsonNumber(value)) return value.rawJSON;
throw new Error(`not a JSON number: ${JSON.stringify(value)}`);
}

function snoozeCommand(metadata: JsonObject) {
const now = Temporal.Instant.from("2026-09-01T00:00:00Z");
return completionCommand(
job(metadata),
"client",
{ outcome: snooze({ seconds: 30 }), status: "succeeded" },
now,
now,
now,
5_000
);
}
56 changes: 13 additions & 43 deletions js/src/runtime/completion-command.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,8 @@ import { toMilliseconds } from "../internal/duration.js";
import type { JobRow } from "../job.js";
import {
exactJsonNumber,
isExactJsonNumber,
isJsonNumber,
jsonNumberToBigInt,
type ExactJsonNumber,
type JsonValue,
} from "../json.js";
Expand Down Expand Up @@ -165,56 +166,25 @@ export function completionCommand(
}

/**
* The snooze count after one more snooze, like River for Go's executor,
* which writes `int(gjson.GetBytes(metadata, "snoozes").Int()) + 1`: `true`
* counts as 1, a decimal integer string as its value, a number truncated
* toward zero, and anything else as 0, with int64 wraparound.
* Increment an integer counter, restarting missing, invalid, or overflowing
* counters at one. Invalid-value recovery is not a cross-language contract.
*/
function nextSnoozeCount(
value: JsonValue | undefined
): ExactJsonNumber | number {
const next = BigInt.asIntN(64, gjsonInt(value) + 1n);
return next >= BigInt(Number.MIN_SAFE_INTEGER) &&
next <= BigInt(Number.MAX_SAFE_INTEGER)
let count: bigint;
try {
count = isJsonNumber(value) ? jsonNumberToBigInt(value) : 0n;
} catch {
count = 0n;
}
const next =
count >= 0n && count < 9_223_372_036_854_775_807n ? count + 1n : 1n;
return next <= BigInt(Number.MAX_SAFE_INTEGER)
? Number(next)
: exactJsonNumber(next.toString(10));
}

/** gjson's `Result.Int()` for a JSON value. */
function gjsonInt(value: JsonValue | undefined): bigint {
if (value === true) return 1n;
if (typeof value === "string") return gjsonParseInt(value) ?? 0n;
if (typeof value !== "number" && !isExactJsonNumber(value)) return 0n;
const raw = typeof value === "number" ? String(value) : value.rawJSON;
const float = typeof value === "number" ? value : Number(value.rawJSON);
// gjson's safeInt, then its parse of the raw integer text, then Go's
// float conversion, which saturates out of range on arm64.
if (Math.abs(float) <= Number.MAX_SAFE_INTEGER) {
return BigInt(Math.trunc(float));
}
const parsed = gjsonParseInt(raw);
if (parsed !== undefined) return parsed;
if (Number.isNaN(float)) return 0n;
if (float >= 2 ** 63) return BigInt.asIntN(64, (1n << 63n) - 1n);
if (float <= -(2 ** 63)) return -(1n << 63n);
return BigInt(Math.trunc(float));
}

/**
* gjson's `parseInt`: an optional `-` then decimal digits only, wrapping
* like int64 arithmetic; undefined for anything else.
*/
function gjsonParseInt(text: string): bigint | undefined {
const negative = text.startsWith("-");
const digits = negative ? text.slice(1) : text;
if (!/^[0-9]+$/.test(digits)) return undefined;
let result = 0n;
for (const digit of digits) {
result = BigInt.asIntN(64, result * 10n + BigInt(digit));
}
return negative ? BigInt.asIntN(64, -result) : result;
}

/** The event announcing a committed completion, by the job's new state. */
export function completionEventKind(
requested: JobCompletionCommand["kind"],
Expand Down
4 changes: 2 additions & 2 deletions riverdriver/go.sum
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
github.com/riverqueue/river/rivertype v0.48.0 h1:t9giVes2Y2w9pGaeuRyUyGC7/QoJgTdr2vhY5W1PazU=
github.com/riverqueue/river/rivertype v0.48.0/go.mod h1:XKkcRQR6zm8RR/JQa1Q2ywpj8uXQu21quPa4Lpw1Xhw=
github.com/riverqueue/river/rivertype v0.49.0 h1:3up3P2DtOqnM2yEYSC1xeK9AjrljOtIfl10yZlaF4aM=
github.com/riverqueue/river/rivertype v0.49.0/go.mod h1:XKkcRQR6zm8RR/JQa1Q2ywpj8uXQu21quPa4Lpw1Xhw=
github.com/stretchr/testify v1.12.1 h1:EuwCh5fleGS7H32xRwO3wRGT7DxrDhLAT6FF8MpWDWE=
github.com/stretchr/testify v1.12.1/go.mod h1:MDEgiDPPsNp5cuIrHPPCyornHKgEVbtFUmoNlxoYthg=
go.yaml.in/yaml/v3 v3.0.5 h1:N6y/pJk8buWs9NY5ERU2HSMfm+IuD/OtfdAnq6kESPw=
Expand Down
Loading
Loading