Skip to content

[bug] Worker handlers are not idempotent under lease reclaim — duplicate writes, lost cost accounting, and inflated/erased counters on retry #36

Description

@joaovictor712

Summary

The Wave 6 durable job queue (lib/jobs/queue.ts) is designed for "crashed-worker reclaim": claimNextJob re-claims any row whose lease_expires_at <= now, and the default lease is DEFAULT_LEASE_MS = 10 * 60 * 1000 (10 min). That is correct queue-level behavior — it was hardened in merged PR #5 — but none of the worker handlers are written to be idempotent under that reclaim. They are written as if they execute exactly once, then mutate extraction_results, parse_runs, extraction_runs, and on-disk extraction.json files unguarded.

The Anthropic-batch handler is the worst offender:

worker/handlers/anthropic-batch.ts:55–149

for await (const entry of streamBatchResults(batchId)) {
  const productSlug = entry.custom_id;
  const resultRow = db.prepare(
    `SELECT id, output_path FROM extraction_results WHERE run_id = ? AND product_slug = ?`
  ).get(runId, productSlug) as { id: number; output_path: string | null } | undefined;

  if (!resultRow) {
    console.warn(`[batch-poll] no extraction_results row for run=${runId} slug=${productSlug}; skipping`);
    continue;                                              // ← silently drops cost + counter
  }

  if (entry.result.type === "succeeded") {
    // ... no idempotency check ...
    const outPath = writeExtractionJson(productDir, parsed);   // ← non-atomic, overwrites
    db.prepare(`UPDATE extraction_results SET status='completed', ...`).run(...);   // ← no status guard
    successCount++;
  }
  ...
}

const finalStatus = failCount === 0 ? "completed" : "completed-with-errors";
db.prepare(
  `UPDATE extraction_runs SET status=?, cost_actual_usd=?, success_count=?, fail_count=?, completed_at=? WHERE id=?`
).run(finalStatus, totalCostUsd, successCount, failCount, ..., runId);

Failure modes that hit production today:

  1. Lease expiry during streaming. Streaming + parsing + writing N results for a 50-product batch routinely exceeds 10 min on slower hosts / network. extendLease is never called. While the handler is mid-loop, another poller (or the same worker after restart) wins the next claimNextJob race and re-enters the same handler. The loop restarts from the top of streamBatchResults — every result is processed twice.
  2. No per-result idempotency. The UPDATE extraction_results SET status='completed', ... and markFailed paths have no WHERE status IN ('queued','running') guard. On reclaim, a row already marked 'completed' is updated again — completed_at, output_tokens, cache_*_tokens, and output_path are all overwritten with the second attempt's values, hiding the first attempt's audit trail.
  3. writeExtractionJson overwrites the prior good file non-atomically. lib/pipeline/extract.ts:425-432 uses raw fs.writeFileSync(out, …) instead of the project's atomicWriteJson helper at lib/pipeline/extraction-paths.ts:121. A reclaim-during-write — or a SIGKILL during a single attempt — leaves a truncated extraction.json as the canonical output. The spot-fix flow already uses atomicWriteJson to write the same file (lib/pipeline/spotfix.ts:605); the main extract path is the inconsistent one.
  4. In-memory counters reset every attempt. successCount, failCount, totalCostUsd are local variables that start at 0 per invocation. The final aggregate UPDATE extraction_runs SET success_count=?, fail_count=?, cost_actual_usd=? therefore reflects only the last attempt's totals — not the cumulative work the user was actually billed for. If a reclaim happens during the loop and the second attempt errors on, say, 5 of 50 entries that the first attempt had already written successfully, the run is finalized with success_count=45 fail_count=5 and cost_actual_usd set to only the second attempt's cost. The first attempt's spend is lost to accounting.
  5. Missing result rows silently dropped (line 65–70). When entry.custom_id does not match any extraction_results row (which can happen on schema drift, manual DB edits, or — once [bug] batch worker writes extraction.json to {DATA_DIR}/{slug}/ instead of nested product path #26 is fixed — any path mismatch between submission and poll), the entry is continue'd with a console.warn. No cost is recorded; neither successCount nor failCount is incremented. The run can finish status='completed' with success_count + fail_count < request_count, and no surfaced signal that products were lost.
  6. markFailed is unguarded too (anthropic-batch.ts:156-161). A reclaim that hits a transient JSON-parse error on the second attempt can flip a row that was previously 'completed' to 'failed', erasing the prior good extraction state.

The Reducto handler has the same shape (worker/handlers/reducto.ts:82-104): the only idempotency check is for status='completed' AND output_md_path NOT NULL at the top of the function; once that gate is passed, every UPDATE parse_runs SET status='failed' … or call to processParseResult is unguarded against reclaim.

Per CONTRIBUTING and README, extraction.json is the headline output of the four-stage pipeline and extraction_results.cost_actual_usd is what lib/usage.ts and the /usage dashboard (merged in PR #3) read for spend tracking. Both are silently corrupted by this bug today.

Steps to reproduce

The reclaim path is hard to trigger by accident but trivial to force:

  1. Submit a small batch:
    npm run extract-one -- --product server/dell/poweredge/r770 --mode batch
  2. Start the worker, let it call streamBatchResults and process the first few entries:
    npm run worker
  3. Shorten the lease so reclaim happens immediately. In a separate process, either run UPDATE jobs SET lease_expires_at = '1970-01-01T00:00:00.000Z' WHERE type='anthropic-batch-poll' AND status='running' against data/studio.db, or in test, call claimNextJob({ leaseMs: 1 }).
  4. Start a second worker (npm run worker again — concurrent workers are explicitly supported by the queue design). The second worker re-claims the same job; the handler streams the batch results again. Observe:
    • extraction_results rows have their output_tokens / cache_read_tokens / completed_at rewritten.
    • extraction.json files have new mtimes (second write).
    • Final extraction_runs.cost_actual_usd reflects only the second attempt.
    • If the second attempt has any failures the first didn't, the run lands in completed-with-errors even though every result was actually successful on at least one attempt.

For (3) — the non-atomic write specifically — point a SIGKILL at the worker mid-writeFileSync on a multi-MB extraction (instrument with a process.kill after 1ms of writing). The on-disk extraction.json is now a truncated JSON file — the prior good extraction is unrecoverable.

Expected behavior

  • handleAnthropicBatchPoll and handleReductoPoll are safe to invoke any number of times for the same (runId, productSlug) or (parseRunId) pair. Every state transition uses a guarded UPDATE (WHERE status IN ('queued','running') for terminal writes; WHERE status IS NOT 'completed' for markFailed).
  • Per-result idempotency: a result row already in a terminal state (completed / failed) is skipped during streaming. It is not re-written, its output_path is not re-rewritten, and its cost is taken from the DB row (not from the recomputed usage of the second attempt).
  • Cost and counter accounting comes from a single SQL aggregate over extraction_results, not from in-memory accumulators that reset per attempt:
    SELECT
      COUNT(*) FILTER (WHERE status='completed') AS success_count,
      COUNT(*) FILTER (WHERE status='failed')    AS fail_count,
      COALESCE(SUM(cost_usd), 0)                 AS cost_actual_usd
      FROM extraction_results WHERE run_id = ?
  • writeExtractionJson writes through atomicWriteJson (tmp + rename) — matching the spot-fix path and every other canonical-file writer (sources-yaml.ts, schemas-fs.ts, option-matrices.ts, product-md-skel.ts, program-docs.ts, sources.ts).
  • Long-running handlers call extendLease periodically (e.g. every K results) so legitimate work doesn't get pre-empted mid-stream.
  • A custom_id with no matching extraction_results row is treated as a fail (an orphan insert with status='failed', error_message='no result row at submission') rather than silently dropped — so the run's final counter math always sums to request_count.

Actual behavior

Every reclaim corrupts state. Even without reclaim, a single SIGKILL during the batch-loop's writeFileSync leaves a partial extraction.json as canonical output, and missing-custom_id entries silently disappear.

Suggested fix

This is a structural fix spanning three layers — the per-result DB helpers, the long-running handlers themselves, and the atomic-write convention for extraction.json. None of these have been touched by the existing open issues / PRs (cross-checked: see "Cross-check" section below).

1. New idempotent helpers in lib/pipeline/runs.ts:

/** Transition queued → running. Returns true iff the row was actually moved. */
export function markResultRunningIdempotent(resultId: number): boolean { ... }

/** Transition queued/running → completed, recording cost + tokens.
 *  Returns false if the row is already in a terminal state. */
export function markResultCompletedIdempotent(
  resultId: number,
  payload: { outputPath: string; inputTokens: number; outputTokens: number;
             cacheReadTokens: number; cacheCreationTokens: number; costUsd: number }
): boolean { ... }

/** Transition queued/running → failed. Returns false if the row was already completed. */
export function markResultFailedIdempotent(resultId: number, error: string): boolean { ... }

/** Single source of truth for run-level aggregates — read directly from extraction_results. */
export function aggregateRunTotalsFromResults(runId: number): {
  successCount: number; failCount: number; costActualUsd: number;
} { ... }

/** Finalize: derives status, counts, and cost from the DB, with a single UPDATE. */
export function finalizeRunFromResults(runId: number, requestCount: number): void { ... }

The DB needs one new column on extraction_results to make per-result cost auditable: cost_usd REAL. Add migration 0003_extraction_result_cost.sql.

2. worker/handlers/anthropic-batch.ts is rewritten to:

  • Skip rows already in 'completed'/'failed' state (per-result idempotency).
  • Call extendLease(jobId, LEASE_RENEWAL_MS) every LEASE_RENEW_EVERY_N_RESULTS (e.g. 10) results.
  • Use the new guarded-markResult*Idempotent helpers — every transition is WHERE status IN (…) constrained.
  • After the loop, call finalizeRunFromResults(runId, requestCount) — the in-memory successCount/failCount/totalCostUsd accumulators are gone.
  • On missing-result-row: synthesize a status='failed' row via insertResult + markResultFailedIdempotent with error_message="no result row at submission; custom_id mismatch", so the request_count math always balances.

3. worker/handlers/reducto.ts gets parallel guards on its UPDATE parse_runs SET status='failed' … and processParseResult paths.

4. lib/pipeline/extract.ts:writeExtractionJson switches to atomicWriteJson.

5. Unit tests covering:

  • Reclaim mid-loop: simulate processing 3 of 5 entries, force lease expiry, second invocation processes the remaining 2 (and the first 3 are skipped).
  • Per-result idempotent transitions (queued→running→completed, repeated calls are no-ops).
  • Aggregate computed from DB, not accumulators (verified by deliberately mutating the in-memory totals and observing the final UPDATE uses DB-derived figures).
  • writeExtractionJson crash-safety (renameSync fails → prior good file preserved, tmp sibling left behind for inspection).
  • Missing custom_id orphan handling: a batch response with a custom_id not in extraction_results produces a failed row, not a silent drop; final success+fail == request_count.

PR scope

~250–400 lines of source across:

  • lib/pipeline/runs.ts (~120 new lines: 5 new exported functions + the aggregate query)
  • worker/handlers/anthropic-batch.ts (~100 lines rewritten to use guarded helpers + lease renewal)
  • worker/handlers/reducto.ts (~30 lines: guarded UPDATEs)
  • lib/pipeline/extract.ts (~3 lines: swap to atomicWriteJson)
  • lib/db/migrations/0003_extraction_result_cost.sql (~10 lines)

Plus ~250–400 lines of vitest coverage in tests/jobs/, tests/worker/, tests/pipeline/runs.test.ts.

The source-side change is intentionally structural — new exported functions, guarded SQL writes, lease-extension control flow, and a migration — so density should sit at or near the cap and the contribution bonus from total work clears the high-token threshold rather than just the floor.

Cross-check against existing issues / PRs (open + closed)

# Subject Why this is distinct
#5 (merged) Make job queue durable and race-safe Queue-level (claim/lease/backoff). This bug is handler-level idempotency that the durable queue itself cannot enforce — handlers must opt in.
#26/#27 Batch worker writes extraction.json to wrong product directory where the file is written. This bug is about how (atomic) and how often (reclaim safety).
#28/#29 Extraction index cache invalidation Read-side stale cache. This bug is write-path correctness.
#32/#33 Schema save/rollback writes filesystem before DB Schema-save transactional ordering. Different table, different file, different invariant.
#22/#23, #24/#25, #10/#11, #30/#31 Path traversal / orphan sweep / data-root bounds / MD frontmatter silent skip Unrelated code paths.
#3 (merged) /usage cost & cache dashboard Reads extraction_runs.cost_actual_usd and extraction_results.cost_*_tokens — i.e. the columns this bug silently corrupts on reclaim. The dashboard is the user-visible symptom, but the fix is upstream.

Activity

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

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions