diff --git a/app/actions/leads.ts b/app/actions/leads.ts index 9a2fe8c1..7a31e350 100644 --- a/app/actions/leads.ts +++ b/app/actions/leads.ts @@ -15,7 +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 { recordDiscoveredPeople } from "@/lib/outreach/contacts"; import { generatePitch } from "@/lib/outreach/generatePitch"; import { leadRunBilling, manualRunPrice } from "@/lib/outreach/billing"; import { runEmailCampaignTick, CAMPAIGN_COLUMNS, summarize, type CampaignRow } from "@/lib/outreach/runner"; @@ -105,26 +105,12 @@ export async function findLeadsAction(input: { // 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; - } - } + // The same recorder the campaign runner uses. Two copies of this is what + // let the runner quietly stop recording people at all. + const peopleRecorded = await recordDiscoveredPeople({ + organizationId: orgId, + people: found.people ?? [], + }); if (!found.prospects.length && !peopleRecorded) { // Nothing to show for it, so nothing to charge for it. diff --git a/lib/outreach/contacts.ts b/lib/outreach/contacts.ts index 596aa433..c9570d25 100644 --- a/lib/outreach/contacts.ts +++ b/lib/outreach/contacts.ts @@ -206,3 +206,60 @@ export async function isContactBlocked( } export { identityKey }; + +/** + * Record everybody a discovery run named. + * + * Shared because it was not, and the two callers drifted exactly as far apart + * as you would expect: the one-shot finder recorded people, the campaign + * runner read `prospects` and `errors` off the same result and dropped + * `people` on the floor. Every person found by a campaign — which is how a + * directory actually gets scraped — was discarded at the last step, after + * being rendered, paginated and parsed for. + * + * A person is worth recording without an address. The directory gives a name, + * a title and a profile; the address is what the pipeline goes looking for + * afterwards, and discarding the rest until it turns up means rediscovering + * the same human on every run. + */ +export async function recordDiscoveredPeople(input: { + organizationId: string | null; + people: Array<{ + fullName: string; + jobTitle?: string | null; + company?: string | null; + companySite?: string | null; + linkedinUrl?: string | null; + location?: string | null; + sourceUrl?: string | null; + socials?: Record; + source?: string; + }>; + /** What the campaign was looking for. The only niche available here. */ + niche?: string | null; +}): Promise { + if (!input.organizationId || !input.people?.length) return 0; + + let recorded = 0; + for (const person of input.people) { + const res = await upsertContact({ + organizationId: input.organizationId, + // Structured markup is a stronger claim than text scraped off a page, + // and the ranking is what decides who wins when the two disagree. + source: person.source === "json-ld" ? "json-ld" : "page", + fields: { + fullName: person.fullName, + title: person.jobTitle ?? null, + companyName: person.company ?? null, + companySite: person.companySite ?? null, + linkedinUrl: person.linkedinUrl ?? null, + country: person.location ?? null, + sourceUrl: person.sourceUrl ?? null, + socials: person.socials, + niche: input.niche ?? null, + }, + }); + if (res) recorded += 1; + } + return recorded; +} diff --git a/lib/outreach/runner.ts b/lib/outreach/runner.ts index ec35f0ab..28da64a5 100644 --- a/lib/outreach/runner.ts +++ b/lib/outreach/runner.ts @@ -25,6 +25,7 @@ import { } from "./pipeline"; import { nextStepReadyAt, type OutreachStep } from "./cold"; import { leadRunBilling, outOfCreditsNote } from "./billing"; +import { recordDiscoveredPeople } from "./contacts"; import { LEAD_RUN_CREDITS } from "@/lib/credits"; export type CampaignRow = { @@ -72,6 +73,8 @@ export type TickResult = { awaitingAuth: string[]; /** Credits this tick actually charged. Zero when it found nothing to do. */ creditsSpent: number; + /** People named by this tick and written to the shared contact record. */ + peopleRecorded: number; skipped: string[]; errors: string[]; }; @@ -99,6 +102,7 @@ export async function runEmailCampaignTick(campaign: CampaignRow): Promise ({ + serviceClient: () => ({ + from: () => { + // Tracks whether this chain has reached an insert, so the terminal + // maybeSingle can answer "no existing row" on a lookup and "here is the + // new id" on a write — which is what upsertContact distinguishes. + let inserted = false; + const chain: Record = { + select: () => chain, + eq: () => chain, + maybeSingle: async () => ({ data: inserted ? { id: "contact-1" } : null }), + insert: (row: unknown) => { + upserts.push(row); + inserted = true; + return chain; + }, + update: () => chain, + }; + return chain; + }, + }), +})); + +const { recordDiscoveredPeople } = await import("@/lib/outreach/contacts"); + +beforeEach(() => { + upserts.length = 0; +}); + +const cto = { + fullName: "Jane Doe", + jobTitle: "CTO", + company: "Acme", + linkedinUrl: "https://linkedin.com/in/janedoe", + location: "Austin", + sourceUrl: "https://ctodirectory.test/jane", +}; + +describe("recordDiscoveredPeople", () => { + it("records a person who has no address yet", async () => { + // The directory gives a name, a title and a profile and withholds the + // email. Waiting for an address before recording anything means finding + // the same human again on every run. + const n = await recordDiscoveredPeople({ organizationId: "org-1", people: [cto] }); + expect(n).toBe(1); + expect(upserts).toHaveLength(1); + expect(upserts[0]).toMatchObject({ + full_name: "Jane Doe", + title: "CTO", + company_name: "Acme", + linkedin_url: "https://linkedin.com/in/janedoe", + }); + }); + + it("stamps the campaign's niche onto the people it found", async () => { + await recordDiscoveredPeople({ organizationId: "org-1", people: [cto], niche: "CTOs" }); + expect(upserts[0]).toMatchObject({ niche: "CTOs" }); + }); + + it("does nothing without an organization to scope to", async () => { + expect(await recordDiscoveredPeople({ organizationId: null, people: [cto] })).toBe(0); + expect(upserts).toHaveLength(0); + }); + + it("does nothing when a run named nobody", async () => { + expect(await recordDiscoveredPeople({ organizationId: "org-1", people: [] })).toBe(0); + }); + + it("ranks structured markup above scraped text", async () => { + await recordDiscoveredPeople({ + organizationId: "org-1", + people: [{ ...cto, source: "json-ld" }], + }); + expect((upserts[0] as { field_sources: Record }).field_sources.full_name).toBe( + "json-ld", + ); + }); +}); + +describe("both discovery paths record people", () => { + // The runner read `prospects` and `errors` off the discovery result and + // dropped `people` entirely, so every person found by a campaign was + // discarded after being rendered, paginated and parsed for. The finder did + // it correctly, which is precisely why nobody noticed. + const runner = readFileSync(new URL("../lib/outreach/runner.ts", import.meta.url), "utf8"); + const action = readFileSync(new URL("../app/actions/leads.ts", import.meta.url), "utf8"); + + it("the campaign runner records the people it names", () => { + expect(runner).toContain("recordDiscoveredPeople"); + expect(runner).toContain("found.people"); + }); + + it("the one-shot finder records them too", () => { + expect(action).toContain("recordDiscoveredPeople"); + }); + + it("both go through the one shared recorder", () => { + // Two copies is what let them drift the first time. + expect(runner).not.toMatch(/upsertContact\(/); + expect(action).not.toMatch(/upsertContact\(/); + }); +});