Explorer
/tmp/restic-stage/lead-engine/outreach_orchestrator.py
← Zurück ↓ Download
"""
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()