fix iot endpoint

This commit is contained in:
anlicheng 2026-05-11 23:21:19 +08:00
parent 59fa7cf822
commit ae37cbe3d8
13 changed files with 64 additions and 64 deletions

View File

@ -43,7 +43,7 @@ WebSocket server
- `cache_model`DETS 离线缓存。 - `cache_model`DETS 离线缓存。
- `efka_service_model`DETS 服务状态表。 - `efka_service_model`DETS 服务状态表。
- `efka_subscription`:本地 topic 订阅中心。 - `efka_subscription`:本地 topic 订阅中心。
- `efka_client`:连接上游 TLS server 的状态机。 - `efka_iot_client`:连接上游 TLS server 的状态机。
`docker` application 启动的子进程: `docker` application 启动的子进程:
@ -91,12 +91,12 @@ WebSocket server
1. 微服务发送 `ServiceCast.MetricData` 1. 微服务发送 `ServiceCast.MetricData`
2. `efka_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_iot_client:metric_data(RouteKey, Metric)`。
4. `efka_client` 如果处于 activated 状态,直接发给上游;否则写入 `cache_model` 离线缓存。 4. `efka_iot_client` 如果处于 activated 状态,直接发给上游;否则写入 `cache_model` 离线缓存。
### 3.2 EFKA 到上游TLS 长连接 ### 3.2 EFKA 到上游TLS 长连接
上游连接逻辑在 `efka_client`,它是一个 `gen_statem` 上游连接逻辑在 `efka_iot_client`,它是一个 `gen_statem`
状态包括: 状态包括:
@ -122,13 +122,13 @@ WebSocket server
但 protobuf 使用的是 `message_pb`,不是微服务侧的 `efka_service_pb` 但 protobuf 使用的是 `message_pb`,不是微服务侧的 `efka_service_pb`
`efka_client` 上报的内容: `efka_iot_client` 上报的内容:
- `metric_data`业务指标数据activated 时实时发送,否则进入 DETS 缓存。 - `metric_data`业务指标数据activated 时实时发送,否则进入 DETS 缓存。
- `task_event_stream`Docker 部署任务流式日志,只在 activated 时发送。 - `task_event_stream`Docker 部署任务流式日志,只在 activated 时发送。
- `close_task_event_stream`:任务结束事件,只在 activated 时发送。 - `close_task_event_stream`:任务结束事件,只在 activated 时发送。
`efka_client` 接收的内容: `efka_iot_client` 接收的内容:
- `RequestFrame.container_request`:远程容器管理请求,交给 `docker_container_service` - `RequestFrame.container_request`:远程容器管理请求,交给 `docker_container_service`
- `CastFrame.command`:授权命令,`COMMAND_AUTH` 用于切换鉴权/受限状态。 - `CastFrame.command`:授权命令,`COMMAND_AUTH` 用于切换鉴权/受限状态。
@ -187,7 +187,7 @@ topic 匹配规则:
3. 它启动一个独立 `docker_deployer` 进程,并 monitor 该部署进程。 3. 它启动一个独立 `docker_deployer` 进程,并 monitor 该部署进程。
4. `docker_deployer` 执行实际部署步骤。 4. `docker_deployer` 执行实际部署步骤。
5. 部署过程事件写入 `docker_task_reporter` 5. 部署过程事件写入 `docker_task_reporter`
6. `docker_task_reporter` 在上游连接 activated 时把事件转给 `efka_client`;未激活时保留队列并定时重试。 6. `docker_task_reporter` 在上游连接 activated 时把事件转给 `efka_iot_client`;未激活时保留队列并定时重试。
`docker_deployer` 当前部署步骤: `docker_deployer` 当前部署步骤:
@ -246,8 +246,8 @@ topic 匹配规则:
作用: 作用:
- 当 `efka_client` 不在 activated 状态时,把待上报 packet 缓存下来。 - 当 `efka_iot_client` 不在 activated 状态时,把待上报 packet 缓存下来。
- `efka_client` 激活后循环 `fetch_next -> send -> delete` 刷缓存。 - `efka_iot_client` 激活后循环 `fetch_next -> send -> delete` 刷缓存。
缓存 id 使用 `os:system_time(microsecond)` 生成。 缓存 id 使用 `os:system_time(microsecond)` 生成。
@ -279,7 +279,7 @@ topic 匹配规则:
- `README.md` 描述的是 JSON-RPC 风格 WebSocket但当前代码实际处理的是 binary protobuf 帧README 可能已经过期。 - `README.md` 描述的是 JSON-RPC 风格 WebSocket但当前代码实际处理的是 binary protobuf 帧README 可能已经过期。
- `efka_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_iot_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`,但匹配广播时没有使用优先级排序。
- `cache_model` 使用 DETS bag但 id 由微秒时间生成,理论上极端并发下可能碰撞。 - `cache_model` 使用 DETS bag但 id 由微秒时间生成,理论上极端并发下可能碰撞。
@ -288,4 +288,4 @@ topic 匹配规则:
## 10. 一句话主流程 ## 10. 一句话主流程
微服务通过 WebSocket 注册到 EFKA本地 channel 把指标交给对应 `efka_service`,再由 `efka_client` 通过 TLS 上报给上游;上游通过同一条 TLS 连接下发 pub/sub 消息和容器管理请求pub/sub 再广播回本地微服务,容器请求则通过 Docker Unix Socket 在本机执行。 微服务通过 WebSocket 注册到 EFKA本地 channel 把指标交给对应 `efka_service`,再由 `efka_iot_client` 通过 TLS 上报给上游;上游通过同一条 TLS 连接下发 pub/sub 消息和容器管理请求pub/sub 再广播回本地微服务,容器请求则通过 Docker Unix Socket 在本机执行。

