repositories / will
will
owned by admin
src/agent/service-client.ts
Rawimport 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
}
},
}
}