"""Read-only LTC/EUR study, no accounts/orders; standard library.
python liquidity_collector.py PRIVATE_OUTPUT_DIRECTORY
15-minute scheduled captures, 30-day raw/metric retention.
"""
from pathlib import Path
from concurrent.futures import ThreadPoolExecutor
from decimal import Decimal
import urllib.request,json,time,hashlib,os,sys,email.utils
URLS={'kraken':'https://api.kraken.com/0/public/Depth?pair=LTCEUR&count=100','coinbase':'https://api.exchange.coinbase.com/products/LTC-EUR/book?level=2','okx':'https://eea.okx.com/api/v5/market/books?instId=LTC-EUR&sz=100'}
BUDGETS=[100,1000,10000]
VERSION='lw-ltceur-depth-2026-09-30-v1'
def dec(x):
 d=Decimal(str(x))
 if not d.is_finite() or d<=0:raise ValueError('Invalid price/size')
 return d
def normalize(data,venue):
 stamp=None;sequence=None
 if venue=='kraken':
  if data.get('error') or len(data['result'])!=1:raise ValueError('Provider error')
  book=next(iter(data['result'].values()))
 elif venue=='okx':
  if data.get('code')!='0' or len(data['data'])!=1:raise ValueError('Provider error')
  book=data['data'][0];stamp=int(book['ts'])/1000
 elif venue=='coinbase':
  if data.get('auction_mode') is not False:raise ValueError('No verified continuous book')
  book=data;sequence=book.get('sequence')
 else:raise ValueError('Unknown venue')
 sides={}
 for side in ['bids','asks']:
  levels=book[side]
  if not isinstance(levels,list) or not levels or len(levels)>8000:raise ValueError('Invalid depth')
  merged={}
  for row in levels:
   if not isinstance(row,list) or len(row)<2:raise ValueError('Invalid level')
   p=dec(row[0]);size=dec(row[1]);merged[p]=merged.get(p,Decimal(0))+size
  sides[side]=sorted(merged.items(),reverse=side=='bids')
 if sides['bids'][0][0]>=sides['asks'][0][0]:raise ValueError('Locked/crossed book')
 return sides,stamp,sequence
def replay(levels,target,buy):
 remaining=target;quantity=Decimal(0);cash=Decimal(0)
 for p,size in levels:
  take=min(size,remaining/p if buy else remaining);quantity+=take;cash+=take*p;remaining-=take*p if buy else take
  if remaining<=Decimal('1e-20'):break
 full=remaining<=max(Decimal('1e-20'),target*Decimal('1e-18'))
 return {'full':full,'filled_fraction':float((target-remaining)/target),'ltc':str(quantity),'eur':str(cash),'vwap':str(cash/quantity) if quantity else None}
def metrics(sides):
 bid=sides['bids'][0][0];ask=sides['asks'][0][0];mid=(bid+ask)/2;orders=[]
 for budget in BUDGETS:
  for side in ['buy','sell']:
   buy=side=='buy';r=replay(sides['asks' if buy else 'bids'],Decimal(budget) if buy else Decimal(budget)/bid,buy);v=Decimal(r['vwap'])
   r.update(side=side,budget_eur=budget,target_ltc=None if buy else str(Decimal(budget)/bid),impact_bps=float(max(0,(v/ask-1 if buy else 1-v/bid)*10000)) if r['full'] else None,mid_cost_bps=float(max(0,(v/mid-1 if buy else 1-v/mid)*10000)) if r['full'] else None);orders.append(r)
 return {'bid':str(bid),'ask':str(ask),'mid':str(mid),'spread_bps':float((ask-bid)/mid*10000),'bid_levels':len(sides['bids']),'ask_levels':len(sides['asks']),'depth_1pct_eur':{'buy':str(sum(p*q for p,q in sides['asks'] if p<=ask*Decimal('1.01'))),'sell':str(sum(p*q for p,q in sides['bids'] if p>=bid*Decimal('.99')))},'orders':orders}
def fetch(item):
 venue,url=item;start=time.time()
 try:
  request=urllib.request.Request(url,headers={'User-Agent':'Litecoin.watch liquidity research/1.0','Accept':'application/json'})
  with urllib.request.urlopen(request,timeout=15) as response:
   raw=response.read(2000001)
   if len(raw)>2000000:raise ValueError('Body too large')
   headers={k:response.headers.get(k) for k in ['Date','Age','Content-Type']}
  end=time.time()
  if end-start>20:raise ValueError('Slow capture')
  if headers.get('Age') and int(headers['Age'])>120:raise ValueError('Old HTTP cache')
  if headers.get('Date'):
   http_at=email.utils.parsedate_to_datetime(headers['Date']).timestamp()
   if end-http_at>120 or http_at>end+60:raise ValueError('Old/future HTTP response')
  parsed=json.loads(raw);sides,stamp,sequence=normalize(parsed,venue)
  if stamp is not None and (end-stamp>120 or stamp>end+60):raise ValueError('Old/future timestamp')
  return {'venue':venue,'status':'ok','url':url,'request_at':start,'response_at':end,'exchange_at':stamp,'sequence':sequence,'headers':headers,'sha256':hashlib.sha256(raw).hexdigest(),'metrics':metrics(sides),'raw':raw,**({'auction_mode':parsed['auction_mode']} if venue=='coinbase' else {})}
 except Exception:return {'venue':venue,'status':'unavailable','url':url,'request_at':start,'response_at':time.time(),'reason':'No validated current book in this capture'}
