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/simple_module_core/events.py b/framework/core/simple_module_core/events.py index 7ca945ca..fc5751cd 100644 --- a/framework/core/simple_module_core/events.py +++ b/framework/core/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,30 @@ 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: + 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._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 +60,12 @@ 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. + """ + handlers = self._emitter.listeners(self._event_key(type(event))) if not handlers: return results = await asyncio.gather( @@ -66,6 +83,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/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/dashboard/endpoints/api.py b/modules/dashboard/dashboard/endpoints/api.py new file mode 100644 index 00000000..e387651b --- /dev/null +++ b/modules/dashboard/dashboard/endpoints/api.py @@ -0,0 +1,17 @@ +"""REST API endpoints for the Dashboard module.""" + +from __future__ import annotations + +from fastapi import APIRouter + +from 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/dashboard/handlers.py b/modules/dashboard/dashboard/handlers.py new file mode 100644 index 00000000..c4611fbd --- /dev/null +++ b/modules/dashboard/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 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/dashboard/module.py b/modules/dashboard/dashboard/module.py index 1554dc47..54facc62 100644 --- a/modules/dashboard/dashboard/module.py +++ b/modules/dashboard/dashboard/module.py @@ -3,19 +3,27 @@ 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( name="Dashboard", + route_prefix="/api/dashboard", view_prefix="", + depends_on=["Products"], ) def register_routes(self, api_router: APIRouter, view_router: APIRouter) -> None: + from dashboard.endpoints.api import router as api from 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 +36,8 @@ def register_menu_items(self, registry: MenuRegistry) -> None: section=MenuSection.SIDEBAR, ) ) + + def register_event_handlers(self, bus: EventBus) -> None: + bus.subscribe(ProductCreated, on_product_created) + bus.subscribe(ProductUpdated, on_product_updated) + bus.subscribe(ProductDeleted, on_product_deleted) diff --git a/modules/dashboard/pyproject.toml b/modules/dashboard/pyproject.toml index d4500a32..37047cc3 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", + "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 } +products = { workspace = true } diff --git a/modules/dashboard/tests/test_dashboard.py b/modules/dashboard/tests/test_dashboard.py new file mode 100644 index 00000000..45e05f0f --- /dev/null +++ b/modules/dashboard/tests/test_dashboard.py @@ -0,0 +1,205 @@ +"""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 dashboard.handlers import ( + get_product_event_counts, + on_product_created, + on_product_deleted, + on_product_updated, + reset_product_event_counts, +) +from dashboard.module import DashboardModule +from products.contracts.events import ProductCreated, ProductDeleted, ProductUpdated +from simple_module_core.events import EventBus + + +@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 + 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): + bus = EventBus() + mod = DashboardModule() + mod.register_event_handlers(bus) + + 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: + 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() + + 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 = 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/modules/products/products/contracts/__init__.py b/modules/products/products/contracts/__init__.py index e0e3a461..9125870d 100644 --- a/modules/products/products/contracts/__init__.py +++ b/modules/products/products/contracts/__init__.py @@ -1,5 +1,10 @@ """Products contracts — public interface for other modules.""" +from products.contracts.events import ( + ProductCreated, + ProductDeleted, + ProductUpdated, +) from products.contracts.schemas import ( ProductCreate, ProductOut, @@ -7,4 +12,12 @@ ) from products.contracts.service import IProductService -__all__ = ["ProductCreate", "ProductOut", "ProductUpdate", "IProductService"] +__all__ = [ + "ProductCreate", + "ProductCreated", + "ProductDeleted", + "ProductOut", + "ProductUpdate", + "ProductUpdated", + "IProductService", +] diff --git a/modules/products/products/contracts/events.py b/modules/products/products/contracts/events.py new file mode 100644 index 00000000..17b63bf9 --- /dev/null +++ b/modules/products/products/contracts/events.py @@ -0,0 +1,24 @@ +"""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): + product_id: int + name: str + + +@dataclass +class ProductUpdated(Event): + product_id: int + name: str + + +@dataclass +class ProductDeleted(Event): + product_id: int diff --git a/modules/products/products/deps.py b/modules/products/products/deps.py index ae77301ae9..d3ad502b 100644 --- a/modules/products/products/deps.py +++ b/modules/products/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/products/endpoints/api.py b/modules/products/products/endpoints/api.py index 9f970ed9..ae527299 100644 --- a/modules/products/products/endpoints/api.py +++ b/modules/products/products/endpoints/api.py @@ -3,10 +3,12 @@ from __future__ import annotations from fastapi import APIRouter, Depends, HTTPException +from simple_module_core.events import EventBus from simple_module_hosting.permissions import RequiresPermission +from products.contracts.events import ProductCreated, ProductDeleted, ProductUpdated from products.contracts.schemas import ProductCreate, ProductOut, ProductUpdate -from products.deps import get_product_service +from products.deps import get_event_bus, get_product_service from products.service import ProductService router = APIRouter() @@ -40,8 +42,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( @@ -53,10 +58,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 @@ -68,7 +75,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)) diff --git a/pyproject.toml b/pyproject.toml index 3ece15d1..5ea726db 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", "tests/integration"] +testpaths = ["framework/core/tests", "framework/db/tests", "framework/hosting/tests", "modules/auth/tests", "modules/dashboard/tests", "modules/products/tests", "tests/integration"] markers = ["e2e: end-to-end tests requiring live services (Keycloak, browser)"] addopts = "-m 'not e2e'"