mirror of
https://github.com/misenhower/splatoon2.ink.git
synced 2026-09-29 04:37:00 -05:00
Simplify run helpers and clarify scheduler control flow
This commit is contained in:
@@ -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++;
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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));
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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) };
|
||||
|
||||
Reference in New Issue
Block a user