aboutsummaryrefslogtreecommitdiffhomepage
path: root/src/core/iomgr/udp_server.c
diff options
context:
space:
mode:
authorGravatar Craig Tiller <ctiller@google.com>2015-09-22 10:42:19 -0700
committerGravatar Craig Tiller <ctiller@google.com>2015-09-22 10:42:19 -0700
commit45724b35e411fef7c5da66a74c78428c11d56843 (patch)
tree9264034aca675c89444e02f72ef58e67d7043604 /src/core/iomgr/udp_server.c
parent298751c1195523ef6228595043b583c3a6270e08 (diff)
indent pass to get logical source lines on one physical line
Diffstat (limited to 'src/core/iomgr/udp_server.c')
-rw-r--r--src/core/iomgr/udp_server.c459
1 files changed, 255 insertions, 204 deletions
diff --git a/src/core/iomgr/udp_server.c b/src/core/iomgr/udp_server.c
index 0f0bea82ab..490b3de205 100644
--- a/src/core/iomgr/udp_server.c
+++ b/src/core/iomgr/udp_server.c
@@ -69,11 +69,13 @@
#define INIT_PORT_CAP 2
/* one listening port */
-typedef struct {
+typedef struct
+{
int fd;
grpc_fd *emfd;
grpc_udp_server *server;
- union {
+ union
+ {
gpr_uint8 untyped[GRPC_MAX_SOCKADDR_SIZE];
struct sockaddr sockaddr;
struct sockaddr_un un;
@@ -84,16 +86,20 @@ typedef struct {
grpc_udp_server_read_cb read_cb;
} 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)
+{
struct stat st;
- if (stat(un->sun_path, &st) == 0 && (st.st_mode & S_IFMT) == S_IFSOCK) {
- unlink(un->sun_path);
- }
+ if (stat (un->sun_path, &st) == 0 && (st.st_mode & S_IFMT) == S_IFSOCK)
+ {
+ unlink (un->sun_path);
+ }
}
/* the overall server */
-struct grpc_udp_server {
+struct grpc_udp_server
+{
gpr_mu mu;
gpr_cv cv;
@@ -119,206 +125,237 @@ struct grpc_udp_server {
size_t pollset_count;
};
-grpc_udp_server *grpc_udp_server_create(void) {
- grpc_udp_server *s = gpr_malloc(sizeof(grpc_udp_server));
- gpr_mu_init(&s->mu);
- gpr_cv_init(&s->cv);
+grpc_udp_server *
+grpc_udp_server_create (void)
+{
+ grpc_udp_server *s = gpr_malloc (sizeof (grpc_udp_server));
+ gpr_mu_init (&s->mu);
+ gpr_cv_init (&s->cv);
s->active_ports = 0;
s->destroyed_ports = 0;
s->shutdown = 0;
- s->ports = gpr_malloc(sizeof(server_port) * INIT_PORT_CAP);
+ s->ports = gpr_malloc (sizeof (server_port) * INIT_PORT_CAP);
s->nports = 0;
s->port_capacity = INIT_PORT_CAP;
return s;
}
-static void finish_shutdown(grpc_udp_server *s,
- grpc_closure_list *closure_list) {
- grpc_closure_list_add(closure_list, s->shutdown_complete, 1);
+static void
+finish_shutdown (grpc_udp_server * s, grpc_closure_list * closure_list)
+{
+ grpc_closure_list_add (closure_list, s->shutdown_complete, 1);
- gpr_mu_destroy(&s->mu);
- gpr_cv_destroy(&s->cv);
+ gpr_mu_destroy (&s->mu);
+ gpr_cv_destroy (&s->cv);
- gpr_free(s->ports);
- gpr_free(s);
+ gpr_free (s->ports);
+ gpr_free (s);
}
-static void destroyed_port(void *server, int success,
- grpc_closure_list *closure_list) {
+static void
+destroyed_port (void *server, int success, grpc_closure_list * closure_list)
+{
grpc_udp_server *s = server;
- gpr_mu_lock(&s->mu);
+ gpr_mu_lock (&s->mu);
s->destroyed_ports++;
- if (s->destroyed_ports == s->nports) {
- gpr_mu_unlock(&s->mu);
- finish_shutdown(s, closure_list);
- } else {
- gpr_mu_unlock(&s->mu);
- }
+ if (s->destroyed_ports == s->nports)
+ {
+ gpr_mu_unlock (&s->mu);
+ finish_shutdown (s, closure_list);
+ }
+ else
+ {
+ gpr_mu_unlock (&s->mu);
+ }
}
/* called when all listening endpoints have been shutdown, so no further
events will be received on them - at this point it's safe to destroy
things */
-static void deactivated_all_ports(grpc_udp_server *s,
- grpc_closure_list *closure_list) {
+static void
+deactivated_all_ports (grpc_udp_server * s, grpc_closure_list * closure_list)
+{
size_t i;
/* delete ALL the things */
- gpr_mu_lock(&s->mu);
-
- if (!s->shutdown) {
- gpr_mu_unlock(&s->mu);
- return;
- }
-
- if (s->nports) {
- for (i = 0; i < s->nports; i++) {
- server_port *sp = &s->ports[i];
- if (sp->addr.sockaddr.sa_family == AF_UNIX) {
- unlink_if_unix_domain_socket(&sp->addr.un);
- }
- sp->destroyed_closure.cb = destroyed_port;
- sp->destroyed_closure.cb_arg = s;
- grpc_fd_orphan(sp->emfd, &sp->destroyed_closure, "udp_listener_shutdown",
- closure_list);
+ gpr_mu_lock (&s->mu);
+
+ if (!s->shutdown)
+ {
+ gpr_mu_unlock (&s->mu);
+ return;
+ }
+
+ if (s->nports)
+ {
+ for (i = 0; i < s->nports; i++)
+ {
+ server_port *sp = &s->ports[i];
+ if (sp->addr.sockaddr.sa_family == AF_UNIX)
+ {
+ unlink_if_unix_domain_socket (&sp->addr.un);
+ }
+ sp->destroyed_closure.cb = destroyed_port;
+ sp->destroyed_closure.cb_arg = s;
+ grpc_fd_orphan (sp->emfd, &sp->destroyed_closure, "udp_listener_shutdown", closure_list);
+ }
+ gpr_mu_unlock (&s->mu);
+ }
+ else
+ {
+ gpr_mu_unlock (&s->mu);
+ finish_shutdown (s, closure_list);
}
- gpr_mu_unlock(&s->mu);
- } else {
- gpr_mu_unlock(&s->mu);
- finish_shutdown(s, closure_list);
- }
}
-void grpc_udp_server_destroy(grpc_udp_server *s, grpc_closure *on_done,
- grpc_closure_list *closure_list) {
+void
+grpc_udp_server_destroy (grpc_udp_server * s, grpc_closure * on_done, grpc_closure_list * closure_list)
+{
size_t i;
- gpr_mu_lock(&s->mu);
+ gpr_mu_lock (&s->mu);
- GPR_ASSERT(!s->shutdown);
+ GPR_ASSERT (!s->shutdown);
s->shutdown = 1;
s->shutdown_complete = on_done;
/* shutdown all fd's */
- if (s->active_ports) {
- for (i = 0; i < s->nports; i++) {
- grpc_fd_shutdown(s->ports[i].emfd, closure_list);
+ if (s->active_ports)
+ {
+ for (i = 0; i < s->nports; i++)
+ {
+ grpc_fd_shutdown (s->ports[i].emfd, closure_list);
+ }
+ gpr_mu_unlock (&s->mu);
+ }
+ else
+ {
+ gpr_mu_unlock (&s->mu);
+ deactivated_all_ports (s, closure_list);
}
- gpr_mu_unlock(&s->mu);
- } else {
- gpr_mu_unlock(&s->mu);
- deactivated_all_ports(s, closure_list);
- }
}
/* Prepare a recently-created socket for listening. */
-static int prepare_socket(int fd, const struct sockaddr *addr,
- size_t addr_len) {
+static int
+prepare_socket (int fd, const struct sockaddr *addr, size_t addr_len)
+{
struct sockaddr_storage sockname_temp;
socklen_t sockname_len;
int get_local_ip;
int rc;
- if (fd < 0) {
- goto error;
- }
+ if (fd < 0)
+ {
+ goto error;
+ }
- if (!grpc_set_socket_nonblocking(fd, 1) || !grpc_set_socket_cloexec(fd, 1)) {
- gpr_log(GPR_ERROR, "Unable to configure socket %d: %s", fd,
- strerror(errno));
- }
+ if (!grpc_set_socket_nonblocking (fd, 1) || !grpc_set_socket_cloexec (fd, 1))
+ {
+ gpr_log (GPR_ERROR, "Unable to configure socket %d: %s", fd, strerror (errno));
+ }
get_local_ip = 1;
- rc = setsockopt(fd, IPPROTO_IP, IP_PKTINFO, &get_local_ip,
- sizeof(get_local_ip));
- if (rc == 0 && addr->sa_family == AF_INET6) {
+ rc = setsockopt (fd, IPPROTO_IP, IP_PKTINFO, &get_local_ip, sizeof (get_local_ip));
+ if (rc == 0 && addr->sa_family == AF_INET6)
+ {
#if !TARGET_OS_IPHONE
- rc = setsockopt(fd, IPPROTO_IPV6, IPV6_RECVPKTINFO, &get_local_ip,
- sizeof(get_local_ip));
+ rc = setsockopt (fd, IPPROTO_IPV6, IPV6_RECVPKTINFO, &get_local_ip, sizeof (get_local_ip));
#endif
- }
+ }
- GPR_ASSERT(addr_len < ~(socklen_t)0);
- if (bind(fd, addr, (socklen_t)addr_len) < 0) {
- char *addr_str;
- grpc_sockaddr_to_string(&addr_str, addr, 0);
- gpr_log(GPR_ERROR, "bind addr=%s: %s", addr_str, strerror(errno));
- gpr_free(addr_str);
- goto error;
- }
+ GPR_ASSERT (addr_len < ~(socklen_t) 0);
+ if (bind (fd, addr, (socklen_t) addr_len) < 0)
+ {
+ char *addr_str;
+ grpc_sockaddr_to_string (&addr_str, addr, 0);
+ gpr_log (GPR_ERROR, "bind addr=%s: %s", addr_str, strerror (errno));
+ gpr_free (addr_str);
+ goto error;
+ }
- sockname_len = sizeof(sockname_temp);
- if (getsockname(fd, (struct sockaddr *)&sockname_temp, &sockname_len) < 0) {
- goto error;
- }
+ sockname_len = sizeof (sockname_temp);
+ if (getsockname (fd, (struct sockaddr *) &sockname_temp, &sockname_len) < 0)
+ {
+ goto error;
+ }
- return grpc_sockaddr_get_port((struct sockaddr *)&sockname_temp);
+ return grpc_sockaddr_get_port ((struct sockaddr *) &sockname_temp);
error:
- if (fd >= 0) {
- close(fd);
- }
+ if (fd >= 0)
+ {
+ close (fd);
+ }
return -1;
}
/* event manager callback when reads are ready */
-static void on_read(void *arg, int success, grpc_closure_list *closure_list) {
+static void
+on_read (void *arg, int success, grpc_closure_list * closure_list)
+{
server_port *sp = arg;
- if (success == 0) {
- gpr_mu_lock(&sp->server->mu);
- if (0 == --sp->server->active_ports) {
- gpr_mu_unlock(&sp->server->mu);
- deactivated_all_ports(sp->server, closure_list);
- } else {
- gpr_mu_unlock(&sp->server->mu);
+ if (success == 0)
+ {
+ gpr_mu_lock (&sp->server->mu);
+ if (0 == --sp->server->active_ports)
+ {
+ gpr_mu_unlock (&sp->server->mu);
+ deactivated_all_ports (sp->server, closure_list);
+ }
+ else
+ {
+ gpr_mu_unlock (&sp->server->mu);
+ }
+ return;
}
- return;
- }
/* Tell the registered callback that data is available to read. */
- GPR_ASSERT(sp->read_cb);
- sp->read_cb(sp->fd);
+ GPR_ASSERT (sp->read_cb);
+ sp->read_cb (sp->fd);
/* Re-arm the notification event so we get another chance to read. */
- grpc_fd_notify_on_read(sp->emfd, &sp->read_closure, closure_list);
+ grpc_fd_notify_on_read (sp->emfd, &sp->read_closure, closure_list);
}
-static int add_socket_to_server(grpc_udp_server *s, int fd,
- const struct sockaddr *addr, size_t addr_len,
- grpc_udp_server_read_cb read_cb) {
+static int
+add_socket_to_server (grpc_udp_server * s, int fd, const struct sockaddr *addr, size_t addr_len, grpc_udp_server_read_cb read_cb)
+{
server_port *sp;
int port;
char *addr_str;
char *name;
- port = prepare_socket(fd, addr, addr_len);
- if (port >= 0) {
- grpc_sockaddr_to_string(&addr_str, (struct sockaddr *)&addr, 1);
- gpr_asprintf(&name, "udp-server-listener:%s", addr_str);
- gpr_mu_lock(&s->mu);
- /* append it to the list under a lock */
- if (s->nports == s->port_capacity) {
- s->port_capacity *= 2;
- s->ports = gpr_realloc(s->ports, sizeof(server_port) * s->port_capacity);
+ port = prepare_socket (fd, addr, addr_len);
+ if (port >= 0)
+ {
+ grpc_sockaddr_to_string (&addr_str, (struct sockaddr *) &addr, 1);
+ gpr_asprintf (&name, "udp-server-listener:%s", addr_str);
+ gpr_mu_lock (&s->mu);
+ /* append it to the list under a lock */
+ if (s->nports == s->port_capacity)
+ {
+ s->port_capacity *= 2;
+ s->ports = gpr_realloc (s->ports, sizeof (server_port) * s->port_capacity);
+ }
+ sp = &s->ports[s->nports++];
+ sp->server = s;
+ sp->fd = fd;
+ sp->emfd = grpc_fd_create (fd, name);
+ memcpy (sp->addr.untyped, addr, addr_len);
+ sp->addr_len = addr_len;
+ sp->read_cb = read_cb;
+ GPR_ASSERT (sp->emfd);
+ gpr_mu_unlock (&s->mu);
}
- sp = &s->ports[s->nports++];
- sp->server = s;
- sp->fd = fd;
- sp->emfd = grpc_fd_create(fd, name);
- memcpy(sp->addr.untyped, addr, addr_len);
- sp->addr_len = addr_len;
- sp->read_cb = read_cb;
- GPR_ASSERT(sp->emfd);
- gpr_mu_unlock(&s->mu);
- }
return port;
}
-int grpc_udp_server_add_port(grpc_udp_server *s, const void *addr,
- size_t addr_len, grpc_udp_server_read_cb read_cb) {
+int
+grpc_udp_server_add_port (grpc_udp_server * s, const void *addr, size_t addr_len, grpc_udp_server_read_cb read_cb)
+{
int allocated_port1 = -1;
int allocated_port2 = -1;
unsigned i;
@@ -333,103 +370,117 @@ int grpc_udp_server_add_port(grpc_udp_server *s, const void *addr,
socklen_t sockname_len;
int port;
- if (((struct sockaddr *)addr)->sa_family == AF_UNIX) {
- unlink_if_unix_domain_socket(addr);
- }
+ if (((struct sockaddr *) addr)->sa_family == AF_UNIX)
+ {
+ unlink_if_unix_domain_socket (addr);
+ }
/* Check if this is a wildcard port, and if so, try to keep the port the same
as some previously created listener. */
- if (grpc_sockaddr_get_port(addr) == 0) {
- for (i = 0; i < s->nports; i++) {
- sockname_len = sizeof(sockname_temp);
- if (0 == getsockname(s->ports[i].fd, (struct sockaddr *)&sockname_temp,
- &sockname_len)) {
- port = grpc_sockaddr_get_port((struct sockaddr *)&sockname_temp);
- if (port > 0) {
- allocated_addr = malloc(addr_len);
- memcpy(allocated_addr, addr, addr_len);
- grpc_sockaddr_set_port(allocated_addr, port);
- addr = allocated_addr;
- break;
- }
- }
+ if (grpc_sockaddr_get_port (addr) == 0)
+ {
+ for (i = 0; i < s->nports; i++)
+ {
+ sockname_len = sizeof (sockname_temp);
+ if (0 == getsockname (s->ports[i].fd, (struct sockaddr *) &sockname_temp, &sockname_len))
+ {
+ port = grpc_sockaddr_get_port ((struct sockaddr *) &sockname_temp);
+ if (port > 0)
+ {
+ allocated_addr = malloc (addr_len);
+ memcpy (allocated_addr, addr, addr_len);
+ grpc_sockaddr_set_port (allocated_addr, port);
+ addr = allocated_addr;
+ break;
+ }
+ }
+ }
}
- }
- if (grpc_sockaddr_to_v4mapped(addr, &addr6_v4mapped)) {
- addr = (const struct sockaddr *)&addr6_v4mapped;
- addr_len = sizeof(addr6_v4mapped);
- }
+ if (grpc_sockaddr_to_v4mapped (addr, &addr6_v4mapped))
+ {
+ addr = (const struct sockaddr *) &addr6_v4mapped;
+ addr_len = sizeof (addr6_v4mapped);
+ }
/* Treat :: or 0.0.0.0 as a family-agnostic wildcard. */
- if (grpc_sockaddr_is_wildcard(addr, &port)) {
- grpc_sockaddr_make_wildcards(port, &wild4, &wild6);
-
- /* Try listening on IPv6 first. */
- addr = (struct sockaddr *)&wild6;
- addr_len = sizeof(wild6);
- fd = grpc_create_dualstack_socket(addr, SOCK_DGRAM, IPPROTO_UDP, &dsmode);
- allocated_port1 = add_socket_to_server(s, fd, addr, addr_len, read_cb);
- if (fd >= 0 && dsmode == GRPC_DSMODE_DUALSTACK) {
- goto done;
+ if (grpc_sockaddr_is_wildcard (addr, &port))
+ {
+ grpc_sockaddr_make_wildcards (port, &wild4, &wild6);
+
+ /* Try listening on IPv6 first. */
+ addr = (struct sockaddr *) &wild6;
+ addr_len = sizeof (wild6);
+ fd = grpc_create_dualstack_socket (addr, SOCK_DGRAM, IPPROTO_UDP, &dsmode);
+ allocated_port1 = add_socket_to_server (s, fd, addr, addr_len, read_cb);
+ if (fd >= 0 && dsmode == GRPC_DSMODE_DUALSTACK)
+ {
+ goto done;
+ }
+
+ /* If we didn't get a dualstack socket, also listen on 0.0.0.0. */
+ if (port == 0 && allocated_port1 > 0)
+ {
+ grpc_sockaddr_set_port ((struct sockaddr *) &wild4, allocated_port1);
+ }
+ addr = (struct sockaddr *) &wild4;
+ addr_len = sizeof (wild4);
}
- /* If we didn't get a dualstack socket, also listen on 0.0.0.0. */
- if (port == 0 && allocated_port1 > 0) {
- grpc_sockaddr_set_port((struct sockaddr *)&wild4, allocated_port1);
+ fd = grpc_create_dualstack_socket (addr, SOCK_DGRAM, IPPROTO_UDP, &dsmode);
+ if (fd < 0)
+ {
+ gpr_log (GPR_ERROR, "Unable to create socket: %s", strerror (errno));
}
- addr = (struct sockaddr *)&wild4;
- addr_len = sizeof(wild4);
- }
-
- fd = grpc_create_dualstack_socket(addr, SOCK_DGRAM, IPPROTO_UDP, &dsmode);
- if (fd < 0) {
- gpr_log(GPR_ERROR, "Unable to create socket: %s", strerror(errno));
- }
- if (dsmode == GRPC_DSMODE_IPV4 &&
- grpc_sockaddr_is_v4mapped(addr, &addr4_copy)) {
- addr = (struct sockaddr *)&addr4_copy;
- addr_len = sizeof(addr4_copy);
- }
- allocated_port2 = add_socket_to_server(s, fd, addr, addr_len, read_cb);
+ if (dsmode == GRPC_DSMODE_IPV4 && grpc_sockaddr_is_v4mapped (addr, &addr4_copy))
+ {
+ addr = (struct sockaddr *) &addr4_copy;
+ addr_len = sizeof (addr4_copy);
+ }
+ allocated_port2 = add_socket_to_server (s, fd, addr, addr_len, read_cb);
done:
- gpr_free(allocated_addr);
+ gpr_free (allocated_addr);
return allocated_port1 >= 0 ? allocated_port1 : allocated_port2;
}
-int grpc_udp_server_get_fd(grpc_udp_server *s, unsigned index) {
+int
+grpc_udp_server_get_fd (grpc_udp_server * s, unsigned index)
+{
return (index < s->nports) ? s->ports[index].fd : -1;
}
-void grpc_udp_server_start(grpc_udp_server *s, grpc_pollset **pollsets,
- size_t pollset_count,
- grpc_closure_list *closure_list) {
+void
+grpc_udp_server_start (grpc_udp_server * s, grpc_pollset ** pollsets, size_t pollset_count, grpc_closure_list * closure_list)
+{
size_t i, j;
- gpr_mu_lock(&s->mu);
- GPR_ASSERT(s->active_ports == 0);
+ gpr_mu_lock (&s->mu);
+ GPR_ASSERT (s->active_ports == 0);
s->pollsets = pollsets;
- for (i = 0; i < s->nports; i++) {
- for (j = 0; j < pollset_count; j++) {
- grpc_pollset_add_fd(pollsets[j], s->ports[i].emfd, closure_list);
+ for (i = 0; i < s->nports; i++)
+ {
+ for (j = 0; j < pollset_count; j++)
+ {
+ grpc_pollset_add_fd (pollsets[j], s->ports[i].emfd, closure_list);
+ }
+ s->ports[i].read_closure.cb = on_read;
+ s->ports[i].read_closure.cb_arg = &s->ports[i];
+ grpc_fd_notify_on_read (s->ports[i].emfd, &s->ports[i].read_closure, closure_list);
+ s->active_ports++;
}
- s->ports[i].read_closure.cb = on_read;
- s->ports[i].read_closure.cb_arg = &s->ports[i];
- grpc_fd_notify_on_read(s->ports[i].emfd, &s->ports[i].read_closure,
- closure_list);
- s->active_ports++;
- }
- gpr_mu_unlock(&s->mu);
+ gpr_mu_unlock (&s->mu);
}
/* TODO(rjshade): Add a test for this method. */
-void grpc_udp_server_write(server_port *sp, const char *buffer, size_t buf_len,
- const struct sockaddr *peer_address) {
+void
+grpc_udp_server_write (server_port * sp, const char *buffer, size_t buf_len, const struct sockaddr *peer_address)
+{
ssize_t rc;
- rc = sendto(sp->fd, buffer, buf_len, 0, peer_address, sizeof(peer_address));
- if (rc < 0) {
- gpr_log(GPR_ERROR, "Unable to send data: %s", strerror(errno));
- }
+ rc = sendto (sp->fd, buffer, buf_len, 0, peer_address, sizeof (peer_address));
+ if (rc < 0)
+ {
+ gpr_log (GPR_ERROR, "Unable to send data: %s", strerror (errno));
+ }
}
#endif