import uuid from typing import Annotated, Any from fastapi import APIRouter, Depends, Query from sqlalchemy import exists from resonance_engine.adapters.db import db from resonance_engine.api import schemas, security from resonance_engine.tasks.models import CollectTask, FanoutTask router = APIRouter( prefix="/tasks", tags=["Tasks"], dependencies=[Depends(security.authenticate)], ) @router.get( "/stats", response_model=schemas.TaskStats, ) def get_task_stats( params: Annotated[schemas.GetTaskStatsInput, Query()], ) -> Any: with db.autocommit(): fanout_agg = FanoutTask.query.aggregate_stats() collect_agg = CollectTask.query.aggregate_stats( dsp_client_name=params.dsp_client_name ) return schemas.TaskStats( fanout=schemas.FanoutRunStats( total=fanout_agg.total, running=fanout_agg.running, last_started_at=fanout_agg.last_started_at, last_fans_dispatched=fanout_agg.last_fans_dispatched, ), collect=schemas.CollectRunStats( total=collect_agg.total, running=collect_agg.running, stale=collect_agg.stale, last_started_at=collect_agg.last_started_at, ), ) @router.get( "/activity", response_model=list[schemas.TaskActivityBucket], ) def get_task_activity( params: Annotated[schemas.GetTaskActivityInput, Query()], ) -> Any: with db.autocommit(): return CollectTask.query.activity( days=params.days, granularity=params.granularity, dsp_client_name=params.dsp_client_name, ) @router.get( "/fanout", response_model=schemas.FanoutTaskPaginated, ) def list_fanout_tasks( params: Annotated[schemas.GetFanoutTasksInput, Query()], ) -> Any: with db.autocommit(): q = FanoutTask.query if params.status is not None: q = q.where(FanoutTask.status == params.status) if params.source is not None: q = q.where(FanoutTask.source == params.source) if params.dsp_client_name is not None: q = q.where( exists().where( CollectTask.fanout_task_id == FanoutTask.id, CollectTask.dsp_client_name == params.dsp_client_name, ) ) return q.paginate(cursor=params.cursor, limit=params.limit) @router.get( "/fanout/{task_id}", response_model=schemas.FanoutTask, ) def get_fanout_task(task_id: uuid.UUID) -> Any: with db.autocommit(): return FanoutTask.query.get_one(task_id) @router.get( "/fanout/{task_id}/collect", response_model=list[schemas.CollectTask], ) def list_collect_tasks(task_id: uuid.UUID) -> Any: with db.autocommit(): return ( CollectTask.query.where(CollectTask.fanout_task_id == task_id) .order_by(CollectTask.started_at) .all() )