OUTREACH-STARTER — vollstaendiges Paket in einer Datei Geruest fuer ein B2B-Kaltakquise-Tool: Leads finden, Mailadresse ermitteln, getaktet anschreiben, Antworten einlesen, Abmeldungen behandeln. Node.js + Postgres. Keine API-Keys, keine Daten enthalten. So verwendest du diese Datei: gib sie deinem Coding-Assistenten mit der Anweisung, daraus die Ordnerstruktur anzulegen. Jede Datei beginnt mit einer Zeile '===== DATEI: ====='. ================================================================== ===== DATEI: README.md ===== # Outreach-Starter Gerüst für ein eigenes B2B-Kaltakquise-Tool: Leads finden, E-Mail-Adresse ermitteln, getaktet anschreiben, Antworten einlesen und klassifizieren, Abmeldungen sauber behandeln. Node.js + Postgres, keine fremde SaaS-Plattform dazwischen. Gedacht als Startpunkt zum Weiterbauen (auch mit einem Coding-Assistenten), nicht als fertiges Produkt — es fehlt bewusst ein Admin-UI und jede Authentifizierung. ## Die Pipeline Fünf Stufen, jede schreibt ihr Ergebnis in die DB und setzt einen Status. Dadurch ist jede Stufe einzeln wiederholbar, ohne die vorige neu zu laufen: ``` harvest Branchenverzeichnis abfragen → prospects (status 'new') enrich Website nach Impressum/Kontakt scrapen → 'enriched' | 'invalid' campaign getaktet versenden → messages, 'contacted' replyMon IMAP pollen, Antwort zuordnen → replies, 'replied' | 'unsubscribed' optOut Widerspruch erkennen + sperren → suppression (dauerhaft) ``` | Datei | Aufgabe | |---|---| | `src/jobs/harvest.js` | Google Places Text Search + Details → neue Prospects | | `src/jobs/enrich.js` | `/impressum`, `/kontakt` … abklappern, Mailadresse parsen | | `src/jobs/campaign.js` | Versand mit Drosselung, Dry-Run, Doppelversand-Schutz | | `src/jobs/replyMonitor.js` | IMAP-Poll, Absender → Prospect, Klassifikation | | `src/jobs/scheduler.js` | Zeitsteuerung ohne Cron-Lib | | `src/services/impressumScraper.js` | Adress-Extraktion + Rausfiltern von Müll-Treffern | | `src/services/optOut.js` | Widerspruchserkennung, Sperrliste, Löschfristen | | `src/services/replyClassifier.js` | positive / negative / Abwesenheit / Abmeldung | | `src/services/mxCheck.js` | MX-Prüfung per DNS → unzustellbare Domains vorab aussortieren | | `src/services/mailer.js` | SMTP-Versand inkl. `List-Unsubscribe` | | `src/services/templates.js` | `{{platzhalter}}` in Betreff/Body | | `migrations/*.sql` | Schema: prospects, campaigns, messages, replies, suppression | ## Loslegen ```bash npm install cp .env.example .env # ausfüllen psql "$DATABASE_URL" -f migrations/001_outreach_schema.sql psql "$DATABASE_URL" -f migrations/002_sent_emails.sql npm start # /health, POST /jobs/harvest|enrich|campaign ``` `CAMPAIGN_DRY_RUN=true` ist Absicht. Erst auf `false`, wenn ein Testversand an die eigene Adresse sauber aussah — Betreff, Umlaute, Abmeldelink, Fußzeile. Eine Kampagne ist eine Zeile in `outreach.campaigns` (Betreff + Body-Template mit `{{name}}`, `{{city}}`, `{{website}}`). `active = TRUE` schaltet sie scharf. ## Was man sich hier spart (teuer gelernt) Die Stellen, an denen so ein Tool in der Praxis kaputtgeht — alle im Code kommentiert: 1. **Die eigene Fußzeile löst Abmeldungen aus.** Wer in der Mail um ein „kein Interesse" bittet, hat diese Wörter in jeder zitierten Antwort. Ohne Herausfiltern der eigenen Textbausteine (`EIGENE_BAUSTEINE` in `optOut.js`) wird jeder Interessent, der mit Zitat antwortet, gesperrt. Das ist der mit Abstand teuerste Fehler des ganzen Systems. 2. **Abmeldungen kommen nicht als „unsubscribe".** Echte Leute schreiben „Nein, danke", „Bitte löschen", „kein Interesse". Eine Kurzliste aus `unsubscribe|opt out|remove me` erkennt davon fast nichts. Die Musterliste in `optOut.js` ist das Ergebnis mehrerer Durchläufe mit echten Antworten. 3. **Im Zweifel sperren.** Eine falsch gesperrte Adresse kostet einen Kontakt. Eine falsch *nicht* gesperrte kostet Bußgeld und die Domain-Reputation. Deshalb werden Prospects nach Widerspruch erst mit Verzögerung gelöscht, die Sperrliste greift aber sofort — und sie überlebt das Löschen, sonst landet dieselbe Adresse beim nächsten Harvest wieder im Versand. 4. **Ein Webserver ohne MX-Eintrag ist ein toter Briefkasten.** `mxCheck.js` sortiert die vorab aus, statt Bounces zu sammeln. Bewusst kein Rückfall auf den A-Eintrag, auch wenn RFC 5321 ihn erlaubt. 5. **Scraper finden den falschen Treffer.** CMS-Platzhalter, `noreply@`, Agentur-Adressen im Footer, Tracking-Domains. Der Scraper bevorzugt Adressen auf der Website-Domain und davon `info@`/`kontakt@`. 6. **Drosseln.** 30 Mails/Stunde aus einer frischen Domain ist schon viel. Vorher SPF, DKIM und DMARC einrichten, sonst landet alles im Spam — das ist Voraussetzung, nicht Optimierung. 7. **Doppelversand.** Das `UNIQUE (prospect_id, campaign_id, step_no)` in `messages` ist die eigentliche Absicherung. Ein Skript, das zweimal läuft, schreibt dann nicht zweimal raus. ## Rechtliches (DE/EU, kurz) Kaltakquise per E-Mail an Firmen ist nach § 7 UWG nur in engen Grenzen zulässig, und die DSGVO gilt zusätzlich. Was dieses Gerüst technisch unterstützt — und was man selbst ergänzen muss: - funktionierender Abmeldeweg in **jeder** Mail (`List-Unsubscribe` ist gesetzt; der Link selbst muss gebaut werden, siehe `unsubToken.js`) - Sperrliste, die Abmeldungen dauerhaft respektiert (`suppression`) - Löschfristen für Kontakte ohne Reaktion (`scheduleRetention` / `purgeExpired`, Standard 12 Monate) - Hinweis nach Art. 14 DSGVO, woher die Daten stammen — gehört in die Mail - Impressum und Datenschutzerklärung, deren Fristen zum Code passen Das ist keine Rechtsberatung, und die Bewertung ist je nach Branche unterschiedlich. Vor dem ersten echten Versand kurz prüfen lassen. ## Was absichtlich fehlt - **Auth**: Die `/jobs/*`-Routen sind offen. Vor dem Deployment hinter ein Passwort legen oder rauswerfen und per CLI starten. - **Admin-UI**, Bounce-Verarbeitung über Webhooks, Follow-up-Stufen (`step_no = 2`), Mehrsprachigkeit, Reporting. - **Eine zweite Lead-Quelle.** Google Places ist kostenpflichtig und liefert nicht jede Branche gut. Ein Branchenverzeichnis-Scraper ist oft die bessere Ergänzung — `harvest.js` ist die Stelle dafür. Die Opt-out- und Scraper-Module sind der Teil, der wirklich Zeit gekostet hat. Der Rest ist bewusst einfach gehalten und darf ersetzt werden. ===== DATEI: package.json ===== { "name": "outreach-starter", "version": "0.1.0", "description": "Starter: Lead-Harvest, E-Mail-Anreicherung, getakteter Versand, Antwort-Erkennung, Opt-out", "main": "src/index.js", "scripts": { "start": "node src/index.js", "dev": "node --watch src/index.js" }, "dependencies": { "@googlemaps/google-maps-services-js": "^3.4.0", "axios": "^1.17.0", "cheerio": "^1.0.0", "dotenv": "^16.4.5", "express": "^4.19.2", "helmet": "^8.2.0", "imapflow": "^1.0.171", "mailparser": "^3.7.2", "nodemailer": "^8.0.8", "p-throttle": "^6.2.0", "pg": "8.12.0" } } ===== DATEI: .gitignore ===== node_modules/ .env *.log tmp/ ===== DATEI: .env.example ===== # --- Server --- PORT=3001 NODE_ENV=production # true = Jobs laufen NICHT automatisch (zum Entwickeln empfohlen) SCHEDULER_DISABLED=true # --- Datenbank (Postgres) --- DATABASE_URL=postgres://user:password@localhost:5432/outreach # true bei Managed-Postgres mit TLS DB_SSL=false # --- Lead-Quelle: Google Places --- # ACHTUNG: kostenpflichtig. Billing-Alarm im Google-Cloud-Projekt setzen, # sonst laeuft ein Harvest-Lauf im Hintergrund ins Geld. GOOGLE_PLACES_API_KEY= # JSON-Array von Suchanfragen. Pro Eintrag ein Text-Search-Aufruf. HARVEST_QUERIES=["Sicherheitsschulung Berlin","Sachkundepruefung Hamburg"] HARVEST_MAX_PER_QUERY=60 # --- SMTP (Versand) --- SMTP_HOST=smtp.example.com SMTP_PORT=465 SMTP_USER=info@deine-domain.de SMTP_PASS= SMTP_FROM_NAME=Deine Firma SMTP_FROM_EMAIL=info@deine-domain.de REPLY_TO=info@deine-domain.de # Mails pro Stunde. Niedrig anfangen (20-30), sonst leidet die Domain-Reputation. CAMPAIGN_THROTTLE_PER_HOUR=30 # true = nichts versenden, nur loggen. ERST auf false stellen, wenn ein # Testversand an die eigene Adresse sauber aussah. CAMPAIGN_DRY_RUN=true # --- IMAP (Antworten einlesen) --- IMAP_HOST=imap.example.com IMAP_PORT=993 IMAP_USER=info@deine-domain.de IMAP_PASS= IMAP_MAILBOX=INBOX # --- Zeitplan --- SCHEDULE_HARVEST_AT=06:00 SCHEDULE_REPLY_POLL_MINUTES=15 # --- Abmeldelinks (optional, siehe unsubToken.js) --- UNSUB_SECRET= ===== DATEI: migrations/001_outreach_schema.sql ===== -- Outreach service schema — replaces n8n outreach workflows -- Lives in a separate schema to keep chatbot data cleanly isolated. CREATE SCHEMA IF NOT EXISTS outreach; -- ── prospects ─────────────────────────────────────────────────────────────── -- A prospect is a business we discovered (via Google Places) and may contact. CREATE TABLE IF NOT EXISTS outreach.prospects ( id BIGSERIAL PRIMARY KEY, source TEXT NOT NULL DEFAULT 'google_places', place_id TEXT UNIQUE, -- Google Places ID, NULL for other sources name TEXT NOT NULL, city TEXT, address TEXT, phone TEXT, website TEXT, email TEXT, -- 'new' → harvested, no enrichment yet -- 'enriched' → email found, ready for campaign -- 'contacted' → at least one outreach message sent -- 'replied' → reply received -- 'invalid' → enrichment failed (no email/website), do not contact -- 'unsubscribed'→ explicit opt-out, never contact again status TEXT NOT NULL DEFAULT 'new' CHECK (status IN ('new','enriched','contacted','replied','invalid','unsubscribed')), harvested_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), enriched_at TIMESTAMPTZ, contacted_at TIMESTAMPTZ, replied_at TIMESTAMPTZ, notes TEXT ); CREATE INDEX IF NOT EXISTS prospects_status_idx ON outreach.prospects (status); CREATE INDEX IF NOT EXISTS prospects_city_idx ON outreach.prospects (city); -- Unique email when set, to prevent duplicate contact across harvests CREATE UNIQUE INDEX IF NOT EXISTS prospects_email_unique_idx ON outreach.prospects (LOWER(email)) WHERE email IS NOT NULL; -- ── campaigns ─────────────────────────────────────────────────────────────── -- A campaign bundles a subject + body template + throttling rules. CREATE TABLE IF NOT EXISTS outreach.campaigns ( id BIGSERIAL PRIMARY KEY, name TEXT NOT NULL UNIQUE, subject TEXT NOT NULL, body_template_plain TEXT NOT NULL, body_template_html TEXT, throttle_per_hour INT NOT NULL DEFAULT 30, active BOOLEAN NOT NULL DEFAULT FALSE, created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() ); -- ── messages ──────────────────────────────────────────────────────────────── -- One row per actual SMTP send (or attempted send in dry-run mode). CREATE TABLE IF NOT EXISTS outreach.messages ( id BIGSERIAL PRIMARY KEY, prospect_id BIGINT NOT NULL REFERENCES outreach.prospects(id) ON DELETE CASCADE, campaign_id BIGINT NOT NULL REFERENCES outreach.campaigns(id) ON DELETE RESTRICT, step_no INT NOT NULL DEFAULT 1, -- 1 = initial, 2 = follow-up, … sent_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), smtp_message_id TEXT, -- value of Message-ID header status TEXT NOT NULL DEFAULT 'sent' CHECK (status IN ('sent','dry_run','bounced','failed')), error TEXT, UNIQUE (prospect_id, campaign_id, step_no) ); CREATE INDEX IF NOT EXISTS messages_prospect_idx ON outreach.messages (prospect_id); CREATE INDEX IF NOT EXISTS messages_sent_at_idx ON outreach.messages (sent_at); -- ── replies ───────────────────────────────────────────────────────────────── -- Inbound mail that matched a prospect. CREATE TABLE IF NOT EXISTS outreach.replies ( id BIGSERIAL PRIMARY KEY, message_id BIGINT REFERENCES outreach.messages(id) ON DELETE SET NULL, prospect_id BIGINT REFERENCES outreach.prospects(id) ON DELETE SET NULL, from_addr TEXT NOT NULL, subject TEXT, body_plain TEXT, received_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), -- 'positive' → interested -- 'negative' → not interested (but no opt-out) -- 'ooo' → auto-reply / out-of-office -- 'unsubscribe' → asks to be removed -- 'unknown' → could not classify, needs manual review classification TEXT NOT NULL DEFAULT 'unknown' CHECK (classification IN ('positive','negative','ooo','unsubscribe','unknown')), processed BOOLEAN NOT NULL DEFAULT FALSE ); CREATE INDEX IF NOT EXISTS replies_prospect_idx ON outreach.replies (prospect_id); CREATE INDEX IF NOT EXISTS replies_received_idx ON outreach.replies (received_at); -- ── kv ────────────────────────────────────────────────────────────────────── -- Tiny key/value table for job state (e.g. last IMAP poll timestamp). CREATE TABLE IF NOT EXISTS outreach.kv ( key TEXT PRIMARY KEY, value TEXT NOT NULL, updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() ); ===== DATEI: migrations/002_sent_emails.sql ===== -- Speichert JEDEN über den Outreach-Service versendeten (oder versuchten) E-Mail-Versand -- inklusive vollständigem Body — damit nachvollziehbar ist, was genau rausging. CREATE TABLE IF NOT EXISTS outreach.sent_emails ( id BIGSERIAL PRIMARY KEY, to_addr TEXT NOT NULL, subject TEXT, body_text TEXT, body_html TEXT, smtp_message_id TEXT, in_reply_to TEXT, headers JSONB NOT NULL DEFAULT '{}'::jsonb, meta JSONB NOT NULL DEFAULT '{}'::jsonb, -- frei: prospect_id, campaign, source, school_id ... status TEXT NOT NULL DEFAULT 'sent' CHECK (status IN ('sent','failed','dry_run')), error TEXT, sent_at TIMESTAMPTZ NOT NULL DEFAULT NOW() ); CREATE INDEX IF NOT EXISTS sent_emails_sent_at_idx ON outreach.sent_emails (sent_at); CREATE INDEX IF NOT EXISTS sent_emails_to_idx ON outreach.sent_emails (LOWER(to_addr)); ===== DATEI: src/index.js ===== require('dotenv').config(); const express = require('express'); const helmet = require('helmet'); const log = require('./logger'); const { startScheduler } = require('./jobs/scheduler'); const app = express(); app.use(helmet()); app.use(express.json({ limit: '256kb' })); app.set('trust proxy', 1); app.get('/health', (_req, res) => res.json({ ok: true, service: 'outreach-service', time: new Date().toISOString() })); // Manuelle Auslöser — praktisch zum Testen einzelner Stufen. // ACHTUNG: vor dem Deployment hinter Auth legen (hier bewusst weggelassen). app.post('/jobs/harvest', async (_req, res) => { try { res.json(await require('./jobs/harvest').runHarvest()); } catch (e) { res.status(500).json({ error: e.message }); } }); app.post('/jobs/enrich', async (_req, res) => { try { res.json(await require('./jobs/enrich').runEnrich()); } catch (e) { res.status(500).json({ error: e.message }); } }); app.post('/jobs/campaign', async (_req, res) => { try { res.json(await require('./jobs/campaign').runCampaign()); } catch (e) { res.status(500).json({ error: e.message }); } }); const port = parseInt(process.env.PORT || '3001'); app.listen(port, '0.0.0.0', () => { log.info('server', 'listening', { port }); if (process.env.SCHEDULER_DISABLED !== 'true') { try { startScheduler(); } catch (err) { log.error('server', 'scheduler failed to start', { err: err.message }); } } else { log.warn('server', 'scheduler disabled via env'); } }); process.on('unhandledRejection', (err) => { log.error('server', 'unhandledRejection', { err: err?.message || String(err) }); }); process.on('uncaughtException', (err) => { log.error('server', 'uncaughtException', { err: err?.message || String(err) }); process.exit(1); }); ===== DATEI: src/db.js ===== require('dotenv').config(); const { Pool } = require('pg'); // Ein einziger Pool fuer den ganzen Prozess. Postgres (lokal, Supabase, RDS …). // // Hinweis zu Managed-Postgres mit eigener Root-CA (z. B. Supabase): dort // verifiziert der System-Truststore das Zertifikat nicht. Dann entweder die // CA-Kette des Anbieters hier als `ca` hinterlegen (sauber) oder // `?sslmode=require` in der URL nutzen. NICHT rejectUnauthorized:false setzen — // damit ist die Verbindung verschluesselt, aber nicht mehr authentifiziert. const useSsl = /^(1|true|yes)$/i.test(process.env.DB_SSL || ''); const pool = new Pool({ connectionString: process.env.DATABASE_URL, ssl: useSsl ? { rejectUnauthorized: true } : false, max: 10, idleTimeoutMillis: 30000, connectionTimeoutMillis: 5000, // Erfahrungswert: in Docker-Containern ohne IPv6-Route zum DB-Host sonst // minutenlange Timeouts. Dann IPv4 erzwingen: // family: 4, }); module.exports = pool; ===== DATEI: src/logger.js ===== function log(level, module, msg, ctx) { const entry = Object.assign({ ts: new Date().toISOString(), level, module, msg }, ctx || {}); const line = JSON.stringify(entry); if (level === 'error') process.stderr.write(line + '\n'); else process.stdout.write(line + '\n'); } module.exports = { info: (module, msg, ctx) => log('info', module, msg, ctx), warn: (module, msg, ctx) => log('warn', module, msg, ctx), error: (module, msg, ctx) => log('error', module, msg, ctx), }; ===== DATEI: src/jobs/scheduler.js ===== // Scheduler — runs the automatic outreach jobs on simple time-of-day triggers. // (no node-cron lib — we keep dependencies minimal). // // Only harvest + reply monitoring are automatic. Enrichment, demo-building and sending // // HARVEST_AT: "HH:MM" in Europe/Berlin // REPLY_POLL_MINUTES: poll interval in minutes require('dotenv').config(); const log = require('../logger'); const { runHarvest } = require('./harvest'); const { runReplyMonitor } = require('./replyMonitor'); const TZ = 'Europe/Berlin'; const DAY_MS = 24 * 60 * 60 * 1000; function parseHHMM(str, defaultHHMM) { const m = /^(\d{1,2}):(\d{2})$/.exec(str || defaultHHMM); if (!m) throw new Error(`bad HH:MM: ${str}`); return { h: parseInt(m[1]), m: parseInt(m[2]) }; } // Compute ms until next occurrence of "HH:MM" in the given timezone. function msUntilNext(hhmm) { const now = new Date(); // Translate "now" into Europe/Berlin to find the next H:M. // Cheap approach: use toLocaleString and reparse. const localStr = now.toLocaleString('en-US', { timeZone: TZ, hour12: false }); // localStr like "10/26/2026, 14:37:21" const [datePart, timePart] = localStr.split(', '); const [mm, dd, yyyy] = datePart.split('/').map(Number); const [hh, mi, ss] = timePart.split(':').map(Number); const localNow = new Date(Date.UTC(yyyy, mm - 1, dd, hh, mi, ss)); const localTarget = new Date(Date.UTC(yyyy, mm - 1, dd, hhmm.h, hhmm.m, 0)); let diff = localTarget - localNow; if (diff <= 0) diff += DAY_MS; return diff; } function scheduleDaily(name, hhmmStr, defaultHHMM, fn) { const hhmm = parseHHMM(hhmmStr, defaultHHMM); function arm() { const delay = msUntilNext(hhmm); log.info('scheduler', 'next run scheduled', { job: name, in_minutes: Math.round(delay / 60000), at: `${hhmm.h}:${String(hhmm.m).padStart(2, '0')} ${TZ}` }); setTimeout(async () => { try { await fn(); } catch (err) { log.error('scheduler', 'job failed', { job: name, err: err.message }); } arm(); // re-arm for tomorrow }, delay); } arm(); } function scheduleInterval(name, minutes, fn) { const ms = Math.max(1, minutes) * 60 * 1000; log.info('scheduler', 'interval scheduled', { job: name, every_minutes: minutes }); // Run after 30s startup delay, then every `minutes` setTimeout(async () => { try { await fn(); } catch (err) { log.error('scheduler', 'job failed', { job: name, err: err.message }); } setInterval(async () => { try { await fn(); } catch (err) { log.error('scheduler', 'job failed', { job: name, err: err.message }); } }, ms); }, 30 * 1000); } function startScheduler() { // Only harvest + reply monitoring run automatically. Enrichment, demo-building and the // outreach send sind manuelle Schritte (manuell bzw. per Skript). scheduleDaily('harvest', process.env.SCHEDULE_HARVEST_AT, '06:00', runHarvest); scheduleInterval('replyMonitor', parseInt(process.env.SCHEDULE_REPLY_POLL_MINUTES || '15'), runReplyMonitor); log.info('scheduler', 'started'); } module.exports = { startScheduler }; ===== DATEI: src/jobs/harvest.js ===== // Harvest job — Teil der Outreach-Pipeline // Reads HARVEST_QUERIES from env, queries Google Places Text Search + Details, // and inserts new prospects into outreach.prospects. require('dotenv').config(); const db = require('../db'); const log = require('../logger'); const { textSearch, placeDetails, extractCity } = require('../services/googlePlaces'); function parseQueryEnv(envName) { const raw = process.env[envName] || '[]'; try { const arr = JSON.parse(raw); return Array.isArray(arr) ? arr.filter(q => typeof q === 'string' && q.trim()) : []; } catch { log.error('harvest', `${envName} not valid JSON`, { raw }); return []; } } // Zwei Quellen (Migration 003): HARVEST_QUERIES = Deutschland (Bestand), // HARVEST_QUERIES_FR = Frankreich-Kampagne ("auto-école Lyon" …, language/region fr). function parseQueries() { return [ ...parseQueryEnv('HARVEST_QUERIES').map(q => ({ query: q, country: 'DE', language: 'de', region: 'de' })), ...parseQueryEnv('HARVEST_QUERIES_FR').map(q => ({ query: q, country: 'FR', language: 'fr', region: 'fr' })), ]; } async function harvestOne({ query, country, language, region }) { const maxPerQuery = parseInt(process.env.HARVEST_MAX_PER_QUERY || '60'); const results = (await textSearch(query, maxPerQuery, { language, region })).slice(0, maxPerQuery); let inserted = 0; let skipped = 0; for (const r of results) { if (!r.place_id) { skipped++; continue; } // Details API gives us website + phone (Text Search doesn't include these reliably) let details = null; try { details = await placeDetails(r.place_id, { language }); } catch (err) { log.warn('harvest', 'placeDetails failed', { place_id: r.place_id, err: err.message }); } const name = details?.name || r.name; const address = details?.formatted_address || r.formatted_address || null; const phone = details?.international_phone_number || null; const website = details?.website || null; const city = extractCity(details?.address_components); const { rowCount } = await db.query( `INSERT INTO outreach.prospects (source, place_id, name, city, address, phone, website, status, country) VALUES ('google_places', $1, $2, $3, $4, $5, $6, 'new', $7) ON CONFLICT (place_id) DO NOTHING`, [r.place_id, name, city, address, phone, website, country] ); if (rowCount > 0) inserted++; else skipped++; } return { query, country, total: results.length, inserted, skipped }; } async function runHarvest() { const queries = parseQueries(); if (queries.length === 0) { log.warn('harvest', 'no queries configured, skipping'); return { runs: [] }; } log.info('harvest', 'start', { queries: queries.map(q => q.query) }); const runs = []; for (const q of queries) { try { const r = await harvestOne(q); runs.push(r); log.info('harvest', 'query done', r); } catch (err) { log.error('harvest', 'query failed', { query: q.query, err: err.message }); runs.push({ query: q.query, error: err.message }); } } log.info('harvest', 'complete', { runs }); return { runs }; } module.exports = { runHarvest }; ===== DATEI: src/jobs/enrich.js ===== // Enrich job — Teil der Outreach-Pipeline // For each prospect with status='new' and a website, fetches /impressum, /kontakt // and parses an email address. Sets status='enriched' or 'invalid'. require('dotenv').config(); const db = require('../db'); const log = require('../logger'); const { findEmail } = require('../services/impressumScraper'); const BATCH_SIZE = 50; async function enrichOne(prospect) { const email = await findEmail(prospect.website); if (!email) { await db.query( `UPDATE outreach.prospects SET status = 'invalid', enriched_at = NOW(), notes = COALESCE(notes,'') || 'no email found via scrape; ' WHERE id = $1`, [prospect.id] ); return { id: prospect.id, ok: false }; } try { await db.query( `UPDATE outreach.prospects SET email = $1, status = 'enriched', enriched_at = NOW() WHERE id = $2`, [email, prospect.id] ); return { id: prospect.id, ok: true, email }; } catch (err) { if (err.code === '23505') { // duplicate email (unique index) — another prospect already owns it await db.query( `UPDATE outreach.prospects SET status = 'invalid', enriched_at = NOW(), notes = COALESCE(notes,'') || 'duplicate email; ' WHERE id = $1`, [prospect.id] ); return { id: prospect.id, ok: false, reason: 'duplicate' }; } throw err; } } async function runEnrich() { const { rows } = await db.query( `SELECT id, name, website FROM outreach.prospects WHERE status = 'new' AND website IS NOT NULL ORDER BY harvested_at ASC LIMIT $1`, [BATCH_SIZE] ); log.info('enrich', 'start', { batch: rows.length }); let enriched = 0, invalid = 0; for (const p of rows) { try { const r = await enrichOne(p); if (r.ok) enriched++; else invalid++; } catch (err) { log.error('enrich', 'failed', { id: p.id, err: err.message }); invalid++; } } // Also mark new prospects without website as invalid (no enrichment path) await db.query( `UPDATE outreach.prospects SET status = 'invalid', enriched_at = NOW(), notes = COALESCE(notes,'') || 'no website; ' WHERE status = 'new' AND website IS NULL` ); const result = { processed: rows.length, enriched, invalid }; log.info('enrich', 'complete', result); return result; } module.exports = { runEnrich }; ===== DATEI: src/jobs/campaign.js ===== // Campaign job — Teil der Outreach-Pipeline // Sends step_no=1 messages for each active campaign to enriched prospects // who haven't been contacted by that campaign yet. Honors: // - CAMPAIGN_THROTTLE_PER_HOUR (env, default 30) // - CAMPAIGN_DRY_RUN (env=true → logs only, marks status='dry_run') // - prospect.status MUST be 'enriched' (not 'unsubscribed' / 'replied' / etc.) require('dotenv').config(); const pThrottle = require('p-throttle'); const db = require('../db'); const log = require('../logger'); const { sendOutreachMail } = require('../services/mailer'); const { buildMail } = require('../services/templates'); const DRY_RUN = String(process.env.CAMPAIGN_DRY_RUN || 'true').toLowerCase() === 'true'; const THROTTLE_PER_HOUR = parseInt(process.env.CAMPAIGN_THROTTLE_PER_HOUR || '30'); async function getActiveCampaigns() { const { rows } = await db.query( `SELECT id, name, subject, body_template_plain, body_template_html, throttle_per_hour FROM outreach.campaigns WHERE active = TRUE ORDER BY id` ); return rows; } async function getCandidates(campaignId, limit) { // prospects that are 'enriched', have email, and no step_no=1 message for this campaign yet const { rows } = await db.query( `SELECT p.id, p.name, p.email, p.city, p.website FROM outreach.prospects p WHERE p.status = 'enriched' AND p.email IS NOT NULL AND NOT EXISTS ( SELECT 1 FROM outreach.messages m WHERE m.prospect_id = p.id AND m.campaign_id = $1 AND m.step_no = 1 ) ORDER BY p.enriched_at ASC LIMIT $2`, [campaignId, limit] ); return rows; } async function recordSend(prospect, campaign, status, smtpMessageId, errMsg) { await db.query( `INSERT INTO outreach.messages (prospect_id, campaign_id, step_no, status, smtp_message_id, error) VALUES ($1, $2, 1, $3, $4, $5) ON CONFLICT (prospect_id, campaign_id, step_no) DO NOTHING`, [prospect.id, campaign.id, status, smtpMessageId || null, errMsg || null] ); if (status === 'sent' || status === 'dry_run') { await db.query( `UPDATE outreach.prospects SET status = 'contacted', contacted_at = NOW() WHERE id = $1`, [prospect.id] ); } } async function runCampaign() { const campaigns = await getActiveCampaigns(); if (campaigns.length === 0) { log.info('campaign', 'no active campaigns, skipping'); return { sent: 0, dry: 0, failed: 0 }; } log.info('campaign', 'start', { campaigns: campaigns.map(c => c.name), dry_run: DRY_RUN }); let sent = 0, dry = 0, failed = 0; for (const c of campaigns) { const limit = c.throttle_per_hour || THROTTLE_PER_HOUR; const throttle = pThrottle({ limit: 1, interval: Math.ceil(3600 * 1000 / limit) }); const candidates = await getCandidates(c.id, limit); log.info('campaign', 'candidates', { campaign: c.name, count: candidates.length }); const sendOne = throttle(async (prospect) => { const { subject, text, html } = buildMail(c, prospect); if (DRY_RUN) { log.info('campaign', 'DRY-RUN send', { to: prospect.email, subject }); await recordSend(prospect, c, 'dry_run', null, null); dry++; return; } try { const info = await sendOutreachMail({ to: prospect.email, subject, text, html }); await recordSend(prospect, c, 'sent', info.messageId, null); sent++; } catch (err) { log.error('campaign', 'send failed', { to: prospect.email, err: err.message }); await recordSend(prospect, c, 'failed', null, err.message); failed++; } }); await Promise.all(candidates.map(sendOne)); } const result = { sent, dry, failed }; log.info('campaign', 'complete', result); return result; } module.exports = { runCampaign }; ===== DATEI: src/jobs/replyMonitor.js ===== // Reply Monitor — Teil der Outreach-Pipeline // Polls the IMAP inbox (read-only) for new mail since the last poll, // matches sender against outreach.prospects, classifies the reply, and // updates prospect status. Last-poll timestamp is stored in outreach.kv. require('dotenv').config(); const db = require('../db'); const log = require('../logger'); const { fetchMessagesSince } = require('../services/imap'); const { classify } = require('../services/replyClassifier'); const KV_KEY = 'reply_monitor_last_poll'; const DEFAULT_LOOKBACK_DAYS = 1; async function getLastPoll() { const { rows } = await db.query(`SELECT value FROM outreach.kv WHERE key = $1`, [KV_KEY]); if (rows.length === 0 || !rows[0].value) { const d = new Date(); d.setDate(d.getDate() - DEFAULT_LOOKBACK_DAYS); return d; } return new Date(rows[0].value); } async function setLastPoll(ts) { await db.query( `INSERT INTO outreach.kv (key, value, updated_at) VALUES ($1, $2, NOW()) ON CONFLICT (key) DO UPDATE SET value = $2, updated_at = NOW()`, [KV_KEY, ts.toISOString()] ); } async function findProspectByEmail(addr) { if (!addr) return null; const { rows } = await db.query( `SELECT id, status FROM outreach.prospects WHERE LOWER(email) = LOWER($1) LIMIT 1`, [addr] ); return rows[0] || null; } async function findLatestMessage(prospectId) { const { rows } = await db.query( `SELECT id FROM outreach.messages WHERE prospect_id = $1 ORDER BY sent_at DESC LIMIT 1`, [prospectId] ); return rows[0]?.id || null; } async function runReplyMonitor() { const since = await getLastPoll(); const pollStart = new Date(); let messages; try { messages = await fetchMessagesSince(since); } catch (err) { log.error('replyMonitor', 'imap fetch failed', { err: err.message }); return { matched: 0, error: err.message }; } let matched = 0; for (const m of messages) { if (!m.from) continue; const prospect = await findProspectByEmail(m.from); if (!prospect) continue; // not from a tracked prospect — ignore const classification = classify({ subject: m.subject || '', body: m.text || '', headers: m.headers }); const messageId = await findLatestMessage(prospect.id); await db.query( `INSERT INTO outreach.replies (message_id, prospect_id, from_addr, subject, body_plain, received_at, classification, processed) VALUES ($1, $2, $3, $4, $5, $6, $7, FALSE)`, [messageId, prospect.id, m.from, m.subject || null, (m.text || '').slice(0, 8000), m.date, classification] ); // Status-Update: // - unsubscribe → 'unsubscribed' (final, never contact again) // - any reply → 'replied' // - ooo bleibt formal 'replied' damit Mensch es sieht const newStatus = classification === 'unsubscribe' ? 'unsubscribed' : 'replied'; await db.query( `UPDATE outreach.prospects SET status = $1, replied_at = NOW() WHERE id = $2`, [newStatus, prospect.id] ); matched++; } await setLastPoll(pollStart); const result = { since: since.toISOString(), checked: messages.length, matched }; log.info('replyMonitor', 'complete', result); return result; } module.exports = { runReplyMonitor }; ===== DATEI: src/services/googlePlaces.js ===== require('dotenv').config(); const axios = require('axios'); const log = require('../logger'); const API_KEY = process.env.GOOGLE_PLACES_API_KEY; const BASE = 'https://maps.googleapis.com/maps/api/place'; // language/region als Options-Parameter (Default 'de' = Bestandsverhalten); // Frankreich-Harvest ruft mit { language: 'fr', region: 'fr' } auf. async function textSearch(query, maxResults = 60, { language = 'de', region = 'de' } = {}) { if (!API_KEY) throw new Error('GOOGLE_PLACES_API_KEY missing'); const results = []; let pageToken = null; let page = 0; do { const params = pageToken ? { pagetoken: pageToken, key: API_KEY } : { query, key: API_KEY, language, region }; let res; try { res = await axios.get(`${BASE}/textsearch/json`, { params, timeout: 10000 }); } catch (e) { log.warn('googlePlaces', 'textSearch page error', { query, page, err: e.message }); break; } const st = res.data && res.data.status; if (st === 'ZERO_RESULTS') break; if (st && st !== 'OK') { log.warn('googlePlaces', 'textSearch status', { query, page, status: st }); break; } results.push(...((res.data && res.data.results) || [])); pageToken = (res.data && res.data.next_page_token) || null; page++; if (results.length >= maxResults) break; if (pageToken) await new Promise(r => setTimeout(r, 2500)); } while (pageToken && page < 3); log.info('googlePlaces', 'textSearch done', { query, count: results.length }); return results; } async function placeDetails(placeId, { language = 'de' } = {}) { if (!API_KEY) throw new Error('GOOGLE_PLACES_API_KEY missing'); const res = await axios.get(`${BASE}/details/json`, { params: { place_id: placeId, key: API_KEY, language, fields: 'name,formatted_address,international_phone_number,website,address_components', }, timeout: 10000, }); return (res.data && res.data.result) || null; } function extractCity(addressComponents = []) { const c = (addressComponents || []).find(x => x.types.includes('locality') || x.types.includes('postal_town')); return c ? c.long_name : null; } module.exports = { textSearch, placeDetails, extractCity }; ===== DATEI: src/services/impressumScraper.js ===== const cheerio = require('cheerio'); const log = require('../logger'); const FETCH_TIMEOUT_MS = 8000; const UA = 'Mozilla/5.0 (compatible; OutreachBot/1.0; +https://example.com/bot)'; // Pages to probe in order. First hit wins. const CANDIDATE_PATHS = ['/impressum', '/impressum/', '/kontakt', '/kontakt/', '/contact', '/']; // Reject email addresses that are clearly noise (CMS placeholders, agencies, etc.) const REJECT_PATTERNS = [ /@example\.(com|de|org)$/i, /@(domain|yourdomain|test|placeholder)\./i, /@(google|googlemail|gmail)\.com$/i, // generic personal /@(sentry|cloudflare|cdn|gstatic|googleusercontent)\./i, /no-?reply@/i, ]; const EMAIL_RE = /[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2,}/g; async function fetchWithTimeout(url) { const ctrl = new AbortController(); const t = setTimeout(() => ctrl.abort(), FETCH_TIMEOUT_MS); try { const res = await fetch(url, { headers: { 'User-Agent': UA, 'Accept': 'text/html,application/xhtml+xml' }, signal: ctrl.signal, redirect: 'follow', }); if (!res.ok) return null; return await res.text(); } catch (err) { return null; } finally { clearTimeout(t); } } function extractEmails(html) { if (!html) return []; const $ = cheerio.load(html); // Drop scripts/styles before stringifying — they often contain noisy fake emails $('script, style, noscript').remove(); const cleanedHtml = $.html(); // Collect mailto hrefs explicitly (highest signal) plus the raw HTML body with tags // collapsed to whitespace — that way "info@x.deHans" doesn't fuse into "info@x.deHans". const mailtos = $('a[href^="mailto:"]').map((_, a) => ($(a).attr('href') || '').replace(/^mailto:/, '').split('?')[0]).get(); const text = mailtos.join(' ') + ' ' + cleanedHtml.replace(/<[^>]*>/g, ' '); const matches = text.match(EMAIL_RE) || []; // Dedupe + filter const seen = new Set(); const out = []; for (const raw of matches) { const e = raw.toLowerCase(); if (seen.has(e)) continue; seen.add(e); if (REJECT_PATTERNS.some(re => re.test(e))) continue; out.push(e); } return out; } function pickBestEmail(emails, baseHost) { if (emails.length === 0) return null; const host = (baseHost || '').toLowerCase(); // 1. Prefer emails whose domain matches the website host const sameDomain = emails.filter(e => host && (e.endsWith('@' + host) || e.endsWith('.' + host))); if (sameDomain.length) { // info@/kontakt@ first const generic = sameDomain.find(e => /^(info|kontakt|hallo|office|mail|service)@/.test(e)); return generic || sameDomain[0]; } // 2. Otherwise first generic, else first const generic = emails.find(e => /^(info|kontakt|hallo|office|mail|service)@/.test(e)); return generic || emails[0]; } function normalizeBaseUrl(website) { if (!website) return null; try { const u = new URL(website); return { origin: u.origin, host: u.hostname.replace(/^www\./, '') }; } catch { return null; } } async function findEmail(website) { const base = normalizeBaseUrl(website); if (!base) return null; for (const path of CANDIDATE_PATHS) { const url = base.origin + path; const html = await fetchWithTimeout(url); if (!html) continue; const emails = extractEmails(html); const best = pickBestEmail(emails, base.host); if (best) { log.info('impressumScraper', 'email found', { url, email: best }); return best; } } return null; } module.exports = { findEmail, extractEmails, pickBestEmail }; ===== DATEI: src/services/mxCheck.js ===== /** * MX-Pruefung einer Maildomain per DNS. Kostenlos, kein Fremddienst, und * fuer den Gegenueber unsichtbar — anders als eine SMTP-Einzelpruefung, die * wie Adress-Abklopfen aussieht. * * Rueckgabe: * true kein MX-Eintrag -> sicher unzustellbar * false MX vorhanden * null unklar (Timeout, SERVFAIL) -> Problem des Aufloesers, nicht der * Domain; der Aufrufer soll die Adresse im Zweifel behalten * * Bewusst KEIN Rueckfall auf den A-Eintrag (RFC 5321 erlaubt ihn): in der * Praxis hatten mehrere Domains einen A-Eintrag und keinen MX — alle * bounced. Ein Webserver ohne MX ist ein toter Briefkasten. */ const dns = require('dns').promises; async function ohneMx(domain, timeoutMs = 5000) { if (!domain) return null; let timer; const timeout = new Promise((_, rej) => { timer = setTimeout(() => rej(Object.assign(new Error('timeout'), { code: 'ETIMEOUT' })), timeoutMs); }); try { const mx = await Promise.race([dns.resolveMx(domain), timeout]); return !mx || !mx.length; } catch (e) { if (e.code === 'ENOTFOUND' || e.code === 'ENODATA') return true; return null; } finally { clearTimeout(timer); } } // Mehrere Domains, jede nur einmal, mit begrenzter Parallelitaet. // Liefert Map domain -> true/false/null. async function pruefeDomains(domains, { parallel = 10 } = {}) { const liste = [...new Set(domains.filter(Boolean).map(d => d.toLowerCase()))]; const ergebnis = new Map(); let i = 0; async function arbeiter() { while (i < liste.length) { const d = liste[i++]; ergebnis.set(d, await ohneMx(d)); } } await Promise.all(Array.from({ length: Math.min(parallel, liste.length || 1) }, arbeiter)); return ergebnis; } const domainVon = email => String(email || '').toLowerCase().split('@')[1] || ''; module.exports = { ohneMx, pruefeDomains, domainVon }; ===== DATEI: src/services/mailer.js ===== require('dotenv').config(); const nodemailer = require('nodemailer'); const log = require('../logger'); let transporter = null; function getTransporter() { if (transporter) return transporter; transporter = nodemailer.createTransport({ host: process.env.SMTP_HOST, port: parseInt(process.env.SMTP_PORT || '465'), secure: process.env.SMTP_PORT === '465', auth: { user: process.env.SMTP_USER, pass: process.env.SMTP_PASS, }, }); return transporter; } async function sendOutreachMail({ to, subject, text, html }) { const fromEmail = process.env.SMTP_FROM_EMAIL || process.env.SMTP_USER; const fromName = process.env.SMTP_FROM_NAME || 'Your Company'; const replyTo = process.env.REPLY_TO || fromEmail; const info = await getTransporter().sendMail({ from: `"${fromName}" <${fromEmail}>`, to, replyTo, subject, text, html: html || undefined, headers: { // Helps inbox-providers recognize legitimate B2B mail. Required by RFC 8058 for one-click unsubscribe. 'List-Unsubscribe': ``, 'List-Unsubscribe-Post': 'List-Unsubscribe=One-Click', }, }); log.info('mailer', 'sent', { to, messageId: info.messageId }); return info; } module.exports = { sendOutreachMail }; ===== DATEI: src/services/templates.js ===== // Simple {{var}} template renderer. Unknown vars resolve to empty string. function render(template, vars = {}) { return String(template || '').replace(/\{\{\s*([a-zA-Z0-9_]+)\s*\}\}/g, (_, key) => { const v = vars[key]; return v == null ? '' : String(v); }); } function buildMail(campaign, prospect) { const vars = { name: prospect.name || '', city: prospect.city || '', website: prospect.website || '', email: prospect.email || '', unsubscribe_email: process.env.REPLY_TO || process.env.SMTP_FROM_EMAIL || '', }; return { subject: render(campaign.subject, vars), text: render(campaign.body_template_plain, vars), html: campaign.body_template_html ? render(campaign.body_template_html, vars) : null, }; } module.exports = { render, buildMail }; ===== DATEI: src/services/imap.js ===== require('dotenv').config(); const { ImapFlow } = require('imapflow'); const { simpleParser } = require('mailparser'); const log = require('../logger'); function buildClient() { return new ImapFlow({ host: process.env.IMAP_HOST, port: parseInt(process.env.IMAP_PORT || '993'), secure: true, auth: { user: process.env.IMAP_USER, pass: process.env.IMAP_PASS, }, logger: false, }); } // Fetch all messages received since `sinceDate` (Date object). // Returns array of { uid, from, subject, date, text, html, inReplyTo, references }. // IMPORTANT: read-only — never deletes or moves messages. async function fetchMessagesSince(sinceDate) { const client = buildClient(); await client.connect(); const mailbox = process.env.IMAP_MAILBOX || 'INBOX'; const out = []; try { const lock = await client.getMailboxLock(mailbox, { readOnly: true }); try { // IMAP SEARCH SINCE — coarse-grained (date only), so we filter again client-side const uids = await client.search({ since: sinceDate }); if (!uids || uids.length === 0) return []; for await (const msg of client.fetch(uids, { envelope: true, source: true, internalDate: true })) { if (msg.internalDate < sinceDate) continue; const parsed = await simpleParser(msg.source); out.push({ uid: msg.uid, from: parsed.from?.value?.[0]?.address?.toLowerCase() || null, fromName: parsed.from?.value?.[0]?.name || null, subject: parsed.subject || null, date: parsed.date || msg.internalDate, text: parsed.text || '', html: parsed.html || null, inReplyTo: parsed.inReplyTo || null, references: parsed.references || null, headers: parsed.headers, }); } } finally { lock.release(); } } finally { await client.logout().catch(() => {}); } log.info('imap', 'fetched', { count: out.length, since: sinceDate.toISOString() }); return out; } module.exports = { fetchMessagesSince }; ===== DATEI: src/services/replyClassifier.js ===== // Lightweight rule-based classifier for inbound replies. // Returns one of: 'positive' | 'negative' | 'ooo' | 'unsubscribe' | 'unknown' const UNSUBSCRIBE_PATTERNS = [ /\bunsubscribe\b/i, /\babbestellen\b/i, /\bbitte (entfernen|streichen|löschen)\b/i, /\baustragen\b/i, /\bkeine (weiteren )?(e-?mails?|nachrichten|werbung)\b/i, /\bnicht (mehr )?kontaktieren\b/i, ]; const NEGATIVE_PATTERNS = [ /\bnicht interessiert\b/i, /\bkein interesse\b/i, /\bdanke,? aber\b/i, /\bpassen?\s+(nicht|wir nicht)\b/i, /\bwir lehnen ab\b/i, ]; const POSITIVE_PATTERNS = [ /\bgerne\b/i, /\b(rufen sie|telefon|telko|zoom|teams)\b/i, /\binteressiert\b/i, /\b(wann|wie viel|kosten|preis|preise|angebot)\b/i, /\bja,?\s/i, ]; function isOutOfOffice(headers, subject = '', body = '') { // X-Auto-Response-Suppress / Auto-Submitted: auto-replied const autoSubmitted = headers?.get?.('auto-submitted'); if (autoSubmitted && /auto-(replied|generated)/i.test(autoSubmitted)) return true; const precedence = headers?.get?.('precedence'); if (precedence && /(auto_reply|bulk|junk)/i.test(precedence)) return true; if (/\b(out of office|abwesend|abwesenheit|urlaub|auto-?reply|automatische antwort)\b/i.test(subject)) return true; if (/\b(bin (im urlaub|abwesend)|out of office|currently away|until [a-z]+ \d+)\b/i.test(body)) return true; return false; } function classify({ subject = '', body = '', headers = null }) { const text = (subject + '\n' + body).trim(); if (isOutOfOffice(headers, subject, body)) return 'ooo'; if (UNSUBSCRIBE_PATTERNS.some(re => re.test(text))) return 'unsubscribe'; if (NEGATIVE_PATTERNS.some(re => re.test(text))) return 'negative'; if (POSITIVE_PATTERNS.some(re => re.test(text))) return 'positive'; return 'unknown'; } module.exports = { classify }; ===== DATEI: src/services/optOut.js ===== // Widerspruchs- und Loeschlogik fuer die Kaltakquise. // // Grundsatz: im Zweifel sperren. Eine faelschlich gesperrte Adresse kostet einen // verlorenen Kontakt. Eine faelschlich NICHT gesperrte Adresse kostet nach // Art. L34-5 CPCE bis zu 375 EUR pro weiterer Nachricht — und die Reputation // der Sendedomain gleich mit. require('dotenv').config(); const db = require('../db'); const log = require('../logger'); // Aufbewahrungsfrist ohne Reaktion. Dieser Wert und die Angabe in der // eigenen Datenschutzerklaerung muessen uebereinstimmen. const PURGE_MONTHS = 12; // Nach einem Widerspruch wird der Prospect-Datensatz nicht sofort geloescht, // sondern nach 7 Tagen. Die Sperrliste blockt den Versand ab der Sekunde des // Widerspruchs und dauerhaft — daran aendert die Frist nichts. Sie existiert, // weil die Erkennung im Zweifel sperrt: ein faelschlich gesperrter Interessent // waere sonst beim naechsten Retention-Lauf unwiederbringlich weg. const OPTOUT_GRACE_DAYS = 7; const OPT_OUT_PATTERNS = [ // Franzoesisch /\bd[ée]sabonn/i, /\bd[ée]sinscri/i, // "ne plus recevoir", aber auch "ne souhaite plus recevoir", // "ne veux plus recevoir", "ne plus etre contacte" … /\bne\s+.{0,30}?\bplus\s+.{0,25}?\brecevoir\b/i, // Kein \b vor "être": JavaScript zaehlt "ê" nicht als Wortzeichen, // dort entsteht also gar keine Wortgrenze und das Muster liefe ins Leere. /\bne\s+.{0,30}?\bplus\s+.{0,25}?[êe]tre\s+contact/i, /\bplus\s+de\s+(mails?|e-?mails?|messages?|sollicitations?)\b/i, /\bretire[zr]\b.{0,30}\bliste\b/i, /\bretirez[- ]moi\b/i, /\bsupprim(ez|er)\b.{0,30}\b(donn[ée]es|adresse|coordonn[ée]es)\b/i, /\bne\s+me\s+(re)?contactez\s+plus\b/i, /\bpas\s+int[ée]ress[ée]/i, /\bcessez\b.{0,25}\b(envo|mail|message)/i, // Deutsch /\babmelden\b/i, /\baustragen\b/i, // Erfahrungswert: "Informationen/Infos/Angebote/Werbung" ergaenzt — "Bitte Keine // weiteren Informationen" (lag im Junk-Ordner) rutschte als // gewoehnliche Antwort durch. Nur MIT "weitere(n)": "ich finde keine // Informationen zum Preis" ist eine Frage, keine Absage. /\bkeine\s+weiteren?\s+(mails?|e-?mails?|nachrichten|informationen|infos?|angebote?|werbung|werbemails?)/i, /\bnicht\s+mehr\s+(schreiben|kontaktieren|anschreiben)/i, // Englisch. Die Kurzliste (unsubscribe/opt out/remove me) fiel in der Abnahme // bei 12 von 12 realistischen Absagen durch — echte Empfaenger schreiben // "Not interested, thanks." und nicht "I hereby unsubscribe". // Negativer Lookbehind: "an/our/the unsubscribe link" ist eine Erwaehnung // des eigenen Systems, kein Widerspruch an uns. /(?") t = t.split('\n').filter(z => !/^\s*>/.test(z)).join('\n'); // Eigene Bausteine, egal ob zitiert oder nicht — faengt auch die Programme, // die ohne ">" zitieren (Outlook). for (const re of EIGENE_BAUSTEINE) t = t.replace(re, ''); return t; } // "STOP" nur als eigenstaendige Nachricht werten, nicht als Wort im Fliesstext // ("stop" kommt in franzoesischen Saetzen sonst kaum vor, aber sicher ist sicher). function isBareStop(text) { return /^[\s>*_-]*stop[\s.!]*$/i.test((text || '').trim().slice(0, 40)); } function detectOptOut(subject = '', body = '') { const eigen = stripEigenes(body); const haystack = `${subject}\n${eigen}`; if (isBareStop(eigen) || isBareStop(subject)) return { hit: true, marker: 'STOP' }; for (const re of OPT_OUT_PATTERNS) { const m = haystack.match(re); if (m) return { hit: true, marker: m[0].slice(0, 40) }; } return { hit: false, marker: null }; } function normalize(email) { return String(email || '').trim().toLowerCase(); } // Vor JEDEM Versand aufrufen. async function isSuppressed(email) { const e = normalize(email); if (!e) return true; // keine Adresse => nicht senden const { rows } = await db.query( `SELECT 1 FROM outreach.suppression WHERE email = $1 LIMIT 1`, [e]); return rows.length > 0; } // Mehrere Adressen auf einmal pruefen — fuer die Vorauswahl einer Sendecharge. async function filterSuppressed(emails) { const list = [...new Set((emails || []).map(normalize).filter(Boolean))]; if (!list.length) return new Set(); const { rows } = await db.query( `SELECT email FROM outreach.suppression WHERE email = ANY($1::text[])`, [list]); return new Set(rows.map(r => r.email)); } // Sperrt die Adresse dauerhaft und markiert den Prospect. async function suppress(email, { reason = 'opt_out', source = null, note = null } = {}) { const e = normalize(email); if (!e) return false; await db.query( `INSERT INTO outreach.suppression (email, reason, source, note) VALUES ($1, $2, $3, $4) ON CONFLICT (email) DO NOTHING`, [e, reason, source, note]); const { rowCount } = await db.query( `UPDATE outreach.prospects SET opted_out_at = COALESCE(opted_out_at, now()), status = 'unsubscribed', purge_after = now() + interval '7 days' WHERE lower(email) = $1`, [e]); log.info('optOut', 'suppressed', { email: e, reason, source, prospects: rowCount }); return true; } // Prueft eine eingegangene Antwort und sperrt bei Bedarf. // Rueckgabe: true, wenn gesperrt wurde. async function handleReply({ from, subject, body }) { const { hit, marker } = detectOptOut(subject, body); if (!hit) return false; await suppress(from, { reason: 'opt_out', source: 'reply', note: `Marker: ${marker}` }); return true; } // Setzt die Loeschfrist fuer alle kontaktierten Prospects ohne Antwort. async function scheduleRetention() { const { rowCount } = await db.query( `UPDATE outreach.prospects SET purge_after = contacted_at + ($1 || ' months')::interval WHERE contacted_at IS NOT NULL AND replied_at IS NULL AND opted_out_at IS NULL AND purge_after IS NULL`, [String(PURGE_MONTHS)]); if (rowCount) log.info('optOut', 'retention scheduled', { rows: rowCount, months: PURGE_MONTHS }); return rowCount; } // Loescht abgelaufene Prospect-Daten. Die Sperrliste bleibt bestehen — // sonst wuerde eine abgemeldete Adresse spaeter erneut angeschrieben. async function purgeExpired() { const { rows } = await db.query( `DELETE FROM outreach.prospects WHERE purge_after IS NOT NULL AND purge_after < now() RETURNING id`); if (rows.length) log.info('optOut', 'purged', { count: rows.length }); return rows.length; } module.exports = { stripEigenes, detectOptOut, isSuppressed, filterSuppressed, suppress, handleReply, scheduleRetention, purgeExpired, PURGE_MONTHS, }; ===== DATEI: src/services/unsubToken.js ===== // Abmelde-Token fuer die Kaltakquise. // // Der Link in der Mail muss ohne Anmeldung funktionieren und trotzdem darf // niemand fremde Adressen abmelden. Deshalb ein HMAC ueber Datensatz-ID UND // E-Mail: aendert sich die Adresse, verfaellt der alte Link automatisch. const crypto = require('crypto'); const SECRET = process.env.UNSUB_SECRET || ''; if (!SECRET) { // Ohne Geheimnis waeren alle Token faelschbar — dann lieber gar keine bauen. console.error('UNSUB_SECRET fehlt — Abmeldelinks sind deaktiviert.'); } function sign(id, email) { return crypto.createHmac('sha256', SECRET) .update(`${id}:${String(email || '').trim().toLowerCase()}`) .digest('base64url').slice(0, 22); } function make(id, email) { if (!SECRET) return null; return `${id}.${sign(id, email)}`; } // Gibt die Datensatz-ID zurueck oder null. Der Vergleich laeuft zeitkonstant, // damit sich das Geheimnis nicht ueber Laufzeitunterschiede abtasten laesst. function parse(token) { if (!SECRET || typeof token !== 'string') return null; const m = token.match(/^(\d{1,12})\.([A-Za-z0-9_-]{22})$/); return m ? { id: parseInt(m[1], 10), sig: m[2] } : null; } function verify(token, email) { const p = parse(token); if (!p) return null; const erwartet = Buffer.from(sign(p.id, email)); const bekommen = Buffer.from(p.sig); if (erwartet.length !== bekommen.length) return null; return crypto.timingSafeEqual(erwartet, bekommen) ? p.id : null; } module.exports = { make, parse, verify };