"""ASV benchmarks for Event System Domain Event Taxonomy. Measures: - DomainEvent construction throughput (minimal, with plan_id, with details) - ReactiveEventBus emit with audit_log overhead - ReactiveEventBus subscribe registration - Fan-out: emit to N subscribers with audit persistence - Audit log retrieval cost - Log correlation (filtering events by correlation_id) """ from __future__ import annotations import importlib import logging import sys from pathlib import Path from typing import ClassVar _SRC = str(Path(__file__).resolve().parents[1] / "src") if _SRC not in sys.path: sys.path.insert(0, _SRC) import cleveragents # noqa: E402 importlib.reload(cleveragents) from cleveragents.infrastructure.events import ( # noqa: E402 DomainEvent, EventType, LoggingEventBus, ReactiveEventBus, ) # --------------------------------------------------------------------------- # DomainEvent construction # --------------------------------------------------------------------------- class TaxonomyDomainEventSuite: """Benchmark DomainEvent construction and serialisation.""" timeout = 60 def setup(self) -> None: self.event_type = EventType.PLAN_CREATED self.plan_id = "01ARZ3NDEKTSV4RRFFQ69G5FAV" def time_create_minimal(self) -> None: DomainEvent(event_type=self.event_type) def time_create_with_all_fields(self) -> None: DomainEvent( event_type=self.event_type, plan_id=self.plan_id, root_plan_id=self.plan_id, session_id=self.plan_id, actor_name="local/strategist", project_name="my-project", details={"phase": "strategize", "actor": "strategists/planner"}, ) def time_json_roundtrip(self) -> None: event = DomainEvent( event_type=self.event_type, plan_id=self.plan_id, details={"phase": "execute"}, ) DomainEvent.model_validate_json(event.model_dump_json()) # --------------------------------------------------------------------------- # ReactiveEventBus emit with audit_log # --------------------------------------------------------------------------- class TaxonomyReactiveEmitSuite: """Benchmark ReactiveEventBus.emit() with audit_log persistence.""" timeout = 60 def setup(self) -> None: self.bus = ReactiveEventBus() self.event = DomainEvent(event_type=EventType.PLAN_CREATED) self.decision_event = DomainEvent(event_type=EventType.DECISION_CREATED) self._received: list[DomainEvent] = [] self.bus.subscribe(EventType.PLAN_CREATED, self._received.append) def time_emit_no_subscribers(self) -> None: self._received.clear() self.bus.clear_audit_log() self.bus.emit(self.decision_event) def time_emit_with_one_subscriber(self) -> None: self._received.clear() self.bus.clear_audit_log() self.bus.emit(self.event) def time_emit_100_events(self) -> None: self._received.clear() self.bus.clear_audit_log() for _ in range(100): self.bus.emit(self.event) # --------------------------------------------------------------------------- # Audit log retrieval # --------------------------------------------------------------------------- class TaxonomyAuditLogSuite: """Benchmark audit_log property access (defensive copy).""" timeout = 60 params: ClassVar[list[int]] = [100, 1000] param_names: ClassVar[list[str]] = ["num_events"] def setup(self, num_events: int) -> None: self.bus = ReactiveEventBus() for _ in range(num_events): self.bus.emit(DomainEvent(event_type=EventType.PLAN_CREATED)) def time_audit_log_snapshot(self, num_events: int) -> None: _ = self.bus.audit_log # --------------------------------------------------------------------------- # Fan-out with audit # --------------------------------------------------------------------------- class TaxonomyFanOutSuite: """Benchmark event fan-out to multiple subscribers with audit persistence.""" timeout = 60 params: ClassVar[list[int]] = [1, 5, 10, 50] param_names: ClassVar[list[str]] = ["num_subscribers"] def setup(self, num_subscribers: int) -> None: self.bus = ReactiveEventBus() for _ in range(num_subscribers): self.bus.subscribe(EventType.PLAN_CREATED, lambda e: None) self.event = DomainEvent(event_type=EventType.PLAN_CREATED) def time_emit(self, num_subscribers: int) -> None: self.bus.clear_audit_log() self.bus.emit(self.event) # --------------------------------------------------------------------------- # LoggingEventBus emit # --------------------------------------------------------------------------- class TaxonomyLoggingEmitSuite: """Benchmark LoggingEventBus.emit() overhead.""" timeout = 60 def setup(self) -> None: logging.disable(logging.CRITICAL) self.bus = LoggingEventBus() self.event = DomainEvent(event_type=EventType.DECISION_CREATED) def teardown(self) -> None: logging.disable(logging.NOTSET) def time_emit_single(self) -> None: self.bus.emit(self.event) def time_emit_100(self) -> None: for _ in range(100): self.bus.emit(self.event)