fix efka_client
This commit is contained in:
parent
e7e3746486
commit
4bc88e5022
@ -54,7 +54,7 @@ deploy(TaskId, ContainerDir, Params = #'ContainerDeployParams'{
|
|||||||
case docker_commands:check_container_exist(ContainerName) of
|
case docker_commands:check_container_exist(ContainerName) of
|
||||||
true ->
|
true ->
|
||||||
trace_log(TaskId, <<"info">>, <<"本地容器已经存在:"/utf8, ContainerName/binary>>),
|
trace_log(TaskId, <<"info">>, <<"本地容器已经存在:"/utf8, ContainerName/binary>>),
|
||||||
efka_remote_agent:close_task_event_stream(TaskId, ?TASK_FAIL);
|
efka_client:close_task_event_stream(TaskId, ?TASK_FAIL);
|
||||||
false ->
|
false ->
|
||||||
Image = normalize_image(Image0),
|
Image = normalize_image(Image0),
|
||||||
|
|
||||||
@ -92,15 +92,15 @@ deploy(TaskId, ContainerDir, Params = #'ContainerDeployParams'{
|
|||||||
ShortContainerId = binary:part(ContainerId, 1, 12),
|
ShortContainerId = binary:part(ContainerId, 1, 12),
|
||||||
trace_log(TaskId, <<"info">>, <<"容器创建成功: "/utf8, ShortContainerId/binary>>),
|
trace_log(TaskId, <<"info">>, <<"容器创建成功: "/utf8, ShortContainerId/binary>>),
|
||||||
trace_log(TaskId, <<"info">>, <<"任务完成"/utf8>>),
|
trace_log(TaskId, <<"info">>, <<"任务完成"/utf8>>),
|
||||||
efka_remote_agent:close_task_event_stream(TaskId, ?TASK_SUCCESS);
|
efka_client:close_task_event_stream(TaskId, ?TASK_SUCCESS);
|
||||||
{error, Reason} ->
|
{error, Reason} ->
|
||||||
trace_log(TaskId, <<"error">>, <<"容器创建失败: "/utf8, Reason/binary>>),
|
trace_log(TaskId, <<"error">>, <<"容器创建失败: "/utf8, Reason/binary>>),
|
||||||
trace_log(TaskId, <<"error">>, <<"任务失败"/utf8>>),
|
trace_log(TaskId, <<"error">>, <<"任务失败"/utf8>>),
|
||||||
efka_remote_agent:close_task_event_stream(TaskId, ?TASK_FAIL)
|
efka_client:close_task_event_stream(TaskId, ?TASK_FAIL)
|
||||||
end;
|
end;
|
||||||
{error, Reason} ->
|
{error, Reason} ->
|
||||||
trace_log(TaskId, <<"error">>, <<"镜像拉取失败: "/utf8, Reason/binary>>),
|
trace_log(TaskId, <<"error">>, <<"镜像拉取失败: "/utf8, Reason/binary>>),
|
||||||
efka_remote_agent:close_task_event_stream(TaskId, ?TASK_FAIL)
|
efka_client:close_task_event_stream(TaskId, ?TASK_FAIL)
|
||||||
end
|
end
|
||||||
end.
|
end.
|
||||||
|
|
||||||
@ -116,6 +116,6 @@ normalize_image(Image) when is_binary(Image) ->
|
|||||||
|
|
||||||
-spec trace_log(TaskId :: integer(), Level :: binary(), Msg :: binary()) -> no_return().
|
-spec trace_log(TaskId :: integer(), Level :: binary(), Msg :: binary()) -> no_return().
|
||||||
trace_log(TaskId, Level, Msg) when is_integer(TaskId), is_binary(Level), is_binary(Msg) ->
|
trace_log(TaskId, Level, Msg) when is_integer(TaskId), is_binary(Level), is_binary(Msg) ->
|
||||||
efka_remote_agent:task_event_stream(TaskId, Level, Msg),
|
efka_client:task_event_stream(TaskId, Level, Msg),
|
||||||
Info = iolist_to_binary([<<"task_id=">>, integer_to_binary(TaskId), <<" ">>, Level, <<" ">>, Msg]),
|
Info = iolist_to_binary([<<"task_id=">>, integer_to_binary(TaskId), <<" ">>, Level, <<" ">>, Msg]),
|
||||||
efka_logger:write(Info).
|
efka_logger:write(Info).
|
||||||
|
|||||||
@ -208,8 +208,8 @@ handle_info({'DOWN', _Ref, process, TaskPid, Reason}, State = #state{task_map =
|
|||||||
ok;
|
ok;
|
||||||
Error0 ->
|
Error0 ->
|
||||||
Error = iolist_to_binary(io_lib:format("~p", [Error0])),
|
Error = iolist_to_binary(io_lib:format("~p", [Error0])),
|
||||||
efka_remote_agent:task_event_stream(TaskId, <<"error">>, <<"任务失败: "/utf8, Error/binary>>),
|
efka_client:task_event_stream(TaskId, <<"error">>, <<"任务失败: "/utf8, Error/binary>>),
|
||||||
efka_remote_agent:close_task_event_stream(TaskId, <<"task exited">>),
|
efka_client:close_task_event_stream(TaskId, <<"task exited">>),
|
||||||
logger:notice("[docker_manager] task_id: ~p, exit with error: ~p", [TaskId, Error]),
|
logger:notice("[docker_manager] task_id: ~p, exit with error: ~p", [TaskId, Error]),
|
||||||
ok
|
ok
|
||||||
end,
|
end,
|
||||||
|
|||||||
@ -4,9 +4,9 @@
|
|||||||
%%% @doc
|
%%% @doc
|
||||||
%%%
|
%%%
|
||||||
%%% @end
|
%%% @end
|
||||||
%%% Created : 21. 5月 2025 18:38
|
%%% Created : 20. 4月 2026 00:00
|
||||||
%%%-------------------------------------------------------------------
|
%%%-------------------------------------------------------------------
|
||||||
-module(efka_remote_agent).
|
-module(efka_client).
|
||||||
-author("anlicheng").
|
-author("anlicheng").
|
||||||
-include("message.hrl").
|
-include("message.hrl").
|
||||||
-include("message_pb.hrl").
|
-include("message_pb.hrl").
|
||||||
@ -25,7 +25,6 @@
|
|||||||
|
|
||||||
%% 标记当前agent的状态,只有在 activated 状态下才可以正常的发送数据
|
%% 标记当前agent的状态,只有在 activated 状态下才可以正常的发送数据
|
||||||
-define(STATE_DENIED, denied).
|
-define(STATE_DENIED, denied).
|
||||||
-define(STATE_CONNECTING, connecting).
|
|
||||||
-define(STATE_AUTH, auth).
|
-define(STATE_AUTH, auth).
|
||||||
%% 不能推送消息到服务,但是可以接受服务器的部分指令
|
%% 不能推送消息到服务,但是可以接受服务器的部分指令
|
||||||
-define(STATE_RESTRICTED, restricted).
|
-define(STATE_RESTRICTED, restricted).
|
||||||
@ -33,8 +32,7 @@
|
|||||||
-define(STATE_ACTIVATED, activated).
|
-define(STATE_ACTIVATED, activated).
|
||||||
|
|
||||||
-record(state, {
|
-record(state, {
|
||||||
transport_pid :: undefined | pid(),
|
socket :: undefined | ssl:sslsocket()
|
||||||
transport_ref :: undefined | reference()
|
|
||||||
}).
|
}).
|
||||||
|
|
||||||
%%%===================================================================
|
%%%===================================================================
|
||||||
@ -57,9 +55,6 @@ close_task_event_stream(TaskId, Reason) when is_integer(TaskId), is_binary(Reaso
|
|||||||
ping(AdCode, BootTime, Province, City, EfkaVersion, KernelArch, Ips, CpuCore, CpuLoad, CpuTemperature, Disk, Memory, Interfaces) ->
|
ping(AdCode, BootTime, Province, City, EfkaVersion, KernelArch, Ips, CpuCore, CpuLoad, CpuTemperature, Disk, Memory, Interfaces) ->
|
||||||
gen_statem:cast(?SERVER, {ping, AdCode, BootTime, Province, City, EfkaVersion, KernelArch, Ips, CpuCore, CpuLoad, CpuTemperature, Disk, Memory, Interfaces}).
|
gen_statem:cast(?SERVER, {ping, AdCode, BootTime, Province, City, EfkaVersion, KernelArch, Ips, CpuCore, CpuLoad, CpuTemperature, Disk, Memory, Interfaces}).
|
||||||
|
|
||||||
%% @doc Creates a gen_statem process which calls Module:init/1 to
|
|
||||||
%% initialize. To ensure a synchronized start-up procedure, this
|
|
||||||
%% function does not return until Module:init/1 has returned.
|
|
||||||
start_link() ->
|
start_link() ->
|
||||||
gen_statem:start_link({local, ?SERVER}, ?MODULE, [], []).
|
gen_statem:start_link({local, ?SERVER}, ?MODULE, [], []).
|
||||||
|
|
||||||
@ -67,31 +62,19 @@ start_link() ->
|
|||||||
%%% gen_statem callbacks
|
%%% gen_statem callbacks
|
||||||
%%%===================================================================
|
%%%===================================================================
|
||||||
|
|
||||||
%% @private
|
|
||||||
%% @doc Whenever a gen_statem is started using gen_statem:start/[3,4] or
|
|
||||||
%% gen_statem:start_link/[3,4], this function is called by the new
|
|
||||||
%% process to initialize.
|
|
||||||
init([]) ->
|
init([]) ->
|
||||||
erlang:start_timer(0, self(), create_transport),
|
erlang:start_timer(0, self(), create_transport),
|
||||||
{ok, ?STATE_DENIED, #state{}}.
|
{ok, ?STATE_DENIED, #state{socket = undefined}}.
|
||||||
|
|
||||||
%% @private
|
|
||||||
%% @doc This function is called by a gen_statem when it needs to find out
|
|
||||||
%% the callback mode of the callback module.
|
|
||||||
callback_mode() ->
|
callback_mode() ->
|
||||||
handle_event_function.
|
handle_event_function.
|
||||||
|
|
||||||
%% @private
|
|
||||||
%% @doc If callback_mode is handle_event_function, then whenever a
|
|
||||||
%% gen_statem receives an event from call/2, cast/2, or as a normal
|
|
||||||
%% process message, this function is called.
|
|
||||||
|
|
||||||
%% 异步发送数据, 连接存在时候直接发送;否则缓存到mnesia
|
%% 异步发送数据, 连接存在时候直接发送;否则缓存到mnesia
|
||||||
handle_event(cast, {metric_data, RouteKey, Metric}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
|
handle_event(cast, {metric_data, RouteKey, Metric}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
Packet = message_pb:encode_msg(#'CastFrame'{
|
Packet = message_pb:encode_msg(#'CastFrame'{
|
||||||
body = {data, #'Data'{route_key = RouteKey, metric = Metric}}
|
body = {data, #'Data'{route_key = RouteKey, metric = Metric}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet),
|
send_packet(Socket, Packet),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
handle_event(cast, {metric_data, RouteKey, Metric}, _, State) ->
|
handle_event(cast, {metric_data, RouteKey, Metric}, _, State) ->
|
||||||
@ -101,98 +84,73 @@ handle_event(cast, {metric_data, RouteKey, Metric}, _, State) ->
|
|||||||
ok = cache_model:insert(Packet),
|
ok = cache_model:insert(Packet),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
handle_event(cast, {task_event_stream, TaskId, Type, Stream}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
|
handle_event(cast, {task_event_stream, TaskId, Type, Stream}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
logger:debug("[efka_remote_agent] event_stream task_id: ~p, stream: ~ts", [TaskId, Stream]),
|
logger:debug("[efka_client] event_stream task_id: ~p, stream: ~ts", [TaskId, Stream]),
|
||||||
EventPacket = message_pb:encode_msg(#'CastFrame'{
|
EventPacket = message_pb:encode_msg(#'CastFrame'{
|
||||||
body = {event_stream, #'TaskEventStream'{task_id = TaskId, type = Type, stream = Stream}}
|
body = {event_stream, #'TaskEventStream'{task_id = TaskId, type = Type, stream = Stream}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, EventPacket),
|
send_packet(Socket, EventPacket),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
handle_event(cast, {close_task_event_stream, TaskId, Reason}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
|
handle_event(cast, {close_task_event_stream, TaskId, Reason}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
EventPacket = message_pb:encode_msg(#'CastFrame'{
|
EventPacket = message_pb:encode_msg(#'CastFrame'{
|
||||||
body = {event_stream, #'TaskEventStream'{task_id = TaskId, type = <<"close">>, stream = Reason}}
|
body = {event_stream, #'TaskEventStream'{task_id = TaskId, type = <<"close">>, stream = Reason}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, EventPacket),
|
send_packet(Socket, EventPacket),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
%% 其他情况下直接忽略
|
%% 其他情况下直接忽略
|
||||||
handle_event(cast, {task_event_stream, _TaskId, _Stream}, _, State = #state{}) ->
|
handle_event(cast, {task_event_stream, _TaskId, _Stream}, _, State = #state{}) ->
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
%handle_event(cast, {ping, AdCode, BootTime, Province, City, EfkaVersion, KernelArch, Ips, CpuCore, CpuLoad, CpuTemperature, Disk, Memory, Interfaces}, ?STATE_ACTIVATED,
|
|
||||||
% State = #state{transport_pid = TransportPid}) ->
|
|
||||||
%
|
|
||||||
% Ping = message_pb:encode_msg(#ping{
|
|
||||||
% adcode = AdCode,
|
|
||||||
% boot_time = BootTime,
|
|
||||||
% province = Province,
|
|
||||||
% city = City,
|
|
||||||
% efka_version = EfkaVersion,
|
|
||||||
% kernel_arch = KernelArch,
|
|
||||||
% ips = Ips,
|
|
||||||
% cpu_core = CpuCore,
|
|
||||||
% cpu_load = CpuLoad,
|
|
||||||
% cpu_temperature = CpuTemperature,
|
|
||||||
% disk = Disk,
|
|
||||||
% memory = Memory,
|
|
||||||
% interfaces = Interfaces
|
|
||||||
% }),
|
|
||||||
% efka_transport:send(TransportPid, Ping),
|
|
||||||
% {keep_state, State};
|
|
||||||
|
|
||||||
%% 异步建立到服务器的连接
|
%% 异步建立到服务器的连接
|
||||||
handle_event(info, {timeout, _, create_transport}, ?STATE_DENIED, State) ->
|
handle_event(info, {timeout, _, create_transport}, ?STATE_DENIED, State) ->
|
||||||
{ok, Props} = application:get_env(efka, tls_server),
|
case connect_socket() of
|
||||||
Host = proplists:get_value(host, Props),
|
{ok, Socket} ->
|
||||||
Port = proplists:get_value(port, Props),
|
|
||||||
{ok, {TransportPid, TransportRef}} = efka_transport:start_monitor(self(), Host, Port),
|
|
||||||
efka_transport:connect(TransportPid),
|
|
||||||
|
|
||||||
{next_state, ?STATE_CONNECTING, State#state{transport_pid = TransportPid, transport_ref = TransportRef}};
|
|
||||||
|
|
||||||
handle_event(info, {connect_reply, Reply}, ?STATE_CONNECTING, State = #state{transport_pid = TransportPid}) ->
|
|
||||||
case Reply of
|
|
||||||
ok ->
|
|
||||||
AuthPacket = auth_packet(),
|
AuthPacket = auth_packet(),
|
||||||
efka_transport:send(TransportPid, AuthPacket),
|
send_packet(Socket, AuthPacket),
|
||||||
{next_state, ?STATE_AUTH, State, [{state_timeout, 5000, auth_timeout}]};
|
{next_state, ?STATE_AUTH, State#state{socket = Socket}, [{state_timeout, 5000, auth_timeout}]};
|
||||||
{error, Reason} ->
|
{error, Reason} ->
|
||||||
logger:debug("[efka_remote_agent] connect failed, error: ~p, pid: ~p", [Reason, TransportPid]),
|
logger:debug("[efka_client] connect failed, error: ~p", [Reason]),
|
||||||
efka_transport:stop(TransportPid),
|
schedule_reconnect(),
|
||||||
{next_state, ?STATE_DENIED, State#state{transport_pid = undefined}}
|
{keep_state, State#state{socket = undefined}}
|
||||||
end;
|
end;
|
||||||
|
handle_event(info, {timeout, _, create_transport}, _, State) ->
|
||||||
|
{keep_state, State};
|
||||||
|
|
||||||
handle_event(info, {server_packet, PacketBin}, ?STATE_AUTH, State = #state{transport_pid = TransportPid}) ->
|
handle_event(info, {ssl, Socket, PacketBin}, ?STATE_AUTH, State = #state{socket = Socket}) ->
|
||||||
#'ResponseFrame'{packet_id = 1, body = {auth_reply, #'AuthReply'{code = Code, payload = Message}}} =
|
#'ResponseFrame'{packet_id = 1, body = {auth_reply, #'AuthReply'{code = Code, payload = Message}}} =
|
||||||
message_pb:decode_msg(PacketBin, 'ResponseFrame'),
|
message_pb:decode_msg(PacketBin, 'ResponseFrame'),
|
||||||
case Code of
|
case Code of
|
||||||
0 ->
|
0 ->
|
||||||
logger:debug("[efka_remote_agent] auth success, message: ~p", [Message]),
|
logger:debug("[efka_client] auth success, message: ~p", [Message]),
|
||||||
{next_state, ?STATE_ACTIVATED, State, [{next_event, info, flush_cache}]};
|
{next_state, ?STATE_ACTIVATED, State, [{next_event, info, flush_cache}]};
|
||||||
1 ->
|
1 ->
|
||||||
logger:debug("[efka_remote_agent] auth denied, message: ~p", [Message]),
|
logger:debug("[efka_client] auth denied, message: ~p", [Message]),
|
||||||
{next_state, ?STATE_RESTRICTED, State};
|
{next_state, ?STATE_RESTRICTED, State};
|
||||||
2 ->
|
2 ->
|
||||||
logger:debug("[efka_remote_agent] auth failed, message: ~p", [Message]),
|
logger:debug("[efka_client] auth failed, message: ~p", [Message]),
|
||||||
efka_transport:stop(TransportPid),
|
disconnect(Socket),
|
||||||
{next_state, ?STATE_DENIED, State#state{transport_pid = undefined}};
|
schedule_reconnect(),
|
||||||
|
{next_state, ?STATE_DENIED, State#state{socket = undefined}};
|
||||||
_ ->
|
_ ->
|
||||||
logger:debug("[efka_remote_agent] auth failed, invalid message"),
|
logger:debug("[efka_client] auth failed, invalid message"),
|
||||||
efka_transport:stop(TransportPid),
|
disconnect(Socket),
|
||||||
{next_state, ?STATE_DENIED, State#state{transport_pid = undefined}}
|
schedule_reconnect(),
|
||||||
|
{next_state, ?STATE_DENIED, State#state{socket = undefined}}
|
||||||
end;
|
end;
|
||||||
|
|
||||||
handle_event(state_timeout, auth_timeout, ?STATE_AUTH, State = #state{transport_pid = TransportPid}) ->
|
handle_event(state_timeout, auth_timeout, ?STATE_AUTH, State = #state{socket = Socket}) ->
|
||||||
logger:debug("[efka_remote_agent] auth request timeout"),
|
logger:debug("[efka_client] auth request timeout"),
|
||||||
efka_transport:stop(TransportPid),
|
disconnect(Socket),
|
||||||
{next_state, ?STATE_DENIED, State#state{transport_pid = undefined}};
|
schedule_reconnect(),
|
||||||
|
{next_state, ?STATE_DENIED, State#state{socket = undefined}};
|
||||||
|
|
||||||
%% 将缓存中的数据推送到服务器端
|
%% 将缓存中的数据推送到服务器端
|
||||||
handle_event(info, flush_cache, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
|
handle_event(info, flush_cache, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
case cache_model:fetch_next() of
|
case cache_model:fetch_next() of
|
||||||
{ok, {Id, Packet}} ->
|
{ok, {Id, Packet}} ->
|
||||||
efka_transport:send(TransportPid, Packet),
|
send_packet(Socket, Packet),
|
||||||
cache_model:delete(Id),
|
cache_model:delete(Id),
|
||||||
{keep_state, State, [{next_event, info, flush_cache}]};
|
{keep_state, State, [{next_event, info, flush_cache}]};
|
||||||
error ->
|
error ->
|
||||||
@ -201,7 +159,7 @@ handle_event(info, flush_cache, ?STATE_ACTIVATED, State = #state{transport_pid =
|
|||||||
handle_event(info, flush_cache, _, State) ->
|
handle_event(info, flush_cache, _, State) ->
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
handle_event(info, {server_packet, <<8, _/binary>> = PacketBin}, StateName, State)
|
handle_event(info, {ssl, Socket, <<8, _/binary>> = PacketBin}, StateName, State = #state{socket = Socket})
|
||||||
when StateName =:= ?STATE_ACTIVATED; StateName =:= ?STATE_RESTRICTED ->
|
when StateName =:= ?STATE_ACTIVATED; StateName =:= ?STATE_RESTRICTED ->
|
||||||
#'RequestFrame'{packet_id = PacketId, body = Body} = message_pb:decode_msg(PacketBin, 'RequestFrame'),
|
#'RequestFrame'{packet_id = PacketId, body = Body} = message_pb:decode_msg(PacketBin, 'RequestFrame'),
|
||||||
true = is_integer(PacketId) andalso PacketId > 0,
|
true = is_integer(PacketId) andalso PacketId > 0,
|
||||||
@ -211,25 +169,22 @@ handle_event(info, {server_packet, <<8, _/binary>> = PacketBin}, StateName, Stat
|
|||||||
{container_request, Request} ->
|
{container_request, Request} ->
|
||||||
{keep_state, State, [{next_event, info, {container_request, PacketId, Request}}]}
|
{keep_state, State, [{next_event, info, {container_request, PacketId, Request}}]}
|
||||||
end;
|
end;
|
||||||
handle_event(info, {server_packet, <<10, _/binary>> = PacketBin}, ?STATE_ACTIVATED, State) ->
|
handle_event(info, {ssl, Socket, <<10, _/binary>> = PacketBin}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
#'CastFrame'{body = {pub, Pub}} = message_pb:decode_msg(PacketBin, 'CastFrame'),
|
#'CastFrame'{body = {pub, Pub}} = message_pb:decode_msg(PacketBin, 'CastFrame'),
|
||||||
{keep_state, State, [{next_event, info, {server_cast, Pub}}]};
|
{keep_state, State, [{next_event, info, {server_cast, Pub}}]};
|
||||||
handle_event(info, {server_packet, <<18, _/binary>> = PacketBin}, ?STATE_ACTIVATED, State) ->
|
handle_event(info, {ssl, Socket, <<18, _/binary>> = PacketBin}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
#'CastFrame'{body = {command, Command}} = message_pb:decode_msg(PacketBin, 'CastFrame'),
|
#'CastFrame'{body = {command, Command}} = message_pb:decode_msg(PacketBin, 'CastFrame'),
|
||||||
{keep_state, State, [{next_event, info, {server_cast, Command}}]};
|
{keep_state, State, [{next_event, info, {server_cast, Command}}]};
|
||||||
handle_event(info, {server_packet, <<10, _/binary>> = PacketBin}, ?STATE_RESTRICTED, State) ->
|
handle_event(info, {ssl, Socket, <<10, _/binary>> = PacketBin}, ?STATE_RESTRICTED, State = #state{socket = Socket}) ->
|
||||||
#'CastFrame'{body = {pub, Pub}} = message_pb:decode_msg(PacketBin, 'CastFrame'),
|
#'CastFrame'{body = {pub, Pub}} = message_pb:decode_msg(PacketBin, 'CastFrame'),
|
||||||
{keep_state, State, [{next_event, info, {server_cast, Pub}}]};
|
{keep_state, State, [{next_event, info, {server_cast, Pub}}]};
|
||||||
handle_event(info, {server_packet, <<18, _/binary>> = PacketBin}, ?STATE_RESTRICTED, State) ->
|
handle_event(info, {ssl, Socket, <<18, _/binary>> = PacketBin}, ?STATE_RESTRICTED, State = #state{socket = Socket}) ->
|
||||||
#'CastFrame'{body = {command, Command}} = message_pb:decode_msg(PacketBin, 'CastFrame'),
|
#'CastFrame'{body = {command, Command}} = message_pb:decode_msg(PacketBin, 'CastFrame'),
|
||||||
{keep_state, State, [{next_event, info, {server_cast, Command}}]};
|
{keep_state, State, [{next_event, info, {server_cast, Command}}]};
|
||||||
|
|
||||||
%% 云端服务器推送了消息
|
|
||||||
%% 激活消息
|
|
||||||
|
|
||||||
%% 微服务部署
|
%% 微服务部署
|
||||||
handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {list, #'ContainerRequest.List'{all = _All}}}},
|
handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {list, #'ContainerRequest.List'{all = _All}}}},
|
||||||
?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
|
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
case docker_manager:get_containers() of
|
case docker_manager:get_containers() of
|
||||||
{ok, Containers} ->
|
{ok, Containers} ->
|
||||||
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
||||||
@ -238,7 +193,7 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {l
|
|||||||
reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(Containers)}}
|
reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(Containers)}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet);
|
send_packet(Socket, Packet);
|
||||||
{error, Reason} when is_binary(Reason) ->
|
{error, Reason} when is_binary(Reason) ->
|
||||||
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
||||||
packet_id = PacketId,
|
packet_id = PacketId,
|
||||||
@ -246,14 +201,14 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {l
|
|||||||
reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}}
|
reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet)
|
send_packet(Socket, Packet)
|
||||||
end,
|
end,
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {deploy, #'ContainerRequest.Deploy'{
|
handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {deploy, #'ContainerRequest.Deploy'{
|
||||||
task_id = TaskId,
|
task_id = TaskId,
|
||||||
params = Params
|
params = Params
|
||||||
}}}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
|
}}}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
case docker_manager:deploy(TaskId, Params) of
|
case docker_manager:deploy(TaskId, Params) of
|
||||||
ok ->
|
ok ->
|
||||||
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
||||||
@ -262,7 +217,7 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {d
|
|||||||
reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(<<"ok">>)}}
|
reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(<<"ok">>)}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet);
|
send_packet(Socket, Packet);
|
||||||
{error, Reason} when is_binary(Reason) ->
|
{error, Reason} when is_binary(Reason) ->
|
||||||
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
||||||
packet_id = PacketId,
|
packet_id = PacketId,
|
||||||
@ -270,12 +225,12 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {d
|
|||||||
reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}}
|
reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet)
|
send_packet(Socket, Packet)
|
||||||
end,
|
end,
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {start, #'ContainerRequest.Start'{target = Target}}}},
|
handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {start, #'ContainerRequest.Start'{target = Target}}}},
|
||||||
?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
|
?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
ContainerTarget = container_target(Target),
|
ContainerTarget = container_target(Target),
|
||||||
case docker_manager:start_container(ContainerTarget) of
|
case docker_manager:start_container(ContainerTarget) of
|
||||||
ok ->
|
ok ->
|
||||||
@ -285,7 +240,7 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {s
|
|||||||
reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(<<"ok">>)}}
|
reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(<<"ok">>)}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet);
|
send_packet(Socket, Packet);
|
||||||
{error, Reason} when is_binary(Reason) ->
|
{error, Reason} when is_binary(Reason) ->
|
||||||
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
||||||
packet_id = PacketId,
|
packet_id = PacketId,
|
||||||
@ -293,14 +248,14 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {s
|
|||||||
reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}}
|
reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet)
|
send_packet(Socket, Packet)
|
||||||
end,
|
end,
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {stop, #'ContainerRequest.Stop'{
|
handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {stop, #'ContainerRequest.Stop'{
|
||||||
target = Target,
|
target = Target,
|
||||||
timeout_seconds = TimeoutSeconds
|
timeout_seconds = TimeoutSeconds
|
||||||
}}}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
|
}}}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
ContainerTarget = container_target(Target),
|
ContainerTarget = container_target(Target),
|
||||||
case docker_manager:stop_container(ContainerTarget, TimeoutSeconds) of
|
case docker_manager:stop_container(ContainerTarget, TimeoutSeconds) of
|
||||||
ok ->
|
ok ->
|
||||||
@ -310,7 +265,7 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {s
|
|||||||
reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(<<"ok">>)}}
|
reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(<<"ok">>)}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet);
|
send_packet(Socket, Packet);
|
||||||
{error, Reason} when is_binary(Reason) ->
|
{error, Reason} when is_binary(Reason) ->
|
||||||
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
||||||
packet_id = PacketId,
|
packet_id = PacketId,
|
||||||
@ -318,14 +273,14 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {s
|
|||||||
reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}}
|
reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet)
|
send_packet(Socket, Packet)
|
||||||
end,
|
end,
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {kill, #'ContainerRequest.Kill'{
|
handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {kill, #'ContainerRequest.Kill'{
|
||||||
target = Target,
|
target = Target,
|
||||||
signal = Signal
|
signal = Signal
|
||||||
}}}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
|
}}}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
ContainerTarget = container_target(Target),
|
ContainerTarget = container_target(Target),
|
||||||
case docker_manager:kill_container(ContainerTarget, to_binary(Signal)) of
|
case docker_manager:kill_container(ContainerTarget, to_binary(Signal)) of
|
||||||
ok ->
|
ok ->
|
||||||
@ -335,7 +290,7 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {k
|
|||||||
reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(<<"ok">>)}}
|
reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(<<"ok">>)}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet);
|
send_packet(Socket, Packet);
|
||||||
{error, Reason} when is_binary(Reason) ->
|
{error, Reason} when is_binary(Reason) ->
|
||||||
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
||||||
packet_id = PacketId,
|
packet_id = PacketId,
|
||||||
@ -343,7 +298,7 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {k
|
|||||||
reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}}
|
reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet)
|
send_packet(Socket, Packet)
|
||||||
end,
|
end,
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
@ -351,7 +306,7 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {r
|
|||||||
target = Target,
|
target = Target,
|
||||||
force = Force,
|
force = Force,
|
||||||
remove_volumes = RemoveVolumes
|
remove_volumes = RemoveVolumes
|
||||||
}}}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
|
}}}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
ContainerTarget = container_target(Target),
|
ContainerTarget = container_target(Target),
|
||||||
case docker_manager:remove_container(ContainerTarget, to_bool(Force), to_bool(RemoveVolumes)) of
|
case docker_manager:remove_container(ContainerTarget, to_bool(Force), to_bool(RemoveVolumes)) of
|
||||||
ok ->
|
ok ->
|
||||||
@ -361,7 +316,7 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {r
|
|||||||
reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(<<"ok">>)}}
|
reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(<<"ok">>)}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet);
|
send_packet(Socket, Packet);
|
||||||
{error, Reason} when is_binary(Reason) ->
|
{error, Reason} when is_binary(Reason) ->
|
||||||
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
||||||
packet_id = PacketId,
|
packet_id = PacketId,
|
||||||
@ -369,14 +324,14 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {r
|
|||||||
reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}}
|
reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet)
|
send_packet(Socket, Packet)
|
||||||
end,
|
end,
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {config, #'ContainerRequest.Config'{
|
handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {config, #'ContainerRequest.Config'{
|
||||||
target = Target,
|
target = Target,
|
||||||
config = Config
|
config = Config
|
||||||
}}}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
|
}}}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
ContainerTarget = container_target(Target),
|
ContainerTarget = container_target(Target),
|
||||||
case docker_manager:config_container(ContainerTarget, iolist_to_binary(Config)) of
|
case docker_manager:config_container(ContainerTarget, iolist_to_binary(Config)) of
|
||||||
ok ->
|
ok ->
|
||||||
@ -386,7 +341,7 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {c
|
|||||||
reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(<<"ok">>)}}
|
reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(<<"ok">>)}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet);
|
send_packet(Socket, Packet);
|
||||||
{error, Reason} ->
|
{error, Reason} ->
|
||||||
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
||||||
packet_id = PacketId,
|
packet_id = PacketId,
|
||||||
@ -394,111 +349,89 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {c
|
|||||||
reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}}
|
reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet)
|
send_packet(Socket, Packet)
|
||||||
end,
|
end,
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
handle_event(info, {container_request, PacketId, _Request}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
|
handle_event(info, {container_request, PacketId, _Request}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
||||||
packet_id = PacketId,
|
packet_id = PacketId,
|
||||||
body = {rpc_reply, #'RpcReply'{
|
body = {rpc_reply, #'RpcReply'{
|
||||||
reply = {error, #'RpcReply.RpcError'{code = -1, message = <<"unsupported container request">>}}
|
reply = {error, #'RpcReply.RpcError'{code = -1, message = <<"unsupported container request">>}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet),
|
send_packet(Socket, Packet),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
handle_event(info, {container_request, PacketId, _Request}, ?STATE_RESTRICTED, State = #state{transport_pid = TransportPid}) ->
|
handle_event(info, {container_request, PacketId, _Request}, ?STATE_RESTRICTED, State = #state{socket = Socket}) ->
|
||||||
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
||||||
packet_id = PacketId,
|
packet_id = PacketId,
|
||||||
body = {rpc_reply, #'RpcReply'{
|
body = {rpc_reply, #'RpcReply'{
|
||||||
reply = {error, #'RpcReply.RpcError'{code = -1, message = <<"agent restricted">>}}
|
reply = {error, #'RpcReply.RpcError'{code = -1, message = <<"agent restricted">>}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet),
|
send_packet(Socket, Packet),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
handle_event(info, {server_rpc, PacketId, #'RpcRequest'{method = Method}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
|
handle_event(info, {server_rpc, PacketId, #'RpcRequest'{method = Method}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
||||||
packet_id = PacketId,
|
packet_id = PacketId,
|
||||||
body = {rpc_reply, #'RpcReply'{
|
body = {rpc_reply, #'RpcReply'{
|
||||||
reply = {error, #'RpcReply.RpcError'{code = -1, message = <<"unsupported rpc request: ", Method/binary>>}}
|
reply = {error, #'RpcReply.RpcError'{code = -1, message = <<"unsupported rpc request: ", Method/binary>>}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet),
|
send_packet(Socket, Packet),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
handle_event(info, {server_rpc, PacketId, _Request}, ?STATE_RESTRICTED, State = #state{transport_pid = TransportPid}) ->
|
handle_event(info, {server_rpc, PacketId, _Request}, ?STATE_RESTRICTED, State = #state{socket = Socket}) ->
|
||||||
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
Packet = message_pb:encode_msg(#'ResponseFrame'{
|
||||||
packet_id = PacketId,
|
packet_id = PacketId,
|
||||||
body = {rpc_reply, #'RpcReply'{
|
body = {rpc_reply, #'RpcReply'{
|
||||||
reply = {error, #'RpcReply.RpcError'{code = -1, message = <<"agent restricted">>}}
|
reply = {error, #'RpcReply.RpcError'{code = -1, message = <<"agent restricted">>}}
|
||||||
}}
|
}}
|
||||||
}),
|
}),
|
||||||
efka_transport:send(TransportPid, Packet),
|
send_packet(Socket, Packet),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
%% 处理task_log
|
|
||||||
%handle_event(info, {server_async_call, PacketId, <<?PUSH_TASK_LOG:8, TaskLogBin/binary>>}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) ->
|
|
||||||
% #fetch_task_log{task_id = TaskId} = message_pb:decode_msg(TaskLogBin, fetch_task_log),
|
|
||||||
% logger:debug("[efka_remote_agent] get task_log request: ~p", [TaskId]),
|
|
||||||
% {ok, Logs} = efka_inetd_task_log:get_logs(TaskId),
|
|
||||||
% Reply = case length(Logs) > 0 of
|
|
||||||
% true ->
|
|
||||||
% Result = iolist_to_binary(jiffy:encode(Logs, [force_utf8])),
|
|
||||||
% #async_call_reply{code = 1, result = Result};
|
|
||||||
% false ->
|
|
||||||
% #async_call_reply{code = 1, result = <<"[]">>}
|
|
||||||
% end,
|
|
||||||
% efka_transport:send(TransportPid, message_pb:encode_msg(Reply)),
|
|
||||||
%
|
|
||||||
% {keep_state, State};
|
|
||||||
|
|
||||||
%% 处理命令
|
%% 处理命令
|
||||||
handle_event(info, {server_cast, #'Command'{command_type = ?COMMAND_AUTH, command = Auth0}}, StateName, State = #state{transport_pid = TransportPid}) ->
|
handle_event(info, {server_cast, #'Command'{command_type = ?COMMAND_AUTH, command = Auth0}}, StateName,
|
||||||
|
State = #state{socket = Socket}) ->
|
||||||
Auth = binary_to_integer(Auth0),
|
Auth = binary_to_integer(Auth0),
|
||||||
case {Auth, StateName} of
|
case {Auth, StateName} of
|
||||||
{1, ?STATE_ACTIVATED} ->
|
{1, ?STATE_ACTIVATED} ->
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
{1, ?STATE_DENIED} ->
|
{1, _} ->
|
||||||
%% 重新激活, 需要重新校验
|
|
||||||
AuthPacket = auth_packet(),
|
AuthPacket = auth_packet(),
|
||||||
efka_transport:send(TransportPid, AuthPacket),
|
send_packet(Socket, AuthPacket),
|
||||||
{next_state, ?STATE_AUTH, State, [{state_timeout, 5000, auth_timeout}]};
|
{next_state, ?STATE_AUTH, State, [{state_timeout, 5000, auth_timeout}]};
|
||||||
{0, _} ->
|
{0, _} ->
|
||||||
%% 这个时候的主机应该是受限制的状态,不允许发送消息;但是能够接受服务器推送的消息
|
|
||||||
{next_state, ?STATE_RESTRICTED, State}
|
{next_state, ?STATE_RESTRICTED, State}
|
||||||
end;
|
end;
|
||||||
|
|
||||||
%% 处理Pub/Sub机制
|
%% 处理Pub/Sub机制
|
||||||
handle_event(info, {server_cast, #'Pub'{topic = Topic, qos = Qos, content = Content}}, ?STATE_ACTIVATED, State) ->
|
handle_event(info, {server_cast, #'Pub'{topic = Topic, qos = Qos, content = Content}}, ?STATE_ACTIVATED, State) ->
|
||||||
logger:debug("[efka_remote_agent] get pub topic: ~p, qos: ~p, content: ~p", [Topic, Qos, Content]),
|
logger:debug("[efka_client] get pub topic: ~p, qos: ~p, content: ~p", [Topic, Qos, Content]),
|
||||||
%% 消息发送到订阅系统
|
|
||||||
efka_subscription:publish(Topic, Qos, Content),
|
efka_subscription:publish(Topic, Qos, Content),
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
|
||||||
%% transport进程退出
|
handle_event(info, {ssl_error, Socket, Reason}, _, State = #state{socket = Socket}) ->
|
||||||
handle_event(info, {'DOWN', MRef, process, TransportPid, Reason}, _, State = #state{transport_ref = MRef}) ->
|
logger:debug("[efka_client] ssl error: ~p", [Reason]),
|
||||||
logger:debug("[efka_remote_agent] transport pid: ~p, exit with reason: ~p", [TransportPid, Reason]),
|
disconnect(Socket),
|
||||||
erlang:start_timer(5000, self(), create_transport),
|
schedule_reconnect(),
|
||||||
{next_state, ?STATE_DENIED, State#state{transport_pid = undefined, transport_ref = undefined}}.
|
{next_state, ?STATE_DENIED, State#state{socket = undefined}};
|
||||||
|
|
||||||
%% @private
|
handle_event(info, {ssl_closed, Socket}, _, State = #state{socket = Socket}) ->
|
||||||
%% @doc This function is called by a gen_statem when it is about to
|
schedule_reconnect(),
|
||||||
%% terminate. It should be the opposite of Module:init/1 and do any
|
{next_state, ?STATE_DENIED, State#state{socket = undefined}};
|
||||||
%% necessary cleaning up. When it returns, the gen_statem terminates with
|
|
||||||
%% Reason. The return value is ignored.
|
handle_event(info, Info, _, State = #state{}) ->
|
||||||
terminate(_Reason, _StateName, _State = #state{transport_pid = TransportPid}) ->
|
logger:notice("[efka_client] get unknown info: ~p", [Info]),
|
||||||
case is_pid(TransportPid) andalso is_process_alive(TransportPid) of
|
{keep_state, State}.
|
||||||
true ->
|
|
||||||
efka_transport:stop(TransportPid);
|
terminate(_Reason, _StateName, _State = #state{socket = Socket}) ->
|
||||||
false ->
|
disconnect(Socket),
|
||||||
ok
|
|
||||||
end,
|
|
||||||
ok.
|
ok.
|
||||||
|
|
||||||
%% @private
|
|
||||||
%% @doc Convert process state when code is changed
|
|
||||||
code_change(_OldVsn, StateName, State = #state{}, _Extra) ->
|
code_change(_OldVsn, StateName, State = #state{}, _Extra) ->
|
||||||
{ok, StateName, State}.
|
{ok, StateName, State}.
|
||||||
|
|
||||||
@ -508,7 +441,6 @@ code_change(_OldVsn, StateName, State = #state{}, _Extra) ->
|
|||||||
|
|
||||||
-spec auth_packet() -> binary().
|
-spec auth_packet() -> binary().
|
||||||
auth_packet() ->
|
auth_packet() ->
|
||||||
%% 重新激活, 需要重新校验
|
|
||||||
{ok, AuthInfo} = application:get_env(efka, auth),
|
{ok, AuthInfo} = application:get_env(efka, auth),
|
||||||
UUID = proplists:get_value(uuid, AuthInfo),
|
UUID = proplists:get_value(uuid, AuthInfo),
|
||||||
Username = proplists:get_value(username, AuthInfo),
|
Username = proplists:get_value(username, AuthInfo),
|
||||||
@ -525,6 +457,34 @@ auth_packet() ->
|
|||||||
}}
|
}}
|
||||||
}).
|
}).
|
||||||
|
|
||||||
|
-spec connect_socket() -> {ok, ssl:sslsocket()} | {error, term()}.
|
||||||
|
connect_socket() ->
|
||||||
|
{ok, Props} = application:get_env(efka, tls_server),
|
||||||
|
Host = proplists:get_value(host, Props),
|
||||||
|
Port = proplists:get_value(port, Props),
|
||||||
|
SslOptions = [
|
||||||
|
binary,
|
||||||
|
{active, true},
|
||||||
|
{packet, 4},
|
||||||
|
{verify, verify_none}
|
||||||
|
],
|
||||||
|
ssl:connect(Host, Port, SslOptions, 5000).
|
||||||
|
|
||||||
|
-spec send_packet(ssl:sslsocket(), binary()) -> ok.
|
||||||
|
send_packet(Socket, Packet) when is_binary(Packet) ->
|
||||||
|
ok = ssl:send(Socket, Packet).
|
||||||
|
|
||||||
|
-spec disconnect(undefined | ssl:sslsocket()) -> ok.
|
||||||
|
disconnect(undefined) ->
|
||||||
|
ok;
|
||||||
|
disconnect(Socket) ->
|
||||||
|
catch ssl:close(Socket),
|
||||||
|
ok.
|
||||||
|
|
||||||
|
-spec schedule_reconnect() -> reference().
|
||||||
|
schedule_reconnect() ->
|
||||||
|
erlang:start_timer(5000, self(), create_transport).
|
||||||
|
|
||||||
-spec encode_rpc_payload(any()) -> binary().
|
-spec encode_rpc_payload(any()) -> binary().
|
||||||
encode_rpc_payload(Payload) ->
|
encode_rpc_payload(Payload) ->
|
||||||
jiffy:encode(Payload, [force_utf8]).
|
jiffy:encode(Payload, [force_utf8]).
|
||||||
@ -104,7 +104,7 @@ handle_call(_Request, _From, State = #state{}) ->
|
|||||||
{stop, Reason :: term(), NewState :: #state{}}).
|
{stop, Reason :: term(), NewState :: #state{}}).
|
||||||
handle_cast({metric_data, RouteKey, Metric}, State = #state{service_id = ServiceId}) ->
|
handle_cast({metric_data, RouteKey, Metric}, State = #state{service_id = ServiceId}) ->
|
||||||
logger:debug("[efka_service] metric_data service_id: ~p, route_key: ~p, metric data: ~p", [ServiceId, RouteKey, Metric]),
|
logger:debug("[efka_service] metric_data service_id: ~p, route_key: ~p, metric data: ~p", [ServiceId, RouteKey, Metric]),
|
||||||
efka_remote_agent:metric_data(RouteKey, Metric),
|
efka_client:metric_data(RouteKey, Metric),
|
||||||
{noreply, State};
|
{noreply, State};
|
||||||
|
|
||||||
handle_cast(_Request, State = #state{}) ->
|
handle_cast(_Request, State = #state{}) ->
|
||||||
|
|||||||
@ -92,12 +92,12 @@ init([]) ->
|
|||||||
},
|
},
|
||||||
|
|
||||||
#{
|
#{
|
||||||
id => 'efka_remote_agent',
|
id => 'efka_client',
|
||||||
start => {'efka_remote_agent', start_link, []},
|
start => {'efka_client', start_link, []},
|
||||||
restart => permanent,
|
restart => permanent,
|
||||||
shutdown => 2000,
|
shutdown => 2000,
|
||||||
type => worker,
|
type => worker,
|
||||||
modules => ['efka_remote_agent']
|
modules => ['efka_client']
|
||||||
}
|
}
|
||||||
|
|
||||||
],
|
],
|
||||||
@ -105,4 +105,3 @@ init([]) ->
|
|||||||
{ok, {SupFlags, ChildSpecs}}.
|
{ok, {SupFlags, ChildSpecs}}.
|
||||||
|
|
||||||
%% internal functions
|
%% internal functions
|
||||||
|
|
||||||
|
|||||||
@ -1,153 +0,0 @@
|
|||||||
%%%-------------------------------------------------------------------
|
|
||||||
%%% @author anlicheng
|
|
||||||
%%% @copyright (C) 2025, <COMPANY>
|
|
||||||
%%% @doc
|
|
||||||
%%%
|
|
||||||
%%% @end
|
|
||||||
%%% Created : 20. 4月 2025 18:47
|
|
||||||
%%%-------------------------------------------------------------------
|
|
||||||
-module(efka_transport).
|
|
||||||
-author("anlicheng").
|
|
||||||
|
|
||||||
-behaviour(gen_server).
|
|
||||||
|
|
||||||
%% API
|
|
||||||
-export([start_monitor/3]).
|
|
||||||
-export([connect/1, send/2, stop/1]).
|
|
||||||
|
|
||||||
%% gen_server callbacks
|
|
||||||
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]).
|
|
||||||
|
|
||||||
-define(SERVER, ?MODULE).
|
|
||||||
|
|
||||||
-record(state, {
|
|
||||||
parent_pid :: pid(),
|
|
||||||
host :: string(),
|
|
||||||
port :: integer(),
|
|
||||||
socket :: undefined | ssl:sslsocket()
|
|
||||||
}).
|
|
||||||
|
|
||||||
-spec connect(Pid :: pid()) -> no_return().
|
|
||||||
connect(Pid) when is_pid(Pid) ->
|
|
||||||
gen_server:cast(Pid, connect).
|
|
||||||
|
|
||||||
-spec send(Pid :: pid(), Packet :: binary()) -> no_return().
|
|
||||||
send(Pid, Packet) when is_pid(Pid), is_binary(Packet) ->
|
|
||||||
gen_server:cast(Pid, {send, Packet}).
|
|
||||||
|
|
||||||
%% 关闭的时候不一定能成功,可能关闭的时候;transport进程已经退出了
|
|
||||||
-spec stop(Pid :: pid() | undefined) -> ok.
|
|
||||||
stop(undefined) ->
|
|
||||||
ok;
|
|
||||||
stop(Pid) when is_pid(Pid) ->
|
|
||||||
catch gen_server:stop(Pid, normal, 2000).
|
|
||||||
|
|
||||||
%% @doc Spawns the server and registers the local name (unique)
|
|
||||||
-spec(start_monitor(ParentPid :: pid(), Host :: string(), Port :: integer()) ->
|
|
||||||
{ok, {Pid :: pid(), MRef :: reference()}} | ignore | {error, Reason :: term()}).
|
|
||||||
start_monitor(ParentPid, Host, Port) when is_pid(ParentPid), is_list(Host), is_integer(Port) ->
|
|
||||||
gen_server:start_monitor(?MODULE, [ParentPid, Host, Port], []).
|
|
||||||
|
|
||||||
%%%===================================================================
|
|
||||||
%%% gen_server callbacks
|
|
||||||
%%%===================================================================
|
|
||||||
|
|
||||||
%% @private
|
|
||||||
%% @doc Initializes the server
|
|
||||||
-spec(init(Args :: term()) ->
|
|
||||||
{ok, State :: #state{}} | {ok, State :: #state{}, timeout() | hibernate} |
|
|
||||||
{stop, Reason :: term()} | ignore).
|
|
||||||
init([ParentPid, Host, Port]) ->
|
|
||||||
{ok, #state{parent_pid = ParentPid, host = Host, port = Port, socket = undefined}}.
|
|
||||||
|
|
||||||
%% @private
|
|
||||||
%% @doc Handling call messages
|
|
||||||
-spec(handle_call(Request :: term(), From :: {pid(), Tag :: term()}, State :: #state{}) ->
|
|
||||||
{reply, Reply :: term(), NewState :: #state{}} |
|
|
||||||
{reply, Reply :: term(), NewState :: #state{}, timeout() | hibernate} |
|
|
||||||
{noreply, NewState :: #state{}} |
|
|
||||||
{noreply, NewState :: #state{}, timeout() | hibernate} |
|
|
||||||
{stop, Reason :: term(), Reply :: term(), NewState :: #state{}} |
|
|
||||||
{stop, Reason :: term(), NewState :: #state{}}).
|
|
||||||
handle_call(_Req, _From, State = #state{}) ->
|
|
||||||
{reply, ok, State#state{}}.
|
|
||||||
|
|
||||||
%% @private
|
|
||||||
%% @doc Handling cast messages
|
|
||||||
-spec(handle_cast(Request :: term(), State :: #state{}) ->
|
|
||||||
{noreply, NewState :: #state{}} |
|
|
||||||
{noreply, NewState :: #state{}, timeout() | hibernate} |
|
|
||||||
{stop, Reason :: term(), NewState :: #state{}}).
|
|
||||||
%% 建立到目标服务器的连接
|
|
||||||
handle_cast(connect, State = #state{host = Host, port = Port, parent_pid = ParentPid}) ->
|
|
||||||
SslOptions = [
|
|
||||||
binary,
|
|
||||||
{packet, 4},
|
|
||||||
{verify, verify_none}
|
|
||||||
],
|
|
||||||
case ssl:connect(Host, Port, SslOptions, 5000) of
|
|
||||||
{ok, Socket} ->
|
|
||||||
ok = ssl:controlling_process(Socket, self()),
|
|
||||||
ParentPid ! {connect_reply, ok},
|
|
||||||
ping_ticker(),
|
|
||||||
{noreply, State#state{socket = Socket}};
|
|
||||||
{error, Reason} ->
|
|
||||||
ParentPid ! {connect_reply, {error, Reason}},
|
|
||||||
{noreply, State#state{socket = undefined}}
|
|
||||||
end;
|
|
||||||
|
|
||||||
handle_cast({send, Packet}, State = #state{socket = Socket}) ->
|
|
||||||
ok = ssl:send(Socket, Packet),
|
|
||||||
{noreply, State}.
|
|
||||||
|
|
||||||
%% @private
|
|
||||||
%% @doc Handling all non call/cast messages
|
|
||||||
-spec(handle_info(Info :: timeout() | term(), State :: #state{}) ->
|
|
||||||
{noreply, NewState :: #state{}} |
|
|
||||||
{noreply, NewState :: #state{}, timeout() | hibernate} |
|
|
||||||
{stop, Reason :: term(), NewState :: #state{}}).
|
|
||||||
%% 服务器主动推送的数据
|
|
||||||
handle_info({ssl, Socket, PacketBin}, State = #state{socket = Socket, parent_pid = ParentPid}) ->
|
|
||||||
ParentPid ! {server_packet, PacketBin},
|
|
||||||
{noreply, State};
|
|
||||||
|
|
||||||
handle_info({ssl_error, Socket, Reason}, State = #state{socket = Socket}) ->
|
|
||||||
logger:debug("[efka_transport] ssl error: ~p", [Reason]),
|
|
||||||
{stop, normal, State};
|
|
||||||
|
|
||||||
handle_info({ssl_closed, Socket}, State = #state{socket = Socket}) ->
|
|
||||||
{stop, normal, State};
|
|
||||||
|
|
||||||
handle_info({timeout, _, ping_ticker}, State) ->
|
|
||||||
ping_ticker(),
|
|
||||||
{noreply, State};
|
|
||||||
|
|
||||||
handle_info(Info, State = #state{}) ->
|
|
||||||
logger:notice("[efka_transport] get unknown info: ~p", [Info]),
|
|
||||||
{noreply, State}.
|
|
||||||
|
|
||||||
%% @private
|
|
||||||
%% @doc This function is called by a gen_server when it is about to
|
|
||||||
%% terminate. It should be the opposite of Module:init/1 and do any
|
|
||||||
%% necessary cleaning up. When it returns, the gen_server terminates
|
|
||||||
%% with Reason. The return value is ignored.
|
|
||||||
-spec(terminate(Reason :: (normal | shutdown | {shutdown, term()} | term()),
|
|
||||||
State :: #state{}) -> term()).
|
|
||||||
terminate(Reason, #state{}) ->
|
|
||||||
logger:notice("[efka_transport] terminate with reason: ~p", [Reason]),
|
|
||||||
ok.
|
|
||||||
|
|
||||||
%% @private
|
|
||||||
%% @doc Convert process state when code is changed
|
|
||||||
-spec(code_change(OldVsn :: term() | {down, term()}, State :: #state{},
|
|
||||||
Extra :: term()) ->
|
|
||||||
{ok, NewState :: #state{}} | {error, Reason :: term()}).
|
|
||||||
code_change(_OldVsn, State = #state{}, _Extra) ->
|
|
||||||
{ok, State}.
|
|
||||||
|
|
||||||
%%%===================================================================
|
|
||||||
%%% Internal functions
|
|
||||||
%%%===================================================================
|
|
||||||
|
|
||||||
ping_ticker() ->
|
|
||||||
erlang:start_timer(5000, self(), ping_ticker).
|
|
||||||
Loading…
x
Reference in New Issue
Block a user