diff --git a/README.md b/README.md index 40975d2..811f16d 100644 --- a/README.md +++ b/README.md @@ -17,7 +17,7 @@ _| _| _| _| _| _| _| _| _| _| _| _| [![CI](https://github.com/halaprix/domino/actions/workflows/ci.yml/badge.svg)](https://github.com/halaprix/domino/actions/workflows/ci.yml) [![npm version](https://img.shields.io/npm/v/@halaprix/domino)](https://www.npmjs.com/package/@halaprix/domino) -[![bundle size](https://img.shields.io/badge/gzip-12.2KB-brightgreen)](https://www.npmjs.com/package/@halaprix/domino) +[![bundle size](https://img.shields.io/badge/gzip-13.0KB-brightgreen)](https://www.npmjs.com/package/@halaprix/domino) [![TypeScript](https://img.shields.io/badge/TypeScript-5.5-blue)](https://www.typescriptlang.org/) [![MIT License](https://img.shields.io/badge/License-MIT-yellow.svg)](LICENSE) diff --git a/docs/benchmarks.md b/docs/benchmarks.md index 7a8abd0..aa46258 100644 --- a/docs/benchmarks.md +++ b/docs/benchmarks.md @@ -52,6 +52,56 @@ Benchmark measures network round-trips for various token/vault combinations. All - At batchSize ∞: always 2 round-trips (assuming vaults > 0) - At batchSize 100 (default): each 100 calls = 1 round-trip +## Dedup Hit Rate (F7) + +`{ dedupe: true }` (or `Presets.throughput`, which bundles it with `maxConcurrentBatches`/`adaptiveBatching`) merges calls that share the same `(target, calldata, output shape)` **within one step, across tasks** — before batching/bisection ever sees the wire list. This is the same "same-token portfolio" shape as the round-trip table above, but the benefit compounds differently: two entries holding the *same* token no longer duplicate that token's `symbol()`/`decimals()`/`totalSupply()` calls at all, regardless of how many entries share it. + +**Same-token portfolio example** — a counting `StepExecutor` (counts calls actually received, never real RPCs) resolving 3 view calls per entry, entries spread evenly across a small set of distinct tokens: + +| Entries | Distinct tokens | Naive calls | Wire calls, `dedupe: true` | Hit rate | +|---|---|---|---|---| +| 100 | 10 | 300 | 30 | 90.0% | +| 1,000 | 50 | 3,000 | 150 | 95.0% | + +Hit rate = `1 − (unique wire calls / naive calls)`. It tracks `1 − distinctTokens/entries` for this shape (repetition per token, not batch size) — the more entries share a token, the higher the hit rate, independent of `batchSize`. + +```ts +import { createPublicClient, http } from "viem" +import { mainnet } from "viem/chains" +import { Eip1193Executor, defineTask, runMultistepTasks } from "@halaprix/domino" +import type { Address } from "@halaprix/domino" + +const provider = createPublicClient({ chain: mainnet, transport: http() }) +const executor = new Eip1193Executor(provider) + +const erc20LikeAbi = [ + { type: "function", name: "symbol", stateMutability: "view", inputs: [], outputs: [{ type: "string" }] }, + { type: "function", name: "decimals", stateMutability: "view", inputs: [], outputs: [{ type: "uint8" }] }, + { type: "function", name: "totalSupply", stateMutability: "view", inputs: [], outputs: [{ type: "uint256" }] }, +] as const + +// 100 portfolio entries, only 10 DISTINCT tokens among them — a realistic +// shape: many holders/positions referencing the same small set of tokens. +declare const portfolioTokens: Address[] // length 100, drawn from 10 distinct addresses + +const tasks = portfolioTokens.map((token) => + defineTask((t) => ({ + symbol: t.call({ target: token, abi: erc20LikeAbi, functionName: "symbol" }), + decimals: t.call({ target: token, abi: erc20LikeAbi, functionName: "decimals" }), + totalSupply: t.call({ target: token, abi: erc20LikeAbi, functionName: "totalSupply" }), + })), +) + +// dedupe: true merges identical (target, calldata, output-shape) calls +// within this one step, across all 100 tasks, BEFORE batching/bisection — +// 300 naive calls collapse to 30 unique wire calls (90% hit rate). +await runMultistepTasks(executor, tasks, { dedupe: true }) +``` + +**Dedup eligibility is per-call, not per-run:** only calls built via `t.call` (a `TypedCallSpec`, eligible by default — opt out per call with `dedupe: false`) are ever merged. A hand-authored legacy `StepCall` carries no eligibility stamp and is never merged, `dedupe: true` or not — turning this option on can never change a legacy task's semantics. + +**Conflicting output ABIs never merge:** two subscribers declaring different output shapes for the identical calldata are always kept as separate wire calls (each decodes correctly against its own ABI) — dedup only merges calls that would also decode identically. + ## Live Benchmark — Real RPC Timing The mock benchmark above measures RPC call-count reduction. The live benchmark measures **real wall time** and finds the **practical batchSize ceiling** for your specific RPC endpoint. diff --git a/src/__tests__/dedup.test.ts b/src/__tests__/dedup.test.ts new file mode 100644 index 0000000..3211d19 --- /dev/null +++ b/src/__tests__/dedup.test.ts @@ -0,0 +1,613 @@ +import { describe, it, expect, expectTypeOf } from 'vitest' +import { defineTask } from '../core/defineTask' +import { runMultistepTasks } from '../core/runMultistepTasks' +import { runSettled } from '../core/runSettled' +import { dedupeKeyFor } from '../core/dedupe' +import { DominoCallError } from '../core/errors' +import { Presets } from '../core/presets' +import { encodeAbiParameters, decodeAbiParameters } from '../core/abi' +import type { Address, MultistepTask, StepCall, StepExecutor, RawResult } from '../core/types' + +/** + * F7 — within-step, cross-task call dedup. See `src/core/dedupe.ts` (key + * computation) and `src/core/engine.ts` (grouping + result fan-out, done + * strictly PRE-bisection). The global unhandled-rejection guard + * (`src/__tests__/setup/unhandled-rejections.ts`) fails any test here that + * leaks a rejection. + */ + +const ADDR = '0xA0b86991c6218b36c1d19D4a2e9Eb004C35d5Cc4' as Address +const ADDR_MIXED_CASE = ('0x' + ADDR.slice(2).toUpperCase()) as Address + +/** Minimal single-arg view function — target for the "identical call" tests. */ +const testAbi = [ + { + type: 'function', + name: 'getVal', + stateMutability: 'view', + inputs: [{ name: 'x', type: 'uint256' }], + outputs: [{ type: 'uint256' }], + }, +] as const + +/** + * Records every physical batch it's invoked with (as `StepCall[]` snapshots) + * and resolves every call as success, with a value derived deterministically + * from target+functionName+args — so two subscribers of a MERGED call can be + * asserted to receive the exact same (correct) value. + */ +function makeEchoExecutor(): { executor: StepExecutor; invocations: () => StepCall[][] } { + const invocations: StepCall[][] = [] + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + invocations.push(calls.map((c) => ({ ...c }))) + return calls.map( + (c): RawResult => ({ + status: 'success', + value: `${c.target.toLowerCase()}:${c.functionName}:${(c.args ?? []).map(String).join(',')}`, + }), + ) + }, + } + return { executor, invocations: () => invocations } +} + +// ───────────────────────────────────────────────────────────────────────── +// 1. Identical calls merge across tasks (dedupe on) / don't (dedupe off). +// ───────────────────────────────────────────────────────────────────────── + +describe('F7 dedup — identical calls across tasks', () => { + it('two defineTask tasks with identical target+function+args -> ONE wire call under dedupe: true; both get the same value', async () => { + const { executor, invocations } = makeEchoExecutor() + const taskA = defineTask((t) => ({ v: t.call({ target: ADDR, abi: testAbi, functionName: 'getVal', args: [1n] }) })) + const taskB = defineTask((t) => ({ v: t.call({ target: ADDR, abi: testAbi, functionName: 'getVal', args: [1n] }) })) + + const [resultA, resultB] = await runMultistepTasks(executor, [taskA, taskB], { dedupe: true }) + + expect(invocations()).toHaveLength(1) // one batch + expect(invocations()[0]).toHaveLength(1) // ONE wire call, not two + expect(resultA!.v).toBe(resultB!.v) + }) + + it('the same identical call issues TWO wire calls when dedupe is off (default)', async () => { + const { executor, invocations } = makeEchoExecutor() + const taskA = defineTask((t) => ({ v: t.call({ target: ADDR, abi: testAbi, functionName: 'getVal', args: [1n] }) })) + const taskB = defineTask((t) => ({ v: t.call({ target: ADDR, abi: testAbi, functionName: 'getVal', args: [1n] }) })) + + const [resultA, resultB] = await runMultistepTasks(executor, [taskA, taskB]) + + expect(invocations()).toHaveLength(1) + expect(invocations()[0]).toHaveLength(2) + expect(resultA!.v).toBe(resultB!.v) + }) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 2. Spec's corruption test, verbatim scenario: conflicting output ABIs for +// identical calldata must never merge. +// ───────────────────────────────────────────────────────────────────────── + +const abiSingleOut = [ + { type: 'function', name: 'getVals', stateMutability: 'view', inputs: [], outputs: [{ type: 'uint256' }] }, +] as const +const abiPairOut = [ + { + type: 'function', + name: 'getVals', + stateMutability: 'view', + inputs: [], + outputs: [{ type: 'uint256' }, { type: 'uint256' }], + }, +] as const + +describe('F7 dedup — conflicting output ABIs never merge (spec corruption test)', () => { + it('same calldata, conflicting output ABIs (returns uint256 vs returns (uint256,uint256)) -> TWO wire calls, each decodes correctly per its own ABI', async () => { + const invocations: StepCall[][] = [] + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + invocations.push(calls.map((c) => ({ ...c }))) + return calls.map((c): RawResult => { + // Real encode/decode round-trip per call, using ITS OWN abi — proves + // no cross-contamination between the two conflicting declarations + // (each wire call is handled entirely independently since they were + // never merged). + if (c.abi === abiSingleOut) { + const data = encodeAbiParameters(abiSingleOut[0].outputs, [42n]) + const [value] = decodeAbiParameters(abiSingleOut[0].outputs, data) + return { status: 'success', value } + } + const data = encodeAbiParameters(abiPairOut[0].outputs, [42n, 43n]) + const decoded = decodeAbiParameters(abiPairOut[0].outputs, data) + return { status: 'success', value: decoded } + }) + }, + } + + const taskA = defineTask((t) => ({ v: t.call({ target: ADDR, abi: abiSingleOut, functionName: 'getVals' }) })) + const taskB = defineTask((t) => ({ v: t.call({ target: ADDR, abi: abiPairOut, functionName: 'getVals' }) })) + + // taskA/taskB have deliberately DIFFERENT result shapes (that's the + // whole point — conflicting output ABIs) — `unknown` sidesteps forcing + // runMultistepTasks' single TResult onto two genuinely different shapes. + const [resultA, resultB] = await runMultistepTasks(executor, [taskA, taskB] as MultistepTask[], { + dedupe: true, + }) + + expect(invocations).toHaveLength(1) + expect(invocations[0]).toHaveLength(2) // never merged + expect((resultA as { v: unknown }).v).toBe(42n) + expect((resultB as { v: unknown }).v).toEqual([42n, 43n]) + }) + + it('external review (P1) regression: overloads with the same output MULTISET but different input->output pairing never merge', async () => { + // ABI A: f(uint256) -> uint256, f(address) -> bool + // ABI B: f(address) -> uint256, f(uint256) -> bool + // Same two output SHAPES in both ABIs (uint256, bool) — only the pairing + // with inputs differs. The old arity-based fallback (both overloads take + // exactly 1 input, so arity alone never disambiguated) serialized ALL + // same-arity candidates' outputs together, in ABI order — producing an + // IDENTICAL combined signature for A and B regardless of pairing, so + // this scenario used to merge and silently hand taskB back taskA's + // uint256 decode (or vice versa). Selector-based resolution (current + // implementation) recovers the exact matched overload per ABI, so the + // two calls key differently and must never merge. + const abiOverloadA = [ + { type: 'function', name: 'f', stateMutability: 'view', inputs: [{ type: 'uint256' }], outputs: [{ type: 'uint256' }] }, + { type: 'function', name: 'f', stateMutability: 'view', inputs: [{ type: 'address' }], outputs: [{ type: 'bool' }] }, + ] as const + const abiOverloadB = [ + { type: 'function', name: 'f', stateMutability: 'view', inputs: [{ type: 'address' }], outputs: [{ type: 'uint256' }] }, + { type: 'function', name: 'f', stateMutability: 'view', inputs: [{ type: 'uint256' }], outputs: [{ type: 'bool' }] }, + ] as const + + const invocations: StepCall[][] = [] + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + invocations.push(calls.map((c) => ({ ...c }))) + return calls.map((c): RawResult => + c.abi === abiOverloadA ? { status: 'success', value: 111n } : { status: 'success', value: true }, + ) + }, + } + + // Both call `f` with a uint256 arg -> resolves to the `f(uint256)` + // overload in EACH abi -> IDENTICAL calldata (a selector is fixed by a + // function's own name+input-types alone, independent of which ABI array + // declares it or what its sibling overloads are) — but `f(uint256)` + // means uint256-out in ABI A and bool-out in ABI B. + const taskA = defineTask((t) => ({ + v: t.call({ target: ADDR, abi: abiOverloadA, functionName: 'f', args: [1n] }), + })) + const taskB = defineTask((t) => ({ + v: t.call({ target: ADDR, abi: abiOverloadB, functionName: 'f', args: [1n] }), + })) + + const [resultA, resultB] = await runMultistepTasks(executor, [taskA, taskB] as MultistepTask[], { + dedupe: true, + }) + + expect(invocations).toHaveLength(1) + expect(invocations[0]).toHaveLength(2) // must NOT merge + expect((resultA as { v: unknown }).v).toBe(111n) + expect((resultB as { v: unknown }).v).toBe(true) + }) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 3. Case-insensitive target merge. +// ───────────────────────────────────────────────────────────────────────── + +describe('F7 dedup — case-insensitive target merge', () => { + it('a call with target differing only by case still merges under dedupe: true', async () => { + const { executor, invocations } = makeEchoExecutor() + const taskA = defineTask((t) => ({ v: t.call({ target: ADDR, abi: testAbi, functionName: 'getVal', args: [1n] }) })) + const taskB = defineTask((t) => ({ + v: t.call({ target: ADDR_MIXED_CASE, abi: testAbi, functionName: 'getVal', args: [1n] }), + })) + + const [resultA, resultB] = await runMultistepTasks(executor, [taskA, taskB], { dedupe: true }) + + expect(invocations()).toHaveLength(1) + expect(invocations()[0]).toHaveLength(1) + expect(resultA!.v).toBe(resultB!.v) + }) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 3b. Calldata hex-case normalization (bytes/bytesN args). +// ───────────────────────────────────────────────────────────────────────── + +describe('F7 dedup — calldata hex-case normalization (bytes/bytesN args)', () => { + it('external review (P2) regression: a bytes32 arg with the same VALUE but different hex-string casing still merges under dedupe: true', async () => { + const bytesAbi = [ + { + type: 'function', + name: 'getByHash', + stateMutability: 'view', + inputs: [{ name: 'h', type: 'bytes32' }], + outputs: [{ type: 'uint256' }], + }, + ] as const + + // Same bytes VALUE, different hex-string casing. viem's ABI encoder + // preserves a bytes/bytesN arg's casing verbatim in the resulting + // calldata (unlike an address, which is never checksum-cased in + // calldata to begin with) — confirmed empirically: encoding these two + // args produces byte-for-byte-equal calldata once lowercased, but + // DIFFERENT raw strings before that. + const mixedCaseHash = '0xaAbBaAbBaAbBaAbBaAbBaAbBaAbBaAbBaAbBaAbBaAbBaAbBaAbBaAbBaAbBaAbB' as `0x${string}` + const lowerCaseHash = mixedCaseHash.toLowerCase() as `0x${string}` + + const { executor, invocations } = makeEchoExecutor() + const taskA = defineTask((t) => ({ + v: t.call({ target: ADDR, abi: bytesAbi, functionName: 'getByHash', args: [mixedCaseHash] }), + })) + const taskB = defineTask((t) => ({ + v: t.call({ target: ADDR, abi: bytesAbi, functionName: 'getByHash', args: [lowerCaseHash] }), + })) + + const [resultA, resultB] = await runMultistepTasks(executor, [taskA, taskB], { dedupe: true }) + + expect(invocations()).toHaveLength(1) + expect(invocations()[0]).toHaveLength(1) // merged despite differing hex-string casing + expect(resultA!.v).toBe(resultB!.v) + }) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 4. Legacy StepCall is never merged — not even under Presets.throughput. +// ───────────────────────────────────────────────────────────────────────── + +describe('F7 dedup — legacy StepCall is never merged (eligibility is per-call, not per-run)', () => { + it('a legacy hand-authored task + a defineTask task with identical calldata, under the FULL Presets.throughput spread -> two wire calls', async () => { + const { executor, invocations } = makeEchoExecutor() + + let legacyValue: unknown + const legacyTask: MultistepTask<{ v: unknown }> = { + maxStep: 1, + buildStepCalls(step) { + return step === 1 ? [{ key: 'legacy', target: ADDR, abi: testAbi, functionName: 'getVal', args: [1n] }] : [] + }, + consumeStepResults(_step, results) { + const r = results.find((res) => res.key === 'legacy') + legacyValue = r?.status === 'success' ? r.value : undefined + }, + finalize() { + return { v: legacyValue } + }, + } + const typedTask = defineTask((t) => ({ v: t.call({ target: ADDR, abi: testAbi, functionName: 'getVal', args: [1n] }) })) + + const [legacyResult, typedResult] = await runMultistepTasks(executor, [legacyTask, typedTask], { + ...Presets.throughput, + }) + + expect(invocations()).toHaveLength(1) + expect(invocations()[0]).toHaveLength(2) // legacy call never merges, preset or not + expect(legacyResult!.v).toBe(typedResult!.v) + }) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 5. Per-call `dedupe: false` override. +// ───────────────────────────────────────────────────────────────────────── + +describe('F7 dedup — dedupe: false per-call override', () => { + it('a TypedCallSpec call with dedupe: false is never merged even when the run itself has dedupe: true', async () => { + const { executor, invocations } = makeEchoExecutor() + const taskA = defineTask((t) => ({ v: t.call({ target: ADDR, abi: testAbi, functionName: 'getVal', args: [1n] }) })) + const taskB = defineTask((t) => ({ + v: t.call({ target: ADDR, abi: testAbi, functionName: 'getVal', args: [1n], dedupe: false }), + })) + + const [resultA, resultB] = await runMultistepTasks(executor, [taskA, taskB], { dedupe: true }) + + expect(invocations()).toHaveLength(1) + expect(invocations()[0]).toHaveLength(2) + expect(resultA!.v).toBe(resultB!.v) // still correct — just via two independent wire calls + }) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 6. Failure fan-out — a merged group's failure reaches EVERY subscriber. +// ───────────────────────────────────────────────────────────────────────── + +describe('F7 dedup — failure fan-out', () => { + it('a merged group whose wire call resolves with a plain per-call (revert-style) failure fans out to every subscriber, each with its OWN error instance (runSettled)', async () => { + const revertError = new DominoCallError('reverted', { kind: 'revert', data: '0x' }) + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + return calls.map((): RawResult => ({ status: 'failure', error: revertError })) + }, + } + const taskA = defineTask((t) => ({ v: t.call({ target: ADDR, abi: testAbi, functionName: 'getVal', args: [1n] }) })) + const taskB = defineTask((t) => ({ v: t.call({ target: ADDR, abi: testAbi, functionName: 'getVal', args: [1n] }) })) + + const [settledA, settledB] = await runSettled(executor, [taskA, taskB], { dedupe: true }) + + expect(settledA!.status).toBe('rejected') + expect(settledB!.status).toBe('rejected') + const errA = (settledA as { status: 'rejected'; error: unknown }).error as DominoCallError + const errB = (settledB as { status: 'rejected'; error: unknown }).error as DominoCallError + + // External review (P2): a merged group's failure is never the SAME + // `DominoCallError` object handed to two or more subscribers (its own + // `.key` would then read as whichever subscriber's call happened to be + // the wire representative) — each gets its OWN instance... + expect(errA).not.toBe(revertError) + expect(errB).not.toBe(revertError) + expect(errA).not.toBe(errB) + // ...but describing the exact same underlying failure: same kind/data. + expect(errA.kind).toBe('revert') + expect(errB.kind).toBe('revert') + expect(errA.data).toBe('0x') + expect(errB.data).toBe('0x') + }) + + it('merged failure fan-out: each subscriber gets its OWN routing key on its own error clone (metadata regression, external review P2)', async () => { + // taskA has an extra LEADING call so its shared call lands at internal + // key "1" — taskB's shared call (its only call) is key "0". Distinct + // keys make the per-subscriber `.key` assertion below meaningful, + // instead of the two subscriber keys coincidentally matching. + const DUMMY_ADDR = '0x4444444444444444444444444444444444444444' as Address + const revertError = new DominoCallError('reverted', { kind: 'revert', data: '0x' }) + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + return calls.map((c): RawResult => + c.target.toLowerCase() === ADDR.toLowerCase() + ? { status: 'failure', error: revertError } + : { status: 'success', value: 'dummy-ok' }, + ) + }, + } + + const taskA = defineTask((t) => { + t.call({ target: DUMMY_ADDR, abi: testAbi, functionName: 'getVal', args: [0n] }) // internal key "0" + const v = t.call({ target: ADDR, abi: testAbi, functionName: 'getVal', args: [1n] }) // internal key "1" + return { v } + }) + const taskB = defineTask((t) => ({ v: t.call({ target: ADDR, abi: testAbi, functionName: 'getVal', args: [1n] }) })) // internal key "0" + + const [settledA, settledB] = await runSettled(executor, [taskA, taskB], { dedupe: true }) + + expect(settledA!.status).toBe('rejected') + expect(settledB!.status).toBe('rejected') + const errA = (settledA as { status: 'rejected'; error: unknown }).error as DominoCallError + const errB = (settledB as { status: 'rejected'; error: unknown }).error as DominoCallError + + expect(errA.key).toBe('1') // taskA's shared call, NOT the representative's own key + expect(errB.key).toBe('0') // taskB's shared call + expect(errA).not.toBe(errB) + }) + + it('a merged group whose wire call becomes a bisection TERMINAL (transport rejection) fans the same synthesized failure out to every subscriber (runSettled)', async () => { + const NOISE_1 = '0x1111111111111111111111111111111111111111' as Address + const NOISE_2 = '0x2222222222222222222222222222222222222222' as Address + const NOISE_3 = '0x3333333333333333333333333333333333333333' as Address + + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + if (calls.some((c) => c.target.toLowerCase() === ADDR.toLowerCase())) { + throw new Error('transport failure') + } + return calls.map((): RawResult => ({ status: 'success', value: 'ok' })) + }, + } + + const taskA = defineTask((t) => ({ v: t.call({ target: ADDR, abi: testAbi, functionName: 'getVal', args: [1n] }) })) + const taskB = defineTask((t) => ({ v: t.call({ target: ADDR, abi: testAbi, functionName: 'getVal', args: [1n] }) })) + const noiseTaskA = defineTask((t) => ({ + v: t.call({ target: NOISE_1, abi: testAbi, functionName: 'getVal', args: [1n] }), + })) + const noiseTaskB = defineTask((t) => ({ + v: t.call({ target: NOISE_2, abi: testAbi, functionName: 'getVal', args: [1n] }), + })) + const noiseTaskC = defineTask((t) => ({ + v: t.call({ target: NOISE_3, abi: testAbi, functionName: 'getVal', args: [1n] }), + })) + + const [settledA, settledB, settledC1, settledC2, settledC3] = await runSettled( + executor, + [taskA, taskB, noiseTaskA, noiseTaskB, noiseTaskC], + { dedupe: true, batchSize: 4, maxConcurrentBatches: 1, adaptiveBatching: true }, + ) + + expect(settledA!.status).toBe('rejected') + expect(settledB!.status).toBe('rejected') + const errA = (settledA as { status: 'rejected'; error: unknown }).error as DominoCallError + const errB = (settledB as { status: 'rejected'; error: unknown }).error as DominoCallError + expect(errA).toBeInstanceOf(DominoCallError) + expect(errA.kind).toBe('batch') + expect(errB.kind).toBe('batch') + // External review (P2): each subscriber gets its OWN `DominoCallError` + // instance (not the SAME object) — but both wrap the identical + // underlying transport error as `cause`, since it's still "the one + // failure", just described once per recipient. + expect(errA).not.toBe(errB) + expect(errA.cause).toBeDefined() + expect(errA.cause).toBe(errB.cause) + + expect(settledC1!.status).toBe('fulfilled') + expect(settledC2!.status).toBe('fulfilled') + expect(settledC3!.status).toBe('fulfilled') + }) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 7. Dedup + bisection interplay: a poisoned UNIQUE call bisects out of a +// wire list that also contains several merged groups — the merged +// groups' subscribers are unaffected and all succeed. +// ───────────────────────────────────────────────────────────────────────── + +describe('F7 dedup + bisection interplay', () => { + it('merged wire list bisects around a poisoned unique call; every subscriber of the poisoned call fails, every subscriber of a merged group succeeds', async () => { + const POISON = '0x9999999999999999999999999999999999999999' as Address + const TOKEN_X = '0x1111111111111111111111111111111111111111' as Address + const TOKEN_Y = '0x2222222222222222222222222222222222222222' as Address + const TOKEN_Z = '0x3333333333333333333333333333333333333333' as Address + + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + if (calls.some((c) => c.target.toLowerCase() === POISON.toLowerCase())) { + throw new Error('poison') + } + return calls.map((c): RawResult => ({ status: 'success', value: `${c.target.toLowerCase()}-ok` })) + }, + } + + const mk = (target: Address) => + defineTask((t) => ({ v: t.call({ target, abi: testAbi, functionName: 'getVal', args: [1n] }) })) + + const x1 = mk(TOKEN_X) + const x2 = mk(TOKEN_X) + const y1 = mk(TOKEN_Y) + const y2 = mk(TOKEN_Y) + const z1 = mk(TOKEN_Z) + const z2 = mk(TOKEN_Z) + const poisoned = mk(POISON) + + const [rx1, rx2, ry1, ry2, rz1, rz2, rp] = await runSettled(executor, [x1, x2, y1, y2, z1, z2, poisoned], { + dedupe: true, + batchSize: 4, + maxConcurrentBatches: 1, + adaptiveBatching: true, + }) + + expect(rp!.status).toBe('rejected') + for (const r of [rx1, rx2, ry1, ry2, rz1, rz2]) { + expect(r!.status).toBe('fulfilled') + } + + // The mock executor resolves with a string (`"-ok"`), not the + // `bigint` `getVal` is statically typed to return — irrelevant here, + // this test only cares about routing, not decoding — hence the + // `unknown` bounce through before re-asserting the concrete shape. + const valueOf = (r: unknown): string => (r as { status: 'fulfilled'; value: { v: string } }).value.v + expect(valueOf(rx1)).toBe(`${TOKEN_X.toLowerCase()}-ok`) + expect(valueOf(rx1)).toBe(valueOf(rx2)) + expect(valueOf(ry1)).toBe(valueOf(ry2)) + expect(valueOf(rz1)).toBe(valueOf(rz2)) + }) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 8. Hit-rate counting — fixture with known duplicates. +// ───────────────────────────────────────────────────────────────────────── + +describe('F7 dedup — hit-rate counting', () => { + it('100 portfolio entries x 3 calls each, over 10 distinct tokens -> exactly 30 unique wire calls (90% dedup hit rate)', async () => { + const erc20LikeAbi = [ + { type: 'function', name: 'symbol', stateMutability: 'view', inputs: [], outputs: [{ type: 'string' }] }, + { type: 'function', name: 'decimals', stateMutability: 'view', inputs: [], outputs: [{ type: 'uint8' }] }, + { type: 'function', name: 'totalSupply', stateMutability: 'view', inputs: [], outputs: [{ type: 'uint256' }] }, + ] as const + + const distinctTokens: Address[] = Array.from( + { length: 10 }, + (_, i) => `0x${(i + 1).toString(16).padStart(40, '0')}` as Address, + ) + + let wireCallCount = 0 + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + wireCallCount += calls.length + return calls.map((): RawResult => ({ status: 'success', value: 1 })) + }, + } + + const tasks = Array.from({ length: 100 }, (_, i) => { + const token = distinctTokens[i % 10]! + return defineTask((t) => ({ + symbol: t.call({ target: token, abi: erc20LikeAbi, functionName: 'symbol' }), + decimals: t.call({ target: token, abi: erc20LikeAbi, functionName: 'decimals' }), + totalSupply: t.call({ target: token, abi: erc20LikeAbi, functionName: 'totalSupply' }), + })) + }) + + await runMultistepTasks(executor, tasks, { dedupe: true }) + + const naiveCallCount = 100 * 3 + expect(wireCallCount).toBe(30) // 10 distinct tokens x 3 calls each + expect(1 - wireCallCount / naiveCallCount).toBeCloseTo(0.9, 5) + }) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 9. Named vs unnamed tuple outputs -> DIFFERENT keys. +// ───────────────────────────────────────────────────────────────────────── + +describe('F7 dedup — named vs unnamed tuple outputs produce different keys (canon includes names)', () => { + it('two calls with identical calldata but tuple outputs differing only in component names produce DIFFERENT dedupeKeyFor results', () => { + const namedTupleAbi = [ + { + type: 'function', + name: 'getPair', + stateMutability: 'view', + inputs: [], + outputs: [ + { + type: 'tuple', + components: [ + { name: 'a', type: 'uint256' }, + { name: 'b', type: 'uint256' }, + ], + }, + ], + }, + ] as const + const unnamedTupleAbi = [ + { + type: 'function', + name: 'getPair', + stateMutability: 'view', + inputs: [], + outputs: [{ type: 'tuple', components: [{ type: 'uint256' }, { type: 'uint256' }] }], + }, + ] as const + + const namedTask = defineTask((t) => ({ v: t.call({ target: ADDR, abi: namedTupleAbi, functionName: 'getPair' }) })) + const unnamedTask = defineTask((t) => ({ + v: t.call({ target: ADDR, abi: unnamedTupleAbi, functionName: 'getPair' }), + })) + + const [namedCall] = namedTask.buildStepCalls(1) + const [unnamedCall] = unnamedTask.buildStepCalls(1) + + const namedKey = dedupeKeyFor(namedCall!) + const unnamedKey = dedupeKeyFor(unnamedCall!) + + expect(namedKey).toBeDefined() + expect(unnamedKey).toBeDefined() + expect(namedKey).not.toBe(unnamedKey) + }) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 10. Presets.throughput. +// ───────────────────────────────────────────────────────────────────────── + +describe('F7 — Presets.throughput', () => { + it('deep-equals the spec literal', () => { + expect(Presets.throughput).toEqual({ maxConcurrentBatches: 5, adaptiveBatching: true, dedupe: true }) + }) + + it('is a readonly literal at the type level (as const)', () => { + // Type-only check: this function is NEVER invoked (see below) — its body + // exists purely so `tsc` (via `npm run typecheck`, which checks all of + // `src`) verifies the assignment is a compile-time error regardless of + // runtime execution. Never letting this actually run means there is no + // risk of mutating the shared `Presets.throughput` singleton other tests + // (and consumers) rely on. + function typeOnlyReadonlyCheck(): void { + // @ts-expect-error — `as const` makes every property of `Presets.throughput` readonly + Presets.throughput.dedupe = false + } + expect(typeof typeOnlyReadonlyCheck).toBe('function') + + expectTypeOf(Presets.throughput).toEqualTypeOf<{ + readonly maxConcurrentBatches: 5 + readonly adaptiveBatching: true + readonly dedupe: true + }>() + }) +}) diff --git a/src/core/abi.ts b/src/core/abi.ts index d79a38e..3073152 100644 --- a/src/core/abi.ts +++ b/src/core/abi.ts @@ -18,6 +18,7 @@ export { encodeAbiParameters, decodeAbiParameters, encodeDeployData, + toFunctionSelector, } from 'viem/utils' export { parseAbi } from 'viem' diff --git a/src/core/dedupe.ts b/src/core/dedupe.ts new file mode 100644 index 0000000..d408f57 --- /dev/null +++ b/src/core/dedupe.ts @@ -0,0 +1,154 @@ +/** + * F7 — within-step, cross-task call dedup key computation. + * + * Scope (controller decision 3, `src/core/engine.ts`): dedup groups calls + * WITHIN one step, ACROSS tasks, PRE-bisection — i.e. before the wire list + * (the deduped call list) ever reaches `runBatchPool` (`src/core/pool.ts`). + * This module owns only the KEY computation; `engine.ts` owns the actual + * grouping/expansion (see its doc comment for the wire-list/subscriber- + * fan-out data structure). + * + * **Key** = `(target.toLowerCase(), calldata, canonicalOutputSignature)`. + * Calldata already captures selector+inputs, so it alone would merge two + * subscribers that declare DIFFERENT output ABIs for the identical calldata + * — corrupting whichever one didn't "win" the merge (first-decoder-wins). + * `canonicalOutputSignature` is the one extra piece of information that + * prevents that: two calls only merge when they'd also decode identically. + * + * **Eligibility** is a separate, per-call concern (`isDedupeEligible`) — + * `dedupeKeyFor` folds both together and returns `undefined` whenever `call` + * must never be merged with anything, for either reason: + * - not dedup-eligible (a legacy hand-authored `StepCall` never carries + * the `DEDUPE_ELIGIBLE` stamp at all; a `TypedCallSpec` with `dedupe: + * false` is stamped `false` explicitly) — "eligible", not "safe": `view`/ + * `pure` alone do not guarantee referential transparency, so this is an + * opt-in the CALLER makes, not something domino infers from ABI + * mutability. + * - `encodeFunctionData` throws while computing the calldata portion of + * the key (e.g. args that don't match the ABI's input types) — the spec + * requires dedup to never introduce a NEW failure mode, so a + * keying-time failure simply falls back to "never merge this call"; + * the executor still runs it (as its own wire call) and produces + * whatever error it normally would for bad args, downstream, unchanged. + */ + +import type { Abi, AbiFunction } from 'abitype' +import type { StepCall } from './types' +import { encodeFunctionData, toFunctionSelector } from './abi' +import { DEDUPE_ELIGIBLE } from './internal' + +/** + * Structural supertype every real `AbiParameter` already satisfies — + * `AbiParameter` (abitype) is a discriminated union where `components` + * exists ONLY on the tuple/tuple-array member, so accessing it unconditionally + * (as the spec's `canon`, below, does) does not type-check without first + * narrowing on `type`. Typing `canon`'s parameter against this looser + * structural shape instead of `AbiParameter` directly keeps the function + * body — the actual canonicalization logic — character-identical to the + * spec text; only the parameter's TYPE ANNOTATION differs. + */ +type CanonParam = { + readonly name?: string | undefined + readonly type: string + readonly components?: readonly CanonParam[] | undefined +} + +/** + * Spec text, VERBATIM (order-preserving, names included — see the module + * doc's "Key" section for why names matter: viem decodes a named tuple to an + * object and an unnamed one to an array, so a component's `name` affects the + * DECODED SHAPE, not just cosmetics). Only the object-key order of the + * produced `{ name, type, components }` representation is normalized (that's + * an inherent property of building a fresh object literal here) — the + * OUTPUT-ARRAY order and TUPLE-COMPONENT order themselves are never sorted, + * because both are semantic. + */ +const canon = (p: CanonParam): unknown => ({ name: p.name ?? '', type: p.type, components: p.components?.map(canon) }) + +/** + * True iff `call` opted into dedup eligibility. Reads the internal + * `DEDUPE_ELIGIBLE` symbol stamped by `defineTask.ts` — absent entirely on + * any hand-authored legacy `StepCall` (never eligible, no mutability promise + * was ever made for it), `true` by default on a compiled `TypedCallSpec` + * call, `false` when that spec set `dedupe: false`. + */ +export function isDedupeEligible(call: StepCall): boolean { + return (call as unknown as Record)[DEDUPE_ELIGIBLE] === true +} + +/** + * The ABI function item that `calldata` was ACTUALLY encoded against — + * resolved by SELECTOR, not by name+arity (external review, P1: arity alone + * cannot disambiguate two same-arity overloads, e.g. `f(uint256)` and + * `f(address)` both take exactly one input; the old arity-based fallback + * conflated such overloads' outputs together, which could produce IDENTICAL + * serialized signatures for two ABIs that pair inputs to outputs + * differently — exactly the corruption the key exists to prevent). + * + * A function selector is the first 4 bytes of `keccak256(signature)`, fixed + * by the function's OWN name+input-types alone — never by which ABI array + * it's declared in, nor by that ABI's OTHER overloads. Since `calldata` was + * produced by `encodeFunctionData` from this exact `(abi, functionName, + * args)` triple, its first 4 bytes are the selector of whichever single + * item `encodeFunctionData` actually resolved `functionName`/`args` to — + * finding the same-named candidate whose OWN computed selector matches + * therefore recovers that EXACT item, unambiguously, with no arity + * heuristic and no risk of conflating two overloads' outputs. + * + * Returns `undefined` if no same-named candidate's selector matches (should + * not happen given `calldata` was just encoded from this same abi, but + * handled defensively — see `dedupeKeyFor`'s catch-all "never merge" + * fallback). + */ +function matchedFunctionFor(abi: Abi, functionName: string, calldata: `0x${string}`): AbiFunction | undefined { + const selector = calldata.slice(0, 10).toLowerCase() + for (const item of abi) { + if (item.type !== 'function' || item.name !== functionName) continue + if (toFunctionSelector(item).toLowerCase() === selector) return item + } + return undefined +} + +/** `JSON.stringify(matchedItem.outputs?.map(canon) ?? [])` per the spec — + * now always the true single matched item (see `matchedFunctionFor`), so + * this is exactly the spec's literal formula with no ambiguous-fallback + * branch needed. */ +function canonicalOutputSignature(matched: AbiFunction): string { + return JSON.stringify(matched.outputs?.map(canon) ?? []) +} + +/** + * Dedup key for `call`, or `undefined` when it must never be merged with + * anything (see the module doc's "Eligibility" section for both reasons, + * plus `matchedFunctionFor`'s defensive `undefined` case). Callers + * (`src/core/engine.ts`) treat `undefined` identically regardless of WHICH + * reason produced it: the call simply gets its own wire-list entry. + */ +export function dedupeKeyFor(call: StepCall): string | undefined { + if (!isDedupeEligible(call)) return undefined + try { + const calldata = encodeFunctionData({ + abi: call.abi, + functionName: call.functionName, + // StepCall.args is untyped by design; viem validates at runtime — same + // pattern as Eip1193Executor's own encodeFunctionData calls. + // eslint-disable-next-line @typescript-eslint/no-explicit-any + args: call.args as any, + }) + const matched = matchedFunctionFor(call.abi, call.functionName, calldata) + if (!matched) return undefined + const signature = canonicalOutputSignature(matched) + // External review (P2): lowercase `calldata` itself before keying — a + // `bytes`/`bytesN` arg's ENCODED segment preserves the caller's own hex + // casing verbatim (unlike an `address`, which viem/Solidity ABI-encodes + // without checksum casing to begin with), so two calls with the same + // bytes VALUE but different hex-string casing (`0xaAbB` vs `0xaabb`) + // would otherwise key differently despite being byte-for-byte identical + // calldata. Only the KEY is normalized here — the call's own `calldata` + // value (and whatever the executor actually sends on the wire) is + // completely untouched. + return JSON.stringify([call.target.toLowerCase(), calldata.toLowerCase(), signature]) + } catch { + return undefined + } +} diff --git a/src/core/defineTask.ts b/src/core/defineTask.ts index e98c80d..ebe3dc2 100644 --- a/src/core/defineTask.ts +++ b/src/core/defineTask.ts @@ -26,7 +26,7 @@ import { parseAbiMemoized } from './abi' import { DominoCallError as E } from './errors' // alias only shortens THIS file's source; esbuild resolves imports to the shared top-level binding when bundling, so it costs nothing either way import { DIAGNOSTICS as DG } from './runSettled' import type { TaskDiagnostics, DiagnosticsCarrier } from './runSettled' -import { SINGLE_USE } from './internal' +import { SINGLE_USE, DEDUPE_ELIGIBLE } from './internal' import type { SingleUseCarrier } from './internal' // ─── Public types ─────────────────────────────────────────────────────────── @@ -124,12 +124,6 @@ export interface TaskBuilder { // ─── Internal graph representation ───────────────────────────────────────── -/** Internal, non-exported marker stamped on every compiled StepCall. `true` - * unless the spec had `dedupe: false`. Nothing reads it yet — dedup (1.2) - * does. Deliberately NOT exported (not even for tests): presence/value are - * verified via `Object.getOwnPropertySymbols` + `Symbol.description`. */ -const DEDUPE_ELIGIBLE = Symbol('domino.dedupeEligible') - /** * Internal graph node — call and derive nodes share ONE flat (mostly * optional-field) shape rather than a discriminated union: `fn` present ⟺ diff --git a/src/core/engine.ts b/src/core/engine.ts index 82d2e49..4a94d09 100644 --- a/src/core/engine.ts +++ b/src/core/engine.ts @@ -19,10 +19,70 @@ * pipeline (`prepareRun`/`resolvePinnedBlock`, `internal.ts`) and the final * `finalize()` pass — both runners diverge sharply there (throw vs settle * per task) and gain nothing from sharing. + * + * ## Within-step, cross-task dedup (F7) + * + * Added after call collection, strictly PRE-bisection: the wire list handed + * to `runBatchPool` (batching/bisection, `src/core/pool.ts`) is the DEDUPED + * list, never the raw per-task call list. `src/core/dedupe.ts` owns the key + * computation (`dedupeKeyFor`) — this module owns the grouping and, on the + * way back, the result FAN-OUT. + * + * **Data structure:** `mapping` used to be one `{ taskIndex, key }` entry + * per WIRE call (1:1, T14-era). It is now one `{ subscribers: { taskIndex, + * key }[] }` entry per wire call — a 1:N "who receives this wire call's + * result" list. With `options.dedupe` off (default), every entry's + * `subscribers` array has exactly one element and the shape degrades to the + * old 1:1 behavior byte-for-byte (same iteration count, same routing). With + * it on, a merged group's entry carries one subscriber per (taskIndex, key) + * that asked for that exact call — every call whose computed key already + * has a group gets appended to that entry's `subscribers` list INSTEAD of + * contributing its own slot to the wire list; a call that is ineligible, or + * whose key computation itself failed (`dedupeKeyFor` returns `undefined` + * either way — see its doc comment), always gets its own singleton entry, + * exactly like the "off" case. + * + * **Wire-list order:** built by iterating tasks/calls in their original + * collection order and pushing a NEW wire entry only the first time a given + * key (or an ineligible/unkeyed call) is seen — so the wire list is "first + * representative of each group + every ineligible call, in original + * relative order of first occurrence", per spec. Batching/bisection then + * operate on exactly that list, with no awareness dedup ever happened. + * + * **Fan-out on the way back:** after the pool settles, each wire call's + * ONE `RawResult` (success or failure — a bisection-terminal synthesized + * failure from `recordTerminal` is just another `RawResult` by the time it + * reaches this routing loop, indistinguishable from an ordinary one) is + * routed to EVERY subscriber in that entry's list, not just one — so a + * merged group's failure (transport-terminal via bisection, or a plain + * per-call revert) reaches every task that asked for it, and a success + * reaches all of them too. Perf: `dedupeKeyFor` is only ever called when + * `options.dedupe` is true (see the call-collection loop below) — the + * default path never computes a key at all. + * + * **Per-subscriber error identity (external review, P2):** a shared failure + * is never hung on the exact same `DominoCallError` OBJECT for two or more + * subscribers — that object's own `.key` field would then read as whichever + * subscriber's call happened to be the wire representative, wrong for every + * OTHER subscriber reading it off their own `StepResult`. When an entry has + * more than one subscriber, each gets its OWN `DominoCallError` instance — + * same `message`/`kind`/`data`/`target`/`functionName`/`cause` (the cause + * REFERENCE is shared, never re-wrapped — it's still "the one transport/ + * revert failure", just described once per recipient), `key` set to THAT + * subscriber's own routing key. A single-subscriber entry (the overwhelming + * common case — dedup off, or a merge of exactly nobody-else) keeps routing + * the original object untouched, so identity pins elsewhere in the codebase + * (e.g. bisection's error-identity determinism tests) are unaffected. A + * non-`DominoCallError` failure (a custom `StepExecutor` resolving a + * `RawResult` failure with some other thrown value, which carries no `.key` + * field to begin with) is always shared as-is — there is no per-call + * metadata on it to correct. */ import type { MultistepTask, StepCall, StepResult, RawResult } from './types' import { runBatchPool } from './pool' +import { dedupeKeyFor } from './dedupe' +import { DominoCallError } from './errors' /** * The hooks that fully capture `run` vs `runSettled`'s divergence — 3 as of @@ -91,6 +151,37 @@ export interface StepEngineOptions { * when `adaptiveBatching` is true; otherwise any rejection is terminal * on its first occurrence regardless of this value (T14 behavior). */ maxBatchAttempts: number + /** F7 — see `BatchOptions.dedupe`. Gates whether this step's call + * collection groups eligible calls before building the wire list — see + * the module doc's "Within-step, cross-task dedup" section. `false` + * (default) never calls `dedupeKeyFor` at all. */ + dedupe: boolean +} + +/** One wire-list entry's "who receives this call's result" list — see the + * module doc's "Data structure" section. Exactly one element unless F7 + * dedup actually merged two or more subscribers into this entry. */ +interface WireMappingEntry { + subscribers: { taskIndex: number; key: string }[] +} + +/** + * A fresh `DominoCallError` describing the SAME failure as `original`, but + * addressed to `key` — see the module doc's "Per-subscriber error identity" + * section. `cause` is passed through as the SAME reference (never re-wrapped + * — `original.cause`, not `original`), so every clone of one merged failure + * still shares one underlying transport/revert cause identity even though + * each now has its own wrapper object and its own `key`. + */ +function retargetError(original: DominoCallError, key: string): DominoCallError { + return new DominoCallError(original.message, { + kind: original.kind, + ...(original.data !== undefined ? { data: original.data } : {}), + ...(original.target !== undefined ? { target: original.target } : {}), + ...(original.functionName !== undefined ? { functionName: original.functionName } : {}), + ...(original.cause !== undefined ? { cause: original.cause } : {}), + key, + }) } /** @@ -111,14 +202,45 @@ export async function runSteps( for (let step = 1; step <= maxStep; step++) { const calls: StepCall[] = [] - const mapping: { taskIndex: number; key: string }[] = [] + const mapping: WireMappingEntry[] = [] + + // F7 dedup: only allocated (and only ever consulted) when the option is + // on — dedup key string -> index of that group's entry in `mapping`/ + // `calls`. Left `undefined` on the default path, which is what makes + // `dedupeKeyFor` provably never run below (see the module doc's "Perf" + // note). + const groupIndexByKey = options.dedupe ? new Map() : undefined for (let taskIndex = 0; taskIndex < ts.length; taskIndex++) { const stepCalls = policy.buildStepCalls(taskIndex, step) if (stepCalls === undefined) continue for (const call of stepCalls) { + const subscriber = { taskIndex, key: call.key } + + // `dedupeKeyFor` returns `undefined` for BOTH reasons a call must + // never be merged (ineligible, or keying itself failed) — either + // way it falls straight through to the "own wire entry" branch + // below, identically to the dedup-off path. + const dedupeKey = groupIndexByKey ? dedupeKeyFor(call) : undefined + + if (dedupeKey !== undefined) { + const existingIndex = groupIndexByKey!.get(dedupeKey) + if (existingIndex !== undefined) { + // Already have a representative wire call for this key — this + // call contributes no new wire-list entry, just another + // subscriber to the existing one. + mapping[existingIndex]!.subscribers.push(subscriber) + continue + } + // First call seen for this key — it becomes the group's + // representative. Record its (about-to-be-pushed) index BEFORE + // pushing, so it equals `calls.length` at the moment of the push + // below. + groupIndexByKey!.set(dedupeKey, calls.length) + } + calls.push(call) - mapping.push({ taskIndex, key: call.key }) + mapping.push({ subscribers: [subscriber] }) } } @@ -155,21 +277,43 @@ export async function runSteps( const mappingEntry = mapping[globalIndex] globalIndex++ if (!mappingEntry) continue - const { taskIndex, key } = mappingEntry const result = batchResults[i] as RawResult - const list = perTaskResults[taskIndex]! - if (result.status === 'success') { - list.push({ status: 'success', key, value: result.value }) - } else { - // Forward the SAME error object — never wrap or discard it. - // exactOptionalPropertyTypes-safe: only include `error` when the - // RawResult actually carried one. - list.push({ - status: 'failure', - key, - ...('error' in result && result.error !== undefined ? { error: result.error } : {}), - }) + // F7 fan-out: route this ONE wire result to EVERY subscriber — + // success and failure alike, and regardless of whether the + // failure came from an ordinary per-call revert or a bisection + // TERMINAL synthesis (`recordTerminal`, `src/core/pool.ts`): by + // the time it's a `RawResult` here, both are indistinguishable, + // so every subscriber of a merged group sees the same outcome. A + // non-merged call's entry has exactly one subscriber, so this + // degrades to the pre-F7 single-route behavior identically. + const failureError = result.status === 'failure' && 'error' in result ? result.error : undefined + const isSharedFailure = failureError !== undefined && mappingEntry.subscribers.length > 1 + + for (const { taskIndex, key } of mappingEntry.subscribers) { + const list = perTaskResults[taskIndex]! + if (result.status === 'success') { + list.push({ status: 'success', key, value: result.value }) + } else { + // Per-subscriber error identity (external review, P2): a + // failure shared by more than one subscriber never hands out + // the SAME `DominoCallError` object twice — each subscriber + // gets its own instance addressed to ITS OWN key (see + // `retargetError` / the module doc's "Per-subscriber error + // identity" section). A single-subscriber entry (the default, + // dedup-off path) forwards the original object completely + // unchanged — exactOptionalPropertyTypes-safe: only include + // `error` when the RawResult actually carried one. + const error = + isSharedFailure && failureError instanceof DominoCallError + ? retargetError(failureError, key) + : failureError + list.push({ + status: 'failure', + key, + ...(error !== undefined ? { error } : {}), + }) + } } } } diff --git a/src/core/internal.ts b/src/core/internal.ts index e5d4d24..5ad8627 100644 --- a/src/core/internal.ts +++ b/src/core/internal.ts @@ -63,24 +63,41 @@ export interface SingleUseCarrier { * vice versa. */ const consumed = new WeakSet() +/** + * Internal marker stamped on every compiled `StepCall` by `defineTask()`'s + * `buildStepCalls` (F7 — relocated here from `defineTask.ts` so the engine + * can read it without a cross-layer import back into the task-builder + * module). `true` unless the originating `TypedCallSpec` had `dedupe: + * false`; a hand-authored legacy `StepCall` never carries this symbol at + * all, so `call[DEDUPE_ELIGIBLE] === true` is the one true "is this call + * eligible for within-step dedup" check — see `src/core/dedupe.ts`. + * + * Deliberately NOT exported from `src/index.ts` (not even for tests): + * presence/value are verified via `Object.getOwnPropertySymbols` + + * `Symbol.description`, exactly as `SINGLE_USE` above. + */ +export const DEDUPE_ELIGIBLE: unique symbol = Symbol('domino.dedupeEligible') + /** TS-only convenience for the brand-check casts below — erased at compile * time, so using it at 3 call sites (instead of a real `isBranded()` * function) costs zero extra runtime bytes over one. */ type Branded = MultistepTask & SingleUseCarrier /** - * Numeric (+ one boolean) options validated + defaulted by `validateOptions` - * (F6a/F6b). All four ride through `prepareRun`'s return value. + * Numeric (+ two boolean) options validated + defaulted by `validateOptions` + * (F6a/F6b/F7). All five ride through `prepareRun`'s return value. * `maxBatchAttempts` and `adaptiveBatching` are both consumed by the engine * (`src/core/engine.ts`/`src/core/pool.ts`, F6b) — `adaptiveBatching` gates * whether bisection ever runs at all; `maxBatchAttempts` bounds it once it - * does. + * does. `dedupe` (F7) gates the engine's within-step, cross-task call + * dedup — see `src/core/dedupe.ts`. */ export interface ValidatedRunOptions { batchSize: number maxConcurrentBatches: number maxBatchAttempts: number adaptiveBatching: boolean + dedupe: boolean } /** Options `validateOptions` reads — a structural subset of `BatchOptions`. */ @@ -89,6 +106,7 @@ export interface NumericOptionsInput { maxConcurrentBatches?: number maxBatchAttempts?: number adaptiveBatching?: boolean + dedupe?: boolean } /** @@ -134,7 +152,12 @@ export function validateOptions(o: NumericOptionsInput | undefined): ValidatedRu // caller has opted in with knowledge of their transport's failure modes). const adaptiveBatching = o?.adaptiveBatching ?? false - return { batchSize, maxConcurrentBatches, maxBatchAttempts, adaptiveBatching } + // Plain boolean flag (F7), default `false` — see `BatchOptions.dedupe`'s + // doc comment. Off by default so dedup key computation never runs on the + // hot path unless a caller opts in (`Presets.throughput` does). + const dedupe = o?.dedupe ?? false + + return { batchSize, maxConcurrentBatches, maxBatchAttempts, adaptiveBatching, dedupe } } /** diff --git a/src/core/presets.ts b/src/core/presets.ts new file mode 100644 index 0000000..bff4289 --- /dev/null +++ b/src/core/presets.ts @@ -0,0 +1,31 @@ +/** + * F7 — off-the-shelf `BatchOptions` bundles. + * + * `Presets.throughput` turns on every concurrency/dedup knob this release + * train has shipped for the common "many independent contract reads, one + * chain, one block" workload (e.g. resolving a portfolio of tokens/vaults): + * concurrent batch dispatch (F6a, `maxConcurrentBatches`), adaptive + * bisection (F6b, `adaptiveBatching`), and within-step cross-task dedup + * (F7, `dedupe`). + * + * It deliberately does NOT set `batchSize`, `maxBatchAttempts`, or `block` + * — spread it ahead of your own overrides: + * + * ```ts + * await resolver.run(tasks, { ...Presets.throughput, batchSize: 200 }) + * ``` + * + * (A future `pinBlock` capability will compose the same way once it ships — + * `{ ...Presets.throughput, pinBlock: true }` — but that option does not + * exist yet as of this release.) + * + * `dedupe: true` here can never change legacy-task semantics: dedup only + * ever merges calls stamped dedup-ELIGIBLE (a compiled `TypedCallSpec` call, + * eligible by default unless its spec set `dedupe: false`) — a + * hand-authored legacy `MultistepTask`'s `StepCall`s carry no such stamp and + * are therefore never merged, preset or not. See `BatchOptions.dedupe`'s + * doc comment (`src/core/runMultistepTasks.ts`) for the full contract. + */ +export const Presets = { + throughput: { maxConcurrentBatches: 5, adaptiveBatching: true, dedupe: true } as const, +} diff --git a/src/core/runMultistepTasks.ts b/src/core/runMultistepTasks.ts index d9703a5..26a97c3 100644 --- a/src/core/runMultistepTasks.ts +++ b/src/core/runMultistepTasks.ts @@ -97,6 +97,33 @@ export interface BatchOptions { /** Block to query at (defaults to 'latest'). Same block used for ALL steps. */ block?: BlockParam + + /** + * Enable within-step, cross-task call dedup (F7): before the wire list for + * a step reaches batching/bisection, calls that are dedup-ELIGIBLE and + * share the same `(target.toLowerCase(), calldata, canonicalOutputSignature)` + * key are merged into a single physical call — its result (success or + * failure) is then fanned out to every subscriber. Default: `false`. + * + * **Eligibility is per-call**, not per-run: a hand-authored legacy + * `StepCall` is never eligible (no mutability promise was ever made for + * it), so turning this on can never change legacy-task semantics — see + * `TypedCallSpec.dedupe` (default `true`; set `dedupe: false` on an + * individual call to opt it out). `view`/`pure` alone do not guarantee + * referential transparency, hence "eligible", not "safe" — this is a + * caller opt-in, not something inferred from ABI mutability. + * + * **Conflicting output ABIs never merge:** calldata alone (selector + + * inputs) does not capture how a caller intends to DECODE the result — two + * subscribers declaring different output shapes for identical calldata are + * kept as separate wire calls, each decoding correctly against its own + * ABI, so dedup can never silently corrupt one decoder with another's + * shape. + * + * See `Presets.throughput` for a ready-made bundle that turns this on + * together with `maxConcurrentBatches`/`adaptiveBatching`. + */ + dedupe?: boolean } /** @@ -150,7 +177,7 @@ export async function runMultistepTasks( // -> mark-consumed), see `src/core/internal.ts`. Only branded tasks // (`defineTask`/`buildErc20Task`/`buildErc4626Task` output) are affected — // legacy `MultistepTask`s pass through every step as a no-op. - const { batchSize, maxConcurrentBatches, adaptiveBatching, maxBatchAttempts } = prepareRun(ts, options) + const { batchSize, maxConcurrentBatches, adaptiveBatching, maxBatchAttempts, dedupe } = prepareRun(ts, options) await resolvePinnedBlock() // `run`'s fail-fast policy (F6a/F6b) — see `src/core/engine.ts`'s doc @@ -186,7 +213,7 @@ export async function runMultistepTasks( }, } - await runSteps(ts, { batchSize, maxConcurrentBatches, adaptiveBatching, maxBatchAttempts }, policy) + await runSteps(ts, { batchSize, maxConcurrentBatches, adaptiveBatching, maxBatchAttempts, dedupe }, policy) return ts.map((task) => task.finalize()) } diff --git a/src/core/runSettled.ts b/src/core/runSettled.ts index 075a86a..058bbe4 100644 --- a/src/core/runSettled.ts +++ b/src/core/runSettled.ts @@ -156,7 +156,7 @@ export async function runSettled( // `runSettled(executor, [], { batchSize: 0 })` rejects rather than // silently resolving to `[]` (unchanged 1.0 ordering — contrast with // `runMultistepTasks`, which checks the empty-tasks shortcut first). - const { batchSize, maxConcurrentBatches, adaptiveBatching, maxBatchAttempts } = prepareRun(ts, options) + const { batchSize, maxConcurrentBatches, adaptiveBatching, maxBatchAttempts, dedupe } = prepareRun(ts, options) if (ts.length === 0) return [] @@ -243,7 +243,7 @@ export async function runSettled( }, } - await runSteps(ts, { batchSize, maxConcurrentBatches, adaptiveBatching, maxBatchAttempts }, policy) + await runSteps(ts, { batchSize, maxConcurrentBatches, adaptiveBatching, maxBatchAttempts, dedupe }, policy) return ts.map((task, i): SettledTaskResult => { // Read the task's OWN live diagnostics if it carries the channel (defineTask, diff --git a/src/index.ts b/src/index.ts index 1e3069d..2b7dafe 100644 --- a/src/index.ts +++ b/src/index.ts @@ -13,6 +13,9 @@ export type { } from './core/types' export type { BatchOptions } from './core/runMultistepTasks' +// Presets (F7) — off-the-shelf BatchOptions bundles +export { Presets } from './core/presets' + // Per-task settlement (F5) export { runSettled } from './core/runSettled' export type { TaskDiagnostics, SettledTaskResult } from './core/runSettled'