Run the social posters in the Worker after the updaters

The Scheduler's hourly alarm now runs two jobs in order, updaters then
posters, matching the container; a failure in either schedules the
retry. POST /post runs the posters by hand and wake?job=posters through
the alarm. The posters read process.env, which Workers populate from the
bindings: SITE_URL and CLOUDFLARE_ACCOUNT_ID are vars, the Browser
Rendering token and the social credentials are secrets. Without social
credentials the posters only render and save the public images, which is
the shadow mode until cutover.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Matt Isenhower
2026-09-06 19:06:48 -07:00
parent 55c1d3ae81
commit 9258447df7
7 changed files with 143 additions and 22 deletions

View File

@@ -12,10 +12,11 @@ fetch handlers. So Cron Triggers never do work here. They call the `Scheduler`
Durable Object, which runs jobs from its own alarm, in place, next to the R2
buckets (the stub is created with a `wnam` location hint).
- **Hourly job** (`updaters`): the full updater run from an alarm at :00:10,
re-armed for the next hour after each run. A run with a failed updater is
retried after a minute, up to three times, before falling back to the next
hour. Alarms have measured about 1 ms of drift.
- **Hourly jobs** (`updaters`, then `posters`): the full updater run followed by
the social posters, from an alarm at :00:10, re-armed for the next hour after
each run. A run with a failed job is retried after a minute, up to three
times, before falling back to the next hour. Alarms have measured about 1 ms
of drift.
- **Woken jobs**: `wake(job)` asks the object to run a job as soon as possible.
The scheduled handler maps each cron expression in `CRON_ACTIONS` to a call on
the object; this is how a cron-driven job is expressed.
@@ -27,6 +28,15 @@ Every run logs a structured summary (`driftMs`, `runMs`, per-updater results,
the colo) under `updater: "alarm"`, visible in the Worker's Observability tab.
`GET /status` returns the last run of each job.
## Social posters
The posters (`src/app/twitter`) read the published data from the public bucket,
keep their state (last post times per platform, the previous Salmon Run shift)
in the private bucket, render screenshots of the deployed site through
Cloudflare Browser Rendering, and post to Bluesky and Twitter. With no social
credentials set they only render and save the public images
(`twitter-images/`), which is the shadow mode used before cutover.
## Errors
The shared updater code reports through `@sentry/core`; this Worker wraps its
@@ -48,7 +58,10 @@ npx wrangler secret put NINTENDO_SESSION_ID_EU --config workers/updater/wrangler
npx wrangler secret put NINTENDO_SESSION_ID_JP --config workers/updater/wrangler.jsonc
npx wrangler secret put SPLATNET_USER_AGENT --config workers/updater/wrangler.jsonc
npx wrangler secret put RUN_TOKEN --config workers/updater/wrangler.jsonc
npx wrangler secret put CLOUDFLARE_BROWSER_RUN_API_TOKEN --config workers/updater/wrangler.jsonc # screenshots
npx wrangler secret put SENTRY_DSN --config workers/updater/wrangler.jsonc # optional
# At cutover, the social credentials: BLUESKY_SERVICE, BLUESKY_IDENTIFIER, BLUESKY_PASSWORD,
# TWITTER_CONSUMER_KEY, TWITTER_CONSUMER_SECRET, TWITTER_ACCESS_TOKEN_KEY, TWITTER_ACCESS_TOKEN_SECRET
```
The updaters read the SplatNet secrets through `process.env`, which Workers
@@ -72,7 +85,8 @@ Authenticated operator endpoints, all requiring `Authorization: Bearer $UPDATER_
BASE=https://splatoon2-ink-updater.<subdomain>.workers.dev
curl -X POST -H "Authorization: Bearer $UPDATER_RUN_TOKEN" "$BASE/run" # run every updater now
curl -X POST -H "Authorization: Bearer $UPDATER_RUN_TOKEN" "$BASE/run?only=Schedules,Timeline" # or some of them
curl -X POST -H "Authorization: Bearer $UPDATER_RUN_TOKEN" "$BASE/wake?job=updaters" # ask the Scheduler to run now
curl -X POST -H "Authorization: Bearer $UPDATER_RUN_TOKEN" "$BASE/post" # run the social posters now
curl -X POST -H "Authorization: Bearer $UPDATER_RUN_TOKEN" "$BASE/wake?job=updaters" # ask the Scheduler to run a job now (or job=posters)
curl -X POST -H "Authorization: Bearer $UPDATER_RUN_TOKEN" "$BASE/arm" # schedule the hourly job if it is not scheduled
curl -H "Authorization: Bearer $UPDATER_RUN_TOKEN" "$BASE/status" # alarm time, hourly/retry state, last run per job
curl -H "Authorization: Bearer $UPDATER_RUN_TOKEN" "$BASE/list?prefix=data/" # keys in the public bucket under a prefix (wrangler cannot list objects)

