Skip to content
← Public packages

@kentcdodds/sentry-triage

Sentry triage wakes Cole (Grok Bot) per issue: one active wake per repo (lease + queue), loop-safe Discord. Cole may spawn Cursor for isolated repo work.

src/triage-core.ts

695 lines · 21.0 KB · TypeScript
/**
 * Shared triage primitives used by both the issue-webhook handler and the
 * Seer-webhook handler so hourly caps, claims, and Cole wake cannot
 * drift between the two ingress paths. Both paths share the same storage
 * counter keys, so the hourly wake cap is enforced across them jointly.
 *
 * Also owns the per-repo concurrency slots + queue: up to
 * MAX_CONCURRENT_WAKES_PER_REPO (default 3) active Cole wakes per repository.
 * Extra issues enqueue while slots are full; when Cole records an outcome
 * (or a slot is released / swept stale), free slots pull queued issues and
 * wake one Cole per issue so work parallelizes. Caps still prevent the old
 * unbounded fan-out that collided on shared PRs.
 */

import { COLE_AGENT_ID, COLE_AGENT_URL, wakeCole } from './cole-handoff.ts'
import { editSentryReport } from './format-discord-report.ts'
import sentryRequest from 'kody:@kentcdodds/sentry/request'
import { buildTriageAgentPrompt } from './agent-prompt.ts'
import { gatherSpawnContext } from './spawn-snapshot.ts'
import {
	hourBucket,
	issueStorageKey,
	spawnedCounterKey,
	truncate,
} from './shared.ts'
import {
	promptOwnerFields,
	readTriageSettings,
	resolveTriageProjectForSlug,
} from './settings.ts'

export {
	awaitingSeerGraceMs,
	awaitingSeerSinceMs,
	findStaleAwaitingIssues,
	isAwaitingSeerStale,
} from './awaiting-seer.ts'

/** Matches the in-flight stale window in handle-sentry-webhook (interrupted wakes). */
export const repoLeaseStaleMs = 3 * 60 * 60 * 1000

/** Default max concurrent Cole wakes per project slug (overridable via settings). */
export const MAX_CONCURRENT_WAKES_PER_REPO = 3

/** Wake Cole (Grok Bot). Cole may choose to escalate repo work to Cursor. */
export { wakeCole }

/**
 * Request a Seer root-cause analysis for an issue (the "pull" trigger). Free
 * Seer allows unlimited RCA runs. The completed RCA is delivered back via the
 * `seer.root_cause_completed` webhook, which seeds Cole. Returns
 * `{ ok: false }` (with a status/reason) when Seer cannot run for this issue,
 * so the caller can fall back to from-scratch triage.
 */
export async function triggerSeerRca(issueId) {
	try {
		const json = await sentryRequest({
			path: `/issues/${issueId}/autofix/`,
			init: {
				method: 'POST',
				headers: { 'content-type': 'application/json' },
				body: JSON.stringify({}),
			},
		})
		return { ok: true, runId: json?.run_id ?? null }
	} catch (error) {
		// SentryApiError carries the HTTP status; keep the { ok, status } shape
		// the callers use to decide on the from-scratch fallback.
		const status = typeof error?.status === 'number' ? error.status : null
		return {
			ok: false,
			...(status != null ? { status } : {}),
			error: error instanceof Error ? error.message : String(error),
		}
	}
}

/**
 * Whether Seer's automated runs for this project continue past RCA all the
 * way to opening a draft PR (project setting: Seer stopping point
 * "open_pr"). Decides per delivery whether root_cause_completed should hold
 * for the PR (ship-agent flow) or wake Cole for investigation immediately —
 * reading it live means flipping the Sentry setting per project is the only
 * rollout step. Fails closed to the investigation flow.
 */
export async function seerRunsToPr(projectSlug) {
	try {
		const prefs = await sentryRequest({
			path: `/projects/${(await readTriageSettings()).sentryOrgSlug}/${projectSlug}/seer/preferences/`,
		})
		return prefs?.preference?.automated_run_stopping_point === 'open_pr'
	} catch {
		return false
	}
}

