from __future__ import annotations import argparse import logging from typing import Any, Callable from app import create_app, db from app.core.events.consumer import dispatch_event from app.core.events.outbox import ( apply_dispatch_result, claim_outbox, consumed, reset_stale_processing, ) logger = logging.getLogger(__name__) def process_available_events( session: Any, handlers: dict[str, Callable[[dict[str, Any]], None]], *, consumer_name: str = "dataops-default", batch_size: int = 50, max_attempts: int = 5, ) -> int: if not handlers: logger.info("No outbox handlers registered; leaving events pending") return 0 events = claim_outbox(session, limit=batch_size) for event in events: handler = handlers.get(event.event_type) if handler is None: handler = lambda _payload: (_ for _ in ()).throw( ValueError(f"no handler registered for {event.event_type}") ) result = dispatch_event( event, handler, already_processed=consumed(session, consumer_name, event.event_id), max_attempts=max_attempts, ) apply_dispatch_result( session, consumer_name=consumer_name, event=event, result=result, ) session.commit() return len(events) def main() -> None: parser = argparse.ArgumentParser(description="Process DataOps outbox events") parser.add_argument("--batch-size", type=int, default=50) parser.add_argument("--reset-stale-seconds", type=int, default=300) args = parser.parse_args() app = create_app() with app.app_context(): reset_stale_processing( db.session, older_than_seconds=args.reset_stale_seconds ) processed = process_available_events( db.session, handlers={}, batch_size=args.batch_size, ) logger.info("Processed %s outbox events", processed) if __name__ == "__main__": main()