152 lines
6.1 KiB
Python
152 lines
6.1 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.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"},
|
|
},
|
|
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),
|
|
},
|
|
},
|
|
)
|
|
|
|
|
|
@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
|
|
|
|
|
|
@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
|