export async function bumpCounter(storage, key) {
	const current = Number((await storage.get(key)) ?? 0)
	const next = current + 1
	await storage.set(key, next)
	return next
}

/**
 * Atomic one-winner claim on an issue id via the shared triage_claims table.
 * Returns true when this caller won the claim (should proceed), false when a
 * concurrent delivery already holds it.
 */
export async function claimIssueOnce(storage, issueId) {
	await storage.sql(
		`CREATE TABLE IF NOT EXISTS triage_claims (
			issue_id TEXT PRIMARY KEY,
			claimed_at TEXT NOT NULL
		)`,
	)
	try {
		await storage.sql(
			`INSERT INTO triage_claims (issue_id, claimed_at) VALUES (?1, ?2)`,
			[String(issueId), new Date().toISOString()],
		)
		return true
	} catch {
		return false
	}
}

async function ensureRepoLeaseTable(storage) {
	await migrateRepoLeaseTable(storage)
}

/**
 * Multi-slot lease table: one row per active wake (project_slug, holder).
 * Migrates the legacy single-holder schema (project_slug PRIMARY KEY) in place.
 */
async function migrateRepoLeaseTable(storage) {
	// Detect legacy single-holder schema (project_slug PRIMARY KEY only).
	let needsMigrate = false
	let hasLeases = false
	try {
		const info = await storage.sql(`PRAGMA table_info(repo_leases)`)
		const cols = info?.results ?? info?.rows ?? []
		hasLeases = Array.isArray(cols) && cols.length > 0
		if (hasLeases) {
			const pkCols = cols.filter(
				(row) => Number(row?.pk ?? row?.PK ?? 0) > 0,
			)
			// Legacy: sole PK column is project_slug. Multi-slot: composite PK.
			needsMigrate =
				pkCols.length === 1 &&
				String(pkCols[0]?.name ?? '') === 'project_slug'
		}
	} catch {
		hasLeases = false
	}

	if (!needsMigrate && hasLeases) {
		// Multi-slot table already in place.
		return
	}

	await storage.sql(
		`CREATE TABLE IF NOT EXISTS repo_lease_slots (
			project_slug TEXT NOT NULL,
			holder TEXT NOT NULL,
			claimed_at TEXT NOT NULL,
			PRIMARY KEY (project_slug, holder)
		)`,
	)

	if (needsMigrate) {
		try {
			const legacy = await storage.sql(
				`SELECT project_slug, holder, claimed_at FROM repo_leases`,
			)
			const rows = legacy?.results ?? legacy?.rows ?? []
			if (Array.isArray(rows)) {
				for (const row of rows) {
					if (!row?.project_slug || !row?.holder) continue
					try {
						await storage.sql(
							`INSERT INTO repo_lease_slots (project_slug, holder, claimed_at)
							VALUES (?1, ?2, ?3)`,
							[
								String(row.project_slug),
								String(row.holder),
								String(row.claimed_at ?? new Date().toISOString()),
							],
						)
					} catch {
						// duplicate / already copied
					}
				}
			}
			await storage.sql(`DROP TABLE IF EXISTS repo_leases`)
		} catch {
			// ignore
		}
	}

	if (!hasLeases || needsMigrate) {
		try {
			await storage.sql(`ALTER TABLE repo_lease_slots RENAME TO repo_leases`)
		} catch {
			// repo_leases already exists as multi-slot
		}
	}
	try {
		await storage.sql(`DROP TABLE IF EXISTS repo_lease_slots`)
	} catch {
		// ignore
	}
}

async function ensureRepoQueueTable(storage) {
	await storage.sql(
		`CREATE TABLE IF NOT EXISTS repo_queue (
			issue_id TEXT PRIMARY KEY,
			project_slug TEXT NOT NULL,
			payload TEXT NOT NULL,
			queued_at TEXT NOT NULL
		)`,
	)
}

/** Resolve the per-repo concurrency limit (settings override, else default 3). */
export async function getMaxConcurrentWakesPerRepo() {
	try {
		const settings = await readTriageSettings()
		const raw = settings?.maxConcurrentWakesPerRepo
		const n = typeof raw === 'number' ? raw : Number(raw)
		if (Number.isFinite(n) && n >= 1 && n <= 20) return Math.floor(n)
	} catch {
		// Settings read failures fall back to the named default.
	}
	return MAX_CONCURRENT_WAKES_PER_REPO
}