def atomic(path,body):
 temp=path.with_suffix(path.suffix+'.tmp');temp.write_bytes(body);os.chmod(temp,0o600);os.replace(temp,path)
def encode(x):return (json.dumps(x,sort_keys=True,separators=(',',':'))+'\n').encode()
def continuous_metadata(observation,raw_dir=None):
 """Retain a verified Coinbase mode; unknown legacy state stays unknown."""
 result=dict(observation)
 if result.get('venue')!='coinbase' or result.get('status')!='ok':return result
 mode=result.get('auction_mode')
 if type(mode) is not bool:
  mode=None
  try:
   name=result['raw_file'];root=Path(raw_dir).resolve()
   import re
   if not re.fullmatch(r'\d{10}-coinbase\.json',name):raise ValueError('Invalid raw name')
   path=root/name
   if path.resolve().parent!=root or path.stat().st_size>2000000:raise ValueError('Invalid raw book')
   raw=path.read_bytes()
   if hashlib.sha256(raw).hexdigest()!=result['sha256']:raise ValueError('Raw book hash mismatch')
   value=json.loads(raw).get('auction_mode')
   if type(value) is bool:mode=value
  except (OSError,ValueError,TypeError,KeyError):pass
 result['auction_mode']=mode
 if mode is not False:
  result['status']='unavailable';result['reason']='No verified continuous Coinbase book in this capture'
 return result
def execution_series(entries,raw_dir=None):
 """A compact view of retained calculations; no new request or replay method."""
 rows=[]
 for e in entries:
  venues=[]
  for original in e['venues']:
   v=continuous_metadata(original,raw_dir)
   item={k:v[k] for k in ['venue','status','response_at','exchange_at','auction_mode'] if k in v}
   if v['status']=='ok':
    m=v['metrics']
    item.update({k:m[k] for k in ['bid','ask','mid','spread_bps','depth_1pct_eur']})
    item['orders']=[{k:o.get(k) for k in ['side','budget_eur','full','filled_fraction','impact_bps','mid_cost_bps']} for o in m['orders']]
   venues.append(item)
  rows.append({'captured_at':e['captured_at'],'venues':venues})
 return {'schema':1,'protocol':VERSION,'view':'execution-history-v1','entries':rows}
def collect(root):
 root=Path(root);root.mkdir(parents=True,exist_ok=True,mode=0o700)
 try:
  import fcntl
  lock=open(root/'collector.lock','a');fcntl.flock(lock,fcntl.LOCK_EX|fcntl.LOCK_NB)
 except ImportError:lock=None
 now=int(time.time());captures=root/'captures';captures.mkdir(exist_ok=True,mode=0o700)
 if (captures/(str(now)+'-summary.json')).exists():raise ValueError('Capture timestamp already exists')
 with ThreadPoolExecutor(max_workers=3) as pool:observations=list(pool.map(fetch,URLS.items()))
 snap={'schema':1,'protocol':VERSION,'captured_at':now,'pair':'LTC/EUR','venues':observations}
 for o in observations:
  if o['status']=='ok':
   filename=str(now)+'-'+o['venue']+'.json';atomic(captures/filename,o.pop('raw'));o['raw_file']=filename
 atomic(captures/(str(now)+'-summary.json'),encode(snap));atomic(root/'latest.json',encode(snap));entries=[]
 for p in sorted(captures.glob('*-summary.json')):
  if int(p.name.split('-')[0])<now-30*86400:continue
  d=json.loads(p.read_text());entries.append({'captured_at':d['captured_at'],'venues':d['venues']})
 atomic(root/'history.json',encode({'schema':1,'protocol':VERSION,'entries':entries}))
 atomic(root/'series.json',encode(execution_series(entries,captures)))
 for p in captures.glob('*.json'):
  if int(p.name.split('-')[0])<now-30*86400:p.unlink()
 print(json.dumps({'captured_at':now,'venues':{o['venue']:o['status'] for o in observations},'stored_captures':len(entries)}))
 if lock:lock.close()
if __name__=='__main__':collect(sys.argv[1])
