This commit is contained in:
anlicheng 2026-07-07 13:13:20 +08:00
parent 9dae7b5b60
commit e4a59ed236

View File

@ -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