Simplify hourly scheduling and serialize all updater runs

This commit is contained in:
Matt Isenhower
2026-09-07 13:00:14 -07:00
parent a0badb5f39
commit 1b430d172a
7 changed files with 302 additions and 303 deletions

View File

@@ -31,7 +31,13 @@ export function createUpdaters(storage) {
export default async function updateAll(storage, { only } = {}) {
let results = [];
for (let updater of createUpdaters(storage)) {
let updaters = createUpdaters(storage);
if (only) {
let unknown = only.filter(name => !updaters.some(updater => updater.options.name === name));
if (!only.length || unknown.length)
throw new Error(`Unknown or empty updater selection: ${unknown.join(', ')}`);
}
for (let updater of updaters) {
let name = updater.options.name;
if (only && !only.includes(name))
continue;

View File

@@ -1,163 +1,205 @@
import { afterEach, describe, expect, it, vi } from 'vitest';
import { createExecutionContext, createScheduledController, env, runDurableObjectAlarm } from 'cloudflare:test';
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
import { createExecutionContext, createScheduledController, env, runDurableObjectAlarm, runInDurableObject } from 'cloudflare:test';
import worker from './src/index.mjs';
import { nextRunAt, GRACE_MS, HOUR_MS } from './src/schedule.mjs';
import { fakeSplatNet, setSessionEnvironment } from './fakeSplatNet.mjs';
function stub() {
return env.SCHEDULER.get(env.SCHEDULER.idFromName(`test-${crypto.randomUUID()}`));
}
// In the test runtime a due alarm may fire on its own (the fake Date is visible to the
// runtime too), so trigger it explicitly and then wait for the outcome either way.
async function runAlarmUntil(scheduler, done) {
await runDurableObjectAlarm(scheduler);
await vi.waitFor(async () => expect(done(await scheduler.status())).toBe(true), { timeout: 5000 });
return scheduler.status();
}
import { fakeSplatNet, setSessionEnvironment } from './fakeSplatNet.mjs';
setSessionEnvironment();
// The Durable Object runs in the test isolate, so stubbing global fetch (and Date) reaches it.
function splatnetDown() {
vi.stubGlobal('fetch', async () => new Response('down', { status: 503 }));
function network({ down = false, renderFails = false, beforeRequest } = {}) {
let splatnet = fakeSplatNet();
let renders = [];
vi.stubGlobal('fetch', async (input, init) => {
const url = new URL(input);
if (url.hostname === 'site.test') {
let object = await env.ASSETS.get(url.pathname.slice(1));
return object ? new Response(object.body) : new Response('missing', { status: 404 });
}
if (url.hostname === 'api.cloudflare.com') {
renders.push(JSON.parse(init.body));
return new Response(renderFails ? 'render failed' : new Uint8Array([1, 2]), { status: renderFails ? 503 : 200 });
}
await beforeRequest?.();
return down ? new Response('down', { status: 503 }) : splatnet(input, init);
});
return renders;
}
function splatnetUp() {
vi.stubGlobal('fetch', fakeSplatNet());
}
beforeEach(() => {
setSessionEnvironment();
process.env.SITE_URL = 'https://site.test';
process.env.CLOUDFLARE_ACCOUNT_ID = 'test';
process.env.CLOUDFLARE_BROWSER_RUN_API_TOKEN = 'test';
for (let name of ['BLUESKY_SERVICE', 'BLUESKY_IDENTIFIER', 'BLUESKY_PASSWORD'])
delete process.env[name];
});
afterEach(() => {
vi.unstubAllGlobals();
vi.useRealTimers();
});
describe('Scheduler', () => {
afterEach(() => {
vi.unstubAllGlobals();
vi.useRealTimers();
});
it('schedules the hourly job for the next :00:10 and is idempotent', async () => {
it('arms the hourly job for :00:10 and preserves an existing schedule', async () => {
let scheduler = stub();
let before = Date.now();
let first = await scheduler.ensureArmed();
expect(first.armed).toBe(true);
expect(first.hourlyAt).toBeGreaterThanOrEqual(nextRunAt(before));
expect((first.hourlyAt - GRACE_MS) % HOUR_MS).toBe(0);
expect(first.alarmAt).toBe(first.hourlyAt);
let second = await scheduler.ensureArmed();
expect(second).toEqual({ armed: false, hourlyAt: first.hourlyAt, alarmAt: first.hourlyAt });
let status = await scheduler.status();
expect(status.alarmAt).toBe(first.hourlyAt);
expect(status.pending).toEqual([]);
expect(status.lastRuns).toEqual({});
expect(await scheduler.ensureArmed()).toEqual({ ...first, armed: false });
});
it('runs a woken job immediately and then goes back to the hourly schedule', async () => {
splatnetUp();
it('persists a pause through watchdog calls and only resumes explicitly', async () => {
network();
let scheduler = stub();
let { hourlyAt } = await scheduler.ensureArmed();
let woke = await scheduler.wake('updaters');
expect(woke.accepted).toBe(true);
expect(woke.alarmAt).toBeLessThan(hourlyAt);
expect(woke.alarmAt - Date.now()).toBeLessThan(1000);
await scheduler.wake('updaters'); // duplicate collapses into the same run
let status = await runAlarmUntil(scheduler, s => s.lastRuns.updaters !== undefined);
expect(status.pending).toEqual([]);
expect(status.lastRuns.updaters.reason).toBe('wake');
expect(status.lastRuns.updaters.ok).toBe(true);
expect(status.lastRuns.updaters.driftMs).toBeNull();
expect(status.lastRuns.updaters.result.updaters).toHaveLength(8);
expect(status.hourlyAt).toBe(hourlyAt); // a wake run does not move the hourly schedule
expect(status.alarmAt).toBe(hourlyAt);
});
it('restores the hourly schedule if a wake runs before the object was ever armed', async () => {
splatnetUp();
let scheduler = stub();
await scheduler.wake('updaters');
let status = await runAlarmUntil(scheduler, s => s.lastRuns.updaters !== undefined);
expect(status.lastRuns.updaters.ok).toBe(true);
expect(status.hourlyAt).toBe(nextRunAt(status.lastRuns.updaters.firedAt));
expect(status.alarmAt).toBe(status.hourlyAt);
});
it('rejects unknown jobs without scheduling anything', async () => {
let scheduler = stub();
expect(await scheduler.wake('nope')).toEqual({ job: 'nope', accepted: false, error: 'Unknown job: nope' });
await scheduler.ensureArmed();
expect(await scheduler.pause()).toEqual({ ok: true, paused: true });
expect(await scheduler.ensureArmed()).toEqual({ armed: false, paused: true });
expect((await scheduler.status()).alarmAt).toBeNull();
expect(await scheduler.run()).toMatchObject({ ok: false, paused: true });
await scheduler.resume();
expect((await scheduler.status()).alarmAt).not.toBeNull();
expect((await scheduler.status()).paused).toBe(false);
});
it('retries a failed hourly run after a minute, then gives up until the next hour', async () => {
splatnetDown();
it('repairs a missing alarm without moving its due time', async () => {
let scheduler = stub();
let { hourlyAt } = await scheduler.ensureArmed();
let initial = await scheduler.ensureArmed();
await runInDurableObject(scheduler, async (instance, state) => state.storage.deleteAlarm());
expect(await scheduler.ensureArmed()).toEqual({ ...initial, armed: true });
});
vi.useFakeTimers({ toFake: ['Date'] });
vi.setSystemTime(hourlyAt + 5);
let status = await runAlarmUntil(scheduler, s => s.retries === 1);
expect(status.lastRuns.updaters).toMatchObject({ reason: 'hourly', ok: false, retries: 0 });
expect(status.lastRuns.updaters.driftMs).toBeGreaterThanOrEqual(5);
expect(status.lastRuns.updaters.driftMs).toBeLessThan(1000); // the faked clock still creeps a little
expect(status.lastRuns.updaters.error).toContain('updaters failed');
expect(status.retries).toBe(1);
expect(status.retryAt - (hourlyAt + 5 + 60 * 1000)).toBeGreaterThanOrEqual(0);
expect(status.retryAt - (hourlyAt + 5 + 60 * 1000)).toBeLessThan(2000);
expect(status.alarmAt).toBe(status.retryAt);
it('converts pending requests from the previous deployment into a full run', async () => {
network();
let scheduler = stub();
let hourlyAt = nextRunAt();
await runInDurableObject(scheduler, async (instance, state) => {
await state.storage.put('state', { hourlyAt, pending: ['updaters', 'posters'] });
});
await scheduler.ensureArmed();
let status = await runAlarmUntil(scheduler, s => s.lastRun !== null);
expect(status.lastRun.ok).toBe(true);
expect(status.hourlyAt).toBe(hourlyAt);
for (let attempt = 2; attempt <= 3; attempt++) {
vi.setSystemTime(status.retryAt);
status = await runAlarmUntil(scheduler, s => s.retries === attempt);
expect(status.lastRuns.updaters.reason).toBe('retry');
}
// Fourth failure exhausts retries: back to the next hour
vi.setSystemTime(status.retryAt);
status = await runAlarmUntil(scheduler, s => s.retryAt === null);
expect(status.retries).toBe(0);
expect(status.hourlyAt).toBe(hourlyAt + HOUR_MS);
expect(status.alarmAt).toBe(hourlyAt + HOUR_MS);
expect(status.retryAt).toBeNull();
});
it('moves to the next hour after a successful hourly run', async () => {
splatnetUp();
it('runs a full hourly pipeline, records completion timing, and schedules the next hour', async () => {
network();
let scheduler = stub();
let { hourlyAt } = await scheduler.ensureArmed();
vi.useFakeTimers({ toFake: ['Date'] });
vi.setSystemTime(hourlyAt + 1);
let status = await runAlarmUntil(scheduler, s => s.hourlyAt !== hourlyAt);
expect(status.lastRuns.updaters).toMatchObject({ reason: 'hourly', ok: true, colo: 'TEST' });
expect(status.lastRuns.posters).toMatchObject({ reason: 'hourly', ok: true });
expect(status.lastRuns.updaters.driftMs).toBeGreaterThanOrEqual(1);
expect(status.lastRuns.updaters.driftMs).toBeLessThan(1000);
expect(status.hourlyAt).toBe(hourlyAt + HOUR_MS);
expect(status.lastRun.ok).toBe(true);
expect(status.lastRun.updaters.updaters).toHaveLength(8);
expect(status.lastRun.social.ok).toBe(true);
expect(status.lastRun.driftMs).toBeGreaterThanOrEqual(1);
expect(status.alarmAt).toBe(hourlyAt + HOUR_MS);
});
});
describe('cron watchdog', () => {
it('arms the scheduler and does nothing else', async () => {
let calls = [];
let fakeEnv = {
SCHEDULER: {
idFromName: name => ({ name }),
get: (id, options) => {
calls.push({ id, options });
return { ensureArmed: async () => ({ armed: true, hourlyAt: 1, alarmAt: 1 }) };
},
},
};
await worker.scheduled(createScheduledController({ cron: '30 * * * *', scheduledTime: new Date }), fakeEnv, createExecutionContext());
expect(calls).toEqual([{ id: { name: 'schedules' }, options: { locationHint: 'wnam' } }]);
it('skips social after updater failure and gives up after three retries', async () => {
let renders = network({ down: true });
let scheduler = stub();
let { hourlyAt } = await scheduler.ensureArmed();
vi.useFakeTimers({ toFake: ['Date'] });
vi.setSystemTime(hourlyAt + 5);
let status = await runAlarmUntil(scheduler, s => s.retries === 1);
expect(status.lastRun.updaters.ok).toBe(false);
expect(status.lastRun.social).toMatchObject({ skipped: true, reason: 'updater-failed' });
expect(renders).toEqual([]);
for (let attempt = 2; attempt <= 3; attempt++) {
vi.setSystemTime(status.retryAt);
status = await runAlarmUntil(scheduler, s => s.retries === attempt);
}
vi.setSystemTime(status.retryAt);
status = await runAlarmUntil(scheduler, s => s.retryAt === null);
expect(status.retries).toBe(0);
expect(status.alarmAt).toBe(hourlyAt + HOUR_MS);
});
it('ignores cron expressions it has no action for', async () => {
let fakeEnv = { SCHEDULER: { idFromName: () => ({}), get: () => { throw new Error('should not be called'); } } };
await expect(worker.scheduled(createScheduledController({ cron: '0 0 1 1 *', scheduledTime: new Date }), fakeEnv, createExecutionContext())).resolves.toBeUndefined();
it('retries screenshot failures instead of reporting a successful social run', async () => {
network({ renderFails: true });
let scheduler = stub();
let { hourlyAt } = await scheduler.ensureArmed();
vi.useFakeTimers({ toFake: ['Date'] });
vi.setSystemTime(hourlyAt + 1);
let status = await runAlarmUntil(scheduler, s => s.retries === 1);
expect(status.lastRun.updaters.ok).toBe(true);
expect(status.lastRun.social.ok).toBe(false);
expect(status.lastRun.social.posts.some(post => post.ok === false)).toBe(true);
});
it('manual targeted repairs do not post or move the hourly schedule', async () => {
let renders = network();
let scheduler = stub();
let { hourlyAt } = await scheduler.ensureArmed();
let result = await scheduler.run({ only: ['Schedules'] });
expect(result.ok).toBe(true);
expect(result.updaters.updaters.map(u => u.name)).toEqual(['Schedules']);
expect(result.social.reason).toBe('targeted-update');
expect(renders).toEqual([]);
expect((await scheduler.status()).hourlyAt).toBe(hourlyAt);
expect((await scheduler.run({ only: ['typo'] })).ok).toBe(false);
expect((await scheduler.status()).busy).toBe(false);
});
it('rejects overlapping manual work explicitly and preserves an hourly alarm due during it', async () => {
let release;
let gate = new Promise(resolve => { release = resolve; });
let entered = false;
network({ beforeRequest: async () => { entered = true; await gate; } });
let scheduler = stub();
let { hourlyAt } = await scheduler.ensureArmed();
let first = scheduler.run({ only: ['Schedules'] });
try {
await vi.waitFor(() => expect(entered).toBe(true));
expect(await scheduler.run()).toMatchObject({ ok: false, busy: true });
expect(await scheduler.ensureArmed()).toMatchObject({ busy: true });
expect(await scheduler.pause()).toMatchObject({ ok: false, busy: true });
vi.useFakeTimers({ toFake: ['Date'] });
vi.setSystemTime(hourlyAt + 1);
await runDurableObjectAlarm(scheduler);
let status = await scheduler.status();
expect(status.hourlyAt).toBe(hourlyAt);
expect(status.alarmAt).not.toBeNull();
} finally {
release();
}
await first;
let status = await runAlarmUntil(scheduler, s => s.lastRun !== null);
expect(status.lastRun.updaters.ok).toBe(true);
expect(status.hourlyAt).toBe(hourlyAt + HOUR_MS);
});
});
describe('Worker routing', () => {
it('cron only calls the watchdog', async () => {
const ensureArmed = vi.fn(async () => ({ armed: true }));
const fakeEnv = { SCHEDULER: { idFromName: () => 'id', get: () => ({ ensureArmed }) } };
await worker.scheduled(createScheduledController({ cron: '30 * * * *' }), fakeEnv, createExecutionContext());
expect(ensureArmed).toHaveBeenCalledOnce();
});
it('authenticates manual runs and sends them through the DO with busy/failure HTTP statuses', async () => {
const run = vi.fn(async () => ({ ok: false, busy: true }));
const fakeEnv = { RUN_TOKEN: 'test-token', SCHEDULER: { idFromName: () => 'id', get: () => ({ run }) } };
const request = headers => new Request('https://worker.test/run?only=Schedules,Timeline', { method: 'POST', headers });
expect((await worker.fetch(request(), fakeEnv, createExecutionContext())).status).toBe(401);
expect(run).not.toHaveBeenCalled();
const headers = { Authorization: 'Bearer test-token' };
expect((await worker.fetch(request(headers), fakeEnv, createExecutionContext())).status).toBe(409);
expect(run).toHaveBeenCalledWith({ only: ['Schedules', 'Timeline'] });
run.mockResolvedValue({ ok: false });
expect((await worker.fetch(request(headers), fakeEnv, createExecutionContext())).status).toBe(502);
run.mockResolvedValue({ ok: true });
expect((await worker.fetch(request(headers), fakeEnv, createExecutionContext())).status).toBe(200);
});
});

View File

@@ -29,8 +29,6 @@ export function fakeSplatNet(routes = ROUTES) {
const url = new URL(input);
const headers = new Headers(init.headers);
requests.push({ path: url.pathname, language: headers.get('Accept-Language'), cookie: headers.get('Cookie') });
if (url.hostname === 'www.cloudflare.com')
return new Response('colo=TEST\n');
if (url.pathname.startsWith('/images/'))
return new Response(new Uint8Array([0x89, 0x50, 0x4e, 0x47]), { headers: { 'content-type': 'image/png' } });
const route = routes[url.pathname];

View File

@@ -1,163 +1,149 @@
// Durable Object that owns *when* and *where* updater jobs run.
//
// Why a Durable Object: Cron Triggers fire anywhere inside their minute and execute in
// whichever colo Cloudflare picks (placement hints only apply to fetch handlers). An
// object's alarm fires within milliseconds, and the object stays in the colo it was
// created in, next to the R2 buckets. So the cron trigger only wakes the object; the
// object does the work.
//
// Two ways work gets scheduled:
// - the hourly job re-arms itself for the next :00:10 after every run;
// - wake(job) asks for a job to run as soon as possible (used by cron-driven jobs).
// Both share one alarm: it is always set to the earliest thing that is due.
// One owner for the hourly update → social pipeline and authenticated manual runs.
// The alarm targets :00:10; the cron watchdog repairs a missing alarm. Alarms can be late.
import { DurableObject } from 'cloudflare:workers';
import { runUpdaters } from './updaters.mjs';
import { runPosters } from './posters.mjs';
import { nextRunAt } from './schedule.mjs';
import { createLogger, describeError } from './log.mjs';
import { currentColo } from './colo.mjs';
const RETRY_DELAY_MS = 60 * 1000;
const MAX_RETRIES = 3;
// What runs every hour, in this order: publish the data, then post about it.
export const HOURLY_JOBS = ['updaters', 'posters'];
// Jobs the object can run.
const JOBS = {
updaters: env => runUpdaters(env),
posters: env => runPosters(env),
};
const EMPTY_STATE = {
hourlyAt: null, // next scheduled run of the hourly job
retryAt: null, // when set, a failed hourly run is retried at this time instead
retries: 0,
pending: [], // jobs requested through wake(), run at the next alarm
lastRuns: {}, // per job: timing and outcome of the most recent run
};
export class Scheduler extends DurableObject {
/** Schedule the hourly job if it is not scheduled yet. Safe to call repeatedly (the cron watchdog does). */
async ensureArmed() {
let state = await this.#state();
let armed = state.hourlyAt === null;
if (armed) {
state.hourlyAt = nextRunAt(Date.now());
await this.#save(state);
}
return { armed, hourlyAt: state.hourlyAt, alarmAt: await this.#rearm(state) };
// Set synchronously before any await. RPCs may interleave with an alarm during external
// I/O. This is a lock, not durable job state: interrupted RPC callers receive an error,
// and interrupted alarms are retried by Cloudflare using the persisted schedule.
#running = false;
async #state() {
let saved = await this.ctx.storage.get('state') ?? {};
// Retain the existing hourly/retry schedule across deployment. Any old pending wake
// becomes one immediate full run, then the generic queue is retired.
return {
paused: saved.paused ?? false,
hourlyAt: saved.hourlyAt ?? null,
retryAt: saved.retryAt ?? (saved.pending?.length ? Date.now() : null),
retries: saved.retries ?? 0,
lastRun: saved.lastRun ?? null,
};
}
/** Run a job as soon as possible. Duplicate requests before the run collapse into one. */
async wake(job) {
if (!Object.hasOwn(JOBS, job))
return { job, accepted: false, error: `Unknown job: ${job}` };
async ensureArmed() {
if (this.#running)
return { armed: false, busy: true };
let state = await this.#state();
if (!state.pending.includes(job))
state.pending.push(job);
await this.#save(state);
return { job, accepted: true, alarmAt: await this.#rearm(state) };
if (state.paused)
return { armed: false, paused: true };
let armed = await this.ctx.storage.getAlarm() === null;
state.hourlyAt ??= nextRunAt();
await this.ctx.storage.put('state', state);
let alarmAt = state.retryAt ?? state.hourlyAt;
await this.ctx.storage.setAlarm(alarmAt);
return { armed, hourlyAt: state.hourlyAt, alarmAt };
}
async pause() {
if (this.#running)
return { ok: false, busy: true, error: 'Wait for the current run to finish before pausing.' };
let state = await this.#state();
state.paused = true;
await this.ctx.storage.put('state', state);
await this.ctx.storage.deleteAlarm();
return { ok: true, paused: true };
}
async resume() {
if (this.#running)
return { ok: false, busy: true };
let state = await this.#state();
state.paused = false;
await this.ctx.storage.put('state', state);
return { ok: true, ...await this.ensureArmed() };
}
async status() {
return { alarmAt: await this.ctx.storage.getAlarm(), ...await this.#state() };
return {
...await this.#state(),
alarmAt: await this.ctx.storage.getAlarm(),
lastManualRun: await this.ctx.storage.get('lastManualRun') ?? null,
busy: this.#running,
};
}
async #state() {
return { ...EMPTY_STATE, ...await this.ctx.storage.get('state') ?? {} };
// No detached work or persisted generic queue. A busy caller gets an explicit response
// and can retry; a successful response means the requested work finished.
async run({ only } = {}) {
if (this.#running)
return { ok: false, busy: true, error: 'An update is already running; retry later.' };
this.#running = true;
try {
if ((await this.#state()).paused)
return { ok: false, paused: true, error: 'Scheduler is paused; use /arm to resume.' };
let result = await this.#execute(only);
await this.ctx.storage.put('lastManualRun', result);
return result;
} finally {
this.#running = false;
await this.ensureArmed();
}
}
async #save(state) {
await this.ctx.storage.put('state', state);
}
/** Point the single alarm at the earliest due time. */
async #rearm(state) {
let candidates = [state.retryAt ?? state.hourlyAt, state.pending.length ? Date.now() : null]
.filter(time => time !== null);
if (!candidates.length)
return null;
let alarmAt = Math.min(...candidates);
await this.ctx.storage.setAlarm(alarmAt);
return alarmAt;
async #execute(only) {
let startedAt = Date.now();
let result;
try {
let updaters = await runUpdaters(this.env, { only });
// A targeted repair does not publish social posts from a partially refreshed dataset.
let social = !updaters.ok || only
? { ok: true, skipped: true, reason: only ? 'targeted-update' : 'updater-failed' }
: await runPosters(this.env);
result = { ok: updaters.ok && social.ok, updaters, social };
} catch (error) {
result = { ok: false, ...describeError(error) };
}
result = { ...result, startedAt, runMs: Date.now() - startedAt };
createLogger('pipeline')[result.ok ? 'info' : 'error']('Run finished', result);
return result;
}
async alarm(alarmInfo) {
let firedAt = Date.now();
let log = createLogger('alarm');
let colo = await currentColo();
let state = await this.#state();
// Work out what is due. The hourly jobs are due when their time (or the retry time) has
// come; pending jobs are due now. A pending request for an hourly job merges into it.
let due = [];
let hourlyScheduledFor = state.retryAt ?? state.hourlyAt;
let hourlyDue = hourlyScheduledFor !== null && firedAt >= hourlyScheduledFor;
if (hourlyDue)
for (let job of HOURLY_JOBS)
due.push({ job, reason: state.retryAt ? 'retry' : 'hourly', scheduledFor: hourlyScheduledFor });
for (let job of state.pending)
if (Object.hasOwn(JOBS, job) && !due.some(entry => entry.job === job))
due.push({ job, reason: 'wake', scheduledFor: null });
state.pending = [];
let hourlyOk = true;
for (let { job, reason, scheduledFor } of due) {
let startedAt = Date.now();
let run = {
job,
reason,
scheduledFor,
firedAt,
driftMs: scheduledFor === null ? null : firedAt - scheduledFor,
retryCount: alarmInfo?.retryCount ?? 0,
retries: state.retries,
colo,
};
// Errors are caught so the alarm is always re-armed; the platform's own alarm retries are
// capped and only cover the latest setAlarm(), so hourly retries are managed here.
try {
let result = await JOBS[job](this.env);
// A job reports partial failure by returning { ok: false } rather than throwing
run = { ...run, ok: result?.ok !== false, runMs: Date.now() - startedAt, result };
if (!run.ok)
run.error = `${job} failed: ${result.updaters?.filter(u => !u.ok).map(u => u.name).join(', ')}`;
} catch (error) {
run = { ...run, ok: false, runMs: Date.now() - startedAt, ...describeError(error) };
if (this.#running) {
// A manual run must not make us skip this hour. Leave the original due time intact.
await this.ctx.storage.setAlarm(Date.now() + RETRY_DELAY_MS);
return;
}
this.#running = true;
try {
let state = await this.#state();
if (state.paused) {
await this.ctx.storage.deleteAlarm();
return;
}
if (reason !== 'wake')
hourlyOk &&= run.ok;
state.lastRuns[job] = run;
log[run.ok ? 'info' : 'error']('Alarm run finished', run);
}
// Schedule the next hourly run, or a retry if any hourly job failed
if (hourlyDue) {
if (hourlyOk || state.retries >= MAX_RETRIES) {
state.hourlyAt = nextRunAt(Date.now());
state.retryAt = null;
state.retries = 0;
} else {
state.retryAt = Date.now() + RETRY_DELAY_MS;
state.retries += 1;
state.hourlyAt ??= nextRunAt();
let scheduledFor = state.retryAt ?? state.hourlyAt;
if (Date.now() >= scheduledFor) {
let result = await this.#execute();
state.lastRun = {
...result, scheduledFor,
driftMs: result.startedAt - scheduledFor,
retries: state.retries,
retryCount: alarmInfo?.retryCount ?? 0,
};
if (result.ok || state.retries >= MAX_RETRIES) {
state.hourlyAt = nextRunAt();
state.retryAt = null;
state.retries = 0;
} else {
state.retryAt = Date.now() + RETRY_DELAY_MS;
state.retries++;
}
}
// Storage failures escape so the platform retries. Schedule and state are saved
// together without external I/O in between.
await this.ctx.storage.put('state', state);
await this.ctx.storage.setAlarm(state.retryAt ?? state.hourlyAt);
} finally {
this.#running = false;
}
// This instance always owns the hourly job. If it has somehow been lost (for example the
// alarm was consumed by a wake before ensureArmed() ever ran), restore it here rather than
// waiting for the cron watchdog.
if (state.hourlyAt === null) {
state.hourlyAt = nextRunAt(Date.now());
log.warn('Hourly schedule was missing; restored', { hourlyAt: state.hourlyAt });
}
await this.#save(state);
let alarmAt = await this.#rearm(state);
log.info('Alarm re-armed', { alarmAt, hourlyAt: state.hourlyAt, retryAt: state.retryAt, ran: due.map(entry => entry.job) });
}
}

View File

@@ -1,12 +0,0 @@
// Which Cloudflare colo this invocation is running in, for comparing scheduling paths.
// A subrequest to a Cloudflare-fronted host is answered from the local colo, so the
// trace endpoint reports where this Worker's code is executing.
export async function currentColo() {
try {
let response = await fetch('https://www.cloudflare.com/cdn-cgi/trace');
let text = await response.text();
return text.match(/^colo=(\w+)$/m)?.[1] ?? null;
} catch {
return null;
}
}

View File

@@ -1,6 +1,4 @@
import { withSentry, instrumentDurableObjectWithSentry } from '@sentry/cloudflare';
import { runUpdaters } from './updaters.mjs';
import { runPosters } from './posters.mjs';
import { Scheduler as SchedulerClass } from './Scheduler.mjs';
import { createLogger, describeError } from './log.mjs';
@@ -10,19 +8,11 @@ const sentryOptions = env => ({ dsn: env.SENTRY_DSN, tracesSampleRate: 0 });
export const Scheduler = instrumentDurableObjectWithSentry(sentryOptions, SchedulerClass);
// The object is created next to the R2 buckets (western North America). Only the first
// get() for an object honors the hint; after that it stays where it is.
// Best-effort initial placement near the buckets; correctness does not depend on it.
function scheduler(env) {
return env.SCHEDULER.get(env.SCHEDULER.idFromName('schedules'), { locationHint: 'wnam' });
}
// Cron Triggers never do work themselves: they wake the Scheduler, which runs jobs in
// place. Each entry maps a cron expression from wrangler.jsonc to a call on the object.
const CRON_ACTIONS = {
'30 * * * *': scheduler => scheduler.ensureArmed(), // watchdog: the hourly alarm must always be armed
// Example of a cron-driven job: '15 3 * * *': scheduler => scheduler.wake('some-daily-job'),
};
function timingSafeEqual(a, b) {
let encoder = new TextEncoder;
let left = encoder.encode(a);
@@ -43,31 +33,19 @@ function isAuthorized(request, env) {
export default withSentry(sentryOptions, {
async scheduled(controller, env, ctx) {
let log = createLogger('cron');
let action = CRON_ACTIONS[controller.cron];
if (!action) {
log.warn('No action for cron expression', { cron: controller.cron });
return;
}
try {
log.info('Cron action finished', { cron: controller.cron, result: await action(scheduler(env)) });
log.info('Cron action finished', { cron: controller.cron, result: await scheduler(env).ensureArmed() });
} catch (error) {
log.error('Cron action failed', { cron: controller.cron, ...describeError(error) });
throw error; // Mark the invocation as failed in Workers metrics
}
},
// Authenticated operator endpoints, "Authorization: Bearer <RUN_TOKEN>":
// POST /run[?only=Name,Name] run the updaters in this invocation and return the summary
// POST /post run the social posters in this invocation
// POST /wake?job=NAME ask the Scheduler to run a job as soon as possible
// POST /arm schedule the hourly job if it is not scheduled
// GET /status alarm state and the last run of each job
// GET /list?prefix=data/ keys in the public bucket under a prefix (first 1000)
// Manual runs use the same owner as alarms. Targeted runs refresh data only.
async fetch(request, env, ctx) {
let url = new URL(request.url);
let route = `${request.method} ${url.pathname}`;
if (!['POST /run', 'POST /post', 'POST /wake', 'POST /arm', 'GET /status', 'GET /list'].includes(route))
if (!['POST /run', 'POST /arm', 'POST /pause', 'GET /status', 'GET /list'].includes(route))
return new Response('Not found', { status: 404 });
if (!isAuthorized(request, env))
return new Response('Unauthorized', { status: 401 });
@@ -76,16 +54,17 @@ export default withSentry(sentryOptions, {
switch (route) {
case 'POST /run': {
let only = url.searchParams.get('only')?.split(',').map(name => name.trim()).filter(Boolean);
return Response.json(await runUpdaters(env, { only }));
let result = await scheduler(env).run({ only });
return Response.json(result, { status: result.busy || result.paused ? 409 : result.ok ? 200 : 502 });
}
case 'POST /post':
return Response.json(await runPosters(env));
case 'POST /wake': {
let woke = await scheduler(env).wake(url.searchParams.get('job') ?? '');
return Response.json({ ok: woke.accepted, ...woke }, { status: woke.accepted ? 200 : 400 });
case 'POST /arm': {
let result = await scheduler(env).resume();
return Response.json(result, { status: result.busy ? 409 : 200 });
}
case 'POST /pause': {
let result = await scheduler(env).pause();
return Response.json(result, { status: result.busy ? 409 : 200 });
}
case 'POST /arm':
return Response.json({ ok: true, ...await scheduler(env).ensureArmed() });
case 'GET /status':
return Response.json({ ok: true, ...await scheduler(env).status() });
case 'GET /list': {

View File

@@ -12,7 +12,7 @@
"workers_dev": true,
"preview_urls": false,
"triggers": {
// Cron Triggers only wake the Scheduler Durable Object (see src/index.mjs CRON_ACTIONS).
// Cron only checks the Scheduler alarm; it does not execute updater work.
// The hourly update runs from the object's own alarm at :00:10; this cron is the watchdog
// that re-arms that alarm if it is ever lost. Minute 30 keeps it clear of the hourly run.
"crons": ["30 * * * *"]
@@ -33,7 +33,7 @@
"r2_buckets": [
{
// Shadow mode: writes go to the dev bucket while the container still owns production.
// Cutover = point this at "splatoon2-ink-assets".
// See README for the full cutover procedure before switching to "splatoon2-ink-assets".
"binding": "ASSETS",
"bucket_name": "splatoon2-ink-dev-assets"
},