import psql_connection as db
import psycopg
from psycopg import sql
import pandas as pd

SRC_DSN = db.create_psycopg_dsn("opendata.kurtosis.nl/opendata/PostgreSQL")
DST_DSN = db.create_psycopg_dsn("dev/iv3/postgres")


COPY_JOBS = [
    {
        "src_schema": "iv3",
        "dst_schema": "public",
        "tables": [
            "gemeenten_iv3",
            "gemeenten_iv3_meta",
        ],
    },
    {
        "src_schema": "gemeenten",
        "dst_schema": "public",
        "tables": [
            "gemeenten",
        ],
    },
]


def copy_table(src_conn, dst_conn, src_schema, dst_schema, table_name):
    print(f"Start copy {src_schema}.{table_name} -> {dst_schema}.{table_name}...", flush=True)

    try:
        with src_conn.cursor() as src_cur, dst_conn.cursor() as dst_cur:
            dst_cur.execute(
                sql.SQL("TRUNCATE TABLE {}.{}").format(
                    sql.Identifier(dst_schema),
                    sql.Identifier(table_name),
                )
            )

            copy_out_sql = sql.SQL(
                "COPY {}.{} TO STDOUT WITH (FORMAT binary)"
            ).format(
                sql.Identifier(src_schema),
                sql.Identifier(table_name),
            )

            copy_in_sql = sql.SQL(
                "COPY {}.{} FROM STDIN WITH (FORMAT binary)"
            ).format(
                sql.Identifier(dst_schema),
                sql.Identifier(table_name),
            )

            with src_cur.copy(copy_out_sql) as out:
                with dst_cur.copy(copy_in_sql) as inn:
                    for buf in out:
                        inn.write(buf)

            dst_cur.execute(
                sql.SQL("ANALYZE {}.{}").format(
                    sql.Identifier(dst_schema),
                    sql.Identifier(table_name),
                )
            )

        dst_conn.commit()
        print(f"Klaar: {table_name}", flush=True)

    except Exception:
        dst_conn.rollback()
        print(f"FOUT bij kopiëren van {table_name}", flush=True)
        raise

'''
with psycopg.connect(SRC_DSN) as src_conn, psycopg.connect(DST_DSN) as dst_conn:
    for job in COPY_JOBS:
        for table in job["tables"]:
            copy_table(
                src_conn=src_conn,
                dst_conn=dst_conn,
                src_schema=job["src_schema"],
                dst_schema=job["dst_schema"],
                table_name=table,
            )

print("Alle tabellen gekopieerd.")
'''

src_engine = db.create_db_connection("opendata.kurtosis.nl/opendata/PostgreSQL")
dst_engine = db.create_db_connection("dev/iv3/postgres")

df = pd.read_sql("SELECT gemeente, jaar, aantal_mannen + aantal_vrouwen AS aantal_inwoners FROM cbs.demografie WHERE jaar >= 2017 AND leeftijd = 'Totaal' AND LEFT(gemeente,2) = 'GM'",src_engine)

df.to_sql('inwoners', dst_engine, index=False, if_exists='replace', schema='public')