Skip to content
File

Blob: src/worker/durable/project-do/reconciliation/watchdog.ts

typescript99 lines
1import { eq } from "drizzle-orm";
2 
3import { type ProjectId, RunId } from "@/contracts";
4import { D1SyncStatus, type ProjectRunTerminalStatus, type RunMetaState, expectTrusted } from "@/worker/contracts";
5import * as projectSchema from "@/worker/db/durable/schema/project-do";
6 
7import { HEARTBEAT_STALE_AFTER_MS } from "../constants";
8import { getProjectState, getRunRow, getRunRowByRunId } from "../repo";
9import { setRunDoTerminal } from "../run-do-sync";
10import { getHeartbeatAt, setHeartbeatAt } from "../sidecar-state";
11import { attemptSandboxCleanup, clearSandboxCleanupRetry, seedSandboxCleanupRetry } from "./sandbox-cleanup";
12import { getRunStub } from "./shared";
13import { isTerminalStatus, nextTerminalD1SyncStatus, promoteNextPendingRun } from "../transitions/shared";
14import type { ProjectDoContext } from "../types";
15 
16export const reconcileActiveRunWatchdog = async (
17 context: ProjectDoContext,
18 projectId: ProjectId,
19): Promise<RunId | null> => {
20 const stateRow = await getProjectState(context, projectId);
21 if (!stateRow?.activeRunId) {
22 return null;
23 }
24 
25 const runId = expectTrusted(RunId, stateRow.activeRunId, "RunId");
26 const heartbeatAt = await getHeartbeatAt(context, runId);
27 const observedAt = heartbeatAt ?? stateRow.updatedAt;
28 if (observedAt + HEARTBEAT_STALE_AFTER_MS > Date.now()) {
29 return null;
30 }
31 
32 const runStub = getRunStub(context, runId);
33 let runMeta: RunMetaState | null;
34 try {
35 runMeta = await runStub.getRunSummary(runId);
36 } catch (error) {
37 context.logger.error("watchdog_run_summary_failed", {
38 projectId,
39 runId,
40 error: error instanceof Error ? error.message : String(error),
41 });
42 throw error;
43 }
44 
45 const terminalStatus: ProjectRunTerminalStatus =
46 runMeta && isTerminalStatus(runMeta.status) && runMeta.finishedAt !== null ? runMeta.status : "failed";
47 const lastError = terminalStatus === "failed" ? (runMeta?.errorMessage ?? "runner_lost") : runMeta?.errorMessage;
48 
49 await context.db.transaction((tx) => {
50 const row = getRunRow(tx, projectId, runId);
51 if (!row || (row.status !== "active" && row.status !== "cancel_requested")) {
52 return;
53 }
54 
55 tx.update(projectSchema.projectRuns)
56 .set({
57 status: terminalStatus,
58 position: null,
59 dispatchStatus: "terminal",
60 d1SyncStatus: nextTerminalD1SyncStatus(expectTrusted(D1SyncStatus, row.d1SyncStatus, "D1SyncStatus")),
61 lastError: lastError ?? null,
62 })
63 .where(eq(projectSchema.projectRuns.runId, row.runId))
64 .run();
65 tx.update(projectSchema.projectState)
66 .set({
67 activeRunId: null,
68 updatedAt: Date.now(),
69 })
70 .where(eq(projectSchema.projectState.projectId, projectId))
71 .run();
72 promoteNextPendingRun(tx, projectId);
73 });
74 
75 await setHeartbeatAt(context, runId, null);
76 if (!(runMeta && isTerminalStatus(runMeta.status) && runMeta.finishedAt !== null)) {
77 try {
78 const updatedRow = await getRunRowByRunId(context, runId);
79 if (updatedRow) {
80 await setRunDoTerminal(context, updatedRow, "failed", "runner_lost");
81 }
82 } catch (error) {
83 context.logger.error("watchdog_run_do_terminalize_failed", {
84 projectId,
85 runId,
86 error: error instanceof Error ? error.message : String(error),
87 });
88 }
89 }
90 
91 if (await attemptSandboxCleanup(context, projectId, runId, "watchdog")) {
92 await clearSandboxCleanupRetry(context, projectId, runId);
93 } else {
94 await seedSandboxCleanupRetry(context, projectId, runId);
95 }
96 
97 return runId;
98};