""" progress_reporter.py - Reporting and alerting system. Generates: - Daily summaries (videos processed, API usage, pipeline performance) - Alert detection (API down, budget low, consecutive failures) - Pipeline statistics (avg time per stage, success rate) - Channel growth tracking All data sourced from StateManager. """ import json import os import threading from datetime import datetime, timedelta from typing import Any, Dict, List, Optional class ProgressReporter: """Reporting and alerting system for AutoDub HerStory.""" def __init__(self, state=None): self.state = state self._lock = threading.Lock() # ================================================================ # Daily Reports # ================================================================ def generate_daily_report(self) -> Dict[str, Any]: """Generate a comprehensive daily summary report.""" if not self.state or not hasattr(self.state, "_state"): return {"error": "No state available"} state = self.state._state today = datetime.now().strftime("%Y-%m-%d") # Videos stats processed = state.get("processed_videos", {}) failed = state.get("failed_videos", {}) stats = state.get("stats", {}) # Count today's activity today_processed = 0 today_failed = 0 for vid_id, info in processed.items(): date = info.get("date_processed", "") if date.startswith(today): today_processed += 1 for vid_id, info in failed.items(): date = info.get("last_failed_at", "") if date.startswith(today): today_failed += 1 # Queue stats task_queue = state.get("task_queue", []) queue_pending = sum(1 for t in task_queue if t.get("status") == "pending") queue_active = sum(1 for t in task_queue if t.get("status") == "claimed") queue_completed = sum(1 for t in task_queue if t.get("status") == "completed") queue_failed = sum(1 for t in task_queue if t.get("status") == "failed") # API budget budget = state.get("provider_budget", {}) groq_daily = budget.get("groq", {}).get("daily", {}).get(today, {}) or_daily = budget.get("openrouter", {}).get("daily", {}).get(today, {}) # Health health_history = state.get("health_history", []) last_health = health_history[-1] if health_history else {} # QA qa_history = state.get("qa_history", []) today_qa = [q for q in qa_history if q.get("checked_at", "").startswith(today)] qa_pass_rate = (sum(1 for q in today_qa if q.get("passed")) / len(today_qa) * 100) if today_qa else 100 report = { "date": today, "generated_at": datetime.now().isoformat(), "videos": { "processed_today": today_processed, "failed_today": today_failed, "total_processed": len(processed), "total_failed": len(failed), }, "queue": { "pending": queue_pending, "active": queue_active, "completed_today": queue_completed, "failed": queue_failed, }, "api_usage": { "openrouter_requests_today": sum( k.get("requests", 0) for k in or_daily.values() ), "groq_requests_today": groq_daily.get("requests", 0), }, "quality": { "qa_checks_today": len(today_qa), "qa_pass_rate": round(qa_pass_rate, 1), "avg_score": round(sum(q.get("score", 0) for q in today_qa) / len(today_qa), 1) if today_qa else None, }, "health": { "last_status": last_health.get("status", "unknown"), "last_check": last_health.get("timestamp", "never"), }, "stats": stats, } # Save report self._save_report(report) return report def generate_video_report(self, video_id: str) -> Dict[str, Any]: """Generate a report for a specific video.""" if not self.state or not hasattr(self.state, "_state"): return {"error": "No state available"} state = self.state._state report = {"video_id": video_id, "generated_at": datetime.now().isoformat()} # Check processed if video_id in state.get("processed_videos", {}): report["status"] = "published" report["details"] = state["processed_videos"][video_id] elif video_id in state.get("failed_videos", {}): report["status"] = "failed" report["details"] = state["failed_videos"][video_id] else: # Check task queue for task in state.get("task_queue", []): if task.get("video_id") == video_id: report["status"] = f"in_pipeline (stage: {task.get('stage', 'unknown')})" report["details"] = task break else: report["status"] = "not_found" # QA history for this video qa_history = state.get("qa_history", []) video_qa = [q for q in qa_history if video_id in q.get("video_path", "")] if video_qa: report["qa"] = video_qa[-1] # Latest QA check # Analytics analytics = state.get("video_analytics", {}).get(video_id) if analytics: report["analytics"] = analytics return report # ================================================================ # Alerts # ================================================================ def should_alert(self) -> Optional[Dict[str, Any]]: """Check if any alert condition is met. Returns alert dict or None.""" if not self.state or not hasattr(self.state, "_state"): return None state = self.state._state alerts = [] # CRITICAL: All API providers down health_history = state.get("health_history", []) if health_history: last_health = health_history[-1] if last_health.get("status") == "error": error_checks = [ c for c in last_health.get("summary", {}) ] alerts.append({ "level": "CRITICAL", "type": "api_down", "message": "System health check shows errors - API providers may be down", "details": last_health, }) # CRITICAL: YouTube auth expired tokens = state.get("youtube_tokens", {}) if not tokens or not (tokens.get("access_token") or tokens.get("token")): alerts.append({ "level": "CRITICAL", "type": "youtube_auth", "message": "YouTube not authenticated - cannot upload videos", }) # WARNING: Budget >80% used budget = state.get("provider_budget", {}) today = datetime.now().strftime("%Y-%m-%d") groq_daily = budget.get("groq", {}).get("daily", {}).get(today, {}) groq_used = groq_daily.get("requests", 0) groq_limit = 14400 if groq_used > groq_limit * 0.8: alerts.append({ "level": "WARNING", "type": "budget_groq", "message": f"Groq API {groq_used}/{groq_limit} requests used ({groq_used/groq_limit*100:.0f}%)", }) # OpenRouter budget or_daily = budget.get("openrouter", {}).get("daily", {}).get(today, {}) or_total = sum(k.get("requests", 0) for k in or_daily.values()) n_keys = len([k.strip() for k in os.getenv("OPENROUTER_API_KEYS", "").split(",") if k.strip()]) or_limit = n_keys * 45 if or_limit > 0 and or_total > or_limit * 0.8: alerts.append({ "level": "WARNING", "type": "budget_openrouter", "message": f"OpenRouter {or_total}/{or_limit} requests used ({or_total/or_limit*100:.0f}%)", }) # WARNING: 3+ consecutive failures failed = state.get("failed_videos", {}) recent_failures = 0 now = datetime.now() for vid_id, info in failed.items(): last_failed = info.get("last_failed_at", "") try: failed_dt = datetime.fromisoformat(last_failed) if (now - failed_dt).total_seconds() < 3600: # Last hour recent_failures += 1 except (ValueError, TypeError): pass if recent_failures >= 3: alerts.append({ "level": "WARNING", "type": "consecutive_failures", "message": f"{recent_failures} failures in the last hour", "details": {"recent_failures": recent_failures}, }) # INFO: Video published processed = state.get("processed_videos", {}) for vid_id, info in processed.items(): date = info.get("date_processed", "") try: if (now - datetime.fromisoformat(date)).total_seconds() < 600: # Last 10 minutes alerts.append({ "level": "INFO", "type": "video_published", "message": f"Video published: {info.get('es_title', vid_id)}", }) break # Only one except (ValueError, TypeError): pass # Return highest priority alert if not alerts: return None # Priority: CRITICAL > WARNING > INFO priority = {"CRITICAL": 0, "WARNING": 1, "INFO": 2} alerts.sort(key=lambda a: priority.get(a["level"], 99)) return alerts[0] # ================================================================ # Formatting # ================================================================ def format_report(self, report: Dict[str, Any], format: str = "text") -> str: """Format a report as text or markdown.""" if format == "markdown": return self._format_markdown(report) return self._format_text(report) def _format_text(self, report: Dict[str, Any]) -> str: """Format report as plain text.""" lines = [] lines.append(f"=== AutoDub Daily Report ({report.get('date', 'N/A')}) ===") lines.append("") videos = report.get("videos", {}) lines.append(f"Videos: {videos.get('processed_today', 0)} processed, {videos.get('failed_today', 0)} failed today") lines.append(f"Total: {videos.get('total_processed', 0)} processed, {videos.get('total_failed', 0)} failed all-time") queue = report.get("queue", {}) lines.append(f"Queue: {queue.get('pending', 0)} pending, {queue.get('active', 0)} active, {queue.get('failed', 0)} failed") api = report.get("api_usage", {}) lines.append(f"API: OpenRouter {api.get('openrouter_requests_today', 0)} req, Groq {api.get('groq_requests_today', 0)} req") quality = report.get("quality", {}) lines.append(f"QA: {quality.get('qa_checks_today', 0)} checks, {quality.get('qa_pass_rate', 0)}% pass rate") health = report.get("health", {}) lines.append(f"Health: {health.get('last_status', 'unknown')} (last: {health.get('last_check', 'never')})") return "\n".join(lines) def _format_markdown(self, report: Dict[str, Any]) -> str: """Format report as markdown.""" lines = [] lines.append(f"# AutoDub Daily Report - {report.get('date', 'N/A')}") lines.append("") videos = report.get("videos", {}) lines.append("## Videos") lines.append(f"- **Today**: {videos.get('processed_today', 0)} processed, {videos.get('failed_today', 0)} failed") lines.append(f"- **All-time**: {videos.get('total_processed', 0)} processed, {videos.get('total_failed', 0)} failed") lines.append("") queue = report.get("queue", {}) lines.append("## Queue") lines.append(f"- Pending: {queue.get('pending', 0)}") lines.append(f"- Active: {queue.get('active', 0)}") lines.append(f"- Failed: {queue.get('failed', 0)}") lines.append("") api = report.get("api_usage", {}) lines.append("## API Usage") lines.append(f"- OpenRouter: {api.get('openrouter_requests_today', 0)} requests") lines.append(f"- Groq: {api.get('groq_requests_today', 0)} requests") lines.append("") quality = report.get("quality", {}) lines.append("## Quality") lines.append(f"- Checks today: {quality.get('qa_checks_today', 0)}") lines.append(f"- Pass rate: {quality.get('qa_pass_rate', 0)}%") if quality.get("avg_score"): lines.append(f"- Avg score: {quality['avg_score']}/10") lines.append("") return "\n".join(lines) # ================================================================ # Pipeline Statistics # ================================================================ def get_pipeline_stats(self) -> Dict[str, Any]: """Get pipeline performance statistics.""" if not self.state or not hasattr(self.state, "_state"): return {"error": "No state available"} state = self.state._state processed = state.get("processed_videos", {}) failed = state.get("failed_videos", {}) task_queue = state.get("task_queue", []) # Success rate total_attempts = len(processed) + sum( info.get("retry_count", 1) for info in failed.values() ) success_rate = (len(processed) / total_attempts * 100) if total_attempts > 0 else 0 # Most common failure reasons failure_reasons = {} for vid_id, info in failed.items(): stage = info.get("stage_failed", "unknown") failure_reasons[stage] = failure_reasons.get(stage, 0) + 1 # Stage distribution of active tasks stage_distribution = {} for task in task_queue: stage = task.get("stage", "unknown") stage_distribution[stage] = stage_distribution.get(stage, 0) + 1 return { "total_processed": len(processed), "total_failed": len(failed), "success_rate": round(success_rate, 1), "failure_reasons": failure_reasons, "stage_distribution": stage_distribution, "active_tasks": sum(1 for t in task_queue if t.get("status") == "claimed"), "pending_tasks": sum(1 for t in task_queue if t.get("status") == "pending"), } def get_channel_growth(self) -> Dict[str, Any]: """Basic channel growth tracking.""" if not self.state or not hasattr(self.state, "_state"): return {"error": "No state available"} state = self.state._state processed = state.get("processed_videos", {}) # Group by date by_date = {} for vid_id, info in processed.items(): date = info.get("date_processed", "")[:10] # YYYY-MM-DD if date: by_date[date] = by_date.get(date, 0) + 1 # Sort by date dates = sorted(by_date.keys()) return { "total_videos": len(processed), "by_date": {d: by_date[d] for d in dates}, "first_video": dates[0] if dates else None, "last_video": dates[-1] if dates else None, "channel_analytics": state.get("channel_analytics", {}), } # ================================================================ # Internal Methods # ================================================================ def _save_report(self, report: Dict[str, Any]): """Save daily report to state history.""" if not self.state or not hasattr(self.state, "_state"): return try: reports = self.state._state.get("reports", []) reports.append({ "date": report.get("date"), "generated_at": report.get("generated_at"), "videos_processed_today": report.get("videos", {}).get("processed_today", 0), "videos_failed_today": report.get("videos", {}).get("failed_today", 0), }) # Keep last 30 daily reports if len(reports) > 30: reports = reports[-30:] self.state._state["reports"] = reports except Exception as e: print(f"[REPORTER] Failed to save report: {e}")