"""
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,
}