From 9de145efd4c4e159bca779ba445a0fef08315700 Mon Sep 17 00:00:00 2001 From: Claude Date: Mon, 13 Apr 2026 21:55:51 +0000 Subject: [PATCH 1/4] feat: implement event bus with pyee for inter-module communication Replace the custom EventBus with a pyee-backed AsyncIOEventEmitter. Products module now publishes ProductCreated/Updated/Deleted events, and Dashboard subscribes to them to track activity stats via a new /api/dashboard/stats endpoint. https://claude.ai/code/session_01KLjCnJoNbCbk8is6pVZcom --- framework/core/pyproject.toml | 1 + .../core/src/simple_module_core/events.py | 43 +++++++++++++---- modules/dashboard/pyproject.toml | 2 + .../src/sm_dashboard/endpoints/api.py | 17 +++++++ .../dashboard/src/sm_dashboard/handlers.py | 46 +++++++++++++++++++ modules/dashboard/src/sm_dashboard/module.py | 22 +++++++++ .../src/sm_products/contracts/__init__.py | 15 +++++- .../src/sm_products/contracts/events.py | 30 ++++++++++++ modules/products/src/sm_products/deps.py | 7 ++- .../products/src/sm_products/endpoints/api.py | 13 +++++- 10 files changed, 182 insertions(+), 14 deletions(-) create mode 100644 modules/dashboard/src/sm_dashboard/endpoints/api.py create mode 100644 modules/dashboard/src/sm_dashboard/handlers.py create mode 100644 modules/products/src/sm_products/contracts/events.py diff --git a/framework/core/pyproject.toml b/framework/core/pyproject.toml index b55de92b..6055f7c6 100644 --- a/framework/core/pyproject.toml +++ b/framework/core/pyproject.toml @@ -10,6 +10,7 @@ dependencies = [ "fastapi>=0.115", "pydantic>=2.0", "pydantic-settings>=2.0", + "pyee>=12.0", ] [build-system] diff --git a/framework/core/src/simple_module_core/events.py b/framework/core/src/simple_module_core/events.py index 7ca945ca..eb60f1f4 100644 --- a/framework/core/src/simple_module_core/events.py +++ b/framework/core/src/simple_module_core/events.py @@ -1,14 +1,15 @@ -"""Async in-process event bus for inter-module communication.""" +"""Async in-process event bus backed by pyee for inter-module communication.""" from __future__ import annotations import asyncio import logging -from collections import defaultdict from collections.abc import Callable, Coroutine from dataclasses import dataclass from typing import Any +from pyee.asyncio import AsyncIOEventEmitter + logger = logging.getLogger(__name__) EventHandler = Callable[..., Coroutine[Any, Any, None]] @@ -28,18 +29,32 @@ class ProductCreated(Event): class EventBus: - """Simple async event bus. + """Async event bus backed by pyee's ``AsyncIOEventEmitter``. Modules subscribe to event types in ``register_event_handlers``. Publishing dispatches to all subscribers concurrently. + + * ``publish`` — awaits all handlers via ``asyncio.gather`` (error-isolated). + * ``publish_nowait`` — fire-and-forget via pyee's event-loop scheduling. """ def __init__(self) -> None: - self._handlers: dict[type[Event], list[EventHandler]] = defaultdict(list) + self._emitter = AsyncIOEventEmitter() + self._emitter.on("error", self._on_emitter_error) + + @staticmethod + def _on_emitter_error(error: Exception) -> None: + """Handle errors from fire-and-forget dispatch.""" + logger.error("EventBus background handler error: %s", error, exc_info=error) + + @staticmethod + def _event_key(event_type: type[Event]) -> str: + """Derive a unique string key from an event class for pyee registration.""" + return f"{event_type.__module__}.{event_type.__qualname__}" def subscribe(self, event_type: type[Event], handler: EventHandler) -> None: """Register a handler for an event type.""" - self._handlers[event_type].append(handler) + self._emitter.on(self._event_key(event_type), handler) logger.debug( "Subscribed %s to %s", getattr(handler, "__qualname__", repr(handler)), @@ -47,8 +62,13 @@ def subscribe(self, event_type: type[Event], handler: EventHandler) -> None: ) async def publish(self, event: Event) -> None: - """Dispatch event to all registered handlers (awaited).""" - handlers = self._handlers.get(type(event), []) + """Dispatch event to all registered handlers (awaited). + + All handlers run concurrently via ``asyncio.gather``. + Individual handler failures are logged but do not propagate. + """ + key = self._event_key(type(event)) + handlers = self._emitter.listeners(key) if not handlers: return results = await asyncio.gather( @@ -66,6 +86,9 @@ async def publish(self, event: Event) -> None: ) def publish_nowait(self, event: Event) -> None: - """Fire-and-forget: schedule event dispatch on the current event loop.""" - loop = asyncio.get_event_loop() - loop.create_task(self.publish(event)) + """Fire-and-forget: schedule event dispatch on the current event loop. + + Uses pyee's ``AsyncIOEventEmitter.emit`` which schedules async + handlers as tasks on the running loop. + """ + self._emitter.emit(self._event_key(type(event)), event) diff --git a/modules/dashboard/pyproject.toml b/modules/dashboard/pyproject.toml index f0ca5165..985eb54a 100644 --- a/modules/dashboard/pyproject.toml +++ b/modules/dashboard/pyproject.toml @@ -10,6 +10,7 @@ dependencies = [ "simple-module-core", "simple-module-db", "simple-module-hosting", + "sm-products", ] [project.entry-points.simple_module] @@ -23,3 +24,4 @@ build-backend = "hatchling.build" simple-module-core = { workspace = true } simple-module-db = { workspace = true } simple-module-hosting = { workspace = true } +sm-products = { workspace = true } diff --git a/modules/dashboard/src/sm_dashboard/endpoints/api.py b/modules/dashboard/src/sm_dashboard/endpoints/api.py new file mode 100644 index 00000000..8765b715 --- /dev/null +++ b/modules/dashboard/src/sm_dashboard/endpoints/api.py @@ -0,0 +1,17 @@ +"""REST API endpoints for the Dashboard module.""" + +from __future__ import annotations + +from fastapi import APIRouter + +from sm_dashboard.handlers import get_product_event_counts + +router = APIRouter() + + +@router.get("/stats") +async def dashboard_stats() -> dict: + """Return dashboard statistics including product event counts.""" + return { + "product_events": get_product_event_counts(), + } diff --git a/modules/dashboard/src/sm_dashboard/handlers.py b/modules/dashboard/src/sm_dashboard/handlers.py new file mode 100644 index 00000000..a7e86f04 --- /dev/null +++ b/modules/dashboard/src/sm_dashboard/handlers.py @@ -0,0 +1,46 @@ +"""Event handlers for the Dashboard module. + +Subscribes to product domain events to maintain real-time stats +without direct coupling to the Products module's internals. +""" + +from __future__ import annotations + +import logging + +from sm_products.contracts.events import ProductCreated, ProductDeleted, ProductUpdated + +logger = logging.getLogger(__name__) + +_product_event_counts: dict[str, int] = { + "created": 0, + "updated": 0, + "deleted": 0, +} + + +async def on_product_created(event: ProductCreated) -> None: + _product_event_counts["created"] += 1 + logger.info("Dashboard received ProductCreated: %s (id=%d)", event.name, event.product_id) + + +async def on_product_updated(event: ProductUpdated) -> None: + _product_event_counts["updated"] += 1 + logger.info("Dashboard received ProductUpdated: %s (id=%d)", event.name, event.product_id) + + +async def on_product_deleted(event: ProductDeleted) -> None: + _product_event_counts["deleted"] += 1 + logger.info("Dashboard received ProductDeleted: id=%d", event.product_id) + + +def get_product_event_counts() -> dict[str, int]: + """Return a snapshot of product event counts.""" + return dict(_product_event_counts) + + +def reset_product_event_counts() -> None: + """Reset counters — useful for testing.""" + _product_event_counts["created"] = 0 + _product_event_counts["updated"] = 0 + _product_event_counts["deleted"] = 0 diff --git a/modules/dashboard/src/sm_dashboard/module.py b/modules/dashboard/src/sm_dashboard/module.py index c486e5bc..89c08993 100644 --- a/modules/dashboard/src/sm_dashboard/module.py +++ b/modules/dashboard/src/sm_dashboard/module.py @@ -3,6 +3,7 @@ from __future__ import annotations from fastapi import APIRouter +from simple_module_core.events import EventBus from simple_module_core.menu import MenuItem, MenuRegistry, MenuSection from simple_module_core.module import ModuleBase, ModuleMeta @@ -10,12 +11,16 @@ class DashboardModule(ModuleBase): meta = ModuleMeta( name="Dashboard", + route_prefix="/api/dashboard", view_prefix="", + depends_on=["Products"], ) def register_routes(self, api_router: APIRouter, view_router: APIRouter) -> None: + from sm_dashboard.endpoints.api import router as api from sm_dashboard.endpoints.views import router as views + api_router.include_router(api) view_router.include_router(views) def register_menu_items(self, registry: MenuRegistry) -> None: @@ -28,3 +33,20 @@ def register_menu_items(self, registry: MenuRegistry) -> None: section=MenuSection.SIDEBAR, ) ) + + def register_event_handlers(self, bus: EventBus) -> None: + from sm_products.contracts.events import ( + ProductCreated, + ProductDeleted, + ProductUpdated, + ) + + from sm_dashboard.handlers import ( + on_product_created, + on_product_deleted, + on_product_updated, + ) + + bus.subscribe(ProductCreated, on_product_created) + bus.subscribe(ProductUpdated, on_product_updated) + bus.subscribe(ProductDeleted, on_product_deleted) diff --git a/modules/products/src/sm_products/contracts/__init__.py b/modules/products/src/sm_products/contracts/__init__.py index fb12cbed..6d3e23a1 100644 --- a/modules/products/src/sm_products/contracts/__init__.py +++ b/modules/products/src/sm_products/contracts/__init__.py @@ -1,5 +1,10 @@ """Products contracts — public interface for other modules.""" +from sm_products.contracts.events import ( + ProductCreated, + ProductDeleted, + ProductUpdated, +) from sm_products.contracts.schemas import ( ProductCreate, ProductOut, @@ -7,4 +12,12 @@ ) from sm_products.contracts.service import IProductService -__all__ = ["ProductCreate", "ProductOut", "ProductUpdate", "IProductService"] +__all__ = [ + "ProductCreate", + "ProductCreated", + "ProductDeleted", + "ProductOut", + "ProductUpdate", + "ProductUpdated", + "IProductService", +] diff --git a/modules/products/src/sm_products/contracts/events.py b/modules/products/src/sm_products/contracts/events.py new file mode 100644 index 00000000..d73aefc0 --- /dev/null +++ b/modules/products/src/sm_products/contracts/events.py @@ -0,0 +1,30 @@ +"""Product domain events — published when product state changes.""" + +from __future__ import annotations + +from dataclasses import dataclass + +from simple_module_core.events import Event + + +@dataclass +class ProductCreated(Event): + """Published after a new product is persisted.""" + + product_id: int + name: str + + +@dataclass +class ProductUpdated(Event): + """Published after an existing product is modified.""" + + product_id: int + name: str + + +@dataclass +class ProductDeleted(Event): + """Published after a product is removed.""" + + product_id: int diff --git a/modules/products/src/sm_products/deps.py b/modules/products/src/sm_products/deps.py index 7adf4cb9..9cdb07e6 100644 --- a/modules/products/src/sm_products/deps.py +++ b/modules/products/src/sm_products/deps.py @@ -2,7 +2,8 @@ from __future__ import annotations -from fastapi import Depends +from fastapi import Depends, Request +from simple_module_core.events import EventBus from simple_module_db.deps import get_db from sqlalchemy.ext.asyncio import AsyncSession @@ -13,3 +14,7 @@ async def get_product_service( db: AsyncSession = Depends(get_db), ) -> ProductService: return ProductService(db) + + +def get_event_bus(request: Request) -> EventBus: + return request.app.state.event_bus diff --git a/modules/products/src/sm_products/endpoints/api.py b/modules/products/src/sm_products/endpoints/api.py index 654b5edb..2f54fbc6 100644 --- a/modules/products/src/sm_products/endpoints/api.py +++ b/modules/products/src/sm_products/endpoints/api.py @@ -3,9 +3,11 @@ from __future__ import annotations from fastapi import APIRouter, Depends, HTTPException +from simple_module_core.events import EventBus +from sm_products.contracts.events import ProductCreated, ProductDeleted, ProductUpdated from sm_products.contracts.schemas import ProductCreate, ProductOut, ProductUpdate -from sm_products.deps import get_product_service +from sm_products.deps import get_event_bus, get_product_service from sm_products.service import ProductService router = APIRouter() @@ -33,8 +35,11 @@ async def get_product( async def create_product( data: ProductCreate, service: ProductService = Depends(get_product_service), + bus: EventBus = Depends(get_event_bus), ) -> ProductOut: - return await service.create(data) + product = await service.create(data) + await bus.publish(ProductCreated(product_id=product.id, name=product.name)) + return product @router.put("/{product_id}", response_model=ProductOut) @@ -42,10 +47,12 @@ async def update_product( product_id: int, data: ProductUpdate, service: ProductService = Depends(get_product_service), + bus: EventBus = Depends(get_event_bus), ) -> ProductOut: product = await service.update(product_id, data) if product is None: raise HTTPException(status_code=404, detail="Product not found") + await bus.publish(ProductUpdated(product_id=product.id, name=product.name)) return product @@ -53,7 +60,9 @@ async def update_product( async def delete_product( product_id: int, service: ProductService = Depends(get_product_service), + bus: EventBus = Depends(get_event_bus), ) -> None: deleted = await service.delete(product_id) if not deleted: raise HTTPException(status_code=404, detail="Product not found") + await bus.publish(ProductDeleted(product_id=product_id)) From e9e57c4006c0378a64f027a0113b277724cc84dd Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 14 Apr 2026 12:44:07 +0000 Subject: [PATCH 2/4] test: add dashboard event bus tests and EventBus edge cases MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Covers dashboard handlers, /api/dashboard/stats, and end-to-end Product API → EventBus → Dashboard handler wiring. Adds EventBus tests for subclass isolation, orphan events, and concurrent dispatch. https://claude.ai/code/session_01KLjCnJoNbCbk8is6pVZcom --- framework/core/tests/test_core.py | 65 +++++++ modules/dashboard/tests/__init__.py | 0 modules/dashboard/tests/test_dashboard.py | 217 ++++++++++++++++++++++ pyproject.toml | 2 +- 4 files changed, 283 insertions(+), 1 deletion(-) create mode 100644 modules/dashboard/tests/__init__.py create mode 100644 modules/dashboard/tests/test_dashboard.py diff --git a/framework/core/tests/test_core.py b/framework/core/tests/test_core.py index 7063d3c9..d15b1929 100644 --- a/framework/core/tests/test_core.py +++ b/framework/core/tests/test_core.py @@ -479,6 +479,71 @@ async def handler(e): assert len(received) == 1 assert received[0].order_id == 99 + async def test_subclass_events_do_not_match_parent_subscription(self): + """Subscribing to a base Event class should not receive subclass events.""" + bus = EventBus() + + @dataclass + class Parent(Event): + pass + + @dataclass + class Child(Parent): + pass + + calls: list = [] + + async def parent_handler(e): + calls.append(("parent", e)) + + bus.subscribe(Parent, parent_handler) + await bus.publish(Child()) + + # Child events should not trigger Parent handlers — strict type match. + assert calls == [] + + async def test_publish_with_no_subscribers_returns_none(self): + """publish() should resolve to None when nothing is listening.""" + bus = EventBus() + + @dataclass + class Orphan(Event): + pass + + result = await bus.publish(Orphan()) + assert result is None + + async def test_publish_nowait_with_no_subscribers_is_noop(self): + """publish_nowait() on an unheard event should not raise.""" + bus = EventBus() + + @dataclass + class Orphan(Event): + pass + + bus.publish_nowait(Orphan()) # must not raise + + async def test_handlers_dispatched_concurrently(self): + """All handlers for an event should run concurrently via gather.""" + import asyncio + + bus = EventBus() + order: list[str] = [] + + async def slow(e): + await asyncio.sleep(0.02) + order.append("slow") + + async def fast(e): + order.append("fast") + + bus.subscribe(OrderCreated, slow) + bus.subscribe(OrderCreated, fast) + await bus.publish(OrderCreated(order_id=1)) + + # "fast" should complete before "slow" because they run concurrently. + assert order == ["fast", "slow"] + # ── MenuRegistry Advanced ─────────────────────────────────────────── diff --git a/modules/dashboard/tests/__init__.py b/modules/dashboard/tests/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/modules/dashboard/tests/test_dashboard.py b/modules/dashboard/tests/test_dashboard.py new file mode 100644 index 00000000..140ae8c7 --- /dev/null +++ b/modules/dashboard/tests/test_dashboard.py @@ -0,0 +1,217 @@ +"""Tests for the Dashboard module: event handlers, stats endpoint, and +end-to-end event-bus wiring between Products and Dashboard.""" + +from __future__ import annotations + +import httpx +import pytest +from simple_module_core.events import EventBus +from sm_dashboard.handlers import ( + get_product_event_counts, + on_product_created, + on_product_deleted, + on_product_updated, + reset_product_event_counts, +) +from sm_dashboard.module import DashboardModule +from sm_products.contracts.events import ProductCreated, ProductDeleted, ProductUpdated + + +@pytest.fixture(autouse=True) +def _reset_counts(): + """Ensure every test starts with zeroed product event counters.""" + reset_product_event_counts() + yield + reset_product_event_counts() + + +# ── Handler unit tests ─────────────────────────────────────────────── + + +class TestDashboardHandlers: + async def test_on_product_created_increments_counter(self): + await on_product_created(ProductCreated(product_id=1, name="Widget")) + counts = get_product_event_counts() + assert counts["created"] == 1 + assert counts["updated"] == 0 + assert counts["deleted"] == 0 + + async def test_on_product_updated_increments_counter(self): + await on_product_updated(ProductUpdated(product_id=1, name="Widget")) + counts = get_product_event_counts() + assert counts["updated"] == 1 + assert counts["created"] == 0 + + async def test_on_product_deleted_increments_counter(self): + await on_product_deleted(ProductDeleted(product_id=1)) + counts = get_product_event_counts() + assert counts["deleted"] == 1 + + async def test_multiple_events_accumulate(self): + await on_product_created(ProductCreated(product_id=1, name="A")) + await on_product_created(ProductCreated(product_id=2, name="B")) + await on_product_updated(ProductUpdated(product_id=1, name="A2")) + counts = get_product_event_counts() + assert counts["created"] == 2 + assert counts["updated"] == 1 + assert counts["deleted"] == 0 + + async def test_get_product_event_counts_returns_snapshot(self): + """Returned dict should be a copy, not the internal store.""" + counts = get_product_event_counts() + counts["created"] = 999 + # Mutating returned dict should not affect internal state + assert get_product_event_counts()["created"] == 0 + + async def test_reset_clears_all_counts(self): + await on_product_created(ProductCreated(product_id=1, name="X")) + await on_product_deleted(ProductDeleted(product_id=1)) + reset_product_event_counts() + counts = get_product_event_counts() + assert counts == {"created": 0, "updated": 0, "deleted": 0} + + +# ── Module registration tests ──────────────────────────────────────── + + +class TestDashboardModuleRegistration: + async def test_module_meta(self): + mod = DashboardModule() + assert mod.meta.name == "Dashboard" + assert mod.meta.route_prefix == "/api/dashboard" + assert "Products" in mod.meta.depends_on + + async def test_register_event_handlers_subscribes_to_all_product_events(self): + """DashboardModule should wire handlers for all three product events.""" + bus = EventBus() + mod = DashboardModule() + mod.register_event_handlers(bus) + + # Publishing each event should update the counters via the subscribed handler. + await bus.publish(ProductCreated(product_id=1, name="Widget")) + await bus.publish(ProductUpdated(product_id=1, name="Widget v2")) + await bus.publish(ProductDeleted(product_id=1)) + + counts = get_product_event_counts() + assert counts == {"created": 1, "updated": 1, "deleted": 1} + + +# ── Stats API endpoint ────────────────────────────────────────────── + + +class TestDashboardStatsEndpoint: + async def test_stats_returns_zero_counts_initially( + self, authenticated_client: httpx.AsyncClient + ): + resp = await authenticated_client.get("/api/dashboard/stats") + assert resp.status_code == 200 + body = resp.json() + assert body == {"product_events": {"created": 0, "updated": 0, "deleted": 0}} + + async def test_stats_reflects_handler_activity( + self, authenticated_client: httpx.AsyncClient + ): + await on_product_created(ProductCreated(product_id=1, name="X")) + await on_product_updated(ProductUpdated(product_id=1, name="X")) + + resp = await authenticated_client.get("/api/dashboard/stats") + assert resp.status_code == 200 + body = resp.json() + assert body["product_events"]["created"] == 1 + assert body["product_events"]["updated"] == 1 + assert body["product_events"]["deleted"] == 0 + + async def test_stats_requires_authentication(self, client: httpx.AsyncClient): + """Unauthenticated requests should be redirected by AuthMiddleware.""" + resp = await client.get("/api/dashboard/stats", follow_redirects=False) + assert resp.status_code in (302, 401, 403) + + +# ── End-to-end: Product API → EventBus → Dashboard handler ────────── + + +class TestProductEventIntegration: + """Prove the modules actually communicate through the event bus.""" + + async def test_create_product_increments_dashboard_counter( + self, authenticated_client: httpx.AsyncClient + ): + resp = await authenticated_client.post( + "/api/products/", + json={"name": "EventTestWidget", "price": "12.34"}, + ) + assert resp.status_code == 201 + + stats = await authenticated_client.get("/api/dashboard/stats") + assert stats.json()["product_events"]["created"] == 1 + + async def test_update_product_increments_dashboard_counter( + self, authenticated_client: httpx.AsyncClient + ): + create = await authenticated_client.post( + "/api/products/", + json={"name": "Original", "price": "1.00"}, + ) + product_id = create.json()["id"] + reset_product_event_counts() # isolate the update count + + resp = await authenticated_client.put( + f"/api/products/{product_id}", + json={"name": "Updated"}, + ) + assert resp.status_code == 200 + + stats = await authenticated_client.get("/api/dashboard/stats") + assert stats.json()["product_events"]["updated"] == 1 + + async def test_delete_product_increments_dashboard_counter( + self, authenticated_client: httpx.AsyncClient + ): + create = await authenticated_client.post( + "/api/products/", + json={"name": "Doomed", "price": "1.00"}, + ) + product_id = create.json()["id"] + reset_product_event_counts() + + resp = await authenticated_client.delete(f"/api/products/{product_id}") + assert resp.status_code == 204 + + stats = await authenticated_client.get("/api/dashboard/stats") + assert stats.json()["product_events"]["deleted"] == 1 + + async def test_failed_update_does_not_emit_event( + self, authenticated_client: httpx.AsyncClient + ): + """404s should not publish ProductUpdated — handler logic must be after the lookup.""" + resp = await authenticated_client.put( + "/api/products/999999", + json={"name": "ghost"}, + ) + assert resp.status_code == 404 + + stats = await authenticated_client.get("/api/dashboard/stats") + assert stats.json()["product_events"]["updated"] == 0 + + async def test_failed_delete_does_not_emit_event( + self, authenticated_client: httpx.AsyncClient + ): + resp = await authenticated_client.delete("/api/products/999999") + assert resp.status_code == 404 + + stats = await authenticated_client.get("/api/dashboard/stats") + assert stats.json()["product_events"]["deleted"] == 0 + + async def test_full_lifecycle_counters(self, authenticated_client: httpx.AsyncClient): + """Create → update → delete should all increment their respective counters.""" + create = await authenticated_client.post( + "/api/products/", + json={"name": "Lifecycle", "price": "1.00"}, + ) + pid = create.json()["id"] + await authenticated_client.put(f"/api/products/{pid}", json={"name": "L2"}) + await authenticated_client.delete(f"/api/products/{pid}") + + stats = await authenticated_client.get("/api/dashboard/stats") + counts = stats.json()["product_events"] + assert counts == {"created": 1, "updated": 1, "deleted": 1} diff --git a/pyproject.toml b/pyproject.toml index 5a8d3060..744c9427 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -47,6 +47,6 @@ extra-paths = [ [tool.pytest.ini_options] asyncio_mode = "auto" -testpaths = ["framework/core/tests", "framework/db/tests", "framework/hosting/tests", "modules/auth/tests", "modules/products/tests"] +testpaths = ["framework/core/tests", "framework/db/tests", "framework/hosting/tests", "modules/auth/tests", "modules/dashboard/tests", "modules/products/tests"] markers = ["e2e: end-to-end tests requiring live services (Keycloak, browser)"] addopts = "-m 'not e2e'" From d5e4972e54bb2becdb1f918d7d2622ebd8fdc0e8 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 14 Apr 2026 13:34:41 +0000 Subject: [PATCH 3/4] refactor: apply simplify-review clear wins MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Use event class objects directly as pyee keys; drop _event_key() (pyee accepts any hashable, saves an f-string on every subscribe/publish) - Hoist register_event_handlers imports to dashboard/module.py top (no cycle risk, no FastAPI side effects) - Drop trivial docstrings on Product event dataclasses (module docstring + class names already convey the meaning) - Drop narrating test comments and redundant class docstrings Architectural findings (await publish → outbox, service-layer publish, SQL-derived /stats) deferred per prior instruction. https://claude.ai/code/session_01KLjCnJoNbCbk8is6pVZcom --- framework/core/simple_module_core/events.py | 13 +++---------- modules/dashboard/dashboard/module.py | 15 +++------------ modules/dashboard/tests/test_dashboard.py | 8 +------- modules/products/products/contracts/events.py | 6 ------ 4 files changed, 7 insertions(+), 35 deletions(-) diff --git a/framework/core/simple_module_core/events.py b/framework/core/simple_module_core/events.py index eb60f1f4..30c1924e 100644 --- a/framework/core/simple_module_core/events.py +++ b/framework/core/simple_module_core/events.py @@ -44,17 +44,11 @@ def __init__(self) -> None: @staticmethod def _on_emitter_error(error: Exception) -> None: - """Handle errors from fire-and-forget dispatch.""" logger.error("EventBus background handler error: %s", error, exc_info=error) - @staticmethod - def _event_key(event_type: type[Event]) -> str: - """Derive a unique string key from an event class for pyee registration.""" - return f"{event_type.__module__}.{event_type.__qualname__}" - def subscribe(self, event_type: type[Event], handler: EventHandler) -> None: """Register a handler for an event type.""" - self._emitter.on(self._event_key(event_type), handler) + self._emitter.on(event_type, handler) logger.debug( "Subscribed %s to %s", getattr(handler, "__qualname__", repr(handler)), @@ -67,8 +61,7 @@ async def publish(self, event: Event) -> None: All handlers run concurrently via ``asyncio.gather``. Individual handler failures are logged but do not propagate. """ - key = self._event_key(type(event)) - handlers = self._emitter.listeners(key) + handlers = self._emitter.listeners(type(event)) if not handlers: return results = await asyncio.gather( @@ -91,4 +84,4 @@ def publish_nowait(self, event: Event) -> None: Uses pyee's ``AsyncIOEventEmitter.emit`` which schedules async handlers as tasks on the running loop. """ - self._emitter.emit(self._event_key(type(event)), event) + self._emitter.emit(type(event), event) diff --git a/modules/dashboard/dashboard/module.py b/modules/dashboard/dashboard/module.py index 7f348e61..54facc62 100644 --- a/modules/dashboard/dashboard/module.py +++ b/modules/dashboard/dashboard/module.py @@ -3,10 +3,13 @@ from __future__ import annotations from fastapi import APIRouter +from products.contracts.events import ProductCreated, ProductDeleted, ProductUpdated from simple_module_core.events import EventBus from simple_module_core.menu import MenuItem, MenuRegistry, MenuSection from simple_module_core.module import ModuleBase, ModuleMeta +from dashboard.handlers import on_product_created, on_product_deleted, on_product_updated + class DashboardModule(ModuleBase): meta = ModuleMeta( @@ -35,18 +38,6 @@ def register_menu_items(self, registry: MenuRegistry) -> None: ) def register_event_handlers(self, bus: EventBus) -> None: - from products.contracts.events import ( - ProductCreated, - ProductDeleted, - ProductUpdated, - ) - - from dashboard.handlers import ( - on_product_created, - on_product_deleted, - on_product_updated, - ) - bus.subscribe(ProductCreated, on_product_created) bus.subscribe(ProductUpdated, on_product_updated) bus.subscribe(ProductDeleted, on_product_deleted) diff --git a/modules/dashboard/tests/test_dashboard.py b/modules/dashboard/tests/test_dashboard.py index de928432..706ac6c7 100644 --- a/modules/dashboard/tests/test_dashboard.py +++ b/modules/dashboard/tests/test_dashboard.py @@ -60,7 +60,6 @@ async def test_get_product_event_counts_returns_snapshot(self): """Returned dict should be a copy, not the internal store.""" counts = get_product_event_counts() counts["created"] = 999 - # Mutating returned dict should not affect internal state assert get_product_event_counts()["created"] == 0 async def test_reset_clears_all_counts(self): @@ -82,12 +81,10 @@ async def test_module_meta(self): assert "Products" in mod.meta.depends_on async def test_register_event_handlers_subscribes_to_all_product_events(self): - """DashboardModule should wire handlers for all three product events.""" bus = EventBus() mod = DashboardModule() mod.register_event_handlers(bus) - # Publishing each event should update the counters via the subscribed handler. await bus.publish(ProductCreated(product_id=1, name="Widget")) await bus.publish(ProductUpdated(product_id=1, name="Widget v2")) await bus.publish(ProductDeleted(product_id=1)) @@ -131,8 +128,6 @@ async def test_stats_requires_authentication(self, client: httpx.AsyncClient): class TestProductEventIntegration: - """Prove the modules actually communicate through the event bus.""" - async def test_create_product_increments_dashboard_counter( self, authenticated_client: httpx.AsyncClient ): @@ -153,7 +148,7 @@ async def test_update_product_increments_dashboard_counter( json={"name": "Original", "price": "1.00"}, ) product_id = create.json()["id"] - reset_product_event_counts() # isolate the update count + reset_product_event_counts() resp = await authenticated_client.put( f"/api/products/{product_id}", @@ -203,7 +198,6 @@ async def test_failed_delete_does_not_emit_event( assert stats.json()["product_events"]["deleted"] == 0 async def test_full_lifecycle_counters(self, authenticated_client: httpx.AsyncClient): - """Create → update → delete should all increment their respective counters.""" create = await authenticated_client.post( "/api/products/", json={"name": "Lifecycle", "price": "1.00"}, diff --git a/modules/products/products/contracts/events.py b/modules/products/products/contracts/events.py index d73aefc0..17b63bf9 100644 --- a/modules/products/products/contracts/events.py +++ b/modules/products/products/contracts/events.py @@ -9,22 +9,16 @@ @dataclass class ProductCreated(Event): - """Published after a new product is persisted.""" - product_id: int name: str @dataclass class ProductUpdated(Event): - """Published after an existing product is modified.""" - product_id: int name: str @dataclass class ProductDeleted(Event): - """Published after a product is removed.""" - product_id: int From 10ad1974a970d3121ac6e096da1e8dfa3b574776 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 14 Apr 2026 15:41:45 +0000 Subject: [PATCH 4/4] fix: restore _event_key str helper for pyee type compatibility MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Reinstate ``_event_key()`` so pyee calls receive str, matching its stubs (ty caught the stub mismatch on CI). The class-as-key simplification was a false-positive — runtime accepts it but the typed API does not. - Apply ruff format to test_dashboard.py (CI check). https://claude.ai/code/session_01KLjCnJoNbCbk8is6pVZcom --- framework/core/simple_module_core/events.py | 10 +++++++--- modules/dashboard/tests/test_dashboard.py | 12 +++--------- 2 files changed, 10 insertions(+), 12 deletions(-) diff --git a/framework/core/simple_module_core/events.py b/framework/core/simple_module_core/events.py index 30c1924e..fc5751cd 100644 --- a/framework/core/simple_module_core/events.py +++ b/framework/core/simple_module_core/events.py @@ -46,9 +46,13 @@ def __init__(self) -> None: def _on_emitter_error(error: Exception) -> None: logger.error("EventBus background handler error: %s", error, exc_info=error) + @staticmethod + def _event_key(event_type: type[Event]) -> str: + return f"{event_type.__module__}.{event_type.__qualname__}" + def subscribe(self, event_type: type[Event], handler: EventHandler) -> None: """Register a handler for an event type.""" - self._emitter.on(event_type, handler) + self._emitter.on(self._event_key(event_type), handler) logger.debug( "Subscribed %s to %s", getattr(handler, "__qualname__", repr(handler)), @@ -61,7 +65,7 @@ async def publish(self, event: Event) -> None: All handlers run concurrently via ``asyncio.gather``. Individual handler failures are logged but do not propagate. """ - handlers = self._emitter.listeners(type(event)) + handlers = self._emitter.listeners(self._event_key(type(event))) if not handlers: return results = await asyncio.gather( @@ -84,4 +88,4 @@ def publish_nowait(self, event: Event) -> None: Uses pyee's ``AsyncIOEventEmitter.emit`` which schedules async handlers as tasks on the running loop. """ - self._emitter.emit(type(event), event) + self._emitter.emit(self._event_key(type(event)), event) diff --git a/modules/dashboard/tests/test_dashboard.py b/modules/dashboard/tests/test_dashboard.py index 706ac6c7..45e05f0f 100644 --- a/modules/dashboard/tests/test_dashboard.py +++ b/modules/dashboard/tests/test_dashboard.py @@ -105,9 +105,7 @@ async def test_stats_returns_zero_counts_initially( body = resp.json() assert body == {"product_events": {"created": 0, "updated": 0, "deleted": 0}} - async def test_stats_reflects_handler_activity( - self, authenticated_client: httpx.AsyncClient - ): + async def test_stats_reflects_handler_activity(self, authenticated_client: httpx.AsyncClient): await on_product_created(ProductCreated(product_id=1, name="X")) await on_product_updated(ProductUpdated(product_id=1, name="X")) @@ -175,9 +173,7 @@ async def test_delete_product_increments_dashboard_counter( stats = await authenticated_client.get("/api/dashboard/stats") assert stats.json()["product_events"]["deleted"] == 1 - async def test_failed_update_does_not_emit_event( - self, authenticated_client: httpx.AsyncClient - ): + async def test_failed_update_does_not_emit_event(self, authenticated_client: httpx.AsyncClient): """404s should not publish ProductUpdated — handler logic must be after the lookup.""" resp = await authenticated_client.put( "/api/products/999999", @@ -188,9 +184,7 @@ async def test_failed_update_does_not_emit_event( stats = await authenticated_client.get("/api/dashboard/stats") assert stats.json()["product_events"]["updated"] == 0 - async def test_failed_delete_does_not_emit_event( - self, authenticated_client: httpx.AsyncClient - ): + async def test_failed_delete_does_not_emit_event(self, authenticated_client: httpx.AsyncClient): resp = await authenticated_client.delete("/api/products/999999") assert resp.status_code == 404