Skip to content
Edge Functions

支持恢复的 WebSockets 与边缘函数

这个例子展示了如何在 Supabase Edge Functions 上构建一个可安全重新连接的聊天流,使用:

🌐 This example shows how to build a reconnect-safe chat stream on Supabase Edge Functions using:

  • WebSocket 升级 + JWT 认证
  • 基于 Postgres 的会话和事件持久化
  • 使用 lastEventId 重播事件
  • 带有 idempotency_key 的幂等用户消息
  • 在工作进程重启时客户端优雅地重新连接

参考实现:使用 Supabase Edge Functions 和 Postgres 构建可恢复的 WebSockets

🌐 Reference implementation: Building Resumable WebSockets with Supabase Edge Functions and Postgres

架构 #

🌐 Architecture

  1. 客户端使用用户 JWT 连接,还可以选择性地使用 sessionIdlastEventId
  2. 这个函数验证身份,然后要么恢复会话,要么创建一个新会话。
  3. 每条消息都是写给 ws_events 的,并且 id 会递增。
  4. 重新连接时,服务器会重放 id > lastEventId 的事件。
  5. 客户端更新本地 lastEventId 后继续运行,不会丢失消息。

数据库架构 #

🌐 Database schema

1
create extension if not exists pgcrypto;
2
3
create table ws_sessions (
4
id uuid primary key default gen_random_uuid(),
5
user_id uuid not null,
6
created_at timestamptz default now(),
7
updated_at timestamptz default now(),
8
last_event_id bigint default 0
9
);
10
11
create table ws_events (
12
id bigint generated by default as identity primary key,
13
session_id uuid not null references ws_sessions(id) on delete cascade,
14
event_type text not null,
15
payload jsonb not null,
16
created_at timestamptz default now()
17
);
18
create index ws_events_session_id_id_idx on ws_events(session_id, id);
19
20
create table ws_idempotency_keys (
21
session_id uuid not null references ws_sessions(id) on delete cascade,
22
idempotency_key uuid not null,
23
primary key(session_id, idempotency_key)
24
);
25
26
create unlogged table ws_live_connections (
27
session_id uuid primary key,
28
connected_at timestamptz default now(),
29
last_seen_at timestamptz default now(),
30
edge_region text
31
);

边缘函数(WebSocket 代理) #

🌐 Edge Function (WebSocket proxy)

在函数里使用 supabase functions serve --no-verify-jwt 并验证 JWT。

🌐 Use supabase functions serve --no-verify-jwt and validate JWT inside the function.

