← 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
M apps/workers/src/pipeline.ts
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 →