fix efka_client

This commit is contained in:
anlicheng 2026-04-20 17:19:02 +08:00
parent 045b32d3a1
commit db5a5e9071
2 changed files with 18 additions and 31 deletions

View File

@ -17,7 +17,7 @@
%% API %% API
-export([start_link/0]). -export([start_link/0]).
-export([metric_data/2, ping/13, task_event_stream/3, close_task_event_stream/2]). -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]). -export([is_activated/0]).
%% gen_statem callbacks %% gen_statem callbacks
-export([init/1, handle_event/4, terminate/3, code_change/4, callback_mode/0]). -export([init/1, handle_event/4, terminate/3, code_change/4, callback_mode/0]).
@ -56,13 +56,9 @@ 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) -> close_task_event_stream(TaskId, Reason) when is_integer(TaskId), is_binary(Reason) ->
gen_statem:cast(?SERVER, {close_task_event_stream, TaskId, 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. -spec is_activated() -> boolean().
send_task_event_stream(TaskId, Type, Stream) when is_integer(TaskId), is_binary(Type), is_binary(Stream) -> is_activated() ->
gen_statem:call(?SERVER, {send_task_event_stream, TaskId, Type, Stream}). gen_statem:call(?SERVER, is_activated).
-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) -> 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}).
@ -118,23 +114,10 @@ handle_event(cast, {task_event_stream, _TaskId, _Stream}, _, State = #state{}) -
handle_event(cast, {close_task_event_stream, _TaskId, _Reason}, _, State = #state{}) -> handle_event(cast, {close_task_event_stream, _TaskId, _Reason}, _, State = #state{}) ->
{keep_state, State}; {keep_state, State};
handle_event({call, From}, {send_task_event_stream, TaskId, Type, Stream}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> handle_event({call, From}, is_activated, ?STATE_ACTIVATED, State = #state{}) ->
EventPacket = message_pb:encode_msg(#'CastFrame'{ {keep_state, State, [{reply, From, true}]};
body = {event_stream, #'TaskEventStream'{task_id = TaskId, type = Type, stream = Stream}} handle_event({call, From}, is_activated, _StateName, State = #state{}) ->
}), {keep_state, State, [{reply, From, false}]};
send_packet(Socket, <<?FRAME_CAST, EventPacket/binary>>),
{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, <<?FRAME_CAST, EventPacket/binary>>),
{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}) -> handle_event(info, {timeout, _, create_transport}, ?STATE_DISCONNECTED, State = #state{next_packet_id = PacketId}) ->

View File

@ -97,11 +97,15 @@ flush_pending(State = #state{pending = Pending0}) ->
end. end.
send_event({stream, TaskId, Type, Stream}) -> send_event({stream, TaskId, Type, Stream}) ->
send_result(catch efka_client:send_task_event_stream(TaskId, Type, Stream)); maybe_send(fun() -> efka_client:task_event_stream(TaskId, Type, Stream) end);
send_event({close, TaskId, Reason}) -> send_event({close, TaskId, Reason}) ->
send_result(catch efka_client:send_close_task_event_stream(TaskId, Reason)). maybe_send(fun() -> efka_client:close_task_event_stream(TaskId, Reason) end).
send_result(ok) -> maybe_send(SendFun) ->
case catch efka_client:is_activated() of
true ->
_ = catch SendFun(),
ok; ok;
send_result(_) -> _ ->
not_ready. not_ready
end.