Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 37 additions & 2 deletions app/actions/leads.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<T = Record<string, never>> = { ok: true } & T;
Expand Down Expand Up @@ -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." };
}

Expand Down Expand Up @@ -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.`,
};
}

Expand Down
208 changes: 208 additions & 0 deletions lib/outreach/contacts.ts
Original file line number Diff line number Diff line change
@@ -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<string, number> = {
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<string, string>;
};

/** Column name per field, so the merge can be written once. */
const COLUMN: Record<keyof ContactFields, string> = {
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<UpsertResult> {
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<string, unknown> = {};
const conflicts: string[] = [];
const row = (existing ?? null) as Record<string, unknown> | null;
const sources = ((row?.field_sources as Record<string, string>) ?? {}) as Record<string, string>;
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<string, string> | 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<string, string>) ?? {}), ...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<boolean> {
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 };
14 changes: 13 additions & 1 deletion lib/outreach/discover.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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. */
Expand All @@ -658,6 +660,7 @@ export async function discoverProspects(input: {
const errors: string[] = [];
const loginRequiredSeeds: string[] = [];
const mineable = new Set<string>();
const people = new Map<string, DiscoveredPerson>();
let serpCalls = 0;

const queries = (input.queries ?? []).slice(0, 5);
Expand Down Expand Up @@ -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,
Expand All @@ -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,
};
}
52 changes: 52 additions & 0 deletions lib/outreach/pipeline.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<string | null> {
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;
Expand Down Expand Up @@ -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(
Expand All @@ -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
Expand Down
Loading
Loading