Compare commits
10 Commits
3cf38f3e89
...
c106302d99
| Author | SHA1 | Date | |
|---|---|---|---|
| c106302d99 | |||
| 4791b66086 | |||
| 22db3e1f46 | |||
| 94dc39805e | |||
| df3696cd50 | |||
| 77df81f8c3 | |||
| 7fba178acf | |||
| 9d5bc114ed | |||
| d6fb2a480f | |||
| 78c1b9a970 |
39
README.md
39
README.md
@ -13,10 +13,39 @@ Run
|
||||
|
||||
$ rebar3 shell
|
||||
|
||||
The UDP listener is configured in `config/sys.config` and listens on port `7000`
|
||||
by default. It echoes received datagrams.
|
||||
The UDP listener is configured in `config/sys.config`; the checked-in config
|
||||
listens on port `1443`, while the code fallback default is `16380`. It accepts
|
||||
RelayKit UDP relay frames.
|
||||
|
||||
Test
|
||||
----
|
||||
Protocol
|
||||
--------
|
||||
|
||||
$ echo "ping" | nc -u -w1 127.0.0.1 7000
|
||||
Each UDP datagram is encrypted with ChaCha20-Poly1305. The encrypted datagram
|
||||
uses the same combined format as CryptoKit `ChaChaPoly`:
|
||||
|
||||
nonce(12), ciphertext, tag(16)
|
||||
|
||||
The decrypted plaintext is one `RKP1` frame:
|
||||
|
||||
magic(4) = "RKP1", version(1) = 1, type(1), reserved(2),
|
||||
stream_id(8), payload_len(4), payload(payload_len)
|
||||
|
||||
Frame types are `open = 1`, `open_ack = 2`, `data = 3`, `close = 4`, and
|
||||
`error = 5`.
|
||||
`open` payload is binary encoded as:
|
||||
|
||||
host_len(2), host(host_len), port(2),
|
||||
username_len(2), username(username_len),
|
||||
password_len(2), password(password_len)
|
||||
|
||||
Integer fields are unsigned big-endian values. `host`, `username`, and
|
||||
`password` are binary strings. The server is a transparent TCP relay, so TLS
|
||||
negotiation, HTTP methods, and request bytes are carried inside subsequent
|
||||
`data` frames without being interpreted by the server.
|
||||
|
||||
`username` and `password` are checked against the `relay_server` application
|
||||
`users` config:
|
||||
|
||||
{users, [
|
||||
{admin, "password"}
|
||||
]}
|
||||
|
||||
38
apps/relay_server/src/chacha20_cipher.erl
Normal file
38
apps/relay_server/src/chacha20_cipher.erl
Normal file
@ -0,0 +1,38 @@
|
||||
%%%-------------------------------------------------------------------
|
||||
%% @doc ChaCha20-Poly1305 helpers compatible with CryptoKit ChaChaPoly.
|
||||
%% 1. 只解决需要加解密数据的逻辑处理
|
||||
%% @end
|
||||
%%%-------------------------------------------------------------------
|
||||
|
||||
-module(chacha20_cipher).
|
||||
|
||||
-export([encrypt/2, decrypt/4]).
|
||||
|
||||
-define(KEY_BYTES, 32).
|
||||
-define(NONCE_BYTES, 12).
|
||||
-define(TAG_BYTES, 16).
|
||||
-define(AAD, <<>>).
|
||||
|
||||
encrypt(PlainText, Key) when is_binary(PlainText), byte_size(Key) =:= ?KEY_BYTES ->
|
||||
Nonce = crypto:strong_rand_bytes(?NONCE_BYTES),
|
||||
try crypto:crypto_one_time_aead(chacha20_poly1305, Key, Nonce, PlainText, ?AAD, true) of
|
||||
{CipherText, Tag} ->
|
||||
{ok, <<Nonce/binary, CipherText/binary, Tag/binary>>}
|
||||
catch
|
||||
_Class:_Reason ->
|
||||
{error, crypto_failed}
|
||||
end;
|
||||
|
||||
encrypt(_PlainText, _Key) ->
|
||||
{error, invalid_plaintext}.
|
||||
|
||||
decrypt(Nonce, CipherText, Tag, Key) when is_binary(Nonce), is_binary(CipherText), is_binary(Tag), is_binary(Key) ->
|
||||
try crypto:crypto_one_time_aead(chacha20_poly1305, Key, Nonce, CipherText, ?AAD, Tag, false) of
|
||||
error ->
|
||||
{error, authentication_failed};
|
||||
PlainText ->
|
||||
{ok, PlainText}
|
||||
catch
|
||||
_Class:_Reason ->
|
||||
{error, crypto_failed}
|
||||
end.
|
||||
@ -9,6 +9,7 @@
|
||||
{applications, [
|
||||
kernel,
|
||||
stdlib,
|
||||
crypto,
|
||||
esockd,
|
||||
sync
|
||||
]},
|
||||
|
||||
@ -1,32 +0,0 @@
|
||||
%%%-------------------------------------------------------------------
|
||||
%% @doc Per-peer UDP packet handler.
|
||||
%% @end
|
||||
%%%-------------------------------------------------------------------
|
||||
|
||||
-module(relay_server_udp_handler).
|
||||
|
||||
-export([start_link/3, init/3, loop/3, handle_packet/2]).
|
||||
|
||||
start_link(Transport, Peer, IdleTimeout) ->
|
||||
{ok, spawn_link(?MODULE, init, [Transport, Peer, IdleTimeout])}.
|
||||
|
||||
init(Transport, Peer, IdleTimeout) ->
|
||||
logger:debug("UDP peer connected: ~s", [esockd:format(Peer)]),
|
||||
loop(Transport, Peer, IdleTimeout).
|
||||
|
||||
loop(Transport = {udp, Server, _Sock}, Peer, IdleTimeout) ->
|
||||
receive
|
||||
{datagram, Server, <<"stop">>} ->
|
||||
logger:debug("UDP peer stopped: ~s", [esockd:format(Peer)]),
|
||||
exit(normal);
|
||||
{datagram, Server, Packet} ->
|
||||
Reply = handle_packet(Peer, Packet),
|
||||
Server ! {datagram, Peer, Reply},
|
||||
loop(Transport, Peer, IdleTimeout)
|
||||
after IdleTimeout ->
|
||||
logger:debug("UDP peer idle timeout: ~s", [esockd:format(Peer)]),
|
||||
exit(normal)
|
||||
end.
|
||||
|
||||
handle_packet(_Peer, Packet) ->
|
||||
Packet.
|
||||
@ -13,8 +13,9 @@
|
||||
|
||||
-define(SERVER, ?MODULE).
|
||||
-define(DEFAULT_LISTENER, 'relay_server/udp').
|
||||
-define(DEFAULT_LISTEN_ON, 7000).
|
||||
-define(DEFAULT_IDLE_TIMEOUT, 30000).
|
||||
-define(DEFAULT_LISTEN_ON, 1443).
|
||||
-define(DEFAULT_IDLE_TIMEOUT, 60000).
|
||||
-define(DEFAULT_CONNECT_TIMEOUT, 5000).
|
||||
-define(DEFAULT_MAX_CONNECTIONS, 1024).
|
||||
-define(DEFAULT_UDP_OPTIONS, [binary, {reuseaddr, true}]).
|
||||
|
||||
@ -34,9 +35,15 @@ init([]) ->
|
||||
Props = application:get_env(relay_server, udp_server, []),
|
||||
|
||||
Listener = proplists:get_value(listener, Props, ?DEFAULT_LISTENER),
|
||||
ListenOn = proplists:get_value(port, Props, ?DEFAULT_LISTEN_ON),
|
||||
ListenOn = proplists:get_value(
|
||||
listen_on,
|
||||
Props,
|
||||
proplists:get_value(port, Props, ?DEFAULT_LISTEN_ON)
|
||||
),
|
||||
|
||||
IdleTimeout = proplists:get_value(idle_timeout, Props, ?DEFAULT_IDLE_TIMEOUT),
|
||||
ConnectTimeout = proplists:get_value(connect_timeout, Props, ?DEFAULT_CONNECT_TIMEOUT),
|
||||
Users = application:get_env(relay_server, users, []),
|
||||
UdpOptions = proplists:get_value(udp_options, Props, ?DEFAULT_UDP_OPTIONS),
|
||||
MaxConnections = proplists:get_value(max_connections, Props, ?DEFAULT_MAX_CONNECTIONS),
|
||||
AccessRules = proplists:get_value(access_rules, Props, [{allow, all}]),
|
||||
@ -48,7 +55,7 @@ init([]) ->
|
||||
],
|
||||
Opts = maybe_add(max_conn_rate, Props, BaseOpts),
|
||||
|
||||
MFA = {relay_server_udp_handler, start_link, [IdleTimeout]},
|
||||
MFA = {relay_server_udp_peer, start_link, [IdleTimeout, ConnectTimeout, Users]},
|
||||
case esockd:open_udp(Listener, ListenOn, Opts, MFA) of
|
||||
{ok, Pid} ->
|
||||
logger:info("UDP listener ~p started on ~s",
|
||||
|
||||
288
apps/relay_server/src/relay_server_udp_peer.erl
Normal file
288
apps/relay_server/src/relay_server_udp_peer.erl
Normal file
@ -0,0 +1,288 @@
|
||||
%%%-------------------------------------------------------------------
|
||||
%% @doc Per-peer UDP relay router.
|
||||
%% @end
|
||||
%%%-------------------------------------------------------------------
|
||||
|
||||
-module(relay_server_udp_peer).
|
||||
|
||||
-behaviour(gen_server).
|
||||
|
||||
-export([start_link/4, start_link/5]).
|
||||
|
||||
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]).
|
||||
|
||||
-record(state, {
|
||||
server :: pid(),
|
||||
peer :: {inet:ip_address(), inet:port_number()},
|
||||
idle_timeout :: timeout(),
|
||||
connect_timeout :: timeout(),
|
||||
users = #{} :: #{binary() => binary()},
|
||||
streams = #{} :: #{non_neg_integer() => pid()},
|
||||
stream_ids = #{} :: #{pid() => non_neg_integer()}
|
||||
}).
|
||||
|
||||
start_link(Transport, Peer, IdleTimeout, ConnectTimeout) ->
|
||||
start_link(Transport, Peer, IdleTimeout, ConnectTimeout, load_users()).
|
||||
|
||||
start_link(Transport, Peer, IdleTimeout, ConnectTimeout, Users) ->
|
||||
gen_server:start_link(?MODULE, [Transport, Peer, IdleTimeout, ConnectTimeout, Users], []).
|
||||
|
||||
init([{udp, Server, _Sock}, Peer, IdleTimeout, ConnectTimeout, Users]) ->
|
||||
process_flag(trap_exit, true),
|
||||
logger:debug("UDP peer connected: ~s", [esockd:format(Peer)]),
|
||||
{ok, #state{
|
||||
server = Server,
|
||||
peer = Peer,
|
||||
idle_timeout = IdleTimeout,
|
||||
connect_timeout = ConnectTimeout,
|
||||
users = normalize_users(Users)
|
||||
}, IdleTimeout}.
|
||||
|
||||
handle_call(_Request, _From, State = #state{idle_timeout = IdleTimeout}) ->
|
||||
{reply, {error, bad_request}, State, IdleTimeout}.
|
||||
|
||||
handle_cast(_Request, State = #state{idle_timeout = IdleTimeout}) ->
|
||||
{noreply, State, IdleTimeout}.
|
||||
|
||||
handle_info({datagram, Server, Packet}, State = #state{server = Server}) ->
|
||||
handle_datagram(Packet, State);
|
||||
|
||||
handle_info({stream_opened, StreamId, StreamPid}, State = #state{idle_timeout = IdleTimeout}) ->
|
||||
NState = case stream_pid(StreamId, State) of
|
||||
StreamPid ->
|
||||
send_frame(open_ack, StreamId, <<>>, State),
|
||||
State;
|
||||
_Other ->
|
||||
relay_server_udp_stream:close(StreamPid),
|
||||
State
|
||||
end,
|
||||
{noreply, NState, IdleTimeout};
|
||||
|
||||
handle_info({stream_data, StreamId, StreamPid, Data}, State = #state{idle_timeout = IdleTimeout}) ->
|
||||
case stream_pid(StreamId, State) of
|
||||
StreamPid ->
|
||||
send_frame(data, StreamId, Data, State),
|
||||
{noreply, State, IdleTimeout};
|
||||
_Other ->
|
||||
{noreply, State, IdleTimeout}
|
||||
end;
|
||||
|
||||
handle_info({stream_closed, StreamId, StreamPid}, State = #state{idle_timeout = IdleTimeout}) ->
|
||||
NState = case stream_pid(StreamId, State) of
|
||||
StreamPid ->
|
||||
send_frame(close, StreamId, <<>>, State),
|
||||
remove_stream(StreamId, StreamPid, State);
|
||||
_Other ->
|
||||
State
|
||||
end,
|
||||
{noreply, NState, IdleTimeout};
|
||||
|
||||
handle_info({stream_error, StreamId, StreamPid, Reason}, State = #state{idle_timeout = IdleTimeout}) ->
|
||||
NState = case stream_pid(StreamId, State) of
|
||||
StreamPid ->
|
||||
send_error(StreamId, format_stream_error(Reason), State),
|
||||
remove_stream(StreamId, StreamPid, State);
|
||||
_Other ->
|
||||
State
|
||||
end,
|
||||
{noreply, NState, IdleTimeout};
|
||||
|
||||
handle_info({'EXIT', StreamPid, Reason}, State = #state{idle_timeout = IdleTimeout}) ->
|
||||
NState = case stream_id(StreamPid, State) of
|
||||
undefined ->
|
||||
State;
|
||||
StreamId ->
|
||||
logger:debug("UDP relay stream ~p worker exited: ~p", [StreamId, Reason]),
|
||||
remove_stream(StreamId, StreamPid, State)
|
||||
end,
|
||||
{noreply, NState, IdleTimeout};
|
||||
|
||||
handle_info(timeout, State = #state{peer = Peer}) ->
|
||||
logger:debug("UDP peer idle timeout: ~s", [esockd:format(Peer)]),
|
||||
{stop, normal, State};
|
||||
handle_info(_Info, State = #state{idle_timeout = IdleTimeout}) ->
|
||||
{noreply, State, IdleTimeout}.
|
||||
|
||||
terminate(_Reason, #state{streams = Streams}) ->
|
||||
maps:foreach(fun(_StreamId, StreamPid) -> relay_server_udp_stream:close(StreamPid) end, Streams),
|
||||
ok.
|
||||
|
||||
code_change(_OldVsn, State, _Extra) ->
|
||||
{ok, State}.
|
||||
|
||||
handle_datagram(Packet, State = #state{idle_timeout = IdleTimeout}) ->
|
||||
case relay_server_udp_protocol:decode(Packet) of
|
||||
{ok, connection_close} ->
|
||||
{stop, normal, State};
|
||||
{ok, {open, StreamId, Payload}} ->
|
||||
logger:debug("UDP peer open stream: ~p", [StreamId]),
|
||||
{noreply, handle_open(StreamId, Payload, State), IdleTimeout};
|
||||
{ok, {data, StreamId, Payload}} ->
|
||||
{noreply, handle_data(StreamId, Payload, State), IdleTimeout};
|
||||
{ok, {close, StreamId, _Payload}} ->
|
||||
{noreply, close_stream(StreamId, State), IdleTimeout};
|
||||
{ok, {error, StreamId, Payload}} ->
|
||||
logger:warning("UDP relay stream ~p client error: ~ts", [StreamId, Payload]),
|
||||
{noreply, close_stream(StreamId, State), IdleTimeout};
|
||||
{error, Reason} ->
|
||||
logger:debug("UDP relay ignored invalid frame: ~p", [Reason]),
|
||||
{noreply, State, IdleTimeout}
|
||||
end.
|
||||
|
||||
handle_open(StreamId, Payload, State = #state{streams = Streams, users = Users}) ->
|
||||
case maps:is_key(StreamId, Streams) of
|
||||
true ->
|
||||
send_error(StreamId, <<"stream already open">>, State),
|
||||
State;
|
||||
false ->
|
||||
case relay_server_udp_protocol:decode_open_request(Payload) of
|
||||
{ok, Request} ->
|
||||
case authenticate(Request, Users) of
|
||||
ok ->
|
||||
start_stream(StreamId, Request, State);
|
||||
error ->
|
||||
send_error(StreamId, <<"authentication failed">>, State),
|
||||
State
|
||||
end;
|
||||
{error, Reason} ->
|
||||
send_error(StreamId, format_error(Reason), State),
|
||||
State
|
||||
end
|
||||
end.
|
||||
|
||||
start_stream(StreamId, Request, State = #state{connect_timeout = ConnectTimeout}) ->
|
||||
case relay_server_udp_stream:start_link(self(), StreamId, Request, ConnectTimeout) of
|
||||
{ok, StreamPid} ->
|
||||
put_stream(StreamId, StreamPid, State);
|
||||
{error, Reason} ->
|
||||
send_error(StreamId, format_stream_error({start_failed, Reason}), State),
|
||||
State
|
||||
end.
|
||||
|
||||
handle_data(StreamId, Payload, State) ->
|
||||
case stream_pid(StreamId, State) of
|
||||
undefined ->
|
||||
send_error(StreamId, <<"unknown stream">>, State),
|
||||
State;
|
||||
StreamPid ->
|
||||
relay_server_udp_stream:send(StreamPid, Payload),
|
||||
State
|
||||
end.
|
||||
|
||||
close_stream(StreamId, State) ->
|
||||
case stream_pid(StreamId, State) of
|
||||
undefined ->
|
||||
State;
|
||||
StreamPid ->
|
||||
relay_server_udp_stream:close(StreamPid),
|
||||
remove_stream(StreamId, StreamPid, State)
|
||||
end.
|
||||
|
||||
stream_pid(StreamId, #state{streams = Streams}) ->
|
||||
maps:get(StreamId, Streams, undefined).
|
||||
|
||||
stream_id(StreamPid, #state{stream_ids = StreamIds}) ->
|
||||
maps:get(StreamPid, StreamIds, undefined).
|
||||
|
||||
put_stream(StreamId, StreamPid, State = #state{streams = Streams, stream_ids = StreamIds}) ->
|
||||
State#state{
|
||||
streams = maps:put(StreamId, StreamPid, Streams),
|
||||
stream_ids = maps:put(StreamPid, StreamId, StreamIds)
|
||||
}.
|
||||
|
||||
remove_stream(StreamId, StreamPid, State = #state{streams = Streams, stream_ids = StreamIds}) ->
|
||||
State#state{
|
||||
streams = maps:remove(StreamId, Streams),
|
||||
stream_ids = maps:remove(StreamPid, StreamIds)
|
||||
}.
|
||||
|
||||
send_frame(Type, StreamId, <<Part:1200/binary, Rest/binary>>, State)
|
||||
when byte_size(Rest) > 0 ->
|
||||
case send_frame0(Type, StreamId, Part, State) of
|
||||
ok ->
|
||||
send_frame(Type, StreamId, Rest, State);
|
||||
Error ->
|
||||
Error
|
||||
end;
|
||||
send_frame(Type, StreamId, Payload, State) ->
|
||||
send_frame0(Type, StreamId, Payload, State).
|
||||
|
||||
send_frame0(Type, StreamId, Part, #state{server = Server, peer = Peer}) ->
|
||||
case relay_server_udp_protocol:encode(Type, StreamId, Part) of
|
||||
{ok, Packet} ->
|
||||
logger:debug("UDP relay encrypted packet size: ~p", [byte_size(Packet)]),
|
||||
Server ! {datagram, Peer, Packet},
|
||||
ok;
|
||||
{error, Reason} ->
|
||||
logger:error("UDP relay failed to encode ~p frame for stream ~p: ~p",
|
||||
[Type, StreamId, Reason]),
|
||||
{error, Reason}
|
||||
end.
|
||||
|
||||
send_error(StreamId, Message, State) ->
|
||||
send_frame(error, StreamId, Message, State).
|
||||
|
||||
format_stream_error({connect_failed, Reason}) ->
|
||||
iolist_to_binary(io_lib:format("connect failed: ~p", [Reason]));
|
||||
format_stream_error({send_failed, Reason}) ->
|
||||
iolist_to_binary(io_lib:format("send failed: ~p", [Reason]));
|
||||
format_stream_error({start_failed, Reason}) ->
|
||||
iolist_to_binary(io_lib:format("stream start failed: ~p", [Reason]));
|
||||
format_stream_error(Reason) ->
|
||||
format_error(Reason).
|
||||
|
||||
format_error(Reason) when is_binary(Reason) ->
|
||||
Reason;
|
||||
format_error(Reason) ->
|
||||
iolist_to_binary(io_lib:format("~p", [Reason])).
|
||||
|
||||
load_users() ->
|
||||
application:get_env(relay_server, users, []).
|
||||
|
||||
normalize_users(Users) when is_list(Users) ->
|
||||
lists:foldl(fun normalize_user/2, #{}, Users);
|
||||
normalize_users(_Users) ->
|
||||
#{}.
|
||||
|
||||
normalize_user({Username0, Password0}, Users) ->
|
||||
case {to_credential_binary(Username0), to_credential_binary(Password0)} of
|
||||
{{ok, Username}, {ok, Password}} when byte_size(Username) > 0 ->
|
||||
maps:put(Username, Password, Users);
|
||||
_ ->
|
||||
Users
|
||||
end;
|
||||
normalize_user(_User, Users) ->
|
||||
Users.
|
||||
|
||||
to_credential_binary(Value) when is_binary(Value) ->
|
||||
{ok, Value};
|
||||
to_credential_binary(Value) when is_atom(Value) ->
|
||||
{ok, atom_to_binary(Value, utf8)};
|
||||
to_credential_binary(Value) when is_list(Value) ->
|
||||
try unicode:characters_to_binary(Value) of
|
||||
Binary when is_binary(Binary) -> {ok, Binary};
|
||||
_Other -> error
|
||||
catch
|
||||
_Class:_Reason -> error
|
||||
end;
|
||||
to_credential_binary(_Value) ->
|
||||
error.
|
||||
|
||||
authenticate(#{username := Username, password := Password}, Users) ->
|
||||
case maps:get(Username, Users, undefined) of
|
||||
undefined ->
|
||||
error;
|
||||
ExpectedPassword ->
|
||||
case secure_equal(Password, ExpectedPassword) of
|
||||
true -> ok;
|
||||
false -> error
|
||||
end
|
||||
end.
|
||||
|
||||
secure_equal(Left, Right) when is_binary(Left), is_binary(Right) ->
|
||||
(byte_size(Left) =:= byte_size(Right)) andalso secure_equal(Left, Right, 0) =:= 0.
|
||||
|
||||
secure_equal(<<L:8, LRest/binary>>, <<R:8, RRest/binary>>, Diff) ->
|
||||
secure_equal(LRest, RRest, Diff bor (L bxor R));
|
||||
secure_equal(<<>>, <<>>, Diff) ->
|
||||
Diff.
|
||||
112
apps/relay_server/src/relay_server_udp_protocol.erl
Normal file
112
apps/relay_server/src/relay_server_udp_protocol.erl
Normal file
@ -0,0 +1,112 @@
|
||||
%%%-------------------------------------------------------------------
|
||||
%% @doc RelayKit UDP relay frame codec.
|
||||
%% @end
|
||||
%%%-------------------------------------------------------------------
|
||||
|
||||
-module(relay_server_udp_protocol).
|
||||
|
||||
-export([decode/1, encode/3, decode_open_request/1, encode_open_request/2,
|
||||
encode_open_request/4]).
|
||||
|
||||
-define(MAGIC, <<"RKP1">>).
|
||||
-define(VERSION, 1).
|
||||
-define(HEADER_LENGTH, 20).
|
||||
|
||||
-define(CONNECTION_CLOSE, 16#FF).
|
||||
|
||||
-define(CHACHA20_KEY, <<
|
||||
16#9F, 16#4A, 16#6C, 16#0D,
|
||||
16#8B, 16#21, 16#7E, 16#3F,
|
||||
16#42, 16#D9, 16#AA, 16#78,
|
||||
16#13, 16#C5, 16#E2, 16#B6,
|
||||
16#F0, 16#A9, 16#27, 16#BC,
|
||||
16#6D, 16#31, 16#E8, 16#4C,
|
||||
16#55, 16#FA, 16#10, 16#2A,
|
||||
16#7E, 16#9D, 16#3C, 16#BB
|
||||
>>).
|
||||
|
||||
decode(<<Nonce:12/binary, CipherHead:20/binary, Tag:16/binary, Body/binary>>) ->
|
||||
case chacha20_cipher:decrypt(Nonce, CipherHead, Tag, ?CHACHA20_KEY) of
|
||||
{ok, PlainText} ->
|
||||
decode0(PlainText, Body);
|
||||
Error ->
|
||||
Error
|
||||
end;
|
||||
decode(_Packet) ->
|
||||
{error, invalid_packet}.
|
||||
|
||||
%% 全局控制指令也要加密处理
|
||||
decode0(<<"RKP1", ?VERSION:8, ?CONNECTION_CLOSE:8, 0:16, 0:64, 0:32>>, _Body) ->
|
||||
{ok, connection_close};
|
||||
decode0(<<"RKP1", ?VERSION:8, TypeNo:8, 0:16, StreamId:64, PayloadLen:32>>, Body) when is_binary(Body) ->
|
||||
case {type(TypeNo), byte_size(Body) >= PayloadLen} of
|
||||
{undefined, _} ->
|
||||
{error, unknown_frame_type};
|
||||
{_, false} ->
|
||||
{error, truncated_frame};
|
||||
{Type, true} ->
|
||||
<<Payload:PayloadLen/binary, _Trailing/binary>> = Body,
|
||||
{ok, {Type, StreamId, Payload}}
|
||||
end;
|
||||
decode0(Header, _Body) when byte_size(Header) < ?HEADER_LENGTH ->
|
||||
{error, frame_too_short};
|
||||
decode0(_Header, _Body) ->
|
||||
{error, bad_frame_header}.
|
||||
|
||||
encode(Type, StreamId, Payload) when is_atom(Type), is_integer(StreamId), is_binary(Payload) ->
|
||||
TypeNo = type_no(Type),
|
||||
PayloadLen = byte_size(Payload),
|
||||
Head = <<"RKP1", ?VERSION:8, TypeNo:8, 0:16, StreamId:64, PayloadLen:32>>,
|
||||
case chacha20_cipher:encrypt(Head, ?CHACHA20_KEY) of
|
||||
{ok, EncryptHead} ->
|
||||
{ok, <<EncryptHead/binary, Payload/binary>>};
|
||||
Error ->
|
||||
Error
|
||||
end.
|
||||
|
||||
decode_open_request(<<HostLen:16, Host:HostLen/binary, Port:16,
|
||||
UsernameLen:16, Username:UsernameLen/binary,
|
||||
PasswordLen:16, Password:PasswordLen/binary>>)
|
||||
when HostLen > 0, Port > 0 ->
|
||||
{ok, #{
|
||||
host => Host,
|
||||
port => Port,
|
||||
username => Username,
|
||||
password => Password
|
||||
}};
|
||||
decode_open_request(_Payload) ->
|
||||
{error, invalid_open_request}.
|
||||
|
||||
encode_open_request(Host, Port) ->
|
||||
encode_open_request(Host, Port, <<>>, <<>>).
|
||||
|
||||
encode_open_request(Host, Port, Username, Password) when
|
||||
is_binary(Host),
|
||||
byte_size(Host) > 0,
|
||||
byte_size(Host) =< 65535,
|
||||
is_integer(Port),
|
||||
Port > 0,
|
||||
Port =< 65535,
|
||||
is_binary(Username),
|
||||
byte_size(Username) =< 65535,
|
||||
is_binary(Password),
|
||||
byte_size(Password) =< 65535 ->
|
||||
HostLen = byte_size(Host),
|
||||
UsernameLen = byte_size(Username),
|
||||
PasswordLen = byte_size(Password),
|
||||
<<HostLen:16, Host/binary, Port:16,
|
||||
UsernameLen:16, Username/binary,
|
||||
PasswordLen:16, Password/binary>>.
|
||||
|
||||
type(1) -> open;
|
||||
type(2) -> open_ack;
|
||||
type(3) -> data;
|
||||
type(4) -> close;
|
||||
type(5) -> error;
|
||||
type(_) -> undefined.
|
||||
|
||||
type_no(open) -> 1;
|
||||
type_no(open_ack) -> 2;
|
||||
type_no(data) -> 3;
|
||||
type_no(close) -> 4;
|
||||
type_no(error) -> 5.
|
||||
122
apps/relay_server/src/relay_server_udp_stream.erl
Normal file
122
apps/relay_server/src/relay_server_udp_stream.erl
Normal file
@ -0,0 +1,122 @@
|
||||
%%%-------------------------------------------------------------------
|
||||
%% @doc Per-stream TCP relay worker.
|
||||
%% @end
|
||||
%%%-------------------------------------------------------------------
|
||||
|
||||
-module(relay_server_udp_stream).
|
||||
|
||||
-behaviour(gen_server).
|
||||
|
||||
-export([start_link/4, send/2, close/1]).
|
||||
|
||||
-export([init/1, handle_continue/2, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]).
|
||||
|
||||
-define(TCP_OPTIONS(Timeout), [
|
||||
binary,
|
||||
{packet, raw},
|
||||
{active, false},
|
||||
{nodelay, true},
|
||||
{send_timeout, Timeout},
|
||||
{send_timeout_close, true}
|
||||
]).
|
||||
|
||||
-record(state, {
|
||||
peer :: pid(),
|
||||
stream_id :: non_neg_integer(),
|
||||
host :: binary(),
|
||||
port :: inet:port_number(),
|
||||
connect_timeout :: timeout(),
|
||||
socket = undefined :: undefined | inet:socket()
|
||||
}).
|
||||
|
||||
start_link(Peer, StreamId, #{host := Host, port := Port}, ConnectTimeout) ->
|
||||
gen_server:start_link(?MODULE, [Peer, StreamId, Host, Port, ConnectTimeout], []).
|
||||
|
||||
send(StreamPid, Payload) when is_pid(StreamPid), is_binary(Payload) ->
|
||||
gen_server:cast(StreamPid, {send, Payload}).
|
||||
|
||||
close(StreamPid) when is_pid(StreamPid) ->
|
||||
gen_server:cast(StreamPid, close).
|
||||
|
||||
init([Peer, StreamId, Host, Port, ConnectTimeout]) ->
|
||||
{ok, #state{
|
||||
peer = Peer,
|
||||
stream_id = StreamId,
|
||||
host = Host,
|
||||
port = Port,
|
||||
connect_timeout = ConnectTimeout
|
||||
}, {continue, connect}}.
|
||||
|
||||
handle_call(_Request, _From, State) ->
|
||||
{reply, {error, bad_request}, State}.
|
||||
|
||||
handle_cast({send, _Payload}, State = #state{socket = undefined}) ->
|
||||
{noreply, State};
|
||||
handle_cast({send, Payload}, State = #state{socket = Socket}) ->
|
||||
case gen_tcp:send(Socket, Payload) of
|
||||
ok ->
|
||||
{noreply, State};
|
||||
{error, Reason} ->
|
||||
notify_error({send_failed, Reason}, State),
|
||||
{stop, normal, State}
|
||||
end;
|
||||
handle_cast(close, State) ->
|
||||
{stop, normal, State};
|
||||
handle_cast(_Request, State) ->
|
||||
{noreply, State}.
|
||||
|
||||
handle_info({tcp, Socket, Data}, State = #state{socket = Socket}) ->
|
||||
notify_data(Data, State),
|
||||
case inet:setopts(Socket, [{active, once}]) of
|
||||
ok ->
|
||||
{noreply, State};
|
||||
{error, Reason} ->
|
||||
notify_error(Reason, State),
|
||||
{stop, normal, State}
|
||||
end;
|
||||
handle_info({tcp_closed, Socket}, State = #state{socket = Socket}) ->
|
||||
notify_closed(State),
|
||||
{stop, normal, State};
|
||||
handle_info({tcp_error, Socket, Reason}, State = #state{socket = Socket}) ->
|
||||
notify_error(Reason, State),
|
||||
{stop, normal, State};
|
||||
handle_info(_Info, State) ->
|
||||
{noreply, State}.
|
||||
|
||||
terminate(_Reason, #state{socket = undefined}) ->
|
||||
ok;
|
||||
terminate(_Reason, #state{socket = Socket}) ->
|
||||
gen_tcp:close(Socket),
|
||||
ok.
|
||||
|
||||
code_change(_OldVsn, State, _Extra) ->
|
||||
{ok, State}.
|
||||
|
||||
handle_continue(connect, State = #state{host = Host, port = Port, connect_timeout = ConnectTimeout}) ->
|
||||
case gen_tcp:connect(binary_to_list(Host), Port, ?TCP_OPTIONS(ConnectTimeout), ConnectTimeout) of
|
||||
{ok, Socket} ->
|
||||
case inet:setopts(Socket, [{active, once}]) of
|
||||
ok ->
|
||||
notify_opened(State),
|
||||
{noreply, State#state{socket = Socket}};
|
||||
{error, Reason} ->
|
||||
gen_tcp:close(Socket),
|
||||
notify_error(Reason, State),
|
||||
{stop, normal, State}
|
||||
end;
|
||||
{error, Reason} ->
|
||||
notify_error({connect_failed, Reason}, State),
|
||||
{stop, normal, State}
|
||||
end.
|
||||
|
||||
notify_opened(#state{peer = Peer, stream_id = StreamId}) ->
|
||||
Peer ! {stream_opened, StreamId, self()}.
|
||||
|
||||
notify_data(Data, #state{peer = Peer, stream_id = StreamId}) ->
|
||||
Peer ! {stream_data, StreamId, self(), Data}.
|
||||
|
||||
notify_closed(#state{peer = Peer, stream_id = StreamId}) ->
|
||||
Peer ! {stream_closed, StreamId, self()}.
|
||||
|
||||
notify_error(Reason, #state{peer = Peer, stream_id = StreamId}) ->
|
||||
Peer ! {stream_error, StreamId, self(), Reason}.
|
||||
@ -3,10 +3,15 @@
|
||||
{udp_server, [
|
||||
{enabled, true},
|
||||
{listener, 'relay_server/udp'},
|
||||
{port, 7000},
|
||||
{idle_timeout, 30000},
|
||||
{listen_on, 1443},
|
||||
{idle_timeout, 60000},
|
||||
{connect_timeout, 5000},
|
||||
{max_connections, 1024},
|
||||
{udp_options, [binary, {reuseaddr, true}]}
|
||||
]},
|
||||
|
||||
{users, [
|
||||
{admin, "v7@Qm!2z#R8$pL4^xT?K"}
|
||||
]}
|
||||
]},
|
||||
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user