| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475 |
- 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()
|