diff --git a/apps/server/package.json b/apps/server/package.json index 1c29a36..797dc91 100644 --- a/apps/server/package.json +++ b/apps/server/package.json @@ -24,12 +24,14 @@ "better-auth": "latest", "cron-parser": "latest", "hono": "latest", + "nodemailer": "^9.0.3", "playwright": "^1.61.1", "web-push": "^3.6.7", "zod": "latest" }, "devDependencies": { "@types/bun": "latest", + "@types/nodemailer": "^8.0.1", "@types/web-push": "^3.6.4", "typescript": "^6.0.3" } diff --git a/apps/server/src/extensions.ts b/apps/server/src/extensions.ts index a36a299..d668f66 100644 --- a/apps/server/src/extensions.ts +++ b/apps/server/src/extensions.ts @@ -1,26 +1,26 @@ export interface ExtensionCatalogEntry { - /** Stable key used to match an installed extension row back to its - * catalog entry — never renamed once shipped, since it's stored in the - * `extension.key` column. */ - key: string; - name: string; - description: string; - category: string; - /** Sidebar icon name — must be a key in ICON_MAP on the web client (see - * apps/web/src/components/app-sidebar.tsx). */ - icon: string; - /** Route segment under /workspace/{id}/extensions/{key} — normally equal - * to `key`, kept separate in case a future extension needs a nicer URL. */ - route: string; - /** GitHub repo URL of a companion plugin (see plugins.ts) installed - * automatically alongside this extension — its skills/sub-agents are - * what the extension's own agents actually use at runtime, so a bare - * `extension` row without it would be a shell with nothing behind it. - * Best-effort: a failed plugin install never blocks extension activation - * (see ensureExtensionPlugin), the extension just falls back to its - * built-in heuristics until the plugin is installed (retry from the - * Plugins page, or by disabling/re-enabling the extension). */ - pluginRepoUrl?: string; + /** Stable key used to match an installed extension row back to its + * catalog entry — never renamed once shipped, since it's stored in the + * `extension.key` column. */ + key: string; + name: string; + description: string; + category: string; + /** Sidebar icon name — must be a key in ICON_MAP on the web client (see + * apps/web/src/components/app-sidebar.tsx). */ + icon: string; + /** Route segment under /workspace/{id}/extensions/{key} — normally equal + * to `key`, kept separate in case a future extension needs a nicer URL. */ + route: string; + /** GitHub repo URL of a companion plugin (see plugins.ts) installed + * automatically alongside this extension — its skills/sub-agents are + * what the extension's own agents actually use at runtime, so a bare + * `extension` row without it would be a shell with nothing behind it. + * Best-effort: a failed plugin install never blocks extension activation + * (see ensureExtensionPlugin), the extension just falls back to its + * built-in heuristics until the plugin is installed (retry from the + * Plugins page, or by disabling/re-enabling the extension). */ + pluginRepoUrl?: string; } // The marketplace of extensions installable from workspace settings. Mirrors @@ -28,29 +28,36 @@ export interface ExtensionCatalogEntry { // extension creates an `extension` row scoped to the workspace, same as // "connecting" a catalog entry creates an `mcpServer` row. export const EXTENSION_CATALOG: ExtensionCatalogEntry[] = [ - { - key: "seo-analyzer", - name: "SEO/GEO/AEO Analyzer", - description: - "Links a domain to a local repo, audits search/answer-engine readiness, dispatches an AI agent to fix findings, and drafts blog posts for keyword gaps.", - category: "Growth", - icon: "TrendingUp", - route: "seo-analyzer", - pluginRepoUrl: "https://github.com/AgricIDaniel/claude-seo", - }, - { - key: "video-studio", - name: "Video Studio", - description: - "Describe a clip and generate it with Sora, then play back, edit, and organize everything it renders — every result lands in the Library too.", - category: "Creative", - icon: "Film", - route: "video-studio", - }, + { + key: "seo-analyzer", + name: "SEO/GEO/AEO Analyzer", + description: + "Links a domain to a local repo, audits search/answer-engine readiness, dispatches an AI agent to fix findings, and drafts blog posts for keyword gaps.", + category: "Growth", + icon: "TrendingUp", + route: "seo-analyzer", + pluginRepoUrl: "https://github.com/AgricIDaniel/claude-seo", + }, + { + key: "video-studio", + name: "Video Studio", + description: + "Describe a clip and generate it with Sora, then play back, edit, and organize everything it renders — every result lands in the Library too.", + category: "Creative", + icon: "Film", + route: "video-studio", + }, + { + key: "local-lead-scout", + name: "Local Lead Scout", + description: + "Finds local businesses with no website via compliant sources (CSV import, OSM Overpass, official Places API), then drafts a prototype and outreach email for you to review and approve before anything is sent.", + category: "Growth", + icon: "MapPin", + route: "local-lead-scout", + }, ]; -export function getExtensionCatalogEntry( - key: string, -): ExtensionCatalogEntry | undefined { - return EXTENSION_CATALOG.find((entry) => entry.key === key); +export function getExtensionCatalogEntry(key: string): ExtensionCatalogEntry | undefined { + return EXTENSION_CATALOG.find((entry) => entry.key === key); } diff --git a/apps/server/src/lead-scout-email.test.ts b/apps/server/src/lead-scout-email.test.ts new file mode 100644 index 0000000..a9bf59e --- /dev/null +++ b/apps/server/src/lead-scout-email.test.ts @@ -0,0 +1,86 @@ +import { afterEach, beforeEach, describe, expect, test } from "bun:test"; +import { getDb } from "@nyxel/db"; +import { installTestDb } from "@nyxel/db/test-utils"; +import { + sendLeadScoutEmail, + toClientSafeLeadScoutEmailSettings, + upsertLeadScoutEmailSettings, +} from "./lead-scout-email"; + +let ctx: Awaited>; + +beforeEach(async () => { + ctx = await installTestDb(); +}); +afterEach(async () => { + await ctx.cleanup(); +}); + +async function withWorkspace(): Promise { + const user = await getDb().getOrCreateDemoUser(); + const workspace = await getDb().createWorkspace({ userId: user.id, name: "email-test" }); + return workspace.id; +} + +describe("toClientSafeLeadScoutEmailSettings", () => { + test("strips credentials, exposes hasCredentials boolean", async () => { + const workspaceId = await withWorkspace(); + const settings = await upsertLeadScoutEmailSettings(workspaceId, { + fromName: "Acme", + fromEmail: "acme@example.com", + credentials: { host: "smtp.example.com", username: "u", password: "s3cret-password" }, + }); + + const safe = toClientSafeLeadScoutEmailSettings(settings); + expect((safe as Record).credentials).toBeUndefined(); + expect(safe.hasCredentials).toBe(true); + expect(JSON.stringify(safe)).not.toContain("s3cret-password"); + }); + + test("hasCredentials is false when none are configured", async () => { + const workspaceId = await withWorkspace(); + const settings = await upsertLeadScoutEmailSettings(workspaceId, { + fromName: "Acme", + fromEmail: "acme@example.com", + }); + expect(toClientSafeLeadScoutEmailSettings(settings).hasCredentials).toBe(false); + }); +}); + +describe("sendLeadScoutEmail dry-run gating", () => { + test("dry run mode never dispatches, even with no credentials configured", async () => { + const workspaceId = await withWorkspace(); + await upsertLeadScoutEmailSettings(workspaceId, { + fromName: "Acme", + fromEmail: "acme@example.com", + dryRunMode: true, + }); + + const result = await sendLeadScoutEmail(workspaceId, { + to: "lead@business.com", + subject: "Hi", + text: "body", + }); + expect(result).toEqual({ sent: false, dryRun: true }); + }); + + test("live mode without credentials throws instead of silently succeeding", async () => { + const workspaceId = await withWorkspace(); + await upsertLeadScoutEmailSettings(workspaceId, { + fromName: "Acme", + fromEmail: "acme@example.com", + dryRunMode: false, + }); + + await expect( + sendLeadScoutEmail(workspaceId, { to: "lead@business.com", subject: "Hi", text: "body" }), + ).rejects.toThrow(/configured/); + }); + + test("throws when email settings were never configured", async () => { + const workspaceId = await withWorkspace(); + await expect( + sendLeadScoutEmail(workspaceId, { to: "lead@business.com", subject: "Hi", text: "body" }), + ).rejects.toThrow(/configured/); + }); +}); diff --git a/apps/server/src/lead-scout-email.ts b/apps/server/src/lead-scout-email.ts new file mode 100644 index 0000000..ca88e70 --- /dev/null +++ b/apps/server/src/lead-scout-email.ts @@ -0,0 +1,321 @@ +import type { + LeadScoutEmailCredentials, + LeadScoutEmailProvider, + LeadScoutEmailSettingsRecord, +} from "@nyxel/db"; +import { getDb } from "@nyxel/db"; +import nodemailer from "nodemailer"; +import { logAudit } from "./audit"; + +/** + * Strips `credentials` before a settings row leaves the server — the client + * only needs to know a provider is configured, not its value (same + * convention as toClientSafeInstallation/toClientSafeMcpServer in + * trpc/router.ts). + */ +export function toClientSafeLeadScoutEmailSettings(settings: LeadScoutEmailSettingsRecord) { + const { credentials, ...rest } = settings; + return { + ...rest, + hasCredentials: credentials !== null && Object.keys(credentials).length > 0, + }; +} +export type LeadScoutEmailSettingsSummary = ReturnType; + +export async function getLeadScoutEmailSettings( + workspaceId: string, +): Promise { + return getDb().getLeadScoutEmailSettings(workspaceId); +} + +export async function upsertLeadScoutEmailSettings( + workspaceId: string, + input: { + provider?: LeadScoutEmailProvider; + fromName: string; + fromEmail: string; + replyTo?: string | null; + credentials?: LeadScoutEmailCredentials | null; + dailySendLimit?: number; + perCampaignSendLimit?: number; + dryRunMode?: boolean; + legalFooter?: string | null; + unsubscribeText?: string; + }, +): Promise { + const settings = await getDb().upsertLeadScoutEmailSettings({ workspaceId, ...input }); + await logAudit({ + workspaceId, + actor: "extension", + toolLabel: "local_lead_scout.email_settings.update", + input: { provider: input.provider, fromEmail: input.fromEmail, dryRunMode: input.dryRunMode }, + output: { hasCredentials: input.credentials !== undefined && input.credentials !== null }, + status: "success", + }); + return settings; +} + +function buildTransportForSmtp(credentials: LeadScoutEmailCredentials) { + if (!credentials.host) throw new Error("SMTP settings are missing a host."); + return nodemailer.createTransport({ + host: credentials.host, + port: Number(credentials.port ?? 587), + secure: credentials.secure === "true", + auth: credentials.username + ? { user: credentials.username, pass: credentials.password } + : undefined, + }); +} + +interface LeadScoutEmailMessage { + from: string; + replyTo?: string; + to: string; + subject: string; + text: string; + html?: string; +} + +async function sendViaSmtp( + credentials: LeadScoutEmailCredentials, + message: LeadScoutEmailMessage, +): Promise { + const transport = buildTransportForSmtp(credentials); + await transport.sendMail(message); +} + +async function sendViaResend( + credentials: LeadScoutEmailCredentials, + message: LeadScoutEmailMessage, +): Promise { + if (!credentials.apiKey) throw new Error("Resend settings are missing an API key."); + const res = await fetch("https://api.resend.com/emails", { + method: "POST", + headers: { Authorization: `Bearer ${credentials.apiKey}`, "Content-Type": "application/json" }, + body: JSON.stringify({ + from: message.from, + reply_to: message.replyTo, + to: [message.to], + subject: message.subject, + text: message.text, + html: message.html, + }), + }); + if (!res.ok) throw new Error(`Resend API error: HTTP ${res.status}`); +} + +async function sendViaMailgun( + credentials: LeadScoutEmailCredentials, + message: LeadScoutEmailMessage, +): Promise { + if (!credentials.apiKey || !credentials.domain) { + throw new Error("Mailgun settings are missing an API key or domain."); + } + const form = new URLSearchParams({ + from: message.from, + to: message.to, + subject: message.subject, + text: message.text, + }); + if (message.html) form.set("html", message.html); + if (message.replyTo) form.set("h:Reply-To", message.replyTo); + const res = await fetch(`https://api.mailgun.net/v3/${credentials.domain}/messages`, { + method: "POST", + headers: { + Authorization: `Basic ${Buffer.from(`api:${credentials.apiKey}`).toString("base64")}`, + "Content-Type": "application/x-www-form-urlencoded", + }, + body: form, + }); + if (!res.ok) throw new Error(`Mailgun API error: HTTP ${res.status}`); +} + +/** Generic compliant API adapter placeholder — posts to a user-configured + * webhook. Kept minimal per the extension's scope (no bespoke provider + * integrations beyond SMTP/Resend/Mailgun). */ +async function sendViaCustom( + credentials: LeadScoutEmailCredentials, + message: LeadScoutEmailMessage, +): Promise { + if (!credentials.webhookUrl) + throw new Error("Custom email provider has no webhookUrl configured."); + const res = await fetch(credentials.webhookUrl, { + method: "POST", + headers: { + "Content-Type": "application/json", + ...(credentials.secret ? { Authorization: `Bearer ${credentials.secret}` } : {}), + }, + body: JSON.stringify(message), + }); + if (!res.ok) throw new Error(`Custom email provider error: HTTP ${res.status}`); +} + +async function dispatchByProvider( + provider: LeadScoutEmailProvider, + credentials: LeadScoutEmailCredentials, + message: LeadScoutEmailMessage, +): Promise { + if (provider === "smtp") return sendViaSmtp(credentials, message); + if (provider === "resend") return sendViaResend(credentials, message); + if (provider === "mailgun") return sendViaMailgun(credentials, message); + return sendViaCustom(credentials, message); +} + +/** + * Real outbound send through the workspace's configured provider — always + * gated by dryRunMode as a defense-in-depth check (the approval-gated + * outreach flow in lead-scout.ts is the primary gate; this is a second, + * cheaper backstop). Body content is deliberately excluded from what's + * logged to the audit trail (only recipient/subject), per the "never log + * full message bodies" requirement. + */ +export async function sendLeadScoutEmail( + workspaceId: string, + message: { to: string; subject: string; text: string; html?: string }, +): Promise<{ sent: boolean; dryRun: boolean }> { + const settings = await getDb().getLeadScoutEmailSettings(workspaceId); + if (!settings) throw new Error("Email settings aren't configured for this workspace yet."); + + if (settings.dryRunMode) { + await logAudit({ + workspaceId, + actor: "extension", + toolLabel: "local_lead_scout.send_email", + input: { to: message.to, subject: message.subject, dryRun: true }, + output: { sent: false, dryRun: true }, + status: "success", + }); + return { sent: false, dryRun: true }; + } + + if (!settings.credentials && settings.provider !== "custom") { + throw new Error( + `${settings.provider} isn't configured — add credentials in email settings first.`, + ); + } + const from = `${settings.fromName} <${settings.fromEmail}>`; + + try { + await dispatchByProvider(settings.provider, settings.credentials ?? {}, { + from, + replyTo: settings.replyTo ?? undefined, + to: message.to, + subject: message.subject, + text: message.text, + html: message.html, + }); + await logAudit({ + workspaceId, + actor: "extension", + toolLabel: "local_lead_scout.send_email", + input: { to: message.to, subject: message.subject, dryRun: false }, + output: { sent: true }, + status: "success", + }); + return { sent: true, dryRun: false }; + } catch (err) { + const errorMessage = err instanceof Error ? err.message : String(err); + await logAudit({ + workspaceId, + actor: "extension", + toolLabel: "local_lead_scout.send_email", + input: { to: message.to, subject: message.subject, dryRun: false }, + output: { sent: false, error: errorMessage }, + status: "error", + }); + throw err; + } +} + +export async function sendLeadScoutTestEmail( + workspaceId: string, + toEmail: string, +): Promise<{ sent: boolean; dryRun: boolean }> { + const settings = await getDb().getLeadScoutEmailSettings(workspaceId); + if (!settings) throw new Error("Email settings aren't configured for this workspace yet."); + const footer = [settings.unsubscribeText, settings.legalFooter].filter(Boolean).join("\n\n"); + return sendLeadScoutEmail(workspaceId, { + to: toEmail, + subject: "Local Lead Scout — test email", + text: `This is a test email from ${settings.fromName} confirming your Local Lead Scout email settings are working.\n\n${footer}`, + }); +} + +/** + * Lightweight connectivity check for the "test connection" action — verifies + * credentials work without sending an actual message where the provider's + * API supports that (SMTP handshake, Resend/Mailgun auth ping). `custom` has + * no generic way to test a webhook without invoking it, so it's reported as + * unsupported rather than guessed at. + */ +export async function testLeadScoutEmailConnection( + workspaceId: string, +): Promise<{ ok: boolean; message: string }> { + const settings = await getDb().getLeadScoutEmailSettings(workspaceId); + if (!settings) + return { ok: false, message: "Email settings aren't configured for this workspace yet." }; + const credentials = settings.credentials ?? {}; + + try { + if (settings.provider === "smtp") { + await buildTransportForSmtp(credentials).verify(); + } else if (settings.provider === "resend") { + if (!credentials.apiKey) throw new Error("Missing API key."); + const res = await fetch("https://api.resend.com/domains", { + headers: { Authorization: `Bearer ${credentials.apiKey}` }, + }); + if (!res.ok) throw new Error(`HTTP ${res.status}`); + } else if (settings.provider === "mailgun") { + if (!credentials.apiKey) throw new Error("Missing API key."); + const res = await fetch("https://api.mailgun.net/v3/domains", { + headers: { + Authorization: `Basic ${Buffer.from(`api:${credentials.apiKey}`).toString("base64")}`, + }, + }); + if (!res.ok) throw new Error(`HTTP ${res.status}`); + } else { + return { + ok: false, + message: "Custom providers have no built-in connection test — send a test email instead.", + }; + } + await logAudit({ + workspaceId, + actor: "extension", + toolLabel: "local_lead_scout.email_settings.test_connection", + input: { provider: settings.provider }, + output: { ok: true }, + status: "success", + }); + return { ok: true, message: "Connection succeeded." }; + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + await logAudit({ + workspaceId, + actor: "extension", + toolLabel: "local_lead_scout.email_settings.test_connection", + input: { provider: settings.provider }, + output: { ok: false, error: message }, + status: "error", + }); + return { ok: false, message }; + } +} + +/** + * Daily send-limit enforcement — called by the approval-gated send flow + * (lead-scout.ts) immediately before dispatching, so the counter and the + * actual send stay consistent. Throws rather than returning false so a + * blocked send can't be silently ignored by a caller that forgets to check + * a boolean. + */ +export async function reserveLeadScoutDailySendSlot(workspaceId: string): Promise { + const settings = await getDb().getLeadScoutEmailSettings(workspaceId); + if (!settings) throw new Error("Email settings aren't configured for this workspace yet."); + const today = new Date().toISOString().slice(0, 10); + const currentCount = settings.sendCountDate === today ? settings.sendCountToday : 0; + if (currentCount >= settings.dailySendLimit) { + throw new Error(`Daily send limit of ${settings.dailySendLimit} reached for this workspace.`); + } + await getDb().incrementLeadScoutEmailSendCount(workspaceId, today); +} diff --git a/apps/server/src/lead-scout-geocode.ts b/apps/server/src/lead-scout-geocode.ts new file mode 100644 index 0000000..b940b60 --- /dev/null +++ b/apps/server/src/lead-scout-geocode.ts @@ -0,0 +1,68 @@ +/** + * Postal code -> coordinates via Nominatim (OpenStreetMap's own geocoder, + * free, no API key) — used by both the google_places_api provider (needs a + * center point for Nearby Search) and the osm_overpass provider (needs a + * bounding box). Respects Nominatim's usage policy: max 1 request/second + * and a descriptive User-Agent identifying the app. + * https://operations.osmfoundation.org/policies/nominatim/ + */ + +const NOMINATIM_USER_AGENT = "NyxelOS-LocalLeadScout/1.0 (local business discovery extension)"; +const NOMINATIM_MIN_INTERVAL_MS = 1100; + +// Process-local throttle, not cross-replica — matches the rest of the +// codebase's "single-process is the deployment target today" assumption +// (see scheduler.ts's automationsInFlight). +let lastNominatimCallAt = 0; + +async function throttleNominatim(): Promise { + const elapsed = Date.now() - lastNominatimCallAt; + if (elapsed < NOMINATIM_MIN_INTERVAL_MS) { + await new Promise((resolve) => setTimeout(resolve, NOMINATIM_MIN_INTERVAL_MS - elapsed)); + } + lastNominatimCallAt = Date.now(); +} + +export interface GeocodeResult { + lat: number; + lon: number; +} + +export async function geocodePostalCode( + postalCode: string, + country: string, +): Promise { + await throttleNominatim(); + const url = new URL("https://nominatim.openstreetmap.org/search"); + url.searchParams.set("postalcode", postalCode); + url.searchParams.set("country", country); + url.searchParams.set("format", "json"); + url.searchParams.set("limit", "1"); + + const res = await fetch(url, { headers: { "User-Agent": NOMINATIM_USER_AGENT } }); + if (!res.ok) throw new Error(`Nominatim geocoding failed: HTTP ${res.status}`); + const results = (await res.json()) as { lat: string; lon: string }[]; + const [first] = results; + if (!first) throw new Error(`Could not geocode postal code "${postalCode}, ${country}".`); + return { lat: Number(first.lat), lon: Number(first.lon) }; +} + +export interface BoundingBox { + south: number; + west: number; + north: number; + east: number; +} + +/** Rough equirectangular approximation — plenty accurate for a "search + * within N km of this postal code" radius, not a routing/distance product. */ +export function boundingBoxFromCenter(lat: number, lon: number, radiusKm: number): BoundingBox { + const latDelta = radiusKm / 111; + const lonDelta = radiusKm / (111 * Math.cos((lat * Math.PI) / 180)); + return { + south: lat - latDelta, + west: lon - lonDelta, + north: lat + latDelta, + east: lon + lonDelta, + }; +} diff --git a/apps/server/src/lead-scout-providers/custom-api.ts b/apps/server/src/lead-scout-providers/custom-api.ts new file mode 100644 index 0000000..e782e9a --- /dev/null +++ b/apps/server/src/lead-scout-providers/custom-api.ts @@ -0,0 +1,76 @@ +import { classifyWebsite } from "./shared"; +import type { LeadSourceProvider, LeadSourceSearchInput } from "./types"; + +/** Minimal generic REST adapter placeholder — expects a compliant endpoint + * (configured per-workspace as sourceConfig.config.baseUrl) that returns a + * flat JSON array shaped like the manual_csv row fields. Real + * provider-specific behavior (auth scheme, pagination, rate limits) is left + * to whichever compliant API a workspace wires up here; this adapter only + * defines the minimal request/response contract. */ +interface CustomApiRow { + sourceId: string; + businessName: string; + category?: string; + formattedAddress?: string; + postalCode?: string; + city?: string; + phone?: string; + email?: string; + website?: string; +} + +export const customApiProvider: LeadSourceProvider = { + sourcePolicy: + "Generic compliant API adapter placeholder — only queries the endpoint you configure; " + + "you are responsible for that endpoint's own terms of service and rate limits.", + + async searchBusinesses(input: LeadSourceSearchInput) { + const baseUrl = input.sourceConfig?.config?.baseUrl; + if (!baseUrl || typeof baseUrl !== "string") { + throw new Error( + "custom_api requires a baseUrl configured in workspace source settings first.", + ); + } + const apiKey = input.sourceConfig?.apiKey; + const url = new URL(baseUrl); + url.searchParams.set("postalCode", input.postalCode); + url.searchParams.set("country", input.country); + url.searchParams.set("radiusKm", String(input.radiusKm)); + url.searchParams.set("niches", input.niches.join(",")); + url.searchParams.set("maxResults", String(input.maxResults)); + + const res = await fetch(url, { + headers: apiKey ? { Authorization: `Bearer ${apiKey}` } : undefined, + }); + if (!res.ok) throw new Error(`custom_api provider error: HTTP ${res.status}`); + const rows = (await res.json()) as CustomApiRow[]; + return Array.isArray(rows) ? rows : []; + }, + + async getBusinessDetails(raw) { + return raw; + }, + + normalizeBusiness(row) { + const { website, websiteStatus } = classifyWebsite(row.website); + return { + sourceId: row.sourceId, + businessName: row.businessName, + category: row.category ?? null, + formattedAddress: row.formattedAddress ?? null, + postalCode: row.postalCode ?? null, + city: row.city ?? null, + phone: row.phone ?? null, + email: row.email ?? null, + website, + websiteStatus, + confidence: 60, + evidenceSummary: + websiteStatus === "missing_website" + ? "custom_api response had no website value for this business." + : "custom_api reported a website for this business.", + missingReason: + websiteStatus === "missing_website" ? "No website field in the custom API response." : null, + }; + }, +}; diff --git a/apps/server/src/lead-scout-providers/google-places.ts b/apps/server/src/lead-scout-providers/google-places.ts new file mode 100644 index 0000000..157db8e --- /dev/null +++ b/apps/server/src/lead-scout-providers/google-places.ts @@ -0,0 +1,134 @@ +import { geocodePostalCode } from "../lead-scout-geocode"; +import type { LeadSourceProvider, LeadSourceSearchInput } from "./types"; + +/** + * Official Places API (New) Text Search only — no browser automation, no + * Google Maps DOM scraping. FieldMask is explicit and minimal (never `*`), + * per the extension's compliance requirements: only what the workflow needs + * to judge a lead and dedupe it (`id`, name, address, location, status, + * types, phone, website, maps link). Reviews/photos/ratings are never + * requested and never stored. + */ +const FIELD_MASK = [ + "places.id", + "places.displayName", + "places.formattedAddress", + "places.location", + "places.businessStatus", + "places.types", + "places.nationalPhoneNumber", + "places.internationalPhoneNumber", + "places.websiteUri", + "places.googleMapsUri", +].join(","); + +const MAX_RESULTS_PER_QUERY = 20; +const QUERY_DELAY_MS = 250; + +interface GooglePlace { + id: string; + displayName?: { text: string }; + formattedAddress?: string; + location?: { latitude: number; longitude: number }; + businessStatus?: string; + types?: string[]; + nationalPhoneNumber?: string; + internationalPhoneNumber?: string; + websiteUri?: string; + googleMapsUri?: string; + /** Not part of the API response — stamped on by searchBusinesses so + * normalizeBusiness knows which niche query surfaced this result. */ + __niche?: string; +} + +async function searchTextOnce( + apiKey: string, + textQuery: string, + center: { lat: number; lon: number }, + radiusMeters: number, + maxResultCount: number, +): Promise { + const res = await fetch("https://places.googleapis.com/v1/places:searchText", { + method: "POST", + headers: { + "Content-Type": "application/json", + "X-Goog-Api-Key": apiKey, + "X-Goog-FieldMask": FIELD_MASK, + }, + body: JSON.stringify({ + textQuery, + locationBias: { + circle: { + center: { latitude: center.lat, longitude: center.lon }, + radius: radiusMeters, + }, + }, + maxResultCount: Math.min(maxResultCount, MAX_RESULTS_PER_QUERY), + }), + }); + if (!res.ok) { + throw new Error(`Google Places API error: HTTP ${res.status} ${await res.text()}`); + } + const body = (await res.json()) as { places?: GooglePlace[] }; + return body.places ?? []; +} + +export const googlePlacesProvider: LeadSourceProvider = { + sourcePolicy: + "Official Google Places API (Text Search, New) only — no Maps scraping, no DOM automation, no bulk harvesting. " + + "You are responsible for complying with the Google Maps Platform Terms of Service for any data retrieved this way.", + + async searchBusinesses(input: LeadSourceSearchInput) { + const apiKey = input.sourceConfig?.apiKey; + if (!apiKey) { + throw new Error( + "google_places_api requires an API key configured in workspace source settings first.", + ); + } + const center = await geocodePostalCode(input.postalCode, input.country); + const radiusMeters = Math.min(Math.max(input.radiusKm, 0.5), 50) * 1000; + const niches = input.niches.length > 0 ? input.niches : ["local business"]; + + const results: GooglePlace[] = []; + for (const niche of niches) { + if (results.length >= input.maxResults) break; + const remaining = input.maxResults - results.length; + const textQuery = `${niche} near ${input.postalCode}, ${input.country}`; + const places = await searchTextOnce(apiKey, textQuery, center, radiusMeters, remaining); + for (const place of places) results.push({ ...place, __niche: niche }); + if (niches.indexOf(niche) < niches.length - 1) { + await new Promise((resolve) => setTimeout(resolve, QUERY_DELAY_MS)); + } + } + return results; + }, + + async getBusinessDetails(raw) { + // Text Search (New) with the FieldMask above already returns everything + // the workflow needs — no separate Place Details call, which would be + // extra quota/latency for data already in hand. + return raw; + }, + + normalizeBusiness(place) { + const website = place.websiteUri?.trim() || null; + return { + sourceId: place.id, + businessName: place.displayName?.text ?? "Unknown business", + category: place.types?.[0] ?? null, + niche: place.__niche ?? null, + formattedAddress: place.formattedAddress ?? null, + postalCode: null, + city: null, + phone: place.nationalPhoneNumber ?? place.internationalPhoneNumber ?? null, + email: null, + website, + websiteStatus: website ? "has_website" : "missing_website", + confidence: 90, + evidenceSummary: website + ? `Google Places API returned a websiteUri (${website}).` + : "Google Places API (Text Search) returned no websiteUri for this business.", + missingReason: website ? null : "No websiteUri field in the Places API response.", + }; + }, +}; diff --git a/apps/server/src/lead-scout-providers/index.ts b/apps/server/src/lead-scout-providers/index.ts new file mode 100644 index 0000000..076f3c2 --- /dev/null +++ b/apps/server/src/lead-scout-providers/index.ts @@ -0,0 +1,20 @@ +import type { LeadScoutProvider } from "@nyxel/db"; +import { customApiProvider } from "./custom-api"; +import { googlePlacesProvider } from "./google-places"; +import { manualCsvProvider } from "./manual-csv"; +import { osmOverpassProvider } from "./osm-overpass"; +import type { LeadSourceProvider } from "./types"; + +export type { LeadSourceProvider, LeadSourceSearchInput, NormalizedBusiness } from "./types"; +export { collectNormalizedBusinesses } from "./types"; + +const PROVIDERS: Record = { + manual_csv: manualCsvProvider, + google_places_api: googlePlacesProvider, + osm_overpass: osmOverpassProvider, + custom_api: customApiProvider, +}; + +export function getLeadSourceProvider(provider: LeadScoutProvider): LeadSourceProvider { + return PROVIDERS[provider]; +} diff --git a/apps/server/src/lead-scout-providers/manual-csv.test.ts b/apps/server/src/lead-scout-providers/manual-csv.test.ts new file mode 100644 index 0000000..5fb4e64 --- /dev/null +++ b/apps/server/src/lead-scout-providers/manual-csv.test.ts @@ -0,0 +1,85 @@ +import { describe, expect, test } from "bun:test"; +import { manualCsvProvider } from "./manual-csv"; +import { parseCsv } from "./shared"; +import { collectNormalizedBusinesses } from "./types"; + +describe("parseCsv", () => { + test("parses quoted fields with embedded commas and escaped quotes", () => { + const csv = 'businessName,address,notes\n"Joe, Inc.","123 Main St","Great ""spot"""\n'; + const rows = parseCsv(csv); + expect(rows).toEqual([ + { businessName: "Joe, Inc.", address: "123 Main St", notes: 'Great "spot"' }, + ]); + }); + + test("pads ragged rows with empty strings for missing trailing columns", () => { + const csv = "businessName,address,phone\nJoe's Pizza,123 Main St\n"; + const rows = parseCsv(csv); + expect(rows).toEqual([{ businessName: "Joe's Pizza", address: "123 Main St", phone: "" }]); + }); +}); + +describe("manualCsvProvider", () => { + test("flags missing website, assigns maximum confidence to user-supplied data", async () => { + const csvText = + "businessName,address,postalCode,city,category,phone,email,website,notes\n" + + "Joe's Pizza,123 Main St,94103,SF,restaurant,555-1234,,,\n" + + "Acme Plumbing,55 Elm St,94103,SF,plumber,555-9999,acme@x.com,https://acme.com,has site\n"; + + const results = await collectNormalizedBusinesses(manualCsvProvider, { + workspaceId: "w1", + postalCode: "94103", + country: "US", + radiusKm: 5, + niches: [], + maxResults: 10, + sourceConfig: null, + csvText, + }); + + expect(results).toHaveLength(2); + const pizza = results.find((r) => r.businessName === "Joe's Pizza"); + expect(pizza?.websiteStatus).toBe("missing_website"); + expect(pizza?.website).toBeNull(); + expect(pizza?.confidence).toBe(100); + expect(pizza?.missingReason).toContain("No website column value in the CSV row"); + expect(pizza?.evidenceSummary).toContain("no website value was provided"); + + const acme = results.find((r) => r.businessName === "Acme Plumbing"); + expect(acme?.websiteStatus).toBe("has_website"); + expect(acme?.website).toBe("https://acme.com"); + }); + + test("re-importing the same row produces the same sourceId (dedupe key)", async () => { + const csvText = + "businessName,address,postalCode,city,category,phone,email,website,notes\n" + + "Joe's Pizza,123 Main St,94103,SF,restaurant,555-1234,,,\n"; + const input = { + workspaceId: "w1", + postalCode: "94103", + country: "US", + radiusKm: 5, + niches: [], + maxResults: 10, + sourceConfig: null, + csvText, + }; + const first = await collectNormalizedBusinesses(manualCsvProvider, input); + const second = await collectNormalizedBusinesses(manualCsvProvider, input); + expect(first[0]?.sourceId).toBe(second[0]?.sourceId); + }); + + test("throws when csvText is missing", async () => { + await expect( + collectNormalizedBusinesses(manualCsvProvider, { + workspaceId: "w1", + postalCode: "94103", + country: "US", + radiusKm: 5, + niches: [], + maxResults: 10, + sourceConfig: null, + }), + ).rejects.toThrow(/csvText/); + }); +}); diff --git a/apps/server/src/lead-scout-providers/manual-csv.ts b/apps/server/src/lead-scout-providers/manual-csv.ts new file mode 100644 index 0000000..ec0b936 --- /dev/null +++ b/apps/server/src/lead-scout-providers/manual-csv.ts @@ -0,0 +1,81 @@ +import { createHash } from "node:crypto"; +import { classifyWebsite, parseCsv } from "./shared"; +import type { LeadSourceProvider, LeadSourceSearchInput } from "./types"; + +/** businessName is the only strictly required column — everything else is + * optional per the extension's field list (businessName, address, + * postalCode, city, category, phone, email, website, notes). */ +interface CsvLeadRow { + businessName: string; + address: string; + postalCode: string; + city: string; + category: string; + phone: string; + email: string; + website: string; + notes: string; +} + +/** Stable per-row identity: a CSV has no natural id, so hash the row's own + * content — re-importing the same file reconciles against existing leads + * instead of duplicating them (see lead_scout_lead's unique constraint). */ +function csvRowSourceId(row: Record): string { + const key = [row.businessName, row.address, row.postalCode, row.phone, row.email] + .join("|") + .toLowerCase(); + return createHash("sha256").update(key).digest("hex").slice(0, 32); +} + +export const manualCsvProvider: LeadSourceProvider = { + sourcePolicy: + "User-supplied CSV import — fully within the user's own control, no external data source involved.", + + async searchBusinesses(input: LeadSourceSearchInput) { + if (!input.csvText) { + throw new Error("manual_csv requires csvText — upload a CSV to import leads."); + } + const rows = parseCsv(input.csvText); + return rows + .filter((row) => row.businessName) + .map((row) => ({ + sourceId: csvRowSourceId(row), + businessName: row.businessName ?? "", + address: row.address ?? "", + postalCode: row.postalCode ?? "", + city: row.city ?? "", + category: row.category ?? "", + phone: row.phone ?? "", + email: row.email ?? "", + website: row.website ?? "", + notes: row.notes ?? "", + })); + }, + + async getBusinessDetails(raw) { + return raw; + }, + + normalizeBusiness(row) { + const { website, websiteStatus } = classifyWebsite(row.website); + return { + sourceId: row.sourceId, + businessName: row.businessName, + category: row.category || null, + formattedAddress: row.address || null, + postalCode: row.postalCode || null, + city: row.city || null, + phone: row.phone || null, + email: row.email || null, + website, + websiteStatus, + confidence: 100, + evidenceSummary: + websiteStatus === "missing_website" + ? "Imported via manual CSV; no website value was provided." + : `Imported via manual CSV.${row.notes ? ` Notes: ${row.notes}` : ""}`, + missingReason: + websiteStatus === "missing_website" ? "No website column value in the CSV row." : null, + }; + }, +}; diff --git a/apps/server/src/lead-scout-providers/osm-overpass.test.ts b/apps/server/src/lead-scout-providers/osm-overpass.test.ts new file mode 100644 index 0000000..773d5bb --- /dev/null +++ b/apps/server/src/lead-scout-providers/osm-overpass.test.ts @@ -0,0 +1,48 @@ +import { describe, expect, test } from "bun:test"; +import { osmOverpassProvider } from "./osm-overpass"; + +describe("osmOverpassProvider.normalizeBusiness", () => { + test("flags missing_website and a separate missing_email note when neither tag is present", () => { + const result = osmOverpassProvider.normalizeBusiness({ + type: "node", + id: 123, + tags: { name: "Village Cafe", shop: "cafe", "addr:postcode": "94103", "addr:city": "SF" }, + }); + expect(result.sourceId).toBe("node/123"); + expect(result.websiteStatus).toBe("missing_website"); + expect(result.email).toBeNull(); + expect(result.evidenceSummary).toContain("missing_email"); + // Lower default than google_places_api's 90, per the extension's design + // (OSM tag completeness varies a lot by region/contributor). + expect(result.confidence).toBe(70); + }); + + test("has_website when a contact:website tag is present", () => { + const result = osmOverpassProvider.normalizeBusiness({ + type: "way", + id: 456, + tags: { name: "Acme Hardware", shop: "hardware", "contact:website": "https://acme.example" }, + }); + expect(result.websiteStatus).toBe("has_website"); + expect(result.website).toBe("https://acme.example"); + }); + + test("has_website when a plain website tag is present (contact:website takes priority when both exist)", () => { + const result = osmOverpassProvider.normalizeBusiness({ + type: "node", + id: 789, + tags: { name: "Bike Shop", shop: "bicycle", website: "https://bikes.example" }, + }); + expect(result.websiteStatus).toBe("has_website"); + expect(result.website).toBe("https://bikes.example"); + }); + + test("category picked from whichever of shop/amenity/office/craft is present", () => { + const result = osmOverpassProvider.normalizeBusiness({ + type: "node", + id: 1, + tags: { name: "Law Office", office: "lawyer" }, + }); + expect(result.category).toBe("lawyer"); + }); +}); diff --git a/apps/server/src/lead-scout-providers/osm-overpass.ts b/apps/server/src/lead-scout-providers/osm-overpass.ts new file mode 100644 index 0000000..c20bfef --- /dev/null +++ b/apps/server/src/lead-scout-providers/osm-overpass.ts @@ -0,0 +1,122 @@ +import { boundingBoxFromCenter, geocodePostalCode } from "../lead-scout-geocode"; +import type { LeadSourceProvider, LeadSourceSearchInput } from "./types"; + +/** Descriptive User-Agent as required by the Overpass/OSM usage policy — + * https://operations.osmfoundation.org/policies/overpass/ */ +const OVERPASS_USER_AGENT = "NyxelOS-LocalLeadScout/1.0 (local business discovery extension)"; +const OVERPASS_ENDPOINT = "https://overpass-api.de/api/interpreter"; +const QUERY_TIMEOUT_S = 25; + +interface OverpassElement { + type: "node" | "way" | "relation"; + id: number; + tags?: Record; +} + +function buildOverpassQuery(bbox: { + south: number; + west: number; + north: number; + east: number; +}): string { + const bounds = `${bbox.south},${bbox.west},${bbox.north},${bbox.east}`; + // nwr matches node/way/relation in one clause per OSM Overpass QL — one + // clause per top-level key the extension is allowed to query (shop, + // amenity, office, craft), scoped to the bounding box. + return `[out:json][timeout:${QUERY_TIMEOUT_S}]; +( + nwr["shop"](${bounds}); + nwr["amenity"](${bounds}); + nwr["office"](${bounds}); + nwr["craft"](${bounds}); +); +out center tags;`; +} + +function elementCategory(tags: Record): string | null { + for (const key of ["shop", "amenity", "office", "craft"]) { + if (tags[key]) return tags[key]; + } + return null; +} + +function elementMatchesNiches(tags: Record, niches: string[]): boolean { + if (niches.length === 0) return true; + const haystack = [tags.name, tags.shop, tags.amenity, tags.office, tags.craft] + .filter(Boolean) + .join(" ") + .toLowerCase(); + return niches.some((niche) => haystack.includes(niche.toLowerCase())); +} + +function elementAddress(tags: Record): string | null { + if (tags["addr:full"]) return tags["addr:full"]; + const parts = [tags["addr:housenumber"], tags["addr:street"]].filter(Boolean); + return parts.length > 0 ? parts.join(" ") : null; +} + +export const osmOverpassProvider: LeadSourceProvider = { + sourcePolicy: + "Public Overpass API over OpenStreetMap data (ODbL license) — a compliant, key-free public API, not scraping. " + + "Fair-use: bounded bounding box per scan, capped result count, descriptive User-Agent.", + + async searchBusinesses(input: LeadSourceSearchInput) { + const center = await geocodePostalCode(input.postalCode, input.country); + const bbox = boundingBoxFromCenter(center.lat, center.lon, input.radiusKm); + const query = buildOverpassQuery(bbox); + + const res = await fetch(OVERPASS_ENDPOINT, { + method: "POST", + headers: { + "Content-Type": "application/x-www-form-urlencoded", + "User-Agent": OVERPASS_USER_AGENT, + }, + body: `data=${encodeURIComponent(query)}`, + }); + if (!res.ok) throw new Error(`Overpass API error: HTTP ${res.status}`); + const body = (await res.json()) as { elements?: OverpassElement[] }; + const elements = (body.elements ?? []).filter((el) => el.tags?.name); + + const matching = elements.filter((el) => elementMatchesNiches(el.tags ?? {}, input.niches)); + return matching.slice(0, input.maxResults); + }, + + async getBusinessDetails(raw) { + return raw; + }, + + normalizeBusiness(element) { + const tags = element.tags ?? {}; + const website = (tags.website || tags["contact:website"] || "").trim() || null; + const email = (tags.email || tags["contact:email"] || "").trim() || null; + const phone = (tags.phone || tags["contact:phone"] || "").trim() || null; + const missingEmail = !email; + + return { + sourceId: `${element.type}/${element.id}`, + businessName: tags.name ?? "Unknown business", + category: elementCategory(tags), + formattedAddress: elementAddress(tags), + postalCode: tags["addr:postcode"] ?? null, + city: tags["addr:city"] ?? null, + phone, + email, + website, + websiteStatus: website ? "has_website" : "missing_website", + // Lower than google_places_api's default (90) — OSM tag completeness + // varies a lot by region/contributor. + confidence: 70, + evidenceSummary: [ + website + ? `OSM tags include a website/contact:website value (${website}).` + : "No website or contact:website tag present on this OSM element.", + missingEmail + ? "No email or contact:email tag present (missing_email) — outreach needs a manually supplied address." + : null, + ] + .filter(Boolean) + .join(" "), + missingReason: website ? null : "No website/contact:website tag on the OSM node/way.", + }; + }, +}; diff --git a/apps/server/src/lead-scout-providers/shared.ts b/apps/server/src/lead-scout-providers/shared.ts new file mode 100644 index 0000000..05ec4b1 --- /dev/null +++ b/apps/server/src/lead-scout-providers/shared.ts @@ -0,0 +1,82 @@ +import type { LeadScoutWebsiteStatus } from "@nyxel/db"; + +/** A lead qualifies as missing_website only when no website value is present + * at all (per every provider's compliance requirements — never inferred + * from scraping). A present-but-unparseable value is reported separately + * (invalid_website) so a reviewer can tell "no site" from "bad data" apart. */ +export function classifyWebsite(raw: string | null | undefined): { + website: string | null; + websiteStatus: LeadScoutWebsiteStatus; +} { + const trimmed = raw?.trim(); + if (!trimmed) return { website: null, websiteStatus: "missing_website" }; + const withProtocol = /^https?:\/\//i.test(trimmed) ? trimmed : `https://${trimmed}`; + try { + new URL(withProtocol); + return { website: trimmed, websiteStatus: "has_website" }; + } catch { + return { website: trimmed, websiteStatus: "invalid_website" }; + } +} + +/** Minimal RFC4180-ish CSV parser (quoted fields, escaped `""`, embedded + * commas/newlines) — avoids a new dependency for the one thing manual_csv + * needs. Returns each row as a header-keyed record; ragged rows are padded + * with empty strings for missing trailing columns. */ +export function parseCsv(text: string): Record[] { + const rows: string[][] = []; + let row: string[] = []; + let field = ""; + let inQuotes = false; + + const pushField = () => { + row.push(field); + field = ""; + }; + const pushRow = () => { + pushField(); + rows.push(row); + row = []; + }; + + for (let i = 0; i < text.length; i++) { + const char = text[i]; + if (inQuotes) { + if (char === '"') { + if (text[i + 1] === '"') { + field += '"'; + i++; + } else { + inQuotes = false; + } + } else { + field += char; + } + continue; + } + if (char === '"') { + inQuotes = true; + } else if (char === ",") { + pushField(); + } else if (char === "\n") { + pushRow(); + } else if (char === "\r") { + // swallow, \n (or end of input) closes the row + } else { + field += char; + } + } + if (field.length > 0 || row.length > 0) pushRow(); + + const [header, ...dataRows] = rows.filter( + (r) => r.length > 0 && !(r.length === 1 && r[0] === ""), + ); + if (!header) return []; + return dataRows.map((cells) => { + const record: Record = {}; + header.forEach((key, index) => { + record[key.trim()] = (cells[index] ?? "").trim(); + }); + return record; + }); +} diff --git a/apps/server/src/lead-scout-providers/types.ts b/apps/server/src/lead-scout-providers/types.ts new file mode 100644 index 0000000..7085fe3 --- /dev/null +++ b/apps/server/src/lead-scout-providers/types.ts @@ -0,0 +1,69 @@ +import type { LeadScoutSourceConfigRecord, LeadScoutWebsiteStatus } from "@nyxel/db"; + +/** What a scan needs to run one query — the campaign's own config, plus the + * workspace's per-provider source config (API keys, provider-specific + * settings) and, for manual_csv only, the raw CSV text to import. */ +export interface LeadSourceSearchInput { + workspaceId: string; + postalCode: string; + country: string; + radiusKm: number; + niches: string[]; + maxResults: number; + sourceConfig: LeadScoutSourceConfigRecord | null; + csvText?: string; +} + +/** The output shape every provider normalizes into — exactly what a + * reviewer needs to judge a lead (no reviews/photos/ratings/bulk metadata, + * per each provider's compliance requirements). */ +export interface NormalizedBusiness { + sourceId: string; + businessName: string; + category?: string | null; + niche?: string | null; + formattedAddress?: string | null; + postalCode?: string | null; + city?: string | null; + phone?: string | null; + email?: string | null; + website?: string | null; + websiteStatus: LeadScoutWebsiteStatus; + confidence: number; + evidenceSummary: string; + missingReason?: string | null; +} + +/** + * Provider adapter contract (see the extension's design doc). `TRaw` is + * whatever shape `searchBusinesses` returns per result (a search hit); the + * details/normalize split lets a provider fetch more per-result info before + * normalizing without every provider needing a separate details call — + * providers whose search response is already complete (Google Places New, + * manual CSV rows) just return the raw item unchanged from + * `getBusinessDetails`. + */ +export interface LeadSourceProvider { + readonly sourcePolicy: string; + searchBusinesses(input: LeadSourceSearchInput): Promise; + getBusinessDetails(raw: TRaw, input: LeadSourceSearchInput): Promise; + normalizeBusiness(details: TRaw): NormalizedBusiness; +} + +/** Runs one provider end-to-end and caps the result count — the one place + * every scan's maxResultsPerRun limit is actually enforced, regardless of + * how many raw results a provider's search returned. */ +export async function collectNormalizedBusinesses( + provider: LeadSourceProvider, + input: LeadSourceSearchInput, +): Promise { + const raws = await provider.searchBusinesses(input); + const normalized: NormalizedBusiness[] = []; + for (const raw of raws) { + if (normalized.length >= input.maxResults) break; + const details = await provider.getBusinessDetails(raw, input); + if (!details) continue; + normalized.push(provider.normalizeBusiness(details)); + } + return normalized; +} diff --git a/apps/server/src/lead-scout-send.test.ts b/apps/server/src/lead-scout-send.test.ts new file mode 100644 index 0000000..bd0aafb --- /dev/null +++ b/apps/server/src/lead-scout-send.test.ts @@ -0,0 +1,141 @@ +import { afterEach, beforeEach, describe, expect, test } from "bun:test"; +import { getDb } from "@nyxel/db"; +import { installTestDb } from "@nyxel/db/test-utils"; +import { upsertLeadScoutEmailSettings } from "./lead-scout-email"; +import { + addLeadScoutSuppression, + approveLeadScoutOutreachDraft, + rejectLeadScoutOutreachDraft, + resetLeadScoutLeadForResend, + sendLeadScoutOutreachDraft, +} from "./lead-scout-send"; + +let ctx: Awaited>; +beforeEach(async () => { + ctx = await installTestDb(); +}); +afterEach(async () => { + await ctx.cleanup(); +}); + +async function setupWorkspace(outreachMode: "draft_only" | "review_and_send" = "review_and_send") { + const db = getDb(); + const user = await db.getOrCreateDemoUser(); + const workspace = await db.createWorkspace({ userId: user.id, name: "ws" }); + const ext = await db.installExtension({ workspaceId: workspace.id, key: "local-lead-scout" }); + const campaign = await db.createLeadScoutCampaign({ + workspaceId: workspace.id, + extensionId: ext.id, + name: "c", + postalCode: "94103", + provider: "manual_csv", + outreachMode, + }); + await upsertLeadScoutEmailSettings(workspace.id, { + fromName: "Acme", + fromEmail: "acme@example.com", + dryRunMode: true, + perCampaignSendLimit: 1, + dailySendLimit: 10, + }); + return { db, workspace, campaign }; +} + +async function makeApprovedDraft( + db: ReturnType, + workspaceId: string, + campaignId: string, + email: string | null, +) { + const lead = await db.createLeadScoutLead({ + workspaceId, + campaignId, + sourceProvider: "manual_csv", + sourceId: `row-${email ?? "none"}`, + businessName: "Joe's Pizza", + email, + websiteStatus: "missing_website", + }); + await db.claimLeadScoutLeadStatus({ id: lead.id, fromStatus: "new", toStatus: "email_drafted" }); + const draft = await db.createLeadScoutOutreachDraft({ workspaceId, leadId: lead.id }); + await db.updateLeadScoutOutreachDraft(draft.id, { subject: "Hi", bodyText: "Body" }); + return { lead, draft }; +} + +describe("approval-gated send flow", () => { + test("draft-only campaigns block sending entirely", async () => { + const { db, workspace, campaign } = await setupWorkspace("draft_only"); + const { draft } = await makeApprovedDraft(db, workspace.id, campaign.id, "lead@business.com"); + await approveLeadScoutOutreachDraft(draft.id); + await expect(sendLeadScoutOutreachDraft(draft.id)).rejects.toThrow(/draft-only/); + }); + + test("rejected approval never sends", async () => { + const { db, workspace, campaign } = await setupWorkspace(); + const { draft } = await makeApprovedDraft(db, workspace.id, campaign.id, "lead2@business.com"); + await rejectLeadScoutOutreachDraft(draft.id); + await expect(sendLeadScoutOutreachDraft(draft.id)).rejects.toThrow(); + }); + + test("no email on the lead blocks send (never guessed)", async () => { + const { db, workspace, campaign } = await setupWorkspace(); + const { draft } = await makeApprovedDraft(db, workspace.id, campaign.id, null); + await approveLeadScoutOutreachDraft(draft.id); + await expect(sendLeadScoutOutreachDraft(draft.id)).rejects.toThrow(/no email/i); + }); + + test("a suppressed email/domain blocks send and marks the lead suppressed", async () => { + const { db, workspace, campaign } = await setupWorkspace(); + await addLeadScoutSuppression({ + workspaceId: workspace.id, + email: "blocked@business.com", + reason: "opted out", + }); + const { lead, draft } = await makeApprovedDraft( + db, + workspace.id, + campaign.id, + "blocked@business.com", + ); + await approveLeadScoutOutreachDraft(draft.id); + await expect(sendLeadScoutOutreachDraft(draft.id)).rejects.toThrow(/suppression/); + + const leadAfter = await db.getLeadScoutLead(lead.id); + expect(leadAfter?.status).toBe("suppressed"); + }); + + test("send succeeds once; duplicate send is blocked; explicit reset allows resend", async () => { + const { db, workspace, campaign } = await setupWorkspace(); + const { lead, draft } = await makeApprovedDraft( + db, + workspace.id, + campaign.id, + "good@business.com", + ); + await approveLeadScoutOutreachDraft(draft.id); + const sent = await sendLeadScoutOutreachDraft(draft.id); + expect(sent.status).toBe("sent"); + + const leadAfter = await db.getLeadScoutLead(lead.id); + expect(leadAfter?.status).toBe("sent"); + + // Duplicate send attempt on the same (already-sent) draft is rejected. + await expect(sendLeadScoutOutreachDraft(draft.id)).rejects.toThrow(); + + // An explicit reset is required before the lead can be emailed again. + const reset = await resetLeadScoutLeadForResend(lead.id); + expect(reset.status).toBe("reviewed"); + await expect(resetLeadScoutLeadForResend(lead.id)).rejects.toThrow(); + }); + + test("per-campaign send limit is enforced", async () => { + const { db, workspace, campaign } = await setupWorkspace(); + const first = await makeApprovedDraft(db, workspace.id, campaign.id, "first@business.com"); + await approveLeadScoutOutreachDraft(first.draft.id); + await sendLeadScoutOutreachDraft(first.draft.id); + + const second = await makeApprovedDraft(db, workspace.id, campaign.id, "second@business.com"); + await approveLeadScoutOutreachDraft(second.draft.id); + await expect(sendLeadScoutOutreachDraft(second.draft.id)).rejects.toThrow(/send limit/); + }); +}); diff --git a/apps/server/src/lead-scout-send.ts b/apps/server/src/lead-scout-send.ts new file mode 100644 index 0000000..b38e4f6 --- /dev/null +++ b/apps/server/src/lead-scout-send.ts @@ -0,0 +1,270 @@ +import type { + LeadScoutLeadRecord, + LeadScoutOutreachDraftRecord, + LeadScoutSuppressionRecord, +} from "@nyxel/db"; +import { getDb } from "@nyxel/db"; +import { logAudit } from "./audit"; +import { reserveLeadScoutDailySendSlot, sendLeadScoutEmail } from "./lead-scout-email"; +import { notifyWorkspaceOwner } from "./push"; + +function domainFromEmail(email: string | null): string | null { + if (!email) return null; + const at = email.lastIndexOf("@"); + return at === -1 ? null : email.slice(at + 1).toLowerCase(); +} + +/** + * Approves a drafted outreach email — the one human decision that unlocks + * `sendLeadScoutOutreachDraft`. Nothing is sent here; this only moves the + * draft/lead into the state that permits a send attempt. Atomic on both the + * draft and the lead so a duplicate click can't approve twice. + */ +export async function approveLeadScoutOutreachDraft( + draftId: string, +): Promise { + const db = getDb(); + const draft = await db.getLeadScoutOutreachDraft(draftId); + if (!draft) throw new Error(`Unknown outreach draft: ${draftId}`); + + const claimed = await db.claimLeadScoutOutreachDraftStatus({ + id: draftId, + fromStatus: "draft", + toStatus: "approved", + }); + if (!claimed) throw new Error("This draft was already approved, rejected, or sent."); + + const updated = await db.updateLeadScoutOutreachDraft(draftId, { approvedAt: new Date() }); + await db.claimLeadScoutLeadStatus({ + id: draft.leadId, + fromStatus: "email_drafted", + toStatus: "approved_to_send", + }); + + await logAudit({ + workspaceId: draft.workspaceId, + actor: "extension", + toolLabel: "local_lead_scout.request_send_approval", + input: { draftId }, + output: { decision: "approved" }, + status: "success", + }); + return updated; +} + +/** Rejecting a draft is final for that draft — nothing sends, and the lead + * moves to `rejected` rather than looping back for another draft attempt + * automatically (a human should decide to re-draft explicitly). */ +export async function rejectLeadScoutOutreachDraft( + draftId: string, +): Promise { + const db = getDb(); + const draft = await db.getLeadScoutOutreachDraft(draftId); + if (!draft) throw new Error(`Unknown outreach draft: ${draftId}`); + + const claimed = await db.claimLeadScoutOutreachDraftStatus({ + id: draftId, + fromStatus: "draft", + toStatus: "rejected", + }); + if (!claimed) throw new Error("This draft was already approved, rejected, or sent."); + + await db.claimLeadScoutLeadStatus({ + id: draft.leadId, + fromStatus: "email_drafted", + toStatus: "rejected", + }); + + await logAudit({ + workspaceId: draft.workspaceId, + actor: "extension", + toolLabel: "local_lead_scout.request_send_approval", + input: { draftId }, + output: { decision: "rejected" }, + status: "rejected", + }); + return claimed; +} + +async function countSentForCampaign(campaignId: string): Promise { + const leads = await getDb().listLeadScoutLeadsByCampaign(campaignId); + return leads.filter((lead) => lead.status === "sent").length; +} + +/** + * Sends an approved outreach draft — the single irreversible action in this + * whole extension. Every compliance gate lives here: campaign outreach mode, + * no-guessed-email, suppression list, per-campaign and daily send limits, + * and an atomic claim (draft `approved` -> `sending`) so a duplicate click + * or a concurrent request can never send the same draft twice. On any + * failure past the claim, both the draft and lead revert to their + * pre-send state so the user can fix the issue and retry. + */ +export async function sendLeadScoutOutreachDraft( + draftId: string, +): Promise { + const db = getDb(); + const draft = await db.getLeadScoutOutreachDraft(draftId); + if (!draft) throw new Error(`Unknown outreach draft: ${draftId}`); + const lead = await db.getLeadScoutLead(draft.leadId); + if (!lead) throw new Error(`Unknown lead: ${draft.leadId}`); + const campaign = await db.getLeadScoutCampaign(lead.campaignId); + if (!campaign) throw new Error(`Unknown lead scout campaign: ${lead.campaignId}`); + const emailSettings = await db.getLeadScoutEmailSettings(lead.workspaceId); + if (!emailSettings) throw new Error("Configure email settings before sending outreach."); + + if (campaign.outreachMode !== "review_and_send") { + throw new Error("This campaign is in draft-only mode — sending is disabled."); + } + if (!lead.email) { + throw new Error("This lead has no email address — outreach is blocked until one is supplied."); + } + const suppressed = await db.getLeadScoutSuppressionMatch( + lead.workspaceId, + lead.email, + domainFromEmail(lead.email), + ); + if (suppressed) { + await db.claimLeadScoutOutreachDraftStatus({ + id: draftId, + fromStatus: "approved", + toStatus: "rejected", + }); + await db.claimLeadScoutLeadStatus({ + id: lead.id, + fromStatus: "approved_to_send", + toStatus: "suppressed", + }); + throw new Error(`This lead's email/domain is on the suppression list (${suppressed.reason}).`); + } + if (lead.status === "sent") { + throw new Error( + "This lead was already sent — use the explicit reset action before sending again.", + ); + } + const sentSoFar = await countSentForCampaign(campaign.id); + if (sentSoFar >= emailSettings.perCampaignSendLimit) { + throw new Error( + `This campaign's send limit (${emailSettings.perCampaignSendLimit}) has been reached.`, + ); + } + + const claimedDraft = await db.claimLeadScoutOutreachDraftStatus({ + id: draftId, + fromStatus: "approved", + toStatus: "sending", + }); + if (!claimedDraft) throw new Error("This draft is no longer approved (already sent or reset)."); + const claimedLead = await db.claimLeadScoutLeadStatus({ + id: lead.id, + fromStatus: "approved_to_send", + toStatus: "sending", + }); + if (!claimedLead) { + // Lead moved on some other way (e.g. reset) between our reads and the + // claim above — revert the draft claim so nothing is left half-sent. + await db.claimLeadScoutOutreachDraftStatus({ + id: draftId, + fromStatus: "sending", + toStatus: "approved", + }); + throw new Error("Lead status changed concurrently — refresh and try again."); + } + + try { + await reserveLeadScoutDailySendSlot(lead.workspaceId); + await sendLeadScoutEmail(lead.workspaceId, { + to: lead.email, + subject: draft.subject ?? "A quick idea for your business", + text: draft.bodyText ?? "", + html: draft.bodyHtml ?? undefined, + }); + + const sent = await db.updateLeadScoutOutreachDraft(draftId, { + status: "sent", + sentAt: new Date(), + }); + await db.claimLeadScoutLeadStatus({ id: lead.id, fromStatus: "sending", toStatus: "sent" }); + + await logAudit({ + workspaceId: lead.workspaceId, + actor: "extension", + toolLabel: "local_lead_scout.send_email", + input: { draftId, leadId: lead.id }, + output: { sent: true }, + status: "success", + }); + await notifyWorkspaceOwner(lead.workspaceId, { + title: "Outreach email sent", + body: `${lead.businessName}: outreach email sent.`, + url: `/workspace/${lead.workspaceId}/extensions/local-lead-scout`, + tag: `lead-scout-send-${draftId}`, + }); + return sent; + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + const failed = await db.updateLeadScoutOutreachDraft(draftId, { + status: "failed", + errorMessage: message, + }); + await db.claimLeadScoutLeadStatus({ + id: lead.id, + fromStatus: "sending", + toStatus: "approved_to_send", + }); + await logAudit({ + workspaceId: lead.workspaceId, + actor: "extension", + toolLabel: "local_lead_scout.send_email", + input: { draftId, leadId: lead.id }, + output: message, + status: "error", + }); + return failed; + } +} + +/** Explicit, manual escape hatch for "email this lead again" — a sent lead + * is otherwise permanently done (see the `lead.status === "sent"` check in + * sendLeadScoutOutreachDraft). Resets it to `reviewed` so the whole + * prototype -> draft -> approve -> send cycle can run again from scratch, + * rather than resurrecting old (possibly stale) prototype/draft rows. */ +export async function resetLeadScoutLeadForResend(leadId: string): Promise { + const db = getDb(); + const claimed = await db.claimLeadScoutLeadStatus({ + id: leadId, + fromStatus: "sent", + toStatus: "reviewed", + }); + if (!claimed) throw new Error("Only a previously-sent lead can be reset for resend."); + await logAudit({ + workspaceId: claimed.workspaceId, + actor: "extension", + toolLabel: "local_lead_scout.reset_for_resend", + input: { leadId }, + output: { status: "reviewed" }, + status: "success", + }); + return claimed; +} + +export async function addLeadScoutSuppression(input: { + workspaceId: string; + email?: string | null; + domain?: string | null; + reason: string; +}): Promise { + if (!input.email && !input.domain) { + throw new Error("A suppression entry needs at least an email or a domain."); + } + const record = await getDb().createLeadScoutSuppression(input); + await logAudit({ + workspaceId: input.workspaceId, + actor: "extension", + toolLabel: "local_lead_scout.add_suppression", + input: { email: input.email, domain: input.domain, reason: input.reason }, + output: { id: record.id }, + status: "success", + }); + return record; +} diff --git a/apps/server/src/lead-scout.test.ts b/apps/server/src/lead-scout.test.ts new file mode 100644 index 0000000..f48bb10 --- /dev/null +++ b/apps/server/src/lead-scout.test.ts @@ -0,0 +1,170 @@ +import { afterEach, beforeEach, describe, expect, test } from "bun:test"; +import { getDb } from "@nyxel/db"; +import { installTestDb } from "@nyxel/db/test-utils"; +import { + approveLeadScoutPrototype, + dispatchLeadScoutPrototype, + markLeadScoutLeadReviewed, + parseEmailDraftOutput, + parsePrototypeOutput, + runLeadScoutScan, +} from "./lead-scout"; + +describe("parsePrototypeOutput", () => { + test("extracts structured fields and an artifact block", () => { + const output = [ + "CONCEPT: A clean single-page site for a local bakery.", + "HERO_COPY: Fresh bread daily | Baked with love since 1990", + "SECTIONS: Hero, Menu, About, Contact", + "CTA: Order now", + "STYLE: warm, rustic, earthy tones", + "---ARTIFACT---", + "# Village Bakery\n\nFresh bread daily.", + ].join("\n"); + const parsed = parsePrototypeOutput(output); + expect(parsed.concept).toBe("A clean single-page site for a local bakery."); + expect(parsed.sections).toEqual(["Hero", "Menu", "About", "Contact"]); + expect(parsed.callToAction).toBe("Order now"); + expect(parsed.artifactMarkdown).toContain("Village Bakery"); + }); + + test("treats the literal word none as no artifact", () => { + const output = + "CONCEPT: x\nHERO_COPY: y\nSECTIONS: a, b\nCTA: z\nSTYLE: s\n---ARTIFACT---\nnone"; + expect(parsePrototypeOutput(output).artifactMarkdown).toBeNull(); + }); +}); + +describe("parseEmailDraftOutput", () => { + test("extracts subject and body", () => { + const output = + "SUBJECT: A quick idea\n---BODY---\nHi there, noticed you don't have a website..."; + const { subject, body } = parseEmailDraftOutput(output); + expect(subject).toBe("A quick idea"); + expect(body).toContain("noticed you don't have a website"); + }); + + test("falls back to a default subject when SUBJECT: is missing", () => { + const { subject } = parseEmailDraftOutput("---BODY---\njust a body"); + expect(subject).toBe("A quick idea for your business"); + }); +}); + +let ctx: Awaited>; +beforeEach(async () => { + ctx = await installTestDb(); +}); +afterEach(async () => { + await ctx.cleanup(); +}); + +async function setupCampaign(provider: "manual_csv" = "manual_csv") { + const db = getDb(); + const user = await db.getOrCreateDemoUser(); + const workspace = await db.createWorkspace({ userId: user.id, name: "ws" }); + const ext = await db.installExtension({ workspaceId: workspace.id, key: "local-lead-scout" }); + const campaign = await db.createLeadScoutCampaign({ + workspaceId: workspace.id, + extensionId: ext.id, + name: "SF test", + postalCode: "94103", + provider, + }); + return { db, workspace, campaign }; +} + +describe("runLeadScoutScan", () => { + test("ingests CSV leads with correct missing-website counts", async () => { + const { db, campaign } = await setupCampaign(); + const csvText = + "businessName,address,postalCode,city,category,phone,email,website,notes\n" + + "Joe's Pizza,123 Main St,94103,SF,restaurant,555-1234,,,\n" + + "Acme Plumbing,55 Elm St,94103,SF,plumber,555-9999,acme@x.com,https://acme.com,has site\n"; + + const run = await runLeadScoutScan(campaign.id, { csvText }); + expect(run.status).toBe("completed"); + expect(run.newLeadCount).toBe(2); + expect(run.missingWebsiteCount).toBe(1); + + const leads = await db.listLeadScoutLeadsByCampaign(campaign.id); + expect(leads).toHaveLength(2); + expect(leads.find((l) => l.businessName === "Joe's Pizza")?.websiteStatus).toBe( + "missing_website", + ); + }); + + test("re-scanning the same CSV reconciles instead of duplicating leads", async () => { + const { db, campaign } = await setupCampaign(); + const csvText = + "businessName,address,postalCode,city,category,phone,email,website,notes\n" + + "Joe's Pizza,123 Main St,94103,SF,restaurant,555-1234,,,\n"; + + const first = await runLeadScoutScan(campaign.id, { csvText }); + expect(first.newLeadCount).toBe(1); + const second = await runLeadScoutScan(campaign.id, { csvText }); + expect(second.newLeadCount).toBe(0); + + const leads = await db.listLeadScoutLeadsByCampaign(campaign.id); + expect(leads).toHaveLength(1); + }); + + test("marks the scan run failed instead of throwing when the provider errors", async () => { + const { campaign } = await setupCampaign(); + // manual_csv without csvText throws inside the provider. + const run = await runLeadScoutScan(campaign.id); + expect(run.status).toBe("failed"); + expect(run.errorMessage).toMatch(/csvText/); + }); +}); + +describe("lead status gating for prototype dispatch", () => { + test("requireApprovalBeforePrototype (default true) blocks dispatch until reviewed", async () => { + const { db, workspace, campaign } = await setupCampaign(); + const lead = await db.createLeadScoutLead({ + workspaceId: workspace.id, + campaignId: campaign.id, + sourceProvider: "manual_csv", + sourceId: "row-1", + businessName: "Joe's Pizza", + websiteStatus: "missing_website", + }); + + await expect(dispatchLeadScoutPrototype(lead.id)).rejects.toThrow(/reviewed/); + + await markLeadScoutLeadReviewed(lead.id); + const reviewed = await db.getLeadScoutLead(lead.id); + expect(reviewed?.status).toBe("reviewed"); + }); + + test("marking an already-reviewed lead reviewed again fails (no double transition)", async () => { + const { workspace, campaign } = await setupCampaign(); + const db = getDb(); + const lead = await db.createLeadScoutLead({ + workspaceId: workspace.id, + campaignId: campaign.id, + sourceProvider: "manual_csv", + sourceId: "row-2", + businessName: "Acme", + websiteStatus: "missing_website", + }); + await markLeadScoutLeadReviewed(lead.id); + await expect(markLeadScoutLeadReviewed(lead.id)).rejects.toThrow(); + }); + + test("approving a prototype that isn't ready is rejected", async () => { + const { db, workspace, campaign } = await setupCampaign(); + const lead = await db.createLeadScoutLead({ + workspaceId: workspace.id, + campaignId: campaign.id, + sourceProvider: "manual_csv", + sourceId: "row-3", + businessName: "Pending Co", + websiteStatus: "missing_website", + }); + const prototype = await db.createLeadScoutPrototype({ + workspaceId: workspace.id, + leadId: lead.id, + }); + await expect(approveLeadScoutPrototype(prototype.id)).rejects.toThrow(/ready/); + }); +}); diff --git a/apps/server/src/lead-scout.ts b/apps/server/src/lead-scout.ts new file mode 100644 index 0000000..41819a6 --- /dev/null +++ b/apps/server/src/lead-scout.ts @@ -0,0 +1,484 @@ +import type { + AgentRecord, + LeadScoutCampaignRecord, + LeadScoutOutreachDraftRecord, + LeadScoutPrototypeRecord, + LeadScoutScanRunRecord, +} from "@nyxel/db"; +import { getDb } from "@nyxel/db"; +import { listAvailableModels } from "@nyxel/model-providers"; +import { executeManagedTask } from "./agent-runtime"; +import { logAudit } from "./audit"; +import { collectNormalizedBusinesses, getLeadSourceProvider } from "./lead-scout-providers"; +import { getInstalledProvidersForWorkspace } from "./models"; +import { notifyWorkspaceOwner } from "./push"; + +/** + * Runs one scan for a campaign: dispatches to its configured provider, + * reconciles results against already-ingested leads (same-source businesses + * update in place instead of duplicating — see lead_scout_lead's unique + * constraint), and records a scan_run with the outcome. Mirrors + * runSeoAnalysis's shape (create run -> do the work -> update run -> + * audit -> notify), synchronous end-to-end so both the manual tRPC mutation + * and the scheduler poll can just await it. + */ +export async function runLeadScoutScan( + campaignId: string, + options?: { csvText?: string }, +): Promise { + const db = getDb(); + const campaign = await db.getLeadScoutCampaign(campaignId); + if (!campaign) throw new Error(`Unknown lead scout campaign: ${campaignId}`); + + const sourceConfig = await db.getLeadScoutSourceConfig(campaign.workspaceId, campaign.provider); + const run = await db.createLeadScoutScanRun({ + campaignId, + workspaceId: campaign.workspaceId, + provider: campaign.provider, + }); + + try { + const provider = getLeadSourceProvider(campaign.provider); + const businesses = await collectNormalizedBusinesses(provider, { + workspaceId: campaign.workspaceId, + postalCode: campaign.postalCode, + country: campaign.country, + radiusKm: campaign.radiusKm, + niches: campaign.niches, + maxResults: campaign.maxResultsPerRun, + sourceConfig, + csvText: options?.csvText, + }); + + let newLeadCount = 0; + let missingWebsiteCount = 0; + for (const business of businesses) { + if (business.confidence < campaign.minConfidence) continue; + if (business.websiteStatus === "missing_website") missingWebsiteCount++; + + const existing = await db.getLeadScoutLeadBySource( + campaignId, + campaign.provider, + business.sourceId, + ); + if (existing) { + await db.updateLeadScoutLead(existing.id, { + scanRunId: run.id, + category: business.category ?? existing.category, + niche: business.niche ?? existing.niche, + formattedAddress: business.formattedAddress ?? existing.formattedAddress, + postalCode: business.postalCode ?? existing.postalCode, + city: business.city ?? existing.city, + phone: business.phone ?? existing.phone, + email: business.email ?? existing.email, + website: business.website ?? existing.website, + websiteStatus: business.websiteStatus, + confidence: business.confidence, + evidenceSummary: business.evidenceSummary, + missingReason: business.missingReason ?? null, + }); + continue; + } + + await db.createLeadScoutLead({ + workspaceId: campaign.workspaceId, + campaignId, + scanRunId: run.id, + sourceProvider: campaign.provider, + sourceId: business.sourceId, + businessName: business.businessName, + category: business.category, + niche: business.niche, + formattedAddress: business.formattedAddress, + postalCode: business.postalCode, + city: business.city, + phone: business.phone, + email: business.email, + website: business.website, + websiteStatus: business.websiteStatus, + confidence: business.confidence, + evidenceSummary: business.evidenceSummary, + missingReason: business.missingReason, + }); + newLeadCount++; + } + + const completed = await db.updateLeadScoutScanRun(run.id, { + status: "completed", + resultCount: businesses.length, + newLeadCount, + missingWebsiteCount, + summary: `${businesses.length} result(s) — ${newLeadCount} new, ${missingWebsiteCount} missing a website.`, + completedAt: new Date(), + }); + await db.updateLeadScoutCampaign(campaignId, { lastScanAt: new Date() }); + + await logAudit({ + workspaceId: campaign.workspaceId, + actor: "extension", + toolLabel: "local_lead_scout.scan", + input: { campaignId, provider: campaign.provider }, + output: { resultCount: businesses.length, newLeadCount, missingWebsiteCount }, + status: "success", + }); + await notifyWorkspaceOwner(campaign.workspaceId, { + title: "Lead scan complete", + body: `${campaign.name}: ${newLeadCount} new lead(s), ${missingWebsiteCount} missing a website.`, + url: `/workspace/${campaign.workspaceId}/extensions/local-lead-scout`, + tag: `lead-scout-scan-${run.id}`, + }); + + return completed; + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + const failed = await db.updateLeadScoutScanRun(run.id, { + status: "failed", + errorMessage: message, + completedAt: new Date(), + }); + await logAudit({ + workspaceId: campaign.workspaceId, + actor: "extension", + toolLabel: "local_lead_scout.scan", + input: { campaignId, provider: campaign.provider }, + output: message, + status: "error", + }); + return failed; + } +} + +/** Marks a new lead reviewed — the human "I looked at this" gate that + * requireApprovalBeforePrototype checks for before letting a prototype be + * requested. */ +export async function markLeadScoutLeadReviewed(leadId: string): Promise { + const claimed = await getDb().claimLeadScoutLeadStatus({ + id: leadId, + fromStatus: "new", + toStatus: "reviewed", + }); + if (!claimed) throw new Error("Lead isn't in a reviewable state (already reviewed or moved on)."); +} + +async function pickDefaultModelId(workspaceId: string): Promise { + const db = getDb(); + const workspace = await db.getWorkspace(workspaceId); + if (workspace?.defaultModelId) return workspace.defaultModelId; + const providers = await getInstalledProvidersForWorkspace(workspaceId); + const models = await listAvailableModels(providers); + const [first] = models; + if (!first) { + throw new Error( + "No models are installed for this workspace — add one in Settings before generating content.", + ); + } + return first.id; +} + +const LEAD_SCOUT_AGENT_SYSTEM_PROMPT = + "You are the Local Lead Scout content agent. You draft lightweight website concepts and short, " + + "respectful outreach emails for local businesses that currently have no website. You never make " + + "misleading claims, never invent a prior relationship with the business, and never use pressure " + + "tactics. Follow the requested output format exactly."; + +/** Lazily provisions (once) and thereafter reuses the campaign's content + * agent — mirrors configureSeoFixerAgent's "auto-provision unless the user + * pinned one" shape, but this agent needs no file tools (it only writes + * text), so there's no sandboxed tool setup to redo per dispatch. */ +async function configureLeadScoutAgent(campaign: LeadScoutCampaignRecord): Promise { + const db = getDb(); + if (campaign.prototypeAgentId) { + const existing = await db.getAgent(campaign.prototypeAgentId); + if (existing) return existing; + } + const modelId = await pickDefaultModelId(campaign.workspaceId); + const agent = await db.createAgent({ + workspaceId: campaign.workspaceId, + name: `Lead Scout Content Agent — ${campaign.name}`, + systemPrompt: LEAD_SCOUT_AGENT_SYSTEM_PROMPT, + modelId, + autonomyLevel: "assisted", + toolIds: [], + skillIds: [], + }); + await db.updateLeadScoutCampaign(campaign.id, { prototypeAgentId: agent.id }); + return agent; +} + +export function parsePrototypeOutput(output: string): { + concept: string | null; + heroCopy: string | null; + sections: string[]; + callToAction: string | null; + styleDirection: string | null; + artifactMarkdown: string | null; +} { + const concept = output.match(/CONCEPT:\s*(.+)/i)?.[1]?.trim() ?? null; + const heroCopy = output.match(/HERO_COPY:\s*(.+)/i)?.[1]?.trim() ?? null; + const sectionsRaw = output.match(/SECTIONS:\s*(.+)/i)?.[1]?.trim() ?? ""; + const sections = sectionsRaw + ? sectionsRaw + .split(",") + .map((s) => s.trim()) + .filter(Boolean) + : []; + const callToAction = output.match(/CTA:\s*(.+)/i)?.[1]?.trim() ?? null; + const styleDirection = output.match(/STYLE:\s*(.+)/i)?.[1]?.trim() ?? null; + const artifactRaw = output.split(/---ARTIFACT---/i)[1]?.trim() ?? null; + const artifactMarkdown = artifactRaw && artifactRaw.toLowerCase() !== "none" ? artifactRaw : null; + return { concept, heroCopy, sections, callToAction, styleDirection, artifactMarkdown }; +} + +const PROTOTYPE_OUTPUT_FORMAT = [ + "Respond in exactly this structure (plain text):", + "CONCEPT: ", + "HERO_COPY: ", + 'SECTIONS: ', + "CTA: ", + "STYLE: ", + "---ARTIFACT---", + '', +].join("\n"); + +/** + * Dispatches the campaign's content agent to draft a lightweight website + * concept for one lead — an existing NyxelOS agent/task run, not a bespoke + * LLM pipeline (see ADR-style rationale in seo-analyzer.ts's + * configureSeoFixerAgent). Gated on requireApprovalBeforePrototype: if set, + * the lead must already be "reviewed" (a human looked at it) before a + * prototype can be requested at all. + */ +export async function dispatchLeadScoutPrototype( + leadId: string, +): Promise { + const db = getDb(); + const lead = await db.getLeadScoutLead(leadId); + if (!lead) throw new Error(`Unknown lead: ${leadId}`); + const campaign = await db.getLeadScoutCampaign(lead.campaignId); + if (!campaign) throw new Error(`Unknown lead scout campaign: ${lead.campaignId}`); + + if (campaign.requireApprovalBeforePrototype && lead.status !== "reviewed") { + throw new Error('This lead needs to be marked "reviewed" before generating a prototype.'); + } + if (lead.status !== "new" && lead.status !== "reviewed") { + throw new Error(`Lead is "${lead.status}" and can't have a prototype requested right now.`); + } + const claimed = await db.claimLeadScoutLeadStatus({ + id: leadId, + fromStatus: lead.status, + toStatus: "prototype_requested", + }); + if (!claimed) throw new Error("Lead status changed concurrently — refresh and try again."); + + const agent = await configureLeadScoutAgent(campaign); + const instruction = [ + "Generate a lightweight website concept for a local business that currently has no website.", + `Business: ${lead.businessName}`, + lead.category ? `Category/niche: ${lead.category}` : null, + lead.formattedAddress ? `Region: ${lead.formattedAddress}` : null, + `Evidence this business has no website: ${lead.evidenceSummary ?? "none recorded"}`, + "Target style: clean, modern, mobile-first, trustworthy for a local service business.", + "Requested prototype type: single-page marketing site concept.", + "", + PROTOTYPE_OUTPUT_FORMAT, + ] + .filter(Boolean) + .join("\n"); + + const prototype = await db.createLeadScoutPrototype({ workspaceId: lead.workspaceId, leadId }); + const task = await db.createTask({ + workspaceId: lead.workspaceId, + assignedAgentId: agent.id, + title: `Generate prototype — ${lead.businessName}`, + instruction, + input: { leadId, prototypeId: prototype.id }, + }); + + try { + const result = await executeManagedTask({ taskId: task.id, agent, trigger: "extension" }); + const parsed = parsePrototypeOutput(result.output); + const ready = await db.updateLeadScoutPrototype(prototype.id, { + status: "ready", + taskId: task.id, + ...parsed, + }); + await db.claimLeadScoutLeadStatus({ + id: leadId, + fromStatus: "prototype_requested", + toStatus: "prototype_ready", + }); + await logAudit({ + workspaceId: lead.workspaceId, + agentId: agent.id, + actor: "extension", + toolLabel: "local_lead_scout.generate_prototype", + input: { leadId }, + output: { prototypeId: prototype.id }, + status: "success", + }); + return ready; + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + const failed = await db.updateLeadScoutPrototype(prototype.id, { + status: "failed", + taskId: task.id, + errorMessage: message, + }); + // Revert to reviewed so the user can retry instead of being stuck. + await db.claimLeadScoutLeadStatus({ + id: leadId, + fromStatus: "prototype_requested", + toStatus: "reviewed", + }); + await logAudit({ + workspaceId: lead.workspaceId, + agentId: agent.id, + actor: "extension", + toolLabel: "local_lead_scout.generate_prototype", + input: { leadId }, + output: message, + status: "error", + }); + return failed; + } +} + +/** Requires a `ready` prototype — approving a failed/pending one makes no + * sense, there'd be nothing reviewed. */ +export async function approveLeadScoutPrototype( + prototypeId: string, +): Promise { + const db = getDb(); + const prototype = await db.getLeadScoutPrototype(prototypeId); + if (!prototype) throw new Error(`Unknown prototype: ${prototypeId}`); + if (prototype.status !== "ready") throw new Error("Only a ready prototype can be approved."); + const approved = await db.updateLeadScoutPrototype(prototypeId, { approved: true }); + await logAudit({ + workspaceId: prototype.workspaceId, + actor: "extension", + toolLabel: "local_lead_scout.approve_prototype", + input: { prototypeId }, + output: { approved: true }, + status: "success", + }); + return approved; +} + +export function parseEmailDraftOutput(output: string): { subject: string; body: string } { + const subject = output.match(/SUBJECT:\s*(.+)/i)?.[1]?.trim() || "A quick idea for your business"; + const body = output.split(/---BODY---/i)[1]?.trim() || output.trim(); + return { subject, body }; +} + +/** + * Dispatches the campaign's content agent to draft an outreach email for a + * lead, using its approved prototype (never an unapproved one — the + * compliance requirement that a human reviews the concept before it's + * referenced in outreach). Sender identity and opt-out text are appended + * programmatically from email settings after generation, rather than left + * to the model, so they're always present regardless of what it wrote. + */ +export async function dispatchLeadScoutEmailDraft( + leadId: string, +): Promise { + const db = getDb(); + const lead = await db.getLeadScoutLead(leadId); + if (!lead) throw new Error(`Unknown lead: ${leadId}`); + const campaign = await db.getLeadScoutCampaign(lead.campaignId); + if (!campaign) throw new Error(`Unknown lead scout campaign: ${lead.campaignId}`); + + const prototypes = await db.listLeadScoutPrototypesByLead(leadId); + const prototype = prototypes.find((p) => p.approved); + if (!prototype) throw new Error("Approve a prototype for this lead before drafting an email."); + + const emailSettings = await db.getLeadScoutEmailSettings(lead.workspaceId); + if (!emailSettings) throw new Error("Configure email settings before drafting outreach emails."); + + if (lead.status !== "prototype_ready") { + throw new Error(`Lead is "${lead.status}" and can't have an email drafted right now.`); + } + const claimed = await db.claimLeadScoutLeadStatus({ + id: leadId, + fromStatus: "prototype_ready", + toStatus: "email_drafted", + }); + if (!claimed) throw new Error("Lead status changed concurrently — refresh and try again."); + + const agent = await configureLeadScoutAgent(campaign); + const instruction = [ + "Draft a short, respectful, non-pushy cold outreach email to a local business about a website concept it might consider.", + `Business: ${lead.businessName}`, + `Website concept: ${prototype.concept ?? "N/A"}`, + `Hero copy: ${prototype.heroCopy ?? "N/A"}`, + `Call to action: ${prototype.callToAction ?? "N/A"}`, + `Sender: ${emailSettings.fromName}`, + "Rules: no misleading claims, no fake prior relationship, no pressure tactics, under 150 words, one clear next step.", + "Do not include a signature or footer — those are added separately.", + "", + "Respond in exactly this format:", + "SUBJECT: ", + "---BODY---", + "", + ].join("\n"); + + const draft = await db.createLeadScoutOutreachDraft({ + workspaceId: lead.workspaceId, + leadId, + prototypeId: prototype.id, + }); + const task = await db.createTask({ + workspaceId: lead.workspaceId, + assignedAgentId: agent.id, + title: `Draft outreach email — ${lead.businessName}`, + instruction, + input: { leadId, draftId: draft.id }, + }); + + try { + const result = await executeManagedTask({ taskId: task.id, agent, trigger: "extension" }); + const { subject, body } = parseEmailDraftOutput(result.output); + const footer = [emailSettings.unsubscribeText, emailSettings.legalFooter] + .filter(Boolean) + .join("\n\n"); + const bodyText = [body, "", `— ${emailSettings.fromName}`, footer].filter(Boolean).join("\n"); + + const ready = await db.updateLeadScoutOutreachDraft(draft.id, { + status: "draft", + taskId: task.id, + subject, + bodyText, + }); + await logAudit({ + workspaceId: lead.workspaceId, + agentId: agent.id, + actor: "extension", + toolLabel: "local_lead_scout.draft_email", + input: { leadId }, + output: { draftId: draft.id }, + status: "success", + }); + return ready; + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + const failed = await db.updateLeadScoutOutreachDraft(draft.id, { + status: "failed", + taskId: task.id, + errorMessage: message, + }); + await db.claimLeadScoutLeadStatus({ + id: leadId, + fromStatus: "email_drafted", + toStatus: "prototype_ready", + }); + await logAudit({ + workspaceId: lead.workspaceId, + agentId: agent.id, + actor: "extension", + toolLabel: "local_lead_scout.draft_email", + input: { leadId }, + output: message, + status: "error", + }); + return failed; + } +} diff --git a/apps/server/src/scheduler.ts b/apps/server/src/scheduler.ts index 96ad085..b6150b4 100644 --- a/apps/server/src/scheduler.ts +++ b/apps/server/src/scheduler.ts @@ -8,6 +8,7 @@ import { logAudit } from "./audit"; import { emitNyxelEvent } from "./event-bus"; import { NyxelEvent } from "./events"; import { reviewDueGoals } from "./goal-orchestrator"; +import { runLeadScoutScan } from "./lead-scout"; import { notifyWorkspaceOwner } from "./push"; import { runSeoAnalysis } from "./seo-analyzer"; import { runWorkflowAndWait } from "./workflow-runner"; @@ -374,6 +375,47 @@ async function checkDueSeoProjects(): Promise { } } +/** + * Checks every Local Lead Scout campaign with scheduling enabled whose + * nextScanAt has passed and runs its scan — same rationale as + * checkDueSeoProjects (a scan isn't an agent chat turn, so it's polled here + * rather than through the automation table). manual_csv campaigns can't run + * unattended (a scan needs freshly uploaded CSV text each time), so those + * are skipped even if scheduling was left on for one — the UI disables + * enabling a schedule for manual_csv campaigns, but a campaign switched to + * manual_csv after scheduling was enabled would otherwise be stuck retrying + * forever. + */ +async function checkDueLeadScoutCampaigns(): Promise { + const db = getDb(); + let due: Awaited>; + try { + due = await db.listDueLeadScoutCampaigns(new Date()); + } catch (err) { + console.error("Scheduler: failed to query due lead scout campaigns:", err); + return; + } + for (const campaign of due) { + if (campaign.provider === "manual_csv") { + await db.updateLeadScoutCampaign(campaign.id, { scheduleEnabled: false }); + continue; + } + try { + await runLeadScoutScan(campaign.id); + const now = new Date(); + const nextScanAt = campaign.scheduleCronExpression + ? computeNextRunAt(campaign.scheduleCronExpression, now) + : null; + await db.updateLeadScoutCampaign(campaign.id, { nextScanAt }); + } catch (err) { + console.error( + `Scheduler: lead scout scan for "${campaign.name}" (${campaign.id}) failed:`, + err, + ); + } + } +} + /** * Recovers agent runs left stuck in "running" by a process that died or * restarted before it could mark them cancelled/completed/failed itself — @@ -484,6 +526,7 @@ export function startScheduler(): () => void { await checkFileWatchAutomations(); await checkDueSeoProjects(); + await checkDueLeadScoutCampaigns(); await checkGoalsForReview(); await checkStaleAgentRuns(); }, POLL_INTERVAL_MS); diff --git a/apps/server/src/trpc/router.ts b/apps/server/src/trpc/router.ts index b78d09a..98a2b69 100644 --- a/apps/server/src/trpc/router.ts +++ b/apps/server/src/trpc/router.ts @@ -3,6 +3,7 @@ import { DEFAULT_CHAT_WORKING_DIRECTORY, getDb, type KnowledgeBaseConfigRecord, + type LeadScoutSourceConfigRecord, type McpServerRecord, type ModelInstallationRecord, } from "@nyxel/db"; @@ -42,6 +43,27 @@ import { listKnowledgeBaseDocuments, runDocsAgentForWorkspace, } from "../knowledge-base"; +import { + approveLeadScoutPrototype, + dispatchLeadScoutEmailDraft, + dispatchLeadScoutPrototype, + markLeadScoutLeadReviewed, + runLeadScoutScan, +} from "../lead-scout"; +import { + getLeadScoutEmailSettings, + sendLeadScoutTestEmail, + testLeadScoutEmailConnection, + toClientSafeLeadScoutEmailSettings, + upsertLeadScoutEmailSettings, +} from "../lead-scout-email"; +import { + addLeadScoutSuppression, + approveLeadScoutOutreachDraft, + rejectLeadScoutOutreachDraft, + resetLeadScoutLeadForResend, + sendLeadScoutOutreachDraft, +} from "../lead-scout-send"; import { createLibraryFolder, deleteLibraryFile, @@ -439,6 +461,13 @@ export function toClientSafeKnowledgeBaseConfig(config: KnowledgeBaseConfigRecor return { ...rest, obsidianApiKeySet: obsidianApiKey !== null && obsidianApiKey.length > 0 }; } +// Same rationale — a lead scout source config's apiKey (google_places_api / +// custom_api) never needs to leave the server. +export function toClientSafeLeadScoutSourceConfig(config: LeadScoutSourceConfigRecord) { + const { apiKey, ...rest } = config; + return { ...rest, hasApiKey: apiKey !== null && apiKey.length > 0 }; +} + export const appRouter = router({ health: publicProcedure.query(() => ({ ok: true, name: "nyxel-server" })), @@ -2710,6 +2739,327 @@ export const appRouter = router({ }), }), + leadScout: router({ + listCampaigns: workspaceProcedure + .input(z.object({ workspaceId: z.string() })) + .query(({ input }) => getDb().listLeadScoutCampaignsByWorkspace(input.workspaceId)), + createCampaign: workspaceProcedure + .input( + z.object({ + workspaceId: z.string(), + name: z.string().min(1), + postalCode: z.string().min(1), + country: z.string().min(1).optional(), + radiusKm: z.number().positive().max(200).optional(), + niches: z.array(z.string()).optional(), + maxResultsPerRun: z.number().int().min(1).max(200).optional(), + provider: z.enum(["manual_csv", "google_places_api", "osm_overpass", "custom_api"]), + minConfidence: z.number().int().min(0).max(100).optional(), + outreachMode: z.enum(["draft_only", "review_and_send"]).optional(), + }), + ) + .mutation(async ({ input }) => { + const db = getDb(); + const ext = await db.getExtensionByKey(input.workspaceId, "local-lead-scout"); + if (!ext) + throw new Error("The Local Lead Scout extension isn't installed in this workspace."); + const { workspaceId, ...rest } = input; + return db.createLeadScoutCampaign({ workspaceId, extensionId: ext.id, ...rest }); + }), + updateCampaign: protectedProcedure + .input( + z.object({ + id: z.string(), + name: z.string().min(1).optional(), + postalCode: z.string().min(1).optional(), + country: z.string().min(1).optional(), + radiusKm: z.number().positive().max(200).optional(), + niches: z.array(z.string()).optional(), + maxResultsPerRun: z.number().int().min(1).max(200).optional(), + provider: z + .enum(["manual_csv", "google_places_api", "osm_overpass", "custom_api"]) + .optional(), + minConfidence: z.number().int().min(0).max(100).optional(), + outreachMode: z.enum(["draft_only", "review_and_send"]).optional(), + autoGeneratePrototype: z.boolean().optional(), + autoDraftEmail: z.boolean().optional(), + autoSendAfterApproval: z.boolean().optional(), + requireApprovalBeforePrototype: z.boolean().optional(), + requireApprovalBeforeEmailSend: z.boolean().optional(), + prototypeAgentId: z.string().nullable().optional(), + }), + ) + .mutation(async ({ input, ctx }) => { + await requireEntityWorkspaceOwner( + ctx.user.id, + () => getDb().getLeadScoutCampaign(input.id), + "Lead scout campaign not found", + ); + const { id, ...patch } = input; + return getDb().updateLeadScoutCampaign(id, patch); + }), + deleteCampaign: protectedProcedure + .input(z.object({ id: z.string() })) + .mutation(async ({ input, ctx }) => { + await requireEntityWorkspaceOwner( + ctx.user.id, + () => getDb().getLeadScoutCampaign(input.id), + "Lead scout campaign not found", + ); + return getDb().deleteLeadScoutCampaign(input.id); + }), + setCampaignSchedule: protectedProcedure + .input(z.object({ id: z.string(), cronExpression: z.string().nullable() })) + .mutation(async ({ input, ctx }) => { + const campaign = await requireEntityWorkspaceOwner( + ctx.user.id, + () => getDb().getLeadScoutCampaign(input.id), + "Lead scout campaign not found", + ); + if (input.cronExpression === null) { + return getDb().updateLeadScoutCampaign(input.id, { + scheduleEnabled: false, + scheduleCronExpression: null, + nextScanAt: null, + }); + } + if (campaign.provider === "manual_csv") { + throw new Error( + "manual_csv campaigns can't run on a schedule — there's no CSV to read automatically.", + ); + } + const nextScanAt = computeNextRunAt(input.cronExpression, new Date()); + if (!nextScanAt) + throw new Error(`"${input.cronExpression}" is not a valid cron expression.`); + return getDb().updateLeadScoutCampaign(input.id, { + scheduleEnabled: true, + scheduleCronExpression: input.cronExpression, + nextScanAt, + }); + }), + + listSourceConfigs: workspaceProcedure + .input(z.object({ workspaceId: z.string() })) + .query(async ({ input }) => { + const configs = await getDb().listLeadScoutSourceConfigsByWorkspace(input.workspaceId); + return configs.map(toClientSafeLeadScoutSourceConfig); + }), + upsertSourceConfig: workspaceProcedure + .input( + z.object({ + workspaceId: z.string(), + provider: z.enum(["manual_csv", "google_places_api", "osm_overpass", "custom_api"]), + config: z.record(z.string(), z.unknown()).optional(), + apiKey: z.string().nullable().optional(), + enabled: z.boolean().optional(), + }), + ) + .mutation(async ({ input }) => { + const record = await getDb().upsertLeadScoutSourceConfig(input); + await logAudit({ + workspaceId: input.workspaceId, + actor: "extension", + toolLabel: "local_lead_scout.source_config.update", + input: { provider: input.provider }, + output: { hasApiKey: Boolean(input.apiKey) }, + status: "success", + }); + return toClientSafeLeadScoutSourceConfig(record); + }), + + runScan: protectedProcedure + .input(z.object({ campaignId: z.string(), csvText: z.string().optional() })) + .mutation(async ({ input, ctx }) => { + await requireEntityWorkspaceOwner( + ctx.user.id, + () => getDb().getLeadScoutCampaign(input.campaignId), + "Lead scout campaign not found", + ); + return runLeadScoutScan(input.campaignId, { csvText: input.csvText }); + }), + listScanRuns: protectedProcedure + .input(z.object({ campaignId: z.string() })) + .query(async ({ input, ctx }) => { + await requireEntityWorkspaceOwner( + ctx.user.id, + () => getDb().getLeadScoutCampaign(input.campaignId), + "Lead scout campaign not found", + ); + return getDb().listLeadScoutScanRunsByCampaign(input.campaignId); + }), + + listLeads: protectedProcedure + .input(z.object({ campaignId: z.string() })) + .query(async ({ input, ctx }) => { + await requireEntityWorkspaceOwner( + ctx.user.id, + () => getDb().getLeadScoutCampaign(input.campaignId), + "Lead scout campaign not found", + ); + return getDb().listLeadScoutLeadsByCampaign(input.campaignId); + }), + getLead: protectedProcedure + .input(z.object({ id: z.string() })) + .query(async ({ input, ctx }) => { + return requireEntityWorkspaceOwner( + ctx.user.id, + () => getDb().getLeadScoutLead(input.id), + "Lead not found", + ); + }), + markLeadReviewed: protectedProcedure + .input(z.object({ id: z.string() })) + .mutation(async ({ input, ctx }) => { + await requireEntityWorkspaceOwner( + ctx.user.id, + () => getDb().getLeadScoutLead(input.id), + "Lead not found", + ); + await markLeadScoutLeadReviewed(input.id); + return getDb().getLeadScoutLead(input.id); + }), + resetLeadForResend: protectedProcedure + .input(z.object({ id: z.string() })) + .mutation(async ({ input, ctx }) => { + await requireEntityWorkspaceOwner( + ctx.user.id, + () => getDb().getLeadScoutLead(input.id), + "Lead not found", + ); + return resetLeadScoutLeadForResend(input.id); + }), + + listPrototypes: protectedProcedure + .input(z.object({ leadId: z.string() })) + .query(async ({ input, ctx }) => { + await requireEntityWorkspaceOwner( + ctx.user.id, + () => getDb().getLeadScoutLead(input.leadId), + "Lead not found", + ); + return getDb().listLeadScoutPrototypesByLead(input.leadId); + }), + generatePrototype: protectedProcedure + .input(z.object({ leadId: z.string() })) + .mutation(async ({ input, ctx }) => { + await requireEntityWorkspaceOwner( + ctx.user.id, + () => getDb().getLeadScoutLead(input.leadId), + "Lead not found", + ); + return dispatchLeadScoutPrototype(input.leadId); + }), + approvePrototype: protectedProcedure + .input(z.object({ id: z.string() })) + .mutation(async ({ input, ctx }) => { + await requireEntityWorkspaceOwner( + ctx.user.id, + () => getDb().getLeadScoutPrototype(input.id), + "Prototype not found", + ); + return approveLeadScoutPrototype(input.id); + }), + + listDrafts: protectedProcedure + .input(z.object({ leadId: z.string() })) + .query(async ({ input, ctx }) => { + await requireEntityWorkspaceOwner( + ctx.user.id, + () => getDb().getLeadScoutLead(input.leadId), + "Lead not found", + ); + return getDb().listLeadScoutOutreachDraftsByLead(input.leadId); + }), + generateDraft: protectedProcedure + .input(z.object({ leadId: z.string() })) + .mutation(async ({ input, ctx }) => { + await requireEntityWorkspaceOwner( + ctx.user.id, + () => getDb().getLeadScoutLead(input.leadId), + "Lead not found", + ); + return dispatchLeadScoutEmailDraft(input.leadId); + }), + approveDraft: protectedProcedure + .input(z.object({ id: z.string() })) + .mutation(async ({ input, ctx }) => { + await requireEntityWorkspaceOwner( + ctx.user.id, + () => getDb().getLeadScoutOutreachDraft(input.id), + "Outreach draft not found", + ); + return approveLeadScoutOutreachDraft(input.id); + }), + rejectDraft: protectedProcedure + .input(z.object({ id: z.string() })) + .mutation(async ({ input, ctx }) => { + await requireEntityWorkspaceOwner( + ctx.user.id, + () => getDb().getLeadScoutOutreachDraft(input.id), + "Outreach draft not found", + ); + return rejectLeadScoutOutreachDraft(input.id); + }), + sendDraft: protectedProcedure + .input(z.object({ id: z.string() })) + .mutation(async ({ input, ctx }) => { + await requireEntityWorkspaceOwner( + ctx.user.id, + () => getDb().getLeadScoutOutreachDraft(input.id), + "Outreach draft not found", + ); + return sendLeadScoutOutreachDraft(input.id); + }), + + listSuppressions: workspaceProcedure + .input(z.object({ workspaceId: z.string() })) + .query(({ input }) => getDb().listLeadScoutSuppressionsByWorkspace(input.workspaceId)), + addSuppression: workspaceProcedure + .input( + z.object({ + workspaceId: z.string(), + email: z.string().email().nullable().optional(), + domain: z.string().nullable().optional(), + reason: z.string().min(1), + }), + ) + .mutation(({ input }) => addLeadScoutSuppression(input)), + + getEmailSettings: workspaceProcedure + .input(z.object({ workspaceId: z.string() })) + .query(async ({ input }) => { + const settings = await getLeadScoutEmailSettings(input.workspaceId); + return settings ? toClientSafeLeadScoutEmailSettings(settings) : null; + }), + upsertEmailSettings: workspaceProcedure + .input( + z.object({ + workspaceId: z.string(), + provider: z.enum(["smtp", "resend", "mailgun", "custom"]).optional(), + fromName: z.string().min(1), + fromEmail: z.string().email(), + replyTo: z.string().email().nullable().optional(), + credentials: z.record(z.string(), z.string()).nullable().optional(), + dailySendLimit: z.number().int().min(0).max(1000).optional(), + perCampaignSendLimit: z.number().int().min(0).max(1000).optional(), + dryRunMode: z.boolean().optional(), + legalFooter: z.string().nullable().optional(), + unsubscribeText: z.string().min(1).optional(), + }), + ) + .mutation(async ({ input }) => { + const { workspaceId, ...rest } = input; + const settings = await upsertLeadScoutEmailSettings(workspaceId, rest); + return toClientSafeLeadScoutEmailSettings(settings); + }), + testEmailConnection: workspaceProcedure + .input(z.object({ workspaceId: z.string() })) + .mutation(({ input }) => testLeadScoutEmailConnection(input.workspaceId)), + sendTestEmail: workspaceProcedure + .input(z.object({ workspaceId: z.string(), toEmail: z.string().email() })) + .mutation(({ input }) => sendLeadScoutTestEmail(input.workspaceId, input.toEmail)), + }), + automations: router({ list: workspaceProcedure .input(z.object({ workspaceId: z.string() })) diff --git a/apps/web/src/app/workspace/[workspaceId]/extensions/[key]/page.tsx b/apps/web/src/app/workspace/[workspaceId]/extensions/[key]/page.tsx index 244e480..25a952e 100644 --- a/apps/web/src/app/workspace/[workspaceId]/extensions/[key]/page.tsx +++ b/apps/web/src/app/workspace/[workspaceId]/extensions/[key]/page.tsx @@ -3,10 +3,11 @@ import { useQuery } from "@tanstack/react-query"; import { useParams } from "next/navigation"; import type { ComponentType } from "react"; -import { PageHeaderSkeleton } from "@/components/loading"; -import { PageHeader } from "@/components/page-header"; +import { LocalLeadScoutExtensionPage } from "@/components/extensions/local-lead-scout-page"; import { SeoAnalyzerExtensionPage } from "@/components/extensions/seo-analyzer-page"; import { VideoStudioExtensionPage } from "@/components/extensions/video-studio-page"; +import { PageHeaderSkeleton } from "@/components/loading"; +import { PageHeader } from "@/components/page-header"; import { trpcClient } from "@/lib/trpc"; /** Per-extension page component, keyed by ExtensionCatalogEntry.route (see @@ -15,6 +16,7 @@ import { trpcClient } from "@/lib/trpc"; const EXTENSION_PAGES: Record> = { "seo-analyzer": SeoAnalyzerExtensionPage, "video-studio": VideoStudioExtensionPage, + "local-lead-scout": LocalLeadScoutExtensionPage, }; export default function ExtensionDetailPage() { diff --git a/apps/web/src/components/app-sidebar.tsx b/apps/web/src/components/app-sidebar.tsx index 5475ec3..b038cae 100644 --- a/apps/web/src/components/app-sidebar.tsx +++ b/apps/web/src/components/app-sidebar.tsx @@ -15,6 +15,7 @@ import { Images, LayoutDashboard, Library, + MapPin, MessageSquare, Package, Plug, @@ -56,6 +57,7 @@ import { useInstallation } from "@/lib/use-installation"; const EXTENSION_ICON_MAP: Record = { TrendingUp, Film, + MapPin, }; export function AppSidebar() { diff --git a/apps/web/src/components/extensions/local-lead-scout-page.tsx b/apps/web/src/components/extensions/local-lead-scout-page.tsx new file mode 100644 index 0000000..00a7436 --- /dev/null +++ b/apps/web/src/components/extensions/local-lead-scout-page.tsx @@ -0,0 +1,934 @@ +"use client"; + +import { useMutation, useQuery, useQueryClient } from "@tanstack/react-query"; +import { useEffect, useState } from "react"; +import { PageHeader } from "@/components/page-header"; +import { Badge } from "@/components/ui/badge"; +import { Button } from "@/components/ui/button"; +import { Dialog, DialogContent, DialogHeader, DialogTitle } from "@/components/ui/dialog"; +import { Input } from "@/components/ui/input"; +import { Label } from "@/components/ui/label"; +import { Switch } from "@/components/ui/switch"; +import { Tabs, TabsContent, TabsList, TabsTrigger } from "@/components/ui/tabs"; +import { Textarea } from "@/components/ui/textarea"; +import { + type LeadScoutCampaignSummary, + type LeadScoutDraftStatus, + type LeadScoutEmailProvider, + type LeadScoutLeadStatus, + type LeadScoutLeadSummary, + type LeadScoutProvider, + type LeadScoutWebsiteStatus, + trpcClient, +} from "@/lib/trpc"; + +const PROVIDER_LABEL: Record = { + manual_csv: "Manual CSV import", + google_places_api: "Google Places API", + osm_overpass: "OSM Overpass (free)", + custom_api: "Custom API", +}; + +const WEBSITE_STATUS_BADGE: Record = { + unknown: "bg-muted text-muted-foreground", + has_website: "bg-emerald-500/15 text-emerald-600", + missing_website: "bg-amber-500/15 text-amber-600", + invalid_website: "bg-destructive/15 text-destructive", +}; + +const LEAD_STATUS_LABEL: Record = { + new: "New", + reviewed: "Reviewed", + prototype_requested: "Prototype pending", + prototype_ready: "Prototype ready", + email_drafted: "Email draft pending", + approved_to_send: "Approval pending", + sending: "Sending", + sent: "Sent", + rejected: "Rejected", + suppressed: "Suppressed", +}; + +const DRAFT_STATUS_LABEL: Record = { + draft: "Draft", + approved: "Approved — ready to send", + rejected: "Rejected", + sending: "Sending", + sent: "Sent", + failed: "Failed", +}; + +function CreateCampaignForm({ + workspaceId, + onCreated, +}: { + workspaceId: string; + onCreated: () => void; +}) { + const [name, setName] = useState(""); + const [postalCode, setPostalCode] = useState(""); + const [country, setCountry] = useState("US"); + const [radiusKm, setRadiusKm] = useState("10"); + const [niches, setNiches] = useState(""); + const [maxResults, setMaxResults] = useState("25"); + const [provider, setProvider] = useState("manual_csv"); + const [minConfidence, setMinConfidence] = useState("50"); + + const createCampaign = useMutation({ + mutationFn: () => + trpcClient.leadScout.createCampaign.mutate({ + workspaceId, + name, + postalCode, + country, + radiusKm: Number(radiusKm), + niches: niches + .split(",") + .map((n) => n.trim()) + .filter(Boolean), + maxResultsPerRun: Number(maxResults), + provider, + minConfidence: Number(minConfidence), + }), + onSuccess: () => { + setName(""); + setPostalCode(""); + setNiches(""); + onCreated(); + }, + }); + + return ( +
+
+

