Luigit
repositories / will

will

owned by admin

src/service/main.ts

Raw
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<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())
}