async function sweepStaleRepoLeases(storage, projectSlug) {
	await ensureRepoLeaseTable(storage)
	const staleBefore = new Date(Date.now() - repoLeaseStaleMs).toISOString()
	await storage.sql(
		`DELETE FROM repo_leases WHERE project_slug = ?1 AND claimed_at < ?2`,
		[String(projectSlug), staleBefore],
	)
}

export async function countActiveRepoLeases(storage, projectSlug) {
	await ensureRepoLeaseTable(storage)
	await sweepStaleRepoLeases(storage, projectSlug)
	const result = await storage.sql(
		`SELECT COUNT(*) AS n FROM repo_leases WHERE project_slug = ?1`,
		[String(projectSlug)],
	)
	const rows = result?.results ?? result?.rows ?? []
	const row = Array.isArray(rows) ? rows[0] : null
	return Number(row?.n ?? row?.['COUNT(*)'] ?? 0)
}

/**
 * Claim one concurrency slot for this project. Holder should be the issue id
 * so record-outcome can release that exact slot. Returns true when a slot was
 * claimed (or this holder already has one); false when the repo is at capacity.
 */
export async function acquireRepoLease(storage, projectSlug, holder) {
	await ensureRepoLeaseTable(storage)
	await sweepStaleRepoLeases(storage, projectSlug)
	const limit = await getMaxConcurrentWakesPerRepo()
	const active = await countActiveRepoLeases(storage, projectSlug)
	const now = new Date().toISOString()
	const slug = String(projectSlug)
	const hold = String(holder)
	// Re-entrant: same issue already holds a slot.
	try {
		const existing = await storage.sql(
			`SELECT holder FROM repo_leases WHERE project_slug = ?1 AND holder = ?2 LIMIT 1`,
			[slug, hold],
		)
		const rows = existing?.results ?? existing?.rows ?? []
		if (Array.isArray(rows) && rows.length > 0) {
			await storage.sql(
				`UPDATE repo_leases SET claimed_at = ?3 WHERE project_slug = ?1 AND holder = ?2`,
				[slug, hold, now],
			)
			return true
		}
	} catch {
		// continue to insert path
	}
	if (active >= limit) return false
	try {
		await storage.sql(
			`INSERT INTO repo_leases (project_slug, holder, claimed_at) VALUES (?1, ?2, ?3)`,
			[slug, hold, now],
		)
		return true
	} catch {
		// Race: another writer filled the last slot or this holder just inserted.
		const existing = await storage.sql(
			`SELECT holder FROM repo_leases WHERE project_slug = ?1 AND holder = ?2 LIMIT 1`,
			[slug, hold],
		)
		const rows = existing?.results ?? existing?.rows ?? []
		return Array.isArray(rows) && rows.length > 0
	}
}

/**
 * Release one slot. When `holder` is provided, only that slot is freed so
 * sibling wakes keep their slots. Omit `holder` to clear every slot for the
 * project (force / maintenance).
 */
export async function releaseRepoLease(storage, projectSlug, holder) {
	await ensureRepoLeaseTable(storage)
	if (holder != null && holder !== '') {
		await storage.sql(
			`DELETE FROM repo_leases WHERE project_slug = ?1 AND holder = ?2`,
			[String(projectSlug), String(holder)],
		)
		return
	}
	await storage.sql(`DELETE FROM repo_leases WHERE project_slug = ?1`, [
		String(projectSlug),
	])
}

/**
 * Heartbeat a slot's claimed_at. Holder must be the issue id that acquired the
 * slot — do not retarget to a shared agent name (that would collide under
 * multi-slot concurrency).
 */
export async function updateRepoLeaseHolder(storage, projectSlug, holder) {
	await ensureRepoLeaseTable(storage)
	await storage.sql(
		`UPDATE repo_leases SET claimed_at = ?3 WHERE project_slug = ?1 AND holder = ?2`,
		[String(projectSlug), String(holder), new Date().toISOString()],
	)
}

