#!/usr/bin/env python3
import argparse,datetime,hashlib,json,os,sqlite3,sys,time
from pathlib import Path
sys.path.insert(0,'/opt/struktur/youtube-research')
import e2e_worker as w
from supadata_native import NativeTranscriptUnavailable, QuotaWait, TransientSupadataError, account_status, fetch_native
DB='/opt/struktur/youtube-research/knowledge.db'; ENV='/opt/struktur/youtube-research/.env.supadata'; REPORT=Path('/opt/struktur/youtube-research/reports/p29'); REPORT.mkdir(parents=True,exist_ok=True)
def now(): return datetime.datetime.now(datetime.timezone.utc).isoformat()
def key():
for line in Path(ENV).read_text().splitlines():
if line.startswith('SUPADATA_API_KEY='): return line.split('=',1)[1].strip()
raise RuntimeError('SUPADATA_API_KEY_missing')
def account(k): return account_status(k)
def conn():
c=sqlite3.connect(DB,timeout=30); c.row_factory=sqlite3.Row; c.execute('PRAGMA busy_timeout=30000'); return c
def candidates(c,limit,exclude_job_ids=()):
q="""SELECT e.job_id,e.video_id,e.youtube_id,e.status e2e_status,e.attempts,v.transcript_status,v.transcript_source,v.created_at,e.updated_at,length(trim(coalesce(v.transcript,''))) transcript_len FROM e2e_jobs e JOIN videos v ON v.id=e.video_id WHERE v.youtube_id IS NOT NULL AND length(trim(coalesce(v.transcript,'')))=0 AND coalesce(v.transcript_status,'')<>'done' AND coalesce(v.transcript_source,'') NOT IN ('supadata_native','windows_local') AND e.status IN ('blocked_transcript_fetch','local_transcript_retry','local_transcript_pending','supadata_retry')"""
params=[]
if exclude_job_ids:
q += " AND e.job_id NOT IN (" + ','.join('?' for _ in exclude_job_ids) + ")"
params.extend(exclude_job_ids)
q += " ORDER BY CASE e.status WHEN 'blocked_transcript_fetch' THEN 0 WHEN 'local_transcript_retry' THEN 1 WHEN 'local_transcript_pending' THEN 2 ELSE 3 END,e.job_id LIMIT ?"
params.append(limit)
return [dict(r) for r in c.execute(q,params)]
def transcript(k,yid):
try:
x=fetch_native(yid,key=k); return 200,{'content':x.content,'lang':x.language},x.billable
except NativeTranscriptUnavailable as ex: return ex.status or 206,{},ex.billable
except QuotaWait as ex: return ex.status or 429,{},ex.billable
except TransientSupadataError as ex: return ex.status or 599,{},ex.billable
def import_one(c,j,payload):
v=c.execute('SELECT * FROM videos WHERE id=?',(j['video_id'],)).fetchone()
if not v: raise RuntimeError('video_missing')
if v['transcript_status']=='done' and (v['transcript'] or '').strip(): return {'result':'skip_already_done'}
content=payload.get('content'); lang=payload.get('lang')
if not isinstance(content,list) or not content or not lang: raise RuntimeError('native_transcript_unavailable')
segs=[]
for i,s in enumerate(content):
if not isinstance(s,dict) or not str(s.get('text','')).strip(): continue
start=float(s.get('offset',0))/1000; dur=float(s.get('duration',0))/1000
segs.append((i,start,start+max(dur,0),str(s['text']).strip()))
text=' '.join(s[3] for s in segs).strip()
if not text or not segs: raise RuntimeError('empty_native_transcript')
th=w.sha(text); extracted=w.extract(dict(v),segs); raw=json.dumps(extracted,ensure_ascii=False,sort_keys=True); ch=w.sha(raw)
c.execute('BEGIN'); c.execute("UPDATE e2e_jobs SET status='supadata_processing',step='supadata_import',updated_at=? WHERE job_id=? AND status IN ('blocked_transcript_fetch','local_transcript_retry','local_transcript_pending','supadata_retry')",(now(),j['job_id'])); c.execute('DELETE FROM transcript_segments WHERE video_id=?',(v['id'],)); c.executemany('INSERT INTO transcript_segments(video_id,start_seconds,text) VALUES(?,?,?)',[(v['id'],round(s[1]),s[3]) for s in segs]); c.execute("UPDATE videos SET transcript=?,transcript_hash=?,transcript_status='done',transcript_source='supadata_native',transcript_updated_at=? WHERE id=?",(text,th,now(),v['id'])); c.commit()
row=dict(v); row['transcript']=text; path=w.write_obsidian(row,extracted,th,ch); w.sync_ce(row,path,th,ch); w.register_graphiti(path,j['youtube_id'])
c.execute("INSERT OR REPLACE INTO e2e_extractions(video_id,youtube_id,extraction_version,transcript_hash,content_hash,payload_json,obsidian_path,processed_at) VALUES(?,?,?,?,?,?,?,?)",(v['id'],j['youtube_id'],w.VERSION,th,ch,raw,path,now())); c.execute("UPDATE e2e_jobs SET status='complete',step='complete',completed_at=?,updated_at=?,lease_until=NULL,error_code=NULL,error_message=NULL WHERE job_id=?",(now(),now(),j['job_id'])); c.commit()
return {'result':'imported','language':lang,'segments':len(segs),'transcript_length':len(text),'transcript_hash':th,'obsidian_path':path}
def main():
ap=argparse.ArgumentParser(); ap.add_argument('--limit',type=int,default=10); ap.add_argument('--reserve-credits',type=int,default=5); ap.add_argument('--exclude-job-ids',default=''); ap.add_argument('--dry-run',action='store_true'); a=ap.parse_args(); assert 1<=a.limit<=10 and a.reserve_credits>=1; exclude=tuple(int(x) for x in a.exclude_job_ids.split(',') if x.strip())
k=key(); before=account(k); assert before['http']==200 and before['remaining'] is not None and before['remaining']>=a.limit+a.reserve_credits
c=conn(); rows=candidates(c,a.limit,exclude); json.dump({'created_at':now(),'account_before':before,'candidates':rows},open(REPORT/'batch-before.json','w'),indent=2,default=str)
if a.dry_run:
print(json.dumps({'before':before,'candidate_count':len(rows),'candidates':rows},ensure_ascii=False)); return
out=[]
for j in rows:
rec={'job_id':j['job_id'],'youtube_id':j['youtube_id'],'started_at':now()}
try:
# Recheck master state immediately before billing request.
r=c.execute("SELECT transcript_status,transcript,transcript_hash FROM videos WHERE id=?",(j['video_id'],)).fetchone()
if not r or (r['transcript_status']=='done' and (r['transcript'] or '').strip()): rec['result']='skip_already_done'; out.append(rec); continue
status,data,bill=transcript(k,j['youtube_id']); rec.update({'http':status,'billable_requests':bill,'response_lang':data.get('lang') if isinstance(data,dict) else None,'available_langs':data.get('availableLangs') if isinstance(data,dict) else None})
if status==200: rec.update(import_one(c,j,data))
elif status in (206,): rec['result']='native_transcript_unavailable'; c.execute("UPDATE e2e_jobs SET status='supadata_native_unavailable',step='supadata_terminal',updated_at=?,error_code='native_unavailable',error_message='native transcript unavailable' WHERE job_id=?",(now(),j['job_id'])); c.commit()
elif status in (401,402,429): rec['result']='api_or_quota_error'; c.execute("UPDATE e2e_jobs SET status='supadata_retry',step='supadata_retry',updated_at=?,error_code=?,error_message=? WHERE job_id=?",(now(),f'http_{status}',f'supadata_http_{status}',j['job_id'])); c.commit(); out.append(rec); break
else: rec['result']='api_error'; c.execute("UPDATE e2e_jobs SET status='supadata_retry',step='supadata_retry',updated_at=?,error_code=?,error_message=? WHERE job_id=?",(now(),f'http_{status}',f'supadata_http_{status}',j['job_id'])); c.commit()
except Exception as ex:
rec.update({'result':'error','error_class':type(ex).__name__}); c.execute("UPDATE e2e_jobs SET status='supadata_retry',step='supadata_retry',updated_at=?,error_code='supadata_processing_error',error_message=? WHERE job_id=?",(now(),type(ex).__name__,j['job_id'])); c.commit()
rec['finished_at']=now(); out.append(rec)
if len([x for x in out if x.get('result')=='imported'])>=a.limit: break
after=account(k); json.dump({'finished_at':now(),'account_after':after,'results':out},open(REPORT/'batch-result.json','w'),indent=2,default=str)
print(json.dumps({'before':before,'after':after,'candidate_count':len(rows),'results':out},ensure_ascii=False))
if __name__=='__main__': main()