Skip to content

Commit 23f3002

Browse files
author
Bolyki György
committed
fix(mqtt): avoid blocking startup on retained cleanup
1 parent 29b1711 commit 23f3002

2 files changed

Lines changed: 48 additions & 2 deletions

File tree

lib/teslamate/mqtt/pubsub/vehicle_subscriber.ex

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -49,9 +49,17 @@ defmodule TeslaMate.Mqtt.PubSub.VehicleSubscriber do
4949
}
5050

5151
:ok = call(deps.vehicles, :subscribe_to_summary, [car_id])
52-
:ok = clear_retained(car_id, namespace, deps.publisher)
5352

54-
{:ok, %State{car_id: car_id, namespace: namespace, deps: deps}}
53+
{:ok, %State{car_id: car_id, namespace: namespace, deps: deps}, {:continue, :clear_retained}}
54+
end
55+
56+
@impl true
57+
def handle_continue(
58+
:clear_retained,
59+
%State{car_id: car_id, namespace: namespace, deps: deps} = state
60+
) do
61+
:ok = clear_retained(car_id, namespace, deps.publisher)
62+
{:noreply, state}
5563
end
5664

5765
@publish_if_nil ~w(charge_energy_added charger_actual_current charger_phases

test/teslamate/mqtt/pubsub/vehicle_subscriber_test.exs

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,16 @@ defmodule TeslaMate.Mqtt.PubSub.VehicleSubscriberTest do
55
alias TeslaMate.Vehicles.Vehicle.Summary
66
alias TeslaMate.Locations.GeoFence
77

8+
defmodule BlockingPublisher do
9+
def publish(test_pid, topic, message, opts) do
10+
send(test_pid, {__MODULE__, {:publish, topic, message, opts}, self()})
11+
12+
receive do
13+
:continue -> :ok
14+
end
15+
end
16+
end
17+
818
defp start_subscriber(name, car_id, namespace \\ nil, publisher_responses \\ %{}) do
919
publisher_name = :"mqtt_publisher_#{name}"
1020
vehicles_name = :"vehicles_#{name}"
@@ -28,6 +38,34 @@ defmodule TeslaMate.Mqtt.PubSub.VehicleSubscriberTest do
2838
)
2939
end
3040

41+
test "starts before retained topic cleanup completes", %{test: name} do
42+
vehicles_name = :"vehicles_#{name}"
43+
{:ok, _pid} = start_supervised({VehiclesMock, name: vehicles_name, pid: self()})
44+
45+
test_pid = self()
46+
47+
task =
48+
Task.async(fn ->
49+
VehicleSubscriber.start_link(
50+
car_id: 0,
51+
namespace: nil,
52+
deps_publisher: {BlockingPublisher, test_pid},
53+
deps_vehicles: {VehiclesMock, vehicles_name}
54+
)
55+
end)
56+
57+
assert_receive {VehiclesMock, {:subscribe_to_summary, 0}}
58+
59+
assert_receive {BlockingPublisher,
60+
{:publish, "teslamate/cars/0/healthy", "", [retain: true, qos: 1]},
61+
subscriber_pid}
62+
63+
assert {:ok, {:ok, ^subscriber_pid}} = Task.yield(task, 500)
64+
65+
send(subscriber_pid, :continue)
66+
:ok = GenServer.stop(subscriber_pid)
67+
end
68+
3169
test "publishes vehicle data", %{test: name} do
3270
{:ok, pid} = start_subscriber(name, 0)
3371

0 commit comments

Comments
 (0)