1
import { createAdminClient, createContextClient, verifyCredentials } from '@supabase/server/core'
2
3
const PREEMPTIVE_RESTART_MS = 340_000
4
5
function send(socket: WebSocket, payload: unknown) {
6
if (socket.readyState === WebSocket.OPEN) {
7
socket.send(JSON.stringify(payload))
8
}
9
}
10
11
Deno.serve(async (req) => {
12
const url = new URL(req.url)
13
const token = url.searchParams.get('token')
14
if (!token) return new Response('Missing token', { status: 401 })
15
16
const { data: auth, error } = await verifyCredentials({ token, apikey: null }, { auth: 'user' })
17
if (error || !auth?.userClaims?.id) {
18
return new Response('Unauthorized', { status: 401 })
19
}
20
21
const admin = createAdminClient()
22
const { socket, response } = Deno.upgradeWebSocket(req, { idleTimeout: 0 })
23
24
// Prevent EarlyDrop by keeping a pending promise until socket close.
25
let resolveClosed!: () => void
26
const closed = new Promise<void>((resolve) => {
27
resolveClosed = resolve
28
})
29
// @ts-ignore
30
EdgeRuntime.waitUntil(closed)
31
32
const requestedSessionId = url.searchParams.get('sessionId')
33
const lastEventId = Number(url.searchParams.get('lastEventId') || 0)
34
const sessionId = requestedSessionId ?? crypto.randomUUID()
35
36
socket.onclose = () => {
37
resolveClosed()
38
}
39
40
socket.onmessage = async (event) => {
41
const msg = JSON.parse(event.data)
42
43
if (msg.type === 'user_message') {
44
const { error: idempotencyError } = await admin.from('ws_idempotency_keys').upsert(
45
{
46
session_id: sessionId,
47
idempotency_key: msg.idempotency_key,
48
},
49
{ onConflict: 'session_id,idempotency_key', ignoreDuplicates: true }
50
)
51
52
let userEvent
53
54
if (idempotencyError) {
55
// Conflict detected - this is a retry, fetch the existing event
56
const { data: existingEvent } = await admin
57
.from('ws_events')
58
.select()
59
.eq('session_id', sessionId)
60
.eq('idempotency_key', msg.idempotency_key)
61
.single()
62
63
userEvent = existingEvent
64
} else {
65
// New idempotency key - insert the event
66
const { data: newEvent } = await admin
67
.from('ws_events')
68
.insert({
69
session_id: sessionId,
70
event_type: 'user_message',
71
payload: { content: msg.content },
72
})
73
.select()
74
.single()
75
76
userEvent = newEvent
77
}
78
79
send(socket, {
80
type: 'user_message',
81
payload: userEvent?.payload,
82
event_id: userEvent?.id,
83
})
84
}
85
}
86
87
send(socket, { type: 'session_init', session_id: sessionId })
88
89
queueMicrotask(async () => {
90
const { data: replayEvents } = await admin
91
.from('ws_events')
92
.select('*')
93
.eq('session_id', sessionId)
94
.gt('id', lastEventId)
95
.order('id')
96
97
for (const event of replayEvents ?? []) {
98
send(socket, {
99
type: event.event_type,
100
payload: event.payload,
101
event_id: event.id,
102
replay: true,
103
})
104
}
105
})
106
107
setTimeout(() => {
108
send(socket, { type: 'server_restarting' })
109
socket.close(1012, 'Service restart')
110
}, PREEMPTIVE_RESTART_MS)
111
112
return response
113
})

浏览器客户端 #

🌐 Browser client

客户端将 sessionIdlastEventId 存储在会话存储中,然后以指数退避方式重新连接。

🌐 The client stores sessionId and lastEventId in session storage, then reconnects with exponential backoff.

1
let sessionId = sessionStorage.getItem('ws_session_id')
2
let lastEventId = Number(sessionStorage.getItem('last_event_id') || 0)
3
4
function connect(token: string) {
5
const url =
6
`wss://YOUR_PROJECT.functions.supabase.co/websocket-proxy` +
7
`?token=${encodeURIComponent(token)}` +
8
`&lastEventId=${lastEventId}` +
9
(sessionId ? `&sessionId=${sessionId}` : '')
10
11
const ws = new WebSocket(url)
12
13
ws.onmessage = (e) => {
14
const msg = JSON.parse(e.data)
15
16
if (msg.event_id) {
17
lastEventId = Math.max(lastEventId, msg.event_id)
18
sessionStorage.setItem('last_event_id', String(lastEventId))
19
}
20
21
if (msg.type === 'session_init') {
22
sessionId = msg.session_id
23
sessionStorage.setItem('ws_session_id', sessionId)
24
}
25
}
26
}

为什么这个模式有效 #

🌐 Why this pattern works

  • 如果工作者重启,客户端会用同一个会话重新连接。
  • Replay 会关闭由重连窗口引起的交付差距。
  • 幂等键可以防止客户端重试时重复插入。
  • EdgeRuntime.waitUntil() 可以防止看起来空闲的 WebSocket 工作线程意外提前终止。

下一步 #

🌐 Next steps

  • 为所有 ws_* 表添加行级安全策略。
  • 为过时的会话添加心跳和清理策略。
  • 添加结构化事件负载类型和输入验证。
  • 为断开率和重放延迟添加可观测性仪表板。