import ast,hashlib,json,re
from pathlib import Path
from datetime import datetime,timezone
import pandas as pd
from sqlalchemy import select
from sqlalchemy.orm import Session
from app.models.entities import ETLPipeline,ETLPipelineRun,ETLNodeRun,TenantConnectorInstance,ConnectorDiscoveredObject
from app.services.connectivity_hub import reveal_config

SAFE_NAME=re.compile(r'^[A-Za-z_][A-Za-z0-9_]*$')
ALLOWED_NODES={'source','select','rename','filter','derive','deduplicate','sort','validate','destination'}
ALLOWED_AST=(ast.Expression,ast.BinOp,ast.UnaryOp,ast.BoolOp,ast.Compare,ast.Name,ast.Load,ast.Constant,ast.Add,ast.Sub,ast.Mult,ast.Div,ast.Mod,ast.Pow,ast.And,ast.Or,ast.Not,ast.Eq,ast.NotEq,ast.Gt,ast.GtE,ast.Lt,ast.LtE,ast.USub,ast.UAdd)

def validate_pipeline(nodes:list,edges:list):
    if not nodes: raise ValueError('O pipeline precisa de ao menos um nó')
    ids=[str(n.get('id','')) for n in nodes]
    if len(ids)!=len(set(ids)) or any(not x for x in ids): raise ValueError('IDs de nós inválidos ou duplicados')
    for n in nodes:
        if n.get('type') not in ALLOWED_NODES: raise ValueError(f"Tipo de nó não permitido: {n.get('type')}")
    graph={i:[] for i in ids}; indeg={i:0 for i in ids}
    for e in edges:
        a,b=str(e.get('source','')),str(e.get('target',''))
        if a not in graph or b not in graph: raise ValueError('Conexão aponta para nó inexistente')
        graph[a].append(b); indeg[b]+=1
    queue=[i for i,v in indeg.items() if v==0]; order=[]
    while queue:
        x=queue.pop(0); order.append(x)
        for y in graph[x]:
            indeg[y]-=1
            if indeg[y]==0: queue.append(y)
    if len(order)!=len(ids): raise ValueError('O pipeline contém ciclo')
    if not any(n.get('type')=='source' for n in nodes): raise ValueError('Inclua um nó de origem')
    if not any(n.get('type')=='destination' for n in nodes): raise ValueError('Inclua um nó de destino')
    return order

def _safe_expr(expr:str,columns:list[str]):
    tree=ast.parse(expr,mode='eval')
    for node in ast.walk(tree):
        if not isinstance(node,ALLOWED_AST): raise ValueError('Expressão contém operação não permitida')
        if isinstance(node,ast.Name) and node.id not in columns: raise ValueError(f'Coluna desconhecida: {node.id}')
    return compile(tree,'<etl-expression>','eval')

def _load_source(db:Session,tenant_id:str,cfg:dict,limit:int=100000):
    instance=db.get(TenantConnectorInstance,cfg.get('instance_id'))
    if not instance or instance.tenant_id!=tenant_id: raise ValueError('Conexão de origem não encontrada')
    obj=db.get(ConnectorDiscoveredObject,cfg.get('object_id'))
    if not obj or obj.tenant_id!=tenant_id or obj.connector_instance_id!=instance.id: raise ValueError('Objeto de origem inválido')
    config=reveal_config(instance); code=instance.connector_code
    if code in {'postgresql','mysql','sqlserver','oracle'}:
        from app.connectors.sqlalchemy_connector import SQLAlchemyConnector
        c=SQLAlchemyConnector(code,config.get('host',''),int(config.get('port') or 0),config.get('database',''),config.get('username',''),config.get('password',''),config.get('options',{}))
        preview=c.preview(obj.qualified_name,limit=min(limit,5000))
        return pd.DataFrame(preview['rows'],columns=preview['columns'])
    if code=='csv': return pd.read_csv(config['path'],nrows=limit)
    if code=='excel': return pd.read_excel(config['path'],sheet_name=obj.object_name,nrows=limit)
    if code=='json': return pd.read_json(config['path'],lines=config.get('lines',False)).head(limit)
    if code=='parquet': return pd.read_parquet(config['path']).head(limit)
    raise ValueError(f'Execução de origem ainda não disponível para {code}')

