[object Object]

← back to Homesonspec

pipeline: make evidence delete+create atomic (contrarian fix) — kills concurrent-reextract race

11a5e5facf98c8ebe43b1b14cf3ebdea3c96bd1a · 2026-07-29 01:48:49 -0700 · Steve Abrams

Cody caught a live race: a FORCE_REEXTRACT backfill running alongside the loop's dr-horton sweep both
did deleteMany+createMany on the same stagedRecord's SourceEvidence non-atomically, interleaving to
DOUBLE evidence rows (verified 3 dup groups). Wrapped both in prisma.$transaction so evidence
replacement is atomic — fixes it for ALL adapters + concurrent runs. (3 cosmetic dup provenance rows
remain; harmless, cleanable later — they don't affect home/geo data.)

Files touched

Diff

commit 11a5e5facf98c8ebe43b1b14cf3ebdea3c96bd1a
Author: Steve Abrams <steve@designerwallcoverings.com>
Date:   Wed Jul 29 01:48:49 2026 -0700

    pipeline: make evidence delete+create atomic (contrarian fix) — kills concurrent-reextract race
    
    Cody caught a live race: a FORCE_REEXTRACT backfill running alongside the loop's dr-horton sweep both
    did deleteMany+createMany on the same stagedRecord's SourceEvidence non-atomically, interleaving to
    DOUBLE evidence rows (verified 3 dup groups). Wrapped both in prisma.$transaction so evidence
    replacement is atomic — fixes it for ALL adapters + concurrent runs. (3 cosmetic dup provenance rows
    remain; harmless, cleanable later — they don't affect home/geo data.)
---
 apps/workers/src/pipeline.ts | 37 +++++++++++++++++++++----------------
 1 file changed, 21 insertions(+), 16 deletions(-)

diff --git a/apps/workers/src/pipeline.ts b/apps/workers/src/pipeline.ts
index f11a8441..631c9976 100644
--- a/apps/workers/src/pipeline.ts
+++ b/apps/workers/src/pipeline.ts
@@ -126,22 +126,27 @@ export async function extractStage(
     });
     stagedIds.push(staged.id);
 
-    // Replace this row's evidence rather than append (idempotent re-extract).
-    await prisma.sourceEvidence.deleteMany({ where: { stagedRecordId: staged.id } });
-    await prisma.sourceEvidence.createMany({
-      data: Object.entries(normalized.fields).map(([field, envelope]) => ({
-        stagedRecordId: staged.id,
-        field,
-        sourceUrl: envelope.sourceUrl,
-        retrievedAt: new Date(page.retrievedAt),
-        contentHash: page.contentHash,
-        extractorVersion: adapter.version,
-        rawValue: envelope.raw,
-        normalizedValue: envelope.value === null ? null : String(envelope.value),
-        evidenceText: envelope.evidenceText,
-        confidence: envelope.confidence,
-      })),
-    });
+    // Replace this row's evidence rather than append (idempotent re-extract). Wrapped in a
+    // transaction so a concurrent re-extract of the same source (e.g. a FORCE_REEXTRACT backfill
+    // running alongside a normal loop sweep) can't interleave delete/create and DOUBLE the evidence
+    // rows — the delete+create is atomic. [yolo iter-5 contrarian fix]
+    await prisma.$transaction([
+      prisma.sourceEvidence.deleteMany({ where: { stagedRecordId: staged.id } }),
+      prisma.sourceEvidence.createMany({
+        data: Object.entries(normalized.fields).map(([field, envelope]) => ({
+          stagedRecordId: staged.id,
+          field,
+          sourceUrl: envelope.sourceUrl,
+          retrievedAt: new Date(page.retrievedAt),
+          contentHash: page.contentHash,
+          extractorVersion: adapter.version,
+          rawValue: envelope.raw,
+          normalizedValue: envelope.value === null ? null : String(envelope.value),
+          evidenceText: envelope.evidenceText,
+          confidence: envelope.confidence,
+        })),
+      }),
+    ]);
   }
   return { stagedIds, errors };
 }

← ee8b7a7c geocode: shea per-home coord backfill from detail pages (yol  ·  back to Homesonspec  ·  david-weekley: capture city-gated community MapPoint coord ( de09b351 →