import { createServer, type IncomingMessage, type Server, type ServerResponse } from 'node:http' import { join } from 'node:path' import { pathToFileURL } from 'node:url' import type { ArchiveCycleResult, MaintenanceOp, ServiceStatus } from '../shared/api.ts' import { MAINTENANCE_OPS } from '../shared/api.ts' import { casReadStream } from '../shared/cas.ts' import type { AnalysisRecord, AuditEvent, ChatMessageIn, DeploymentEventMirror, TombstoneIn, } from '../shared/events.ts' import { createRouter, HttpError, readJsonBody, readRawBody, sendError, sendJson, sseEvent, sseInit, } from '../shared/http.ts' import { createLogger } from '../shared/log.ts' import { systemClock, utcDayKey } from '../shared/time.ts' import { loadServiceConfig, type ServiceConfig } from './config.ts' import { readJsonFile, serveFile } from './static.ts' import { startWorkerClient, type WorkerClient } from './worker-client.ts' const log = createLogger({ component: 'service' }) interface SseClient { res: ServerResponse lastEventId: number } export interface RunningService { config: ServiceConfig worker: WorkerClient publicServer: Server internalServer: Server close(): Promise } function parseAddr(addr: string): { host: string; port: number } { const idx = addr.lastIndexOf(':') if (idx === -1) throw new Error(`bad address ${addr}`) return { host: addr.slice(0, idx), port: Number.parseInt(addr.slice(idx + 1), 10) } } async function callOps(worker: WorkerClient, op: string, ...args: unknown[]): Promise { return (await worker.call(op, ...args)) as T } export async function startService(config: ServiceConfig): Promise { const worker = startWorkerClient(config.workerUrl, config.dataDir) await worker.ready // Deployment configuration wins over persisted values (spec: thresholds are env-owned). await callOps(worker, 'setStorageThresholds', config.storageWarnPct, config.storageStopPct).catch( () => undefined, ) const sseClients = new Set() let shuttingDown = false let lastStatus: ServiceStatus | undefined const broadcastAudit = async (seqs: number[]) => { if (seqs.length === 0 || sseClients.size === 0) return const minSeq = Math.min(...seqs) const rows = await callOps<{ seq: number; id: string }[]>( worker, 'eventsSince', minSeq - 1, 1000, ) for (const client of sseClients) { for (const row of rows) { if (row.seq > client.lastEventId) { client.lastEventId = row.seq sseEvent(client.res, row.seq, 'audit', row) } } } } worker.onNotify((kind, seqs) => { if (kind === 'audit') void broadcastAudit(seqs) else for (const client of sseClients) sseEvent(client.res, 'chat', 'chat', { seqs }) }) const internal = createRouter() internal.post('/v1/audit/ingest', async (req, res) => { const body = await readJsonBody<{ events: AuditEvent[] }>(req) sendJson(res, 200, await callOps(worker, 'ingestAudit', body.events)) }) internal.post('/v1/chat/ingest', async (req, res) => { const body = await readJsonBody<{ messages: ChatMessageIn[] }>(req) sendJson(res, 200, await callOps(worker, 'ingestChat', body.messages)) }) internal.get('/v1/chat/by-message-id/:id', async (_req, res, params) => { const id = Number.parseInt(params.id ?? '', 10) if (Number.isNaN(id)) throw new HttpError(400, 'bad message id') const row = await callOps(worker, 'chatByMessageId', id) if (!row) throw new HttpError(404, 'unknown message') sendJson(res, 200, row) }) internal.post('/v1/chat/hide', async (req, res) => { const body = await readJsonBody(req) if (typeof body.messageId !== 'number' || typeof body.requesterId !== 'string') { throw new HttpError(400, 'messageId and requesterId are required') } sendJson(res, 200, await callOps(worker, 'hideMessage', body)) }) internal.post('/v1/deploy/events', async (req, res) => { const body = await readJsonBody<{ events: DeploymentEventMirror[] }>(req) sendJson(res, 200, await callOps(worker, 'ingestDeploy', body.events)) }) internal.post('/v1/media', async (req, res) => { if (lastStatus?.storageState === 'stop') throw new HttpError(503, 'storage below stop threshold') const bytes = await readRawBody(req, config.mediaMaxBytes) const type = req.headers['x-will-media-type'] if (typeof type !== 'string') throw new HttpError(400, 'x-will-media-type header required') const meta = await callOps(worker, 'putMedia', bytes, { type, originalName: header(req, 'x-will-media-name'), telegramFileId: header(req, 'x-will-telegram-file-id'), }) sendJson(res, 201, meta) }) internal.get('/v1/media/:digest', async (_req, res, params) => { await streamMedia(res, worker, config, params.digest) }) internal.get('/v1/media/:digest/meta', async (_req, res, params) => { const meta = await callOps(worker, 'getMediaMeta', params.digest) if (!meta) throw new HttpError(404, 'unknown media') sendJson(res, 200, meta) }) internal.put('/v1/media/:digest/analysis', async (req, res, params) => { const body = await readJsonBody>(req) if (typeof body.capability !== 'string' || typeof body.model !== 'string') { throw new HttpError(400, 'capability and model are required') } const rec = await callOps(worker, 'putAnalysis', { digest: params.digest, capability: body.capability, model: body.model, schemaVersion: body.schemaVersion ?? 1, result: body.result ?? null, usage: body.usage, }) sendJson(res, 200, rec) }) internal.get('/v1/media/:digest/analysis', async (req, res, params) => { const query = new URL(req.url ?? '/', 'http://localhost').searchParams const key = { digest: params.digest, capability: query.get('capability') ?? '', model: query.get('model') ?? '', schemaVersion: Number.parseInt(query.get('schemaVersion') ?? '1', 10), } const found = await callOps(worker, 'getAnalysis', key) if (!found) throw new HttpError(404, 'analysis cache miss') sendJson(res, 200, found) }) internal.post('/v1/recall', async (req, res) => { const body = await readJsonBody>(req) sendJson(res, 200, await callOps(worker, 'recall', body)) }) internal.get('/v1/status', async (_req, res) => { lastStatus = await callOps(worker, 'status') sendJson(res, 200, lastStatus) }) internal.post('/v1/maintenance/:op', async (_req, res, params) => { const op = params.op as MaintenanceOp if (!(MAINTENANCE_OPS as readonly string[]).includes(op)) { throw new HttpError(400, `unknown maintenance op ${op}`) } if (op === 'quick_check') sendJson(res, 200, { result: await callOps(worker, 'quickCheck') }) else if (op === 'archive_cycle') sendJson(res, 200, await callOps(worker, 'archiveCycle')) else if (op === 'vacuum') { await callOps(worker, 'vacuum') sendJson(res, 200, { done: true }) } else { sendJson(res, 200, await callOps<{ removedTmpFiles: number }>(worker, 'reconcile')) } }) internal.get('/readyz', async (_req, res) => { sendJson(res, shuttingDown ? 503 : 200, { ready: !shuttingDown }) }) const publicRouter = createRouter() publicRouter.get('/healthz', async (_req, res) => { if (shuttingDown) return sendJson(res, 503, { ok: false }) try { await callOps(worker, 'lastSeq') sendJson(res, 200, { ok: true }) } catch { sendJson(res, 503, { ok: false }) } }) publicRouter.get('/readyz', async (_req, res) => { sendJson(res, shuttingDown ? 503 : 200, { ready: !shuttingDown }) }) publicRouter.get('/api/timeline', async (_req, res, _params, query) => { const before = query.get('before') const page = await callOps<{ rows: unknown[]; cold: boolean }>(worker, 'timeline', { before: before !== null ? Number.parseInt(before, 10) : undefined, limit: Number.parseInt(query.get('limit') ?? '50', 10), }) sendJson(res, 200, page) }) publicRouter.get('/api/stream', async (req, res, _params, query) => { sseInit(res) const headerId = req.headers['last-event-id'] let lastEventId = Number.parseInt( (typeof headerId === 'string' ? headerId : query.get('lastEventId')) ?? '0', 10, ) if (Number.isNaN(lastEventId) || lastEventId < 0) lastEventId = 0 const client: SseClient = { res, lastEventId } sseClients.add(client) res.on('close', () => sseClients.delete(client)) // Replay missed durable events, then live notifications take over. const missed = await callOps<{ seq: number; id: string }[]>( worker, 'eventsSince', lastEventId, 500, ) for (const row of missed) { if (row.seq > client.lastEventId) { client.lastEventId = row.seq sseEvent(res, row.seq, 'audit', row) } } }) publicRouter.get('/api/chat', async (_req, res, _params, query) => { sendJson( res, 200, await callOps(worker, 'chatPage', { before: query.get('before') ?? undefined, limit: Number.parseInt(query.get('limit') ?? '50', 10), }), ) }) publicRouter.get('/api/media/:digest', async (_req, res, params) => { await streamMedia(res, worker, config, params.digest) }) publicRouter.get('/api/deployments', async (_req, res, _params, query) => { sendJson( res, 200, await callOps(worker, 'recall', { kind: 'deploy', limit: Number.parseInt(query.get('limit') ?? '50', 10), }), ) }) publicRouter.get('/api/status', async (_req, res) => { lastStatus = await callOps(worker, 'status') sendJson(res, 200, lastStatus) }) publicRouter.get('/api/usage', async (_req, res) => { sendJson(res, 200, await callOps(worker, 'usage')) }) publicRouter.get('/api/docs', async (_req, res) => { const manifest = await readJsonFile(join(config.docsDist, 'manifest.json')) sendJson(res, 200, manifest) }) const dispatch = async ( serverName: string, router: ReturnType, req: IncomingMessage, res: ServerResponse, ) => { try { await router.dispatch(req, res) } catch (err) { if (err instanceof HttpError) sendError(res, err.status, err.message) else { log.error('request failed', { server: serverName, path: req.url, error: err instanceof Error ? err.message : String(err), }) if (!res.headersSent) sendError(res, 500, 'internal error') else res.end() } } } const publicServer = createServer((req, res) => { const url = new URL(req.url ?? '/', 'http://localhost') if (url.pathname === '/' || url.pathname === '/index.html') { void serveFile(res, config.webDist, 'index.html').catch(() => sendError(res, 404, 'no frontend build'), ) return } if (url.pathname.startsWith('/assets/')) { void serveFile(res, config.webDist, url.pathname, { immutable: true }).catch( (err: unknown) => { if (err instanceof HttpError) sendError(res, err.status, err.message) else sendError(res, 500, 'static error') }, ) return } if (url.pathname === '/docs' || url.pathname.startsWith('/docs/')) { const rel = url.pathname.replace(/^\/docs\/?/, '') void serveFile(res, config.docsDist, rel === '' ? 'index.html' : `${rel}.html`).catch( (err: unknown) => { if (err instanceof HttpError) sendError(res, err.status, err.message) else sendError(res, 500, 'static error') }, ) return } void dispatch('public', publicRouter, req, res) }) const internalServer = createServer((req, res) => void dispatch('internal', internal, req, res)) const heartbeat = setInterval(() => { for (const client of sseClients) client.res.write(':hb\n\n') }, 25_000) heartbeat.unref() // Maintenance scheduler: hourly archive+vacuum+reconcile, daily quick_check (spec AC12/AC13). let lastArchiveHour = '' let lastCheckDay = '' const maintenance = setInterval(() => { void (async () => { try { lastStatus = await callOps(worker, 'status') const now = systemClock.now() const hourKey = now.toISOString().slice(0, 13) const dayKey = utcDayKey(now) if (lastArchiveHour !== hourKey) { lastArchiveHour = hourKey const cycle = await callOps(worker, 'archiveCycle') await callOps(worker, 'vacuum') await callOps(worker, 'reconcile') log.info('maintenance cycle', { exportedAudit: cycle.exportedAudit, exportedChat: cycle.exportedChat, storage: lastStatus.storageState, }) } if (lastCheckDay !== dayKey) { lastCheckDay = dayKey const check = await callOps(worker, 'quickCheck') if (check !== 'ok') log.error('sqlite quick_check failed', { result: check }) } } catch (err) { log.error('maintenance failed', { error: err instanceof Error ? err.message : String(err) }) } })() }, 60_000) maintenance.unref() await new Promise((resolve, reject) => { publicServer.once('error', reject) internalServer.once('error', reject) const pub = parseAddr(config.publicAddr) const int = parseAddr(config.internalAddr) publicServer.listen(pub.port, pub.host, () => { internalServer.listen(int.port, int.host, () => resolve()) }) }) log.info('service listening', { public: config.publicAddr, internal: config.internalAddr }) return { config, worker, publicServer, internalServer, async close() { shuttingDown = true clearInterval(maintenance) clearInterval(heartbeat) for (const client of sseClients) client.res.end() await Promise.all([ new Promise((res) => publicServer.close(() => res())), new Promise((res) => internalServer.close(() => res())), ]) await worker.close(5000) }, } } function header(req: IncomingMessage, name: string): string | undefined { const v = req.headers[name] return typeof v === 'string' ? v : undefined } async function streamMedia( res: ServerResponse, worker: WorkerClient, config: ServiceConfig, digest: string | undefined, ): Promise { if (!digest || !/^[0-9a-f]{64}$/.test(digest)) throw new HttpError(400, 'bad digest') const meta = await callOps<{ type: string; size: number } | undefined>( worker, 'getMediaMeta', digest, ) if (!meta) throw new HttpError(404, 'unknown media') res.writeHead(200, { 'content-type': meta.type, 'content-length': meta.size, 'cache-control': 'public, max-age=31536000, immutable', }) await new Promise((done) => { const stream = casReadStream(join(config.dataDir, 'media'), digest) stream.on('error', () => { res.end() done() }) stream.on('end', () => done()) stream.pipe(res) }) } const isEntryPoint = process.argv[1] !== undefined && import.meta.url === pathToFileURL(process.argv[1]).href if (isEntryPoint) { const distRoot = new URL('.', import.meta.url).pathname const config = loadServiceConfig({ distRoot }) const service = await startService(config) const shutdown = async () => { log.info('shutdown starting', {}) const hard = setTimeout(() => process.exit(1), 9500) hard.unref() await service.close() clearTimeout(hard) process.exit(0) } process.on('SIGTERM', () => void shutdown()) process.on('SIGINT', () => void shutdown()) }