diff --git a/src/docker/docker_deployer.erl b/src/docker/docker_deployer.erl index 3381f73..69eac0c 100644 --- a/src/docker/docker_deployer.erl +++ b/src/docker/docker_deployer.erl @@ -54,7 +54,7 @@ deploy(TaskId, ContainerDir, Params = #'ContainerDeployParams'{ case docker_commands:check_container_exist(ContainerName) of true -> 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 -> Image = normalize_image(Image0), @@ -92,15 +92,15 @@ deploy(TaskId, ContainerDir, Params = #'ContainerDeployParams'{ ShortContainerId = binary:part(ContainerId, 1, 12), trace_log(TaskId, <<"info">>, <<"容器创建成功: "/utf8, ShortContainerId/binary>>), 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} -> trace_log(TaskId, <<"error">>, <<"容器创建失败: "/utf8, Reason/binary>>), 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; {error, Reason} -> 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. @@ -116,6 +116,6 @@ normalize_image(Image) when is_binary(Image) -> -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) -> - 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]), efka_logger:write(Info). diff --git a/src/docker/docker_manager.erl b/src/docker/docker_manager.erl index 3d0d3e8..80d3c84 100644 --- a/src/docker/docker_manager.erl +++ b/src/docker/docker_manager.erl @@ -208,8 +208,8 @@ handle_info({'DOWN', _Ref, process, TaskPid, Reason}, State = #state{task_map = ok; Error0 -> Error = iolist_to_binary(io_lib:format("~p", [Error0])), - efka_remote_agent:task_event_stream(TaskId, <<"error">>, <<"任务失败: "/utf8, Error/binary>>), - efka_remote_agent:close_task_event_stream(TaskId, <<"task exited">>), + efka_client:task_event_stream(TaskId, <<"error">>, <<"任务失败: "/utf8, Error/binary>>), + efka_client:close_task_event_stream(TaskId, <<"task exited">>), logger:notice("[docker_manager] task_id: ~p, exit with error: ~p", [TaskId, Error]), ok end, diff --git a/src/efka_remote_agent.erl b/src/efka_client.erl similarity index 65% rename from src/efka_remote_agent.erl rename to src/efka_client.erl index 8d3e3c2..696c57e 100644 --- a/src/efka_remote_agent.erl +++ b/src/efka_client.erl @@ -4,9 +4,9 @@ %%% @doc %%% %%% @end -%%% Created : 21. 5月 2025 18:38 +%%% Created : 20. 4月 2026 00:00 %%%------------------------------------------------------------------- --module(efka_remote_agent). +-module(efka_client). -author("anlicheng"). -include("message.hrl"). -include("message_pb.hrl"). @@ -25,7 +25,6 @@ %% 标记当前agent的状态,只有在 activated 状态下才可以正常的发送数据 -define(STATE_DENIED, denied). --define(STATE_CONNECTING, connecting). -define(STATE_AUTH, auth). %% 不能推送消息到服务,但是可以接受服务器的部分指令 -define(STATE_RESTRICTED, restricted). @@ -33,8 +32,7 @@ -define(STATE_ACTIVATED, activated). -record(state, { - transport_pid :: undefined | pid(), - transport_ref :: undefined | reference() + socket :: undefined | ssl:sslsocket() }). %%%=================================================================== @@ -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) -> 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() -> gen_statem:start_link({local, ?SERVER}, ?MODULE, [], []). @@ -67,31 +62,19 @@ start_link() -> %%% 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([]) -> 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() -> 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 -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'{ body = {data, #'Data'{route_key = RouteKey, metric = Metric}} }), - efka_transport:send(TransportPid, Packet), + send_packet(Socket, Packet), {keep_state, 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), {keep_state, State}; -handle_event(cast, {task_event_stream, TaskId, Type, Stream}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) -> - logger:debug("[efka_remote_agent] event_stream task_id: ~p, stream: ~ts", [TaskId, Stream]), +handle_event(cast, {task_event_stream, TaskId, Type, Stream}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> + logger:debug("[efka_client] event_stream task_id: ~p, stream: ~ts", [TaskId, Stream]), EventPacket = message_pb:encode_msg(#'CastFrame'{ body = {event_stream, #'TaskEventStream'{task_id = TaskId, type = Type, stream = Stream}} }), - efka_transport:send(TransportPid, EventPacket), + send_packet(Socket, EventPacket), {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'{ body = {event_stream, #'TaskEventStream'{task_id = TaskId, type = <<"close">>, stream = Reason}} }), - efka_transport:send(TransportPid, EventPacket), + send_packet(Socket, EventPacket), {keep_state, State}; %% 其他情况下直接忽略 handle_event(cast, {task_event_stream, _TaskId, _Stream}, _, 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) -> - {ok, Props} = application:get_env(efka, tls_server), - Host = proplists:get_value(host, Props), - 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 -> + case connect_socket() of + {ok, Socket} -> AuthPacket = auth_packet(), - efka_transport:send(TransportPid, AuthPacket), - {next_state, ?STATE_AUTH, State, [{state_timeout, 5000, auth_timeout}]}; + send_packet(Socket, AuthPacket), + {next_state, ?STATE_AUTH, State#state{socket = Socket}, [{state_timeout, 5000, auth_timeout}]}; {error, Reason} -> - logger:debug("[efka_remote_agent] connect failed, error: ~p, pid: ~p", [Reason, TransportPid]), - efka_transport:stop(TransportPid), - {next_state, ?STATE_DENIED, State#state{transport_pid = undefined}} + logger:debug("[efka_client] connect failed, error: ~p", [Reason]), + schedule_reconnect(), + {keep_state, State#state{socket = undefined}} 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}}} = message_pb:decode_msg(PacketBin, 'ResponseFrame'), case Code of 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}]}; 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}; 2 -> - logger:debug("[efka_remote_agent] auth failed, message: ~p", [Message]), - efka_transport:stop(TransportPid), - {next_state, ?STATE_DENIED, State#state{transport_pid = undefined}}; + logger:debug("[efka_client] auth failed, message: ~p", [Message]), + disconnect(Socket), + schedule_reconnect(), + {next_state, ?STATE_DENIED, State#state{socket = undefined}}; _ -> - logger:debug("[efka_remote_agent] auth failed, invalid message"), - efka_transport:stop(TransportPid), - {next_state, ?STATE_DENIED, State#state{transport_pid = undefined}} + logger:debug("[efka_client] auth failed, invalid message"), + disconnect(Socket), + schedule_reconnect(), + {next_state, ?STATE_DENIED, State#state{socket = undefined}} end; -handle_event(state_timeout, auth_timeout, ?STATE_AUTH, State = #state{transport_pid = TransportPid}) -> - logger:debug("[efka_remote_agent] auth request timeout"), - efka_transport:stop(TransportPid), - {next_state, ?STATE_DENIED, State#state{transport_pid = undefined}}; +handle_event(state_timeout, auth_timeout, ?STATE_AUTH, State = #state{socket = Socket}) -> + logger:debug("[efka_client] auth request timeout"), + disconnect(Socket), + 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 {ok, {Id, Packet}} -> - efka_transport:send(TransportPid, Packet), + send_packet(Socket, Packet), cache_model:delete(Id), {keep_state, State, [{next_event, info, flush_cache}]}; error -> @@ -201,7 +159,7 @@ handle_event(info, flush_cache, ?STATE_ACTIVATED, State = #state{transport_pid = handle_event(info, flush_cache, _, 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 -> #'RequestFrame'{packet_id = PacketId, body = Body} = message_pb:decode_msg(PacketBin, 'RequestFrame'), true = is_integer(PacketId) andalso PacketId > 0, @@ -211,25 +169,22 @@ handle_event(info, {server_packet, <<8, _/binary>> = PacketBin}, StateName, Stat {container_request, Request} -> {keep_state, State, [{next_event, info, {container_request, PacketId, Request}}]} 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'), {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'), {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'), {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'), {keep_state, State, [{next_event, info, {server_cast, Command}}]}; -%% 云端服务器推送了消息 -%% 激活消息 - %% 微服务部署 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 {ok, Containers} -> 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)}} }} }), - efka_transport:send(TransportPid, Packet); + send_packet(Socket, Packet); {error, Reason} when is_binary(Reason) -> Packet = message_pb:encode_msg(#'ResponseFrame'{ packet_id = PacketId, @@ -246,14 +201,14 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {l reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}} }} }), - efka_transport:send(TransportPid, Packet) + send_packet(Socket, Packet) end, {keep_state, State}; handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {deploy, #'ContainerRequest.Deploy'{ task_id = TaskId, params = Params -}}}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) -> +}}}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> case docker_manager:deploy(TaskId, Params) of ok -> 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">>)}} }} }), - efka_transport:send(TransportPid, Packet); + send_packet(Socket, Packet); {error, Reason} when is_binary(Reason) -> Packet = message_pb:encode_msg(#'ResponseFrame'{ packet_id = PacketId, @@ -270,12 +225,12 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {d reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}} }} }), - efka_transport:send(TransportPid, Packet) + send_packet(Socket, Packet) end, {keep_state, State}; 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), case docker_manager:start_container(ContainerTarget) of ok -> @@ -285,7 +240,7 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {s reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(<<"ok">>)}} }} }), - efka_transport:send(TransportPid, Packet); + send_packet(Socket, Packet); {error, Reason} when is_binary(Reason) -> Packet = message_pb:encode_msg(#'ResponseFrame'{ packet_id = PacketId, @@ -293,14 +248,14 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {s reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}} }} }), - efka_transport:send(TransportPid, Packet) + send_packet(Socket, Packet) end, {keep_state, State}; handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {stop, #'ContainerRequest.Stop'{ target = Target, timeout_seconds = TimeoutSeconds -}}}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) -> +}}}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> ContainerTarget = container_target(Target), case docker_manager:stop_container(ContainerTarget, TimeoutSeconds) of ok -> @@ -310,7 +265,7 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {s reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(<<"ok">>)}} }} }), - efka_transport:send(TransportPid, Packet); + send_packet(Socket, Packet); {error, Reason} when is_binary(Reason) -> Packet = message_pb:encode_msg(#'ResponseFrame'{ packet_id = PacketId, @@ -318,14 +273,14 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {s reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}} }} }), - efka_transport:send(TransportPid, Packet) + send_packet(Socket, Packet) end, {keep_state, State}; handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {kill, #'ContainerRequest.Kill'{ target = Target, signal = Signal -}}}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) -> +}}}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> ContainerTarget = container_target(Target), case docker_manager:kill_container(ContainerTarget, to_binary(Signal)) of ok -> @@ -335,7 +290,7 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {k reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(<<"ok">>)}} }} }), - efka_transport:send(TransportPid, Packet); + send_packet(Socket, Packet); {error, Reason} when is_binary(Reason) -> Packet = message_pb:encode_msg(#'ResponseFrame'{ packet_id = PacketId, @@ -343,7 +298,7 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {k reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}} }} }), - efka_transport:send(TransportPid, Packet) + send_packet(Socket, Packet) end, {keep_state, State}; @@ -351,7 +306,7 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {r target = Target, force = Force, remove_volumes = RemoveVolumes -}}}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) -> +}}}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> ContainerTarget = container_target(Target), case docker_manager:remove_container(ContainerTarget, to_bool(Force), to_bool(RemoveVolumes)) of ok -> @@ -361,7 +316,7 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {r reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(<<"ok">>)}} }} }), - efka_transport:send(TransportPid, Packet); + send_packet(Socket, Packet); {error, Reason} when is_binary(Reason) -> Packet = message_pb:encode_msg(#'ResponseFrame'{ packet_id = PacketId, @@ -369,14 +324,14 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {r reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}} }} }), - efka_transport:send(TransportPid, Packet) + send_packet(Socket, Packet) end, {keep_state, State}; handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {config, #'ContainerRequest.Config'{ target = Target, config = Config -}}}}, ?STATE_ACTIVATED, State = #state{transport_pid = TransportPid}) -> +}}}}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> ContainerTarget = container_target(Target), case docker_manager:config_container(ContainerTarget, iolist_to_binary(Config)) of ok -> @@ -386,7 +341,7 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {c reply = {result, #'RpcReply.RpcResult'{data = encode_rpc_payload(<<"ok">>)}} }} }), - efka_transport:send(TransportPid, Packet); + send_packet(Socket, Packet); {error, Reason} -> Packet = message_pb:encode_msg(#'ResponseFrame'{ packet_id = PacketId, @@ -394,111 +349,89 @@ handle_event(info, {container_request, PacketId, #'ContainerRequest'{action = {c reply = {error, #'RpcReply.RpcError'{code = -1, message = Reason}} }} }), - efka_transport:send(TransportPid, Packet) + send_packet(Socket, Packet) end, {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_id = PacketId, body = {rpc_reply, #'RpcReply'{ reply = {error, #'RpcReply.RpcError'{code = -1, message = <<"unsupported container request">>}} }} }), - efka_transport:send(TransportPid, Packet), + send_packet(Socket, Packet), {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_id = PacketId, body = {rpc_reply, #'RpcReply'{ reply = {error, #'RpcReply.RpcError'{code = -1, message = <<"agent restricted">>}} }} }), - efka_transport:send(TransportPid, Packet), + send_packet(Socket, Packet), {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_id = PacketId, body = {rpc_reply, #'RpcReply'{ reply = {error, #'RpcReply.RpcError'{code = -1, message = <<"unsupported rpc request: ", Method/binary>>}} }} }), - efka_transport:send(TransportPid, Packet), + send_packet(Socket, Packet), {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_id = PacketId, body = {rpc_reply, #'RpcReply'{ reply = {error, #'RpcReply.RpcError'{code = -1, message = <<"agent restricted">>}} }} }), - efka_transport:send(TransportPid, Packet), + send_packet(Socket, Packet), {keep_state, State}; -%% 处理task_log -%handle_event(info, {server_async_call, PacketId, <>}, ?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), case {Auth, StateName} of {1, ?STATE_ACTIVATED} -> {keep_state, State}; - {1, ?STATE_DENIED} -> - %% 重新激活, 需要重新校验 + {1, _} -> AuthPacket = auth_packet(), - efka_transport:send(TransportPid, AuthPacket), + send_packet(Socket, AuthPacket), {next_state, ?STATE_AUTH, State, [{state_timeout, 5000, auth_timeout}]}; {0, _} -> - %% 这个时候的主机应该是受限制的状态,不允许发送消息;但是能够接受服务器推送的消息 {next_state, ?STATE_RESTRICTED, State} end; %% 处理Pub/Sub机制 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), {keep_state, State}; -%% transport进程退出 -handle_event(info, {'DOWN', MRef, process, TransportPid, Reason}, _, State = #state{transport_ref = MRef}) -> - logger:debug("[efka_remote_agent] transport pid: ~p, exit with reason: ~p", [TransportPid, Reason]), - erlang:start_timer(5000, self(), create_transport), - {next_state, ?STATE_DENIED, State#state{transport_pid = undefined, transport_ref = undefined}}. +handle_event(info, {ssl_error, Socket, Reason}, _, State = #state{socket = Socket}) -> + logger:debug("[efka_client] ssl error: ~p", [Reason]), + disconnect(Socket), + schedule_reconnect(), + {next_state, ?STATE_DENIED, State#state{socket = undefined}}; -%% @private -%% @doc This function is called by a gen_statem 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_statem terminates with -%% Reason. The return value is ignored. -terminate(_Reason, _StateName, _State = #state{transport_pid = TransportPid}) -> - case is_pid(TransportPid) andalso is_process_alive(TransportPid) of - true -> - efka_transport:stop(TransportPid); - false -> - ok - end, +handle_event(info, {ssl_closed, Socket}, _, State = #state{socket = Socket}) -> + schedule_reconnect(), + {next_state, ?STATE_DENIED, State#state{socket = undefined}}; + +handle_event(info, Info, _, State = #state{}) -> + logger:notice("[efka_client] get unknown info: ~p", [Info]), + {keep_state, State}. + +terminate(_Reason, _StateName, _State = #state{socket = Socket}) -> + disconnect(Socket), ok. -%% @private -%% @doc Convert process state when code is changed code_change(_OldVsn, StateName, State = #state{}, _Extra) -> {ok, StateName, State}. @@ -508,7 +441,6 @@ code_change(_OldVsn, StateName, State = #state{}, _Extra) -> -spec auth_packet() -> binary(). auth_packet() -> - %% 重新激活, 需要重新校验 {ok, AuthInfo} = application:get_env(efka, auth), UUID = proplists:get_value(uuid, 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(). encode_rpc_payload(Payload) -> jiffy:encode(Payload, [force_utf8]). diff --git a/src/efka_service.erl b/src/efka_service.erl index 70ad1ab..0080063 100644 --- a/src/efka_service.erl +++ b/src/efka_service.erl @@ -104,7 +104,7 @@ handle_call(_Request, _From, State = #state{}) -> {stop, Reason :: term(), NewState :: #state{}}). 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]), - efka_remote_agent:metric_data(RouteKey, Metric), + efka_client:metric_data(RouteKey, Metric), {noreply, State}; handle_cast(_Request, State = #state{}) -> diff --git a/src/efka_sup.erl b/src/efka_sup.erl index 37b1e52..0094370 100644 --- a/src/efka_sup.erl +++ b/src/efka_sup.erl @@ -92,12 +92,12 @@ init([]) -> }, #{ - id => 'efka_remote_agent', - start => {'efka_remote_agent', start_link, []}, + id => 'efka_client', + start => {'efka_client', start_link, []}, restart => permanent, shutdown => 2000, type => worker, - modules => ['efka_remote_agent'] + modules => ['efka_client'] } ], @@ -105,4 +105,3 @@ init([]) -> {ok, {SupFlags, ChildSpecs}}. %% internal functions - diff --git a/src/efka_transport.erl b/src/efka_transport.erl deleted file mode 100644 index 22ded0a..0000000 --- a/src/efka_transport.erl +++ /dev/null @@ -1,153 +0,0 @@ -%%%------------------------------------------------------------------- -%%% @author anlicheng -%%% @copyright (C) 2025, -%%% @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).