blob: 173fa79d96fa9e383c70c3f860fa12f167ff6970 [file] [log] [blame]
David Reiss80862312008-06-11 00:59:55 +00001%%%-------------------------------------------------------------------
2%%% File : thrift_socket_server.erl
3%%% Author : eugene letuchy <eletuchy@facebook.com>
4%%% Description : A rewrite of thrift_server, based quite heavily
5%%% on the mochiweb_socket_server module of mochiweb
6%%% Created : 3 Mar 2008 by eugene letuchy <eletuchy@facebook.com>
7%%%-------------------------------------------------------------------
8-module(thrift_socket_server).
9
10-author('eletuchy@facebook.com').
11-author('todd@lipcon.org').
12
13-behaviour(gen_server).
14
15-export([start/1, stop/1]).
16
17-export([init/1, handle_call/3, handle_cast/2, terminate/2, code_change/3,
David Reiss1a2f2182008-06-11 01:14:01 +000018 handle_info/2]).
David Reiss80862312008-06-11 00:59:55 +000019
20-export([acceptor_loop/1]).
21
22-record(thrift_socket_server,
David Reiss1a2f2182008-06-11 01:14:01 +000023 {port,
David Reiss80862312008-06-11 00:59:55 +000024 service,
25 handler,
David Reiss1a2f2182008-06-11 01:14:01 +000026 name,
27 max=2048,
28 ip=any,
29 listen=null,
30 acceptor=null,
David Reissb7c88022008-06-11 01:00:20 +000031 socket_opts=[{recv_timeout, 500}]
32 }).
David Reiss80862312008-06-11 00:59:55 +000033
34start(State=#thrift_socket_server{}) ->
David Reiss80862312008-06-11 00:59:55 +000035 start_server(State);
36start(Options) ->
37 start(parse_options(Options)).
38
39stop(Name) when is_atom(Name) ->
40 gen_server:cast(Name, stop);
41stop(Pid) when is_pid(Pid) ->
42 gen_server:cast(Pid, stop);
43stop({local, Name}) ->
44 stop(Name);
45stop({global, Name}) ->
46 stop(Name);
47stop(Options) ->
48 State = parse_options(Options),
49 stop(State#thrift_socket_server.name).
50
51%% Internal API
52
53parse_options(Options) ->
54 parse_options(Options, #thrift_socket_server{}).
55
56parse_options([], State) ->
57 State;
58parse_options([{name, L} | Rest], State) when is_list(L) ->
59 Name = {local, list_to_atom(L)},
60 parse_options(Rest, State#thrift_socket_server{name=Name});
61parse_options([{name, A} | Rest], State) when is_atom(A) ->
62 Name = {local, A},
63 parse_options(Rest, State#thrift_socket_server{name=Name});
64parse_options([{name, Name} | Rest], State) ->
65 parse_options(Rest, State#thrift_socket_server{name=Name});
66parse_options([{port, L} | Rest], State) when is_list(L) ->
67 Port = list_to_integer(L),
68 parse_options(Rest, State#thrift_socket_server{port=Port});
69parse_options([{port, Port} | Rest], State) ->
70 parse_options(Rest, State#thrift_socket_server{port=Port});
71parse_options([{ip, Ip} | Rest], State) ->
72 ParsedIp = case Ip of
David Reiss1a2f2182008-06-11 01:14:01 +000073 any ->
74 any;
75 Ip when is_tuple(Ip) ->
76 Ip;
77 Ip when is_list(Ip) ->
78 {ok, IpTuple} = inet_parse:address(Ip),
79 IpTuple
80 end,
David Reiss80862312008-06-11 00:59:55 +000081 parse_options(Rest, State#thrift_socket_server{ip=ParsedIp});
David Reissb7c88022008-06-11 01:00:20 +000082parse_options([{socket_opts, L} | Rest], State) when is_list(L), length(L) > 0 ->
83 parse_options(Rest, State#thrift_socket_server{socket_opts=L});
David Reiss80862312008-06-11 00:59:55 +000084parse_options([{handler, Handler} | Rest], State) ->
85 parse_options(Rest, State#thrift_socket_server{handler=Handler});
86parse_options([{service, Service} | Rest], State) ->
87 parse_options(Rest, State#thrift_socket_server{service=Service});
88parse_options([{max, Max} | Rest], State) ->
89 MaxInt = case Max of
David Reiss1a2f2182008-06-11 01:14:01 +000090 Max when is_list(Max) ->
91 list_to_integer(Max);
92 Max when is_integer(Max) ->
93 Max
94 end,
David Reiss80862312008-06-11 00:59:55 +000095 parse_options(Rest, State#thrift_socket_server{max=MaxInt}).
96
97start_server(State=#thrift_socket_server{name=Name}) ->
David Reiss80862312008-06-11 00:59:55 +000098 case Name of
David Reiss1a2f2182008-06-11 01:14:01 +000099 undefined ->
100 gen_server:start_link(?MODULE, State, []);
101 _ ->
102 gen_server:start_link(Name, ?MODULE, State, [])
David Reiss80862312008-06-11 00:59:55 +0000103 end.
104
105init(State=#thrift_socket_server{ip=Ip, port=Port}) ->
David Reissd74b0232008-06-11 01:02:55 +0000106 process_flag(trap_exit, true),
David Reiss80862312008-06-11 00:59:55 +0000107 BaseOpts = [binary,
David Reiss1a2f2182008-06-11 01:14:01 +0000108 {reuseaddr, true},
109 {packet, 0},
110 {backlog, 4096},
111 {recbuf, 8192},
112 {active, false}],
David Reiss80862312008-06-11 00:59:55 +0000113 Opts = case Ip of
David Reiss1a2f2182008-06-11 01:14:01 +0000114 any ->
David Reiss80862312008-06-11 00:59:55 +0000115 BaseOpts;
David Reiss1a2f2182008-06-11 01:14:01 +0000116 Ip ->
117 [{ip, Ip} | BaseOpts]
118 end,
David Reiss80862312008-06-11 00:59:55 +0000119 case gen_tcp_listen(Port, Opts, State) of
120 {stop, eacces} ->
121 %% fdsrv module allows another shot to bind
122 %% ports which require root access
123 case Port < 1024 of
124 true ->
125 case fdsrv:start() of
126 {ok, _} ->
127 case fdsrv:bind_socket(tcp, Port) of
128 {ok, Fd} ->
129 gen_tcp_listen(Port, [{fd, Fd} | Opts], State);
130 _ ->
131 {stop, fdsrv_bind_failed}
132 end;
133 _ ->
134 {stop, fdsrv_start_failed}
135 end;
136 false ->
137 {stop, eacces}
138 end;
139 Other ->
140 error_logger:info_msg("thrift service listening on port ~p", [Port]),
141 Other
142 end.
143
144gen_tcp_listen(Port, Opts, State) ->
145 case gen_tcp:listen(Port, Opts) of
146 {ok, Listen} ->
David Reiss1a2f2182008-06-11 01:14:01 +0000147 {ok, ListenPort} = inet:port(Listen),
148 {ok, new_acceptor(State#thrift_socket_server{listen=Listen,
David Reiss80862312008-06-11 00:59:55 +0000149 port=ListenPort})};
David Reiss1a2f2182008-06-11 01:14:01 +0000150 {error, Reason} ->
151 {stop, Reason}
David Reiss80862312008-06-11 00:59:55 +0000152 end.
153
154new_acceptor(State=#thrift_socket_server{max=0}) ->
155 error_logger:error_msg("Not accepting new connections"),
156 State#thrift_socket_server{acceptor=null};
David Reissb7c88022008-06-11 01:00:20 +0000157new_acceptor(State=#thrift_socket_server{acceptor=OldPid, listen=Listen,
158 service=Service, handler=Handler,
159 socket_opts=Opts
160 }) ->
David Reiss80862312008-06-11 00:59:55 +0000161 Pid = proc_lib:spawn_link(?MODULE, acceptor_loop,
David Reissb7c88022008-06-11 01:00:20 +0000162 [{self(), Listen, Service, Handler, Opts}]),
David Reiss919a8012008-06-11 01:00:12 +0000163%% error_logger:info_msg("Spawning new acceptor: ~p => ~p", [OldPid, Pid]),
David Reiss80862312008-06-11 00:59:55 +0000164 State#thrift_socket_server{acceptor=Pid}.
165
David Reissb7c88022008-06-11 01:00:20 +0000166acceptor_loop({Server, Listen, Service, Handler, SocketOpts})
167 when is_pid(Server), is_list(SocketOpts) ->
David Reissd74b0232008-06-11 01:02:55 +0000168 case catch gen_tcp:accept(Listen) of % infinite timeout
David Reiss1a2f2182008-06-11 01:14:01 +0000169 {ok, Socket} ->
170 gen_server:cast(Server, {accepted, self()}),
David Reiss80862312008-06-11 00:59:55 +0000171 ProtoGen = fun() ->
David Reissb7c88022008-06-11 01:00:20 +0000172 {ok, SocketTransport} = thrift_socket_transport:new(Socket, SocketOpts),
David Reiss80862312008-06-11 00:59:55 +0000173 {ok, BufferedTransport} = thrift_buffered_transport:new(SocketTransport),
174 {ok, Protocol} = thrift_binary_protocol:new(BufferedTransport),
175 {ok, IProt=Protocol, OProt=Protocol}
176 end,
177 thrift_processor:init({Server, ProtoGen, Service, Handler});
David Reiss1a2f2182008-06-11 01:14:01 +0000178 {error, closed} ->
179 exit({error, closed});
180 Other ->
181 error_logger:error_report(
182 [{application, thrift},
183 "Accept failed error",
184 lists:flatten(io_lib:format("~p", [Other]))]),
185 exit({error, accept_failed})
David Reiss80862312008-06-11 00:59:55 +0000186 end.
187
188handle_call({get, port}, _From, State=#thrift_socket_server{port=Port}) ->
189 {reply, Port, State};
190handle_call(_Message, _From, State) ->
191 Res = error,
192 {reply, Res, State}.
193
194handle_cast({accepted, Pid},
David Reiss1a2f2182008-06-11 01:14:01 +0000195 State=#thrift_socket_server{acceptor=Pid, max=Max}) ->
David Reiss80862312008-06-11 00:59:55 +0000196 % io:format("accepted ~p~n", [Pid]),
197 State1 = State#thrift_socket_server{max=Max - 1},
198 {noreply, new_acceptor(State1)};
199handle_cast(stop, State) ->
200 {stop, normal, State}.
201
202terminate(_Reason, #thrift_socket_server{listen=Listen, port=Port}) ->
203 gen_tcp:close(Listen),
204 case Port < 1024 of
205 true ->
206 catch fdsrv:stop(),
207 ok;
208 false ->
209 ok
210 end.
211
212code_change(_OldVsn, State, _Extra) ->
213 State.
214
215handle_info({'EXIT', Pid, normal},
David Reiss1a2f2182008-06-11 01:14:01 +0000216 State=#thrift_socket_server{acceptor=Pid}) ->
David Reiss80862312008-06-11 00:59:55 +0000217 {noreply, new_acceptor(State)};
218handle_info({'EXIT', Pid, Reason},
David Reiss1a2f2182008-06-11 01:14:01 +0000219 State=#thrift_socket_server{acceptor=Pid}) ->
David Reiss80862312008-06-11 00:59:55 +0000220 error_logger:error_report({?MODULE, ?LINE,
David Reiss1a2f2182008-06-11 01:14:01 +0000221 {acceptor_error, Reason}}),
David Reiss80862312008-06-11 00:59:55 +0000222 timer:sleep(100),
223 {noreply, new_acceptor(State)};
224handle_info({'EXIT', _LoopPid, Reason},
David Reiss1a2f2182008-06-11 01:14:01 +0000225 State=#thrift_socket_server{acceptor=Pid, max=Max}) ->
David Reiss80862312008-06-11 00:59:55 +0000226 case Reason of
David Reiss1a2f2182008-06-11 01:14:01 +0000227 normal -> ok;
David Reiss4cf5a6a2008-06-11 01:00:59 +0000228 shutdown -> ok;
David Reiss1a2f2182008-06-11 01:14:01 +0000229 _ -> error_logger:error_report({?MODULE, ?LINE,
David Reiss80862312008-06-11 00:59:55 +0000230 {child_error, Reason, erlang:get_stacktrace()}})
231 end,
232 State1 = State#thrift_socket_server{max=Max + 1},
233 State2 = case Pid of
David Reiss1a2f2182008-06-11 01:14:01 +0000234 null -> new_acceptor(State1);
235 _ -> State1
236 end,
David Reiss80862312008-06-11 00:59:55 +0000237 {noreply, State2};
238handle_info(Info, State) ->
239 error_logger:info_report([{'INFO', Info}, {'State', State}]),
240 {noreply, State}.