|
|
@ -74,7 +74,7 @@ static gpr_once s_init_max_accept_queue_size; |
|
|
|
static int s_max_accept_queue_size; |
|
|
|
static int s_max_accept_queue_size; |
|
|
|
|
|
|
|
|
|
|
|
/* one listening port */ |
|
|
|
/* one listening port */ |
|
|
|
typedef struct { |
|
|
|
typedef struct server_port { |
|
|
|
int fd; |
|
|
|
int fd; |
|
|
|
grpc_fd *emfd; |
|
|
|
grpc_fd *emfd; |
|
|
|
grpc_tcp_server *server; |
|
|
|
grpc_tcp_server *server; |
|
|
@ -84,8 +84,13 @@ typedef struct { |
|
|
|
struct sockaddr_un un; |
|
|
|
struct sockaddr_un un; |
|
|
|
} addr; |
|
|
|
} addr; |
|
|
|
size_t addr_len; |
|
|
|
size_t addr_len; |
|
|
|
|
|
|
|
int port; |
|
|
|
grpc_closure read_closure; |
|
|
|
grpc_closure read_closure; |
|
|
|
grpc_closure destroyed_closure; |
|
|
|
grpc_closure destroyed_closure; |
|
|
|
|
|
|
|
gpr_refcount refs; |
|
|
|
|
|
|
|
struct server_port *next; |
|
|
|
|
|
|
|
struct server_port *dual_stack_second_port; |
|
|
|
|
|
|
|
int is_dual_stack_second_port; |
|
|
|
} server_port; |
|
|
|
} server_port; |
|
|
|
|
|
|
|
|
|
|
|
static void unlink_if_unix_domain_socket(const struct sockaddr_un *un) { |
|
|
|
static void unlink_if_unix_domain_socket(const struct sockaddr_un *un) { |
|
|
@ -112,10 +117,9 @@ struct grpc_tcp_server { |
|
|
|
/* is this server shutting down? (boolean) */ |
|
|
|
/* is this server shutting down? (boolean) */ |
|
|
|
int shutdown; |
|
|
|
int shutdown; |
|
|
|
|
|
|
|
|
|
|
|
/* all listening ports */ |
|
|
|
/* linked list of server ports */ |
|
|
|
server_port *ports; |
|
|
|
server_port *head; |
|
|
|
size_t nports; |
|
|
|
unsigned nports; |
|
|
|
size_t port_capacity; |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
/* shutdown callback */ |
|
|
|
/* shutdown callback */ |
|
|
|
grpc_closure *shutdown_complete; |
|
|
|
grpc_closure *shutdown_complete; |
|
|
@ -134,18 +138,22 @@ grpc_tcp_server *grpc_tcp_server_create(void) { |
|
|
|
s->shutdown = 0; |
|
|
|
s->shutdown = 0; |
|
|
|
s->on_accept_cb = NULL; |
|
|
|
s->on_accept_cb = NULL; |
|
|
|
s->on_accept_cb_arg = NULL; |
|
|
|
s->on_accept_cb_arg = NULL; |
|
|
|
s->ports = gpr_malloc(sizeof(server_port) * INIT_PORT_CAP); |
|
|
|
s->head = NULL; |
|
|
|
s->nports = 0; |
|
|
|
s->nports = 0; |
|
|
|
s->port_capacity = INIT_PORT_CAP; |
|
|
|
|
|
|
|
return s; |
|
|
|
return s; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
static void finish_shutdown(grpc_exec_ctx *exec_ctx, grpc_tcp_server *s) { |
|
|
|
static void finish_shutdown(grpc_exec_ctx *exec_ctx, grpc_tcp_server *s) { |
|
|
|
|
|
|
|
server_port *sp; |
|
|
|
|
|
|
|
|
|
|
|
grpc_exec_ctx_enqueue(exec_ctx, s->shutdown_complete, 1); |
|
|
|
grpc_exec_ctx_enqueue(exec_ctx, s->shutdown_complete, 1); |
|
|
|
|
|
|
|
|
|
|
|
gpr_mu_destroy(&s->mu); |
|
|
|
gpr_mu_destroy(&s->mu); |
|
|
|
|
|
|
|
|
|
|
|
gpr_free(s->ports); |
|
|
|
for (sp = s->head; sp; sp = sp->next) { |
|
|
|
|
|
|
|
grpc_tcp_listener_unref((grpc_tcp_listener *)sp); |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
gpr_free(s); |
|
|
|
gpr_free(s); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
@ -166,8 +174,6 @@ static void destroyed_port(grpc_exec_ctx *exec_ctx, void *server, int success) { |
|
|
|
events will be received on them - at this point it's safe to destroy |
|
|
|
events will be received on them - at this point it's safe to destroy |
|
|
|
things */ |
|
|
|
things */ |
|
|
|
static void deactivated_all_ports(grpc_exec_ctx *exec_ctx, grpc_tcp_server *s) { |
|
|
|
static void deactivated_all_ports(grpc_exec_ctx *exec_ctx, grpc_tcp_server *s) { |
|
|
|
size_t i; |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
/* delete ALL the things */ |
|
|
|
/* delete ALL the things */ |
|
|
|
gpr_mu_lock(&s->mu); |
|
|
|
gpr_mu_lock(&s->mu); |
|
|
|
|
|
|
|
|
|
|
@ -176,9 +182,9 @@ static void deactivated_all_ports(grpc_exec_ctx *exec_ctx, grpc_tcp_server *s) { |
|
|
|
return; |
|
|
|
return; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
if (s->nports) { |
|
|
|
if (s->head) { |
|
|
|
for (i = 0; i < s->nports; i++) { |
|
|
|
server_port *sp; |
|
|
|
server_port *sp = &s->ports[i]; |
|
|
|
for (sp = s->head; sp; sp = sp->next) { |
|
|
|
if (sp->addr.sockaddr.sa_family == AF_UNIX) { |
|
|
|
if (sp->addr.sockaddr.sa_family == AF_UNIX) { |
|
|
|
unlink_if_unix_domain_socket(&sp->addr.un); |
|
|
|
unlink_if_unix_domain_socket(&sp->addr.un); |
|
|
|
} |
|
|
|
} |
|
|
@ -196,7 +202,6 @@ static void deactivated_all_ports(grpc_exec_ctx *exec_ctx, grpc_tcp_server *s) { |
|
|
|
|
|
|
|
|
|
|
|
void grpc_tcp_server_destroy(grpc_exec_ctx *exec_ctx, grpc_tcp_server *s, |
|
|
|
void grpc_tcp_server_destroy(grpc_exec_ctx *exec_ctx, grpc_tcp_server *s, |
|
|
|
grpc_closure *closure) { |
|
|
|
grpc_closure *closure) { |
|
|
|
size_t i; |
|
|
|
|
|
|
|
gpr_mu_lock(&s->mu); |
|
|
|
gpr_mu_lock(&s->mu); |
|
|
|
|
|
|
|
|
|
|
|
GPR_ASSERT(!s->shutdown); |
|
|
|
GPR_ASSERT(!s->shutdown); |
|
|
@ -206,8 +211,9 @@ void grpc_tcp_server_destroy(grpc_exec_ctx *exec_ctx, grpc_tcp_server *s, |
|
|
|
|
|
|
|
|
|
|
|
/* shutdown all fd's */ |
|
|
|
/* shutdown all fd's */ |
|
|
|
if (s->active_ports) { |
|
|
|
if (s->active_ports) { |
|
|
|
for (i = 0; i < s->nports; i++) { |
|
|
|
server_port *sp; |
|
|
|
grpc_fd_shutdown(exec_ctx, s->ports[i].emfd); |
|
|
|
for (sp = s->head; sp; sp = sp->next) { |
|
|
|
|
|
|
|
grpc_fd_shutdown(exec_ctx, sp->emfd); |
|
|
|
} |
|
|
|
} |
|
|
|
gpr_mu_unlock(&s->mu); |
|
|
|
gpr_mu_unlock(&s->mu); |
|
|
|
} else { |
|
|
|
} else { |
|
|
@ -364,9 +370,10 @@ error: |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
static int add_socket_to_server(grpc_tcp_server *s, int fd, |
|
|
|
static server_port *add_socket_to_server(grpc_tcp_server *s, int fd, |
|
|
|
const struct sockaddr *addr, size_t addr_len) { |
|
|
|
const struct sockaddr *addr, |
|
|
|
server_port *sp; |
|
|
|
size_t addr_len) { |
|
|
|
|
|
|
|
server_port *sp = NULL; |
|
|
|
int port; |
|
|
|
int port; |
|
|
|
char *addr_str; |
|
|
|
char *addr_str; |
|
|
|
char *name; |
|
|
|
char *name; |
|
|
@ -376,32 +383,35 @@ static int add_socket_to_server(grpc_tcp_server *s, int fd, |
|
|
|
grpc_sockaddr_to_string(&addr_str, (struct sockaddr *)&addr, 1); |
|
|
|
grpc_sockaddr_to_string(&addr_str, (struct sockaddr *)&addr, 1); |
|
|
|
gpr_asprintf(&name, "tcp-server-listener:%s", addr_str); |
|
|
|
gpr_asprintf(&name, "tcp-server-listener:%s", addr_str); |
|
|
|
gpr_mu_lock(&s->mu); |
|
|
|
gpr_mu_lock(&s->mu); |
|
|
|
|
|
|
|
s->nports++; |
|
|
|
GPR_ASSERT(!s->on_accept_cb && "must add ports before starting server"); |
|
|
|
GPR_ASSERT(!s->on_accept_cb && "must add ports before starting server"); |
|
|
|
/* append it to the list under a lock */ |
|
|
|
sp = gpr_malloc(sizeof(server_port)); |
|
|
|
if (s->nports == s->port_capacity) { |
|
|
|
sp->next = s->head; |
|
|
|
s->port_capacity *= 2; |
|
|
|
s->head = sp; |
|
|
|
s->ports = gpr_realloc(s->ports, sizeof(server_port) * s->port_capacity); |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
sp = &s->ports[s->nports++]; |
|
|
|
|
|
|
|
sp->server = s; |
|
|
|
sp->server = s; |
|
|
|
sp->fd = fd; |
|
|
|
sp->fd = fd; |
|
|
|
sp->emfd = grpc_fd_create(fd, name); |
|
|
|
sp->emfd = grpc_fd_create(fd, name); |
|
|
|
memcpy(sp->addr.untyped, addr, addr_len); |
|
|
|
memcpy(sp->addr.untyped, addr, addr_len); |
|
|
|
sp->addr_len = addr_len; |
|
|
|
sp->addr_len = addr_len; |
|
|
|
|
|
|
|
sp->port = port; |
|
|
|
|
|
|
|
sp->is_dual_stack_second_port = 0; |
|
|
|
|
|
|
|
sp->dual_stack_second_port = NULL; |
|
|
|
|
|
|
|
gpr_ref_init(&sp->refs, 1); |
|
|
|
GPR_ASSERT(sp->emfd); |
|
|
|
GPR_ASSERT(sp->emfd); |
|
|
|
gpr_mu_unlock(&s->mu); |
|
|
|
gpr_mu_unlock(&s->mu); |
|
|
|
gpr_free(addr_str); |
|
|
|
gpr_free(addr_str); |
|
|
|
gpr_free(name); |
|
|
|
gpr_free(name); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
return port; |
|
|
|
return sp; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
int grpc_tcp_server_add_port(grpc_tcp_server *s, const void *addr, |
|
|
|
grpc_tcp_listener *grpc_tcp_server_add_port(grpc_tcp_server *s, |
|
|
|
size_t addr_len) { |
|
|
|
const void *addr, |
|
|
|
int allocated_port1 = -1; |
|
|
|
size_t addr_len) { |
|
|
|
int allocated_port2 = -1; |
|
|
|
int allocated_port = -1; |
|
|
|
unsigned i; |
|
|
|
server_port *sp; |
|
|
|
|
|
|
|
server_port *sp2 = NULL; |
|
|
|
int fd; |
|
|
|
int fd; |
|
|
|
grpc_dualstack_mode dsmode; |
|
|
|
grpc_dualstack_mode dsmode; |
|
|
|
struct sockaddr_in6 addr6_v4mapped; |
|
|
|
struct sockaddr_in6 addr6_v4mapped; |
|
|
@ -420,9 +430,9 @@ int grpc_tcp_server_add_port(grpc_tcp_server *s, const void *addr, |
|
|
|
/* Check if this is a wildcard port, and if so, try to keep the port the same
|
|
|
|
/* Check if this is a wildcard port, and if so, try to keep the port the same
|
|
|
|
as some previously created listener. */ |
|
|
|
as some previously created listener. */ |
|
|
|
if (grpc_sockaddr_get_port(addr) == 0) { |
|
|
|
if (grpc_sockaddr_get_port(addr) == 0) { |
|
|
|
for (i = 0; i < s->nports; i++) { |
|
|
|
for (sp = s->head; sp; sp = sp->next) { |
|
|
|
sockname_len = sizeof(sockname_temp); |
|
|
|
sockname_len = sizeof(sockname_temp); |
|
|
|
if (0 == getsockname(s->ports[i].fd, (struct sockaddr *)&sockname_temp, |
|
|
|
if (0 == getsockname(sp->fd, (struct sockaddr *)&sockname_temp, |
|
|
|
&sockname_len)) { |
|
|
|
&sockname_len)) { |
|
|
|
port = grpc_sockaddr_get_port((struct sockaddr *)&sockname_temp); |
|
|
|
port = grpc_sockaddr_get_port((struct sockaddr *)&sockname_temp); |
|
|
|
if (port > 0) { |
|
|
|
if (port > 0) { |
|
|
@ -436,6 +446,8 @@ int grpc_tcp_server_add_port(grpc_tcp_server *s, const void *addr, |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
sp = NULL; |
|
|
|
|
|
|
|
|
|
|
|
if (grpc_sockaddr_to_v4mapped(addr, &addr6_v4mapped)) { |
|
|
|
if (grpc_sockaddr_to_v4mapped(addr, &addr6_v4mapped)) { |
|
|
|
addr = (const struct sockaddr *)&addr6_v4mapped; |
|
|
|
addr = (const struct sockaddr *)&addr6_v4mapped; |
|
|
|
addr_len = sizeof(addr6_v4mapped); |
|
|
|
addr_len = sizeof(addr6_v4mapped); |
|
|
@ -449,14 +461,16 @@ int grpc_tcp_server_add_port(grpc_tcp_server *s, const void *addr, |
|
|
|
addr = (struct sockaddr *)&wild6; |
|
|
|
addr = (struct sockaddr *)&wild6; |
|
|
|
addr_len = sizeof(wild6); |
|
|
|
addr_len = sizeof(wild6); |
|
|
|
fd = grpc_create_dualstack_socket(addr, SOCK_STREAM, 0, &dsmode); |
|
|
|
fd = grpc_create_dualstack_socket(addr, SOCK_STREAM, 0, &dsmode); |
|
|
|
allocated_port1 = add_socket_to_server(s, fd, addr, addr_len); |
|
|
|
sp = add_socket_to_server(s, fd, addr, addr_len); |
|
|
|
|
|
|
|
allocated_port = sp->port; |
|
|
|
if (fd >= 0 && dsmode == GRPC_DSMODE_DUALSTACK) { |
|
|
|
if (fd >= 0 && dsmode == GRPC_DSMODE_DUALSTACK) { |
|
|
|
goto done; |
|
|
|
goto done; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
/* If we didn't get a dualstack socket, also listen on 0.0.0.0. */ |
|
|
|
/* If we didn't get a dualstack socket, also listen on 0.0.0.0. */ |
|
|
|
if (port == 0 && allocated_port1 > 0) { |
|
|
|
if (port == 0 && allocated_port > 0) { |
|
|
|
grpc_sockaddr_set_port((struct sockaddr *)&wild4, allocated_port1); |
|
|
|
grpc_sockaddr_set_port((struct sockaddr *)&wild4, allocated_port); |
|
|
|
|
|
|
|
sp2 = sp; |
|
|
|
} |
|
|
|
} |
|
|
|
addr = (struct sockaddr *)&wild4; |
|
|
|
addr = (struct sockaddr *)&wild4; |
|
|
|
addr_len = sizeof(wild4); |
|
|
|
addr_len = sizeof(wild4); |
|
|
@ -471,22 +485,31 @@ int grpc_tcp_server_add_port(grpc_tcp_server *s, const void *addr, |
|
|
|
addr = (struct sockaddr *)&addr4_copy; |
|
|
|
addr = (struct sockaddr *)&addr4_copy; |
|
|
|
addr_len = sizeof(addr4_copy); |
|
|
|
addr_len = sizeof(addr4_copy); |
|
|
|
} |
|
|
|
} |
|
|
|
allocated_port2 = add_socket_to_server(s, fd, addr, addr_len); |
|
|
|
sp = add_socket_to_server(s, fd, addr, addr_len); |
|
|
|
|
|
|
|
sp->dual_stack_second_port = sp2; |
|
|
|
|
|
|
|
if (sp2) sp2->is_dual_stack_second_port = 1; |
|
|
|
|
|
|
|
|
|
|
|
done: |
|
|
|
done: |
|
|
|
gpr_free(allocated_addr); |
|
|
|
gpr_free(allocated_addr); |
|
|
|
return allocated_port1 >= 0 ? allocated_port1 : allocated_port2; |
|
|
|
return (grpc_tcp_listener *)sp; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
int grpc_tcp_server_get_fd(grpc_tcp_server *s, unsigned port_index) { |
|
|
|
int grpc_tcp_server_get_fd(grpc_tcp_server *s, unsigned port_index) { |
|
|
|
return (port_index < s->nports) ? s->ports[port_index].fd : -1; |
|
|
|
server_port *sp; |
|
|
|
|
|
|
|
for (sp = s->head; sp && port_index != 0; sp = sp->next, port_index--); |
|
|
|
|
|
|
|
if (port_index == 0 && sp) { |
|
|
|
|
|
|
|
return sp->fd; |
|
|
|
|
|
|
|
} else { |
|
|
|
|
|
|
|
return -1; |
|
|
|
|
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
void grpc_tcp_server_start(grpc_exec_ctx *exec_ctx, grpc_tcp_server *s, |
|
|
|
void grpc_tcp_server_start(grpc_exec_ctx *exec_ctx, grpc_tcp_server *s, |
|
|
|
grpc_pollset **pollsets, size_t pollset_count, |
|
|
|
grpc_pollset **pollsets, size_t pollset_count, |
|
|
|
grpc_tcp_server_cb on_accept_cb, |
|
|
|
grpc_tcp_server_cb on_accept_cb, |
|
|
|
void *on_accept_cb_arg) { |
|
|
|
void *on_accept_cb_arg) { |
|
|
|
size_t i, j; |
|
|
|
size_t i; |
|
|
|
|
|
|
|
server_port *sp; |
|
|
|
GPR_ASSERT(on_accept_cb); |
|
|
|
GPR_ASSERT(on_accept_cb); |
|
|
|
gpr_mu_lock(&s->mu); |
|
|
|
gpr_mu_lock(&s->mu); |
|
|
|
GPR_ASSERT(!s->on_accept_cb); |
|
|
|
GPR_ASSERT(!s->on_accept_cb); |
|
|
@ -495,17 +518,36 @@ void grpc_tcp_server_start(grpc_exec_ctx *exec_ctx, grpc_tcp_server *s, |
|
|
|
s->on_accept_cb_arg = on_accept_cb_arg; |
|
|
|
s->on_accept_cb_arg = on_accept_cb_arg; |
|
|
|
s->pollsets = pollsets; |
|
|
|
s->pollsets = pollsets; |
|
|
|
s->pollset_count = pollset_count; |
|
|
|
s->pollset_count = pollset_count; |
|
|
|
for (i = 0; i < s->nports; i++) { |
|
|
|
for (sp = s->head; sp; sp = sp->next) { |
|
|
|
for (j = 0; j < pollset_count; j++) { |
|
|
|
for (i = 0; i < pollset_count; i++) { |
|
|
|
grpc_pollset_add_fd(exec_ctx, pollsets[j], s->ports[i].emfd); |
|
|
|
grpc_pollset_add_fd(exec_ctx, pollsets[i], sp->emfd); |
|
|
|
} |
|
|
|
} |
|
|
|
s->ports[i].read_closure.cb = on_read; |
|
|
|
sp->read_closure.cb = on_read; |
|
|
|
s->ports[i].read_closure.cb_arg = &s->ports[i]; |
|
|
|
sp->read_closure.cb_arg = sp; |
|
|
|
grpc_fd_notify_on_read(exec_ctx, s->ports[i].emfd, |
|
|
|
grpc_fd_notify_on_read(exec_ctx, sp->emfd, |
|
|
|
&s->ports[i].read_closure); |
|
|
|
&sp->read_closure); |
|
|
|
s->active_ports++; |
|
|
|
s->active_ports++; |
|
|
|
} |
|
|
|
} |
|
|
|
gpr_mu_unlock(&s->mu); |
|
|
|
gpr_mu_unlock(&s->mu); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
int grpc_tcp_listener_get_port(grpc_tcp_listener *listener) { |
|
|
|
|
|
|
|
server_port *sp = (server_port *)listener; |
|
|
|
|
|
|
|
return sp->port; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
void grpc_tcp_listener_ref(grpc_tcp_listener *listener) { |
|
|
|
|
|
|
|
server_port *sp = (server_port *)listener; |
|
|
|
|
|
|
|
gpr_ref(&sp->refs); |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
void grpc_tcp_listener_unref(grpc_tcp_listener *listener) { |
|
|
|
|
|
|
|
server_port *sp = (server_port *)listener; |
|
|
|
|
|
|
|
if (sp->is_dual_stack_second_port) return; |
|
|
|
|
|
|
|
if (gpr_unref(&sp->refs)) { |
|
|
|
|
|
|
|
if (sp->dual_stack_second_port) gpr_free(sp->dual_stack_second_port); |
|
|
|
|
|
|
|
gpr_free(listener); |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
#endif |
|
|
|
#endif |
|
|
|