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
5 changes: 4 additions & 1 deletion apps/api/src/app.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3205,6 +3205,9 @@ export function createApp(options: AppOptions): Hono<AppEnv> {
company: z.string().optional(),
title: z.string().optional(),
location: z.string().optional(),
companyDomain: z.string().max(253).optional(),
linkedinUrl: z.string().max(500).optional(),
updatedAt: z.string().max(64).optional(),
}),
)
.max(IMPORT_CHUNK_MAX),
Expand Down Expand Up @@ -3296,7 +3299,7 @@ export function createApp(options: AppOptions): Hono<AppEnv> {
const batch = await queryOne<Record<string, unknown>>(
db,
`SELECT id, filename, consent_basis, consent_source, total_rows, imported, merged,
rejected, status, created_at
updated, rejected, status, created_at
FROM contact_imports WHERE id = ? AND workspace_id = ?`,
[importId, actor.workspaceId],
);
Expand Down
7 changes: 7 additions & 0 deletions apps/api/src/autogtm-docs.ts
Original file line number Diff line number Diff line change
Expand Up @@ -441,6 +441,7 @@ export const OPERATIONS: readonly Operation[] = [
company: { type: 'string' },
job_title: { type: 'string' },
location: { type: 'string' },
linkedin_url: { type: 'string' },
},
},
},
Expand All @@ -459,6 +460,11 @@ export const OPERATIONS: readonly Operation[] = [
type: 'integer',
description: 'Already on file; updated rather than duplicated.',
},
updated: {
type: 'integer',
description:
'Of the merged, how many people this import changed. Newer data wins: a value in the import replaces a different stored one, unless the row is dated older.',
},
rejected: { type: 'integer' },
crawls_queued: { type: 'integer' },
},
Expand All @@ -480,6 +486,7 @@ export const OPERATIONS: readonly Operation[] = [
total_rows: { type: 'integer' },
imported: { type: 'integer' },
merged: { type: 'integer' },
updated: { type: 'integer' },
rejected: { type: 'integer' },
},
},
Expand Down
11 changes: 10 additions & 1 deletion apps/api/src/autogtm.ts
Original file line number Diff line number Diff line change
Expand Up @@ -153,6 +153,7 @@ const importLead = z.object({
company: z.string().trim().max(200).optional(),
job_title: z.string().trim().max(200).optional(),
location: z.string().trim().max(200).optional(),
linkedin_url: z.string().trim().max(500).optional(),
});