New campaign

+

+ Finds local businesses with no website in a region, using a compliant source only. You are + responsible for lawful outreach and consent requirements for whatever you contact. +

+
+
+
+ + setName(e.target.value)} /> +
+
+ + setPostalCode(e.target.value)} + /> +
+
+ + setCountry(e.target.value)} /> +
+
+ + setRadiusKm(e.target.value)} + /> +
+
+ + setNiches(e.target.value)} + /> +
+
+ + +
+
+ + setMaxResults(e.target.value)} + /> +
+
+ + setMinConfidence(e.target.value)} + /> +
+
+ {provider === "google_places_api" && ( +

+ You are responsible for complying with the Google Maps Platform Terms of Service for any + data retrieved through this provider. Configure an API key in Source Settings first. +

+ )} + {createCampaign.error && ( +

{(createCampaign.error as Error).message}

+ )} + +
+ ); +} + +function RunScanControl({ campaign }: { campaign: LeadScoutCampaignSummary }) { + const queryClient = useQueryClient(); + const [csvText, setCsvText] = useState(""); + + const runScan = useMutation({ + mutationFn: () => + trpcClient.leadScout.runScan.mutate({ + campaignId: campaign.id, + csvText: campaign.provider === "manual_csv" ? csvText : undefined, + }), + onSuccess: () => { + queryClient.invalidateQueries({ queryKey: ["leadScout", "scanRuns", campaign.id] }); + queryClient.invalidateQueries({ queryKey: ["leadScout", "leads", campaign.id] }); + }, + }); + + return ( +
+

Run scan — {PROVIDER_LABEL[campaign.provider]}

+ {campaign.provider === "manual_csv" ? ( + <> +

+ CSV columns: businessName, address, postalCode, city, category, phone, email, website, + notes. +

+