process_outbox.py 2.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475
  1. from __future__ import annotations
  2. import argparse
  3. import logging
  4. from typing import Any, Callable
  5. from app import create_app, db
  6. from app.core.events.consumer import dispatch_event
  7. from app.core.events.outbox import (
  8. apply_dispatch_result,
  9. claim_outbox,
  10. consumed,
  11. reset_stale_processing,
  12. )
  13. logger = logging.getLogger(__name__)
  14. def process_available_events(
  15. session: Any,
  16. handlers: dict[str, Callable[[dict[str, Any]], None]],
  17. *,
  18. consumer_name: str = "dataops-default",
  19. batch_size: int = 50,
  20. max_attempts: int = 5,
  21. ) -> int:
  22. if not handlers:
  23. logger.info("No outbox handlers registered; leaving events pending")
  24. return 0
  25. events = claim_outbox(session, limit=batch_size)
  26. for event in events:
  27. handler = handlers.get(event.event_type)
  28. if handler is None:
  29. handler = lambda _payload: (_ for _ in ()).throw(
  30. ValueError(f"no handler registered for {event.event_type}")
  31. )
  32. result = dispatch_event(
  33. event,
  34. handler,
  35. already_processed=consumed(session, consumer_name, event.event_id),
  36. max_attempts=max_attempts,
  37. )
  38. apply_dispatch_result(
  39. session,
  40. consumer_name=consumer_name,
  41. event=event,
  42. result=result,
  43. )
  44. session.commit()
  45. return len(events)
  46. def main() -> None:
  47. parser = argparse.ArgumentParser(description="Process DataOps outbox events")
  48. parser.add_argument("--batch-size", type=int, default=50)
  49. parser.add_argument("--reset-stale-seconds", type=int, default=300)
  50. args = parser.parse_args()
  51. app = create_app()
  52. with app.app_context():
  53. reset_stale_processing(
  54. db.session, older_than_seconds=args.reset_stale_seconds
  55. )
  56. processed = process_available_events(
  57. db.session,
  58. handlers={},
  59. batch_size=args.batch_size,
  60. )
  61. logger.info("Processed %s outbox events", processed)
  62. if __name__ == "__main__":
  63. main()