diff --git a/docs/reference/event_bus.md b/docs/reference/event_bus.md index 283a85f79..9e43934af 100644 --- a/docs/reference/event_bus.md +++ b/docs/reference/event_bus.md @@ -10,7 +10,7 @@ one another. ## Overview -The system consists of four cooperating components: +The system consists of five cooperating components: | Component | Role | |-----------|------| @@ -87,8 +87,8 @@ after construction raises a `ValidationError`. ## EventBus Protocol `cleveragents.infrastructure.events.EventBus` is a `@runtime_checkable` -structural protocol. Any object with `emit()` and `subscribe()` methods -satisfies it: +structural protocol. Any object with `emit()`, `subscribe()`, and `unsubscribe()` +methods satisfies it: ```python from cleveragents.infrastructure.events import EventBus, ReactiveEventBus @@ -104,6 +104,20 @@ class EventBus(Protocol): def emit(self, event: DomainEvent) -> None: ... def subscribe(self, event_type: EventType, handler: Callable[[DomainEvent], None]) -> None: ... + def unsubscribe(self, event_type: EventType, + handler: Callable[[DomainEvent], None]) -> bool: ... +``` + +**Memory management:** Handlers registered via `subscribe()` are stored as strong +references inside the bus. Failing to call `unsubscribe()` after subscribing means +the handler is retained in memory for the lifetime of the bus (which, in production, +is the entire application lifetime). Always call `unsubscribe()` when a handler is +no longer needed: + +```python +bus.subscribe(EventType.PLAN_CREATED, on_plan_created) +# ... later ... +bus.unsubscribe(EventType.PLAN_CREATED, on_plan_created) # frees the reference ``` --- @@ -131,11 +145,16 @@ bus.stream.pipe( # Emit an event bus.emit(DomainEvent(event_type=EventType.PLAN_CREATED, plan_id="01ABC...")) + +# Unsubscribe when the handler is no longer needed +bus.unsubscribe(EventType.PLAN_CREATED, lambda e: print("...")) ``` - `emit(event)` pushes to the RxPY Subject **and** calls all type-specific handlers synchronously. - `subscribe(event_type, handler)` registers a callback for one event type. +- `unsubscribe(event_type, handler) -> bool` removes a callback; returns ``True`` + if the subscription was found and removed, ``False`` otherwise. - `stream` exposes a read-only `rx.Observable` view for RxPY operators. - `clear_audit_log()` clears retained in-memory events when needed. @@ -155,6 +174,9 @@ bus.emit(DomainEvent(event_type=EventType.DECISION_CREATED, plan_id="01ABC...")) # → structured log: {"event": "domain_event", "event_type": "decision.created", ...} ``` +The logging bus also supports `unsubscribe()` for the same memory-management +reasons as ``ReactiveEventBus``. + --- ## Dependency Injection @@ -218,4 +240,24 @@ bus.subscribe(EventType.PLAN_CREATED, on_plan_created) # All subsequent plan.create operations will call on_plan_created svc = container.plan_lifecycle_service() + +# Later — clean up to prevent memory leaks: +bus.unsubscribe(EventType.PLAN_CREATED, on_plan_created) +``` + +### Example: Lifecycle-aware subscribers (AuditEventSubscriber) + +The built-in :class:`~cleveragents.application.services.audit_event_subscriber.AuditEventSubscriber` +subscribes to all security-relevant event types and supports clean teardown +via :meth:`~cleveragents.application.services.audit_event_subscriber.AuditEventSubscriber.stop`: + +```python +from cleveragents.application.services.audit_service import AuditService +from cleveragents.application.services.audit_event_subscriber import ( + AuditEventSubscriber, +) + +subscriber = AuditEventSubscriber(audit_service, bus) +# ... application runs ... +subscriber.stop() # unsubscribes from all event types ``` diff --git a/features/event_bus.feature b/features/event_bus.feature index 2875b4f58..67a684511 100644 --- a/features/event_bus.feature +++ b/features/event_bus.feature @@ -89,6 +89,42 @@ Feature: EventBus protocol and domain event emission Given a ReactiveEventBus Then subscribing with a non-EventType should raise TypeError + # --------------------------------------------------------------------------- + # ReactiveEventBus unsubscribe (issue #10356) + # --------------------------------------------------------------------------- + + Scenario: ReactiveEventBus.unsubscribe() removes handler from list + Given a ReactiveEventBus + When I subscribe to "plan.created" events + And I unsubscribe the handler for "plan.created" + Then unsubcribing should return True + + Scenario: ReactiveEventBus.unsubscribe() returns False when handler absent + Given a ReactiveEventBus + Then unsubcribing a non-subscribed handler should return False + + Scenario: ReactiveEventBus.unsubscribe() releases handler reference (no event received after) + Given a ReactiveEventBus + When I subscribe to "plan.created" events + And I unsubscribe the handler for "plan.created" + And I emit a "plan.created" DomainEvent + Then the handler should have received 0 event + + Scenario: ReactiveEventBus.unsubscribe() is selective (other handlers preserved) + Given a ReactiveEventBus + When I subscribe two handlers to "plan.created" events + And I unsubscribe the first handler for "plan.created" + And I emit a "plan.created" DomainEvent + Then each remaining handler should have received 1 event + + Scenario: ReactiveEventBus.unsubscribe() rejects invalid event_type + Given a ReactiveEventBus + Then unsubcribing with a non-EventType should raise TypeError + + Scenario: ReactiveEventBus.unsubscribe() rejects non-callable handler + Given a ReactiveEventBus + Then unsubcribing with a non-callable handler should raise TypeError + # --------------------------------------------------------------------------- # LoggingEventBus # --------------------------------------------------------------------------- @@ -115,6 +151,39 @@ Feature: EventBus protocol and domain event emission Given a ReactiveEventBus Then subscribing with a non-callable handler should raise TypeError + # --------------------------------------------------------------------------- + # LoggingEventBus unsubscribe (issue #10356) + # --------------------------------------------------------------------------- + + Scenario: LoggingEventBus.unsubscribe() removes handler from list + Given a LoggingEventBus + When I subscribe to "plan.created" events + And I unsubscribe the handler for "plan.created" + Then unsubcribing should return True + + Scenario: LoggingEventBus.unsubscribe() returns False when handler absent + Given a LoggingEventBus + Then unsubcribing a non-subscribed handler should return False + + Scenario: LoggingEventBus.unsubscribe() releases handler reference (no event received after) + Given a LoggingEventBus + When I subscribe to "decision.created" events + And I unsubscribe the handler for "decision.created" + And I emit a "decision.created" DomainEvent + Then the handler should have received 0 event + + Scenario: LoggingEventBus.unsubscribe() rejects invalid event_type + Given a LoggingEventBus + Then unsubcribing with a non-EventType should raise TypeError + + Scenario: LoggingEventBus.unsubscribe() rejects non-callable handler + Given a LoggingEventBus + Then unsubcribing with a non-callable handler should raise TypeError + + # --------------------------------------------------------------------------- + # EventBus Protocol stubs + # --------------------------------------------------------------------------- + Scenario: EventBus Protocol emit and subscribe stubs are executable Then the EventBus Protocol stubs should be callable diff --git a/features/steps/event_bus_steps.py b/features/steps/event_bus_steps.py index 35c1abd93..d6918dc91 100644 --- a/features/steps/event_bus_steps.py +++ b/features/steps/event_bus_steps.py @@ -1,7 +1,8 @@ """Step definitions for event_bus.feature. Tests EventType, DomainEvent, EventBus protocol, ReactiveEventBus, -LoggingEventBus, service event emission, and DI container registration. +LoggingEventBus, service event emission, DI container registration, +and unsubscribe scenarios (issue #10356). """ from __future__ import annotations @@ -144,16 +145,24 @@ def step_domain_event_immutable(ctx: Context) -> None: def step_reactive_satisfies_protocol(ctx: Context) -> None: bus = ReactiveEventBus() assert isinstance(bus, EventBus), "ReactiveEventBus does not satisfy EventBus" + # Verify the unsubscribe method exists on the protocol. + assert hasattr( + bus, "unsubscribe" + ), "ReactiveEventBus missing required unsubscribe method" @then("LoggingEventBus should satisfy the EventBus protocol") def step_logging_satisfies_protocol(ctx: Context) -> None: bus = LoggingEventBus() assert isinstance(bus, EventBus), "LoggingEventBus does not satisfy EventBus" + # Verify the unsubscribe method exists on the protocol. + assert hasattr( + bus, "unsubscribe" + ), "LoggingEventBus missing required unsubscribe method" # --------------------------------------------------------------------------- -# ReactiveEventBus given +# ReactiveEventBus / LoggingEventBus given # --------------------------------------------------------------------------- @@ -178,6 +187,8 @@ def step_given_logging_bus(ctx: Context) -> None: def step_subscribe_single(ctx: Context, et: str) -> None: collector = _EventCollector() ctx.collectors.append(collector) + # Store reference for later unsubscribe. + ctx.subscribed_handler = collector ctx.bus.subscribe(EventType(et), collector) @@ -186,7 +197,10 @@ def step_subscribe_two(ctx: Context, et: str) -> None: for _ in range(2): collector = _EventCollector() ctx.collectors.append(collector) - ctx.bus.subscribe(EventType(et), collector) + # Store references in order for selective unsubscribe. + ctx.subscribed_handlers: list[_EventCollector] = [] + ctx.bus.subscribe(EventType(et), ctx.collectors[0]) + ctx.bus.subscribe(EventType(et), ctx.collectors[1]) @when('I emit a "{et}" DomainEvent') @@ -194,6 +208,23 @@ def step_emit_event(ctx: Context, et: str) -> None: ctx.bus.emit(_make_event(et)) +# --------------------------------------------------------------------------- +# Unsubscribe steps +# --------------------------------------------------------------------------- + + +@when("I unsubscribe the handler for \"{et}\"") +def step_unsubscribe_handler(ctx: Context, et: str) -> None: + handler = getattr(ctx, "subscribed_handler", None) + if handler is not None: + ctx.unsubscribe_result = ctx.bus.unsubscribe(EventType(et), handler) + elif hasattr(ctx, "subscribed_handlers") and len(ctx.subscribed_handlers) >= 2: + # Selective unsubscribe — remove the first handler. + ctx.unsubscribe_result = ctx.bus.unsubscribe( + EventType(et), ctx.subscribed_handlers[0] + ) + + # --------------------------------------------------------------------------- # ReactiveEventBus then # --------------------------------------------------------------------------- @@ -267,6 +298,60 @@ def step_protocol_stubs_callable(ctx: Context) -> None: event = _make_event("plan.created") obj.emit(event) obj.subscribe(EventType.PLAN_CREATED, lambda e: None) + obj.unsubscribe(EventType.PLAN_CREATED, lambda e: None) + + +# --------------------------------------------------------------------------- +# Unsubscribe then steps +# --------------------------------------------------------------------------- + + +@then('unsubcribing should return {should}') +def step_unsubscribe_returns(ctx: Context, should: str) -> None: + parsed = should.strip().lower() == "true" + if not hasattr(ctx, "unsubscribe_result"): + result = ctx.bus.unsubscribe( + EventType.PLAN_CREATED, # type: ignore[arg-type] + lambda e: None, # type: ignore[arg-type] + ) + else: + result = ctx.unsubscribe_result + assert result is parsed, f"Expected unsubscribe to return {parsed}, got {result}" + + +@then("unsubcribing a non-subscribed handler should return False") +def step_unsubscribe_not_subscribed(ctx: Context) -> None: + result = ctx.bus.unsubscribe( + EventType.PLAN_CREATED, # type: ignore[arg-type] + lambda e: None, # type: ignore[arg-type] + ) + assert result is False, f"Expected False, got {result}" + + +@then("unsubcribing with a non-EventType should raise TypeError") +def step_unsubscribe_bad_type(ctx: Context) -> None: + raised = False + try: + ctx.bus.unsubscribe( # type: ignore[arg-type] + "plan.created", # type: ignore[arg-type] + _EventCollector(), # type: ignore[arg-type] + ) + except TypeError: + raised = True + assert raised, "Expected TypeError for non-EventType event_type" + + +@then("unsubcribing with a non-callable handler should raise TypeError") +def step_unsubscribe_non_callable(ctx: Context) -> None: + raised = False + try: + ctx.bus.unsubscribe( # type: ignore[arg-type] + EventType.PLAN_CREATED, + "not-callable", # type: ignore[arg-type] + ) + except TypeError: + raised = True + assert raised, "Expected TypeError for non-callable handler" # --------------------------------------------------------------------------- diff --git a/src/cleveragents/a2a/events.py b/src/cleveragents/a2a/events.py index f17ae21e7..2a1d84342 100644 --- a/src/cleveragents/a2a/events.py +++ b/src/cleveragents/a2a/events.py @@ -11,6 +11,9 @@ the A2A protocol specification. :class:`EventBusBridge` subscribes to the internal ``EventBus`` and publishes translated :class:`A2aEvent` instances to an event queue. + +Based on: + - Forgejo issue #10356 (unsubscribe support in EventBus) """ from __future__ import annotations @@ -232,6 +235,15 @@ class EventBusBridge: bridge = EventBusBridge(event_bus, event_queue) bridge.start() # subscribes to bus bridge.stop() # unsubscribes + + The bridge registers a single handler covering all domain event types it + needs. :meth:`stop` calls :meth:`EventBus.unsubscribe` with that handler + so that the handler reference is released from the bus (see Forgejo issue + **#10356**). + + Args: + event_bus: The :class:`EventBus` bridge subscribes to. + event_queue: The target queue for translated A2A events. """ # Domain event types that map to SSE TaskStatusUpdateEvent @@ -260,21 +272,32 @@ class EventBusBridge: ) -> None: self._event_bus = event_bus self._event_queue = event_queue - self._subscription: Any | None = None + self._subscription: bool = False # Whether bridge is currently started def start(self) -> None: """Subscribe to the event bus and begin forwarding.""" if hasattr(self._event_bus, "subscribe"): - self._subscription = self._event_bus.subscribe(self._on_domain_event) + self._event_bus.subscribe(self._on_domain_event) + self._subscription = True logger.info("a2a.event_bridge.started") def stop(self) -> None: - """Unsubscribe from the event bus.""" - if self._subscription is not None: - if hasattr(self._subscription, "dispose"): - self._subscription.dispose() - self._subscription = None + """Unsubscribe from the event bus. + + Calls :meth:`EventBus.unsubscribe` to remove the previously + registered handler. This prevents memory leaks caused by lingering + handler references in the bus (see Forgejo issue #10356). + """ + if self._subscription: + # The bridge subscribes without event_type filtering, so we try + # all known event types to remove our handler from every one. + from cleveragents.infrastructure.events.types import EventType + + if hasattr(self._event_bus, "unsubscribe"): + for event_type in EventType: + self._event_bus.unsubscribe(event_type, self._on_domain_event) logger.info("a2a.event_bridge.stopped") + self._subscription = False def _on_domain_event(self, domain_event: Any) -> None: """Translate a domain event to an A2A event and publish.""" diff --git a/src/cleveragents/application/services/audit_event_subscriber.py b/src/cleveragents/application/services/audit_event_subscriber.py index 215cacf89..a1b106ee6 100644 --- a/src/cleveragents/application/services/audit_event_subscriber.py +++ b/src/cleveragents/application/services/audit_event_subscriber.py @@ -12,6 +12,7 @@ Sensitive data in event ``details`` is redacted via Based on: - docs/specification.md §Audit Logging (SEC7) - Forgejo issue #581 + - Forgejo issue #10356 (unsubscribe support) """ from __future__ import annotations @@ -71,6 +72,10 @@ class AuditEventSubscriber: secrets in the details dict. 3. Calls ``audit_service.record()`` with the mapped event type string. + The subscriber can be cleanly torn down via :meth:`stop`, which + unsubscribes the handler from every event type it originally + subscribed to, preventing memory leaks (see Forgejo issue #10356). + Args: audit_service: The :class:`AuditService` used for persistence. event_bus: The :class:`EventBus` to subscribe to. @@ -79,9 +84,31 @@ class AuditEventSubscriber: def __init__(self, audit_service: AuditService, event_bus: EventBus) -> None: self._audit_service = audit_service self._event_bus = event_bus - for event_type in SECURITY_EVENT_TYPES: + self._subscribed_types: frozenset[EventType] = SECURITY_EVENT_TYPES + for event_type in self._subscribed_types: event_bus.subscribe(event_type, self._handle_event) + def stop(self) -> None: + """Unsubscribe from all event types previously subscribed to. + + Calls :meth:`EventBus.unsubscribe` for the handler on every + :class:`EventType` registered during :meth:`__init__`. The method is + safe to call multiple times — unsubscription is idempotent. + + Call this when the subscriber is no longer needed (e.g. application + shutdown or test teardown) to release handler references held by the + event bus and prevent memory leaks. + """ + for event_type in self._subscribed_types: + removed = self._event_bus.unsubscribe( + event_type, self._handle_event + ) + if removed: + _logger.debug( + "audit_unsubscribed", + event_type=event_type.value, + ) + def _handle_event(self, event: DomainEvent) -> None: """Redact sensitive data and persist an audit log entry. diff --git a/src/cleveragents/infrastructure/events/logging_bus.py b/src/cleveragents/infrastructure/events/logging_bus.py index 307bf4ed3..2aebe3578 100644 --- a/src/cleveragents/infrastructure/events/logging_bus.py +++ b/src/cleveragents/infrastructure/events/logging_bus.py @@ -9,6 +9,7 @@ testing, and environments where reactive streaming is not needed. Based on: - docs/specification.md §Event-Driven Architecture - Forgejo issue #473 + - Forgejo issue #10356 (unsubscribe support) """ from __future__ import annotations @@ -105,5 +106,51 @@ class LoggingEventBus: raise TypeError("handler must be callable") self._subscriptions.setdefault(event_type, []).append(handler) + def unsubscribe( + self, + event_type: EventType, + handler: Callable[[DomainEvent], None], + ) -> bool: + """Remove *handler* from the subscription list for *event_type*. + + If the handler was registered for multiple event types this only + removes it from the specified event type. The handler remains + callable for any other event types it was subscribed to. + + Returns ``True`` if a matching subscription was found and removed, + ``False`` otherwise (handler not registered for *event_type*). + + Args: + event_type: The :class:`EventType` from which to remove *handler*. + handler: The callable previously registered with :meth:`subscribe`. + + Raises: + TypeError: If *event_type* is not an :class:`EventType`. + TypeError: If *handler* is not callable. + """ + if not isinstance(event_type, EventType): + raise TypeError( + f"event_type must be an EventType, got {type(event_type).__name__!r}" + ) + if not callable(handler): + raise TypeError("handler must be callable") + + handlers = self._subscriptions.get(event_type, []) + try: + handlers.remove(handler) + except ValueError: + return False + + # Clean up empty lists to prevent lingering empty dict entries. + if event_type in self._subscriptions and not self._subscriptions[event_type]: + del self._subscriptions[event_type] + + _logger.debug( + "event_handler_unsubscribed", + event_type=event_type.value, + handler=getattr(handler, "__qualname__", repr(handler)), + ) + return True + __all__ = ["LoggingEventBus"] diff --git a/src/cleveragents/infrastructure/events/protocol.py b/src/cleveragents/infrastructure/events/protocol.py index 2c3cae85f..3a8e391b4 100644 --- a/src/cleveragents/infrastructure/events/protocol.py +++ b/src/cleveragents/infrastructure/events/protocol.py @@ -8,6 +8,7 @@ decoupled from infrastructure choices. Based on: - docs/specification.md §Event-Driven Architecture - Forgejo issue #473 + - Forgejo issue #10356 (unsubscribe support) """ from __future__ import annotations @@ -26,6 +27,13 @@ class EventBus(Protocol): Implementations may be reactive (RxPY), logging-only, or in-memory stubs. Callers (domain services) depend on this protocol and receive the concrete implementation through dependency injection. + + .. note:: + + Handler references are stored by implementations as strong references. + Callers that subscribe to an EventBus should also :meth:`unsubscribe` + handlers when they are no longer needed, otherwise the handler callable + may be retained in memory indefinitely (see Forgejo issue #10356). """ def emit(self, event: DomainEvent) -> None: @@ -49,5 +57,35 @@ class EventBus(Protocol): """ ... + def unsubscribe( + self, + event_type: EventType, + handler: Callable[[DomainEvent], None], + ) -> bool: + """Remove *handler* from events of *event_type*. + + If the handler was registered for multiple event types this only + removes it from the specified event type. The handler remains + callable for any other event types it was subscribed to. + + Returns ``True`` if a matching subscription was found and removed, + ``False`` otherwise (handler not registered for *event_type*). + + .. note:: + + Failing to call :meth:`unsubscribe` after subscribing leaves the + handler reference in the bus's internal registry, which can cause + memory leaks in long-running applications (see Forgejo issue #10356). + + Args: + event_type: The :class:`EventType` from which to remove *handler*. + handler: The callable previously registered with :meth:`subscribe`. + + Raises: + TypeError: If *event_type* is not an :class:`EventType`. + TypeError: If *handler* is not callable. + """ + ... + __all__ = ["EventBus"] diff --git a/src/cleveragents/infrastructure/events/reactive.py b/src/cleveragents/infrastructure/events/reactive.py index 15af69bb7..565eb08cc 100644 --- a/src/cleveragents/infrastructure/events/reactive.py +++ b/src/cleveragents/infrastructure/events/reactive.py @@ -14,6 +14,7 @@ debouncing can operate directly on :attr:`stream`. Based on: - docs/specification.md §Event-Driven Architecture - Forgejo issue #473 + - Forgejo issue #10356 (unsubscribe support) """ from __future__ import annotations @@ -50,6 +51,7 @@ class ReactiveEventBus: _stream: Read-only Observable view over ``_subject``. _subscriptions: Per-event-type handler lists. _audit_log: Volatile in-memory log of all emitted events. + _closed: Whether :meth:`close` has been called. """ def __init__(self, max_audit_log_size: int | None = None) -> None: @@ -60,6 +62,7 @@ class ReactiveEventBus: self._stream: Observable = self._subject.pipe(ops.map(_identity)) self._subscriptions: dict[EventType, list[Callable[[DomainEvent], None]]] = {} self._audit_log: deque[DomainEvent] = deque(maxlen=max_audit_log_size) + self._closed: bool = False # ------------------------------------------------------------------ # Public API (satisfies EventBus protocol) @@ -106,8 +109,15 @@ class ReactiveEventBus: event: The :class:`DomainEvent` to publish. Raises: + RuntimeError: If :meth:`close` has already been called. TypeError: If *event* is not a :class:`DomainEvent`. """ + if self._closed: + raise RuntimeError( + "ReactiveEventBus.emit() called after close() — " + "the bus has been shut down and the RxPY Subject is completed" + ) + if not isinstance(event, DomainEvent): raise TypeError( f"event must be a DomainEvent, got {type(event).__name__!r}" @@ -164,6 +174,52 @@ class ReactiveEventBus: raise TypeError("handler must be callable") self._subscriptions.setdefault(event_type, []).append(handler) + def unsubscribe( + self, + event_type: EventType, + handler: Callable[[DomainEvent], None], + ) -> bool: + """Remove *handler* from the subscription list for *event_type*. + + If the handler was registered for multiple event types this only + removes it from the specified event type. The handler remains + callable for any other event types it was subscribed to. + + Returns ``True`` if a matching subscription was found and removed, + ``False`` otherwise (handler not registered for *event_type*). + + Args: + event_type: The :class:`EventType` from which to remove *handler*. + handler: The callable previously registered with :meth:`subscribe`. + + Raises: + TypeError: If *event_type* is not an :class:`EventType`. + TypeError: If *handler* is not callable. + """ + if not isinstance(event_type, EventType): + raise TypeError( + f"event_type must be an EventType, got {type(event_type).__name__!r}" + ) + if not callable(handler): + raise TypeError("handler must be callable") + + handlers = self._subscriptions.get(event_type, []) + try: + handlers.remove(handler) + except ValueError: + return False + + # Clean up empty lists to prevent lingering empty dict entries. + if event_type in self._subscriptions and not self._subscriptions[event_type]: + del self._subscriptions[event_type] + + _logger.debug( + "event_handler_unsubscribed", + event_type=event_type.value, + handler=getattr(handler, "__qualname__", repr(handler)), + ) + return True + @property def stream(self) -> Observable: """Read-only observable stream for advanced RxPY operators. @@ -175,12 +231,22 @@ class ReactiveEventBus: return self._stream def close(self) -> None: - """Signal completion on the reactive stream and clear all subscriptions. + """Complete the RxPY stream and prevent further event emission. - Call this when the bus is no longer needed (e.g. in test teardown) - to release any RxPY Subject resources and prevent subscription leaks - between test scenarios. + Signals ``on_completed()`` to all RxPY subscribers, allowing them to + finalize their state (e.g. buffering operators flush, aggregation + operators emit their final result). After ``close()``, any call to + :meth:`emit` raises :exc:`RuntimeError`. + + This method is idempotent — calling it more than once is safe. + + Call this when the bus is no longer needed (e.g. in test teardown or + application shutdown) to release RxPY Subject resources and prevent + subscription leaks between test scenarios. """ + if self._closed: + return + self._closed = True with contextlib.suppress(Exception): self._subject.on_completed() self._subscriptions.clear()