diff --git a/apps/efka/src/iot/efka_iot_client.erl b/apps/efka/src/iot/efka_iot_client.erl index d6ed883..a680234 100644 --- a/apps/efka/src/iot/efka_iot_client.erl +++ b/apps/efka/src/iot/efka_iot_client.erl @@ -451,9 +451,10 @@ handle_container_command(Ref, Request, Socket) -> send_container_response(Socket, Ref, {error, <<"agent invalid">>}), ok. +%% iot 发起 open 时不携带参数;efka 固定按本地 stream_target 建连。 -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) -> +handle_stream_frame(StreamId, <<"open">>, State = #state{streams = Streams}) + when is_integer(StreamId), StreamId > 0 -> case valid_iot_stream_id(StreamId) andalso not maps:is_key(StreamId, Streams) of true -> {ok, {WorkerPid, MonitorRef}} = efka_iot_stream:start_stream(StreamId), @@ -463,6 +464,8 @@ handle_stream_frame(StreamId, {<<"open">>, Params}, State = #state{streams = Str send_stream(StreamId, {<<"reset">>, <<"invalid stream open">>}), {keep_state, State} end; +handle_stream_frame(StreamId, {<<"open">>, _Params}, State) -> + handle_stream_frame(StreamId, <<"open">>, State); handle_stream_frame(StreamId, Body, State = #state{streams = Streams}) when is_integer(StreamId), StreamId > 0 -> case maps:get(StreamId, Streams, undefined) of