#!/usr/bin/env python3
"""Prepare, apply, inspect-return and publish a bounded OpendTect log edit.
Requires installed Connectors and the public application_access.py companion on PYTHONPATH.
"""
import argparse
import base64
from contextlib import contextmanager
import fcntl
import hashlib
import json
import os
from pathlib import Path
import tempfile
import urllib.error
import urllib.parse
import urllib.request
import uuid
import io
import subprocess
from asset_connectors.opendtect import Host, prepare, compare, changes, digest


def load(path):
    return json.loads(path.read_text())


def save(path,value):
    fd,tmp=tempfile.mkstemp(prefix='.ophiolite-',dir=path.parent)
    try:
        with os.fdopen(fd,'w') as f:
            json.dump(value,f,indent=2,allow_nan=False);f.flush();os.fsync(f.fileno())
        os.replace(tmp,path)
    finally:
        if os.path.exists(tmp):os.unlink(tmp)


def private_bytes(path,raw):
    fd,tmp=tempfile.mkstemp(prefix='.ophiolite-',dir=path.parent)
    try:
        with os.fdopen(fd,'wb') as f:f.write(raw);f.flush();os.fsync(f.fileno())
        os.replace(tmp,path)
    finally:
        if os.path.exists(tmp):os.unlink(tmp)


def verify_result(original,result,curve,patches):
    import lasio
    import numpy as np
    if hashlib.sha256(original).hexdigest()!=curve['source_sha256']:raise ValueError('Retained original checksum changed')
    before=lasio.read(io.StringIO(original.decode('utf-8-sig')))
    after=lasio.read(io.StringIO(result.decode('utf-8-sig')))
    if before.keys()!=after.keys() or [c.unit for c in before.curves]!=[c.unit for c in after.curves]:
        raise ValueError('Published LAS curve/unit mismatch')
    expected=before.data.copy();column=before.keys().index(curve['curve'])
    for patch in patches:expected[patch['index'],column]=np.nan if patch['value'] is None else patch['value']
    if not np.array_equal(expected,after.data,equal_nan=True):raise ValueError('Published LAS differs from the exact reviewed changes')


@contextmanager
def lock(path,mode=0o600):
    fd=os.open(path,os.O_CREAT|os.O_RDWR|os.O_NOFOLLOW,mode)
    if os.fstat(fd).st_uid==os.getuid():os.fchmod(fd,mode)
    try:
        fcntl.flock(fd,fcntl.LOCK_EX|fcntl.LOCK_NB);yield
    finally:os.close(fd)


class NoRedirect(urllib.request.HTTPRedirectHandler):
    def redirect_request(self,*args,**kwargs):raise ValueError('Redirect refused')


class Client:
    def __init__(self,url,project,credentials):
        from application_access import Session,origin
        origin(url)
        if urllib.parse.urlsplit(url).path not in ('','/'):raise ValueError('Use the gateway origin without a path')
        self.url=url.rstrip('/');self.project=project;self.access=Session(credentials)
    def call(self,op,**body):
        url=self.url+'/api/v1/projects/'+urllib.parse.quote(self.project,safe='')+'/applications/'+op
        req=urllib.request.Request(url,data=json.dumps({'project_id':self.project,**body},allow_nan=False).encode(),
            headers={'Content-Type':'application/json',**self.access.headers(self.url,self.project)})
        try:
            with urllib.request.build_opener(NoRedirect).open(req,timeout=60) as r:
                raw=r.read(64*1024*1024+1)
                if len(raw)>64*1024*1024:raise ValueError('Response exceeds transfer limit')
                return json.loads(raw)
        except urllib.error.HTTPError as e:
            raise ValueError(f'{op}: HTTP {e.code}. Check permissions or log in again; retain this directory and retry.') from None


