added support for sending data using pollfd
Costa Tsaousis (ktsaou) committed
Apr 28, 2017 at 01:59 UTC
924210942da2382b36917ada2fffabe4a5041902
3 files changed
+57
-7
src/socket.c
+46
-7
@@ -746,6 +746,7 @@ struct poll {
746
void *(*add_callback)(int fd, short int *events);
747
void (*del_callback)(int fd, void *data);
748
int (*rcv_callback)(int fd, int socktype, void *data, short int *events);
749
+ int (*snd_callback)(int fd, int socktype, void *data, short int *events);
750
};
751
752
static inline struct pollinfo *poll_add_fd(struct poll *p, int fd, int socktype, short int events, uint32_t flags) {
@@ -762,7 +763,7 @@ static inline struct pollinfo *poll_add_fd(struct poll *p, int fd, int socktype,
763
764
ssize_t i;
765
for(i = new_slots - 1; i >= (ssize_t)p->slots ; i--) {
765
- debug(D_POLLFD, "POLLFD: ADD: reseting new slot %zd", i);
766
+ debug(D_POLLFD, "POLLFD: ADD: resetting new slot %zd", i);
767
p->fds[i].fd = -1;
768
p->fds[i].events = 0;
769
p->fds[i].revents = 0;
@@ -854,6 +855,7 @@ void poll_events(LISTEN_SOCKETS *sockets
855
, void *(*add_callback)(int fd, short int *events)
856
, void (*del_callback)(int fd, void *data)
857
, int (*rcv_callback)(int fd, int socktype, void *data, short int *events)
858
+ , int (*snd_callback)(int fd, int socktype, void *data, short int *events)
859
) {
860
int retval;
861
@@ -867,7 +869,8 @@ void poll_events(LISTEN_SOCKETS *sockets
869
870
.add_callback = add_callback,
871
.del_callback = del_callback,
870
- .rcv_callback = rcv_callback
872
+ .rcv_callback = rcv_callback,
873
+ .snd_callback = snd_callback
874
};
875
876
size_t i;
@@ -905,9 +908,10 @@ void poll_events(LISTEN_SOCKETS *sockets
908
continue;
909
}
910
908
- // check for new incoming connections
911
if(pf->revents & POLLIN || pf->revents & POLLPRI) {
910
- debug(D_POLLFD, "POLLFD: LISTENER: processing events for slot %zu (events = %d, revents = %d)", i, pf->events, pf->revents);
912
+ // receiving data
913
+
914
+ debug(D_POLLFD, "POLLFD: LISTENER: processing POLLIN events for slot %zu (events = %d, revents = %d)", i, pf->events, pf->revents);
915
916
pf->revents = 0;
917
@@ -916,8 +920,10 @@ void poll_events(LISTEN_SOCKETS *sockets
920
921
debug(D_POLLFD, "POLLFD: LISTENER: reading data from TCP client slot %zu (fd %d)", i, fd);
922
919
- if (p.rcv_callback(fd, pi->socktype, pi->data, &pf->events) == -1)
923
+ if (p.rcv_callback(fd, pi->socktype, pi->data, &pf->events) == -1) {
924
poll_close_fd(&p, pi);
925
+ continue;
926
+ }
927
}
928
929
if(likely(pi->flags & POLLINFO_FLAG_SERVER_SOCKET)) {
@@ -973,8 +979,41 @@ void poll_events(LISTEN_SOCKETS *sockets
979
}
980
}
981
}
976
- else {
977
- debug(D_POLLFD, "POLLFD: LISTENER: no events for slot %zu (events = %d, revents = %d)", i, pf->events, pf->revents);
982
+
983
+ if(unlikely(pf->revents & POLLOUT)) {
984
+ // sending data
985
+
986
+ debug(D_POLLFD, "POLLFD: LISTENER: processing POLLOUT events for slot %zu (events = %d, revents = %d)", i, pf->events, pf->revents);
987
+
988
+ pf->revents = 0;
989
+
990
+ debug(D_POLLFD, "POLLFD: LISTENER: sending data to socket on slot %zu (fd %d)", i, fd);
991
+
992
+ if (p.snd_callback(fd, pi->socktype, pi->data, &pf->events) == -1) {
993
+ poll_close_fd(&p, pi);
994
+ continue;
995
+ }
996
+ }
997
+
998
+ if(unlikely(pf->revents & POLLERR)) {
999
+ debug(D_POLLFD, "POLLFD: LISTENER: processing POLLERR events for slot %zu (events = %d, revents = %d)", i, pf->events, pf->revents);
1000
+ error("POLLFD: LISTENER: processing POLLERR events for slot %zu (events = %d, revents = %d)", i, pf->events, pf->revents);
1001
+ poll_close_fd(&p, pi);
1002
+ continue;
1003
+ }
1004
+
1005
+ if(unlikely(pf->revents & POLLHUP)) {
1006
+ debug(D_POLLFD, "POLLFD: LISTENER: processing POLLHUP events for slot %zu (events = %d, revents = %d)", i, pf->events, pf->revents);
1007
+ error("POLLFD: LISTENER: processing POLLHUP events for slot %zu (events = %d, revents = %d)", i, pf->events, pf->revents);
1008
+ poll_close_fd(&p, pi);
1009
+ continue;
1010
+ }
1011
+
1012
+ if(unlikely(pf->revents & POLLNVAL)) {
1013
+ debug(D_POLLFD, "POLLFD: LISTENER: processing POLLNVAP events for slot %zu (events = %d, revents = %d)", i, pf->events, pf->revents);
1014
+ error("POLLFD: LISTENER: processing POLLNVAP events for slot %zu (events = %d, revents = %d)", i, pf->events, pf->revents);
1015
+ poll_close_fd(&p, pi);
1016
+ continue;
1017
}
1018
}
1019
}
src/socket.h
+1
@@ -47,6 +47,7 @@ extern void poll_events(LISTEN_SOCKETS *sockets
47
, void *(*add_callback)(int fd, short int *events)
48
, void (*del_callback)(int fd, void *data)
49
, int (*rcv_callback)(int fd, int socktype, void *data, short int *events)
50
+ , int (*snd_callback)(int fd, int socktype, void *data, short int *events)
51
);
52
53
#endif //NETDATA_SOCKET_H
src/statsd.c
+10
@@ -682,6 +682,15 @@ static int statsd_rcv_callback(int fd, int socktype, void *data, short int *even
682
return 0;
683
}
684
685
+static int statsd_snd_callback(int fd, int socktype, void *data, short int *events) {
686
+ (void)fd;
687
+ (void)socktype;
688
+ (void)data;
689
+ (void)events;
690
+
691
+ error("STATSD: snd_callback() called, but we never requested to send data to statsd clients.");
692
+ return -1;
693
+}
694
695
// --------------------------------------------------------------------------------------------------------------------
696
// statsd child thread to collect metrics from network
@@ -701,6 +710,7 @@ void *statsd_collector_thread(void *ptr) {
710
, statsd_add_callback
711
, statsd_del_callback
712
, statsd_rcv_callback
713
+ , statsd_snd_callback
714
);
715
716
debug(D_WEB_CLIENT, "STATSD: exit!");