jevstrudel.git / worker / src / relay.test.ts
1import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
2import { memoryBudgets } from '../test/budget';
3import { testD1 } from '../test/d1';
4import { memoryKV } from '../test/kv';
5import { d1Accounts } from './accounts-store';
6import { CACHE_HEADER, cacheKey } from './answer-cache';
7import { dataPoint } from './events';
8import { ATTEMPT_HEADER, BUDGET_HEADER, LEAD_HEADER, pageTiming, relay, whyNot, type RelayLog } from './relay';
9import { newSessionToken, SESSION_TTL_MS } from './session';
10import {
11  ATTEMPT_HEADER as PAGE_ATTEMPT_HEADER,
12  BUDGET_HEADER as PAGE_BUDGET_HEADER,
13  CACHE_HEADER as PAGE_CACHE_HEADER,
14  LEAD_HEADER as PAGE_LEAD_HEADER,
15} from '../../website/src/jev/ask.mjs';
16
17const choice = (options: string[]) => ({
18  type: 'choice',
19  instructions: 'pick one',
20  criteria: Object.fromEntries(options.map((o) => [o, `means ${o}`])),
21});
22const ok = { model: 'jev-1.13.0', questions: { move: choice(['a', 'b']) }, state: { song: 'x' } };
23
24const score = { type: 'score', instructions: 'how intense', criteria: ['calm', 'driving', 'all out'] };
25const noul = { type: 'noul', instructions: 'is it art?', criteria: { true: 'art', false: 'not art' } };
26
27describe('whyNot: the relay forwards only what the page sends', () => {
28  it('accepts the booth, the critic and the mood picker', () => {
29    expect(whyNot(ok)).toBeNull();
30    // the form: section and move, with the section's intensity
31    expect(
32      whyNot({ ...ok, questions: { move: choice(['a', 'b']), section: choice(['c', 'd', 'e']), intensity: score } }),
33    ).toBeNull();
34    expect(whyNot({ ...ok, questions: { art: noul } })).toBeNull();
35    // the art critic's rubric: six scores in one call
36    expect(whyNot({ ...ok, questions: Object.fromEntries([...'abcdef'].map((k) => [k, score])) })).toBeNull();
37    expect(whyNot({ ...ok, questions: { art: { type: 'noul', instructions: 'x' } } })).toBeNull();
38    expect(whyNot({ ...ok, questions: { song: choice([...'abcdefghijklmnopqrstuvwxyz']) } })).toBeNull();
39  });
40
41  it('refuses anything else before it reaches TypeSafe', () => {
42    expect(whyNot(null)).toMatch(/object/);
43    expect(whyNot({ ...ok, model: 'jev-2' })).toMatch(/model/);
44    expect(whyNot({ ...ok, extra: 1 })).toMatch(/unexpected/);
45    expect(whyNot({ ...ok, state: 'x' })).toMatch(/state/);
46    expect(whyNot({ ...ok, questions: {} })).toMatch(/at least one/);
47    const twelve = Object.fromEntries([...'abcdefghijkl'].map((k) => [k, choice(['a', 'b'])]));
48    expect(whyNot({ ...ok, questions: twelve })).toBeNull();
49    expect(whyNot({ ...ok, questions: { q: { type: 'extract', instructions: 'x' } } })).toMatch(
50      /choice, score or noul/,
51    );
52    expect(whyNot({ ...ok, questions: { q: choice(['a']) } })).toMatch(/2 to 255/);
53    expect(whyNot({ ...ok, questions: { q: choice([...Array(255).keys()].map(String)) } })).toBeNull();
54    expect(whyNot({ ...ok, questions: { q: choice([...Array(256).keys()].map(String)) } })).toMatch(/2 to 255/);
55    expect(whyNot({ ...ok, questions: { q: { ...choice(['a', 'b']), criteria: { a: { x: 1 }, b: 'y' } } } })).toMatch(
56      /strings/,
57    );
58    expect(whyNot({ ...ok, questions: { q: { ...score, criteria: ['one'] } } })).toMatch(/2 to 10/);
59    expect(whyNot({ ...ok, questions: { q: { ...score, criteria: [...'abcdefghijk'] } } })).toMatch(/2 to 10/);
60    expect(whyNot({ ...ok, questions: { q: { ...noul, criteria: { true: 'x', false: 'y', maybe: 'z' } } } })).toMatch(
61      /true, false/,
62    );
63    expect(whyNot({ ...ok, questions: { q: { ...choice(['a', 'b']), instructions: 7 } } })).toMatch(/instructions/);
64  });
65});
66
67describe('relay: one structured log line per call, with no visitor in it', () => {
68  const IP = '203.0.113.7';
69  let points: AnalyticsEngineDataPoint[];
70  const env = (success = true) => ({
71    JEVSTRUDEL_TYPESAFE_API_KEY: 'test-key',
72    JEV_LIMIT: { limit: async () => ({ success }) },
73    EVENTS: { writeDataPoint: (p: AnalyticsEngineDataPoint) => void points.push(p) },
74    CACHE: memoryKV(),
75  });
76  const ctx = { waitUntil: () => {} };
77  const call = (body: string, headers: Record<string, string> = {}) =>
78    new Request('https://jevstrudel.example/jev/v1/systemone', {
79      method: 'POST',
80      headers: { 'Content-Type': 'application/json', 'CF-Connecting-IP': IP, ...headers },
81      body,
82    });
83  let lines: unknown[];
84  beforeEach(() => {
85    lines = [];
86    points = [];
87    vi.spyOn(console, 'log').mockImplementation((line: unknown) => void lines.push(line));
88  });
89  afterEach(() => {
90    vi.restoreAllMocks();
91    vi.unstubAllGlobals();
92  });
93  const logged = () => {
94    expect(lines).toHaveLength(1);
95    // the same call as a jev.relay event, and nothing of the visitor in it either
96    expect(points).toHaveLength(1);
97    const line = lines[0] as RelayLog;
98    expect(points[0]).toEqual(
99      dataPoint('jev.relay', {
100        blobs: { outcome: line.outcome, colo: line.colo },
101        doubles: {
102          status: line.status,
103          upstreamMs: line.upstreamMs,
104          requestBytes: line.requestBytes,
105          sectionInS: line.sectionInS,
106          attempt: line.attempt,
107        },
108      }),
109    );
110    expect(JSON.stringify(points[0])).not.toContain(IP);
111    const text = JSON.stringify(lines[0]);
112    expect(text).not.toContain(IP);
113    expect(text).not.toContain('pick one');
114    return lines[0] as Record<string, unknown>;
115  };
116
117  it('logs status, upstream latency, body size and the page timing, and forwards none of the page headers', async () => {
118    const upstream = vi.fn(async (_url: string, _init: RequestInit) => Response.json({ answers: {} }));
119    vi.stubGlobal('fetch', upstream);
120    const body = JSON.stringify(ok);
121    const res = await relay(call(body, { [LEAD_HEADER]: '5.8349', [ATTEMPT_HEADER]: '1' }), env() as never, ctx);
122    expect(res.status).toBe(200);
123    const [, init] = upstream.mock.calls[0];
124    expect(new TextDecoder().decode(init.body as ArrayBuffer)).toBe(body);
125    expect(Object.keys(init.headers as Record<string, string>).sort()).toEqual(['Authorization', 'Content-Type']);
126    expect(logged()).toEqual({
127      event: 'jev.relay',
128      outcome: 'forwarded',
129      status: 200,
130      upstreamMs: expect.any(Number),
131      requestBytes: new TextEncoder().encode(body).byteLength,
132      sectionInS: 5.83,
133      attempt: 1,
134      colo: null,
135    });
136  });
137
138  it("passes TypeSafe's status through to the log", async () => {
139    vi.stubGlobal('fetch', async () => new Response('slow down', { status: 429 }));
140    const res = await relay(call(JSON.stringify(ok)), env() as never, ctx);
141    expect(res.status).toBe(429);
142    expect(logged()).toMatchObject({
143      outcome: 'forwarded',
144      status: 429,
145      upstreamMs: expect.any(Number),
146      sectionInS: null,
147      attempt: null,
148    });
149  });
150
151  it('logs a call TypeSafe could not answer as 502', async () => {
152    vi.stubGlobal('fetch', async () => {
153      throw new TypeError('network');
154    });
155    const res = await relay(call(JSON.stringify(ok)), env() as never, ctx);
156    expect(res.status).toBe(502);
157    expect(logged()).toMatchObject({ outcome: 'unreachable', status: 502, upstreamMs: expect.any(Number) });
158  });
159
160  it('logs refusals without asking TypeSafe, and without what was asked', async () => {
161    const upstream = vi.fn();
162    vi.stubGlobal('fetch', upstream);
163    await relay(call(JSON.stringify(ok)), env(false) as never, ctx);
164    expect(logged()).toMatchObject({ outcome: 'refused', status: 429, upstreamMs: null, requestBytes: null });
165    lines = [];
166    points = [];
167    const bad = JSON.stringify({ ...ok, model: 'pick one' });
168    await relay(call(bad, { [LEAD_HEADER]: '-1.5' }), env() as never, ctx);
169    expect(logged()).toMatchObject({
170      outcome: 'refused',
171      status: 400,
172      upstreamMs: null,
173      requestBytes: bad.length,
174      sectionInS: -1.5,
175    });
176    expect(upstream).not.toHaveBeenCalled();
177  });
178
179  it('keeps only bounded numbers from the page headers', () => {
180    const t = (lead: string, attempt: string) =>
181      pageTiming(new Headers({ [LEAD_HEADER]: lead, [ATTEMPT_HEADER]: attempt }));
182    expect(t('2.345', '0')).toEqual({ sectionInS: 2.35, attempt: 0 });
183    expect(t('203.0.113.7', 'x')).toEqual({ sectionInS: null, attempt: null });
184    expect(t('1e9', '1.5')).toEqual({ sectionInS: null, attempt: null });
185    expect(t('Infinity', '99')).toEqual({ sectionInS: null, attempt: null });
186    expect(pageTiming(new Headers())).toEqual({ sectionInS: null, attempt: null });
187  });
188
189  it('reads the headers the page sends', () => {
190    expect(LEAD_HEADER).toBe(PAGE_LEAD_HEADER);
191    expect(ATTEMPT_HEADER).toBe(PAGE_ATTEMPT_HEADER);
192  });
193});
194
195describe('relay: the answer cache', () => {
196  const answer = { answers: { move: 'a' } };
197  let upstream: ReturnType<typeof vi.fn>;
198  let pending: Promise<unknown>[];
199  let lines: RelayLog[];
200  const ctx = { waitUntil: (p: Promise<unknown>) => void pending.push(p) };
201  const env = (cache = memoryKV()) => ({
202    JEVSTRUDEL_TYPESAFE_API_KEY: 'test-key',
203    JEV_LIMIT: { limit: async () => ({ success: true }) },
204    EVENTS: { writeDataPoint: () => {} },
205    CACHE: cache,
206  });
207  const call = (body: string) => new Request('https://jevstrudel.example/jev/v1/systemone', { method: 'POST', body });
208  // a call, and then whatever it left to finish after answering
209  const ask = async (e: ReturnType<typeof env>, body: string) => {
210    const res = await relay(call(body), e as never, ctx);
211    const text = await res.text();
212    await Promise.all(pending.splice(0));
213    return { res, text };
214  };
215  const bytes = (s: string) => new TextEncoder().encode(s).buffer;
216  beforeEach(() => {
217    pending = [];
218    lines = [];
219    upstream = vi.fn(async () => Response.json(answer));
220    vi.stubGlobal('fetch', upstream);
221    vi.spyOn(console, 'log').mockImplementation((line: RelayLog) => void lines.push(line));
222    vi.spyOn(console, 'error').mockImplementation(() => {});
223  });
224  afterEach(() => {
225    vi.restoreAllMocks();
226    vi.unstubAllGlobals();
227  });
228
229  it('answers the same request again from KV, for an hour, without asking TypeSafe', async () => {
230    const e = env();
231    const body = JSON.stringify(ok);
232    const first = await ask(e, body);
233    expect(first.res.headers.get(CACHE_HEADER)).toBe('miss');
234    expect(e.CACHE.puts).toEqual([
235      {
236        key: await cacheKey('jev-1.13.0', bytes(body)),
237        expirationTtl: 3600,
238        metadata: { contentType: 'application/json' },
239      },
240    ]);
241    const second = await ask(e, body);
242    expect(second.res.status).toBe(200);
243    expect(second.res.headers.get(CACHE_HEADER)).toBe('hit');
244    expect(second.res.headers.get('Content-Type')).toBe('application/json');
245    expect(second.text).toBe(first.text);
246    expect(upstream).toHaveBeenCalledTimes(1);
247    expect(lines.map((l) => l.outcome)).toEqual(['forwarded', 'cached']);
248    expect(lines[1]).toMatchObject({ status: 200, upstreamMs: null, requestBytes: body.length });
249  });
250
251  it('never answers a request that differs in any byte from the one it kept', async () => {
252    const e = env();
253    const body = JSON.stringify(ok);
254    await ask(e, body);
255    // the same JSON as other bytes, and one character of state apart
256    const variants = [
257      JSON.stringify(ok, null, 1),
258      JSON.stringify({ state: ok.state, questions: ok.questions, model: ok.model }),
259      JSON.stringify({ ...ok, state: { song: 'y' } }),
260      `${body} `,
261    ];
262    for (const v of variants) expect((await ask(e, v)).res.headers.get(CACHE_HEADER)).toBe('miss');
263    expect(upstream).toHaveBeenCalledTimes(1 + variants.length);
264    expect(new Set(e.CACHE.puts.map((p) => p.key)).size).toBe(1 + variants.length);
265  });
266
267  it('keeps only a 200', async () => {
268    const e = env();
269    for (const status of [400, 401, 429, 500, 529]) {
270      upstream.mockImplementationOnce(async () => Response.json({ error: 'no' }, { status }));
271      const { res } = await ask(e, JSON.stringify(ok));
272      expect(res.status).toBe(status);
273      expect(res.headers.get(CACHE_HEADER)).toBe('miss');
274    }
275    expect(e.CACHE.puts).toEqual([]);
276  });
277
278  it('does not keep an answer larger than it bounds, and the page still gets all of it', async () => {
279    const e = env();
280    upstream.mockImplementationOnce(async () => new Response('x'.repeat(64 * 1024 + 1), { status: 200 }));
281    const { res, text } = await ask(e, JSON.stringify(ok));
282    expect(res.status).toBe(200);
283    expect(text).toHaveLength(64 * 1024 + 1);
284    expect(e.CACHE.puts).toEqual([]);
285  });
286
287  it('forwards as if there were no cache when KV fails', async () => {
288    const fail = async () => {
289      throw new Error('kv down');
290    };
291    const e = env(Object.assign(memoryKV(), { getWithMetadata: fail, put: fail }));
292    const { res } = await ask(e, JSON.stringify(ok));
293    expect(res.status).toBe(200);
294    expect(res.headers.get(CACHE_HEADER)).toBe('miss');
295    expect(upstream).toHaveBeenCalledTimes(1);
296  });
297
298  it('keys by the pinned model and the body, and the page reads its header', async () => {
299    const body = bytes(JSON.stringify(ok));
300    expect(await cacheKey('jev-1.13.0', body)).toMatch(/^jev\.relay\/v1\/jev-1\.13\.0\/[0-9a-f]{64}$/);
301    expect(await cacheKey('jev-1.14.0', body)).not.toBe(await cacheKey('jev-1.13.0', body));
302    expect(CACHE_HEADER).toBe(PAGE_CACHE_HEADER);
303  });
304});
305
306describe('relay: a signed-in caller’s daily budget', () => {
307  const noon = Date.UTC(2026, 8, 25, 12);
308  const USER = 'u'.repeat(22);
309  let db: D1Database;
310  let points: AnalyticsEngineDataPoint[];
311  let keys: string[];
312  let pending: Promise<unknown>[];
313  let upstream: ReturnType<typeof vi.fn>;
314  let token: string;
315  const ctx = { waitUntil: (p: Promise<unknown>) => void pending.push(p) };
316  const env = (limit = 3) => ({
317    JEVSTRUDEL_TYPESAFE_API_KEY: 'test-key',
318    JEV_LIMIT: { limit: async ({ key }: { key: string }) => (keys.push(key), { success: true }) },
319    EVENTS: { writeDataPoint: (p: AnalyticsEngineDataPoint) => void points.push(p) },
320    CACHE: memoryKV(),
321    DB: db,
322    BUDGET: memoryBudgets(limit, () => Date.now()),
323  });
324  const call = (body: string, signedIn = true) =>
325    new Request('https://jevstrudel.example/jev/v1/systemone', {
326      method: 'POST',
327      headers: { 'CF-Connecting-IP': '203.0.113.7', ...(signedIn ? { Cookie: `jev_session=${token}` } : {}) },
328      body,
329    });
330  // a different body each time, so the answer cache never answers
331  const fresh = (i: number) => JSON.stringify({ ...ok, state: { song: `x${i}` } });
332  const used = (e: ReturnType<typeof env>) => e.BUDGET.get(e.BUDGET.idFromName(USER)).balance().then((b) => b.used);
333
334  beforeEach(async () => {
335    vi.useFakeTimers({ now: noon, toFake: ['Date'] });
336    db = testD1();
337    points = [];
338    keys = [];
339    pending = [];
340    upstream = vi.fn(async () => Response.json({ answers: {} }));
341    vi.stubGlobal('fetch', upstream);
342    vi.spyOn(console, 'log').mockImplementation(() => {});
343    vi.spyOn(console, 'error').mockImplementation(() => {});
344    token = newSessionToken();
345    await d1Accounts(db).createUser(
346      { id: USER, displayName: 'Link' },
347      { id: 'c'.repeat(16), publicKey: new Uint8Array([1]), signCount: 0, transports: [] },
348      token,
349      SESSION_TTL_MS,
350    );
351  });
352  afterEach(() => {
353    vi.useRealTimers();
354    vi.restoreAllMocks();
355    vi.unstubAllGlobals();
356  });
357
358  it('charges each forwarded call, says how much is used, and refuses once spent until UTC midnight', async () => {
359    const e = env(3);
360    for (let i = 1; i <= 3; i++) {
361      const res = await relay(call(fresh(i)), e as never, ctx);
362      expect(res.status).toBe(200);
363      expect(res.headers.get(BUDGET_HEADER)).toBe(`${i}/3`);
364    }
365    const spent = await relay(call(fresh(4)), e as never, ctx);
366    expect(spent.status).toBe(429);
367    expect(spent.headers.get(BUDGET_HEADER)).toBe('spent');
368    expect(spent.headers.get('Retry-After')).toBe(String(12 * 3600));
369    expect(await spent.text()).toMatch(/daily Jev budget \(3 calls\) is spent; it resets at 00:00 UTC/);
370    expect(upstream).toHaveBeenCalledTimes(3);
371    expect(BUDGET_HEADER).toBe(PAGE_BUDGET_HEADER);
372  });
373
374  it('writes one event, with no one in it, on the call that spends the last of it', async () => {
375    const e = env(2);
376    for (let i = 0; i < 4; i++) await relay(call(fresh(i)), e as never, ctx);
377    const exhausted = points.filter((p) => p.indexes?.[0] === 'jev.budget.exhausted');
378    expect(exhausted).toEqual([
379      dataPoint('jev.budget.exhausted', { blobs: { colo: null }, doubles: { limit: 2, secondsToReset: 12 * 3600 } }),
380    ]);
381    const text = JSON.stringify(points);
382    expect(text).not.toContain(USER);
383    expect(text).not.toContain('203.0.113.7');
384  });
385
386  it('does not charge an answer from the cache, nor a request it refuses', async () => {
387    const e = env(5);
388    const body = fresh(1);
389    await (await relay(call(body), e as never, ctx)).text();
390    await Promise.all(pending.splice(0));
391    const again = await relay(call(body), e as never, ctx);
392    expect(again.headers.get(CACHE_HEADER)).toBe('hit');
393    expect((await relay(call('not json'), e as never, ctx)).status).toBe(400);
394    expect(await used(e)).toBe(1);
395  });
396
397  it('leaves anonymous callers as they were: per address, never charged', async () => {
398    const e = env(1);
399    for (let i = 0; i < 3; i++) {
400      const res = await relay(call(fresh(i), false), e as never, ctx);
401      expect(res.status).toBe(200);
402      expect(res.headers.get(BUDGET_HEADER)).toBeNull();
403    }
404    expect(keys).toEqual(['203.0.113.7', '203.0.113.7', '203.0.113.7']);
405    expect(await used(e)).toBe(0);
406    await relay(call(fresh(9)), e as never, ctx);
407    expect(keys.at(-1)).toBe(`user:${USER}`);
408  });
409
410  it('never takes a session from another origin: the sandbox a listener’s song plays in is anonymous', async () => {
411    const e = env(5);
412    for (const origin of ['null', 'https://evil.example']) {
413      const req = new Request('https://jevstrudel.example/jev/v1/systemone', {
414        method: 'POST',
415        headers: { 'CF-Connecting-IP': '203.0.113.7', Origin: origin, Cookie: `jev_session=${token}` },
416        body: fresh(origin.length),
417      });
418      const res = await relay(req, e as never, ctx);
419      expect(res.headers.get(BUDGET_HEADER)).toBeNull();
420      expect(keys.at(-1)).toBe('203.0.113.7');
421    }
422    expect(await used(e)).toBe(0);
423    // this site's own page, with the same cookie, is charged
424    const own = new Request('https://jevstrudel.example/jev/v1/systemone', {
425      method: 'POST',
426      headers: { 'CF-Connecting-IP': '203.0.113.7', Origin: 'https://jevstrudel.example', Cookie: `jev_session=${token}` },
427      body: fresh(99),
428    });
429    expect((await relay(own, e as never, ctx)).headers.get(BUDGET_HEADER)).toBe('1/5');
430  });
431
432  it('refuses, unforwarded, when the budget cannot be counted', async () => {
433    const broken = {
434      idFromName: () => ({}),
435      get: () => ({
436        spend: async () => {
437          throw new Error('down');
438        },
439      }),
440    };
441    const res = await relay(call(fresh(1)), { ...env(), BUDGET: broken } as never, ctx);
442    expect(res.status).toBe(503);
443    expect(upstream).not.toHaveBeenCalled();
444  });
445});