模块整理
This commit is contained in:
parent
651e7c81a5
commit
59fa7cf822
@ -34,14 +34,14 @@ WebSocket server:
|
|||||||
|
|
||||||
- 从 `config/sys.config` 读取 `efka.websocket_server`。
|
- 从 `config/sys.config` 读取 `efka.websocket_server`。
|
||||||
- 监听 `/ws`。
|
- 监听 `/ws`。
|
||||||
- 请求处理模块是 `service_channel`。
|
- 请求处理模块是 `efka_service_channel`。
|
||||||
|
|
||||||
`efka_sup` 启动的子进程:
|
`efka_sup` 启动的子进程:
|
||||||
|
|
||||||
- `efka_logger`:部署日志落盘。
|
- `efka_logger`:部署日志落盘。
|
||||||
- `efka_service_sup`:动态管理每个已注册微服务对应的 `efka_service` 进程。
|
- `efka_service_sup`:动态管理每个已注册微服务对应的 `efka_service` 进程。
|
||||||
- `cache_model`:DETS 离线缓存。
|
- `cache_model`:DETS 离线缓存。
|
||||||
- `service_model`:DETS 服务状态表。
|
- `efka_service_model`:DETS 服务状态表。
|
||||||
- `efka_subscription`:本地 topic 订阅中心。
|
- `efka_subscription`:本地 topic 订阅中心。
|
||||||
- `efka_client`:连接上游 TLS server 的状态机。
|
- `efka_client`:连接上游 TLS server 的状态机。
|
||||||
|
|
||||||
@ -60,12 +60,12 @@ WebSocket server:
|
|||||||
|
|
||||||
本机微服务连接 `ws://<host>:18080/ws`。
|
本机微服务连接 `ws://<host>:18080/ws`。
|
||||||
|
|
||||||
协议处理在 `service_channel`:
|
协议处理在 `efka_service_channel`:
|
||||||
|
|
||||||
- WebSocket 收到 `ping` 直接回复 `pong`。
|
- WebSocket 收到 `ping` 直接回复 `pong`。
|
||||||
- 收到 binary 帧时,第一个字节是帧类型。
|
- 收到 binary 帧时,第一个字节是帧类型。
|
||||||
- `FRAME_REQUEST = 0x01`,用 `service_pb` 解码为 `ServiceRequest`。
|
- `FRAME_REQUEST = 0x01`,用 `efka_service_pb` 解码为 `ServiceRequest`。
|
||||||
- `FRAME_CAST = 0x03`,用 `service_pb` 解码为 `ServiceCast`。
|
- `FRAME_CAST = 0x03`,用 `efka_service_pb` 解码为 `ServiceCast`。
|
||||||
- 回复时使用 `FRAME_REPLY = 0x02`。
|
- 回复时使用 `FRAME_REPLY = 0x02`。
|
||||||
|
|
||||||
支持的微服务请求:
|
支持的微服务请求:
|
||||||
@ -80,17 +80,17 @@ WebSocket server:
|
|||||||
注册流程:
|
注册流程:
|
||||||
|
|
||||||
1. 微服务发送 `ServiceRequest.Register`,包含 `service_id`。
|
1. 微服务发送 `ServiceRequest.Register`,包含 `service_id`。
|
||||||
2. `service_channel` 调用 `efka_service_sup:start_service(ServiceId)`。
|
2. `efka_service_channel` 调用 `efka_service_sup:start_service(ServiceId)`。
|
||||||
3. 如果该服务进程不存在,则动态启动一个 `efka_service`;如果已存在,复用已有 pid。
|
3. 如果该服务进程不存在,则动态启动一个 `efka_service`;如果已存在,复用已有 pid。
|
||||||
4. `service_channel` 调用 `efka_service:attach_channel(ServicePid, self())` 绑定当前 WebSocket channel。
|
4. `efka_service_channel` 调用 `efka_service:attach_channel(ServicePid, self())` 绑定当前 WebSocket channel。
|
||||||
5. `efka_service` 只允许绑定一个存活 channel;已有 channel 存活时返回 `channel exists`。
|
5. `efka_service` 只允许绑定一个存活 channel;已有 channel 存活时返回 `channel exists`。
|
||||||
6. 注册成功后写入 `service_model`,服务状态设为 running。
|
6. 注册成功后写入 `efka_service_model`,服务状态设为 running。
|
||||||
7. channel 退出时,取消全部订阅,并把服务状态改为 stopped。
|
7. channel 退出时,取消全部订阅,并把服务状态改为 stopped。
|
||||||
|
|
||||||
指标上报流程:
|
指标上报流程:
|
||||||
|
|
||||||
1. 微服务发送 `ServiceCast.MetricData`。
|
1. 微服务发送 `ServiceCast.MetricData`。
|
||||||
2. `service_channel` 调用 `efka_service:metric_data(ServicePid, RouteKey, Metric)`。
|
2. `efka_service_channel` 调用 `efka_service:metric_data(ServicePid, RouteKey, Metric)`。
|
||||||
3. `efka_service` 转发给 `efka_client:metric_data(RouteKey, Metric)`。
|
3. `efka_service` 转发给 `efka_client:metric_data(RouteKey, Metric)`。
|
||||||
4. `efka_client` 如果处于 activated 状态,直接发给上游;否则写入 `cache_model` 离线缓存。
|
4. `efka_client` 如果处于 activated 状态,直接发给上游;否则写入 `cache_model` 离线缓存。
|
||||||
|
|
||||||
@ -120,7 +120,7 @@ WebSocket server:
|
|||||||
- `FRAME_REPLY = 0x02`
|
- `FRAME_REPLY = 0x02`
|
||||||
- `FRAME_CAST = 0x03`
|
- `FRAME_CAST = 0x03`
|
||||||
|
|
||||||
但 protobuf 使用的是 `message_pb`,不是微服务侧的 `service_pb`。
|
但 protobuf 使用的是 `message_pb`,不是微服务侧的 `efka_service_pb`。
|
||||||
|
|
||||||
`efka_client` 上报的内容:
|
`efka_client` 上报的内容:
|
||||||
|
|
||||||
@ -140,7 +140,7 @@ WebSocket server:
|
|||||||
|
|
||||||
订阅:
|
订阅:
|
||||||
|
|
||||||
- `service_channel` 收到微服务 `subscribe` 后调用 `efka_subscription:subscribe(Topic, self())`。
|
- `efka_service_channel` 收到微服务 `subscribe` 后调用 `efka_subscription:subscribe(Topic, self())`。
|
||||||
- 同一个 pid 对同一个 topic 重复订阅会被忽略。
|
- 同一个 pid 对同一个 topic 重复订阅会被忽略。
|
||||||
- 首次看到某个 subscriber pid 时,会 monitor 该 pid;pid 退出后自动移除订阅。
|
- 首次看到某个 subscriber pid 时,会 monitor 该 pid;pid 退出后自动移除订阅。
|
||||||
|
|
||||||
@ -154,8 +154,8 @@ topic 匹配规则:
|
|||||||
发布:
|
发布:
|
||||||
|
|
||||||
- 上游发来的 `CastFrame.Pub` 会调用 `efka_subscription:publish(Topic, Qos, Content)`。
|
- 上游发来的 `CastFrame.Pub` 会调用 `efka_subscription:publish(Topic, Qos, Content)`。
|
||||||
- 如果匹配到订阅者,向对应 `service_channel` 发送 `{topic_broadcast, Topic, Content}`。
|
- 如果匹配到订阅者,向对应 `efka_service_channel` 发送 `{topic_broadcast, Topic, Content}`。
|
||||||
- `service_channel` 再编码为 `ServiceCast.TopicEvent` 推送给微服务。
|
- `efka_service_channel` 再编码为 `ServiceCast.TopicEvent` 推送给微服务。
|
||||||
- 如果没有订阅者且 `Qos != 0`,消息会暂存在 `remaining_messages`。
|
- 如果没有订阅者且 `Qos != 0`,消息会暂存在 `remaining_messages`。
|
||||||
- 新订阅建立时,会尝试把匹配的 remaining messages 补发给该订阅者。
|
- 新订阅建立时,会尝试把匹配的 remaining messages 补发给该订阅者。
|
||||||
|
|
||||||
@ -215,7 +215,7 @@ topic 匹配规则:
|
|||||||
|
|
||||||
## 6. 本地持久化
|
## 6. 本地持久化
|
||||||
|
|
||||||
### 6.1 `service_model`
|
### 6.1 `efka_service_model`
|
||||||
|
|
||||||
使用 DETS 表 `service`。
|
使用 DETS 表 `service`。
|
||||||
|
|
||||||
@ -272,13 +272,13 @@ topic 匹配规则:
|
|||||||
|
|
||||||
注意:
|
注意:
|
||||||
|
|
||||||
- `service_model` 和 `cache_model` 打开 DETS 前假设 `dets_dir` 已存在,代码里没有显式创建目录。
|
- `efka_service_model` 和 `cache_model` 打开 DETS 前假设 `dets_dir` 已存在,代码里没有显式创建目录。
|
||||||
- `docker_client` 固定使用 `/var/run/docker.sock`,每次请求内部由短生命周期普通进程执行 open/request/close。
|
- `docker_client` 固定使用 `/var/run/docker.sock`,每次请求内部由短生命周期普通进程执行 open/request/close。
|
||||||
|
|
||||||
## 9. 当前看到的几个注意点
|
## 9. 当前看到的几个注意点
|
||||||
|
|
||||||
- `README.md` 描述的是 JSON-RPC 风格 WebSocket,但当前代码实际处理的是 binary protobuf 帧;README 可能已经过期。
|
- `README.md` 描述的是 JSON-RPC 风格 WebSocket,但当前代码实际处理的是 binary protobuf 帧;README 可能已经过期。
|
||||||
- `service_channel` 里 `register` 只使用 `service_id`,没有使用 README 中提到的 `meta_data/container_name`。
|
- `efka_service_channel` 里 `register` 只使用 `service_id`,没有使用 README 中提到的 `meta_data/container_name`。
|
||||||
- `efka_client:send_result_reply/3` 和 `send_error_reply/3` 编码 `ReplyFrame` 后没有加 `FRAME_REPLY` 前缀;接收侧是否期望裸 protobuf 需要确认。
|
- `efka_client:send_result_reply/3` 和 `send_error_reply/3` 编码 `ReplyFrame` 后没有加 `FRAME_REPLY` 前缀;接收侧是否期望裸 protobuf 需要确认。
|
||||||
- `docker_deployer:ensure_container_absent/2` 当前没有真正确保旧容器不存在,只是上报日志。
|
- `docker_deployer:ensure_container_absent/2` 当前没有真正确保旧容器不存在,只是上报日志。
|
||||||
- `efka_subscription` 计算了 topic `order`,但匹配广播时没有使用优先级排序。
|
- `efka_subscription` 计算了 topic `order`,但匹配广播时没有使用优先级排序。
|
||||||
|
|||||||
@ -2,10 +2,10 @@
|
|||||||
%% Automatically generated, do not edit
|
%% Automatically generated, do not edit
|
||||||
%% Generated by gpb_compile version 4.21.7
|
%% Generated by gpb_compile version 4.21.7
|
||||||
|
|
||||||
-ifndef(service_pb).
|
-ifndef(efka_service_pb).
|
||||||
-define(service_pb, true).
|
-define(efka_service_pb, true).
|
||||||
|
|
||||||
-define(service_pb_gpb_version, "4.21.7").
|
-define(efka_service_pb_gpb_version, "4.21.7").
|
||||||
|
|
||||||
|
|
||||||
-ifndef('SERVICEREQUEST.REGISTER_PB_H').
|
-ifndef('SERVICEREQUEST.REGISTER_PB_H').
|
||||||
@ -26,7 +26,7 @@
|
|||||||
-define('SERVICEREQUEST_PB_H', true).
|
-define('SERVICEREQUEST_PB_H', true).
|
||||||
-record('ServiceRequest',
|
-record('ServiceRequest',
|
||||||
{packet_id = 0 :: non_neg_integer() | undefined, % = 1, optional, 32 bits
|
{packet_id = 0 :: non_neg_integer() | undefined, % = 1, optional, 32 bits
|
||||||
request :: {register, service_pb:'ServiceRequest.Register'()} | {subscribe, service_pb:'ServiceRequest.Subscribe'()} | undefined % oneof
|
request :: {register, efka_service_pb:'ServiceRequest.Register'()} | {subscribe, efka_service_pb:'ServiceRequest.Subscribe'()} | undefined % oneof
|
||||||
}).
|
}).
|
||||||
-endif.
|
-endif.
|
||||||
|
|
||||||
@ -42,7 +42,7 @@
|
|||||||
-define('SERVICEREPLY_PB_H', true).
|
-define('SERVICEREPLY_PB_H', true).
|
||||||
-record('ServiceReply',
|
-record('ServiceReply',
|
||||||
{packet_id = 0 :: non_neg_integer() | undefined, % = 1, optional, 32 bits
|
{packet_id = 0 :: non_neg_integer() | undefined, % = 1, optional, 32 bits
|
||||||
reply :: {result, iodata()} | {error, service_pb:'ServiceReply.Error'()} | undefined % oneof
|
reply :: {result, iodata()} | {error, efka_service_pb:'ServiceReply.Error'()} | undefined % oneof
|
||||||
}).
|
}).
|
||||||
-endif.
|
-endif.
|
||||||
|
|
||||||
@ -65,7 +65,7 @@
|
|||||||
-ifndef('SERVICECAST_PB_H').
|
-ifndef('SERVICECAST_PB_H').
|
||||||
-define('SERVICECAST_PB_H', true).
|
-define('SERVICECAST_PB_H', true).
|
||||||
-record('ServiceCast',
|
-record('ServiceCast',
|
||||||
{body :: {topic_event, service_pb:'ServiceCast.TopicEvent'()} | {metric_data, service_pb:'ServiceCast.MetricData'()} | undefined % oneof
|
{body :: {topic_event, efka_service_pb:'ServiceCast.TopicEvent'()} | {metric_data, efka_service_pb:'ServiceCast.MetricData'()} | undefined % oneof
|
||||||
}).
|
}).
|
||||||
-endif.
|
-endif.
|
||||||
|
|
||||||
@ -4,6 +4,7 @@
|
|||||||
{i, "proto"},
|
{i, "proto"},
|
||||||
{f, ["service.proto"]},
|
{f, ["service.proto"]},
|
||||||
recursive,
|
recursive,
|
||||||
|
{module_name_prefix, "efka_"},
|
||||||
{module_name_suffix, "_pb"},
|
{module_name_suffix, "_pb"},
|
||||||
{o_erl, "src/protobuf"},
|
{o_erl, "src/protobuf"},
|
||||||
{o_hrl, "include"},
|
{o_hrl, "include"},
|
||||||
|
|||||||
@ -32,7 +32,7 @@ start_websocket_server() ->
|
|||||||
Port = proplists:get_value(port, Props),
|
Port = proplists:get_value(port, Props),
|
||||||
|
|
||||||
Dispatcher = cowboy_router:compile([
|
Dispatcher = cowboy_router:compile([
|
||||||
{'_', [{"/ws", service_channel, []}]}
|
{'_', [{"/ws", efka_service_channel, []}]}
|
||||||
]),
|
]),
|
||||||
|
|
||||||
TransOpts = [
|
TransOpts = [
|
||||||
|
|||||||
@ -49,12 +49,12 @@ init([]) ->
|
|||||||
},
|
},
|
||||||
|
|
||||||
#{
|
#{
|
||||||
id => service_model,
|
id => efka_service_model,
|
||||||
start => {service_model, start_link, []},
|
start => {efka_service_model, start_link, []},
|
||||||
restart => permanent,
|
restart => permanent,
|
||||||
shutdown => 5000,
|
shutdown => 5000,
|
||||||
type => worker,
|
type => worker,
|
||||||
modules => ['service_model']
|
modules => ['efka_service_model']
|
||||||
},
|
},
|
||||||
|
|
||||||
#{
|
#{
|
||||||
|
|||||||
@ -3,7 +3,7 @@
|
|||||||
%% Automatically @generated, do not edit
|
%% Automatically @generated, do not edit
|
||||||
%% Generated by gpb_compile version 4.21.7
|
%% Generated by gpb_compile version 4.21.7
|
||||||
%% Version source: file
|
%% Version source: file
|
||||||
-module(service_pb).
|
-module(efka_service_pb).
|
||||||
|
|
||||||
-export([encode_msg/1, encode_msg/2, encode_msg/3]).
|
-export([encode_msg/1, encode_msg/2, encode_msg/3]).
|
||||||
-export([decode_msg/2, decode_msg/3]).
|
-export([decode_msg/2, decode_msg/3]).
|
||||||
@ -46,7 +46,7 @@
|
|||||||
-export([gpb_version_as_string/0, gpb_version_as_list/0]).
|
-export([gpb_version_as_string/0, gpb_version_as_list/0]).
|
||||||
-export([gpb_version_source/0]).
|
-export([gpb_version_source/0]).
|
||||||
|
|
||||||
-include("service_pb.hrl").
|
-include("efka_service_pb.hrl").
|
||||||
-include_lib("gpb/include/gpb.hrl").
|
-include_lib("gpb/include/gpb.hrl").
|
||||||
|
|
||||||
%% enumerated types
|
%% enumerated types
|
||||||
@ -6,10 +6,10 @@
|
|||||||
%%% @end
|
%%% @end
|
||||||
%%% Created : 11. 1月 2021 上午12:17
|
%%% Created : 11. 1月 2021 上午12:17
|
||||||
%%%-------------------------------------------------------------------
|
%%%-------------------------------------------------------------------
|
||||||
-module(service_channel).
|
-module(efka_service_channel).
|
||||||
-author("licheng5").
|
-author("licheng5").
|
||||||
-include("efka_tables.hrl").
|
-include("efka_tables.hrl").
|
||||||
-include("service_pb.hrl").
|
-include("efka_service_pb.hrl").
|
||||||
|
|
||||||
%% 一级帧类型
|
%% 一级帧类型
|
||||||
%% REQUEST: 需要响应的请求帧
|
%% REQUEST: 需要响应的请求帧
|
||||||
@ -40,7 +40,7 @@ init(Req, Opts) ->
|
|||||||
|
|
||||||
-spec websocket_init(term()) -> {ok, #state{}}.
|
-spec websocket_init(term()) -> {ok, #state{}}.
|
||||||
websocket_init(_State) ->
|
websocket_init(_State) ->
|
||||||
logger:debug("[service_channel] get a new connection"),
|
logger:debug("[efka_service_channel] get a new connection"),
|
||||||
%% 初始状态为true
|
%% 初始状态为true
|
||||||
{ok, #state{}}.
|
{ok, #state{}}.
|
||||||
|
|
||||||
@ -50,41 +50,41 @@ websocket_handle(ping, State) ->
|
|||||||
{reply, pong, State};
|
{reply, pong, State};
|
||||||
|
|
||||||
websocket_handle({binary, <<?FRAME_REQUEST, PacketBin/binary>>}, State) ->
|
websocket_handle({binary, <<?FRAME_REQUEST, PacketBin/binary>>}, State) ->
|
||||||
Request = service_pb:decode_msg(PacketBin, 'ServiceRequest'),
|
Request = efka_service_pb:decode_msg(PacketBin, 'ServiceRequest'),
|
||||||
logger:debug("[service_channel] get request: ~p", [Request]),
|
logger:debug("[efka_service_channel] get request: ~p", [Request]),
|
||||||
handle_request(Request, State);
|
handle_request(Request, State);
|
||||||
websocket_handle({binary, <<?FRAME_CAST, PacketBin/binary>>}, State) ->
|
websocket_handle({binary, <<?FRAME_CAST, PacketBin/binary>>}, State) ->
|
||||||
Cast = service_pb:decode_msg(PacketBin, 'ServiceCast'),
|
Cast = efka_service_pb:decode_msg(PacketBin, 'ServiceCast'),
|
||||||
logger:debug("[service_channel] get cast: ~p", [Cast]),
|
logger:debug("[efka_service_channel] get cast: ~p", [Cast]),
|
||||||
handle_cast(Cast, State);
|
handle_cast(Cast, State);
|
||||||
|
|
||||||
websocket_handle(Info, State) ->
|
websocket_handle(Info, State) ->
|
||||||
logger:error("[service_channel] get a unknown message: ~p, channel will closed", [Info]),
|
logger:error("[efka_service_channel] get a unknown message: ~p, channel will closed", [Info]),
|
||||||
{ok, State}.
|
{ok, State}.
|
||||||
|
|
||||||
%% 订阅的消息
|
%% 订阅的消息
|
||||||
-spec websocket_info(term(), #state{}) ->
|
-spec websocket_info(term(), #state{}) ->
|
||||||
{reply, term(), #state{}} | {stop, #state{}} | {ok, #state{}}.
|
{reply, term(), #state{}} | {stop, #state{}} | {ok, #state{}}.
|
||||||
websocket_info({topic_broadcast, Topic, Content}, State = #state{}) ->
|
websocket_info({topic_broadcast, Topic, Content}, State = #state{}) ->
|
||||||
Packet = service_pb:encode_msg(#'ServiceCast'{
|
Packet = efka_service_pb:encode_msg(#'ServiceCast'{
|
||||||
body = {topic_event, #'ServiceCast.TopicEvent'{topic = Topic, content = Content}}
|
body = {topic_event, #'ServiceCast.TopicEvent'{topic = Topic, content = Content}}
|
||||||
}),
|
}),
|
||||||
logger:debug("[service_channel] will publish topic: ~p", [Topic]),
|
logger:debug("[efka_service_channel] will publish topic: ~p", [Topic]),
|
||||||
{reply, {binary, <<?FRAME_CAST, Packet/binary>>}, State};
|
{reply, {binary, <<?FRAME_CAST, Packet/binary>>}, 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("[service_channel] container_pid: ~p, exited: ~p", [ServicePid, Reason]),
|
logger:debug("[efka_service_channel] container_pid: ~p, exited: ~p", [ServicePid, Reason]),
|
||||||
{stop, State#state{service_pid = undefined}};
|
{stop, State#state{service_pid = undefined}};
|
||||||
|
|
||||||
%% 处理关闭信号
|
%% 处理关闭信号
|
||||||
websocket_info({stop, Reason}, State) ->
|
websocket_info({stop, Reason}, State) ->
|
||||||
logger:debug("[service_channel] the channel will be closed with reason: ~p", [Reason]),
|
logger:debug("[efka_service_channel] the channel will be closed with reason: ~p", [Reason]),
|
||||||
{stop, State};
|
{stop, State};
|
||||||
|
|
||||||
%% 处理其他未知消息
|
%% 处理其他未知消息
|
||||||
websocket_info(Info, State) ->
|
websocket_info(Info, State) ->
|
||||||
logger:debug("[service_channel] channel get unknown info: ~p", [Info]),
|
logger:debug("[efka_service_channel] channel get unknown info: ~p", [Info]),
|
||||||
{ok, State}.
|
{ok, State}.
|
||||||
|
|
||||||
%% 进程关闭事件
|
%% 进程关闭事件
|
||||||
@ -93,11 +93,11 @@ terminate(Reason, _Req, State = #state{service_id = ServiceId, is_registered = I
|
|||||||
ok = efka_subscription:unsubscribe_all(self()),
|
ok = efka_subscription:unsubscribe_all(self()),
|
||||||
case IsRegistered of
|
case IsRegistered of
|
||||||
true ->
|
true ->
|
||||||
ok = service_model:change_status(ServiceId, 0);
|
ok = efka_service_model:change_status(ServiceId, 0);
|
||||||
false ->
|
false ->
|
||||||
ok
|
ok
|
||||||
end,
|
end,
|
||||||
logger:debug("[service_channel] channel close with reason: ~p, state is: ~p", [Reason, State]),
|
logger:debug("[efka_service_channel] channel close with reason: ~p, state is: ~p", [Reason, State]),
|
||||||
ok.
|
ok.
|
||||||
|
|
||||||
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
|
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
|
||||||
@ -105,7 +105,7 @@ terminate(Reason, _Req, State = #state{service_id = ServiceId, is_registered = I
|
|||||||
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
|
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
|
||||||
|
|
||||||
%% 注册, 要建立程序和容器之间的关系
|
%% 注册, 要建立程序和容器之间的关系
|
||||||
-spec handle_request(service_pb:'ServiceRequest'(), #state{}) -> {reply, {binary, binary()}, #state{}}.
|
-spec handle_request(efka_service_pb:'ServiceRequest'(), #state{}) -> {reply, {binary, binary()}, #state{}}.
|
||||||
handle_request(#'ServiceRequest'{packet_id = PacketId, request = {register, #'ServiceRequest.Register'{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
|
||||||
@ -113,7 +113,7 @@ handle_request(#'ServiceRequest'{packet_id = PacketId, request = {register, #'Se
|
|||||||
erlang:monitor(process, ServicePid),
|
erlang:monitor(process, ServicePid),
|
||||||
|
|
||||||
%% 更新微服务的状态
|
%% 更新微服务的状态
|
||||||
ok = service_model:insert(#service{
|
ok = efka_service_model:insert(#service{
|
||||||
service_id = ServiceId,
|
service_id = ServiceId,
|
||||||
container_name = <<>>,
|
container_name = <<>>,
|
||||||
status = ?SERVICE_RUNNING,
|
status = ?SERVICE_RUNNING,
|
||||||
@ -125,7 +125,7 @@ handle_request(#'ServiceRequest'{packet_id = PacketId, request = {register, #'Se
|
|||||||
{reply, {binary, result_reply_packet(PacketId, <<"ok">>)},
|
{reply, {binary, result_reply_packet(PacketId, <<"ok">>)},
|
||||||
State#state{service_id = ServiceId, service_pid = ServicePid, is_registered = true}};
|
State#state{service_id = ServiceId, service_pid = ServicePid, is_registered = true}};
|
||||||
{error, Error} ->
|
{error, Error} ->
|
||||||
logger:warning("[service_channel] service_id: ~p, attach_channel get error: ~p", [ServiceId, Error]),
|
logger:warning("[efka_service_channel] service_id: ~p, attach_channel get error: ~p", [ServiceId, Error]),
|
||||||
{reply, {binary, error_reply_packet(PacketId, -1, <<"attach channel failed">>)}, State}
|
{reply, {binary, error_reply_packet(PacketId, -1, <<"attach channel failed">>)}, State}
|
||||||
end;
|
end;
|
||||||
|
|
||||||
@ -144,7 +144,7 @@ handle_request(#'ServiceRequest'{packet_id = PacketId, request = {subscribe, #'S
|
|||||||
handle_request(#'ServiceRequest'{packet_id = PacketId}, State) ->
|
handle_request(#'ServiceRequest'{packet_id = PacketId}, State) ->
|
||||||
{reply, {binary, error_reply_packet(PacketId, -1, <<"invalid request">>)}, State}.
|
{reply, {binary, error_reply_packet(PacketId, -1, <<"invalid request">>)}, State}.
|
||||||
|
|
||||||
-spec handle_cast(service_pb:'ServiceCast'(), #state{}) -> {ok, #state{}}.
|
-spec handle_cast(efka_service_pb:'ServiceCast'(), #state{}) -> {ok, #state{}}.
|
||||||
handle_cast(#'ServiceCast'{body = {metric_data, #'ServiceCast.MetricData'{route_key = RouteKey, metric = Metric}}},
|
handle_cast(#'ServiceCast'{body = {metric_data, #'ServiceCast.MetricData'{route_key = RouteKey, metric = Metric}}},
|
||||||
State = #state{service_pid = ServicePid, is_registered = true}) ->
|
State = #state{service_pid = ServicePid, is_registered = true}) ->
|
||||||
efka_service:metric_data(ServicePid, RouteKey, Metric),
|
efka_service:metric_data(ServicePid, RouteKey, Metric),
|
||||||
@ -154,7 +154,7 @@ handle_cast(#'ServiceCast'{body = _Body}, State) ->
|
|||||||
|
|
||||||
-spec result_reply_packet(integer(), binary()) -> binary().
|
-spec result_reply_packet(integer(), binary()) -> binary().
|
||||||
result_reply_packet(PacketId, Result) when is_integer(PacketId), is_binary(Result) ->
|
result_reply_packet(PacketId, Result) when is_integer(PacketId), is_binary(Result) ->
|
||||||
Reply = service_pb:encode_msg(#'ServiceReply'{
|
Reply = efka_service_pb:encode_msg(#'ServiceReply'{
|
||||||
packet_id = PacketId,
|
packet_id = PacketId,
|
||||||
reply = {result, Result}
|
reply = {result, Result}
|
||||||
}),
|
}),
|
||||||
@ -162,7 +162,7 @@ result_reply_packet(PacketId, Result) when is_integer(PacketId), is_binary(Resul
|
|||||||
|
|
||||||
-spec error_reply_packet(integer(), integer(), binary()) -> 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) ->
|
error_reply_packet(PacketId, Code, Message) when is_integer(PacketId), is_integer(Code), is_binary(Message) ->
|
||||||
Reply = service_pb:encode_msg(#'ServiceReply'{
|
Reply = efka_service_pb:encode_msg(#'ServiceReply'{
|
||||||
packet_id = PacketId,
|
packet_id = PacketId,
|
||||||
reply = {error, #'ServiceReply.Error'{code = Code, message = Message}}
|
reply = {error, #'ServiceReply.Error'{code = Code, message = Message}}
|
||||||
}),
|
}),
|
||||||
@ -6,7 +6,7 @@
|
|||||||
%%% @end
|
%%% @end
|
||||||
%%% Created : 13. 8月 2025 16:41
|
%%% Created : 13. 8月 2025 16:41
|
||||||
%%%-------------------------------------------------------------------
|
%%%-------------------------------------------------------------------
|
||||||
-module(service_model).
|
-module(efka_service_model).
|
||||||
-author("anlicheng").
|
-author("anlicheng").
|
||||||
-include("efka_tables.hrl").
|
-include("efka_tables.hrl").
|
||||||
|
|
||||||
@ -2,7 +2,7 @@
|
|||||||
|
|
||||||
本文档描述部署在 `efka` 主机上的 service 与 `efka` 之间的 WebSocket 通讯逻辑。当前接入实现位于:
|
本文档描述部署在 `efka` 主机上的 service 与 `efka` 之间的 WebSocket 通讯逻辑。当前接入实现位于:
|
||||||
|
|
||||||
- [service_channel.erl](/usr/local/code/cloudkit/efka/apps/efka/src/service/service_channel.erl)
|
- [efka_service_channel.erl](/usr/local/code/cloudkit/efka/apps/efka/src/service/efka_service_channel.erl)
|
||||||
- [service.proto](/usr/local/code/cloudkit/efka/apps/efka/proto/service.proto)
|
- [service.proto](/usr/local/code/cloudkit/efka/apps/efka/proto/service.proto)
|
||||||
- [efka_subscription.erl](/usr/local/code/cloudkit/efka/apps/efka/src/service/efka_subscription.erl)
|
- [efka_subscription.erl](/usr/local/code/cloudkit/efka/apps/efka/src/service/efka_subscription.erl)
|
||||||
- [efka_service.erl](/usr/local/code/cloudkit/efka/apps/efka/src/service/efka_service.erl)
|
- [efka_service.erl](/usr/local/code/cloudkit/efka/apps/efka/src/service/efka_service.erl)
|
||||||
@ -26,11 +26,11 @@ WebSocket 路由固定为:
|
|||||||
GET /ws
|
GET /ws
|
||||||
```
|
```
|
||||||
|
|
||||||
连接建立后,每个 WebSocket 连接对应一个 `service_channel` 进程。
|
连接建立后,每个 WebSocket 连接对应一个 `efka_service_channel` 进程。
|
||||||
|
|
||||||
当前协议没有独立的认证字段。service 必须先发送 `register` request 完成注册;未注册前只能处理 `register`,其它 request 会返回 `invalid request`,cast 会被忽略。
|
当前协议没有独立的认证字段。service 必须先发送 `register` request 完成注册;未注册前只能处理 `register`,其它 request 会返回 `invalid request`,cast 会被忽略。
|
||||||
|
|
||||||
Cowboy WebSocket ping 帧由 `service_channel` 直接回复 pong:
|
Cowboy WebSocket ping 帧由 `efka_service_channel` 直接回复 pong:
|
||||||
|
|
||||||
```erlang
|
```erlang
|
||||||
websocket_handle(ping, State) -> {reply, pong, State}
|
websocket_handle(ping, State) -> {reply, pong, State}
|
||||||
@ -52,7 +52,7 @@ service 与 `efka` 之间使用 WebSocket binary frame。每个业务 frame 的
|
|||||||
| `REPLY` | `0x02` | efka -> service | request 的响应 | `ServiceReply` |
|
| `REPLY` | `0x02` | efka -> service | request 的响应 | `ServiceReply` |
|
||||||
| `CAST` | `0x03` | 双向 | 单向消息,不要求回复 | `ServiceCast` |
|
| `CAST` | `0x03` | 双向 | 单向消息,不要求回复 | `ServiceCast` |
|
||||||
|
|
||||||
`service_channel` 当前只处理 service 发来的 `REQUEST` 和 `CAST`。收到未知格式时只记录日志,不主动发送错误响应。
|
`efka_service_channel` 当前只处理 service 发来的 `REQUEST` 和 `CAST`。收到未知格式时只记录日志,不主动发送错误响应。
|
||||||
|
|
||||||
## 3. Protobuf 定义
|
## 3. Protobuf 定义
|
||||||
|
|
||||||
@ -142,10 +142,10 @@ WebSocket binary frame:
|
|||||||
|
|
||||||
`efka` 处理逻辑:
|
`efka` 处理逻辑:
|
||||||
|
|
||||||
1. `service_channel` 解码 `ServiceRequest`。
|
1. `efka_service_channel` 解码 `ServiceRequest`。
|
||||||
2. 调用 `efka_service_sup:start_service(ServiceId)` 启动或复用 service 进程。
|
2. 调用 `efka_service_sup:start_service(ServiceId)` 启动或复用 service 进程。
|
||||||
3. 调用 `efka_service:attach_channel(ServicePid, ChannelPid)` 绑定当前 WebSocket channel。
|
3. 调用 `efka_service:attach_channel(ServicePid, ChannelPid)` 绑定当前 WebSocket channel。
|
||||||
4. 通过 `service_model:insert/1` 写入 service 记录:
|
4. 通过 `efka_service_model:insert/1` 写入 service 记录:
|
||||||
|
|
||||||
```erlang
|
```erlang
|
||||||
#service{
|
#service{
|
||||||
@ -158,9 +158,9 @@ WebSocket binary frame:
|
|||||||
}
|
}
|
||||||
```
|
```
|
||||||
|
|
||||||
如果该 `service_id` 是新记录,`status` 会写入 `?SERVICE_RUNNING`。如果该 `service_id` 已存在,当前 `service_model:insert/1` 只更新 `meta_data`、`container_name` 和 `update_ts`,不会覆盖已有 `status`。
|
如果该 `service_id` 是新记录,`status` 会写入 `?SERVICE_RUNNING`。如果该 `service_id` 已存在,当前 `efka_service_model:insert/1` 只更新 `meta_data`、`container_name` 和 `update_ts`,不会覆盖已有 `status`。
|
||||||
|
|
||||||
5. `service_channel` 状态更新为:
|
5. `efka_service_channel` 状态更新为:
|
||||||
|
|
||||||
```erlang
|
```erlang
|
||||||
#state{
|
#state{
|
||||||
@ -219,7 +219,7 @@ ServiceRequest {
|
|||||||
|
|
||||||
1. 要求当前 channel 已注册,即 `is_registered = true`。
|
1. 要求当前 channel 已注册,即 `is_registered = true`。
|
||||||
2. 调用 `efka_subscription:subscribe(Topic, ChannelPid)`。
|
2. 调用 `efka_subscription:subscribe(Topic, ChannelPid)`。
|
||||||
3. 订阅成功后,把 topic 加入 `service_channel` 的 `subscribed_topics` 集合。
|
3. 订阅成功后,把 topic 加入 `efka_service_channel` 的 `subscribed_topics` 集合。
|
||||||
|
|
||||||
成功响应:
|
成功响应:
|
||||||
|
|
||||||
@ -277,7 +277,7 @@ WebSocket binary frame:
|
|||||||
|
|
||||||
`efka` 处理逻辑:
|
`efka` 处理逻辑:
|
||||||
|
|
||||||
1. `service_channel` 解码 `ServiceCast`。
|
1. `efka_service_channel` 解码 `ServiceCast`。
|
||||||
2. 要求当前 channel 已注册。
|
2. 要求当前 channel 已注册。
|
||||||
3. 调用:
|
3. 调用:
|
||||||
|
|
||||||
@ -310,13 +310,13 @@ efka_client:metric_data(RouteKey, Metric)
|
|||||||
efka_subscription:publish(Topic, Qos, Content)
|
efka_subscription:publish(Topic, Qos, Content)
|
||||||
```
|
```
|
||||||
|
|
||||||
`efka_subscription` 根据 topic 匹配已订阅的 `service_channel`。匹配成功后向 channel 进程发送:
|
`efka_subscription` 根据 topic 匹配已订阅的 `efka_service_channel`。匹配成功后向 channel 进程发送:
|
||||||
|
|
||||||
```erlang
|
```erlang
|
||||||
{topic_broadcast, Topic, Content}
|
{topic_broadcast, Topic, Content}
|
||||||
```
|
```
|
||||||
|
|
||||||
`service_channel` 再推送 WebSocket `CAST` 给 service:
|
`efka_service_channel` 再推送 WebSocket `CAST` 给 service:
|
||||||
|
|
||||||
```protobuf
|
```protobuf
|
||||||
ServiceCast {
|
ServiceCast {
|
||||||
@ -361,7 +361,7 @@ binary:split(Topic, <<$/>>, [global])
|
|||||||
| `device/+` | `device/a/temp` | 是 |
|
| `device/+` | `device/a/temp` | 是 |
|
||||||
| `device/+` | `device` | 否 |
|
| `device/+` | `device` | 否 |
|
||||||
|
|
||||||
同一个 `service_channel` 订阅多个 topic 时,同一次 publish 最多投递一次。
|
同一个 `efka_service_channel` 订阅多个 topic 时,同一次 publish 最多投递一次。
|
||||||
|
|
||||||
## 7. QoS 与遗留消息
|
## 7. QoS 与遗留消息
|
||||||
|
|
||||||
@ -378,20 +378,20 @@ binary:split(Topic, <<$/>>, [global])
|
|||||||
|
|
||||||
### WebSocket channel 关闭
|
### WebSocket channel 关闭
|
||||||
|
|
||||||
`service_channel:terminate/3` 会执行:
|
`efka_service_channel:terminate/3` 会执行:
|
||||||
|
|
||||||
1. 调用 `efka_subscription:unsubscribe_all(ChannelPid)` 清理该 channel 的所有订阅。
|
1. 调用 `efka_subscription:unsubscribe_all(ChannelPid)` 清理该 channel 的所有订阅。
|
||||||
2. 如果 channel 已完成 register,则调用:
|
2. 如果 channel 已完成 register,则调用:
|
||||||
|
|
||||||
```erlang
|
```erlang
|
||||||
service_model:change_status(ServiceId, ?SERVICE_STOPPED)
|
efka_service_model:change_status(ServiceId, ?SERVICE_STOPPED)
|
||||||
```
|
```
|
||||||
|
|
||||||
也就是把 service 状态更新为停止。
|
也就是把 service 状态更新为停止。
|
||||||
|
|
||||||
### service 进程退出
|
### service 进程退出
|
||||||
|
|
||||||
`service_channel` monitor 绑定的 `efka_service` 进程。如果 service 进程退出,channel 会停止:
|
`efka_service_channel` monitor 绑定的 `efka_service` 进程。如果 service 进程退出,channel 会停止:
|
||||||
|
|
||||||
```erlang
|
```erlang
|
||||||
websocket_info({'DOWN', _Ref, process, ServicePid, Reason}, State) ->
|
websocket_info({'DOWN', _Ref, process, ServicePid, Reason}, State) ->
|
||||||
@ -405,7 +405,7 @@ websocket_info({'DOWN', _Ref, process, ServicePid, Reason}, State) ->
|
|||||||
典型流程:
|
典型流程:
|
||||||
|
|
||||||
```text
|
```text
|
||||||
service efka/service_channel efka_service efka_subscription
|
service efka/efka_service_channel efka_service efka_subscription
|
||||||
| | | |
|
| | | |
|
||||||
|-- WebSocket connect /ws ------------>| | |
|
|-- WebSocket connect /ws ------------>| | |
|
||||||
| | | |
|
| | | |
|
||||||
@ -435,5 +435,5 @@ service efka/service_channel efka_service
|
|||||||
- `register` 可重复连接同一个 `service_id`,但同一时间只允许一个活跃 channel 绑定到 `efka_service`。
|
- `register` 可重复连接同一个 `service_id`,但同一时间只允许一个活跃 channel 绑定到 `efka_service`。
|
||||||
- request 只支持 `register` 和 `subscribe`。
|
- request 只支持 `register` 和 `subscribe`。
|
||||||
- cast 只支持 `metric_data` 和 `topic_event`。
|
- cast 只支持 `metric_data` 和 `topic_event`。
|
||||||
- `subscribed_topics` 目前只在 `service_channel` 内记录,关闭时实际清理由 `efka_subscription:unsubscribe_all/1` 完成。
|
- `subscribed_topics` 目前只在 `efka_service_channel` 内记录,关闭时实际清理由 `efka_subscription:unsubscribe_all/1` 完成。
|
||||||
- `efka_subscription` 当前基于列表扫描匹配,订阅规模很大时需要后续优化为索引结构。
|
- `efka_subscription` 当前基于列表扫描匹配,订阅规模很大时需要后续优化为索引结构。
|
||||||
|
|||||||
@ -36,7 +36,7 @@
|
|||||||
### 中优先级
|
### 中优先级
|
||||||
|
|
||||||
- 继续收敛 `efka_client` 的职责,将传输层状态管理与 request/cast 协议处理进一步拆开。
|
- 继续收敛 `efka_client` 的职责,将传输层状态管理与 request/cast 协议处理进一步拆开。
|
||||||
- 评估并落实 `service_channel` 中 `subscribed_topics` 集合的用途;如果只是被动保存状态,则应简化。
|
- 评估并落实 `efka_service_channel` 中 `subscribed_topics` 集合的用途;如果只是被动保存状态,则应简化。
|
||||||
|
|
||||||
### 低优先级
|
### 低优先级
|
||||||
|
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user