"""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':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') sides,stamp,sequence=normalize(json.loads(raw),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} 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 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])