diff --git a/src/docker/docker_deployer.erl b/src/docker/docker_deployer.erl index 69eac0c..4065884 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_client:close_task_event_stream(TaskId, ?TASK_FAIL); + efka_task_reporter:close(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_client:close_task_event_stream(TaskId, ?TASK_SUCCESS); + efka_task_reporter:close(TaskId, ?TASK_SUCCESS); {error, Reason} -> trace_log(TaskId, <<"error">>, <<"容器创建失败: "/utf8, Reason/binary>>), trace_log(TaskId, <<"error">>, <<"任务失败"/utf8>>), - efka_client:close_task_event_stream(TaskId, ?TASK_FAIL) + efka_task_reporter:close(TaskId, ?TASK_FAIL) end; {error, Reason} -> trace_log(TaskId, <<"error">>, <<"镜像拉取失败: "/utf8, Reason/binary>>), - efka_client:close_task_event_stream(TaskId, ?TASK_FAIL) + efka_task_reporter:close(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_client:task_event_stream(TaskId, Level, Msg), + efka_task_reporter: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 80d3c84..2c792ed 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_client:task_event_stream(TaskId, <<"error">>, <<"任务失败: "/utf8, Error/binary>>), - efka_client:close_task_event_stream(TaskId, <<"task exited">>), + efka_task_reporter:stream(TaskId, <<"error">>, <<"任务失败: "/utf8, Error/binary>>), + efka_task_reporter:close(TaskId, <<"task exited">>), logger:notice("[docker_manager] task_id: ~p, exit with error: ~p", [TaskId, Error]), ok end, diff --git a/src/efka_client.erl b/src/efka_client.erl index 5edd6fc..c0c7b06 100644 --- a/src/efka_client.erl +++ b/src/efka_client.erl @@ -17,6 +17,7 @@ %% API -export([start_link/0]). -export([metric_data/2, ping/13, task_event_stream/3, close_task_event_stream/2]). +-export([send_task_event_stream/3, send_close_task_event_stream/2]). %% gen_statem callbacks -export([init/1, handle_event/4, terminate/3, code_change/4, callback_mode/0]). @@ -55,6 +56,14 @@ task_event_stream(TaskId, Type, Stream) when is_integer(TaskId), is_binary(Type) close_task_event_stream(TaskId, Reason) when is_integer(TaskId), is_binary(Reason) -> gen_statem:cast(?SERVER, {close_task_event_stream, TaskId, Reason}). +-spec send_task_event_stream(TaskId :: integer(), Type :: binary(), Stream :: binary()) -> ok | not_ready. +send_task_event_stream(TaskId, Type, Stream) when is_integer(TaskId), is_binary(Type), is_binary(Stream) -> + gen_statem:call(?SERVER, {send_task_event_stream, TaskId, Type, Stream}). + +-spec send_close_task_event_stream(TaskId :: integer(), Reason :: binary()) -> ok | not_ready. +send_close_task_event_stream(TaskId, Reason) when is_integer(TaskId), is_binary(Reason) -> + gen_statem:call(?SERVER, {send_close_task_event_stream, TaskId, Reason}). + 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}). @@ -106,6 +115,26 @@ handle_event(cast, {close_task_event_stream, TaskId, Reason}, ?STATE_ACTIVATED, %% 其他情况下直接忽略 handle_event(cast, {task_event_stream, _TaskId, _Stream}, _, State = #state{}) -> {keep_state, State}; +handle_event(cast, {close_task_event_stream, _TaskId, _Reason}, _, State = #state{}) -> + {keep_state, State}; + +handle_event({call, From}, {send_task_event_stream, TaskId, Type, Stream}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> + EventPacket = message_pb:encode_msg(#'CastFrame'{ + body = {event_stream, #'TaskEventStream'{task_id = TaskId, type = Type, stream = Stream}} + }), + send_packet(Socket, <>), + {keep_state, State, [{reply, From, ok}]}; +handle_event({call, From}, {send_task_event_stream, _TaskId, _Type, _Stream}, _StateName, State = #state{}) -> + {keep_state, State, [{reply, From, not_ready}]}; + +handle_event({call, From}, {send_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}} + }), + send_packet(Socket, <>), + {keep_state, State, [{reply, From, ok}]}; +handle_event({call, From}, {send_close_task_event_stream, _TaskId, _Reason}, _StateName, State = #state{}) -> + {keep_state, State, [{reply, From, not_ready}]}; %% 异步建立到服务器的连接 handle_event(info, {timeout, _, create_transport}, ?STATE_DISCONNECTED, State = #state{next_packet_id = PacketId}) -> @@ -302,10 +331,10 @@ handle_container_request(#'ContainerRequest'{action = {deploy, #'ContainerReques handle_container_request(#'ContainerRequest'{action = {start, #'ContainerRequest.Start'{target = Target}}}) -> ContainerTarget = container_target(Target), docker_manager:start_container(ContainerTarget); -handle_container_request(#'ContainerRequest'{action = {stop, #'ContainerRequest.Stop'{ target = Target, timeout_seconds = TimeoutSeconds }}}) -> +handle_container_request(#'ContainerRequest'{action = {stop, #'ContainerRequest.Stop'{target = Target, timeout_seconds = TimeoutSeconds }}}) -> ContainerTarget = container_target(Target), docker_manager:stop_container(ContainerTarget, TimeoutSeconds); -handle_container_request(#'ContainerRequest'{action = {kill, #'ContainerRequest.Kill'{ target = Target, signal = Signal}}}) -> +handle_container_request(#'ContainerRequest'{action = {kill, #'ContainerRequest.Kill'{target = Target, signal = Signal}}}) -> ContainerTarget = container_target(Target), docker_manager:kill_container(ContainerTarget, to_binary(Signal)); handle_container_request(#'ContainerRequest'{action = {remove, #'ContainerRequest.Remove'{target = Target, force = Force, remove_volumes = RemoveVolumes}}}) -> diff --git a/src/efka_sup.erl b/src/efka_sup.erl index 0094370..0d1a6c6 100644 --- a/src/efka_sup.erl +++ b/src/efka_sup.erl @@ -82,15 +82,6 @@ init([]) -> modules => ['efka_subscription'] }, - #{ - id => 'docker_manager', - start => {'docker_manager', start_link, []}, - restart => permanent, - shutdown => 2000, - type => worker, - modules => ['docker_manager'] - }, - #{ id => 'efka_client', start => {'efka_client', start_link, []}, @@ -98,6 +89,24 @@ init([]) -> shutdown => 2000, type => worker, modules => ['efka_client'] + }, + + #{ + id => 'efka_task_reporter', + start => {'efka_task_reporter', start_link, []}, + restart => permanent, + shutdown => 2000, + type => worker, + modules => ['efka_task_reporter'] + }, + + #{ + id => 'docker_manager', + start => {'docker_manager', start_link, []}, + restart => permanent, + shutdown => 2000, + type => worker, + modules => ['docker_manager'] } ], diff --git a/src/efka_task_reporter.erl b/src/efka_task_reporter.erl new file mode 100644 index 0000000..dfc9875 --- /dev/null +++ b/src/efka_task_reporter.erl @@ -0,0 +1,107 @@ +%%%------------------------------------------------------------------- +%%% @author anlicheng +%%% @copyright (C) 2026, +%%% @doc +%%% +%%% @end +%%% Created : 20. 4月 2026 +%%%------------------------------------------------------------------- +-module(efka_task_reporter). +-author("anlicheng"). + +-behaviour(gen_server). + +%% API +-export([start_link/0]). +-export([stream/3, close/2]). + +%% gen_server callbacks +-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]). + +-define(SERVER, ?MODULE). +-define(FLUSH_INTERVAL, 1000). + +-record(state, { + pending = queue:new(), + flush_ref = undefined +}). + +%%%=================================================================== +%%% API +%%%=================================================================== + +-spec stream(TaskId :: integer(), Type :: binary(), Stream :: binary()) -> ok. +stream(TaskId, Type, Stream) when is_integer(TaskId), is_binary(Type), is_binary(Stream) -> + gen_server:cast(?SERVER, {stream, TaskId, Type, Stream}). + +-spec close(TaskId :: integer(), Reason :: binary()) -> ok. +close(TaskId, Reason) when is_integer(TaskId), is_binary(Reason) -> + gen_server:cast(?SERVER, {close, TaskId, Reason}). + +-spec start_link() -> {ok, pid()} | ignore | {error, term()}. +start_link() -> + gen_server:start_link({local, ?SERVER}, ?MODULE, [], []). + +%%%=================================================================== +%%% gen_server callbacks +%%%=================================================================== + +init([]) -> + {ok, #state{}}. + +handle_call(_Request, _From, State) -> + {reply, ok, State}. + +handle_cast({stream, TaskId, Type, Stream}, State0 = #state{pending = Pending0}) -> + Pending = queue:in({stream, TaskId, Type, Stream}, Pending0), + {noreply, ensure_flush_timer(State0#state{pending = Pending}, 0)}; +handle_cast({close, TaskId, Reason}, State0 = #state{pending = Pending0}) -> + Pending = queue:in({close, TaskId, Reason}, Pending0), + {noreply, ensure_flush_timer(State0#state{pending = Pending}, 0)}; +handle_cast(_Request, State) -> + {noreply, State}. + +handle_info(flush, State0 = #state{}) -> + State1 = State0#state{flush_ref = undefined}, + {noreply, flush_pending(State1)}; +handle_info(_Info, State) -> + {noreply, State}. + +terminate(_Reason, _State) -> + ok. + +code_change(_OldVsn, State, _Extra) -> + {ok, State}. + +%%%=================================================================== +%%% Internal functions +%%%=================================================================== + +ensure_flush_timer(State = #state{flush_ref = undefined}, Delay) when is_integer(Delay), Delay >= 0 -> + Ref = erlang:send_after(Delay, self(), flush), + State#state{flush_ref = Ref}; +ensure_flush_timer(State = #state{}, _Delay) -> + State. + +flush_pending(State = #state{pending = Pending0}) -> + case queue:out(Pending0) of + {empty, _} -> + State; + {{value, Event}, Pending1} -> + case send_event(Event) of + ok -> + flush_pending(State#state{pending = Pending1}); + not_ready -> + ensure_flush_timer(State#state{pending = Pending0}, ?FLUSH_INTERVAL) + end + end. + +send_event({stream, TaskId, Type, Stream}) -> + send_result(catch efka_client:send_task_event_stream(TaskId, Type, Stream)); +send_event({close, TaskId, Reason}) -> + send_result(catch efka_client:send_close_task_event_stream(TaskId, Reason)). + +send_result(ok) -> + ok; +send_result(_) -> + not_ready.