import hashlib
import os
import tempfile
from datetime import datetime, timezone
import pandas as pd
from sqlalchemy.orm import Session
from app.connectors.sqlalchemy_connector import SQLAlchemyConnector
from app.core.config import settings
from app.core.security import decrypt_secret
from app.models.entities import Dataset, IngestionRun, JobStatus
from app.services.storage import ObjectStorage


def _mask_dataframe(df: pd.DataFrame, policy: dict) -> pd.DataFrame:
    result = df.copy()
    for column, method in policy.items():
        if column not in result.columns:
            continue
        if method == "null":
            result[column] = None
        elif method == "hash":
            result[column] = result[column].astype(str).map(lambda v: hashlib.sha256(v.encode()).hexdigest())
        elif method == "partial":
            result[column] = result[column].astype(str).map(lambda v: (v[:2] + "***" + v[-2:]) if len(v) > 4 else "***")
    return result


def run_ingestion(db: Session, dataset_id: str) -> str:
    dataset = db.get(Dataset, dataset_id)
    if not dataset or not dataset.enabled:
        raise ValueError("Dataset não encontrado ou desabilitado")

    run = IngestionRun(dataset_id=dataset.id, status=JobStatus.running, started_at=datetime.now(timezone.utc))
    db.add(run)
    db.commit()
    db.refresh(run)

    source = dataset.source
    connector = SQLAlchemyConnector(
        source.db_type.value, source.host, source.port, source.database,
        source.username, decrypt_secret(source.encrypted_password), source.options,
    )
    storage = ObjectStorage()
    total_rows = 0
    latest_watermark = dataset.watermark_value
    tmpdir = tempfile.mkdtemp(prefix="datalake_")

    try:
        part = 0
        uploaded_keys = []
        for chunk in connector.extract(dataset.source_object, dataset.incremental_column, dataset.watermark_value):
            total_rows += len(chunk)
            if total_rows > settings.max_extract_rows:
                raise RuntimeError("Limite máximo de linhas por execução excedido")
            if dataset.masking_policy:
                chunk = _mask_dataframe(chunk, dataset.masking_policy)
            if dataset.incremental_column and not chunk.empty:
                latest_watermark = str(chunk[dataset.incremental_column].max())

            date_path = datetime.now(timezone.utc).strftime("%Y/%m/%d")
            object_key = f"{dataset.destination_zone}/{dataset.name}/{date_path}/{run.id}/part-{part:05d}.parquet"
            local_path = os.path.join(tmpdir, f"part-{part:05d}.parquet")
            chunk.to_parquet(local_path, index=False, compression="snappy")
            storage.upload(local_path, object_key)
            uploaded_keys.append(object_key)
            part += 1

        manifest = "\n".join(uploaded_keys).encode()
        run.checksum = hashlib.sha256(manifest).hexdigest()
        run.object_key = uploaded_keys[0] if uploaded_keys else None
        run.row_count = total_rows
        run.status = JobStatus.success
        run.finished_at = datetime.now(timezone.utc)
        if latest_watermark is not None:
            dataset.watermark_value = latest_watermark
        db.commit()
        return run.id
    except Exception as exc:
        run.status = JobStatus.failed
        run.error_message = str(exc)[:4000]
        run.finished_at = datetime.now(timezone.utc)
        db.commit()
        raise
