diff --git a/apps/api/src/app.ts b/apps/api/src/app.ts index 2a8b776..4c6b303 100644 --- a/apps/api/src/app.ts +++ b/apps/api/src/app.ts @@ -16,14 +16,16 @@ import { connectEmailAccountSchema, createSuppressionSchema, executeActionSchema, + backfillDraftsSchema, + recordReplySchema, loginSchema, privacyRequestSchema, registerSchema, snoozeRecommendationSchema, workspaceProfileSchema, } from '@outreachgraph/contracts'; -import { newId, type ActionKind, type Network } from '@outreachgraph/domain'; -import { now, queryOne, type Client } from '@outreachgraph/db'; +import { newId, OUTBOUND_ACTION_KINDS, type ActionKind, type Network } from '@outreachgraph/domain'; +import { now, queryAll, queryOne, type Client } from '@outreachgraph/db'; import { actorFromSession, clearedCookie, @@ -932,6 +934,60 @@ export function createApp(options: AppOptions): Hono { return c.json({ signals }); }); + /** + * Records that this person replied, which takes them out of cold outreach. + * + * The product sends over SMTP and reads no mailbox, so it cannot notice a + * reply by itself. Until inbound polling exists, this route is how a reply + * becomes something the policy engine can act on: `conversation_open` then + * refuses further outreach, and a follow-up has to be an explicit human + * decision rather than the queue's next tick. + * + * Idempotent by intent rather than by constraint — recording a second reply + * is a real event, and the gate only asks whether any exist. + */ + api.post('/people/:id/replied', async (c) => { + const actor = c.get('actor'); + const db = c.get('db'); + const personId = c.req.param('id'); + const body = await parseBody(c.req.raw, recordReplySchema); + + const person = await repo.getPerson(db, personId); + if (!person) throw ApiError.notFound('person'); + + const contact = await repo.resolveContactAddress(db, personId); + const at = body.occurredAt ?? now(); + const address = body.fromAddress?.trim().toLowerCase() ?? contact?.address ?? null; + + await db.execute({ + sql: `INSERT INTO interactions (id, workspace_id, person_id, network, direction, + state, body, contact_address, shared_inbox, occurred_at, recorded_at) + VALUES (?, ?, ?, 'email', 'inbound', 'replied', ?, ?, ?, ?, ?)`, + args: [ + newId('interaction'), + actor.workspaceId, + personId, + body.body ?? null, + address, + contact?.shared ? 1 : 0, + at, + now(), + ], + }); + + await repo.audit(db, { + workspaceId: actor.workspaceId, + actorKind: 'user', + actorId: actor.userId, + eventType: 'interaction.reply_recorded', + entityKind: 'person', + entityId: personId, + detail: { address }, + }); + + return c.json({ recorded: true, personId, conversationOpen: true }); + }); + /** * Deletes a person and everything derived from them, leaving a suppression * tombstone so a later provider lookup cannot re-ingest them (PRD §17.3, @@ -1125,6 +1181,83 @@ export function createApp(options: AppOptions): Hono { }); }); + /** + * Writes the missing drafts for a queue that has gone unreviewed. + * + * Drafting has only ever happened one card at a time, on request, so a queue + * that filled faster than anyone clicked ends up as a wall of cards with + * nothing written on them — which is what the approvals page had become: 74 + * pending recommendations, zero drafts. + * + * Bounded rather than exhaustive, and sequential rather than parallel. Each + * card is a model call, so an unbounded version of this route is an + * unbounded invoice; `limit` is what the caller is willing to spend and the + * response says exactly what was written, refused and skipped. + */ + api.post('/recommendations/drafts/backfill', async (c) => { + const actor = c.get('actor'); + const db = c.get('db'); + const body = await parseBody(c.req.raw, backfillDraftsSchema); + + if (!options.model) { + throw new ApiError( + 503, + 'composer_unavailable', + 'no language model is configured; set ANTHROPIC_API_KEY, OPENAI_API_KEY or GEMINI_API_KEY to enable drafting', + ); + } + + // Only actions that put a message in front of a human. Most of a real + // queue is `refresh_research` and friends, which have no message by + // definition — production had 75 of those against 2 genuinely undrafted + // emails, so an unfiltered version of this route would have spent 75 model + // calls to write nothing and reported them all as refusals. + const outbound = OUTBOUND_ACTION_KINDS.map(() => '?').join(', '); + + const pending = await queryAll<{ id: string }>( + db, + `SELECT r.id FROM recommendations r + WHERE r.workspace_id = ? AND r.status IN ('pending', 'approved') + AND r.action IN (${outbound}) + AND NOT EXISTS (SELECT 1 FROM drafts d WHERE d.recommendation_id = r.id) + ORDER BY r.created_at ASC + LIMIT ?`, + [actor.workspaceId, ...OUTBOUND_ACTION_KINDS, body.limit], + ); + + const written: string[] = []; + const withheld: { id: string; reason: string }[] = []; + + for (const row of pending) { + // One failure is not the batch's failure. A model that refuses this card + // will happily write the next, and stopping here would leave the queue + // as empty as it was found. + try { + const result = await draftForRecommendation(db, options.model, row.id); + if (result.ok) written.push(row.id); + else withheld.push({ id: row.id, reason: result.reason ?? 'withheld' }); + } catch (error) { + withheld.push({ id: row.id, reason: error instanceof Error ? error.message : 'failed' }); + } + } + + await repo.audit(db, { + workspaceId: actor.workspaceId, + actorKind: 'user', + actorId: actor.userId, + eventType: 'drafts.backfilled', + entityKind: 'workspace', + entityId: actor.workspaceId, + detail: { considered: pending.length, written: written.length, withheld: withheld.length }, + }); + + return c.json({ + considered: pending.length, + written: written.length, + withheld, + }); + }); + /** * Compose (or recompose) the message for a recommendation. * @@ -1303,6 +1436,16 @@ export function createApp(options: AppOptions): Hono { } const stamp = now(); + + // A message a human sent by hand still landed in someone's mailbox, so it + // is recorded against the same address the automated path would have used. + // Leaving it null here would let the queue offer that inbox again tomorrow + // on the grounds that the product had never written to it. + const manualContact = + action.network === 'email' + ? await repo.resolveContactAddress(db, action.person_id) + : undefined; + await db.batch([ { sql: `UPDATE actions SET status = 'completed', mode = ?, external_url = ?, executed_at = ? @@ -1311,14 +1454,16 @@ export function createApp(options: AppOptions): Hono { }, { sql: `INSERT INTO interactions (id, workspace_id, person_id, action_id, network, direction, - state, occurred_at, recorded_at) - VALUES (?, ?, ?, ?, ?, 'outbound', 'contacted', ?, ?)`, + state, contact_address, shared_inbox, occurred_at, recorded_at) + VALUES (?, ?, ?, ?, ?, 'outbound', 'contacted', ?, ?, ?, ?)`, args: [ newId('interaction'), actor.workspaceId, action.person_id, action.id, action.network, + manualContact?.address ?? null, + manualContact?.shared ? 1 : 0, stamp, stamp, ], @@ -1893,13 +2038,16 @@ async function recheckPolicy( /** The action being executed, which must not count against its own limit. */ excludeActionId?: string, ) { - const [person, workspace, campaign, counts, flags, connected] = await Promise.all([ + const [person, workspace, campaign, counts, flags, connected, contact] = await Promise.all([ repo.getPerson(db, recommendation.person_id), repo.getWorkspace(db, actor.workspaceId), repo.getCampaign(db, actor.workspaceId, recommendation.campaign_id), repo.actionCounts(db, actor.workspaceId, recommendation.person_id, excludeActionId), repo.featureFlags(db, actor.workspaceId), repo.hasConnectedAccount(db, actor.workspaceId, recommendation.network as Network), + recommendation.network === 'email' + ? repo.resolveContactAddress(db, recommendation.person_id) + : Promise.resolve(undefined), ]); if (!person) throw ApiError.notFound('person'); @@ -1909,6 +2057,15 @@ async function recheckPolicy( const suppressed = person.status === 'suppressed' || (await repo.isSuppressed(db, actor.workspaceId, matchKeys)); + // Counted against the mailbox rather than the person. Skipped entirely when + // no address resolves — there is nothing to protect and nothing to count. + const [addressUsage, replied] = await Promise.all([ + contact + ? repo.addressCounts(db, actor.workspaceId, contact.address, excludeActionId) + : Promise.resolve(undefined), + repo.conversationOpen(db, actor.workspaceId, recommendation.person_id, contact?.address), + ]); + const budget = safeJson(campaign.budget_json); return evaluatePolicy({ @@ -1928,6 +2085,17 @@ async function recheckPolicy( ...(counts.hoursSinceLast === undefined ? {} : { hoursSinceLastActionToProspect: counts.hoursSinceLast }), + ...(addressUsage === undefined + ? {} + : { + actionsToThisAddressThisWeek: addressUsage.thisWeek, + maxActionsPerAddressPerWeek: numberOr(budget.maxActionsPerAddressPerWeek, 1), + addressShared: contact?.shared === true, + ...(addressUsage.hoursSinceLast === undefined + ? {} + : { hoursSinceLastActionToAddress: addressUsage.hoursSinceLast }), + }), + conversationOpen: replied, featureFlags: flags, }); } diff --git a/apps/api/src/drafts-backfill.test.ts b/apps/api/src/drafts-backfill.test.ts new file mode 100644 index 0000000..5795a95 --- /dev/null +++ b/apps/api/src/drafts-backfill.test.ts @@ -0,0 +1,148 @@ +/** + * Filling in the drafts a queue is missing, without paying for the ones it is + * not missing. + * + * The approvals page looked like a wall of cards with nothing written on them, + * and the obvious reading — "drafting is broken" — was wrong. Production had + * 77 recommendations without a draft, of which 75 were `refresh_research`: + * internal actions that have no message by definition and never will. Only 2 + * were genuinely undrafted emails. + * + * So the expensive mistake here is not failing to draft. It is drafting + * everything: an unfiltered backfill spends a model call per research card to + * produce nothing, and reports each one as a refusal. + */ + +import { afterEach, describe, expect, test } from 'bun:test'; +import { now, type Client } from '@outreachgraph/db'; +import type { GenerateInput, GenerateResult, TextModel } from '@outreachgraph/ai'; +import { createApp } from './app'; +import type { RequestActor } from './context'; +import { seedDatabase, SEED, type SeededDatabase } from './test-seed'; + +const ACTOR: RequestActor = { + userId: SEED.userId, + workspaceId: SEED.workspaceId, + organizationId: SEED.organizationId, + role: 'owner', +}; + +let active: SeededDatabase | undefined; + +afterEach(() => { + active?.cleanup(); + active = undefined; +}); + +/** Counts calls, so "did not draft it" is a measurable claim rather than a hope. */ +function countingModel(): { calls: GenerateInput[]; model: TextModel } { + const calls: GenerateInput[] = []; + return { + calls, + model: { + generate: async (input): Promise => { + calls.push(input); + return { + text: 'A short grounded note.', + model: 'test', + inputTokens: 0, + outputTokens: 0, + cachedTokens: 0, + refused: false, + }; + }, + }, + }; +} + +/** A recommendation with no draft, of whichever kind the caller asks for. */ +async function addUndrafted( + db: Client, + id: string, + action: string, + network: string, +): Promise { + await db.execute({ + sql: `INSERT INTO recommendations (id, workspace_id, campaign_id, person_id, action, + network, priority, reason, trigger_signal_id, policy_status, policy_version, + expected_goal, status, created_at) + VALUES (?, ?, ?, ?, ?, ?, 80, 'Queued.', ?, 'allow_with_approval', '2026-08-11', + 'start_conversation', 'pending', ?)`, + args: [ + id, + SEED.workspaceId, + SEED.campaignId, + SEED.personId, + action, + network, + SEED.signalId, + now(), + ], + }); +} + +describe('backfilling missing drafts', () => { + test('ignores internal research cards, which have no message to write', async () => { + const seeded = await seedDatabase('backfill-skips-research'); + active = seeded; + const { calls, model } = countingModel(); + const app = createApp({ db: seeded.db, authenticate: async () => ACTOR, model }); + + for (let i = 0; i < 5; i += 1) { + await addUndrafted(seeded.db, `rec_research_${i}`, 'refresh_research', 'website'); + } + + const response = await app.request('/api/v1/recommendations/drafts/backfill', { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ limit: 25 }), + }); + + expect(response.status).toBe(200); + const payload = (await response.json()) as { considered: number; written: number }; + + // Five cards the page shows as blank, and not one model call for them. + expect(payload.considered).toBe(0); + expect(calls).toHaveLength(0); + }); + + test('writes the outbound ones that really are missing a draft', async () => { + const seeded = await seedDatabase('backfill-writes-email'); + active = seeded; + const { calls, model } = countingModel(); + const app = createApp({ db: seeded.db, authenticate: async () => ACTOR, model }); + + await addUndrafted(seeded.db, 'rec_email_1', 'send_email', 'email'); + await addUndrafted(seeded.db, 'rec_research_1', 'refresh_research', 'website'); + + const response = await app.request('/api/v1/recommendations/drafts/backfill', { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ limit: 25 }), + }); + + const payload = (await response.json()) as { considered: number; written: number }; + expect(payload.considered).toBe(1); + expect(calls.length).toBeGreaterThan(0); + }); + + test('respects the limit, because every card is a model call', async () => { + const seeded = await seedDatabase('backfill-limit'); + active = seeded; + const { model } = countingModel(); + const app = createApp({ db: seeded.db, authenticate: async () => ACTOR, model }); + + for (let i = 0; i < 6; i += 1) { + await addUndrafted(seeded.db, `rec_email_${i}`, 'send_email', 'email'); + } + + const response = await app.request('/api/v1/recommendations/drafts/backfill', { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ limit: 2 }), + }); + + const payload = (await response.json()) as { considered: number }; + expect(payload.considered).toBe(2); + }); +}); diff --git a/apps/api/src/repository.ts b/apps/api/src/repository.ts index 475dd7b..aa54d12 100644 --- a/apps/api/src/repository.ts +++ b/apps/api/src/repository.ts @@ -234,6 +234,123 @@ export async function actionCounts( }; } +export interface ContactAddress { + readonly address: string; + /** True when this is the company's inbox rather than the person's own. */ + readonly shared: boolean; +} + +/** + * The address an email to this person would actually be delivered to. + * + * Mirrors `pickEmailRecipient` in `@outreachgraph/pipeline` deliberately: the + * policy engine has to count against the same mailbox the sender will use, and + * if the two ever disagree the limits are protecting an address nobody writes + * to. Personal identity first, company inbox second, nothing third. + */ +export async function resolveContactAddress( + db: Client, + personId: string, +): Promise { + const personal = await queryOne<{ handle: string }>( + db, + `SELECT handle FROM social_identities + WHERE person_id = ? AND network = 'email' AND handle IS NOT NULL AND trim(handle) <> '' + ORDER BY confidence DESC LIMIT 1`, + [personId], + ); + if (personal?.handle) return { address: personal.handle.trim().toLowerCase(), shared: false }; + + const company = await queryOne<{ contact_email: string }>( + db, + `SELECT co.contact_email FROM people p + JOIN companies co ON co.id = p.current_company_id + WHERE p.id = ? AND co.contact_email IS NOT NULL AND trim(co.contact_email) <> ''`, + [personId], + ); + if (company?.contact_email) { + return { address: company.contact_email.trim().toLowerCase(), shared: true }; + } + + return undefined; +} + +/** + * How much mail one address has had, regardless of who it was addressed to. + * + * Counted from `interactions` rather than `actions` because an interaction is + * the record of something that actually went out. An approved action that has + * not been sent has not reached anyone's inbox, and refusing on account of it + * would block the very send it is waiting for. + */ +export async function addressCounts( + db: Client, + workspaceId: string, + address: string, + /** The action being executed, which must not count against its own limit. */ + excludeActionId?: string, +): Promise<{ thisWeek: number; hoursSinceLast?: number }> { + const weekAgo = new Date(Date.now() - 7 * 86_400_000).toISOString(); + const exclude = excludeActionId ? 'AND (action_id IS NULL OR action_id != ?)' : ''; + const excludeArgs = excludeActionId ? [excludeActionId] : []; + + const week = await queryOne<{ n: number }>( + db, + `SELECT count(*) AS n FROM interactions + WHERE workspace_id = ? AND contact_address = ? AND direction = 'outbound' + AND occurred_at >= ? ${exclude}`, + [workspaceId, address, weekAgo, ...excludeArgs], + ); + + const last = await queryOne<{ occurred_at: string }>( + db, + `SELECT occurred_at FROM interactions + WHERE workspace_id = ? AND contact_address = ? AND direction = 'outbound' ${exclude} + ORDER BY occurred_at DESC LIMIT 1`, + [workspaceId, address, ...excludeArgs], + ); + + const hoursSinceLast = last ? (Date.now() - Date.parse(last.occurred_at)) / 3_600_000 : undefined; + + return { + thisWeek: Number(week?.n ?? 0), + ...(hoursSinceLast === undefined ? {} : { hoursSinceLast }), + }; +} + +/** + * Whether this contact has written back. + * + * Checked by person *and* by address: a reply from a shared inbox answers on + * behalf of everyone who was written to there, and continuing to mail + * colleagues after someone at the company has replied is the same mistake in + * a thinner disguise. + */ +export async function conversationOpen( + db: Client, + workspaceId: string, + personId: string, + address?: string, +): Promise { + const byPerson = await queryOne<{ n: number }>( + db, + `SELECT count(*) AS n FROM interactions + WHERE workspace_id = ? AND person_id = ? AND direction = 'inbound'`, + [workspaceId, personId], + ); + if (Number(byPerson?.n ?? 0) > 0) return true; + + if (!address) return false; + + const byAddress = await queryOne<{ n: number }>( + db, + `SELECT count(*) AS n FROM interactions + WHERE workspace_id = ? AND contact_address = ? AND direction = 'inbound'`, + [workspaceId, address], + ); + return Number(byAddress?.n ?? 0) > 0; +} + /** True when any suppression key matches this person (PRD §17.3). */ export async function isSuppressed( db: Client, diff --git a/apps/api/src/shared-inbox.test.ts b/apps/api/src/shared-inbox.test.ts new file mode 100644 index 0000000..9c24286 --- /dev/null +++ b/apps/api/src/shared-inbox.test.ts @@ -0,0 +1,278 @@ +/** + * The bug that made us look like a bot. + * + * Every rate limit is keyed on `person_id`, which is the correct key for "how + * often do we contact this human" and the wrong one for "how much mail does + * this mailbox get". A prospect with no personal address falls back to their + * company's shared inbox, so N prospects at one company are N separate people + * — each comfortably inside its own weekly limit — and one `support@` mailbox + * receives N messages. + * + * Production: 24 outbound emails reached 6 distinct addresses, and + * `support@canny.io` alone received 14 of them. Nothing was ever contacted + * twice by its own key, which is why no existing gate saw it. + * + * These tests are written against the delivered address rather than the + * person, because that is the thing the recipient experiences. + */ + +import { afterEach, describe, expect, test } from 'bun:test'; +import type { Hono } from 'hono'; +import { now, queryOne, type Client } from '@outreachgraph/db'; +import type { Mailer, Message, SendResult } from '@outreachgraph/email'; +import { createApp } from './app'; +import type { AppEnv, RequestActor } from './context'; +import { seedDatabase, SEED, type SeededDatabase } from './test-seed'; + +const ACTOR: RequestActor = { + userId: SEED.userId, + workspaceId: SEED.workspaceId, + organizationId: SEED.organizationId, + role: 'owner', +}; + +const SHARED_INBOX = 'support@acme.com'; + +let active: SeededDatabase | undefined; + +afterEach(() => { + active?.cleanup(); + active = undefined; +}); + +function recordingMailer(): { sent: Message[]; mailer: Mailer } { + const sent: Message[] = []; + return { + sent, + mailer: { + send: async (message): Promise => { + sent.push(message); + return { id: `msg_${sent.length}` }; + }, + }, + }; +} + +async function harness(label: string, mailer: Mailer): Promise<{ app: Hono; db: Client }> { + const seeded = await seedDatabase(label); + active = seeded; + return { + app: createApp({ db: seeded.db, authenticate: async () => ACTOR, mailer }), + db: seeded.db, + }; +} + +/** + * Gives the seeded company a shared inbox and turns Jane's card into an email + * one. Deliberately no personal address: falling back to the company mailbox + * is the whole subject. + */ +async function giveCompanySharedInbox(db: Client): Promise { + await db.execute({ + sql: 'UPDATE companies SET contact_email = ? WHERE id = ?', + args: [SHARED_INBOX, SEED.companyId], + }); + await db.execute({ + sql: `UPDATE recommendations SET action = 'send_email', network = 'email' WHERE id = ?`, + args: [SEED.recommendationId], + }); +} + +/** A second prospect at the same company, with their own card and no address. */ +async function addColleague( + db: Client, + suffix: string, +): Promise<{ personId: string; recommendationId: string }> { + const stamp = now(); + const personId = `per_${suffix}`; + const recommendationId = `rec_${suffix}`; + + await db.batch([ + { + sql: `INSERT INTO people (id, display_name, first_name, last_name, current_company_id, + current_title, location, identity_confidence, status, outreach_eligible, + believed_minor, created_at, updated_at) + VALUES (?, ?, ?, 'Colleague', ?, 'Director', 'Remote', 0.97, 'active', 1, 0, ?, ?)`, + args: [personId, `${suffix} Colleague`, suffix, SEED.companyId, stamp, stamp], + }, + { + sql: `INSERT INTO campaign_people (campaign_id, person_id, workspace_id, status, + interaction_state, discovered_at, updated_at) + VALUES (?, ?, ?, 'recommended', 'never_contacted', ?, ?)`, + args: [SEED.campaignId, personId, SEED.workspaceId, stamp, stamp], + }, + { + sql: `INSERT INTO recommendations (id, workspace_id, campaign_id, person_id, action, + network, priority, reason, trigger_signal_id, policy_status, policy_version, + expected_goal, status, created_at) + VALUES (?, ?, ?, ?, 'send_email', 'email', 90, 'Colleague at the same company.', + ?, 'allow_with_approval', '2026-08-11', 'start_conversation', 'pending', ?)`, + args: [recommendationId, SEED.workspaceId, SEED.campaignId, personId, SEED.signalId, stamp], + }, + { + sql: `INSERT INTO drafts (id, workspace_id, recommendation_id, body, grounded_signal_ids, + checks_json, created_at, updated_at) + VALUES (?, ?, ?, 'We ran into a similar cross-border settlement issue...', ?, '[]', ?, ?)`, + args: [ + `drf_${suffix}`, + SEED.workspaceId, + recommendationId, + JSON.stringify([SEED.signalId]), + stamp, + stamp, + ], + }, + ]); + + return { personId, recommendationId }; +} + +async function approve(app: Hono, recommendationId: string): Promise { + return await app.request(`/api/v1/recommendations/${recommendationId}/approve`, { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({}), + }); +} + +describe('a shared company inbox', () => { + test('records the address a message was actually delivered to', async () => { + const { mailer } = recordingMailer(); + const { app, db } = await harness('shared-records-address', mailer); + await giveCompanySharedInbox(db); + + await approve(app, SEED.recommendationId); + + const interaction = await queryOne<{ contact_address: string; shared_inbox: number }>( + db, + `SELECT contact_address, shared_inbox FROM interactions + WHERE workspace_id = ? AND direction = 'outbound'`, + [SEED.workspaceId], + ); + + // Without this column the next policy check cannot see the send at all. + expect(interaction?.contact_address).toBe(SHARED_INBOX); + expect(interaction?.shared_inbox).toBe(1); + }); + + test('refuses a second colleague whose mail lands in the same inbox', async () => { + const { sent, mailer } = recordingMailer(); + const { app, db } = await harness('shared-blocks-second', mailer); + await giveCompanySharedInbox(db); + const colleague = await addColleague(db, 'bob'); + + const first = await approve(app, SEED.recommendationId); + expect(first.status).toBe(200); + expect(sent).toHaveLength(1); + expect(sent[0]?.to).toBe(SHARED_INBOX); + + // Bob has never been contacted, is inside every per-person limit, and + // would still have put a second message in the same human's inbox. + const second = await approve(app, colleague.recommendationId); + expect(second.status).toBe(409); + + // Both address gates fire here — the inbox is over its weekly count *and* + // inside the cooldown — and the engine reports the last one to restrict. + // Which of the two is named matters less than that the refusal is about + // the address rather than about Bob, who is a stranger to us. + const payload = (await second.json()) as { + error?: { message?: string; details?: { gate?: string } }; + }; + expect(['rate_limit_address', 'cooldown']).toContain(payload.error?.details?.gate ?? ''); + expect(payload.error?.message).toMatch(/address|inbox/i); + + expect(sent).toHaveLength(1); + }); + + test('names the shared inbox when the count is what stops it', async () => { + const { sent, mailer } = recordingMailer(); + const { app, db } = await harness('shared-names-inbox', mailer); + await giveCompanySharedInbox(db); + const colleague = await addColleague(db, 'erin'); + + await approve(app, SEED.recommendationId); + + // Push the send outside the 72h cooldown so only the weekly count is left + // to refuse it, and the operator-facing reason is the one under test. + await db.execute({ + sql: `UPDATE interactions SET occurred_at = ? WHERE workspace_id = ? AND direction = 'outbound'`, + args: [new Date(Date.now() - 96 * 3_600_000).toISOString(), SEED.workspaceId], + }); + + const second = await approve(app, colleague.recommendationId); + expect(second.status).toBe(409); + + const payload = (await second.json()) as { + error?: { message?: string; details?: { gate?: string } }; + }; + expect(payload.error?.details?.gate).toBe('rate_limit_address'); + expect(payload.error?.message).toContain('shares a company inbox'); + expect(sent).toHaveLength(1); + }); + + test('a colleague with their own address is unaffected', async () => { + const { sent, mailer } = recordingMailer(); + const { app, db } = await harness('shared-allows-personal', mailer); + await giveCompanySharedInbox(db); + const colleague = await addColleague(db, 'carol'); + + await db.execute({ + sql: `INSERT INTO social_identities (id, person_id, network, handle, platform_user_id, + confidence, source_type, verified_by, first_seen_at) + VALUES ('sid_carol_email', ?, 'email', 'carol@acme.com', 'carol@acme.com', + 0.9, 'public_web', '[]', ?)`, + args: [colleague.personId, now()], + }); + + await approve(app, SEED.recommendationId); + const second = await approve(app, colleague.recommendationId); + + // A personal mailbox is a different human reading it, so the shared-inbox + // limit has nothing to say about it. + expect(second.status).toBe(200); + expect(sent).toHaveLength(2); + expect(sent[1]?.to).toBe('carol@acme.com'); + }); +}); + +describe('a contact who has replied', () => { + test('is not written to again', async () => { + const { sent, mailer } = recordingMailer(); + const { app, db } = await harness('reply-blocks', mailer); + await giveCompanySharedInbox(db); + + // The reply arrives against Jane, who has not yet been mailed. + const recorded = await app.request(`/api/v1/people/${SEED.personId}/replied`, { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ body: 'Thanks — already sorted, please stop.' }), + }); + expect(recorded.status).toBe(200); + + const response = await approve(app, SEED.recommendationId); + expect(response.status).toBe(409); + + const payload = (await response.json()) as { error?: { details?: { gate?: string } } }; + expect(payload.error?.details?.gate).toBe('conversation_open'); + expect(sent).toHaveLength(0); + }); + + test('a reply from a shared inbox also protects their colleagues', async () => { + const { sent, mailer } = recordingMailer(); + const { app, db } = await harness('reply-blocks-colleagues', mailer); + await giveCompanySharedInbox(db); + const colleague = await addColleague(db, 'dave'); + + await app.request(`/api/v1/people/${SEED.personId}/replied`, { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ body: 'Not interested.' }), + }); + + // Someone at that mailbox has answered on behalf of everyone written to + // there; carrying on down the staff list is the same mistake in disguise. + const response = await approve(app, colleague.recommendationId); + expect(response.status).toBe(409); + expect(sent).toHaveLength(0); + }); +}); diff --git a/apps/server/src/index.ts b/apps/server/src/index.ts index db186d2..4afa5ee 100644 --- a/apps/server/src/index.ts +++ b/apps/server/src/index.ts @@ -102,40 +102,70 @@ if (process.env.RUN_MIGRATIONS !== 'false') { * message. Refusing to boot over a missing key would take the whole system * down for a feature that is, by design, allowed to produce nothing. */ -const chain: FallbackEntry[] = []; - -if (process.env.ANTHROPIC_API_KEY) { - chain.push({ - name: 'anthropic', - model: new ClaudeModel({ - ...(process.env.ANTHROPIC_MODEL ? { model: process.env.ANTHROPIC_MODEL } : {}), - }), - }); +/** + * Which provider leads, and which ones back it up. + * + * `MODEL_CHAIN` is a comma list of provider names in preference order, e.g. + * `openai,anthropic`. It exists because the right leader is a billing + * decision, not an architectural one: drafting one short email per prospect is + * the highest-volume model call the product makes, and a frontier model is + * simply the wrong tool for it. Naming the order in the environment lets that + * be changed without a deploy, and a provider left out of the list is not + * built at all — the cheap chain stays cheap rather than quietly falling + * through to an expensive one. + * + * Unset preserves the original order, so nothing changes for a deployment that + * has not thought about it. + */ +const DEFAULT_CHAIN_ORDER = ['anthropic', 'openai', 'gemini'] as const; + +const builders: Record FallbackEntry | undefined> = { + anthropic: () => + process.env.ANTHROPIC_API_KEY + ? { + name: 'anthropic', + model: new ClaudeModel({ + ...(process.env.ANTHROPIC_MODEL ? { model: process.env.ANTHROPIC_MODEL } : {}), + }), + } + : undefined, + openai: () => + process.env.OPENAI_API_KEY + ? { + name: 'openai', + model: new OpenAIModel({ + ...(process.env.OPENAI_MODEL ? { model: process.env.OPENAI_MODEL } : {}), + }), + } + : undefined, + gemini: () => + process.env.GEMINI_API_KEY + ? { + name: 'gemini', + model: new GeminiModel({ + ...(process.env.GEMINI_MODEL ? { model: process.env.GEMINI_MODEL } : {}), + }), + } + : undefined, +}; + +const requestedOrder = (process.env.MODEL_CHAIN ?? '') + .split(',') + .map((name) => name.trim().toLowerCase()) + .filter(Boolean); + +// An unrecognised name is louder than it is fatal: booting with a silently +// shorter chain is how a deployment ends up with no composer and no clue why. +for (const name of requestedOrder) { + if (!(name in builders)) console.log(`MODEL_CHAIN: ignoring unknown provider "${name}"`); } -// Second, not instead of. A capped Anthropic key regains access on its own, and -// when it does the chain returns to Claude with no redeploy — which is what a -// fallback should do rather than becoming a permanent quiet substitution. -// -// OpenAI comes before Gemini because it is the closer substitute: the prompts -// and the deterministic gates were built around Claude, and a same-shaped model -// is likelier to produce a draft that still passes them. -if (process.env.OPENAI_API_KEY) { - chain.push({ - name: 'openai', - model: new OpenAIModel({ - ...(process.env.OPENAI_MODEL ? { model: process.env.OPENAI_MODEL } : {}), - }), - }); -} +const order = requestedOrder.length > 0 ? requestedOrder : [...DEFAULT_CHAIN_ORDER]; -if (process.env.GEMINI_API_KEY) { - chain.push({ - name: 'gemini', - model: new GeminiModel({ - ...(process.env.GEMINI_MODEL ? { model: process.env.GEMINI_MODEL } : {}), - }), - }); +const chain: FallbackEntry[] = []; +for (const name of order) { + const entry = builders[name]?.(); + if (entry && !chain.some((existing) => existing.name === entry.name)) chain.push(entry); } const model = diff --git a/migrations/0011_contact_address.sql b/migrations/0011_contact_address.sql new file mode 100644 index 0000000..550a1fd --- /dev/null +++ b/migrations/0011_contact_address.sql @@ -0,0 +1,68 @@ +-- 0011_contact_address.sql +-- Record the address a message actually went to, not just the person it was for. +-- +-- Every rate limit in `packages/policy` is keyed on `person_id`, which is the +-- right key for "how often do we contact this human" and the wrong key for +-- "how much mail does this mailbox get". A prospect with no personal address +-- falls back to their company's shared inbox, so fourteen prospects at one +-- company are fourteen separate people — each comfortably inside its own +-- weekly limit — and one `support@` mailbox receives fourteen messages. +-- +-- That is exactly what happened in production: 24 outbound emails reached 6 +-- distinct addresses, and `support@canny.io` alone got 14 of them. No +-- person-keyed limit could have seen it, because by its own key nothing was +-- ever contacted twice. +-- +-- Storing the resolved address on the interaction makes the mailbox countable. +-- The backfill below reconstructs it for rows written before this migration, +-- using the same resolution order the sender uses, so the limits protect the +-- addresses already contacted rather than starting from a clean slate. + +ALTER TABLE interactions ADD COLUMN contact_address TEXT; + +-- Whether that address belongs to the person or to their company. A shared +-- inbox is worth naming in a refusal: "weekly limit reached" is baffling for a +-- prospect who has never been contacted, until you know who shares the mailbox. +ALTER TABLE interactions ADD COLUMN shared_inbox INTEGER NOT NULL DEFAULT 0; + +CREATE INDEX idx_interactions_address + ON interactions(workspace_id, contact_address, occurred_at DESC); + +-- Backfill, in the sender's own order of preference: a personal email identity +-- first, the company contact address second. Rows that resolve to neither keep +-- a NULL address and are simply not counted, which is the honest answer — we +-- cannot say where they went. +UPDATE interactions + SET contact_address = ( + SELECT lower(trim(si.handle)) + FROM social_identities si + WHERE si.person_id = interactions.person_id + AND si.network = 'email' + ORDER BY si.confidence DESC + LIMIT 1 + ) + WHERE network = 'email' + AND contact_address IS NULL; + +UPDATE interactions + SET contact_address = ( + SELECT lower(trim(co.contact_email)) + FROM people p + JOIN companies co ON co.id = p.current_company_id + WHERE p.id = interactions.person_id + AND co.contact_email IS NOT NULL + AND trim(co.contact_email) <> '' + ), + shared_inbox = 1 + WHERE network = 'email' + AND contact_address IS NULL + -- Only where a company address actually exists, so a row with neither + -- keeps a NULL address instead of being marked as a shared inbox it has. + AND EXISTS ( + SELECT 1 + FROM people p + JOIN companies co ON co.id = p.current_company_id + WHERE p.id = interactions.person_id + AND co.contact_email IS NOT NULL + AND trim(co.contact_email) <> '' + ); diff --git a/packages/contracts/src/index.ts b/packages/contracts/src/index.ts index 412a93f..06f9677 100644 --- a/packages/contracts/src/index.ts +++ b/packages/contracts/src/index.ts @@ -253,6 +253,38 @@ export const executeActionSchema = z.object({ note: z.string().max(1_000).optional(), }); +/** + * Filling in the drafts a queue is missing. + * + * `limit` is capped because every card is a model call: the ceiling is what + * stops "catch the queue up" from being an open-ended invoice. + */ +export const backfillDraftsSchema = z.object({ + limit: z.number().int().min(1).max(200).default(25), +}); + +/** + * Recording that a contact wrote back. + * + * The product cannot see replies on its own — it sends through SMTP and reads + * no mailbox — so until inbound polling exists this is how a reply becomes a + * fact the policy engine can act on. It matters more than it looks: an + * unrecorded reply leaves the contact in the cold-outreach pool, and mailing + * someone who has already answered is the most bot-like thing we can do. + */ +export const recordReplySchema = z.object({ + /** When they replied. Defaults to now; accepted so backfills stay honest. */ + occurredAt: z.string().datetime().optional(), + /** The reply itself, when the caller has it. */ + body: z.string().max(20_000).optional(), + /** + * The address that replied, for a shared inbox where the person who answered + * is not necessarily the person written to. Defaults to the address the + * outreach was delivered to. + */ + fromAddress: z.string().email().optional(), +}); + /** * Connecting the mailbox outreach is sent from. * diff --git a/packages/pipeline/src/outreach-email.ts b/packages/pipeline/src/outreach-email.ts index 90d5a8d..479abaf 100644 --- a/packages/pipeline/src/outreach-email.ts +++ b/packages/pipeline/src/outreach-email.ts @@ -112,10 +112,14 @@ export async function recordEmailSent(db: Client, record: SentEmailRecord): Prom args: [record.externalId ?? null, at, record.actionId], }); + // `contact_address` is the mailbox this actually reached, which is not + // always a fact about the person: with no personal address it is their + // company's shared inbox. The rate limits count this column, so a send that + // does not write it is a send the next policy check cannot see. await db.execute({ sql: `INSERT INTO interactions (id, workspace_id, person_id, campaign_id, action_id, - network, direction, state, body, occurred_at, recorded_at) - VALUES (?, ?, ?, ?, ?, 'email', 'outbound', 'contacted', ?, ?, ?)`, + network, direction, state, body, contact_address, shared_inbox, occurred_at, recorded_at) + VALUES (?, ?, ?, ?, ?, 'email', 'outbound', 'contacted', ?, ?, ?, ?, ?)`, args: [ newId('interaction'), record.workspaceId, @@ -123,6 +127,8 @@ export async function recordEmailSent(db: Client, record: SentEmailRecord): Prom record.campaignId, record.actionId, record.body, + record.to.trim().toLowerCase(), + record.sharedInbox ? 1 : 0, at, at, ], diff --git a/packages/policy/src/engine.test.ts b/packages/policy/src/engine.test.ts index 5f54ef3..f0a17f8 100644 --- a/packages/policy/src/engine.test.ts +++ b/packages/policy/src/engine.test.ts @@ -240,6 +240,106 @@ describe('rate limits and budget (PRD §7.7, §18)', () => { expect(result.gate).toBe('budget_exhausted'); }); + test('deny a second message to an address, even for an untouched prospect', () => { + // The production bug: fourteen prospects at one company, each contacted + // for the first time, all delivering to the same support@ inbox. + const result = evaluatePolicy( + request({ + actionsToThisProspectThisWeek: 0, + maxActionsPerProspectPerWeek: 1, + actionsToThisAddressThisWeek: 1, + addressShared: true, + }), + ); + + expect(result.decision).toBe('deny'); + expect(result.gate).toBe('rate_limit_address'); + expect(result.reason).toContain('shares a company inbox'); + }); + + test('deny inside the cooldown for the address as well as the person', () => { + const result = evaluatePolicy( + request({ hoursSinceLastActionToAddress: 12, minHoursBetweenActions: 72 }), + ); + expect(result.decision).toBe('deny'); + expect(result.gate).toBe('cooldown'); + }); + + test('a personal address the prospect has to itself is unaffected', () => { + const result = evaluatePolicy( + request({ actionsToThisAddressThisWeek: 0, addressShared: false }), + ); + expect(result.decision).toBe('allow_with_approval'); + }); + + test('no resolved address means the address gates simply do not fire', () => { + const result = evaluatePolicy(request({ actionsToThisAddressThisWeek: undefined })); + expect(result.decision).toBe('allow_with_approval'); + }); + + test('the address cap defaults to one when the caller sets none', () => { + const result = evaluatePolicy(request({ actionsToThisAddressThisWeek: 1 })); + expect(result.decision).toBe('deny'); + expect(result.gate).toBe('rate_limit_address'); + }); + + test('a raised address cap is respected', () => { + const result = evaluatePolicy( + request({ actionsToThisAddressThisWeek: 1, maxActionsPerAddressPerWeek: 3 }), + ); + expect(result.decision).toBe('allow_with_approval'); + }); +}); + +describe('a contact who has replied', () => { + test('denies further cold outreach outright', () => { + const result = evaluatePolicy(request({ conversationOpen: true })); + + expect(result.decision).toBe('deny'); + expect(result.gate).toBe('conversation_open'); + expect(result.reason).toContain('already replied'); + }); + + test('downgrades an explicit follow-up to human approval rather than allowing it', () => { + const result = evaluatePolicy(request({ conversationOpen: true, isFollowUp: true })); + + expect(result.decision).toBe('allow_with_approval'); + expect(isExecutable(result.decision, true)).toBe(true); + expect(isExecutable(result.decision, false)).toBe(false); + }); + + test('a follow-up is still refused when a rate limit already denied it', () => { + // `deny` outranks `allow_with_approval`, so the follow-up downgrade must + // not be able to loosen a decision another gate has already tightened. + const result = evaluatePolicy( + request({ + conversationOpen: true, + isFollowUp: true, + actionsToThisAddressThisWeek: 5, + }), + ); + expect(result.decision).toBe('deny'); + }); + + test('trusted automation may not continue an open thread unattended', () => { + const result = evaluatePolicy( + request({ + network: 'email', + action: 'send_email', + approvalMode: 'trusted_automation', + conversationOpen: true, + }), + ); + + expect(result.decision).toBe('deny'); + expect(isExecutable(result.decision, false)).toBe(false); + }); + + test('leaves a contact who has not replied alone', () => { + const result = evaluatePolicy(request({ conversationOpen: false })); + expect(result.decision).toBe('allow_with_approval'); + }); + test('exempt internal bookkeeping from the daily cap', () => { const result = evaluatePolicy( request({ action: 'observe', actionsToday: 999, maxActionsPerDay: 50 }), diff --git a/packages/policy/src/engine.ts b/packages/policy/src/engine.ts index 9d67db7..3f81725 100644 --- a/packages/policy/src/engine.ts +++ b/packages/policy/src/engine.ts @@ -57,7 +57,9 @@ export const POLICY_GATES = [ 'budget_exhausted', 'rate_limit_daily', 'rate_limit_prospect', + 'rate_limit_address', 'cooldown', + 'conversation_open', 'no_connected_account', 'capability_mode', 'outbound_requires_approval', @@ -92,6 +94,42 @@ export interface PolicyRequest { readonly minHoursBetweenActions?: number; readonly budgetExhausted?: boolean; + /** + * The same limits again, counted against the address a message would + * actually be delivered to rather than the person it is addressed to. + * + * These are not redundant. Every limit above is keyed on `person_id`, but a + * person with no personal address falls back to their company's shared inbox + * — so fourteen different prospects at one company are fourteen separate + * people, each within its own weekly limit, and one `support@` mailbox + * receives fourteen messages. That is what "we look like a bot" is made of, + * and no person-keyed limit can see it. + * + * Omitted when the caller could not resolve an address, in which case these + * gates simply do not fire. + */ + readonly actionsToThisAddressThisWeek?: number; + readonly maxActionsPerAddressPerWeek?: number; + readonly hoursSinceLastActionToAddress?: number; + /** True when the address is a shared company inbox rather than a person's. */ + readonly addressShared?: boolean; + + /** + * Whether this contact has written back. + * + * A reply ends the cold-outreach phase. Anything further is a reply or a + * follow-up a human decided to send, never something a queue emits on its + * own — mailing someone who already answered is the single most bot-like + * thing the product can do. + */ + readonly conversationOpen?: boolean; + /** + * A deliberate follow-up to an open conversation, which downgrades the + * `conversation_open` gate from a refusal to a human approval rather than + * bypassing it. + */ + readonly isFollowUp?: boolean; + /** Explicit flag overrides. A missing key means enabled (PRD §37). */ readonly featureFlags?: Readonly>; @@ -118,6 +156,15 @@ export interface PolicyTraceEntry { const DEFAULT_COOLDOWN_HOURS = 72; +/** + * How many messages one delivery address may receive in a week. + * + * One, like the per-prospect limit it mirrors. A shared inbox is read by a + * human who does not care that the fourteen messages were addressed to + * fourteen different colleagues. + */ +const DEFAULT_ADDRESS_CAP_PER_WEEK = 1; + /** * Evaluates a single (network, action) request. * @@ -227,6 +274,53 @@ export function evaluatePolicy(request: PolicyRequest): PolicyResult { `Only ${round(elapsed)}h since the last contact; the cooldown is ${cooldown}h.`, ); } + + // The same two limits, counted against the delivery address. A shared + // inbox is the case they exist for, so it is named in the refusal — + // "weekly limit reached" is baffling for a prospect who has never been + // contacted, until you know fourteen colleagues share their mailbox. + const addressCount = request.actionsToThisAddressThisWeek; + const addressCap = request.maxActionsPerAddressPerWeek ?? DEFAULT_ADDRESS_CAP_PER_WEEK; + if (addressCount !== undefined && addressCount >= addressCap) { + restrict( + 'rate_limit_address', + 'deny', + request.addressShared === true + ? `This prospect shares a company inbox that has already had ` + + `${addressCount} message(s) this week; the limit is ${addressCap}.` + : `Weekly limit for this address reached (${addressCount}/${addressCap}).`, + ); + } + + const addressElapsed = request.hoursSinceLastActionToAddress; + if (addressElapsed !== undefined && addressElapsed < cooldown) { + restrict( + 'cooldown', + 'deny', + `Only ${round(addressElapsed)}h since this address was last contacted; ` + + `the cooldown is ${cooldown}h.`, + ); + } + } + + // 7b. A contact who has replied is out of the cold-outreach flow for good. + // Denied outright unless the sender says this is a follow-up, and even + // then a human approves it — an open thread is never autopilot's to + // continue. + if (isOutboundAction(request.action) && request.conversationOpen === true) { + if (request.isFollowUp === true) { + restrict( + 'conversation_open', + 'allow_with_approval', + 'This contact has replied. A follow-up needs a human to approve it.', + ); + } else { + restrict( + 'conversation_open', + 'deny', + 'This contact has already replied; reply to them instead of sending new outreach.', + ); + } } // 8. What the platform rule permits.