
    Rmj                         d Z ddlmZ ddlmZ ddlmZ ddlmZ	 ddlm
Z erddlmZ ded	d
ddfdZdedefdZddZg dZdS )a]  Global event subscriber registration.

Subscribers observe all lifecycle events emitted by the current process,
including scope, tool, LLM, and mark events. They are typically used for
logging, metrics, tracing, and custom observability pipelines.

Example::

    import nemo_relay

    def log_event(event):
        print(f"{event.kind}: {event.name}")

    nemo_relay.subscribers.register("logger", log_event)
    try:
        with nemo_relay.scope.scope("demo", nemo_relay.ScopeType.Agent):
            nemo_relay.scope.event("started")
    finally:
        nemo_relay.subscribers.deregister("logger")
    )Callable)TYPE_CHECKING)deregister_subscriber)flush_subscribers)register_subscriber)EventnamecallbackzCallable[[Event], None]returnNc                 "    t          | |          S )a  Register a global event subscriber.

    Args:
        name: Unique subscriber name.
        callback: Callable invoked as ``callback(event)`` for every emitted
            lifecycle event.

    Returns:
        None: This function returns after the subscriber is registered.

    Raises:
        RuntimeError: If a subscriber with the same name already exists.

    Example::

        import nemo_relay

        nemo_relay.subscribers.register("printer", lambda event: print(event.kind))
    )_native_register)r	   r
   s     ^/home/thesage/.hermes/hermes-agent/venv/lib/python3.11/site-packages/nemo_relay/subscribers.pyregisterr   *   s    ( D(+++    c                      t          |           S )av  Remove a previously registered global subscriber.

    Args:
        name: Subscriber name passed to ``register()``.

    Returns:
        ``True`` if a subscriber was removed, otherwise ``False``.

    Notes:
        Deregistering a subscriber affects only future event delivery. Events
        already emitted before removal carry a subscriber snapshot, so queued
        callbacks from that snapshot may still run.

    Example::

        import nemo_relay

        nemo_relay.subscribers.register("printer", lambda event: None)
        removed = nemo_relay.subscribers.deregister("printer")
        assert removed is True
    )_native_deregister)r	   s    r   
deregisterr   A   s    , d###r   c                      t                      S )a  Wait for subscriber callbacks already queued by native event emission.

    Native NeMo Relay event APIs enqueue subscriber callbacks and return without
    waiting for observer work. Use this barrier in tests and shutdown paths when
    captured subscriber output must be complete before continuing.

    Call this function outside subscriber callbacks. A re-entrant call returns
    without waiting to avoid blocking the dispatcher, so callbacks later in the
    same dispatch snapshot can still run.
    )_native_flush r   r   flushr   Z   s     ??r   )r   r   r   )r   N)__doc__collections.abcr   typingr   nemo_relay._nativer   r   r   r   r   r   
nemo_relayr   strr   boolr   r   __all__r   r   r   <module>r       s   * % $ $ $ $ $                             !      ,3 ,"; , , , , ,.$S $T $ $ $ $2    .
-
-r   