File
Blob: src/worker/api/private/projects/read-handlers.ts
| 1 | import { type GetProjectRunsResponse, type GetProjectsResponse, toRunStatusOrNull } from "@/contracts"; |
| 2 | import type { AppContext } from "@/worker/hono"; |
| 3 | import { |
| 4 | listLatestRunStatusByProjectIds, |
| 5 | listProjectRunsPage, |
| 6 | listProjectsByOwnerUserId, |
| 7 | } from "@/worker/db/d1/repositories"; |
| 8 | import { |
| 9 | serializePendingRunSummary, |
| 10 | serializeProjectDetail, |
| 11 | serializeProjectSummary, |
| 12 | serializeRunSummary, |
| 13 | } from "@/worker/presentation/serializers"; |
| 14 | import { queueProjectReconciliation } from "@/worker/api/private/reconciliation"; |
| 15 | import { requireOwnedProject } from "@/worker/api/private/shared"; |
| 16 | |
| 17 | import { |
| 18 | decodeRunCursor, |
| 19 | encodeRunCursor, |
| 20 | getProjectStub, |
| 21 | getRunStub, |
| 22 | mergeRunSummaryWithMeta, |
| 23 | parseProjectRunsQuery, |
| 24 | } from "./shared"; |
| 25 | export const handleGetProjects = async (c: AppContext): Promise<Response> => { |
| 26 | const user = c.get("user"); |
| 27 | const db = c.get("db"); |
| 28 | const projectRows = await listProjectsByOwnerUserId(db, user.id); |
| 29 | const latestRunStatusByProjectId = await listLatestRunStatusByProjectIds( |
| 30 | db, |
| 31 | projectRows.map((project) => project.id), |
| 32 | ); |
| 33 | |
| 34 | const projects = projectRows.map((project) => |
| 35 | serializeProjectSummary(project, toRunStatusOrNull(latestRunStatusByProjectId.get(project.id))), |
| 36 | ); |
| 37 | |
| 38 | const response: GetProjectsResponse = { projects }; |
| 39 | return c.json(response, 200); |
| 40 | }; |
| 41 | |
| 42 | export const handleGetProjectDetail = async (c: AppContext): Promise<Response> => { |
| 43 | const { projectId, project: projectIndex } = await requireOwnedProject(c); |
| 44 | const db = c.get("db"); |
| 45 | const projectStub = getProjectStub(c.env, projectId); |
| 46 | const [projectConfig, coordination, latestRunStatusByProjectId] = await Promise.all([ |
| 47 | projectStub.getProjectConfig(projectId), |
| 48 | projectStub.getProjectDetailState(projectId), |
| 49 | listLatestRunStatusByProjectIds(db, [projectIndex.id]), |
| 50 | ]); |
| 51 | if (!projectConfig) { |
| 52 | throw new Error(`Project config ${projectId} is missing.`); |
| 53 | } |
| 54 | |
| 55 | let activeRun = null; |
| 56 | if (coordination.activeRunId) { |
| 57 | const runMeta = await getRunStub(c.env, coordination.activeRunId).getRunSummary(coordination.activeRunId); |
| 58 | if (runMeta) { |
| 59 | activeRun = mergeRunSummaryWithMeta(coordination.activeRunId, runMeta, null); |
| 60 | } |
| 61 | } |
| 62 | |
| 63 | const detail = serializeProjectDetail({ |
| 64 | project: { |
| 65 | ...projectIndex, |
| 66 | ...projectConfig, |
| 67 | }, |
| 68 | lastRunStatus: |
| 69 | activeRun?.status ?? |
| 70 | (coordination.pendingRuns.length > 0 |
| 71 | ? "queued" |
| 72 | : toRunStatusOrNull(latestRunStatusByProjectId.get(projectIndex.id))), |
| 73 | activeRun, |
| 74 | pendingRuns: coordination.pendingRuns.map(serializePendingRunSummary), |
| 75 | }); |
| 76 | |
| 77 | if (coordination.activeRunId || coordination.pendingRuns.length > 0) { |
| 78 | queueProjectReconciliation(c, projectId, "get_project_detail"); |
| 79 | } |
| 80 | |
| 81 | return c.json(detail, 200); |
| 82 | }; |
| 83 | |
| 84 | export const handleGetProjectRuns = async (c: AppContext): Promise<Response> => { |
| 85 | const { projectId } = await requireOwnedProject(c); |
| 86 | const db = c.get("db"); |
| 87 | |
| 88 | const query = parseProjectRunsQuery(c); |
| 89 | const cursor = query.cursor ? decodeRunCursor(query.cursor) : undefined; |
| 90 | const rows = await listProjectRunsPage(db, projectId, query.limit + 1, cursor); |
| 91 | let runs = rows.slice(0, query.limit).map(serializeRunSummary); |
| 92 | |
| 93 | if (!query.cursor) { |
| 94 | const coordination = await getProjectStub(c.env, projectId).getProjectDetailState(projectId); |
| 95 | if (coordination.activeRunId) { |
| 96 | const activeMeta = await getRunStub(c.env, coordination.activeRunId).getRunSummary(coordination.activeRunId); |
| 97 | if (activeMeta) { |
| 98 | const existingIndex = runs.findIndex((run) => run.id === coordination.activeRunId); |
| 99 | if (existingIndex !== -1) { |
| 100 | const activeSummary = mergeRunSummaryWithMeta(coordination.activeRunId, activeMeta, runs[existingIndex]); |
| 101 | runs.splice(existingIndex, 1, activeSummary); |
| 102 | } |
| 103 | } |
| 104 | } |
| 105 | } |
| 106 | |
| 107 | const nextCursor = |
| 108 | rows.length > query.limit && runs.length > 0 |
| 109 | ? encodeRunCursor({ |
| 110 | queuedAt: Date.parse(runs[runs.length - 1].queuedAt), |
| 111 | runId: runs[runs.length - 1].id, |
| 112 | }) |
| 113 | : null; |
| 114 | |
| 115 | const response: GetProjectRunsResponse = { runs, nextCursor }; |
| 116 | return c.json(response, 200); |
| 117 | }; |