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:/); + }); +});