diff --git a/.changeset/tidy-local-shutdown.md b/.changeset/tidy-local-shutdown.md new file mode 100644 index 0000000000..3b4fef7d76 --- /dev/null +++ b/.changeset/tidy-local-shutdown.md @@ -0,0 +1,5 @@ +--- +'@workflow/world-local': patch +--- + +Abort active local queue deliveries when the World closes, including when transport timeouts are disabled. diff --git a/packages/world-local/src/queue.test.ts b/packages/world-local/src/queue.test.ts index acf51ca08e..a243eaca77 100644 --- a/packages/world-local/src/queue.test.ts +++ b/packages/world-local/src/queue.test.ts @@ -557,6 +557,40 @@ describe('queue transport timeouts', () => { }); }); + it('aborts an in-flight delivery when the queue closes', async () => { + let requestReceived = false; + server = createServer(() => { + requestReceived = true; + }); + await new Promise((resolve) => { + server?.listen(0, '127.0.0.1', resolve); + }); + const { port } = server.address() as AddressInfo; + + process.env.WORKFLOW_LOCAL_HEADERS_TIMEOUT_MS = '0'; + process.env.WORKFLOW_LOCAL_BODY_TIMEOUT_MS = '0'; + const localQueue = createQueue({ + baseUrl: `http://127.0.0.1:${port}`, + }); + + await localQueue.queue('__wkf_workflow_test' as any, workflowPayload); + await vi.waitFor(() => expect(requestReceived).toBe(true)); + const closePromise = localQueue.close(); + let timeout: ReturnType | undefined; + const outcome = await Promise.race([ + closePromise.then(() => 'closed' as const), + new Promise<'timed-out'>((resolve) => { + timeout = setTimeout(() => resolve('timed-out'), 1_000); + }), + ]); + clearTimeout(timeout); + if (outcome === 'timed-out') { + server.closeAllConnections(); + await closePromise; + } + expect(outcome).toBe('closed'); + }); + it('redelivers when a handler accepts a request but never responds', async () => { let requests = 0; server = createServer((_request, response) => { diff --git a/packages/world-local/src/queue.ts b/packages/world-local/src/queue.ts index 84ddefde8d..62ba4bf25d 100644 --- a/packages/world-local/src/queue.ts +++ b/packages/world-local/src/queue.ts @@ -259,6 +259,7 @@ export function createQueue(config: Partial): LocalQueue { headers: new Headers(headers), body, agents: nodeHttpAgents, + signal: closeSignal, headersTimeoutMs: agentOptions.headersTimeout, bodyTimeoutMs: agentOptions.bodyTimeout, }) @@ -269,6 +270,7 @@ export function createQueue(config: Partial): LocalQueue { dispatcher: httpAgent, headers, body, + signal: closeSignal, } as any); } delivery++;