How Phoenix PubSub and the MQTT Publisher Coordinate Data Updates in TeslaMate
Phoenix PubSub distributes internal vehicle state updates via Elixir message passing, while the MQTT Publisher forwards those updates to an external broker, creating a clean separation between state distribution and transport concerns.
TeslaMate maintains the real-time state of each vehicle in a %Summary{} struct fetched from the Tesla API. To expose this data to home-automation systems, the application coordinates between an internal Phoenix PubSub event bus and the MQTT Publisher module that interfaces with the Tortoise311 MQTT client. This design keeps the core vehicle logic decoupled from network transport implementation details.
Architecture Overview
The coordination relies on three primary components working in sequence. First, TeslaMate.Application initializes the infrastructure by starting both the Phoenix PubSub server named TeslaMate.PubSub and the TeslaMate.Mqtt.PubSub supervisor. This setup establishes the internal messaging fabric and the process supervision tree required to handle MQTT publication.
Second, TeslaMate.Mqtt.PubSub acts as a supervisor that spawns a dedicated VehicleSubscriber process for each known car via VehicleSubscriber.start_link/1. Each subscriber process isolates MQTT publication concerns per vehicle, ensuring that failures in one car's publication stream do not affect others.
Third, the TeslaMate.Mqtt.Publisher GenServer provides a thin wrapper around the Tortoise311 MQTT client. It handles the actual socket communication, QoS acknowledgments, and connection state management, presenting a simple publish/3 API to the rest of the application.
Step-by-Step Data Flow
1. Subscription Initialization
When a VehicleSubscriber process starts, its init/1 function calls TeslaMate.Vehicles.subscribe_to_summary(car_id), which internally executes Phoenix.PubSub.subscribe/2 on the topic specific to that vehicle. According to the source code in lib/teslamate/mqtt/pubsub/vehicle_subscriber.ex, the subscriber also clears any stale retained MQTT messages to ensure subscribers on the broker receive fresh state immediately.
2. Broadcasting Vehicle Updates
The TeslaMate.Vehicles module periodically polls the Tesla API and constructs %Summary{} structs containing the latest vehicle data. Upon receiving new data, Vehicles.broadcast_summary/1 publishes the struct to the PubSub topic returned by summary_topic(car_id), making the update available to all local subscribers without blocking on network I/O.
3. Processing and Filtering
The VehicleSubscriber process receives the %Summary{} struct in its handle_info/2 callback. It extracts the fields designated for external exposure, builds a key → value mapping, and filters out fields listed in the @do_not_retain module attribute. This filtering occurs in lib/teslamate/mqtt/pubsub/vehicle_subscriber.ex and ensures that transient values are not incorrectly retained by the MQTT broker.
4. Topic Construction and Delegation
For each selected field, the subscriber constructs an MQTT topic following the pattern teslamate/<namespace>/cars/<car_id>/<key> and calls VehicleSubscriber.publish/3. This function delegates to TeslaMate.Mqtt.Publisher.publish/3, which forwards the message to Tortoise311. The Publisher module handles QoS 0 messages as fire-and-forget operations, while QoS 1 messages generate references that are tracked in the GenServer's handle_info/2 callback to await acknowledgment from the broker.
Code Examples
Subscribing to Vehicle Updates via PubSub
Any process can receive vehicle state changes by subscribing to the Phoenix PubSub topic, mirroring the pattern used by VehicleSubscriber:
defmodule MyApp.CarWatcher do
use GenServer
def start_link(car_id) do
GenServer.start_link(__MODULE__, car_id, name: __MODULE__)
end
@impl true
def init(car_id) do
# Same call that VehicleSubscriber uses internally
:ok = TeslaMate.Vehicles.subscribe_to_summary(car_id)
{:ok, %{car_id: car_id}}
end
@impl true
def handle_info(%TeslaMate.Vehicles.Vehicle.Summary{} = summary, state) do
IO.inspect(summary, label: "Received summary for #{state.car_id}")
{:noreply, state}
end
end
Publishing Custom MQTT Messages
To publish arbitrary data to the MQTT broker from within the application:
# Build a topic like "teslamate/home/cars/1/custom_event"
topic = ["teslamate", "home", "cars", "1", "custom_event"]
|> Enum.join("/")
# Send a retained QoS-1 message
TeslaMate.Mqtt.Publisher.publish(topic, "door_opened", retain: true, qos: 1)
Resulting MQTT Message Structure
External subscribers to the MQTT broker receive messages structured as follows:
topic: teslamate/home/cars/1/state
payload: "online"
retain: true
qos: 1
Key Files and Responsibilities
lib/teslamate/application.ex– Starts the Phoenix PubSub server (TeslaMate.PubSub) and theTeslaMate.Mqtt.PubSubsupervisor during application boot.lib/teslamate/mqtt/pubsub.ex– Supervisor module that manages the lifecycle ofVehicleSubscriberprocesses, creating one per vehicle.lib/teslamate/mqtt/pubsub/vehicle_subscriber.ex– Listens to%Summary{}updates via PubSub, transforms them into the MQTT topic hierarchy, and delegates publication to the Publisher module.lib/teslamate/mqtt/publisher.ex– GenServer wrapper around Tortoise311 that manages the MQTT connection, handles QoS handshakes, and provides the primarypublish/3interface.lib/teslamate/vehicles.ex– Fetches data from the Tesla API, constructs%Summary{}structs, and broadcasts them on the internal PubSub bus.
Summary
- Phoenix PubSub acts as the internal event bus, allowing processes to subscribe to vehicle-specific
%Summary{}updates without coupling to network code. VehicleSubscriberprocesses bridge the internal PubSub system and the external MQTT broker, filtering data and constructing topic hierarchies.- MQTT Publisher isolates the Tortoise311 client implementation, handling connection state, QoS levels, and transmission acknowledgments.
- The architecture separates concerns by keeping Tesla API polling logic in
TeslaMate.Vehicles, state distribution in PubSub, and network transport in the MQTT-specific modules.
Frequently Asked Questions
How does TeslaMate ensure MQTT messages are not lost during connection drops?
The TeslaMate.Mqtt.Publisher GenServer tracks QoS 1 message references in its state. When publish/3 is called with QoS 1, it stores the reference and waits for the acknowledgment via handle_info/2. QoS 0 messages are fire-and-forget and do not guarantee delivery. For critical state that must persist, the subscriber sets the retain flag based on the @do_not_retain list logic, ensuring the broker stores the last known good value for new subscribers.
Can I subscribe to vehicle updates without using MQTT?
Yes. Any Elixir process can subscribe directly to the Phoenix PubSub topic by calling TeslaMate.Vehicles.subscribe_to_summary(car_id), which invokes Phoenix.PubSub.subscribe/2 on the TeslaMate.PubSub instance. This allows internal modules or custom extensions to react to vehicle state changes without the overhead of MQTT protocol encoding or network transmission.
Why does each vehicle have a dedicated VehicleSubscriber process?
The TeslaMate.Mqtt.PubSub supervisor spawns a VehicleSubscriber per car to isolate failure domains and manage subscription lifecycle. If one vehicle's data stream encounters an error, the supervisor can restart that specific subscriber without affecting MQTT publication for other vehicles. This per-vehicle architecture also simplifies topic construction and filtering logic, as each process only handles a single car's state.
What determines whether an MQTT message is retained?
According to the implementation in lib/teslamate/mqtt/pubsub/vehicle_subscriber.ex, the subscriber checks field names against the @do_not_retain module attribute. Fields in this list publish without the retain flag, while all others publish with retain: true. This ensures transient values (like current speed) are not retained by the broker, while persistent state (like charge level or locked status) remains available for new subscribers connecting to the topic.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →