@cryptotaxi247 / netdata-1 / commits / 3efad7148

more general poll_events() api

Costa Tsaousis (ktsaou) committed Jan 8, 2018 at 02:15 UTC 3efad714866c36199bddfcd22dfdbc930baf1e5c
4 files changed +248 -244
src/socket.c
+158 -195
@@ -961,42 +961,7 @@ int accept_socket(int fd, int flags, char *client_ip, size_t ipsize, char *clien
961
962 #define POLL_FDS_INCREASE_STEP 10
963
964 -#define POLLINFO_FLAG_SERVER_SOCKET 0x00000001
965 -#define POLLINFO_FLAG_CLIENT_SOCKET 0x00000002
966 -
967 -struct pollinfo {
968 - struct poll *p; // the parent
969 -
970 - size_t slot;
971 - char *client_ip;
972 - char *client_port;
973 - struct pollinfo *next;
974 - uint32_t flags;
975 - int socktype;
976 -
977 - void (*del_callback)(int fd, int socktype, void *data);
978 - int (*rcv_callback)(int fd, int socktype, void *data, short int *events);
979 - int (*snd_callback)(int fd, int socktype, void *data, short int *events);
980 -
981 - void *data;
982 -};
983 -
984 -struct poll {
985 - size_t slots;
986 - size_t used;
987 - size_t min;
988 - size_t max;
989 - struct pollfd *fds;
990 - struct pollinfo *inf;
991 - struct pollinfo *first_free;
992 -
993 - void *(*add_callback)(int fd, int socktype, short int *events, const char *client_ip, const char *client_port);
994 - void (*del_callback)(int fd, int socktype, void *data);
995 - int (*rcv_callback)(int fd, int socktype, void *data, short int *events);
996 - int (*snd_callback)(int fd, int socktype, void *data, short int *events);
997 -};
998 -
999 -static inline struct pollinfo *poll_add_fd(struct poll *p, int fd, int socktype, short int events, uint32_t flags, const char *client_ip, const char *client_port) {
964 +static inline POLLINFO *poll_add_fd(POLLJOB *p, int fd, int socktype, short int events, uint32_t flags, const char *client_ip, const char *client_port) {
965 debug(D_POLLFD, "POLLFD: ADD: request to add fd %d, slots = %zu, used = %zu, min = %zu, max = %zu, next free = %zd", fd, p->slots, p->used, p->min, p->max, p->first_free?(ssize_t)p->first_free->slot:(ssize_t)-1);
966
967 if(unlikely(fd < 0)) return NULL;
@@ -1006,7 +971,7 @@ static inline struct pollinfo *poll_add_fd(struct poll *p, int fd, int socktype,
971 debug(D_POLLFD, "POLLFD: ADD: increasing size (current = %zu, new = %zu, used = %zu, min = %zu, max = %zu)", p->slots, new_slots, p->used, p->min, p->max);
972
973 p->fds = reallocz(p->fds, sizeof(struct pollfd) * new_slots);
1009 - p->inf = reallocz(p->inf, sizeof(struct pollinfo) * new_slots);
974 + p->inf = reallocz(p->inf, sizeof(POLLINFO) * new_slots);
975
976 // reset all the newly added slots
977 ssize_t i;
@@ -1036,7 +1001,7 @@ static inline struct pollinfo *poll_add_fd(struct poll *p, int fd, int socktype,
1001 p->slots = new_slots;
1002 }
1003
1039 - struct pollinfo *pi = p->first_free;
1004 + POLLINFO *pi = p->first_free;
1005 p->first_free = p->first_free->next;
1006
1007 debug(D_POLLFD, "POLLFD: ADD: selected slot %zu, next free is %zd", pi->slot, p->first_free?(ssize_t)p->first_free->slot:(ssize_t)-1);
@@ -1046,6 +1011,7 @@ static inline struct pollinfo *poll_add_fd(struct poll *p, int fd, int socktype,
1011 pf->events = events;
1012 pf->revents = 0;
1013
1014 + pi->fd = fd;
1015 pi->p = p;
1016 pi->socktype = socktype;
1017 pi->flags = flags;
@@ -1062,7 +1028,7 @@ static inline struct pollinfo *poll_add_fd(struct poll *p, int fd, int socktype,
1028 p->max = pi->slot;
1029
1030 if(pi->flags & POLLINFO_FLAG_CLIENT_SOCKET) {
1065 - pi->data = p->add_callback(fd, pi->socktype, &pf->events, client_ip, client_port);
1031 + pi->data = p->add_callback(pi, &pf->events);
1032 }
1033
1034 if(pi->flags & POLLINFO_FLAG_SERVER_SOCKET) {
@@ -1074,14 +1040,16 @@ static inline struct pollinfo *poll_add_fd(struct poll *p, int fd, int socktype,
1040 return pi;
1041 }
1042
1077 -static inline void poll_close_fd(struct poll *p, struct pollinfo *pi) {
1043 +static inline void poll_close_fd(POLLINFO *pi) {
1044 + POLLJOB *p = pi->p;
1045 +
1046 struct pollfd *pf = &p->fds[pi->slot];
1047 debug(D_POLLFD, "POLLFD: DEL: request to clear slot %zu (fd %d), old next free was %zd", pi->slot, pf->fd, p->first_free?(ssize_t)p->first_free->slot:(ssize_t)-1);
1048
1049 if(unlikely(pf->fd == -1)) return;
1050
1051 if(pi->flags & POLLINFO_FLAG_CLIENT_SOCKET) {
1084 - pi->del_callback(pf->fd, pi->socktype, pi->data);
1052 + pi->del_callback(pi);
1053 }
1054
1055 // info("POLLFD: closing fd %d", pf->fd);
@@ -1090,6 +1058,7 @@ static inline void poll_close_fd(struct poll *p, struct pollinfo *pi) {
1058 pf->events = 0;
1059 pf->revents = 0;
1060
1061 + pi->fd = -1;
1062 pi->socktype = -1;
1063 pi->flags = 0;
1064 pi->data = NULL;
@@ -1122,34 +1091,26 @@ static inline void poll_close_fd(struct poll *p, struct pollinfo *pi) {
1091 debug(D_POLLFD, "POLLFD: DEL: completed, slots = %zu, used = %zu, min = %zu, max = %zu, next free = %zd", p->slots, p->used, p->min, p->max, p->first_free?(ssize_t)p->first_free->slot:(ssize_t)-1);
1092 }
1093
1125 -static void *add_callback_default(int fd, int socktype, short int *events, const char *client_ip, const char *client_port) {
1126 - (void)fd;
1127 - (void)socktype;
1094 +void *poll_default_add_callback(POLLINFO *pi, short int *events) {
1095 + (void)pi;
1096 (void)events;
1129 - (void)client_ip;
1130 - (void)client_port;
1097
1098 return NULL;
1099 }
1134 -static void del_callback_default(int fd, int socktype, void *data) {
1135 - (void)fd;
1136 - (void)socktype;
1137 - (void)data;
1100
1139 - if(data)
1101 +void poll_default_del_callback(POLLINFO *pi) {
1102 + if(pi->data)
1103 error("POLLFD: internal error: del_callback_default() called with data pointer - possible memory leak");
1104 }
1105
1143 -static int rcv_callback_default(int fd, int socktype, void *data, short int *events) {
1144 - (void)socktype;
1145 - (void)data;
1106 +int poll_default_rcv_callback(POLLINFO *pi, short int *events) {
1107 (void)events;
1108
1109 char buffer[1024 + 1];
1110
1111 ssize_t rc;
1112 do {
1152 - rc = recv(fd, buffer, 1024, MSG_DONTWAIT);
1113 + rc = recv(pi->fd, buffer, 1024, MSG_DONTWAIT);
1114 if (rc < 0) {
1115 // read failed
1116 if (errno != EWOULDBLOCK && errno != EAGAIN) {
@@ -1158,48 +1119,164 @@ static int rcv_callback_default(int fd, int socktype, void *data, short int *eve
1119 }
1120 } else if (rc) {
1121 // data received
1161 - info("POLLFD: internal error: discarding %zd bytes received on socket %d", rc, fd);
1122 + info("POLLFD: internal error: discarding %zd bytes received on socket %d", rc, pi->fd);
1123 }
1124 } while (rc != -1);
1125
1126 return 0;
1127 }
1128
1168 -static int snd_callback_default(int fd, int socktype, void *data, short int *events) {
1169 - (void)socktype;
1170 - (void)data;
1171 - (void)events;
1172 -
1129 +int poll_default_snd_callback(POLLINFO *pi, short int *events) {
1130 *events &= ~POLLOUT;
1131
1175 - info("POLLFD: internal error: nothing to send on socket %d", fd);
1132 + info("POLLFD: internal error: nothing to send on socket %d", pi->fd);
1133 return 0;
1134 }
1135
1136 static void poll_events_cleanup(void *data) {
1180 - struct poll *p = (struct poll *)data;
1137 + POLLJOB *p = (POLLJOB *)data;
1138
1139 size_t i;
1140 for(i = 0 ; i <= p->max ; i++) {
1184 - struct pollinfo *pi = &p->inf[i];
1185 - poll_close_fd(p, pi);
1141 + POLLINFO *pi = &p->inf[i];
1142 + poll_close_fd(pi);
1143 }
1144
1145 freez(p->fds);
1146 freez(p->inf);
1147 }
1148
1149 +static void poll_events_process(POLLJOB *p, POLLINFO *pi, struct pollfd *pf, short int revents) {
1150 + short int events = pf->events;
1151 + int fd = pf->fd;
1152 + pf->revents = 0;
1153 + size_t i = pi->slot;
1154 +
1155 + if(unlikely(fd == -1)) {
1156 + debug(D_POLLFD, "POLLFD: LISTENER: ignoring slot %zu, it does not have an fd", i);
1157 + return;
1158 + }
1159 +
1160 + debug(D_POLLFD, "POLLFD: LISTENER: processing events for slot %zu (events = %d, revents = %d)", i, events, revents);
1161 +
1162 + if(revents & POLLIN || revents & POLLPRI) {
1163 + // receiving data
1164 +
1165 + if(likely(pi->flags & POLLINFO_FLAG_SERVER_SOCKET)) {
1166 + // new connection
1167 + // debug(D_POLLFD, "POLLFD: LISTENER: accepting connections from slot %zu (fd %d)", i, fd);
1168 +
1169 + switch(pi->socktype) {
1170 + case SOCK_STREAM: {
1171 + // a TCP socket
1172 + // we accept the connection
1173 +
1174 + int nfd;
1175 + do {
1176 + char client_ip[NI_MAXHOST + 1];
1177 + char client_port[NI_MAXSERV + 1];
1178 +
1179 + debug(D_POLLFD, "POLLFD: LISTENER: calling accept4() slot %zu (fd %d)", i, fd);
1180 + nfd = accept_socket(fd, SOCK_NONBLOCK, client_ip, NI_MAXHOST + 1, client_port, NI_MAXSERV + 1, p->access_list);
1181 + if (unlikely(nfd < 0)) {
1182 + // accept failed
1183 +
1184 + debug(D_POLLFD, "POLLFD: LISTENER: accept4() slot %zu (fd %d) failed.", i, fd);
1185 +
1186 + if(errno != EWOULDBLOCK && errno != EAGAIN)
1187 + error("POLLFD: LISTENER: accept() failed.");
1188 +
1189 + break;
1190 + }
1191 + else {
1192 + // accept ok
1193 + // info("POLLFD: LISTENER: client '[%s]:%s' connected to '%s' on fd %d", client_ip, client_port, sockets->fds_names[i], nfd);
1194 + poll_add_fd(p, nfd, SOCK_STREAM, POLLIN, POLLINFO_FLAG_CLIENT_SOCKET, client_ip, client_port);
1195 +
1196 + // it may have reallocated them, so refresh our pointers
1197 + pf = &p->fds[i];
1198 + pi = &p->inf[i];
1199 + }
1200 + } while (nfd >= 0);
1201 + break;
1202 + }
1203 +
1204 + case SOCK_DGRAM: {
1205 + // a UDP socket
1206 + // we read data from the server socket
1207 +
1208 + debug(D_POLLFD, "POLLFD: LISTENER: reading data from UDP slot %zu (fd %d)", i, fd);
1209 +
1210 + // FIXME: access_list is not applied to UDP
1211 +
1212 + pf->events = 0;
1213 + pi->rcv_callback(pi, &pf->events);
1214 + break;
1215 + }
1216 +
1217 + default: {
1218 + error("POLLFD: LISTENER: Unknown socktype %d on slot %zu", pi->socktype, pi->slot);
1219 + break;
1220 + }
1221 + }
1222 + }
1223 +
1224 + if(likely(pi->flags & POLLINFO_FLAG_CLIENT_SOCKET)) {
1225 + // read data from client TCP socket
1226 + debug(D_POLLFD, "POLLFD: LISTENER: reading data from TCP client slot %zu (fd %d)", i, fd);
1227 +
1228 + pf->events = 0;
1229 + if (pi->rcv_callback(pi, &pf->events) == -1) {
1230 + poll_close_fd(pi);
1231 + return;
1232 + }
1233 + }
1234 + }
1235 +
1236 + if(unlikely(revents & POLLOUT)) {
1237 + // sending data
1238 + debug(D_POLLFD, "POLLFD: LISTENER: sending data to socket on slot %zu (fd %d)", i, fd);
1239 +
1240 + pf->events = 0;
1241 + if (pi->snd_callback(pi, &pf->events) == -1) {
1242 + poll_close_fd(pi);
1243 + return;
1244 + }
1245 + }
1246 +
1247 + if(unlikely(revents & POLLERR)) {
1248 + error("POLLFD: LISTENER: processing POLLERR events for slot %zu fd %d (events = %d, revents = %d)", i, events, revents, fd);
1249 + pf->events = 0;
1250 + poll_close_fd(pi);
1251 + return;
1252 + }
1253 +
1254 + if(unlikely(revents & POLLHUP)) {
1255 + error("POLLFD: LISTENER: processing POLLHUP events for slot %zu fd %d (events = %d, revents = %d)", i, events, revents, fd);
1256 + pf->events = 0;
1257 + poll_close_fd(pi);
1258 + return;
1259 + }
1260 +
1261 + if(unlikely(revents & POLLNVAL)) {
1262 + error("POLLFD: LISTENER: processing POLLNVAL events for slot %zu fd %d (events = %d, revents = %d)", i, events, revents, fd);
1263 + pf->events = 0;
1264 + poll_close_fd(pi);
1265 + return;
1266 + }
1267 +}
1268 +
1269 void poll_events(LISTEN_SOCKETS *sockets
1193 - , void *(*add_callback)(int fd, int socktype, short int *events, const char *client_ip, const char *client_port)
1194 - , void (*del_callback)(int fd, int socktype, void *data)
1195 - , int (*rcv_callback)(int fd, int socktype, void *data, short int *events)
1196 - , int (*snd_callback)(int fd, int socktype, void *data, short int *events)
1270 + , void *(*add_callback)(POLLINFO *pi, short int *events)
1271 + , void (*del_callback)(POLLINFO *pi)
1272 + , int (*rcv_callback)(POLLINFO *pi, short int *events)
1273 + , int (*snd_callback)(POLLINFO *pi, short int *events)
1274 , SIMPLE_PATTERN *access_list
1275 , void *data
1276 ) {
1277 int retval;
1278
1202 - struct poll p = {
1279 + POLLJOB p = {
1280 .slots = 0,
1281 .used = 0,
1282 .max = 0,
@@ -1207,15 +1284,17 @@ void poll_events(LISTEN_SOCKETS *sockets
1284 .inf = NULL,
1285 .first_free = NULL,
1286
1210 - .add_callback = add_callback?add_callback:add_callback_default,
1211 - .del_callback = del_callback?del_callback:del_callback_default,
1212 - .rcv_callback = rcv_callback?rcv_callback:rcv_callback_default,
1213 - .snd_callback = snd_callback?snd_callback:snd_callback_default
1287 + .access_list = access_list,
1288 +
1289 + .add_callback = add_callback?add_callback:poll_default_add_callback,
1290 + .del_callback = del_callback?del_callback:poll_default_del_callback,
1291 + .rcv_callback = rcv_callback?rcv_callback:poll_default_rcv_callback,
1292 + .snd_callback = snd_callback?snd_callback:poll_default_snd_callback
1293 };
1294
1295 size_t i;
1296 for(i = 0; i < sockets->opened ;i++) {
1218 - struct pollinfo *pi = poll_add_fd(&p, sockets->fds[i], sockets->fds_types[i], POLLIN, POLLINFO_FLAG_SERVER_SOCKET, (sockets->fds_names[i])?sockets->fds_names[i]:"UNKNOWN", "");
1297 + POLLINFO *pi = poll_add_fd(&p, sockets->fds[i], sockets->fds_types[i], POLLIN, POLLINFO_FLAG_SERVER_SOCKET, (sockets->fds_names[i])?sockets->fds_names[i]:"UNKNOWN", "");
1298 pi->data = data;
1299 info("POLLFD: LISTENER: listening on '%s'", (sockets->fds_names[i])?sockets->fds_names[i]:"UNKNOWN");
1300 }
@@ -1224,9 +1303,7 @@ void poll_events(LISTEN_SOCKETS *sockets
1303
1304 netdata_thread_cleanup_push(poll_events_cleanup, &p);
1305
1227 - for(;;) {
1228 - if(unlikely(netdata_exit)) break;
1229 -
1306 + while(!netdata_exit) {
1307 debug(D_POLLFD, "POLLFD: LISTENER: Waiting on %zu sockets...", p.max + 1);
1308 retval = poll(p.fds, p.max + 1, timeout);
1309
@@ -1243,123 +1320,9 @@ void poll_events(LISTEN_SOCKETS *sockets
1320
1321 for(i = 0 ; i <= p.max ; i++) {
1322 struct pollfd *pf = &p.fds[i];
1246 - struct pollinfo *pi = &p.inf[i];
1247 - int fd = pf->fd;
1248 - short int events = pf->events, revents = pf->revents;
1249 - pf->revents = 0;
1250 -
1251 - if(unlikely(fd == -1)) {
1252 - debug(D_POLLFD, "POLLFD: LISTENER: ignoring slot %zu, it does not have an fd", i);
1253 - continue;
1254 - }
1255 -
1256 - debug(D_POLLFD, "POLLFD: LISTENER: processing events for slot %zu (events = %d, revents = %d)", i, events, revents);
1257 -
1258 - if(revents & POLLIN || revents & POLLPRI) {
1259 - // receiving data
1260 -
1261 - if(likely(pi->flags & POLLINFO_FLAG_SERVER_SOCKET)) {
1262 - // new connection
1263 - // debug(D_POLLFD, "POLLFD: LISTENER: accepting connections from slot %zu (fd %d)", i, fd);
1264 -
1265 - switch(pi->socktype) {
1266 - case SOCK_STREAM: {
1267 - // a TCP socket
1268 - // we accept the connection
1269 -
1270 - int nfd;
1271 - do {
1272 - char client_ip[NI_MAXHOST + 1];
1273 - char client_port[NI_MAXSERV + 1];
1274 -
1275 - debug(D_POLLFD, "POLLFD: LISTENER: calling accept4() slot %zu (fd %d)", i, fd);
1276 - nfd = accept_socket(fd, SOCK_NONBLOCK, client_ip, NI_MAXHOST + 1, client_port, NI_MAXSERV + 1, access_list);
1277 - if (unlikely(nfd < 0)) {
1278 - // accept failed
1279 -
1280 - debug(D_POLLFD, "POLLFD: LISTENER: accept4() slot %zu (fd %d) failed.", i, fd);
1281 -
1282 - if(errno != EWOULDBLOCK && errno != EAGAIN)
1283 - error("POLLFD: LISTENER: accept() failed.");
1284 -
1285 - break;
1286 - }
1287 - else {
1288 - // accept ok
1289 - // info("POLLFD: LISTENER: client '[%s]:%s' connected to '%s' on fd %d", client_ip, client_port, sockets->fds_names[i], nfd);
1290 - poll_add_fd(&p, nfd, SOCK_STREAM, POLLIN, POLLINFO_FLAG_CLIENT_SOCKET, client_ip, client_port);
1291 -
1292 - // it may have reallocated them, so refresh our pointers
1293 - pf = &p.fds[i];
1294 - pi = &p.inf[i];
1295 - }
1296 - } while (nfd >= 0);
1297 - break;
1298 - }
1299 -
1300 - case SOCK_DGRAM: {
1301 - // a UDP socket
1302 - // we read data from the server socket
1303 -
1304 - debug(D_POLLFD, "POLLFD: LISTENER: reading data from UDP slot %zu (fd %d)", i, fd);
1305 -
1306 - // FIXME: access_list is not applied to UDP
1307 -
1308 - pf->events = 0;
1309 - pi->rcv_callback(fd, pi->socktype, pi->data, &pf->events);
1310 - break;
1311 - }
1312 -
1313 - default: {
1314 - error("POLLFD: LISTENER: Unknown socktype %d on slot %zu", pi->socktype, pi->slot);
1315 - break;
1316 - }
1317 - }
1318 - }
1319 -
1320 - if(likely(pi->flags & POLLINFO_FLAG_CLIENT_SOCKET)) {
1321 - // read data from client TCP socket
1322 - debug(D_POLLFD, "POLLFD: LISTENER: reading data from TCP client slot %zu (fd %d)", i, fd);
1323 -
1324 - pf->events = 0;
1325 - if (pi->rcv_callback(fd, pi->socktype, pi->data, &pf->events) == -1) {
1326 - poll_close_fd(&p, pi);
1327 - continue;
1328 - }
1329 - }
1330 - }
1331 -
1332 - if(unlikely(revents & POLLOUT)) {
1333 - // sending data
1334 - debug(D_POLLFD, "POLLFD: LISTENER: sending data to socket on slot %zu (fd %d)", i, fd);
1335 -
1336 - pf->events = 0;
1337 - if (pi->snd_callback(fd, pi->socktype, pi->data, &pf->events) == -1) {
1338 - poll_close_fd(&p, pi);
1339 - continue;
1340 - }
1341 - }
1342 -
1343 - if(unlikely(revents & POLLERR)) {
1344 - error("POLLFD: LISTENER: processing POLLERR events for slot %zu fd %d (events = %d, revents = %d)", i, events, revents, fd);
1345 - pf->events = 0;
1346 - poll_close_fd(&p, pi);
1347 - continue;
1348 - }
1349 -
1350 - if(unlikely(revents & POLLHUP)) {
1351 - error("POLLFD: LISTENER: processing POLLHUP events for slot %zu fd %d (events = %d, revents = %d)", i, events, revents, fd);
1352 - pf->events = 0;
1353 - poll_close_fd(&p, pi);
1354 - continue;
1355 - }
1356 -
1357 - if(unlikely(revents & POLLNVAL)) {
1358 - error("POLLFD: LISTENER: processing POLLNVAL events for slot %zu fd %d (events = %d, revents = %d)", i, events, revents, fd);
1359 - pf->events = 0;
1360 - poll_close_fd(&p, pi);
1361 - continue;
1362 - }
1323 + short int revents = pf->revents;
1324 + if(unlikely(revents))
1325 + poll_events_process(&p, &p.inf[i], pf, revents);
1326 }
1327 }
1328
src/socket.h
+58 -4
@@ -53,11 +53,65 @@ extern int accept4(int sock, struct sockaddr *addr, socklen_t *addrlen, int flag
53 #endif /* #ifndef HAVE_ACCEPT4 */
54
55
56 +// ----------------------------------------------------------------------------
57 +// poll() based listener
58 +
59 +#define POLLINFO_FLAG_SERVER_SOCKET 0x00000001
60 +#define POLLINFO_FLAG_CLIENT_SOCKET 0x00000002
61 +
62 +typedef struct pollinfo {
63 + struct poll *p; // the parent
64 + size_t slot; // the slot id
65 +
66 + int fd; // the file descriptor
67 + int socktype; // the client socket type
68 + char *client_ip; // the connected client IP
69 + char *client_port; // the connected client port
70 +
71 + uint32_t flags; // internal flags
72 +
73 + // callbacks for this socket
74 + void (*del_callback)(struct pollinfo *pi);
75 + int (*rcv_callback)(struct pollinfo *pi, short int *events);
76 + int (*snd_callback)(struct pollinfo *pi, short int *events);
77 +
78 + // the user data
79 + void *data;
80 +
81 + // linking of free pollinfo structures
82 + // for quickly finding the next available
83 + // this is like a stack, it grows and shrinks
84 + // (with gaps - lower empty slots are preferred)
85 + struct pollinfo *next;
86 +} POLLINFO;
87 +
88 +typedef struct poll {
89 + size_t slots;
90 + size_t used;
91 + size_t min;
92 + size_t max;
93 + struct pollfd *fds;
94 + struct pollinfo *inf;
95 + struct pollinfo *first_free;
96 +
97 + SIMPLE_PATTERN *access_list;
98 +
99 + void *(*add_callback)(POLLINFO *pi, short int *events);
100 + void (*del_callback)(POLLINFO *pi);
101 + int (*rcv_callback)(POLLINFO *pi, short int *events);
102 + int (*snd_callback)(POLLINFO *pi, short int *events);
103 +} POLLJOB;
104 +
105 +extern int poll_default_snd_callback(POLLINFO *pi, short int *events);
106 +extern int poll_default_rcv_callback(POLLINFO *pi, short int *events);
107 +extern void poll_default_del_callback(POLLINFO *pi);
108 +extern void *poll_default_add_callback(POLLINFO *pi, short int *events);
109 +
110 extern void poll_events(LISTEN_SOCKETS *sockets
57 - , void *(*add_callback)(int fd, int socktype, short int *events, const char *client_ip, const char *client_port)
58 - , void (*del_callback)(int fd, int socktype, void *data)
59 - , int (*rcv_callback)(int fd, int socktype, void *data, short int *events)
60 - , int (*snd_callback)(int fd, int socktype, void *data, short int *events)
111 + , void *(*add_callback)(POLLINFO *pi, short int *events)
112 + , void (*del_callback)(POLLINFO *pi)
113 + , int (*rcv_callback)(POLLINFO *pi, short int *events)
114 + , int (*snd_callback)(POLLINFO *pi, short int *events)
115 , SIMPLE_PATTERN *access_list
116 , void *data
117 );
src/statsd.c
+19 -24
@@ -697,28 +697,23 @@ struct statsd_udp {
697 #endif
698
699 // new TCP client connected
700 -static void *statsd_add_callback(int fd, int socktype, short int *events, const char *client_ip, const char *client_port) {
701 - (void)fd;
702 - (void)socktype;
703 - (void)client_ip;
704 - (void)client_port;
700 +static void *statsd_add_callback(POLLINFO *pi, short int *events) {
701 + (void)pi;
702
703 *events = POLLIN;
704
708 - struct statsd_tcp *data = (struct statsd_tcp *)callocz(sizeof(struct statsd_tcp) + STATSD_TCP_BUFFER_SIZE, 1);
709 - data->type = STATSD_SOCKET_DATA_TYPE_TCP;
710 - data->size = STATSD_TCP_BUFFER_SIZE - 1;
705 + struct statsd_tcp *t = (struct statsd_tcp *)callocz(sizeof(struct statsd_tcp) + STATSD_TCP_BUFFER_SIZE, 1);
706 + t->type = STATSD_SOCKET_DATA_TYPE_TCP;
707 + t->size = STATSD_TCP_BUFFER_SIZE - 1;
708
712 - return data;
709 + return t;
710 }
711
712 // TCP client disconnected
716 -static void statsd_del_callback(int fd, int socktype, void *data) {
717 - (void)fd;
718 - (void)socktype;
713 +static void statsd_del_callback(POLLINFO *pi) {
714 + struct statsd_tcp *t = pi->data;
715
720 - if(data) {
721 - struct statsd_tcp *t = data;
716 + if(likely(t)) {
717 if(t->type == STATSD_SOCKET_DATA_TYPE_TCP) {
718 if(t->len != 0) {
719 statsd.socket_errors++;
@@ -729,17 +724,19 @@ static void statsd_del_callback(int fd, int socktype, void *data) {
724 else
725 error("STATSD: internal error: received socket data type is %d, but expected %d", (int)t->type, (int)STATSD_SOCKET_DATA_TYPE_TCP);
726
732 - freez(data);
727 + freez(t);
728 }
729 }
730
731 // Receive data
737 -static int statsd_rcv_callback(int fd, int socktype, void *data, short int *events) {
732 +static int statsd_rcv_callback(POLLINFO *pi, short int *events) {
733 *events = POLLIN;
734
740 - switch(socktype) {
735 + int fd = pi->fd;
736 +
737 + switch(pi->socktype) {
738 case SOCK_STREAM: {
742 - struct statsd_tcp *d = (struct statsd_tcp *)data;
739 + struct statsd_tcp *d = (struct statsd_tcp *)pi->data;
740 if(unlikely(!d)) {
741 error("STATSD: internal error: expected TCP data pointer is NULL");
742 statsd.socket_errors++;
@@ -791,7 +788,7 @@ static int statsd_rcv_callback(int fd, int socktype, void *data, short int *even
788 }
789
790 case SOCK_DGRAM: {
794 - struct statsd_udp *d = (struct statsd_udp *)data;
791 + struct statsd_udp *d = (struct statsd_udp *)pi->data;
792 if(unlikely(!d)) {
793 error("STATSD: internal error: expected UDP data pointer is NULL");
794 statsd.socket_errors++;
@@ -856,7 +853,7 @@ static int statsd_rcv_callback(int fd, int socktype, void *data, short int *even
853 }
854
855 default: {
859 - error("STATSD: internal error: unknown socktype %d on socket %d", socktype, fd);
856 + error("STATSD: internal error: unknown socktype %d on socket %d", pi->socktype, fd);
857 statsd.socket_errors++;
858 return -1;
859 }
@@ -865,10 +862,8 @@ static int statsd_rcv_callback(int fd, int socktype, void *data, short int *even
862 return 0;
863 }
864
868 -static int statsd_snd_callback(int fd, int socktype, void *data, short int *events) {
869 - (void)fd;
870 - (void)socktype;
871 - (void)data;
865 +static int statsd_snd_callback(POLLINFO *pi, short int *events) {
866 + (void)pi;
867 (void)events;
868
869 error("STATSD: snd_callback() called, but we never requested to send data to statsd clients.");
src/web_server.c
+13 -21
@@ -422,9 +422,7 @@ static struct web_server_static_threaded_worker *static_workers_private_data = N
422 static __thread struct web_server_static_threaded_worker *worker_private = NULL;
423
424 // new TCP client connected
425 -static void *web_server_add_callback(int fd, int socktype, short int *events, const char *client_ip, const char *client_port) {
426 - (void)fd;
427 - (void)socktype;
425 +static void *web_server_add_callback(POLLINFO *pi, short int *events) {
426
427 worker_private->connected++;
428 size_t concurrent = worker_private->connected - worker_private->disconnected;
@@ -433,10 +431,10 @@ static void *web_server_add_callback(int fd, int socktype, short int *events, co
431
432 *events = POLLIN;
433
436 - debug(D_WEB_CLIENT_ACCESS, "LISTENER on %d: new connection.", fd);
437 - struct web_client *w = web_client_create_on_fd(fd, client_ip, client_port);
434 + debug(D_WEB_CLIENT_ACCESS, "LISTENER on %d: new connection.", pi->fd);
435 + struct web_client *w = web_client_create_on_fd(pi->fd, pi->client_ip, pi->client_port);
436
439 - if(unlikely(socktype == AF_UNIX))
437 + if(unlikely(pi->socktype == AF_UNIX))
438 web_client_set_unix(w);
439 else
440 web_client_set_tcp(w);
@@ -445,13 +443,11 @@ static void *web_server_add_callback(int fd, int socktype, short int *events, co
443 }
444
445 // TCP client disconnected
448 -static void web_server_del_callback(int fd, int socktype, void *data) {
449 - (void)fd;
450 - (void)socktype;
451 -
446 +static void web_server_del_callback(POLLINFO *pi) {
447 worker_private->disconnected++;
448
454 - struct web_client *w = (struct web_client *)data;
449 + struct web_client *w = (struct web_client *)pi->data;
450 + int fd = pi->fd;
451
452 if(likely(w)) {
453 if(w->ofd == -1 || fd == w->ofd) {
@@ -474,13 +470,11 @@ static inline int web_server_check_client_status(struct web_client *w) {
470 }
471
472 // Receive data
477 -static int web_server_rcv_callback(int fd, int socktype, void *data, short int *events) {
478 - (void)fd;
479 - (void)socktype;
480 -
473 +static int web_server_rcv_callback(POLLINFO *pi, short int *events) {
474 worker_private->receptions++;
475
483 - struct web_client *w = (struct web_client *)data;
476 + struct web_client *w = (struct web_client *)pi->data;
477 + int fd = pi->fd;
478
479 if(unlikely(!web_client_has_wait_receive(w)))
480 return -1;
@@ -513,13 +507,11 @@ static int web_server_rcv_callback(int fd, int socktype, void *data, short int *
507 return web_server_check_client_status(w);
508 }
509
516 -static int web_server_snd_callback(int fd, int socktype, void *data, short int *events) {
517 - (void)fd;
518 - (void)socktype;
519 -
510 +static int web_server_snd_callback(POLLINFO *pi, short int *events) {
511 worker_private->sends++;
512
522 - struct web_client *w = (struct web_client *)data;
513 + struct web_client *w = (struct web_client *)pi->data;
514 + int fd = pi->fd;
515
516 if(unlikely(!web_client_has_wait_send(w)))
517 return -1;