"""Run: python examples/http_client.py http://127.0.0.1:8000
No model calls. Credentials/handles stay in client memory, never printed.
"""
import json
import os
from pathlib import Path
import random
import sys
import time
from urllib.request import Request, HTTPRedirectHandler, build_opener
from urllib.error import HTTPError, URLError
from uuid import uuid4


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


# A redirect is not permission to disclose credentials or private handles to
# another endpoint. Canonical API paths used below need no redirect.
urlopen=build_opener(NoRedirect()).open


class Client:
    def __init__(self, base, state_path=None):
        self.base, self.credential = base.rstrip('/'), None
        self.state_path = Path(state_path) if state_path else None
        self.state = {}
        if self.state_path and self.state_path.exists():
            self.state=json.loads(self.state_path.read_text())
            if self.state.get('base') != self.base:
                raise ValueError('Private state belongs to another endpoint')
            self.credential=self.state.get('credential')

    def save(self, **values):
        self.state.update(values,base=self.base,credential=self.credential)
        if self.state_path:
            self.state_path.parent.mkdir(parents=True,exist_ok=True)
            temporary=self.state_path.with_name(self.state_path.name+'.'+str(os.getpid())+'.tmp')
            descriptor=os.open(temporary,os.O_WRONLY|os.O_CREAT|os.O_TRUNC,0o600)
            with os.fdopen(descriptor,'w') as output:
                json.dump(self.state,output)
            os.chmod(temporary,0o600)
            os.replace(temporary,self.state_path)

    def request(self, method, path, body=None, key=None, deadline=None, raw=False):
        deadline = time.monotonic() + 180 if deadline is None else deadline
        delay = 0.5
        retry_safe = method in ('GET','DELETE') or (method == 'POST' and (bool(key) or path == '/v1/quote'))
        payload = json.dumps(body,ensure_ascii=False,allow_nan=False,separators=(',',':')).encode() if body is not None else None
        while True:
            if time.monotonic() >= deadline:
                raise TimeoutError('Wait limit reached; retain the same job ID and idempotency key.')
            headers = {'Content-Type': 'application/json', 'User-Agent': 'MarginNook-Python/0.1'}
            if os.environ.get('MARGINNOOK_STAGING_TOKEN'):
                headers['X-MarginNook-Staging']=os.environ['MARGINNOOK_STAGING_TOKEN']
            if self.credential:
                headers['X-MarginNook-Client'] = self.credential
            if key:
                headers['Idempotency-Key'] = key
            request = Request(self.base + path, data=payload, headers=headers, method=method)
            try:
                with urlopen(request, timeout=min(30, max(1,deadline-time.monotonic()))) as response:
                    self.credential = response.headers.get('X-MarginNook-Client', self.credential)
                    self.save()
                    content=response.read()
                    return content if raw else json.loads(content)
            except HTTPError as error:
                try:
                    details = json.loads(error.read()).get('error', {})
                    if not isinstance(details,dict):raise ValueError('Invalid error envelope')
                except (ValueError,AttributeError):
                    details={'code':'HTTP_ERROR','status':error.code,'category':'transport','retryable':False,
                             'message':'Unexpected HTTP response; check the configured trusted endpoint.'}
                if not retry_safe or details.get('retryable') is False or error.code not in (429,503) or details.get('code') == 'FREE_QUOTA_EXHAUSTED':
                    raise RuntimeError(json.dumps(details)) from None
                pause = float(error.headers.get('Retry-After',delay))
            except URLError:
                # Upload has no idempotency contract: never retry an uncertain upload.
                if not retry_safe:
                    raise
                pause = delay
            if time.monotonic() + pause >= deadline:
                raise TimeoutError('Wait limit reached; retain the same job ID and idempotency key.')
            time.sleep(pause + random.random()*0.1)
            delay = min(5,delay*1.5)

    def run(self, resource, arguments, key):
        body = {'tool': 'data.aggregate', 'resource_id': resource['id'], 'arguments': arguments}
        self.request('POST','/v1/quote',body)
        self.save(resource=resource,arguments=arguments,key=key,body=body,job_id=None)
        deadline = time.monotonic()+180
        job = self.request('POST','/v1/jobs',body,key,deadline)
        self.save(job_id=job['id'])
        job=self.wait(job,deadline)
        if job['status']!='completed' or job['result_expired']:
            raise RuntimeError(json.dumps(job.get('error') or {'code':'RESULT_EXPIRED'}))
        result=self.request('GET','/v1/resources/'+job['result']['result_resource']['id'])
        replay=self.request('POST','/v1/jobs',body,key)
        assert replay['id']==job['id'] and replay['reused']
        return result

    def resume(self):
        """Resume the saved logical task; never create an upload or a new key."""
        deadline=time.monotonic()+180
        if self.state.get('job_id'):
            job=self.request('GET','/v1/jobs/'+self.state['job_id'],deadline=deadline)
        else:
            body=self.state.get('body')
            if not body and self.state.get('resource') and self.state.get('arguments'):
                body={'tool':'data.aggregate','resource_id':self.state['resource']['id'],'arguments':self.state['arguments']}
            if not body or not self.state.get('key'):
                raise ValueError('No saved task to resume; an uncertain upload cannot be retried automatically.')
            job=self.request('POST','/v1/jobs',body,self.state['key'],deadline)
        job=self.wait(job,deadline)
        if job['status']!='completed' or job.get('result_expired'):
            raise RuntimeError(json.dumps(job.get('error') or {'code':'RESULT_EXPIRED'}))
        return job

    def wait(self,job,deadline=None):
        deadline=time.monotonic()+180 if deadline is None else deadline
        self.save(job_id=job['id'])
        delay=0.5
        while job['status']=='pending':
            if time.monotonic()+delay>=deadline:
                raise TimeoutError('Retain job ID and key; timeout does not cancel work.')
            time.sleep(delay)
            job=self.request('GET','/v1/jobs/'+job['id'],deadline=deadline)
            delay=min(5,delay*1.5)
        return job


if __name__=='__main__':
    client=Client(sys.argv[1] if len(sys.argv)>1 else 'http://127.0.0.1:8000', '.runtime/python-client.json')
    if '--resume' in sys.argv[2:]:
        job=client.resume()
        client.request('GET','/v1/resources/'+job['result']['result_resource']['id'])
        print('HTTP Python: saved task resumed and result read; no new upload or key')
        sys.exit(0)
    content='category,amount\nbooks,12.50\nbooks,7.25\ntools,3.00\n'
    resource=client.request('POST','/v1/resources',{'format':'csv','content':content})
    result=client.run(resource,{'group_by':'category','value_field':'amount'},str(uuid4()))
    assert [(g['group'],g['count'],g['sum']) for g in result['groups']]==[('books',2,'19.75'),('tools',1,'3.00')]
    assert client.request('GET','/v1/resources/'+resource['id'],raw=True)==content.encode()
    client.request('DELETE','/v1/resources/'+resource['id'])
    print('HTTP Python: result, evidence download, idempotent replay and management deletion verified')
