/
gr.ev.vl
/
TestGen
Обзор
Документация
Войти
/
gr.ev.vl
/
TestGen
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master2
backend/src/testgen/tasks/analytics.py
119 строк
4 KB
gr.ev.vl
Initial commit
06 июн 2026, 02:41
06 июн 2026, 02:41
02fb1fb
Код
Авторство
О чём код?
""" Периодические задачи аналитики. """ import asyncio import logging from celery.schedules import crontab from ..domain.repositories.operator_decision import OperatorDecisionRepository from ..domain.repositories.user import UserRepository from ..infrastructure.database.session import async_session_factory from .celery_app import celery_app logger = logging.getLogger(__name__) def _run_async(coro): """Запускает асинхронную корутину в синхронном контексте Celery.""" loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) try: return loop.run_until_complete(coro) finally: loop.close() @celery_app.task( bind=True, queue="llm_low_priority", ) def update_operator_statistics(self) -> dict: """ Задача: обновление агрегированной статистики операторов. Выполняется ежедневно. """ async def _run(): session = async_session_factory() try: user_repo = UserRepository(session) decision_repo = OperatorDecisionRepository(session) users = await user_repo.list(limit=1000) stats = {} for user in users: decisions = await decision_repo.get_by_operator(user.id) total = len(decisions) accepted = sum(1 for d in decisions if "accept" in d.decision_type) rejected = sum(1 for d in decisions if "reject" in d.decision_type) edited = sum(1 for d in decisions if "edit" in d.decision_type) stats[str(user.id)] = { "total_decisions": total, "accepted": accepted, "rejected": rejected, "edited": edited, "acceptance_rate": (accepted / total * 100) if total > 0 else 0, } return stats finally: await session.close() try: result = _run_async(_run()) logger.info(f"Operator statistics updated for {len(result)} users") return {"status": "completed", "users_processed": len(result)} except Exception as exc: logger.error(f"Failed to update operator statistics: {exc}") raise @celery_app.task( bind=True, queue="llm_low_priority", ) def cleanup_old_data(self, days: int = 90) -> dict: """ Задача: очистка устаревших данных. Args: days: Возраст данных в днях для удаления """ async def _run(): session = async_session_factory() try: # Здесь будет логика очистки старых записей logger.info(f"Cleaning data older than {days} days") return {"deleted_records": 0} finally: await session.close() try: result = _run_async(_run()) logger.info("Data cleanup completed") return result except Exception as exc: logger.error(f"Data cleanup failed: {exc}") raise # Настройка периодических задач celery_app.conf.beat_schedule = { "update-operator-statistics": { "task": "src.testgen.tasks.analytics.update_operator_statistics", "schedule": crontab(hour=2, minute=0), # Каждый день в 02:00 }, "cleanup-old-data": { "task": "src.testgen.tasks.analytics.cleanup_old_data", "schedule": crontab( hour=3, minute=0, day_of_week=1 ), # Каждый понедельник в 03:00 "kwargs": {"days": 90}, }, }