71 lines
1.8 KiB
Python
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")
|