Merge a faster deploy plan: Cloudflare's API read directly and in parallel (2 s against 4 min), and the plan runs beside the check
6 files+365−370/6 viewed
| 4 | 4 | # | |
| 5 | 5 | # check the deploy manifest is consistent, and the tool's tests pass | |
| 6 | 6 | # plan what changed since each Worker's live commit, and pending migrations | |
| 7 | − | # migrate pending D1 migrations, before any code | |
| 7 | + | # (beside check, not after it: neither waits for the other) | |
| 8 | + | # migrate pending D1 migrations, before any code, once check and plan pass | |
| 8 | 9 | # core, edge, front the units of each stage, in jobs that share a build; | |
| 9 | 10 | # a stage starts only when the one before it succeeded | |
| 10 | 11 | # smoke sign-in, sign-up and the waitlist still work on g1t.sh | |
| ⋯ | |||
| 76 | 77 | - name: The deploy tool's tests | |
| 77 | 78 | run: npm run test:deploy | |
| 78 | 79 | ||
| 80 | + | # Runs beside check: nothing deploys until both have succeeded. | |
| 79 | 81 | plan: | |
| 80 | 82 | name: Plan | |
| 81 | − | needs: check | |
| 82 | 83 | runs-on: ubuntu-latest | |
| 83 | 84 | # Production's secrets, without a deployment: planning deploys nothing. | |
| 84 | 85 | environment: | |
| ⋯ | |||
| 99 | 100 | with: | |
| 100 | 101 | # Each Worker's live commit is compared with this one. | |
| 101 | 102 | fetch-depth: 0 | |
| 102 | − | - name: Install Wrangler | |
| 103 | − | run: npm ci --workspaces=false --no-audit --no-fund | |
| 103 | + | # No npm ci: with CLOUDFLARE_API_TOKEN the plan reads Cloudflare's API | |
| 104 | + | # itself (scripts/deploy/cloudflare.mjs), and needs no Wrangler. | |
| 104 | 105 | - name: Plan | |
| 105 | 106 | id: plan | |
| 106 | 107 | env: | |
| ⋯ | |||
| 115 | 116 | ||
| 116 | 117 | migrate: | |
| 117 | 118 | name: Migrations | |
| 118 | − | needs: plan | |
| 119 | + | needs: [check, plan] | |
| 119 | 120 | if: ${{ needs.plan.outputs.migrate == 'true' && inputs.dry_run != true }} | |
| 120 | 121 | runs-on: ubuntu-latest | |
| 121 | 122 | environment: | |
| ⋯ | |||
| 133 | 134 | ||
| 134 | 135 | core: | |
| 135 | 136 | name: core (${{ matrix.group }}) | |
| 136 | − | needs: [plan, migrate] | |
| 137 | + | needs: [check, plan, migrate] | |
| 137 | 138 | # Runs when nothing before it failed: a migrate job skipped for having | |
| 138 | 139 | # nothing to apply is not a failure. | |
| 139 | 140 | if: ${{ !failure() && !cancelled() && needs.plan.outputs.has_core == 'true' && inputs.dry_run != true }} | |
| ⋯ | |||
| 226 | 227 | ||
| 227 | 228 | edge: | |
| 228 | 229 | name: edge (${{ matrix.group }}) | |
| 229 | − | needs: [plan, migrate, core] | |
| 230 | + | needs: [check, plan, migrate, core] | |
| 230 | 231 | if: ${{ !failure() && !cancelled() && needs.plan.outputs.has_edge == 'true' && inputs.dry_run != true }} | |
| 231 | 232 | runs-on: ${{ (matrix.rust || matrix.image) && 'g1t-4core' || 'ubuntu-latest' }} | |
| 232 | 233 | environment: | |
| ⋯ | |||
| 241 | 242 | ||
| 242 | 243 | front: | |
| 243 | 244 | name: front (${{ matrix.group }}) | |
| 244 | − | needs: [plan, migrate, core, edge] | |
| 245 | + | needs: [check, plan, migrate, core, edge] | |
| 245 | 246 | if: ${{ !failure() && !cancelled() && needs.plan.outputs.has_front == 'true' && inputs.dry_run != true }} | |
| 246 | 247 | runs-on: ${{ (matrix.rust || matrix.image) && 'g1t-4core' || 'ubuntu-latest' }} | |
| 247 | 248 | environment: | |
| ⋯ | |||
| 260 | 261 | # real waitlist is never touched. scripts/ops/smoke.mjs. | |
| 261 | 262 | smoke: | |
| 262 | 263 | name: Smoke | |
| 263 | − | needs: [plan, core, edge, front] | |
| 264 | + | needs: [check, plan, core, edge, front] | |
| 264 | 265 | if: ${{ !failure() && !cancelled() && inputs.dry_run != true && (needs.plan.outputs.has_core == 'true' || needs.plan.outputs.has_edge == 'true' || needs.plan.outputs.has_front == 'true') }} | |
| 265 | 266 | runs-on: ubuntu-latest | |
| 266 | 267 | timeout-minutes: 5 | |
| 161 | 161 | assert!(!starts(&deploy, "front", &[("migrate", "success"), ("core", "failure"), ("edge", "skipped")], &all, push.clone(), false)); | |
| 162 | 162 | // A cancelled run starts nothing more. | |
| 163 | 163 | assert!(!starts(&deploy, "edge", &[("migrate", "success"), ("core", "success")], &all, push.clone(), true)); | |
| 164 | − | // A failed check: plan, and every stage after it, is skipped, and | |
| 165 | − | // failure() still sees the check's failure through them. | |
| 166 | − | let skipped = [("plan", "skipped"), ("migrate", "skipped"), ("core", "skipped"), ("edge", "skipped")]; | |
| 167 | − | assert!(!starts_after(&deploy, "plan", &[("check", "failure")], &all, push.clone(), false, false)); | |
| 168 | − | assert!(!starts_after(&deploy, "core", &skipped, &all, push.clone(), false, true)); | |
| 164 | + | // Check and plan run side by side: plan waits for nothing. | |
| 165 | + | let plan = deploy.jobs.iter().find(|j| j.id == "plan").unwrap(); | |
| 166 | + | assert!(plan.needs.is_empty(), "{:?}", plan.needs); | |
| 167 | + | for id in ["migrate", "core", "edge", "front", "smoke"] { | |
| 168 | + | let job = deploy.jobs.iter().find(|j| j.id == id).unwrap(); | |
| 169 | + | assert!(job.needs.iter().any(|n| n == "check") && job.needs.iter().any(|n| n == "plan"), "{id}: {:?}", job.needs); | |
| 170 | + | } | |
| 171 | + | // A failed check, though plan succeeded: nothing migrates or deploys, | |
| 172 | + | // and failure() still sees the check's failure through skipped jobs. | |
| 173 | + | assert!(!starts(&deploy, "migrate", &[("check", "failure"), ("plan", "success")], &all, push.clone(), false)); | |
| 174 | + | assert!(!starts(&deploy, "core", &[("check", "failure"), ("plan", "success"), ("migrate", "skipped")], &all, push.clone(), false)); | |
| 175 | + | assert!(!starts(&deploy, "front", &[("check", "failure"), ("plan", "success"), ("migrate", "skipped"), ("core", "skipped"), ("edge", "skipped")], &all, push.clone(), false)); | |
| 176 | + | let skipped = [("migrate", "skipped"), ("core", "skipped"), ("edge", "skipped")]; | |
| 169 | 177 | assert!(!starts_after(&deploy, "front", &skipped, &all, push.clone(), false, true)); | |
| 178 | + | assert!(!starts(&deploy, "smoke", &[("check", "failure"), ("core", "skipped"), ("edge", "skipped"), ("front", "skipped")], &all, push.clone(), false)); | |
| 179 | + | // A failed plan stops everything too. | |
| 180 | + | assert!(!starts(&deploy, "migrate", &[("plan", "failure")], &all, push.clone(), false)); | |
| 181 | + | assert!(!starts(&deploy, "core", &[("plan", "failure"), ("migrate", "skipped")], &all, push.clone(), false)); | |
| 170 | 182 | ||
| 171 | 183 | // Smoke follows the last stage that ran, and not a failed one. | |
| 172 | 184 | assert!(starts(&deploy, "smoke", &[("core", "success"), ("edge", "success"), ("front", "success")], &all, push.clone(), false)); |
| 124 | 124 | On the Worker itself. Every deploy runs `wrangler deploy --message | |
| 125 | 125 | "g1t-deploy <40-char sha> <subject>" --tag g1t-<12-char sha>`, which | |
| 126 | 126 | Cloudflare keeps as the version's `workers/message` and `workers/tag` | |
| 127 | − | annotations. `plan` reads them back with `wrangler deployments status` | |
| 128 | − | (the live version) and `wrangler versions list` (its annotations): two | |
| 129 | − | read-only calls per unit, in parallel; a plan of all 22 units takes about | |
| 130 | − | 10 seconds. No KV namespace or other infrastructure is needed. | |
| 127 | + | annotations. `plan` reads them back with two read-only requests per unit | |
| 128 | + | to Cloudflare's API, the ones `wrangler deployments status` and `wrangler | |
| 129 | + | versions list` make: `GET /accounts/{account_id}/workers/scripts/{worker}/deployments` | |
| 130 | + | (the live version) and `GET .../versions?deployable=true` (its | |
| 131 | + | annotations). Pending migrations are one more per database: `POST | |
| 132 | + | /accounts/{account_id}/d1/database/{database_id}/query` with `SELECT name | |
| 133 | + | FROM "d1_migrations"` (the table Wrangler keeps; a unit's | |
| 134 | + | `migrations_table` if it names one), compared with the `.sql` files in the | |
| 135 | + | unit's migrations folder. Every request starts at once, at most 16 in | |
| 136 | + | flight, each tried once more after a 403, 429, 5xx or network error, so a | |
| 137 | + | plan of every unit takes a few seconds and needs no Wrangler. No KV | |
| 138 | + | namespace or other infrastructure is needed. | |
| 139 | + | ||
| 140 | + | The API is used when the tool has the token Wrangler would be given (in CI, | |
| 141 | + | `CLOUDFLARE_API_TOKEN`; on a laptop, `CLOUDFLARE_DEPLOY_TOKEN`) and | |
| 142 | + | `CLOUDFLARE_ACCOUNT_ID` (or g1t's account by default). Without a token | |
| 143 | + | (your `wrangler login`) the plan asks Wrangler instead, with the same | |
| 144 | + | answers, a few units at a time; that takes minutes. Either way a Worker | |
| 145 | + | that does not exist (404, or Cloudflare's code 10007) is "never deployed", | |
| 146 | + | and any other failure is a reason in the plan, never a crash. | |
| 131 | 147 | ||
| 132 | 148 | - A version made by `wrangler secret put` keeps the code of the one before | |
| 133 | 149 | it, so the tool looks through those to the deploy before. | |
| ⋯ | |||
| 211 | 227 | On a laptop the tool uses your `wrangler login` (or `CLOUDFLARE_DEPLOY_TOKEN` | |
| 212 | 228 | if set), as `scripts/deploy.sh` always did: a `CLOUDFLARE_API_TOKEN` or | |
| 213 | 229 | global API key in your shell, or in the repository's `.env`, is ignored. | |
| 214 | − | With `CI=true` it uses `CLOUDFLARE_API_TOKEN`. | |
| 230 | + | With `CI=true` it uses `CLOUDFLARE_API_TOKEN`. With a token, `plan` reads | |
| 231 | + | Cloudflare's API itself; with only `wrangler login`, it asks Wrangler (see | |
| 232 | + | "Where the deployed commit is kept"). | |
| 215 | 233 | ||
| 216 | 234 | ### The runner's images | |
| 217 | 235 | ||
| ⋯ | |||
| 427 | 445 | | Job | Does | Needs | | |
| 428 | 446 | | --- | --- | --- | | |
| 429 | 447 | | `check` | `manifest --check` and `npm run test:deploy` | — | | |
| 430 | − | | `plan` | `plan --github-output`: outputs per stage, the plan in the run's summary | `check` | | |
| 431 | − | | `migrate` | `migrate --only <units with pending migrations>` | `plan`; skipped when none are pending | | |
| 432 | − | | `core`, `edge`, `front` | `deploy --only <units> --force --no-migrations`, one job per build group | the stages before; skipped when empty | | |
| 433 | − | | `smoke` | `node scripts/ops/smoke.mjs`: the landing page, sign-in, sign-up and pricing load, and the waitlist form reaches identity (sent an address identity refuses before keeping or counting anything, so the real waitlist is never touched) | every stage; skipped when nothing deployed | | |
| 448 | + | | `plan` | `plan --github-output`: outputs per stage, the plan in the run's summary. Reads Cloudflare's API itself, so it installs nothing. | — (runs beside `check`) | | |
| 449 | + | | `migrate` | `migrate --only <units with pending migrations>` | `check` and `plan`; skipped when none are pending | | |
| 450 | + | | `core`, `edge`, `front` | `deploy --only <units> --force --no-migrations`, one job per build group | `check`, `plan` and the stages before; skipped when empty | | |
| 451 | + | | `smoke` | `node scripts/ops/smoke.mjs`: the landing page, sign-in, sign-up and pricing load, and the waitlist form reaches identity (sent an address identity refuses before keeping or counting anything, so the real waitlist is never touched) | `check`, `plan` and every stage; skipped when nothing deployed | | |
| 434 | 452 | ||
| 435 | 453 | - **One at a time:** `concurrency: deploy-production`, never cancelled in | |
| 436 | 454 | progress; a second push waits. | |
| 455 | + | - **Check and plan side by side:** neither waits for the other, so the | |
| 456 | + | plan's few seconds overlap the check's tests; nothing migrates or deploys | |
| 457 | + | until both have succeeded. The plan job keeps `fetch-depth: 0`: it diffs | |
| 458 | + | from each Worker's live commit, which may be any commit, and tells a | |
| 459 | + | rollback by ancestry, which a shallow clone cannot answer. | |
| 437 | 460 | - **Build groups:** a stage's units are split so each job shares a build: | |
| 438 | 461 | Rust workers at most four to a job (each a 4-vCPU `g1t-4core` machine), the | |
| 439 | 462 | TypeScript Workers together, each site alone, and a unit whose image must | |
| ⋯ | |||
| 442 | 465 | another off mid-upload; the next stage then does not start. | |
| 443 | 466 | - **Tests:** there is no CI workflow on g1t yet; `main` is kept passing by | |
| 444 | 467 | the merge queue's checks. `check` runs the deploy tool's own tests. When a | |
| 445 | − | CI workflow is added, make `plan` wait for it (`workflow_run`, or a job in | |
| 468 | + | CI workflow is added, make `migrate` and the stages wait for it (`workflow_run`, or a job in | |
| 446 | 469 | this file). | |
| 447 | 470 | - **Machines:** Rust jobs and the runner's image run on `g1t-4core` (4 vCPUs, | |
| 448 | 471 | 12 GiB, 20 GB), the others on the standard machine | |
| ⋯ | |||
| 462 | 485 | - **Conditions:** each stage runs with `!failure() && !cancelled()`, which | |
| 463 | 486 | on g1t (as on GitHub) is true when no job before it failed, however far | |
| 464 | 487 | back: a `migrate` job skipped for having nothing to apply does not stop | |
| 465 | − | the stages after it, and a failed `check` stops all of them. | |
| 488 | + | the stages after it, and a failed `check` or `plan` stops all of them | |
| 489 | + | (`migrate`'s own condition needs both to have succeeded). | |
| 466 | 490 | - `crates/actions/tests/repository_workflows.rs` reads the workflow with | |
| 467 | 491 | g1t's own parser and expressions, and checks the jobs start, wait and | |
| 468 | 492 | stop as above (`cargo test -p g1t-actions --test repository_workflows`). | |
| ⋯ | |||
| 505 | 529 | status.g1t.sh | deploy.yml | production | |
| 506 | 530 | ``` | |
| 507 | 531 | ||
| 508 | − | `api.cloudflare.com` is Wrangler's API; `registry.cloudflare.com` is where | |
| 532 | + | `api.cloudflare.com` is Cloudflare's API, which Wrangler and the plan call; `registry.cloudflare.com` is where | |
| 509 | 533 | the deploy asks whether the runner's image is already built, and where the | |
| 510 | 534 | `runner-image` job pulls the base from and pushes the runner's image to | |
| 511 | 535 | (as `runner-base.yml` pushes the base); `status.g1t.sh` hears the deploy | |
| ⋯ | |||
| 525 | 549 | ||
| 526 | 550 | | Scope | Permission | Why | | |
| 527 | 551 | | --- | --- | --- | | |
| 528 | − | | Account | Workers Scripts: Edit | Upload, versions, deployments, crons, bindings, `secret list` (doctor) | | |
| 529 | − | | Account | D1: Edit | `d1 migrations list` and `apply` | | |
| 552 | + | | Account | Workers Scripts: Edit | Upload, versions, deployments (the plan reads both), crons, bindings, `secret list` (doctor) | | |
| 553 | + | | Account | D1: Edit | The plan's query of each database's `d1_migrations`, and `d1 migrations apply` | | |
| 530 | 554 | | Account | Queues: Edit | Attaching each unit's queue consumers on deploy | | |
| 531 | 555 | | Account | Workers R2 Storage: Read | Wrangler checks `og`'s bucket binding | | |
| 532 | 556 | | Account | Account Settings: Read | Wrangler reads the account | | |
| 29 | 29 | import { ensureWorkerBuild } from "./build-rust-worker.mjs"; | |
| 30 | 30 | import { | |
| 31 | 31 | annotation, | |
| 32 | + | apiAuth, | |
| 32 | 33 | applyMigrations, | |
| 33 | 34 | dockerAvailable, | |
| 34 | 35 | exec, | |
| ⋯ | |||
| 157 | 158 | migrations[unit.id] = await pendingMigrations(unit); | |
| 158 | 159 | }), | |
| 159 | 160 | ]; | |
| 160 | − | await pool(tasks, Math.max(8, opts.concurrency * 2), (task) => task()); | |
| 161 | + | // Through the API every read starts at once (the requests themselves are | |
| 162 | + | // held to API_CONCURRENCY); through Wrangler, a few processes at a time. | |
| 163 | + | await pool(tasks, apiAuth() ? tasks.length : Math.max(8, opts.concurrency * 2), (task) => task()); | |
| 161 | 164 | return { live, migrations }; | |
| 162 | 165 | } | |
| 163 | 166 | ||
| 1 | − | // Wrangler, as the deploy tool uses it: reading which commit each Worker | |
| 2 | − | // runs, D1 migrations, and deploying. Every call runs in the unit's own | |
| 3 | − | // folder, so Wrangler reads that unit's config (and not a .env at the | |
| 4 | − | // repository root, which may hold a token meant for something else). | |
| 1 | + | // Cloudflare, as the deploy tool uses it: reading which commit each Worker | |
| 2 | + | // runs and which D1 migrations are pending (Cloudflare's REST API when a | |
| 3 | + | // token is set, Wrangler otherwise), applying migrations, and deploying. | |
| 4 | + | // Every Wrangler call runs in the unit's own folder, so Wrangler reads that | |
| 5 | + | // unit's config (and not a .env at the repository root, which may hold a | |
| 6 | + | // token meant for something else). | |
| 5 | 7 | ||
| 6 | 8 | import { spawn } from "node:child_process"; | |
| 9 | + | import { readdirSync } from "node:fs"; | |
| 7 | 10 | import { join } from "node:path"; | |
| 8 | 11 | ||
| 9 | 12 | import { ROOT } from "./stack.mjs"; | |
| ⋯ | |||
| 134 | 137 | return { sha: null, version: version.id, why: "its live version was not deployed by scripts/deploy.mjs" }; | |
| 135 | 138 | } | |
| 136 | 139 | ||
| 140 | + | // ── Reading production ─────────────────────────────────────────────────── | |
| 141 | + | // | |
| 142 | + | // The plan reads Cloudflare's REST API directly when it has a token (always | |
| 143 | + | // in CI): one request per question, all in parallel, where a Wrangler | |
| 144 | + | // process would boot Node and check its login before each. Without one (a | |
| 145 | + | // laptop's `wrangler login`) it asks Wrangler, as it always did. The | |
| 146 | + | // requests are the ones Wrangler makes: `deployments status` prints the | |
| 147 | + | // first of GET .../deployments, `versions list` the items of GET | |
| 148 | + | // .../versions?deployable=true, and `d1 migrations list` compares the names | |
| 149 | + | // in the database's d1_migrations table with the files in its migrations | |
| 150 | + | // folder. | |
| 151 | + | ||
| 152 | + | const API = "https://api.cloudflare.com/client/v4"; | |
| 153 | + | /** At most this many requests to Cloudflare at once. */ | |
| 154 | + | export const API_CONCURRENCY = 16; | |
| 155 | + | ||
| 156 | + | /** | |
| 157 | + | * The token and account the plan reads Cloudflare's API with: the token | |
| 158 | + | * Wrangler would be given (see wranglerEnv), or null, and then Wrangler is | |
| 159 | + | * asked instead, with your `wrangler login`. | |
| 160 | + | */ | |
| 161 | + | export function apiAuth(base = process.env) { | |
| 162 | + | const env = wranglerEnv(base); | |
| 163 | + | if (!env.CLOUDFLARE_API_TOKEN) return null; | |
| 164 | + | return { token: env.CLOUDFLARE_API_TOKEN, account: env.CLOUDFLARE_ACCOUNT_ID }; | |
| 165 | + | } | |
| 166 | + | ||
| 167 | + | let inFlight = 0; | |
| 168 | + | const waiting = []; | |
| 169 | + | async function limited(task) { | |
| 170 | + | while (inFlight >= API_CONCURRENCY) await new Promise((resolve) => waiting.push(resolve)); | |
| 171 | + | inFlight++; | |
| 172 | + | try { | |
| 173 | + | return await task(); | |
| 174 | + | } finally { | |
| 175 | + | inFlight--; | |
| 176 | + | waiting.shift()?.(); | |
| 177 | + | } | |
| 178 | + | } | |
| 179 | + | ||
| 180 | + | /** Cloudflare's errors in an answer, as one line. */ | |
| 181 | + | export function apiError(status, body) { | |
| 182 | + | const errors = (body?.errors ?? []).map((e) => (e.code ? `${e.message} (${e.code})` : e.message)).filter(Boolean); | |
| 183 | + | return `Cloudflare API ${status || "request"} failed${errors.length ? `: ${errors.join("; ")}` : ""}`; | |
| 184 | + | } | |
| 185 | + | ||
| 186 | + | /** | |
| 187 | + | * One request to Cloudflare's API. Resolves with { status, result } or | |
| 188 | + | * { status, error, codes }; never throws. Retried once after a refusal | |
| 189 | + | * that may pass (a 403 while a token propagates, 429, 5xx, the network). | |
| 190 | + | */ | |
| 191 | + | export async function cloudflareApi(auth, path, { method = "GET", body, fetchImpl = fetch, retries = 1, timeoutMs = 30_000 } = {}) { | |
| 192 | + | for (let attempt = 0; ; attempt++) { | |
| 193 | + | let status = 0; | |
| 194 | + | let data = null; | |
| 195 | + | let error = null; | |
| 196 | + | try { | |
| 197 | + | const response = await limited(async () => { | |
| 198 | + | const res = await fetchImpl(`${API}${path}`, { | |
| 199 | + | method, | |
| 200 | + | headers: { authorization: `Bearer ${auth.token}`, ...(body === undefined ? {} : { "content-type": "application/json" }) }, | |
| 201 | + | body: body === undefined ? undefined : JSON.stringify(body), | |
| 202 | + | signal: AbortSignal.timeout(timeoutMs), | |
| 203 | + | }); | |
| 204 | + | return { status: res.status, text: await res.text() }; | |
| 205 | + | }); | |
| 206 | + | status = response.status; | |
| 207 | + | try { | |
| 208 | + | data = JSON.parse(response.text); | |
| 209 | + | } catch { | |
| 210 | + | error = `Cloudflare API ${status}: ${response.text.trim().slice(0, 300) || "an empty answer"}`; | |
| 211 | + | } | |
| 212 | + | } catch (failure) { | |
| 213 | + | error = `Cloudflare API request failed: ${failure?.message ?? failure}`; | |
| 214 | + | } | |
| 215 | + | if (!error && status < 400 && data?.success !== false) return { status, result: data?.result ?? null }; | |
| 216 | + | error ??= apiError(status, data); | |
| 217 | + | const codes = (data?.errors ?? []).map((e) => e.code); | |
| 218 | + | const mayPass = status === 0 || status === 403 || status === 429 || status >= 500; | |
| 219 | + | if (attempt >= retries || !mayPass) return { status, error, codes }; | |
| 220 | + | await new Promise((resolve) => setTimeout(resolve, 500 * (attempt + 1))); | |
| 221 | + | } | |
| 222 | + | } | |
| 223 | + | ||
| 224 | + | /** Cloudflare's code for a Worker that does not exist. */ | |
| 225 | + | const SCRIPT_NOT_FOUND = 10007; | |
| 226 | + | ||
| 227 | + | /** | |
| 228 | + | * The live commit from the API's answers: `deployments` is the result of | |
| 229 | + | * GET .../deployments ({ deployments: [newest first] }), `versions` that of | |
| 230 | + | * GET .../versions?deployable=true ({ items }), or null if it could not be | |
| 231 | + | * read. `wrangler deployments status --json` prints the first deployment, | |
| 232 | + | * and `versions list --json` those items. | |
| 233 | + | */ | |
| 234 | + | export function liveFromApi(worker, deployments, versions) { | |
| 235 | + | const latest = deployments?.deployments?.[0]; | |
| 236 | + | if (!latest) return { sha: null, error: `The Worker ${worker} has no deployments.` }; | |
| 237 | + | return liveCommit(latest, versions?.items ?? []); | |
| 238 | + | } | |
| 239 | + | ||
| 137 | 240 | /** Reads the commit a unit's Worker runs. Never throws. */ | |
| 138 | − | export async function readLive(unit) { | |
| 241 | + | export async function readLive(unit, { auth = apiAuth(), fetchImpl = fetch } = {}) { | |
| 242 | + | if (!auth) return readLiveWithWrangler(unit); | |
| 243 | + | const script = `/accounts/${auth.account}/workers/scripts/${encodeURIComponent(unit.worker)}`; | |
| 244 | + | const [deployments, versions] = await Promise.all([ | |
| 245 | + | cloudflareApi(auth, `${script}/deployments`, { fetchImpl }), | |
| 246 | + | cloudflareApi(auth, `${script}/versions?deployable=true`, { fetchImpl }), | |
| 247 | + | ]); | |
| 248 | + | if (deployments.error) { | |
| 249 | + | if (deployments.status === 404 || deployments.codes.includes(SCRIPT_NOT_FOUND)) return { sha: null, missing: true, why: "never deployed" }; | |
| 250 | + | return { sha: null, error: deployments.error }; | |
| 251 | + | } | |
| 252 | + | try { | |
| 253 | + | return liveFromApi(unit.worker, deployments.result, versions.error ? null : versions.result); | |
| 254 | + | } catch (error) { | |
| 255 | + | return { sha: null, error: String(error.message ?? error) }; | |
| 256 | + | } | |
| 257 | + | } | |
| 258 | + | ||
| 259 | + | /** readLive through Wrangler: `deployments status` and `versions list`. */ | |
| 260 | + | async function readLiveWithWrangler(unit) { | |
| 139 | 261 | const cwd = join(ROOT, unit.path); | |
| 140 | 262 | const [status, versions] = await Promise.all([ | |
| 141 | 263 | wrangler(["deployments", "status", "--name", unit.worker, "--json"], { cwd }), | |
| ⋯ | |||
| 159 | 281 | return [...new Set(names)]; | |
| 160 | 282 | } | |
| 161 | 283 | ||
| 162 | − | /** Pending migrations of a unit's database: { pending } or { error }. */ | |
| 163 | − | export async function pendingMigrations(unit) { | |
| 284 | + | /** Wrangler's default name for the table of applied migrations. */ | |
| 285 | + | export const MIGRATIONS_TABLE = "d1_migrations"; | |
| 286 | + | ||
| 287 | + | /** | |
| 288 | + | * A unit's database as its wrangler.jsonc names it: { id, table, dir } | |
| 289 | + | * (dir relative to the repository), or null. | |
| 290 | + | */ | |
| 291 | + | export function databaseOf(unit) { | |
| 292 | + | const db = (unit.config?.d1_databases ?? []).find((d) => d.database_name === unit.d1?.database); | |
| 293 | + | if (!db?.database_id) return null; | |
| 294 | + | return { id: db.database_id, table: db.migrations_table || MIGRATIONS_TABLE, dir: join(unit.path, unit.d1.migrations) }; | |
| 295 | + | } | |
| 296 | + | ||
| 297 | + | /** The migration files in a folder, as Wrangler finds them: its *.sql files. */ | |
| 298 | + | export function migrationFiles(dir) { | |
| 299 | + | return readdirSync(dir, { withFileTypes: true }) | |
| 300 | + | .filter((entry) => entry.isFile() && entry.name.endsWith(".sql")) | |
| 301 | + | .map((entry) => entry.name) | |
| 302 | + | .sort(); | |
| 303 | + | } | |
| 304 | + | ||
| 305 | + | /** | |
| 306 | + | * The files not yet applied, in the files' order, given the D1 query API's | |
| 307 | + | * result for `SELECT name FROM d1_migrations` ([{ results: [{ name }] }]). | |
| 308 | + | */ | |
| 309 | + | export function pendingAgainst(files, result) { | |
| 310 | + | const applied = new Set((result?.[0]?.results ?? []).map((row) => row.name)); | |
| 311 | + | return files.filter((file) => !applied.has(file)); | |
| 312 | + | } | |
| 313 | + | ||
| 314 | + | const quoteIdentifier = (name) => `"${name.replaceAll('"', '""')}"`; | |
| 315 | + | ||
| 316 | + | /** Pending migrations of a unit's database: { pending } or { error }. Never throws. */ | |
| 317 | + | export async function pendingMigrations(unit, { auth = apiAuth(), fetchImpl = fetch, root = ROOT } = {}) { | |
| 318 | + | const db = databaseOf(unit); | |
| 319 | + | if (!auth || !db) return pendingMigrationsWithWrangler(unit); | |
| 320 | + | let files; | |
| 321 | + | try { | |
| 322 | + | files = migrationFiles(join(root, db.dir)); | |
| 323 | + | } catch (error) { | |
| 324 | + | return { error: `Could not read ${db.dir}: ${error.message ?? error}` }; | |
| 325 | + | } | |
| 326 | + | const answer = await cloudflareApi(auth, `/accounts/${auth.account}/d1/database/${db.id}/query`, { | |
| 327 | + | method: "POST", | |
| 328 | + | body: { sql: `SELECT name FROM ${quoteIdentifier(db.table)} ORDER BY id` }, | |
| 329 | + | fetchImpl, | |
| 330 | + | }); | |
| 331 | + | if (answer.error) { | |
| 332 | + | // A database no migration was ever applied to has no table yet. | |
| 333 | + | if (/no such table/i.test(answer.error)) return { pending: files }; | |
| 334 | + | return { error: answer.error }; | |
| 335 | + | } | |
| 336 | + | return { pending: pendingAgainst(files, answer.result) }; | |
| 337 | + | } | |
| 338 | + | ||
| 339 | + | /** pendingMigrations through Wrangler: `d1 migrations list --remote`. */ | |
| 340 | + | async function pendingMigrationsWithWrangler(unit) { | |
| 164 | 341 | const list = () => wrangler(["d1", "migrations", "list", unit.d1.database, "--remote"], { cwd: join(ROOT, unit.path) }); | |
| 165 | 342 | // Once more after a failure: Cloudflare's API sometimes answers 403 | |
| 166 | 343 | // while Wrangler's login refreshes (seen on 2026-10-06). | |
| 24 | 24 | writeDeployConfig, | |
| 25 | 25 | } from "./image.mjs"; | |
| 26 | 26 | ||
| 27 | − | import { annotation, commitFrom, liveCommit, pendingFrom, versionFrom } from "./cloudflare.mjs"; | |
| 27 | + | import { | |
| 28 | + | annotation, | |
| 29 | + | apiAuth, | |
| 30 | + | commitFrom, | |
| 31 | + | databaseOf, | |
| 32 | + | liveCommit, | |
| 33 | + | liveFromApi, | |
| 34 | + | migrationFiles, | |
| 35 | + | pendingAgainst, | |
| 36 | + | pendingFrom, | |
| 37 | + | pendingMigrations, | |
| 38 | + | readLive, | |
| 39 | + | versionFrom, | |
| 40 | + | } from "./cloudflare.mjs"; | |
| 28 | 41 | import { changedNames, parseCargoLock, parseNpmLock, reaches } from "./lockfiles.mjs"; | |
| 29 | 42 | import { decide, planJson, pool } from "./plan.mjs"; | |
| 30 | 43 | import { | |
| ⋯ | |||
| 534 | 547 | assert.equal(versionFrom("Deployed g1t-events triggers\nCurrent Version ID: 2c7fc93a-82d9-4850-8a77-ba887d157a4f\n"), "2c7fc93a-82d9-4850-8a77-ba887d157a4f"); | |
| 535 | 548 | }); | |
| 536 | 549 | ||
| 550 | + | // Cloudflare's REST answers, as the API sends them (Wrangler's --json | |
| 551 | + | // prints the first deployment, and the versions' items). | |
| 552 | + | const AUTH = { token: "t", account: "acct" }; | |
| 553 | + | const ok = (result) => ({ success: true, errors: [], messages: [], result }); | |
| 554 | + | const refused = (code, message) => ({ success: false, errors: [{ code, message }], messages: [], result: null }); | |
| 555 | + | /** A fetch that answers by path: { "GET /path": [status, body] }. */ | |
| 556 | + | function fakeFetch(answers, calls = []) { | |
| 557 | + | return async (url, init) => { | |
| 558 | + | const path = url.replace("https://api.cloudflare.com/client/v4", ""); | |
| 559 | + | calls.push({ path, init }); | |
| 560 | + | const found = answers[`${init.method} ${path}`]; | |
| 561 | + | if (!found) throw new Error(`unexpected ${init.method} ${path}`); | |
| 562 | + | const [status, body] = found; | |
| 563 | + | return { status, text: async () => (typeof body === "string" ? body : JSON.stringify(body)) }; | |
| 564 | + | }; | |
| 565 | + | } | |
| 566 | + | const SCRIPT = "/accounts/acct/workers/scripts/g1t-events"; | |
| 567 | + | const deploymentsAnswer = (id) => | |
| 568 | + | ok({ | |
| 569 | + | deployments: [ | |
| 570 | + | { id: "d2", source: "wrangler", strategy: "percentage", versions: [{ version_id: id, percentage: 100 }], annotations: { "workers/triggered_by": "deployment" } }, | |
| 571 | + | { id: "d1", source: "wrangler", strategy: "percentage", versions: [{ version_id: "v1", percentage: 100 }] }, | |
| 572 | + | ], | |
| 573 | + | }); | |
| 574 | + | const versionsAnswer = ok({ items: [version("v1", 1, annotation(OLD).message), version("v2", 2, null, "secret")] }); | |
| 575 | + | ||
| 576 | + | test("the API's answers: the live commit as Wrangler's would give it", async () => { | |
| 577 | + | const calls = []; | |
| 578 | + | const fetchImpl = fakeFetch({ [`GET ${SCRIPT}/deployments`]: [200, deploymentsAnswer("v2")], [`GET ${SCRIPT}/versions?deployable=true`]: [200, versionsAnswer] }, calls); | |
| 579 | + | const found = await readLive(unit("events"), { auth: AUTH, fetchImpl }); | |
| 580 | + | assert.equal(found.sha, OLD); | |
| 581 | + | assert.equal(found.version, "v2"); | |
| 582 | + | assert.equal(found.at, "2026-10-05T00:00:00Z"); | |
| 583 | + | assert.equal(calls[0].init.headers.authorization, "Bearer t"); | |
| 584 | + | // The same as liveCommit over what Wrangler printed. | |
| 585 | + | assert.deepEqual(liveFromApi("w", deploymentsAnswer("v2").result, versionsAnswer.result), liveCommit(deploymentsAnswer("v2").result.deployments[0], versionsAnswer.result.items)); | |
| 586 | + | // Versions that cannot be read: the live one is not among them. | |
| 587 | + | assert.match(liveFromApi("w", deploymentsAnswer("v2").result, null).why, /not among/); | |
| 588 | + | // No deployment at all is an error, as `deployments status` made it. | |
| 589 | + | assert.match(liveFromApi("w", { deployments: [] }, versionsAnswer.result).error, /no deployments/); | |
| 590 | + | }); | |
| 591 | + | ||
| 592 | + | test("the API's answers: a Worker never deployed, and a refusal, never throw", async () => { | |
| 593 | + | const missing = fakeFetch({ | |
| 594 | + | [`GET ${SCRIPT}/deployments`]: [404, refused(10007, "workers.api.error.script_not_found")], | |
| 595 | + | [`GET ${SCRIPT}/versions?deployable=true`]: [404, refused(10007, "workers.api.error.script_not_found")], | |
| 596 | + | }); | |
| 597 | + | assert.deepEqual(await readLive(unit("events"), { auth: AUTH, fetchImpl: missing }), { sha: null, missing: true, why: "never deployed" }); | |
| 598 | + | const denied = fakeFetch({ | |
| 599 | + | [`GET ${SCRIPT}/deployments`]: [400, refused(10000, "Authentication error")], | |
| 600 | + | [`GET ${SCRIPT}/versions?deployable=true`]: [400, refused(10000, "Authentication error")], | |
| 601 | + | }); | |
| 602 | + | const found = await readLive(unit("events"), { auth: AUTH, fetchImpl: denied }); | |
| 603 | + | assert.equal(found.sha, null); | |
| 604 | + | assert.match(found.error, /Authentication error \(10000\)/); | |
| 605 | + | const broken = async () => { | |
| 606 | + | throw new Error("connect ECONNREFUSED"); | |
| 607 | + | }; | |
| 608 | + | const offline = await readLive(unit("events"), { auth: AUTH, fetchImpl: broken }); | |
| 609 | + | assert.match(offline.error, /ECONNREFUSED/); | |
| 610 | + | const garbled = fakeFetch({ [`GET ${SCRIPT}/deployments`]: [200, "<html>"], [`GET ${SCRIPT}/versions?deployable=true`]: [200, versionsAnswer] }); | |
| 611 | + | assert.match((await readLive(unit("events"), { auth: AUTH, fetchImpl: garbled })).error, /<html>/); | |
| 612 | + | }); | |
| 613 | + | ||
| 614 | + | test("the API is used only with the token Wrangler would be given", () => { | |
| 615 | + | assert.equal(apiAuth({ CLOUDFLARE_API_TOKEN: "shell" }), null); | |
| 616 | + | assert.deepEqual(apiAuth({ CI: "true", CLOUDFLARE_API_TOKEN: "ci", CLOUDFLARE_ACCOUNT_ID: "a" }), { token: "ci", account: "a" }); | |
| 617 | + | assert.equal(apiAuth({ CLOUDFLARE_DEPLOY_TOKEN: "mine" }).token, "mine"); | |
| 618 | + | assert.equal(apiAuth({ CLOUDFLARE_DEPLOY_TOKEN: "mine" }).account, "1e6f2cffa3f445920836e8ebe446bb58"); | |
| 619 | + | }); | |
| 620 | + | ||
| 621 | + | test("pending migrations: the folder's files less the names in d1_migrations", async () => { | |
| 622 | + | const events = unit("events"); | |
| 623 | + | const db = databaseOf(events); | |
| 624 | + | assert.equal(db.id, events.config.d1_databases[0].database_id); | |
| 625 | + | assert.equal(db.table, "d1_migrations"); | |
| 626 | + | const files = migrationFiles(join(ROOT, db.dir)); | |
| 627 | + | assert.ok(files.length > 2 && files.every((f) => f.endsWith(".sql"))); | |
| 628 | + | assert.deepEqual(files, [...files].sort()); | |
| 629 | + | ||
| 630 | + | assert.deepEqual(pendingAgainst(["0001_a.sql", "0002_b.sql", "0003_c.sql"], [{ results: [{ name: "0001_a.sql" }, { name: "0002_b.sql" }], success: true, meta: {} }]), ["0003_c.sql"]); | |
| 631 | + | assert.deepEqual(pendingAgainst(["0001_a.sql"], [{ results: [], success: true }]), ["0001_a.sql"]); | |
| 632 | + | ||
| 633 | + | const query = `POST /accounts/acct/d1/database/${db.id}/query`; | |
| 634 | + | const calls = []; | |
| 635 | + | const applied = files.slice(0, -1).map((name, i) => ({ name, id: i + 1 })); | |
| 636 | + | const found = await pendingMigrations(events, { auth: AUTH, fetchImpl: fakeFetch({ [query]: [200, ok([{ results: applied, success: true, meta: {} }])] }, calls) }); | |
| 637 | + | assert.deepEqual(found, { pending: files.slice(-1) }); | |
| 638 | + | assert.match(JSON.parse(calls[0].init.body).sql, /^SELECT name FROM "d1_migrations"/); | |
| 639 | + | const upToDate = await pendingMigrations(events, { auth: AUTH, fetchImpl: fakeFetch({ [query]: [200, ok([{ results: files.map((name) => ({ name })) }])] }) }); | |
| 640 | + | assert.deepEqual(upToDate, { pending: [] }); | |
| 641 | + | // A database never migrated has no table: everything is pending. | |
| 642 | + | const fresh = await pendingMigrations(events, { auth: AUTH, fetchImpl: fakeFetch({ [query]: [400, refused(7500, "no such table: d1_migrations: SQLITE_ERROR")] }) }); | |
| 643 | + | assert.deepEqual(fresh, { pending: files }); | |
| 644 | + | const denied = await pendingMigrations(events, { auth: AUTH, fetchImpl: fakeFetch({ [query]: [400, refused(10000, "Authentication error")] }) }); | |
| 645 | + | assert.match(denied.error, /Authentication error/); | |
| 646 | + | }); | |
| 647 | + | ||
| 537 | 648 | test("the pool runs everything, no more than its limit at once", async () => { | |
| 538 | 649 | let running = 0; | |
| 539 | 650 | let most = 0; | |