From 1cf007f946096b936c1b695c23cf57c7145bfbb8 Mon Sep 17 00:00:00 2001 From: anlicheng <244108715@qq.com> Date: Sun, 3 May 2026 14:06:05 +0800 Subject: [PATCH] fix transport --- config/sys-dev.config | 1 - config/sys-prod.config | 1 - src/quic/sdlan_quic_transport.erl | 18 +------ src/quic/sdlan_session.erl | 21 ++++++++ src/ssl/sdlan_ssl_transport.erl | 84 ++++--------------------------- 5 files changed, 32 insertions(+), 93 deletions(-) diff --git a/config/sys-dev.config b/config/sys-dev.config index d4f8350..e237178 100644 --- a/config/sys-dev.config +++ b/config/sys-dev.config @@ -31,7 +31,6 @@ {keyfile, "server.key"}, {limits, [ {max_packet_size, 16384}, - {socket_active_n, 100}, %% 单位为秒 {heartbeat_sec, 15} ]} diff --git a/config/sys-prod.config b/config/sys-prod.config index 7b62d8b..7607f2b 100644 --- a/config/sys-prod.config +++ b/config/sys-prod.config @@ -31,7 +31,6 @@ {keyfile, "server.key"}, {limits, [ {max_packet_size, 16384}, - {socket_active_n, 100}, %% 单位为秒 {heartbeat_sec, 15} ]} diff --git a/src/quic/sdlan_quic_transport.erl b/src/quic/sdlan_quic_transport.erl index 595720c..3c04a34 100644 --- a/src/quic/sdlan_quic_transport.erl +++ b/src/quic/sdlan_quic_transport.erl @@ -118,24 +118,8 @@ handle_event(info, {quic, dgram_state_changed, Conn, Opts = #{dgram_send_enabled handle_event(info, {quic, new_stream, Stream, Opts}, waiting_stream, State = #state{max_packet_size = MaxPacketSize, heartbeat_sec = HeartbeatSec}) -> logger:debug("[sdlan_quic_transport] call new_stream: ~p, opts: ~p", [Stream, Opts]), - Ipv6Assist = case application:get_env(sdlan, ipv6_assist_info) of - {ok, {V6Bytes, Port}} -> - #'SDLV6Info' { - v6 = V6Bytes, - port = Port - }; - _ -> - undefined - end, %% 发送欢迎消息 - WelcomePkt = sdlan_pb:encode_msg(#'SDLWelcome'{ - version = 1, - max_bidi_streams = 1, - max_packet_size = MaxPacketSize, - heartbeat_sec = HeartbeatSec, - ipv6_assist = Ipv6Assist - }), - quic_send(Stream, <>), + quic_send(Stream, sdlan_session:welcome_packet(MaxPacketSize, HeartbeatSec)), logger:debug("[sdlan_quic_transport] get stream: ~p, send welcome", [Stream]), {next_state, initialized, State#state{stream = Stream}}; diff --git a/src/quic/sdlan_session.erl b/src/quic/sdlan_session.erl index 35b29b6..8842931 100644 --- a/src/quic/sdlan_session.erl +++ b/src/quic/sdlan_session.erl @@ -13,6 +13,7 @@ %% API -export([new/1, state_name/1, handle_frame/2, handle_timeout/2]). -export([send_event/2, command/4, close/1, debug_info/1]). +-export([welcome_packet/2]). -export([test_rules/2]). -export_type([state/0]). @@ -71,6 +72,26 @@ test_rules(SrcIdentityId, DstIdentityId) when is_integer(SrcIdentityId), is_inte state_name(#state{status = Status}) -> Status. +-spec welcome_packet(MaxPacketSize :: integer(), HeartbeatSec :: integer()) -> binary(). +welcome_packet(MaxPacketSize, HeartbeatSec) -> + Ipv6Assist = case application:get_env(sdlan, ipv6_assist_info) of + {ok, {V6Bytes, Port}} -> + #'SDLV6Info' { + v6 = V6Bytes, + port = Port + }; + _ -> + undefined + end, + WelcomePkt = sdlan_pb:encode_msg(#'SDLWelcome'{ + version = 1, + max_bidi_streams = 1, + max_packet_size = MaxPacketSize, + heartbeat_sec = HeartbeatSec, + ipv6_assist = Ipv6Assist + }), + <>. + -spec handle_frame(Frame :: binary(), State :: #state{}) -> {ok, NewState :: #state{}, Packets :: [binary()]} | {stop, Reason :: term(), NewState :: #state{}, Packets :: [binary()]}. diff --git a/src/ssl/sdlan_ssl_transport.erl b/src/ssl/sdlan_ssl_transport.erl index 341ee8a..884bb6d 100644 --- a/src/ssl/sdlan_ssl_transport.erl +++ b/src/ssl/sdlan_ssl_transport.erl @@ -13,8 +13,6 @@ -behaviour(gen_statem). -behaviour(ranch_protocol). --define(SOCKET_ACTIVE_N, 100). - %% Ranch protocol callback -export([start_link/4]). @@ -25,17 +23,10 @@ ref :: ranch:ref(), socket :: undefined | ssl:sslsocket(), transport :: module(), - ok_msg :: atom(), - closed_msg :: atom(), - error_msg :: atom(), - passive_msg :: atom(), max_packet_size = 16384, heartbeat_sec = 15, - socket_active_n = ?SOCKET_ACTIVE_N, - %% 累积器,用于处理协议framing的解析 - buf = <<>>, session :: sdlan_session:state(), frames_recv = 0, @@ -58,20 +49,13 @@ start_link(Ref, Socket, Transport, Limits) -> init([Ref, Socket, Transport, Limits]) -> MaxPacketSize = proplists:get_value(max_packet_size, Limits, 16384), HeartbeatSec = proplists:get_value(heartbeat_sec, Limits, 15), - SocketActiveN = proplists:get_value(socket_active_n, Limits, proplists:get_value(stream_active_n, Limits, ?SOCKET_ACTIVE_N)), - {OkMsg, ClosedMsg, ErrorMsg} = Transport:messages(), Session = sdlan_session:new(HeartbeatSec), {ok, handshaking, #state{ ref = Ref, socket = Socket, transport = Transport, - ok_msg = OkMsg, - closed_msg = ClosedMsg, - error_msg = ErrorMsg, - passive_msg = passive_msg(OkMsg), max_packet_size = MaxPacketSize, heartbeat_sec = HeartbeatSec, - socket_active_n = SocketActiveN, session = Session }}. @@ -80,38 +64,26 @@ callback_mode() -> handle_event(info, {handshake, Ref, Transport, Socket, Timeout}, handshaking, State = #state{ref = Ref, transport = Transport, max_packet_size = MaxPacketSize, - heartbeat_sec = HeartbeatSec, socket_active_n = SocketActiveN}) -> + heartbeat_sec = HeartbeatSec}) -> case Transport:handshake(Socket, [], Timeout) of {ok, SslSocket} -> - ok = Transport:setopts(SslSocket, [{mode, binary}, {active, SocketActiveN}]), - WelcomePkt = welcome_packet(MaxPacketSize, HeartbeatSec), - ssl_send(Transport, SslSocket, <>), + ok = Transport:setopts(SslSocket, [{mode, binary}, {packet, 2}, {active, true}]), + ssl_send(Transport, SslSocket, sdlan_session:welcome_packet(MaxPacketSize, HeartbeatSec)), logger:debug("[sdlan_ssl_transport] ssl handshake ok, send welcome"), {next_state, initialized, State#state{socket = SslSocket}}; {error, Reason} -> {stop, {ssl_handshake_failed, Reason}, State} end; -handle_event(info, {OkMsg, Socket, Data}, _StateName, - State = #state{socket = Socket, ok_msg = OkMsg, buf = Buf, max_packet_size = MaxPacketSize, - bytes_recv = BytesRecv, frames_recv = FramesRecv}) when is_binary(Data) -> - case decode_frames(<>, MaxPacketSize) of - {error, Reason} -> - {stop, Reason, State}; - {ok, NBuf, Frames} -> - Actions = [{next_event, internal, {frame, Frame}} || Frame <- Frames], - {keep_state, State#state{buf = NBuf, bytes_recv = BytesRecv + byte_size(Data), frames_recv = FramesRecv + length(Frames)}, Actions} - end; +handle_event(info, {ssl, Socket, Data}, _StateName, + State = #state{socket = Socket, bytes_recv = BytesRecv, frames_recv = FramesRecv}) when is_binary(Data) -> + {keep_state, State#state{bytes_recv = BytesRecv + byte_size(Data), frames_recv = FramesRecv + 1}, + [{next_event, internal, {frame, Data}}]}; -handle_event(info, {PassiveMsg, Socket}, _StateName, - State = #state{socket = Socket, transport = Transport, passive_msg = PassiveMsg, socket_active_n = SocketActiveN}) -> - ok = Transport:setopts(Socket, [{active, SocketActiveN}]), - {keep_state, State}; - -handle_event(info, {ClosedMsg, Socket}, _StateName, State = #state{socket = Socket, closed_msg = ClosedMsg}) -> +handle_event(info, {ssl_closed, Socket}, _StateName, State = #state{socket = Socket}) -> expected_stop(socket_closed, State); -handle_event(info, {ErrorMsg, Socket, Reason}, _StateName, State = #state{socket = Socket, error_msg = ErrorMsg}) -> +handle_event(info, {ssl_error, Socket, Reason}, _StateName, State = #state{socket = Socket}) -> expected_stop({socket_error, Reason}, State); %% 处理内部的包消息 @@ -173,41 +145,10 @@ code_change(_OldVsn, StateName, State = #state{}, _Extra) -> %%% Internal functions %%%=================================================================== -welcome_packet(MaxPacketSize, HeartbeatSec) -> - Ipv6Assist = case application:get_env(sdlan, ipv6_assist_info) of - {ok, {V6Bytes, Port}} -> - #'SDLV6Info' { - v6 = V6Bytes, - port = Port - }; - _ -> - undefined - end, - sdlan_pb:encode_msg(#'SDLWelcome'{ - version = 1, - max_bidi_streams = 1, - max_packet_size = MaxPacketSize, - heartbeat_sec = HeartbeatSec, - ipv6_assist = Ipv6Assist - }). - -%% 有2种情况 -%% 1. 收到了多个完整的请求 -%% 2. 不完整,则不处理 --spec decode_frames(Buf :: binary(), MaxPacketSize :: integer()) -> {ok, RestBin::binary(), Frames :: list()} | {error, Reason :: any()}. -decode_frames(Buf, MaxPacketSize) when is_binary(Buf) -> - decode_frames0(Buf, MaxPacketSize, []). -decode_frames0(<>, MaxPacketSize, _Frames) when Len > MaxPacketSize -> - {error, frame_too_large}; -decode_frames0(<>, MaxPacketSize, Frames) -> - decode_frames0(Rest, MaxPacketSize, [Frame|Frames]); -decode_frames0(Rest, _MaxPacketSize, Frames) -> - {ok, Rest, lists:reverse(Frames)}. - ssl_send(Transport, Socket, Packet) when is_binary(Packet) -> Len = byte_size(Packet), true = Len =< 65535, - case Transport:send(Socket, <>) of + case Transport:send(Socket, Packet) of ok -> incr_counter(ssl_frames_sent, 1), incr_counter(ssl_bytes_sent, Len + 2), @@ -226,14 +167,10 @@ expected_stop(Reason, State) -> next_state_name(_StateName, Session) -> sdlan_session:state_name(Session). -passive_msg(OkMsg) -> - list_to_atom(atom_to_list(OkMsg) ++ "_passive"). - debug_info(StateName, #state{ session = Session, frames_recv = FramesRecv, bytes_recv = BytesRecv, - socket_active_n = SocketActiveN, heartbeat_sec = HeartbeatSec }) -> ProcInfo = maps:from_list(process_info(self(), [message_queue_len, memory, reductions])), @@ -245,7 +182,6 @@ debug_info(StateName, #state{ frames_sent => get_counter(ssl_frames_sent), bytes_recv => BytesRecv, bytes_sent => get_counter(ssl_bytes_sent), - socket_active_n => SocketActiveN, heartbeat_sec => HeartbeatSec }).