支持恢复的 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
- 客户端使用用户 JWT 连接,还可以选择性地使用
sessionId和lastEventId。 - 这个函数验证身份,然后要么恢复会话,要么创建一个新会话。
- 每条消息都是写给
ws_events的,并且id会递增。 - 重新连接时,服务器会重放
id > lastEventId的事件。 - 客户端更新本地
lastEventId后继续运行,不会丢失消息。
数据库架构 #
🌐 Database schema
1create extension if not exists pgcrypto;23create 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 09);1011create 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);18create index ws_events_session_id_id_idx on ws_events(session_id, id);1920create 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);2526create 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 text31);边缘函数(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.
1import { createAdminClient, createContextClient, verifyCredentials } from '@supabase/server/core'23const PREEMPTIVE_RESTART_MS = 340_00045function send(socket: WebSocket, payload: unknown) {6 if (socket.readyState === WebSocket.OPEN) {7 socket.send(JSON.stringify(payload))8 }9}1011Deno.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 })1516 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 }2021 const admin = createAdminClient()22 const { socket, response } = Deno.upgradeWebSocket(req, { idleTimeout: 0 })2324 // Prevent EarlyDrop by keeping a pending promise until socket close.25 let resolveClosed!: () => void26 const closed = new Promise<void>((resolve) => {27 resolveClosed = resolve28 })29 // @ts-ignore30 EdgeRuntime.waitUntil(closed)3132 const requestedSessionId = url.searchParams.get('sessionId')33 const lastEventId = Number(url.searchParams.get('lastEventId') || 0)34 const sessionId = requestedSessionId ?? crypto.randomUUID()3536 socket.onclose = () => {37 resolveClosed()38 }3940 socket.onmessage = async (event) => {41 const msg = JSON.parse(event.data)4243 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 )5152 let userEvent5354 if (idempotencyError) {55 // Conflict detected - this is a retry, fetch the existing event56 const { data: existingEvent } = await admin57 .from('ws_events')58 .select()59 .eq('session_id', sessionId)60 .eq('idempotency_key', msg.idempotency_key)61 .single()6263 userEvent = existingEvent64 } else {65 // New idempotency key - insert the event66 const { data: newEvent } = await admin67 .from('ws_events')68 .insert({69 session_id: sessionId,70 event_type: 'user_message',71 payload: { content: msg.content },72 })73 .select()74 .single()7576 userEvent = newEvent77 }7879 send(socket, {80 type: 'user_message',81 payload: userEvent?.payload,82 event_id: userEvent?.id,83 })84 }85 }8687 send(socket, { type: 'session_init', session_id: sessionId })8889 queueMicrotask(async () => {90 const { data: replayEvents } = await admin91 .from('ws_events')92 .select('*')93 .eq('session_id', sessionId)94 .gt('id', lastEventId)95 .order('id')9697 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 })106107 setTimeout(() => {108 send(socket, { type: 'server_restarting' })109 socket.close(1012, 'Service restart')110 }, PREEMPTIVE_RESTART_MS)111112 return response113})浏览器客户端 #
🌐 Browser client
客户端将 sessionId 和 lastEventId 存储在会话存储中,然后以指数退避方式重新连接。
🌐 The client stores sessionId and lastEventId in session storage, then reconnects with exponential backoff.
1let sessionId = sessionStorage.getItem('ws_session_id')2let lastEventId = Number(sessionStorage.getItem('last_event_id') || 0)34function 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}` : '')1011 const ws = new WebSocket(url)1213 ws.onmessage = (e) => {14 const msg = JSON.parse(e.data)1516 if (msg.event_id) {17 lastEventId = Math.max(lastEventId, msg.event_id)18 sessionStorage.setItem('last_event_id', String(lastEventId))19 }2021 if (msg.type === 'session_init') {22 sessionId = msg.session_id23 sessionStorage.setItem('ws_session_id', sessionId)24 }25 }26}为什么这个模式有效 #
🌐 Why this pattern works
- 如果工作者重启,客户端会用同一个会话重新连接。
- Replay 会关闭由重连窗口引起的交付差距。
- 幂等键可以防止客户端重试时重复插入。
EdgeRuntime.waitUntil()可以防止看起来空闲的 WebSocket 工作线程意外提前终止。
下一步 #
🌐 Next steps
- 为所有
ws_*表添加行级安全策略。 - 为过时的会话添加心跳和清理策略。
- 添加结构化事件负载类型和输入验证。
- 为断开率和重放延迟添加可观测性仪表板。