def _transform(df:pd.DataFrame,node:dict):
    cfg=node.get('config') or {}; kind=node['type']
    if kind=='select':
        cols=cfg.get('columns') or []
        missing=[c for c in cols if c not in df.columns]
        if missing: raise ValueError(f'Colunas ausentes: {missing}')
        return df[cols].copy()
    if kind=='rename': return df.rename(columns=cfg.get('mapping') or {})
    if kind=='filter':
        expr=str(cfg.get('expression','')).strip(); code=_safe_expr(expr,list(df.columns))
        mask=df.apply(lambda r: bool(eval(code,{'__builtins__':{}},r.to_dict())),axis=1)
        return df[mask].copy()
    if kind=='derive':
        name=cfg.get('name'); expr=str(cfg.get('expression',''))
        if not name or not SAFE_NAME.fullmatch(name): raise ValueError('Nome do campo calculado inválido')
        code=_safe_expr(expr,list(df.columns)); out=df.copy(); out[name]=out.apply(lambda r: eval(code,{'__builtins__':{}},r.to_dict()),axis=1); return out
    if kind=='deduplicate': return df.drop_duplicates(subset=cfg.get('columns') or None,keep=cfg.get('keep','first'))
    if kind=='sort': return df.sort_values(by=cfg.get('columns') or [],ascending=cfg.get('ascending',True))
    if kind=='validate':
        required=cfg.get('required') or []; missing=[c for c in required if c not in df.columns]
        if missing: raise ValueError(f'Campos obrigatórios não existem: {missing}')
        invalid=df[required].isna().any(axis=1) if required else pd.Series(False,index=df.index)
        if invalid.any() and cfg.get('on_error','fail')=='fail': raise ValueError(f'{int(invalid.sum())} linhas falharam na validação')
        return df[~invalid].copy() if cfg.get('on_error')=='drop' else df
    return df

def execute_pipeline(db:Session,pipeline:ETLPipeline,run:ETLPipelineRun,debug:bool=False):
    order=validate_pipeline(pipeline.nodes,pipeline.edges); nodes={str(n['id']):n for n in pipeline.nodes}
    run.status='running'; run.started_at=datetime.now(timezone.utc); db.commit(); data=None
    try:
        for idx,node_id in enumerate(order):
            node=nodes[node_id]; nr=ETLNodeRun(tenant_id=pipeline.tenant_id,pipeline_run_id=run.id,node_id=node_id,node_type=node['type'],status='running',started_at=datetime.now(timezone.utc)); db.add(nr); db.commit()
            before=0 if data is None else len(data)
            if node['type']=='source': data=_load_source(db,pipeline.tenant_id,node.get('config') or {},limit=500 if debug else 100000)
            elif node['type']=='destination':
                if data is None: raise ValueError('Nenhum dado disponível para destino')
                root=Path('/tmp/cedp-etl')/pipeline.tenant_id; root.mkdir(parents=True,exist_ok=True)
                path=root/f'{run.id}.parquet'; data.to_parquet(path,index=False)
                run.output_key=str(path); run.checksum=hashlib.sha256(path.read_bytes()).hexdigest()
            else:
                if data is None: raise ValueError('Transformação sem dados de origem')
                data=_transform(data,node)
            nr.rows_in=before; nr.rows_out=0 if data is None else len(data); nr.sample=[] if data is None else json.loads(data.head(10).to_json(orient='records',date_format='iso')); nr.status='success'; nr.finished_at=datetime.now(timezone.utc)
            run.progress=int(((idx+1)/len(order))*100); db.commit()
        run.rows_in=db.scalar(select(ETLNodeRun.rows_out).where(ETLNodeRun.pipeline_run_id==run.id,ETLNodeRun.node_type=='source')) or 0
        run.rows_out=0 if data is None else len(data); run.status='success'; run.finished_at=datetime.now(timezone.utc); pipeline.last_run_at=run.finished_at; db.commit(); return run
    except Exception as exc:
        run.status='failed'; run.error=str(exc)[:3000]; run.finished_at=datetime.now(timezone.utc); db.commit(); raise
