"""Pinned Qwen acquisition/verification. Download is an explicit command, never an import effect."""
import argparse
import hashlib
import json
import multiprocessing
import os
from pathlib import Path
import stat
import threading
import time
import urllib.error
import urllib.request
from tiny_agent_boundary import Denied
from tiny_agent_workspace import fresh_directory

REPOSITORY='Qwen/Qwen2.5-Coder-0.5B-Instruct'
REVISION='ea3f2471cf1b1f0db85067f1ef93848e38e88c25'
FILES={
 'config.json':(659,'e2aa18293b0dd66539341467fa454878a5adb2be'),
 'generation_config.json':(243,'c28f9c697cbc09f047434efec57556505af31111'),
 'tokenizer_config.json':(7305,'acee076f49bf3c0298e15de0909d1da7b392f0c3'),
 'tokenizer.json':(7031645,'443909a61d429dff23010e5bddd28ff530edda00'),
 'vocab.json':(2776833,'4783fe10ac3adce15ac8f358ef5462739852c569'),
 'merges.txt':(1671839,'20024bfe7c83998e9aeaf98a0cd6a2ce6306c2f0'),
 'model.safetensors':(988097824,'f9523886352217ded3aeeef552b381af79d568c6d49a4b9e423288cea56b0a44'),
 'LICENSE':(11343,'6634c8cc3133b3848ec74b9f275acaaa1ea618ab')}
LIMIT=1_200_000_000
PREFIX=f'https://huggingface.co/{REPOSITORY}/resolve/{REVISION}/'


def hash_file(path,size):
    sha=hashlib.sha256();blob=hashlib.sha1(b'blob '+str(size).encode()+b'\0');observed=0
    with path.open('rb') as handle:
        while data:=handle.read(1024*1024):sha.update(data);blob.update(data);observed+=len(data)
    return observed,sha.hexdigest(),blob.hexdigest()


def plan():
    return {'repository':REPOSITORY,'revision':REVISION,'expected_payload_bytes':sum(n for n,_ in FILES.values()),
            'payload_ceiling':LIMIT,'acquisition_seconds':1200,'execution_seconds':600,'whole_seconds':1800,
            'maximum_retries_per_file':1,'model_parameters':494032768,'load':'CPU_float32_eager',
            'files':{name:{'url':PREFIX+name,'size':size,'identity':identity,
                           'identity_kind':'sha256' if name=='model.safetensors' else 'git_blob_sha1'}
                     for name,(size,identity) in FILES.items()}}


def verify(directory,*,require_approval=True):
    directory=Path(directory).absolute()
    if any(p.is_symlink() for p in (directory,*directory.parents)):raise Denied('artifact_symlink')
    if not directory.is_dir():raise Denied('artifact_missing')
    names={p.name for p in directory.iterdir()}
    approved={'supervisor-approved.json'} if require_approval else set()
    if names!=set(FILES)|{'acquisition.json'}|approved:raise Denied('unexpected_artifact')
    ledger_info=(directory/'acquisition.json').lstat()
    if not stat.S_ISREG(ledger_info.st_mode) or ledger_info.st_nlink!=1:raise Denied('artifact_type')
    ledger=json.loads((directory/'acquisition.json').read_text())
    if (ledger.get('status')!='complete' or ledger.get('repository')!=REPOSITORY
        or ledger.get('revision')!=REVISION or type(ledger.get('payload_bytes')) is not int
        or not sum(n for n,_ in FILES.values())<=ledger['payload_bytes']<=LIMIT):
        raise Denied('artifact_ledger')
    if require_approval:
        marker=directory/'supervisor-approved.json';info=marker.lstat()
        if not stat.S_ISREG(info.st_mode) or info.st_nlink!=1:raise Denied('artifact_type')
        approval=json.loads(marker.read_text())
        if (approval.get('status')!='approved' or approval.get('repository')!=REPOSITORY
            or approval.get('revision')!=REVISION
            or approval.get('ledger_sha256')!=hashlib.sha256((directory/'acquisition.json').read_bytes()).hexdigest()):
            raise Denied('supervisor_approval')
    observations={}
    for name,(size,identity) in FILES.items():
        path=directory/name;info=path.lstat()
        if not stat.S_ISREG(info.st_mode) or info.st_nlink!=1:raise Denied('artifact_type')
        observed,sha,blob=hash_file(path,size)
        if observed!=size or (sha if name=='model.safetensors' else blob)!=identity:raise Denied('artifact_identity')
        observations[name]={'bytes':observed,'sha256':sha,'git_blob_sha1':blob}
    config=json.loads((directory/'config.json').read_text())
    if ('auto_map' in config or config.get('architectures')!=['Qwen2ForCausalLM'] or config.get('model_type')!='qwen2'):
        raise Denied('model_architecture')
    for name in ('generation_config.json','tokenizer_config.json'):
        if 'auto_map' in json.loads((directory/name).read_text()):raise Denied('remote_code_metadata')
    return observations


