89 lines
2.5 KiB
Python
89 lines
2.5 KiB
Python
|
|
"""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...")
|
||
|
|
|
||
|
|
|
||
|
|
def _collect_job_functions() -> dict[str, Any]:
|
||
|
|
"""Collect all job functions from plugins and core.
|
||
|
|
|
||
|
|
Each plugin's jobs module may define async functions with the
|
||
|
|
signature (ctx: dict, *args) -> None. We import them all and
|
||
|
|
register by name.
|
||
|
|
"""
|
||
|
|
functions: dict[str, Any] = {}
|
||
|
|
|
||
|
|
# Unified search jobs
|
||
|
|
try:
|
||
|
|
from app.plugins.builtins.unified_search.jobs import (
|
||
|
|
index_mails,
|
||
|
|
index_file,
|
||
|
|
index_contact,
|
||
|
|
index_company,
|
||
|
|
index_event,
|
||
|
|
reindex,
|
||
|
|
embedding_batch,
|
||
|
|
)
|
||
|
|
functions.update({
|
||
|
|
"index_mails": index_mails,
|
||
|
|
"index_file": index_file,
|
||
|
|
"index_contact": index_contact,
|
||
|
|
"index_company": index_company,
|
||
|
|
"index_event": index_event,
|
||
|
|
"reindex": reindex,
|
||
|
|
"embedding_batch": embedding_batch,
|
||
|
|
})
|
||
|
|
except Exception:
|
||
|
|
logger.warning("Failed to import unified_search jobs", exc_info=True)
|
||
|
|
|
||
|
|
# AI proactive jobs
|
||
|
|
try:
|
||
|
|
from app.plugins.builtins.ai_proactive.jobs import deep_analysis
|
||
|
|
functions["deep_analysis"] = deep_analysis
|
||
|
|
except Exception:
|
||
|
|
logger.warning("Failed to import ai_proactive jobs", exc_info=True)
|
||
|
|
|
||
|
|
# Calendar jobs (calendar_reminder is enqueued but no handler exists yet)
|
||
|
|
# try:
|
||
|
|
# from app.plugins.builtins.calendar.jobs import calendar_reminder
|
||
|
|
# functions["calendar_reminder"] = calendar_reminder
|
||
|
|
# except Exception:
|
||
|
|
# logger.warning("Failed to import calendar jobs", exc_info=True)
|
||
|
|
|
||
|
|
return functions
|
||
|
|
|
||
|
|
|
||
|
|
class WorkerSettings:
|
||
|
|
"""ARQ worker settings."""
|
||
|
|
functions = _collect_job_functions()
|
||
|
|
redis_settings = _get_redis_settings()
|
||
|
|
on_startup = on_startup
|
||
|
|
on_shutdown = on_shutdown
|
||
|
|
max_jobs = 10
|
||
|
|
job_timeout = 300
|
||
|
|
queue_name = "arq:queue"
|