#!/usr/bin/env python3
"""Recover stale processing leases and restart writer."""
import sqlite3, datetime
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()
print(f"Processing rows: {len(rows)}")
recovered = []
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,))
recovered.append(rid)
k.commit()
print(f"Recovered to lease_expired: {len(recovered)} [IDs: {recovered[:5]}{'...' if len(recovered)>5 else ''}]")
print("\nFinal 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}")
# Show first few processed items for verification
print("\nSample done items:")
for qid, path in k.execute("""select id, obsidian_path
from graphiti_import_queue
where graphiti_status='done'
order by id limit 5""").fetchall():
print(f" {qid}: {path}")