From 74ccbfd8c8780e5ca9b81ed3ee5924f8c87e0375 Mon Sep 17 00:00:00 2001 From: Dev M <294291171+devtechedge@users.noreply.github.com> Date: Thu, 8 Oct 2026 23:41:13 +0530 Subject: [PATCH] Clear the task identifier cache when a lookup fails --- .changeset/kind-otters-task-cache.md | 6 +++++ __tests__/taskIdentifiers.test.ts | 37 ++++++++++++++++++++++++++++ src/taskIdentifiers.ts | 8 +++++- 3 files changed, 50 insertions(+), 1 deletion(-) create mode 100644 .changeset/kind-otters-task-cache.md create mode 100644 __tests__/taskIdentifiers.test.ts diff --git a/.changeset/kind-otters-task-cache.md b/.changeset/kind-otters-task-cache.md new file mode 100644 index 00000000..c3f35584 --- /dev/null +++ b/.changeset/kind-otters-task-cache.md @@ -0,0 +1,6 @@ +--- +"graphile-worker": patch +--- + +Clear the task identifier cache when a lookup fails, so a later poll retries +instead of reusing the rejected promise (#638). diff --git a/__tests__/taskIdentifiers.test.ts b/__tests__/taskIdentifiers.test.ts new file mode 100644 index 00000000..5bb244e2 --- /dev/null +++ b/__tests__/taskIdentifiers.test.ts @@ -0,0 +1,37 @@ +import type { EnhancedWithPgClient, TaskList } from "../src/interfaces.ts"; +import type { CompiledSharedOptions } from "../src/lib.ts"; +import { getTaskDetails } from "../src/taskIdentifiers.ts"; + +const tasks: TaskList = { + send_email: () => undefined, +}; + +test("retries a task identifier lookup after the cached lookup rejects", async () => { + let attempts = 0; + const withPgClient = Object.assign(async () => undefined, { + withRetries: async () => { + attempts += 1; + if (attempts === 1) { + const error = new Error("admin shutdown"); + Object.assign(error, { code: "57P01" }); + throw error; + } + return { rows: [{ id: 7, identifier: "send_email" }] }; + }, + }) as EnhancedWithPgClient; + const compiled = { + escapedWorkerSchema: "graphile_worker", + } as CompiledSharedOptions; + + await expect(getTaskDetails(compiled, withPgClient, tasks)).rejects.toThrow( + "admin shutdown", + ); + + const details = await getTaskDetails(compiled, withPgClient, tasks); + expect(attempts).toBe(2); + expect(details.taskIds).toEqual([7]); + expect(details.supportedTaskIdentifierByTaskId[7]).toBe("send_email"); + + await getTaskDetails(compiled, withPgClient, tasks); + expect(attempts).toBe(2); +}); diff --git a/src/taskIdentifiers.ts b/src/taskIdentifiers.ts index eb5d630a..40901b87 100644 --- a/src/taskIdentifiers.ts +++ b/src/taskIdentifiers.ts @@ -40,7 +40,7 @@ export function getTaskDetails( const { escapedWorkerSchema } = compiledSharedOptions; assert.ok(supportedTaskNames.length, "No runnable tasks!"); cache.lastStr = str; - cache.lastDigest = (async () => { + const digestPromise = (async () => { const { rows } = await withPgClient.withRetries(async (client) => { await client.query({ text: `insert into ${escapedWorkerSchema}._private_tasks as tasks (identifier) select i from unnest($1::text[]) as u(i) where not exists (select 1 from ${escapedWorkerSchema}._private_tasks as existing where existing.identifier = u.i) on conflict do nothing`, @@ -69,6 +69,12 @@ export function getTaskDetails( cache.lastStr = str; return cache.lastDigest; })(); + cache.lastDigest = digestPromise; + void digestPromise.catch(() => { + if (cache.lastStr === str) { + cache.lastStr = ""; + } + }); } return cache.lastDigest; }