process_outbox.py 2.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102
  1. from __future__ import annotations
  2. import argparse
  3. import logging
  4. from collections.abc import Callable
  5. from typing import Any
  6. from app import create_app, db
  7. from app.core.events.consumer import dispatch_event
  8. from app.core.events.outbox import (
  9. apply_dispatch_result,
  10. claim_outbox,
  11. consumed,
  12. reset_stale_processing,
  13. )
  14. logger = logging.getLogger(__name__)
  15. def _missing_handler(event_type):
  16. def handler(_payload):
  17. raise ValueError(f"no handler registered for {event_type}")
  18. return handler
  19. def build_governance_knowledge_handlers(app, session):
  20. from app.core.data_research.semantic_knowledge import (
  21. build_governance_knowledge_handlers as build_handlers,
  22. )
  23. return build_handlers(app, session)
  24. def build_default_handlers(app, session):
  25. handlers = build_governance_knowledge_handlers(app, session)
  26. if not handlers:
  27. logger.warning(
  28. "Governance knowledge sync is not configured; "
  29. "publication events will remain pending"
  30. )
  31. return handlers
  32. def process_available_events(
  33. session: Any,
  34. handlers: dict[str, Callable[[dict[str, Any]], None]],
  35. *,
  36. consumer_name: str = "dataops-default",
  37. batch_size: int = 50,
  38. max_attempts: int = 5,
  39. ) -> int:
  40. if not handlers:
  41. logger.info("No outbox handlers registered; leaving events pending")
  42. return 0
  43. events = claim_outbox(
  44. session,
  45. limit=batch_size,
  46. event_types=sorted(handlers),
  47. )
  48. for event in events:
  49. handler = handlers.get(event.event_type)
  50. if handler is None:
  51. handler = _missing_handler(event.event_type)
  52. result = dispatch_event(
  53. event,
  54. handler,
  55. already_processed=consumed(session, consumer_name, event.event_id),
  56. max_attempts=max_attempts,
  57. )
  58. apply_dispatch_result(
  59. session,
  60. consumer_name=consumer_name,
  61. event=event,
  62. result=result,
  63. )
  64. session.commit()
  65. return len(events)
  66. def main() -> None:
  67. parser = argparse.ArgumentParser(description="Process DataOps outbox events")
  68. parser.add_argument("--batch-size", type=int, default=50)
  69. parser.add_argument("--reset-stale-seconds", type=int, default=300)
  70. args = parser.parse_args()
  71. app = create_app()
  72. with app.app_context():
  73. reset_stale_processing(
  74. db.session, older_than_seconds=args.reset_stale_seconds
  75. )
  76. processed = process_available_events(
  77. db.session,
  78. handlers=build_default_handlers(app, db.session),
  79. batch_size=args.batch_size,
  80. )
  81. logger.info("Processed %s outbox events", processed)
  82. if __name__ == "__main__":
  83. main()