-
Notifications
You must be signed in to change notification settings - Fork 162
Expand file tree
/
Copy pathstream.ts
More file actions
82 lines (66 loc) · 1.91 KB
/
Copy pathstream.ts
File metadata and controls
82 lines (66 loc) · 1.91 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
export interface ServerSentEvent {
event?: string;
data: string;
id?: string;
}
export async function* parseSSE(response: Response): AsyncGenerator<ServerSentEvent> {
const reader = response.body?.getReader();
if (!reader) return;
const decoder = new TextDecoder();
let buffer = '';
let event: Partial<ServerSentEvent> = {};
const processLine = (rawLine: string): ServerSentEvent | undefined => {
const line = rawLine.endsWith('\r') ? rawLine.slice(0, -1) : rawLine;
if (line === '') {
const completed = event.data !== undefined
? { data: event.data, event: event.event, id: event.id }
: undefined;
event = {};
return completed;
}
if (line.startsWith(':')) return undefined;
const colonIndex = line.indexOf(':');
if (colonIndex === -1) return undefined;
const field = line.slice(0, colonIndex);
const value = line.slice(colonIndex + 1).trimStart();
switch (field) {
case 'data':
event.data = event.data !== undefined ? `${event.data}\n${value}` : value;
break;
case 'event':
event.event = value;
break;
case 'id':
event.id = value;
break;
}
return undefined;
};
try {
while (true) {
const { done, value } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
const lines = buffer.split('\n');
buffer = lines.pop() || '';
for (const line of lines) {
const completed = processLine(line);
if (completed) {
yield completed;
}
}
}
buffer += decoder.decode();
if (buffer.length > 0) {
const completed = processLine(buffer);
if (completed) {
yield completed;
}
}
if (event.data !== undefined) {
yield { data: event.data, event: event.event, id: event.id };
}
} finally {
reader.releaseLock();
}
}