Skip to content

Rebroadcasting data

nebra.rebroadcast(data_source, max_retries=7, initial_retry_delay=1.0, send_function=send, **send_kwargs)

Rebroadcast events from a compatible DataSource object onto the AT Protocol. Includes threading for speed, robust error handling that automatically reboots clients that crash, a Queue system to deal with busy networks or rate limited PDSs, and more!

This function is the recommended way to rebroadcast events onto the AT Protocol from an existing source in a robust way.

Parameters:

Name Type Description Default
data_source DataSource

A DataSource instance that provides events to rebroadcast.

required
max_retries int

Maximum number of retry attempts for failed sends. Defaults to 7.

7
initial_retry_delay float

Initial delay in seconds for retry attempts (exponential backoff). Defaults to 1.0.

1.0
send_function Callable

Function to use for sending events. Defaults to nebra.client.send.

send
**send_kwargs dict

Additional keyword arguments to pass to the send function.

{}
Source code in nebra/broadcast.py
def rebroadcast(
    data_source: DataSource,
    max_retries: int = 7,
    initial_retry_delay: float = 1.0,
    send_function=send,
    **send_kwargs,
):
    """Rebroadcast events from a compatible DataSource object onto the AT Protocol.
    Includes threading for speed, robust error handling that automatically reboots
    clients that crash, a Queue system to deal with busy networks or rate limited PDSs,
    and more!

    This function is the recommended way to rebroadcast events onto the AT Protocol from
    an existing source in a robust way.

    Parameters
    ----------
    data_source : DataSource
        A DataSource instance that provides events to rebroadcast.
    max_retries : int, optional
        Maximum number of retry attempts for failed sends. Defaults to 7.
    initial_retry_delay : float, optional
        Initial delay in seconds for retry attempts (exponential backoff). Defaults to 1.0.
    send_function : Callable, optional
        Function to use for sending events. Defaults to nebra.client.send.
    **send_kwargs : dict
        Additional keyword arguments to pass to the send function.
    """
    client = RebroadcastClient(
        data_source=data_source,
        max_retries=max_retries,
        initial_retry_delay=initial_retry_delay,
        send_function=send_function,
        **send_kwargs,
    )
    client.start()

nebra.DataSource

Bases: ABC

Abstract base class for event sources.

Users should subclass this and implement the run() method to provide events. The run() method should handle its own error recovery.

Attributes:

Name Type Description
event_queue Queue

A thread-safe queue for storing events.

stop_event Event

An event to signal when the data source should stop.

Examples:

The following implementation of a DataSource would periodically create new Bluesky posts.

import time
import nebra

class PostDataSource(nebra.DataSource):
    def run(self):
        counter = 1
        while not self.stop_event.is_set():
            new_post = {
                "$type": "app.bsky.feed.post",
                "text": f"This is test post {counter}.",
                "createdAt": nebra.get_atproto_utc_time()
            }
            self.add_event(new_post)
            counter += 1
            time.sleep(1)
