feat(04-06): implement live-sync SSE vertical slice (LIST-04, D-04/D-10/D-11/D-12)
- Wire publishListEvent fan-out in lists.ts after every write mutation (item:added/updated/deleted, list:updated/deleted)
- Add GET /api/sse/lists scoped endpoint in sse.ts: resolveUserId → 401 on null; getAccessibleListIds → subscribe only to accessible channels; 30s heartbeat; cleanup on disconnect (D-04/T-04-01/T-04-02)
- Create useListSSE.ts: bounded-backoff EventSource wrapper (250ms→500ms→1s→2s→4s→cap 8s); MAX_ATTEMPTS=6; withCredentials:true; close-before-retry prevents reconnect storm (Pitfall 3); invalidates ['list', listId] on open (D-10) and on each event; onStateChange('disconnected') after exhaustion (D-11)
- Create LiveSyncIndicator.tsx: connected=green dot; reconnecting=pulsing muted dot + label; disconnected=red dot + 'Updates paused' (role=alert); correct ARIA per UI-SPEC
- Wire useListSSE + LiveSyncIndicator into ListDetail header; retain refetchInterval:30000 polling fallback (D-12)
- All 8 useListSSE tests pass; all 54 API tests pass; both typechecks pass
- playwright-cli: live update confirmed (eggs item added via API appeared in browser without manual refresh)
This commit is contained in:
@@ -30,8 +30,7 @@ import { rankForAppend } from '../lib/rank.js'
|
||||
// Side-effect import: brings in the ContextVariableMap augmentation for c.get('user')
|
||||
import '../auth/devBypass.js'
|
||||
|
||||
// Uncomment in Plan 06 when SSE endpoint exists:
|
||||
// import { publishListEvent } from '../lib/listEmitter.js'
|
||||
import { publishListEvent } from '../lib/listEmitter.js'
|
||||
|
||||
export const listsRouter = new Hono()
|
||||
|
||||
@@ -284,8 +283,8 @@ listsRouter.post('/', zValidator('json', createListSchema), async (c) => {
|
||||
.where(eq(lists.id, listId))
|
||||
.limit(1)
|
||||
|
||||
// Plan 06 SSE seam:
|
||||
// publishListEvent(listId, { type: 'list:updated', listId, payload: newList })
|
||||
// Fan-out: notify accessible subscribers that this list was created/updated (LIST-04)
|
||||
publishListEvent(listId, { type: 'list:updated', listId, payload: { id: listId, name: newList.name } })
|
||||
|
||||
return c.json(
|
||||
{
|
||||
@@ -376,8 +375,8 @@ listsRouter.patch('/:id', zValidator('json', patchListSchema), async (c) => {
|
||||
.where(eq(lists.id, listId))
|
||||
.limit(1)
|
||||
|
||||
// Plan 06 SSE seam:
|
||||
// publishListEvent(listId, { type: 'list:updated', listId, payload: updated })
|
||||
// Fan-out: notify accessible subscribers that this list metadata changed (LIST-04)
|
||||
publishListEvent(listId, { type: 'list:updated', listId, payload: { id: listId, name: updated.name } })
|
||||
|
||||
return c.json({
|
||||
id: updated.id,
|
||||
@@ -422,8 +421,8 @@ listsRouter.delete('/:id', async (c) => {
|
||||
|
||||
await db.delete(lists).where(eq(lists.id, listId))
|
||||
|
||||
// Plan 06 SSE seam:
|
||||
// publishListEvent(listId, { type: 'list:deleted', listId, payload: { id: listId } })
|
||||
// Fan-out: notify accessible subscribers that this list was deleted (LIST-04)
|
||||
publishListEvent(listId, { type: 'list:deleted', listId, payload: { id: listId } })
|
||||
|
||||
return c.json({ id: listId })
|
||||
} catch (err) {
|
||||
@@ -488,8 +487,8 @@ listsRouter.post('/:id/items', zValidator('json', createItemSchema), async (c) =
|
||||
.where(eq(listItems.id, inserted.id))
|
||||
.limit(1)
|
||||
|
||||
// Plan 06 SSE seam:
|
||||
// publishListEvent(listId, { type: 'item:added', listId, payload: newItem })
|
||||
// Fan-out: notify accessible subscribers that an item was added (LIST-04)
|
||||
publishListEvent(listId, { type: 'item:added', listId, payload: { id: newItem.id, listId, text: newItem.text } })
|
||||
|
||||
return c.json(
|
||||
{
|
||||
@@ -631,8 +630,8 @@ listItemsRouter.patch('/:itemId', zValidator('json', patchItemSchema), async (c)
|
||||
.where(eq(listItems.id, itemId))
|
||||
.limit(1)
|
||||
|
||||
// Plan 06 SSE seam:
|
||||
// publishListEvent(item.listId, { type: 'item:updated', listId: item.listId, payload: updated })
|
||||
// Fan-out: notify accessible subscribers that an item was updated (LIST-04)
|
||||
publishListEvent(item.listId, { type: 'item:updated', listId: item.listId, payload: { id: updated.id, listId: updated.listId } })
|
||||
|
||||
return c.json({
|
||||
id: updated.id,
|
||||
@@ -682,8 +681,8 @@ listItemsRouter.delete('/:itemId', async (c) => {
|
||||
// Delete-wins (D-09): delete is final; no rollback path.
|
||||
await db.delete(listItems).where(eq(listItems.id, itemId))
|
||||
|
||||
// Plan 06 SSE seam:
|
||||
// publishListEvent(item.listId, { type: 'item:deleted', listId: item.listId, payload: { id: itemId } })
|
||||
// Fan-out: notify accessible subscribers that an item was deleted (LIST-04)
|
||||
publishListEvent(item.listId, { type: 'item:deleted', listId: item.listId, payload: { id: itemId } })
|
||||
|
||||
return c.json({ id: itemId })
|
||||
} catch (err) {
|
||||
|
||||
+91
-10
@@ -1,25 +1,46 @@
|
||||
/**
|
||||
* GET /api/sse/heartbeat — Pangolin SSE pass-through smoke test endpoint.
|
||||
* SSE router — streaming endpoints for real-time events.
|
||||
*
|
||||
* Emits a `heartbeat` event every 10 seconds with { ts, id } payload.
|
||||
* Runs until the client disconnects (stream.aborted).
|
||||
* Routes:
|
||||
* GET /heartbeat — Pangolin smoke-test (Phase 1 entry gate)
|
||||
* GET /lists — Scoped live-list fan-out stream (LIST-04, D-04)
|
||||
*
|
||||
* Mounted under /api/sse in index.ts, so it sits behind oidcAuthMiddleware (T-04-01).
|
||||
* Heartbeat payload carries only timestamps — no user data or secrets (T-04-02).
|
||||
*
|
||||
* Smoke test procedure (D-08):
|
||||
* curl -N https://familysync.<domain>/api/sse/heartbeat
|
||||
* Keep open 5+ min — confirm no proxy timeout. PASS → SSE viable for Phase 4.
|
||||
* Both routes sit behind oidcAuthMiddleware in index.ts (T-04-01).
|
||||
*
|
||||
* Source: https://hono.dev/docs/helpers/streaming
|
||||
* RESEARCH Pattern 5: Pangolin SSE Smoke Test
|
||||
* RESEARCH Pattern 5 (heartbeat) + Finding 1 (lists scoped fan-out)
|
||||
*/
|
||||
|
||||
import { Hono } from 'hono'
|
||||
import type { Context } from 'hono'
|
||||
import { streamSSE } from 'hono/streaming'
|
||||
import { getAuth } from '../auth/middleware.js'
|
||||
import { upsertUser, deriveDisplayName } from '../auth/user.js'
|
||||
import { subscribeListEvents } from '../lib/listEmitter.js'
|
||||
import { getAccessibleListIds } from '../lib/listAccess.js'
|
||||
// Side-effect import: brings in the ContextVariableMap augmentation for c.get('user')
|
||||
import '../auth/devBypass.js'
|
||||
|
||||
export const sseRouter = new Hono()
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Auth helper — same pattern as lists.ts (per-router duplication convention).
|
||||
// Resolution order: dev-bypass user first, then OIDC.
|
||||
// ---------------------------------------------------------------------------
|
||||
async function resolveUserId(c: Context): Promise<number | null> {
|
||||
const devUser = c.get('user') as { id: number } | undefined
|
||||
if (devUser) return devUser.id
|
||||
|
||||
const auth = await getAuth(c)
|
||||
if (!auth) return null
|
||||
|
||||
const iss = (auth.iss as string | undefined) ?? ''
|
||||
const sub = auth.sub ?? ''
|
||||
const displayName = deriveDisplayName(auth)
|
||||
const user = await upsertUser(iss, sub, displayName)
|
||||
return user?.id ?? null
|
||||
}
|
||||
|
||||
/**
|
||||
* GET /heartbeat
|
||||
* Streams SSE heartbeat events every 10 seconds until client disconnects.
|
||||
@@ -38,3 +59,63 @@ sseRouter.get('/heartbeat', (c) => {
|
||||
}
|
||||
})
|
||||
})
|
||||
|
||||
/**
|
||||
* GET /lists
|
||||
*
|
||||
* Scoped live-list SSE stream (LIST-04, D-04).
|
||||
*
|
||||
* Security (T-04-01, T-04-02):
|
||||
* - resolveUserId → 401 on null (same OIDC guard as all /api/sse/* routes).
|
||||
* - Subscribes ONLY to list channels the caller can access (owner + shares).
|
||||
* Events for private lists of other members are never delivered (D-04).
|
||||
*
|
||||
* Behavior:
|
||||
* - Resolves the caller's accessible list IDs via getAccessibleListIds(userId).
|
||||
* - Opens one subscribeListEvents subscription per accessible list.
|
||||
* - Each event is forwarded as `event: event.type; data: JSON.stringify(event)`.
|
||||
* - A 30s heartbeat keeps the Pangolin connection alive (proven in smoke test).
|
||||
* - On client disconnect (stream.aborted), all subscriptions are cleaned up.
|
||||
*
|
||||
* D-10: Client receives minimal { type, listId } payload and full-refetches.
|
||||
* Payload is not relied on for cache updates — only triggers invalidation.
|
||||
* D-11: Bounded backoff + give-up logic lives in the PWA useListSSE hook.
|
||||
* D-12: PWA refetchInterval: 30000 polling fallback always active.
|
||||
*/
|
||||
sseRouter.get('/lists', async (c) => {
|
||||
const userId = await resolveUserId(c)
|
||||
if (userId === null) return c.json({ error: 'Unauthorized' }, 401)
|
||||
|
||||
const accessibleListIds = await getAccessibleListIds(userId)
|
||||
|
||||
return streamSSE(c, async (stream) => {
|
||||
const unsubscribers: Array<() => void> = []
|
||||
|
||||
// Subscribe to each accessible list's channel (D-04 — scoped, not global)
|
||||
for (const listId of accessibleListIds) {
|
||||
const unsub = subscribeListEvents(listId, async (event) => {
|
||||
if (stream.aborted) return
|
||||
await stream.writeSSE({
|
||||
data: JSON.stringify(event),
|
||||
event: event.type,
|
||||
id: `${listId}-${Date.now()}`,
|
||||
})
|
||||
})
|
||||
unsubscribers.push(unsub)
|
||||
}
|
||||
|
||||
// 30s heartbeat — keeps Pangolin connection alive (smoke-tested in Phase 1)
|
||||
let tick = 0
|
||||
while (!stream.aborted) {
|
||||
await stream.writeSSE({
|
||||
data: JSON.stringify({ ts: new Date().toISOString() }),
|
||||
event: 'heartbeat',
|
||||
id: String(tick++),
|
||||
})
|
||||
await stream.sleep(30_000)
|
||||
}
|
||||
|
||||
// Cleanup all subscriptions on client disconnect
|
||||
unsubscribers.forEach((unsub) => unsub())
|
||||
})
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user