Pre-launch — Gluecron is in final validation. Public signups and git hosting for non-owner users open after launch review.
CodeIssuesPull RequestsActionsSecurityInsightsSettings
✨ AI
More
Blame · Line-by-line history

autopilot.ts

Each line is annotated with the commit that last touched it. Click any SHA to jump to that commit and see the surrounding change.

autopilot.tsBlame1637 lines · 3 contributors
2b821b7Claude1/**
2 * Autopilot — self-sufficiency loop.
3 *
4 * Runs existing platform-maintenance tasks (mirror sync, merge queue progress,
5 * weekly digests, advisory rescans) on an interval so the host runs itself
6 * without an external cron. All sub-tasks are injected so tests can stub them
7 * without touching the DB; the default task set wires real helpers from the
8 * locked libs. Nothing here throws — every sub-task and the outer tick are
9 * try/caught so a single failure never blocks the others.
10 */
11
2b9055eClaude12import { and, eq, gte, sql } from "drizzle-orm";
2b821b7Claude13import { db } from "../db";
2b9055eClaude14import {
15 mergeQueueEntries,
16 prComments,
17 pullRequests,
18 repoDependencies,
19 repositories,
20 users,
21} from "../db/schema";
2b821b7Claude22import { syncAllDue } from "./mirrors";
23import { peekHead } from "./merge-queue";
24import { sendDigestsToAll } from "./email-digest";
25import { scanRepositoryForAlerts } from "./advisories";
a76d984Claude26import { releaseExpiredWaitTimers } from "./environments";
665c8bfClaude27import { runScheduledWorkflowsTick } from "./scheduled-workflows";
2b9055eClaude28import {
29 evaluateAutoMerge,
30 recordAutoMergeAttempt,
31 type AutoMergeContext,
32 type AutoMergeDecision,
33} from "./auto-merge";
34import { matchProtection } from "./branch-protection";
35import { performMerge, type PerformMergeResult } from "./pr-merge";
36import { audit } from "./notify";
37import { runAiBuildTaskOnce } from "./ai-build-tasks";
46d6165Claude38import {
39 sendSleepModeDigestForUser,
40 SLEEP_MODE_USER_CAP_PER_TICK,
41 SLEEP_MODE_COOLDOWN_HOURS,
42} from "./sleep-mode";
534f04aClaude43import {
44 runStalePrSweepOnce,
45 runStaleIssueSweepOnce,
46} from "./stale-sweep";
47import { computePrRiskForPullRequest } from "./pr-risk";
48import { prRiskScores } from "../db/schema";
c63b860Claude49import { purgeScheduledAccounts } from "./account-deletion";
cd4f63bTest User50import { purgeExpiredPlaygroundAccounts } from "./playground";
b1be050CC LABS App51import {
52 runSyntheticChecks,
53 persistChecks,
54 latestStatusByCheck,
55 type SyntheticCheckResult,
56} from "./synthetic-monitor";
e9aa4d8Claude57import { aiProactiveMonitorTick } from "./ai-proactive-monitor";
9336f45Claude58import { runCiHealerTick } from "./ai-ci-healer";
a686079Claude59import {
60 runDailyStandupTaskOnce,
61 runWeeklyStandupTaskOnce,
62} from "./ai-standup";
950ef90Claude63import { runSpecToPrTaskOnce } from "./autopilot-spec-to-pr";
45f3b73Claude64import { runMigrationWatcherTaskOnce } from "./migration-assistant";
3c03977Claude65import { sweepStale as sweepStalePrLive } from "./pr-live";
424eb72Claude66import { runAutoReleaseNotesTaskOnce } from "./ai-release-notes";
1d4ff60Claude67import { runPrTestGeneratorTaskOnce } from "./autopilot-pr-test-generator";
79ed944Claude68import { runAdvancementScan } from "./advancement-scanner";
69import { expireOldSandboxes } from "./pr-sandbox";
9b3a183Claude70import { expireIdleEnvs } from "./dev-env";
44ed968Claude71import { expireOldPreviews } from "./branch-previews";
a7460bfClaude72import { getBotUserIdOrFallback } from "./bot-user";
f65f600Claude73import { runOnboardingDripTaskOnce } from "./onboarding-drip";
f5ad215Claude74import { runDepUpdateSweepOnce } from "./dep-updater-sweep";
cc34156Claude75import { sendSmartDigestsToAll } from "./smart-digest";
2b821b7Claude76
77export interface AutopilotTaskResult {
78 name: string;
79 ok: boolean;
80 durationMs: number;
81 error?: string;
82}
83
84export interface AutopilotTickResult {
85 startedAt: string;
86 finishedAt: string;
87 tasks: AutopilotTaskResult[];
88}
89
90export interface AutopilotTask {
91 name: string;
92 run: () => Promise<void>;
93}
94
95export interface StartAutopilotOpts {
96 intervalMs?: number;
97 now?: () => number;
98 tasks?: AutopilotTask[];
99}
100
101export interface RunTickOpts {
102 tasks?: AutopilotTask[];
103 now?: () => number;
104}
105
106const DEFAULT_INTERVAL_MS = 5 * 60 * 1000;
107const ADVISORY_RESCAN_BATCH = 5;
2b9055eClaude108/** K3 — recency window for auto-merge candidate selection. */
109const AUTO_MERGE_LOOKBACK_HOURS = 24;
110/** K3 — hard cap on PRs evaluated per tick (runaway protection). */
111const AUTO_MERGE_MAX_PER_TICK = 50;
112/** K3 — stable marker for the auto-merge audit comment. */
113const AUTO_MERGE_COMMENT_MARKER = "<!-- gluecron:auto-merge:v1 -->";
534f04aClaude114/** M3 — hard cap on PRs scored per tick (runaway protection). */
115const PR_RISK_RESCORE_MAX_PER_TICK = 20;
116/** M3 — recency window for the pr-risk-rescore sweep. */
117const PR_RISK_RESCORE_LOOKBACK_HOURS = 1;
e9aa4d8Claude118/** Proactive monitor cadence — Claude scans platform telemetry hourly. */
119const PROACTIVE_MONITOR_INTERVAL_MS = 60 * 60 * 1000;
120let _lastProactiveMonitorAt = 0;
950ef90Claude121/** Spec-to-PR cadence — autopilot scans `.gluecron/specs/*.md` every 2 minutes. */
122const SPEC_TO_PR_INTERVAL_MS = 2 * 60 * 1000;
123let _lastSpecToPrAt = 0;
45f3b73Claude124/**
125 * Migration watcher cadence. The lookup is cheap (registry calls per
126 * declared dep) but we still throttle to every 6 hours so we don't hammer
127 * npm and don't propose more than ~one PR per repo per day.
128 */
129const MIGRATION_WATCHER_INTERVAL_MS = 6 * 60 * 60 * 1000;
130let _lastMigrationWatcherAt = 0;
424eb72Claude131/**
132 * Auto-release-notes cadence. Cheap once tags are rare; we still throttle
133 * to every 10 minutes so freshly-pushed tags (whose release row was just
134 * created by `POST /:owner/:repo/releases`) get notes within ~one tick
135 * without us scanning the table every 5 minutes.
136 */
137const AUTO_RELEASE_NOTES_INTERVAL_MS = 10 * 60 * 1000;
138let _lastAutoReleaseNotesAt = 0;
1d4ff60Claude139/**
140 * PR test generator cadence. Cheap when opted-in repos are quiet; the
141 * task itself short-circuits via the per-PR `ai:added-tests` marker so
142 * we never re-process the same PR. 5-minute cadence aligns with the
143 * "freshly opened PR" window the task uses for candidate selection.
144 */
145const PR_TEST_GENERATOR_INTERVAL_MS = 5 * 60 * 1000;
146let _lastPrTestGeneratorAt = 0;
79ed944Claude147/**
148 * PR sandbox cleanup cadence (migration 0067). Runs every 30 minutes —
149 * sandboxes default to a 4h TTL so finer-grained cleanup isn't needed,
150 * and skipping most of the 5-min outer ticks keeps the loop cheap.
151 */
152const PR_SANDBOX_CLEANUP_INTERVAL_MS = 30 * 60 * 1000;
153let _lastPrSandboxCleanupAt = 0;
9b3a183Claude154/**
155 * Dev-env idle sweep cadence (migration 0072). Runs on every autopilot
156 * tick (5 minutes) because per-env idle timeouts can be as short as a
157 * few minutes — we don't want a 30-min cleanup window stranding a user
158 * who set idle_minutes=5. The lib itself is a single SQL UPDATE so this
159 * is cheap regardless of repo count.
160 */
161const DEV_ENV_IDLE_SWEEP_INTERVAL_MS = 5 * 60 * 1000;
162let _lastDevEnvIdleSweepAt = 0;
f5ad215Claude163/**
164 * AI dependency auto-updater cadence (migration 0077). Once per day is
165 * plenty — the task itself caps at 10 repos and 2 candidates per repo,
166 * keeping npm registry traffic low. Skips when DEP_UPDATER_ENABLED env
167 * flag is not set to "1".
168 */
169const DEP_UPDATE_SWEEP_INTERVAL_MS = 24 * 60 * 60 * 1000;
170let _lastDepUpdateSweepAt = 0;
79ed944Claude171/**
172 * Advancement scanner cadence. Designed to run weekly on Mondays at
173 * 08:00 UTC. The task itself is the cheap gate (checks both day-of-week
174 * and minimum interval since last run) so we don't bake any cron-style
175 * triggers into the autopilot loop. Each tick (5min) probes the gate;
176 * the gate is satisfied at most once per 6 days, keeping the cadence
177 * effectively weekly even if Monday-08:00 happens to be missed.
178 */
179const ADVANCEMENT_SCAN_MIN_INTERVAL_MS = 6 * 24 * 60 * 60 * 1000;
180let _lastAdvancementScanAt = 0;
181/** Hour of day (UTC) the advancement scan prefers. Configurable via env. */
182function advancementScanHourUtc(): number {
183 const raw = process.env.ADVANCEMENT_SCAN_HOUR_UTC;
184 if (!raw) return 8;
185 const n = Number(raw);
186 if (!Number.isFinite(n) || n < 0 || n > 23) return 8;
187 return Math.floor(n);
188}
189/** Day of week (0=Sun..6=Sat, UTC) the advancement scan prefers. */
190function advancementScanDayOfWeek(): number {
191 const raw = process.env.ADVANCEMENT_SCAN_DOW_UTC;
192 if (!raw) return 1; // Monday
193 const n = Number(raw);
194 if (!Number.isFinite(n) || n < 0 || n > 6) return 1;
195 return Math.floor(n);
196}
2b821b7Claude197
cc34156Claude198/**
199 * Smart-digest cadence. Designed to run once per day at 07:00 UTC.
200 * The task itself checks `lastSmartDigestSentAt` per user (20h cooldown),
201 * so even if the outer loop fires multiple times near 07:00, only one
202 * digest is ever sent per day per user.
203 */
204const SMART_DIGEST_INTERVAL_MS = 22 * 60 * 60 * 1000; // 22h between outer checks
205let _lastSmartDigestAt = 0;
206
2b821b7Claude207/**
208 * Default task set. Each task is a thin wrapper around an existing locked
209 * helper — no gate/merge logic is duplicated here.
210 */
211export function defaultTasks(): AutopilotTask[] {
212 return [
213 {
214 name: "mirror-sync",
215 run: async () => {
216 await syncAllDue();
217 },
218 },
219 {
220 name: "merge-queue",
221 run: async () => {
222 await processMergeQueues();
223 },
224 },
225 {
226 name: "weekly-digest",
227 run: async () => {
228 await sendDigestsToAll();
229 },
230 },
231 {
232 name: "advisory-rescan",
233 run: async () => {
234 await rescanAdvisoriesBatch(ADVISORY_RESCAN_BATCH);
235 },
236 },
a76d984Claude237 {
238 name: "wait-timer-release",
239 run: async () => {
240 await releaseExpiredWaitTimers();
241 },
242 },
665c8bfClaude243 {
244 name: "scheduled-workflows",
245 run: async () => {
246 await runScheduledWorkflowsTick();
247 },
248 },
2b9055eClaude249 {
250 name: "auto-merge-sweep",
251 run: async () => {
252 await runAutoMergeSweep();
253 },
254 },
255 {
256 name: "ai-build-from-issues",
257 run: async () => {
258 const summary = await runAiBuildTaskOnce();
259 console.log(
260 `[autopilot] ai-build: queued=${summary.queued} skipped=${summary.skipped}`
261 );
262 },
263 },
46d6165Claude264 {
265 name: "sleep-mode-digest",
266 run: async () => {
267 const summary = await runSleepModeDigestTaskOnce();
268 console.log(
269 `[autopilot] sleep-mode-digest: sent=${summary.sent} skipped=${summary.skipped}`
270 );
271 },
272 },
534f04aClaude273 {
274 name: "stale-pr-sweep",
275 run: async () => {
276 // Two-stage gate: poke at 7d stale, close at 14d after poke
277 // (when the repo opts in via `auto_close_stale_prs`).
278 // Wrapped in try/catch so a finder crash never wedges the tick.
279 try {
280 const summary = await runStalePrSweepOnce();
281 console.log(
282 `[autopilot] stale-pr-sweep: poked=${summary.poked} closed=${summary.closed}`
283 );
284 } catch (err) {
285 console.error("[autopilot] stale-pr-sweep: threw:", err);
286 }
287 },
288 },
289 {
290 name: "stale-issue-sweep",
291 run: async () => {
292 // Mirror of stale-pr-sweep with the issue thresholds (30d/60d).
293 try {
294 const summary = await runStaleIssueSweepOnce();
295 console.log(
296 `[autopilot] stale-issue-sweep: poked=${summary.poked} closed=${summary.closed}`
297 );
298 } catch (err) {
299 console.error("[autopilot] stale-issue-sweep: threw:", err);
300 }
301 },
302 },
303 {
304 name: "pr-risk-rescore",
305 run: async () => {
306 const summary = await runPrRiskRescoreTaskOnce();
307 console.log(
308 `[autopilot] pr-risk-rescore: scored=${summary.scored} skipped=${summary.skipped}`
309 );
310 },
311 },
c63b860Claude312
313 {
314 // Block P5 — Hard-delete users whose 30-day grace period expired.
315 name: "account-purge",
316 run: async () => {
317 try {
318 const summary = await purgeScheduledAccounts({ cap: 50 });
319 console.log(
320 `[autopilot] account-purge: purged=${summary.purged} errors=${summary.errors}`
321 );
322 } catch (err) {
323 console.error("[autopilot] account-purge: threw:", err);
324 }
325 },
326 },
cd4f63bTest User327 {
328 // Block Q3 — Hard-delete anonymous playground accounts past their
329 // 24h TTL. CASCADE handles repos, sessions, issues. Per-user
330 // try/catch in the lib so one FK violation can't stall the queue.
331 name: "playground-purge",
332 run: async () => {
333 try {
334 const summary = await purgeExpiredPlaygroundAccounts({ cap: 50 });
335 console.log(
336 `[autopilot] playground-purge: purged=${summary.purged} errors=${summary.errors}`
337 );
338 } catch (err) {
339 console.error("[autopilot] playground-purge: threw:", err);
340 }
341 },
342 },
e9aa4d8Claude343 {
344 // Proactive AI monitor — hourly. Claude reads 24h of audit_log +
345 // platformDeploys + workflowRuns and opens issues on anomalies
346 // (degraded deploy times, recurring failures, suspicious audit
347 // patterns). Skips when ANTHROPIC_API_KEY is unset (the lib
348 // itself short-circuits, but the cadence gate avoids redundant
349 // work on every 5-min tick too).
350 name: "ai-proactive-monitor",
351 run: async () => {
352 if (!process.env.ANTHROPIC_API_KEY) return;
353 const now = Date.now();
354 if (now - _lastProactiveMonitorAt < PROACTIVE_MONITOR_INTERVAL_MS) {
355 return;
356 }
357 _lastProactiveMonitorAt = now;
358 try {
359 const summary = await aiProactiveMonitorTick();
360 console.log(
361 `[autopilot] ai-proactive-monitor: opened=${summary.opened} considered=${summary.considered} dedup=${summary.skippedDedupe}`
362 );
363 } catch (err) {
364 console.error("[autopilot] ai-proactive-monitor: threw:", err);
365 }
366 },
367 },
9336f45Claude368 {
369 // AI CI Healer — autonomous CI failure → root-cause → patch PR loop.
370 // Polls every tick (5 min) for failed workflow_runs that finished
371 // at least HEAL_MIN_AGE_MS ago and haven't been processed yet.
950ef90Claude372 // Skips when ANTHROPIC_API_KEY is unset.
9336f45Claude373 name: "ci-healer",
374 run: async () => {
375 if (!process.env.ANTHROPIC_API_KEY) return;
376 try {
377 const summary = await runCiHealerTick();
378 console.log(
379 `[autopilot] ci-healer: considered=${summary.considered} healed=${summary.healed} gaveUp=${summary.gaveUp} skipped=${summary.skipped}`
380 );
381 } catch (err) {
382 console.error("[autopilot] ci-healer: threw:", err);
383 }
384 },
385 },
950ef90Claude386 {
387 // Spec-to-PR autopilot — picks up `.gluecron/specs/*.md` files whose
388 // front-matter status is `ready`, asks Claude to implement the spec,
389 // opens a draft PR tagged `ai:spec-implementation`. Cadence-gated
390 // to every 2 minutes.
391 name: "spec-to-pr",
392 run: async () => {
393 if (!process.env.ANTHROPIC_API_KEY) return;
394 const now = Date.now();
395 if (now - _lastSpecToPrAt < SPEC_TO_PR_INTERVAL_MS) return;
396 _lastSpecToPrAt = now;
397 try {
398 const summary = await runSpecToPrTaskOnce();
399 console.log(
400 `[autopilot] spec-to-pr: considered=${summary.considered} dispatched=${summary.dispatched} skipped=${summary.skipped} failed=${summary.failed}`
401 );
402 } catch (err) {
403 console.error("[autopilot] spec-to-pr: threw:", err);
404 }
405 },
406 },
b1be050CC LABS App407 {
408 // BLOCK S4 — Synthetic monitor.
409 //
410 // Runs the URL-only smoke suite (see src/lib/synthetic-monitor.ts),
411 // records the outcome into `synthetic_checks`, and on a
412 // green->red transition fires a webhook to MONITOR_ALERT_WEBHOOK_URL
413 // (when configured) so the owner finds out instantly that the live
414 // site is broken. Wrapped in try/catch — the monitor must never
415 // wedge the tick.
416 name: "synthetic-monitor",
417 run: async () => {
418 try {
419 const summary = await runSyntheticMonitorTaskOnce();
420 console.log(
421 `[autopilot] synthetic-monitor: green=${summary.green} red=${summary.red} transitions=${summary.transitions}`
422 );
423 } catch (err) {
424 console.error("[autopilot] synthetic-monitor: threw:", err);
425 }
426 },
427 },
45f3b73Claude428 {
429 // Migration watcher — scans each repo's package.json for deps that
430 // are at least one major version behind and asks Claude to draft an
431 // upgrade PR. Cadence-gated to every 6 hours; the lib itself
432 // enforces a per-repo + per-{dep,version} 7-day dedupe so we never
433 // re-propose the same migration twice in a single window. Skips
434 // entirely when ANTHROPIC_API_KEY is unset OR the
435 // MIGRATION_WATCHER_ENABLED env flag is off.
436 name: "migration-watcher",
437 run: async () => {
438 if (!process.env.ANTHROPIC_API_KEY) return;
439 const now = Date.now();
440 if (now - _lastMigrationWatcherAt < MIGRATION_WATCHER_INTERVAL_MS) {
441 return;
442 }
443 _lastMigrationWatcherAt = now;
444 try {
445 const summary = await runMigrationWatcherTaskOnce();
446 console.log(
447 `[autopilot] migration-watcher: considered=${summary.considered} proposed=${summary.proposed} throttled=${summary.skippedThrottle} disabled=${summary.skippedNotEnabled} errors=${summary.errors}`
448 );
449 } catch (err) {
450 console.error("[autopilot] migration-watcher: threw:", err);
451 }
452 },
453 },
a686079Claude454 {
455 // AI Standup — daily Claude-generated team brief.
456 // Fires at the user's configured UTC hour (default 09:00). Skips
457 // entirely when ANTHROPIC_API_KEY is unset (the lib still has a
458 // deterministic fallback, but we keep this task quiet unless the
459 // operator has wired AI). Per-user dedupe via `hasStandupForToday`.
460 name: "daily-standup",
461 run: async () => {
462 if (!process.env.ANTHROPIC_API_KEY) return;
463 try {
464 const summary = await runDailyStandupTaskOnce();
465 console.log(
466 `[autopilot] daily-standup: sent=${summary.sent} skipped=${summary.skipped} errors=${summary.errors}`
467 );
468 } catch (err) {
469 console.error("[autopilot] daily-standup: threw:", err);
470 }
471 },
472 },
3c03977Claude473 {
474 // PR live co-editing — transition stale `pr_live_sessions` rows
475 // to 'idle' (>60s) and 'left' (>5m) so the presence pill on the
476 // PR detail page never claims a ghost user is still editing.
477 // Cheap pure-SQL UPDATE; runs every tick.
478 name: "pr-live-cleanup",
479 run: async () => {
480 try {
481 const summary = await sweepStalePrLive();
482 if (summary.idled > 0 || summary.left > 0) {
483 console.log(
484 `[autopilot] pr-live-cleanup: idled=${summary.idled} left=${summary.left}`
485 );
486 }
487 } catch (err) {
488 console.error("[autopilot] pr-live-cleanup: threw:", err);
489 }
490 },
491 },
a686079Claude492 {
493 // AI Standup — weekly Claude-generated team brief. Mondays only.
494 name: "weekly-standup",
495 run: async () => {
496 if (!process.env.ANTHROPIC_API_KEY) return;
497 try {
498 const summary = await runWeeklyStandupTaskOnce();
499 console.log(
500 `[autopilot] weekly-standup: sent=${summary.sent} skipped=${summary.skipped} errors=${summary.errors}`
501 );
502 } catch (err) {
503 console.error("[autopilot] weekly-standup: threw:", err);
504 }
505 },
506 },
1d4ff60Claude507 {
508 // PR test generator — when a fresh PR opens against a repo that's
509 // opted in (`autoGenerateTests=true`) and is not itself AI-generated,
510 // ask Claude to write tests for the new code and push a commit onto
511 // the same branch. Skips PRs without source-file changes; idempotent
512 // via the `ai:added-tests` marker comment. Skips entirely when
513 // ANTHROPIC_API_KEY is unset.
514 name: "pr-test-generator",
515 run: async () => {
516 if (!process.env.ANTHROPIC_API_KEY) return;
517 const now = Date.now();
518 if (now - _lastPrTestGeneratorAt < PR_TEST_GENERATOR_INTERVAL_MS) {
519 return;
520 }
521 _lastPrTestGeneratorAt = now;
522 try {
523 const summary = await runPrTestGeneratorTaskOnce();
524 if (summary.considered > 0) {
525 console.log(
526 `[autopilot] pr-test-generator: considered=${summary.considered} dispatched=${summary.dispatched} skipped=${summary.skipped} failed=${summary.failed}`
527 );
528 }
529 } catch (err) {
530 console.error("[autopilot] pr-test-generator: threw:", err);
531 }
532 },
533 },
79ed944Claude534 {
535 // Advancement scanner — weekly Claude-driven scan for "what we
536 // should ship next". Probes:
537 // 1. Newer Claude models vs the one wired in ai-client.ts
538 // 2. Stack dependencies that are at least one major behind
539 // 3. Self-improvement opportunities in the last 7d of telemetry
540 // 4. Trending dev-platform features competitors shipped
541 // Cadence-gated to Mondays 08:00 UTC (configurable via
542 // ADVANCEMENT_SCAN_HOUR_UTC + ADVANCEMENT_SCAN_DOW_UTC) and
543 // throttled to at most one scan per 6 days. Skips entirely when
544 // ANTHROPIC_API_KEY is unset — the offline probes are still
545 // useful but most of the value comes from the Claude calls.
546 // The lib itself enforces per-finding (sha256 of title) dedupe
547 // for 30 days so repeated runs never re-file the same advancement.
548 name: "advancement-scanner",
549 run: async () => {
550 if (!process.env.ANTHROPIC_API_KEY) return;
551 if (process.env.ADVANCEMENT_SCAN_DISABLED === "1") return;
552 const now = new Date();
553 const nowMs = now.getTime();
554 const targetHour = advancementScanHourUtc();
555 const targetDow = advancementScanDayOfWeek();
556 const dueByCadence =
557 nowMs - _lastAdvancementScanAt >= ADVANCEMENT_SCAN_MIN_INTERVAL_MS;
558 const dueByClock =
559 now.getUTCDay() === targetDow &&
560 now.getUTCHours() === targetHour;
561 if (!dueByCadence || !dueByClock) return;
562 _lastAdvancementScanAt = nowMs;
563 try {
564 const summary = await runAdvancementScan();
565 console.log(
566 `[autopilot] advancement-scanner: findings=${summary.findings.length} issues=${summary.openedIssues} prs=${summary.openedPrs} dedup=${summary.skippedDedupe} errors=${summary.errors}`
567 );
568 } catch (err) {
569 console.error("[autopilot] advancement-scanner: threw:", err);
570 }
571 },
572 },
424eb72Claude573 {
574 // Auto-release-notes — backfills `releases.body` with Claude-generated
575 // polished changelogs for any semver-tagged release whose body is
576 // empty / too short. Cadence-gated to every 10 minutes; the lib
577 // itself caps the per-tick batch and is no-op when no candidates
578 // exist. Falls back to a deterministic bucketed summary when
579 // ANTHROPIC_API_KEY is unset, so the body still ends up populated.
580 name: "auto-release-notes",
581 run: async () => {
582 const now = Date.now();
583 if (now - _lastAutoReleaseNotesAt < AUTO_RELEASE_NOTES_INTERVAL_MS) {
584 return;
585 }
586 _lastAutoReleaseNotesAt = now;
587 try {
588 const summary = await runAutoReleaseNotesTaskOnce();
589 if (summary.considered > 0 || summary.filled > 0) {
590 console.log(
591 `[autopilot] auto-release-notes: considered=${summary.considered} filled=${summary.filled} skipped=${summary.skipped} errors=${summary.errors}`
592 );
593 }
594 } catch (err) {
595 console.error("[autopilot] auto-release-notes: threw:", err);
596 }
597 },
598 },
79ed944Claude599 {
600 // PR sandbox cleanup — migration 0067. Tears down every PR sandbox
601 // whose `expires_at` has passed. Cadence-gated to every 30 minutes
602 // so the every-5-min outer loop doesn't pay the cost on most ticks.
603 // The lib itself is a pure SQL UPDATE — safe + cheap to call even
604 // when there's nothing to do.
605 name: "pr-sandbox-cleanup",
606 run: async () => {
607 const now = Date.now();
608 if (now - _lastPrSandboxCleanupAt < PR_SANDBOX_CLEANUP_INTERVAL_MS) {
609 return;
610 }
611 _lastPrSandboxCleanupAt = now;
612 try {
613 const expired = await expireOldSandboxes();
614 if (expired > 0) {
615 console.log(`[autopilot] pr-sandbox-cleanup: expired=${expired}`);
616 }
617 } catch (err) {
618 console.error("[autopilot] pr-sandbox-cleanup: threw:", err);
619 }
620 },
621 },
9b3a183Claude622 {
623 // Dev-env idle sweep — migration 0072. Stops every dev env whose
624 // `last_active_at + idle_minutes` has slipped into the past. Runs
625 // on every tick (cadence-gated to 5min so a faster outer loop
626 // doesn't hammer it). The lib itself is fail-soft + per-row idle
627 // aware, so this is cheap + correct even when half the rows are
628 // already 'stopped'.
629 name: "dev-env-idle-sweep",
630 run: async () => {
631 const now = Date.now();
632 if (now - _lastDevEnvIdleSweepAt < DEV_ENV_IDLE_SWEEP_INTERVAL_MS) {
633 return;
634 }
635 _lastDevEnvIdleSweepAt = now;
636 try {
637 const stopped = await expireIdleEnvs();
638 if (stopped > 0) {
639 console.log(`[autopilot] dev-env-idle-sweep: stopped=${stopped}`);
640 }
641 } catch (err) {
642 console.error("[autopilot] dev-env-idle-sweep: threw:", err);
643 }
644 },
645 },
44ed968Claude646 {
647 // Preview expiry — migration 0062. Flips every branch-preview row
648 // whose `expires_at` is in the past to status='expired'. Runs on
649 // every tick; the lib itself is a cheap SQL UPDATE and is a no-op
650 // when there's nothing to expire, so this never adds meaningful
651 // overhead.
652 name: "preview-expiry",
653 run: async () => {
654 try {
655 const expired = await expireOldPreviews();
656 if (expired > 0) {
657 console.log(`[autopilot] preview-expiry: expired=${expired}`);
658 }
659 } catch (err) {
660 console.error("[autopilot] preview-expiry: threw:", err);
661 }
662 },
663 },
f65f600Claude664 {
665 // Onboarding drip — sends pending T+1d and T+3d drip emails to new
666 // users. The T+0 "welcome" email is sent immediately at registration
667 // via src/routes/auth.tsx. This task handles the delayed emails.
668 // Idempotent via the per-user `onboarding_emails_sent` jsonb column
b12f20dClaude669 // (migration 0081). Silently skips when email is not configured.
f65f600Claude670 name: "onboarding-drip",
671 run: async () => {
672 try {
673 const summary = await runOnboardingDripTaskOnce();
674 if (summary.sent > 0 || summary.errors > 0) {
675 console.log(
676 `[autopilot] onboarding-drip: sent=${summary.sent} skipped=${summary.skipped} errors=${summary.errors}`
677 );
678 }
679 } catch (err) {
680 console.error("[autopilot] onboarding-drip: threw:", err);
681 }
682 },
683 },
f5ad215Claude684 {
685 // AI dependency auto-updater (migration 0077). Once per day: scans
686 // up to 10 repos with depUpdaterEnabled=true, checks for patch/minor
687 // npm updates, applies them, runs GateTest, and either auto-merges
688 // (green) or opens a PR with an AI migration guide (red). Skips
689 // when DEP_UPDATER_ENABLED env flag is not set to "1".
690 name: "dep-update-sweep",
691 run: async () => {
692 if (process.env.DEP_UPDATER_ENABLED !== "1") return;
693 const now = Date.now();
694 if (now - _lastDepUpdateSweepAt < DEP_UPDATE_SWEEP_INTERVAL_MS) {
695 return;
696 }
697 _lastDepUpdateSweepAt = now;
698 try {
699 const summary = await runDepUpdateSweepOnce();
700 console.log(
701 `[autopilot] dep-update-sweep: repos=${summary.repos} runs=${summary.runs} merged=${summary.merged} prs=${summary.prs} skipped=${summary.skipped} errors=${summary.errors}`
702 );
703 } catch (err) {
704 console.error("[autopilot] dep-update-sweep: threw:", err);
b271465Claude705 }
706 },
707 },
708 {
cc34156Claude709 // Smart morning digest — AI-curated daily developer queue.
710 // Fires once per day at 07:00 UTC (or whenever the 22h outer gate
711 // next opens after 07:00). Per-user 20h cooldown in sendSmartDigest
712 // ensures no user receives more than one digest per day even if the
713 // outer gate fires multiple times. Requires ANTHROPIC_API_KEY but
714 // degrades gracefully to a rule-based prioritisation when unset.
715 name: "smart-digest",
716 run: async () => {
717 const now = Date.now();
718 const nowDate = new Date(now);
719 // Only fire at hour=7 UTC or if last run was >22h ago (catch-up)
720 const isDigestHour = nowDate.getUTCHours() === 7;
721 const pastCooldown = now - _lastSmartDigestAt >= SMART_DIGEST_INTERVAL_MS;
722 if (!isDigestHour && !pastCooldown) return;
723 _lastSmartDigestAt = now;
724 try {
725 await sendSmartDigestsToAll();
726 console.log("[autopilot] smart-digest: completed");
727 } catch (err) {
728 console.error("[autopilot] smart-digest: threw:", err);
f5ad215Claude729 }
730 },
731 },
2b821b7Claude732 ];
733}
734
b1be050CC LABS App735// ---------------------------------------------------------------------------
736// BLOCK S4 — synthetic-monitor task
737// ---------------------------------------------------------------------------
738
739export interface SyntheticMonitorTaskDeps {
740 /** Override the suite runner (DI for tests). */
741 runChecks?: () => Promise<SyntheticCheckResult[]>;
742 /** Override the persistence step (DI for tests). */
743 persist?: (results: SyntheticCheckResult[]) => Promise<void>;
744 /** Override the previous-state loader (DI for tests). */
745 loadPrevious?: () => Promise<Record<string, SyntheticCheckResult>>;
746 /** Override the webhook poster (DI for tests). */
747 postAlert?: (url: string, payload: unknown) => Promise<void>;
748 /** Override the alert-webhook URL lookup (defaults to env). */
749 alertUrl?: () => string;
750}
751
752export interface SyntheticMonitorTaskSummary {
753 green: number;
754 red: number;
755 yellow: number;
756 transitions: number;
757}
758
759async function defaultPostAlert(
760 url: string,
761 payload: unknown
762): Promise<void> {
763 try {
764 await fetch(url, {
765 method: "POST",
766 headers: { "Content-Type": "application/json" },
767 body: JSON.stringify(payload),
768 });
769 } catch (err) {
770 console.error("[autopilot] synthetic-monitor: alert webhook failed:", err);
771 }
772}
773
774/**
775 * One iteration of the synthetic-monitor task. Runs the checks, persists
776 * them, compares against the prior state, and fires a webhook on each
777 * green->red transition (red->red repeats stay quiet so we don't spam
778 * the channel). Never throws.
779 */
780export async function runSyntheticMonitorTaskOnce(
781 deps: SyntheticMonitorTaskDeps = {}
782): Promise<SyntheticMonitorTaskSummary> {
783 const runChecks = deps.runChecks ?? (() => runSyntheticChecks());
784 const persist = deps.persist ?? persistChecks;
785 const loadPrevious =
786 deps.loadPrevious ??
787 (async () => {
788 const latest = await latestStatusByCheck();
789 // Strip the `checkedAt` from the result shape so the diff loop
790 // compares the canonical SyntheticCheckResult fields only.
791 const out: Record<string, SyntheticCheckResult> = {};
792 for (const [k, v] of Object.entries(latest)) {
793 const { checkedAt: _unused, ...rest } = v;
794 void _unused;
795 out[k] = rest;
796 }
797 return out;
798 });
799 const postAlert = deps.postAlert ?? defaultPostAlert;
800 const alertUrl =
801 deps.alertUrl ?? (() => process.env.MONITOR_ALERT_WEBHOOK_URL || "");
802
803 let previous: Record<string, SyntheticCheckResult> = {};
804 try {
805 previous = await loadPrevious();
806 } catch (err) {
807 console.error(
808 "[autopilot] synthetic-monitor: loadPrevious threw:",
809 err
810 );
811 previous = {};
812 }
813
814 const results = await runChecks();
815 await persist(results);
816
817 let green = 0;
818 let red = 0;
819 let yellow = 0;
820 let transitions = 0;
821 const url = alertUrl();
822
823 for (const r of results) {
824 if (r.status === "green") green += 1;
825 else if (r.status === "red") red += 1;
826 else yellow += 1;
827
828 const prior = previous[r.name];
829 // green->red transition: prior was green (or absent and current is red
830 // after a green is also a transition — but absent-before is treated as
831 // green to avoid spamming on a fresh DB). We only alert on the
832 // green->red edge so red->red doesn't re-fire.
833 const priorWasGreen = !prior || prior.status === "green";
834 if (priorWasGreen && r.status === "red") {
835 transitions += 1;
836 if (url) {
837 await postAlert(url, {
838 check: r.name,
839 status: r.status,
840 statusCode: r.statusCode ?? null,
841 durationMs: r.durationMs,
842 error: r.error ?? null,
843 checkedAt: new Date().toISOString(),
844 });
845 }
846 }
847 }
848
849 return { green, red, yellow, transitions };
850}
851
46d6165Claude852// ---------------------------------------------------------------------------
853// L1 — sleep-mode-digest
854// ---------------------------------------------------------------------------
855
856export interface SleepModeDigestCandidate {
857 userId: string;
858 digestHourUtc: number;
e1fc7dbClaude859 /** Independent cooldown anchor for the sleep-mode digest (migration 0077). */
860 lastSleepDigestSentAt: Date | null;
46d6165Claude861}
862
863export interface SleepModeDigestTaskDeps {
864 /** Override the candidate finder. */
865 findCandidates?: (cap: number) => Promise<SleepModeDigestCandidate[]>;
866 /** Override the send-one-user helper (DI for tests). */
867 sendOne?: (userId: string) => Promise<{ ok: boolean; reason?: string }>;
868 /** Override the wall clock (DI for tests). */
869 now?: () => Date;
870 /** Override the per-tick cap. */
871 cap?: number;
872 /** Override the cooldown hours. */
873 cooldownHours?: number;
874}
875
876export interface SleepModeDigestTaskSummary {
877 sent: number;
878 skipped: number;
879}
880
881/**
882 * Default candidate-finder. Returns enabled users whose
e1fc7dbClaude883 * `lastSleepDigestSentAt` is older than the cooldown OR null. The hour-match
46d6165Claude884 * filter is applied in JS by `runSleepModeDigestTaskOnce` so it stays
885 * timezone-independent of any SQL `extract(hour ...)` behaviour.
e1fc7dbClaude886 * Uses the dedicated `last_sleep_digest_sent_at` column (migration 0077) so
887 * the cooldown is independent of the weekly digest timer.
46d6165Claude888 */
889async function defaultFindSleepModeCandidates(
890 cap: number
891): Promise<SleepModeDigestCandidate[]> {
892 try {
893 const rows = await db
894 .select({
895 userId: users.id,
896 digestHourUtc: users.sleepModeDigestHourUtc,
e1fc7dbClaude897 lastSleepDigestSentAt: users.lastSleepDigestSentAt,
46d6165Claude898 })
899 .from(users)
900 .where(eq(users.sleepModeEnabled, true))
901 .limit(cap);
902 return rows.map((r) => ({
903 userId: r.userId,
904 digestHourUtc: r.digestHourUtc,
e1fc7dbClaude905 lastSleepDigestSentAt: r.lastSleepDigestSentAt,
46d6165Claude906 }));
907 } catch (err) {
908 console.error("[autopilot] sleep-mode-digest: candidate query failed:", err);
909 return [];
910 }
911}
912
913/**
914 * One iteration of the sleep-mode-digest task. Never throws.
915 *
916 * Per-user filters (applied in JS so we can DI a clock):
e1fc7dbClaude917 * 1. `lastSleepDigestSentAt` is null OR older than cooldown (23h).
46d6165Claude918 * 2. `now.getUTCHours() === digestHourUtc` — fires once at the user's
919 * configured local UTC hour.
920 *
921 * Caps at `SLEEP_MODE_USER_CAP_PER_TICK` (100) users per tick.
922 */
923export async function runSleepModeDigestTaskOnce(
924 deps: SleepModeDigestTaskDeps = {}
925): Promise<SleepModeDigestTaskSummary> {
926 const findCandidates =
927 deps.findCandidates ?? defaultFindSleepModeCandidates;
928 const sendOne = deps.sendOne ?? sendSleepModeDigestForUser;
929 const now = deps.now ?? (() => new Date());
930 const cap = deps.cap ?? SLEEP_MODE_USER_CAP_PER_TICK;
931 const cooldownHours = deps.cooldownHours ?? SLEEP_MODE_COOLDOWN_HOURS;
932
933 let candidates: SleepModeDigestCandidate[] = [];
934 try {
935 candidates = await findCandidates(cap);
936 } catch (err) {
937 console.error("[autopilot] sleep-mode-digest: findCandidates threw:", err);
938 return { sent: 0, skipped: 0 };
939 }
940
941 const nowDate = now();
942 const currentHour = nowDate.getUTCHours();
943 const cooldownMs = cooldownHours * 60 * 60 * 1000;
944
945 let sent = 0;
946 let skipped = 0;
947
948 for (const cand of candidates) {
949 try {
950 // Hour-match: must equal the user's configured UTC delivery hour.
951 if (cand.digestHourUtc !== currentHour) {
952 skipped += 1;
953 continue;
954 }
955 // Cooldown: skip if we sent within the last cooldown window.
956 if (
e1fc7dbClaude957 cand.lastSleepDigestSentAt &&
958 nowDate.getTime() - new Date(cand.lastSleepDigestSentAt).getTime() <
46d6165Claude959 cooldownMs
960 ) {
961 skipped += 1;
962 continue;
963 }
964 const result = await sendOne(cand.userId);
965 if (result.ok) sent += 1;
966 else skipped += 1;
967 } catch (err) {
968 skipped += 1;
969 console.error(
970 `[autopilot] sleep-mode-digest: per-user failure for user=${cand.userId}:`,
971 err
972 );
973 }
974 }
975
976 return { sent, skipped };
977}
978
534f04aClaude979// ---------------------------------------------------------------------------
980// M3 — pr-risk-rescore
981// ---------------------------------------------------------------------------
982
983export interface PrRiskRescoreCandidate {
984 pullRequestId: string;
985 headBranch: string;
986 updatedAt: Date;
987}
988
989export interface PrRiskRescoreTaskDeps {
990 /** Override candidate finder for tests. */
991 findCandidates?: (
992 lookbackHours: number,
993 cap: number
994 ) => Promise<PrRiskRescoreCandidate[]>;
995 /** Override score computation for tests. */
996 scoreOne?: (prId: string) => Promise<{ ok: boolean }>;
997 /** Override per-tick cap. */
998 cap?: number;
999 /** Override lookback. */
1000 lookbackHours?: number;
1001}
1002
1003export interface PrRiskRescoreTaskSummary {
1004 scored: number;
1005 skipped: number;
1006}
1007
1008/**
1009 * Default candidate-finder. Returns open, non-draft PRs from non-archived
1010 * repos whose `updated_at` falls inside the lookback window. The "scored
1011 * at all" filter is applied as a second pass via `defaultFilterNeedsScoring`.
1012 */
1013async function defaultFindPrRiskCandidates(
1014 lookbackHours: number,
1015 cap: number
1016): Promise<PrRiskRescoreCandidate[]> {
1017 const cutoff = new Date(Date.now() - lookbackHours * 60 * 60 * 1000);
1018 try {
1019 const rows = await db
1020 .select({
1021 pullRequestId: pullRequests.id,
1022 headBranch: pullRequests.headBranch,
1023 updatedAt: pullRequests.updatedAt,
1024 })
1025 .from(pullRequests)
1026 .innerJoin(
1027 repositories,
1028 eq(repositories.id, pullRequests.repositoryId)
1029 )
1030 .where(
1031 and(
1032 eq(pullRequests.state, "open"),
1033 eq(pullRequests.isDraft, false),
1034 eq(repositories.isArchived, false),
1035 gte(pullRequests.updatedAt, cutoff)
1036 )
1037 )
1038 .orderBy(sql`${pullRequests.updatedAt} DESC`)
1039 .limit(cap);
1040 return rows.map((r) => ({
1041 pullRequestId: r.pullRequestId,
1042 headBranch: r.headBranch,
1043 updatedAt: r.updatedAt,
1044 }));
1045 } catch (err) {
1046 console.error("[autopilot] pr-risk-rescore: candidate query failed:", err);
1047 return [];
1048 }
1049}
1050
1051/**
1052 * Drop candidates that already have ANY cached score row. The unique
1053 * constraint on (pull_request_id, commit_sha) handles the "score-the-
1054 * same-SHA-twice" case at persist time; this filter just keeps the work
1055 * list small enough to fit under the per-tick cap when many PRs are
1056 * being pushed concurrently.
1057 */
1058async function defaultFilterNeedsScoring(
1059 candidates: PrRiskRescoreCandidate[]
1060): Promise<PrRiskRescoreCandidate[]> {
1061 if (candidates.length === 0) return [];
1062 try {
1063 const rows = await db
1064 .select({ pullRequestId: prRiskScores.pullRequestId })
1065 .from(prRiskScores);
1066 const scoredIds = new Set(rows.map((r) => r.pullRequestId));
1067 return candidates.filter((c) => !scoredIds.has(c.pullRequestId));
1068 } catch (err) {
1069 console.error("[autopilot] pr-risk-rescore: filter query failed:", err);
1070 // Fail-open: better to score everything than silently skip.
1071 return candidates;
1072 }
1073}
1074
1075/**
1076 * One iteration of the pr-risk-rescore task. Never throws. Compute risk
1077 * for up to `cap` recently-touched open PRs that have no cached score
1078 * yet, so reviewers usually see a populated card on first visit.
1079 */
1080export async function runPrRiskRescoreTaskOnce(
1081 deps: PrRiskRescoreTaskDeps = {}
1082): Promise<PrRiskRescoreTaskSummary> {
1083 const findCandidates = deps.findCandidates ?? defaultFindPrRiskCandidates;
1084 const scoreOne =
1085 deps.scoreOne ??
1086 (async (prId: string) => {
1087 try {
1088 const result = await computePrRiskForPullRequest(prId);
1089 return { ok: result !== null };
1090 } catch {
1091 return { ok: false };
1092 }
1093 });
1094 const cap = deps.cap ?? PR_RISK_RESCORE_MAX_PER_TICK;
1095 const lookbackHours =
1096 deps.lookbackHours ?? PR_RISK_RESCORE_LOOKBACK_HOURS;
1097
1098 let candidates: PrRiskRescoreCandidate[] = [];
1099 try {
1100 candidates = await findCandidates(lookbackHours, cap);
1101 } catch (err) {
1102 console.error("[autopilot] pr-risk-rescore: findCandidates threw:", err);
1103 return { scored: 0, skipped: 0 };
1104 }
1105
1106 // Only score PRs missing a cached row. Skip filter when the caller
1107 // injected a custom finder (tests pass already-filtered lists).
1108 const needsScoring =
1109 deps.findCandidates === undefined
1110 ? await defaultFilterNeedsScoring(candidates)
1111 : candidates;
1112
1113 let scored = 0;
1114 let skipped = 0;
1115 for (const cand of needsScoring.slice(0, cap)) {
1116 try {
1117 const result = await scoreOne(cand.pullRequestId);
1118 if (result.ok) scored += 1;
1119 else skipped += 1;
1120 } catch (err) {
1121 skipped += 1;
1122 console.error(
1123 `[autopilot] pr-risk-rescore: per-PR failure for pr=${cand.pullRequestId}:`,
1124 err
1125 );
1126 }
1127 }
1128 if (needsScoring.length > cap) {
1129 skipped += needsScoring.length - cap;
1130 }
1131
1132 return { scored, skipped };
1133}
1134
2b9055eClaude1135// ---------------------------------------------------------------------------
1136// K3 — auto-merge-sweep
1137// ---------------------------------------------------------------------------
1138
1139interface SweepCandidate {
1140 prId: string;
1141 prNumber: number;
1142 prTitle: string;
1143 prBody: string | null;
1144 baseBranch: string;
1145 headBranch: string;
1146 isDraft: boolean;
1147 repositoryId: string;
1148 authorUserId: string;
1149 ownerUsername: string | null;
1150 repoName: string;
1151 state: string;
1152}
1153
1154export interface AutoMergeSweepDeps {
1155 /** Inject candidate-finder for tests. */
1156 findCandidates?: (lookbackHours: number, limit: number) => Promise<SweepCandidate[]>;
1157 /** Inject evaluator for tests. */
1158 evaluate?: (ctx: AutoMergeContext) => Promise<AutoMergeDecision>;
1159 /** Inject the merge executor for tests. */
1160 merge?: (cand: SweepCandidate) => Promise<PerformMergeResult>;
1161 /** Inject the audit-recording side-effect for tests. */
1162 recordAttempt?: (
1163 repoId: string,
1164 prId: string,
1165 decision: AutoMergeDecision
1166 ) => Promise<void>;
1167 /** Inject the audit/comment side-effects for the merged path (tests). */
1168 onMerged?: (
1169 cand: SweepCandidate,
1170 result: PerformMergeResult
1171 ) => Promise<void>;
1172 /** Inject the audit side-effect for the merge-failed path (tests). */
1173 onMergeFailed?: (cand: SweepCandidate, error: string) => Promise<void>;
1174 /** Inject the AI-key short-circuit signal for tests. */
1175 shouldShortCircuitAi?: (cand: SweepCandidate) => Promise<boolean>;
1176}
1177
1178export interface AutoMergeSweepSummary {
1179 evaluated: number;
1180 merged: number;
1181 blocked: number;
1182}
1183
1184/**
1185 * Default candidate-finder. Selects open, non-draft PRs from non-archived
1186 * repos whose `updated_at` is within the lookback window. Joins repo +
1187 * owner so the merge executor doesn't need extra round trips. Cap is
1188 * enforced at the SQL layer.
1189 */
1190async function defaultFindAutoMergeCandidates(
1191 lookbackHours: number,
1192 limit: number
1193): Promise<SweepCandidate[]> {
1194 const cutoff = new Date(Date.now() - lookbackHours * 60 * 60 * 1000);
1195 try {
1196 const rows = await db
1197 .select({
1198 prId: pullRequests.id,
1199 prNumber: pullRequests.number,
1200 prTitle: pullRequests.title,
1201 prBody: pullRequests.body,
1202 baseBranch: pullRequests.baseBranch,
1203 headBranch: pullRequests.headBranch,
1204 isDraft: pullRequests.isDraft,
1205 repositoryId: pullRequests.repositoryId,
1206 authorUserId: pullRequests.authorId,
1207 ownerUsername: users.username,
1208 repoName: repositories.name,
1209 state: pullRequests.state,
1210 })
1211 .from(pullRequests)
1212 .innerJoin(
1213 repositories,
1214 eq(repositories.id, pullRequests.repositoryId)
1215 )
1216 .leftJoin(users, eq(users.id, repositories.ownerId))
1217 .where(
1218 and(
1219 eq(pullRequests.state, "open"),
1220 eq(pullRequests.isDraft, false),
1221 eq(repositories.isArchived, false),
1222 gte(pullRequests.updatedAt, cutoff)
1223 )
1224 )
1225 .limit(limit);
1226 return rows.map((r) => ({
1227 prId: r.prId,
1228 prNumber: r.prNumber,
1229 prTitle: r.prTitle,
1230 prBody: r.prBody,
1231 baseBranch: r.baseBranch,
1232 headBranch: r.headBranch,
1233 isDraft: r.isDraft,
1234 repositoryId: r.repositoryId,
1235 authorUserId: r.authorUserId,
1236 ownerUsername: r.ownerUsername ?? null,
1237 repoName: r.repoName,
1238 state: r.state,
1239 }));
1240 } catch (err) {
1241 console.error("[autopilot] auto-merge: candidate query failed:", err);
1242 return [];
1243 }
1244}
1245
1246/**
1247 * Determine whether the matched branch_protection rule on this PR
1248 * requires AI approval but no `ANTHROPIC_API_KEY` is configured. In that
1249 * case the AI-approval check would inevitably fail downstream, so we
1250 * short-circuit to a "blocked" decision without invoking `evaluateAutoMerge`
1251 * — keeps the log readable and prevents misleading "AI review unavailable"
1252 * lines in the audit trail.
1253 */
1254async function defaultShouldShortCircuitAi(
1255 cand: SweepCandidate
1256): Promise<boolean> {
1257 if (process.env.ANTHROPIC_API_KEY) return false;
1258 try {
1259 const rule = await matchProtection(cand.repositoryId, cand.baseBranch);
1260 return !!(rule && rule.requireAiApproval);
1261 } catch {
1262 return false;
1263 }
1264}
1265
1266/**
1267 * Default success-path: post an `auto_merge.merged` audit row + a stable
1268 * marker comment on the PR so a partial-merge retry doesn't double-post.
1269 * Both are best-effort; failures are logged not thrown.
1270 */
1271async function defaultOnMerged(
1272 cand: SweepCandidate,
1273 result: PerformMergeResult
1274): Promise<void> {
1275 try {
1276 await audit({
1277 repositoryId: cand.repositoryId,
1278 action: "auto_merge.merged",
1279 targetType: "pull_request",
1280 targetId: cand.prId,
1281 metadata: {
1282 prNumber: cand.prNumber,
1283 baseBranch: cand.baseBranch,
1284 headBranch: cand.headBranch,
1285 closedIssueNumbers: result.closedIssueNumbers,
1286 resolvedFiles: result.resolvedFiles,
1287 },
1288 });
1289 } catch (err) {
1290 console.error("[autopilot] auto-merge: merged audit failed:", err);
1291 }
1292 try {
a7460bfClaude1293 const commentAuthorId = await getBotUserIdOrFallback(cand.authorUserId);
2b9055eClaude1294 await db.insert(prComments).values({
1295 pullRequestId: cand.prId,
a7460bfClaude1296 authorId: commentAuthorId,
2b9055eClaude1297 isAiReview: true,
1298 body: `${AUTO_MERGE_COMMENT_MARKER}\nAuto-merged by Gluecron autopilot — branch protection conditions satisfied.`,
1299 });
1300 } catch (err) {
1301 console.error("[autopilot] auto-merge: comment insert failed:", err);
1302 }
1303}
1304
1305/** Default failure-path: only an audit row; no comment (we may retry). */
1306async function defaultOnMergeFailed(
1307 cand: SweepCandidate,
1308 error: string
1309): Promise<void> {
1310 try {
1311 await audit({
1312 repositoryId: cand.repositoryId,
1313 action: "auto_merge.merge_failed",
1314 targetType: "pull_request",
1315 targetId: cand.prId,
1316 metadata: {
1317 prNumber: cand.prNumber,
1318 baseBranch: cand.baseBranch,
1319 headBranch: cand.headBranch,
1320 error,
1321 },
1322 });
1323 } catch (err) {
1324 console.error("[autopilot] auto-merge: merge_failed audit failed:", err);
1325 }
1326}
1327
1328/**
1329 * Execute one sweep over recently-updated open PRs. For each, evaluate
1330 * with K2's `evaluateAutoMerge`; on `merge: true`, call `performMerge` and
1331 * record the merged/merge-failed audit row + comment. Always record the
1332 * `auto_merge.evaluated` audit row via `recordAutoMergeAttempt`.
1333 *
1334 * Returns a counts summary that the autopilot prints as the tick log line.
1335 * Never throws.
1336 */
1337export async function runAutoMergeSweep(
1338 deps: AutoMergeSweepDeps = {}
1339): Promise<AutoMergeSweepSummary> {
1340 const findCandidates = deps.findCandidates ?? defaultFindAutoMergeCandidates;
1341 const evaluate =
1342 deps.evaluate ?? ((ctx) => evaluateAutoMerge(ctx, {}));
1343 const merge =
1344 deps.merge ??
1345 (async (cand) => {
1346 if (!cand.ownerUsername) {
1347 return {
1348 ok: false,
1349 error: "owner username unresolved",
1350 closedIssueNumbers: [],
1351 resolvedFiles: [],
1352 };
1353 }
1354 return performMerge({
1355 pr: {
1356 id: cand.prId,
1357 number: cand.prNumber,
1358 title: cand.prTitle,
1359 body: cand.prBody,
1360 baseBranch: cand.baseBranch,
1361 headBranch: cand.headBranch,
1362 repositoryId: cand.repositoryId,
1363 authorId: cand.authorUserId,
1364 state: cand.state as "open",
1365 isDraft: cand.isDraft,
1366 },
1367 ownerName: cand.ownerUsername,
1368 repoName: cand.repoName,
1369 actorUserId: cand.authorUserId,
1370 });
1371 });
1372 const recordAttempt = deps.recordAttempt ?? recordAutoMergeAttempt;
1373 const onMerged = deps.onMerged ?? defaultOnMerged;
1374 const onMergeFailed = deps.onMergeFailed ?? defaultOnMergeFailed;
1375 const shouldShortCircuitAi =
1376 deps.shouldShortCircuitAi ?? defaultShouldShortCircuitAi;
1377
1378 let candidates: SweepCandidate[] = [];
1379 try {
1380 candidates = await findCandidates(
1381 AUTO_MERGE_LOOKBACK_HOURS,
1382 AUTO_MERGE_MAX_PER_TICK
1383 );
1384 } catch (err) {
1385 console.error("[autopilot] auto-merge: findCandidates threw:", err);
1386 return { evaluated: 0, merged: 0, blocked: 0 };
1387 }
1388
1389 let evaluated = 0;
1390 let merged = 0;
1391 let blocked = 0;
1392
1393 for (const cand of candidates) {
1394 try {
1395 evaluated += 1;
1396
1397 // AI-key short-circuit: if the rule requires AI approval and we have
1398 // no key, treat as blocked without calling the evaluator (which would
1399 // log a misleading "AI review unavailable").
1400 let decision: AutoMergeDecision;
1401 if (await shouldShortCircuitAi(cand)) {
1402 decision = {
1403 merge: false,
1404 reason:
1405 "Branch protection requires AI approval but ANTHROPIC_API_KEY is unset.",
1406 blocking: [
1407 "ANTHROPIC_API_KEY missing; AI approval cannot be sourced.",
1408 ],
1409 };
1410 } else {
1411 decision = await evaluate({
1412 pullRequestId: cand.prId,
1413 repositoryId: cand.repositoryId,
1414 baseBranch: cand.baseBranch,
1415 isDraft: cand.isDraft,
1416 authorUserId: cand.authorUserId,
1417 });
1418 }
1419
1420 // Always record the evaluation, regardless of outcome.
1421 try {
1422 await recordAttempt(cand.repositoryId, cand.prId, decision);
1423 } catch (err) {
1424 console.error(
1425 `[autopilot] auto-merge: recordAttempt failed for pr=${cand.prId}:`,
1426 err
1427 );
1428 }
1429
1430 if (!decision.merge) {
1431 blocked += 1;
1432 continue;
1433 }
1434
1435 // Perform the actual merge.
1436 const result = await merge(cand);
1437 if (result.ok) {
1438 merged += 1;
1439 await onMerged(cand, result);
1440 } else {
1441 blocked += 1;
1442 await onMergeFailed(cand, result.error || "unknown merge error");
1443 }
1444 } catch (err) {
1445 blocked += 1;
1446 console.error(
1447 `[autopilot] auto-merge: per-PR failure for pr=${cand.prId}:`,
1448 err
1449 );
1450 }
1451 }
1452
1453 console.log(
1454 `[autopilot] auto-merge: evaluated=${evaluated} merged=${merged} blocked=${blocked}`
1455 );
1456
1457 return { evaluated, merged, blocked };
1458}
1459
2b821b7Claude1460/**
1461 * Visits each distinct (repo, base_branch) that has queued rows and logs a
1462 * stub depth line. The actual gate-running + merge happens in the pulls
1463 * route; this tick is just a heartbeat so we can wire per-queue progress
1464 * through without duplicating merge logic.
1465 */
1466async function processMergeQueues(): Promise<void> {
1467 let distinct: Array<{ repositoryId: string; baseBranch: string }> = [];
1468 try {
1469 const rows = await db
1470 .selectDistinct({
1471 repositoryId: mergeQueueEntries.repositoryId,
1472 baseBranch: mergeQueueEntries.baseBranch,
1473 })
1474 .from(mergeQueueEntries)
1475 .where(sql`${mergeQueueEntries.state} IN ('queued','running')`);
1476 distinct = rows;
1477 } catch (err) {
1478 console.error("[autopilot] merge-queue: distinct query failed:", err);
1479 return;
1480 }
1481 for (const d of distinct) {
1482 try {
1483 const head = await peekHead(d.repositoryId, d.baseBranch);
1484 if (head) {
1485 console.log(
1486 `[autopilot] merge queue depth head=${head.id.slice(0, 8)} repo=${d.repositoryId.slice(0, 8)} base=${d.baseBranch}`
1487 );
1488 }
1489 } catch (err) {
1490 console.error(
1491 `[autopilot] merge-queue: peek failed for repo=${d.repositoryId}:`,
1492 err
1493 );
1494 }
1495 }
1496}
1497
1498/**
1499 * Pick a small batch of repos that actually have dep rows and re-run
1500 * advisory scan against them. Cheap — one SELECT DISTINCT with LIMIT.
1501 */
1502async function rescanAdvisoriesBatch(limit: number): Promise<void> {
1503 let repoIds: string[] = [];
1504 try {
1505 const rows = await db
1506 .selectDistinct({ repositoryId: repoDependencies.repositoryId })
1507 .from(repoDependencies)
1508 .limit(limit);
1509 repoIds = rows.map((r) => r.repositoryId);
1510 } catch (err) {
1511 console.error("[autopilot] advisory-rescan: query failed:", err);
1512 return;
1513 }
1514 for (const id of repoIds) {
1515 try {
1516 await scanRepositoryForAlerts(id);
1517 } catch (err) {
1518 console.error(
1519 `[autopilot] advisory-rescan: scan failed for repo=${id}:`,
1520 err
1521 );
1522 }
1523 }
1524}
1525
1526/** Resolve the tick interval from env → opts → default. */
1527function resolveIntervalMs(optsMs?: number): number {
1528 if (typeof optsMs === "number" && optsMs > 0) return optsMs;
1529 const raw = process.env.AUTOPILOT_INTERVAL_MS;
1530 if (raw) {
1531 const parsed = Number(raw);
1532 if (Number.isFinite(parsed) && parsed > 0) return parsed;
1533 }
1534 return DEFAULT_INTERVAL_MS;
1535}
1536
1537/**
1538 * Start the recurring autopilot loop. No-op when AUTOPILOT_DISABLED=1.
1539 * The first tick fires after `intervalMs`, not immediately, to keep boot
1540 * fast. Returns a `stop()` that clears the interval.
1541 */
1542export function startAutopilot(opts?: StartAutopilotOpts): { stop: () => void } {
1543 if (process.env.AUTOPILOT_DISABLED === "1") {
1544 return { stop: () => {} };
1545 }
1546 const intervalMs = resolveIntervalMs(opts?.intervalMs);
1547 const tasks = opts?.tasks ?? defaultTasks();
1548 let running = false;
1549 const handle = setInterval(() => {
1550 if (running) return;
1551 running = true;
1552 void runAutopilotTick({ tasks, now: opts?.now })
1553 .catch(() => {
1554 // runAutopilotTick already never throws, but belt-and-braces.
1555 })
1556 .finally(() => {
1557 running = false;
1558 });
1559 }, intervalMs);
1560 return {
1561 stop: () => clearInterval(handle),
1562 };
1563}
1564
8e9f1d9Claude1565/** Last tick snapshot for observability. Module-level, swap-on-complete. */
1566let lastTick: AutopilotTickResult | null = null;
1567let tickCount = 0;
1568
1569/** Return the most recent completed tick, or null if autopilot hasn't run yet. */
1570export function getLastTick(): AutopilotTickResult | null {
1571 return lastTick;
1572}
1573
1574/** Return the total number of completed ticks in this process. */
1575export function getTickCount(): number {
1576 return tickCount;
1577}
1578
2b821b7Claude1579/**
1580 * Run one tick: invokes every sub-task with its own try/catch, records a
1581 * per-task result, and emits a single summary line. Never throws.
1582 */
1583export async function runAutopilotTick(
1584 opts?: RunTickOpts
1585): Promise<AutopilotTickResult> {
1586 const now = opts?.now ?? Date.now;
1587 const tasks = opts?.tasks ?? defaultTasks();
1588 const startedAt = new Date(now()).toISOString();
1589 const results: AutopilotTaskResult[] = [];
1590 for (const t of tasks) {
1591 const t0 = now();
1592 try {
1593 await t.run();
1594 results.push({ name: t.name, ok: true, durationMs: now() - t0 });
1595 } catch (err) {
1596 const message =
1597 err instanceof Error ? err.message : String(err ?? "unknown error");
1598 console.error(`[autopilot] ${t.name}: ${message}`);
1599 results.push({
1600 name: t.name,
1601 ok: false,
1602 durationMs: now() - t0,
1603 error: message,
1604 });
1605 }
1606 }
1607 const finishedAt = new Date(now()).toISOString();
1608 const totalMs = results.reduce((a, r) => a + r.durationMs, 0);
1609 const okCount = results.filter((r) => r.ok).length;
1610 console.log(
1611 `[autopilot] tick ok tasks=${okCount}/${results.length} ms=${totalMs}`
1612 );
8e9f1d9Claude1613 const result: AutopilotTickResult = { startedAt, finishedAt, tasks: results };
1614 lastTick = result;
1615 tickCount += 1;
1616 return result;
2b821b7Claude1617}
1618
1619/** Exposed for unit tests. */
1620export const __test = {
1621 resolveIntervalMs,
1622 processMergeQueues,
1623 rescanAdvisoriesBatch,
1624 DEFAULT_INTERVAL_MS,
1625 ADVISORY_RESCAN_BATCH,
2b9055eClaude1626 AUTO_MERGE_LOOKBACK_HOURS,
1627 AUTO_MERGE_MAX_PER_TICK,
1628 AUTO_MERGE_COMMENT_MARKER,
534f04aClaude1629 PR_RISK_RESCORE_MAX_PER_TICK,
1630 PR_RISK_RESCORE_LOOKBACK_HOURS,
2b9055eClaude1631 defaultFindAutoMergeCandidates,
1632 defaultOnMerged,
1633 defaultOnMergeFailed,
1634 defaultShouldShortCircuitAi,
534f04aClaude1635 defaultFindPrRiskCandidates,
1636 defaultFilterNeedsScoring,
2b821b7Claude1637};