Skip to content
File

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

typescript91 lines
1import { type ProjectId, RunId } from "@/contracts";
2import { D1SyncStatus, expectTrusted } from "@/worker/contracts";
3import { type NewRunIndexRow } from "@/worker/db/d1/repositories";
4import * as projectSchema from "@/worker/db/durable/schema/project-do";
5 
6import { getRunRowByRunId } from "../repo";
7import { ensureProjectState, isTerminalStatus, promoteNextPendingRun } from "../transitions/shared";
8import type { D1RetryPhase, D1RetryState, ProjectDoContext } from "../types";
9 
10export type ProjectRunTableRow = typeof projectSchema.projectRuns.$inferSelect;
11 
12export const getRunStub = (context: ProjectDoContext, runId: RunId) => context.env.RUN_DO.getByName(runId);
13 
14export const getRetryDelay = (attempt: number, delays: readonly number[]): number =>
15 delays[Math.min(attempt - 1, delays.length - 1)];
16 
17export const buildQueuedRunIndexTruthRow = (row: ProjectRunTableRow): NewRunIndexRow => ({
18 id: row.runId,
19 projectId: row.projectId,
20 triggeredByUserId: row.triggeredByUserId,
21 triggerType: row.triggerType,
22 branch: row.branch,
23 commitSha: row.commitSha,
24 dispatchMode: row.dispatchMode,
25 executionRuntime: row.executionRuntime,
26 status: "queued",
27 queuedAt: row.createdAt,
28 startedAt: null,
29 finishedAt: null,
30 exitCode: null,
31});
32 
33export const getNextD1RetryAttempt = (retryState: D1RetryState | null, phase: D1RetryPhase): number => {
34 if (!retryState || retryState.phase !== phase) {
35 return 1;
36 }
37 
38 return retryState.attempt + 1;
39};
40 
41export const resolvePostQueuedSyncD1Status = (
42 latestRow: ProjectRunTableRow,
43 syncedTruthRow: NewRunIndexRow,
44 preserveMetadataRetry = false,
45): D1SyncStatus => {
46 const latestD1SyncStatus = expectTrusted(D1SyncStatus, latestRow.d1SyncStatus, "D1SyncStatus");
47 
48 if (latestD1SyncStatus === "done") {
49 return "done";
50 }
51 
52 if (latestD1SyncStatus === "needs_terminal_update" || isTerminalStatus(latestRow.status)) {
53 return "needs_terminal_update";
54 }
55 
56 if (latestD1SyncStatus === "needs_update") {
57 if (preserveMetadataRetry) {
58 return "needs_update";
59 }
60 
61 return latestRow.commitSha === syncedTruthRow.commitSha ? "current" : "needs_update";
62 }
63 
64 return "current";
65};
66 
67export const buildLatestQueuedRunIndexTruth = async (
68 context: ProjectDoContext,
69 runId: RunId,
70 rowOverride?: ProjectRunTableRow,
71): Promise<{
72 latestRow: ProjectRunTableRow;
73 truthRow: NewRunIndexRow;
74} | null> => {
75 const row = rowOverride ?? (await getRunRowByRunId(context, runId));
76 if (!row) {
77 return null;
78 }
79 
80 return {
81 latestRow: row,
82 truthRow: buildQueuedRunIndexTruthRow(row),
83 };
84};
85 
86export const promotePendingRun = (context: ProjectDoContext, projectId: ProjectId): RunId | null =>
87 context.db.transaction((tx) => {
88 ensureProjectState(context, tx, projectId);
89 return promoteNextPendingRun(tx, projectId);
90 });