3 * Copyright 2015 gRPC authors.
5 * Licensed under the Apache License, Version 2.0 (the "License");
6 * you may not use this file except in compliance with the License.
7 * You may obtain a copy of the License at
9 * http://www.apache.org/licenses/LICENSE-2.0
11 * Unless required by applicable law or agreed to in writing, software
12 * distributed under the License is distributed on an "AS IS" BASIS,
13 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14 * See the License for the specific language governing permissions and
15 * limitations under the License.
19 #include <grpc/support/port_platform.h>
21 #include "src/core/ext/transport/chttp2/server/chttp2_server.h"
27 #include <grpc/grpc.h>
28 #include <grpc/impl/codegen/grpc_types.h>
29 #include <grpc/support/alloc.h>
30 #include <grpc/support/log.h>
31 #include <grpc/support/string_util.h>
32 #include <grpc/support/sync.h>
34 #include "absl/strings/str_format.h"
36 #include "src/core/ext/filters/http/server/http_server_filter.h"
37 #include "src/core/ext/transport/chttp2/transport/chttp2_transport.h"
38 #include "src/core/ext/transport/chttp2/transport/internal.h"
39 #include "src/core/lib/channel/channel_args.h"
40 #include "src/core/lib/channel/handshaker.h"
41 #include "src/core/lib/channel/handshaker_registry.h"
42 #include "src/core/lib/iomgr/endpoint.h"
43 #include "src/core/lib/iomgr/resolve_address.h"
44 #include "src/core/lib/iomgr/resource_quota.h"
45 #include "src/core/lib/iomgr/tcp_server.h"
46 #include "src/core/lib/slice/slice_internal.h"
47 #include "src/core/lib/surface/api_trace.h"
48 #include "src/core/lib/surface/server.h"
52 grpc_tcp_server* tcp_server;
53 grpc_channel_args* args;
56 grpc_closure tcp_server_shutdown_complete;
57 grpc_closure* server_destroy_listener_done;
58 grpc_core::HandshakeManager* pending_handshake_mgrs;
59 grpc_core::RefCountedPtr<grpc_core::channelz::ListenSocketNode>
60 channelz_listen_socket;
63 struct server_connection_state {
65 server_state* svr_state;
66 grpc_pollset* accepting_pollset;
67 grpc_tcp_server_acceptor* acceptor;
68 grpc_core::RefCountedPtr<grpc_core::HandshakeManager> handshake_mgr;
69 // State for enforcing handshake timeout on receiving HTTP/2 settings.
70 grpc_chttp2_transport* transport;
73 grpc_closure on_timeout;
74 grpc_closure on_receive_settings;
75 grpc_pollset_set* interested_parties;
78 static void server_connection_state_unref(
79 server_connection_state* connection_state) {
80 if (gpr_unref(&connection_state->refs)) {
81 if (connection_state->transport != nullptr) {
82 GRPC_CHTTP2_UNREF_TRANSPORT(connection_state->transport,
83 "receive settings timeout");
85 grpc_pollset_set_del_pollset(connection_state->interested_parties,
86 connection_state->accepting_pollset);
87 grpc_pollset_set_destroy(connection_state->interested_parties);
88 gpr_free(connection_state);
92 static void on_timeout(void* arg, grpc_error* error) {
93 server_connection_state* connection_state =
94 static_cast<server_connection_state*>(arg);
95 // Note that we may be called with GRPC_ERROR_NONE when the timer fires
96 // or with an error indicating that the timer system is being shut down.
97 if (error != GRPC_ERROR_CANCELLED) {
98 grpc_transport_op* op = grpc_make_transport_op(nullptr);
99 op->disconnect_with_error = GRPC_ERROR_CREATE_FROM_STATIC_STRING(
100 "Did not receive HTTP/2 settings before handshake timeout");
101 grpc_transport_perform_op(&connection_state->transport->base, op);
103 server_connection_state_unref(connection_state);
106 static void on_receive_settings(void* arg, grpc_error* error) {
107 server_connection_state* connection_state =
108 static_cast<server_connection_state*>(arg);
109 if (error == GRPC_ERROR_NONE) {
110 grpc_timer_cancel(&connection_state->timer);
112 server_connection_state_unref(connection_state);
115 static void on_handshake_done(void* arg, grpc_error* error) {
116 auto* args = static_cast<grpc_core::HandshakerArgs*>(arg);
117 server_connection_state* connection_state =
118 static_cast<server_connection_state*>(args->user_data);
119 gpr_mu_lock(&connection_state->svr_state->mu);
120 grpc_resource_user* resource_user = grpc_server_get_default_resource_user(
121 connection_state->svr_state->server);
122 if (error != GRPC_ERROR_NONE || connection_state->svr_state->shutdown) {
123 const char* error_str = grpc_error_string(error);
124 gpr_log(GPR_DEBUG, "Handshaking failed: %s", error_str);
125 grpc_resource_user* resource_user = grpc_server_get_default_resource_user(
126 connection_state->svr_state->server);
127 if (resource_user != nullptr) {
128 grpc_resource_user_free(resource_user, GRPC_RESOURCE_QUOTA_CHANNEL_SIZE);
130 if (error == GRPC_ERROR_NONE && args->endpoint != nullptr) {
131 // We were shut down after handshaking completed successfully, so
132 // destroy the endpoint here.
133 // TODO(ctiller): It is currently necessary to shutdown endpoints
134 // before destroying them, even if we know that there are no
135 // pending read/write callbacks. This should be fixed, at which
136 // point this can be removed.
137 grpc_endpoint_shutdown(args->endpoint, GRPC_ERROR_NONE);
138 grpc_endpoint_destroy(args->endpoint);
139 grpc_channel_args_destroy(args->args);
140 grpc_slice_buffer_destroy_internal(args->read_buffer);
141 gpr_free(args->read_buffer);
144 // If the handshaking succeeded but there is no endpoint, then the
145 // handshaker may have handed off the connection to some external
146 // code, so we can just clean up here without creating a transport.
147 if (args->endpoint != nullptr) {
148 grpc_transport* transport = grpc_create_chttp2_transport(
149 args->args, args->endpoint, false, resource_user);
150 grpc_server_setup_transport(
151 connection_state->svr_state->server, transport,
152 connection_state->accepting_pollset, args->args,
153 grpc_chttp2_transport_get_socket_node(transport), resource_user);
154 // Use notify_on_receive_settings callback to enforce the
155 // handshake deadline.
156 connection_state->transport =
157 reinterpret_cast<grpc_chttp2_transport*>(transport);
158 gpr_ref(&connection_state->refs);
159 GRPC_CLOSURE_INIT(&connection_state->on_receive_settings,
160 on_receive_settings, connection_state,
161 grpc_schedule_on_exec_ctx);
162 grpc_chttp2_transport_start_reading(
163 transport, args->read_buffer, &connection_state->on_receive_settings);
164 grpc_channel_args_destroy(args->args);
165 gpr_ref(&connection_state->refs);
166 GRPC_CHTTP2_REF_TRANSPORT((grpc_chttp2_transport*)transport,
167 "receive settings timeout");
168 GRPC_CLOSURE_INIT(&connection_state->on_timeout, on_timeout,
169 connection_state, grpc_schedule_on_exec_ctx);
170 grpc_timer_init(&connection_state->timer, connection_state->deadline,
171 &connection_state->on_timeout);
173 if (resource_user != nullptr) {
174 grpc_resource_user_free(resource_user,
175 GRPC_RESOURCE_QUOTA_CHANNEL_SIZE);
179 connection_state->handshake_mgr->RemoveFromPendingMgrList(
180 &connection_state->svr_state->pending_handshake_mgrs);
181 gpr_mu_unlock(&connection_state->svr_state->mu);
182 connection_state->handshake_mgr.reset();
183 gpr_free(connection_state->acceptor);
184 grpc_tcp_server_unref(connection_state->svr_state->tcp_server);
185 server_connection_state_unref(connection_state);
188 static void on_accept(void* arg, grpc_endpoint* tcp,
189 grpc_pollset* accepting_pollset,
190 grpc_tcp_server_acceptor* acceptor) {
191 server_state* state = static_cast<server_state*>(arg);
192 gpr_mu_lock(&state->mu);
193 if (state->shutdown) {
194 gpr_mu_unlock(&state->mu);
195 grpc_endpoint_shutdown(tcp, GRPC_ERROR_NONE);
196 grpc_endpoint_destroy(tcp);
200 grpc_resource_user* resource_user =
201 grpc_server_get_default_resource_user(state->server);
202 if (resource_user != nullptr &&
203 !grpc_resource_user_safe_alloc(resource_user,
204 GRPC_RESOURCE_QUOTA_CHANNEL_SIZE)) {
207 "Memory quota exhausted, rejecting the connection, no handshaking.");
208 gpr_mu_unlock(&state->mu);
209 grpc_endpoint_shutdown(tcp, GRPC_ERROR_NONE);
210 grpc_endpoint_destroy(tcp);
214 auto handshake_mgr = grpc_core::MakeRefCounted<grpc_core::HandshakeManager>();
215 handshake_mgr->AddToPendingMgrList(&state->pending_handshake_mgrs);
216 grpc_tcp_server_ref(state->tcp_server);
217 gpr_mu_unlock(&state->mu);
218 server_connection_state* connection_state =
219 static_cast<server_connection_state*>(
220 gpr_zalloc(sizeof(*connection_state)));
221 gpr_ref_init(&connection_state->refs, 1);
222 connection_state->svr_state = state;
223 connection_state->accepting_pollset = accepting_pollset;
224 connection_state->acceptor = acceptor;
225 connection_state->handshake_mgr = handshake_mgr;
226 connection_state->interested_parties = grpc_pollset_set_create();
227 grpc_pollset_set_add_pollset(connection_state->interested_parties,
228 connection_state->accepting_pollset);
229 grpc_core::HandshakerRegistry::AddHandshakers(
230 grpc_core::HANDSHAKER_SERVER, state->args,
231 connection_state->interested_parties,
232 connection_state->handshake_mgr.get());
233 const grpc_arg* timeout_arg =
234 grpc_channel_args_find(state->args, GRPC_ARG_SERVER_HANDSHAKE_TIMEOUT_MS);
235 connection_state->deadline =
236 grpc_core::ExecCtx::Get()->Now() +
237 grpc_channel_arg_get_integer(timeout_arg,
238 {120 * GPR_MS_PER_SEC, 1, INT_MAX});
239 connection_state->handshake_mgr->DoHandshake(
240 tcp, state->args, connection_state->deadline, acceptor, on_handshake_done,
244 /* Server callback: start listening on our ports */
245 static void server_start_listener(grpc_server* /*server*/, void* arg,
246 grpc_pollset** pollsets,
247 size_t pollset_count) {
248 server_state* state = static_cast<server_state*>(arg);
249 gpr_mu_lock(&state->mu);
250 state->shutdown = false;
251 gpr_mu_unlock(&state->mu);
252 grpc_tcp_server_start(state->tcp_server, pollsets, pollset_count, on_accept,
256 static void tcp_server_shutdown_complete(void* arg, grpc_error* error) {
257 server_state* state = static_cast<server_state*>(arg);
258 /* ensure all threads have unlocked */
259 gpr_mu_lock(&state->mu);
260 grpc_closure* destroy_done = state->server_destroy_listener_done;
261 GPR_ASSERT(state->shutdown);
262 if (state->pending_handshake_mgrs != nullptr) {
263 state->pending_handshake_mgrs->ShutdownAllPending(GRPC_ERROR_REF(error));
265 state->channelz_listen_socket.reset();
266 gpr_mu_unlock(&state->mu);
267 // Flush queued work before destroying handshaker factory, since that
268 // may do a synchronous unref.
269 grpc_core::ExecCtx::Get()->Flush();
270 if (destroy_done != nullptr) {
271 grpc_core::ExecCtx::Run(DEBUG_LOCATION, destroy_done,
272 GRPC_ERROR_REF(error));
273 grpc_core::ExecCtx::Get()->Flush();
275 grpc_channel_args_destroy(state->args);
276 gpr_mu_destroy(&state->mu);
280 /* Server callback: destroy the tcp listener (so we don't generate further
282 static void server_destroy_listener(grpc_server* /*server*/, void* arg,
283 grpc_closure* destroy_done) {
284 server_state* state = static_cast<server_state*>(arg);
285 gpr_mu_lock(&state->mu);
286 state->shutdown = true;
287 state->server_destroy_listener_done = destroy_done;
288 grpc_tcp_server* tcp_server = state->tcp_server;
289 gpr_mu_unlock(&state->mu);
290 grpc_tcp_server_shutdown_listeners(tcp_server);
291 grpc_tcp_server_unref(tcp_server);
294 static grpc_error* chttp2_server_add_acceptor(grpc_server* server,
296 grpc_channel_args* args) {
297 grpc_tcp_server* tcp_server = nullptr;
298 grpc_error* err = GRPC_ERROR_NONE;
299 server_state* state = nullptr;
300 const grpc_arg* arg = nullptr;
301 grpc_core::TcpServerFdHandler** arg_val = nullptr;
302 state = static_cast<server_state*>(gpr_zalloc(sizeof(*state)));
303 GRPC_CLOSURE_INIT(&state->tcp_server_shutdown_complete,
304 tcp_server_shutdown_complete, state,
305 grpc_schedule_on_exec_ctx);
306 err = grpc_tcp_server_create(&state->tcp_server_shutdown_complete, args,
308 if (err != GRPC_ERROR_NONE) {
311 state->server = server;
312 state->tcp_server = tcp_server;
314 state->shutdown = true;
315 gpr_mu_init(&state->mu);
316 // TODO(yangg) channelz
317 arg = grpc_channel_args_find(args, name);
318 GPR_ASSERT(arg->type == GRPC_ARG_POINTER);
319 arg_val = static_cast<grpc_core::TcpServerFdHandler**>(arg->value.pointer.p);
320 *arg_val = grpc_tcp_server_create_fd_handler(tcp_server);
322 grpc_server_add_listener(server, state, server_start_listener,
323 server_destroy_listener, /* node */ nullptr);
326 /* Error path: cleanup and return */
328 GPR_ASSERT(err != GRPC_ERROR_NONE);
330 grpc_tcp_server_unref(tcp_server);
332 grpc_channel_args_destroy(args);
338 grpc_error* grpc_chttp2_server_add_port(grpc_server* server, const char* addr,
339 grpc_channel_args* args,
341 grpc_resolved_addresses* resolved = nullptr;
342 grpc_tcp_server* tcp_server = nullptr;
346 grpc_error* err = GRPC_ERROR_NONE;
347 server_state* state = nullptr;
348 grpc_error** errors = nullptr;
350 const grpc_arg* arg = nullptr;
354 if (strncmp(addr, "external:", 9) == 0) {
355 return chttp2_server_add_acceptor(server, addr, args);
358 /* resolve address */
359 err = grpc_blocking_resolve_address(addr, "https", &resolved);
360 if (err != GRPC_ERROR_NONE) {
363 state = static_cast<server_state*>(gpr_zalloc(sizeof(*state)));
364 GRPC_CLOSURE_INIT(&state->tcp_server_shutdown_complete,
365 tcp_server_shutdown_complete, state,
366 grpc_schedule_on_exec_ctx);
367 err = grpc_tcp_server_create(&state->tcp_server_shutdown_complete, args,
369 if (err != GRPC_ERROR_NONE) {
373 state->server = server;
374 state->tcp_server = tcp_server;
376 state->shutdown = true;
377 gpr_mu_init(&state->mu);
379 naddrs = resolved->naddrs;
380 errors = static_cast<grpc_error**>(gpr_malloc(sizeof(*errors) * naddrs));
381 for (i = 0; i < naddrs; i++) {
383 grpc_tcp_server_add_port(tcp_server, &resolved->addrs[i], &port_temp);
384 if (errors[i] == GRPC_ERROR_NONE) {
385 if (*port_num == -1) {
386 *port_num = port_temp;
388 GPR_ASSERT(*port_num == port_temp);
395 gpr_asprintf(&msg, "No address added out of total %" PRIuPTR " resolved",
397 err = GRPC_ERROR_CREATE_REFERENCING_FROM_COPIED_STRING(msg, errors, naddrs);
400 } else if (count != naddrs) {
403 "Only %" PRIuPTR " addresses added out of total %" PRIuPTR
406 err = GRPC_ERROR_CREATE_REFERENCING_FROM_COPIED_STRING(msg, errors, naddrs);
409 const char* warning_message = grpc_error_string(err);
410 gpr_log(GPR_INFO, "WARNING: %s", warning_message);
412 /* we managed to bind some addresses: continue */
414 grpc_resolved_addresses_destroy(resolved);
416 arg = grpc_channel_args_find(args, GRPC_ARG_ENABLE_CHANNELZ);
417 if (grpc_channel_arg_get_bool(arg, GRPC_ENABLE_CHANNELZ_DEFAULT)) {
418 state->channelz_listen_socket =
419 grpc_core::MakeRefCounted<grpc_core::channelz::ListenSocketNode>(
420 addr, absl::StrFormat("chttp2 listener %s", addr));
423 /* Register with the server only upon success */
424 grpc_server_add_listener(server, state, server_start_listener,
425 server_destroy_listener,
426 state->channelz_listen_socket);
429 /* Error path: cleanup and return */
431 GPR_ASSERT(err != GRPC_ERROR_NONE);
433 grpc_resolved_addresses_destroy(resolved);
436 grpc_tcp_server_unref(tcp_server);
438 grpc_channel_args_destroy(args);
444 if (errors != nullptr) {
445 for (i = 0; i < naddrs; i++) {
446 GRPC_ERROR_UNREF(errors[i]);