Skip to content
File

Blob: src/worker/api/private/runs.ts

typescript260 lines
1import { type LogStreamTicketResponse, ProjectId, RunId, type RunStatus, toRunStatusOrNull, UserId } from "@/contracts";
2import { expectTrusted, isoDateTimeFromTimestamp, isTerminalStatus, type RunMetaState } from "@/worker/contracts";
3import { toCodecIssueDetails } from "@/lib/codec-errors";
4import { queueProjectReconciliation } from "@/worker/api/private/reconciliation";
5import { findOwnedProjectForCurrentUser } from "@/worker/api/private/shared";
6import { findRunIndexById, findUserById, type ProjectRow, type RunIndexRow } from "@/worker/db/d1/repositories";
7import type { AppContext } from "@/worker/hono";
8import { HttpError } from "@/worker/http";
9import {
10 serializeLogEvent,
11 serializeRunDetail,
12 serializeRunStep,
13 serializeRunSummary,
14} from "@/worker/presentation/serializers";
15import { createInternalRunLogHeaders } from "@/worker/run-logs/auth";
16import { extractTimestampFromDurableEntityId, generateOpaqueToken } from "@/worker/services";
17 
18const LOG_TICKET_TTL_SECONDS = 60;
19const LOG_TICKET_KEY_PREFIX = "run-log-ticket:";
20 
21const getProjectStub = (env: Env, projectId: ProjectId) => env.PROJECT_DO.getByName(projectId);
22const getRunStub = (env: Env, runId: RunId) => env.RUN_DO.getByName(runId);
23 
24interface OwnedRunResolution {
25 projectId: ProjectId;
26 project: ProjectRow;
27 d1Run: RunIndexRow | undefined;
28 summary: RunMetaState | null;
29}
30 
31interface LogTicketRecord {
32 runId: RunId;
33 userId: UserId;
34 expiresAt: number;
35}
36 
37const decodeRunIdParam = (runId: string): RunId => {
38 try {
39 return RunId.assertDecode(runId);
40 } catch (error) {
41 throw new HttpError(404, "run_not_found", "Run was not found.", toCodecIssueDetails(error));
42 }
43};
44 
45const getLogTicketKey = (ticket: string): string => `${LOG_TICKET_KEY_PREFIX}${ticket}`;
46 
47const parseLogTicketRecord = (value: string | null): LogTicketRecord | null => {
48 if (!value) {
49 return null;
50 }
51 
52 try {
53 const parsed = JSON.parse(value) as unknown;
54 if (
55 !parsed ||
56 typeof parsed !== "object" ||
57 Array.isArray(parsed) ||
58 !("runId" in parsed) ||
59 !("userId" in parsed) ||
60 !("expiresAt" in parsed)
61 ) {
62 return null;
63 }
64 
65 return {
66 runId: RunId.assertDecode(parsed.runId),
67 userId: UserId.assertDecode(parsed.userId),
68 expiresAt: Number(parsed.expiresAt),
69 };
70 } catch {
71 return null;
72 }
73};
74 
75const requireRunIdParam = (c: AppContext): RunId => {
76 const runId = c.req.param("runId");
77 if (!runId) {
78 throw new HttpError(404, "run_not_found", "Run was not found.");
79 }
80 
81 return decodeRunIdParam(runId);
82};
83 
84const resolveOwnedRun = async (c: AppContext, runId: RunId): Promise<OwnedRunResolution> => {
85 const db = c.get("db");
86 const d1Run = await findRunIndexById(db, runId);
87 const summary = d1Run ? null : await getRunStub(c.env, runId).getRunSummary(runId);
88 const projectId = d1Run?.projectId ?? summary?.projectId;
89 
90 if (!projectId) {
91 throw new HttpError(404, "run_not_found", "Run was not found.");
92 }
93 
94 const decodedProjectId = expectTrusted(ProjectId, projectId, "ProjectId");
95 
96 const project = await findOwnedProjectForCurrentUser(c, decodedProjectId);
97 if (!project || (!d1Run && !summary)) {
98 throw new HttpError(404, "run_not_found", "Run was not found.");
99 }
100 
101 return { projectId: decodedProjectId, project, d1Run, summary };
102};
103 
104interface RunDetailResponseState {
105 response: Response;
106 runStatus: RunStatus;
107}
108 
109const buildRunDetailResponse = async (
110 c: AppContext,
111 runId: RunId,
112 resolution: OwnedRunResolution,
113 statusOverride: RunStatus | null = null,
114): Promise<RunDetailResponseState> => {
115 const runDetail = await getRunStub(c.env, runId).getRunDetail(runId);
116 const d1RunStatus = toRunStatusOrNull(resolution.d1Run?.status);
117 const runStatus = statusOverride ?? runDetail.meta?.status ?? resolution.summary?.status ?? d1RunStatus ?? "queued";
118 
119 const run = serializeRunSummary({
120 id: runId,
121 projectId: resolution.project.id,
122 triggeredByUserId: resolution.d1Run?.triggeredByUserId ?? null,
123 triggerType:
124 runDetail.meta?.triggerType ?? resolution.summary?.triggerType ?? resolution.d1Run?.triggerType ?? "manual",
125 branch:
126 runDetail.meta?.branch ??
127 resolution.summary?.branch ??
128 resolution.d1Run?.branch ??
129 resolution.project.defaultBranch,
130 commitSha: runDetail.meta?.commitSha ?? resolution.summary?.commitSha ?? resolution.d1Run?.commitSha ?? null,
131 status: runStatus,
132 queuedAt: resolution.d1Run?.queuedAt ?? extractTimestampFromDurableEntityId(runId) ?? Date.now(),
133 startedAt: runDetail.meta?.startedAt ?? resolution.summary?.startedAt ?? resolution.d1Run?.startedAt ?? null,
134 finishedAt: runDetail.meta?.finishedAt ?? resolution.summary?.finishedAt ?? resolution.d1Run?.finishedAt ?? null,
135 exitCode: runDetail.meta?.exitCode ?? resolution.summary?.exitCode ?? resolution.d1Run?.exitCode ?? null,
136 });
137 
138 return {
139 response: c.json(
140 serializeRunDetail({
141 run,
142 currentStep: runDetail.meta?.currentStep ?? null,
143 errorMessage: runDetail.meta?.errorMessage ?? null,
144 steps: runDetail.steps.map(serializeRunStep),
145 recentLogs: runDetail.recentLogs.map(serializeLogEvent),
146 detailAvailable: runDetail.meta !== null,
147 }),
148 200,
149 ),
150 runStatus,
151 };
152};
153 
154const requireWebSocketUpgrade = (request: Request): void => {
155 if (request.headers.get("upgrade")?.toLowerCase() !== "websocket") {
156 throw new HttpError(426, "upgrade_required", "WebSocket upgrade required.");
157 }
158};
159 
160export const handleGetRunDetail = async (c: AppContext): Promise<Response> => {
161 const runId = requireRunIdParam(c);
162 const resolution = await resolveOwnedRun(c, runId);
163 const result = await buildRunDetailResponse(c, runId, resolution);
164 
165 if (!isTerminalStatus(result.runStatus) || !resolution.d1Run) {
166 queueProjectReconciliation(c, resolution.projectId, "get_run_detail");
167 }
168 
169 return result.response;
170};
171 
172export const handleCancelRun = async (c: AppContext): Promise<Response> => {
173 const runId = requireRunIdParam(c);
174 const resolution = await resolveOwnedRun(c, runId);
175 
176 const cancelResult = await getProjectStub(c.env, resolution.projectId).requestRunCancel({
177 projectId: resolution.projectId,
178 runId,
179 });
180 
181 const statusOverride =
182 cancelResult.status === "cancel_requested" ||
183 cancelResult.status === "canceled" ||
184 cancelResult.status === "passed" ||
185 cancelResult.status === "failed"
186 ? cancelResult.status
187 : null;
188 
189 queueProjectReconciliation(c, resolution.projectId, "cancel_run");
190 
191 return (await buildRunDetailResponse(c, runId, await resolveOwnedRun(c, runId), statusOverride)).response;
192};
193 
194export const handleCreateRunLogTicket = async (c: AppContext): Promise<Response> => {
195 const runId = requireRunIdParam(c);
196 await resolveOwnedRun(c, runId);
197 
198 const userId = expectTrusted(UserId, c.get("user").id, "UserId");
199 const ticket = generateOpaqueToken(32);
200 const expiresAt = Date.now() + LOG_TICKET_TTL_SECONDS * 1000;
201 
202 await c.env.LOG_TICKETS.put(
203 getLogTicketKey(ticket),
204 JSON.stringify({
205 runId,
206 userId,
207 expiresAt,
208 } satisfies LogTicketRecord),
209 { expirationTtl: LOG_TICKET_TTL_SECONDS },
210 );
211 
212 const response: LogStreamTicketResponse = {
213 ticket,
214 expiresAt: isoDateTimeFromTimestamp(expiresAt),
215 };
216 
217 return c.json(response, 200);
218};
219 
220export const handleGetRunLogsWebSocket = async (c: AppContext): Promise<Response> => {
221 const runId = requireRunIdParam(c);
222 requireWebSocketUpgrade(c.req.raw);
223 
224 const ticket = c.req.query("ticket")?.trim();
225 if (!ticket) {
226 throw new HttpError(403, "invalid_log_ticket", "Log stream ticket is missing or invalid.");
227 }
228 
229 const ticketKey = getLogTicketKey(ticket);
230 const ticketRecord = parseLogTicketRecord(await c.env.LOG_TICKETS.get(ticketKey));
231 await c.env.LOG_TICKETS.delete(ticketKey);
232 
233 if (!ticketRecord || ticketRecord.expiresAt <= Date.now() || ticketRecord.runId !== runId) {
234 throw new HttpError(403, "invalid_log_ticket", "Log stream ticket is missing or invalid.");
235 }
236 
237 const ticketUser = await findUserById(c.get("db"), ticketRecord.userId);
238 if (!ticketUser) {
239 throw new HttpError(403, "invalid_session", "Session user no longer exists.");
240 }
241 
242 if (ticketUser.disabledAt !== null) {
243 throw new HttpError(403, "user_disabled", "User account is disabled.");
244 }
245 
246 const url = new URL(c.req.url);
247 url.search = "";
248 
249 const headers = createInternalRunLogHeaders(c.req.raw.headers, runId, ticketRecord.userId);
250 headers.delete("authorization");
251 headers.delete("cookie");
252 
253 const internalRequest = new Request(url.toString(), {
254 method: "GET",
255 headers,
256 });
257 
258 return getRunStub(c.env, runId).fetch(internalRequest);
259};