2026-04-27 07:24:32 +00:00
|
|
|
const { assert, requireBundledModule, cleanup } = require('./direct-test-setup');
|
|
|
|
|
|
|
|
|
|
(async () => {
|
|
|
|
|
try {
|
|
|
|
|
const streaming = await requireBundledModule('src/streaming.ts');
|
|
|
|
|
|
|
|
|
|
const openAiExtractor = streaming.deltaExtractorForFormat('openai-chat');
|
|
|
|
|
const anthropicExtractor = streaming.deltaExtractorForFormat('anthropic-messages');
|
|
|
|
|
|
|
|
|
|
assert.ok(openAiExtractor);
|
|
|
|
|
assert.ok(anthropicExtractor);
|
|
|
|
|
assert.strictEqual(streaming.deltaExtractorForFormat('unknown-format'), null);
|
|
|
|
|
assert.strictEqual(streaming.deltaExtractorForFormat('google-generative-ai'), null);
|
|
|
|
|
|
|
|
|
|
assert.strictEqual(openAiExtractor({ choices: [{ delta: { content: 'hello' } }] }), 'hello');
|
|
|
|
|
assert.strictEqual(openAiExtractor({ choices: [{ delta: {} }] }), '');
|
|
|
|
|
assert.strictEqual(openAiExtractor({}), '');
|
|
|
|
|
|
|
|
|
|
assert.strictEqual(anthropicExtractor({ type: 'content_block_delta', delta: { text: 'world' } }), 'world');
|
|
|
|
|
assert.strictEqual(anthropicExtractor({ type: 'content_block_start' }), '');
|
|
|
|
|
assert.strictEqual(anthropicExtractor({}), '');
|
|
|
|
|
|
|
|
|
|
// ── parseSseBuffer ──
|
|
|
|
|
const single = streaming.parseSseBuffer('data: {"choices":[{"delta":{"content":"hi"}}]}\n\n', openAiExtractor);
|
|
|
|
|
assert.deepStrictEqual(single.deltas, ['hi']);
|
|
|
|
|
assert.strictEqual(single.rest, '');
|
|
|
|
|
|
|
|
|
|
const multi = streaming.parseSseBuffer(
|
|
|
|
|
'data: {"choices":[{"delta":{"content":"a"}}]}\n\ndata: {"choices":[{"delta":{"content":"b"}}]}\n\n',
|
|
|
|
|
openAiExtractor,
|
|
|
|
|
);
|
|
|
|
|
assert.deepStrictEqual(multi.deltas, ['a', 'b']);
|
|
|
|
|
|
|
|
|
|
const incomplete = streaming.parseSseBuffer('data: {"choices":[{"delta":{"content":"partial"}}]}', openAiExtractor);
|
|
|
|
|
assert.deepStrictEqual(incomplete.deltas, []);
|
|
|
|
|
assert.ok(incomplete.rest.length > 0);
|
|
|
|
|
|
2026-04-27 11:53:35 +00:00
|
|
|
const withDone = streaming.parseSseBuffer(
|
|
|
|
|
'data: {"choices":[{"delta":{"content":"x"}}]}\n\ndata: [DONE]\n\n',
|
|
|
|
|
openAiExtractor,
|
|
|
|
|
);
|
2026-04-27 07:24:32 +00:00
|
|
|
assert.deepStrictEqual(withDone.deltas, ['x']);
|
|
|
|
|
|
2026-04-27 11:53:35 +00:00
|
|
|
const withComments = streaming.parseSseBuffer(
|
|
|
|
|
': keep-alive\nevent: message\ndata: {"choices":[{"delta":{"content":"ok"}}]}\n\n',
|
|
|
|
|
openAiExtractor,
|
|
|
|
|
);
|
2026-04-27 07:24:32 +00:00
|
|
|
assert.deepStrictEqual(withComments.deltas, ['ok']);
|
|
|
|
|
|
|
|
|
|
const crlf = streaming.parseSseBuffer('data: {"choices":[{"delta":{"content":"crlf"}}]}\r\n\r\n', openAiExtractor);
|
|
|
|
|
assert.deepStrictEqual(crlf.deltas, ['crlf']);
|
|
|
|
|
|
2026-04-27 11:53:35 +00:00
|
|
|
const multiLine = streaming.parseSseBuffer(
|
|
|
|
|
'data: {"choices":[{"delta":\ndata: {"content":"split"}}]}\n\n',
|
|
|
|
|
openAiExtractor,
|
|
|
|
|
);
|
2026-04-27 07:24:32 +00:00
|
|
|
assert.deepStrictEqual(multiLine.deltas, ['split']);
|
|
|
|
|
|
2026-04-27 11:53:35 +00:00
|
|
|
const malformed = streaming.parseSseBuffer(
|
|
|
|
|
'data: not-json\n\ndata: {"choices":[{"delta":{"content":"ok"}}]}\n\n',
|
|
|
|
|
openAiExtractor,
|
|
|
|
|
);
|
2026-04-27 07:24:32 +00:00
|
|
|
assert.deepStrictEqual(malformed.deltas, ['ok']);
|
|
|
|
|
|
|
|
|
|
const anthEvent = streaming.parseSseBuffer(
|
|
|
|
|
'data: {"type":"content_block_delta","delta":{"text":"ant"}}\n\ndata: {"type":"message_stop"}\n\n',
|
|
|
|
|
anthropicExtractor,
|
|
|
|
|
);
|
|
|
|
|
assert.deepStrictEqual(anthEvent.deltas, ['ant']);
|
|
|
|
|
|
|
|
|
|
const noSpace = streaming.parseSseBuffer('data:{"choices":[{"delta":{"content":"ns"}}]}\n\n', openAiExtractor);
|
|
|
|
|
assert.deepStrictEqual(noSpace.deltas, ['ns']);
|
|
|
|
|
|
|
|
|
|
const empty = streaming.parseSseBuffer('', openAiExtractor);
|
|
|
|
|
assert.deepStrictEqual(empty.deltas, []);
|
|
|
|
|
|
2026-05-08 14:57:31 +00:00
|
|
|
// ── streamingRequestUrl ──
|
2026-04-27 07:24:32 +00:00
|
|
|
function trackedSignal() {
|
|
|
|
|
const controller = new AbortController();
|
|
|
|
|
const signal = controller.signal;
|
|
|
|
|
let activeListeners = 0;
|
|
|
|
|
const addEventListener = signal.addEventListener.bind(signal);
|
|
|
|
|
const removeEventListener = signal.removeEventListener.bind(signal);
|
2026-04-27 11:53:35 +00:00
|
|
|
signal.addEventListener = (type, listener, options) => {
|
|
|
|
|
if (type === 'abort') activeListeners++;
|
|
|
|
|
return addEventListener(type, listener, options);
|
|
|
|
|
};
|
|
|
|
|
signal.removeEventListener = (type, listener, options) => {
|
|
|
|
|
if (type === 'abort') activeListeners--;
|
|
|
|
|
return removeEventListener(type, listener, options);
|
|
|
|
|
};
|
2026-04-27 07:24:32 +00:00
|
|
|
return { controller, signal, activeListeners: () => activeListeners };
|
|
|
|
|
}
|
|
|
|
|
|
2026-05-08 14:57:31 +00:00
|
|
|
{
|
2026-04-27 07:24:32 +00:00
|
|
|
const success = trackedSignal();
|
2026-05-08 14:57:31 +00:00
|
|
|
const progress = [];
|
|
|
|
|
const text = await streaming.streamingRequestUrl(
|
|
|
|
|
async (params) => {
|
|
|
|
|
assert.strictEqual(params.method, 'POST');
|
|
|
|
|
assert.strictEqual(params.url, 'https://example.test');
|
|
|
|
|
assert.strictEqual(params.body, '{"stream":true}');
|
|
|
|
|
return {
|
|
|
|
|
status: 200,
|
|
|
|
|
json: null,
|
|
|
|
|
text: 'data: {"choices":[{"delta":{"content":"ok"}}]}\n\n',
|
|
|
|
|
};
|
|
|
|
|
},
|
2026-04-27 11:53:35 +00:00
|
|
|
'https://example.test',
|
|
|
|
|
{},
|
2026-05-08 14:57:31 +00:00
|
|
|
{ stream: true },
|
2026-04-27 11:53:35 +00:00
|
|
|
streaming.deltaExtractorForFormat('openai-chat'),
|
2026-05-08 14:57:31 +00:00
|
|
|
(p) => progress.push(p),
|
2026-04-27 11:53:35 +00:00
|
|
|
success.signal,
|
|
|
|
|
{ streamingTimeoutMs: 1000 },
|
|
|
|
|
);
|
2026-05-08 14:57:31 +00:00
|
|
|
assert.strictEqual(text, 'ok');
|
|
|
|
|
assert.deepStrictEqual(progress, [
|
|
|
|
|
{ accumulated: 'ok', done: false },
|
|
|
|
|
{ accumulated: 'ok', done: true },
|
|
|
|
|
]);
|
2026-04-27 07:24:32 +00:00
|
|
|
assert.strictEqual(success.activeListeners(), 0, 'cleanup after success');
|
2026-05-08 14:57:31 +00:00
|
|
|
}
|
2026-04-27 07:24:32 +00:00
|
|
|
|
2026-05-08 14:57:31 +00:00
|
|
|
{
|
2026-04-27 07:24:32 +00:00
|
|
|
const httpError = trackedSignal();
|
2026-04-27 11:53:35 +00:00
|
|
|
await assert.rejects(
|
|
|
|
|
() =>
|
2026-05-08 14:57:31 +00:00
|
|
|
streaming.streamingRequestUrl(
|
|
|
|
|
async () => ({ status: 500, json: null, text: 'bad' }),
|
2026-04-27 11:53:35 +00:00
|
|
|
'https://example.test',
|
|
|
|
|
{},
|
|
|
|
|
{},
|
|
|
|
|
streaming.deltaExtractorForFormat('openai-chat'),
|
|
|
|
|
undefined,
|
|
|
|
|
httpError.signal,
|
|
|
|
|
{ streamingTimeoutMs: 1000 },
|
|
|
|
|
),
|
|
|
|
|
/HTTP 500|API returned HTTP 500/,
|
|
|
|
|
);
|
2026-04-27 07:24:32 +00:00
|
|
|
assert.strictEqual(httpError.activeListeners(), 0, 'cleanup after HTTP error');
|
2026-05-08 14:57:31 +00:00
|
|
|
}
|
2026-04-27 07:24:32 +00:00
|
|
|
|
2026-05-08 14:57:31 +00:00
|
|
|
{
|
2026-04-27 07:24:32 +00:00
|
|
|
const timeout = trackedSignal();
|
2026-04-27 11:53:35 +00:00
|
|
|
await assert.rejects(
|
|
|
|
|
() =>
|
2026-05-08 14:57:31 +00:00
|
|
|
streaming.streamingRequestUrl(
|
|
|
|
|
async () => new Promise(() => {}),
|
2026-04-27 11:53:35 +00:00
|
|
|
'https://example.test',
|
|
|
|
|
{},
|
|
|
|
|
{},
|
|
|
|
|
streaming.deltaExtractorForFormat('openai-chat'),
|
|
|
|
|
undefined,
|
|
|
|
|
timeout.signal,
|
|
|
|
|
{ streamingTimeoutMs: 1 },
|
|
|
|
|
),
|
|
|
|
|
/Streaming timed out/,
|
|
|
|
|
);
|
2026-04-27 07:24:32 +00:00
|
|
|
assert.strictEqual(timeout.activeListeners(), 0, 'cleanup after timeout');
|
2026-05-08 14:57:31 +00:00
|
|
|
}
|
2026-04-27 07:24:32 +00:00
|
|
|
|
2026-05-08 14:57:31 +00:00
|
|
|
{
|
2026-04-27 07:24:32 +00:00
|
|
|
const preAborted = trackedSignal();
|
|
|
|
|
preAborted.controller.abort();
|
2026-05-08 14:57:31 +00:00
|
|
|
let called = false;
|
2026-04-27 11:53:35 +00:00
|
|
|
await assert.rejects(
|
|
|
|
|
() =>
|
2026-05-08 14:57:31 +00:00
|
|
|
streaming.streamingRequestUrl(
|
|
|
|
|
async () => {
|
|
|
|
|
called = true;
|
|
|
|
|
return { status: 200, json: null, text: '' };
|
|
|
|
|
},
|
2026-04-27 11:53:35 +00:00
|
|
|
'https://example.test',
|
|
|
|
|
{},
|
|
|
|
|
{},
|
|
|
|
|
streaming.deltaExtractorForFormat('openai-chat'),
|
|
|
|
|
undefined,
|
|
|
|
|
preAborted.signal,
|
|
|
|
|
{ streamingTimeoutMs: 5000 },
|
|
|
|
|
),
|
|
|
|
|
/abort/i,
|
|
|
|
|
);
|
2026-05-08 14:57:31 +00:00
|
|
|
assert.strictEqual(called, false, 'pre-aborted request is not started');
|
2026-04-27 07:24:32 +00:00
|
|
|
assert.strictEqual(preAborted.activeListeners(), 0, 'cleanup on pre-aborted');
|
2026-05-08 14:57:31 +00:00
|
|
|
}
|
2026-04-27 07:24:32 +00:00
|
|
|
|
2026-05-08 14:57:31 +00:00
|
|
|
{
|
|
|
|
|
const abortDuringRequest = trackedSignal();
|
2026-04-27 11:53:35 +00:00
|
|
|
await assert.rejects(
|
|
|
|
|
() =>
|
2026-05-08 14:57:31 +00:00
|
|
|
streaming.streamingRequestUrl(
|
|
|
|
|
async () =>
|
|
|
|
|
new Promise((resolve) => {
|
|
|
|
|
setTimeout(() => abortDuringRequest.controller.abort(), 1);
|
|
|
|
|
setTimeout(() => resolve({ status: 200, json: null, text: '' }), 20);
|
|
|
|
|
}),
|
2026-04-27 11:53:35 +00:00
|
|
|
'https://example.test',
|
|
|
|
|
{},
|
|
|
|
|
{},
|
|
|
|
|
streaming.deltaExtractorForFormat('openai-chat'),
|
|
|
|
|
undefined,
|
2026-05-08 14:57:31 +00:00
|
|
|
abortDuringRequest.signal,
|
2026-04-27 11:53:35 +00:00
|
|
|
{ streamingTimeoutMs: 5000 },
|
|
|
|
|
),
|
|
|
|
|
/abort/i,
|
|
|
|
|
);
|
2026-05-08 14:57:31 +00:00
|
|
|
assert.strictEqual(abortDuringRequest.activeListeners(), 0, 'cleanup after mid-request abort');
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
{
|
|
|
|
|
const requestFailure = trackedSignal();
|
|
|
|
|
await assert.rejects(
|
|
|
|
|
() =>
|
|
|
|
|
streaming.streamingRequestUrl(
|
|
|
|
|
async () => {
|
|
|
|
|
throw new Error('network down');
|
|
|
|
|
},
|
|
|
|
|
'https://example.test',
|
|
|
|
|
{},
|
|
|
|
|
{},
|
|
|
|
|
streaming.deltaExtractorForFormat('openai-chat'),
|
|
|
|
|
undefined,
|
|
|
|
|
requestFailure.signal,
|
|
|
|
|
{ streamingTimeoutMs: 1000 },
|
|
|
|
|
),
|
|
|
|
|
/network down/,
|
|
|
|
|
);
|
|
|
|
|
assert.strictEqual(requestFailure.activeListeners(), 0, 'cleanup after transport failure');
|
2026-04-27 07:24:32 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
console.log('direct streaming tests passed');
|
|
|
|
|
} finally {
|
|
|
|
|
cleanup();
|
|
|
|
|
}
|
2026-04-27 11:53:35 +00:00
|
|
|
})().catch((e) => {
|
|
|
|
|
console.error(e);
|
|
|
|
|
process.exit(1);
|
|
|
|
|
});
|