Merge branch 'worktree-agent-ab9543c492a7ed481' into spend-guardrails
# Conflicts: # services/og/src/index.ts
36 files+835−1370/36 viewed
| 56 | 56 | name: String, | |
| 57 | 57 | } | |
| 58 | 58 | ||
| 59 | + | /// Artifacts older runners kept in KV expire 14 days after they were made, | |
| 60 | + | /// and none has been made there since artifacts moved to R2 on 2026-10-08: | |
| 61 | + | /// from 2026-10-22T00:00Z every one is gone, and KV is not asked (a list is | |
| 62 | + | /// the dearest thing KV does). Delete the KV artifact code after that date, | |
| 63 | + | /// with its twin in apps/web/app/lib/artifacts.server.ts. | |
| 64 | + | const LEGACY_KV_UNTIL_MS: u64 = 1_792_627_200_000; | |
| 65 | + | ||
| 66 | + | /// Whether artifacts kept in KV may still be there at `now`. | |
| 67 | + | fn legacy_kv(now: u64) -> bool { | |
| 68 | + | now < LEGACY_KV_UNTIL_MS | |
| 69 | + | } | |
| 70 | + | ||
| 59 | 71 | fn store(env: &Env) -> Result<KvStore> { | |
| 60 | 72 | env.kv("BLOBS") | |
| 61 | 73 | } | |
| ⋯ | |||
| 404 | 416 | // Artifacts older runners kept in KV, for the days they | |
| 405 | 417 | // are still there. | |
| 406 | 418 | let legacy_run = run_id.as_deref().unwrap_or(run); | |
| 407 | − | for (_, meta) in list(kv, &format!("a/{legacy_run}/")).await? { | |
| 408 | − | if !listed.iter().any(|a| a["name"] == meta.name.as_str()) { | |
| 409 | − | listed.push(json!({ "name": meta.name, "size": meta.size, "format": "tgz" })); | |
| 419 | + | if legacy_kv(g1t_kit::now_ms()) { | |
| 420 | + | for (_, meta) in list(kv, &format!("a/{legacy_run}/")).await? { | |
| 421 | + | if !listed.iter().any(|a| a["name"] == meta.name.as_str()) { | |
| 422 | + | listed.push(json!({ "name": meta.name, "size": meta.size, "format": "tgz" })); | |
| 423 | + | } | |
| 410 | 424 | } | |
| 411 | 425 | } | |
| 412 | 426 | crate::reply(&listed) | |
| ⋯ | |||
| 547 | 561 | match found { | |
| 548 | 562 | Outcome::Ok(found) => stream(bucket, &found.object, &found.artifact.format).await, | |
| 549 | 563 | // One an older runner kept in KV. | |
| 550 | − | Outcome::Fail(_) => match get(kv, &format!("a/{run}/{name}")).await? { | |
| 564 | + | Outcome::Fail(_) if legacy_kv(g1t_kit::now_ms()) => match get(kv, &format!("a/{run}/{name}")).await? { | |
| 551 | 565 | Some(bytes) => Response::from_bytes(bytes), | |
| 552 | 566 | None => error(404, "No such artifact."), | |
| 553 | 567 | }, | |
| 568 | + | Outcome::Fail(_) => error(404, "No such artifact."), | |
| 554 | 569 | } | |
| 555 | 570 | } | |
| 556 | 571 | _ => error(404, "No such endpoint."), | |
| ⋯ | |||
| 612 | 627 | return Response::redirect_with_status(Url::parse(&crate::artifacts::blob_url(&services.addresses.api, &found.blob))?, 302); | |
| 613 | 628 | } | |
| 614 | 629 | // Kept in KV by an older runner. | |
| 630 | + | if !legacy_kv(g1t_kit::now_ms()) { | |
| 631 | + | return error(404, "No such artifact, or it has expired."); | |
| 632 | + | } | |
| 615 | 633 | match get(&store(env)?, &format!("a/{run}/{name}")).await? { | |
| 616 | 634 | Some(bytes) => { | |
| 617 | 635 | let mut response = Response::from_bytes(bytes)?; | |
| ⋯ | |||
| 624 | 642 | } | |
| 625 | 643 | } | |
| 626 | 644 | ||
| 645 | + | #[cfg(test)] | |
| 646 | + | mod tests { | |
| 647 | + | use super::*; | |
| 648 | + | ||
| 649 | + | #[test] | |
| 650 | + | fn kv_artifacts_are_not_asked_for_after_they_have_all_expired() { | |
| 651 | + | assert_eq!(g1t_contracts::time::rfc3339(LEGACY_KV_UNTIL_MS), "2026-10-22T00:00:00.000Z"); | |
| 652 | + | assert!(legacy_kv(LEGACY_KV_UNTIL_MS - 1)); | |
| 653 | + | assert!(!legacy_kv(LEGACY_KV_UNTIL_MS)); | |
| 654 | + | } | |
| 655 | + | } | |
| 78 | 78 | Why each of these is missing, and what to use instead, is on | |
| 79 | 79 | [What g1t can't do yet](/about/limitations/#actions-and-runners). | |
| 80 | 80 | ||
| 81 | + | ## Schedules | |
| 82 | + | ||
| 83 | + | A workflow with `on: schedule` runs on the default branch's latest commit, | |
| 84 | + | at each time its cron lines name, in UTC: | |
| 85 | + | ||
| 86 | + | ```yaml | |
| 87 | + | on: | |
| 88 | + | schedule: | |
| 89 | + | - cron: "0 9 * * mon" # Mondays at 09:00 UTC | |
| 90 | + | - cron: "*/30 * * * *" # every 30 minutes | |
| 91 | + | ``` | |
| 92 | + | ||
| 93 | + | Each line has five fields: minute, hour, day of month, month and day of | |
| 94 | + | the week. A field takes `*`, a number, a range `1-5`, a list `1,15` and a | |
| 95 | + | step `*/10`; days take `mon` to `sun`, and months `jan` to `dec`. | |
| 96 | + | ||
| 97 | + | | Rule | What happens | | |
| 98 | + | | --- | --- | | |
| 99 | + | | Every 5 minutes at most | A schedule more frequent than every 5 minutes, such as `* * * * *` or `*/2 * * * *`, runs every 5 minutes instead, on the five-minute marks (:00, :05, :10 and so on), at each mark that ends five minutes in which it would have run. A mark outside the hours, days or months the schedule names never runs: `* 9 * * *` runs from 09:00 to 09:55. | | |
| 100 | + | | No push for 60 days | Schedules pause in a repository that has had no push for 60 days. The next push to any branch resumes them. Other events and `workflow_dispatch` still start the workflow. | | |
| 101 | + | | Actions not paid for | When a scheduled run's job could not start because the workspace's plan, spend limit or the open-source pool does not cover it, that run fails and says why, and the workflow's schedule waits an hour before it tries again. | | |
| 102 | + | | Archived repository | Schedules wait until the repository is unarchived. | | |
| 103 | + | ||
| 81 | 104 | ## Actions and workflows from other repositories | |
| 82 | 105 | ||
| 83 | 106 | A step's `uses: owner/repo@ref` (or `owner/repo/path@ref`) and a job's |
| 124 | 124 | are served under the binding name your config gives them. Cron triggers | |
| 125 | 125 | (`triggers.crons`) are not scheduled, so a `scheduled` handler never | |
| 126 | 126 | runs; the deployment says so in its warnings. | |
| 127 | + | - Each request your Worker answers may use up to 50 ms of CPU time | |
| 128 | + | (waiting on the network does not count) and make up to 50 requests of | |
| 129 | + | its own (`fetch` calls and the like). A request that goes over either | |
| 130 | + | is stopped and answered with a 503 page that says so. Static assets are | |
| 131 | + | served without running your Worker, and count toward neither. | |
| 127 | 132 | - `vars` are deployed as plain-text bindings (or JSON, for objects). Rows | |
| 128 | 133 | of the project's [secrets and variables](/guides/secrets-and-variables/) | |
| 129 | 134 | available to Deployments are bound too, and replace a `var` of the same |
| 2 | 2 | ||
| 3 | 3 | import type { RepoPath, Viewer } from "@g1t/contracts"; | |
| 4 | 4 | ||
| 5 | + | import { LEGACY_KV_UNTIL } from "./artifacts"; | |
| 5 | 6 | import { actions } from "./services.server"; | |
| 6 | 7 | ||
| 7 | 8 | /** | |
| 8 | 9 | * A run's artifacts, for its page. The actions service lists them (their | |
| 9 | 10 | * bytes are in R2, downloaded through the API's signed links); artifacts an | |
| 10 | 11 | * older runner kept in KV (`a/{run}/{name}`, its bytes in chunks `…#0`, | |
| 11 | − | * `…#1`) are listed too until KV expires them. | |
| 12 | + | * `…#1`) are listed too until KV expires them, for runs the caller has | |
| 13 | + | * already been shown (`legacyArtifactsWorthAsking` in artifacts.ts). | |
| 12 | 14 | */ | |
| 13 | 15 | export type ArtifactRow = { | |
| 14 | 16 | /** The artifact's number; null for one kept in KV. */ | |
| ⋯ | |||
| 35 | 37 | })); | |
| 36 | 38 | } | |
| 37 | 39 | ||
| 40 | + | /** A run's artifacts the actions service keeps, which checks who may see them. */ | |
| 38 | 41 | export async function listArtifacts(repo: RepoPath, viewer: Viewer, run: string): Promise<ArtifactRow[]> { | |
| 39 | − | const [kept, legacy] = await Promise.all([ | |
| 40 | − | actions.artifacts(repo, viewer, { run, per_page: 100 }), | |
| 41 | − | kvArtifacts(run).catch(() => []), | |
| 42 | − | ]); | |
| 42 | + | const kept = await actions.artifacts(repo, viewer, { run, per_page: 100 }); | |
| 43 | 43 | const rows: ArtifactRow[] = kept.ok | |
| 44 | 44 | ? kept.value.artifacts.map((a) => ({ id: a.id, name: a.name, size: a.size, expiresAt: a.expires_at, createdAt: a.created_at })) | |
| 45 | 45 | : []; | |
| 46 | + | return rows.sort((a, b) => a.name.localeCompare(b.name)); | |
| 47 | + | } | |
| 48 | + | ||
| 49 | + | /** | |
| 50 | + | * `rows` with the artifacts an older runner kept in KV for `run` added. | |
| 51 | + | * Only for a run the viewer was allowed to see. Delete after | |
| 52 | + | * 2026-10-22T00:00Z (`LEGACY_KV_UNTIL`). | |
| 53 | + | */ | |
| 54 | + | export async function withLegacyArtifacts(rows: ArtifactRow[], run: string): Promise<ArtifactRow[]> { | |
| 55 | + | if (Date.now() >= LEGACY_KV_UNTIL) return rows; | |
| 56 | + | const legacy = await kvArtifacts(run).catch(() => []); | |
| 57 | + | const all = [...rows]; | |
| 46 | 58 | for (const row of legacy) { | |
| 47 | − | if (!rows.some((r) => r.name === row.name)) rows.push(row); | |
| 59 | + | if (!all.some((r) => r.name === row.name)) all.push(row); | |
| 48 | 60 | } | |
| 49 | − | return rows.sort((a, b) => a.name.localeCompare(b.name)); | |
| 61 | + | return all.sort((a, b) => a.name.localeCompare(b.name)); | |
| 50 | 62 | } | |
| 51 | 63 | ||
| 52 | 64 | export async function readArtifact(run: string, name: string): Promise<Uint8Array | null> { | |
| 65 | + | // Every artifact kept in KV has expired (artifacts.ts). | |
| 66 | + | if (Date.now() >= LEGACY_KV_UNTIL) return null; | |
| 53 | 67 | const base = `a/${run}/${name}`; | |
| 54 | 68 | const meta = await env.BLOBS.get<Meta>(base, "json"); | |
| 55 | 69 | if (!meta) return null; | |
| 1 | 1 | import assert from "node:assert/strict"; | |
| 2 | 2 | import { test } from "node:test"; | |
| 3 | 3 | ||
| 4 | − | import { expiresIn, formatBytes } from "./artifacts.ts"; | |
| 4 | + | import { expiresIn, formatBytes, legacyArtifactsWorthAsking } from "./artifacts.ts"; | |
| 5 | 5 | ||
| 6 | 6 | test("sizes read as KB, MB or GB", () => { | |
| 7 | 7 | assert.equal(formatBytes(10), "1 KB"); | |
| ⋯ | |||
| 18 | 18 | assert.equal(expiresIn("2026-10-08T11:00:00Z", now), "expired"); | |
| 19 | 19 | assert.equal(expiresIn(null, now), ""); | |
| 20 | 20 | }); | |
| 21 | + | ||
| 22 | + | test("KV is asked for an old, finished run's artifacts only until they have expired", () => { | |
| 23 | + | const old = { createdAt: "2026-10-07T12:00:00Z", status: "completed" }; | |
| 24 | + | const now = Date.parse("2026-10-10T00:00:00Z"); | |
| 25 | + | assert.ok(legacyArtifactsWorthAsking(old, now)); | |
| 26 | + | // Still going: its artifacts are in R2, and its page refreshes. | |
| 27 | + | assert.ok(!legacyArtifactsWorthAsking({ ...old, status: "in_progress" }, now)); | |
| 28 | + | // Made after the move to R2. | |
| 29 | + | assert.ok(!legacyArtifactsWorthAsking({ ...old, createdAt: "2026-10-09T08:00:00Z" }, now)); | |
| 30 | + | // Every KV artifact has expired. | |
| 31 | + | assert.ok(!legacyArtifactsWorthAsking(old, Date.parse("2026-10-22T00:00:00Z"))); | |
| 32 | + | }); | |
| 16 | 16 | if (days === 1) return "expires tomorrow"; | |
| 17 | 17 | return `expires in ${days} days`; | |
| 18 | 18 | } | |
| 19 | + | ||
| 20 | + | /** | |
| 21 | + | * Artifacts older runners kept in Workers KV. None has been made there | |
| 22 | + | * since artifacts moved to R2 on 2026-10-08, and KV expires each 14 days | |
| 23 | + | * after it was made: from 2026-10-22T00:00Z every one is gone. Delete the | |
| 24 | + | * KV artifact code (here, artifacts.server.ts, and apps/api/src/blobs.rs) | |
| 25 | + | * after that date. | |
| 26 | + | */ | |
| 27 | + | export const LEGACY_KV_UNTIL = Date.parse("2026-10-22T00:00:00Z"); | |
| 28 | + | /** Runs made from this time on kept their artifacts in R2 only. */ | |
| 29 | + | export const LEGACY_KV_BEFORE = Date.parse("2026-10-09T00:00:00Z"); | |
| 30 | + | ||
| 31 | + | /** | |
| 32 | + | * Whether a run's page asks KV for artifacts an older runner kept there: | |
| 33 | + | * only for a finished run made before the move, and only until they have | |
| 34 | + | * all expired. A KV list is the dearest thing KV does, and a run still | |
| 35 | + | * going refreshes its page every few seconds. | |
| 36 | + | */ | |
| 37 | + | export function legacyArtifactsWorthAsking(run: { createdAt: string; status: string }, now: number = Date.now()): boolean { | |
| 38 | + | return now < LEGACY_KV_UNTIL && run.status === "completed" && Date.parse(run.createdAt) < LEGACY_KV_BEFORE; | |
| 39 | + | } |
| 4 | 4 | import { | |
| 5 | 5 | type Merged, | |
| 6 | 6 | placePushes, | |
| 7 | + | historyCovers, | |
| 7 | 8 | change, | |
| 8 | 9 | checksFact, | |
| 9 | 10 | confidenceAsk, | |
| ⋯ | |||
| 204 | 205 | assert.deepEqual(placePushes(history, pushes), [{ hash: "m1", at: "2026-10-05T12:00:00Z" }]); | |
| 205 | 206 | }); | |
| 206 | 207 | ||
| 208 | + | test("a short read of the branch is enough when it reaches the oldest push", () => { | |
| 209 | + | const h = (...hashes: string[]) => hashes.map((hash) => ({ hash })); | |
| 210 | + | const pushes = [ | |
| 211 | + | { time: "2026-10-07T12:00:00Z", data: { after: "c1", before: "c2" } }, | |
| 212 | + | { time: "2026-10-05T12:00:00Z", data: { after: "c2", before: "c3" } }, | |
| 213 | + | ]; | |
| 214 | + | // Holds the oldest push's before. | |
| 215 | + | assert.ok(historyCovers(h("c1", "c2", "c3"), pushes, 3)); | |
| 216 | + | // Full, and the oldest push reaches past it: read further. | |
| 217 | + | assert.ok(!historyCovers(h("c1", "c2", "x"), pushes, 3)); | |
| 218 | + | // Shorter than asked: the whole branch. | |
| 219 | + | assert.ok(historyCovers(h("c1", "c2"), pushes, 3)); | |
| 220 | + | // The oldest push made the branch: only the whole branch will do. | |
| 221 | + | assert.ok(!historyCovers(h("c1", "c2", "c3"), [{ time: "", data: { after: "c1" } }], 3)); | |
| 222 | + | }); | |
| 223 | + | ||
| 207 | 224 | test("a change landed without a person when g1t merged it", () => { | |
| 208 | 225 | assert.ok(landedByAgents({ mergedBy: "g1t" })); | |
| 209 | 226 | assert.ok(landedByAgents({ mergedBy: "g1t" })); | |
| 503 | 503 | return [...seen].flatMap(([hash, at]) => (at ? [{ hash, at }] : [])); | |
| 504 | 504 | } | |
| 505 | 505 | ||
| 506 | + | /** | |
| 507 | + | * Whether a read of the default branch reaches back far enough to place | |
| 508 | + | * every push, newest first: it holds the oldest push's `before`, or it is | |
| 509 | + | * the whole branch. A shorter read than `asked` is the whole branch. | |
| 510 | + | */ | |
| 511 | + | export function historyCovers(history: { hash: string }[], pushes: PushRecord[], asked: number): boolean { | |
| 512 | + | if (history.length < asked) return true; | |
| 513 | + | const oldest = pushes.at(-1)?.data.before; | |
| 514 | + | return oldest != null && history.some((commit) => commit.hash === oldest); | |
| 515 | + | } | |
| 516 | + | ||
| 506 | 517 | /** A merged pull request, as the week counts it. */ | |
| 507 | 518 | export type Merged = { | |
| 508 | 519 | repo: RepoPath; |
| 57 | 57 | usd, | |
| 58 | 58 | waitingRows, | |
| 59 | 59 | weekOf, | |
| 60 | + | historyCovers, | |
| 60 | 61 | who, | |
| 61 | 62 | whyFor, | |
| 62 | 63 | withConfidence, | |
| ⋯ | |||
| 161 | 162 | */ | |
| 162 | 163 | const PUSH_PAGE = 200; | |
| 163 | 164 | const PUSH_PAGES = 5; | |
| 164 | − | /** Commits of the default branch read once, to place each push's commits. */ | |
| 165 | + | /** The first page is smaller: most projects push far less in two weeks. */ | |
| 166 | + | const FIRST_PUSH_PAGE = 50; | |
| 167 | + | /** | |
| 168 | + | * Commits of the default branch read to place each push's commits: a | |
| 169 | + | * short read first, which covers most projects' fortnight, and the longer | |
| 170 | + | * one only when it does not reach the oldest push. | |
| 171 | + | */ | |
| 172 | + | const HISTORY_FIRST = 200; | |
| 165 | 173 | const HISTORY_READ = 1000; | |
| 166 | 174 | ||
| 167 | 175 | /** | |
| ⋯ | |||
| 174 | 182 | const pushes: G1tEvent<"git.push">[] = []; | |
| 175 | 183 | let before: string | undefined; | |
| 176 | 184 | for (let page = 0; page < PUSH_PAGES; page++) { | |
| 177 | − | const batch = await eventLog.list({ repoId: repo.id, types: ["git.push"], limit: PUSH_PAGE, ...(before ? { before } : {}) }); | |
| 185 | + | const limit = page === 0 ? FIRST_PUSH_PAGE : PUSH_PAGE; | |
| 186 | + | const batch = await eventLog.list({ repoId: repo.id, types: ["git.push"], limit, ...(before ? { before } : {}) }); | |
| 178 | 187 | for (const event of batch) { | |
| 179 | 188 | if (event.type === "git.push" && event.data.defaultBranch && Date.parse(event.time) >= since) pushes.push(event as G1tEvent<"git.push">); | |
| 180 | 189 | } | |
| 181 | 190 | const oldest = batch.at(-1); | |
| 182 | − | if (batch.length < PUSH_PAGE || !oldest || Date.parse(oldest.time) < since) break; | |
| 191 | + | if (batch.length < limit || !oldest || Date.parse(oldest.time) < since) break; | |
| 183 | 192 | before = oldest.id; | |
| 184 | 193 | } | |
| 185 | 194 | if (pushes.length === 0) return []; | |
| 186 | − | // One read of the branch from the newest push back; each push brought | |
| 195 | + | // A read of the branch from the newest push back; each push brought | |
| 187 | 196 | // what lies between its `after` and its `before` in that history. | |
| 188 | 197 | const path = { namespace: repo.namespace, name: repo.name }; | |
| 189 | − | const history = await reposApi.log(path, viewer, pushes[0].data.after, HISTORY_READ).catch(() => null); | |
| 198 | + | const read = (limit: number) => reposApi.log(path, viewer, pushes[0].data.after, limit).catch(() => null); | |
| 199 | + | let history = await read(HISTORY_FIRST); | |
| 200 | + | if (history?.ok && !historyCovers(history.value, pushes, HISTORY_FIRST)) history = await read(HISTORY_READ); | |
| 190 | 201 | if (!history?.ok) return []; | |
| 191 | 202 | // g1t's own commits are told by their author address, never by name. | |
| 192 | 203 | return placePushes(history.value, pushes, (commit) => G1T_COMMIT_EMAILS.has(normalizeEmail(commit.author.email))); | |
| 39 | 39 | import { DropdownMenu, DropdownMenuContent, DropdownMenuItem, DropdownMenuTrigger } from "../../components/ui/dropdown-menu"; | |
| 40 | 40 | import { Hint } from "../../components/ui/hint"; | |
| 41 | 41 | import { searchLog } from "../../lib/log-lines"; | |
| 42 | − | import { listArtifacts } from "../../lib/artifacts.server"; | |
| 43 | − | import { expiresIn, formatBytes } from "../../lib/artifacts"; | |
| 42 | + | import { listArtifacts, withLegacyArtifacts } from "../../lib/artifacts.server"; | |
| 43 | + | import { expiresIn, formatBytes, legacyArtifactsWorthAsking } from "../../lib/artifacts"; | |
| 44 | 44 | import { actions } from "../../lib/services.server"; | |
| 45 | 45 | import { assertSameOrigin, getViewer, requireUser, roleIn, unwrap } from "../../lib/session.server"; | |
| 46 | 46 | import { accessTo, refusal } from "../../lib/access.server"; | |
| ⋯ | |||
| 57 | 57 | // An earlier attempt, when one is asked for. | |
| 58 | 58 | const asked = Number(new URL(request.url).searchParams.get("attempt") ?? ""); | |
| 59 | 59 | const attempt = Number.isInteger(asked) && asked > 0 ? asked : undefined; | |
| 60 | − | const [detail, artifacts, summaries] = await Promise.all([ | |
| 60 | + | const [detail, kept, summaries] = await Promise.all([ | |
| 61 | 61 | actions.run(repo, viewer, params.id, attempt).then(unwrap), | |
| 62 | 62 | listArtifacts(repo, viewer, params.id).catch(() => []), | |
| 63 | 63 | actions | |
| ⋯ | |||
| 65 | 65 | .then((found) => (found.ok ? found.value : [])) | |
| 66 | 66 | .catch((): JobSummary[] => []), | |
| 67 | 67 | ]); | |
| 68 | + | // Artifacts an older runner kept in KV: asked for only once the service | |
| 69 | + | // has shown the viewer the run, and never while it is still going. | |
| 70 | + | const artifacts = legacyArtifactsWorthAsking(detail.run) ? await withLegacyArtifacts(kept, params.id) : kept; | |
| 68 | 71 | // Cancelling and re-running need Write. | |
| 69 | 72 | return { | |
| 70 | 73 | detail, | |
| 85 | 85 | ||
| 86 | 86 | /// Whether it fires in the minute starting at `ms` since the epoch, UTC. | |
| 87 | 87 | pub fn fires_at(&self, ms: u64) -> bool { | |
| 88 | + | self.minutes[(ms / 60_000 % 60) as usize] && self.hour_and_date_at(ms) | |
| 89 | + | } | |
| 90 | + | ||
| 91 | + | /// Whether the minute starting at `ms` is in its hours, days and | |
| 92 | + | /// months, whatever its minute field says. | |
| 93 | + | fn hour_and_date_at(&self, ms: u64) -> bool { | |
| 88 | 94 | let minutes_total = ms / 60_000; | |
| 89 | − | let minute = (minutes_total % 60) as usize; | |
| 90 | 95 | let hour = (minutes_total / 60 % 24) as usize; | |
| 91 | 96 | let days_since_epoch = (minutes_total / 60 / 24) as i64; | |
| 92 | 97 | // 1970-01-01 was a Thursday. | |
| ⋯ | |||
| 98 | 103 | (true, true) => day_ok || weekday_ok, | |
| 99 | 104 | _ => day_ok && weekday_ok, | |
| 100 | 105 | }; | |
| 101 | − | self.minutes[minute] && self.hours[hour] && self.months[month as usize] && date_ok | |
| 106 | + | self.hours[hour] && self.months[month as usize] && date_ok | |
| 102 | 107 | } | |
| 108 | + | ||
| 109 | + | /// Whether its minutes come closer together than | |
| 110 | + | /// [`MIN_INTERVAL_MINUTES`], counting round the hour. | |
| 111 | + | pub fn too_frequent(&self) -> bool { | |
| 112 | + | let set: Vec<usize> = (0..60).filter(|&m| self.minutes[m]).collect(); | |
| 113 | + | let Some(&first) = set.first() else { return false }; | |
| 114 | + | let wrap = first + 60 - set[set.len() - 1]; | |
| 115 | + | set.windows(2).map(|pair| pair[1] - pair[0]).chain([wrap]).any(|gap| gap < MIN_INTERVAL_MINUTES as usize) | |
| 116 | + | } | |
| 117 | + | ||
| 118 | + | /// Whether a workflow on this schedule runs in the minute starting at | |
| 119 | + | /// `ms`. A schedule no more frequent than every | |
| 120 | + | /// [`MIN_INTERVAL_MINUTES`] runs when it fires. A more frequent one | |
| 121 | + | /// runs on the five-minute marks: at a mark that is itself in its | |
| 122 | + | /// hours, days and months, when it fired in the five minutes the mark | |
| 123 | + | /// ends. At most every five minutes, and never outside its own hours. | |
| 124 | + | pub fn runs_at(&self, ms: u64) -> bool { | |
| 125 | + | if !self.too_frequent() { | |
| 126 | + | return self.fires_at(ms); | |
| 127 | + | } | |
| 128 | + | let minute = ms / 60_000; | |
| 129 | + | minute % MIN_INTERVAL_MINUTES == 0 | |
| 130 | + | && self.hour_and_date_at(ms) | |
| 131 | + | && (0..MIN_INTERVAL_MINUTES).any(|back| minute >= back && self.fires_at((minute - back) * 60_000)) | |
| 132 | + | } | |
| 103 | 133 | } | |
| 104 | 134 | ||
| 135 | + | /// The shortest interval a workflow's schedule runs at, in minutes. | |
| 136 | + | pub const MIN_INTERVAL_MINUTES: u64 = 5; | |
| 137 | + | ||
| 105 | 138 | /// The date of a day counted from 1970-01-01 (Howard Hinnant's algorithm). | |
| 106 | 139 | fn civil_from_days(days: i64) -> (i64, u32, u32) { | |
| 107 | 140 | let z = days + 719_468; | |
| ⋯ | |||
| 168 | 201 | } | |
| 169 | 202 | ||
| 170 | 203 | #[test] | |
| 204 | + | fn schedules_run_at_most_every_five_minutes() { | |
| 205 | + | let runs = |text: &str| -> Vec<u64> { | |
| 206 | + | let schedule = Schedule::parse(text).unwrap(); | |
| 207 | + | (0..30).filter(|&m| schedule.runs_at(at(MONDAY, 3, m))).collect() | |
| 208 | + | }; | |
| 209 | + | assert_eq!(runs("* * * * *"), [0, 5, 10, 15, 20, 25]); | |
| 210 | + | assert_eq!(runs("*/2 * * * *"), [0, 5, 10, 15, 20, 25]); | |
| 211 | + | // Every five minutes or less often: exactly when it fires. | |
| 212 | + | assert_eq!(runs("*/5 * * * *"), [0, 5, 10, 15, 20, 25]); | |
| 213 | + | assert_eq!(runs("7,17 * * * *"), [7, 17]); | |
| 214 | + | assert_eq!(runs("*/10 * * * *"), [0, 10, 20]); | |
| 215 | + | // Two minutes close together: the later one waits for the mark. | |
| 216 | + | assert_eq!(runs("0,3 * * * *"), [0, 5]); | |
| 217 | + | // Close across the hour counts too. | |
| 218 | + | assert!(Schedule::parse("2,58 * * * *").unwrap().too_frequent()); | |
| 219 | + | assert!(!Schedule::parse("0 9 * * mon").unwrap().too_frequent()); | |
| 220 | + | // Every minute of one hour: its own marks, 09:00 to 09:55, and | |
| 221 | + | // nothing after. | |
| 222 | + | let nine = Schedule::parse("* 9 * * *").unwrap(); | |
| 223 | + | let day: Vec<(u64, u64)> = | |
| 224 | + | (0..24 * 60).filter(|&m| nine.runs_at(at(MONDAY, m / 60, m % 60))).map(|m| (m / 60, m % 60)).collect(); | |
| 225 | + | assert_eq!(day, (0..12).map(|i| (9, i * 5)).collect::<Vec<_>>()); | |
| 226 | + | assert!(!nine.runs_at(at(MONDAY, 10, 0))); | |
| 227 | + | // Every minute of Mondays: the last run is 23:55, none on Tuesday. | |
| 228 | + | let mondays = Schedule::parse("*/1 * * * 1").unwrap(); | |
| 229 | + | assert!(mondays.runs_at(at(MONDAY, 0, 0))); | |
| 230 | + | assert!(mondays.runs_at(at(MONDAY, 23, 55))); | |
| 231 | + | assert!(!mondays.runs_at(at(MONDAY + 1, 0, 0))); | |
| 232 | + | assert!(!(0..24 * 60).any(|m| mondays.runs_at(at(MONDAY + 1, m / 60, m % 60)))); | |
| 233 | + | // Close minutes: one run for each window they fall in, at its mark. | |
| 234 | + | let close = Schedule::parse("0,3 * * * *").unwrap(); | |
| 235 | + | let hour: Vec<u64> = (0..60).filter(|&m| close.runs_at(at(MONDAY, 3, m))).collect(); | |
| 236 | + | assert_eq!(hour, [0, 5]); | |
| 237 | + | assert!(close.runs_at(at(MONDAY, 4, 0))); | |
| 238 | + | } | |
| 239 | + | ||
| 240 | + | #[test] | |
| 171 | 241 | fn nonsense_is_refused() { | |
| 172 | 242 | for text in ["", "* * * *", "61 * * * *", "* * * * funday", "*/0 * * * *", "5-1 * * * *"] { | |
| 173 | 243 | assert!(Schedule::parse(text).is_err(), "{text}"); | |
| 121 | 121 | buildSeconds: number; | |
| 122 | 122 | /** Charged so far this month for builds, in millionths of a dollar. */ | |
| 123 | 123 | buildMicros: number; | |
| 124 | − | /** RFC 3339: when requests and CPU time were last counted. */ | |
| 124 | + | /** RFC 3339: when requests and CPU time last went up. */ | |
| 125 | 125 | countedAt: string | null; | |
| 126 | 126 | }; | |
| 127 | 127 |
| 1 | + | // node --test "scripts/deploy/*.test.mjs" (npm run test:deploy) | |
| 2 | + | // | |
| 3 | + | // The queries that run on a timer or on every page load, planned by SQLite | |
| 4 | + | // against each service's migrations as they apply in order: each must be a | |
| 5 | + | // search of an index, never a scan of the table. D1 is SQLite, so its | |
| 6 | + | // planner chooses the same way. | |
| 7 | + | import assert from "node:assert/strict"; | |
| 8 | + | import { readdirSync, readFileSync } from "node:fs"; | |
| 9 | + | import { join } from "node:path"; | |
| 10 | + | import { DatabaseSync } from "node:sqlite"; | |
| 11 | + | import { test } from "node:test"; | |
| 12 | + | ||
| 13 | + | import { ROOT } from "./stack.mjs"; | |
| 14 | + | ||
| 15 | + | /** A database with every one of `service`'s migrations applied, in order. */ | |
| 16 | + | function migrated(service) { | |
| 17 | + | const db = new DatabaseSync(":memory:"); | |
| 18 | + | const dir = join(ROOT, "services", service, "migrations"); | |
| 19 | + | for (const file of readdirSync(dir).filter((name) => name.endsWith(".sql")).sort()) { | |
| 20 | + | db.exec(readFileSync(join(dir, file), "utf8")); | |
| 21 | + | } | |
| 22 | + | return db; | |
| 23 | + | } | |
| 24 | + | ||
| 25 | + | /** How SQLite would run `sql`, one line per step. */ | |
| 26 | + | function plan(db, sql, ...params) { | |
| 27 | + | return db | |
| 28 | + | .prepare(`EXPLAIN QUERY PLAN ${sql}`) | |
| 29 | + | .all(...params) | |
| 30 | + | .map((row) => row.detail); | |
| 31 | + | } | |
| 32 | + | ||
| 33 | + | function assertSearches(steps, index) { | |
| 34 | + | assert.ok( | |
| 35 | + | steps.some((step) => step.includes(`INDEX ${index}`)), | |
| 36 | + | `expected a search of ${index}, planned:\n ${steps.join("\n ")}`, | |
| 37 | + | ); | |
| 38 | + | } | |
| 39 | + | ||
| 40 | + | test("webhooks: the hourly purge deletes old deliveries by time", () => { | |
| 41 | + | const steps = plan( | |
| 42 | + | migrated("webhooks"), | |
| 43 | + | "DELETE FROM deliveries WHERE rowid IN (SELECT rowid FROM deliveries WHERE created_at < ? ORDER BY created_at LIMIT ?)", | |
| 44 | + | "2026-09-24T00:00:00.000Z", | |
| 45 | + | 1000, | |
| 46 | + | ); | |
| 47 | + | assertSearches(steps, "deliveries_by_time"); | |
| 48 | + | }); | |
| 49 | + | ||
| 50 | + | test("billing: this month's users are read from an index, not the ledger", () => { | |
| 51 | + | const steps = plan( | |
| 52 | + | migrated("billing"), | |
| 53 | + | "SELECT DISTINCT workspace FROM ledger WHERE kind = 'usage' AND created_at >= ?", | |
| 54 | + | "2026-10-01", | |
| 55 | + | ); | |
| 56 | + | assertSearches(steps, "ledger_usage_by_time"); | |
| 57 | + | assert.ok(steps.some((step) => step.includes("COVERING INDEX")), steps.join("\n")); | |
| 58 | + | }); | |
| 59 | + | ||
| 60 | + | test("events: a repository's events of one type, newest first, page by page", () => { | |
| 61 | + | const db = migrated("events"); | |
| 62 | + | const first = plan(db, "SELECT * FROM events WHERE repo_id = ? AND type IN (?) ORDER BY id DESC LIMIT ?", "rep_1", "git.push", 50); | |
| 63 | + | assertSearches(first, "events_repo_type"); | |
| 64 | + | const next = plan(db, "SELECT * FROM events WHERE repo_id = ? AND type IN (?) AND id < ? ORDER BY id DESC LIMIT ?", "rep_1", "git.push", "evt_9", 200); | |
| 65 | + | assertSearches(next, "events_repo_type"); | |
| 66 | + | // Read in order: no sort of every match. | |
| 67 | + | assert.ok(!next.some((step) => step.includes("TEMP B-TREE")), next.join("\n")); | |
| 68 | + | }); | |
| 69 | + | ||
| 70 | + | test("actions: a self-hosted runner's poll finds queued jobs by namespace, whatever its case", () => { | |
| 71 | + | const steps = plan( | |
| 72 | + | migrated("actions"), | |
| 73 | + | `SELECT jobs.*, runs.repo AS run_repo FROM jobs JOIN runs ON runs.id = jobs.run_id | |
| 74 | + | WHERE jobs.status = 'queued' AND jobs.labels IS NOT NULL AND lower(jobs.namespace) = ? | |
| 75 | + | ORDER BY jobs.rowid LIMIT 50`, | |
| 76 | + | "acme", | |
| 77 | + | ); | |
| 78 | + | assertSearches(steps, "jobs_self_hosted_lower"); | |
| 79 | + | }); |
| 1 | + | -- What a self-hosted runner can take, asked on every poll: queued jobs with | |
| 2 | + | -- labels in its workspace, matched without regard to case. By the | |
| 3 | + | -- lowercased namespace, so the match is a range of the index rather than | |
| 4 | + | -- every labelled job. | |
| 5 | + | CREATE INDEX IF NOT EXISTS jobs_self_hosted_lower ON jobs (lower(namespace), status) WHERE labels IS NOT NULL; |
| 1 | + | -- Schedules pause in a repository with no push for 60 days, and wait an | |
| 2 | + | -- hour after billing refused to start one of the workflow's scheduled jobs. | |
| 3 | + | ||
| 4 | + | -- When the repository last had a push, at a day's resolution; written on | |
| 5 | + | -- pushes to repositories with a scheduled workflow. Workflows already here | |
| 6 | + | -- start counting now. | |
| 7 | + | ALTER TABLE workflows ADD COLUMN pushed_at TEXT; | |
| 8 | + | UPDATE workflows SET pushed_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now') WHERE crons != '[]'; | |
| 9 | + | ||
| 10 | + | -- Until when its schedule waits: set when billing refused a job of one of | |
| 11 | + | -- its scheduled runs. | |
| 12 | + | ALTER TABLE workflows ADD COLUMN schedule_refused_until TEXT; |
| 59 | 59 | pub const SELF_HOSTED_MAX_TIMEOUT_MINUTES: u32 = 24 * 60; | |
| 60 | 60 | /// A running job that has said nothing for this long is taken as lost. | |
| 61 | 61 | pub const SILENT_MS: u64 = 10 * 60 * 1000; | |
| 62 | + | /// Schedules pause in a repository that has had no push for this long; the | |
| 63 | + | /// next push resumes them. | |
| 64 | + | pub const SCHEDULE_IDLE_MS: u64 = 60 * 24 * 60 * 60 * 1000; | |
| 65 | + | /// How long a workflow's schedule waits after billing refused to start a | |
| 66 | + | /// job of one of its scheduled runs, before it tries again. | |
| 67 | + | pub const SCHEDULE_REFUSED_MS: u64 = 60 * 60 * 1000; | |
| 62 | 68 | pub const SITE: &str = "https://g1t.sh"; | |
| 63 | 69 | pub const API: &str = "https://api.g1t.sh"; | |
| 64 | 70 |
| 1348 | 1348 | .await | |
| 1349 | 1349 | .unwrap_or_else(|error| fail(FailureCode::Conflict, format!("The runner could not be reached: {error}"))); | |
| 1350 | 1350 | if let Outcome::Fail(refused) = started { | |
| 1351 | + | // Billing would not pay for a scheduled run's job: the | |
| 1352 | + | // schedule waits rather than making runs only to refuse them. | |
| 1353 | + | if refused.code == FailureCode::PaymentRequired | |
| 1354 | + | && let Some(run) = run.as_ref().filter(|run| run.event == "schedule") | |
| 1355 | + | { | |
| 1356 | + | self.db | |
| 1357 | + | .prepare("UPDATE workflows SET schedule_refused_until = ? WHERE id = ?") | |
| 1358 | + | .bind(&[rfc3339(now_ms() + crate::SCHEDULE_REFUSED_MS).into(), run.workflow_id.as_str().into()])? | |
| 1359 | + | .run() | |
| 1360 | + | .await?; | |
| 1361 | + | } | |
| 1351 | 1362 | Box::pin(self.finish_job(&job.id, "failure", Some(&refused.message), None)).await?; | |
| 1352 | 1363 | } else { | |
| 1353 | 1364 | self.job_started(&job.id).await?; |
| 36 | 36 | const POLL_EVERY_MS: u64 = 2_500; | |
| 37 | 37 | /// A runner that has been offline this long is removed, as on GitHub. | |
| 38 | 38 | const FORGET_OFFLINE_MS: u64 = 14 * 24 * 60 * 60 * 1000; | |
| 39 | + | /// How often a poll writes down that a runner, or the work it holds, is | |
| 40 | + | /// still there. Polls come every twenty seconds or so; a write each time | |
| 41 | + | /// would be most of the service's writes for nothing. Well inside | |
| 42 | + | /// [`model::ONLINE_WITHIN_MS`] and [`crate::SILENT_MS`]. | |
| 43 | + | const SEEN_RESOLUTION_MS: u64 = 60 * 1000; | |
| 44 | + | ||
| 45 | + | /// Whether `seen` is old enough, at `at`, to be written down again. | |
| 46 | + | fn seen_is_stale(seen: Option<&str>, at: u64) -> bool { | |
| 47 | + | seen.is_none_or(|seen| seen < rfc3339(at.saturating_sub(SEEN_RESOLUTION_MS)).as_str()) | |
| 48 | + | } | |
| 39 | 49 | ||
| 40 | 50 | #[derive(Clone, Debug, Deserialize)] | |
| 41 | 51 | pub struct RunnerRow { | |
| ⋯ | |||
| 124 | 134 | pub runner_id: Option<String>, | |
| 125 | 135 | pub runner_name: Option<String>, | |
| 126 | 136 | pub started_at: Option<String>, | |
| 137 | + | #[serde(default)] | |
| 138 | + | pub seen_at: Option<String>, | |
| 127 | 139 | } | |
| 128 | 140 | ||
| 129 | 141 | /// Who may do what with a set of runners: the workspace, and the | |
| ⋯ | |||
| 912 | 924 | for id in a.running.iter().take(10) { | |
| 913 | 925 | if let Some(job) = self.db.prepare("SELECT * FROM jobs WHERE id = ?").bind(&[id.as_str().into()])?.first::<JobRow>(None).await? { | |
| 914 | 926 | if job.status == "in_progress" && job.runner_id.as_deref() == Some(row.id.as_str()) { | |
| 915 | − | self.db.prepare("UPDATE jobs SET seen_at = ? WHERE id = ?").bind(&[seen.as_str().into(), id.as_str().into()])?.run().await?; | |
| 927 | + | if seen_is_stale(job.seen_at.as_deref(), at) { | |
| 928 | + | self.db.prepare("UPDATE jobs SET seen_at = ? WHERE id = ?").bind(&[seen.as_str().into(), id.as_str().into()])?.run().await?; | |
| 929 | + | } | |
| 916 | 930 | active = Some((id.clone(), "workflow")); | |
| 917 | 931 | } else { | |
| 918 | 932 | cancel.push(id.clone()); | |
| ⋯ | |||
| 921 | 935 | } | |
| 922 | 936 | match self.db.prepare("SELECT * FROM runner_tasks WHERE id = ?").bind(&[id.as_str().into()])?.first::<TaskRow>(None).await? { | |
| 923 | 937 | Some(task) if task.status == "in_progress" && task.runner_id.as_deref() == Some(row.id.as_str()) => { | |
| 924 | − | self.db.prepare("UPDATE runner_tasks SET seen_at = ? WHERE id = ?").bind(&[seen.as_str().into(), id.as_str().into()])?.run().await?; | |
| 938 | + | if seen_is_stale(task.seen_at.as_deref(), at) { | |
| 939 | + | self.db.prepare("UPDATE runner_tasks SET seen_at = ? WHERE id = ?").bind(&[seen.as_str().into(), id.as_str().into()])?.run().await?; | |
| 940 | + | } | |
| 925 | 941 | active = Some((id.clone(), "agent")); | |
| 926 | 942 | } | |
| 927 | 943 | _ => cancel.push(id.clone()), | |
| 928 | 944 | } | |
| 929 | 945 | } | |
| 930 | − | self.db | |
| 931 | − | .prepare("UPDATE runners SET last_seen_at = ?, version = ?, work_id = ?, work_kind = ? WHERE id = ?") | |
| 932 | − | .bind(&[ | |
| 933 | − | seen.as_str().into(), | |
| 934 | − | (if a.version.is_empty() { row.version.clone() } else { a.version.chars().take(40).collect() }).into(), | |
| 935 | − | optional(active.as_ref().map(|(id, _)| id.as_str())), | |
| 936 | − | optional(active.as_ref().map(|(_, kind)| *kind)), | |
| 937 | − | row.id.as_str().into(), | |
| 938 | − | ])? | |
| 939 | − | .run() | |
| 940 | − | .await?; | |
| 946 | + | // Written when something changed, or once a minute to keep it online. | |
| 947 | + | let version: String = if a.version.is_empty() { row.version.clone() } else { a.version.chars().take(40).collect() }; | |
| 948 | + | let work_id = active.as_ref().map(|(id, _)| id.as_str()); | |
| 949 | + | let work_kind = active.as_ref().map(|(_, kind)| *kind); | |
| 950 | + | let changed = version != row.version || work_id != row.work_id.as_deref() || work_kind != row.work_kind.as_deref(); | |
| 951 | + | if changed || seen_is_stale(row.last_seen_at.as_deref(), at) { | |
| 952 | + | self.db | |
| 953 | + | .prepare("UPDATE runners SET last_seen_at = ?, version = ?, work_id = ?, work_kind = ? WHERE id = ?") | |
| 954 | + | .bind(&[seen.as_str().into(), version.into(), optional(work_id), optional(work_kind), row.id.as_str().into()])? | |
| 955 | + | .run() | |
| 956 | + | .await?; | |
| 957 | + | } | |
| 941 | 958 | let mut poll = Poll { assignment: None, cancel, credential, removed: false }; | |
| 942 | 959 | // Busy, or an ephemeral runner that has had its one job. | |
| 943 | 960 | if active.is_some() || !a.running.is_empty() { | |
| ⋯ | |||
| 1543 | 1560 | } | |
| 1544 | 1561 | ||
| 1545 | 1562 | #[test] | |
| 1563 | + | fn being_seen_is_written_once_a_minute() { | |
| 1564 | + | let at = 2_000_000_000_000; | |
| 1565 | + | assert!(seen_is_stale(None, at)); | |
| 1566 | + | assert!(!seen_is_stale(Some(&rfc3339(at - 59_000)), at)); | |
| 1567 | + | assert!(seen_is_stale(Some(&rfc3339(at - 61_000)), at)); | |
| 1568 | + | // A runner written down at the last moment is still online until the | |
| 1569 | + | // next write, a poll later. | |
| 1570 | + | assert!(SEEN_RESOLUTION_MS + model::MAX_POLL_WAIT_MS < model::ONLINE_WITHIN_MS); | |
| 1571 | + | assert!(SEEN_RESOLUTION_MS < crate::SILENT_MS); | |
| 1572 | + | } | |
| 1573 | + | ||
| 1574 | + | #[test] | |
| 1546 | 1575 | fn a_group_lets_its_repositories_use_it() { | |
| 1547 | 1576 | assert!(group(&[]).allows("acme/web")); | |
| 1548 | 1577 | assert!(group(&["web", "api"]).allows("acme/WEB")); | |
| 11 | 11 | use g1t_contracts::repos::{Commit, CompareArgs, Comparison, LogArgs, Repo, RepoPath}; | |
| 12 | 12 | use g1t_contracts::work::{IssueDetail, PullDetail, ViewArgs}; | |
| 13 | 13 | use g1t_contracts::{FailureCode, Outcome, User, new_id}; | |
| 14 | + | use g1t_contracts::time::rfc3339; | |
| 14 | 15 | use g1t_kit::now_ms; | |
| 15 | 16 | use serde_json::{Map, Value, json}; | |
| 16 | 17 | use worker::Result; | |
| ⋯ | |||
| 511 | 512 | }) | |
| 512 | 513 | } | |
| 513 | 514 | ||
| 515 | + | /// A push keeps the repository's schedules going (`run_schedules`). | |
| 516 | + | /// Written at most once a day, and only on scheduled workflows. | |
| 517 | + | async fn note_push(&self, repo_id: &str) -> Result<()> { | |
| 518 | + | let at = now_ms(); | |
| 519 | + | self.db | |
| 520 | + | .prepare("UPDATE workflows SET pushed_at = ? WHERE repo_id = ? AND crons != '[]' AND (pushed_at IS NULL OR pushed_at < ?)") | |
| 521 | + | .bind(&[rfc3339(at).into(), repo_id.into(), rfc3339(at.saturating_sub(24 * 60 * 60 * 1000)).into()])? | |
| 522 | + | .run() | |
| 523 | + | .await?; | |
| 524 | + | Ok(()) | |
| 525 | + | } | |
| 526 | + | ||
| 514 | 527 | pub async fn on_event(&self, event: &Event) -> Result<()> { | |
| 515 | 528 | let Some(repo_id) = event.repo_id.as_deref() else { return Ok(()) }; | |
| 516 | 529 | let mut mapped = github_events(&event.kind); | |
| ⋯ | |||
| 519 | 532 | mapped.push(("create", None)); | |
| 520 | 533 | } | |
| 521 | 534 | let pushed_default = event.kind == "git.push" && event.data["defaultBranch"].as_bool() == Some(true); | |
| 535 | + | if event.kind == "git.push" { | |
| 536 | + | self.note_push(repo_id).await?; | |
| 537 | + | } | |
| 522 | 538 | if mapped.is_empty() && !pushed_default { | |
| 523 | 539 | return Ok(()); | |
| 524 | 540 | } | |
| ⋯ | |||
| 717 | 733 | self.record_failed_run(&row, subject.git_ref.as_str(), &subject.sha, event_key, actor_id, sender, problem).await | |
| 718 | 734 | } | |
| 719 | 735 | ||
| 720 | − | /// Scheduled workflows whose cron fires this minute, on the default branch. | |
| 736 | + | /// Scheduled workflows that run this minute, on the default branch: at | |
| 737 | + | /// most every five minutes (`cron::Schedule::runs_at`). Not in a | |
| 738 | + | /// repository with no push for [`crate::SCHEDULE_IDLE_MS`], nor for an | |
| 739 | + | /// hour after billing refused one of the workflow's scheduled jobs | |
| 740 | + | /// ([`crate::SCHEDULE_REFUSED_MS`]): no run is made only to be refused. | |
| 721 | 741 | pub async fn run_schedules(&self, minute: u64) -> Result<()> { | |
| 722 | 742 | let rows = self | |
| 723 | 743 | .db | |
| 724 | − | .prepare("SELECT * FROM workflows WHERE state = 'active' AND crons != '[]' AND error IS NULL") | |
| 744 | + | .prepare( | |
| 745 | + | "SELECT * FROM workflows WHERE state = 'active' AND crons != '[]' AND error IS NULL | |
| 746 | + | AND COALESCE(pushed_at, updated_at) >= ? AND (schedule_refused_until IS NULL OR schedule_refused_until <= ?)", | |
| 747 | + | ) | |
| 748 | + | .bind(&[rfc3339(minute.saturating_sub(crate::SCHEDULE_IDLE_MS)).into(), rfc3339(minute).into()])? | |
| 725 | 749 | .all() | |
| 726 | 750 | .await? | |
| 727 | 751 | .results::<WorkflowRow>()?; | |
| 728 | 752 | for row in rows { | |
| 729 | 753 | let crons: Vec<String> = serde_json::from_str(&row.crons).unwrap_or_default(); | |
| 730 | − | let Some(cron) = crons.iter().find(|cron| g1t_actions::cron::Schedule::parse(cron).is_ok_and(|s| s.fires_at(minute))) else { | |
| 754 | + | let Some(cron) = crons.iter().find(|cron| g1t_actions::cron::Schedule::parse(cron).is_ok_and(|s| s.runs_at(minute))) else { | |
| 731 | 755 | continue; | |
| 732 | 756 | }; | |
| 733 | 757 | let Ok(workflow) = workflow::parse(&row.source) else { continue }; | |
| 1 | + | -- This month's users, read every fifteen minutes for autopay and the limit | |
| 2 | + | -- warnings: usage entries from a date on, by kind and time, and the | |
| 3 | + | -- workspace in the index so the read never touches the ledger itself. | |
| 4 | + | CREATE INDEX IF NOT EXISTS ledger_usage_by_time ON ledger (kind, created_at, workspace); |
| 1154 | 1154 | } | |
| 1155 | 1155 | } | |
| 1156 | 1156 | ||
| 1157 | − | #[event(scheduled)] | |
| 1158 | − | async fn scheduled(event: ScheduledEvent, env: Env, _ctx: ScheduleContext) { | |
| 1159 | − | let Ok(billing) = Billing::from_env(&env) else { | |
| 1160 | − | return; | |
| 1161 | − | }; | |
| 1162 | − | let keeper = keeper::Keeper::from_env(&env); | |
| 1163 | − | if let Err(error) = billing.settle_runs(&keeper).await { | |
| 1157 | + | /// The fifteen-minute steps, each on its own: one failing never skips the | |
| 1158 | + | /// rest. | |
| 1159 | + | async fn tick(billing: &Billing, env: &Env, keeper: &keeper::Keeper) { | |
| 1160 | + | if let Err(error) = billing.settle_runs(keeper).await { | |
| 1164 | 1161 | worker::console_error!("settling runs failed: {error}"); | |
| 1165 | 1162 | } | |
| 1166 | 1163 | if let Err(error) = billing.settle_own_runs().await { | |
| ⋯ | |||
| 1170 | 1167 | match billing.replay_events().await { | |
| 1171 | 1168 | Ok(done) => worker::console_log!("stripe replay: {done}"), | |
| 1172 | 1169 | Err(error) => worker::console_error!("replaying Stripe events failed: {error}"), | |
| 1173 | − | } | |
| 1174 | − | if let Err(error) = billing.autopay().await { | |
| 1175 | − | worker::console_error!("paying at the limit failed: {error}"); | |
| 1176 | 1170 | } | |
| 1171 | + | // This month's users, read once for autopay and the warnings. | |
| 1172 | + | let users = if billing.stripe.is_some() { | |
| 1173 | + | match billing.month_users().await { | |
| 1174 | + | Ok(users) => Some(users), | |
| 1175 | + | Err(error) => { | |
| 1176 | + | worker::console_error!("reading this month's users failed: {error}"); | |
| 1177 | + | None | |
| 1178 | + | } | |
| 1179 | + | } | |
| 1180 | + | } else { | |
| 1181 | + | None | |
| 1182 | + | }; | |
| 1183 | + | if let Some(users) = &users | |
| 1184 | + | && let Err(error) = billing.autopay(users).await { | |
| 1185 | + | worker::console_error!("paying at the limit failed: {error}"); | |
| 1186 | + | } | |
| 1177 | 1187 | // AI credit below a workspace's auto-reload threshold, reloaded (ai.rs). | |
| 1178 | 1188 | match billing.reload_ai_credit().await { | |
| 1179 | 1189 | Ok(0) => {} | |
| ⋯ | |||
| 1191 | 1201 | if let Err(error) = billing.invoice_enterprises().await { | |
| 1192 | 1202 | worker::console_error!("invoicing enterprises failed: {error}"); | |
| 1193 | 1203 | } | |
| 1204 | + | if let Some(users) = &users | |
| 1205 | + | && let Ok(identity) = env.service("IDENTITY") | |
| 1206 | + | && let Err(error) = billing.warn_limits(&identity, users).await { | |
| 1207 | + | worker::console_error!("warning owners failed: {error}"); | |
| 1208 | + | } | |
| 1209 | + | } | |
| 1210 | + | ||
| 1211 | + | /// Every fifteen minutes: runs settled, Stripe replayed, cards charged at | |
| 1212 | + | /// the limit, months closed, owners warned. Once a day, at | |
| 1213 | + | /// [`keeper::DAILY`], in an invocation of its own: the slower checks | |
| 1214 | + | /// against Stripe and Cloudflare, and measuring. | |
| 1215 | + | /// | |
| 1216 | + | /// Each step logs its own failure and the next still runs. The spend | |
| 1217 | + | /// breaker and budget alerts go first: a step that runs out of CPU ends | |
| 1218 | + | /// the invocation, and must never take the alerts with it. | |
| 1219 | + | #[event(scheduled)] | |
| 1220 | + | async fn scheduled(event: ScheduledEvent, env: Env, _ctx: ScheduleContext) { | |
| 1221 | + | let Ok(billing) = Billing::from_env(&env) else { | |
| 1222 | + | return; | |
| 1223 | + | }; | |
| 1194 | 1224 | // Comped budgets' alerts, and a tripped breaker staff were not told of. | |
| 1195 | 1225 | if let Err(error) = billing.watch_spend().await { | |
| 1196 | 1226 | worker::console_error!("watching g1t's own spend failed: {error}"); | |
| 1197 | 1227 | } | |
| 1198 | − | if let Ok(identity) = env.service("IDENTITY") | |
| 1199 | − | && let Err(error) = billing.warn_limits(&identity).await { | |
| 1200 | − | worker::console_error!("warning owners failed: {error}"); | |
| 1201 | − | } | |
| 1228 | + | let keeper = keeper::Keeper::from_env(&env); | |
| 1229 | + | let daily = event.cron() == keeper::DAILY; | |
| 1230 | + | if !daily { | |
| 1231 | + | tick(&billing, &env, &keeper).await; | |
| 1232 | + | } | |
| 1202 | 1233 | // Once a day: Stripe's endpoint kept listening to billing's events and | |
| 1203 | 1234 | // enabled, and saved cards and plans not read in a while read again. | |
| 1204 | − | if event.cron() == keeper::DAILY { | |
| 1235 | + | if daily { | |
| 1205 | 1236 | match billing.keep_endpoint("billing").await { | |
| 1206 | 1237 | Ok(done) => worker::console_log!("stripe endpoint: {done}"), | |
| 1207 | 1238 | Err(error) => worker::console_error!("keeping Stripe's endpoint failed: {error}"), | |
| ⋯ | |||
| 1213 | 1244 | } | |
| 1214 | 1245 | // Once a day, and at once if the costs were never checked: check every | |
| 1215 | 1246 | // cost against what Cloudflare billed. | |
| 1216 | − | if (event.cron() == keeper::DAILY || billing.never_checked().await.unwrap_or(false)) | |
| 1247 | + | if (daily || billing.never_checked().await.unwrap_or(false)) | |
| 1217 | 1248 | && let Err(error) = billing.reconcile(&keeper).await { | |
| 1218 | 1249 | worker::console_error!("checking costs against Cloudflare failed: {error}"); | |
| 1219 | 1250 | } | |
| 1220 | 1251 | // Once a day: credit from g1t past its expiry stops counting | |
| 1221 | 1252 | // (grants.rs), before the day is reconciled. | |
| 1222 | − | if event.cron() == keeper::DAILY { | |
| 1253 | + | if daily { | |
| 1223 | 1254 | match billing.expire_credits().await { | |
| 1224 | 1255 | Ok(closed) => worker::console_log!("credits expired: {closed}"), | |
| 1225 | 1256 | Err(error) => worker::console_error!("expiring credits failed: {error}"), | |
| ⋯ | |||
| 1228 | 1259 | // Once a day: what Cloudflare charged, reconciled against what g1t | |
| 1229 | 1260 | // counted and charged; prices whose day has come; margin alerts | |
| 1230 | 1261 | // (margin.rs). After the keeper, so its proposals are in. | |
| 1231 | − | if event.cron() == keeper::DAILY { | |
| 1262 | + | if daily { | |
| 1232 | 1263 | match billing.costs_daily(&env, &keeper).await { | |
| 1233 | 1264 | Ok(run) => worker::console_log!("costs: {} lines, {} days, {} proposals, {} alerts", run.lines, run.days, run.proposals, run.alerts), | |
| 1234 | 1265 | Err(error) => worker::console_error!("reconciling costs failed: {error}"), | |
| ⋯ | |||
| 1237 | 1268 | // Once a day: what each workspace's private repositories hold, its git | |
| 1238 | 1269 | // operations, Deployments plans from before the g1t plan set to end, | |
| 1239 | 1270 | // and old reservations cleared. | |
| 1240 | − | if event.cron() == keeper::DAILY { | |
| 1271 | + | if daily { | |
| 1241 | 1272 | if let Err(error) = billing.measure_packages().await { | |
| 1242 | 1273 | worker::console_error!("measuring package storage failed: {error}"); | |
| 1243 | 1274 | } | |
| 37 | 37 | //! Usage counts at what it cost g1t or what it is charged, whichever is | |
| 38 | 38 | //! more. Test-mode payments are not money, so they do not raise trust. | |
| 39 | 39 | ||
| 40 | + | use std::collections::BTreeSet; | |
| 41 | + | ||
| 40 | 42 | use futures_util::future::{try_join, try_join5, try_join_all}; | |
| 41 | 43 | use g1t_contracts::billing::{ | |
| 42 | 44 | BillingAccount, CheckLimitArgs, Limit, LimitArgs, LimitState, NotePendingArgs, PlanKind, SetBudgetArgs, SetSpendLimitArgs, Trust, | |
| ⋯ | |||
| 714 | 716 | Ok(true) | |
| 715 | 717 | } | |
| 716 | 718 | ||
| 719 | + | /// The workspaces with usage on the ledger this month: autopay's and | |
| 720 | + | /// the limit warnings' candidates, read once a tick for both. Served by | |
| 721 | + | /// `ledger_usage_by_time` (migration 0050), not a scan of the ledger. | |
| 722 | + | pub(crate) async fn month_users(&self) -> Result<BTreeSet<String>> { | |
| 723 | + | #[derive(Deserialize)] | |
| 724 | + | struct User { | |
| 725 | + | workspace: String, | |
| 726 | + | } | |
| 727 | + | let month_start = format!("{}-01", &rfc3339(now_ms())[..7]); | |
| 728 | + | Ok(self | |
| 729 | + | .db | |
| 730 | + | .prepare("SELECT DISTINCT workspace FROM ledger WHERE kind = 'usage' AND created_at >= ?") | |
| 731 | + | .bind(&[month_start.into()])? | |
| 732 | + | .all() | |
| 733 | + | .await? | |
| 734 | + | .results::<User>()? | |
| 735 | + | .into_iter() | |
| 736 | + | .map(|user| user.workspace) | |
| 737 | + | .collect()) | |
| 738 | + | } | |
| 739 | + | ||
| 717 | 740 | /// Charges the saved card of each workspace nearing its limit, for what | |
| 718 | 741 | /// it owes, so that a workspace that pays never has its work stopped. | |
| 719 | 742 | /// Only with live payments: test-mode payments are not money and lower | |
| 720 | 743 | /// nothing. Not for a workspace's own spend limit, which means stop, nor | |
| 721 | 744 | /// for enterprises, which are invoiced. A charge at the limit always | |
| 722 | 745 | /// goes through, whatever the minimum charge. | |
| 723 | − | pub(crate) async fn autopay(&self) -> Result<()> { | |
| 746 | + | /// | |
| 747 | + | /// `users` are the workspaces with usage this month ([`Self::month_users`]). | |
| 748 | + | pub(crate) async fn autopay(&self, users: &BTreeSet<String>) -> Result<()> { | |
| 724 | 749 | if !self.stripe.as_ref().is_some_and(crate::stripe::Stripe::live) { | |
| 725 | 750 | return Ok(()); | |
| 726 | 751 | } | |
| ⋯ | |||
| 728 | 753 | struct Candidate { | |
| 729 | 754 | workspace: String, | |
| 730 | 755 | } | |
| 731 | − | let month_start = format!("{}-01", &rfc3339(now_ms())[..7]); | |
| 732 | 756 | // With a card, and not already declined: a declined card waits for | |
| 733 | 757 | // the owners, rather than being tried again every few minutes. | |
| 734 | 758 | let candidates = self | |
| 735 | 759 | .db | |
| 736 | 760 | .prepare( | |
| 737 | − | "SELECT DISTINCT ledger.workspace AS workspace | |
| 738 | − | FROM ledger JOIN accounts ON accounts.workspace = ledger.workspace | |
| 739 | − | LEFT JOIN limits ON limits.workspace = ledger.workspace | |
| 740 | − | WHERE ledger.kind = 'usage' AND ledger.created_at >= ? AND accounts.customer_id IS NOT NULL | |
| 741 | − | AND limits.autopay_failed_at IS NULL", | |
| 761 | + | "SELECT accounts.workspace AS workspace | |
| 762 | + | FROM accounts LEFT JOIN limits ON limits.workspace = accounts.workspace | |
| 763 | + | WHERE accounts.customer_id IS NOT NULL AND limits.autopay_failed_at IS NULL", | |
| 742 | 764 | ) | |
| 743 | − | .bind(&[month_start.as_str().into()])? | |
| 744 | 765 | .all() | |
| 745 | 766 | .await? | |
| 746 | − | .results::<Candidate>()?; | |
| 767 | + | .results::<Candidate>()? | |
| 768 | + | .into_iter() | |
| 769 | + | .filter(|candidate| users.contains(&candidate.workspace)); | |
| 747 | 770 | for candidate in candidates { | |
| 748 | 771 | let limit = self.limit_of(&candidate.workspace).await?; | |
| 749 | 772 | // Near g1t's ceiling on what is unpaid; the spend limit is the | |
| ⋯ | |||
| 850 | 873 | /// plan's included usage, its spend limit and g1t's ceiling, once each a | |
| 851 | 874 | /// month; when its card was declined; and when a spend spike paused it. | |
| 852 | 875 | /// The same alerts show in the app (`entitlements`). | |
| 853 | − | pub(crate) async fn warn_limits(&self, identity: &worker::Fetcher) -> Result<()> { | |
| 876 | + | /// | |
| 877 | + | /// `users` are the workspaces with usage this month ([`Self::month_users`]). | |
| 878 | + | pub(crate) async fn warn_limits(&self, identity: &worker::Fetcher, users: &BTreeSet<String>) -> Result<()> { | |
| 854 | 879 | if self.stripe.is_none() { | |
| 855 | 880 | return Ok(()); | |
| 856 | 881 | } | |
| ⋯ | |||
| 860 | 885 | struct Candidate { | |
| 861 | 886 | workspace: String, | |
| 862 | 887 | } | |
| 863 | − | let candidates = self | |
| 864 | − | .db | |
| 865 | − | .prepare( | |
| 866 | − | "SELECT DISTINCT workspace FROM ledger WHERE kind = 'usage' AND created_at >= ?1 | |
| 867 | − | UNION SELECT workspace FROM limits WHERE autopay_failed_at IS NOT NULL", | |
| 868 | − | ) | |
| 869 | − | .bind(&[format!("{month}-01").into()])? | |
| 870 | − | .all() | |
| 871 | − | .await? | |
| 872 | − | .results::<Candidate>()?; | |
| 888 | + | let mut candidates = users.clone(); | |
| 889 | + | candidates.extend( | |
| 890 | + | self.db | |
| 891 | + | .prepare("SELECT workspace FROM limits WHERE autopay_failed_at IS NOT NULL") | |
| 892 | + | .all() | |
| 893 | + | .await? | |
| 894 | + | .results::<Candidate>()? | |
| 895 | + | .into_iter() | |
| 896 | + | .map(|candidate| candidate.workspace), | |
| 897 | + | ); | |
| 873 | 898 | #[derive(Deserialize)] | |
| 874 | 899 | struct Told { | |
| 875 | 900 | autopay_failed_at: Option<String>, | |
| ⋯ | |||
| 880 | 905 | meter: String, | |
| 881 | 906 | level: Option<i64>, | |
| 882 | 907 | } | |
| 883 | − | for Candidate { workspace } in candidates { | |
| 908 | + | for workspace in candidates { | |
| 884 | 909 | let told = self | |
| 885 | 910 | .db | |
| 886 | 911 | .prepare("SELECT autopay_failed_at, declined_told_at FROM limits WHERE workspace = ?") | |
| 79 | 79 | import { Cloudflare, type BuiltWorker, type Manifest } from "./cloudflare"; | |
| 80 | 80 | import { CustomHostnames } from "./custom-hostnames"; | |
| 81 | 81 | import { Domains, NOT_ENABLED_NOTICE, toDomain } from "./domains"; | |
| 82 | − | import { monthCost } from "./metering"; | |
| 82 | + | import { monthCost, movesMeter } from "./metering"; | |
| 83 | 83 | import { moveTargets, ownerOf, rebuildOutcome, type DroppedBuild, type MoveTarget } from "./moves"; | |
| 84 | 84 | import { commitMissing, MAX_IDENTICAL_FAILURES, missingCommitMessage, retryDecision, type PastBuild } from "./retries"; | |
| 85 | 85 | import { appHost, appUrl, label, uniqueLabel } from "./names"; | |
| ⋯ | |||
| 2039 | 2039 | * previews, the apps of workspaces whose plan ended, and scripts no app | |
| 2040 | 2040 | * holds come down; and a month that is over is charged past its | |
| 2041 | 2041 | * allowance. | |
| 2042 | + | * | |
| 2043 | + | * Each step stands alone: one that fails is logged and the rest still | |
| 2044 | + | * run. Taking idle previews down and charging months, which only read | |
| 2045 | + | * g1t's own tables, go before the steps that wait on Cloudflare's API, | |
| 2046 | + | * so a slow or failing API never holds them up. | |
| 2042 | 2047 | */ | |
| 2043 | 2048 | async sweep(): Promise<void> { | |
| 2044 | − | const cutoff = new Date(Date.now() - BUILD_TIMEOUT_MS).toISOString(); | |
| 2045 | − | const stuck = await this.db | |
| 2046 | − | .prepare("SELECT id FROM deployments WHERE status IN ('queued', 'building') AND created_at < ?") | |
| 2047 | − | .bind(cutoff) | |
| 2048 | − | .all<{ id: string }>(); | |
| 2049 | − | for (const { id } of stuck.results) await this.finishFailed(id, "The build did not finish in 45 minutes.", null, null); | |
| 2049 | + | const step = (name: string, work: () => Promise<unknown>) => work().catch((error) => console.error(`sweep: ${name} failed`, error)); | |
| 2050 | + | await step("failing stuck builds", async () => { | |
| 2051 | + | const cutoff = new Date(Date.now() - BUILD_TIMEOUT_MS).toISOString(); | |
| 2052 | + | const stuck = await this.db | |
| 2053 | + | .prepare("SELECT id FROM deployments WHERE status IN ('queued', 'building') AND created_at < ?") | |
| 2054 | + | .bind(cutoff) | |
| 2055 | + | .all<{ id: string }>(); | |
| 2056 | + | for (const { id } of stuck.results) await this.finishFailed(id, "The build did not finish in 45 minutes.", null, null); | |
| 2057 | + | }); | |
| 2058 | + | await step("taking down idle previews", () => this.takeDownIdle()); | |
| 2059 | + | await step("charging months", () => this.chargeMonths()); | |
| 2050 | 2060 | ||
| 2051 | 2061 | let apps = (await this.db.prepare("SELECT * FROM apps").all<AppRow>()).results; | |
| 2052 | 2062 | // Each app is its project's workspace's, as deployments has it now: an | |
| ⋯ | |||
| 2066 | 2076 | const billing = billingClient(this.env.BILLING); | |
| 2067 | 2077 | await this.holdToLimits(live, settings).catch((error) => console.error("could not apply limits", error)); | |
| 2068 | 2078 | for (const workspace of workspaces) { | |
| 2069 | − | const plan = await billing.hasFeature(workspace, "deployments"); | |
| 2070 | − | if (!plan.ok && plan.error.code === "payment_required") { | |
| 2071 | − | for (const app of live.filter((a) => ownerOf(a, owners) === workspace)) await this.removeApp(app.script); | |
| 2072 | − | // Custom domains cost g1t by the month: they go with the plan. | |
| 2073 | − | await this.domains.removeWhere("workspace", workspace).catch((error) => console.error("could not remove domains", error)); | |
| 2074 | − | } | |
| 2079 | + | await step(`ending ${workspace}'s apps`, async () => { | |
| 2080 | + | const plan = await billing.hasFeature(workspace, "deployments"); | |
| 2081 | + | if (!plan.ok && plan.error.code === "payment_required") { | |
| 2082 | + | for (const app of live.filter((a) => ownerOf(a, owners) === workspace)) await this.removeApp(app.script); | |
| 2083 | + | // Custom domains cost g1t by the month: they go with the plan. | |
| 2084 | + | await this.domains.removeWhere("workspace", workspace).catch((error) => console.error("could not remove domains", error)); | |
| 2085 | + | } | |
| 2086 | + | }); | |
| 2075 | 2087 | } | |
| 2076 | 2088 | ||
| 2077 | 2089 | // Apps whose project moved and are not up under the new name yet: a | |
| ⋯ | |||
| 2091 | 2103 | ||
| 2092 | 2104 | await this.removeOrphans(apps).catch((error) => console.error("could not remove orphans", error)); | |
| 2093 | 2105 | await this.count(apps).catch((error) => console.error("could not count usage", error)); | |
| 2094 | − | await this.takeDownIdle(); | |
| 2095 | − | await this.chargeMonths(); | |
| 2096 | 2106 | } | |
| 2097 | 2107 | ||
| 2098 | 2108 | /** | |
| ⋯ | |||
| 2187 | 2197 | perWorkspace.set(app.workspace, sum); | |
| 2188 | 2198 | } | |
| 2189 | 2199 | const at = now(); | |
| 2200 | + | // Written only when the count moves the meter: most sweeps it does not. | |
| 2201 | + | const stored = new Map( | |
| 2202 | + | ( | |
| 2203 | + | await this.db | |
| 2204 | + | .prepare("SELECT namespace, requests, cpu_ms FROM meters WHERE month = ?") | |
| 2205 | + | .bind(month()) | |
| 2206 | + | .all<{ namespace: string; requests: number; cpu_ms: number }>() | |
| 2207 | + | ).results.map((row) => [row.namespace, row]), | |
| 2208 | + | ); | |
| 2190 | 2209 | for (const [workspace, used] of perWorkspace) { | |
| 2210 | + | if (!movesMeter(stored.get(workspace), used)) continue; | |
| 2191 | 2211 | await this.db | |
| 2192 | 2212 | .prepare( | |
| 2193 | 2213 | `INSERT INTO meters (namespace, month, requests, cpu_ms, counted_at) VALUES (?1, ?2, ?3, ?4, ?5) | |
| ⋯ | |||
| 2203 | 2223 | new Date(Date.now() - 24 * 60 * 60 * 1000).toISOString(), | |
| 2204 | 2224 | at, | |
| 2205 | 2225 | ); | |
| 2226 | + | // At an hour's resolution: previews idle for days are what it finds. | |
| 2227 | + | const hourAgo = new Date(Date.now() - 60 * 60 * 1000).toISOString(); | |
| 2228 | + | const lastRequest = new Map(apps.map((app) => [app.script, app.last_request_at])); | |
| 2206 | 2229 | for (const [script, used] of recent) { | |
| 2207 | − | if (used.requests > 0) { | |
| 2230 | + | const last = lastRequest.get(script); | |
| 2231 | + | if (used.requests > 0 && (!last || last < hourAgo)) { | |
| 2208 | 2232 | await this.db.prepare("UPDATE apps SET last_request_at = ? WHERE script = ?").bind(at, script).run(); | |
| 2209 | 2233 | } | |
| 2210 | 2234 | } | |
| 1 | 1 | import assert from "node:assert/strict"; | |
| 2 | 2 | import { test } from "node:test"; | |
| 3 | 3 | ||
| 4 | − | import { amount, monthCost } from "./metering.ts"; | |
| 4 | + | import { amount, monthCost, movesMeter } from "./metering.ts"; | |
| 5 | 5 | ||
| 6 | 6 | /** Cloudflare's prices, at cost: $0.30 per million requests, $0.02 per million CPU ms, $0.10 a custom domain a month. */ | |
| 7 | 7 | const COSTS = { millionRequests: 300_000, millionCpuMs: 20_000, domainMonth: 100_000 }; | |
| ⋯ | |||
| 44 | 44 | assert.equal(amount(1_240_000), "1.2 million"); | |
| 45 | 45 | assert.equal(amount(0), "0"); | |
| 46 | 46 | }); | |
| 47 | + | ||
| 48 | + | test("a count that does not move the meter is not written", () => { | |
| 49 | + | assert.ok(movesMeter(undefined, { requests: 0, cpuMs: 0 })); | |
| 50 | + | assert.ok(movesMeter({ requests: 10, cpu_ms: 5 }, { requests: 11, cpuMs: 5 })); | |
| 51 | + | assert.ok(movesMeter({ requests: 10, cpu_ms: 5 }, { requests: 10, cpuMs: 6 })); | |
| 52 | + | assert.ok(!movesMeter({ requests: 10, cpu_ms: 5 }, { requests: 10, cpuMs: 5 })); | |
| 53 | + | // Analytics forgot an app that came down: the meter keeps what it had. | |
| 54 | + | assert.ok(!movesMeter({ requests: 10, cpu_ms: 5 }, { requests: 4, cpuMs: 2 })); | |
| 55 | + | }); | |
| 62 | 62 | description: parts.join(", "), | |
| 63 | 63 | }; | |
| 64 | 64 | } | |
| 65 | + | ||
| 66 | + | /** | |
| 67 | + | * Whether a count from analytics moves a workspace's meter for the month. | |
| 68 | + | * The meter keeps the most it has seen (analytics forgets apps that came | |
| 69 | + | * down), so a count at or below it changes nothing and is not written. | |
| 70 | + | */ | |
| 71 | + | export function movesMeter(stored: { requests: number; cpu_ms: number } | undefined, used: { requests: number; cpuMs: number }): boolean { | |
| 72 | + | return !stored || used.requests > stored.requests || used.cpuMs > stored.cpu_ms; | |
| 73 | + | } |
| 1 | + | -- A repository's events of some types, newest first: the home page's pushes, | |
| 2 | + | -- the overview's feed and the Activity tab. By repository, type and id | |
| 3 | + | -- those are ranges, rather than a walk back through every event the | |
| 4 | + | -- repository ever had. | |
| 5 | + | CREATE INDEX IF NOT EXISTS events_repo_type ON events (repo_id, type, id); |
| 1 | 1 | import { OG_RENDER_VERSION } from "@g1t/contracts/og"; | |
| 2 | 2 | ||
| 3 | + | import { cardPath, clean } from "./resolve.ts"; | |
| 4 | + | ||
| 3 | 5 | /** | |
| 4 | 6 | * The version of the cards' design, part of every cache key. Bump it (in | |
| 5 | 7 | * `packages/contracts/src/og.ts`, which the site and the docs read too) | |
| ⋯ | |||
| 9 | 11 | export const RENDER_VERSION = OG_RENDER_VERSION; | |
| 10 | 12 | ||
| 11 | 13 | /** | |
| 14 | + | * A `v` as the site and the docs write it (`ogVersion` in | |
| 15 | + | * packages/contracts): the render version, then, for the site, a | |
| 16 | + | * fingerprint of what the card shows, at most seven base-36 characters. | |
| 17 | + | */ | |
| 18 | + | const VERSION = new RegExp(`^${RENDER_VERSION.replace(/\./g, "\\.")}(\\.[0-9a-z]{1,7})?$`); | |
| 19 | + | ||
| 20 | + | /** | |
| 12 | 21 | * The address a card is cached under: the render version, then only the | |
| 13 | − | * parameters that change the card, in order, including the page's own | |
| 14 | − | * content version `v`, so stray ones (trackers, cache busters) cannot fill | |
| 15 | − | * the cache. | |
| 22 | + | * parameters that change the card, in order, so stray ones (trackers, | |
| 23 | + | * cache busters) cannot fill the cache. A page's address is the page whose | |
| 24 | + | * card it shows (`cardPath`), the docs' text as the card draws it, and the | |
| 25 | + | * page's own content version `v` only in the form the site writes it. | |
| 16 | 26 | */ | |
| 17 | 27 | export function cacheKey(url: URL): string { | |
| 18 | − | const wanted = | |
| 19 | − | url.pathname === "/docs" ? ["title", "section", "description", "v"] : url.pathname === "/image" ? ["path", "v"] : []; | |
| 20 | 28 | const key = new URL(url.pathname, url.origin); | |
| 21 | 29 | key.searchParams.set("render", RENDER_VERSION); | |
| 22 | − | for (const name of wanted) { | |
| 23 | − | const value = url.searchParams.get(name); | |
| 30 | + | const set = (name: string, value: string | null) => { | |
| 24 | 31 | if (value !== null) key.searchParams.set(name, value); | |
| 32 | + | }; | |
| 33 | + | if (url.pathname === "/docs") { | |
| 34 | + | // As docsCard reads them. | |
| 35 | + | set("title", clean(url.searchParams.get("title"), 160)); | |
| 36 | + | set("section", clean(url.searchParams.get("section"), 60)); | |
| 37 | + | set("description", clean(url.searchParams.get("description"), 300)); | |
| 38 | + | } else if (url.pathname === "/image") { | |
| 39 | + | set("path", cardPath(url.searchParams.get("path") ?? "/")); | |
| 40 | + | } else { | |
| 41 | + | return key.toString(); | |
| 25 | 42 | } | |
| 43 | + | const version = url.searchParams.get("v"); | |
| 44 | + | if (version !== null && VERSION.test(version)) set("v", version); | |
| 45 | + | return key.toString(); | |
| 46 | + | } | |
| 47 | + | ||
| 48 | + | /** | |
| 49 | + | * The address a drawn card is also kept under, by what it shows: two | |
| 50 | + | * addresses for the same card (an old `v`, a made-up one) share one | |
| 51 | + | * drawing rather than each being drawn. | |
| 52 | + | */ | |
| 53 | + | export async function drawnKey(origin: string, card: unknown): Promise<string> { | |
| 54 | + | const digest = await crypto.subtle.digest("SHA-256", new TextEncoder().encode(JSON.stringify(card))); | |
| 55 | + | const hex = [...new Uint8Array(digest)].map((byte) => byte.toString(16).padStart(2, "0")).join(""); | |
| 56 | + | const key = new URL("/drawn", origin); | |
| 57 | + | key.searchParams.set("render", RENDER_VERSION); | |
| 58 | + | key.searchParams.set("card", hex); | |
| 26 | 59 | return key.toString(); | |
| 27 | 60 | } | |
| 32 | 32 | import mono400 from "./fonts/ibm-plex-mono-400.ttf"; | |
| 33 | 33 | import mono500 from "./fonts/ibm-plex-mono-500.ttf"; | |
| 34 | 34 | import { cardPng } from "./render.ts"; | |
| 35 | − | import { cacheKey } from "./cache.ts"; | |
| 35 | + | import { cacheKey, drawnKey } from "./cache.ts"; | |
| 36 | 36 | import { type Shot, screenshotOf, sweep, take } from "./capture.ts"; | |
| 37 | 37 | import { parseShot } from "./screenshot.ts"; | |
| 38 | − | import { BRAND, type Card, docsCard, resolve } from "./resolve.ts"; | |
| 38 | + | import { BRAND, type Card, cardPath, docsCard, resolve } from "./resolve.ts"; | |
| 39 | 39 | ||
| 40 | 40 | interface Env { | |
| 41 | 41 | /** Uploaded avatars by hash, with `{ contentType }`; written by identity. */ | |
| ⋯ | |||
| 138 | 138 | return drawn.ok ? withCacheControl(drawn, true) : drawn; | |
| 139 | 139 | } | |
| 140 | 140 | ||
| 141 | − | const card = await cardFor(url, env); | |
| 142 | − | const failed = card.kind === "brand" && card.failed === true; | |
| 143 | − | if (card.kind === "brand" && key !== brandKey) { | |
| 141 | + | const found = await lookUpCard(url, env); | |
| 142 | + | const failed = found.kind === "brand" && found.failed === true; | |
| 143 | + | if (found.kind === "brand" && key !== brandKey) { | |
| 144 | 144 | const brand = await cache.match(brandKey); | |
| 145 | 145 | if (brand) return withCacheControl(brand, failed); | |
| 146 | 146 | } | |
| 147 | + | // The same card under another address (another `v`): drawn once. | |
| 148 | + | const drawn = found.kind === "brand" ? null : await drawnKey(url.origin, found); | |
| 149 | + | const same = drawn ? await cache.match(drawn) : undefined; | |
| 150 | + | if (same) { | |
| 151 | + | ctx.waitUntil(cache.put(key, same.clone())); | |
| 152 | + | return same; | |
| 153 | + | } | |
| 154 | + | const card = await withIcon(found, env); | |
| 147 | 155 | ||
| 148 | 156 | let png: Uint8Array; | |
| 149 | 157 | let fellBack = false; | |
| ⋯ | |||
| 180 | 188 | ctx.waitUntil(cache.put(brandKey, response.clone())); | |
| 181 | 189 | } else { | |
| 182 | 190 | ctx.waitUntil(cache.put(key, response.clone())); | |
| 191 | + | if (drawn) ctx.waitUntil(cache.put(drawn, response.clone())); | |
| 183 | 192 | } | |
| 184 | 193 | return withCacheControl(response, failed); | |
| 185 | 194 | }, | |
| ⋯ | |||
| 219 | 228 | async function iconFor(env: Env, avatar: string | null | undefined): Promise<string | undefined> { | |
| 220 | 229 | if (!avatar || !/^[0-9a-f]{64}$/.test(avatar)) return undefined; | |
| 221 | 230 | try { | |
| 231 | + | // An avatar's key is its hash: what is under it never changes. | |
| 222 | 232 | const { value, metadata } = await env.AVATARS.getWithMetadata<{ contentType?: string }>(avatar, { | |
| 223 | 233 | type: "arrayBuffer", | |
| 234 | + | cacheTtl: 86_400, | |
| 224 | 235 | }); | |
| 225 | 236 | const type = metadata?.contentType; | |
| 226 | 237 | if (!value || !type || !DRAWABLE.has(type)) return undefined; | |
| ⋯ | |||
| 233 | 244 | } | |
| 234 | 245 | } | |
| 235 | 246 | ||
| 236 | − | async function cardFor(url: URL, env: Env): Promise<Card> { | |
| 237 | − | const card = await lookUpCard(url, env); | |
| 247 | + | async function withIcon(card: Card, env: Env): Promise<Card> { | |
| 238 | 248 | if (card.kind === "workspace" || card.kind === "person") return { ...card, icon: await iconFor(env, card.avatar) }; | |
| 239 | 249 | return card; | |
| 240 | 250 | } | |
| ⋯ | |||
| 242 | 252 | async function lookUpCard(url: URL, env: Env): Promise<Card> { | |
| 243 | 253 | if (url.pathname === "/docs") return docsCard(url.searchParams); | |
| 244 | 254 | if (url.pathname === "/") return BRAND; | |
| 245 | − | return resolve(url.searchParams.get("path") ?? "/", { | |
| 255 | + | // The page whose card it is, as the cache key has it (cache.ts). | |
| 256 | + | return resolve(cardPath(url.searchParams.get("path") ?? "/"), { | |
| 246 | 257 | identity: identityClient(env.IDENTITY), | |
| 247 | 258 | repos: reposClient(env.REPOS), | |
| 248 | 259 | work: workClient(env.WORK), | |
| 4 | 4 | import type { Issue, Project, Pull, Repo, Result, Workspace } from "@g1t/contracts"; | |
| 5 | 5 | ||
| 6 | 6 | import { RENDER_VERSION, cacheKey } from "./cache.ts"; | |
| 7 | − | import { type Sources, clean, docsCard, resolve, segments } from "./resolve.ts"; | |
| 7 | + | import { type Sources, cardPath, clean, docsCard, resolve, segments } from "./resolve.ts"; | |
| 8 | 8 | ||
| 9 | 9 | const ok = <T>(value: T): Result<T> => ({ ok: true, value }); | |
| 10 | 10 | const notFound: Result<never> = { ok: false, error: { code: "not_found", message: "Not found" } }; | |
| ⋯ | |||
| 227 | 227 | }); | |
| 228 | 228 | ||
| 229 | 229 | test("the cache key keeps only what changes the card", () => { | |
| 230 | − | const a = cacheKey(new URL("https://og.g1t.sh/image?utm=1&path=/acme/web&v=3")); | |
| 231 | − | const b = cacheKey(new URL("https://og.g1t.sh/image?path=/acme/web&v=3&x=2")); | |
| 230 | + | const v = `${RENDER_VERSION}.1z141z3`; | |
| 231 | + | const a = cacheKey(new URL(`https://og.g1t.sh/image?utm=1&path=/acme/web&v=${v}`)); | |
| 232 | + | const b = cacheKey(new URL(`https://og.g1t.sh/image?path=/acme/web&v=${v}&x=2`)); | |
| 232 | 233 | assert.equal(a, b); | |
| 233 | − | assert.notEqual(a, cacheKey(new URL("https://og.g1t.sh/image?path=/acme/web&v=4"))); | |
| 234 | + | assert.notEqual(a, cacheKey(new URL(`https://og.g1t.sh/image?path=/acme/web&v=${RENDER_VERSION}.abc`))); | |
| 234 | 235 | assert.match(a, new RegExp(`render=${RENDER_VERSION}`)); | |
| 235 | 236 | // The keys cards were kept under before the render version was. | |
| 236 | 237 | assert.notEqual(a, "https://og.g1t.sh/image?design=1&path=%2Facme%2Fweb&v=3"); | |
| 237 | − | const docs = cacheKey(new URL("https://og.g1t.sh/docs?title=Quickstart&v=2&utm=x")); | |
| 238 | − | assert.equal(docs, `https://og.g1t.sh/docs?render=${RENDER_VERSION}&title=Quickstart&v=2`); | |
| 238 | + | const docs = cacheKey(new URL(`https://og.g1t.sh/docs?title=Quickstart&v=${RENDER_VERSION}&utm=x`)); | |
| 239 | + | assert.equal(docs, `https://og.g1t.sh/docs?render=${RENDER_VERSION}&title=Quickstart&v=${RENDER_VERSION}`); | |
| 240 | + | // The docs' text as the card draws it. | |
| 241 | + | assert.equal(cacheKey(new URL("https://og.g1t.sh/docs?title=%20Quick%0Astart%20")), cacheKey(new URL("https://og.g1t.sh/docs?title=Quick%20start"))); | |
| 242 | + | }); | |
| 243 | + | ||
| 244 | + | test("a v the site never writes is left out of the key", () => { | |
| 245 | + | const plain = cacheKey(new URL("https://og.g1t.sh/image?path=/acme/web")); | |
| 246 | + | for (const v of ["random", "999", `${RENDER_VERSION}.TOOLONG12`, `${RENDER_VERSION}.ab-c`, `x${RENDER_VERSION}`]) { | |
| 247 | + | assert.equal(cacheKey(new URL(`https://og.g1t.sh/image?path=/acme/web&v=${encodeURIComponent(v)}`)), plain, v); | |
| 248 | + | } | |
| 249 | + | }); | |
| 250 | + | ||
| 251 | + | test("a path is keyed by the page whose card it shows", () => { | |
| 252 | + | const key = (path: string) => cacheKey(new URL(`https://og.g1t.sh/image?path=${encodeURIComponent(path)}`)); | |
| 253 | + | assert.equal(key("/acme/web/tree/main/src?x=1"), key("/acme/web")); | |
| 254 | + | assert.equal(key("/acme/web/pull/3/files"), key("/acme/web/pull/3")); | |
| 255 | + | assert.notEqual(key("/acme/web/pull/3"), key("/acme/web")); | |
| 256 | + | assert.equal(key("/Acme/-/settings"), key("/acme")); | |
| 257 | + | assert.equal(key("/pricing/"), key("/pricing")); | |
| 258 | + | assert.equal(key("/not a name/web"), key("/")); | |
| 259 | + | assert.equal(key("/login/extra/parts"), key("/")); | |
| 260 | + | assert.equal(cardPath("/u/Ada"), "/u/ada"); | |
| 261 | + | assert.equal(cardPath("/acme/web/issues/12"), "/acme/web/issues/12"); | |
| 262 | + | assert.equal(cardPath("/acme/web/issues/012"), "/acme/web"); | |
| 263 | + | assert.equal(cardPath("/acme/web/soon/no-such-thing"), "/acme/web"); | |
| 264 | + | assert.equal(cardPath("//evil"), "/"); | |
| 265 | + | }); | |
| 266 | + | ||
| 267 | + | test("a path resolves to the same card as the page it is keyed by", async () => { | |
| 268 | + | const { sources: s } = sources(); | |
| 269 | + | for (const path of ["/acme/web/tree/main", "/acme/web/pull/7/files", "/acme/-/settings", "/acme/web/issues/3"]) { | |
| 270 | + | assert.deepEqual(await resolve(path, s), await resolve(cardPath(path), s), path); | |
| 271 | + | } | |
| 239 | 272 | }); | |
| 234 | 234 | } | |
| 235 | 235 | } | |
| 236 | 236 | ||
| 237 | + | /** | |
| 238 | + | * The page whose card `path` shows, asking no one: the parts of an address | |
| 239 | + | * a card does not read are dropped, and an address no card is for becomes | |
| 240 | + | * `/`, the brand card. `/acme/web/tree/main/src` and `/acme/web?tab=1` are | |
| 241 | + | * both `/acme/web`. The cache key and the lookup both use it, so made-up | |
| 242 | + | * addresses cannot make the service draw the same card again and again. | |
| 243 | + | */ | |
| 244 | + | export function cardPath(path: string): string { | |
| 245 | + | const parts = segments(path); | |
| 246 | + | if (!parts || parts.length === 0) return "/"; | |
| 247 | + | const join = (...kept: string[]) => `/${kept.map(encodeURIComponent).join("/")}`; | |
| 248 | + | const [owner, repo, section, item] = parts; | |
| 249 | + | const lower = owner.toLowerCase(); | |
| 250 | + | if (lower === "u") return repo !== undefined && parts.length === 2 && NAME.test(repo) ? join("u", repo.toLowerCase()) : "/"; | |
| 251 | + | if (RESERVED.has(lower)) { | |
| 252 | + | if (parts.length === 1) return lower in PAGES ? join(lower) : "/"; | |
| 253 | + | const page = `${lower}/${parts[1].toLowerCase()}`; | |
| 254 | + | return parts.length === 2 && page in PAGES ? join(lower, parts[1].toLowerCase()) : "/"; | |
| 255 | + | } | |
| 256 | + | if (!NAME.test(owner)) return "/"; | |
| 257 | + | if (repo === undefined || repo === "-") return join(lower); | |
| 258 | + | if (!NAME.test(repo)) return "/"; | |
| 259 | + | if (section === "issues" && item !== undefined && isNumber(item) && parts.length === 4) return join(owner, repo, "issues", item); | |
| 260 | + | if (section === "pull" && item !== undefined && isNumber(item)) return join(owner, repo, "pull", item); | |
| 261 | + | if (section === "soon" && item !== undefined && parts.length === 4 && roadmapItem(item)) return join(owner, repo, "soon", item); | |
| 262 | + | return join(owner, repo); | |
| 263 | + | } | |
| 264 | + | ||
| 237 | 265 | /** The card for a page of g1t.sh, as an anonymous visitor would see it. */ | |
| 238 | 266 | export async function resolve(path: string, sources: Sources): Promise<Card> { | |
| 239 | 267 | const parts = segments(path); |
| 5 | 5 | * Platforms namespace (`web-git-fix-login-acme.g1t.page` is the fix-login | |
| 6 | 6 | * branch's preview of acme's web project), so the app is fetched by name | |
| 7 | 7 | * and runs only for as long as it answers. An app no one visits runs | |
| 8 | − | * nothing and costs nothing. The one lookup is for an address an app | |
| 8 | + | * nothing and costs nothing. Each request it serves is held to | |
| 9 | + | * `APP_LIMITS`, so one app cannot spend without bound. The one lookup is for an address an app | |
| 9 | 10 | * had before its project moved: `DOMAINS` holds a redirect under the old | |
| 10 | 11 | * hostname for as long as the old name is held, and it is followed before | |
| 11 | 12 | * anything the old script would answer (such as a paused notice). | |
| ⋯ | |||
| 18 | 19 | * people sign in to. | |
| 19 | 20 | */ | |
| 20 | 21 | ||
| 21 | − | import { DOMAIN, FALLBACK, parseEntry, redirectTo, route, type DomainEntry } from "./route.ts"; | |
| 22 | + | import { APP_LIMITS, DOMAIN, FALLBACK, overLimits, parseEntry, redirectTo, route, type DomainEntry } from "./route.ts"; | |
| 22 | 23 | ||
| 23 | 24 | type Env = { | |
| 24 | 25 | APPS: DispatchNamespace; | |
| ⋯ | |||
| 76 | 77 | async function dispatch(env: Env, request: Request, script: string, shown: string, missing: () => Response): Promise<Response> { | |
| 77 | 78 | let app: Fetcher; | |
| 78 | 79 | try { | |
| 79 | − | app = env.APPS.get(script); | |
| 80 | + | app = env.APPS.get(script, {}, { limits: APP_LIMITS }); | |
| 80 | 81 | } catch { | |
| 81 | 82 | return missing(); | |
| 82 | 83 | } | |
| ⋯ | |||
| 84 | 85 | return await app.fetch(request); | |
| 85 | 86 | } catch (error) { | |
| 86 | 87 | if (/worker not found|script not found|does not exist/i.test(String(error))) return missing(); | |
| 88 | + | if (overLimits(error)) { | |
| 89 | + | return page( | |
| 90 | + | 503, | |
| 91 | + | "The app went over its limits", | |
| 92 | + | `<p>The app at <code>${escape(shown)}</code> used more than ${APP_LIMITS.cpuMs} ms of CPU time, or made more than ${APP_LIMITS.subRequests} requests of its own, answering this request.</p>`, | |
| 93 | + | shown, | |
| 94 | + | ); | |
| 95 | + | } | |
| 87 | 96 | return page( | |
| 88 | 97 | 502, | |
| 89 | 98 | "The app failed", | |
| 1 | 1 | import assert from "node:assert/strict"; | |
| 2 | 2 | import { test } from "node:test"; | |
| 3 | 3 | ||
| 4 | − | import { parseEntry, redirectTo, route } from "./route.ts"; | |
| 4 | + | import { APP_LIMITS, overLimits, parseEntry, redirectTo, route } from "./route.ts"; | |
| 5 | 5 | import worker from "./index.ts"; | |
| 6 | 6 | ||
| 7 | 7 | test("hosts on g1t.page are the home page, the fallback origin, or an app", () => { | |
| ⋯ | |||
| 124 | 124 | }; | |
| 125 | 125 | assert.equal(await (await call("https://shop-acme.g1t.page/", e)).text(), "shop"); | |
| 126 | 126 | }); | |
| 127 | + | ||
| 128 | + | test("every app runs within the limits, and one that goes over them says so", async () => { | |
| 129 | + | const asked: unknown[] = []; | |
| 130 | + | const e = { | |
| 131 | + | APPS: { | |
| 132 | + | get(_name: string, _args: unknown, options: unknown) { | |
| 133 | + | asked.push(options); | |
| 134 | + | return { | |
| 135 | + | async fetch() { | |
| 136 | + | throw new Error("Worker exceeded CPU time limit."); | |
| 137 | + | }, | |
| 138 | + | }; | |
| 139 | + | }, | |
| 140 | + | }, | |
| 141 | + | } as any; | |
| 142 | + | const response = await call("https://web-acme.g1t.page/", e); | |
| 143 | + | assert.deepEqual(asked, [{ limits: APP_LIMITS }]); | |
| 144 | + | assert.equal(response.status, 503); | |
| 145 | + | assert.match(await response.text(), /went over its limits/); | |
| 146 | + | assert.ok(overLimits(new Error("Too many subrequests."))); | |
| 147 | + | assert.ok(!overLimits(new Error("TypeError: x is undefined"))); | |
| 148 | + | }); | |
| 47 | 47 | const url = new URL(requestUrl); | |
| 48 | 48 | return `https://${target}${url.pathname}${url.search}`; | |
| 49 | 49 | } | |
| 50 | + | ||
| 51 | + | /** | |
| 52 | + | * What one request to an app may use: CPU time, not counting time spent | |
| 53 | + | * waiting on the network, and requests it makes of its own. Static assets | |
| 54 | + | * are served without running the app, and use neither. No plan sets | |
| 55 | + | * others; the Deployments guide names these. | |
| 56 | + | */ | |
| 57 | + | export const APP_LIMITS = { cpuMs: 50, subRequests: 50 } as const; | |
| 58 | + | ||
| 59 | + | /** Whether an app's failure was going over `APP_LIMITS`, as the runtime words it. */ | |
| 60 | + | export function overLimits(error: unknown): boolean { | |
| 61 | + | return /exceeded (its )?cpu|cpu time limit|too many subrequests|subrequest limit/i.test(String(error)); | |
| 62 | + | } |
| 1 | + | -- Deliveries older than a fortnight are forgotten each hour, oldest first: | |
| 2 | + | -- by time, a range rather than a scan of every delivery. | |
| 3 | + | CREATE INDEX IF NOT EXISTS deliveries_by_time ON deliveries (created_at); |
| 40 | 40 | const KEPT_DAYS: u64 = 14; | |
| 41 | 41 | /// How many due retries one sweep makes. | |
| 42 | 42 | const SWEEP: u32 = 50; | |
| 43 | + | /// The hourly cron that forgets old deliveries (wrangler.jsonc); every | |
| 44 | + | /// other minute only retries. | |
| 45 | + | const PURGE_CRON: &str = "37 * * * *"; | |
| 46 | + | /// How many old deliveries one statement deletes, and how many statements | |
| 47 | + | /// one purge makes: a busy hour's worth, in bites D1 answers quickly. | |
| 48 | + | const PURGE_BATCH: u32 = 1_000; | |
| 49 | + | const PURGE_BATCHES: u32 = 50; | |
| 43 | 50 | ||
| 44 | 51 | #[derive(Deserialize)] | |
| 45 | 52 | struct HookRow { | |
| ⋯ | |||
| 661 | 668 | Ok(()) | |
| 662 | 669 | } | |
| 663 | 670 | ||
| 664 | − | /// Tries again what is due, and forgets what is old. | |
| 671 | + | /// Tries again what is due. | |
| 665 | 672 | async fn sweep(&self) -> Result<()> { | |
| 666 | 673 | let now = now_ms(); | |
| 667 | 674 | let due = self | |
| ⋯ | |||
| 690 | 697 | } | |
| 691 | 698 | } | |
| 692 | 699 | } | |
| 693 | − | self.db | |
| 694 | − | .prepare("DELETE FROM deliveries WHERE created_at < ?") | |
| 695 | − | .bind(&[rfc3339(now.saturating_sub(KEPT_DAYS * 24 * 60 * 60 * 1000)).into()])? | |
| 696 | − | .run() | |
| 697 | − | .await?; | |
| 698 | 700 | Ok(()) | |
| 699 | 701 | } | |
| 702 | + | ||
| 703 | + | /// Forgets deliveries older than [`KEPT_DAYS`], oldest first, a batch | |
| 704 | + | /// at a time (`deliveries_by_time`, migration 0002). What one purge | |
| 705 | + | /// leaves, the next hour's takes. Returns how many it deleted. | |
| 706 | + | async fn purge(&self) -> Result<u32> { | |
| 707 | + | let before = rfc3339(now_ms().saturating_sub(KEPT_DAYS * 24 * 60 * 60 * 1000)); | |
| 708 | + | let mut deleted = 0; | |
| 709 | + | for _ in 0..PURGE_BATCHES { | |
| 710 | + | let result = self | |
| 711 | + | .db | |
| 712 | + | .prepare(PURGE_SQL) | |
| 713 | + | .bind(&[before.as_str().into(), PURGE_BATCH.into()])? | |
| 714 | + | .run() | |
| 715 | + | .await?; | |
| 716 | + | let changed = result.meta()?.and_then(|meta| meta.changes).unwrap_or(0) as u32; | |
| 717 | + | deleted += changed; | |
| 718 | + | if changed < PURGE_BATCH { | |
| 719 | + | break; | |
| 720 | + | } | |
| 721 | + | } | |
| 722 | + | Ok(deleted) | |
| 723 | + | } | |
| 700 | 724 | } | |
| 701 | 725 | ||
| 702 | 726 | /// One HTTPS POST of a delivery's payload, signed, given ten seconds. | |
| ⋯ | |||
| 810 | 834 | Ok(()) | |
| 811 | 835 | } | |
| 812 | 836 | ||
| 813 | − | /// Every minute: retries that are due, and deliveries old enough to forget. | |
| 837 | + | /// One batch of old deliveries, by time. | |
| 838 | + | const PURGE_SQL: &str = "DELETE FROM deliveries WHERE rowid IN (SELECT rowid FROM deliveries WHERE created_at < ? ORDER BY created_at LIMIT ?)"; | |
| 839 | + | ||
| 840 | + | /// Every minute: retries that are due. Once an hour, at [`PURGE_CRON`]: | |
| 841 | + | /// deliveries old enough to forget. | |
| 814 | 842 | #[event(scheduled)] | |
| 815 | − | async fn scheduled(_event: ScheduledEvent, env: Env, _ctx: ScheduleContext) { | |
| 816 | − | match Webhooks::new(&env) { | |
| 817 | − | Ok(service) => { | |
| 818 | − | if let Err(error) = service.sweep().await { | |
| 819 | − | worker::console_error!("webhooks: the sweep failed: {error}"); | |
| 820 | − | } | |
| 843 | + | async fn scheduled(event: ScheduledEvent, env: Env, _ctx: ScheduleContext) { | |
| 844 | + | let service = match Webhooks::new(&env) { | |
| 845 | + | Ok(service) => service, | |
| 846 | + | Err(error) => { | |
| 847 | + | worker::console_error!("webhooks: could not start: {error}"); | |
| 848 | + | return; | |
| 849 | + | } | |
| 850 | + | }; | |
| 851 | + | if event.cron() == PURGE_CRON { | |
| 852 | + | match service.purge().await { | |
| 853 | + | Ok(0) => {} | |
| 854 | + | Ok(deleted) => worker::console_log!("webhooks: forgot {deleted} old deliveries"), | |
| 855 | + | Err(error) => worker::console_error!("webhooks: the purge failed: {error}"), | |
| 821 | 856 | } | |
| 822 | − | Err(error) => worker::console_error!("webhooks: could not start: {error}"), | |
| 857 | + | return; | |
| 858 | + | } | |
| 859 | + | if let Err(error) = service.sweep().await { | |
| 860 | + | worker::console_error!("webhooks: the sweep failed: {error}"); | |
| 861 | + | } | |
| 862 | + | } | |
| 863 | + | ||
| 864 | + | #[cfg(test)] | |
| 865 | + | mod tests { | |
| 866 | + | use super::*; | |
| 867 | + | ||
| 868 | + | #[test] | |
| 869 | + | fn the_purge_has_a_cron_of_its_own() { | |
| 870 | + | let wrangler = include_str!("../wrangler.jsonc"); | |
| 871 | + | assert!(wrangler.contains(&format!("\"{PURGE_CRON}\"")), "{PURGE_CRON} is not among the crons"); | |
| 872 | + | assert_ne!(PURGE_CRON, "* * * * *"); | |
| 873 | + | } | |
| 874 | + | ||
| 875 | + | #[test] | |
| 876 | + | fn the_purge_deletes_in_batches_by_time() { | |
| 877 | + | assert!(PURGE_SQL.contains("WHERE created_at < ? ORDER BY created_at LIMIT ?")); | |
| 878 | + | assert!(PURGE_BATCH * PURGE_BATCHES >= 10_000); | |
| 823 | 879 | } | |
| 824 | 880 | } | |
| 825 | 881 | ||
| 25 | 25 | "queues": { | |
| 26 | 26 | "consumers": [{ "queue": "g1t-events-webhooks", "max_batch_size": 20, "max_batch_timeout": 1, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" }] | |
| 27 | 27 | }, | |
| 28 | − | // Retries that are due, and deliveries old enough to forget. | |
| 29 | − | "triggers": { "crons": ["* * * * *"] }, | |
| 28 | + | // Every minute: retries. At 37 past each hour: deliveries older than a | |
| 29 | + | // fortnight forgotten (PURGE_CRON in src/lib.rs). | |
| 30 | + | "triggers": { "crons": ["* * * * *", "37 * * * *"] }, | |
| 30 | 31 | // Secret: WEBHOOKS_KEY, 64 hex characters. Webhooks' signing secrets are | |
| 31 | 32 | // sealed with it; without it none can be made. | |
| 32 | 33 | // Logs of a tenth of requests: every event is a delivery to look up. |