@cryptotaxi247 / netdata-1 / commits / a1d450d2b

Plugins.d doubles (#21349)

* Plugins.d protocol support DIMENSION type=float and SET accepts floats * docs about type=float * missed pluginsd parsing * cubix fixes * cubix fixes * fixed unittest * fixed copilot comments * test: make go.d/ping using floats * Check for STREAM_CAP_FLOAT_BASELINE when dim is float. Fallback to int64 as needed * Revert "test: make go.d/ping using floats" This reverts commit 44faac80be0982e1dcb7509717d3638e32db1eda. * - Use STREAM_CAP_FLOAT_BASELINE to determine sender's float capability. - Ensure proper fallback to int64 for older senders. - Optimize dimension collection/reset logic based on data type. * Add support for handling float and int dimensions in exporting and pluginsd - Use `STREAM_CAP_FLOAT_BASELINE` to distinguish float dimensions from int. - Update exporting code (Graphite, OpenTSDB, JSON) to handle both types. - Reset collected values appropriately when switching data types. - Refactor parsing logic to align with sender capabilities on float/int. - Ensure proper fallback to int64 for older senders without float support. --------- Co-authored-by: ilyam8 <ilya@netdata.cloud> Co-authored-by: Stelios Fragkakis <52996999+stelfrag@users.noreply.github.com>

Costa Tsaousis committed Feb 26, 2026 at 09:10 UTC a1d450d2b1444d018c835fa74ee419836a26a016
22 files changed +386 -159
src/collectors/proc.plugin/proc_diskstats.c
+7 -7
@@ -1560,8 +1560,8 @@ int do_proc_diskstats(int update_every, usec_t dt) {
1560 if (d->do_io == CONFIG_BOOLEAN_YES || d->do_io == CONFIG_BOOLEAN_AUTO) {
1561 d->do_io = CONFIG_BOOLEAN_YES;
1562
1563 - last_readsectors = d->disk_io.rd_io_reads ? d->disk_io.rd_io_reads->collector.last_collected_value / SECTOR_SIZE : 0;
1564 - last_writesectors = d->disk_io.rd_io_writes ? d->disk_io.rd_io_writes->collector.last_collected_value / SECTOR_SIZE : 0;
1563 + last_readsectors = d->disk_io.rd_io_reads ? rrddim_last_collected_raw_int(d->disk_io.rd_io_reads) / SECTOR_SIZE : 0;
1564 + last_writesectors = d->disk_io.rd_io_writes ? rrddim_last_collected_raw_int(d->disk_io.rd_io_writes) / SECTOR_SIZE : 0;
1565
1566 common_disk_io(&d->disk_io,
1567 d->chart_id,
@@ -1602,8 +1602,8 @@ int do_proc_diskstats(int update_every, usec_t dt) {
1602 if (d->do_ops == CONFIG_BOOLEAN_YES || d->do_ops == CONFIG_BOOLEAN_AUTO) {
1603 d->do_ops = CONFIG_BOOLEAN_YES;
1604
1605 - last_rd_ios = d->disk_ops.rd_ops_reads ? d->disk_ops.rd_ops_reads->collector.last_collected_value : 0;
1606 - last_wr_ios = d->disk_ops.rd_ops_writes ? d->disk_ops.rd_ops_writes->collector.last_collected_value : 0;
1605 + last_rd_ios = d->disk_ops.rd_ops_reads ? rrddim_last_collected_raw_int(d->disk_ops.rd_ops_reads) : 0;
1606 + last_wr_ios = d->disk_ops.rd_ops_writes ? rrddim_last_collected_raw_int(d->disk_ops.rd_ops_writes) : 0;
1607
1608 common_disk_ops(&d->disk_ops,
1609 d->chart_id,
@@ -1687,7 +1687,7 @@ int do_proc_diskstats(int update_every, usec_t dt) {
1687 if (d->do_util == CONFIG_BOOLEAN_YES || d->do_util == CONFIG_BOOLEAN_AUTO) {
1688 d->do_util = CONFIG_BOOLEAN_YES;
1689
1690 - last_busy_ms = d->disk_busy.rd_busy ? d->disk_busy.rd_busy->collector.last_collected_value : 0;
1690 + last_busy_ms = d->disk_busy.rd_busy ? rrddim_last_collected_raw_int(d->disk_busy.rd_busy) : 0;
1691
1692 common_disk_busy(&d->disk_busy,
1693 d->chart_id,
@@ -1772,8 +1772,8 @@ int do_proc_diskstats(int update_every, usec_t dt) {
1772 if (d->do_iotime == CONFIG_BOOLEAN_YES || d->do_iotime == CONFIG_BOOLEAN_AUTO) {
1773 d->do_iotime = CONFIG_BOOLEAN_YES;
1774
1775 - last_readms = d->disk_iotime.rd_reads_ms ? d->disk_iotime.rd_reads_ms->collector.last_collected_value : 0;
1776 - last_writems = d->disk_iotime.rd_writes_ms ? d->disk_iotime.rd_writes_ms->collector.last_collected_value : 0;
1775 + last_readms = d->disk_iotime.rd_reads_ms ? rrddim_last_collected_raw_int(d->disk_iotime.rd_reads_ms) : 0;
1776 + last_writems = d->disk_iotime.rd_writes_ms ? rrddim_last_collected_raw_int(d->disk_iotime.rd_writes_ms) : 0;
1777
1778 common_disk_iotime(
1779 &d->disk_iotime,
src/database/engine/dbengine-stresstest.c
+2 -2
@@ -37,13 +37,13 @@ static RRDHOST *dbengine_rrdhost_find_or_create(char *name) {
37 static inline void rrddim_set_by_pointer_fake_time(RRDDIM *rd, collected_number value, time_t now) {
38 rd->collector.last_collected_time.tv_sec = now;
39 rd->collector.last_collected_time.tv_usec = 0;
40 - rd->collector.collected_value = value;
40 + rrddim_set_collected_int(rd, value);
41 rrddim_set_updated(rd);
42
43 rd->collector.counter++;
44
45 collected_number v = (value >= 0) ? value : -value;
46 - if(unlikely(v > rd->collector.collected_value_max)) rd->collector.collected_value_max = v;
46 + if(unlikely(v > rd->collector.collected.i.collected_value_max)) rrddim_set_collected_max_int(rd, v);
47 }
48
49 struct dbengine_chart_thread {
src/database/engine/dbengine-unittest.c
+2 -2
@@ -81,13 +81,13 @@ static inline void storage_point_check(size_t region, size_t chart, size_t dim,
81 static inline void rrddim_set_by_pointer_fake_time(RRDDIM *rd, collected_number value, time_t now) {
82 rd->collector.last_collected_time.tv_sec = now;
83 rd->collector.last_collected_time.tv_usec = 0;
84 - rd->collector.collected_value = value;
84 + rrddim_set_collected_int(rd, value);
85 rrddim_set_updated(rd);
86
87 rd->collector.counter++;
88
89 collected_number v = (value >= 0) ? value : -value;
90 - if(unlikely(v > rd->collector.collected_value_max)) rd->collector.collected_value_max = v;
90 + if(unlikely(v > rd->collector.collected.i.collected_value_max)) rrddim_set_collected_max_int(rd, v);
91 }
92
93 static RRDHOST *dbengine_rrdhost_find_or_create(char *name) {
src/database/rrddim.c
+38 -6
@@ -611,11 +611,37 @@ inline collected_number rrddim_set_by_pointer(RRDSET *st, RRDDIM *rd, collected_
611 return rrddim_timed_set_by_pointer(st, rd, now, value);
612 }
613
614 +// Sets the collected value for this dimension.
615 +// Returns the *previous* cycle's collected value (last_collected_value), not the current one.
616 +NETDATA_DOUBLE rrddim_set_by_pointer_double(RRDSET *st, RRDDIM *rd, NETDATA_DOUBLE value) {
617 + struct timeval now;
618 + now_realtime_timeval(&now);
619 +
620 + // For float dims, store exact; for int dims, keep int semantics
621 + if(rrddim_is_float(rd)) {
622 + rd->collector.last_collected_time = now;
623 + rrddim_set_collected_float(rd, value);
624 + rrddim_set_updated(rd);
625 + rd->collector.counter++;
626 +
627 + NETDATA_DOUBLE v = value >= 0 ? value : -value;
628 + if(unlikely(v > rrddim_collected_max_as_double(rd)))
629 + rrddim_set_collected_max_float(rd, v);
630 +
631 + return rrddim_last_collected_as_double(rd);
632 + }
633 +
634 + return (NETDATA_DOUBLE)rrddim_timed_set_by_pointer(st, rd, now, (collected_number)value);
635 +}
636 +
637 collected_number rrddim_timed_set_by_pointer(RRDSET *st __maybe_unused, RRDDIM *rd, struct timeval collected_time, collected_number value) {
638 netdata_log_debug(D_RRD_CALLS, "rrddim_set_by_pointer() for chart %s, dimension %s, value " COLLECTED_NUMBER_FORMAT, rrdset_name(st), rrddim_name(rd), value);
639
640 rd->collector.last_collected_time = collected_time;
618 - rd->collector.collected_value = value;
641 + if(rrddim_is_float(rd))
642 + rrddim_set_collected_float(rd, (NETDATA_DOUBLE)value);
643 + else
644 + rrddim_set_collected_int(rd, value);
645 rrddim_set_updated(rd);
646 rd->collector.counter++;
647
@@ -625,11 +651,17 @@ collected_number rrddim_timed_set_by_pointer(RRDSET *st __maybe_unused, RRDDIM *
651 // *((int64_t *)Pvalue) = *((int64_t *)Pvalue) + 1;
652 // spinlock_unlock(&st->rrdhost->accounting.spinlock);
653
628 - collected_number v = (value >= 0) ? value : -value;
629 - if (unlikely(v > rd->collector.collected_value_max))
630 - rd->collector.collected_value_max = v;
631 -
632 - return rd->collector.last_collected_value;
654 + NETDATA_DOUBLE v = value >= 0 ? (NETDATA_DOUBLE)value : (NETDATA_DOUBLE)(-value);
655 + if (unlikely(v > rrddim_collected_max_as_double(rd))) {
656 + if(rrddim_is_float(rd))
657 + rrddim_set_collected_max_float(rd, v);
658 + else
659 + rrddim_set_collected_max_int(rd, (int64_t)v);
660 + }
661 + // For int dims return the last collected int; for float dims the integer return is meaningless, so return 0 to avoid truncation misuse.
662 + if(rrddim_is_float(rd))
663 + return 0;
664 + return rd->collector.collected.i.last_collected_value;
665 }
666
667
src/database/rrddim.h
+63 -3
@@ -27,6 +27,7 @@ typedef enum __attribute__ ((__packed__)) rrddim_options {
27 RRDDIM_OPTION_DONT_DETECT_RESETS_OR_OVERFLOWS = (1 << 1), // do not offer RESET or OVERFLOW info to callers
28 RRDDIM_OPTION_BACKFILLED_HIGH_TIERS = (1 << 2), // when set, we have backfilled higher tiers
29 RRDDIM_OPTION_UPDATED = (1 << 3), // single-threaded collector updated flag
30 + RRDDIM_OPTION_VALUE_FLOAT = (1 << 4), // collected_value lane should be treated as float
31
32 // this is 8-bit
33 } RRDDIM_OPTIONS;
@@ -34,6 +35,8 @@ typedef enum __attribute__ ((__packed__)) rrddim_options {
35 #define rrddim_option_check(rd, option) ((rd)->collector.options & (option))
36 #define rrddim_option_set(rd, option) (rd)->collector.options |= (option)
37 #define rrddim_option_clear(rd, option) (rd)->collector.options &= ~(option)
38 +#define rrddim_is_float(rd) ((rd)->collector.options & RRDDIM_OPTION_VALUE_FLOAT)
39 +#define rrddim_is_int(rd) (!rrddim_is_float(rd))
40
41 // flags are runtime changing status flags (atomics are required to alter/access them)
42 typedef enum __attribute__ ((__packed__)) rrddim_flags {
@@ -119,9 +122,18 @@ struct rrddim {
122
123 uint32_t counter; // the number of times we added values to this rrddim
124
122 - collected_number collected_value; // the current value, as collected - resets to 0 after being used
123 - collected_number collected_value_max; // the absolute maximum of the collected value
124 - collected_number last_collected_value; // the last value that was collected, after being processed
125 + union {
126 + struct {
127 + collected_number collected_value; // legacy int lane
128 + collected_number collected_value_max; // legacy int lane
129 + collected_number last_collected_value; // legacy int lane
130 + } i;
131 + struct {
132 + NETDATA_DOUBLE collected_value;
133 + NETDATA_DOUBLE collected_value_max;
134 + NETDATA_DOUBLE last_collected_value;
135 + } f;
136 + } collected; // type-aware storage (option bit selects lane)
137
138 struct timeval last_collected_time; // when was this dimension last updated
139 // this is actual date time we updated the last_collected_value
@@ -189,6 +201,53 @@ static inline bool rrddim_check_upstream_exposed_collector(RRDDIM *rd) {
201 return rd->rrdset->version == rd->stream.snd.sent_version;
202 }
203
204 +// ------------------------------------------------------------------------
205 +// Type-aware accessors (Phase 1: int-only behavior preserved)
206 +
207 +static inline void rrddim_set_collected_int(RRDDIM *rd, int64_t v) {
208 + rd->collector.collected.i.collected_value = (collected_number)v;
209 +}
210 +
211 +static inline void rrddim_set_collected_float(RRDDIM *rd, NETDATA_DOUBLE v) {
212 + rd->collector.collected.f.collected_value = v;
213 +}
214 +
215 +static inline void rrddim_set_collected_max_int(RRDDIM *rd, int64_t v) {
216 + rd->collector.collected.i.collected_value_max = (collected_number)v;
217 +}
218 +
219 +static inline void rrddim_set_collected_max_float(RRDDIM *rd, NETDATA_DOUBLE v) {
220 + rd->collector.collected.f.collected_value_max = v;
221 +}
222 +
223 +static inline void rrddim_set_last_collected_int(RRDDIM *rd, int64_t v) {
224 + rd->collector.collected.i.last_collected_value = (collected_number)v;
225 +}
226 +
227 +static inline void rrddim_set_last_collected_float(RRDDIM *rd, NETDATA_DOUBLE v) {
228 + rd->collector.collected.f.last_collected_value = v;
229 +}
230 +
231 +static inline NETDATA_DOUBLE rrddim_collected_as_double(RRDDIM *rd) {
232 + return rrddim_is_float(rd) ? rd->collector.collected.f.collected_value
233 + : (NETDATA_DOUBLE)rd->collector.collected.i.collected_value;
234 +}
235 +
236 +static inline NETDATA_DOUBLE rrddim_last_collected_as_double(RRDDIM *rd) {
237 + return rrddim_is_float(rd) ? rd->collector.collected.f.last_collected_value
238 + : (NETDATA_DOUBLE)rd->collector.collected.i.last_collected_value;
239 +}
240 +
241 +static inline NETDATA_DOUBLE rrddim_collected_max_as_double(RRDDIM *rd) {
242 + return rrddim_is_float(rd) ? rd->collector.collected.f.collected_value_max
243 + : (NETDATA_DOUBLE)rd->collector.collected.i.collected_value_max;
244 +}
245 +
246 +static inline int64_t rrddim_last_collected_raw_int(RRDDIM *rd) {
247 + return rrddim_is_float(rd) ? (int64_t)rd->collector.collected.f.last_collected_value
248 + : (int64_t)rd->collector.collected.i.last_collected_value;
249 +}
250 +
251 void rrddim_index_init(RRDSET *st);
252 void rrddim_index_destroy(RRDSET *st);
253
@@ -228,6 +287,7 @@ void rrddim_isnot_obsolete___safe_from_collector_thread(RRDSET *st, RRDDIM *rd);
287
288 collected_number rrddim_timed_set_by_pointer(RRDSET *st, RRDDIM *rd, struct timeval collected_time, collected_number value);
289 collected_number rrddim_set_by_pointer(RRDSET *st, RRDDIM *rd, collected_number value);
290 +NETDATA_DOUBLE rrddim_set_by_pointer_double(RRDSET *st, RRDDIM *rd, NETDATA_DOUBLE value);
291 collected_number rrddim_set(RRDSET *st, const char *id, collected_number value);
292
293 bool rrddim_finalize_collection_and_check_retention(RRDDIM *rd);
src/database/rrdset-collection.c
+91 -82
@@ -647,8 +647,8 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
647 size_t dim_id;
648 size_t dimensions = 0;
649 struct rda_item *rda = rda_base;
650 - total_number collected_total = 0;
651 - total_number last_collected_total = 0;
650 + NETDATA_DOUBLE collected_total = 0.0;
651 + NETDATA_DOUBLE last_collected_total = 0.0;
652 rrddim_foreach_read(rd, st) {
653 if(rd_dfe.counter >= rda_slots)
654 break;
@@ -664,22 +664,25 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
664 if(likely(rrddim_check_updated(rd))) {
665 // if the new is smaller than the old (an overflow, or reset), set the old equal to the new
666 // to reset the calculation (it will give zero as the calculation for this second)
667 - if(unlikely(rd->algorithm == RRD_ALGORITHM_PCENT_OVER_DIFF_TOTAL && rd->collector.last_collected_value > rd->collector.collected_value)) {
668 - netdata_log_debug(D_RRD_STATS, "'%s' / '%s': RESET or OVERFLOW. Last collected value = " COLLECTED_NUMBER_FORMAT ", current = " COLLECTED_NUMBER_FORMAT
667 + if(unlikely(rd->algorithm == RRD_ALGORITHM_PCENT_OVER_DIFF_TOTAL && rrddim_last_collected_as_double(rd) > rrddim_collected_as_double(rd))) {
668 + netdata_log_debug(D_RRD_STATS, "'%s' / '%s': RESET or OVERFLOW. Last collected value = " NETDATA_DOUBLE_FORMAT ", current = " NETDATA_DOUBLE_FORMAT
669 , rrdset_id(st)
670 , rrddim_name(rd)
671 - , rd->collector.last_collected_value
672 - , rd->collector.collected_value
671 + , rrddim_last_collected_as_double(rd)
672 + , rrddim_collected_as_double(rd)
673 );
674
675 if(!(rrddim_option_check(rd, RRDDIM_OPTION_DONT_DETECT_RESETS_OR_OVERFLOWS)))
676 rda->reset_or_overflow = true;
677
678 - rd->collector.last_collected_value = rd->collector.collected_value;
678 + if(rrddim_is_float(rd))
679 + rrddim_set_last_collected_float(rd, rrddim_collected_as_double(rd));
680 + else
681 + rrddim_set_last_collected_int(rd, rd->collector.collected.i.collected_value);
682 }
683
681 - last_collected_total += rd->collector.last_collected_value;
682 - collected_total += rd->collector.collected_value;
684 + last_collected_total += rrddim_last_collected_as_double(rd);
685 + collected_total += rrddim_collected_as_double(rd);
686
687 if(unlikely(rrddim_flag_check(rd, RRDDIM_FLAG_OBSOLETE))) {
688 netdata_log_error("Dimension %s in chart '%s' has the OBSOLETE flag set, but it is collected.", rrddim_name(rd), rrdset_id(st));
@@ -712,30 +715,30 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
715 }
716
717 rrdset_debug(st, "%s: START "
715 - " last_collected_value = " COLLECTED_NUMBER_FORMAT
716 - " collected_value = " COLLECTED_NUMBER_FORMAT
718 + " last_collected_value = " NETDATA_DOUBLE_FORMAT
719 + " collected_value = " NETDATA_DOUBLE_FORMAT
720 " last_calculated_value = " NETDATA_DOUBLE_FORMAT
721 " calculated_value = " NETDATA_DOUBLE_FORMAT
722 , rrddim_name(rd)
720 - , rd->collector.last_collected_value
721 - , rd->collector.collected_value
723 + , rrddim_last_collected_as_double(rd)
724 + , rrddim_collected_as_double(rd)
725 , rd->collector.last_calculated_value
726 , rd->collector.calculated_value
727 );
728
729 switch(rd->algorithm) {
730 case RRD_ALGORITHM_ABSOLUTE:
728 - rd->collector.calculated_value = (NETDATA_DOUBLE)rd->collector.collected_value
731 + rd->collector.calculated_value = rrddim_collected_as_double(rd)
732 * (NETDATA_DOUBLE)rd->multiplier
733 / (NETDATA_DOUBLE)rd->divisor;
734
735 rrdset_debug(st, "%s: CALC ABS/ABS-NO-IN " NETDATA_DOUBLE_FORMAT " = "
733 - COLLECTED_NUMBER_FORMAT
736 + NETDATA_DOUBLE_FORMAT
737 " * " NETDATA_DOUBLE_FORMAT
738 " / " NETDATA_DOUBLE_FORMAT
739 , rrddim_name(rd)
737 - , rd->collector.calculated_value
738 - , rd->collector.collected_value
740 + , rd->collector.calculated_value
741 + , rrddim_collected_as_double(rd)
742 , (NETDATA_DOUBLE)rd->multiplier
743 , (NETDATA_DOUBLE)rd->divisor
744 );
@@ -749,16 +752,16 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
752 // over the total of all dimensions
753 rd->collector.calculated_value =
754 (NETDATA_DOUBLE)100
752 - * (NETDATA_DOUBLE)rd->collector.collected_value
755 + * rrddim_collected_as_double(rd)
756 / (NETDATA_DOUBLE)collected_total;
757
758 rrdset_debug(st, "%s: CALC PCENT-ROW " NETDATA_DOUBLE_FORMAT " = 100"
756 - " * " COLLECTED_NUMBER_FORMAT
757 - " / " COLLECTED_NUMBER_FORMAT
759 + " * " NETDATA_DOUBLE_FORMAT
760 + " / " NETDATA_DOUBLE_FORMAT
761 , rrddim_name(rd)
762 , rd->collector.calculated_value
760 - , rd->collector.collected_value
761 - , collected_total
763 + , rrddim_collected_as_double(rd)
764 + , (NETDATA_DOUBLE)collected_total
765 );
766 break;
767
@@ -768,64 +771,64 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
771 continue;
772 }
773
771 - // If the new is smaller than the old (an overflow, or reset), set the old equal to the new
772 - // to reset the calculation (it will give zero as the calculation for this second).
773 - // It is imperative to set the comparison to uint64_t since type collected_number is signed and
774 - // produces wrong results as far as incremental counters are concerned.
775 - if(unlikely((uint64_t)rd->collector.last_collected_value > (uint64_t)rd->collector.collected_value)) {
776 - netdata_log_debug(D_RRD_STATS, "'%s' / '%s': RESET or OVERFLOW. Last collected value = " COLLECTED_NUMBER_FORMAT ", current = " COLLECTED_NUMBER_FORMAT
777 - , rrdset_id(st)
778 - , rrddim_name(rd)
779 - , rd->collector.last_collected_value
780 - , rd->collector.collected_value);
781 -
782 - if(!(rrddim_option_check(rd, RRDDIM_OPTION_DONT_DETECT_RESETS_OR_OVERFLOWS)))
783 - rda->reset_or_overflow = true;
784 -
785 - uint64_t last = (uint64_t)rd->collector.last_collected_value;
786 - uint64_t new = (uint64_t)rd->collector.collected_value;
787 - uint64_t max = (uint64_t)rd->collector.collected_value_max;
788 - uint64_t cap = 0;
789 -
790 - // Signed values are handled by exploiting two's complement which will produce positive deltas
791 - if (max > 0x00000000FFFFFFFFULL)
792 - cap = 0xFFFFFFFFFFFFFFFFULL; // handles signed and unsigned 64-bit counters
793 - else
794 - cap = 0x00000000FFFFFFFFULL; // handles signed and unsigned 32-bit counters
795 -
796 - uint64_t delta = cap - last + new;
797 - uint64_t max_acceptable_rate = (cap / 100) * MAX_INCREMENTAL_PERCENT_RATE;
798 -
799 - // If the delta is less than the maximum acceptable rate and the previous value was near the cap
800 - // then this is an overflow. There can be false positives such that a reset is detected as an
801 - // overflow.
802 - // TODO: remember recent history of rates and compare with current rate to reduce this chance.
803 - if (delta < max_acceptable_rate) {
774 + if(rrddim_is_int(rd)) {
775 + uint64_t last = (uint64_t)rd->collector.collected.i.last_collected_value;
776 + uint64_t new = (uint64_t)rd->collector.collected.i.collected_value;
777 + uint64_t max = (uint64_t)rd->collector.collected.i.collected_value_max;
778 +
779 + // If the new is smaller than the old (overflow/reset), handle wrap
780 + if(unlikely(last > new)) {
781 + netdata_log_debug(D_RRD_STATS, "'%s' / '%s': RESET or OVERFLOW. Last collected value = " NETDATA_DOUBLE_FORMAT ", current = " NETDATA_DOUBLE_FORMAT
782 + , rrdset_id(st)
783 + , rrddim_name(rd)
784 + , rrddim_last_collected_as_double(rd)
785 + , rrddim_collected_as_double(rd));
786 +
787 + if(!(rrddim_option_check(rd, RRDDIM_OPTION_DONT_DETECT_RESETS_OR_OVERFLOWS)))
788 + rda->reset_or_overflow = true;
789 +
790 + uint64_t cap = (max > 0x00000000FFFFFFFFULL) ? 0xFFFFFFFFFFFFFFFFULL : 0x00000000FFFFFFFFULL;
791 + uint64_t delta = cap - last + new;
792 + uint64_t max_acceptable_rate = (cap / 100) * MAX_INCREMENTAL_PERCENT_RATE;
793 +
794 + if (delta < max_acceptable_rate) {
795 + rd->collector.calculated_value +=
796 + (NETDATA_DOUBLE) delta
797 + * (NETDATA_DOUBLE) rd->multiplier
798 + / (NETDATA_DOUBLE) rd->divisor;
799 + } else {
800 + rd->collector.calculated_value += 0;
801 + }
802 + }
803 + else {
804 rd->collector.calculated_value +=
805 - (NETDATA_DOUBLE) delta
805 + (NETDATA_DOUBLE)((int64_t)(rd->collector.collected.i.collected_value - rd->collector.collected.i.last_collected_value))
806 * (NETDATA_DOUBLE) rd->multiplier
807 / (NETDATA_DOUBLE) rd->divisor;
808 - } else {
809 - // This is a reset. Any overflow with a rate greater than MAX_INCREMENTAL_PERCENT_RATE will also
810 - // be detected as a reset instead.
811 - rd->collector.calculated_value += (NETDATA_DOUBLE)0;
808 }
809 }
810 else {
815 - rd->collector.calculated_value +=
816 - (NETDATA_DOUBLE) (rd->collector.collected_value - rd->collector.last_collected_value)
817 - * (NETDATA_DOUBLE) rd->multiplier
818 - / (NETDATA_DOUBLE) rd->divisor;
811 + NETDATA_DOUBLE last = rrddim_last_collected_as_double(rd);
812 + NETDATA_DOUBLE cur = rrddim_collected_as_double(rd);
813 + if(unlikely(cur < last)) {
814 + if(!(rrddim_option_check(rd, RRDDIM_OPTION_DONT_DETECT_RESETS_OR_OVERFLOWS)))
815 + rda->reset_or_overflow = true;
816 + rd->collector.calculated_value += 0;
817 + }
818 + else {
819 + rd->collector.calculated_value +=
820 + (cur - last) * (NETDATA_DOUBLE) rd->multiplier / (NETDATA_DOUBLE) rd->divisor;
821 + }
822 }
823
824 rrdset_debug(st, "%s: CALC INC PRE " NETDATA_DOUBLE_FORMAT " = ("
822 - COLLECTED_NUMBER_FORMAT " - " COLLECTED_NUMBER_FORMAT
825 + NETDATA_DOUBLE_FORMAT " - " NETDATA_DOUBLE_FORMAT
826 ")"
827 " * " NETDATA_DOUBLE_FORMAT
828 " / " NETDATA_DOUBLE_FORMAT
829 , rrddim_name(rd)
830 , rd->collector.calculated_value
828 - , rd->collector.collected_value, rd->collector.last_collected_value
831 + , rrddim_collected_as_double(rd), rrddim_last_collected_as_double(rd)
832 , (NETDATA_DOUBLE)rd->multiplier
833 , (NETDATA_DOUBLE)rd->divisor
834 );
@@ -844,16 +847,16 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
847 else
848 rd->collector.calculated_value =
849 (NETDATA_DOUBLE)100
847 - * (NETDATA_DOUBLE)(rd->collector.collected_value - rd->collector.last_collected_value)
850 + * (NETDATA_DOUBLE)(rrddim_collected_as_double(rd) - rrddim_last_collected_as_double(rd))
851 / (NETDATA_DOUBLE)(collected_total - last_collected_total);
852
853 rrdset_debug(st, "%s: CALC PCENT-DIFF " NETDATA_DOUBLE_FORMAT " = 100"
851 - " * (" COLLECTED_NUMBER_FORMAT " - " COLLECTED_NUMBER_FORMAT ")"
852 - " / (" COLLECTED_NUMBER_FORMAT " - " COLLECTED_NUMBER_FORMAT ")"
854 + " * (" NETDATA_DOUBLE_FORMAT " - " NETDATA_DOUBLE_FORMAT ")"
855 + " / (" NETDATA_DOUBLE_FORMAT " - " NETDATA_DOUBLE_FORMAT ")"
856 , rrddim_name(rd)
857 , rd->collector.calculated_value
855 - , rd->collector.collected_value, rd->collector.last_collected_value
856 - , collected_total, last_collected_total
858 + , rrddim_collected_as_double(rd), rrddim_last_collected_as_double(rd)
859 + , (NETDATA_DOUBLE)collected_total, (NETDATA_DOUBLE)last_collected_total
860 );
861 break;
862
@@ -870,13 +873,13 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
873 }
874
875 rrdset_debug(st, "%s: PHASE2 "
873 - " last_collected_value = " COLLECTED_NUMBER_FORMAT
874 - " collected_value = " COLLECTED_NUMBER_FORMAT
876 + " last_collected_value = " NETDATA_DOUBLE_FORMAT
877 + " collected_value = " NETDATA_DOUBLE_FORMAT
878 " last_calculated_value = " NETDATA_DOUBLE_FORMAT
879 " calculated_value = " NETDATA_DOUBLE_FORMAT
880 , rrddim_name(rd)
878 - , rd->collector.last_collected_value
879 - , rd->collector.collected_value
881 + , rrddim_last_collected_as_double(rd)
882 + , rrddim_collected_as_double(rd)
883 , rd->collector.last_calculated_value
884 , rd->collector.calculated_value
885 );
@@ -913,9 +916,12 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
916 if(unlikely(!rrddim_check_updated(rd)))
917 continue;
918
916 - 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);
919 + rrdset_debug(st, "%s: setting last_collected_value (old: " NETDATA_DOUBLE_FORMAT ") to last_collected_value (new: " NETDATA_DOUBLE_FORMAT ")", rrddim_name(rd), rrddim_last_collected_as_double(rd), rrddim_collected_as_double(rd));
920
918 - rd->collector.last_collected_value = rd->collector.collected_value;
921 + if(rrddim_is_float(rd))
922 + rrddim_set_last_collected_float(rd, rrddim_collected_as_double(rd));
923 + else
924 + rrddim_set_last_collected_int(rd, rd->collector.collected.i.collected_value);
925
926 switch(rd->algorithm) {
927 case RRD_ALGORITHM_INCREMENTAL:
@@ -947,17 +953,20 @@ void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next)
953 }
954
955 rd->collector.calculated_value = 0;
950 - rd->collector.collected_value = 0;
956 + if(rrddim_is_float(rd))
957 + rrddim_set_collected_float(rd, 0.0);
958 + else
959 + rrddim_set_collected_int(rd, 0);
960 rrddim_clear_updated(rd);
961
962 rrdset_debug(st, "%s: END "
954 - " last_collected_value = " COLLECTED_NUMBER_FORMAT
955 - " collected_value = " COLLECTED_NUMBER_FORMAT
963 + " last_collected_value = " NETDATA_DOUBLE_FORMAT
964 + " collected_value = " NETDATA_DOUBLE_FORMAT
965 " last_calculated_value = " NETDATA_DOUBLE_FORMAT
966 " calculated_value = " NETDATA_DOUBLE_FORMAT
967 , rrddim_name(rd)
959 - , rd->collector.last_collected_value
960 - , rd->collector.collected_value
968 + , rrddim_last_collected_as_double(rd)
969 + , rrddim_collected_as_double(rd)
970 , rd->collector.last_calculated_value
971 , rd->collector.calculated_value
972 );
src/exporting/graphite/graphite.c
+22 -10
@@ -129,16 +129,28 @@ int format_dimension_collected_graphite_plaintext(struct instance *instance, RRD
129 (instance->config.options & EXPORTING_OPTION_SEND_NAMES && rd->name) ? rrddim_name(rd) : rrddim_id(rd),
130 RRD_ID_LENGTH_MAX);
131
132 - buffer_sprintf(
133 - instance->buffer,
134 - "%s.%s.%s.%s%s " COLLECTED_NUMBER_FORMAT " %llu\n",
135 - instance->config.prefix,
136 - (host == localhost) ? instance->config.hostname : rrdhost_hostname(host),
137 - chart_name,
138 - dimension_name,
139 - (instance->labels_buffer) ? buffer_tostring(instance->labels_buffer) : "",
140 - rd->collector.last_collected_value,
141 - (unsigned long long)rd->collector.last_collected_time.tv_sec);
132 + if(rrddim_is_float(rd))
133 + buffer_sprintf(
134 + instance->buffer,
135 + "%s.%s.%s.%s%s " NETDATA_DOUBLE_FORMAT " %llu\n",
136 + instance->config.prefix,
137 + (host == localhost) ? instance->config.hostname : rrdhost_hostname(host),
138 + chart_name,
139 + dimension_name,
140 + (instance->labels_buffer) ? buffer_tostring(instance->labels_buffer) : "",
141 + rrddim_last_collected_as_double(rd),
142 + (unsigned long long)rd->collector.last_collected_time.tv_sec);
143 + else
144 + buffer_sprintf(
145 + instance->buffer,
146 + "%s.%s.%s.%s%s " COLLECTED_NUMBER_FORMAT " %llu\n",
147 + instance->config.prefix,
148 + (host == localhost) ? instance->config.hostname : rrdhost_hostname(host),
149 + chart_name,
150 + dimension_name,
151 + (instance->labels_buffer) ? buffer_tostring(instance->labels_buffer) : "",
152 + rrddim_last_collected_raw_int(rd),
153 + (unsigned long long)rd->collector.last_collected_time.tv_sec);
154
155 return 0;
156 }
src/exporting/json/json.c
+8 -5
@@ -164,9 +164,7 @@ int format_dimension_collected_json_plaintext(struct instance *instance, RRDDIM
164
165 "\"id\":\"%s\","
166 "\"name\":\"%s\","
167 - "\"value\":" COLLECTED_NUMBER_FORMAT ","
168 -
169 - "\"timestamp\":%llu}",
167 + "\"value\":",
168
169 instance->config.prefix,
170 (host == localhost) ? instance->config.hostname : rrdhost_hostname(host),
@@ -179,9 +177,14 @@ int format_dimension_collected_json_plaintext(struct instance *instance, RRDDIM
177 rrdset_parts_type(st),
178 rrdset_units(st),
179 rrddim_id(rd),
182 - rrddim_name(rd),
183 - rd->collector.last_collected_value,
180 + rrddim_name(rd));
181 +
182 + if(rrddim_is_float(rd))
183 + buffer_sprintf(instance->buffer, NETDATA_DOUBLE_FORMAT, rrddim_last_collected_as_double(rd));
184 + else
185 + buffer_sprintf(instance->buffer, COLLECTED_NUMBER_FORMAT, rrddim_last_collected_raw_int(rd));
186
187 + buffer_sprintf(instance->buffer, ",\"timestamp\":%llu}",
188 (unsigned long long)rd->collector.last_collected_time.tv_sec);
189
190 if (instance->config.type != EXPORTING_CONNECTOR_TYPE_JSON_HTTP) {
src/exporting/opentsdb/opentsdb.c
+36 -16
@@ -180,16 +180,28 @@ int format_dimension_collected_opentsdb_telnet(struct instance *instance, RRDDIM
180 (instance->config.options & EXPORTING_OPTION_SEND_NAMES && rd->name) ? rrddim_name(rd) : rrddim_id(rd),
181 RRD_ID_LENGTH_MAX);
182
183 - buffer_sprintf(
184 - instance->buffer,
185 - "put %s.%s.%s %llu " COLLECTED_NUMBER_FORMAT " host=%s%s\n",
186 - instance->config.prefix,
187 - chart_name,
188 - dimension_name,
189 - (unsigned long long)rd->collector.last_collected_time.tv_sec,
190 - rd->collector.last_collected_value,
191 - (host == localhost) ? instance->config.hostname : rrdhost_hostname(host),
192 - (instance->labels_buffer) ? buffer_tostring(instance->labels_buffer) : "");
183 + if(rrddim_is_float(rd))
184 + buffer_sprintf(
185 + instance->buffer,
186 + "put %s.%s.%s %llu " NETDATA_DOUBLE_FORMAT " host=%s%s\n",
187 + instance->config.prefix,
188 + chart_name,
189 + dimension_name,
190 + (unsigned long long)rd->collector.last_collected_time.tv_sec,
191 + rrddim_last_collected_as_double(rd),
192 + (host == localhost) ? instance->config.hostname : rrdhost_hostname(host),
193 + (instance->labels_buffer) ? buffer_tostring(instance->labels_buffer) : "");
194 + else
195 + buffer_sprintf(
196 + instance->buffer,
197 + "put %s.%s.%s %llu " COLLECTED_NUMBER_FORMAT " host=%s%s\n",
198 + instance->config.prefix,
199 + chart_name,
200 + dimension_name,
201 + (unsigned long long)rd->collector.last_collected_time.tv_sec,
202 + rrddim_last_collected_raw_int(rd),
203 + (host == localhost) ? instance->config.hostname : rrdhost_hostname(host),
204 + (instance->labels_buffer) ? buffer_tostring(instance->labels_buffer) : "");
205
206 return 0;
207 }
@@ -316,16 +328,24 @@ int format_dimension_collected_opentsdb_http(struct instance *instance, RRDDIM *
328 "{"
329 "\"metric\":\"%s.%s.%s\","
330 "\"timestamp\":%llu,"
319 - "\"value\":"COLLECTED_NUMBER_FORMAT","
331 + "\"value\":",
332 + instance->config.prefix,
333 + chart_name,
334 + dimension_name,
335 + (unsigned long long)rd->collector.last_collected_time.tv_sec);
336 +
337 + if(rrddim_is_float(rd))
338 + buffer_sprintf(instance->buffer, NETDATA_DOUBLE_FORMAT, rrddim_last_collected_as_double(rd));
339 + else
340 + buffer_sprintf(instance->buffer, COLLECTED_NUMBER_FORMAT, rrddim_last_collected_raw_int(rd));
341 +
342 + buffer_sprintf(
343 + instance->buffer,
344 + ","
345 "\"tags\":{"
346 "\"host\":\"%s\"%s"
347 "}"
348 "}",
324 - instance->config.prefix,
325 - chart_name,
326 - dimension_name,
327 - (unsigned long long)rd->collector.last_collected_time.tv_sec,
328 - rd->collector.last_collected_value,
349 (host == localhost) ? instance->config.hostname : rrdhost_hostname(host),
350 instance->labels_buffer ? buffer_tostring(instance->labels_buffer) : "");
351
src/exporting/prometheus/prometheus.c
+3 -3
@@ -495,12 +495,12 @@ static void generate_as_collected_from_metric(BUFFER *wb,
495 buffer_putc(wb, '}');
496 buffer_putc(wb, ' ');
497
498 - if (prometheus_collector)
498 + if (prometheus_collector || rrddim_is_float(p->rd))
499 buffer_print_netdata_double(wb,
500 - (NETDATA_DOUBLE)p->rd->collector.last_collected_value * (NETDATA_DOUBLE)p->rd->multiplier /
500 + rrddim_last_collected_as_double(p->rd) * (NETDATA_DOUBLE)p->rd->multiplier /
501 (NETDATA_DOUBLE)p->rd->divisor);
502 else
503 - buffer_print_int64(wb, p->rd->collector.last_collected_value);
503 + buffer_print_int64(wb, p->rd->collector.collected.i.last_collected_value);
504
505 if (p->output_options & PROMETHEUS_OUTPUT_TIMESTAMPS) {
506 buffer_putc(wb, ' ');
src/exporting/prometheus/remote_write/remote_write.c
+2 -2
@@ -273,7 +273,7 @@ int format_dimension_prometheus_remote_write(struct instance *instance, RRDDIM *
273 connector_specific_data->write_request,
274 name, chart, family, dimension,
275 (host == localhost) ? instance->config.hostname : rrdhost_hostname(host),
276 - rd->collector.last_collected_value, timeval_msec(&rd->collector.last_collected_time));
276 + rrddim_last_collected_as_double(rd), timeval_msec(&rd->collector.last_collected_time));
277 } else {
278 // the dimensions of the chart, do not have the same algorithm, multiplier or divisor
279 // we create a metric per dimension
@@ -290,7 +290,7 @@ int format_dimension_prometheus_remote_write(struct instance *instance, RRDDIM *
290 connector_specific_data->write_request,
291 name, chart, family, NULL,
292 (host == localhost) ? instance->config.hostname : rrdhost_hostname(host),
293 - rd->collector.last_collected_value, timeval_msec(&rd->collector.last_collected_time));
293 + rrddim_last_collected_as_double(rd), timeval_msec(&rd->collector.last_collected_time));
294 }
295 } else {
296 // we need average or sum of the data
src/health/health_variable.c
+1 -1
@@ -75,7 +75,7 @@ static bool variable_lookup_in_chart(struct variable_lookup_job *vbd, RRDSET *st
75 variable_lookup_add_result_with_score(vbd, (NETDATA_DOUBLE)rd->collector.last_stored_value, st, "last stored value of dimension");
76 break;
77 case DIM_SELECT_RAW:
78 - variable_lookup_add_result_with_score(vbd, (NETDATA_DOUBLE)rd->collector.last_collected_value, st, "last collected value of dimension");
78 + variable_lookup_add_result_with_score(vbd, rrddim_last_collected_as_double(rd), st, "last collected value of dimension");
79 break;
80 case DIM_SELECT_LAST_COLLECTED:
81 variable_lookup_add_result_with_score(vbd, (NETDATA_DOUBLE)rd->collector.last_collected_time.tv_sec, st, "last collected time of dimension");
src/health/rrdvar.c
+8 -2
@@ -209,10 +209,16 @@ void health_api_v1_chart_variables2json(RRDSET *st, BUFFER *wb) {
209 RRDDIM *rd;
210 dfe_start_read(st->rrddim_root_index, rd) {
211 snprintfz(name, sizeof(name), "%s_raw", string2str(rd->id));
212 - buffer_json_member_add_int64(wb, name, rd->collector.last_collected_value);
212 + if(rrddim_is_float(rd))
213 + buffer_json_member_add_double(wb, name, rrddim_last_collected_as_double(rd));
214 + else
215 + buffer_json_member_add_int64(wb, name, rd->collector.collected.i.last_collected_value);
216 if(rd->name != rd->id) {
217 snprintfz(name, sizeof(name), "%s_raw", string2str(rd->name));
215 - buffer_json_member_add_int64(wb, name, rd->collector.last_collected_value);
218 + if(rrddim_is_float(rd))
219 + buffer_json_member_add_double(wb, name, rrddim_last_collected_as_double(rd));
220 + else
221 + buffer_json_member_add_int64(wb, name, rd->collector.collected.i.last_collected_value);
222 }
223 }
224 dfe_done(rd);
src/plugins.d/README.md
+9 -2
@@ -332,6 +332,15 @@ the template is:
332
333 a space separated list of options, enclosed in quotes. The following options are currently supported: `obsolete` to mark a chart as obsolete (Netdata will hide it and delete it after some time), `store_first` to make Netdata store the first collected value, assuming there was an invisible previous value set to zero (this is used by statsd charts - if the first data collected value of incremental dimensions is not zero based, unrealistic spikes will appear with this option set) and `hidden` to perform all operations on a chart, but do not offer it on dashboards (the chart will be send to external databases). `CHART` options have been added in Netdata v1.7 and the `hidden` option was added in 1.10.
334
335 + (for CHART options see above; DIMENSION-specific options are described below)
336 +
337 +#### DIMENSION options
338 +
339 +Additional option: `type=float` (default `type=int`).
340 +- `type=int` (default): values parsed as 64-bit integers; wrap/reset detection applies for incremental counters.
341 +- `type=float`: values parsed as double; incremental/delta-incremental allowed but wrap detection uses simple drop detection (no uint64 wrap math).
342 + Older parents without `FLOATBASELINE` capability will truncate baselines to int when streaming/replicating.
343 +
344 - `plugin` and `module`
345
346 both are just names that are used to let the user identify the plugin and the module that generated the chart. If `plugin` is unset or empty, Netdata will automatically set the filename of the plugin that generated the chart. `module` has not default.
@@ -897,5 +906,3 @@ There are a few rules for writing plugins properly:
906 3. If you are not sure of memory leaks, exit every one hour. Netdata will re-start your process.
907
908 4. If possible, try to autodetect if your plugin should be enabled, without any configuration.
900 -
901 -
src/plugins.d/pluginsd_parser.c
+53 -9
@@ -27,8 +27,12 @@ static inline PARSER_RC pluginsd_set(char **words, size_t num_words, PARSER *par
27 netdata_log_debug(D_PLUGINSD, "PLUGINSD: 'host:%s/chart:%s/dim:%s' SET is setting value to '%s'",
28 rrdhost_hostname(host), rrdset_id(st), dimension, value && *value ? value : "UNSET");
29
30 - if (value && *value)
31 - rrddim_set_by_pointer(st, rd, str2ll_encoded(value));
30 + if (value && *value) {
31 + if(rrddim_is_float(rd))
32 + rrddim_set_by_pointer_double(st, rd, str2ndd_encoded(value, NULL));
33 + else
34 + rrddim_set_by_pointer(st, rd, str2ll_encoded(value));
35 + }
36
37 return PARSER_RC_OK;
38 }
@@ -519,6 +523,17 @@ static inline PARSER_RC pluginsd_dimension(char **words, size_t num_words, PARSE
523 rrddim_option_set(rd, RRDDIM_OPTION_DONT_DETECT_RESETS_OR_OVERFLOWS);
524 if (strstr(options, "nooverflow") != NULL)
525 rrddim_option_set(rd, RRDDIM_OPTION_DONT_DETECT_RESETS_OR_OVERFLOWS);
526 +
527 + if (strstr(options, "type=float") != NULL) {
528 + if(!rrddim_is_float(rd))
529 + memset(&rd->collector.collected, 0, sizeof(rd->collector.collected));
530 + rrddim_option_set(rd, RRDDIM_OPTION_VALUE_FLOAT);
531 + }
532 + else if (strstr(options, "type=int") != NULL) {
533 + if(rrddim_is_float(rd))
534 + memset(&rd->collector.collected, 0, sizeof(rd->collector.collected));
535 + rrddim_option_clear(rd, RRDDIM_OPTION_VALUE_FLOAT);
536 + }
537 }
538 else
539 rrddim_isnot_obsolete___safe_from_collector_thread(st, rd);
@@ -949,11 +964,20 @@ static ALWAYS_INLINE PARSER_RC pluginsd_set_v2(char **words, size_t num_words, P
964 // ------------------------------------------------------------------------
965 // parse the parameters
966
952 - collected_number collected_value = (collected_number) str2ll_encoded(collected_str);
967 + // The sender only sends float baselines when it has STREAM_CAP_FLOAT_BASELINE;
968 + // older senders always send int64, even for float dimensions.
969 + bool sender_sent_float = rrddim_is_float(rd) && stream_has_capability(&parser->user, STREAM_CAP_FLOAT_BASELINE);
970 +
971 + collected_number collected_value = 0;
972 + NETDATA_DOUBLE collected_value_d = 0.0;
973 + if(sender_sent_float)
974 + collected_value_d = str2ndd_encoded(collected_str, NULL);
975 + else
976 + collected_value = (collected_number) str2ll_encoded(collected_str);
977
978 NETDATA_DOUBLE value;
979 if(*value_str == '#')
956 - value = (NETDATA_DOUBLE)collected_value;
980 + value = sender_sent_float ? collected_value_d : (NETDATA_DOUBLE)collected_value;
981 else
982 value = str2ndd_encoded(value_str, NULL);
983
@@ -1002,7 +1026,12 @@ static ALWAYS_INLINE PARSER_RC pluginsd_set_v2(char **words, size_t num_words, P
1026 // check if receiver and sender have the same number parsing capabilities
1027 bool can_copy = stream_has_capability(&parser->user, STREAM_CAP_IEEE754) == stream_has_capability(&parser->user.v2.stream_buffer, STREAM_CAP_IEEE754);
1028
1005 - // check the sender capabilities
1029 + // check if the float baseline capability matches between incoming and outgoing
1030 + bool downstream_float = rrddim_is_float(rd) && stream_has_capability(&parser->user.v2.stream_buffer, STREAM_CAP_FLOAT_BASELINE);
1031 + if(sender_sent_float != downstream_float)
1032 + can_copy = false;
1033 +
1034 + // check the downstream parent capabilities
1035 bool with_slots = stream_has_capability(&parser->user.v2.stream_buffer, STREAM_CAP_SLOTS) ? true : false;
1036 NUMBER_ENCODING integer_encoding = stream_has_capability(&parser->user.v2.stream_buffer, STREAM_CAP_IEEE754) ? NUMBER_ENCODING_BASE64 : NUMBER_ENCODING_HEX;
1037 NUMBER_ENCODING doubles_encoding = stream_has_capability(&parser->user.v2.stream_buffer, STREAM_CAP_IEEE754) ? NUMBER_ENCODING_BASE64 : NUMBER_ENCODING_DECIMAL;
@@ -1021,8 +1050,10 @@ static ALWAYS_INLINE PARSER_RC pluginsd_set_v2(char **words, size_t num_words, P
1050 buffer_fast_strcat(wb, "' ", 2);
1051 if(can_copy)
1052 buffer_strcat(wb, collected_str);
1053 + else if(downstream_float)
1054 + buffer_print_netdata_double_encoded(wb, doubles_encoding, sender_sent_float ? collected_value_d : (NETDATA_DOUBLE)collected_value);
1055 else
1025 - buffer_print_int64_encoded(wb, integer_encoding, collected_value); // original v2 had hex
1056 + buffer_print_int64_encoded(wb, integer_encoding, sender_sent_float ? (int64_t)collected_value_d : collected_value);
1057 buffer_fast_strcat(wb, " ", 1);
1058 if(can_copy)
1059 buffer_strcat(wb, value_str);
@@ -1041,7 +1072,12 @@ static ALWAYS_INLINE PARSER_RC pluginsd_set_v2(char **words, size_t num_words, P
1072 rrddim_store_metric(rd, parser->user.v2.end_time * USEC_PER_SEC, value, flags);
1073 rd->collector.last_collected_time.tv_sec = parser->user.v2.end_time;
1074 rd->collector.last_collected_time.tv_usec = 0;
1044 - rd->collector.last_collected_value = collected_value;
1075 + if(sender_sent_float)
1076 + rrddim_set_last_collected_float(rd, collected_value_d);
1077 + else if(rrddim_is_float(rd))
1078 + rrddim_set_last_collected_float(rd, (NETDATA_DOUBLE)collected_value);
1079 + else
1080 + rrddim_set_last_collected_int(rd, collected_value);
1081 rd->collector.last_stored_value = value;
1082 rd->collector.last_calculated_value = value;
1083 rd->collector.counter++;
@@ -1100,7 +1136,11 @@ static ALWAYS_INLINE PARSER_RC pluginsd_end_v2(char **words __maybe_unused, size
1136 continue;
1137
1138 rd->collector.calculated_value = 0;
1103 - rd->collector.collected_value = 0;
1139 + if(rrddim_is_float(rd)) {
1140 + rrddim_set_collected_float(rd, 0.0);
1141 + }
1142 + else
1143 + rrddim_set_collected_int(rd, 0);
1144 rrddim_clear_updated(rd);
1145 }
1146 }
@@ -1108,7 +1148,11 @@ static ALWAYS_INLINE PARSER_RC pluginsd_end_v2(char **words __maybe_unused, size
1148 RRDDIM *rd;
1149 rrddim_foreach_read(rd, st){
1150 rd->collector.calculated_value = 0;
1111 - rd->collector.collected_value = 0;
1151 + if(rrddim_is_float(rd)) {
1152 + rrddim_set_collected_float(rd, 0.0);
1153 + }
1154 + else
1155 + rrddim_set_collected_int(rd, 0);
1156 rrddim_clear_updated(rd);
1157 }
1158 rrddim_foreach_done(rd);
src/plugins.d/pluginsd_replication.c
+9 -1
@@ -315,7 +315,15 @@ ALWAYS_INLINE PARSER_RC pluginsd_replay_rrddim_collection_state(char **words, si
315 rd->collector.last_collected_time.tv_usec = (suseconds_t)(last_collected_ut % USEC_PER_SEC);
316 }
317
318 - rd->collector.last_collected_value = last_collected_value_str ? str2ll_encoded(last_collected_value_str) : 0;
318 + // The sender only sends float baselines when it has STREAM_CAP_FLOAT_BASELINE;
319 + // older senders always send int64, even for float dimensions.
320 + bool sender_sent_float = rrddim_is_float(rd) && (parser->user.capabilities & STREAM_CAP_FLOAT_BASELINE);
321 + if(sender_sent_float)
322 + rrddim_set_last_collected_float(rd, last_collected_value_str ? str2ndd_encoded(last_collected_value_str, NULL) : 0.0);
323 + else if(rrddim_is_float(rd))
324 + rrddim_set_last_collected_float(rd, last_collected_value_str ? (NETDATA_DOUBLE)str2ll_encoded(last_collected_value_str) : 0.0);
325 + else
326 + rrddim_set_last_collected_int(rd, last_collected_value_str ? str2ll_encoded(last_collected_value_str) : 0);
327 rd->collector.last_calculated_value = last_calculated_value_str ? str2ndd_encoded(last_calculated_value_str, NULL) : 0;
328 rd->collector.last_stored_value = last_stored_value_str ? str2ndd_encoded(last_stored_value_str, NULL) : 0.0;
329
src/streaming/protocol/command-begin-set-end-v1.c
+8 -1
@@ -29,7 +29,14 @@ void stream_send_rrdset_metrics_v1(RRDSET_STREAM_BUFFER *rsb, RRDSET *st) {
29 buffer_fast_strcat(wb, PLUGINSD_KEYWORD_SET " \"", 5);
30 buffer_fast_strcat(wb, rrddim_id(rd), string_strlen(rd->id));
31 buffer_fast_strcat(wb, "\" = ", 4);
32 - buffer_print_int64(wb, rd->collector.collected_value);
32 + if(rrddim_is_float(rd)) {
33 + if(stream_has_capability(rsb, STREAM_CAP_FLOAT_BASELINE))
34 + buffer_print_netdata_double(wb, rrddim_collected_as_double(rd));
35 + else
36 + buffer_print_int64(wb, (int64_t)rrddim_collected_as_double(rd));
37 + }
38 + else
39 + buffer_print_int64(wb, rd->collector.collected.i.collected_value);
40 buffer_fast_strcat(wb, "\n", 1);
41 }
42 else {
src/streaming/protocol/command-begin-set-end-v2.c
+12 -3
@@ -52,10 +52,20 @@ void stream_send_rrddim_metrics_v2(RRDSET_STREAM_BUFFER *rsb, RRDDIM *rd, usec_t
52 buffer_fast_strcat(wb, " '", 2);
53 buffer_fast_strcat(wb, rrddim_id(rd), string_strlen(rd->id));
54 buffer_fast_strcat(wb, "' ", 2);
55 - buffer_print_int64_encoded(wb, integer_encoding, rd->collector.last_collected_value);
55 + bool send_double_baseline = rrddim_is_float(rd) && stream_has_capability(rsb, STREAM_CAP_FLOAT_BASELINE);
56 + NETDATA_DOUBLE baseline_d = rrddim_last_collected_as_double(rd);
57 + int64_t baseline_i = rrddim_last_collected_raw_int(rd);
58 +
59 + if(send_double_baseline)
60 + buffer_print_netdata_double_encoded(wb, doubles_encoding, baseline_d);
61 + else
62 + buffer_print_int64_encoded(wb, integer_encoding, baseline_i);
63 +
64 buffer_fast_strcat(wb, " ", 1);
65
58 - if((NETDATA_DOUBLE)rd->collector.last_collected_value == n)
66 + NETDATA_DOUBLE baseline_cmp = send_double_baseline ? baseline_d : (NETDATA_DOUBLE)baseline_i;
67 +
68 + if(baseline_cmp == n)
69 buffer_fast_strcat(wb, "#", 1);
70 else
71 buffer_print_netdata_double_encoded(wb, doubles_encoding, n);
@@ -80,4 +90,3 @@ ALWAYS_INLINE void stream_send_rrdset_metrics_finished(RRDSET_STREAM_BUFFER *rsb
90
91 *rsb = (RRDSET_STREAM_BUFFER){ .wb = NULL, };
92 }
83 -
src/streaming/protocol/command-chart-definition.c
+2 -1
@@ -85,7 +85,7 @@ bool stream_sender_send_rrdset_definition(BUFFER *wb, RRDSET *st) {
85
86 buffer_sprintf(
87 wb
88 - , " \"%s\" \"%s\" \"%s\" %d %d \"%s %s %s\"\n"
88 + , " \"%s\" \"%s\" \"%s\" %d %d \"%s %s %s %s\"\n"
89 , rrddim_id(rd)
90 , rrddim_name(rd)
91 , rrd_algorithm_name(rd->algorithm)
@@ -94,6 +94,7 @@ bool stream_sender_send_rrdset_definition(BUFFER *wb, RRDSET *st) {
94 , rrddim_flag_check(rd, RRDDIM_FLAG_OBSOLETE)?"obsolete":""
95 , rrddim_option_check(rd, RRDDIM_OPTION_HIDDEN)?"hidden":""
96 , rrddim_option_check(rd, RRDDIM_OPTION_DONT_DETECT_RESETS_OR_OVERFLOWS)?"noreset":""
97 + , rrddim_is_float(rd)?"type=float":"type=int"
98 );
99 }
100 rrddim_foreach_done(rd);
src/streaming/stream-capabilities.c
+2
@@ -36,6 +36,7 @@ static struct {
36 {STREAM_CAP_PROGRESS, "PROGRESS" },
37 {STREAM_CAP_NODE_ID, "NODEID" },
38 {STREAM_CAP_PATHS, "PATHS" },
39 + {STREAM_CAP_FLOAT_BASELINE, "FLOATBASELINE" },
40
41 // terminator
42 {0 , NULL },
@@ -139,6 +140,7 @@ STREAM_CAPABILITIES stream_our_capabilities(RRDHOST *host, bool sender) {
140 STREAM_CAP_PATHS |
141 STREAM_CAP_IEEE754 |
142 STREAM_CAP_ML_MODELS |
143 + STREAM_CAP_FLOAT_BASELINE |
144 0) & ~disabled_capabilities;
145 }
146
src/streaming/stream-capabilities.h
+1
@@ -49,6 +49,7 @@ typedef enum {
49 STREAM_CAP_NODE_ID = (1 << 24), // support for sending NODE_ID back to the child
50 STREAM_CAP_PATHS = (1 << 25), // support for sending PATHS upstream and downstream
51 STREAM_CAP_ML_MODELS = (1 << 26), // support for sending MODELS upstream
52 + STREAM_CAP_FLOAT_BASELINE = (1 << 27), // support float baselines for dimensions
53
54 STREAM_CAP_INVALID = (1 << 30), // used as an invalid value for capabilities when this is set
55 // this must be signed int, so don't use the last bit
src/streaming/stream-replication-sender.c
+7 -1
@@ -236,7 +236,13 @@ static void replication_send_chart_collection_state(BUFFER *wb, RRDSET *st, STRE
236 buffer_print_uint64_encoded(wb, integer_encoding, (usec_t) rd->collector.last_collected_time.tv_sec * USEC_PER_SEC +
237 (usec_t) rd->collector.last_collected_time.tv_usec);
238 buffer_fast_strcat(wb, " ", 1);
239 - buffer_print_int64_encoded(wb, integer_encoding, rd->collector.last_collected_value);
239 +
240 + bool send_double_baseline = rrddim_is_float(rd) && (capabilities & STREAM_CAP_FLOAT_BASELINE);
241 + if(send_double_baseline)
242 + buffer_print_netdata_double_encoded(wb, integer_encoding, rrddim_last_collected_as_double(rd));
243 + else
244 + buffer_print_int64_encoded(wb, integer_encoding, rrddim_last_collected_raw_int(rd));
245 +
246 buffer_fast_strcat(wb, " ", 1);
247 buffer_print_netdata_double_encoded(wb, integer_encoding, rd->collector.last_calculated_value);
248 buffer_fast_strcat(wb, " ", 1);