26
size_t flushes_running;
27
size_t evictions_running;
28
size_t cleanup_running;
29
+
30
+ struct {
31
+ ARAL *ar;
32
+
33
+ struct {
34
+ SPINLOCK spinlock;
35
+
36
+ size_t waiting;
37
+ struct rrdeng_cmd *waiting_items_by_priority[STORAGE_PRIORITY_INTERNAL_MAX_DONT_USE];
38
+ size_t executed_by_priority[STORAGE_PRIORITY_INTERNAL_MAX_DONT_USE];
39
+ } unsafe;
40
+ } cmd_queue;
41
+
42
+ struct {
43
+ ARAL *ar;
44
+
45
+ struct {
46
+ size_t dispatched;
47
+ size_t executing;
48
+ size_t pending_cb;
49
+ } atomics;
50
+ } work_cmd;
51
+
52
+ struct {
53
+ ARAL *ar;
54
+ } handles;
55
+
56
+ struct {
57
+ ARAL *ar;
58
+ } descriptors;
59
+
60
+ struct {
61
+ ARAL *ar;
62
+ } xt_io_descr;
63
+
64
} rrdeng_main = {
65
.thread = 0,
66
.loop = {},
69
.flushes_running = 0,
70
.evictions_running = 0,
71
.cleanup_running = 0,
72
+
73
+ .cmd_queue = {
74
+ .unsafe = {
75
+ .spinlock = NETDATA_SPINLOCK_INITIALIZER,
76
+ },
77
+ }
78
};
79
80
static void sanity_check(void)
120
work_cb work_cb;
121
after_work_cb after_work_cb;
122
enum rrdeng_opcode opcode;
82
-
83
- struct {
84
- struct rrdeng_work *prev;
85
- struct rrdeng_work *next;
86
- } cache;
123
};
124
89
-static struct {
90
- struct {
91
- SPINLOCK spinlock;
92
- struct rrdeng_work *available_items;
93
- size_t available;
94
- } protected;
95
-
96
- struct {
97
- size_t allocated;
98
- size_t dispatched;
99
- size_t executing;
100
- size_t pending_cb;
101
- } atomics;
102
-} work_request_globals = {
103
- .protected = {
104
- .spinlock = NETDATA_SPINLOCK_INITIALIZER,
105
- .available_items = NULL,
106
- .available = 0,
107
- },
108
- .atomics = {
109
- .allocated = 0,
110
- .dispatched = 0,
111
- .executing = 0,
112
- },
113
-};
114
-
115
-static inline bool work_request_full(void) {
116
- return __atomic_load_n(&work_request_globals.atomics.dispatched, __ATOMIC_RELAXED) >= (size_t)(libuv_worker_threads - RESERVED_LIBUV_WORKER_THREADS);
125
+static void work_request_init(void) {
126
+ rrdeng_main.work_cmd.ar = aral_create(
127
+ "dbengine-work-cmd",
128
+ sizeof(struct rrdeng_work),
129
+ 0,
130
+ 65536, NULL,
131
+ NULL, NULL, false, false
132
+ );
133
}
134
119
-static void work_request_cleanup1(void) {
120
- struct rrdeng_work *item = NULL;
121
-
122
- if(!netdata_spinlock_trylock(&work_request_globals.protected.spinlock))
123
- return;
124
-
125
- if(work_request_globals.protected.available_items && work_request_globals.protected.available > (size_t)libuv_worker_threads) {
126
- item = work_request_globals.protected.available_items;
127
- DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(work_request_globals.protected.available_items, item, cache.prev, cache.next);
128
- work_request_globals.protected.available--;
129
- }
130
- netdata_spinlock_unlock(&work_request_globals.protected.spinlock);
131
-
132
- if(item) {
133
- freez(item);
134
- __atomic_sub_fetch(&work_request_globals.atomics.allocated, 1, __ATOMIC_RELAXED);
135
- }
135
+static inline bool work_request_full(void) {
136
+ return __atomic_load_n(&rrdeng_main.work_cmd.atomics.dispatched, __ATOMIC_RELAXED) >= (size_t)(libuv_worker_threads - RESERVED_LIBUV_WORKER_THREADS);
137
}
138
139
static inline void work_done(struct rrdeng_work *work_request) {
139
- netdata_spinlock_lock(&work_request_globals.protected.spinlock);
140
- DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(work_request_globals.protected.available_items, work_request, cache.prev, cache.next);
141
- work_request_globals.protected.available++;
142
- netdata_spinlock_unlock(&work_request_globals.protected.spinlock);
140
+ aral_freez(rrdeng_main.work_cmd.ar, work_request);
141
}
142
143
static void work_standard_worker(uv_work_t *req) {
146
- __atomic_add_fetch(&work_request_globals.atomics.executing, 1, __ATOMIC_RELAXED);
144
+ __atomic_add_fetch(&rrdeng_main.work_cmd.atomics.executing, 1, __ATOMIC_RELAXED);
145
146
register_libuv_worker_jobs();
147
worker_is_busy(UV_EVENT_WORKER_INIT);
150
work_request->data = work_request->work_cb(work_request->ctx, work_request->data, work_request->completion, req);
151
worker_is_idle();
152
155
- __atomic_sub_fetch(&work_request_globals.atomics.dispatched, 1, __ATOMIC_RELAXED);
156
- __atomic_sub_fetch(&work_request_globals.atomics.executing, 1, __ATOMIC_RELAXED);
157
- __atomic_add_fetch(&work_request_globals.atomics.pending_cb, 1, __ATOMIC_RELAXED);
153
+ __atomic_sub_fetch(&rrdeng_main.work_cmd.atomics.dispatched, 1, __ATOMIC_RELAXED);
154
+ __atomic_sub_fetch(&rrdeng_main.work_cmd.atomics.executing, 1, __ATOMIC_RELAXED);
155
+ __atomic_add_fetch(&rrdeng_main.work_cmd.atomics.pending_cb, 1, __ATOMIC_RELAXED);
156
157
// signal the event loop a worker is available
158
fatal_assert(0 == uv_async_send(&rrdeng_main.async));
167
work_request->after_work_cb(work_request->ctx, work_request->data, work_request->completion, req, status);
168
169
work_done(work_request);
172
- __atomic_sub_fetch(&work_request_globals.atomics.pending_cb, 1, __ATOMIC_RELAXED);
170
+ __atomic_sub_fetch(&rrdeng_main.work_cmd.atomics.pending_cb, 1, __ATOMIC_RELAXED);
171
172
worker_is_idle();
173
}
177
178
internal_fatal(rrdeng_main.tid != gettid(), "work_dispatch() can only be run from the event loop thread");
179
182
- netdata_spinlock_lock(&work_request_globals.protected.spinlock);
183
-
184
- if(likely(work_request_globals.protected.available_items)) {
185
- work_request = work_request_globals.protected.available_items;
186
- DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(work_request_globals.protected.available_items, work_request, cache.prev, cache.next);
187
- work_request_globals.protected.available--;
188
- }
189
-
190
- netdata_spinlock_unlock(&work_request_globals.protected.spinlock);
191
-
192
- if(unlikely(!work_request)) {
193
- work_request = mallocz(sizeof(struct rrdeng_work));
194
- __atomic_add_fetch(&work_request_globals.atomics.allocated, 1, __ATOMIC_RELAXED);
195
- }
196
-
180
+ work_request = aral_mallocz(rrdeng_main.work_cmd.ar);
181
memset(work_request, 0, sizeof(struct rrdeng_work));
182
work_request->req.data = work_request;
183
work_request->ctx = ctx;
193
return false;
194
}
195
212
- __atomic_add_fetch(&work_request_globals.atomics.dispatched, 1, __ATOMIC_RELAXED);
196
+ __atomic_add_fetch(&rrdeng_main.work_cmd.atomics.dispatched, 1, __ATOMIC_RELAXED);
197
198
return true;
199
}
201
// ----------------------------------------------------------------------------
202
// page descriptor cache
203
220
-static struct {
221
- struct {
222
- SPINLOCK spinlock;
223
- struct page_descr_with_data *available_items;
224
- size_t available;
225
- } protected;
226
-
227
- struct {
228
- size_t allocated;
229
- } atomics;
230
-} page_descriptor_globals = {
231
- .protected = {
232
- .spinlock = NETDATA_SPINLOCK_INITIALIZER,
233
- .available_items = NULL,
234
- .available = 0,
235
- },
236
- .atomics = {
237
- .allocated = 0,
238
- },
239
-};
240
-
241
-static void page_descriptor_cleanup1(void) {
242
- struct page_descr_with_data *item = NULL;
243
-
244
- if(!netdata_spinlock_trylock(&page_descriptor_globals.protected.spinlock))
245
- return;
246
-
247
- if(page_descriptor_globals.protected.available_items && page_descriptor_globals.protected.available > MAX_PAGES_PER_EXTENT) {
248
- item = page_descriptor_globals.protected.available_items;
249
- DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(page_descriptor_globals.protected.available_items, item, cache.prev, cache.next);
250
- page_descriptor_globals.protected.available--;
251
- }
252
-
253
- netdata_spinlock_unlock(&page_descriptor_globals.protected.spinlock);
254
-
255
- if(item) {
256
- freez(item);
257
- __atomic_sub_fetch(&page_descriptor_globals.atomics.allocated, 1, __ATOMIC_RELAXED);
258
- }
204
+void page_descriptors_init(void) {
205
+ rrdeng_main.descriptors.ar = aral_create(
206
+ "dbengine-descriptors",
207
+ sizeof(struct page_descr_with_data),
208
+ 0,
209
+ 65536 * 4,
210
+ NULL,
211
+ NULL, NULL, false, false);
212
}
213
214
struct page_descr_with_data *page_descriptor_get(void) {
262
- struct page_descr_with_data *descr = NULL;
263
-
264
- netdata_spinlock_lock(&page_descriptor_globals.protected.spinlock);
265
-
266
- if(likely(page_descriptor_globals.protected.available_items)) {
267
- descr = page_descriptor_globals.protected.available_items;
268
- DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(page_descriptor_globals.protected.available_items, descr, cache.prev, cache.next);
269
- page_descriptor_globals.protected.available--;
270
- }
271
-
272
- netdata_spinlock_unlock(&page_descriptor_globals.protected.spinlock);
273
-
274
- if(unlikely(!descr)) {
275
- descr = mallocz(sizeof(struct page_descr_with_data));
276
- __atomic_add_fetch(&page_descriptor_globals.atomics.allocated, 1, __ATOMIC_RELAXED);
277
- }
278
-
215
+ struct page_descr_with_data *descr = aral_mallocz(rrdeng_main.descriptors.ar);
216
memset(descr, 0, sizeof(struct page_descr_with_data));
217
return descr;
218
}
219
220
static inline void page_descriptor_release(struct page_descr_with_data *descr) {
284
- if(unlikely(!descr)) return;
285
-
286
- netdata_spinlock_lock(&page_descriptor_globals.protected.spinlock);
287
- DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(page_descriptor_globals.protected.available_items, descr, cache.prev, cache.next);
288
- page_descriptor_globals.protected.available++;
289
- netdata_spinlock_unlock(&page_descriptor_globals.protected.spinlock);
221
+ aral_freez(rrdeng_main.descriptors.ar, descr);
222
}
223
224
// ----------------------------------------------------------------------------
225
// extent io descriptor cache
226
295
-static struct {
296
- struct {
297
- SPINLOCK spinlock;
298
- struct extent_io_descriptor *available_items;
299
- size_t available;
300
- } protected;
301
-
302
- struct {
303
- size_t allocated;
304
- } atomics;
305
-
306
-} extent_io_descriptor_globals = {
307
- .protected = {
308
- .spinlock = NETDATA_SPINLOCK_INITIALIZER,
309
- .available_items = NULL,
310
- .available = 0,
311
- },
312
- .atomics = {
313
- .allocated = 0,
314
- },
315
-};
316
-
317
-static void extent_io_descriptor_cleanup1(void) {
318
- struct extent_io_descriptor *item = NULL;
319
-
320
- if(!netdata_spinlock_trylock(&extent_io_descriptor_globals.protected.spinlock))
321
- return;
322
-
323
- if(extent_io_descriptor_globals.protected.available_items && extent_io_descriptor_globals.protected.available > (size_t)libuv_worker_threads) {
324
- item = extent_io_descriptor_globals.protected.available_items;
325
- DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(extent_io_descriptor_globals.protected.available_items, item, cache.prev, cache.next);
326
- extent_io_descriptor_globals.protected.available--;
327
- }
328
- netdata_spinlock_unlock(&extent_io_descriptor_globals.protected.spinlock);
329
-
330
- if(item) {
331
- freez(item);
332
- __atomic_sub_fetch(&extent_io_descriptor_globals.atomics.allocated, 1, __ATOMIC_RELAXED);
333
- }
227
+static void extent_io_descriptor_init(void) {
228
+ rrdeng_main.xt_io_descr.ar = aral_create(
229
+ "dbengine-extent-io",
230
+ sizeof(struct extent_io_descriptor),
231
+ 0,
232
+ 65536,
233
+ NULL,
234
+ NULL, NULL, false, false
235
+ );
236
}
237
238
static struct extent_io_descriptor *extent_io_descriptor_get(void) {
337
- struct extent_io_descriptor *xt_io_descr = NULL;
338
-
339
- netdata_spinlock_lock(&extent_io_descriptor_globals.protected.spinlock);
340
-
341
- if(likely(extent_io_descriptor_globals.protected.available_items)) {
342
- xt_io_descr = extent_io_descriptor_globals.protected.available_items;
343
- DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(extent_io_descriptor_globals.protected.available_items, xt_io_descr, cache.prev, cache.next);
344
- extent_io_descriptor_globals.protected.available--;
345
- }
346
-
347
- netdata_spinlock_unlock(&extent_io_descriptor_globals.protected.spinlock);
348
-
349
- if(unlikely(!xt_io_descr)) {
350
- xt_io_descr = mallocz(sizeof(struct extent_io_descriptor));
351
- __atomic_add_fetch(&extent_io_descriptor_globals.atomics.allocated, 1, __ATOMIC_RELAXED);
352
- }
353
-
239
+ struct extent_io_descriptor *xt_io_descr = aral_mallocz(rrdeng_main.xt_io_descr.ar);
240
memset(xt_io_descr, 0, sizeof(struct extent_io_descriptor));
241
return xt_io_descr;
242
}
243
244
static inline void extent_io_descriptor_release(struct extent_io_descriptor *xt_io_descr) {
359
- if(unlikely(!xt_io_descr)) return;
360
-
361
- netdata_spinlock_lock(&extent_io_descriptor_globals.protected.spinlock);
362
- DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(extent_io_descriptor_globals.protected.available_items, xt_io_descr, cache.prev, cache.next);
363
- extent_io_descriptor_globals.protected.available++;
364
- netdata_spinlock_unlock(&extent_io_descriptor_globals.protected.spinlock);
245
+ aral_freez(rrdeng_main.xt_io_descr.ar, xt_io_descr);
246
}
247
248
// ----------------------------------------------------------------------------
249
// query handle cache
250
370
-static struct {
371
- struct {
372
- SPINLOCK spinlock;
373
- struct rrdeng_query_handle *available_items;
374
- size_t available;
375
- } protected;
376
-
377
- struct {
378
- size_t allocated;
379
- } atomics;
380
-} rrdeng_query_handle_globals = {
381
- .protected = {
382
- .spinlock = NETDATA_SPINLOCK_INITIALIZER,
383
- .available_items = NULL,
384
- .available = 0,
385
- },
386
- .atomics = {
387
- .allocated = 0,
388
- },
389
-};
390
-
391
-static void rrdeng_query_handle_cleanup1(void) {
392
- struct rrdeng_query_handle *item = NULL;
393
-
394
- if(!netdata_spinlock_trylock(&rrdeng_query_handle_globals.protected.spinlock))
395
- return;
396
-
397
- if(rrdeng_query_handle_globals.protected.available_items && rrdeng_query_handle_globals.protected.available > 10) {
398
- item = rrdeng_query_handle_globals.protected.available_items;
399
- DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(rrdeng_query_handle_globals.protected.available_items, item, cache.prev, cache.next);
400
- rrdeng_query_handle_globals.protected.available--;
401
- }
402
-
403
- netdata_spinlock_unlock(&rrdeng_query_handle_globals.protected.spinlock);
404
-
405
- if(item) {
406
- freez(item);
407
- __atomic_sub_fetch(&rrdeng_query_handle_globals.atomics.allocated, 1, __ATOMIC_RELAXED);
408
- }
251
+void rrdeng_query_handle_init(void) {
252
+ rrdeng_main.handles.ar = aral_create(
253
+ "dbengine-query-handles",
254
+ sizeof(struct rrdeng_query_handle),
255
+ 0,
256
+ 65536,
257
+ NULL,
258
+ NULL, NULL, false, false);
259
}
260
261
struct rrdeng_query_handle *rrdeng_query_handle_get(void) {
412
- struct rrdeng_query_handle *handle = NULL;
413
-
414
- netdata_spinlock_lock(&rrdeng_query_handle_globals.protected.spinlock);
415
-
416
- if(likely(rrdeng_query_handle_globals.protected.available_items)) {
417
- handle = rrdeng_query_handle_globals.protected.available_items;
418
- DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(rrdeng_query_handle_globals.protected.available_items, handle, cache.prev, cache.next);
419
- rrdeng_query_handle_globals.protected.available--;
420
- }
421
-
422
- netdata_spinlock_unlock(&rrdeng_query_handle_globals.protected.spinlock);
423
-
424
- if(unlikely(!handle)) {
425
- handle = mallocz(sizeof(struct rrdeng_query_handle));
426
- __atomic_add_fetch(&rrdeng_query_handle_globals.atomics.allocated, 1, __ATOMIC_RELAXED);
427
- }
428
-
262
+ struct rrdeng_query_handle *handle = aral_mallocz(rrdeng_main.handles.ar);
263
memset(handle, 0, sizeof(struct rrdeng_query_handle));
264
return handle;
265
}
266
267
void rrdeng_query_handle_release(struct rrdeng_query_handle *handle) {
434
- if(unlikely(!handle)) return;
435
-
436
- netdata_spinlock_lock(&rrdeng_query_handle_globals.protected.spinlock);
437
- DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(rrdeng_query_handle_globals.protected.available_items, handle, cache.prev, cache.next);
438
- rrdeng_query_handle_globals.protected.available++;
439
- netdata_spinlock_unlock(&rrdeng_query_handle_globals.protected.spinlock);
268
+ aral_freez(rrdeng_main.handles.ar, handle);
269
}
270
271
// ----------------------------------------------------------------------------
380
struct {
381
struct rrdeng_cmd *prev;
382
struct rrdeng_cmd *next;
554
- } cache;
555
-};
556
-
557
-static struct {
558
- struct {
559
- SPINLOCK spinlock;
560
- struct rrdeng_cmd *available_items;
561
- size_t available;
562
-
563
- struct {
564
- size_t allocated;
565
- } atomics;
566
- } cache;
567
-
568
- struct {
569
- SPINLOCK spinlock;
570
- size_t waiting;
571
- struct rrdeng_cmd *waiting_items_by_priority[STORAGE_PRIORITY_INTERNAL_MAX_DONT_USE];
572
- size_t executed_by_priority[STORAGE_PRIORITY_INTERNAL_MAX_DONT_USE];
383
} queue;
574
-
575
-
576
-} rrdeng_cmd_globals = {
577
- .cache = {
578
- .spinlock = NETDATA_SPINLOCK_INITIALIZER,
579
- .available_items = NULL,
580
- .available = 0,
581
- .atomics = {
582
- .allocated = 0,
583
- },
584
- },
585
- .queue = {
586
- .spinlock = NETDATA_SPINLOCK_INITIALIZER,
587
- .waiting = 0,
588
- },
384
};
385
591
-static void rrdeng_cmd_cleanup1(void) {
592
- struct rrdeng_cmd *item = NULL;
593
-
594
- if(!netdata_spinlock_trylock(&rrdeng_cmd_globals.cache.spinlock))
595
- return;
596
-
597
- if(rrdeng_cmd_globals.cache.available_items && rrdeng_cmd_globals.cache.available > 100) {
598
- item = rrdeng_cmd_globals.cache.available_items;
599
- DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(rrdeng_cmd_globals.cache.available_items, item, cache.prev, cache.next);
600
- rrdeng_cmd_globals.cache.available--;
601
- }
602
- netdata_spinlock_unlock(&rrdeng_cmd_globals.cache.spinlock);
603
-
604
- if(item) {
605
- freez(item);
606
- __atomic_sub_fetch(&rrdeng_cmd_globals.cache.atomics.allocated, 1, __ATOMIC_RELAXED);
607
- }
386
+static void rrdeng_cmd_queue_init(void) {
387
+ rrdeng_main.cmd_queue.ar = aral_create("dbengine-opcodes",
388
+ sizeof(struct rrdeng_cmd),
389
+ 0,
390
+ 65536,
391
+ NULL,
392
+ NULL, NULL, false, false);
393
}
394
395
static inline STORAGE_PRIORITY rrdeng_enq_cmd_map_opcode_to_priority(enum rrdeng_opcode opcode, STORAGE_PRIORITY priority) {
417
}
418
419
void rrdeng_req_cmd(requeue_callback_t get_cmd_cb, void *data, STORAGE_PRIORITY priority) {
635
- netdata_spinlock_lock(&rrdeng_cmd_globals.queue.spinlock);
420
+ netdata_spinlock_lock(&rrdeng_main.cmd_queue.unsafe.spinlock);
421
422
struct rrdeng_cmd *cmd = get_cmd_cb(data);
423
if(cmd) {
424
priority = rrdeng_enq_cmd_map_opcode_to_priority(cmd->opcode, priority);
425
426
if (cmd->priority > priority) {
642
- DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(rrdeng_cmd_globals.queue.waiting_items_by_priority[cmd->priority], cmd, cache.prev, cache.next);
643
- DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(rrdeng_cmd_globals.queue.waiting_items_by_priority[priority], cmd, cache.prev, cache.next);
427
+ DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(rrdeng_main.cmd_queue.unsafe.waiting_items_by_priority[cmd->priority], cmd, queue.prev, queue.next);
428
+ DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(rrdeng_main.cmd_queue.unsafe.waiting_items_by_priority[priority], cmd, queue.prev, queue.next);
429
cmd->priority = priority;
430
}
431
}
432
648
- netdata_spinlock_unlock(&rrdeng_cmd_globals.queue.spinlock);
433
+ netdata_spinlock_unlock(&rrdeng_main.cmd_queue.unsafe.spinlock);
434
}
435
436
void rrdeng_enq_cmd(struct rrdengine_instance *ctx, enum rrdeng_opcode opcode, void *data, struct completion *completion,
437
enum storage_priority priority, enqueue_callback_t enqueue_cb, dequeue_callback_t dequeue_cb) {
653
- struct rrdeng_cmd *cmd = NULL;
438
439
priority = rrdeng_enq_cmd_map_opcode_to_priority(opcode, priority);
440
657
- netdata_spinlock_lock(&rrdeng_cmd_globals.cache.spinlock);
658
- if(likely(rrdeng_cmd_globals.cache.available_items)) {
659
- cmd = rrdeng_cmd_globals.cache.available_items;
660
- DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(rrdeng_cmd_globals.cache.available_items, cmd, cache.prev, cache.next);
661
- rrdeng_cmd_globals.cache.available--;
662
- }
663
- netdata_spinlock_unlock(&rrdeng_cmd_globals.cache.spinlock);
664
-
665
- if(unlikely(!cmd)) {
666
- cmd = mallocz(sizeof(struct rrdeng_cmd));
667
- __atomic_add_fetch(&rrdeng_cmd_globals.cache.atomics.allocated, 1, __ATOMIC_RELAXED);
668
- }
669
-
441
+ struct rrdeng_cmd *cmd = aral_mallocz(rrdeng_main.cmd_queue.ar);
442
memset(cmd, 0, sizeof(struct rrdeng_cmd));
443
cmd->ctx = ctx;
444
cmd->opcode = opcode;
447
cmd->priority = priority;
448
cmd->dequeue_cb = dequeue_cb;
449
678
- netdata_spinlock_lock(&rrdeng_cmd_globals.queue.spinlock);
679
- DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(rrdeng_cmd_globals.queue.waiting_items_by_priority[priority], cmd, cache.prev, cache.next);
680
- rrdeng_cmd_globals.queue.waiting++;
450
+ netdata_spinlock_lock(&rrdeng_main.cmd_queue.unsafe.spinlock);
451
+ DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(rrdeng_main.cmd_queue.unsafe.waiting_items_by_priority[priority], cmd, queue.prev, queue.next);
452
+ rrdeng_main.cmd_queue.unsafe.waiting++;
453
if(enqueue_cb)
454
enqueue_cb(cmd);
683
- netdata_spinlock_unlock(&rrdeng_cmd_globals.queue.spinlock);
455
+ netdata_spinlock_unlock(&rrdeng_main.cmd_queue.unsafe.spinlock);
456
457
fatal_assert(0 == uv_async_send(&rrdeng_main.async));
458
}
459
460
static inline bool rrdeng_cmd_has_waiting_opcodes_in_lower_priorities(STORAGE_PRIORITY priority, STORAGE_PRIORITY max_priority) {
461
for(; priority <= max_priority ; priority++)
690
- if(rrdeng_cmd_globals.queue.waiting_items_by_priority[priority])
462
+ if(rrdeng_main.cmd_queue.unsafe.waiting_items_by_priority[priority])
463
return true;
464
465
return false;
471
STORAGE_PRIORITY max_priority = work_request_full() ? STORAGE_PRIORITY_INTERNAL_DBENGINE : STORAGE_PRIORITY_BEST_EFFORT;
472
473
// find an opcode to execute from the queue
702
- netdata_spinlock_lock(&rrdeng_cmd_globals.queue.spinlock);
474
+ netdata_spinlock_lock(&rrdeng_main.cmd_queue.unsafe.spinlock);
475
for(STORAGE_PRIORITY priority = STORAGE_PRIORITY_INTERNAL_DBENGINE; priority <= max_priority ; priority++) {
704
- cmd = rrdeng_cmd_globals.queue.waiting_items_by_priority[priority];
476
+ cmd = rrdeng_main.cmd_queue.unsafe.waiting_items_by_priority[priority];
477
if(cmd) {
478
479
// avoid starvation of lower priorities
480
if(unlikely(priority >= STORAGE_PRIORITY_HIGH &&
481
priority < STORAGE_PRIORITY_BEST_EFFORT &&
710
- ++rrdeng_cmd_globals.queue.executed_by_priority[priority] % 50 == 0 &&
482
+ ++rrdeng_main.cmd_queue.unsafe.executed_by_priority[priority] % 50 == 0 &&
483
rrdeng_cmd_has_waiting_opcodes_in_lower_priorities(priority + 1, max_priority))) {
484
// let the others run 2% of the requests
485
cmd = NULL;
487
}
488
489
// remove it from the queue
718
- DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(rrdeng_cmd_globals.queue.waiting_items_by_priority[priority], cmd, cache.prev, cache.next);
719
- rrdeng_cmd_globals.queue.waiting--;
490
+ DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(rrdeng_main.cmd_queue.unsafe.waiting_items_by_priority[priority], cmd, queue.prev, queue.next);
491
+ rrdeng_main.cmd_queue.unsafe.waiting--;
492
break;
493
}
494
}
498
cmd->dequeue_cb = NULL;
499
}
500
729
- netdata_spinlock_unlock(&rrdeng_cmd_globals.queue.spinlock);
501
+ netdata_spinlock_unlock(&rrdeng_main.cmd_queue.unsafe.spinlock);
502
503
struct rrdeng_cmd ret;
504
if(cmd) {
505
// copy it, to return it
506
ret = *cmd;
507
736
- // put it in the cache
737
- netdata_spinlock_lock(&rrdeng_cmd_globals.cache.spinlock);
738
- DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(rrdeng_cmd_globals.cache.available_items, cmd, cache.prev, cache.next);
739
- rrdeng_cmd_globals.cache.available++;
740
- netdata_spinlock_unlock(&rrdeng_cmd_globals.cache.spinlock);
508
+ aral_freez(rrdeng_main.cmd_queue.ar, cmd);
509
}
510
else
511
ret = (struct rrdeng_cmd) {
522
523
// ----------------------------------------------------------------------------
524
757
-#define MAX_PAGE_SIZES_TO_KEEP 3
758
-#define MIN_PAGES_PER_SIZE_TO_KEEP 100
759
-
760
-struct dbengine_page_size {
761
- SPINLOCK spinlock;
762
- size_t page_size; // read-only, no lock required to read it
763
-
764
- size_t demand;
765
- size_t supply;
766
-
767
- size_t hit;
768
- size_t miss;
769
-
770
- size_t used;
771
- size_t array_size;
772
- void **array;
773
-};
774
-
525
struct {
776
- struct {
777
- size_t hit;
778
- size_t miss_wrong_size;
779
- size_t miss_short_supply;
780
- size_t cached_size;
781
- size_t struct_size;
782
- } atomic;
783
-
784
- struct dbengine_page_size slots[MAX_PAGE_SIZES_TO_KEEP];
785
-} dbengine_page_alloc_globals = {
786
- .atomic = {
787
- .struct_size = sizeof(dbengine_page_alloc_globals),
788
- }
789
-};
526
+ ARAL *aral[RRD_STORAGE_TIERS];
527
+} dbengine_page_alloc_globals = {};
528
791
-__attribute__((constructor)) void initialize_sizes_to_slots(void) {
792
- uint8_t found[RRDENG_BLOCK_SIZE + 1];
793
- memset(found, 0, RRDENG_BLOCK_SIZE + 1);
794
-
795
- for(int i = 0; i < MAX_PAGE_SIZES_TO_KEEP ; i++) {
796
- struct dbengine_page_size *dps = &dbengine_page_alloc_globals.slots[i];
797
- memset(dps, 0, sizeof(struct dbengine_page_size));
798
- netdata_spinlock_init(&dps->spinlock);
799
- }
800
-
801
- for(int tier = 0; tier < MAX_PAGE_SIZES_TO_KEEP && tier < RRD_STORAGE_TIERS ; tier++) {
802
- size_t size = tier_page_size[tier];
803
-
804
- if(size <= RRDENG_BLOCK_SIZE && !found[size]) {
805
- struct dbengine_page_size *dps = &dbengine_page_alloc_globals.slots[tier];
806
- dps->page_size = size;
807
- found[size] = 1;
808
- }
809
- }
810
-}
529
+static inline ARAL *page_size_lookup(size_t size) {
530
+ for(size_t tier = 0; tier < storage_tiers ;tier++)
531
+ if(size == tier_page_size[tier])
532
+ return dbengine_page_alloc_globals.aral[tier];
533
812
-static inline struct dbengine_page_size *page_size_lookup(size_t size) {
813
- for(int i = 0; i < MAX_PAGE_SIZES_TO_KEEP ; i++) {
814
- if(size == dbengine_page_alloc_globals.slots[i].page_size)
815
- return &dbengine_page_alloc_globals.slots[i];
816
- }
534
return NULL;
535
}
536
820
-static void dbengine_page_alloc_cleanup1(void) {
821
- for(int i = 0; i < MAX_PAGE_SIZES_TO_KEEP ; i++) {
822
- void *page = NULL;
537
+static void dbengine_page_alloc_init(void) {
538
+ for(size_t i = storage_tiers; i > 0 ;i--) {
539
+ size_t tier = storage_tiers - i;
540
824
- struct dbengine_page_size *dps = &dbengine_page_alloc_globals.slots[i];
825
- netdata_spinlock_lock(&dps->spinlock);
826
- if(dps->used > MIN_PAGES_PER_SIZE_TO_KEEP) {
827
- dps->used--;
828
- internal_fatal(!dps->array[dps->used], "DBENGINE: slot should have a page but is empty");
829
- page = dps->array[dps->used];
830
- dps->array[dps->used] = NULL;
831
- __atomic_sub_fetch(&dbengine_page_alloc_globals.atomic.cached_size, dps->page_size, __ATOMIC_RELAXED);
832
- }
833
- netdata_spinlock_unlock(&dps->spinlock);
541
+ char buf[20 + 1];
542
+ snprintfz(buf, 20, "tier%zu-pages", tier);
543
835
- if(page)
836
- freez(page);
544
+ dbengine_page_alloc_globals.aral[tier] = aral_create(
545
+ buf,
546
+ tier_page_size[tier],
547
+ 64,
548
+ 512 * tier_page_size[tier],
549
+ pgc_aral_statistics(),
550
+ NULL, NULL, false, false);
551
}
552
}
553
554
void *dbengine_page_alloc(size_t size) {
841
- void *page = NULL;
842
-
843
- struct dbengine_page_size *dps = page_size_lookup(size);
844
- if(dps) {
845
- netdata_spinlock_lock(&dps->spinlock);
846
- dps->demand++;
847
-
848
- if(dps->used > 0) {
849
- dps->hit++;
850
- dps->used--;
851
- internal_fatal(!dps->array[dps->used], "DBENGINE: slot should have a page but is empty");
852
- page = dps->array[dps->used];
853
- dps->array[dps->used] = NULL;
854
- __atomic_add_fetch(&dbengine_page_alloc_globals.atomic.hit, 1, __ATOMIC_RELAXED);
855
- __atomic_sub_fetch(&dbengine_page_alloc_globals.atomic.cached_size, dps->page_size, __ATOMIC_RELAXED);
856
- }
857
- else {
858
- dps->miss++;
859
- __atomic_add_fetch(&dbengine_page_alloc_globals.atomic.miss_short_supply, 1, __ATOMIC_RELAXED);
860
- }
555
+ ARAL *ar = page_size_lookup(size);
556
+ if(ar) return aral_mallocz(ar);
557
862
- netdata_spinlock_unlock(&dps->spinlock);
863
- }
864
- else
865
- __atomic_add_fetch(&dbengine_page_alloc_globals.atomic.miss_wrong_size, 1, __ATOMIC_RELAXED);
866
-
867
- if(!page)
868
- page = mallocz(size);
869
-
870
- return page;
558
+ return mallocz(size);
559
}
560
561
void dbengine_page_free(void *page, size_t size __maybe_unused) {
562
if(unlikely(!page || page == DBENGINE_EMPTY_PAGE))
563
return;
564
877
- struct dbengine_page_size *dps = page_size_lookup(size);
878
- if(dps) {
879
- netdata_spinlock_lock(&dps->spinlock);
880
- dps->supply++;
881
-
882
- if(dps->used == dps->array_size) {
883
- size_t new_array_size = dps->array_size ? dps->array_size * 2 : MIN_PAGES_PER_SIZE_TO_KEEP;
884
- dps->array = reallocz(dps->array, new_array_size * sizeof(void *));
885
-
886
- __atomic_add_fetch(&dbengine_page_alloc_globals.atomic.struct_size,
887
- (new_array_size - dps->array_size) * sizeof(void *), __ATOMIC_RELAXED);
888
-
889
- dps->array_size = new_array_size;
890
- }
891
-
892
- if(dps->used < dps->array_size) {
893
- dps->array[dps->used] = page;
894
- dps->used++;
895
- page = NULL;
896
- __atomic_add_fetch(&dbengine_page_alloc_globals.atomic.cached_size, dps->page_size, __ATOMIC_RELAXED);
897
- }
898
-
899
- netdata_spinlock_unlock(&dps->spinlock);
900
- }
901
-
902
- if(page)
565
+ ARAL *ar = page_size_lookup(size);
566
+ if(ar)
567
+ aral_freez(ar, page);
568
+ else
569
freez(page);
570
}
571
1483
1484
struct rrdeng_buffer_sizes rrdeng_get_buffer_sizes(void) {
1485
return (struct rrdeng_buffer_sizes) {
1820
- .opcodes = __atomic_load_n(&rrdeng_cmd_globals.cache.atomics.allocated, __ATOMIC_RELAXED) * sizeof(struct rrdeng_cmd),
1821
- .handles = __atomic_load_n(&rrdeng_query_handle_globals.atomics.allocated, __ATOMIC_RELAXED) * sizeof(struct rrdeng_query_handle),
1822
- .descriptors = __atomic_load_n(&page_descriptor_globals.atomics.allocated, __ATOMIC_RELAXED) * sizeof(struct page_descr_with_data),
1486
+ .pgc = pgc_aral_overhead() + pgc_aral_structures(),
1487
+ .mrg = mrg_aral_overhead() + mrg_aral_structures(),
1488
+ .opcodes = aral_overhead(rrdeng_main.cmd_queue.ar) + aral_structures(rrdeng_main.cmd_queue.ar),
1489
+ .handles = aral_overhead(rrdeng_main.handles.ar) + aral_structures(rrdeng_main.handles.ar),
1490
+ .descriptors = aral_overhead(rrdeng_main.descriptors.ar) + aral_structures(rrdeng_main.descriptors.ar),
1491
.wal = __atomic_load_n(&wal_globals.atomics.allocated, __ATOMIC_RELAXED) * (sizeof(WAL) + RRDENG_BLOCK_SIZE),
1824
- .workers = __atomic_load_n(&work_request_globals.atomics.allocated, __ATOMIC_RELAXED) * sizeof(struct rrdeng_work),
1492
+ .workers = aral_overhead(rrdeng_main.work_cmd.ar),
1493
.pdc = pdc_cache_size(),
1826
- .xt_io = __atomic_load_n(&extent_io_descriptor_globals.atomics.allocated, __ATOMIC_RELAXED) * sizeof(struct extent_io_descriptor),
1494
+ .xt_io = aral_overhead(rrdeng_main.xt_io_descr.ar) + aral_structures(rrdeng_main.xt_io_descr.ar),
1495
.xt_buf = extent_buffer_cache_size(),
1496
.epdl = epdl_cache_size(),
1497
.deol = deol_cache_size(),
1498
.pd = pd_cache_size(),
1831
- .pages = __atomic_load_n(&dbengine_page_alloc_globals.atomic.cached_size, __ATOMIC_RELAXED) +
1832
- __atomic_load_n(&dbengine_page_alloc_globals.atomic.struct_size, __ATOMIC_RELAXED),
1499
1500
#ifdef PDC_USE_JULYL
1501
.julyl = julyl_cache_size(),
1510
static void *cleanup_tp_worker(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t *uv_work_req __maybe_unused) {
1511
worker_is_busy(UV_EVENT_DBENGINE_BUFFERS_CLEANUP);
1512
1847
- rrdeng_cmd_cleanup1();
1848
- work_request_cleanup1();
1849
- page_descriptor_cleanup1();
1850
- extent_io_descriptor_cleanup1();
1851
- pdc_cleanup1();
1852
- page_details_cleanup1();
1853
- rrdeng_query_handle_cleanup1();
1513
wal_cleanup1();
1514
extent_buffer_cleanup1();
1856
- epdl_cleanup1();
1857
- deol_cleanup1();
1858
- dbengine_page_alloc_cleanup1();
1515
1516
{
1517
static time_t last_run_s = 0;
1534
uv_stop(handle->loop);
1535
uv_update_time(handle->loop);
1536
1881
- worker_set_metric(RRDENG_OPCODES_WAITING, (NETDATA_DOUBLE)rrdeng_cmd_globals.queue.waiting);
1882
- worker_set_metric(RRDENG_WORKS_DISPATCHED, (NETDATA_DOUBLE)__atomic_load_n(&work_request_globals.atomics.dispatched, __ATOMIC_RELAXED));
1883
- worker_set_metric(RRDENG_WORKS_EXECUTING, (NETDATA_DOUBLE)__atomic_load_n(&work_request_globals.atomics.executing, __ATOMIC_RELAXED));
1537
+ worker_set_metric(RRDENG_OPCODES_WAITING, (NETDATA_DOUBLE)rrdeng_main.cmd_queue.unsafe.waiting);
1538
+ worker_set_metric(RRDENG_WORKS_DISPATCHED, (NETDATA_DOUBLE)__atomic_load_n(&rrdeng_main.work_cmd.atomics.dispatched, __ATOMIC_RELAXED));
1539
+ worker_set_metric(RRDENG_WORKS_EXECUTING, (NETDATA_DOUBLE)__atomic_load_n(&rrdeng_main.work_cmd.atomics.executing, __ATOMIC_RELAXED));
1540
1541
rrdeng_enq_cmd(NULL, RRDENG_OPCODE_FLUSH_INIT, NULL, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
1542
rrdeng_enq_cmd(NULL, RRDENG_OPCODE_EVICT_INIT, NULL, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
1545
worker_is_idle();
1546
}
1547
1548
+static void dbengine_initialize_structures(void) {
1549
+ pgc_and_mrg_initialize();
1550
+
1551
+ pdc_init();
1552
+ page_details_init();
1553
+ epdl_init();
1554
+ deol_init();
1555
+ rrdeng_cmd_queue_init();
1556
+ work_request_init();
1557
+ rrdeng_query_handle_init();
1558
+ page_descriptors_init();
1559
+ extent_buffer_init();
1560
+ dbengine_page_alloc_init();
1561
+ extent_io_descriptor_init();
1562
+}
1563
+
1564
bool rrdeng_dbengine_spawn(struct rrdengine_instance *ctx __maybe_unused) {
1565
static bool spawned = false;
1566
+ static SPINLOCK spinlock = NETDATA_SPINLOCK_INITIALIZER;
1567
+
1568
+ netdata_spinlock_lock(&spinlock);
1569
1570
if(!spawned) {
1571
int ret;
1594
}
1595
rrdeng_main.timer.data = &rrdeng_main;
1596
1597
+ dbengine_initialize_structures();
1598
+
1599
fatal_assert(0 == uv_thread_create(&rrdeng_main.thread, dbengine_event_loop, &rrdeng_main));
1600
spawned = true;
1601
}
1602
1603
+ netdata_spinlock_unlock(&spinlock);
1604
return true;
1605
}
1606
1646
worker_register_job_custom_metric(RRDENG_WORKS_DISPATCHED, "works dispatched", "works", WORKER_METRIC_ABSOLUTE);
1647
worker_register_job_custom_metric(RRDENG_WORKS_EXECUTING, "works executing", "works", WORKER_METRIC_ABSOLUTE);
1648
1971
- extent_buffer_init();
1972
-
1649
struct rrdeng_main *main = arg;
1650
enum rrdeng_opcode opcode;
1651
struct rrdeng_cmd cmd;