Skip to content

Commit 77615c0

Browse files
committed
Add shackle_socket, a socket NIF based protocol
An alternative to shackle_tcp built on the socket module. Accepted sockets run with {otp, select_read}, so the server gets one $socket message plus one recv per wake instead of the inet driver's active mode delivery, and send is a single NIF call instead of a port command with its monitor and inet_reply round trip. Benchmarked at 64 callers / pool of 16 against the arithmetic test server: 168-169k requests/s with shackle_tcp, 172-174k with shackle_socket (+2.5%). Requires OTP 27.3 for select_read; the default protocol is unchanged.
1 parent 60e996b commit 77615c0

4 files changed

Lines changed: 263 additions & 2 deletions

File tree

src/shackle.erl

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -22,10 +22,10 @@
2222
-type external_request_id() :: term().
2323
-type inet_address() :: inet:ip_address() | inet:hostname().
2424
-type inet_port() :: inet:port_number().
25-
-type protocol() :: shackle_ssl| shackle_tcp | shackle_udp.
25+
-type protocol() :: shackle_socket | shackle_ssl | shackle_tcp | shackle_udp.
2626
-type request_id() :: {shackle_server:name(), reference()}.
2727
-type response() :: {external_request_id(), term()}.
28-
-type socket() :: inet:socket() | ssl:sslsocket().
28+
-type socket() :: inet:socket() | socket:socket() | ssl:sslsocket().
2929
-type socket_option() :: gen_tcp:connect_option() | gen_udp:option() | ssl:tls_client_option().
3030
-type socket_options() :: [socket_option()].
3131
-type table() :: atom().

