1
+// SPDX-License-Identifier: GPL-3.0-or-later
2
+#define NETDATA_RRD_INTERNALS
3
+
4
+#include "rrdengine.h"
5
+
6
+void sanity_check(void)
7
+{
8
+ /* Magic numbers must fit in the super-blocks */
9
+ BUILD_BUG_ON(strlen(RRDENG_DF_MAGIC) > RRDENG_MAGIC_SZ);
10
+ BUILD_BUG_ON(strlen(RRDENG_JF_MAGIC) > RRDENG_MAGIC_SZ);
11
+
12
+ /* Version strings must fit in the super-blocks */
13
+ BUILD_BUG_ON(strlen(RRDENG_DF_VER) > RRDENG_VER_SZ);
14
+ BUILD_BUG_ON(strlen(RRDENG_JF_VER) > RRDENG_VER_SZ);
15
+
16
+ /* Data file super-block cannot be larger than RRDENG_BLOCK_SIZE */
17
+ BUILD_BUG_ON(RRDENG_DF_SB_PADDING_SZ < 0);
18
+
19
+ BUILD_BUG_ON(sizeof(uuid_t) != UUID_SZ); /* check UUID size */
20
+
21
+ /* page count must fit in 8 bits */
22
+ BUILD_BUG_ON(MAX_PAGES_PER_EXTENT > 255);
23
+}
24
+
25
+void read_extent_cb(uv_fs_t* req)
26
+{
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;
31
+ int ret;
32
+ unsigned i, j, count;
33
+ void *page, *uncompressed_buf = NULL;
34
+ uint32_t payload_length, payload_offset, page_offset, uncompressed_payload_length;
35
+ struct rrdengine_datafile *datafile;
36
+ /* persistent structures */
37
+ struct rrdeng_df_extent_header *header;
38
+ struct rrdeng_df_extent_trailer *trailer;
39
+ uLong crc;
40
+
41
+ xt_io_descr = req->data;
42
+ if (req->result < 0) {
43
+ error("%s: uv_fs_read: %s", __func__, uv_strerror((int)req->result));
44
+ goto cleanup;
45
+ }
46
+
47
+ header = xt_io_descr->buf;
48
+ payload_length = header->payload_length;
49
+ count = header->number_of_pages;
50
+
51
+ payload_offset = sizeof(*header) + sizeof(header->descr[0]) * count;
52
+
53
+ trailer = xt_io_descr->buf + xt_io_descr->bytes - sizeof(*trailer);
54
+ crc = crc32(0L, Z_NULL, 0);
55
+ crc = crc32(crc, xt_io_descr->buf, xt_io_descr->bytes - sizeof(*trailer));
56
+ ret = crc32cmp(trailer->checksum, crc);
57
+ datafile = xt_io_descr->descr_array[0]->extent->datafile;
58
+ debug(D_RRDENGINE, "%s: Extent at offset %"PRIu64"(%u) was read from datafile %u-%u. CRC32 check: %s", __func__,
59
+ xt_io_descr->pos, xt_io_descr->bytes, datafile->tier, datafile->fileno, ret ? "FAILED" : "SUCCEEDED");
60
+ if (unlikely(ret)) {
61
+ /* TODO: handle errors */
62
+ exit(UV_EIO);
63
+ goto cleanup;
64
+ }
65
+
66
+ if (RRD_NO_COMPRESSION != header->compression_algorithm) {
67
+ uncompressed_payload_length = 0;
68
+ for (i = 0 ; i < count ; ++i) {
69
+ uncompressed_payload_length += header->descr[i].page_length;
70
+ }
71
+ uncompressed_buf = mallocz(uncompressed_payload_length);
72
+ ret = LZ4_decompress_safe(xt_io_descr->buf + payload_offset, uncompressed_buf,
73
+ payload_length, uncompressed_payload_length);
74
+ ctx->stats.before_decompress_bytes += payload_length;
75
+ ctx->stats.after_decompress_bytes += ret;
76
+ debug(D_RRDENGINE, "LZ4 decompressed %u bytes to %d bytes.", payload_length, ret);
77
+ /* care, we don't hold the descriptor mutex */
78
+ }
79
+
80
+ for (i = 0 ; i < xt_io_descr->descr_count; ++i) {
81
+ page = mallocz(RRDENG_BLOCK_SIZE);
82
+ descr = xt_io_descr->descr_array[i];
83
+ for (j = 0, page_offset = 0; j < count; ++j) {
84
+ /* care, we don't hold the descriptor mutex */
85
+ if (!uuid_compare(*(uuid_t *) header->descr[j].uuid, *descr->id) &&
86
+ header->descr[j].page_length == descr->page_length &&
87
+ header->descr[j].start_time == descr->start_time &&
88
+ header->descr[j].end_time == descr->end_time) {
89
+ break;
90
+ }
91
+ page_offset += header->descr[j].page_length;
92
+ }
93
+ /* care, we don't hold the descriptor mutex */
94
+ if (RRD_NO_COMPRESSION == header->compression_algorithm) {
95
+ (void) memcpy(page, xt_io_descr->buf + payload_offset + page_offset, descr->page_length);
96
+ } else {
97
+ (void) memcpy(page, uncompressed_buf + page_offset, descr->page_length);
98
+ }
99
+ 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;
104
+ debug(D_RRDENGINE, "%s: Waking up waiters.", __func__);
105
+ if (xt_io_descr->release_descr) {
106
+ pg_cache_put_unsafe(descr);
107
+ } else {
108
+ pg_cache_wake_up_waiters_unsafe(descr);
109
+ }
110
+ uv_mutex_unlock(&descr->mutex);
111
+ }
112
+ if (RRD_NO_COMPRESSION != header->compression_algorithm) {
113
+ free(uncompressed_buf);
114
+ }
115
+ if (xt_io_descr->completion)
116
+ complete(xt_io_descr->completion);
117
+cleanup:
118
+ uv_fs_req_cleanup(req);
119
+ free(xt_io_descr->buf);
120
+ free(xt_io_descr);
121
+}
122
+
123
+
124
+static void do_read_extent(struct rrdengine_worker_config* wc,
125
+ struct rrdeng_page_cache_descr **descr,
126
+ unsigned count,
127
+ uint8_t release_descr)
128
+{
129
+ struct rrdengine_instance *ctx = wc->ctx;
130
+ int ret;
131
+ unsigned i, size_bytes, pos, real_io_size;
132
+// uint32_t payload_length;
133
+ struct extent_io_descriptor *xt_io_descr;
134
+ struct rrdengine_datafile *datafile;
135
+
136
+ datafile = descr[0]->extent->datafile;
137
+ pos = descr[0]->extent->offset;
138
+ size_bytes = descr[0]->extent->size;
139
+
140
+ xt_io_descr = mallocz(sizeof(*xt_io_descr));
141
+ ret = posix_memalign((void *)&xt_io_descr->buf, RRDFILE_ALIGNMENT, ALIGN_BYTES_CEILING(size_bytes));
142
+ if (unlikely(ret)) {
143
+ fatal("posix_memalign:%s", strerror(ret));
144
+ /* free(xt_io_descr);
145
+ return;*/
146
+ }
147
+ for (i = 0 ; i < count; ++i) {
148
+ uv_mutex_lock(&descr[i]->mutex);
149
+ descr[i]->flags |= RRD_PAGE_READ_PENDING;
150
+// payload_length = descr[i]->page_length;
151
+ uv_mutex_unlock(&descr[i]->mutex);
152
+
153
+ xt_io_descr->descr_array[i] = descr[i];
154
+ }
155
+ xt_io_descr->descr_count = count;
156
+ xt_io_descr->bytes = size_bytes;
157
+ xt_io_descr->pos = pos;
158
+ xt_io_descr->req.data = xt_io_descr;
159
+ xt_io_descr->completion = NULL;
160
+ /* xt_io_descr->descr_commit_idx_array[0] */
161
+ xt_io_descr->release_descr = release_descr;
162
+
163
+ real_io_size = ALIGN_BYTES_CEILING(size_bytes);
164
+ xt_io_descr->iov = uv_buf_init((void *)xt_io_descr->buf, real_io_size);
165
+ ret = uv_fs_read(wc->loop, &xt_io_descr->req, datafile->file, &xt_io_descr->iov, 1, pos, read_extent_cb);
166
+ assert (-1 != ret);
167
+ ctx->stats.io_read_bytes += real_io_size;
168
+ ++ctx->stats.io_read_requests;
169
+ ctx->stats.io_read_extent_bytes += real_io_size;
170
+ ++ctx->stats.io_read_extents;
171
+ ctx->stats.pg_cache_backfills += count;
172
+}
173
+
174
+static void commit_data_extent(struct rrdengine_worker_config* wc, struct extent_io_descriptor *xt_io_descr)
175
+{
176
+ struct rrdengine_instance *ctx = wc->ctx;
177
+ unsigned count, payload_length, descr_size, size_bytes;
178
+ void *buf;
179
+ /* persistent structures */
180
+ struct rrdeng_df_extent_header *df_header;
181
+ struct rrdeng_jf_transaction_header *jf_header;
182
+ struct rrdeng_jf_store_data *jf_metric_data;
183
+ struct rrdeng_jf_transaction_trailer *jf_trailer;
184
+ uLong crc;
185
+
186
+ df_header = xt_io_descr->buf;
187
+ count = df_header->number_of_pages;
188
+ descr_size = sizeof(*jf_metric_data->descr) * count;
189
+ payload_length = sizeof(*jf_metric_data) + descr_size;
190
+ size_bytes = sizeof(*jf_header) + payload_length + sizeof(*jf_trailer);
191
+
192
+ buf = wal_get_transaction_buffer(wc, size_bytes);
193
+
194
+ jf_header = buf;
195
+ jf_header->type = STORE_DATA;
196
+ jf_header->reserved = 0;
197
+ jf_header->id = ctx->commit_log.transaction_id++;
198
+ jf_header->payload_length = payload_length;
199
+
200
+ jf_metric_data = buf + sizeof(*jf_header);
201
+ jf_metric_data->extent_offset = xt_io_descr->pos;
202
+ jf_metric_data->extent_size = xt_io_descr->bytes;
203
+ jf_metric_data->number_of_pages = count;
204
+ memcpy(jf_metric_data->descr, df_header->descr, descr_size);
205
+
206
+ jf_trailer = buf + sizeof(*jf_header) + payload_length;
207
+ crc = crc32(0L, Z_NULL, 0);
208
+ crc = crc32(crc, buf, sizeof(*jf_header) + payload_length);
209
+ crc32set(jf_trailer->checksum, crc);
210
+}
211
+
212
+static void do_commit_transaction(struct rrdengine_worker_config* wc, uint8_t type, void *data)
213
+{
214
+ switch (type) {
215
+ case STORE_DATA:
216
+ commit_data_extent(wc, (struct extent_io_descriptor *)data);
217
+ break;
218
+ default:
219
+ assert(type == STORE_DATA);
220
+ break;
221
+ }
222
+}
223
+
224
+void flush_pages_cb(uv_fs_t* req)
225
+{
226
+ struct rrdengine_worker_config* wc = req->loop->data;
227
+ struct rrdengine_instance *ctx = wc->ctx;
228
+ struct page_cache *pg_cache = &ctx->pg_cache;
229
+ struct extent_io_descriptor *xt_io_descr;
230
+ struct rrdeng_page_cache_descr *descr;
231
+ struct rrdengine_datafile *datafile;
232
+ int ret;
233
+ unsigned i, count;
234
+ Word_t commit_id;
235
+
236
+ xt_io_descr = req->data;
237
+ if (req->result < 0) {
238
+ error("%s: uv_fs_write: %s", __func__, uv_strerror((int)req->result));
239
+ goto cleanup;
240
+ }
241
+ datafile = xt_io_descr->descr_array[0]->extent->datafile;
242
+ debug(D_RRDENGINE, "%s: Extent at offset %"PRIu64"(%u) was written to datafile %u-%u. Waking up waiters.",
243
+ __func__, xt_io_descr->pos, xt_io_descr->bytes, datafile->tier, datafile->fileno);
244
+
245
+ count = xt_io_descr->descr_count;
246
+ for (i = 0 ; i < count ; ++i) {
247
+ /* care, we don't hold the descriptor mutex */
248
+ descr = xt_io_descr->descr_array[i];
249
+
250
+ uv_rwlock_wrlock(&pg_cache->commited_page_index.lock);
251
+ commit_id = xt_io_descr->descr_commit_idx_array[i];
252
+ ret = JudyLDel(&pg_cache->commited_page_index.JudyL_array, commit_id, PJE0);
253
+ assert(1 == ret);
254
+ --pg_cache->commited_page_index.nr_commited_pages;
255
+ uv_rwlock_wrunlock(&pg_cache->commited_page_index.lock);
256
+
257
+ pg_cache_replaceQ_insert(ctx, descr);
258
+
259
+ uv_mutex_lock(&descr->mutex);
260
+ descr->flags &= ~(RRD_PAGE_DIRTY | RRD_PAGE_WRITE_PENDING);
261
+ /* wake up waiters, care no reference being held */
262
+ pg_cache_wake_up_waiters_unsafe(descr);
263
+ uv_mutex_unlock(&descr->mutex);
264
+ }
265
+ if (xt_io_descr->completion)
266
+ complete(xt_io_descr->completion);
267
+cleanup:
268
+ uv_fs_req_cleanup(req);
269
+ free(xt_io_descr->buf);
270
+ free(xt_io_descr);
271
+}
272
+
273
+/*
274
+ * completion must be NULL or valid.
275
+ * Returns 0 when no flushing can take place.
276
+ * Returns datafile bytes to be written on successful flushing initiation.
277
+ */
278
+static int do_flush_pages(struct rrdengine_worker_config* wc, int force, struct completion *completion)
279
+{
280
+ struct rrdengine_instance *ctx = wc->ctx;
281
+ struct page_cache *pg_cache = &ctx->pg_cache;
282
+ int ret;
283
+ int compressed_size, max_compressed_size = 0;
284
+ unsigned i, count, size_bytes, pos, real_io_size;
285
+ uint32_t uncompressed_payload_length, payload_offset;
286
+ struct rrdeng_page_cache_descr *descr, *eligible_pages[MAX_PAGES_PER_EXTENT];
287
+ struct extent_io_descriptor *xt_io_descr;
288
+ void *compressed_buf = NULL;
289
+ Word_t descr_commit_idx_array[MAX_PAGES_PER_EXTENT];
290
+ Pvoid_t *PValue;
291
+ Word_t Index;
292
+ uint8_t compression_algorithm = ctx->global_compress_alg;
293
+ struct extent_info *extent;
294
+ struct rrdengine_datafile *datafile;
295
+ /* persistent structures */
296
+ struct rrdeng_df_extent_header *header;
297
+ struct rrdeng_df_extent_trailer *trailer;
298
+ uLong crc;
299
+
300
+ if (force) {
301
+ debug(D_RRDENGINE, "Asynchronous flushing of extent has been forced by page pressure.");
302
+ }
303
+ uv_rwlock_rdlock(&pg_cache->commited_page_index.lock);
304
+ for (Index = 0, count = 0, uncompressed_payload_length = 0,
305
+ PValue = JudyLFirst(pg_cache->commited_page_index.JudyL_array, &Index, PJE0),
306
+ descr = unlikely(NULL == PValue) ? NULL : *PValue ;
307
+
308
+ descr != NULL && count != MAX_PAGES_PER_EXTENT ;
309
+
310
+ PValue = JudyLNext(pg_cache->commited_page_index.JudyL_array, &Index, PJE0),
311
+ descr = unlikely(NULL == PValue) ? NULL : *PValue) {
312
+ assert(0 != descr->page_length);
313
+
314
+ uv_mutex_lock(&descr->mutex);
315
+ if (!(descr->flags & RRD_PAGE_WRITE_PENDING)) {
316
+ /* care, no reference being held */
317
+ descr->flags |= RRD_PAGE_WRITE_PENDING;
318
+ uncompressed_payload_length += descr->page_length;
319
+ descr_commit_idx_array[count] = Index;
320
+ eligible_pages[count++] = descr;
321
+ }
322
+ uv_mutex_unlock(&descr->mutex);
323
+ }
324
+ uv_rwlock_rdunlock(&pg_cache->commited_page_index.lock);
325
+
326
+ if (!count) {
327
+ debug(D_RRDENGINE, "%s: no pages eligible for flushing.", __func__);
328
+ if (completion)
329
+ complete(completion);
330
+ return 0;
331
+ }
332
+ xt_io_descr = mallocz(sizeof(*xt_io_descr));
333
+ payload_offset = sizeof(*header) + count * sizeof(header->descr[0]);
334
+ switch (compression_algorithm) {
335
+ case RRD_NO_COMPRESSION:
336
+ size_bytes = payload_offset + uncompressed_payload_length + sizeof(*trailer);
337
+ break;
338
+ default: /* Compress */
339
+ assert(uncompressed_payload_length < LZ4_MAX_INPUT_SIZE);
340
+ max_compressed_size = LZ4_compressBound(uncompressed_payload_length);
341
+ compressed_buf = mallocz(max_compressed_size);
342
+ size_bytes = payload_offset + MAX(uncompressed_payload_length, (unsigned)max_compressed_size) + sizeof(*trailer);
343
+ break;
344
+ }
345
+ ret = posix_memalign((void *)&xt_io_descr->buf, RRDFILE_ALIGNMENT, ALIGN_BYTES_CEILING(size_bytes));
346
+ if (unlikely(ret)) {
347
+ fatal("posix_memalign:%s", strerror(ret));
348
+ /* free(xt_io_descr);*/
349
+ }
350
+ (void) memcpy(xt_io_descr->descr_array, eligible_pages, sizeof(struct rrdeng_page_cache_descr *) * count);
351
+ xt_io_descr->descr_count = count;
352
+
353
+ pos = 0;
354
+ header = xt_io_descr->buf;
355
+ header->compression_algorithm = compression_algorithm;
356
+ header->number_of_pages = count;
357
+ pos += sizeof(*header);
358
+
359
+ extent = mallocz(sizeof(*extent) + count * sizeof(extent->pages[0]));
360
+ datafile = ctx->datafiles.last; /* TODO: check for exceeded size quota */
361
+ extent->offset = datafile->pos;
362
+ extent->number_of_pages = count;
363
+ extent->datafile = datafile;
364
+ extent->next = NULL;
365
+
366
+ for (i = 0 ; i < count ; ++i) {
367
+ /* This is here for performance reasons */
368
+ xt_io_descr->descr_commit_idx_array[i] = descr_commit_idx_array[i];
369
+
370
+ descr = xt_io_descr->descr_array[i];
371
+ header->descr[i].type = PAGE_METRICS;
372
+ uuid_copy(*(uuid_t *)header->descr[i].uuid, *descr->id);
373
+ header->descr[i].page_length = descr->page_length;
374
+ header->descr[i].start_time = descr->start_time;
375
+ header->descr[i].end_time = descr->end_time;
376
+ pos += sizeof(header->descr[i]);
377
+ }
378
+ for (i = 0 ; i < count ; ++i) {
379
+ descr = xt_io_descr->descr_array[i];
380
+ /* care, we don't hold the descriptor mutex */
381
+ (void) memcpy(xt_io_descr->buf + pos, descr->page, descr->page_length);
382
+ descr->extent = extent;
383
+ extent->pages[i] = descr;
384
+
385
+ pos += descr->page_length;
386
+ }
387
+ df_extent_insert(extent);
388
+
389
+ switch (compression_algorithm) {
390
+ case RRD_NO_COMPRESSION:
391
+ header->payload_length = uncompressed_payload_length;
392
+ break;
393
+ default: /* Compress */
394
+ compressed_size = LZ4_compress_default(xt_io_descr->buf + payload_offset, compressed_buf,
395
+ uncompressed_payload_length, max_compressed_size);
396
+ ctx->stats.before_compress_bytes += uncompressed_payload_length;
397
+ ctx->stats.after_compress_bytes += compressed_size;
398
+ debug(D_RRDENGINE, "LZ4 compressed %"PRIu32" bytes to %d bytes.", uncompressed_payload_length, compressed_size);
399
+ (void) memcpy(xt_io_descr->buf + payload_offset, compressed_buf, compressed_size);
400
+ free(compressed_buf);
401
+ size_bytes = payload_offset + compressed_size + sizeof(*trailer);
402
+ header->payload_length = compressed_size;
403
+ break;
404
+ }
405
+ extent->size = size_bytes;
406
+ xt_io_descr->bytes = size_bytes;
407
+ xt_io_descr->pos = datafile->pos;
408
+ xt_io_descr->req.data = xt_io_descr;
409
+ xt_io_descr->completion = completion;
410
+
411
+ trailer = xt_io_descr->buf + size_bytes - sizeof(*trailer);
412
+ crc = crc32(0L, Z_NULL, 0);
413
+ crc = crc32(crc, xt_io_descr->buf, size_bytes - sizeof(*trailer));
414
+ crc32set(trailer->checksum, crc);
415
+
416
+ real_io_size = ALIGN_BYTES_CEILING(size_bytes);
417
+ xt_io_descr->iov = uv_buf_init((void *)xt_io_descr->buf, real_io_size);
418
+ ret = uv_fs_write(wc->loop, &xt_io_descr->req, datafile->file, &xt_io_descr->iov, 1, datafile->pos, flush_pages_cb);
419
+ assert (-1 != ret);
420
+ ctx->stats.io_write_bytes += real_io_size;
421
+ ++ctx->stats.io_write_requests;
422
+ ctx->stats.io_write_extent_bytes += real_io_size;
423
+ ++ctx->stats.io_write_extents;
424
+ do_commit_transaction(wc, STORE_DATA, xt_io_descr);
425
+ datafile->pos += ALIGN_BYTES_CEILING(size_bytes);
426
+ ctx->disk_space += ALIGN_BYTES_CEILING(size_bytes);
427
+ rrdeng_test_quota(wc);
428
+
429
+ return ALIGN_BYTES_CEILING(size_bytes);
430
+}
431
+
432
+static void after_delete_old_data(uv_work_t *req, int status)
433
+{
434
+ struct rrdengine_instance *ctx = req->data;
435
+ struct rrdengine_worker_config* wc = &ctx->worker_config;
436
+ struct rrdengine_datafile *datafile;
437
+ struct rrdengine_journalfile *journalfile;
438
+ unsigned bytes;
439
+
440
+ (void)status;
441
+ datafile = ctx->datafiles.first;
442
+ journalfile = datafile->journalfile;
443
+ bytes = datafile->pos + journalfile->pos;
444
+
445
+ datafile_list_delete(ctx, datafile);
446
+ destroy_journal_file(journalfile, datafile);
447
+ destroy_data_file(datafile);
448
+ info("Deleted data file \""DATAFILE_PREFIX RRDENG_FILE_NUMBER_PRINT_TMPL DATAFILE_EXTENSION"\".",
449
+ datafile->tier, datafile->fileno);
450
+ free(journalfile);
451
+ free(datafile);
452
+
453
+ ctx->disk_space -= bytes;
454
+ info("Reclaimed %u bytes of disk space.", bytes);
455
+
456
+ /* unfreeze command processing */
457
+ wc->now_deleting.data = NULL;
458
+ /* wake up event loop */
459
+ assert(0 == uv_async_send(&wc->async));
460
+}
461
+
462
+static void delete_old_data(uv_work_t *req)
463
+{
464
+ struct rrdengine_instance *ctx = req->data;
465
+ struct rrdengine_datafile *datafile;
466
+ struct extent_info *extent, *next;
467
+ struct rrdeng_page_cache_descr *descr;
468
+ unsigned count, i;
469
+
470
+ /* Safe to use since it will be deleted after we are done */
471
+ datafile = ctx->datafiles.first;
472
+
473
+ for (extent = datafile->extents.first ; extent != NULL ; extent = next) {
474
+ count = extent->number_of_pages;
475
+ for (i = 0 ; i < count ; ++i) {
476
+ descr = extent->pages[i];
477
+ pg_cache_punch_hole(ctx, descr);
478
+ }
479
+ next = extent->next;
480
+ free(extent);
481
+ }
482
+}
483
+
484
+void rrdeng_test_quota(struct rrdengine_worker_config* wc)
485
+{
486
+ struct rrdengine_instance *ctx = wc->ctx;
487
+ struct rrdengine_datafile *datafile;
488
+ unsigned current_size, target_size;
489
+ uint8_t out_of_space, only_one_datafile;
490
+
491
+ out_of_space = 0;
492
+ if (unlikely(ctx->disk_space > ctx->max_disk_space)) {
493
+ out_of_space = 1;
494
+ }
495
+ datafile = ctx->datafiles.last;
496
+ current_size = datafile->pos;
497
+ target_size = ctx->max_disk_space / TARGET_DATAFILES;
498
+ target_size = MIN(target_size, MAX_DATAFILE_SIZE);
499
+ target_size = MAX(target_size, MIN_DATAFILE_SIZE);
500
+ only_one_datafile = (datafile == ctx->datafiles.first) ? 1 : 0;
501
+ if (unlikely(current_size >= target_size || (out_of_space && only_one_datafile))) {
502
+ /* Finalize data and journal file and create a new pair */
503
+ wal_flush_transaction_buffer(wc);
504
+ create_new_datafile_pair(ctx, 1, datafile->fileno + 1);
505
+ }
506
+ if (unlikely(out_of_space)) {
507
+ /* delete old data */
508
+ if (wc->now_deleting.data) {
509
+ /* already deleting data */
510
+ return;
511
+ }
512
+ info("Deleting data file \""DATAFILE_PREFIX RRDENG_FILE_NUMBER_PRINT_TMPL DATAFILE_EXTENSION"\".",
513
+ ctx->datafiles.first->tier, ctx->datafiles.first->fileno);
514
+ wc->now_deleting.data = ctx;
515
+ uv_queue_work(wc->loop, &wc->now_deleting, delete_old_data, after_delete_old_data);
516
+ }
517
+}
518
+
519
+int init_rrd_files(struct rrdengine_instance *ctx)
520
+{
521
+ return init_data_files(ctx);
522
+}
523
+
524
+void rrdeng_init_cmd_queue(struct rrdengine_worker_config* wc)
525
+{
526
+ wc->cmd_queue.head = wc->cmd_queue.tail = 0;
527
+ wc->queue_size = 0;
528
+ assert(0 == uv_cond_init(&wc->cmd_cond));
529
+ assert(0 == uv_mutex_init(&wc->cmd_mutex));
530
+}
531
+
532
+void rrdeng_enq_cmd(struct rrdengine_worker_config* wc, struct rrdeng_cmd *cmd)
533
+{
534
+ unsigned queue_size;
535
+
536
+ /* wait for free space in queue */
537
+ uv_mutex_lock(&wc->cmd_mutex);
538
+ while ((queue_size = wc->queue_size) == RRDENG_CMD_Q_MAX_SIZE) {
539
+ uv_cond_wait(&wc->cmd_cond, &wc->cmd_mutex);
540
+ }
541
+ assert(queue_size < RRDENG_CMD_Q_MAX_SIZE);
542
+ /* enqueue command */
543
+ wc->cmd_queue.cmd_array[wc->cmd_queue.tail] = *cmd;
544
+ wc->cmd_queue.tail = wc->cmd_queue.tail != RRDENG_CMD_Q_MAX_SIZE - 1 ?
545
+ wc->cmd_queue.tail + 1 : 0;
546
+ wc->queue_size = queue_size + 1;
547
+ uv_mutex_unlock(&wc->cmd_mutex);
548
+
549
+ /* wake up event loop */
550
+ assert(0 == uv_async_send(&wc->async));
551
+}
552
+
553
+struct rrdeng_cmd rrdeng_deq_cmd(struct rrdengine_worker_config* wc)
554
+{
555
+ struct rrdeng_cmd ret;
556
+ unsigned queue_size;
557
+
558
+ uv_mutex_lock(&wc->cmd_mutex);
559
+ queue_size = wc->queue_size;
560
+ if (queue_size == 0) {
561
+ ret.opcode = RRDENG_NOOP;
562
+ } else {
563
+ /* dequeue command */
564
+ ret = wc->cmd_queue.cmd_array[wc->cmd_queue.head];
565
+ if (queue_size == 1) {
566
+ wc->cmd_queue.head = wc->cmd_queue.tail = 0;
567
+ } else {
568
+ wc->cmd_queue.head = wc->cmd_queue.head != RRDENG_CMD_Q_MAX_SIZE - 1 ?
569
+ wc->cmd_queue.head + 1 : 0;
570
+ }
571
+ wc->queue_size = queue_size - 1;
572
+
573
+ /* wake up producers */
574
+ uv_cond_signal(&wc->cmd_cond);
575
+ }
576
+ uv_mutex_unlock(&wc->cmd_mutex);
577
+
578
+ return ret;
579
+}
580
+
581
+void async_cb(uv_async_t *handle)
582
+{
583
+ uv_stop(handle->loop);
584
+ uv_update_time(handle->loop);
585
+ debug(D_RRDENGINE, "%s called, active=%d.", __func__, uv_is_active((uv_handle_t *)handle));
586
+}
587
+
588
+void timer_cb(uv_timer_t* handle)
589
+{
590
+ struct rrdengine_worker_config* wc = handle->data;
591
+ struct rrdengine_instance *ctx = wc->ctx;
592
+
593
+ uv_stop(handle->loop);
594
+ uv_update_time(handle->loop);
595
+ rrdeng_test_quota(wc);
596
+ debug(D_RRDENGINE, "%s: timeout reached.", __func__);
597
+ if (likely(!wc->now_deleting.data)) {
598
+ unsigned total_bytes, bytes_written;
599
+
600
+ /* There is free space so we can write to disk */
601
+ debug(D_RRDENGINE, "Flushing pages to disk.");
602
+ for (total_bytes = bytes_written = do_flush_pages(wc, 0, NULL) ;
603
+ bytes_written && (total_bytes < DATAFILE_IDEAL_IO_SIZE) ;
604
+ total_bytes += bytes_written) {
605
+ bytes_written = do_flush_pages(wc, 0, NULL);
606
+ }
607
+ }
608
+#ifdef NETDATA_INTERNAL_CHECKS
609
+ {
610
+ char buf[4096];
611
+ debug(D_RRDENGINE, "%s", get_rrdeng_statistics(ctx, buf, sizeof(buf)));
612
+ }
613
+#endif
614
+}
615
+
616
+/* Flushes dirty pages when timer expires */
617
+#define TIMER_PERIOD_MS (1000)
618
+
619
+#define CMD_BATCH_SIZE (256)
620
+
621
+void rrdeng_worker(void* arg)
622
+{
623
+ struct rrdengine_worker_config* wc = arg;
624
+ struct rrdengine_instance *ctx = wc->ctx;
625
+ uv_loop_t* loop;
626
+ int shutdown;
627
+ enum rrdeng_opcode opcode;
628
+ uv_timer_t timer_req;
629
+ struct rrdeng_cmd cmd;
630
+
631
+ rrdeng_init_cmd_queue(wc);
632
+
633
+ loop = wc->loop = mallocz(sizeof(uv_loop_t));
634
+ uv_loop_init(loop);
635
+ loop->data = wc;
636
+
637
+ uv_async_init(wc->loop, &wc->async, async_cb);
638
+ wc->async.data = wc;
639
+
640
+ wc->now_deleting.data = NULL;
641
+
642
+ /* dirty page flushing timer */
643
+ uv_timer_init(loop, &timer_req);
644
+ timer_req.data = wc;
645
+
646
+ /* wake up initialization thread */
647
+ complete(&ctx->rrdengine_completion);
648
+
649
+ uv_timer_start(&timer_req, timer_cb, TIMER_PERIOD_MS, TIMER_PERIOD_MS);
650
+ shutdown = 0;
651
+ while (shutdown == 0 || uv_loop_alive(loop)) {
652
+ uv_run(loop, UV_RUN_DEFAULT);
653
+ /* wait for commands */
654
+ do {
655
+ cmd = rrdeng_deq_cmd(wc);
656
+ opcode = cmd.opcode;
657
+
658
+ switch (opcode) {
659
+ case RRDENG_NOOP:
660
+ /* the command queue was empty, do nothing */
661
+ break;
662
+ case RRDENG_SHUTDOWN:
663
+ shutdown = 1;
664
+ if (unlikely(wc->now_deleting.data)) {
665
+ /* postpone shutdown until after deletion */
666
+ info("Postponing shutting RRD engine event loop down until after datafile deletion is finished.");
667
+ rrdeng_enq_cmd(wc, &cmd);
668
+ break;
669
+ }
670
+ /*
671
+ * uv_async_send after uv_close does not seem to crash in linux at the moment,
672
+ * it is however undocumented behaviour and we need to be aware if this becomes
673
+ * an issue in the future.
674
+ */
675
+ uv_close((uv_handle_t *)&wc->async, NULL);
676
+ assert(0 == uv_timer_stop(&timer_req));
677
+ uv_close((uv_handle_t *)&timer_req, NULL);
678
+ info("Shutting down RRD engine event loop.");
679
+ while (do_flush_pages(wc, 1, NULL)) {
680
+ ; /* Force flushing of all commited pages. */
681
+ }
682
+ break;
683
+ case RRDENG_READ_PAGE:
684
+ do_read_extent(wc, &cmd.read_page.page_cache_descr, 1, 0);
685
+ break;
686
+ case RRDENG_READ_EXTENT:
687
+ do_read_extent(wc, cmd.read_extent.page_cache_descr, cmd.read_extent.page_count, 1);
688
+ break;
689
+ case RRDENG_COMMIT_PAGE:
690
+ do_commit_transaction(wc, STORE_DATA, NULL);
691
+ break;
692
+ case RRDENG_FLUSH_PAGES: {
693
+ unsigned total_bytes, bytes_written;
694
+
695
+ /* First I/O should be enough to call completion */
696
+ 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);
701
+ }
702
+ break;
703
+ }
704
+ default:
705
+ debug(D_RRDENGINE, "%s: default.", __func__);
706
+ break;
707
+ }
708
+ } while (opcode != RRDENG_NOOP);
709
+ }
710
+ /* cleanup operations of the event loop */
711
+ wal_flush_transaction_buffer(wc);
712
+ uv_run(loop, UV_RUN_DEFAULT);
713
+
714
+ info("Shutting down RRD engine event loop complete.");
715
+ /* TODO: don't let the API block by waiting to enqueue commands */
716
+ uv_cond_destroy(&wc->cmd_cond);
717
+/* uv_mutex_destroy(&wc->cmd_mutex); */
718
+ assert(0 == uv_loop_close(loop));
719
+ free(loop);
720
+}
721
+
722
+
723
+#define NR_PAGES (256)
724
+static void basic_functional_test(struct rrdengine_instance *ctx)
725
+{
726
+ int i, j, failed_validations;
727
+ uuid_t uuid[NR_PAGES];
728
+ 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 */
732
+
733
+ for (i = 0 ; i < NR_PAGES ; ++i) {
734
+ uuid_generate(uuid[i]);
735
+ uuid_unparse_lower(uuid[i], uuid_str);
736
+// fprintf(stderr, "Generated uuid[%d]=%s\n", i, uuid_str);
737
+ buf = rrdeng_create_page(&uuid[i], &handle[i]);
738
+ /* Each page contains 10 times its own UUID stringified */
739
+ for (j = 0 ; j < 100 ; ++j) {
740
+ strcpy(buf + 37 * j, uuid_str);
741
+ strcpy(backup[i] + 37 * j, uuid_str);
742
+ }
743
+ rrdeng_commit_page(ctx, handle[i], (Word_t)i);
744
+ }
745
+ fprintf(stderr, "\n********** CREATED %d METRIC PAGES ***********\n\n", NR_PAGES);
746
+ failed_validations = 0;
747
+ for (i = 0 ; i < NR_PAGES ; ++i) {
748
+ buf = rrdeng_get_latest_page(ctx, &uuid[i], (void **)&handle[i]);
749
+ if (NULL == buf) {
750
+ ++failed_validations;
751
+ fprintf(stderr, "Page %d was LOST.\n", i);
752
+ }
753
+ if (memcmp(backup[i], buf, 37 * 100)) {
754
+ ++failed_validations;
755
+ fprintf(stderr, "Page %d data comparison with backup FAILED validation.\n", i);
756
+ }
757
+ rrdeng_put_page(ctx, handle[i]);
758
+ }
759
+ fprintf(stderr, "\n********** CORRECTLY VALIDATED %d/%d METRIC PAGES ***********\n\n",
760
+ NR_PAGES - failed_validations, NR_PAGES);
761
+
762
+}
763
+/* C entry point for development purposes
764
+ * make "LDFLAGS=-errdengine_main"
765
+ */
766
+void rrdengine_main(void)
767
+{
768
+ int ret;
769
+ struct rrdengine_instance *ctx;
770
+
771
+ ret = rrdeng_init(&ctx, "/tmp", RRDENG_MIN_PAGE_CACHE_SIZE_MB, RRDENG_MIN_DISK_SPACE_MB);
772
+ if (ret) {
773
+ exit(ret);
774
+ }
775
+ basic_functional_test(ctx);
776
+
777
+ rrdeng_exit(ctx);
778
+ fprintf(stderr, "Hello world!");
779
+ exit(0);
780
+}
\ No newline at end of file