#!/usr/bin/env python3
"""Upload one complete CSV as one private dataset, using complete-file streaming (or automatic batches).

ISYNTH_BASE_URL=https://your-site ISYNTH_AUTOMATION_TOKEN=... \
  python upload_csv.py source.csv metadata.json

metadata.json contains source identity, authors, attribution_source,
dataset.description, dataset/license, optional description_details, column_mapping and source_units.
Rerun the identical command after a network failure. Do not split the source file.
"""
import argparse
import hashlib
import csv
import json
import os
import time
import random
from datetime import datetime, timezone
from email.utils import parsedate_to_datetime
from pathlib import Path
from urllib.error import HTTPError, URLError
from urllib.parse import urlsplit
from urllib.request import Request, HTTPRedirectHandler, build_opener

BATCH_BYTES = 900_000  # Leaves room below the server's 1 MiB rows budget.
RETRYABLE_HTTP = {408, 429, 500, 502, 503, 504}
MAX_ATTEMPTS = 6


def retry_delay(headers, attempt):
    """Honor server cooldown (seconds or HTTP date), plus bounded jittered backoff."""
    delay = min(2 ** attempt, 60) + random.uniform(0, 1)
    value = (headers or {}).get('Retry-After', '')
    try:
        wait = float(value)
    except (ValueError, TypeError):
        try:
            wait = (parsedate_to_datetime(value) - datetime.now(timezone.utc)).total_seconds()
        except (ValueError, TypeError, OverflowError):
            wait = 0
    return max(delay, wait) if wait < float('inf') else delay


def batches(path):
    with open(path, encoding='utf-8-sig', newline='') as source:
        reader=csv.DictReader(source)
        pending=[];size=2
        for row in reader:
            if None in row or any(v is None for v in row.values()):raise ValueError('Every source row must match the CSV header')
            weight=len(json.dumps(row,ensure_ascii=False).encode())+2
            if weight+2>BATCH_BYTES:raise ValueError('One source row exceeds the transport budget; do not truncate or split its fields')
            if pending and size+weight>BATCH_BYTES:
                yield pending;pending=[];size=2
            pending.append(row);size+=weight
        if pending:yield pending


class NoRedirect(HTTPRedirectHandler):
    def redirect_request(self,*args,**kwargs):return None


def main():
    parser=argparse.ArgumentParser(description=__doc__)
    parser.add_argument('csv_file',type=Path);parser.add_argument('metadata',type=Path)
    parser.add_argument('--batches',action='store_true',help='Compatibility transport for proxies that reject large requests; still produces one file')
    args=parser.parse_args()
    base=os.environ['ISYNTH_BASE_URL'].rstrip('/');url=urlsplit(base)
    if url.scheme!='https' and not (url.scheme=='http' and url.hostname in {'localhost','127.0.0.1'}):parser.error('Use HTTPS outside localhost')
    if url.username or url.password or url.query or url.fragment or url.path:parser.error('Use a site origin only')
    opener=build_opener(NoRedirect())
    def call(path,body=None):
        data=json.dumps(body,ensure_ascii=False).encode() if body is not None else None
        for attempt in range(MAX_ATTEMPTS):
            delay=retry_delay({},attempt)
            request=Request(base+'/api/automation/v1'+path,data=data,headers={'Authorization':'Bearer '+os.environ['ISYNTH_AUTOMATION_TOKEN'],'Content-Type':'application/json'})
            try:
                with opener.open(request,timeout=180) as response:return json.load(response)
            except HTTPError as exc:
                if exc.code not in RETRYABLE_HTTP or attempt==MAX_ATTEMPTS-1:raise RuntimeError(f'HTTP {exc.code}: {exc.read().decode()}. Resume the same source identity; do not rename or create another dataset.') from None
                delay=retry_delay(exc.headers,attempt)
            except (URLError,TimeoutError):
                if attempt==MAX_ATTEMPTS-1:raise
            time.sleep(delay)
    def send_file(upload_id):
        # Read from disk for each retry; never load the complete CSV into memory.
        sha=hashlib.sha256()
        with args.csv_file.open('rb') as source:
            for block in iter(lambda:source.read(1024*1024),b''):sha.update(block)
        for attempt in range(MAX_ATTEMPTS):
            delay=retry_delay({},attempt)
            with args.csv_file.open('rb') as source:
                request=Request(base+'/api/automation/v1/uploads/'+upload_id+'/file',data=source,method='PUT',headers={
                    'Authorization':'Bearer '+os.environ['ISYNTH_AUTOMATION_TOKEN'],
                    'Content-Type':'text/csv','Content-Length':str(args.csv_file.stat().st_size),
                    'X-Content-SHA256':sha.hexdigest()})
                try:
                    with opener.open(request,timeout=180) as response:return json.load(response)
                except HTTPError as exc:
                    if exc.code not in RETRYABLE_HTTP or attempt==MAX_ATTEMPTS-1:
                        raise RuntimeError(f'HTTP {exc.code}: {exc.read().decode()}. If a proxy rejected the request size, rerun the SAME command with --batches; do not split files.') from None
                    delay=retry_delay(exc.headers,attempt)
                except (URLError,TimeoutError):
                    if attempt==MAX_ATTEMPTS-1:raise
            time.sleep(delay)
    identity=call('/me')
    print('Destination:',identity['username'],'— private draft, publication requires review.')
    metadata=json.loads(args.metadata.read_text(encoding='utf-8'))
    metadata.pop('rows',None)
    with args.csv_file.open(encoding='utf-8-sig',newline='') as source:
        reader=csv.DictReader(source);columns=reader.fieldnames;count=sum(1 for _ in reader)
    if not columns or len(set(columns))!=len(columns) or not count:parser.error('A nonempty CSV with unique columns is required')
    metadata.update(filename=args.csv_file.name,columns=columns,expected_rows=count)
    session=call('/uploads',metadata);upload_id=session['upload_id']
    print('Upload',upload_id,'— resume this same source identity after interruptions.')
    if not args.batches and not session['next_batch_index']:
        receipt=send_file(upload_id)
        print(json.dumps(receipt,ensure_ascii=False,indent=2))
        return
    # Replay even previously acknowledged batches: hashes verify the source and
    # stable batch boundaries, without appending duplicate rows.
    for index,rows in enumerate(batches(args.csv_file)):
        status=call('/uploads/'+upload_id+'/batches',{'batch_index':index,'rows':rows})
        print(status['received_rows'],'/',count,'rows')
    receipt=call('/uploads/'+upload_id+'/finish',{})
    print(json.dumps(receipt,ensure_ascii=False,indent=2))


if __name__=='__main__':main()
