"""
Shared utilities for dlt connector DAGs.

Provides helpers for reading dlt configuration from Airflow Variables
and running connectors via uv subprocess.
"""

import os
import subprocess

from airflow.models import Variable


CONNECTORS_DIR = "/opt/airflow/connectors"


def get_dlt_destination():
    """Read the DLT destination from Airflow Variable."""
    return Variable.get("dlt_destination")


def get_dlt_secrets_toml():
    """Read DLT secrets TOML from Airflow Variable."""
    dlt_secrets_toml = Variable.get("dlt_secrets_toml")

    if not dlt_secrets_toml or not dlt_secrets_toml.strip():
        raise ValueError("Airflow Variable 'dlt_secrets_toml' is empty")

    return dlt_secrets_toml


def run_dlt_connector(connector_name, source_name, extra_env=None):
    """
    Run a dlt connector via uv in a subprocess.

    Writes secrets.toml directly into the connector's .dlt/ directory,
    then runs the pipeline with uv.

    Args:
        connector_name: Directory name under connectors/
        source_name: DLT_SOURCE_NAME value
        extra_env: Optional dict of additional environment variables
    """
    connector_dir = os.path.join(CONNECTORS_DIR, connector_name)

    if not os.path.isdir(connector_dir):
        raise FileNotFoundError(
            f"Connector directory not found: {connector_dir}"
        )

    # Write secrets.toml into the connector's .dlt/ directory
    dlt_dir = os.path.join(connector_dir, ".dlt")
    os.makedirs(dlt_dir, exist_ok=True)
    with open(os.path.join(dlt_dir, "secrets.toml"), "w") as f:
        f.write(get_dlt_secrets_toml())

    # Pipeline entry point: <connector_name>_pipeline.py
    pipeline_file = f"{connector_name}_pipeline.py"

    env = {
        "PATH": os.environ.get("PATH", "/usr/local/bin:/usr/bin:/bin"),
        "HOME": os.environ.get("HOME", "/home/airflow"),
        "DLT_DESTINATION": get_dlt_destination(),
        "DLT_SOURCE_NAME": source_name,
    }
    if extra_env:
        env.update(extra_env)

    print(f"Running connector: {connector_name} (source: {source_name})")
    print(f"Working directory: {connector_dir}")

    result = subprocess.run(
        ["uv", "run", "python", pipeline_file],
        cwd=connector_dir,
        env=env,
        capture_output=True,
        text=True,
    )

    # Print stdout/stderr to Airflow logs
    if result.stdout:
        print(result.stdout)
    if result.stderr:
        print(result.stderr)

    if result.returncode != 0:
        raise RuntimeError(
            f"Connector {connector_name} failed "
            f"with exit code {result.returncode}"
        )
