repositories / will
will
owned by admin
src/service/main.ts
Rawimport { 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<void>
}
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<T>(worker: WorkerClient, op: string, ...args: unknown[]): Promise<T> {
return (await worker.call(op, ...args)) as T
}
export async function startService(config: ServiceConfig): Promise<RunningService> {
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<SseClient>()
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<unknown>(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<TombstoneIn>(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<Omit<AnalysisRecord, 'digest' | 'createdAt'>>(req)
if (typeof body.capability !== 'string' || typeof body.model !== 'string') {
throw new HttpError(400, 'capability and model are required')
}
const rec = await callOps<AnalysisRecord>(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<Record<string, unknown>>(req)
sendJson(res, 200, await callOps(worker, 'recall', body))
})
internal.get('/v1/status', async (_req, res) => {
lastStatus = await callOps<ServiceStatus>(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<string>(worker, 'quickCheck') })
else if (op === 'archive_cycle')
sendJson(res, 200, await callOps<ArchiveCycleResult>(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<ServiceStatus>(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<typeof createRouter>,
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<ServiceStatus>(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<ArchiveCycleResult>(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<string>(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<void>((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<void>((res) => publicServer.close(() => res())),
new Promise<void>((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<void> {
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<void>((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())
}