Compare commits
No commits in common. "c106302d993010839d84a2ad856f3fcd6376f5b9" and "3cf38f3e89b0a0dc9ac9cc204cabc8151803d197" have entirely different histories.
c106302d99
...
3cf38f3e89
39
README.md
39
README.md
@ -13,39 +13,10 @@ Run
|
|||||||
|
|
||||||
$ rebar3 shell
|
$ rebar3 shell
|
||||||
|
|
||||||
The UDP listener is configured in `config/sys.config`; the checked-in config
|
The UDP listener is configured in `config/sys.config` and listens on port `7000`
|
||||||
listens on port `1443`, while the code fallback default is `16380`. It accepts
|
by default. It echoes received datagrams.
|
||||||
RelayKit UDP relay frames.
|
|
||||||
|
|
||||||
Protocol
|
Test
|
||||||
--------
|
----
|
||||||
|
|
||||||
Each UDP datagram is encrypted with ChaCha20-Poly1305. The encrypted datagram
|
$ echo "ping" | nc -u -w1 127.0.0.1 7000
|
||||||
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"}
|
|
||||||
]}
|
|
||||||
|
|||||||
@ -1,38 +0,0 @@
|
|||||||
%%%-------------------------------------------------------------------
|
|
||||||
%% @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,7 +9,6 @@
|
|||||||
{applications, [
|
{applications, [
|
||||||
kernel,
|
kernel,
|
||||||
stdlib,
|
stdlib,
|
||||||
crypto,
|
|
||||||
esockd,
|
esockd,
|
||||||
sync
|
sync
|
||||||
]},
|
]},
|
||||||
|
|||||||
32
apps/relay_server/src/relay_server_udp_handler.erl
Normal file
32
apps/relay_server/src/relay_server_udp_handler.erl
Normal file
@ -0,0 +1,32 @@
|
|||||||
|
%%%-------------------------------------------------------------------
|
||||||
|
%% @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,9 +13,8 @@
|
|||||||
|
|
||||||
-define(SERVER, ?MODULE).
|
-define(SERVER, ?MODULE).
|
||||||
-define(DEFAULT_LISTENER, 'relay_server/udp').
|
-define(DEFAULT_LISTENER, 'relay_server/udp').
|
||||||
-define(DEFAULT_LISTEN_ON, 1443).
|
-define(DEFAULT_LISTEN_ON, 7000).
|
||||||
-define(DEFAULT_IDLE_TIMEOUT, 60000).
|
-define(DEFAULT_IDLE_TIMEOUT, 30000).
|
||||||
-define(DEFAULT_CONNECT_TIMEOUT, 5000).
|
|
||||||
-define(DEFAULT_MAX_CONNECTIONS, 1024).
|
-define(DEFAULT_MAX_CONNECTIONS, 1024).
|
||||||
-define(DEFAULT_UDP_OPTIONS, [binary, {reuseaddr, true}]).
|
-define(DEFAULT_UDP_OPTIONS, [binary, {reuseaddr, true}]).
|
||||||
|
|
||||||
@ -35,15 +34,9 @@ init([]) ->
|
|||||||
Props = application:get_env(relay_server, udp_server, []),
|
Props = application:get_env(relay_server, udp_server, []),
|
||||||
|
|
||||||
Listener = proplists:get_value(listener, Props, ?DEFAULT_LISTENER),
|
Listener = proplists:get_value(listener, Props, ?DEFAULT_LISTENER),
|
||||||
ListenOn = proplists:get_value(
|
ListenOn = proplists:get_value(port, Props, ?DEFAULT_LISTEN_ON),
|
||||||
listen_on,
|
|
||||||
Props,
|
|
||||||
proplists:get_value(port, Props, ?DEFAULT_LISTEN_ON)
|
|
||||||
),
|
|
||||||
|
|
||||||
IdleTimeout = proplists:get_value(idle_timeout, Props, ?DEFAULT_IDLE_TIMEOUT),
|
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),
|
UdpOptions = proplists:get_value(udp_options, Props, ?DEFAULT_UDP_OPTIONS),
|
||||||
MaxConnections = proplists:get_value(max_connections, Props, ?DEFAULT_MAX_CONNECTIONS),
|
MaxConnections = proplists:get_value(max_connections, Props, ?DEFAULT_MAX_CONNECTIONS),
|
||||||
AccessRules = proplists:get_value(access_rules, Props, [{allow, all}]),
|
AccessRules = proplists:get_value(access_rules, Props, [{allow, all}]),
|
||||||
@ -55,7 +48,7 @@ init([]) ->
|
|||||||
],
|
],
|
||||||
Opts = maybe_add(max_conn_rate, Props, BaseOpts),
|
Opts = maybe_add(max_conn_rate, Props, BaseOpts),
|
||||||
|
|
||||||
MFA = {relay_server_udp_peer, start_link, [IdleTimeout, ConnectTimeout, Users]},
|
MFA = {relay_server_udp_handler, start_link, [IdleTimeout]},
|
||||||
case esockd:open_udp(Listener, ListenOn, Opts, MFA) of
|
case esockd:open_udp(Listener, ListenOn, Opts, MFA) of
|
||||||
{ok, Pid} ->
|
{ok, Pid} ->
|
||||||
logger:info("UDP listener ~p started on ~s",
|
logger:info("UDP listener ~p started on ~s",
|
||||||
|
|||||||
@ -1,288 +0,0 @@
|
|||||||
%%%-------------------------------------------------------------------
|
|
||||||
%% @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.
|
|
||||||
@ -1,112 +0,0 @@
|
|||||||
%%%-------------------------------------------------------------------
|
|
||||||
%% @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.
|
|
||||||
@ -1,122 +0,0 @@
|
|||||||
%%%-------------------------------------------------------------------
|
|
||||||
%% @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,15 +3,10 @@
|
|||||||
{udp_server, [
|
{udp_server, [
|
||||||
{enabled, true},
|
{enabled, true},
|
||||||
{listener, 'relay_server/udp'},
|
{listener, 'relay_server/udp'},
|
||||||
{listen_on, 1443},
|
{port, 7000},
|
||||||
{idle_timeout, 60000},
|
{idle_timeout, 30000},
|
||||||
{connect_timeout, 5000},
|
|
||||||
{max_connections, 1024},
|
{max_connections, 1024},
|
||||||
{udp_options, [binary, {reuseaddr, true}]}
|
{udp_options, [binary, {reuseaddr, true}]}
|
||||||
]},
|
|
||||||
|
|
||||||
{users, [
|
|
||||||
{admin, "v7@Qm!2z#R8$pL4^xT?K"}
|
|
||||||
]}
|
]}
|
||||||
]},
|
]},
|
||||||
|
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user