ekfa/src/transport/service_channel.erl
2026-04-21 16:24:59 +08:00

163 lines
6.7 KiB
Erlang

%%%-------------------------------------------------------------------
%%% @author licheng5
%%% @copyright (C) 2021, <COMPANY>
%%% @doc
%%%
%%% @end
%%% Created : 11. 1月 2021 上午12:17
%%%-------------------------------------------------------------------
-module(service_channel).
-author("licheng5").
-include("efka_tables.hrl").
-include("protocol.hrl").
-include("service_pb.hrl").
%% API
-export([init/2]).
-export([websocket_init/1, websocket_handle/2, websocket_info/2, terminate/3]).
-record(state, {
service_id :: undefined | binary(),
service_pid :: undefined | pid(),
subscribed_topics = sets:new(),
is_registered = false :: boolean()
}).
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
%% 逻辑处理方法
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
-spec init(term(), term()) -> {cowboy_websocket, term(), term()}.
init(Req, Opts) ->
{cowboy_websocket, Req, Opts}.
-spec websocket_init(term()) -> {ok, #state{}}.
websocket_init(_State) ->
logger:debug("[service_channel] get a new connection"),
%% 初始状态为true
{ok, #state{}}.
-spec websocket_handle(term(), #state{}) ->
{reply, term(), #state{}} | {ok, #state{}}.
websocket_handle(ping, State) ->
{reply, pong, State};
websocket_handle({binary, <<?FRAME_REQUEST, PacketBin/binary>>}, State) ->
Request = service_pb:decode_msg(PacketBin, 'ServiceRequest'),
logger:debug("[service_channel] get request: ~p", [Request]),
handle_request(Request, State);
websocket_handle({binary, <<?FRAME_CAST, PacketBin/binary>>}, State) ->
Cast = service_pb:decode_msg(PacketBin, 'ServiceCast'),
logger:debug("[service_channel] get cast: ~p", [Cast]),
handle_cast(Cast, State);
websocket_handle(Info, State) ->
logger:error("[service_channel] get a unknown message: ~p, channel will closed", [Info]),
{ok, State}.
%% 订阅的消息
-spec websocket_info(term(), #state{}) ->
{reply, term(), #state{}} | {stop, #state{}} | {ok, #state{}}.
websocket_info({topic_broadcast, Topic, Content}, State = #state{}) ->
Packet = service_pb:encode_msg(#'ServiceCast'{
body = {topic_event, #'ServiceCast.TopicEvent'{topic = Topic, content = Content}}
}),
logger:debug("[service_channel] will publish topic: ~p", [Topic]),
{reply, {binary, <<?FRAME_CAST, Packet/binary>>}, State};
%% service进程关闭
websocket_info({'DOWN', _Ref, process, ServicePid, Reason}, State = #state{service_pid = ServicePid}) ->
logger:debug("[service_channel] container_pid: ~p, exited: ~p", [ServicePid, Reason]),
{stop, State#state{service_pid = undefined}};
%% 处理关闭信号
websocket_info({stop, Reason}, State) ->
logger:debug("[service_channel] the channel will be closed with reason: ~p", [Reason]),
{stop, State};
%% 处理其他未知消息
websocket_info(Info, State) ->
logger:debug("[service_channel] channel get unknown info: ~p", [Info]),
{ok, State}.
%% 进程关闭事件
-spec terminate(term(), term(), #state{}) -> ok.
terminate(Reason, _Req, State = #state{service_id = ServiceId, is_registered = IsRegistered}) ->
ok = efka_subscription:unsubscribe_all(self()),
case IsRegistered of
true ->
ok = service_model:change_status(ServiceId, 0);
false ->
ok
end,
logger:debug("[service_channel] channel close with reason: ~p, state is: ~p", [Reason, State]),
ok.
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
%% helper methods
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
%% 注册, 要建立程序和容器之间的关系
-spec handle_request(service_pb:'ServiceRequest'(), #state{}) -> {reply, {binary, binary()}, #state{}}.
handle_request(#'ServiceRequest'{packet_id = PacketId, request = {register, #'ServiceRequest.Register'{service_id = ServiceId}}}, State) ->
{ok, ServicePid} = efka_service_sup:start_service(ServiceId),
case efka_service:attach_channel(ServicePid, self()) of
ok ->
erlang:monitor(process, ServicePid),
%% 更新微服务的状态
ok = service_model:insert(#service{
service_id = ServiceId,
container_name = <<>>,
status = ?SERVICE_RUNNING,
meta_data = #{},
create_ts = efka_util:timestamp(),
update_ts = efka_util:timestamp()
}),
{reply, {binary, result_reply_packet(PacketId, <<"ok">>)},
State#state{service_id = ServiceId, service_pid = ServicePid, is_registered = true}};
{error, Error} ->
logger:warning("[service_channel] service_id: ~p, attach_channel get error: ~p", [ServiceId, Error]),
{reply, {binary, error_reply_packet(PacketId, -1, <<"attach channel failed">>)}, State}
end;
%% 订阅事件
handle_request(#'ServiceRequest'{packet_id = PacketId, request = {subscribe, #'ServiceRequest.Subscribe'{topic = Topic}}},
State = #state{ subscribed_topics = SubscribedTopics, is_registered = true}) ->
case efka_subscription:subscribe(Topic, self()) of
ok ->
Packet = result_reply_packet(PacketId, <<"ok">>),
{reply, {binary, Packet}, State#state{subscribed_topics = sets:add_element(Topic, SubscribedTopics)}};
{error, Reason} ->
Packet = error_reply_packet(PacketId, -1, Reason),
{reply, {binary, Packet}, State}
end;
handle_request(#'ServiceRequest'{packet_id = PacketId}, State) ->
{reply, {binary, error_reply_packet(PacketId, -1, <<"invalid request">>)}, State}.
-spec handle_cast(service_pb:'ServiceCast'(), #state{}) -> {ok, #state{}}.
handle_cast(#'ServiceCast'{body = {metric_data, #'ServiceCast.MetricData'{route_key = RouteKey, metric = Metric}}},
State = #state{service_pid = ServicePid, is_registered = true}) ->
efka_service:metric_data(ServicePid, RouteKey, Metric),
{ok, State};
handle_cast(#'ServiceCast'{body = _Body}, State) ->
{ok, State}.
-spec result_reply_packet(integer(), binary()) -> binary().
result_reply_packet(PacketId, Result) when is_integer(PacketId), is_binary(Result) ->
Reply = service_pb:encode_msg(#'ServiceReply'{
packet_id = PacketId,
reply = {result, Result}
}),
<<?FRAME_REPLY, Reply/binary>>.
-spec error_reply_packet(integer(), integer(), binary()) -> binary().
error_reply_packet(PacketId, Code, Message) when is_integer(PacketId), is_integer(Code), is_binary(Message) ->
Reply = service_pb:encode_msg(#'ServiceReply'{
packet_id = PacketId,
reply = {error, #'ServiceReply.Error'{code = Code, message = Message}}
}),
<<?FRAME_REPLY, Reply/binary>>.