from __future__ import annotations 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_TRIGGER_DISPATCHER, DataflowTriggerDispatcher, ) 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.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_triggers": {"queue": "dataflow"}, }, 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,), }, }, ) @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() -> DataflowTriggerDispatcher | None: registry = _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 @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