File
Blob: src/worker/api/private/projects/write-handlers.ts
| 1 | import { |
| 2 | CreateProjectRequest, |
| 3 | DEFAULT_DISPATCH_MODE, |
| 4 | DEFAULT_EXECUTION_RUNTIME, |
| 5 | ProjectId, |
| 6 | type ProjectResponse, |
| 7 | TriggerRunRequest, |
| 8 | type TriggerRunAcceptedResponse, |
| 9 | toRunStatusOrNull, |
| 10 | UpdateProjectRequest, |
| 11 | } from "@/contracts"; |
| 12 | import { type AcceptManualRunInput, expectTrusted } from "@/worker/contracts"; |
| 13 | import type { AppContext } from "@/worker/hono"; |
| 14 | import { |
| 15 | deleteProjectIndexById, |
| 16 | insertProjectIndex, |
| 17 | listLatestRunStatusByProjectIds, |
| 18 | } from "@/worker/db/d1/repositories"; |
| 19 | import { HttpError, parseJson } from "@/worker/http"; |
| 20 | import { serializeProjectConfigSummary } from "@/worker/presentation/serializers"; |
| 21 | import { encryptSecret } from "@/worker/security/secrets"; |
| 22 | import { generateDurableEntityId } from "@/worker/services"; |
| 23 | import { queueProjectReconciliation } from "@/worker/api/private/reconciliation"; |
| 24 | import { requireOwnedProject } from "@/worker/api/private/shared"; |
| 25 | import { |
| 26 | assertValidSlug, |
| 27 | normalizeBranchName, |
| 28 | normalizeConfigPath, |
| 29 | normalizeProjectName, |
| 30 | normalizeRepositoryUrl, |
| 31 | } from "@/worker/validation"; |
| 32 | |
| 33 | import { |
| 34 | getProjectStub, |
| 35 | isConstraintError, |
| 36 | logger, |
| 37 | toTriggeredByUserId, |
| 38 | UNIQUE_PROJECT_SLUG_CONSTRAINT, |
| 39 | } from "./shared"; |
| 40 | |
| 41 | export const handleCreateProject = async (c: AppContext): Promise<Response> => { |
| 42 | const user = c.get("user"); |
| 43 | const db = c.get("db"); |
| 44 | const payload = await parseJson(c.req.raw, CreateProjectRequest); |
| 45 | assertValidSlug(payload.projectSlug, "projectSlug"); |
| 46 | |
| 47 | const now = Date.now(); |
| 48 | const encryptedToken = typeof payload.repoToken === "string" ? await encryptSecret(c.env, payload.repoToken) : null; |
| 49 | const projectId = expectTrusted(ProjectId, generateDurableEntityId("prj", now), "ProjectId"); |
| 50 | |
| 51 | const project = { |
| 52 | id: projectId, |
| 53 | ownerUserId: user.id, |
| 54 | ownerSlug: user.slug, |
| 55 | projectSlug: payload.projectSlug, |
| 56 | name: normalizeProjectName(payload.name), |
| 57 | repoUrl: normalizeRepositoryUrl(payload.repoUrl), |
| 58 | defaultBranch: normalizeBranchName(payload.defaultBranch), |
| 59 | configPath: normalizeConfigPath(payload.configPath ?? ".anvil.yml"), |
| 60 | createdAt: now, |
| 61 | updatedAt: now, |
| 62 | } as const; |
| 63 | |
| 64 | try { |
| 65 | await insertProjectIndex(db, project); |
| 66 | } catch (error) { |
| 67 | if (isConstraintError(error, UNIQUE_PROJECT_SLUG_CONSTRAINT)) { |
| 68 | throw new HttpError(409, "project_slug_taken", "Project slug is already in use for this owner."); |
| 69 | } |
| 70 | |
| 71 | throw error; |
| 72 | } |
| 73 | |
| 74 | try { |
| 75 | await getProjectStub(c.env, projectId).initializeProject({ |
| 76 | projectId, |
| 77 | name: project.name, |
| 78 | repoUrl: project.repoUrl, |
| 79 | defaultBranch: project.defaultBranch, |
| 80 | configPath: project.configPath, |
| 81 | encryptedRepoToken: encryptedToken, |
| 82 | dispatchMode: payload.dispatchMode ?? DEFAULT_DISPATCH_MODE, |
| 83 | executionRuntime: DEFAULT_EXECUTION_RUNTIME, |
| 84 | createdAt: now, |
| 85 | updatedAt: now, |
| 86 | }); |
| 87 | } catch (error) { |
| 88 | try { |
| 89 | await deleteProjectIndexById(db, projectId); |
| 90 | } catch (cleanupError) { |
| 91 | logger.error("project_create_compensation_failed", { |
| 92 | projectId, |
| 93 | userId: user.id, |
| 94 | error: cleanupError instanceof Error ? cleanupError.message : String(cleanupError), |
| 95 | }); |
| 96 | } |
| 97 | |
| 98 | throw error; |
| 99 | } |
| 100 | |
| 101 | logger.info("project_created", { |
| 102 | projectId: project.id, |
| 103 | userId: user.id, |
| 104 | }); |
| 105 | |
| 106 | const response: ProjectResponse = { |
| 107 | project: serializeProjectConfigSummary( |
| 108 | { |
| 109 | ...project, |
| 110 | dispatchMode: payload.dispatchMode ?? DEFAULT_DISPATCH_MODE, |
| 111 | }, |
| 112 | null, |
| 113 | ), |
| 114 | }; |
| 115 | |
| 116 | return c.json(response, 201); |
| 117 | }; |
| 118 | |
| 119 | export const handleUpdateProject = async (c: AppContext): Promise<Response> => { |
| 120 | const { projectId, project: projectIndex } = await requireOwnedProject(c); |
| 121 | const db = c.get("db"); |
| 122 | const user = c.get("user"); |
| 123 | |
| 124 | const payload = await parseJson(c.req.raw, UpdateProjectRequest); |
| 125 | if ( |
| 126 | payload.name === undefined && |
| 127 | payload.repoUrl === undefined && |
| 128 | payload.defaultBranch === undefined && |
| 129 | payload.configPath === undefined && |
| 130 | payload.dispatchMode === undefined && |
| 131 | payload.repoToken === undefined |
| 132 | ) { |
| 133 | throw new HttpError(400, "empty_update", "At least one project field must be provided."); |
| 134 | } |
| 135 | |
| 136 | const encryptedToken = |
| 137 | typeof payload.repoToken === "string" |
| 138 | ? await encryptSecret(c.env, payload.repoToken) |
| 139 | : payload.repoToken === null |
| 140 | ? null |
| 141 | : undefined; |
| 142 | const projectStub = getProjectStub(c.env, projectId); |
| 143 | const result = await projectStub.updateProjectConfig({ |
| 144 | projectId, |
| 145 | name: payload.name === undefined ? undefined : normalizeProjectName(payload.name), |
| 146 | repoUrl: payload.repoUrl === undefined ? undefined : normalizeRepositoryUrl(payload.repoUrl), |
| 147 | defaultBranch: payload.defaultBranch === undefined ? undefined : normalizeBranchName(payload.defaultBranch), |
| 148 | configPath: payload.configPath === undefined ? undefined : normalizeConfigPath(payload.configPath), |
| 149 | dispatchMode: payload.dispatchMode, |
| 150 | encryptedRepoToken: encryptedToken, |
| 151 | now: Date.now(), |
| 152 | }); |
| 153 | switch (result.kind) { |
| 154 | case "invalid": |
| 155 | throw new HttpError(result.status as 400 | 404 | 409 | 500, result.code, result.message, result.details); |
| 156 | case "not_found": |
| 157 | throw new HttpError(500, "project_config_missing", "Project configuration is missing."); |
| 158 | case "applied": |
| 159 | break; |
| 160 | } |
| 161 | |
| 162 | queueProjectReconciliation(c, projectId, "update_project"); |
| 163 | const latestRunStatusByProjectId = await listLatestRunStatusByProjectIds(db, [projectIndex.id]); |
| 164 | |
| 165 | logger.info("project_updated", { |
| 166 | projectId: projectIndex.id, |
| 167 | userId: user.id, |
| 168 | }); |
| 169 | |
| 170 | const response: ProjectResponse = { |
| 171 | project: serializeProjectConfigSummary( |
| 172 | { |
| 173 | ...projectIndex, |
| 174 | ...result.config, |
| 175 | }, |
| 176 | toRunStatusOrNull(latestRunStatusByProjectId.get(projectIndex.id)), |
| 177 | ), |
| 178 | }; |
| 179 | |
| 180 | return c.json(response, 200); |
| 181 | }; |
| 182 | |
| 183 | export const handleTriggerProjectRun = async (c: AppContext): Promise<Response> => { |
| 184 | const { projectId } = await requireOwnedProject(c); |
| 185 | const user = c.get("user"); |
| 186 | |
| 187 | const payload = await parseJson(c.req.raw, TriggerRunRequest); |
| 188 | const branch = payload.branch === undefined ? null : normalizeBranchName(payload.branch); |
| 189 | const triggeredByUserId = toTriggeredByUserId(user.id); |
| 190 | |
| 191 | const projectStub = getProjectStub(c.env, projectId); |
| 192 | const acceptInput: AcceptManualRunInput = { |
| 193 | projectId, |
| 194 | triggeredByUserId, |
| 195 | branch, |
| 196 | }; |
| 197 | const accepted = await projectStub.acceptManualRun(acceptInput); |
| 198 | if (accepted.kind === "rejected") { |
| 199 | throw new HttpError(409, "project_queue_full", "Project already has the maximum number of queued runs."); |
| 200 | } |
| 201 | |
| 202 | queueProjectReconciliation(c, projectId, "trigger_run"); |
| 203 | |
| 204 | logger.info("run_triggered", { |
| 205 | projectId, |
| 206 | runId: accepted.runId, |
| 207 | userId: user.id, |
| 208 | }); |
| 209 | |
| 210 | const response: TriggerRunAcceptedResponse = { runId: accepted.runId }; |
| 211 | return c.json(response, 202); |
| 212 | }; |