Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
152 changes: 13 additions & 139 deletions apps/sim/app/api/table/[tableId]/events/stream/route.ts
Original file line number Diff line number Diff line change
@@ -1,26 +1,13 @@
import { createLogger } from '@sim/logger'
import { toError } from '@sim/utils/errors'
import { sleep } from '@sim/utils/helpers'
import { type NextRequest, NextResponse } from 'next/server'
import { tableEventStreamContract } from '@/lib/api/contracts/tables'
import { parseRequest } from '@/lib/api/server'
import { checkSessionOrInternalAuth } from '@/lib/auth/hybrid'
import { generateRequestId } from '@/lib/core/utils/request'
import { SSE_HEADERS } from '@/lib/core/utils/sse'
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
import {
getLatestTableEventId,
readTableEventsSince,
type TableEventEntry,
} from '@/lib/table/events'
import { createEventStreamResponse } from '@/lib/realtime/event-stream-route'
import { getLatestTableEventId, readTableEventsSince } from '@/lib/table/events'
import { accessError, checkAccess } from '@/app/api/table/utils'

const logger = createLogger('TableEventStreamAPI')

const POLL_INTERVAL_MS = 500
const HEARTBEAT_INTERVAL_MS = 15_000
const MAX_STREAM_DURATION_MS = 4 * 60 * 60 * 1000 // 4 hours; client reconnects past this

export const runtime = 'nodejs'
export const dynamic = 'force-dynamic'

