Event Bus API

Event Bus Client

class neoruntime_ipc_sdk.events.Event(topic: 'str', payload: 'dict[str, Any]', source: 'str' = '', event_id: 'str' = '', timestamp_ns: 'int' = 0, metadata: 'dict[str, str]'=<factory>)[source]

Bases: object

topic: str
payload: dict[str, Any]
source: str = ''
event_id: str = ''
timestamp_ns: int = 0
metadata: dict[str, str]
to_json()[source]
classmethod from_proto(msg)[source]
__init__(topic, payload, source='', event_id='', timestamp_ns=0, metadata=<factory>)
class neoruntime_ipc_sdk.events.TopicInfo(topic: 'str', subscriber_count: 'int', total_messages: 'int', last_message_ts: 'int')[source]

Bases: object

topic: str
subscriber_count: int
total_messages: int
last_message_ts: int
__init__(topic, subscriber_count, total_messages, last_message_ts)
class neoruntime_ipc_sdk.events.EventClient(endpoint=None)[source]

Bases: GrpcClient

Event Bus Client

Usage:

events = EventClient()

events.publish("app/alert", {"type": "person_detected"})

for event in events.subscribe("model/*/detections"):
    print(f"Received: {event.topic}")
__init__(endpoint=None)[source]
close()[source]
publish(topic, payload, persistent=False, ttl_ms=None, metadata=None, compact=False)[source]
publish_batch(events, persistent=False)[source]
subscribe(topic, filters=None, queue_size=100, drop_old=True)[source]
on_event(topic, callback, filters=None)[source]
unsubscribe(topic)[source]
list_topics()[source]
get_topic_info(topic)[source]
get_stats()[source]
get_topic_stats(topic)[source]

EventClient

class neoruntime_ipc_sdk.EventClient(endpoint=None)[source]

Bases: GrpcClient

Event Bus Client

Usage:

events = EventClient()

events.publish("app/alert", {"type": "person_detected"})

for event in events.subscribe("model/*/detections"):
    print(f"Received: {event.topic}")
__init__(endpoint=None)[source]
close()[source]
publish(topic, payload, persistent=False, ttl_ms=None, metadata=None, compact=False)[source]
publish_batch(events, persistent=False)[source]
subscribe(topic, filters=None, queue_size=100, drop_old=True)[source]
on_event(topic, callback, filters=None)[source]
unsubscribe(topic)[source]
list_topics()[source]
get_topic_info(topic)[source]
get_stats()[source]
get_topic_stats(topic)[source]

Data Types

Event

class neoruntime_ipc_sdk.Event(topic: 'str', payload: 'dict[str, Any]', source: 'str' = '', event_id: 'str' = '', timestamp_ns: 'int' = 0, metadata: 'dict[str, str]'=<factory>)[source]
topic: str
payload: dict[str, Any]
source: str = ''
event_id: str = ''
timestamp_ns: int = 0
metadata: dict[str, str]
to_json()[source]
classmethod from_proto(msg)[source]
__init__(topic, payload, source='', event_id='', timestamp_ns=0, metadata=<factory>)

TopicInfo

class neoruntime_ipc_sdk.TopicInfo(topic: 'str', subscriber_count: 'int', total_messages: 'int', last_message_ts: 'int')[source]
topic: str
subscriber_count: int
total_messages: int
last_message_ts: int
__init__(topic, subscriber_count, total_messages, last_message_ts)

Usage Examples

Publishing Events

from neoruntime_ipc_sdk import EventClient

events = EventClient()

# Publish simple event
events.publish("app/status", {"status": "running"})

# Publish event with metadata
events.publish("app/detection", {
    "timestamp": 1234567890,
    "objects": [
        {"label": "person", "score": 0.95, "bbox": [100, 200, 50, 100]},
        {"label": "car", "score": 0.88, "bbox": [300, 400, 80, 120]}
    ]
}, metadata={"camera": "cam0", "model": "person_vehicle_v1"})

# Publish persistent event
events.publish("app/alert", {"type": "person_detected"}, persistent=True)

# Batch publish
events.publish_batch([
    {"topic": "app/event1", "payload": {"data": "value1"}},
    {"topic": "app/event2", "payload": {"data": "value2"}}
])

Subscribing to Events

# Subscribe to a single topic
for event in events.subscribe("system/temperature"):
    temp = event.payload.get("value")
    print(f"Current temperature: {temp}°C")

# Subscribe to multiple topics (wildcards)
for event in events.subscribe("model/*/detections"):
    model_name = event.topic.split('/')[1]
    count = len(event.payload.get("objects", []))
    print(f"Model {model_name} detected {count} object(s)")

# Subscribe with filters
for event in events.subscribe("app/alerts", filters={"severity": "high"}):
    print(f"High severity alert: {event.payload}")

Using Callback Functions

def on_alert(event):
    alert_type = event.payload.get("type")
    print(f"Alert received: {alert_type}")

def on_detection(event):
    objects = event.payload.get("objects", [])
    print(f"Detected {len(objects)} object(s)")

# Register callbacks
events.on_event("app/alert", on_alert)
events.on_event("model/*/detections", on_detection)

# Keep running
import time
while True:
    time.sleep(1)

Topic Wildcards

The event bus supports wildcards:

# Matches model/person_v1/detections, model/car_v1/detections
events.subscribe("model/*/detections")

# Matches app/my_app/alert, app/my_app/status/running
events.subscribe("app/my_app/**")

Listing Topics

# List all active topics
topics = events.list_topics()
for topic in topics:
    print(f"{topic.topic}: {topic.subscriber_count} subscribers, {topic.total_messages} messages")

# Get specific topic info
info = events.get_topic_info("app/alerts")
if info:
    print(f"Subscribers: {info.subscriber_count}")

Getting Statistics

# Get global statistics
stats = events.get_stats()
print(f"Total subscribers: {stats['total_subscribers']}")
print(f"Total topics: {stats['total_topics']}")
print(f"Uptime: {stats['uptime_ms']}ms")

# Get specific topic statistics
topic_stats = events.get_topic_stats("app/alerts")
print(f"Published: {topic_stats['published_count']}")
print(f"Delivered: {topic_stats['delivered_count']}")
print(f"Dropped: {topic_stats['dropped_count']}")
print(f"Average latency: {topic_stats['avg_latency_us']}us")

Unsubscribing

# Unsubscribe
events.unsubscribe("app/alert")

Event Filtering

# Subscribe and filter
for event in events.subscribe("model/*/detections"):
    # Only process high-confidence detections
    objects = event.payload.get("objects", [])
    high_conf = [obj for obj in objects if obj.get("score", 0) > 0.9]

    if high_conf:
        print(f"High confidence detections: {len(high_conf)} object(s)")

Context Manager

# Use context manager for automatic connection management
with EventClient() as events:
    events.publish("app/status", {"status": "running"})
    for event in events.subscribe("app/*"):
        print(f"Received: {event.topic}")

Error Handling

from grpc import RpcError

try:
    events.publish("app/status", {"status": "running"})
except RpcError as e:
    print(f"Publish failed: {e.details()}")

try:
    for event in events.subscribe("invalid/topic"):
        process_event(event)
except RpcError as e:
    print(f"Subscription failed: {e.details()}")