""" Celery async tasks para scrapers. Usadas para procesamiento background de informes pesados. En Fase 1 los scrapers corren inline en el endpoint (más simple). En Fase 2 estos tasks permiten colas y reintentos automáticos. """ import asyncio import logging from app.tasks.celery_app import celery_app from app.scrapers.arca_afip import ArcaAfipScraper from app.scrapers.bcra import BcraScraper from app.scrapers.boletin_oficial import BoletinOficialScraper logger = logging.getLogger(__name__) def run_async(coro): """Helper para correr corutinas async dentro de tasks síncronas de Celery.""" loop = asyncio.new_event_loop() try: return loop.run_until_complete(coro) finally: loop.close() @celery_app.task(bind=True, max_retries=3, default_retry_delay=5, name="tasks.fetch_arca") def fetch_arca_task(self, cuit: str) -> dict: try: scraper = ArcaAfipScraper() return run_async(scraper.safe_fetch(cuit)) except Exception as exc: logger.error(f"ARCA task error for {cuit}: {exc}") self.retry(exc=exc) @celery_app.task(bind=True, max_retries=3, default_retry_delay=5, name="tasks.fetch_bcra") def fetch_bcra_task(self, cuil: str) -> dict: try: scraper = BcraScraper() return run_async(scraper.safe_fetch(cuil)) except Exception as exc: logger.error(f"BCRA task error for {cuil}: {exc}") self.retry(exc=exc) @celery_app.task(bind=True, max_retries=2, default_retry_delay=10, name="tasks.fetch_boletin") def fetch_boletin_task(self, query: str) -> dict: try: scraper = BoletinOficialScraper() return run_async(scraper.safe_fetch(query)) except Exception as exc: logger.error(f"BO task error for {query}: {exc}") self.retry(exc=exc)