From 5d91fef1117e89e91024c995af89828c5a15656f Mon Sep 17 00:00:00 2001 From: anlicheng <244108715@qq.com> Date: Sat, 18 Apr 2026 17:34:50 +0800 Subject: [PATCH] fix --- src/transport/tcp/tcp_channel.erl | 19 +++++++++++++------ 1 file changed, 13 insertions(+), 6 deletions(-) diff --git a/src/transport/tcp/tcp_channel.erl b/src/transport/tcp/tcp_channel.erl index 5b1b7e1..71609ec 100644 --- a/src/transport/tcp/tcp_channel.erl +++ b/src/transport/tcp/tcp_channel.erl @@ -36,6 +36,12 @@ inflight = #{} }). +-record(inflight_request, { + receiver_pid :: pid(), + ref :: reference(), + timer_ref :: reference() +}). + %% 向通道中写入消息 -spec pub(Pid :: pid(), Topic :: binary(), Qos :: integer(), Content :: binary()) -> no_return(). pub(Pid, Topic, Qos, Content) when is_pid(Pid), is_binary(Topic), is_integer(Qos), is_binary(Content) -> @@ -81,7 +87,7 @@ init(Ref, Transport, _Opts = []) -> handle_call({cancel_jsonrpc_call, Ref}, _From, State = #state{inflight = Inflight}) -> case take_inflight_by_ref(Ref, Inflight) of - {ok, _PacketId, {_ReceiverPid, _Ref, TimerRef}, NInflight} -> + {ok, #inflight_request{timer_ref = TimerRef}, NInflight} -> erlang:cancel_timer(TimerRef), {reply, ok, State#state{inflight = NInflight}}; error -> @@ -115,7 +121,8 @@ handle_cast({jsonrpc_call, ReceiverPid, Ref, {Method, Params}}, State = #state{t TimerRef = erlang:start_timer(?INFLIGHT_TIMEOUT, self(), {jsonrpc_timeout, NPacketId}), Transport:send(Socket, <>), - {noreply, State#state{packet_id = NextPacketId, inflight = maps:put(NPacketId, {ReceiverPid, Ref, TimerRef}, Inflight)}}; + RequestInfo = #inflight_request{receiver_pid = ReceiverPid, ref = Ref, timer_ref = TimerRef}, + {noreply, State#state{packet_id = NextPacketId, inflight = maps:put(NPacketId, RequestInfo, Inflight)}}; {error, inflight_full} -> logger:warning("[ws_channel] uuid: ~p, inflight requests exhausted", [State#state.uuid]), {noreply, State} @@ -201,7 +208,7 @@ handle_info({tcp, Socket, < {noreply, State}; - {{ReceiverPid, Ref, TimerRef}, NInflight} -> + {#inflight_request{receiver_pid = ReceiverPid, ref = Ref, timer_ref = TimerRef}, NInflight} -> erlang:cancel_timer(TimerRef), case is_pid(ReceiverPid) andalso is_process_alive(ReceiverPid) of true -> @@ -214,7 +221,7 @@ handle_info({tcp, Socket, < case maps:get(PacketId, Inflight, undefined) of - {_ReceiverPid, Ref, TimerRef} -> + #inflight_request{ref = Ref, timer_ref = TimerRef} -> logger:warning("[ws_channel] jsonrpc request timeout, packet_id: ~p, ref: ~p", [PacketId, Ref]), {noreply, State#state{inflight = maps:remove(PacketId, Inflight)}}; _ -> @@ -250,10 +257,10 @@ code_change(_OldVsn, State, _Extra) -> {ok, State}. take_inflight_by_ref(Ref, Inflight) -> - maps:fold(fun(PacketId, Value = {_ReceiverPid, PacketRef, _TimerRef}, Acc) -> + maps:fold(fun(PacketId, Request = #inflight_request{ref = PacketRef}, Acc) -> case Acc of error when PacketRef =:= Ref -> - {ok, PacketId, Value, maps:remove(PacketId, Inflight)}; + {ok, Request, maps:remove(PacketId, Inflight)}; _ -> Acc end