Skip to content

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#

MQTT Client Example
# /// 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.

MQTT Write Command Example
# /// 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()