/**
 * Enqueue an issue for later flush when all concurrency slots are busy.
 * Payload is JSON with everything needed to wake later. INSERT-on-primary-key
 * makes re-enqueue of the same issue a no-op for the first writer.
 */
export async function enqueueRepoIssue(storage, projectSlug, issueId, payload) {
	await ensureRepoQueueTable(storage)
	try {
		await storage.sql(
			`INSERT INTO repo_queue (issue_id, project_slug, payload, queued_at)
			VALUES (?1, ?2, ?3, ?4)`,
			[
				String(issueId),
				String(projectSlug),
				JSON.stringify(payload),
				new Date().toISOString(),
			],
		)
		return true
	} catch {
		return false
	}
}

/**
 * Side-effect-free capacity probe for dry-run: true when active wakes are at
 * the per-repo limit (new issues would queue).
 */
export async function isRepoLeaseHeld(storage, projectSlug) {
	const limit = await getMaxConcurrentWakesPerRepo()
	const active = await countActiveRepoLeases(storage, projectSlug)
	return active >= limit
}

/** Side-effect-free queue depth probe for dry-run. */
export async function getRepoQueueDepth(storage, projectSlug) {
	await ensureRepoQueueTable(storage)
	const result = await storage.sql(
		`SELECT COUNT(*) AS n FROM repo_queue WHERE project_slug = ?1`,
		[String(projectSlug)],
	)
	const rows = result?.results ?? result?.rows ?? []
	const row = Array.isArray(rows) ? rows[0] : null
	return Number(row?.n ?? row?.['COUNT(*)'] ?? 0)
}

/**
 * Find project slugs whose queue can make progress: free capacity after sweeping
 * stale slots, or a force-worthy stale/missing lease with queued work. Caps to
 * a few repos per call so a webhook never becomes a long sweeper.
 */
export async function findReposNeedingQueueFlush(storage) {
	await ensureRepoLeaseTable(storage)
	await ensureRepoQueueTable(storage)
	const limit = await getMaxConcurrentWakesPerRepo()
	const staleBefore = new Date(Date.now() - repoLeaseStaleMs).toISOString()
	const queued = await storage.sql(
		`SELECT DISTINCT project_slug AS project_slug FROM repo_queue LIMIT 20`,
	)
	const rows = queued?.results ?? queued?.rows ?? []
	if (!Array.isArray(rows)) return []
	const needing = []
	for (const row of rows) {
		const slug = row?.project_slug
		if (typeof slug !== 'string' || !slug) continue
		await sweepStaleRepoLeases(storage, slug)
		const active = await countActiveRepoLeases(storage, slug)
		// Free capacity, or stale rows were just swept leaving work to drain.
		if (active < limit) {
			needing.push(slug)
			if (needing.length >= 5) break
			continue
		}
		// Still at capacity but has stale rows that sweep missed? Re-check stale.
		const stale = await storage.sql(
			`SELECT COUNT(*) AS n FROM repo_leases
			WHERE project_slug = ?1 AND claimed_at < ?2`,
			[slug, staleBefore],
		)
		const staleRows = stale?.results ?? stale?.rows ?? []
		const staleN = Number(staleRows?.[0]?.n ?? staleRows?.[0]?.['COUNT(*)'] ?? 0)
		if (staleN > 0) {
			needing.push(slug)
			if (needing.length >= 5) break
		}
	}
	return needing
}

async function drainRepoQueue(storage, projectSlug, maxItems = Infinity) {
	await ensureRepoQueueTable(storage)
	const result = await storage.sql(
		`SELECT issue_id, payload, queued_at FROM repo_queue
		WHERE project_slug = ?1
		ORDER BY queued_at ASC`,
		[String(projectSlug)],
	)
	const rows = result?.results ?? result?.rows ?? []
	const all = Array.isArray(rows) ? rows : []
	if (all.length === 0) return []
	const capped =
		Number.isFinite(maxItems) && maxItems >= 0
			? all.slice(0, Math.max(0, Math.floor(maxItems)))
			: all
	// Delete only the rows we selected — items enqueued between SELECT and
	// DELETE stay for the next flush (two statements are not atomic together).
	for (const row of capped) {
		await storage.sql(`DELETE FROM repo_queue WHERE issue_id = ?1`, [
			String(row.issue_id),
		])
	}
	return capped.map((row) => ({
		issueId: String(row.issue_id),
		payload: parseQueuePayload(row.payload),
		queuedAt: row.queued_at ?? null,
	}))
}

