fix channel

This commit is contained in:
anlicheng 2026-04-20 20:12:49 +08:00
parent bbc4b81d2e
commit 1bb9c33ded
3 changed files with 1619 additions and 157 deletions

72
include/service_pb.hrl Normal file
View File

@ -0,0 +1,72 @@
%% -*- coding: utf-8 -*-
%% Automatically generated, do not edit
%% Generated by gpb_compile version 4.21.7
-ifndef(service_pb).
-define(service_pb, true).
-define(service_pb_gpb_version, "4.21.7").
-ifndef('SERVICEREQUEST.REGISTER_PB_H').
-define('SERVICEREQUEST.REGISTER_PB_H', true).
-record('ServiceRequest.Register',
{service_id = <<>> :: unicode:chardata() | undefined % = 1, optional
}).
-endif.
-ifndef('SERVICEREQUEST.SUBSCRIBE_PB_H').
-define('SERVICEREQUEST.SUBSCRIBE_PB_H', true).
-record('ServiceRequest.Subscribe',
{topic = <<>> :: unicode:chardata() | undefined % = 1, optional
}).
-endif.
-ifndef('SERVICEREQUEST_PB_H').
-define('SERVICEREQUEST_PB_H', true).
-record('ServiceRequest',
{packet_id = 0 :: non_neg_integer() | undefined, % = 1, optional, 32 bits
request :: {register, service_pb:'ServiceRequest.Register'()} | {subscribe, service_pb:'ServiceRequest.Subscribe'()} | undefined % oneof
}).
-endif.
-ifndef('SERVICEREPLY.ERROR_PB_H').
-define('SERVICEREPLY.ERROR_PB_H', true).
-record('ServiceReply.Error',
{code = 0 :: integer() | undefined, % = 1, optional, 32 bits
message = <<>> :: unicode:chardata() | undefined % = 2, optional
}).
-endif.
-ifndef('SERVICEREPLY_PB_H').
-define('SERVICEREPLY_PB_H', true).
-record('ServiceReply',
{packet_id = 0 :: non_neg_integer() | undefined, % = 1, optional, 32 bits
reply :: {result, iodata()} | {error, service_pb:'ServiceReply.Error'()} | undefined % oneof
}).
-endif.
-ifndef('SERVICECAST.METRICDATA_PB_H').
-define('SERVICECAST.METRICDATA_PB_H', true).
-record('ServiceCast.MetricData',
{route_key = <<>> :: iodata() | undefined, % = 1, optional
metric = <<>> :: iodata() | undefined % = 2, optional
}).
-endif.
-ifndef('SERVICECAST.TOPICEVENT_PB_H').
-define('SERVICECAST.TOPICEVENT_PB_H', true).
-record('ServiceCast.TopicEvent',
{topic = <<>> :: unicode:chardata() | undefined, % = 1, optional
content = <<>> :: iodata() | undefined % = 2, optional
}).
-endif.
-ifndef('SERVICECAST_PB_H').
-define('SERVICECAST_PB_H', true).
-record('ServiceCast',
{body :: {topic_event, service_pb:'ServiceCast.TopicEvent'()} | {metric_data, service_pb:'ServiceCast.MetricData'()} | undefined % oneof
}).
-endif.
-endif.

1498
src/protobuf/service_pb.erl Normal file

File diff suppressed because it is too large Load Diff

View File

