From 7cff348d0de25f68cac8c8c958579df70b1c95cd Mon Sep 17 00:00:00 2001 From: Artemiy Solopov Date: Thu, 16 Jul 2026 19:28:11 +0300 Subject: [PATCH] Refactoring --- lib/neon_vertigo_service/message_processor.ex | 33 ++++++++++++++++--- lib/neon_vertigo_service/mqtt/server.ex | 23 ++++++------- 2 files changed, 38 insertions(+), 18 deletions(-) diff --git a/lib/neon_vertigo_service/message_processor.ex b/lib/neon_vertigo_service/message_processor.ex index 83a739e..929d4f6 100644 --- a/lib/neon_vertigo_service/message_processor.ex +++ b/lib/neon_vertigo_service/message_processor.ex @@ -1,17 +1,40 @@ defmodule NeonVertigoService.MessageProcessor do - def process_message(["auctions", "create"], payload, _properties) do + alias NeonVertigoService.Repo + + def process_message(mqtt_pid, {"auctions", "create"}, payload, properties) do changeset = Auction.changeset(%Auction{}, payload) - reply = case NeonVertigoService.Repo.insert(changeset) do + reply = case Repo.insert(changeset) do {:ok, auction} -> %{ok: true, gid: auction.gid} {:error, changeset} -> %{ok: false, gid: changeset |> Ecto.Changeset.get_field(:gid)} end - {:reply, reply} + reply(reply, mqtt_pid, properties) end - def process_message(_, _, _) do - :noreply + def process_message(_, _, _, _) do + {:ok, :unknown_message} end + + defp reply(data, mqtt_pid, orig_properties, message_properties \\ []) do + response_topic = orig_properties[:"Response-Topic"] + if is_nil(response_topic) do + {:ok, :no_response_topic} + else + send_mqtt(data, mqtt_pid, response_topic, message_properties) + end + end + + defp send_mqtt(data, mqtt_pid, response_topic, message_properties \\ []) + + defp send_mqtt(data, mqtt_pid, response_topic, message_properties) when is_binary(response_topic) and byte_size(response_topic) > 0 do + if String.valid?(response_topic) && String.length(response_topic) > 0 do + :emqtt.publish(mqtt_pid, response_topic, JSON.encode!(data), message_properties) + else + {:error, :invalid_topic} + end + end + + defp send_mqtt(_, _, _, _), do: {:error, :invalid_topic} end diff --git a/lib/neon_vertigo_service/mqtt/server.ex b/lib/neon_vertigo_service/mqtt/server.ex index b119167..85d4238 100644 --- a/lib/neon_vertigo_service/mqtt/server.ex +++ b/lib/neon_vertigo_service/mqtt/server.ex @@ -1,16 +1,20 @@ defmodule NeonVertigoService.Mqtt.Server do use GenServer + alias __MODULE__, as: MqttServer + + defstruct [:mqtt_pid] + def start_link(_) do GenServer.start_link(__MODULE__, [], name: __MODULE__) end def init(_) do - {:ok, {}, {:continue, :post_init}} + {:ok, %MqttServer{}, {:continue, :post_init}} end # Callbacks - def handle_continue(:post_init, {}) do + def handle_continue(:post_init, %MqttServer{} = state) do mqtt_config = Application.fetch_env!(:neon_vertigo_service, :mqtt_params) {:ok, mqtt_pid} = Supervisor.start_child( NeonVertigoService.Mqtt.Supervisor, @@ -23,25 +27,18 @@ defmodule NeonVertigoService.Mqtt.Server do :emqtt.subscribe(mqtt_pid, %{}, [{"$share/neon-vertigo-service/auctions/+", [{:qos, 1}]}]) - {:noreply, {mqtt_pid}} + {:noreply, %{state | mqtt_pid: mqtt_pid}} end # Handling MQTT messages - def handle_info({:publish, msg}, {mqtt_pid} = state) do + def handle_info({:publish, msg}, %{mqtt_pid: 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"] + topic = String.split(topic_str, "/") |> List.to_tuple() 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 + NeonVertigoService.MessageProcessor.process_message(mqtt_pid, topic, payload, properties) {:noreply, state} end