@cryptotaxi247 / netdata-1 / commits / caf7b1919

DB engine optimize RAM usage (#6134)

* Optimize memory footprint of DB engine * Update documentation with the new memory requirements of dbengine * Fixed code style * Fix code style * Fix compile error

Markos Fountoulakis committed May 30, 2019 at 12:09 UTC caf7b1919496c3e8806546dbfcf940668c901955
17 files changed +595 -247
CMakeLists.txt
+2
@@ -456,6 +456,8 @@ set(RRD_PLUGIN_FILES
456 database/engine/rrdengineapi.h
457 database/engine/pagecache.c
458 database/engine/pagecache.h
459 + database/engine/rrdenglocking.c
460 + database/engine/rrdenglocking.h
461 )
462
463 set(WEB_PLUGIN_FILES
Makefile.am
+2
@@ -326,6 +326,8 @@ if ENABLE_DBENGINE
326 database/engine/rrdengineapi.h \
327 database/engine/pagecache.c \
328 database/engine/pagecache.h \
329 + database/engine/rrdenglocking.c \
330 + database/engine/rrdenglocking.h \
331 $(NULL)
332 endif
333
daemon/global_statistics.c
+5 -2
@@ -535,10 +535,10 @@ void global_statistics_charts(void) {
535
536 #ifdef ENABLE_DBENGINE
537 if (localhost->rrd_memory_mode == RRD_MEMORY_MODE_DBENGINE) {
538 - unsigned long long stats_array[27];
538 + unsigned long long stats_array[RRDENG_NR_STATS];
539
540 /* get localhost's DB engine's statistics */
541 - rrdeng_get_27_statistics(localhost->rrdeng_ctx, stats_array);
541 + rrdeng_get_28_statistics(localhost->rrdeng_ctx, stats_array);
542
543 // ----------------------------------------------------------------
544
@@ -637,6 +637,7 @@ void global_statistics_charts(void) {
637
638 {
639 static RRDSET *st_pg_cache_pages = NULL;
640 + static RRDDIM *rd_descriptors = NULL;
641 static RRDDIM *rd_populated = NULL;
642 static RRDDIM *rd_commited = NULL;
643 static RRDDIM *rd_insertions = NULL;
@@ -660,6 +661,7 @@ void global_statistics_charts(void) {
661 , RRDSET_TYPE_LINE
662 );
663
664 + rd_descriptors = rrddim_add(st_pg_cache_pages, "descriptors", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
665 rd_populated = rrddim_add(st_pg_cache_pages, "populated", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
666 rd_commited = rrddim_add(st_pg_cache_pages, "commited", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
667 rd_insertions = rrddim_add(st_pg_cache_pages, "insertions", NULL, 1, 1, RRD_ALGORITHM_INCREMENTAL);
@@ -670,6 +672,7 @@ void global_statistics_charts(void) {
672 else
673 rrdset_next(st_pg_cache_pages);
674
675 + rrddim_set_by_pointer(st_pg_cache_pages, rd_descriptors, (collected_number)stats_array[27]);
676 rrddim_set_by_pointer(st_pg_cache_pages, rd_populated, (collected_number)stats_array[3]);
677 rrddim_set_by_pointer(st_pg_cache_pages, rd_commited, (collected_number)stats_array[4]);
678 rrddim_set_by_pointer(st_pg_cache_pages, rd_insertions, (collected_number)stats_array[5]);
database/engine/README.md
+2 -2
@@ -96,9 +96,9 @@ There are explicit memory requirements **per** DB engine **instance**, meaning *
96
97 - `page cache size` must be at least `#dimensions-being-collected x 4096 x 2` bytes.
98
99 -- an additional `#pages-on-disk x 4096 x 0.06` bytes of RAM are allocated for metadata.
99 +- an additional `#pages-on-disk x 4096 x 0.03` bytes of RAM are allocated for metadata.
100
101 - - roughly speaking this is 6% of the uncompressed disk space taken by the DB files.
101 + - roughly speaking this is 3% of the uncompressed disk space taken by the DB files.
102
103 - for very highly compressible data (compression ratio > 90%) this RAM overhead
104 is comparable to the disk space footprint.
database/engine/datafile.h
+1 -1
@@ -26,7 +26,7 @@ struct extent_info {
26 uint8_t number_of_pages;
27 struct rrdengine_datafile *datafile;
28 struct extent_info *next;
29 - struct rrdeng_page_cache_descr *pages[];
29 + struct rrdeng_page_descr *pages[];
30 };
31
32 struct rrdengine_df_extents {
database/engine/journalfile.c
+1 -1
@@ -226,7 +226,7 @@ static void restore_extent_metadata(struct rrdengine_instance *ctx, struct rrden
226 {
227 struct page_cache *pg_cache = &ctx->pg_cache;
228 unsigned i, count, payload_length, descr_size, valid_pages;
229 - struct rrdeng_page_cache_descr *descr;
229 + struct rrdeng_page_descr *descr;
230 struct extent_info *extent;
231 /* persistent structures */
232 struct rrdeng_jf_store_data *jf_metric_data;
database/engine/pagecache.c
+146 -121
@@ -8,28 +8,29 @@ static int pg_cache_try_evict_one_page_unsafe(struct rrdengine_instance *ctx);
8
9 /* always inserts into tail */
10 static inline void pg_cache_replaceQ_insert_unsafe(struct rrdengine_instance *ctx,
11 - struct rrdeng_page_cache_descr *descr)
11 + struct rrdeng_page_descr *descr)
12 {
13 struct page_cache *pg_cache = &ctx->pg_cache;
14 + struct page_cache_descr *pg_cache_descr = descr->pg_cache_descr;
15
16 if (likely(NULL != pg_cache->replaceQ.tail)) {
16 - descr->prev = pg_cache->replaceQ.tail;
17 - pg_cache->replaceQ.tail->next = descr;
17 + pg_cache_descr->prev = pg_cache->replaceQ.tail;
18 + pg_cache->replaceQ.tail->next = pg_cache_descr;
19 }
20 if (unlikely(NULL == pg_cache->replaceQ.head)) {
20 - pg_cache->replaceQ.head = descr;
21 + pg_cache->replaceQ.head = pg_cache_descr;
22 }
22 - pg_cache->replaceQ.tail = descr;
23 + pg_cache->replaceQ.tail = pg_cache_descr;
24 }
25
26 static inline void pg_cache_replaceQ_delete_unsafe(struct rrdengine_instance *ctx,
26 - struct rrdeng_page_cache_descr *descr)
27 + struct rrdeng_page_descr *descr)
28 {
29 struct page_cache *pg_cache = &ctx->pg_cache;
29 - struct rrdeng_page_cache_descr *prev, *next;
30 + struct page_cache_descr *pg_cache_descr = descr->pg_cache_descr, *prev, *next;
31
31 - prev = descr->prev;
32 - next = descr->next;
32 + prev = pg_cache_descr->prev;
33 + next = pg_cache_descr->next;
34
35 if (likely(NULL != prev)) {
36 prev->next = next;
@@ -37,17 +38,17 @@ static inline void pg_cache_replaceQ_delete_unsafe(struct rrdengine_instance *ct
38 if (likely(NULL != next)) {
39 next->prev = prev;
40 }
40 - if (unlikely(descr == pg_cache->replaceQ.head)) {
41 + if (unlikely(pg_cache_descr == pg_cache->replaceQ.head)) {
42 pg_cache->replaceQ.head = next;
43 }
43 - if (unlikely(descr == pg_cache->replaceQ.tail)) {
44 + if (unlikely(pg_cache_descr == pg_cache->replaceQ.tail)) {
45 pg_cache->replaceQ.tail = prev;
46 }
46 - descr->prev = descr->next = NULL;
47 + pg_cache_descr->prev = pg_cache_descr->next = NULL;
48 }
49
50 void pg_cache_replaceQ_insert(struct rrdengine_instance *ctx,
50 - struct rrdeng_page_cache_descr *descr)
51 + struct rrdeng_page_descr *descr)
52 {
53 struct page_cache *pg_cache = &ctx->pg_cache;
54
@@ -57,7 +58,7 @@ void pg_cache_replaceQ_insert(struct rrdengine_instance *ctx,
58 }
59
60 void pg_cache_replaceQ_delete(struct rrdengine_instance *ctx,
60 - struct rrdeng_page_cache_descr *descr)
61 + struct rrdeng_page_descr *descr)
62 {
63 struct page_cache *pg_cache = &ctx->pg_cache;
64
@@ -66,7 +67,7 @@ void pg_cache_replaceQ_delete(struct rrdengine_instance *ctx,
67 uv_rwlock_wrunlock(&pg_cache->replaceQ.lock);
68 }
69 void pg_cache_replaceQ_set_hot(struct rrdengine_instance *ctx,
69 - struct rrdeng_page_cache_descr *descr)
70 + struct rrdeng_page_descr *descr)
71 {
72 struct page_cache *pg_cache = &ctx->pg_cache;
73
@@ -76,40 +77,28 @@ void pg_cache_replaceQ_set_hot(struct rrdengine_instance *ctx,
77 uv_rwlock_wrunlock(&pg_cache->replaceQ.lock);
78 }
79
79 -struct rrdeng_page_cache_descr *pg_cache_create_descr(void)
80 +struct rrdeng_page_descr *pg_cache_create_descr(void)
81 {
81 - struct rrdeng_page_cache_descr *descr;
82 + struct rrdeng_page_descr *descr;
83
84 descr = mallocz(sizeof(*descr));
84 - descr->page = NULL;
85 descr->page_length = 0;
86 descr->start_time = INVALID_TIME;
87 descr->end_time = INVALID_TIME;
88 descr->id = NULL;
89 descr->extent = NULL;
90 - descr->flags = 0;
91 - descr->prev = descr->next = descr->private = NULL;
92 - descr->refcnt = 0;
93 - descr->waiters = 0;
94 - descr->handle = NULL;
95 - assert(0 == uv_cond_init(&descr->cond));
96 - assert(0 == uv_mutex_init(&descr->mutex));
90 + descr->pg_cache_descr_state = 0;
91 + descr->pg_cache_descr = NULL;
92
93 return descr;
94 }
95
101 -void pg_cache_destroy_descr(struct rrdeng_page_cache_descr *descr)
102 -{
103 - uv_cond_destroy(&descr->cond);
104 - uv_mutex_destroy(&descr->mutex);
105 - free(descr);
106 -}
107 -
96 /* The caller must hold page descriptor lock. */
109 -void pg_cache_wake_up_waiters_unsafe(struct rrdeng_page_cache_descr *descr)
97 +void pg_cache_wake_up_waiters_unsafe(struct rrdeng_page_descr *descr)
98 {
111 - if (descr->waiters)
112 - uv_cond_broadcast(&descr->cond);
99 + struct page_cache_descr *pg_cache_descr = descr->pg_cache_descr;
100 + if (pg_cache_descr->waiters)
101 + uv_cond_broadcast(&pg_cache_descr->cond);
102 }
103
104 /*
@@ -117,11 +106,13 @@ void pg_cache_wake_up_waiters_unsafe(struct rrdeng_page_cache_descr *descr)
106 * The lock will be released and re-acquired. The descriptor is not guaranteed
107 * to exist after this function returns.
108 */
120 -void pg_cache_wait_event_unsafe(struct rrdeng_page_cache_descr *descr)
109 +void pg_cache_wait_event_unsafe(struct rrdeng_page_descr *descr)
110 {
122 - ++descr->waiters;
123 - uv_cond_wait(&descr->cond, &descr->mutex);
124 - --descr->waiters;
111 + struct page_cache_descr *pg_cache_descr = descr->pg_cache_descr;
112 +
113 + ++pg_cache_descr->waiters;
114 + uv_cond_wait(&pg_cache_descr->cond, &pg_cache_descr->mutex);
115 + --pg_cache_descr->waiters;
116 }
117
118 /*
@@ -129,14 +120,15 @@ void pg_cache_wait_event_unsafe(struct rrdeng_page_cache_descr *descr)
120 * The lock will be released and re-acquired. The descriptor is not guaranteed
121 * to exist after this function returns.
122 */
132 -unsigned long pg_cache_wait_event(struct rrdeng_page_cache_descr *descr)
123 +unsigned long pg_cache_wait_event(struct rrdengine_instance *ctx, struct rrdeng_page_descr *descr)
124 {
125 + struct page_cache_descr *pg_cache_descr = descr->pg_cache_descr;
126 unsigned long flags;
127
136 - uv_mutex_lock(&descr->mutex);
128 + rrdeng_page_descr_mutex_lock(ctx, descr);
129 pg_cache_wait_event_unsafe(descr);
138 - flags = descr->flags;
139 - uv_mutex_unlock(&descr->mutex);
130 + flags = pg_cache_descr->flags;
131 + rrdeng_page_descr_mutex_unlock(ctx, descr);
132
133 return flags;
134 }
@@ -146,15 +138,17 @@ unsigned long pg_cache_wait_event(struct rrdeng_page_cache_descr *descr)
138 * Gets a reference to the page descriptor.
139 * Returns 1 on success and 0 on failure.
140 */
149 -int pg_cache_try_get_unsafe(struct rrdeng_page_cache_descr *descr, int exclusive_access)
141 +int pg_cache_try_get_unsafe(struct rrdeng_page_descr *descr, int exclusive_access)
142 {
151 - if ((descr->flags & (RRD_PAGE_LOCKED | RRD_PAGE_READ_PENDING)) ||
152 - (exclusive_access && descr->refcnt)) {
143 + struct page_cache_descr *pg_cache_descr = descr->pg_cache_descr;
144 +
145 + if ((pg_cache_descr->flags & (RRD_PAGE_LOCKED | RRD_PAGE_READ_PENDING)) ||
146 + (exclusive_access && pg_cache_descr->refcnt)) {
147 return 0;
148 }
149 if (exclusive_access)
156 - descr->flags |= RRD_PAGE_LOCKED;
157 - ++descr->refcnt;
150 + pg_cache_descr->flags |= RRD_PAGE_LOCKED;
151 + ++pg_cache_descr->refcnt;
152
153 return 1;
154 }
@@ -163,10 +157,12 @@ int pg_cache_try_get_unsafe(struct rrdeng_page_cache_descr *descr, int exclusive
157 * The caller must hold page descriptor lock.
158 * Same return values as pg_cache_try_get_unsafe() without doing anything.
159 */
166 -int pg_cache_can_get_unsafe(struct rrdeng_page_cache_descr *descr, int exclusive_access)
160 +int pg_cache_can_get_unsafe(struct rrdeng_page_descr *descr, int exclusive_access)
161 {
168 - if ((descr->flags & (RRD_PAGE_LOCKED | RRD_PAGE_READ_PENDING)) ||
169 - (exclusive_access && descr->refcnt)) {
162 + struct page_cache_descr *pg_cache_descr = descr->pg_cache_descr;
163 +
164 + if ((pg_cache_descr->flags & (RRD_PAGE_LOCKED | RRD_PAGE_READ_PENDING)) ||
165 + (exclusive_access && pg_cache_descr->refcnt)) {
166 return 0;
167 }
168
@@ -177,23 +173,24 @@ int pg_cache_can_get_unsafe(struct rrdeng_page_cache_descr *descr, int exclusive
173 * The caller must hold the page descriptor lock.
174 * This function may block doing cleanup.
175 */
180 -void pg_cache_put_unsafe(struct rrdeng_page_cache_descr *descr)
176 +void pg_cache_put_unsafe(struct rrdeng_page_descr *descr)
177 {
182 - descr->flags &= ~RRD_PAGE_LOCKED;
183 - if (0 == --descr->refcnt) {
178 + struct page_cache_descr *pg_cache_descr = descr->pg_cache_descr;
179 +
180 + pg_cache_descr->flags &= ~RRD_PAGE_LOCKED;
181 + if (0 == --pg_cache_descr->refcnt) {
182 pg_cache_wake_up_waiters_unsafe(descr);
183 }
186 - /* TODO: perform cleanup */
184 }
185
186 /*
187 * This function may block doing cleanup.
188 */
192 -void pg_cache_put(struct rrdeng_page_cache_descr *descr)
189 +void pg_cache_put(struct rrdengine_instance *ctx, struct rrdeng_page_descr *descr)
190 {
194 - uv_mutex_lock(&descr->mutex);
191 + rrdeng_page_descr_mutex_lock(ctx, descr);
192 pg_cache_put_unsafe(descr);
196 - uv_mutex_unlock(&descr->mutex);
193 + rrdeng_page_descr_mutex_unlock(ctx, descr);
194 }
195
196 /* The caller must hold the page cache lock */
@@ -286,18 +283,20 @@ static int pg_cache_try_reserve_pages(struct rrdengine_instance *ctx, unsigned n
283 }
284
285 /* The caller must hold the page cache and the page descriptor locks in that order */
289 -static void pg_cache_evict_unsafe(struct rrdengine_instance *ctx, struct rrdeng_page_cache_descr *descr)
286 +static void pg_cache_evict_unsafe(struct rrdengine_instance *ctx, struct rrdeng_page_descr *descr)
287 {
291 - free(descr->page);
292 - descr->page = NULL;
293 - descr->flags &= ~RRD_PAGE_POPULATED;
288 + struct page_cache_descr *pg_cache_descr = descr->pg_cache_descr;
289 +
290 + free(pg_cache_descr->page);
291 + pg_cache_descr->page = NULL;
292 + pg_cache_descr->flags &= ~RRD_PAGE_POPULATED;
293 pg_cache_release_pages_unsafe(ctx, 1);
294 ++ctx->stats.pg_cache_evictions;
295 }
296
297 /*
298 * The caller must hold the page cache lock.
300 - * Lock order: page cache -> replaceQ -> descriptor
299 + * Lock order: page cache -> replaceQ -> page descriptor
300 * This function iterates all pages and tries to evict one.
301 * If it fails it sets in_flight_descr to the oldest descriptor that has write-back in progress,
302 * or it sets it to NULL if no write-back is in progress.
@@ -308,23 +307,29 @@ static int pg_cache_try_evict_one_page_unsafe(struct rrdengine_instance *ctx)
307 {
308 struct page_cache *pg_cache = &ctx->pg_cache;
309 unsigned long old_flags;
311 - struct rrdeng_page_cache_descr *descr;
310 + struct rrdeng_page_descr *descr;
311 + struct page_cache_descr *pg_cache_descr = NULL;
312
313 uv_rwlock_wrlock(&pg_cache->replaceQ.lock);
314 - for (descr = pg_cache->replaceQ.head ; NULL != descr ; descr = descr->next) {
315 - uv_mutex_lock(&descr->mutex);
316 - old_flags = descr->flags;
314 + for (pg_cache_descr = pg_cache->replaceQ.head ; NULL != pg_cache_descr ; pg_cache_descr = pg_cache_descr->next) {
315 + descr = pg_cache_descr->descr;
316 +
317 + rrdeng_page_descr_mutex_lock(ctx, descr);
318 + old_flags = pg_cache_descr->flags;
319 if ((old_flags & RRD_PAGE_POPULATED) && !(old_flags & RRD_PAGE_DIRTY) && pg_cache_try_get_unsafe(descr, 1)) {
320 /* must evict */
321 pg_cache_evict_unsafe(ctx, descr);
322 pg_cache_put_unsafe(descr);
321 - uv_mutex_unlock(&descr->mutex);
323 pg_cache_replaceQ_delete_unsafe(ctx, descr);
324 +
325 + rrdeng_page_descr_mutex_unlock(ctx, descr);
326 uv_rwlock_wrunlock(&pg_cache->replaceQ.lock);
327
328 + rrdeng_try_deallocate_pg_cache_descr(ctx, descr);
329 +
330 return 1;
331 }
327 - uv_mutex_unlock(&descr->mutex);
332 + rrdeng_page_descr_mutex_unlock(ctx, descr);
333 };
334 uv_rwlock_wrunlock(&pg_cache->replaceQ.lock);
335
@@ -335,9 +340,10 @@ static int pg_cache_try_evict_one_page_unsafe(struct rrdengine_instance *ctx)
340 /*
341 * TODO: last waiter frees descriptor ?
342 */
338 -void pg_cache_punch_hole(struct rrdengine_instance *ctx, struct rrdeng_page_cache_descr *descr)
343 +void pg_cache_punch_hole(struct rrdengine_instance *ctx, struct rrdeng_page_descr *descr)
344 {
345 struct page_cache *pg_cache = &ctx->pg_cache;
346 + struct page_cache_descr *pg_cache_descr = NULL;
347 Pvoid_t *PValue;
348 struct pg_cache_page_index *page_index;
349 int ret;
@@ -353,8 +359,9 @@ void pg_cache_punch_hole(struct rrdengine_instance *ctx, struct rrdeng_page_cach
359 uv_rwlock_wrunlock(&page_index->lock);
360 if (unlikely(0 == ret)) {
361 error("Page under deletion was not in index.");
356 - if (unlikely(debug_flags & D_RRDENGINE))
357 - print_page_cache_descr(descr);
362 + if (unlikely(debug_flags & D_RRDENGINE)) {
363 + print_page_descr(descr);
364 + }
365 goto destroy;
366 }
367 assert(1 == ret);
@@ -364,37 +371,44 @@ void pg_cache_punch_hole(struct rrdengine_instance *ctx, struct rrdeng_page_cach
371 --pg_cache->page_descriptors;
372 uv_rwlock_wrunlock(&pg_cache->pg_cache_rwlock);
373
367 - uv_mutex_lock(&descr->mutex);
368 - while (!pg_cache_try_get_unsafe(descr, 1)) {
369 - debug(D_RRDENGINE, "%s: Waiting for locked page:", __func__);
370 - if(unlikely(debug_flags & D_RRDENGINE))
371 - print_page_cache_descr(descr);
372 - pg_cache_wait_event_unsafe(descr);
373 - }
374 - /* even a locked page could be dirty */
375 - while (unlikely(descr->flags & RRD_PAGE_DIRTY)) {
376 - debug(D_RRDENGINE, "%s: Found dirty page, waiting for it to be flushed:", __func__);
377 - if (unlikely(debug_flags & D_RRDENGINE))
378 - print_page_cache_descr(descr);
379 - pg_cache_wait_event_unsafe(descr);
380 - }
381 - uv_mutex_unlock(&descr->mutex);
374 + if (0 != descr->pg_cache_descr_state) {
375 + /* there is page cache descriptor state under possible contention */
376
383 - if (descr->flags & RRD_PAGE_POPULATED) {
384 - /* only after locking can it be safely deleted from LRU */
385 - pg_cache_replaceQ_delete(ctx, descr);
377 + rrdeng_page_descr_mutex_lock(ctx, descr);
378 + pg_cache_descr = descr->pg_cache_descr;
379 + while (!pg_cache_try_get_unsafe(descr, 1)) {
380 + debug(D_RRDENGINE, "%s: Waiting for locked page:", __func__);
381 + if (unlikely(debug_flags & D_RRDENGINE))
382 + print_page_cache_descr(descr);
383 + pg_cache_wait_event_unsafe(descr);
384 + }
385 + /* even a locked page could be dirty */
386 + while (unlikely(pg_cache_descr->flags & RRD_PAGE_DIRTY)) {
387 + debug(D_RRDENGINE, "%s: Found dirty page, waiting for it to be flushed:", __func__);
388 + if (unlikely(debug_flags & D_RRDENGINE))
389 + print_page_cache_descr(descr);
390 + pg_cache_wait_event_unsafe(descr);
391 + }
392 + if (pg_cache_descr->flags & RRD_PAGE_POPULATED) {
393 + /* only after locking can it be safely deleted from LRU */
394 + pg_cache_replaceQ_delete(ctx, descr);
395
387 - uv_rwlock_wrlock(&pg_cache->pg_cache_rwlock);
388 - pg_cache_evict_unsafe(ctx, descr);
389 - uv_rwlock_wrunlock(&pg_cache->pg_cache_rwlock);
396 + uv_rwlock_wrlock(&pg_cache->pg_cache_rwlock);
397 + pg_cache_evict_unsafe(ctx, descr);
398 + uv_rwlock_wrunlock(&pg_cache->pg_cache_rwlock);
399 + }
400 + pg_cache_put_unsafe(descr);
401 + rrdeng_page_descr_mutex_unlock(ctx, descr);
402 +
403 + rrdeng_destroy_pg_cache_descr(ctx, pg_cache_descr);
404 }
391 - pg_cache_put(descr);
405 destroy:
393 - pg_cache_destroy_descr(descr);
406 + assert(0 == descr->pg_cache_descr_state);
407 + freez(descr);
408 pg_cache_update_metric_times(page_index);
409 }
410
397 -static inline int is_page_in_time_range(struct rrdeng_page_cache_descr *descr, usec_t start_time, usec_t end_time)
411 +static inline int is_page_in_time_range(struct rrdeng_page_descr *descr, usec_t start_time, usec_t end_time)
412 {
413 usec_t pg_start, pg_end;
414
@@ -405,13 +419,13 @@ static inline int is_page_in_time_range(struct rrdeng_page_cache_descr *descr, u
419 (pg_start >= start_time && pg_start <= end_time);
420 }
421
408 -static inline int is_point_in_time_in_page(struct rrdeng_page_cache_descr *descr, usec_t point_in_time)
422 +static inline int is_point_in_time_in_page(struct rrdeng_page_descr *descr, usec_t point_in_time)
423 {
424 return (point_in_time >= descr->start_time && point_in_time <= descr->end_time);
425 }
426
427 /* Update metric oldest and latest timestamps efficiently when adding new values */
414 -void pg_cache_add_new_metric_time(struct pg_cache_page_index *page_index, struct rrdeng_page_cache_descr *descr)
428 +void pg_cache_add_new_metric_time(struct pg_cache_page_index *page_index, struct rrdeng_page_descr *descr)
429 {
430 usec_t oldest_time = page_index->oldest_time;
431 usec_t latest_time = page_index->latest_time;
@@ -429,7 +443,7 @@ void pg_cache_update_metric_times(struct pg_cache_page_index *page_index)
443 {
444 Pvoid_t *firstPValue, *lastPValue;
445 Word_t firstIndex, lastIndex;
432 - struct rrdeng_page_cache_descr *descr;
446 + struct rrdeng_page_descr *descr;
447 usec_t oldest_time = INVALID_TIME;
448 usec_t latest_time = INVALID_TIME;
449
@@ -460,16 +474,23 @@ void pg_cache_update_metric_times(struct pg_cache_page_index *page_index)
474
475 /* If index is NULL lookup by UUID (descr->id) */
476 void pg_cache_insert(struct rrdengine_instance *ctx, struct pg_cache_page_index *index,
463 - struct rrdeng_page_cache_descr *descr)
477 + struct rrdeng_page_descr *descr)
478 {
479 struct page_cache *pg_cache = &ctx->pg_cache;
480 Pvoid_t *PValue;
481 struct pg_cache_page_index *page_index;
482 + unsigned long pg_cache_descr_state = descr->pg_cache_descr_state;
483 +
484 + if (0 != pg_cache_descr_state) {
485 + /* there is page cache descriptor pre-allocated state */
486 + struct page_cache_descr *pg_cache_descr = descr->pg_cache_descr;
487
469 - if (descr->flags & RRD_PAGE_POPULATED) {
470 - pg_cache_reserve_pages(ctx, 1);
471 - if (!(descr->flags & RRD_PAGE_DIRTY))
472 - pg_cache_replaceQ_insert(ctx, descr);
488 + assert(pg_cache_descr_state & PG_CACHE_DESCR_ALLOCATED);
489 + if (pg_cache_descr->flags & RRD_PAGE_POPULATED) {
490 + pg_cache_reserve_pages(ctx, 1);
491 + if (!(pg_cache_descr->flags & RRD_PAGE_DIRTY))
492 + pg_cache_replaceQ_insert(ctx, descr);
493 + }
494 }
495
496 if (unlikely(NULL == index)) {
@@ -503,7 +524,8 @@ struct pg_cache_page_index *
524 pg_cache_preload(struct rrdengine_instance *ctx, uuid_t *id, usec_t start_time, usec_t end_time)
525 {
526 struct page_cache *pg_cache = &ctx->pg_cache;
506 - struct rrdeng_page_cache_descr *descr = NULL, *preload_array[PAGE_CACHE_MAX_PRELOAD_PAGES];
527 + struct rrdeng_page_descr *descr = NULL, *preload_array[PAGE_CACHE_MAX_PRELOAD_PAGES];
528 + struct page_cache_descr *pg_cache_descr = NULL;
529 int i, j, k, count, found;
530 unsigned long flags;
531 Pvoid_t *PValue;
@@ -557,12 +579,13 @@ struct pg_cache_page_index *
579
580 if (unlikely(0 == descr->page_length))
581 continue;
560 - uv_mutex_lock(&descr->mutex);
561 - flags = descr->flags;
582 + rrdeng_page_descr_mutex_lock(ctx, descr);
583 + pg_cache_descr = descr->pg_cache_descr;
584 + flags = pg_cache_descr->flags;
585 if (pg_cache_can_get_unsafe(descr, 0)) {
586 if (flags & RRD_PAGE_POPULATED) {
587 /* success */
565 - uv_mutex_unlock(&descr->mutex);
588 + rrdeng_page_descr_mutex_unlock(ctx, descr);
589 debug(D_RRDENGINE, "%s: Page was found in memory.", __func__);
590 continue;
591 }
@@ -570,11 +593,11 @@ struct pg_cache_page_index *
593 if (!(flags & RRD_PAGE_POPULATED) && pg_cache_try_get_unsafe(descr, 1)) {
594 preload_array[count++] = descr;
595 if (PAGE_CACHE_MAX_PRELOAD_PAGES == count) {
573 - uv_mutex_unlock(&descr->mutex);
596 + rrdeng_page_descr_mutex_unlock(ctx, descr);
597 break;
598 }
599 }
577 - uv_mutex_unlock(&descr->mutex);
600 + rrdeng_page_descr_mutex_unlock(ctx, descr);
601
602 };
603 uv_rwlock_rdunlock(&page_index->lock);
@@ -582,7 +605,7 @@ struct pg_cache_page_index *
605 failed_to_reserve = 0;
606 for (i = 0 ; i < count && !failed_to_reserve ; ++i) {
607 struct rrdeng_cmd cmd;
585 - struct rrdeng_page_cache_descr *next;
608 + struct rrdeng_page_descr *next;
609
610 descr = preload_array[i];
611 if (NULL == descr) {
@@ -622,7 +645,7 @@ struct pg_cache_page_index *
645 if (NULL == descr) {
646 continue;
647 }
625 - pg_cache_put(descr);
648 + pg_cache_put(ctx, descr);
649 }
650 }
651 if (!count) {
@@ -637,12 +660,13 @@ struct pg_cache_page_index *
660 * When point_in_time is INVALID_TIME get any page.
661 * If index is NULL lookup by UUID (id).
662 */
640 -struct rrdeng_page_cache_descr *
663 +struct rrdeng_page_descr *
664 pg_cache_lookup(struct rrdengine_instance *ctx, struct pg_cache_page_index *index, uuid_t *id,
665 usec_t point_in_time)
666 {
667 struct page_cache *pg_cache = &ctx->pg_cache;
645 - struct rrdeng_page_cache_descr *descr = NULL;
668 + struct rrdeng_page_descr *descr = NULL;
669 + struct page_cache_descr *pg_cache_descr = NULL;
670 unsigned long flags;
671 Pvoid_t *PValue;
672 struct pg_cache_page_index *page_index;
@@ -682,11 +706,12 @@ struct rrdeng_page_cache_descr *
706 pg_cache_release_pages(ctx, 1);
707 return NULL;
708 }
685 - uv_mutex_lock(&descr->mutex);
686 - flags = descr->flags;
709 + rrdeng_page_descr_mutex_lock(ctx, descr);
710 + pg_cache_descr = descr->pg_cache_descr;
711 + flags = pg_cache_descr->flags;
712 if ((flags & RRD_PAGE_POPULATED) && pg_cache_try_get_unsafe(descr, 0)) {
713 /* success */
689 - uv_mutex_unlock(&descr->mutex);
714 + rrdeng_page_descr_mutex_unlock(ctx, descr);
715 debug(D_RRDENGINE, "%s: Page was found in memory.", __func__);
716 break;
717 }
@@ -702,14 +727,14 @@ struct rrdeng_page_cache_descr *
727 debug(D_RRDENGINE, "%s: Waiting for page to be asynchronously read from disk:", __func__);
728 if(unlikely(debug_flags & D_RRDENGINE))
729 print_page_cache_descr(descr);
705 - while (!(descr->flags & RRD_PAGE_POPULATED)) {
730 + while (!(pg_cache_descr->flags & RRD_PAGE_POPULATED)) {
731 pg_cache_wait_event_unsafe(descr);
732 }
733 /* success */
734 /* Downgrade exclusive reference to allow other readers */
710 - descr->flags &= ~RRD_PAGE_LOCKED;
735 + pg_cache_descr->flags &= ~RRD_PAGE_LOCKED;
736 pg_cache_wake_up_waiters_unsafe(descr);
712 - uv_mutex_unlock(&descr->mutex);
737 + rrdeng_page_descr_mutex_unlock(ctx, descr);
738 rrd_stat_atomic_add(&ctx->stats.pg_cache_misses, 1);
739 return descr;
740 }
@@ -720,7 +745,7 @@ struct rrdeng_page_cache_descr *
745 if (!(flags & RRD_PAGE_POPULATED))
746 page_not_in_cache = 1;
747 pg_cache_wait_event_unsafe(descr);
723 - uv_mutex_unlock(&descr->mutex);
748 + rrdeng_page_descr_mutex_unlock(ctx, descr);
749
750 /* reset scan to find again */
751 uv_rwlock_rdlock(&page_index->lock);
database/engine/pagecache.h
+56 -30
@@ -5,9 +5,10 @@
5
6 #include "rrdengine.h"
7
8 -/* Forward declerations */
8 +/* Forward declarations */
9 struct rrdengine_instance;
10 struct extent_info;
11 +struct rrdeng_page_descr;
12
13 #define INVALID_TIME (0)
14
@@ -18,24 +19,46 @@ struct extent_info;
19 #define RRD_PAGE_WRITE_PENDING (1LU << 3)
20 #define RRD_PAGE_POPULATED (1LU << 4)
21
21 -struct rrdeng_page_cache_descr {
22 +struct page_cache_descr {
23 + struct rrdeng_page_descr *descr; /* parent descriptor */
24 void *page;
23 - uint32_t page_length;
24 - usec_t start_time;
25 - usec_t end_time;
26 - uuid_t *id; /* never changes */
27 - struct extent_info *extent;
25 unsigned long flags;
29 - void *private;
30 - struct rrdeng_page_cache_descr *prev;
31 - struct rrdeng_page_cache_descr *next;
26 + struct page_cache_descr *prev; /* LRU */
27 + struct page_cache_descr *next; /* LRU */
28
33 - /* TODO: move waiter logic to concurrency table */
29 unsigned refcnt;
30 uv_mutex_t mutex; /* always take it after the page cache lock or after the commit lock */
31 uv_cond_t cond;
32 unsigned waiters;
38 - struct rrdeng_collect_handle *handle; /* API user */
33 +};
34 +
35 +/* Page cache descriptor flags, state = 0 means no descriptor */
36 +#define PG_CACHE_DESCR_ALLOCATED (1LU << 0)
37 +#define PG_CACHE_DESCR_DESTROY (1LU << 1)
38 +#define PG_CACHE_DESCR_LOCKED (1LU << 2)
39 +#define PG_CACHE_DESCR_SHIFT (3)
40 +#define PG_CACHE_DESCR_USERS_MASK (((unsigned long)-1) << PG_CACHE_DESCR_SHIFT)
41 +#define PG_CACHE_DESCR_FLAGS_MASK (((unsigned long)-1) >> (BITS_PER_ULONG - PG_CACHE_DESCR_SHIFT))
42 +
43 +/*
44 + * Page cache descriptor state bits (works for both 32-bit and 64-bit architectures):
45 + *
46 + * 63 ... 31 ... 3 | 2 | 1 | 0|
47 + * -----------------------------+------------+------------+-----------|
48 + * number of descriptor users | DESTROY | LOCKED | ALLOCATED |
49 + */
50 +struct rrdeng_page_descr {
51 + uint32_t page_length;
52 + usec_t start_time;
53 + usec_t end_time;
54 + uuid_t *id; /* never changes */
55 + struct extent_info *extent;
56 +
57 + /* points to ephemeral page cache descriptor if the page resides in the cache */
58 + struct page_cache_descr *pg_cache_descr;
59 +
60 + /* Compare-And-Swap target for page cache descriptor allocation algorithm */
61 + volatile unsigned long pg_cache_descr_state;
62 };
63
64 #define PAGE_CACHE_MAX_PRELOAD_PAGES (256)
@@ -85,12 +108,15 @@ struct pg_cache_commited_page_index {
108 unsigned nr_commited_pages;
109 };
110
88 -/* gathers populated pages to be evicted */
111 +/*
112 + * Gathers populated pages to be evicted.
113 + * Relies on page cache descriptors being there as it uses their memory.
114 + */
115 struct pg_cache_replaceQ {
116 uv_rwlock_t lock; /* LRU lock */
117
92 - struct rrdeng_page_cache_descr *head; /* LRU */
93 - struct rrdeng_page_cache_descr *tail; /* MRU */
118 + struct page_cache_descr *head; /* LRU */
119 + struct page_cache_descr *tail; /* MRU */
120 };
121
122 struct page_cache { /* TODO: add statistics */
@@ -104,30 +130,30 @@ struct page_cache { /* TODO: add statistics */
130 unsigned populated_pages;
131 };
132
107 -extern void pg_cache_wake_up_waiters_unsafe(struct rrdeng_page_cache_descr *descr);
108 -extern void pg_cache_wait_event_unsafe(struct rrdeng_page_cache_descr *descr);
109 -extern unsigned long pg_cache_wait_event(struct rrdeng_page_cache_descr *descr);
133 +extern void pg_cache_wake_up_waiters_unsafe(struct rrdeng_page_descr *descr);
134 +extern void pg_cache_wait_event_unsafe(struct rrdeng_page_descr *descr);
135 +extern unsigned long pg_cache_wait_event(struct rrdengine_instance *ctx, struct rrdeng_page_descr *descr);
136 extern void pg_cache_replaceQ_insert(struct rrdengine_instance *ctx,
111 - struct rrdeng_page_cache_descr *descr);
137 + struct rrdeng_page_descr *descr);
138 extern void pg_cache_replaceQ_delete(struct rrdengine_instance *ctx,
113 - struct rrdeng_page_cache_descr *descr);
139 + struct rrdeng_page_descr *descr);
140 extern void pg_cache_replaceQ_set_hot(struct rrdengine_instance *ctx,
115 - struct rrdeng_page_cache_descr *descr);
116 -extern struct rrdeng_page_cache_descr *pg_cache_create_descr(void);
117 -extern int pg_cache_try_get_unsafe(struct rrdeng_page_cache_descr *descr, int exclusive_access);
118 -extern void pg_cache_put_unsafe(struct rrdeng_page_cache_descr *descr);
119 -extern void pg_cache_put(struct rrdeng_page_cache_descr *descr);
141 + struct rrdeng_page_descr *descr);
142 +extern struct rrdeng_page_descr *pg_cache_create_descr(void);
143 +extern int pg_cache_try_get_unsafe(struct rrdeng_page_descr *descr, int exclusive_access);
144 +extern void pg_cache_put_unsafe(struct rrdeng_page_descr *descr);
145 +extern void pg_cache_put(struct rrdengine_instance *ctx, struct rrdeng_page_descr *descr);
146 extern void pg_cache_insert(struct rrdengine_instance *ctx, struct pg_cache_page_index *index,
121 - struct rrdeng_page_cache_descr *descr);
122 -extern void pg_cache_punch_hole(struct rrdengine_instance *ctx, struct rrdeng_page_cache_descr *descr);
147 + struct rrdeng_page_descr *descr);
148 +extern void pg_cache_punch_hole(struct rrdengine_instance *ctx, struct rrdeng_page_descr *descr);
149 extern struct pg_cache_page_index *
150 pg_cache_preload(struct rrdengine_instance *ctx, uuid_t *id, usec_t start_time, usec_t end_time);
125 -extern struct rrdeng_page_cache_descr *
151 +extern struct rrdeng_page_descr *
152 pg_cache_lookup(struct rrdengine_instance *ctx, struct pg_cache_page_index *index, uuid_t *id,
153 usec_t point_in_time);
154 extern struct pg_cache_page_index *create_page_index(uuid_t *id);
155 extern void init_page_cache(struct rrdengine_instance *ctx);
130 -extern void pg_cache_add_new_metric_time(struct pg_cache_page_index *page_index, struct rrdeng_page_cache_descr *descr);
156 +extern void pg_cache_add_new_metric_time(struct pg_cache_page_index *page_index, struct rrdeng_page_descr *descr);
157 extern void pg_cache_update_metric_times(struct pg_cache_page_index *page_index);
158
133 -#endif /* NETDATA_PAGECACHE_H */
\ No newline at end of file
159 +#endif /* NETDATA_PAGECACHE_H */
database/engine/rrdengine.c
+42 -34
@@ -27,7 +27,8 @@ void read_extent_cb(uv_fs_t* req)
27 struct rrdengine_worker_config* wc = req->loop->data;
28 struct rrdengine_instance *ctx = wc->ctx;
29 struct extent_io_descriptor *xt_io_descr;
30 - struct rrdeng_page_cache_descr *descr;
30 + struct rrdeng_page_descr *descr;
31 + struct page_cache_descr *pg_cache_descr;
32 int ret;
33 unsigned i, j, count;
34 void *page, *uncompressed_buf = NULL;
@@ -97,17 +98,18 @@ void read_extent_cb(uv_fs_t* req)
98 (void) memcpy(page, uncompressed_buf + page_offset, descr->page_length);
99 }
100 pg_cache_replaceQ_insert(ctx, descr);
100 - uv_mutex_lock(&descr->mutex);
101 - descr->page = page;
102 - descr->flags |= RRD_PAGE_POPULATED;
103 - descr->flags &= ~RRD_PAGE_READ_PENDING;
101 + rrdeng_page_descr_mutex_lock(ctx, descr);
102 + pg_cache_descr = descr->pg_cache_descr;
103 + pg_cache_descr->page = page;
104 + pg_cache_descr->flags |= RRD_PAGE_POPULATED;
105 + pg_cache_descr->flags &= ~RRD_PAGE_READ_PENDING;
106 debug(D_RRDENGINE, "%s: Waking up waiters.", __func__);
107 if (xt_io_descr->release_descr) {
108 pg_cache_put_unsafe(descr);
109 } else {
110 pg_cache_wake_up_waiters_unsafe(descr);
111 }
110 - uv_mutex_unlock(&descr->mutex);
112 + rrdeng_page_descr_mutex_unlock(ctx, descr);
113 }
114 if (RRD_NO_COMPRESSION != header->compression_algorithm) {
115 free(uncompressed_buf);
@@ -122,11 +124,12 @@ cleanup:
124
125
126 static void do_read_extent(struct rrdengine_worker_config* wc,
125 - struct rrdeng_page_cache_descr **descr,
127 + struct rrdeng_page_descr **descr,
128 unsigned count,
129 uint8_t release_descr)
130 {
131 struct rrdengine_instance *ctx = wc->ctx;
132 + struct page_cache_descr *pg_cache_descr;
133 int ret;
134 unsigned i, size_bytes, pos, real_io_size;
135 // uint32_t payload_length;
@@ -145,10 +148,11 @@ static void do_read_extent(struct rrdengine_worker_config* wc,
148 return;*/
149 }
150 for (i = 0 ; i < count; ++i) {
148 - uv_mutex_lock(&descr[i]->mutex);
149 - descr[i]->flags |= RRD_PAGE_READ_PENDING;
151 + rrdeng_page_descr_mutex_lock(ctx, descr[i]);
152 + pg_cache_descr = descr[i]->pg_cache_descr;
153 + pg_cache_descr->flags |= RRD_PAGE_READ_PENDING;
154 // payload_length = descr[i]->page_length;
151 - uv_mutex_unlock(&descr[i]->mutex);
155 + rrdeng_page_descr_mutex_unlock(ctx, descr[i]);
156
157 xt_io_descr->descr_array[i] = descr[i];
158 }
@@ -227,7 +231,8 @@ void flush_pages_cb(uv_fs_t* req)
231 struct rrdengine_instance *ctx = wc->ctx;
232 struct page_cache *pg_cache = &ctx->pg_cache;
233 struct extent_io_descriptor *xt_io_descr;
230 - struct rrdeng_page_cache_descr *descr;
234 + struct rrdeng_page_descr *descr;
235 + struct page_cache_descr *pg_cache_descr;
236 struct rrdengine_datafile *datafile;
237 int ret;
238 unsigned i, count;
@@ -256,11 +261,12 @@ void flush_pages_cb(uv_fs_t* req)
261
262 pg_cache_replaceQ_insert(ctx, descr);
263
259 - uv_mutex_lock(&descr->mutex);
260 - descr->flags &= ~(RRD_PAGE_DIRTY | RRD_PAGE_WRITE_PENDING);
264 + rrdeng_page_descr_mutex_lock(ctx, descr);
265 + pg_cache_descr = descr->pg_cache_descr;
266 + pg_cache_descr->flags &= ~(RRD_PAGE_DIRTY | RRD_PAGE_WRITE_PENDING);
267 /* wake up waiters, care no reference being held */
268 pg_cache_wake_up_waiters_unsafe(descr);
263 - uv_mutex_unlock(&descr->mutex);
269 + rrdeng_page_descr_mutex_unlock(ctx, descr);
270 }
271 if (xt_io_descr->completion)
272 complete(xt_io_descr->completion);
@@ -283,7 +289,8 @@ static int do_flush_pages(struct rrdengine_worker_config* wc, int force, struct
289 int compressed_size, max_compressed_size = 0;
290 unsigned i, count, size_bytes, pos, real_io_size;
291 uint32_t uncompressed_payload_length, payload_offset;
286 - struct rrdeng_page_cache_descr *descr, *eligible_pages[MAX_PAGES_PER_EXTENT];
292 + struct rrdeng_page_descr *descr, *eligible_pages[MAX_PAGES_PER_EXTENT];
293 + struct page_cache_descr *pg_cache_descr;
294 struct extent_io_descriptor *xt_io_descr;
295 void *compressed_buf = NULL;
296 Word_t descr_commit_idx_array[MAX_PAGES_PER_EXTENT];
@@ -311,15 +318,16 @@ static int do_flush_pages(struct rrdengine_worker_config* wc, int force, struct
318 descr = unlikely(NULL == PValue) ? NULL : *PValue) {
319 assert(0 != descr->page_length);
320
314 - uv_mutex_lock(&descr->mutex);
315 - if (!(descr->flags & RRD_PAGE_WRITE_PENDING)) {
321 + rrdeng_page_descr_mutex_lock(ctx, descr);
322 + pg_cache_descr = descr->pg_cache_descr;
323 + if (!(pg_cache_descr->flags & RRD_PAGE_WRITE_PENDING)) {
324 /* care, no reference being held */
317 - descr->flags |= RRD_PAGE_WRITE_PENDING;
325 + pg_cache_descr->flags |= RRD_PAGE_WRITE_PENDING;
326 uncompressed_payload_length += descr->page_length;
327 descr_commit_idx_array[count] = Index;
328 eligible_pages[count++] = descr;
329 }
322 - uv_mutex_unlock(&descr->mutex);
330 + rrdeng_page_descr_mutex_unlock(ctx, descr);
331 }
332 uv_rwlock_rdunlock(&pg_cache->commited_page_index.lock);
333
@@ -347,7 +355,7 @@ static int do_flush_pages(struct rrdengine_worker_config* wc, int force, struct
355 fatal("posix_memalign:%s", strerror(ret));
356 /* free(xt_io_descr);*/
357 }
350 - (void) memcpy(xt_io_descr->descr_array, eligible_pages, sizeof(struct rrdeng_page_cache_descr *) * count);
358 + (void) memcpy(xt_io_descr->descr_array, eligible_pages, sizeof(struct rrdeng_page_descr *) * count);
359 xt_io_descr->descr_count = count;
360
361 pos = 0;
@@ -378,7 +386,7 @@ static int do_flush_pages(struct rrdengine_worker_config* wc, int force, struct
386 for (i = 0 ; i < count ; ++i) {
387 descr = xt_io_descr->descr_array[i];
388 /* care, we don't hold the descriptor mutex */
381 - (void) memcpy(xt_io_descr->buf + pos, descr->page, descr->page_length);
389 + (void) memcpy(xt_io_descr->buf + pos, descr->pg_cache_descr->page, descr->page_length);
390 descr->extent = extent;
391 extent->pages[i] = descr;
392
@@ -464,7 +472,7 @@ static void delete_old_data(uv_work_t *req)
472 struct rrdengine_instance *ctx = req->data;
473 struct rrdengine_datafile *datafile;
474 struct extent_info *extent, *next;
467 - struct rrdeng_page_cache_descr *descr;
475 + struct rrdeng_page_descr *descr;
476 unsigned count, i;
477
478 /* Safe to use since it will be deleted after we are done */
@@ -690,14 +698,14 @@ void rrdeng_worker(void* arg)
698 do_commit_transaction(wc, STORE_DATA, NULL);
699 break;
700 case RRDENG_FLUSH_PAGES: {
693 - unsigned total_bytes, bytes_written;
701 + unsigned bytes_written;
702
703 /* First I/O should be enough to call completion */
704 bytes_written = do_flush_pages(wc, 1, cmd.completion);
697 - for (total_bytes = bytes_written ;
698 - bytes_written && (total_bytes < DATAFILE_IDEAL_IO_SIZE) ;
699 - total_bytes += bytes_written) {
700 - bytes_written = do_flush_pages(wc, 1, NULL);
705 + if (bytes_written) {
706 + while (do_flush_pages(wc, 1, NULL)) {
707 + ; /* Force flushing of all commited pages. */
708 + }
709 }
710 break;
711 }
@@ -726,19 +734,19 @@ static void basic_functional_test(struct rrdengine_instance *ctx)
734 int i, j, failed_validations;
735 uuid_t uuid[NR_PAGES];
736 void *buf;
729 - struct rrdeng_page_cache_descr *handle[NR_PAGES];
730 - char uuid_str[37];
731 - char backup[NR_PAGES][37 * 100]; /* backup storage for page data verification */
737 + struct rrdeng_page_descr *handle[NR_PAGES];
738 + char uuid_str[UUID_STR_LEN];
739 + char backup[NR_PAGES][UUID_STR_LEN * 100]; /* backup storage for page data verification */
740
741 for (i = 0 ; i < NR_PAGES ; ++i) {
742 uuid_generate(uuid[i]);
743 uuid_unparse_lower(uuid[i], uuid_str);
744 // fprintf(stderr, "Generated uuid[%d]=%s\n", i, uuid_str);
737 - buf = rrdeng_create_page(&uuid[i], &handle[i]);
745 + buf = rrdeng_create_page(ctx, &uuid[i], &handle[i]);
746 /* Each page contains 10 times its own UUID stringified */
747 for (j = 0 ; j < 100 ; ++j) {
740 - strcpy(buf + 37 * j, uuid_str);
741 - strcpy(backup[i] + 37 * j, uuid_str);
748 + strcpy(buf + UUID_STR_LEN * j, uuid_str);
749 + strcpy(backup[i] + UUID_STR_LEN * j, uuid_str);
750 }
751 rrdeng_commit_page(ctx, handle[i], (Word_t)i);
752 }
@@ -750,7 +758,7 @@ static void basic_functional_test(struct rrdengine_instance *ctx)
758 ++failed_validations;
759 fprintf(stderr, "Page %d was LOST.\n", i);
760 }
753 - if (memcmp(backup[i], buf, 37 * 100)) {
761 + if (memcmp(backup[i], buf, UUID_STR_LEN * 100)) {
762 ++failed_validations;
763 fprintf(stderr, "Page %d data comparison with backup FAILED validation.\n", i);
764 }
database/engine/rrdengine.h
+5 -3
@@ -22,6 +22,7 @@
22 #include "journalfile.h"
23 #include "rrdengineapi.h"
24 #include "pagecache.h"
25 +#include "rrdenglocking.h"
26
27 #ifdef NETDATA_RRD_INTERNALS
28
@@ -59,10 +60,10 @@ struct rrdeng_cmd {
60 enum rrdeng_opcode opcode;
61 union {
62 struct rrdeng_read_page {
62 - struct rrdeng_page_cache_descr *page_cache_descr;
63 + struct rrdeng_page_descr *page_cache_descr;
64 } read_page;
65 struct rrdeng_read_extent {
65 - struct rrdeng_page_cache_descr *page_cache_descr[MAX_PAGES_PER_EXTENT];
66 + struct rrdeng_page_descr *page_cache_descr[MAX_PAGES_PER_EXTENT];
67 int page_count;
68 } read_extent;
69 struct completion *completion;
@@ -85,7 +86,7 @@ struct extent_io_descriptor {
86 struct completion *completion;
87 unsigned descr_count;
88 int release_descr;
88 - struct rrdeng_page_cache_descr *descr_array[MAX_PAGES_PER_EXTENT];
89 + struct rrdeng_page_descr *descr_array[MAX_PAGES_PER_EXTENT];
90 Word_t descr_commit_idx_array[MAX_PAGES_PER_EXTENT];
91 };
92
@@ -142,6 +143,7 @@ struct rrdengine_statistics {
143 rrdeng_stats_t datafile_deletions;
144 rrdeng_stats_t journalfile_creations;
145 rrdeng_stats_t journalfile_deletions;
146 + rrdeng_stats_t page_cache_descriptors;
147 };
148
149 struct rrdengine_instance {
database/engine/rrdengineapi.c
+41 -32
@@ -65,7 +65,7 @@ void rrdeng_store_metric_next(RRDDIM *rd, usec_t point_in_time, storage_number n
65 struct rrdeng_collect_handle *handle;
66 struct rrdengine_instance *ctx;
67 struct page_cache *pg_cache;
68 - struct rrdeng_page_cache_descr *descr;
68 + struct rrdeng_page_descr *descr;
69 storage_number *page;
70
71 handle = &rd->state->handle.rrdeng;
@@ -74,7 +74,6 @@ void rrdeng_store_metric_next(RRDDIM *rd, usec_t point_in_time, storage_number n
74 descr = handle->descr;
75 if (unlikely(NULL == descr || descr->page_length + sizeof(number) > RRDENG_BLOCK_SIZE)) {
76 if (descr) {
77 - descr->handle = NULL;
77 if (descr->page_length) {
78 int ret;
79
@@ -82,33 +81,33 @@ void rrdeng_store_metric_next(RRDDIM *rd, usec_t point_in_time, storage_number n
81 rrd_stat_atomic_add(&ctx->stats.metric_API_producers, -1);
82 #endif
83 /* added 1 extra reference to keep 2 dirty pages pinned per metric, expected refcnt = 2 */
85 - uv_mutex_lock(&descr->mutex);
84 + rrdeng_page_descr_mutex_lock(ctx, descr);
85 ret = pg_cache_try_get_unsafe(descr, 0);
87 - uv_mutex_unlock(&descr->mutex);
86 + rrdeng_page_descr_mutex_unlock(ctx, descr);
87 assert (1 == ret);
88
89 rrdeng_commit_page(ctx, descr, handle->page_correlation_id);
90 if (handle->prev_descr) {
91 /* unpin old second page */
93 - pg_cache_put(handle->prev_descr);
92 + pg_cache_put(ctx, handle->prev_descr);
93 }
94 handle->prev_descr = descr;
95 } else {
97 - free(descr->page);
96 + free(descr->pg_cache_descr->page);
97 + rrdeng_destroy_pg_cache_descr(ctx, descr->pg_cache_descr);
98 free(descr);
99 handle->descr = NULL;
100 }
101 }
102 - page = rrdeng_create_page(&handle->page_index->id, &descr);
102 + page = rrdeng_create_page(ctx, &handle->page_index->id, &descr);
103 assert(page);
104 handle->prev_descr = handle->descr;
105 handle->descr = descr;
106 - descr->handle = handle;
106 uv_rwlock_wrlock(&pg_cache->commited_page_index.lock);
107 handle->page_correlation_id = pg_cache->commited_page_index.latest_corr_id++;
108 uv_rwlock_wrunlock(&pg_cache->commited_page_index.lock);
109 }
111 - page = descr->page;
110 + page = descr->pg_cache_descr->page;
111
112 page[descr->page_length / sizeof(number)] = number;
113 descr->end_time = point_in_time;
@@ -132,13 +131,12 @@ void rrdeng_store_metric_finalize(RRDDIM *rd)
131 {
132 struct rrdeng_collect_handle *handle;
133 struct rrdengine_instance *ctx;
135 - struct rrdeng_page_cache_descr *descr;
134 + struct rrdeng_page_descr *descr;
135
136 handle = &rd->state->handle.rrdeng;
137 ctx = handle->ctx;
138 descr = handle->descr;
139 if (descr) {
141 - descr->handle = NULL;
140 if (descr->page_length) {
141 #ifdef NETDATA_INTERNAL_CHECKS
142 rrd_stat_atomic_add(&ctx->stats.metric_API_producers, -1);
@@ -146,10 +144,11 @@ void rrdeng_store_metric_finalize(RRDDIM *rd)
144 rrdeng_commit_page(ctx, descr, handle->page_correlation_id);
145 if (handle->prev_descr) {
146 /* unpin old second page */
149 - pg_cache_put(handle->prev_descr);
147 + pg_cache_put(ctx, handle->prev_descr);
148 }
149 } else {
152 - free(descr->page);
150 + free(descr->pg_cache_descr->page);
151 + rrdeng_destroy_pg_cache_descr(ctx, descr->pg_cache_descr);
152 free(descr);
153 }
154 }
@@ -180,7 +179,7 @@ storage_number rrdeng_load_metric_next(struct rrddim_query_handle *rrdimm_handle
179 {
180 struct rrdeng_query_handle *handle;
181 struct rrdengine_instance *ctx;
183 - struct rrdeng_page_cache_descr *descr;
182 + struct rrdeng_page_descr *descr;
183 storage_number *page, ret;
184 unsigned position;
185 usec_t point_in_time;
@@ -204,7 +203,7 @@ storage_number rrdeng_load_metric_next(struct rrddim_query_handle *rrdimm_handle
203 #ifdef NETDATA_INTERNAL_CHECKS
204 rrd_stat_atomic_add(&ctx->stats.metric_API_consumers, -1);
205 #endif
207 - pg_cache_put(descr);
206 + pg_cache_put(ctx, descr);
207 handle->descr = NULL;
208 }
209 descr = pg_cache_lookup(ctx, handle->page_index, &handle->page_index->id, point_in_time);
@@ -222,7 +221,7 @@ storage_number rrdeng_load_metric_next(struct rrddim_query_handle *rrdimm_handle
221 ret = SN_EMPTY_SLOT;
222 goto out;
223 }
225 - page = descr->page;
224 + page = descr->pg_cache_descr->page;
225 if (unlikely(descr->start_time == descr->end_time)) {
226 ret = page[0];
227 goto out;
@@ -254,7 +253,7 @@ void rrdeng_load_metric_finalize(struct rrddim_query_handle *rrdimm_handle)
253 {
254 struct rrdeng_query_handle *handle;
255 struct rrdengine_instance *ctx;
257 - struct rrdeng_page_cache_descr *descr;
256 + struct rrdeng_page_descr *descr;
257
258 handle = &rrdimm_handle->rrdeng;
259 ctx = handle->ctx;
@@ -263,7 +262,7 @@ void rrdeng_load_metric_finalize(struct rrddim_query_handle *rrdimm_handle)
262 #ifdef NETDATA_INTERNAL_CHECKS
263 rrd_stat_atomic_add(&ctx->stats.metric_API_consumers, -1);
264 #endif
266 - pg_cache_put(descr);
265 + pg_cache_put(ctx, descr);
266 }
267 }
268
@@ -289,28 +288,32 @@ time_t rrdeng_metric_oldest_time(RRDDIM *rd)
288 }
289
290 /* Also gets a reference for the page */
292 -void *rrdeng_create_page(uuid_t *id, struct rrdeng_page_cache_descr **ret_descr)
291 +void *rrdeng_create_page(struct rrdengine_instance *ctx, uuid_t *id, struct rrdeng_page_descr **ret_descr)
292 {
294 - struct rrdeng_page_cache_descr *descr;
293 + struct rrdeng_page_descr *descr;
294 + struct page_cache_descr *pg_cache_descr;
295 void *page;
296 /* TODO: check maximum number of pages in page cache limit */
297
298 - page = mallocz(RRDENG_BLOCK_SIZE); /*TODO: add page size */
298 descr = pg_cache_create_descr();
300 - descr->page = page;
299 descr->id = id; /* TODO: add page type: metric, log, something? */
302 - descr->flags = RRD_PAGE_DIRTY /*| RRD_PAGE_LOCKED */ | RRD_PAGE_POPULATED /* | BEING_COLLECTED */;
303 - descr->refcnt = 1;
300 + page = mallocz(RRDENG_BLOCK_SIZE); /*TODO: add page size */
301 + rrdeng_page_descr_mutex_lock(ctx, descr);
302 + pg_cache_descr = descr->pg_cache_descr;
303 + pg_cache_descr->page = page;
304 + pg_cache_descr->flags = RRD_PAGE_DIRTY /*| RRD_PAGE_LOCKED */ | RRD_PAGE_POPULATED /* | BEING_COLLECTED */;
305 + pg_cache_descr->refcnt = 1;
306
307 debug(D_RRDENGINE, "-----------------\nCreated new page:\n-----------------");
308 if(unlikely(debug_flags & D_RRDENGINE))
309 print_page_cache_descr(descr);
310 + rrdeng_page_descr_mutex_unlock(ctx, descr);
311 *ret_descr = descr;
312 return page;
313 }
314
315 /* The page must not be empty */
313 -void rrdeng_commit_page(struct rrdengine_instance *ctx, struct rrdeng_page_cache_descr *descr,
316 +void rrdeng_commit_page(struct rrdengine_instance *ctx, struct rrdeng_page_descr *descr,
317 Word_t page_correlation_id)
318 {
319 struct page_cache *pg_cache = &ctx->pg_cache;
@@ -328,13 +331,14 @@ void rrdeng_commit_page(struct rrdengine_instance *ctx, struct rrdeng_page_cache
331 ++pg_cache->commited_page_index.nr_commited_pages;
332 uv_rwlock_wrunlock(&pg_cache->commited_page_index.lock);
333
331 - pg_cache_put(descr);
334 + pg_cache_put(ctx, descr);
335 }
336
337 /* Gets a reference for the page */
338 void *rrdeng_get_latest_page(struct rrdengine_instance *ctx, uuid_t *id, void **handle)
339 {
337 - struct rrdeng_page_cache_descr *descr;
340 + struct rrdeng_page_descr *descr;
341 + struct page_cache_descr *pg_cache_descr;
342
343 debug(D_RRDENGINE, "----------------------\nReading existing page:\n----------------------");
344 descr = pg_cache_lookup(ctx, NULL, id, INVALID_TIME);
@@ -344,14 +348,16 @@ void *rrdeng_get_latest_page(struct rrdengine_instance *ctx, uuid_t *id, void **
348 return NULL;
349 }
350 *handle = descr;
351 + pg_cache_descr = descr->pg_cache_descr;
352
348 - return descr->page;
353 + return pg_cache_descr->page;
354 }
355
356 /* Gets a reference for the page */
357 void *rrdeng_get_page(struct rrdengine_instance *ctx, uuid_t *id, usec_t point_in_time, void **handle)
358 {
354 - struct rrdeng_page_cache_descr *descr;
359 + struct rrdeng_page_descr *descr;
360 + struct page_cache_descr *pg_cache_descr;
361
362 debug(D_RRDENGINE, "----------------------\nReading existing page:\n----------------------");
363 descr = pg_cache_lookup(ctx, NULL, id, point_in_time);
@@ -361,11 +367,12 @@ void *rrdeng_get_page(struct rrdengine_instance *ctx, uuid_t *id, usec_t point_i
367 return NULL;
368 }
369 *handle = descr;
370 + pg_cache_descr = descr->pg_cache_descr;
371
365 - return descr->page;
372 + return pg_cache_descr->page;
373 }
374
368 -void rrdeng_get_27_statistics(struct rrdengine_instance *ctx, unsigned long long *array)
375 +void rrdeng_get_28_statistics(struct rrdengine_instance *ctx, unsigned long long *array)
376 {
377 struct page_cache *pg_cache = &ctx->pg_cache;
378
@@ -396,13 +403,15 @@ void rrdeng_get_27_statistics(struct rrdengine_instance *ctx, unsigned long long
403 array[24] = (uint64_t)ctx->stats.datafile_deletions;
404 array[25] = (uint64_t)ctx->stats.journalfile_creations;
405 array[26] = (uint64_t)ctx->stats.journalfile_deletions;
406 + array[27] = (uint64_t)ctx->stats.page_cache_descriptors;
407 + assert(RRDENG_NR_STATS == 28);
408 }
409
410 /* Releases reference to page */
411 void rrdeng_put_page(struct rrdengine_instance *ctx, void *handle)
412 {
413 (void)ctx;
405 - pg_cache_put((struct rrdeng_page_cache_descr *)handle);
414 + pg_cache_put(ctx, (struct rrdeng_page_descr *)handle);
415 }
416
417 /*
database/engine/rrdengineapi.h
+6 -3
@@ -7,11 +7,14 @@
7
8 #define RRDENG_MIN_PAGE_CACHE_SIZE_MB (32)
9 #define RRDENG_MIN_DISK_SPACE_MB (256)
10 +
11 +#define RRDENG_NR_STATS (28)
12 +
13 extern int default_rrdeng_page_cache_mb;
14 extern int default_rrdeng_disk_quota_mb;
15
13 -extern void *rrdeng_create_page(uuid_t *id, struct rrdeng_page_cache_descr **ret_descr);
14 -extern void rrdeng_commit_page(struct rrdengine_instance *ctx, struct rrdeng_page_cache_descr *descr,
16 +extern void *rrdeng_create_page(struct rrdengine_instance *ctx, uuid_t *id, struct rrdeng_page_descr **ret_descr);
17 +extern void rrdeng_commit_page(struct rrdengine_instance *ctx, struct rrdeng_page_descr *descr,
18 Word_t page_correlation_id);
19 extern void *rrdeng_get_latest_page(struct rrdengine_instance *ctx, uuid_t *id, void **handle);
20 extern void *rrdeng_get_page(struct rrdengine_instance *ctx, uuid_t *id, usec_t point_in_time, void **handle);
@@ -26,7 +29,7 @@ extern int rrdeng_load_metric_is_finished(struct rrddim_query_handle *rrdimm_han
29 extern void rrdeng_load_metric_finalize(struct rrddim_query_handle *rrdimm_handle);
30 extern time_t rrdeng_metric_latest_time(RRDDIM *rd);
31 extern time_t rrdeng_metric_oldest_time(RRDDIM *rd);
29 -extern void rrdeng_get_27_statistics(struct rrdengine_instance *ctx, unsigned long long *array);
32 +extern void rrdeng_get_28_statistics(struct rrdengine_instance *ctx, unsigned long long *array);
33
34 /* must call once before using anything */
35 extern int rrdeng_init(struct rrdengine_instance **ctxp, char *dbfiles_path, unsigned page_cache_mb,
database/engine/rrdenginelib.c
+41 -13
@@ -1,25 +1,51 @@
1 // SPDX-License-Identifier: GPL-3.0-or-later
2 #include "rrdengine.h"
3
4 -void print_page_cache_descr(struct rrdeng_page_cache_descr *page_cache_descr)
4 +#define BUFSIZE (512)
5 +
6 +/* Caller must hold descriptor lock */
7 +void print_page_cache_descr(struct rrdeng_page_descr *descr)
8 {
6 - char uuid_str[37];
7 - char str[512];
9 + struct page_cache_descr *pg_cache_descr = descr->pg_cache_descr;
10 + char uuid_str[UUID_STR_LEN];
11 + char str[BUFSIZE];
12 int pos = 0;
13
10 - uuid_unparse_lower(*page_cache_descr->id, uuid_str);
11 - pos += snprintfz(str, 512 - pos, "page(%p) id=%s\n"
14 + uuid_unparse_lower(*descr->id, uuid_str);
15 + pos += snprintfz(str, BUFSIZE - pos, "page(%p) id=%s\n"
16 "--->len:%"PRIu32" time:%"PRIu64"->%"PRIu64" xt_offset:",
13 - page_cache_descr->page, uuid_str,
14 - page_cache_descr->page_length,
15 - (uint64_t)page_cache_descr->start_time,
16 - (uint64_t)page_cache_descr->end_time);
17 - if (!page_cache_descr->extent) {
18 - pos += snprintfz(str + pos, 512 - pos, "N/A");
17 + pg_cache_descr->page, uuid_str,
18 + descr->page_length,
19 + (uint64_t)descr->start_time,
20 + (uint64_t)descr->end_time);
21 + if (!descr->extent) {
22 + pos += snprintfz(str + pos, BUFSIZE - pos, "N/A");
23 + } else {
24 + pos += snprintfz(str + pos, BUFSIZE - pos, "%"PRIu64, descr->extent->offset);
25 + }
26 + snprintfz(str + pos, BUFSIZE - pos, " flags:0x%2.2lX refcnt:%u\n\n", pg_cache_descr->flags, pg_cache_descr->refcnt);
27 + fputs(str, stderr);
28 +}
29 +
30 +void print_page_descr(struct rrdeng_page_descr *descr)
31 +{
32 + char uuid_str[UUID_STR_LEN];
33 + char str[BUFSIZE];
34 + int pos = 0;
35 +
36 + uuid_unparse_lower(*descr->id, uuid_str);
37 + pos += snprintfz(str, BUFSIZE - pos, "id=%s\n"
38 + "--->len:%"PRIu32" time:%"PRIu64"->%"PRIu64" xt_offset:",
39 + uuid_str,
40 + descr->page_length,
41 + (uint64_t)descr->start_time,
42 + (uint64_t)descr->end_time);
43 + if (!descr->extent) {
44 + pos += snprintfz(str + pos, BUFSIZE - pos, "N/A");
45 } else {
20 - pos += snprintfz(str + pos, 512 - pos, "%"PRIu64, page_cache_descr->extent->offset);
46 + pos += snprintfz(str + pos, BUFSIZE - pos, "%"PRIu64, descr->extent->offset);
47 }
22 - snprintfz(str + pos, 512 - pos, " flags:0x%2.2lX refcnt:%u\n\n", page_cache_descr->flags, page_cache_descr->refcnt);
48 + snprintfz(str + pos, BUFSIZE - pos, "\n\n");
49 fputs(str, stderr);
50 }
51
@@ -60,6 +86,7 @@ char *get_rrdeng_statistics(struct rrdengine_instance *ctx, char *str, size_t si
86 "metric_API_producers: %ld\n"
87 "metric_API_consumers: %ld\n"
88 "page_cache_total_pages: %ld\n"
89 + "page_cache_descriptors: %ld\n"
90 "page_cache_populated_pages: %ld\n"
91 "page_cache_commited_pages: %ld\n"
92 "page_cache_insertions: %ld\n"
@@ -87,6 +114,7 @@ char *get_rrdeng_statistics(struct rrdengine_instance *ctx, char *str, size_t si
114 (long)ctx->stats.metric_API_producers,
115 (long)ctx->stats.metric_API_consumers,
116 (long)pg_cache->page_descriptors,
117 + (long)ctx->stats.page_cache_descriptors,
118 (long)pg_cache->populated_pages,
119 (long)pg_cache->commited_page_index.nr_commited_pages,
120 (long)ctx->stats.pg_cache_insertions,
database/engine/rrdenginelib.h
+16 -2
@@ -6,11 +6,17 @@
6 #include "rrdengine.h"
7
8 /* Forward declarations */
9 -struct rrdeng_page_cache_descr;
9 +struct rrdeng_page_descr;
10
11 #define STR_HELPER(x) #x
12 #define STR(x) STR_HELPER(x)
13
14 +#define BITS_PER_ULONG (sizeof(unsigned long) * 8)
15 +
16 +#ifndef UUID_STR_LEN
17 +#define UUID_STR_LEN (37)
18 +#endif
19 +
20 /* Taken from linux kernel */
21 #define BUILD_BUG_ON(condition) ((void)sizeof(char[1 - 2*!!(condition)]))
22
@@ -25,6 +31,13 @@ typedef uintptr_t rrdeng_stats_t;
31 #define rrd_stat_atomic_add(p, n) do {(void) __sync_fetch_and_add(p, n);} while(0)
32 #endif
33
34 +/* returns old *ptr value */
35 +static inline unsigned long ulong_compare_and_swap(volatile unsigned long *ptr,
36 + unsigned long oldval, unsigned long newval)
37 +{
38 + return __sync_val_compare_and_swap(ptr, oldval, newval);
39 +}
40 +
41 #ifndef O_DIRECT
42 /* Workaround for OS X */
43 #define O_DIRECT (0)
@@ -77,7 +90,8 @@ static inline void crc32set(void *crcp, uLong crc)
90 *(uint32_t *)crcp = crc;
91 }
92
80 -extern void print_page_cache_descr(struct rrdeng_page_cache_descr *page_cache_descr);
93 +extern void print_page_cache_descr(struct rrdeng_page_descr *page_cache_descr);
94 +extern void print_page_descr(struct rrdeng_page_descr *descr);
95 extern int check_file_properties(uv_file file, uint64_t *file_size, size_t min_size);
96 extern char *get_rrdeng_statistics(struct rrdengine_instance *ctx, char *str, size_t size);
97
database/engine/rrdenglocking.c new
+209
@@ -0,0 +1,209 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +#include "rrdengine.h"
3 +
4 +struct page_cache_descr *rrdeng_create_pg_cache_descr(struct rrdengine_instance *ctx)
5 +{
6 + struct page_cache_descr *pg_cache_descr;
7 +
8 + pg_cache_descr = mallocz(sizeof(*pg_cache_descr));
9 + rrd_stat_atomic_add(&ctx->stats.page_cache_descriptors, 1);
10 + pg_cache_descr->page = NULL;
11 + pg_cache_descr->flags = 0;
12 + pg_cache_descr->prev = pg_cache_descr->next = NULL;
13 + pg_cache_descr->refcnt = 0;
14 + pg_cache_descr->waiters = 0;
15 + assert(0 == uv_cond_init(&pg_cache_descr->cond));
16 + assert(0 == uv_mutex_init(&pg_cache_descr->mutex));
17 +
18 + return pg_cache_descr;
19 +}
20 +
21 +void rrdeng_destroy_pg_cache_descr(struct rrdengine_instance *ctx, struct page_cache_descr *pg_cache_descr)
22 +{
23 + uv_cond_destroy(&pg_cache_descr->cond);
24 + uv_mutex_destroy(&pg_cache_descr->mutex);
25 + free(pg_cache_descr);
26 + rrd_stat_atomic_add(&ctx->stats.page_cache_descriptors, -1);
27 +}
28 +
29 +/* also allocates page cache descriptor if missing */
30 +void rrdeng_page_descr_mutex_lock(struct rrdengine_instance *ctx, struct rrdeng_page_descr *descr)
31 +{
32 + unsigned long old_state, old_users, new_state, ret_state;
33 + struct page_cache_descr *pg_cache_descr = NULL;
34 + uint8_t we_locked;
35 +
36 + we_locked = 0;
37 + while (1) { /* spin */
38 + old_state = descr->pg_cache_descr_state;
39 + old_users = old_state >> PG_CACHE_DESCR_SHIFT;
40 +
41 + if (unlikely(we_locked)) {
42 + assert(old_state & PG_CACHE_DESCR_LOCKED);
43 + new_state = (1 << PG_CACHE_DESCR_SHIFT) | (old_state & PG_CACHE_DESCR_FLAGS_MASK);
44 + new_state &= ~PG_CACHE_DESCR_LOCKED;
45 + new_state |= PG_CACHE_DESCR_ALLOCATED;
46 + ret_state = ulong_compare_and_swap(&descr->pg_cache_descr_state, old_state, new_state);
47 + if (old_state == ret_state) {
48 + /* success */
49 + break;
50 + }
51 + continue; /* spin */
52 + }
53 + if (old_state & PG_CACHE_DESCR_LOCKED) {
54 + assert(0 == old_users);
55 + continue; /* spin */
56 + }
57 + if (0 == old_state) {
58 + /* no page cache descriptor has been allocated */
59 +
60 + if (NULL == pg_cache_descr) {
61 + pg_cache_descr = rrdeng_create_pg_cache_descr(ctx);
62 + }
63 + new_state = PG_CACHE_DESCR_LOCKED;
64 + ret_state = ulong_compare_and_swap(&descr->pg_cache_descr_state, 0, new_state);
65 + if (0 == ret_state) {
66 + we_locked = 1;
67 + descr->pg_cache_descr = pg_cache_descr;
68 + pg_cache_descr->descr = descr;
69 + pg_cache_descr = NULL; /* make sure we don't free pg_cache_descr */
70 + /* retry */
71 + continue;
72 + }
73 + continue; /* spin */
74 + }
75 + /* page cache descriptor is already allocated */
76 + assert(old_state & PG_CACHE_DESCR_ALLOCATED);
77 +
78 + new_state = (old_users + 1) << PG_CACHE_DESCR_SHIFT;
79 + new_state |= old_state & PG_CACHE_DESCR_FLAGS_MASK;
80 +
81 + ret_state = ulong_compare_and_swap(&descr->pg_cache_descr_state, old_state, new_state);
82 + if (old_state == ret_state) {
83 + /* success */
84 + break;
85 + }
86 + /* spin */
87 + }
88 +
89 + if (pg_cache_descr) {
90 + free(pg_cache_descr);
91 + }
92 + pg_cache_descr = descr->pg_cache_descr;
93 + uv_mutex_lock(&pg_cache_descr->mutex);
94 +}
95 +
96 +void rrdeng_page_descr_mutex_unlock(struct rrdengine_instance *ctx, struct rrdeng_page_descr *descr)
97 +{
98 + unsigned long old_state, new_state, ret_state, old_users;
99 + struct page_cache_descr *pg_cache_descr;
100 + uint8_t we_locked;
101 +
102 + uv_mutex_unlock(&descr->pg_cache_descr->mutex);
103 +
104 + we_locked = 0;
105 + while (1) { /* spin */
106 + old_state = descr->pg_cache_descr_state;
107 + assert(old_state & PG_CACHE_DESCR_ALLOCATED);
108 + old_users = old_state >> PG_CACHE_DESCR_SHIFT;
109 +
110 + if (unlikely(we_locked)) {
111 + assert(0 == old_users);
112 +
113 + ret_state = ulong_compare_and_swap(&descr->pg_cache_descr_state, old_state, 0);
114 + if (old_state == ret_state) {
115 + /* success */
116 + break;
117 + }
118 + continue; /* spin */
119 + }
120 + if (old_state & PG_CACHE_DESCR_LOCKED) {
121 + assert(0 == old_users);
122 + continue; /* spin */
123 + }
124 + pg_cache_descr = descr->pg_cache_descr;
125 + /* caller is the only page cache descriptor user and there are no pending references on the page */
126 + if ((old_state & PG_CACHE_DESCR_DESTROY) && (1 == old_users) &&
127 + !pg_cache_descr->flags && !pg_cache_descr->refcnt) {
128 + new_state = PG_CACHE_DESCR_LOCKED;
129 + ret_state = ulong_compare_and_swap(&descr->pg_cache_descr_state, old_state, new_state);
130 + if (old_state == ret_state) {
131 + we_locked = 1;
132 + rrdeng_destroy_pg_cache_descr(ctx, pg_cache_descr);
133 + /* retry */
134 + continue;
135 + }
136 + continue; /* spin */
137 + }
138 + assert(old_users > 0);
139 + new_state = (old_users - 1) << PG_CACHE_DESCR_SHIFT;
140 + new_state |= old_state & PG_CACHE_DESCR_FLAGS_MASK;
141 +
142 + ret_state = ulong_compare_and_swap(&descr->pg_cache_descr_state, old_state, new_state);
143 + if (old_state == ret_state) {
144 + /* success */
145 + break;
146 + }
147 + /* spin */
148 + }
149 +
150 +}
151 +
152 +/*
153 + * Tries to deallocate page cache descriptor. If it fails, it postpones deallocation by setting the
154 + * PG_CACHE_DESCR_DESTROY flag which will be eventually cleared by a different context after doing
155 + * the deallocation.
156 + */
157 +void rrdeng_try_deallocate_pg_cache_descr(struct rrdengine_instance *ctx, struct rrdeng_page_descr *descr)
158 +{
159 + unsigned long old_state, new_state, ret_state, old_users;
160 + struct page_cache_descr *pg_cache_descr;
161 + uint8_t we_locked;
162 +
163 + we_locked = 0;
164 + while (1) { /* spin */
165 + old_state = descr->pg_cache_descr_state;
166 + old_users = old_state >> PG_CACHE_DESCR_SHIFT;
167 +
168 + if (unlikely(we_locked)) {
169 + assert(0 == old_users);
170 +
171 + ret_state = ulong_compare_and_swap(&descr->pg_cache_descr_state, old_state, 0);
172 + if (old_state == ret_state) {
173 + /* success */
174 + break;
175 + }
176 + continue; /* spin */
177 + }
178 + if (!(old_state & PG_CACHE_DESCR_ALLOCATED) || (old_state & PG_CACHE_DESCR_DESTROY)) {
179 + /* don't do anything */
180 + return;
181 + }
182 + if (old_state & PG_CACHE_DESCR_LOCKED) {
183 + assert(0 == old_users);
184 + continue; /* spin */
185 + }
186 + pg_cache_descr = descr->pg_cache_descr;
187 + /* caller is the only page cache descriptor user and there are no pending references on the page */
188 + if ((0 == old_users) && !pg_cache_descr->flags && !pg_cache_descr->refcnt) {
189 + new_state = PG_CACHE_DESCR_LOCKED;
190 + ret_state = ulong_compare_and_swap(&descr->pg_cache_descr_state, old_state, new_state);
191 + if (old_state == ret_state) {
192 + we_locked = 1;
193 + rrdeng_destroy_pg_cache_descr(ctx, pg_cache_descr);
194 + /* retry */
195 + continue;
196 + }
197 + continue; /* spin */
198 + }
199 + /* plant PG_CACHE_DESCR_DESTROY so that other contexts eventually free the page cache descriptor */
200 + new_state = old_state | PG_CACHE_DESCR_DESTROY;
201 +
202 + ret_state = ulong_compare_and_swap(&descr->pg_cache_descr_state, old_state, new_state);
203 + if (old_state == ret_state) {
204 + /* success */
205 + break;
206 + }
207 + /* spin */
208 + }
209 +}
\ No newline at end of file
database/engine/rrdenglocking.h new
+17
@@ -0,0 +1,17 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +#ifndef NETDATA_RRDENGLOCKING_H
4 +#define NETDATA_RRDENGLOCKING_H
5 +
6 +#include "rrdengine.h"
7 +
8 +/* Forward declarations */
9 +struct page_cache_descr;
10 +
11 +extern struct page_cache_descr *rrdeng_create_pg_cache_descr(struct rrdengine_instance *ctx);
12 +extern void rrdeng_destroy_pg_cache_descr(struct rrdengine_instance *ctx, struct page_cache_descr *pg_cache_descr);
13 +extern void rrdeng_page_descr_mutex_lock(struct rrdengine_instance *ctx, struct rrdeng_page_descr *descr);
14 +extern void rrdeng_page_descr_mutex_unlock(struct rrdengine_instance *ctx, struct rrdeng_page_descr *descr);
15 +extern void rrdeng_try_deallocate_pg_cache_descr(struct rrdengine_instance *ctx, struct rrdeng_page_descr *descr);
16 +
17 +#endif /* NETDATA_RRDENGLOCKING_H */
\ No newline at end of file
database/rrd.h
+3 -3
@@ -17,7 +17,7 @@ typedef struct alarm_entry ALARM_ENTRY;
17 // forward declarations
18 struct rrddim_volatile;
19 #ifdef ENABLE_DBENGINE
20 -struct rrdeng_page_cache_descr;
20 +struct rrdeng_page_descr;
21 struct rrdengine_instance;
22 struct pg_cache_page_index;
23 #endif
@@ -246,7 +246,7 @@ union rrddim_collect_handle {
246 } slotted; // state the legacy code uses
247 #ifdef ENABLE_DBENGINE
248 struct rrdeng_collect_handle {
249 - struct rrdeng_page_cache_descr *descr, *prev_descr;
249 + struct rrdeng_page_descr *descr, *prev_descr;
250 unsigned long page_correlation_id;
251 struct rrdengine_instance *ctx;
252 struct pg_cache_page_index *page_index;
@@ -268,7 +268,7 @@ struct rrddim_query_handle {
268 } slotted; // state the legacy code uses
269 #ifdef ENABLE_DBENGINE
270 struct rrdeng_query_handle {
271 - struct rrdeng_page_cache_descr *descr;
271 + struct rrdeng_page_descr *descr;
272 struct rrdengine_instance *ctx;
273 struct pg_cache_page_index *page_index;
274 time_t now; //TODO: remove now to implement next point iteration