import base64, hashlib, json
from datetime import datetime, timezone, timedelta
from sqlalchemy import select, func
from app.models.entities import FederatedQueryJob,FederatedJoinDefinition,QueryQuota,QueryUsageLedger,SemanticModel
from app.services.semantic_execution import execute_semantic

def role_value(user): return user.role.value if hasattr(user.role,'value') else str(user.role)
def _quota(db,user):
    return db.scalar(select(QueryQuota).where(QueryQuota.enabled.is_(True),QueryQuota.subject_type=='user',QueryQuota.subject_value==user.id)) or db.scalar(select(QueryQuota).where(QueryQuota.enabled.is_(True),QueryQuota.subject_type=='role',QueryQuota.subject_value==role_value(user)))
def enforce_quota(db,user,estimated_rows,estimated_cost=0):
    q=_quota(db,user); max_concurrent=q.max_concurrent if q else 3; row_limit=q.daily_row_limit if q else 1_000_000; cost_limit=q.daily_cost_limit if q else 1_000_000
    running=db.scalar(select(func.count()).select_from(FederatedQueryJob).where(FederatedQueryJob.user_id==user.id,FederatedQueryJob.status.in_(['queued','running']))) or 0
    since=datetime.now(timezone.utc)-timedelta(days=1)
    rows=db.scalar(select(func.coalesce(func.sum(QueryUsageLedger.rows_processed),0)).where(QueryUsageLedger.user_id==user.id,QueryUsageLedger.created_at>=since)) or 0
    cost=db.scalar(select(func.coalesce(func.sum(QueryUsageLedger.cost_units),0)).where(QueryUsageLedger.user_id==user.id,QueryUsageLedger.created_at>=since)) or 0
    if running>=max_concurrent: raise PermissionError('Limite de consultas concorrentes atingido')
    if rows+estimated_rows>row_limit: raise PermissionError('Cota diária de linhas excedida')
    if cost+estimated_cost>cost_limit: raise PermissionError('Cota diária de custo excedida')
    return q

def cursor_for(job_id,offset):
    raw=json.dumps({'j':job_id,'o':offset},separators=(',',':')).encode(); return base64.urlsafe_b64encode(raw).decode().rstrip('=')
def decode_cursor(token):
    try:
        raw=base64.urlsafe_b64decode(token+'='*(-len(token)%4)); data=json.loads(raw); return str(data['j']),int(data['o'])
    except Exception as exc: raise ValueError('Cursor inválido') from exc

def create_job(db,user,request):
    limit=int(request.get('page_size',500)); enforce_quota(db,user,limit,limit)
    plan={'strategy':'single_source' if not request.get('join_id') else 'application_hash_join','pushdown':True,'page_size':limit}
    if request.get('join_id'):
        join=db.get(FederatedJoinDefinition,request['join_id'])
        if not join or not join.enabled: raise ValueError('Definição de join não encontrada')
        plan.update({'join_id':join.id,'join_type':join.join_type,'cardinality':join.cardinality,'left_key':join.left_key,'right_key':join.right_key})
    job=FederatedQueryJob(user_id=user.id,request_json=request,execution_plan=plan,status='queued',cursor_token='')
    db.add(job); db.commit(); db.refresh(job); job.cursor_token=cursor_for(job.id,0); db.commit(); db.refresh(job); return job

def cancel_job(db,user,job):
    if job.user_id!=user.id and role_value(user) not in {'admin','director','data_engineer'}: raise PermissionError('Sem permissão para cancelar esta consulta')
    if job.status in {'success','failed','cancelled'}: return job
    job.cancel_requested=True; job.status='cancelled'; job.finished_at=datetime.now(timezone.utc); db.commit(); db.refresh(job); return job

def execute_job(db,user,job):
    if job.cancel_requested: return cancel_job(db,user,job)
    req=job.request_json; job.status='running'; job.started_at=datetime.now(timezone.utc); job.progress_percent=10; db.commit()
    if req.get('join_id'):
        raise ValueError('Execução física de join multi-fonte requer worker assíncrono configurado; o plano governado foi criado')
    model=db.get(SemanticModel,req.get('model_id'))
    if not model: raise ValueError('Modelo semântico não encontrado')
    result=execute_semantic(db,user,model,req.get('metrics',[]),req.get('dimensions',[]),req.get('filters',{}),int(req.get('page_size',500)),req.get('use_cache',True))
    if job.cancel_requested: return cancel_job(db,user,job)
    payload={'columns':result.get('columns',[]),'rows':result.get('rows',[])}
    job.execution_plan={**job.execution_plan,'result':payload}; job.row_count=result.get('row_count',0); job.progress_percent=100; job.status='success'; job.finished_at=datetime.now(timezone.utc)
    db.add(QueryUsageLedger(user_id=user.id,job_id=job.id,rows_processed=job.row_count,cost_units=max(1,job.row_count),cache_hit=result.get('cache_hit',False)))
    db.commit(); db.refresh(job); return job

def page_job(job,offset,page_size):
    result=(job.execution_plan or {}).get('result',{}); rows=result.get('rows',[]); end=min(offset+page_size,len(rows)); next_cursor=cursor_for(job.id,end) if end<len(rows) else None
    return {'job_id':job.id,'status':job.status,'columns':result.get('columns',[]),'rows':rows[offset:end],'offset':offset,'page_size':page_size,'next_cursor':next_cursor,'total_rows':job.row_count}
