import asyncio import sys from collections import defaultdict from typing import Any from fastapi import BackgroundTasks from sqlalchemy import select, Result from sqlalchemy.ext.asyncio import AsyncSession sys.path.insert(0, "/app/src") from utils.flight_plan_helpers import refresh_markers_weather_info from database.models import FlightPlan from database import models from database.transaction import get_session async def get_plans(db: AsyncSession) -> list[Any] | Result[tuple[FlightPlan, Any]]: markers_without_weather = (await db.scalars( select(models.FlightPlanMarker).filter(models.FlightPlanMarker.weather_info_id.is_(None)) )) flight_plan_ids = {m.flight_plan_id for m in markers_without_weather} flight_plan_ids.add(11) if not flight_plan_ids: return [] return (await db.execute( select(models.FlightPlan, models.FlightPlanMarker) .join(models.FlightPlan.markers) .filter(models.FlightPlan.planned_takeoff_datetime.is_not(None)) .filter(models.FlightPlan.id.in_(flight_plan_ids)) )) async def download_flight_plan_weather(): tasks = BackgroundTasks() async with get_session() as db: plans_with_markers_without_weather = await get_plans(db) markers_by_plan = defaultdict(list) for plan, marker in plans_with_markers_without_weather: markers_by_plan[plan].append(marker) for plan, markers in markers_by_plan.items(): await refresh_markers_weather_info( plan.planned_takeoff_datetime, plan.planned_speed, markers, background_tasks=tasks ) await tasks() async def run_all(): await asyncio.gather(download_flight_plan_weather()) if __name__ == "__main__": loop = asyncio.get_event_loop() loop.run_until_complete(run_all())