← back to Homesonspec

apps/workers/src/payload-guard.itest.ts

179 lines

import { beforeAll, describe, expect, it } from "vitest";
import { prisma } from "@homesonspec/database";
import { canonicalKey } from "@homesonspec/shared";
import type { ExtractedRecord } from "@homesonspec/schemas";
import type { RawPage, SourceAdapter } from "@homesonspec/collectors-common";
import { extractStage } from "./pipeline";

/**
 * Integration: the TK-11125 forward-only spurious-generation guard, proven at
 * the StagedRecord + SourceEvidence data boundary (e2e-proof R3). A new raw
 * snapshot with byte-identical extracted data must NOT mint a new generation;
 * a real field change must; history is preserved throughout.
 */

const SOURCE_KEY = "payload-guard-fixtures";
// One physical home — canonicalKey is stable across every snapshot below.
const HINTS = { builderSlug: "guard-test", communityName: "Guard Bend", address: "1 Guard Way" };
const KEY = canonicalKey(HINTS);

const fv = (value: unknown, raw: string | null = null) => ({
  value,
  raw,
  evidenceText: raw,
  sourceUrl: "fixture://guard",
  confidence: value === null ? 0 : 1,
});

// A synthetic adapter whose extract() returns whatever record the test set.
let currentRecord: ExtractedRecord;
const adapter: SourceAdapter = {
  key: SOURCE_KEY,
  version: "guard-test-1",
  // eslint-disable-next-line require-yield
  async *fetch() {
    /* unused: tests call extractStage directly */
  },
  extract() {
    return { records: [currentRecord], errors: [] };
  },
};

function record(fields: Record<string, ReturnType<typeof fv>>): ExtractedRecord {
  return { entityType: "sales_office", canonicalHints: HINTS, fields } as ExtractedRecord;
}

// A distinct RawSnapshot per call = "volatile page bytes changed".
let snapCounter = 0;
async function newSnapshot(sourceId: string): Promise<string> {
  snapCounter += 1;
  const contentHash = `guardhash-${snapCounter}-${"0".repeat(50)}`.slice(0, 64);
  const snap = await prisma.rawSnapshot.create({
    data: {
      sourceId,
      url: "https://guard.test/home/1",
      retrievedAt: new Date(),
      contentHash,
      storagePath: `payload-guard/${contentHash}.json`,
      contentType: "application/json",
      httpStatus: 200,
    },
  });
  return snap.id;
}

const gens = () => prisma.stagedRecord.count({ where: { canonicalKey: KEY } });
const evidenceRows = async () => {
  const staged = await prisma.stagedRecord.findMany({
    where: { canonicalKey: KEY },
    select: { id: true },
  });
  return prisma.sourceEvidence.count({
    where: { stagedRecordId: { in: staged.map((s) => s.id) } },
  });
};

let sourceId: string;

beforeAll(async () => {
  expect(process.env.DATABASE_URL).toMatch(/homesonspec_test/);
  // Clean only this test's rows (co-exists with the other itest file).
  const staged = await prisma.stagedRecord.findMany({
    where: { canonicalKey: KEY },
    select: { id: true },
  });
  await prisma.sourceEvidence.deleteMany({ where: { stagedRecordId: { in: staged.map((s) => s.id) } } });
  await prisma.stagedRecord.deleteMany({ where: { canonicalKey: KEY } });
  await prisma.rawSnapshot.deleteMany({ where: { source: { key: SOURCE_KEY } } });
  await prisma.sourceRegistry.deleteMany({ where: { key: SOURCE_KEY } });
  await prisma.builder.deleteMany({ where: { slug: "guard-test" } });

  const builder = await prisma.builder.create({
    data: { slug: "guard-test", name: "Guard Test Homes", isDemo: true },
  });
  const source = await prisma.sourceRegistry.create({
    data: {
      key: SOURCE_KEY,
      name: "Payload Guard (synthetic)",
      builderId: builder.id,
      collectionMethod: "SYNTHETIC",
      mediaRights: "NONE",
    },
  });
  sourceId = source.id;
});

