fix cache
This commit is contained in:
parent
2e5ddbe087
commit
97f9bc7a55
@ -17,13 +17,15 @@
|
|||||||
%% 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([is_activated/0]).
|
-export([is_activated/0, dropped_message_count/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]).
|
||||||
|
|
||||||
-define(SERVER, ?MODULE).
|
-define(SERVER, ?MODULE).
|
||||||
-define(CACHE_TAB, cache).
|
-define(CACHE_TAB, cache).
|
||||||
|
-define(MAX_CACHE_ITEMS, 5000000).
|
||||||
|
-define(MAX_CACHE_FILE_SIZE, 1073741824).
|
||||||
|
|
||||||
%% 标记当前agent的状态,只有在 activated 状态下才可以正常的发送数据
|
%% 标记当前agent的状态,只有在 activated 状态下才可以正常的发送数据
|
||||||
-define(STATE_DISCONNECTED, disconnected).
|
-define(STATE_DISCONNECTED, disconnected).
|
||||||
@ -37,7 +39,8 @@
|
|||||||
socket :: undefined | ssl:sslsocket(),
|
socket :: undefined | ssl:sslsocket(),
|
||||||
next_packet_id = 1,
|
next_packet_id = 1,
|
||||||
%% 保存当前auth请求的packet_id,用来建立auth请求和响应的对应关系
|
%% 保存当前auth请求的packet_id,用来建立auth请求和响应的对应关系
|
||||||
auth_packet_id = 1
|
auth_packet_id = 1,
|
||||||
|
dropped_message_count = 0 :: non_neg_integer()
|
||||||
}).
|
}).
|
||||||
|
|
||||||
%%%===================================================================
|
%%%===================================================================
|
||||||
@ -61,6 +64,10 @@ close_task_event_stream(TaskId, Reason) when is_integer(TaskId), is_binary(Reaso
|
|||||||
is_activated() ->
|
is_activated() ->
|
||||||
gen_statem:call(?SERVER, is_activated).
|
gen_statem:call(?SERVER, is_activated).
|
||||||
|
|
||||||
|
-spec dropped_message_count() -> non_neg_integer().
|
||||||
|
dropped_message_count() ->
|
||||||
|
gen_statem:call(?SERVER, dropped_message_count).
|
||||||
|
|
||||||
-spec ping(term(), term(), term(), term(), term(), term(), term(), term(), term(), term(), term(), term(), term()) -> ok.
|
-spec ping(term(), term(), term(), term(), term(), term(), term(), term(), term(), term(), term(), term(), term()) -> ok.
|
||||||
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}).
|
||||||
@ -95,11 +102,12 @@ handle_event(cast, {metric_data, RouteKey, Metric}, StateName, State = #state{so
|
|||||||
}),
|
}),
|
||||||
case StateName of
|
case StateName of
|
||||||
?STATE_ACTIVATED ->
|
?STATE_ACTIVATED ->
|
||||||
send_packet(Socket, [?FRAME_CAST, CastFrame]);
|
send_packet(Socket, [?FRAME_CAST, CastFrame]),
|
||||||
_ ->
|
|
||||||
ok = cache_insert(<<?FRAME_CAST, CastFrame/binary>>)
|
|
||||||
end,
|
|
||||||
{keep_state, State};
|
{keep_state, State};
|
||||||
|
_ ->
|
||||||
|
{ok, DroppedCount} = cache_insert(<<?FRAME_CAST, CastFrame/binary>>),
|
||||||
|
{keep_state, State#state{dropped_message_count = State#state.dropped_message_count + DroppedCount}}
|
||||||
|
end;
|
||||||
|
|
||||||
%% Task的stream流,只做实时的
|
%% Task的stream流,只做实时的
|
||||||
handle_event(cast, {task_event_stream, TaskId, Type, Stream}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
handle_event(cast, {task_event_stream, TaskId, Type, Stream}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
|
||||||
@ -125,6 +133,8 @@ handle_event({call, From}, is_activated, ?STATE_ACTIVATED, State = #state{}) ->
|
|||||||
{keep_state, State, [{reply, From, true}]};
|
{keep_state, State, [{reply, From, true}]};
|
||||||
handle_event({call, From}, is_activated, _StateName, State = #state{}) ->
|
handle_event({call, From}, is_activated, _StateName, State = #state{}) ->
|
||||||
{keep_state, State, [{reply, From, false}]};
|
{keep_state, State, [{reply, From, false}]};
|
||||||
|
handle_event({call, From}, dropped_message_count, _StateName, State = #state{dropped_message_count = DroppedCount}) ->
|
||||||
|
{keep_state, State, [{reply, From, DroppedCount}]};
|
||||||
|
|
||||||
%% 异步建立到服务器的连接
|
%% 异步建立到服务器的连接
|
||||||
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}) ->
|
||||||
@ -327,9 +337,14 @@ close_cache_table() ->
|
|||||||
ok
|
ok
|
||||||
end.
|
end.
|
||||||
|
|
||||||
-spec cache_insert(binary()) -> ok | {error, term()}.
|
-spec cache_insert(binary()) -> {ok, non_neg_integer()} | {error, term()}.
|
||||||
cache_insert(Data) when is_binary(Data) ->
|
cache_insert(Data) when is_binary(Data) ->
|
||||||
dets:insert(?CACHE_TAB, {generate_cache_id(), Data}).
|
case dets:insert(?CACHE_TAB, {generate_cache_id(), Data}) of
|
||||||
|
ok ->
|
||||||
|
trim_cache_limits(0);
|
||||||
|
{error, Reason} ->
|
||||||
|
{error, Reason}
|
||||||
|
end.
|
||||||
|
|
||||||
-spec cache_fetch_next() -> error | {ok, {integer(), binary()}}.
|
-spec cache_fetch_next() -> error | {ok, {integer(), binary()}}.
|
||||||
cache_fetch_next() ->
|
cache_fetch_next() ->
|
||||||
@ -337,8 +352,12 @@ cache_fetch_next() ->
|
|||||||
'$end_of_table' ->
|
'$end_of_table' ->
|
||||||
error;
|
error;
|
||||||
Key ->
|
Key ->
|
||||||
[Entry] = dets:lookup(?CACHE_TAB, Key),
|
case dets:lookup(?CACHE_TAB, Key) of
|
||||||
{ok, Entry}
|
[Entry | _] ->
|
||||||
|
{ok, Entry};
|
||||||
|
[] ->
|
||||||
|
cache_fetch_next()
|
||||||
|
end
|
||||||
end.
|
end.
|
||||||
|
|
||||||
-spec cache_delete(integer()) -> ok | {error, term()}.
|
-spec cache_delete(integer()) -> ok | {error, term()}.
|
||||||
@ -347,7 +366,44 @@ cache_delete(Id) when is_integer(Id) ->
|
|||||||
|
|
||||||
-spec generate_cache_id() -> integer().
|
-spec generate_cache_id() -> integer().
|
||||||
generate_cache_id() ->
|
generate_cache_id() ->
|
||||||
os:system_time(microsecond).
|
erlang:unique_integer([monotonic, positive]).
|
||||||
|
|
||||||
|
-spec cache_over_limit() -> boolean().
|
||||||
|
cache_over_limit() ->
|
||||||
|
cache_item_count() > ?MAX_CACHE_ITEMS orelse cache_file_size() > ?MAX_CACHE_FILE_SIZE.
|
||||||
|
|
||||||
|
-spec cache_item_count() -> non_neg_integer().
|
||||||
|
cache_item_count() ->
|
||||||
|
dets:info(?CACHE_TAB, size).
|
||||||
|
|
||||||
|
-spec cache_file_size() -> non_neg_integer().
|
||||||
|
cache_file_size() ->
|
||||||
|
dets:info(?CACHE_TAB, file_size).
|
||||||
|
|
||||||
|
-spec trim_cache_limits(non_neg_integer()) -> {ok, non_neg_integer()} | {error, term()}.
|
||||||
|
trim_cache_limits(DroppedCount) ->
|
||||||
|
case cache_over_limit() of
|
||||||
|
true ->
|
||||||
|
case delete_oldest_cache_entry() of
|
||||||
|
ok ->
|
||||||
|
trim_cache_limits(DroppedCount + 1);
|
||||||
|
error ->
|
||||||
|
{ok, DroppedCount};
|
||||||
|
{error, Reason} ->
|
||||||
|
{error, Reason}
|
||||||
|
end;
|
||||||
|
false ->
|
||||||
|
{ok, DroppedCount}
|
||||||
|
end.
|
||||||
|
|
||||||
|
-spec delete_oldest_cache_entry() -> ok | error | {error, term()}.
|
||||||
|
delete_oldest_cache_entry() ->
|
||||||
|
case dets:first(?CACHE_TAB) of
|
||||||
|
'$end_of_table' ->
|
||||||
|
error;
|
||||||
|
Key ->
|
||||||
|
dets:delete(?CACHE_TAB, Key)
|
||||||
|
end.
|
||||||
|
|
||||||
-spec send_result_reply(ssl:sslsocket(), integer(), binary()) -> ok.
|
-spec send_result_reply(ssl:sslsocket(), integer(), binary()) -> ok.
|
||||||
send_result_reply(Socket, PacketId, Payload) when is_binary(Payload) ->
|
send_result_reply(Socket, PacketId, Payload) when is_binary(Payload) ->
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user