Skip to content

extract.ai: replace first-row preflight with bounded concurrent scheduling #1209

Description

@ebhills

Description

Replace the first-row preflight in extract.ai with a bounded, continuously refilled worker pool. Use the available concurrency immediately for valid work, while stopping new requests promptly when an invocation-wide fatal error makes the remaining work pointless.

The current safeguard was introduced in PR #1171: complete the first row/unique request before submitting the rest so an invalid model produces one failure. That serial wait also applies to valid models, including all of the first row's retries and backoff. It can consume most of a hosted recipe's execution window before the other workers start. This is backend scheduling behavior; Excel selection/header handling is separate.

Simply removing preflight_first=True is insufficient: the existing parallel path submits the entire input and consumes futures in submission order. It can queue unnecessary calls and hide a fast fatal error behind a slow earlier row.

Desired behavior

  • Start up to the effective worker count immediately, without waiting for a validation row or making an extra provider-validation request.
  • Refill freed slots continuously while the invocation remains healthy. Do not add a barrier between fixed batches or wait for every member of the first batch.
  • Detect fatal failures from any completed request promptly, stop admitting new requests/retries, cancel work that has not started, and surface the original error.
  • Preserve input/output association, output shapes, cache behavior, and existing concurrency/timeout/retry configuration.

Example current / desired

For five distinct uncached inputs and threads >= 5:

  • Current: row 1 completes, then rows 2–5 start.
  • Desired: rows 1–5 can start together. A slow or retrying row 1 does not hold the others.

For 1,000 inputs and threads=32:

  • Start up to 32 work items, then refill individual slots as they complete.
  • If a request reports an invalid model, stop filling slots immediately and cancel work that has not begun. Do not enqueue all 1,000 items first.
  • When every provider request fails fatally on its first attempt, no more than the initial 32 requests may reach the provider, and fewer is acceptable if the failure is detected while filling the initial window.

Full initial concurrency necessarily allows multiple requests to be sent before the first failure arrives. Calls already sent cannot be recalled. The general guarantee is bounded outstanding work and no new admissions after the fatal state is recorded, not a one-request failure guarantee or a lifetime limit of 32 calls when successful completions preceded a delayed failure.

Scheduling design

  1. Keep the existing worker setting. Resolve threads / default_concurrency as today. Let W be the effective worker count, bounded by available work. Preserve empty-input and threads=1 behavior. Do not turn the default of 32 into a hard cap or introduce a separate public batch-size setting.
  2. Use a rolling window of at most W outstanding work items. Retain unsubmitted work in a local iterator; do not eagerly create a future for every row. Fill the initial window without waiting for successful results, unless a fatal error is already known. With result caching enabled, retain existing grouping of equivalent requests; with caching disabled, retain per-row execution.
  3. Observe completion order; return input order. Use completion notifications or wait(..., FIRST_COMPLETED) to inspect whichever workers finish. Check all available completions for failures before refilling. Store successful results against their original indices and expand duplicate groups using defensive copies. Never block the scheduler on the first submitted future.
  4. Record fatal state in the worker, not only in the collector. Use invocation-local, thread-safe state containing the first fatal cause and a stop signal. The worker records the failure before publishing completion, so another successful completion cannot keep feeding the pool while the collector has yet to inspect the fatal future.
  5. Coordinate stop and admission. Check the stop state when scheduling work and immediately before each provider attempt, including retries. Admission and recording the stop state must share synchronization so their order is defined. Release the lock before network I/O; do not serialize HTTP requests. An attempt admitted before the stop signal may already be in flight. No attempt may be admitted afterward.
  6. Cancel promptly and preserve the cause. Stop refilling, cancel pending futures, and interrupt retry/backoff waits for work owned only by the aborted invocation. Cancellation is internal control flow: it must not become a cell error, blank result, retry, or replacement for the original fatal exception. Avoid executor cleanup that silently waits for all other rows before exposing the fatal error. Already-running HTTP calls may finish their current configured attempt; clean up their futures/cache coordination without starting further work for the aborted invocation.
  7. Keep ordinary retries unchanged. While the invocation is healthy, retain per-attempt timeouts, retry counts, and backoff. A row timeout, 429, transient provider failure, or exhausted ordinary row-level retry does not automatically become an invocation-wide fatal error. Do not restore a shared batch deadline.

Fatal error classification

