do not try to reconnect too soon
Costa Tsaousis committed
Feb 21, 2017 at 14:28 UTC
0e4b7907f2bf29a86c0e5245812394d2f0fc048a
3 files changed
+53
-39
src/log.c
+1
-1
@@ -257,7 +257,7 @@ void info_int( const char *file, const char *function, const unsigned long line,
257
log_date(stderr);
258
259
va_start( args, fmt );
260
- if(debug_flags) fprintf(stderr, "%s: INFO: (%04lu@%-10.10s:%-15.15s):", program_name, line, file, function);
260
+ if(debug_flags) fprintf(stderr, "%s: INFO: (%04lu@%-10.10s:%-15.15s): ", program_name, line, file, function);
261
else fprintf(stderr, "%s: INFO: ", program_name);
262
vfprintf( stderr, fmt, args );
263
va_end( args );
src/rrdpush.c
+43
-30
@@ -104,40 +104,43 @@ void rrdset_done_push(RRDSET *st) {
104
if(unlikely(!rrdset_flag_check(st, RRDSET_FLAG_ENABLED)))
105
return;
106
107
+
108
+ rrdpush_lock();
109
+
110
if(unlikely(!rrdpush_buffer || !rrdpush_connected)) {
111
if(!error_shown)
109
- error("PUSH: not ready - discarding collected metrics.");
112
+ error("STREAM: not ready - discarding collected metrics.");
113
114
error_shown = 1;
115
+
116
+ rrdpush_unlock();
117
return;
118
}
119
error_shown = 0;
120
116
- rrdpush_lock();
117
- rrdset_rdlock(st);
118
-
121
if(st->rrdhost != last_host) {
122
buffer_sprintf(rrdpush_buffer, "HOST '%s' '%s'\n", st->rrdhost->machine_guid, st->rrdhost->hostname);
123
last_host = st->rrdhost;
124
}
125
126
+ rrdset_rdlock(st);
127
if(need_to_send_chart_definition(st))
128
send_chart_definition(st);
129
130
send_chart_metrics(st);
131
+ rrdset_unlock(st);
132
133
// signal the sender there are more data
134
if(write(rrdpush_pipe[PIPE_WRITE], " ", 1) == -1)
131
- error("Cannot write to internal pipe");
135
+ error("STREAM: cannot write to internal pipe");
136
133
- rrdset_unlock(st);
137
rrdpush_unlock();
138
}
139
140
static inline void rrdpush_flush(void) {
141
rrdpush_lock();
142
if(buffer_strlen(rrdpush_buffer))
140
- error("PUSH: discarding %zu bytes of metrics data already in the buffer.", buffer_strlen(rrdpush_buffer));
143
+ error("STREAM: discarding %zu bytes of metrics data already in the buffer.", buffer_strlen(rrdpush_buffer));
144
145
buffer_flush(rrdpush_buffer);
146
reset_all_charts();
@@ -148,19 +151,19 @@ static inline void rrdpush_flush(void) {
151
void *central_netdata_push_thread(void *ptr) {
152
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
153
151
- info("Central netdata push thread created with task id %d", gettid());
154
+ info("STREAM: central netdata push thread created with task id %d", gettid());
155
156
if(pthread_setcanceltype(PTHREAD_CANCEL_DEFERRED, NULL) != 0)
154
- error("Cannot set pthread cancel type to DEFERRED.");
157
+ error("STREAM: cannot set pthread cancel type to DEFERRED.");
158
159
if(pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
157
- error("Cannot set pthread cancel state to ENABLE.");
160
+ error("STREAM: cannot set pthread cancel state to ENABLE.");
161
162
163
rrdpush_buffer = buffer_create(1);
164
165
if(pipe(rrdpush_pipe) == -1)
163
- fatal("Cannot create required pipe.");
166
+ fatal("STREAM: cannot create required pipe.");
167
168
struct timeval tv = {
169
.tv_sec = 60,
@@ -176,6 +179,7 @@ void *central_netdata_push_thread(void *ptr) {
179
int sock = -1;
180
181
struct pollfd fds[2], *ifd, *ofd;
182
+ nfds_t fdmax;
183
184
ifd = &fds[0];
185
ofd = &fds[1];
@@ -184,32 +188,37 @@ void *central_netdata_push_thread(void *ptr) {
188
if(netdata_exit) break;
189
190
if(unlikely(sock == -1)) {
191
+ // stop appending data into rrdpush_buffer
192
+ // they will be lost, so there is no point to do it
193
rrdpush_connected = 0;
194
189
- info("PUSH: connecting to central netdata at: %s", central_netdata_to_push_data);
195
+ info("STREAM: connecting to central netdata at: %s", central_netdata_to_push_data);
196
sock = connect_to_one_of(central_netdata_to_push_data, 19999, &tv, &reconnects_counter);
197
198
if(unlikely(sock == -1)) {
193
- error("PUSH: failed to connect to central netdata at: %s", central_netdata_to_push_data);
199
+ error("STREAM: failed to connect to central netdata at: %s", central_netdata_to_push_data);
200
+ sleep(5);
201
continue;
202
}
203
197
- info("PUSH: connected to central netdata at: %s", central_netdata_to_push_data);
204
+ info("STREAM: initializing communication to central netdata at: %s", central_netdata_to_push_data);
205
206
char http[1000 + 1];
207
snprintfz(http, 1000, "GET /stream?key=%s HTTP/1.1\r\nUser-Agent: netdata-push-service/%s\r\nAccept: */*\r\n\r\n", config_get("global", "central netdata api key", ""), program_version);
208
if(send_timeout(sock, http, strlen(http), 0, 60) == -1) {
209
close(sock);
210
sock = -1;
204
- error("PUSH: failed to send http header to netdata at: %s", central_netdata_to_push_data);
211
+ error("STREAM: failed to send http header to netdata at: %s", central_netdata_to_push_data);
212
sleep(5);
213
continue;
214
}
215
216
+ info("STREAM: Waiting for STREAM from central netdata at: %s", central_netdata_to_push_data);
217
+
218
if(recv_timeout(sock, http, 1000, 0, 60) == -1) {
219
close(sock);
220
sock = -1;
212
- error("PUSH: failed to receive OK from netdata at: %s", central_netdata_to_push_data);
221
+ error("STREAM: failed to receive STREAM from netdata at: %s", central_netdata_to_push_data);
222
sleep(5);
223
continue;
224
}
@@ -217,16 +226,20 @@ void *central_netdata_push_thread(void *ptr) {
226
if(strncmp(http, "STREAM", 6)) {
227
close(sock);
228
sock = -1;
220
- error("PUSH: netdata servers at %s, did not send STREAM", central_netdata_to_push_data);
229
+ error("STREAM: netdata servers at %s, did not send STREAM", central_netdata_to_push_data);
230
sleep(5);
231
continue;
232
}
233
234
+ info("STREAM: Established STREAM with central netdata at: %s - sending metrics...", central_netdata_to_push_data);
235
+
236
if(fcntl(sock, F_SETFL, O_NONBLOCK) < 0)
226
- error("PUSH: cannot set non-blocking mode for socket.");
237
+ error("STREAM: cannot set non-blocking mode for socket.");
238
239
rrdpush_flush();
240
sent_connection = 0;
241
+
242
+ // allow appending data into rrdpush_buffer
243
rrdpush_connected = 1;
244
}
245
@@ -235,15 +248,15 @@ void *central_netdata_push_thread(void *ptr) {
248
ifd->revents = 0;
249
250
ofd->fd = sock;
238
- ofd->events = POLLOUT;
251
ofd->revents = 0;
240
-
241
- nfds_t fdmax = 2;
242
-
243
- if(begin < buffer_strlen(rrdpush_buffer))
252
+ if(begin < buffer_strlen(rrdpush_buffer)) {
253
ofd->events = POLLOUT;
245
- else
254
+ fdmax = 2;
255
+ }
256
+ else {
257
ofd->events = 0;
258
+ fdmax = 1;
259
+ }
260
261
if(netdata_exit) break;
262
int retval = poll(fds, fdmax, 60 * 1000);
@@ -253,7 +266,7 @@ void *central_netdata_push_thread(void *ptr) {
266
if(errno == EAGAIN || errno == EINTR)
267
continue;
268
256
- error("PUSH: Failed to poll().");
269
+ error("STREAM: Failed to poll().");
270
close(sock);
271
sock = -1;
272
break;
@@ -266,11 +279,11 @@ void *central_netdata_push_thread(void *ptr) {
279
if(ifd->revents & POLLIN) {
280
char buffer[1000 + 1];
281
if(read(rrdpush_pipe[PIPE_READ], buffer, 1000) == -1)
269
- error("PUSH: Cannot read from internal pipe.");
282
+ error("STREAM: Cannot read from internal pipe.");
283
}
284
285
if(ofd->revents & POLLOUT && begin < buffer_strlen(rrdpush_buffer)) {
273
- // info("PUSH: send buffer is ready, sending %zu bytes starting at %zu", buffer_strlen(rrdpush_buffer) - begin, begin);
286
+ // info("STREAM: send buffer is ready, sending %zu bytes starting at %zu", buffer_strlen(rrdpush_buffer) - begin, begin);
287
288
// fprintf(stderr, "PUSH BEGIN\n");
289
// fwrite(&rrdpush_buffer->buffer[begin], 1, buffer_strlen(rrdpush_buffer) - begin, stderr);
@@ -280,7 +293,7 @@ void *central_netdata_push_thread(void *ptr) {
293
ssize_t ret = send(sock, &rrdpush_buffer->buffer[begin], buffer_strlen(rrdpush_buffer) - begin, MSG_DONTWAIT);
294
if(ret == -1) {
295
if(errno != EAGAIN && errno != EINTR) {
283
- error("PUSH: failed to send metrics to central netdata at %s. We have sent %zu bytes on this connection.", central_netdata_to_push_data, sent_connection);
296
+ error("STREAM: failed to send metrics to central netdata at %s. We have sent %zu bytes on this connection.", central_netdata_to_push_data, sent_connection);
297
close(sock);
298
sock = -1;
299
}
@@ -300,7 +313,7 @@ void *central_netdata_push_thread(void *ptr) {
313
// protection from overflow
314
if(rrdpush_buffer->len > max_size) {
315
errno = 0;
303
- error("PUSH: too many data pending. Buffer is %zu bytes long, %zu unsent. We have sent %zu bytes in total, %zu on this connection. Closing connection to flush the data.", rrdpush_buffer->len, rrdpush_buffer->len - begin, sent_bytes, sent_connection);
316
+ error("STREAM: too many data pending. Buffer is %zu bytes long, %zu unsent. We have sent %zu bytes in total, %zu on this connection. Closing connection to flush the data.", rrdpush_buffer->len, rrdpush_buffer->len - begin, sent_bytes, sent_connection);
317
if(sock != -1) {
318
close(sock);
319
sock = -1;
@@ -308,7 +321,7 @@ void *central_netdata_push_thread(void *ptr) {
321
}
322
}
323
311
- debug(D_WEB_CLIENT, "Central netdata push thread exits.");
324
+ debug(D_WEB_CLIENT, "STREAM: central netdata push thread exits.");
325
if(sock != -1) {
326
close(sock);
327
}
src/web_client.c
+9
-8
@@ -1689,14 +1689,14 @@ int web_client_stream_request(RRDHOST *host, struct web_client *w, char *url) {
1689
}
1690
1691
if(!key || !*key) {
1692
- error("STREAM request from client '%s:%s', without an API key. Forbidding access.", w->client_ip, w->client_port);
1692
+ error("STREAM [%s]:%s: request without an API key. Forbidding access.", w->client_ip, w->client_port);
1693
buffer_flush(w->response.data);
1694
buffer_sprintf(w->response.data, "You need an API key for this request.");
1695
return 401;
1696
}
1697
1698
if(!validate_stream_api_key(key)) {
1699
- error("STREAM request from client '%s:%s': API key '%s' is not allowed. Forbidding access.", w->client_ip, w->client_port, key);
1699
+ error("STREAM [%s]:%s: API key '%s' is not allowed. Forbidding access.", w->client_ip, w->client_port, key);
1700
buffer_flush(w->response.data);
1701
buffer_sprintf(w->response.data, "Your API key is not permitted access.");
1702
return 401;
@@ -1719,16 +1719,17 @@ int web_client_stream_request(RRDHOST *host, struct web_client *w, char *url) {
1719
snprintfz(cd.fullfilename, FILENAME_MAX, "%s:%s", w->client_ip, w->client_port);
1720
snprintfz(cd.cmd, PLUGINSD_CMD_MAX, "%s:%s", w->client_ip, w->client_port);
1721
1722
+ info("STREAM [%s]:%s: sending STREAM to initiate streaming...", w->client_ip, w->client_port);
1723
if(send_timeout(w->ifd, "STREAM", 6, 0, 60) != 6) {
1723
- error("Cannot send STREAM to netdata at %s:%s", w->client_ip, w->client_port);
1724
+ error("STREAM [%s]:%s: cannot send STREAM.", w->client_ip, w->client_port);
1725
buffer_flush(w->response.data);
1725
- buffer_sprintf(w->response.data, "STREAM failed to reply back with STREAM");
1726
+ buffer_sprintf(w->response.data, "Failed to reply back with STREAM");
1727
return 400;
1728
}
1729
1730
// remove the non-blocking flag from the socket
1731
if(fcntl(w->ifd, F_SETFL, fcntl(w->ifd, F_GETFL, 0) & ~O_NONBLOCK) == -1)
1731
- error("STREAM from '%s:%s': cannot remove the non-blocking flag from socket %d", w->client_ip, w->client_port, w->ifd);
1732
+ error("STREAM [%s]:%s: cannot remove the non-blocking flag from socket %d", w->client_ip, w->client_port, w->ifd);
1733
1734
/*
1735
char buffer[1000 + 1];
@@ -1742,16 +1743,16 @@ int web_client_stream_request(RRDHOST *host, struct web_client *w, char *url) {
1743
// convert the socket to a FILE *
1744
FILE *fp = fdopen(w->ifd, "r");
1745
if(!fp) {
1745
- error("STREAM from '%s:%s': failed to get a FILE for FD %d.", w->client_ip, w->client_port, w->ifd);
1746
+ error("STREAM [%s]:%s: failed to get a FILE for FD %d.", w->client_ip, w->client_port, w->ifd);
1747
buffer_flush(w->response.data);
1748
buffer_sprintf(w->response.data, "Failed to get a FILE for an FD.");
1749
return 500;
1750
}
1751
1752
// call the plugins.d processor to receive the metrics
1752
- info("STREAM connecting client '%s:%s' to plugins.d.", w->client_ip, w->client_port);
1753
+ info("STREAM [%s]:%s: connecting client to plugins.d.", w->client_ip, w->client_port);
1754
size_t count = pluginsd_process(host, &cd, fp, 1);
1754
- error("STREAM from '%s:%s': client disconnected.", w->client_ip, w->client_port);
1755
+ error("STREAM [%s]:%s: client disconnected.", w->client_ip, w->client_port);
1756
1757
// close all sockets, to let the socket worker we are done
1758
fclose(fp);