mirror of
https://github.com/supabase/supabase.git
synced 2026-10-05 09:25:06 +03:00
feat(docs): add resumable WebSockets + Edge Functions troubleshooting guides (#46178)
## I have read the [CONTRIBUTING.md](https://github.com/supabase/supabase/blob/master/CONTRIBUTING.md) file. YES ## What kind of change does this PR introduce? Docs update (new guides + follow-up documentation fix from review feedback). ## What is the current behavior? There was no consolidated docs example for resumable WebSockets with Edge Functions, and no dedicated troubleshooting guide for worker timeouts / WebSocket drops. ## What is the new behavior? - Adds a resumable WebSockets guide for Edge Functions, including: - session persistence - event replay - idempotency pattern and schema examples - client/server example flow - Adds an Edge Functions troubleshooting guide for worker timeouts and WebSocket drops. - Updates docs navigation to surface the new guides. - Follow-up fix from review feedback: the browser client example now stores `sessionId` and `lastEventId` in `sessionStorage` (instead of `localStorage`). ## Additional context - Branch has been updated with latest `origin/master`. - This PR remains documentation-focused; no production runtime code changes were introduced. <!-- This is an auto-generated comment: release notes by coderabbit.ai --> ## Summary by CodeRabbit * **Documentation** * Added a guide on resumable WebSockets covering session persistence, event replay, idempotency patterns, SQL schema examples, and client/server usage. * Added a troubleshooting guide on Edge Functions worker timeouts and WebSocket drops with scenarios, symptoms, and practical workarounds. * Enhanced WebSocket docs with a production note on worker lifecycle and keeping runtime promises open to avoid premature shutdown. * Navigation updated to surface the new guides. <!-- review_stack_entry_start --> [](https://app.coderabbit.ai/change-stack/supabase/supabase/pull/46178?utm_source=github_walkthrough&utm_medium=github&utm_campaign=change_stack) <!-- review_stack_entry_end --> <!-- end of auto-generated comment: release notes by coderabbit.ai --> --------- Co-authored-by: Lakshan Perera <lakshan@supabase.io> Co-authored-by: coderabbitai[bot] <136622811+coderabbitai[bot]@users.noreply.github.com> Co-authored-by: CodeRabbit <noreply@coderabbit.ai> Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com>
This commit is contained in:
4 files changed
+406
No files matched your search
@@ -1656,6 +1656,10 @@ export const functions: NavMenuConstant = {
|
||||
name: 'Troubleshooting',
|
||||
url: '/guides/functions/troubleshooting' as `/${string}`,
|
||||
},
|
||||
{
|
||||
name: 'Worker timeouts and WebSocket drops',
|
||||
url: '/troubleshooting/edge-functions-worker-timeouts-and-websocket-drops' as `/${string}`,
|
||||
},
|
||||
],
|
||||
},
|
||||
{
|
||||
@@ -1783,6 +1787,10 @@ export const functions: NavMenuConstant = {
|
||||
name: 'Image Transformation & Optimization',
|
||||
url: '/guides/functions/examples/image-manipulation' as `/${string}`,
|
||||
},
|
||||
{
|
||||
name: 'Resumable WebSockets with replay',
|
||||
url: '/guides/functions/examples/resumable-websockets' as `/${string}`,
|
||||
},
|
||||
],
|
||||
},
|
||||
{
|
||||
|
||||
@@ -0,0 +1,225 @@
|
||||
---
|
||||
title: 'Resumable WebSockets with Edge Functions'
|
||||
description: 'Build reconnect-safe WebSockets with event replay, idempotency keys, and graceful restarts.'
|
||||
---
|
||||
|
||||
This example shows how to build a reconnect-safe chat stream on Supabase Edge Functions using:
|
||||
|
||||
- WebSocket upgrade + JWT auth
|
||||
- Postgres-backed session and event persistence
|
||||
- Event replay with `lastEventId`
|
||||
- Idempotent user messages with `idempotency_key`
|
||||
- Graceful client reconnects during worker restarts
|
||||
|
||||
Reference implementation: [Building Resumable WebSockets with Supabase Edge Functions and Postgres](https://blog.mansueli.com/building-resumable-websockets-with-supabase-edge-functions-and-postgres)
|
||||
|
||||
## Architecture
|
||||
|
||||
1. Client connects with a user JWT, plus optional `sessionId` and `lastEventId`.
|
||||
2. Function verifies auth and either resumes or creates a session.
|
||||
3. Every message is written to `ws_events` with an incrementing `id`.
|
||||
4. On reconnect, server replays events where `id > lastEventId`.
|
||||
5. Client updates local `lastEventId` and resumes without losing messages.
|
||||
|
||||
## Database schema
|
||||
|
||||
```sql
|
||||
create extension if not exists pgcrypto;
|
||||
|
||||
create table ws_sessions (
|
||||
id uuid primary key default gen_random_uuid(),
|
||||
user_id uuid not null,
|
||||
created_at timestamptz default now(),
|
||||
updated_at timestamptz default now(),
|
||||
last_event_id bigint default 0
|
||||
);
|
||||
|
||||
create table ws_events (
|
||||
id bigint generated by default as identity primary key,
|
||||
session_id uuid not null references ws_sessions(id) on delete cascade,
|
||||
event_type text not null,
|
||||
payload jsonb not null,
|
||||
created_at timestamptz default now()
|
||||
);
|
||||
create index ws_events_session_id_id_idx on ws_events(session_id, id);
|
||||
|
||||
create table ws_idempotency_keys (
|
||||
session_id uuid not null references ws_sessions(id) on delete cascade,
|
||||
idempotency_key uuid not null,
|
||||
primary key(session_id, idempotency_key)
|
||||
);
|
||||
|
||||
create unlogged table ws_live_connections (
|
||||
session_id uuid primary key,
|
||||
connected_at timestamptz default now(),
|
||||
last_seen_at timestamptz default now(),
|
||||
edge_region text
|
||||
);
|
||||
```
|
||||
|
||||
## Edge Function (WebSocket proxy)
|
||||
|
||||
Use `supabase functions serve --no-verify-jwt` and validate JWT inside the function.
|
||||
|
||||
```ts
|
||||
import { createAdminClient, createContextClient, verifyCredentials } from '@supabase/server/core'
|
||||
|
||||
const PREEMPTIVE_RESTART_MS = 340_000
|
||||
|
||||
function send(socket: WebSocket, payload: unknown) {
|
||||
if (socket.readyState === WebSocket.OPEN) {
|
||||
socket.send(JSON.stringify(payload))
|
||||
}
|
||||
}
|
||||
|
||||
Deno.serve(async (req) => {
|
||||
const url = new URL(req.url)
|
||||
const token = url.searchParams.get('token')
|
||||
if (!token) return new Response('Missing token', { status: 401 })
|
||||
|
||||
const { data: auth, error } = await verifyCredentials({ token, apikey: null }, { auth: 'user' })
|
||||
if (error || !auth?.userClaims?.id) {
|
||||
return new Response('Unauthorized', { status: 401 })
|
||||
}
|
||||
|
||||
const admin = createAdminClient()
|
||||
const { socket, response } = Deno.upgradeWebSocket(req, { idleTimeout: 0 })
|
||||
|
||||
// Prevent EarlyDrop by keeping a pending promise until socket close.
|
||||
let resolveClosed!: () => void
|
||||
const closed = new Promise<void>((resolve) => {
|
||||
resolveClosed = resolve
|
||||
})
|
||||
// @ts-ignore
|
||||
EdgeRuntime.waitUntil(closed)
|
||||
|
||||
const requestedSessionId = url.searchParams.get('sessionId')
|
||||
const lastEventId = Number(url.searchParams.get('lastEventId') || 0)
|
||||
const sessionId = requestedSessionId ?? crypto.randomUUID()
|
||||
|
||||
socket.onclose = () => {
|
||||
resolveClosed()
|
||||
}
|
||||
|
||||
socket.onmessage = async (event) => {
|
||||
const msg = JSON.parse(event.data)
|
||||
|
||||
if (msg.type === 'user_message') {
|
||||
const { error: idempotencyError } = await admin.from('ws_idempotency_keys').upsert(
|
||||
{
|
||||
session_id: sessionId,
|
||||
idempotency_key: msg.idempotency_key,
|
||||
},
|
||||
{ onConflict: 'session_id,idempotency_key', ignoreDuplicates: true }
|
||||
)
|
||||
|
||||
let userEvent
|
||||
|
||||
if (idempotencyError) {
|
||||
// Conflict detected - this is a retry, fetch the existing event
|
||||
const { data: existingEvent } = await admin
|
||||
.from('ws_events')
|
||||
.select()
|
||||
.eq('session_id', sessionId)
|
||||
.eq('idempotency_key', msg.idempotency_key)
|
||||
.single()
|
||||
|
||||
userEvent = existingEvent
|
||||
} else {
|
||||
// New idempotency key - insert the event
|
||||
const { data: newEvent } = await admin
|
||||
.from('ws_events')
|
||||
.insert({
|
||||
session_id: sessionId,
|
||||
event_type: 'user_message',
|
||||
payload: { content: msg.content },
|
||||
})
|
||||
.select()
|
||||
.single()
|
||||
|
||||
userEvent = newEvent
|
||||
}
|
||||
|
||||
send(socket, {
|
||||
type: 'user_message',
|
||||
payload: userEvent?.payload,
|
||||
event_id: userEvent?.id,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
send(socket, { type: 'session_init', session_id: sessionId })
|
||||
|
||||
queueMicrotask(async () => {
|
||||
const { data: replayEvents } = await admin
|
||||
.from('ws_events')
|
||||
.select('*')
|
||||
.eq('session_id', sessionId)
|
||||
.gt('id', lastEventId)
|
||||
.order('id')
|
||||
|
||||
for (const event of replayEvents ?? []) {
|
||||
send(socket, {
|
||||
type: event.event_type,
|
||||
payload: event.payload,
|
||||
event_id: event.id,
|
||||
replay: true,
|
||||
})
|
||||
}
|
||||
})
|
||||
|
||||
setTimeout(() => {
|
||||
send(socket, { type: 'server_restarting' })
|
||||
socket.close(1012, 'Service restart')
|
||||
}, PREEMPTIVE_RESTART_MS)
|
||||
|
||||
return response
|
||||
})
|
||||
```
|
||||
|
||||
## Browser client
|
||||
|
||||
The client stores `sessionId` and `lastEventId` in session storage, then reconnects with exponential backoff.
|
||||
|
||||
```ts
|
||||
let sessionId = sessionStorage.getItem('ws_session_id')
|
||||
let lastEventId = Number(sessionStorage.getItem('last_event_id') || 0)
|
||||
|
||||
function connect(token: string) {
|
||||
const url =
|
||||
`wss://YOUR_PROJECT.functions.supabase.co/websocket-proxy` +
|
||||
`?token=${encodeURIComponent(token)}` +
|
||||
`&lastEventId=${lastEventId}` +
|
||||
(sessionId ? `&sessionId=${sessionId}` : '')
|
||||
|
||||
const ws = new WebSocket(url)
|
||||
|
||||
ws.onmessage = (e) => {
|
||||
const msg = JSON.parse(e.data)
|
||||
|
||||
if (msg.event_id) {
|
||||
lastEventId = Math.max(lastEventId, msg.event_id)
|
||||
sessionStorage.setItem('last_event_id', String(lastEventId))
|
||||
}
|
||||
|
||||
if (msg.type === 'session_init') {
|
||||
sessionId = msg.session_id
|
||||
sessionStorage.setItem('ws_session_id', sessionId)
|
||||
}
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
## Why this pattern works
|
||||
|
||||
- If the worker restarts, the client reconnects with the same session.
|
||||
- Replay closes delivery gaps caused by reconnect windows.
|
||||
- Idempotency keys prevent duplicate inserts when clients retry.
|
||||
- `EdgeRuntime.waitUntil()` prevents unexpected early termination of idle-looking WebSocket workers.
|
||||
|
||||
## Next steps
|
||||
|
||||
- Add row-level security policies for all `ws_*` tables.
|
||||
- Add a heartbeat and cleanup policy for stale sessions.
|
||||
- Add structured event payload types and input validation.
|
||||
- Add observability dashboards for disconnect rate and replay lag.
|
||||
@@ -13,6 +13,8 @@ This allows you to:
|
||||
- Create WebSocket relay servers for external APIs
|
||||
- Establish both incoming and outgoing WebSocket connections
|
||||
|
||||
For a production-ready reconnect pattern with session persistence and replay, see [Resumable WebSockets with Edge Functions](/docs/guides/functions/examples/resumable-websockets).
|
||||
|
||||
---
|
||||
|
||||
## Creating WebSocket servers
|
||||
@@ -249,6 +251,8 @@ The maximum duration is capped based on the wall-clock, CPU, and memory limits.
|
||||
|
||||
</Admonition>
|
||||
|
||||
When using WebSockets, keep in mind that the HTTP request is considered complete after `Deno.upgradeWebSocket(req)` returns the response. To prevent early worker retirement while the socket is still open, keep an unresolved `EdgeRuntime.waitUntil()` promise that resolves in `socket.onclose`.
|
||||
|
||||
---
|
||||
|
||||
## Testing WebSockets locally
|
||||
|
||||
+169
@@ -0,0 +1,169 @@
|
||||
---
|
||||
title = "Edge Functions worker timeouts and WebSocket drops"
|
||||
topics = [ "functions" ]
|
||||
keywords = [ "websocket", "timeout", "earlydrop", "wall clock", "cpu limit", "streaming", "cold start" ]
|
||||
teams = [ "team-functions", "team-support" ]
|
||||
types = [ "support" ]
|
||||
---
|
||||
|
||||
## Background
|
||||
|
||||
Edge Functions run inside V8 isolates managed by a supervisor in `edge-runtime`.
|
||||
|
||||
The supervisor enforces resource limits:
|
||||
|
||||
- Wall clock time (`worker_timeout_ms`)
|
||||
- CPU time (soft + hard limits)
|
||||
- Memory usage
|
||||
|
||||
It may retire early (`EarlyDrop` event) if the isolate is idle Isolate is considered idle if following conditions are met:
|
||||
|
||||
- The HTTP response has already been returned.
|
||||
- All `EdgeRuntime.waitUntil()` promises have resolved.
|
||||
|
||||
If both are true during a resource check, the isolate can be terminated even with open WebSocket connections.
|
||||
|
||||
## Scenario 1: WebSocket drops around half of the wall clock limit
|
||||
|
||||
### Symptoms
|
||||
|
||||
- WebSocket closes at a consistent interval (often around half of the configured wall clock limit).
|
||||
- Logs include `EarlyDrop` or wall clock warning events.
|
||||
|
||||
### Root cause
|
||||
|
||||
After `Deno.upgradeWebSocket(req)` returns a response, the HTTP request is considered acknowledged. If there is no unresolved `waitUntil` work, the worker may look idle and be retired early.
|
||||
|
||||
### What to check
|
||||
|
||||
1. Compare connection lifetime to your wall clock limit.
|
||||
2. Inspect logs for `EarlyDrop` around disconnect time.
|
||||
3. Confirm there is no unresolved `EdgeRuntime.waitUntil()` promise tied to socket lifecycle.
|
||||
|
||||
### Workaround
|
||||
|
||||
Keep a promise pending until the socket closes.
|
||||
|
||||
```ts
|
||||
Deno.serve((req) => {
|
||||
const { socket, response } = Deno.upgradeWebSocket(req)
|
||||
|
||||
const socketClosedPromise = new Promise<void>((resolve) => {
|
||||
socket.onclose = () => resolve()
|
||||
})
|
||||
|
||||
EdgeRuntime.waitUntil(socketClosedPromise)
|
||||
|
||||
socket.onmessage = (event) => {
|
||||
socket.send(event.data)
|
||||
}
|
||||
|
||||
return response
|
||||
})
|
||||
```
|
||||
|
||||
`EdgeRuntime.waitUntil()` prevents early retirement, but it does not extend the hard wall clock limit.
|
||||
|
||||
## Scenario 2: Function killed at a consistent duration
|
||||
|
||||
### Symptoms
|
||||
|
||||
- Function fails at a predictable runtime (for example, always around the same second mark).
|
||||
- Logs include wall clock shutdown reasons.
|
||||
- Clients may receive status `546` or cancellation errors.
|
||||
|
||||
### Root cause
|
||||
|
||||
The function exceeded the configured wall clock budget.
|
||||
|
||||
### Workarounds
|
||||
|
||||
- Split work into smaller units.
|
||||
- Move long work to async/background processing and return early.
|
||||
- Use streaming from upstream APIs where possible.
|
||||
- Use queues (`pg_net`, `pgmq`, or webhooks) for chunked processing.
|
||||
|
||||
## Scenario 3: Function killed by CPU limit
|
||||
|
||||
### Symptoms
|
||||
|
||||
- Failures during compute-heavy tasks.
|
||||
- Logs include CPU soft/hard limit events.
|
||||
|
||||
### Root cause
|
||||
|
||||
CPU budget and wall clock budget are independent. A function can run out of CPU time long before wall clock is exhausted.
|
||||
|
||||
### Workarounds
|
||||
|
||||
- Break large synchronous loops into async chunks.
|
||||
- Optimize expensive paths and avoid repeated recalculation.
|
||||
- Move heavy compute to systems designed for long CPU-bound workloads.
|
||||
|
||||
## Scenario 4: SSE or AI streams end before completion
|
||||
|
||||
### Symptoms
|
||||
|
||||
- Streaming starts but ends prematurely.
|
||||
- No final `[DONE]` token or normal close marker.
|
||||
|
||||
### Root cause
|
||||
|
||||
The worker hits wall clock or early retirement conditions while forwarding a long stream.
|
||||
|
||||
### Workaround
|
||||
|
||||
Keep the isolate alive for the stream piping lifecycle:
|
||||
|
||||
```ts
|
||||
Deno.serve(async (_req) => {
|
||||
const upstream = await fetch('https://api.openai.com/v1/chat/completions', {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
'Content-Type': 'application/json',
|
||||
Authorization: `Bearer ${Deno.env.get('OPENAI_API_KEY')}`,
|
||||
},
|
||||
body: JSON.stringify({ stream: true }),
|
||||
})
|
||||
|
||||
const { readable, writable } = new TransformStream()
|
||||
|
||||
EdgeRuntime.waitUntil(upstream.body!.pipeTo(writable))
|
||||
|
||||
return new Response(readable, {
|
||||
headers: { 'Content-Type': 'text/event-stream' },
|
||||
})
|
||||
})
|
||||
```
|
||||
|
||||
## Scenario 5: Cold starts fail before first response
|
||||
|
||||
### Symptoms
|
||||
|
||||
- First request after idle fails (for example, `504` or worker creation timeout).
|
||||
- Subsequent requests may succeed.
|
||||
|
||||
### Root cause
|
||||
|
||||
Large dependency trees or expensive top-level initialization can exceed startup budget.
|
||||
|
||||
### Workarounds
|
||||
|
||||
- Avoid slow top-level `await` work.
|
||||
- Lazy-initialize heavy clients inside request handlers.
|
||||
- Reduce bundle and dependency size.
|
||||
- Keep critical functions warm if needed.
|
||||
|
||||
## Key distinctions
|
||||
|
||||
- `EarlyDrop`: early retirement when worker appears idle.
|
||||
- `WallClockTime`: hard runtime ceiling reached.
|
||||
- CPU and wall clock limits are independent.
|
||||
- WebSockets are not request-tracked by the supervisor after upgrade.
|
||||
|
||||
## Related resources
|
||||
|
||||
- [Edge Functions limits](/docs/guides/functions/limits)
|
||||
- [Edge Function shutdown reasons explained](./edge-function-shutdown-reasons-explained)
|
||||
- [Monitoring Edge Function resource usage](./edge-function-monitoring-resource-usage)
|
||||
- [Handling WebSockets in Edge Functions](/docs/guides/functions/websockets)
|
||||
Reference in new issue
Block a user