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
- source: str = ''
- event_id: str = ''
- timestamp_ns: int = 0
- 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:
GrpcClientEvent 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:
GrpcClientEvent 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
TopicInfo
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()}")