View File

@ -68,10 +68,10 @@ update_container_config(ContainerName, Config) when is_binary(ContainerName), is
ConfigFile = get_config_file(ContainerDir), ConfigFile = get_config_file(ContainerDir),
case file:write_file(ConfigFile, Config, [write, binary]) of case file:write_file(ConfigFile, Config, [write, binary]) of
ok -> ok ->
logger:warning("[efka_client] write config file: ~p success", [ConfigFile]), logger:warning("[efka_iot_client] write config file: ~p success", [ConfigFile]),
ok; ok;
{error, Reason} -> {error, Reason} ->
logger:warning("[efka_client] write config file: ~p, get error: ~p", [ConfigFile, Reason]), logger:warning("[efka_iot_client] write config file: ~p, get error: ~p", [ConfigFile, Reason]),
{error, <<"write config failed">>} {error, <<"write config failed">>}
end; end;
error -> error ->

View File

@ -106,13 +106,13 @@ flush_pending(State = #state{pending = Pending0}) ->
-spec send_event({stream, integer(), binary(), binary()} | {close, integer(), binary()}) -> ok | not_ready. -spec send_event({stream, integer(), binary(), binary()} | {close, integer(), binary()}) -> ok | not_ready.
send_event({stream, TaskId, Type, Stream}) -> send_event({stream, TaskId, Type, Stream}) ->
maybe_send(fun() -> efka_client:task_event_stream(TaskId, Type, Stream) end); maybe_send(fun() -> efka_iot_client:task_event_stream(TaskId, Type, Stream) end);
send_event({close, TaskId, Reason}) -> send_event({close, TaskId, Reason}) ->
maybe_send(fun() -> efka_client:close_task_event_stream(TaskId, Reason) end). maybe_send(fun() -> efka_iot_client:close_task_event_stream(TaskId, Reason) end).
-spec maybe_send(fun(() -> any())) -> ok | not_ready. -spec maybe_send(fun(() -> any())) -> ok | not_ready.
maybe_send(SendFun) -> maybe_send(SendFun) ->
case catch efka_client:is_activated() of case catch efka_iot_client:is_activated() of
true -> true ->
_ = catch SendFun(), _ = catch SendFun(),
ok; ok;

View File

@ -67,21 +67,21 @@ init([]) ->
}, },
#{ #{
id => 'efka_client', id => 'efka_iot_client',
start => {'efka_client', start_link, []}, start => {'efka_iot_client', start_link, []},
restart => permanent, restart => permanent,
shutdown => 2000, shutdown => 2000,
type => worker, type => worker,
modules => ['efka_client'] modules => ['efka_iot_client']
}, },
#{ #{
id => 'efka_heartbeat', id => 'efka_iot_heartbeat',
start => {'efka_heartbeat', start_link, []}, start => {'efka_iot_heartbeat', start_link, []},
restart => permanent, restart => permanent,
shutdown => 2000, shutdown => 2000,
type => worker, type => worker,
modules => ['efka_heartbeat'] modules => ['efka_iot_heartbeat']
} }
], ],

