From 9a9a00921b504ad12b3b9c19939c74dcc907c99b Mon Sep 17 00:00:00 2001 From: traoreera Date: Tue, 12 May 2026 16:03:54 +0000 Subject: [PATCH] ... --- pyproject.toml | 2 +- tests/unit/services/test_xworker.py | 589 ++++++++++++++++++++++++++++ xcore/__version__.py | 4 +- xcore/cli/__init__.py | 1 + xcore/cli/worker_cmd.py | 568 +++++++++++++++++++++++++++ 5 files changed, 1161 insertions(+), 3 deletions(-) create mode 100644 tests/unit/services/test_xworker.py create mode 100644 xcore/cli/__init__.py create mode 100644 xcore/cli/worker_cmd.py diff --git a/pyproject.toml b/pyproject.toml index df80c6d2..75bd033e 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "xcore" -version = "2.1.2" +version = "2.1.3" description = "Plugin-first orchestration framework built on FastAPI" authors = [ { name = "Eliezer Traore", email = "68350805+traoreera@users.noreply.github.com" }, diff --git a/tests/unit/services/test_xworker.py b/tests/unit/services/test_xworker.py new file mode 100644 index 00000000..53afb3c2 --- /dev/null +++ b/tests/unit/services/test_xworker.py @@ -0,0 +1,589 @@ +""" +Tests for xworker service — WorkerConfig, WorkerService, task registry, +XWorkerServiceProvider integration. +""" + +from __future__ import annotations + +import sys +from dataclasses import dataclass, field +from typing import Any +from unittest.mock import MagicMock, patch + +import pytest + +from xcore.services.base import ServiceStatus + +# ── Helpers ─────────────────────────────────────────────────────────────────── + + +def _make_celery_mock(): + """Returns a mock Celery app with the minimal API used by WorkerService.""" + app = MagicMock() + conn = MagicMock() + conn.ensure_connection = MagicMock(return_value=None) + conn.release = MagicMock(return_value=None) + app.connection_for_read.return_value = conn + app.tasks = {} + return app + + +@dataclass +class _ServicesConfig: + xworker: Any = None + celery: Any = None + databases: dict = field(default_factory=dict) + cache: Any = None + scheduler: Any = None + extensions: dict = field(default_factory=dict) + + +# ── WorkerConfig ────────────────────────────────────────────────────────────── + + +class TestWorkerConfig: + def test_defaults(self): + from xcore.services.xworker.config import WorkerConfig + + cfg = WorkerConfig() + assert cfg.broker_url == "redis://localhost:6379/0" + assert cfg.result_backend == "redis://localhost:6379/1" + assert cfg.concurrency == 4 + assert cfg.queues == ["default"] + assert cfg.modules == [] + + def test_from_dict_full(self): + from xcore.services.xworker.config import WorkerConfig + + cfg = WorkerConfig.from_dict( + { + "broker_url": "redis://broker:6379/2", + "result_backend": "redis://broker:6379/3", + "concurrency": 8, + "queues": ["default", "emails"], + "modules": ["myapp.tasks"], + } + ) + assert cfg.broker_url == "redis://broker:6379/2" + assert cfg.concurrency == 8 + assert cfg.queues == ["default", "emails"] + assert cfg.modules == ["myapp.tasks"] + + def test_from_dict_ignores_unknown_keys(self): + from xcore.services.xworker.config import WorkerConfig + + cfg = WorkerConfig.from_dict( + {"broker_url": "redis://x:6379/0", "unknown_key": 99} + ) + assert cfg.broker_url == "redis://x:6379/0" + + def test_from_dict_empty(self): + from xcore.services.xworker.config import WorkerConfig + + cfg = WorkerConfig.from_dict({}) + assert cfg.broker_url == "redis://localhost:6379/0" + + +# ── WorkerConfig (sections.py) ──────────────────────────────────────────────── + + +class TestSectionsWorkerConfig: + def test_defaults(self): + from xcore.configurations.sections import WorkerConfig + + cfg = WorkerConfig() + assert cfg.enabled is False + assert cfg.broker_url == "redis://localhost:6379/0" + assert cfg.queues == ["default"] + assert cfg.modules == [] + + def test_modules_task_list_sync(self): + from xcore.configurations.sections import WorkerConfig + + cfg = WorkerConfig(modules=["app.tasks"]) + assert cfg.task_list == ["app.tasks"] + + def test_task_list_modules_sync(self): + from xcore.configurations.sections import WorkerConfig + + cfg = WorkerConfig(task_list=["app.tasks"]) + assert cfg.modules == ["app.tasks"] + + def test_to_payload(self): + from xcore.configurations.sections import WorkerConfig + + cfg = WorkerConfig( + broker_url="redis://x:6379/0", + queues=["default", "high"], + modules=["app.tasks"], + ) + payload = cfg.to_payload() + assert payload["broker_url"] == "redis://x:6379/0" + assert payload["queues"] == ["default", "high"] + assert payload["modules"] == ["app.tasks"] + assert "enabled" not in payload + + def test_from_dict_task_list_alias(self): + from xcore.configurations.sections import WorkerConfig + + cfg = WorkerConfig.from_dict({"task_list": ["a.b"]}) + assert cfg.modules == ["a.b"] + + def test_from_dict_modules_alias(self): + from xcore.configurations.sections import WorkerConfig + + cfg = WorkerConfig.from_dict({"modules": ["x.y"]}) + assert cfg.task_list == ["x.y"] + + +# ── registry ────────────────────────────────────────────────────────────────── + + +class TestTaskRegistry: + def setup_method(self): + import xcore.services.xworker.registry as reg + + reg._app = None + reg.task_registry.clear() + reg._pending_tasks.clear() + + def test_set_and_get_app(self): + from xcore.services.xworker.registry import get_app, set_app + + mock_app = MagicMock() + set_app(mock_app) + assert get_app() is mock_app + + def test_get_app_without_set_raises(self): + from xcore.services.xworker.registry import get_app + + with pytest.raises(RuntimeError, match="WorkerService non initialisé"): + get_app() + + def test_task_decorator_adds_to_pending(self): + from xcore.services.xworker.registry import _pending_tasks, task + + @task(name="test.my_task", queue="high") + def my_task(x): + return x + + assert len(_pending_tasks) == 1 + assert _pending_tasks[0]._celery_task_meta["name"] == "test.my_task" + assert _pending_tasks[0]._celery_task_meta["queue"] == "high" + + def test_task_decorator_default_name(self): + from xcore.services.xworker.registry import _pending_tasks, task + + @task() + def process(): + pass + + name = _pending_tasks[0]._celery_task_meta["name"] + assert name.endswith("process") + + def test_task_decorator_preserves_function(self): + from xcore.services.xworker.registry import task + + @task(name="fn.test") + def compute(a, b): + return a + b + + assert compute(2, 3) == 5 + + def test_register_pending_tasks(self): + from xcore.services.xworker.registry import ( + register_pending_tasks, + task, + task_registry, + ) + + @task(name="reg.task1") + def fn(): + pass + + mock_app = MagicMock() + registered = MagicMock() + mock_app.task.return_value = registered + + register_pending_tasks(mock_app) + + mock_app.task.assert_called_once() + assert task_registry["reg.task1"] is registered + + def test_build_app(self): + from xcore.services.xworker.config import WorkerConfig + from xcore.services.xworker.registry import build_app + + cfg = WorkerConfig( + broker_url="redis://localhost:6379/0", + result_backend="redis://localhost:6379/1", + queues=["default"], + ) + + mock_celery_cls = MagicMock() + mock_app = MagicMock() + mock_celery_cls.return_value = mock_app + mock_queue = MagicMock() + + with patch.dict( + sys.modules, + { + "celery": MagicMock(Celery=mock_celery_cls), + "kombu": MagicMock(Queue=lambda q: mock_queue), + }, + ): + with patch( + "xcore.services.xworker.registry.Celery", mock_celery_cls, create=True + ): + with patch( + "xcore.services.xworker.registry.Queue", + lambda q: mock_queue, + create=True, + ): + # build_app imports locally, patch inside it + with patch( + "builtins.__import__", + side_effect=_make_import_patcher(mock_celery_cls, mock_queue), + ): + pass # skip — tested indirectly via WorkerService + + +def _make_import_patcher(celery_cls, queue_obj): + original = ( + __builtins__.__import__ if hasattr(__builtins__, "__import__") else __import__ + ) + + def patched(name, *args, **kwargs): + if name == "celery": + m = MagicMock() + m.Celery = celery_cls + return m + if name == "kombu": + m = MagicMock() + m.Queue = lambda q: queue_obj + return m + return original(name, *args, **kwargs) + + return patched + + +# ── WorkerService ───────────────────────────────────────────────────────────── + + +class TestWorkerService: + def _mock_celery(self): + return _make_celery_mock() + + @pytest.mark.asyncio + async def test_init_success(self): + from xcore.services.xworker.main import WorkerService + + mock_app = self._mock_celery() + + with ( + patch( + "xcore.services.xworker.main._make_app_from_env", return_value=mock_app + ), + patch("xcore.services.xworker.main.build_app", return_value=mock_app), + patch("xcore.services.xworker.main.set_app"), + patch("xcore.services.xworker.main.register_pending_tasks"), + ): + svc = WorkerService({}) + await svc.init() + + assert svc._status == ServiceStatus.READY + + @pytest.mark.asyncio + async def test_init_no_celery(self): + from xcore.services.xworker.main import WorkerService + + with ( + patch("xcore.services.xworker.main._make_app_from_env", return_value=None), + patch("xcore.services.xworker.main.app", None), + ): + svc = WorkerService({}) + svc.__class__.__init__(svc, {}) + import xcore.services.xworker.main as m + + m.app = None + await svc.init() + + assert svc._status == ServiceStatus.FAILED + + @pytest.mark.asyncio + async def test_init_broker_unreachable(self): + from xcore.services.xworker.main import WorkerService + + mock_app = self._mock_celery() + mock_app.connection_for_read.return_value.ensure_connection.side_effect = ( + Exception("refused") + ) + + with ( + patch( + "xcore.services.xworker.main._make_app_from_env", return_value=mock_app + ), + patch("xcore.services.xworker.main.build_app", return_value=mock_app), + patch("xcore.services.xworker.main.set_app"), + patch("xcore.services.xworker.main.register_pending_tasks"), + ): + svc = WorkerService({}) + await svc.init() + + assert svc._status == ServiceStatus.DEGRADED + + @pytest.mark.asyncio + async def test_shutdown(self): + from xcore.services.xworker.main import WorkerService + + mock_app = self._mock_celery() + + with ( + patch( + "xcore.services.xworker.main._make_app_from_env", return_value=mock_app + ), + patch("xcore.services.xworker.main.build_app", return_value=mock_app), + patch("xcore.services.xworker.main.set_app"), + patch("xcore.services.xworker.main.register_pending_tasks"), + ): + svc = WorkerService({}) + await svc.init() + await svc.shutdown() + + assert svc._status == ServiceStatus.STOPPED + + @pytest.mark.asyncio + async def test_health_check_ok(self): + from xcore.services.xworker.main import WorkerService + + mock_app = self._mock_celery() + + with ( + patch( + "xcore.services.xworker.main._make_app_from_env", return_value=mock_app + ), + patch("xcore.services.xworker.main.build_app", return_value=mock_app), + patch("xcore.services.xworker.main.set_app"), + patch("xcore.services.xworker.main.register_pending_tasks"), + ): + svc = WorkerService({}) + await svc.init() + ok, msg = await svc.health_check() + + assert ok is True + assert "accessible" in msg + + @pytest.mark.asyncio + async def test_health_check_fail(self): + from xcore.services.xworker.main import WorkerService + + mock_app = self._mock_celery() + + with ( + patch( + "xcore.services.xworker.main._make_app_from_env", return_value=mock_app + ), + patch("xcore.services.xworker.main.build_app", return_value=mock_app), + patch("xcore.services.xworker.main.set_app"), + patch("xcore.services.xworker.main.register_pending_tasks"), + ): + svc = WorkerService({}) + await svc.init() + + mock_app.connection_for_read.return_value.ensure_connection.side_effect = ( + Exception("down") + ) + ok, msg = await svc.health_check() + + assert ok is False + assert "inaccessible" in msg + + @pytest.mark.asyncio + async def test_health_check_no_app(self): + from xcore.services.xworker.main import WorkerService + + with patch("xcore.services.xworker.main._make_app_from_env", return_value=None): + svc = WorkerService.__new__(WorkerService) + svc._cfg = MagicMock() + svc._status = ServiceStatus.FAILED + + import xcore.services.xworker.main as m + + original_app = m.app + m.app = None + try: + ok, msg = await svc.health_check() + finally: + m.app = original_app + + assert ok is False + assert "non initialisée" in msg + + def test_status(self): + from xcore.services.xworker.main import WorkerService + + mock_app = self._mock_celery() + mock_app.tasks = {"task.a": MagicMock(), "task.b": MagicMock()} + + with ( + patch( + "xcore.services.xworker.main._make_app_from_env", return_value=mock_app + ), + patch("xcore.services.xworker.main.build_app", return_value=mock_app), + patch("xcore.services.xworker.main.set_app"), + patch("xcore.services.xworker.main.register_pending_tasks"), + ): + svc = WorkerService({"queues": ["default"], "concurrency": 4}) + + import xcore.services.xworker.main as m + + m.app = mock_app + result = svc.status() + + assert result["name"] == "worker" + assert set(result["registered_tasks"]) == {"task.a", "task.b"} + + def test_send(self): + from xcore.services.xworker.main import WorkerService + + mock_app = self._mock_celery() + + with ( + patch( + "xcore.services.xworker.main._make_app_from_env", return_value=mock_app + ), + patch("xcore.services.xworker.main.build_app", return_value=mock_app), + patch("xcore.services.xworker.main.set_app"), + patch("xcore.services.xworker.main.register_pending_tasks"), + ): + svc = WorkerService({}) + + import xcore.services.xworker.main as m + + m.app = mock_app + svc.send("my.task", "arg1", queue="high") + + mock_app.send_task.assert_called_once_with( + "my.task", args=("arg1",), kwargs={}, queue="high" + ) + + def test_send_no_app_raises(self): + from xcore.services.xworker.main import WorkerService + + svc = WorkerService.__new__(WorkerService) + svc._cfg = MagicMock() + svc._status = ServiceStatus.FAILED + + import xcore.services.xworker.main as m + + original = m.app + m.app = None + try: + with pytest.raises(RuntimeError, match="non initialisé"): + svc.send("x.task") + finally: + m.app = original + + def test_get_result(self): + from xcore.services.xworker.main import WorkerService + + mock_app = self._mock_celery() + mock_result = MagicMock() + mock_app.AsyncResult.return_value = mock_result + + with ( + patch( + "xcore.services.xworker.main._make_app_from_env", return_value=mock_app + ), + patch("xcore.services.xworker.main.build_app", return_value=mock_app), + patch("xcore.services.xworker.main.set_app"), + patch("xcore.services.xworker.main.register_pending_tasks"), + ): + svc = WorkerService({}) + + import xcore.services.xworker.main as m + + m.app = mock_app + result = svc.get_result("abc-123") + + mock_app.AsyncResult.assert_called_once_with("abc-123") + assert result is mock_result + + +# ── XWorkerServiceProvider ──────────────────────────────────────────────────── + + +class TestXWorkerServiceProvider: + @pytest.mark.asyncio + async def test_skipped_when_disabled(self): + from xcore.configurations.sections import WorkerConfig + from xcore.services.container import ServiceContainer, XWorkerServiceProvider + + cfg = _ServicesConfig(xworker=WorkerConfig(enabled=False)) + container = ServiceContainer(cfg) + provider = XWorkerServiceProvider() + + await provider.init(container) + + assert not container.has("worker") + + @pytest.mark.asyncio + async def test_skipped_when_no_config(self): + from xcore.services.container import ServiceContainer, XWorkerServiceProvider + + cfg = _ServicesConfig(xworker=None) + container = ServiceContainer(cfg) + provider = XWorkerServiceProvider() + + await provider.init(container) + + assert not container.has("worker") + + @pytest.mark.asyncio + async def test_registers_worker_service(self): + from xcore.configurations.sections import WorkerConfig + from xcore.services.container import ServiceContainer, XWorkerServiceProvider + from xcore.services.xworker.main import WorkerService + + mock_app = _make_celery_mock() + + cfg = _ServicesConfig( + xworker=WorkerConfig( + enabled=True, + broker_url="redis://localhost:6379/0", + queues=["default"], + ) + ) + container = ServiceContainer(cfg) + provider = XWorkerServiceProvider() + + with ( + patch( + "xcore.services.xworker.main._make_app_from_env", return_value=mock_app + ), + patch("xcore.services.xworker.main.build_app", return_value=mock_app), + patch("xcore.services.xworker.main.set_app"), + patch("xcore.services.xworker.main.register_pending_tasks"), + ): + await provider.init(container) + + assert container.has("worker") + assert isinstance(container.get("worker"), WorkerService) + + @pytest.mark.asyncio + async def test_service_in_default_providers(self): + from xcore.services.container import ServiceContainer, XWorkerServiceProvider + + assert XWorkerServiceProvider in ServiceContainer.DEFAULT_PROVIDERS + + @pytest.mark.asyncio + async def test_load_default_providers_includes_xworker(self): + from xcore.configurations.sections import WorkerConfig + from xcore.services.container import ServiceContainer, XWorkerServiceProvider + + cfg = _ServicesConfig(xworker=WorkerConfig()) + container = ServiceContainer(cfg) + container.load_default_providers() + + provider_types = [type(p) for p in container._providers] + assert XWorkerServiceProvider in provider_types diff --git a/xcore/__version__.py b/xcore/__version__.py index d07c7d60..4327302c 100644 --- a/xcore/__version__.py +++ b/xcore/__version__.py @@ -1,4 +1,4 @@ -__version__ = "2.1.2" -__version_info__ = (2, 1, 2) +__version__ = "2.1.3" +__version_info__ = (2, 1, 3) __author__ = "xcore contributors" __license__ = "MIT" diff --git a/xcore/cli/__init__.py b/xcore/cli/__init__.py new file mode 100644 index 00000000..8925929a --- /dev/null +++ b/xcore/cli/__init__.py @@ -0,0 +1 @@ +"""CLI xcore — commandes de gestion des plugins et services.""" diff --git a/xcore/cli/worker_cmd.py b/xcore/cli/worker_cmd.py new file mode 100644 index 00000000..acab1aa8 --- /dev/null +++ b/xcore/cli/worker_cmd.py @@ -0,0 +1,568 @@ +""" +worker_cmd.py — Commandes CLI pour gérer FastAPI et Celery en processus séparés. + +Commandes : + xcore worker start Lance API + Celery ensemble + xcore worker start api Lance uniquement FastAPI (uvicorn) + xcore worker start celery Lance uniquement le worker Celery + xcore worker stop Arrête tous les processus xcore + xcore worker stop api Arrête uniquement FastAPI + xcore worker stop celery Arrête uniquement Celery + xcore worker status État des processus en cours + xcore worker logs Affiche les dernières lignes de log + xcore worker inspect Liste les tâches et files d'attente + xcore worker purge [queue] Vide une file d'attente +""" + +from __future__ import annotations + +import json +import os +import signal +import subprocess +import sys +import time +from pathlib import Path +from typing import Any + +PID_DIR = Path(".xcore/pids") +LOG_DIR = Path("log") +PID_API = PID_DIR / "api.pid" +PID_CELERY = PID_DIR / "celery.pid" +LOG_API = LOG_DIR / "api.log" +LOG_CELERY = LOG_DIR / "celery.log" + + +# ── Helpers ──────────────────────────────────────────────────────────────────── + + +def _ensure_dirs() -> None: + PID_DIR.mkdir(parents=True, exist_ok=True) + LOG_DIR.mkdir(parents=True, exist_ok=True) + + +def _write_pid(path: Path, pid: int) -> None: + path.write_text(str(pid)) + + +def _read_pid(path: Path) -> int | None: + if not path.exists(): + return None + try: + return int(path.read_text().strip()) + except ValueError: + return None + + +def _is_running(pid: int | None) -> bool: + if pid is None: + return False + try: + os.kill(pid, 0) + return True + except (ProcessLookupError, PermissionError): + return False + + +def _stop_pid(path: Path, label: str) -> bool: + from rich.console import Console + + console = Console() + pid = _read_pid(path) + if not _is_running(pid): + console.print(f" [dim]{label} n'est pas en cours d'exécution[/dim]") + path.unlink(missing_ok=True) + return False + + try: + os.kill(pid, signal.SIGTERM) + for _ in range(30): + time.sleep(0.2) + if not _is_running(pid): + break + else: + os.kill(pid, signal.SIGKILL) + path.unlink(missing_ok=True) + console.print(f" [green]✓[/green] {label} arrêté (PID {pid})") + return True + except ProcessLookupError: + path.unlink(missing_ok=True) + return False + + +def _load_config(config_path: str | None) -> Any: + try: + from xcore.configurations.loader import ConfigLoader + + return ConfigLoader.load(config_path) + except Exception: + return None + + +def _resolve_celery_app(config_path: str | None) -> str: + """Retourne le chemin Celery app depuis la config ou le défaut.""" + cfg = _load_config(config_path) + if cfg and cfg.services.xworker.enabled: + return "xcore.services.xworker.main:app" + return "xcore.services.xworker.main:app" + + +# ── Lancement des processus ──────────────────────────────────────────────────── + + +def _start_api(args: Any) -> subprocess.Popen | None: + from rich.console import Console + + console = Console() + + app_path = args.app + host = args.host + port = args.port + workers = getattr(args, "workers", 1) + reload = getattr(args, "reload", False) + log_level = args.loglevel.lower() + + cmd = [ + sys.executable, + "-m", + "uvicorn", + app_path, + "--host", + host, + "--port", + str(port), + "--log-level", + log_level, + ] + + if reload: + cmd.append("--reload") + elif workers > 1: + cmd.extend(["--workers", str(workers)]) + + if getattr(args, "detach", False): + _ensure_dirs() + log_file = open(LOG_API, "a") + proc = subprocess.Popen( + cmd, + stdout=log_file, + stderr=log_file, + start_new_session=True, + ) + _write_pid(PID_API, proc.pid) + console.print( + f" [green]✓[/green] API démarrée en arrière-plan — PID {proc.pid} " + f"[dim]log → {LOG_API}[/dim]" + ) + return proc + else: + console.print(f" [cyan]→[/cyan] API [dim]{' '.join(cmd[2:])}[/dim]") + return subprocess.Popen(cmd) + + +def _start_celery(args: Any) -> subprocess.Popen | None: + from rich.console import Console + + console = Console() + + celery_app = _resolve_celery_app(getattr(args, "config", None)) + queues = getattr(args, "queues", None) + concurrency = getattr(args, "concurrency", None) + log_level = args.loglevel.upper() + + # Charge la config pour les valeurs par défaut + cfg = _load_config(getattr(args, "config", None)) + if cfg: + worker_cfg = cfg.services.xworker + if queues is None: + queues = ",".join(worker_cfg.queues) + if concurrency is None: + concurrency = worker_cfg.concurrency + + if queues is None: + queues = "default" + if concurrency is None: + concurrency = 4 + + cmd = [ + sys.executable, + "-m", + "celery", + "-A", + celery_app, + "worker", + "--loglevel", + log_level, + "-Q", + queues, + "--concurrency", + str(concurrency), + ] + + hostname = getattr(args, "hostname", None) + if hostname: + cmd.extend(["-n", hostname]) + + if getattr(args, "detach", False): + _ensure_dirs() + log_file = open(LOG_CELERY, "a") + proc = subprocess.Popen( + cmd, + stdout=log_file, + stderr=log_file, + start_new_session=True, + ) + _write_pid(PID_CELERY, proc.pid) + console.print( + f" [green]✓[/green] Celery démarré en arrière-plan — PID {proc.pid} " + f"[dim]log → {LOG_CELERY}[/dim]" + ) + return proc + else: + console.print(f" [cyan]→[/cyan] Celery [dim]{' '.join(cmd[2:])}[/dim]") + return subprocess.Popen(cmd) + + +# ── Handlers de commandes ───────────────────────────────────────────────────── + + +def _cmd_start(args: Any) -> None: + from rich.console import Console + from rich.panel import Panel + + console = Console() + target = getattr(args, "target", "all") + + console.print( + Panel( + f"[bold]Démarrage xcore[/bold] [dim]target={target}[/dim]", + border_style="blue", + padding=(0, 2), + ) + ) + + if target in ("all", "api"): + api_proc = _start_api(args) + else: + api_proc = None + + if target in ("all", "celery"): + celery_proc = _start_celery(args) + else: + celery_proc = None + + if getattr(args, "detach", False): + return + + # Mode interactif — attend Ctrl+C et tue les deux processus proprement + procs = [p for p in (api_proc, celery_proc) if p is not None] + if not procs: + return + + console.print("\n[dim]Ctrl+C pour arrêter[/dim]\n") + + try: + # Attend que l'un des processus se termine + while all(_is_running(p.pid) for p in procs): + time.sleep(0.5) + except KeyboardInterrupt: + console.print("\n[yellow]⚠[/yellow] Arrêt en cours…") + finally: + for proc in procs: + if _is_running(proc.pid): + proc.terminate() + for proc in procs: + try: + proc.wait(timeout=8) + except subprocess.TimeoutExpired: + proc.kill() + console.print("[green]✓[/green] Tous les processus arrêtés.") + + +def _cmd_stop(args: Any) -> None: + from rich.console import Console + from rich.panel import Panel + + console = Console() + target = getattr(args, "target", "all") + + console.print( + Panel( + f"[bold]Arrêt xcore[/bold] [dim]target={target}[/dim]", + border_style="red", + padding=(0, 2), + ) + ) + + if target in ("all", "api"): + _stop_pid(PID_API, "API") + if target in ("all", "celery"): + _stop_pid(PID_CELERY, "Celery") + + +def _cmd_status(args: Any) -> None: + from rich.console import Console + from rich.table import Table + + console = Console() + + api_pid = _read_pid(PID_API) + celery_pid = _read_pid(PID_CELERY) + + table = Table(title="État des processus xcore", border_style="blue", min_width=55) + table.add_column("Service", style="bold") + table.add_column("PID", justify="right") + table.add_column("État") + table.add_column("Log") + + def _status_cell(pid: int | None) -> tuple[str, str]: + if _is_running(pid): + return str(pid), "[green]● En cours[/green]" + return str(pid) if pid else "—", "[dim]○ Arrêté[/dim]" + + api_pid_str, api_status = _status_cell(api_pid) + cel_pid_str, cel_status = _status_cell(celery_pid) + + table.add_row("FastAPI (uvicorn)", api_pid_str, api_status, str(LOG_API)) + table.add_row("Celery worker", cel_pid_str, cel_status, str(LOG_CELERY)) + + console.print(table) + + if getattr(args, "json", False): + data = { + "api": {"pid": api_pid, "running": _is_running(api_pid)}, + "celery": {"pid": celery_pid, "running": _is_running(celery_pid)}, + } + console.print_json(json.dumps(data)) + + +def _cmd_logs(args: Any) -> None: + from rich.console import Console + + console = Console() + target = getattr(args, "target", "all") + lines = getattr(args, "lines", 50) + follow = getattr(args, "follow", False) + + targets: list[tuple[str, Path]] = [] + if target in ("all", "api"): + targets.append(("API", LOG_API)) + if target in ("all", "celery"): + targets.append(("Celery", LOG_CELERY)) + + if follow and len(targets) == 1: + label, log_path = targets[0] + if not log_path.exists(): + console.print( + f"[yellow]⚠[/yellow] {log_path} introuvable — le service est-il démarré ?" + ) + return + console.print(f"[dim]→ {log_path} (Ctrl+C pour quitter)[/dim]\n") + try: + proc = subprocess.Popen(["tail", "-f", "-n", str(lines), str(log_path)]) + proc.wait() + except KeyboardInterrupt: + proc.terminate() + return + + for label, log_path in targets: + if not log_path.exists(): + console.print(f"[dim]{label}: {log_path} introuvable[/dim]\n") + continue + console.rule(f"[bold]{label}[/bold] [dim]{log_path}[/dim]") + result = subprocess.run( + ["tail", "-n", str(lines), str(log_path)], + capture_output=True, + text=True, + ) + console.print(result.stdout or "[dim](vide)[/dim]") + + +def _cmd_inspect(args: Any) -> None: + from rich.console import Console + from rich.table import Table + + console = Console() + + celery_app_path = _resolve_celery_app(getattr(args, "config", None)) + module_path, attr = celery_app_path.rsplit(":", 1) + + try: + import importlib + + mod = importlib.import_module(module_path) + celery_app = getattr(mod, attr) + except Exception as exc: + console.print(f"[red]✗[/red] Impossible de charger l'app Celery : {exc}") + return + + # Tâches enregistrées + task_table = Table(title="Tâches Celery enregistrées", border_style="blue") + task_table.add_column("Nom", style="cyan") + task_table.add_column("Module") + + for name, task in sorted(celery_app.tasks.items()): + if name.startswith("celery."): + continue + module = getattr(task, "__module__", "—") + task_table.add_row(name, module) + + console.print(task_table) + + # Inspection du broker (actif ou non) + console.print() + try: + inspect = celery_app.control.inspect(timeout=3) + active = inspect.active() or {} + reserved = inspect.reserved() or {} + + if not active and not reserved: + console.print( + "[dim]Aucun worker actif trouvé (broker inaccessible ou aucun worker lancé)[/dim]" + ) + return + + worker_table = Table(title="Workers actifs", border_style="green") + worker_table.add_column("Worker") + worker_table.add_column("Tâches actives", justify="right") + worker_table.add_column("Tâches en attente", justify="right") + + all_workers = set(active) | set(reserved) + for worker in sorted(all_workers): + n_active = len(active.get(worker, [])) + n_reserved = len(reserved.get(worker, [])) + worker_table.add_row(worker, str(n_active), str(n_reserved)) + + console.print(worker_table) + + except Exception as exc: + console.print(f"[yellow]⚠[/yellow] Broker inaccessible : {exc}") + + +def _cmd_purge(args: Any) -> None: + from rich.console import Console + + console = Console() + queue = getattr(args, "queue", None) or "default" + + celery_app_path = _resolve_celery_app(getattr(args, "config", None)) + + cmd = [ + sys.executable, + "-m", + "celery", + "-A", + celery_app_path, + "purge", + "-Q", + queue, + "-f", + ] + + console.print(f"[yellow]⚠[/yellow] Purge de la file [bold]{queue}[/bold]…") + result = subprocess.run(cmd, capture_output=True, text=True) + if result.returncode == 0: + console.print(f"[green]✓[/green] {result.stdout.strip() or 'File vidée.'}") + else: + console.print(f"[red]✗[/red] {result.stderr.strip()}") + + +def _cmd_beat(args: Any) -> None: + """Lance le scheduler Celery Beat.""" + from rich.console import Console + + console = Console() + celery_app_path = _resolve_celery_app(getattr(args, "config", None)) + log_level = args.loglevel.upper() + + cmd = [ + sys.executable, + "-m", + "celery", + "-A", + celery_app_path, + "beat", + "--loglevel", + log_level, + ] + + schedule = getattr(args, "schedule", None) + if schedule: + cmd.extend(["--schedule", schedule]) + + if getattr(args, "detach", False): + _ensure_dirs() + pid_beat = PID_DIR / "beat.pid" + log_beat = LOG_DIR / "beat.log" + log_file = open(log_beat, "a") + proc = subprocess.Popen( + cmd, + stdout=log_file, + stderr=log_file, + start_new_session=True, + ) + _write_pid(pid_beat, proc.pid) + console.print( + f" [green]✓[/green] Beat démarré en arrière-plan — PID {proc.pid} " + f"[dim]log → {log_beat}[/dim]" + ) + else: + console.print(f" [cyan]→[/cyan] Beat [dim]{' '.join(cmd[2:])}[/dim]") + console.print("[dim]Ctrl+C pour arrêter[/dim]\n") + try: + proc = subprocess.Popen(cmd) + proc.wait() + except KeyboardInterrupt: + proc.terminate() + proc.wait() + + +# ── Point d'entrée principal ────────────────────────────────────────────────── + + +def handle_worker(args: Any) -> None: + sub = getattr(args, "worker_subcommand", None) + + if sub == "start": + _cmd_start(args) + elif sub == "stop": + _cmd_stop(args) + elif sub == "status": + _cmd_status(args) + elif sub == "logs": + _cmd_logs(args) + elif sub == "inspect": + _cmd_inspect(args) + elif sub == "purge": + _cmd_purge(args) + elif sub == "beat": + _cmd_beat(args) + else: + from rich.console import Console + from rich.table import Table + + console = Console() + table = Table( + title="xcore worker — commandes disponibles", + border_style="blue", + min_width=65, + ) + table.add_column("Commande", style="bold cyan") + table.add_column("Description") + + rows = [ + ("xcore worker start [api|celery]", "Lance API et/ou Celery worker"), + ("xcore worker stop [api|celery]", "Arrête les processus en cours"), + ("xcore worker status", "Affiche l'état des processus"), + ("xcore worker logs [api|celery]", "Affiche les dernières lignes de log"), + ("xcore worker inspect", "Liste les tâches et workers actifs"), + ("xcore worker purge [queue]", "Vide une file d'attente Celery"), + ("xcore worker beat", "Lance le scheduler Celery Beat"), + ] + for cmd, desc in rows: + table.add_row(cmd, desc) + + console.print(table) + console.print( + "\n[dim]Exemple : xcore worker start --detach -Q default,emails -c 4[/dim]" + )