#!/usr/bin/env bash
# Load delta CSVs from ~/migration/delta_csv into self-hosted Postgres.
# Filters each CSV to shared non-generated columns (handles schema drift).
set -euo pipefail
source ~/migration/db-env.sh
cd ~/migration/delta_csv
mkdir -p filtered
: > ../load_delta.log
set +e   # continue after one table fails

for csv in *.csv; do
  [[ -f "$csv" ]] || continue
  schema=${csv%%.*}
  table=${csv#*.}
  table=${table%.csv}
  tbl="${schema}.${table}"

  pk_cols=$(psql -At -c "
    SELECT string_agg(a.attname, ', ' ORDER BY x.n)
    FROM pg_index i
    JOIN LATERAL unnest(i.indkey) WITH ORDINALITY AS x(attnum, n) ON true
    JOIN pg_attribute a ON a.attrelid = i.indrelid AND a.attnum = x.attnum
    WHERE i.indrelid = '${tbl}'::regclass AND i.indisprimary;
  ")
  if [[ -z "$pk_cols" ]]; then
    echo "SKIP $tbl — no primary key" | tee -a ../load_delta.log
    continue
  fi

  # Target non-generated columns (one per line)
  psql -At -c "
    SELECT a.attname
    FROM pg_attribute a
    JOIN pg_class c ON c.oid = a.attrelid
    JOIN pg_namespace n ON n.oid = c.relnamespace
    WHERE n.nspname='${schema}' AND c.relname='${table}'
      AND a.attnum > 0 AND NOT a.attisdropped
      AND a.attgenerated = ''
    ORDER BY a.attnum;
  " > "/tmp/tgt_cols_${schema}_${table}.txt"

  if [[ ! -s "/tmp/tgt_cols_${schema}_${table}.txt" ]]; then
    echo "SKIP $tbl — table missing on target?" | tee -a ../load_delta.log
    continue
  fi

  filtered="filtered/${csv}"
  # Intersect CSV header with target cols in Python (avoids psql -v / CRLF bugs)
  shared_cols=$(python3 - "$csv" "$filtered" "/tmp/tgt_cols_${schema}_${table}.txt" <<'PY'
import csv, sys
src, dst, tgt_path = sys.argv[1], sys.argv[2], sys.argv[3]
with open(tgt_path, encoding="utf-8") as f:
    tgt = {line.strip() for line in f if line.strip()}
with open(src, newline="", encoding="utf-8-sig") as f:  # utf-8-sig strips BOM
    reader = csv.DictReader(f)
    if not reader.fieldnames:
        print("", end="")
        sys.exit(1)
    # Preserve CSV header order; only keep cols that exist on target
    cols = [c.strip() for c in reader.fieldnames if c and c.strip() in tgt]
    if not cols:
        print("", end="")
        sys.exit(2)
    rows = list(reader)
with open(dst, "w", newline="", encoding="utf-8") as f:
    w = csv.DictWriter(f, fieldnames=cols, extrasaction="ignore")
    w.writeheader()
    for r in rows:
        w.writerow({c: r.get(c, "") for c in cols})
print(",".join(cols), end="")
PY
)
  rc=$?
  if [[ $rc -ne 0 || -z "$shared_cols" ]]; then
    echo "SKIP $tbl — no shared columns (python rc=$rc)" | tee -a ../load_delta.log
    echo "  CSV header: $(head -1 "$csv" | tr -d '\r' | cut -c1-120)..." | tee -a ../load_delta.log
    echo "  Target cols sample: $(head -5 /tmp/tgt_cols_${schema}_${table}.txt | tr '\n' ' ')" | tee -a ../load_delta.log
    continue
  fi

  quoted_cols=$(psql -At -c "
    SELECT string_agg(quote_ident(x), ', ')
    FROM unnest(string_to_array('${shared_cols}', ',')) AS x;
  ")

  updates=$(psql -At -c "
    WITH pk AS (
      SELECT trim(x) AS col FROM unnest(string_to_array('${pk_cols}', ',')) AS x
    )
    SELECT coalesce(
      string_agg(format('%I = EXCLUDED.%I', x, x), ', '),
      ''
    )
    FROM unnest(string_to_array('${shared_cols}', ',')) AS x
    WHERE trim(x) NOT IN (SELECT col FROM pk);
  ")

  echo "load $tbl (cols=$(echo "$shared_cols" | tr ',' ' ' | wc -w)) ..." | tee -a ../load_delta.log
  if [[ -n "$updates" ]]; then
    conflict_sql="ON CONFLICT (${pk_cols}) DO UPDATE SET ${updates}"
  else
    conflict_sql="ON CONFLICT (${pk_cols}) DO NOTHING"
  fi

  psql --variable ON_ERROR_STOP=1 <<EOSQL 2>&1 | tee -a ../load_delta.log
SET session_replication_role = replica;
CREATE TEMP TABLE staging_delta (LIKE ${tbl} INCLUDING DEFAULTS);
\\copy staging_delta (${quoted_cols}) FROM '${PWD}/${filtered}' CSV HEADER
INSERT INTO ${tbl} (${quoted_cols})
SELECT ${quoted_cols} FROM staging_delta
${conflict_sql};
DROP TABLE staging_delta;
EOSQL

  if [[ $? -ne 0 ]]; then
    echo "FAILED $tbl" | tee -a ../load_delta.log
  else
    echo "OK $tbl" | tee -a ../load_delta.log
  fi
done

echo "Done. Check ~/migration/load_delta.log" | tee -a ../load_delta.log
