ocado-grocy/ocado_grocy/enrichment_retry.py
2026-09-19 11:32:53 +01:00

53 lines
2.8 KiB
Python

"""Durable product-level retries, independent of completed stock orders."""
import json
import sqlite3
class RetryQueue:
def __init__(self, state, server):
self.server = server
self.db = sqlite3.connect(state)
with self.db:
self.db.execute('''CREATE TABLE IF NOT EXISTS enrichment_retries (
server TEXT, stage TEXT, product_id INTEGER, error TEXT,
PRIMARY KEY(server,stage,product_id))''')
self.db.execute('CREATE TABLE IF NOT EXISTS enrichment_retry_migrations (server TEXT PRIMARY KEY)')
def start(self, stage, products):
self.results(stage, [{'product_id':p['id'], 'status':'error', 'error':'Interrupted enrichment'} for p in products])
def results(self, stage, rows):
with self.db:
for row in rows:
key = (self.server,stage,int(row['product_id']))
if row['status'] == 'error':
self.db.execute('INSERT OR REPLACE INTO enrichment_retries VALUES (?,?,?,?)', (*key,row.get('error','Enrichment failed')))
else:
self.db.execute('DELETE FROM enrichment_retries WHERE server=? AND stage=? AND product_id=?',key)
def pending(self, stage):
return {r[0] for r in self.db.execute('SELECT product_id FROM enrichment_retries WHERE server=? AND stage=?',(self.server,stage))}
def migrate_reports(self, archive, allowed_ids):
"""Import old errors once; newer successful reports supersede old failures."""
if self.db.execute('SELECT 1 FROM enrichment_retry_migrations WHERE server=?',(self.server,)).fetchone():
return
events = []
for path in (archive/'enrichment').glob('*/summary.json'):
summary = json.loads(path.read_text())
if summary.get('server',self.server) != self.server:
continue
for stage in ('images','openfoodfacts'):
report_path = path.parent / f'{stage}.json'
if report_path.exists():
data = json.loads(report_path.read_text())
events.append((report_path.stat().st_mtime,stage,data.get('products',[])))
elif any(error.lower().startswith('images:' if stage=='images' else 'open food facts:') for error in summary.get('errors',[])):
events.append((path.stat().st_mtime,stage,[{'product_id':pid,'status':'error','error':'Previous enrichment failed'} for pid in summary.get('product_ids',[])]))
for _,stage,rows in sorted(events,key=lambda event:event[0]):
self.results(stage,[r for r in rows if int(r['product_id']) in allowed_ids])
with self.db:
self.db.execute('INSERT INTO enrichment_retry_migrations VALUES (?)',(self.server,))
def close(self):
self.db.close()