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
28 changes: 7 additions & 21 deletions app/actions/leads.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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.
Expand Down
57 changes: 57 additions & 0 deletions lib/outreach/contacts.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, string>;
source?: string;
}>;
/** What the campaign was looking for. The only niche available here. */
niche?: string | null;
}): Promise<number> {
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;
}
15 changes: 15 additions & 0 deletions lib/outreach/runner.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 = {
Expand Down Expand Up @@ -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[];
};
Expand Down Expand Up @@ -99,6 +102,7 @@ export async function runEmailCampaignTick(campaign: CampaignRow): Promise<TickR
autoSend: campaign.auto_send,
awaitingAuth: [],
creditsSpent: 0,
peopleRecorded: 0,
researched: 0,
drafted: 0,
sent: 0,
Expand Down Expand Up @@ -257,6 +261,16 @@ export async function runEmailCampaignTick(campaign: CampaignRow): Promise<TickR
result.errors.push(...found.errors);
result.awaitingAuth = found.loginRequiredSeeds;

// The people the run named, not just the companies. Reading prospects off
// this result and ignoring `people` is what threw away every name, title
// and profile link a directory gave up — after paying to render, paginate
// and parse for them.
result.peopleRecorded = await recordDiscoveredPeople({
organizationId,
people: found.people,
niche: campaign.name,
});

// Park the gated hosts on the campaign so the UI can say what it is
// waiting for, and offer the form that unblocks it, instead of leaving
// the reason buried in an error string.
Expand Down Expand Up @@ -450,6 +464,7 @@ export function summarize(r: TickResult): string {
? `${r.dryRuns} drafted, 0 sent`
: `${r.dryRuns} drafted (auto_send off)`,
];
if (r.peopleRecorded) parts.push(`${r.peopleRecorded} people`);
if (r.awaitingAuth.length) parts.push(`${r.awaitingAuth.length} waiting_for_auth`);
// Only when something was charged. Printing "0 credits" on every idle tick
// would bury the line that matters under the ones that cost nothing.
Expand Down
106 changes: 106 additions & 0 deletions tests/people-recording.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,106 @@
import { describe, it, expect, vi, beforeEach } from "vitest";
import { readFileSync } from "node:fs";

const upserts: unknown[] = [];
vi.mock("@/lib/supabase/service", () => ({
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<string, unknown> = {
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<string, string> }).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\(/);
});
});
Loading