Добавлена оперативная и историческая статистика нод и узлов mnesia #22
This commit is contained in:
@@ -252,6 +252,25 @@
|
|||||||
timestamp :: calendar:datetime()
|
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, {
|
-record(schema_migration, {
|
||||||
version :: string(),
|
version :: string(),
|
||||||
applied_at :: calendar:datetime()
|
applied_at :: calendar:datetime()
|
||||||
|
|||||||
@@ -26,6 +26,7 @@ create(AdminId, RefreshToken) ->
|
|||||||
type = refresh
|
type = refresh
|
||||||
},
|
},
|
||||||
mnesia:dirty_write(Session),
|
mnesia:dirty_write(Session),
|
||||||
|
core_counters:inc(admin_active_sessions),
|
||||||
ok.
|
ok.
|
||||||
|
|
||||||
%%%-------------------------------------------------------------------
|
%%%-------------------------------------------------------------------
|
||||||
@@ -44,6 +45,7 @@ validate(Token) ->
|
|||||||
true -> {ok, Session#admin_session.admin_id};
|
true -> {ok, Session#admin_session.admin_id};
|
||||||
false ->
|
false ->
|
||||||
mnesia:dirty_delete({admin_session, Token}),
|
mnesia:dirty_delete({admin_session, Token}),
|
||||||
|
core_counters:dec(admin_active_sessions),
|
||||||
{error, expired}
|
{error, expired}
|
||||||
end;
|
end;
|
||||||
[] -> {error, not_found}
|
[] -> {error, not_found}
|
||||||
@@ -56,4 +58,5 @@ validate(Token) ->
|
|||||||
-spec delete(Token :: binary()) -> ok.
|
-spec delete(Token :: binary()) -> ok.
|
||||||
delete(Token) ->
|
delete(Token) ->
|
||||||
mnesia:dirty_delete({admin_session, Token}),
|
mnesia:dirty_delete({admin_session, Token}),
|
||||||
|
core_counters:dec(admin_active_sessions),
|
||||||
ok.
|
ok.
|
||||||
@@ -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.
|
||||||
@@ -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.
|
||||||
@@ -26,6 +26,7 @@ create(UserId, RefreshToken) ->
|
|||||||
type = refresh
|
type = refresh
|
||||||
},
|
},
|
||||||
mnesia:dirty_write(Session),
|
mnesia:dirty_write(Session),
|
||||||
|
core_counters:inc(active_sessions),
|
||||||
ok.
|
ok.
|
||||||
|
|
||||||
%%%-------------------------------------------------------------------
|
%%%-------------------------------------------------------------------
|
||||||
@@ -49,6 +50,7 @@ validate(Token) ->
|
|||||||
end;
|
end;
|
||||||
false ->
|
false ->
|
||||||
mnesia:dirty_delete({session, Token}),
|
mnesia:dirty_delete({session, Token}),
|
||||||
|
core_counters:dec(active_sessions),
|
||||||
{error, expired}
|
{error, expired}
|
||||||
end;
|
end;
|
||||||
[] -> {error, not_found}
|
[] -> {error, not_found}
|
||||||
@@ -61,4 +63,5 @@ validate(Token) ->
|
|||||||
-spec delete(Token :: binary()) -> ok.
|
-spec delete(Token :: binary()) -> ok.
|
||||||
delete(Token) ->
|
delete(Token) ->
|
||||||
mnesia:dirty_delete({session, Token}),
|
mnesia:dirty_delete({session, Token}),
|
||||||
|
core_counters:dec(active_sessions),
|
||||||
ok.
|
ok.
|
||||||
@@ -3,7 +3,6 @@
|
|||||||
-export([start/2, stop/1]).
|
-export([start/2, stop/1]).
|
||||||
|
|
||||||
start(_StartType, _StartArgs) ->
|
start(_StartType, _StartArgs) ->
|
||||||
pg:start_link(),
|
|
||||||
case infra_sup:start_link() of
|
case infra_sup:start_link() of
|
||||||
{ok, Pid} ->
|
{ok, Pid} ->
|
||||||
% Определяем список узлов кластера, если режим CLUSTER_MODE=true
|
% Определяем список узлов кластера, если режим CLUSTER_MODE=true
|
||||||
@@ -109,6 +108,7 @@ start_admin_http() ->
|
|||||||
% ================== БАЗОВЫЕ ==================
|
% ================== БАЗОВЫЕ ==================
|
||||||
{"/admin/health", admin_handler_health, []},
|
{"/admin/health", admin_handler_health, []},
|
||||||
{"/v1/admin/stats", admin_handler_stats, []},
|
{"/v1/admin/stats", admin_handler_stats, []},
|
||||||
|
{"/v1/admin/nodes/metrics", admin_handler_node_metrics, []},
|
||||||
{"/v1/admin/login", admin_handler_login, []},
|
{"/v1/admin/login", admin_handler_login, []},
|
||||||
{"/v1/admin/refresh", admin_handler_refresh, []},
|
{"/v1/admin/refresh", admin_handler_refresh, []},
|
||||||
% ================== ПОЛЬЗОВАТЕЛИ ==================
|
% ================== ПОЛЬЗОВАТЕЛИ ==================
|
||||||
@@ -169,7 +169,10 @@ start_admin_http() ->
|
|||||||
cowboy:start_clear(ws, [{port, PortWs}], #{env => #{dispatch => WsDispatch}}),
|
cowboy:start_clear(ws, [{port, PortWs}], #{env => #{dispatch => WsDispatch}}),
|
||||||
|
|
||||||
% WebSocket для админов
|
% 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),
|
PortAdminWs = get_env_int(admin_ws_port, 8446),
|
||||||
cowboy:start_clear(admin_ws, [{port, PortAdminWs}], #{env => #{dispatch => AdminWsDispatch}}),
|
cowboy:start_clear(admin_ws, [{port, PortAdminWs}], #{env => #{dispatch => AdminWsDispatch}}),
|
||||||
|
|
||||||
|
|||||||
@@ -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
|
||||||
|
}.
|
||||||
@@ -1,11 +1,13 @@
|
|||||||
%%%-------------------------------------------------------------------
|
%%%-------------------------------------------------------------------
|
||||||
%%% @doc Административный WebSocket-обработчик.
|
%%% @doc Административный WebSocket-обработчик.
|
||||||
%%% Устанавливает WebSocket-соединение после проверки JWT-токена
|
%%% Устанавливает WebSocket-соединение после проверки JWT-токена
|
||||||
%%% и подписывает администратора на каналы уведомлений.
|
%%% и подписывает администратора на каналы уведомлений, включая
|
||||||
|
%%% канал метрик узлов `node_metrics`.
|
||||||
%%% @end
|
%%% @end
|
||||||
%%%-------------------------------------------------------------------
|
%%%-------------------------------------------------------------------
|
||||||
-module(admin_ws_handler).
|
-module(admin_ws_handler).
|
||||||
-behaviour(cowboy_websocket).
|
-behaviour(cowboy_websocket).
|
||||||
|
-include("records.hrl").
|
||||||
|
|
||||||
-export([init/2]).
|
-export([init/2]).
|
||||||
-export([websocket_init/1]).
|
-export([websocket_init/1]).
|
||||||
@@ -30,7 +32,7 @@ init(Req, _Opts) ->
|
|||||||
Resp = cowboy_req:reply(401, #{}, <<"Missing token">>, Req),
|
Resp = cowboy_req:reply(401, #{}, <<"Missing token">>, Req),
|
||||||
{ok, Resp, undefined};
|
{ok, Resp, undefined};
|
||||||
Token ->
|
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
|
case logic_auth:verify_jwt(Token) of
|
||||||
{ok, UserId, Role} ->
|
{ok, UserId, Role} ->
|
||||||
io:format("[ADMIN_WS] UserId: ~s, Role: ~s~n", [UserId, Role]),
|
io:format("[ADMIN_WS] UserId: ~s, Role: ~s~n", [UserId, Role]),
|
||||||
@@ -57,8 +59,11 @@ init(Req, _Opts) ->
|
|||||||
%% @doc Вызывается после установки WebSocket-соединения.
|
%% @doc Вызывается после установки WebSocket-соединения.
|
||||||
-spec websocket_init(#state{}) -> {ok, #state{}}.
|
-spec websocket_init(#state{}) -> {ok, #state{}}.
|
||||||
websocket_init(State) ->
|
websocket_init(State) ->
|
||||||
|
core_counters:inc(admin_ws_connections),
|
||||||
io:format("[ADMIN_WS] WebSocket initialized for admin ~s~n", [State#state.admin_id]),
|
io:format("[ADMIN_WS] WebSocket initialized for admin ~s~n", [State#state.admin_id]),
|
||||||
pg:join(eventhub_admin_ws, self()),
|
pg:join(eventhub_admin_ws, self()),
|
||||||
|
% Подписываемся на метрики узлов
|
||||||
|
pg:join(node_metrics, self()),
|
||||||
{ok, State}.
|
{ok, State}.
|
||||||
|
|
||||||
%% @doc Обрабатывает входящие текстовые сообщения (subscribe/unsubscribe/ping).
|
%% @doc Обрабатывает входящие текстовые сообщения (subscribe/unsubscribe/ping).
|
||||||
@@ -84,7 +89,7 @@ websocket_handle({text, Msg}, State) ->
|
|||||||
websocket_handle(_Frame, State) ->
|
websocket_handle(_Frame, State) ->
|
||||||
{ok, State}.
|
{ok, State}.
|
||||||
|
|
||||||
%% @doc Отправляет административное уведомление через WebSocket.
|
%% @doc Отправляет административное уведомление или метрику узла через WebSocket.
|
||||||
-spec websocket_info(term(), #state{}) -> {reply, {text, binary()}, #state{}} | {ok, #state{}}.
|
-spec websocket_info(term(), #state{}) -> {reply, {text, binary()}, #state{}} | {ok, #state{}}.
|
||||||
websocket_info({admin_notification, Type, Data}, State) ->
|
websocket_info({admin_notification, Type, Data}, State) ->
|
||||||
Msg = jsx:encode(#{
|
Msg = jsx:encode(#{
|
||||||
@@ -93,6 +98,12 @@ websocket_info({admin_notification, Type, Data}, State) ->
|
|||||||
timestamp => os:system_time(seconds)
|
timestamp => os:system_time(seconds)
|
||||||
}),
|
}),
|
||||||
{reply, {text, Msg}, State};
|
{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) ->
|
websocket_info(_Info, State) ->
|
||||||
{ok, State}.
|
{ok, State}.
|
||||||
|
|
||||||
@@ -100,4 +111,37 @@ websocket_info(_Info, State) ->
|
|||||||
-spec terminate(term(), cowboy_req:req(), #state{}) -> ok.
|
-spec terminate(term(), cowboy_req:req(), #state{}) -> ok.
|
||||||
terminate(_Reason, _Req, _State) ->
|
terminate(_Reason, _Req, _State) ->
|
||||||
pg:leave(eventhub_admin_ws, self()),
|
pg:leave(eventhub_admin_ws, self()),
|
||||||
|
core_counters:dec(admin_ws_connections),
|
||||||
ok.
|
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
|
||||||
|
}.
|
||||||
@@ -55,6 +55,7 @@ init(Req, _Opts) ->
|
|||||||
websocket_init(#state{user_id = UserId} = State) ->
|
websocket_init(#state{user_id = UserId} = State) ->
|
||||||
pg:join(eventhub_ws, self()),
|
pg:join(eventhub_ws, self()),
|
||||||
io:format("[WS] User ~s connected~n", [UserId]),
|
io:format("[WS] User ~s connected~n", [UserId]),
|
||||||
|
core_counters:inc(ws_connections),
|
||||||
{ok, State#state{subscriptions = []}}.
|
{ok, State#state{subscriptions = []}}.
|
||||||
|
|
||||||
-spec websocket_handle(term(), #state{}) ->
|
-spec websocket_handle(term(), #state{}) ->
|
||||||
@@ -99,6 +100,7 @@ websocket_info(_Info, State) ->
|
|||||||
-spec terminate(term(), cowboy_req:req(), #state{}) -> ok.
|
-spec terminate(term(), cowboy_req:req(), #state{}) -> ok.
|
||||||
terminate(_Reason, _Req, #state{user_id = UserId}) ->
|
terminate(_Reason, _Req, #state{user_id = UserId}) ->
|
||||||
pg:leave(eventhub_ws, self()),
|
pg:leave(eventhub_ws, self()),
|
||||||
|
core_counters:dec(ws_connections),
|
||||||
io:format("[WS] User ~s disconnected~n", [UserId]),
|
io:format("[WS] User ~s disconnected~n", [UserId]),
|
||||||
ok.
|
ok.
|
||||||
|
|
||||||
|
|||||||
@@ -19,10 +19,10 @@
|
|||||||
review, report, banned_word,
|
review, report, banned_word,
|
||||||
ticket, subscription,
|
ticket, subscription,
|
||||||
admin_audit, notification,
|
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(TABLE_WAIT_TIMEOUT, 5000).
|
||||||
-define(CLEANUP_INTERVAL, 30000). % 30 секунд
|
-define(CLEANUP_INTERVAL, 30000). % 30 секунд
|
||||||
|
|
||||||
@@ -31,6 +31,8 @@
|
|||||||
%% ===================================================================
|
%% ===================================================================
|
||||||
|
|
||||||
start_link() ->
|
start_link() ->
|
||||||
|
% Счётчики для метрик (сессии, ws-соединения)
|
||||||
|
ets:new(eventhub_counters, [named_table, public, set, {write_concurrency, true}]),
|
||||||
gen_server:start_link({local, ?MODULE}, ?MODULE, [], []).
|
gen_server:start_link({local, ?MODULE}, ?MODULE, [], []).
|
||||||
|
|
||||||
init_tables() ->
|
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(schema_migration) -> [{disc_copies, [node()]}, {attributes, record_info(fields, schema_migration)}];
|
||||||
table_opts(session) -> [{ram_copies, [node()]}, {attributes, record_info(fields, session)}];
|
table_opts(session) -> [{ram_copies, [node()]}, {attributes, record_info(fields, session)}];
|
||||||
table_opts(verification) -> [{ram_copies, [node()]}, {attributes, record_info(fields, verification)}];
|
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)}].
|
||||||
|
|
||||||
%% ===================================================================
|
%% ===================================================================
|
||||||
%% Индексы
|
%% Индексы
|
||||||
|
|||||||
@@ -21,12 +21,20 @@ init([]) ->
|
|||||||
shutdown => 5000,
|
shutdown => 5000,
|
||||||
type => worker,
|
type => worker,
|
||||||
modules => [infra_mnesia]},
|
modules => [infra_mnesia]},
|
||||||
|
#{id => pg,
|
||||||
|
start => {pg, start_link, []},
|
||||||
|
restart => permanent,
|
||||||
|
type => worker},
|
||||||
#{id => stats_collector,
|
#{id => stats_collector,
|
||||||
start => {stats_collector, start_link, []},
|
start => {stats_collector, start_link, []},
|
||||||
restart => permanent,
|
restart => permanent,
|
||||||
shutdown => 5000,
|
shutdown => 5000,
|
||||||
type => worker,
|
type => worker,
|
||||||
modules => [stats_collector]},
|
modules => [stats_collector]},
|
||||||
|
#{id => node_monitor,
|
||||||
|
start => {node_monitor, start_link, []},
|
||||||
|
restart => permanent,
|
||||||
|
type => worker},
|
||||||
#{id => archive_manager,
|
#{id => archive_manager,
|
||||||
start => {archive_manager, start_link, []},
|
start => {archive_manager, start_link, []},
|
||||||
restart => permanent,
|
restart => permanent,
|
||||||
|
|||||||
@@ -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.
|
||||||
@@ -11,6 +11,7 @@ admin() ->
|
|||||||
% ================== БАЗОВЫЕ ==================
|
% ================== БАЗОВЫЕ ==================
|
||||||
admin_handler_health,
|
admin_handler_health,
|
||||||
admin_handler_stats,
|
admin_handler_stats,
|
||||||
|
admin_handler_node_metrics,
|
||||||
admin_handler_login,
|
admin_handler_login,
|
||||||
admin_handler_refresh,
|
admin_handler_refresh,
|
||||||
% ================== ПОЛЬЗОВАТЕЛИ ==================
|
% ================== ПОЛЬЗОВАТЕЛИ ==================
|
||||||
|
|||||||
@@ -69,6 +69,7 @@ test() ->
|
|||||||
test_stats_for_role("Support", SupportToken, loose),
|
test_stats_for_role("Support", SupportToken, loose),
|
||||||
|
|
||||||
test_stats_with_dates(SuperToken),
|
test_stats_with_dates(SuperToken),
|
||||||
|
test_node_metrics_history(SuperToken),
|
||||||
|
|
||||||
% Детальная статистика
|
% Детальная статистика
|
||||||
test_user_stats(SuperToken),
|
test_user_stats(SuperToken),
|
||||||
@@ -220,3 +221,20 @@ test_calendar_stats(Token) ->
|
|||||||
?assert(maps:is_key(<<"top_calendars_by_negative_reviews">>, Stats)),
|
?assert(maps:is_key(<<"top_calendars_by_negative_reviews">>, Stats)),
|
||||||
?assert(maps:is_key(<<"top_calendars_by_rating">>, Stats)),
|
?assert(maps:is_key(<<"top_calendars_by_rating">>, Stats)),
|
||||||
ct:pal(" OK: total=~p", [maps:get(<<"total_calendars">>, Stats)]).
|
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)]).
|
||||||
@@ -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).
|
-module(admin_websocket_tests).
|
||||||
|
-include_lib("eunit/include/eunit.hrl").
|
||||||
-export([test/0]).
|
-export([test/0]).
|
||||||
|
|
||||||
|
%%%===================================================================
|
||||||
|
%%% Главная тестовая функция
|
||||||
|
%%%===================================================================
|
||||||
|
-spec test() -> ok.
|
||||||
test() ->
|
test() ->
|
||||||
ct:pal("Testing WebSocket API..."),
|
ct:pal("=== Admin WebSocket Tests ==="),
|
||||||
application:ensure_all_started(gun),
|
application:ensure_all_started(gun),
|
||||||
|
|
||||||
AdminToken = api_test_runner:get_admin_token(),
|
AdminToken = api_test_runner:get_admin_token(),
|
||||||
UserToken = api_test_runner:get_user_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(
|
#{<<"id">> := CalId} = api_test_runner:client_post(
|
||||||
<<"/v1/calendars">>, UserToken,
|
<<"/v1/calendars">>, UserToken,
|
||||||
#{title => <<"WS Test Calendar">>, type => <<"commercial">>}),
|
#{title => <<"WS Test Calendar">>, type => <<"commercial">>}),
|
||||||
ct:pal(" CalId: ~s", [CalId]),
|
|
||||||
|
|
||||||
#{<<"id">> := EventId} = api_test_runner:client_post(
|
#{<<"id">> := EventId} = api_test_runner:client_post(
|
||||||
<<"/v1/calendars/", CalId/binary, "/events">>, UserToken,
|
<<"/v1/calendars/", CalId/binary, "/events">>, UserToken,
|
||||||
#{title => <<"WS Test Event">>,
|
#{title => <<"WS Test Event">>,
|
||||||
start_time => <<"2026-06-01T10:00:00Z">>,
|
start_time => <<"2026-06-01T10:00:00Z">>,
|
||||||
duration => 60}),
|
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",
|
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
|
% ── Пользовательские WebSocket тесты ──
|
||||||
ct:pal(" TEST 1: Connect WebSocket with valid token..."),
|
{ok, UserWS} = test_user_ws_connect(WsUrl, UserToken),
|
||||||
ct:pal(" URL: ~s", [WsUrl]),
|
test_user_ws_subscribe(UserWS, CalId),
|
||||||
ct:pal(" Token: ~s...", [binary_part(UserToken, 0, 30)]),
|
test_ws_close(UserWS),
|
||||||
case test_ws_connect_debug(WsUrl, UserToken) of
|
|
||||||
{ok, WS} ->
|
|
||||||
ct:pal(" OK - Connected"),
|
|
||||||
|
|
||||||
%% TEST 2: Subscribe to calendar updates
|
ct:pal("~n✅ User WebSocket tests passed!"),
|
||||||
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"),
|
|
||||||
|
|
||||||
|
% ── Административные 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_ws_close(AdminWS),
|
||||||
|
|
||||||
%% TEST 12: Admin WebSocket with user token (should fail)
|
% ── Метрики узлов ──
|
||||||
ct:pal(" TEST 12: Admin WS with user token..."),
|
test_admin_ws_node_metrics(MetricsWsUrl, AdminToken),
|
||||||
{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..."),
|
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">>,
|
Chars = <<"abcdefghijklmnopqrstuvwxyz0123456789">>,
|
||||||
InvalidToken = << <<(binary:at(Chars, rand:uniform(byte_size(Chars)) - 1))>>
|
InvalidToken = << <<(binary:at(Chars, rand:uniform(byte_size(Chars)) - 1))>>
|
||||||
|| _ <- lists:seq(1, 30) >>,
|
|| _ <- lists:seq(1, 30) >>,
|
||||||
{error, {401, _}} = test_ws_connect_debug(AdminWsUrl, InvalidToken),
|
?assertMatch({error, {401, _}}, test_ws_connect_debug(Url, InvalidToken)).
|
||||||
ct:pal(" OK - Rejected"),
|
|
||||||
|
|
||||||
ct:pal("~n✅ Admin WebSocket API tests passed!"),
|
|
||||||
{?MODULE, ok}.
|
|
||||||
|
|
||||||
%% ============ WebSocket хелперы с отладкой ============
|
%% ============ WebSocket хелперы с отладкой ============
|
||||||
test_ws_connect_debug(Url, Token) ->
|
test_ws_connect_debug(Url, Token) ->
|
||||||
|
|||||||
@@ -1,117 +1,71 @@
|
|||||||
|
%%%-------------------------------------------------------------------
|
||||||
|
%%% @doc Тесты пользовательского WebSocket API.
|
||||||
|
%%%
|
||||||
|
%%% Покрывает эндпоинты:
|
||||||
|
%%% ws://localhost:8081/ws
|
||||||
|
%%%
|
||||||
|
%%% Проверяет:
|
||||||
|
%%% - подключение с валидным пользовательским токеном
|
||||||
|
%%% - подписку на календарь и получение подтверждения
|
||||||
|
%%% @end
|
||||||
|
%%%-------------------------------------------------------------------
|
||||||
-module(user_websocket_tests).
|
-module(user_websocket_tests).
|
||||||
|
-include_lib("eunit/include/eunit.hrl").
|
||||||
-export([test/0]).
|
-export([test/0]).
|
||||||
|
|
||||||
|
%%%===================================================================
|
||||||
|
%%% Главная тестовая функция
|
||||||
|
%%%===================================================================
|
||||||
|
-spec test() -> ok.
|
||||||
test() ->
|
test() ->
|
||||||
ct:pal("Testing WebSocket API..."),
|
ct:pal("=== User WebSocket Tests ==="),
|
||||||
application:ensure_all_started(gun),
|
application:ensure_all_started(gun),
|
||||||
|
|
||||||
AdminToken = api_test_runner:get_admin_token(),
|
|
||||||
UserToken = api_test_runner:get_user_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
|
% Создаём календарь и событие
|
||||||
CalId = api_test_runner:create_calendar(UserToken, #{title => <<"WS Test Calendar">>, type => <<"commercial">>}),
|
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">>,
|
title => <<"WS Test Event">>,
|
||||||
start_time => <<"2026-06-01T10:00:00Z">>,
|
start_time => <<"2026-06-01T10:00:00Z">>,
|
||||||
duration => 60
|
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",
|
|
||||||
|
|
||||||
%% TEST 1: Connect to WebSocket with valid token
|
% ── Пользовательские WebSocket тесты ──
|
||||||
ct:pal(" TEST 1: Connect WebSocket with valid token..."),
|
{ok, UserWS} = test_user_ws_connect(WsUrl, UserToken),
|
||||||
ct:pal(" URL: ~s", [WsUrl]),
|
test_user_ws_subscribe(UserWS, CalId),
|
||||||
ct:pal(" Token: ~s...", [binary_part(UserToken, 0, 30)]),
|
test_ws_close(UserWS),
|
||||||
case test_ws_connect_debug(WsUrl, UserToken) of
|
|
||||||
{ok, WS} ->
|
|
||||||
ct:pal(" OK - Connected"),
|
|
||||||
|
|
||||||
%% TEST 2: Subscribe to calendar updates
|
ct:pal("~n=== All user WebSocket tests passed! ==="),
|
||||||
ct:pal(" TEST 2: Subscribe to calendar..."),
|
{?MODULE, ok}.
|
||||||
SubMsg = #{action => <<"subscribe">>, calendar_id => CalId},
|
|
||||||
ct:pal(" Sending: ~p", [SubMsg]),
|
%%%===================================================================
|
||||||
|
%%% Тестовые функции
|
||||||
|
%%%===================================================================
|
||||||
|
|
||||||
|
%% @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),
|
ok = test_ws_send(WS, SubMsg),
|
||||||
case test_ws_recv(WS) of
|
case test_ws_recv(WS) of
|
||||||
{ok, #{<<"status">> := <<"subscribed">>}} ->
|
{ok, #{<<"status">> := <<"subscribed">>}} -> ok;
|
||||||
ct:pal(" OK - Subscribed");
|
{ok, Other} -> error({unexpected_response, Other});
|
||||||
{ok, Other} ->
|
{error, timeout} -> error(timeout)
|
||||||
ct:pal(" ERROR: Unexpected response: ~p", [Other]),
|
end.
|
||||||
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!"),
|
|
||||||
{?MODULE, ok}.
|
|
||||||
|
|
||||||
%% ============ WebSocket хелперы с отладкой ============
|
%% ============ WebSocket хелперы с отладкой ============
|
||||||
test_ws_connect_debug(Url, Token) ->
|
test_ws_connect_debug(Url, Token) ->
|
||||||
|
|||||||
Reference in New Issue
Block a user