Repository navigation
Add Java implementation of River - #1424
Conversation
0aadfc0 to
521e5b5
Compare
|
@codex review |
Codex Review SummaryThis comment shows the latest Codex review activity on this pull request.
ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings. |
There was a problem hiding this comment.
Is this file meant to be committed?
There was a problem hiding this comment.
Is this file meant to be committed? Are these differences still in place? Most seem like they're just describing Java-specific design choices and not really behavioral differences or anything deserving explicit documentation.
There was a problem hiding this comment.
Is this file meant to be committed?
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 521e5b50fa
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| -- name: notify | ||
| INSERT INTO river_notification (topic, payload) VALUES (?, ?) |
There was a problem hiding this comment.
Purge expired SQLite notification rows
Every SQLite insert and control notification appends a durable row here, but the worker maintenance path only cleans jobs and queues; a repo-wide search finds no deletion of old river_notification rows outside reset. A long-running, insert-heavy SQLite deployment therefore grows this table and its indexes without bound, consuming disk indefinitely. Add a retention-based, bounded cleanup comparable to the existing Go SQLite notification cleaner.
Useful? React with 👍 / 👎.
| for (var job : list(c, query).jobs()) | ||
| if (job.state() != Job.State.RUNNING) result.add(delete(c, job.id())); |
There was a problem hiding this comment.
Ignore jobs claimed during bulk deletion
When a worker claims a matched job after list returns but before this call, the cached state is non-running, while delete locks the current row, sees RUNNING, and throws. That rolls back the entire bulk deletion instead of honoring the documented behavior that running jobs are ignored; this makes deletion against an active queue fail under a normal claim race. Delete with a state predicate atomically, or treat the transition to running as a skipped row.
Useful? React with 👍 / 👎.
| if (!all | ||
| && query.ids().isEmpty() | ||
| && query.kinds().isEmpty() | ||
| && query.queues().isEmpty() | ||
| && query.states().isEmpty()) |
There was a problem hiding this comment.
Recognize priority filters in the bulk-delete guard
A query filtered only by priorities is rejected as an unfiltered deletion even though JobQuery.sql applies that predicate and the advertised JobDeleteMany surface supports priority filtering. The same mismatch affects metadata and tag-only queries. Include every supported deletion predicate in this safety check so callers do not have to pass all=true merely to use those filters.
Useful? React with 👍 / 👎.
| var sorted = | ||
| paths.stream() | ||
| .distinct() | ||
| .sorted((a, b) -> Json.UTF8_ORDER.compare(String.join(".", a), String.join(".", b))) |
There was a problem hiding this comment.
Order unique paths without flattening their components
Joining path components with . makes distinct paths such as List.of("a.b") and List.of("a", "b") compare equal. Because the sort is then stable, reversing those two entries changes the insertion order of the selected JSON members and produces a different uniqueness hash for the same selected field set, allowing duplicate jobs depending solely on caller-supplied path order. Compare components unambiguously instead of using a dotted flattened string.
Useful? React with 👍 / 👎.
|
@codex review |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: d7976548e9
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| if (state == Job.State.RETRYABLE) { | ||
| Duration delay = retryDelay(original, now); | ||
| scheduled = now.plus(delay); | ||
| if (delay.compareTo(config.maintenanceInterval) <= 0) state = Job.State.AVAILABLE; |
There was a problem hiding this comment.
Keep short retries scheduled until they are due
When pollInterval is longer than a retry delay that is at most maintenanceInterval (for example, a one-minute poll and a one-second retry), this changes the job to AVAILABLE with a future scheduled_at. The completion wake-up is consumed immediately while the job is still ineligible, and the scheduler ignores AVAILABLE jobs, so the retry can remain idle until the next poll—almost a minute in this example. Keep the job in a scheduler-managed state or arrange a wake-up at scheduled_at; the analogous short-snooze branch has the same problem.
Useful? React with 👍 / 👎.
| -- name: output | ||
| UPDATE river_job SET metadata=jsonb_set(metadata,'$.output',jsonb(?)) WHERE id=? RETURNING {columns} | ||
| -- name: filter_metadata | ||
| NOT EXISTS(SELECT 1 FROM json_tree(?) wanted WHERE wanted.type NOT IN ('object','array') AND NOT EXISTS(SELECT 1 FROM json_tree(metadata) actual WHERE actual.fullkey=wanted.fullkey AND actual.value IS wanted.value)) |
There was a problem hiding this comment.
Preserve empty containers in SQLite metadata filtering
For a filter containing an empty nested container, such as {"tenant": {}}, json_tree(?) produces only object rows, all of which this predicate excludes. The outer NOT EXISTS therefore succeeds for every job, including rows with no tenant key, unlike the PostgreSQL containment query. The SQLite predicate needs to check the presence of requested empty objects and arrays as well as primitive leaves.
Useful? React with 👍 / 👎.
| if (!kind.matches("[a-zA-Z0-9_][a-zA-Z0-9_\\-\\[\\]<>/.·:+]{1,126}")) | ||
| throw new IllegalArgumentException("Invalid job kind: " + kind); |
There was a problem hiding this comment.
Accept every River-valid job kind
River's shared job-kind validation permits any length that matches its character syntax, and the database stores kinds as unbounded text, but this quantifier rejects kinds longer than 127 characters. Consequently, a Java client cannot construct a JobType to insert, decode, or register a handler for a valid long kind created by another River implementation. Remove the Java-only upper bound so cross-language job kinds remain usable.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
This one is incorrect. 127 character maximum is enforced at the database level.
fd1cc7c to
f532852
Compare
|
@codex review |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: f532852ed6
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| VALUES (?::jsonb, ?, ?, ?, ?::jsonb, ?, ?, ?, ?::{schema}river_job_state, ?, ?, ?::bit(8)) | ||
| ON CONFLICT (unique_key) WHERE unique_key IS NOT NULL AND unique_states IS NOT NULL | ||
| AND {schema}river_job_state_in_bitmask(unique_states, state) | ||
| DO UPDATE SET kind = EXCLUDED.kind |
There was a problem hiding this comment.
Preserve the incumbent kind on unique conflicts
When Unique.excludeKind(true) makes two different job kinds share a key, this conflict handler rewrites the existing job's kind to the attempted kind even though the insertion is reported as skipped. The persisted arguments and other fields still belong to the incumbent, so a later worker can dispatch them to the wrong handler; make the update a no-op that retains the existing kind.
Useful? React with 👍 / 👎.
| } else if (failure instanceof InterruptedException | ||
| && context != null | ||
| && context.cancellation() == WorkContext.Cancellation.SHUTDOWN) { |
There was a problem hiding this comment.
Requeue cooperative shutdowns without requiring an exception
When stopAndCancel() signals SHUTDOWN and a cooperative handler observes the cancellation and returns normally, failure is null, so this condition is skipped and the attempt is persisted as completed. That can lose partially processed work during hard shutdown; classify any attempt whose context has SHUTDOWN cancellation as interrupted and requeue it, rather than requiring the handler to throw InterruptedException.
Useful? React with 👍 / 👎.
| properties.setProperty("user", URLDecoder.decode(auth[0], StandardCharsets.UTF_8)); | ||
| if (auth.length > 1) | ||
| properties.setProperty("password", URLDecoder.decode(auth[1], StandardCharsets.UTF_8)); |
There was a problem hiding this comment.
Preserve literal plus signs in URI credentials
For a standard PostgreSQL URI whose username or password contains a literal +, getRawUserInfo() retains that character but URLDecoder applies form-decoding rules and converts it to a space. The client therefore authenticates with different credentials and fails to connect; percent-decode URI user-info without treating + specially.
Useful? React with 👍 / 👎.
|
@codex review |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: a04c4702d7
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| DELETE FROM {schema}river_notification | ||
|
|
||
| -- name: create_schema | ||
| CREATE SCHEMA IF NOT EXISTS {schema_name} |
There was a problem hiding this comment.
Quote mixed-case schemas when creating them
When withSchema("RiverJobs") or --schema RiverJobs is used, the validation in Database accepts the name and all table references preserve its case by quoting it, but this statement is unquoted, so PostgreSQL creates the lowercased riverjobs schema. Migration 1 then targets "RiverJobs".river_migration and fails because that schema does not exist. Quote the schema identifier here, or restrict accepted names to lowercase.
Useful? React with 👍 / 👎.
|
@codex review |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 66911b6a8d
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| Sql.prepare( | ||
| c, | ||
| Sql.query(river.database(), "leader_renew"), | ||
| river.database().timestamp(now.plusSeconds(10)), |
There was a problem hiding this comment.
Tie the leadership lease to the service interval
When serviceInterval is configured to 10 seconds or longer (the builder accepts any positive duration, and tests use one hour), this fixed ten-second expiry elapses before the next renewal. Another runtime can then acquire leadership while this instance still reports isLeader() and may still have an asynchronous maintenance pass running, defeating the single-leader guarantee. Either bound serviceInterval below the lease TTL or derive the expiry from the configured renewal interval with sufficient margin.
Useful? React with 👍 / 👎.
|
@codex review |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: e169e95a0c
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| AND {schema}river_job_state_in_bitmask(unique_states, state) | ||
| -- Keep the existing kind, which may differ when uniqueness excludes kind. | ||
| DO UPDATE SET kind = river_job.kind | ||
| RETURNING {columns}, (xmax != 0) AS duplicate |
There was a problem hiding this comment.
Select a nonce-based RETURNING clause on YugabyteDB
On YugabyteDB, which Database.supportsNotifications() explicitly detects and routes through metadata-nonce duplicate detection, the xmax system column is unavailable. This unconditional expression therefore makes every PostgreSQL-style insert fail before the nonce comparison can run. Choose the RETURNING duplicate expression based on server capabilities, using a constant such as false for the nonce path.
Useful? React with 👍 / 👎.
| if (state == Job.State.RETRYABLE) { | ||
| Duration delay = retryDelay(original, now); | ||
| scheduled = now.plus(delay); | ||
| if (delay.compareTo(config.maintenanceInterval) <= 0) state = Job.State.AVAILABLE; |
There was a problem hiding this comment.
Base retry routing on the scheduler's real cadence
When serviceInterval exceeds maintenanceInterval, delays just over maintenanceInterval are persisted as RETRYABLE, but schedule() only runs from the election loop at serviceInterval. For example, with an hourly service interval, a second default failure has a roughly 16-second delay yet can remain retryable for nearly an hour. Either run scheduling at maintenanceInterval or classify delays shorter than the actual scheduler cadence as AVAILABLE so the timed fetch wake-up handles them.
Useful? React with 👍 / 👎.
| periodic | ||
| .schedule | ||
| .next(now.atOffset(java.time.ZoneOffset.UTC)) | ||
| .map(java.time.OffsetDateTime::toInstant) |
There was a problem hiding this comment.
Advance interval schedules from their planned occurrence
When maintenance observes a periodic job after its stored next time, calculating the following occurrence from now re-anchors Schedule.every to the late scan. With a one-hour schedule and a seven-minute service interval, runs can become 63 minutes apart and continue drifting instead of retaining the configured hourly cadence. Compute the next occurrence from the previous planned next value, as is already naturally required for fixed-rate periodic scheduling.
Useful? React with 👍 / 👎.
|
@codex review |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 7f4cbfa6bc
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| : rows.getBoolean("duplicate"); | ||
| if (!duplicate) plugin.afterInsert(connection, driver(), job); | ||
| if (!duplicate && job.state() == Job.State.AVAILABLE) notifyInsert(connection, job.queue()); | ||
| return new Job.InsertResult<>(typed(job, type), duplicate); |
There was a problem hiding this comment.
Avoid decoding cross-kind duplicates as the requested type
When excludeKind(true) causes an insertion to conflict with an existing job of another kind, job contains the incumbent's arguments but this call decodes them using the attempted JobType<A>. Different argument schemas can therefore make a valid duplicate insertion throw during decoding, or silently return a Job<A> whose kind and decoded arguments disagree. Preserve the incumbent as an untyped result or otherwise use a result shape that can represent cross-kind conflicts safely.
Useful? React with 👍 / 👎.
| : Job.State.RETRYABLE; | ||
| Instant scheduled = | ||
| state.isFinalized() ? job.scheduledAt() : now.plus(retryDelay(job, now)); |
There was a problem hiding this comment.
Wake short rescued retries when their delay expires
When rescuing an abandoned job with a retry delay shorter than serviceInterval, this always persists RETRYABLE. In services, scheduling runs before rescue, so the new retry cannot become available until the next maintenance pass; with an hourly service interval and a one-second retry policy, the rescued job waits nearly an hour. Apply the same scheduler-bypass/timed-fetch handling used for normal failed attempts, or schedule again when this retry becomes due.
Useful? React with 👍 / 👎.
d1ae2fd to
06d599a
Compare
|
@codex review |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 06d599ab42
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| unsubscribeInsert.run(); | ||
| leadershipWake.release(); | ||
| } | ||
| if (cancel) cancelOnStop = true; |
There was a problem hiding this comment.
Synchronize cancellation with in-flight claims
When stopAndCancel() runs while claim() is still committing, the active-map scan can finish before launch() registers the claimed job. The only remaining safeguard is launch() reading cancelOnStop, but that field is neither volatile nor protected by synchronization, so the dispatch thread is allowed to observe its old false value. The newly claimed handler then runs without a shutdown cancellation and can cause stopAndCancel() to time out; publish the flag safely or synchronize claim registration with the final cancellation scan.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
From Codex:
The reviewed commit itself also has
volatile, so this is a false positive. The new test holds a claim open until the stop call has finished its cancellation scan, then checks both modes:stopAndCancel()cancels the late job, while gracefulstop()lets it finish.
|
@codex review |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 06d599ab42
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| if (id == null || id.isBlank() || id.length() > 127) | ||
| throw new IllegalArgumentException("Invalid periodic job ID"); |
There was a problem hiding this comment.
Enforce the shared syntax for periodic IDs
A periodic ID such as "daily report" or "a" is accepted here and persisted under the reserved river:periodic_job_id metadata key, but the corresponding Go validator in internal/maintenance/periodic_job_enqueuer.go rejects IDs that do not match the shared job-kind syntax. In a mixed Java/Go deployment, the Java runtime can therefore start with a periodic definition that Go peers cannot register consistently. Apply the same syntax validation already used by JobType in addition to the length check.
Useful? React with 👍 / 👎.
|
@codex review |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 815383a187
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| var rows = statement.executeQuery()) { | ||
| var job = rows.next() ? read(rows) : before; | ||
| plugin.afterStateChange(c, driver(), job); | ||
| return job; |
There was a problem hiding this comment.
Notify workers after making a job available
When retry changes a scheduled, retryable, cancelled, completed, or discarded job to AVAILABLE, this path neither emits an insert notification nor wakes this client's dispatcher. If workers use a long pollInterval and no unrelated notification arrives, the explicitly retried job remains idle until the next poll because the scheduler ignores already-available jobs. Notify the job's queue when the update succeeds so both owned and caller-managed transactions wake workers after commit.
Useful? React with 👍 / 👎.
|
@codex review |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: a087f2cd29
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| Sql.prepare( | ||
| c, | ||
| Sql.query(river.database(), "clean_queues"), | ||
| river.database().timestamp(now.minus(Duration.ofDays(1))))) { |
There was a problem hiding this comment.
Bound heartbeats below the queue retention window
When serviceInterval is configured above one day, an actively consumed queue can satisfy this fixed cleanup horizon before its next heartbeat. The dispatcher heartbeat and leader cleanup run independently at the interval boundary, so cleanup may delete the queue first (or do so after a transient heartbeat failure); the subsequent heartbeat recreates it with default metadata and no paused_at, unexpectedly resuming a paused queue. Either reject service intervals at or above the retention period or exclude currently configured queues from cleanup.
Useful? React with 👍 / 👎.
|
@codex review |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: cf96473ba7
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| if (interval != null) return Optional.of(after.truncatedTo(ChronoUnit.SECONDS).plus(interval)); | ||
| var current = | ||
| after | ||
| .atZoneSameInstant(zone == null ? after.getOffset() : zone) |
There was a problem hiding this comment.
Evaluate unqualified cron schedules in the process time zone
For an unprefixed cron expression, this uses the reference value's fixed offset; Workers.services supplies a UTC value, whereas River's documented robfig ParseStandard semantics evaluate schedules without CRON_TZ in the process-local time zone. On any non-UTC deployment, expressions such as 0 9 * * * therefore run at 09:00 UTC instead of 09:00 local time and ignore local DST transitions. Use the JVM's system zone when no explicit zone was parsed.
Useful? React with 👍 / 👎.
|
@codex review |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: b5120a1b1a
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| for (var job : list(c, query).jobs()) | ||
| try (var statement = Sql.prepare(c, Sql.query(database, "delete"), job.id()); | ||
| var rows = statement.executeQuery()) { |
There was a problem hiding this comment.
Recheck bulk-delete filters before deleting cached IDs
On PostgreSQL, another transaction can change a listed job to a different non-running state before this per-ID delete executes. For example, a job selected by a CANCELLED filter can be explicitly retried to AVAILABLE, after which this statement still deletes it because the delete query checks only state != 'running'. This can remove runnable work that no longer matches the caller's filter; lock the selected rows or reapply the original predicates atomically when deleting.
Useful? React with 👍 / 👎.
| try { | ||
| long value = (negative ? nanos.negate() : nanos).longValueExact(); | ||
| return Duration.ofSeconds(Math.max(1, value / 1000000000)); |
There was a problem hiding this comment.
Preserve fractional seconds in @every durations
For a valid expression such as @every 1.5s, integer division converts the parsed 1.5 billion nanoseconds to a one-second Duration, so periodic jobs run 50% more often than configured and differ from Go's constant-delay schedule. Keep the parsed nanoseconds (while rounding only durations below one second up to one second) instead of truncating every duration to whole seconds.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
False positive:
- The
@every 1.5sfinding was a false positive: Go also truncates it to one second. Added explicit regression coverage.
|
@codex review |
|
Codex Review: Didn't find any major issues. Keep them coming! Reviewed commit: ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
If Codex has suggestions, it will comment; otherwise it will react with 👍. Codex can also answer questions or update the PR. Try commenting "@codex address that feedback". |
|
Oof, finally. Going to pull this in so it's more easily kept up to date, but not sure I'm going to go for a release straight away. |
No description provided.