fix(event-stream): clean up reconnect delay listeners (#1180)
* fix(event-stream): clean up reconnect delay listeners * test(event-stream): clarify listener count --------- Co-authored-by: Isaac Sanchez <isanchez-hawkins@arize.com>
This commit is contained in:
committed by
GitHub
co-authored by
Isaac Sanchez
parent
7506b736dd
commit
2a9c7e6bca
@@ -15,11 +15,19 @@ function waitForReconnectDelay(ms, signal) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
return new Promise((resolve) => {
|
return new Promise((resolve) => {
|
||||||
const timeout = setTimeout(resolve, Math.max(0, ms));
|
let settled = false;
|
||||||
signal?.addEventListener('abort', () => {
|
const finish = () => {
|
||||||
clearTimeout(timeout);
|
if (settled) return;
|
||||||
|
settled = true;
|
||||||
|
signal?.removeEventListener('abort', onAbort);
|
||||||
resolve();
|
resolve();
|
||||||
}, { once: true });
|
};
|
||||||
|
const timeout = setTimeout(finish, Math.max(0, ms));
|
||||||
|
const onAbort = () => {
|
||||||
|
clearTimeout(timeout);
|
||||||
|
finish();
|
||||||
|
};
|
||||||
|
signal?.addEventListener('abort', onAbort, { once: true });
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -37,6 +37,28 @@ function createSseResponse({ blocks = [], signal, holdOpen = false }) {
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function createTrackedSignal() {
|
||||||
|
const listeners = new Set();
|
||||||
|
return {
|
||||||
|
signal: {
|
||||||
|
aborted: false,
|
||||||
|
addEventListener(type, listener) {
|
||||||
|
if (type === 'abort') {
|
||||||
|
listeners.add(listener);
|
||||||
|
}
|
||||||
|
},
|
||||||
|
removeEventListener(type, listener) {
|
||||||
|
if (type === 'abort') {
|
||||||
|
listeners.delete(listener);
|
||||||
|
}
|
||||||
|
},
|
||||||
|
},
|
||||||
|
getListenerCount() {
|
||||||
|
return listeners.size;
|
||||||
|
},
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
describe('createUpstreamSseReader', () => {
|
describe('createUpstreamSseReader', () => {
|
||||||
it('emits parsed events and tracks the latest event id', async () => {
|
it('emits parsed events and tracks the latest event id', async () => {
|
||||||
const events = [];
|
const events = [];
|
||||||
@@ -211,4 +233,44 @@ describe('createUpstreamSseReader', () => {
|
|||||||
expect(unavailableBodyCanceled).toBe(true);
|
expect(unavailableBodyCanceled).toBe(true);
|
||||||
expect(attempt).toBe(2);
|
expect(attempt).toBe(2);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it('removes reconnect delay abort listeners after normal timeout completion', async () => {
|
||||||
|
const tracked = createTrackedSignal();
|
||||||
|
let attempt = 0;
|
||||||
|
let reader;
|
||||||
|
|
||||||
|
reader = createUpstreamSseReader({
|
||||||
|
buildUrl: () => 'http://127.0.0.1:4096/global/event',
|
||||||
|
reconnectDelayMs: 1,
|
||||||
|
signal: tracked.signal,
|
||||||
|
fetchImpl: async (_url, options) => {
|
||||||
|
attempt += 1;
|
||||||
|
if (attempt === 1) {
|
||||||
|
return {
|
||||||
|
ok: false,
|
||||||
|
status: 503,
|
||||||
|
body: {
|
||||||
|
cancel: async () => {},
|
||||||
|
},
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
return createSseResponse({
|
||||||
|
signal: options.signal,
|
||||||
|
blocks: [
|
||||||
|
'id: evt-1\ndata: {"type":"server.connected","properties":{}}\n\n',
|
||||||
|
],
|
||||||
|
});
|
||||||
|
},
|
||||||
|
onEvent() {
|
||||||
|
reader.stop();
|
||||||
|
},
|
||||||
|
});
|
||||||
|
|
||||||
|
await reader.start();
|
||||||
|
|
||||||
|
expect(attempt).toBe(2);
|
||||||
|
// The top-level stop listener remains; the reconnect-delay listener should be removed.
|
||||||
|
expect(tracked.getListenerCount()).toBe(1);
|
||||||
|
});
|
||||||
});
|
});
|
||||||
|
|||||||
Reference in New Issue
Block a user