@cryptotaxi247 / netdata-1 / commits / fdfc8fa0b

Optimizations part 3 (#15293)

* use madvise to speed up indexing * collect all rrddim members into a collector structure * use tier 0 virtual point for storing last stored value * reorganize key fields in rrddim * remove fgets from pluginsd and replace it with read() * properly uncork the web server sockets * Revert "reorganize key fields in rrddim" This reverts commit 2d45fa3959087e05462d387ff115a260f3a04b60. * Revert "use tier 0 virtual point for storing last stored value" This reverts commit a576cdd377ad4778a3b8608cabbb7ea7bb19a3a8. * fix cork names * fix compilation warnings

Costa Tsaousis committed Jul 1, 2023 at 01:13 UTC fdfc8fa0b13414898d1ac7d6e51808b418b951de
23 files changed +302 -297
collectors/plugins.d/pluginsd_parser.c
+51 -64
@@ -1307,9 +1307,9 @@ static inline PARSER_RC pluginsd_replay_set(char **words, size_t num_words, PARS
1307 }
1308
1309 rrddim_store_metric(rd, parser->user.replay.end_time_ut, value, flags);
1310 - rd->last_collected_time.tv_sec = parser->user.replay.end_time;
1311 - rd->last_collected_time.tv_usec = 0;
1312 - rd->collections_counter++;
1310 + rd->collector.last_collected_time.tv_sec = parser->user.replay.end_time;
1311 + rd->collector.last_collected_time.tv_usec = 0;
1312 + rd->collector.counter++;
1313 }
1314 else {
1315 error_limit_static_global_var(erl, 1, 0);
@@ -1340,16 +1340,16 @@ static inline PARSER_RC pluginsd_replay_rrddim_collection_state(char **words, si
1340 RRDDIM *rd = pluginsd_acquire_dimension(host, st, dimension, PLUGINSD_KEYWORD_REPLAY_RRDDIM_STATE);
1341 if(!rd) return PLUGINSD_DISABLE_PLUGIN(parser, NULL, NULL);
1342
1343 - usec_t dim_last_collected_ut = (usec_t)rd->last_collected_time.tv_sec * USEC_PER_SEC + (usec_t)rd->last_collected_time.tv_usec;
1343 + usec_t dim_last_collected_ut = (usec_t)rd->collector.last_collected_time.tv_sec * USEC_PER_SEC + (usec_t)rd->collector.last_collected_time.tv_usec;
1344 usec_t last_collected_ut = last_collected_ut_str ? str2ull_encoded(last_collected_ut_str) : 0;
1345 if(last_collected_ut > dim_last_collected_ut) {
1346 - rd->last_collected_time.tv_sec = (time_t)(last_collected_ut / USEC_PER_SEC);
1347 - rd->last_collected_time.tv_usec = (last_collected_ut % USEC_PER_SEC);
1346 + rd->collector.last_collected_time.tv_sec = (time_t)(last_collected_ut / USEC_PER_SEC);
1347 + rd->collector.last_collected_time.tv_usec = (last_collected_ut % USEC_PER_SEC);
1348 }
1349
1350 - rd->last_collected_value = last_collected_value_str ? str2ll_encoded(last_collected_value_str) : 0;
1351 - rd->last_calculated_value = last_calculated_value_str ? str2ndd_encoded(last_calculated_value_str, NULL) : 0;
1352 - rd->last_stored_value = last_stored_value_str ? str2ndd_encoded(last_stored_value_str, NULL) : 0.0;
1350 + rd->collector.last_collected_value = last_collected_value_str ? str2ll_encoded(last_collected_value_str) : 0;
1351 + rd->collector.last_calculated_value = last_calculated_value_str ? str2ndd_encoded(last_calculated_value_str, NULL) : 0;
1352 + rd->collector.last_stored_value = last_stored_value_str ? str2ndd_encoded(last_stored_value_str, NULL) : 0.0;
1353
1354 return PARSER_RC_OK;
1355 }
@@ -1718,12 +1718,12 @@ static inline PARSER_RC pluginsd_set_v2(char **words, size_t num_words, PARSER *
1718 // store it
1719
1720 rrddim_store_metric(rd, parser->user.v2.end_time * USEC_PER_SEC, value, flags);
1721 - rd->last_collected_time.tv_sec = parser->user.v2.end_time;
1722 - rd->last_collected_time.tv_usec = 0;
1723 - rd->last_collected_value = collected_value;
1724 - rd->last_stored_value = value;
1725 - rd->last_calculated_value = value;
1726 - rd->collections_counter++;
1721 + rd->collector.last_collected_time.tv_sec = parser->user.v2.end_time;
1722 + rd->collector.last_collected_time.tv_usec = 0;
1723 + rd->collector.last_collected_value = collected_value;
1724 + rd->collector.last_stored_value = value;
1725 + rd->collector.last_calculated_value = value;
1726 + rd->collector.counter++;
1727 rrddim_set_updated(rd);
1728
1729 timing_step(TIMING_STEP_SET2_STORE);
@@ -1778,8 +1778,8 @@ static inline PARSER_RC pluginsd_end_v2(char **words __maybe_unused, size_t num_
1778
1779 RRDDIM *rd;
1780 rrddim_foreach_read(rd, st) {
1781 - rd->calculated_value = 0;
1782 - rd->collected_value = 0;
1781 + rd->collector.calculated_value = 0;
1782 + rd->collector.collected_value = 0;
1783 rrddim_clear_updated(rd);
1784 }
1785 rrddim_foreach_done(rd);
@@ -1849,76 +1849,59 @@ static inline PARSER_RC streaming_claimed_id(char **words, size_t num_words, PAR
1849
1850 // ----------------------------------------------------------------------------
1851
1852 -typedef enum {
1853 - PARSER_FGETS_RESULT_OK,
1854 - PARSER_FGETS_RESULT_TIMEOUT,
1855 - PARSER_FGETS_RESULT_ERROR,
1856 - PARSER_FGETS_RESULT_EOF,
1857 -} PARSER_FGETS_RESULT;
1852 +static inline bool buffered_reader_read(struct buffered_reader *reader, int fd) {
1853 +#ifdef NETDATA_INTERNAL_CHECKS
1854 + if(reader->read_buffer[reader->read_len] != '\0')
1855 + fatal("%s(): read_buffer does not start with zero", __FUNCTION__ );
1856 +#endif
1857
1859 -static inline PARSER_FGETS_RESULT parser_fgets(char *s, int size, FILE *stream) {
1860 - errno = 0;
1858 + ssize_t bytes_read = read(fd, reader->read_buffer + reader->read_len, sizeof(reader->read_buffer) - reader->read_len - 1);
1859 + if(unlikely(bytes_read <= 0))
1860 + return false;
1861 +
1862 + reader->read_len += bytes_read;
1863 + reader->read_buffer[reader->read_len] = '\0';
1864 +
1865 + return true;
1866 +}
1867
1868 +static inline bool buffered_reader_read_timeout(struct buffered_reader *reader, int fd, int timeout_ms) {
1869 + errno = 0;
1870 struct pollfd fds[1];
1863 - int timeout_msecs = 2 * 60 * MSEC_PER_SEC;
1871
1865 - fds[0].fd = fileno(stream);
1872 + fds[0].fd = fd;
1873 fds[0].events = POLLIN;
1874
1868 - int ret = poll(fds, 1, timeout_msecs);
1875 + int ret = poll(fds, 1, timeout_ms);
1876
1877 if (ret > 0) {
1878 /* There is data to read */
1872 - if (fds[0].revents & POLLIN) {
1873 - char *tmp = fgets(s, size, stream);
1874 -
1875 - if(unlikely(!tmp)) {
1876 - if (feof(stream)) {
1877 - error("PARSER: read failed: end of file.");
1878 - return PARSER_FGETS_RESULT_EOF;
1879 - }
1880 -
1881 - else if (ferror(stream)) {
1882 - error("PARSER: read failed: input error.");
1883 - return PARSER_FGETS_RESULT_ERROR;
1884 - }
1885 -
1886 - error("PARSER: read failed: unknown error.");
1887 - return PARSER_FGETS_RESULT_ERROR;
1888 - }
1879 + if (fds[0].revents & POLLIN)
1880 + return buffered_reader_read(reader, fd);
1881
1890 - return PARSER_FGETS_RESULT_OK;
1891 - }
1882 else if(fds[0].revents & POLLERR) {
1883 error("PARSER: read failed: POLLERR.");
1894 - return PARSER_FGETS_RESULT_ERROR;
1884 + return false;
1885 }
1886 else if(fds[0].revents & POLLHUP) {
1887 error("PARSER: read failed: POLLHUP.");
1898 - return PARSER_FGETS_RESULT_ERROR;
1888 + return false;
1889 }
1890 else if(fds[0].revents & POLLNVAL) {
1891 error("PARSER: read failed: POLLNVAL.");
1902 - return PARSER_FGETS_RESULT_ERROR;
1892 + return false;
1893 }
1894
1895 error("PARSER: poll() returned positive number, but POLLIN|POLLERR|POLLHUP|POLLNVAL are not set.");
1906 - return PARSER_FGETS_RESULT_ERROR;
1896 + return false;
1897 }
1898 else if (ret == 0) {
1899 error("PARSER: timeout while waiting for data.");
1910 - return PARSER_FGETS_RESULT_TIMEOUT;
1900 + return false;
1901 }
1902
1903 error("PARSER: poll() failed with code %d.", ret);
1914 - return PARSER_FGETS_RESULT_ERROR;
1915 -}
1916 -
1917 -static int parser_next(PARSER *parser, char *buffer, size_t buffer_size) {
1918 - if(likely(parser_fgets(buffer, (int)buffer_size, (FILE *)parser->fp_input) == PARSER_FGETS_RESULT_OK))
1919 - return 0;
1920 -
1921 - return 1;
1904 + return false;
1905 }
1906
1907 void pluginsd_process_thread_cleanup(void *ptr) {
@@ -1979,10 +1962,14 @@ inline size_t pluginsd_process(RRDHOST *host, struct plugind *cd, FILE *fp_plugi
1962 // so, parser needs to be allocated before pushing it
1963 netdata_thread_cleanup_push(pluginsd_process_thread_cleanup, parser);
1964
1982 - char buffer[PLUGINSD_LINE_MAX + 1];
1983 -
1984 - while (likely(!parser_next(parser, buffer, PLUGINSD_LINE_MAX))) {
1985 - if (unlikely(!service_running(SERVICE_COLLECTORS) || parser_action(parser, buffer)))
1965 + buffered_reader_init(&parser->reader);
1966 + char buffer[PLUGINSD_LINE_MAX + 2];
1967 + while(likely(service_running(SERVICE_COLLECTORS))) {
1968 + if (unlikely(!buffered_reader_next_line(&parser->reader, buffer, PLUGINSD_LINE_MAX + 2))) {
1969 + if(unlikely(!buffered_reader_read_timeout(&parser->reader, fileno((FILE *)parser->fp_input), 2 * 60 * MSEC_PER_SEC)))
1970 + break;
1971 + }
1972 + else if(unlikely(parser_action(parser, buffer)))
1973 break;
1974 }
1975
collectors/plugins.d/pluginsd_parser.h
+2
@@ -95,6 +95,8 @@ typedef struct parser {
95
96 PARSER_USER_OBJECT user; // User defined structure to hold extra state between calls
97
98 + struct buffered_reader reader;
99 +
100 struct {
101 const char *end_keyword;
102 BUFFER *response;
collectors/statsd.plugin/statsd.c
+2 -2
@@ -2189,13 +2189,13 @@ static inline RRDDIM *statsd_add_dim_to_app_chart(STATSD_APP *app, STATSD_APP_CH
2189
2190 dim->rd = rrddim_add(chart->st, metric, dim->name, dim->multiplier, dim->divisor, dim->algorithm);
2191 if(dim->flags != RRDDIM_FLAG_NONE) dim->rd->flags |= dim->flags;
2192 - if(dim->options != RRDDIM_OPTION_NONE) dim->rd->options |= dim->options;
2192 + if(dim->options != RRDDIM_OPTION_NONE) dim->rd->collector.options |= dim->options;
2193 return dim->rd;
2194 }
2195
2196 dim->rd = rrddim_add(chart->st, dim->metric, dim->name, dim->multiplier, dim->divisor, dim->algorithm);
2197 if(dim->flags != RRDDIM_FLAG_NONE) dim->rd->flags |= dim->flags;
2198 - if(dim->options != RRDDIM_OPTION_NONE) dim->rd->options |= dim->options;
2198 + if(dim->options != RRDDIM_OPTION_NONE) dim->rd->collector.options |= dim->options;
2199 return dim->rd;
2200 }
2201
daemon/service.c
+1 -1
@@ -87,7 +87,7 @@ static bool svc_rrdset_archive_obsolete_dimensions(RRDSET *st, bool all_dimensio
87 dfe_start_write(st->rrddim_root_index, rd) {
88 if(unlikely(
89 all_dimensions ||
90 - (rrddim_flag_check(rd, RRDDIM_FLAG_OBSOLETE) && (rd->last_collected_time.tv_sec + rrdset_free_obsolete_time_s < now))
90 + (rrddim_flag_check(rd, RRDDIM_FLAG_OBSOLETE) && (rd->collector.last_collected_time.tv_sec + rrdset_free_obsolete_time_s < now))
91 )) {
92
93 if(dictionary_acquired_item_references(rd_dfe.item) == 1) {
daemon/unit_test.c
+13 -13
@@ -1325,7 +1325,7 @@ int run_test(struct test *test)
1325 // align the first entry to second boundary
1326 if(!c) {
1327 fprintf(stderr, " > %s: fixing first collection time to be %llu microseconds to second boundary\n", test->name, test->feed[c].microseconds);
1328 - rd->last_collected_time.tv_usec = st->last_collected_time.tv_usec = st->last_updated.tv_usec = test->feed[c].microseconds;
1328 + rd->collector.last_collected_time.tv_usec = st->last_collected_time.tv_usec = st->last_updated.tv_usec = test->feed[c].microseconds;
1329 // time_start = st->last_collected_time.tv_sec;
1330 }
1331 }
@@ -1585,7 +1585,7 @@ int unit_test(long delay, long shift)
1585
1586 // prevent it from deleting the dimensions
1587 rrddim_foreach_read(rd, st) {
1588 - rd->last_collected_time.tv_sec = st->last_collected_time.tv_sec;
1588 + rd->collector.last_collected_time.tv_sec = st->last_collected_time.tv_sec;
1589 }
1590 rrddim_foreach_done(rd);
1591
@@ -1831,15 +1831,15 @@ int unit_test_bitmap256(void) {
1831 #ifdef ENABLE_DBENGINE
1832 static inline void rrddim_set_by_pointer_fake_time(RRDDIM *rd, collected_number value, time_t now)
1833 {
1834 - rd->last_collected_time.tv_sec = now;
1835 - rd->last_collected_time.tv_usec = 0;
1836 - rd->collected_value = value;
1834 + rd->collector.last_collected_time.tv_sec = now;
1835 + rd->collector.last_collected_time.tv_usec = 0;
1836 + rd->collector.collected_value = value;
1837 rrddim_set_updated(rd);
1838
1839 - rd->collections_counter++;
1839 + rd->collector.counter++;
1840
1841 collected_number v = (value >= 0) ? value : -value;
1842 - if(unlikely(v > rd->collected_value_max)) rd->collected_value_max = v;
1842 + if(unlikely(v > rd->collector.collected_value_max)) rd->collector.collected_value_max = v;
1843 }
1844
1845 static RRDHOST *dbengine_rrdhost_find_or_create(char *name)
@@ -1911,9 +1911,9 @@ static void test_dbengine_create_charts(RRDHOST *host, RRDSET *st[CHARTS], RRDDI
1911 // Initialize DB with the very first entries
1912 for (i = 0 ; i < CHARTS ; ++i) {
1913 for (j = 0 ; j < DIMS ; ++j) {
1914 - rd[i][j]->last_collected_time.tv_sec =
1914 + rd[i][j]->collector.last_collected_time.tv_sec =
1915 st[i]->last_collected_time.tv_sec = st[i]->last_updated.tv_sec = 2 * API_RELATIVE_TIME_MAX - 1;
1916 - rd[i][j]->last_collected_time.tv_usec =
1916 + rd[i][j]->collector.last_collected_time.tv_usec =
1917 st[i]->last_collected_time.tv_usec = st[i]->last_updated.tv_usec = 0;
1918 }
1919 }
@@ -1952,9 +1952,9 @@ static time_t test_dbengine_create_metrics(RRDSET *st[CHARTS], RRDDIM *rd[CHARTS
1952 for (j = 0 ; j < DIMS ; ++j) {
1953 storage_engine_store_change_collection_frequency(rd[i][j]->tiers[0].db_collection_handle, update_every);
1954
1955 - rd[i][j]->last_collected_time.tv_sec =
1955 + rd[i][j]->collector.last_collected_time.tv_sec =
1956 st[i]->last_collected_time.tv_sec = st[i]->last_updated.tv_sec = time_now;
1957 - rd[i][j]->last_collected_time.tv_usec =
1957 + rd[i][j]->collector.last_collected_time.tv_usec =
1958 st[i]->last_collected_time.tv_usec = st[i]->last_updated.tv_usec = 0;
1959 }
1960 }
@@ -2318,9 +2318,9 @@ static void generate_dbengine_chart(void *arg)
2318 // feed it with the test data
2319 time_current = time_present - history_seconds;
2320 for (j = 0 ; j < DSET_DIMS ; ++j) {
2321 - rd[j]->last_collected_time.tv_sec =
2321 + rd[j]->collector.last_collected_time.tv_sec =
2322 st->last_collected_time.tv_sec = st->last_updated.tv_sec = time_current - update_every;
2323 - rd[j]->last_collected_time.tv_usec =
2323 + rd[j]->collector.last_collected_time.tv_usec =
2324 st->last_collected_time.tv_usec = st->last_updated.tv_usec = 0;
2325 }
2326 for( ; !thread_info->done && time_current < time_present ; time_current += update_every) {
database/contexts/metric.c
+1 -1
@@ -33,7 +33,7 @@ inline NETDATA_DOUBLE rrdmetric_acquired_last_stored_value(RRDMETRIC_ACQUIRED *r
33 RRDMETRIC *rm = rrdmetric_acquired_value(rma);
34
35 if(rm->rrddim)
36 - return rm->rrddim->last_stored_value;
36 + return rm->rrddim->collector.last_stored_value;
37
38 return NAN;
39 }
database/engine/journalfile.c
+12 -7
@@ -216,15 +216,20 @@ static struct journal_v2_header *journalfile_v2_mounted_data_get(struct rrdengin
216 madvise_dontfork(journalfile->mmap.data, journalfile->mmap.size);
217 madvise_dontdump(journalfile->mmap.data, journalfile->mmap.size);
218
219 - // let the kernel know that we don't want read-ahead on this file
220 - madvise_random(journalfile->mmap.data, journalfile->mmap.size);
221 -
222 -// madvise_willneed(journalfile->mmap.data, journalfile->v2.size_of_directory);
223 -// madvise_dontneed(journalfile->mmap.data, journalfile->mmap.size);
224 -
219 spinlock_lock(&journalfile->v2.spinlock);
220 journalfile->v2.flags |= JOURNALFILE_FLAG_IS_AVAILABLE | JOURNALFILE_FLAG_IS_MOUNTED;
221 + JOURNALFILE_FLAGS flags = journalfile->v2.flags;
222 spinlock_unlock(&journalfile->v2.spinlock);
223 +
224 + if(flags & JOURNALFILE_FLAG_MOUNTED_FOR_RETENTION) {
225 + // we need the entire metrics directory into memory to process it
226 + madvise_willneed(journalfile->mmap.data, journalfile->v2.size_of_directory);
227 + }
228 + else {
229 + // let the kernel know that we don't want read-ahead on this file
230 + madvise_random(journalfile->mmap.data, journalfile->mmap.size);
231 + // madvise_dontneed(journalfile->mmap.data, journalfile->mmap.size);
232 + }
233 }
234 }
235
@@ -413,7 +418,7 @@ void journalfile_v2_data_set(struct rrdengine_journalfile *journalfile, int fd,
418 struct journal_v2_header *j2_header = journalfile->mmap.data;
419 journalfile->v2.first_time_s = (time_t)(j2_header->start_time_ut / USEC_PER_SEC);
420 journalfile->v2.last_time_s = (time_t)(j2_header->end_time_ut / USEC_PER_SEC);
416 - // journalfile->v2.size_of_directory = j2_header->metric_offset + j2_header->metric_count * sizeof(struct journal_metric_list);
421 + journalfile->v2.size_of_directory = j2_header->metric_offset + j2_header->metric_count * sizeof(struct journal_metric_list);
422
423 journalfile_v2_mounted_data_unmount(journalfile, true, true);
424
database/engine/journalfile.h
+1 -1
@@ -40,7 +40,7 @@ struct rrdengine_journalfile {
40 time_t first_time_s;
41 time_t last_time_s;
42 time_t not_needed_since_s;
43 - // uint32_t size_of_directory;
43 + uint32_t size_of_directory;
44 } v2;
45
46 struct {
database/rrd.h
+32 -33
@@ -241,9 +241,9 @@ typedef enum __attribute__ ((__packed__)) rrddim_options {
241 // this is 8-bit
242 } RRDDIM_OPTIONS;
243
244 -#define rrddim_option_check(rd, option) ((rd)->options & (option))
245 -#define rrddim_option_set(rd, option) (rd)->options |= (option)
246 -#define rrddim_option_clear(rd, option) (rd)->options &= ~(option)
244 +#define rrddim_option_check(rd, option) ((rd)->collector.options & (option))
245 +#define rrddim_option_set(rd, option) (rd)->collector.options |= (option)
246 +#define rrddim_option_clear(rd, option) (rd)->collector.options &= ~(option)
247
248 // flags are runtime changing status flags (atomics are required to alter/access them)
249 typedef enum __attribute__ ((__packed__)) rrddim_flags {
@@ -352,42 +352,19 @@ struct rrddim {
352 STRING *name; // the name of this dimension (as presented to user)
353
354 RRD_ALGORITHM algorithm; // the algorithm that is applied to add new collected values
355 - RRDDIM_OPTIONS options; // permanent configuration options
355 RRD_MEMORY_MODE rrd_memory_mode; // the memory mode for this dimension
356 RRDDIM_FLAGS flags; // run time changing status flags
357
358 int32_t multiplier; // the multiplier of the collected values
359 int32_t divisor; // the divider of the collected values
360
362 - uint32_t collections_counter; // the number of times we added values to this rrddim
363 -
361 // ------------------------------------------------------------------------
362 // operational state members
363
367 - rrd_ml_dimension_t *ml_dimension; // machine learning data about this dimension
368 -
369 - // ------------------------------------------------------------------------
370 - // linking to siblings and parents
371 -
364 struct rrdset *rrdset;
365 + rrd_ml_dimension_t *ml_dimension; // machine learning data about this dimension
366 RRDMETRIC_ACQUIRED *rrdmetric; // the rrdmetric of this dimension
367
375 - // ------------------------------------------------------------------------
376 - // data collection members
377 -
378 - struct timeval last_collected_time; // when was this dimension last updated
379 - // this is actual date time we updated the last_collected_value
380 - // THIS IS DIFFERENT FROM THE SAME MEMBER OF RRDSET
381 -
382 - collected_number collected_value_max; // the absolute maximum of the collected value
383 -
384 - NETDATA_DOUBLE calculated_value; // the current calculated value, after applying the algorithm - resets to zero after being used
385 - NETDATA_DOUBLE last_calculated_value; // the last calculated value processed
386 - NETDATA_DOUBLE last_stored_value; // the last value as stored in the database (after interpolation)
387 -
388 - collected_number collected_value; // the current value, as collected - resets to 0 after being used
389 - collected_number last_collected_value; // the last value that was collected, after being processed
390 -
368 #ifdef NETDATA_LOG_COLLECTION_ERRORS
369 usec_t rrddim_store_metric_last_ut; // the timestamp we last called rrddim_store_metric()
370 size_t rrddim_store_metric_count; // the rrddim_store_metric() counter
@@ -405,6 +382,28 @@ struct rrddim {
382 storage_number *data; // the array of values
383 } db;
384
385 + // ------------------------------------------------------------------------
386 + // data collection members
387 +
388 + struct {
389 + RRDDIM_OPTIONS options; // permanent configuration options
390 +
391 + uint32_t counter; // the number of times we added values to this rrddim
392 +
393 + collected_number collected_value; // the current value, as collected - resets to 0 after being used
394 + collected_number collected_value_max; // the absolute maximum of the collected value
395 + collected_number last_collected_value; // the last value that was collected, after being processed
396 +
397 + struct timeval last_collected_time; // when was this dimension last updated
398 + // this is actual date time we updated the last_collected_value
399 + // THIS IS DIFFERENT FROM THE SAME MEMBER OF RRDSET
400 +
401 + NETDATA_DOUBLE calculated_value; // the current calculated value, after applying the algorithm - resets to zero after being used
402 + NETDATA_DOUBLE last_calculated_value; // the last calculated value processed
403 +
404 + NETDATA_DOUBLE last_stored_value; // the last value as stored in the database (after interpolation)
405 + } collector;
406 +
407 // ------------------------------------------------------------------------
408
409 struct rrddim_tier tiers[]; // our tiers of databases
@@ -415,13 +414,13 @@ size_t rrddim_size(void);
414 #define rrddim_id(rd) string2str((rd)->id)
415 #define rrddim_name(rd) string2str((rd) ->name)
416
418 -#define rrddim_check_updated(rd) ((rd)->options & RRDDIM_OPTION_UPDATED)
419 -#define rrddim_set_updated(rd) (rd)->options |= RRDDIM_OPTION_UPDATED
420 -#define rrddim_clear_updated(rd) (rd)->options &= ~RRDDIM_OPTION_UPDATED
417 +#define rrddim_check_updated(rd) ((rd)->collector.options & RRDDIM_OPTION_UPDATED)
418 +#define rrddim_set_updated(rd) (rd)->collector.options |= RRDDIM_OPTION_UPDATED
419 +#define rrddim_clear_updated(rd) (rd)->collector.options &= ~RRDDIM_OPTION_UPDATED
420
422 -#define rrddim_check_exposed(rd) ((rd)->options & RRDDIM_OPTION_EXPOSED)
423 -#define rrddim_set_exposed(rd) (rd)->options |= RRDDIM_OPTION_EXPOSED
424 -#define rrddim_clear_exposed(rd) (rd)->options &= ~RRDDIM_OPTION_EXPOSED
421 +#define rrddim_check_exposed(rd) ((rd)->collector.options & RRDDIM_OPTION_EXPOSED)
422 +#define rrddim_set_exposed(rd) (rd)->collector.options |= RRDDIM_OPTION_EXPOSED
423 +#define rrddim_clear_exposed(rd) (rd)->collector.options &= ~RRDDIM_OPTION_EXPOSED
424
425 // returns the RRDDIM cache filename, or NULL if it does not exist
426 const char *rrddim_cache_filename(RRDDIM *rd);
database/rrddim.c
+11 -11
@@ -49,7 +49,7 @@ static void rrddim_insert_callback(const DICTIONARY_ITEM *item __maybe_unused, v
49 rd->rrdset = st;
50
51 if(rrdset_flag_check(st, RRDSET_FLAG_STORE_FIRST))
52 - rd->collections_counter = 1;
52 + rd->collector.counter = 1;
53
54 if(ctr->memory_mode == RRD_MEMORY_MODE_MAP || ctr->memory_mode == RRD_MEMORY_MODE_SAVE) {
55 if(!rrddim_memory_load_or_create_map_save(st, rd, ctr->memory_mode)) {
@@ -565,16 +565,16 @@ inline collected_number rrddim_set_by_pointer(RRDSET *st, RRDDIM *rd, collected_
565 collected_number rrddim_timed_set_by_pointer(RRDSET *st __maybe_unused, RRDDIM *rd, struct timeval collected_time, collected_number value) {
566 debug(D_RRD_CALLS, "rrddim_set_by_pointer() for chart %s, dimension %s, value " COLLECTED_NUMBER_FORMAT, rrdset_name(st), rrddim_name(rd), value);
567
568 - rd->last_collected_time = collected_time;
569 - rd->collected_value = value;
568 + rd->collector.last_collected_time = collected_time;
569 + rd->collector.collected_value = value;
570 rrddim_set_updated(rd);
571 - rd->collections_counter++;
571 + rd->collector.counter++;
572
573 collected_number v = (value >= 0) ? value : -value;
574 - if (unlikely(v > rd->collected_value_max))
575 - rd->collected_value_max = v;
574 + if (unlikely(v > rd->collector.collected_value_max))
575 + rd->collector.collected_value_max = v;
576
577 - return rd->last_collected_value;
577 + return rd->collector.last_collected_value;
578 }
579
580
@@ -644,9 +644,9 @@ void rrddim_memory_file_update(RRDDIM *rd) {
644 if(!rd || !rd->db.rd_on_file) return;
645 struct rrddim_map_save_v019 *rd_on_file = rd->db.rd_on_file;
646
647 - rd_on_file->last_collected_time.tv_sec = rd->last_collected_time.tv_sec;
648 - rd_on_file->last_collected_time.tv_usec = rd->last_collected_time.tv_usec;
649 - rd_on_file->last_collected_value = rd->last_collected_value;
647 + rd_on_file->last_collected_time.tv_sec = rd->collector.last_collected_time.tv_sec;
648 + rd_on_file->last_collected_time.tv_usec = rd->collector.last_collected_time.tv_usec;
649 + rd_on_file->last_collected_value = rd->collector.last_collected_value;
650 }
651
652 void rrddim_memory_file_free(RRDDIM *rd) {
@@ -727,7 +727,7 @@ bool rrddim_memory_load_or_create_map_save(RRDSET *st, RRDDIM *rd, RRD_MEMORY_MO
727 }
728
729 if(!reset) {
730 - rd->last_collected_value = rd_on_file->last_collected_value;
730 + rd->collector.last_collected_value = rd_on_file->last_collected_value;
731
732 if(rd_on_file->algorithm != rd->algorithm)
733 netdata_log_info("File %s does not have the expected algorithm (expected %u '%s', found %u '%s'). Previous values may be wrong.",
database/rrdset.c
+78 -78
@@ -759,9 +759,9 @@ void rrdset_reset(RRDSET *st) {
759
760 RRDDIM *rd;
761 rrddim_foreach_read(rd, st) {
762 - rd->last_collected_time.tv_sec = 0;
763 - rd->last_collected_time.tv_usec = 0;
764 - rd->collections_counter = 0;
762 + rd->collector.last_collected_time.tv_sec = 0;
763 + rd->collector.last_collected_time.tv_usec = 0;
764 + rd->collector.counter = 0;
765
766 if(!rrddim_flag_check(rd, RRDDIM_FLAG_ARCHIVED)) {
767 for(size_t tier = 0; tier < storage_tiers ;tier++)
@@ -1361,7 +1361,7 @@ static inline size_t rrdset_done_interpolate(
1361 switch(rd->algorithm) {
1362 case RRD_ALGORITHM_INCREMENTAL:
1363 new_value = (NETDATA_DOUBLE)
1364 - ( rd->calculated_value
1364 + ( rd->collector.calculated_value
1365 * (NETDATA_DOUBLE)(next_store_ut - last_collect_ut)
1366 / (NETDATA_DOUBLE)(now_collect_ut - last_collect_ut)
1367 );
@@ -1372,14 +1372,14 @@ static inline size_t rrdset_done_interpolate(
1372 " / (%llu - %llu)"
1373 , rrddim_name(rd)
1374 , new_value
1375 - , rd->calculated_value
1375 + , rd->collector.calculated_value
1376 , next_store_ut, last_collect_ut
1377 , now_collect_ut, last_collect_ut
1378 );
1379
1380 - rd->calculated_value -= new_value;
1381 - new_value += rd->last_calculated_value;
1382 - rd->last_calculated_value = 0;
1380 + rd->collector.calculated_value -= new_value;
1381 + new_value += rd->collector.last_calculated_value;
1382 + rd->collector.last_calculated_value = 0;
1383 new_value /= (NETDATA_DOUBLE)st->update_every;
1384
1385 if(unlikely(next_store_ut - last_stored_ut < update_every_ut)) {
@@ -1402,18 +1402,18 @@ static inline size_t rrdset_done_interpolate(
1402 // do not interpolate
1403 // just show the calculated value
1404
1405 - new_value = rd->calculated_value;
1405 + new_value = rd->collector.calculated_value;
1406 }
1407 else {
1408 // we have missed an update
1409 // interpolate in the middle values
1410
1411 new_value = (NETDATA_DOUBLE)
1412 - ( ( (rd->calculated_value - rd->last_calculated_value)
1412 + ( ( (rd->collector.calculated_value - rd->collector.last_calculated_value)
1413 * (NETDATA_DOUBLE)(next_store_ut - last_collect_ut)
1414 / (NETDATA_DOUBLE)(now_collect_ut - last_collect_ut)
1415 )
1416 - + rd->last_calculated_value
1416 + + rd->collector.last_calculated_value
1417 );
1418
1419 rrdset_debug(st, "%s: CALC2 DEF " NETDATA_DOUBLE_FORMAT " = ((("
@@ -1421,9 +1421,9 @@ static inline size_t rrdset_done_interpolate(
1421 " * %llu"
1422 " / %llu) + " NETDATA_DOUBLE_FORMAT, rrddim_name(rd)
1423 , new_value
1424 - , rd->calculated_value, rd->last_calculated_value
1424 + , rd->collector.calculated_value, rd->collector.last_calculated_value
1425 , (next_store_ut - first_ut)
1426 - , (now_collect_ut - first_ut), rd->last_calculated_value
1426 + , (now_collect_ut - first_ut), rd->collector.last_calculated_value
1427 );
1428 }
1429 break;
@@ -1441,7 +1441,7 @@ static inline size_t rrdset_done_interpolate(
1441 continue;
1442 }
1443
1444 - if(likely(rrddim_check_updated(rd) && rd->collections_counter > 1 && iterations < gap_when_lost_iterations_above)) {
1444 + if(likely(rrddim_check_updated(rd) && rd->collector.counter > 1 && iterations < gap_when_lost_iterations_above)) {
1445 uint32_t dim_storage_flags = storage_flags;
1446
1447 if (ml_dimension_is_anomalous(rd, current_time_s, new_value, true)) {
@@ -1453,7 +1453,7 @@ static inline size_t rrdset_done_interpolate(
1453 rrddim_push_metrics_v2(rsb, rd, next_store_ut, new_value, dim_storage_flags);
1454
1455 rrddim_store_metric(rd, next_store_ut, new_value, dim_storage_flags);
1456 - rd->last_stored_value = new_value;
1456 + rd->collector.last_stored_value = new_value;
1457 }
1458 else {
1459 (void) ml_dimension_is_anomalous(rd, current_time_s, 0, false);
@@ -1464,7 +1464,7 @@ static inline size_t rrdset_done_interpolate(
1464 rrddim_push_metrics_v2(rsb, rd, next_store_ut, NAN, SN_FLAG_NONE);
1465
1466 rrddim_store_metric(rd, next_store_ut, NAN, SN_FLAG_NONE);
1467 - rd->last_stored_value = NAN;
1467 + rd->collector.last_stored_value = NAN;
1468 }
1469
1470 stored_entries++;
@@ -1667,22 +1667,22 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
1667 if(likely(rrddim_check_updated(rd))) {
1668 // if the new is smaller than the old (an overflow, or reset), set the old equal to the new
1669 // to reset the calculation (it will give zero as the calculation for this second)
1670 - if(unlikely(rd->algorithm == RRD_ALGORITHM_PCENT_OVER_DIFF_TOTAL && rd->last_collected_value > rd->collected_value)) {
1670 + if(unlikely(rd->algorithm == RRD_ALGORITHM_PCENT_OVER_DIFF_TOTAL && rd->collector.last_collected_value > rd->collector.collected_value)) {
1671 debug(D_RRD_STATS, "'%s' / '%s': RESET or OVERFLOW. Last collected value = " COLLECTED_NUMBER_FORMAT ", current = " COLLECTED_NUMBER_FORMAT
1672 , rrdset_id(st)
1673 , rrddim_name(rd)
1674 - , rd->last_collected_value
1675 - , rd->collected_value
1674 + , rd->collector.last_collected_value
1675 + , rd->collector.collected_value
1676 );
1677
1678 if(!(rrddim_option_check(rd, RRDDIM_OPTION_DONT_DETECT_RESETS_OR_OVERFLOWS)))
1679 has_reset_value = 1;
1680
1681 - rd->last_collected_value = rd->collected_value;
1681 + rd->collector.last_collected_value = rd->collector.collected_value;
1682 }
1683
1684 - last_collected_total += rd->last_collected_value;
1685 - collected_total += rd->collected_value;
1684 + last_collected_total += rd->collector.last_collected_value;
1685 + collected_total += rd->collector.collected_value;
1686
1687 if(unlikely(rrddim_flag_check(rd, RRDDIM_FLAG_OBSOLETE))) {
1688 error("Dimension %s in chart '%s' has the OBSOLETE flag set, but it is collected.", rrddim_name(rd), rrdset_id(st));
@@ -1706,7 +1706,7 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
1706 if(unlikely(!rd)) continue;
1707
1708 if(unlikely(!rrddim_check_updated(rd))) {
1709 - rd->calculated_value = 0;
1709 + rd->collector.calculated_value = 0;
1710 continue;
1711 }
1712
@@ -1716,25 +1716,25 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
1716 " last_calculated_value = " NETDATA_DOUBLE_FORMAT
1717 " calculated_value = " NETDATA_DOUBLE_FORMAT
1718 , rrddim_name(rd)
1719 - , rd->last_collected_value
1720 - , rd->collected_value
1721 - , rd->last_calculated_value
1722 - , rd->calculated_value
1719 + , rd->collector.last_collected_value
1720 + , rd->collector.collected_value
1721 + , rd->collector.last_calculated_value
1722 + , rd->collector.calculated_value
1723 );
1724
1725 switch(rd->algorithm) {
1726 case RRD_ALGORITHM_ABSOLUTE:
1727 - rd->calculated_value = (NETDATA_DOUBLE)rd->collected_value
1728 - * (NETDATA_DOUBLE)rd->multiplier
1729 - / (NETDATA_DOUBLE)rd->divisor;
1727 + rd->collector.calculated_value = (NETDATA_DOUBLE)rd->collector.collected_value
1728 + * (NETDATA_DOUBLE)rd->multiplier
1729 + / (NETDATA_DOUBLE)rd->divisor;
1730
1731 rrdset_debug(st, "%s: CALC ABS/ABS-NO-IN " NETDATA_DOUBLE_FORMAT " = "
1732 COLLECTED_NUMBER_FORMAT
1733 " * " NETDATA_DOUBLE_FORMAT
1734 " / " NETDATA_DOUBLE_FORMAT
1735 , rrddim_name(rd)
1736 - , rd->calculated_value
1737 - , rd->collected_value
1736 + , rd->collector.calculated_value
1737 + , rd->collector.collected_value
1738 , (NETDATA_DOUBLE)rd->multiplier
1739 , (NETDATA_DOUBLE)rd->divisor
1740 );
@@ -1742,28 +1742,28 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
1742
1743 case RRD_ALGORITHM_PCENT_OVER_ROW_TOTAL:
1744 if(unlikely(!collected_total))
1745 - rd->calculated_value = 0;
1745 + rd->collector.calculated_value = 0;
1746 else
1747 // the percentage of the current value
1748 // over the total of all dimensions
1749 - rd->calculated_value =
1749 + rd->collector.calculated_value =
1750 (NETDATA_DOUBLE)100
1751 - * (NETDATA_DOUBLE)rd->collected_value
1751 + * (NETDATA_DOUBLE)rd->collector.collected_value
1752 / (NETDATA_DOUBLE)collected_total;
1753
1754 rrdset_debug(st, "%s: CALC PCENT-ROW " NETDATA_DOUBLE_FORMAT " = 100"
1755 " * " COLLECTED_NUMBER_FORMAT
1756 " / " COLLECTED_NUMBER_FORMAT
1757 , rrddim_name(rd)
1758 - , rd->calculated_value
1759 - , rd->collected_value
1758 + , rd->collector.calculated_value
1759 + , rd->collector.collected_value
1760 , collected_total
1761 );
1762 break;
1763
1764 case RRD_ALGORITHM_INCREMENTAL:
1765 - if(unlikely(rd->collections_counter <= 1)) {
1766 - rd->calculated_value = 0;
1765 + if(unlikely(rd->collector.counter <= 1)) {
1766 + rd->collector.calculated_value = 0;
1767 continue;
1768 }
1769
@@ -1771,19 +1771,19 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
1771 // to reset the calculation (it will give zero as the calculation for this second).
1772 // It is imperative to set the comparison to uint64_t since type collected_number is signed and
1773 // produces wrong results as far as incremental counters are concerned.
1774 - if(unlikely((uint64_t)rd->last_collected_value > (uint64_t)rd->collected_value)) {
1774 + if(unlikely((uint64_t)rd->collector.last_collected_value > (uint64_t)rd->collector.collected_value)) {
1775 debug(D_RRD_STATS, "'%s' / '%s': RESET or OVERFLOW. Last collected value = " COLLECTED_NUMBER_FORMAT ", current = " COLLECTED_NUMBER_FORMAT
1776 , rrdset_id(st)
1777 , rrddim_name(rd)
1778 - , rd->last_collected_value
1779 - , rd->collected_value);
1778 + , rd->collector.last_collected_value
1779 + , rd->collector.collected_value);
1780
1781 if(!(rrddim_option_check(rd, RRDDIM_OPTION_DONT_DETECT_RESETS_OR_OVERFLOWS)))
1782 has_reset_value = 1;
1783
1784 - uint64_t last = (uint64_t)rd->last_collected_value;
1785 - uint64_t new = (uint64_t)rd->collected_value;
1786 - uint64_t max = (uint64_t)rd->collected_value_max;
1784 + uint64_t last = (uint64_t)rd->collector.last_collected_value;
1785 + uint64_t new = (uint64_t)rd->collector.collected_value;
1786 + uint64_t max = (uint64_t)rd->collector.collected_value_max;
1787 uint64_t cap = 0;
1788
1789 // Signed values are handled by exploiting two's complement which will produce positive deltas
@@ -1800,19 +1800,19 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
1800 // overflow.
1801 // TODO: remember recent history of rates and compare with current rate to reduce this chance.
1802 if (delta < max_acceptable_rate) {
1803 - rd->calculated_value +=
1803 + rd->collector.calculated_value +=
1804 (NETDATA_DOUBLE) delta
1805 * (NETDATA_DOUBLE) rd->multiplier
1806 / (NETDATA_DOUBLE) rd->divisor;
1807 } else {
1808 // This is a reset. Any overflow with a rate greater than MAX_INCREMENTAL_PERCENT_RATE will also
1809 // be detected as a reset instead.
1810 - rd->calculated_value += (NETDATA_DOUBLE)0;
1810 + rd->collector.calculated_value += (NETDATA_DOUBLE)0;
1811 }
1812 }
1813 else {
1814 - rd->calculated_value +=
1815 - (NETDATA_DOUBLE) (rd->collected_value - rd->last_collected_value)
1814 + rd->collector.calculated_value +=
1815 + (NETDATA_DOUBLE) (rd->collector.collected_value - rd->collector.last_collected_value)
1816 * (NETDATA_DOUBLE) rd->multiplier
1817 / (NETDATA_DOUBLE) rd->divisor;
1818 }
@@ -1823,35 +1823,35 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
1823 " * " NETDATA_DOUBLE_FORMAT
1824 " / " NETDATA_DOUBLE_FORMAT
1825 , rrddim_name(rd)
1826 - , rd->calculated_value
1827 - , rd->collected_value, rd->last_collected_value
1826 + , rd->collector.calculated_value
1827 + , rd->collector.collected_value, rd->collector.last_collected_value
1828 , (NETDATA_DOUBLE)rd->multiplier
1829 , (NETDATA_DOUBLE)rd->divisor
1830 );
1831 break;
1832
1833 case RRD_ALGORITHM_PCENT_OVER_DIFF_TOTAL:
1834 - if(unlikely(rd->collections_counter <= 1)) {
1835 - rd->calculated_value = 0;
1834 + if(unlikely(rd->collector.counter <= 1)) {
1835 + rd->collector.calculated_value = 0;
1836 continue;
1837 }
1838
1839 // the percentage of the current increment
1840 // over the increment of all dimensions together
1841 if(unlikely(collected_total == last_collected_total))
1842 - rd->calculated_value = 0;
1842 + rd->collector.calculated_value = 0;
1843 else
1844 - rd->calculated_value =
1844 + rd->collector.calculated_value =
1845 (NETDATA_DOUBLE)100
1846 - * (NETDATA_DOUBLE)(rd->collected_value - rd->last_collected_value)
1846 + * (NETDATA_DOUBLE)(rd->collector.collected_value - rd->collector.last_collected_value)
1847 / (NETDATA_DOUBLE)(collected_total - last_collected_total);
1848
1849 rrdset_debug(st, "%s: CALC PCENT-DIFF " NETDATA_DOUBLE_FORMAT " = 100"
1850 " * (" COLLECTED_NUMBER_FORMAT " - " COLLECTED_NUMBER_FORMAT ")"
1851 " / (" COLLECTED_NUMBER_FORMAT " - " COLLECTED_NUMBER_FORMAT ")"
1852 , rrddim_name(rd)
1853 - , rd->calculated_value
1854 - , rd->collected_value, rd->last_collected_value
1853 + , rd->collector.calculated_value
1854 + , rd->collector.collected_value, rd->collector.last_collected_value
1855 , collected_total, last_collected_total
1856 );
1857 break;
@@ -1859,11 +1859,11 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
1859 default:
1860 // make the default zero, to make sure
1861 // it gets noticed when we add new types
1862 - rd->calculated_value = 0;
1862 + rd->collector.calculated_value = 0;
1863
1864 rrdset_debug(st, "%s: CALC " NETDATA_DOUBLE_FORMAT " = 0"
1865 , rrddim_name(rd)
1866 - , rd->calculated_value
1866 + , rd->collector.calculated_value
1867 );
1868 break;
1869 }
@@ -1874,10 +1874,10 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
1874 " last_calculated_value = " NETDATA_DOUBLE_FORMAT
1875 " calculated_value = " NETDATA_DOUBLE_FORMAT
1876 , rrddim_name(rd)
1877 - , rd->last_collected_value
1878 - , rd->collected_value
1879 - , rd->last_calculated_value
1880 - , rd->calculated_value
1877 + , rd->collector.last_collected_value
1878 + , rd->collector.collected_value
1879 + , rd->collector.last_calculated_value
1880 + , rd->collector.calculated_value
1881 );
1882 }
1883
@@ -1913,9 +1913,9 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
1913 if(unlikely(!rrddim_check_updated(rd)))
1914 continue;
1915
1916 - rrdset_debug(st, "%s: setting last_collected_value (old: " COLLECTED_NUMBER_FORMAT ") to last_collected_value (new: " COLLECTED_NUMBER_FORMAT ")", rrddim_name(rd), rd->last_collected_value, rd->collected_value);
1916 + rrdset_debug(st, "%s: setting last_collected_value (old: " COLLECTED_NUMBER_FORMAT ") to last_collected_value (new: " COLLECTED_NUMBER_FORMAT ")", rrddim_name(rd), rd->collector.last_collected_value, rd->collector.collected_value);
1917
1918 - rd->last_collected_value = rd->collected_value;
1918 + rd->collector.last_collected_value = rd->collector.collected_value;
1919
1920 switch(rd->algorithm) {
1921 case RRD_ALGORITHM_INCREMENTAL:
@@ -1923,10 +1923,10 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
1923 rrdset_debug(st, "%s: setting last_calculated_value (old: " NETDATA_DOUBLE_FORMAT ") to "
1924 "last_calculated_value (new: " NETDATA_DOUBLE_FORMAT ")"
1925 , rrddim_name(rd)
1926 - , rd->last_calculated_value + rd->calculated_value
1927 - , rd->calculated_value);
1926 + , rd->collector.last_calculated_value + rd->collector.calculated_value
1927 + , rd->collector.calculated_value);
1928
1929 - rd->last_calculated_value += rd->calculated_value;
1929 + rd->collector.last_calculated_value += rd->collector.calculated_value;
1930 }
1931 else {
1932 rrdset_debug(st, "THIS IS THE FIRST POINT");
@@ -1939,15 +1939,15 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
1939 rrdset_debug(st, "%s: setting last_calculated_value (old: " NETDATA_DOUBLE_FORMAT ") to "
1940 "last_calculated_value (new: " NETDATA_DOUBLE_FORMAT ")"
1941 , rrddim_name(rd)
1942 - , rd->last_calculated_value
1943 - , rd->calculated_value);
1942 + , rd->collector.last_calculated_value
1943 + , rd->collector.calculated_value);
1944
1945 - rd->last_calculated_value = rd->calculated_value;
1945 + rd->collector.last_calculated_value = rd->collector.calculated_value;
1946 break;
1947 }
1948
1949 - rd->calculated_value = 0;
1950 - rd->collected_value = 0;
1949 + rd->collector.calculated_value = 0;
1950 + rd->collector.collected_value = 0;
1951 rrddim_clear_updated(rd);
1952
1953 rrdset_debug(st, "%s: END "
@@ -1956,10 +1956,10 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
1956 " last_calculated_value = " NETDATA_DOUBLE_FORMAT
1957 " calculated_value = " NETDATA_DOUBLE_FORMAT
1958 , rrddim_name(rd)
1959 - , rd->last_collected_value
1960 - , rd->collected_value
1961 - , rd->last_calculated_value
1962 - , rd->calculated_value
1959 + , rd->collector.last_collected_value
1960 + , rd->collector.collected_value
1961 + , rd->collector.last_calculated_value
1962 + , rd->collector.calculated_value
1963 );
1964 }
1965
database/rrdvar.c
+3 -3
@@ -272,9 +272,9 @@ void rrdvar_store_for_chart(RRDHOST *host, RRDSET *st) {
272
273 RRDDIM *rd;
274 rrddim_foreach_read(rd, st) {
275 - rrddimvar_add_and_leave_released(rd, RRDVAR_TYPE_CALCULATED, NULL, NULL, &rd->last_stored_value, RRDVAR_FLAG_NONE);
276 - rrddimvar_add_and_leave_released(rd, RRDVAR_TYPE_COLLECTED, NULL, "_raw", &rd->last_collected_value, RRDVAR_FLAG_NONE);
277 - rrddimvar_add_and_leave_released(rd, RRDVAR_TYPE_TIME_T, NULL, "_last_collected_t", &rd->last_collected_time.tv_sec, RRDVAR_FLAG_NONE);
275 + rrddimvar_add_and_leave_released(rd, RRDVAR_TYPE_CALCULATED, NULL, NULL, &rd->collector.last_stored_value, RRDVAR_FLAG_NONE);
276 + rrddimvar_add_and_leave_released(rd, RRDVAR_TYPE_COLLECTED, NULL, "_raw", &rd->collector.last_collected_value, RRDVAR_FLAG_NONE);
277 + rrddimvar_add_and_leave_released(rd, RRDVAR_TYPE_TIME_T, NULL, "_last_collected_t", &rd->collector.last_collected_time.tv_sec, RRDVAR_FLAG_NONE);
278 }
279 rrddim_foreach_done(rd);
280 }
exporting/graphite/graphite.c
+2 -2
@@ -141,8 +141,8 @@ int format_dimension_collected_graphite_plaintext(struct instance *instance, RRD
141 (host->tags) ? ";" : "",
142 (host->tags) ? rrdhost_tags(host) : "",
143 (instance->labels_buffer) ? buffer_tostring(instance->labels_buffer) : "",
144 - rd->last_collected_value,
145 - (unsigned long long)rd->last_collected_time.tv_sec);
144 + rd->collector.last_collected_value,
145 + (unsigned long long)rd->collector.last_collected_time.tv_sec);
146
147 return 0;
148 }
exporting/json/json.c
+2 -2
@@ -200,9 +200,9 @@ int format_dimension_collected_json_plaintext(struct instance *instance, RRDDIM
200 rrdset_units(st),
201 rrddim_id(rd),
202 rrddim_name(rd),
203 - rd->last_collected_value,
203 + rd->collector.last_collected_value,
204
205 - (unsigned long long)rd->last_collected_time.tv_sec);
205 + (unsigned long long)rd->collector.last_collected_time.tv_sec);
206
207 if (instance->config.type != EXPORTING_CONNECTOR_TYPE_JSON_HTTP) {
208 buffer_strcat(instance->buffer, "\n");
exporting/opentsdb/opentsdb.c
+4 -4
@@ -190,8 +190,8 @@ int format_dimension_collected_opentsdb_telnet(struct instance *instance, RRDDIM
190 instance->config.prefix,
191 chart_name,
192 dimension_name,
193 - (unsigned long long)rd->last_collected_time.tv_sec,
194 - rd->last_collected_value,
193 + (unsigned long long)rd->collector.last_collected_time.tv_sec,
194 + rd->collector.last_collected_value,
195 (host == localhost) ? instance->config.hostname : rrdhost_hostname(host),
196 (host->tags) ? " " : "",
197 (host->tags) ? rrdhost_tags(host) : "",
@@ -332,8 +332,8 @@ int format_dimension_collected_opentsdb_http(struct instance *instance, RRDDIM *
332 instance->config.prefix,
333 chart_name,
334 dimension_name,
335 - (unsigned long long)rd->last_collected_time.tv_sec,
336 - rd->last_collected_value,
335 + (unsigned long long)rd->collector.last_collected_time.tv_sec,
336 + rd->collector.last_collected_value,
337 (host == localhost) ? instance->config.hostname : rrdhost_hostname(host),
338 (host->tags) ? " " : "",
339 (host->tags) ? rrdhost_tags(host) : "",
exporting/prometheus/prometheus.c
+6 -6
@@ -540,13 +540,13 @@ static void generate_as_collected_prom_metric(BUFFER *wb,
540 buffer_sprintf(
541 wb,
542 NETDATA_DOUBLE_FORMAT,
543 - (NETDATA_DOUBLE)p->rd->last_collected_value * (NETDATA_DOUBLE)p->rd->multiplier /
544 - (NETDATA_DOUBLE)p->rd->divisor);
543 + (NETDATA_DOUBLE)p->rd->collector.last_collected_value * (NETDATA_DOUBLE)p->rd->multiplier /
544 + (NETDATA_DOUBLE)p->rd->divisor);
545 else
546 - buffer_sprintf(wb, COLLECTED_NUMBER_FORMAT, p->rd->last_collected_value);
546 + buffer_sprintf(wb, COLLECTED_NUMBER_FORMAT, p->rd->collector.last_collected_value);
547
548 if (p->output_options & PROMETHEUS_OUTPUT_TIMESTAMPS)
549 - buffer_sprintf(wb, " %llu\n", timeval_msec(&p->rd->last_collected_time));
549 + buffer_sprintf(wb, " %llu\n", timeval_msec(&p->rd->collector.last_collected_time));
550 else
551 buffer_sprintf(wb, "\n");
552 }
@@ -675,7 +675,7 @@ static void rrd_stats_api_v1_charts_allmetrics_prometheus(
675 // for each dimension
676 RRDDIM *rd;
677 rrddim_foreach_read(rd, st) {
678 - if (rd->collections_counter && !rrddim_flag_check(rd, RRDDIM_FLAG_OBSOLETE)) {
678 + if (rd->collector.counter && !rrddim_flag_check(rd, RRDDIM_FLAG_OBSOLETE)) {
679 char dimension[PROMETHEUS_ELEMENT_MAX + 1];
680 char *suffix = "";
681
@@ -694,7 +694,7 @@ static void rrd_stats_api_v1_charts_allmetrics_prometheus(
694 p.st = st;
695 p.rd = rd;
696
697 - if (unlikely(rd->last_collected_time.tv_sec < instance->after))
697 + if (unlikely(rd->collector.last_collected_time.tv_sec < instance->after))
698 continue;
699
700 p.type = "gauge";
exporting/prometheus/remote_write/remote_write.c
+9 -9
@@ -234,7 +234,7 @@ int format_dimension_prometheus_remote_write(struct instance *instance, RRDDIM *
234 struct prometheus_remote_write_specific_data *connector_specific_data =
235 (struct prometheus_remote_write_specific_data *)simple_connector_data->connector_specific_data;
236
237 - if (rd->collections_counter && !rrddim_flag_check(rd, RRDDIM_FLAG_OBSOLETE)) {
237 + if (rd->collector.counter && !rrddim_flag_check(rd, RRDDIM_FLAG_OBSOLETE)) {
238 char name[PROMETHEUS_LABELS_MAX + 1];
239 char dimension[PROMETHEUS_ELEMENT_MAX + 1];
240 char *suffix = "";
@@ -243,14 +243,14 @@ int format_dimension_prometheus_remote_write(struct instance *instance, RRDDIM *
243 if (as_collected) {
244 // we need as-collected / raw data
245
246 - if (unlikely(rd->last_collected_time.tv_sec < instance->after)) {
246 + if (unlikely(rd->collector.last_collected_time.tv_sec < instance->after)) {
247 debug(
248 D_EXPORTING,
249 "EXPORTING: not sending dimension '%s' of chart '%s' from host '%s', "
250 "its last data collection (%lu) is not within our timeframe (%lu to %lu)",
251 rrddim_id(rd), rrdset_id(rd->rrdset),
252 (host == localhost) ? instance->config.hostname : rrdhost_hostname(host),
253 - (unsigned long)rd->last_collected_time.tv_sec,
253 + (unsigned long)rd->collector.last_collected_time.tv_sec,
254 (unsigned long)instance->after,
255 (unsigned long)instance->before);
256 return 0;
@@ -272,10 +272,10 @@ int format_dimension_prometheus_remote_write(struct instance *instance, RRDDIM *
272 snprintf(name, PROMETHEUS_LABELS_MAX, "%s_%s%s", instance->config.prefix, context, suffix);
273
274 add_metric(
275 - connector_specific_data->write_request,
276 - name, chart, family, dimension,
275 + connector_specific_data->write_request,
276 + name, chart, family, dimension,
277 (host == localhost) ? instance->config.hostname : rrdhost_hostname(host),
278 - rd->last_collected_value, timeval_msec(&rd->last_collected_time));
278 + rd->collector.last_collected_value, timeval_msec(&rd->collector.last_collected_time));
279 } else {
280 // the dimensions of the chart, do not have the same algorithm, multiplier or divisor
281 // we create a metric per dimension
@@ -289,10 +289,10 @@ int format_dimension_prometheus_remote_write(struct instance *instance, RRDDIM *
289 suffix);
290
291 add_metric(
292 - connector_specific_data->write_request,
293 - name, chart, family, NULL,
292 + connector_specific_data->write_request,
293 + name, chart, family, NULL,
294 (host == localhost) ? instance->config.hostname : rrdhost_hostname(host),
295 - rd->last_collected_value, timeval_msec(&rd->last_collected_time));
295 + rd->collector.last_collected_value, timeval_msec(&rd->collector.last_collected_time));
296 }
297 } else {
298 // we need average or sum of the data
streaming/receiver.c
+39 -41
@@ -94,19 +94,19 @@ static inline int read_stream(struct receiver_state *r, char* buffer, size_t siz
94
95 static inline bool receiver_read_uncompressed(struct receiver_state *r) {
96 #ifdef NETDATA_INTERNAL_CHECKS
97 - if(r->read_buffer[r->read_len] != '\0')
97 + if(r->reader.read_buffer[r->reader.read_len] != '\0')
98 fatal("%s(): read_buffer does not start with zero", __FUNCTION__ );
99 #endif
100
101 - int bytes_read = read_stream(r, r->read_buffer + r->read_len, sizeof(r->read_buffer) - r->read_len - 1);
101 + int bytes_read = read_stream(r, r->reader.read_buffer + r->reader.read_len, sizeof(r->reader.read_buffer) - r->reader.read_len - 1);
102 if(unlikely(bytes_read <= 0))
103 return false;
104
105 worker_set_metric(WORKER_RECEIVER_JOB_BYTES_READ, (NETDATA_DOUBLE)bytes_read);
106 worker_set_metric(WORKER_RECEIVER_JOB_BYTES_UNCOMPRESSED, (NETDATA_DOUBLE)bytes_read);
107
108 - r->read_len += bytes_read;
109 - r->read_buffer[r->read_len] = '\0';
108 + r->reader.read_len += bytes_read;
109 + r->reader.read_buffer[r->reader.read_len] = '\0';
110
111 return true;
112 }
@@ -114,24 +114,24 @@ static inline bool receiver_read_uncompressed(struct receiver_state *r) {
114 #ifdef ENABLE_COMPRESSION
115 static inline bool receiver_read_compressed(struct receiver_state *r) {
116
117 - internal_fatal(r->read_buffer[r->read_len] != '\0',
117 + internal_fatal(r->reader.read_buffer[r->reader.read_len] != '\0',
118 "%s: read_buffer does not start with zero #2", __FUNCTION__ );
119
120 // first use any available uncompressed data
121 if (likely(rrdpush_decompressed_bytes_in_buffer(&r->decompressor))) {
122 - size_t available = sizeof(r->read_buffer) - r->read_len - 1;
122 + size_t available = sizeof(r->reader.read_buffer) - r->reader.read_len - 1;
123 if (likely(available)) {
124 - size_t len = rrdpush_decompressor_get(&r->decompressor, r->read_buffer + r->read_len, available);
124 + size_t len = rrdpush_decompressor_get(&r->decompressor, r->reader.read_buffer + r->reader.read_len, available);
125 if (unlikely(!len)) {
126 internal_error(true, "decompressor returned zero length #1");
127 return false;
128 }
129
130 - r->read_len += (int)len;
131 - r->read_buffer[r->read_len] = '\0';
130 + r->reader.read_len += (int)len;
131 + r->reader.read_buffer[r->reader.read_len] = '\0';
132 }
133 else
134 - internal_fatal(true, "The line to read is too big! Already have %d bytes in read_buffer.", r->read_len);
134 + internal_fatal(true, "The line to read is too big! Already have %zd bytes in read_buffer.", r->reader.read_len);
135
136 return true;
137 }
@@ -139,9 +139,9 @@ static inline bool receiver_read_compressed(struct receiver_state *r) {
139 // no decompressed data available
140 // read the compression signature of the next block
141
142 - if(unlikely(r->read_len + r->decompressor.signature_size > sizeof(r->read_buffer) - 1)) {
142 + if(unlikely(r->reader.read_len + r->decompressor.signature_size > sizeof(r->reader.read_buffer) - 1)) {
143 internal_error(true, "The last incomplete line does not leave enough room for the next compression header! "
144 - "Already have %d bytes in read_buffer.", r->read_len);
144 + "Already have %zd bytes in read_buffer.", r->reader.read_len);
145 return false;
146 }
147
@@ -149,7 +149,7 @@ static inline bool receiver_read_compressed(struct receiver_state *r) {
149 // we have to do a loop here, because read_stream() may return less than the data we need
150 int bytes_read = 0;
151 do {
152 - int ret = read_stream(r, r->read_buffer + r->read_len + bytes_read, r->decompressor.signature_size - bytes_read);
152 + int ret = read_stream(r, r->reader.read_buffer + r->reader.read_len + bytes_read, r->decompressor.signature_size - bytes_read);
153 if (unlikely(ret <= 0))
154 return false;
155
@@ -161,11 +161,11 @@ static inline bool receiver_read_compressed(struct receiver_state *r) {
161 if(unlikely(bytes_read != (int)r->decompressor.signature_size))
162 fatal("read %d bytes, but expected compression signature of size %zu", bytes_read, r->decompressor.signature_size);
163
164 - size_t compressed_message_size = rrdpush_decompressor_start(&r->decompressor, r->read_buffer + r->read_len, bytes_read);
164 + size_t compressed_message_size = rrdpush_decompressor_start(&r->decompressor, r->reader.read_buffer + r->reader.read_len, bytes_read);
165 if (unlikely(!compressed_message_size)) {
166 internal_error(true, "multiplexed uncompressed data in compressed stream!");
167 - r->read_len += bytes_read;
168 - r->read_buffer[r->read_len] = '\0';
167 + r->reader.read_len += bytes_read;
168 + r->reader.read_buffer[r->reader.read_len] = '\0';
169 return true;
170 }
171
@@ -176,7 +176,7 @@ static inline bool receiver_read_compressed(struct receiver_state *r) {
176 }
177
178 // delete compression header from our read buffer
179 - r->read_buffer[r->read_len] = '\0';
179 + r->reader.read_buffer[r->reader.read_len] = '\0';
180
181 // Read the entire compressed block of compressed data
182 char compressed[compressed_message_size];
@@ -207,13 +207,13 @@ static inline bool receiver_read_compressed(struct receiver_state *r) {
207 worker_set_metric(WORKER_RECEIVER_JOB_BYTES_UNCOMPRESSED, (NETDATA_DOUBLE)bytes_to_parse);
208
209 // fill read buffer with decompressed data
210 - size_t len = (int) rrdpush_decompressor_get(&r->decompressor, r->read_buffer + r->read_len, sizeof(r->read_buffer) - r->read_len - 1);
210 + size_t len = (int) rrdpush_decompressor_get(&r->decompressor, r->reader.read_buffer + r->reader.read_len, sizeof(r->reader.read_buffer) - r->reader.read_len - 1);
211 if (unlikely(!len)) {
212 internal_error(true, "decompressor returned zero length #2");
213 return false;
214 }
215 - r->read_len += (int)len;
216 - r->read_buffer[r->read_len] = '\0';
215 + r->reader.read_len += (int)len;
216 + r->reader.read_buffer[r->reader.read_len] = '\0';
217
218 return true;
219 }
@@ -226,19 +226,19 @@ static inline bool receiver_read_compressed(struct receiver_state *r) {
226 /* Produce a full line if one exists, statefully return where we start next time.
227 * When we hit the end of the buffer with a partial line move it to the beginning for the next fill.
228 */
229 -static inline char *receiver_next_line(struct receiver_state *r, char *buffer, size_t buffer_length, size_t *pos) {
230 - size_t start = *pos;
229 +inline char *buffered_reader_next_line(struct buffered_reader *reader, char *dst, size_t dst_size) {
230 + size_t start = reader->pos;
231
232 - char *ss = &r->read_buffer[start];
233 - char *se = &r->read_buffer[r->read_len];
234 - char *ds = buffer;
235 - char *de = &buffer[buffer_length - 2];
232 + char *ss = &reader->read_buffer[start];
233 + char *se = &reader->read_buffer[reader->read_len];
234 + char *ds = dst;
235 + char *de = &dst[dst_size - 2];
236
237 if(ss >= se) {
238 *ds = '\0';
239 - *pos = 0;
240 - r->read_len = 0;
241 - r->read_buffer[r->read_len] = '\0';
239 + reader->pos = 0;
240 + reader->read_len = 0;
241 + reader->read_buffer[reader->read_len] = '\0';
242 return NULL;
243 }
244
@@ -253,25 +253,25 @@ static inline char *receiver_next_line(struct receiver_state *r, char *buffer, s
253 *ds++ = *ss++; // copy the newline too
254 *ds = '\0';
255
256 - *pos = ss - r->read_buffer;
257 - return buffer;
256 + reader->pos = ss - reader->read_buffer;
257 + return dst;
258 }
259
260 // if the destination is full, oops!
261 if(ds == de) {
262 error("STREAM: received line exceeds %d bytes. Truncating it.", PLUGINSD_LINE_MAX);
263 *ds = '\0';
264 - *pos = ss - r->read_buffer;
265 - return buffer;
264 + reader->pos = ss - reader->read_buffer;
265 + return dst;
266 }
267
268 // no newline found in the r->read_buffer
269 // move everything to the beginning
270 - memmove(r->read_buffer, &r->read_buffer[start], r->read_len - start);
271 - r->read_len -= (int)start;
272 - r->read_buffer[r->read_len] = '\0';
270 + memmove(reader->read_buffer, &reader->read_buffer[start], reader->read_len - start);
271 + reader->read_len -= (int)start;
272 + reader->read_buffer[reader->read_len] = '\0';
273 *ds = '\0';
274 - *pos = 0;
274 + reader->pos = 0;
275 return NULL;
276 }
277
@@ -340,14 +340,12 @@ static size_t streaming_parser(struct receiver_state *rpt, struct plugind *cd, i
340 rrdpush_decompressor_destroy(&rpt->decompressor);
341 #endif
342
343 - rpt->read_buffer[0] = '\0';
344 - rpt->read_len = 0;
343 + buffered_reader_init(&rpt->reader);
344
346 - size_t read_buffer_start = 0;
345 char buffer[PLUGINSD_LINE_MAX + 2] = "";
346 while(!receiver_should_stop(rpt)) {
347
350 - if(!receiver_next_line(rpt, buffer, PLUGINSD_LINE_MAX + 2, &read_buffer_start)) {
348 + if(!buffered_reader_next_line(&rpt->reader, buffer, PLUGINSD_LINE_MAX + 2)) {
349 bool have_new_data = compressed_connection ? receiver_read_compressed(rpt) : receiver_read_uncompressed(rpt);
350
351 if(unlikely(!have_new_data)) {
streaming/replication.c
+5 -5
@@ -222,14 +222,14 @@ static void replication_send_chart_collection_state(BUFFER *wb, RRDSET *st, STRE
222 sizeof(PLUGINSD_KEYWORD_REPLAY_RRDDIM_STATE) - 1 + 2);
223 buffer_fast_strcat(wb, rrddim_id(rd), string_strlen(rd->id));
224 buffer_fast_strcat(wb, "' ", 2);
225 - buffer_print_uint64_encoded(wb, encoding, (usec_t) rd->last_collected_time.tv_sec * USEC_PER_SEC +
226 - (usec_t) rd->last_collected_time.tv_usec);
225 + buffer_print_uint64_encoded(wb, encoding, (usec_t) rd->collector.last_collected_time.tv_sec * USEC_PER_SEC +
226 + (usec_t) rd->collector.last_collected_time.tv_usec);
227 buffer_fast_strcat(wb, " ", 1);
228 - buffer_print_int64_encoded(wb, encoding, rd->last_collected_value);
228 + buffer_print_int64_encoded(wb, encoding, rd->collector.last_collected_value);
229 buffer_fast_strcat(wb, " ", 1);
230 - buffer_print_netdata_double_encoded(wb, encoding, rd->last_calculated_value);
230 + buffer_print_netdata_double_encoded(wb, encoding, rd->collector.last_calculated_value);
231 buffer_fast_strcat(wb, " ", 1);
232 - buffer_print_netdata_double_encoded(wb, encoding, rd->last_stored_value);
232 + buffer_print_netdata_double_encoded(wb, encoding, rd->collector.last_stored_value);
233 buffer_fast_strcat(wb, "\n", 1);
234 }
235 rrddim_foreach_done(rd);
streaming/rrdpush.c
+3 -3
@@ -362,7 +362,7 @@ static void rrdpush_send_chart_metrics(BUFFER *wb, RRDSET *st, struct sender_sta
362 buffer_fast_strcat(wb, "SET \"", 5);
363 buffer_fast_strcat(wb, rrddim_id(rd), string_strlen(rd->id));
364 buffer_fast_strcat(wb, "\" = ", 4);
365 - buffer_print_int64(wb, rd->collected_value);
365 + buffer_print_int64(wb, rd->collector.collected_value);
366 buffer_fast_strcat(wb, "\n", 1);
367 }
368 else {
@@ -436,10 +436,10 @@ void rrddim_push_metrics_v2(RRDSET_STREAM_BUFFER *rsb, RRDDIM *rd, usec_t point_
436 buffer_fast_strcat(wb, PLUGINSD_KEYWORD_SET_V2 " '", sizeof(PLUGINSD_KEYWORD_SET_V2) - 1 + 2);
437 buffer_fast_strcat(wb, rrddim_id(rd), string_strlen(rd->id));
438 buffer_fast_strcat(wb, "' ", 2);
439 - buffer_print_int64_encoded(wb, integer_encoding, rd->last_collected_value);
439 + buffer_print_int64_encoded(wb, integer_encoding, rd->collector.last_collected_value);
440 buffer_fast_strcat(wb, " ", 1);
441
442 - if((NETDATA_DOUBLE)rd->last_collected_value == n)
442 + if((NETDATA_DOUBLE)rd->collector.last_collected_value == n)
443 buffer_fast_strcat(wb, "#", 1);
444 else
445 buffer_print_netdata_double_encoded(wb, doubles_encoding, n);
streaming/rrdpush.h
+15 -2
@@ -350,6 +350,19 @@ typedef struct stream_node_instance {
350 } STREAM_NODE_INSTANCE;
351 */
352
353 +struct buffered_reader {
354 + ssize_t read_len;
355 + ssize_t pos;
356 + char read_buffer[PLUGINSD_LINE_MAX + 1];
357 +};
358 +
359 +char *buffered_reader_next_line(struct buffered_reader *reader, char *dst, size_t dst_size);
360 +static inline void buffered_reader_init(struct buffered_reader *reader) {
361 + reader->read_buffer[0] = '\0';
362 + reader->read_len = 0;
363 + reader->pos = 0;
364 +}
365 +
366 struct receiver_state {
367 RRDHOST *host;
368 pid_t tid;
@@ -371,8 +384,8 @@ struct receiver_state {
384 struct rrdhost_system_info *system_info;
385 STREAM_CAPABILITIES capabilities;
386 time_t last_msg_t;
374 - char read_buffer[PLUGINSD_LINE_MAX + 1];
375 - int read_len;
387 +
388 + struct buffered_reader reader;
389
390 uint16_t hops;
391
web/api/exporters/shell/allmetrics_shell.c
+5 -5
@@ -41,11 +41,11 @@ void rrd_stats_api_v1_charts_allmetrics_shell(RRDHOST *host, const char *filter_
41 // for each dimension
42 RRDDIM *rd;
43 rrddim_foreach_read(rd, st) {
44 - if(rd->collections_counter && !rrddim_flag_check(rd, RRDDIM_FLAG_OBSOLETE)) {
44 + if(rd->collector.counter && !rrddim_flag_check(rd, RRDDIM_FLAG_OBSOLETE)) {
45 char dimension[SHELL_ELEMENT_MAX + 1];
46 shell_name_copy(dimension, rd->name?rrddim_name(rd):rrddim_id(rd), SHELL_ELEMENT_MAX);
47
48 - NETDATA_DOUBLE n = rd->last_stored_value;
48 + NETDATA_DOUBLE n = rd->collector.last_stored_value;
49
50 if(isnan(n) || isinf(n))
51 buffer_sprintf(wb, "NETDATA_%s_%s=\"\" # %s\n", chart, dimension, rrdset_units(st));
@@ -135,7 +135,7 @@ void rrd_stats_api_v1_charts_allmetrics_json(RRDHOST *host, const char *filter_s
135 // for each dimension
136 RRDDIM *rd;
137 rrddim_foreach_read(rd, st) {
138 - if(rd->collections_counter && !rrddim_flag_check(rd, RRDDIM_FLAG_OBSOLETE)) {
138 + if(rd->collector.counter && !rrddim_flag_check(rd, RRDDIM_FLAG_OBSOLETE)) {
139 buffer_sprintf(
140 wb,
141 "%s\n"
@@ -146,10 +146,10 @@ void rrd_stats_api_v1_charts_allmetrics_json(RRDHOST *host, const char *filter_s
146 rrddim_id(rd),
147 rrddim_name(rd));
148
149 - if(isnan(rd->last_stored_value))
149 + if(isnan(rd->collector.last_stored_value))
150 buffer_strcat(wb, "null");
151 else
152 - buffer_sprintf(wb, NETDATA_DOUBLE_FORMAT, rd->last_stored_value);
152 + buffer_sprintf(wb, NETDATA_DOUBLE_FORMAT, rd->collector.last_stored_value);
153
154 buffer_strcat(wb, "\n\t\t\t}");
155
web/server/web_client.c
+5 -4
@@ -26,7 +26,7 @@ static inline int bad_request_multiple_dashboard_versions(struct web_client *w)
26 return HTTP_RESP_BAD_REQUEST;
27 }
28
29 -static inline int web_client_crock_socket(struct web_client *w __maybe_unused) {
29 +static inline int web_client_cork_socket(struct web_client *w __maybe_unused) {
30 #ifdef TCP_CORK
31 if(likely(web_client_is_corkable(w) && !w->tcp_cork && w->ofd != -1)) {
32 w->tcp_cork = true;
@@ -53,9 +53,10 @@ static inline void web_client_enable_wait_from_ssl(struct web_client *w) {
53 }
54 }
55
56 -static inline int web_client_uncrock_socket(struct web_client *w __maybe_unused) {
56 +static inline int web_client_uncork_socket(struct web_client *w __maybe_unused) {
57 #ifdef TCP_CORK
58 if(likely(w->tcp_cork && w->ofd != -1)) {
59 + w->tcp_cork = false;
60 if(unlikely(setsockopt(w->ofd, IPPROTO_TCP, TCP_CORK, (char *) &w->tcp_cork, sizeof(int)) != 0)) {
61 error("%llu: failed to disable TCP_CORK on socket.", w->id);
62 w->tcp_cork = true;
@@ -153,7 +154,7 @@ static void web_client_reset_allocations(struct web_client *w, bool free_all) {
154 }
155
156 void web_client_request_done(struct web_client *w) {
156 - web_client_uncrock_socket(w);
157 + web_client_uncork_socket(w);
158
159 debug(D_WEB_CLIENT, "%llu: Resetting client.", w->id);
160
@@ -1265,7 +1266,7 @@ static inline void web_client_send_http_header(struct web_client *w) {
1266 , buffer_tostring(w->response.header_output)
1267 );
1268
1268 - web_client_crock_socket(w);
1269 + web_client_cork_socket(w);
1270
1271 size_t count = 0;
1272 ssize_t bytes;