diff --git a/src/app/updater/updateAll.js b/src/app/updater/updateAll.js index 2d4bf37..abfaf33 100644 --- a/src/app/updater/updateAll.js +++ b/src/app/updater/updateAll.js @@ -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; diff --git a/workers/updater/Scheduler.spec.mjs b/workers/updater/Scheduler.spec.mjs index 00d7bce..315895c 100644 --- a/workers/updater/Scheduler.spec.mjs +++ b/workers/updater/Scheduler.spec.mjs @@ -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); }); }); diff --git a/workers/updater/fakeSplatNet.mjs b/workers/updater/fakeSplatNet.mjs index 357b311..252fb4d 100644 --- a/workers/updater/fakeSplatNet.mjs +++ b/workers/updater/fakeSplatNet.mjs @@ -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]; diff --git a/workers/updater/src/Scheduler.mjs b/workers/updater/src/Scheduler.mjs index 6264cc4..40715d7 100644 --- a/workers/updater/src/Scheduler.mjs +++ b/workers/updater/src/Scheduler.mjs @@ -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) }); } } diff --git a/workers/updater/src/colo.mjs b/workers/updater/src/colo.mjs deleted file mode 100644 index a1af155..0000000 --- a/workers/updater/src/colo.mjs +++ /dev/null @@ -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; - } -} diff --git a/workers/updater/src/index.mjs b/workers/updater/src/index.mjs index 0c10ff5..24dac66 100644 --- a/workers/updater/src/index.mjs +++ b/workers/updater/src/index.mjs @@ -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 ": - // 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': { diff --git a/workers/updater/wrangler.jsonc b/workers/updater/wrangler.jsonc index 627a45e..85ec317 100644 --- a/workers/updater/wrangler.jsonc +++ b/workers/updater/wrangler.jsonc @@ -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" },