View File

@@ -132,6 +132,7 @@ describe('Scheduler', () => {
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);

View File

@@ -0,0 +1,56 @@
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
import { env } from 'cloudflare:test';
import { runPosters } from './src/posters.mjs';
import { fakeSplatNet, ROUTES } from './fakeSplatNet.mjs';
import { getTopOfCurrentHour } from '../../src/common/time.js';
// The posters (src/app/twitter) running inside workerd: data from R2, screenshots from a
// stubbed Browser Rendering endpoint, no social credentials (shadow mode).
const PNG = new Uint8Array([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a]);
describe('runPosters', () => {
let renders;
beforeEach(async () => {
process.env.SITE_URL = 'https://example.test';
process.env.CLOUDFLARE_ACCOUNT_ID = 'acct';
process.env.CLOUDFLARE_BROWSER_RUN_API_TOKEN = 'token';
for (let name of ['BLUESKY_SERVICE', 'BLUESKY_IDENTIFIER', 'BLUESKY_PASSWORD', 'TWITTER_CONSUMER_KEY'])
delete process.env[name];
let now = getTopOfCurrentHour();
let rotation = { ...ROUTES['/api/schedules']().regular[0], start_time: now, end_time: now + 7200 };
await env.ASSETS.put('data/schedules.json', JSON.stringify({ regular: [rotation], gachi: [rotation], league: [rotation] }));
await env.ASSETS.put('data/festivals.json', JSON.stringify({ na: { festivals: [], results: [] }, eu: { festivals: [], results: [] }, jp: { festivals: [], results: [] } }));
await env.ASSETS.put('data/coop-schedules.json', JSON.stringify({ schedules: [], details: [] }));
await env.ASSETS.put('data/timeline.json', JSON.stringify({ coop: null, weapon_availability: null }));
await env.ASSETS.put('data/merchandises.json', JSON.stringify({ merchandises: [{ end_time: now + 3600, gear: { name: 'Hat' }, skill: { name: 'Skill' } }] }));
renders = [];
let splatnet = fakeSplatNet();
vi.stubGlobal('fetch', async (input, init) => {
let url = new URL(input);
if (url.hostname === 'api.cloudflare.com') {
renders.push(JSON.parse(init.body));
return new Response(PNG, { headers: { 'content-type': 'image/png' } });
}
return splatnet(input, init);
});
});
afterEach(() => vi.unstubAllGlobals());
it('renders the public images for the hour without posting anywhere', async () => {
let summary = await runPosters(env);
expect(summary.ok).toBe(true);
expect(summary.clients).toEqual([]);
expect(renders.map(r => r.url)).toEqual([
`https://example.test/screenshots.html#/schedules/${getTopOfCurrentHour()}`,
`https://example.test/screenshots.html#/splatNetGear/${getTopOfCurrentHour()}`,
]);
expect(renders.every(r => r.screenshotOptions.type === 'png')).toBe(true);
expect(await env.ASSETS.get('twitter-images/schedule.png')).not.toBeNull();
expect(await env.ASSETS.get('twitter-images/gear.png')).not.toBeNull();
expect((await env.ASSETS.head('twitter-images/schedule.png')).httpMetadata.contentType).toBe('image/png');
expect(await env.PRIVATE.get('bluesky-lastPostTimes.json')).toBeNull(); // nothing was posted
});
});

View File

@@ -13,6 +13,7 @@
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';
@@ -20,11 +21,13 @@ import { currentColo } from './colo.mjs';
const RETRY_DELAY_MS = 60 * 1000;
const MAX_RETRIES = 3;
export const HOURLY_JOB = 'updaters';
// 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. The hourly job is the full updater run; more can be added here.
// Jobs the object can run.
const JOBS = {
updaters: env => runUpdaters(env),
posters: env => runPosters(env),
};
const EMPTY_STATE = {
@@ -87,17 +90,20 @@ export class Scheduler extends DurableObject {
let colo = await currentColo();
let state = await this.#state();
// Work out what is due. The hourly job is due when its time (or its retry time) has come;
// pending jobs are due now. A pending request for the hourly job merges into the hourly run.
// 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;
if (hourlyScheduledFor !== null && firedAt >= hourlyScheduledFor)
due.push({ job: HOURLY_JOB, reason: state.retryAt ? 'retry' : 'hourly', scheduledFor: hourlyScheduledFor });
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 = {
@@ -123,21 +129,25 @@ export class Scheduler extends DurableObject {
run = { ...run, ok: false, runMs: Date.now() - startedAt, ...describeError(error) };
}
if (job === HOURLY_JOB && reason !== 'wake') {
if (run.ok || 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;
}
}
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;
}
}
// 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.

View File

@@ -1,5 +1,6 @@
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';
@@ -58,6 +59,7 @@ export default withSentry(sentryOptions, {
// 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
@@ -65,7 +67,7 @@ export default withSentry(sentryOptions, {
async fetch(request, env, ctx) {
let url = new URL(request.url);
let route = `${request.method} ${url.pathname}`;
if (!['POST /run', 'POST /wake', 'POST /arm', 'GET /status', 'GET /list'].includes(route))
if (!['POST /run', 'POST /post', 'POST /wake', 'POST /arm', 'GET /status', 'GET /list'].includes(route))
return new Response('Not found', { status: 404 });
if (!isAuthorized(request, env))
return new Response('Unauthorized', { status: 401 });
@@ -76,6 +78,8 @@ export default withSentry(sentryOptions, {
let only = url.searchParams.get('only')?.split(',').map(name => name.trim()).filter(Boolean);
return Response.json(await runUpdaters(env, { only }));
}
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 });

View File

@@ -0,0 +1,31 @@
// Runs the social posters (src/app/twitter) against this Worker's R2 bindings, after the
// updaters have published the hour's data. The Worker's counterpart of postLocally().
//
// The posters read their credentials and the screenshot configuration from process.env,
// which Workers populate from the bindings: BLUESKY_*, TWITTER_*, SITE_URL,
// CLOUDFLARE_ACCOUNT_ID, CLOUDFLARE_BROWSER_RUN_API_TOKEN. With no social credentials set
// they only render and save the public images (shadow mode).
import { maybePostTweets, createClients } from '../../../src/app/twitter/index.js';
import { bucketStorage } from './updaters.mjs';
import { createLogger } from './log.mjs';
/**
* @param {{ ASSETS: R2Bucket, PRIVATE: R2Bucket }} env
* @returns {Promise<{ ok: boolean, ms: number, clients: string[] }>}
*/
export async function runPosters(env) {
let log = createLogger('posters');
let started = Date.now();
let clients = createClients();
let enabled = [];
for (let client of clients)
if (await client.canSend())
enabled.push(client.key);
await maybePostTweets(bucketStorage(env), clients);
let summary = { ok: true, ms: Date.now() - started, clients: enabled };
log.info('Posters finished', summary);
return summary;
}

View File

@@ -25,6 +25,11 @@
"migrations": [
{ "tag": "v1", "new_sqlite_classes": ["Scheduler"] }
],
"vars": {
// Screenshots for social posts are rendered from the deployed site by Browser Rendering.
// The API token is a secret (CLOUDFLARE_BROWSER_RUN_API_TOKEN).
"SITE_URL": "https://splatoon2.ink"
},
"r2_buckets": [
{
// Shadow mode: writes go to the dev bucket while the container still owns production.