Explorer
/tmp/restic-stage/lead-engine/cold_outreach_db.py
← Zurück ↓ Download
"""
Cold Outreach Database
Verwaltet Kontakte und Email-Sequenzen für ausgehende Kaltakquise.
Ergänzt die bestehende EmailDatabase um Outbound-Tabellen.
"""

import sqlite3
import csv
from datetime import datetime, timedelta
from typing import List, Optional, Dict, Any
from contextlib import contextmanager


class ColdOutreachDB:
    """Datenbank für Kaltakquise-Kontakte und Sequenzen"""

    SEQUENCE_DELAYS = {1: 0, 2: 4, 3: 8}  # Tage nach Erstkontakt

    def __init__(self, db_path: str = 'email_agent.db'):
        self.db_path = db_path
        self._init_schema()
        self._migrate_schema()

    @contextmanager
    def get_connection(self):
        conn = sqlite3.connect(self.db_path)
        conn.row_factory = sqlite3.Row
        try:
            yield conn
        finally:
            conn.close()

    def _init_schema(self):
        with self.get_connection() as conn:
            # Kontakte-Tabelle
            conn.execute("""
                CREATE TABLE IF NOT EXISTS cold_contacts (
                    id INTEGER PRIMARY KEY AUTOINCREMENT,
                    first_name TEXT NOT NULL,
                    last_name TEXT,
                    email TEXT NOT NULL UNIQUE,
                    company TEXT,
                    role TEXT,
                    industry TEXT,
                    pain_point TEXT,
                    our_offer TEXT,
                    social_proof TEXT,
                    icp TEXT DEFAULT 'AG',          -- AG / LP / beide
                    status TEXT DEFAULT 'active',    -- active / paused / unsubscribed / replied / converted
                    notes TEXT,
                    created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
                    updated_at DATETIME DEFAULT CURRENT_TIMESTAMP
                )
            """)

            # Sequenz-Tabelle — eine Zeile pro geplanter Email
            conn.execute("""
                CREATE TABLE IF NOT EXISTS cold_sequences (
                    id INTEGER PRIMARY KEY AUTOINCREMENT,
                    contact_id INTEGER NOT NULL REFERENCES cold_contacts(id),
                    sequence_step INTEGER NOT NULL,  -- 1, 2, 3
                    scheduled_for DATETIME NOT NULL,
                    sent_at DATETIME,
                    status TEXT DEFAULT 'pending',   -- pending / sent / failed / skipped / replied
                    subject TEXT,
                    body TEXT,
                    error_message TEXT,
                    created_at DATETIME DEFAULT CURRENT_TIMESTAMP
                )
            """)

            # Kontakt-Varianten (Produkte, Abteilungen, etc.)
            conn.execute("""
                CREATE TABLE IF NOT EXISTS contact_variants (
                    id INTEGER PRIMARY KEY AUTOINCREMENT,
                    contact_id INTEGER NOT NULL REFERENCES cold_contacts(id),
                    offer_title TEXT,
                    industry_segment TEXT,
                    buyer_relevance TEXT,
                    outreach_angle TEXT,
                    priority_tier TEXT,
                    contact_page TEXT,
                    is_primary BOOLEAN DEFAULT 0,
                    created_at DATETIME DEFAULT CURRENT_TIMESTAMP
                )
            """)

            # Neue Tabellen
            conn.execute("""
                CREATE TABLE IF NOT EXISTS outreach_campaigns (
                    id INTEGER PRIMARY KEY AUTOINCREMENT,
                    name TEXT NOT NULL,
                    icp TEXT NOT NULL,
                    source_file TEXT,
                    created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
                    contacts_total INTEGER DEFAULT 0,
                    contacts_enriched INTEGER DEFAULT 0,
                    contacts_sent INTEGER DEFAULT 0,
                    contacts_replied INTEGER DEFAULT 0,
                    contacts_converted INTEGER DEFAULT 0
                )
            """)
            conn.execute("""
                CREATE TABLE IF NOT EXISTS enrichment_cache (
                    domain TEXT PRIMARY KEY,
                    scraped_text TEXT,
                    meta_title TEXT,
                    meta_description TEXT,
                    scraped_at DATETIME DEFAULT CURRENT_TIMESTAMP,
                    status TEXT DEFAULT 'success'
                )
            """)

            conn.execute("CREATE INDEX IF NOT EXISTS idx_cold_sequences_due ON cold_sequences(scheduled_for, status)")
            conn.execute("CREATE INDEX IF NOT EXISTS idx_cold_contacts_status ON cold_contacts(status)")
            conn.commit()

    def _migrate_schema(self):
        """Fuegt neue Spalten hinzu falls sie fehlen. Idempotent."""
        contact_migrations = {
            'website':                  'TEXT',
            'phone':                    'TEXT',
            'company_description':      'TEXT',
            'enrichment_source':        'TEXT',
            'enrichment_date':          'DATETIME',
            'icp_confidence':           'REAL DEFAULT 0.0',
            'icp_recommended_product':  'TEXT',
            'language':                 "TEXT DEFAULT 'de'",
            'campaign_id':              'INTEGER',
            'source_id':                'TEXT',
            # Erweiterte Enrichment-Felder
            'company_size':             'TEXT',          # z.B. "50-100 Mitarbeiter"
            'headquarters':             'TEXT',          # Sitz/Hauptquartier
            'founding_year':            'INTEGER',       # Gründungsjahr
            'technology_stack':         'TEXT',          # z.B. "Node.js, React, PostgreSQL"
            'key_products':             'TEXT',          # Hauptprodukte/Services
            'linkedin_url':             'TEXT',          # LinkedIn Company Page
        }
        seq_migrations = {
            'personalization_score':    'REAL',
            'personalization_elements': 'TEXT',
            'review_status':            "TEXT DEFAULT 'pending'",
            'reviewed_at':              'DATETIME',
            'reviewed_by':              "TEXT DEFAULT 'karlo'",
            'knowledge_snippet':        'TEXT',
        }
        with self.get_connection() as conn:
            existing_contacts = {row[1] for row in conn.execute("PRAGMA table_info(cold_contacts)")}
            for col, col_type in contact_migrations.items():
                if col not in existing_contacts:
                    conn.execute(f"ALTER TABLE cold_contacts ADD COLUMN {col} {col_type}")

            existing_seqs = {row[1] for row in conn.execute("PRAGMA table_info(cold_sequences)")}
            for col, col_type in seq_migrations.items():
                if col not in existing_seqs:
                    conn.execute(f"ALTER TABLE cold_sequences ADD COLUMN {col} {col_type}")
            conn.commit()

    # ── Kontakte ─────────────────────────────────────────────────────────────

    def add_contact(self, data: Dict[str, Any], create_sequence: bool = False) -> int:
        """
        Einzelnen Kontakt hinzufügen. Gibt ID zurück.
        create_sequence=False: Sequenz wird erst nach Enrichment/Freigabe angelegt.
        """
        with self.get_connection() as conn:
            cursor = conn.execute("""
                INSERT INTO cold_contacts
                    (first_name, last_name, email, company, role, industry,
                     pain_point, our_offer, social_proof, icp, notes,
                     website, phone, company_description, language,
                     campaign_id, status, source_id)
                VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)
            """, (
                data.get('first_name', ''),
                data.get('last_name', ''),
                data.get('email', ''),
                data.get('company', ''),
                data.get('role', ''),
                data.get('industry', ''),
                data.get('pain_point', ''),
                data.get('our_offer', ''),
                data.get('social_proof', ''),
                data.get('icp', 'AG'),
                data.get('notes', ''),
                data.get('website', ''),
                data.get('phone', ''),
                data.get('company_description', ''),
                data.get('language', 'de'),
                data.get('campaign_id'),
                data.get('status', 'raw'),
                data.get('source_id', ''),
            ))
            conn.commit()
            contact_id = cursor.lastrowid

        if create_sequence:
            self._create_sequence(contact_id)
        return contact_id

    def update_contact(self, contact_id: int, data: Dict[str, Any]):
        """Kontaktfelder aktualisieren."""
        allowed = {
            'industry', 'pain_point', 'company_description', 'website', 'phone',
            'icp', 'icp_confidence', 'icp_recommended_product', 'language',
            'enrichment_source', 'enrichment_date', 'status', 'notes', 'our_offer',
            'company_size', 'headquarters', 'founding_year', 'technology_stack',
            'key_products', 'linkedin_url',
        }
        updates = {k: v for k, v in data.items() if k in allowed}
        if not updates:
            return
        updates['updated_at'] = datetime.now().isoformat()
        cols = ', '.join(f"{k} = ?" for k in updates)
        with self.get_connection() as conn:
            conn.execute(
                f"UPDATE cold_contacts SET {cols} WHERE id = ?",
                list(updates.values()) + [contact_id]
            )
            conn.commit()

    def get_contacts_by_status(self, status: str, campaign_id: int = None) -> List[Dict]:
        """Kontakte nach Status filtern, optional nach Kampagne."""
        query = "SELECT * FROM cold_contacts WHERE status = ?"
        params = [status]
        if campaign_id:
            query += " AND campaign_id = ?"
            params.append(campaign_id)
        with self.get_connection() as conn:
            return [dict(r) for r in conn.execute(query, params).fetchall()]

    def email_exists(self, email: str) -> bool:
        """Prüft ob Email-Adresse bereits in der DB vorhanden ist."""
        with self.get_connection() as conn:
            row = conn.execute(
                "SELECT 1 FROM cold_contacts WHERE email = ?", (email.strip().lower(),)
            ).fetchone()
            return row is not None

    # ── Kampagnen ────────────────────────────────────────────────────────────

    def create_campaign(self, name: str, icp: str, source_file: str = None) -> int:
        """Neue Kampagne anlegen. Gibt ID zurück."""
        with self.get_connection() as conn:
            cursor = conn.execute("""
                INSERT INTO outreach_campaigns (name, icp, source_file)
                VALUES (?,?,?)
            """, (name, icp, source_file))
            conn.commit()
            return cursor.lastrowid

    def update_campaign_stats(self, campaign_id: int):
        """Statistiken einer Kampagne aus DB-Daten neu berechnen."""
        with self.get_connection() as conn:
            stats = conn.execute("""
                SELECT
                    COUNT(*) as total,
                    SUM(CASE WHEN status IN ('enriched','active','replied','converted') THEN 1 ELSE 0 END) as enriched,
                    SUM(CASE WHEN status IN ('active','replied','converted') THEN 1 ELSE 0 END) as sent,
                    SUM(CASE WHEN status = 'replied' THEN 1 ELSE 0 END) as replied,
                    SUM(CASE WHEN status = 'converted' THEN 1 ELSE 0 END) as converted
                FROM cold_contacts WHERE campaign_id = ?
            """, (campaign_id,)).fetchone()
            conn.execute("""
                UPDATE outreach_campaigns
                SET contacts_total=?, contacts_enriched=?, contacts_sent=?,
                    contacts_replied=?, contacts_converted=?
                WHERE id=?
            """, (stats[0], stats[1], stats[2], stats[3], stats[4], campaign_id))
            conn.commit()

    def get_campaign(self, campaign_id: int) -> Optional[Dict]:
        with self.get_connection() as conn:
            row = conn.execute(
                "SELECT * FROM outreach_campaigns WHERE id = ?", (campaign_id,)
            ).fetchone()
            return dict(row) if row else None

    def list_campaigns(self) -> List[Dict]:
        with self.get_connection() as conn:
            return [dict(r) for r in conn.execute(
                "SELECT * FROM outreach_campaigns ORDER BY created_at DESC"
            ).fetchall()]

    # ── Enrichment Cache ─────────────────────────────────────────────────────

    def cache_enrichment(self, domain: str, scraped_text: str,
                         meta_title: str = '', meta_description: str = '',
                         status: str = 'success'):
        """Website-Scraping-Ergebnis cachen."""
        with self.get_connection() as conn:
            conn.execute("""
                INSERT OR REPLACE INTO enrichment_cache
                    (domain, scraped_text, meta_title, meta_description, scraped_at, status)
                VALUES (?,?,?,?,?,?)
            """, (domain, scraped_text, meta_title, meta_description,
                  datetime.now().isoformat(), status))
            conn.commit()

    def get_cached_enrichment(self, domain: str) -> Optional[Dict]:
        """Gecachtes Enrichment-Ergebnis für eine Domain holen."""
        with self.get_connection() as conn:
            row = conn.execute(
                "SELECT * FROM enrichment_cache WHERE domain = ?", (domain,)
            ).fetchone()
            return dict(row) if row else None

    # ── Sequenzen (Review) ───────────────────────────────────────────────────

    def save_generated_email(self, seq_id: int, subject: str, body: str,
                             personalization_score: float = None,
                             personalization_elements: list = None,
                             knowledge_snippet: str = None):
        """Generierten Emailinhalt speichern (erweitert)."""
        with self.get_connection() as conn:
            conn.execute("""
                UPDATE cold_sequences
                SET subject = ?, body = ?,
                    personalization_score = ?,
                    personalization_elements = ?,
                    knowledge_snippet = ?,
                    review_status = 'pending'
                WHERE id = ?
            """, (
                subject, body,
                personalization_score,
                str(personalization_elements) if personalization_elements else None,
                knowledge_snippet,
                seq_id,
            ))
            conn.commit()

    def approve_email(self, seq_id: int, edited_body: str = None, edited_subject: str = None):
        """Email für Versand freigeben."""
        with self.get_connection() as conn:
            updates = {
                'review_status': 'approved',
                'reviewed_at': datetime.now().isoformat(),
            }
            if edited_body:
                updates['body'] = edited_body
            if edited_subject:
                updates['subject'] = edited_subject
            cols = ', '.join(f"{k} = ?" for k in updates)
            conn.execute(
                f"UPDATE cold_sequences SET {cols} WHERE id = ?",
                list(updates.values()) + [seq_id]
            )
            conn.commit()

    def reject_email(self, seq_id: int):
        """Email ablehnen (wird nicht gesendet)."""
        with self.get_connection() as conn:
            conn.execute(
                "UPDATE cold_sequences SET review_status = 'rejected', reviewed_at = ? WHERE id = ?",
                (datetime.now().isoformat(), seq_id)
            )
            conn.commit()

    def get_pending_review(self, campaign_id: int = None) -> List[Dict]:
        """Alle generierten Emails die noch auf Review warten."""
        query = """
            SELECT s.*, c.first_name, c.last_name, c.email, c.company,
                   c.role, c.industry, c.pain_point, c.icp,
                   c.icp_recommended_product, c.language,
                   c.status AS contact_status
            FROM cold_sequences s
            JOIN cold_contacts c ON s.contact_id = c.id
            WHERE s.review_status = 'pending'
              AND s.body IS NOT NULL
              AND c.status NOT IN ('unsubscribed', 'bounced')
        """
        params = []
        if campaign_id:
            query += " AND c.campaign_id = ?"
            params.append(campaign_id)
        query += " ORDER BY s.sequence_step ASC, s.scheduled_for ASC"
        with self.get_connection() as conn:
            return [dict(r) for r in conn.execute(query, params).fetchall()]

    def get_approved_emails(self) -> List[Dict]:
        """Alle freigegebenen Emails die noch nicht gesendet wurden."""
        with self.get_connection() as conn:
            rows = conn.execute("""
                SELECT s.*, c.first_name, c.last_name, c.email, c.company,
                       c.role, c.industry, c.icp, c.language,
                       c.status AS contact_status
                FROM cold_sequences s
                JOIN cold_contacts c ON s.contact_id = c.id
                WHERE s.review_status = 'approved'
                  AND s.status = 'pending'
                  AND c.status NOT IN ('unsubscribed', 'bounced', 'replied')
                ORDER BY s.scheduled_for ASC
            """).fetchall()
            return [dict(r) for r in rows]

    def import_from_csv(self, csv_path: str) -> Dict[str, int]:
        """
        Kontakte aus CSV importieren.

        Erwartete Spalten (Mindestzahl):
            first_name, email
        Optional:
            last_name, company, role, industry, pain_point,
            our_offer, social_proof, icp, notes
        """
        imported, skipped = 0, 0
        with open(csv_path, newline='', encoding='utf-8') as f:
            reader = csv.DictReader(f)
            for row in reader:
                try:
                    self.add_contact(row)
                    imported += 1
                except sqlite3.IntegrityError:
                    skipped += 1  # Email bereits vorhanden
        return {'imported': imported, 'skipped': skipped}

    def get_contact(self, contact_id: int) -> Optional[Dict]:
        with self.get_connection() as conn:
            row = conn.execute("SELECT * FROM cold_contacts WHERE id = ?", (contact_id,)).fetchone()
            return dict(row) if row else None

    def set_contact_status(self, contact_id: int, status: str):
        """Status setzen: active / paused / unsubscribed / replied / converted"""
        with self.get_connection() as conn:
            conn.execute(
                "UPDATE cold_contacts SET status = ?, updated_at = ? WHERE id = ?",
                (status, datetime.now().isoformat(), contact_id)
            )
            conn.commit()

    def get_contacts(self, status: str = 'active', icp: str = None) -> List[Dict]:
        query = "SELECT * FROM cold_contacts WHERE status = ?"
        params = [status]
        if icp:
            query += " AND icp = ?"
            params.append(icp)
        with self.get_connection() as conn:
            return [dict(r) for r in conn.execute(query, params).fetchall()]

    # ── Sequenzen ────────────────────────────────────────────────────────────

    def add_variant(self, contact_id: int, variant_data: Dict[str, Any], is_primary: bool = False):
        """Fügt eine Produkt-Variante zu einem Kontakt hinzu."""
        with self.get_connection() as conn:
            conn.execute("""
                INSERT INTO contact_variants
                    (contact_id, offer_title, industry_segment, buyer_relevance,
                     outreach_angle, priority_tier, contact_page, is_primary)
                VALUES (?, ?, ?, ?, ?, ?, ?, ?)
            """, (
                contact_id,
                variant_data.get('offer_title'),
                variant_data.get('industry_segment'),
                variant_data.get('buyer_relevance'),
                variant_data.get('outreach_angle'),
                variant_data.get('priority_tier'),
                variant_data.get('contact_page'),
                1 if is_primary else 0,
            ))
            conn.commit()

    def get_variants(self, contact_id: int) -> List[Dict]:
        """Gibt alle Varianten eines Kontakts zurück."""
        with self.get_connection() as conn:
            rows = conn.execute(
                "SELECT * FROM contact_variants WHERE contact_id = ? ORDER BY is_primary DESC, id ASC",
                (contact_id,)
            ).fetchall()
            return [dict(r) for r in rows]

    def activate_enriched_contacts(self, campaign_id: int = None) -> int:
        """
        Aktiviert alle 'enriched' Kontakte: erstellt Sequenzen und setzt Status auf 'active'.
        Wird vor der Email-Generierung aufgerufen.

        Returns:
            Anzahl aktivierter Kontakte
        """
        contacts = self.get_contacts_by_status('enriched', campaign_id=campaign_id)
        activated = 0
        for c in contacts:
            contact_id = c['id']
            # Prüfen ob bereits Sequenzen existieren
            with self.get_connection() as conn:
                existing = conn.execute(
                    "SELECT COUNT(*) FROM cold_sequences WHERE contact_id = ?",
                    (contact_id,)
                ).fetchone()[0]
            if existing == 0:
                self._create_sequence(contact_id)
            self.update_contact(contact_id, {'status': 'active'})
            activated += 1
        return activated

    def _create_sequence(self, contact_id: int):
        """Legt 3 Sequenz-Einträge (pending) für einen Kontakt an."""
        now = datetime.now()
        with self.get_connection() as conn:
            for step, delay_days in self.SEQUENCE_DELAYS.items():
                scheduled = now + timedelta(days=delay_days)
                conn.execute("""
                    INSERT INTO cold_sequences (contact_id, sequence_step, scheduled_for)
                    VALUES (?,?,?)
                """, (contact_id, step, scheduled.isoformat()))
            conn.commit()

    def get_due_emails(self) -> List[Dict]:
        """Alle fälligen Emails (scheduled_for <= jetzt, status=pending)."""
        with self.get_connection() as conn:
            rows = conn.execute("""
                SELECT s.*, c.first_name, c.last_name, c.email, c.company,
                       c.role, c.industry, c.pain_point, c.our_offer,
                       c.social_proof, c.icp, c.notes, c.company_description,
                       c.icp_recommended_product, c.language,
                       c.status AS contact_status
                FROM cold_sequences s
                JOIN cold_contacts c ON s.contact_id = c.id
                WHERE s.status = 'pending'
                  AND s.scheduled_for <= ?
                  AND c.status = 'active'
                ORDER BY s.scheduled_for ASC
            """, (datetime.now().isoformat(),)).fetchall()
            return [dict(r) for r in rows]

    def save_generated_email(self, seq_id: int, subject: str, body: str,
                             personalization_score: float = None,
                             personalization_elements: list = None,
                             knowledge_snippet: str = None):
        """Generierten Emailinhalt speichern."""
        import json as _json
        elements_json = _json.dumps(personalization_elements) if personalization_elements else None
        with self.get_connection() as conn:
            conn.execute("""
                UPDATE cold_sequences
                SET subject = ?, body = ?,
                    personalization_score = ?,
                    personalization_elements = ?,
                    knowledge_snippet = ?
                WHERE id = ?
            """, (subject, body, personalization_score, elements_json, knowledge_snippet, seq_id))
            conn.commit()

    def mark_sent(self, seq_id: int):
        with self.get_connection() as conn:
            conn.execute(
                "UPDATE cold_sequences SET status = 'sent', sent_at = ? WHERE id = ?",
                (datetime.now().isoformat(), seq_id)
            )
            conn.commit()

    def mark_failed(self, seq_id: int, error: str):
        with self.get_connection() as conn:
            conn.execute(
                "UPDATE cold_sequences SET status = 'failed', error_message = ? WHERE id = ?",
                (error, seq_id)
            )
            conn.commit()

    def mark_replied(self, contact_email: str):
        """Wenn Kontakt antwortet: alle pending Sequenz-Steps überspringen."""
        with self.get_connection() as conn:
            conn.execute("""
                UPDATE cold_contacts SET status = 'replied', updated_at = ?
                WHERE email = ?
            """, (datetime.now().isoformat(), contact_email))
            conn.execute("""
                UPDATE cold_sequences SET status = 'skipped'
                WHERE contact_id = (SELECT id FROM cold_contacts WHERE email = ?)
                  AND status = 'pending'
            """, (contact_email,))
            conn.commit()

    # ── Statistik ────────────────────────────────────────────────────────────

    def get_stats(self) -> Dict[str, Any]:
        with self.get_connection() as conn:
            total     = conn.execute("SELECT COUNT(*) FROM cold_contacts").fetchone()[0]
            active    = conn.execute("SELECT COUNT(*) FROM cold_contacts WHERE status='active'").fetchone()[0]
            replied   = conn.execute("SELECT COUNT(*) FROM cold_contacts WHERE status='replied'").fetchone()[0]
            converted = conn.execute("SELECT COUNT(*) FROM cold_contacts WHERE status='converted'").fetchone()[0]
            sent      = conn.execute("SELECT COUNT(*) FROM cold_sequences WHERE status='sent'").fetchone()[0]
            pending   = conn.execute("SELECT COUNT(*) FROM cold_sequences WHERE status='pending'").fetchone()[0]
            return {
                'contacts_total': total,
                'contacts_active': active,
                'contacts_replied': replied,
                'contacts_converted': converted,
                'reply_rate': round(replied / total * 100, 1) if total else 0,
                'emails_sent': sent,
                'emails_pending': pending,
            }