Source code in nebra/broadcast.py
class DataSource(ABC):
    """Abstract base class for event sources.

    Users should subclass this and implement the run() method to provide events.
    The run() method should handle its own error recovery.

    Attributes
    ----------
    event_queue : queue.Queue
        A thread-safe queue for storing events.
    stop_event : threading.Event
        An event to signal when the data source should stop.

    Examples
    --------
    The following implementation of a DataSource would periodically create new Bluesky
    posts.

    ```python
    import time
    import nebra

    class PostDataSource(nebra.DataSource):
        def run(self):
            counter = 1
            while not self.stop_event.is_set():
                new_post = {
                    "$type": "app.bsky.feed.post",
                    "text": f"This is test post {counter}.",
                    "createdAt": nebra.get_atproto_utc_time()
                }
                self.add_event(new_post)
                counter += 1
                time.sleep(1)
    ```
    """

    def __init__(self, max_queue_size: int = 1000):
        """Initialize the DataSource with an event queue.

        Parameters
        ----------
        max_queue_size : int, optional
            Maximum size of the event queue. Defaults to 1000.
        """
        self.event_queue = queue.Queue(maxsize=max_queue_size)
        self.stop_event = threading.Event()

    def add_event(self, event: dict[str, Any]) -> bool:
        """Add an event to the queue.

        Parameters
        ----------
        event : dict[str, Any]
            The event to add to the queue.

        Returns
        -------
        bool
            True if the event was added, False if the queue was full.
        """
        try:
            self.event_queue.put_nowait(event)
            return True
        except queue.Full:
            print("Event queue full, dropping oldest event")
            try:
                # Remove oldest event and try again
                self.event_queue.get_nowait()
                self.event_queue.put_nowait(event)
                return True
            except queue.Empty:
                # Queue was empty after all, just put the event
                self.event_queue.put_nowait(event)
                return True

    def stop(self) -> None:
        """Signal the data source to stop.

        This method sets the stop_event, which should be checked periodically
        in the run() method to allow for clean shutdown.
        """
        self.stop_event.set()

    @abstractmethod
    def run(self) -> None:
        """Run the event source.

        This method should:
        1. Generate events and add them to the queue using add_event()
        2. Handle its own error recovery
        3. Check self.stop_event.is_set() periodically to allow clean shutdown

        Notes
        -----
        This is an abstract method that must be implemented by subclasses.
        """

__init__(max_queue_size=1000)

Initialize the DataSource with an event queue.

Parameters:

Name Type Description Default
max_queue_size int

Maximum size of the event queue. Defaults to 1000.

1000
Source code in nebra/broadcast.py
def __init__(self, max_queue_size: int = 1000):
    """Initialize the DataSource with an event queue.

    Parameters
    ----------
    max_queue_size : int, optional
        Maximum size of the event queue. Defaults to 1000.
    """
    self.event_queue = queue.Queue(maxsize=max_queue_size)
    self.stop_event = threading.Event()

add_event(event)

Add an event to the queue.

Parameters:

Name Type Description Default
event dict[str, Any]

The event to add to the queue.

required

Returns:

Type Description
bool

True if the event was added, False if the queue was full.

Source code in nebra/broadcast.py
def add_event(self, event: dict[str, Any]) -> bool:
    """Add an event to the queue.

    Parameters
    ----------
    event : dict[str, Any]
        The event to add to the queue.

    Returns
    -------
    bool
        True if the event was added, False if the queue was full.
    """
    try:
        self.event_queue.put_nowait(event)
        return True
    except queue.Full:
        print("Event queue full, dropping oldest event")
        try:
            # Remove oldest event and try again
            self.event_queue.get_nowait()
            self.event_queue.put_nowait(event)
            return True
        except queue.Empty:
            # Queue was empty after all, just put the event
            self.event_queue.put_nowait(event)
            return True

run() abstractmethod

Run the event source.

This method should: 1. Generate events and add them to the queue using add_event() 2. Handle its own error recovery 3. Check self.stop_event.is_set() periodically to allow clean shutdown

Notes

This is an abstract method that must be implemented by subclasses.

Source code in nebra/broadcast.py
@abstractmethod
def run(self) -> None:
    """Run the event source.

    This method should:
    1. Generate events and add them to the queue using add_event()
    2. Handle its own error recovery
    3. Check self.stop_event.is_set() periodically to allow clean shutdown

    Notes
    -----
    This is an abstract method that must be implemented by subclasses.
    """

stop()

Signal the data source to stop.

This method sets the stop_event, which should be checked periodically in the run() method to allow for clean shutdown.

Source code in nebra/broadcast.py
def stop(self) -> None:
    """Signal the data source to stop.

    This method sets the stop_event, which should be checked periodically
    in the run() method to allow for clean shutdown.
    """
    self.stop_event.set()