From c5a7866b53cde1a0ceed3800089863baa852d6da Mon Sep 17 00:00:00 2001 From: traoreera Date: Mon, 28 Sep 2026 13:54:52 +0000 Subject: [PATCH 1/2] =?UTF-8?q?feat(plugins):=20table=20de=20v=C3=A9rit?= =?UTF-8?q?=C3=A9=20persistante,=20ramasse-miette=20forc=C3=A9=20et=20supe?= =?UTF-8?q?rvision=20IPC/events?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Activation/désactivation des plugins : - PluginStateStore (JSON) persiste l'état actif/inactif, consulté par PluginLoader.load_all() pour sauter les plugins désactivés au boot - PluginSupervisor.enable()/disable(reason=) togglent en live + persistent - registry_table() expose la table de vérité complète (disque + persisté + état live), y compris les plugins désactivés ou jamais chargés Ramasse-miette forcé au unload/disable (PluginResourceTracker) : - jobs scheduler, health checks, abonnements events/hooks désinscrits de force, indépendamment de la qualité des hooks on_stop/on_unload du plugin - PluginRegistry.unregister() enfin appelé (existait, jamais invoqué) - ctx.spawn_task() pour des tâches de fond suivies et annulées au unload - routes HTTP démontées au unload (pas seulement au reload) Centre de contrôle HTTP sur /plugins/ipc/* : - GET /registry, POST /{name}/enable, POST /{name}/disable - GET /audit (PluginSupervisor.ipc_audit()/ipc_stats()) - GET /events (EventBus.recent_emissions()/stats(), HookManager.recent_emissions()) Corrections (vérification reports/technical_debt_remediation_verification_2026-08-11.md) : - propagate_services() cassait le reload d'un plugin TrustedBase (collision sur un service protégé reçu en injection, pas exporté) - SandboxProcessManager.call() ne recyclait pas le subprocess sur IPCTimeoutError, seulement sur IPCProcessDead - PermissionEngine(audit_cache_hits=False) pour sauter l'audit log sur cache hit (optionnel, défaut inchangé) - tests/conftest.py : préfixe xcore_test_ uniforme + fixture autouse de nettoyage des dossiers temporaires orphelins en cas de crash dur --- CHANGELOG.md | 16 ++ doc/changelog.md | 13 + tests/conftest.py | 29 +- .../integration/test_plugin_enable_disable.py | 125 +++++++++ tests/test_conftest_cleanup.py | 31 +++ tests/unit/kernel/test_api_router.py | 94 ++++++- tests/unit/kernel/test_events.py | 86 +++++- tests/unit/kernel/test_health.py | 16 ++ tests/unit/kernel/test_hooks.py | 48 ++++ tests/unit/kernel/test_lifecycle.py | 164 +++++++++++ tests/unit/kernel/test_loader.py | 59 ++++ tests/unit/kernel/test_permissions.py | 31 +++ tests/unit/kernel/test_process_manager.py | 32 ++- tests/unit/kernel/test_supervisor.py | 149 ++++++++++ tests/unit/test_registry_state_store.py | 49 ++++ xcore/__init__.py | 33 ++- xcore/kernel/api/context.py | 22 +- xcore/kernel/api/router.py | 54 ++++ xcore/kernel/events/bus.py | 91 ++++++- xcore/kernel/events/hooks.py | 40 ++- xcore/kernel/observability/health.py | 4 + xcore/kernel/permissions/engine.py | 20 +- xcore/kernel/runtime/lifecycle.py | 101 ++++++- xcore/kernel/runtime/loader.py | 45 +++- xcore/kernel/runtime/plugin_gc.py | 254 ++++++++++++++++++ xcore/kernel/runtime/supervisor.py | 155 ++++++++++- xcore/kernel/sandbox/process_manager.py | 9 +- xcore/registry/__init__.py | 9 +- xcore/registry/index.py | 8 + xcore/registry/state_store.py | 79 ++++++ 30 files changed, 1828 insertions(+), 38 deletions(-) create mode 100644 tests/integration/test_plugin_enable_disable.py create mode 100644 tests/test_conftest_cleanup.py create mode 100644 tests/unit/test_registry_state_store.py create mode 100644 xcore/kernel/runtime/plugin_gc.py create mode 100644 xcore/registry/state_store.py diff --git a/CHANGELOG.md b/CHANGELOG.md index e4d188a5..29db1fe8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,22 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [2.6.0] - 2026-09-28 + +### Added +- **Persistent plugin enable/disable state** (`PluginStateStore`, `xcore/registry/state_store.py`): until now there was no way to disable a plugin short of deleting its folder — the manifest schema's `enabled` field lived under `runtime.health_check`, not on the plugin itself. `PluginLoader.load_all()` now consults a JSON file (`/../.xcore/plugins_state.json`) and skips disabled plugins at boot; `PluginSupervisor.enable()`/`disable(reason=...)` toggle the state live *and* persist it, so a process restart honors the same active/inactive set. +- **Forced garbage collection on unload/disable** (`PluginResourceTracker`, `xcore/kernel/runtime/plugin_gc.py`): `LifecycleManager._do_unload()` used to trust only the plugin's own `on_stop`/`on_unload` hooks plus `sys.modules` cleanup — no scheduler job, health check, or event/hook subscription was ever unregistered, and `PluginRegistry.unregister()` existed but was never called. The kernel now wraps scheduler/health/events/hooks at load time to track what a plugin registers, and forces their release on unload regardless of how well the plugin's own hooks behave. New `ctx.spawn_task()` for background tasks that are tracked and cancelled automatically. +- **HTTP routes unmounted on unload**: `xcore/__init__.py` only stripped a plugin's FastAPI routes from `app.routes` on reload — never on unload/disable, leaving a disabled plugin's endpoints reachable indefinitely. Extracted into `_unmount_plugin_router()`, now also subscribed to `plugin.*.unloaded`. +- **HTTP control center** on `/plugins/ipc/*`: `GET /registry` (the full plugin truth table — including disabled or never-loaded plugins, which `status()` never exposed), `POST /{name}/enable`, `POST /{name}/disable`. +- **IPC call supervision**: `PluginSupervisor.ipc_audit()`/`ipc_stats()` log every call (`plugin`, `action`, `caller`, `tenant_id`, status, duration) to a bounded audit trail, mirroring the existing `PermissionEngine.audit_log()` pattern. Exposed via `GET /plugins/ipc/audit`. +- **Event/hook supervision**: `EventBus.recent_emissions()`/`.stats()` and `HookManager.recent_emissions()` keep a record of recent emissions (event, handlers/hooks matched, errors, duration) — `EventBus` previously had no metrics at all. Exposed via `GET /plugins/ipc/events`. + +### Fixed +- **`propagate_services()` broke reload/re-enable of a `TrustedBase` plugin**: `self._services` exposes the entire `ctx.services` dict for backward compatibility (including db/cache/scheduler), and `propagate_services()` tried to re-register those as the plugin's own exports. This passed on first boot (the registry doesn't protect core services until after `load_all()` runs), but any later reload raised `PermissionError: Impossible d'écraser le service protégé`. A collision on an object identical to the one already protected (received via injection, not exported) is now ignored; a genuinely different object (an actual override attempt) still raises. +- **Sandboxed subprocess wasn't recycled on IPC timeout**: `SandboxProcessManager.call()` only caught `IPCProcessDead`, not `IPCTimeoutError` — a subprocess that stopped responding without actually dying stayed unusable until the next periodic `_health_loop` check caught up with it. `call()` now recycles the subprocess immediately on either exception, on the failing request's own path, instead of waiting for the next health-check interval. +- **`PermissionEngine` audit log had no way to skip cache-hit entries**: every `allows()`/`check()` cache hit still appended to `_audit_log` unconditionally — the expensive part (event emission) was already skipped on cache hits, but the log append wasn't. New optional `PermissionEngine(audit_cache_hits=False)` skips it; default (`True`) keeps the existing behavior (complete audit trail) unchanged. +- **Stray temp directories from crashed test runs**: `tests/conftest.py`'s `plugins_dir`/`temp_dir` fixtures already clean up via `yield` + `shutil.rmtree`, but that teardown never runs if a test crashes hard (e.g. `SIGKILL`) before reaching it. `temp_dir` now uses the same distinguishing `xcore_test_` prefix as `plugins_dir`, and a new session-scoped autouse fixture sweeps any `xcore_test_*` directories left behind in the system temp dir at the end of the run — scoped to that exact prefix only, never a broader temp-dir sweep. + ## [2.5.1] - 2026-08-20 ### Fixed diff --git a/doc/changelog.md b/doc/changelog.md index e4d188a5..017e541e 100644 --- a/doc/changelog.md +++ b/doc/changelog.md @@ -5,6 +5,19 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [2.6.0] - 2026-09-28 + +### Added +- **Persistent plugin enable/disable state** (`PluginStateStore`, `xcore/registry/state_store.py`): until now there was no way to disable a plugin short of deleting its folder — the manifest schema's `enabled` field lived under `runtime.health_check`, not on the plugin itself. `PluginLoader.load_all()` now consults a JSON file (`/../.xcore/plugins_state.json`) and skips disabled plugins at boot; `PluginSupervisor.enable()`/`disable(reason=...)` toggle the state live *and* persist it, so a process restart honors the same active/inactive set. +- **Forced garbage collection on unload/disable** (`PluginResourceTracker`, `xcore/kernel/runtime/plugin_gc.py`): `LifecycleManager._do_unload()` used to trust only the plugin's own `on_stop`/`on_unload` hooks plus `sys.modules` cleanup — no scheduler job, health check, or event/hook subscription was ever unregistered, and `PluginRegistry.unregister()` existed but was never called. The kernel now wraps scheduler/health/events/hooks at load time to track what a plugin registers, and forces their release on unload regardless of how well the plugin's own hooks behave. New `ctx.spawn_task()` for background tasks that are tracked and cancelled automatically. +- **HTTP routes unmounted on unload**: `xcore/__init__.py` only stripped a plugin's FastAPI routes from `app.routes` on reload — never on unload/disable, leaving a disabled plugin's endpoints reachable indefinitely. Extracted into `_unmount_plugin_router()`, now also subscribed to `plugin.*.unloaded`. +- **HTTP control center** on `/plugins/ipc/*`: `GET /registry` (the full plugin truth table — including disabled or never-loaded plugins, which `status()` never exposed), `POST /{name}/enable`, `POST /{name}/disable`. +- **IPC call supervision**: `PluginSupervisor.ipc_audit()`/`ipc_stats()` log every call (`plugin`, `action`, `caller`, `tenant_id`, status, duration) to a bounded audit trail, mirroring the existing `PermissionEngine.audit_log()` pattern. Exposed via `GET /plugins/ipc/audit`. +- **Event/hook supervision**: `EventBus.recent_emissions()`/`.stats()` and `HookManager.recent_emissions()` keep a record of recent emissions (event, handlers/hooks matched, errors, duration) — `EventBus` previously had no metrics at all. Exposed via `GET /plugins/ipc/events`. + +### Fixed +- **`propagate_services()` broke reload/re-enable of a `TrustedBase` plugin**: `self._services` exposes the entire `ctx.services` dict for backward compatibility (including db/cache/scheduler), and `propagate_services()` tried to re-register those as the plugin's own exports. This passed on first boot (the registry doesn't protect core services until after `load_all()` runs), but any later reload raised `PermissionError: Impossible d'écraser le service protégé`. A collision on an object identical to the one already protected (received via injection, not exported) is now ignored; a genuinely different object (an actual override attempt) still raises. + ## [2.5.1] - 2026-08-20 ### Fixed diff --git a/tests/conftest.py b/tests/conftest.py index bac4fa90..621d9709 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -131,11 +131,38 @@ def fake_plugin_dir(plugins_dir: Path) -> Path: @pytest.fixture def temp_dir() -> Generator[Path, None, None]: """Fixture utilitaire pour un dossier temporaire propre.""" - tmp = tempfile.mkdtemp() + # Préfixe distinctif : permet à _cleanup_stray_tmp_dirs (ci-dessous) de + # cibler précisément ces dossiers sans risquer de toucher aux tmp d'un + # autre processus si ce yield+rmtree n'a pas pu s'exécuter (crash dur). + tmp = tempfile.mkdtemp(prefix="xcore_test_") yield Path(tmp) shutil.rmtree(tmp, ignore_errors=True) +def sweep_stray_tmp_dirs() -> None: + """ + Balaie tout ce qui porte le préfixe `xcore_test_` restant dans le dossier + temp système — jamais plus large, pour ne pas toucher aux fichiers + temporaires d'un autre processus. Fonction plain (pas une fixture) pour + rester testable directement. + """ + tmp_root = Path(tempfile.gettempdir()) + for leftover in tmp_root.glob("xcore_test_*"): + shutil.rmtree(leftover, ignore_errors=True) + + +@pytest.fixture(scope="session", autouse=True) +def _cleanup_stray_tmp_dirs() -> Generator[None, None, None]: + """ + Filet de sécurité en fin de session : `plugins_dir`/`temp_dir` nettoient + déjà systématiquement leur propre dossier via yield+rmtree, mais ce + teardown ne s'exécute pas si un test crashe durement (SIGKILL) avant + d'y arriver. + """ + yield + sweep_stray_tmp_dirs() + + @pytest.fixture def invalid_plugin_dir(plugins_dir: Path) -> Path: """Plugin invalide — manque PLUGIN_INFO.""" diff --git a/tests/integration/test_plugin_enable_disable.py b/tests/integration/test_plugin_enable_disable.py new file mode 100644 index 00000000..8757b0a0 --- /dev/null +++ b/tests/integration/test_plugin_enable_disable.py @@ -0,0 +1,125 @@ +""" +Integration tests — table de vérité persistante (enable/disable) + centre de +contrôle (registry_table). Vérifie que désactiver un plugin survit à un +redémarrage complet du process, pas seulement au unload en mémoire. +""" + +import pytest + +from xcore import Xcore + + +@pytest.fixture +def config_with_plugin(tmp_path): + """Config xcore pointant vers un plugin de test réel sur disque.""" + plugins_dir = tmp_path / "plugins" + plugins_dir.mkdir() + + test_plugin = plugins_dir / "test_plugin" + test_plugin.mkdir() + src_dir = test_plugin / "src" + src_dir.mkdir() + + (test_plugin / "plugin.yaml").write_text(""" +name: test_plugin +version: 1.0.0 +execution_mode: trusted +permissions: + - resource: "*" + actions: ["*"] +""") + (src_dir / "main.py").write_text(""" +from xcore.sdk import TrustedBase, ok + +class Plugin(TrustedBase): + async def handle(self, action, payload): + return ok(message="pong") +""") + + config_content = f""" +app: + name: test-app + secret_key: test-secret-key-32-chars-long!!! + +plugins: + directory: {plugins_dir} + strict_trusted: false + +services: + databases: {{}} + cache: + backend: memory + ttl: 300 +""" + config_path = tmp_path / "config.yaml" + config_path.write_text(config_content) + return str(config_path) + + +class TestRegistryTable: + @pytest.mark.asyncio + async def test_loaded_plugin_appears_enabled_and_ready(self, config_with_plugin): + xcore = Xcore(config_path=config_with_plugin) + try: + await xcore.boot() + table = {row["name"]: row for row in xcore.plugins.registry_table()} + assert table["test_plugin"]["enabled"] is True + assert table["test_plugin"]["state"] == "ready" + assert table["test_plugin"]["loaded"] is True + finally: + await xcore.shutdown() + + +class TestEnableDisable: + @pytest.mark.asyncio + async def test_disable_unloads_and_persists(self, config_with_plugin): + xcore = Xcore(config_path=config_with_plugin) + try: + await xcore.boot() + assert "test_plugin" in xcore.plugins.list_plugins() + + await xcore.plugins.disable("test_plugin", reason="maintenance") + + assert "test_plugin" not in xcore.plugins.list_plugins() + table = {row["name"]: row for row in xcore.plugins.registry_table()} + assert table["test_plugin"]["enabled"] is False + assert table["test_plugin"]["state"] == "not_loaded" + assert table["test_plugin"]["reason"] == "maintenance" + finally: + await xcore.shutdown() + + @pytest.mark.asyncio + async def test_disabled_state_survives_restart(self, config_with_plugin): + """Le cœur de la demande : désactiver doit tenir après un redémarrage du process.""" + first = Xcore(config_path=config_with_plugin) + await first.boot() + await first.plugins.disable("test_plugin") + await first.shutdown() + + # Nouveau process xcore, même config/dossier de plugins. + second = Xcore(config_path=config_with_plugin) + try: + await second.boot() + assert "test_plugin" not in second.plugins.list_plugins() + table = {row["name"]: row for row in second.plugins.registry_table()} + assert table["test_plugin"]["enabled"] is False + assert table["test_plugin"]["state"] == "not_loaded" + finally: + await second.shutdown() + + @pytest.mark.asyncio + async def test_enable_reloads_plugin(self, config_with_plugin): + xcore = Xcore(config_path=config_with_plugin) + try: + await xcore.boot() + await xcore.plugins.disable("test_plugin") + assert "test_plugin" not in xcore.plugins.list_plugins() + + await xcore.plugins.enable("test_plugin") + + assert "test_plugin" in xcore.plugins.list_plugins() + table = {row["name"]: row for row in xcore.plugins.registry_table()} + assert table["test_plugin"]["enabled"] is True + assert table["test_plugin"]["state"] == "ready" + finally: + await xcore.shutdown() diff --git a/tests/test_conftest_cleanup.py b/tests/test_conftest_cleanup.py new file mode 100644 index 00000000..cda6f56d --- /dev/null +++ b/tests/test_conftest_cleanup.py @@ -0,0 +1,31 @@ +""" +Tests for the tmp-dir cleanup safety net in tests/conftest.py. +""" + +import tempfile +from pathlib import Path + +from conftest import sweep_stray_tmp_dirs + + +def test_sweeps_stray_xcore_test_dirs(): + stray = Path(tempfile.mkdtemp(prefix="xcore_test_leftover_")) + (stray / "marker.txt").write_text("leftover") + assert stray.exists() + + sweep_stray_tmp_dirs() + + assert not stray.exists() + + +def test_does_not_touch_unrelated_tmp_dirs(): + unrelated = Path(tempfile.mkdtemp(prefix="not_xcore_related_")) + try: + sweep_stray_tmp_dirs() + assert unrelated.exists() + finally: + unrelated.rmdir() + + +def test_temp_dir_fixture_uses_distinguishing_prefix(temp_dir): + assert temp_dir.name.startswith("xcore_test_") diff --git a/tests/unit/kernel/test_api_router.py b/tests/unit/kernel/test_api_router.py index 96dd988b..25e240c0 100644 --- a/tests/unit/kernel/test_api_router.py +++ b/tests/unit/kernel/test_api_router.py @@ -7,7 +7,6 @@ from xcore.kernel.api.router import CallRequest, CallResponse, _hash_key, build_router - SECRET_KEY = b"test-secret" SERVER_KEY = b"test-server" @@ -37,6 +36,7 @@ def _client_with_key(app, key: str = "test-secret"): # ── _hash_key ───────────────────────────────────────────────────────────────── + class TestHashKey: def test_returns_bytes(self): result = _hash_key("mykey", "salt") @@ -67,6 +67,7 @@ def test_different_keys_differ(self): # ── Auth ────────────────────────────────────────────────────────────────────── + class TestRouterAuth: def test_missing_api_key_returns_401(self): app, _ = _make_app() @@ -97,14 +98,34 @@ def test_correct_api_key_passes(self): # ── Endpoints ───────────────────────────────────────────────────────────────── + class TestRouterEndpoints: def setup_method(self): self.supervisor = MagicMock() - self.supervisor.call = AsyncMock(return_value={"status": "ok", "result": "pong"}) + self.supervisor.call = AsyncMock( + return_value={"status": "ok", "result": "pong"} + ) self.supervisor.status = MagicMock(return_value={"plugins": ["auth"]}) self.supervisor.reload = AsyncMock() self.supervisor.load = AsyncMock() self.supervisor.unload = AsyncMock() + self.supervisor.enable = AsyncMock() + self.supervisor.disable = AsyncMock() + self.supervisor.registry_table = MagicMock( + return_value=[{"name": "auth", "enabled": True, "state": "ready"}] + ) + self.supervisor.ipc_audit = MagicMock( + return_value=[{"plugin": "auth", "action": "login"}] + ) + self.supervisor.ipc_stats = MagicMock( + return_value={"entries": 1, "by_plugin": {}} + ) + self.supervisor.events_activity = MagicMock( + return_value={"recent": [], "stats": {}} + ) + self.supervisor.hooks_activity = MagicMock( + return_value={"recent": [], "metrics": {}} + ) self.app, _ = _make_app(self.supervisor) self.client = TestClient(self.app, raise_server_exceptions=False) self.headers = {"X-Plugin-Key": "test-secret"} @@ -122,7 +143,11 @@ def test_call_plugin_success(self): def test_call_plugin_not_found(self): self.supervisor.call = AsyncMock( - return_value={"status": "error", "code": "not_found", "msg": "Plugin 'x' not found"} + return_value={ + "status": "error", + "code": "not_found", + "msg": "Plugin 'x' not found", + } ) resp = self.client.post( "/ipc/x/ping", @@ -167,6 +192,63 @@ def test_unload_plugin(self): assert resp.status_code == 200 self.supervisor.unload.assert_called_once_with("auth") + def test_registry_endpoint(self): + resp = self.client.get("/ipc/registry", headers=self.headers) + assert resp.status_code == 200 + assert resp.json() == { + "plugins": [{"name": "auth", "enabled": True, "state": "ready"}] + } + + def test_enable_plugin(self): + resp = self.client.post("/ipc/auth/enable", headers=self.headers) + assert resp.status_code == 200 + self.supervisor.enable.assert_called_once_with("auth") + + def test_disable_plugin_no_body(self): + resp = self.client.post("/ipc/auth/disable", headers=self.headers) + assert resp.status_code == 200 + self.supervisor.disable.assert_called_once_with("auth", reason=None) + + def test_disable_plugin_with_reason(self): + resp = self.client.post( + "/ipc/auth/disable", + json={"reason": "maintenance"}, + headers=self.headers, + ) + assert resp.status_code == 200 + self.supervisor.disable.assert_called_once_with("auth", reason="maintenance") + + def test_audit_endpoint(self): + resp = self.client.get("/ipc/audit", headers=self.headers) + assert resp.status_code == 200 + data = resp.json() + assert data["calls"] == [{"plugin": "auth", "action": "login"}] + assert data["stats"] == {"entries": 1, "by_plugin": {}} + + def test_audit_endpoint_with_query_params(self): + resp = self.client.get( + "/ipc/audit", params={"plugin": "auth", "limit": 10}, headers=self.headers + ) + assert resp.status_code == 200 + self.supervisor.ipc_audit.assert_called_once_with("auth", 10) + + def test_events_endpoint(self): + resp = self.client.get("/ipc/events", headers=self.headers) + assert resp.status_code == 200 + data = resp.json() + assert data["events"] == {"recent": [], "stats": {}} + assert data["hooks"] == {"recent": [], "metrics": {}} + + def test_events_endpoint_with_query_params(self): + resp = self.client.get( + "/ipc/events", + params={"event": "user.created", "limit": 5}, + headers=self.headers, + ) + assert resp.status_code == 200 + self.supervisor.events_activity.assert_called_once_with("user.created", 5) + self.supervisor.hooks_activity.assert_called_once_with("user.created", 5) + def test_health_no_checker(self): resp = self.client.get("/ipc/health", headers=self.headers) assert resp.status_code == 200 @@ -187,6 +269,7 @@ def test_metrics_no_registry(self): def test_metrics_with_registry(self): from xcore.kernel.observability.metrics import MetricsRegistry + metrics = MetricsRegistry() metrics.counter("calls").inc(5) app, _ = _make_app(self.supervisor, metrics_registry=metrics) @@ -197,6 +280,7 @@ def test_metrics_with_registry(self): # ── Models ──────────────────────────────────────────────────────────────────── + class TestModels: def test_call_request_default(self): req = CallRequest() @@ -207,6 +291,8 @@ def test_call_request_with_payload(self): assert req.payload == {"key": "value"} def test_call_response(self): - resp = CallResponse(status="ok", plugin="auth", action="login", result={"token": "abc"}) + resp = CallResponse( + status="ok", plugin="auth", action="login", result={"token": "abc"} + ) assert resp.status == "ok" assert resp.plugin == "auth" diff --git a/tests/unit/kernel/test_events.py b/tests/unit/kernel/test_events.py index a94e935d..410cf771 100644 --- a/tests/unit/kernel/test_events.py +++ b/tests/unit/kernel/test_events.py @@ -222,7 +222,8 @@ async def error_handler(event): assert len(results) == 0 assert any( - r.getMessage() == "event handler error" and r.xcore_ctx.get("error") == "Test error" + r.getMessage() == "event handler error" + and r.xcore_ctx.get("error") == "Test error" for r in caplog.records ) @@ -260,3 +261,86 @@ async def handler(event): assert "event1" in events assert "event2" in events + + +class TestEventBusSupervision: + """Supervision : recent_emissions()/stats() — la visibilité qui manquait.""" + + @pytest.fixture + def event_bus(self): + return EventBus() + + @pytest.mark.asyncio + async def test_recent_emissions_records_successful_emit(self, event_bus): + async def handler(event): + return "ok" + + event_bus.subscribe("test.event", handler) + await event_bus.emit("test.event", {"key": "value"}, source="my_plugin") + + recent = event_bus.recent_emissions() + assert len(recent) == 1 + assert recent[0]["event"] == "test.event" + assert recent[0]["source"] == "my_plugin" + assert recent[0]["handlers_matched"] == 1 + assert recent[0]["errors"] == 0 + + @pytest.mark.asyncio + async def test_recent_emissions_records_no_handlers(self, event_bus): + await event_bus.emit("nobody.listens", {}) + + recent = event_bus.recent_emissions() + assert len(recent) == 1 + assert recent[0]["handlers_matched"] == 0 + + @pytest.mark.asyncio + async def test_recent_emissions_counts_handler_errors(self, event_bus): + async def boom(event): + raise ValueError("boom") + + event_bus.subscribe("test.event", boom) + await event_bus.emit("test.event", {}) + + recent = event_bus.recent_emissions() + assert recent[0]["errors"] == 1 + + @pytest.mark.asyncio + async def test_recent_emissions_filters_by_event_name(self, event_bus): + await event_bus.emit("event.a", {}) + await event_bus.emit("event.b", {}) + + recent = event_bus.recent_emissions("event.a") + assert len(recent) == 1 + assert recent[0]["event"] == "event.a" + + @pytest.mark.asyncio + async def test_recent_emissions_respects_limit_and_order(self, event_bus): + for i in range(5): + await event_bus.emit(f"event.{i}", {}) + + recent = event_bus.recent_emissions(limit=2) + assert [e["event"] for e in recent] == ["event.3", "event.4"] + + @pytest.mark.asyncio + async def test_stats_aggregates_per_event(self, event_bus): + async def handler(event): + return "ok" + + event_bus.subscribe("test.event", handler) + await event_bus.emit("test.event", {}) + await event_bus.emit("test.event", {}) + + stats = event_bus.stats("test.event") + assert stats["emissions"] == 2 + assert stats["errors"] == 0 + + def test_stats_unknown_event_returns_empty(self, event_bus): + assert event_bus.stats("never.emitted") == {} + + @pytest.mark.asyncio + async def test_stats_all_events(self, event_bus): + await event_bus.emit("a", {}) + await event_bus.emit("b", {}) + + stats = event_bus.stats() + assert set(stats.keys()) == {"a", "b"} diff --git a/tests/unit/kernel/test_health.py b/tests/unit/kernel/test_health.py index a426727f..ebbab4e8 100644 --- a/tests/unit/kernel/test_health.py +++ b/tests/unit/kernel/test_health.py @@ -38,6 +38,22 @@ def check_cache(): assert result["status"] == "healthy" assert result["checks"]["cache"]["status"] == "healthy" + @pytest.mark.asyncio + async def test_unregister_removes_check(self): + hc = HealthChecker() + + @hc.register("database") + async def check_db(): + return True, "ok" + + assert hc.unregister("database") is True + result = await hc.run_all() + assert result["checks"] == {} + + def test_unregister_unknown_returns_false(self): + hc = HealthChecker() + assert hc.unregister("does_not_exist") is False + @pytest.mark.asyncio async def test_degraded_check(self): hc = HealthChecker() diff --git a/tests/unit/kernel/test_hooks.py b/tests/unit/kernel/test_hooks.py index af9a39f3..9caee83c 100644 --- a/tests/unit/kernel/test_hooks.py +++ b/tests/unit/kernel/test_hooks.py @@ -333,3 +333,51 @@ def handler(event): hooks = hook_manager.list_hooks() assert len(hooks) == 0 + + +class TestHookManagerSupervision: + """Supervision : recent_emissions() — le détail récent derrière get_metrics().""" + + @pytest.fixture + def hook_manager(self): + return HookManager() + + @pytest.mark.asyncio + async def test_recent_emissions_records_emit(self, hook_manager): + def handler(event): + return "ok" + + hook_manager.register("test.event", handler) + await hook_manager.emit("test.event", {"key": "value"}) + + recent = hook_manager.recent_emissions() + assert len(recent) == 1 + assert recent[0]["event"] == "test.event" + assert recent[0]["hooks_matched"] == 1 + assert recent[0]["errors"] == 0 + + @pytest.mark.asyncio + async def test_recent_emissions_no_handlers(self, hook_manager): + await hook_manager.emit("nobody.listens", {}) + recent = hook_manager.recent_emissions() + assert recent[0]["hooks_matched"] == 0 + + @pytest.mark.asyncio + async def test_recent_emissions_counts_errors(self, hook_manager): + def boom(event): + raise ValueError("boom") + + hook_manager.register("test.event", boom) + await hook_manager.emit("test.event", {}) + + recent = hook_manager.recent_emissions() + assert recent[0]["errors"] == 1 + + @pytest.mark.asyncio + async def test_recent_emissions_filters_by_event_name(self, hook_manager): + await hook_manager.emit("event.a", {}) + await hook_manager.emit("event.b", {}) + + recent = hook_manager.recent_emissions("event.a") + assert len(recent) == 1 + assert recent[0]["event"] == "event.a" diff --git a/tests/unit/kernel/test_lifecycle.py b/tests/unit/kernel/test_lifecycle.py index 4315825d..2e66c301 100644 --- a/tests/unit/kernel/test_lifecycle.py +++ b/tests/unit/kernel/test_lifecycle.py @@ -332,6 +332,170 @@ async def on_unload(self): assert lifecycle_manager._instance is None assert lifecycle_manager.state == PluginState.UNLOADED + @pytest.mark.asyncio + async def test_unload_unregisters_from_registry(self, lifecycle_manager, tmp_path): + """Le ramasse-miette forcé doit toujours retirer le plugin du PluginRegistry.""" + src_dir = tmp_path / "src" + src_dir.mkdir() + (src_dir / "main.py").write_text(""" +from xcore.kernel.api.contract import BasePlugin + +class Plugin(BasePlugin): + async def handle(self, action, payload): + return {"status": "ok"} +""") + + await lifecycle_manager.load() + await lifecycle_manager.unload() + + lifecycle_manager._registry.unregister.assert_called_once_with("test_plugin") + + @pytest.mark.asyncio + async def test_unload_forces_cleanup_even_if_on_unload_raises( + self, lifecycle_manager, tmp_path + ): + """Un on_unload buggé ne doit pas empêcher le ramasse-miette forcé.""" + src_dir = tmp_path / "src" + src_dir.mkdir() + (src_dir / "main.py").write_text(""" +from xcore.kernel.api.contract import BasePlugin + +class Plugin(BasePlugin): + async def handle(self, action, payload): + return {"status": "ok"} + + async def on_unload(self): + raise RuntimeError("boom") +""") + + await lifecycle_manager.load() + await lifecycle_manager.unload() + + assert lifecycle_manager.state == PluginState.UNLOADED + assert lifecycle_manager._instance is None + lifecycle_manager._registry.unregister.assert_called_once_with("test_plugin") + + @pytest.mark.asyncio + async def test_unload_removes_scheduler_jobs(self, mock_manifest): + """Un job planifié par le plugin doit être désinscrit du vrai scheduler au unload.""" + from xcore.kernel.context import KernelContext + + real_scheduler = MagicMock() + real_scheduler.add_job.side_effect = lambda func, trigger="cron", job_id=None, **kw: job_id + + ctx = KernelContext( + # tenancy=None explicite : sinon un MagicMock() auto-généré rend + # `config.tenancy.enabled` truthy et le scheduler passe aussi par + # TenantAwareScheduler, hors sujet de ce test. + config=MagicMock(tenancy=None), + services=MagicMock(), + events=MagicMock(), + hooks=MagicMock(), + # Comme en production, le scheduler n'est pas encore exporté dans le + # registry pendant load_all() (register_core_service() n'a lieu + # qu'après) : get_service() doit retomber sur le dict de services. + registry=MagicMock(get_service=MagicMock(side_effect=KeyError("scheduler"))), + metrics=MagicMock(), + tracer=MagicMock(), + health=MagicMock(), + ) + ctx.services.as_dict.return_value = {"scheduler": real_scheduler} + + manager = LifecycleManager(manifest=mock_manifest, ctx=ctx) + + src_dir = mock_manifest.plugin_dir / "src" + src_dir.mkdir() + (src_dir / "main.py").write_text(""" +from xcore.kernel.api.contract import BasePlugin + +class Plugin(BasePlugin): + async def handle(self, action, payload): + return {"status": "ok"} + + async def _inject_context(self, ctx): + self.ctx = ctx + ctx.get_service("scheduler").add_job(self.handle, trigger="cron", job_id="nightly") +""") + + await manager.load() + real_scheduler.add_job.assert_called_once() + await manager.unload() + real_scheduler.remove_job.assert_called_once_with("nightly") + + @pytest.mark.asyncio + async def test_unload_removes_health_check(self, mock_manifest): + """Un health check enregistré par le plugin doit être désinscrit au unload.""" + from xcore.kernel.context import KernelContext + + real_health = MagicMock() + + ctx = KernelContext( + config=MagicMock(), + services=MagicMock(), + events=MagicMock(), + hooks=MagicMock(), + registry=MagicMock(), + metrics=MagicMock(), + tracer=MagicMock(), + health=real_health, + ) + ctx.services.as_dict.return_value = {} + + manager = LifecycleManager(manifest=mock_manifest, ctx=ctx) + + src_dir = mock_manifest.plugin_dir / "src" + src_dir.mkdir() + (src_dir / "main.py").write_text(""" +from xcore.kernel.api.contract import BasePlugin + +class Plugin(BasePlugin): + async def handle(self, action, payload): + return {"status": "ok"} + + async def _inject_context(self, ctx): + self.ctx = ctx + + @ctx.health.register("test_plugin.db") + async def check(): + return True, "ok" +""") + + await manager.load() + real_health.register.assert_called_once_with("test_plugin.db") + await manager.unload() + real_health.unregister.assert_called_once_with("test_plugin.db") + + @pytest.mark.asyncio + async def test_spawn_task_cancelled_on_unload(self, lifecycle_manager, tmp_path): + """Une tâche créée via ctx.spawn_task() doit être annulée au unload.""" + src_dir = tmp_path / "src" + src_dir.mkdir() + (src_dir / "main.py").write_text(""" +import asyncio +from xcore.kernel.api.contract import BasePlugin + +class Plugin(BasePlugin): + async def handle(self, action, payload): + return {"status": "ok"} + + async def _inject_context(self, ctx): + self.ctx = ctx + + async def _forever(): + await asyncio.sleep(3600) + + self.task = ctx.spawn_task(_forever(), name="forever") +""") + + await lifecycle_manager.load() + task = lifecycle_manager._instance.task + assert not task.done() + + await lifecycle_manager.unload() + await asyncio.sleep(0) # laisse la cancellation se propager + + assert task.cancelled() + @pytest.mark.asyncio async def test_collect_router(self, lifecycle_manager, tmp_path): """Test router collection.""" diff --git a/tests/unit/kernel/test_loader.py b/tests/unit/kernel/test_loader.py index ec7b7c71..d279c168 100644 --- a/tests/unit/kernel/test_loader.py +++ b/tests/unit/kernel/test_loader.py @@ -145,6 +145,65 @@ def test_loader_topo_sort(loader): assert [m.name for m in ordered] == ["p1", "p2"] +@pytest.fixture +def mock_ctx_tmp(tmp_path): + """Comme mock_ctx, mais avec un vrai dossier isolé (pour la state_store).""" + ctx = MagicMock() + ctx.config.directory = str(tmp_path / "plugins") + ctx.services.as_dict.return_value = {} + ctx.events = MagicMock() + ctx.hooks = MagicMock() + ctx.registry = MagicMock() + ctx.metrics = MagicMock() + ctx.tracer = MagicMock() + ctx.health = MagicMock() + return ctx + + +def test_loader_state_store_default_path(mock_ctx_tmp, tmp_path): + loader = PluginLoader(mock_ctx_tmp) + assert loader.state_store._path == tmp_path / ".xcore" / "plugins_state.json" + + +def test_loader_discover_names(mock_ctx_tmp, tmp_path): + plugins_dir = tmp_path / "plugins" + (plugins_dir / "shop").mkdir(parents=True) + (plugins_dir / "_disabled_by_prefix").mkdir() + (plugins_dir / "not_a_dir.txt").write_text("x") + + loader = PluginLoader(mock_ctx_tmp) + assert loader.discover_names() == ["shop"] + + +@pytest.mark.asyncio +async def test_loader_load_all_skips_disabled_plugin(mock_ctx_tmp): + m1 = MagicMock() + m1.name = "p1" + m1.requires = [] + m1.execution_mode = ExecutionMode.TRUSTED + m1.version = "1.0.0" + + p1 = MagicMock(spec=Path) + p1.is_dir.return_value = True + p1.name = "p1" + + loader = PluginLoader(mock_ctx_tmp) + loader.state_store.set_enabled("p1", False, reason="maintenance") + + with patch("pathlib.Path.exists", return_value=True): + with patch("pathlib.Path.iterdir", return_value=[p1]): + with patch.object( + loader._validator, "load_and_validate", return_value=(m1, True, "2.3.2") + ): + with patch.object( + loader, "_activate", new_callable=AsyncMock + ) as mock_activate: + res = await loader.load_all() + assert res["loaded"] == [] + assert res["disabled"] == ["p1"] + mock_activate.assert_not_called() + + def test_loader_topo_sort_circular(loader): m1 = MagicMock() m1.name = "p1" diff --git a/tests/unit/kernel/test_permissions.py b/tests/unit/kernel/test_permissions.py index a9b5b36e..ab848984 100644 --- a/tests/unit/kernel/test_permissions.py +++ b/tests/unit/kernel/test_permissions.py @@ -116,3 +116,34 @@ def emit_sync(self, event_name: str, data: dict) -> None: # Check allows second call (cache hit for db.posts) assert engine.allows("p1", "db.posts", "read") is True assert len(events.emitted) == 2 + + def test_audit_cache_hits_default_keeps_full_log(self, engine): + """Comportement historique inchangé : audit_cache_hits=True par défaut.""" + engine.load_from_manifest("p1", [{"resource": "*", "actions": ["*"]}]) + + engine.allows("p1", "res", "act") # cache miss -> logged + engine.allows("p1", "res", "act") # cache hit -> logged (défaut) + + assert len(engine.audit_log(plugin_name="p1")) == 2 + + def test_audit_cache_hits_disabled_skips_hit_entries(self): + engine = PermissionEngine(audit_cache_hits=False) + engine.load_from_manifest("p1", [{"resource": "*", "actions": ["*"]}]) + + engine.allows("p1", "res", "act") # cache miss -> logged + for _ in range(5): + engine.allows("p1", "res", "act") # cache hits -> not logged + + assert len(engine.audit_log(plugin_name="p1")) == 1 + + def test_audit_cache_hits_disabled_does_not_change_check_behavior(self): + """Le flag ne touche que l'audit log — check()/allows() restent corrects.""" + engine = PermissionEngine(audit_cache_hits=False) + engine.load_from_manifest( + "p1", [{"resource": "db.*", "actions": ["read"], "effect": "deny"}] + ) + + with pytest.raises(PermissionDenied): + engine.check("p1", "db.users", "read") # cache miss + with pytest.raises(PermissionDenied): + engine.check("p1", "db.users", "read") # cache hit, still denies diff --git a/tests/unit/kernel/test_process_manager.py b/tests/unit/kernel/test_process_manager.py index fc4dbc60..66ed14c1 100644 --- a/tests/unit/kernel/test_process_manager.py +++ b/tests/unit/kernel/test_process_manager.py @@ -7,7 +7,7 @@ from unittest.mock import MagicMock, AsyncMock, patch from pathlib import Path from xcore.kernel.sandbox.process_manager import SandboxProcessManager, SandboxConfig, ProcessState -from xcore.kernel.sandbox.ipc import IPCResponse, IPCProcessDead +from xcore.kernel.sandbox.ipc import IPCResponse, IPCProcessDead, IPCTimeoutError @pytest.fixture def mock_manifest(tmp_path): @@ -105,6 +105,36 @@ async def test_manager_call_not_available(manager): with pytest.raises(RuntimeError, match="non disponible"): await manager.call("ping", {}) + +@pytest.mark.asyncio +async def test_manager_call_process_dead_triggers_recycle(manager): + """IPCProcessDead sur un appel doit déclencher le recyclage immédiatement.""" + manager._state = ProcessState.RUNNING + manager._channel = MagicMock() + manager._channel.call = AsyncMock(side_effect=IPCProcessDead("dead")) + + with patch.object(manager, "_handle_crash", new_callable=AsyncMock) as mock_crash: + with pytest.raises(IPCProcessDead): + await manager.call("ping", {}) + mock_crash.assert_called_once() + + +@pytest.mark.asyncio +async def test_manager_call_ipc_timeout_triggers_recycle(manager): + """ + IPCTimeoutError (subprocess bloqué, pas mort) doit aussi déclencher le + recyclage immédiat — avant ce fix, seul IPCProcessDead était catché ici + et il fallait attendre le prochain _health_loop pour recycler. + """ + manager._state = ProcessState.RUNNING + manager._channel = MagicMock() + manager._channel.call = AsyncMock(side_effect=IPCTimeoutError("no response")) + + with patch.object(manager, "_handle_crash", new_callable=AsyncMock) as mock_crash: + with pytest.raises(IPCTimeoutError): + await manager.call("ping", {}) + mock_crash.assert_called_once() + @pytest.mark.asyncio async def test_manager_stop(manager): manager._state = ProcessState.RUNNING diff --git a/tests/unit/kernel/test_supervisor.py b/tests/unit/kernel/test_supervisor.py index 4ac034d1..c32aa1dd 100644 --- a/tests/unit/kernel/test_supervisor.py +++ b/tests/unit/kernel/test_supervisor.py @@ -25,6 +25,7 @@ def _make_ctx(): class TestPluginSupervisorPreBoot: def _make(self): from xcore.kernel.runtime.supervisor import PluginSupervisor + return PluginSupervisor(_make_ctx()) @pytest.mark.asyncio @@ -74,6 +75,7 @@ async def test_shutdown_before_boot(self): def test_err_static(self): from xcore.kernel.runtime.supervisor import PluginSupervisor + result = PluginSupervisor._err("something", "my_code") assert result["status"] == "error" assert result["code"] == "my_code" @@ -84,10 +86,42 @@ def test_register_middleware_before_boot_raises(self): with pytest.raises(RuntimeError, match="boot"): sup.register_middleware(mw) + @pytest.mark.asyncio + async def test_enable_before_boot_raises(self): + sup = self._make() + with pytest.raises(RuntimeError): + await sup.enable("shop") + + @pytest.mark.asyncio + async def test_disable_before_boot_raises(self): + sup = self._make() + with pytest.raises(RuntimeError): + await sup.disable("shop") + + def test_registry_table_before_boot_empty(self): + sup = self._make() + assert sup.registry_table() == [] + + def test_ipc_audit_before_boot_empty(self): + sup = self._make() + assert sup.ipc_audit() == [] + assert sup.ipc_stats() == {"entries": 0, "by_plugin": {}} + + def test_events_activity_no_bus(self): + ctx = _make_ctx() + ctx.events = None + ctx.hooks = None + from xcore.kernel.runtime.supervisor import PluginSupervisor + + sup = PluginSupervisor(ctx) + assert sup.events_activity() == {"recent": [], "stats": {}} + assert sup.hooks_activity() == {"recent": [], "metrics": {}} + @pytest.mark.asyncio async def test_boot_empty_plugin_dir(self): import tempfile, os from xcore.kernel.runtime.supervisor import PluginSupervisor + ctx = _make_ctx() with tempfile.TemporaryDirectory() as tmp: ctx.config.directory = tmp @@ -99,6 +133,7 @@ async def test_boot_empty_plugin_dir(self): async def test_call_plugin_not_found_after_boot(self): import tempfile from xcore.kernel.runtime.supervisor import PluginSupervisor + ctx = _make_ctx() with tempfile.TemporaryDirectory() as tmp: ctx.config.directory = tmp @@ -108,10 +143,124 @@ async def test_call_plugin_not_found_after_boot(self): assert result["status"] == "error" assert result["code"] == "not_found" + +class TestIPCSupervision: + """L'appel IPC 'xcore' (KernelHandler) est déjà monté au boot — pas besoin + d'un vrai plugin sur disque pour observer l'audit trail.""" + + @pytest.mark.asyncio + async def test_call_is_audited(self): + import tempfile + from xcore.kernel.runtime.supervisor import PluginSupervisor + + ctx = _make_ctx() + with tempfile.TemporaryDirectory() as tmp: + ctx.config.directory = tmp + sup = PluginSupervisor(ctx) + await sup.boot() + + await sup.call( + "xcore", "plugin.list", {}, caller="test_caller", tenant_id="acme" + ) + + audit = sup.ipc_audit() + assert len(audit) == 1 + entry = audit[0] + assert entry["plugin"] == "xcore" + assert entry["action"] == "plugin.list" + assert entry["caller"] == "test_caller" + assert entry["tenant_id"] == "acme" + assert entry["status"] == "ok" + assert entry["duration_ms"] >= 0 + + @pytest.mark.asyncio + async def test_not_found_call_is_audited(self): + import tempfile + from xcore.kernel.runtime.supervisor import PluginSupervisor + + ctx = _make_ctx() + with tempfile.TemporaryDirectory() as tmp: + ctx.config.directory = tmp + sup = PluginSupervisor(ctx) + await sup.boot() + + await sup.call("ghost", "ping", {}) + + audit = sup.ipc_audit() + assert audit[-1]["plugin"] == "ghost" + assert audit[-1]["status"] == "error" + assert audit[-1]["code"] == "not_found" + + @pytest.mark.asyncio + async def test_ipc_audit_filters_by_plugin(self): + import tempfile + from xcore.kernel.runtime.supervisor import PluginSupervisor + + ctx = _make_ctx() + with tempfile.TemporaryDirectory() as tmp: + ctx.config.directory = tmp + sup = PluginSupervisor(ctx) + await sup.boot() + + await sup.call("xcore", "plugin.list", {}) + await sup.call("ghost", "ping", {}) + + audit = sup.ipc_audit(plugin_name="xcore") + assert len(audit) == 1 + assert audit[0]["plugin"] == "xcore" + + @pytest.mark.asyncio + async def test_ipc_stats_aggregates(self): + import tempfile + from xcore.kernel.runtime.supervisor import PluginSupervisor + + ctx = _make_ctx() + with tempfile.TemporaryDirectory() as tmp: + ctx.config.directory = tmp + sup = PluginSupervisor(ctx) + await sup.boot() + + await sup.call("xcore", "plugin.list", {}) + await sup.call("xcore", "plugin.list", {}) + await sup.call("ghost", "ping", {}) + + stats = sup.ipc_stats() + assert stats["entries"] == 3 + assert stats["by_plugin"]["xcore"]["calls"] == 2 + assert stats["by_plugin"]["xcore"]["errors"] == 0 + assert stats["by_plugin"]["ghost"]["errors"] == 1 + + +class TestEventsAndHooksSupervision: + @pytest.mark.asyncio + async def test_events_activity_reflects_real_bus(self): + import tempfile + from xcore.kernel.events.bus import EventBus + from xcore.kernel.events.hooks import HookManager + from xcore.kernel.runtime.supervisor import PluginSupervisor + + ctx = _make_ctx() + ctx.events = EventBus() + ctx.hooks = HookManager() + with tempfile.TemporaryDirectory() as tmp: + ctx.config.directory = tmp + sup = PluginSupervisor(ctx) + await sup.boot() + + await ctx.events.emit("custom.event", {"x": 1}, source="test") + await ctx.hooks.emit("custom.hook", {"x": 1}) + + events = sup.events_activity() + assert events["stats"]["custom.event"]["emissions"] >= 1 + + hooks = sup.hooks_activity() + assert any(e["event"] == "custom.hook" for e in hooks["recent"]) + @pytest.mark.asyncio async def test_load_unload_after_boot(self): import tempfile from xcore.kernel.runtime.supervisor import PluginSupervisor + ctx = _make_ctx() with tempfile.TemporaryDirectory() as tmp: ctx.config.directory = tmp diff --git a/tests/unit/test_registry_state_store.py b/tests/unit/test_registry_state_store.py new file mode 100644 index 00000000..faff432f --- /dev/null +++ b/tests/unit/test_registry_state_store.py @@ -0,0 +1,49 @@ +"""Tests for PluginStateStore — table de vérité actif/inactif persistante.""" + +import json + +from xcore.registry.state_store import PluginStateStore + + +class TestPluginStateStore: + def test_default_enabled_when_never_toggled(self, tmp_path): + store = PluginStateStore(tmp_path / ".xcore" / "plugins_state.json") + assert store.is_enabled("shop") is True + + def test_set_enabled_false_then_true(self, tmp_path): + path = tmp_path / ".xcore" / "plugins_state.json" + store = PluginStateStore(path) + + store.set_enabled("shop", False, reason="maintenance") + assert store.is_enabled("shop") is False + assert store.all()["shop"]["reason"] == "maintenance" + + store.set_enabled("shop", True) + assert store.is_enabled("shop") is True + + def test_persists_across_instances(self, tmp_path): + path = tmp_path / ".xcore" / "plugins_state.json" + PluginStateStore(path).set_enabled("shop", False) + + # Nouvelle instance (simule un redémarrage du process) : doit relire le fichier. + reloaded = PluginStateStore(path) + assert reloaded.is_enabled("shop") is False + + def test_writes_valid_json_atomically(self, tmp_path): + path = tmp_path / ".xcore" / "plugins_state.json" + store = PluginStateStore(path) + store.set_enabled("shop", False) + + assert path.exists() + assert not path.with_suffix(".json.tmp").exists() + data = json.loads(path.read_text()) + assert data["shop"]["enabled"] is False + + def test_corrupt_file_starts_empty_instead_of_crashing(self, tmp_path): + path = tmp_path / ".xcore" / "plugins_state.json" + path.parent.mkdir(parents=True) + path.write_text("{ not valid json") + + store = PluginStateStore(path) + assert store.is_enabled("shop") is True + assert store.all() == {} diff --git a/xcore/__init__.py b/xcore/__init__.py index 02eab6de..5d1589cc 100644 --- a/xcore/__init__.py +++ b/xcore/__init__.py @@ -246,6 +246,14 @@ async def _on_plugin_reloaded(event): self.events.subscribe("plugin.*.reloaded", _on_plugin_reloaded) + # Retire les routes FastAPI d'un plugin désactivé/déchargé — sans elles, + # les routes restaient montées indéfiniment après un unload ou disable(). + async def _on_plugin_unloaded(event): + plugin_name = event.name.split(".")[1] + self._unmount_plugin_router(plugin_name) + + self.events.subscribe("plugin.*.unloaded", _on_plugin_unloaded) + # 5. Attache le router FastAPI si une app est fournie if app is not None: self._app = app @@ -272,21 +280,38 @@ async def shutdown(self) -> None: self._booted = False self._logger.info("xcore stopped") - def _remount_plugin_router(self, plugin_name: str) -> None: - """Re-monte les routes FastAPI d'un plugin après un hot-reload.""" + def _unmount_plugin_router(self, plugin_name: str) -> None: + """ + Retire de l'app FastAPI toutes les routes montées pour ce plugin. + + FastAPI n'offre pas de désinscription native d'un routeur — on filtre + donc `app.routes`, seule approche possible. Utilisé au unload/disable + et en première étape du remount lors d'un reload. + """ app = self._app - if app is None or self.plugins is None: + if app is None: return prefix = self._config.app.plugin_prefix or "/plugins" plugin_prefix = f"{prefix}/{plugin_name}" - # Retire toutes les routes qui appartiennent à ce plugin app.routes = [ r for r in app.routes if not getattr(r, "path", "").startswith(plugin_prefix) ] + app.openapi_schema = None # force regen du schéma OpenAPI + + def _remount_plugin_router(self, plugin_name: str) -> None: + """Re-monte les routes FastAPI d'un plugin après un hot-reload.""" + app = self._app + if app is None or self.plugins is None: + return + + self._unmount_plugin_router(plugin_name) + + prefix = self._config.app.plugin_prefix or "/plugins" + plugin_prefix = f"{prefix}/{plugin_name}" # Récupère le nouveau router depuis le handler rechargé try: diff --git a/xcore/kernel/api/context.py b/xcore/kernel/api/context.py index 18ff30b9..eaf74d2a 100644 --- a/xcore/kernel/api/context.py +++ b/xcore/kernel/api/context.py @@ -8,8 +8,9 @@ from __future__ import annotations +import asyncio from dataclasses import dataclass, field -from typing import TYPE_CHECKING, Any, Awaitable, Callable +from typing import TYPE_CHECKING, Any, Awaitable, Callable, Coroutine if TYPE_CHECKING: from ...registry import PluginRegistry @@ -46,6 +47,25 @@ class PluginContext: health: HealthChecker = None # HealthChecker registry: PluginRegistry = None # PluginRegistry + # Rempli par LifecycleManager : les tâches créées via spawn_task() y sont + # ajoutées pour pouvoir être annulées de force au unload, sans dépendre + # du plugin pour les traquer lui-même. + _task_sink: list[asyncio.Task] | None = field(default=None, repr=False) + + def spawn_task(self, coro: Coroutine, name: str | None = None) -> asyncio.Task: + """ + Crée une tâche de fond suivie par le kernel. + + À utiliser à la place d'un `asyncio.create_task()` brut pour tout ce + qui ne doit pas survivre au unload du plugin : le kernel annule + automatiquement les tâches encore en cours au moment du unload, + qu'on_unload/on_stop l'ait fait ou non. + """ + task = asyncio.create_task(coro, name=name) + if self._task_sink is not None: + self._task_sink.append(task) + return task + def get_service(self, name: str) -> Any: """ Accès sécurisé à un service avec vérification de scoping via le registry diff --git a/xcore/kernel/api/router.py b/xcore/kernel/api/router.py index dee30cc6..feb62977 100644 --- a/xcore/kernel/api/router.py +++ b/xcore/kernel/api/router.py @@ -23,6 +23,12 @@ class CallRequest(BaseModel): payload: dict[str, Any] = Field(default_factory=dict) +class DisableRequest(BaseModel): + """optional body for POST /{plugin_name}/disable.""" + + reason: Optional[str] = None + + class CallResponse(BaseModel): status: str plugin: str @@ -115,6 +121,34 @@ async def verify_api_key( # ========================= # Routes # ========================= + # + # NB : ces routes spécifiques (registry/enable/disable) doivent être + # déclarées AVANT le catch-all POST /{plugin_name}/{action} ci-dessous — + # Starlette matche dans l'ordre de déclaration, et un POST sur + # "/{plugin_name}/enable" correspondrait sinon toujours à call_plugin() + # en premier (action="enable"), qui échoue faute de corps CallRequest. + + @router.get("/registry") + async def plugins_registry() -> dict[str, Any]: + """ + Table de vérité complète : tous les plugins présents sur disque, avec + leur flag persisté (enabled) et leur état live — y compris les + plugins désactivés ou jamais chargés. Le vrai « centre de contrôle ». + """ + return {"plugins": supervisor.registry_table()} + + @router.post("/{plugin_name}/enable") + async def enable_plugin(plugin_name: str) -> dict[str, str]: + await supervisor.enable(plugin_name) + return {"status": "ok", "msg": f"Plugin '{plugin_name}' enabled"} + + @router.post("/{plugin_name}/disable") + async def disable_plugin( + plugin_name: str, body: DisableRequest | None = None + ) -> dict[str, str]: + reason = body.reason if body else None + await supervisor.disable(plugin_name, reason=reason) + return {"status": "ok", "msg": f"Plugin '{plugin_name}' disabled"} @router.post( "/{plugin_name}/{action}", @@ -155,6 +189,26 @@ async def call_plugin( async def plugins_status() -> dict[str, Any]: return supervisor.status() + @router.get("/audit") + async def ipc_audit( + plugin: Optional[str] = None, limit: int = 100 + ) -> dict[str, Any]: + """Journal des appels IPC (qui a appelé quoi, quand) + stats par plugin.""" + return { + "calls": supervisor.ipc_audit(plugin, limit), + "stats": supervisor.ipc_stats(), + } + + @router.get("/events") + async def events_activity( + event: Optional[str] = None, limit: int = 100 + ) -> dict[str, Any]: + """Activité récente de l'EventBus et du HookManager — plus de zone d'ombre sur les events.""" + return { + "events": supervisor.events_activity(event, limit), + "hooks": supervisor.hooks_activity(event, limit), + } + @router.post("/{plugin_name}/reload") async def reload_plugin(plugin_name: str) -> dict[str, str]: await supervisor.reload(plugin_name) diff --git a/xcore/kernel/events/bus.py b/xcore/kernel/events/bus.py index c1200934..2e63c0bd 100644 --- a/xcore/kernel/events/bus.py +++ b/xcore/kernel/events/bus.py @@ -13,6 +13,8 @@ import fnmatch import inspect import re +import time +from collections import deque from typing import TYPE_CHECKING, Any, Callable, Pattern from .section import Event, _HandlerEntry @@ -44,10 +46,16 @@ async def welcome(event: Event): ``` """ - def __init__(self, cache: "CacheService" = None) -> None: + def __init__(self, cache: "CacheService" = None, max_audit: int = 10_000) -> None: self._handlers: dict[str, list[_HandlerEntry]] = {} self._wildcard_patterns: dict[str, Pattern] = {} + # Supervision : journal borné des émissions + métriques agrégées par + # event, même pattern que PermissionEngine._audit_log — pour voir ce + # qui s'est réellement passé plutôt que de deviner depuis les logs. + self._emission_log: deque[dict] = deque(maxlen=max_audit) + self._stats: dict[str, dict] = {} + # ── Enregistrement ──────────────────────────────────────── def on( @@ -129,6 +137,7 @@ async def emit( await bus.emit("user.created", {"email": "alice@example.com"}) ``` """ + t0 = time.monotonic() event = Event(name=event_name, data=data or {}, source=source) # 1. Exact match lookup (O(1)) @@ -145,6 +154,9 @@ async def emit( matched_handlers.extend(self._handlers[pattern]) if not matched_handlers: + self._audit_emission( + event_name, source, matched=0, errors=0, duration_ms=0.0 + ) return [] # Sort by priority across all matched patterns @@ -152,6 +164,7 @@ async def emit( results: list[Any] = [] to_remove: list[_HandlerEntry] = [] + error_count = 0 if gather: # Fast-path: only one handler @@ -165,6 +178,7 @@ async def emit( ) results.append(result) except Exception as e: + error_count += 1 logger.error( "event handler error", handler=entry.name, @@ -190,6 +204,7 @@ async def _call_sync(h, e): raw = await asyncio.gather(*tasks, return_exceptions=True) for entry, result in zip(matched_handlers, raw): if isinstance(result, Exception): + error_count += 1 logger.error( "event handler error", handler=entry.name, @@ -212,6 +227,7 @@ async def _call_sync(h, e): ) results.append(result) except Exception as e: + error_count += 1 logger.error( "event handler error", handler=entry.name, error=str(e) ) @@ -227,6 +243,15 @@ async def _call_sync(h, e): if not entries: self._handlers.pop(entry.pattern, None) self._wildcard_patterns.pop(entry.pattern, None) + + duration_ms = (time.monotonic() - t0) * 1000 + self._audit_emission( + event_name, + source, + matched=len(matched_handlers), + errors=error_count, + duration_ms=duration_ms, + ) return results def emit_sync(self, event_name: str, data: dict[str, Any] | None = None) -> None: @@ -237,6 +262,70 @@ def emit_sync(self, event_name: str, data: dict[str, Any] | None = None) -> None except RuntimeError: asyncio.run(self.emit(event_name, data)) + # ── Supervision ─────────────────────────────────────────── + + def _audit_emission( + self, + event_name: str, + source: str | None, + matched: int, + errors: int, + duration_ms: float, + ) -> None: + entry = { + "event": event_name, + "source": source, + "handlers_matched": matched, + "errors": errors, + "duration_ms": round(duration_ms, 2), + "timestamp": time.time(), + } + self._emission_log.append(entry) + + stats = self._stats.setdefault( + event_name, {"emissions": 0, "errors": 0, "total_duration_ms": 0.0} + ) + stats["emissions"] += 1 + stats["errors"] += errors + stats["total_duration_ms"] += duration_ms + + def recent_emissions( + self, event_name: str | None = None, limit: int = 100 + ) -> list[dict]: + """Journal des émissions récentes — le plus récent en dernier.""" + from itertools import islice + + it = reversed(self._emission_log) + if event_name: + it = (e for e in it if e["event"] == event_name) + results = list(islice(it, limit)) + results.reverse() + return results + + def stats(self, event_name: str | None = None) -> dict: + """Statistiques agrégées par event : émissions, erreurs, durée moyenne.""" + if event_name: + s = self._stats.get(event_name) + if not s: + return {} + return { + "emissions": s["emissions"], + "errors": s["errors"], + "avg_duration_ms": round(s["total_duration_ms"] / s["emissions"], 2), + } + return { + name: { + "emissions": s["emissions"], + "errors": s["errors"], + "avg_duration_ms": ( + round(s["total_duration_ms"] / s["emissions"], 2) + if s["emissions"] + else 0.0 + ), + } + for name, s in self._stats.items() + } + # ── Introspection ───────────────────────────────────────── def list_events(self) -> dict[str, list[str]]: diff --git a/xcore/kernel/events/hooks.py b/xcore/kernel/events/hooks.py index af37c3f8..44cfae2c 100644 --- a/xcore/kernel/events/hooks.py +++ b/xcore/kernel/events/hooks.py @@ -9,6 +9,7 @@ import fnmatch import inspect import time +from collections import deque from typing import Any, Callable, Dict, List, Optional, Tuple from ..observability import get_logger @@ -31,13 +32,17 @@ class HookManager: Identical to v1 but relocated to kernel/events. """ - def __init__(self): + def __init__(self, max_audit: int = 10_000): self._hooks: Dict[str, List[HookInfo]] = {} self._pre_interceptors: Dict[str, List[Tuple[Callable, int]]] = {} self._post_interceptors: Dict[str, List[Tuple[Callable, int]]] = {} self._metrics: Dict[str, Dict[str, Any]] = {} self._result_processors: Dict[str, List[Callable]] = {} + # Supervision : journal borné des émissions, même pattern que EventBus — + # get_metrics() donne déjà l'agrégat, ceci donne le détail récent. + self._emission_log: deque[dict] = deque(maxlen=max_audit) + def register( self, event_name: str, @@ -186,6 +191,7 @@ async def emit( event = Event(name=event_name, data={**(data or {}), **kwargs}) matching = self._get_matching_hooks(event_name) if not matching: + self._log_emission(event_name, matched=0, errors=0, duration_ms=0.0) return [] results: List[HookResult] = [] @@ -203,6 +209,38 @@ async def emit( self.unregister(pattern, func) self._update_metrics(event_name, results) + self._log_emission( + event_name, + matched=len(matching), + errors=sum(1 for r in results if r.error), + duration_ms=sum(r.execution_time_ms for r in results), + ) + return results + + def _log_emission( + self, event_name: str, matched: int, errors: int, duration_ms: float + ) -> None: + self._emission_log.append( + { + "event": event_name, + "hooks_matched": matched, + "errors": errors, + "duration_ms": round(duration_ms, 2), + "timestamp": time.time(), + } + ) + + def recent_emissions( + self, event_name: Optional[str] = None, limit: int = 100 + ) -> List[dict]: + """Journal des émissions récentes — le plus récent en dernier.""" + from itertools import islice + + it = reversed(self._emission_log) + if event_name: + it = (e for e in it if e["event"] == event_name) + results = list(islice(it, limit)) + results.reverse() return results def _update_metrics(self, event_name: str, results: List[HookResult]) -> None: diff --git a/xcore/kernel/observability/health.py b/xcore/kernel/observability/health.py index 18973c62..4fc9c877 100644 --- a/xcore/kernel/observability/health.py +++ b/xcore/kernel/observability/health.py @@ -60,6 +60,10 @@ def decorator(fn: Callable) -> Callable: return decorator + def unregister(self, name: str) -> bool: + """Retire un health check précédemment enregistré. Renvoie True s'il existait.""" + return self._checks.pop(name, None) is not None + async def run_all(self, timeout: float = 5.0) -> dict[str, Any]: results: list[CheckResult] = [] for name, (fn, is_async) in self._checks.items(): diff --git a/xcore/kernel/permissions/engine.py b/xcore/kernel/permissions/engine.py index 65da6bd1..c3ba0c39 100644 --- a/xcore/kernel/permissions/engine.py +++ b/xcore/kernel/permissions/engine.py @@ -38,11 +38,17 @@ class PermissionEngine: ``` """ - def __init__(self, events=None, max_audit=100_000) -> None: + def __init__( + self, events=None, max_audit=100_000, audit_cache_hits: bool = True + ) -> None: self._policies: dict[str, PolicySet] = {} self._events = events self._audit_log: deque[dict] = deque(maxlen=max_audit) self._cache: dict[tuple[str, str, str], PolicyEffect] = {} + # audit_cache_hits=False : n'ajoute plus d'entrée au journal pour les + # cache hits (le check() reste identique, seul le append() est sauté). + # Défaut True = comportement historique inchangé (journal complet). + self._audit_cache_hits = audit_cache_hits def load_from_manifest( self, plugin_name: str, raw_permissions: list[dict] | None @@ -78,7 +84,9 @@ def check(self, plugin_name: str, resource: str, action: str) -> None: else: # Cache hit: minimal audit (log entry only, no events) # This keeps audit_log complete while being fast - self._audit(plugin_name, resource, action, effect, emit_event=False) + self._audit( + plugin_name, resource, action, effect, emit_event=False, cache_hit=True + ) if effect == PolicyEffect.DENY: raise PermissionDenied( @@ -92,7 +100,9 @@ def allows(self, plugin_name: str, resource: str, action: str) -> bool: if effect is not None: # Audit even on cache hit for allows() to keep log complete - self._audit(plugin_name, resource, action, effect, emit_event=False) + self._audit( + plugin_name, resource, action, effect, emit_event=False, cache_hit=True + ) return effect == PolicyEffect.ALLOW try: @@ -128,6 +138,7 @@ def _audit( action: str, effect: PolicyEffect, emit_event: bool = True, + cache_hit: bool = False, ) -> None: entry = { "plugin": plugin_name, @@ -135,7 +146,8 @@ def _audit( "action": action, "effect": effect.value, } - self._audit_log.append(entry) + if not cache_hit or self._audit_cache_hits: + self._audit_log.append(entry) # Only log warning or emit events on miss or deny if effect == PolicyEffect.DENY: diff --git a/xcore/kernel/runtime/lifecycle.py b/xcore/kernel/runtime/lifecycle.py index 9ce11f01..9d6c3a3a 100644 --- a/xcore/kernel/runtime/lifecycle.py +++ b/xcore/kernel/runtime/lifecycle.py @@ -24,6 +24,7 @@ from ..api.context import PluginContext from ..api.contract import BasePlugin from ..observability import get_logger +from .plugin_gc import PluginResourceTracker from .state_machine import PluginState, StateMachine logger = get_logger("xcore.runtime.lifecycle") @@ -64,6 +65,13 @@ def __init__( self.plugin_router: Any | None = None self.plugin_middlewares: dict[Any] = {} + # Ramasse-miette forcé : ce que le plugin a enregistré (jobs, health + # checks, abonnements events/hooks) via le PluginContext, et les + # tâches de fond créées via ctx.spawn_task(). Nettoyés de force au + # unload, indépendamment de ce que fait on_unload/on_stop. + self._resource_tracker: PluginResourceTracker | None = None + self._spawned_tasks: list[asyncio.Task] = [] + self._sm = StateMachine( manifest.name, on_change=self._on_state_change, @@ -184,6 +192,21 @@ async def _do_load(self) -> None: isolate_scheduler=tenancy.isolate_scheduler, ) + # Ramasse-miette forcé : on intercepte scheduler/health/events/hooks + # pour mémoriser ce que CE plugin y enregistre pendant on_load, afin + # de tout désinscrire au unload sans dépendre de on_unload/on_stop. + tracker = PluginResourceTracker(self.manifest.name) + ctx.events = tracker.wrap_events(ctx.events) + ctx.hooks = tracker.wrap_hooks(ctx.hooks) + ctx.health = tracker.wrap_health(ctx.health) + if ctx.services.get("scheduler") is not None: + ctx.services = dict(ctx.services) + ctx.services["scheduler"] = tracker.wrap_scheduler( + ctx.services["scheduler"] + ) + self._resource_tracker = tracker + ctx._task_sink = self._spawned_tasks + if hasattr(self._instance, "_inject_context"): await self._instance._inject_context(ctx) elif hasattr(self._instance, "env_variable"): @@ -309,7 +332,32 @@ async def unload(self) -> None: async def _do_unload(self) -> None: if self._instance: - await self._invoke_hooks(["on_stop", "on_unload"]) + # Best-effort : on essaie les hooks du plugin, mais une erreur ici + # ne doit pas empêcher le ramasse-miette forcé ci-dessous — c'est + # justement le filet de sécurité pour un on_unload/on_stop bâclé. + try: + await self._invoke_hooks(["on_stop", "on_unload"]) + except Exception as e: + logger.error( + "on_stop/on_unload hook failed, forcing cleanup anyway", + plugin=self.manifest.name, + error=str(e), + ) + + # Ramasse-miette forcé : libère tout ce que le plugin a enregistré, + # qu'il l'ait fait proprement dans ses hooks ou pas. + if self._resource_tracker is not None: + self._resource_tracker.cleanup() + self._resource_tracker = None + + for task in self._spawned_tasks: + if not task.done(): + task.cancel() + self._spawned_tasks = [] + + if self._registry is not None: + self._registry.unregister(self.manifest.name) + module_name = f"xcore_plugin_{self.manifest.name}" # Nettoie le module principal et le package namespace sys.modules.pop(f"{module_name}.main", None) @@ -408,17 +456,46 @@ def propagate_services(self, *, is_reload: bool = False) -> dict: svc_meta = manifest_services_config.get(name, {}) scope = svc_meta.get("scope", "public") - # register_service lèvera une PermissionError si le service est protégé - self._registry.register_service( - plugin_name=self.manifest.name, - service_name=name, - service_obj=obj, - metadata={ - "reloaded": is_reload, - "scope": scope, - "description": svc_meta.get("description", ""), - }, - ) + try: + # register_service lève PermissionError si le service est protégé + self._registry.register_service( + plugin_name=self.manifest.name, + service_name=name, + service_obj=obj, + metadata={ + "reloaded": is_reload, + "scope": scope, + "description": svc_meta.get("description", ""), + }, + ) + except PermissionError: + # TrustedBase expose tout `ctx.services` (y compris db/cache/ + # scheduler) via `self._services` pour la rétro-compatibilité — + # ce ne sont pas forcément des services que CE plugin exporte, + # potentiellement juste ceux qu'il a reçus en injection. Au + # premier boot le registre ne les protège pas encore + # (register_core_service() n'a lieu qu'après load_all()), donc + # ça passe ; mais dès qu'on recharge/réactive le même plugin + # plus tard, le nom est déjà protégé par le noyau. Si c'est + # bien le même objet reçu en injection, ce n'est pas une + # tentative d'écrasement — on l'ignore. Si c'est un objet + # différent, c'est une vraie tentative malveillante : on + # relève l'erreur telle quelle. `obj` peut être un proxy de + # ramasse-miette (_ScopedScheduler, etc.) posé par ce même + # LifecycleManager autour du service réel — on déballe un + # niveau (`_real`) avant de comparer, sinon l'identité ne + # matcherait jamais pour un service ainsi enveloppé. + underlying = getattr(obj, "_real", obj) + if self._registry.is_registered_as( + name, obj + ) or self._registry.is_registered_as(name, underlying): + logger.debug( + "skip re-registering kernel-protected service", + plugin=self.manifest.name, + service=name, + ) + else: + raise else: # Fallback de sécurité si le registre est absent (pour les tests ou configs minimales) # On définit une liste minimale de services à protéger diff --git a/xcore/kernel/runtime/loader.py b/xcore/kernel/runtime/loader.py index 0aa9925a..123baf21 100644 --- a/xcore/kernel/runtime/loader.py +++ b/xcore/kernel/runtime/loader.py @@ -25,6 +25,7 @@ from ...kernel.observability import get_logger from ...kernel.security.validation import ManifestValidator +from ...registry import PluginStateStore from ..api.contract import PluginHandler from .activator import ( ActivatorRegistry, @@ -83,6 +84,27 @@ def __init__( self._validator = ManifestValidator() + # Table de vérité persistante actif/inactif — survit aux redémarrages. + # Défaut : /../.xcore/plugins_state.json, aucune config requise. + plugins_dir = Path(self._config.directory) + state_path = plugins_dir.parent / ".xcore" / "plugins_state.json" + self._state_store = PluginStateStore(state_path) + + @property + def state_store(self) -> PluginStateStore: + return self._state_store + + def discover_names(self) -> list[str]: + """Liste tous les dossiers de plugins présents sur disque (chargés ou non).""" + plugin_dir = Path(self._config.directory) + if not plugin_dir.exists(): + return [] + return sorted( + d.name + for d in plugin_dir.iterdir() + if d.is_dir() and not d.name.startswith("_") + ) + # ── Chargement global ───────────────────────────────────── async def load_all(self) -> dict[str, list[str]]: @@ -95,16 +117,21 @@ async def load_all(self) -> dict[str, list[str]]: loaded: list[str] = [] failed: list[str] = [] skipped: list[str] = [] + disabled: list[str] = [] manifests = [] plugin_dir = Path(self._config.directory) if not plugin_dir.exists(): logger.warning("plugin directory not found", path=str(plugin_dir)) - return {"loaded": [], "failed": [], "skipped": []} + return {"loaded": [], "failed": [], "skipped": [], "disabled": []} for d in sorted(plugin_dir.iterdir()): if not d.is_dir() or d.name.startswith("_"): continue + if not self._state_store.is_enabled(d.name): + logger.info("plugin disabled, skipped", plugin=d.name) + disabled.append(d.name) + continue try: manifest, validate_version, frameversion = ( self._validator.load_and_validate(d) @@ -123,7 +150,12 @@ async def load_all(self) -> dict[str, list[str]]: skipped.append(d.name) if not manifests: - return {"loaded": [], "failed": [], "skipped": skipped} + return { + "loaded": [], + "failed": [], + "skipped": skipped, + "disabled": disabled, + } try: ordered = _topo_sort(manifests) @@ -133,6 +165,7 @@ async def load_all(self) -> dict[str, list[str]]: "loaded": [], "failed": [m.name for m in manifests], "skipped": skipped, + "disabled": disabled, } # FIX #2 : deux ensembles distincts — "chargé avec succès" vs "traité" @@ -211,8 +244,14 @@ async def load_all(self) -> dict[str, list[str]]: loaded=len(loaded), failed=len(failed), skipped=len(skipped), + disabled=len(disabled), ) - return {"loaded": loaded, "failed": failed, "skipped": skipped} + return { + "loaded": loaded, + "failed": failed, + "skipped": skipped, + "disabled": disabled, + } async def _try_load(self, manifest: Any) -> tuple[Any, bool]: try: diff --git a/xcore/kernel/runtime/plugin_gc.py b/xcore/kernel/runtime/plugin_gc.py new file mode 100644 index 00000000..7df83234 --- /dev/null +++ b/xcore/kernel/runtime/plugin_gc.py @@ -0,0 +1,254 @@ +""" +plugin_gc.py — Ramasse-miette forcé pour les plugins Trusted. + +On ne peut pas se fier à ce qu'un plugin nettoie correctement ce qu'il a +enregistré pendant on_load, dans son `on_unload`/`on_stop` — ces hooks sont +écrits par des auteurs tiers, avec des garanties variables. Ce module +intercepte, au moment où le plugin obtient scheduler/health/events/hooks +via son PluginContext, ce qu'il y enregistre (job scheduler, health check, +abonnement event/hook) — pour pouvoir tout désinscrire de force au unload, +que le plugin l'ait fait proprement ou pas. + +Usage (dans LifecycleManager._do_load) : + tracker = PluginResourceTracker(manifest.name) + ctx.events = tracker.wrap_events(ctx.events) + ctx.hooks = tracker.wrap_hooks(ctx.hooks) + ctx.health = tracker.wrap_health(ctx.health) + if "scheduler" in ctx.services: + ctx.services = dict(ctx.services) + ctx.services["scheduler"] = tracker.wrap_scheduler(ctx.services["scheduler"]) + +Puis, dans LifecycleManager._do_unload : + tracker.cleanup() +""" + +from __future__ import annotations + +import contextlib +from typing import Any, Callable + +from ..observability import get_logger + +logger = get_logger("xcore.runtime.plugin_gc") + + +class _ScopedScheduler: + """Proxy autour de SchedulerService : mémorise les job_id créés par CE plugin.""" + + def __init__(self, real: Any) -> None: + self._real = real + self._job_ids: set[str] = set() + + def add_job( + self, + func: Callable, + trigger: str = "cron", + job_id: str | None = None, + **kwargs: Any, + ) -> Any: + effective_id = job_id or getattr(func, "__name__", repr(func)) + self._job_ids.add(effective_id) + return self._real.add_job(func, trigger=trigger, job_id=job_id, **kwargs) + + def cron(self, expression: str, job_id: str | None = None) -> Callable: + def decorator(fn: Callable) -> Callable: + self._job_ids.add(job_id or fn.__name__) + return self._real.cron(expression, job_id=job_id)(fn) + + return decorator + + def interval(self, **kwargs: Any) -> Callable: + def decorator(fn: Callable) -> Callable: + self._job_ids.add(fn.__name__) + return self._real.interval(**kwargs)(fn) + + return decorator + + def cleanup(self, plugin_name: str) -> None: + for job_id in self._job_ids: + with contextlib.suppress(Exception): + self._real.remove_job(job_id) + if self._job_ids: + logger.debug( + "scheduler jobs released", + plugin=plugin_name, + jobs=sorted(self._job_ids), + ) + + def __getattr__(self, name: str) -> Any: + return getattr(self._real, name) + + +class _ScopedHealth: + """Proxy autour de HealthChecker : mémorise les checks enregistrés par CE plugin.""" + + def __init__(self, real: Any) -> None: + self._real = real + self._names: set[str] = set() + + def register(self, name: str) -> Callable: + def decorator(fn: Callable) -> Callable: + self._names.add(name) + return self._real.register(name)(fn) + + return decorator + + def cleanup(self, plugin_name: str) -> None: + for name in self._names: + with contextlib.suppress(Exception): + self._real.unregister(name) + if self._names: + logger.debug( + "health checks released", plugin=plugin_name, checks=sorted(self._names) + ) + + def __getattr__(self, name: str) -> Any: + return getattr(self._real, name) + + +class _ScopedEvents: + """Proxy autour de l'EventBus (.on/.subscribe/.once) : mémorise (event_name, handler).""" + + def __init__(self, real: Any) -> None: + self._real = real + self._pairs: list[tuple[str, Callable]] = [] + + def on( + self, event_name: str, priority: int = 50, name: str | None = None + ) -> Callable: + def decorator(fn: Callable) -> Callable: + self._pairs.append((event_name, fn)) + return self._real.on(event_name, priority=priority, name=name)(fn) + + return decorator + + def once(self, event_name: str, priority: int = 50) -> Callable: + def decorator(fn: Callable) -> Callable: + self._pairs.append((event_name, fn)) + return self._real.once(event_name, priority=priority)(fn) + + return decorator + + def subscribe( + self, + event_name: str, + handler: Callable, + priority: int = 50, + once: bool = False, + name: str | None = None, + ) -> None: + self._pairs.append((event_name, handler)) + self._real.subscribe( + event_name, handler, priority=priority, once=once, name=name + ) + + def cleanup(self, plugin_name: str) -> None: + for event_name, handler in self._pairs: + with contextlib.suppress(Exception): + self._real.unsubscribe(event_name, handler) + if self._pairs: + logger.debug( + "event subscriptions released", + plugin=plugin_name, + count=len(self._pairs), + ) + + def __getattr__(self, name: str) -> Any: + return getattr(self._real, name) + + +class _ScopedHooks: + """Proxy autour du HookManager (.on/.register/.once) : mémorise (event_name, func).""" + + def __init__(self, real: Any) -> None: + self._real = real + self._pairs: list[tuple[str, Callable]] = [] + + def register( + self, + event_name: str, + func: Callable, + priority: int = 50, + once: bool = False, + timeout: float | None = None, + ) -> Callable: + self._pairs.append((event_name, func)) + return self._real.register( + event_name, func, priority=priority, once=once, timeout=timeout + ) + + def on( + self, + event_name: str, + priority: int = 50, + once: bool = False, + timeout: float | None = None, + ) -> Callable: + def wrapper(func: Callable) -> Callable: + self._pairs.append((event_name, func)) + return self._real.on( + event_name, priority=priority, once=once, timeout=timeout + )(func) + + return wrapper + + def once( + self, event_name: str, priority: int = 50, timeout: float | None = None + ) -> Callable: + return self.on(event_name, priority=priority, once=True, timeout=timeout) + + def cleanup(self, plugin_name: str) -> None: + for event_name, func in self._pairs: + with contextlib.suppress(Exception): + self._real.unregister(event_name, func) + if self._pairs: + logger.debug("hooks released", plugin=plugin_name, count=len(self._pairs)) + + def __getattr__(self, name: str) -> Any: + return getattr(self._real, name) + + +class PluginResourceTracker: + """ + Regroupe les proxies de scoping-par-plugin et leur nettoyage forcé au unload. + + Une instance par (re)chargement de plugin — jetée après le cleanup qui + suit chaque unload/reload. + """ + + def __init__(self, plugin_name: str) -> None: + self._plugin_name = plugin_name + self._scoped: list[Any] = [] + + def wrap_scheduler(self, real: Any) -> Any: + if real is None: + return real + proxy = _ScopedScheduler(real) + self._scoped.append(proxy) + return proxy + + def wrap_health(self, real: Any) -> Any: + if real is None: + return real + proxy = _ScopedHealth(real) + self._scoped.append(proxy) + return proxy + + def wrap_events(self, real: Any) -> Any: + if real is None: + return real + proxy = _ScopedEvents(real) + self._scoped.append(proxy) + return proxy + + def wrap_hooks(self, real: Any) -> Any: + if real is None: + return real + proxy = _ScopedHooks(real) + self._scoped.append(proxy) + return proxy + + def cleanup(self) -> None: + for proxy in self._scoped: + proxy.cleanup(self._plugin_name) + self._scoped.clear() diff --git a/xcore/kernel/runtime/supervisor.py b/xcore/kernel/runtime/supervisor.py index 5e54c66a..c76051af 100644 --- a/xcore/kernel/runtime/supervisor.py +++ b/xcore/kernel/runtime/supervisor.py @@ -8,6 +8,8 @@ from __future__ import annotations import contextlib +import time +from collections import deque from typing import TYPE_CHECKING, Any if TYPE_CHECKING: @@ -59,6 +61,11 @@ def __init__(self, ctx: "KernelContext") -> None: self._rate = RateLimiterRegistry() self._permissions = PermissionEngine(events=self._events) + # Supervision des appels IPC — même pattern que PermissionEngine._audit_log : + # un journal borné pour voir QUI a appelé QUOI, quand, avec quel résultat. + self._ipc_audit: deque[dict] = deque(maxlen=100_000) + self._ipc_stats: dict[str, dict] = {} + self._loader: PluginLoader | None = None self._pipeline: MiddlewarePipeline | None = None @@ -153,6 +160,7 @@ async def boot(self) -> None: loaded=len(report["loaded"]), failed=len(report["failed"]), skipped=len(report["skipped"]), + disabled=len(report.get("disabled", [])), ) if self._events: @@ -241,11 +249,14 @@ async def call( return self._err("Supervisor non démarré", "not_ready") if not self._loader.has(plugin_name): - return self._err(f"Plugin '{plugin_name}' introuvable", "not_found") + result = self._err(f"Plugin '{plugin_name}' introuvable", "not_found") + self._audit_ipc_call(plugin_name, action, caller, tenant_id, result, 0.0) + return result handler = self._loader.get(plugin_name) - return await self._pipeline.execute( + t0 = time.monotonic() + result = await self._pipeline.execute( plugin_name, action, payload, @@ -254,6 +265,43 @@ async def call( caller=caller, tenant_id=tenant_id, ) + duration_ms = (time.monotonic() - t0) * 1000 + self._audit_ipc_call( + plugin_name, action, caller, tenant_id, result, duration_ms + ) + return result + + def _audit_ipc_call( + self, + plugin_name: str, + action: str, + caller: str | None, + tenant_id: str, + result: dict, + duration_ms: float, + ) -> None: + """Journalise chaque appel IPC — qui a appelé quoi, quand, avec quel résultat.""" + status = result.get("status", "ok") if isinstance(result, dict) else "ok" + code = result.get("code") if isinstance(result, dict) else None + entry = { + "plugin": plugin_name, + "action": action, + "caller": caller, + "tenant_id": tenant_id, + "status": status, + "code": code, + "duration_ms": round(duration_ms, 2), + "timestamp": time.time(), + } + self._ipc_audit.append(entry) + + stats = self._ipc_stats.setdefault( + plugin_name, {"calls": 0, "errors": 0, "total_duration_ms": 0.0} + ) + stats["calls"] += 1 + stats["total_duration_ms"] += duration_ms + if status == "error": + stats["errors"] += 1 async def _dispatch( self, plugin_name: str, action: str, payload: dict, handler, **kwargs @@ -326,6 +374,59 @@ async def unload(self, plugin_name: str) -> None: if self._events: await self._events.emit(f"plugin.{plugin_name}.unloaded", {}) + # ── Table de vérité actif/inactif ────────────────────────── + + async def enable(self, plugin_name: str) -> None: + """ + Active un plugin de façon persistante : il sera aussi chargé aux + prochains redémarrages. Le charge immédiatement s'il ne l'est pas déjà. + """ + if self._loader is None: + raise RuntimeError("Le superviseur n'est pas encore démarré.") + self._loader.state_store.set_enabled(plugin_name, True) + if not self._loader.has(plugin_name): + await self.load(plugin_name) + + async def disable(self, plugin_name: str, reason: str | None = None) -> None: + """ + Désactive un plugin de façon persistante : il ne sera plus chargé aux + prochains redémarrages. Le décharge immédiatement s'il est chargé + (ramasse-miette forcé via LifecycleManager._do_unload). + """ + if self._loader is None: + raise RuntimeError("Le superviseur n'est pas encore démarré.") + self._loader.state_store.set_enabled(plugin_name, False, reason=reason) + if self._loader.has(plugin_name): + await self.unload(plugin_name) + + def registry_table(self) -> list[dict]: + """ + Table de vérité complète : tous les plugins présents sur disque, avec + leur flag persisté (enabled) et leur état live (ready/unloaded/...). + Contrairement à status(), inclut aussi les plugins désactivés ou + jamais chargés — la vraie vue « centre de contrôle ». + """ + if self._loader is None: + return [] + + live_status = {s["name"]: s for s in self._loader.status()} + persisted = self._loader.state_store.all() + + rows = [] + for name in self._loader.discover_names(): + live = live_status.get(name) + rows.append( + { + "name": name, + "enabled": self._loader.state_store.is_enabled(name), + "state": live["state"] if live else "not_loaded", + "mode": live.get("mode") if live else None, + "loaded": live is not None, + "reason": persisted.get(name, {}).get("reason"), + } + ) + return rows + # ── Observabilité ───────────────────────────────────────── def status(self) -> dict: @@ -353,6 +454,56 @@ def permissions_audit( """Retourne le journal d'audit des permissions.""" return self._permissions.audit_log(plugin_name, limit) + def ipc_audit(self, plugin_name: str | None = None, limit: int = 100) -> list[dict]: + """ + Journal des appels IPC — qui a appelé quoi, quand, avec quel résultat. + Même pattern que permissions_audit(), le plus récent en dernier. + """ + from itertools import islice + + it = reversed(self._ipc_audit) + if plugin_name: + it = (e for e in it if e["plugin"] == plugin_name) + results = list(islice(it, limit)) + results.reverse() + return results + + def ipc_stats(self) -> dict: + """Statistiques agrégées par plugin : appels, erreurs, latence moyenne.""" + return { + "entries": len(self._ipc_audit), + "by_plugin": { + name: { + "calls": s["calls"], + "errors": s["errors"], + "avg_duration_ms": ( + round(s["total_duration_ms"] / s["calls"], 2) + if s["calls"] + else 0.0 + ), + } + for name, s in self._ipc_stats.items() + }, + } + + def events_activity(self, event_name: str | None = None, limit: int = 100) -> dict: + """Activité récente + métriques de l'EventBus (ctx.events).""" + if self._events is None: + return {"recent": [], "stats": {}} + return { + "recent": self._events.recent_emissions(event_name, limit), + "stats": self._events.stats(event_name), + } + + def hooks_activity(self, event_name: str | None = None, limit: int = 100) -> dict: + """Activité récente + métriques du HookManager (ctx.hooks).""" + if self._hooks is None: + return {"recent": [], "metrics": {}} + return { + "recent": self._hooks.recent_emissions(event_name, limit), + "metrics": self._hooks.get_metrics(event_name), + } + # ── Arrêt ───────────────────────────────────────────────── async def shutdown(self) -> None: diff --git a/xcore/kernel/sandbox/process_manager.py b/xcore/kernel/sandbox/process_manager.py index 3a088ed4..a06ba9c4 100644 --- a/xcore/kernel/sandbox/process_manager.py +++ b/xcore/kernel/sandbox/process_manager.py @@ -19,7 +19,7 @@ if TYPE_CHECKING: from ..runtime.loader import PluginLoader -from .ipc import IPCChannel, IPCProcessDead +from .ipc import IPCChannel, IPCProcessDead, IPCTimeoutError from .isolation import DiskQuotaExceeded, DiskWatcher logger = get_logger("xcore.sandbox.process_manager") @@ -178,7 +178,12 @@ async def call(self, action: str, payload: dict) -> dict: try: resp = await self._channel.call(action, payload) return resp.data - except IPCProcessDead: + except (IPCProcessDead, IPCTimeoutError): + # IPCTimeoutError : le subprocess n'est pas forcément mort, juste + # bloqué/sans réponse — sans ce catch, il fallait attendre le + # prochain cycle de _health_loop pour que le recyclage se déclenche. + # On recycle immédiatement sur le chemin de l'appel en échec plutôt + # que de laisser un process inutilisable jusqu'au prochain health check. await self._handle_crash() raise diff --git a/xcore/registry/__init__.py b/xcore/registry/__init__.py index 53789d1f..cd8b7c8b 100644 --- a/xcore/registry/__init__.py +++ b/xcore/registry/__init__.py @@ -1,5 +1,12 @@ from .index import PluginRegistry from .resolver import DependencyResolver +from .state_store import PluginStateStore from .versioning import VersionConstraint, satisfies -__all__ = ["PluginRegistry", "DependencyResolver", "VersionConstraint", "satisfies"] +__all__ = [ + "PluginRegistry", + "DependencyResolver", + "PluginStateStore", + "VersionConstraint", + "satisfies", +] diff --git a/xcore/registry/index.py b/xcore/registry/index.py index 7be41c4c..d9ad112f 100644 --- a/xcore/registry/index.py +++ b/xcore/registry/index.py @@ -135,6 +135,14 @@ def list_services(self) -> list[dict]: for name, meta in self._exported_services.items() ] + def is_registered_as(self, service_name: str, obj: Any) -> bool: + """Vrai si `service_name` est déjà exporté et pointe vers CE MÊME objet + (identité, pas égalité) — utile pour distinguer une simple ré-injection + du service reçu (bénin) d'une véritable tentative d'écrasement par un + objet différent (à bloquer).""" + existing = self._exported_services.get(service_name) + return existing is not None and existing.get("obj") is obj + def has(self, name: str) -> bool: return name in self._entries diff --git a/xcore/registry/state_store.py b/xcore/registry/state_store.py new file mode 100644 index 00000000..01d944ea --- /dev/null +++ b/xcore/registry/state_store.py @@ -0,0 +1,79 @@ +""" +state_store.py — Table de vérité persistante pour l'état actif/inactif des plugins. + +Contrairement au registre en mémoire (PluginRegistry), cet état survit aux +redémarrages : un plugin désactivé via `PluginSupervisor.disable()` reste +désactivé après un restart du process, car `PluginLoader.load_all()` le +consulte avant de charger quoi que ce soit. + +Stockage : un unique fichier JSON, écrit de façon atomique (fichier temporaire ++ remplacement) pour éviter toute corruption en cas de crash pendant l'écriture. +""" + +from __future__ import annotations + +import json +import os +import time +from pathlib import Path +from typing import Any + +from ..kernel.observability import get_logger + +logger = get_logger("xcore.registry.state_store") + + +class PluginStateStore: + """ + Table de vérité `nom_plugin -> {enabled, reason, updated_at}`. + + Usage: + store = PluginStateStore(Path("plugins").parent / ".xcore" / "plugins_state.json") + store.is_enabled("shop") # True par défaut si jamais togglé + store.set_enabled("shop", False, reason="maintenance") + store.all() + """ + + def __init__(self, path: Path) -> None: + self._path = path + self._data: dict[str, dict[str, Any]] = {} + self._load() + + def _load(self) -> None: + if not self._path.exists(): + return + try: + self._data = json.loads(self._path.read_text(encoding="utf-8")) or {} + except (json.JSONDecodeError, OSError) as e: + logger.warning( + "plugin state file unreadable, starting empty", + path=str(self._path), + error=str(e), + ) + self._data = {} + + def _persist(self) -> None: + self._path.parent.mkdir(parents=True, exist_ok=True) + tmp_path = self._path.with_suffix(f"{self._path.suffix}.tmp") + tmp_path.write_text( + json.dumps(self._data, indent=2, sort_keys=True), encoding="utf-8" + ) + os.replace(tmp_path, self._path) + + def is_enabled(self, name: str, default: bool = True) -> bool: + entry = self._data.get(name) + return default if entry is None else bool(entry.get("enabled", default)) + + def set_enabled(self, name: str, enabled: bool, reason: str | None = None) -> None: + self._data[name] = { + "enabled": enabled, + "reason": reason, + "updated_at": time.time(), + } + self._persist() + logger.info( + "plugin state persisted", plugin=name, enabled=enabled, reason=reason + ) + + def all(self) -> dict[str, dict[str, Any]]: + return dict(self._data) From 2ee8814339c78acf41e066defae15a204cc49b9a Mon Sep 17 00:00:00 2001 From: traoreera Date: Mon, 28 Sep 2026 14:12:55 +0000 Subject: [PATCH 2/2] =?UTF-8?q?chore(release):=20v2.6.0=20=E2=80=94=20bump?= =?UTF-8?q?,=20d=C3=A9pendances=20all=C3=A9g=C3=A9es,=20roadmap=20resynchr?= =?UTF-8?q?onis=C3=A9e?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - xcore/__version__.py était figé à 2.3.3 (désynchronisé de pyproject.toml depuis plusieurs releases) — aligné sur 2.6.0 - pyproject.toml : retrait de uvicorn, pydantic-settings, rich (zéro usage dans xcore/, déjà couverts transitivement par fastapi[standard] pour les deux premiers) ; pydantic[email] → pydantic (EmailStr jamais utilisé). poetry.lock régénéré en conséquence, aucun changement de comportement. - roadmap/executed_roadmap.md (FR) affichait V3 60% / V4 15% alors que son propre tableau détaillé (identique à la version EN) ne justifie que 25% / 5% — recompté et aligné sur ROADMAP_PROGRESS.md --- CHANGELOG.md | 3 +++ doc/changelog.md | 6 ++++++ poetry.lock | 2 +- pyproject.toml | 7 ++----- roadmap/executed_roadmap.md | 4 ++-- xcore/__version__.py | 4 ++-- 6 files changed, 16 insertions(+), 10 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 29db1fe8..dbf22ce2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -21,6 +21,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - **`PermissionEngine` audit log had no way to skip cache-hit entries**: every `allows()`/`check()` cache hit still appended to `_audit_log` unconditionally — the expensive part (event emission) was already skipped on cache hits, but the log append wasn't. New optional `PermissionEngine(audit_cache_hits=False)` skips it; default (`True`) keeps the existing behavior (complete audit trail) unchanged. - **Stray temp directories from crashed test runs**: `tests/conftest.py`'s `plugins_dir`/`temp_dir` fixtures already clean up via `yield` + `shutil.rmtree`, but that teardown never runs if a test crashes hard (e.g. `SIGKILL`) before reaching it. `temp_dir` now uses the same distinguishing `xcore_test_` prefix as `plugins_dir`, and a new session-scoped autouse fixture sweeps any `xcore_test_*` directories left behind in the system temp dir at the end of the run — scoped to that exact prefix only, never a broader temp-dir sweep. +### Changed +- **Trimmed unused core dependencies**: `uvicorn` and `pydantic-settings` were never imported anywhere in `xcore` — both are already pulled in transitively by `fastapi[standard]` (confirmed against its own metadata) for anyone who needs them. `rich` had zero usage in the core package (it belongs to `xcorecli`, a separate install). `pydantic[email]` is now plain `pydantic` — `EmailStr`/`email-validator` were never used, and `fastapi[standard]` already brings in `email-validator` regardless. No behavior change; `poetry.lock` regenerated to match. + ## [2.5.1] - 2026-08-20 ### Fixed diff --git a/doc/changelog.md b/doc/changelog.md index 017e541e..dbf22ce2 100644 --- a/doc/changelog.md +++ b/doc/changelog.md @@ -17,6 +17,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed - **`propagate_services()` broke reload/re-enable of a `TrustedBase` plugin**: `self._services` exposes the entire `ctx.services` dict for backward compatibility (including db/cache/scheduler), and `propagate_services()` tried to re-register those as the plugin's own exports. This passed on first boot (the registry doesn't protect core services until after `load_all()` runs), but any later reload raised `PermissionError: Impossible d'écraser le service protégé`. A collision on an object identical to the one already protected (received via injection, not exported) is now ignored; a genuinely different object (an actual override attempt) still raises. +- **Sandboxed subprocess wasn't recycled on IPC timeout**: `SandboxProcessManager.call()` only caught `IPCProcessDead`, not `IPCTimeoutError` — a subprocess that stopped responding without actually dying stayed unusable until the next periodic `_health_loop` check caught up with it. `call()` now recycles the subprocess immediately on either exception, on the failing request's own path, instead of waiting for the next health-check interval. +- **`PermissionEngine` audit log had no way to skip cache-hit entries**: every `allows()`/`check()` cache hit still appended to `_audit_log` unconditionally — the expensive part (event emission) was already skipped on cache hits, but the log append wasn't. New optional `PermissionEngine(audit_cache_hits=False)` skips it; default (`True`) keeps the existing behavior (complete audit trail) unchanged. +- **Stray temp directories from crashed test runs**: `tests/conftest.py`'s `plugins_dir`/`temp_dir` fixtures already clean up via `yield` + `shutil.rmtree`, but that teardown never runs if a test crashes hard (e.g. `SIGKILL`) before reaching it. `temp_dir` now uses the same distinguishing `xcore_test_` prefix as `plugins_dir`, and a new session-scoped autouse fixture sweeps any `xcore_test_*` directories left behind in the system temp dir at the end of the run — scoped to that exact prefix only, never a broader temp-dir sweep. + +### Changed +- **Trimmed unused core dependencies**: `uvicorn` and `pydantic-settings` were never imported anywhere in `xcore` — both are already pulled in transitively by `fastapi[standard]` (confirmed against its own metadata) for anyone who needs them. `rich` had zero usage in the core package (it belongs to `xcorecli`, a separate install). `pydantic[email]` is now plain `pydantic` — `EmailStr`/`email-validator` were never used, and `fastapi[standard]` already brings in `email-validator` regardless. No behavior change; `poetry.lock` regenerated to match. ## [2.5.1] - 2026-08-20 diff --git a/poetry.lock b/poetry.lock index 4d1635e4..54fe27b2 100644 --- a/poetry.lock +++ b/poetry.lock @@ -4525,4 +4525,4 @@ xcli = ["xcorecli"] [metadata] lock-version = "2.1" python-versions = ">=3.12,<4.0" -content-hash = "4351f42149ca3bd14affc6c5088c32478f5f29f23056f768ac5da484f5174331" +content-hash = "b6f25778d5af0083617a5b268cc677107c5a99b374f613fcb568aff41e247bb0" diff --git a/pyproject.toml b/pyproject.toml index b6471382..60e03387 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "XCoreRuntime" -version = "2.5.2" +version = "2.6.0" description = "Plugin-first orchestration framework built on FastAPI" authors = [ { name = "Eliezer Traore", email = "68350805+traoreera@users.noreply.github.com" }, @@ -18,9 +18,7 @@ classifiers = [ requires-python = ">=3.12,<4.0" dependencies = [ "fastapi[standard]>=0.135.1,<1.0.0", - "uvicorn>=0.38.0,<1.0.0", - "pydantic[email]>=2.11.7,<3.0.0", - "pydantic-settings>=2.14.2,<3.0.0", + "pydantic>=2.11.7,<3.0.0", "sqlalchemy>=2.0.41,<3.0.0", "alembic>=1.16.1,<2.0.0", "psycopg2>=2.9.11,<3.0.0", @@ -29,7 +27,6 @@ dependencies = [ "apscheduler>=3.11.0,<4.0.0", "python-dotenv>=1.1.0,<2.0.0", "pyyaml>=6.0.3,<7.0.0", - "rich>=14.0.0,<15.0.0", "celery>=5.6.3,<6.0.0", "opentelemetry-api>=1.27.0,<2.0.0", "opentelemetry-sdk>=1.27.0,<2.0.0", diff --git a/roadmap/executed_roadmap.md b/roadmap/executed_roadmap.md index 01f8bafd..2e1204a9 100644 --- a/roadmap/executed_roadmap.md +++ b/roadmap/executed_roadmap.md @@ -8,8 +8,8 @@ Ce document présente l'état actuel du framework XCore par rapport aux objectif | :--- | :--- | :--- | :--- | | **V1** | Fondation Kernel | **Terminé** | 100% | | **V2** | Industrialisation | **Terminé** | 100% | -| **V3** | Distribution | **Avancé** | 60% | -| **V4** | Cloud Native | **Démarré** | 15% | +| **V3** | Distribution | **Avancé** | 25% | +| **V4** | Cloud Native | **Démarré** | 5% | | **V5** | Intelligence Native | **Concept** | 0% | --- diff --git a/xcore/__version__.py b/xcore/__version__.py index 1d85f72b..0be1fae4 100644 --- a/xcore/__version__.py +++ b/xcore/__version__.py @@ -1,4 +1,4 @@ -__version__ = "2.3.3" -__version_info__ = (2, 3, 3) +__version__ = "2.6.0" +__version_info__ = (2, 6, 0) __author__ = "xcore contributors" __license__ = "MIT"