Explorer
/tmp/recover_leases.py
← Zurück ↓ Download
#!/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}")