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
3 changes: 2 additions & 1 deletion .github/workflows/deploy-dev2.yml
Original file line number Diff line number Diff line change
Expand Up @@ -52,9 +52,10 @@ jobs:
DEV2_HOST: ${{ secrets.DEV2_HOST }}
VALUESERP_API_KEY: ${{ secrets.VALUESERP_API_KEY }}
CHOVY_CAMPAIGN_SECRET: ${{ secrets.CHOVY_CAMPAIGN_SECRET }}
PDL_API_KEY: ${{ secrets.PDL_API_KEY }}
run: |
lines=""
for key in VALUESERP_API_KEY CHOVY_CAMPAIGN_SECRET; do
for key in VALUESERP_API_KEY CHOVY_CAMPAIGN_SECRET PDL_API_KEY; do
value=$(printenv "$key" || true)
[ -n "$value" ] && lines="$lines$key=$value"$'\n'
done
Expand Down
8 changes: 8 additions & 0 deletions apps/api/src/app.ts
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,7 @@ import {
budgetStatus,
crawlDedupeKey,
screenHold,
type LeadEnrichDeps,
startContactImport,
importContactChunk,
intakeSocialPeople,
Expand Down Expand Up @@ -336,6 +337,12 @@ export interface AppOptions {
* people by company. Absent without a key, which leaves pasting URLs.
*/
readonly jobSearcher?: WebSearcher | undefined;
/**
* Lead enrichment (names from addresses, ValueSERP LinkedIn search, People
* Data Labs), for the on-demand run a campaign page can start. Absent: the
* route still fills names from addresses.
*/
readonly leadEnrichment?: Omit<LeadEnrichDeps, 'db'> | undefined;
/** Test seam for the job boards and company sites a posting is read from. */
readonly jobReader?: JobReaderOptions | undefined;
/**
Expand Down Expand Up @@ -1375,6 +1382,7 @@ export function createApp(options: AppOptions): Hono<AppEnv> {
approve: (db, actor, recommendation) =>
approveRecommendation(db, options, actor, recommendation, {}),
requireVerifiedEmail: (db, actor) => requireVerifiedEmail(db, actor),
...(options.leadEnrichment ? { leadEnrichment: options.leadEnrichment } : {}),
}),
);

Expand Down
55 changes: 55 additions & 0 deletions apps/api/src/append-leads.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -159,6 +159,61 @@ describe('POST /autogtm/campaigns/:id/leads', () => {
T,
);

test(
'enrich runs in the background and the status reports what was filled',
async () => {
const seeded = await seedDatabase('append-enrich');
active = seeded;
const asked: string[] = [];
const app = createApp({
db: seeded.db,
authenticate: async () => ACTOR,
leadEnrichment: {
searcher: {
async search(query) {
asked.push(query);
return query.includes('Scott Perry')
? [
{
title: 'Scott Perry - VP Sales - Northwind | LinkedIn',
link: 'https://www.linkedin.com/in/scott-perry-9',
},
]
: [];
},
},
},
});

await send(app, 'POST', `/autogtm/campaigns/${SEED.campaignId}/leads`, {
leads: [{ email: 'scott.perry@northwind.io', company_domain: 'northwind.io' }],
});

const started = await send(app, 'POST', `/autogtm/campaigns/${SEED.campaignId}/enrich`, {});
expect(started.status).toBe(202);

let status: Record<string, unknown> = { running: true };
for (let i = 0; i < 100 && status.running; i += 1) {
await Bun.sleep(100);
status = (await (
await app.request(`/api/v1/autogtm/campaigns/${SEED.campaignId}/enrichment`)
).json()) as Record<string, unknown>;
}
expect(status.running).toBe(false);
expect(status.last_run).toMatchObject({ titles: 1, profiles: 1 });

const scott = await queryOne<{ current_title: string; first_name: string }>(
seeded.db,
`SELECT p.current_title, p.first_name FROM people p
JOIN person_emails pe ON pe.person_id = p.id
WHERE pe.address = 'scott.perry@northwind.io'`,
);
expect(scott).toEqual({ current_title: 'VP Sales', first_name: 'Scott' });
expect(asked.some((q) => q.includes('"Scott Perry" northwind'))).toBe(true);
},
T,
);

test(
'takes the CSV itself and names a missing email column',
async () => {
Expand Down
45 changes: 45 additions & 0 deletions apps/api/src/autogtm-docs.ts
Original file line number Diff line number Diff line change
Expand Up @@ -601,6 +601,51 @@ export const OPERATIONS: readonly Operation[] = [
properties: { allow: { type: 'boolean' } },
},
},
{
method: 'post',
path: '/autogtm/campaigns/{campaign_id}/enrich',
id: 'enrichLeads',
tag: 'Campaigns',
summary: 'Fill in leads’ missing name, job title and LinkedIn',
description:
'Starts a run and answers 202 at once. Names come from unambiguous addresses ' +
'(first.last@), free. Job title and LinkedIn profile come from a Google search of LinkedIn, ' +
'taken only when the result carries every part of the name and the company; the company’s ' +
'LinkedIn page is the fallback. People Data Labs is used first when configured. Only blanks ' +
'are filled, every search is cached so a rerun costs nothing, and searches are capped per ' +
'run (`max_searches`) and per day. Read the outcome from GET .../enrichment.',
params: [idParam('campaign_id', 'The campaign.')],
body: {
type: 'object',
properties: {
max_searches: { type: 'integer', minimum: 0, maximum: 1000 },
max_leads: { type: 'integer', minimum: 1, maximum: 5000 },
},
},
status: 202,
},
{
method: 'get',
path: '/autogtm/campaigns/{campaign_id}/enrichment',
id: 'getEnrichment',
tag: 'Campaigns',
summary: 'What a campaign’s leads are missing, and the last enrichment run',
params: [idParam('campaign_id', 'The campaign.')],
response: {
type: 'object',
properties: {
leads: { type: 'integer' },
missing_name: { type: 'integer' },
missing_title: { type: 'integer' },
missing_linkedin: { type: 'integer' },
looked_up: { type: 'integer' },
providers: { type: 'object' },
today: { type: 'object' },
running: { type: 'boolean' },
last_run: { type: 'object' },
},
},
},
{
method: 'get',
path: '/autogtm/campaigns/import/{task_id}',
Expand Down
87 changes: 85 additions & 2 deletions apps/api/src/autogtm.ts
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,9 @@ import {
domainMatchKey,
emailMatchKey,
enqueue,
enrichLeads,
enrichmentStatus,
type LeadEnrichDeps,
finishContactImport,
importContactChunk,
normaliseDomain,
Expand Down Expand Up @@ -101,8 +104,13 @@ export interface AutogtmDeps {
) => Promise<ApproveResult>;
/** Refuses an unverified account anything that reaches a stranger. */
readonly requireVerifiedEmail: (db: Client, actor: RequestActor) => Promise<void>;
/** The enrichment providers; absent means names from addresses only. */
readonly leadEnrichment?: Omit<LeadEnrichDeps, 'db'>;
}

/** Campaigns with an on-demand enrichment run in flight, and the last result of each. */
const enrichRuns = new Map<string, { running: boolean; startedAt: string; last?: unknown }>();

// ------------------------------------------------------------------ schemas

const budgetBody = z.object({
Expand Down Expand Up @@ -219,6 +227,12 @@ const appendBody = z

const allowBody = z.object({ allow: z.boolean() });

const enrichBody = z.object({
/** New (uncached) searches this run may spend; the daily cap still applies. */
max_searches: z.number().int().min(0).max(1_000).optional(),
max_leads: z.number().int().min(1).max(5_000).optional(),
});

// ------------------------------------------------------------------ rows

interface CampaignRow {
Expand Down Expand Up @@ -877,7 +891,7 @@ export function autogtmRoutes(deps: AutogtmDeps): Hono<AppEnv> {
if (c.req.query('format') === 'csv') return importReportCsv(c, report);
return c.json({
...report,
rows: report.rows.map((row) => ({ ...row, why: reasonText(row.reason) })),
rows: report.rows.map((row) => ({ ...row, why: reasonText(row.reason, row.outcome) })),
});
});

Expand Down Expand Up @@ -957,6 +971,75 @@ export function autogtmRoutes(deps: AutogtmDeps): Hono<AppEnv> {
});
});

