"""Google Sheet → FormData import Taskiq tasks (dedicated sheet_import stream).""" from __future__ import annotations import logging import os from datetime import datetime,timezone import redis.asyncio as redis from dotenv import load_dotenv from db_setup import session_scope from g_sheet.models import SheetImportRun from g_sheet.views import SheetImport from taskiq_management.broker_setup import MAX_RETRIES,RETRY_DELAY from taskiq_management.g_sheet_broker_setup import sheet_broker from taskiq_management.middleware import PermanentTaskError load_dotenv() logger=logging.getLogger("g_sheet.tasks") REDIS_URL=os.getenv("REDIS_URL","redis://localhost:6379/0") _LOCK_KEY="g_sheet:import:lock" _LOCK_TTL=3600 async def _fail(run_id:str,error:str) -> dict: async with session_scope() as session: await SheetImportRun.update_run(session,run_id,{ "status":"failed", "error":error, "finished_at":datetime.now(timezone.utc), }) return {"status":"failed","error":error} @sheet_broker.task( task_name="g_sheet.import_sheets", retry_on_error=True, max_retries=MAX_RETRIES, delay=RETRY_DELAY, ) async def import_sheets(run_id:str) -> dict: if not run_id or not str(run_id).strip(): raise PermanentTaskError("run_id is required") run_id=str(run_id).strip() client=redis.from_url(REDIS_URL,decode_responses=True) try: acquired=await client.set(_LOCK_KEY,run_id,nx=True,ex=_LOCK_TTL) if not acquired: return await _fail(run_id,"another sheet import is already running") try: async with session_scope() as session: row=await SheetImportRun.get_by_id(session,run_id) if not row: raise PermanentTaskError(f"import run {run_id} not found") await SheetImportRun.update_run(session,run_id,{ "status":"running", "started_at":datetime.now(timezone.utc), "error":None, }) tab=row.tab async with session_scope() as session: service=SheetImport(session=session) try: if tab: report=await service.import_sheet(tab) else: report=await service.import_all() except Exception as e: logger.exception("sheet import failed for run %s",run_id) # Bad tab names and permanent Sheets 4xx — do not burn retries. from fastapi import HTTPException if isinstance(e,HTTPException) and e.status_code in (400,404,422): await _fail(run_id,str(e.detail)) raise PermanentTaskError(str(e.detail)) from e return await _fail(run_id,str(e)) await SheetImportRun.update_run(session,run_id,{ "status":"completed", "report":report, "error":None, "finished_at":datetime.now(timezone.utc), }) return {"status":"completed","report":report} finally: current=await client.get(_LOCK_KEY) if current==run_id: await client.delete(_LOCK_KEY) finally: await client.aclose()