fix
This commit is contained in:
parent
d9833e6340
commit
ca24499da4
@ -1,3 +1,3 @@
|
|||||||
FROM erlang:25.3
|
FROM erlang:25.3
|
||||||
|
|
||||||
CMD /data/power_gateway/bin/power_gateway foreground
|
CMD /data/light/bin/light foreground
|
||||||
@ -1,4 +1,4 @@
|
|||||||
power_gateway
|
light
|
||||||
=====
|
=====
|
||||||
|
|
||||||
An OTP application
|
An OTP application
|
||||||
|
|||||||
@ -389,7 +389,7 @@ unpack(<<PacketId:16, Type:8, Body/binary>>) ->
|
|||||||
%%%===================================================================
|
%%%===================================================================
|
||||||
|
|
||||||
handle_poll_command(#{<<"device_uuid">> := DeviceUUID, <<"command">> := <<"query_status">>}) when is_binary(DeviceUUID) ->
|
handle_poll_command(#{<<"device_uuid">> := DeviceUUID, <<"command">> := <<"query_status">>}) when is_binary(DeviceUUID) ->
|
||||||
case power_device:get_pid(DeviceUUID) of
|
case light_device:get_pid(DeviceUUID) of
|
||||||
undefined ->
|
undefined ->
|
||||||
#{
|
#{
|
||||||
<<"c">> => 1,
|
<<"c">> => 1,
|
||||||
@ -403,7 +403,7 @@ handle_poll_command(#{<<"device_uuid">> := DeviceUUID, <<"command">> := <<"query
|
|||||||
0 => <<"离线"/utf8>>,
|
0 => <<"离线"/utf8>>,
|
||||||
1 => <<"在线"/utf8>>
|
1 => <<"在线"/utf8>>
|
||||||
},
|
},
|
||||||
{ok, Status} = power_device:poll_status(Pid),
|
{ok, Status} = light_device:poll_status(Pid),
|
||||||
#{
|
#{
|
||||||
<<"c">> => 1,
|
<<"c">> => 1,
|
||||||
<<"r">> => #{
|
<<"r">> => #{
|
||||||
@ -415,12 +415,12 @@ handle_poll_command(#{<<"device_uuid">> := DeviceUUID, <<"command">> := <<"query
|
|||||||
|
|
||||||
-spec handle_param(Params :: map()) -> ok | {error, Reason :: binary()}.
|
-spec handle_param(Params :: map()) -> ok | {error, Reason :: binary()}.
|
||||||
handle_param(Params) when is_map(Params) ->
|
handle_param(Params) when is_map(Params) ->
|
||||||
power_gateway_args:push_param(Params),
|
light_args:push_param(Params),
|
||||||
ok.
|
ok.
|
||||||
|
|
||||||
-spec handle_metric(Metric :: list()) -> ok | {error, Reason :: binary()}.
|
-spec handle_metric(Metric :: list()) -> ok | {error, Reason :: binary()}.
|
||||||
handle_metric(Metric) when is_list(Metric) ->
|
handle_metric(Metric) when is_list(Metric) ->
|
||||||
power_gateway_args:push_metric(Metric),
|
light_args:push_metric(Metric),
|
||||||
ok.
|
ok.
|
||||||
|
|
||||||
-spec handle_stream_call(ServiceName :: binary(), Fields :: list(), Tag :: map()) ->
|
-spec handle_stream_call(ServiceName :: binary(), Fields :: list(), Tag :: map()) ->
|
||||||
@ -45,20 +45,20 @@
|
|||||||
|
|
||||||
-export_type([host/0 , option/0 , properties/0 , payload/0 , pubopt/0 , subopt/0 , mqtt_msg/0 , client/0]).
|
-export_type([host/0 , option/0 , properties/0 , payload/0 , pubopt/0 , subopt/0 , mqtt_msg/0 , client/0]).
|
||||||
|
|
||||||
-type(host() :: inet:ip_address() | inet:hostname()).
|
-type host() :: inet:ip_address() | inet:hostname().
|
||||||
|
|
||||||
%% Message handler is a set of callbacks defined to handle MQTT messages
|
%% Message handler is a set of callbacks defined to handle MQTT messages
|
||||||
%% as well as the disconnect event.
|
%% as well as the disconnect event.
|
||||||
-define(NO_MSG_HDLR, undefined).
|
-define(NO_MSG_HDLR, undefined).
|
||||||
|
|
||||||
-type(mfas() :: {module(), atom(), list()} | {function(), list()}).
|
-type mfas() :: {module(), atom(), list()} | {function(), list()}.
|
||||||
|
|
||||||
-type(msg_handler() :: #{puback := fun((_) -> any()) | mfas(),
|
-type msg_handler() :: #{puback := fun((_) -> any()) | mfas(),
|
||||||
publish := fun((emqx_types:message()) -> any()) | mfas(),
|
publish := fun((emqx_types:message()) -> any()) | mfas(),
|
||||||
disconnected := fun(({reason_code(), _Properties :: term()}) -> any()) | mfas()
|
disconnected := fun(({reason_code(), _Properties :: term()}) -> any()) | mfas()
|
||||||
}).
|
}.
|
||||||
|
|
||||||
-type(option() :: {name, atom()}
|
-type option() :: {name, atom()}
|
||||||
| {owner, pid()}
|
| {owner, pid()}
|
||||||
| {msg_handler, msg_handler()}
|
| {msg_handler, msg_handler()}
|
||||||
| {host, host()}
|
| {host, host()}
|
||||||
@ -86,34 +86,34 @@
|
|||||||
| {auto_ack, boolean()}
|
| {auto_ack, boolean()}
|
||||||
| {ack_timeout, pos_integer()}
|
| {ack_timeout, pos_integer()}
|
||||||
| {force_ping, boolean()}
|
| {force_ping, boolean()}
|
||||||
| {properties, properties()}).
|
| {properties, properties()}.
|
||||||
|
|
||||||
-type(maybe(T) :: undefined | T).
|
-type maybe_t(T) :: undefined | T.
|
||||||
-type(topic() :: binary()).
|
-type topic() :: binary().
|
||||||
-type(payload() :: iodata()).
|
-type payload() :: iodata().
|
||||||
-type(packet_id() :: 0..16#FFFF).
|
-type packet_id() :: 0..16#FFFF.
|
||||||
-type(reason_code() :: 0..16#FF).
|
-type reason_code() :: 0..16#FF.
|
||||||
-type(properties() :: #{atom() => term()}).
|
-type properties() :: #{atom() => term()}.
|
||||||
-type(version() :: ?MQTT_PROTO_V3
|
-type version() :: ?MQTT_PROTO_V3
|
||||||
| ?MQTT_PROTO_V4
|
| ?MQTT_PROTO_V4
|
||||||
| ?MQTT_PROTO_V5).
|
| ?MQTT_PROTO_V5.
|
||||||
-type(qos() :: ?QOS_0 | ?QOS_1 | ?QOS_2).
|
-type qos() :: ?QOS_0 | ?QOS_1 | ?QOS_2.
|
||||||
-type(qos_name() :: qos0 | at_most_once |
|
-type qos_name() :: qos0 | at_most_once |
|
||||||
qos1 | at_least_once |
|
qos1 | at_least_once |
|
||||||
qos2 | exactly_once).
|
qos2 | exactly_once.
|
||||||
-type(pubopt() :: {retain, boolean()}
|
-type pubopt() :: {retain, boolean()}
|
||||||
| {qos, qos() | qos_name()}).
|
| {qos, qos() | qos_name()}.
|
||||||
-type(subopt() :: {rh, 0 | 1 | 2}
|
-type subopt() :: {rh, 0 | 1 | 2}
|
||||||
| {rap, boolean()}
|
| {rap, boolean()}
|
||||||
| {nl, boolean()}
|
| {nl, boolean()}
|
||||||
| {qos, qos() | qos_name()}).
|
| {qos, qos() | qos_name()}.
|
||||||
|
|
||||||
-type(subscribe_ret() ::
|
-type subscribe_ret() ::
|
||||||
{ok, properties(), [reason_code()]} | {error, term()}).
|
{ok, properties(), [reason_code()]} | {error, term()}.
|
||||||
|
|
||||||
-type(client() :: pid() | atom()).
|
-type client() :: pid() | atom().
|
||||||
|
|
||||||
-opaque(mqtt_msg() :: #mqtt_msg{}).
|
-opaque mqtt_msg() :: #mqtt_msg{}.
|
||||||
|
|
||||||
-record(state, {
|
-record(state, {
|
||||||
name :: atom(),
|
name :: atom(),
|
||||||
@ -128,12 +128,12 @@
|
|||||||
bridge_mode :: boolean(),
|
bridge_mode :: boolean(),
|
||||||
clientid :: binary(),
|
clientid :: binary(),
|
||||||
clean_start :: boolean(),
|
clean_start :: boolean(),
|
||||||
username :: maybe(binary()),
|
username :: maybe_t(binary()),
|
||||||
password :: maybe(binary()),
|
password :: maybe_t(binary()),
|
||||||
proto_ver :: version(),
|
proto_ver :: version(),
|
||||||
proto_name :: iodata(),
|
proto_name :: iodata(),
|
||||||
keepalive :: non_neg_integer(),
|
keepalive :: non_neg_integer(),
|
||||||
keepalive_timer :: maybe(reference()),
|
keepalive_timer :: maybe_t(reference()),
|
||||||
force_ping :: boolean(),
|
force_ping :: boolean(),
|
||||||
paused :: boolean(),
|
paused :: boolean(),
|
||||||
will_flag :: boolean(),
|
will_flag :: boolean(),
|
||||||
@ -24,22 +24,22 @@
|
|||||||
|
|
||||||
-export_type([options/0, parse_state/0, parse_result/0, serialize_fun/0]).
|
-export_type([options/0, parse_state/0, parse_result/0, serialize_fun/0]).
|
||||||
|
|
||||||
-type(version() :: ?MQTT_PROTO_V3
|
-type version() :: ?MQTT_PROTO_V3
|
||||||
| ?MQTT_PROTO_V4
|
| ?MQTT_PROTO_V4
|
||||||
| ?MQTT_PROTO_V5).
|
| ?MQTT_PROTO_V5.
|
||||||
|
|
||||||
-type(options() :: #{strict_mode => boolean(),
|
-type options() :: #{strict_mode => boolean(),
|
||||||
max_size => 1..?MAX_PACKET_SIZE,
|
max_size => 1..?MAX_PACKET_SIZE,
|
||||||
version => version()}).
|
version => version()}.
|
||||||
|
|
||||||
-opaque(parse_state() :: {none, options()} | cont_fun()).
|
-opaque parse_state() :: {none, options()} | cont_fun().
|
||||||
|
|
||||||
-opaque(parse_result() :: {more, cont_fun()}
|
-opaque parse_result() :: {more, cont_fun()}
|
||||||
| {ok, #mqtt_packet{}, binary(), parse_state()}).
|
| {ok, #mqtt_packet{}, binary(), parse_state()}.
|
||||||
|
|
||||||
-type(cont_fun() :: fun((binary()) -> parse_result())).
|
-type cont_fun() :: fun((binary()) -> parse_result()).
|
||||||
|
|
||||||
-type(serialize_fun() :: fun((emqx_types:packet()) -> iodata())).
|
-type serialize_fun() :: fun((emqx_types:packet()) -> iodata()).
|
||||||
|
|
||||||
-define(none(Options), {none, Options}).
|
-define(none(Options), {none, Options}).
|
||||||
|
|
||||||
@ -24,8 +24,8 @@
|
|||||||
%% For tests
|
%% For tests
|
||||||
-export([all/0]).
|
-export([all/0]).
|
||||||
|
|
||||||
-type(prop_name() :: atom()).
|
-type prop_name() :: atom().
|
||||||
-type(prop_id() :: pos_integer()).
|
-type prop_id() :: pos_integer().
|
||||||
|
|
||||||
-define(PROPS_TABLE,
|
-define(PROPS_TABLE,
|
||||||
#{16#01 => {'Payload-Format-Indicator', 'Byte', [?PUBLISH]},
|
#{16#01 => {'Payload-Format-Indicator', 'Byte', [?PUBLISH]},
|
||||||
@ -25,11 +25,11 @@
|
|||||||
ssl
|
ssl
|
||||||
}).
|
}).
|
||||||
|
|
||||||
-type(socket() :: inet:socket() | #ssl_socket{}).
|
-type socket() :: inet:socket() | #ssl_socket{}.
|
||||||
|
|
||||||
-type(sockname() :: {inet:ip_address(), inet:port_number()}).
|
-type sockname() :: {inet:ip_address(), inet:port_number()}.
|
||||||
|
|
||||||
-type(option() :: gen_tcp:connect_option() | {ssl_opts, [ssl:ssl_option()]}).
|
-type option() :: gen_tcp:connect_option() | {ssl_opts, [ssl:ssl_option()]}.
|
||||||
|
|
||||||
-export_type([socket/0, option/0]).
|
-export_type([socket/0, option/0]).
|
||||||
|
|
||||||
@ -1,8 +1,8 @@
|
|||||||
{application, power_gateway,
|
{application, light,
|
||||||
[{description, "An OTP application"},
|
[{description, "An OTP application"},
|
||||||
{vsn, "1.0"},
|
{vsn, "1.0"},
|
||||||
{registered, []},
|
{registered, []},
|
||||||
{mod, {power_gateway_app, []}},
|
{mod, {light_app, []}},
|
||||||
{applications,
|
{applications,
|
||||||
[
|
[
|
||||||
hackney,
|
hackney,
|
||||||
@ -1,16 +1,16 @@
|
|||||||
%%%-------------------------------------------------------------------
|
%%%-------------------------------------------------------------------
|
||||||
%% @doc power_gateway public API
|
%% @doc light public API
|
||||||
%% @end
|
%% @end
|
||||||
%%%-------------------------------------------------------------------
|
%%%-------------------------------------------------------------------
|
||||||
|
|
||||||
-module(power_gateway_app).
|
-module(light_app).
|
||||||
|
|
||||||
-behaviour(application).
|
-behaviour(application).
|
||||||
|
|
||||||
-export([start/2, stop/1]).
|
-export([start/2, stop/1]).
|
||||||
|
|
||||||
start(_StartType, _StartArgs) ->
|
start(_StartType, _StartArgs) ->
|
||||||
SupResult = power_gateway_sup:start_link(),
|
SupResult = light_sup:start_link(),
|
||||||
start_http_server(),
|
start_http_server(),
|
||||||
SupResult.
|
SupResult.
|
||||||
|
|
||||||
@ -20,7 +20,7 @@ stop(_State) ->
|
|||||||
%% internal functions
|
%% internal functions
|
||||||
|
|
||||||
start_http_server() ->
|
start_http_server() ->
|
||||||
{ok, Props} = application:get_env(power_gateway, http_server),
|
{ok, Props} = application:get_env(light, http_server),
|
||||||
Acceptors = proplists:get_value(acceptors, Props, 50),
|
Acceptors = proplists:get_value(acceptors, Props, 50),
|
||||||
MaxConnections = proplists:get_value(max_connections, Props, 10240),
|
MaxConnections = proplists:get_value(max_connections, Props, 10240),
|
||||||
Backlog = proplists:get_value(backlog, Props, 1024),
|
Backlog = proplists:get_value(backlog, Props, 1024),
|
||||||
@ -41,4 +41,4 @@ start_http_server() ->
|
|||||||
|
|
||||||
{ok, Pid} = cowboy:start_clear(http_listener, TransOpts, #{env => #{dispatch => Dispatcher}}),
|
{ok, Pid} = cowboy:start_clear(http_listener, TransOpts, #{env => #{dispatch => Dispatcher}}),
|
||||||
|
|
||||||
lager:debug("[iot_app] the http server start at: ~p, pid is: ~p", [Port, Pid]).
|
lager:debug("[light_app] the http server start at: ~p, pid is: ~p", [Port, Pid]).
|
||||||
@ -6,7 +6,7 @@
|
|||||||
%%% @end
|
%%% @end
|
||||||
%%% Created : 06. 9月 2023 16:37
|
%%% Created : 06. 9月 2023 16:37
|
||||||
%%%-------------------------------------------------------------------
|
%%%-------------------------------------------------------------------
|
||||||
-module(power_gateway_args).
|
-module(light_args).
|
||||||
-author("aresei").
|
-author("aresei").
|
||||||
|
|
||||||
-behaviour(gen_server).
|
-behaviour(gen_server).
|
||||||
@ -69,13 +69,13 @@ init([]) ->
|
|||||||
{ok, Metrics} = efka_client:request_metric(),
|
{ok, Metrics} = efka_client:request_metric(),
|
||||||
try convert_metric(Metrics) of
|
try convert_metric(Metrics) of
|
||||||
{ok, MetricMap} ->
|
{ok, MetricMap} ->
|
||||||
lager:debug("[power_gateway_args] init load metric_map: ~p", [MetricMap]),
|
lager:debug("[light_args] init load metric_map: ~p", [MetricMap]),
|
||||||
{ok, Param} = efka_client:request_param(),
|
{ok, Param} = efka_client:request_param(),
|
||||||
|
|
||||||
{ok, #state{metrics = MetricMap, param = Param}}
|
{ok, #state{metrics = MetricMap, param = Param}}
|
||||||
catch
|
catch
|
||||||
_:Error:Stack->
|
_:Error:Stack->
|
||||||
lager:warning("[power_gateway_args] request_metric get error: ~p, stack: ~p", [Error, Stack]),
|
lager:warning("[light_args] request_metric get error: ~p, stack: ~p", [Error, Stack]),
|
||||||
{ok, #state{metrics = #{}, param = #{}}}
|
{ok, #state{metrics = #{}, param = #{}}}
|
||||||
end.
|
end.
|
||||||
|
|
||||||
@ -111,7 +111,7 @@ handle_cast({push_param, Param}, State = #state{}) ->
|
|||||||
handle_cast({push_metric, Metrics}, State = #state{}) ->
|
handle_cast({push_metric, Metrics}, State = #state{}) ->
|
||||||
try convert_metric(Metrics) of
|
try convert_metric(Metrics) of
|
||||||
{ok, MetricMap} ->
|
{ok, MetricMap} ->
|
||||||
lager:debug("[power_gateway_args] push metric_map: ~p", [MetricMap]),
|
lager:debug("[light_args] push metric_map: ~p", [MetricMap]),
|
||||||
{noreply, State#state{metrics = MetricMap}}
|
{noreply, State#state{metrics = MetricMap}}
|
||||||
catch _:_ ->
|
catch _:_ ->
|
||||||
{noreply, State}
|
{noreply, State}
|
||||||
@ -6,7 +6,7 @@
|
|||||||
%%% @end
|
%%% @end
|
||||||
%%% Created : 04. 9月 2024 16:13
|
%%% Created : 04. 9月 2024 16:13
|
||||||
%%%-------------------------------------------------------------------
|
%%%-------------------------------------------------------------------
|
||||||
-module(power_device).
|
-module(light_device).
|
||||||
-author("anlicheng").
|
-author("anlicheng").
|
||||||
|
|
||||||
-behaviour(gen_server).
|
-behaviour(gen_server).
|
||||||
@ -40,7 +40,7 @@ get_pid(DeviceUUID) when is_binary(DeviceUUID) ->
|
|||||||
|
|
||||||
-spec get_name(DeviceUUID :: binary()) -> atom().
|
-spec get_name(DeviceUUID :: binary()) -> atom().
|
||||||
get_name(DeviceUUID) when is_binary(DeviceUUID) ->
|
get_name(DeviceUUID) when is_binary(DeviceUUID) ->
|
||||||
binary_to_atom(<<"power_device:", DeviceUUID/binary>>).
|
binary_to_atom(<<"light_device:", DeviceUUID/binary>>).
|
||||||
|
|
||||||
-spec metric_data(Pid :: pid(), Message :: binary()) -> no_return().
|
-spec metric_data(Pid :: pid(), Message :: binary()) -> no_return().
|
||||||
metric_data(Pid, Message) when is_pid(Pid), is_binary(Message) ->
|
metric_data(Pid, Message) when is_pid(Pid), is_binary(Message) ->
|
||||||
@ -67,7 +67,7 @@ start_link(Name, DeviceUUID) when is_atom(Name), is_binary(DeviceUUID) ->
|
|||||||
{ok, State :: #state{}} | {ok, State :: #state{}, timeout() | hibernate} |
|
{ok, State :: #state{}} | {ok, State :: #state{}, timeout() | hibernate} |
|
||||||
{stop, Reason :: term()} | ignore).
|
{stop, Reason :: term()} | ignore).
|
||||||
init([DeviceUUID]) ->
|
init([DeviceUUID]) ->
|
||||||
{ok, HeartbeatTicker0} = application:get_env(power_gateway, heartbeat_ticker),
|
{ok, HeartbeatTicker0} = application:get_env(light, heartbeat_ticker),
|
||||||
HeartbeatTicker = HeartbeatTicker0 * 1000,
|
HeartbeatTicker = HeartbeatTicker0 * 1000,
|
||||||
|
|
||||||
erlang:start_timer(HeartbeatTicker, self(), heartbeat_ticker),
|
erlang:start_timer(HeartbeatTicker, self(), heartbeat_ticker),
|
||||||
@ -101,9 +101,9 @@ handle_cast({metric_data, Message}, State = #state{device_uuid = DeviceUUID, dat
|
|||||||
Info = iolist_to_binary(jiffy:encode(Props, [force_utf8])),
|
Info = iolist_to_binary(jiffy:encode(Props, [force_utf8])),
|
||||||
case catch efka_client:send_metric_data(Props, #{}) of
|
case catch efka_client:send_metric_data(Props, #{}) of
|
||||||
{ok, _} ->
|
{ok, _} ->
|
||||||
power_logger:write([<<"OK">>, Info]);
|
light_logger:write([<<"OK">>, Info]);
|
||||||
_ ->
|
_ ->
|
||||||
power_logger:write([<<"ERROR">>, Info])
|
light_logger:write([<<"ERROR">>, Info])
|
||||||
end,
|
end,
|
||||||
|
|
||||||
%% 如果设备当前是离线状态,则需要发送上线消息
|
%% 如果设备当前是离线状态,则需要发送上线消息
|
||||||
@ -111,9 +111,9 @@ handle_cast({metric_data, Message}, State = #state{device_uuid = DeviceUUID, dat
|
|||||||
|
|
||||||
{noreply, State#state{data_counter = DataCounter + 1, status = 1}};
|
{noreply, State#state{data_counter = DataCounter + 1, status = 1}};
|
||||||
M when is_map(M) ->
|
M when is_map(M) ->
|
||||||
lager:notice("[power_device] invalid map: ~p", [M]);
|
lager:notice("[light_device] invalid map: ~p", [M]);
|
||||||
Error ->
|
Error ->
|
||||||
lager:notice("[power_device] jiffy decode error: ~p", [Error]),
|
lager:notice("[light_device] jiffy decode error: ~p", [Error]),
|
||||||
{noreply, State}
|
{noreply, State}
|
||||||
end.
|
end.
|
||||||
|
|
||||||
@ -6,7 +6,7 @@
|
|||||||
%%% @end
|
%%% @end
|
||||||
%%% Created : 04. 9月 2024 15:39
|
%%% Created : 04. 9月 2024 15:39
|
||||||
%%%-------------------------------------------------------------------
|
%%%-------------------------------------------------------------------
|
||||||
-module(power_device_sup).
|
-module(light_device_sup).
|
||||||
-author("anlicheng").
|
-author("anlicheng").
|
||||||
|
|
||||||
-behaviour(supervisor).
|
-behaviour(supervisor).
|
||||||
@ -49,7 +49,7 @@ init([]) ->
|
|||||||
|
|
||||||
-spec ensure_device_started(DeviceUUID :: binary()) -> {ok, Pid :: pid()} | {error, Reason :: any()}.
|
-spec ensure_device_started(DeviceUUID :: binary()) -> {ok, Pid :: pid()} | {error, Reason :: any()}.
|
||||||
ensure_device_started(DeviceUUID) when is_binary(DeviceUUID) ->
|
ensure_device_started(DeviceUUID) when is_binary(DeviceUUID) ->
|
||||||
case power_device:get_pid(DeviceUUID) of
|
case light_device:get_pid(DeviceUUID) of
|
||||||
DevicePid when is_pid(DevicePid) ->
|
DevicePid when is_pid(DevicePid) ->
|
||||||
{ok, DevicePid};
|
{ok, DevicePid};
|
||||||
undefined ->
|
undefined ->
|
||||||
@ -64,17 +64,17 @@ ensure_device_started(DeviceUUID) when is_binary(DeviceUUID) ->
|
|||||||
end.
|
end.
|
||||||
|
|
||||||
delete_device(DeviceUUID) when is_binary(DeviceUUID) ->
|
delete_device(DeviceUUID) when is_binary(DeviceUUID) ->
|
||||||
Id = power_device:get_name(DeviceUUID),
|
Id = light_device:get_name(DeviceUUID),
|
||||||
ok = supervisor:terminate_child(?MODULE, Id),
|
ok = supervisor:terminate_child(?MODULE, Id),
|
||||||
supervisor:delete_child(?MODULE, Id).
|
supervisor:delete_child(?MODULE, Id).
|
||||||
|
|
||||||
child_spec(DeviceUUID) when is_binary(DeviceUUID) ->
|
child_spec(DeviceUUID) when is_binary(DeviceUUID) ->
|
||||||
Name = power_device:get_name(DeviceUUID),
|
Name = light_device:get_name(DeviceUUID),
|
||||||
#{
|
#{
|
||||||
id => Name,
|
id => Name,
|
||||||
start => {power_device, start_link, [Name, DeviceUUID]},
|
start => {light_device, start_link, [Name, DeviceUUID]},
|
||||||
restart => permanent,
|
restart => permanent,
|
||||||
shutdown => 2000,
|
shutdown => 2000,
|
||||||
type => worker,
|
type => worker,
|
||||||
modules => ['power_device']
|
modules => ['light_device']
|
||||||
}.
|
}.
|
||||||
@ -6,7 +6,7 @@
|
|||||||
%%% @end
|
%%% @end
|
||||||
%%% Created : 07. 9月 2023 17:07
|
%%% Created : 07. 9月 2023 17:07
|
||||||
%%%-------------------------------------------------------------------
|
%%%-------------------------------------------------------------------
|
||||||
-module(power_logger).
|
-module(light_logger).
|
||||||
-author("aresei").
|
-author("aresei").
|
||||||
|
|
||||||
-behaviour(gen_server).
|
-behaviour(gen_server).
|
||||||
@ -7,7 +7,7 @@
|
|||||||
%%% @end
|
%%% @end
|
||||||
%%% Created : 12. 3月 2023 21:27
|
%%% Created : 12. 3月 2023 21:27
|
||||||
%%%-------------------------------------------------------------------
|
%%%-------------------------------------------------------------------
|
||||||
-module(power_gateway_mqtt_subscriber).
|
-module(light_mqtt_subscriber).
|
||||||
-author("aresei").
|
-author("aresei").
|
||||||
|
|
||||||
-behaviour(gen_server).
|
-behaviour(gen_server).
|
||||||
@ -50,7 +50,7 @@ start_link() ->
|
|||||||
{stop, Reason :: term()} | ignore).
|
{stop, Reason :: term()} | ignore).
|
||||||
init([]) ->
|
init([]) ->
|
||||||
%% 建立到emqx服务器的连接
|
%% 建立到emqx服务器的连接
|
||||||
Opts = emqx_opts(<<"power-subscriber">>),
|
Opts = emqx_opts(<<"light-subscriber">>),
|
||||||
lager:debug("[opts] is: ~p", [Opts]),
|
lager:debug("[opts] is: ~p", [Opts]),
|
||||||
case emqtt:start_link(Opts) of
|
case emqtt:start_link(Opts) of
|
||||||
{ok, ConnPid} ->
|
{ok, ConnPid} ->
|
||||||
@ -135,10 +135,10 @@ terminate(Reason, _State = #state{conn_pid = ConnPid}) when is_pid(ConnPid) ->
|
|||||||
{ok, _Props, _ReasonCode} = emqtt:unsubscribe(ConnPid, #{}, TopicNames),
|
{ok, _Props, _ReasonCode} = emqtt:unsubscribe(ConnPid, #{}, TopicNames),
|
||||||
|
|
||||||
ok = emqtt:disconnect(ConnPid),
|
ok = emqtt:disconnect(ConnPid),
|
||||||
lager:debug("[iot_mqtt_subscriber] terminate with reason: ~p", [Reason]),
|
lager:debug("[light_mqtt_subscriber] terminate with reason: ~p", [Reason]),
|
||||||
ok;
|
ok;
|
||||||
terminate(Reason, _State) ->
|
terminate(Reason, _State) ->
|
||||||
lager:debug("[iot_mqtt_subscriber] terminate with reason: ~p", [Reason]),
|
lager:debug("[light_mqtt_subscriber] terminate with reason: ~p", [Reason]),
|
||||||
ok.
|
ok.
|
||||||
|
|
||||||
%% @private
|
%% @private
|
||||||
@ -156,7 +156,7 @@ code_change(_OldVsn, State = #state{}, _Extra) ->
|
|||||||
-spec emqx_opts(ClientSuffix :: binary()) -> list().
|
-spec emqx_opts(ClientSuffix :: binary()) -> list().
|
||||||
emqx_opts(ClientSuffix) when is_binary(ClientSuffix) ->
|
emqx_opts(ClientSuffix) when is_binary(ClientSuffix) ->
|
||||||
%% 建立到emqx服务器的连接
|
%% 建立到emqx服务器的连接
|
||||||
{ok, Props} = application:get_env(power_gateway, emqx_server),
|
{ok, Props} = application:get_env(light, emqx_server),
|
||||||
EMQXHost = proplists:get_value(host, Props),
|
EMQXHost = proplists:get_value(host, Props),
|
||||||
EMQXPort = proplists:get_value(port, Props, 1883),
|
EMQXPort = proplists:get_value(port, Props, 1883),
|
||||||
Username = proplists:get_value(username, Props),
|
Username = proplists:get_value(username, Props),
|
||||||
@ -182,13 +182,13 @@ emqx_opts(ClientSuffix) when is_binary(ClientSuffix) ->
|
|||||||
|
|
||||||
-spec dispatch(DeviceMac :: binary(), Message :: binary()) -> no_return().
|
-spec dispatch(DeviceMac :: binary(), Message :: binary()) -> no_return().
|
||||||
dispatch(DeviceMac, Message) when is_binary(DeviceMac), is_binary(Message) ->
|
dispatch(DeviceMac, Message) when is_binary(DeviceMac), is_binary(Message) ->
|
||||||
case power_gateway_args:get_device_uuid(DeviceMac) of
|
case light_args:get_device_uuid(DeviceMac) of
|
||||||
error ->
|
error ->
|
||||||
lager:notice("[mqtt_subscriber] device_mac: ~p, device_uuid not found", [DeviceMac]);
|
lager:notice("[mqtt_subscriber] device_mac: ~p, device_uuid not found", [DeviceMac]);
|
||||||
{ok, DeviceUUID} ->
|
{ok, DeviceUUID} ->
|
||||||
case power_device_sup:ensure_device_started(DeviceUUID) of
|
case light_device_sup:ensure_device_started(DeviceUUID) of
|
||||||
{ok, DevicePid} ->
|
{ok, DevicePid} ->
|
||||||
power_device:metric_data(DevicePid, Message);
|
light_device:metric_data(DevicePid, Message);
|
||||||
{error, Reason} ->
|
{error, Reason} ->
|
||||||
lager:notice("[mqtt_subscriber] start device get error: ~p", [Reason])
|
lager:notice("[mqtt_subscriber] start device get error: ~p", [Reason])
|
||||||
end
|
end
|
||||||
@ -1,9 +1,9 @@
|
|||||||
%%%-------------------------------------------------------------------
|
%%%-------------------------------------------------------------------
|
||||||
%% @doc power_gateway top level supervisor.
|
%% @doc light top level supervisor.
|
||||||
%% @end
|
%% @end
|
||||||
%%%-------------------------------------------------------------------
|
%%%-------------------------------------------------------------------
|
||||||
|
|
||||||
-module(power_gateway_sup).
|
-module(light_sup).
|
||||||
|
|
||||||
-behaviour(supervisor).
|
-behaviour(supervisor).
|
||||||
|
|
||||||
@ -30,7 +30,7 @@ init([]) ->
|
|||||||
intensity => 1000,
|
intensity => 1000,
|
||||||
period => 3600},
|
period => 3600},
|
||||||
|
|
||||||
{ok, EfkaServer} = application:get_env(power_gateway, efka_server),
|
{ok, EfkaServer} = application:get_env(light, efka_server),
|
||||||
Host = proplists:get_value(host, EfkaServer),
|
Host = proplists:get_value(host, EfkaServer),
|
||||||
Port = proplists:get_value(port, EfkaServer),
|
Port = proplists:get_value(port, EfkaServer),
|
||||||
|
|
||||||
@ -48,39 +48,39 @@ init([]) ->
|
|||||||
},
|
},
|
||||||
|
|
||||||
#{
|
#{
|
||||||
id => 'power_logger',
|
id => 'light_logger',
|
||||||
start => {'power_logger', start_link, ["device_data"]},
|
start => {'light_logger', start_link, ["device_data"]},
|
||||||
restart => permanent,
|
restart => permanent,
|
||||||
shutdown => 2000,
|
shutdown => 2000,
|
||||||
type => worker,
|
type => worker,
|
||||||
modules => ['power_logger']
|
modules => ['light_logger']
|
||||||
},
|
},
|
||||||
|
|
||||||
#{
|
#{
|
||||||
id => 'power_gateway_args',
|
id => 'light_args',
|
||||||
start => {'power_gateway_args', start_link, []},
|
start => {'light_args', start_link, []},
|
||||||
restart => permanent,
|
restart => permanent,
|
||||||
shutdown => 2000,
|
shutdown => 2000,
|
||||||
type => worker,
|
type => worker,
|
||||||
modules => ['power_gateway_args']
|
modules => ['light_args']
|
||||||
},
|
},
|
||||||
|
|
||||||
#{
|
#{
|
||||||
id => 'power_device_sup',
|
id => 'light_device_sup',
|
||||||
start => {'power_device_sup', start_link, []},
|
start => {'light_device_sup', start_link, []},
|
||||||
restart => permanent,
|
restart => permanent,
|
||||||
shutdown => 2000,
|
shutdown => 2000,
|
||||||
type => worker,
|
type => worker,
|
||||||
modules => ['power_device_sup']
|
modules => ['light_device_sup']
|
||||||
},
|
},
|
||||||
|
|
||||||
#{
|
#{
|
||||||
id => 'power_gateway_mqtt_subscriber',
|
id => 'light_mqtt_subscriber',
|
||||||
start => {'power_gateway_mqtt_subscriber', start_link, []},
|
start => {'light_mqtt_subscriber', start_link, []},
|
||||||
restart => permanent,
|
restart => permanent,
|
||||||
shutdown => 2000,
|
shutdown => 2000,
|
||||||
type => worker,
|
type => worker,
|
||||||
modules => ['power_gateway_mqtt_subscriber']
|
modules => ['light_mqtt_subscriber']
|
||||||
}
|
}
|
||||||
],
|
],
|
||||||
|
|
||||||
@ -98,6 +98,6 @@ read_service_name() ->
|
|||||||
{ok, RegisterName0} ->
|
{ok, RegisterName0} ->
|
||||||
string:trim(RegisterName0);
|
string:trim(RegisterName0);
|
||||||
{error, Reason} ->
|
{error, Reason} ->
|
||||||
lager:warning("[power_app] read .version file get error: ~p", [Reason]),
|
lager:warning("[light_app] read .version file get error: ~p", [Reason]),
|
||||||
<<"power_gateway">>
|
<<"light">>
|
||||||
end.
|
end.
|
||||||
2
boot.sh
2
boot.sh
@ -1,3 +1,3 @@
|
|||||||
#!/bin/bash
|
#!/bin/bash
|
||||||
|
|
||||||
docker run -e TZ=Asia/Shanghai --hostname=power_gateway --net=host --restart=always -v /data/docker/power_gateway/:/data/power_gateway/ power_gateway:latest
|
docker run -e TZ=Asia/Shanghai --hostname=light --net=host --restart=always -v /data/docker/light/:/data/light/ light:latest
|
||||||
2
build.sh
2
build.sh
@ -1,4 +1,4 @@
|
|||||||
#!/bin/bash
|
#!/bin/bash
|
||||||
rebar3 compile && rebar3 release && rebar3 tar
|
rebar3 compile && rebar3 release && rebar3 tar
|
||||||
|
|
||||||
tar -czvf power_gateway.tgz boot.sh stop.sh Dockerfile .args.yaml
|
tar -czvf light.tgz boot.sh stop.sh Dockerfile .args.yaml
|
||||||
@ -1,5 +1,5 @@
|
|||||||
[
|
[
|
||||||
{power_gateway, [
|
{light, [
|
||||||
|
|
||||||
{http_server, [
|
{http_server, [
|
||||||
{port, 18082},
|
{port, 18082},
|
||||||
@ -31,7 +31,7 @@
|
|||||||
{lager, [
|
{lager, [
|
||||||
{colored, true},
|
{colored, true},
|
||||||
%% Whether to write a crash log, and where. Undefined means no crash logger.
|
%% Whether to write a crash log, and where. Undefined means no crash logger.
|
||||||
{crash_log, "trade_hub.crash.log"},
|
{crash_log, "light.crash.log"},
|
||||||
%% Maximum size in bytes of events in the crash log - defaults to 65536
|
%% Maximum size in bytes of events in the crash log - defaults to 65536
|
||||||
{crash_log_msg_size, 65536},
|
{crash_log_msg_size, 65536},
|
||||||
%% Maximum size of the crash log in bytes, before its rotated, set
|
%% Maximum size of the crash log in bytes, before its rotated, set
|
||||||
|
|||||||
@ -1,5 +1,5 @@
|
|||||||
[
|
[
|
||||||
{power_gateway, [
|
{light, [
|
||||||
|
|
||||||
{http_server, [
|
{http_server, [
|
||||||
{port, 18082},
|
{port, 18082},
|
||||||
@ -15,7 +15,7 @@
|
|||||||
{host, "172.30.37.212"},
|
{host, "172.30.37.212"},
|
||||||
{port, 1883},
|
{port, 1883},
|
||||||
{tcp_opts, []},
|
{tcp_opts, []},
|
||||||
{username, "power_gateway"},
|
{username, "light"},
|
||||||
{password, "PqRsTuVwXyZ!@#xy"},
|
{password, "PqRsTuVwXyZ!@#xy"},
|
||||||
{keepalive, 86400},
|
{keepalive, 86400},
|
||||||
{retry_interval, 5}
|
{retry_interval, 5}
|
||||||
@ -31,7 +31,7 @@
|
|||||||
{lager, [
|
{lager, [
|
||||||
{colored, true},
|
{colored, true},
|
||||||
%% Whether to write a crash log, and where. Undefined means no crash logger.
|
%% Whether to write a crash log, and where. Undefined means no crash logger.
|
||||||
{crash_log, "trade_hub.crash.log"},
|
{crash_log, "light.crash.log"},
|
||||||
%% Maximum size in bytes of events in the crash log - defaults to 65536
|
%% Maximum size in bytes of events in the crash log - defaults to 65536
|
||||||
{crash_log_msg_size, 65536},
|
{crash_log_msg_size, 65536},
|
||||||
%% Maximum size of the crash log in bytes, before its rotated, set
|
%% Maximum size of the crash log in bytes, before its rotated, set
|
||||||
|
|||||||
@ -1,6 +1,6 @@
|
|||||||
-sname power_gateway
|
-sname light
|
||||||
|
|
||||||
-setcookie power_gateway_cookie
|
-setcookie light_cookie
|
||||||
|
|
||||||
+K true
|
+K true
|
||||||
+A30
|
+A30
|
||||||
|
|||||||
@ -10,8 +10,8 @@
|
|||||||
{lager, ".*", {git,"https://github.com/erlang-lager/lager.git", {tag, "3.9.2"}}}
|
{lager, ".*", {git,"https://github.com/erlang-lager/lager.git", {tag, "3.9.2"}}}
|
||||||
]}.
|
]}.
|
||||||
|
|
||||||
{relx, [{release, {power_gateway, "1.0"},
|
{relx, [{release, {light, "1.0"},
|
||||||
[power_gateway,
|
[light,
|
||||||
sasl]},
|
sasl]},
|
||||||
|
|
||||||
{mode, dev},
|
{mode, dev},
|
||||||
|
|||||||
2
run
2
run
@ -2,4 +2,4 @@
|
|||||||
rebar3 compile
|
rebar3 compile
|
||||||
rebar3 release
|
rebar3 release
|
||||||
|
|
||||||
_build/default/rel/power_gateway/bin/power_gateway console
|
_build/default/rel/light/bin/light console
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user