#!/usr/bin/env node import fs from 'node:fs'; import http from 'node:http'; import https from 'node:https'; import type { IncomingMessage, ServerResponse } from 'node:http'; import type { Socket } from 'node:net'; import { pathToFileURL } from 'node:url'; import { authenticate } from './auth.ts'; import { loadConfig, readSecret, validateServerConfig } from './config.ts'; import { MailboxStore } from './storage.ts'; import { acceptUpgrade } from './websocket.ts'; import { WebSocketPeer } from './websocket.ts'; import type { ServerConfig } from './types.ts'; const errors = { invalid_request: [400, 'Invalid request'], unauthenticated: [401, 'Authentication required'], not_found: [404, 'Not found'], conflict: [409, 'Message ID conflict'] } as const; type ErrorCode = keyof typeof errors; function respond(res: ServerResponse, status: number, body: unknown): void { res.writeHead(status, { 'content-type': 'application/json; charset=utf-8', 'cache-control': 'no-store' }); res.end(JSON.stringify(body)); } function fail(res: ServerResponse, code: ErrorCode): void { const [status, message] = errors[code]; respond(res, status, { error: { code, message } }); } function uuid(value: unknown): value is string { return typeof value === 'string' && /^[0-9a-f]{8}-[0-9a-f]{4}-[1-8][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i.test(value); } async function jsonBody(req: IncomingMessage): Promise { if (!req.headers['content-type']?.toLowerCase().startsWith('application/json')) throw new Error('content-type'); let size = 0, chunks = []; for await (const chunk of req) { size += chunk.length; // JSON escapes can be much larger than the decoded 32 KiB UTF-8 body. if (size > 256 * 1024) throw new Error('too large'); chunks.push(chunk); } return JSON.parse(Buffer.concat(chunks).toString('utf8')); } export function createMailboxServer(config: ServerConfig) { const registry = validateServerConfig(config); const store = new MailboxStore(config.database_path, registry.keys.values()); const subscribers = new Set(); const server = registry.listener.tls ? https.createServer({ cert: fs.readFileSync(registry.listener.tls.cert_file), key: readSecret(registry.listener.tls.key) }) : http.createServer(); server.on('request', async (req: IncomingMessage, res: ServerResponse) => { const auth = authenticate(req.headers.authorization, registry); if (!auth) return fail(res, 'unauthenticated'); let url; try { url = new URL(req.url ?? '/', 'http://localhost'); } catch { return fail(res, 'invalid_request'); } const path = url.pathname; const messageRoute = /^\/v1\/messages\/([^/]+)$/.exec(path); try { if (req.method === 'POST' && path === '/v1/messages') { const body = await jsonBody(req); if (!body || typeof body !== 'object' || Array.isArray(body)) return fail(res, 'invalid_request'); const input = body as Record; if (Object.keys(input).sort().join(',') !== 'body,id' || !uuid(input.id) || typeof input.body !== 'string' || Buffer.byteLength(input.body, 'utf8') < 1 || Buffer.byteLength(input.body, 'utf8') > 32768) return fail(res, 'invalid_request'); const result = store.publish({ id: input.id, body: input.body }, auth.principal.id, auth.kid); if (!result) return fail(res, 'conflict'); respond(res, result.created ? 201 : 200, result.message); if (result.created) for (const peer of subscribers) peer.sendJSON({ type: 'message.available' }); return; } if (req.method === 'GET' && messageRoute) { const message = store.get(messageRoute[1]); return message ? respond(res, 200, message) : fail(res, 'not_found'); } if (req.method === 'GET' && path === '/v1/messages') { const rawLimit = url.searchParams.get('limit'); if ([...url.searchParams.keys()].some(key => key !== 'limit') || url.searchParams.getAll('limit').length > 1 || (rawLimit !== null && !/^(?:[1-9]|[1-9][0-9]|100)$/.test(rawLimit))) return fail(res, 'invalid_request'); return respond(res, 200, { messages: store.list(auth.kid, rawLimit === null ? 100 : Number(rawLimit)) }); } if (req.method === 'GET' && path === '/v1/cursor') return respond(res, 200, { sequence: store.cursor(auth.kid) }); if (req.method === 'POST' && path === '/v1/cursor') { const body = await jsonBody(req); if (!body || typeof body !== 'object' || Array.isArray(body)) return fail(res, 'invalid_request'); const input = body as Record; if (Object.keys(input).join(',') !== 'sequence' || typeof input.sequence !== 'number' || !Number.isSafeInteger(input.sequence) || input.sequence < 0) return fail(res, 'invalid_request'); const result = store.advance(auth.kid, input.sequence); return result ? respond(res, 200, result) : fail(res, 'invalid_request'); } return fail(res, 'not_found'); } catch (error) { if (error instanceof SyntaxError || (error instanceof Error && ['content-type', 'too large'].includes(error.message))) return fail(res, 'invalid_request'); console.error('request failed:', error); if (!res.headersSent) respond(res, 500, { error: { code: 'internal_error', message: 'Internal error' } }); } }); server.on('upgrade', (req: IncomingMessage, socket: Socket, head: Buffer) => { const auth = authenticate(req.headers.authorization, registry); if (!auth) { socket.end('HTTP/1.1 401 Unauthorized\r\nCache-Control: no-store\r\nConnection: close\r\n\r\n'); return; } if (req.url !== '/v1/events') { socket.end('HTTP/1.1 404 Not Found\r\nCache-Control: no-store\r\nConnection: close\r\n\r\n'); return; } const peer = acceptUpgrade(req, socket, head); if (!peer) { socket.end('HTTP/1.1 400 Bad Request\r\nCache-Control: no-store\r\nConnection: close\r\n\r\n'); return; } subscribers.add(peer); const expiry = setTimeout(() => peer.close(4001, 'token_expired'), Math.max(0, auth.expiresAt * 1000 - Date.now())); let pongDeadline: NodeJS.Timeout | undefined; const pinger = setInterval(() => { peer.ping(); pongDeadline = setTimeout(() => peer.close(1001, 'pong_timeout'), 15000); }, 30000); peer.on('pong', () => clearTimeout(pongDeadline)); peer.on('close', () => { clearTimeout(expiry); clearInterval(pinger); clearTimeout(pongDeadline); subscribers.delete(peer); }); peer.on('error', () => socket.destroy()); peer.sendJSON({ type: 'ready' }); }); return { server, store, listen: () => new Promise(resolve => server.listen(registry.listener.port, registry.listener.host, () => resolve())), close: async () => { for (const peer of subscribers) peer.close(1001); await new Promise(resolve => server.close(() => resolve())); store.close(); } }; } if (process.argv[1] && import.meta.url === pathToFileURL(process.argv[1]).href) { if (!process.argv[2]) throw new Error('usage: orpheus-server CONFIG.json'); const config = loadConfig(process.argv[2]); const app = createMailboxServer(config); app.listen().then(() => console.log(`orpheus listening on ${config.listen.host}:${config.listen.port}`)); process.on('SIGTERM', () => app.close().then(() => process.exit(0))); process.on('SIGINT', () => app.close().then(() => process.exit(0))); }