Skip to content
← 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 · TypeScript
import { 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
}