@cryptotaxi247 / netdata-1 / commits / 4355e50bc

updated dbengine unittest (#17232)

* when a page cannot be acquired, repeat the call until it can or does not exist * updated dbengine unittest * Update dbengine-unittest.c

Costa Tsaousis committed Mar 23, 2024 at 12:41 UTC 4355e50bceb7c6aab409cb671c04577a31c6a868
4 files changed +877 -848
CMakeLists.txt
+2
@@ -1049,6 +1049,8 @@ if(ENABLE_DBENGINE)
1049 src/database/engine/metric.h
1050 src/database/engine/pdc.c
1051 src/database/engine/pdc.h
1052 + src/database/engine/dbengine-unittest.c
1053 + src/database/engine/dbengine-stresstest.c
1054 )
1055 endif()
1056
src/daemon/unit_test.c
-848
@@ -1804,851 +1804,3 @@ int unit_test_bitmaps(void) {
1804 fprintf(stderr, "%s() %d errors\n", __FUNCTION__, errors);
1805 return errors;
1806 }
1807 -
1808 -#ifdef ENABLE_DBENGINE
1809 -static inline void rrddim_set_by_pointer_fake_time(RRDDIM *rd, collected_number value, time_t now)
1810 -{
1811 - rd->collector.last_collected_time.tv_sec = now;
1812 - rd->collector.last_collected_time.tv_usec = 0;
1813 - rd->collector.collected_value = value;
1814 - rrddim_set_updated(rd);
1815 -
1816 - rd->collector.counter++;
1817 -
1818 - collected_number v = (value >= 0) ? value : -value;
1819 - if(unlikely(v > rd->collector.collected_value_max)) rd->collector.collected_value_max = v;
1820 -}
1821 -
1822 -static RRDHOST *dbengine_rrdhost_find_or_create(char *name)
1823 -{
1824 - /* We don't want to drop metrics when generating load, we prefer to block data generation itself */
1825 -
1826 - return rrdhost_find_or_create(
1827 - name,
1828 - name,
1829 - name,
1830 - os_type,
1831 - netdata_configured_timezone,
1832 - netdata_configured_abbrev_timezone,
1833 - netdata_configured_utc_offset,
1834 - program_name,
1835 - program_version,
1836 - default_rrd_update_every,
1837 - default_rrd_history_entries,
1838 - RRD_MEMORY_MODE_DBENGINE,
1839 - health_plugin_enabled(),
1840 - default_rrdpush_enabled,
1841 - default_rrdpush_destination,
1842 - default_rrdpush_api_key,
1843 - default_rrdpush_send_charts_matching,
1844 - default_rrdpush_enable_replication,
1845 - default_rrdpush_seconds_to_replicate,
1846 - default_rrdpush_replication_step,
1847 - NULL,
1848 - 0);
1849 -}
1850 -
1851 -// constants for test_dbengine
1852 -static const int CHARTS = 64;
1853 -static const int DIMS = 16; // That gives us 64 * 16 = 1024 metrics
1854 -#define REGIONS (3) // 3 regions of update_every
1855 -// first region update_every is 2, second is 3, third is 1
1856 -static const int REGION_UPDATE_EVERY[REGIONS] = {2, 3, 1};
1857 -static const int REGION_POINTS[REGIONS] = {
1858 - 16384, // This produces 64MiB of metric data for the first region: update_every = 2
1859 - 16384, // This produces 64MiB of metric data for the second region: update_every = 3
1860 - 16384, // This produces 64MiB of metric data for the third region: update_every = 1
1861 -};
1862 -static const int QUERY_BATCH = 4096;
1863 -
1864 -static void test_dbengine_create_charts(RRDHOST *host, RRDSET *st[CHARTS], RRDDIM *rd[CHARTS][DIMS],
1865 - int update_every)
1866 -{
1867 - fprintf(stderr, "%s() running...\n", __FUNCTION__ );
1868 - int i, j;
1869 - char name[101];
1870 -
1871 - for (i = 0 ; i < CHARTS ; ++i) {
1872 - snprintfz(name, sizeof(name) - 1, "dbengine-chart-%d", i);
1873 -
1874 - // create the chart
1875 - st[i] = rrdset_create(host, "netdata", name, name, "netdata", NULL, "Unit Testing", "a value", "unittest",
1876 - NULL, 1, update_every, RRDSET_TYPE_LINE);
1877 - rrdset_flag_set(st[i], RRDSET_FLAG_DEBUG);
1878 - rrdset_flag_set(st[i], RRDSET_FLAG_STORE_FIRST);
1879 - for (j = 0 ; j < DIMS ; ++j) {
1880 - snprintfz(name, sizeof(name) - 1, "dim-%d", j);
1881 -
1882 - rd[i][j] = rrddim_add(st[i], name, NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
1883 - }
1884 - }
1885 -
1886 - // Initialize DB with the very first entries
1887 - for (i = 0 ; i < CHARTS ; ++i) {
1888 - for (j = 0 ; j < DIMS ; ++j) {
1889 - rd[i][j]->collector.last_collected_time.tv_sec =
1890 - st[i]->last_collected_time.tv_sec = st[i]->last_updated.tv_sec = 2 * API_RELATIVE_TIME_MAX - 1;
1891 - rd[i][j]->collector.last_collected_time.tv_usec =
1892 - st[i]->last_collected_time.tv_usec = st[i]->last_updated.tv_usec = 0;
1893 - }
1894 - }
1895 - for (i = 0 ; i < CHARTS ; ++i) {
1896 - st[i]->usec_since_last_update = USEC_PER_SEC;
1897 -
1898 - for (j = 0; j < DIMS; ++j) {
1899 - rrddim_set_by_pointer_fake_time(rd[i][j], 69, 2 * API_RELATIVE_TIME_MAX); // set first value to 69
1900 - }
1901 -
1902 - struct timeval now;
1903 - now_realtime_timeval(&now);
1904 - rrdset_timed_done(st[i], now, false);
1905 - }
1906 - // Flush pages for subsequent real values
1907 - for (i = 0 ; i < CHARTS ; ++i) {
1908 - for (j = 0; j < DIMS; ++j) {
1909 - rrdeng_store_metric_flush_current_page((rd[i][j])->tiers[0].sch);
1910 - }
1911 - }
1912 -}
1913 -
1914 -// Feeds the database region with test data, returns last timestamp of region
1915 -static time_t test_dbengine_create_metrics(RRDSET *st[CHARTS], RRDDIM *rd[CHARTS][DIMS],
1916 - int current_region, time_t time_start)
1917 -{
1918 - fprintf(stderr, "%s() running...\n", __FUNCTION__ );
1919 - time_t time_now;
1920 - int i, j, c, update_every;
1921 - collected_number next;
1922 -
1923 - update_every = REGION_UPDATE_EVERY[current_region];
1924 - time_now = time_start;
1925 - // feed it with the test data
1926 - for (i = 0 ; i < CHARTS ; ++i) {
1927 - for (j = 0 ; j < DIMS ; ++j) {
1928 - storage_engine_store_change_collection_frequency(rd[i][j]->tiers[0].sch, update_every);
1929 -
1930 - rd[i][j]->collector.last_collected_time.tv_sec =
1931 - st[i]->last_collected_time.tv_sec = st[i]->last_updated.tv_sec = time_now;
1932 - rd[i][j]->collector.last_collected_time.tv_usec =
1933 - st[i]->last_collected_time.tv_usec = st[i]->last_updated.tv_usec = 0;
1934 - }
1935 - }
1936 - for (c = 0; c < REGION_POINTS[current_region] ; ++c) {
1937 - time_now += update_every; // time_now = start + (c + 1) * update_every
1938 -
1939 - for (i = 0 ; i < CHARTS ; ++i) {
1940 - st[i]->usec_since_last_update = USEC_PER_SEC * update_every;
1941 -
1942 - for (j = 0; j < DIMS; ++j) {
1943 - next = ((collected_number)i * DIMS) * REGION_POINTS[current_region] +
1944 - j * REGION_POINTS[current_region] + c;
1945 - rrddim_set_by_pointer_fake_time(rd[i][j], next, time_now);
1946 - }
1947 -
1948 - struct timeval now;
1949 - now.tv_sec = time_now;
1950 - now.tv_usec = 0;
1951 -
1952 - rrdset_timed_done(st[i], now, false);
1953 - }
1954 - }
1955 - return time_now; //time_end
1956 -}
1957 -
1958 -// Checks the metric data for the given region, returns number of errors
1959 -static int test_dbengine_check_metrics(RRDSET *st[CHARTS], RRDDIM *rd[CHARTS][DIMS],
1960 - int current_region, time_t time_start)
1961 -{
1962 - fprintf(stderr, "%s() running...\n", __FUNCTION__ );
1963 - uint8_t same;
1964 - time_t time_now, time_retrieved, end_time;
1965 - int i, j, k, c, errors, update_every;
1966 - collected_number last;
1967 - NETDATA_DOUBLE value, expected;
1968 - struct storage_engine_query_handle seqh;
1969 - size_t value_errors = 0, time_errors = 0;
1970 -
1971 - update_every = REGION_UPDATE_EVERY[current_region];
1972 - errors = 0;
1973 -
1974 - // check the result
1975 - for (c = 0; c < REGION_POINTS[current_region] ; c += QUERY_BATCH) {
1976 - time_now = time_start + (c + 1) * update_every;
1977 - for (i = 0 ; i < CHARTS ; ++i) {
1978 - for (j = 0; j < DIMS; ++j) {
1979 - storage_engine_query_init(rd[i][j]->tiers[0].seb, rd[i][j]->tiers[0].smh, &seqh, time_now, time_now + QUERY_BATCH * update_every, STORAGE_PRIORITY_NORMAL);
1980 - for (k = 0; k < QUERY_BATCH; ++k) {
1981 - last = ((collected_number)i * DIMS) * REGION_POINTS[current_region] +
1982 - j * REGION_POINTS[current_region] + c + k;
1983 - expected = unpack_storage_number(pack_storage_number((NETDATA_DOUBLE)last, SN_DEFAULT_FLAGS));
1984 -
1985 - STORAGE_POINT sp = storage_engine_query_next_metric(&seqh);
1986 - value = sp.sum;
1987 - time_retrieved = sp.start_time_s;
1988 - end_time = sp.end_time_s;
1989 -
1990 - same = (roundndd(value) == roundndd(expected)) ? 1 : 0;
1991 - if(!same) {
1992 - if(!value_errors)
1993 - fprintf(stderr, " DB-engine unittest %s/%s: at %lu secs, expecting value " NETDATA_DOUBLE_FORMAT
1994 - ", found " NETDATA_DOUBLE_FORMAT ", ### E R R O R ###\n",
1995 - rrdset_name(st[i]), rrddim_name(rd[i][j]), (unsigned long)time_now + k * update_every, expected, value);
1996 - value_errors++;
1997 - errors++;
1998 - }
1999 - if(end_time != time_now + k * update_every) {
2000 - if(!time_errors)
2001 - fprintf(stderr, " DB-engine unittest %s/%s: at %lu secs, found timestamp %lu ### E R R O R ###\n",
2002 - rrdset_name(st[i]), rrddim_name(rd[i][j]), (unsigned long)time_now + k * update_every, (unsigned long)time_retrieved);
2003 - time_errors++;
2004 - errors++;
2005 - }
2006 - }
2007 - storage_engine_query_finalize(&seqh);
2008 - }
2009 - }
2010 - }
2011 -
2012 - if(value_errors)
2013 - fprintf(stderr, "%zu value errors encountered\n", value_errors);
2014 -
2015 - if(time_errors)
2016 - fprintf(stderr, "%zu time errors encountered\n", time_errors);
2017 -
2018 - return errors;
2019 -}
2020 -
2021 -// Check rrdr transformations
2022 -static int test_dbengine_check_rrdr(RRDSET *st[CHARTS], RRDDIM *rd[CHARTS][DIMS],
2023 - int current_region, time_t time_start, time_t time_end)
2024 -{
2025 - int update_every = REGION_UPDATE_EVERY[current_region];
2026 - fprintf(stderr, "%s() running on region %d, start time %lld, end time %lld, update every %d, on %d dimensions...\n",
2027 - __FUNCTION__, current_region, (long long)time_start, (long long)time_end, update_every, CHARTS * DIMS);
2028 - uint8_t same;
2029 - time_t time_now, time_retrieved;
2030 - int i, j, errors, value_errors = 0, time_errors = 0, value_right = 0, time_right = 0;
2031 - long c;
2032 - collected_number last;
2033 - NETDATA_DOUBLE value, expected;
2034 -
2035 - errors = 0;
2036 - long points = (time_end - time_start) / update_every;
2037 - for (i = 0 ; i < CHARTS ; ++i) {
2038 - ONEWAYALLOC *owa = onewayalloc_create(0);
2039 - RRDR *r = rrd2rrdr_legacy(owa, st[i], points, time_start, time_end,
2040 - RRDR_GROUPING_AVERAGE, 0, RRDR_OPTION_NATURAL_POINTS,
2041 - NULL, NULL, 0, 0,
2042 - QUERY_SOURCE_UNITTEST, STORAGE_PRIORITY_NORMAL);
2043 - if (!r) {
2044 - fprintf(stderr, " DB-engine unittest %s: empty RRDR on region %d ### E R R O R ###\n", rrdset_name(st[i]), current_region);
2045 - return ++errors;
2046 - } else {
2047 - assert(r->internal.qt->request.st == st[i]);
2048 - for (c = 0; c != (long)rrdr_rows(r) ; ++c) {
2049 - RRDDIM *d;
2050 - time_now = time_start + (c + 1) * update_every;
2051 - time_retrieved = r->t[c];
2052 -
2053 - // for each dimension
2054 - rrddim_foreach_read(d, r->internal.qt->request.st) {
2055 - if(unlikely(d_dfe.counter >= r->d)) break; // d_counter is provided by the dictionary dfe
2056 -
2057 - j = (int)d_dfe.counter;
2058 -
2059 - NETDATA_DOUBLE *cn = &r->v[ c * r->d ];
2060 - value = cn[j];
2061 - assert(rd[i][j] == d);
2062 -
2063 - last = i * DIMS * REGION_POINTS[current_region] + j * REGION_POINTS[current_region] + c;
2064 - expected = unpack_storage_number(pack_storage_number((NETDATA_DOUBLE)last, SN_DEFAULT_FLAGS));
2065 -
2066 - same = (roundndd(value) == roundndd(expected)) ? 1 : 0;
2067 - if(!same) {
2068 - if(value_errors < 20)
2069 - fprintf(stderr, " DB-engine unittest %s/%s: point #%ld, at %lu secs, expecting value " NETDATA_DOUBLE_FORMAT
2070 - ", RRDR found " NETDATA_DOUBLE_FORMAT ", ### E R R O R ###\n",
2071 - rrdset_name(st[i]), rrddim_name(rd[i][j]), (long) c+1, (unsigned long)time_now, expected, value);
2072 - value_errors++;
2073 - }
2074 - else
2075 - value_right++;
2076 -
2077 - if(time_retrieved != time_now) {
2078 - if(time_errors < 20)
2079 - fprintf(stderr, " DB-engine unittest %s/%s: point #%ld at %lu secs, found RRDR timestamp %lu ### E R R O R ###\n",
2080 - rrdset_name(st[i]), rrddim_name(rd[i][j]), (long)c+1, (unsigned long)time_now, (unsigned long)time_retrieved);
2081 - time_errors++;
2082 - }
2083 - else
2084 - time_right++;
2085 - }
2086 - rrddim_foreach_done(d);
2087 - }
2088 - rrdr_free(owa, r);
2089 - }
2090 - onewayalloc_destroy(owa);
2091 - }
2092 -
2093 - if(value_errors)
2094 - fprintf(stderr, "%d value errors encountered (%d were ok)\n", value_errors, value_right);
2095 -
2096 - if(time_errors)
2097 - fprintf(stderr, "%d time errors encountered (%d were ok)\n", time_errors, value_right);
2098 -
2099 - return errors + value_errors + time_errors;
2100 -}
2101 -
2102 -void test_dbengine_charts_and_dims_are_not_collected(RRDSET *st[CHARTS], RRDDIM *rd[CHARTS][DIMS]) {
2103 - for(int c = 0; c < CHARTS ; c++) {
2104 - st[c]->rrdcontexts.collected = false;
2105 - for(int d = 0; d < DIMS ; d++)
2106 - rd[c][d]->rrdcontexts.collected = false;
2107 - }
2108 -}
2109 -
2110 -int test_dbengine(void)
2111 -{
2112 - fprintf(stderr, "%s() running...\n", __FUNCTION__ );
2113 - int i, j, errors = 0, value_errors = 0, time_errors = 0, update_every, current_region;
2114 - RRDHOST *host = NULL;
2115 - RRDSET *st[CHARTS];
2116 - RRDDIM *rd[CHARTS][DIMS];
2117 - time_t time_start[REGIONS], time_end[REGIONS];
2118 -
2119 - nd_log_limits_unlimited();
2120 - fprintf(stderr, "\nRunning DB-engine test\n");
2121 -
2122 - default_rrd_memory_mode = RRD_MEMORY_MODE_DBENGINE;
2123 -
2124 - fprintf(stderr, "Initializing localhost with hostname 'unittest-dbengine'");
2125 - host = dbengine_rrdhost_find_or_create("unittest-dbengine");
2126 - if (NULL == host)
2127 - return 1;
2128 -
2129 - current_region = 0; // this is the first region of data
2130 - update_every = REGION_UPDATE_EVERY[current_region]; // set data collection frequency to 2 seconds
2131 - test_dbengine_create_charts(host, st, rd, update_every);
2132 -
2133 - time_start[current_region] = 2 * API_RELATIVE_TIME_MAX;
2134 - time_end[current_region] = test_dbengine_create_metrics(st,rd, current_region, time_start[current_region]);
2135 -
2136 - errors += test_dbengine_check_metrics(st, rd, current_region, time_start[current_region]);
2137 - test_dbengine_charts_and_dims_are_not_collected(st, rd);
2138 -
2139 - current_region = 1; //this is the second region of data
2140 - update_every = REGION_UPDATE_EVERY[current_region]; // set data collection frequency to 3 seconds
2141 - // Align pages for frequency change
2142 - for (i = 0 ; i < CHARTS ; ++i) {
2143 - st[i]->update_every = update_every;
2144 - for (j = 0; j < DIMS; ++j) {
2145 - rrdeng_store_metric_flush_current_page((rd[i][j])->tiers[0].sch);
2146 - }
2147 - }
2148 -
2149 - time_start[current_region] = time_end[current_region - 1] + update_every;
2150 - if (0 != time_start[current_region] % update_every) // align to update_every
2151 - time_start[current_region] += update_every - time_start[current_region] % update_every;
2152 - time_end[current_region] = test_dbengine_create_metrics(st,rd, current_region, time_start[current_region]);
2153 -
2154 - errors += test_dbengine_check_metrics(st, rd, current_region, time_start[current_region]);
2155 - test_dbengine_charts_and_dims_are_not_collected(st, rd);
2156 -
2157 - current_region = 2; //this is the third region of data
2158 - update_every = REGION_UPDATE_EVERY[current_region]; // set data collection frequency to 1 seconds
2159 - // Align pages for frequency change
2160 - for (i = 0 ; i < CHARTS ; ++i) {
2161 - st[i]->update_every = update_every;
2162 - for (j = 0; j < DIMS; ++j) {
2163 - rrdeng_store_metric_flush_current_page((rd[i][j])->tiers[0].sch);
2164 - }
2165 - }
2166 -
2167 - time_start[current_region] = time_end[current_region - 1] + update_every;
2168 - if (0 != time_start[current_region] % update_every) // align to update_every
2169 - time_start[current_region] += update_every - time_start[current_region] % update_every;
2170 - time_end[current_region] = test_dbengine_create_metrics(st,rd, current_region, time_start[current_region]);
2171 -
2172 - errors += test_dbengine_check_metrics(st, rd, current_region, time_start[current_region]);
2173 - test_dbengine_charts_and_dims_are_not_collected(st, rd);
2174 -
2175 - for (current_region = 0 ; current_region < REGIONS ; ++current_region) {
2176 - errors += test_dbengine_check_rrdr(st, rd, current_region, time_start[current_region], time_end[current_region]);
2177 - }
2178 -
2179 - current_region = 1;
2180 - update_every = REGION_UPDATE_EVERY[current_region]; // use the maximum update_every = 3
2181 - long points = (time_end[REGIONS - 1] - time_start[0]) / update_every; // cover all time regions with RRDR
2182 - long point_offset = (time_start[current_region] - time_start[0]) / update_every;
2183 - for (i = 0 ; i < CHARTS ; ++i) {
2184 - ONEWAYALLOC *owa = onewayalloc_create(0);
2185 - RRDR *r = rrd2rrdr_legacy(owa, st[i], points, time_start[0] + update_every,
2186 - time_end[REGIONS - 1], RRDR_GROUPING_AVERAGE, 0,
2187 - RRDR_OPTION_NATURAL_POINTS, NULL, NULL, 0, 0,
2188 - QUERY_SOURCE_UNITTEST, STORAGE_PRIORITY_NORMAL);
2189 -
2190 - if (!r) {
2191 - fprintf(stderr, " DB-engine unittest %s: empty RRDR ### E R R O R ###\n", rrdset_name(st[i]));
2192 - ++errors;
2193 - } else {
2194 - long c;
2195 -
2196 - assert(r->internal.qt->request.st == st[i]);
2197 - // test current region values only, since they must be left unchanged
2198 - for (c = point_offset ; c < (long)(point_offset + rrdr_rows(r) / REGIONS / 2) ; ++c) {
2199 - RRDDIM *d;
2200 - time_t time_now = time_start[current_region] + (c - point_offset + 2) * update_every;
2201 - time_t time_retrieved = r->t[c];
2202 -
2203 - // for each dimension
2204 - rrddim_foreach_read(d, r->internal.qt->request.st) {
2205 - if(unlikely(d_dfe.counter >= r->d)) break; // d_counter is provided by the dictionary dfe
2206 -
2207 - j = (int)d_dfe.counter;
2208 -
2209 - NETDATA_DOUBLE *cn = &r->v[ c * r->d ];
2210 - NETDATA_DOUBLE value = cn[j];
2211 - assert(rd[i][j] == d);
2212 -
2213 - collected_number last = i * DIMS * REGION_POINTS[current_region] + j * REGION_POINTS[current_region] + c - point_offset + 1;
2214 - NETDATA_DOUBLE expected = unpack_storage_number(pack_storage_number((NETDATA_DOUBLE)last, SN_DEFAULT_FLAGS));
2215 -
2216 - uint8_t same = (roundndd(value) == roundndd(expected)) ? 1 : 0;
2217 - if(!same) {
2218 - if(!value_errors)
2219 - fprintf(stderr, " DB-engine unittest %s/%s: at %lu secs, expecting value " NETDATA_DOUBLE_FORMAT
2220 - ", RRDR found " NETDATA_DOUBLE_FORMAT ", ### E R R O R ###\n",
2221 - rrdset_name(st[i]), rrddim_name(rd[i][j]), (unsigned long)time_now, expected, value);
2222 - value_errors++;
2223 - }
2224 - if(time_retrieved != time_now) {
2225 - if(!time_errors)
2226 - fprintf(stderr, " DB-engine unittest %s/%s: at %lu secs, found RRDR timestamp %lu ### E R R O R ###\n",
2227 - rrdset_name(st[i]), rrddim_name(rd[i][j]), (unsigned long)time_now, (unsigned long)time_retrieved);
2228 - time_errors++;
2229 - }
2230 - }
2231 - rrddim_foreach_done(d);
2232 - }
2233 - rrdr_free(owa, r);
2234 - }
2235 - onewayalloc_destroy(owa);
2236 - }
2237 -
2238 - rrd_wrlock();
2239 - rrdeng_prepare_exit((struct rrdengine_instance *)host->db[0].si);
2240 - rrdeng_exit((struct rrdengine_instance *)host->db[0].si);
2241 - rrdeng_enq_cmd(NULL, RRDENG_OPCODE_SHUTDOWN_EVLOOP, NULL, NULL, STORAGE_PRIORITY_BEST_EFFORT, NULL, NULL);
2242 - rrd_unlock();
2243 -
2244 - return errors + value_errors + time_errors;
2245 -}
2246 -
2247 -struct dbengine_chart_thread {
2248 - uv_thread_t thread;
2249 - RRDHOST *host;
2250 - char *chartname; /* Will be prefixed by type, e.g. "example_local1.", "example_local2." etc */
2251 - unsigned dset_charts; /* number of charts */
2252 - unsigned dset_dims; /* dimensions per chart */
2253 - unsigned chart_i; /* current chart offset */
2254 - time_t time_present; /* current virtual time of the benchmark */
2255 - volatile time_t time_max; /* latest timestamp of stored values */
2256 - unsigned history_seconds; /* how far back in the past to go */
2257 -
2258 - volatile long done; /* initialize to 0, set to 1 to stop thread */
2259 - struct completion charts_initialized;
2260 - unsigned long errors, stored_metrics_nr; /* statistics */
2261 -
2262 - RRDSET *st;
2263 - RRDDIM *rd[]; /* dset_dims elements */
2264 -};
2265 -
2266 -collected_number generate_dbengine_chart_value(int chart_i, int dim_i, time_t time_current)
2267 -{
2268 - collected_number value;
2269 -
2270 - value = ((collected_number)time_current) * (chart_i + 1);
2271 - value += ((collected_number)time_current) * (dim_i + 1);
2272 - value %= 1024LLU;
2273 -
2274 - return value;
2275 -}
2276 -
2277 -static void generate_dbengine_chart(void *arg)
2278 -{
2279 - fprintf(stderr, "%s() running...\n", __FUNCTION__ );
2280 - struct dbengine_chart_thread *thread_info = (struct dbengine_chart_thread *)arg;
2281 - RRDHOST *host = thread_info->host;
2282 - char *chartname = thread_info->chartname;
2283 - const unsigned DSET_DIMS = thread_info->dset_dims;
2284 - unsigned history_seconds = thread_info->history_seconds;
2285 - time_t time_present = thread_info->time_present;
2286 -
2287 - unsigned j, update_every = 1;
2288 - RRDSET *st;
2289 - RRDDIM *rd[DSET_DIMS];
2290 - char name[RRD_ID_LENGTH_MAX + 1];
2291 - time_t time_current;
2292 -
2293 - // create the chart
2294 - snprintfz(name, RRD_ID_LENGTH_MAX, "example_local%u", thread_info->chart_i + 1);
2295 - thread_info->st = st = rrdset_create(host, name, chartname, chartname, "example", NULL, chartname, chartname,
2296 - chartname, NULL, 1, update_every, RRDSET_TYPE_LINE);
2297 - for (j = 0 ; j < DSET_DIMS ; ++j) {
2298 - snprintfz(name, RRD_ID_LENGTH_MAX, "%s%u", chartname, j + 1);
2299 -
2300 - thread_info->rd[j] = rd[j] = rrddim_add(st, name, NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
2301 - }
2302 - completion_mark_complete(&thread_info->charts_initialized);
2303 -
2304 - // feed it with the test data
2305 - time_current = time_present - history_seconds;
2306 - for (j = 0 ; j < DSET_DIMS ; ++j) {
2307 - rd[j]->collector.last_collected_time.tv_sec =
2308 - st->last_collected_time.tv_sec = st->last_updated.tv_sec = time_current - update_every;
2309 - rd[j]->collector.last_collected_time.tv_usec =
2310 - st->last_collected_time.tv_usec = st->last_updated.tv_usec = 0;
2311 - }
2312 - for( ; !thread_info->done && time_current < time_present ; time_current += update_every) {
2313 - st->usec_since_last_update = USEC_PER_SEC * update_every;
2314 -
2315 - for (j = 0; j < DSET_DIMS; ++j) {
2316 - collected_number value;
2317 -
2318 - value = generate_dbengine_chart_value(thread_info->chart_i, j, time_current);
2319 - rrddim_set_by_pointer_fake_time(rd[j], value, time_current);
2320 - ++thread_info->stored_metrics_nr;
2321 - }
2322 - rrdset_done(st);
2323 - thread_info->time_max = time_current;
2324 - }
2325 - for (j = 0; j < DSET_DIMS; ++j) {
2326 - rrdeng_store_metric_finalize((rd[j])->tiers[0].sch);
2327 - }
2328 -}
2329 -
2330 -void generate_dbengine_dataset(unsigned history_seconds)
2331 -{
2332 - fprintf(stderr, "%s() running...\n", __FUNCTION__ );
2333 - const int DSET_CHARTS = 16;
2334 - const int DSET_DIMS = 128;
2335 - const uint64_t EXPECTED_COMPRESSION_RATIO = 20;
2336 - RRDHOST *host = NULL;
2337 - struct dbengine_chart_thread **thread_info;
2338 - int i;
2339 - time_t time_present;
2340 -
2341 - default_rrd_memory_mode = RRD_MEMORY_MODE_DBENGINE;
2342 - default_rrdeng_page_cache_mb = 128;
2343 - // Worst case for uncompressible data
2344 - default_rrdeng_disk_quota_mb = (((uint64_t)DSET_DIMS * DSET_CHARTS) * sizeof(storage_number) * history_seconds) /
2345 - (1024 * 1024);
2346 - default_rrdeng_disk_quota_mb -= default_rrdeng_disk_quota_mb * EXPECTED_COMPRESSION_RATIO / 100;
2347 -
2348 - nd_log_limits_unlimited();
2349 - fprintf(stderr, "Initializing localhost with hostname 'dbengine-dataset'");
2350 -
2351 - host = dbengine_rrdhost_find_or_create("dbengine-dataset");
2352 - if (NULL == host)
2353 - return;
2354 -
2355 - thread_info = mallocz(sizeof(*thread_info) * DSET_CHARTS);
2356 - for (i = 0 ; i < DSET_CHARTS ; ++i) {
2357 - thread_info[i] = mallocz(sizeof(*thread_info[i]) + sizeof(RRDDIM *) * DSET_DIMS);
2358 - }
2359 - fprintf(stderr, "\nRunning DB-engine workload generator\n");
2360 -
2361 - time_present = now_realtime_sec();
2362 - for (i = 0 ; i < DSET_CHARTS ; ++i) {
2363 - thread_info[i]->host = host;
2364 - thread_info[i]->chartname = "random";
2365 - thread_info[i]->dset_charts = DSET_CHARTS;
2366 - thread_info[i]->chart_i = i;
2367 - thread_info[i]->dset_dims = DSET_DIMS;
2368 - thread_info[i]->history_seconds = history_seconds;
2369 - thread_info[i]->time_present = time_present;
2370 - thread_info[i]->time_max = 0;
2371 - thread_info[i]->done = 0;
2372 - completion_init(&thread_info[i]->charts_initialized);
2373 - fatal_assert(0 == uv_thread_create(&thread_info[i]->thread, generate_dbengine_chart, thread_info[i]));
2374 - completion_wait_for(&thread_info[i]->charts_initialized);
2375 - completion_destroy(&thread_info[i]->charts_initialized);
2376 - }
2377 - for (i = 0 ; i < DSET_CHARTS ; ++i) {
2378 - fatal_assert(0 == uv_thread_join(&thread_info[i]->thread));
2379 - }
2380 -
2381 - for (i = 0 ; i < DSET_CHARTS ; ++i) {
2382 - freez(thread_info[i]);
2383 - }
2384 - freez(thread_info);
2385 - rrd_wrlock();
2386 - rrdhost_free___while_having_rrd_wrlock(localhost, true);
2387 - rrd_unlock();
2388 -}
2389 -
2390 -struct dbengine_query_thread {
2391 - uv_thread_t thread;
2392 - RRDHOST *host;
2393 - char *chartname; /* Will be prefixed by type, e.g. "example_local1.", "example_local2." etc */
2394 - unsigned dset_charts; /* number of charts */
2395 - unsigned dset_dims; /* dimensions per chart */
2396 - time_t time_present; /* current virtual time of the benchmark */
2397 - unsigned history_seconds; /* how far back in the past to go */
2398 - volatile long done; /* initialize to 0, set to 1 to stop thread */
2399 - unsigned long errors, queries_nr, queried_metrics_nr; /* statistics */
2400 - uint8_t delete_old_data; /* if non zero then data are deleted when disk space is exhausted */
2401 -
2402 - struct dbengine_chart_thread *chart_threads[]; /* dset_charts elements */
2403 -};
2404 -
2405 -static void query_dbengine_chart(void *arg)
2406 -{
2407 - fprintf(stderr, "%s() running...\n", __FUNCTION__ );
2408 - struct dbengine_query_thread *thread_info = (struct dbengine_query_thread *)arg;
2409 - const int DSET_CHARTS = thread_info->dset_charts;
2410 - const int DSET_DIMS = thread_info->dset_dims;
2411 - time_t time_after, time_before, time_min, time_approx_min, time_max, duration;
2412 - int i, j, update_every = 1;
2413 - RRDSET *st;
2414 - RRDDIM *rd;
2415 - uint8_t same;
2416 - time_t time_now, time_retrieved, end_time;
2417 - collected_number generatedv;
2418 - NETDATA_DOUBLE value, expected;
2419 - struct storage_engine_query_handle seqh;
2420 - size_t value_errors = 0, time_errors = 0;
2421 -
2422 - do {
2423 - // pick a chart and dimension
2424 - i = random() % DSET_CHARTS;
2425 - st = thread_info->chart_threads[i]->st;
2426 - j = random() % DSET_DIMS;
2427 - rd = thread_info->chart_threads[i]->rd[j];
2428 -
2429 - time_min = thread_info->time_present - thread_info->history_seconds + 1;
2430 - time_max = thread_info->chart_threads[i]->time_max;
2431 -
2432 - if (thread_info->delete_old_data) {
2433 - /* A time window of twice the disk space is sufficient for compression space savings of up to 50% */
2434 - time_approx_min = time_max - (default_rrdeng_disk_quota_mb * 2 * 1024 * 1024) /
2435 - (((uint64_t) DSET_DIMS * DSET_CHARTS) * sizeof(storage_number));
2436 - time_min = MAX(time_min, time_approx_min);
2437 - }
2438 - if (!time_max) {
2439 - time_before = time_after = time_min;
2440 - } else {
2441 - time_after = time_min + random() % (MAX(time_max - time_min, 1));
2442 - duration = random() % 3600;
2443 - time_before = MIN(time_after + duration, time_max); /* up to 1 hour queries */
2444 - }
2445 -
2446 - storage_engine_query_init(rd->tiers[0].seb, rd->tiers[0].smh, &seqh, time_after, time_before, STORAGE_PRIORITY_NORMAL);
2447 - ++thread_info->queries_nr;
2448 - for (time_now = time_after ; time_now <= time_before ; time_now += update_every) {
2449 - generatedv = generate_dbengine_chart_value(i, j, time_now);
2450 - expected = unpack_storage_number(pack_storage_number((NETDATA_DOUBLE) generatedv, SN_DEFAULT_FLAGS));
2451 -
2452 - if (unlikely(storage_engine_query_is_finished(&seqh))) {
2453 - if (!thread_info->delete_old_data) { /* data validation only when we don't delete */
2454 - fprintf(stderr, " DB-engine stresstest %s/%s: at %lu secs, expecting value " NETDATA_DOUBLE_FORMAT
2455 - ", found data gap, ### E R R O R ###\n",
2456 - rrdset_name(st), rrddim_name(rd), (unsigned long) time_now, expected);
2457 - ++thread_info->errors;
2458 - }
2459 - break;
2460 - }
2461 -
2462 - STORAGE_POINT sp = storage_engine_query_next_metric(&seqh);
2463 - value = sp.sum;
2464 - time_retrieved = sp.start_time_s;
2465 - end_time = sp.end_time_s;
2466 -
2467 - if (!netdata_double_isnumber(value)) {
2468 - if (!thread_info->delete_old_data) { /* data validation only when we don't delete */
2469 - fprintf(stderr, " DB-engine stresstest %s/%s: at %lu secs, expecting value " NETDATA_DOUBLE_FORMAT
2470 - ", found data gap, ### E R R O R ###\n",
2471 - rrdset_name(st), rrddim_name(rd), (unsigned long) time_now, expected);
2472 - ++thread_info->errors;
2473 - }
2474 - break;
2475 - }
2476 - ++thread_info->queried_metrics_nr;
2477 -
2478 - same = (roundndd(value) == roundndd(expected)) ? 1 : 0;
2479 - if (!same) {
2480 - if (!thread_info->delete_old_data) { /* data validation only when we don't delete */
2481 - if(!value_errors)
2482 - fprintf(stderr, " DB-engine stresstest %s/%s: at %lu secs, expecting value " NETDATA_DOUBLE_FORMAT
2483 - ", found " NETDATA_DOUBLE_FORMAT ", ### E R R O R ###\n",
2484 - rrdset_name(st), rrddim_name(rd), (unsigned long) time_now, expected, value);
2485 - value_errors++;
2486 - thread_info->errors++;
2487 - }
2488 - }
2489 - if (end_time != time_now) {
2490 - if (!thread_info->delete_old_data) { /* data validation only when we don't delete */
2491 - if(!time_errors)
2492 - fprintf(stderr,
2493 - " DB-engine stresstest %s/%s: at %lu secs, found timestamp %lu ### E R R O R ###\n",
2494 - rrdset_name(st), rrddim_name(rd), (unsigned long) time_now, (unsigned long) time_retrieved);
2495 - time_errors++;
2496 - thread_info->errors++;
2497 - }
2498 - }
2499 - }
2500 - storage_engine_query_finalize(&seqh);
2501 - } while(!thread_info->done);
2502 -
2503 - if(value_errors)
2504 - fprintf(stderr, "%zu value errors encountered\n", value_errors);
2505 -
2506 - if(time_errors)
2507 - fprintf(stderr, "%zu time errors encountered\n", time_errors);
2508 -}
2509 -
2510 -void dbengine_stress_test(unsigned TEST_DURATION_SEC, unsigned DSET_CHARTS, unsigned QUERY_THREADS,
2511 - unsigned RAMP_UP_SECONDS, unsigned PAGE_CACHE_MB, unsigned DISK_SPACE_MB)
2512 -{
2513 - fprintf(stderr, "%s() running...\n", __FUNCTION__ );
2514 - const unsigned DSET_DIMS = 128;
2515 - const uint64_t EXPECTED_COMPRESSION_RATIO = 20;
2516 - const unsigned HISTORY_SECONDS = 3600 * 24 * 365 * 50; /* 50 year of history */
2517 - RRDHOST *host = NULL;
2518 - struct dbengine_chart_thread **chart_threads;
2519 - struct dbengine_query_thread **query_threads;
2520 - unsigned i, j;
2521 - time_t time_start, test_duration;
2522 -
2523 - nd_log_limits_unlimited();
2524 -
2525 - if (!TEST_DURATION_SEC)
2526 - TEST_DURATION_SEC = 10;
2527 - if (!DSET_CHARTS)
2528 - DSET_CHARTS = 1;
2529 - if (!QUERY_THREADS)
2530 - QUERY_THREADS = 1;
2531 - if (PAGE_CACHE_MB < RRDENG_MIN_PAGE_CACHE_SIZE_MB)
2532 - PAGE_CACHE_MB = RRDENG_MIN_PAGE_CACHE_SIZE_MB;
2533 -
2534 - default_rrd_memory_mode = RRD_MEMORY_MODE_DBENGINE;
2535 - default_rrdeng_page_cache_mb = PAGE_CACHE_MB;
2536 - if (DISK_SPACE_MB) {
2537 - fprintf(stderr, "By setting disk space limit data are allowed to be deleted. "
2538 - "Data validation is turned off for this run.\n");
2539 - default_rrdeng_disk_quota_mb = DISK_SPACE_MB;
2540 - } else {
2541 - // Worst case for uncompressible data
2542 - default_rrdeng_disk_quota_mb =
2543 - (((uint64_t) DSET_DIMS * DSET_CHARTS) * sizeof(storage_number) * HISTORY_SECONDS) / (1024 * 1024);
2544 - default_rrdeng_disk_quota_mb -= default_rrdeng_disk_quota_mb * EXPECTED_COMPRESSION_RATIO / 100;
2545 - }
2546 -
2547 - fprintf(stderr, "Initializing localhost with hostname 'dbengine-stress-test'\n");
2548 -
2549 - (void)sql_init_meta_database(DB_CHECK_NONE, 1);
2550 - host = dbengine_rrdhost_find_or_create("dbengine-stress-test");
2551 - if (NULL == host)
2552 - return;
2553 -
2554 - chart_threads = mallocz(sizeof(*chart_threads) * DSET_CHARTS);
2555 - for (i = 0 ; i < DSET_CHARTS ; ++i) {
2556 - chart_threads[i] = mallocz(sizeof(*chart_threads[i]) + sizeof(RRDDIM *) * DSET_DIMS);
2557 - }
2558 - query_threads = mallocz(sizeof(*query_threads) * QUERY_THREADS);
2559 - for (i = 0 ; i < QUERY_THREADS ; ++i) {
2560 - query_threads[i] = mallocz(sizeof(*query_threads[i]) + sizeof(struct dbengine_chart_thread *) * DSET_CHARTS);
2561 - }
2562 - fprintf(stderr, "\nRunning DB-engine stress test, %u seconds writers ramp-up time,\n"
2563 - "%u seconds of concurrent readers and writers, %u writer threads, %u reader threads,\n"
2564 - "%u MiB of page cache.\n",
2565 - RAMP_UP_SECONDS, TEST_DURATION_SEC, DSET_CHARTS, QUERY_THREADS, PAGE_CACHE_MB);
2566 -
2567 - time_start = now_realtime_sec() + HISTORY_SECONDS; /* move history to the future */
2568 - for (i = 0 ; i < DSET_CHARTS ; ++i) {
2569 - chart_threads[i]->host = host;
2570 - chart_threads[i]->chartname = "random";
2571 - chart_threads[i]->dset_charts = DSET_CHARTS;
2572 - chart_threads[i]->chart_i = i;
2573 - chart_threads[i]->dset_dims = DSET_DIMS;
2574 - chart_threads[i]->history_seconds = HISTORY_SECONDS;
2575 - chart_threads[i]->time_present = time_start;
2576 - chart_threads[i]->time_max = 0;
2577 - chart_threads[i]->done = 0;
2578 - chart_threads[i]->errors = chart_threads[i]->stored_metrics_nr = 0;
2579 - completion_init(&chart_threads[i]->charts_initialized);
2580 - fatal_assert(0 == uv_thread_create(&chart_threads[i]->thread, generate_dbengine_chart, chart_threads[i]));
2581 - }
2582 - /* barrier so that subsequent queries can access valid chart data */
2583 - for (i = 0 ; i < DSET_CHARTS ; ++i) {
2584 - completion_wait_for(&chart_threads[i]->charts_initialized);
2585 - completion_destroy(&chart_threads[i]->charts_initialized);
2586 - }
2587 - sleep(RAMP_UP_SECONDS);
2588 - /* at this point data have already began being written to the database */
2589 - for (i = 0 ; i < QUERY_THREADS ; ++i) {
2590 - query_threads[i]->host = host;
2591 - query_threads[i]->chartname = "random";
2592 - query_threads[i]->dset_charts = DSET_CHARTS;
2593 - query_threads[i]->dset_dims = DSET_DIMS;
2594 - query_threads[i]->history_seconds = HISTORY_SECONDS;
2595 - query_threads[i]->time_present = time_start;
2596 - query_threads[i]->done = 0;
2597 - query_threads[i]->errors = query_threads[i]->queries_nr = query_threads[i]->queried_metrics_nr = 0;
2598 - for (j = 0 ; j < DSET_CHARTS ; ++j) {
2599 - query_threads[i]->chart_threads[j] = chart_threads[j];
2600 - }
2601 - query_threads[i]->delete_old_data = DISK_SPACE_MB ? 1 : 0;
2602 - fatal_assert(0 == uv_thread_create(&query_threads[i]->thread, query_dbengine_chart, query_threads[i]));
2603 - }
2604 - sleep(TEST_DURATION_SEC);
2605 - /* stop workload */
2606 - for (i = 0 ; i < DSET_CHARTS ; ++i) {
2607 - chart_threads[i]->done = 1;
2608 - }
2609 - for (i = 0 ; i < QUERY_THREADS ; ++i) {
2610 - query_threads[i]->done = 1;
2611 - }
2612 - for (i = 0 ; i < DSET_CHARTS ; ++i) {
2613 - assert(0 == uv_thread_join(&chart_threads[i]->thread));
2614 - }
2615 - for (i = 0 ; i < QUERY_THREADS ; ++i) {
2616 - assert(0 == uv_thread_join(&query_threads[i]->thread));
2617 - }
2618 - test_duration = now_realtime_sec() - (time_start - HISTORY_SECONDS);
2619 - if (!test_duration)
2620 - test_duration = 1;
2621 - fprintf(stderr, "\nDB-engine stress test finished in %lld seconds.\n", (long long)test_duration);
2622 - unsigned long stored_metrics_nr = 0;
2623 - for (i = 0 ; i < DSET_CHARTS ; ++i) {
2624 - stored_metrics_nr += chart_threads[i]->stored_metrics_nr;
2625 - }
2626 - unsigned long queried_metrics_nr = 0;
2627 - for (i = 0 ; i < QUERY_THREADS ; ++i) {
2628 - queried_metrics_nr += query_threads[i]->queried_metrics_nr;
2629 - }
2630 - fprintf(stderr, "%u metrics were stored (dataset size of %lu MiB) in %u charts by 1 writer thread per chart.\n",
2631 - DSET_CHARTS * DSET_DIMS, stored_metrics_nr * sizeof(storage_number) / (1024 * 1024), DSET_CHARTS);
2632 - fprintf(stderr, "Metrics were being generated per 1 emulated second and time was accelerated.\n");
2633 - fprintf(stderr, "%lu metric data points were queried by %u reader threads.\n", queried_metrics_nr, QUERY_THREADS);
2634 - fprintf(stderr, "Query starting time is randomly chosen from the beginning of the time-series up to the time of\n"
2635 - "the latest data point, and ending time from 1 second up to 1 hour after the starting time.\n");
2636 - fprintf(stderr, "Performance is %lld written data points/sec and %lld read data points/sec.\n",
2637 - (long long)(stored_metrics_nr / test_duration), (long long)(queried_metrics_nr / test_duration));
2638 -
2639 - for (i = 0 ; i < DSET_CHARTS ; ++i) {
2640 - freez(chart_threads[i]);
2641 - }
2642 - freez(chart_threads);
2643 - for (i = 0 ; i < QUERY_THREADS ; ++i) {
2644 - freez(query_threads[i]);
2645 - }
2646 - freez(query_threads);
2647 - rrd_wrlock();
2648 - rrdeng_prepare_exit((struct rrdengine_instance *)host->db[0].si);
2649 - rrdeng_exit((struct rrdengine_instance *)host->db[0].si);
2650 - rrdeng_enq_cmd(NULL, RRDENG_OPCODE_SHUTDOWN_EVLOOP, NULL, NULL, STORAGE_PRIORITY_BEST_EFFORT, NULL, NULL);
2651 - rrd_unlock();
2652 -}
2653 -
2654 -#endif
src/database/engine/dbengine-stresstest.c new
+456
@@ -0,0 +1,456 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +#include "../../daemon/common.h"
4 +
5 +#ifdef ENABLE_DBENGINE
6 +
7 +static RRDHOST *dbengine_rrdhost_find_or_create(char *name) {
8 + /* We don't want to drop metrics when generating load,
9 + * we prefer to block data generation itself */
10 +
11 + return rrdhost_find_or_create(
12 + name,
13 + name,
14 + name,
15 + os_type,
16 + netdata_configured_timezone,
17 + netdata_configured_abbrev_timezone,
18 + netdata_configured_utc_offset,
19 + program_name,
20 + program_version,
21 + default_rrd_update_every,
22 + default_rrd_history_entries,
23 + RRD_MEMORY_MODE_DBENGINE,
24 + health_plugin_enabled(),
25 + default_rrdpush_enabled,
26 + default_rrdpush_destination,
27 + default_rrdpush_api_key,
28 + default_rrdpush_send_charts_matching,
29 + default_rrdpush_enable_replication,
30 + default_rrdpush_seconds_to_replicate,
31 + default_rrdpush_replication_step,
32 + NULL,
33 + 0
34 + );
35 +}
36 +
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;
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;
47 +}
48 +
49 +struct dbengine_chart_thread {
50 + uv_thread_t thread;
51 + RRDHOST *host;
52 + char *chartname; /* Will be prefixed by type, e.g. "example_local1.", "example_local2." etc */
53 + unsigned dset_charts; /* number of charts */
54 + unsigned dset_dims; /* dimensions per chart */
55 + unsigned chart_i; /* current chart offset */
56 + time_t time_present; /* current virtual time of the benchmark */
57 + volatile time_t time_max; /* latest timestamp of stored values */
58 + unsigned history_seconds; /* how far back in the past to go */
59 +
60 + volatile long done; /* initialize to 0, set to 1 to stop thread */
61 + struct completion charts_initialized;
62 + unsigned long errors, stored_metrics_nr; /* statistics */
63 +
64 + RRDSET *st;
65 + RRDDIM *rd[]; /* dset_dims elements */
66 +};
67 +
68 +collected_number generate_dbengine_chart_value(int chart_i, int dim_i, time_t time_current)
69 +{
70 + collected_number value;
71 +
72 + value = ((collected_number)time_current) * (chart_i + 1);
73 + value += ((collected_number)time_current) * (dim_i + 1);
74 + value %= 1024LLU;
75 +
76 + return value;
77 +}
78 +
79 +static void generate_dbengine_chart(void *arg)
80 +{
81 + fprintf(stderr, "%s() running...\n", __FUNCTION__ );
82 + struct dbengine_chart_thread *thread_info = (struct dbengine_chart_thread *)arg;
83 + RRDHOST *host = thread_info->host;
84 + char *chartname = thread_info->chartname;
85 + const unsigned DSET_DIMS = thread_info->dset_dims;
86 + unsigned history_seconds = thread_info->history_seconds;
87 + time_t time_present = thread_info->time_present;
88 +
89 + unsigned j, update_every = 1;
90 + RRDSET *st;
91 + RRDDIM *rd[DSET_DIMS];
92 + char name[RRD_ID_LENGTH_MAX + 1];
93 + time_t time_current;
94 +
95 + // create the chart
96 + snprintfz(name, RRD_ID_LENGTH_MAX, "example_local%u", thread_info->chart_i + 1);
97 + thread_info->st = st = rrdset_create(host, name, chartname, chartname, "example", NULL, chartname, chartname,
98 + chartname, NULL, 1, update_every, RRDSET_TYPE_LINE);
99 + for (j = 0 ; j < DSET_DIMS ; ++j) {
100 + snprintfz(name, RRD_ID_LENGTH_MAX, "%s%u", chartname, j + 1);
101 +
102 + thread_info->rd[j] = rd[j] = rrddim_add(st, name, NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
103 + }
104 + completion_mark_complete(&thread_info->charts_initialized);
105 +
106 + // feed it with the test data
107 + time_current = time_present - history_seconds;
108 + for (j = 0 ; j < DSET_DIMS ; ++j) {
109 + rd[j]->collector.last_collected_time.tv_sec =
110 + st->last_collected_time.tv_sec = st->last_updated.tv_sec = time_current - update_every;
111 + rd[j]->collector.last_collected_time.tv_usec =
112 + st->last_collected_time.tv_usec = st->last_updated.tv_usec = 0;
113 + }
114 + for( ; !thread_info->done && time_current < time_present ; time_current += update_every) {
115 + st->usec_since_last_update = USEC_PER_SEC * update_every;
116 +
117 + for (j = 0; j < DSET_DIMS; ++j) {
118 + collected_number value;
119 +
120 + value = generate_dbengine_chart_value(thread_info->chart_i, j, time_current);
121 + rrddim_set_by_pointer_fake_time(rd[j], value, time_current);
122 + ++thread_info->stored_metrics_nr;
123 + }
124 + rrdset_done(st);
125 + thread_info->time_max = time_current;
126 + }
127 + for (j = 0; j < DSET_DIMS; ++j) {
128 + rrdeng_store_metric_finalize((rd[j])->tiers[0].sch);
129 + }
130 +}
131 +
132 +void generate_dbengine_dataset(unsigned history_seconds)
133 +{
134 + fprintf(stderr, "%s() running...\n", __FUNCTION__ );
135 + const int DSET_CHARTS = 16;
136 + const int DSET_DIMS = 128;
137 + const uint64_t EXPECTED_COMPRESSION_RATIO = 20;
138 + RRDHOST *host = NULL;
139 + struct dbengine_chart_thread **thread_info;
140 + int i;
141 + time_t time_present;
142 +
143 + default_rrd_memory_mode = RRD_MEMORY_MODE_DBENGINE;
144 + default_rrdeng_page_cache_mb = 128;
145 + // Worst case for uncompressible data
146 + default_rrdeng_disk_quota_mb = (((uint64_t)DSET_DIMS * DSET_CHARTS) * sizeof(storage_number) * history_seconds) /
147 + (1024 * 1024);
148 + default_rrdeng_disk_quota_mb -= default_rrdeng_disk_quota_mb * EXPECTED_COMPRESSION_RATIO / 100;
149 +
150 + nd_log_limits_unlimited();
151 + fprintf(stderr, "Initializing localhost with hostname 'dbengine-dataset'");
152 +
153 + host = dbengine_rrdhost_find_or_create("dbengine-dataset");
154 + if (NULL == host)
155 + return;
156 +
157 + thread_info = mallocz(sizeof(*thread_info) * DSET_CHARTS);
158 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
159 + thread_info[i] = mallocz(sizeof(*thread_info[i]) + sizeof(RRDDIM *) * DSET_DIMS);
160 + }
161 + fprintf(stderr, "\nRunning DB-engine workload generator\n");
162 +
163 + time_present = now_realtime_sec();
164 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
165 + thread_info[i]->host = host;
166 + thread_info[i]->chartname = "random";
167 + thread_info[i]->dset_charts = DSET_CHARTS;
168 + thread_info[i]->chart_i = i;
169 + thread_info[i]->dset_dims = DSET_DIMS;
170 + thread_info[i]->history_seconds = history_seconds;
171 + thread_info[i]->time_present = time_present;
172 + thread_info[i]->time_max = 0;
173 + thread_info[i]->done = 0;
174 + completion_init(&thread_info[i]->charts_initialized);
175 + fatal_assert(0 == uv_thread_create(&thread_info[i]->thread, generate_dbengine_chart, thread_info[i]));
176 + completion_wait_for(&thread_info[i]->charts_initialized);
177 + completion_destroy(&thread_info[i]->charts_initialized);
178 + }
179 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
180 + fatal_assert(0 == uv_thread_join(&thread_info[i]->thread));
181 + }
182 +
183 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
184 + freez(thread_info[i]);
185 + }
186 + freez(thread_info);
187 + rrd_wrlock();
188 + rrdhost_free___while_having_rrd_wrlock(localhost, true);
189 + rrd_unlock();
190 +}
191 +
192 +struct dbengine_query_thread {
193 + uv_thread_t thread;
194 + RRDHOST *host;
195 + char *chartname; /* Will be prefixed by type, e.g. "example_local1.", "example_local2." etc */
196 + unsigned dset_charts; /* number of charts */
197 + unsigned dset_dims; /* dimensions per chart */
198 + time_t time_present; /* current virtual time of the benchmark */
199 + unsigned history_seconds; /* how far back in the past to go */
200 + volatile long done; /* initialize to 0, set to 1 to stop thread */
201 + unsigned long errors, queries_nr, queried_metrics_nr; /* statistics */
202 + uint8_t delete_old_data; /* if non zero then data are deleted when disk space is exhausted */
203 +
204 + struct dbengine_chart_thread *chart_threads[]; /* dset_charts elements */
205 +};
206 +
207 +static void query_dbengine_chart(void *arg)
208 +{
209 + fprintf(stderr, "%s() running...\n", __FUNCTION__ );
210 + struct dbengine_query_thread *thread_info = (struct dbengine_query_thread *)arg;
211 + const int DSET_CHARTS = thread_info->dset_charts;
212 + const int DSET_DIMS = thread_info->dset_dims;
213 + time_t time_after, time_before, time_min, time_approx_min, time_max, duration;
214 + int i, j, update_every = 1;
215 + RRDSET *st;
216 + RRDDIM *rd;
217 + uint8_t same;
218 + time_t time_now, time_retrieved, end_time;
219 + collected_number generatedv;
220 + NETDATA_DOUBLE value, expected;
221 + struct storage_engine_query_handle seqh;
222 + size_t value_errors = 0, time_errors = 0;
223 +
224 + do {
225 + // pick a chart and dimension
226 + i = random() % DSET_CHARTS;
227 + st = thread_info->chart_threads[i]->st;
228 + j = random() % DSET_DIMS;
229 + rd = thread_info->chart_threads[i]->rd[j];
230 +
231 + time_min = thread_info->time_present - thread_info->history_seconds + 1;
232 + time_max = thread_info->chart_threads[i]->time_max;
233 +
234 + if (thread_info->delete_old_data) {
235 + /* A time window of twice the disk space is sufficient for compression space savings of up to 50% */
236 + time_approx_min = time_max - (default_rrdeng_disk_quota_mb * 2 * 1024 * 1024) /
237 + (((uint64_t) DSET_DIMS * DSET_CHARTS) * sizeof(storage_number));
238 + time_min = MAX(time_min, time_approx_min);
239 + }
240 + if (!time_max) {
241 + time_before = time_after = time_min;
242 + } else {
243 + time_after = time_min + random() % (MAX(time_max - time_min, 1));
244 + duration = random() % 3600;
245 + time_before = MIN(time_after + duration, time_max); /* up to 1 hour queries */
246 + }
247 +
248 + storage_engine_query_init(rd->tiers[0].seb, rd->tiers[0].smh, &seqh, time_after, time_before, STORAGE_PRIORITY_NORMAL);
249 + ++thread_info->queries_nr;
250 + for (time_now = time_after ; time_now <= time_before ; time_now += update_every) {
251 + generatedv = generate_dbengine_chart_value(i, j, time_now);
252 + expected = unpack_storage_number(pack_storage_number((NETDATA_DOUBLE) generatedv, SN_DEFAULT_FLAGS));
253 +
254 + if (unlikely(storage_engine_query_is_finished(&seqh))) {
255 + if (!thread_info->delete_old_data) { /* data validation only when we don't delete */
256 + fprintf(stderr, " DB-engine stresstest %s/%s: at %lu secs, expecting value " NETDATA_DOUBLE_FORMAT
257 + ", found data gap, ### ERROR 12 ###\n",
258 + rrdset_name(st), rrddim_name(rd), (unsigned long) time_now, expected);
259 + ++thread_info->errors;
260 + }
261 + break;
262 + }
263 +
264 + STORAGE_POINT sp = storage_engine_query_next_metric(&seqh);
265 + value = sp.sum;
266 + time_retrieved = sp.start_time_s;
267 + end_time = sp.end_time_s;
268 +
269 + if (!netdata_double_isnumber(value)) {
270 + if (!thread_info->delete_old_data) { /* data validation only when we don't delete */
271 + fprintf(stderr, " DB-engine stresstest %s/%s: at %lu secs, expecting value " NETDATA_DOUBLE_FORMAT
272 + ", found data gap, ### ERROR 13 ###\n",
273 + rrdset_name(st), rrddim_name(rd), (unsigned long) time_now, expected);
274 + ++thread_info->errors;
275 + }
276 + break;
277 + }
278 + ++thread_info->queried_metrics_nr;
279 +
280 + same = (roundndd(value) == roundndd(expected)) ? 1 : 0;
281 + if (!same) {
282 + if (!thread_info->delete_old_data) { /* data validation only when we don't delete */
283 + if(!value_errors)
284 + fprintf(stderr, " DB-engine stresstest %s/%s: at %lu secs, expecting value " NETDATA_DOUBLE_FORMAT
285 + ", found " NETDATA_DOUBLE_FORMAT ", ### ERROR 14 ###\n",
286 + rrdset_name(st), rrddim_name(rd), (unsigned long) time_now, expected, value);
287 + value_errors++;
288 + thread_info->errors++;
289 + }
290 + }
291 + if (end_time != time_now) {
292 + if (!thread_info->delete_old_data) { /* data validation only when we don't delete */
293 + if(!time_errors)
294 + fprintf(stderr,
295 + " DB-engine stresstest %s/%s: at %lu secs, found timestamp %lu ### ERROR 15 ###\n",
296 + rrdset_name(st), rrddim_name(rd), (unsigned long) time_now, (unsigned long) time_retrieved);
297 + time_errors++;
298 + thread_info->errors++;
299 + }
300 + }
301 + }
302 + storage_engine_query_finalize(&seqh);
303 + } while(!thread_info->done);
304 +
305 + if(value_errors)
306 + fprintf(stderr, "%zu value errors encountered\n", value_errors);
307 +
308 + if(time_errors)
309 + fprintf(stderr, "%zu time errors encountered\n", time_errors);
310 +}
311 +
312 +void dbengine_stress_test(unsigned TEST_DURATION_SEC, unsigned DSET_CHARTS, unsigned QUERY_THREADS,
313 + unsigned RAMP_UP_SECONDS, unsigned PAGE_CACHE_MB, unsigned DISK_SPACE_MB)
314 +{
315 + fprintf(stderr, "%s() running...\n", __FUNCTION__ );
316 + const unsigned DSET_DIMS = 128;
317 + const uint64_t EXPECTED_COMPRESSION_RATIO = 20;
318 + const unsigned HISTORY_SECONDS = 3600 * 24 * 365 * 50; /* 50 year of history */
319 + RRDHOST *host = NULL;
320 + struct dbengine_chart_thread **chart_threads;
321 + struct dbengine_query_thread **query_threads;
322 + unsigned i, j;
323 + time_t time_start, test_duration;
324 +
325 + nd_log_limits_unlimited();
326 +
327 + if (!TEST_DURATION_SEC)
328 + TEST_DURATION_SEC = 10;
329 + if (!DSET_CHARTS)
330 + DSET_CHARTS = 1;
331 + if (!QUERY_THREADS)
332 + QUERY_THREADS = 1;
333 + if (PAGE_CACHE_MB < RRDENG_MIN_PAGE_CACHE_SIZE_MB)
334 + PAGE_CACHE_MB = RRDENG_MIN_PAGE_CACHE_SIZE_MB;
335 +
336 + default_rrd_memory_mode = RRD_MEMORY_MODE_DBENGINE;
337 + default_rrdeng_page_cache_mb = PAGE_CACHE_MB;
338 + if (DISK_SPACE_MB) {
339 + fprintf(stderr, "By setting disk space limit data are allowed to be deleted. "
340 + "Data validation is turned off for this run.\n");
341 + default_rrdeng_disk_quota_mb = DISK_SPACE_MB;
342 + } else {
343 + // Worst case for uncompressible data
344 + default_rrdeng_disk_quota_mb =
345 + (((uint64_t) DSET_DIMS * DSET_CHARTS) * sizeof(storage_number) * HISTORY_SECONDS) / (1024 * 1024);
346 + default_rrdeng_disk_quota_mb -= default_rrdeng_disk_quota_mb * EXPECTED_COMPRESSION_RATIO / 100;
347 + }
348 +
349 + fprintf(stderr, "Initializing localhost with hostname 'dbengine-stress-test'\n");
350 +
351 + (void)sql_init_meta_database(DB_CHECK_NONE, 1);
352 + host = dbengine_rrdhost_find_or_create("dbengine-stress-test");
353 + if (NULL == host)
354 + return;
355 +
356 + chart_threads = mallocz(sizeof(*chart_threads) * DSET_CHARTS);
357 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
358 + chart_threads[i] = mallocz(sizeof(*chart_threads[i]) + sizeof(RRDDIM *) * DSET_DIMS);
359 + }
360 + query_threads = mallocz(sizeof(*query_threads) * QUERY_THREADS);
361 + for (i = 0 ; i < QUERY_THREADS ; ++i) {
362 + query_threads[i] = mallocz(sizeof(*query_threads[i]) + sizeof(struct dbengine_chart_thread *) * DSET_CHARTS);
363 + }
364 + fprintf(stderr, "\nRunning DB-engine stress test, %u seconds writers ramp-up time,\n"
365 + "%u seconds of concurrent readers and writers, %u writer threads, %u reader threads,\n"
366 + "%u MiB of page cache.\n",
367 + RAMP_UP_SECONDS, TEST_DURATION_SEC, DSET_CHARTS, QUERY_THREADS, PAGE_CACHE_MB);
368 +
369 + time_start = now_realtime_sec() + HISTORY_SECONDS; /* move history to the future */
370 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
371 + chart_threads[i]->host = host;
372 + chart_threads[i]->chartname = "random";
373 + chart_threads[i]->dset_charts = DSET_CHARTS;
374 + chart_threads[i]->chart_i = i;
375 + chart_threads[i]->dset_dims = DSET_DIMS;
376 + chart_threads[i]->history_seconds = HISTORY_SECONDS;
377 + chart_threads[i]->time_present = time_start;
378 + chart_threads[i]->time_max = 0;
379 + chart_threads[i]->done = 0;
380 + chart_threads[i]->errors = chart_threads[i]->stored_metrics_nr = 0;
381 + completion_init(&chart_threads[i]->charts_initialized);
382 + fatal_assert(0 == uv_thread_create(&chart_threads[i]->thread, generate_dbengine_chart, chart_threads[i]));
383 + }
384 + /* barrier so that subsequent queries can access valid chart data */
385 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
386 + completion_wait_for(&chart_threads[i]->charts_initialized);
387 + completion_destroy(&chart_threads[i]->charts_initialized);
388 + }
389 + sleep(RAMP_UP_SECONDS);
390 + /* at this point data have already began being written to the database */
391 + for (i = 0 ; i < QUERY_THREADS ; ++i) {
392 + query_threads[i]->host = host;
393 + query_threads[i]->chartname = "random";
394 + query_threads[i]->dset_charts = DSET_CHARTS;
395 + query_threads[i]->dset_dims = DSET_DIMS;
396 + query_threads[i]->history_seconds = HISTORY_SECONDS;
397 + query_threads[i]->time_present = time_start;
398 + query_threads[i]->done = 0;
399 + query_threads[i]->errors = query_threads[i]->queries_nr = query_threads[i]->queried_metrics_nr = 0;
400 + for (j = 0 ; j < DSET_CHARTS ; ++j) {
401 + query_threads[i]->chart_threads[j] = chart_threads[j];
402 + }
403 + query_threads[i]->delete_old_data = DISK_SPACE_MB ? 1 : 0;
404 + fatal_assert(0 == uv_thread_create(&query_threads[i]->thread, query_dbengine_chart, query_threads[i]));
405 + }
406 + sleep(TEST_DURATION_SEC);
407 + /* stop workload */
408 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
409 + chart_threads[i]->done = 1;
410 + }
411 + for (i = 0 ; i < QUERY_THREADS ; ++i) {
412 + query_threads[i]->done = 1;
413 + }
414 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
415 + assert(0 == uv_thread_join(&chart_threads[i]->thread));
416 + }
417 + for (i = 0 ; i < QUERY_THREADS ; ++i) {
418 + assert(0 == uv_thread_join(&query_threads[i]->thread));
419 + }
420 + test_duration = now_realtime_sec() - (time_start - HISTORY_SECONDS);
421 + if (!test_duration)
422 + test_duration = 1;
423 + fprintf(stderr, "\nDB-engine stress test finished in %lld seconds.\n", (long long)test_duration);
424 + unsigned long stored_metrics_nr = 0;
425 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
426 + stored_metrics_nr += chart_threads[i]->stored_metrics_nr;
427 + }
428 + unsigned long queried_metrics_nr = 0;
429 + for (i = 0 ; i < QUERY_THREADS ; ++i) {
430 + queried_metrics_nr += query_threads[i]->queried_metrics_nr;
431 + }
432 + fprintf(stderr, "%u metrics were stored (dataset size of %lu MiB) in %u charts by 1 writer thread per chart.\n",
433 + DSET_CHARTS * DSET_DIMS, stored_metrics_nr * sizeof(storage_number) / (1024 * 1024), DSET_CHARTS);
434 + fprintf(stderr, "Metrics were being generated per 1 emulated second and time was accelerated.\n");
435 + fprintf(stderr, "%lu metric data points were queried by %u reader threads.\n", queried_metrics_nr, QUERY_THREADS);
436 + fprintf(stderr, "Query starting time is randomly chosen from the beginning of the time-series up to the time of\n"
437 + "the latest data point, and ending time from 1 second up to 1 hour after the starting time.\n");
438 + fprintf(stderr, "Performance is %lld written data points/sec and %lld read data points/sec.\n",
439 + (long long)(stored_metrics_nr / test_duration), (long long)(queried_metrics_nr / test_duration));
440 +
441 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
442 + freez(chart_threads[i]);
443 + }
444 + freez(chart_threads);
445 + for (i = 0 ; i < QUERY_THREADS ; ++i) {
446 + freez(query_threads[i]);
447 + }
448 + freez(query_threads);
449 + rrd_wrlock();
450 + rrdeng_prepare_exit((struct rrdengine_instance *)host->db[0].si);
451 + rrdeng_exit((struct rrdengine_instance *)host->db[0].si);
452 + rrdeng_enq_cmd(NULL, RRDENG_OPCODE_SHUTDOWN_EVLOOP, NULL, NULL, STORAGE_PRIORITY_BEST_EFFORT, NULL, NULL);
453 + rrd_unlock();
454 +}
455 +
456 +#endif
\ No newline at end of file
src/database/engine/dbengine-unittest.c new
+419
@@ -0,0 +1,419 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +#include "../../daemon/common.h"
4 +
5 +#ifdef ENABLE_DBENGINE
6 +
7 +#define CHARTS 64
8 +#define DIMS 8 // CHARTS * DIMS dimensions
9 +#define REGIONS 11
10 +#define POINTS_PER_REGION 4096
11 +static const int REGION_UPDATE_EVERY[REGIONS] = {1, 15, 3, 20, 2, 6, 30, 12, 5, 4, 10};
12 +
13 +#define START_TIMESTAMP MAX(2 * API_RELATIVE_TIME_MAX, 200000000)
14 +
15 +static time_t region_start_time(time_t previous_region_end_time, time_t update_every) {
16 + // leave a small gap between regions
17 + // but keep them close together, so that cross-region queries will be fast
18 +
19 + time_t rc = previous_region_end_time + update_every;
20 + rc += update_every - (rc % update_every);
21 + rc += update_every;
22 + return rc;
23 +}
24 +
25 +static inline collected_number point_value_get(size_t region, size_t chart, size_t dim, size_t point) {
26 + // calculate the value to be stored for each point in the database
27 +
28 + collected_number r = (collected_number)region;
29 + collected_number c = (collected_number)chart;
30 + collected_number d = (collected_number)dim;
31 + collected_number p = (collected_number)point;
32 +
33 + return (r * CHARTS * DIMS * POINTS_PER_REGION +
34 + c * DIMS * POINTS_PER_REGION +
35 + d * POINTS_PER_REGION +
36 + p) % 10000000;
37 +}
38 +
39 +static inline void storage_point_check(size_t region, size_t chart, size_t dim, size_t point, time_t now, time_t update_every, STORAGE_POINT sp, size_t *value_errors, size_t *time_errors, size_t *update_every_errors) {
40 + // check the supplied STORAGE_POINT retrieved from the database
41 + // against the computed timestamp, update_every and expected value
42 +
43 + if(storage_point_is_gap(sp)) sp.min = sp.max = sp.sum = NAN;
44 +
45 + collected_number expected = point_value_get(region, chart, dim, point);
46 +
47 + if(roundndd(expected) != roundndd(sp.sum)) {
48 + if(*value_errors < DIMS * 2) {
49 + fprintf(stderr, " >>> DBENGINE: VALUE DOES NOT MATCH: "
50 + "region %zu, chart %zu, dimension %zu, point %zu, time %ld: "
51 + "expected %lld, found %f\n",
52 + region, chart, dim, point, now, expected, sp.sum);
53 + }
54 +
55 + (*value_errors)++;
56 + }
57 +
58 + if(sp.start_time_s > now || sp.end_time_s < now) {
59 + if(*time_errors < DIMS * 2) {
60 + fprintf(stderr, " >>> DBENGINE: TIMESTAMP DOES NOT MATCH: "
61 + "region %zu, chart %zu, dimension %zu, point %zu, timestamp %ld: "
62 + "expected %ld, found %ld - %ld\n",
63 + region, chart, dim, point, now, now, sp.start_time_s, sp.end_time_s);
64 + }
65 +
66 + (*time_errors)++;
67 + }
68 +
69 + if(update_every != sp.end_time_s - sp.start_time_s) {
70 + if(*update_every_errors < DIMS * 2) {
71 + fprintf(stderr, " >>> DBENGINE: UPDATE EVERY DOES NOT MATCH: "
72 + "region %zu, chart %zu, dimension %zu, point %zu, timestamp %ld: "
73 + "expected %ld, found %ld\n",
74 + region, chart, dim, point, now, update_every, sp.end_time_s - sp.start_time_s);
75 + }
76 +
77 + (*update_every_errors)++;
78 + }
79 +}
80 +
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;
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;
91 +}
92 +
93 +static RRDHOST *dbengine_rrdhost_find_or_create(char *name) {
94 + /* We don't want to drop metrics when generating load,
95 + * we prefer to block data generation itself */
96 +
97 + return rrdhost_find_or_create(
98 + name,
99 + name,
100 + name,
101 + os_type,
102 + netdata_configured_timezone,
103 + netdata_configured_abbrev_timezone,
104 + netdata_configured_utc_offset,
105 + program_name,
106 + program_version,
107 + default_rrd_update_every,
108 + default_rrd_history_entries,
109 + RRD_MEMORY_MODE_DBENGINE,
110 + health_plugin_enabled(),
111 + default_rrdpush_enabled,
112 + default_rrdpush_destination,
113 + default_rrdpush_api_key,
114 + default_rrdpush_send_charts_matching,
115 + default_rrdpush_enable_replication,
116 + default_rrdpush_seconds_to_replicate,
117 + default_rrdpush_replication_step,
118 + NULL,
119 + 0
120 + );
121 +}
122 +
123 +static void test_dbengine_create_charts(RRDHOST *host, RRDSET *st[CHARTS], RRDDIM *rd[CHARTS][DIMS],
124 + int update_every) {
125 + fprintf(stderr, "DBENGINE Creating Test Charts...\n");
126 +
127 + int i, j;
128 + char name[101];
129 +
130 + for (i = 0 ; i < CHARTS ; ++i) {
131 + snprintfz(name, sizeof(name) - 1, "dbengine-chart-%d", i);
132 +
133 + // create the chart
134 + st[i] = rrdset_create(host, "netdata", name, name, "netdata", NULL, "Unit Testing", "a value", "unittest",
135 + NULL, 1, update_every, RRDSET_TYPE_LINE);
136 + rrdset_flag_set(st[i], RRDSET_FLAG_DEBUG);
137 + rrdset_flag_set(st[i], RRDSET_FLAG_STORE_FIRST);
138 + for (j = 0 ; j < DIMS ; ++j) {
139 + snprintfz(name, sizeof(name) - 1, "dim-%d", j);
140 +
141 + rd[i][j] = rrddim_add(st[i], name, NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
142 + }
143 + }
144 +
145 + // Initialize DB with the very first entries
146 + for (i = 0 ; i < CHARTS ; ++i) {
147 + for (j = 0 ; j < DIMS ; ++j) {
148 + rd[i][j]->collector.last_collected_time.tv_sec =
149 + st[i]->last_collected_time.tv_sec = st[i]->last_updated.tv_sec = START_TIMESTAMP - 1;
150 + rd[i][j]->collector.last_collected_time.tv_usec =
151 + st[i]->last_collected_time.tv_usec = st[i]->last_updated.tv_usec = 0;
152 + }
153 + }
154 + for (i = 0 ; i < CHARTS ; ++i) {
155 + st[i]->usec_since_last_update = USEC_PER_SEC;
156 +
157 + for (j = 0; j < DIMS; ++j) {
158 + rrddim_set_by_pointer_fake_time(rd[i][j], 69, START_TIMESTAMP); // set first value to 69
159 + }
160 +
161 + struct timeval now;
162 + now_realtime_timeval(&now);
163 + rrdset_timed_done(st[i], now, false);
164 + }
165 + // Flush pages for subsequent real values
166 + for (i = 0 ; i < CHARTS ; ++i) {
167 + for (j = 0; j < DIMS; ++j) {
168 + rrdeng_store_metric_flush_current_page((rd[i][j])->tiers[0].sch);
169 + }
170 + }
171 +}
172 +
173 +static time_t test_dbengine_create_metrics(
174 + RRDSET *st[CHARTS],
175 + RRDDIM *rd[CHARTS][DIMS],
176 + size_t current_region,
177 + time_t time_start) {
178 +
179 + time_t update_every = REGION_UPDATE_EVERY[current_region];
180 + fprintf(stderr, "DBENGINE Single Region Write to "
181 + "region %zu, from %ld to %ld, with update every %ld...\n",
182 + current_region, time_start, time_start + POINTS_PER_REGION * update_every, update_every);
183 +
184 + // for the database to save the metrics at the right time, we need to set
185 + // the last data collection time to be just before the first data collection.
186 + time_t time_now = time_start;
187 + for (size_t c = 0 ; c < CHARTS ; ++c) {
188 + for (size_t d = 0 ; d < DIMS ; ++d) {
189 + storage_engine_store_change_collection_frequency(rd[c][d]->tiers[0].sch, (int)update_every);
190 +
191 + // setting these timestamps, to the data collection time, prevents interpolation
192 + // during data collection, so that our value will be written as-is to the
193 + // database.
194 +
195 + rd[c][d]->collector.last_collected_time.tv_sec =
196 + st[c]->last_collected_time.tv_sec = st[c]->last_updated.tv_sec = time_now;
197 +
198 + rd[c][d]->collector.last_collected_time.tv_usec =
199 + st[c]->last_collected_time.tv_usec = st[c]->last_updated.tv_usec = 0;
200 + }
201 + }
202 +
203 + // set the samples to the database
204 + for (size_t p = 0; p < POINTS_PER_REGION ; ++p) {
205 + for (size_t c = 0 ; c < CHARTS ; ++c) {
206 + st[c]->usec_since_last_update = USEC_PER_SEC * update_every;
207 +
208 + for (size_t d = 0; d < DIMS; ++d)
209 + rrddim_set_by_pointer_fake_time(rd[c][d], point_value_get(current_region, c, d, p), time_now);
210 +
211 + rrdset_timed_done(st[c], (struct timeval){ .tv_sec = time_now, .tv_usec = 0 }, false);
212 + }
213 +
214 + time_now += update_every;
215 + }
216 +
217 + return time_now;
218 +}
219 +
220 +// Checks the metric data for the given region, returns number of errors
221 +static size_t test_dbengine_check_metrics(
222 + RRDSET *st[CHARTS] __maybe_unused,
223 + RRDDIM *rd[CHARTS][DIMS],
224 + size_t current_region,
225 + time_t time_start,
226 + time_t time_end) {
227 +
228 + time_t update_every = REGION_UPDATE_EVERY[current_region];
229 + fprintf(stderr, "DBENGINE Single Region Read from "
230 + "region %zu, from %ld to %ld, with update every %ld...\n",
231 + current_region, time_start, time_end, update_every);
232 +
233 + // initialize all queries
234 + struct storage_engine_query_handle handles[CHARTS * DIMS] = { 0 };
235 + for (size_t c = 0 ; c < CHARTS ; ++c) {
236 + for (size_t d = 0; d < DIMS; ++d) {
237 + storage_engine_query_init(rd[c][d]->tiers[0].seb,
238 + rd[c][d]->tiers[0].smh,
239 + &handles[c * DIMS + d],
240 + time_start,
241 + time_end,
242 + STORAGE_PRIORITY_NORMAL);
243 + }
244 + }
245 +
246 + // check the stored samples
247 + size_t value_errors = 0, time_errors = 0, update_every_errors = 0;
248 + time_t time_now = time_start;
249 + for(size_t p = 0; p < POINTS_PER_REGION ;p++) {
250 + for (size_t c = 0 ; c < CHARTS ; ++c) {
251 + for (size_t d = 0; d < DIMS; ++d) {
252 + STORAGE_POINT sp = storage_engine_query_next_metric(&handles[c * DIMS + d]);
253 + storage_point_check(current_region, c, d, p, time_now, update_every, sp,
254 + &value_errors, &time_errors, &update_every_errors);
255 + }
256 + }
257 +
258 + time_now += update_every;
259 + }
260 +
261 + // finalize the queries
262 + for (size_t c = 0 ; c < CHARTS ; ++c) {
263 + for (size_t d = 0; d < DIMS; ++d) {
264 + storage_engine_query_finalize(&handles[c * DIMS + d]);
265 + }
266 + }
267 +
268 + if(value_errors)
269 + fprintf(stderr, "%zu value errors encountered (out of %d checks)\n", value_errors, POINTS_PER_REGION * CHARTS * DIMS);
270 +
271 + if(time_errors)
272 + fprintf(stderr, "%zu time errors encountered (out of %d checks)\n", time_errors, POINTS_PER_REGION * CHARTS * DIMS);
273 +
274 + if(update_every_errors)
275 + fprintf(stderr, "%zu update every errors encountered (out of %d checks)\n", update_every_errors, POINTS_PER_REGION * CHARTS * DIMS);
276 +
277 + return value_errors + time_errors + update_every_errors;
278 +}
279 +
280 +static size_t dbengine_test_rrdr_single_region(
281 + RRDSET *st[CHARTS],
282 + RRDDIM *rd[CHARTS][DIMS],
283 + size_t current_region,
284 + time_t time_start,
285 + time_t time_end) {
286 +
287 + time_t update_every = REGION_UPDATE_EVERY[current_region];
288 + fprintf(stderr, "RRDR Single Region Test on "
289 + "region %zu, start time %lld, end time %lld, update every %ld, on %d dimensions...\n",
290 + current_region, (long long)time_start, (long long)time_end, update_every, CHARTS * DIMS);
291 +
292 + size_t errors = 0, value_errors = 0, time_errors = 0, update_every_errors = 0;
293 + long points = (time_end - time_start) / update_every;
294 + for(size_t c = 0; c < CHARTS ;c++) {
295 + ONEWAYALLOC *owa = onewayalloc_create(0);
296 + RRDR *r = rrd2rrdr_legacy(owa, st[c], points, time_start, time_end,
297 + RRDR_GROUPING_AVERAGE, 0, RRDR_OPTION_NATURAL_POINTS,
298 + NULL, NULL, 0, 0,
299 + QUERY_SOURCE_UNITTEST, STORAGE_PRIORITY_NORMAL);
300 + if (!r) {
301 + fprintf(stderr, " >>> DBENGINE: %s: empty RRDR on region %zu\n", rrdset_name(st[c]), current_region);
302 + onewayalloc_destroy(owa);
303 + errors++;
304 + continue;
305 + }
306 +
307 + if(r->internal.qt->request.st != st[c])
308 + fatal("queried wrong chart");
309 +
310 + if(rrdr_rows(r) != POINTS_PER_REGION)
311 + fatal("query returned wrong number of points (expected %d, got %zu)", POINTS_PER_REGION, rrdr_rows(r));
312 +
313 + time_t time_now = time_start;
314 + for (size_t p = 0; p < rrdr_rows(r); p++) {
315 + size_t d = 0;
316 + RRDDIM *dim;
317 + rrddim_foreach_read(dim, r->internal.qt->request.st) {
318 + if(unlikely(d >= r->d))
319 + fatal("got more dimensions (%zu) than expected (%zu)", d, r->d);
320 +
321 + if(rd[c][d] != dim)
322 + fatal("queried wrong dimension");
323 +
324 + RRDR_VALUE_FLAGS *co = &r->o[ p * r->d ];
325 + NETDATA_DOUBLE *cn = &r->v[ p * r->d ];
326 +
327 + STORAGE_POINT sp = STORAGE_POINT_UNSET;
328 + sp.min = sp.max = sp.sum = (co[d] & RRDR_VALUE_EMPTY) ? NAN :cn[d];
329 + sp.count = 1;
330 + sp.end_time_s = r->t[p];
331 + sp.start_time_s = sp.end_time_s - r->view.update_every;
332 +
333 + storage_point_check(current_region, c, d, p, time_now, update_every, sp, &value_errors, &time_errors, &update_every_errors);
334 + d++;
335 + }
336 + rrddim_foreach_done(dim);
337 + time_now += update_every;
338 + }
339 +
340 + rrdr_free(owa, r);
341 + onewayalloc_destroy(owa);
342 + }
343 +
344 + if(value_errors)
345 + fprintf(stderr, "%zu value errors encountered (out of %d checks)\n", value_errors, POINTS_PER_REGION * CHARTS * DIMS);
346 +
347 + if(time_errors)
348 + fprintf(stderr, "%zu time errors encountered (out of %d checks)\n", time_errors, POINTS_PER_REGION * CHARTS * DIMS);
349 +
350 + if(update_every_errors)
351 + fprintf(stderr, "%zu update every errors encountered (out of %d checks)\n", update_every_errors, POINTS_PER_REGION * CHARTS * DIMS);
352 +
353 + return errors + value_errors + time_errors + update_every_errors;
354 +}
355 +
356 +int test_dbengine(void) {
357 + // provide enough threads to dbengine
358 + setenv("UV_THREADPOOL_SIZE", "48", 1);
359 +
360 + size_t errors = 0, value_errors = 0, time_errors = 0;
361 +
362 + nd_log_limits_unlimited();
363 + fprintf(stderr, "\nRunning DB-engine test\n");
364 +
365 + default_rrd_memory_mode = RRD_MEMORY_MODE_DBENGINE;
366 + fprintf(stderr, "Initializing localhost with hostname 'unittest-dbengine'");
367 + RRDHOST *host = dbengine_rrdhost_find_or_create("unittest-dbengine");
368 + if(!host)
369 + fatal("Failed to initialize host");
370 +
371 + RRDSET *st[CHARTS] = { 0 };
372 + RRDDIM *rd[CHARTS][DIMS] = { 0 };
373 + time_t time_start[REGIONS] = { 0 }, time_end[REGIONS] = { 0 };
374 +
375 + // create the charts and dimensions we need
376 + test_dbengine_create_charts(host, st, rd, REGION_UPDATE_EVERY[0]);
377 +
378 + time_t now = START_TIMESTAMP;
379 + time_t update_every_old = REGION_UPDATE_EVERY[0];
380 + for(size_t current_region = 0; current_region < REGIONS ;current_region++) {
381 + time_t update_every = REGION_UPDATE_EVERY[current_region];
382 +
383 + if(update_every != update_every_old) {
384 + for (size_t c = 0 ; c < CHARTS ; ++c)
385 + rrdset_set_update_every_s(st[c], update_every);
386 + }
387 +
388 + time_start[current_region] = region_start_time(now, update_every);
389 + now = time_end[current_region] = test_dbengine_create_metrics(st,rd, current_region, time_start[current_region]);
390 +
391 + errors += test_dbengine_check_metrics(st, rd, current_region, time_start[current_region], time_end[current_region]);
392 + }
393 +
394 + // check everything again
395 + for(size_t current_region = 0; current_region < REGIONS ;current_region++)
396 + errors += test_dbengine_check_metrics(st, rd, current_region, time_start[current_region], time_end[current_region]);
397 +
398 + // check again in reverse order
399 + for(size_t current_region = 0; current_region < REGIONS ;current_region++) {
400 + size_t region = REGIONS - 1 - current_region;
401 + errors += test_dbengine_check_metrics(st, rd, region, time_start[region], time_end[region]);
402 + }
403 +
404 + // check all the regions using RRDR
405 + // this also checks the query planner and the query engine of Netdata
406 + for (size_t current_region = 0 ; current_region < REGIONS ; current_region++) {
407 + errors += dbengine_test_rrdr_single_region(st, rd, current_region, time_start[current_region], time_end[current_region]);
408 + }
409 +
410 + rrd_wrlock();
411 + rrdeng_prepare_exit((struct rrdengine_instance *)host->db[0].si);
412 + rrdeng_exit((struct rrdengine_instance *)host->db[0].si);
413 + rrdeng_enq_cmd(NULL, RRDENG_OPCODE_SHUTDOWN_EVLOOP, NULL, NULL, STORAGE_PRIORITY_BEST_EFFORT, NULL, NULL);
414 + rrd_unlock();
415 +
416 + return (int)(errors + value_errors + time_errors);
417 +}
418 +
419 +#endif