jevstrudel.git / worker / src / content-jobs.ts

The work Jev owes listeners' content: screening every item before it is public (screen.ts), and scoring every public song revision with the art critic (art-rubric.mjs). Each is a row in content_jobs, made by the schema's triggers when an item is written and deleted by them when its verdict or score is (migrations/0004_listener_content.sql), so the owed work is the database's state, not a message that could be lost.

A job runs when claimed (content-store.ts's claimDue: a lease, so no two workers run it at once, and a dead worker's lease expires). Three things claim jobs:

  • the write itself: publishing, commenting and pitching screen the new item at once, so the author sees the verdict in the answer, and a song's score follows in waitUntil (content.ts);
  • the author's own view of their content (GET /jev/listeners/mine) runs their due jobs in waitUntil, so their pending items move while they watch;
  • the cron trigger (wrangler.json, triggers.crons, every minute) runs everyone's due jobs, SWEEP_BATCH at a time: the guarantee, since it needs nobody to come back.

A job that cannot run now is rescheduled with why, and its author sees both: the budget spent waits for the budget's reset at 00:00 UTC; TypeSafe unreachable, answering an error, or answering without a verdict is asked again after a backoff from a minute, doubling, to an hour, for as long as it takes. Nothing is ever dropped as "failed": pending until done.

Every attempt that reaches TypeSafe is charged to the author's daily Jev budget (budget.ts), a screening one call and a score SCORE_RUNS: charged before asking, as the relay charges, so a call that then fails is still spent. A budget that cannot be counted is not charged and not asked.

31import { artOf, artRequest, CRITIC_MODEL, meanVerdict, RUBRIC } from './art-rubric.mjs';
32import type { ContentStore, Job, Why } from './content-store';
33import type { Env } from './env';
34import { event } from './events';
35import { upstream, whyNot } from './relay';
36import { screenRequest, verdictOf } from './screen';

A job's lease: longer than three TypeSafe calls in parallel take, so a live worker keeps it; short enough that a dead one's job waits little.

40export const LEASE_MS = 2 * 60 * 1000;

The art is the mean of three runs, as tools/critic records the site's own songs: one call moves by up to 0.06 on unchanged code (critic.mjs), which would reorder a list sorted by art, and listeners' songs are listed and compared beside the site's. Three calls of at most ~24k input tokens (a 64 KiB song and a 32 KiB spec) each, charged to the publisher.

46export const SCORE_RUNS = 3;

Jobs a cron run takes. On the free plan an invocation may make 50 subrequests, D1 and Durable Object calls included: claiming is one, and a job at most six (its content, the budget, three TypeSafe calls, its result), so six jobs are 37. Due jobs past that wait a minute for the next run.

51export const SWEEP_BATCH = 6;
52const BACKOFF_FIRST_MS = 60 * 1000;
53const BACKOFF_MAX_MS = 60 * 60 * 1000;
55export const backoff = (attempts: number) => Math.min(BACKOFF_FIRST_MS * 2 ** attempts, BACKOFF_MAX_MS);
56
57type JobEnv = Pick<Env, 'JEVSTRUDEL_TYPESAFE_API_KEY' | 'BUDGET' | 'EVENTS'>;
58export type Outcome = 'done' | 'pending' | 'gone';
59
60type Asked = { answers: unknown[]; status: number; upstreamMs: number; bytes: number } | { why: Why; status: number | null; upstreamMs: number | null; retryMs: number | null };

The retry TypeSafe asked for, in ms, if any (as ask.mjs reads it).

63function retryAfterMs(res: Response): number | null {
64  const ms = Number(res.headers.get('retry-after-ms'));
65  if (res.headers.has('retry-after-ms') && Number.isFinite(ms) && ms >= 0) return ms;
66  const s = Number(res.headers.get('Retry-After'));
67  return res.headers.has('Retry-After') && Number.isFinite(s) && s >= 0 ? s * 1000 : null;
68}

Asks TypeSafe runs times the same body (each run is its own answer: nothing here goes through the relay's answer cache, which would hand the later runs the first one's).

73async function ask(env: JobEnv, body: string, runs: number): Promise<Asked> {
74  const asked = Date.now();
75  let responses: Response[];
76  try {
77    responses = await Promise.all(Array.from({ length: runs }, () => upstream(env, body)));
78  } catch {
79    return { why: 'unreachable', status: null, upstreamMs: Date.now() - asked, retryMs: null };
80  }
81  const upstreamMs = Date.now() - asked;
82  const failed = responses.find((r) => !r.ok);
83  if (failed) {
84    await Promise.all(responses.map((r) => r.body?.cancel()));
85    return { why: 'upstream', status: failed.status, upstreamMs, retryMs: retryAfterMs(failed) };
86  }
87  try {
88    const answers = await Promise.all(responses.map(async (r) => ((await r.json()) as { answers?: unknown }).answers));
89    return { answers, status: 200, upstreamMs, bytes: body.length };
90  } catch {
91    return { why: 'unanswered', status: 200, upstreamMs, retryMs: null };
92  }
93}

