Files
leocrm/app/core/job_registry.py

71 lines
1.8 KiB
Python

"""
Job Registry — decouples ARQ worker from direct plugin imports.
Plugins register their job functions via register_job() at import time.
The worker retrieves all registered jobs via get_all_jobs() instead of
importing from plugin modules directly.
"""
from __future__ import annotations
import logging
from collections.abc import Callable, Coroutine
from typing import Any
logger = logging.getLogger(__name__)
# Type alias for an async job function
JobFunc = Callable[..., Coroutine[Any, Any, Any]]
# Internal registry: name -> job function
_registry: dict[str, JobFunc] = {}
def register_job(name: str, func: JobFunc) -> None:
"""Register a job function under the given name.
Args:
name: Unique job name (e.g. 'index_mails').
func: The async callable to register.
"""
if name in _registry:
logger.warning("Job '%s' is being re-registered — overwriting", name)
_registry[name] = func
logger.debug("Registered job: %s", name)
def get_job(name: str) -> JobFunc | None:
"""Retrieve a registered job function by name.
Args:
name: The job name to look up.
Returns:
The registered callable, or None if not found.
"""
return _registry.get(name)
def unregister_job(name: str) -> None:
"""Remove a registered job function (plugin deactivation lifecycle).
Args:
name: The job name to remove.
"""
_registry.pop(name, None)
def get_all_jobs() -> list[JobFunc]:
"""Return all registered job functions (order is insertion order).
Returns:
List of all registered async callables.
"""
return list(_registry.values())
def clear_registry() -> None:
"""Clear all registered jobs. Useful for testing."""
_registry.clear()
logger.debug("Job registry cleared")