diff --git a/chatops/jobs_commands.py b/chatops/jobs_commands.py new file mode 100644 index 000000000..58416c656 --- /dev/null +++ b/chatops/jobs_commands.py @@ -0,0 +1,159 @@ +# chatops/jobs_commands.py +""" +פקודות ChatOps לניהול Background Jobs. +""" + +import os +from typing import Dict, List +from services.job_registry import JobRegistry, JobCategory +from services.job_tracker import get_job_tracker + + +def handle_jobs_command(args: str) -> str: + """ + /jobs [category|status|] + + דוגמאות: + - /jobs - רשימת כל ה-jobs + - /jobs backup - jobs בקטגוריית גיבויים + - /jobs active - הרצות פעילות + - /jobs failed - הרצות שנכשלו לאחרונה + - /jobs cache_warming - פרטי job ספציפי + """ + args = args.strip().lower() + registry = JobRegistry() + tracker = get_job_tracker() + + # URL בסיס למוניטור (ניתן לקנפג דרך ENV) + monitor_base_url = os.getenv("WEBAPP_URL", "http://localhost") + + # Active runs + if args == "active": + runs = tracker.get_active_runs() + if not runs: + return "✅ אין הרצות פעילות כרגע" + + lines = ["⚡ **הרצות פעילות:**\n"] + for run in runs: + status_icon = {"running": "🔄", "pending": "⏳"}.get(run.status.value, "❓") + # 🔗 קישור ישיר ללוגים של ההרצה + logs_link = f"{monitor_base_url}/jobs/monitor?run_id={run.run_id}" + lines.append( + f"{status_icon} `{run.job_id}` - {run.progress}% " + f"({run.processed_items}/{run.total_items})\n" + f" [📋 לוגים]({logs_link})" + ) + return "\n".join(lines) + + # Failed runs + if args == "failed": + runs = tracker.get_failed_runs(limit=10) + if not runs: + return "✅ אין הרצות שנכשלו לאחרונה" + + lines = ["❌ **הרצות שנכשלו:**\n"] + for run in runs: + time_str = run.ended_at.strftime('%d/%m %H:%M') if run.ended_at else "-" + error_short = (run.error_message[:50] + "...") if run.error_message and len(run.error_message) > 50 else (run.error_message or "") + logs_link = f"{monitor_base_url}/jobs/monitor?run_id={run.run_id}" + lines.append( + f"❌ `{run.job_id}` - {time_str}\n" + f" {error_short}\n" + f" [📋 ראה לוגים]({logs_link})" + ) + return "\n".join(lines) + + # By category + try: + category = JobCategory(args) + jobs = registry.list_by_category(category) + if not jobs: + return f"אין jobs בקטגוריה `{args}`" + + lines = [f"📋 **Jobs בקטגוריית {args}:**\n"] + for j in jobs: + status = "✅" if registry.is_enabled(j.job_id) else "❌" + lines.append(f"{status} `{j.job_id}` - {j.name}") + return "\n".join(lines) + except ValueError: + pass + + # Specific job + if args: + job = registry.get(args) + if not job: + return f"❌ Job `{args}` לא נמצא" + + history = tracker.get_job_history(args, limit=5) + status = "✅ פעיל" if registry.is_enabled(args) else "❌ מושבת" + + lines = [ + f"📋 **{job.name}**\n", + f"• מזהה: `{job.job_id}`", + f"• סטטוס: {status}", + f"• קטגוריה: {job.category.value}", + f"• סוג: {job.job_type.value}", + ] + + if job.interval_seconds: + lines.append(f"• אינטרוול: {_format_interval(job.interval_seconds)}") + + if history: + lines.append("\n**5 הרצות אחרונות:**") + for run in history[:5]: + icon = { + "completed": "✅", "failed": "❌", + "running": "🔄", "skipped": "⏭️" + }.get(run.status.value, "❓") + dur = "" + if run.ended_at and run.started_at: + dur = f" ({(run.ended_at - run.started_at).total_seconds():.1f}s)" + + line = f" {icon} {run.started_at.strftime('%d/%m %H:%M')}{dur}" + + # 🔗 אם נכשל, הוסף קישור ללוגים + if run.status.value == "failed": + logs_link = f"{monitor_base_url}/jobs/monitor?run_id={run.run_id}" + line += f"\n └─ [📋 ראה לוגים]({logs_link})" + + lines.append(line) + + return "\n".join(lines) + + # All jobs summary + jobs = registry.list_all() + if not jobs: + return "📋 אין jobs רשומים במערכת" + + categories: Dict[str, List[str]] = {} + for job in jobs: + cat = job.category.value + if cat not in categories: + categories[cat] = [] + status = "✅" if registry.is_enabled(job.job_id) else "❌" + categories[cat].append(f"{status} {job.name}") + + lines = ["🔄 **Background Jobs:**\n"] + for cat, items in categories.items(): + icon = { + "backup": "💾", "cache": "🗄️", "sync": "☁️", "cleanup": "🧹", + "monitoring": "📊", "batch": "📦", "other": "📋" + }.get(cat, "📋") + lines.append(f"**{icon} {cat}:**") + for item in items: + lines.append(f" {item}") + lines.append("") + + lines.append("_השתמש ב-`/jobs active` לצפייה בהרצות פעילות_") + lines.append("_השתמש ב-`/jobs failed` לצפייה בשגיאות אחרונות_") + return "\n".join(lines) + + +def _format_interval(seconds: int) -> str: + if seconds >= 86400: + return f"{seconds // 86400} ימים" + if seconds >= 3600: + return f"{seconds // 3600} שעות" + if seconds >= 60: + return f"{seconds // 60} דקות" + return f"{seconds} שניות" diff --git a/services/job_registry.py b/services/job_registry.py new file mode 100644 index 000000000..533ffd936 --- /dev/null +++ b/services/job_registry.py @@ -0,0 +1,120 @@ +# services/job_registry.py +""" +מודול לרישום מרכזי של כל ה-Background Jobs במערכת. + +JobRegistry הוא Singleton שמנהל את הגדרות ה-Jobs הידועים במערכת. +""" + +import threading +import os +import logging +from dataclasses import dataclass, field +from typing import Dict, List, Optional, Any +from enum import Enum + +logger = logging.getLogger(__name__) + + +class JobType(Enum): + """סוג ה-Job""" + REPEATING = "repeating" # חוזר לפי אינטרוול + ONCE = "once" # חד-פעמי + ON_DEMAND = "on_demand" # לפי דרישה (ידני/API) + + +class JobCategory(Enum): + """קטגוריית ה-Job""" + BACKUP = "backup" + CACHE = "cache" + SYNC = "sync" + CLEANUP = "cleanup" + MONITORING = "monitoring" + BATCH = "batch" + OTHER = "other" + + +@dataclass +class JobDefinition: + """הגדרת Job במערכת""" + job_id: str # מזהה ייחודי + name: str # שם תצוגה + description: str # תיאור + category: JobCategory # קטגוריה + job_type: JobType # סוג (חוזר/חד-פעמי/on-demand) + interval_seconds: Optional[int] = None # אינטרוול (ל-repeating) + enabled: bool = True # האם מופעל + env_toggle: Optional[str] = None # משתנה סביבה להפעלה/כיבוי + callback_name: str = "" # שם הפונקציה המופעלת + source_file: str = "" # קובץ מקור + metadata: Dict[str, Any] = field(default_factory=dict) + + +class JobRegistry: + """Singleton לרישום כל ה-Jobs במערכת""" + + _instance: Optional["JobRegistry"] = None + _lock = threading.Lock() + _jobs: Dict[str, JobDefinition] # Declared for mypy; initialized in __new__ + + def __new__(cls) -> "JobRegistry": + if cls._instance is None: + with cls._lock: + if cls._instance is None: + # חשוב: מאתחלים את _jobs לפני שחושפים את ה-instance + # כדי למנוע race condition שבו thread אחר רואה instance + # אבל _jobs עדיין לא קיים + new_instance = super().__new__(cls) + new_instance._jobs = {} + cls._instance = new_instance + return cls._instance + + def register(self, job: JobDefinition) -> None: + """רישום Job חדש""" + self._jobs[job.job_id] = job + logger.info(f"Registered job: {job.job_id} ({job.name})") + + def get(self, job_id: str) -> Optional[JobDefinition]: + """קבלת Job לפי ID""" + return self._jobs.get(job_id) + + def list_all(self) -> List[JobDefinition]: + """רשימת כל ה-Jobs""" + return list(self._jobs.values()) + + def list_by_category(self, category: JobCategory) -> List[JobDefinition]: + """רשימת Jobs לפי קטגוריה""" + return [j for j in self._jobs.values() if j.category == category] + + def is_enabled(self, job_id: str) -> bool: + """בדיקה האם Job מופעל""" + job = self._jobs.get(job_id) + if not job: + return False + if job.env_toggle: + return os.getenv(job.env_toggle, "").lower() in ("1", "true", "yes", "on") + return job.enabled + + def clear(self) -> None: + """מחיקת כל ה-Jobs (לשימוש בטסטים)""" + self._jobs.clear() + + +def register_job( + job_id: str, + name: str, + description: str, + category: JobCategory, + job_type: JobType, + **kwargs: Any +) -> JobDefinition: + """רישום Job חדש במערכת""" + job = JobDefinition( + job_id=job_id, + name=name, + description=description, + category=category, + job_type=job_type, + **kwargs + ) + JobRegistry().register(job) + return job diff --git a/services/job_tracker.py b/services/job_tracker.py new file mode 100644 index 000000000..d311c9b2f --- /dev/null +++ b/services/job_tracker.py @@ -0,0 +1,454 @@ +# services/job_tracker.py +""" +מודול למעקב אחרי הרצות של Background Jobs. + +JobTracker מנהל את מחזור החיים של הרצות: התחלה, עדכון התקדמות, סיום/כישלון. +שומר היסטוריה ב-MongoDB עם TTL Index לניקוי אוטומטי. +""" + +import uuid +import logging +import threading +from contextlib import contextmanager +from dataclasses import dataclass, field +from datetime import datetime, timezone +from enum import Enum +from typing import Any, Dict, Generator, List, Optional + +logger = logging.getLogger(__name__) + + +class JobStatus(Enum): + """סטטוס הרצה""" + PENDING = "pending" + RUNNING = "running" + COMPLETED = "completed" + FAILED = "failed" + CANCELLED = "cancelled" + SKIPPED = "skipped" + + +@dataclass +class JobLogEntry: + """רשומת לוג בודדת""" + timestamp: datetime + level: str # info/warning/error + message: str + details: Optional[Dict[str, Any]] = None + + +@dataclass +class JobRun: + """הרצה בודדת של Job""" + run_id: str # מזהה הרצה ייחודי + job_id: str # מזהה ה-Job + started_at: datetime + ended_at: Optional[datetime] = None + status: JobStatus = JobStatus.PENDING + progress: int = 0 # 0-100 + total_items: int = 0 # סה"כ פריטים לעיבוד + processed_items: int = 0 # פריטים שעובדו + error_message: Optional[str] = None + logs: List[JobLogEntry] = field(default_factory=list) + result: Optional[Dict[str, Any]] = None # תוצאה סופית + trigger: str = "scheduled" # scheduled/manual/api + user_id: Optional[int] = None # אם רלוונטי למשתמש + + +class JobAlreadyRunningError(Exception): + """נזרק כאשר מנסים להפעיל Job שכבר רץ""" + pass + + +class JobTracker: + """מעקב אחרי הרצות Jobs""" + + def __init__(self, db_manager: Any = None): + self._db = db_manager + self._active_runs: Dict[str, JobRun] = {} + + @property + def db(self) -> Any: + """Lazy loading של DB manager""" + if self._db is None: + try: + from database import db + self._db = db + except Exception: + logger.warning("Failed to initialize database connection for job tracking", exc_info=True) + self._db = None + return self._db + + @property + def db_name(self) -> str: + """שם ה-database""" + if self.db is None: + return "test" + # תמיכה ב-mock DB עם db_name כ-attribute + if hasattr(self.db, "db_name"): + return self.db.db_name + # תמיכה ב-DatabaseManager אמיתי + if hasattr(self.db, "db") and hasattr(self.db.db, "name"): + return self.db.db.name + return "code_keeper_bot" + + def start_run( + self, + job_id: str, + trigger: str = "scheduled", + user_id: Optional[int] = None, + total_items: int = 0, + allow_concurrent: bool = False + ) -> JobRun: + """התחלת הרצה חדשה + + Args: + job_id: מזהה ה-Job + trigger: מה הפעיל את ההרצה (scheduled/manual/api) + user_id: מזהה משתמש (אם רלוונטי) + total_items: סה"כ פריטים לעיבוד + allow_concurrent: האם לאפשר הרצות מקבילות של אותו Job + + Raises: + JobAlreadyRunningError: אם Job כבר רץ ו-allow_concurrent=False + """ + # 🔒 מניעת הרצות מקבילות (Singleton Jobs) + if not allow_concurrent: + existing = [ + r for r in self._active_runs.values() + if r.job_id == job_id and r.status == JobStatus.RUNNING + ] + if existing: + raise JobAlreadyRunningError( + f"Job '{job_id}' is already running (run_id: {existing[0].run_id})" + ) + + run = JobRun( + run_id=str(uuid.uuid4())[:12], + job_id=job_id, + started_at=datetime.now(timezone.utc), + status=JobStatus.RUNNING, + trigger=trigger, + user_id=user_id, + total_items=total_items + ) + self._active_runs[run.run_id] = run + self._persist_run(run) + + # Emit observability event (best-effort, fail-open) + try: + from observability import emit_event + emit_event("job_started", severity="info", job_id=job_id, run_id=run.run_id) + except Exception: + # Observability is optional - don't fail job tracking if unavailable + logger.warning("Failed to emit job_started event for %s", job_id, exc_info=True) + + return run + + def update_progress( + self, + run_id: str, + processed: int, + total: Optional[int] = None, + message: Optional[str] = None + ) -> None: + """עדכון התקדמות""" + run = self._active_runs.get(run_id) + if not run: + return + + run.processed_items = processed + if total is not None: + run.total_items = total + + if run.total_items > 0: + run.progress = int((processed / run.total_items) * 100) + + if message: + self.add_log(run_id, "info", message) + + self._persist_run(run) + + def add_log( + self, + run_id: str, + level: str, + message: str, + details: Optional[Dict[str, Any]] = None + ) -> None: + """הוספת לוג להרצה""" + run = self._active_runs.get(run_id) + if not run: + return + + entry = JobLogEntry( + timestamp=datetime.now(timezone.utc), + level=level, + message=message, + details=details + ) + run.logs.append(entry) + + # שמירה ל-DB רק כל 10 לוגים או ב-error + if level == "error" or len(run.logs) % 10 == 0: + self._persist_run(run) + + def complete_run( + self, + run_id: str, + result: Optional[Dict[str, Any]] = None + ) -> None: + """סיום הרצה בהצלחה""" + run = self._active_runs.get(run_id) + if not run: + return + + run.status = JobStatus.COMPLETED + run.ended_at = datetime.now(timezone.utc) + run.progress = 100 + run.result = result + + self._persist_run(run) + self._active_runs.pop(run_id, None) + + duration = (run.ended_at - run.started_at).total_seconds() + try: + from observability import emit_event + emit_event( + "job_completed", + severity="info", + job_id=run.job_id, + run_id=run_id, + duration_seconds=duration + ) + except Exception: + # Observability is optional - don't fail job tracking if unavailable + logger.warning("Failed to emit job_completed event for %s", run_id, exc_info=True) + + def fail_run( + self, + run_id: str, + error_message: str + ) -> None: + """סיום הרצה בכישלון""" + run = self._active_runs.get(run_id) + if not run: + return + + run.status = JobStatus.FAILED + run.ended_at = datetime.now(timezone.utc) + run.error_message = error_message + self.add_log(run_id, "error", error_message) + + self._persist_run(run) + self._active_runs.pop(run_id, None) + + try: + from observability import emit_event + emit_event( + "job_failed", + severity="error", + job_id=run.job_id, + run_id=run_id, + error=error_message + ) + except Exception: + # Observability is optional - don't fail job tracking if unavailable + logger.warning("Failed to emit job_failed event for %s", run_id, exc_info=True) + + @contextmanager + def track( + self, + job_id: str, + trigger: str = "scheduled", + user_id: Optional[int] = None, + allow_concurrent: bool = False + ) -> Generator[JobRun, None, None]: + """Context manager למעקב אחרי הרצה + + שימוש: + with tracker.track("my_job") as run: + # ... לוגיקה ... + tracker.add_log(run.run_id, "info", "Processing...") + + ⚠️ חשוב: השתמש ב-run.run_id ולא ב-get_active_runs()[0]! + """ + run = self.start_run(job_id, trigger, user_id, allow_concurrent=allow_concurrent) + try: + yield run + self.complete_run(run.run_id) + except Exception as e: + self.fail_run(run.run_id, str(e)) + raise + + def _persist_run(self, run: JobRun) -> None: + """שמירת הרצה ל-DB""" + try: + if self.db is None: + return + + doc = { + "run_id": run.run_id, + "job_id": run.job_id, + "started_at": run.started_at, + "ended_at": run.ended_at, + "status": run.status.value, + "progress": run.progress, + "total_items": run.total_items, + "processed_items": run.processed_items, + "error_message": run.error_message, + "logs": [ + { + "timestamp": log.timestamp, + "level": log.level, + "message": log.message, + "details": log.details + } + for log in run.logs[-50:] # שמירת 50 לוגים אחרונים + ], + "result": run.result, + "trigger": run.trigger, + "user_id": run.user_id + } + + # תמיכה ב-mock DB + if hasattr(self.db, "client"): + self.db.client[self.db_name]["job_runs"].update_one( + {"run_id": run.run_id}, + {"$set": doc}, + upsert=True + ) + elif hasattr(self.db, "db") and self.db.db is not None: + self.db.db["job_runs"].update_one( + {"run_id": run.run_id}, + {"$set": doc}, + upsert=True + ) + except Exception as e: + logger.error(f"Failed to persist job run: {e}") + + def get_run(self, run_id: str) -> Optional[JobRun]: + """קבלת הרצה לפי ID""" + if run_id in self._active_runs: + return self._active_runs[run_id] + + try: + if self.db is None: + return None + + doc = None + if hasattr(self.db, "client"): + doc = self.db.client[self.db_name]["job_runs"].find_one( + {"run_id": run_id} + ) + elif hasattr(self.db, "db") and self.db.db is not None: + doc = self.db.db["job_runs"].find_one({"run_id": run_id}) + + if doc: + return self._doc_to_run(doc) + except Exception: + # DB lookup failed - return None gracefully + logger.error("Failed to get run %s from DB", run_id, exc_info=True) + return None + + def get_job_history( + self, + job_id: str, + limit: int = 20 + ) -> List[JobRun]: + """היסטוריית הרצות של Job""" + try: + if self.db is None: + return [] + + cursor = None + if hasattr(self.db, "client"): + cursor = self.db.client[self.db_name]["job_runs"].find( + {"job_id": job_id} + ).sort("started_at", -1).limit(limit) + elif hasattr(self.db, "db") and self.db.db is not None: + cursor = self.db.db["job_runs"].find( + {"job_id": job_id} + ).sort("started_at", -1).limit(limit) + + if cursor: + return [self._doc_to_run(doc) for doc in cursor] + except Exception: + # DB lookup failed - return empty list gracefully + logger.error("Failed to get job history for %s", job_id, exc_info=True) + return [] + + def get_active_runs(self) -> List[JobRun]: + """רשימת הרצות פעילות""" + return list(self._active_runs.values()) + + def get_failed_runs(self, limit: int = 10) -> List[JobRun]: + """רשימת הרצות שנכשלו""" + try: + if self.db is None: + return [] + + cursor = None + if hasattr(self.db, "client"): + cursor = self.db.client[self.db_name]["job_runs"].find( + {"status": "failed"} + ).sort("ended_at", -1).limit(limit) + elif hasattr(self.db, "db") and self.db.db is not None: + cursor = self.db.db["job_runs"].find( + {"status": "failed"} + ).sort("ended_at", -1).limit(limit) + + if cursor: + return [self._doc_to_run(doc) for doc in cursor] + except Exception: + # DB lookup failed - return empty list gracefully + logger.error("Failed to get failed runs from DB", exc_info=True) + return [] + + def _doc_to_run(self, doc: dict) -> JobRun: + """המרת מסמך DB ל-JobRun""" + logs = [ + JobLogEntry( + timestamp=log["timestamp"], + level=log["level"], + message=log["message"], + details=log.get("details") + ) + for log in doc.get("logs", []) + ] + return JobRun( + run_id=doc["run_id"], + job_id=doc["job_id"], + started_at=doc["started_at"], + ended_at=doc.get("ended_at"), + status=JobStatus(doc.get("status", "completed")), + progress=doc.get("progress", 100), + total_items=doc.get("total_items", 0), + processed_items=doc.get("processed_items", 0), + error_message=doc.get("error_message"), + logs=logs, + result=doc.get("result"), + trigger=doc.get("trigger", "scheduled"), + user_id=doc.get("user_id") + ) + + +# Singleton instance with thread-safe initialization +_tracker: Optional[JobTracker] = None +_tracker_lock = threading.Lock() + + +def get_job_tracker() -> JobTracker: + """קבלת Singleton instance של JobTracker (thread-safe)""" + global _tracker + with _tracker_lock: + if _tracker is None: + _tracker = JobTracker() + return _tracker + + +def reset_job_tracker() -> None: + """איפוס ה-tracker (לשימוש בטסטים)""" + global _tracker + with _tracker_lock: + _tracker = None diff --git a/services/register_jobs.py b/services/register_jobs.py new file mode 100644 index 000000000..11ce58858 --- /dev/null +++ b/services/register_jobs.py @@ -0,0 +1,147 @@ +# services/register_jobs.py +""" +רישום כל ה-Background Jobs במערכת. +יש לייבא קובץ זה ב-main.py לאחר אתחול ה-Application. +""" + +from services.job_registry import ( + register_job, + JobCategory, + JobType, +) + + +def register_all_jobs() -> None: + """רישום כל ה-Jobs המוכרים""" + + # === Backup Jobs === + register_job( + job_id="backups_cleanup", + name="ניקוי גיבויים", + description="מחיקת גיבויים ישנים לפי מדיניות retention", + category=JobCategory.CLEANUP, + job_type=JobType.REPEATING, + interval_seconds=86400, + env_toggle="BACKUPS_CLEANUP_ENABLED", + callback_name="_backups_cleanup_job", + source_file="main.py" + ) + + # === Cache Jobs === + register_job( + job_id="cache_maintenance", + name="תחזוקת קאש", + description="ניקוי רשומות קאש שפגו תוקפן", + category=JobCategory.CACHE, + job_type=JobType.REPEATING, + interval_seconds=600, + enabled=True, + callback_name="_cache_maintenance_job", + source_file="main.py" + ) + + register_job( + job_id="cache_warming", + name="חימום קאש", + description="טעינה מראש של נתונים נפוצים לקאש", + category=JobCategory.CACHE, + job_type=JobType.REPEATING, + interval_seconds=900, + env_toggle="CACHE_WARMING_ENABLED", + callback_name="_cache_warming_job", + source_file="main.py" + ) + + # === Drive Sync Jobs === + register_job( + job_id="drive_reschedule", + name="תזמון Drive", + description="שמירה על תזמוני גיבוי אוטומטי ל-Google Drive", + category=JobCategory.SYNC, + job_type=JobType.REPEATING, + interval_seconds=900, + enabled=True, + callback_name="_reschedule_drive_jobs", + source_file="main.py" + ) + + # === Monitoring Jobs === + register_job( + job_id="sentry_poll", + name="סקירת Sentry", + description="משיכת אירועים חדשים מ-Sentry", + category=JobCategory.MONITORING, + job_type=JobType.REPEATING, + interval_seconds=300, + env_toggle="SENTRY_POLL_ENABLED", + callback_name="_sentry_poll_job", + source_file="main.py" + ) + + register_job( + job_id="predictive_sampler", + name="דגימה חזויה", + description="איסוף מדדים לזיהוי אנומליות", + category=JobCategory.MONITORING, + job_type=JobType.REPEATING, + interval_seconds=60, + enabled=True, + callback_name="_predictive_sampler_job", + source_file="main.py" + ) + + register_job( + job_id="weekly_admin_report", + name='דו"ח שבועי', + description='שליחת דו"ח סיכום שבועי לאדמינים', + category=JobCategory.MONITORING, + job_type=JobType.REPEATING, + interval_seconds=7 * 24 * 3600, + enabled=True, + callback_name="_weekly_admin_report", + source_file="main.py" + ) + + # === Reminders === + register_job( + job_id="recurring_reminders_check", + name="בדיקת תזכורות", + description="עיבוד תזכורות חוזרות", + category=JobCategory.OTHER, + job_type=JobType.REPEATING, + interval_seconds=3600, + enabled=True, + callback_name="_check_recurring_reminders", + source_file="reminders/scheduler.py" + ) + + # === Batch Processing === + register_job( + job_id="batch_analyze", + name="ניתוח קבצים", + description="ניתוח batch של קבצים", + category=JobCategory.BATCH, + job_type=JobType.ON_DEMAND, + callback_name="analyze_files_batch", + source_file="batch_processor.py" + ) + + register_job( + job_id="batch_validate", + name="בדיקת תקינות", + description="בדיקת תקינות batch של קבצים", + category=JobCategory.BATCH, + job_type=JobType.ON_DEMAND, + callback_name="validate_files_batch", + source_file="batch_processor.py" + ) + + register_job( + job_id="batch_export", + name="ייצוא קבצים", + description="ייצוא batch של קבצים", + category=JobCategory.BATCH, + job_type=JobType.ON_DEMAND, + callback_name="export_files_batch", + source_file="batch_processor.py" + ) diff --git a/services/webserver.py b/services/webserver.py index 7422e9c55..1a65fa850 100644 --- a/services/webserver.py +++ b/services/webserver.py @@ -1173,6 +1173,154 @@ async def db_health_summary_view(request: web.Request) -> web.Response: app.router.add_get("/api/db/collections", db_health_collections_view) app.router.add_get("/api/db/health", db_health_summary_view) + # === Background Jobs Monitor API === + async def get_jobs_list(request: web.Request) -> web.Response: + """GET /api/jobs - רשימת כל ה-jobs""" + try: + from services.job_registry import JobRegistry + registry = JobRegistry() + jobs = [] + + for job in registry.list_all(): + jobs.append({ + "job_id": job.job_id, + "name": job.name, + "description": job.description, + "category": job.category.value, + "type": job.job_type.value, + "interval_seconds": job.interval_seconds, + "enabled": registry.is_enabled(job.job_id), + "env_toggle": job.env_toggle, + }) + + return web.json_response({"jobs": jobs}) + except Exception as e: + logger.error(f"get_jobs_list error: {e}") + return web.json_response({"error": "failed", "jobs": []}, status=500) + + async def get_job_detail(request: web.Request) -> web.Response: + """GET /api/jobs/{job_id} - פרטי job ספציפי""" + try: + job_id = request.match_info.get("job_id") + from services.job_registry import JobRegistry + from services.job_tracker import get_job_tracker + registry = JobRegistry() + tracker = get_job_tracker() + + job = registry.get(job_id) + if not job: + return web.json_response({"error": "Job not found"}, status=404) + + history = tracker.get_job_history(job_id, limit=20) + active = [r for r in tracker.get_active_runs() if r.job_id == job_id] + + return web.json_response({ + "job": { + "job_id": job.job_id, + "name": job.name, + "description": job.description, + "category": job.category.value, + "type": job.job_type.value, + "interval_seconds": job.interval_seconds, + "enabled": registry.is_enabled(job.job_id), + "source_file": job.source_file, + }, + "active_runs": [_run_to_dict(r) for r in active], + "history": [_run_to_dict(r) for r in history], + }) + except Exception as e: + logger.error(f"get_job_detail error: {e}") + return web.json_response({"error": "failed"}, status=500) + + async def get_run_detail(request: web.Request) -> web.Response: + """GET /api/jobs/runs/{run_id} - פרטי הרצה""" + try: + run_id = request.match_info.get("run_id") + from services.job_tracker import get_job_tracker + tracker = get_job_tracker() + + run = tracker.get_run(run_id) + if not run: + return web.json_response({"error": "Run not found"}, status=404) + + return web.json_response({"run": _run_to_dict(run, include_logs=True)}) + except Exception as e: + logger.error(f"get_run_detail error: {e}") + return web.json_response({"error": "failed"}, status=500) + + async def get_active_runs_view(request: web.Request) -> web.Response: + """GET /api/jobs/active - הרצות פעילות""" + try: + from services.job_tracker import get_job_tracker + tracker = get_job_tracker() + runs = tracker.get_active_runs() + + return web.json_response({ + "active_runs": [_run_to_dict(r) for r in runs] + }) + except Exception as e: + logger.error(f"get_active_runs error: {e}") + return web.json_response({"error": "failed", "active_runs": []}, status=500) + + async def trigger_job_view(request: web.Request) -> web.Response: + """POST /api/jobs/{job_id}/trigger - הפעלה ידנית""" + try: + job_id = request.match_info.get("job_id") + from services.job_registry import JobRegistry + registry = JobRegistry() + + job = registry.get(job_id) + if not job: + return web.json_response({"error": "Job not found"}, status=404) + + # TODO: implement actual trigger via job_queue + # כרגע מחזירים 501 Not Implemented כדי לא להטעות את המשתמש + return web.json_response({ + "ok": False, + "error": "not_implemented", + "message": f"הפעלה ידנית של Job עדיין לא מומשה. Job: {job_id}" + }, status=501) + except Exception as e: + logger.error(f"trigger_job error: {e}") + return web.json_response({"error": "failed"}, status=500) + + def _run_to_dict(run, include_logs: bool = False) -> dict: + """המרת JobRun ל-dict""" + d = { + "run_id": run.run_id, + "job_id": run.job_id, + "started_at": run.started_at.isoformat() if run.started_at else None, + "ended_at": run.ended_at.isoformat() if run.ended_at else None, + "status": run.status.value, + "progress": run.progress, + "total_items": run.total_items, + "processed_items": run.processed_items, + "error_message": run.error_message, + "trigger": run.trigger, + "user_id": run.user_id, + "duration_seconds": ( + (run.ended_at - run.started_at).total_seconds() + if run.ended_at and run.started_at else None + ), + } + if include_logs: + d["logs"] = [ + { + "timestamp": log.timestamp.isoformat(), + "level": log.level, + "message": log.message, + } + for log in run.logs + ] + return d + + # Register Jobs Monitor routes + app.router.add_get("/api/jobs", get_jobs_list) + app.router.add_get("/api/jobs/active", get_active_runs_view) + app.router.add_get("/api/jobs/{job_id}", get_job_detail) + app.router.add_get("/api/jobs/runs/{run_id}", get_run_detail) + app.router.add_post("/api/jobs/{job_id}/trigger", trigger_job_view) + return app diff --git a/tests/test_job_tracker.py b/tests/test_job_tracker.py new file mode 100644 index 000000000..24fd14d99 --- /dev/null +++ b/tests/test_job_tracker.py @@ -0,0 +1,404 @@ +# tests/test_job_tracker.py +""" +Unit Tests עבור JobTracker ו-JobRegistry. +""" + +import pytest +from unittest.mock import MagicMock + +from services.job_tracker import ( + JobTracker, + JobStatus, + JobAlreadyRunningError, + get_job_tracker, + reset_job_tracker, +) +from services.job_registry import ( + JobRegistry, + register_job, + JobCategory, + JobType, +) + + +# === Mock DB classes === + +class MockCollection: + """Mock MongoDB collection""" + def __init__(self): + self.docs: dict = {} + + def update_one(self, query, update, upsert=False): + run_id = query.get("run_id") + if run_id: + self.docs[run_id] = update.get("$set", {}) + return MagicMock(acknowledged=True, modified_count=1) + + def find_one(self, query): + run_id = query.get("run_id") + return self.docs.get(run_id) + + def find(self, query): + job_id = query.get("job_id") + status = query.get("status") + results = [] + for doc in self.docs.values(): + if job_id and doc.get("job_id") != job_id: + continue + if status and doc.get("status") != status: + continue + results.append(doc) + return MockCursor(results) + + +class MockCursor: + """Mock MongoDB cursor""" + def __init__(self, docs): + self._docs = docs + + def sort(self, *args): + return self + + def limit(self, n): + self._docs = self._docs[:n] + return self + + def __iter__(self): + return iter(self._docs) + + +class MockDB: + """Mock database manager""" + def __init__(self): + self._collections: dict = {} + self.db_name = "test" + + @property + def client(self): + return {self.db_name: self} + + def __getitem__(self, name: str): + if name not in self._collections: + self._collections[name] = MockCollection() + return self._collections[name] + + +# === Fixtures === + +@pytest.fixture +def tracker(): + """Tracker with mock DB""" + mock_db = MockDB() + return JobTracker(mock_db) + + +@pytest.fixture +def registry(): + """Clean registry for each test""" + reg = JobRegistry() + reg.clear() + return reg + + +# === JobTracker Tests === + +class TestJobTracker: + """Tests for JobTracker""" + + def test_start_and_complete_run(self, tracker): + """Test starting and completing a job run""" + run = tracker.start_run("test_job") + + assert run.status == JobStatus.RUNNING + assert run.run_id in [r.run_id for r in tracker.get_active_runs()] + + tracker.complete_run(run.run_id, result={"count": 5}) + + assert run.run_id not in [r.run_id for r in tracker.get_active_runs()] + + def test_fail_run(self, tracker): + """Test failing a job run""" + run = tracker.start_run("test_job") + tracker.fail_run(run.run_id, "Test error") + + assert run.status == JobStatus.FAILED + assert run.error_message == "Test error" + + def test_track_context_manager_success(self, tracker): + """Test track context manager on success""" + with tracker.track("test_job") as run: + tracker.add_log(run.run_id, "info", "Processing...") + + assert run.status == JobStatus.COMPLETED + assert run.progress == 100 + + def test_track_context_manager_failure(self, tracker): + """Test track context manager on error""" + with pytest.raises(ValueError): + with tracker.track("test_job") as run: + raise ValueError("Oops") + + assert run.status == JobStatus.FAILED + assert "Oops" in run.error_message + + def test_update_progress(self, tracker): + """Test updating progress""" + run = tracker.start_run("test_job", total_items=100) + + tracker.update_progress(run.run_id, processed=50) + assert run.progress == 50 + assert run.processed_items == 50 + + tracker.update_progress(run.run_id, processed=100) + assert run.progress == 100 + + def test_add_log(self, tracker): + """Test adding logs""" + run = tracker.start_run("test_job") + + tracker.add_log(run.run_id, "info", "Step 1") + tracker.add_log(run.run_id, "warning", "Step 2") + tracker.add_log(run.run_id, "error", "Step 3 failed") + + assert len(run.logs) == 3 + assert run.logs[0].level == "info" + assert run.logs[2].level == "error" + + def test_prevent_concurrent_runs(self, tracker): + """Test that concurrent runs are prevented by default""" + run1 = tracker.start_run("singleton_job") + + with pytest.raises(JobAlreadyRunningError): + tracker.start_run("singleton_job") + + # After completing, should allow new run + tracker.complete_run(run1.run_id) + run2 = tracker.start_run("singleton_job") + assert run2.run_id != run1.run_id + + def test_allow_concurrent_runs(self, tracker): + """Test that concurrent runs can be allowed""" + run1 = tracker.start_run("parallel_job", allow_concurrent=True) + run2 = tracker.start_run("parallel_job", allow_concurrent=True) + + assert run1.run_id != run2.run_id + assert len(tracker.get_active_runs()) == 2 + + def test_get_active_runs(self, tracker): + """Test getting active runs""" + run1 = tracker.start_run("job1", allow_concurrent=True) + run2 = tracker.start_run("job2", allow_concurrent=True) + + active = tracker.get_active_runs() + assert len(active) == 2 + assert run1.run_id in [r.run_id for r in active] + assert run2.run_id in [r.run_id for r in active] + + tracker.complete_run(run1.run_id) + active = tracker.get_active_runs() + assert len(active) == 1 + + def test_run_with_user_id(self, tracker): + """Test run with user_id""" + run = tracker.start_run("user_job", user_id=12345, trigger="manual") + + assert run.user_id == 12345 + assert run.trigger == "manual" + + +class TestJobTrackerPersistence: + """Tests for JobTracker persistence""" + + def test_persist_run_to_db(self, tracker): + """Test that runs are persisted to DB""" + run = tracker.start_run("test_job") + tracker.add_log(run.run_id, "info", "Test message") + tracker.complete_run(run.run_id) + + # Check that doc was saved - access through the mock structure + collection = tracker._db._collections.get("job_runs") + assert collection is not None + doc = collection.docs.get(run.run_id) + assert doc is not None + assert doc["job_id"] == "test_job" + assert doc["status"] == "completed" + + def test_get_run_from_db(self, tracker): + """Test getting a run from DB""" + run = tracker.start_run("test_job") + tracker.complete_run(run.run_id) + + # Remove from active runs to force DB lookup + tracker._active_runs.clear() + + retrieved = tracker.get_run(run.run_id) + assert retrieved is not None + assert retrieved.job_id == "test_job" + + def test_get_job_history(self, tracker): + """Test getting job history""" + # Create multiple runs + for i in range(3): + run = tracker.start_run("history_job", allow_concurrent=True) + tracker.complete_run(run.run_id) + + history = tracker.get_job_history("history_job", limit=10) + assert len(history) == 3 + + +# === JobRegistry Tests === + +class TestJobRegistry: + """Tests for JobRegistry""" + + def test_registry_singleton(self, registry): + """Test that JobRegistry is a singleton""" + reg1 = JobRegistry() + reg2 = JobRegistry() + assert reg1 is reg2 + + def test_register_and_get_job(self, registry): + """Test registering and getting a job""" + register_job( + job_id="test_backup", + name="Test Backup", + description="A test job", + category=JobCategory.BACKUP, + job_type=JobType.REPEATING, + interval_seconds=3600 + ) + + retrieved = registry.get("test_backup") + assert retrieved is not None + assert retrieved.job_id == "test_backup" + assert retrieved.name == "Test Backup" + assert retrieved.interval_seconds == 3600 + + def test_list_all_jobs(self, registry): + """Test listing all jobs""" + register_job( + job_id="job1", + name="Job 1", + description="First job", + category=JobCategory.CACHE, + job_type=JobType.REPEATING + ) + register_job( + job_id="job2", + name="Job 2", + description="Second job", + category=JobCategory.BACKUP, + job_type=JobType.ONCE + ) + + jobs = registry.list_all() + assert len(jobs) == 2 + + def test_list_by_category(self, registry): + """Test listing jobs by category""" + register_job( + job_id="cache1", + name="Cache 1", + description="Cache job", + category=JobCategory.CACHE, + job_type=JobType.REPEATING + ) + register_job( + job_id="backup1", + name="Backup 1", + description="Backup job", + category=JobCategory.BACKUP, + job_type=JobType.REPEATING + ) + + cache_jobs = registry.list_by_category(JobCategory.CACHE) + assert len(cache_jobs) == 1 + assert cache_jobs[0].job_id == "cache1" + + def test_is_enabled_default(self, registry): + """Test is_enabled with default enabled=True""" + register_job( + job_id="enabled_job", + name="Enabled Job", + description="Should be enabled", + category=JobCategory.OTHER, + job_type=JobType.ONCE, + enabled=True + ) + + assert registry.is_enabled("enabled_job") is True + + def test_is_enabled_with_env_toggle(self, registry, monkeypatch): + """Test is_enabled with env_toggle""" + register_job( + job_id="toggle_job", + name="Toggle Job", + description="Controlled by env", + category=JobCategory.OTHER, + job_type=JobType.ONCE, + env_toggle="MY_JOB_ENABLED" + ) + + # Not set - should be disabled + monkeypatch.delenv("MY_JOB_ENABLED", raising=False) + assert registry.is_enabled("toggle_job") is False + + # Set to true + monkeypatch.setenv("MY_JOB_ENABLED", "true") + assert registry.is_enabled("toggle_job") is True + + # Set to false + monkeypatch.setenv("MY_JOB_ENABLED", "false") + assert registry.is_enabled("toggle_job") is False + + def test_is_enabled_nonexistent_job(self, registry): + """Test is_enabled for nonexistent job""" + assert registry.is_enabled("nonexistent") is False + + +# === Integration Tests === + +class TestJobTrackerRegistryIntegration: + """Integration tests for JobTracker and JobRegistry""" + + def test_tracker_with_registered_job(self, tracker, registry): + """Test using tracker with a registered job""" + register_job( + job_id="integrated_job", + name="Integrated Job", + description="Test integration", + category=JobCategory.MONITORING, + job_type=JobType.REPEATING, + interval_seconds=60 + ) + + job = registry.get("integrated_job") + assert job is not None + + with tracker.track(job.job_id) as run: + tracker.add_log(run.run_id, "info", f"Running {job.name}") + + assert run.status == JobStatus.COMPLETED + + +# === Global Singleton Tests === + +class TestGlobalSingleton: + """Tests for global singleton functions""" + + def test_get_job_tracker_singleton(self): + """Test that get_job_tracker returns singleton""" + reset_job_tracker() + + tracker1 = get_job_tracker() + tracker2 = get_job_tracker() + + assert tracker1 is tracker2 + + def test_reset_job_tracker(self): + """Test resetting the tracker singleton""" + tracker1 = get_job_tracker() + reset_job_tracker() + tracker2 = get_job_tracker() + + assert tracker1 is not tracker2 diff --git a/webapp/app.py b/webapp/app.py index cb50a596b..98546f92f 100644 --- a/webapp/app.py +++ b/webapp/app.py @@ -6794,6 +6794,32 @@ def dashboard(): error="אירעה שגיאה בטעינת הנתונים. אנא נסה שוב.", bot_username=BOT_USERNAME_CLEAN) + +@app.route('/jobs/monitor') +@login_required +def jobs_monitor(): + """דף ניטור Background Jobs""" + job_registration_error = None + try: + # רישום ה-Jobs אם טרם נרשמו + try: + from services.register_jobs import register_all_jobs + from services.job_registry import JobRegistry + if not JobRegistry().list_all(): + register_all_jobs() + except Exception as e: + logger.exception("Failed to register background jobs in jobs_monitor") + job_registration_error = "אירעה שגיאה ברישום משימות הרקע. ייתכן שחלק מהמשימות אינן פעילות." + return render_template('jobs_monitor.html', + user=session.get('user_data', {}), + job_registration_error=job_registration_error) + except Exception as e: + logger.error(f"jobs_monitor error: {e}") + return render_template('jobs_monitor.html', + user=session.get('user_data', {}), + error="שגיאה בטעינת הדף") + + @app.route('/files') @app.route('/files', endpoint='files_page') @login_required diff --git a/webapp/templates/jobs_monitor.html b/webapp/templates/jobs_monitor.html new file mode 100644 index 000000000..1202054a9 --- /dev/null +++ b/webapp/templates/jobs_monitor.html @@ -0,0 +1,646 @@ +{% extends "base.html" %} +{% block title %}Background Jobs Monitor{% endblock %} + +{% block content %} +
+ + + +
+

⚡ הרצות פעילות

+
+
טוען...
+
+
+ + +
+

📋 כל ה-Jobs

+ +
+ + + + + + + +
+ +
+
טוען...
+
+
+ + + +
+ + + + +{% endblock %} diff --git a/webapp/templates/settings.html b/webapp/templates/settings.html index 9e2a27461..51442b4a4 100644 --- a/webapp/templates/settings.html +++ b/webapp/templates/settings.html @@ -437,6 +437,24 @@

פתח דשבורד +
+

🔄 Background Jobs Monitor

+

+ ניטור ובקרה של כל התהליכים שרצים ברקע: גיבויים, ניקוי קאש, סנכרונים ותזכורות +

+
+ + + פתח דשבורד + +
+