Metadata-Version: 2.4
Name: ai-box-lib
Version: 1.2.2
Summary: Python library for NXP Edge AI Industrial Platform
Author: Cedric
Requires-Python: >=3.11
Description-Content-Type: text/markdown
License-File: LICENSE
Requires-Dist: nats-py==2.15.0
Requires-Dist: protobuf==7.35.1
Requires-Dist: uuid6==2025.0.1
Dynamic: license-file

# AI Box Library

Python library for the NXP Edge AI Industrial Platform.

## Overview

`ai-box-lib` provides thin clients for component-to-component communication:

- `DataCollectorClient` for publishing data collector payloads
- `PreProcessorClient` for subscribing data collector payloads and publishing features
- `ContextEngineClient` for publishing context messages
- `ChannelClient` for generic pub/sub communication between any components

## Installation

Install from PyPI:

```bash
pip install ai-box-lib
```

## Quick Start

### 1) Publish data from a data collector

```python
from ai_box_lib.data_collector_client import DataCollectorClient

client = DataCollectorClient[dict]()
client.connect()
client.publish_timestream({"random_number": "42"})
```

### 2) Subscribe and publish from a pre-processor

```python
from ai_box_lib.pre_processor_client import PreProcessorClient

client = PreProcessorClient[dict, dict]()
client.connect()

def handle_raw(message: dict) -> None:
	processed = {"feature_a": [1.0, 2.0, 3.0]}
	client.publish_data(processed)

unsubscribe = client.subscribe_timestream(handle_raw)
```

### 3) Publish and subscribe to context data

```python
from ai_box_lib.context_engine_client import ContextEngineClient

client = ContextEngineClient[dict]()
client.connect()
client.publish_data({"state": "ok"})
```

In a multi-asset setup, a single context engine is often responsible for all assets. Pass `custom_asset_id` to publish to a different asset's context engine topic:

```python
client.publish_data({"state": "ok"}, custom_asset_id="asset-abc123")
```

```python
from ai_box_lib.context_engine_client import ContextEngineClient
client = ContextEngineClient[dict]()
client.connect()
client.subscribe()

client.context  # Access the latest context value at any time

def handle_context(message: dict) -> None:
	print(f"Context update: {message}")

client.subscribe(handle_context) # Subscribe with a handler to receive real-time updates

```

Optionally you can pass `custom_asset_id` to subscribe to a Context Engine on a different asset, as long as it's connected to the same Box.

### 4) Communicate over a generic channel

A channel lets any two components exchange messages without being tied to a specific pipeline stage. You define the channel by providing a `channel_id` string. Both publisher and subscriber must use the same `channel_id`.

Messages are delivered on the topic `{asset_id}/channel/{channel_id}`.

**Publisher:**

```python
from ai_box_lib.channel_client import ChannelClient

client: ChannelClient[None, dict] = ChannelClient("my-alerts")
client.connect()
client.publish({"severity": "high", "value": 42.0})
```

**Subscriber:**

```python
from ai_box_lib.channel_client import ChannelClient

client: ChannelClient[dict, None] = ChannelClient("my-alerts")
client.connect()

def handle_alert(message: dict) -> None:
    print(f"Alert received: {message}")

unsubscribe = client.subscribe(handle_alert)
# call unsubscribe() when done
```

A single client instance can both publish and subscribe on the same channel.

## Correlation IDs

Every message carries a **correlation id** so you can trace one piece of data through the pipeline (e.g. link an ML result back to its source frame). It is minted automatically at the origin and forwarded automatically when you publish from inside a subscription handler — no extra code needed.

Read the id of the message being handled with `current_correlation_id()`:

```python
def handle_raw(message: dict) -> None:
    frame_id = client.current_correlation_id()
    client.publish_data(process(message))  # inherits frame_id automatically
```

Automatic forwarding only works when you publish from *within* the handler. If you publish later — from another thread, a timer, or an external trigger such as an API callback — capture the id and pass it back explicitly. Every `publish*` method accepts an optional `correlation_id`:

```python
client.publish_data(process(message), correlation_id=frame_id)
```

## API Summary

All clients must call `connect()` before any publish or subscribe operation.

- `DataCollectorClient.connect()`
- `DataCollectorClient.publish_timestream(data, correlation_id?)`
- `DataCollectorClient.publish_audio(data, correlation_id?)`
- `DataCollectorClient.publish_image(data, correlation_id?)`
- `PreProcessorClient.connect()`
- `PreProcessorClient.subscribe_timestream(handler)`
- `PreProcessorClient.subscribe_audio(handler)`
- `PreProcessorClient.subscribe_image(handler)`
- `PreProcessorClient.publish_data(data, correlation_id?)`
- `ContextEngineClient.connect()`
- `ContextEngineClient.publish_data(data, custom_asset_id?, correlation_id?)`
- `ContextEngineClient.subscribe(handler?, custom_asset_id?)`
- `ContextEngineClient.context` — latest received context value (read-only)
- `ChannelClient(channel_id).connect()`
- `ChannelClient(channel_id).publish(data, correlation_id?)`
- `ChannelClient(channel_id).subscribe(handler)`
- `current_correlation_id()` — id of the message currently being handled (any client; `None` outside a handler)

## Validation and Limits

- `DataCollectorClient` validates message keys against `CHANNELS`.
- `PreProcessorClient` validates feature keys and feature shapes against `FEATURES`.
- Maximum publish payload size is 2 MB.

## License

See `LICENSE`.