def acquire(directory,*,started=None,clock=time.monotonic,opener=urllib.request.urlopen):
    start=clock() if started is None else started;deadline=min(start+1200,start+1800)
    directory=fresh_directory(directory)
    manifest={**plan(),'status':'in_progress','started_monotonic':start,'payload_bytes':0,'attempts':[]}
    # Trusted metadata/ledger lives outside the model's eight immutable artifacts.
    def save():
        tmp=directory/'acquisition.tmp'
        with tmp.open('x') as handle:json.dump(manifest,handle,indent=2);handle.write('\n');handle.flush();os.fsync(handle.fileno())
        os.replace(tmp,directory/'acquisition.json')
    save()
    try:
        for name,(size,identity) in FILES.items():
            for attempt in range(2):
                record={'file':name,'attempt':attempt+1,'url':PREFIX+name,'received_bytes':0,'status':'started'}
                manifest['attempts'].append(record);save()
                partial=directory/f'{name}.part-{attempt+1}'
                sha=hashlib.sha256();blob=hashlib.sha1(b'blob '+str(size).encode()+b'\0')
                try:
                    remaining=deadline-clock()
                    if remaining<=0:raise Denied('acquisition_time_limit')
                    request=urllib.request.Request(PREFIX+name,headers={'User-Agent':'LLMCourse-pinned-lab16/1'})
                    with opener(request,timeout=min(30,remaining)) as response,partial.open('xb') as handle:
                        record['observed_url']=response.geturl()
                        length=response.headers.get('Content-Length')
                        if length is not None and int(length)!=size:raise Denied('artifact_size_metadata')
                        while True:
                            if clock()>=deadline:raise Denied('acquisition_time_limit')
                            # Never consume a byte beyond the shared payload allowance.
                            if length is not None and record['received_bytes']==size:break
                            remaining_payload=LIMIT-manifest['payload_bytes']
                            if remaining_payload<=0:raise Denied('artifact_payload_limit')
                            chunk=response.read(min(65536,size-record['received_bytes']+1,remaining_payload))
                            if not chunk:break
                            record['received_bytes']+=len(chunk);manifest['payload_bytes']+=len(chunk)
                            if record['received_bytes']>size or manifest['payload_bytes']>LIMIT:raise Denied('artifact_payload_limit')
                            handle.write(chunk);sha.update(chunk);blob.update(chunk)
                        handle.flush();os.fsync(handle.fileno())
                    record.update(sha256=sha.hexdigest(),git_blob_sha1=blob.hexdigest())
                    if record['received_bytes']!=size:raise OSError('Incomplete artifact transfer')
                    if (sha.hexdigest() if name=='model.safetensors' else blob.hexdigest())!=identity:raise Denied('artifact_identity')
                    os.replace(partial,directory/name);record['status']='verified';save();break
                except (OSError,urllib.error.URLError) as error:
                    record['status']='transfer_failed';record['diagnostic_type']=type(error).__name__
                    if isinstance(error,urllib.error.HTTPError):record['http_status']=error.code
                    save()
                    # Retain partial evidence outside the loadable artifact directory, count every byte.
                    if partial.exists():os.replace(partial,directory.parent/(directory.name+'-'+partial.name))
                    if attempt==1 or clock()>=deadline:raise
                except BaseException:
                    record['status']='aborted';save();raise
        manifest['status']='complete';manifest['elapsed_seconds']=clock()-start;save()
        verify(directory,require_approval=False)
        return manifest
    except BaseException as error:
        manifest['status']='aborted';manifest['diagnostic_type']=type(error).__name__
        manifest['elapsed_seconds']=clock()-start;save();raise


def acquisition_worker(connection,directory,started):
    try:
        result=acquire(directory,started=started)
        connection.send_bytes(json.dumps({'status':'complete','manifest':result}).encode())
    except BaseException as error:
        connection.send_bytes(json.dumps({'status':'aborted','code':getattr(error,'code','transfer_failed'),
                                        'diagnostic_type':type(error).__name__}).encode())
    finally:connection.close()


def cleanup_acquisition(process,parent,reader):
    parent.close()
    if process.is_alive():process.terminate()
    process.join(timeout=1)
    if process.is_alive():process.kill();process.join(timeout=1)
    reader.join(timeout=1)
    if process.is_alive():raise Denied('acquisition_worker_termination_failed')
    if reader.is_alive():raise Denied('acquisition_reader_termination_failed')


