From 6c44b548b65007bb2e00219315a088a26d6f9800 Mon Sep 17 00:00:00 2001 From: Matt Isenhower Date: Mon, 7 Sep 2026 13:48:06 -0700 Subject: [PATCH] Simplify run helpers and clarify scheduler control flow --- scripts/compare-data.mjs | 38 ++++++++++++++------------ src/app/log.js | 41 ++++++++++++++++++++++------- src/common/storage/BucketStorage.js | 4 --- test/social/support.js | 21 ++++++++++++--- workers/updater/README.md | 5 ++-- workers/updater/src/Scheduler.mjs | 22 +++++++++++----- 6 files changed, 87 insertions(+), 44 deletions(-) diff --git a/scripts/compare-data.mjs b/scripts/compare-data.mjs index 5c91cd2..bf9b7ae 100644 --- a/scripts/compare-data.mjs +++ b/scripts/compare-data.mjs @@ -8,6 +8,7 @@ import fs from 'node:fs'; import path from 'node:path'; +import stringify from 'json-stable-stringify'; const [base, current] = process.argv.slice(2); if (!base || !current) { @@ -15,32 +16,28 @@ if (!base || !current) { process.exit(2); } -const canon = value => Array.isArray(value) - ? value.map(canon) - : (value && typeof value === 'object') - ? Object.fromEntries(Object.keys(value).sort().map(key => [key, canon(value[key])])) - : value; - function firstDifference(a, b, at = '$') { if (Array.isArray(a) && Array.isArray(b)) { if (a.length !== b.length) return `${at}.length ${a.length} vs ${b.length}`; // Festival result lists come back in random order if (a.length && a[0] && typeof a[0] === 'object' && 'festival_id' in a[0]) { - let key = item => JSON.stringify(canon(item)); - let sa = a.map(key).sort(), sb = b.map(key).sort(); - return sa.every((v, i) => v === sb[i]) ? null : `${at}: set of items differs`; + let leftItems = a.map(item => stringify(item)).sort(); + let rightItems = b.map(item => stringify(item)).sort(); + return leftItems.every((item, index) => item === rightItems[index]) ? null : `${at}: set of items differs`; } for (let i = 0; i < a.length; i++) { - let d = firstDifference(a[i], b[i], `${at}[${i}]`); - if (d) return d; + let difference = firstDifference(a[i], b[i], `${at}[${i}]`); + if (difference) + return difference; } return null; } if (a && b && typeof a === 'object' && typeof b === 'object') { for (let key of new Set([...Object.keys(a), ...Object.keys(b)])) { - let d = firstDifference(a[key], b[key], `${at}.${key}`); - if (d) return d; + let difference = firstDifference(a[key], b[key], `${at}.${key}`); + if (difference) + return difference; } return null; } @@ -59,15 +56,22 @@ for (let file of walk(base)) { problems++; continue; } - let a = fs.readFileSync(file), b = fs.readFileSync(other); + let a = fs.readFileSync(file); + let b = fs.readFileSync(other); if (a.equals(b)) continue; if (relative.endsWith('.json')) { - let d = firstDifference(JSON.parse(a), JSON.parse(b)); - if (d) { console.log(`${relative}: ${d}`); problems++; } + let difference = firstDifference(JSON.parse(a), JSON.parse(b)); + if (difference) { + console.log(`${relative}: ${difference}`); + problems++; + } } else if (relative.endsWith('.ics')) { let strip = s => s.toString().split(/\r?\n/).filter(line => !line.startsWith('DTSTAMP')).join('\n'); - if (strip(a) !== strip(b)) { console.log(`${relative}: calendar content differs`); problems++; } + if (strip(a) !== strip(b)) { + console.log(`${relative}: calendar content differs`); + problems++; + } } else { console.log(`${relative}: binary content differs`); problems++; diff --git a/src/app/log.js b/src/app/log.js index 33b0987..8fb2dce 100644 --- a/src/app/log.js +++ b/src/app/log.js @@ -9,27 +9,48 @@ const encoder = new TextEncoder(); export function createRunLog(secrets = []) { let snapshot = { lines: [], omitted: 0 }; + let context = { + snapshot, + bytes: 0, + secrets: secrets.filter(value => typeof value === 'string' && value.length >= 6), + }; + return { snapshot, - run: callback => currentRun.run({ snapshot, bytes: 0, secrets: secrets.filter(value => typeof value === 'string' && value.length >= 6) }, callback), + run: callback => currentRun.run(context, callback), }; } export function logMessage(level, message, fields) { - if (fields === undefined) console[level](message); - else if (fields.updater) console[level]({ message, ...fields }); - else console[level](message, fields); + if (fields === undefined) { + console[level](message); + } else if (fields.updater) { + console[level]({ message, ...fields }); + } else { + console[level](message, fields); + } + let context = currentRun.getStore(); - if (!context) return; + if (!context) + return; + let text = message instanceof Error ? message.message : String(message); // Keep readable progress messages; the full structured result has its own JSON view. - if (fields?.updater) text = `[${fields.updater}] ${text}`; - if (fields?.error) text += `: ${fields.error}`; - if (fields?.attempt) text += ` (retry ${fields.attempt})`; - for (let secret of context.secrets) text = text.replaceAll(secret, '[redacted]'); + if (fields?.updater) + text = `[${fields.updater}] ${text}`; + if (fields?.error) + text += `: ${fields.error}`; + if (fields?.attempt) + text += ` (retry ${fields.attempt})`; + for (let secret of context.secrets) + text = text.replaceAll(secret, '[redacted]'); text = text.replace(/Bearer\s+[^\s,;]+/gi, 'Bearer [redacted]'); const { snapshot } = context; - snapshot.lines.push({ at: Date.now(), level, text: text.length > MAX_LINE_LENGTH ? text.slice(0, MAX_LINE_LENGTH) + '…' : text }); + snapshot.lines.push({ + at: Date.now(), + level, + text: text.length > MAX_LINE_LENGTH ? text.slice(0, MAX_LINE_LENGTH) + '…' : text, + }); context.bytes += encoder.encode(JSON.stringify(snapshot.lines.at(-1))).length; while (snapshot.lines.length > MAX_LINES || context.bytes > MAX_BYTES) { context.bytes -= encoder.encode(JSON.stringify(snapshot.lines.shift())).length; diff --git a/src/common/storage/BucketStorage.js b/src/common/storage/BucketStorage.js index d704eb1..30f188c 100644 --- a/src/common/storage/BucketStorage.js +++ b/src/common/storage/BucketStorage.js @@ -22,10 +22,6 @@ export default class BucketStorage { this.#bucket = bucket; } - get bucket() { - return this.#bucket; - } - async #listing(directory) { if (!this.#listings.has(directory)) { let keys = new Set; diff --git a/test/social/support.js b/test/social/support.js index 4a8e1f5..35b6c87 100644 --- a/test/social/support.js +++ b/test/social/support.js @@ -1,8 +1,14 @@ import { MemoryBucket, BucketStorage } from '../../src/common/storage/index.js'; export function storage() { - const publicBucket = new MemoryBucket, privateBucket = new MemoryBucket; - return { publicBucket, privateBucket, publicStorage: new BucketStorage(publicBucket), privateStorage: new BucketStorage(privateBucket) }; + const publicBucket = new MemoryBucket; + const privateBucket = new MemoryBucket; + return { + publicBucket, + privateBucket, + publicStorage: new BucketStorage(publicBucket), + privateStorage: new BucketStorage(privateBucket), + }; } /** A fresh per-run view over the same buckets, as each real run gets (the cache is per run). */ @@ -18,9 +24,16 @@ export function fakeClient(key, { canSend = true, fail = false } = {}) { name: key[0].toUpperCase() + key.slice(1), sent, canSend: async () => canSend, - send: async status => { if (fail) throw new Error(`${key} is down`); sent.push(status); }, + send: async status => { + if (fail) + throw new Error(`${key} is down`); + sent.push(status); + }, }; } -export const json = async (bucket, key) => { const object = await bucket.get(key); return object ? object.json() : null; }; +export async function json(bucket, key) { + const object = await bucket.get(key); + return object ? object.json() : null; +} export const seed = (bucket, key, value) => bucket.put(key, JSON.stringify(value)); diff --git a/workers/updater/README.md b/workers/updater/README.md index 02649ca..2930430 100644 --- a/workers/updater/README.md +++ b/workers/updater/README.md @@ -19,8 +19,9 @@ not depend on where the object runs. There is no per-run colo probe. - The minute-30 cron checks/re-arms the alarm. The object may sleep between alarms; the watchdog repairs scheduling, rather than keeping a process alive. - Manual runs use the same object. A request during an active run returns HTTP - 409; retry it later. There is no background operator queue and no `/wake` or - `/post` endpoint. An hourly alarm that encounters a manual run remains due. + 409; retry it later. The admin panel can persist one background manual request. + There is no `/wake` or `/post` endpoint. An hourly alarm that encounters a + manual run remains due. - `POST /run` runs the complete pipeline and waits for its result. With `only`, it repairs the named updaters and skips social posting. Unknown names fail. - The existing object name and hourly/retry state survive deployment. Pending diff --git a/workers/updater/src/Scheduler.mjs b/workers/updater/src/Scheduler.mjs index c810528..a244ef6 100644 --- a/workers/updater/src/Scheduler.mjs +++ b/workers/updater/src/Scheduler.mjs @@ -13,8 +13,8 @@ export const MANUAL_MODES = ['data', 'social', 'both']; export class Scheduler extends DurableObject { // 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. + // I/O. Direct /run callers wait for completion; admin requests are stored separately + // and run from an alarm. Both use this same lock. #running = false; #activeRun = null; @@ -69,13 +69,14 @@ export class Scheduler extends DurableObject { } async status() { + let pendingManual = await this.#pendingManual(); return { ...await this.#state(), alarmAt: await this.ctx.storage.getAlarm(), lastManualRun: await this.ctx.storage.get('lastManualRun') ?? null, - pendingManual: await this.#pendingManual(), + pendingManual, activeRun: this.#activeRun, - busy: this.#running || !!await this.#pendingManual(), + busy: this.#running || !!pendingManual, }; } @@ -147,9 +148,16 @@ export class Scheduler extends DurableObject { ? { ok: true, skipped: true } : await runUpdaters(this.env, { only }); // A targeted repair does not publish social posts from a partially refreshed dataset. - let social = !updaters.ok || only || mode === 'data' - ? { ok: true, skipped: true, reason: mode === 'data' ? 'data-only' : only ? 'targeted-update' : 'updater-failed' } - : await runPosters(this.env); + let social; + if (mode === 'data') { + social = { ok: true, skipped: true, reason: 'data-only' }; + } else if (only) { + social = { ok: true, skipped: true, reason: 'targeted-update' }; + } else if (!updaters.ok) { + social = { ok: true, skipped: true, reason: 'updater-failed' }; + } else { + social = await runPosters(this.env); + } result = { ok: updaters.ok && social.ok, updaters, social }; } catch (error) { result = { ok: false, ...describeError(error) };