@ -9,6 +9,7 @@
-module(ws_channel). -module(ws_channel).
-author("licheng5"). -author("licheng5").
-include("efka_tables.hrl"). -include("efka_tables.hrl").
-include("service_pb.hrl").
%% REQUEST: %% REQUEST:
%% RESPONSE: REQUEST %% RESPONSE: REQUEST
@ -21,17 +22,9 @@
-export([init/2]). -export([init/2]).
-export([websocket_init/1, websocket_handle/2, websocket_info/2, terminate/3]). -export([websocket_init/1, websocket_handle/2, websocket_info/2, terminate/3]).
%%
-define(PENDING_TIMEOUT, 10 * 1000).
-record(state, { -record(state, {
service_id :: undefined | binary(), service_id :: undefined | binary(),
service_pid :: undefined | pid(), service_pid :: undefined | pid(),
stream_id = 1,
%% #{stream_id => {StreamPid, StreamRef}}
stream_map = #{},
is_registered = false :: boolean() is_registered = false :: boolean()
}). }).
@ -50,10 +43,14 @@ websocket_init(_State) ->
websocket_handle(ping, State) -> websocket_handle(ping, State) ->
{reply, pong, State}; {reply, pong, State};
websocket_handle({text, Data}, State) -> websocket_handle({binary, <<?FRAME_REQUEST, PacketBin/binary>>}, State) ->
Request = jiffy:decode(Data, [return_maps]), Request = service_pb:decode_msg(PacketBin, 'ServiceRequest'),
logger:debug("[ws_channle] get request: ~p", [Request]), logger:debug("[ws_channel] get request: ~p", [Request]),
handle_request(Request, State); handle_request(Request, State);
websocket_handle({binary, <<?FRAME_CAST, PacketBin/binary>>}, State) ->
Cast = service_pb:decode_msg(PacketBin, 'ServiceCast'),
logger:debug("[ws_channel] get cast: ~p", [Cast]),
handle_cast(Cast, State);
websocket_handle(Info, State) -> websocket_handle(Info, State) ->
logger:error("[ws_channel] get a unknown message: ~p, channel will closed", [Info]), logger:error("[ws_channel] get a unknown message: ~p, channel will closed", [Info]),
@ -61,55 +58,17 @@ websocket_handle(Info, State) ->
%% %%
websocket_info({topic_broadcast, Topic, Content}, State = #state{}) -> websocket_info({topic_broadcast, Topic, Content}, State = #state{}) ->
Req = iolist_to_binary(jiffy:encode(#{ Packet = service_pb:encode_msg(#'ServiceCast'{
<<"method">> => <<"publish">>, body = {topic_event, #'ServiceCast.TopicEvent'{topic = Topic, content = Content}}
<<"params">> => #{<<"topic">> => Topic, <<"content">> => Content} }),
}, [force_utf8])), logger:debug("[ws_channel] will publish topic: ~p", [Topic]),
{reply, {binary, <<?FRAME_CAST, Packet/binary>>}, State};
logger:debug("[ws_channel] will publish topic: ~p, message: ~p", [Topic, Req]),
{reply, {text, Req}, State};
%% service进程关闭 %% service进程关闭
websocket_info({'DOWN', _Ref, process, ServicePid, Reason}, State = #state{service_pid = ServicePid}) -> websocket_info({'DOWN', _Ref, process, ServicePid, Reason}, State = #state{service_pid = ServicePid}) ->
logger:debug("[ws_channel] container_pid: ~p, exited: ~p", [ServicePid, Reason]), logger:debug("[ws_channel] container_pid: ~p, exited: ~p", [ServicePid, Reason]),
{stop, State#state{service_pid = undefined}}; {stop, State#state{service_pid = undefined}};
%% stream进程关闭
websocket_info({'DOWN', _Ref, process, StreamPid, Reason}, State = #state{stream_map = StreamMap}) ->
case search_stream_id(StreamPid, StreamMap) of
error ->
{ok, State};
{ok, StreamId} ->
case Reason of
normal ->
{ok, State#state{stream_map = maps:remove(StreamId, StreamMap)}};
_ ->
PushReply = json_push(#{
<<"stream_reply">> => #{
<<"stream_id">> => StreamId,
<<"result">> => <<"task failed">>
}
}),
{reply, {text, PushReply}, State#state{stream_map = maps:remove(StreamId, StreamMap)}}
end
end;
%% stream任务完成
websocket_info({stream_reply, StreamPid, Reply}, State = #state{stream_map = StreamMap}) ->
case search_stream_id(StreamPid, StreamMap) of
error ->
{ok, State};
{ok, StreamId} ->
PushReply = json_push(#{
<<"stream_reply">> => #{
<<"stream_id">> => StreamId,
<<"result">> => Reply
}
}),
{reply, {text, PushReply}, State}
end;
%% %%
websocket_info({stop, Reason}, State) -> websocket_info({stop, Reason}, State) ->
logger:debug("[ws_channel] the channel will be closed with reason: ~p", [Reason]), logger:debug("[ws_channel] the channel will be closed with reason: ~p", [Reason]),
@ -136,129 +95,62 @@ terminate(Reason, _Req, State = #state{service_id = ServiceId, is_registered = I
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%% %%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
%% , %% ,
handle_request(#{<<"id">> := Id, <<"method">> := <<"register">>, <<"params">> := Params = #{<<"service_id">> := ServiceId}}, State) -> handle_request(#'ServiceRequest'{packet_id = PacketId, request = {register, #'ServiceRequest.Register'{service_id = ServiceId}}}, State) ->
{ok, ServicePid} = efka_service_sup:start_service(ServiceId), {ok, ServicePid} = efka_service_sup:start_service(ServiceId),
case efka_service:attach_channel(ServicePid, self()) of case efka_service:attach_channel(ServicePid, self()) of
ok -> ok ->
Reply = json_result(Id, <<"ok">>),
erlang:monitor(process, ServicePid), erlang:monitor(process, ServicePid),
%% %%
MetaData = maps:get(<<"meta_data">>, Params, #{}),
ContainerName = maps:get(<<"container_name">>, Params, <<>>),
ok = service_model:insert(#service{ ok = service_model:insert(#service{
service_id = ServiceId, service_id = ServiceId,
container_name = ContainerName, container_name = <<>>,
status = ?SERVICE_RUNNING, status = ?SERVICE_RUNNING,
meta_data = MetaData, meta_data = #{},
create_ts = efka_util:timestamp(), create_ts = efka_util:timestamp(),
update_ts = efka_util:timestamp() update_ts = efka_util:timestamp()
}), }),
{reply, {text, Reply}, State#state{service_id = ServiceId, service_pid = ServicePid, is_registered = true}}; {reply, {binary, result_reply_packet(PacketId, <<"ok">>)},
State#state{service_id = ServiceId, service_pid = ServicePid, is_registered = true}};
{error, Error} -> {error, Error} ->
logger:warning("[ws_channel] service_id: ~p, attach_channel get error: ~p", [ServiceId, Error]), logger:warning("[ws_channel] service_id: ~p, attach_channel get error: ~p", [ServiceId, Error]),
{stop, State} {reply, {binary, error_reply_packet(PacketId, -1, <<"attach channel failed">>)}, State}
end; end;
%% %%
handle_request(#{<<"id">> := Id, <<"method">> := <<"subscribe">>, <<"params">> := #{<<"topic">> := Topic}}, State = #state{is_registered = true}) -> handle_request(#'ServiceRequest'{packet_id = PacketId, request = {subscribe, #'ServiceRequest.Subscribe'{topic = Topic}}},
Reply = case efka_subscription:subscribe(Topic, self()) of State = #state{is_registered = true}) ->
Packet = case efka_subscription:subscribe(Topic, self()) of
ok -> ok ->
json_result(Id, <<"ok">>); result_reply_packet(PacketId, <<"ok">>);
{error, Reason} -> {error, Reason} ->
json_error(Id, -1, Reason) error_reply_packet(PacketId, -1, Reason)
end, end,
{reply, {text, Reply}, State}; {reply, {binary, Packet}, State};
handle_request(#'ServiceRequest'{packet_id = PacketId}, State) ->
{reply, {binary, error_reply_packet(PacketId, -1, <<"invalid request">>)}, State}.
%% handle_cast(#'ServiceCast'{body = {metric_data, #'ServiceCast.MetricData'{route_key = RouteKey, metric = Metric}}},
handle_request(#{<<"id">> := Id, <<"method">> := <<"new_stream">>, State = #state{service_pid = ServicePid, is_registered = true}) ->
<<"params">> := #{<<"file_name">> := Filename0, <<"file_size">> := FileSize}}, State = #state{stream_id = StreamId, stream_map = StreamMap, is_registered = true}) -> efka_service:metric_data(ServicePid, RouteKey, Metric),
Filename = filename:basename(binary_to_list(Filename0)),
{ok, {StreamPid, StreamRef}} = efka_stream:start_monitor(self()),
{ok, Path} = efka_stream:setup(StreamPid, Filename, FileSize),
Reply = json_result(Id, #{
<<"stream_id">> => StreamId,
<<"path">> => Path
}),
{reply, {text, Reply}, State#state{stream_id = StreamId + 1, stream_map = maps:put(StreamId, {StreamPid, StreamRef}, StreamMap)}};
handle_request(#{<<"method">> := <<"stream_chunk">>,
<<"params">> := #{<<"stream_id">> := StreamId, <<"chunk_data">> := ChunkData}}, State = #state{stream_map = StreamMap, is_registered = true}) ->
case maps:find(StreamId, StreamMap) of
error ->
{ok, State}; {ok, State};
{ok, {StreamPid, _}} -> handle_cast(#'ServiceCast'{body = _Body}, State) ->
case ChunkData =:= <<>> of
true ->
efka_stream:finish(StreamPid);
false ->
efka_stream:data(StreamPid, ChunkData)
end,
{ok, State}
end;
%%
handle_request(#{<<"method">> := <<"metric_data">>,
<<"params">> := #{<<"route_key">> := RouteKey, <<"metric">> := Metric0}}, State = #state{service_pid = ServicePid, is_registered = true}) ->
case map_metric(Metric0) of
{ok, Metric} ->
efka_service:metric_data(ServicePid, RouteKey, Metric);
error ->
logger:debug("[ws_channel] metric_data get invalid metric: ~p", [Metric0])
end,
{ok, State}. {ok, State}.
-spec json_result(Id :: integer(), Result :: term()) -> binary(). -spec result_reply_packet(integer(), binary()) -> binary().
json_result(Id, Result) when is_integer(Id) -> result_reply_packet(PacketId, Result) when is_integer(PacketId), is_binary(Result) ->
Response = #{ Reply = service_pb:encode_msg(#'ServiceReply'{
<<"id">> => Id, packet_id = PacketId,
<<"result">> => Result reply = {result, Result}
}, }),
jiffy:encode(Response, [force_utf8]). <<?FRAME_RESPONSE, Reply/binary>>.
-spec json_error(Id :: integer(), Code :: integer(), Message :: binary()) -> binary(). -spec error_reply_packet(integer(), integer(), binary()) -> binary().
json_error(Id, Code, Message) when is_integer(Id), is_integer(Code), is_binary(Message) -> error_reply_packet(PacketId, Code, Message) when is_integer(PacketId), is_integer(Code), is_binary(Message) ->
Response = #{ Reply = service_pb:encode_msg(#'ServiceReply'{
<<"id">> => Id, packet_id = PacketId,
<<"error">> => #{<<"code">> => Code, <<"message">> => Message} reply = {error, #'ServiceReply.Error'{code = Code, message = Message}}
}, }),
jiffy:encode(Response, [force_utf8]). <<?FRAME_RESPONSE, Reply/binary>>.
-spec json_push(Result :: term()) -> binary().
json_push(Result) ->
Response = #{
<<"push">> => Result
},
jiffy:encode(Response, [force_utf8]).
-spec search_stream_id(StreamPid :: pid(), StreamMap :: map()) -> error | {ok, StreamId :: integer()}.
search_stream_id(StreamPid, StreamMap) when is_pid(StreamPid), is_map(StreamMap) ->
StreamIds = lists:filtermap(fun({StreamId, {StreamPid0, _}}) ->
case StreamPid0 =:= StreamPid of
true ->
{true, StreamId};
false ->
false
end
end, maps:to_list(StreamMap)),
case StreamIds of
[] ->
error;
[StreamId|_] ->
{ok, StreamId}
end.
-spec map_metric(Metric :: any()) -> {ok, binary()} | error.
map_metric(Metric) when is_binary(Metric) ->
{ok, Metric};
map_metric(Metric) when is_map(Metric) orelse is_list(Metric) ->
{ok, jiffy:encode(Metric, [force_utf8])};
map_metric(Metric) when is_integer(Metric) ->
{ok, integer_to_binary(Metric)};
map_metric(Metric) when is_float(Metric) ->
{ok, erlang:float_to_binary(Metric, [compact, {decimals, 10}])};
map_metric(_) ->
error.