def validate_receipt(result,directory,started):
    if type(result) is not dict or set(result)!={'status','manifest'} or result['status']!='complete':
        raise Denied('acquisition_receipt')
    manifest=result['manifest']
    if type(manifest) is not dict or any(manifest.get(k)!=v for k,v in plan().items()):
        raise Denied('acquisition_receipt')
    if (manifest.get('status')!='complete' or manifest.get('started_monotonic')!=started
        or type(manifest.get('payload_bytes')) is not int
        or not sum(n for n,_ in FILES.values())<=manifest['payload_bytes']<=LIMIT):
        raise Denied('acquisition_receipt')
    if manifest!=json.loads((directory/'acquisition.json').read_text()):raise Denied('acquisition_receipt')
    return manifest


def invalidate_acquisition(directory,started):
    # Approval is never created on failure; even a surviving worker cannot create this
    # supervisor-owned authorization through the fixed acquisition-worker implementation.
    ledger=directory/'acquisition.json'
    if ledger.is_file() and not ledger.is_symlink():
        record=json.loads(ledger.read_text())
        record.update(previous_status=record.get('status'),status='aborted',
                      supervisor_status='unsuccessful',elapsed_seconds=time.monotonic()-started,
                      payload_accounting='interrupted; recorded attempts plus retained partial evidence; no load permitted')
        evidence=directory/'supervisor-failure.json'
        with evidence.open('x') as handle:json.dump(record,handle,indent=2);handle.write('\n')
        replacement=directory/'supervisor-ledger.tmp'
        with replacement.open('x') as handle:
            json.dump(record,handle,indent=2);handle.write('\n');handle.flush();os.fsync(handle.fileno())
        os.replace(replacement,ledger)


def supervised_acquire(directory,*,target=acquisition_worker,seconds=1200):
    candidate=Path(directory).absolute()
    if candidate.exists() or candidate.is_symlink():raise Denied('run_exists')
    for ancestor in (candidate,*candidate.parents):
        if ancestor.is_symlink():raise Denied('unsafe_run_path')
        if (ancestor/'.git').exists() or (ancestor/'.git').is_symlink():raise Denied('run_inside_repository')
    if not candidate.parent.is_dir():raise Denied('run_parent_missing')
    start=time.monotonic();context=multiprocessing.get_context('spawn')
    parent,child=context.Pipe(duplex=False)
    process=context.Process(target=target,args=(child,str(candidate),start),daemon=True)
    process.start();child.close();deadline=start+seconds
    received=[];errors=[]
    def receive_frame():
        try:received.append(parent.recv_bytes(65536))
        except (EOFError,OSError,ValueError) as error:errors.append(type(error).__name__)
    reader=threading.Thread(target=receive_frame,daemon=True);reader.start()
    succeeded=False
    try:
        try:
            result=None
            while time.monotonic()<deadline:
                reader.join(timeout=min(.05,max(0,deadline-time.monotonic())))
                if not reader.is_alive():
                    if errors or not received:raise Denied('acquisition_worker_error')
                    result=json.loads(received[0]);break
                if not process.is_alive():raise Denied('acquisition_worker_error')
            if result is None or time.monotonic()>=deadline:raise Denied('acquisition_time_limit')
            if type(result) is dict and result.get('status')=='aborted':
                raise Denied(result.get('code','acquisition_worker_error'))
            manifest=validate_receipt(result,candidate,start)
        finally:
            # Cleanup exceptions still reach the OUTER failure handler below.
            cleanup_acquisition(process,parent,reader)
        if time.monotonic()>=deadline:raise Denied('acquisition_time_limit')
        approval={'status':'approved','repository':REPOSITORY,'revision':REVISION,
                  'ledger_sha256':hashlib.sha256((candidate/'acquisition.json').read_bytes()).hexdigest(),
                  'started_monotonic':start,'approved_monotonic':time.monotonic()}
        with (candidate/'supervisor-approved.json').open('x') as handle:
            json.dump(approval,handle,indent=2);handle.write('\n');handle.flush();os.fsync(handle.fileno())
        if time.monotonic()>=deadline:raise Denied('acquisition_time_limit')
        succeeded=True
        return manifest
    finally:
        if not succeeded:invalidate_acquisition(candidate,start)


def main():
    parser=argparse.ArgumentParser(description=__doc__)
    parser.add_argument('command',choices=('plan','verify','acquire'))
    parser.add_argument('--directory',type=Path)
    args=parser.parse_args()
    if args.command=='plan':print(json.dumps(plan(),indent=2));return
    if args.directory is None:parser.error('--directory required')
    result=verify(args.directory) if args.command=='verify' else supervised_acquire(args.directory)
    print(json.dumps(result,indent=2))

if __name__=='__main__':main()
