Skip to content

Commit 7e8e5dd

Browse files
fix(world-postgres): stream reader lifecycle cleanup and offset cursor (#4125)
1 parent 03455a2 commit 7e8e5dd

3 files changed

Lines changed: 380 additions & 12 deletions

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
'@workflow/world-postgres': patch
3+
---
4+
5+
Fix stream readers leaking EventEmitter listeners on EOF, initial query failure, and World close, and fail pending readers when the World is closed.
Lines changed: 311 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,311 @@
1+
import { EventEmitter } from 'node:events';
2+
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
3+
import type { Drizzle } from './drizzle/index.js';
4+
import { createStreamer } from './streamer.js';
5+
6+
vi.mock('pg', () => ({
7+
Client: vi.fn(function Client() {
8+
return {
9+
connect: vi.fn(async () => {}),
10+
query: vi.fn(async () => {}),
11+
on: vi.fn(),
12+
removeListener: vi.fn(),
13+
end: vi.fn(async () => {}),
14+
};
15+
}),
16+
Pool: vi.fn(),
17+
}));
18+
19+
const fakePool = {
20+
options: {},
21+
query: vi.fn(async () => ({ rows: [] })),
22+
} as any;
23+
24+
/**
25+
* Minimal fake for the drizzle query chain used by `streams.get()`:
26+
* `drizzle.select(...).from(...).where(...).orderBy(...)`.
27+
*/
28+
function createFakeDrizzle(
29+
rows: () => Promise<Array<{ id: string; eof: boolean; data: Buffer }>>
30+
): Drizzle {
31+
return {
32+
select: () => ({
33+
from: () => ({
34+
where: () => ({
35+
orderBy: () => rows(),
36+
limit: () => rows(),
37+
}),
38+
}),
39+
}),
40+
} as unknown as Drizzle;
41+
}
42+
43+
describe('postgres streamer reader listener cleanup', () => {
44+
// `createStreamer()` keeps its EventEmitter private, so capture any
45+
// emitter that receives a `strm:*` listener via a prototype spy.
46+
let streamEmitters: Set<EventEmitter>;
47+
48+
const listenerCount = (name: string) => {
49+
let count = 0;
50+
for (const emitter of streamEmitters) {
51+
count += emitter.listenerCount(`strm:${name}`);
52+
}
53+
return count;
54+
};
55+
56+
beforeEach(() => {
57+
streamEmitters = new Set();
58+
const originalOn = EventEmitter.prototype.on;
59+
vi.spyOn(EventEmitter.prototype, 'on').mockImplementation(function (
60+
this: EventEmitter,
61+
eventName,
62+
listener
63+
) {
64+
if (typeof eventName === 'string' && eventName.startsWith('strm:')) {
65+
streamEmitters.add(this);
66+
}
67+
return originalOn.call(this, eventName, listener);
68+
});
69+
});
70+
71+
afterEach(() => {
72+
vi.restoreAllMocks();
73+
});
74+
75+
it('removes the listener when a persisted EOF closes the stream', async () => {
76+
const drizzle = createFakeDrizzle(async () => [
77+
{
78+
id: 'chnk_00000000000000000000000001',
79+
eof: false,
80+
data: Buffer.from('hello'),
81+
},
82+
{
83+
id: 'chnk_00000000000000000000000002',
84+
eof: true,
85+
data: Buffer.from([]),
86+
},
87+
]);
88+
const streamer = createStreamer(fakePool, drizzle);
89+
90+
const warnings: Error[] = [];
91+
const onWarning = (warning: Error) => warnings.push(warning);
92+
process.on('warning', onWarning);
93+
94+
try {
95+
// Repeated reads of the same completed stream must not accumulate
96+
// listeners (issue reproduction uses 12 readers to trip the default
97+
// max-listeners threshold of 10).
98+
for (let i = 0; i < 12; i++) {
99+
const stream = await streamer.streams.get('run_1', 'stream-eof');
100+
const reader = stream.getReader();
101+
let done = false;
102+
while (!done) {
103+
({ done } = await reader.read());
104+
}
105+
expect(listenerCount('stream-eof')).toBe(0);
106+
}
107+
108+
// Warnings are emitted via process.nextTick; flush before asserting.
109+
await new Promise((resolve) => setImmediate(resolve));
110+
expect(
111+
warnings.filter((w) => w.name === 'MaxListenersExceededWarning')
112+
).toEqual([]);
113+
} finally {
114+
process.removeListener('warning', onWarning);
115+
await streamer.close();
116+
}
117+
});
118+
119+
it('removes the listener when the initial chunk query rejects', async () => {
120+
const drizzle = createFakeDrizzle(async () => {
121+
throw new Error('initial query failed');
122+
});
123+
const streamer = createStreamer(fakePool, drizzle);
124+
125+
try {
126+
for (let i = 0; i < 12; i++) {
127+
const stream = await streamer.streams.get('run_1', 'stream-err');
128+
const reader = stream.getReader();
129+
await expect(reader.read()).rejects.toThrow('initial query failed');
130+
}
131+
expect(listenerCount('stream-err')).toBe(0);
132+
} finally {
133+
await streamer.close();
134+
}
135+
});
136+
137+
it('removes the listener when a live EOF arrives via notification', async () => {
138+
// No persisted EOF: the reader tails live events and is closed by an
139+
// EOF chunk delivered through the LISTEN emitter.
140+
const drizzle = createFakeDrizzle(async () => [
141+
{
142+
id: 'chnk_00000000000000000000000001',
143+
eof: false,
144+
data: Buffer.from('hello'),
145+
},
146+
]);
147+
const streamer = createStreamer(fakePool, drizzle);
148+
149+
try {
150+
for (let i = 0; i < 12; i++) {
151+
const stream = await streamer.streams.get('run_1', 'stream-live-eof');
152+
const reader = stream.getReader();
153+
await reader.read();
154+
expect(listenerCount('stream-live-eof')).toBe(1);
155+
156+
for (const emitter of streamEmitters) {
157+
emitter.emit('strm:stream-live-eof', {
158+
id: 'chnk_00000000000000000000000002',
159+
eof: true,
160+
data: Buffer.from([]),
161+
});
162+
}
163+
expect(await reader.read()).toEqual({ done: true, value: undefined });
164+
expect(listenerCount('stream-live-eof')).toBe(0);
165+
}
166+
} finally {
167+
await streamer.close();
168+
}
169+
});
170+
171+
it('detaches and fails pending readers when the streamer is closed', async () => {
172+
// No EOF chunk: readers stay open, tailing live events.
173+
const drizzle = createFakeDrizzle(async () => [
174+
{
175+
id: 'chnk_00000000000000000000000001',
176+
eof: false,
177+
data: Buffer.from('hello'),
178+
},
179+
]);
180+
const streamer = createStreamer(fakePool, drizzle);
181+
182+
const readers: ReadableStreamDefaultReader<Uint8Array>[] = [];
183+
for (let i = 0; i < 10; i++) {
184+
const stream = await streamer.streams.get('run_1', 'stream-open');
185+
readers.push(stream.getReader());
186+
// Consume the persisted chunk; the reader keeps waiting for more.
187+
await readers[i].read();
188+
}
189+
expect(listenerCount('stream-open')).toBe(10);
190+
191+
// An outstanding read() must settle once the streamer shuts down: the
192+
// LISTEN client is gone, so nothing could ever wake it otherwise.
193+
const pending = readers[0].read();
194+
195+
await streamer.close();
196+
expect(listenerCount('stream-open')).toBe(0);
197+
await expect(pending).rejects.toThrow('streamer has been closed');
198+
for (const reader of readers) {
199+
await expect(reader.read()).rejects.toThrow('streamer has been closed');
200+
}
201+
});
202+
203+
it('rejects streams.get() after the streamer is closed', async () => {
204+
const drizzle = createFakeDrizzle(async () => []);
205+
const streamer = createStreamer(fakePool, drizzle);
206+
await streamer.close();
207+
208+
await expect(
209+
streamer.streams.get('run_1', 'stream-after-close')
210+
).rejects.toThrow('streamer has been closed');
211+
// A rejected get() must not leave a listener behind on the dead emitter.
212+
expect(listenerCount('stream-after-close')).toBe(0);
213+
});
214+
215+
it('cleans up when cancelled while the initial query is in flight', async () => {
216+
let resolveQuery!: (
217+
rows: Array<{ id: string; eof: boolean; data: Buffer }>
218+
) => void;
219+
const drizzle = createFakeDrizzle(
220+
() => new Promise((resolve) => (resolveQuery = resolve))
221+
);
222+
const streamer = createStreamer(fakePool, drizzle);
223+
224+
try {
225+
const stream = await streamer.streams.get('run_1', 'stream-cancel-early');
226+
const reader = stream.getReader();
227+
expect(listenerCount('stream-cancel-early')).toBe(1);
228+
await reader.cancel();
229+
expect(listenerCount('stream-cancel-early')).toBe(0);
230+
231+
// The query resolving into a cancelled stream is dropped; the listener
232+
// count must not bounce back.
233+
resolveQuery([
234+
{
235+
id: 'chnk_00000000000000000000000001',
236+
eof: false,
237+
data: Buffer.from('hello'),
238+
},
239+
{
240+
id: 'chnk_00000000000000000000000002',
241+
eof: true,
242+
data: Buffer.from([]),
243+
},
244+
]);
245+
await new Promise((resolve) => setImmediate(resolve));
246+
expect(listenerCount('stream-cancel-early')).toBe(0);
247+
} finally {
248+
await streamer.close();
249+
}
250+
});
251+
252+
it('does not count a redelivered skipped chunk against the start offset', async () => {
253+
const drizzle = createFakeDrizzle(async () => []);
254+
const streamer = createStreamer(fakePool, drizzle);
255+
256+
try {
257+
// startIndex 2: skip 'first' and 'second', deliver from 'third'.
258+
const stream = await streamer.streams.get('run_1', 'stream-offset', 2);
259+
const reader = stream.getReader();
260+
const first = {
261+
id: 'chnk_00000000000000000000000001',
262+
eof: false,
263+
data: Buffer.from('first'),
264+
};
265+
const emit = (chunk: typeof first) => {
266+
for (const emitter of streamEmitters) {
267+
emitter.emit('strm:stream-offset', chunk);
268+
}
269+
};
270+
// Let start() finish its (empty) initial query so notifications are
271+
// enqueued directly rather than buffered.
272+
await new Promise((resolve) => setImmediate(resolve));
273+
274+
emit(first);
275+
// A redelivered NOTIFY for the chunk that was just skipped must not
276+
// decrement the remaining offset a second time.
277+
emit(first);
278+
emit({
279+
id: 'chnk_00000000000000000000000002',
280+
eof: false,
281+
data: Buffer.from('second'),
282+
});
283+
emit({
284+
id: 'chnk_00000000000000000000000003',
285+
eof: false,
286+
data: Buffer.from('third'),
287+
});
288+
289+
const result = await reader.read();
290+
expect(result.done).toBe(false);
291+
expect(Buffer.from(result.value as Uint8Array).toString()).toBe('third');
292+
} finally {
293+
await streamer.close();
294+
}
295+
});
296+
297+
it('still cleans up on explicit consumer cancellation', async () => {
298+
const drizzle = createFakeDrizzle(async () => []);
299+
const streamer = createStreamer(fakePool, drizzle);
300+
301+
try {
302+
const stream = await streamer.streams.get('run_1', 'stream-cancel');
303+
const reader = stream.getReader();
304+
expect(listenerCount('stream-cancel')).toBe(1);
305+
await reader.cancel();
306+
expect(listenerCount('stream-cancel')).toBe(0);
307+
} finally {
308+
await streamer.close();
309+
}
310+
});
311+
});

0 commit comments

Comments
 (0)