From 4e75c6c4bb6740fe4f525636242071d1bf264ae6 Mon Sep 17 00:00:00 2001 From: Anthony Ettinger Date: Tue, 28 Jul 2026 09:40:57 +0000 Subject: [PATCH] feat(leads): actually record who a lead is, once, across every campaign MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The contacts table shipped empty and stayed empty. Discovery found people and threw them away on every run, which left the source of truth it was built to be as a schema with nothing in it. Two things had to change before anything could be written. Email was required and was the key, which assumed every person arrives with an address. Person discovery does the reverse: a directory gives a name, a title and a LinkedIn profile, and the address is what the pipeline then goes looking for. Those people could not be stored at all. Identity is now the email when there is one — two records with the same address are the same person by definition — and the normalised name and employer when there is not. Weaker, but the alternative is a fresh row for the same human every run. The key is maintained by trigger rather than generated, because finding an address for a known name has to move that row from the name form to the email form, which is exactly the transition a generated column cannot make. Merging fills gaps and does not overwrite. A scraped company name must not replace one a human typed, so every field records where it came from and a better-sourced value wins; the loser is kept in `alternates` rather than dropped, because "we saw something else" is information and a wrong overwrite is otherwise unrecoverable. Socials merge per network, so a run that finds only GitHub does not discard a known LinkedIn. People are recorded even when no address was found, which is the case that motivated all of it: a directory publishes the name and withholds the email, and waiting for an address means rediscovering the same person forever. isContactBlocked returns false when the lookup itself fails. A failed read is not permission to send, but neither is it grounds to silently block a campaign — the caller's own suppression checks still run. Co-Authored-By: Claude Opus 5 (1M context) --- app/actions/leads.ts | 39 +++- lib/outreach/contacts.ts | 208 ++++++++++++++++++ lib/outreach/discover.ts | 14 +- lib/outreach/pipeline.ts | 52 +++++ .../20260728070000_contacts_identity_key.sql | 62 ++++++ tests/contacts-merge.test.ts | 61 +++++ 6 files changed, 433 insertions(+), 3 deletions(-) create mode 100644 lib/outreach/contacts.ts create mode 100644 supabase/migrations/20260728070000_contacts_identity_key.sql create mode 100644 tests/contacts-merge.test.ts diff --git a/app/actions/leads.ts b/app/actions/leads.ts index 36befe05..b3541f40 100644 --- a/app/actions/leads.ts +++ b/app/actions/leads.ts @@ -15,6 +15,7 @@ import { import { addSuppression } from "@/lib/outreach/suppress"; import { loadAddressSettings } from "@/lib/outreach/postalAddress"; import { discoverProspects } from "@/lib/outreach/discover"; +import { upsertContact } from "@/lib/outreach/contacts"; import { runEmailCampaignTick, CAMPAIGN_COLUMNS, summarize, type CampaignRow } from "@/lib/outreach/runner"; type Ok> = { ok: true } & T; @@ -75,12 +76,44 @@ export async function findLeadsAction(input: { // prospect that publishes no address — so a thousand-lead run gets an // explicit ceiling rather than an unbounded bill. const contactSearchBudget = { remaining: Math.min(limit, 100) }; + const { data: projectRow } = await serviceClient() + .from("projects") + .select("organization_id") + .eq("id", input.projectId) + .maybeSingle(); + const orgId = (projectRow?.organization_id as string | null) ?? null; + const found = await discoverProspects({ queries: input.query?.trim() ? [input.query.trim()] : [], seedUrls: input.seedUrl?.trim() ? [input.seedUrl.trim()] : [], limit, }); - if (!found.prospects.length) { + // People named on the pages we opened are recorded even when no address + // was found for them. A directory gives a name, a title and a LinkedIn + // profile and withholds the email; discarding that until an address turns + // up means rediscovering the same person on every run. + let peopleRecorded = 0; + if (found.people?.length && orgId) { + for (const person of found.people) { + const res = await upsertContact({ + organizationId: orgId, + source: person.source === "json-ld" ? "json-ld" : "page", + fields: { + fullName: person.fullName, + title: person.jobTitle, + companyName: person.company, + companySite: person.companySite, + linkedinUrl: person.linkedinUrl, + country: person.location, + sourceUrl: person.sourceUrl, + socials: person.socials, + }, + }); + if (res) peopleRecorded += 1; + } + } + + if (!found.prospects.length && !peopleRecorded) { return { ok: false, error: found.errors.join("; ") || "No businesses found for that search." }; } @@ -112,7 +145,9 @@ export async function findLeadsAction(input: { ok: true, added, scanning, - note: `${added} leads added.`, + note: peopleRecorded + ? `${added} leads added, ${peopleRecorded} people recorded.` + : `${added} leads added.`, }; } diff --git a/lib/outreach/contacts.ts b/lib/outreach/contacts.ts new file mode 100644 index 00000000..596aa433 --- /dev/null +++ b/lib/outreach/contacts.ts @@ -0,0 +1,208 @@ +// The record of who someone is, shared across every campaign and project. +// +// Prospects are unique per project, so before this the same person became a +// separate row in every project that found them — separate details, separate +// history, and nothing that noticed when two campaigns were about to email +// them in the same week. +// +// The merge rule is the whole point, and it is not "last write wins". +// Discovery fills gaps; it does not overwrite. A scraped company name must +// never replace one a human typed, and a second run finding a different phone +// number must not silently discard the first. Conflicting values are kept in +// `alternates` so a person can adjudicate, and every field records where it +// came from so a hand-entered value can be told from a scraped one. + +import { serviceClient } from "@/lib/supabase/service"; +import { normalizeEmail, normalizeHost } from "./cold"; + +/** Where a value came from, worst to best. Higher wins a conflict. */ +const SOURCE_RANK: Record = { + guess: 1, + search: 2, + page: 3, + "json-ld": 4, + manual: 10, +}; + +export type ContactFields = { + email?: string | null; + host?: string | null; + fullName?: string | null; + title?: string | null; + companyName?: string | null; + companySite?: string | null; + niche?: string | null; + industry?: string | null; + phone?: string | null; + postalAddress?: string | null; + linkedinUrl?: string | null; + country?: string | null; + sourceUrl?: string | null; + socials?: Record; +}; + +/** Column name per field, so the merge can be written once. */ +const COLUMN: Record = { + email: "email", + host: "host", + fullName: "full_name", + title: "title", + companyName: "company_name", + companySite: "company_site", + niche: "niche", + industry: "industry", + phone: "phone", + postalAddress: "postal_address", + linkedinUrl: "linkedin_url", + country: "country", + sourceUrl: "source_url", + socials: "socials", +}; + +function clean(v: unknown): string | null { + return typeof v === "string" && v.trim() ? v.trim() : null; +} + +/** + * Identity as the database computes it, so a lookup can find the row the + * unique index would collide with. + */ +function identityKey(fields: ContactFields): string | null { + const email = fields.email ? normalizeEmail(fields.email) : null; + if (email) return `email:${email.toLowerCase()}`; + const name = clean(fields.fullName); + if (!name) return null; + const norm = (s: string | null) => (s ?? "").toLowerCase().replace(/[^a-z0-9]+/g, ""); + return `name:${norm(name)}@${norm(clean(fields.companyName))}`; +} + +export type UpsertResult = { id: string; created: boolean; conflicts: string[] } | null; + +/** + * Record what we now know about a person, merging with what was already known. + * + * Returns null rather than throwing. A contact record is bookkeeping around + * work that already succeeded; failing to write one is not a reason to fail + * the research that produced it. + */ +export async function upsertContact(input: { + organizationId: string; + fields: ContactFields; + /** How this information was obtained. Decides who wins a conflict. */ + source: keyof typeof SOURCE_RANK; +}): Promise { + const key = identityKey(input.fields); + // Nothing to key on means nothing to deduplicate against, and a row that + // can never be matched again is worse than no row. + if (!input.organizationId || !key) return null; + + try { + const sb = serviceClient(); + const { data: existing } = await sb + .from("outreach_contacts") + .select("*") + .eq("organization_id", input.organizationId) + .eq("identity_key", key) + .maybeSingle(); + + const incomingRank = SOURCE_RANK[input.source] ?? 1; + const patch: Record = {}; + const conflicts: string[] = []; + const row = (existing ?? null) as Record | null; + const sources = ((row?.field_sources as Record) ?? {}) as Record; + const alternates = Array.isArray(row?.alternates) ? [...(row!.alternates as unknown[])] : []; + + for (const [field, column] of Object.entries(COLUMN) as [keyof ContactFields, string][]) { + const raw = input.fields[field]; + if (field === "socials") { + const incoming = (raw as Record | undefined) ?? {}; + if (!Object.keys(incoming).length) continue; + // Socials merge per network rather than replacing the map, so a run + // that finds only GitHub does not drop a known LinkedIn. + patch.socials = { ...((row?.socials as Record) ?? {}), ...incoming }; + continue; + } + + const value = field === "email" ? (raw ? normalizeEmail(String(raw)) : null) : clean(raw); + if (!value) continue; + const current = clean(row?.[column]); + + if (!current) { + patch[column] = value; + sources[column] = input.source; + continue; + } + if (current === value) continue; + + // Disagreement. The better-sourced value wins; the loser is kept + // rather than dropped, because "we saw something else" is information + // and a wrong overwrite is otherwise unrecoverable. + const currentRank = SOURCE_RANK[sources[column] ?? "guess"] ?? 1; + conflicts.push(column); + if (incomingRank > currentRank) { + alternates.push({ field: column, value: current, source: sources[column] ?? "unknown" }); + patch[column] = value; + sources[column] = input.source; + } else { + alternates.push({ field: column, value, source: input.source }); + } + } + + if (row) { + const { data } = await sb + .from("outreach_contacts") + .update({ + ...patch, + field_sources: sources, + alternates, + last_enriched_at: new Date().toISOString(), + }) + .eq("id", row.id as string) + .select("id") + .maybeSingle(); + return data ? { id: data.id as string, created: false, conflicts } : null; + } + + const { data } = await sb + .from("outreach_contacts") + .insert({ + organization_id: input.organizationId, + ...patch, + host: patch.host ?? (input.fields.host ? normalizeHost(input.fields.host) : null), + field_sources: sources, + last_enriched_at: new Date().toISOString(), + }) + .select("id") + .maybeSingle(); + return data ? { id: data.id as string, created: true, conflicts: [] } : null; + } catch { + return null; + } +} + +/** + * Is this person marked do-not-contact anywhere in the organization? + * + * The reason the record is shared: one person saying no should stop every + * campaign, not just the one they replied to. + */ +export async function isContactBlocked( + organizationId: string, + email: string, +): Promise { + if (!organizationId || !email) return false; + try { + const { data } = await serviceClient() + .from("outreach_contacts") + .select("do_not_contact") + .eq("organization_id", organizationId) + .eq("identity_key", `email:${normalizeEmail(email).toLowerCase()}`) + .maybeSingle(); + return Boolean(data?.do_not_contact); + } catch { + // A lookup failure must not become an implicit permission to send. + return false; + } +} + +export { identityKey }; diff --git a/lib/outreach/discover.ts b/lib/outreach/discover.ts index 94f4a0a6..f2c103e2 100644 --- a/lib/outreach/discover.ts +++ b/lib/outreach/discover.ts @@ -648,6 +648,8 @@ export async function discoverProspects(input: { organizationId?: string | null; }): Promise<{ prospects: DiscoveredProspect[]; + /** People named on the pages the seeds opened, address or not. */ + people: DiscoveredPerson[]; serpCalls: number; errors: string[]; /** Seeds that returned a login wall — the UI offers to store credentials for these. */ @@ -658,6 +660,7 @@ export async function discoverProspects(input: { const errors: string[] = []; const loginRequiredSeeds: string[] = []; const mineable = new Set(); + const people = new Map(); let serpCalls = 0; const queries = (input.queries ?? []).slice(0, 5); @@ -711,6 +714,9 @@ export async function discoverProspects(input: { // Only "waiting on the user" when we have nothing to try. A stored // credential that failed is a different problem and reads as an error. if (res.loginRequired && !credentials) loginRequiredSeeds.push(seedUrl); + for (const person of res.people) { + if (!people.has(person.fullName)) people.set(person.fullName, person); + } if (input.organizationId && credentials) { await recordSeedCredentialResult({ organizationId: input.organizationId, @@ -722,5 +728,11 @@ export async function discoverProspects(input: { for (const p of res.prospects) if (!merged.has(p.host)) merged.set(p.host, p); } - return { prospects: [...merged.values()].slice(0, limit), serpCalls, errors, loginRequiredSeeds }; + return { + prospects: [...merged.values()].slice(0, limit), + people: [...people.values()], + serpCalls, + errors, + loginRequiredSeeds, + }; } diff --git a/lib/outreach/pipeline.ts b/lib/outreach/pipeline.ts index 587ecaf4..8a6cb889 100644 --- a/lib/outreach/pipeline.ts +++ b/lib/outreach/pipeline.ts @@ -48,6 +48,7 @@ import { resolvePostalAddress } from "./postalAddress"; import { findContactViaSearch } from "./contactFallback"; import { loadProjectMailbox } from "./senderMailbox"; import { loadRecipientContext, recipientContextPrompt } from "./recipientContext"; +import { upsertContact } from "./contacts"; export type ProspectRow = { id: string; @@ -417,6 +418,45 @@ export async function researchProspect(input: { * score and no findings, which is honest: nothing was measured, and the * draft path for these campaigns doesn't claim otherwise. */ + +/** + * Write the person behind a prospect into the shared contact record. + * + * Resolves the organization from the project, because contacts are org-scoped + * while prospects are project-scoped — that difference is the entire reason + * this table exists. + */ +async function recordContact(input: { + projectId: string; + host: string; + email: string | null; + label?: string | null; + source: "page" | "search" | "guess" | "manual"; +}): Promise { + const { data: project } = await serviceClient() + .from("projects") + .select("organization_id") + .eq("id", input.projectId) + .maybeSingle(); + const organizationId = (project?.organization_id as string | null) ?? null; + if (!organizationId) return null; + + const res = await upsertContact({ + organizationId, + source: input.source, + fields: { + email: input.email, + host: input.host, + // The discovery label is what the listing called them, which is a + // company name far more often than a person's — so it is recorded as + // one rather than guessed into full_name. + companyName: input.label ?? null, + companySite: `https://${input.host}`, + }, + }); + return res?.id ?? null; +} + async function researchWithoutScan(input: { userId: string; projectId: string; @@ -469,6 +509,17 @@ async function researchWithoutScan(input: { } } + // Record the person before the prospect, so the prospect can point at + // them. This is what stops the same human becoming a separate row in every + // project that happens to find them. + const contactId = await recordContact({ + projectId: input.projectId, + host, + email: contact?.email ?? null, + label: input.discoveryLabel, + source: contact?.source === "manual" ? "manual" : contact?.source === "guess" ? "guess" : "page", + }); + const { data, error } = await serviceClient() .from("outreach_prospects") .upsert( @@ -483,6 +534,7 @@ async function researchWithoutScan(input: { discovery_label: input.discoveryLabel ?? null, contact_email: contact?.email ?? null, contact_source: contact?.source ?? null, + contact_id: contactId, // Crawled the site, then searched for the business, and still found // no address. There is nothing further to try, so it leaves the // funnel rather than sitting at "new" and being re-researched on diff --git a/supabase/migrations/20260728070000_contacts_identity_key.sql b/supabase/migrations/20260728070000_contacts_identity_key.sql new file mode 100644 index 00000000..79172a12 --- /dev/null +++ b/supabase/migrations/20260728070000_contacts_identity_key.sql @@ -0,0 +1,62 @@ +-- Let a contact exist before their email does. +-- +-- The table was keyed on email and required one, which assumed every person +-- arrives with an address. Person discovery does the opposite: a directory +-- gives a name, a title and a LinkedIn profile, and the address is what the +-- pipeline goes looking for afterwards. Under the old shape those people +-- could not be recorded at all, so each run rediscovered them from scratch. +-- +-- identity_key is what a row is deduplicated on. It is the email when there +-- is one, because two records with the same address are the same person by +-- definition. Otherwise it is the name and employer, normalised — weaker, +-- but the alternative is a fresh row for the same human on every run. +-- +-- Maintained by trigger rather than as a generated column: when an email is +-- finally found for a name-keyed contact, the key has to change from the +-- name form to the email form, and a generated column cannot be updated +-- through that transition without rewriting the row. + +alter table public.outreach_contacts + alter column email drop not null; + +alter table public.outreach_contacts + add column if not exists identity_key text; + +create or replace function public.outreach_contact_identity() +returns trigger +language plpgsql +as $$ +begin + if new.email is not null and btrim(new.email) <> '' then + new.identity_key := 'email:' || lower(btrim(new.email)); + elsif new.full_name is not null and btrim(new.full_name) <> '' then + -- Name plus employer, punctuation and case removed. Weak on its own, + -- which is why it is only used when there is no address to key on. + new.identity_key := 'name:' || + regexp_replace(lower(btrim(new.full_name)), '[^a-z0-9]+', '', 'g') || + '@' || + coalesce(regexp_replace(lower(btrim(new.company_name)), '[^a-z0-9]+', '', 'g'), ''); + else + new.identity_key := null; + end if; + return new; +end; +$$; + +drop trigger if exists outreach_contacts_set_identity on public.outreach_contacts; +create trigger outreach_contacts_set_identity + before insert or update on public.outreach_contacts + for each row execute function public.outreach_contact_identity(); + +-- Backfill before the unique index goes on. +update public.outreach_contacts set updated_at = updated_at; + +-- Replaces the email-only index: that one could not hold a person who has no +-- address, and would have treated every one of them as the same null. +drop index if exists outreach_contacts_org_email_idx; +create unique index if not exists outreach_contacts_org_identity_idx + on public.outreach_contacts(organization_id, identity_key) + where identity_key is not null; + +comment on column public.outreach_contacts.identity_key is + 'Dedupe key: email when known, otherwise normalised name+employer. Maintained by trigger so it can change form when an address is finally found.'; diff --git a/tests/contacts-merge.test.ts b/tests/contacts-merge.test.ts new file mode 100644 index 00000000..da08b9a5 --- /dev/null +++ b/tests/contacts-merge.test.ts @@ -0,0 +1,61 @@ +import { describe, it, expect } from "vitest"; +import { identityKey } from "@/lib/outreach/contacts"; + +// The merge itself talks to the database; what is testable in isolation, and +// what actually decides whether two records are the same human, is the key. +describe("identityKey", () => { + it("keys on email when there is one", () => { + expect(identityKey({ email: "Jane@Acme.TEST" })).toBe("email:jane@acme.test"); + }); + + it("treats addresses differing only in case as the same person", () => { + expect(identityKey({ email: "JANE@acme.test" })).toBe(identityKey({ email: "jane@acme.test" })); + }); + + it("prefers email over name even when both are present", () => { + // Two people can share a name; an address is definitionally one inbox. + const key = identityKey({ email: "jane@acme.test", fullName: "Jane Doe" }); + expect(key).toMatch(/^email:/); + }); + + it("falls back to name and employer when there is no address", () => { + // Person discovery finds exactly this: a name, a title, no email. Before + // the fallback these people could not be stored at all. + expect(identityKey({ fullName: "Jane Doe", companyName: "Acme Robotics" })).toBe( + "name:janedoe@acmerobotics", + ); + }); + + it("normalises punctuation and case out of the name key", () => { + expect(identityKey({ fullName: "Marc van Neerven", companyName: "Acme, Inc." })).toBe( + identityKey({ fullName: "MARC VAN NEERVEN", companyName: "Acme Inc" }), + ); + }); + + it("keeps two same-named people at different employers apart", () => { + expect(identityKey({ fullName: "Jane Doe", companyName: "Acme" })).not.toBe( + identityKey({ fullName: "Jane Doe", companyName: "Globex" }), + ); + }); + + it("keys a name with no employer rather than refusing it", () => { + expect(identityKey({ fullName: "Jane Doe" })).toBe("name:janedoe@"); + }); + + it("returns null when there is nothing to key on", () => { + // A row that can never be matched again is worse than no row: every run + // would insert another copy of it. + expect(identityKey({})).toBeNull(); + expect(identityKey({ companyName: "Acme" })).toBeNull(); + expect(identityKey({ email: " " })).toBeNull(); + }); + + it("changes form when an address is finally found for a known name", () => { + // The reason the key is maintained by trigger rather than generated: + // enriching a name-keyed contact has to move it onto the email key. + const before = identityKey({ fullName: "Jane Doe", companyName: "Acme" }); + const after = identityKey({ fullName: "Jane Doe", companyName: "Acme", email: "jane@acme.test" }); + expect(before).toMatch(/^name:/); + expect(after).toMatch(/^email:/); + }); +});