r.get('/campaigns/:id/enrichment', async (c) => {
const actor = c.get('actor');
const db = c.get('db');
const campaign = await ownedCampaign(db, actor.workspaceId, c.req.param('id'));
const status = await enrichmentStatus(
{ db, ...deps.leadEnrichment },
{ workspaceId: actor.workspaceId, campaignId: campaign.id },
);
const run = enrichRuns.get(campaign.id);
return c.json({
campaign_id: campaign.id,
...status,
running: run?.running === true,
...(run?.last ? { last_run: run.last } : {}),
});
});

/**
* Starts a run over this campaign's leads and answers at once: searches take
* 10-45 s each, so the run continues after the response and its outcome is
* read back from GET .../enrichment.
*/
r.post('/campaigns/:id/enrich', async (c) => {
const actor = c.get('actor');
const db = c.get('db');
requireWriter(actor, 'enriching leads');
const campaign = await ownedCampaign(db, actor.workspaceId, c.req.param('id'));
const body = await parse(c.req.raw, enrichBody);

const current = enrichRuns.get(campaign.id);
if (current?.running) {
return c.json({ campaign_id: campaign.id, started: false, running: true }, 202);
}
const startedAt = now();
enrichRuns.set(campaign.id, {
running: true,
startedAt,
...(current?.last ? { last: current.last } : {}),
});
void enrichLeads(
{ db, ...deps.leadEnrichment },
{
workspaceId: actor.workspaceId,
campaignId: campaign.id,
limit: body.max_leads ?? 200,
maxSearches: body.max_searches ?? 100,
},
)
.then((result) => {
enrichRuns.set(campaign.id, {
running: false,
startedAt,
last: { ...result, started_at: startedAt, finished_at: now() },
});
})
.catch((error: unknown) => {
enrichRuns.set(campaign.id, {
running: false,
startedAt,
last: {
error: error instanceof Error ? error.message : String(error),
started_at: startedAt,
},
});
});

return c.json({ campaign_id: campaign.id, started: true, running: true }, 202);
});

