import sqlite3,json,datetime,subprocess,re,urllib.request
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()}
# service preflight
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
# source DB read only
c=sqlite3.connect('file:'+DB+'?mode=ro',uri=True); 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')}
ids=[274,275,276,277,278,279]; out['jobs']=[]
for jid in ids:
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,e.local_retry_after,v.transcript_status,v.transcript_source,length(trim(coalesce(v.transcript,''))) transcript_len from e2e_jobs e left join videos v on v.id=e.video_id where e.job_id=?",(jid,)).fetchone()
if not r: out['jobs'].append({'job_id':jid,'missing':True}); continue
d=dict(r); d['queue_rows']=c.execute("select count(*) from graphiti_import_queue where video_id=?",(d['youtube_id'],)).fetchone()[0]; out['jobs'].append(d)
# exact reserve predicate, no account call
q=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").fetchall(); out['reserve_candidates']=[dict(r) for r in q]
# terminal exclusion simulated against known terminal populations
term=c.execute("select count(*) from e2e_jobs e join videos v on v.id=e.video_id where ((e.status='supadata_native_unavailable' and v.transcript_status='unavailable') or (e.status in ('video_unavailable','access_restricted') and v.transcript_status='unavailable')) and e.status in ('queued','supadata_retry','supadata_quota_wait')").fetchone()[0]; out['terminal_exclusion_pass']=(term==0); out['terminal_selectable']=term
# queue
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]
# trigger and defaults
out['triggers']=[dict(r) for r in c.execute("select name,sql from sqlite_master where type='trigger'")]; out['job_default']=c.execute("select sql from sqlite_master where type='table' and name='e2e_jobs'").fetchone()[0]
# source code markers
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,'reserve_statuses':"status IN ('queued','supadata_retry','supadata_quota_wait')" in src}
# account /me only
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__}
# journal error classes since last start
j=subprocess.run(['journalctl','-u','youtube-research-e2e-worker.service','--since','2026-08-27 15:47:50','--no-pager'],capture_output=True,text=True).stdout
out['journal_error_lines']=[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_supadata_lines']=[l for l in j.splitlines() if 'reserved for Supadata Native' in l]
print(json.dumps(out,ensure_ascii=False,default=str))