← Public packages
@kentcdodds/kody-issue-triage
Loop-safe triage for Kody run errors and fleet package-runtime error-rate elevations. Wakes Cole (Grok Bot) to decide; Cursor agents only when Cole escalates.
src/storage.ts
474 lines · 13.6 KB · TypeScriptimport { packageStorage } from 'kody:runtime'
import type { FingerprintRecord, RunSnippet } from './shared.ts'
import {
activityUrlFor,
fingerprintFor,
leaseStaleMs,
nextFingerprintStatusAfterEvent,
normalizeErrorFamily,
ownerFor,
} from './shared.ts'
function storage() {
return packageStorage()
}
export async function sql(query: string, bindings: unknown[] = []) {
return await storage().sql(query, bindings)
}
function rowsOf(result: { rows?: unknown[] } | unknown[] | null | undefined) {
if (!result) return []
if (Array.isArray(result)) return result
return Array.isArray(result.rows) ? result.rows : []
}
export async function initSchema() {
await sql(`
CREATE TABLE IF NOT EXISTS fingerprints (
fingerprint TEXT PRIMARY KEY,
surface TEXT NOT NULL,
owner TEXT NOT NULL,
package_id TEXT,
kody_id TEXT,
error_family TEXT NOT NULL,
sample_message TEXT,
sample_run_id TEXT,
activity_url TEXT,
sample_run_ids_json TEXT NOT NULL,
count INTEGER NOT NULL,
first_seen TEXT NOT NULL,
last_seen TEXT NOT NULL,
status TEXT NOT NULL,
agent_id TEXT,
agent_url TEXT,
kody_agent_id TEXT,
kody_agent_url TEXT,
classification TEXT,
discord_message_id TEXT,
spawned_at TEXT,
completed_at TEXT,
outcome TEXT,
summary TEXT
)
`)
await sql(`
CREATE TABLE IF NOT EXISTS counters (
key TEXT PRIMARY KEY,
value INTEGER NOT NULL,
updated_at TEXT NOT NULL
)
`)
await sql(`
CREATE TABLE IF NOT EXISTS lease (
id TEXT PRIMARY KEY,
holder TEXT,
fingerprint TEXT,
acquired_at TEXT
)
`)
await sql(`
CREATE TABLE IF NOT EXISTS config (
key TEXT PRIMARY KEY,
value TEXT NOT NULL
)
`)
await ensureFingerprintColumn('kody_agent_id', 'TEXT')
await ensureFingerprintColumn('kody_agent_url', 'TEXT')
await ensureFingerprintColumn('classification', 'TEXT')
await ensureFingerprintColumn('pr_url', 'TEXT')
await ensureFingerprintColumn('merge_commit_url', 'TEXT')
await ensureFingerprintColumn('deploy_url', 'TEXT')
await ensureFingerprintColumn('card_stage', 'TEXT')
}
async function ensureFingerprintColumn(column: string, type: string) {
try {
await sql(`ALTER TABLE fingerprints ADD COLUMN ${column} ${type}`)
} catch {
// Column already exists on previously published schemas.
}
}
function parseRecord(row: Record<string, unknown>): FingerprintRecord {
let sampleRunIds: string[] = []
try {
const parsed = JSON.parse(String(row.sample_run_ids_json || '[]'))
if (Array.isArray(parsed)) sampleRunIds = parsed.map(String)
} catch {
sampleRunIds = []
}
return {
fingerprint: String(row.fingerprint),
surface: String(row.surface),
owner: String(row.owner),
package_id: row.package_id ? String(row.package_id) : null,
kody_id: row.kody_id ? String(row.kody_id) : null,
error_family: String(row.error_family),
sample_message: String(row.sample_message || ''),
sample_run_id: row.sample_run_id ? String(row.sample_run_id) : null,
activity_url: row.activity_url ? String(row.activity_url) : null,
sample_run_ids: sampleRunIds,
count: Number(row.count || 0),
first_seen: String(row.first_seen),
last_seen: String(row.last_seen),
status: String(row.status),
agent_id: row.agent_id ? String(row.agent_id) : null,
agent_url: row.agent_url ? String(row.agent_url) : null,
kody_agent_id: row.kody_agent_id ? String(row.kody_agent_id) : null,
kody_agent_url: row.kody_agent_url ? String(row.kody_agent_url) : null,
classification: row.classification ? String(row.classification) : null,
discord_message_id: row.discord_message_id
? String(row.discord_message_id)
: null,
spawned_at: row.spawned_at ? String(row.spawned_at) : null,
completed_at: row.completed_at ? String(row.completed_at) : null,
outcome: row.outcome ? String(row.outcome) : null,
summary: row.summary ? String(row.summary) : null,
pr_url: row.pr_url ? String(row.pr_url) : null,
merge_commit_url: row.merge_commit_url
? String(row.merge_commit_url)
: null,
deploy_url: row.deploy_url ? String(row.deploy_url) : null,
card_stage: row.card_stage ? String(row.card_stage) : null,
}
}
export async function getFingerprint(fingerprint: string) {
const result = await sql('SELECT * FROM fingerprints WHERE fingerprint = ?', [
fingerprint,
])
const row = rowsOf(result)[0] as Record<string, unknown> | undefined
return row ? parseRecord(row) : null
}
export async function listFingerprintsByStatus(status: string, limit = 20) {
const result = await sql(
'SELECT * FROM fingerprints WHERE status = ? ORDER BY first_seen ASC LIMIT ?',
[status, limit],
)
return rowsOf(result).map((row) => parseRecord(row as Record<string, unknown>))
}
export async function listFingerprintsByStatuses(statuses: string[], limit = 80) {
if (!statuses.length) return []
const placeholders = statuses.map(() => '?').join(',')
const result = await sql(
`SELECT * FROM fingerprints WHERE status IN (${placeholders}) ORDER BY last_seen DESC LIMIT ?`,
[...statuses, limit],
)
return rowsOf(result).map((row) => parseRecord(row as Record<string, unknown>))
}
export async function listRecentFingerprints(limit = 20) {
const result = await sql(
'SELECT * FROM fingerprints ORDER BY last_seen DESC LIMIT ?',
[limit],
)
return rowsOf(result).map((row) => parseRecord(row as Record<string, unknown>))
}
export async function listSiblingFingerprints(
record: Pick<FingerprintRecord, 'fingerprint' | 'owner' | 'kody_id' | 'error_family'>,
limit = 20,
) {
const result = await sql(
`SELECT * FROM fingerprints
WHERE fingerprint != ?
AND (owner = ? OR error_family = ? OR kody_id = ?)
ORDER BY last_seen DESC
LIMIT ?`,
[
record.fingerprint,
record.owner,
record.error_family,
record.kody_id,
limit,
],
)
return rowsOf(result).map((row) => parseRecord(row as Record<string, unknown>))
}
export async function listStaleSpawned(now = Date.now()) {
const cutoff = new Date(now - leaseStaleMs).toISOString()
const result = await sql(
`SELECT * FROM fingerprints
WHERE status = 'spawned'
AND (spawned_at IS NULL OR spawned_at < ?)
ORDER BY spawned_at ASC
LIMIT 20`,
[cutoff],
)
return rowsOf(result).map((row) => parseRecord(row as Record<string, unknown>))
}
export async function upsertFingerprintFromRun(run: RunSnippet) {
const fingerprint = fingerprintFor(run)
const now = new Date().toISOString()
const existing = await getFingerprint(fingerprint)
const runId = run.id || null
const sampleIds = existing?.sample_run_ids ? [...existing.sample_run_ids] : []
if (runId && !sampleIds.includes(runId)) {
sampleIds.unshift(runId)
if (sampleIds.length > 20) sampleIds.length = 20
}
if (!existing) {
await sql(
`INSERT INTO fingerprints (
fingerprint, surface, owner, package_id, kody_id, error_family,
sample_message, sample_run_id, activity_url, sample_run_ids_json,
count, first_seen, last_seen, status
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1, ?, ?, 'queued')
ON CONFLICT(fingerprint) DO UPDATE SET
count = fingerprints.count + 1,
last_seen = excluded.last_seen,
sample_message = COALESCE(excluded.sample_message, fingerprints.sample_message),
sample_run_id = COALESCE(excluded.sample_run_id, fingerprints.sample_run_id),
activity_url = COALESCE(excluded.activity_url, fingerprints.activity_url),
sample_run_ids_json = excluded.sample_run_ids_json`,
[
fingerprint,
run.surface || 'unknown',
ownerFor(run),
run.package_id || null,
run.kody_id || null,
normalizeErrorFamily(run.error_message),
String(run.error_message || '').slice(0, 500),
runId,
run.activity_url || activityUrlFor(runId),
JSON.stringify(sampleIds),
now,
now,
],
)
const after = await getFingerprint(fingerprint)
const count = after?.count || 1
return { fingerprint, created: count === 1, count }
}
const nextStatus = nextFingerprintStatusAfterEvent({
status: existing.status,
completedAt: existing.completed_at,
kodyAgentId: existing.kody_agent_id,
runStartedAt: run.started_at,
})
await sql(
`UPDATE fingerprints SET
count = count + 1,
last_seen = ?,
sample_message = ?,
sample_run_id = ?,
activity_url = ?,
sample_run_ids_json = ?,
status = ?
WHERE fingerprint = ?`,
[
now,
String(run.error_message || existing.sample_message || '').slice(0, 500),
runId || existing.sample_run_id,
run.activity_url || activityUrlFor(runId) || existing.activity_url,
JSON.stringify(sampleIds),
nextStatus,
fingerprint,
],
)
return { fingerprint, created: false, count: existing.count + 1, nextStatus }
}
export async function upsertFleetRateFingerprint(input: {
fingerprint: string
window: 'hour' | 'day'
message: string
insightsUrl: string | null
}) {
const now = new Date().toISOString()
const existing = await getFingerprint(input.fingerprint)
const message = String(input.message || '').slice(0, 500)
if (!existing) {
await sql(
`INSERT INTO fingerprints (
fingerprint, surface, owner, package_id, kody_id, error_family,
sample_message, sample_run_id, activity_url, sample_run_ids_json,
count, first_seen, last_seen, status
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1, ?, ?, 'queued')
ON CONFLICT(fingerprint) DO UPDATE SET
count = fingerprints.count + 1,
last_seen = excluded.last_seen,
sample_message = COALESCE(excluded.sample_message, fingerprints.sample_message),
activity_url = COALESCE(excluded.activity_url, fingerprints.activity_url)`,
[
input.fingerprint,
'fleet-rate',
'kody-platform',
null,
null,
`package-error-rate-${input.window}`,
message,
null,
input.insightsUrl,
'[]',
now,
now,
],
)
const after = await getFingerprint(input.fingerprint)
const count = after?.count || 1
return { fingerprint: input.fingerprint, created: count === 1, count }
}
const nextStatus = nextFingerprintStatusAfterEvent({
status: existing.status,
completedAt: existing.completed_at,
kodyAgentId: existing.kody_agent_id,
})
await sql(
`UPDATE fingerprints SET
count = count + 1,
last_seen = ?,
sample_message = ?,
activity_url = ?,
status = ?
WHERE fingerprint = ?`,
[
now,
message || existing.sample_message,
input.insightsUrl || existing.activity_url,
nextStatus,
input.fingerprint,
],
)
return {
fingerprint: input.fingerprint,
created: false,
count: existing.count + 1,
nextStatus,
}
}
export async function updateFingerprint(
fingerprint: string,
fields: Partial<FingerprintRecord>,
) {
const allowed: Array<keyof FingerprintRecord> = [
'status',
'agent_id',
'agent_url',
'kody_agent_id',
'kody_agent_url',
'classification',
'discord_message_id',
'spawned_at',
'completed_at',
'outcome',
'summary',
'sample_run_id',
'activity_url',
'pr_url',
'merge_commit_url',
'deploy_url',
'card_stage',
]
const sets: string[] = []
const bindings: unknown[] = []
for (const key of allowed) {
if (fields[key] !== undefined) {
sets.push(`${key} = ?`)
bindings.push(fields[key])
}
}
if (fields.sample_run_ids) {
sets.push('sample_run_ids_json = ?')
bindings.push(JSON.stringify(fields.sample_run_ids))
}
if (!sets.length) return
bindings.push(fingerprint)
await sql(
`UPDATE fingerprints SET ${sets.join(', ')} WHERE fingerprint = ?`,
bindings,
)
}
export async function incrementCounter(key: string, by = 1) {
const now = new Date().toISOString()
await sql(
`INSERT INTO counters (key, value, updated_at) VALUES (?, ?, ?)
ON CONFLICT(key) DO UPDATE SET
value = counters.value + excluded.value,
updated_at = excluded.updated_at`,
[key, by, now],
)
return await getCounter(key)
}
export async function getCounter(key: string) {
const result = await sql('SELECT value FROM counters WHERE key = ?', [key])
const row = rowsOf(result)[0] as { value?: number } | undefined
return Number(row?.value || 0)
}
export async function getConfig(key: string, fallback = '') {
const result = await sql('SELECT value FROM config WHERE key = ?', [key])
const row = rowsOf(result)[0] as { value?: string } | undefined
return row?.value != null ? String(row.value) : fallback
}
export async function setConfig(key: string, value: string) {
await sql('INSERT OR REPLACE INTO config (key, value) VALUES (?, ?)', [
key,
value,
])
}
export async function isEnabled() {
const value = await getConfig('enabled', 'true')
return value !== 'false'
}
export async function tryAcquireLease(fingerprint: string, holder: string) {
const now = Date.now()
const existing = await sql('SELECT * FROM lease WHERE id = ?', ['global'])
const row = rowsOf(existing)[0] as
| { holder?: string; fingerprint?: string; acquired_at?: string }
| undefined
if (row?.acquired_at) {
const age = now - Date.parse(row.acquired_at)
const sameHolder = row.fingerprint === fingerprint
if (age < 2 * 60 * 60 * 1000 && !sameHolder) {
return {
ok: false as const,
heldBy: row.fingerprint || row.holder || 'unknown',
}
}
}
await sql(
'INSERT OR REPLACE INTO lease (id, holder, fingerprint, acquired_at) VALUES (?, ?, ?, ?)',
['global', holder, fingerprint, new Date(now).toISOString()],
)
return { ok: true as const }
}
export async function releaseLease(fingerprint?: string) {
if (fingerprint) {
await sql('DELETE FROM lease WHERE id = ? AND fingerprint = ?', [
'global',
fingerprint,
])
return
}
await sql('DELETE FROM lease WHERE id = ?', ['global'])
}
export async function getLease() {
const result = await sql('SELECT * FROM lease WHERE id = ?', ['global'])
const row = rowsOf(result)[0] as
| { holder?: string; fingerprint?: string; acquired_at?: string }
| undefined
return row
? {
holder: row.holder || null,
fingerprint: row.fingerprint || null,
acquired_at: row.acquired_at || null,
}
: null
}