From 7c5f2d26b8ca83f54c9480bc7af331b3163dbfe8 Mon Sep 17 00:00:00 2001 From: Aleksey Sabilin Date: Tue, 7 Jul 2026 17:54:29 +0300 Subject: [PATCH] Mnesia migrations with global lock on startup. Refs EventHub/EventHubBack#24 --- src/eventhub_app.erl | 463 +++++++++++------------ src/infra/infra_mnesia.erl | 542 +++++++++++++-------------- src/infra/migration_engine.erl | 319 ++++++++++------ src/migrations/README.md | 66 ++-- test/unit/migration_engine_tests.erl | 63 ++++ 5 files changed, 806 insertions(+), 647 deletions(-) create mode 100644 test/unit/migration_engine_tests.erl diff --git a/src/eventhub_app.erl b/src/eventhub_app.erl index 14af3e5..fb9a914 100644 --- a/src/eventhub_app.erl +++ b/src/eventhub_app.erl @@ -1,232 +1,233 @@ --module(eventhub_app). --behaviour(application). --export([start/2, stop/1]). - -start(_StartType, _StartArgs) -> - case infra_sup:start_link() of - {ok, Pid} -> - % Определяем список узлов кластера, если режим CLUSTER_MODE=true - Nodes = case os:getenv("CLUSTER_MODE", "false") of - "true" -> - DnsName = os:getenv("DNS_NAME", "eventhub-node"), - try inet:getaddrs(DnsName, inet) of - {ok, IPs} when is_list(IPs), IPs /= [] -> - % Получаем имена всех узлов Erlang, зарегистрированных в EPMD - AllNodes = lists:flatmap(fun(IP) -> - case erl_epmd:names(IP) of - {ok, Names} -> - [list_to_atom(Name ++ "@" ++ Name) || {Name, _Port} <- Names, - lists:prefix("eventhub-node", Name)]; - _ -> [] - end - end, IPs), - % Исключаем свой узел, чтобы не подключаться к самому себе - AllNodes -- [node()]; - _ -> [] - catch - _:_ -> - io:format("DNS lookup failed, starting as first node~n"), - [] - end; - _ -> [] - end, - case Nodes of - [] -> - io:format("~nCluster: no nodes found or first node~n"); - _ -> - io:format("~nCluster: discovered nodes ~p, joining cluster~n", [Nodes]), - application:set_env(eventhub, extra_db_nodes, Nodes) - end, - application:ensure_all_started(mnesia), - ok = infra_mnesia:init_tables(), - ok = infra_mnesia:wait_for_tables(), - calendar_html_renderer:init_cache(), - application:ensure_all_started(cowboy), - start_http(), % Пользовательский API (8080) - start_admin_http(), % Административный API (8445) - start_swagger_http(), % Swagger UI и спецификация (8447) - application:ensure_all_started(prometheus), - application:ensure_all_started(prometheus_cowboy), - init_default_admins(), - {ok, Pid}; - Error -> - Error - end. - -stop(_State) -> ok. - -%% =================================================================== -%% Пользовательский HTTP (порт 8080) — только публичные эндпоинты -%% =================================================================== -start_http() -> - Port = get_env_int(http_port, 8080), - Dispatch = cowboy_router:compile([ - {'_', [ - {"/metrics/[:registry]", prometheus_cowboy2_handler, []}, - {"/health", handler_health, []}, - {"/v1/register", handler_register, []}, - {"/v1/verify", handler_verify, []}, - {"/v1/login", handler_login, []}, - {"/v1/refresh", handler_refresh, []}, - {"/v1/user/me", handler_user_me, []}, - {"/v1/user/bookings", handler_user_bookings, []}, - {"/v1/user/reviews", handler_user_reviews, []}, - {"/v1/search", handler_search, []}, - {"/v1/calendars", handler_calendars, []}, - {"/v1/calendars/:id", handler_calendar_by_id, []}, - {"/v1/calendars/:calendar_id/view", handler_calendar_view, []}, - {"/v1/calendars/:calendar_id/events", handler_events, []}, - {"/v1/events/:id", handler_event_by_id, []}, - {"/v1/events/:id/occurrences", handler_event_occurrences, []}, - {"/v1/events/:id/occurrences/:start_time", handler_event_occurrences, []}, - {"/v1/events/:id/bookings", handler_bookings, []}, - {"/v1/bookings/:id", handler_booking_by_id, []}, - {"/v1/reviews", handler_reviews, []}, - {"/v1/reviews/:id", handler_review_by_id, []}, - {"/v1/reports", handler_reports, []}, - {"/v1/tickets", handler_tickets, []}, - {"/v1/tickets/:id", handler_ticket_by_id, []}, - {"/v1/subscription", handler_subscription, []} - ]} %% 23 - ]), - Middlewares = [cowboy_router, cowboy_handler], - Env = #{dispatch => Dispatch}, - cowboy:start_clear(http, [{port, Port}], - #{env => Env, middlewares => Middlewares, - metrics_callback => fun prometheus_cowboy2_instrumenter:observe/1, - stream_handlers => [cowboy_metrics_h, cowboy_stream_h] - }), - io:format("HTTP server started on port ~p~n", [Port]). - -%% =================================================================== -%% Административный HTTP (порт 8445) — все админские эндпоинты -%% =================================================================== -start_admin_http() -> - PortAdmin = get_env_int(admin_http_port, 8445), - Dispatch = cowboy_router:compile([ - {'_', [ - % ================== БАЗОВЫЕ ================== - {"/admin/health", admin_handler_health, []}, - {"/v1/admin/stats", admin_handler_stats, []}, - {"/v1/admin/nodes/metrics", admin_handler_node_metrics, []}, - {"/v1/admin/login", admin_handler_login, []}, - {"/v1/admin/refresh", admin_handler_refresh, []}, - % ================== ПОЛЬЗОВАТЕЛИ ================== - {"/v1/admin/users", admin_handler_users, []}, - {"/v1/admin/users/stats", admin_handler_user_stats, []}, - {"/v1/admin/users/:id", admin_handler_user_by_id, []}, - {"/v1/admin/users/:id/verification-token", admin_handler_user_verification_token, []}, - % ================== КАЛЕНДАРИ ================== - {"/v1/admin/calendars", admin_handler_calendars, []}, - {"/v1/admin/calendars/stats", admin_handler_calendar_stats, []}, - {"/v1/admin/calendars/:id", admin_handler_calendar_by_id, []}, - % ================== СОБЫТИЯ ================== - {"/v1/admin/events", admin_handler_events, []}, - {"/v1/admin/events/stats", admin_handler_event_stats, []}, - {"/v1/admin/events/:id", admin_handler_event_by_id, []}, - % ================== ОТЧЁТЫ ================== - {"/v1/admin/reports", admin_handler_reports, []}, - {"/v1/admin/reports/stats", admin_handler_report_stats, []}, - {"/v1/admin/reports/:id", admin_handler_report_by_id, []}, - % ================== ОТЗЫВЫ ================== - {"/v1/admin/reviews", admin_handler_reviews, []}, - {"/v1/admin/reviews/stats", admin_handler_reviews_stats, []}, - {"/v1/admin/reviews/:id", admin_handler_reviews_by_id, []}, - % ================== БАН-СЛОВА ================== - {"/v1/admin/banned-words/batch", admin_handler_banned_words, []}, - {"/v1/admin/banned-words", admin_handler_banned_words, []}, - {"/v1/admin/banned-words/:word", admin_handler_banned_words, []}, - % ================== ТИКЕТЫ ================== - {"/v1/admin/tickets/stats", admin_handler_ticket_stats, []}, - {"/v1/admin/tickets/:id", admin_handler_ticket_by_id, []}, - {"/v1/admin/tickets", admin_handler_tickets, []}, - % ================== ПОДПИСКИ ================== - {"/v1/admin/subscriptions", admin_handler_subscriptions, []}, - {"/v1/admin/subscriptions/stats", admin_handler_subscription_stats, []}, - {"/v1/admin/subscriptions/:id", admin_handler_subscriptions_by_id, []}, - % ================== Управление ролями (только для superadmin) ================== - {"/v1/admin/me", admin_handler_me, []}, - {"/v1/admin/admins", admin_handler_admins, []}, - {"/v1/admin/admins/:id", admin_handler_admins_by_id, []}, - {"/v1/admin/audit", admin_handler_audit, []}, - % ================== МОДЕРАЦИЯ (общий маршрут) ================== - {"/v1/admin/:target_type/:id", admin_handler_moderation, []} - ]} - ]), - - Middlewares = [cowboy_router, cowboy_handler], - Env = #{dispatch => Dispatch}, - cowboy:start_clear(admin_http, [{port, PortAdmin}], - #{env => Env, middlewares => Middlewares, - metrics_callback => fun prometheus_cowboy2_instrumenter:observe/1, - stream_handlers => [cowboy_metrics_h, cowboy_stream_h] - }), - io:format("Admin HTTP server started on port ~p~n", [PortAdmin]), - - % WebSocket для пользователей - WsDispatch = cowboy_router:compile([{'_', [{"/ws", ws_handler, []}]}]), - PortWs = get_env_int(ws_port, 8081), - cowboy:start_clear(ws, [{port, PortWs}], #{env => #{dispatch => WsDispatch}}), - - % WebSocket для админов - AdminWsDispatch = cowboy_router:compile([{'_', [ - {"/admin/ws/metrics", admin_ws_handler, []}, - {"/admin/ws", admin_ws_handler, []} - ]}]), - PortAdminWs = get_env_int(admin_ws_port, 8446), - cowboy:start_clear(admin_ws, [{port, PortAdminWs}], #{env => #{dispatch => AdminWsDispatch}}), - - io:format("WebSocket started on ports ~p (user) and ~p (admin)~n", [PortWs, PortAdminWs]). - -%% =================================================================== -%% Swagger HTTP (порт 8447) — документация API -%% =================================================================== -start_swagger_http() -> - PortSwagger = get_env_int(swagger_http_port, 8447), - Dispatch = cowboy_router:compile([ - {'_', [ - {"/", swagger_docs_handler, []}, - {"/[...]", swagger_docs_handler, []} - ]} - ]), - Middlewares = [cowboy_router, cowboy_handler], - Env = #{dispatch => Dispatch}, - cowboy:start_clear(swagger_http, [{port, PortSwagger}], #{env => Env, middlewares => Middlewares}), - io:format("Swagger HTTP server started on port ~p~n", [PortSwagger]). - -%% ---------- Инициализация администраторов ---------- -init_default_admins() -> - case core_admin:list_all() of - [] -> - % Суперадмин - SuperEmail = list_to_binary(os:getenv("ADMIN_SUPER_EMAIL", "superadmin@eventhub.local")), - SuperPass = list_to_binary(os:getenv("ADMIN_SUPER_PASSWORD", "123456")), - {ok, _} = core_admin:create(SuperEmail, SuperPass, superadmin), - io:format("Default superadmin created: ~s~n", [SuperEmail]), - - % Админ - AdminEmail = list_to_binary(os:getenv("ADMIN_EMAIL", "admin@eventhub.local")), - AdminPass = list_to_binary(os:getenv("ADMIN_PASSWORD", "123456")), - {ok, _} = core_admin:create(AdminEmail, AdminPass, admin), - io:format("Default admin created: ~s~n", [AdminEmail]), - - % Модератор - ModerEmail = list_to_binary(os:getenv("ADMIN_MODER_EMAIL", "moderator@eventhub.local")), - ModerPass = list_to_binary(os:getenv("ADMIN_MODER_PASSWORD", "123456")), - {ok, _} = core_admin:create(ModerEmail, ModerPass, moderator), - io:format("Default moderator created: ~s~n", [ModerEmail]), - - % Поддержка - SupportEmail = list_to_binary(os:getenv("ADMIN_SUPPORT_EMAIL", "support@eventhub.local")), - SupportPass = list_to_binary(os:getenv("ADMIN_SUPPORT_PASSWORD", "123456")), - {ok, _} = core_admin:create(SupportEmail, SupportPass, support), - io:format("Default support created: ~s~n", [SupportEmail]); - _ -> - io:format("Admins already exist. Skipping creation.~n") - end. - -get_env_int(Key, Default) -> - case application:get_env(eventhub, Key, Default) of - Val when is_list(Val) -> list_to_integer(Val); - Val when is_integer(Val) -> Val +-module(eventhub_app). +-behaviour(application). +-export([start/2, stop/1]). + +start(_StartType, _StartArgs) -> + case infra_sup:start_link() of + {ok, Pid} -> + % Определяем список узлов кластера, если режим CLUSTER_MODE=true + Nodes = case os:getenv("CLUSTER_MODE", "false") of + "true" -> + DnsName = os:getenv("DNS_NAME", "eventhub-node"), + try inet:getaddrs(DnsName, inet) of + {ok, IPs} when is_list(IPs), IPs /= [] -> + % Получаем имена всех узлов Erlang, зарегистрированных в EPMD + AllNodes = lists:flatmap(fun(IP) -> + case erl_epmd:names(IP) of + {ok, Names} -> + [list_to_atom(Name ++ "@" ++ Name) || {Name, _Port} <- Names, + lists:prefix("eventhub-node", Name)]; + _ -> [] + end + end, IPs), + % Исключаем свой узел, чтобы не подключаться к самому себе + AllNodes -- [node()]; + _ -> [] + catch + _:_ -> + io:format("DNS lookup failed, starting as first node~n"), + [] + end; + _ -> [] + end, + case Nodes of + [] -> + io:format("~nCluster: no nodes found or first node~n"); + _ -> + io:format("~nCluster: discovered nodes ~p, joining cluster~n", [Nodes]), + application:set_env(eventhub, extra_db_nodes, Nodes) + end, + application:ensure_all_started(mnesia), + ok = infra_mnesia:init_tables(), + ok = infra_mnesia:wait_for_tables(), + ok = migration_engine:ensure_applied(), + calendar_html_renderer:init_cache(), + application:ensure_all_started(cowboy), + start_http(), % Пользовательский API (8080) + start_admin_http(), % Административный API (8445) + start_swagger_http(), % Swagger UI и спецификация (8447) + application:ensure_all_started(prometheus), + application:ensure_all_started(prometheus_cowboy), + init_default_admins(), + {ok, Pid}; + Error -> + Error + end. + +stop(_State) -> ok. + +%% =================================================================== +%% Пользовательский HTTP (порт 8080) — только публичные эндпоинты +%% =================================================================== +start_http() -> + Port = get_env_int(http_port, 8080), + Dispatch = cowboy_router:compile([ + {'_', [ + {"/metrics/[:registry]", prometheus_cowboy2_handler, []}, + {"/health", handler_health, []}, + {"/v1/register", handler_register, []}, + {"/v1/verify", handler_verify, []}, + {"/v1/login", handler_login, []}, + {"/v1/refresh", handler_refresh, []}, + {"/v1/user/me", handler_user_me, []}, + {"/v1/user/bookings", handler_user_bookings, []}, + {"/v1/user/reviews", handler_user_reviews, []}, + {"/v1/search", handler_search, []}, + {"/v1/calendars", handler_calendars, []}, + {"/v1/calendars/:id", handler_calendar_by_id, []}, + {"/v1/calendars/:calendar_id/view", handler_calendar_view, []}, + {"/v1/calendars/:calendar_id/events", handler_events, []}, + {"/v1/events/:id", handler_event_by_id, []}, + {"/v1/events/:id/occurrences", handler_event_occurrences, []}, + {"/v1/events/:id/occurrences/:start_time", handler_event_occurrences, []}, + {"/v1/events/:id/bookings", handler_bookings, []}, + {"/v1/bookings/:id", handler_booking_by_id, []}, + {"/v1/reviews", handler_reviews, []}, + {"/v1/reviews/:id", handler_review_by_id, []}, + {"/v1/reports", handler_reports, []}, + {"/v1/tickets", handler_tickets, []}, + {"/v1/tickets/:id", handler_ticket_by_id, []}, + {"/v1/subscription", handler_subscription, []} + ]} %% 23 + ]), + Middlewares = [cowboy_router, cowboy_handler], + Env = #{dispatch => Dispatch}, + cowboy:start_clear(http, [{port, Port}], + #{env => Env, middlewares => Middlewares, + metrics_callback => fun prometheus_cowboy2_instrumenter:observe/1, + stream_handlers => [cowboy_metrics_h, cowboy_stream_h] + }), + io:format("HTTP server started on port ~p~n", [Port]). + +%% =================================================================== +%% Административный HTTP (порт 8445) — все админские эндпоинты +%% =================================================================== +start_admin_http() -> + PortAdmin = get_env_int(admin_http_port, 8445), + Dispatch = cowboy_router:compile([ + {'_', [ + % ================== БАЗОВЫЕ ================== + {"/admin/health", admin_handler_health, []}, + {"/v1/admin/stats", admin_handler_stats, []}, + {"/v1/admin/nodes/metrics", admin_handler_node_metrics, []}, + {"/v1/admin/login", admin_handler_login, []}, + {"/v1/admin/refresh", admin_handler_refresh, []}, + % ================== ПОЛЬЗОВАТЕЛИ ================== + {"/v1/admin/users", admin_handler_users, []}, + {"/v1/admin/users/stats", admin_handler_user_stats, []}, + {"/v1/admin/users/:id", admin_handler_user_by_id, []}, + {"/v1/admin/users/:id/verification-token", admin_handler_user_verification_token, []}, + % ================== КАЛЕНДАРИ ================== + {"/v1/admin/calendars", admin_handler_calendars, []}, + {"/v1/admin/calendars/stats", admin_handler_calendar_stats, []}, + {"/v1/admin/calendars/:id", admin_handler_calendar_by_id, []}, + % ================== СОБЫТИЯ ================== + {"/v1/admin/events", admin_handler_events, []}, + {"/v1/admin/events/stats", admin_handler_event_stats, []}, + {"/v1/admin/events/:id", admin_handler_event_by_id, []}, + % ================== ОТЧЁТЫ ================== + {"/v1/admin/reports", admin_handler_reports, []}, + {"/v1/admin/reports/stats", admin_handler_report_stats, []}, + {"/v1/admin/reports/:id", admin_handler_report_by_id, []}, + % ================== ОТЗЫВЫ ================== + {"/v1/admin/reviews", admin_handler_reviews, []}, + {"/v1/admin/reviews/stats", admin_handler_reviews_stats, []}, + {"/v1/admin/reviews/:id", admin_handler_reviews_by_id, []}, + % ================== БАН-СЛОВА ================== + {"/v1/admin/banned-words/batch", admin_handler_banned_words, []}, + {"/v1/admin/banned-words", admin_handler_banned_words, []}, + {"/v1/admin/banned-words/:word", admin_handler_banned_words, []}, + % ================== ТИКЕТЫ ================== + {"/v1/admin/tickets/stats", admin_handler_ticket_stats, []}, + {"/v1/admin/tickets/:id", admin_handler_ticket_by_id, []}, + {"/v1/admin/tickets", admin_handler_tickets, []}, + % ================== ПОДПИСКИ ================== + {"/v1/admin/subscriptions", admin_handler_subscriptions, []}, + {"/v1/admin/subscriptions/stats", admin_handler_subscription_stats, []}, + {"/v1/admin/subscriptions/:id", admin_handler_subscriptions_by_id, []}, + % ================== Управление ролями (только для superadmin) ================== + {"/v1/admin/me", admin_handler_me, []}, + {"/v1/admin/admins", admin_handler_admins, []}, + {"/v1/admin/admins/:id", admin_handler_admins_by_id, []}, + {"/v1/admin/audit", admin_handler_audit, []}, + % ================== МОДЕРАЦИЯ (общий маршрут) ================== + {"/v1/admin/:target_type/:id", admin_handler_moderation, []} + ]} + ]), + + Middlewares = [cowboy_router, cowboy_handler], + Env = #{dispatch => Dispatch}, + cowboy:start_clear(admin_http, [{port, PortAdmin}], + #{env => Env, middlewares => Middlewares, + metrics_callback => fun prometheus_cowboy2_instrumenter:observe/1, + stream_handlers => [cowboy_metrics_h, cowboy_stream_h] + }), + io:format("Admin HTTP server started on port ~p~n", [PortAdmin]), + + % WebSocket для пользователей + WsDispatch = cowboy_router:compile([{'_', [{"/ws", ws_handler, []}]}]), + PortWs = get_env_int(ws_port, 8081), + cowboy:start_clear(ws, [{port, PortWs}], #{env => #{dispatch => WsDispatch}}), + + % WebSocket для админов + AdminWsDispatch = cowboy_router:compile([{'_', [ + {"/admin/ws/metrics", admin_ws_handler, []}, + {"/admin/ws", admin_ws_handler, []} + ]}]), + PortAdminWs = get_env_int(admin_ws_port, 8446), + cowboy:start_clear(admin_ws, [{port, PortAdminWs}], #{env => #{dispatch => AdminWsDispatch}}), + + io:format("WebSocket started on ports ~p (user) and ~p (admin)~n", [PortWs, PortAdminWs]). + +%% =================================================================== +%% Swagger HTTP (порт 8447) — документация API +%% =================================================================== +start_swagger_http() -> + PortSwagger = get_env_int(swagger_http_port, 8447), + Dispatch = cowboy_router:compile([ + {'_', [ + {"/", swagger_docs_handler, []}, + {"/[...]", swagger_docs_handler, []} + ]} + ]), + Middlewares = [cowboy_router, cowboy_handler], + Env = #{dispatch => Dispatch}, + cowboy:start_clear(swagger_http, [{port, PortSwagger}], #{env => Env, middlewares => Middlewares}), + io:format("Swagger HTTP server started on port ~p~n", [PortSwagger]). + +%% ---------- Инициализация администраторов ---------- +init_default_admins() -> + case core_admin:list_all() of + [] -> + % Суперадмин + SuperEmail = list_to_binary(os:getenv("ADMIN_SUPER_EMAIL", "superadmin@eventhub.local")), + SuperPass = list_to_binary(os:getenv("ADMIN_SUPER_PASSWORD", "123456")), + {ok, _} = core_admin:create(SuperEmail, SuperPass, superadmin), + io:format("Default superadmin created: ~s~n", [SuperEmail]), + + % Админ + AdminEmail = list_to_binary(os:getenv("ADMIN_EMAIL", "admin@eventhub.local")), + AdminPass = list_to_binary(os:getenv("ADMIN_PASSWORD", "123456")), + {ok, _} = core_admin:create(AdminEmail, AdminPass, admin), + io:format("Default admin created: ~s~n", [AdminEmail]), + + % Модератор + ModerEmail = list_to_binary(os:getenv("ADMIN_MODER_EMAIL", "moderator@eventhub.local")), + ModerPass = list_to_binary(os:getenv("ADMIN_MODER_PASSWORD", "123456")), + {ok, _} = core_admin:create(ModerEmail, ModerPass, moderator), + io:format("Default moderator created: ~s~n", [ModerEmail]), + + % Поддержка + SupportEmail = list_to_binary(os:getenv("ADMIN_SUPPORT_EMAIL", "support@eventhub.local")), + SupportPass = list_to_binary(os:getenv("ADMIN_SUPPORT_PASSWORD", "123456")), + {ok, _} = core_admin:create(SupportEmail, SupportPass, support), + io:format("Default support created: ~s~n", [SupportEmail]); + _ -> + io:format("Admins already exist. Skipping creation.~n") + end. + +get_env_int(Key, Default) -> + case application:get_env(eventhub, Key, Default) of + Val when is_list(Val) -> list_to_integer(Val); + Val when is_integer(Val) -> Val end. \ No newline at end of file diff --git a/src/infra/infra_mnesia.erl b/src/infra/infra_mnesia.erl index 0497c65..94c2cbf 100644 --- a/src/infra/infra_mnesia.erl +++ b/src/infra/infra_mnesia.erl @@ -1,273 +1,271 @@ -%% =================================================================== -%% 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([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, - ticket, subscription, - admin_audit, notification, - stats, node_metric, schema_migration -]). - --define(DISC_TABLES, ?TABLES -- [session, verification, admin_session, node_metric]). --define(TABLE_WAIT_TIMEOUT, 5000). --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). - -wait_for_tables() -> - gen_server:call(?MODULE, wait_for_tables). - -add_cluster_nodes(Nodes) -> - gen_server:call(?MODULE, {add_nodes, Nodes}). - -%% =================================================================== -%% 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 = migration_engine:init_migrations_table(), -%% _ = migration_engine:apply_pending(); //todo выключил - обваливает кластер, нужно разбираться - _ -> - ok = join_cluster(ExtraNodes) - end, - lists:foreach(fun create_table/1, ?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(), - ok = stats_collector:subscribe(), - 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) -> - mnesia:wait_for_tables(?TABLES, ?TABLE_WAIT_TIMEOUT), - {reply, ok, State}. - -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). - -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). - -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. - -wait_for_tables_available() -> - lists:foreach(fun(Tab) -> wait_for_table(Tab) end, ?DISC_TABLES). - -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), - 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(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) -> [{disc_copies, [node()]}, {attributes, record_info(fields, stats)}]; -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(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), +%% =================================================================== +%% 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([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, + ticket, subscription, + admin_audit, notification, + stats, node_metric, schema_migration +]). + +-define(DISC_TABLES, ?TABLES -- [session, verification, admin_session, node_metric]). +-define(TABLE_WAIT_TIMEOUT, 5000). +-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). + +wait_for_tables() -> + gen_server:call(?MODULE, wait_for_tables). + +add_cluster_nodes(Nodes) -> + gen_server:call(?MODULE, {add_nodes, Nodes}). + +%% =================================================================== +%% 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), + % Принудительное создание 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(), + ok = stats_collector:subscribe(), + 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) -> + mnesia:wait_for_tables(?TABLES, ?TABLE_WAIT_TIMEOUT), + {reply, ok, State}. + +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). + +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). + +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. + +wait_for_tables_available() -> + lists:foreach(fun(Tab) -> wait_for_table(Tab) end, ?DISC_TABLES). + +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), + 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(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) -> [{disc_copies, [node()]}, {attributes, record_info(fields, stats)}]; +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(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), ok. \ No newline at end of file diff --git a/src/infra/migration_engine.erl b/src/infra/migration_engine.erl index acf2a92..bfee2a8 100644 --- a/src/infra/migration_engine.erl +++ b/src/infra/migration_engine.erl @@ -1,121 +1,198 @@ --module(migration_engine). --behaviour(gen_server). - --include("records.hrl"). - -%% API --export([start_link/0, init_migrations_table/0, apply_pending/0, - rollback/1, status/0]). - -%% gen_server callbacks --export([init/1, handle_call/3, handle_cast/2, handle_info/2, - terminate/2, code_change/3]). - --define(TABLE, schema_migration). - -%% ------------------------------ -%% API -%% ------------------------------ - -start_link() -> - gen_server:start_link({local, ?MODULE}, ?MODULE, [], []). - -init_migrations_table() -> - gen_server:call(?MODULE, init_table). - -apply_pending() -> - gen_server:call(?MODULE, apply_pending). - -rollback(Version) -> - gen_server:call(?MODULE, {rollback, Version}). - -status() -> - gen_server:call(?MODULE, status). - -%% ------------------------------ -%% gen_server callbacks -%% ------------------------------ - -init([]) -> - {ok, #{}}. - -handle_call(init_table, _From, State) -> - case lists:member(?TABLE, mnesia:system_info(tables)) of - true -> ok; - false -> - mnesia:create_table(?TABLE, [ - {disc_copies, [node()]}, - {attributes, record_info(fields, schema_migration)}, - {type, set} - ]) - end, - infra_mnesia:wait_for_table(?TABLE), - {reply, ok, State}; - -handle_call(apply_pending, _From, State) -> - Result = do_apply_pending(), - {reply, Result, State}; - -handle_call({rollback, Version}, _From, State) -> - Result = do_rollback(Version), - {reply, Result, State}; - -handle_call(status, _From, State) -> - Applied = applied_versions(), - Pending = pending_versions() -- Applied, - {reply, #{applied => lists:map(fun atom_to_list/1, Applied), - pending => lists:map(fun atom_to_list/1, Pending)}, State}. - -handle_cast(_Msg, State) -> {noreply, State}. -handle_info(_Msg, State) -> {noreply, State}. -terminate(_Reason, _State) -> ok. -code_change(_OldVsn, State, _Extra) -> {ok, State}. - -%% ------------------------------ -%% Внутренняя логика -%% ------------------------------ - -do_apply_pending() -> - Pending = pending_versions() -- applied_versions(), - lists:foreach(fun(Module) -> code:ensure_loaded(Module) end, Pending), - lists:foldl(fun(Version, Acc) -> - try Version:up() of - _ -> mark_applied(Version), Acc - catch _:Reason -> - [{error, atom_to_list(Version), Reason} | Acc] - end - end, [], Pending). - -do_rollback(TargetStr) -> - Target = list_to_atom(TargetStr), - Applied = applied_versions(), - ToRollback = lists:sort(fun(A,B) -> A > B end, - [V || V <- Applied, V > Target]), - lists:foreach(fun(Version) -> - code:ensure_loaded(Version), - try Version:down() of - _ -> unmark_applied(Version) - catch _:Reason -> - throw({rollback_failed, atom_to_list(Version), Reason}) - end - end, ToRollback). - -applied_versions() -> - [list_to_atom(V) || #schema_migration{version = V} <- - mnesia:dirty_match_object(#schema_migration{_ = '_'})]. - -pending_versions() -> - AllMods = code:all_available(), - [list_to_atom(Module) || Module <- extract_module_names(AllMods), lists:prefix("20", Module)]. - -extract_module_names(ModInfoList) -> - [Name || {Name, _, _} <- ModInfoList]. - -mark_applied(Version) -> - mnesia:dirty_write(#schema_migration{ - version = atom_to_list(Version), - applied_at = calendar:local_time() - }). - -unmark_applied(Version) -> - mnesia:dirty_delete({?TABLE, atom_to_list(Version)}). \ No newline at end of file +-module(migration_engine). +-behaviour(gen_server). + +-include("records.hrl"). + +%% API +-export([start_link/0, ensure_applied/0, apply_pending/0, + rollback/1, status/0]). + +%% gen_server callbacks +-export([init/1, handle_call/3, handle_cast/2, handle_info/2, + terminate/2, code_change/3]). + +-define(TABLE, schema_migration). +-define(LOCK_VERSION, "__migration_lock__"). +-define(LOCK_STALE_SECONDS, 300). +-define(WAIT_TIMEOUT_MS, 120000). +-define(WAIT_INTERVAL_MS, 200). + +%% Упорядоченный реестр миграций (добавлять новые модули в конец списка). +-define(ALL_MIGRATIONS, [ + '20260501120000_base_schema', + '20260504150000_test_migration' +]). + +%% ------------------------------ +%% API +%% ------------------------------ + +start_link() -> + gen_server:start_link({local, ?MODULE}, ?MODULE, [], []). + +%% @doc Безопасно применить pending-миграции при старте (с global lock в Mnesia). +-spec ensure_applied() -> ok | {error, term()}. +ensure_applied() -> + gen_server:call(?MODULE, ensure_applied, infinity). + +apply_pending() -> + gen_server:call(?MODULE, apply_pending). + +rollback(Version) -> + gen_server:call(?MODULE, {rollback, Version}). + +status() -> + gen_server:call(?MODULE, status). + +%% ------------------------------ +%% gen_server callbacks +%% ------------------------------ + +init([]) -> + {ok, #{}}. + +handle_call(ensure_applied, _From, State) -> + Result = do_ensure_applied(), + {reply, Result, State}; + +handle_call(apply_pending, _From, State) -> + Result = do_apply_pending(), + {reply, Result, State}; + +handle_call({rollback, Version}, _From, State) -> + Result = do_rollback(Version), + {reply, Result, State}; + +handle_call(status, _From, State) -> + Applied = applied_versions(), + Pending = pending_versions() -- Applied, + {reply, #{applied => [atom_to_list(V) || V <- Applied], + pending => [atom_to_list(V) || V <- Pending]}, State}. + +handle_cast(_Msg, State) -> {noreply, State}. +handle_info(_Msg, State) -> {noreply, State}. +terminate(_Reason, _State) -> ok. +code_change(_OldVsn, State, _Extra) -> {ok, State}. + +%% ------------------------------ +%% Координация при старте +%% ------------------------------ + +do_ensure_applied() -> + case all_applied() of + true -> ok; + false -> + case acquire_lock() of + ok -> + try + case do_apply_pending() of + ok -> ok; + {error, Reason} -> {error, Reason} + end + after + release_lock() + end; + {error, locked} -> + wait_until_applied(?WAIT_TIMEOUT_MS) + end + end. + +all_applied() -> + pending_versions() -- applied_versions() =:= []. + +wait_until_applied(TimeoutMs) when TimeoutMs =< 0 -> + {error, migration_timeout}; +wait_until_applied(TimeoutMs) -> + case all_applied() of + true -> ok; + false -> + timer:sleep(?WAIT_INTERVAL_MS), + wait_until_applied(TimeoutMs - ?WAIT_INTERVAL_MS) + end. + +acquire_lock() -> + Now = calendar:universal_time(), + case mnesia:sync_transaction(fun() -> try_acquire_lock(Now) end) of + {atomic, ok} -> ok; + {aborted, locked} -> {error, locked}; + {aborted, Reason} -> {error, Reason} + end. + +try_acquire_lock(Now) -> + case mnesia:read(?TABLE, ?LOCK_VERSION, read) of + [] -> + mnesia:write(#schema_migration{version = ?LOCK_VERSION, applied_at = Now}), + ok; + [#schema_migration{applied_at = At}] -> + case lock_stale(At, Now) of + true -> + mnesia:write(#schema_migration{version = ?LOCK_VERSION, applied_at = Now}), + ok; + false -> + mnesia:abort(locked) + end + end. + +lock_stale(At, Now) -> + SecAt = calendar:datetime_to_gregorian_seconds(At), + SecNow = calendar:datetime_to_gregorian_seconds(Now), + SecNow - SecAt > ?LOCK_STALE_SECONDS. + +release_lock() -> + mnesia:sync_transaction(fun() -> + mnesia:delete({?TABLE, ?LOCK_VERSION}) + end), + ok. + +%% ------------------------------ +%% Применение / откат +%% ------------------------------ + +do_apply_pending() -> + Pending = pending_versions() -- applied_versions(), + lists:foreach(fun(Module) -> code:ensure_loaded(Module) end, Pending), + lists:foldl(fun(Version, ok) -> + try Version:up() of + _ -> + mark_applied(Version), + ok + catch + _:Reason -> + {error, {migration_failed, atom_to_list(Version), Reason}} + end; + (_, Acc) -> Acc + end, ok, Pending). + +do_rollback(TargetStr) -> + Target = list_to_atom(TargetStr), + Applied = applied_versions(), + ToRollback = lists:sort(fun(A, B) -> A > B end, + [V || V <- Applied, V > Target]), + lists:foreach(fun(Version) -> + code:ensure_loaded(Version), + try Version:down() of + _ -> unmark_applied(Version) + catch _:Reason -> + throw({rollback_failed, atom_to_list(Version), Reason}) + end + end, ToRollback). + +applied_versions() -> + [list_to_atom(V) || #schema_migration{version = V} <- + mnesia:dirty_match_object(#schema_migration{_ = '_'}), + not is_lock_version(V)]. + +pending_versions() -> + ?ALL_MIGRATIONS. + +is_lock_version(?LOCK_VERSION) -> true; +is_lock_version(_) -> false. + +mark_applied(Version) -> + mnesia:dirty_write(#schema_migration{ + version = atom_to_list(Version), + applied_at = calendar:universal_time() + }). + +unmark_applied(Version) -> + mnesia:dirty_delete({?TABLE, atom_to_list(Version)}). diff --git a/src/migrations/README.md b/src/migrations/README.md index 8a9dff9..1f72df9 100644 --- a/src/migrations/README.md +++ b/src/migrations/README.md @@ -1,23 +1,43 @@ -# Миграции схемы данных EventHub - -## Применение миграций -При старте приложения автоматически выполняются все неприменённые миграции. -Для ручного запуска можно вызвать: -migration_engine:apply_pending(). - -## Создание новой миграции -1. Создайте файл в `priv/migrations/` с именем вида `YYYYMMDDHHMMSS_описание.erl`. -2. Реализуйте поведение `db_migration` (функции `up/0` и `down/0`). -3. При очередном запуске приложения миграция будет применена. - -## Откат миграций - migration_engine:rollback("20260501120000_base_schema"). - -## Плавающее обновление (rolling update) -1. Переведите узел в режим обслуживания. -2. Выполните бэкап: `mnesia:backup("backup_node.bak")`. -3. Обновите код приложения (git pull / rsync). -4. Перезапустите узел – миграции применятся автоматически. -5. Убедитесь в согласованности данных. -6. Верните узел в работу. -7. Повторите для остальных узлов. \ No newline at end of file +# Миграции схемы данных EventHub + +## Когда применяются + +После `infra_mnesia:init_tables()` и `wait_for_tables()` приложение вызывает +`migration_engine:ensure_applied/0`: + +1. Узел, захвативший **глобальную блокировку** в Mnesia (`schema_migration`, ключ `__migration_lock__`), применяет все pending-миграции. +2. Остальные узлы (join кластера) **ждут**, пока pending-список станет пустым, и не стартуют HTTP до синхронизации. + +Базовые таблицы создаёт `infra_mnesia` при старте. Миграции — только для **инкрементальных** изменений (индексы, трансформация данных, новые поля). + +## Создание новой миграции + +1. Создайте файл в `src/migrations/` с именем `YYYYMMDDHHMMSS_описание.erl`. +2. Реализуйте `up/0` и `down/0`. +3. Добавьте модуль в список `?ALL_MIGRATIONS` в `src/infra/migration_engine.erl` (в конец, по возрастанию версии). +4. При следующем старте миграция применится автоматически (на узле с lock). + +## Ручной запуск + +```erlang +migration_engine:apply_pending(). +migration_engine:status(). +migration_engine:rollback("20260501120000_base_schema"). +``` + +## Плавающее обновление (rolling update) + +1. Переведите узел в режим обслуживания. +2. Выполните бэкап: `mnesia:backup("backup_node.bak")`. +3. Обновите код приложения (git pull / rsync). +4. Перезапустите узел — `ensure_applied/0` применит только новые миграции (один узел с lock). +5. Убедитесь: `migration_engine:status()` — `pending => []` на всех узлах. +6. Верните узел в работу. +7. Повторите для остальных узлов. + +## Восстановление после сбоя + +Если узел упал во время миграции, lock в `schema_migration` может остаться. +Он считается устаревшим через **5 минут** (`?LOCK_STALE_SECONDS`) — следующий стартующий узел перехватит lock и продолжит. + +При ошибке в `up/0` приложение не стартует (`{error, {migration_failed, ...}}`). Откатите код или исправьте миграцию, при необходимости восстановите из `mnesia:backup/1`. diff --git a/test/unit/migration_engine_tests.erl b/test/unit/migration_engine_tests.erl new file mode 100644 index 0000000..224de15 --- /dev/null +++ b/test/unit/migration_engine_tests.erl @@ -0,0 +1,63 @@ +-module(migration_engine_tests). +-include_lib("eunit/include/eunit.hrl"). +-include("records.hrl"). + +-define(LOCK_VERSION, "__migration_lock__"). + +setup() -> + mnesia:stop(), + mnesia:delete_schema([node()]), + mnesia:create_schema([node()]), + mnesia:start(), + mnesia:create_table(schema_migration, [ + {disc_copies, [node()]}, + {attributes, record_info(fields, schema_migration)}, + {type, set} + ]), + {ok, _} = migration_engine:start_link(), + ok. + +cleanup(_) -> + catch gen_server:stop(migration_engine), + catch mnesia:delete_table(schema_migration), + mnesia:stop(), + mnesia:delete_schema([node()]), + ok. + +migration_engine_test_() -> + {foreach, fun setup/0, fun cleanup/1, [ + {"ensure_applied applies pending migrations", fun test_ensure_applied/0}, + {"ensure_applied is idempotent", fun test_ensure_applied_idempotent/0}, + {"lock record is not treated as migration", fun test_lock_not_applied/0}, + {"join node waits until migrations applied", fun test_join_wait/0} + ]}. + +test_ensure_applied() -> + ?assertEqual(ok, migration_engine:ensure_applied()), + Status = migration_engine:status(), + ?assertEqual([], maps:get(pending, Status)), + Applied = maps:get(applied, Status), + ?assert(lists:member("20260501120000_base_schema", Applied)), + ?assert(lists:member("20260504150000_test_migration", Applied)). + +test_ensure_applied_idempotent() -> + ?assertEqual(ok, migration_engine:ensure_applied()), + ?assertEqual(ok, migration_engine:ensure_applied()). + +test_lock_not_applied() -> + ?assertEqual(ok, migration_engine:ensure_applied()), + Status = migration_engine:status(), + ?assertNot(lists:member(?LOCK_VERSION, maps:get(applied, Status))). + +test_join_wait() -> + Now = calendar:universal_time(), + mnesia:dirty_write(#schema_migration{version = ?LOCK_VERSION, applied_at = Now}), + spawn(fun() -> + timer:sleep(100), + mnesia:dirty_write(#schema_migration{ + version = "20260501120000_base_schema", applied_at = Now}), + mnesia:dirty_write(#schema_migration{ + version = "20260504150000_test_migration", applied_at = Now}), + mnesia:dirty_delete({schema_migration, ?LOCK_VERSION}) + end), + ?assertEqual(ok, migration_engine:ensure_applied()).