|
| 1 | +import { Hono } from 'hono' |
| 2 | +import { streamSSE } from 'hono/streaming' |
| 3 | +import { sseAggregator } from '../services/sse-aggregator' |
| 4 | +import { SSESubscribeSchema } from '@opencode-manager/shared/schemas' |
| 5 | +import { logger } from '../utils/logger' |
| 6 | + |
| 7 | +export function createSSERoutes() { |
| 8 | + const app = new Hono() |
| 9 | + |
| 10 | + app.get('/stream', async (c) => { |
| 11 | + const directoriesParam = c.req.query('directories') |
| 12 | + const directories = directoriesParam ? directoriesParam.split(',').filter(Boolean) : [] |
| 13 | + |
| 14 | + return streamSSE(c, async (stream) => { |
| 15 | + const clientId = `client_${Date.now()}_${Math.random().toString(36).slice(2)}` |
| 16 | + |
| 17 | + const cleanup = sseAggregator.addClient( |
| 18 | + clientId, |
| 19 | + (event, data) => { |
| 20 | + stream.writeSSE({ event, data }) |
| 21 | + }, |
| 22 | + directories |
| 23 | + ) |
| 24 | + |
| 25 | + stream.onAbort(() => { |
| 26 | + cleanup() |
| 27 | + }) |
| 28 | + |
| 29 | + try { |
| 30 | + await stream.writeSSE({ |
| 31 | + event: 'connected', |
| 32 | + data: JSON.stringify({ clientId, directories, ...sseAggregator.getConnectionStatus() }) |
| 33 | + }) |
| 34 | + } catch (err) { |
| 35 | + logger.error(`Failed to send SSE connected event for ${clientId}:`, err) |
| 36 | + } |
| 37 | + |
| 38 | + await new Promise(() => {}) |
| 39 | + }) |
| 40 | + }) |
| 41 | + |
| 42 | + app.post('/subscribe', async (c) => { |
| 43 | + const body = await c.req.json() |
| 44 | + const result = SSESubscribeSchema.safeParse(body) |
| 45 | + if (!result.success) { |
| 46 | + return c.json({ success: false, error: 'Invalid request', details: result.error.issues }, 400) |
| 47 | + } |
| 48 | + const success = sseAggregator.addDirectories(result.data.clientId, result.data.directories) |
| 49 | + if (!success) { |
| 50 | + return c.json({ success: false, error: 'Client not found' }, 404) |
| 51 | + } |
| 52 | + return c.json({ success: true }) |
| 53 | + }) |
| 54 | + |
| 55 | + app.post('/unsubscribe', async (c) => { |
| 56 | + const body = await c.req.json() |
| 57 | + const result = SSESubscribeSchema.safeParse(body) |
| 58 | + if (!result.success) { |
| 59 | + return c.json({ success: false, error: 'Invalid request', details: result.error.issues }, 400) |
| 60 | + } |
| 61 | + const success = sseAggregator.removeDirectories(result.data.clientId, result.data.directories) |
| 62 | + if (!success) { |
| 63 | + return c.json({ success: false, error: 'Client not found' }, 404) |
| 64 | + } |
| 65 | + return c.json({ success: true }) |
| 66 | + }) |
| 67 | + |
| 68 | + app.get('/status', (c) => { |
| 69 | + return c.json({ |
| 70 | + ...sseAggregator.getConnectionStatus(), |
| 71 | + clients: sseAggregator.getClientCount(), |
| 72 | + directories: sseAggregator.getActiveDirectories(), |
| 73 | + activeSessions: sseAggregator.getActiveSessions() |
| 74 | + }) |
| 75 | + }) |
| 76 | + |
| 77 | + return app |
| 78 | +} |
0 commit comments