from __future__ import annotations
import hashlib, json, os, socket, subprocess, tempfile
from datetime import datetime, timezone, timedelta
from pathlib import Path
from croniter import croniter
from sqlalchemy import select, func
from sqlalchemy.engine import make_url
from sqlalchemy.orm import Session
from app.core.config import settings
from app.core.database import SessionLocal
from app.core.security import encrypt_secret
from app.models.entities import ClusterNode, ClusterLease, FailoverEvent, BackupPolicy, BackupRun, DisasterRecoveryPlan, DisasterRecoveryDrill
from app.services.storage import ObjectStorage

def now(): return datetime.now(timezone.utc)

def register_or_heartbeat(db:Session,node_name:str|None=None,node_type:str|None=None,address:str|None=None,capabilities:list|None=None):
    name=node_name or settings.cluster_node_name or socket.gethostname()
    ntype=node_type or settings.cluster_node_type
    node=db.scalar(select(ClusterNode).where(ClusterNode.node_name==name,ClusterNode.node_type==ntype))
    if not node:
        node=ClusterNode(node_name=name,node_type=ntype,region=settings.cluster_region,zone=settings.cluster_zone,address=address or settings.cluster_advertise_address,version='0.14.0',status='healthy',capabilities=capabilities or [])
        db.add(node)
    node.status='healthy'; node.last_heartbeat_at=now(); node.address=address or node.address; node.capabilities=capabilities or node.capabilities
    db.commit(); db.refresh(node); return node

def mark_stale_nodes(db:Session):
    cutoff=now()-timedelta(seconds=settings.cluster_node_timeout_seconds)
    stale=db.scalars(select(ClusterNode).where(ClusterNode.last_heartbeat_at<cutoff,ClusterNode.status!='offline')).all()
    for node in stale: node.status='offline'; node.role='follower'
    db.commit(); return len(stale)

def acquire_lease(db:Session,name:str,node:ClusterNode,ttl_seconds:int|None=None):
    ttl=ttl_seconds or settings.cluster_lease_seconds
    lease=db.get(ClusterLease,name)
    current=now()
    if not lease:
        lease=ClusterLease(name=name,holder_node_id=node.id,fencing_token=1,acquired_at=current,expires_at=current+timedelta(seconds=ttl)); db.add(lease); won=True
    elif lease.holder_node_id==node.id or not lease.expires_at or lease.expires_at<=current:
        previous=lease.holder_node_id
        lease.holder_node_id=node.id; lease.fencing_token+=1; lease.acquired_at=current; lease.expires_at=current+timedelta(seconds=ttl); won=True
        if previous and previous!=node.id:
            db.add(FailoverEvent(component=name,event_type='lease_failover',from_node_id=previous,to_node_id=node.id,reason='Lease expired or previous holder unavailable',evidence={'fencing_token':lease.fencing_token}))
    else: won=False
    if won:
        db.query(ClusterNode).filter(ClusterNode.node_type==node.node_type).update({'role':'follower'})
        node.role='leader'
    db.commit(); return won,lease

def cluster_summary(db:Session):
    mark_stale_nodes(db)
    nodes=db.scalars(select(ClusterNode).order_by(ClusterNode.node_type,ClusterNode.node_name)).all()
    leases=db.scalars(select(ClusterLease)).all()
    return {'total_nodes':len(nodes),'healthy_nodes':sum(1 for n in nodes if n.status=='healthy'),'offline_nodes':sum(1 for n in nodes if n.status=='offline'),'leaders':sum(1 for n in nodes if n.role=='leader'),'nodes':nodes,'leases':leases}

def next_cron(expr:str): return croniter(expr,now()).get_next(datetime)

def backup_policy_due(policy:BackupPolicy): return policy.enabled and (not policy.next_run_at or policy.next_run_at<=now())

def _safe_pg_dump(target:Path):
    # Build arguments without exposing the password in the process list.
    url=make_url(settings.metadata_database_url)
    if not url.drivername.startswith('postgresql'):
        raise ValueError('Backup lógico automático requer PostgreSQL')
    env=os.environ.copy()
    if url.password: env['PGPASSWORD']=url.password
    cmd=['pg_dump','--format=custom','--no-owner','--no-acl','--file',str(target)]
    if url.host: cmd += ['--host',url.host]
    if url.port: cmd += ['--port',str(url.port)]
    if url.username: cmd += ['--username',url.username]
    cmd += [url.database or 'datalake']
    subprocess.run(cmd,check=True,timeout=3600,env=env,capture_output=True)

def execute_backup(db:Session,policy:BackupPolicy,triggered_by:str|None=None):
    run=BackupRun(policy_id=policy.id,status='running',triggered_by=triggered_by,started_at=now()); db.add(run); db.commit(); db.refresh(run)
    work=Path(settings.backup_work_dir); work.mkdir(parents=True,exist_ok=True)
    raw=work/f'{run.id}.dump'; enc=work/f'{run.id}.dump.enc'
    try:
        if policy.resource_type!='metadata_database': raise ValueError('A V014 executa backup automático apenas do metadata_database')
        _safe_pg_dump(raw)
        data=raw.read_bytes()
        if len(data)>settings.backup_max_bytes: raise ValueError('Backup excede o limite configurado')
        encrypted=encrypt_secret(data.hex()).encode('utf-8') if policy.encryption_required else data
        enc.write_bytes(encrypted)
        checksum=hashlib.sha256(encrypted).hexdigest()
        key=f"{policy.destination_prefix.rstrip('/')}/{now().strftime('%Y/%m/%d')}/{run.id}.dump.enc"
        storage=ObjectStorage(); storage.upload(str(enc),key)
        run.status='success'; run.object_key=key; run.checksum=checksum; run.size_bytes=len(encrypted); run.encryption='fernet-hex' if policy.encryption_required else 'none'
        policy.last_run_at=now(); policy.next_run_at=next_cron(policy.schedule_cron)
    except Exception as exc:
        run.status='failed'; run.error=str(exc)[:2000]
    finally:
        run.finished_at=now(); db.commit()
        for f in (raw,enc):
            try:f.unlink(missing_ok=True)
            except OSError:pass
    return run

def execute_due_backups():
    db=SessionLocal(); result=[]
    try:
        node=register_or_heartbeat(db,node_type='scheduler',capabilities=['backup_scheduler'])
        won,_=acquire_lease(db,'backup-scheduler',node)
        if not won:return {'leader':False,'executed':[]}
        for p in db.scalars(select(BackupPolicy).where(BackupPolicy.enabled.is_(True))).all():
            if backup_policy_due(p): result.append(execute_backup(db,p).id)
        return {'leader':True,'executed':result}
    finally: db.close()
