blob: c554b62ec9bf8acb88ce0d3706f7a47bfc9a521e [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,
18 handle_info/2]).
19
20-export([acceptor_loop/1]).
21
22-record(thrift_socket_server,
23 {port,
24 service,
25 handler,
26 name,
27 max=2048,
28 ip=any,
29 listen=null,
30 acceptor=null}).
31
32start(State=#thrift_socket_server{}) ->
33 io:format("~p~n", [State]),
34 start_server(State);
35start(Options) ->
36 start(parse_options(Options)).
37
38stop(Name) when is_atom(Name) ->
39 gen_server:cast(Name, stop);
40stop(Pid) when is_pid(Pid) ->
41 gen_server:cast(Pid, stop);
42stop({local, Name}) ->
43 stop(Name);
44stop({global, Name}) ->
45 stop(Name);
46stop(Options) ->
47 State = parse_options(Options),
48 stop(State#thrift_socket_server.name).
49
50%% Internal API
51
52parse_options(Options) ->
53 parse_options(Options, #thrift_socket_server{}).
54
55parse_options([], State) ->
56 State;
57parse_options([{name, L} | Rest], State) when is_list(L) ->
58 Name = {local, list_to_atom(L)},
59 parse_options(Rest, State#thrift_socket_server{name=Name});
60parse_options([{name, A} | Rest], State) when is_atom(A) ->
61 Name = {local, A},
62 parse_options(Rest, State#thrift_socket_server{name=Name});
63parse_options([{name, Name} | Rest], State) ->
64 parse_options(Rest, State#thrift_socket_server{name=Name});
65parse_options([{port, L} | Rest], State) when is_list(L) ->
66 Port = list_to_integer(L),
67 parse_options(Rest, State#thrift_socket_server{port=Port});
68parse_options([{port, Port} | Rest], State) ->
69 parse_options(Rest, State#thrift_socket_server{port=Port});
70parse_options([{ip, Ip} | Rest], State) ->
71 ParsedIp = case Ip of
72 any ->
73 any;
74 Ip when is_tuple(Ip) ->
75 Ip;
76 Ip when is_list(Ip) ->
77 {ok, IpTuple} = inet_parse:address(Ip),
78 IpTuple
79 end,
80 parse_options(Rest, State#thrift_socket_server{ip=ParsedIp});
81parse_options([{handler, Handler} | Rest], State) ->
82 parse_options(Rest, State#thrift_socket_server{handler=Handler});
83parse_options([{service, Service} | Rest], State) ->
84 parse_options(Rest, State#thrift_socket_server{service=Service});
85parse_options([{max, Max} | Rest], State) ->
86 MaxInt = case Max of
87 Max when is_list(Max) ->
88 list_to_integer(Max);
89 Max when is_integer(Max) ->
90 Max
91 end,
92 parse_options(Rest, State#thrift_socket_server{max=MaxInt}).
93
94start_server(State=#thrift_socket_server{name=Name}) ->
95 io:format("starting~n"),
96 case Name of
97 undefined ->
98 gen_server:start_link(?MODULE, State, []);
99 _ ->
100 gen_server:start_link(Name, ?MODULE, State, [])
101 end.
102
103init(State=#thrift_socket_server{ip=Ip, port=Port}) ->
104 process_flag(trap_exit, true), %% only temporary
105 BaseOpts = [binary,
106 {reuseaddr, true},
107 {packet, 0},
108 {backlog, 30},
109 {recbuf, 8192},
110 {active, false}],
111 Opts = case Ip of
112 any ->
113 BaseOpts;
114 Ip ->
115 [{ip, Ip} | BaseOpts]
116 end,
117 case gen_tcp_listen(Port, Opts, State) of
118 {stop, eacces} ->
119 %% fdsrv module allows another shot to bind
120 %% ports which require root access
121 case Port < 1024 of
122 true ->
123 case fdsrv:start() of
124 {ok, _} ->
125 case fdsrv:bind_socket(tcp, Port) of
126 {ok, Fd} ->
127 gen_tcp_listen(Port, [{fd, Fd} | Opts], State);
128 _ ->
129 {stop, fdsrv_bind_failed}
130 end;
131 _ ->
132 {stop, fdsrv_start_failed}
133 end;
134 false ->
135 {stop, eacces}
136 end;
137 Other ->
138 error_logger:info_msg("thrift service listening on port ~p", [Port]),
139 Other
140 end.
141
142gen_tcp_listen(Port, Opts, State) ->
143 case gen_tcp:listen(Port, Opts) of
144 {ok, Listen} ->
145 {ok, ListenPort} = inet:port(Listen),
146 {ok, new_acceptor(State#thrift_socket_server{listen=Listen,
147 port=ListenPort})};
148 {error, Reason} ->
149 {stop, Reason}
150 end.
151
152new_acceptor(State=#thrift_socket_server{max=0}) ->
153 error_logger:error_msg("Not accepting new connections"),
154 State#thrift_socket_server{acceptor=null};
155new_acceptor(State=#thrift_socket_server{acceptor=OldPid, listen=Listen,service=Service, handler=Handler}) ->
156 Pid = proc_lib:spawn_link(?MODULE, acceptor_loop,
157 [{self(), Listen, Service, Handler}]),
158 error_logger:info_msg("Spawning new acceptor: ~p => ~p", [OldPid, Pid]),
159 State#thrift_socket_server{acceptor=Pid}.
160
161acceptor_loop({Server, Listen, Service, Handler})
162 when is_pid(Server) ->
163 case catch gen_tcp:accept(Listen) of
164 {ok, Socket} ->
165 gen_server:cast(Server, {accepted, self()}),
166 ProtoGen = fun() ->
167 {ok, SocketTransport} = thrift_socket_transport:new(Socket),
168 {ok, BufferedTransport} = thrift_buffered_transport:new(SocketTransport),
169 {ok, Protocol} = thrift_binary_protocol:new(BufferedTransport),
170 {ok, IProt=Protocol, OProt=Protocol}
171 end,
172 thrift_processor:init({Server, ProtoGen, Service, Handler});
173 {error, closed} ->
174 exit({error, closed});
175 Other ->
176 error_logger:error_report(
177 [{application, thrift},
178 "Accept failed error",
179 lists:flatten(io_lib:format("~p", [Other]))]),
180 exit({error, accept_failed})
181 end.
182
183handle_call({get, port}, _From, State=#thrift_socket_server{port=Port}) ->
184 {reply, Port, State};
185handle_call(_Message, _From, State) ->
186 Res = error,
187 {reply, Res, State}.
188
189handle_cast({accepted, Pid},
190 State=#thrift_socket_server{acceptor=Pid, max=Max}) ->
191 % io:format("accepted ~p~n", [Pid]),
192 State1 = State#thrift_socket_server{max=Max - 1},
193 {noreply, new_acceptor(State1)};
194handle_cast(stop, State) ->
195 {stop, normal, State}.
196
197terminate(_Reason, #thrift_socket_server{listen=Listen, port=Port}) ->
198 gen_tcp:close(Listen),
199 case Port < 1024 of
200 true ->
201 catch fdsrv:stop(),
202 ok;
203 false ->
204 ok
205 end.
206
207code_change(_OldVsn, State, _Extra) ->
208 State.
209
210handle_info({'EXIT', Pid, normal},
211 State=#thrift_socket_server{acceptor=Pid}) ->
212 {noreply, new_acceptor(State)};
213handle_info({'EXIT', Pid, Reason},
214 State=#thrift_socket_server{acceptor=Pid}) ->
215 error_logger:error_report({?MODULE, ?LINE,
216 {acceptor_error, Reason}}),
217 timer:sleep(100),
218 {noreply, new_acceptor(State)};
219handle_info({'EXIT', _LoopPid, Reason},
220 State=#thrift_socket_server{acceptor=Pid, max=Max}) ->
221 case Reason of
222 normal -> ok;
223 protocol_closed -> ok;
224 _ -> error_logger:error_report({?MODULE, ?LINE,
225 {child_error, Reason, erlang:get_stacktrace()}})
226 end,
227 State1 = State#thrift_socket_server{max=Max + 1},
228 State2 = case Pid of
229 null -> new_acceptor(State1);
230 _ -> State1
231 end,
232 {noreply, State2};
233handle_info(Info, State) ->
234 error_logger:info_report([{'INFO', Info}, {'State', State}]),
235 {noreply, State}.