diff --git a/packages/cli/src/bench/scenarios/inbox.test.ts b/packages/cli/src/bench/scenarios/inbox.test.ts index 57df0e970..ac10a8912 100644 --- a/packages/cli/src/bench/scenarios/inbox.test.ts +++ b/packages/cli/src/bench/scenarios/inbox.test.ts @@ -9,6 +9,7 @@ import test from "node:test"; import { serve } from "srvx"; import { getContextLoader, getDocumentLoader } from "../../docloader.ts"; import { buildFleet } from "../actor/fleet.ts"; +import type { Clock } from "../load/clock.ts"; import { normalizeSuite } from "../scenario/normalize.ts"; import type { Suite } from "../scenario/types.ts"; import { spawnSyntheticServer } from "../server/synthetic.ts"; @@ -60,16 +61,17 @@ async function spawnBenchmarkTarget(usernames: string[] = ["alice"]) { received++; }); - // Record every inbox path that was POSTed to, so a test can confirm that - // deliveries were spread across multiple recipients' personal inboxes. - const inboxHits = new Set(); + // Count POSTs per inbox path, so a test can confirm how deliveries were + // allocated across multiple recipients' personal inboxes. + const inboxHits = new Map(); const server = serve({ port: 0, hostname: "127.0.0.1", silent: true, fetch: (request: Request) => { if (request.method === "POST") { - inboxHits.add(new URL(request.url).pathname); + const path = new URL(request.url).pathname; + inboxHits.set(path, (inboxHits.get(path) ?? 0) + 1); } return federation.fetch(request, { contextData: undefined }); }, @@ -192,7 +194,13 @@ test("inboxRunner - reports server metrics scoped past the warm-up", async () => } }); -test("inboxRunner - rotates deliveries across multiple recipients", async () => { +// Recipients are allocated round-robin by signing index, so the run makes a +// fixed number of attempts rather than however many fit in a real-time window: +// four constant arrivals under a fake clock, each consuming one allocation +// index, split exactly two per recipient whatever order they complete in. +async function assertAllocatesAcrossRecipients( + signing: "jit" | "presign", +): Promise { const target = await spawnBenchmarkTarget(["alice", "bob"]); let fleet: Awaited> | undefined; try { @@ -214,11 +222,22 @@ test("inboxRunner - rotates deliveries across multiple recipients", async () => ], // Personal inboxes so each recipient's deliveries hit a distinct path. inbox: "personal", - load: { concurrency: 2 }, - duration: "300ms", + // 10/s over 400ms schedules exactly four arrivals (0, 100, 200, and + // 300ms), which is also the batch size presign signs up front. + load: { rate: 10, maxInFlight: 2 }, + duration: "400ms", + signing, }], }; const scenario = normalizeSuite(suite).scenarios[0]; + let now = 0; + const clock: Clock = { + now: () => now, + sleepUntil: (timeMs) => { + now = Math.max(now, timeMs); + return Promise.resolve(); + }, + }; const measurement = await inboxRunner.run({ scenario, target: target.url, @@ -226,8 +245,10 @@ test("inboxRunner - rotates deliveries across multiple recipients", async () => contextLoader: await getContextLoader({ allowPrivateAddress: true }), allowPrivateAddress: true, fleet, + clock, }); + assert.strictEqual(measurement.requests.total, 4); assert.strictEqual( measurement.requests.successRate, 1, @@ -235,15 +256,10 @@ test("inboxRunner - rotates deliveries across multiple recipients", async () => JSON.stringify(measurement.errors) }`, ); - // Both recipients' personal inboxes received deliveries. - const hits = target.inboxHits(); - assert.ok( - hits.has("/users/alice/inbox"), - `expected alice's inbox to be hit; hits: ${JSON.stringify([...hits])}`, - ); - assert.ok( - hits.has("/users/bob/inbox"), - `expected bob's inbox to be hit; hits: ${JSON.stringify([...hits])}`, + // Each recipient's personal inbox received exactly half of the deliveries. + assert.deepStrictEqual( + Object.fromEntries(target.inboxHits()), + { "/users/alice/inbox": 2, "/users/bob/inbox": 2 }, ); } finally { try { @@ -252,7 +268,19 @@ test("inboxRunner - rotates deliveries across multiple recipients", async () => await target.close(); } } -}); +} + +test( + "inboxRunner - allocates jit deliveries across multiple recipients", + { timeout: 60_000 }, + () => assertAllocatesAcrossRecipients("jit"), +); + +test( + "inboxRunner - allocates presigned deliveries across multiple recipients", + { timeout: 60_000 }, + () => assertAllocatesAcrossRecipients("presign"), +); test("inboxRunner.validate - rejects activity options it cannot honor", () => { function resolve(activity: Record) {