From 81d00464598fc3ddc04fc85e8b4da4e790a650dc Mon Sep 17 00:00:00 2001 From: Brandur Date: Tue, 6 Oct 2026 12:36:51 -0500 Subject: [PATCH] Bring TypeScript conformance suite in line with Rust's Codex noticed as it was reviewing #1451 that a more exhaustive set of checks are applied to Rust than TypeScript when it comes to conformance. This is probably due to more recent changes in the conformance suite that hadn't yet made their way to TypeScript. Here, do one more pass to bring TypeScript as inline with Rust as we can make it. --- Makefile | 2 +- js/docs/development.md | 8 +- js/src/conformance.test.ts | 344 ++++++++++++++++++ js/src/runtime/completion-command.test.ts | 1 + .../notification-pump.conformance.test.ts | 196 ++++++++++ 5 files changed, 548 insertions(+), 3 deletions(-) create mode 100644 js/src/conformance.test.ts create mode 100644 js/src/runtime/notification-pump.conformance.test.ts 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); +}