View File

@ -2,10 +2,10 @@
%%% @author anlicheng %%% @author anlicheng
%%% @copyright (C) 2026, <COMPANY> %%% @copyright (C) 2026, <COMPANY>
%%% @doc %%% @doc
%%% DETS backed outbound packet cache for efka_client. %%% DETS backed outbound packet cache for efka_iot_client.
%%% @end %%% @end
%%%------------------------------------------------------------------- %%%-------------------------------------------------------------------
-module(efka_client_cache). -module(efka_iot_cache).
-author("anlicheng"). -author("anlicheng").
-define(CACHE_TAB, cache). -define(CACHE_TAB, cache).

View File

@ -6,7 +6,7 @@
%%% @end %%% @end
%%% Created : 20. 4 2026 00:00 %%% Created : 20. 4 2026 00:00
%%%------------------------------------------------------------------- %%%-------------------------------------------------------------------
-module(efka_client). -module(efka_iot_client).
-author("anlicheng"). -author("anlicheng").
-include("efka_tables.hrl"). -include("efka_tables.hrl").
@ -77,7 +77,7 @@ start_link() ->
-spec init(list()) -> {ok, atom(), #state{}}. -spec init(list()) -> {ok, atom(), #state{}}.
init([]) -> init([]) ->
case efka_client_cache:open() of case efka_iot_cache:open() of
ok -> ok ->
erlang:start_timer(0, self(), create_transport), erlang:start_timer(0, self(), create_transport),
{ok, ?STATE_DISCONNECTED, #state{socket = undefined}}; {ok, ?STATE_DISCONNECTED, #state{socket = undefined}};
@ -98,13 +98,13 @@ handle_event(cast, {metric_data, RouteKey, Metric}, StateName, State = #state{so
ok = ssl:send(Socket, Packet), ok = ssl:send(Socket, Packet),
{keep_state, State}; {keep_state, State};
_ -> _ ->
{ok, DroppedCount} = efka_client_cache:insert(Packet), {ok, DroppedCount} = efka_iot_cache:insert(Packet),
{keep_state, State#state{dropped_message_count = State#state.dropped_message_count + DroppedCount}} {keep_state, State#state{dropped_message_count = State#state.dropped_message_count + DroppedCount}}
end; end;
%% Task的stream流 %% Task的stream流
handle_event(cast, {task_event_stream, TaskId, Type, Stream}, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> handle_event(cast, {task_event_stream, TaskId, Type, Stream}, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
logger:debug("[efka_client] event_stream task_id: ~p, stream: ~ts", [TaskId, Stream]), logger:debug("[efka_iot_client] event_stream task_id: ~p, stream: ~ts", [TaskId, Stream]),
Packet = term_to_binary({<<"message">>, {<<"task_event">>, #{<<"task_id">> => TaskId, <<"type">> => Type, <<"stream">> => Stream}}}), Packet = term_to_binary({<<"message">>, {<<"task_event">>, #{<<"task_id">> => TaskId, <<"type">> => Type, <<"stream">> => Stream}}}),
ok = ssl:send(Socket, Packet), ok = ssl:send(Socket, Packet),
{keep_state, State}; {keep_state, State};
@ -132,7 +132,7 @@ handle_event(info, {timeout, _, create_transport}, ?STATE_DISCONNECTED, State) -
Ref = request_ref(), Ref = request_ref(),
AuthPacket = auth_packet(Ref), AuthPacket = auth_packet(Ref),
ok = ssl:send(Socket, AuthPacket), ok = ssl:send(Socket, AuthPacket),
logger:debug("[efka_client] send auth request, ref: ~p", [Ref]), logger:debug("[efka_iot_client] send auth request, ref: ~p", [Ref]),
{next_state, ?STATE_AUTH, State#state{socket = Socket, auth_ref = Ref}, [{state_timeout, 5000, auth_timeout}]}; {next_state, ?STATE_AUTH, State#state{socket = Socket, auth_ref = Ref}, [{state_timeout, 5000, auth_timeout}]};
{error, _Reason} -> {error, _Reason} ->
schedule_reconnect(), schedule_reconnect(),
@ -140,7 +140,7 @@ handle_event(info, {timeout, _, create_transport}, ?STATE_DISCONNECTED, State) -
end; end;
handle_event(state_timeout, auth_timeout, ?STATE_AUTH, State = #state{socket = Socket}) -> handle_event(state_timeout, auth_timeout, ?STATE_AUTH, State = #state{socket = Socket}) ->
logger:debug("[efka_client] auth request timeout"), logger:debug("[efka_iot_client] auth request timeout"),
disconnect(Socket), disconnect(Socket),
schedule_reconnect(), schedule_reconnect(),
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined}}; {next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined}};
@ -151,7 +151,7 @@ handle_event(info, {timeout, TimerRef, ssl_ping}, ?STATE_ACTIVATED, State = #sta
ok -> ok ->
{keep_state, schedule_ssl_ping(State)}; {keep_state, schedule_ssl_ping(State)};
{error, Reason} -> {error, Reason} ->
logger:warning("[efka_client] send ssl ping failed, reason: ~p", [Reason]), logger:warning("[efka_iot_client] send ssl ping failed, reason: ~p", [Reason]),
disconnect(Socket), disconnect(Socket),
schedule_reconnect(), schedule_reconnect(),
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}} {next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}}
@ -161,10 +161,10 @@ handle_event(info, {timeout, _TimerRef, ssl_ping}, _StateName, State) ->
%% %%
handle_event(info, flush_cache, ?STATE_ACTIVATED, State = #state{socket = Socket}) -> handle_event(info, flush_cache, ?STATE_ACTIVATED, State = #state{socket = Socket}) ->
case efka_client_cache:fetch_next() of case efka_iot_cache:fetch_next() of
{ok, {Id, Packet}} -> {ok, {Id, Packet}} ->
ok = ssl:send(Socket, Packet), ok = ssl:send(Socket, Packet),
ok = efka_client_cache:delete(Id), ok = efka_iot_cache:delete(Id),
{keep_state, State, [{next_event, info, flush_cache}]}; {keep_state, State, [{next_event, info, flush_cache}]};
error -> error ->
{keep_state, State} {keep_state, State}
@ -179,20 +179,20 @@ handle_event(info, {ssl, Socket, PacketBin}, _, State = #state{socket = Socket})
{keep_state, State, [{next_event, internal, Packet}]} {keep_state, State, [{next_event, internal, Packet}]}
catch catch
error:Error -> error:Error ->
logger:warning("[efka_client] binary_to_term get error: ~p, packet_size: ~p", [Error, byte_size(PacketBin)]), logger:warning("[efka_iot_client] binary_to_term get error: ~p, packet_size: ~p", [Error, byte_size(PacketBin)]),
disconnect(Socket), disconnect(Socket),
cancel_ssl_ping(State), cancel_ssl_ping(State),
schedule_reconnect(), schedule_reconnect(),
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}} {next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}}
end; end;
handle_event(info, {ssl_error, Socket, Reason}, _, State = #state{}) -> handle_event(info, {ssl_error, Socket, Reason}, _, State = #state{}) ->
logger:debug("[efka_client] ssl error: ~p", [Reason]), logger:debug("[efka_iot_client] ssl error: ~p", [Reason]),
disconnect(Socket), disconnect(Socket),
cancel_ssl_ping(State), cancel_ssl_ping(State),
schedule_reconnect(), schedule_reconnect(),
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}}; {next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}};
handle_event(info, {ssl_closed, Socket}, _, State = #state{}) -> handle_event(info, {ssl_closed, Socket}, _, State = #state{}) ->
logger:debug("[efka_client] ssl closed"), logger:debug("[efka_iot_client] ssl closed"),
disconnect(Socket), disconnect(Socket),
cancel_ssl_ping(State), cancel_ssl_ping(State),
schedule_reconnect(), schedule_reconnect(),
@ -205,40 +205,40 @@ handle_event(internal, {<<"command">>, Ref, {<<"container">>, Request}}, ?STATE_
handle_container_command(Ref, Request, Socket), handle_container_command(Ref, Request, Socket),
{keep_state, State}; {keep_state, State};
handle_event(internal, {<<"command">>, Ref, {<<"container">>, Request}}, _StateName, State = #state{socket = Socket}) -> handle_event(internal, {<<"command">>, Ref, {<<"container">>, Request}}, _StateName, State = #state{socket = Socket}) ->
logger:notice("[efka_client] get an invalid command: ~p, agent invalid", [Request]), logger:notice("[efka_iot_client] get an invalid command: ~p, agent invalid", [Request]),
send_container_response(Socket, Ref, {error, <<"agent invalid">>}), send_container_response(Socket, Ref, {error, <<"agent invalid">>}),
{keep_state, State}; {keep_state, State};
%% response %% response
handle_event(internal, {<<"response">>, AuthRef, {<<"auth_response">>, <<"ok">>}}, ?STATE_AUTH, State = #state{auth_ref = AuthRef}) -> handle_event(internal, {<<"response">>, AuthRef, {<<"auth_response">>, <<"ok">>}}, ?STATE_AUTH, State = #state{auth_ref = AuthRef}) ->
logger:debug("[efka_client] auth success"), logger:debug("[efka_iot_client] auth success"),
State1 = schedule_ssl_ping(State#state{auth_ref = undefined}), State1 = schedule_ssl_ping(State#state{auth_ref = undefined}),
{next_state, ?STATE_ACTIVATED, State1, [{next_event, info, flush_cache}]}; {next_state, ?STATE_ACTIVATED, State1, [{next_event, info, flush_cache}]};
handle_event(internal, {<<"response">>, AuthRef, {<<"auth_response">>, {<<"error">>, Reason}}}, ?STATE_AUTH, State = #state{socket = Socket, auth_ref = AuthRef}) -> handle_event(internal, {<<"response">>, AuthRef, {<<"auth_response">>, {<<"error">>, Reason}}}, ?STATE_AUTH, State = #state{socket = Socket, auth_ref = AuthRef}) ->
logger:debug("[efka_client] auth failed, reason: ~p", [Reason]), logger:debug("[efka_iot_client] auth failed, reason: ~p", [Reason]),
disconnect(Socket), disconnect(Socket),
schedule_reconnect(), schedule_reconnect(),
{next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined}}; {next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined}};
handle_event(internal, {<<"response">>, _Ref, Reply}, StateName, State) -> handle_event(internal, {<<"response">>, _Ref, Reply}, StateName, State) ->
logger:warning("[efka_client] ignore unexpected response in state ~p: ~p", [StateName, Reply]), logger:warning("[efka_iot_client] ignore unexpected response in state ~p: ~p", [StateName, Reply]),
{keep_state, State}; {keep_state, State};
handle_event(internal, {<<"command_response">>, _Ref, Reply}, StateName, State) -> handle_event(internal, {<<"command_response">>, _Ref, Reply}, StateName, State) ->
logger:warning("[efka_client] ignore unexpected command_response in state ~p: ~p", [StateName, Reply]), logger:warning("[efka_iot_client] ignore unexpected command_response in state ~p: ~p", [StateName, Reply]),
{keep_state, State}; {keep_state, State};
%% Pub/Sub机制 %% Pub/Sub机制
handle_event(internal, {<<"message">>, <<"pong">>}, ?STATE_ACTIVATED, State) -> handle_event(internal, {<<"message">>, <<"pong">>}, ?STATE_ACTIVATED, State) ->
{keep_state, State}; {keep_state, State};
handle_event(internal, {<<"message">>, {<<"pub">>, #{<<"topic">> := Topic, <<"qos">> := Qos, <<"content">> := Content}}}, ?STATE_ACTIVATED, State) -> handle_event(internal, {<<"message">>, {<<"pub">>, #{<<"topic">> := Topic, <<"qos">> := Qos, <<"content">> := Content}}}, ?STATE_ACTIVATED, State) ->
logger:debug("[efka_client] get pub topic: ~p, qos: ~p, content: ~p", [Topic, Qos, Content]), logger:debug("[efka_iot_client] get pub topic: ~p, qos: ~p, content: ~p", [Topic, Qos, Content]),
efka_subscription:publish(Topic, Qos, Content), efka_subscription:publish(Topic, Qos, Content),
{keep_state, State}; {keep_state, State};
handle_event(internal, Packet, _StateName, State) -> handle_event(internal, Packet, _StateName, State) ->
logger:warning("[efka_client] ignore unknown packet: ~p", [Packet]), logger:warning("[efka_iot_client] ignore unknown packet: ~p", [Packet]),
{keep_state, State}; {keep_state, State};
handle_event(info, Info, _, State = #state{}) -> handle_event(info, Info, _, State = #state{}) ->
logger:notice("[efka_client] get unknown info: ~p", [Info]), logger:notice("[efka_iot_client] get unknown info: ~p", [Info]),
{keep_state, State}. {keep_state, State}.
-spec handle_container_command(binary(), term(), ssl:sslsocket()) -> ok. -spec handle_container_command(binary(), term(), ssl:sslsocket()) -> ok.
@ -271,7 +271,7 @@ handle_container_command(Ref, #{<<"action">> := <<"config">>, <<"target">> := Ta
send_container_response(Socket, Ref, Reply), send_container_response(Socket, Ref, Reply),
ok; ok;
handle_container_command(Ref, Request, Socket) -> handle_container_command(Ref, Request, Socket) ->
logger:notice("[efka_client] get an invalid command: ~p, agent invalid", [Request]), logger:notice("[efka_iot_client] get an invalid command: ~p, agent invalid", [Request]),
send_container_response(Socket, Ref, {error, <<"agent invalid">>}), send_container_response(Socket, Ref, {error, <<"agent invalid">>}),
ok. ok.
@ -279,8 +279,8 @@ handle_container_command(Ref, Request, Socket) ->
terminate(Reason, _StateName, State = #state{socket = Socket}) -> terminate(Reason, _StateName, State = #state{socket = Socket}) ->
cancel_ssl_ping(State), cancel_ssl_ping(State),
disconnect(Socket), disconnect(Socket),
efka_client_cache:close(), efka_iot_cache:close(),
logger:notice("[efka_client] terminate with reason: ~p", [Reason]), logger:notice("[efka_iot_client] terminate with reason: ~p", [Reason]),
ok. ok.
-spec code_change(term(), atom(), #state{}, term()) -> {ok, atom(), #state{}}. -spec code_change(term(), atom(), #state{}, term()) -> {ok, atom(), #state{}}.

View File

@ -2,7 +2,7 @@
%%% @doc UDP heartbeat sender for iot host liveness. %%% @doc UDP heartbeat sender for iot host liveness.
%%% @end %%% @end
%%%------------------------------------------------------------------- %%%-------------------------------------------------------------------
-module(efka_heartbeat). -module(efka_iot_heartbeat).
-behaviour(gen_server). -behaviour(gen_server).
@ -82,7 +82,7 @@ handle_info(heartbeat, State = #state{interval = Interval}) ->
erlang:send_after(Interval, self(), heartbeat), erlang:send_after(Interval, self(), heartbeat),
{noreply, State}; {noreply, State};
handle_info(Info, State) -> handle_info(Info, State) ->
logger:warning("[efka_heartbeat] ignore unknown info: ~p", [Info]), logger:warning("[efka_iot_heartbeat] ignore unknown info: ~p", [Info]),
{noreply, State}. {noreply, State}.
-spec terminate(term(), #state{}) -> ok. -spec terminate(term(), #state{}) -> ok.
@ -105,7 +105,7 @@ send_heartbeat(#state{socket = Socket, host = Host, port = Port, uuid = UUID, he
ok -> ok ->
ok; ok;
{error, Reason} -> {error, Reason} ->
logger:warning("[efka_heartbeat] send heartbeat failed, reason: ~p", [Reason]), logger:warning("[efka_iot_heartbeat] send heartbeat failed, reason: ~p", [Reason]),
ok ok
end. end.

View File

@ -104,7 +104,7 @@ handle_call(_Request, _From, State = #state{}) ->
{stop, Reason :: term(), NewState :: #state{}}). {stop, Reason :: term(), NewState :: #state{}}).
handle_cast({metric_data, RouteKey, Metric}, State = #state{service_id = ServiceId}) -> handle_cast({metric_data, RouteKey, Metric}, State = #state{service_id = ServiceId}) ->
logger:debug("[efka_service] metric_data service_id: ~p, route_key: ~p, metric data: ~p", [ServiceId, RouteKey, Metric]), logger:debug("[efka_service] metric_data service_id: ~p, route_key: ~p, metric data: ~p", [ServiceId, RouteKey, Metric]),
efka_client:metric_data(RouteKey, Metric), efka_iot_client:metric_data(RouteKey, Metric),
{noreply, State}; {noreply, State};
handle_cast(_Request, State = #state{}) -> handle_cast(_Request, State = #state{}) ->

View File

@ -4,7 +4,7 @@
对应代码: 对应代码:
- command 接收入口:[src/transport/efka_client.erl](/usr/local/code/cloudkit/efka/apps/efka/src/transport/efka_client.erl:177) - command 接收入口:[src/iot/efka_iot_client.erl](/usr/local/code/cloudkit/efka/apps/efka/src/iot/efka_iot_client.erl:177)
- 部署任务管理:[docker_deploy_manager.erl](/usr/local/code/cloudkit/efka/apps/docker/src/docker_deploy_manager.erl:36) - 部署任务管理:[docker_deploy_manager.erl](/usr/local/code/cloudkit/efka/apps/docker/src/docker_deploy_manager.erl:36)
- 部署执行:[docker_deployer.erl](/usr/local/code/cloudkit/efka/apps/docker/src/docker_deployer.erl:39) - 部署执行:[docker_deployer.erl](/usr/local/code/cloudkit/efka/apps/docker/src/docker_deployer.erl:39)
- Docker JSON 构造:[docker_container_builder.erl](/usr/local/code/cloudkit/efka/apps/docker/src/docker_container_builder.erl:14) - Docker JSON 构造:[docker_container_builder.erl](/usr/local/code/cloudkit/efka/apps/docker/src/docker_container_builder.erl:14)
@ -34,11 +34,11 @@
`Ref``crypto:strong_rand_bytes(16)` 生成的 16 字节 binary。网络帧只使用 `binary_to_term(PacketBin, [safe])` 可解码的 safe term协议 label、业务 label、map key 和 action 都使用 binary。 `Ref``crypto:strong_rand_bytes(16)` 生成的 16 字节 binary。网络帧只使用 `binary_to_term(PacketBin, [safe])` 可解码的 safe term协议 label、业务 label、map key 和 action 都使用 binary。
只有 `efka_client` 处于 `activated` 状态时,容器命令才会正常执行;处于非 activated 状态时会返回错误。 只有 `efka_iot_client` 处于 `activated` 状态时,容器命令才会正常执行;处于非 activated 状态时会返回错误。
## 2. 容器命令 map ## 2. 容器命令 map
`efka_client` 收到网络协议里的 binary-key `CommandMap` 后,直接按 binary key 和 binary action 做函数参数匹配,不再做整包 atom-key 转换。 `efka_iot_client` 收到网络协议里的 binary-key `CommandMap` 后,直接按 binary key 和 binary action 做函数参数匹配,不再做整包 atom-key 转换。
- command 顶层、`target`、deploy `params``create` 中间结构都使用 binary key。 - command 顶层、`target`、deploy `params``create` 中间结构都使用 binary key。
- `<<"action">>` 的取值使用 binary`<<"list">>``<<"deploy">>``<<"start">>``<<"stop">>``<<"kill">>``<<"remove">>``<<"config">>` - `<<"action">>` 的取值使用 binary`<<"list">>``<<"deploy">>``<<"start">>``<<"stop">>``<<"kill">>``<<"remove">>``<<"config">>`

View File

@ -257,7 +257,7 @@ GET /event_stream?uuid=<host_uuid>&task_id=<task_id>
## UDP 心跳 ## UDP 心跳
`efka` 通过独立的 `efka_heartbeat` 进程向 `iot` 发送 UDP 心跳。TLS control channel 和 UDP 心跳共用 `iot_server.host`,分别使用 `tls_port``udp_port`。UDP 心跳包使用 HMAC-SHA256 校验HMAC key 为 `SHA256(auth.token)` `efka` 通过独立的 `efka_iot_heartbeat` 进程向 `iot` 发送 UDP 心跳。TLS control channel 和 UDP 心跳共用 `iot_server.host`,分别使用 `tls_port``udp_port`。UDP 心跳包使用 HMAC-SHA256 校验HMAC key 为 `SHA256(auth.token)`
详细格式见 [heartbeat.md](heartbeat.md)。 详细格式见 [heartbeat.md](heartbeat.md)。

View File

@ -1,6 +1,6 @@
# UDP 心跳 # UDP 心跳
`efka` 通过独立的 `efka_heartbeat` 进程向 `iot` 发送 UDP 心跳。该心跳只表示主机存活,不依赖 TLS control channel 是否在线。 `efka` 通过独立的 `efka_iot_heartbeat` 进程向 `iot` 发送 UDP 心跳。该心跳只表示主机存活,不依赖 TLS control channel 是否在线。
## 配置 ## 配置
@ -23,8 +23,8 @@ TLS 和 UDP 共用同一个 `iot_server.host`
]} ]}
``` ```
- `tls_port` 用于 `efka_client` 建立 TLS 连接。 - `tls_port` 用于 `efka_iot_client` 建立 TLS 连接。
- `udp_port` 用于 `efka_heartbeat` 发送 UDP 心跳。 - `udp_port` 用于 `efka_iot_heartbeat` 发送 UDP 心跳。
- `heartbeat.interval` 是发送间隔,单位毫秒,默认 5000。 - `heartbeat.interval` 是发送间隔,单位毫秒,默认 5000。
- UDP 心跳和 TLS 鉴权复用 `auth.uuid``auth.token` - UDP 心跳和 TLS 鉴权复用 `auth.uuid``auth.token`

View File

@ -288,10 +288,10 @@ efka_service:metric_data(ServicePid, RouteKey, Metric)
4. `efka_service` 再调用: 4. `efka_service` 再调用:
```erlang ```erlang
efka_client:metric_data(RouteKey, Metric) efka_iot_client:metric_data(RouteKey, Metric)
``` ```
5. `efka_client` 通过 efka 与 iot 的 TLS 长连接把数据上报给 iot 5. `efka_iot_client` 通过 efka 与 iot 的 TLS 长连接把数据上报给 iot
```erlang ```erlang
{<<"message">>, {<<"data">>, #{ {<<"message">>, {<<"data">>, #{
@ -304,7 +304,7 @@ efka_client:metric_data(RouteKey, Metric)
### 5.2 efka -> service: topic_event ### 5.2 efka -> service: topic_event
`iot` 通过 efka 与 iot 的 TLS channel 下发 pub 消息到 `efka` 时,`efka_client` 会调用: `iot` 通过 efka 与 iot 的 TLS channel 下发 pub 消息到 `efka` 时,`efka_iot_client` 会调用:
```erlang ```erlang
efka_subscription:publish(Topic, Qos, Content) efka_subscription:publish(Topic, Qos, Content)
@ -420,7 +420,7 @@ service efka/efka_service_channel efka_serv
| | | | | | | |
|-- CAST metric_data ----------------->| | | |-- CAST metric_data ----------------->| | |
| |-- metric_data -------------->| | | |-- metric_data -------------->| |
| | |-- efka_client:metric_data -> iot | | |-- efka_iot_client:metric_data -> iot
| | | | | | | |
|<------------- CAST topic_event ------|<--------------- topic_broadcast -----------------------| |<------------- CAST topic_event ------|<--------------- topic_broadcast -----------------------|
| | | | | | | |

View File

@ -35,12 +35,12 @@
### 中优先级 ### 中优先级
- 继续收敛 `efka_client` 的职责,将传输层状态管理与 request/cast 协议处理进一步拆开。 - 继续收敛 `efka_iot_client` 的职责,将传输层状态管理与 request/cast 协议处理进一步拆开。
- 评估并落实 `efka_service_channel``subscribed_topics` 集合的用途;如果只是被动保存状态,则应简化。 - 评估并落实 `efka_service_channel``subscribed_topics` 集合的用途;如果只是被动保存状态,则应简化。
### 低优先级 ### 低优先级
- 如果 `efka_client` 的缓存指标量继续增长,可以为缓存刷出增加批量发送或节流机制。 - 如果 `efka_iot_client` 的缓存指标量继续增长,可以为缓存刷出增加批量发送或节流机制。
## 基础设施层 ## 基础设施层