blob: 62bdfdaf6320aaf5f3c8db9b69f23cf9f4124761 [file] [log] [blame]
Gavin McDonald0b75e1a2010-10-28 02:12:01 +00001%%
2%% Licensed to the Apache Software Foundation (ASF) under one
3%% or more contributor license agreements. See the NOTICE file
4%% distributed with this work for additional information
5%% regarding copyright ownership. The ASF licenses this file
6%% to you under the Apache License, Version 2.0 (the
7%% "License"); you may not use this file except in compliance
8%% with the License. You may obtain a copy of the License at
9%%
10%% http://www.apache.org/licenses/LICENSE-2.0
11%%
12%% Unless required by applicable law or agreed to in writing,
13%% software distributed under the License is distributed on an
14%% "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15%% KIND, either express or implied. See the License for the
16%% specific language governing permissions and limitations
17%% under the License.
18%%
19
20-module(thrift_socket_server).
21
22-behaviour(gen_server).
23
24-export([start/1, stop/1]).
25
26-export([init/1, handle_call/3, handle_cast/2, terminate/2, code_change/3,
27 handle_info/2]).
28
29-export([acceptor_loop/1]).
30
31-record(thrift_socket_server,
32 {port,
33 service,
34 handler,
35 name,
36 max=2048,
37 ip=any,
38 listen=null,
39 acceptor=null,
40 socket_opts=[{recv_timeout, 500}]
41 }).
42
43start(State=#thrift_socket_server{}) ->
44 start_server(State);
45start(Options) ->
46 start(parse_options(Options)).
47
48stop(Name) when is_atom(Name) ->
49 gen_server:cast(Name, stop);
50stop(Pid) when is_pid(Pid) ->
51 gen_server:cast(Pid, stop);
52stop({local, Name}) ->
53 stop(Name);
54stop({global, Name}) ->
55 stop(Name);
56stop(Options) ->
57 State = parse_options(Options),
58 stop(State#thrift_socket_server.name).
59
60%% Internal API
61
62parse_options(Options) ->
63 parse_options(Options, #thrift_socket_server{}).
64
65parse_options([], State) ->
66 State;
67parse_options([{name, L} | Rest], State) when is_list(L) ->
68 Name = {local, list_to_atom(L)},
69 parse_options(Rest, State#thrift_socket_server{name=Name});
70parse_options([{name, A} | Rest], State) when is_atom(A) ->
71 Name = {local, A},
72 parse_options(Rest, State#thrift_socket_server{name=Name});
73parse_options([{name, Name} | Rest], State) ->
74 parse_options(Rest, State#thrift_socket_server{name=Name});
75parse_options([{port, L} | Rest], State) when is_list(L) ->
76 Port = list_to_integer(L),
77 parse_options(Rest, State#thrift_socket_server{port=Port});
78parse_options([{port, Port} | Rest], State) ->
79 parse_options(Rest, State#thrift_socket_server{port=Port});
80parse_options([{ip, Ip} | Rest], State) ->
81 ParsedIp = case Ip of
82 any ->
83 any;
84 Ip when is_tuple(Ip) ->
85 Ip;
86 Ip when is_list(Ip) ->
87 {ok, IpTuple} = inet_parse:address(Ip),
88 IpTuple
89 end,
90 parse_options(Rest, State#thrift_socket_server{ip=ParsedIp});
91parse_options([{socket_opts, L} | Rest], State) when is_list(L), length(L) > 0 ->
92 parse_options(Rest, State#thrift_socket_server{socket_opts=L});
93parse_options([{handler, Handler} | Rest], State) ->
94 parse_options(Rest, State#thrift_socket_server{handler=Handler});
95parse_options([{service, Service} | Rest], State) ->
96 parse_options(Rest, State#thrift_socket_server{service=Service});
97parse_options([{max, Max} | Rest], State) ->
98 MaxInt = case Max of
99 Max when is_list(Max) ->
100 list_to_integer(Max);
101 Max when is_integer(Max) ->
102 Max
103 end,
104 parse_options(Rest, State#thrift_socket_server{max=MaxInt}).
105
106start_server(State=#thrift_socket_server{name=Name}) ->
107 case Name of
108 undefined ->
109 gen_server:start_link(?MODULE, State, []);
110 _ ->
111 gen_server:start_link(Name, ?MODULE, State, [])
112 end.
113
114init(State=#thrift_socket_server{ip=Ip, port=Port}) ->
115 process_flag(trap_exit, true),
116 BaseOpts = [binary,
117 {reuseaddr, true},
118 {packet, 0},
119 {backlog, 4096},
120 {recbuf, 8192},
121 {active, false}],
122 Opts = case Ip of
123 any ->
124 BaseOpts;
125 Ip ->
126 [{ip, Ip} | BaseOpts]
127 end,
128 case gen_tcp_listen(Port, Opts, State) of
129 {stop, eacces} ->
130 %% fdsrv module allows another shot to bind
131 %% ports which require root access
132 case Port < 1024 of
133 true ->
134 case fdsrv:start() of
135 {ok, _} ->
136 case fdsrv:bind_socket(tcp, Port) of
137 {ok, Fd} ->
138 gen_tcp_listen(Port, [{fd, Fd} | Opts], State);
139 _ ->
140 {stop, fdsrv_bind_failed}
141 end;
142 _ ->
143 {stop, fdsrv_start_failed}
144 end;
145 false ->
146 {stop, eacces}
147 end;
148 Other ->
149 error_logger:info_msg("thrift service listening on port ~p", [Port]),
150 Other
151 end.
152
153gen_tcp_listen(Port, Opts, State) ->
154 case gen_tcp:listen(Port, Opts) of
155 {ok, Listen} ->
156 {ok, ListenPort} = inet:port(Listen),
157 {ok, new_acceptor(State#thrift_socket_server{listen=Listen,
158 port=ListenPort})};
159 {error, Reason} ->
160 {stop, Reason}
161 end.
162
163new_acceptor(State=#thrift_socket_server{max=0}) ->
164 error_logger:error_msg("Not accepting new connections"),
165 State#thrift_socket_server{acceptor=null};
166new_acceptor(State=#thrift_socket_server{acceptor=OldPid, listen=Listen,
167 service=Service, handler=Handler,
168 socket_opts=Opts
169 }) ->
170 Pid = proc_lib:spawn_link(?MODULE, acceptor_loop,
171 [{self(), Listen, Service, Handler, Opts}]),
172%% error_logger:info_msg("Spawning new acceptor: ~p => ~p", [OldPid, Pid]),
173 State#thrift_socket_server{acceptor=Pid}.
174
175acceptor_loop({Server, Listen, Service, Handler, SocketOpts})
176 when is_pid(Server), is_list(SocketOpts) ->
177 case catch gen_tcp:accept(Listen) of % infinite timeout
178 {ok, Socket} ->
179 gen_server:cast(Server, {accepted, self()}),
180 ProtoGen = fun() ->
181 {ok, SocketTransport} = thrift_socket_transport:new(Socket, SocketOpts),
182 {ok, BufferedTransport} = thrift_buffered_transport:new(SocketTransport),
183 {ok, Protocol} = thrift_binary_protocol:new(BufferedTransport),
184 {ok, IProt=Protocol, OProt=Protocol}
185 end,
186 thrift_processor:init({Server, ProtoGen, Service, Handler});
187 {error, closed} ->
188 exit({error, closed});
189 Other ->
190 error_logger:error_report(
191 [{application, thrift},
192 "Accept failed error",
193 lists:flatten(io_lib:format("~p", [Other]))]),
194 exit({error, accept_failed})
195 end.
196
197handle_call({get, port}, _From, State=#thrift_socket_server{port=Port}) ->
198 {reply, Port, State};
199handle_call(_Message, _From, State) ->
200 Res = error,
201 {reply, Res, State}.
202
203handle_cast({accepted, Pid},
204 State=#thrift_socket_server{acceptor=Pid, max=Max}) ->
205 % io:format("accepted ~p~n", [Pid]),
206 State1 = State#thrift_socket_server{max=Max - 1},
207 {noreply, new_acceptor(State1)};
208handle_cast(stop, State) ->
209 {stop, normal, State}.
210
211terminate(_Reason, #thrift_socket_server{listen=Listen, port=Port}) ->
212 gen_tcp:close(Listen),
213 case Port < 1024 of
214 true ->
215 catch fdsrv:stop(),
216 ok;
217 false ->
218 ok
219 end.
220
221code_change(_OldVsn, State, _Extra) ->
222 State.
223
224handle_info({'EXIT', Pid, normal},
225 State=#thrift_socket_server{acceptor=Pid}) ->
226 {noreply, new_acceptor(State)};
227handle_info({'EXIT', Pid, Reason},
228 State=#thrift_socket_server{acceptor=Pid}) ->
229 error_logger:error_report({?MODULE, ?LINE,
230 {acceptor_error, Reason}}),
231 timer:sleep(100),
232 {noreply, new_acceptor(State)};
233handle_info({'EXIT', _LoopPid, Reason},
234 State=#thrift_socket_server{acceptor=Pid, max=Max}) ->
235 case Reason of
236 normal -> ok;
237 shutdown -> ok;
238 _ -> error_logger:error_report({?MODULE, ?LINE,
239 {child_error, Reason, erlang:get_stacktrace()}})
240 end,
241 State1 = State#thrift_socket_server{max=Max + 1},
242 State2 = case Pid of
243 null -> new_acceptor(State1);
244 _ -> State1
245 end,
246 {noreply, State2};
247handle_info(Info, State) ->
248 error_logger:info_report([{'INFO', Info}, {'State', State}]),
249 {noreply, State}.