How data flows
Source to destination, with the class names used in the code. All of it runs in one process under the default role (all); see Roles for the split.
The whole picture
flowchart LR
subgraph SRC["Source integrations"]
PDU["`Vertiv rPDU
VertivDeviceReader`"]
DP["`Device plugins
IDeviceSourcePlugin`"]
MQ["MQTT"]
MB["Modbus TCP"]
EC["EmonCMS"]
HA["Home Assistant"]
PL["`Value plugins
IValueSourcePlugin`"]
end
subgraph CORE["rPDU2MQTT"]
POLL["DevicePollService"]
BUS["Message bus"]
SNAP["`SnapshotCache
latest per PDU`"]
LIVE["`Live values
IFlowValueSource`"]
AGG["`EnergyAggregationService
kWh from W, daily totals`"]
PASS["`DestinationHost
one ExportPass per poll`"]
GRAPH["`Flow graphs
power, energy, energy_d`"]
end
subgraph DST["Destinations"]
MQTTOUT["`MQTT
+ Home Assistant discovery`"]
PROM["`Prometheus /metrics
and Pushgateway`"]
EMON["EmonCMS"]
HAED["HA Energy Dashboard"]
HIST["Local history store"]
end
PDU --> POLL
DP --> POLL
POLL --> BUS --> SNAP
MQ --> LIVE
MB --> LIVE
EC --> LIVE
HA --> LIVE
PL --> LIVE
LIVE --> AGG
SNAP -- outlet energy --> AGG
AGG --> LIVE
SNAP --> PASS
LIVE --> PASS
PASS --> GRAPH
GRAPH --> MQTTOUT
GRAPH --> PROM
GRAPH --> EMON
LIVE --> HIST
SNAP --> HAED
LIVE --> HAED
SNAP --> UI["GUI and REST API"]
LIVE --> UI
- PDU snapshots — a whole PDU's devices, outlets and measurements, read in one poll.
- Live values — one reading per node and metric, from any bound source.
- EnergyAggregationService — daily totals from live values and from PDU outlet energy counters.
Both feed the flow graph, handed to the destinations once per poll.
PDU readings
sequenceDiagram
participant PDU as PDU (HTTP API)
participant Poll as DevicePollService
participant Bus as Message bus
participant Cache as SnapshotCache
participant Host as DestinationHost
participant Dest as Destinations
loop every PollInterval
Poll->>PDU: read devices, outlets, measurements
PDU-->>Poll: JSON
Poll->>Bus: PduSnapshot (instance, time, data)
Bus->>Cache: keep latest per instance
end
loop every poll interval of the primary PDU
Host->>Cache: fresh snapshots only
Host->>Host: ExportPass.Build
Host->>Dest: SendAsync(pass), each on its own
end
DevicePollServicepolls every configured PDU instance, plus any device a plugin supplies, and publishes each snapshot on the in-process bus (ChannelMessageBus).- Each bus subscriber has its own bounded channel. A slow consumer drops its oldest snapshot.
SnapshotCachekeeps the latest snapshot per instance. The GUI and the REST API read it, never the PDU, so a page load does not cause a poll.DestinationHostskips stale snapshots.
Live values for energy-flow nodes
Every node binding (EnergyFlow.Nodes[].Sources) is read through one seam, IFlowValueSource. The
sources are asked in order and the first one with a fresh reading wins:
flowchart TD Q["Value for node + metric?"] --> D["`DerivedFlowValueSource V×I, P÷PF … when bound as derived`"] D --> C["CompositeFlowValueSource"] C --> S1["`1. MQTT ingest EnergyFlowMqttSourceService`"] S1 -->|nothing fresh| S2["`2. Live value cache Modbus and plugin sources write here`"] S2 -->|nothing fresh| S3["3. Home Assistant entities"] S3 -->|nothing fresh| S4["4. EmonCMS feeds"] S4 -->|nothing fresh| S5["5. Other plugin sources"] S5 -->|nothing fresh| S6["`6. EnergyAggregationService kWh integrated from W, daily totals`"] S6 -->|nothing fresh| S7["`7. History fallback only with History.ValueFallback`"] S7 -->|nothing| N["`No data never 0`"]
- MQTT (
EnergyFlowMqttSourceService) subscribes to the bound topics and keeps the latest value. Subscriptions are reconciled on a timer, so a topic bound in the GUI works without a restart. - Modbus (
ModbusPollService) polls each physical device (host:port:unitId) once per cycle, however many connections name it. - Home Assistant, EmonCMS and plugin sources are kept in step with their bindings by
ValueSourcePluginHost. - Aggregation reads the measured sources only (never its own output), integrates watts into kWh for
nodes with no energy counter when
Aggregation.Enabledis on, and keeps each node's total since the period boundary whenTrackPeriodsis on. Those totals are saved to the cache (Valkey/Redis) whenCache.Enabled, otherwise toenergy-totals.json. See Totals and counters. - A reading older than its binding's
StaleAfterSecondsis skipped; the next source is tried. With none left, the node has no value.
The export pass
Once per poll, DestinationHost builds one ExportPass and offers it to every destination that is
switched on:
flowchart LR SNAP["Fresh PduSnapshots"] --> M["`Merge devices across instances`"] M --> R["`Readings every PDU measurement`"] M --> G["FlowTiers.Graphs"] LIVE["IFlowValueSource"] --> G CFG["`EnergyFlow config nodes, links, groups, panels, sites`"] --> G G --> T1["Graph: realpower"] G --> T2["Graph: energy"] G --> T3["Graph: energy_d"] R --> P["ExportPass"] T1 --> P T2 --> P T3 --> P P --> MQ["MQTT"] P --> PR["Prometheus"] P --> EM["EmonCMS"] P --> PLG["Plugin destinations"]
- The graphs are built once per pass and shared by every destination.
- PDU devices and outlets become nodes (
pdu:<device>,outlet:<device>:<n>) automatically; your nodes, links, groups, breakers and rooms are added fromEnergyFlow. - Each tier's value comes from a live value, a PDU measurement, or the sum of its children. See Energy flow › Accuracy.
- Each destination applies its own node-tag filter (
TiersFor). - A failing destination is shown on the Status board; the others continue.
What each destination does with it
| Destination | PDU readings | Flow tiers |
|---|---|---|
MQTT (MqttPduIntegration, MqttIntegration) |
<parent>/<serial>/… topics |
<parent>/energyflow/<id> JSON, when EnergyFlow.MqttExport is on |
| Home Assistant discovery | a device per PDU, outlet and group | a device per exported tier |
| Prometheus | rpdu2mqtt_<type> gauges |
rpdu2mqtt_flow_* gauges |
| EmonCMS | inputs / feeds | flow-node feeds |
Topic and metric names are listed under MQTT topics and on the GUI's Paths page.
Slower paths
flowchart LR CPH["`ConfigurationPublisherHost own cadence, default 5 min`"] --> DISC["`Home Assistant discovery MQTT and EmonCMS config`"] HAS["HA Energy Dashboard sync"] --> HAWS["`Home Assistant WebSocket API`"] LHW["`LocalHistoryWriterService timer, leader only`"] --> STORE["`LocalSeriesStore fixed-interval files`"] STORE --> ROUTER["`FlowHistoryRouter History.Provider`"] PROMR["Prometheus"] --> ROUTER EMR["EmonCMS"] --> ROUTER HAR["Home Assistant"] --> ROUTER ROUTER --> PAGES["`Trends, Node Trends, history sheets`"]
- Configuration.
ConfigurationPublisherHostrepublishes discovery and other configuration on its own interval. - History.
LocalHistoryWriterServicerecords every node's readings into the local store whateverHistory.Providersays; the pages read through whichever backendProvidernames. See History.
Commands back to the PDU
Only when ActionsEnabled is on:
flowchart RL HAB["Home Assistant switch / button"] -->|"…/outlets/n/set, /reboot"| OCS["OutletCommandService"] GUI["GUI: PDU Control"] --> IOC["`IOutletControl DeviceOutletControl`"] OCS --> IOC API["`REST API X-Api-Key`"] --> PDUOBJ["PDU client"] IOC --> PDUOBJ PDUOBJ -->|HTTP| DEV["PDU"]
The next poll reads the new state back.
Roles
RPDU2MQTT_ROLE / --role (Command line):
| Part | worker | api | ui |
|---|---|---|---|
| PDU polling, MQTT and Modbus ingest, plugin sources | ✓ | ||
| Energy aggregation | writes | reads | reads |
| Destinations, discovery, outlet commands | ✓ | ||
| REST API | ✓ | ||
| GUI | ✓ |
The health endpoints run in every role.
Leader lease
With RPDU2MQTT_LEADER_LEASE=true (set by the chart for graceful rollouts), only the lease holder runs destinations and writes local history. Every process serves and refreshes its own Prometheus /metrics.