crowdata / app /tasks /monitoring.py
YOSOYYONOSOYOTRO's picture
Upload folder using huggingface_hub (part 2)
4223796 verified
Raw
History Blame Contribute Delete
4.95 kB
"""
Monitoring Tasks — CrowData (Celery version)
Distributed monitoring tasks instead of daemon loop.
"""
import asyncio
import logging
from datetime import datetime, timezone
from sqlalchemy import select
from celery import shared_task
from celery.schedules import crontab
from app.database import AsyncSessionLocal
from app.reports.models import MonitorTask
from app.reports.service import get_person_report, get_company_report
from app.cache.redis_client import cache_get
from app.database import AsyncSessionLocal
logger = logging.getLogger(__name__)
try:
from app.utils.email_service import send_monitoring_alert
except ImportError:
send_monitoring_alert = None
@shared_task(bind=True, max_retries=3, default_retry_delay=60)
def check_monitor_task(self, task_id: int):
"""
Celery task to check a single monitor task.
Runs in isolation, can be scaled horizontally.
"""
async def _check():
async with AsyncSessionLocal() as db:
task = await db.get(MonitorTask, task_id)
if not task or not task.active:
return {"status": "skipped", "reason": "task not found or inactive"}
logger.info(f"Monitoring {task.type} -> {task.identifier} (User: {task.user_id})")
# Get cached previous report
cache_key = f"persona:{task.identifier}" if task.type == "persona" else f"empresa:{task.identifier}"
old_data = await cache_get(cache_key)
# Fetch fresh report
new_report = None
try:
if task.type == "persona":
new_report = await get_person_report(task.identifier)
else:
new_report = await get_company_report(task.identifier)
except Exception as e:
logger.error(f"Error fetching report for {task.identifier}: {e}")
raise self.retry(exc=e, countdown=60 * (2 ** self.request.retries))
if not old_data or not new_report:
return {"status": "no_data", "task_id": task_id}
# Change detection
changes = []
# 1. Judicial causes
old_causas = len(old_data.get("judicial", {}).get("causas", []))
new_causas = len(new_report.judicial.causas) if hasattr(new_report.judicial, 'causas') else 0
if new_causas > old_causas:
changes.append(f"Nueva causa judicial detectada ({new_causas - old_causas} adicional/es)")
# 2. BCRA situation
old_bcra = old_data.get("financiero", {}).get("bcra_situacion_actual", 1)
new_bcra = getattr(new_report.financiero, 'bcra_situacion_actual', None) or 1
if new_bcra != old_bcra:
changes.append(f"Cambio BCRA: {old_bcra}{new_bcra}")
# 3. INPI marcas
old_marcas = len(old_data.get("marcas_inpi", []))
new_marcas = len(new_report.marcas_inpi) if hasattr(new_report, 'marcas_inpi') else 0
if new_marcas > old_marcas:
changes.append(f"Nueva marca INPI registrada")
# Trigger alerts if changes
if changes:
logger.warning(f"🚨 ALERTA para {task.identifier}: {changes}")
return {"status": "alert_sent", "task_id": task_id, "changes": changes}
return {"status": "no_changes", "task_id": task_id}
try:
return asyncio.run(_check())
except Exception as e:
logger.error(f"Monitor task {task_id} failed: {e}")
raise self.retry(exc=e, countdown=60 * (2 ** self.request.retries))
@shared_task
def check_all_monitors():
"""
Periodic task to queue all active monitor tasks.
Runs via Celery Beat (cron: daily at 3 AM).
"""
async def _queue_all():
async with AsyncSessionLocal() as db:
stmt = select(MonitorTask).where(MonitorTask.active == True)
result = await db.execute(stmt)
tasks = result.scalars().all()
for task in tasks:
check_monitor_task.delay(task.id)
logger.info(f"Queued {len(tasks)} monitor tasks")
return {"queued": len(tasks)}
return asyncio.run(_queue_all())
# Celery Beat schedule
CELERY_BEAT_SCHEDULE = {
"check-all-monitors": {
"task": "app.tasks.monitoring.check_all_monitors",
"schedule": crontab(hour=3, minute=0), # 3 AM daily
},
"report-cache-cleanup": {
"task": "app.tasks.reports.cleanup_old_cache",
"schedule": crontab(hour=4, minute=30), # 4:30 AM daily
},
"scraper-health-check": {
"task": "app.tasks.scrapers.health_check_all",
"schedule": crontab(minute="*/15"), # Every 15 minutes
},
}
# Import crontab
from celery.schedules import crontab