diff --git a/.github/actions/setup-java/action.yaml b/.github/actions/setup-java/action.yaml
new file mode 100644
index 000000000..709f5b3cf
--- /dev/null
+++ b/.github/actions/setup-java/action.yaml
@@ -0,0 +1,17 @@
+name: Set up Java workspace
+description: Install a supported JDK and cache Maven dependencies for the Java workspace.
+
+inputs:
+ java-version:
+ default: "25"
+ description: JDK version to use.
+
+runs:
+ using: composite
+ steps:
+ - uses: actions/setup-java@v6
+ with:
+ distribution: temurin
+ java-version: ${{ inputs.java-version }}
+ cache: maven
+ cache-dependency-path: "java/**/pom.xml"
diff --git a/.github/workflows/conformance.yaml b/.github/workflows/conformance.yaml
index 04899d01f..4b2abb030 100644
--- a/.github/workflows/conformance.yaml
+++ b/.github/workflows/conformance.yaml
@@ -1,6 +1,6 @@
name: Conformance
-# The Rust and JavaScript workflows run only when their own files change, but
+# The port workflows run only when their own files or dependencies change, but
# their fixture tests compare against values generated from River's Go code.
# This job closes that gap: when Go code, SQL queries the generator reads, or
# the generator itself changes, it regenerates the fixtures and runs only the
@@ -11,6 +11,7 @@ on:
- master
paths:
- ".github/workflows/conformance.yaml"
+ - "Makefile"
- "**.go"
- "**/go.mod"
- "**/go.sum"
@@ -20,6 +21,7 @@ on:
pull_request:
paths:
- ".github/workflows/conformance.yaml"
+ - "Makefile"
- "**.go"
- "**/go.mod"
- "**/go.sum"
@@ -69,9 +71,16 @@ jobs:
- uses: ./.github/actions/setup-js
+ - uses: ./.github/actions/setup-java
+ with:
+ java-version: "21"
+
# Each target generates the fixtures before running its tests.
- name: Rust fixture tests
run: make test/rust/conformance
- name: JavaScript fixture tests
run: make test/js/conformance
+
+ - name: Java fixture tests
+ run: make test/java/conformance
diff --git a/.github/workflows/java.yaml b/.github/workflows/java.yaml
new file mode 100644
index 000000000..ff482e0ca
--- /dev/null
+++ b/.github/workflows/java.yaml
@@ -0,0 +1,132 @@
+name: Java
+
+# Filter the workflow so unrelated changes don't add skipped jobs to CI.
+on:
+ push:
+ branches:
+ - master
+ paths:
+ - ".github/actions/setup-java/**"
+ - ".github/workflows/java.yaml"
+ - "Makefile"
+ - "conformance/**"
+ - "java/**"
+ - "riverdriver/riverpgxv5/migration/**"
+ - "riverdriver/riversqlite/migration/**"
+ pull_request:
+ paths:
+ - ".github/actions/setup-java/**"
+ - ".github/workflows/java.yaml"
+ - "Makefile"
+ - "conformance/**"
+ - "java/**"
+ - "riverdriver/riverpgxv5/migration/**"
+ - "riverdriver/riversqlite/migration/**"
+ workflow_dispatch:
+
+concurrency:
+ cancel-in-progress: ${{ github.event_name == 'pull_request' }}
+ group: ${{ github.workflow }}-${{ github.ref }}
+
+permissions:
+ contents: read
+
+jobs:
+ quality:
+ name: Quality and package archives
+ runs-on: ubuntu-latest
+ timeout-minutes: 15
+
+ steps:
+ - uses: actions/checkout@v6
+ with:
+ persist-credentials: false
+
+ - uses: ./.github/actions/setup-java
+
+ - uses: actions/setup-go@v6
+ with:
+ go-version-file: go.work
+
+ # Install once; lint/java invokes the linter for each standalone Go tool.
+ - name: Set up Go linter
+ uses: golangci/golangci-lint-action@v9
+ with:
+ args: --help
+ version: v2.13.1
+
+ - name: Compile, check formatting, and lint maintenance tools
+ run: make lint/java
+
+ - name: Test Go maintenance tools
+ run: make test/java/tools
+
+ - name: Verify canonical migration mirror
+ run: make verify/java-migrations
+
+ - name: Check publishable archives
+ run: make check/java/package
+
+ java_versions:
+ name: Test (Java ${{ matrix.java-version }}, SQLite)
+ runs-on: ubuntu-latest
+ timeout-minutes: 15
+ strategy:
+ fail-fast: false
+ matrix:
+ java-version: ["21", "25"]
+
+ steps:
+ - uses: actions/checkout@v6
+ with:
+ persist-credentials: false
+
+ - uses: ./.github/actions/setup-java
+ with:
+ java-version: ${{ matrix.java-version }}
+
+ - uses: actions/setup-go@v6
+ with:
+ go-version-file: go.work
+
+ - name: Unit, SQLite, and executable CLI tests
+ run: make test/java/sqlite
+
+ postgres:
+ name: Test (PostgreSQL ${{ matrix.postgres-version }})
+ runs-on: ubuntu-latest
+ timeout-minutes: 15
+ strategy:
+ fail-fast: false
+ matrix:
+ postgres-version: [14, 15, 16, 17, 18]
+ env:
+ RIVER_TEST_DATABASE_URL: postgres://postgres:postgres@127.0.0.1:5432/river_test?sslmode=disable
+
+ services:
+ postgres:
+ image: postgres:${{ matrix.postgres-version }}
+ env:
+ POSTGRES_DB: river_test
+ POSTGRES_PASSWORD: postgres
+ options: >-
+ --health-cmd "pg_isready -U postgres -d river_test"
+ --health-interval 2s
+ --health-timeout 5s
+ --health-retries 5
+ ports:
+ - 5432:5432
+
+ steps:
+ - uses: actions/checkout@v6
+ with:
+ persist-credentials: false
+
+ - uses: ./.github/actions/setup-java
+
+ - uses: actions/setup-go@v6
+ with:
+ go-version-file: go.work
+
+ - name: PostgreSQL client, worker, and migration tests
+ run: make test/java/postgres
diff --git a/Makefile b/Makefile
index 170de5979..4cfeb5593 100644
--- a/Makefile
+++ b/Makefile
@@ -1,9 +1,10 @@
.DEFAULT_GOAL := help
SQLC ?= sqlc
+MVN ?= mvn
.PHONY: check/modzip
-check/modzip: ## Check that no Go module zip includes fixtures, testdata, or the Rust or JS ports
+check/modzip: ## Check that no Go module zip includes fixtures, testdata, or another language's port
go run ./conformance/cmd/checkmodzip ./go.work
.PHONY: db/reset
@@ -24,6 +25,7 @@ db/reset/test: ## Drop, create, and migrate test databases
.PHONY: generate
generate: ## Generate generated artifacts
generate: generate/fixtures
+generate: generate/java-migrations
generate: generate/js-migrations
generate: generate/migrations
generate: generate/rust-migrations
@@ -35,6 +37,10 @@ generate: generate/sqlc
generate/fixtures: ## Generate cross-language conformance fixtures from River's Go implementation
go run ./conformance/cmd/generatefixtures
+.PHONY: generate/java-migrations
+generate/java-migrations: ## Sync database migrations to Java
+ go run ./java/bin/sync-migrations/main.go
+
.PHONY: generate/js-migrations
generate/js-migrations: ## Sync database migrations to JavaScript
pnpm -C js run generate:migrations
@@ -107,6 +113,23 @@ lint/rust: ## Run Rust formatting and clippy checks, including single-backend bu
cd rust && cargo clippy -p riverqueue -p riverqueue-migrate -p riverqueue-cli -p riverqueue-test --no-default-features --features sqlite --all-targets --locked -- -D warnings
cd rust && $(RUST_POSTGRES_TESTS_ENV) cargo clippy -p riverqueue -p riverqueue-migrate --all-targets --all-features --locked -- -D warnings
+# Java targets stay separate so Go-only contributors do not need a JDK or Maven.
+.PHONY: build/java
+build/java: ## Build the Java library, CLI, and development adapter (JDK 21+)
+ $(MVN) --batch-mode --no-transfer-progress -f java/pom.xml package -DskipTests
+
+.PHONY: lint/java
+lint/java: ## Compile Java, check formatting, and lint Go maintenance tools
+lint/java: lint/java/tools
+ $(MVN) --batch-mode --no-transfer-progress -f java/pom.xml verify -DskipTests
+
+# Java's module boundary excludes it from Go package discovery. Pass the tool
+# files explicitly so they use this workspace's Go version and test dependencies.
+.PHONY: lint/java/tools
+lint/java/tools: ## Lint Java's Go maintenance tools
+ golangci-lint run --fix ./java/bin/check-packages/*.go
+ golangci-lint run --fix ./java/bin/sync-migrations/*.go
+
# JavaScript targets, like the Rust ones, are separate from `lint` and `test`
# and need Node.js 26 and pnpm; they delegate to the workspace's own scripts.
# Run `pnpm -C js install` first.
@@ -144,6 +167,34 @@ RUST_POSTGRES_TESTS_ENV = RUSTFLAGS="$$RUSTFLAGS --cfg river_postgres_tests" \
RUSTDOCFLAGS="$$RUSTDOCFLAGS --cfg river_postgres_tests" \
CARGO_TARGET_DIR="$${CARGO_TARGET_DIR:-$(CURDIR)/rust/target}/postgres-tests"
+.PHONY: test/java
+test/java: ## Run Java tests, executable CLI tests, and formatting checks
+test/java: generate/fixtures verify/java-migrations
+test/java: test/java/tools
+ $(MVN) --batch-mode --no-transfer-progress -f java/pom.xml verify
+
+# Fixture comparisons use temporary SQLite databases, with no PostgreSQL or legacy adapter.
+.PHONY: test/java/conformance
+test/java/conformance: ## Run Java tests that check Go-generated conformance fixtures
+test/java/conformance: generate/fixtures
+ $(MVN) --batch-mode --no-transfer-progress -f java/pom.xml -pl river test -Dgroups=conformance -Driver.test.database=sqlite
+
+.PHONY: test/java/postgres
+test/java/postgres: ## Run Java tests with PostgreSQL client and worker coverage (requires RIVER_TEST_DATABASE_URL)
+test/java/postgres: generate/fixtures verify/java-migrations
+ @test -n "$$RIVER_TEST_DATABASE_URL" || { echo "RIVER_TEST_DATABASE_URL is required" >&2; exit 1; }
+ $(MVN) --batch-mode --no-transfer-progress -f java/pom.xml verify -Driver.test.database=postgres
+
+.PHONY: test/java/sqlite
+test/java/sqlite: ## Run Java unit, SQLite, and executable CLI tests without PostgreSQL
+test/java/sqlite: generate/fixtures verify/java-migrations
+ $(MVN) --batch-mode --no-transfer-progress -f java/pom.xml verify -Driver.test.database=sqlite -DexcludedGroups=postgres
+
+.PHONY: test/java/tools
+test/java/tools: ## Test Java's Go maintenance tools
+ go test ./java/bin/check-packages/*.go
+ go test ./java/bin/sync-migrations/*.go
+
# PostgreSQL integration tests need RIVER_RUST_DATABASE_URL. Without it
# test/rust still runs unit, doc, and SQLite integration tests, and fails in CI
# so a missing URL cannot turn the PostgreSQL suite into a silent pass.
@@ -218,6 +269,11 @@ doc/rust: ## Build Rust API documentation, compiled examples, and doctests for e
doc/rust/docsrs: ## Build Rust API documentation as docs.rs does (nightly toolchain, `--cfg docsrs`)
cd rust && RUSTDOCFLAGS="--cfg docsrs -D warnings" CARGO_TARGET_DIR="$${CARGO_TARGET_DIR:-target}/docsrs" cargo +nightly doc -p riverqueue -p riverqueue-migrate -p riverqueue-test --all-features --no-deps --locked
+.PHONY: check/java/package
+check/java/package: ## Build and verify publishable Maven archives without publishing
+check/java/package: build/java
+ go run ./java/bin/check-packages/main.go
+
.PHONY: check/js/dependencies
check/js/dependencies: ## Audit JavaScript advisories and production dependency licenses
pnpm -C js audit
@@ -304,11 +360,16 @@ update-mod-version: ## Update River packages in all submodules to $VERSION
.PHONY: verify
verify: ## Verify generated artifacts
+verify: verify/java-migrations
verify: verify/js-migrations
verify: verify/migrations
verify: verify/rust-migrations
verify: verify/sqlc
+.PHONY: verify/java-migrations
+verify/java-migrations: ## Verify Java migrations match the canonical migrations
+ go run ./java/bin/sync-migrations/main.go -check
+
.PHONY: verify/js-migrations
verify/js-migrations: ## Verify JavaScript migrations match the canonical migrations
pnpm -C js run verify:migrations
diff --git a/conformance/cmd/checkmodzip/main.go b/conformance/cmd/checkmodzip/main.go
index d051602ae..a002d83ff 100644
--- a/conformance/cmd/checkmodzip/main.go
+++ b/conformance/cmd/checkmodzip/main.go
@@ -27,8 +27,8 @@ const conformanceModulePath = "github.com/riverqueue/river/conformance"
// disallowedPathPattern matches paths, relative to a module's root, that must
// never be published: fixture and testdata directories, JSON files, and the
-// conformance, JavaScript, and Rust trees.
-var disallowedPathPattern = regexp.MustCompile(`(^|/)(fixtures?|testdata)/|\.json$|^(conformance|js|rust)/`)
+// conformance, Java, JavaScript, and Rust trees.
+var disallowedPathPattern = regexp.MustCompile(`(^|/)(fixtures?|testdata)/|\.json$|^(conformance|java|js|rust)/`)
func main() {
if len(os.Args) != 2 {
diff --git a/java/.gitignore b/java/.gitignore
new file mode 100644
index 000000000..b9752d435
--- /dev/null
+++ b/java/.gitignore
@@ -0,0 +1,5 @@
+target/
+.idea/
+*.iml
+.conformance/
+__pycache__/
diff --git a/java/API.md b/java/API.md
new file mode 100644
index 000000000..e047f08eb
--- /dev/null
+++ b/java/API.md
@@ -0,0 +1,100 @@
+# Java API review
+
+Reviewed against the exported APIs in this River checkout and the pinned Go
+conformance reference. This is an application API comparison, not a claim that
+Java exposes every Go helper. Stored values and protocol behavior still follow Go.
+
+| Go surface | Java surface | Reason for the shape |
+| --- | --- | --- |
+| `NewClient`, `Config`, `Start`, `Stop`, `StopAndCancel` | `Client`, `Workers.Builder`, `start`, `stop`, `stopAndCancel` | A client can insert without starting anything. A running `Workers` owns threads and implements `AutoCloseable`. |
+| `JobArgs`, `Job[T]`, worker interfaces and registration | `JobType`, `Job`, `add(type, lambda)` | Records need no River interface; explicit kinds stay stable when Java classes are renamed. |
+| `Insert`, `InsertTx`, `InsertMany`, `InsertManyTx` | `insert`, `insertMany`, connection-first overloads | Arguments are typed; results contain persisted `Job` rows, including cross-kind duplicates. Mixed batches use `JobType.submission`. |
+| `InsertOpts`, `UniqueOpts`, struct tags | `InsertOptions`, `Unique`, `JobType.uniqueBy` | Named builders and fluent options avoid positional flags. Explicit JSON paths replace struct tags. |
+| `JobGet`, `JobList`, cursor/order/filter params | `get`, `JobQuery`, `JobQuery.Page` | Raw reads use `JsonNode`; `get(id, type)` checks the kind and decodes a known argument type. Pagination remains explicit. |
+| `JobCancel`, `JobDelete`, `JobDeleteMany`, `JobRetry` and transaction variants | `cancel`, `delete`, `deleteMany`, `retry`, connection-first overloads | Short operation names are unambiguous on `Client`; transaction ownership stays visible in the argument list. |
+| `JobUpdate` output parameter, `RecordOutput` | `Client.output`, `WorkContext.output` | Names describe the supported operation; attempt output is buffered until completion. |
+| `JobCompleteTx` | `Client.complete(connection, job)`, `WorkContext.complete(connection)` | Completion keeps `Job` and commits with application changes. Only a running attempt can be completed. A committed completion remains authoritative if the handler subsequently throws. |
+| Context client, cancellation, cancel/snooze errors | `WorkContext.client`, cancellation methods, `cancel`, `discard`, `snooze` | Java has no Go context parameter convention. Control methods end the attempt; virtual-thread interruption supplements cooperative cancellation. |
+| Hooks and middleware | `Extension` | Default methods group related callbacks. JDBC hooks and middleware may throw checked exceptions; insert middleware receives a standard `Callable`. |
+| Retry policy | `RetryPolicy` | A lambda receives the full job and returns `Duration`; the default uses the same error-count backoff and jitter. |
+| `Queues`, queue CRUD/control APIs | `Workers.addQueue/removeQueue`, `Client.queues()` | Runtime concurrency and persisted queue state have separate owners. Reads and mutations support caller transactions. |
+| `Subscribe`, `EventKind`, `Event` | `Workers.subscribe`, `EventKind`, `Event`, `Subscription` | Enum filters and an `AutoCloseable` subscription replace channels. Callbacks are synchronous; queue events include their queue. |
+| Periodic jobs and schedules | `Workers.Builder.periodic`, `Schedule` | Java time types and a functional scheduling interface; cron implementation details are package-private. |
+| Leadership notifications and maintenance config | `requestResign`, worker builder settings | Coordination stays in the database; operational durations use `Duration`. |
+| `rivermigrate` | `Migrator`, `Direction`, `Options`, executable CLI | Named options replace booleans and sentinel target versions. Down defaults to one migration; explicit target zero removes the line. |
+| Driver, pool, transaction integration | `Database`, JDBC `DataSource`/`Connection` | Use existing Java database pools. Database drivers remain optional library dependencies. |
+| Errors and row types | `RiverException.Code`, `Job.State`, nested result records | Exceptions carry failures; enums model closed sets; records keep related values together. Missing rows throw `NOT_FOUND`. |
+
+## Changes made before release
+
+- Removed no-op `Client.close()`. Close `Workers` and the application-owned pool.
+- Added typed retrieval and typed transactional completion, including running-state validation.
+- Unified bulk insertion around typed homogeneous lists and mixed-kind submissions,
+ with both owned and caller-owned transactions.
+- Made insertion results contain persisted JSON rows so a duplicate can safely
+ refer to a different kind or argument schema, matching Go's insertion result.
+- Added savepoint protection to queue controls and connection-based reads, and a
+ public transaction callback overload for grouping operations in a savepoint.
+- Allowed checked failures from insert hooks and middleware; SQL failures retain
+ their cause under `RiverException.Code.DATABASE`.
+- Replaced string event kinds with enums, added subscription filters and queue
+ payloads, and made subscription closure free of checked exceptions.
+- Made repeated stop calls continue to wait after a timeout. A graceful-stop
+ timeout leaves attempts running, as in Go; `stopAndCancel` explicitly escalates it.
+- Changed retry callbacks to receive the full job, enabling kind- and metadata-specific policies.
+- Added named migration options and fluent uniqueness state/kind overrides.
+- Kept destructive reset and retention helpers behind the internal driver seam.
+ `Plugin`, `Client.Driver`, protocol JSON helpers, and explicitly marked companion
+ methods are not stable application extension APIs.
+
+Options copy collections and metadata on construction; `InsertOptions.metadata()`
+and `JobQuery.metadata()` also return copies. Job/result JSON values are snapshots,
+not database-backed objects. Builders are mutable and should be confined to
+configuration code.
+
+## Runtime and option review
+
+- Lists now default to ascending job ID order, matching the Go implementation and
+ keeping pagination independent of rescheduling. Explicit time ordering remains
+ available with `JobQuery.Order.TIME`.
+- Attempt timeouts are supervised independently of fetching and continue during
+ graceful stop. Finished handlers are no longer cancelled while their completion
+ is awaiting a database commit; an interrupted handler cannot abandon completion
+ retries or strand its worker slot.
+- Worker durations and snoozes must fit in signed 64-bit nanoseconds, matching Go's
+ duration range. Invalid worker settings fail during configuration, and an invalid
+ snooze fails the attempt normally instead of trapping it in completion retries.
+- An explicit `rescueAfter` cannot be shorter than `jobTimeout`. The default is one
+ hour, plus a positive explicitly configured job timeout, following Go and Rust.
+- `JobType.uniqueBy("account_id", "address")` selects top-level JSON fields without
+ nested lists. The existing component-list overload handles nested paths and
+ retains literal dots in field names.
+- Reusing a worker builder generates a separate runtime ID on each `start()`;
+ an explicit `id` remains the caller's responsibility to keep unique.
+- Observer and subscriber failures are reported without losing a worker slot,
+ changing the attempt outcome, or preventing delivery to the remaining callbacks.
+ Failures in the error handler fall back to the system logger. These callbacks
+ remain synchronous and must finish quickly.
+- Retry policies run only when the attempt will be retried, matching Go's terminal
+ state handling. Throwing policies and invalid delays are reported and use the
+ default backoff; they cannot prevent completion from being persisted.
+- Query metadata is copied on access so reusable filters cannot change through
+ an exposed JSON node.
+
+## Remaining Go API gaps
+
+These are implementation limits, not claims that Java conventions require a
+smaller feature set:
+
+- OSS periodic definitions cannot be added or removed on a running runtime and
+ currently carry fixed arguments rather than a constructor invoked each time.
+- Events do not yet include Go's per-job timing statistics. Java supplies callback
+ delivery rather than a channel buffer and its associated subscription options.
+- Resumable progress is buffered until the attempt ends; transactional checkpoint
+ helpers are not exposed. The Go logging middleware and worker-test harness have
+ no packaged Java equivalents.
+- Some maintenance configuration is coarser: reindexing accepts an interval,
+ polling policy belongs to a runtime, and insert defaults belong to a job type.
+
+See [differences](DIFFERENCES.md), the [feature inventory](conformance/feature-inventory.json),
+and [validation](VALIDATION.md) for protocol coverage and existing limitations.
diff --git a/java/DIFFERENCES.md b/java/DIFFERENCES.md
new file mode 100644
index 000000000..fd82ea440
--- /dev/null
+++ b/java/DIFFERENCES.md
@@ -0,0 +1,74 @@
+# Intentional differences from Go
+
+- Java 21 is the minimum version. Records, lambdas, `Instant`, `Duration`, JDBC
+ transactions, and virtual threads are the native API; no Go-shaped worker
+ inheritance hierarchy is required.
+- Job kinds are explicit `JobType` values, independent of Java class names.
+ Unique argument fields are explicit paths instead of Go struct tags.
+- JSON uses Jackson 3 and snake_case record properties. Applications must use
+ compatible JSON field names in every implementation sharing a job kind.
+- A `DataSource` and its pool belong to the application. `Database.connect`
+ uses DriverManager without pooling; callers wanting pooled connections supply
+ a data source. River-owned JDBC connections are closed after each transaction.
+- Caller-owned transactions use JDBC savepoints. Hooks run within those
+ savepoints, so a failed River operation leaves the outer transaction usable.
+- Cancellation is cooperative and followed by interruption after the configured
+ stuck threshold. The JVM has no safe way to kill an arbitrary uncooperative
+ thread. Such a handler retains its slot until it returns; handlers should use
+ interruptible I/O or inspect their `WorkContext` cancellation signal.
+- Transaction variants use connection-first overloads instead of `Tx` suffixes.
+ Storage failures are unchecked `RiverException` values with a stable code and
+ original cause; JDBC callbacks may throw checked exceptions.
+- Runtime events use enum kinds and are process-local callbacks. They are not a durable delivery
+ mechanism. A callback should finish quickly, and callback failures are sent to
+ the runtime error handler.
+- Retry policies return relative `Duration` values. Exceptions and null, negative,
+ or overflowing retry delays are reported and fall back to the default backoff.
+- Explicit retries that update a job send a standard insert notification in the
+ same transaction, so workers wake after commit. Go currently leaves discovery
+ of explicitly retried jobs to polling.
+- Short rescued retries become `available` with their retry timestamp and notify
+ the queue once due. This avoids waiting for Java's combined maintenance cadence
+ when `serviceInterval` is long; Go rescues into `retryable` for its scheduler.
+- Java's `serviceInterval` controls elections, maintenance, and queue heartbeats.
+ It must be shorter than the one-day queue retention period; Go reports queue
+ heartbeats separately, every ten minutes by default.
+- Cron schedules follow Go's DST gap/overlap choices. Java uses the JVM's
+ time-zone database, which must agree with Go's for identical results. When a
+ zone skips an entire calendar date, Java advances past it instead of
+ reproducing Go's nonterminating date loop.
+- SQL is kept in dialect-specific resource catalogs. The Go sqlc driver seam is
+ not part of the Java API.
+- SQLite follows Go's JSONB storage and millisecond timestamp representation;
+ PostgreSQL uses microsecond timestamps. These are shared storage contracts,
+ not configurable Java serialization choices.
+
+No intentional differences in stored uniqueness hashes, job state values, reserved
+metadata, migration versions, or notification payloads are allowed. Any remaining
+conformance failure in these areas is a defect, not an API design choice.
+
+The following are current API limitations, not differences required by Java
+conventions. The [API review](API.md) maps the complete application surface.
+
+OSS periodic registrations and polling policies belong to a worker configuration;
+start a new runtime to change its periodic definitions, or use separate runtimes
+for queues with different polling intervals. Fetch wakeups are coalesced without
+a separate cooldown setting. Reindexing supports the daily UTC default or a
+duration override, rather than an arbitrary cron expression. Job-type defaults
+replace process-wide insert defaults.
+
+Insertion is independent of worker registration so a client can enqueue jobs
+owned by another language. New job kinds always use Go's validated kind syntax;
+the legacy bypass is omitted. `Workers.stop()` (also called by `close()`) and
+`Workers.stopAndCancel()` wait for attempts to finish instead of exposing a
+Go-style stopped channel. `Workers.Builder.stopTimeout` configures their waits.
+
+The [feature inventory](conformance/feature-inventory.json) accounts for all 304
+entries in the pinned upstream inventory, including internal driver details
+that are not application APIs.
+
+Java has no bundled `riverlog` middleware or `rivertest.Worker` harness. Applications
+can write the shared log metadata format explicitly and test with JUnit against
+an isolated database. Resumable steps buffer progress until the attempt ends;
+Go's transactional resumable checkpoint helpers are not yet exposed. See the
+[README examples](README.md#features) for these API limitations.
diff --git a/java/Makefile b/java/Makefile
new file mode 100644
index 000000000..66c39bc8c
--- /dev/null
+++ b/java/Makefile
@@ -0,0 +1,19 @@
+.PHONY: test test/conformance test/postgres test/sqlite test/tools lint lint/tools format package
+test:
+ $(MAKE) -C .. test/java
+test/conformance:
+ $(MAKE) -C .. test/java/conformance
+test/postgres:
+ $(MAKE) -C .. test/java/postgres
+test/sqlite:
+ $(MAKE) -C .. test/java/sqlite
+test/tools:
+ $(MAKE) -C .. test/java/tools
+lint:
+ $(MAKE) -C .. lint/java
+lint/tools:
+ $(MAKE) -C .. lint/java/tools
+format:
+ mvn spotless:apply
+package:
+ $(MAKE) -C .. check/java/package
diff --git a/java/README.md b/java/README.md
new file mode 100644
index 000000000..e6914eca5
--- /dev/null
+++ b/java/README.md
@@ -0,0 +1,788 @@
+# River for Java
+
+Prerelease Java 21 client and worker runtime for River's PostgreSQL and SQLite
+schemas. Jobs are ordinary River jobs: another language can insert, cancel,
+retry, or work them using the same database.
+
+## Quick start
+
+Define job arguments as a record, give the job a stable kind, and register a
+worker lambda. This complete `Example.java` starts a client, inserts a job,
+waits for its completion, and stops the workers. Set `DATABASE_URL` to a
+PostgreSQL URL or a file-backed SQLite JDBC URL.
+
+```java
+import com.riverqueue.*;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.TimeUnit;
+
+public class Example {
+ record SendEmail(String address, String subject) {}
+
+ public static void main(String[] args) throws Exception {
+ var database = Database.connect(System.getenv("DATABASE_URL"));
+ var client = new Client(database);
+ new Migrator(database).migrate();
+
+ var completed = new CompletableFuture();
+ try (var workers = client.workers()
+ .queue("default", 10)
+ .add(JobType.of("send_email", SendEmail.class), context -> {
+ System.out.println("Sending " + context.args().subject() + " to " + context.args().address());
+ context.output("sent");
+ })
+ .start();
+ var subscription = workers.subscribe(completed::complete, Workers.EventKind.JOB_COMPLETED)) {
+ client.insert(JobType.of("send_email", SendEmail.class),
+ new SendEmail("hello@example.com", "Welcome"));
+ System.out.println("Completed job " + completed.get(30, TimeUnit.SECONDS).job().id());
+ }
+ }
+}
+```
+
+Run this demonstration against a development database with no other workers.
+The worker prints a message; replace its body with your mail service. In a
+service, keep `Workers` open until the application stops. Each active job runs
+on a virtual thread; queue limits bound concurrent jobs. In production, run
+migrations as a deployment step using the CLI below.
+
+## Installation
+
+The Maven artifact is `com.riverqueue:river:0.48.0-alpha.1`. To build from source,
+run `make generate/fixtures` from the repository root, then `mvn install` from
+`java/`, with `JAVA_HOME` pointing to JDK 21 or 25. Source tests need Go to generate
+their reference fixtures; applications do not need Go. The pinned formatter does
+not run on JDK 27. Add the JDBC driver for your database to your application.
+Jackson 3 is included transitively.
+
+```xml
+
+ com.riverqueue
+ river
+ 0.48.0-alpha.1
+
+```
+
+Job kinds are explicit, stable wire names independent of Java class names.
+Record properties use snake_case in JSON by default, so `accountId` becomes
+`account_id`. Match kinds and JSON fields across languages sharing a job.
+
+Use a pooled `DataSource` in applications. `Database.connect(url)` is a small
+DriverManager convenience for scripts and tests and opens a physical connection
+for each operation. River never closes an application-owned data source. An
+[insert-only client](https://riverqueue.com/docs/insert-only-clients) is simply a
+`Client` instance without a worker runtime; inserting needs no start or stop call.
+
+## Migration CLI
+
+The `river-cli` executable JAR includes PostgreSQL and SQLite JDBC drivers and
+runs with Java 21 or newer. It needs no Go installation or application classpath.
+Build it as described under [Installation](#installation), then run:
+
+```sh
+RIVER_CLI=cli/target/river-cli-0.48.0-alpha.1-all.jar
+export DATABASE_URL=postgres://localhost/myapp
+# For SQLite: export DATABASE_URL=jdbc:sqlite:/absolute/path/myapp.db
+
+java -jar "$RIVER_CLI" --help
+java -jar "$RIVER_CLI" migrate-list
+java -jar "$RIVER_CLI" migrate-up --dry-run
+java -jar "$RIVER_CLI" migrate-up
+```
+
+The Maven executable artifact is `com.riverqueue:river-cli` with version
+`0.48.0-alpha.1` and classifier `all`. Copy that single JAR into your deployment image to run
+migrations before starting workers. `--database-url URL` overrides `DATABASE_URL`.
+`--schema background_jobs` selects a PostgreSQL schema; an actual migration up
+creates it if needed. Listing and dry runs do not create schemas or migration tables.
+
+| Command | Behavior |
+| --- | --- |
+| `migrate-up` | Applies all pending migrations, committing each separately. |
+| `migrate-down` | Reverses one migration by default; may remove tables and data. |
+| `migrate-list` | Lists bundled versions and whether each is applied. |
+| `migrate-get` | Exports canonical SQL without connecting to a database. |
+| `version` | Prints the CLI version; `--version` works too. |
+
+Use `--target-version N` to stop at a version, `--max-steps N` to limit the count,
+and `--dry-run --show-sql` to inspect a planned change. On `migrate-down`, target
+version `0` explicitly removes the entire migration line. A zero `--max-steps`
+uses the command's default: all pending up migrations or one down migration.
+
+```sh
+java -jar "$RIVER_CLI" migrate-down --max-steps 2 --dry-run --show-sql
+java -jar "$RIVER_CLI" migrate-up --target-version 8
+java -jar "$RIVER_CLI" migrate-get --version 3 --up > river_3.up.sql
+java -jar "$RIVER_CLI" migrate-get --driver sqlite --all --up > river_sqlite.up.sql
+```
+
+SQL export defaults to PostgreSQL without a database URL. Use `--driver sqlite`
+for SQLite, or select the dialect through the URL. Export accepts comma-separated
+`--version` values or `--all`; `--exclude-version 1` excludes migration history
+setup when exporting for another migration framework. Exit codes are `0` for
+success, `1` for an operation failure, and `2` for invalid arguments.
+
+The [migration history](https://riverqueue.com/docs/migrations) is shared with Go
+and the other ports. Use the CLI version matching your Java library. The Java
+API remains available for embedded use:
+
+```java
+var result = new Migrator(database).migrate();
+System.out.println("Applied migrations: " + result.applied());
+var preview = new Migrator(database).migrate(Migrator.Direction.DOWN,
+ Migrator.Options.defaults().maxSteps(2).dryRun(true));
+```
+
+## Transactional enqueueing
+
+[Enqueue in the application's transaction](https://riverqueue.com/docs/transactional-enqueueing)
+so the job and application writes commit or roll back together. Pass a JDBC
+connection with auto-commit disabled. River uses savepoints and never commits,
+rolls back the outer transaction, or closes that connection. Insert hooks and
+middleware run inside the same transaction.
+
+```java
+try (var connection = applicationDataSource.getConnection()) {
+ connection.setAutoCommit(false);
+ try {
+ accounts.create(connection, account);
+ client.insert(connection, email, new SendEmail(account.email(), "Welcome"));
+ connection.commit();
+ } catch (Exception failure) {
+ connection.rollback();
+ throw failure;
+ }
+}
+```
+
+`client.transaction(connection -> ...)` manages a new transaction for River-owned
+operations. `client.transaction(existingConnection, connection -> ...)` groups
+multiple operations under one savepoint in an existing transaction. Neither
+`Client` nor `Database` needs closing; the application owns its pool. Workers should be [idempotent](https://riverqueue.com/docs/reliable-workers):
+process failure can cause an attempt to run again after its external effects
+have succeeded.
+
+## Inserting many jobs
+
+[Bulk insertion](https://riverqueue.com/docs/inserting-many-jobs) inserts a list atomically.
+For an existing JDBC transaction, use the connection-first overload of either
+form below. Homogeneous batches take a common job type. Mixed batches use
+submissions with their own kind and options. Both return insertion results
+containing the persisted `Job` values.
+
+```java
+var inserted = client.insertMany(email, List.of(
+ new SendEmail("one@example.com", "Welcome"),
+ new SendEmail("two@example.com", "Welcome")));
+
+record GenerateReport(long accountId) {}
+var report = JobType.of("generate_report", GenerateReport.class);
+var mixed = client.insertMany(List.of(
+ email.submission(new SendEmail("one@example.com", "Welcome")),
+ report.submission(new GenerateReport(42), InsertOptions.builder().queue("reports").build())));
+```
+
+## Reading typed jobs
+
+`client.get(id)` returns `Job` because an ID does not identify a Java
+argument class. Supply a job type when it is known:
+
+```java
+Job job = client.get(jobId, email);
+System.out.println(job.args().address());
+```
+
+Typed retrieval checks the kind and decodes the arguments.
+`context.complete(connection)` also retains its `Job` argument type.
+Insertion results contain `Job` because uniqueness can return an
+existing job with another kind or argument schema. The returned arguments are
+the persisted values, including any additional fields written by another language.
+See the [API comparison](API.md) for the Go mapping and prerelease API changes.
+
+Lists default to ascending job ID order, matching Go. Keep the filters and order
+unchanged when continuing from a page's cursor:
+
+```java
+var query = JobQuery.builder().kinds("send_email").limit(50).build();
+var first = client.list(query);
+if (!first.jobs().isEmpty()) {
+ var next = client.list(query.after(first.cursor()));
+}
+```
+
+Use `order(JobQuery.Order.TIME)` for ordering by the timestamp appropriate to the
+selected state, or `SCHEDULED_AT` or `FINALIZED_AT` for an explicit timestamp.
+
+## Job retries
+
+A thrown exception triggers [retries](https://riverqueue.com/docs/job-retries)
+until `maxAttempts` is exhausted, when the job becomes discarded. The default
+policy uses River's quartic backoff with jitter. Configure a different policy
+on the worker builder, or explicitly retry an existing job with `client.retry(id)`.
+An explicit retry that updates a job notifies its queue after commit. Use
+`client.retry(connection, id)` to make the retry part of an existing transaction.
+A retry policy receives the full job snapshot before the current failure is
+appended to `errors`; use `job.errors().size() + 1` to include that failure.
+It runs only when another attempt is possible. Return a nonnegative `Duration`;
+if the policy throws, returns null or a negative delay, or overflows the retry
+time, River reports the problem and uses its default backoff.
+
+```java
+var retryingEmail = email.withDefaults(InsertOptions.builder().maxAttempts(10).build());
+var workerConfig = client.workers()
+ .queue("default", 10)
+ .retryPolicy(job -> Duration.ofSeconds(Math.min(300, 5L * (job.errors().size() + 1))))
+ .add(retryingEmail, context ->
+ mailer.send(context.args().address(), context.args().subject()));
+client.insert(retryingEmail, new SendEmail("hello@example.com", "Welcome"));
+```
+
+Call `workerConfig.start()` and keep the runtime open as in the first example.
+Job-type defaults apply when inserting through that `JobType`; another producer
+must set its own compatible options.
+
+## Features
+
+These sections follow the non-Pro entries in the River documentation's
+[Features sidebar](https://riverqueue.com/docs). Snippets are independent and
+reuse `client` and `database` from the quick start, with imports from
+`com.riverqueue`, `java.time`, and `java.util`. Handler examples configure a builder;
+start it after registering your handlers. `mailer`, `accounts`, and `deliveries`
+are application services, `applicationDataSource` is your JDBC connection pool,
+`application.awaitStop()` waits for your service to stop, and `jobId` is an
+existing job's ID.
+
+```java
+var email = JobType.of("send_email", SendEmail.class);
+var workerConfig = client.workers().queue("default", 10);
+```
+
+### Cancelling jobs
+
+[Cancel](https://riverqueue.com/docs/cancelling-jobs) an enqueued or running job
+with `client.cancel(id)`. Running attempts receive a cooperative cancellation
+signal. A handler can cancel itself permanently with `context.cancel(reason)`;
+this ends the attempt without further retries.
+
+```java
+client.cancel(jobId);
+
+workerConfig.add(email, context -> {
+ context.checkCancelled();
+ if (context.args().address().isBlank()) context.cancel("Missing address");
+ mailer.send(context.args().address(), context.args().subject());
+});
+```
+
+Long-running handlers should check cancellation between operations and use
+interruptible I/O. `context.awaitCancellation()` waits for the signal.
+As in Go, a normal return reports success even if cancellation or a timeout was
+requested. Call `context.checkCancelled()` when abandoning unfinished work.
+
+### Getting the client within workers
+
+[`context.client()`](https://riverqueue.com/docs/context-client) gives a handler
+its client, including configured extensions. Use its transaction overloads when
+enqueueing follow-up work that must commit with other writes or job completion.
+
+```java
+record Welcome(String address) {}
+var welcome = JobType.of("welcome", Welcome.class);
+workerConfig.add(welcome, context ->
+ context.client().insert(email, new SendEmail(context.args().address(), "Welcome")));
+```
+
+### Error and panic handling
+
+Java exceptions and uncaught `Error`s are handled as failed attempts; River
+persists failures and retries according to policy. For
+[application error reporting](https://riverqueue.com/docs/error-handling), attach
+an `Extension.afterWork` hook. `Workers.Builder.errorHandler` separately receives
+runtime failures, such as database or subscriber errors.
+Failures in observers or event subscribers are reported without changing the
+job's outcome or preventing other callbacks from running. If the runtime error
+handler itself throws, River logs both failures and continues.
+
+```java
+var monitoredClient = client.withExtension(new Extension() {
+ @Override
+ public void afterWork(WorkContext> context, Throwable failure) {
+ if (failure != null) {
+ System.getLogger("jobs").log(System.Logger.Level.ERROR,
+ "Job " + context.job().id() + " failed", failure);
+ }
+ }
+});
+var workerConfig = monitoredClient.workers()
+ .queue("default", 10)
+ .errorHandler(failure ->
+ System.getLogger("river").log(System.Logger.Level.ERROR, "Runtime failure", failure))
+ .add(email, context -> mailer.send(context.args().address(), context.args().subject()));
+```
+
+Hooks also see the control exceptions used for snoozing and cancellation. Keep
+reporting hooks fast and avoid throwing from them.
+
+### Job-persisted logging
+
+[Job logs](https://riverqueue.com/docs/job-logging) live in `river:log` metadata
+as an array of `{attempt, log}` objects. Java has no bundled equivalent of Go's
+`riverlog` middleware yet. You can explicitly write that shared format; this
+example appends a fixed message and retains the last ten attempt entries.
+
+```java
+workerConfig.add(email, context -> {
+ var logs = new ArrayList