"""
outreach_orchestrator.py -- Steuert den gesamten Outbound-Workflow.
Aufgaben:
- Genehmigte Emails versenden (via SendGrid)
- Bounce-Handling
- Faellige Step-2/3 Emails triggern
- Kampagnen-Statistiken anzeigen
- Anti-Spam: Max 50 Emails/Tag, Zeitfenster, Delay zwischen Emails
Verwendung:
python outreach_orchestrator.py --send-approved
python outreach_orchestrator.py --check-due
python outreach_orchestrator.py --stats
python outreach_orchestrator.py --stats --campaign 1
"""
import os
import time
import logging
import argparse
from datetime import datetime, timedelta
from typing import Dict, List, Optional
from cold_outreach_db import ColdOutreachDB
from sendgrid_integration import SendGridClient
from outreach_generator import OutreachGenerator
logging.basicConfig(level=logging.INFO, format='%(asctime)s %(levelname)s %(message)s')
logger = logging.getLogger(__name__)
# ---------------------------------------------------------------------------
# Anti-Spam Konfiguration
# ---------------------------------------------------------------------------
MAX_EMAILS_PER_DAY = 50
DELAY_BETWEEN_SENDS = 30 # Sekunden zwischen Emails
SEND_WINDOW_START_HOUR = 9 # 09:00 Uhr
SEND_WINDOW_END_HOUR = 17 # 17:00 Uhr
SEND_DAYS = {1, 2, 3, 4, 5} # Mo-Fr (1=Mo, 7=So in Python isoweekday)
# Absender-Konfiguration
FROM_EMAIL_AG = os.getenv('OUTREACH_FROM_EMAIL_AG', 'karlo@agentsolutions.tech')
FROM_EMAIL_LP = os.getenv('OUTREACH_FROM_EMAIL_LP', 'karlo@mallorcaairservices.com')
FROM_NAME = 'Karlo'
# ---------------------------------------------------------------------------
# Orchestrator
# ---------------------------------------------------------------------------
class OutreachOrchestrator:
"""Steuert den gesamten Outbound-Email-Workflow."""
def __init__(self, db_path: str = 'email_agent.db',
sendgrid_api_key: str = None):
self.db = ColdOutreachDB(db_path)
self.sendgrid = SendGridClient(api_key=sendgrid_api_key)
self.generator = OutreachGenerator(db_path=db_path)
# ── Hauptfunktionen ──────────────────────────────────────────────────────
def send_approved(self, dry_run: bool = False) -> Dict:
"""
Sendet alle genehmigten Emails.
Respektiert Anti-Spam-Einstellungen.
Args:
dry_run: Wenn True, werden Emails nicht wirklich gesendet.
Returns:
Dict: {sent, failed, skipped_window, skipped_limit}
"""
approved = self.db.get_approved_emails()
stats = {'sent': 0, 'failed': 0, 'skipped_window': 0, 'skipped_limit': 0}
if not approved:
logger.info("Keine genehmigten Emails zu senden.")
return stats
logger.info(f"{len(approved)} genehmigte Emails gefunden.")
# Tageslimit pruefen
sent_today = self._count_sent_today()
remaining = MAX_EMAILS_PER_DAY - sent_today
if remaining <= 0:
logger.warning(f"Tageslimit erreicht ({MAX_EMAILS_PER_DAY}/Tag). Abbruch.")
stats['skipped_limit'] = len(approved)
return stats
# Sendefenster pruefen
if not self._in_send_window():
logger.warning(
f"Ausserhalb des Sendefensters "
f"({SEND_WINDOW_START_HOUR}:00-{SEND_WINDOW_END_HOUR}:00, Mo-Fr)."
)
stats['skipped_window'] = len(approved)
return stats
to_send = approved[:remaining]
logger.info(f"Sende {len(to_send)} Emails (Tageslimit: {MAX_EMAILS_PER_DAY}, verbleibend: {remaining}).")
for i, seq in enumerate(to_send, start=1):
try:
from_email = self._get_from_email(seq.get('icp', 'AG'))
if dry_run:
logger.info(
f"[DRY RUN] Wuerde senden: {seq['email']} | "
f"{seq.get('subject', '')} | ICP: {seq.get('icp', '?')}"
)
stats['sent'] += 1
else:
success = self.sendgrid.send_email(
to_email=seq['email'],
to_name=f"{seq.get('first_name', '')} {seq.get('last_name', '')}".strip(),
subject=seq.get('subject', ''),
body=seq.get('body', ''),
from_email=from_email,
from_name=FROM_NAME,
)
if success:
self.db.mark_sent(seq['id'])
self.db.set_contact_status(seq['contact_id'], 'active')
stats['sent'] += 1
logger.info(
f"[{i}/{len(to_send)}] Gesendet: {seq['email']} "
f"(Step {seq.get('sequence_step', '?')})"
)
else:
self.db.mark_failed(seq['id'], 'SendGrid returned error')
stats['failed'] += 1
logger.error(f"Fehlgeschlagen: {seq['email']}")
# Delay zwischen Emails (nicht nach letzter Email)
if i < len(to_send):
time.sleep(DELAY_BETWEEN_SENDS)
except Exception as e:
self.db.mark_failed(seq['id'], str(e))
stats['failed'] += 1
logger.error(f"Exception bei {seq.get('email', '?')}: {e}")
# Kampagnen-Stats aktualisieren
self._refresh_all_campaign_stats()
logger.info(f"Versand abgeschlossen: {stats}")
return stats
def check_due_and_generate(self, step: int = None) -> Dict:
"""
Prueft welche Step-2/3 Emails faellig sind und generiert sie.
Wird typischerweise taeglich ausgefuehrt.
Returns:
Dict mit Generierungs-Statistiken
"""
# Faellige Emails aus DB (scheduled_for <= jetzt, status=pending, noch kein body)
due = self.db.get_due_emails()
if step:
due = [d for d in due if d['sequence_step'] == step]
# Nur Faellige ohne generierten Body
need_generation = [d for d in due if not d.get('body')]
if not need_generation:
logger.info("Keine faelligen Emails zur Generierung.")
return {'generated': 0}
logger.info(f"{len(need_generation)} faellige Emails zur Generierung.")
# Nach Step gruppieren und generieren
stats = {}
for s in (1, 2, 3):
step_emails = [d for d in need_generation if d['sequence_step'] == s]
if step_emails:
logger.info(f"Generiere {len(step_emails)} Emails fuer Step {s}...")
step_stats = self.generator.generate_step(step=s)
stats[f'step_{s}'] = step_stats
return stats
def get_stats(self, campaign_id: int = None) -> Dict:
"""
Gibt Statistiken zurueck.
Wenn campaign_id gegeben: nur diese Kampagne.
Sonst: alle Kampagnen + Gesamt.
"""
if campaign_id:
campaign = self.db.get_campaign(campaign_id)
if not campaign:
return {'error': f'Kampagne {campaign_id} nicht gefunden'}
self.db.update_campaign_stats(campaign_id)
campaign = self.db.get_campaign(campaign_id)
outreach_stats = self.db.get_stats()
return {
'campaign': campaign,
'outreach': outreach_stats,
}
campaigns = self.db.list_campaigns()
for c in campaigns:
self.db.update_campaign_stats(c['id'])
campaigns = self.db.list_campaigns() # Aktualisierte Werte
outreach_stats = self.db.get_stats()
return {
'campaigns': campaigns,
'outreach': outreach_stats,
}
# ── Hilfsmethoden ────────────────────────────────────────────────────────
def _get_from_email(self, icp: str) -> str:
"""Gibt die passende Absender-Email je nach ICP zurueck."""
if icp == 'LP':
return FROM_EMAIL_LP
return FROM_EMAIL_AG
def _in_send_window(self) -> bool:
"""Prueft ob wir uns im erlaubten Sendefenster befinden."""
now = datetime.now()
if now.isoweekday() not in SEND_DAYS:
return False
if now.hour < SEND_WINDOW_START_HOUR or now.hour >= SEND_WINDOW_END_HOUR:
return False
return True
def _count_sent_today(self) -> int:
"""Zaehlt die bereits heute gesendeten Emails."""
# Vereinfacht: Statistiken aus der DB
# In Produktion wuerde man die cold_sequences nach sent_at filtern
stats = self.db.get_stats()
# Grobe Schaetzung -- fuer praezise Zaehlung wuerde man sent_at pruefen
return 0 # TODO: Praezise Zaehlung via sent_at Feld
def _refresh_all_campaign_stats(self):
"""Aktualisiert Statistiken aller Kampagnen."""
for campaign in self.db.list_campaigns():
self.db.update_campaign_stats(campaign['id'])
# ---------------------------------------------------------------------------
# Anzeige-Hilfsfunktionen
# ---------------------------------------------------------------------------
def print_stats(stats: Dict):
"""Gibt Statistiken formatiert aus."""
sep = '=' * 60
if 'campaigns' in stats:
print(f"\n{sep}")
print(" KAMPAGNEN-UEBERSICHT")
print(sep)
campaigns = stats['campaigns']
if not campaigns:
print(" Keine Kampagnen gefunden.")
else:
for c in campaigns:
total = c.get('contacts_total', 0)
replied = c.get('contacts_replied', 0)
reply_rate = (replied / total * 100) if total > 0 else 0
converted = c.get('contacts_converted', 0)
print(f"\n [{c['id']}] {c['name']}")
print(f" ICP: {c['icp']} | Quelle: {c.get('source_file', '?')}")
print(f" Kontakte: {total} | "
f"Angereichert: {c.get('contacts_enriched', 0)} | "
f"Gesendet: {c.get('contacts_sent', 0)}")
print(f" Antworten: {replied} ({reply_rate:.1f}%) | "
f"Konvertiert: {converted}")
print(f" Erstellt: {c.get('created_at', '?')[:10]}")
elif 'campaign' in stats:
c = stats['campaign']
total = c.get('contacts_total', 0)
replied = c.get('contacts_replied', 0)
reply_rate = (replied / total * 100) if total > 0 else 0
print(f"\n{sep}")
print(f" KAMPAGNE: {c['name']}")
print(sep)
print(f" Kontakte gesamt: {total}")
print(f" Angereichert: {c.get('contacts_enriched', 0)}")
print(f" Emails gesendet: {c.get('contacts_sent', 0)}")
print(f" Antworten: {replied} ({reply_rate:.1f}%)")
print(f" Konvertiert: {c.get('contacts_converted', 0)}")
outreach = stats.get('outreach', {})
if outreach:
print(f"\n{sep}")
print(" GESAMT-STATISTIK")
print('-' * 60)
print(f" Alle Kontakte: {outreach.get('contacts_total', 0)}")
print(f" Aktiv (Sequenz): {outreach.get('contacts_active', 0)}")
print(f" Haben geantwortet: {outreach.get('contacts_replied', 0)}")
print(f" Konvertiert: {outreach.get('contacts_converted', 0)}")
print(f" Reply-Rate: {outreach.get('reply_rate', 0)}%")
print(f" Emails gesendet: {outreach.get('emails_sent', 0)}")
print(f" Emails ausstehend: {outreach.get('emails_pending', 0)}")
print(sep)
# ---------------------------------------------------------------------------
# CLI
# ---------------------------------------------------------------------------
def main():
parser = argparse.ArgumentParser(
description='Outreach Orchestrator -- Versand und Workflow-Steuerung.'
)
group = parser.add_mutually_exclusive_group(required=True)
group.add_argument('--send-approved', action='store_true',
help='Alle genehmigten Emails versenden')
group.add_argument('--check-due', action='store_true',
help='Faellige Step-2/3 Emails pruefen und generieren')
group.add_argument('--stats', action='store_true',
help='Statistiken anzeigen')
parser.add_argument('--campaign', type=int, metavar='ID',
help='Nur fuer diese Kampagne')
parser.add_argument('--step', type=int, choices=[1, 2, 3],
help='Nur fuer diesen Step')
parser.add_argument('--dry-run', action='store_true',
help='Emails nicht wirklich senden (Test)')
parser.add_argument('--db', default='email_agent.db', help='Pfad zur DB-Datei')
args = parser.parse_args()
orchestrator = OutreachOrchestrator(db_path=args.db)
if args.send_approved:
dry_run = args.dry_run
if dry_run:
print("\n[DRY RUN] Emails werden NICHT wirklich gesendet.")
print("\nSende genehmigte Emails...")
stats = orchestrator.send_approved(dry_run=dry_run)
print(f"\nErgebnis:")
print(f" Gesendet: {stats['sent']}")
print(f" Fehlgeschlagen: {stats['failed']}")
if stats.get('skipped_window'):
print(f" Ausserhalb Zeitfenster: {stats['skipped_window']}")
if stats.get('skipped_limit'):
print(f" Tageslimit erreicht: {stats['skipped_limit']}")
if stats['sent'] > 0:
print(f"\nAntworten werden automatisch erkannt (agent.py laeuft).")
print(f"Naechste Steps faellig in 4 Tagen -- dann:")
print(f" python outreach_orchestrator.py --check-due")
elif args.check_due:
print("\nPruefe faellige Follow-Up Emails...")
stats = orchestrator.check_due_and_generate(step=args.step)
if stats.get('generated', 0) == 0 and not any(
v for k, v in stats.items() if k.startswith('step_')
):
print("Keine faelligen Emails.")
else:
for step_key, step_stats in stats.items():
if step_key.startswith('step_'):
step = step_key.split('_')[1]
print(f"\nStep {step}: {step_stats}")
print(f"\nNaechster Schritt:")
print(f" python review_cli.py")
elif args.stats:
stats = orchestrator.get_stats(campaign_id=args.campaign)
print_stats(stats)
if __name__ == '__main__':
main()