Luigit
repositories / will

will

owned by admin

src/agent/service-client.ts

Raw
import type { RecallQuery, RecallResponse, ServiceStatus } from '../shared/api.ts'
import type {
  AnalysisRecord,
  AuditEvent,
  ChatMessageIn,
  DeploymentEventMirror,
  MediaMeta,
  TombstoneIn,
} from '../shared/events.ts'

/**
 * Typed client for the service's pod-local API.
 * All agent persistence (audit, chat, media, analyses, deployment mirror)
 * crosses this boundary (spec: service is the only SQLite writer).
 */

export interface ServiceClient {
  ingestAudit(events: AuditEvent[]): Promise<{ accepted: number; duplicates: number }>
  ingestChat(messages: ChatMessageIn[]): Promise<{ accepted: number; duplicates: number }>
  hideMessage(tombstone: TombstoneIn): Promise<{ ok: boolean; alreadyHidden: boolean }>
  ingestDeploy(events: DeploymentEventMirror[]): Promise<{ accepted: number; duplicates: number }>
  putMedia(
    bytes: Uint8Array,
    meta: { type: string; originalName?: string; telegramFileId?: string },
  ): Promise<MediaMeta>
  getMedia(digest: string): Promise<Uint8Array>
  putAnalysis(
    rec: Omit<AnalysisRecord, 'digest' | 'createdAt'> & { digest: string },
  ): Promise<AnalysisRecord>
  getAnalysis(key: {
    digest: string
    capability: string
    model: string
    schemaVersion: number
  }): Promise<AnalysisRecord | undefined>
  recall(query: RecallQuery): Promise<RecallResponse>
  status(): Promise<ServiceStatus>
  chatByMessageId(messageId: number): Promise<ChatMessageIn | undefined>
  maintenance(op: 'quick_check' | 'archive_cycle' | 'vacuum' | 'reconcile'): Promise<unknown>
}

export function createServiceClient(baseUrl: string): ServiceClient {
  const json = async <T>(path: string, init?: RequestInit): Promise<T> => {
    const res = await fetch(`${baseUrl}${path}`, {
      ...init,
      headers: { 'content-type': 'application/json', ...(init?.headers ?? {}) },
    })
    if (!res.ok) throw new Error(`service ${path} failed: ${res.status} ${await res.text()}`)
    return (await res.json()) as T
  }
  return {
    ingestAudit: (events) =>
      json('/v1/audit/ingest', { method: 'POST', body: JSON.stringify({ events }) }),
    ingestChat: (messages) =>
      json('/v1/chat/ingest', { method: 'POST', body: JSON.stringify({ messages }) }),
    hideMessage: (tombstone) =>
      json('/v1/chat/hide', { method: 'POST', body: JSON.stringify(tombstone) }),
    ingestDeploy: (events) =>
      json('/v1/deploy/events', { method: 'POST', body: JSON.stringify({ events }) }),
    async putMedia(bytes, meta) {
      const res = await fetch(`${baseUrl}/v1/media`, {
        method: 'POST',
        headers: {
          'content-type': 'application/octet-stream',
          'x-will-media-type': meta.type,
          ...(meta.originalName ? { 'x-will-media-name': meta.originalName } : {}),
          ...(meta.telegramFileId ? { 'x-will-telegram-file-id': meta.telegramFileId } : {}),
        },
        body: bytes,
      })
      if (!res.ok) throw new Error(`media put failed: ${res.status}`)
      return (await res.json()) as MediaMeta
    },
    async getMedia(digest) {
      const res = await fetch(`${baseUrl}/v1/media/${digest}`)
      if (!res.ok) throw new Error(`media get failed: ${res.status}`)
      return new Uint8Array(await res.arrayBuffer())
    },
    putAnalysis: (rec) =>
      json(`/v1/media/${rec.digest}/analysis`, {
        method: 'PUT',
        body: JSON.stringify({
          capability: rec.capability,
          model: rec.model,
          schemaVersion: rec.schemaVersion,
          result: rec.result,
          usage: rec.usage,
        }),
      }),
    async getAnalysis(key) {
      try {
        return await json<AnalysisRecord>(
          `/v1/media/${key.digest}/analysis?capability=${encodeURIComponent(key.capability)}&model=${encodeURIComponent(key.model)}&schemaVersion=${key.schemaVersion}`,
          { method: 'GET' },
        )
      } catch {
        return undefined
      }
    },
    recall: (query) => json('/v1/recall', { method: 'POST', body: JSON.stringify(query) }),
    status: () => json('/v1/status', { method: 'GET' }),
    maintenance: (op) => json(`/v1/maintenance/${op}`, { method: 'POST', body: '{}' }),
    async chatByMessageId(messageId) {
      try {
        return await json<ChatMessageIn>(`/v1/chat/by-message-id/${messageId}`, { method: 'GET' })
      } catch {
        return undefined
      }
    },
  }
}