from pathlib import Path
import requests,json,time,datetime as dt,fcntl
import pandas as pd
R=Path(__file__).resolve().parent
def utc():return dt.datetime.now(dt.timezone.utc).isoformat()
def dump(p,d):
    tmp=p.with_suffix(p.suffix+'.tmp');tmp.write_text(json.dumps(d,ensure_ascii=False));tmp.replace(p)
policy=R/'collection_policy.json';spec=json.loads(policy.read_text());assert spec['active']
lock=(R/'collection.lock').open('w');fcntl.flock(lock,fcntl.LOCK_EX|fcntl.LOCK_NB)
s=pd.read_parquet(R/'sample.parquet');pages=R/'pages';pages.mkdir(exist_ok=True);log=R/'requests.jsonl';statuses=[];last=0.;nreq=len(log.read_text().splitlines()) if log.exists() else 0;nreviews=0
session=requests.Session();session.headers['User-Agent']='SteamCatalogResearch/0.4 (bounded surviving review-arrival study)'
try:
 for a in s.itertuples():
    cursor='*';seen=set();rows=[];page=0;summary=None
    while True:
        p=pages/f'{a.appid}_{page:03}.json'
        if p.exists():data=json.loads(p.read_text());assert data['requested_cursor']==cursor
        else:
            nowpol=json.loads(policy.read_text())
            if not nowpol['active'] or (R/'STOP').exists():raise RuntimeError('Collection stopped')
            if nreq>=spec['max_requests'] or nreviews+len(rows)>=spec['max_reviews']:raise RuntimeError('Collection cap reached')
            time.sleep(max(0,spec['min_interval_seconds']-(time.monotonic()-last)))
            params={'json':1,'filter':'recent','language':'all','purchase_type':'steam','review_type':'all','num_per_page':100,'filter_offtopic_activity':1,'cursor':cursor}
            rec={'sequence':nreq+1,'appid':a.appid,'page':page,'started_at':utc(),'epoch':time.time(),'params':params};last=time.monotonic();nreq+=1
            try:
                response=session.get(f'https://store.steampowered.com/appreviews/{a.appid}',params=params,timeout=(10,35),stream=True);rec['status']=response.status_code;response.raise_for_status();chunks=[];size=0
                for chunk in response.iter_content(65536):
                    size+=len(chunk)
                    if size>spec['max_response_bytes']:raise RuntimeError('Response too large')
                    chunks.append(chunk)
                raw=json.loads(b''.join(chunks));rec['bytes']=size;assert raw.get('success')==1
                safe=[{k:r.get(k) for k in ['recommendationid','timestamp_created','timestamp_updated','steam_purchase','written_during_early_access']} for r in raw.get('reviews',[])]
                data={'appid':a.appid,'at':rec['started_at'],'requested_cursor':cursor,'next_cursor':raw.get('cursor'),'query_summary':raw.get('query_summary'),'reviews':safe}
                dump(p,data)
            except Exception as e:rec['error']=str(e);raise
            finally:
                rec['finished_at']=utc()
                with log.open('a') as f:f.write(json.dumps(rec)+'\n')
        if page==0:summary=data['query_summary']
        batch=data['reviews'];new=[r for r in batch if r['recommendationid'] not in seen]
        for r in new:seen.add(r['recommendationid']);rows.append(dict(appid=a.appid,**r))
        assert all(r['steam_purchase'] for r in rows)
        if not batch:break
        nxt=data.get('next_cursor')
        if not nxt or nxt==cursor or not new:raise RuntimeError('Nonadvancing review cursor')
        cursor=nxt;page+=1
    nreviews+=len(rows);dest=R/'completed';dest.mkdir(exist_ok=True)
    pd.DataFrame(rows,columns=['appid','recommendationid','timestamp_created','timestamp_updated','steam_purchase','written_during_early_access']).to_parquet(dest/f'{a.appid}.parquet',index=False)
    statuses.append(dict(appid=a.appid,name=a.name,order=a.order,complete=True,collected=len(rows),summary_total=summary['total_reviews'],pages=page+1,snapshot_total=a.reviews))
    dump(R/'status.json',statuses)
    print(json.dumps(statuses[-1],ensure_ascii=False),flush=True)
finally:
 spec['active']=False;spec['requests_used']=nreq;spec['reviews_written']=nreviews;spec['finished_at']=utc();dump(policy,spec)