function parseQueuePayload(raw) {
	if (raw && typeof raw === 'object') return raw
	try {
		return JSON.parse(String(raw ?? '{}'))
	} catch {
		return {}
	}
}

function issueLinkFromPayload(payload, issueId) {
	return (
		payload?.issue?.permalink ??
		payload?.link ??
		`https://sentry.io/organizations/issues/${issueId}/`
	)
}

function issueTitleFromPayload(payload) {
	return truncate(payload?.issue?.title ?? 'Unknown issue', 200)
}

/**
 * Wake Cole for a single drained queue entry after its slot is already claimed.
 * Returns { ok, issueId } or { ok: false, error, issueId }.
 */
async function wakeColeForQueuedEntry(storage, project, entry) {
	const projectSlug = project.slug
	const payload = entry.payload
	const title = issueTitleFromPayload(payload)
	const link = issueLinkFromPayload(payload, entry.issueId)
	const discordId = payload?.discordMessageId ?? null
	const issueKey = issueStorageKey(entry.issueId)
	const existing = (await storage.get(issueKey)) ?? {}
	await storage.set(issueKey, {
		...existing,
		issueId: entry.issueId,
		projectSlug,
		shortId: payload?.issue?.shortId ?? existing.shortId ?? null,
		title,
		link,
		status: 'spawn-failed',
		wakePending: true,
		discordMessageId: discordId ?? existing.discordMessageId ?? null,
		seenAt: existing.seenAt ?? new Date().toISOString(),
	})
	if (discordId) {
		await editSentryReport(discordId, {
			status: 'queued',
			title,
			shortId: payload?.issue?.shortId ?? entry.issueId,
			issueId: entry.issueId,
			projectSlug,
			link,
			summary: 'Waking Cole…',
		}).catch(() => {})
	}

	const issue = {
		id: entry.issueId,
		shortId: payload?.issue?.shortId ?? null,
		title: payload?.issue?.title ?? 'Unknown issue',
		culprit: payload?.issue?.culprit ?? null,
		level: payload?.issue?.level ?? 'error',
		permalink: payload?.issue?.permalink ?? null,
	}
	const spawnContext = await gatherSpawnContext(issue.id, project)
	const promptText = buildTriageAgentPrompt({
		issue,
		project,
		discordMessageId: discordId ?? 'PENDING',
		context: spawnContext.sentryContext,
		issueState: spawnContext.issueState,
		priorRecord: payload?.isRetriage
			? { status: 'queued', note: 're-triage flushed from per-repo queue' }
			: null,
		seerRootCause: payload?.seerRootCause ?? null,
		backfill: null,
		...promptOwnerFields(await readTriageSettings()),
	})

	try {
		await wakeCole({
			promptText,
			discordMessageId: discordId,
			discordChannelId: (await readTriageSettings()).discordChannelId,
		})
	} catch (error) {
		await releaseRepoLease(storage, projectSlug, entry.issueId)
		const errText = truncate(
			error instanceof Error ? error.message : String(error),
			300,
		)
		await storage.set(issueKey, {
			...existing,
			issueId: entry.issueId,
			projectSlug,
			shortId: payload?.issue?.shortId ?? existing.shortId ?? null,
			title,
			link,
			status: 'spawn-failed',
			discordMessageId: discordId ?? existing.discordMessageId ?? null,
			seenAt: existing.seenAt ?? new Date().toISOString(),
			error: errText,
		})
		if (discordId) {
			await editSentryReport(discordId, {
				status: 'queued',
				title,
				shortId: payload?.issue?.shortId ?? entry.issueId,
				issueId: entry.issueId,
				projectSlug,
				link,
				summary: `Cole wake failed — will retry: ${errText}`,
			}).catch(() => {})
		}
		return { ok: false, error: 'cole-wake-failed', issueId: entry.issueId }
	}

	await updateRepoLeaseHolder(storage, projectSlug, entry.issueId)
	await bumpCounter(storage, spawnedCounterKey())
	const spawnedAt = new Date().toISOString()
	const agentUrl = COLE_AGENT_URL
	await storage.set(issueKey, {
		...existing,
		issueId: entry.issueId,
		projectSlug,
		shortId: payload?.issue?.shortId ?? existing.shortId ?? null,
		title,
		link,
		status: 'agent-spawned',
		wakePending: false,
		agentId: COLE_AGENT_ID,
		agentUrl,
		discordMessageId: discordId ?? existing.discordMessageId ?? null,
		spawnedAt,
		hourBucket: hourBucket(),
		queueFlush: true,
		seededBy: payload?.seerRootCause
			? existing.seededBy ?? 'seer'
			: existing.seededBy,
	})
	if (discordId) {
		await editSentryReport(discordId, {
			status: 'investigating',
			title,
			shortId: payload?.issue?.shortId ?? entry.issueId,
			issueId: entry.issueId,
			projectSlug,
			link,
			summary: 'Cole is triaging',
			agentUrl,
		}).catch(() => {})
	}
	return { ok: true, issueId: entry.issueId, agentId: COLE_AGENT_ID }
}

