diff --git a/Makefile b/Makefile index a54a2b176..cb3a676cb 100644 --- a/Makefile +++ b/Makefile @@ -159,7 +159,7 @@ test/js: generate/fixtures .PHONY: test/js/conformance test/js/conformance: ## Run JavaScript tests that check Go-generated conformance fixtures test/js/conformance: generate/fixtures - pnpm -C js exec vitest run src/cron.test.ts src/runtime/completion-command.test.ts + pnpm -C js exec vitest run src/conformance.test.ts src/cron.test.ts src/runtime/completion-command.test.ts src/runtime/notification-pump.conformance.test.ts # Integration tests use TEST_DATABASE_URL (default # postgres://localhost:5432/river_test), migrated with diff --git a/js/docs/development.md b/js/docs/development.md index 894cb94f5..e0e99de37 100644 --- a/js/docs/development.md +++ b/js/docs/development.md @@ -197,12 +197,16 @@ formatting, licenses, packed archives and examples, unit tests on Node 26.0.0 and the current Node 26 release, and integration tests on PostgreSQL 14 through 18. -Unit tests compare cron schedules and snooze counting with fixtures that +Unit tests compare unique keys, protocol values, notification dispatch, +retry timing, cron schedules, and snooze counting with fixtures that River's Go implementation generates into `conformance/testdata`, the same files the Rust port reads. They aren't committed: `make test/js` generates them first, so Go is needed to run the unit tests, and `pnpm run test` needs a prior `make generate/fixtures` from the repository root. A missing fixture -fails its test. +fails its test. `make test/js/conformance` runs just these checks, including +both drivers' notification adapters without a PostgreSQL server. The two +raw JSON unique-key cases involving duplicate keys or integer-key insertion +order are Rust-only because JavaScript objects cannot preserve them. ## Preparing a release diff --git a/js/src/conformance.test.ts b/js/src/conformance.test.ts new file mode 100644 index 000000000..2ee61a886 --- /dev/null +++ b/js/src/conformance.test.ts @@ -0,0 +1,344 @@ +import { Buffer } from "node:buffer"; +import { readFile } from "node:fs/promises"; +import { fileURLToPath } from "node:url"; + +import { describe, expect, it, onTestFinished, vi } from "vitest"; + +import { testSqliteMemory } from "../driver/sqlite/src/driver.js"; +import { createMigrator } from "../migrate/src/index.js"; +import { buildUniqueKey, type Client } from "./client.js"; +import { decodeAttemptError, decodeJobState } from "./driver-codecs.js"; +import { ValidationError } from "./errors.js"; +import type { UniqueOptions } from "./insert-options.js"; +import { UNIQUE_INSERT_NONCE_KEY } from "./internal/postgres-capabilities.js"; +import { defineJob } from "./job-definition.js"; +import { + JOB_STATE, + jobToJsonValue, + type JobRow, + type JobState, +} from "./job.js"; +import { + jsonNumberToBigInt, + parseJson, + stringifyJson, + type ExactJsonNumber, + type JsonObject, +} from "./json.js"; +import { buildPeriodicInsert, periodicJob } from "./periodic.js"; +import { Resumable } from "./resumable.js"; +import { defaultNextRetry } from "./runtime/completion-command.js"; +import { + uniqueBitmaskFromStates, + uniqueBitmaskToStates, +} from "./unique-bitmask.js"; + +interface UniqueKeyCase { + readonly args: JsonObject; + readonly expected_error?: string; + readonly expected_sha256?: string; + readonly expected_state_mask: number; + readonly kind: string; + readonly name: string; + readonly now: string; + readonly options: { + readonly by_args: boolean; + readonly by_period_nanos: number; + readonly by_queue: boolean; + readonly by_state?: readonly JobState[]; + readonly exclude_kind: boolean; + }; + readonly queue: string; + readonly scheduled_at: string | null; + readonly selected_unique_components?: readonly (readonly string[])[]; +} + +/** Keep large argument integers and retry bounds exact when reading Go JSON. */ +async function readFixture(name: string) { + const url = new URL( + `../../conformance/testdata/${name}.json`, + import.meta.url + ); + try { + return parseJson(await readFile(url, "utf8")); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") { + throw new Error( + `missing conformance fixture ${fileURLToPath(url)}; run \`make generate/fixtures\` from the repository root`, + { cause: error } + ); + } + throw error; + } +} + +const uniqueKeys = (await readFixture("unique_keys")) as unknown as { + readonly cases: readonly UniqueKeyCase[]; + readonly typed_only_cases: readonly UniqueKeyCase[]; +}; +const protocol = (await readFixture("protocol_values")) as unknown as { + readonly attempt_error: JsonObject; + readonly job_states: readonly { state: string; unique_bit: number }[]; + readonly metadata_keys: { + readonly output: string; + readonly periodic_job_id: string; + readonly rescue_count: string; + readonly resumable_cursor: string; + readonly resumable_step: string; + readonly unique_nonce: string; + }; + readonly retry_cases: readonly { + error_count: number; + job_id: ExactJsonNumber | number; + max_delay_ns: ExactJsonNumber | number; + min_delay_ns: ExactJsonNumber | number; + now: string; + }[]; +}; + +describe("Go unique-key fixtures", () => { + it.each(uniqueKeys.cases)("$name", (fixture) => { + const now = Temporal.Instant.from(fixture.now); + const options: UniqueOptions = { + ...(fixture.options.by_args + ? { + byArgs: fixture.selected_unique_components?.length + ? fixture.selected_unique_components.map((segments) => + segments + .map((segment) => segment.replace(/[.\\]/g, "\\$&")) + .join(".") + ) + : true, + } + : {}), + ...(fixture.options.by_period_nanos > 0 + ? { byPeriod: { nanoseconds: fixture.options.by_period_nanos } } + : {}), + ...(fixture.options.by_state === undefined + ? {} + : { byState: fixture.options.by_state }), + byQueue: fixture.options.by_queue, + excludeKind: fixture.options.exclude_kind, + }; + const build = () => + buildUniqueKey( + { + args: fixture.args, + kind: fixture.kind, + queue: fixture.queue, + // Client insertion resolves an absent scheduled time to now. + scheduledAt: + fixture.scheduled_at === null + ? now + : Temporal.Instant.from(fixture.scheduled_at), + }, + options + ); + + if (fixture.expected_error !== undefined) { + expect(fixture.expected_error).toBe("rejected"); + expect(build).toThrow(ValidationError); + } else { + const [key, states] = build(); + expect(Buffer.from(key).toString("hex")).toBe(fixture.expected_sha256); + expect(Number.parseInt(uniqueBitmaskFromStates(states), 2)).toBe( + fixture.expected_state_mask + ); + } + }); + + it("keeps the unsupported raw JSON cases explicit", () => { + // JavaScript objects cannot preserve duplicate keys or insertion order + // for integer-like keys. Rust can test these via its raw JSON API. + expect(uniqueKeys.typed_only_cases.map(({ name }) => name).sort()).toEqual([ + "typed_duplicate_top_level_keys", + "typed_integer_like_map_keys", + ]); + }); +}); + +describe("Go protocol fixtures", () => { + it("decodes every state and preserves its uniqueness bit in both directions", () => { + expect(protocol.job_states.map(({ state }) => state).sort()).toEqual( + Object.values(JOB_STATE).sort() + ); + for (const { state, unique_bit: bit } of protocol.job_states) { + const decoded = decodeJobState(state); + expect(Number.parseInt(uniqueBitmaskFromStates([decoded]), 2)).toBe(bit); + expect(uniqueBitmaskToStates(bit)).toEqual([decoded]); + } + }); + + it("reads and writes Go's attempt error without losing fields or timestamp precision", () => { + const error = decodeAttemptError(stringifyJson(protocol.attempt_error)); + // Fixed semantic expectations also catch a renamed field being silently ignored. + expect(error).toEqual({ + at: Temporal.Instant.from("2026-01-02T03:04:05.6789Z"), + attempt: 3, + error: 'worker failed: escaped "detail"', + trace: "frame one\nframe two", + }); + expect(jobToJsonValue({ ...job(), errors: [error] }).errors).toEqual([ + protocol.attempt_error, + ]); + }); + + // Go and JS have different random generators. Check the shared bounds, + // including both jitter extremes and the maximum time.Duration cap. + describe.each(protocol.retry_cases)( + "retry with $error_count errors", + (fixture) => { + it.each([0, 0.5, 1 - Number.EPSILON])( + "jitter %s stays inside Go's bounds", + (random) => { + const now = Temporal.Instant.from(fixture.now); + const row = { + ...job(), + id: jsonNumberToBigInt(fixture.job_id), + errors: Array.from({ length: fixture.error_count - 1 }, () => + decodeAttemptError(stringifyJson(protocol.attempt_error)) + ), + }; + const delay = + defaultNextRetry(row, now, () => random).epochNanoseconds - + now.epochNanoseconds; + expect(delay).toBeGreaterThanOrEqual( + jsonNumberToBigInt(fixture.min_delay_ns) + ); + expect(delay).toBeLessThanOrEqual( + jsonNumberToBigInt(fixture.max_delay_ns) + ); + } + ); + } + ); + + it("uses Go's periodic job ID and unique insert nonce keys", async () => { + const periodic = periodicJob({ + args: {}, + every: { hours: 1 }, + id: "conformance_periodic", + job: defineJob({ kind: "periodic" }), + }); + const insert = await buildPeriodicInsert({ + job: periodic, + scheduledAt: job().scheduledAt, + }); + expect(insert?.options?.metadata).toEqual({ + periodic: true, + [protocol.metadata_keys.periodic_job_id]: "conformance_periodic", + }); + expect(UNIQUE_INSERT_NONCE_KEY).toBe(protocol.metadata_keys.unique_nonce); + }); + + it("resumes from and writes Go's resumable step and cursor keys", async () => { + const keys = protocol.metadata_keys; + const resumable = new Resumable({} as Client, { + ...job(), + metadata: { + [keys.resumable_cursor]: { process: { offset: 2 } }, + [keys.resumable_step]: "process", + }, + }); + const before = vi.fn(); + await resumable.step("before", before); + expect(before).not.toHaveBeenCalled(); + await expect( + resumable.stepWithCursor("process", (cursor) => { + expect(cursor).toEqual({ offset: 2 }); + resumable.setCursor({ offset: 3 }); + throw new Error("retry"); + }) + ).rejects.toThrow("failed"); + expect(resumable.finish(true).metadata).toEqual({ + [keys.resumable_cursor]: { process: { offset: 3 } }, + [keys.resumable_step]: "process", + }); + }); + + it("persists output and increments Go's rescue counter", async () => { + const driver = testSqliteMemory(); + const database = driver.connect(); + onTestFinished(() => { + database.close(); + driver.close(); + }); + await createMigrator({ database }).migrateUp(); + const now = Temporal.Now.instant(); + const inserted = await driver.jobInsert({ + args: {}, + kind: "conformance", + metadata: { [protocol.metadata_keys.rescue_count]: 2 }, + scheduledAt: now, + }); + expect(inserted.job.metadata[protocol.metadata_keys.unique_nonce]).toEqual( + expect.any(String) + ); + await driver.jobClaim({ + attemptedBy: "client", + kinds: [], + queues: [{ limit: 1, name: "default" }], + }); + const rescued = await driver.jobRescueMany( + [ + { + error: decodeAttemptError(stringifyJson(protocol.attempt_error)), + id: inserted.job.id, + scheduledAt: now, + state: "retryable", + }, + ], + Temporal.Now.instant().add({ seconds: 1 }) + ); + expect(rescued[0]?.metadata).toMatchObject({ + [protocol.metadata_keys.rescue_count]: 3, + }); + await driver.jobSchedule({ limit: 1, now }); + const claimed = await driver.jobClaim({ + attemptedBy: "client", + kinds: [], + queues: [{ limit: 1, name: "default" }], + }); + const completed = await driver.jobCompleteMany([ + { + attempt: claimed.jobs[0]!.attempt, + attemptedBy: "client", + error: null, + finalizedAt: now, + id: inserted.job.id, + kind: "complete", + output: { ok: true }, + outputSet: true, + scheduledAt: null, + }, + ]); + expect(completed[0]?.job?.metadata).toMatchObject({ + [protocol.metadata_keys.output]: { ok: true }, + [protocol.metadata_keys.rescue_count]: 3, + }); + }); +}); + +function job(): JobRow { + const now = Temporal.Instant.from("2026-01-02T03:04:05.6789Z"); + return { + args: {}, + attempt: 1, + attemptedAt: now, + attemptedBy: ["client"], + createdAt: now, + errors: [], + finalizedAt: null, + id: 42n, + kind: "conformance", + maxAttempts: 25, + metadata: {}, + priority: 1, + queue: "default", + scheduledAt: now, + state: "running", + tags: [], + uniqueKey: null, + uniqueStates: null, + }; +} diff --git a/js/src/runtime/completion-command.test.ts b/js/src/runtime/completion-command.test.ts index 898df5fa4..cc6861c9c 100644 --- a/js/src/runtime/completion-command.test.ts +++ b/js/src/runtime/completion-command.test.ts @@ -49,6 +49,7 @@ describe("completionCommand", () => { const golden = parseJson(await readFixture(GOLDENS)) as unknown as { 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 }) => { diff --git a/js/src/runtime/notification-pump.conformance.test.ts b/js/src/runtime/notification-pump.conformance.test.ts new file mode 100644 index 000000000..05b7c2c88 --- /dev/null +++ b/js/src/runtime/notification-pump.conformance.test.ts @@ -0,0 +1,196 @@ +import { EventEmitter } from "node:events"; +import { readFile } from "node:fs/promises"; +import { fileURLToPath } from "node:url"; + +import type { Pool } from "pg"; +import { describe, expect, it, onTestFinished, vi } from "vitest"; + +import { PgDatabase } from "../../driver/pg/src/database.js"; +import { runtimeNotificationSubscribe } from "../../driver/pg/src/sql/notify.js"; +import { testSqliteMemory } from "../../driver/sqlite/src/driver.js"; +import type { RuntimeDriver, RuntimeNotification } from "../driver.js"; +import { parseJson, stringifyJson, type JsonObject } from "../json.js"; +import type { AttemptRunner } from "./attempt-runner.js"; +import type { RuntimeContext } from "./context.js"; +import { NotificationPump } from "./notification-pump.js"; +import type { QueueProducer } from "./queue-producer.js"; + +interface NotificationFixture { + readonly name: string; + readonly payload: JsonObject; + readonly topic: string; +} + +const url = new URL( + "../../../conformance/testdata/protocol_values.json", + import.meta.url +); +const protocol = parseJson( + await readFile(url, "utf8").catch((error: unknown) => { + if ((error as NodeJS.ErrnoException).code === "ENOENT") { + throw new Error( + `missing conformance fixture ${fileURLToPath(url)}; run \`make generate/fixtures\` from the repository root`, + { cause: error } + ); + } + throw error; + }) +) as unknown as { + readonly notifications: readonly NotificationFixture[]; + readonly topics: Readonly>; +}; + +it("covers every Go notification action", () => { + expect(protocol.notifications.map(({ name }) => name).sort()).toEqual([ + "cancel", + "insert", + "metadata_changed", + "pause", + "request_resign", + "resigned", + "resume", + ]); +}); + +describe.each(["postgres", "sqlite"] as const)( + "Go notifications through %s", + (backend) => { + describe.each(["client-1", "observer"])("client %s", (clientId) => { + it.each(protocol.notifications)("dispatches $name", async (fixture) => { + const stop = new AbortController(); + const tasks: Promise[] = []; + onTestFinished(async () => { + stop.abort(); + await Promise.allSettled(tasks); + }); + const subscribe = + backend === "postgres" + ? postgresStream(fixture) + : sqliteStream(fixture); + const context = { + claimSignal: stop.signal, + clientId, + driver: { + async *runtimeNotificationSubscribe(topics, signal, ready) { + for await (const notification of subscribe( + topics, + signal, + ready + )) { + yield notification; + // The pump has handled this one fixture. Stop without timers + // or resubscribing; failures before delivery reject start(). + stop.abort(); + } + }, + } satisfies Pick, + guard: (task: Promise) => task, + trackTask: (task: Promise) => { + tasks.push(task); + }, + } as unknown as RuntimeContext; + + const effects: unknown[][] = []; + const producer = { + refresh: (queue: string) => { + effects.push(["refresh", queue]); + }, + refreshAll: () => { + effects.push(["refreshAll"]); + }, + wake: (queue: string) => { + effects.push(["wake", queue]); + }, + wakeAll: () => { + effects.push(["wakeAll"]); + }, + } as unknown as QueueProducer; + const runner = { + cancelAttempt: (id: bigint, owner: string) => { + effects.push(["cancel", id, owner]); + }, + } as unknown as AttemptRunner; + const pump = new NotificationPump(context, { + leaderResigned: () => { + effects.push(["leaderResigned"]); + }, + producer, + resignLeadership: () => { + effects.push(["resignLeadership"]); + return undefined; + }, + runner, + }); + + await pump.start(true); + await Promise.all(tasks); + + // Semantic expectations are independent of the fixture's payload. + // Checking all effects also rules out unrelated cancellations/wakeups. + const expected: Record = { + cancel: [ + ["cancel", 42n, clientId], + ["refresh", "priority"], + ], + insert: [["wake", "priority"]], + metadata_changed: [["refresh", "priority"]], + pause: [["refresh", "priority"]], + request_resign: [["resignLeadership"]], + resigned: clientId === "client-1" ? [] : [["leaderResigned"]], + resume: [["refresh", "priority"]], + }; + expect(effects).toEqual(expected[fixture.name]); + }); + }); + } +); + +type Subscribe = NonNullable; + +/** Exercise LISTEN channel naming and decoding with only the socket faked. */ +function postgresStream(fixture: NotificationFixture): Subscribe { + const client = Object.assign(new EventEmitter(), { + query: vi.fn().mockResolvedValue({ rows: [] }), + release: vi.fn(), + }); + const pool = { + connect: async () => client, + idleCount: 0, + totalCount: 0, + } as unknown as Pool; + const database = new PgDatabase(pool, "conformance"); + return (topics, signal, ready) => + runtimeNotificationSubscribe(database, topics, signal, () => { + expect(client.query.mock.calls.map(([sql]) => sql).sort()).toEqual( + Object.values(protocol.topics) + .map((topic) => `LISTEN "conformance.${topic}"`) + .sort() + ); + ready(); + client.emit("notification", { + channel: `conformance.${fixture.topic}`, + payload: stringifyJson(fixture.payload), + }); + }); +} + +/** Exercise outbox topic filtering and decoding with database reads faked. */ +function sqliteStream(fixture: NotificationFixture): Subscribe { + const driver = testSqliteMemory(); + onTestFinished(() => driver.close()); + vi.spyOn(driver, "notificationLastId").mockResolvedValue(0n); + vi.spyOn(driver, "notificationPoll").mockImplementation(async (params) => { + expect([...(params?.topics ?? [])].sort()).toEqual( + Object.values(protocol.topics).sort() + ); + return [ + { + createdAt: Temporal.Now.instant(), + id: 1n, + topic: fixture.topic, + payload: stringifyJson(fixture.payload), + }, + ]; + }); + return driver.runtimeNotificationSubscribe.bind(driver); +}