from app.core.database import SessionLocal
from app.services.ingestion import run_ingestion
from app.workers.celery_app import celery_app


@celery_app.task(name="ingest_dataset", autoretry_for=(Exception,), retry_backoff=True, max_retries=3)
def ingest_dataset_task(dataset_id: str):
    db = SessionLocal()
    try:
        return run_ingestion(db, dataset_id)
    finally:
        db.close()


@celery_app.task(name="scan_dataset_schedules")
def scan_dataset_schedules():
    from datetime import datetime
    from zoneinfo import ZoneInfo
    from croniter import croniter
    from sqlalchemy import select
    from app.models.entities import Dataset

    db = SessionLocal()
    try:
        now = datetime.now(ZoneInfo("America/Sao_Paulo")).replace(second=0, microsecond=0)
        datasets = db.scalars(select(Dataset).where(Dataset.enabled.is_(True), Dataset.schedule_cron.is_not(None))).all()
        queued = []
        for dataset in datasets:
            try:
                if croniter.match(dataset.schedule_cron, now):
                    ingest_dataset_task.delay(dataset.id)
                    queued.append(dataset.id)
            except (ValueError, KeyError):
                continue
        return {"queued": queued}
    finally:
        db.close()

@celery_app.task(name='run_automations')
def run_automations_task():
    from app.services.automation import execute_due_automations
    return execute_due_automations()

@celery_app.task(name='run_report_schedules')
def run_report_schedules_task():
    from app.services.report_delivery import execute_due_report_schedules
    return execute_due_report_schedules()

@celery_app.task(name='evaluate_alert_rules')
def evaluate_alert_rules_task():
    from app.services.alerts import evaluate_alerts
    return evaluate_alerts()

@celery_app.task(name='advance_workflow_timers')
def advance_workflow_timers_task():
    from datetime import datetime, timezone
    from sqlalchemy import select
    from app.models.entities import WorkflowTask, WorkflowInstance, WorkflowDefinition
    from app.services.workflow import advance
    db=SessionLocal(); count=0
    try:
        tasks=db.scalars(select(WorkflowTask).where(WorkflowTask.status=='waiting',WorkflowTask.due_at<=datetime.now(timezone.utc))).all()
        for task in tasks:
            inst=db.get(WorkflowInstance,task.instance_id); wf=db.get(WorkflowDefinition,inst.workflow_id)
            edges=[e for e in (wf.definition or {}).get('edges',[]) if e.get('source')==task.node_id]
            task.status='completed'; task.completed_at=datetime.now(timezone.utc)
            if edges: inst.current_node_id=edges[0]['target']; advance(db,inst,None)
            count+=1
        db.commit(); return {'advanced':count}
    finally: db.close()

@celery_app.task(name='dispatch-mobile-push')
def dispatch_mobile_push_task():
    from app.services.mobile_push import dispatch_pending_push
    return dispatch_pending_push()

@celery_app.task(name='run-quality-schedules')
def run_quality_schedules_task():
    from datetime import datetime
    from zoneinfo import ZoneInfo
    from croniter import croniter
    from sqlalchemy import select
    from app.models.entities import Dataset, EnterpriseQualityRule, User
    from app.services.data_quality import execute_quality_run
    db=SessionLocal(); executed=[]
    try:
        current=datetime.now(ZoneInfo('America/Sao_Paulo')).replace(second=0,microsecond=0)
        rules=db.scalars(select(EnterpriseQualityRule).where(EnterpriseQualityRule.enabled.is_(True),EnterpriseQualityRule.schedule_cron.is_not(None))).all()
        due={}
        for rule in rules:
            try:
                if croniter.match(rule.schedule_cron,current): due.setdefault(rule.dataset_id,rule.owner_id)
            except (ValueError,KeyError): continue
        for dataset_id,owner_id in due.items():
            dataset=db.get(Dataset,dataset_id); owner=db.get(User,owner_id)
            if dataset and owner:
                run=execute_quality_run(db,dataset,owner,'schedule'); executed.append(run.id)
        return {'executed':executed}
    finally: db.close()


@celery_app.task(name='collect-platform-health')
def collect_platform_health_task():
    from app.services.observability import collect_health
    return {'checks': collect_health()}

@celery_app.task(name='evaluate-observability-policies')
def evaluate_observability_policies_task():
    from app.services.observability import evaluate_policies
    return evaluate_policies()


@celery_app.task(name='cluster-heartbeat')
def cluster_heartbeat_task():
    from app.services.cluster import register_or_heartbeat, mark_stale_nodes
    db=SessionLocal()
    try:
        node=register_or_heartbeat(db,node_type='worker',capabilities=['celery','ingestion','workflow'])
        stale=mark_stale_nodes(db)
        return {'node_id':node.id,'stale_nodes':stale}
    finally: db.close()

@celery_app.task(name='run-backup-schedules')
def run_backup_schedules_task():
    from app.services.cluster import execute_due_backups
    return execute_due_backups()

@celery_app.task(name='execute-federated-job-v018',bind=True,acks_late=True)
def execute_federated_job_v018_task(self,job_id,user_id,memory_limit_mb=256,worker_name='federated-worker'):
    from app.models.entities import FederatedQueryJob,User
    from app.services.federated_worker import execute_federated_job
    db=SessionLocal()
    try:
        job=db.get(FederatedQueryJob,job_id); user=db.get(User,user_id)
        if not job or not user: raise ValueError('Job ou usuário não encontrado')
        return {'job_id':execute_federated_job(db,user,job,memory_limit_mb,worker_name).id}
    finally: db.close()