/**
 * Fill free per-repo slots from the queue — one Cole wake per claimed slot.
 * With `force`, clears every slot first (maintenance). Otherwise sweeps stale
 * slots, then drains up to (limit - active) queued issues. Used by
 * record-outcome, the lazy webhook sweeper, and `./flush-queue`. Does NOT
 * apply the hourly wake cap — queued issues already passed their gates.
 */
export async function flushRepoQueue(storage, projectSlug, options = {}) {
	const force = Boolean(options.force)
	await ensureRepoLeaseTable(storage)
	await ensureRepoQueueTable(storage)

	if (force) {
		await releaseRepoLease(storage, projectSlug)
	} else {
		await sweepStaleRepoLeases(storage, projectSlug)
	}

	const limit = await getMaxConcurrentWakesPerRepo()
	const active = await countActiveRepoLeases(storage, projectSlug)
	const free = Math.max(0, limit - active)
	if (free === 0) {
		return { ok: true, skipped: 'lease-held', projectSlug, flushed: 0 }
	}

	const queued = await drainRepoQueue(storage, projectSlug, free)
	if (queued.length === 0) {
		return { ok: true, flushed: 0, projectSlug }
	}

	const project = await resolveTriageProjectForSlug(projectSlug)
	if (!project) {
		for (const entry of queued) {
			const issueKey = issueStorageKey(entry.issueId)
			const existing = (await storage.get(issueKey)) ?? {}
			await storage.set(issueKey, {
				...existing,
				issueId: entry.issueId,
				projectSlug,
				status: 'spawn-failed',
				seenAt: existing.seenAt ?? new Date().toISOString(),
				error: 'unconfigured-project-on-flush',
			})
		}
		return { ok: false, error: 'unconfigured-project', projectSlug, flushed: 0 }
	}

	const issueIds = []
	const errors = []
	for (const entry of queued) {
		const acquired = await acquireRepoLease(
			storage,
			projectSlug,
			entry.issueId,
		)
		if (!acquired) {
			// Slot filled between drain and claim — re-enqueue and stop.
			await enqueueRepoIssue(
				storage,
				projectSlug,
				entry.issueId,
				entry.payload,
			)
			break
		}
		const result = await wakeColeForQueuedEntry(storage, project, entry)
		if (result.ok) {
			issueIds.push(result.issueId)
		} else {
			errors.push(result)
			// Slot already released inside wakeColeForQueuedEntry on failure.
		}
	}

	return {
		ok: errors.length === 0,
		flushed: issueIds.length,
		projectSlug,
		agentId: COLE_AGENT_ID,
		issueIds,
		...(errors.length > 0 ? { errors } : {}),
	}
}