Basic routing + account creation
This commit is contained in:
@@ -0,0 +1,23 @@
|
|||||||
|
defmodule NeonVertigoService.MessageProcessor do
|
||||||
|
@ets_auctions NeonVertigoService.Auctions
|
||||||
|
@ets_bids NeonVertigoService.Bids
|
||||||
|
|
||||||
|
def init_tables() do
|
||||||
|
@ets_auctions = :ets.new(@ets_auctions, [:named_table])
|
||||||
|
@ets_bids = :ets.new(@ets_bids, [:named_table])
|
||||||
|
:ok
|
||||||
|
end
|
||||||
|
|
||||||
|
def process_message(["auctions", "create"], payload, _properties) do
|
||||||
|
gid = payload["gid"]
|
||||||
|
data = Map.take(payload, ["author_gid", "min_bet"]) |> Enum.map(fn {k, v} -> {String.to_atom(k), v} end) |> Map.new
|
||||||
|
{min_bet, _} = Decimal.parse(data[:min_bet])
|
||||||
|
data = Map.put(data, :min_bet, min_bet)
|
||||||
|
:ets.insert_new(@ets_auctions, {gid, data})
|
||||||
|
:noreply
|
||||||
|
end
|
||||||
|
|
||||||
|
def process_message(_, _, _) do
|
||||||
|
:noreply
|
||||||
|
end
|
||||||
|
end
|
||||||
@@ -1,4 +1,4 @@
|
|||||||
defmodule NeonVertigoService.Mqtt.Service do
|
defmodule NeonVertigoService.Mqtt.Server do
|
||||||
use GenServer
|
use GenServer
|
||||||
|
|
||||||
def start_link(_) do
|
def start_link(_) do
|
||||||
@@ -6,30 +6,41 @@ defmodule NeonVertigoService.Mqtt.Service do
|
|||||||
end
|
end
|
||||||
|
|
||||||
def init(_) do
|
def init(_) do
|
||||||
{:ok, {}, {:continue, :init_mqtt}}
|
{:ok, {}, {:continue, :post_init}}
|
||||||
end
|
end
|
||||||
|
|
||||||
# Callbacks
|
# Callbacks
|
||||||
def handle_continue(:init_mqtt, {}) do
|
def handle_continue(:post_init, {}) do
|
||||||
|
:ok = NeonVertigoService.MessageProcessor.init_tables()
|
||||||
|
|
||||||
mqtt_config = Application.fetch_env!(:neon_vertigo_service, :mqtt_params)
|
mqtt_config = Application.fetch_env!(:neon_vertigo_service, :mqtt_params)
|
||||||
{:ok, mqtt_pid} = Supervisor.start_child(
|
{:ok, mqtt_pid} = Supervisor.start_child(
|
||||||
NeonVertigoService.Mqtt.Supervisor,
|
NeonVertigoService.Mqtt.Supervisor,
|
||||||
%{
|
%{
|
||||||
id: :neon_vertigo_emqtt,
|
id: :neon_vertigo_emqtt,
|
||||||
start: {:emqtt, :start_link, [[owner: self(), name: :neon_vertigo_service_emqtt] ++ mqtt_config]}
|
start: {:emqtt, :start_link, [[owner: self(), name: :neon_vertigo_service_emqtt, proto_ver: :v5] ++ mqtt_config]}
|
||||||
}
|
}
|
||||||
)
|
)
|
||||||
{:ok, _mqtt_props} = :emqtt.connect(mqtt_pid)
|
{:ok, _mqtt_props} = :emqtt.connect(mqtt_pid)
|
||||||
|
|
||||||
:emqtt.subscribe(mqtt_pid, %{}, [{"$share/neon-vertigo-service/test/+", [{:qos, 1}]}])
|
:emqtt.subscribe(mqtt_pid, %{}, [{"$share/neon-vertigo-service/auctions/+", [{:qos, 1}]}])
|
||||||
|
|
||||||
{:noreply, {mqtt_pid}}
|
{:noreply, {mqtt_pid}}
|
||||||
end
|
end
|
||||||
|
|
||||||
# Handling MQTT messages
|
# Handling MQTT messages
|
||||||
def handle_info({:publish, msg}, state) do
|
def handle_info({:publish, msg}, {mqtt_pid} = state) do
|
||||||
IO.puts("Received message")
|
IO.puts("Received message")
|
||||||
IO.inspect(msg)
|
IO.inspect(msg)
|
||||||
|
%{payload: payload_str, topic: topic_str, properties: properties} = msg
|
||||||
|
topic = String.split(topic_str, "/")
|
||||||
|
payload = JSON.decode!(payload_str)
|
||||||
|
|
||||||
|
case NeonVertigoService.MessageProcessor.process_message(topic, payload, properties) do
|
||||||
|
{:reply, topic, data} -> :emqtt.publish(mqtt_pid, topic, data)
|
||||||
|
:noreply -> nil
|
||||||
|
end
|
||||||
|
|
||||||
{:noreply, state}
|
{:noreply, state}
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
|
|||||||
@@ -10,7 +10,7 @@ defmodule NeonVertigoService.Mqtt.Supervisor do
|
|||||||
@impl true
|
@impl true
|
||||||
def init(_init_arg) do
|
def init(_init_arg) do
|
||||||
children = [
|
children = [
|
||||||
NeonVertigoService.Mqtt.Service
|
NeonVertigoService.Mqtt.Server
|
||||||
]
|
]
|
||||||
|
|
||||||
Supervisor.init(children, strategy: :one_for_all)
|
Supervisor.init(children, strategy: :one_for_all)
|
||||||
|
|||||||
@@ -23,6 +23,7 @@ defmodule NeonVertigoService.MixProject do
|
|||||||
defp deps do
|
defp deps do
|
||||||
[
|
[
|
||||||
{:emqtt, "~> 1.14"},
|
{:emqtt, "~> 1.14"},
|
||||||
|
{:decimal, "~> 2.3"},
|
||||||
# {:dep_from_hexpm, "~> 0.3.0"},
|
# {:dep_from_hexpm, "~> 0.3.0"},
|
||||||
# {:dep_from_git, git: "https://github.com/elixir-lang/my_dep.git", tag: "0.1.0"}
|
# {:dep_from_git, git: "https://github.com/elixir-lang/my_dep.git", tag: "0.1.0"}
|
||||||
]
|
]
|
||||||
|
|||||||
@@ -11,6 +11,7 @@
|
|||||||
"exqlite": {:hex, :exqlite, "0.33.1", "0465fdb997be174edeba6a27496fa27dfe8bc79ef1324a723daa8f0e8579da24", [:make, :mix], [{:cc_precompiler, "~> 0.1", [hex: :cc_precompiler, repo: "hexpm", optional: false]}, {:db_connection, "~> 2.1", [hex: :db_connection, repo: "hexpm", optional: false]}, {:elixir_make, "~> 0.8", [hex: :elixir_make, repo: "hexpm", optional: false]}, {:table, "~> 0.1.0", [hex: :table, repo: "hexpm", optional: true]}], "hexpm", "b3db0c9ae6e5ee7cf84dd0a1b6dc7566b80912eb7746d45370f5666ed66700f9"},
|
"exqlite": {:hex, :exqlite, "0.33.1", "0465fdb997be174edeba6a27496fa27dfe8bc79ef1324a723daa8f0e8579da24", [:make, :mix], [{:cc_precompiler, "~> 0.1", [hex: :cc_precompiler, repo: "hexpm", optional: false]}, {:db_connection, "~> 2.1", [hex: :db_connection, repo: "hexpm", optional: false]}, {:elixir_make, "~> 0.8", [hex: :elixir_make, repo: "hexpm", optional: false]}, {:table, "~> 0.1.0", [hex: :table, repo: "hexpm", optional: true]}], "hexpm", "b3db0c9ae6e5ee7cf84dd0a1b6dc7566b80912eb7746d45370f5666ed66700f9"},
|
||||||
"getopt": {:hex, :getopt, "1.0.3", "4f3320c1f6f26b2bec0f6c6446b943eb927a1e6428ea279a1c6c534906ee79f1", [:rebar3], [], "hexpm", "7e01de90ac540f21494ff72792b1e3162d399966ebbfc674b4ce52cb8f49324f"},
|
"getopt": {:hex, :getopt, "1.0.3", "4f3320c1f6f26b2bec0f6c6446b943eb927a1e6428ea279a1c6c534906ee79f1", [:rebar3], [], "hexpm", "7e01de90ac540f21494ff72792b1e3162d399966ebbfc674b4ce52cb8f49324f"},
|
||||||
"gun": {:hex, :gun, "2.1.0", "b4e4cbbf3026d21981c447e9e7ca856766046eff693720ba43114d7f5de36e87", [:make, :rebar3], [{:cowlib, "2.13.0", [hex: :cowlib, repo: "hexpm", optional: false]}], "hexpm", "52fc7fc246bfc3b00e01aea1c2854c70a366348574ab50c57dfe796d24a0101d"},
|
"gun": {:hex, :gun, "2.1.0", "b4e4cbbf3026d21981c447e9e7ca856766046eff693720ba43114d7f5de36e87", [:make, :rebar3], [{:cowlib, "2.13.0", [hex: :cowlib, repo: "hexpm", optional: false]}], "hexpm", "52fc7fc246bfc3b00e01aea1c2854c70a366348574ab50c57dfe796d24a0101d"},
|
||||||
|
"jason": {:hex, :jason, "1.4.4", "b9226785a9aa77b6857ca22832cffa5d5011a667207eb2a0ad56adb5db443b8a", [:mix], [{:decimal, "~> 1.0 or ~> 2.0", [hex: :decimal, repo: "hexpm", optional: true]}], "hexpm", "c5eb0cab91f094599f94d55bc63409236a8ec69a21a67814529e8d5f6cc90b3b"},
|
||||||
"quicer": {:hex, :quicer, "0.2.15", "775fcfa09c9ce5a4a00d23c1c8b1e81cd361ebc847e43bb766efac4373cdb1b8", [:rebar3], [{:snabbkaffe, "1.0.10", [hex: :snabbkaffe, repo: "hexpm", optional: false]}], "hexpm", "52236d9384a541341706bcae438e81b209fa21d2045d3df317f8ba5beda6a8e5"},
|
"quicer": {:hex, :quicer, "0.2.15", "775fcfa09c9ce5a4a00d23c1c8b1e81cd361ebc847e43bb766efac4373cdb1b8", [:rebar3], [{:snabbkaffe, "1.0.10", [hex: :snabbkaffe, repo: "hexpm", optional: false]}], "hexpm", "52236d9384a541341706bcae438e81b209fa21d2045d3df317f8ba5beda6a8e5"},
|
||||||
"snabbkaffe": {:hex, :snabbkaffe, "1.0.10", "9be2f54f61fc6862391b666b2b5b76c3fa53598e2989a17cef1b48cf347a8a63", [:rebar3], [], "hexpm", "70a98df36ae756908d55b5770891d443d63c903833e3e87d544036e13d4fac26"},
|
"snabbkaffe": {:hex, :snabbkaffe, "1.0.10", "9be2f54f61fc6862391b666b2b5b76c3fa53598e2989a17cef1b48cf347a8a63", [:rebar3], [], "hexpm", "70a98df36ae756908d55b5770891d443d63c903833e3e87d544036e13d4fac26"},
|
||||||
"telemetry": {:hex, :telemetry, "1.3.0", "fedebbae410d715cf8e7062c96a1ef32ec22e764197f70cda73d82778d61e7a2", [:rebar3], [], "hexpm", "7015fc8919dbe63764f4b4b87a95b7c0996bd539e0d499be6ec9d7f3875b79e6"},
|
"telemetry": {:hex, :telemetry, "1.3.0", "fedebbae410d715cf8e7062c96a1ef32ec22e764197f70cda73d82778d61e7a2", [:rebar3], [], "hexpm", "7015fc8919dbe63764f4b4b87a95b7c0996bd539e0d499be6ec9d7f3875b79e6"},
|
||||||
|
|||||||
Reference in New Issue
Block a user