MQTT API#
Page under construction
We are currently working on this documentation, so this page is not complete yet.
Overview#
The MQTT API allows you to interact with the different assets on a site in a streaming manner. A persistent connection is setup from you to our MQTT broker, allowing continuous data exchange. You can receive live metrics and send commands to control assets on site. Data is exchanged using the Protocol Buffers format.
The MQTT API is designed for low-latency interactions compared to the HTTP API. While data from the HTTP API may have a delay of several minutes due to cleaning, processing and aggregation, data on the MQTT API is much more raw, direct, and is transported with a latency in the order of seconds.
The message stream only provides live data and does not include historical data. It is not possible to query previous metrics. If a message is published and you are not connected at that moment, this message is lost.
Our MQTT broker is reachable at:
# Recommended: this is the default port for encrypted MQTT
mqtt.octave.energy:8883
# In case port 8883 is blocked on your network, you can use 443 but this needs some additional configuration, see below
mqtt.octave.energy:443
When using port 443, your client must implement the Application Layer Protocol Negotiation (ALPN) TLS extension. Protocol name x-amzn-mqtt-ca must be used in the ClientHello message.
Security#
The MQTT API uses mTLS for secure communication. Each client must authenticate using a certificate. See the Credentials setup section for more details.
Clients connecting to the MQTT broker must use a specific client ID format, which will be provided during the onboarding process. This is typically in the format <partner_name>-{integration,production}-*, where the suffix is a wildcard which you can choose freely.
Data conventions#
Protocol Buffers (protobuf)#
Messages that are transferred over MQTT are serialized via the Protocol Buffers (protobuf) protocol, describing the schema of each message. You need this schema in order to decode the messages (as opposed to the HTTP API which uses JSON as schemaless format). Based on the protobuf (.proto) files available in the following zipfile, you will need to compile them into a package for your appropriate programming language.
As an example, the documentation for python can be found here: Python protobuf documentation.
Zipfile containing the protobuf (.proto) schema: proto_20250910.zip
Data batching#
All metrics and states are batched together by category when emitted to their respective topic. At plant level, a batch of events is guaranteed to belong to one single plant. At unit level, all events in a batch belong to the same unit.
Under normal circumstances, each metric and states are emitted at every level at least once per minute.
For example, a batch of metrics received at octave/site/e64e0ba7-bbfa-4540-950c-526e24301b37/site_meter might contain the following events:
[
{
timestamp_ms: 1751517515239
metric_name: ACTIVE_POWER_PHASE1_KW
value: 0.701
},
{
timestamp_ms: 1751517513258
metric_name: ACIVE_POWER_PHASE2_KW
value: -1.214
},
{
timestamp_ms: 1751517511247
metric_name: ACTIVE_POWER_PHASE3_KW
value: 0.411
},
{
timestamp_ms: 1751517519404
metric_name: ACTIVE_POWER_KW
value: 0.149
},
{
timestamp_ms: 1751517519232
metric_name: ACTIVE_POWER_PHASE1_KW
value: 1.01
}
]
Parallel clients#
You can use multiple readers to consume messages from the same topic. Every reader must use a separate clientID to connect. You will receive the allowed clientID format you can use from us, which you can suffix with your own identifier. In this case, every client listening to the same topic will get every message on the topic.
We support the MQTTv5 feature shared subscriptions, so you can scale multiple subscribers horizontally by putting them in a group. Each message will only be delivered to one of the subscribers, the traffic is load-balanced by the broker and you don't need to implement your own mechanism to deduplicate processing of incoming messages. Note that you still need to use different clientIDs when connecting.
MQTT topics#
Different messages are published to different topics. Each topic corresponds to a specific type of data or event. By subscribing to the appropriate topics, clients can receive only the messages they are interested in.
You can subscribe to all messages, or to a subset of messages via MQTT. This is configured by subscribing to a specific topic pattern. The topic structure below allows the client application to selectively register to specific topics, or to all metrics with a path similar to octave/site/site123/bess_plant/+/metrics/# , or any other valid MQTT subscription pattern.
Reading metrics#
| Asset type | Plant-level | Unit-level |
|---|---|---|
| Site power meter | Topic name: octave/site/{site_id}/site_meter Available metrics: ACTIVE_POWER_KW ACTIVE_POWER_PHASE1_KW ACTIVE_POWER_PHASE2_KW ACTIVE_POWER_PHASE3_KW REACTIVE_POWER_KVAR REACTIVE_POWER_PHASE1_KVAR REACTIVE_POWER_PHASE2_KVAR REACTIVE_POWER_PHASE3_KVAR | |
| Site settings | Topic name: octave/site/{site_id}/site/settings Available metrics: MAX_GRID_OFFTAKE_KW MAX_GRID_INJECTION_KW MAX_CURRENT_A | |
| BESS PCS | Topic name: octave/site/{site_id}/bess_plant/pcs Available metrics: SETPOINT_KW ACTIVE_POWER_KW REACTIVE_POWER_KVAR MAX_CHARGE_POWER_KW MAX_DISCHARGE_POWER_KW | Topic name: octave/site/{site_id}/bess_unit/{unit_id}/pcs Available metrics: SETPOINT_KW ACTIVE_POWER_KW REACTIVE_POWER_KVAR ACTIVE_POWER_PHASE1_KW ACTIVE_POWER_PHASE2_KW ACTIVE_POWER_PHASE3_KW REACTIVE_POWER_PHASE1_KVAR REACTIVE_POWER_PHASE2_KVAR REACTIVE_POWER_PHASE3_KVAR |
| PV | Topic name: octave/site/{site_id}/pv_plant Available metrics: AVAILABLE_INVERTER_POWER_KW ACTIVE_POWER_KW SITE_TARGET_POWER_KW (if set) | Topic name: octave/site/{site_id}/pv_unit/{unit_id} Available metrics: AVAILABLE_INVERTER_POWER_KW ACTIVE_POWER_KW MAX_ACTIVE_POWER_KW (if set) |
| PV meter | Topic name: octave/site/{site_id}/pv_meter Available metrics: ACTIVE_POWER_KW ACTIVE_POWER_PHASE1_KW ACTIVE_POWER_PHASE2_KW ACTIVE_POWER_PHASE3_KW REACTIVE_POWER_KVAR REACTIVE_POWER_PHASE1_KVAR REACTIVE_POWER_PHASE2_KVAR REACTIVE_POWER_PHASE3_KVAR AC_CURRENT_PHASE1_A AC_CURRENT_PHASE2_A AC_CURRENT_PHASE3_A AC_PHASE_VOLTAGE_PHASE1_V AC_PHASE_VOLTAGE_PHASE2_V AC_PHASE_VOLTAGE_PHASE3_V |
Reading state#
| Asset type | Plant-level | Unit-level |
|---|---|---|
| BESS PCS status | Topic name: octave/site/{site_id}/bess_plant/pcs_status Available metrics: OFF PRECHARGE STANDBY STARTING GRID_CONNECTED THROTTLED ERROR | Topic name: octave/site/{site_id}/bess_unit/{unit_id}/pcs_status Available metrics: OFF PRECHARGE STANDBY STARTING GRID_CONNECTED THROTTLED ERROR |
| BESS BMS status | Topic name: octave/site/{site_id}/bess_plant/bms_status Available metrics: STARTUP ALERT ALERT_NO_CHARGE ALERT_NO_DISCHARGE ALERT_STOP EMERGENCY BATTERY_ERROR SHUTDOWN | Topic name: octave/site/{site_id}/bess_unit/{unit_id}/bms_status Available metrics: STARTUP ALERT ALERT_NO_CHARGE ALERT_NO_DISCHARGE ALERT_STOP EMERGENCY BATTERY_ERROR SHUTDOWN |
Sending commands#
A plant can be controlled by sending commands to the topics below.
Commands sent at any level fully replace the previously sent command at the same destination. For example, sending a power set-point command to a battery unit followed by a self-consumption command to the same unit will set the battery in self-consumption.
Commands are never batched but need to be sent one by one.
| Asset type | Commands | Plant-level |
|---|---|---|
| BESS | STEERED_SET_POINT AUTO_CONTROLLER | Topic name: octave/site/{site_id}/bess_plant_commands/inbox |
| PV | PLANT_MAX_POWER SITE_TARGET_POWER | Topic name: octave/site/{site_id}/pv_plant_commands/inbox |
Via our portal it is possible to see the future setpoints commands we have received and to follow-up on the BESS behaviour following them.
Getting started#
Below is an example python script to demonstrate usage of the real-time API. It configures an MQTT client with the necessary host and certificates, subscribes to a topic, and logs all messages that come in. The messages are parsed from their protobuf format, for this you will need to have the python-compiled version of the protobufs available on your python path.
Reading#
# /// script
# dependencies = [
# "paho-mqtt",
# "click",
# "mh-structlog",
# "protobuf",
# ]
# ///
import click
import paho.mqtt.client as mqtt
from bess_messages_v1_pb2 import BessBmsEvents, BessPcsEvents, BmsStatusEvents, PcsStateEvents
from google.protobuf.message import Message
from mh_structlog import get_logger, setup
from paho.mqtt.client import CallbackAPIVersion, MQTTProtocolVersion
from pv_messages_v1_pb2 import PvEvents
from pv_meter_messages_v1_pb2 import PvMeterEvents
from site_meter_messages_v1_pb2 import SiteMeterEvents
setup()
logger = get_logger('demo_client')
@click.command()
@click.option("--topic", default="octave/site/#", help="MQTT topic to subscribe to.")
@click.option("--api_endpoint", default="mqtt.octave.energy", help="Endpoint to connect to")
@click.option("--client_id", required=True, help="Client ID to use")
@click.option("--cert_key_path", required=True, help="Path to client key")
@click.option("--cert_pem_path", required=True, help="Path to client pem")
def main( # noqa: C901
topic: str, api_endpoint: str, client_id: str, cert_key_path: str, cert_pem_path: str
):
logger.info(
"Starting mqtt client...",
broker=api_endpoint,
topic=topic,
client_id=client_id,
cert_key_path=cert_key_path,
cert_pem_path=cert_pem_path,
)
client = mqtt.Client(
client_id=client_id,
protocol=MQTTProtocolVersion.MQTTv5,
callback_api_version=CallbackAPIVersion.VERSION2,
)
client.tls_set(
certfile=str(cert_pem_path),
keyfile=str(cert_key_path),
tls_version=2,
)
@client.log_callback()
def on_log(client, userdata, level, buf):
get_logger('mqtt').debug(buf)
@client.connect_callback()
def on_connect(client, userdata, flags, reason_code, properties):
logger.info("Connected to MQTT Broker", result_code=reason_code, client_id=client._client_id.decode('utf-8'))
client.subscribe(topic, qos=1)
@client.message_callback()
def on_message(client, userdata, msg): # noqa: C901
def parse_message(topic: str, payload: bytes) -> Message: # noqa: PLR0912
match topic.split('/'):
case ["octave", "site", _site_id, "bess_plant", "pcs"]:
kls = BessPcsEvents
case ["octave", "site", _site_id, "bess_unit", _unit_id, "pcs"]:
kls = BessPcsEvents
case ["octave", "site", _site_id, "bess_plant", "bms"]:
kls = BessBmsEvents
case ["octave", "site", _site_id, "bess_unit", _unit_id, "bms"]:
kls = BessBmsEvents
case ["octave", "site", _site_id, "bess_plant", "bms_status"]:
kls = BmsStatusEvents
case ["octave", "site", _site_id, "bess_unit", _unit_id, "bms_status"]:
kls = BmsStatusEvents
case ["octave", "site", _site_id, "bess_plant", "pcs_status"]:
kls = PcsStateEvents
case ["octave", "site", _site_id, "bess_unit", _unit_id, "pcs_status"]:
kls = PcsStateEvents
case ["octave", "site", _site_id, "pv_plant"]:
kls = PvEvents
case ["octave", "site", _site_id, "pv_plant_meter"]:
kls = PvMeterEvents
case ["octave", "site", _site_id, "pv_unit", _unit_id]:
kls = PvEvents
case ["octave", "site", _site_id, "site_meter"]:
kls = SiteMeterEvents
case _:
logger.warning("Unrecognized topic to parse payload", topic=topic)
return None
obj = kls()
obj.ParseFromString(msg.payload)
if not obj.events:
logger.warning("No events found in message, probably wrongly deserialized")
return None
return obj
if event := parse_message(msg.topic, msg.payload):
logger.info(
"Received message on topic %s:\n%s",
msg.topic,
event,
metric_type=event.__class__.__name__,
topic=msg.topic,
)
client.connect(host=api_endpoint, port=8883, keepalive=60)
client.loop_forever()
if __name__ == '__main__':
main()
Writing#
The script below can be used to write a command via the mqtt api. There are different kinds of commands implemented.
# /// script
# dependencies = [
# "paho-mqtt",
# "click",
# "mh-structlog",
# "protobuf",
# ]
# ///
import datetime as dt
import time
import click
import paho.mqtt.client as mqtt
from bess_messages_v1_pb2 import (
AutoControllerBessCommand,
BessCommand,
SteeredControllerSetPointBessCommand,
SteeredControllerSetPointBessCommandPeriod,
)
from mh_structlog import get_logger, setup
from paho.mqtt.client import CallbackAPIVersion, MQTTProtocolVersion
from pv_messages_v1_pb2 import (
PvCommand,
PvPlantMaxPower,
PvPlantMaxPowerCommand,
PvSiteTarget,
PvSiteTargetPowerCommand,
)
setup()
logger = get_logger('demo_client')
def to_milliseconds(dt_obj: dt.datetime) -> int:
"""Convert datetime to milliseconds since epoch."""
return int(dt_obj.timestamp() * 1000)
@click.command()
@click.option(
"--command_type",
required=True,
help="Command to send. bess_power_profile / bess_auto / pv_plant_max_power / pv_site_target",
)
@click.option("--api_endpoint", default="mqtt.octave.energy", help="Endpoint to connect to")
@click.option("--client_id", required=True, help="Client ID to use")
@click.option("--cert_key_path", required=True, help="Path to client key")
@click.option("--cert_pem_path", required=True, help="Path to client pem")
def main( # noqa: C901
command_type: str, api_endpoint: str, client_id: str, cert_key_path: str, cert_pem_path: str
):
mqtt_broker_url = api_endpoint
logger.info(
"Starting mqtt client...",
broker=mqtt_broker_url,
client_id=client_id,
cert_key_path=cert_key_path,
cert_pem_path=cert_pem_path,
)
client = mqtt.Client(
client_id=client_id,
protocol=MQTTProtocolVersion.MQTTv5,
callback_api_version=CallbackAPIVersion.VERSION2,
)
client.tls_set(
certfile=str(cert_pem_path),
keyfile=str(cert_key_path),
tls_version=2,
)
@client.log_callback()
def on_log(client, userdata, level, buf):
get_logger('mqtt').debug(buf)
@client.connect_callback()
def on_connect(client, userdata, flags, reason_code, properties):
logger.info("Connected to MQTT Broker", result_code=reason_code, client_id=client._client_id.decode('utf-8'))
@client.publish_callback()
def on_publish(client, userdata, mid, reason_code, properties):
logger.info("Message published", result_code=reason_code, client_id=client._client_id.decode('utf-8'))
client.connect(host=mqtt_broker_url, port=8883, keepalive=60)
start = dt.datetime.now(tz=dt.UTC).replace(second=0, microsecond=0)
match command_type:
case "bess_power_profile":
topic = "octave/site/4b0e6f41-b6ec-45db-9a58-38370164c1ff/bess_plant_commands/inbox"
command = BessCommand(
set_point=SteeredControllerSetPointBessCommand(
periods=[
SteeredControllerSetPointBessCommandPeriod(
from_timestamp_ms=to_milliseconds(start),
to_timestamp_ms=to_milliseconds(start + dt.timedelta(minutes=1)),
active_power_kw=-50,
),
SteeredControllerSetPointBessCommandPeriod(
from_timestamp_ms=to_milliseconds(start + dt.timedelta(minutes=1)),
to_timestamp_ms=to_milliseconds(start + dt.timedelta(minutes=2)),
active_power_kw=10,
),
SteeredControllerSetPointBessCommandPeriod(
from_timestamp_ms=to_milliseconds(start + dt.timedelta(minutes=2)),
to_timestamp_ms=to_milliseconds(start + dt.timedelta(minutes=3)),
active_power_kw=-90,
),
]
)
)
case "bess_auto":
command = BessCommand(auto=AutoControllerBessCommand())
topic = "octave/site/4b0e6f41-b6ec-45db-9a58-38370164c1ff/bess_plant_commands/inbox"
case "pv_plant_max_power":
topic = "octave/site/4b0e6f41-b6ec-45db-9a58-38370164c1ff/pv_plant_commands/inbox"
command = PvCommand(
plant_max_power=PvPlantMaxPowerCommand(
max_power_periods=[
PvPlantMaxPower(
from_timestamp_ms=to_milliseconds(start),
to_timestamp_ms=to_milliseconds(start + dt.timedelta(minutes=1)),
max_active_power_kw=10,
),
PvPlantMaxPower(
from_timestamp_ms=to_milliseconds(start + dt.timedelta(minutes=1)),
to_timestamp_ms=to_milliseconds(start + dt.timedelta(minutes=2)),
max_active_power_kw=20,
),
PvPlantMaxPower(
from_timestamp_ms=to_milliseconds(start + dt.timedelta(minutes=2)),
to_timestamp_ms=to_milliseconds(start + dt.timedelta(minutes=3)),
max_active_power_kw=30,
),
]
)
)
case "pv_site_target":
topic = "octave/site/4b0e6f41-b6ec-45db-9a58-38370164c1ff/pv_plant_commands/inbox"
command = PvCommand(
site_target_power=PvSiteTargetPowerCommand(
site_target_periods=[
PvSiteTarget(
from_timestamp_ms=to_milliseconds(start),
to_timestamp_ms=to_milliseconds(start + dt.timedelta(minutes=1)),
site_target_power_kw=1,
),
PvSiteTarget(
from_timestamp_ms=to_milliseconds(start + dt.timedelta(minutes=1)),
to_timestamp_ms=to_milliseconds(start + dt.timedelta(minutes=2)),
site_target_power_kw=0,
),
]
)
)
case _:
logger.error("Unsupported command type", command_type=command_type)
return
client.loop_start()
logger.info("Publishing command", topic=topic, command=command)
time.sleep(1) # Sleep briefly to ensure connection is established before publishing
msg_info = client.publish(topic, command.SerializeToString(), qos=1)
msg_info.wait_for_publish()
client.disconnect()
client.loop_stop()
if __name__ == '__main__':
main()