r.post('/leads/:person_id/screening', async (c) => {
const actor = c.get('actor');
const db = c.get('db');
Expand Down Expand Up @@ -1828,7 +1911,7 @@ async function importIntoCampaign(db: Client, actor: RequestActor, input: Import
crawls_queued: crawls,
report: reportRows
.slice(0, INLINE_REPORT_ROWS)
.map((row) => ({ ...row, why: reasonText(row.reason) })),
.map((row) => ({ ...row, why: reasonText(row.reason, row.outcome) })),
report_truncated: reportRows.length > INLINE_REPORT_ROWS,
report_url: `/api/v1/autogtm/campaigns/import/${importId}/report?format=csv`,
};
Expand Down
20 changes: 17 additions & 3 deletions apps/api/src/import-report.ts
Original file line number Diff line number Diff line change
Expand Up @@ -55,11 +55,25 @@ const REASON_TEXT: Readonly<Record<string, string>> = {
suppressed: 'on a suppress list',
};

/**
* Screening's words for its flags. `role_address` means something different
* there: the cleaner rejects a mailbox nobody reads (noreply@), screening
* holds back one a team reads (info@).
*/
const FLAG_TEXT: Readonly<Record<string, string>> = {
generated_name: 'a generated name, held back from sending',
relay_address: 'a relay address that hides the person, held back from sending',
temp_mail_domain: 'a temp-mail domain, held back from sending',
agent_account: 'a bot, agent or test account, held back from sending',
role_address: 'a team inbox, not a person, held back from sending',
};

/** A reason as words, for the UI and the CSV `why` column. */
export function reasonText(reason: string): string {
export function reasonText(reason: string, outcome?: string): string {
const words = outcome === 'flagged' ? { ...REASON_TEXT, ...FLAG_TEXT } : REASON_TEXT;
return reason
.split('+')
.map((part) => REASON_TEXT[part] ?? part.replace(/_/g, ' '))
.map((part) => words[part] ?? part.replace(/_/g, ' '))
.join('; ');
}

Expand Down Expand Up @@ -130,7 +144,7 @@ export function importReportCsv(c: Context, report: ImportReport): Response {
email: row.email,
outcome: row.outcome,
reason: row.reason,
why: reasonText(row.reason),
why: reasonText(row.reason, row.outcome),
detail: row.detail,
})),
['row', 'email', 'outcome', 'reason', 'why', 'detail'],
Expand Down
37 changes: 35 additions & 2 deletions apps/cli/src/commands.ts
Original file line number Diff line number Diff line change
Expand Up @@ -273,6 +273,39 @@ async function runLeads({ client, args, flags }: CommandContext): Promise<string
].join('\n');
}

