import sqlite3, datetime, subprocess
k = sqlite3.connect('/opt/struktur/youtube-research/knowledge.db')
now = datetime.datetime.now(datetime.timezone.utc)
rows = k.execute(
"select id, lease_expires_at from graphiti_import_queue "
"where reconciliation_record_id like 'aa043:%' "
"and graphiti_status='processing'"
).fetchall()
for rid, lease in rows:
if not lease or datetime.datetime.fromisoformat(lease) < now:
k.execute(
"update graphiti_import_queue set graphiti_status='lease_expired', "
"last_error_class='lease_expired', updated_at=CURRENT_TIMESTAMP where id=?",
(rid,)
)
print(f'Recovered {rid}')
k.commit()
print("\nStatus distribution:")
for st, c in k.execute(
"select graphiti_status, count(*) from graphiti_import_queue "
"where reconciliation_record_id like 'aa043:%' group by 1"
).fetchall():
print(f" {st}: {c}")
# Starte den Writer-Loop neu um die recovered Queue zu verarbeiten
subprocess.run(['systemctl', 'restart', 'aa043-single-writer-loop.service'], check=True)
print("\nWriter-Loop restarted")