@cryptotaxi247 / netdata-1 / commits / 6312080b6

Handle file descriptors running out (#6303)

* Handle file descriptors running out * Added alarm for dbengine FS and I/O errors * more verbose alarm message * * Added File-Descriptor budget to Database Engine instances. * Changed FD budget of the web server from 50% to 25%. * Allocated 25% of FDs to dbengine. * Created a new dbengine global FD utilization chart.

Markos Fountoulakis committed Jun 21, 2019 at 12:52 UTC 6312080b697dd3b3415a504dd59c351815155e47
18 files changed +578 -158
configs.signatures
+1
@@ -381,6 +381,7 @@ declare -A configs_signatures=(
381 ['7deb236ec68a512b9bdd18e6a51d76f7']='python.d/mysql.conf'
382 ['7e5fc1644aa7a54f9dbb1bd102521b09']='health.d/memcached.conf'
383 ['7f13631183fbdf79c21c8e5a171e9b34']='health.d/zfs.conf'
384 + ['ce285c90747428ee5da4efb547418dda']='health.d/dbengine.conf'
385 ['7fb8184d56a27040e73261ed9c6fc76f']='health_alarm_notify.conf'
386 ['80266bddd3df374923c750a6de91d120']='health.d/apache.conf'
387 ['803a7f9dcb942eeac0fd764b9e3e38ca']='fping.conf'
daemon/global_statistics.c
+70 -1
@@ -538,7 +538,7 @@ void global_statistics_charts(void) {
538 unsigned long long stats_array[RRDENG_NR_STATS];
539
540 /* get localhost's DB engine's statistics */
541 - rrdeng_get_28_statistics(localhost->rrdeng_ctx, stats_array);
541 + rrdeng_get_33_statistics(localhost->rrdeng_ctx, stats_array);
542
543 // ----------------------------------------------------------------
544
@@ -749,6 +749,75 @@ void global_statistics_charts(void) {
749 rrddim_set_by_pointer(st_io_stats, rd_writes, (collected_number)stats_array[16]);
750 rrdset_done(st_io_stats);
751 }
752 +
753 + // ----------------------------------------------------------------
754 +
755 + {
756 + static RRDSET *st_errors = NULL;
757 + static RRDDIM *rd_fs_errors = NULL;
758 + static RRDDIM *rd_io_errors = NULL;
759 +
760 + if (unlikely(!st_errors)) {
761 + st_errors = rrdset_create_localhost(
762 + "netdata"
763 + , "dbengine_global_errors"
764 + , NULL
765 + , "dbengine"
766 + , NULL
767 + , "NetData DB engine errors"
768 + , "errors/s"
769 + , "netdata"
770 + , "stats"
771 + , 130507
772 + , localhost->rrd_update_every
773 + , RRDSET_TYPE_LINE
774 + );
775 +
776 + rd_io_errors = rrddim_add(st_errors, "I/O errors", NULL, 1, 1, RRD_ALGORITHM_INCREMENTAL);
777 + rd_fs_errors = rrddim_add(st_errors, "FS errors", NULL, 1, 1, RRD_ALGORITHM_INCREMENTAL);
778 + }
779 + else
780 + rrdset_next(st_errors);
781 +
782 + rrddim_set_by_pointer(st_errors, rd_io_errors, (collected_number)stats_array[30]);
783 + rrddim_set_by_pointer(st_errors, rd_fs_errors, (collected_number)stats_array[31]);
784 + rrdset_done(st_errors);
785 + }
786 +
787 + // ----------------------------------------------------------------
788 +
789 + {
790 + static RRDSET *st_fd = NULL;
791 + static RRDDIM *rd_fd_current = NULL;
792 + static RRDDIM *rd_fd_max = NULL;
793 +
794 + if (unlikely(!st_fd)) {
795 + st_fd = rrdset_create_localhost(
796 + "netdata"
797 + , "dbengine_global_file_descriptors"
798 + , NULL
799 + , "dbengine"
800 + , NULL
801 + , "NetData DB engine File Descriptors"
802 + , "descriptors"
803 + , "netdata"
804 + , "stats"
805 + , 130508
806 + , localhost->rrd_update_every
807 + , RRDSET_TYPE_LINE
808 + );
809 +
810 + rd_fd_current = rrddim_add(st_fd, "current", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
811 + rd_fd_max = rrddim_add(st_fd, "max", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
812 + }
813 + else
814 + rrdset_next(st_fd);
815 +
816 + rrddim_set_by_pointer(st_fd, rd_fd_current, (collected_number)stats_array[32]);
817 + /* Careful here, modify this accordingly if the File-Descriptor budget ever changes */
818 + rrddim_set_by_pointer(st_fd, rd_fd_max, (collected_number)rlimit_nofile.rlim_cur / 4);
819 + rrdset_done(st_fd);
820 + }
821 }
822 #endif
823
database/engine/datafile.c
+155 -49
@@ -49,44 +49,69 @@ static void datafile_init(struct rrdengine_datafile *datafile, struct rrdengine_
49 datafile->ctx = ctx;
50 }
51
52 -static void generate_datafilepath(struct rrdengine_datafile *datafile, char *str, size_t maxlen)
52 +void generate_datafilepath(struct rrdengine_datafile *datafile, char *str, size_t maxlen)
53 {
54 (void) snprintf(str, maxlen, "%s/" DATAFILE_PREFIX RRDENG_FILE_NUMBER_PRINT_TMPL DATAFILE_EXTENSION,
55 datafile->ctx->dbfiles_path, datafile->tier, datafile->fileno);
56 }
57
58 +int close_data_file(struct rrdengine_datafile *datafile)
59 +{
60 + struct rrdengine_instance *ctx = datafile->ctx;
61 + uv_fs_t req;
62 + int ret;
63 + char path[RRDENG_PATH_MAX];
64 +
65 + generate_datafilepath(datafile, path, sizeof(path));
66 +
67 + ret = uv_fs_close(NULL, &req, datafile->file, NULL);
68 + if (ret < 0) {
69 + error("uv_fs_close(%s): %s", path, uv_strerror(ret));
70 + ++ctx->stats.fs_errors;
71 + rrd_stat_atomic_add(&global_fs_errors, 1);
72 + }
73 + uv_fs_req_cleanup(&req);
74 +
75 + return ret;
76 +}
77 +
78 +
79 int destroy_data_file(struct rrdengine_datafile *datafile)
80 {
81 struct rrdengine_instance *ctx = datafile->ctx;
82 uv_fs_t req;
62 - int ret, fd;
63 - char path[1024];
83 + int ret;
84 + char path[RRDENG_PATH_MAX];
85 +
86 + generate_datafilepath(datafile, path, sizeof(path));
87
88 ret = uv_fs_ftruncate(NULL, &req, datafile->file, 0, NULL);
89 if (ret < 0) {
67 - fatal("uv_fs_ftruncate: %s", uv_strerror(ret));
90 + error("uv_fs_ftruncate(%s): %s", path, uv_strerror(ret));
91 + ++ctx->stats.fs_errors;
92 + rrd_stat_atomic_add(&global_fs_errors, 1);
93 }
69 - assert(0 == req.result);
94 uv_fs_req_cleanup(&req);
95
96 ret = uv_fs_close(NULL, &req, datafile->file, NULL);
97 if (ret < 0) {
74 - fatal("uv_fs_close: %s", uv_strerror(ret));
98 + error("uv_fs_close(%s): %s", path, uv_strerror(ret));
99 + ++ctx->stats.fs_errors;
100 + rrd_stat_atomic_add(&global_fs_errors, 1);
101 }
76 - assert(0 == req.result);
102 uv_fs_req_cleanup(&req);
103
79 - generate_datafilepath(datafile, path, sizeof(path));
80 - fd = uv_fs_unlink(NULL, &req, path, NULL);
81 - if (fd < 0) {
82 - fatal("uv_fs_fsunlink: %s", uv_strerror(fd));
104 + ret = uv_fs_unlink(NULL, &req, path, NULL);
105 + if (ret < 0) {
106 + error("uv_fs_fsunlink(%s): %s", path, uv_strerror(ret));
107 + ++ctx->stats.fs_errors;
108 + rrd_stat_atomic_add(&global_fs_errors, 1);
109 }
84 - assert(0 == req.result);
110 uv_fs_req_cleanup(&req);
111
112 ++ctx->stats.datafile_deletions;
113
89 - return 0;
114 + return ret;
115 }
116
117 int create_data_file(struct rrdengine_datafile *datafile)
@@ -97,13 +122,17 @@ int create_data_file(struct rrdengine_datafile *datafile)
122 int ret, fd;
123 struct rrdeng_df_sb *superblock;
124 uv_buf_t iov;
100 - char path[1024];
125 + char path[RRDENG_PATH_MAX];
126
127 generate_datafilepath(datafile, path, sizeof(path));
128 fd = open_file_direct_io(path, O_CREAT | O_RDWR | O_TRUNC, &file);
129 if (fd < 0) {
105 - fatal("uv_fs_fsopen: %s", uv_strerror(fd));
130 + ++ctx->stats.fs_errors;
131 + rrd_stat_atomic_add(&global_fs_errors, 1);
132 + return fd;
133 }
134 + datafile->file = file;
135 + ++ctx->stats.datafile_creations;
136
137 ret = posix_memalign((void *)&superblock, RRDFILE_ALIGNMENT, sizeof(*superblock));
138 if (unlikely(ret)) {
@@ -117,19 +146,21 @@ int create_data_file(struct rrdengine_datafile *datafile)
146
147 ret = uv_fs_write(NULL, &req, file, &iov, 1, 0, NULL);
148 if (ret < 0) {
120 - fatal("uv_fs_write: %s", uv_strerror(ret));
121 - }
122 - if (req.result < 0) {
123 - fatal("uv_fs_write: %s", uv_strerror((int)req.result));
149 + assert(req.result < 0);
150 + error("uv_fs_write: %s", uv_strerror(ret));
151 + ++ctx->stats.io_errors;
152 + rrd_stat_atomic_add(&global_io_errors, 1);
153 }
154 uv_fs_req_cleanup(&req);
155 free(superblock);
156 + if (ret < 0) {
157 + destroy_data_file(datafile);
158 + return ret;
159 + }
160
128 - datafile->file = file;
161 datafile->pos = sizeof(*superblock);
162 ctx->stats.io_write_bytes += sizeof(*superblock);
163 ++ctx->stats.io_write_requests;
132 - ++ctx->stats.datafile_creations;
164
165 return 0;
166 }
@@ -174,15 +205,15 @@ static int load_data_file(struct rrdengine_datafile *datafile)
205 struct rrdengine_instance *ctx = datafile->ctx;
206 uv_fs_t req;
207 uv_file file;
177 - int ret, fd;
208 + int ret, fd, error;
209 uint64_t file_size;
179 - char path[1024];
210 + char path[RRDENG_PATH_MAX];
211
212 generate_datafilepath(datafile, path, sizeof(path));
213 fd = open_file_direct_io(path, O_RDWR, &file);
214 if (fd < 0) {
184 - /* if (UV_ENOENT != fd) */
185 - error("uv_fs_fsopen: %s", uv_strerror(fd));
215 + ++ctx->stats.fs_errors;
216 + rrd_stat_atomic_add(&global_fs_errors, 1);
217 return fd;
218 }
219 info("Initializing data file \"%s\".", path);
@@ -205,15 +236,21 @@ static int load_data_file(struct rrdengine_datafile *datafile)
236 return 0;
237
238 error:
208 - (void) uv_fs_close(NULL, &req, file, NULL);
239 + error = ret;
240 + ret = uv_fs_close(NULL, &req, file, NULL);
241 + if (ret < 0) {
242 + error("uv_fs_close(%s): %s", path, uv_strerror(ret));
243 + ++ctx->stats.fs_errors;
244 + rrd_stat_atomic_add(&global_fs_errors, 1);
245 + }
246 uv_fs_req_cleanup(&req);
210 - return ret;
247 + return error;
248 }
249
250 static int scan_data_files_cmp(const void *a, const void *b)
251 {
252 struct rrdengine_datafile *file1, *file2;
216 - char path1[1024], path2[1024];
253 + char path1[RRDENG_PATH_MAX], path2[RRDENG_PATH_MAX];
254
255 file1 = *(struct rrdengine_datafile **)a;
256 file2 = *(struct rrdengine_datafile **)b;
@@ -222,7 +259,7 @@ static int scan_data_files_cmp(const void *a, const void *b)
259 return strcmp(path1, path2);
260 }
261
225 -/* Returns number of datafiles that were loaded */
262 +/* Returns number of datafiles that were loaded or < 0 on error */
263 static int scan_data_files(struct rrdengine_instance *ctx)
264 {
265 int ret;
@@ -233,16 +270,22 @@ static int scan_data_files(struct rrdengine_instance *ctx)
270 struct rrdengine_journalfile *journalfile;
271
272 ret = uv_fs_scandir(NULL, &req, ctx->dbfiles_path, 0, NULL);
236 - assert(ret >= 0);
237 - assert(req.result >= 0);
273 + if (ret < 0) {
274 + assert(req.result < 0);
275 + uv_fs_req_cleanup(&req);
276 + error("uv_fs_scandir(%s): %s", ctx->dbfiles_path, uv_strerror(ret));
277 + ++ctx->stats.fs_errors;
278 + rrd_stat_atomic_add(&global_fs_errors, 1);
279 + return ret;
280 + }
281 info("Found %d files in path %s", ret, ctx->dbfiles_path);
282
283 datafiles = callocz(MIN(ret, MAX_DATAFILES), sizeof(*datafiles));
284 for (matched_files = 0 ; UV_EOF != uv_fs_scandir_next(&req, &dent) && matched_files < MAX_DATAFILES ; ) {
242 - info("Scanning file \"%s\"", dent.name);
285 + info("Scanning file \"%s/%s\"", ctx->dbfiles_path, dent.name);
286 ret = sscanf(dent.name, DATAFILE_PREFIX RRDENG_FILE_NUMBER_SCAN_TMPL DATAFILE_EXTENSION, &tier, &no);
287 if (2 == ret) {
245 - info("Matched file \"%s\"", dent.name);
288 + info("Matched file \"%s/%s\"", ctx->dbfiles_path, dent.name);
289 datafile = mallocz(sizeof(*datafile));
290 datafile_init(datafile, ctx, tier, no);
291 datafiles[matched_files++] = datafile;
@@ -250,70 +293,133 @@ static int scan_data_files(struct rrdengine_instance *ctx)
293 }
294 uv_fs_req_cleanup(&req);
295
296 + if (0 == matched_files) {
297 + freez(datafiles);
298 + return 0;
299 + }
300 if (matched_files == MAX_DATAFILES) {
301 error("Warning: hit maximum database engine file limit of %d files", MAX_DATAFILES);
302 }
303 qsort(datafiles, matched_files, sizeof(*datafiles), scan_data_files_cmp);
304 + /* TODO: change this when tiering is implemented */
305 + ctx->last_fileno = datafiles[matched_files - 1]->fileno;
306 +
307 for (failed_to_load = 0, i = 0 ; i < matched_files ; ++i) {
308 datafile = datafiles[i];
309 ret = load_data_file(datafile);
310 if (0 != ret) {
261 - free(datafile);
311 + freez(datafile);
312 ++failed_to_load;
263 - continue;
313 + break;
314 }
315 journalfile = mallocz(sizeof(*journalfile));
316 datafile->journalfile = journalfile;
317 journalfile_init(journalfile, datafile);
318 ret = load_journal_file(ctx, journalfile, datafile);
319 if (0 != ret) {
270 - free(datafile);
271 - free(journalfile);
320 + close_data_file(datafile);
321 + freez(datafile);
322 + freez(journalfile);
323 ++failed_to_load;
273 - continue;
324 + break;
325 }
326 datafile_list_insert(ctx, datafile);
327 ctx->disk_space += datafile->pos + journalfile->pos;
328 }
329 + freez(datafiles);
330 if (failed_to_load) {
279 - error("%u files failed to load.", failed_to_load);
331 + error("%u datafiles failed to load.", failed_to_load);
332 + finalize_data_files(ctx);
333 + return UV_EIO;
334 }
281 - free(datafiles);
335
283 - return matched_files - failed_to_load;
336 + return matched_files;
337 }
338
339 /* Creates a datafile and a journalfile pair */
287 -void create_new_datafile_pair(struct rrdengine_instance *ctx, unsigned tier, unsigned fileno)
340 +int create_new_datafile_pair(struct rrdengine_instance *ctx, unsigned tier, unsigned fileno)
341 {
342 struct rrdengine_datafile *datafile;
343 struct rrdengine_journalfile *journalfile;
344 int ret;
345 + char path[RRDENG_PATH_MAX];
346
293 - info("Creating new data and journal files.");
347 + info("Creating new data and journal files in path %s", ctx->dbfiles_path);
348 datafile = mallocz(sizeof(*datafile));
349 datafile_init(datafile, ctx, tier, fileno);
350 ret = create_data_file(datafile);
297 - assert(!ret);
351 + if (!ret) {
352 + generate_datafilepath(datafile, path, sizeof(path));
353 + info("Created data file \"%s\".", path);
354 + } else {
355 + goto error_after_datafile;
356 + }
357
358 journalfile = mallocz(sizeof(*journalfile));
359 datafile->journalfile = journalfile;
360 journalfile_init(journalfile, datafile);
361 ret = create_journal_file(journalfile, datafile);
303 - assert(!ret);
362 + if (!ret) {
363 + generate_journalfilepath(datafile, path, sizeof(path));
364 + info("Created journal file \"%s\".", path);
365 + } else {
366 + goto error_after_journalfile;
367 + }
368 datafile_list_insert(ctx, datafile);
369 ctx->disk_space += datafile->pos + journalfile->pos;
370 +
371 + return 0;
372 +
373 +error_after_journalfile:
374 + destroy_data_file(datafile);
375 + freez(journalfile);
376 +error_after_datafile:
377 + freez(datafile);
378 + return ret;
379 }
380
308 -/* Page cache must already be initialized. */
381 +/* Page cache must already be initialized.
382 + * Return 0 on success.
383 + */
384 int init_data_files(struct rrdengine_instance *ctx)
385 {
386 int ret;
387
388 ret = scan_data_files(ctx);
314 - if (0 == ret) {
315 - info("Data files not found, creating.");
316 - create_new_datafile_pair(ctx, 1, 1);
389 + if (ret < 0) {
390 + error("Failed to scan path \"%s\".", ctx->dbfiles_path);
391 + return ret;
392 + } else if (0 == ret) {
393 + info("Data files not found, creating in path \"%s\".", ctx->dbfiles_path);
394 + ret = create_new_datafile_pair(ctx, 1, 1);
395 + if (ret) {
396 + error("Failed to create data and journal files in path \"%s\".", ctx->dbfiles_path);
397 + return ret;
398 + }
399 + ctx->last_fileno = 1;
400 }
401 +
402 return 0;
403 +}
404 +
405 +void finalize_data_files(struct rrdengine_instance *ctx)
406 +{
407 + struct rrdengine_datafile *datafile, *next_datafile;
408 + struct rrdengine_journalfile *journalfile;
409 + struct extent_info *extent, *next_extent;
410 +
411 + for (datafile = ctx->datafiles.first ; datafile != NULL ; datafile = next_datafile) {
412 + journalfile = datafile->journalfile;
413 + next_datafile = datafile->next;
414 +
415 + for (extent = datafile->extents.first ; extent != NULL ; extent = next_extent) {
416 + next_extent = extent->next;
417 + freez(extent);
418 + }
419 + close_journal_file(journalfile, datafile);
420 + close_data_file(datafile);
421 + freez(journalfile);
422 + freez(datafile);
423 +
424 + }
425 }
\ No newline at end of file
database/engine/datafile.h
+4 -1
@@ -55,9 +55,12 @@ struct rrdengine_datafile_list {
55 extern void df_extent_insert(struct extent_info *extent);
56 extern void datafile_list_insert(struct rrdengine_instance *ctx, struct rrdengine_datafile *datafile);
57 extern void datafile_list_delete(struct rrdengine_instance *ctx, struct rrdengine_datafile *datafile);
58 +extern void generate_datafilepath(struct rrdengine_datafile *datafile, char *str, size_t maxlen);
59 +extern int close_data_file(struct rrdengine_datafile *datafile);
60 extern int destroy_data_file(struct rrdengine_datafile *datafile);
61 extern int create_data_file(struct rrdengine_datafile *datafile);
60 -extern void create_new_datafile_pair(struct rrdengine_instance *ctx, unsigned tier, unsigned fileno);
62 +extern int create_new_datafile_pair(struct rrdengine_instance *ctx, unsigned tier, unsigned fileno);
63 extern int init_data_files(struct rrdengine_instance *ctx);
64 +extern void finalize_data_files(struct rrdengine_instance *ctx);
65
66 #endif /* NETDATA_DATAFILE_H */
\ No newline at end of file
database/engine/journalfile.c
+66 -29
@@ -13,7 +13,7 @@ static void flush_transaction_buffer_cb(uv_fs_t* req)
13
14 uv_fs_req_cleanup(req);
15 free(io_descr->buf);
16 - free(io_descr);
16 + freez(io_descr);
17 }
18
19 /* Careful to always call this before creating a new journal file */
@@ -87,7 +87,7 @@ void * wal_get_transaction_buffer(struct rrdengine_worker_config* wc, unsigned s
87 return ctx->commit_log.buf + buf_pos;
88 }
89
90 -static void generate_journalfilepath(struct rrdengine_datafile *datafile, char *str, size_t maxlen)
90 +void generate_journalfilepath(struct rrdengine_datafile *datafile, char *str, size_t maxlen)
91 {
92 (void) snprintf(str, maxlen, "%s/" WALFILE_PREFIX RRDENG_FILE_NUMBER_PRINT_TMPL WALFILE_EXTENSION,
93 datafile->ctx->dbfiles_path, datafile->tier, datafile->fileno);
@@ -100,39 +100,62 @@ void journalfile_init(struct rrdengine_journalfile *journalfile, struct rrdengin
100 journalfile->datafile = datafile;
101 }
102
103 +int close_journal_file(struct rrdengine_journalfile *journalfile, struct rrdengine_datafile *datafile)
104 +{
105 + struct rrdengine_instance *ctx = datafile->ctx;
106 + uv_fs_t req;
107 + int ret;
108 + char path[RRDENG_PATH_MAX];
109 +
110 + generate_journalfilepath(datafile, path, sizeof(path));
111 +
112 + ret = uv_fs_close(NULL, &req, journalfile->file, NULL);
113 + if (ret < 0) {
114 + error("uv_fs_close(%s): %s", path, uv_strerror(ret));
115 + ++ctx->stats.fs_errors;
116 + rrd_stat_atomic_add(&global_fs_errors, 1);
117 + }
118 + uv_fs_req_cleanup(&req);
119 +
120 + return ret;
121 +}
122 +
123 int destroy_journal_file(struct rrdengine_journalfile *journalfile, struct rrdengine_datafile *datafile)
124 {
125 struct rrdengine_instance *ctx = datafile->ctx;
126 uv_fs_t req;
107 - int ret, fd;
108 - char path[1024];
127 + int ret;
128 + char path[RRDENG_PATH_MAX];
129 +
130 + generate_journalfilepath(datafile, path, sizeof(path));
131
132 ret = uv_fs_ftruncate(NULL, &req, journalfile->file, 0, NULL);
133 if (ret < 0) {
112 - fatal("uv_fs_ftruncate: %s", uv_strerror(ret));
134 + error("uv_fs_ftruncate(%s): %s", path, uv_strerror(ret));
135 + ++ctx->stats.fs_errors;
136 + rrd_stat_atomic_add(&global_fs_errors, 1);
137 }
114 - assert(0 == req.result);
138 uv_fs_req_cleanup(&req);
139
140 ret = uv_fs_close(NULL, &req, journalfile->file, NULL);
141 if (ret < 0) {
119 - fatal("uv_fs_close: %s", uv_strerror(ret));
120 - exit(ret);
142 + error("uv_fs_close(%s): %s", path, uv_strerror(ret));
143 + ++ctx->stats.fs_errors;
144 + rrd_stat_atomic_add(&global_fs_errors, 1);
145 }
122 - assert(0 == req.result);
146 uv_fs_req_cleanup(&req);
147
125 - generate_journalfilepath(datafile, path, sizeof(path));
126 - fd = uv_fs_unlink(NULL, &req, path, NULL);
127 - if (fd < 0) {
128 - fatal("uv_fs_fsunlink: %s", uv_strerror(fd));
148 + ret = uv_fs_unlink(NULL, &req, path, NULL);
149 + if (ret < 0) {
150 + error("uv_fs_fsunlink(%s): %s", path, uv_strerror(ret));
151 + ++ctx->stats.fs_errors;
152 + rrd_stat_atomic_add(&global_fs_errors, 1);
153 }
130 - assert(0 == req.result);
154 uv_fs_req_cleanup(&req);
155
156 ++ctx->stats.journalfile_deletions;
157
135 - return 0;
158 + return ret;
159 }
160
161 int create_journal_file(struct rrdengine_journalfile *journalfile, struct rrdengine_datafile *datafile)
@@ -143,13 +166,17 @@ int create_journal_file(struct rrdengine_journalfile *journalfile, struct rrdeng
166 int ret, fd;
167 struct rrdeng_jf_sb *superblock;
168 uv_buf_t iov;
146 - char path[1024];
169 + char path[RRDENG_PATH_MAX];
170
171 generate_journalfilepath(datafile, path, sizeof(path));
172 fd = open_file_direct_io(path, O_CREAT | O_RDWR | O_TRUNC, &file);
173 if (fd < 0) {
151 - fatal("uv_fs_fsopen: %s", uv_strerror(fd));
174 + ++ctx->stats.fs_errors;
175 + rrd_stat_atomic_add(&global_fs_errors, 1);
176 + return fd;
177 }
178 + journalfile->file = file;
179 + ++ctx->stats.journalfile_creations;
180
181 ret = posix_memalign((void *)&superblock, RRDFILE_ALIGNMENT, sizeof(*superblock));
182 if (unlikely(ret)) {
@@ -162,19 +189,21 @@ int create_journal_file(struct rrdengine_journalfile *journalfile, struct rrdeng
189
190 ret = uv_fs_write(NULL, &req, file, &iov, 1, 0, NULL);
191 if (ret < 0) {
165 - fatal("uv_fs_write: %s", uv_strerror(ret));
166 - }
167 - if (req.result < 0) {
168 - fatal("uv_fs_write: %s", uv_strerror((int)req.result));
192 + assert(req.result < 0);
193 + error("uv_fs_write: %s", uv_strerror(ret));
194 + ++ctx->stats.io_errors;
195 + rrd_stat_atomic_add(&global_io_errors, 1);
196 }
197 uv_fs_req_cleanup(&req);
198 free(superblock);
199 + if (ret < 0) {
200 + destroy_journal_file(journalfile, datafile);
201 + return ret;
202 + }
203
173 - journalfile->file = file;
204 journalfile->pos = sizeof(*superblock);
205 ctx->stats.io_write_bytes += sizeof(*superblock);
206 ++ctx->stats.io_write_requests;
177 - ++ctx->stats.journalfile_creations;
207
208 return 0;
209 }
@@ -263,6 +292,8 @@ static void restore_extent_metadata(struct rrdengine_instance *ctx, struct rrden
292 PValue = JudyHSIns(&pg_cache->metrics_index.JudyHS_array, temp_id, sizeof(uuid_t), PJE0);
293 assert(NULL == *PValue); /* TODO: figure out concurrency model */
294 *PValue = page_index = create_page_index(temp_id);
295 + page_index->prev = pg_cache->metrics_index.last_page_index;
296 + pg_cache->metrics_index.last_page_index = page_index;
297 uv_rwlock_wrunlock(&pg_cache->metrics_index.lock);
298 }
299
@@ -398,15 +429,15 @@ int load_journal_file(struct rrdengine_instance *ctx, struct rrdengine_journalfi
429 {
430 uv_fs_t req;
431 uv_file file;
401 - int ret, fd;
432 + int ret, fd, error;
433 uint64_t file_size, max_id;
403 - char path[1024];
434 + char path[RRDENG_PATH_MAX];
435
436 generate_journalfilepath(datafile, path, sizeof(path));
437 fd = open_file_direct_io(path, O_RDWR, &file);
438 if (fd < 0) {
408 - /* if (UV_ENOENT != fd) */
409 - error("uv_fs_fsopen: %s", uv_strerror(fd));
439 + ++ctx->stats.fs_errors;
440 + rrd_stat_atomic_add(&global_fs_errors, 1);
441 return fd;
442 }
443 info("Loading journal file \"%s\".", path);
@@ -433,9 +464,15 @@ int load_journal_file(struct rrdengine_instance *ctx, struct rrdengine_journalfi
464 return 0;
465
466 error:
436 - (void) uv_fs_close(NULL, &req, file, NULL);
467 + error = ret;
468 + ret = uv_fs_close(NULL, &req, file, NULL);
469 + if (ret < 0) {
470 + error("uv_fs_close(%s): %s", path, uv_strerror(ret));
471 + ++ctx->stats.fs_errors;
472 + rrd_stat_atomic_add(&global_fs_errors, 1);
473 + }
474 uv_fs_req_cleanup(&req);
438 - return ret;
475 + return error;
476 }
477
478 void init_commit_log(struct rrdengine_instance *ctx)
database/engine/journalfile.h
+2
@@ -33,9 +33,11 @@ struct transaction_commit_log {
33 unsigned buf_size;
34 };
35
36 +extern void generate_journalfilepath(struct rrdengine_datafile *datafile, char *str, size_t maxlen);
37 extern void journalfile_init(struct rrdengine_journalfile *journalfile, struct rrdengine_datafile *datafile);
38 extern void *wal_get_transaction_buffer(struct rrdengine_worker_config* wc, unsigned size);
39 extern void wal_flush_transaction_buffer(struct rrdengine_worker_config* wc);
40 +extern int close_journal_file(struct rrdengine_journalfile *journalfile, struct rrdengine_datafile *datafile);
41 extern int destroy_journal_file(struct rrdengine_journalfile *journalfile, struct rrdengine_datafile *datafile);
42 extern int create_journal_file(struct rrdengine_journalfile *journalfile, struct rrdengine_datafile *datafile);
43 extern int load_journal_file(struct rrdengine_instance *ctx, struct rrdengine_journalfile *journalfile,
database/engine/pagecache.c
+66 -3
@@ -287,7 +287,7 @@ static void pg_cache_evict_unsafe(struct rrdengine_instance *ctx, struct rrdeng_
287 {
288 struct page_cache_descr *pg_cache_descr = descr->pg_cache_descr;
289
290 - free(pg_cache_descr->page);
290 + freez(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);
@@ -330,7 +330,7 @@ static int pg_cache_try_evict_one_page_unsafe(struct rrdengine_instance *ctx)
330 return 1;
331 }
332 rrdeng_page_descr_mutex_unlock(ctx, descr);
333 - };
333 + }
334 uv_rwlock_wrunlock(&pg_cache->replaceQ.lock);
335
336 /* failed to evict */
@@ -594,7 +594,7 @@ struct pg_cache_page_index *
594 }
595 rrdeng_page_descr_mutex_unlock(ctx, descr);
596
597 - };
597 + }
598 uv_rwlock_rdunlock(&page_index->lock);
599
600 failed_to_reserve = 0;
@@ -767,6 +767,7 @@ struct pg_cache_page_index *create_page_index(uuid_t *id)
767 assert(0 == uv_rwlock_init(&page_index->lock));
768 page_index->oldest_time = INVALID_TIME;
769 page_index->latest_time = INVALID_TIME;
770 + page_index->prev = NULL;
771
772 return page_index;
773 }
@@ -776,6 +777,7 @@ static void init_metrics_index(struct rrdengine_instance *ctx)
777 struct page_cache *pg_cache = &ctx->pg_cache;
778
779 pg_cache->metrics_index.JudyHS_array = (Pvoid_t) NULL;
780 + pg_cache->metrics_index.last_page_index = NULL;
781 assert(0 == uv_rwlock_init(&pg_cache->metrics_index.lock));
782 }
783
@@ -809,4 +811,65 @@ void init_page_cache(struct rrdengine_instance *ctx)
811 init_metrics_index(ctx);
812 init_replaceQ(ctx);
813 init_commited_page_index(ctx);
814 +}
815 +
816 +void free_page_cache(struct rrdengine_instance *ctx)
817 +{
818 + struct page_cache *pg_cache = &ctx->pg_cache;
819 + Word_t ret_Judy, bytes_freed = 0;
820 + Pvoid_t *PValue;
821 + struct pg_cache_page_index *page_index, *prev_page_index;
822 + Word_t Index;
823 + struct rrdeng_page_descr *descr;
824 + struct page_cache_descr *pg_cache_descr;
825 +
826 + /* Free commited page index */
827 + ret_Judy = JudyLFreeArray(&pg_cache->commited_page_index.JudyL_array, PJE0);
828 + assert(NULL == pg_cache->commited_page_index.JudyL_array);
829 + bytes_freed += ret_Judy;
830 +
831 + for (page_index = pg_cache->metrics_index.last_page_index ;
832 + page_index != NULL ;
833 + page_index = prev_page_index) {
834 + prev_page_index = page_index->prev;
835 +
836 + /* Find first page in range */
837 + Index = (Word_t) 0;
838 + PValue = JudyLFirst(page_index->JudyL_array, &Index, PJE0);
839 + if (likely(NULL != PValue)) {
840 + descr = *PValue;
841 + }
842 + while (descr != NULL) {
843 + /* Iterate all page descriptors of this metric */
844 +
845 + if (descr->pg_cache_descr_state & PG_CACHE_DESCR_ALLOCATED) {
846 + /* Check rrdenglocking.c */
847 + pg_cache_descr = descr->pg_cache_descr;
848 + if (pg_cache_descr->flags & RRD_PAGE_POPULATED) {
849 + freez(pg_cache_descr->page);
850 + bytes_freed += RRDENG_BLOCK_SIZE;
851 + }
852 + rrdeng_destroy_pg_cache_descr(ctx, pg_cache_descr);
853 + bytes_freed += sizeof(*pg_cache_descr);
854 + }
855 + freez(descr);
856 + bytes_freed += sizeof(*descr);
857 +
858 + PValue = JudyLNext(page_index->JudyL_array, &Index, PJE0);
859 + descr = unlikely(NULL == PValue) ? NULL : *PValue;
860 + }
861 +
862 + /* Free page index */
863 + ret_Judy = JudyLFreeArray(&page_index->JudyL_array, PJE0);
864 + assert(NULL == page_index->JudyL_array);
865 + bytes_freed += ret_Judy;
866 + freez(page_index);
867 + bytes_freed += sizeof(*page_index);
868 + }
869 + /* Free metrics index */
870 + ret_Judy = JudyHSFreeArray(&pg_cache->metrics_index.JudyHS_array, PJE0);
871 + assert(NULL == pg_cache->metrics_index.JudyHS_array);
872 + bytes_freed += ret_Judy;
873 +
874 + info("Freed %lu bytes of memory from page cache.", bytes_freed);
875 }
\ No newline at end of file
database/engine/pagecache.h
+4
@@ -84,12 +84,15 @@ struct pg_cache_page_index {
84 * It's also written by the data deletion workqueue when data collection is disabled for this metric.
85 */
86 usec_t latest_time;
87 +
88 + struct pg_cache_page_index *prev;
89 };
90
91 /* maps UUIDs to page indices */
92 struct pg_cache_metrics_index {
93 uv_rwlock_t lock;
94 Pvoid_t JudyHS_array;
95 + struct pg_cache_page_index *last_page_index;
96 };
97
98 /* gathers dirty pages to be written on disk */
@@ -153,6 +156,7 @@ extern struct rrdeng_page_descr *
156 usec_t point_in_time);
157 extern struct pg_cache_page_index *create_page_index(uuid_t *id);
158 extern void init_page_cache(struct rrdengine_instance *ctx);
159 +extern void free_page_cache(struct rrdengine_instance *ctx);
160 extern void pg_cache_add_new_metric_time(struct pg_cache_page_index *page_index, struct rrdeng_page_descr *descr);
161 extern void pg_cache_update_metric_times(struct pg_cache_page_index *page_index);
162
database/engine/rrdengine.c
+108 -48
@@ -3,6 +3,10 @@
3
4 #include "rrdengine.h"
5
6 +rrdeng_stats_t global_io_errors = 0;
7 +rrdeng_stats_t global_fs_errors = 0;
8 +rrdeng_stats_t rrdeng_reserved_file_descriptors = 0;
9 +
10 void sanity_check(void)
11 {
12 /* Magic numbers must fit in the super-blocks */
@@ -33,7 +37,6 @@ void read_extent_cb(uv_fs_t* req)
37 unsigned i, j, count;
38 void *page, *uncompressed_buf = NULL;
39 uint32_t payload_length, payload_offset, page_offset, uncompressed_payload_length;
36 - struct rrdengine_datafile *datafile;
40 /* persistent structures */
41 struct rrdeng_df_extent_header *header;
42 struct rrdeng_df_extent_trailer *trailer;
@@ -55,9 +58,13 @@ void read_extent_cb(uv_fs_t* req)
58 crc = crc32(0L, Z_NULL, 0);
59 crc = crc32(crc, xt_io_descr->buf, xt_io_descr->bytes - sizeof(*trailer));
60 ret = crc32cmp(trailer->checksum, crc);
58 - datafile = xt_io_descr->descr_array[0]->extent->datafile;
59 - debug(D_RRDENGINE, "%s: Extent at offset %"PRIu64"(%u) was read from datafile %u-%u. CRC32 check: %s", __func__,
60 - xt_io_descr->pos, xt_io_descr->bytes, datafile->tier, datafile->fileno, ret ? "FAILED" : "SUCCEEDED");
61 +#ifdef NETDATA_INTERNAL_CHECKS
62 + {
63 + struct rrdengine_datafile *datafile = xt_io_descr->descr_array[0]->extent->datafile;
64 + debug(D_RRDENGINE, "%s: Extent at offset %"PRIu64"(%u) was read from datafile %u-%u. CRC32 check: %s", __func__,
65 + xt_io_descr->pos, xt_io_descr->bytes, datafile->tier, datafile->fileno, ret ? "FAILED" : "SUCCEEDED");
66 + }
67 +#endif
68 if (unlikely(ret)) {
69 /* TODO: handle errors */
70 exit(UV_EIO);
@@ -112,14 +119,14 @@ void read_extent_cb(uv_fs_t* req)
119 rrdeng_page_descr_mutex_unlock(ctx, descr);
120 }
121 if (RRD_NO_COMPRESSION != header->compression_algorithm) {
115 - free(uncompressed_buf);
122 + freez(uncompressed_buf);
123 }
124 if (xt_io_descr->completion)
125 complete(xt_io_descr->completion);
126 cleanup:
127 uv_fs_req_cleanup(req);
128 free(xt_io_descr->buf);
122 - free(xt_io_descr);
129 + freez(xt_io_descr);
130 }
131
132
@@ -144,7 +151,7 @@ static void do_read_extent(struct rrdengine_worker_config* wc,
151 ret = posix_memalign((void *)&xt_io_descr->buf, RRDFILE_ALIGNMENT, ALIGN_BYTES_CEILING(size_bytes));
152 if (unlikely(ret)) {
153 fatal("posix_memalign:%s", strerror(ret));
147 - /* free(xt_io_descr);
154 + /* freez(xt_io_descr);
155 return;*/
156 }
157 for (i = 0 ; i < count; ++i) {
@@ -233,7 +240,6 @@ void flush_pages_cb(uv_fs_t* req)
240 struct extent_io_descriptor *xt_io_descr;
241 struct rrdeng_page_descr *descr;
242 struct page_cache_descr *pg_cache_descr;
236 - struct rrdengine_datafile *datafile;
243 int ret;
244 unsigned i, count;
245 Word_t commit_id;
@@ -243,10 +249,13 @@ void flush_pages_cb(uv_fs_t* req)
249 error("%s: uv_fs_write: %s", __func__, uv_strerror((int)req->result));
250 goto cleanup;
251 }
246 - datafile = xt_io_descr->descr_array[0]->extent->datafile;
247 - debug(D_RRDENGINE, "%s: Extent at offset %"PRIu64"(%u) was written to datafile %u-%u. Waking up waiters.",
248 - __func__, xt_io_descr->pos, xt_io_descr->bytes, datafile->tier, datafile->fileno);
249 -
252 +#ifdef NETDATA_INTERNAL_CHECKS
253 + {
254 + struct rrdengine_datafile *datafile = xt_io_descr->descr_array[0]->extent->datafile;
255 + debug(D_RRDENGINE, "%s: Extent at offset %"PRIu64"(%u) was written to datafile %u-%u. Waking up waiters.",
256 + __func__, xt_io_descr->pos, xt_io_descr->bytes, datafile->tier, datafile->fileno);
257 + }
258 +#endif
259 count = xt_io_descr->descr_count;
260 for (i = 0 ; i < count ; ++i) {
261 /* care, we don't hold the descriptor mutex */
@@ -273,7 +282,7 @@ void flush_pages_cb(uv_fs_t* req)
282 cleanup:
283 uv_fs_req_cleanup(req);
284 free(xt_io_descr->buf);
276 - free(xt_io_descr);
285 + freez(xt_io_descr);
286 }
287
288 /*
@@ -353,7 +362,7 @@ static int do_flush_pages(struct rrdengine_worker_config* wc, int force, struct
362 ret = posix_memalign((void *)&xt_io_descr->buf, RRDFILE_ALIGNMENT, ALIGN_BYTES_CEILING(size_bytes));
363 if (unlikely(ret)) {
364 fatal("posix_memalign:%s", strerror(ret));
356 - /* free(xt_io_descr);*/
365 + /* freez(xt_io_descr);*/
366 }
367 (void) memcpy(xt_io_descr->descr_array, eligible_pages, sizeof(struct rrdeng_page_descr *) * count);
368 xt_io_descr->descr_count = count;
@@ -405,7 +414,7 @@ static int do_flush_pages(struct rrdengine_worker_config* wc, int force, struct
414 ctx->stats.after_compress_bytes += compressed_size;
415 debug(D_RRDENGINE, "LZ4 compressed %"PRIu32" bytes to %d bytes.", uncompressed_payload_length, compressed_size);
416 (void) memcpy(xt_io_descr->buf + payload_offset, compressed_buf, compressed_size);
408 - free(compressed_buf);
417 + freez(compressed_buf);
418 size_bytes = payload_offset + compressed_size + sizeof(*trailer);
419 header->payload_length = compressed_size;
420 break;
@@ -443,23 +452,36 @@ static void after_delete_old_data(uv_work_t *req, int status)
452 struct rrdengine_worker_config* wc = &ctx->worker_config;
453 struct rrdengine_datafile *datafile;
454 struct rrdengine_journalfile *journalfile;
446 - unsigned bytes;
455 + unsigned deleted_bytes, journalfile_bytes, datafile_bytes;
456 + int ret;
457 + char path[RRDENG_PATH_MAX];
458
459 (void)status;
460 datafile = ctx->datafiles.first;
461 journalfile = datafile->journalfile;
451 - bytes = datafile->pos + journalfile->pos;
462 + datafile_bytes = datafile->pos;
463 + journalfile_bytes = journalfile->pos;
464 + deleted_bytes = 0;
465
466 + info("Deleting data and journal file pair.");
467 datafile_list_delete(ctx, datafile);
454 - destroy_journal_file(journalfile, datafile);
455 - destroy_data_file(datafile);
456 - info("Deleted data file \""DATAFILE_PREFIX RRDENG_FILE_NUMBER_PRINT_TMPL DATAFILE_EXTENSION"\".",
457 - datafile->tier, datafile->fileno);
458 - free(journalfile);
459 - free(datafile);
468 + ret = destroy_journal_file(journalfile, datafile);
469 + if (!ret) {
470 + generate_journalfilepath(datafile, path, sizeof(path));
471 + info("Deleted journal file \"%s\".", path);
472 + deleted_bytes += journalfile_bytes;
473 + }
474 + ret = destroy_data_file(datafile);
475 + if (!ret) {
476 + generate_datafilepath(datafile, path, sizeof(path));
477 + info("Deleted data file \"%s\".", path);
478 + deleted_bytes += datafile_bytes;
479 + }
480 + freez(journalfile);
481 + freez(datafile);
482
461 - ctx->disk_space -= bytes;
462 - info("Reclaimed %u bytes of disk space.", bytes);
483 + ctx->disk_space -= deleted_bytes;
484 + info("Reclaimed %u bytes of disk space.", deleted_bytes);
485
486 /* unfreeze command processing */
487 wc->now_deleting.data = NULL;
@@ -485,7 +507,7 @@ static void delete_old_data(uv_work_t *req)
507 pg_cache_punch_hole(ctx, descr, 0);
508 }
509 next = extent->next;
488 - free(extent);
510 + freez(extent);
511 }
512 }
513
@@ -495,6 +517,7 @@ void rrdeng_test_quota(struct rrdengine_worker_config* wc)
517 struct rrdengine_datafile *datafile;
518 unsigned current_size, target_size;
519 uint8_t out_of_space, only_one_datafile;
520 + int ret;
521
522 out_of_space = 0;
523 if (unlikely(ctx->disk_space > ctx->max_disk_space)) {
@@ -509,7 +532,10 @@ void rrdeng_test_quota(struct rrdengine_worker_config* wc)
532 if (unlikely(current_size >= target_size || (out_of_space && only_one_datafile))) {
533 /* Finalize data and journal file and create a new pair */
534 wal_flush_transaction_buffer(wc);
512 - create_new_datafile_pair(ctx, 1, datafile->fileno + 1);
535 + ret = create_new_datafile_pair(ctx, 1, ctx->last_fileno + 1);
536 + if (likely(!ret)) {
537 + ++ctx->last_fileno;
538 + }
539 }
540 if (unlikely(out_of_space)) {
541 /* delete old data */
@@ -517,18 +543,30 @@ void rrdeng_test_quota(struct rrdengine_worker_config* wc)
543 /* already deleting data */
544 return;
545 }
520 - info("Deleting data file \""DATAFILE_PREFIX RRDENG_FILE_NUMBER_PRINT_TMPL DATAFILE_EXTENSION"\".",
521 - ctx->datafiles.first->tier, ctx->datafiles.first->fileno);
546 + if (NULL == ctx->datafiles.first->next) {
547 + error("Cannot delete data file \"%s/"DATAFILE_PREFIX RRDENG_FILE_NUMBER_PRINT_TMPL DATAFILE_EXTENSION"\""
548 + " to reclaim space, there are no other file pairs left.",
549 + ctx->dbfiles_path, ctx->datafiles.first->tier, ctx->datafiles.first->fileno);
550 + return;
551 + }
552 + info("Deleting data file \"%s/"DATAFILE_PREFIX RRDENG_FILE_NUMBER_PRINT_TMPL DATAFILE_EXTENSION"\".",
553 + ctx->dbfiles_path, ctx->datafiles.first->tier, ctx->datafiles.first->fileno);
554 wc->now_deleting.data = ctx;
523 - uv_queue_work(wc->loop, &wc->now_deleting, delete_old_data, after_delete_old_data);
555 + assert(0 == uv_queue_work(wc->loop, &wc->now_deleting, delete_old_data, after_delete_old_data));
556 }
557 }
558
559 +/* return 0 on success */
560 int init_rrd_files(struct rrdengine_instance *ctx)
561 {
562 return init_data_files(ctx);
563 }
564
565 +void finalize_rrd_files(struct rrdengine_instance *ctx)
566 +{
567 + return finalize_data_files(ctx);
568 +}
569 +
570 void rrdeng_init_cmd_queue(struct rrdengine_worker_config* wc)
571 {
572 wc->cmd_queue.head = wc->cmd_queue.tail = 0;
@@ -596,7 +634,6 @@ void async_cb(uv_async_t *handle)
634 void timer_cb(uv_timer_t* handle)
635 {
636 struct rrdengine_worker_config* wc = handle->data;
599 - struct rrdengine_instance *ctx = wc->ctx;
637
638 uv_stop(handle->loop);
639 uv_update_time(handle->loop);
@@ -616,7 +653,7 @@ void timer_cb(uv_timer_t* handle)
653 #ifdef NETDATA_INTERNAL_CHECKS
654 {
655 char buf[4096];
619 - debug(D_RRDENGINE, "%s", get_rrdeng_statistics(ctx, buf, sizeof(buf)));
656 + debug(D_RRDENGINE, "%s", get_rrdeng_statistics(wc->ctx, buf, sizeof(buf)));
657 }
658 #endif
659 }
@@ -631,7 +668,7 @@ void rrdeng_worker(void* arg)
668 struct rrdengine_worker_config* wc = arg;
669 struct rrdengine_instance *ctx = wc->ctx;
670 uv_loop_t* loop;
634 - int shutdown;
671 + int shutdown, ret;
672 enum rrdeng_opcode opcode;
673 uv_timer_t timer_req;
674 struct rrdeng_cmd cmd;
@@ -639,22 +676,35 @@ void rrdeng_worker(void* arg)
676 rrdeng_init_cmd_queue(wc);
677
678 loop = wc->loop = mallocz(sizeof(uv_loop_t));
642 - uv_loop_init(loop);
679 + ret = uv_loop_init(loop);
680 + if (ret) {
681 + error("uv_loop_init(): %s", uv_strerror(ret));
682 + goto error_after_loop_init;
683 + }
684 loop->data = wc;
685
645 - uv_async_init(wc->loop, &wc->async, async_cb);
686 + ret = uv_async_init(wc->loop, &wc->async, async_cb);
687 + if (ret) {
688 + error("uv_async_init(): %s", uv_strerror(ret));
689 + goto error_after_async_init;
690 + }
691 wc->async.data = wc;
692
693 wc->now_deleting.data = NULL;
694
695 /* dirty page flushing timer */
651 - uv_timer_init(loop, &timer_req);
696 + ret = uv_timer_init(loop, &timer_req);
697 + if (ret) {
698 + error("uv_timer_init(): %s", uv_strerror(ret));
699 + goto error_after_timer_init;
700 + }
701 timer_req.data = wc;
702
703 + wc->error = 0;
704 /* wake up initialization thread */
705 complete(&ctx->rrdengine_completion);
706
657 - uv_timer_start(&timer_req, timer_cb, TIMER_PERIOD_MS, TIMER_PERIOD_MS);
707 + assert(0 == uv_timer_start(&timer_req, timer_cb, TIMER_PERIOD_MS, TIMER_PERIOD_MS));
708 shutdown = 0;
709 while (shutdown == 0 || uv_loop_alive(loop)) {
710 uv_run(loop, UV_RUN_DEFAULT);
@@ -669,12 +719,6 @@ void rrdeng_worker(void* arg)
719 break;
720 case RRDENG_SHUTDOWN:
721 shutdown = 1;
672 - if (unlikely(wc->now_deleting.data)) {
673 - /* postpone shutdown until after deletion */
674 - info("Postponing shutting RRD engine event loop down until after datafile deletion is finished.");
675 - rrdeng_enq_cmd(wc, &cmd);
676 - break;
677 - }
722 /*
723 * uv_async_send after uv_close does not seem to crash in linux at the moment,
724 * it is however undocumented behaviour and we need to be aware if this becomes
@@ -683,10 +727,6 @@ void rrdeng_worker(void* arg)
727 uv_close((uv_handle_t *)&wc->async, NULL);
728 assert(0 == uv_timer_stop(&timer_req));
729 uv_close((uv_handle_t *)&timer_req, NULL);
686 - info("Shutting down RRD engine event loop.");
687 - while (do_flush_pages(wc, 1, NULL)) {
688 - ; /* Force flushing of all commited pages. */
689 - }
730 break;
731 case RRDENG_READ_PAGE:
732 do_read_extent(wc, &cmd.read_page.page_cache_descr, 1, 0);
@@ -716,6 +756,13 @@ void rrdeng_worker(void* arg)
756 } while (opcode != RRDENG_NOOP);
757 }
758 /* cleanup operations of the event loop */
759 + if (unlikely(wc->now_deleting.data)) {
760 + info("Postponing shutting RRD engine event loop down until after datafile deletion is finished.");
761 + }
762 + info("Shutting down RRD engine event loop.");
763 + while (do_flush_pages(wc, 1, NULL)) {
764 + ; /* Force flushing of all commited pages. */
765 + }
766 wal_flush_transaction_buffer(wc);
767 uv_run(loop, UV_RUN_DEFAULT);
768
@@ -724,7 +771,20 @@ void rrdeng_worker(void* arg)
771 uv_cond_destroy(&wc->cmd_cond);
772 /* uv_mutex_destroy(&wc->cmd_mutex); */
773 assert(0 == uv_loop_close(loop));
727 - free(loop);
774 + freez(loop);
775 +
776 + return;
777 +
778 +error_after_timer_init:
779 + uv_close((uv_handle_t *)&wc->async, NULL);
780 +error_after_async_init:
781 + assert(0 == uv_loop_close(loop));
782 +error_after_loop_init:
783 + freez(loop);
784 +
785 + wc->error = UV_EAGAIN;
786 + /* wake up initialization thread */
787 + complete(&ctx->rrdengine_completion);
788 }
789
790
database/engine/rrdengine.h
+13 -1
@@ -112,6 +112,8 @@ struct rrdengine_worker_config {
112 uv_cond_t cmd_cond;
113 volatile unsigned queue_size;
114 struct rrdeng_cmdqueue cmd_queue;
115 +
116 + int error;
117 };
118
119 /*
@@ -144,10 +146,18 @@ struct rrdengine_statistics {
146 rrdeng_stats_t journalfile_creations;
147 rrdeng_stats_t journalfile_deletions;
148 rrdeng_stats_t page_cache_descriptors;
149 + rrdeng_stats_t io_errors;
150 + rrdeng_stats_t fs_errors;
151 };
152
153 +/* I/O errors global counter */
154 +extern rrdeng_stats_t global_io_errors;
155 +/* File-System errors global counter */
156 +extern rrdeng_stats_t global_fs_errors;
157 +/* number of File-Descriptors that have been reserved by dbengine */
158 +extern rrdeng_stats_t rrdeng_reserved_file_descriptors;
159 +
160 struct rrdengine_instance {
150 - rrdengine_state_t rrdengine_state;
161 struct rrdengine_worker_config worker_config;
162 struct completion rrdengine_completion;
163 struct page_cache pg_cache;
@@ -157,6 +167,7 @@ struct rrdengine_instance {
167 char dbfiles_path[FILENAME_MAX+1];
168 uint64_t disk_space;
169 uint64_t max_disk_space;
170 + unsigned last_fileno; /* newest index of datafile and journalfile */
171 unsigned long max_cache_pages;
172 unsigned long cache_pages_low_watermark;
173
@@ -165,6 +176,7 @@ struct rrdengine_instance {
176
177 extern void sanity_check(void);
178 extern int init_rrd_files(struct rrdengine_instance *ctx);
179 +extern void finalize_rrd_files(struct rrdengine_instance *ctx);
180 extern void rrdeng_test_quota(struct rrdengine_worker_config* wc);
181 extern void rrdeng_worker(void* arg);
182 extern void rrdeng_enq_cmd(struct rrdengine_worker_config* wc, struct rrdeng_cmd *cmd);
database/engine/rrdengineapi.c
+52 -20
@@ -55,6 +55,8 @@ void rrdeng_store_metric_init(RRDDIM *rd)
55 PValue = JudyHSIns(&pg_cache->metrics_index.JudyHS_array, &temp_id, sizeof(uuid_t), PJE0);
56 assert(NULL == *PValue); /* TODO: figure out concurrency model */
57 *PValue = page_index = create_page_index(&temp_id);
58 + page_index->prev = pg_cache->metrics_index.last_page_index;
59 + pg_cache->metrics_index.last_page_index = page_index;
60 uv_rwlock_wrunlock(&pg_cache->metrics_index.lock);
61 }
62 rd->state->rrdeng_uuid = &page_index->id;
@@ -119,9 +121,9 @@ void rrdeng_store_metric_flush_current_page(RRDDIM *rd)
121 handle->prev_descr = descr;
122 }
123 } else {
122 - free(descr->pg_cache_descr->page);
124 + freez(descr->pg_cache_descr->page);
125 rrdeng_destroy_pg_cache_descr(ctx, descr->pg_cache_descr);
124 - free(descr);
126 + freez(descr);
127 }
128 handle->descr = NULL;
129 }
@@ -434,7 +436,13 @@ void *rrdeng_get_page(struct rrdengine_instance *ctx, uuid_t *id, usec_t point_i
436 return pg_cache_descr->page;
437 }
438
437 -void rrdeng_get_28_statistics(struct rrdengine_instance *ctx, unsigned long long *array)
439 +/*
440 + * Gathers Database Engine statistics.
441 + * Careful when modifying this function.
442 + * You must not change the indices of the statistics or user code will break.
443 + * You must not exceed RRDENG_NR_STATS or it will crash.
444 + */
445 +void rrdeng_get_33_statistics(struct rrdengine_instance *ctx, unsigned long long *array)
446 {
447 struct page_cache *pg_cache = &ctx->pg_cache;
448
@@ -466,7 +474,12 @@ void rrdeng_get_28_statistics(struct rrdengine_instance *ctx, unsigned long long
474 array[25] = (uint64_t)ctx->stats.journalfile_creations;
475 array[26] = (uint64_t)ctx->stats.journalfile_deletions;
476 array[27] = (uint64_t)ctx->stats.page_cache_descriptors;
469 - assert(RRDENG_NR_STATS == 28);
477 + array[28] = (uint64_t)ctx->stats.io_errors;
478 + array[29] = (uint64_t)ctx->stats.fs_errors;
479 + array[30] = (uint64_t)global_io_errors;
480 + array[31] = (uint64_t)global_fs_errors;
481 + array[32] = (uint64_t)rrdeng_reserved_file_descriptors;
482 + assert(RRDENG_NR_STATS == 33);
483 }
484
485 /* Releases reference to page */
@@ -477,14 +490,29 @@ void rrdeng_put_page(struct rrdengine_instance *ctx, void *handle)
490 }
491
492 /*
480 - * Returns 0 on success, 1 on error
493 + * Returns 0 on success, negative on error
494 */
495 int rrdeng_init(struct rrdengine_instance **ctxp, char *dbfiles_path, unsigned page_cache_mb, unsigned disk_space_mb)
496 {
497 struct rrdengine_instance *ctx;
498 int error;
499 + uint32_t max_open_files;
500
501 sanity_check();
502 +
503 + max_open_files = rlimit_nofile.rlim_cur / 4;
504 +
505 + /* reserve RRDENG_FD_BUDGET_PER_INSTANCE file descriptors for this instance */
506 + rrd_stat_atomic_add(&rrdeng_reserved_file_descriptors, RRDENG_FD_BUDGET_PER_INSTANCE);
507 + if (rrdeng_reserved_file_descriptors > max_open_files) {
508 + error("Exceeded the budget of available file descriptors (%u/%u), cannot create new dbengine instance.",
509 + (unsigned)rrdeng_reserved_file_descriptors, (unsigned)max_open_files);
510 +
511 + rrd_stat_atomic_add(&global_fs_errors, 1);
512 + rrd_stat_atomic_add(&rrdeng_reserved_file_descriptors, -RRDENG_FD_BUDGET_PER_INSTANCE);
513 + return UV_EMFILE;
514 + }
515 +
516 if (NULL == ctxp) {
517 /* for testing */
518 ctx = &default_global_ctx;
@@ -492,10 +520,6 @@ int rrdeng_init(struct rrdengine_instance **ctxp, char *dbfiles_path, unsigned p
520 } else {
521 *ctxp = ctx = callocz(1, sizeof(*ctx));
522 }
495 - if (ctx->rrdengine_state != RRDENGINE_STATUS_UNINITIALIZED) {
496 - return 1;
497 - }
498 - ctx->rrdengine_state = RRDENGINE_STATUS_INITIALIZING;
523 ctx->global_compress_alg = RRD_LZ4;
524 if (page_cache_mb < RRDENG_MIN_PAGE_CACHE_SIZE_MB)
525 page_cache_mb = RRDENG_MIN_PAGE_CACHE_SIZE_MB;
@@ -514,11 +538,7 @@ int rrdeng_init(struct rrdengine_instance **ctxp, char *dbfiles_path, unsigned p
538 init_commit_log(ctx);
539 error = init_rrd_files(ctx);
540 if (error) {
517 - ctx->rrdengine_state = RRDENGINE_STATUS_UNINITIALIZED;
518 - if (ctx != &default_global_ctx) {
519 - freez(ctx);
520 - }
521 - return 1;
541 + goto error_after_init_rrd_files;
542 }
543
544 init_completion(&ctx->rrdengine_completion);
@@ -526,9 +546,21 @@ int rrdeng_init(struct rrdengine_instance **ctxp, char *dbfiles_path, unsigned p
546 /* wait for worker thread to initialize */
547 wait_for_completion(&ctx->rrdengine_completion);
548 destroy_completion(&ctx->rrdengine_completion);
529 -
530 - ctx->rrdengine_state = RRDENGINE_STATUS_INITIALIZED;
549 + if (ctx->worker_config.error) {
550 + goto error_after_rrdeng_worker;
551 + }
552 return 0;
553 +
554 +error_after_rrdeng_worker:
555 + finalize_rrd_files(ctx);
556 +error_after_init_rrd_files:
557 + free_page_cache(ctx);
558 + if (ctx != &default_global_ctx) {
559 + freez(ctx);
560 + *ctxp = NULL;
561 + }
562 + rrd_stat_atomic_add(&rrdeng_reserved_file_descriptors, -RRDENG_FD_BUDGET_PER_INSTANCE);
563 + return UV_EIO;
564 }
565
566 /*
@@ -539,10 +571,6 @@ int rrdeng_exit(struct rrdengine_instance *ctx)
571 struct rrdeng_cmd cmd;
572
573 if (NULL == ctx) {
542 - /* TODO: move to per host basis */
543 - ctx = &default_global_ctx;
544 - }
545 - if (ctx->rrdengine_state != RRDENGINE_STATUS_INITIALIZED) {
574 return 1;
575 }
576
@@ -552,8 +580,12 @@ int rrdeng_exit(struct rrdengine_instance *ctx)
580
581 assert(0 == uv_thread_join(&ctx->worker_config.thread));
582
583 + finalize_rrd_files(ctx);
584 + free_page_cache(ctx);
585 +
586 if (ctx != &default_global_ctx) {
587 freez(ctx);
588 }
589 + rrd_stat_atomic_add(&rrdeng_reserved_file_descriptors, -RRDENG_FD_BUDGET_PER_INSTANCE);
590 return 0;
591 }
\ No newline at end of file
database/engine/rrdengineapi.h
+4 -2
@@ -8,7 +8,9 @@
8 #define RRDENG_MIN_PAGE_CACHE_SIZE_MB (32)
9 #define RRDENG_MIN_DISK_SPACE_MB (256)
10
11 -#define RRDENG_NR_STATS (28)
11 +#define RRDENG_NR_STATS (33)
12 +
13 +#define RRDENG_FD_BUDGET_PER_INSTANCE (50)
14
15 extern int default_rrdeng_page_cache_mb;
16 extern int default_rrdeng_disk_quota_mb;
@@ -30,7 +32,7 @@ extern int rrdeng_load_metric_is_finished(struct rrddim_query_handle *rrdimm_han
32 extern void rrdeng_load_metric_finalize(struct rrddim_query_handle *rrdimm_handle);
33 extern time_t rrdeng_metric_latest_time(RRDDIM *rd);
34 extern time_t rrdeng_metric_oldest_time(RRDDIM *rd);
33 -extern void rrdeng_get_28_statistics(struct rrdengine_instance *ctx, unsigned long long *array);
35 +extern void rrdeng_get_33_statistics(struct rrdengine_instance *ctx, unsigned long long *array);
36
37 /* must call once before using anything */
38 extern int rrdeng_init(struct rrdengine_instance **ctxp, char *dbfiles_path, unsigned page_cache_mb,
database/engine/rrdenginelib.c
+1 -1
@@ -103,7 +103,7 @@ int open_file_direct_io(char *path, int flags, uv_file *file)
103 error("File \"%s\" does not support direct I/O, falling back to buffered I/O.", path);
104 } else {
105 error("Failed to open file \"%s\".", path);
106 - return fd;
106 + --direct; /* break the loop */
107 }
108 } else {
109 assert(req.result >= 0);
database/engine/rrdenginelib.h
+2
@@ -31,6 +31,8 @@ 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 +#define RRDENG_PATH_MAX (4096)
35 +
36 /* returns old *ptr value */
37 static inline unsigned long ulong_compare_and_swap(volatile unsigned long *ptr,
38 unsigned long oldval, unsigned long newval)
database/engine/rrdenglocking.c
+2 -2
@@ -22,7 +22,7 @@ void rrdeng_destroy_pg_cache_descr(struct rrdengine_instance *ctx, struct page_c
22 {
23 uv_cond_destroy(&pg_cache_descr->cond);
24 uv_mutex_destroy(&pg_cache_descr->mutex);
25 - free(pg_cache_descr);
25 + freez(pg_cache_descr);
26 rrd_stat_atomic_add(&ctx->stats.page_cache_descriptors, -1);
27 }
28
@@ -102,7 +102,6 @@ void rrdeng_page_descr_mutex_unlock(struct rrdengine_instance *ctx, struct rrden
102 we_locked = 0;
103 while (1) { /* spin */
104 old_state = descr->pg_cache_descr_state;
105 - assert(old_state & PG_CACHE_DESCR_ALLOCATED);
105 old_users = old_state >> PG_CACHE_DESCR_SHIFT;
106
107 if (unlikely(we_locked)) {
@@ -119,6 +118,7 @@ void rrdeng_page_descr_mutex_unlock(struct rrdengine_instance *ctx, struct rrden
118 assert(0 == old_users);
119 continue; /* spin */
120 }
121 + assert(old_state & PG_CACHE_DESCR_ALLOCATED);
122 pg_cache_descr = descr->pg_cache_descr;
123 /* caller is the only page cache descriptor user and there are no pending references on the page */
124 if ((old_state & PG_CACHE_DESCR_DESTROY) && (1 == old_users) &&
health/Makefile.am
+1
@@ -86,4 +86,5 @@ dist_healthconfig_DATA = \
86 health.d/wmi.conf \
87 health.d/x509check.conf \
88 health.d/zfs.conf \
89 + health.d/dbengine.conf \
90 $(NULL)
health/health.d/dbengine.conf new
+26
@@ -0,0 +1,26 @@
1 +
2 +# you can disable an alarm notification by setting the 'to' line to: silent
3 +
4 + alarm: 10min_dbengine_global_fs_errors
5 + on: netdata.dbengine_global_errors
6 + os: linux freebsd macos
7 + hosts: *
8 + lookup: sum -10m unaligned of FS errors
9 + units: errors
10 + every: 10s
11 + crit: $this > 0
12 + delay: down 15m multiplier 1.5 max 1h
13 + info: number of File-System errors dbengine came across the last 10 minutes (too many open files, wrong permissions etc)
14 + to: sysadmin
15 +
16 + alarm: 10min_dbengine_global_io_errors
17 + on: netdata.dbengine_global_errors
18 + os: linux freebsd macos
19 + hosts: *
20 + lookup: sum -10m unaligned of I/O errors
21 + units: errors
22 + every: 10s
23 + crit: $this > 0
24 + delay: down 1h multiplier 1.5 max 3h
25 + info: number of IO errors dbengine came across the last 10 minutes (out of space, bad disk etc)
26 + to: sysadmin
web/server/static/static-threaded.c
+1 -1
@@ -474,7 +474,7 @@ void *socket_listen_main_static_threaded(void *ptr) {
474
475 if(static_threaded_workers_count < 1) static_threaded_workers_count = 1;
476
477 - size_t max_sockets = (size_t)config_get_number(CONFIG_SECTION_WEB, "web server max sockets", (long long int)(rlimit_nofile.rlim_cur / 2));
477 + size_t max_sockets = (size_t)config_get_number(CONFIG_SECTION_WEB, "web server max sockets", (long long int)(rlimit_nofile.rlim_cur / 4));
478
479 static_workers_private_data = callocz((size_t)static_threaded_workers_count, sizeof(struct web_server_static_threaded_worker));
480