From 8f801fd29ca0d910909b1a3c43c9e91c8f8c82f7 Mon Sep 17 00:00:00 2001 From: anlicheng <244108715@qq.com> Date: Mon, 6 Jul 2026 13:47:30 +0800 Subject: [PATCH] support stream --- apps/efka/src/iot/efka_iot_client.erl | 215 +++++++++++++++++++++----- apps/efka/src/iot/efka_iot_stream.erl | 150 ++++++++++++++++++ config/sys.config.src | 9 ++ docs/efka_iot_protocol.md | 46 ++++++ 4 files changed, 382 insertions(+), 38 deletions(-) create mode 100644 apps/efka/src/iot/efka_iot_stream.erl diff --git a/apps/efka/src/iot/efka_iot_client.erl b/apps/efka/src/iot/efka_iot_client.erl index eb727f6..565632b 100644 --- a/apps/efka/src/iot/efka_iot_client.erl +++ b/apps/efka/src/iot/efka_iot_client.erl @@ -15,6 +15,7 @@ %% API -export([start_link/0]). -export([metric_data/2, ping/13, task_event_stream/3, close_task_event_stream/2]). +-export([send_stream/2, stream_done/1]). -export([is_activated/0, dropped_message_count/0]). %% gen_statem callbacks @@ -36,12 +37,20 @@ -record(state, { socket :: undefined | ssl:sslsocket(), outbox :: efka_iot_outbox:outbox(), + streams = #{}, %% 保存当前auth请求的ref,用来建立auth请求和响应的对应关系 auth_ref = undefined :: undefined | binary(), ping_timer_ref = undefined :: undefined | reference(), dropped_message_count = 0 :: non_neg_integer() }). +-record(stream_state, { + worker_pid :: pid(), + monitor_ref :: reference() +}). + +-type stream_id() :: pos_integer(). + %%%=================================================================== %%% API %%%=================================================================== @@ -59,6 +68,14 @@ task_event_stream(TaskId, Type, Stream) when is_integer(TaskId), is_binary(Type) close_task_event_stream(TaskId, Reason) when is_integer(TaskId), is_binary(Reason) -> gen_statem:cast(?SERVER, {close_task_event_stream, TaskId, Reason}). +-spec send_stream(StreamId :: stream_id(), Body :: term()) -> ok. +send_stream(StreamId, Body) when is_integer(StreamId), StreamId > 0 -> + gen_statem:cast(?SERVER, {send_stream, StreamId, Body}). + +-spec stream_done(StreamId :: stream_id()) -> ok. +stream_done(StreamId) when is_integer(StreamId), StreamId > 0 -> + gen_statem:cast(?SERVER, {stream_done, StreamId}). + -spec is_activated() -> boolean(). is_activated() -> gen_statem:call(?SERVER, is_activated). @@ -139,6 +156,24 @@ handle_event(cast, {close_task_event_stream, TaskId, Reason}, ?STATE_ACTIVATED, ok = ssl:send(Socket, Packet), {keep_state, State}; +handle_event(cast, {send_stream, StreamId, Body}, ?STATE_ACTIVATED, State = #state{socket = Socket}) + when is_integer(StreamId), StreamId > 0 -> + Packet = term_to_binary({<<"stream">>, StreamId, Body}), + ok = ssl:send(Socket, Packet), + {keep_state, State}; +handle_event(cast, {send_stream, _StreamId, _Body}, _StateName, State) -> + {keep_state, State}; + +handle_event(cast, {stream_done, StreamId}, _StateName, State = #state{streams = Streams}) + when is_integer(StreamId), StreamId > 0 -> + case maps:take(StreamId, Streams) of + error -> + {keep_state, State}; + {StreamState, NStreams} -> + demonitor_stream(StreamState), + {keep_state, State#state{streams = NStreams}} + end; + %% 其他情况下直接忽略 handle_event(cast, _, _, State = #state{}) -> {keep_state, State}; @@ -179,7 +214,8 @@ handle_event(info, {timeout, TimerRef, ssl_ping}, ?STATE_ACTIVATED, State = #sta logger:warning("[efka_iot_client] send ssl ping failed, reason: ~p", [Reason]), disconnect(Socket), schedule_reconnect(), - {next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}} + NState = close_all_streams({send_ping_failed, Reason}, State), + {next_state, ?STATE_DISCONNECTED, NState#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}} end; handle_event(info, {timeout, _TimerRef, ssl_ping}, _StateName, State) -> {keep_state, State}; @@ -216,20 +252,40 @@ handle_event(info, {ssl, Socket, PacketBin}, _, State = #state{socket = Socket}) disconnect(Socket), cancel_ssl_ping(State), schedule_reconnect(), - {next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}} + NState = close_all_streams(bad_packet, State), + {next_state, ?STATE_DISCONNECTED, NState#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}} end; handle_event(info, {ssl_error, Socket, Reason}, _, State = #state{}) -> logger:debug("[efka_iot_client] ssl error: ~p", [Reason]), disconnect(Socket), cancel_ssl_ping(State), schedule_reconnect(), - {next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}}; + NState = close_all_streams(ssl_error, State), + {next_state, ?STATE_DISCONNECTED, NState#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}}; handle_event(info, {ssl_closed, Socket}, _, State = #state{}) -> logger:debug("[efka_iot_client] ssl closed"), disconnect(Socket), cancel_ssl_ping(State), schedule_reconnect(), - {next_state, ?STATE_DISCONNECTED, State#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}}; + NState = close_all_streams(ssl_closed, State), + {next_state, ?STATE_DISCONNECTED, NState#state{socket = undefined, auth_ref = undefined, ping_timer_ref = undefined}}; + +handle_event(info, {'DOWN', MonitorRef, process, WorkerPid, Reason}, _StateName, State = #state{streams = Streams}) -> + case take_stream_by_monitor(MonitorRef, WorkerPid, Streams) of + error -> + {keep_state, State}; + {StreamId, StreamState, NStreams} -> + demonitor_stream(StreamState), + logger:warning("[efka_iot_client] stream worker down, stream_id: ~p, reason: ~p", [StreamId, Reason]), + case State#state.socket of + undefined -> + {keep_state, State#state{streams = NStreams}}; + Socket -> + Packet = term_to_binary({<<"stream">>, StreamId, {<<"reset">>, safe_term({worker_down, Reason})}}), + ok = ssl:send(Socket, Packet), + {keep_state, State#state{streams = NStreams}} + end + end; %%% 处理内部消息,ssl收到的消息会先 binary_to_term,再由这里按协议结构模式匹配 @@ -259,6 +315,14 @@ handle_event(internal, {<<"command_response">>, _Ref, Reply}, StateName, State) logger:warning("[efka_iot_client] ignore unexpected command_response in state ~p: ~p", [StateName, Reply]), {keep_state, State}; +%% 透明 TCP stream 多路复用,StreamId 只出现在新顶层 <<"stream">> 帧里。 +handle_event(internal, {<<"stream">>, StreamId, Body}, ?STATE_ACTIVATED, State) -> + handle_stream_frame(StreamId, Body, State); +handle_event(internal, {<<"stream">>, StreamId, Body}, StateName, State) -> + logger:warning("[efka_iot_client] ignore stream frame in state ~p, stream_id: ~p, body: ~p", + [StateName, StreamId, Body]), + {keep_state, State}; + %% 处理Pub/Sub机制 handle_event(internal, {<<"message">>, <<"pong">>}, ?STATE_ACTIVATED, State) -> {keep_state, State}; @@ -274,43 +338,10 @@ handle_event(info, Info, _, State = #state{}) -> logger:notice("[efka_iot_client] get unknown info: ~p", [Info]), {keep_state, State}. --spec handle_container_command(binary(), term(), ssl:sslsocket()) -> ok. -handle_container_command(Ref, #{<<"action">> := <<"list">>}, Socket) -> - Reply = docker_commands:get_containers(), - send_container_response(Socket, Ref, Reply), - ok; -handle_container_command(Ref, #{<<"action">> := <<"deploy">>, <<"task_id">> := TaskId, <<"params">> := Params}, Socket) -> - Reply = docker_deploy_manager:deploy(TaskId, Params), - send_container_response(Socket, Ref, Reply), - ok; -handle_container_command(Ref, #{<<"action">> := <<"start">>, <<"target">> := Target}, Socket) -> - Reply = docker_commands:start_container(container_target(Target)), - send_container_response(Socket, Ref, Reply), - ok; -handle_container_command(Ref, #{<<"action">> := <<"stop">>, <<"target">> := Target, <<"timeout_seconds">> := TimeoutSeconds}, Socket) -> - Reply = docker_commands:stop_container(container_target(Target), TimeoutSeconds), - send_container_response(Socket, Ref, Reply), - ok; -handle_container_command(Ref, #{<<"action">> := <<"kill">>, <<"target">> := Target, <<"signal">> := Signal}, Socket) -> - Reply = docker_commands:kill_container(container_target(Target), to_binary(Signal)), - send_container_response(Socket, Ref, Reply), - ok; -handle_container_command(Ref, #{<<"action">> := <<"remove">>, <<"target">> := Target, <<"force">> := Force, <<"remove_volumes">> := RemoveVolumes}, Socket) -> - Reply = docker_commands:remove_container(container_target(Target), to_bool(Force), to_bool(RemoveVolumes)), - send_container_response(Socket, Ref, Reply), - ok; -handle_container_command(Ref, #{<<"action">> := <<"config">>, <<"target">> := Target, <<"config">> := Config}, Socket) -> - Reply = docker_helper:update_container_config(container_target(Target), iolist_to_binary(Config)), - send_container_response(Socket, Ref, Reply), - ok; -handle_container_command(Ref, Request, Socket) -> - logger:notice("[efka_iot_client] get an invalid command: ~p, agent invalid", [Request]), - send_container_response(Socket, Ref, {error, <<"agent invalid">>}), - ok. - -spec terminate(term(), atom(), #state{}) -> ok. terminate(Reason, _StateName, State = #state{socket = Socket, outbox = Outbox}) -> cancel_ssl_ping(State), + _ = close_all_streams(Reason, State), disconnect(Socket), efka_iot_outbox:close(Outbox), logger:notice("[efka_iot_client] terminate with reason: ~p", [Reason]), @@ -386,6 +417,114 @@ send_container_response(Socket, Ref, Reply) -> Packet = term_to_binary({<<"command_response">>, Ref, {<<"container">>, safe_reply(Reply)}}), ok = ssl:send(Socket, Packet). +-spec handle_container_command(binary(), term(), ssl:sslsocket()) -> ok. +handle_container_command(Ref, #{<<"action">> := <<"list">>}, Socket) -> + Reply = docker_commands:get_containers(), + send_container_response(Socket, Ref, Reply), + ok; +handle_container_command(Ref, #{<<"action">> := <<"deploy">>, <<"task_id">> := TaskId, <<"params">> := Params}, Socket) -> + Reply = docker_deploy_manager:deploy(TaskId, Params), + send_container_response(Socket, Ref, Reply), + ok; +handle_container_command(Ref, #{<<"action">> := <<"start">>, <<"target">> := Target}, Socket) -> + Reply = docker_commands:start_container(container_target(Target)), + send_container_response(Socket, Ref, Reply), + ok; +handle_container_command(Ref, #{<<"action">> := <<"stop">>, <<"target">> := Target, <<"timeout_seconds">> := TimeoutSeconds}, Socket) -> + Reply = docker_commands:stop_container(container_target(Target), TimeoutSeconds), + send_container_response(Socket, Ref, Reply), + ok; +handle_container_command(Ref, #{<<"action">> := <<"kill">>, <<"target">> := Target, <<"signal">> := Signal}, Socket) -> + Reply = docker_commands:kill_container(container_target(Target), to_binary(Signal)), + send_container_response(Socket, Ref, Reply), + ok; +handle_container_command(Ref, #{<<"action">> := <<"remove">>, <<"target">> := Target, <<"force">> := Force, <<"remove_volumes">> := RemoveVolumes}, Socket) -> + Reply = docker_commands:remove_container(container_target(Target), to_bool(Force), to_bool(RemoveVolumes)), + send_container_response(Socket, Ref, Reply), + ok; +handle_container_command(Ref, #{<<"action">> := <<"config">>, <<"target">> := Target, <<"config">> := Config}, Socket) -> + Reply = docker_helper:update_container_config(container_target(Target), iolist_to_binary(Config)), + send_container_response(Socket, Ref, Reply), + ok; +handle_container_command(Ref, Request, Socket) -> + logger:notice("[efka_iot_client] get an invalid command: ~p, agent invalid", [Request]), + send_container_response(Socket, Ref, {error, <<"agent invalid">>}), + ok. + +-spec handle_stream_frame(term(), term(), #state{}) -> gen_statem:event_handler_result(atom(), #state{}). +handle_stream_frame(StreamId, {<<"open">>, Params}, State = #state{streams = Streams}) + when is_integer(StreamId), StreamId > 0, is_map(Params) -> + case valid_iot_stream_id(StreamId) andalso not maps:is_key(StreamId, Streams) of + true -> + {ok, {WorkerPid, MonitorRef}} = efka_iot_stream:start_stream(StreamId, Params), + StreamState = #stream_state{worker_pid = WorkerPid, monitor_ref = MonitorRef}, + {keep_state, State#state{streams = maps:put(StreamId, StreamState, Streams)}}; + false -> + send_stream(StreamId, {<<"reset">>, <<"invalid stream open">>}), + {keep_state, State} + end; +handle_stream_frame(StreamId, Body, State = #state{streams = Streams}) + when is_integer(StreamId), StreamId > 0 -> + case maps:get(StreamId, Streams, undefined) of + undefined -> + maybe_reset_unknown_stream(StreamId, Body), + {keep_state, State}; + StreamState = #stream_state{worker_pid = WorkerPid} -> + WorkerPid ! {stream, StreamId, Body}, + case Body of + {<<"reset">>, _Reason} -> + demonitor_stream(StreamState), + {keep_state, State#state{streams = maps:remove(StreamId, Streams)}}; + _ -> + {keep_state, State} + end + end; +handle_stream_frame(StreamId, Body, State) -> + logger:warning("[efka_iot_client] invalid stream frame, stream_id: ~p, body: ~p", [StreamId, Body]), + {keep_state, State}. + +-spec maybe_reset_unknown_stream(stream_id(), term()) -> ok. +maybe_reset_unknown_stream(_StreamId, {<<"reset">>, _Reason}) -> + ok; +maybe_reset_unknown_stream(_StreamId, <<"fin">>) -> + ok; +maybe_reset_unknown_stream(_StreamId, {<<"open_error">>, _Reason}) -> + ok; +maybe_reset_unknown_stream(StreamId, _Body) -> + send_stream(StreamId, {<<"reset">>, <<"unknown stream">>}). + +-spec valid_iot_stream_id(stream_id()) -> boolean(). +valid_iot_stream_id(StreamId) -> + StreamId rem 2 =:= 1. + +-spec take_stream_by_monitor(reference(), pid(), map()) -> + {stream_id(), #stream_state{}, map()} | error. +take_stream_by_monitor(MonitorRef, WorkerPid, Streams) -> + take_stream_by_monitor(MonitorRef, WorkerPid, maps:iterator(Streams), Streams). + +take_stream_by_monitor(MonitorRef, WorkerPid, Iter, Streams) -> + case maps:next(Iter) of + none -> + error; + {StreamId, StreamState = #stream_state{worker_pid = WorkerPid, monitor_ref = MonitorRef}, _NextIter} -> + {StreamId, StreamState, maps:remove(StreamId, Streams)}; + {_StreamId, _StreamState, NextIter} -> + take_stream_by_monitor(MonitorRef, WorkerPid, NextIter, Streams) + end. + +-spec close_all_streams(term(), #state{}) -> #state{}. +close_all_streams(Reason, State = #state{streams = Streams}) -> + maps:foreach(fun(StreamId, StreamState = #stream_state{worker_pid = WorkerPid}) -> + demonitor_stream(StreamState), + WorkerPid ! {stream, StreamId, {<<"reset">>, {<<"channel_closed">>, safe_term(Reason)}}} + end, Streams), + State#state{streams = #{}}. + +-spec demonitor_stream(#stream_state{}) -> ok. +demonitor_stream(#stream_state{monitor_ref = MonitorRef}) -> + erlang:demonitor(MonitorRef, [flush]), + ok. + -spec safe_reply(term()) -> term(). safe_reply(ok) -> <<"ok">>; diff --git a/apps/efka/src/iot/efka_iot_stream.erl b/apps/efka/src/iot/efka_iot_stream.erl new file mode 100644 index 0000000..60062fd --- /dev/null +++ b/apps/efka/src/iot/efka_iot_stream.erl @@ -0,0 +1,150 @@ +%%%------------------------------------------------------------------- +%%% @doc One transparent TCP stream from iot to the local manager service. +%%% The iot/efka protocol does not inspect HTTP; all HTTP bytes are carried +%%% in <<"data">> frames. +%%% @end +%%%------------------------------------------------------------------- +-module(efka_iot_stream). + +-export([start_stream/2]). +-export([run/2]). + +-define(DEFAULT_CONNECT_TIMEOUT, 3000). +-define(DEFAULT_IDLE_TIMEOUT, 120000). + +-type stream_id() :: pos_integer(). + +-record(options, { + host :: inet:hostname() | inet:ip_address(), + port :: inet:port_number(), + connect_timeout = 0 :: timeout(), + idle_timeout = 0 :: timeout() +}). + +-spec start_stream(StreamId :: integer(), Params :: map()) -> {ok, {pid(), reference()}}. +start_stream(StreamId, Params) when is_integer(StreamId), StreamId > 0, is_map(Params) -> + {ok, spawn_monitor(?MODULE, run, [StreamId, Params])}. + +-spec run(StreamId :: stream_id(), Params :: map()) -> ok. +run(StreamId, Params) when is_integer(StreamId), StreamId > 0, is_map(Params) -> + try run0(StreamId, Params) of + ok -> + ok + catch + Class:Reason:Stack -> + logger:warning("[efka_iot_stream] stream_id: ~p crashed, class: ~p, reason: ~p, stack: ~p", + [StreamId, Class, Reason, Stack]), + efka_iot_client:send_stream(StreamId, {<<"reset">>, safe_term({Class, Reason})}), + ok + after + efka_iot_client:stream_done(StreamId) + end. + +-spec run0(stream_id(), map()) -> ok. +run0(StreamId, Params) -> + case open_target_socket(Params) of + {ok, Socket, IdleTimeout} -> + efka_iot_client:send_stream(StreamId, <<"opened">>), + ok = inet:setopts(Socket, [{active, once}]), + loop(StreamId, Socket, IdleTimeout); + {error, Reason} -> + efka_iot_client:send_stream(StreamId, {<<"open_error">>, safe_term(Reason)}), + ok + end. + +-spec open_target_socket(map()) -> {ok, gen_tcp:socket(), timeout()} | {error, term()}. +open_target_socket(Params) when is_map(Params) -> + case stream_target_options(Params) of + error -> + {error, <<"invalid_target">>}; + {ok, #options{host = Host, port = Port, connect_timeout = ConnectTimeout, idle_timeout = IdleTimeout}} -> + SocketOpts = [ + binary, + {packet, raw}, + {active, false}, + {nodelay, true} + ], + case gen_tcp:connect(Host, Port, SocketOpts, ConnectTimeout) of + {ok, Socket} -> + {ok, Socket, IdleTimeout}; + {error, Reason} -> + {error, Reason} + end + end. + +-spec stream_target_options(map()) -> error | {ok, #options{}}. +stream_target_options(Params = #{<<"target">> := Target0}) -> + Target = binary_to_atom(Target0), + {ok, StreamTargets} = application:get_env(efka, stream_targets), + case proplists:get_value(Target, StreamTargets, undefined) of + undefined -> + error; + Props -> + Host = proplists:get_value(host, Props), + Port = proplists:get_value(port, Props), + ConnectTimeout0 = proplists:get_value(connect_timeout, Props, ?DEFAULT_CONNECT_TIMEOUT), + IdleTimeout0 = proplists:get_value(idle_timeout, Props, ?DEFAULT_IDLE_TIMEOUT), + {ok, #options{ + host = Host, + port = Port, + connect_timeout = maps:get(<<"connect_timeout">>, Params, ConnectTimeout0), + idle_timeout = maps:get(<<"idle_timeout">>, Params, IdleTimeout0) + }} + end. + +-spec loop(stream_id(), gen_tcp:socket(), timeout()) -> ok. +loop(StreamId, Socket, IdleTimeout) -> + receive + {stream, StreamId, {<<"data">>, Data}} when is_binary(Data) -> + case gen_tcp:send(Socket, Data) of + ok -> + loop(StreamId, Socket, IdleTimeout); + {error, Reason} -> + efka_iot_client:send_stream(StreamId, {<<"reset">>, safe_term(Reason)}), + close_socket(Socket) + end; + {stream, StreamId, <<"fin">>} -> + _ = gen_tcp:shutdown(Socket, write), + loop(StreamId, Socket, IdleTimeout); + {stream, StreamId, {<<"reset">>, _Reason}} -> + close_socket(Socket); + {tcp, Socket, Data} -> + efka_iot_client:send_stream(StreamId, {<<"data">>, Data}), + ok = inet:setopts(Socket, [{active, once}]), + loop(StreamId, Socket, IdleTimeout); + {tcp_closed, Socket} -> + efka_iot_client:send_stream(StreamId, <<"fin">>), + close_socket(Socket); + {tcp_error, Socket, Reason} -> + efka_iot_client:send_stream(StreamId, {<<"reset">>, safe_term(Reason)}), + close_socket(Socket); + Info -> + logger:debug("[efka_iot_stream] stream_id: ~p ignore unknown info: ~p", [StreamId, Info]), + loop(StreamId, Socket, IdleTimeout) + after IdleTimeout -> + efka_iot_client:send_stream(StreamId, {<<"reset">>, <<"idle_timeout">>}), + close_socket(Socket) + end. + +-spec close_socket(gen_tcp:socket()) -> ok. +close_socket(Socket) -> + catch gen_tcp:close(Socket), + ok. + +-spec safe_term(term()) -> term(). +safe_term(true) -> + true; +safe_term(false) -> + false; +safe_term(undefined) -> + undefined; +safe_term(Value) when is_atom(Value) -> + atom_to_binary(Value, utf8); +safe_term(Value) when is_map(Value) -> + maps:from_list([{safe_term(K), safe_term(V)} || {K, V} <- maps:to_list(Value)]); +safe_term(Value) when is_list(Value) -> + [safe_term(Item) || Item <- Value]; +safe_term(Value) when is_tuple(Value) -> + list_to_tuple([safe_term(Item) || Item <- tuple_to_list(Value)]); +safe_term(Value) -> + Value. diff --git a/config/sys.config.src b/config/sys.config.src index 7b64902..be284dd 100644 --- a/config/sys.config.src +++ b/config/sys.config.src @@ -15,6 +15,15 @@ {udp_port, 24000} ]}, + {stream_targets, [ + {manager, [ + {host, "127.0.0.1"}, + {port, 18091}, + {connect_timeout, 3000}, + {idle_timeout, 120000} + ]} + ]}, + {heartbeat, [ {interval, 5000} ]}, diff --git a/docs/efka_iot_protocol.md b/docs/efka_iot_protocol.md index fd39e7c..d7d5845 100644 --- a/docs/efka_iot_protocol.md +++ b/docs/efka_iot_protocol.md @@ -21,6 +21,7 @@ {<<"command">>, Ref, {Domain, Payload}} {<<"command_response">>, Ref, {Domain, Reply}} {<<"message">>, Body} +{<<"stream">>, StreamId, Body} ``` 语义说明: @@ -32,11 +33,56 @@ | `{<<"command">>, Ref, {Domain, Payload}}` | iot -> efka | iot 下发命令,需要 efka 回复 | | `{<<"command_response">>, Ref, {Domain, Reply}}` | efka -> iot | efka 对 iot command 的回复 | | `{<<"message">>, Body}` | 双向 | 异步消息,不要求回复 | +| `{<<"stream">>, StreamId, Body}` | 双向 | 透明 TCP 字节流多路复用 | `command` 和 `command_response` 的 `Domain` 表示业务域,目前支持: - `<<"container">>` +## 透明 TCP Stream + +`stream` 用于在同一条 TLS 长连接上复用多条透明 TCP 字节流。旧的 `request`、`response`、`command`、`command_response`、`message` 帧格式保持不变;`StreamId = 0` 只作为保留概念,不出现在新的 `stream` 帧里。 + +`StreamId` 规则: + +- `0`:保留给现有控制流。 +- 奇数:`iot` 发起。 +- 偶数:`efka` 发起,当前预留。 +- 同一条 TLS 连接内 `StreamId` 不复用。 + +`Body` 取值: + +```erlang +{<<"open">>, #{<<"target">> => <<"manager">>}} +<<"opened">> +{<<"open_error">>, Reason} +{<<"data">>, Chunk} +<<"fin">> +{<<"reset">>, Reason} +``` + +语义: + +- `open`:发起方请求打开一个透明 TCP stream。当前 `efka` 端支持 `target = <<"manager">>`,由本地配置决定实际本机管理程序地址。 +- `opened`:本地 TCP 连接建立成功。 +- `open_error`:本地 TCP 连接建立失败,stream 结束。 +- `data`:透明字节数据,`Chunk` 是 binary。HTTP 请求头、请求体、响应头、响应体都只是该 binary 的内容,iot/efka 中间协议不解析 HTTP。 +- `fin`:半关闭,发送方后续不再发送 `data`,但仍可继续接收对端 `data`。 +- `reset`:异常关闭,双方应立即释放该 `StreamId` 的资源。 + +典型 iot 发起访问 efka 本机监控服务的流程: + +```erlang +iot -> efka: {<<"stream">>, 1, {<<"open">>, #{<<"target">> => <<"manager">>}}} +efka -> iot: {<<"stream">>, 1, <<"opened">>} + +iot -> efka: {<<"stream">>, 1, {<<"data">>, RawHttpRequestBytes}} +iot -> efka: {<<"stream">>, 1, <<"fin">>} + +efka -> iot: {<<"stream">>, 1, {<<"data">>, RawHttpResponseBytes}} +efka -> iot: {<<"stream">>, 1, <<"fin">>} +``` + ## 鉴权请求 初始连接由 `efka` 发起鉴权 request。每条 TLS 连接只允许一次鉴权;`iot` 侧鉴权成功后会在 `ssl_channel` 标记该连接已鉴权,如果同一连接再次发送 `auth_request`,`iot` 会直接关闭连接。