@celery_app.task(name='cleanup-federated-artifacts-v018')
def cleanup_federated_artifacts_v018_task():
    from datetime import datetime,timezone
    from pathlib import Path
    from sqlalchemy import select
    from app.models.entities import QuerySpillSegment,QueryExport
    db=SessionLocal(); removed=0
    try:
        now=datetime.now(timezone.utc)
        for obj in list(db.scalars(select(QuerySpillSegment).where(QuerySpillSegment.expires_at<now)).all())+list(db.scalars(select(QueryExport).where(QueryExport.expires_at<now)).all()):
            path=Path(obj.object_key)
            if path.is_file(): path.unlink(missing_ok=True); removed+=1
            db.delete(obj)
        db.commit(); return {'removed_files':removed}
    finally: db.close()

@celery_app.task(name='execute-data-process-flow-v024',bind=True,acks_late=True)
def execute_data_process_flow_v024_task(self,flow_id:str,run_id:str):
    from datetime import datetime,timezone
    from app.models.entities import DataProcessFlow,DataProcessRun
    db=SessionLocal()
    try:
        flow=db.get(DataProcessFlow,flow_id); run=db.get(DataProcessRun,run_id)
        if not flow or not run: raise ValueError('Fluxo ou execução não encontrado')
        run.status='running'; run.started_at=datetime.now(timezone.utc); db.commit()
        completed={}; rows=0
        for node in flow.nodes:
            nid=str(node.get('id')); ntype=node.get('type')
            completed[nid]='running'; run.node_status=completed; db.commit()
            if ntype in {'dataset','extract','ingestion'} and node.get('dataset_id'):
                result=run_ingestion(db,node['dataset_id']); rows+=int(result.get('row_count',0) if isinstance(result,dict) else 0)
            completed[nid]='success'; run.node_status=dict(completed); db.commit()
        run.rows_processed=rows; run.status='success'; run.finished_at=datetime.now(timezone.utc); flow.last_run_at=run.finished_at; db.commit()
        return {'run_id':run.id,'status':run.status,'rows_processed':rows}
    except Exception as exc:
        if run:
            run.status='failed'; run.error=str(exc)[:2000]; run.finished_at=datetime.now(timezone.utc); db.commit()
        raise
    finally: db.close()

@celery_app.task(name='scan-data-process-flows-v024')
def scan_data_process_flows_v024_task():
    from datetime import datetime,timezone
    from sqlalchemy import select
    from app.models.entities import DataProcessFlow,DataProcessRun
    db=SessionLocal(); queued=[]
    try:
        now=datetime.now(timezone.utc)
        flows=db.scalars(select(DataProcessFlow).where(DataProcessFlow.enabled.is_(True),DataProcessFlow.publish_status=='published',DataProcessFlow.next_run_at.is_not(None),DataProcessFlow.next_run_at<=now)).all()
        from croniter import croniter
        for flow in flows:
            run=DataProcessRun(tenant_id=flow.tenant_id,flow_id=flow.id,status='queued'); db.add(run); db.flush()
            flow.next_run_at=croniter(flow.schedule_cron,now).get_next(datetime) if flow.schedule_cron else None
            execute_data_process_flow_v024_task.delay(flow.id,run.id); queued.append(run.id)
        db.commit(); return {'queued':queued}
    finally: db.close()


@celery_app.task(name='execute-etl-pipeline-v026',bind=True,acks_late=True)
def execute_etl_pipeline_v026_task(self,pipeline_id:str,run_id:str):
    from app.models.entities import ETLPipeline,ETLPipelineRun
    from app.services.etl_studio import execute_pipeline
    db=SessionLocal()
    try:
        pipeline=db.get(ETLPipeline,pipeline_id); run=db.get(ETLPipelineRun,run_id)
        if not pipeline or not run: raise ValueError('Pipeline ou execução não encontrado')
        result=execute_pipeline(db,pipeline,run,False)
        return {'run_id':result.id,'status':result.status,'rows_out':result.rows_out}
    finally: db.close()

@celery_app.task(name='scan-etl-pipelines-v026')
def scan_etl_pipelines_v026_task():
    from datetime import datetime,timezone
    from croniter import croniter
    from sqlalchemy import select
    from app.models.entities import ETLPipeline,ETLPipelineRun
    db=SessionLocal(); queued=[]
    try:
        now=datetime.now(timezone.utc)
        rows=db.scalars(select(ETLPipeline).where(ETLPipeline.enabled.is_(True),ETLPipeline.status=='published',ETLPipeline.next_run_at.is_not(None),ETLPipeline.next_run_at<=now)).all()
        for p in rows:
            r=ETLPipelineRun(tenant_id=p.tenant_id,pipeline_id=p.id,requested_by=p.owner_id,mode='schedule'); db.add(r); db.flush()
            p.next_run_at=croniter(p.schedule_cron,now).get_next(datetime) if p.schedule_cron else None
            execute_etl_pipeline_v026_task.delay(p.id,r.id); queued.append(r.id)
        db.commit(); return {'queued':queued}
    finally: db.close()
