Files
leocrm/app/core/worker.py
T

62 lines
1.4 KiB
Python
Raw Normal View History

"""ARQ worker configuration and entrypoint."""
from __future__ import annotations
import logging
from typing import Any
from arq.connections import RedisSettings
from app.config import get_settings
logger = logging.getLogger(__name__)
def _get_redis_settings() -> RedisSettings:
"""Get Redis settings from app config."""
settings = get_settings()
return RedisSettings.from_dsn(settings.redis_url)
async def on_startup(ctx: dict[str, Any]) -> None:
"""Called when worker starts."""
logger.info("ARQ worker starting...")
async def on_shutdown(ctx: dict[str, Any]) -> None:
"""Called when worker shuts down."""
logger.info("ARQ worker shutting down...")
# Import job functions directly so ARQ registers them by __name__
from app.plugins.builtins.unified_search.jobs import (
index_mails,
index_file,
index_contact,
index_company,
index_event,
reindex,
embedding_batch,
)
from app.plugins.builtins.ai_proactive.jobs import deep_analysis
class WorkerSettings:
"""ARQ worker settings."""
functions = [
index_mails,
index_file,
index_contact,
index_company,
index_event,
reindex,
embedding_batch,
deep_analysis,
]
redis_settings = _get_redis_settings()
on_startup = on_startup
on_shutdown = on_shutdown
max_jobs = 10
job_timeout = 300
queue_name = "arq:queue"