import sqlite3,json,sys
sys.path.insert(0,'/opt/struktur/youtube-research'); import e2e_worker as w
DB='/opt/struktur/youtube-research/knowledge.db'; results=[]
for jid in range(129,144):
c=sqlite3.connect(DB); c.row_factory=sqlite3.Row; r=c.execute("SELECT e.*,v.* FROM e2e_jobs e JOIN videos v ON v.id=e.video_id WHERE e.job_id=?",(jid,)).fetchone(); c.close()
assert r and r['status']=='blocked_transcript_fetch' and r['transcript_status']=='done' and (r['transcript'] or '').strip()
c=sqlite3.connect(DB); segs=[tuple(x) for x in c.execute('SELECT id,start_seconds,start_seconds,text FROM transcript_segments WHERE video_id=? ORDER BY id',(r['video_id'],))]; c.close(); assert segs
payload=w.extract(r,segs); th=r['transcript_hash'] or w.sha(r['transcript']); ch=w.sha(json.dumps(payload,ensure_ascii=False,sort_keys=True))
path=w.write_obsidian(r,payload,th,ch); w.sync_ce(r,path,th,ch); w.register_graphiti(path,r['youtube_id'])
c=sqlite3.connect(DB); c.row_factory=sqlite3.Row
q=c.execute('SELECT COUNT(*) FROM graphiti_import_queue WHERE obsidian_path=?',(path,)).fetchone()[0]
assert q==1,(jid,'queue confirmation',q,path)
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(?,?,?,?,?,?,?,?)',(r['video_id'],r['youtube_id'],w.VERSION,th,ch,json.dumps(payload,ensure_ascii=False),path,w.now()))
cur=c.execute("UPDATE e2e_jobs SET status='complete',step='complete',completed_at=COALESCE(completed_at,?),updated_at=?,lease_until=NULL WHERE job_id=? AND status='blocked_transcript_fetch'",(w.now(),w.now(),jid)); assert cur.rowcount==1; c.commit(); c.close()
results.append({'job_id':jid,'video_id':r['video_id'],'queue':q,'path':path})
# ten exact false done
c=sqlite3.connect(DB); c.row_factory=sqlite3.Row; c.execute('BEGIN IMMEDIATE')
for jid in range(33,43):
r=c.execute("SELECT e.job_id,e.video_id,e.status,v.transcript_status,v.transcript FROM e2e_jobs e JOIN videos v ON v.id=e.video_id WHERE e.job_id=?",(jid,)).fetchone(); assert r and r['status']=='blocked_transcript_fetch' and r['transcript_status']=='done' and not (r['transcript'] or '').strip()
assert c.execute("UPDATE videos SET transcript_status='retry',transcript_updated_at=NULL WHERE id=? AND transcript_status='done' AND length(trim(coalesce(transcript,'')))=0",(r['video_id'],)).rowcount==1
assert c.execute("UPDATE e2e_jobs SET status='supadata_retry',step='supadata_retry',updated_at=? WHERE job_id=? AND status='blocked_transcript_fetch'",(w.now(),jid)).rowcount==1
c.commit(); c.close()
# exact missing-job set after above
c=sqlite3.connect(DB); c.row_factory=sqlite3.Row; missing=[r for r in c.execute("SELECT v.id,v.youtube_id FROM videos v WHERE v.transcript_status='retry' AND length(trim(coalesce(v.transcript,'')))=0 ORDER BY v.id") if c.execute('SELECT COUNT(*) FROM e2e_jobs WHERE video_id=?',(r['id'],)).fetchone()[0]==0]; assert len(missing)==80,(len(missing),)
c.execute('BEGIN IMMEDIATE')
for r in missing: assert c.execute("INSERT INTO e2e_jobs(video_id,youtube_id,status,step,priority,attempts,updated_at) VALUES(?,?, 'supadata_retry','supadata_retry',0,0,?)",(r['id'],r['youtube_id'],w.now())).rowcount==1
c.commit(); c.close()
print(json.dumps({'stale_repaired':15,'false_done_repaired':10,'jobs_created':80,'stale_results':results},ensure_ascii=False))