rrdengine.c 70 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866
  1. // SPDX-License-Identifier: GPL-3.0-or-later
  2. #define NETDATA_RRD_INTERNALS
  3. #include "rrdengine.h"
  4. #include "pdc.h"
  5. rrdeng_stats_t global_io_errors = 0;
  6. rrdeng_stats_t global_fs_errors = 0;
  7. rrdeng_stats_t rrdeng_reserved_file_descriptors = 0;
  8. rrdeng_stats_t global_pg_cache_over_half_dirty_events = 0;
  9. rrdeng_stats_t global_flushing_pressure_page_deletions = 0;
  10. unsigned rrdeng_pages_per_extent = MAX_PAGES_PER_EXTENT;
  11. #if WORKER_UTILIZATION_MAX_JOB_TYPES < (RRDENG_OPCODE_MAX + 2)
  12. #error Please increase WORKER_UTILIZATION_MAX_JOB_TYPES to at least (RRDENG_MAX_OPCODE + 2)
  13. #endif
  14. struct rrdeng_cmd {
  15. struct rrdengine_instance *ctx;
  16. enum rrdeng_opcode opcode;
  17. void *data;
  18. struct completion *completion;
  19. enum storage_priority priority;
  20. dequeue_callback_t dequeue_cb;
  21. struct {
  22. struct rrdeng_cmd *prev;
  23. struct rrdeng_cmd *next;
  24. } queue;
  25. };
  26. static inline struct rrdeng_cmd rrdeng_deq_cmd(bool from_worker);
  27. static inline void worker_dispatch_extent_read(struct rrdeng_cmd cmd, bool from_worker);
  28. static inline void worker_dispatch_query_prep(struct rrdeng_cmd cmd, bool from_worker);
  29. struct rrdeng_main {
  30. uv_thread_t thread;
  31. uv_loop_t loop;
  32. uv_async_t async;
  33. uv_timer_t timer;
  34. pid_t tid;
  35. bool shutdown;
  36. size_t flushes_running;
  37. size_t evictions_running;
  38. size_t cleanup_running;
  39. struct {
  40. ARAL *ar;
  41. struct {
  42. SPINLOCK spinlock;
  43. size_t waiting;
  44. struct rrdeng_cmd *waiting_items_by_priority[STORAGE_PRIORITY_INTERNAL_MAX_DONT_USE];
  45. size_t executed_by_priority[STORAGE_PRIORITY_INTERNAL_MAX_DONT_USE];
  46. } unsafe;
  47. } cmd_queue;
  48. struct {
  49. ARAL *ar;
  50. struct {
  51. size_t dispatched;
  52. size_t executing;
  53. } atomics;
  54. } work_cmd;
  55. struct {
  56. ARAL *ar;
  57. } handles;
  58. struct {
  59. ARAL *ar;
  60. } descriptors;
  61. struct {
  62. ARAL *ar;
  63. } xt_io_descr;
  64. } rrdeng_main = {
  65. .thread = 0,
  66. .loop = {},
  67. .async = {},
  68. .timer = {},
  69. .flushes_running = 0,
  70. .evictions_running = 0,
  71. .cleanup_running = 0,
  72. .cmd_queue = {
  73. .unsafe = {
  74. .spinlock = NETDATA_SPINLOCK_INITIALIZER,
  75. },
  76. }
  77. };
  78. static void sanity_check(void)
  79. {
  80. BUILD_BUG_ON(WORKER_UTILIZATION_MAX_JOB_TYPES < (RRDENG_OPCODE_MAX + 2));
  81. /* Magic numbers must fit in the super-blocks */
  82. BUILD_BUG_ON(strlen(RRDENG_DF_MAGIC) > RRDENG_MAGIC_SZ);
  83. BUILD_BUG_ON(strlen(RRDENG_JF_MAGIC) > RRDENG_MAGIC_SZ);
  84. /* Version strings must fit in the super-blocks */
  85. BUILD_BUG_ON(strlen(RRDENG_DF_VER) > RRDENG_VER_SZ);
  86. BUILD_BUG_ON(strlen(RRDENG_JF_VER) > RRDENG_VER_SZ);
  87. /* Data file super-block cannot be larger than RRDENG_BLOCK_SIZE */
  88. BUILD_BUG_ON(RRDENG_DF_SB_PADDING_SZ < 0);
  89. BUILD_BUG_ON(sizeof(uuid_t) != UUID_SZ); /* check UUID size */
  90. /* page count must fit in 8 bits */
  91. BUILD_BUG_ON(MAX_PAGES_PER_EXTENT > 255);
  92. /* extent cache count must fit in 32 bits */
  93. // BUILD_BUG_ON(MAX_CACHED_EXTENTS > 32);
  94. /* page info scratch space must be able to hold 2 32-bit integers */
  95. BUILD_BUG_ON(sizeof(((struct rrdeng_page_info *)0)->scratch) < 2 * sizeof(uint32_t));
  96. }
  97. // ----------------------------------------------------------------------------
  98. // work request cache
  99. typedef void *(*work_cb)(struct rrdengine_instance *ctx, void *data, struct completion *completion, uv_work_t* req);
  100. typedef void (*after_work_cb)(struct rrdengine_instance *ctx, void *data, struct completion *completion, uv_work_t* req, int status);
  101. struct rrdeng_work {
  102. uv_work_t req;
  103. struct rrdengine_instance *ctx;
  104. void *data;
  105. struct completion *completion;
  106. work_cb work_cb;
  107. after_work_cb after_work_cb;
  108. enum rrdeng_opcode opcode;
  109. };
  110. static void work_request_init(void) {
  111. rrdeng_main.work_cmd.ar = aral_create(
  112. "dbengine-work-cmd",
  113. sizeof(struct rrdeng_work),
  114. 0,
  115. 65536, NULL,
  116. NULL, NULL, false, false
  117. );
  118. }
  119. enum LIBUV_WORKERS_STATUS {
  120. LIBUV_WORKERS_RELAXED,
  121. LIBUV_WORKERS_STRESSED,
  122. LIBUV_WORKERS_CRITICAL,
  123. };
  124. static inline enum LIBUV_WORKERS_STATUS work_request_full(void) {
  125. size_t dispatched = __atomic_load_n(&rrdeng_main.work_cmd.atomics.dispatched, __ATOMIC_RELAXED);
  126. if(dispatched >= (size_t)(libuv_worker_threads))
  127. return LIBUV_WORKERS_CRITICAL;
  128. else if(dispatched >= (size_t)(libuv_worker_threads - RESERVED_LIBUV_WORKER_THREADS))
  129. return LIBUV_WORKERS_STRESSED;
  130. return LIBUV_WORKERS_RELAXED;
  131. }
  132. static inline void work_done(struct rrdeng_work *work_request) {
  133. aral_freez(rrdeng_main.work_cmd.ar, work_request);
  134. }
  135. static void work_standard_worker(uv_work_t *req) {
  136. __atomic_add_fetch(&rrdeng_main.work_cmd.atomics.executing, 1, __ATOMIC_RELAXED);
  137. register_libuv_worker_jobs();
  138. worker_is_busy(UV_EVENT_WORKER_INIT);
  139. struct rrdeng_work *work_request = req->data;
  140. work_request->data = work_request->work_cb(work_request->ctx, work_request->data, work_request->completion, req);
  141. worker_is_idle();
  142. if(work_request->opcode == RRDENG_OPCODE_EXTENT_READ || work_request->opcode == RRDENG_OPCODE_QUERY) {
  143. internal_fatal(work_request->after_work_cb != NULL, "DBENGINE: opcodes with a callback should not boosted");
  144. while(1) {
  145. struct rrdeng_cmd cmd = rrdeng_deq_cmd(true);
  146. if (cmd.opcode == RRDENG_OPCODE_NOOP)
  147. break;
  148. worker_is_busy(UV_EVENT_WORKER_INIT);
  149. switch (cmd.opcode) {
  150. case RRDENG_OPCODE_EXTENT_READ:
  151. worker_dispatch_extent_read(cmd, true);
  152. break;
  153. case RRDENG_OPCODE_QUERY:
  154. worker_dispatch_query_prep(cmd, true);
  155. break;
  156. default:
  157. fatal("DBENGINE: Opcode should not be executed synchronously");
  158. break;
  159. }
  160. worker_is_idle();
  161. }
  162. }
  163. __atomic_sub_fetch(&rrdeng_main.work_cmd.atomics.dispatched, 1, __ATOMIC_RELAXED);
  164. __atomic_sub_fetch(&rrdeng_main.work_cmd.atomics.executing, 1, __ATOMIC_RELAXED);
  165. // signal the event loop a worker is available
  166. fatal_assert(0 == uv_async_send(&rrdeng_main.async));
  167. }
  168. static void after_work_standard_callback(uv_work_t* req, int status) {
  169. struct rrdeng_work *work_request = req->data;
  170. worker_is_busy(RRDENG_OPCODE_MAX + work_request->opcode);
  171. if(work_request->after_work_cb)
  172. work_request->after_work_cb(work_request->ctx, work_request->data, work_request->completion, req, status);
  173. work_done(work_request);
  174. worker_is_idle();
  175. }
  176. static bool work_dispatch(struct rrdengine_instance *ctx, void *data, struct completion *completion, enum rrdeng_opcode opcode, work_cb work_cb, after_work_cb after_work_cb) {
  177. struct rrdeng_work *work_request = NULL;
  178. internal_fatal(rrdeng_main.tid != gettid(), "work_dispatch() can only be run from the event loop thread");
  179. work_request = aral_mallocz(rrdeng_main.work_cmd.ar);
  180. memset(work_request, 0, sizeof(struct rrdeng_work));
  181. work_request->req.data = work_request;
  182. work_request->ctx = ctx;
  183. work_request->data = data;
  184. work_request->completion = completion;
  185. work_request->work_cb = work_cb;
  186. work_request->after_work_cb = after_work_cb;
  187. work_request->opcode = opcode;
  188. if(uv_queue_work(&rrdeng_main.loop, &work_request->req, work_standard_worker, after_work_standard_callback)) {
  189. internal_fatal(true, "DBENGINE: cannot queue work");
  190. work_done(work_request);
  191. return false;
  192. }
  193. __atomic_add_fetch(&rrdeng_main.work_cmd.atomics.dispatched, 1, __ATOMIC_RELAXED);
  194. return true;
  195. }
  196. // ----------------------------------------------------------------------------
  197. // page descriptor cache
  198. void page_descriptors_init(void) {
  199. rrdeng_main.descriptors.ar = aral_create(
  200. "dbengine-descriptors",
  201. sizeof(struct page_descr_with_data),
  202. 0,
  203. 65536 * 4,
  204. NULL,
  205. NULL, NULL, false, false);
  206. }
  207. struct page_descr_with_data *page_descriptor_get(void) {
  208. struct page_descr_with_data *descr = aral_mallocz(rrdeng_main.descriptors.ar);
  209. memset(descr, 0, sizeof(struct page_descr_with_data));
  210. return descr;
  211. }
  212. static inline void page_descriptor_release(struct page_descr_with_data *descr) {
  213. aral_freez(rrdeng_main.descriptors.ar, descr);
  214. }
  215. // ----------------------------------------------------------------------------
  216. // extent io descriptor cache
  217. static void extent_io_descriptor_init(void) {
  218. rrdeng_main.xt_io_descr.ar = aral_create(
  219. "dbengine-extent-io",
  220. sizeof(struct extent_io_descriptor),
  221. 0,
  222. 65536,
  223. NULL,
  224. NULL, NULL, false, false
  225. );
  226. }
  227. static struct extent_io_descriptor *extent_io_descriptor_get(void) {
  228. struct extent_io_descriptor *xt_io_descr = aral_mallocz(rrdeng_main.xt_io_descr.ar);
  229. memset(xt_io_descr, 0, sizeof(struct extent_io_descriptor));
  230. return xt_io_descr;
  231. }
  232. static inline void extent_io_descriptor_release(struct extent_io_descriptor *xt_io_descr) {
  233. aral_freez(rrdeng_main.xt_io_descr.ar, xt_io_descr);
  234. }
  235. // ----------------------------------------------------------------------------
  236. // query handle cache
  237. void rrdeng_query_handle_init(void) {
  238. rrdeng_main.handles.ar = aral_create(
  239. "dbengine-query-handles",
  240. sizeof(struct rrdeng_query_handle),
  241. 0,
  242. 65536,
  243. NULL,
  244. NULL, NULL, false, false);
  245. }
  246. struct rrdeng_query_handle *rrdeng_query_handle_get(void) {
  247. struct rrdeng_query_handle *handle = aral_mallocz(rrdeng_main.handles.ar);
  248. memset(handle, 0, sizeof(struct rrdeng_query_handle));
  249. return handle;
  250. }
  251. void rrdeng_query_handle_release(struct rrdeng_query_handle *handle) {
  252. aral_freez(rrdeng_main.handles.ar, handle);
  253. }
  254. // ----------------------------------------------------------------------------
  255. // WAL cache
  256. static struct {
  257. struct {
  258. SPINLOCK spinlock;
  259. WAL *available_items;
  260. size_t available;
  261. } protected;
  262. struct {
  263. size_t allocated;
  264. } atomics;
  265. } wal_globals = {
  266. .protected = {
  267. .spinlock = NETDATA_SPINLOCK_INITIALIZER,
  268. .available_items = NULL,
  269. .available = 0,
  270. },
  271. .atomics = {
  272. .allocated = 0,
  273. },
  274. };
  275. static void wal_cleanup1(void) {
  276. WAL *wal = NULL;
  277. if(!spinlock_trylock(&wal_globals.protected.spinlock))
  278. return;
  279. if(wal_globals.protected.available_items && wal_globals.protected.available > storage_tiers) {
  280. wal = wal_globals.protected.available_items;
  281. DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(wal_globals.protected.available_items, wal, cache.prev, cache.next);
  282. wal_globals.protected.available--;
  283. }
  284. spinlock_unlock(&wal_globals.protected.spinlock);
  285. if(wal) {
  286. posix_memfree(wal->buf);
  287. freez(wal);
  288. __atomic_sub_fetch(&wal_globals.atomics.allocated, 1, __ATOMIC_RELAXED);
  289. }
  290. }
  291. WAL *wal_get(struct rrdengine_instance *ctx, unsigned size) {
  292. if(!size || size > RRDENG_BLOCK_SIZE)
  293. fatal("DBENGINE: invalid WAL size requested");
  294. WAL *wal = NULL;
  295. spinlock_lock(&wal_globals.protected.spinlock);
  296. if(likely(wal_globals.protected.available_items)) {
  297. wal = wal_globals.protected.available_items;
  298. DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(wal_globals.protected.available_items, wal, cache.prev, cache.next);
  299. wal_globals.protected.available--;
  300. }
  301. uint64_t transaction_id = __atomic_fetch_add(&ctx->atomic.transaction_id, 1, __ATOMIC_RELAXED);
  302. spinlock_unlock(&wal_globals.protected.spinlock);
  303. if(unlikely(!wal)) {
  304. wal = mallocz(sizeof(WAL));
  305. wal->buf_size = RRDENG_BLOCK_SIZE;
  306. int ret = posix_memalign((void *)&wal->buf, RRDFILE_ALIGNMENT, wal->buf_size);
  307. if (unlikely(ret))
  308. fatal("DBENGINE: posix_memalign:%s", strerror(ret));
  309. __atomic_add_fetch(&wal_globals.atomics.allocated, 1, __ATOMIC_RELAXED);
  310. }
  311. // these need to survive
  312. unsigned buf_size = wal->buf_size;
  313. void *buf = wal->buf;
  314. memset(wal, 0, sizeof(WAL));
  315. // put them back
  316. wal->buf_size = buf_size;
  317. wal->buf = buf;
  318. memset(wal->buf, 0, wal->buf_size);
  319. wal->transaction_id = transaction_id;
  320. wal->size = size;
  321. return wal;
  322. }
  323. void wal_release(WAL *wal) {
  324. if(unlikely(!wal)) return;
  325. spinlock_lock(&wal_globals.protected.spinlock);
  326. DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(wal_globals.protected.available_items, wal, cache.prev, cache.next);
  327. wal_globals.protected.available++;
  328. spinlock_unlock(&wal_globals.protected.spinlock);
  329. }
  330. // ----------------------------------------------------------------------------
  331. // command queue cache
  332. static void rrdeng_cmd_queue_init(void) {
  333. rrdeng_main.cmd_queue.ar = aral_create("dbengine-opcodes",
  334. sizeof(struct rrdeng_cmd),
  335. 0,
  336. 65536,
  337. NULL,
  338. NULL, NULL, false, false);
  339. }
  340. static inline STORAGE_PRIORITY rrdeng_enq_cmd_map_opcode_to_priority(enum rrdeng_opcode opcode, STORAGE_PRIORITY priority) {
  341. if(unlikely(priority >= STORAGE_PRIORITY_INTERNAL_MAX_DONT_USE))
  342. priority = STORAGE_PRIORITY_BEST_EFFORT;
  343. switch(opcode) {
  344. case RRDENG_OPCODE_QUERY:
  345. priority = STORAGE_PRIORITY_INTERNAL_QUERY_PREP;
  346. break;
  347. default:
  348. break;
  349. }
  350. return priority;
  351. }
  352. void rrdeng_enqueue_epdl_cmd(struct rrdeng_cmd *cmd) {
  353. epdl_cmd_queued(cmd->data, cmd);
  354. }
  355. void rrdeng_dequeue_epdl_cmd(struct rrdeng_cmd *cmd) {
  356. epdl_cmd_dequeued(cmd->data);
  357. }
  358. void rrdeng_req_cmd(requeue_callback_t get_cmd_cb, void *data, STORAGE_PRIORITY priority) {
  359. spinlock_lock(&rrdeng_main.cmd_queue.unsafe.spinlock);
  360. struct rrdeng_cmd *cmd = get_cmd_cb(data);
  361. if(cmd) {
  362. priority = rrdeng_enq_cmd_map_opcode_to_priority(cmd->opcode, priority);
  363. if (cmd->priority > priority) {
  364. DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(rrdeng_main.cmd_queue.unsafe.waiting_items_by_priority[cmd->priority], cmd, queue.prev, queue.next);
  365. DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(rrdeng_main.cmd_queue.unsafe.waiting_items_by_priority[priority], cmd, queue.prev, queue.next);
  366. cmd->priority = priority;
  367. }
  368. }
  369. spinlock_unlock(&rrdeng_main.cmd_queue.unsafe.spinlock);
  370. }
  371. void rrdeng_enq_cmd(struct rrdengine_instance *ctx, enum rrdeng_opcode opcode, void *data, struct completion *completion,
  372. enum storage_priority priority, enqueue_callback_t enqueue_cb, dequeue_callback_t dequeue_cb) {
  373. priority = rrdeng_enq_cmd_map_opcode_to_priority(opcode, priority);
  374. struct rrdeng_cmd *cmd = aral_mallocz(rrdeng_main.cmd_queue.ar);
  375. memset(cmd, 0, sizeof(struct rrdeng_cmd));
  376. cmd->ctx = ctx;
  377. cmd->opcode = opcode;
  378. cmd->data = data;
  379. cmd->completion = completion;
  380. cmd->priority = priority;
  381. cmd->dequeue_cb = dequeue_cb;
  382. spinlock_lock(&rrdeng_main.cmd_queue.unsafe.spinlock);
  383. DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(rrdeng_main.cmd_queue.unsafe.waiting_items_by_priority[priority], cmd, queue.prev, queue.next);
  384. rrdeng_main.cmd_queue.unsafe.waiting++;
  385. if(enqueue_cb)
  386. enqueue_cb(cmd);
  387. spinlock_unlock(&rrdeng_main.cmd_queue.unsafe.spinlock);
  388. fatal_assert(0 == uv_async_send(&rrdeng_main.async));
  389. }
  390. static inline bool rrdeng_cmd_has_waiting_opcodes_in_lower_priorities(STORAGE_PRIORITY priority, STORAGE_PRIORITY max_priority) {
  391. for(; priority <= max_priority ; priority++)
  392. if(rrdeng_main.cmd_queue.unsafe.waiting_items_by_priority[priority])
  393. return true;
  394. return false;
  395. }
  396. #define opcode_empty (struct rrdeng_cmd) { \
  397. .ctx = NULL, \
  398. .opcode = RRDENG_OPCODE_NOOP, \
  399. .priority = STORAGE_PRIORITY_BEST_EFFORT, \
  400. .completion = NULL, \
  401. .data = NULL, \
  402. }
  403. static inline struct rrdeng_cmd rrdeng_deq_cmd(bool from_worker) {
  404. struct rrdeng_cmd *cmd = NULL;
  405. enum LIBUV_WORKERS_STATUS status = work_request_full();
  406. STORAGE_PRIORITY min_priority, max_priority;
  407. min_priority = STORAGE_PRIORITY_INTERNAL_DBENGINE;
  408. max_priority = (status != LIBUV_WORKERS_RELAXED) ? STORAGE_PRIORITY_INTERNAL_DBENGINE : STORAGE_PRIORITY_INTERNAL_MAX_DONT_USE - 1;
  409. if(from_worker) {
  410. if(status == LIBUV_WORKERS_CRITICAL)
  411. return opcode_empty;
  412. min_priority = STORAGE_PRIORITY_INTERNAL_QUERY_PREP;
  413. max_priority = STORAGE_PRIORITY_BEST_EFFORT;
  414. }
  415. // find an opcode to execute from the queue
  416. spinlock_lock(&rrdeng_main.cmd_queue.unsafe.spinlock);
  417. for(STORAGE_PRIORITY priority = min_priority; priority <= max_priority ; priority++) {
  418. cmd = rrdeng_main.cmd_queue.unsafe.waiting_items_by_priority[priority];
  419. if(cmd) {
  420. // avoid starvation of lower priorities
  421. if(unlikely(priority >= STORAGE_PRIORITY_HIGH &&
  422. priority < STORAGE_PRIORITY_BEST_EFFORT &&
  423. ++rrdeng_main.cmd_queue.unsafe.executed_by_priority[priority] % 50 == 0 &&
  424. rrdeng_cmd_has_waiting_opcodes_in_lower_priorities(priority + 1, max_priority))) {
  425. // let the others run 2% of the requests
  426. cmd = NULL;
  427. continue;
  428. }
  429. // remove it from the queue
  430. DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(rrdeng_main.cmd_queue.unsafe.waiting_items_by_priority[priority], cmd, queue.prev, queue.next);
  431. rrdeng_main.cmd_queue.unsafe.waiting--;
  432. break;
  433. }
  434. }
  435. if(cmd && cmd->dequeue_cb) {
  436. cmd->dequeue_cb(cmd);
  437. cmd->dequeue_cb = NULL;
  438. }
  439. spinlock_unlock(&rrdeng_main.cmd_queue.unsafe.spinlock);
  440. struct rrdeng_cmd ret;
  441. if(cmd) {
  442. // copy it, to return it
  443. ret = *cmd;
  444. aral_freez(rrdeng_main.cmd_queue.ar, cmd);
  445. }
  446. else
  447. ret = opcode_empty;
  448. return ret;
  449. }
  450. // ----------------------------------------------------------------------------
  451. void *dbengine_extent_alloc(size_t size) {
  452. void *extent = mallocz(size);
  453. return extent;
  454. }
  455. void dbengine_extent_free(void *extent, size_t size __maybe_unused) {
  456. freez(extent);
  457. }
  458. static void journalfile_extent_build(struct rrdengine_instance *ctx, struct extent_io_descriptor *xt_io_descr) {
  459. unsigned count, payload_length, descr_size, size_bytes;
  460. void *buf;
  461. /* persistent structures */
  462. struct rrdeng_df_extent_header *df_header;
  463. struct rrdeng_jf_transaction_header *jf_header;
  464. struct rrdeng_jf_store_data *jf_metric_data;
  465. struct rrdeng_jf_transaction_trailer *jf_trailer;
  466. uLong crc;
  467. df_header = xt_io_descr->buf;
  468. count = df_header->number_of_pages;
  469. descr_size = sizeof(*jf_metric_data->descr) * count;
  470. payload_length = sizeof(*jf_metric_data) + descr_size;
  471. size_bytes = sizeof(*jf_header) + payload_length + sizeof(*jf_trailer);
  472. xt_io_descr->wal = wal_get(ctx, size_bytes);
  473. buf = xt_io_descr->wal->buf;
  474. jf_header = buf;
  475. jf_header->type = STORE_DATA;
  476. jf_header->reserved = 0;
  477. jf_header->id = xt_io_descr->wal->transaction_id;
  478. jf_header->payload_length = payload_length;
  479. jf_metric_data = buf + sizeof(*jf_header);
  480. jf_metric_data->extent_offset = xt_io_descr->pos;
  481. jf_metric_data->extent_size = xt_io_descr->bytes;
  482. jf_metric_data->number_of_pages = count;
  483. memcpy(jf_metric_data->descr, df_header->descr, descr_size);
  484. jf_trailer = buf + sizeof(*jf_header) + payload_length;
  485. crc = crc32(0L, Z_NULL, 0);
  486. crc = crc32(crc, buf, sizeof(*jf_header) + payload_length);
  487. crc32set(jf_trailer->checksum, crc);
  488. }
  489. static void after_extent_flushed_to_open(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t* req __maybe_unused, int status __maybe_unused) {
  490. if(completion)
  491. completion_mark_complete(completion);
  492. if(ctx_is_available_for_queries(ctx))
  493. rrdeng_enq_cmd(ctx, RRDENG_OPCODE_DATABASE_ROTATE, NULL, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
  494. }
  495. static void *extent_flushed_to_open_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) {
  496. worker_is_busy(UV_EVENT_DBENGINE_FLUSHED_TO_OPEN);
  497. uv_fs_t *uv_fs_request = data;
  498. struct extent_io_descriptor *xt_io_descr = uv_fs_request->data;
  499. struct page_descr_with_data *descr;
  500. struct rrdengine_datafile *datafile;
  501. unsigned i;
  502. datafile = xt_io_descr->datafile;
  503. bool still_running = ctx_is_available_for_queries(ctx);
  504. for (i = 0 ; i < xt_io_descr->descr_count ; ++i) {
  505. descr = xt_io_descr->descr_array[i];
  506. if (likely(still_running))
  507. pgc_open_add_hot_page(
  508. (Word_t)ctx, descr->metric_id,
  509. (time_t) (descr->start_time_ut / USEC_PER_SEC),
  510. (time_t) (descr->end_time_ut / USEC_PER_SEC),
  511. descr->update_every_s,
  512. datafile,
  513. xt_io_descr->pos, xt_io_descr->bytes, descr->page_length);
  514. page_descriptor_release(descr);
  515. }
  516. uv_fs_req_cleanup(uv_fs_request);
  517. posix_memfree(xt_io_descr->buf);
  518. extent_io_descriptor_release(xt_io_descr);
  519. spinlock_lock(&datafile->writers.spinlock);
  520. datafile->writers.flushed_to_open_running--;
  521. spinlock_unlock(&datafile->writers.spinlock);
  522. if(datafile->fileno != ctx_last_fileno_get(ctx) && still_running)
  523. // we just finished a flushing on a datafile that is not the active one
  524. rrdeng_enq_cmd(ctx, RRDENG_OPCODE_JOURNAL_INDEX, datafile, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
  525. return data;
  526. }
  527. // Main event loop callback
  528. static void after_extent_write_datafile_io(uv_fs_t *uv_fs_request) {
  529. worker_is_busy(RRDENG_OPCODE_MAX + RRDENG_OPCODE_EXTENT_WRITE);
  530. struct extent_io_descriptor *xt_io_descr = uv_fs_request->data;
  531. struct rrdengine_datafile *datafile = xt_io_descr->datafile;
  532. struct rrdengine_instance *ctx = datafile->ctx;
  533. if (uv_fs_request->result < 0) {
  534. ctx_io_error(ctx);
  535. netdata_log_error("DBENGINE: %s: uv_fs_write(): %s", __func__, uv_strerror((int)uv_fs_request->result));
  536. }
  537. journalfile_v1_extent_write(ctx, xt_io_descr->datafile, xt_io_descr->wal, &rrdeng_main.loop);
  538. spinlock_lock(&datafile->writers.spinlock);
  539. datafile->writers.running--;
  540. datafile->writers.flushed_to_open_running++;
  541. spinlock_unlock(&datafile->writers.spinlock);
  542. rrdeng_enq_cmd(xt_io_descr->ctx,
  543. RRDENG_OPCODE_FLUSHED_TO_OPEN,
  544. uv_fs_request,
  545. xt_io_descr->completion,
  546. STORAGE_PRIORITY_INTERNAL_DBENGINE,
  547. NULL,
  548. NULL);
  549. worker_is_idle();
  550. }
  551. static bool datafile_is_full(struct rrdengine_instance *ctx, struct rrdengine_datafile *datafile) {
  552. bool ret = false;
  553. spinlock_lock(&datafile->writers.spinlock);
  554. if(ctx_is_available_for_queries(ctx) && datafile->pos > rrdeng_target_data_file_size(ctx))
  555. ret = true;
  556. spinlock_unlock(&datafile->writers.spinlock);
  557. return ret;
  558. }
  559. static struct rrdengine_datafile *get_datafile_to_write_extent(struct rrdengine_instance *ctx) {
  560. struct rrdengine_datafile *datafile;
  561. // get the latest datafile
  562. uv_rwlock_rdlock(&ctx->datafiles.rwlock);
  563. datafile = ctx->datafiles.first->prev;
  564. // become a writer on this datafile, to prevent it from vanishing
  565. spinlock_lock(&datafile->writers.spinlock);
  566. datafile->writers.running++;
  567. spinlock_unlock(&datafile->writers.spinlock);
  568. uv_rwlock_rdunlock(&ctx->datafiles.rwlock);
  569. if(datafile_is_full(ctx, datafile)) {
  570. // remember the datafile we have become writers to
  571. struct rrdengine_datafile *old_datafile = datafile;
  572. // only 1 datafile creation at a time
  573. static netdata_mutex_t mutex = NETDATA_MUTEX_INITIALIZER;
  574. netdata_mutex_lock(&mutex);
  575. // take the latest datafile again - without this, multiple threads may create multiple files
  576. uv_rwlock_rdlock(&ctx->datafiles.rwlock);
  577. datafile = ctx->datafiles.first->prev;
  578. uv_rwlock_rdunlock(&ctx->datafiles.rwlock);
  579. if(datafile_is_full(ctx, datafile) && create_new_datafile_pair(ctx, true) == 0)
  580. rrdeng_enq_cmd(ctx, RRDENG_OPCODE_JOURNAL_INDEX, datafile, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL,
  581. NULL);
  582. netdata_mutex_unlock(&mutex);
  583. // get the new latest datafile again, like above
  584. uv_rwlock_rdlock(&ctx->datafiles.rwlock);
  585. datafile = ctx->datafiles.first->prev;
  586. // become a writer on this datafile, to prevent it from vanishing
  587. spinlock_lock(&datafile->writers.spinlock);
  588. datafile->writers.running++;
  589. spinlock_unlock(&datafile->writers.spinlock);
  590. uv_rwlock_rdunlock(&ctx->datafiles.rwlock);
  591. // release the writers on the old datafile
  592. spinlock_lock(&old_datafile->writers.spinlock);
  593. old_datafile->writers.running--;
  594. spinlock_unlock(&old_datafile->writers.spinlock);
  595. }
  596. return datafile;
  597. }
  598. /*
  599. * Take a page list in a judy array and write them
  600. */
  601. static struct extent_io_descriptor *datafile_extent_build(struct rrdengine_instance *ctx, struct page_descr_with_data *base, struct completion *completion) {
  602. int ret;
  603. int compressed_size, max_compressed_size = 0;
  604. unsigned i, count, size_bytes, pos, real_io_size;
  605. uint32_t uncompressed_payload_length, payload_offset;
  606. struct page_descr_with_data *descr, *eligible_pages[MAX_PAGES_PER_EXTENT];
  607. struct extent_io_descriptor *xt_io_descr;
  608. struct extent_buffer *eb = NULL;
  609. void *compressed_buf = NULL;
  610. Word_t Index;
  611. uint8_t compression_algorithm = ctx->config.global_compress_alg;
  612. struct rrdengine_datafile *datafile;
  613. /* persistent structures */
  614. struct rrdeng_df_extent_header *header;
  615. struct rrdeng_df_extent_trailer *trailer;
  616. uLong crc;
  617. for(descr = base, Index = 0, count = 0, uncompressed_payload_length = 0;
  618. descr && count != rrdeng_pages_per_extent;
  619. descr = descr->link.next, Index++) {
  620. uncompressed_payload_length += descr->page_length;
  621. eligible_pages[count++] = descr;
  622. }
  623. if (!count) {
  624. if (completion)
  625. completion_mark_complete(completion);
  626. __atomic_sub_fetch(&ctx->atomic.extents_currently_being_flushed, 1, __ATOMIC_RELAXED);
  627. return NULL;
  628. }
  629. xt_io_descr = extent_io_descriptor_get();
  630. xt_io_descr->ctx = ctx;
  631. payload_offset = sizeof(*header) + count * sizeof(header->descr[0]);
  632. switch (compression_algorithm) {
  633. case RRD_NO_COMPRESSION:
  634. size_bytes = payload_offset + uncompressed_payload_length + sizeof(*trailer);
  635. break;
  636. default: /* Compress */
  637. fatal_assert(uncompressed_payload_length < LZ4_MAX_INPUT_SIZE);
  638. max_compressed_size = LZ4_compressBound(uncompressed_payload_length);
  639. eb = extent_buffer_get(max_compressed_size);
  640. compressed_buf = eb->data;
  641. size_bytes = payload_offset + MAX(uncompressed_payload_length, (unsigned)max_compressed_size) + sizeof(*trailer);
  642. break;
  643. }
  644. ret = posix_memalign((void *)&xt_io_descr->buf, RRDFILE_ALIGNMENT, ALIGN_BYTES_CEILING(size_bytes));
  645. if (unlikely(ret)) {
  646. fatal("DBENGINE: posix_memalign:%s", strerror(ret));
  647. /* freez(xt_io_descr);*/
  648. }
  649. memset(xt_io_descr->buf, 0, ALIGN_BYTES_CEILING(size_bytes));
  650. (void) memcpy(xt_io_descr->descr_array, eligible_pages, sizeof(struct page_descr_with_data *) * count);
  651. xt_io_descr->descr_count = count;
  652. pos = 0;
  653. header = xt_io_descr->buf;
  654. header->compression_algorithm = compression_algorithm;
  655. header->number_of_pages = count;
  656. pos += sizeof(*header);
  657. for (i = 0 ; i < count ; ++i) {
  658. descr = xt_io_descr->descr_array[i];
  659. header->descr[i].type = descr->type;
  660. uuid_copy(*(uuid_t *)header->descr[i].uuid, *descr->id);
  661. header->descr[i].page_length = descr->page_length;
  662. header->descr[i].start_time_ut = descr->start_time_ut;
  663. switch (descr->type) {
  664. case PAGE_METRICS:
  665. case PAGE_TIER:
  666. header->descr[i].end_time_ut = descr->end_time_ut;
  667. break;
  668. case PAGE_GORILLA_METRICS:
  669. header->descr[i].gorilla.delta_time_s = (uint32_t) ((descr->end_time_ut - descr->start_time_ut) / USEC_PER_SEC);
  670. header->descr[i].gorilla.entries = pgd_slots_used(descr->pgd);
  671. break;
  672. default:
  673. fatal("Unknown page type: %uc", descr->type);
  674. }
  675. pos += sizeof(header->descr[i]);
  676. }
  677. for (i = 0 ; i < count ; ++i) {
  678. descr = xt_io_descr->descr_array[i];
  679. pgd_copy_to_extent(descr->pgd, xt_io_descr->buf + pos, descr->page_length);
  680. pos += descr->page_length;
  681. }
  682. if(likely(compression_algorithm == RRD_LZ4)) {
  683. compressed_size = LZ4_compress_default(
  684. xt_io_descr->buf + payload_offset,
  685. compressed_buf,
  686. (int)uncompressed_payload_length,
  687. max_compressed_size);
  688. __atomic_add_fetch(&ctx->stats.before_compress_bytes, uncompressed_payload_length, __ATOMIC_RELAXED);
  689. __atomic_add_fetch(&ctx->stats.after_compress_bytes, compressed_size, __ATOMIC_RELAXED);
  690. (void) memcpy(xt_io_descr->buf + payload_offset, compressed_buf, compressed_size);
  691. extent_buffer_release(eb);
  692. size_bytes = payload_offset + compressed_size + sizeof(*trailer);
  693. header->payload_length = compressed_size;
  694. }
  695. else { // RRD_NO_COMPRESSION
  696. header->payload_length = uncompressed_payload_length;
  697. }
  698. real_io_size = ALIGN_BYTES_CEILING(size_bytes);
  699. datafile = get_datafile_to_write_extent(ctx);
  700. spinlock_lock(&datafile->writers.spinlock);
  701. xt_io_descr->datafile = datafile;
  702. xt_io_descr->pos = datafile->pos;
  703. datafile->pos += real_io_size;
  704. spinlock_unlock(&datafile->writers.spinlock);
  705. xt_io_descr->bytes = size_bytes;
  706. xt_io_descr->uv_fs_request.data = xt_io_descr;
  707. xt_io_descr->completion = completion;
  708. trailer = xt_io_descr->buf + size_bytes - sizeof(*trailer);
  709. crc = crc32(0L, Z_NULL, 0);
  710. crc = crc32(crc, xt_io_descr->buf, size_bytes - sizeof(*trailer));
  711. crc32set(trailer->checksum, crc);
  712. xt_io_descr->iov = uv_buf_init((void *)xt_io_descr->buf, real_io_size);
  713. journalfile_extent_build(ctx, xt_io_descr);
  714. ctx_last_flush_fileno_set(ctx, datafile->fileno);
  715. ctx_current_disk_space_increase(ctx, real_io_size);
  716. ctx_io_write_op_bytes(ctx, real_io_size);
  717. return xt_io_descr;
  718. }
  719. static void after_extent_write(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t* uv_work_req __maybe_unused, int status __maybe_unused) {
  720. struct extent_io_descriptor *xt_io_descr = data;
  721. if(xt_io_descr) {
  722. int ret = uv_fs_write(&rrdeng_main.loop,
  723. &xt_io_descr->uv_fs_request,
  724. xt_io_descr->datafile->file,
  725. &xt_io_descr->iov,
  726. 1,
  727. (int64_t) xt_io_descr->pos,
  728. after_extent_write_datafile_io);
  729. fatal_assert(-1 != ret);
  730. }
  731. }
  732. static void *extent_write_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) {
  733. worker_is_busy(UV_EVENT_DBENGINE_EXTENT_WRITE);
  734. struct page_descr_with_data *base = data;
  735. struct extent_io_descriptor *xt_io_descr = datafile_extent_build(ctx, base, completion);
  736. return xt_io_descr;
  737. }
  738. static void after_database_rotate(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t* req __maybe_unused, int status __maybe_unused) {
  739. __atomic_store_n(&ctx->atomic.now_deleting_files, false, __ATOMIC_RELAXED);
  740. }
  741. struct uuid_first_time_s {
  742. uuid_t *uuid;
  743. time_t first_time_s;
  744. METRIC *metric;
  745. size_t pages_found;
  746. size_t df_matched;
  747. size_t df_index_oldest;
  748. };
  749. struct rrdengine_datafile *datafile_release_and_acquire_next_for_retention(struct rrdengine_instance *ctx, struct rrdengine_datafile *datafile) {
  750. uv_rwlock_rdlock(&ctx->datafiles.rwlock);
  751. struct rrdengine_datafile *next_datafile = datafile->next;
  752. while(next_datafile && !datafile_acquire(next_datafile, DATAFILE_ACQUIRE_RETENTION))
  753. next_datafile = next_datafile->next;
  754. uv_rwlock_rdunlock(&ctx->datafiles.rwlock);
  755. datafile_release(datafile, DATAFILE_ACQUIRE_RETENTION);
  756. return next_datafile;
  757. }
  758. time_t find_uuid_first_time(
  759. struct rrdengine_instance *ctx,
  760. struct rrdengine_datafile *datafile,
  761. struct uuid_first_time_s *uuid_first_entry_list,
  762. size_t count)
  763. {
  764. time_t global_first_time_s = LONG_MAX;
  765. // acquire the datafile to work with it
  766. uv_rwlock_rdlock(&ctx->datafiles.rwlock);
  767. while(datafile && !datafile_acquire(datafile, DATAFILE_ACQUIRE_RETENTION))
  768. datafile = datafile->next;
  769. uv_rwlock_rdunlock(&ctx->datafiles.rwlock);
  770. if (unlikely(!datafile))
  771. return global_first_time_s;
  772. unsigned journalfile_count = 0;
  773. size_t binary_match = 0;
  774. size_t not_matching_bsearches = 0;
  775. while (datafile) {
  776. struct journal_v2_header *j2_header = journalfile_v2_data_acquire(datafile->journalfile, NULL, 0, 0);
  777. if (!j2_header) {
  778. datafile = datafile_release_and_acquire_next_for_retention(ctx, datafile);
  779. continue;
  780. }
  781. time_t journal_start_time_s = (time_t) (j2_header->start_time_ut / USEC_PER_SEC);
  782. if(journal_start_time_s < global_first_time_s)
  783. global_first_time_s = journal_start_time_s;
  784. struct journal_metric_list *uuid_list = (struct journal_metric_list *)((uint8_t *) j2_header + j2_header->metric_offset);
  785. struct uuid_first_time_s *uuid_original_entry;
  786. size_t journal_metric_count = j2_header->metric_count;
  787. for (size_t index = 0; index < count; ++index) {
  788. uuid_original_entry = &uuid_first_entry_list[index];
  789. // Check here if we should skip this
  790. if (uuid_original_entry->df_matched > 3 || uuid_original_entry->pages_found > 5)
  791. continue;
  792. struct journal_metric_list *live_entry =
  793. bsearch(uuid_original_entry->uuid,uuid_list,journal_metric_count,
  794. sizeof(*uuid_list), journal_metric_uuid_compare);
  795. if (!live_entry) {
  796. // Not found in this journal
  797. not_matching_bsearches++;
  798. continue;
  799. }
  800. uuid_original_entry->pages_found += live_entry->entries;
  801. uuid_original_entry->df_matched++;
  802. time_t old_first_time_s = uuid_original_entry->first_time_s;
  803. // Calculate first / last for this match
  804. time_t first_time_s = live_entry->delta_start_s + journal_start_time_s;
  805. uuid_original_entry->first_time_s = MIN(uuid_original_entry->first_time_s, first_time_s);
  806. if (uuid_original_entry->first_time_s != old_first_time_s)
  807. uuid_original_entry->df_index_oldest = uuid_original_entry->df_matched;
  808. binary_match++;
  809. }
  810. journalfile_count++;
  811. journalfile_v2_data_release(datafile->journalfile);
  812. datafile = datafile_release_and_acquire_next_for_retention(ctx, datafile);
  813. }
  814. // Let's scan the open cache for almost exact match
  815. size_t open_cache_count = 0;
  816. size_t df_index[10] = { 0 };
  817. size_t without_metric = 0;
  818. size_t open_cache_gave_first_time_s = 0;
  819. size_t metric_count = 0;
  820. size_t without_retention = 0;
  821. size_t not_needed_bsearches = 0;
  822. for (size_t index = 0; index < count; ++index) {
  823. struct uuid_first_time_s *uuid_first_t_entry = &uuid_first_entry_list[index];
  824. metric_count++;
  825. size_t idx = uuid_first_t_entry->df_index_oldest;
  826. if(idx >= 10)
  827. idx = 9;
  828. df_index[idx]++;
  829. not_needed_bsearches += uuid_first_t_entry->df_matched - uuid_first_t_entry->df_index_oldest;
  830. if (unlikely(!uuid_first_t_entry->metric)) {
  831. without_metric++;
  832. continue;
  833. }
  834. PGC_PAGE *page = pgc_page_get_and_acquire(
  835. open_cache, (Word_t)ctx,
  836. (Word_t)uuid_first_t_entry->metric, 0,
  837. PGC_SEARCH_FIRST);
  838. if (page) {
  839. time_t old_first_time_s = uuid_first_t_entry->first_time_s;
  840. time_t first_time_s = pgc_page_start_time_s(page);
  841. uuid_first_t_entry->first_time_s = MIN(uuid_first_t_entry->first_time_s, first_time_s);
  842. pgc_page_release(open_cache, page);
  843. open_cache_count++;
  844. if(uuid_first_t_entry->first_time_s != old_first_time_s) {
  845. open_cache_gave_first_time_s++;
  846. }
  847. }
  848. else {
  849. if(!uuid_first_t_entry->df_index_oldest)
  850. without_retention++;
  851. }
  852. }
  853. internal_error(true,
  854. "DBENGINE: analyzed the retention of %zu rotated metrics of tier %d, "
  855. "did %zu jv2 matching binary searches (%zu not matching, %zu overflown) in %u journal files, "
  856. "%zu metrics with entries in open cache, "
  857. "metrics first time found per datafile index ([not in jv2]:%zu, [1]:%zu, [2]:%zu, [3]:%zu, [4]:%zu, [5]:%zu, [6]:%zu, [7]:%zu, [8]:%zu, [bigger]: %zu), "
  858. "open cache found first time %zu, "
  859. "metrics without any remaining retention %zu, "
  860. "metrics not in MRG %zu",
  861. metric_count,
  862. ctx->config.tier,
  863. binary_match,
  864. not_matching_bsearches,
  865. not_needed_bsearches,
  866. journalfile_count,
  867. open_cache_count,
  868. df_index[0], df_index[1], df_index[2], df_index[3], df_index[4], df_index[5], df_index[6], df_index[7], df_index[8], df_index[9],
  869. open_cache_gave_first_time_s,
  870. without_retention,
  871. without_metric
  872. );
  873. return global_first_time_s;
  874. }
  875. static void update_metrics_first_time_s(struct rrdengine_instance *ctx, struct rrdengine_datafile *datafile_to_delete, struct rrdengine_datafile *first_datafile_remaining, bool worker) {
  876. time_t global_first_time_s = LONG_MAX;
  877. if(worker)
  878. worker_is_busy(UV_EVENT_DBENGINE_FIND_ROTATED_METRICS);
  879. struct rrdengine_journalfile *journalfile = datafile_to_delete->journalfile;
  880. struct journal_v2_header *j2_header = journalfile_v2_data_acquire(journalfile, NULL, 0, 0);
  881. if (unlikely(!j2_header)) {
  882. if (worker)
  883. worker_is_idle();
  884. return;
  885. }
  886. __atomic_add_fetch(&rrdeng_cache_efficiency_stats.metrics_retention_started, 1, __ATOMIC_RELAXED);
  887. struct journal_metric_list *uuid_list = (struct journal_metric_list *)((uint8_t *) j2_header + j2_header->metric_offset);
  888. size_t count = j2_header->metric_count;
  889. struct uuid_first_time_s *uuid_first_t_entry;
  890. struct uuid_first_time_s *uuid_first_entry_list = callocz(count, sizeof(struct uuid_first_time_s));
  891. size_t added = 0;
  892. for (size_t index = 0; index < count; ++index) {
  893. METRIC *metric = mrg_metric_get_and_acquire(main_mrg, &uuid_list[index].uuid, (Word_t) ctx);
  894. if (!metric)
  895. continue;
  896. uuid_first_entry_list[added].metric = metric;
  897. uuid_first_entry_list[added].first_time_s = LONG_MAX;
  898. uuid_first_entry_list[added].df_matched = 0;
  899. uuid_first_entry_list[added].df_index_oldest = 0;
  900. uuid_first_entry_list[added].uuid = mrg_metric_uuid(main_mrg, metric);
  901. added++;
  902. }
  903. netdata_log_info("DBENGINE: recalculating tier %d retention for %zu metrics starting with datafile %u",
  904. ctx->config.tier, count, first_datafile_remaining->fileno);
  905. journalfile_v2_data_release(journalfile);
  906. // Update the first time / last time for all metrics we plan to delete
  907. if(worker)
  908. worker_is_busy(UV_EVENT_DBENGINE_FIND_REMAINING_RETENTION);
  909. global_first_time_s = find_uuid_first_time(ctx, first_datafile_remaining, uuid_first_entry_list, added);
  910. if(worker)
  911. worker_is_busy(UV_EVENT_DBENGINE_POPULATE_MRG);
  912. netdata_log_info("DBENGINE: updating tier %d metrics registry retention for %zu metrics",
  913. ctx->config.tier, added);
  914. size_t deleted_metrics = 0, zero_retention_referenced = 0, zero_disk_retention = 0, zero_disk_but_live = 0;
  915. for (size_t index = 0; index < added; ++index) {
  916. uuid_first_t_entry = &uuid_first_entry_list[index];
  917. if (likely(uuid_first_t_entry->first_time_s != LONG_MAX)) {
  918. mrg_metric_set_first_time_s_if_bigger(main_mrg, uuid_first_t_entry->metric, uuid_first_t_entry->first_time_s);
  919. mrg_metric_release(main_mrg, uuid_first_t_entry->metric);
  920. }
  921. else {
  922. zero_disk_retention++;
  923. // there is no retention for this metric
  924. bool has_retention = mrg_metric_zero_disk_retention(main_mrg, uuid_first_t_entry->metric);
  925. if (!has_retention) {
  926. bool deleted = mrg_metric_release_and_delete(main_mrg, uuid_first_t_entry->metric);
  927. if(deleted)
  928. deleted_metrics++;
  929. else
  930. zero_retention_referenced++;
  931. }
  932. else {
  933. zero_disk_but_live++;
  934. mrg_metric_release(main_mrg, uuid_first_t_entry->metric);
  935. }
  936. }
  937. }
  938. freez(uuid_first_entry_list);
  939. internal_error(zero_disk_retention,
  940. "DBENGINE: deleted %zu metrics, zero retention but referenced %zu (out of %zu total, of which %zu have main cache retention) zero on-disk retention tier %d metrics from metrics registry",
  941. deleted_metrics, zero_retention_referenced, zero_disk_retention, zero_disk_but_live, ctx->config.tier);
  942. if(global_first_time_s != LONG_MAX)
  943. __atomic_store_n(&ctx->atomic.first_time_s, global_first_time_s, __ATOMIC_RELAXED);
  944. if(worker)
  945. worker_is_idle();
  946. }
  947. void datafile_delete(struct rrdengine_instance *ctx, struct rrdengine_datafile *datafile, bool update_retention, bool worker) {
  948. if(worker)
  949. worker_is_busy(UV_EVENT_DBENGINE_DATAFILE_DELETE_WAIT);
  950. bool datafile_got_for_deletion = datafile_acquire_for_deletion(datafile);
  951. if (update_retention)
  952. update_metrics_first_time_s(ctx, datafile, datafile->next, worker);
  953. while (!datafile_got_for_deletion) {
  954. if(worker)
  955. worker_is_busy(UV_EVENT_DBENGINE_DATAFILE_DELETE_WAIT);
  956. datafile_got_for_deletion = datafile_acquire_for_deletion(datafile);
  957. if (!datafile_got_for_deletion) {
  958. netdata_log_info("DBENGINE: waiting for data file '%s/"
  959. DATAFILE_PREFIX RRDENG_FILE_NUMBER_PRINT_TMPL DATAFILE_EXTENSION
  960. "' to be available for deletion, "
  961. "it is in use currently by %u users.",
  962. ctx->config.dbfiles_path, ctx->datafiles.first->tier, ctx->datafiles.first->fileno, datafile->users.lockers);
  963. __atomic_add_fetch(&rrdeng_cache_efficiency_stats.datafile_deletion_spin, 1, __ATOMIC_RELAXED);
  964. sleep_usec(1 * USEC_PER_SEC);
  965. }
  966. }
  967. __atomic_add_fetch(&rrdeng_cache_efficiency_stats.datafile_deletion_started, 1, __ATOMIC_RELAXED);
  968. netdata_log_info("DBENGINE: deleting data file '%s/"
  969. DATAFILE_PREFIX RRDENG_FILE_NUMBER_PRINT_TMPL DATAFILE_EXTENSION
  970. "'.",
  971. ctx->config.dbfiles_path, ctx->datafiles.first->tier, ctx->datafiles.first->fileno);
  972. if(worker)
  973. worker_is_busy(UV_EVENT_DBENGINE_DATAFILE_DELETE);
  974. struct rrdengine_journalfile *journal_file;
  975. unsigned deleted_bytes, journal_file_bytes, datafile_bytes;
  976. int ret;
  977. char path[RRDENG_PATH_MAX];
  978. uv_rwlock_wrlock(&ctx->datafiles.rwlock);
  979. datafile_list_delete_unsafe(ctx, datafile);
  980. uv_rwlock_wrunlock(&ctx->datafiles.rwlock);
  981. journal_file = datafile->journalfile;
  982. datafile_bytes = datafile->pos;
  983. journal_file_bytes = journalfile_current_size(journal_file);
  984. deleted_bytes = journalfile_v2_data_size_get(journal_file);
  985. netdata_log_info("DBENGINE: deleting data and journal files to maintain disk quota");
  986. ret = journalfile_destroy_unsafe(journal_file, datafile);
  987. if (!ret) {
  988. journalfile_v1_generate_path(datafile, path, sizeof(path));
  989. netdata_log_info("DBENGINE: deleted journal file \"%s\".", path);
  990. journalfile_v2_generate_path(datafile, path, sizeof(path));
  991. netdata_log_info("DBENGINE: deleted journal file \"%s\".", path);
  992. deleted_bytes += journal_file_bytes;
  993. }
  994. ret = destroy_data_file_unsafe(datafile);
  995. if (!ret) {
  996. generate_datafilepath(datafile, path, sizeof(path));
  997. netdata_log_info("DBENGINE: deleted data file \"%s\".", path);
  998. deleted_bytes += datafile_bytes;
  999. }
  1000. freez(journal_file);
  1001. freez(datafile);
  1002. ctx_current_disk_space_decrease(ctx, deleted_bytes);
  1003. netdata_log_info("DBENGINE: reclaimed %u bytes of disk space.", deleted_bytes);
  1004. }
  1005. static void *database_rotate_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) {
  1006. datafile_delete(ctx, ctx->datafiles.first, ctx_is_available_for_queries(ctx), true);
  1007. if (rrdeng_ctx_exceeded_disk_quota(ctx))
  1008. rrdeng_enq_cmd(ctx, RRDENG_OPCODE_DATABASE_ROTATE, NULL, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
  1009. rrdcontext_db_rotation();
  1010. return data;
  1011. }
  1012. static void after_flush_all_hot_and_dirty_pages_of_section(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t* req __maybe_unused, int status __maybe_unused) {
  1013. ;
  1014. }
  1015. static void *flush_all_hot_and_dirty_pages_of_section_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) {
  1016. worker_is_busy(UV_EVENT_DBENGINE_QUIESCE);
  1017. pgc_flush_all_hot_and_dirty_pages(main_cache, (Word_t)ctx);
  1018. completion_mark_complete(&ctx->quiesce.completion);
  1019. return data;
  1020. }
  1021. static void after_populate_mrg(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t* req __maybe_unused, int status __maybe_unused) {
  1022. ;
  1023. }
  1024. static void *populate_mrg_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) {
  1025. worker_is_busy(UV_EVENT_DBENGINE_POPULATE_MRG);
  1026. do {
  1027. struct rrdengine_datafile *datafile = NULL;
  1028. // find a datafile to work
  1029. uv_rwlock_rdlock(&ctx->datafiles.rwlock);
  1030. for(datafile = ctx->datafiles.first; datafile ; datafile = datafile->next) {
  1031. if(!spinlock_trylock(&datafile->populate_mrg.spinlock))
  1032. continue;
  1033. if(datafile->populate_mrg.populated) {
  1034. spinlock_unlock(&datafile->populate_mrg.spinlock);
  1035. continue;
  1036. }
  1037. // we have the spinlock and it is not populated
  1038. break;
  1039. }
  1040. uv_rwlock_rdunlock(&ctx->datafiles.rwlock);
  1041. if(!datafile)
  1042. break;
  1043. journalfile_v2_populate_retention_to_mrg(ctx, datafile->journalfile);
  1044. datafile->populate_mrg.populated = true;
  1045. spinlock_unlock(&datafile->populate_mrg.spinlock);
  1046. } while(1);
  1047. completion_mark_complete(completion);
  1048. return data;
  1049. }
  1050. static void after_ctx_shutdown(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t* req __maybe_unused, int status __maybe_unused) {
  1051. ;
  1052. }
  1053. static void *ctx_shutdown_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) {
  1054. worker_is_busy(UV_EVENT_DBENGINE_SHUTDOWN);
  1055. bool logged = false;
  1056. while(__atomic_load_n(&ctx->atomic.extents_currently_being_flushed, __ATOMIC_RELAXED) ||
  1057. __atomic_load_n(&ctx->atomic.inflight_queries, __ATOMIC_RELAXED)) {
  1058. if(!logged) {
  1059. logged = true;
  1060. netdata_log_info("DBENGINE: waiting for %zu inflight queries to finish to shutdown tier %d...",
  1061. __atomic_load_n(&ctx->atomic.inflight_queries, __ATOMIC_RELAXED),
  1062. (ctx->config.legacy) ? -1 : ctx->config.tier);
  1063. }
  1064. sleep_usec(1 * USEC_PER_MS);
  1065. }
  1066. completion_mark_complete(completion);
  1067. return data;
  1068. }
  1069. static void *cache_flush_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) {
  1070. if (!main_cache)
  1071. return data;
  1072. worker_is_busy(UV_EVENT_DBENGINE_FLUSH_MAIN_CACHE);
  1073. pgc_flush_pages(main_cache, 0);
  1074. return data;
  1075. }
  1076. static void *cache_evict_tp_worker(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t *req __maybe_unused) {
  1077. if (!main_cache)
  1078. return data;
  1079. worker_is_busy(UV_EVENT_DBENGINE_EVICT_MAIN_CACHE);
  1080. pgc_evict_pages(main_cache, 0, 0);
  1081. return data;
  1082. }
  1083. static void *query_prep_tp_worker(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t *req __maybe_unused) {
  1084. PDC *pdc = data;
  1085. rrdeng_prep_query(pdc, true);
  1086. return data;
  1087. }
  1088. uint64_t rrdeng_target_data_file_size(struct rrdengine_instance *ctx) {
  1089. uint64_t target_size = ctx->config.max_disk_space / TARGET_DATAFILES;
  1090. target_size = MIN(target_size, MAX_DATAFILE_SIZE);
  1091. target_size = MAX(target_size, MIN_DATAFILE_SIZE);
  1092. return target_size;
  1093. }
  1094. bool rrdeng_ctx_exceeded_disk_quota(struct rrdengine_instance *ctx)
  1095. {
  1096. if(!ctx->datafiles.first)
  1097. // no datafiles available
  1098. return false;
  1099. if(!ctx->datafiles.first->next)
  1100. // only 1 datafile available
  1101. return false;
  1102. uint64_t estimated_disk_space = ctx_current_disk_space_get(ctx) + rrdeng_target_data_file_size(ctx) -
  1103. (ctx->datafiles.first->prev ? ctx->datafiles.first->prev->pos : 0);
  1104. return estimated_disk_space > ctx->config.max_disk_space;
  1105. }
  1106. /* return 0 on success */
  1107. int init_rrd_files(struct rrdengine_instance *ctx)
  1108. {
  1109. return init_data_files(ctx);
  1110. }
  1111. void finalize_rrd_files(struct rrdengine_instance *ctx)
  1112. {
  1113. return finalize_data_files(ctx);
  1114. }
  1115. void async_cb(uv_async_t *handle)
  1116. {
  1117. uv_stop(handle->loop);
  1118. uv_update_time(handle->loop);
  1119. netdata_log_debug(D_RRDENGINE, "%s called, active=%d.", __func__, uv_is_active((uv_handle_t *)handle));
  1120. }
  1121. #define TIMER_PERIOD_MS (1000)
  1122. static void *extent_read_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) {
  1123. EPDL *epdl = data;
  1124. epdl_find_extent_and_populate_pages(ctx, epdl, true);
  1125. return data;
  1126. }
  1127. static void epdl_populate_pages_asynchronously(struct rrdengine_instance *ctx, EPDL *epdl, STORAGE_PRIORITY priority) {
  1128. rrdeng_enq_cmd(ctx, RRDENG_OPCODE_EXTENT_READ, epdl, NULL, priority,
  1129. rrdeng_enqueue_epdl_cmd, rrdeng_dequeue_epdl_cmd);
  1130. }
  1131. void pdc_route_asynchronously(struct rrdengine_instance *ctx, struct page_details_control *pdc) {
  1132. pdc_to_epdl_router(ctx, pdc, epdl_populate_pages_asynchronously, epdl_populate_pages_asynchronously);
  1133. }
  1134. void epdl_populate_pages_synchronously(struct rrdengine_instance *ctx, EPDL *epdl, enum storage_priority priority __maybe_unused) {
  1135. epdl_find_extent_and_populate_pages(ctx, epdl, false);
  1136. }
  1137. void pdc_route_synchronously(struct rrdengine_instance *ctx, struct page_details_control *pdc) {
  1138. pdc_to_epdl_router(ctx, pdc, epdl_populate_pages_synchronously, epdl_populate_pages_synchronously);
  1139. }
  1140. #define MAX_RETRIES_TO_START_INDEX (100)
  1141. static void *journal_v2_indexing_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) {
  1142. unsigned count = 0;
  1143. worker_is_busy(UV_EVENT_DBENGINE_JOURNAL_INDEX_WAIT);
  1144. while (__atomic_load_n(&ctx->atomic.now_deleting_files, __ATOMIC_RELAXED) && count++ < MAX_RETRIES_TO_START_INDEX)
  1145. sleep_usec(100 * USEC_PER_MS);
  1146. if (count == MAX_RETRIES_TO_START_INDEX) {
  1147. worker_is_idle();
  1148. return data;
  1149. }
  1150. struct rrdengine_datafile *datafile = ctx->datafiles.first;
  1151. worker_is_busy(UV_EVENT_DBENGINE_JOURNAL_INDEX);
  1152. count = 0;
  1153. while (datafile && datafile->fileno != ctx_last_fileno_get(ctx) && datafile->fileno != ctx_last_flush_fileno_get(ctx)) {
  1154. if(journalfile_v2_data_available(datafile->journalfile)) {
  1155. // journal file v2 is already there for this datafile
  1156. datafile = datafile->next;
  1157. continue;
  1158. }
  1159. spinlock_lock(&datafile->writers.spinlock);
  1160. bool available = (datafile->writers.running || datafile->writers.flushed_to_open_running) ? false : true;
  1161. spinlock_unlock(&datafile->writers.spinlock);
  1162. if(!available) {
  1163. nd_log(NDLS_DAEMON, NDLP_NOTICE,
  1164. "DBENGINE: journal file %u needs to be indexed, but it has writers working on it - "
  1165. "skipping it for now",
  1166. datafile->fileno);
  1167. datafile = datafile->next;
  1168. continue;
  1169. }
  1170. nd_log(NDLS_DAEMON, NDLP_DEBUG,
  1171. "DBENGINE: journal file %u is ready to be indexed",
  1172. datafile->fileno);
  1173. pgc_open_cache_to_journal_v2(open_cache, (Word_t) ctx, (int) datafile->fileno, ctx->config.page_type,
  1174. journalfile_migrate_to_v2_callback, (void *) datafile->journalfile);
  1175. count++;
  1176. datafile = datafile->next;
  1177. if (unlikely(!ctx_is_available_for_queries(ctx)))
  1178. break;
  1179. }
  1180. errno = 0;
  1181. if(count)
  1182. nd_log(NDLS_DAEMON, NDLP_DEBUG,
  1183. "DBENGINE: journal indexing done; %u files processed",
  1184. count);
  1185. worker_is_idle();
  1186. return data;
  1187. }
  1188. static void after_do_cache_flush(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t* req __maybe_unused, int status __maybe_unused) {
  1189. rrdeng_main.flushes_running--;
  1190. }
  1191. static void after_do_cache_evict(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t* req __maybe_unused, int status __maybe_unused) {
  1192. rrdeng_main.evictions_running--;
  1193. }
  1194. static void after_journal_v2_indexing(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t* req __maybe_unused, int status __maybe_unused) {
  1195. __atomic_store_n(&ctx->atomic.migration_to_v2_running, false, __ATOMIC_RELAXED);
  1196. rrdeng_enq_cmd(ctx, RRDENG_OPCODE_DATABASE_ROTATE, NULL, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
  1197. }
  1198. struct rrdeng_buffer_sizes rrdeng_get_buffer_sizes(void) {
  1199. return (struct rrdeng_buffer_sizes) {
  1200. .pgc = pgc_aral_overhead() + pgc_aral_structures(),
  1201. .mrg = mrg_aral_overhead() + mrg_aral_structures(),
  1202. .opcodes = aral_overhead(rrdeng_main.cmd_queue.ar) + aral_structures(rrdeng_main.cmd_queue.ar),
  1203. .handles = aral_overhead(rrdeng_main.handles.ar) + aral_structures(rrdeng_main.handles.ar),
  1204. .descriptors = aral_overhead(rrdeng_main.descriptors.ar) + aral_structures(rrdeng_main.descriptors.ar),
  1205. .wal = __atomic_load_n(&wal_globals.atomics.allocated, __ATOMIC_RELAXED) * (sizeof(WAL) + RRDENG_BLOCK_SIZE),
  1206. .workers = aral_overhead(rrdeng_main.work_cmd.ar),
  1207. .pdc = pdc_cache_size(),
  1208. .xt_io = aral_overhead(rrdeng_main.xt_io_descr.ar) + aral_structures(rrdeng_main.xt_io_descr.ar),
  1209. .xt_buf = extent_buffer_cache_size(),
  1210. .epdl = epdl_cache_size(),
  1211. .deol = deol_cache_size(),
  1212. .pd = pd_cache_size(),
  1213. #ifdef PDC_USE_JULYL
  1214. .julyl = julyl_cache_size(),
  1215. #endif
  1216. };
  1217. }
  1218. static void after_cleanup(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t* req __maybe_unused, int status __maybe_unused) {
  1219. rrdeng_main.cleanup_running--;
  1220. }
  1221. 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) {
  1222. worker_is_busy(UV_EVENT_DBENGINE_BUFFERS_CLEANUP);
  1223. wal_cleanup1();
  1224. extent_buffer_cleanup1();
  1225. {
  1226. static time_t last_run_s = 0;
  1227. time_t now_s = now_monotonic_sec();
  1228. if(now_s - last_run_s >= 10) {
  1229. last_run_s = now_s;
  1230. journalfile_v2_data_unmount_cleanup(now_s);
  1231. }
  1232. }
  1233. #ifdef PDC_USE_JULYL
  1234. julyl_cleanup1();
  1235. #endif
  1236. return data;
  1237. }
  1238. void timer_cb(uv_timer_t* handle) {
  1239. worker_is_busy(RRDENG_TIMER_CB);
  1240. uv_stop(handle->loop);
  1241. uv_update_time(handle->loop);
  1242. worker_set_metric(RRDENG_OPCODES_WAITING, (NETDATA_DOUBLE)rrdeng_main.cmd_queue.unsafe.waiting);
  1243. worker_set_metric(RRDENG_WORKS_DISPATCHED, (NETDATA_DOUBLE)__atomic_load_n(&rrdeng_main.work_cmd.atomics.dispatched, __ATOMIC_RELAXED));
  1244. worker_set_metric(RRDENG_WORKS_EXECUTING, (NETDATA_DOUBLE)__atomic_load_n(&rrdeng_main.work_cmd.atomics.executing, __ATOMIC_RELAXED));
  1245. rrdeng_enq_cmd(NULL, RRDENG_OPCODE_FLUSH_INIT, NULL, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
  1246. rrdeng_enq_cmd(NULL, RRDENG_OPCODE_EVICT_INIT, NULL, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
  1247. rrdeng_enq_cmd(NULL, RRDENG_OPCODE_CLEANUP, NULL, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
  1248. worker_is_idle();
  1249. }
  1250. static void dbengine_initialize_structures(void) {
  1251. pgc_and_mrg_initialize();
  1252. pdc_init();
  1253. page_details_init();
  1254. epdl_init();
  1255. deol_init();
  1256. rrdeng_cmd_queue_init();
  1257. work_request_init();
  1258. rrdeng_query_handle_init();
  1259. page_descriptors_init();
  1260. extent_buffer_init();
  1261. pgd_init_arals();
  1262. extent_io_descriptor_init();
  1263. }
  1264. bool rrdeng_dbengine_spawn(struct rrdengine_instance *ctx __maybe_unused) {
  1265. static bool spawned = false;
  1266. static SPINLOCK spinlock = NETDATA_SPINLOCK_INITIALIZER;
  1267. spinlock_lock(&spinlock);
  1268. if(!spawned) {
  1269. int ret;
  1270. ret = uv_loop_init(&rrdeng_main.loop);
  1271. if (ret) {
  1272. netdata_log_error("DBENGINE: uv_loop_init(): %s", uv_strerror(ret));
  1273. return false;
  1274. }
  1275. rrdeng_main.loop.data = &rrdeng_main;
  1276. ret = uv_async_init(&rrdeng_main.loop, &rrdeng_main.async, async_cb);
  1277. if (ret) {
  1278. netdata_log_error("DBENGINE: uv_async_init(): %s", uv_strerror(ret));
  1279. fatal_assert(0 == uv_loop_close(&rrdeng_main.loop));
  1280. return false;
  1281. }
  1282. rrdeng_main.async.data = &rrdeng_main;
  1283. ret = uv_timer_init(&rrdeng_main.loop, &rrdeng_main.timer);
  1284. if (ret) {
  1285. netdata_log_error("DBENGINE: uv_timer_init(): %s", uv_strerror(ret));
  1286. uv_close((uv_handle_t *)&rrdeng_main.async, NULL);
  1287. fatal_assert(0 == uv_loop_close(&rrdeng_main.loop));
  1288. return false;
  1289. }
  1290. rrdeng_main.timer.data = &rrdeng_main;
  1291. dbengine_initialize_structures();
  1292. fatal_assert(0 == uv_thread_create(&rrdeng_main.thread, dbengine_event_loop, &rrdeng_main));
  1293. spawned = true;
  1294. }
  1295. spinlock_unlock(&spinlock);
  1296. return true;
  1297. }
  1298. static inline void worker_dispatch_extent_read(struct rrdeng_cmd cmd, bool from_worker) {
  1299. struct rrdengine_instance *ctx = cmd.ctx;
  1300. EPDL *epdl = cmd.data;
  1301. if(from_worker)
  1302. epdl_find_extent_and_populate_pages(ctx, epdl, true);
  1303. else
  1304. work_dispatch(ctx, epdl, NULL, cmd.opcode, extent_read_tp_worker, NULL);
  1305. }
  1306. static inline void worker_dispatch_query_prep(struct rrdeng_cmd cmd, bool from_worker) {
  1307. struct rrdengine_instance *ctx = cmd.ctx;
  1308. PDC *pdc = cmd.data;
  1309. if(from_worker)
  1310. rrdeng_prep_query(pdc, true);
  1311. else
  1312. work_dispatch(ctx, pdc, NULL, cmd.opcode, query_prep_tp_worker, NULL);
  1313. }
  1314. void dbengine_event_loop(void* arg) {
  1315. sanity_check();
  1316. uv_thread_set_name_np(pthread_self(), "DBENGINE");
  1317. service_register(SERVICE_THREAD_TYPE_EVENT_LOOP, NULL, NULL, NULL, true);
  1318. worker_register("DBENGINE");
  1319. // opcode jobs
  1320. worker_register_job_name(RRDENG_OPCODE_NOOP, "noop");
  1321. worker_register_job_name(RRDENG_OPCODE_QUERY, "query");
  1322. worker_register_job_name(RRDENG_OPCODE_EXTENT_WRITE, "extent write");
  1323. worker_register_job_name(RRDENG_OPCODE_EXTENT_READ, "extent read");
  1324. worker_register_job_name(RRDENG_OPCODE_FLUSHED_TO_OPEN, "flushed to open");
  1325. worker_register_job_name(RRDENG_OPCODE_DATABASE_ROTATE, "db rotate");
  1326. worker_register_job_name(RRDENG_OPCODE_JOURNAL_INDEX, "journal index");
  1327. worker_register_job_name(RRDENG_OPCODE_FLUSH_INIT, "flush init");
  1328. worker_register_job_name(RRDENG_OPCODE_EVICT_INIT, "evict init");
  1329. worker_register_job_name(RRDENG_OPCODE_CTX_SHUTDOWN, "ctx shutdown");
  1330. worker_register_job_name(RRDENG_OPCODE_CTX_QUIESCE, "ctx quiesce");
  1331. worker_register_job_name(RRDENG_OPCODE_SHUTDOWN_EVLOOP, "dbengine shutdown");
  1332. worker_register_job_name(RRDENG_OPCODE_MAX, "get opcode");
  1333. worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_QUERY, "query cb");
  1334. worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_EXTENT_WRITE, "extent write cb");
  1335. worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_EXTENT_READ, "extent read cb");
  1336. worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_FLUSHED_TO_OPEN, "flushed to open cb");
  1337. worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_DATABASE_ROTATE, "db rotate cb");
  1338. worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_JOURNAL_INDEX, "journal index cb");
  1339. worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_FLUSH_INIT, "flush init cb");
  1340. worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_EVICT_INIT, "evict init cb");
  1341. worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_CTX_SHUTDOWN, "ctx shutdown cb");
  1342. worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_CTX_QUIESCE, "ctx quiesce cb");
  1343. // special jobs
  1344. worker_register_job_name(RRDENG_TIMER_CB, "timer");
  1345. worker_register_job_name(RRDENG_FLUSH_TRANSACTION_BUFFER_CB, "transaction buffer flush cb");
  1346. worker_register_job_custom_metric(RRDENG_OPCODES_WAITING, "opcodes waiting", "opcodes", WORKER_METRIC_ABSOLUTE);
  1347. worker_register_job_custom_metric(RRDENG_WORKS_DISPATCHED, "works dispatched", "works", WORKER_METRIC_ABSOLUTE);
  1348. worker_register_job_custom_metric(RRDENG_WORKS_EXECUTING, "works executing", "works", WORKER_METRIC_ABSOLUTE);
  1349. struct rrdeng_main *main = arg;
  1350. enum rrdeng_opcode opcode;
  1351. struct rrdeng_cmd cmd;
  1352. main->tid = gettid();
  1353. fatal_assert(0 == uv_timer_start(&main->timer, timer_cb, TIMER_PERIOD_MS, TIMER_PERIOD_MS));
  1354. bool shutdown = false;
  1355. while (likely(!shutdown)) {
  1356. worker_is_idle();
  1357. uv_run(&main->loop, UV_RUN_DEFAULT);
  1358. /* wait for commands */
  1359. do {
  1360. worker_is_busy(RRDENG_OPCODE_MAX);
  1361. cmd = rrdeng_deq_cmd(RRDENG_OPCODE_NOOP);
  1362. opcode = cmd.opcode;
  1363. worker_is_busy(opcode);
  1364. switch (opcode) {
  1365. case RRDENG_OPCODE_EXTENT_READ:
  1366. worker_dispatch_extent_read(cmd, false);
  1367. break;
  1368. case RRDENG_OPCODE_QUERY:
  1369. worker_dispatch_query_prep(cmd, false);
  1370. break;
  1371. case RRDENG_OPCODE_EXTENT_WRITE: {
  1372. struct rrdengine_instance *ctx = cmd.ctx;
  1373. struct page_descr_with_data *base = cmd.data;
  1374. struct completion *completion = cmd.completion; // optional
  1375. work_dispatch(ctx, base, completion, opcode, extent_write_tp_worker, after_extent_write);
  1376. break;
  1377. }
  1378. case RRDENG_OPCODE_FLUSHED_TO_OPEN: {
  1379. struct rrdengine_instance *ctx = cmd.ctx;
  1380. uv_fs_t *uv_fs_request = cmd.data;
  1381. struct extent_io_descriptor *xt_io_descr = uv_fs_request->data;
  1382. struct completion *completion = xt_io_descr->completion;
  1383. work_dispatch(ctx, uv_fs_request, completion, opcode, extent_flushed_to_open_tp_worker, after_extent_flushed_to_open);
  1384. break;
  1385. }
  1386. case RRDENG_OPCODE_FLUSH_INIT: {
  1387. if(rrdeng_main.flushes_running < (size_t)(libuv_worker_threads / 4)) {
  1388. rrdeng_main.flushes_running++;
  1389. work_dispatch(NULL, NULL, NULL, opcode, cache_flush_tp_worker, after_do_cache_flush);
  1390. }
  1391. break;
  1392. }
  1393. case RRDENG_OPCODE_EVICT_INIT: {
  1394. if(!rrdeng_main.evictions_running) {
  1395. rrdeng_main.evictions_running++;
  1396. work_dispatch(NULL, NULL, NULL, opcode, cache_evict_tp_worker, after_do_cache_evict);
  1397. }
  1398. break;
  1399. }
  1400. case RRDENG_OPCODE_CLEANUP: {
  1401. if(!rrdeng_main.cleanup_running) {
  1402. rrdeng_main.cleanup_running++;
  1403. work_dispatch(NULL, NULL, NULL, opcode, cleanup_tp_worker, after_cleanup);
  1404. }
  1405. break;
  1406. }
  1407. case RRDENG_OPCODE_JOURNAL_INDEX: {
  1408. struct rrdengine_instance *ctx = cmd.ctx;
  1409. struct rrdengine_datafile *datafile = cmd.data;
  1410. if(!__atomic_load_n(&ctx->atomic.migration_to_v2_running, __ATOMIC_RELAXED)) {
  1411. __atomic_store_n(&ctx->atomic.migration_to_v2_running, true, __ATOMIC_RELAXED);
  1412. work_dispatch(ctx, datafile, NULL, opcode, journal_v2_indexing_tp_worker, after_journal_v2_indexing);
  1413. }
  1414. break;
  1415. }
  1416. case RRDENG_OPCODE_DATABASE_ROTATE: {
  1417. struct rrdengine_instance *ctx = cmd.ctx;
  1418. if (!__atomic_load_n(&ctx->atomic.now_deleting_files, __ATOMIC_RELAXED) &&
  1419. ctx->datafiles.first->next != NULL &&
  1420. ctx->datafiles.first->next->next != NULL &&
  1421. rrdeng_ctx_exceeded_disk_quota(ctx)) {
  1422. __atomic_store_n(&ctx->atomic.now_deleting_files, true, __ATOMIC_RELAXED);
  1423. work_dispatch(ctx, NULL, NULL, opcode, database_rotate_tp_worker, after_database_rotate);
  1424. }
  1425. break;
  1426. }
  1427. case RRDENG_OPCODE_CTX_POPULATE_MRG: {
  1428. struct rrdengine_instance *ctx = cmd.ctx;
  1429. struct completion *completion = cmd.completion;
  1430. work_dispatch(ctx, NULL, completion, opcode, populate_mrg_tp_worker, after_populate_mrg);
  1431. break;
  1432. }
  1433. case RRDENG_OPCODE_CTX_QUIESCE: {
  1434. // a ctx will shutdown shortly
  1435. struct rrdengine_instance *ctx = cmd.ctx;
  1436. __atomic_store_n(&ctx->quiesce.enabled, true, __ATOMIC_RELEASE);
  1437. work_dispatch(ctx, NULL, NULL, opcode,
  1438. flush_all_hot_and_dirty_pages_of_section_tp_worker,
  1439. after_flush_all_hot_and_dirty_pages_of_section);
  1440. break;
  1441. }
  1442. case RRDENG_OPCODE_CTX_SHUTDOWN: {
  1443. // a ctx is shutting down
  1444. struct rrdengine_instance *ctx = cmd.ctx;
  1445. struct completion *completion = cmd.completion;
  1446. work_dispatch(ctx, NULL, completion, opcode, ctx_shutdown_tp_worker, after_ctx_shutdown);
  1447. break;
  1448. }
  1449. case RRDENG_OPCODE_SHUTDOWN_EVLOOP: {
  1450. uv_close((uv_handle_t *)&main->async, NULL);
  1451. (void) uv_timer_stop(&main->timer);
  1452. uv_close((uv_handle_t *)&main->timer, NULL);
  1453. shutdown = true;
  1454. }
  1455. case RRDENG_OPCODE_NOOP: {
  1456. /* the command queue was empty, do nothing */
  1457. break;
  1458. }
  1459. // not opcodes
  1460. case RRDENG_OPCODE_MAX:
  1461. default: {
  1462. internal_fatal(true, "DBENGINE: unknown opcode");
  1463. break;
  1464. }
  1465. }
  1466. } while (opcode != RRDENG_OPCODE_NOOP);
  1467. }
  1468. nd_log(NDLS_DAEMON, NDLP_DEBUG, "Shutting down dbengine thread");
  1469. uv_loop_close(&main->loop);
  1470. worker_unregister();
  1471. }