Expand All @@ -30,12 +17,9 @@ interface RouteContext {

/** GET /api/table/[tableId]/events/stream?from=<lastEventId>
*
* SSE stream of cell-state transitions. Replay-on-reconnect via `from`;
* absent `from` tails from the latest event id (fresh mount — the client has
* just fetched current state, so replaying history would rewind it).
* Pruning (buffer cap exceeded or TTL expired) sends a `pruned` event and
* closes; the client responds with a full row-query refetch and reconnects
* tailing from latest. */
* SSE stream of cell-state transitions over the shared durable event log. Auth
* and access are checked here; the replay/tail/poll/heartbeat/prune mechanics
* come from `createEventStreamResponse`. */
export const GET = withRouteHandler(async (req: NextRequest, context: RouteContext) => {
const requestId = generateRequestId()
const parsed = await parseRequest(tableEventStreamContract, req, context)
Expand All @@ -51,123 +35,13 @@ export const GET = withRouteHandler(async (req: NextRequest, context: RouteConte
const access = await checkAccess(tableId, auth.userId, 'read')
if (!access.ok) return accessError(access, requestId, tableId)

logger.info(`[${requestId}] Table event stream opened`, { tableId, fromEventId })

const encoder = new TextEncoder()
let closed = false

const stream = new ReadableStream<Uint8Array>({
async start(controller) {
let lastEventId = fromEventId ?? 0
const deadline = Date.now() + MAX_STREAM_DURATION_MS
let nextHeartbeatAt = Date.now() + HEARTBEAT_INTERVAL_MS

const enqueue = (text: string) => {
if (closed) return
try {
controller.enqueue(encoder.encode(text))
} catch {
closed = true
}
}

const sendEvents = (events: TableEventEntry[]) => {
for (const entry of events) {
if (closed) return
enqueue(`data: ${JSON.stringify(entry)}\n\n`)
lastEventId = entry.eventId
}
}

const sendPrunedAndClose = (earliestEventId: number | undefined) => {
enqueue(
`event: pruned\ndata: ${JSON.stringify({ earliestEventId: earliestEventId ?? null })}\n\n`
)
if (!closed) {
closed = true
try {
controller.close()
} catch {}
}
}

const sendHeartbeat = () => {
// SSE comment line — keeps proxies (ALB default 60s idle) from closing
// the connection during quiet periods.
enqueue(`: ping ${Date.now()}\n\n`)
}

try {
// No replay cursor → tail from the latest event id. Resolved inside
// the try so a Redis failure errors the stream (client reconnects
// with backoff) rather than silently replaying the whole buffer.
if (fromEventId === undefined) {
lastEventId = await getLatestTableEventId(tableId)
}
// Initial replay from buffer.
const initial = await readTableEventsSince(tableId, lastEventId)
if (initial.status === 'pruned') {
sendPrunedAndClose(initial.earliestEventId)
return
}
if (initial.status === 'unavailable') {
throw new Error(`Table event buffer unavailable: ${initial.error}`)
}
sendEvents(initial.events)

// Stream loop — poll the buffer and forward new events. Workflow
// execution stream uses the same shape; pub/sub wakeups are an
// optimization we can add later if 500ms polling becomes a problem.
while (!closed && Date.now() < deadline) {
await sleep(POLL_INTERVAL_MS)
if (closed) return

const result = await readTableEventsSince(tableId, lastEventId)
if (result.status === 'pruned') {
sendPrunedAndClose(result.earliestEventId)
return
}
if (result.status === 'unavailable') {
throw new Error(`Table event buffer unavailable: ${result.error}`)
}
if (result.events.length > 0) {
sendEvents(result.events)
}

if (Date.now() >= nextHeartbeatAt) {
sendHeartbeat()
nextHeartbeatAt = Date.now() + HEARTBEAT_INTERVAL_MS
}
}

// Reached the defensive duration ceiling — close cleanly so the client
// reconnects with the latest lastEventId.
if (!closed) {
enqueue(`event: rotate\ndata: {}\n\n`)
closed = true
try {
controller.close()
} catch {}
}
} catch (error) {
logger.error(`[${requestId}] Table event stream error`, {
tableId,
error: toError(error).message,
})
if (!closed) {
try {
controller.error(error)
} catch {}
}
}
},
cancel() {
closed = true
logger.info(`[${requestId}] Client disconnected from table event stream`, { tableId })
},
})

return new NextResponse(stream, {
headers: { ...SSE_HEADERS, 'X-Table-Id': tableId },
return createEventStreamResponse({
requestId,
label: 'table',
streamId: tableId,
fromEventId,
getLatestEventId: getLatestTableEventId,
readEventsSince: readTableEventsSince,
extraHeaders: { 'X-Table-Id': tableId },
})
})
86 changes: 86 additions & 0 deletions apps/sim/lib/realtime/event-log.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,86 @@
/**
* @vitest-environment node
*/
import { beforeEach, describe, expect, it, vi } from 'vitest'

vi.mock('@/lib/core/config/env', () => ({ env: { REDIS_URL: undefined } }))
vi.mock('@/lib/core/config/redis', () => ({ getRedisClient: () => null }))

import {
appendEvent,
type EventLogConfig,
type EventLogEntry,
getLatestEventId,
readEventsSince,
resetEventLogMemoryForTesting,
} from '@/lib/realtime/event-log'

interface TestEntry extends EventLogEntry {
eventId: number
streamId: string
value: string
}

const config: EventLogConfig = { prefix: 'test:stream:', ttlSeconds: 3600, cap: 3, readChunk: 500 }

function serializerFor(streamId: string, value: string) {
return {
entryPrefix: '{"eventId":',
entrySuffix: `,"streamId":${JSON.stringify(streamId)},"value":${JSON.stringify(value)}}`,
buildMemory: (eventId: number): TestEntry => ({ eventId, streamId, value }),
}
}

describe('event-log (memory fallback)', () => {
beforeEach(() => resetEventLogMemoryForTesting())

it('assigns monotonically increasing event ids', async () => {
const first = await appendEvent(config, 's1', serializerFor('s1', 'a'))
const second = await appendEvent(config, 's1', serializerFor('s1', 'b'))
expect(first?.eventId).toBe(1)
expect(second?.eventId).toBe(2)
})

it('isolates streams by id', async () => {
await appendEvent(config, 's1', serializerFor('s1', 'a'))
const other = await appendEvent(config, 's2', serializerFor('s2', 'x'))
expect(other?.eventId).toBe(1)
expect(await getLatestEventId(config, 's1')).toBe(1)
expect(await getLatestEventId(config, 's2')).toBe(1)
})

it('reads only events after the cursor', async () => {
await appendEvent(config, 's1', serializerFor('s1', 'a'))
await appendEvent(config, 's1', serializerFor('s1', 'b'))
const result = await readEventsSince<TestEntry>(config, 's1', 1)
expect(result.status).toBe('ok')
if (result.status === 'ok') {
expect(result.events).toHaveLength(1)
expect(result.events[0].eventId).toBe(2)
expect(result.events[0].value).toBe('b')
}
})

it('tails from the latest id and returns nothing for a fresh cursor', async () => {
await appendEvent(config, 's1', serializerFor('s1', 'a'))
await appendEvent(config, 's1', serializerFor('s1', 'b'))
const latest = await getLatestEventId(config, 's1')
const result = await readEventsSince<TestEntry>(config, 's1', latest)
expect(result).toEqual({ status: 'ok', events: [] })
})

it('reports pruned when the cursor falls behind the cap-trimmed buffer', async () => {
// cap = 3; append 5, so the earliest retained id is 3.
for (const v of ['a', 'b', 'c', 'd', 'e']) {
await appendEvent(config, 's1', serializerFor('s1', v))
}
const result = await readEventsSince<TestEntry>(config, 's1', 1)
expect(result.status).toBe('pruned')
if (result.status === 'pruned') expect(result.earliestEventId).toBe(3)
})

it('reports pruned for a non-zero cursor against a never-seen stream', async () => {
const result = await readEventsSince<TestEntry>(config, 'missing', 5)
expect(result.status).toBe('pruned')
})
})
Loading