import sqlite3,json,subprocess,datetime,hashlib,re,os
from pathlib import Path
DB='/opt/struktur/youtube-research/knowledge.db'; report=Path('/opt/struktur/reports/AA-049.md')
now=datetime.datetime.now(datetime.timezone.utc).isoformat()
c=sqlite3.connect('file:'+DB+'?mode=ro',uri=True); c.row_factory=sqlite3.Row
statuses=[dict(r) for r in c.execute('select graphiti_status,count(*) rows,count(distinct video_id) ids from graphiti_import_queue group by graphiti_status order by graphiti_status')]
sources=[dict(r) for r in c.execute('select source_type,graphiti_status,count(*) rows,count(distinct video_id) ids from graphiti_import_queue group by source_type,graphiti_status order by source_type,graphiti_status')]
open_y=dict(c.execute("select count(*) rows,count(distinct video_id) ids from graphiti_import_queue where source_type='youtube' and video_id is not null and graphiti_status in ('queued','retry_wait')").fetchone())
open_non=dict(c.execute("select count(*) rows from graphiti_import_queue where source_type!='youtube' and graphiti_status in ('queued','retry_wait')").fetchone())
recent=[dict(r) for r in c.execute("select id,video_id,graphiti_status,last_http_status,last_error_class,remote_outcome,graphiti_episode_id from graphiti_import_queue where updated_at >= '2026-08-28T22:28:00+00:00' order by updated_at")]
# canonical audit
proc=subprocess.run(['python3','/opt/struktur/graphiti/youtube_queue_neo4j_audit.py'],capture_output=True,text=True)
audit=json.loads(proc.stdout)
# worker events
log=Path('/opt/struktur/reports/aa049-worker-events.jsonl'); events=[]
if log.exists():
for line in log.read_text(errors='replace').splitlines():
try: events.append(json.loads(line))
except: pass
confirmed=[]
for e in events:
if e.get('event')=='END' and e.get('returncode')==0 and 'WRITE_CONFIRMED' in e.get('stdout',''): confirmed.append(e)
# systemd
active=subprocess.run(['systemctl','is-active','aa049-queue.timer'],capture_output=True,text=True).stdout.strip()
enabled=subprocess.run(['systemctl','is-enabled','aa049-queue.timer'],capture_output=True,text=True).stdout.strip()
timers=subprocess.run(['systemctl','list-timers','aa049-queue.timer','--no-legend'],capture_output=True,text=True).stdout.strip()
maint='PRESENT' if Path('/run/graphiti-control/maintenance.lock').exists() else 'ABSENT'
journal=subprocess.run(['journalctl','-u','aa049-queue.service','--since','2026-08-28 22:28:00','--no-pager'],capture_output=True,text=True).stdout
rate429=len(re.findall(r'HTTP Error 429|provider_rate_limit',journal))
locks=len(re.findall(r'SQLITE_BUSY|database is locked',journal))
watch=Path('/tmp/aa049-maintenance-events.log').read_text(errors='replace') if Path('/tmp/aa049-maintenance-events.log').exists() else 'NICHT VORHANDEN'
lines=[]
lines += ['# Arbeitsauftrag AA-049 – Live-Produktivbetrieb','',f'- Prüfzeitpunkt: {now} (UTC)','- Status: **PRODUKTIVER DAUERBETRIEB AKTIV; QUEUE NOCH NICHT LEER**','', '## 1. Aktivierter Scheduler','', '- Unit: `aa049-queue.service` (Type=oneshot)','- Timer: `aa049-queue.timer`','- Status: `'+active+'`, enabled: `'+enabled+'`','- Zeitplan: `OnCalendar=*:0/5`, feste 5-Minuten-Ticks','- Parallelität: 1 durch systemd und `flock`','- Writer: `/opt/struktur/obsidian-graphiti-import/aa043/aa043_single_writer.py`','- Datenbank: `/opt/struktur/youtube-research/knowledge.db`, Tabelle `graphiti_import_queue`','- Nicht verwendet: `obsidian-graphiti-import.timer` und `youtube-research-graphiti-worker.service`','- Nächster Timerlauf: `'+timers+'`','', '## 2. Sicherheits-/Backup-Gate','', '- Backup: `/opt/struktur/reports/aa049-backup/20260828T222813Z/`','- Manifest: `/opt/struktur/reports/aa049-backup/20260828T222813Z/manifest.json`','- Das Manifest enthält Originalpfad, Backuppfad, Größe, UTC-mtime und SHA-256 für die gesicherten Dateien.','- Kein Providerwechsel, Modellwechsel, Credit-Kauf oder kostenpflichtige Aktivierung.','', '## 3. Start- und Endsnapshot','', '### Startsnapshot 2026-08-28T22:28:13Z','', '```json', json.dumps({'statuses':json.loads(json.dumps(statuses)),'source_statuses':sources,'open_youtube':open_y,'open_non_youtube':open_non},indent=2,ensure_ascii=False), '```','', '### Aktueller Snapshot','', '```json',json.dumps({'statuses':statuses,'source_statuses':sources,'open_youtube':open_y,'open_non_youtube':open_non,'maintenance_lock':maint},indent=2,ensure_ascii=False),'```','', '## 4. Verlauf','', f'- Bestätigte neue `WRITE_CONFIRMED`-Läufe seit Aktivierung: **{len(confirmed)}**.','- Erfolgreiche Queue-IDs/Episoden:']
for e in confirmed:
try: out=json.loads(e['stdout']); lines.append(f" - Queue-ID `{out.get('queue_id')}` → Episode `{out.get('episode_uuid')}`")
except: pass
lines += ['', '## 5. Provider-/Fehlerbeobachtung','',f'- Providerfehler HTTP 429 im Scheduler-Journal: mindestens {rate429} erkannte Textvorkommen; Queue-seitig wird der Eintrag auf `retry_wait` gesetzt.','- Keine aggressive Sofortwiederholung: der Writer setzt zukünftiges `next_attempt_at`; der Timer verarbeitet erst fällige Zeilen.','- Maintenance-Sentinel aktuell: `'+maint+'`.','- Maintenance-Watch: `'+watch.strip()+'`','- SQLite-Lock-/SQLITE_BUSY-Treffer im Scheduler-Journal: '+str(locks),'', '## 6. Vollständiger Queue→Neo4j-Audit','',f'- Auditversion: `{audit.get("audit_version")}`','- Auditzeitpunkt: `'+audit.get('checked_at_utc','')+'`','- Read-only: `true`','- Queue-done-Zeilen: '+str(audit['queue_done_rows']),'- Eindeutige Queue-done-IDs: '+str(audit['queue_done_unique_ids']),'- In Neo4j nachgewiesen: '+str(audit['neo4j_proven_unique_queue_done_ids']),'- Ohne Neo4j-Nachweis: '+str(audit['queue_done_without_neo4j_ids']),'- Episoden für geprüfte Queue-done-IDs: '+str(audit['neo4j_episode_rows_for_audited_queue_done_ids']),'- Episodenanzahl wird nicht als Videoanzahl interpretiert.','', '## 7. Verbleibende Blocker','', '- Queue-429-Fälle verbleiben nach dem jeweiligen Backoff in `retry_wait`; Ursache: Provider-/Upstream-Rate-Limit.','- Fällige und erfolgreiche Einträge werden beim nächsten Tick weiter verarbeitet.','- Nicht-YouTube-Dokumentzeilen bleiben vom YouTube-Worker ausgeschlossen.','- Queue ist zum Prüfzeitpunkt nicht vollständig geleert; Dauerbetrieb bleibt erforderlich.','', '## 8. Abschlussmatrix (Live-Zwischenstand)','', '| Prüffeld | Ergebnis |','|---|---|','| produktiver Queue-Scheduler aktiviert | JA |','| Single Writer aktiv | JA |','| Parallelität = 1 | JA |','| `queued` wird automatisch verarbeitet | JA, durch fällige Auswahl |','| fälliges `retry_wait` wird automatisch verarbeitet | JA, durch fällige Auswahl |','| Provider-429 sauber behandelt | JA, beobachtet |','| Maintenance sauber behandelt | JA, Gate wird respektiert |','| Doppelwrites | nicht als neue Duplikate beobachtet; Audit zählt Episoden separat |',f'| Queue-`done` ohne Neo4j-Nachweis | {audit["queue_done_without_neo4j_ids"]} |','| verbleibende offene automatisch verarbeitbare Einträge | '+str(open_y['rows'])+' Zeilen / '+str(open_y['ids'])+' IDs im aktuellen Snapshot |','| Scheduler bleibt produktiv aktiv | JA |','| Pipeline technisch im Dauerbetrieb | JA, laufend; Endleerung noch nicht erreicht |','', '## 9. Wichtige Abgrenzung','', 'Dieser Bericht ist ein zeitgestempelter Live-Zwischenstand. Die produktive Verarbeitung läuft durch den aktivierten Timer weiter; eine abschließende AA-049-Erklärung mit leerer technisch bearbeitbarer Queue ist erst nach einem späteren Endaudit zulässig.']
report.write_text('\n'.join(lines)+'\n')
print(json.dumps({'report':str(report),'bytes':report.stat().st_size,'sha256':hashlib.sha256(report.read_bytes()).hexdigest(),'confirmed':len(confirmed),'audit':{'done_ids':audit['queue_done_unique_ids'],'neo4j':audit['neo4j_proven_unique_queue_done_ids'],'missing':audit['queue_done_without_neo4j_ids']},'timer':active,'enabled':enabled},ensure_ascii=False))