async function extractOn(snapshotId: string) {
  const source = await prisma.sourceRegistry.findUniqueOrThrow({ where: { key: SOURCE_KEY } });
  const page = {
    url: "https://guard.test/home/1",
    retrievedAt: new Date().toISOString(),
    contentType: "application/json",
    body: Buffer.from("{}"),
    contentHash: "unused-in-extractStage",
  } as RawPage;
  return extractStage(adapter, source, page, snapshotId);
}

describe("TK-11125 spurious-generation guard", () => {
  it("case 4a — a brand-new home → exactly 1 generation + evidence", async () => {
    currentRecord = record({ name: fv("Guard Bend Office"), phone: fv("512-555-0100") });
    const snap1 = await newSnapshot(sourceId);
    const { stagedIds } = await extractOn(snap1);
    expect(stagedIds).toHaveLength(1);
    expect(await gens()).toBe(1);
    expect(await evidenceRows()).toBe(2); // name + phone
  });

  it("case 2 — new snapshot, byte-identical fields → 0 new generations (the fix)", async () => {
    const evBefore = await evidenceRows();
    currentRecord = record({ name: fv("Guard Bend Office"), phone: fv("512-555-0100") });
    const snap2 = await newSnapshot(sourceId); // volatile bytes changed → different snapshot
    const { stagedIds } = await extractOn(snap2);
    expect(stagedIds).toHaveLength(0); // SKIPPED
    expect(await gens()).toBe(1); // still one generation
    expect(await evidenceRows()).toBe(evBefore); // no fresh evidence set minted
  });

  it("case 3 — new snapshot, a real field change → exactly 1 new generation", async () => {
    currentRecord = record({ name: fv("Guard Bend Office"), phone: fv("512-555-0199") }); // phone changed
    const snap3 = await newSnapshot(sourceId);
    const { stagedIds } = await extractOn(snap3);
    expect(stagedIds).toHaveLength(1); // NEW generation
    expect(await gens()).toBe(2); // history preserved: 2 generations now
    expect(await evidenceRows()).toBe(4); // gen1 (2) + gen2 (2)
  });

  it("case 2b — re-seeing the just-changed data on yet another snapshot → 0 new (steady state)", async () => {
    const gensBefore = await gens();
    currentRecord = record({ name: fv("Guard Bend Office"), phone: fv("512-555-0199") });
    const snap4 = await newSnapshot(sourceId);
    const { stagedIds } = await extractOn(snap4);
    expect(stagedIds).toHaveLength(0);
    expect(await gens()).toBe(gensBefore); // no drift once data settles
  });

  it("case 4b — FORCE_REEXTRACT same snapshot adding a field → updates the row, no new generation", async () => {
    // Re-extract the SAME latest snapshot with an added field. payloadHash differs
    // → guard's snapshotId-equality escape lets it fall through to the upsert UPDATE,
    // preserving today's FORCE_REEXTRACT semantics (row updated in place, not forked).
    const latest = await prisma.stagedRecord.findFirstOrThrow({
      where: { canonicalKey: KEY },
      orderBy: { createdAt: "desc" },
      select: { id: true, snapshotId: true },
    });
    const gensBefore = await gens();
    currentRecord = record({
      name: fv("Guard Bend Office"),
      phone: fv("512-555-0199"),
      hours: fv("9-5"), // newly backfilled field
    });
    const { stagedIds } = await extractOn(latest.snapshotId!);
    expect(stagedIds).toHaveLength(1);
    expect(stagedIds[0]).toBe(latest.id); // SAME row updated, not a new generation
    expect(await gens()).toBe(gensBefore);
    const ev = await prisma.sourceEvidence.count({ where: { stagedRecordId: latest.id } });
    expect(ev).toBe(3); // name + phone + hours (evidence replaced for the row)
  });
});