@cryptotaxi247 / netdata-1 / commits / 1d667b145

Fix dbengine consistency when a writer modifies a page concurrently with a reader querying its metrics (#6979)

Markos Fountoulakis committed Oct 2, 2019 at 07:12 UTC 1d667b145c16811639143afa377293d3ef8d62f0
2 files changed +54 -12
database/engine/pagecache.h
+38
@@ -183,4 +183,42 @@ extern void free_page_cache(struct rrdengine_instance *ctx);
183 extern void pg_cache_add_new_metric_time(struct pg_cache_page_index *page_index, struct rrdeng_page_descr *descr);
184 extern void pg_cache_update_metric_times(struct pg_cache_page_index *page_index);
185
186 +static inline void
187 + pg_cache_atomic_get_pg_info(struct rrdeng_page_descr *descr, usec_t *end_timep, uint32_t *page_lengthp)
188 +{
189 + usec_t end_time, old_end_time;
190 + uint32_t page_length;
191 +
192 + if (NULL == descr->extent) {
193 + /* this page is currently being modified, get consistent info locklessly */
194 + do {
195 + end_time = descr->end_time;
196 + __sync_synchronize();
197 + old_end_time = end_time;
198 + page_length = descr->page_length;
199 + __sync_synchronize();
200 + end_time = descr->end_time;
201 + __sync_synchronize();
202 + } while ((end_time != old_end_time || (end_time & 1) != 0));
203 +
204 + *end_timep = end_time;
205 + *page_lengthp = page_length;
206 + } else {
207 + *end_timep = descr->end_time;
208 + *page_lengthp = descr->page_length;
209 + }
210 +}
211 +
212 +/* The caller must hold a reference to the page and must have already set the new data */
213 +static inline void pg_cache_atomic_set_pg_info(struct rrdeng_page_descr *descr, usec_t end_time, uint32_t page_length)
214 +{
215 + assert(!(end_time & 1));
216 + __sync_synchronize();
217 + descr->end_time |= 1; /* mark start of uncertainty period by adding 1 microsecond */
218 + __sync_synchronize();
219 + descr->page_length = page_length;
220 + __sync_synchronize();
221 + descr->end_time = end_time; /* mark end of uncertainty period */
222 +}
223 +
224 #endif /* NETDATA_PAGECACHE_H */
database/engine/rrdengineapi.c
+16 -12
@@ -184,8 +184,8 @@ void rrdeng_store_metric_next(RRDDIM *rd, usec_t point_in_time, storage_number n
184 }
185 page = descr->pg_cache_descr->page;
186 page[descr->page_length / sizeof(number)] = number;
187 - descr->end_time = point_in_time;
188 - descr->page_length += sizeof(number);
187 + pg_cache_atomic_set_pg_info(descr, point_in_time, descr->page_length + sizeof(number));
188 +
189 if (perfect_page_alignment)
190 rd->rrdset->rrddim_page_alignment = descr->page_length;
191 if (unlikely(INVALID_TIME == descr->start_time)) {
@@ -470,7 +470,8 @@ storage_number rrdeng_load_metric_next(struct rrddim_query_handle *rrdimm_handle
470 struct rrdeng_page_descr *descr;
471 storage_number *page, ret;
472 unsigned position, entries;
473 - usec_t next_page_time, current_position_time;
473 + usec_t next_page_time, current_position_time, page_end_time;
474 + uint32_t page_length;
475
476 handle = &rrdimm_handle->rrdeng;
477 if (unlikely(INVALID_TIME == handle->next_page_time)) {
@@ -480,15 +481,17 @@ storage_number rrdeng_load_metric_next(struct rrddim_query_handle *rrdimm_handle
481 if (unlikely(NULL == (descr = handle->descr))) {
482 /* it's the first call */
483 next_page_time = handle->next_page_time * USEC_PER_SEC;
484 + } else {
485 + pg_cache_atomic_get_pg_info(descr, &page_end_time, &page_length);
486 }
487 position = handle->position + 1;
488
489 if (unlikely(NULL == descr ||
487 - position >= (descr->page_length / sizeof(storage_number)))) {
490 + position >= (page_length / sizeof(storage_number)))) {
491 /* We need to get a new page */
492 if (descr) {
493 /* Drop old page's reference */
491 - handle->next_page_time = (descr->end_time / USEC_PER_SEC) + 1;
494 + handle->next_page_time = (page_end_time / USEC_PER_SEC) + 1;
495 if (unlikely(handle->next_page_time > rrdimm_handle->end_time)) {
496 goto no_more_metrics;
497 }
@@ -508,26 +511,27 @@ storage_number rrdeng_load_metric_next(struct rrddim_query_handle *rrdimm_handle
511 rrd_stat_atomic_add(&ctx->stats.metric_API_consumers, 1);
512 #endif
513 handle->descr = descr;
514 + pg_cache_atomic_get_pg_info(descr, &page_end_time, &page_length);
515 if (unlikely(INVALID_TIME == descr->start_time ||
512 - INVALID_TIME == descr->end_time)) {
516 + INVALID_TIME == page_end_time)) {
517 goto no_more_metrics;
518 }
515 - if (unlikely(descr->start_time != descr->end_time && next_page_time > descr->start_time)) {
519 + if (unlikely(descr->start_time != page_end_time && next_page_time > descr->start_time)) {
520 /* we're in the middle of the page somewhere */
517 - entries = descr->page_length / sizeof(storage_number);
518 - position = ((uint64_t)(next_page_time - descr->start_time)) * entries /
519 - (descr->end_time - descr->start_time + 1);
521 + entries = page_length / sizeof(storage_number);
522 + position = ((uint64_t)(next_page_time - descr->start_time)) * (entries - 1) /
523 + (page_end_time - descr->start_time);
524 } else {
525 position = 0;
526 }
527 }
528 page = descr->pg_cache_descr->page;
529 ret = page[position];
526 - entries = descr->page_length / sizeof(storage_number);
530 + entries = page_length / sizeof(storage_number);
531 if (entries > 1) {
532 usec_t dt;
533
530 - dt = (descr->end_time - descr->start_time) / (entries - 1);
534 + dt = (page_end_time - descr->start_time) / (entries - 1);
535 current_position_time = descr->start_time + position * dt;
536 } else {
537 current_position_time = descr->start_time;