repositories / will
will
owned by admin
src/service/ops.ts
Rawimport { randomBytes } from 'node:crypto'
import { stat, statfs } from 'node:fs/promises'
import { join } from 'node:path'
import type { DatabaseSync } from 'node:sqlite'
import type {
ArchiveCycleResult,
ChatViewRow,
IngestResult,
RecallQuery,
RecallResponse,
ServiceStatus,
TimelineRow,
} from '../shared/api.ts'
import { listSidecars, readSegment, writeSegment } from '../shared/archive.ts'
import { casReconcile, casStore } from '../shared/cas.ts'
import { schemaVersion } from '../shared/db.ts'
import type {
AnalysisRecord,
AuditEvent,
ChatMessageIn,
DeploymentEventMirror,
MediaMeta,
TombstoneIn,
} from '../shared/events.ts'
import { safeParse } from '../shared/json.ts'
import {
type ArchivedRecord,
cutoffIso,
DEFAULT_RETENTION,
groupByMonth,
rowsToArchive,
segmentsOverlapping,
} from '../shared/retention.ts'
import { type Clock, nowIso } from '../shared/time.ts'
import { aggregateTurnUsage, type UsageReport } from './usage.ts'
export interface OpsDirs {
dataDir: string
archiveDir: string
mediaDir: string
}
const BATCH_LIMIT = 5000
interface AuditRowDb {
seq: number
id: string
ts: string
session_id: string | null
revision: string | null
agent_digest: string | null
worktree_head: string | null
worktree_dirty: number
kind: string
payload: string
}
interface ChatRowDb {
id: string
update_id: number | null
chat_id: number
message_id: number
direction: 'in' | 'out'
ts: string
sender_id: string | null
sender_name: string | null
reply_to_message_id: number | null
trigger_reason: string | null
text: string | null
media_json: string | null
}
function auditDbToRow(r: AuditRowDb, archived = false): TimelineRow {
const provenance: NonNullable<TimelineRow['provenance']> = {
worktreeDirty: r.worktree_dirty === 1,
}
if (r.revision !== null) provenance.revision = r.revision
if (r.agent_digest !== null) provenance.agentDigest = r.agent_digest
if (r.worktree_head !== null) provenance.worktreeHead = r.worktree_head
return {
seq: r.seq,
id: r.id,
ts: r.ts,
kind: r.kind,
...(r.session_id !== null ? { sessionId: r.session_id } : {}),
provenance,
payload: safeParse<unknown>(r.payload) ?? r.payload,
archived,
}
}
function chatDbToRow(r: ChatRowDb, hidden: boolean, archived = false): ChatViewRow {
return {
id: r.id,
chatId: r.chat_id,
messageId: r.message_id,
direction: r.direction,
ts: r.ts,
...(r.sender_id !== null ? { senderId: r.sender_id } : {}),
...(r.sender_name !== null ? { senderName: r.sender_name } : {}),
...(r.reply_to_message_id !== null ? { replyToMessageId: r.reply_to_message_id } : {}),
...(r.trigger_reason !== null
? { triggerReason: r.trigger_reason as ChatViewRow['triggerReason'] }
: {}),
...(r.text !== null ? { text: r.text } : {}),
media: safeParse<ChatViewRow['media']>(r.media_json ?? 'null'),
hidden,
archived,
}
}
/**
* All SQLite and file operations of the service, framework-free so tests
* can drive it directly on tmpdirs.
*/
export function createOps(db: DatabaseSync, dirs: OpsDirs, clock: Clock) {
const insertAudit = db.prepare(`
INSERT OR IGNORE INTO audit_events
(id, ts, session_id, revision, agent_digest, worktree_head, worktree_dirty, kind, payload)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
`)
const insertChat = db.prepare(`
INSERT OR IGNORE INTO chat_messages
(id, update_id, chat_id, message_id, direction, ts, sender_id, sender_name,
reply_to_message_id, trigger_reason, text, media_json, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`)
const insertDeploy = db.prepare(`
INSERT OR IGNORE INTO deployment_events (id, ts, kind, revision, detail, received_at)
VALUES (?, ?, ?, ?, ?, ?)
`)
const recordJob = <T>(kind: string, run: () => T | Promise<T>): Promise<T> => {
const info = db
.prepare('INSERT INTO jobs (kind, started_at, status) VALUES (?, ?, ?)')
.run(kind, nowIso(clock), 'running')
const finish = (status: string, detail: unknown) =>
db
.prepare('UPDATE jobs SET finished_at = ?, status = ?, detail = ? WHERE id = ?')
.run(
nowIso(clock),
status,
detail === undefined ? null : JSON.stringify(detail),
Number(info.lastInsertRowid),
)
return Promise.resolve()
.then(run)
.then(
(result) => {
finish('done', result)
return result as T
},
(err) => {
finish('failed', err instanceof Error ? err.message : String(err))
throw err
},
)
}
const hiddenMessageIds = (): Set<number> =>
new Set(
(db.prepare('SELECT message_id FROM tombstones').all() as { message_id: number }[]).map(
(r) => r.message_id,
),
)
const loadAuditRows = (ids: string[]): ArchivedRecord[] => {
const out: ArchivedRecord[] = []
const sel = db.prepare('SELECT * FROM audit_events WHERE id = ?')
for (const id of ids) {
const r = sel.get(id) as AuditRowDb | undefined
if (r)
out.push({
id: r.id,
ts: r.ts,
record: {
kind: r.kind,
sessionId: r.session_id ?? undefined,
provenance: {
revision: r.revision ?? undefined,
agentDigest: r.agent_digest ?? undefined,
worktreeHead: r.worktree_head ?? undefined,
worktreeDirty: r.worktree_dirty === 1,
},
payload: safeParse<unknown>(r.payload) ?? r.payload,
},
})
}
return out.sort((a, b) => a.ts.localeCompare(b.ts))
}
const loadChatRows = (ids: string[]): ArchivedRecord[] => {
const out: ArchivedRecord[] = []
const sel = db.prepare('SELECT * FROM chat_messages WHERE id = ?')
for (const id of ids) {
const r = sel.get(id) as ChatRowDb | undefined
if (r)
out.push({
id: r.id,
ts: r.ts,
record: {
id: r.id,
chatId: r.chat_id,
messageId: r.message_id,
direction: r.direction,
ts: r.ts,
senderId: r.sender_id ?? undefined,
senderName: r.sender_name ?? undefined,
replyToMessageId: r.reply_to_message_id ?? undefined,
triggerReason: (r.trigger_reason as ChatViewRow['triggerReason']) ?? undefined,
text: r.text ?? undefined,
media: safeParse<ChatViewRow['media']>(r.media_json ?? 'null') ?? undefined,
hidden: false,
},
})
}
return out.sort((a, b) => a.ts.localeCompare(b.ts))
}
/**
* Archive one table's stale rows into immutable monthly segments, then delete
* the hot copies. Exactly-once: rows already present in an overlapping segment
* (a previous cycle crashed after publishing, before deleting) are skipped and
* only deleted.
*/
const runKind = async (
kind: 'audit' | 'chat',
hotDays: number,
loadRows: (ids: string[]) => ArchivedRecord[],
): Promise<{ exported: number; segments: number }> => {
const cutoff = cutoffIso(clock.now(), hotDays)
const table = kind === 'audit' ? 'audit_events' : 'chat_messages'
const stale = db
.prepare(`SELECT id, ts FROM ${table} WHERE ts < ? ORDER BY ts LIMIT ${BATCH_LIMIT}`)
.all(cutoff) as { id: string; ts: string }[]
const toArchive = rowsToArchive(stale, cutoff)
if (toArchive.length === 0) return { exported: 0, segments: 0 }
const segments = (await listSidecars(dirs.archiveDir, kind)).map((s) => ({
...s.sidecar,
zstPath: s.zstPath,
}))
const oldestWanted = toArchive[toArchive.length - 1]?.ts ?? ''
const overlapping = segmentsOverlapping(segments, { until: oldestWanted })
const alreadyExported = new Set<string>()
for (const seg of overlapping) {
for (const r of await readSegment(seg.zstPath, () => true, BATCH_LIMIT)) {
alreadyExported.add(r.id)
}
}
let exported = 0
let segmentCount = 0
for (const [month, rows] of groupByMonth(toArchive.filter((r) => !alreadyExported.has(r.id)))) {
const records = loadRows(rows.map((r) => r.id))
if (records.length === 0) continue
await writeSegment(
dirs.archiveDir,
kind,
month,
records,
`${Date.now()}-${randomBytes(4).toString('hex')}`,
)
segmentCount++
exported += records.length
}
db.exec('BEGIN')
try {
const del = db.prepare(`DELETE FROM ${table} WHERE id = ?`)
for (const r of toArchive) del.run(r.id)
db.exec('COMMIT')
} catch (err) {
db.exec('ROLLBACK')
throw err
}
return { exported, segments: segmentCount }
}
const coldFill = async <T>(
kind: 'audit' | 'chat',
range: { since?: string; until?: string },
seen: Set<string>,
budget: number,
mapRecord: (r: ArchivedRecord) => T,
matches: (r: ArchivedRecord) => boolean = () => true,
): Promise<T[]> => {
const segments = (await listSidecars(dirs.archiveDir, kind)).map((s) => ({
...s.sidecar,
zstPath: s.zstPath,
}))
const overlapping = segmentsOverlapping(segments, range).sort((a, b) =>
b.firstTs.localeCompare(a.firstTs),
)
const out: T[] = []
outer: for (const seg of overlapping) {
const records = await readSegment(seg.zstPath, matches, BATCH_LIMIT)
const mapped = records
.map(mapRecord)
.sort((a, b) => (b as { ts: string }).ts.localeCompare((a as { ts: string }).ts))
for (const row of mapped) {
const id = (row as { id: string }).id
if (seen.has(id)) continue
seen.add(id)
out.push(row)
if (out.length >= budget) break outer
}
}
return out
}
const coldSearch = async <T extends { id: string; ts: string }>(
kind: 'audit' | 'chat',
range: { since?: string; until?: string },
seen: Set<string>,
mapRecord: (r: ArchivedRecord) => T,
matches: (r: ArchivedRecord) => boolean,
): Promise<T[]> => {
const segments = (await listSidecars(dirs.archiveDir, kind)).map((s) => ({
...s.sidecar,
zstPath: s.zstPath,
}))
const out: T[] = []
for (const seg of segmentsOverlapping(segments, range)) {
for (const record of await readSegment(seg.zstPath, matches, BATCH_LIMIT)) {
if (seen.has(record.id)) continue
seen.add(record.id)
out.push(mapRecord(record))
}
}
return out.sort((a, b) => b.ts.localeCompare(a.ts) || b.id.localeCompare(a.id))
}
const newestFirst = (a: { id: string; ts: string }, b: { id: string; ts: string }) =>
b.ts.localeCompare(a.ts) || Buffer.compare(Buffer.from(b.id), Buffer.from(a.id))
db.function('unicode_lower', { deterministic: true }, (value) =>
String(value ?? '').toLowerCase(),
)
const ops = {
ingestAudit(events: AuditEvent[]): IngestResult {
let accepted = 0
let duplicates = 0
const insertedSeqs: number[] = []
db.exec('BEGIN')
try {
for (const e of events) {
const p = e.provenance ?? {}
const res = insertAudit.run(
e.id,
e.ts,
e.sessionId ?? null,
p.revision ?? null,
p.agentDigest ?? null,
p.worktreeHead ?? null,
p.worktreeDirty === true ? 1 : 0,
e.kind,
JSON.stringify(e.payload ?? null),
)
if (res.changes === 1) {
accepted++
insertedSeqs.push(Number(res.lastInsertRowid))
} else duplicates++
}
db.exec('COMMIT')
} catch (err) {
db.exec('ROLLBACK')
throw err
}
return { accepted, duplicates, insertedSeqs }
},
ingestChat(msgs: ChatMessageIn[]): IngestResult {
let accepted = 0
let duplicates = 0
db.exec('BEGIN')
try {
for (const m of msgs) {
const res = insertChat.run(
m.id,
m.updateId ?? null,
m.chatId,
m.messageId,
m.direction,
m.ts,
m.senderId ?? null,
m.senderName ?? null,
m.replyToMessageId ?? null,
m.triggerReason ?? null,
m.text ?? null,
m.media ? JSON.stringify(m.media) : null,
nowIso(clock),
)
if (res.changes === 1) accepted++
else duplicates++
}
db.exec('COMMIT')
} catch (err) {
db.exec('ROLLBACK')
throw err
}
return { accepted, duplicates, insertedSeqs: [] }
},
chatByMessageId(messageId: number): { senderId?: string } | undefined {
const row = db
.prepare('SELECT sender_id FROM chat_messages WHERE message_id = ?')
.get(messageId) as { sender_id: string | null } | undefined
if (!row || row.sender_id === null) return undefined
return { senderId: row.sender_id }
},
hideMessage(input: TombstoneIn): { ok: boolean; alreadyHidden: boolean } {
const res = db
.prepare(
'INSERT OR IGNORE INTO tombstones (message_id, requester_id, authority, reason, ts) VALUES (?, ?, ?, ?, ?)',
)
.run(input.messageId, input.requesterId, input.authority, input.reason ?? null, input.ts)
return { ok: true, alreadyHidden: res.changes === 0 }
},
ingestDeploy(events: DeploymentEventMirror[]): IngestResult {
let accepted = 0
let duplicates = 0
db.exec('BEGIN')
try {
for (const e of events) {
const res = insertDeploy.run(
e.id,
e.ts,
e.kind,
e.revision ?? null,
JSON.stringify(e.detail),
nowIso(clock),
)
if (res.changes === 1) accepted++
else duplicates++
}
db.exec('COMMIT')
} catch (err) {
db.exec('ROLLBACK')
throw err
}
return { accepted, duplicates, insertedSeqs: [] }
},
async putMedia(
input: Uint8Array,
meta: { type: string; originalName?: string; telegramFileId?: string },
): Promise<MediaMeta> {
const bytes = Buffer.from(input.buffer, input.byteOffset, input.byteLength)
const stored = await casStore(dirs.mediaDir, bytes)
db.prepare(
`INSERT OR IGNORE INTO media (digest, size, type, original_name, telegram_file_id, created_at)
VALUES (?, ?, ?, ?, ?, ?)`,
).run(
stored.digest,
stored.size,
meta.type,
meta.originalName ?? null,
meta.telegramFileId ?? null,
nowIso(clock),
)
const found = ops.getMediaMeta(stored.digest)
if (!found) throw new Error('media row missing after insert')
return found
},
getMediaMeta(digest: string): MediaMeta | undefined {
const row = db.prepare('SELECT * FROM media WHERE digest = ?').get(digest) as
| {
digest: string
size: number
type: string
original_name: string | null
telegram_file_id: string | null
created_at: string
}
| undefined
if (!row) return undefined
const meta: MediaMeta = {
digest: row.digest,
size: row.size,
type: row.type,
createdAt: row.created_at,
}
if (row.original_name !== null) meta.originalName = row.original_name
if (row.telegram_file_id !== null) meta.telegramFileId = row.telegram_file_id
return meta
},
putAnalysis(rec: Omit<AnalysisRecord, 'createdAt'> & { createdAt?: string }): AnalysisRecord {
const createdAt = rec.createdAt ?? nowIso(clock)
db.prepare(
`INSERT OR REPLACE INTO media_analyses
(digest, capability, model, schema_version, result, usage_json, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?)`,
).run(
rec.digest,
rec.capability,
rec.model,
rec.schemaVersion,
JSON.stringify(rec.result ?? null),
rec.usage === undefined ? null : JSON.stringify(rec.usage),
createdAt,
)
return { ...rec, createdAt }
},
getAnalysis(key: {
digest: string
capability: string
model: string
schemaVersion: number
}): AnalysisRecord | undefined {
const row = db
.prepare(
'SELECT * FROM media_analyses WHERE digest = ? AND capability = ? AND model = ? AND schema_version = ?',
)
.get(key.digest, key.capability, key.model, key.schemaVersion) as
| {
digest: string
capability: string
model: string
schema_version: number
result: string
usage_json: string | null
created_at: string
}
| undefined
if (!row) return undefined
return {
digest: row.digest,
capability: row.capability,
model: row.model,
schemaVersion: row.schema_version,
result: safeParse<unknown>(row.result) ?? row.result,
...(row.usage_json !== null
? { usage: safeParse<unknown>(row.usage_json) ?? undefined }
: {}),
createdAt: row.created_at,
}
},
lastSeq(): number {
const row = db.prepare('SELECT COALESCE(MAX(seq), 0) AS m FROM audit_events').get() as {
m: number
}
return row.m
},
eventsSince(seq: number, limit: number): TimelineRow[] {
const rows = db
.prepare('SELECT * FROM audit_events WHERE seq > ? ORDER BY seq LIMIT ?')
.all(seq, limit) as unknown as AuditRowDb[]
return rows.map((r) => auditDbToRow(r))
},
async timeline(query: {
before?: number
limit: number
}): Promise<{ rows: TimelineRow[]; cold: boolean }> {
const limit = Math.min(Math.max(query.limit, 1), 500)
const hot = (
query.before !== undefined
? (db
.prepare('SELECT * FROM audit_events WHERE seq < ? ORDER BY seq DESC LIMIT ?')
.all(query.before, limit) as unknown as AuditRowDb[])
: (db
.prepare('SELECT * FROM audit_events ORDER BY seq DESC LIMIT ?')
.all(limit) as unknown as AuditRowDb[])
).map((r) => auditDbToRow(r))
if (hot.length >= limit) return { rows: hot, cold: false }
const until = hot.length > 0 ? hot[hot.length - 1]?.ts : undefined
const cold = await coldFill<TimelineRow>(
'audit',
until === undefined ? {} : { until },
new Set(hot.map((r) => r.id)),
limit - hot.length,
(r) => {
const rec = r.record as {
kind: string
sessionId?: string
provenance?: TimelineRow['provenance']
payload?: unknown
}
return {
id: r.id,
ts: r.ts,
kind: rec.kind,
sessionId: rec.sessionId,
provenance: rec.provenance,
payload: rec.payload,
archived: true,
} satisfies TimelineRow
},
)
const rows = [...hot, ...cold].slice(0, limit)
return { rows, cold: cold.length > 0 }
},
async chatPage(query: {
before?: string
limit: number
}): Promise<{ rows: ChatViewRow[]; cold: boolean }> {
const limit = Math.min(Math.max(query.limit, 1), 500)
const hidden = hiddenMessageIds()
const hot = (
query.before !== undefined
? (db
.prepare('SELECT * FROM chat_messages WHERE ts < ? ORDER BY ts DESC LIMIT ?')
.all(query.before, limit) as unknown as ChatRowDb[])
: (db
.prepare('SELECT * FROM chat_messages ORDER BY ts DESC LIMIT ?')
.all(limit) as unknown as ChatRowDb[])
).map((r) => chatDbToRow(r, hidden.has(r.message_id)))
if (hot.length >= limit) return { rows: hot, cold: false }
const until = hot.length > 0 ? hot[hot.length - 1]?.ts : undefined
const cold = await coldFill<ChatViewRow>(
'chat',
until === undefined ? {} : { until },
new Set(hot.map((r) => r.id)),
limit - hot.length,
(r) => {
const rec = r.record as ChatViewRow
return { ...rec, archived: true, hidden: hidden.has(rec.messageId) }
},
)
return { rows: [...hot, ...cold].slice(0, limit), cold: cold.length > 0 }
},
async recall(query: RecallQuery): Promise<RecallResponse> {
const limit = Math.min(query.limit ?? 50, 500)
const q = query.q?.toLowerCase() || undefined
const range: { since?: string; until?: string } = {}
if (query.since !== undefined) range.since = query.since
if (query.until !== undefined) range.until = query.until
switch (query.kind) {
case 'audit': {
const hot = (
db
.prepare(
`SELECT * FROM audit_events
WHERE (? IS NULL OR ts >= ?) AND (? IS NULL OR ts <= ?)
AND (? IS NULL OR instr(unicode_lower(kind), ?) > 0 OR instr(unicode_lower(payload), ?) > 0)
ORDER BY ts DESC, id DESC LIMIT ?`,
)
.all(
query.since ?? null,
query.since ?? null,
query.until ?? null,
query.until ?? null,
q ?? null,
q ?? null,
q ?? null,
limit,
) as unknown as AuditRowDb[]
).map((r) => auditDbToRow(r))
const cold = await coldSearch<TimelineRow>(
'audit',
range,
new Set(hot.map((r) => r.id)),
(r) => {
const rec = r.record as {
kind: string
sessionId?: string
provenance?: TimelineRow['provenance']
payload?: unknown
}
return {
id: r.id,
ts: r.ts,
kind: rec.kind,
sessionId: rec.sessionId,
provenance: rec.provenance,
payload: rec.payload,
archived: true,
} satisfies TimelineRow
},
(r) => {
const rec = r.record as { kind: string; payload?: unknown }
return (
(!q ||
rec.kind.toLowerCase().includes(q) ||
JSON.stringify(rec.payload ?? null)
.toLowerCase()
.includes(q)) &&
(!query.since || r.ts >= query.since) &&
(!query.until || r.ts <= query.until)
)
},
)
const rows = [...hot, ...cold].sort(newestFirst).slice(0, limit)
return { rows, cold: rows.some((r) => r.archived === true) }
}
case 'chat': {
const hidden = hiddenMessageIds()
const hot = (
db
.prepare(
`SELECT * FROM chat_messages
WHERE (? IS NULL OR ts >= ?) AND (? IS NULL OR ts <= ?)
AND (? IS NULL OR instr(unicode_lower(text), ?) > 0
OR instr(unicode_lower(sender_name), ?) > 0)
AND NOT EXISTS (
SELECT 1 FROM tombstones WHERE tombstones.message_id = chat_messages.message_id
)
ORDER BY ts DESC, id DESC LIMIT ?`,
)
.all(
query.since ?? null,
query.since ?? null,
query.until ?? null,
query.until ?? null,
q ?? null,
q ?? null,
q ?? null,
limit,
) as unknown as ChatRowDb[]
).map((r) => chatDbToRow(r, false))
const cold = await coldSearch<ChatViewRow>(
'chat',
range,
new Set(hot.map((r) => r.id)),
(r) => ({ ...(r.record as ChatViewRow), archived: true, hidden: false }),
(r) => {
const rec = r.record as ChatViewRow
return (
!hidden.has(rec.messageId) &&
(!q ||
(rec.text ?? '').toLowerCase().includes(q) ||
(rec.senderName ?? '').toLowerCase().includes(q)) &&
(!query.since || r.ts >= query.since) &&
(!query.until || r.ts <= query.until)
)
},
)
const rows = [...hot, ...cold].sort(newestFirst).slice(0, limit)
return { rows, cold: rows.some((r) => r.archived === true) }
}
case 'deploy': {
const rows = db
.prepare(
`SELECT * FROM deployment_events
WHERE (? IS NULL OR ts >= ?) AND (? IS NULL OR ts <= ?)
ORDER BY ts DESC LIMIT ?`,
)
.all(
query.since ?? null,
query.since ?? null,
query.until ?? null,
query.until ?? null,
limit,
) as { id: string; ts: string; kind: string; revision: string | null; detail: string }[]
return {
rows: rows.map((r) => ({
id: r.id,
ts: r.ts,
kind: r.kind,
...(r.revision !== null ? { revision: r.revision } : {}),
detail: safeParse<unknown>(r.detail) ?? r.detail,
})),
cold: false,
}
}
case 'analysis': {
if (!query.digest) return { rows: [], cold: false }
const rows = db
.prepare('SELECT * FROM media_analyses WHERE digest = ? ORDER BY created_at DESC')
.all(query.digest) as {
digest: string
capability: string
model: string
schema_version: number
result: string
usage_json: string | null
created_at: string
}[]
return {
rows: rows.map((r) => ({
digest: r.digest,
capability: r.capability,
model: r.model,
schemaVersion: r.schema_version,
result: safeParse<unknown>(r.result) ?? r.result,
usage: safeParse<unknown>(r.usage_json ?? 'null') ?? undefined,
createdAt: r.created_at,
})),
cold: false,
}
}
}
},
async archiveCycle(): Promise<ArchiveCycleResult> {
return recordJob('archive_cycle', async () => {
const audit = await runKind('audit', DEFAULT_RETENTION.auditHotDays, loadAuditRows)
const chat = await runKind('chat', DEFAULT_RETENTION.chatHotDays, loadChatRows)
return {
exportedAudit: audit.exported,
exportedChat: chat.exported,
segments: audit.segments + chat.segments,
} satisfies ArchiveCycleResult
})
},
quickCheck(): Promise<string> {
return recordJob('quick_check', () => {
const row = db.prepare('PRAGMA quick_check').get() as { quick_check: string }
return row.quick_check
})
},
vacuum(): Promise<void> {
return recordJob('vacuum', () => {
db.exec('PRAGMA incremental_vacuum')
})
},
reconcile(): Promise<{ removedTmpFiles: number }> {
return recordJob('reconcile', async () => ({
removedTmpFiles: await casReconcile(dirs.mediaDir),
}))
},
setStorageThresholds(warnPct: number, stopPct: number): void {
db.prepare(
`INSERT INTO meta (key, value) VALUES ('storage_warn_pct', ?)
ON CONFLICT(key) DO UPDATE SET value = excluded.value`,
).run(String(warnPct))
db.prepare(
`INSERT INTO meta (key, value) VALUES ('storage_stop_pct', ?)
ON CONFLICT(key) DO UPDATE SET value = excluded.value`,
).run(String(stopPct))
},
counts(): { audit: number; chat: number; deployments: number; media: number } {
const n = (t: string) =>
Number((db.prepare(`SELECT COUNT(*) AS c FROM ${t}`).get() as { c: number }).c)
return {
audit: n('audit_events'),
chat: n('chat_messages'),
deployments: n('deployment_events'),
media: n('media'),
}
},
async usage(): Promise<UsageReport> {
const rows = db
.prepare("SELECT ts, payload FROM audit_events WHERE kind = 'turn_usage' ORDER BY ts")
.all() as unknown as { ts: string; payload: string | null }[]
return aggregateTurnUsage(
rows.map((r) => ({
ts: r.ts,
payload: r.payload === null ? null : (safeParse(r.payload) ?? null),
})),
clock.now(),
)
},
async status(): Promise<ServiceStatus> {
const dbStat = await stat(join(dirs.dataDir, 'will.db'))
let freeBytes = 0
let totalBytes = 0
try {
const fsStat = await statfs(dirs.dataDir)
freeBytes = Number(fsStat.bavail) * Number(fsStat.bsize)
totalBytes = Number(fsStat.blocks) * Number(fsStat.bsize)
} catch {
// statfs unavailable on some filesystems; report zeros and stay functional
}
const freePct = totalBytes > 0 ? (freeBytes / totalBytes) * 100 : 100
const warn = db.prepare("SELECT value FROM meta WHERE key = 'storage_warn_pct'").get() as
| { value: string }
| undefined
const stop = db.prepare("SELECT value FROM meta WHERE key = 'storage_stop_pct'").get() as
| { value: string }
| undefined
const warnPct = warn ? Number(warn.value) : 20
const stopPct = stop ? Number(stop.value) : 10
const storageState: ServiceStatus['storageState'] =
freePct < stopPct ? 'stop' : freePct < warnPct ? 'warn' : 'ok'
const lastCheck = db
.prepare(
"SELECT finished_at FROM jobs WHERE kind = 'quick_check' AND status = 'done' ORDER BY id DESC LIMIT 1",
)
.get() as { finished_at: string } | undefined
const status: ServiceStatus = {
ready: true,
schemaVersion: schemaVersion(db),
lastAuditSeq: ops.lastSeq(),
counts: ops.counts(),
dbSizeBytes: dbStat.size,
storage: { freeBytes, totalBytes, freePct: Math.round(freePct * 10) / 10 },
storageState,
}
if (lastCheck?.finished_at !== undefined) status.lastQuickCheck = lastCheck.finished_at
return status
},
close(): void {
db.exec('PRAGMA wal_checkpoint(TRUNCATE)')
db.close()
},
}
return ops
}
export type Ops = ReturnType<typeof createOps>