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
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
21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 | |
__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
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
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
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.