feat: upsert stats counters, daily buckets and ETS tops for admin metrics. Refs EventHub/EventHubBack#42 #43 #44 #45 #46
CI / test (push) Failing after 25m40s
CI / deploy-ift (push) Has been skipped
CI / e2e-ift (push) Has been skipped
CI / deploy-stage (push) Has been skipped
CI / e2e-stage (push) Has been skipped

This commit is contained in:
2026-07-17 23:10:00 +03:00
parent f2445db66c
commit dd754e28cc
20 changed files with 1885 additions and 408 deletions
+5 -2
View File
@@ -20,7 +20,7 @@
review, report, banned_word, automod_settings, automod_hit,
ticket, subscription,
admin_audit, notification,
stats, node_metric, schema_migration
stats_counter, stats_daily, node_metric, schema_migration
]).
-define(DISC_TABLES, ?TABLES -- [session, verification, admin_session, node_metric]).
@@ -303,7 +303,8 @@ table_opts(ticket) -> [{disc_copies, [node()]}, {attributes, record_info(fields,
table_opts(subscription) -> [{disc_copies, [node()]}, {attributes, record_info(fields, subscription)}];
table_opts(admin_audit) -> [{disc_copies, [node()]}, {attributes, record_info(fields, admin_audit)}];
table_opts(notification) -> [{disc_copies, [node()]}, {attributes, record_info(fields, notification)}];
table_opts(stats) -> [{disc_copies, [node()]}, {attributes, record_info(fields, stats)}];
table_opts(stats_counter) -> [{disc_copies, [node()]}, {attributes, record_info(fields, stats_counter)}];
table_opts(stats_daily) -> [{disc_copies, [node()]}, {attributes, record_info(fields, stats_daily)}];
table_opts(schema_migration) -> [{disc_copies, [node()]}, {attributes, record_info(fields, schema_migration)}];
table_opts(session) -> [{ram_copies, [node()]}, {attributes, record_info(fields, session)}];
table_opts(verification) -> [{ram_copies, [node()]}, {attributes, record_info(fields, verification)}];
@@ -340,4 +341,6 @@ create_indices() ->
mnesia:add_table_index(notification, is_read),
mnesia:add_table_index(auth_session, family_id),
mnesia:add_table_index(auth_session, subject_id),
mnesia:add_table_index(report, resolved_by),
mnesia:add_table_index(ticket, assigned_to),
ok.
+3 -1
View File
@@ -21,7 +21,9 @@
-define(ALL_MIGRATIONS, [
'20260501120000_base_schema',
'20260504150000_test_migration',
'20260716230000_ticket_source_and_hash_index'
'20260716230000_ticket_source_and_hash_index',
'20260717180000_stats_counters',
'20260717190000_admin_stats_indexes'
]).
%% ------------------------------
+542 -106
View File
@@ -1,133 +1,569 @@
%%%-------------------------------------------------------------------
%%% @doc Сбор admin-статистики через Mnesia detailed-subscribe.
%%% ETS держит абсолютные upsert-счётчики и дневные бакеты;
%%% периодический flush пишет в stats_counter / stats_daily.
%%% @end
%%%-------------------------------------------------------------------
-module(stats_collector).
-behaviour(gen_server).
-include("records.hrl").
-export([start_link/0, get_stats/0, subscribe/0]).
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]).
-export([start_link/0, subscribe/0, rebuild/0, is_ready/0, call/1,
get_stats/0]).
-export([init/1, handle_call/3, handle_cast/2, handle_info/2,
terminate/2, code_change/3]).
-define(FLUSH_INTERVAL, 300000).
-define(FLUSH_INTERVAL, 300000). %% 5 min
-define(CLEANUP_INTERVAL, 3600000). %% 1 h
-define(DAILY_RETENTION_DAYS, 730).
-define(COUNTER_ETS, stats_counter_ets).
-define(DAILY_ETS, stats_daily_ets).
-define(SUB_TABLES, [
user, event, calendar, booking, review, report, ticket, subscription
]).
start_link() -> gen_server:start_link({local, ?MODULE}, ?MODULE, [], []).
get_stats() -> gen_server:call(?MODULE, get_stats).
subscribe() -> gen_server:call(?MODULE, subscribe).
-record(state, {
ready = false :: boolean(),
subscribed = false :: boolean()
}).
%%%===================================================================
%%% API
%%%===================================================================
start_link() ->
gen_server:start_link({local, ?MODULE}, ?MODULE, [], []).
subscribe() ->
gen_server:call(?MODULE, subscribe, infinity).
rebuild() ->
gen_server:call(?MODULE, rebuild, infinity).
is_ready() ->
gen_server:call(?MODULE, is_ready).
call(Msg) ->
gen_server:call(?MODULE, Msg).
%% @doc Совместимость: снимок counter ETS.
get_stats() ->
gen_server:call(?MODULE, get_stats).
%%%===================================================================
%%% gen_server
%%%===================================================================
init([]) ->
ets:new(stats_ets, [set, public, named_table, {keypos, 1}]),
ets:new(?COUNTER_ETS, [named_table, public, set, {keypos, 1},
{write_concurrency, true}]),
ets:new(?DAILY_ETS, [named_table, public, set, {keypos, 1},
{write_concurrency, true}]),
stats_tops:init_tables(),
erlang:send_after(?FLUSH_INTERVAL, self(), flush),
{ok, #{}}.
erlang:send_after(?CLEANUP_INTERVAL, self(), cleanup_daily),
{ok, #state{}}.
handle_call(subscribe, _From, State) ->
mnesia:subscribe({table, event, simple}),
mnesia:subscribe({table, booking, simple}),
mnesia:subscribe({table, review, simple}),
{reply, ok, State};
NewState = do_subscribe(State),
{reply, ok, NewState};
handle_call(rebuild, _From, State) ->
do_rebuild(),
{reply, ok, State#state{ready = true}};
handle_call(is_ready, _From, State) ->
{reply, State#state.ready, State};
handle_call(get_stats, _From, State) ->
{reply, ets:tab2list(stats_ets), State};
{reply, ets:tab2list(?COUNTER_ETS), State};
handle_call({get_counter, _Key}, _From, #state{ready = false} = State) ->
{reply, {error, not_ready}, State};
handle_call({get_counter, Key}, _From, State) ->
{reply, {ok, ets_get(?COUNTER_ETS, Key)}, State};
handle_call({get_dim, _E, _F}, _From, #state{ready = false} = State) ->
{reply, {error, not_ready}, State};
handle_call({get_dim, Entity, Field}, _From, State) ->
{reply, {ok, read_dim(Entity, Field)}, State};
handle_call({get_daily, _M, _F, _T}, _From, #state{ready = false} = State) ->
{reply, {error, not_ready}, State};
handle_call({get_daily, Metric, From, To}, _From, State) ->
{reply, {ok, read_daily(Metric, From, To)}, State};
handle_call(_Msg, _From, State) ->
{reply, {error, unknown}, State}.
handle_cast(_Msg, State) -> {noreply, State}.
handle_cast(_Msg, State) ->
{noreply, State}.
handle_info(flush, State) ->
flush_stats(),
flush_to_mnesia(),
erlang:send_after(?FLUSH_INTERVAL, self(), flush),
{noreply, State};
handle_info({mnesia_table_event, {write, Record, _ActivityId}}, State) ->
Table = element(1, Record),
process_write(Table, Record, []),
handle_info(cleanup_daily, State) ->
cleanup_old_daily(?DAILY_RETENTION_DAYS),
erlang:send_after(?CLEANUP_INTERVAL, self(), cleanup_daily),
{noreply, State};
handle_info({write, Table, New, Old, _ActivityId}, State) ->
process_write(Table, New, Old),
handle_info({mnesia_table_event, Event}, State) ->
handle_mnesia_event(Event),
{noreply, State};
handle_info(_Info, State) -> {noreply, State}.
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, _State) -> ok.
code_change(_OldVsn, State, _Extra) -> {ok, State}.
terminate(_Reason, _State) ->
flush_to_mnesia(),
ok.
ensure_counter(Key) ->
case ets:member(stats_ets, Key) of
false -> ets:insert(stats_ets, {Key, 0});
true -> ok
code_change(_OldVsn, State, _Extra) ->
{ok, State}.
%%%===================================================================
%%% Bootstrap
%%%===================================================================
do_subscribe(#state{subscribed = true} = State) ->
State;
do_subscribe(State) ->
load_from_mnesia(),
case ets:info(?COUNTER_ETS, size) =:= 0 andalso ets:info(?DAILY_ETS, size) =:= 0 of
true -> do_rebuild();
false -> stats_tops:rebuild()
end,
lists:foreach(fun(Table) ->
mnesia:subscribe({table, Table, detailed})
end, ?SUB_TABLES),
State#state{ready = true, subscribed = true}.
do_rebuild() ->
ets:delete_all_objects(?COUNTER_ETS),
ets:delete_all_objects(?DAILY_ETS),
backfill_all(),
stats_tops:rebuild(),
flush_to_mnesia(),
ok.
load_from_mnesia() ->
case lists:member(stats_counter, mnesia:system_info(tables)) of
true ->
lists:foreach(fun(#stats_counter{key = K, value = V}) ->
ets:insert(?COUNTER_ETS, {K, V})
end, mnesia:dirty_match_object(#stats_counter{_ = '_'}));
false -> ok
end,
case lists:member(stats_daily, mnesia:system_info(tables)) of
true ->
lists:foreach(fun(#stats_daily{key = K, value = V}) ->
ets:insert(?DAILY_ETS, {K, V})
end, mnesia:dirty_match_object(#stats_daily{_ = '_'}));
false -> ok
end,
ok.
%%%===================================================================
%%% Mnesia events
%%%===================================================================
handle_mnesia_event({write, Table, New, Olds, _ActivityId}) ->
Old = case Olds of
[O | _] -> O;
[] -> undefined
end,
process_write(Table, New, Old);
handle_mnesia_event({delete, _Table, _Key, Olds, _ActivityId}) ->
lists:foreach(fun(Old) -> process_delete(element(1, Old), Old) end, Olds);
handle_mnesia_event({delete_object, Table, Old, _Olds, _ActivityId}) ->
process_delete(Table, Old);
handle_mnesia_event(_) ->
ok.
%%%===================================================================
%%% Write / delete handlers
%%%===================================================================
process_write(user, New, undefined) ->
#user{status = St, role = Role, created_at = Created} = New,
add_dim(user, status, St, 1),
add_dim(user, role, Role, 1),
add_daily(users_created, date_of(Created), 1),
case St of
pending -> add_daily(pending_users_created, date_of(Created), 1);
_ -> ok
end;
process_write(user, New, Old) ->
#user{status = StNew, role = RoleNew, created_at = Created} = New,
#user{status = StOld, role = RoleOld} = Old,
move_dim(user, status, StOld, StNew),
move_dim(user, role, RoleOld, RoleNew),
case {StOld, StNew} of
{pending, pending} -> ok;
{pending, _} -> add_daily(pending_users_created, date_of(Created), -1);
{_, pending} -> add_daily(pending_users_created, date_of(Created), 1);
_ -> ok
end;
process_write(event, New, undefined) ->
#event{status = St, event_type = Type, created_at = Created} = New,
add_dim(event, status, St, 1),
add_dim(event, type, Type, 1),
add_daily(events_created, date_of(Created), 1),
stats_tops:on_event(New, undefined);
process_write(event, New, Old) ->
#event{status = StNew, event_type = TypeNew} = New,
#event{status = StOld, event_type = TypeOld} = Old,
move_dim(event, status, StOld, StNew),
move_dim(event, type, TypeOld, TypeNew),
stats_tops:on_event(New, Old);
process_write(calendar, New, undefined) ->
#calendar{status = St, type = Type, created_at = Created} = New,
add_dim(calendar, status, St, 1),
add_dim(calendar, type, Type, 1),
add_daily(calendars_created, date_of(Created), 1),
stats_tops:on_calendar(New, undefined);
process_write(calendar, New, Old) ->
#calendar{status = StNew, type = TypeNew} = New,
#calendar{status = StOld, type = TypeOld} = Old,
move_dim(calendar, status, StOld, StNew),
move_dim(calendar, type, TypeOld, TypeNew),
stats_tops:on_calendar(New, Old);
process_write(booking, New, undefined) ->
#booking{status = St, created_at = Created} = New,
add_dim(booking, status, St, 1),
add_daily(bookings_created, date_of(Created), 1);
process_write(booking, New, Old) ->
move_dim(booking, status, Old#booking.status, New#booking.status);
process_write(review, New, undefined) ->
#review{status = St, target_type = TT, created_at = Created} = New,
add_dim(review, status, St, 1),
add_dim(review, target_type, TT, 1),
add_daily(reviews_created, date_of(Created), 1),
stats_tops:on_review(New, undefined);
process_write(review, New, Old) ->
#review{status = StNew, target_type = TTNew} = New,
#review{status = StOld, target_type = TTOld} = Old,
move_dim(review, status, StOld, StNew),
move_dim(review, target_type, TTOld, TTNew),
stats_tops:on_review(New, Old);
process_write(report, New, undefined) ->
#report{status = St, target_type = TT, created_at = Created} = New,
add_dim(report, status, St, 1),
add_dim(report, target_type, TT, 1),
add_daily(reports_created, date_of(Created), 1),
apply_report_resolution(report_resolution(New)),
stats_tops:on_report(New, undefined);
process_write(report, New, Old) ->
#report{status = StNew, target_type = TTNew} = New,
#report{status = StOld, target_type = TTOld} = Old,
move_dim(report, status, StOld, StNew),
move_dim(report, target_type, TTOld, TTNew),
move_resolution(report_resolution(Old), report_resolution(New)),
stats_tops:on_report(New, Old);
process_write(ticket, New, undefined) ->
#ticket{status = St, count = Count, first_seen = First} = New,
add_dim(ticket, status, St, 1),
add_counter(ticket_errors_sum, Count),
add_daily(tickets_created, date_of(First), 1),
apply_ticket_resolution(ticket_resolution(New));
process_write(ticket, New, Old) ->
#ticket{status = StNew, count = CountNew} = New,
#ticket{status = StOld, count = CountOld} = Old,
move_dim(ticket, status, StOld, StNew),
add_counter(ticket_errors_sum, CountNew - CountOld),
move_resolution(ticket_resolution(Old), ticket_resolution(New));
process_write(subscription, New, undefined) ->
#subscription{status = St, plan = Plan, trial_used = Trial,
created_at = Created} = New,
add_dim(subscription, status, St, 1),
add_dim(subscription, plan, Plan, 1),
case Trial of
false -> add_counter(trial_subscriptions, 1);
_ -> ok
end,
add_daily(subscriptions_created, date_of(Created), 1);
process_write(subscription, New, Old) ->
#subscription{status = StNew, plan = PlanNew, trial_used = TrialNew} = New,
#subscription{status = StOld, plan = PlanOld, trial_used = TrialOld} = Old,
move_dim(subscription, status, StOld, StNew),
move_dim(subscription, plan, PlanOld, PlanNew),
case {TrialOld, TrialNew} of
{false, false} -> ok;
{false, _} -> add_counter(trial_subscriptions, -1);
{_, false} -> add_counter(trial_subscriptions, 1);
_ -> ok
end;
process_write(_, _, _) ->
ok.
process_delete(user, #user{status = St, role = Role, created_at = Created}) ->
add_dim(user, status, St, -1),
add_dim(user, role, Role, -1),
add_daily(users_created, date_of(Created), -1),
case St of
pending -> add_daily(pending_users_created, date_of(Created), -1);
_ -> ok
end;
process_delete(event, #event{status = St, event_type = Type, created_at = Created} = E) ->
add_dim(event, status, St, -1),
add_dim(event, type, Type, -1),
add_daily(events_created, date_of(Created), -1),
stats_tops:remove_event(E);
process_delete(calendar, #calendar{status = St, type = Type, created_at = Created} = C) ->
add_dim(calendar, status, St, -1),
add_dim(calendar, type, Type, -1),
add_daily(calendars_created, date_of(Created), -1),
stats_tops:remove_calendar(C);
process_delete(booking, #booking{status = St, created_at = Created}) ->
add_dim(booking, status, St, -1),
add_daily(bookings_created, date_of(Created), -1);
process_delete(review, #review{status = St, target_type = TT, created_at = Created} = R) ->
add_dim(review, status, St, -1),
add_dim(review, target_type, TT, -1),
add_daily(reviews_created, date_of(Created), -1),
stats_tops:remove_review(R);
process_delete(report, #report{status = St, target_type = TT, created_at = Created} = R) ->
add_dim(report, status, St, -1),
add_dim(report, target_type, TT, -1),
add_daily(reports_created, date_of(Created), -1),
remove_report_resolution(report_resolution(R)),
stats_tops:remove_report(R);
process_delete(ticket, #ticket{status = St, count = Count, first_seen = First} = T) ->
add_dim(ticket, status, St, -1),
add_counter(ticket_errors_sum, -Count),
add_daily(tickets_created, date_of(First), -1),
remove_ticket_resolution(ticket_resolution(T));
process_delete(subscription, #subscription{status = St, plan = Plan,
trial_used = Trial, created_at = Created}) ->
add_dim(subscription, status, St, -1),
add_dim(subscription, plan, Plan, -1),
case Trial of
false -> add_counter(trial_subscriptions, -1);
_ -> ok
end,
add_daily(subscriptions_created, date_of(Created), -1);
process_delete(_, _) ->
ok.
%%%===================================================================
%%% Backfill
%%%===================================================================
backfill_all() ->
lists:foreach(fun(U) -> process_write(user, U, undefined) end,
safe_match(user, #user{_ = '_'})),
lists:foreach(fun(E) -> process_write(event, E, undefined) end,
safe_match(event, #event{_ = '_'})),
lists:foreach(fun(C) -> process_write(calendar, C, undefined) end,
safe_match(calendar, #calendar{_ = '_'})),
lists:foreach(fun(B) -> process_write(booking, B, undefined) end,
safe_match(booking, #booking{_ = '_'})),
lists:foreach(fun(R) -> process_write(review, R, undefined) end,
safe_match(review, #review{_ = '_'})),
lists:foreach(fun(R) -> process_write(report, R, undefined) end,
safe_match(report, #report{_ = '_'})),
lists:foreach(fun(T) -> process_write(ticket, T, undefined) end,
safe_match(ticket, #ticket{_ = '_'})),
lists:foreach(fun(S) -> process_write(subscription, S, undefined) end,
safe_match(subscription, #subscription{_ = '_'})),
ok.
safe_match(Table, Pattern) ->
case lists:member(Table, mnesia:system_info(tables)) of
true -> mnesia:dirty_match_object(Pattern);
false -> []
end.
process_write(event, New, Old) ->
CId = element(#event.calendar_id, New),
StNew = element(#event.status, New),
StOld = if Old =:= [] -> []; true -> element(#event.status, Old) end,
if StNew =:= active, StOld =/= active ->
ensure_counter({events_created, CId}),
ets:update_counter(stats_ets, {events_created, CId}, {2, 1});
StNew =:= cancelled, StOld =/= cancelled ->
ensure_counter({events_cancelled, CId}),
ets:update_counter(stats_ets, {events_cancelled, CId}, {2, 1});
true -> ok
end;
process_write(booking, New, Old) ->
EvId = element(#booking.event_id, New),
StNew = element(#booking.status, New),
StOld = if Old =:= [] -> []; true -> element(#booking.status, Old) end,
CId = case mnesia:dirty_read({event, EvId}) of
[#event{calendar_id = C}] -> C;
_ -> undefined
end,
if CId =/= undefined ->
if StNew =:= confirmed, StOld =/= confirmed ->
ensure_counter({bookings_confirmed, CId}),
ets:update_counter(stats_ets, {bookings_confirmed, CId}, {2, 1});
StNew =:= cancelled, StOld =/= cancelled ->
ensure_counter({bookings_cancelled, CId}),
ets:update_counter(stats_ets, {bookings_cancelled, CId}, {2, 1});
true -> ok
end,
if Old =:= [] ->
ensure_counter({bookings_created, CId}),
ets:update_counter(stats_ets, {bookings_created, CId}, {2, 1});
true -> ok
end;
true -> ok
end;
process_write(review, New, Old) ->
CId = case element(#review.target_type, New) of
event ->
EvId = element(#review.target_id, New),
case mnesia:dirty_read({event, EvId}) of
[#event{calendar_id = C}] -> C;
_ -> undefined
end;
calendar ->
element(#review.target_id, New)
end,
Rating = element(#review.rating, New),
if CId =/= undefined ->
if Old =:= [] ->
ensure_counter({reviews_created, CId}),
ets:update_counter(stats_ets, {reviews_created, CId}, {2, 1}),
ensure_counter({reviews_sum, CId}),
ets:update_counter(stats_ets, {reviews_sum, CId}, {2, Rating}),
ensure_counter({reviews_count, CId}),
ets:update_counter(stats_ets, {reviews_count, CId}, {2, 1});
true ->
OldRating = element(#review.rating, Old),
if Rating =/= OldRating ->
ensure_counter({reviews_sum, CId}),
ets:update_counter(stats_ets, {reviews_sum, CId}, {2, Rating - OldRating});
true -> ok
end
end;
true -> ok
end;
process_write(_, _, _) -> ok.
%%%===================================================================
%%% ETS helpers
%%%===================================================================
flush_stats() ->
All = ets:tab2list(stats_ets),
lists:foreach(fun({{Metric, EntityId}, Value}) ->
mnesia:dirty_write(#stats{
id = list_to_binary(io_lib:format("~p_~p_~p", [Metric, EntityId, os:system_time(millisecond)])),
metric = Metric,
entity_id = EntityId,
value = Value,
timestamp = calendar:local_time()
})
end, All),
ets:match_delete(stats_ets, '_').
dim_key(Entity, Field, Value) ->
{list_to_atom(atom_to_list(Entity) ++ "_by_" ++ atom_to_list(Field)), Value}.
add_dim(Entity, Field, Value, Delta) when is_atom(Value) ->
add_counter(dim_key(Entity, Field, Value), Delta);
add_dim(_, _, _, _) ->
ok.
move_dim(_E, _F, Same, Same) ->
ok;
move_dim(Entity, Field, OldVal, NewVal) ->
add_dim(Entity, Field, OldVal, -1),
add_dim(Entity, Field, NewVal, 1).
add_counter(_Key, 0) ->
ok;
add_counter(Key, Delta) ->
ensure_key(?COUNTER_ETS, Key),
ets:update_counter(?COUNTER_ETS, Key, {2, Delta}),
%% clamp at 0 for safety under races
case ets:lookup(?COUNTER_ETS, Key) of
[{Key, V}] when V < 0 -> ets:insert(?COUNTER_ETS, {Key, 0});
_ -> ok
end.
add_daily(_Metric, undefined, _) ->
ok;
add_daily(_Metric, _Date, 0) ->
ok;
add_daily(Metric, Date, Delta) ->
Key = {Metric, Date},
ensure_key(?DAILY_ETS, Key),
ets:update_counter(?DAILY_ETS, Key, {2, Delta}),
case ets:lookup(?DAILY_ETS, Key) of
[{Key, V}] when V < 0 -> ets:insert(?DAILY_ETS, {Key, 0});
_ -> ok
end.
ensure_key(Tab, Key) ->
case ets:member(Tab, Key) of
true -> ok;
false -> ets:insert(Tab, {Key, 0})
end.
ets_get(Tab, Key) ->
case ets:lookup(Tab, Key) of
[{Key, V}] -> V;
[] -> 0
end.
read_dim(Entity, Field) ->
Prefix = list_to_atom(atom_to_list(Entity) ++ "_by_" ++ atom_to_list(Field)),
lists:filtermap(fun
({{P, Val}, V}) when P =:= Prefix, V > 0, is_atom(Val) -> {true, {Val, V}};
(_) -> false
end, ets:tab2list(?COUNTER_ETS)).
read_daily(Metric, From, To) ->
FromDate = date_of(From),
ToDate = date_of(To),
Matches = lists:filtermap(fun
({{M, Day}, V}) when M =:= Metric, V > 0,
Day >= FromDate, Day =< ToDate ->
{true, {Day, V}};
(_) -> false
end, ets:tab2list(?DAILY_ETS)),
lists:sort(Matches).
date_of({{Y, M, D}, _}) -> {Y, M, D};
date_of({Y, M, D} = Date) when is_integer(Y), is_integer(M), is_integer(D) -> Date;
date_of(_) -> undefined.
%%%===================================================================
%%% Resolution time (avg)
%%%===================================================================
-define(TICKET_EPOCH, {{1970, 1, 1}, {0, 0, 0}}).
%% Ticket: same semantics as core_ticket:avg_resolution_time — closed_at set.
ticket_resolution(#ticket{first_seen = First, closed_at = Closed}) ->
case is_set_datetime(Closed) of
true ->
case resolution_seconds(First, Closed) of
undefined -> undefined;
Sec -> {ticket, Sec}
end;
false -> undefined
end.
%% Report: only reviewed|dismissed with resolved_at.
report_resolution(#report{status = St, created_at = Created, resolved_at = Resolved})
when St =:= reviewed; St =:= dismissed ->
case is_set_datetime(Resolved) of
true ->
case resolution_seconds(Created, Resolved) of
undefined -> undefined;
Sec -> {{report, St}, Sec}
end;
false -> undefined
end;
report_resolution(_) ->
undefined.
move_resolution(Same, Same) ->
ok;
move_resolution(Old, New) ->
remove_resolution(Old),
apply_resolution(New).
apply_resolution(undefined) -> ok;
apply_resolution({ticket, Sec}) -> apply_ticket_resolution({ticket, Sec});
apply_resolution({{report, St}, Sec}) -> apply_report_resolution({{report, St}, Sec}).
remove_resolution(undefined) -> ok;
remove_resolution({ticket, Sec}) -> remove_ticket_resolution({ticket, Sec});
remove_resolution({{report, St}, Sec}) -> remove_report_resolution({{report, St}, Sec}).
apply_ticket_resolution(undefined) -> ok;
apply_ticket_resolution({ticket, Sec}) ->
add_counter(ticket_resolved_count, 1),
add_counter(ticket_resolved_seconds_sum, Sec).
remove_ticket_resolution(undefined) -> ok;
remove_ticket_resolution({ticket, Sec}) ->
add_counter(ticket_resolved_count, -1),
add_counter(ticket_resolved_seconds_sum, -Sec).
apply_report_resolution(undefined) -> ok;
apply_report_resolution({{report, St}, Sec}) ->
add_counter({report_resolved_count, St}, 1),
add_counter({report_resolved_seconds_sum, St}, Sec).
remove_report_resolution(undefined) -> ok;
remove_report_resolution({{report, St}, Sec}) ->
add_counter({report_resolved_count, St}, -1),
add_counter({report_resolved_seconds_sum, St}, -Sec).
is_set_datetime(undefined) -> false;
is_set_datetime(?TICKET_EPOCH) -> false;
is_set_datetime({{_, _, _}, {_, _, _}}) -> true;
is_set_datetime(_) -> false.
resolution_seconds(Start, End) ->
case {is_set_datetime(Start), is_set_datetime(End)} of
{true, true} ->
max(0, calendar:datetime_to_gregorian_seconds(End) -
calendar:datetime_to_gregorian_seconds(Start));
_ -> undefined
end.
%%%===================================================================
%%% Persist / cleanup
%%%===================================================================
flush_to_mnesia() ->
Now = calendar:universal_time(),
case lists:member(stats_counter, mnesia:system_info(tables)) of
true ->
lists:foreach(fun({Key, Value}) ->
mnesia:dirty_write(#stats_counter{key = Key, value = Value, updated_at = Now})
end, ets:tab2list(?COUNTER_ETS));
false -> ok
end,
case lists:member(stats_daily, mnesia:system_info(tables)) of
true ->
lists:foreach(fun({Key, Value}) ->
mnesia:dirty_write(#stats_daily{key = Key, value = Value, updated_at = Now})
end, ets:tab2list(?DAILY_ETS));
false -> ok
end,
ok.
cleanup_old_daily(RetentionDays) ->
case lists:member(stats_daily, mnesia:system_info(tables)) of
false -> ok;
true ->
{Date, _} = calendar:universal_time(),
CutoffSec = calendar:datetime_to_gregorian_seconds({Date, {0, 0, 0}})
- RetentionDays * 86400,
Cutoff = element(1, calendar:gregorian_seconds_to_datetime(CutoffSec)),
lists:foreach(fun
(#stats_daily{key = {_, Day}} = Row) when Day < Cutoff ->
mnesia:dirty_delete_object(Row),
ets:delete(?DAILY_ETS, Row#stats_daily.key);
(_) -> ok
end, mnesia:dirty_match_object(#stats_daily{_ = '_'})),
ok
end.
+329
View File
@@ -0,0 +1,329 @@
%%%-------------------------------------------------------------------
%%% @doc ETS ordered_set tops для admin stats.
%%% Класс A: denormalized rating на event/calendar.
%%% Класс B: агрегаты целей review/report (+ calendar pos/neg).
%%% @end
%%%-------------------------------------------------------------------
-module(stats_tops).
-include("records.hrl").
-export([
init_tables/0,
clear/0,
rebuild/0,
on_event/2,
on_calendar/2,
on_review/2,
on_report/2,
remove_event/1,
remove_calendar/1,
remove_review/1,
remove_report/1,
get_top_events_by_rating/1,
get_top_calendars_by_rating/1,
get_top_calendars_by_reviews/1,
get_top_calendars_by_positive_reviews/1,
get_top_calendars_by_negative_reviews/1,
get_top_targets_by_reviews/1,
get_top_targets_by_positive_reviews/1,
get_top_targets_by_negative_reviews/1,
get_top_targets_by_reports/1
]).
-define(IDX, stats_tops_idx).
-define(ORD_EVENT_RATING, stats_top_event_rating).
-define(ORD_CAL_RATING, stats_top_calendar_rating).
-define(ORD_CAL_REVIEWS, stats_top_calendar_reviews).
-define(ORD_CAL_POS, stats_top_calendar_pos).
-define(ORD_CAL_NEG, stats_top_calendar_neg).
-define(ORD_REVIEWS_ALL, stats_top_reviews_all).
-define(ORD_REVIEWS_POS, stats_top_reviews_pos).
-define(ORD_REVIEWS_NEG, stats_top_reviews_neg).
-define(ORD_REPORTS_ALL, stats_top_reports_all).
-define(ALL_ORD, [
?ORD_EVENT_RATING, ?ORD_CAL_RATING, ?ORD_CAL_REVIEWS,
?ORD_CAL_POS, ?ORD_CAL_NEG,
?ORD_REVIEWS_ALL, ?ORD_REVIEWS_POS, ?ORD_REVIEWS_NEG, ?ORD_REPORTS_ALL
]).
%%%===================================================================
%%% Lifecycle
%%%===================================================================
init_tables() ->
ensure_ets(?IDX, [named_table, public, set, {read_concurrency, true}]),
lists:foreach(fun(T) ->
ensure_ets(T, [named_table, public, ordered_set, {read_concurrency, true}])
end, ?ALL_ORD),
ok.
clear() ->
case ready() of
true ->
ets:delete_all_objects(?IDX),
lists:foreach(fun(T) -> ets:delete_all_objects(T) end, ?ALL_ORD);
false -> ok
end,
ok.
rebuild() ->
clear(),
foreach_table(event, #event{_ = '_'}, fun(E) -> on_event(E, undefined) end),
foreach_table(calendar, #calendar{_ = '_'}, fun(C) -> on_calendar(C, undefined) end),
foreach_table(review, #review{_ = '_'}, fun(R) -> on_review(R, undefined) end),
foreach_table(report, #report{_ = '_'}, fun(R) -> on_report(R, undefined) end),
ok.
ready() ->
ets:info(?IDX) =/= undefined.
foreach_table(Table, Pattern, Fun) ->
case lists:member(Table, mnesia:system_info(tables)) of
true -> lists:foreach(Fun, mnesia:dirty_match_object(Pattern));
false -> ok
end.
%%%===================================================================
%%% Entity rating updates (class A)
%%%===================================================================
on_event(undefined, _) -> ok;
on_event(#event{id = Id, rating_avg = Avg}, Old) ->
when_ready(fun() ->
OldAvg = case Old of #event{rating_avg = A} -> A; _ -> undefined end,
set_score(?ORD_EVENT_RATING, {event_rating, Id}, Id, OldAvg, Avg)
end);
on_event(_, _) -> ok.
on_calendar(undefined, _) -> ok;
on_calendar(#calendar{id = Id, rating_avg = Avg, rating_count = Count}, Old) ->
when_ready(fun() ->
{OldAvg, OldCount} = case Old of
#calendar{rating_avg = A, rating_count = C} -> {A, C};
_ -> {undefined, undefined}
end,
set_score(?ORD_CAL_RATING, {calendar_rating, Id}, Id, OldAvg, Avg),
set_score(?ORD_CAL_REVIEWS, {calendar_reviews, Id}, Id, OldCount, Count)
end);
on_calendar(_, _) -> ok.
remove_event(#event{id = Id, rating_avg = Avg}) ->
when_ready(fun() -> clear_score(?ORD_EVENT_RATING, {event_rating, Id}, Id, Avg) end);
remove_event(_) -> ok.
remove_calendar(#calendar{id = Id, rating_avg = Avg, rating_count = Count}) ->
when_ready(fun() ->
clear_score(?ORD_CAL_RATING, {calendar_rating, Id}, Id, Avg),
clear_score(?ORD_CAL_REVIEWS, {calendar_reviews, Id}, Id, Count)
end);
remove_calendar(_) -> ok.
%%%===================================================================
%%% Review / report target updates (class B)
%%%===================================================================
on_review(undefined, _) -> ok;
on_review(New, undefined) ->
when_ready(fun() -> apply_review_delta(New, 1) end);
on_review(New, Old) when is_tuple(Old) ->
when_ready(fun() ->
apply_review_delta(Old, -1),
apply_review_delta(New, 1)
end);
on_review(_, _) -> ok.
remove_review(Review) ->
when_ready(fun() -> apply_review_delta(Review, -1) end).
on_report(undefined, _) -> ok;
on_report(New, undefined) ->
when_ready(fun() -> apply_report_delta(New, 1) end);
on_report(New, Old) when is_tuple(Old) ->
when_ready(fun() ->
apply_report_delta(Old, -1),
apply_report_delta(New, 1)
end);
on_report(_, _) -> ok.
remove_report(Report) ->
when_ready(fun() -> apply_report_delta(Report, -1) end).
apply_review_delta(#review{target_type = Type, target_id = Id, rating = Rating}, Delta)
when is_atom(Type), is_binary(Id) ->
Target = {Type, Id},
bump(?ORD_REVIEWS_ALL, {reviews_all, Type, Id}, Target, Delta),
case Rating >= 4 of
true ->
bump(?ORD_REVIEWS_POS, {reviews_pos, Type, Id}, Target, Delta),
case Type of
calendar -> bump(?ORD_CAL_POS, {calendar_pos, Id}, Id, Delta);
_ -> ok
end;
false -> ok
end,
case Rating =< 2 of
true ->
bump(?ORD_REVIEWS_NEG, {reviews_neg, Type, Id}, Target, Delta),
case Type of
calendar -> bump(?ORD_CAL_NEG, {calendar_neg, Id}, Id, Delta);
_ -> ok
end;
false -> ok
end;
apply_review_delta(_, _) -> ok.
apply_report_delta(#report{target_type = Type, target_id = Id}, Delta)
when is_atom(Type), is_binary(Id) ->
bump(?ORD_REPORTS_ALL, {reports_all, Type, Id}, {Type, Id}, Delta);
apply_report_delta(_, _) -> ok.
%%%===================================================================
%%% Reads
%%%===================================================================
get_top_events_by_rating(N) ->
get_entities(?ORD_EVENT_RATING, N, fun load_event/1).
get_top_calendars_by_rating(N) ->
get_entities(?ORD_CAL_RATING, N, fun load_calendar/1).
get_top_calendars_by_reviews(N) ->
get_entities(?ORD_CAL_REVIEWS, N, fun load_calendar/1).
get_top_calendars_by_positive_reviews(N) ->
get_entities(?ORD_CAL_POS, N, fun load_calendar/1).
get_top_calendars_by_negative_reviews(N) ->
get_entities(?ORD_CAL_NEG, N, fun load_calendar/1).
get_top_targets_by_reviews(N) ->
get_target_triples(?ORD_REVIEWS_ALL, N).
get_top_targets_by_positive_reviews(N) ->
get_target_triples(?ORD_REVIEWS_POS, N).
get_top_targets_by_negative_reviews(N) ->
get_target_triples(?ORD_REVIEWS_NEG, N).
get_top_targets_by_reports(N) ->
get_target_triples(?ORD_REPORTS_ALL, N).
%%%===================================================================
%%% Internal
%%%===================================================================
when_ready(Fun) ->
case ready() of
true -> Fun();
false -> ok
end.
ensure_ets(Name, Opts) ->
case ets:info(Name) of
undefined -> ets:new(Name, Opts);
_ -> Name
end.
set_score(_Ord, _IdxKey, _Entity, Same, Same) ->
ok;
set_score(Ord, IdxKey, Entity, _OldScore, NewScore) ->
case ets:lookup(?IDX, IdxKey) of
[{IdxKey, Prev}] -> ets:delete(Ord, ord_key(Prev, Entity));
[] -> ok
end,
ets:insert(Ord, {ord_key(NewScore, Entity), Entity}),
ets:insert(?IDX, {IdxKey, NewScore}),
ok.
clear_score(Ord, IdxKey, Entity, Score) ->
case ets:lookup(?IDX, IdxKey) of
[{IdxKey, Prev}] ->
ets:delete(Ord, ord_key(Prev, Entity)),
ets:delete(?IDX, IdxKey);
[] ->
ets:delete(Ord, ord_key(Score, Entity))
end,
ok.
bump(_Ord, _IdxKey, _Entity, 0) ->
ok;
bump(Ord, IdxKey, Entity, Delta) ->
Old = case ets:lookup(?IDX, IdxKey) of
[{IdxKey, C}] -> C;
[] -> 0
end,
case Old > 0 of
true -> ets:delete(Ord, ord_key(Old, Entity));
false -> ok
end,
New = max(0, Old + Delta),
case New > 0 of
true ->
ets:insert(Ord, {ord_key(New, Entity), Entity}),
ets:insert(?IDX, {IdxKey, New});
false ->
ets:delete(?IDX, IdxKey)
end,
ok.
ord_key(Score, Entity) when is_number(Score) ->
{-Score, Entity};
ord_key(_, Entity) ->
{0, Entity}.
get_entities(Ord, N, LoadFun) ->
case ready() of
false -> {error, not_ready};
true ->
Keys = take_entities(Ord, ets:first(Ord), N, []),
{ok, lists:filtermap(LoadFun, Keys)}
end.
get_target_triples(Ord, N) ->
case ready() of
false -> {error, not_ready};
true ->
Items = take_entities(Ord, ets:first(Ord), N, []),
Triples = lists:filtermap(fun
({Type, Id}) ->
case ets:lookup(?IDX, idx_for_ord(Ord, Type, Id)) of
[{_, Count}] when Count > 0 -> {true, {Type, Id, Count}};
_ -> false
end;
(_) -> false
end, Items),
{ok, Triples}
end.
idx_for_ord(?ORD_REVIEWS_ALL, Type, Id) -> {reviews_all, Type, Id};
idx_for_ord(?ORD_REVIEWS_POS, Type, Id) -> {reviews_pos, Type, Id};
idx_for_ord(?ORD_REVIEWS_NEG, Type, Id) -> {reviews_neg, Type, Id};
idx_for_ord(?ORD_REPORTS_ALL, Type, Id) -> {reports_all, Type, Id}.
take_entities(_Ord, '$end_of_table', _N, Acc) ->
lists:reverse(Acc);
take_entities(_Ord, _Key, 0, Acc) ->
lists:reverse(Acc);
take_entities(Ord, Key, N, Acc) ->
case ets:lookup(Ord, Key) of
[{Key, Entity}] ->
take_entities(Ord, ets:next(Ord, Key), N - 1, [Entity | Acc]);
[] ->
take_entities(Ord, ets:next(Ord, Key), N, Acc)
end.
load_event(Id) when is_binary(Id) ->
case mnesia:dirty_read({event, Id}) of
[E] -> {true, E};
_ -> false
end;
load_event(_) -> false.
load_calendar(Id) when is_binary(Id) ->
case mnesia:dirty_read({calendar, Id}) of
[C] -> {true, C};
_ -> false
end;
load_calendar(_) -> false.