def main(argv=None):
    p=argparse.ArgumentParser(description=__doc__)
    p.add_argument('action',choices=['prepare','apply','inspect-return','publish'])
    p.add_argument('--directory',type=Path,required=True);p.add_argument('--credentials',type=Path,required=True)
    p.add_argument('--preset',type=Path);p.add_argument('--url');p.add_argument('--project');p.add_argument('--binding')
    p.add_argument('--method',default='User-declared native GR sample edits')
    args=p.parse_args(argv);d=args.directory.expanduser()
    if d.is_symlink():raise ValueError('Working directory must not be a symlink')
    d.mkdir(mode=0o700,parents=True,exist_ok=True)
    if d.stat().st_mode&0o077 or d.stat().st_uid!=os.getuid():raise ValueError('Working directory must be owner-private')
    with lock(d/'.lock'):
        statepath=d/'state.json'
        if args.action=='prepare':
            if not all((args.preset,args.url,args.project,args.binding)):raise ValueError('Prepare requires preset, url, project and binding')
            config={'url':args.url.rstrip('/'),'project':args.project,'binding':args.binding,'method':args.method,'preset':load(args.preset)}
            if statepath.exists():
                state=load(statepath)
                if state['config']!=config:raise ValueError('Directory belongs to another request')
            else:
                host=Host(config['preset']);preset={**config['preset'],'identity':host.identity}
                label=preset.get('log_name','')
                if not label or len(label)>60 or any(c not in 'abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789_-' for c in label):raise ValueError('Log name: 1–60 letters, digits, hyphen or underscore')
                client=Client(config['url'],config['project'],args.credentials)
                binding=next((b for b in client.call('list')['bindings'] if b['id']==args.binding),None)
                if not binding:raise ValueError('Binding unavailable')
                ident=uuid.uuid4().hex;preset['log_name']=label+'__'+ident[:12]
                state={'config':config,'preset':preset,'generation':binding['generation'],'command_id':ident,'phase':'new'}
                save(statepath,state)
            client=Client(config['url'],config['project'],args.credentials)
            parameters={'method':args.method,'destination':{k:v for k,v in state['preset'].items() if k!='data_root'},'mapping':'opendtect-gr-patch/1',
              'losses':'Float32 depth/values; missing endpoints omitted natively; API label mapping; return preserves untouched source precision'}
            run=client.call('start',id=args.binding,generation=state['generation'],command_id=state['command_id'],application_version='opendtect-gr-patch/1',parameters=parameters)
            curve=client.call('read',id=run['id']);original=client.call('original',id=run['id'])
            raw=base64.b64decode(original['payload_base64'],validate=True)
            if hashlib.sha256(raw).hexdigest()!=curve['source_sha256'] or original['source']!=curve['source']:raise ValueError('Original reference/checksum mismatch')
            baseline=prepare(curve)
            if state['phase']=='new':
                private_bytes(d/'original.las',raw);save(d/'input.json',curve);save(d/'baseline.json',baseline)
                state.update(run_id=run['id'],phase='prepared',input_digest=digest(curve),baseline_digest=digest(baseline));save(statepath,state)
            elif state['input_digest']!=digest(curve):raise ValueError('Prepared input changed')
            print(json.dumps({'destination':{k:v for k,v in state['preset'].items() if k!='data_root'},'transfer':baseline['report']},indent=2));print('Review the mapping/losses. Close native editors before apply. Python runs beside OpendTect.')
            return
        if not statepath.exists():raise ValueError('Prepare first')
        state=load(statepath);config=state['config'];curve=load(d/'input.json');baseline=load(d/'baseline.json')
        if digest(curve)!=state['input_digest'] or digest(baseline)!=state['baseline_digest']:raise ValueError('Local input/baseline integrity failure')
        client=Client(config['url'],config['project'],args.credentials)
        # Reauthorize exact input before every action; downloaded bytes confer no new API grants.
        if digest(client.call('read',id=state['run_id']))!=state['input_digest']:raise ValueError('Input/interpretation changed')
        host=Host(state['preset']);name=state['preset']['log_name']
        with lock(host.file.parent/('.ophiolite-'+host.info['ID']+'.lock'),0o660):
            if args.action=='apply':
                if state['phase']=='prepared':
                    if host.exists(name):raise ValueError('Name collision before apply; prepare a new directory/name')
                    state['phase']='apply-intent';save(statepath,state)
                if state['phase']!='apply-intent':
                    actual=host.read(name);compare(baseline,actual)
                    print('Already applied; existing native data left untouched.');return
                if host.exists(name):
                    actual=host.read(name);compare(baseline,actual)
                    if actual['values']!=baseline['values']:raise ValueError('Uncertain native write differs; do not overwrite. Inspect and prepare another name.')
                else:actual=host.put(name,baseline)
                state.update(phase='native-verified',native_digest=digest(actual));save(statepath,state)
                save(d/'transfer.json',{'target':state['preset'],'input':curve['source'],'report':baseline['report'],'native_readback':actual,'evidence':'adapter-declared'})
                print('Native log created and fully read back:',name);return
            if state['phase'] not in ('native-verified','edits-reviewed','published'):raise ValueError('Apply and verify the native log first')
            if args.action=='inspect-return':
                if state['phase']=='published':raise ValueError('Run already published; prepare a new run for further edits')
                native=host.read(name);patches=changes(curve,baseline,native)
                pending={'changes':patches,'native_digest':digest(native),'run_id':state['run_id']}
                save(d/'pending.json',pending);state.update(phase='edits-reviewed',pending_digest=digest(pending));save(statepath,state)
                print(json.dumps(pending,indent=2));print('Review these source-index changes, then publish explicitly.');return
            if state['phase'] not in ('edits-reviewed','published'):raise ValueError('Inspect and review returned edits before publication')
            pending=load(d/'pending.json')
            if digest(pending)!=state['pending_digest']:raise ValueError('Reviewed patch changed')
            run=client.call('get',id=state['run_id'])
            if run['state']!='published':
                if digest(host.read(name))!=pending['native_digest']:raise ValueError('Native log changed after review; inspect-return again')
            # Always send the frozen patch, including on retry, so server checks its digest.
            run=client.call('publish',id=state['run_id'],changes=pending['changes'])
            artifact=client.call('download',id=state['run_id']);raw=base64.b64decode(artifact['payload_base64'],validate=True)
            if hashlib.sha256(raw).hexdigest()!=run['receipt']['manifest']['sha256']:raise ValueError('Published checksum mismatch')
            verify_result((d/'original.las').read_bytes(),raw,curve,pending['changes']);private_bytes(d/'result.las',raw);save(d/'receipt.json',run)
            state['phase']='published';save(statepath,state)
            print('Published separate LAS. Original unchanged. Find it in Workspace → Data → Results. Run:',run['id'])

if __name__=='__main__':
    try:main()
    except (ValueError,OSError,KeyError,subprocess.TimeoutExpired) as e:
        raise SystemExit(str(e))
