import sqlite3,time,json,datetime,subprocess,glob
from pathlib import Path
DB='/opt/struktur/youtube-research/knowledge.db'; CE='/opt/struktur/content-extraction/data/content_extraction.db'; out={'checked_at_utc':datetime.datetime.now(datetime.timezone.utc).isoformat()}
def db(): return sqlite3.connect('file:'+DB+'?mode=ro',uri=True)
def snap():
c=db(); c.row_factory=sqlite3.Row
cand=[dict(r) for r in c.execute("select e.job_id,e.video_id,e.youtube_id,e.status,e.step,e.attempts,v.transcript_status from e2e_jobs e join videos v on v.id=e.video_id where e.status in ('queued','supadata_retry','supadata_quota_wait') and (e.local_retry_after is null or e.local_retry_after<=CURRENT_TIMESTAMP) order by e.priority,e.job_id")]
return {'candidates':cand,'processing':[dict(r) for r in c.execute("select job_id,status from e2e_jobs where status='supadata_processing'")], 'credits':None}
obs=[]
for i in range(5): obs.append({'at':datetime.datetime.now(datetime.timezone.utc).isoformat(),**snap()}); time.sleep(5)
out['observation']=obs
c=db(); c.row_factory=sqlite3.Row
out['videos']={r['transcript_status']:r['n'] for r in c.execute("select coalesce(transcript_status,'NULL') transcript_status,count(*) n from videos group by transcript_status")}; out['videos']['total']=c.execute('select count(*) from videos').fetchone()[0]
out['e2e']={r['status']:r['n'] for r in c.execute('select status,count(*) n from e2e_jobs group by status')}
jobs=[]
for jid in [275,279]:
r=c.execute("select e.job_id,e.video_id,e.youtube_id,e.status,e.step,e.attempts,e.error_code,e.error_message,e.lease_until,v.transcript_status,v.transcript_source,length(trim(coalesce(v.transcript,''))) transcript_len,v.transcript_hash,v.transcript_updated_at from e2e_jobs e join videos v on v.id=e.video_id where e.job_id=?",(jid,)).fetchone(); d=dict(r); d['segments']=c.execute('select count(*) from transcript_segments where video_id=?',(d['video_id'],)).fetchone()[0]; d['queue_rows']=c.execute('select count(*) from graphiti_import_queue where video_id=?',(d['youtube_id'],)).fetchone()[0]; jobs.append(d)
out['jobs']=jobs; q=c.execute("select * from graphiti_import_queue where video_id='1CLc-VeEivk'").fetchall(); out['queue275']=[dict(r) for r in q]
out['queue']={r['graphiti_status']:r['n'] for r in c.execute('select graphiti_status,count(*) n from graphiti_import_queue group by graphiti_status')}; out['queue']['total']=sum(out['queue'].values()); out['queue']['dup_import_identity']=c.execute('select count(*) from (select import_identity from graphiti_import_queue where import_identity is not null group by import_identity having count(*)>1)').fetchone()[0]; out['queue']['dup_reconciliation_id']=c.execute('select count(*) from (select reconciliation_record_id from graphiti_import_queue where reconciliation_record_id is not null group by reconciliation_record_id having count(*)>1)').fetchone()[0]; out['queue']['release_nonzero']=c.execute('select count(*) from graphiti_import_queue where graphiti_release<>0').fetchone()[0]
out['terminal_exclusion']=c.execute("select count(*) from e2e_jobs e join videos v on v.id=e.video_id where e.status in ('queued','supadata_retry','supadata_quota_wait') and ((e.status='supadata_native_unavailable' and v.transcript_status='unavailable') or (e.status in ('video_unavailable','access_restricted') and v.transcript_status='unavailable'))").fetchone()[0]
out['terminal_rows']=[dict(r) for r in c.execute("select e.job_id,e.status,v.transcript_status from e2e_jobs e join videos v on v.id=e.video_id where e.job_id in (198,245,253,263) or e.status='supadata_native_unavailable' order by e.job_id")]
out['reserve_statuses']=c.execute("select status,count(*) n from e2e_jobs where status in ('queued','supadata_retry','supadata_quota_wait','local_transcript_pending') group by status").fetchall(); c.close()
ce=sqlite3.connect('file:'+CE+'?mode=ro',uri=True); ce.row_factory=sqlite3.Row; out['ce275']=[dict(r) for r in ce.execute("select source_video_id,youtube_id,transcript_status,obsidian_artifact_path,length(transcript) transcript_len,transcript_hash,content_hash from ce_sources where source_video_id=1247 or youtube_id='1CLc-VeEivk'")]; ce.close()
out['obsidian275']=[]
for p in glob.glob('/opt/obsidian-vault/YouTube-Research/**/*.md',recursive=True):
try:
if 'youtube_id: 1CLc-VeEivk' in Path(p).read_text(errors='ignore'): out['obsidian275'].append(p)
except: pass
try:
import sys; sys.path.insert(0,'/opt/struktur/youtube-research'); from supadata_native import account_status; out['account']=account_status()
except Exception as e: out['account']={'error_class':type(e).__name__}
out['service']=subprocess.run(['systemctl','show','youtube-research-e2e-worker.service','-p','ActiveState','-p','SubState','-p','MainPID','-p','NRestarts','-p','ExecMainStatus'],capture_output=True,text=True).stdout
j=subprocess.run(['journalctl','-u','youtube-research-e2e-worker.service','--since','2026-08-27 17:02:00','--no-pager'],capture_output=True,text=True).stdout; out['journal_errors']=[l for l in j.splitlines() if any(x.lower() in l.lower() for x in ['database is locked','sqlite_busy','traceback','ipblocked','requestblocked'])]; out['journal_reservations']=[l for l in j.splitlines() if 'reserved for Supadata Native' in l]
src=Path('/opt/struktur/youtube-research/e2e_worker.py').read_text(errors='ignore'); out['worker_code']={'fetch_native':'fetch_native' in src,'YouTubeTranscriptApi':'YouTubeTranscriptApi' in src,'native_params':all(x in Path('/opt/struktur/youtube-research/supadata_native.py').read_text(errors='ignore') for x in ['mode','native','text','false'])}
print(json.dumps(out,ensure_ascii=False,default=str))