from __future__ import annotations from datetime import datetime, timedelta, timezone from celery import Celery from govoplan_core.core.campaigns import CAPABILITY_CAMPAIGNS_DELIVERY_TASKS, CampaignDeliveryTaskProvider from govoplan_core.core.calendar import CAPABILITY_CALENDAR_OUTBOX, CalendarOutboxProvider from govoplan_core.core.dataflows import ( CAPABILITY_DATAFLOW_RUN_WORKER, CAPABILITY_DATAFLOW_TRIGGER_DISPATCHER, DataflowRunWorker, DataflowTriggerDispatcher, ) from govoplan_core.core.events import ( CAPABILITY_PLATFORM_EVENT_OUTBOX, DurableEventConsumer, PlatformEvent, PlatformEventOutbox, publish_platform_event, ) from govoplan_core.core.module_management import load_startup_enabled_modules, startup_candidate_module_ids from govoplan_core.core.modules import ModuleContext from govoplan_core.core.notifications import CAPABILITY_NOTIFICATIONS_DISPATCH, NotificationDispatchProvider from govoplan_core.core.postbox import ( CAPABILITY_POSTBOX_ROUTING, PostboxRoutingProvider, ) from govoplan_core.core.workflows import ( CAPABILITY_WORKFLOW_RUNTIME_WORKER, WorkflowRuntimeWorker, ) from govoplan_core.core.registry import PlatformRegistry from govoplan_core.core.runtime import configure_runtime from govoplan_core.settings import settings from govoplan_core.db.session import configure_database from govoplan_core.server.registry import available_module_manifests, build_platform_registry configure_database(settings.database_url) celery = Celery( "govoplan", broker=settings.redis_url, backend=settings.redis_url, ) celery.conf.update( task_default_queue="default", task_routes={ "govoplan.campaigns.send_email": {"queue": "send_email"}, "govoplan.campaigns.append_sent": {"queue": "append_sent"}, "govoplan.notifications.deliver": {"queue": "notifications"}, "govoplan.notifications.deliver_pending": {"queue": "notifications"}, "govoplan.calendar.dispatch_outbox": {"queue": "calendar"}, "govoplan.dataflow.dispatch_runs": {"queue": "dataflow"}, "govoplan.dataflow.purge_runs": {"queue": "dataflow"}, "govoplan.dataflow.dispatch_triggers": {"queue": "dataflow"}, "govoplan.workflow.reconcile": {"queue": "workflow"}, "govoplan.postbox.dispatch_routes": {"queue": "postbox"}, "govoplan.events.dispatch_outbox": {"queue": "events"}, "govoplan.events.purge_outbox": {"queue": "events"}, }, worker_prefetch_multiplier=1, task_acks_late=True, task_reject_on_worker_lost=True, beat_schedule={ "calendar-outbox-every-minute": { "task": "govoplan.calendar.dispatch_outbox", "schedule": 60.0, "args": (None, 100), }, "dataflow-triggers-every-minute": { "task": "govoplan.dataflow.dispatch_triggers", "schedule": 60.0, "args": (100,), }, "dataflow-runs-every-five-seconds": { "task": "govoplan.dataflow.dispatch_runs", "schedule": 5.0, "args": (10,), }, "dataflow-run-retention-daily": { "task": "govoplan.dataflow.purge_runs", "schedule": 24 * 60 * 60.0, "args": (500,), }, "workflow-reconcile-every-five-seconds": { "task": "govoplan.workflow.reconcile", "schedule": 5.0, "args": (50,), }, "postbox-routes-every-minute": { "task": "govoplan.postbox.dispatch_routes", "schedule": 60.0, "args": (None, 50), }, "platform-events-every-ten-seconds": { "task": "govoplan.events.dispatch_outbox", "schedule": 10.0, "args": (100,), }, "platform-event-retention-daily": { "task": "govoplan.events.purge_outbox", "schedule": 24 * 60 * 60.0, "args": (500,), }, }, ) @celery.task(name="govoplan.ping") def ping(): return "pong" def _platform_registry() -> PlatformRegistry: raw_enabled_modules = load_startup_enabled_modules(settings.enabled_modules) candidate_modules = startup_candidate_module_ids(settings.enabled_modules, raw_enabled_modules) available_modules = available_module_manifests(enabled_modules=candidate_modules, ignore_load_errors=True) enabled_modules = load_startup_enabled_modules(settings.enabled_modules, available=available_modules) registry = build_platform_registry(enabled_modules) context = ModuleContext(registry=registry, settings=settings) configure_runtime(context) registry.configure_capability_context(context) return registry def _campaign_delivery_tasks() -> CampaignDeliveryTaskProvider: registry = _platform_registry() capability = registry.require_capability(CAPABILITY_CAMPAIGNS_DELIVERY_TASKS) if not isinstance(capability, CampaignDeliveryTaskProvider): raise RuntimeError("Campaign delivery task capability is invalid") return capability def _notification_dispatch() -> NotificationDispatchProvider: registry = _platform_registry() capability = registry.require_capability(CAPABILITY_NOTIFICATIONS_DISPATCH) if not isinstance(capability, NotificationDispatchProvider): raise RuntimeError("Notification dispatch capability is invalid") return capability def _calendar_outbox() -> CalendarOutboxProvider | None: registry = _platform_registry() if not registry.has_capability(CAPABILITY_CALENDAR_OUTBOX): return None capability = registry.require_capability(CAPABILITY_CALENDAR_OUTBOX) if not isinstance(capability, CalendarOutboxProvider): raise RuntimeError("Calendar outbox capability is invalid") return capability def _dataflow_trigger_dispatcher( registry: PlatformRegistry | None = None, ) -> DataflowTriggerDispatcher | None: registry = registry or _platform_registry() if not registry.has_capability(CAPABILITY_DATAFLOW_TRIGGER_DISPATCHER): return None capability = registry.require_capability( CAPABILITY_DATAFLOW_TRIGGER_DISPATCHER ) if not isinstance(capability, DataflowTriggerDispatcher): raise RuntimeError("Dataflow trigger dispatcher capability is invalid") return capability def _dataflow_run_worker( registry: PlatformRegistry | None = None, ) -> DataflowRunWorker | None: registry = registry or _platform_registry() if not registry.has_capability(CAPABILITY_DATAFLOW_RUN_WORKER): return None capability = registry.require_capability(CAPABILITY_DATAFLOW_RUN_WORKER) if not isinstance(capability, DataflowRunWorker): raise RuntimeError("Dataflow run worker capability is invalid") return capability def _workflow_runtime_worker( registry: PlatformRegistry | None = None, ) -> WorkflowRuntimeWorker | None: registry = registry or _platform_registry() if not registry.has_capability(CAPABILITY_WORKFLOW_RUNTIME_WORKER): return None capability = registry.require_capability( CAPABILITY_WORKFLOW_RUNTIME_WORKER ) if not isinstance(capability, WorkflowRuntimeWorker): raise RuntimeError("Workflow runtime worker capability is invalid") return capability def _postbox_routing_provider( registry: PlatformRegistry | None = None, ) -> PostboxRoutingProvider | None: registry = registry or _platform_registry() if not registry.has_capability(CAPABILITY_POSTBOX_ROUTING): return None capability = registry.require_capability(CAPABILITY_POSTBOX_ROUTING) if not isinstance(capability, PostboxRoutingProvider): raise RuntimeError("Postbox routing capability is invalid") return capability def _platform_event_outbox( registry: PlatformRegistry | None = None, ) -> PlatformEventOutbox | None: registry = registry or _platform_registry() if not registry.has_capability(CAPABILITY_PLATFORM_EVENT_OUTBOX): return None capability = registry.require_capability(CAPABILITY_PLATFORM_EVENT_OUTBOX) if not isinstance(capability, PlatformEventOutbox): raise RuntimeError("Platform event outbox capability is invalid") return capability @celery.task(name="govoplan.campaigns.send_email", bind=True, max_retries=0) def send_email(self, job_id: str): """Send one explicitly queued campaign job. SMTP failures are persisted but are not retried implicitly. A worker-loss redelivery is safe because the delivery service converts an unfinished SMTP attempt into ``outcome_unknown`` instead of transmitting again. """ from govoplan_core.db.session import get_database with get_database().SessionLocal() as session: return dict(_campaign_delivery_tasks().send_campaign_job(session, job_id=job_id, enqueue_imap_task=True)) @celery.task(name="govoplan.campaigns.append_sent", bind=True, max_retries=None) def append_sent(self, job_id: str): """Append the exact sent MIME to the configured IMAP Sent folder.""" from govoplan_core.db.session import get_database with get_database().SessionLocal() as session: try: return dict(_campaign_delivery_tasks().append_sent_for_job(session, job_id=job_id)) except Exception as exc: if getattr(exc, "temporary", None) is True: raise self.retry(exc=exc, countdown=300) raise @celery.task(name="govoplan.notifications.deliver", bind=True, max_retries=0) def deliver_notification(self, notification_id: str): from govoplan_core.db.session import get_database with get_database().SessionLocal() as session: result = dict(_notification_dispatch().deliver_notification(session, notification_id=notification_id)) session.commit() return result @celery.task(name="govoplan.notifications.deliver_pending", bind=True, max_retries=0) def deliver_pending_notifications(self, tenant_id: str | None = None, limit: int = 50): from govoplan_core.db.session import get_database with get_database().SessionLocal() as session: result = dict(_notification_dispatch().deliver_pending(session, tenant_id=tenant_id, limit=limit)) session.commit() return result @celery.task(name="govoplan.calendar.dispatch_outbox", bind=True, max_retries=0) def dispatch_calendar_outbox(self, tenant_id: str | None = None, limit: int = 50): """Drain durable Calendar operations; retry timing lives in the database.""" from govoplan_core.db.session import get_database with get_database().SessionLocal() as session: provider = _calendar_outbox() if provider is None: return {"processed": 0, "succeeded": 0, "retrying": 0, "failed": 0, "operations": []} result = dict(provider.dispatch_due(session, tenant_id=tenant_id, limit=limit)) session.commit() return result @celery.task( name="govoplan.dataflow.dispatch_triggers", bind=True, max_retries=0, ) def dispatch_dataflow_triggers(self, limit: int = 100): """Drain durable Dataflow trigger deliveries and due schedules.""" from govoplan_core.db.session import get_database with get_database().SessionLocal() as session: provider = _dataflow_trigger_dispatcher() if provider is None: return { "queued": 0, "processed": 0, "succeeded": 0, "failed": 0, "blocked": 0, "skipped": 0, } result = dict(provider.dispatch_due(session, limit=limit)) session.commit() return result @celery.task( name="govoplan.dataflow.dispatch_runs", bind=True, max_retries=0, ) def dispatch_dataflow_runs(self, limit: int = 10): """Claim and execute durable Dataflow runs outside the API process.""" from govoplan_core.db.session import get_database with get_database().SessionLocal() as session: provider = _dataflow_run_worker() if provider is None: return { "claimed": 0, "succeeded": 0, "retrying": 0, "failed": 0, "cancelled": 0, } result = dict( provider.dispatch_pending( session, limit=limit, worker_id=getattr(self.request, "hostname", None), ) ) session.commit() return result @celery.task( name="govoplan.dataflow.purge_runs", bind=True, max_retries=0, ) def purge_dataflow_runs(self, limit: int = 500): """Apply Dataflow evidence-retention policy without deleting run records.""" from govoplan_core.db.session import get_database with get_database().SessionLocal() as session: provider = _dataflow_run_worker() if provider is None: return {"purged": 0} result = dict(provider.purge_expired(session, limit=limit)) session.commit() return result @celery.task( name="govoplan.workflow.reconcile", bind=True, max_retries=0, ) def reconcile_workflow_instances(self, limit: int = 50): """Resume asynchronous Workflow steps from durable provider state.""" from govoplan_core.db.session import get_database with get_database().SessionLocal() as session: provider = _workflow_runtime_worker() if provider is None: return { "inspected": 0, "advanced": 0, "waiting": 0, "failed": 0, } result = dict(provider.reconcile_pending(session, limit=limit)) session.commit() return result @celery.task( name="govoplan.postbox.dispatch_routes", bind=True, max_retries=0, ) def dispatch_postbox_routes( self, tenant_id: str | None = None, limit: int = 50, ): """Deliver due Postbox vacancy escalations from durable route rows.""" from govoplan_core.db.session import get_database with get_database().SessionLocal() as session: provider = _postbox_routing_provider() if provider is None: return { "selected": 0, "delivered": 0, "vacant": 0, "rescheduled": 0, "cancelled": 0, "failed": 0, "route_ids": [], } result = dict( provider.dispatch_due_routes( session, tenant_id=tenant_id, limit=limit, ) ) session.commit() return result @celery.task( name="govoplan.events.dispatch_outbox", bind=True, max_retries=0, ) def dispatch_platform_events(self, limit: int = 100): """Deliver committed platform events through persistent consumer ledgers.""" from govoplan_core.db.session import get_database with get_database().SessionLocal() as session: registry = _platform_registry() outbox = _platform_event_outbox(registry) if outbox is None: return { "selected": 0, "delivered": 0, "retrying": 0, "quarantined": 0, "dispatched": 0, "observer_failed": 0, } dataflow_dispatcher = _dataflow_trigger_dispatcher(registry) consumers = () if dataflow_dispatcher is not None: def deliver_to_dataflow( event: PlatformEvent, _delivery_key: str, ) -> None: dataflow_dispatcher.ingest_event( session, event=event, ) consumers = ( DurableEventConsumer( consumer_id="dataflow.event-triggers.v1", event_types=frozenset({"*"}), classifications=frozenset({"public", "internal"}), handler=deliver_to_dataflow, ), ) result = dict( outbox.dispatch_pending( session, consumers=consumers, observer=publish_platform_event, limit=limit, ) ) session.commit() return result @celery.task( name="govoplan.events.purge_outbox", bind=True, max_retries=0, ) def purge_platform_events(self, limit: int = 500): """Remove old terminal event envelopes while retaining quarantine evidence.""" from govoplan_core.db.session import get_database with get_database().SessionLocal() as session: outbox = _platform_event_outbox() if outbox is None: return {"deleted": 0} before = datetime.now(timezone.utc) - timedelta( days=settings.platform_event_outbox_terminal_retention_days ) result = dict( outbox.purge_terminal( session, before=before, limit=limit, ) ) session.commit() return result