if (verb === 'enrich' || verb === 'enrichment') {
if (!target) throw new Error(`usage: og leads ${verb} <campaignId>`);
const base = `/autogtm/campaigns/${encodeURIComponent(target)}`;
if (verb === 'enrich') {
const max = flagString(flags, 'max');
await client.post(`${base}/enrich`, max ? { max_searches: Number(max) } : {});
}
const status = (await client.get(`${base}/enrichment`)) as Record<string, unknown>;
const today = (status.today ?? {}) as Record<string, unknown>;
const last = status.last_run as Record<string, unknown> | undefined;
return [
verb === 'enrich'
? 'Started. It runs in the background; check: og leads enrichment ' + target
: '',
`${text(status, 'leads', '0')} leads: missing ${text(status, 'missing_title', '0')} titles, ` +
`${text(status, 'missing_linkedin', '0')} LinkedIn, ${text(status, 'missing_name', '0')} names`,
`searches today: ${text(today, 'searches', '0')} of ${text(today, 'searches_cap', '0')}` +
`${status.running ? ' (a run is in progress)' : ''}`,
...(last
? [
last.error
? `last run failed: ${text(last, 'error')}`
: `last run: +${text(last, 'names', '0')} names, +${text(last, 'titles', '0')} titles, ` +
`+${text(last, 'profiles', '0')} profiles, +${text(last, 'companies', '0')} company pages, ` +
`${text(last, 'searches', '0')} searches (${text(last, 'cached', '0')} cached)` +
`${last.stopped ? `; stopped: ${text(last, 'stopped')}` : ''}`,
]
: []),
]
.filter(Boolean)
.join('\n');
}

if (verb === 'allow' || verb === 'hold') {
if (!target) throw new Error(`usage: og leads ${verb} <personId>`);
await client.post(`/autogtm/leads/${encodeURIComponent(target)}/screening`, {
Expand All @@ -283,7 +316,7 @@ async function runLeads({ client, args, flags }: CommandContext): Promise<string
: `${target} held back from sending again.`;
}

throw new Error('usage: og leads add|report|screened|allow|hold … (og help)');
throw new Error('usage: og leads add|report|screened|enrich|enrichment|allow|hold … (og help)');
}

async function runJobs({ client, args, flags }: CommandContext): Promise<string> {
Expand Down Expand Up @@ -1662,7 +1695,7 @@ export const COMMANDS: readonly Command[] = [
{
name: 'leads',
usage:
'og leads add <campaignId> <file.csv> --consent-source "<where>" [--allow-flagged] [--keep-project-duplicates] [--report out.csv] | report <taskId> [--csv] | screened <campaignId> [--all] | allow <personId> | hold <personId>',
'og leads add <campaignId> <file.csv> --consent-source "<where>" [--allow-flagged] [--keep-project-duplicates] [--report out.csv] | report <taskId> [--csv] | screened <campaignId> [--all] | enrich <campaignId> [--max <searches>] | enrichment <campaignId> | allow <personId> | hold <personId>',
summary: 'Add a CSV of leads to a running campaign, with a per-row report and screening.',
run: runLeads,
},
Expand Down
Loading
Loading