Files
EventHubBack/src/infra/stats_tops.erl
T
aleksey 31c2a1b7d7
CI / test (push) Failing after 13m10s
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
fix: subscribe stats_collector after Mnesia tables are ready
Backfill ran during init_tables before wait_for_tables, so dirty reads hit {no_exists,user} and crash-looped eventhub on IFT.
2026-07-18 23:33:56 +03:00

335 lines
9.9 KiB
Erlang

%%%-------------------------------------------------------------------
%%% @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 ->
try lists:foreach(Fun, mnesia:dirty_match_object(Pattern))
catch
exit:{aborted, {no_exists, _}} -> ok;
error:{aborted, {no_exists, _}} -> ok
end;
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.