g1t/deploy/self-host/scheduler.mjs

133 lines5,769 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member1#!/usr/bin/env node
2// Runs the services' cron triggers in a self-hosted g1t. `wrangler dev`
3// (workerd) never fires a Worker's `scheduled` handler on its own, so this
4// does: once a minute it asks Wrangler's local API to run the handler of
5// each service whose cron matches the time (UTC), the same handler hosted
6// g1t's Cron Triggers run.
7//
8// Which services, and their crons, are in schedules.json, written by
9// configs.mjs from each service's wrangler.jsonc (`triggers.crons`), for
10// the services whose crons are safe to run here (SELF_HOST_CRONS there).
11// Wrangler's local API answers only requests addressed to localhost, so
12// this runs inside the g1t container, beside it (start.sh).
13//
14// Usage: node scheduler.mjs <schedules.json> [base URL, default http://127.0.0.1:8787]
Merge branch 'worktree-agent-aaf03bdceac799c89'15// node scheduler.mjs --once <schedules.json> [base URL]
16//
17// --once runs every service's every cron now, says how each went, and exits
18// non-zero if any failed. A handler still running from the minute before is
19// not started again, and one is given up on after ten minutes.
Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member20
21import { readFileSync } from "node:fs";
22
23const FIELDS = [
24 { min: 0, max: 59 }, // minute
25 { min: 0, max: 23 }, // hour
26 { min: 1, max: 31 }, // day of the month
27 { min: 1, max: 12 }, // month
28 { min: 0, max: 7 }, // day of the week, 0 and 7 both Sunday
29];
30
31/** The values one cron field allows, or null when it is malformed. */
32export function fieldValues(text, { min, max }) {
33 const values = new Set();
34 for (const part of text.split(",")) {
35 const [range, stepText] = part.split("/");
36 const step = stepText === undefined ? 1 : Number(stepText);
37 if (!Number.isInteger(step) || step < 1) return null;
38 let from;
39 let to;
40 if (range === "*") {
41 [from, to] = [min, max];
42 } else if (range.includes("-")) {
43 [from, to] = range.split("-").map(Number);
44 } else {
45 from = Number(range);
46 to = stepText === undefined ? from : max;
47 }
48 if (!Number.isInteger(from) || !Number.isInteger(to) || from < min || to > max || from > to) return null;
49 for (let value = from; value <= to; value += step) values.add(value);
50 }
51 return values;
52}
53
54/** Whether a five-field cron expression matches `date` (UTC, to the minute). */
55export function matches(cron, date) {
56 const fields = cron.trim().split(/\s+/);
57 if (fields.length !== 5) return false;
58 const sets = fields.map((field, i) => fieldValues(field, FIELDS[i]));
59 if (sets.some((set) => set === null)) return false;
60 const [minutes, hours, days, months, weekdays] = sets;
61 const weekday = date.getUTCDay();
62 const dayMatches = days.has(date.getUTCDate());
63 const weekdayMatches = weekdays.has(weekday) || (weekday === 0 && weekdays.has(7));
64 // As cron does: when both day fields are restricted, either may match.
65 const day =
66 fields[2] !== "*" && fields[4] !== "*" ? dayMatches || weekdayMatches : dayMatches && weekdayMatches;
67 return minutes.has(date.getUTCMinutes()) && hours.has(date.getUTCHours()) && months.has(date.getUTCMonth() + 1) && day;
68}
69
70/** The handlers due at `date`: `{ worker, cron }` for each cron that matches. */
71export function due(schedules, date) {
72 return schedules.flatMap(({ worker, crons }) => crons.filter((cron) => matches(cron, date)).map((cron) => ({ worker, cron })));
73}
74
Merge branch 'worktree-agent-aaf03bdceac799c89'75/** Handlers still running, by worker and cron: a slow one is not started again on top of itself. */
76const running = new Set();
77
78/** Runs one handler through Wrangler's local API; whether it ran and said ok. */
79export async function run(base, { worker, cron }, { fetch: send = fetch, timeoutMs = 10 * 60_000 } = {}) {
80 const key = `${worker} ${cron}`;
81 if (running.has(key)) {
82 console.error(`scheduler: ${worker} (${cron}) is still running from the last time; skipped`);
83 return false;
84 }
85 running.add(key);
Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member86 try {
Merge branch 'worktree-agent-aaf03bdceac799c89'87 const answer = await send(`${base}/cdn-cgi/local/explorer/api/local/scheduled?worker=${encodeURIComponent(worker)}`, {
Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member88 method: "POST",
89 headers: { "content-type": "application/json" },
90 body: JSON.stringify({ cron }),
Merge branch 'worktree-agent-aaf03bdceac799c89'91 signal: AbortSignal.timeout(timeoutMs),
Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member92 });
93 const body = await answer.json().catch(() => ({}));
94 if (!answer.ok || !body.success || body.result?.outcome !== "ok") {
95 console.error(`scheduler: ${worker} (${cron}): ${answer.status} ${JSON.stringify(body.errors ?? body.result ?? body)}`);
Merge branch 'worktree-agent-aaf03bdceac799c89'96 return false;
Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member97 }
Merge branch 'worktree-agent-aaf03bdceac799c89'98 return true;
Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member99 } catch (error) {
100 console.error(`scheduler: ${worker} (${cron}) could not be run: ${error.message}`);
Merge branch 'worktree-agent-aaf03bdceac799c89'101 return false;
102 } finally {
103 running.delete(key);
Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member104 }
105}
106
107if (process.argv[1]?.replaceAll("\\", "/").endsWith("deploy/self-host/scheduler.mjs")) {
Merge branch 'worktree-agent-aaf03bdceac799c89'108 const once = process.argv.includes("--once");
109 const [file, given] = process.argv.slice(2).filter((arg) => arg !== "--once");
110 const schedules = JSON.parse(readFileSync(file, "utf8"));
111 const base = (given ?? "http://127.0.0.1:8787").replace(/\/$/, "");
112 if (once) {
113 // Every cron of every service, now, one after another: a check that
114 // each handler runs (smoke.sh), not a schedule.
115 let failed = 0;
116 for (const { worker, crons } of schedules) {
117 for (const cron of crons) {
118 const ok = await run(base, { worker, cron });
119 console.log(`${ok ? "ok " : "FAIL"} ${worker} (${cron})`);
120 if (!ok) failed++;
121 }
122 }
123 process.exit(failed ? 1 : 0);
124 }
Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member125 console.log(`scheduler: ${schedules.map((s) => `${s.worker} ${s.crons.join(", ")}`).join("; ")}`);
126 const tick = () => {
127 const now = new Date();
128 for (const job of due(schedules, now)) run(base, job);
129 // The next whole minute, a second in.
130 setTimeout(tick, 60_000 - (now.getTime() % 60_000) + 1_000);
131 };
132 setTimeout(tick, 60_000 - (Date.now() % 60_000) + 1_000);
133}

This file's history is long; its oldest lines are credited to the oldest commit read.