-
Notifications
You must be signed in to change notification settings - Fork 379
fix(world-postgres): preserve queue error metadata #4117
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,5 @@ | ||
| --- | ||
| '@workflow/world-postgres': patch | ||
| --- | ||
|
|
||
| Preserve error names, messages, stacks, and nested causes in Graphile Worker log metadata. |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,61 @@ | ||
| import { describe, expect, it } from 'vitest'; | ||
| import { serializeGraphileMeta } from './queue.js'; | ||
|
|
||
| describe('serializeGraphileMeta', () => { | ||
| it('expands the non-enumerable Error fields JSON.stringify drops', () => { | ||
| const error = new Error('delivery failed'); | ||
| expect(JSON.parse(JSON.stringify({ error }))).toEqual({ error: {} }); | ||
|
|
||
| const meta = JSON.parse(serializeGraphileMeta({ error, jobId: '7' })); | ||
| expect(meta).toMatchObject({ | ||
| jobId: '7', | ||
| error: { | ||
| name: 'Error', | ||
| message: 'delivery failed', | ||
| stack: expect.stringContaining('Error: delivery failed'), | ||
| }, | ||
| }); | ||
| expect(meta.error).not.toHaveProperty('cause'); | ||
| }); | ||
|
|
||
| it('keeps enumerable diagnostics and recurses through cause and AggregateError', () => { | ||
| const socket = Object.assign(new Error('other side closed'), { | ||
| code: 'UND_ERR_SOCKET', | ||
| }); | ||
| const fetchFailed = new TypeError('fetch failed', { cause: socket }); | ||
| const aggregate = new AggregateError([fetchFailed], 'all attempts failed'); | ||
|
|
||
| const meta = JSON.parse(serializeGraphileMeta({ error: aggregate })); | ||
| expect(meta.error).toMatchObject({ | ||
| name: 'AggregateError', | ||
| message: 'all attempts failed', | ||
| errors: [ | ||
| { | ||
| name: 'TypeError', | ||
| message: 'fetch failed', | ||
| cause: { | ||
| name: 'Error', | ||
| message: 'other side closed', | ||
| code: 'UND_ERR_SOCKET', | ||
| stack: expect.stringContaining('other side closed'), | ||
| }, | ||
| }, | ||
| ], | ||
| }); | ||
| }); | ||
|
|
||
| it('does not throw on a cyclic cause chain', () => { | ||
| const outer = new Error('outer'); | ||
| const inner = new Error('inner', { cause: outer }); | ||
| outer.cause = inner; | ||
|
|
||
| const meta = JSON.parse(serializeGraphileMeta({ error: outer })); | ||
| expect(meta.error).toMatchObject({ | ||
| message: 'outer', | ||
| cause: { | ||
| message: 'inner', | ||
| cause: { name: 'Error', message: 'outer', repeated: true }, | ||
| }, | ||
| }); | ||
| }); | ||
| }); |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,55 @@ | ||
| import assert from 'node:assert/strict'; | ||
| import { randomUUID } from 'node:crypto'; | ||
| import { createServer } from 'node:http'; | ||
| import { setTimeout as sleep } from 'node:timers/promises'; | ||
| import { getQueueTopicPrefix } from '@workflow/world'; | ||
| import { Pool } from 'pg'; | ||
| import { createQueue } from '../../dist/queue.js'; | ||
|
|
||
| assert(process.env.WORKFLOW_POSTGRES_URL, 'WORKFLOW_POSTGRES_URL is required'); | ||
| const connectionString = process.env.WORKFLOW_POSTGRES_URL; | ||
| const pool = new Pool({ connectionString, max: 4 }); | ||
| const server = createServer(async (request, response) => { | ||
| await request.toArray(); | ||
| response.writeHead(503, { 'content-type': 'text/plain' }); | ||
| response.end('test queue delivery failure'); | ||
| }); | ||
| await new Promise((resolve, reject) => { | ||
| server.once('error', reject); | ||
| server.listen(0, '127.0.0.1', resolve); | ||
| }); | ||
| const address = server.address(); | ||
| assert(address && typeof address !== 'string', 'Expected a TCP server'); | ||
| process.env.WORKFLOW_LOCAL_BASE_URL = `http://127.0.0.1:${address.port}`; | ||
| const queue = createQueue( | ||
| { connectionString, queueConcurrency: 1, applicationManagedShutdown: true }, | ||
| pool | ||
| ); | ||
| try { | ||
| const { messageId } = await queue.queue( | ||
| `${getQueueTopicPrefix('workflow')}test`, | ||
| { runId: `run_${randomUUID()}` } | ||
| ); | ||
| const deadline = Date.now() + 10_000; | ||
| for (;;) { | ||
| const jobs = await pool.query( | ||
| 'SELECT attempts, last_error, locked_at FROM graphile_worker._private_jobs WHERE key = $1', | ||
| [messageId] | ||
| ); | ||
| const job = jobs.rows[0]; | ||
| if (job && job.last_error !== null && job.locked_at === null) { | ||
| assert.equal(job.attempts, 1); | ||
| break; | ||
| } | ||
| assert(Date.now() < deadline, 'Queue delivery did not fail'); | ||
| await sleep(20); | ||
| } | ||
| } finally { | ||
| await queue.close(); | ||
| await pool.query('TRUNCATE graphile_worker._private_jobs'); | ||
| await pool.end(); | ||
| await new Promise((resolve, reject) => { | ||
| server.close((error) => (error ? reject(error) : resolve())); | ||
| server.closeAllConnections(); | ||
| }); | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,70 @@ | ||
| import { execFile } from 'node:child_process'; | ||
| import { fileURLToPath } from 'node:url'; | ||
| import { promisify } from 'node:util'; | ||
| import { PostgreSqlContainer } from '@testcontainers/postgresql'; | ||
| import { afterAll, beforeAll, describe, expect, test } from 'vitest'; | ||
|
|
||
| const run = promisify(execFile); | ||
|
|
||
| /** | ||
| * What Graphile Worker writes to stderr when a delivery fails, captured from a | ||
| * real worker in a subprocess (`fixtures/failing-queue.mjs`, which imports the | ||
| * built `dist/queue.js`). The serializer itself is unit-tested in | ||
| * `src/queue-logging.test.ts`; this pins down that Graphile hands the logger | ||
| * the `Error` in `meta` and that the two output controls still hold. | ||
| */ | ||
| describe('Postgres queue error logs (integration)', () => { | ||
| if (process.platform === 'win32') { | ||
| test.skip('skipped on Windows since it relies on a docker container', () => {}); | ||
| return; | ||
| } | ||
|
|
||
| let container: Awaited<ReturnType<PostgreSqlContainer['start']>>; | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. AI Review: Nit Same as on #4114: the |
||
|
|
||
| beforeAll(async () => { | ||
| container = await new PostgreSqlContainer('postgres:15-alpine').start(); | ||
| }, 120_000); | ||
|
|
||
| afterAll(async () => { | ||
| await container.stop(); | ||
| }); | ||
|
|
||
| async function failDelivery(jsonMode: boolean) { | ||
| return await run( | ||
| process.execPath, | ||
| [fileURLToPath(new URL('./fixtures/failing-queue.mjs', import.meta.url))], | ||
| { | ||
| env: { | ||
| ...process.env, | ||
| DEBUG: '', | ||
| WORKFLOW_JSON_MODE: jsonMode ? '1' : '0', | ||
| WORKFLOW_POSTGRES_URL: container.getConnectionUri(), | ||
| }, | ||
| timeout: 15_000, | ||
| } | ||
| ); | ||
| } | ||
|
|
||
| test('preserves the delivery error and stack in worker metadata', async () => { | ||
| const { stdout, stderr } = await failDelivery(false); | ||
| expect(stdout).toBe(''); | ||
| expect(stderr).toContain('[Graphile Worker] Failed task'); | ||
| const metadataStart = stderr.indexOf('{\n'); | ||
| expect(metadataStart).toBeGreaterThan(-1); | ||
| const metadata = JSON.parse(stderr.slice(metadataStart)); | ||
| expect(metadata.error).toMatchObject({ | ||
| name: 'Error', | ||
| message: | ||
| '[postgres world] Queue execution failed (503): test queue delivery failure', | ||
| stack: expect.stringContaining( | ||
| 'Error: [postgres world] Queue execution failed (503): test queue delivery failure' | ||
| ), | ||
| }); | ||
| }); | ||
|
|
||
| test('keeps worker output suppressed in CLI JSON mode', async () => { | ||
| const { stdout, stderr } = await failDelivery(true); | ||
| expect(stdout).toBe(''); | ||
| expect(stderr).toBe(''); | ||
| }); | ||
| }); | ||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
AI Review: Blocking
The original replacer returned
{ ...value, name, message, stack, cause: value.cause }, andJSON.stringifyrecurses into the returned object, so a cyclic cause chain (a.cause = b; b.cause = a) overflows the stack (verified:RangeError: Maximum call stack size exceeded). That throw happens inside Graphile Worker's logger, which has no fallback, so a delivery failure with such an error would take the worker's error path down with it. Rare in practice, but the whole point of this change is the error path, and the guard is cheap. The follow-up tracks expanded errors in aWeakSetand replaces a repeat with{ name, message, repeated: true }; it also carriesAggregateError.errors, which the spread does not copy.serializeGraphileMetais exported so the unit test insrc/queue-logging.test.tscan cover the cycle, thecauserecursion, and enumerablecoderetention directly, without a subprocess.