const importBody = z.object({
Expand Down Expand Up @@ -769,6 +770,7 @@ export function autogtmRoutes(deps: AutogtmDeps): Hono<AppEnv> {

let imported = 0;
let merged = 0;
let updated = 0;
let rejected = 0;
const personIds = new Set<string>();

Expand All @@ -780,11 +782,14 @@ export function autogtmRoutes(deps: AutogtmDeps): Hono<AppEnv> {
...(lead.company ? { company: lead.company } : {}),
...(lead.job_title ? { title: lead.job_title } : {}),
...(lead.location ? { location: lead.location } : {}),
...(lead.company_domain ? { companyDomain: lead.company_domain } : {}),
...(lead.linkedin_url ? { linkedinUrl: lead.linkedin_url } : {}),
}));

const result = await importContactChunk(db, importId, rows, { startRow: offset });
imported += result.imported;
merged += result.merged;
updated += result.updated;
rejected += result.rejected;
for (const id of result.personIds) personIds.add(id);
}
Expand Down Expand Up @@ -846,6 +851,7 @@ export function autogtmRoutes(deps: AutogtmDeps): Hono<AppEnv> {
status: 'completed',
imported,
merged,
updated,
rejected,
crawls_queued: crawls,
},
Expand All @@ -862,12 +868,14 @@ export function autogtmRoutes(deps: AutogtmDeps): Hono<AppEnv> {
total_rows: number;
imported: number;
merged: number;
updated: number;
rejected: number;
created_at: string;
updated_at: string;
}>(
c.get('db'),
`SELECT id, campaign_id, status, total_rows, imported, merged, rejected, created_at, updated_at
`SELECT id, campaign_id, status, total_rows, imported, merged, updated, rejected, created_at,
updated_at
FROM contact_imports WHERE id = ? AND workspace_id = ?`,
[c.req.param('task_id'), actor.workspaceId],
);
Expand All @@ -881,6 +889,7 @@ export function autogtmRoutes(deps: AutogtmDeps): Hono<AppEnv> {
total_rows: row.total_rows,
imported: row.imported,
merged: row.merged,
updated: Number(row.updated ?? 0),
rejected: row.rejected,
created_at: row.created_at,
updated_at: row.updated_at,
Expand Down
113 changes: 112 additions & 1 deletion apps/api/src/contact-import.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@
* import at all.
*/

import { afterEach, describe, expect, test } from 'bun:test';
import { afterEach, describe, expect, setDefaultTimeout, test } from 'bun:test';
import type { Hono } from 'hono';
import { createApp } from './app';
import type { AppEnv, RequestActor } from './context';
Expand All @@ -22,6 +22,9 @@ const ACTOR: RequestActor = {
role: 'owner',
};

// Seeding takes 5-7 s on a loaded box; see crawl.test.ts.
setDefaultTimeout(30_000);

let active: SeededDatabase | undefined;

afterEach(() => {
Expand Down Expand Up @@ -283,4 +286,112 @@ describe('contact import', () => {
// anything like this against a local file, let alone a remote database.
expect(Date.now() - started).toBeLessThan(20_000);
});

test('a re-import patches known people with its newer data: newest wins', async () => {
const { app, seeded } = await harness('import-patch-newest');
const first = await startImport(app);
await post(app, `/contacts/imports/${first}/rows`, {
rows: [
{ email: 'dave.mackenzie@corp.com', name: 'Dave Mackenzie', title: 'Engineer' },
{ email: 'admin@theirstartup.com', name: 'Ana Ruiz' },
{ email: 'kim@corp.com', name: 'Kim Lee', title: 'CTO' },
],
});

// The enriched export, a while later.
const second = await startImport(app);
const response = await post(app, `/contacts/imports/${second}/rows`, {
rows: [
{
email: 'dave.mackenzie@corp.com',
name: 'Dave Mackenzie',
title: 'VP Engineering',
location: 'Oakland, CA',
companyDomain: 'https://www.corp.com/about',
linkedinUrl: 'linkedin.com/in/DaveMackenzie/',
},
// Dated before we last touched Ana: only blanks are filled.
{
email: 'admin@theirstartup.com',
name: 'Ana R',
title: 'Founder',
updatedAt: '2001-01-01',
},
// Nothing new: an empty cell never erases the stored title.
{ email: 'kim@corp.com', name: 'Kim Lee' },
],
});
const body = (await response.json()) as { merged: number; updated: number };
expect(body.merged).toBe(3);
expect(body.updated).toBe(2);

const people = await seeded.db.execute({
sql: `SELECT e.address, p.display_name, p.current_title, p.location, c.domain
FROM person_emails e JOIN people p ON p.id = e.person_id
LEFT JOIN companies c ON c.id = p.current_company_id
WHERE e.workspace_id = ? ORDER BY e.address`,
args: [SEED.workspaceId],
});
const by = Object.fromEntries(people.rows.map((r) => [String(r.address), r]));
expect(by['dave.mackenzie@corp.com']).toMatchObject({
current_title: 'VP Engineering',
location: 'Oakland, CA',
domain: 'corp.com',
});
expect(by['admin@theirstartup.com']).toMatchObject({
display_name: 'Ana Ruiz',
current_title: 'Founder',
});
expect(by['kim@corp.com']).toMatchObject({ current_title: 'CTO' });

const linkedin = await seeded.db.execute({
sql: "SELECT profile_url FROM social_identities WHERE network = 'linkedin'",
args: [],
});
expect(linkedin.rows.map((r) => r.profile_url)).toEqual([
'https://www.linkedin.com/in/davemackenzie',
]);

// Importing the same enriched file again changes nothing more.
const third = await startImport(app);
const again = (await (
await post(app, `/contacts/imports/${third}/rows`, {
rows: [
{
email: 'dave.mackenzie@corp.com',
name: 'Dave Mackenzie',
title: 'VP Engineering',
location: 'Oakland, CA',
companyDomain: 'corp.com',
linkedinUrl: 'https://www.linkedin.com/in/davemackenzie',
},
],
})
).json()) as { updated: number };
expect(again.updated).toBe(0);
});

test('every known row in a chunk is patched, not the first fifty', async () => {
const { app, seeded } = await harness('import-patch-all');
const rows = Array.from({ length: 120 }, (_, i) => ({
email: `person${i}@acme-${i}.com`,
name: `Person Number${i}`,
}));
const first = await startImport(app);
await post(app, `/contacts/imports/${first}/rows`, { rows });

const second = await startImport(app);
const body = (await (
await post(app, `/contacts/imports/${second}/rows`, {
rows: rows.map((r) => ({ ...r, title: 'Head of Growth' })),
})
).json()) as { merged: number; updated: number };
expect(body).toMatchObject({ merged: 120, updated: 120 });

const titled = await seeded.db.execute({
sql: "SELECT count(*) AS n FROM people WHERE current_title = 'Head of Growth'",
args: [],
});
expect(Number(titled.rows[0]?.n)).toBe(120);
});
});
9 changes: 7 additions & 2 deletions apps/web/components/contact-import.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,8 @@ const CHUNK = 500;
interface ImportSummary {
readonly imported: number;
readonly merged: number;
/** Of the merged, how many the file's newer data changed. */
readonly updated?: number;
readonly rejected: number;
}

Expand Down Expand Up @@ -116,7 +118,7 @@ export function ContactImport() {
}

const { importId } = (await started.json()) as { importId: string };
const totals = { imported: 0, merged: 0, rejected: 0 };
const totals = { imported: 0, merged: 0, updated: 0, rejected: 0 };

// Sequentially, so the server sees a steady trickle rather than
// thirty-four simultaneous writes, and so progress means something.
Expand All @@ -135,6 +137,7 @@ export function ContactImport() {
const chunk = (await response.json()) as ImportSummary;
totals.imported += chunk.imported;
totals.merged += chunk.merged;
totals.updated += chunk.updated ?? 0;
totals.rejected += chunk.rejected;

setProgress(Math.min(100, Math.round(((offset + CHUNK) / rows.length) * 100)));
Expand Down Expand Up @@ -238,7 +241,9 @@ export function ContactImport() {
<div className="border-border mt-4 rounded-xl border p-3 text-sm">
<p>
<span className="font-medium">{summary.imported.toLocaleString()}</span> added
{summary.merged > 0 ? `, ${summary.merged.toLocaleString()} already known` : ''}
{summary.merged > 0
? `, ${summary.merged.toLocaleString()} already known (${(summary.updated ?? 0).toLocaleString()} updated with newer data)`
: ''}
{summary.rejected > 0 ? `, ${summary.rejected.toLocaleString()} not usable` : ''}.
</p>

Expand Down
3 changes: 3 additions & 0 deletions apps/web/lib/csv.ts
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,9 @@ export interface MappedRow {
company?: string;
title?: string;
location?: string;
companyDomain?: string;
linkedinUrl?: string;
updatedAt?: string;
}

/** Applies a header mapping to the data rows. */
Expand Down
2 changes: 2 additions & 0 deletions migrations-pg/0050_contact_import_updated.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
-- 0050_contact_import_updated.sql (Postgres). See migrations/0050_contact_import_updated.sql.
ALTER TABLE contact_imports ADD COLUMN IF NOT EXISTS updated bigint NOT NULL DEFAULT 0;
7 changes: 7 additions & 0 deletions migrations/0050_contact_import_updated.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
-- 0050_contact_import_updated.sql: how many already-known people an import changed.
--
-- A re-import used to be a no-op past fifty rows a chunk and only ever filled
-- blanks. It now patches every known person with the row's newer data (see
-- planPatch in packages/pipeline/src/contact-import.ts), and this counts them,
-- so "6,204 already known" can say how many of those were actually updated.
ALTER TABLE contact_imports ADD COLUMN updated INTEGER NOT NULL DEFAULT 0;
48 changes: 48 additions & 0 deletions packages/domain/src/contact-import.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import {
cleanContact,
emailDedupeKey,
isFreemailDomain,
linkedinProfile,
mapHeaders,
nameFromEmail,
} from './contact-import';
Expand Down Expand Up @@ -170,3 +171,50 @@ describe('mapHeaders', () => {
expect(mapHeaders(['email', 'secondary email']).email).toBe(0);
});
});

describe('enriched columns', () => {
test('headers from an enriched export map onto the new fields', () => {
expect(
mapHeaders([
'email',
'first_name',
'company_domain',
'job_title',
'linkedin_url',
'enriched_at',
]),
).toMatchObject({
email: 0,
firstName: 1,
companyDomain: 2,
title: 3,
linkedinUrl: 4,
updatedAt: 5,
});
});

test('LinkedIn profiles are normalised; company pages are not a person', () => {
expect(linkedinProfile('linkedin.com/in/DaveMackenzie/?trk=x')).toBe(
'https://www.linkedin.com/in/davemackenzie',
);
expect(linkedinProfile('https://www.linkedin.com/company/acme')).toBeUndefined();
expect(linkedinProfile('not a url')).toBeUndefined();
});

test('cleanContact keeps a company domain, a profile and a date', () => {
const result = cleanContact({
email: 'dave@corp.com',
name: 'Dave Mackenzie',
companyDomain: 'https://www.corp.com/about',
linkedinUrl: 'https://linkedin.com/in/dave',
updatedAt: '2026-10-01',
});
expect(result.ok && result.contact).toMatchObject({
companyDomain: 'corp.com',
linkedinUrl: 'https://www.linkedin.com/in/dave',
updatedAt: '2026-10-01T00:00:00.000Z',
});
const board = cleanContact({ email: 'x@corp.com', companyDomain: 'linkedin.com' });
expect(board.ok && board.contact.companyDomain).toBeUndefined();
});
});
Loading
Loading