Explorer
/tmp/supadata_native_batch.py
← Zurück ↓ Download
#!/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()