forked from Chatbot/v3-api
74 lines
1.8 KiB
TypeScript
74 lines
1.8 KiB
TypeScript
import { Readable } from 'node:stream';
|
|
|
|
export type SseEvent = {
|
|
event: string;
|
|
data: unknown;
|
|
rawData: string;
|
|
};
|
|
|
|
export function parseSseBlock(block: string): SseEvent | null {
|
|
const lines = block.split('\n');
|
|
let event = 'message';
|
|
const dataLines: string[] = [];
|
|
|
|
for (const line of lines) {
|
|
if (!line || line.startsWith(':')) continue;
|
|
if (line.startsWith('event:')) {
|
|
event = line.slice(6).trim();
|
|
continue;
|
|
}
|
|
if (line.startsWith('data:')) {
|
|
dataLines.push(line.slice(5).trimStart());
|
|
}
|
|
}
|
|
|
|
if (dataLines.length === 0) return null;
|
|
const rawData = dataLines.join('\n');
|
|
let data: unknown = rawData;
|
|
try {
|
|
data = JSON.parse(rawData);
|
|
} catch {
|
|
data = rawData;
|
|
}
|
|
return { event, data, rawData };
|
|
}
|
|
|
|
export async function readStreamToString(stream: Readable): Promise<string> {
|
|
const chunks: Buffer[] = [];
|
|
for await (const chunk of stream) {
|
|
chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk));
|
|
}
|
|
return Buffer.concat(chunks).toString('utf8');
|
|
}
|
|
|
|
export async function* iterateSseEvents(
|
|
stream: Readable,
|
|
): AsyncGenerator<SseEvent> {
|
|
let buffer = '';
|
|
|
|
for await (const chunk of stream) {
|
|
buffer += Buffer.isBuffer(chunk) ? chunk.toString('utf8') : String(chunk);
|
|
const parts = buffer.split(/\r?\n\r?\n/);
|
|
buffer = parts.pop() ?? '';
|
|
for (const part of parts) {
|
|
const parsed = parseSseBlock(part);
|
|
if (parsed) yield parsed;
|
|
}
|
|
}
|
|
|
|
if (buffer.trim()) {
|
|
const parsed = parseSseBlock(buffer);
|
|
if (parsed) yield parsed;
|
|
}
|
|
}
|
|
|
|
export async function collectSseEvents(
|
|
stream: Readable,
|
|
): Promise<SseEvent[]> {
|
|
const events: SseEvent[] = [];
|
|
for await (const event of iterateSseEvents(stream)) {
|
|
events.push(event);
|
|
}
|
|
return events;
|
|
}
|