-
-
Notifications
You must be signed in to change notification settings - Fork 444
/
sse.ts
77 lines (66 loc) · 1.75 KB
/
sse.ts
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
import type { Context } from '../../context.ts'
import { StreamingApi } from '../../utils/stream.ts'
export interface SSEMessage {
data: string
event?: string
id?: string
retry?: number
}
export class SSEStreamingApi extends StreamingApi {
constructor(writable: WritableStream, readable: ReadableStream) {
super(writable, readable)
}
async writeSSE(message: SSEMessage) {
const data = message.data
.split('\n')
.map((line) => {
return `data: ${line}`
})
.join('\n')
const sseData =
[
message.event && `event: ${message.event}`,
data,
message.id && `id: ${message.id}`,
message.retry && `retry: ${message.retry}`,
]
.filter(Boolean)
.join('\n') + '\n\n'
await this.write(sseData)
}
}
const run = async (
stream: SSEStreamingApi,
cb: (stream: SSEStreamingApi) => Promise<void>,
onError?: (e: Error, stream: SSEStreamingApi) => Promise<void>
) => {
try {
await cb(stream)
} catch (e) {
if (e instanceof Error && onError) {
await onError(e, stream)
await stream.writeSSE({
event: 'error',
data: e.message,
})
} else {
console.error(e)
}
} finally {
stream.close()
}
}
export const streamSSE = (
c: Context,
cb: (stream: SSEStreamingApi) => Promise<void>,
onError?: (e: Error, stream: SSEStreamingApi) => Promise<void>
) => {
const { readable, writable } = new TransformStream()
const stream = new SSEStreamingApi(writable, readable)
c.header('Transfer-Encoding', 'chunked')
c.header('Content-Type', 'text/event-stream')
c.header('Cache-Control', 'no-cache')
c.header('Connection', 'keep-alive')
run(stream, cb, onError)
return c.newResponse(stream.responseReadable)
}