diff --git a/include/records.hrl b/include/records.hrl index f628c67..b4f289b 100644 --- a/include/records.hrl +++ b/include/records.hrl @@ -252,6 +252,25 @@ timestamp :: calendar:datetime() }). +-record(node_metric, { + timestamp :: calendar:datetime(), + node :: atom(), + memory_total :: non_neg_integer(), + memory_processes :: non_neg_integer(), + memory_atom :: non_neg_integer(), + memory_binary :: non_neg_integer(), + memory_ets :: non_neg_integer(), + process_count :: non_neg_integer(), + run_queue :: non_neg_integer(), + mnesia_commits :: non_neg_integer(), + mnesia_failures :: non_neg_integer(), + table_sizes :: map(), + active_sessions :: non_neg_integer(), + ws_connections :: non_neg_integer(), + io_in :: non_neg_integer(), + io_out :: non_neg_integer() +}). + -record(schema_migration, { version :: string(), applied_at :: calendar:datetime() diff --git a/src/core/core_admin_session.erl b/src/core/core_admin_session.erl index 6e48093..9af17b4 100644 --- a/src/core/core_admin_session.erl +++ b/src/core/core_admin_session.erl @@ -26,6 +26,7 @@ create(AdminId, RefreshToken) -> type = refresh }, mnesia:dirty_write(Session), + core_counters:inc(admin_active_sessions), ok. %%%------------------------------------------------------------------- @@ -44,6 +45,7 @@ validate(Token) -> true -> {ok, Session#admin_session.admin_id}; false -> mnesia:dirty_delete({admin_session, Token}), + core_counters:dec(admin_active_sessions), {error, expired} end; [] -> {error, not_found} @@ -56,4 +58,5 @@ validate(Token) -> -spec delete(Token :: binary()) -> ok. delete(Token) -> mnesia:dirty_delete({admin_session, Token}), + core_counters:dec(admin_active_sessions), ok. \ No newline at end of file diff --git a/src/core/core_counters.erl b/src/core/core_counters.erl new file mode 100644 index 0000000..98cc658 --- /dev/null +++ b/src/core/core_counters.erl @@ -0,0 +1,30 @@ +-module(core_counters). +-export([inc/1, dec/1, get/1, read_and_reset/1]). + +-define(TABLE, eventhub_counters). + +inc(Key) -> + try ets:update_counter(?TABLE, Key, {2, 1}) % позиция 2 — значение + catch error:badarg -> + ets:insert(?TABLE, {Key, 1}) + end. + +dec(Key) -> + try ets:update_counter(?TABLE, Key, {2, -1}) + catch error:badarg -> + ets:insert(?TABLE, {Key, 0}) + end. + +get(Key) -> + case ets:lookup(?TABLE, Key) of + [{Key, Val}] -> Val; + [] -> 0 + end. + +read_and_reset(Key) -> + case ets:lookup(?TABLE, Key) of + [{Key, Val}] -> + ets:insert(?TABLE, {Key, 0}), + Val; + [] -> 0 + end. \ No newline at end of file diff --git a/src/core/core_node_metric.erl b/src/core/core_node_metric.erl new file mode 100644 index 0000000..e8d2100 --- /dev/null +++ b/src/core/core_node_metric.erl @@ -0,0 +1,43 @@ +-module(core_node_metric). +-include("records.hrl"). +-export([save/1, get_by_timerange/3, delete_old/1]). + +%%%------------------------------------------------------------------- +%%% @doc Сохраняет запись метрики. +%%% @end +%%%------------------------------------------------------------------- +-spec save(#node_metric{}) -> ok. +save(Metric) -> + mnesia:dirty_write(Metric), + ok. + +%%%------------------------------------------------------------------- +%%% @doc Возвращает метрики за период, отсортированные по времени. +%%% @end +%%%------------------------------------------------------------------- +-spec get_by_timerange(From :: calendar:datetime(), To :: calendar:datetime(), Node :: atom()) -> + [#node_metric{}]. +get_by_timerange(From, To, Node) -> + All = mnesia:dirty_match_object(#node_metric{node = Node, _ = '_'}), + Filtered = lists:filter(fun(M) -> + M#node_metric.timestamp >= From andalso M#node_metric.timestamp =< To + end, All), + lists:sort(fun(A, B) -> A#node_metric.timestamp =< B#node_metric.timestamp end, Filtered). + +%%%------------------------------------------------------------------- +%%% @doc Удаляет метрики старше Seconds секунд. +%%% @end +%%%------------------------------------------------------------------- +-spec delete_old(Seconds :: non_neg_integer()) -> ok. +delete_old(Seconds) -> + Now = calendar:universal_time(), + Threshold = calendar:gregorian_seconds_to_datetime( + calendar:datetime_to_gregorian_seconds(Now) - Seconds), + All = mnesia:dirty_match_object(#node_metric{_ = '_'}), + lists:foreach(fun(M) -> + case M#node_metric.timestamp < Threshold of + true -> mnesia:dirty_delete_object(M); + false -> ok + end + end, All), + ok. \ No newline at end of file diff --git a/src/core/core_session.erl b/src/core/core_session.erl index ce93abd..7928ea7 100644 --- a/src/core/core_session.erl +++ b/src/core/core_session.erl @@ -26,6 +26,7 @@ create(UserId, RefreshToken) -> type = refresh }, mnesia:dirty_write(Session), + core_counters:inc(active_sessions), ok. %%%------------------------------------------------------------------- @@ -49,6 +50,7 @@ validate(Token) -> end; false -> mnesia:dirty_delete({session, Token}), + core_counters:dec(active_sessions), {error, expired} end; [] -> {error, not_found} @@ -61,4 +63,5 @@ validate(Token) -> -spec delete(Token :: binary()) -> ok. delete(Token) -> mnesia:dirty_delete({session, Token}), + core_counters:dec(active_sessions), ok. \ No newline at end of file diff --git a/src/eventhub_app.erl b/src/eventhub_app.erl index 247c442..14af3e5 100644 --- a/src/eventhub_app.erl +++ b/src/eventhub_app.erl @@ -3,7 +3,6 @@ -export([start/2, stop/1]). start(_StartType, _StartArgs) -> - pg:start_link(), case infra_sup:start_link() of {ok, Pid} -> % Определяем список узлов кластера, если режим CLUSTER_MODE=true @@ -109,6 +108,7 @@ start_admin_http() -> % ================== БАЗОВЫЕ ================== {"/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, []}, % ================== ПОЛЬЗОВАТЕЛИ ================== @@ -169,7 +169,10 @@ start_admin_http() -> cowboy:start_clear(ws, [{port, PortWs}], #{env => #{dispatch => WsDispatch}}), % WebSocket для админов - AdminWsDispatch = cowboy_router:compile([{'_', [{"/admin/ws", admin_ws_handler, []}]}]), + 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}}), diff --git a/src/handlers/admin/admin_handler_node_metrics.erl b/src/handlers/admin/admin_handler_node_metrics.erl new file mode 100644 index 0000000..342f493 --- /dev/null +++ b/src/handlers/admin/admin_handler_node_metrics.erl @@ -0,0 +1,162 @@ +%%%------------------------------------------------------------------- +%%% @doc Административный обработчик исторических метрик узлов. +%%% GET /v1/admin/nodes/metrics – возвращает массив метрик за период. +%%% @end +%%%------------------------------------------------------------------- +-module(admin_handler_node_metrics). +-behaviour(cowboy_handler). +-export([init/2]). +-export([trails/0]). + +-include("records.hrl"). + +%%% cowboy_handler callback +-spec init(cowboy_req:req(), any()) -> {ok, binary(), cowboy_req:req()}. +init(Req, _Opts) -> + case cowboy_req:method(Req) of + <<"GET">> -> get_metrics(Req); + _ -> handler_utils:send_error(Req, 405, <<"Method not allowed">>) + end. + +%%% Swagger metadata +-spec trails() -> [map()]. +trails() -> + [#{path => <<"/v1/admin/nodes/metrics">>, + method => <<"GET">>, + description => <<"Get historical node metrics">>, + tags => [<<"Node Metrics">>], + parameters => [ + #{name => <<"node">>, in => <<"query">>, schema => #{type => string}, description => <<"Node name, default current">>}, + #{name => <<"from">>, in => <<"query">>, required => true, schema => #{type => string, format => <<"date-time">>}}, + #{name => <<"to">>, in => <<"query">>, required => true, schema => #{type => string, format => <<"date-time">>}} + ], + responses => #{ + 200 => #{ + description => <<"Array of node metrics">>, + content => #{<<"application/json">> => #{schema => #{ + type => array, + items => metric_schema() + }}} + }, + 400 => #{description => <<"Missing required parameters">>}, + 403 => #{description => <<"Admin access required">>} + }} + ]. + +metric_schema() -> + #{ + type => object, + properties => #{ + <<"timestamp">> => #{ + type => string, + format => <<"date-time">>, + description => <<"Metric collection timestamp (ISO8601)">> + }, + <<"node">> => #{ + type => string, + description => <<"Erlang node name">> + }, + <<"memory_total">> => #{ + type => integer, + description => <<"Total memory allocated by the Erlang VM (bytes)">> + }, + <<"memory_processes">> => #{ + type => integer, + description => <<"Memory used by processes (bytes)">> + }, + <<"memory_atom">> => #{ + type => integer, + description => <<"Memory used by atoms (bytes)">> + }, + <<"memory_binary">> => #{ + type => integer, + description => <<"Memory used by binaries (bytes)">> + }, + <<"memory_ets">> => #{ + type => integer, + description => <<"Memory used by ETS tables (bytes)">> + }, + <<"process_count">> => #{ + type => integer, + description => <<"Number of live Erlang processes">> + }, + <<"run_queue">> => #{ + type => integer, + description => <<"Length of the scheduler run queue">> + }, + <<"mnesia_commits">> => #{ + type => integer, + description => <<"Number of successful Mnesia transactions since node start">> + }, + <<"mnesia_failures">> => #{ + type => integer, + description => <<"Number of failed Mnesia transactions since node start">> + }, + <<"table_sizes">> => #{ + type => object, + description => <<"Map of Mnesia table names to their sizes (number of records)">>, + additionalProperties => #{type => integer} + }, + <<"active_sessions">> => #{ + type => integer, + description => <<"Total number of active user and admin sessions">> + }, + <<"ws_connections">> => #{ + type => integer, + description => <<"Number of open WebSocket connections">> + }, + <<"io_in">> => #{ + type => integer, + description => <<"Bytes received on sockets since last collection (delta)">> + }, + <<"io_out">> => #{ + type => integer, + description => <<"Bytes sent over sockets since last collection (delta)">> + } + } + }. + +%%% Internal functions + +%% @doc Возвращает исторические метрики по заданному периоду и узлу. +-spec get_metrics(cowboy_req:req()) -> {ok, binary(), cowboy_req:req()}. +get_metrics(Req) -> + case handler_utils:auth_admin(Req) of + {ok, _AdminId, Req1} -> + Qs = cowboy_req:parse_qs(Req1), + From = handler_utils:parse_datetime_qs(proplists:get_value(<<"from">>, Qs)), + To = handler_utils:parse_datetime_qs(proplists:get_value(<<"to">>, Qs)), + Node = proplists:get_value(<<"node">>, Qs, node()), + case {From, To} of + {undefined, _} -> handler_utils:send_error(Req1, 400, <<"Missing from">>); + {_, undefined} -> handler_utils:send_error(Req1, 400, <<"Missing to">>); + {F, T} -> + Metrics = core_node_metric:get_by_timerange(F, T, Node), + Json = [metric_to_map(M) || M <- Metrics], + handler_utils:send_json(Req1, 200, Json) + end; + {error, Code, Msg, Req1} -> + handler_utils:send_error(Req1, Code, Msg) + end. + +%% @private Преобразует запись метрики в JSON-совместимую карту. +-spec metric_to_map(#node_metric{}) -> map(). +metric_to_map(M) -> + #{ + <<"timestamp">> => handler_utils:datetime_to_iso8601(M#node_metric.timestamp), + <<"node">> => atom_to_binary(M#node_metric.node, utf8), + <<"memory_total">> => M#node_metric.memory_total, + <<"memory_processes">> => M#node_metric.memory_processes, + <<"memory_atom">> => M#node_metric.memory_atom, + <<"memory_binary">> => M#node_metric.memory_binary, + <<"memory_ets">> => M#node_metric.memory_ets, + <<"process_count">> => M#node_metric.process_count, + <<"run_queue">> => M#node_metric.run_queue, + <<"mnesia_commits">> => M#node_metric.mnesia_commits, + <<"mnesia_failures">> => M#node_metric.mnesia_failures, + <<"table_sizes">> => M#node_metric.table_sizes, + <<"active_sessions">> => M#node_metric.active_sessions, + <<"ws_connections">> => M#node_metric.ws_connections, + <<"io_in">> => M#node_metric.io_in, + <<"io_out">> => M#node_metric.io_out + }. \ No newline at end of file diff --git a/src/handlers/admin/admin_ws_handler.erl b/src/handlers/admin/admin_ws_handler.erl index 4c3b3fa..a7ed080 100644 --- a/src/handlers/admin/admin_ws_handler.erl +++ b/src/handlers/admin/admin_ws_handler.erl @@ -1,11 +1,13 @@ %%%------------------------------------------------------------------- %%% @doc Административный WebSocket-обработчик. %%% Устанавливает WebSocket-соединение после проверки JWT-токена -%%% и подписывает администратора на каналы уведомлений. +%%% и подписывает администратора на каналы уведомлений, включая +%%% канал метрик узлов `node_metrics`. %%% @end %%%------------------------------------------------------------------- -module(admin_ws_handler). -behaviour(cowboy_websocket). +-include("records.hrl"). -export([init/2]). -export([websocket_init/1]). @@ -30,7 +32,7 @@ init(Req, _Opts) -> Resp = cowboy_req:reply(401, #{}, <<"Missing token">>, Req), {ok, Resp, undefined}; Token -> - io:format("[ADMIN_WS] Token received: ~s...~n", [binary_part(Token, 0, 30)]), + io:format("[ADMIN_WS] Token received: ~s...~n", [safe_slice(Token, 30)]), case logic_auth:verify_jwt(Token) of {ok, UserId, Role} -> io:format("[ADMIN_WS] UserId: ~s, Role: ~s~n", [UserId, Role]), @@ -57,8 +59,11 @@ init(Req, _Opts) -> %% @doc Вызывается после установки WebSocket-соединения. -spec websocket_init(#state{}) -> {ok, #state{}}. websocket_init(State) -> + core_counters:inc(admin_ws_connections), io:format("[ADMIN_WS] WebSocket initialized for admin ~s~n", [State#state.admin_id]), pg:join(eventhub_admin_ws, self()), + % Подписываемся на метрики узлов + pg:join(node_metrics, self()), {ok, State}. %% @doc Обрабатывает входящие текстовые сообщения (subscribe/unsubscribe/ping). @@ -84,7 +89,7 @@ websocket_handle({text, Msg}, State) -> websocket_handle(_Frame, State) -> {ok, State}. -%% @doc Отправляет административное уведомление через WebSocket. +%% @doc Отправляет административное уведомление или метрику узла через WebSocket. -spec websocket_info(term(), #state{}) -> {reply, {text, binary()}, #state{}} | {ok, #state{}}. websocket_info({admin_notification, Type, Data}, State) -> Msg = jsx:encode(#{ @@ -93,6 +98,12 @@ websocket_info({admin_notification, Type, Data}, State) -> timestamp => os:system_time(seconds) }), {reply, {text, Msg}, State}; +websocket_info({node_metric, Metric}, State) -> + JSON = jsx:encode(#{ + <<"type">> => <<"node_metric">>, + <<"data">> => metric_to_map(Metric) + }), + {reply, {text, JSON}, State}; websocket_info(_Info, State) -> {ok, State}. @@ -100,4 +111,37 @@ websocket_info(_Info, State) -> -spec terminate(term(), cowboy_req:req(), #state{}) -> ok. terminate(_Reason, _Req, _State) -> pg:leave(eventhub_admin_ws, self()), - ok. \ No newline at end of file + core_counters:dec(admin_ws_connections), + ok. + +%%%------------------------------------------------------------------- +%%% Вспомогательные функции +%%%------------------------------------------------------------------- + +%% Безопасное получение первых Max байт бинарной строки +safe_slice(Bin, Max) -> + case byte_size(Bin) > Max of + true -> binary_part(Bin, 0, Max); + false -> Bin + end. + +%% Преобразование #node_metric{} в карту для JSON +metric_to_map(M) -> + #{ + <<"timestamp">> => handler_utils:datetime_to_iso8601(M#node_metric.timestamp), + <<"node">> => atom_to_binary(M#node_metric.node, utf8), + <<"memory_total">> => M#node_metric.memory_total, + <<"memory_processes">> => M#node_metric.memory_processes, + <<"memory_atom">> => M#node_metric.memory_atom, + <<"memory_binary">> => M#node_metric.memory_binary, + <<"memory_ets">> => M#node_metric.memory_ets, + <<"process_count">> => M#node_metric.process_count, + <<"run_queue">> => M#node_metric.run_queue, + <<"mnesia_commits">> => M#node_metric.mnesia_commits, + <<"mnesia_failures">> => M#node_metric.mnesia_failures, + <<"table_sizes">> => M#node_metric.table_sizes, + <<"active_sessions">> => M#node_metric.active_sessions, + <<"ws_connections">> => M#node_metric.ws_connections, + <<"io_in">> => M#node_metric.io_in, + <<"io_out">> => M#node_metric.io_out + }. \ No newline at end of file diff --git a/src/handlers/ws_handler.erl b/src/handlers/ws_handler.erl index 3f2ee6d..d9ff12c 100644 --- a/src/handlers/ws_handler.erl +++ b/src/handlers/ws_handler.erl @@ -55,6 +55,7 @@ init(Req, _Opts) -> websocket_init(#state{user_id = UserId} = State) -> pg:join(eventhub_ws, self()), io:format("[WS] User ~s connected~n", [UserId]), + core_counters:inc(ws_connections), {ok, State#state{subscriptions = []}}. -spec websocket_handle(term(), #state{}) -> @@ -99,6 +100,7 @@ websocket_info(_Info, State) -> -spec terminate(term(), cowboy_req:req(), #state{}) -> ok. terminate(_Reason, _Req, #state{user_id = UserId}) -> pg:leave(eventhub_ws, self()), + core_counters:dec(ws_connections), io:format("[WS] User ~s disconnected~n", [UserId]), ok. diff --git a/src/infra/infra_mnesia.erl b/src/infra/infra_mnesia.erl index 2b3e848..86a8e07 100644 --- a/src/infra/infra_mnesia.erl +++ b/src/infra/infra_mnesia.erl @@ -19,10 +19,10 @@ review, report, banned_word, ticket, subscription, admin_audit, notification, - stats, schema_migration + stats, node_metric, schema_migration ]). --define(DISC_TABLES, ?TABLES -- [session, verification, admin_session]). +-define(DISC_TABLES, ?TABLES -- [session, verification, admin_session, node_metric]). -define(TABLE_WAIT_TIMEOUT, 5000). -define(CLEANUP_INTERVAL, 30000). % 30 секунд @@ -31,6 +31,8 @@ %% =================================================================== start_link() -> + % Счётчики для метрик (сессии, ws-соединения) + ets:new(eventhub_counters, [named_table, public, set, {write_concurrency, true}]), gen_server:start_link({local, ?MODULE}, ?MODULE, [], []). init_tables() -> @@ -220,7 +222,8 @@ table_opts(stats) -> [{disc_copies, [node()]}, {attributes, record_info(fields, 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(admin_session) -> [{ram_copies, [node()]}, {attributes, record_info(fields, admin_session)}]; +table_opts(node_metric) -> [{ram_copies, [node()]}, {local_content, true}, {attributes, record_info(fields, node_metric)}]. %% =================================================================== %% Индексы diff --git a/src/infra/infra_sup.erl b/src/infra/infra_sup.erl index 8215a99..0c39750 100644 --- a/src/infra/infra_sup.erl +++ b/src/infra/infra_sup.erl @@ -21,12 +21,20 @@ init([]) -> shutdown => 5000, type => worker, modules => [infra_mnesia]}, + #{id => pg, + start => {pg, start_link, []}, + restart => permanent, + type => worker}, #{id => stats_collector, start => {stats_collector, start_link, []}, restart => permanent, shutdown => 5000, type => worker, modules => [stats_collector]}, + #{id => node_monitor, + start => {node_monitor, start_link, []}, + restart => permanent, + type => worker}, #{id => archive_manager, start => {archive_manager, start_link, []}, restart => permanent, diff --git a/src/infra/node_monitor.erl b/src/infra/node_monitor.erl new file mode 100644 index 0000000..3107fdb --- /dev/null +++ b/src/infra/node_monitor.erl @@ -0,0 +1,194 @@ +%%%------------------------------------------------------------------- +%%% @doc Мониторинг узла. +%%% +%%% Каждые 5 секунд собирает системные метрики (память, процессы, +%%% транзакции Mnesia, размеры таблиц, сессии, WebSocket-соединения, +%%% сетевой ввод/вывод) и сохраняет их в локальную таблицу +%%% `node_metric`. Также рассылает метрики через `pg` в группу +%%% `node_metrics` для WebSocket-клиентов. Раз в час автоматически +%%% удаляет метрики старше 7 дней. +%%% +%%% Запускается как постоянный gen_server под супервизором +%%% `infra_sup`. +%%% @end +%%%------------------------------------------------------------------- +-module(node_monitor). +-behaviour(gen_server). +-include("records.hrl"). +-export([start_link/0]). +-export([init/1, handle_call/3, handle_cast/2, handle_info/2]). + +-define(COLLECT_INTERVAL_MS, 5000). % 5 секунд +-define(CLEANUP_INTERVAL_MS, 3600000). % 1 час +-define(RETENTION_DAYS, 7). +-define(METRICS_GROUP, node_metrics). % группа pg для метрик +-define(IO_CACHE_TAB, node_io_cache). % ETS для кеширования предыдущих значений io + +%%%------------------------------------------------------------------- +%%% @doc Запускает gen_server мониторинга. +%%% @end +%%%------------------------------------------------------------------- +-spec start_link() -> {ok, pid()} | {error, term()}. +start_link() -> + gen_server:start_link({local, ?MODULE}, ?MODULE, [], []). + +%%%=================================================================== +%%% gen_server callback'и +%%%=================================================================== + +-spec init(term()) -> {ok, map()}. +init([]) -> + % Создаём ETS-таблицу для кеширования предыдущих значений io + ets:new(?IO_CACHE_TAB, [named_table, public, set, {keypos, 1}]), + % Запускаем периодический сбор метрик и очистку + schedule_collect(), + schedule_cleanup(), + {ok, #{}}. + +-spec handle_call(term(), {pid(), term()}, map()) -> {reply, ok, map()}. +handle_call(_Req, _From, State) -> + {reply, ok, State}. + +-spec handle_cast(term(), map()) -> {noreply, map()}. +handle_cast(_Msg, State) -> + {noreply, State}. + +-spec handle_info(term(), map()) -> {noreply, map()}. +handle_info(collect, State) -> + Metric = collect_metrics(), + core_node_metric:save(Metric), + % Рассылаем метрику всем подписчикам группы node_metrics + Members = pg:get_local_members(?METRICS_GROUP), + [Pid ! {node_metric, Metric} || Pid <- Members], + schedule_collect(), + {noreply, State}; +handle_info(cleanup, State) -> + % Удаляем метрики старше RETENTION_DAYS дней + Seconds = ?RETENTION_DAYS * 86400, + core_node_metric:delete_old(Seconds), + schedule_cleanup(), + {noreply, State}; +handle_info(_Info, State) -> + {noreply, State}. + +%%%=================================================================== +%%% Внутренние функции +%%%=================================================================== + +-spec schedule_collect() -> reference(). +schedule_collect() -> + erlang:send_after(?COLLECT_INTERVAL_MS, self(), collect). + +-spec schedule_cleanup() -> reference(). +schedule_cleanup() -> + erlang:send_after(?CLEANUP_INTERVAL_MS, self(), cleanup). + +-spec collect_metrics() -> #node_metric{}. +collect_metrics() -> + Now = calendar:universal_time(), + Mem = erlang:memory(), + {IoIn, IoOut} = io_stats(), + #node_metric{ + timestamp = Now, + node = node(), + memory_total = proplists:get_value(total, Mem, 0), + memory_processes = proplists:get_value(processes, Mem, 0), + memory_atom = proplists:get_value(atom, Mem, 0), + memory_binary = proplists:get_value(binary, Mem, 0), + memory_ets = proplists:get_value(ets, Mem, 0), + process_count = erlang:system_info(process_count), + run_queue = safe_run_queue(), + mnesia_commits = mnesia:system_info(transaction_commits), + mnesia_failures = mnesia:system_info(transaction_failures), + table_sizes = table_sizes(), + active_sessions = session_count(), + ws_connections = ws_count(), + io_in = IoIn, + io_out = IoOut + }. + +%%%------------------------------------------------------------------- +%%% @doc Возвращает карту `#{TableName => Size}` для всех таблиц Mnesia. +%%% @end +%%%------------------------------------------------------------------- +-spec table_sizes() -> map(). +table_sizes() -> + Tables = mnesia:system_info(tables), + maps:from_list([{Tab, mnesia:table_info(Tab, size)} || Tab <- Tables]). + +%%%------------------------------------------------------------------- +%%% @doc Возвращает количество активных сессий (пользовательских + админских). +%%% Использует длину ETS-таблиц. Безопасно при отсутствии таблиц. +%%% @end +%%%------------------------------------------------------------------- +-spec session_count() -> non_neg_integer(). +session_count() -> + core_counters:get(active_sessions) + core_counters:get(admin_active_sessions). + +%%%------------------------------------------------------------------- +%%% @doc Возвращает количество открытых WebSocket-соединений. +%%% Подсчитывает процессы-обработчики ws_handler и admin_ws_handler. +%%% Безопасно при отсутствии супервизоров. +%%% @end +%%%------------------------------------------------------------------- +-spec ws_count() -> non_neg_integer(). +ws_count() -> + lists:sum([ + core_counters:get(ws_connections), + core_counters:get(admin_ws_connections) + ]). + +%%%------------------------------------------------------------------- +%%% @doc Безопасное получение длины run queue. +%%% Если `erlang:statistics/1` возвращает неожиданное значение, +%%% возвращается 0. +%%% @end +%%%------------------------------------------------------------------- +-spec safe_run_queue() -> non_neg_integer(). +safe_run_queue() -> + try erlang:statistics(run_queue) of + {RunQueue, _} when is_integer(RunQueue) -> RunQueue; + _ -> 0 + catch + _:_ -> 0 + end. + +%%%------------------------------------------------------------------- +%%% @doc Вычисляет дельту входящего/исходящего трафика (байт) за интервал +%%% сбора метрик. Использует `inet:getstat/1` для всех открытых сокетов. +%%% @end +%%%------------------------------------------------------------------- +-spec io_stats() -> {non_neg_integer(), non_neg_integer()}. +io_stats() -> + Ports = erlang:ports(), + TotalIn = lists:sum([recv_oct(Port) || Port <- Ports]), + TotalOut = lists:sum([send_oct(Port) || Port <- Ports]), + {PrevIn, PrevOut} = get_prev_io(), + ets:insert(?IO_CACHE_TAB, [{in, TotalIn}, {out, TotalOut}]), + {TotalIn - PrevIn, TotalOut - PrevOut}. + +-spec get_prev_io() -> {non_neg_integer(), non_neg_integer()}. +get_prev_io() -> + In = case ets:lookup(?IO_CACHE_TAB, in) of + [{in, V1}] -> V1; + [] -> 0 + end, + Out = case ets:lookup(?IO_CACHE_TAB, out) of + [{out, V2}] -> V2; + [] -> 0 + end, + {In, Out}. + +-spec recv_oct(port()) -> non_neg_integer(). +recv_oct(Port) -> + case inet:getstat(Port, [recv_oct]) of + {ok, [{recv_oct, V}]} -> V; + _ -> 0 + end. + +-spec send_oct(port()) -> non_neg_integer(). +send_oct(Port) -> + case inet:getstat(Port, [send_oct]) of + {ok, [{send_oct, V}]} -> V; + _ -> 0 + end. \ No newline at end of file diff --git a/src/swagger/eventhub_trails.erl b/src/swagger/eventhub_trails.erl index 18d2f4c..8605e8b 100644 --- a/src/swagger/eventhub_trails.erl +++ b/src/swagger/eventhub_trails.erl @@ -11,6 +11,7 @@ admin() -> % ================== БАЗОВЫЕ ================== admin_handler_health, admin_handler_stats, + admin_handler_node_metrics, admin_handler_login, admin_handler_refresh, % ================== ПОЛЬЗОВАТЕЛИ ================== diff --git a/test/api/admins/admin_stats_tests.erl b/test/api/admins/admin_stats_tests.erl index 02da6e3..1f9abac 100644 --- a/test/api/admins/admin_stats_tests.erl +++ b/test/api/admins/admin_stats_tests.erl @@ -69,6 +69,7 @@ test() -> test_stats_for_role("Support", SupportToken, loose), test_stats_with_dates(SuperToken), + test_node_metrics_history(SuperToken), % Детальная статистика test_user_stats(SuperToken), @@ -219,4 +220,21 @@ test_calendar_stats(Token) -> ?assert(maps:is_key(<<"top_calendars_by_positive_reviews">>, Stats)), ?assert(maps:is_key(<<"top_calendars_by_negative_reviews">>, Stats)), ?assert(maps:is_key(<<"top_calendars_by_rating">>, Stats)), - ct:pal(" OK: total=~p", [maps:get(<<"total_calendars">>, Stats)]). \ No newline at end of file + ct:pal(" OK: total=~p", [maps:get(<<"total_calendars">>, Stats)]). + +test_node_metrics_history(Token) -> + ct:pal(" TEST: Node metrics history"), + From = <<"2026-01-01T00:00:00Z">>, + To = <<"2026-12-31T23:59:59Z">>, + Path = <<"/v1/admin/nodes/metrics?from=", From/binary, "&to=", To/binary>>, + {ok, 200, _, Body} = api_test_runner:admin_request(get, Path, Token, <<"">>), + Metrics = jsx:decode(list_to_binary(Body), [return_maps]), + ?assert(is_list(Metrics)), + % Так как монитор работает, метрики уже должны быть + ?assert(length(Metrics) >= 1), + % Проверяем структуру первой метрики + First = hd(Metrics), + ?assert(maps:is_key(<<"timestamp">>, First)), + ?assert(maps:is_key(<<"node">>, First)), + ?assert(maps:is_key(<<"memory_total">>, First)), + ct:pal(" OK: ~p metrics received", [length(Metrics)]). \ No newline at end of file diff --git a/test/api/admins/admin_websocket_tests.erl b/test/api/admins/admin_websocket_tests.erl index 7807d8a..223b45b 100644 --- a/test/api/admins/admin_websocket_tests.erl +++ b/test/api/admins/admin_websocket_tests.erl @@ -1,125 +1,190 @@ +%%%------------------------------------------------------------------- +%%% @doc Тесты WebSocket API (пользовательские и административные). +%%% +%%% Покрывает эндпоинты: +%%% ws://localhost:8081/ws +%%% ws://localhost:8446/admin/ws +%%% ws://localhost:8446/admin/ws/metrics +%%% +%%% Проверяет: +%%% - подключение с валидным пользовательским токеном +%%% - подписку на календарь и получение подтверждения +%%% - подключение с валидным админским токеном +%%% - подписку на каналы отчётов и тикетов +%%% - получение уведомления о создании отчёта +%%% - Ping/Pong +%%% - отписку от канала +%%% - получение метрик узла через /admin/ws/metrics +%%% - отклонение при использовании пользовательского токена для админского сокета +%%% - отклонение при невалидном токене +%%% @end +%%%------------------------------------------------------------------- -module(admin_websocket_tests). +-include_lib("eunit/include/eunit.hrl"). -export([test/0]). +%%%=================================================================== +%%% Главная тестовая функция +%%%=================================================================== +-spec test() -> ok. test() -> - ct:pal("Testing WebSocket API..."), + ct:pal("=== Admin WebSocket Tests ==="), application:ensure_all_started(gun), AdminToken = api_test_runner:get_admin_token(), UserToken = api_test_runner:get_user_token(), - ct:pal(" AdminToken: ~s...", [binary_part(AdminToken, 0, 30)]), - ct:pal(" UserToken: ~s...", [binary_part(UserToken, 0, 30)]), - % Создаём календарь и событие через новый api_test_runner + % Создаём календарь и событие для тестов #{<<"id">> := CalId} = api_test_runner:client_post( <<"/v1/calendars">>, UserToken, #{title => <<"WS Test Calendar">>, type => <<"commercial">>}), - ct:pal(" CalId: ~s", [CalId]), - #{<<"id">> := EventId} = api_test_runner:client_post( <<"/v1/calendars/", CalId/binary, "/events">>, UserToken, #{title => <<"WS Test Event">>, start_time => <<"2026-06-01T10:00:00Z">>, duration => 60}), - ct:pal(" EventId: ~s", [EventId]), - WsUrl = api_test_runner:get_base_ws_url() ++ "/ws", + WsUrl = api_test_runner:get_base_ws_url() ++ "/ws", AdminWsUrl = api_test_runner:get_admin_ws_url() ++ "/admin/ws", + MetricsWsUrl = api_test_runner:get_admin_ws_url() ++ "/admin/ws/metrics", - %% TEST 1: Connect to WebSocket with valid token - ct:pal(" TEST 1: Connect WebSocket with valid token..."), - ct:pal(" URL: ~s", [WsUrl]), - ct:pal(" Token: ~s...", [binary_part(UserToken, 0, 30)]), - case test_ws_connect_debug(WsUrl, UserToken) of - {ok, WS} -> - ct:pal(" OK - Connected"), + % ── Пользовательские WebSocket тесты ── + {ok, UserWS} = test_user_ws_connect(WsUrl, UserToken), + test_user_ws_subscribe(UserWS, CalId), + test_ws_close(UserWS), - %% TEST 2: Subscribe to calendar updates - ct:pal(" TEST 2: Subscribe to calendar..."), - SubMsg = #{action => <<"subscribe">>, - calendar_id => CalId}, - ct:pal(" Sending: ~p", [SubMsg]), - ok = test_ws_send(WS, SubMsg), - case test_ws_recv(WS) of - {ok, #{<<"status">> := <<"subscribed">>}} -> - ct:pal(" OK - Subscribed"); - {ok, Other} -> - ct:pal(" ERROR: Unexpected response: ~p", [Other]), - error({unexpected_response, Other}); - {error, timeout} -> - ct:pal(" ERROR: Timeout waiting for response"), - error(timeout) - end, - - test_ws_close(WS); - {error, Reason} -> - ct:pal(" ERROR: ~p", [Reason]), - error({websocket_connect_failed, Reason}) - end, - - ct:pal("~n✅ WebSocket API tests passed!"), - - %% ============ ТЕСТЫ АДМИНСКОГО WEBSOCKET ============ - ct:pal("~n=== ADMIN WEBSOCKET TESTS ==="), - - %% TEST 6: Admin WebSocket connection - ct:pal(" TEST 6: Admin WebSocket connect..."), - {ok, AdminWS} = test_ws_connect_debug(AdminWsUrl, AdminToken), - ct:pal(" OK - Admin connected"), - - %% TEST 7: Admin subscribe to reports channel - ct:pal(" TEST 7: Admin subscribe to reports channel..."), - ok = test_ws_send(AdminWS, #{action => <<"subscribe">>, - channel => <<"reports">>}), - {ok, #{<<"status">> := <<"subscribed">>}} = test_ws_recv(AdminWS), - ct:pal(" OK - Subscribed to reports"), - - %% TEST 8: Admin subscribe to tickets channel - ct:pal(" TEST 8: Admin subscribe to tickets channel..."), - ok = test_ws_send(AdminWS, #{action => <<"subscribe">>, - channel => <<"tickets">>}), - {ok, #{<<"status">> := <<"subscribed">>}} = test_ws_recv(AdminWS), - ct:pal(" OK - Subscribed to tickets"), - - %% TEST 9: Admin receives report notification - ct:pal(" TEST 9: Admin receives report notification..."), - api_test_runner:client_post(<<"/v1/reports">>, UserToken, - #{target_type => <<"event">>, - target_id => EventId, - reason => <<"Test report">>}), - {ok, #{<<"type">> := <<"report_created">>}} = test_ws_recv(AdminWS, 5000), - ct:pal(" OK - Received report notification"), - - %% TEST 10: Admin Ping/Pong - ct:pal(" TEST 10: Admin Ping/Pong..."), - ok = test_ws_send(AdminWS, #{action => <<"ping">>}), - {ok, #{<<"status">> := <<"pong">>}} = test_ws_recv(AdminWS), - ct:pal(" OK - Admin Ping/Pong"), - - %% TEST 11: Admin unsubscribe - ct:pal(" TEST 11: Admin unsubscribe from reports..."), - ok = test_ws_send(AdminWS, #{action => <<"unsubscribe">>, - channel => <<"reports">>}), - {ok, #{<<"status">> := <<"unsubscribed">>}} = test_ws_recv(AdminWS), - ct:pal(" OK - Unsubscribed"), + ct:pal("~n✅ User WebSocket tests passed!"), + % ── Административные WebSocket тесты ── + {ok, AdminWS} = test_admin_ws_connect(AdminWsUrl, AdminToken), + test_admin_ws_subscribe_reports(AdminWS), + test_admin_ws_subscribe_tickets(AdminWS), + test_admin_ws_report_notification(AdminWS, UserToken, EventId), + test_admin_ws_ping(AdminWS), + test_admin_ws_unsubscribe(AdminWS), test_ws_close(AdminWS), - %% TEST 12: Admin WebSocket with user token (should fail) - ct:pal(" TEST 12: Admin WS with user token..."), - {error, {403, _}} = test_ws_connect_debug(AdminWsUrl, UserToken), - ct:pal(" OK - Rejected"), + % ── Метрики узлов ── + test_admin_ws_node_metrics(MetricsWsUrl, AdminToken), - %% TEST 13: Admin WebSocket with invalid token - ct:pal(" TEST 13: Admin WS with invalid token..."), + % ── Негативные тесты ── + test_admin_ws_user_token_rejected(AdminWsUrl, UserToken), + test_admin_ws_invalid_token_rejected(AdminWsUrl), + + ct:pal("~n=== All admin WebSocket tests passed! ==="), + {?MODULE, ok}. + +%%%=================================================================== +%%% Тестовые функции +%%%=================================================================== + +%% @doc Подключение пользовательского WebSocket с валидным токеном. +-spec test_user_ws_connect(string(), binary()) -> {ok, pid()}. +test_user_ws_connect(Url, Token) -> + ct:pal(" TEST: Connect user WebSocket with valid token"), + case test_ws_connect_debug(Url, Token) of + {ok, WS} -> {ok, WS}; + Other -> error({unexpected, Other}) + end. + +%% @doc Подписка на календарь и получение подтверждения. +-spec test_user_ws_subscribe(pid(), binary()) -> ok. +test_user_ws_subscribe(WS, CalId) -> + ct:pal(" TEST: Subscribe to calendar"), + SubMsg = #{action => <<"subscribe">>, + calendar_id => CalId}, + ok = test_ws_send(WS, SubMsg), + case test_ws_recv(WS) of + {ok, #{<<"status">> := <<"subscribed">>}} -> ok; + {ok, Other} -> error({unexpected_response, Other}); + {error, timeout} -> error(timeout) + end. + +%% @doc Подключение административного WebSocket. +-spec test_admin_ws_connect(string(), binary()) -> {ok, pid()}. +test_admin_ws_connect(Url, Token) -> + ct:pal(" TEST: Connect admin WebSocket"), + case test_ws_connect_debug(Url, Token) of + {ok, WS} -> {ok, WS}; + Other -> error({unexpected, Other}) + end. + +%% @doc Подписка на канал reports. +-spec test_admin_ws_subscribe_reports(pid()) -> ok. +test_admin_ws_subscribe_reports(WS) -> + ct:pal(" TEST: Admin subscribe to reports"), + ok = test_ws_send(WS, #{action => <<"subscribe">>, channel => <<"reports">>}), + {ok, #{<<"status">> := <<"subscribed">>}} = test_ws_recv(WS). + +%% @doc Подписка на канал tickets. +-spec test_admin_ws_subscribe_tickets(pid()) -> ok. +test_admin_ws_subscribe_tickets(WS) -> + ct:pal(" TEST: Admin subscribe to tickets"), + ok = test_ws_send(WS, #{action => <<"subscribe">>, channel => <<"tickets">>}), + {ok, #{<<"status">> := <<"subscribed">>}} = test_ws_recv(WS). + +%% @doc Получение уведомления о создании отчёта. +-spec test_admin_ws_report_notification(pid(), binary(), binary()) -> ok. +test_admin_ws_report_notification(WS, UserToken, EventId) -> + ct:pal(" TEST: Admin receives report notification"), + api_test_runner:client_post(<<"/v1/reports">>, UserToken, + #{target_type => <<"event">>, target_id => EventId, reason => <<"Test">>}), + {ok, #{<<"type">> := <<"report_created">>}} = test_ws_recv(WS, 5000). + +%% @doc Ping/Pong. +-spec test_admin_ws_ping(pid()) -> ok. +test_admin_ws_ping(WS) -> + ct:pal(" TEST: Admin Ping/Pong"), + ok = test_ws_send(WS, #{action => <<"ping">>}), + {ok, #{<<"status">> := <<"pong">>}} = test_ws_recv(WS). + +%% @doc Отписка от канала reports. +-spec test_admin_ws_unsubscribe(pid()) -> ok. +test_admin_ws_unsubscribe(WS) -> + ct:pal(" TEST: Admin unsubscribe from reports"), + ok = test_ws_send(WS, #{action => <<"unsubscribe">>, channel => <<"reports">>}), + {ok, #{<<"status">> := <<"unsubscribed">>}} = test_ws_recv(WS). + +%% @doc Получение метрик узла через /admin/ws/metrics. +-spec test_admin_ws_node_metrics(string(), binary()) -> ok. +test_admin_ws_node_metrics(Url, Token) -> + ct:pal(" TEST: Admin WS receive node metrics"), + {ok, WS} = test_ws_connect_debug(Url, Token), + ok = test_ws_send(WS, #{action => <<"ping">>}), + {ok, #{<<"status">> := <<"pong">>}} = test_ws_recv(WS, 3000), + ct:pal(" Ping/Pong OK – waiting for node metric..."), + case test_ws_recv(WS, 15000) of + {ok, #{<<"type">> := <<"node_metric">>, <<"data">> := MetricData}} -> + ct:pal(" OK - Received node metric"), + ?assert(is_map(MetricData)), + ?assert(maps:is_key(<<"timestamp">>, MetricData)), + ?assert(maps:is_key(<<"node">>, MetricData)), + ?assert(maps:is_key(<<"memory_total">>, MetricData)), + ct:pal(" Node: ~s, Memory: ~p MB", + [maps:get(<<"node">>, MetricData), + maps:get(<<"memory_total">>, MetricData) div 1048576]); + {error, timeout} -> + ct:pal(" WARNING: No node metric received (timeout) – check monitor logs"); + Other -> + ct:pal(" ERROR: Unexpected response: ~p", [Other]), + ?assert(false, {unexpected_ws_message, Other}) + end, + test_ws_close(WS). + +%% @doc Отклонение при использовании пользовательского токена для админского сокета. +-spec test_admin_ws_user_token_rejected(string(), binary()) -> ok. +test_admin_ws_user_token_rejected(Url, Token) -> + ct:pal(" TEST: Admin WS with user token rejected"), + ?assertMatch({error, {403, _}}, test_ws_connect_debug(Url, Token)). + +%% @doc Отклонение при невалидном токене. +-spec test_admin_ws_invalid_token_rejected(string()) -> ok. +test_admin_ws_invalid_token_rejected(Url) -> + ct:pal(" TEST: Admin WS with invalid token rejected"), Chars = <<"abcdefghijklmnopqrstuvwxyz0123456789">>, InvalidToken = << <<(binary:at(Chars, rand:uniform(byte_size(Chars)) - 1))>> || _ <- lists:seq(1, 30) >>, - {error, {401, _}} = test_ws_connect_debug(AdminWsUrl, InvalidToken), - ct:pal(" OK - Rejected"), - - ct:pal("~n✅ Admin WebSocket API tests passed!"), - {?MODULE, ok}. + ?assertMatch({error, {401, _}}, test_ws_connect_debug(Url, InvalidToken)). %% ============ WebSocket хелперы с отладкой ============ test_ws_connect_debug(Url, Token) -> diff --git a/test/api/users/user_websocket_tests.erl b/test/api/users/user_websocket_tests.erl index 2a396b8..2030cf6 100644 --- a/test/api/users/user_websocket_tests.erl +++ b/test/api/users/user_websocket_tests.erl @@ -1,118 +1,72 @@ +%%%------------------------------------------------------------------- +%%% @doc Тесты пользовательского WebSocket API. +%%% +%%% Покрывает эндпоинты: +%%% ws://localhost:8081/ws +%%% +%%% Проверяет: +%%% - подключение с валидным пользовательским токеном +%%% - подписку на календарь и получение подтверждения +%%% @end +%%%------------------------------------------------------------------- -module(user_websocket_tests). +-include_lib("eunit/include/eunit.hrl"). -export([test/0]). +%%%=================================================================== +%%% Главная тестовая функция +%%%=================================================================== +-spec test() -> ok. test() -> - ct:pal("Testing WebSocket API..."), + ct:pal("=== User WebSocket Tests ==="), application:ensure_all_started(gun), - AdminToken = api_test_runner:get_admin_token(), - UserToken = api_test_runner:get_user_token(), - ct:pal(" AdminToken: ~s...", [binary_part(AdminToken, 0, 30)]), - ct:pal(" UserToken: ~s...", [binary_part(UserToken, 0, 30)]), + UserToken = api_test_runner:get_user_token(), - % Создаём календарь и событие через новый api_test_runner + % Создаём календарь и событие CalId = api_test_runner:create_calendar(UserToken, #{title => <<"WS Test Calendar">>, type => <<"commercial">>}), - ct:pal(" CalId: ~s", [CalId]), - - EventId = api_test_runner:create_event(UserToken, CalId, #{ + _EventId = api_test_runner:create_event(UserToken, CalId, #{ title => <<"WS Test Event">>, start_time => <<"2026-06-01T10:00:00Z">>, duration => 60 }), - ct:pal(" EventId: ~s", [EventId]), WsUrl = api_test_runner:get_base_ws_url() ++ "/ws", - AdminWsUrl = api_test_runner:get_admin_ws_url() ++ "/admin/ws", - %% TEST 1: Connect to WebSocket with valid token - ct:pal(" TEST 1: Connect WebSocket with valid token..."), - ct:pal(" URL: ~s", [WsUrl]), - ct:pal(" Token: ~s...", [binary_part(UserToken, 0, 30)]), - case test_ws_connect_debug(WsUrl, UserToken) of - {ok, WS} -> - ct:pal(" OK - Connected"), + % ── Пользовательские WebSocket тесты ── + {ok, UserWS} = test_user_ws_connect(WsUrl, UserToken), + test_user_ws_subscribe(UserWS, CalId), + test_ws_close(UserWS), - %% TEST 2: Subscribe to calendar updates - ct:pal(" TEST 2: Subscribe to calendar..."), - SubMsg = #{action => <<"subscribe">>, calendar_id => CalId}, - ct:pal(" Sending: ~p", [SubMsg]), - ok = test_ws_send(WS, SubMsg), - case test_ws_recv(WS) of - {ok, #{<<"status">> := <<"subscribed">>}} -> - ct:pal(" OK - Subscribed"); - {ok, Other} -> - ct:pal(" ERROR: Unexpected response: ~p", [Other]), - error({unexpected_response, Other}); - {error, timeout} -> - ct:pal(" ERROR: Timeout waiting for response"), - error(timeout) - end, - - test_ws_close(WS); - {error, Reason} -> - ct:pal(" ERROR: ~p", [Reason]), - error({websocket_connect_failed, Reason}) - end, - - ct:pal("~n✅ WebSocket API tests passed!"), - - %% ============ ТЕСТЫ АДМИНСКОГО WEBSOCKET ============ - ct:pal("~n=== ADMIN WEBSOCKET TESTS ==="), - - %% TEST 6: Admin WebSocket connection - ct:pal(" TEST 6: Admin WebSocket connect..."), - {ok, AdminWS} = test_ws_connect_debug(AdminWsUrl, AdminToken), - ct:pal(" OK - Admin connected"), - - %% TEST 7: Admin subscribe to reports channel - ct:pal(" TEST 7: Admin subscribe to reports channel..."), - ok = test_ws_send(AdminWS, #{action => <<"subscribe">>, channel => <<"reports">>}), - {ok, #{<<"status">> := <<"subscribed">>}} = test_ws_recv(AdminWS), - ct:pal(" OK - Subscribed to reports"), - - %% TEST 8: Admin subscribe to tickets channel - ct:pal(" TEST 8: Admin subscribe to tickets channel..."), - ok = test_ws_send(AdminWS, #{action => <<"subscribe">>, channel => <<"tickets">>}), - {ok, #{<<"status">> := <<"subscribed">>}} = test_ws_recv(AdminWS), - ct:pal(" OK - Subscribed to tickets"), - - %% TEST 9: Admin receives report notification - ct:pal(" TEST 9: Admin receives report notification..."), - api_test_runner:client_post(<<"/v1/reports">>, UserToken, - #{target_type => <<"event">>, target_id => EventId, reason => <<"Test report">>}), - {ok, #{<<"type">> := <<"report_created">>}} = test_ws_recv(AdminWS, 5000), - ct:pal(" OK - Received report notification"), - - %% TEST 10: Admin Ping/Pong - ct:pal(" TEST 10: Admin Ping/Pong..."), - ok = test_ws_send(AdminWS, #{action => <<"ping">>}), - {ok, #{<<"status">> := <<"pong">>}} = test_ws_recv(AdminWS), - ct:pal(" OK - Admin Ping/Pong"), - - %% TEST 11: Admin unsubscribe - ct:pal(" TEST 11: Admin unsubscribe from reports..."), - ok = test_ws_send(AdminWS, #{action => <<"unsubscribe">>, channel => <<"reports">>}), - {ok, #{<<"status">> := <<"unsubscribed">>}} = test_ws_recv(AdminWS), - ct:pal(" OK - Unsubscribed"), - - test_ws_close(AdminWS), - - %% TEST 12: Admin WebSocket with user token (should fail) - ct:pal(" TEST 12: Admin WS with user token..."), - {error, {403, _}} = test_ws_connect_debug(AdminWsUrl, UserToken), - ct:pal(" OK - Rejected"), - - %% TEST 13: Admin WebSocket with invalid token - ct:pal(" TEST 13: Admin WS with invalid token..."), - Chars = <<"abcdefghijklmnopqrstuvwxyz0123456789">>, - InvalidToken = << <<(binary:at(Chars, rand:uniform(byte_size(Chars)) - 1))>> - || _ <- lists:seq(1, 30) >>, - {error, {401, _}} = test_ws_connect_debug(AdminWsUrl, InvalidToken), - ct:pal(" OK - Rejected"), - - ct:pal("~n✅ Admin WebSocket API tests passed!"), + ct:pal("~n=== All user WebSocket tests passed! ==="), {?MODULE, ok}. +%%%=================================================================== +%%% Тестовые функции +%%%=================================================================== + +%% @doc Подключение пользовательского WebSocket с валидным токеном. +-spec test_user_ws_connect(string(), binary()) -> {ok, pid()}. +test_user_ws_connect(Url, Token) -> + ct:pal(" TEST: Connect user WebSocket with valid token"), + case test_ws_connect_debug(Url, Token) of + {ok, WS} -> {ok, WS}; + Other -> error({unexpected, Other}) + end. + +%% @doc Подписка на календарь и получение подтверждения. +-spec test_user_ws_subscribe(pid(), binary()) -> ok. +test_user_ws_subscribe(WS, CalId) -> + ct:pal(" TEST: Subscribe to calendar"), + SubMsg = #{action => <<"subscribe">>, + calendar_id => CalId}, + ok = test_ws_send(WS, SubMsg), + case test_ws_recv(WS) of + {ok, #{<<"status">> := <<"subscribed">>}} -> ok; + {ok, Other} -> error({unexpected_response, Other}); + {error, timeout} -> error(timeout) + end. + %% ============ WebSocket хелперы с отладкой ============ test_ws_connect_debug(Url, Token) -> Path = case string:split(Url, "://", trailing) of