import sqlite3, json, os, re, subprocess, pathlib, hashlib, datetime
DB='/opt/struktur/youtube-research/knowledge.db'
CE='/opt/struktur/content-extraction/data/content_extraction.db'
VAULT='/opt/obsidian-vault'
def conn(p): return sqlite3.connect(f'file:{p}?mode=ro', uri=True)
def q(c,s,args=()):
c.row_factory=sqlite3.Row
try: return [dict(x) for x in c.execute(s,args)]
except Exception as e: return {'error':str(e),'sql':s}
def scalar(c,s,args=()):
r=q(c,s,args); return r[0][list(r[0])[0]] if isinstance(r,list) and r else None
out={'checked_at_utc':datetime.datetime.now(datetime.timezone.utc).isoformat(),'sources':{'knowledge_db':DB,'ce_db':CE,'vault':VAULT}}
c=conn(DB); out['tables']=q(c,"SELECT name,sql FROM sqlite_master WHERE type='table' ORDER BY name")
for t in ['knowledge_units','videos','e2e_jobs','transcript_segments','graphiti_import_queue','source_processing_registry','reconciliation_decisions']:
out.setdefault('schema',{})[t]=q(c,f'PRAGMA table_info({t})')
for t in ['videos','e2e_jobs','transcript_segments','graphiti_import_queue','source_processing_registry','knowledge_units','reconciliation_decisions']:
out.setdefault('counts',{})[t]=scalar(c,f'SELECT COUNT(*) AS n FROM {t}')
out['ku_rows']=q(c,'SELECT * FROM knowledge_units ORDER BY created_at, rowid')
for col in ['source_identity','source_type','source_id','source_version_id','knowledge_unit_identity']:
try: out.setdefault('ku_group',{})[col]=q(c,f'SELECT {col} AS value, COUNT(*) AS n FROM knowledge_units GROUP BY {col} ORDER BY n DESC')
except: pass
for s,name in [("SELECT MIN(created_at) AS v,MAX(created_at) AS v2 FROM knowledge_units",'ku_dates'),("SELECT COUNT(*) AS n FROM knowledge_units WHERE datetime(replace(created_at,'T',' ')) >= datetime('now','-7 day')",'ku_7d'),("SELECT COUNT(*) AS n FROM knowledge_units WHERE datetime(replace(created_at,'T',' ')) >= datetime('now','-1 day')",'ku_24h')]: out[name]=q(c,s)
out['queue_status']=q(c,'SELECT graphiti_status,COUNT(*) n,MIN(created_at) oldest,MAX(updated_at) newest FROM graphiti_import_queue GROUP BY graphiti_status')
out['queue_open']=q(c,"SELECT COUNT(*) n, MIN(COALESCE(updated_at,created_at)) oldest FROM graphiti_import_queue WHERE graphiti_status IN ('queued','retry_wait','processing')")
out['queue_recent_done']=q(c,"SELECT MAX(COALESCE(updated_at,created_at)) last_success FROM graphiti_import_queue WHERE graphiti_status='done'")
out['registry_status']=q(c,'SELECT artifact_status,quality_status,graphiti_status,COUNT(*) n FROM source_processing_registry GROUP BY artifact_status,quality_status,graphiti_status ORDER BY n DESC')
out['registry_columns']=q(c,'PRAGMA table_info(source_processing_registry)')
out['reconciliation_status']=q(c,'SELECT approval_state,promotion_policy,COUNT(*) n FROM reconciliation_decisions GROUP BY approval_state,promotion_policy')
c.close()
ce=conn(CE); out['ce_tables']=q(ce,"SELECT name,sql FROM sqlite_master WHERE type='table' ORDER BY name")
for t in ['ce_sources','ce_items']:
try: out.setdefault('ce_counts',{})[t]=scalar(ce,f'SELECT COUNT(*) AS n FROM {t}')
except: pass
for t in ['ce_sources','ce_items']:
out.setdefault('ce_schema',{})[t]=q(ce,f'PRAGMA table_info({t})')
out['ce_source_status']=q(ce,'SELECT transcript_status,COUNT(*) n FROM ce_sources GROUP BY transcript_status')
out['ce_artifacts']=q(ce,"SELECT COUNT(*) n FROM ce_sources WHERE COALESCE(TRIM(obsidian_artifact_path),'')<>''")
out['ce_transcripts']=q(ce,"SELECT COUNT(*) n FROM ce_sources WHERE COALESCE(LENGTH(transcript),0)>0")
ce.close()
# filesystem
files=[]
for root,ds,fs in os.walk(VAULT):
ds[:]=[d for d in ds if d not in {'.git','.obsidian'}]
for f in fs:
if f.lower().endswith('.md'): files.append(os.path.join(root,f))
out['obsidian_md_count']=len(files)
out['obsidian_youtube_id_hits']=sum(1 for p in files if re.search(r'(?i)youtub',open(p,errors='ignore').read()))
# code search, bounded active project areas
patterns=re.compile(r'INSERT INTO knowledge_units|UPDATE knowledge_units|knowledge_units\s*\(|create_knowledge_unit|KnowledgeUnit|knowledge_unit',re.I)
paths=[]
for base in ['/opt/struktur/knowledge-pipeline-aggregator-staging','/opt/struktur/obsidian-knowledge-pipeline','/opt/struktur/knowledge-curator','/opt/struktur/obsidian-graphiti-import']:
for root,ds,fs in os.walk(base):
ds[:]=[d for d in ds if d not in {'backups','__pycache__','.git','node_modules'}]
for f in fs:
if f.endswith(('.py','.service','.timer','.sh','.md')):
p=os.path.join(root,f)
try:
txt=open(p,errors='ignore').read()
if patterns.search(txt): paths.append(p)
except: pass
out['ku_code_matches']=paths
# service/timer/ports
for cmd in [["ss","-ltnp"],["systemctl","list-units","--type=service","--all","--no-legend"],["systemctl","list-timers","--all","--no-legend"],["crontab","-l"]]:
key=' '.join(cmd)
try: out.setdefault('system',{})[key]=subprocess.check_output(cmd,text=True,stderr=subprocess.STDOUT,timeout=20)
except Exception as e: out.setdefault('system',{})[key]=str(e)
print(json.dumps(out,ensure_ascii=False,default=str))