from __future__ import annotations import argparse import logging from collections.abc import Callable from typing import Any 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 _missing_handler(event_type): def handler(_payload): raise ValueError(f"no handler registered for {event_type}") return handler def build_governance_knowledge_handlers(app, session): from app.core.data_research.semantic_knowledge import ( build_governance_knowledge_handlers as build_handlers, ) return build_handlers(app, session) def build_default_handlers(app, session): handlers = build_governance_knowledge_handlers(app, session) if not handlers: logger.warning( "Governance knowledge sync is not configured; " "publication events will remain pending" ) return handlers 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, event_types=sorted(handlers), ) for event in events: handler = handlers.get(event.event_type) if handler is None: handler = _missing_handler(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=build_default_handlers(app, db.session), batch_size=args.batch_size, ) logger.info("Processed %s outbox events", processed) if __name__ == "__main__": main()