defmodule NeonVertigoService.Mqtt.Server do use GenServer def start_link(_) do GenServer.start_link(__MODULE__, [], name: __MODULE__) end def init(_) do {:ok, {}, {:continue, :post_init}} end # Callbacks def handle_continue(:post_init, {}) do mqtt_config = Application.fetch_env!(:neon_vertigo_service, :mqtt_params) {:ok, mqtt_pid} = Supervisor.start_child( NeonVertigoService.Mqtt.Supervisor, %{ id: :neon_vertigo_emqtt, start: {:emqtt, :start_link, [[owner: self(), name: :neon_vertigo_service_emqtt, proto_ver: :v5] ++ mqtt_config]} } ) {:ok, _mqtt_props} = :emqtt.connect(mqtt_pid) :emqtt.subscribe(mqtt_pid, %{}, [{"$share/neon-vertigo-service/auctions/+", [{:qos, 1}]}]) {:noreply, {mqtt_pid}} end # Handling MQTT messages def handle_info({:publish, msg}, {mqtt_pid} = state) do IO.puts("Received message") IO.inspect(msg) %{payload: payload_str, topic: topic_str, properties: properties} = msg topic = String.split(topic_str, "/") response_topic = properties[:"Response-Topic"] payload = JSON.decode!(payload_str) case NeonVertigoService.MessageProcessor.process_message(topic, payload, properties) do {:reply, data} -> if is_binary(response_topic) && String.length(response_topic) > 0 do :emqtt.publish(mqtt_pid, response_topic, JSON.encode!(data), retain: true) end :noreply -> nil end {:noreply, state} end end