%% =================================================================== %% EventHub – infra_mnesia (финальная версия с автоочисткой кластера) %% =================================================================== -module(infra_mnesia). -behaviour(gen_server). -include("records.hrl"). -export([start_link/0, init_tables/0, wait_for_tables/0, wait_for_table/1]). -export([add_cluster_nodes/1]). -export([configure_dump_log/0, verify_dump_log/0]). -export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]). -define(TABLES, [ user, session, verification, admin, admin_session, auth_session, calendar, calendar_share, calendar_specialist, event, recurrence_exception, booking, review, report, banned_word, automod_settings, automod_hit, ticket, subscription, admin_audit, notification, stats_counter, stats_daily, node_metric, schema_migration ]). -define(DISC_TABLES, ?TABLES -- [session, verification, admin_session, node_metric]). %% ram_copies: joining nodes must add_table_copy — create_table already_exists skips it. -define(RAM_TABLES, [session, verification, admin_session]). %% Disc load on IFT after crash-loop can exceed default gen_server:call 5s. -define(TABLE_WAIT_TIMEOUT, 120000). -define(CLEANUP_INTERVAL, 30000). % 30 секунд %% =================================================================== %% API %% =================================================================== start_link() -> % Счётчики для метрик (сессии, ws-соединения) ets:new(eventhub_counters, [named_table, public, set, {write_concurrency, true}]), gen_server:start_link({local, ?MODULE}, ?MODULE, [], []). init_tables() -> gen_server:call(?MODULE, init_tables, infinity). wait_for_tables() -> gen_server:call(?MODULE, wait_for_tables, infinity). add_cluster_nodes(Nodes) -> gen_server:call(?MODULE, {add_nodes, Nodes}, infinity). %% @doc Soften dump_log overload under write-heavy load (IFT stress). %% OTP warns {dump_log, write_threshold|time_threshold} when a new dump %% is requested while the previous dump still runs. %% Must be set BEFORE mnesia starts — change_config/2 does not accept %% dump_log_write_threshold (returns {error, dump_log_write_threshold}). %% Override: MNESIA_DUMP_LOG_WRITE_THRESHOLD (default 50000; OTP default ~1000). -spec configure_dump_log() -> ok. configure_dump_log() -> Threshold = env_int("MNESIA_DUMP_LOG_WRITE_THRESHOLD", 50000), case application:load(mnesia) of ok -> ok; {error, {already_loaded, mnesia}} -> ok; {error, LoadErr} -> io:format("Mnesia load before dump_log config failed: ~p~n", [LoadErr]) end, ok = application:set_env(mnesia, dump_log_write_threshold, Threshold), io:format("Mnesia dump_log_write_threshold set to ~p (before start)~n", [Threshold]), ok. %% @doc Log active threshold after mnesia:start (sanity check for IFT logs). -spec verify_dump_log() -> ok. verify_dump_log() -> case catch mnesia:system_info(dump_log_write_threshold) of N when is_integer(N) -> io:format("Mnesia dump_log_write_threshold active: ~p~n", [N]), ok; Other -> io:format("Mnesia dump_log_write_threshold verify: ~p~n", [Other]), ok end. env_int(Name, Default) -> case os:getenv(Name) of false -> Default; "" -> Default; Val -> try list_to_integer(Val) of N when N > 0 -> N; _ -> Default catch error:badarg -> Default end end. %% =================================================================== %% gen_server callbacks %% =================================================================== init([]) -> {ok, #{}}. handle_call(init_tables, _From, State) -> ExtraNodes = application:get_env(eventhub, extra_db_nodes, []), case ExtraNodes of [] -> ok = maybe_recreate_schema(); _ -> ok = join_cluster(ExtraNodes) end, lists:foreach(fun create_table/1, ?TABLES), lists:foreach(fun add_local_ram_copy/1, ?RAM_TABLES), % Принудительное создание node_metric на каждом узле case mnesia:create_table(node_metric, [ {disc_copies, [node()]}, % хранить на диске {local_content, true}, % не реплицировать {attributes, record_info(fields, node_metric)} ]) of {atomic, ok} -> ok; {aborted, {already_exists, node_metric}} -> ok; _ -> ok end, % ГАРАНТИРУЕМ, что узел имеет локальную копию node_metric (критично для присоединяющихся узлов) case lists:member(node(), mnesia:table_info(node_metric, disc_copies) ++ mnesia:table_info(node_metric, ram_copies)) of false -> mnesia:add_table_copy(node_metric, node(), disc_copies); true -> ok end, ok = create_indices(), %% stats_collector:subscribe — после wait_for_tables + migrations (eventhub_app) ok = start_cleanup_timer(), {reply, ok, State}; handle_call({add_nodes, Nodes}, _From, State) -> ok = do_add_nodes(Nodes), {reply, ok, State}; handle_call(wait_for_tables, _From, State) -> case do_wait_for_tables() of ok -> {reply, ok, State}; {error, Reason} -> {reply, {error, Reason}, State} end; handle_cast(_Msg, State) -> {noreply, State}. handle_info({cleanup}, State) -> prune_dead_nodes(), erlang:send_after(?CLEANUP_INTERVAL, self(), {cleanup}), {noreply, State}; handle_info(_Info, State) -> {noreply, State}. terminate(_Reason, _State) -> ok. code_change(_OldVsn, State, _Extra) -> {ok, State}. %% =================================================================== %% Управление схемой %% =================================================================== maybe_recreate_schema() -> MnesiaDir = mnesia:system_info(directory), case filelib:is_dir(MnesiaDir) of false -> io:format("Mnesia directory (~s) not found. Creating fresh schema...~n", [MnesiaDir]), mnesia:stop(), mnesia:delete_schema([node()]), mnesia:create_schema([node()]), mnesia:start(), ok; true -> io:format("Mnesia directory exists (~s). Reusing existing schema.~n", [MnesiaDir]), case mnesia:system_info(is_running) of yes -> ok; _ -> mnesia:start() end end. join_cluster(Nodes) -> case mnesia:system_info(is_running) of yes -> mnesia:stop(); no -> ok end, application:set_env(mnesia, extra_db_nodes, Nodes), mnesia:start(), ensure_schema_disc(), wait_for_tables_available(), lists:foreach(fun add_local_disc_copy/1, ?DISC_TABLES), lists:foreach(fun add_local_ram_copy/1, ?RAM_TABLES). do_add_nodes(Nodes) -> ExtraNodes = application:get_env(eventhub, extra_db_nodes, []), application:set_env(eventhub, extra_db_nodes, Nodes ++ ExtraNodes), {ok, _} = mnesia:change_config(extra_db_nodes, Nodes), ensure_schema_disc(), wait_for_tables_available(), lists:foreach(fun add_local_disc_copy/1, ?DISC_TABLES), lists:foreach(fun add_local_ram_copy/1, ?RAM_TABLES). ensure_schema_disc() -> case lists:member(node(), mnesia:table_info(schema, disc_copies)) of false -> io:format("Changing schema copy to disc...~n"), case mnesia:change_table_copy_type(schema, node(), disc_copies) of {atomic, ok} -> ok; {aborted, {already_exists, _, _}} -> ok; {aborted, Reason} -> error({failed_schema_disc, Reason}) end; true -> ok end. add_local_disc_copy(Tab) -> case lists:member(node(), mnesia:table_info(Tab, disc_copies)) of false -> io:format("Adding local disc copy of table ~p...~n", [Tab]), case mnesia:add_table_copy(Tab, node(), disc_copies) of {atomic, ok} -> ok; {aborted, {already_exists, _}} -> ok; {aborted, Reason} -> io:format("Could not add disc copy for ~p: ~p~n", [Tab, Reason]) end; true -> ok end. add_local_ram_copy(Tab) -> case lists:member(node(), mnesia:table_info(Tab, ram_copies)) of false -> io:format("Adding local ram copy of table ~p...~n", [Tab]), case mnesia:add_table_copy(Tab, node(), ram_copies) of {atomic, ok} -> ok; {aborted, {already_exists, _}} -> ok; {aborted, Reason} -> io:format("Could not add ram copy for ~p: ~p~n", [Tab, Reason]) end; true -> ok end. wait_for_tables_available() -> lists:foreach(fun(Tab) -> wait_for_table(Tab) end, ?DISC_TABLES ++ ?RAM_TABLES). do_wait_for_tables() -> case mnesia:wait_for_tables(?TABLES, ?TABLE_WAIT_TIMEOUT) of ok -> ok; {timeout, Remaining} -> io:format("Mnesia wait_for_tables timeout, repairing ~p~n", [Remaining]), lists:foreach(fun(Tab) -> catch create_table(Tab), case lists:member(Tab, ?RAM_TABLES) of true -> catch add_local_ram_copy(Tab); false -> catch add_local_disc_copy(Tab) end end, Remaining), case mnesia:wait_for_tables(?TABLES, ?TABLE_WAIT_TIMEOUT) of ok -> ok; {timeout, Still} -> io:format("Mnesia tables still not ready: ~p~n", [Still]), {error, {tables_not_ready, Still}} end end. wait_for_table(Tab) -> case lists:member(Tab, mnesia:system_info(tables)) of true -> ok; false -> timer:sleep(100), wait_for_table(Tab) end. %% =================================================================== %% Автоматическая очистка мёртвых узлов %% =================================================================== start_cleanup_timer() -> erlang:send_after(?CLEANUP_INTERVAL, self(), {cleanup}), ok. prune_dead_nodes() -> AliveNodes = lists:filter(fun(Node) -> Node =/= node() andalso net_adm:ping(Node) =:= pong end, mnesia:system_info(db_nodes)), DeadNodes = mnesia:system_info(db_nodes) -- [node() | AliveNodes], lists:foreach(fun(Node) -> io:format("Removing dead node ~p from Mnesia schema...~n", [Node]), lists:foreach(fun(Tab) -> case lists:member(Node, mnesia:table_info(Tab, disc_copies)) of true -> catch mnesia:del_table_copy(Tab, Node); false -> ok end end, ?DISC_TABLES), lists:foreach(fun(Tab) -> case lists:member(Node, mnesia:table_info(Tab, ram_copies)) of true -> catch mnesia:del_table_copy(Tab, Node); false -> ok end end, ?RAM_TABLES), catch mnesia:del_table_copy(schema, Node) end, DeadNodes). %% =================================================================== %% Создание / открытие таблиц %% =================================================================== create_table(Table) -> Opts = table_opts(Table), case mnesia:create_table(Table, Opts) of {atomic, ok} -> ok; {aborted, {already_exists, _}} -> ok; {aborted, Reason} -> error({table_creation_failed, Table, Reason}) end. %% =================================================================== %% Опции хранения таблиц %% =================================================================== table_opts(user) -> [{disc_copies, [node()]}, {attributes, record_info(fields, user)}]; table_opts(admin) -> [{disc_copies, [node()]}, {attributes, record_info(fields, admin)}]; table_opts(calendar) -> [{disc_copies, [node()]}, {attributes, record_info(fields, calendar)}]; table_opts(calendar_share) -> [{disc_copies, [node()]}, {attributes, record_info(fields, calendar_share)}]; table_opts(calendar_specialist) -> [{disc_copies, [node()]}, {attributes, record_info(fields, calendar_specialist)}]; table_opts(event) -> [{disc_copies, [node()]}, {attributes, record_info(fields, event)}]; table_opts(recurrence_exception) -> [{disc_copies, [node()]}, {attributes, record_info(fields, recurrence_exception)}]; table_opts(booking) -> [{disc_copies, [node()]}, {attributes, record_info(fields, booking)}]; table_opts(review) -> [{disc_copies, [node()]}, {attributes, record_info(fields, review)}]; table_opts(report) -> [{disc_copies, [node()]}, {attributes, record_info(fields, report)}]; table_opts(banned_word) -> [{disc_copies, [node()]}, {attributes, record_info(fields, banned_word)}]; table_opts(automod_settings) -> [{disc_copies, [node()]}, {attributes, record_info(fields, automod_settings)}]; table_opts(automod_hit) -> [{disc_copies, [node()]}, {attributes, record_info(fields, automod_hit)}]; table_opts(ticket) -> [{disc_copies, [node()]}, {attributes, record_info(fields, ticket)}]; 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_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)}]; table_opts(admin_session) -> [{ram_copies, [node()]}, {attributes, record_info(fields, admin_session)}]; table_opts(auth_session) -> [{disc_copies, [node()]}, {attributes, record_info(fields, auth_session)}]; table_opts(node_metric) -> [{disc_copies, [node()]}, {local_content, true}, {attributes, record_info(fields, node_metric)}]. %% =================================================================== %% Индексы %% =================================================================== create_indices() -> mnesia:add_table_index(event, calendar_id), mnesia:add_table_index(event, title), mnesia:add_table_index(event, created_at), mnesia:add_table_index(event, start_time), mnesia:add_table_index(event, event_type), mnesia:add_table_index(event, master_id), mnesia:add_table_index(event, specialist_id), mnesia:add_table_index(event, status), mnesia:add_table_index(booking, event_id), mnesia:add_table_index(booking, user_id), mnesia:add_table_index(booking, status), mnesia:add_table_index(calendar, owner_id), mnesia:add_table_index(calendar, status), mnesia:add_table_index(calendar, short_name), mnesia:add_table_index(calendar, category), mnesia:add_table_index(calendar_specialist, calendar_id), mnesia:add_table_index(calendar_specialist, user_id), mnesia:add_table_index(user, nickname), mnesia:add_table_index(user, email), mnesia:add_table_index(verification, user_id), mnesia:add_table_index(notification, user_id), 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.