rrdengine.c 70 KB

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