Runs one claimed job to its end or its next attempt.

96export async function runJob(env: JobEnv, store: ContentStore, job: Job): Promise<Outcome> {
97  const { now } = store;
98  const screening = job.task === 'screen';
99  const note = (fields: {
100    outcome: string;
101    why: Why | null;
102    status?: number | null;
103    upstreamMs?: number | null;
104    confidence?: number | null;
105    art?: number | null;
106    bytes?: number | null;
107  }) => {
108    const line = { event: screening ? 'jev.screen' : 'jev.score', kind: job.kind, attempt: job.attempts, ...fields };
109    console.log(line);
110    if (screening) {
111      event(env, 'jev.screen', {
112        blobs: { kind: job.kind, outcome: fields.outcome, why: fields.why },
113        doubles: {
114          status: fields.status ?? null,
115          upstreamMs: fields.upstreamMs ?? null,
116          attempt: job.attempts,
117          confidence: fields.confidence ?? null,
118          requestBytes: fields.bytes ?? null,
119        },
120      });
121    } else {
122      event(env, 'jev.score', {
123        blobs: { outcome: fields.outcome, why: fields.why },
124        doubles: {
125          status: fields.status ?? null,
126          upstreamMs: fields.upstreamMs ?? null,
127          attempt: job.attempts,
128          art: fields.art ?? null,
129          runs: SCORE_RUNS,
130        },
131      });
132    }
133  };
134  const later = async (why: Why, at: number, fields: { status?: number | null; upstreamMs?: number | null } = {}) => {
135    await store.reschedule(job, at, why);
136    note({ outcome: 'pending', why, ...fields });
137    return 'pending' as const;
138  };

what to ask

141  let body: object;
142  const runs = screening ? 1 : SCORE_RUNS;
143  if (screening) {
144    const subject = await store.subject(job.kind, job.item);
145    if (!subject) return 'gone'; // deleted since: its trigger took the job too
146    body = screenRequest(subject);
147  } else {
148    const revision = await store.toScore(job.item);
149    if (!revision) return 'gone';
150    body = artRequest(revision);
151  }
152  const why = whyNot(body);
153  if (why) throw new Error(`content-jobs: a ${job.task} request the relay would refuse: ${why}`);
155  if (!env.JEVSTRUDEL_TYPESAFE_API_KEY) return later('unreachable', now() + backoff(job.attempts));

charged first, as the relay charges: only what is about to reach TypeSafe

158  let charge;
159  try {
160    charge = await env.BUDGET.get(env.BUDGET.idFromName(job.userId)).spend(runs);
161  } catch (e) {
162    console.error({ event: 'jev.budget', op: 'spend', error: (e as Error).message });
163    return later('budget-unavailable', now() + backoff(job.attempts));
164  }
165  if (!charge.ok) return later('budget', charge.resetsAt);
167  const res = await ask(env, JSON.stringify(body), runs);
168  if ('why' in res) {
169    const wait = Math.max(backoff(job.attempts), res.retryMs ?? 0);
170    return later(res.why, now() + wait, { status: res.status, upstreamMs: res.upstreamMs });
171  }
172
173  const asked = { status: res.status, upstreamMs: res.upstreamMs };
174  if (screening) {
175    const verdict = verdictOf(res.answers[0]);
176    if (!verdict) return later('unanswered', now() + backoff(job.attempts), asked);
177    await store.recordScreen(job.kind, job.item, verdict.verdict, verdict.confidence);
178    note({ outcome: verdict.verdict, why: null, ...asked, bytes: res.bytes, confidence: verdict.confidence });
179    return 'done';
180  }
181  const verdicts = res.answers.map((a) => artOf(a));
182  if (verdicts.some((v) => v === null)) return later('unanswered', now() + backoff(job.attempts), asked);
183  const mean = meanVerdict(verdicts as NonNullable<(typeof verdicts)[number]>[]);
184  await store.recordScore(job.item, {
185    art: mean.art,
186    artRuns: mean.artRuns,
187    criteria: mean.criteria,
188    model: CRITIC_MODEL,
189    method: RUBRIC,
190  });
191  note({ outcome: 'scored', why: null, ...asked, art: mean.art });
192  return 'done';
193}

Claims up to limit due jobs (everyone's, or one user's) and runs them. A job that throws keeps its lease until it expires and is then due again.

197export async function sweep(
198  env: JobEnv,
199  store: ContentStore,
200  { limit = SWEEP_BATCH, userId = null }: { limit?: number; userId?: string | null } = {},
201): Promise<Outcome[]> {
202  const jobs = await store.claimDue(limit, LEASE_MS, userId);
203  return Promise.all(
204    jobs.map((job) =>
205      runJob(env, store, job).catch((e) => {
206        console.error({ event: 'jev.content.job', kind: job.kind, task: job.task, error: (e as Error).message });
207        return 'pending' as const;
208      }),
209    ),
210  );
211}