src/shackle_server.erl

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -139,6 +139,24 @@ handle_msg({Request, #cast {
139139
reply({error, client_crash}, Cast, State),
140140
{ok, {State, ClientState}}
141141
end;
142+
handle_msg({'$socket', Socket, select, _Handle}, {#state {
143+
socket = Socket
144+
} = State, ClientState}) ->
145+
146+
case shackle_socket:recv(Socket) of
147+
{ok, Data} ->
148+
handle_msg_data(Socket, Data, State, ClientState);
149+
wait ->
150+
{ok, {State, ClientState}};
151+
{error, closed} ->
152+
handle_msg_close(Socket, State, ClientState);
153+
{error, Reason} ->
154+
handle_msg_error(Socket, Reason, State, ClientState)
155+
end;
156+
handle_msg({'$socket', _Socket, select, _Handle}, {State, ClientState}) ->
157+
{ok, {State, ClientState}};
158+
handle_msg({'$socket', Socket, abort, _Info}, {State, ClientState}) ->
159+
handle_msg_close(Socket, State, ClientState);
142160
handle_msg({ssl, Socket, Data}, {State, ClientState}) ->
143161
handle_msg_data(Socket, Data, State, ClientState);
144162
handle_msg({ssl_closed, Socket}, {State, ClientState}) ->

src/shackle_socket.erl

Lines changed: 142 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,142 @@
1+
-module(shackle_socket).
2+
-include("shackle_internal.hrl").
3+
4+
-compile(inline).
5+
-compile({inline_size, 512}).
6+
7+
%% The socket specs on OTP < 28 predate {otp, select_read} and its
8+
%% recv return shapes, although both work at runtime from OTP 27.3.
9+
-dialyzer({nowarn_function, [recv/1, setopts/2]}).
10+
11+
-behavior(shackle_protocol).
12+
-export([
13+
close/1,
14+
connect/3,
15+
send/2,
16+
setopts/2
17+
]).
18+
19+
%% internal
20+
-export([
21+
recv/1
22+
]).
23+
24+
%% callbacks
25+
-spec close(shackle:socket()) ->
26+
ok.
27+
28+
close(Socket) ->
29+
_ = socket:close(Socket),
30+
ok.
31+
32+
-spec connect(shackle:inet_address(), shackle:inet_port(), shackle:socket_options()) ->
33+
{ok, shackle:socket()} | {error, atom()}.
34+
35+
connect(Address, Port, SocketOptions) ->
36+
case socket:open(inet, stream, tcp) of
37+
{ok, Socket} ->
38+
case connect_opts(Socket, SocketOptions) of
39+
ok ->
40+
SockAddr = #{family => inet, addr => Address, port => Port},
41+
case socket:connect(Socket, SockAddr,
42+
?DEFAULT_CONNECT_TIMEOUT) of
43+
44+
ok ->
45+
{ok, Socket};
46+
{error, Reason} ->
47+
_ = socket:close(Socket),
48+
{error, Reason}
49+
end;
50+
{error, Reason} ->
51+
_ = socket:close(Socket),
52+
{error, Reason}
53+
end;
54+
{error, _} = Error ->
55+
Error
56+
end.
57+
58+
%% Reads whatever is available; with {otp, select_read} enabled the
59+
%% successful recv re-arms the read select inside the same NIF call.
60+
-spec recv(shackle:socket()) ->
61+
{ok, binary()} | wait | {error, atom()}.
62+
63+
recv(Socket) ->
64+
case socket:recv(Socket, 0, [], nowait) of
65+
{select_read, {_SelectInfo, Data}} ->
66+
{ok, Data};
67+
{select, {_SelectInfo, Data}} ->
68+
{ok, Data};
69+
{select, _SelectInfo} ->
70+
wait;
71+
{ok, Data} ->
72+
%% select_read did not re-arm; force another recv round
73+
self() ! {'$socket', Socket, select, undefined},
74+
{ok, Data};
75+
{error, {Reason, _Data}} ->
76+
{error, Reason};
77+
{error, Reason} ->
78+
{error, Reason}
79+
end.
80+
81+
-spec send(shackle:socket(), iodata()) ->
82+
ok | {error, atom()}.
83+
84+
send(Socket, Data) ->
85+
case socket:send(Socket, Data) of
86+
ok ->
87+
ok;
88+
{error, {Reason, _RestData}} ->
89+
{error, Reason};
90+
{error, Reason} ->
91+
{error, Reason}
92+
end.
93+
94+
-spec setopts(shackle:socket(), [gen_tcp:option()]) ->
95+
ok |
96+
{error, atom()}.
97+
98+
setopts(Socket, [{active, false}]) ->
99+
_ = socket:setopt(Socket, {otp, select_read}, false),
100+
ok;
101+
setopts(Socket, [{active, true}]) ->
102+
ok = socket:setopt(Socket, {otp, select_read}, true),
103+
case recv(Socket) of
104+
{ok, Data} ->
105+
%% shackle_server delivers this to the client like any
106+
%% active-mode packet
107+
self() ! {tcp, Socket, Data},
108+
ok;
109+
wait ->
110+
ok;
111+
{error, _} = Error ->
112+
Error
113+
end;
114+
setopts(Socket, Opts) ->
115+
connect_opts(Socket, Opts).
116+
117+
%% private
118+
connect_opts(_Socket, []) ->
119+
ok;
120+
connect_opts(Socket, [{nodelay, Bool} | T]) ->
121+
case socket:setopt(Socket, {tcp, nodelay}, Bool) of
122+
ok ->
123+
connect_opts(Socket, T);
124+
{error, _} = Error ->
125+
Error
126+
end;
127+
connect_opts(Socket, [{recbuf, Size} | T]) ->
128+
case socket:setopt(Socket, {socket, rcvbuf}, Size) of
129+
ok ->
130+
connect_opts(Socket, T);
131+
{error, _} = Error ->
132+
Error
133+
end;
134+
connect_opts(Socket, [{sndbuf, Size} | T]) ->
135+
case socket:setopt(Socket, {socket, sndbuf}, Size) of
136+
ok ->
137+
connect_opts(Socket, T);
138+
{error, _} = Error ->
139+
Error
140+
end;
141+
connect_opts(Socket, [_ | T]) ->
142+
connect_opts(Socket, T).

test/arithmetic_socket_client.erl

Lines changed: 101 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,101 @@
1+
-module(arithmetic_socket_client).
2+
-include("test.hrl").
3+
4+
-export([
5+
add/2,
6+
start/0,
7+
start/1,
8+
stop/0
9+
]).
10+
11+
-behavior(shackle_client).
12+
-export([
13+
init/1,
14+
setup/2,
15+
handle_request/2,
16+
handle_data/2,
17+
terminate/1
18+
]).
19+
20+
-record(state, {
21+
buffer = <<>>,
22+
request_counter = 0
23+
}).
24+
25+
-type tiny_int() :: 0..255.
26+
27+
%% public
28+
-spec add(tiny_int(), tiny_int()) ->
29+
pos_integer().
30+
31+
add(A, B) ->
32+
shackle:call(?POOL_NAME, {add, A, B}, ?TIMEOUT).
33+
34+
-spec start() ->
35+
ok | {error, shackle_not_started | pool_already_started}.
36+
37+
start() ->
38+
start([
39+
{backlog_size, ?BACKLOG_SIZE},
40+
{pool_size, 1}
41+
]).
42+
43+
-spec start(shackle_pool:options()) ->
44+
ok | {error, shackle_not_started | pool_already_started}.
45+
46+
start(PoolOptions) ->
47+
shackle_pool:start(?POOL_NAME, ?MODULE, [
48+
{port, ?PORT},
49+
{protocol, shackle_socket},
50+
{reconnect, true},
51+
{reconnect_time_min, 1},
52+
{socket_options, []}
53+
], PoolOptions).
54+
55+
-spec stop() ->
56+
ok | {error, pool_not_started}.
57+
58+
stop() ->
59+
shackle_pool:stop(?POOL_NAME).
60+
61+
%% shackle_server callbacks
62+
init(_) ->
63+
{ok, #state {}}.
64+
65+
setup(Socket, State) ->
66+
case socket:send(Socket, <<"INIT">>) of
67+
ok ->
68+
case socket:recv(Socket, 0, ?TIMEOUT) of
69+
{ok, <<"OK">>} ->
70+
{ok, State};
71+
{error, Reason} ->
72+
{error, Reason, State}
73+
end;
74+
{error, Reason} ->
75+
{error, Reason, State}
76+
end.
77+
78+
handle_data(Data, #state {
79+
buffer = Buffer
80+
} = State) ->
81+
82+
Data2 = <<Buffer/binary, Data/binary>>,
83+
{Replies, Buffer2} = arithmetic_protocol:parse_replies(Data2),
84+
85+
{ok, Replies, State#state {
86+
buffer = Buffer2
87+
}}.
88+
89+
handle_request({Operation, A, B}, #state {
90+
request_counter = RequestCounter
91+
} = State) ->
92+
93+
RequestId = arithmetic_protocol:request_id(RequestCounter),
94+
Data = arithmetic_protocol:request(RequestId, Operation, A, B),
95+
96+
{ok, RequestId, Data, State#state {
97+
request_counter = RequestCounter + 1
98+
}}.
99+
100+
terminate(_State) ->
101+
ok.

0 commit comments

Comments
 (0)