From 97f9bc7a550db34019f234d4f46fd054f8bcc157 Mon Sep 17 00:00:00 2001 From: anlicheng <244108715@qq.com> Date: Fri, 24 Apr 2026 10:50:26 +0800 Subject: [PATCH] fix cache --- src/transport/efka_client.erl | 78 ++++++++++++++++++++++++++++++----- 1 file changed, 67 insertions(+), 11 deletions(-) diff --git a/src/transport/efka_client.erl b/src/transport/efka_client.erl index 879ac67..2ef0173 100644 --- a/src/transport/efka_client.erl +++ b/src/transport/efka_client.erl @@ -17,13 +17,15 @@ %% API -export([start_link/0]). -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 -export([init/1, handle_event/4, terminate/3, code_change/4, callback_mode/0]). -define(SERVER, ?MODULE). -define(CACHE_TAB, cache). +-define(MAX_CACHE_ITEMS, 5000000). +-define(MAX_CACHE_FILE_SIZE, 1073741824). %% 标记当前agent的状态,只有在 activated 状态下才可以正常的发送数据 -define(STATE_DISCONNECTED, disconnected). @@ -37,7 +39,8 @@ socket :: undefined | ssl:sslsocket(), next_packet_id = 1, %% 保存当前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() -> 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. 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 ?STATE_ACTIVATED -> - send_packet(Socket, [?FRAME_CAST, CastFrame]); + send_packet(Socket, [?FRAME_CAST, CastFrame]), + {keep_state, State}; _ -> - ok = cache_insert(<>) - end, - {keep_state, State}; + {ok, DroppedCount} = cache_insert(<>), + {keep_state, State#state{dropped_message_count = State#state.dropped_message_count + DroppedCount}} + end; %% Task的stream流,只做实时的 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}]}; handle_event({call, From}, is_activated, _StateName, State = #state{}) -> {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}) -> @@ -327,9 +337,14 @@ close_cache_table() -> ok 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) -> - 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()}}. cache_fetch_next() -> @@ -337,8 +352,12 @@ cache_fetch_next() -> '$end_of_table' -> error; Key -> - [Entry] = dets:lookup(?CACHE_TAB, Key), - {ok, Entry} + case dets:lookup(?CACHE_TAB, Key) of + [Entry | _] -> + {ok, Entry}; + [] -> + cache_fetch_next() + end end. -spec cache_delete(integer()) -> ok | {error, term()}. @@ -347,7 +366,44 @@ cache_delete(Id) when is_integer(Id) -> -spec generate_cache_id() -> integer(). 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. send_result_reply(Socket, PacketId, Payload) when is_binary(Payload) ->