From c7d9ac7f9b499794deb94777a9d8cebc7d1d4950 Mon Sep 17 00:00:00 2001 From: halaprix Date: Thu, 23 Jul 2026 19:21:53 +0200 Subject: [PATCH 1/2] feat: adaptive bisection with attempts cap (F6b) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Bisect a rejecting physical batch (length > 1) into two halves and retry both through the pool's own central queue instead of failing the whole batch immediately. Bounded per original batch by maxBatchAttempts (default 2*ceil(log2(batchSize))+1); on exhaustion the unresolved calls become terminal with the last transport error. Splitting never awaits its own children — the permit is released before they're enqueued — so a pool of size 1 can drain an arbitrarily deep bisection tree without deadlocking. Off by default (adaptiveBatching: false) so existing run()/runSettled() behavior is unchanged byte-for-byte; enabling it lets run() rethrow the isolated raw transport error and runSettled() record per-call kind: 'batch' failures instead of coarsening the whole batch, while reverts carried in a resolved RawResult[] are never retried. --- README.md | 2 +- src/__tests__/bisection.test.ts | 470 ++++++++++++++++++++++++++++++++ src/core/engine.ts | 57 +++- src/core/internal.ts | 21 +- src/core/pool.ts | 398 ++++++++++++++++++--------- src/core/runMultistepTasks.ts | 81 ++++-- src/core/runSettled.ts | 120 +++++--- 7 files changed, 934 insertions(+), 215 deletions(-) create mode 100644 src/__tests__/bisection.test.ts diff --git a/README.md b/README.md index 157b3bb..40975d2 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-11.5KB-brightgreen)](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) [![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__/bisection.test.ts b/src/__tests__/bisection.test.ts new file mode 100644 index 0000000..eaf54fc --- /dev/null +++ b/src/__tests__/bisection.test.ts @@ -0,0 +1,470 @@ +import { describe, it, expect } from 'vitest' +import { runMultistepTasks } from '../core/runMultistepTasks' +import { runSettled } from '../core/runSettled' +import { DominoCallError } from '../core/errors' +import type { Address, MultistepTask, StepCall, StepResult, StepExecutor, RawResult } from '../core/types' + +/** + * F6b — adaptive bisection. See `src/core/pool.ts` (the central queue that + * owns splitting + attempts accounting) and `src/core/engine.ts` + * (`StepEnginePolicy.recordTerminal`, the hook that lets `runSettled` record + * a terminal batch failure instead of triggering `run`'s fail-fast + * cancellation). The global unhandled-rejection guard + * (`src/__tests__/setup/unhandled-rejections.ts`) fails any test here that + * leaks a rejection. + */ + +const ADDR = '0xA0b86991c6218b36c1d19D4a2e9Eb004C35d5Cc4' as Address + +function sleep(ms: number): Promise { + return new Promise((resolve) => setTimeout(resolve, ms)) +} + +function call(key: string): StepCall { + return { key, target: ADDR, abi: [], functionName: 'f' } +} + +/** Single-step task with `n` calls keyed `c0..c(n-1)`. */ +function makeCallsTask(n: number): 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 + }, + finalize() { + return captured + }, + } +} + +/** + * An executor that rejects any physical batch CONTAINING a poisoned key (by + * exact key match against `calls`), otherwise resolves every call as + * `success`. Tracks total invocation count and every batch (as key arrays) + * it was actually called with — used to assert the executor was called + * exactly the expected number of times and never with a batch that could not + * possibly need re-execution. + */ +function makePoisonExecutor( + poisonedKeys: Set, + opts?: { delay?: () => number }, +): { executor: StepExecutor; callCount: () => number; invocations: () => string[][] } { + let count = 0 + const invocations: string[][] = [] + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + count++ + const keys = calls.map((c) => c.key) + invocations.push(keys) + if (opts?.delay) await sleep(opts.delay()) + if (keys.some((k) => poisonedKeys.has(k))) { + throw new Error(`transport failure containing poisoned key(s): ${keys.filter((k) => poisonedKeys.has(k)).join(',')}`) + } + return calls.map((c): RawResult => ({ status: 'success', value: 'ok-' + c.key })) + }, + } + return { executor, callCount: () => count, invocations: () => invocations } +} + +function findResult(results: StepResult[], key: string): StepResult | undefined { + return results.find((r) => r.key === key) +} + +// ───────────────────────────────────────────────────────────────────────── +// 1. Isolation — 100 calls, 1 poisoned, runSettled, adaptive on. +// ───────────────────────────────────────────────────────────────────────── + +describe('bisection isolation (runSettled)', () => { + it('99 succeed, 1 kind:"batch" failure with cause identity = the last transport error; call count within default cap', async () => { + const poisoned = 'c42' + const { executor, callCount } = makePoisonExecutor(new Set([poisoned])) + const task = makeCallsTask(100) + + const [settled] = await runSettled(executor, [task], { + batchSize: 100, + maxConcurrentBatches: 1, + adaptiveBatching: true, + }) + + expect(settled!.status).toBe('fulfilled') + const results = (settled as { status: 'fulfilled'; value: StepResult[] }).value + expect(results).toHaveLength(100) + + const successes = results.filter((r) => r.status === 'success') + const failures = results.filter((r) => r.status === 'failure') + expect(successes).toHaveLength(99) + expect(failures).toHaveLength(1) + expect(failures[0]!.key).toBe(poisoned) + + const err = (failures[0] as { status: 'failure'; error?: unknown }).error + expect(err).toBeInstanceOf(DominoCallError) + expect((err as DominoCallError).kind).toBe('batch') + expect((err as DominoCallError).cause).toBeInstanceOf(Error) + expect(((err as DominoCallError).cause as Error).message).toContain(poisoned) + + // default maxBatchAttempts = 2*ceil(log2(100))+1 = 15; single-fault + // isolation should land close to 1 + 2*ceil(log2(100)) and never exceed + // the cap. + const defaultCap = 2 * Math.ceil(Math.log2(100)) + 1 + expect(callCount()).toBeLessThanOrEqual(defaultCap) + expect(callCount()).toBeGreaterThan(1) // must have actually split at least once + + // Every non-poisoned call kept its correct value. + for (const r of successes) { + expect((r as { status: 'success'; value: unknown }).value).toBe('ok-' + r.key) + } + }) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 2. run() single-poisoned — thrown error identity, shuffled runs. +// ───────────────────────────────────────────────────────────────────────── + +describe('bisection identity (run)', () => { + it('throws the exact transport error from the final length-1 execution; stable over >=100 shuffled runs', async () => { + const poisoned = 'c17' + + for (let i = 0; i < 100; i++) { + let poisonError: Error | undefined + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + await sleep(Math.random() * 5) + const keys = calls.map((c) => c.key) + if (keys.includes(poisoned)) { + const err = new Error('poison-' + poisoned) + if (keys.length === 1) poisonError = err + throw err + } + return calls.map((): RawResult => ({ status: 'success', value: 1 })) + }, + } + + const task = makeCallsTask(20) + const conc = (i % 4) + 1 // cycle 1..4 + + let thrown: unknown + try { + await runMultistepTasks(executor, [task], { + batchSize: 20, + maxConcurrentBatches: conc, + adaptiveBatching: true, + }) + throw new Error('expected run() to reject') + } catch (err) { + thrown = err + } + + expect(poisonError).toBeDefined() + expect(thrown).toBe(poisonError) + } + }, 30_000) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 3. Attempts cap. +// ───────────────────────────────────────────────────────────────────────── + +describe('attempts cap', () => { + it('maxBatchAttempts: 3 with a 100-call poisoned batch -> <=3 executions; a coarse (>1) unresolved group fails kind "batch"; resolved calls keep exact values', async () => { + const poisoned = 'c42' + const { executor, callCount } = makePoisonExecutor(new Set([poisoned])) + const task = makeCallsTask(100) + + const [settled] = await runSettled(executor, [task], { + batchSize: 100, + maxConcurrentBatches: 1, + adaptiveBatching: true, + maxBatchAttempts: 3, + }) + + expect(callCount()).toBeLessThanOrEqual(3) + + const results = (settled as { status: 'fulfilled'; value: StepResult[] }).value + const failures = results.filter((r) => r.status === 'failure') + const successes = results.filter((r) => r.status === 'success') + + // With only 3 executions available, full isolation to a single call is + // impossible for a 100-call batch -> a coarse group (>1 call) remains + // unresolved, all as kind:'batch' failures. + expect(failures.length).toBeGreaterThan(1) + for (const f of failures) { + const err = (f as { status: 'failure'; error?: unknown }).error + expect(err).toBeInstanceOf(DominoCallError) + expect((err as DominoCallError).kind).toBe('batch') + } + // The poisoned call itself is always among the unresolved. + expect(failures.some((f) => f.key === poisoned)).toBe(true) + + // Every resolved call kept its correct value (never wrong data). + for (const r of successes) { + expect((r as { status: 'success'; value: unknown }).value).toBe('ok-' + r.key) + } + expect(successes.length + failures.length).toBe(100) + }) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 4. Deadlock. +// ───────────────────────────────────────────────────────────────────────── + +describe('no pool deadlock', () => { + it('pool of 1, poisoned batch of 100, adaptive on -> terminates with full isolation', async () => { + const poisoned = 'c7' + const { executor } = makePoisonExecutor(new Set([poisoned])) + const task = makeCallsTask(100) + + const [settled] = await runSettled(executor, [task], { + batchSize: 100, + maxConcurrentBatches: 1, + adaptiveBatching: true, + }) + + const results = (settled as { status: 'fulfilled'; value: StepResult[] }).value + const failures = results.filter((r) => r.status === 'failure') + expect(failures).toHaveLength(1) + expect(failures[0]!.key).toBe(poisoned) + expect(results.filter((r) => r.status === 'success')).toHaveLength(99) + }) + + it('conc=2, two original batches, one bisects -> the other original batch fully succeeds', async () => { + const poisoned = 'poison-x' + const poisonedBatchCalls = Array.from({ length: 50 }, (_, i) => call('p' + i)) + // Replace one call's key with the poisoned marker so the executor can + // recognize it. + poisonedBatchCalls[13] = call(poisoned) + const healthyBatchCalls = Array.from({ length: 50 }, (_, i) => call('h' + i)) + + const invocationKeys: string[][] = [] + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + const keys = calls.map((c) => c.key) + invocationKeys.push(keys) + await sleep(Math.random() * 5) + if (keys.includes(poisoned)) throw new Error('poison') + return calls.map((c): RawResult => ({ status: 'success', value: 'ok-' + c.key })) + }, + } + + let capturedA: StepResult[] = [] + let capturedB: StepResult[] = [] + const taskA: MultistepTask = { + maxStep: 1, + buildStepCalls(step) { + return step === 1 ? poisonedBatchCalls : [] + }, + consumeStepResults(_step, results) { + capturedA = results + }, + finalize() { + return null + }, + } + const taskB: MultistepTask = { + maxStep: 1, + buildStepCalls(step) { + return step === 1 ? healthyBatchCalls : [] + }, + consumeStepResults(_step, results) { + capturedB = results + }, + finalize() { + return null + }, + } + + const [resultA, resultB] = await runSettled(executor, [taskA, taskB], { + batchSize: 50, + maxConcurrentBatches: 2, + adaptiveBatching: true, + }) + + expect(resultA!.status).toBe('fulfilled') + expect(resultB!.status).toBe('fulfilled') + + // Batch B (healthy, separate original batch) — every call succeeded. + expect(capturedB).toHaveLength(50) + expect(capturedB.every((r) => r.status === 'success')).toBe(true) + for (const c of healthyBatchCalls) { + expect(findResult(capturedB, c.key)).toEqual({ status: 'success', key: c.key, value: 'ok-' + c.key }) + } + + // Batch A isolated its single poisoned call. + expect(capturedA.filter((r) => r.status === 'failure')).toHaveLength(1) + expect(findResult(capturedA, poisoned)?.status).toBe('failure') + }) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 5. Reverts are per-call, never retried as transport failures. +// ───────────────────────────────────────────────────────────────────────── + +describe('reverts not retried', () => { + it('a resolving multicall with per-call (allowFailure-style) failures -> executor called once, zero splits', async () => { + let callCount = 0 + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + callCount++ + return calls.map((c, i): RawResult => + i === 3 ? { status: 'failure', error: new Error('reverted: ' + c.key) } : { status: 'success', value: 'ok-' + c.key }, + ) + }, + } + const task = makeCallsTask(10) + + const results = await runMultistepTasks(executor, [task], { + batchSize: 10, + maxConcurrentBatches: 1, + adaptiveBatching: true, + }) + + expect(callCount).toBe(1) + const [taskResults] = results + expect(taskResults!.filter((r) => r.status === 'failure')).toHaveLength(1) + expect(taskResults!.filter((r) => r.status === 'success')).toHaveLength(9) + expect(findResult(taskResults!, 'c3')?.status).toBe('failure') + }) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 6. Multi-poisoned (run) — some error thrown, no unhandled rejections. +// ───────────────────────────────────────────────────────────────────────── + +describe('multi-poisoned (run)', () => { + it('2 poisoned calls in different halves, adaptive on -> some error thrown; no unhandled rejections', async () => { + for (let i = 0; i < 10; i++) { + const poisonedA = 'c3' + const poisonedB = 'c14' + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + await sleep(Math.random() * 5) + const keys = calls.map((c) => c.key) + if (keys.includes(poisonedA) || keys.includes(poisonedB)) { + throw new Error('poison') + } + return calls.map((): RawResult => ({ status: 'success', value: 1 })) + }, + } + const task = makeCallsTask(20) + + let rejected = false + let thrown: unknown + try { + await runMultistepTasks(executor, [task], { + batchSize: 20, + maxConcurrentBatches: 3, + adaptiveBatching: true, + }) + } catch (err) { + rejected = true + thrown = err + } + expect(rejected).toBe(true) + expect(thrown).toBeInstanceOf(Error) + } + }, 20_000) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 7. Adaptive off — T14 behavior exactly, zero splits. +// ───────────────────────────────────────────────────────────────────────── + +describe('adaptive off (default)', () => { + it('a length>1 transport rejection triggers NO splits — executor called exactly once for that batch (run)', async () => { + let callCount = 0 + const boom = new Error('boom') + const executor: StepExecutor = { + async executeMulticall(): Promise { + callCount++ + throw boom + }, + } + const task = makeCallsTask(10) + + await expect(runMultistepTasks(executor, [task], { batchSize: 10, maxConcurrentBatches: 1 })).rejects.toBe(boom) + expect(callCount).toBe(1) + }) + + it('a length>1 transport rejection triggers NO splits — T14 batch-failure shape exactly (runSettled)', async () => { + let callCount = 0 + const boom = new Error('boom') + const executor: StepExecutor = { + async executeMulticall(): Promise { + callCount++ + throw boom + }, + } + const task = makeCallsTask(10) + + const [settled] = await runSettled(executor, [task], { batchSize: 10, maxConcurrentBatches: 1 }) + expect(callCount).toBe(1) + const results = (settled as { status: 'fulfilled'; value: StepResult[] }).value + expect(results).toHaveLength(10) + for (const r of results) { + expect(r.status).toBe('failure') + const err = (r as { status: 'failure'; error?: unknown }).error + expect(err).toBeInstanceOf(DominoCallError) + expect((err as DominoCallError).kind).toBe('batch') + expect((err as DominoCallError).cause).toBe(boom) + } + }) + + it('explicit adaptiveBatching: false behaves identically to omitted (still no splits)', async () => { + let callCount = 0 + const boom = new Error('boom') + const executor: StepExecutor = { + async executeMulticall(): Promise { + callCount++ + throw boom + }, + } + const task = makeCallsTask(10) + + await expect( + runMultistepTasks(executor, [task], { batchSize: 10, maxConcurrentBatches: 1, adaptiveBatching: false }), + ).rejects.toBe(boom) + expect(callCount).toBe(1) + }) +}) + +// ───────────────────────────────────────────────────────────────────────── +// 8. Length-mismatch (programmer error) — still aborts, never retried. +// ───────────────────────────────────────────────────────────────────────── + +describe('length-mismatch is never retried, even under adaptive', () => { + it('run(): aborts immediately, executor called exactly once', async () => { + let callCount = 0 + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + callCount++ + return calls.slice(0, calls.length - 1).map((): RawResult => ({ status: 'success', value: 1 })) + }, + } + const task = makeCallsTask(10) + + await expect( + runMultistepTasks(executor, [task], { batchSize: 10, maxConcurrentBatches: 1, adaptiveBatching: true }), + ).rejects.toThrow('length mismatch') + expect(callCount).toBe(1) + }) + + it('runSettled(): aborts the entire run (rejects), executor called exactly once', async () => { + let callCount = 0 + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + callCount++ + return calls.slice(0, calls.length - 1).map((): RawResult => ({ status: 'success', value: 1 })) + }, + } + const task = makeCallsTask(10) + + await expect( + runSettled(executor, [task], { batchSize: 10, maxConcurrentBatches: 1, adaptiveBatching: true }), + ).rejects.toThrow('length mismatch') + expect(callCount).toBe(1) + }) +}) diff --git a/src/core/engine.ts b/src/core/engine.ts index 9f85fca..82d2e49 100644 --- a/src/core/engine.ts +++ b/src/core/engine.ts @@ -7,10 +7,13 @@ * 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 + * (record-and-continue) is captured in the `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. + * F6b adds one more (optional) hook, `recordTerminal` — see its doc comment + * below and `src/core/pool.ts`'s "Terminal policy hook" section for the full + * adaptive-bisection design. * * What stays OUTSIDE this module (deliberately): the F2 consumption * pipeline (`prepareRun`/`resolvePinnedBlock`, `internal.ts`) and the final @@ -22,9 +25,10 @@ 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`. + * The hooks that fully capture `run` vs `runSettled`'s divergence — 3 as of + * F6a, plus the optional `recordTerminal` added by F6b. 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 { /** @@ -43,14 +47,19 @@ export interface StepEnginePolicy { /** * 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: + * behavior matters for cancellation/bisection: * - `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). + * — a rejection here is exactly what feeds the pool's bisection (F6b) + * and, once terminal, its fail-fast cancellation. The length-mismatch + * guard now lives in the pool itself (`src/core/pool.ts`), not here — + * it applies uniformly to whole batches AND bisected sub-batches. + * - `runSettled`'s hook, with adaptive bisection OFF (default), catches + * a transport rejection and *resolves* with synthesized per-call + * `kind:'batch'` failures instead — this is what makes `runSettled` + * never cancel on an ordinary transport failure. With adaptive ON, it + * deliberately lets the rejection propagate to the pool instead, so + * bisection can split it — the per-call synthesis then happens once, + * at the TERMINAL point, via `recordTerminal` below. */ executeBatch(batch: StepCall[]): Promise @@ -60,12 +69,28 @@ export interface StepEnginePolicy { * and marks the task dead. */ consumeStepResults(taskIndex: number, step: number, results: StepResult[]): void + + /** + * Adaptive bisection's terminal hook (F6b) — see + * `BisectionPolicy.recordTerminal` in `src/core/pool.ts` for the full + * contract. Present only for `runSettled` (never cancel on a terminal + * batch failure — synthesize a per-call `kind:'batch'` `DominoCallError` + * instead). Absent for `run`, whose terminal handling is the pool's plain + * fail-fast cancellation, unchanged from T14 in shape. + */ + recordTerminal?(calls: StepCall[], error: unknown): RawResult[] } /** The subset of `BatchOptions` the engine itself needs. */ export interface StepEngineOptions { batchSize: number maxConcurrentBatches: number + /** F6b — see `BatchOptions.adaptiveBatching`. */ + adaptiveBatching: boolean + /** F6b — see `BatchOptions.maxBatchAttempts`. Consulted by the pool only + * when `adaptiveBatching` is true; otherwise any rejection is terminal + * on its first occurrence regardless of this value (T14 behavior). */ + maxBatchAttempts: number } /** @@ -107,9 +132,11 @@ export async function runSteps( batches.push(calls.slice(batchStart, batchStart + options.batchSize)) } - const outcome = await runBatchPool(batches, options.maxConcurrentBatches, (batch) => - policy.executeBatch(batch), - ) + const outcome = await runBatchPool(batches, options.maxConcurrentBatches, (batch) => policy.executeBatch(batch), { + adaptive: options.adaptiveBatching, + maxBatchAttempts: options.maxBatchAttempts, + ...(policy.recordTerminal ? { recordTerminal: policy.recordTerminal } : {}), + }) if (outcome.outcome === 'cancelled') { // Spec (d): in-flight results are discarded — we never look at diff --git a/src/core/internal.ts b/src/core/internal.ts index d87c1e4..e5d4d24 100644 --- a/src/core/internal.ts +++ b/src/core/internal.ts @@ -69,16 +69,18 @@ const consumed = new WeakSet() type Branded = MultistepTask & SingleUseCarrier /** - * 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. + * Numeric (+ one boolean) options validated + defaulted by `validateOptions` + * (F6a/F6b). All four 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. */ export interface ValidatedRunOptions { batchSize: number maxConcurrentBatches: number maxBatchAttempts: number + adaptiveBatching: boolean } /** Options `validateOptions` reads — a structural subset of `BatchOptions`. */ @@ -86,6 +88,7 @@ export interface NumericOptionsInput { batchSize?: number maxConcurrentBatches?: number maxBatchAttempts?: number + adaptiveBatching?: boolean } /** @@ -125,7 +128,13 @@ export function validateOptions(o: NumericOptionsInput | undefined): ValidatedRu const maxBatchAttempts = o?.maxBatchAttempts ?? defaultMaxBatchAttempts validatePositiveSafeInteger('maxBatchAttempts', maxBatchAttempts) - return { batchSize, maxConcurrentBatches, maxBatchAttempts } + // Not a positive-safe-integer field — a plain boolean flag (F6b), default + // `false` (see `BatchOptions.adaptiveBatching`'s doc comment for why: rate + // limiting makes bisection's retry amplification actively harmful unless a + // caller has opted in with knowledge of their transport's failure modes). + const adaptiveBatching = o?.adaptiveBatching ?? false + + return { batchSize, maxConcurrentBatches, maxBatchAttempts, adaptiveBatching } } /** diff --git a/src/core/pool.ts b/src/core/pool.ts index 59c5814..610e721 100644 --- a/src/core/pool.ts +++ b/src/core/pool.ts @@ -1,107 +1,153 @@ /** - * Concurrency-limited batch dispatch pool + fail-fast cancellation (F6a). + * Concurrency-limited batch dispatch pool + fail-fast cancellation (F6a), + * extended with adaptive bisection (F6b) — split-and-retry on a + * batch-level transport rejection, bounded by a per-original-batch attempts + * cap, coordinated entirely through this module's own central queue. * * `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 + * area adds. It knows nothing about tasks, steps, or which runner (`run` vs + * `runSettled`) is calling it — only "batches" (`StepCall[]`), an `execute` + * callback that resolves with `RawResult[]` or rejects, and (F6b) an + * optional `BisectionPolicy` describing whether/how hard to retry a + * rejection and what a caller wants done with a TERMINAL one. See + * `src/core/engine.ts` for how the two runners' policies feed both of these. + * + * ## The queue (F6b) + * + * Work items are `{ calls, origBatchIndex, origOffset }` — `origBatchIndex` + * is which of the ORIGINAL `batches` this item's calls ultimately belong to; + * `origOffset` is where in that original batch's call list this item's + * calls start. The queue is seeded with exactly one whole-batch item per + * original batch (`origOffset: 0`) — this degrades to T14's flat batch list + * whenever nothing is ever split. + * + * Workers claim by `nextIndex++` exactly as T14 did for `batches` — reading + * `queue.length` fresh on every loop check (T14's "forward-compat" note + * anticipated exactly this: a future queue that grows mid-run needs no + * dedicated wake-up, because the length check already re-reads live). On a + * retryable rejection (`adaptive` && the item has more than one call), the + * worker `queue.push()`es the two half-sized children and immediately loops + * to claim whatever's next — **it never awaits its own children.** + * + * **This is the whole no-deadlock argument.** A "permit" here is nothing + * more than "being one of the `maxConcurrentBatches` running `worker()` + * loops" — there is no separate semaphore object to release. Splitting is a + * synchronous `queue.push()` followed immediately by `continue`; the worker + * moves on to claim the NEXT queued item (which may be one of the children + * it just pushed, or may belong to an entirely different original batch — + * whichever is lowest-indexed and unclaimed). A pool of size 1 therefore + * just serially drains however many splits the poisoned call(s) need, one + * `execute()` at a time, with nothing ever waiting on anything the SAME + * worker itself hasn't already finished. There is no recursive await for a + * deadlock to hide in. + * + * ## Attempts accounting (F6b) + * + * `attempts`/`lastError`, both indexed by `origBatchIndex` (never by queue + * position — a split child shares its parent's `origBatchIndex`). The cap + * is enforced in exactly one place: at CLAIM time, before every execution — + * the original batch's own first one included, and every bisected child. + * `attempts[i] >= maxBatchAttempts` means "do not execute" — the item goes + * straight to TERMINAL using `lastError[i]` (the most recent transport + * rejection seen for that original, which may have come from a *sibling* + * item, not this one — this one was never executed at all). Otherwise + * `attempts[i]++` then execute. + * + * Because the cap is checked at claim time rather than at split-decision + * time, a rejection always unconditionally pushes both halves when + * `adaptive && calls.length > 1` — the cap decides, independently and per + * child, whether each one actually gets to run. This is what produces the + * spec's "coarse group" behavior for free when the cap runs out mid-tree: + * whichever children can't get a claim-time attempt slot become terminal + * together, with no special-casing for "the cap ran out partway through a + * split". + * + * ## Terminal policy hook (F6b) + * + * `BisectionPolicy.recordTerminal` — present only for a runner that never + * wants the whole pool cancelled on a terminal item (`runSettled`): + * + * - **Absent (`run`):** a terminal item triggers T14's fail-fast + * cancellation verbatim — see (a)-(d) below, now generalized so + * `callIndex` is a REAL offset (`origOffset`) instead of always `0`. + * `compareTerminal`'s 2-key ordering needed no change to support this — + * T14's doc comment called this out as the exact reason it was already + * shaped as a 2-key comparator. + * - **Present (`runSettled`):** called once per terminal item with that + * item's calls and the last transport error for its original batch; + * returns one `RawResult` per call (same order). The pool writes that + * array into the shared per-original results row at + * `[origOffset, origOffset + calls.length)` and does NOT cancel — other + * items (siblings, other originals) keep going exactly like a normal + * success. + * + * Only two runners exist, so "hook present vs absent" is a sufficient policy + * signal; there is no separate cancel/never-cancel enum. + * + * Bisection retries ONLY a rejection from `execute()` itself — it can never + * see or touch a per-call revert/failure carried inside a RESOLVED + * `RawResult[]` (`allowFailure`-style); the entire retry/terminal machinery + * lives inside the `catch` block below, which a resolution never reaches. + * + * ## Length-mismatch (programmer error) — never retried, both runners abort + * + * Checked directly here, right after ANY `execute()` resolves (whole batch + * or bisected child): `batchResults.length !== calls.length`. This is a + * plain post-resolve `if`, not something that arrives through the `catch` + * block — so it can never be mistaken for a retryable transport rejection, + * needs no marker-error/`instanceof` trick to distinguish the two, and + * unconditionally takes the same `cancelWith(...)` path T14 used for every + * rejection, ignoring `recordTerminal` entirely. That reproduces "aborts + * everything, for both runners" for free — `runSettled` still aborts + * wholesale on this, exactly as it did before bisection existed. + * + * ## Cancellation policy — spec (a)-(d), generalized from T14 + * + * (a) **Queued items are not dispatched.** Every worker checks the shared + * `cancelled` flag BEFORE claiming its next item. 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. - * - * **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 - * 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. + * two workers can never claim the same item, and once `cancelled` + * flips, no worker claims a further one. + * + * (b) **In-flight executions are allowed to settle; no unhandled + * rejections.** A worker that has already claimed an item keeps + * `await`-ing it to completion regardless of `cancelled`. Because every + * item's promise is awaited directly inside its own worker's + * `try/catch` — never handed off un-awaited, even when that worker's + * own catch block goes on to `queue.push()` more work — there is never + * a floating rejection for the runtime to report as + * `unhandledRejection`. + * + * (c) **Thrown error = lowest (origBatchIndex, origOffset) among discovered + * terminal errors.** Every worker that reaches a real terminal (no + * `recordTerminal`) pushes `{ batchIndex: origBatchIndex, callIndex: + * origOffset, error }`; once all workers finish, the lowest entry (by + * `compareTerminal`) is returned. For exactly one poisoned call this is + * deterministic regardless of timing/concurrency — bisection always + * eventually isolates it (or the cap forces a coarser but still unique + * terminal group) and it always settles. With two or more + * concurrently-failing regions, which one's error is thrown remains + * explicitly timing-dependent, same as T14. + * + * (d) **Results of in-flight items are discarded on cancel.** A write into + * the shared `results` array only happens `if (!cancelled)` at the + * moment an item settles — see `runSteps` in `src/core/engine.ts`, which + * throws immediately on a `cancelled` outcome instead of routing + * anything. + * + * **Result snapshot at settlement (T14, P1, preserved):** every write into + * `results` clones the array AND every wrapper (`.map(r => ({...r}))`) + * synchronously, with no `await` between it and the call that produced the + * values — see T14's original reasoning, unchanged by bisection. + * + * **Backward-compatible signature:** `runBatchPool`'s first three parameters + * are byte-identical to T14's (`batches`, `maxConcurrentBatches`, `execute`) + * — `bisection` is a new, optional 4th parameter. Omitting it (or passing + * `adaptive: false`) makes `canRetry` always false, so every rejection goes + * straight to terminal on its FIRST occurrence — T14's exact behavior. This + * is what lets every direct `runBatchPool(...)` call already in + * `concurrency.test.ts` keep working, completely unmodified. */ import type { StepCall, RawResult } from './types' @@ -129,21 +175,49 @@ interface TerminalCandidate { /** * 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. + * original batch index, then lowest call offset within that batch. Every + * T14 terminal was whole-batch (`callIndex` always `0`), so this used to + * degrade to plain `batchIndex` ordering; bisection (F6b) now feeds real + * per-call offsets through the same 2-key comparator with no changes here. */ function compareTerminal(a: TerminalCandidate, b: TerminalCandidate): number { if (a.batchIndex !== b.batchIndex) return a.batchIndex - b.batchIndex return a.callIndex - b.callIndex } +/** One unit of queued work — see the module doc's "The queue" section. */ +interface WorkItem { + calls: StepCall[] + /** Which of the ORIGINAL `batches` these calls belong to. */ + origBatchIndex: number + /** Offset of `calls[0]` within that original batch's full call list. */ + origOffset: number +} + +/** + * Adaptive-bisection configuration (F6b) — optional 4th argument to + * `runBatchPool`. See the module doc's "Attempts accounting" and "Terminal + * policy hook" sections. + */ +export interface BisectionPolicy { + /** Enable split-and-retry on a batch-level rejection. `false` (or the + * whole `bisection` argument omitted) reproduces T14 exactly: any + * rejection is terminal on its first occurrence, `maxBatchAttempts` is + * never consulted. */ + adaptive: boolean + /** Total executions allowed per ORIGINAL batch — enforced at claim time, + * covering the original's own first execution and every bisected child. */ + maxBatchAttempts: number + /** Present only for a runner that never wants the whole pool cancelled on + * a terminal item (`runSettled`) — see the module doc's "Terminal policy + * hook" section. Absent means "cancel on terminal", exactly like T14. */ + recordTerminal?: (calls: StepCall[], error: unknown) => RawResult[] +} + /** * 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. + * in flight at once, applying adaptive bisection (if `bisection?.adaptive`) + * and fail-fast cancellation ((a)-(d) above) per the module doc comment. * * `maxConcurrentBatches` is clamped to `batches.length` — requesting more * workers than there are batches would just spin up idle workers that exit @@ -153,44 +227,106 @@ export async function runBatchPool( batches: StepCall[][], maxConcurrentBatches: number, execute: (batch: StepCall[], batchIndex: number) => Promise, + bisection?: BisectionPolicy, ): Promise { if (batches.length === 0) return { outcome: 'completed', results: [] } - const results: RawResult[][] = new Array(batches.length) + const adaptive = bisection?.adaptive ?? false + const maxBatchAttempts = bisection?.maxBatchAttempts ?? Infinity + const recordTerminal = bisection?.recordTerminal + + // Pre-sized PER ORIGINAL BATCH so bisected sub-batches can write into + // their exact offset range. A never-split item covers the whole row + // (offset 0, length === batches[i].length) — the same single write T14 + // did, in the same shape. + const results: RawResult[][] = batches.map((b) => new Array(b.length)) + + const queue: WorkItem[] = batches.map((calls, origBatchIndex) => ({ calls, origBatchIndex, origOffset: 0 })) + + const attempts: number[] = new Array(batches.length).fill(0) + const lastError: unknown[] = new Array(batches.length) + const terminalErrors: TerminalCandidate[] = [] let nextIndex = 0 let cancelled = false + /** T14's fail-fast cancellation, generalized to a real (origBatchIndex, + * origOffset) pair instead of the always-0 callIndex T14 had. Used both + * for a genuine terminal (no `recordTerminal`) and, unconditionally, for + * the length-mismatch programmer-error case regardless of policy. */ + function cancelWith(origBatchIndex: number, origOffset: number, error: unknown): void { + cancelled = true + terminalErrors.push({ batchIndex: origBatchIndex, callIndex: origOffset, error }) + } + + function writeResults(item: WorkItem, batchResults: RawResult[]): void { + if (cancelled) return + // Snapshot (array AND each wrapper) at the moment of settlement — see + // the module doc's "Result snapshot at settlement" note. + const snapshot = batchResults.map((r): RawResult => ({ ...r })) + const row = results[item.origBatchIndex]! + for (let i = 0; i < snapshot.length; i++) { + row[item.origOffset + i] = snapshot[i]! + } + } + + /** An item became terminal: a length-1 rejection, or its original batch's + * attempts cap exhausted before this item could even execute. */ + function handleTerminal(item: WorkItem, error: unknown): void { + if (recordTerminal) { + writeResults(item, recordTerminal(item.calls, error)) + return + } + cancelWith(item.origBatchIndex, item.origOffset, error) + } + async function worker(): Promise { for (;;) { if (cancelled) return - if (nextIndex >= batches.length) return - const batchIndex = nextIndex + if (nextIndex >= queue.length) return + const item = queue[nextIndex]! nextIndex++ - const batch = batches[batchIndex]! + const { calls, origBatchIndex, origOffset } = item + + // Attempts cap (F6b) — checked at CLAIM time, before executing. + if (attempts[origBatchIndex]! >= maxBatchAttempts) { + handleTerminal(item, lastError[origBatchIndex]) + continue + } + attempts[origBatchIndex]!++ + try { - const batchResults = await execute(batch, batchIndex) - // 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 })) + const batchResults = await execute(calls, origBatchIndex) + + if (batchResults.length !== calls.length) { + // Programmer-error path (executor-implementation bug): never + // retried, never counted as a bisection-eligible transport + // failure, aborts the run for BOTH runners regardless of + // `recordTerminal` — see the module doc comment. + cancelWith( + origBatchIndex, + origOffset, + new Error( + `StepExecutor returned ${batchResults.length} results for ${calls.length} calls — length mismatch`, + ), + ) + continue + } + + writeResults(item, batchResults) } catch (error) { - cancelled = true - terminalErrors.push({ batchIndex, callIndex: 0, error }) + lastError[origBatchIndex] = error + + if (adaptive && calls.length > 1) { + // Bisection: split and retry both halves through this SAME + // central queue — no recursive await. See the module doc's "The + // queue" section for why this can never deadlock. + const mid = Math.ceil(calls.length / 2) + queue.push({ calls: calls.slice(0, mid), origBatchIndex, origOffset }) + queue.push({ calls: calls.slice(mid), origBatchIndex, origOffset: origOffset + mid }) + } else { + handleTerminal(item, error) + } } } } diff --git a/src/core/runMultistepTasks.ts b/src/core/runMultistepTasks.ts index 6272203..d9703a5 100644 --- a/src/core/runMultistepTasks.ts +++ b/src/core/runMultistepTasks.ts @@ -51,16 +51,50 @@ export interface BatchOptions { /** * Maximum retry attempts for a failing physical batch before it is - * treated as a terminal failure (adaptive bisection, F6b). Default: + * treated as a terminal failure (adaptive bisection, F6b) — counted PER + * ORIGINAL batch: every execution of that original batch or any of its + * bisected sub-batches counts, successful or not, including the + * original's own first execution. Default: * `2 * Math.ceil(Math.log2(batchSize)) + 1`. Must be a positive safe - * integer ≥ 1 — anything else throws. + * integer ≥ 1 — anything else throws. Only consulted when + * `adaptiveBatching` is `true`. * - * 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. + * Fully isolating every bad call in a batch can require up to `2N − 1` + * executions (`N` = batch size) in the worst case (multiple bad calls + * spread across the batch). The default cap is sized for the common + * single-bad-call case; a batch with multiple failing calls may exhaust + * it before every one is individually isolated — under `runSettled` this + * produces one or more coarser-grained `kind: 'batch'` failures instead of + * fully isolated ones, never wrong data (unaffected calls always keep + * their real values). */ maxBatchAttempts?: number + /** + * Enable adaptive bisection (F6b): when a physical batch's executor call + * rejects (a transport/RPC-level failure — never a per-call revert, which + * is already isolated via `allowFailure`) and the batch has more than one + * call, split it in half and retry both halves independently instead of + * treating the whole batch as failed. Recurses until either every call is + * isolated to its own single-call execution, or `maxBatchAttempts` is + * exhausted for that original batch. Default: `false`. + * + * **Off by default.** Bisection cannot distinguish a genuine per-call + * failure (e.g. an out-of-gas call poisoning its whole batch) from a + * transient transport problem like HTTP 429 rate-limiting — under rate + * limiting, retrying amplifies load into an already-throttled endpoint by + * up to `2N − 1` calls for a single original batch. Enable this once you + * know your transport/RPC failures are dominated by genuinely bad calls, + * not rate limits; leave it off (or handle 429s at the transport layer) + * otherwise. A future failure-cause classification may allow enabling + * this safely by default. + * + * With this off, `run` and `runSettled` are both byte-for-byte identical + * to their pre-F6b (1.1) behavior — a batch rejection is terminal + * immediately, `maxBatchAttempts` is never consulted. + */ + adaptiveBatching?: boolean + /** Block to query at (defaults to 'latest'). Same block used for ALL steps. */ block?: BlockParam } @@ -116,15 +150,23 @@ 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 } = prepareRun(ts, options) + const { batchSize, maxConcurrentBatches, adaptiveBatching, maxBatchAttempts } = prepareRun(ts, options) await resolvePinnedBlock() - // `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. + // `run`'s fail-fast policy (F6a/F6b) — see `src/core/engine.ts`'s doc + // comment for the full hook contract. `buildStepCalls`/`consumeStepResults` + // let a throw propagate untouched (a plain synchronous throw inside this + // async function's body, exactly like the pre-F6a code). `executeBatch` is + // the raw, unwrapped executor call — nothing here branches on + // `adaptiveBatching`, because bisection is entirely the POOL's decision + // (see `src/core/pool.ts`), not this policy's: the exact same call works + // whether the pool hands it a whole batch or a bisected sub-batch. The + // length-mismatch guard also now lives in the pool (checked uniformly + // after any `execute()` resolves) rather than here — see its doc comment + // for why that's also what makes "never retried" free. `recordTerminal` is + // deliberately NOT provided: a terminal batch (single-call rejection, or + // the attempts cap exhausted) always triggers the pool's fail-fast + // cancellation, exactly as T14, whether or not bisection ever ran. const policy: StepEnginePolicy = { buildStepCalls(taskIndex, step) { const task = ts[taskIndex]! @@ -132,17 +174,8 @@ export async function runMultistepTasks( 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`, - ) - } - return results + executeBatch(batch) { + return executor.executeMulticall(batch, options?.block) }, consumeStepResults(taskIndex, step, results) { @@ -153,7 +186,7 @@ export async function runMultistepTasks( }, } - await runSteps(ts, { batchSize, maxConcurrentBatches }, policy) + await runSteps(ts, { batchSize, maxConcurrentBatches, adaptiveBatching, maxBatchAttempts }, policy) return ts.map((task) => task.finalize()) } diff --git a/src/core/runSettled.ts b/src/core/runSettled.ts index 9779593..075a86a 100644 --- a/src/core/runSettled.ts +++ b/src/core/runSettled.ts @@ -17,6 +17,20 @@ * Tasks whose `finalize()` copes with the resulting failures settle as * `fulfilled`. * + * **Adaptive bisection (F6b, `options.adaptiveBatching`):** with it OFF + * (default), the paragraph above is the WHOLE story, byte-for-byte unchanged + * from F5/1.1 — `executeBatch` catches every ordinary transport rejection + * itself and resolves with synthesized failures, so the pool + * (`src/core/pool.ts`) never even observes one. With it ON, `executeBatch` + * deliberately lets a transport rejection propagate to the pool instead of + * catching it, so the pool's bisection can split a `length > 1` batch and + * retry both halves. The per-call `DominoCallError` synthesis above still + * happens — just once per TERMINAL sub-batch (a single-call rejection, or + * the attempts cap exhausted for that original batch) via the + * `recordTerminal` policy hook below, instead of once per whole-batch + * rejection. Either way, the run never cancels on an ordinary transport + * failure: siblings, other steps, and other tasks keep going. + * * **Dead-task rule (controller ruling — not explicit in the spec text):** * if a task's `buildStepCalls()` or `consumeStepResults()` throws, that task * is marked dead: no further `buildStepCalls`, `consumeStepResults`, or @@ -37,7 +51,7 @@ * `consumeStepResults`/`finalize`) throws". */ -import type { MultistepTask, StepExecutor, RawResult, Address } from './types' +import type { MultistepTask, StepExecutor, StepCall, RawResult, Address } from './types' import type { BatchOptions } from './runMultistepTasks' import { DominoCallError } from './errors' import { prepareRun, resolvePinnedBlock } from './internal' @@ -78,6 +92,30 @@ export type SettledTaskResult = | { status: 'fulfilled'; value: T; diagnostics: TaskDiagnostics } | { status: 'rejected'; error: unknown; diagnostics: TaskDiagnostics } +/** + * Synthesize one `DominoCallError` (`kind: 'batch'`) per call, all sharing + * the SAME `cause` — the transport error that failed the physical batch (or + * bisected sub-batch, or coarse terminal group) these calls belong to. Used + * from two call sites that must stay provably identical: the non-adaptive + * `executeBatch` catch branch (whole-batch failure, T14 shape) and the + * adaptive `recordTerminal` policy hook (per-terminal-item failure, F6b) — + * factoring this out is what proves they produce the exact same shape. + */ +function synthesizeBatchFailures(calls: StepCall[], transportError: unknown): RawResult[] { + return calls.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, + }), + }), + ) +} + /** * Execute multiple MultistepTasks against a single StepExecutor, settling * each task independently instead of rejecting the whole call on the first @@ -93,8 +131,9 @@ export type SettledTaskResult = * an invalid `batchSize` makes the returned promise itself reject — it does * not produce a settled-but-rejected array. Same for a `StepExecutor` that * resolves with the wrong number of results for a batch (an implementation - * bug in the executor, not a call failure) — see the length-mismatch guard - * below, which mirrors `run`'s behavior exactly. + * bug in the executor, not a call failure) — the length-mismatch guard lives + * in the pool itself (`src/core/pool.ts`), applied uniformly to `run` and + * `runSettled`, and never retried even when `adaptiveBatching` is on. */ export async function runSettled( executor: StepExecutor, @@ -117,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 } = prepareRun(ts, options) + const { batchSize, maxConcurrentBatches, adaptiveBatching, maxBatchAttempts } = prepareRun(ts, options) if (ts.length === 0) return [] @@ -130,7 +169,7 @@ export async function runSettled( const dead: boolean[] = new Array(ts.length).fill(false) const deadError: unknown[] = new Array(ts.length) - // `runSettled`'s record-and-continue policy (F6a) — see + // `runSettled`'s record-and-continue policy (F6a/F6b) — see // `src/core/engine.ts`'s doc comment for the full hook contract. const policy: StepEnginePolicy = { buildStepCalls(taskIndex, step) { @@ -147,42 +186,47 @@ export async function runSettled( }, async executeBatch(batch) { - let results: RawResult[] + if (adaptiveBatching) { + // F6b: let a transport rejection propagate to the pool so its + // bisection can split this batch and retry both halves — the T14 + // catch-and-resolve-immediately behavior below would otherwise hide + // the rejection from the pool entirely, and bisection could never + // engage. The per-call `DominoCallError` synthesis still happens — + // once per TERMINAL sub-batch, via `recordTerminal` below, instead + // of once per whole-batch rejection here. The length-mismatch guard + // (a programmer error, never a transport failure) lives in the pool + // itself now — see `src/core/pool.ts` — so it applies uniformly to + // this call whether the pool handed it a whole batch or a bisected + // sub-batch. + return executor.executeMulticall(batch, options?.block) + } + + // Adaptive OFF (default) — byte-for-byte T14/F5 behavior: catch a + // transport rejection here and RESOLVE with synthesized per-call + // failures instead of letting it reach the pool. This is what makes + // `runSettled` never cancel by default: `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. (The length-mismatch guard is likewise + // now the pool's job — see `executeBatch` above and `pool.ts` — so + // this branch, like the adaptive one, is just the raw executor call.) try { - results = await executor.executeMulticall(batch, options?.block) + return 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, - }), - }), - ) + return synthesizeBatchFailures(batch, transportError) } + }, - // 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 + // F6b — adaptive bisection's terminal hook: called once per terminal + // sub-batch (a single-call rejection, or the attempts cap exhausted for + // its original batch) instead of the whole-batch catch above. Same + // synthesis, same shape, proven identical by sharing the one helper. + // Never invoked when `adaptiveBatching` is off (no rejection ever + // reaches the pool from `executeBatch` above in that case, except the + // length-mismatch bug, which bypasses this hook entirely and aborts — + // see `pool.ts`). + recordTerminal(calls, error) { + return synthesizeBatchFailures(calls, error) }, consumeStepResults(taskIndex, step, results) { @@ -199,7 +243,7 @@ export async function runSettled( }, } - await runSteps(ts, { batchSize, maxConcurrentBatches }, policy) + await runSteps(ts, { batchSize, maxConcurrentBatches, adaptiveBatching, maxBatchAttempts }, policy) return ts.map((task, i): SettledTaskResult => { // Read the task's OWN live diagnostics if it carries the channel (defineTask, From 0c1031196992cf8ab460cdf34af45ed9b8f8ad85 Mon Sep 17 00:00:00 2001 From: halaprix Date: Thu, 23 Jul 2026 19:45:22 +0200 Subject: [PATCH 2/2] fix: terminalize an exhausted-on-rejection bisection item immediately MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A length>1 rejection that consumed an original batch's LAST allowed attempt was still unconditionally pushing both children onto the queue. Those children only terminalize once claimed, at the claim-time cap check — but if a sibling original's cancellation wins the race first, queued items are never claimed at all, so this item's own terminal error silently never reached terminalErrors. Re-check attempts remaining before splitting; when none remain, terminalize this item synchronously in the same tick as its rejection instead. Regression test: two original batches, conc=2, maxBatchAttempts=2, where the lower-index original's final attempt is still in flight when the higher-index original exhausts and cancels first — the lower index's last transport error must still win the (batchIndex, callIndex) selection. --- src/__tests__/bisection.test.ts | 88 +++++++++++++++++++++++++++++++++ src/core/pool.ts | 49 ++++++++++++++---- 2 files changed, 128 insertions(+), 9 deletions(-) diff --git a/src/__tests__/bisection.test.ts b/src/__tests__/bisection.test.ts index eaf54fc..19f693e 100644 --- a/src/__tests__/bisection.test.ts +++ b/src/__tests__/bisection.test.ts @@ -468,3 +468,91 @@ describe('length-mismatch is never retried, even under adaptive', () => { expect(callCount).toBe(1) }) }) + +// ───────────────────────────────────────────────────────────────────────── +// External review (P1): a rejection that consumes an original batch's LAST +// attempt must terminalize immediately, not push children whose fate then +// depends on being claimed — because once a SIBLING original's cancellation +// wins the race, those children are never claimed, and this item's own +// terminal error would silently never enter `terminalErrors` at all. +// ───────────────────────────────────────────────────────────────────────── + +describe('concurrent cancellation does not drop an in-flight exhaustion terminal (P1)', () => { + it('original B (lower batch index) exhausts while in flight, after original A (higher index) has already cancelled -> B\'s last error wins the (batchIndex, callIndex) ordering', async () => { + // Task order [B, A] with batchSize=4 makes B's 4 calls exactly original + // batch 0 and A's 1 call exactly original batch 1 — deterministic index + // assignment, no reliance on timing for WHICH original gets which index. + const bTask: MultistepTask = { + maxStep: 1, + buildStepCalls(step) { + return step === 1 ? Array.from({ length: 4 }, (_, i) => call('b' + i)) : [] + }, + consumeStepResults() {}, + finalize() { + return null + }, + } + const aTask: MultistepTask = { + maxStep: 1, + buildStepCalls(step) { + return step === 1 ? [call('a0')] : [] + }, + consumeStepResults() {}, + finalize() { + return null + }, + } + + const bFinalError = new Error('B-fail-final') + const aError = new Error('A-fail') + + const executor: StepExecutor = { + async executeMulticall(calls: StepCall[]): Promise { + const keys = calls.map((c) => c.key) + if (keys[0]!.startsWith('a')) { + // A: single call, fast and deterministic — exhausts (terminalizes + // immediately, since length===1) and cancels the whole pool well + // before B's final attempt below ever settles. + await sleep(10) + throw aError + } + // B: 4 calls, maxBatchAttempts=2. + if (calls.length > 2) { + // Attempt 1 (whole batch) — reject fast so it splits immediately, + // long before A's cancellation. + await sleep(1) + throw new Error('B-fail-whole') + } + // Attempt 2 (one of the two length-2 children) — THIS is B's last + // allowed attempt (maxBatchAttempts: 2). Deliberately slow: still + // in flight when A's cancellation (at ~10ms) fires, so this + // rejection is discovered strictly AFTER `cancelled` is already + // true — exactly the race the P1 fix closes. + await sleep(50) + throw bFinalError + }, + } + + let thrown: unknown + let rejected = false + try { + await runMultistepTasks(executor, [bTask, aTask], { + batchSize: 4, + maxConcurrentBatches: 2, + adaptiveBatching: true, + maxBatchAttempts: 2, + }) + } catch (err) { + rejected = true + thrown = err + } + + expect(rejected).toBe(true) + // B is original batch index 0, A is index 1 — (batchIndex, callIndex) + // ordering deterministically selects B's error, REGARDLESS of A having + // cancelled the pool first. Before the P1 fix, B's in-flight final + // rejection was silently dropped (its children were pushed but never + // claimed post-cancellation) and `aError` was thrown instead. + expect(thrown).toBe(bFinalError) + }) +}) diff --git a/src/core/pool.ts b/src/core/pool.ts index 610e721..349ec91 100644 --- a/src/core/pool.ts +++ b/src/core/pool.ts @@ -53,14 +53,33 @@ * item, not this one — this one was never executed at all). Otherwise * `attempts[i]++` then execute. * - * Because the cap is checked at claim time rather than at split-decision - * time, a rejection always unconditionally pushes both halves when - * `adaptive && calls.length > 1` — the cap decides, independently and per - * child, whether each one actually gets to run. This is what produces the - * spec's "coarse group" behavior for free when the cap runs out mid-tree: - * whichever children can't get a claim-time attempt slot become terminal - * together, with no special-casing for "the cap ran out partway through a - * split". + * A rejection ALSO re-checks the cap itself, synchronously, before deciding + * to split (external review, P1): if `attempts[i]` has already reached + * `maxBatchAttempts` — i.e. this execution WAS the original's last allowed + * one — the item terminalizes immediately, right here, instead of pushing + * children that would merely be exhausted later, AT CLAIM TIME, by the check + * above. This is not just an optimization: those children only reach the + * claim-time check if a worker actually claims them, and once `cancelled` is + * true (a *sibling* original's own terminal, discovered concurrently, always + * wins the race unconditionally — no worker claims anything further once it + * flips), a queued item is — correctly, per (a) below — never claimed at + * all. Deferring a KNOWN exhaustion to "push children and hope they get + * claimed" would silently drop this item's own terminal error from + * `terminalErrors` whenever cancellation from elsewhere wins that race, + * corrupting both `cause = last transport error` and the lowest- + * `(batchIndex, callIndex)` selection in (c) below. Catching the exhausted + * case synchronously, in the same tick as the rejection that caused it, + * guarantees this item's own terminal is recorded regardless of what any + * other in-flight item does. + * + * When attempts DO remain, a rejection unconditionally pushes both halves — + * each child's own fate (execute, or find itself exhausted) is then decided + * independently at ITS OWN claim time, per the paragraph above. This is what + * produces the spec's "coarse group" behavior for free when the cap runs + * out mid-tree AFTER this point: whichever children can't get a claim-time + * attempt slot become terminal together, using the same `lastError[i]`, with + * no special-casing needed beyond the exhausted-at-rejection-time check just + * described. * * ## Terminal policy hook (F6b) * @@ -317,7 +336,19 @@ export async function runBatchPool( } catch (error) { lastError[origBatchIndex] = error - if (adaptive && calls.length > 1) { + // External review (P1): re-check the cap HERE, synchronously, before + // deciding to split — not just at each child's own future claim + // time. If this execution already consumed the original's last + // allowed attempt, pushing children would make this item's terminal + // depend on a worker actually claiming them later — which never + // happens once a *sibling* original's cancellation has already won + // the race (see the module doc's "Attempts accounting" section). + // Terminalizing THIS item immediately, covering its own full call + // range, guarantees its error is recorded regardless of what any + // concurrently in-flight item does. + const attemptsRemain = attempts[origBatchIndex]! < maxBatchAttempts + + if (adaptive && calls.length > 1 && attemptsRemain) { // Bisection: split and retry both halves through this SAME // central queue — no recursive await. See the module doc's "The // queue" section for why this can never deadlock.