Preserve the failures that already abort extraction: model_not_found, invalid shared output schema, and missing/invalid API key, plus normal propagation of unexpected unhandled exceptions. Keep existing public exception compatibility and useful messages; a dedicated internal fatal exception may subclass the existing public exception type.

Use explicit exception/classification state, not strings found in returned cell values. Do not classify every HTTP 4xx or every response-validation error as fatal. Perform existing local configuration/schema validation before submitting work wherever possible; do not add a provider probe or reject an otherwise supported model solely because it is absent from the local catalog.

Cache and concurrent-invocation isolation

  • Preserve cache keys, tenant/credential separation, TTL/LRU behavior, warm-cache reuse, duplicate expansion, and cross-invocation single-flight deduplication.
  • Cache hits must not generate provider calls; resolve/refill them without a first-row gate.
  • Never cache cancellation or publish a synthetic cancellation as the outcome of a shared request. Release single-flight waiters on every exit path.
  • Aborting one invocation must not cancel another healthy invocation. If another active invocation still needs a shared computation, it may complete under its normal policy; detach the canceled consumer or safely transfer ownership rather than poisoning shared state.
  • Keep successful completed cache entries usable. The failed invocation still raises its fatal error rather than returning a partial successful batch.

Implementation scope

Apply the new scheduling behavior to extract.ai, covering standalone Python use and recipe execution, including WranglesXL's recipe wrapper. Implement the scheduler in the existing shared execution infrastructure with an explicit internal opt-in for this behavior; preserve other callers' scheduling contracts.

Relevant code:

Cover each extraction protocol still supported on main when implemented; do not restore a removed legacy path. Leave the separate wrangles.ai callers' preflight behavior unchanged.

Excel selected-row/header handling, Columns/JSON formatting, the recipe wrapper, model-catalog changes, and increasing hosted timeouts are outside this issue. Eric is tracking the hidden XL timeout message separately. The current development package's reasoning-setting support is a separate deployment correction.

Acceptance criteria and verification

Use deterministic events/barriers and instrumented mock calls instead of timing-sensitive sleeps.

  • Full first-window utilization: with N >= W distinct inputs, all W workers can enter the mock provider before any request completes; no validation/probe call is added.
  • Continuous refill: with a slow first row and N > W, completing another row admits the next item without waiting for row 1 or the rest of the initial window. Outstanding work never exceeds W.
  • Fatal error from any row: a later row reports a fatal model error while row 1 remains blocked. The invocation exposes the original error without waiting for row 1; no new work is admitted after the fatal state is recorded. Release blocked mocks during test cleanup.
  • Bounded invalid-model waste: for a large all-invalid input, provider attempts are <= W, fatal errors are not retried, and the remaining rows are never submitted. Replace the existing test's universal single-call expectation with this concurrency-aware contract.
  • Admission race: synchronize a successful completion/refill against another worker's fatal failure. Verify the recorded stop state prevents subsequent initial attempts and retries, and document the treatment of attempts admitted just beforehand.
  • Retry cancellation: a sibling waiting in backoff stops promptly after a fatal error; a healthy invocation still follows its existing timeout/retry policy. Ordinary 429s, timeouts, and malformed row responses do not cancel unrelated rows.
  • Result compatibility: preserve ordered results, scalar/list return shapes, recipe row mapping, nested/null values, duplicate handling, and defensive copies with cache both enabled and disabled.
  • Cache isolation and cleanup: cover warm hits, cross-invocation coalescing, cancellation of a shared request's original consumer, and two concurrent invocations where only one fails. No cached cancellation, stranded waiter, leaked in-flight entry, or cancellation of the healthy invocation.
  • Boundary/regression cases: cover empty/single input, threads=1, N < W, explicit worker overrides, supported extraction protocols, and unchanged behavior for other shared-executor callers.
  • Documentation and delivery: describe the new failure/concurrency contract; run focused tests and pytest -c pytest-local.ini. Cover any new test file in that configuration. After merging to main and promoting the package to execute-recipe-dev, run a small valid saved-model smoke test and record request start times, peak concurrency, output correctness, and elapsed time. Keep offline, CI, deployment, and live evidence distinct.

Measure before/after latency using the same effective model, reasoning, retry, and cache settings. Full concurrency is an acceptance requirement; a particular wall-clock speedup is not guaranteed because provider latency and prompt-cache behavior vary.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

Type

No type

Fields

Stage

None yet

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions