Files
LogHive/backend/celery_worker/tasks.py
T
2026-05-09 14:55:14 +08:00

83 lines
2.6 KiB
Python

"""Celery tasks for periodic and async processing."""
import asyncio
import logging
from datetime import datetime, timezone, timedelta
from sqlalchemy import delete
from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker
from app.config import settings
from app.models.log import LogEntry
from app.services.alert_service import AlertService
logger = logging.getLogger(__name__)
from celery_worker.celery_app import celery_app
@celery_app.task(bind=True, max_retries=3)
def evaluate_alerts(self):
"""Evaluate all alert rules and trigger notifications.
Runs periodically (default: every 60 seconds via beat schedule).
"""
try:
loop = asyncio.get_event_loop()
if loop.is_closed():
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
result = loop.run_until_complete(_async_evaluate_alerts())
logger.info("Alert evaluation complete: %d triggers", len(result))
return {"triggered": len(result)}
except Exception as e:
logger.error("Alert evaluation failed: %s", e)
raise self.retry(exc=e, countdown=30)
async def _async_evaluate_alerts():
"""Async body of alert evaluation."""
engine = create_async_engine(settings.DATABASE_URL)
factory = async_sessionmaker(engine, expire_on_commit=False)
try:
async with factory() as db:
service = AlertService(db)
triggered = await service.evaluate_rules()
await db.commit()
return triggered
finally:
await engine.dispose()
@celery_app.task
def cleanup_old_indices():
"""Delete log entries older than the retention period."""
try:
loop = asyncio.get_event_loop()
if loop.is_closed():
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
result = loop.run_until_complete(_async_cleanup())
logger.info("Cleanup complete: %d entries deleted", result)
return {"deleted": result}
except Exception as e:
logger.error("Cleanup failed: %s", e)
return {"error": str(e)}
async def _async_cleanup() -> int:
"""Delete old log entries."""
engine = create_async_engine(settings.DATABASE_URL)
factory = async_sessionmaker(engine, expire_on_commit=False)
try:
cutoff = datetime.now(timezone.utc) - timedelta(days=settings.LOG_RETENTION_DAYS)
async with factory() as db:
result = await db.execute(delete(LogEntry).where(LogEntry.timestamp < cutoff))
await db.commit()
return result.rowcount
finally:
await engine.dispose()