From 3e5eaa52133549c70a294b2259a0e9bac17765cf Mon Sep 17 00:00:00 2001 From: halaprix Date: Thu, 23 Jul 2026 18:44:15 +0200 Subject: [PATCH 1/2] feat: concurrency pool + fail-fast cancellation (F6a) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Unify runMultistepTasks/runSettled's duplicated per-step loop (F5 ledger debt) into one shared engine (src/core/engine.ts) driven by a concurrency- limited batch pool (src/core/pool.ts). Adds maxConcurrentBatches (default 1, genuinely sequential and bit-identical to 1.1) and maxBatchAttempts (validated + default-computed, not yet consumed — bisection lands later) to BatchOptions. run() gains fail-fast cancellation: on the first irrecoverable batch rejection, queued batches are not dispatched, in-flight batches settle without leaking, and the run rejects with the lowest-index discovered terminal error. runSettled() never cancels — its executeBatch policy converts transport rejections into recorded kind:'batch' failures instead of ever rejecting into the pool. Adds a global unhandledRejection test guard (vitest setupFiles, both configs) and src/__tests__/concurrency.test.ts covering wall-clock parallelism, completion-order-independent routing (fuzzed), fail-fast (a)-(d), error-identity determinism over 100 runs, runSettled's never- cancel behavior, numeric validation, and the step barrier under concurrency. --- README.md | 2 +- src/__tests__/concurrency.test.ts | 520 ++++++++++++++++++++ src/__tests__/setup/unhandled-rejections.ts | 45 ++ src/core/engine.ts | 158 ++++++ src/core/internal.ts | 79 ++- src/core/pool.ts | 172 +++++++ src/core/runMultistepTasks.ts | 132 +++-- src/core/runSettled.ts | 141 +++--- vitest.compat-dist.config.ts | 1 + vitest.config.ts | 1 + 10 files changed, 1083 insertions(+), 168 deletions(-) create mode 100644 src/__tests__/concurrency.test.ts create mode 100644 src/__tests__/setup/unhandled-rejections.ts create mode 100644 src/core/engine.ts create mode 100644 src/core/pool.ts diff --git a/README.md b/README.md index af716a5..157b3bb 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-10.9KB-brightgreen)](https://www.npmjs.com/package/@halaprix/domino) +[![bundle size](https://img.shields.io/badge/gzip-11.5KB-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/src/__tests__/concurrency.test.ts b/src/__tests__/concurrency.test.ts new file mode 100644 index 0000000..d1474c5 --- /dev/null +++ b/src/__tests__/concurrency.test.ts @@ -0,0 +1,520 @@ +import { describe, it, expect, vi } from 'vitest' +import { runMultistepTasks } from '../core/runMultistepTasks' +import { runSettled } from '../core/runSettled' +import { runBatchPool } from '../core/pool' +import { defineTask } from '../core/defineTask' +import { DominoCallError } from '../core/errors' +import type { Address, MultistepTask, StepCall, StepResult, StepExecutor, RawResult } from '../core/types' +import type { BatchOptions } from '../core/runMultistepTasks' + +/** + * F6a — concurrency pool + fail-fast cancellation. + * + * See `src/core/pool.ts` (the concurrency primitive) and + * `src/core/engine.ts` (the shared step loop `run`/`runSettled` now both + * drive) for the design this suite pins. The global unhandled-rejection + * guard (`src/__tests__/setup/unhandled-rejections.ts`) fails any test in + * this file (or any other) that leaks a rejection — several assertions + * below rely on it instead of re-implementing their own leak detection. + */ + +const ADDR = '0xA0b86991c6218b36c1d19D4a2e9Eb004C35d5Cc4' as Address + +const testAbi = [ + { + type: 'function', + name: 'getNum', + stateMutability: 'view', + inputs: [], + outputs: [{ type: 'uint256' }], + }, +] as const + +function sleep(ms: number): Promise { + return new Promise((resolve) => setTimeout(resolve, ms)) +} + +function call(key: string): StepCall { + return { key, target: ADDR, abi: [], functionName: 'f' } +} + +/** A single-step task with `n` calls keyed `c0..c(n-1)`. Captures whatever + * it's given via consumeStepResults so tests can assert routing/dead-ness. */ +function makeCallsTask( + n: number, + opts?: { onConsume?: (step: number, results: StepResult[]) => void }, +): MultistepTask { + let captured: StepResult[] = [] + return { + maxStep: 1, + buildStepCalls(step) { + if (step !== 1) return [] + return Array.from({ length: n }, (_, i) => call('c' + i)) + }, + consumeStepResults(step, results) { + captured = results + opts?.onConsume?.(step, results) + }, + finalize() { + return captured + }, + } +} + +function okExecutor(): StepExecutor { + return { + async executeMulticall(calls: StepCall[]): Promise { + return calls.map((): RawResult => ({ status: 'success', value: 10n })) + }, + } +} + +// ───────────────────────────────────────────────────────────────────────── +// 1. Wall-clock — concurrency pool actually parallelizes within a step. +// ───────────────────────────────────────────────────────────────────────── + +describe('wall-clock concurrency', () => { + it('600 calls / bs=100 / conc=5 completes in ~2 sequential-RT (two waves)', async () => { + const latencyMs = 50 + const invocationSizes: number[] = [] + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + invocationSizes.push(calls.length) + await sleep(latencyMs) + return calls.map((): RawResult => ({ status: 'success', value: 1 })) + }, + } + + const task = makeCallsTask(600) + const start = performance.now() + await runMultistepTasks(executor, [task], { batchSize: 100, maxConcurrentBatches: 5 }) + const elapsed = performance.now() - start + + expect(invocationSizes).toHaveLength(6) // 600 / 100 = 6 physical batches + // 6 batches at concurrency 5 -> 2 waves (5 + 1). Generous bounds against + // CI jitter: [1.5x, 3.5x] one batch's latency. + expect(elapsed).toBeGreaterThanOrEqual(latencyMs * 1.5) + expect(elapsed).toBeLessThanOrEqual(latencyMs * 3.5) + }, 20_000) + + it('600 calls / bs=100 / conc=1 is genuinely sequential (~6 sequential-RT)', async () => { + const latencyMs = 50 + const invocationSizes: number[] = [] + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + invocationSizes.push(calls.length) + await sleep(latencyMs) + return calls.map((): RawResult => ({ status: 'success', value: 1 })) + }, + } + + const task = makeCallsTask(600) + const start = performance.now() + await runMultistepTasks(executor, [task], { batchSize: 100, maxConcurrentBatches: 1 }) + const elapsed = performance.now() - start + + expect(invocationSizes).toHaveLength(6) + // Default (maxConcurrentBatches: 1) must be genuinely sequential: 6 + // batches back-to-back. Generous bounds: [5x, 8x]. + expect(elapsed).toBeGreaterThanOrEqual(latencyMs * 5) + expect(elapsed).toBeLessThanOrEqual(latencyMs * 8) + }, 20_000) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 2. Routing fuzz — completion-order-independent, index-based routing. +// ───────────────────────────────────────────────────────────────────────── + +describe('routing fuzz', () => { + /** 3 tasks x 2 steps, distinguishable keys — fresh closures every call. */ + function makeFixtureTasks(): MultistepTask>[] { + return [0, 1, 2].map((taskIdx) => { + const captured: Record = {} + return { + maxStep: 2, + buildStepCalls(step) { + if (step === 1) return [0, 1, 2].map((i) => call(`t${taskIdx}-s1-${i}`)) + if (step === 2) return [0, 1].map((i) => call(`t${taskIdx}-s2-${i}`)) + return [] + }, + consumeStepResults(_step, results) { + for (const r of results) { + if (r.status === 'success') captured[r.key] = r.value as string + } + }, + finalize() { + return captured + }, + } + }) + } + + function fuzzExecutor(maxDelayMs: number): StepExecutor { + return { + async executeMulticall(calls: StepCall[]): Promise { + if (maxDelayMs > 0) await sleep(Math.random() * maxDelayMs) + return calls.map((c): RawResult => ({ status: 'success', value: 'v-' + c.key })) + }, + } + } + + it('20+ runs with random per-batch delays at conc=4 match the serial (conc=1) reference exactly', async () => { + const reference = await runMultistepTasks(fuzzExecutor(0), makeFixtureTasks(), { + batchSize: 2, + maxConcurrentBatches: 1, + }) + + for (let i = 0; i < 25; i++) { + const fuzzed = await runMultistepTasks(fuzzExecutor(20), makeFixtureTasks(), { + batchSize: 2, + maxConcurrentBatches: 4, + }) + expect(fuzzed).toEqual(reference) + } + }, 20_000) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 3-6. Fail-fast cancellation — spec (a)-(d). +// ───────────────────────────────────────────────────────────────────────── + +describe('fail-fast cancellation (run)', () => { + it('(a) queued batches are not dispatched after the failure is discovered', async () => { + const dispatched: string[] = [] + const boom = new Error('boom-c1') + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + const key = calls[0]!.key + dispatched.push(key) + if (key === 'c1') throw boom // rejects with no delay + await sleep(20) // c0 (and any other claimed batch) settles slower + return calls.map((): RawResult => ({ status: 'success', value: 1 })) + }, + } + + const task = makeCallsTask(6) + await expect( + runMultistepTasks(executor, [task], { batchSize: 1, maxConcurrentBatches: 2 }), + ).rejects.toBe(boom) + + // Only batch 0 (already claimed, slower) and batch 1 (the poisoned one) + // are ever dispatched — batches 2-5 were still queued when cancellation + // triggered (c1 rejects long before c0's 20ms delay elapses). + expect(dispatched.length).toBeLessThan(6) + expect(new Set(dispatched)).toEqual(new Set(['c0', 'c1'])) + }) + + it('(b)+(d) in-flight batches settle without leaking, and consumeStepResults is never called for the failing step', async () => { + const boom = new Error('boom-c1') + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + const key = calls[0]!.key + if (key === 'c1') throw boom + await sleep(20) + return calls.map((): RawResult => ({ status: 'success', value: 1 })) + }, + } + + const consumeSpy = vi.fn() + const task: MultistepTask = { + maxStep: 1, + buildStepCalls(step) { + if (step !== 1) return [] + return Array.from({ length: 6 }, (_, i) => call('c' + i)) + }, + consumeStepResults: consumeSpy, + finalize() { + return null + }, + } + + await expect( + runMultistepTasks(executor, [task], { batchSize: 1, maxConcurrentBatches: 2 }), + ).rejects.toBe(boom) + + // (d): the failing step's results are discarded — consumeStepResults is + // never invoked for it. (b)'s "no unhandled rejections" half is proven + // by the global afterEach guard, not by an assertion in this test body. + expect(consumeSpy).not.toHaveBeenCalled() + }) + + it('(c) error-identity determinism: single poisoned batch, random delays, conc=3, >=100 runs -> same object every time', async () => { + const boom = new Error('poison-c2') + + for (let i = 0; i < 100; i++) { + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + await sleep(Math.random() * 20) + if (calls[0]!.key === 'c2') throw boom + return calls.map((): RawResult => ({ status: 'success', value: 1 })) + }, + } + const task = makeCallsTask(6) + + await expect( + runMultistepTasks(executor, [task], { batchSize: 1, maxConcurrentBatches: 3 }), + ).rejects.toBe(boom) + } + }, 30_000) + + it('multi-poisoned fixture: thrown error is one of the poisoned batches\' errors; no unhandled rejections', async () => { + const errC1 = new Error('boom-c1') + const errC4 = new Error('boom-c4') + + for (let i = 0; i < 10; i++) { + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + await sleep(Math.random() * 20) + const key = calls[0]!.key + if (key === 'c1') throw errC1 + if (key === 'c4') throw errC4 + return calls.map((): RawResult => ({ status: 'success', value: 1 })) + }, + } + const task = makeCallsTask(6) + + let thrown: unknown + let rejected = false + try { + await runMultistepTasks(executor, [task], { batchSize: 1, maxConcurrentBatches: 3 }) + } catch (err) { + rejected = true + thrown = err + } + expect(rejected).toBe(true) + // Identity is explicitly UNSPECIFIED for multi-failure — only assert + // membership, never which one. + expect([errC1, errC4]).toContain(thrown) + } + }, 20_000) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 7. runSettled never cancels — record-and-continue under concurrency. +// ───────────────────────────────────────────────────────────────────────── + +describe('runSettled concurrency (never cancels)', () => { + it('poisoned batch among 6 (conc=3): ALL 6 batches dispatched, non-poisoned calls succeed, later step still runs', async () => { + const boom = new Error('transport down') + const dispatched: string[] = [] + + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + const key = calls[0]!.key + dispatched.push(key) + await sleep(Math.random() * 15) + if (key === 'a1') throw boom + return calls.map((): RawResult => ({ status: 'success', value: 'ok-' + key })) + }, + } + + let step2Built = false + let taskAStep1: StepResult[] = [] + let taskAStep2: StepResult[] = [] + let taskBStep1: StepResult[] = [] + + const taskA: MultistepTask<{ step1: StepResult[]; step2: StepResult[] }> = { + maxStep: 2, + buildStepCalls(step) { + if (step === 1) return ['a0', 'a1', 'a2'].map(call) + if (step === 2) { + step2Built = true + return [call('a-final')] + } + return [] + }, + consumeStepResults(step, results) { + if (step === 1) taskAStep1 = results + if (step === 2) taskAStep2 = results + }, + finalize() { + return { step1: taskAStep1, step2: taskAStep2 } + }, + } + + const taskB: MultistepTask<{ step1: StepResult[] }> = { + maxStep: 1, + buildStepCalls(step) { + if (step !== 1) return [] + return ['b0', 'b1', 'b2'].map(call) + }, + consumeStepResults(step, results) { + if (step === 1) taskBStep1 = results + }, + finalize() { + return { step1: taskBStep1 } + }, + } + + const [resultA, resultB] = await runSettled(executor, [taskA, taskB], { + batchSize: 1, + maxConcurrentBatches: 3, + }) + + // ALL 6 step-1 batches dispatched — never cancels. (`dispatched` also + // picks up step 2's single call once that step runs; filter it out to + // check step 1 specifically — step 2 running at all is asserted below.) + const step1Dispatched = dispatched.filter((k) => k !== 'a-final').sort() + expect(step1Dispatched).toEqual(['a0', 'a1', 'a2', 'b0', 'b1', 'b2']) + expect(step2Built).toBe(true) + expect(resultA!.status).toBe('fulfilled') + expect(resultB!.status).toBe('fulfilled') + + // Poisoned call carries a kind:'batch' DominoCallError with cause identity. + const a1Result = taskAStep1.find((r) => r.key === 'a1') + expect(a1Result?.status).toBe('failure') + const a1Error = (a1Result as { status: 'failure'; error?: unknown }).error + expect(a1Error).toBeInstanceOf(DominoCallError) + expect((a1Error as DominoCallError).kind).toBe('batch') + expect((a1Error as DominoCallError).cause).toBe(boom) + + // Non-poisoned calls (same task, sibling task) succeeded normally. + expect(taskAStep1.find((r) => r.key === 'a0')).toEqual({ status: 'success', key: 'a0', value: 'ok-a0' }) + expect(taskAStep1.find((r) => r.key === 'a2')).toEqual({ status: 'success', key: 'a2', value: 'ok-a2' }) + expect(taskBStep1).toEqual([ + { status: 'success', key: 'b0', value: 'ok-b0' }, + { status: 'success', key: 'b1', value: 'ok-b1' }, + { status: 'success', key: 'b2', value: 'ok-b2' }, + ]) + + // Later step (a-final) still ran and succeeded. + expect(taskAStep2).toEqual([{ status: 'success', key: 'a-final', value: 'ok-a-final' }]) + }) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 8. Validation — numeric fields, pre-consumption. +// ───────────────────────────────────────────────────────────────────────── + +describe('validation', () => { + const badValues = [0, -1, 1.5, NaN, 2 ** 53] + + for (const field of ['batchSize', 'maxConcurrentBatches', 'maxBatchAttempts'] as const) { + it(`${field}: 0/-1/1.5/NaN/2^53 all throw pre-consumption; the branded task stays fresh`, async () => { + const executor = okExecutor() + const task = defineTask((t) => t.call({ target: ADDR, abi: testAbi, functionName: 'getNum' })) + + for (const bad of badValues) { + await expect( + runMultistepTasks(executor, [task], { [field]: bad } as BatchOptions), + ).rejects.toThrow(`${field} must be a positive integer`) + } + + // None of the invalid attempts consumed the task — a subsequent valid + // run still succeeds (same pattern as singleUse.test.ts's "does-not- + // consume" cases). + const [result] = await runMultistepTasks(executor, [task]) + expect(result).toBe(10n) + }) + } + + it('validation applies identically through runSettled (shared prepareRun)', async () => { + const executor = okExecutor() + const task = defineTask((t) => t.call({ target: ADDR, abi: testAbi, functionName: 'getNum' })) + + await expect( + runSettled(executor, [task], { maxConcurrentBatches: 0 }), + ).rejects.toThrow('maxConcurrentBatches must be a positive integer') + + const [settled] = await runSettled(executor, [task]) + expect(settled).toEqual({ status: 'fulfilled', value: 10n, diagnostics: { optionalFailures: [] } }) + }) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 9. Step barrier — concurrency window never crosses a step boundary. +// ───────────────────────────────────────────────────────────────────────── + +describe('step barrier under concurrency', () => { + it('conc=5: step-2 buildStepCalls fires only after every step-1 batch has executed AND been consumed', async () => { + const log: string[] = [] + + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + const key = calls[0]!.key + if (key.startsWith('s1-')) await sleep(Math.random() * 20) + log.push('batch-end-' + key) + return calls.map((): RawResult => ({ status: 'success', value: 1 })) + }, + } + + const task: MultistepTask = { + maxStep: 2, + buildStepCalls(step) { + if (step === 1) { + log.push('build-1') + return Array.from({ length: 20 }, (_, i) => call('s1-' + i)) + } + if (step === 2) { + log.push('build-2') + return [call('s2-0')] + } + return [] + }, + consumeStepResults(step) { + log.push('consume-' + step) + }, + finalize() { + return null + }, + } + + await runMultistepTasks(executor, [task], { batchSize: 4, maxConcurrentBatches: 5 }) + + const build2Index = log.indexOf('build-2') + const consume1Index = log.indexOf('consume-1') + expect(build2Index).toBeGreaterThan(-1) + expect(consume1Index).toBeGreaterThan(-1) + expect(consume1Index).toBeLessThan(build2Index) + + // Every step-1 batch-end entry happened before consume-1 (which itself + // happens before build-2, checked above). + const step1BatchEnds = log.filter((e) => e.startsWith('batch-end-s1-')) + expect(step1BatchEnds).toHaveLength(5) // 20 calls / batchSize 4 + for (const entry of step1BatchEnds) { + expect(log.indexOf(entry)).toBeLessThan(consume1Index) + } + }, 10_000) +}) + +// ───────────────────────────────────────────────────────────────────────── +// runBatchPool — direct unit coverage (engine refactor coverage). +// ───────────────────────────────────────────────────────────────────────── + +describe('runBatchPool (direct)', () => { + it('empty batches list resolves completed with an empty results array, without calling execute', async () => { + const execute = vi.fn() + const outcome = await runBatchPool([], 4, execute) + expect(outcome).toEqual({ outcome: 'completed', results: [] }) + expect(execute).not.toHaveBeenCalled() + }) + + it('clamps worker count to the number of batches (no idle over-dispatch)', async () => { + let maxConcurrent = 0 + let current = 0 + const batches: StepCall[][] = [[call('a')], [call('b')]] + const outcome = await runBatchPool(batches, 10, async (batch) => { + current++ + maxConcurrent = Math.max(maxConcurrent, current) + await sleep(5) + current-- + return batch.map((): RawResult => ({ status: 'success', value: 1 })) + }) + expect(outcome.outcome).toBe('completed') + expect(maxConcurrent).toBe(2) // never exceeds batches.length even though maxConcurrentBatches=10 + }) + + it('routes results by original batch index regardless of completion order', async () => { + const batches: StepCall[][] = [[call('slow')], [call('fast')]] + const outcome = await runBatchPool(batches, 2, async (batch) => { + const key = batch[0]!.key + await sleep(key === 'slow' ? 20 : 0) + return [{ status: 'success', value: key }] + }) + expect(outcome.outcome).toBe('completed') + if (outcome.outcome === 'completed') { + expect(outcome.results[0]).toEqual([{ status: 'success', value: 'slow' }]) + expect(outcome.results[1]).toEqual([{ status: 'success', value: 'fast' }]) + } + }) +}) diff --git a/src/__tests__/setup/unhandled-rejections.ts b/src/__tests__/setup/unhandled-rejections.ts new file mode 100644 index 0000000..f9136fb --- /dev/null +++ b/src/__tests__/setup/unhandled-rejections.ts @@ -0,0 +1,45 @@ +/** + * Global `unhandledRejection` guard (F6a, controller decision 6). + * + * Wired into BOTH `vitest.config.ts` and `vitest.compat-dist.config.ts` via + * `test.setupFiles`, so it protects the entire suite — not just the new + * concurrency tests. Any promise that rejects without a handler attached + * anywhere in the process during a test run is collected here; `afterEach` + * asserts nothing was collected and clears the list, failing whichever test + * was running when the leak surfaced. + * + * This is the harness-side half of the fail-fast cancellation contract's + * spec (b) ("in-flight batches ... their rejections attached (`.catch(noop)`) + * so no unhandled rejections escape") — `src/core/pool.ts`'s design makes + * that true by construction (every batch promise is awaited directly inside + * its own worker's try/catch), and this guard is what actually PROVES it + * across the test suite rather than merely asserting it in the design doc. + */ + +import { afterEach } from 'vitest' + +const leaked: unknown[] = [] + +function onUnhandledRejection(reason: unknown): void { + leaked.push(reason) +} + +process.on('unhandledRejection', onUnhandledRejection) + +afterEach(() => { + if (leaked.length === 0) return + + const captured = leaked.slice() + leaked.length = 0 + + const details = captured + .map((reason, i) => { + const label = reason instanceof Error ? (reason.stack ?? reason.message) : String(reason) + return ` [${i}] ${label}` + }) + .join('\n') + + throw new Error( + `${captured.length} unhandled rejection(s) leaked during this test:\n${details}`, + ) +}) diff --git a/src/core/engine.ts b/src/core/engine.ts new file mode 100644 index 0000000..9f85fca --- /dev/null +++ b/src/core/engine.ts @@ -0,0 +1,158 @@ +/** + * Shared step-execution engine (F6a) — unifies the per-step loop that + * `runMultistepTasks` and `runSettled` used to duplicate (ledger debt from + * F5). `runSteps` owns everything that is bit-identical between the two + * runners: per-step call collection, batch slicing, dispatch through the + * concurrency pool (`src/core/pool.ts`), and index-based routing of results + * back into per-task `StepResult[]` arrays. + * + * The ONLY behavior that differs between `run` (fail-fast) and `runSettled` + * (record-and-continue) is captured in the 3-hook `StepEnginePolicy` each + * runner builds for itself, closing over its own state (`ts`, and — for + * `runSettled` only — its `dead`/`deadError` bookkeeping). `runSteps` never + * sees that state; it just calls the hooks and trusts their return values. + * + * What stays OUTSIDE this module (deliberately): the F2 consumption + * 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. + */ + +import type { MultistepTask, StepCall, StepResult, RawResult } from './types' +import { runBatchPool } from './pool' + +/** + * The 3 hooks that fully capture `run` vs `runSettled`'s divergence. Each + * runner builds one instance of this per call, as plain closures over its + * own local state — see `runMultistepTasks.ts`/`runSettled.ts`. + */ +export interface StepEnginePolicy { + /** + * Build one task's calls for `step`. + * + * Return `undefined` to skip this task-step entirely (either it's + * inactive because `step > task.maxStep`, or — `runSettled` only — the + * task is already dead). Error handling is entirely the policy's + * business: `run`'s hook lets a throw propagate straight out of + * `runSteps` (a plain synchronous throw inside this async function's + * body — fatal, matches pre-F6a behavior exactly); `runSettled`'s hook + * catches internally, marks the task dead, and returns `undefined`. + */ + buildStepCalls(taskIndex: number, step: number): StepCall[] | undefined + + /** + * Execute ONE physical batch of calls. This is the only hook + * `runBatchPool` actually calls, so it is the only one whose rejection + * behavior matters for cancellation: + * - `run`'s hook is the raw, unwrapped `executor.executeMulticall` call + * (plus the length-mismatch guard) — a rejection here is exactly what + * triggers fail-fast. + * - `runSettled`'s hook catches a transport rejection and *resolves* + * with synthesized per-call `kind:'batch'` failures instead; it only + * ever rejects for the length-mismatch executor-bug case (which is + * deliberately NOT converted into a recorded failure). + */ + executeBatch(batch: StepCall[]): Promise + + /** + * Consume one task's results for `step`. Same error-handling split as + * `buildStepCalls`: `run` lets a throw propagate; `runSettled` catches + * and marks the task dead. + */ + consumeStepResults(taskIndex: number, step: number, results: StepResult[]): void +} + +/** The subset of `BatchOptions` the engine itself needs. */ +export interface StepEngineOptions { + batchSize: number + maxConcurrentBatches: number +} + +/** + * Drive every step from 1..maxStep: collect calls, slice into batches, + * dispatch through the concurrency pool, route results back by task index, + * then hand each task its per-step results via the policy. + * + * Throws (aborting all remaining steps) iff the pool reports a cancelled + * outcome, or `buildStepCalls`/`consumeStepResults` throws — the latter two + * propagate directly (no pool involvement; they aren't batch executions). + */ +export async function runSteps( + ts: MultistepTask[], + options: StepEngineOptions, + policy: StepEnginePolicy, +): Promise { + const maxStep = ts.reduce((max, task) => (task.maxStep > max ? task.maxStep : max), 0) + + for (let step = 1; step <= maxStep; step++) { + const calls: StepCall[] = [] + const mapping: { taskIndex: number; key: string }[] = [] + + for (let taskIndex = 0; taskIndex < ts.length; taskIndex++) { + const stepCalls = policy.buildStepCalls(taskIndex, step) + if (stepCalls === undefined) continue + for (const call of stepCalls) { + calls.push(call) + mapping.push({ taskIndex, key: call.key }) + } + } + + // Pre-allocate a 2D array indexed by taskIndex for O(1) result grouping + // — same rationale as the pre-F6a code this replaces. + const perTaskResults: StepResult[][] = Array.from({ length: ts.length }, () => []) + + if (calls.length > 0) { + const batches: StepCall[][] = [] + for (let batchStart = 0; batchStart < calls.length; batchStart += options.batchSize) { + batches.push(calls.slice(batchStart, batchStart + options.batchSize)) + } + + const outcome = await runBatchPool(batches, options.maxConcurrentBatches, (batch) => + policy.executeBatch(batch), + ) + + if (outcome.outcome === 'cancelled') { + // Spec (d): in-flight results are discarded — we never look at + // whatever the pool may have partially collected, and we never + // reach the consumeStepResults dispatch loop below for this step. + throw outcome.error + } + + // Route each batch's results back into the shared perTaskResults + // arrays, indexed by ORIGINAL batch index (not completion order) — + // this is what makes routing completion-order-independent. + let globalIndex = 0 + for (let b = 0; b < batches.length; b++) { + const batchResults = outcome.results[b]! + for (let i = 0; i < batchResults.length; i++) { + 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 } : {}), + }) + } + } + } + } + + // Dispatch to every task index — including inactive/dead ones, whose + // hook is expected to no-op — so consumeStepResults is invoked + // consistently each step for every task that's still active. + for (let i = 0; i < ts.length; i++) { + policy.consumeStepResults(i, step, perTaskResults[i]!) + } + } +} diff --git a/src/core/internal.ts b/src/core/internal.ts index 0dbdf0c..d87c1e4 100644 --- a/src/core/internal.ts +++ b/src/core/internal.ts @@ -69,15 +69,63 @@ const consumed = new WeakSet() type Branded = MultistepTask & SingleUseCarrier /** - * Pipeline step 1 — the existing `batchSize` validation (message/behavior - * unchanged from 1.0). A programmer error: failure does NOT consume + * Numeric options validated + defaulted by `validateOptions` (F6a). All + * three ride through `prepareRun`'s return value; `maxBatchAttempts` is + * validated and defaulted here but not yet CONSUMED by either runner — + * bisection (T15) wires it into the engine without touching this function + * again. + */ +export interface ValidatedRunOptions { + batchSize: number + maxConcurrentBatches: number + maxBatchAttempts: number +} + +/** Options `validateOptions` reads — a structural subset of `BatchOptions`. */ +export interface NumericOptionsInput { + batchSize?: number + maxConcurrentBatches?: number + maxBatchAttempts?: number +} + +/** + * Shared positive-safe-integer check for all three F6a numeric fields. + * `Number.isSafeInteger` (not `Number.isInteger`) — a value like `2**53` + * passes `isInteger` but is not exactly representable, so it must still be + * rejected; message wording is unchanged across all three fields (mirrors + * the original `batchSize`-only message from 1.0/1.1). + */ +function validatePositiveSafeInteger(name: string, value: number): void { + if (!Number.isSafeInteger(value) || value < 1) { + throw new Error(`${name} must be a positive integer, got ${value}`) + } +} + +/** + * Pipeline step 1 — numeric option validation (message/behavior for + * `batchSize` alone unchanged from 1.0 in every case any existing test + * exercises; F6a extends its check from `Number.isInteger` to + * `Number.isSafeInteger` and adds the same check for `maxConcurrentBatches`/ + * `maxBatchAttempts`). A programmer error: failure does NOT consume * anything, because nothing has been touched yet. + * + * `maxBatchAttempts`'s default is computed from the RESOLVED `batchSize` + * (so a caller-supplied `batchSize` changes the default), per the spec: + * `2 * Math.ceil(Math.log2(batchSize)) + 1` — `batchSize: 1` gives + * `log2(1) = 0`, `ceil(0) = 0`, default `1`. */ -export function validateOptions(batchSize: number | undefined): number { - const resolved = batchSize ?? 100 - if (!Number.isInteger(resolved) || resolved < 1) - throw new Error(`batchSize must be a positive integer, got ${resolved}`) - return resolved +export function validateOptions(o: NumericOptionsInput | undefined): ValidatedRunOptions { + const batchSize = o?.batchSize ?? 100 + validatePositiveSafeInteger('batchSize', batchSize) + + const maxConcurrentBatches = o?.maxConcurrentBatches ?? 1 + validatePositiveSafeInteger('maxConcurrentBatches', maxConcurrentBatches) + + const defaultMaxBatchAttempts = 2 * Math.ceil(Math.log2(batchSize)) + 1 + const maxBatchAttempts = o?.maxBatchAttempts ?? defaultMaxBatchAttempts + validatePositiveSafeInteger('maxBatchAttempts', maxBatchAttempts) + + return { batchSize, maxConcurrentBatches, maxBatchAttempts } } /** @@ -166,19 +214,20 @@ export async function resolvePinnedBlock(): Promise { } /** - * Shared synchronous prefix of both runners' bodies: validate batchSize → - * reject duplicate branded instances → pin-capability check → mark branded - * tasks consumed. Returns the validated `batchSize` so each runner can drop - * this straight in place of its old inline validation. + * Shared synchronous prefix of both runners' bodies: validate numeric + * options → reject duplicate branded instances → pin-capability check → + * mark branded tasks consumed. Returns the validated options bundle so each + * runner can drop this straight in place of its old inline validation. * * Deliberately stops here (not shared with `resolvePinnedBlock`/the step * loop): `run()`/`runSettled()` diverge in failure handling once execution - * actually starts, so those stay in each runner's own body. + * actually starts, so those stay in each runner's own body (see + * `src/core/engine.ts`'s `runSteps` for the part that IS now shared). */ -export function prepareRun(ts: MultistepTask[], o: { batchSize?: number } | undefined): number { - const batchSize = validateOptions(o?.batchSize) +export function prepareRun(ts: MultistepTask[], o: NumericOptionsInput | undefined): ValidatedRunOptions { + const validated = validateOptions(o) rejectDuplicateInstances(ts) validatePinCapability() markTasksConsumed(ts) - return batchSize + return validated } diff --git a/src/core/pool.ts b/src/core/pool.ts new file mode 100644 index 0000000..3a8cb73 --- /dev/null +++ b/src/core/pool.ts @@ -0,0 +1,172 @@ +/** + * Concurrency-limited batch dispatch pool + fail-fast cancellation (F6a). + * + * `runBatchPool` is the one piece of genuinely new machinery this feature + * adds. It knows nothing about tasks, steps, or which runner (`run` vs + * `runSettled`) is calling it — only "batches" (`StepCall[]`) and an + * `execute` callback that resolves with `RawResult[]` or rejects. See + * `src/core/engine.ts` for how the two runners' policies feed this. + * + * **Why `runSettled` "never cancels" needs no special case here:** + * `runSettled`'s `executeBatch` policy (see `runSettled.ts`) catches every + * ordinary transport rejection itself and *resolves* with synthesized + * per-call `kind:'batch'` `DominoCallError`s — so, from this module's point + * of view, that `execute` callback simply never rejects for a call failure. + * The cancellation machinery below is real code that runs for `runSettled` + * too (nothing is skipped), it just never gets triggered because nothing + * ever throws into it. `runSettled`'s `executeBatch` DOES still reject for + * the length-mismatch executor-bug case — and when it does, this pool's + * ordinary fail-fast handling applies (matches the existing "aborts + * runSettled entirely" behavior for that case, unchanged from 1.0/1.1 at + * the default `maxConcurrentBatches: 1`). + * + * ## Cancellation policy — spec (a)-(d) + * + * (a) **Queued batches are not dispatched.** Every worker checks the shared + * `cancelled` flag BEFORE claiming its next batch index. Claiming + * (`nextIndex++`) and the check happen with no `await` between them, so + * — JavaScript being single-threaded/run-to-completion between awaits — + * two workers can never claim the same index, and once `cancelled` + * flips, no worker claims a further index. + * + * (b) **In-flight batches are allowed to settle; no unhandled rejections.** + * A worker that has already claimed a batch keeps `await`-ing it to + * completion regardless of `cancelled` (there is no way to abort an + * in-flight promise, and the spec requires letting it finish rather + * than abandoning it). Because every batch's promise is awaited + * directly inside its own worker's `try/catch` — never handed off + * un-awaited — there is never a floating rejection for the runtime to + * report as `unhandledRejection`. The spec's "attach `.catch(noop)`" + * requirement falls out of this shape for free; there is no separate + * `.catch()` bolted onto anything. + * + * (c) **Thrown error = lowest (batchIndex, callIndex) among DISCOVERED + * terminal errors.** Every worker that observes a rejection pushes + * `{ batchIndex, callIndex, error }` onto a shared list; once all + * workers finish, the lowest entry (by this module's `compareTerminal`) + * is returned. "Discovered" is exactly "belongs to a batch that was + * actually claimed" — a batch still queued when cancellation triggers + * never contributes an entry. `callIndex` is always `0` in this + * feature slice (T14): every terminal error today is whole-batch, so + * there is no finer-grained call index yet. The shape survives + * bisection (T15) unmodified — a future per-call terminal discovered + * inside a split sub-batch will carry its own real `callIndex`, and the + * comparator needs no redesign. + * + * For exactly one failing batch, the selected error is deterministic + * regardless of timing: dispatch always claims indices in ascending + * order, so the single poisoned batch is always eventually claimed (no + * other rejection exists to stop the pool before that happens), and it + * always settles (in-flight batches are never abandoned, per (b)) — so + * the discovered-error list always ends up with exactly that one entry. + * With two or more concurrently-failing batches, whether a lower-index + * poisoned batch was already claimed by the time a higher-index one's + * rejection flips `cancelled` depends on the concurrency window and + * relative delay — genuinely timing-dependent, as the spec states. + * + * (d) **Results of in-flight batches are discarded.** A batch's result is + * only written into the shared `results` array `if (!cancelled)` at the + * moment it settles; a cancelled run's `results` array is never even + * inspected — see `runSteps` in `src/core/engine.ts`, which throws + * immediately on a `cancelled` outcome instead of routing anything. + * + * **`maxConcurrentBatches: 1` is not special-cased.** `workerCount = min(1, + * batchCount) = 1` naturally yields exactly one worker processing batches + * strictly in order, one at a time — on rejection there is no other worker + * with anything in flight, so the pool immediately reports `cancelled` with + * the untouched original error object. This is bit-identical to the + * pre-F6a sequential path, including error identity (pinned by the compat + * suite and `singleUse.test.ts`). + * + * **Forward-compat note (T15):** `batches.length` is read fresh on every + * loop check rather than captured once up front. This costs nothing today + * (the array is never mutated after `runSteps` builds it), but means a + * future "central per-step queue" that pushes split child batches onto this + * same array mid-run would already be picked up by the existing workers + * without a redesign. + */ + +import type { StepCall, RawResult } from './types' + +/** All batches ran to completion — `results[i]` is batch `i`'s `RawResult[]`. */ +export interface BatchPoolCompleted { + outcome: 'completed' + results: RawResult[][] +} + +/** Fail-fast triggered — `error` is the selected terminal error (see (c) above). */ +export interface BatchPoolCancelled { + outcome: 'cancelled' + error: unknown +} + +export type BatchPoolOutcome = BatchPoolCompleted | BatchPoolCancelled + +/** One discovered terminal error, tagged with where it came from — see (c). */ +interface TerminalCandidate { + batchIndex: number + callIndex: number + error: unknown +} + +/** + * Orders discovered terminal errors for spec (c)'s selection rule: lowest + * original batch index, then lowest call index within that batch. Every + * T14 terminal is whole-batch (`callIndex` always `0`), so this currently + * degrades to plain `batchIndex` ordering — kept as a 2-key comparator so + * bisection (T15) can slot per-call terminals in without redesigning this + * function. + */ +function compareTerminal(a: TerminalCandidate, b: TerminalCandidate): number { + if (a.batchIndex !== b.batchIndex) return a.batchIndex - b.batchIndex + return a.callIndex - b.callIndex +} + +/** + * Dispatch `batches` through `execute` with at most `maxConcurrentBatches` + * in flight at once, in original index order, applying fail-fast + * cancellation ((a)-(d) above) the instant any dispatched batch rejects. + * + * `maxConcurrentBatches` is clamped to `batches.length` — requesting more + * workers than there are batches would just spin up idle workers that exit + * immediately. + */ +export async function runBatchPool( + batches: StepCall[][], + maxConcurrentBatches: number, + execute: (batch: StepCall[], batchIndex: number) => Promise, +): Promise { + if (batches.length === 0) return { outcome: 'completed', results: [] } + + const results: RawResult[][] = new Array(batches.length) + const terminalErrors: TerminalCandidate[] = [] + let nextIndex = 0 + let cancelled = false + + async function worker(): Promise { + for (;;) { + if (cancelled) return + if (nextIndex >= batches.length) return + const batchIndex = nextIndex + nextIndex++ + const batch = batches[batchIndex]! + try { + const batchResults = await execute(batch, batchIndex) + if (!cancelled) results[batchIndex] = batchResults + } catch (error) { + cancelled = true + terminalErrors.push({ batchIndex, callIndex: 0, error }) + } + } + } + + const workerCount = Math.min(maxConcurrentBatches, batches.length) + await Promise.all(Array.from({ length: workerCount }, () => worker())) + + if (terminalErrors.length > 0) { + terminalErrors.sort(compareTerminal) + return { outcome: 'cancelled', error: terminalErrors[0]!.error } + } + + return { outcome: 'completed', results } +} diff --git a/src/core/runMultistepTasks.ts b/src/core/runMultistepTasks.ts index 5b66539..6272203 100644 --- a/src/core/runMultistepTasks.ts +++ b/src/core/runMultistepTasks.ts @@ -15,6 +15,7 @@ import type { MultistepTask, StepCall, StepResult, StepExecutor, RawResult, BlockParam } from './types' import { prepareRun, resolvePinnedBlock } from './internal' +import { runSteps, type StepEnginePolicy } from './engine' /** * Options for runMultistepTasks. @@ -25,10 +26,41 @@ export interface BatchOptions { * * Multicall3 aggregate3 has a per-call gas limit. When a single step has * more than this many calls, it is split into sequential batches. - * Default: 100. Must be a positive integer — anything else throws. + * Default: 100. Must be a positive safe integer ≥ 1 — anything else throws. */ batchSize?: number + /** + * Concurrency-limited pool size for dispatching physical batches within a + * single step (F6a). Batches are dispatched in original index order as + * permits free; routing is index-based and does not depend on completion + * order. Default: 1 (genuinely sequential — batch k+1 is not dispatched + * until batch k settles, bit-identical to the pre-F6a behavior). Must be + * a positive safe integer ≥ 1 — anything else throws. + * + * **Cancellation policy (fail-fast):** on the first batch whose executor + * call rejects, no further queued batches for that step are dispatched; + * batches already in flight are allowed to settle (their results + * discarded); once every in-flight batch has settled, the run rejects + * with the selected terminal error (see `src/core/pool.ts` for the exact + * selection rule). This is deterministic only when exactly one batch + * fails irrecoverably — with multiple concurrently-failing batches, which + * one's error is thrown is explicitly unspecified (timing-dependent). + */ + maxConcurrentBatches?: number + + /** + * Maximum retry attempts for a failing physical batch before it is + * treated as a terminal failure (adaptive bisection, F6b). Default: + * `2 * Math.ceil(Math.log2(batchSize)) + 1`. Must be a positive safe + * integer ≥ 1 — anything else throws. + * + * Validated and defaulted as of F6a, but NOT YET consumed — every batch + * failure is currently terminal on first rejection (no retry/bisection). + * Adaptive bisection lands in a later release. + */ + maxBatchAttempts?: number + /** Block to query at (defaults to 'latest'). Same block used for ALL steps. */ block?: BlockParam } @@ -84,83 +116,45 @@ 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 = prepareRun(ts, options) + const { batchSize, maxConcurrentBatches } = prepareRun(ts, options) await resolvePinnedBlock() - const maxStep = ts.reduce((max, task) => (task.maxStep > max ? task.maxStep : max), 0) - - for (let step = 1; step <= maxStep; step++) { - const calls: StepCall[] = [] - const mapping: { taskIndex: number; key: string }[] = [] - - for (let taskIndex = 0; taskIndex < ts.length; taskIndex++) { + // `run`'s fail-fast policy (F6a) — see `src/core/engine.ts`'s doc comment + // for the full hook contract. Every hook here either lets a throw + // propagate untouched (buildStepCalls/consumeStepResults — a plain + // synchronous throw inside this async function's body, exactly like the + // pre-F6a code) or is the raw unwrapped executor call (executeBatch), + // whose rejection is what `runBatchPool` treats as the fail-fast trigger. + const policy: StepEnginePolicy = { + buildStepCalls(taskIndex, step) { const task = ts[taskIndex]! - if (step > task.maxStep) continue - - const stepCalls = task.buildStepCalls(step) - for (const call of stepCalls) { - calls.push(call) - mapping.push({ taskIndex, key: call.key }) + if (step > task.maxStep) return undefined + return task.buildStepCalls(step) + }, + + async executeBatch(batch) { + const results = await executor.executeMulticall(batch, options?.block) + + // Dev-time guard: a misbehaving executor that returns fewer results than + // calls would silently corrupt routing — fail loudly instead. + if (results.length !== batch.length) { + throw new Error( + `StepExecutor returned ${results.length} results for ${batch.length} calls — length mismatch`, + ) } - } - - // Pre-allocate a 2D array indexed by taskIndex for O(1) result grouping. - // Avoids Map hashing overhead — taskIndex is sequential zero-based, so - // array indexing is both faster and simpler. - const perTaskResults: StepResult[][] = Array.from({ length: ts.length }, () => []) - - // Only hit the network when there are calls; a step where every active task - // built nothing still dispatches empty results below (consistent per-step - // notification regardless of sibling tasks). - if (calls.length > 0) { - // Split calls into batches to stay under per-call gas limits. - // Each batch executes as a separate multicall round-trip. - for (let batchStart = 0; batchStart < calls.length; batchStart += batchSize) { - const batch = calls.slice(batchStart, batchStart + batchSize) - const results = await executor.executeMulticall(batch, options?.block) - - // Dev-time guard: a misbehaving executor that returns fewer results than - // calls would silently corrupt routing — fail loudly instead. - if (results.length !== batch.length) { - throw new Error( - `StepExecutor returned ${results.length} results for ${batch.length} calls — length mismatch`, - ) - } + return results + }, - // Route this batch's results into the shared perTaskResults arrays. - for (let i = 0; i < results.length; i++) { - const mappingEntry = mapping[batchStart + i] - if (!mappingEntry) continue - const { taskIndex, key } = mappingEntry - const result = results[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 } : {}), - }) - } - } - } - } - - // Dispatch to every task active at this step — including those that built no - // calls — so consumeStepResults is invoked consistently each step. - for (let i = 0; i < ts.length; i++) { - const task = ts[i] + consumeStepResults(taskIndex, step, results) { + const task = ts[taskIndex] if (task && step <= task.maxStep) { - task.consumeStepResults(step, perTaskResults[i]!) + task.consumeStepResults(step, results) } - } + }, } + await runSteps(ts, { batchSize, maxConcurrentBatches }, policy) + return ts.map((task) => task.finalize()) } diff --git a/src/core/runSettled.ts b/src/core/runSettled.ts index 232ebc7..9779593 100644 --- a/src/core/runSettled.ts +++ b/src/core/runSettled.ts @@ -37,10 +37,11 @@ * `consumeStepResults`/`finalize`) throws". */ -import type { MultistepTask, StepCall, StepResult, StepExecutor, RawResult, Address } from './types' +import type { MultistepTask, StepExecutor, RawResult, Address } from './types' import type { BatchOptions } from './runMultistepTasks' import { DominoCallError } from './errors' import { prepareRun, resolvePinnedBlock } from './internal' +import { runSteps, type StepEnginePolicy } from './engine' /** * Internal-only diagnostics channel (F2). A compiled `defineTask()` output @@ -116,116 +117,90 @@ 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 = prepareRun(ts, options) + const { batchSize, maxConcurrentBatches } = prepareRun(ts, options) if (ts.length === 0) return [] await resolvePinnedBlock() - const maxStep = ts.reduce((max, task) => (task.maxStep > max ? task.maxStep : max), 0) - // Dead-task bookkeeping: once true, no further buildStepCalls/consumeStepResults // calls happen for that task index, and finalize() is skipped entirely. + // Lives here (not in the shared engine) — it's `runSettled`-specific state + // that its policy hooks close over; the engine never sees it directly. const dead: boolean[] = new Array(ts.length).fill(false) const deadError: unknown[] = new Array(ts.length) - for (let step = 1; step <= maxStep; step++) { - const calls: StepCall[] = [] - const mapping: { taskIndex: number; key: string }[] = [] - - for (let taskIndex = 0; taskIndex < ts.length; taskIndex++) { - if (dead[taskIndex]) continue + // `runSettled`'s record-and-continue policy (F6a) — see + // `src/core/engine.ts`'s doc comment for the full hook contract. + const policy: StepEnginePolicy = { + buildStepCalls(taskIndex, step) { + if (dead[taskIndex]) return undefined const task = ts[taskIndex]! - if (step > task.maxStep) continue - - let stepCalls: StepCall[] + if (step > task.maxStep) return undefined try { - stepCalls = task.buildStepCalls(step) + return task.buildStepCalls(step) } catch (err) { dead[taskIndex] = true deadError[taskIndex] = err - continue + return undefined } + }, - for (const call of stepCalls) { - calls.push(call) - mapping.push({ taskIndex, key: call.key }) - } - } - - const perTaskResults: StepResult[][] = Array.from({ length: ts.length }, () => []) - - if (calls.length > 0) { - for (let batchStart = 0; batchStart < calls.length; batchStart += batchSize) { - const batch = calls.slice(batchStart, batchStart + batchSize) - - let results: RawResult[] - try { - results = await executor.executeMulticall(batch, options?.block) - } catch (transportError) { - // Batch-level failure: every call in THIS physical batch fails with - // its own DominoCallError, all sharing the SAME cause. Do not - // rethrow — later batches/steps still execute. - results = batch.map( - (call): RawResult => ({ - status: 'failure', - error: new DominoCallError(`Call ${call.key} failed: containing batch rejected`, { - kind: 'batch', - cause: transportError, - target: call.target, - functionName: call.functionName, - key: call.key, - }), + async executeBatch(batch) { + let results: RawResult[] + try { + results = await executor.executeMulticall(batch, options?.block) + } catch (transportError) { + // Batch-level failure: every call in THIS physical batch fails with + // its own DominoCallError, all sharing the SAME cause. Resolve + // (never reject) — this is what makes `runSettled` never cancel: + // `runBatchPool`'s fail-fast machinery only triggers on a REJECTION + // from this hook, and an ordinary transport failure never produces + // one here. Later batches/steps still execute. + return batch.map( + (call): RawResult => ({ + status: 'failure', + error: new DominoCallError(`Call ${call.key} failed: containing batch rejected`, { + kind: 'batch', + cause: transportError, + target: call.target, + functionName: call.functionName, + key: call.key, }), - ) - } - - // Dev-time guard (same as `run`): an executor that resolved but - // returned the wrong number of results would silently corrupt - // routing. This is an executor-implementation bug, not a call - // failure, so it is NOT converted into a batch StepResult — it - // aborts `runSettled` entirely, same as an invalid batchSize. - if (results.length !== batch.length) { - throw new Error( - `StepExecutor returned ${results.length} results for ${batch.length} calls — length mismatch`, - ) - } + }), + ) + } - for (let i = 0; i < results.length; i++) { - const mappingEntry = mapping[batchStart + i] - if (!mappingEntry) continue - const { taskIndex, key } = mappingEntry - const result = results[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. - list.push({ - status: 'failure', - key, - ...('error' in result && result.error !== undefined ? { error: result.error } : {}), - }) - } - } + // Dev-time guard (same as `run`): an executor that resolved but + // returned the wrong number of results would silently corrupt + // routing. This is an executor-implementation bug, not a call + // failure, so it is NOT converted into a batch StepResult — it + // rejects this hook, which aborts `runSettled` entirely (via the + // pool's ordinary fail-fast path) same as an invalid batchSize. + if (results.length !== batch.length) { + throw new Error( + `StepExecutor returned ${results.length} results for ${batch.length} calls — length mismatch`, + ) } - } + return results + }, - for (let i = 0; i < ts.length; i++) { - if (dead[i]) continue - const task = ts[i] + consumeStepResults(taskIndex, step, results) { + if (dead[taskIndex]) return + const task = ts[taskIndex] if (task && step <= task.maxStep) { try { - task.consumeStepResults(step, perTaskResults[i]!) + task.consumeStepResults(step, results) } catch (err) { - dead[i] = true - deadError[i] = err + dead[taskIndex] = true + deadError[taskIndex] = err } } - } + }, } + await runSteps(ts, { batchSize, maxConcurrentBatches }, policy) + return ts.map((task, i): SettledTaskResult => { // Read the task's OWN live diagnostics if it carries the channel (defineTask, // F2); every legacy MultistepTask lacks [DIAGNOSTICS] and falls back to an diff --git a/vitest.compat-dist.config.ts b/vitest.compat-dist.config.ts index b09287b..6d23472 100644 --- a/vitest.compat-dist.config.ts +++ b/vitest.compat-dist.config.ts @@ -14,5 +14,6 @@ export default defineConfig({ globals: true, environment: "node", include: ["src/__tests__/compat/**/*.test.ts"], + setupFiles: ["./src/__tests__/setup/unhandled-rejections.ts"], }, }); diff --git a/vitest.config.ts b/vitest.config.ts index 448e0de..c653251 100644 --- a/vitest.config.ts +++ b/vitest.config.ts @@ -14,5 +14,6 @@ export default defineConfig({ globals: true, environment: "node", include: ["src/__tests__/**/*.test.ts", "src/__tests__/**/*.test-d.ts"], + setupFiles: ["./src/__tests__/setup/unhandled-rejections.ts"], }, }); \ No newline at end of file From 129092407e22044aa70bb6391f8824fa422e7eb1 Mon Sep 17 00:00:00 2001 From: halaprix Date: Thu, 23 Jul 2026 18:57:35 +0200 Subject: [PATCH 2/2] fix: snapshot pool results at settlement, harden rejection-guard timing, use observed concurrency in timing tests MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit External review (3 findings), same branch: - pool.ts: clone each batch's result array/wrappers at the moment its own await resolves (results[i] = batchResults.map(r => ({...r}))) instead of retaining the raw executor-owned objects. Pre-F6a sequential code read and converted results immediately after each await; deferring routing until the whole step's pool settles let a later batch's mutation of a reused array/wrapper corrupt an earlier, already-settled batch's entry (reproducible at maxConcurrentBatches: 1 too). New test proves a mutating-buffer executor routes identically to a well-behaved one at conc=1 and conc=2. - unhandled-rejection guard: afterEach now awaits one macrotask (setImmediate) before asserting/clearing, since Node's unhandledRejection emission can land on a later tick than the offending test's own body; afterAll deregisters the listener so re-evaluation can't stack them. - concurrency.test.ts: replaced the two absolute wall-clock-bound tests with instrumented max-in-flight/dispatch-order assertions (deterministic, no timing assumption) plus a single self-calibrated relative wall-clock check (parallel < 0.6x serial for the same fixture, in the same test) — removes the absolute-deadline flake risk on loaded CI. --- src/__tests__/concurrency.test.ts | 161 +++++++++++++++----- src/__tests__/setup/unhandled-rejections.ts | 29 +++- src/core/pool.ts | 37 ++++- 3 files changed, 187 insertions(+), 40 deletions(-) diff --git a/src/__tests__/concurrency.test.ts b/src/__tests__/concurrency.test.ts index d1474c5..74d7e6b 100644 --- a/src/__tests__/concurrency.test.ts +++ b/src/__tests__/concurrency.test.ts @@ -70,54 +70,100 @@ function okExecutor(): StepExecutor { } // ───────────────────────────────────────────────────────────────────────── -// 1. Wall-clock — concurrency pool actually parallelizes within a step. +// 1. Concurrency observation — external review (P2): absolute wall-clock +// deadlines flake on loaded CI. Primary assertions observe concurrency +// directly (max in-flight, dispatch order) via an instrumented executor; +// exactly ONE relative wall-clock sanity check remains, self-calibrated +// against a same-test serial baseline instead of an absolute deadline. // ───────────────────────────────────────────────────────────────────────── -describe('wall-clock concurrency', () => { - it('600 calls / bs=100 / conc=5 completes in ~2 sequential-RT (two waves)', async () => { - const latencyMs = 50 - const invocationSizes: number[] = [] - const executor: StepExecutor = { - async executeMulticall(calls: StepCall[]): Promise { - invocationSizes.push(calls.length) - await sleep(latencyMs) - return calls.map((): RawResult => ({ status: 'success', value: 1 })) - }, - } +/** Parses the flat call index back out of a `makeCallsTask` key ('c123' -> 123). */ +function callIndex(key: string): number { + return Number(key.slice(1)) +} + +/** + * Executor that tracks observed in-flight concurrency and the order batches + * were CLAIMED (start of `executeMulticall`, before its artificial delay) — + * derives each call's original batch index from its flat call index and the + * `batchSize` the test is about to use, so it works for any fixture built + * from `makeCallsTask`. + */ +function makeConcurrencyProbeExecutor( + batchSize: number, + delayMs: number, +): { executor: StepExecutor; maxInFlight: () => number; dispatchOrder: number[] } { + let inFlight = 0 + let observedMax = 0 + const dispatchOrder: number[] = [] + + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + dispatchOrder.push(Math.floor(callIndex(calls[0]!.key) / batchSize)) + inFlight++ + observedMax = Math.max(observedMax, inFlight) + await sleep(delayMs) + inFlight-- + return calls.map((): RawResult => ({ status: 'success', value: 1 })) + }, + } + + return { executor, maxInFlight: () => observedMax, dispatchOrder } +} +describe('observed concurrency', () => { + it('conc=5: max observed in-flight is 5, batches claimed in ascending index order', async () => { + const { executor, maxInFlight, dispatchOrder } = makeConcurrencyProbeExecutor(100, 15) const task = makeCallsTask(600) - const start = performance.now() + await runMultistepTasks(executor, [task], { batchSize: 100, maxConcurrentBatches: 5 }) - const elapsed = performance.now() - start - expect(invocationSizes).toHaveLength(6) // 600 / 100 = 6 physical batches - // 6 batches at concurrency 5 -> 2 waves (5 + 1). Generous bounds against - // CI jitter: [1.5x, 3.5x] one batch's latency. - expect(elapsed).toBeGreaterThanOrEqual(latencyMs * 1.5) - expect(elapsed).toBeLessThanOrEqual(latencyMs * 3.5) - }, 20_000) + expect(maxInFlight()).toBe(5) + // 600/100 = 6 batches; the 6th is only claimed once one of the first 5 + // frees up, so it can only ever be dispatched last — the full order is + // deterministic regardless of which of the 5 finishes first. + expect(dispatchOrder).toEqual([0, 1, 2, 3, 4, 5]) + }) + + it('conc=1 (default): max observed in-flight is 1, batches claimed strictly in order', async () => { + const { executor, maxInFlight, dispatchOrder } = makeConcurrencyProbeExecutor(100, 5) + const task = makeCallsTask(600) + + await runMultistepTasks(executor, [task], { batchSize: 100, maxConcurrentBatches: 1 }) + + expect(maxInFlight()).toBe(1) + expect(dispatchOrder).toEqual([0, 1, 2, 3, 4, 5]) + }) - it('600 calls / bs=100 / conc=1 is genuinely sequential (~6 sequential-RT)', async () => { + it('relative wall-clock sanity: conc=5 is meaningfully faster than conc=1 for the same fixture (self-calibrated, no absolute deadline)', async () => { const latencyMs = 50 - const invocationSizes: number[] = [] - const executor: StepExecutor = { - async executeMulticall(calls: StepCall[]): Promise { - invocationSizes.push(calls.length) - await sleep(latencyMs) - return calls.map((): RawResult => ({ status: 'success', value: 1 })) - }, + function makeLatencyExecutor(): StepExecutor { + return { + async executeMulticall(calls: StepCall[]): Promise { + await sleep(latencyMs) + return calls.map((): RawResult => ({ status: 'success', value: 1 })) + }, + } } - const task = makeCallsTask(600) - const start = performance.now() - await runMultistepTasks(executor, [task], { batchSize: 100, maxConcurrentBatches: 1 }) - const elapsed = performance.now() - start + const serialStart = performance.now() + await runMultistepTasks(makeLatencyExecutor(), [makeCallsTask(600)], { + batchSize: 100, + maxConcurrentBatches: 1, + }) + const serialElapsed = performance.now() - serialStart + + const parallelStart = performance.now() + await runMultistepTasks(makeLatencyExecutor(), [makeCallsTask(600)], { + batchSize: 100, + maxConcurrentBatches: 5, + }) + const parallelElapsed = performance.now() - parallelStart - expect(invocationSizes).toHaveLength(6) - // Default (maxConcurrentBatches: 1) must be genuinely sequential: 6 - // batches back-to-back. Generous bounds: [5x, 8x]. - expect(elapsed).toBeGreaterThanOrEqual(latencyMs * 5) - expect(elapsed).toBeLessThanOrEqual(latencyMs * 8) + // 6 batches: serial ~6 waves, conc=5 ~2 waves — expect roughly a 3x + // speedup. Assert well under that (< 0.6x serial) so ordinary CI jitter + // on either measurement can never flip the comparison. + expect(parallelElapsed).toBeLessThan(serialElapsed * 0.6) }, 20_000) }) @@ -517,4 +563,45 @@ describe('runBatchPool (direct)', () => { expect(outcome.results[1]).toEqual([{ status: 'success', value: 'fast' }]) } }) + + it('snapshots results at settlement — an executor reusing the SAME array/wrapper objects across batches still routes correctly (P1)', async () => { + // Deliberately hostile: one array, one wrapper object, reused and + // MUTATED IN PLACE for every batch. `StepExecutor` never promised a + // fresh/immutable array per call — without the settlement-time clone in + // `runBatchPool`, a later batch's mutation would retroactively corrupt + // an earlier batch's already-"settled" entry once every batch shares + // the same underlying object. + function makeMutatingExecutor(): StepExecutor { + const sharedWrapper: { status: 'success'; value: unknown } = { status: 'success', value: undefined } + const sharedArray = [sharedWrapper] + return { + async executeMulticall(calls: StepCall[]): Promise { + await sleep(Math.random() * 10) + sharedWrapper.value = calls[0]!.key + return sharedArray + }, + } + } + + function makeWellBehavedExecutor(): StepExecutor { + return { + async executeMulticall(calls: StepCall[]): Promise { + await sleep(Math.random() * 10) + return [{ status: 'success', value: calls[0]!.key }] + }, + } + } + + for (const maxConcurrentBatches of [1, 2]) { + const wellBehaved = await runMultistepTasks(makeWellBehavedExecutor(), [makeCallsTask(6)], { + batchSize: 1, + maxConcurrentBatches, + }) + const mutating = await runMultistepTasks(makeMutatingExecutor(), [makeCallsTask(6)], { + batchSize: 1, + maxConcurrentBatches, + }) + expect(mutating).toEqual(wellBehaved) + } + }) }) diff --git a/src/__tests__/setup/unhandled-rejections.ts b/src/__tests__/setup/unhandled-rejections.ts index f9136fb..f2e24e9 100644 --- a/src/__tests__/setup/unhandled-rejections.ts +++ b/src/__tests__/setup/unhandled-rejections.ts @@ -14,9 +14,26 @@ * that true by construction (every batch promise is awaited directly inside * its own worker's try/catch), and this guard is what actually PROVES it * across the test suite rather than merely asserting it in the design doc. + * + * **Timing (external review, P2):** Node emits `unhandledRejection` only + * after the current microtask queue has fully drained with no handler + * attached — that emission can land on a LATER tick than the one in which + * the offending test's body returned. Asserting synchronously in `afterEach` + * races that emission and can miss a genuine leak. `afterEach` is async here + * and waits one macrotask (`setImmediate`) before checking — generous + * enough for the emission to have already landed, without needing fake + * timers or a fixed sleep. + * + * **Listener lifecycle (external review, P2):** the module-level + * `process.on(...)` runs once per evaluation of this setup file. Vitest + * normally isolates each test file's module registry, but nothing here + * should rely on that for correctness — an `afterAll` deregisters the exact + * listener this file added, so re-evaluating this module (watch mode, + * re-imports, etc.) can never stack listeners beyond the current file's + * lifetime. */ -import { afterEach } from 'vitest' +import { afterAll, afterEach } from 'vitest' const leaked: unknown[] = [] @@ -26,7 +43,15 @@ function onUnhandledRejection(reason: unknown): void { process.on('unhandledRejection', onUnhandledRejection) -afterEach(() => { +afterAll(() => { + process.removeListener('unhandledRejection', onUnhandledRejection) +}) + +afterEach(async () => { + // Give a pending 'unhandledRejection' emission one macrotask to land + // before we check — see the timing note above. + await new Promise((resolve) => setImmediate(resolve)) + if (leaked.length === 0) return const captured = leaked.slice() diff --git a/src/core/pool.ts b/src/core/pool.ts index 3a8cb73..59c5814 100644 --- a/src/core/pool.ts +++ b/src/core/pool.ts @@ -70,6 +70,24 @@ * inspected — see `runSteps` in `src/core/engine.ts`, which throws * immediately on a `cancelled` outcome instead of routing anything. * + * **Result snapshot at settlement (external review, P1):** `results[i]` is + * written as `batchResults.map((r) => ({ ...r }))` — a shallow clone of the + * array AND every wrapper object — not the raw value `execute()` resolved + * with. `StepExecutor` never promised its returned array/wrappers are fresh + * or immutable; the pre-F6a sequential code converted each batch's raw + * results into `StepResult`s synchronously right after that batch's own + * `await`, before the next batch's call even started, so an executor that + * reused and mutated the same array/wrapper across calls was already legal. + * Routing now happens once, after the WHOLE step's pool settles — without + * this clone, an executor mutating a reused wrapper for batch N+1 would + * retroactively corrupt batch N's already-"settled" result by the time + * routing reads `results[]`. The clone runs with no `await` between it and + * the `execute()` call that just resolved, so it captures exactly what that + * batch's own settlement saw — true for every `maxConcurrentBatches` value, + * including the default (1), where this bug would otherwise reproduce + * identically to the concurrent case (deferred routing, not concurrency + * itself, is what changed the timing). + * * **`maxConcurrentBatches: 1` is not special-cased.** `workerCount = min(1, * batchCount) = 1` naturally yields exactly one worker processing batches * strictly in order, one at a time — on rejection there is no other worker @@ -152,7 +170,24 @@ export async function runBatchPool( const batch = batches[batchIndex]! try { const batchResults = await execute(batch, batchIndex) - if (!cancelled) results[batchIndex] = batchResults + // Snapshot (array AND each wrapper) at the moment THIS batch's own + // await resolves — external review (P1): the `StepExecutor` + // contract never promised fresh/immutable arrays or wrappers, and + // the pre-F6a sequential code read (converted to `StepResult`) each + // batch's results synchronously right after its own await, before + // the next batch's call even started — so an executor reusing a + // result array/wrapper object across calls (mutating it in place + // for the next batch) was legal and safe. Under the pool, a + // batch's raw results now sit in `results[]` until the WHOLE + // step's pool settles (routing happens after `runBatchPool` + // returns) — without a copy here, a later batch's mutation of a + // reused array/wrapper would retroactively corrupt an earlier, + // already-settled batch's results by the time routing reads them. + // Cloning here — synchronously, with no `await` between this line + // and the `execute()` that just resolved — restores the exact + // same "read at the moment of settlement" semantics for every + // `maxConcurrentBatches` value, including the default (1). + if (!cancelled) results[batchIndex] = batchResults.map((r): RawResult => ({ ...r })) } catch (error) { cancelled = true terminalErrors.push({ batchIndex, callIndex: 0, error })