Files
govoplan-core/src/govoplan_core/celery_app.py

262 lines
9.6 KiB
Python

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.events import (
CAPABILITY_PLATFORM_EVENT_OUTBOX,
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.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"},
"govoplan.events.dispatch_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,),
},
"platform-events-every-ten-seconds": {
"task": "govoplan.events.dispatch_outbox",
"schedule": 10.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(
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 _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.events.dispatch_outbox",
bind=True,
max_retries=0,
)
def dispatch_platform_events(self, limit: int = 100):
"""Deliver committed platform events once across all worker processes."""
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, "dispatched": 0, "failed": 0}
dataflow_dispatcher = _dataflow_trigger_dispatcher(registry)
def dispatch(event: PlatformEvent) -> None:
if (
dataflow_dispatcher is not None
and event.classification in {"public", "internal"}
):
dataflow_dispatcher.ingest_event(session, event=event)
publish_platform_event(event)
result = dict(
outbox.dispatch_pending(
session,
dispatcher=dispatch,
limit=limit,
)
)
session.commit()
return result