76
if(unlikely(!service_running(SERVICE_CONTEXT)))
77
break;
78
79
- const DICTIONARY_ITEM *item = dictionary_get_and_acquire_item(host->rrdctx.contexts, string2str(rc->id));
79
+ STRING *lookup_id = string_dup(rc->id);
80
+ usec_t queued_ut = rc->queue.queued_ut;
81
+ spinlock_unlock(&host->rrdctx.hub_queue.spinlock);
82
+
83
+ const DICTIONARY_ITEM *item = dictionary_get_and_acquire_item(host->rrdctx.contexts, string2str(lookup_id));
84
bool context_matches = item && (dictionary_acquired_item_value(item) == rc);
85
bool stale = !item || !context_matches;
86
bool do_it = context_matches;
87
84
- spinlock_unlock(&host->rrdctx.hub_queue.spinlock);
85
-
88
if(context_matches) {
89
bool drop = force_drop_all;
90
if(!drop && apply_bounds) {
91
bool queue_is_over_limit = queued > RRDCONTEXT_HUB_QUEUE_MAX_OFFLINE_ENTRIES;
90
- bool queue_item_expired = (now_ut > rc->queue.queued_ut) &&
91
- ((now_ut - rc->queue.queued_ut) > RRDCONTEXT_HUB_QUEUE_MAX_OFFLINE_AGE_UT);
92
+ bool queue_item_expired = (now_ut > queued_ut) &&
93
+ ((now_ut - queued_ut) > RRDCONTEXT_HUB_QUEUE_MAX_OFFLINE_AGE_UT);
94
drop = queue_is_over_limit || queue_item_expired;
95
}
96
99
100
if(item)
101
dictionary_acquired_item_release(host->rrdctx.contexts, item);
102
+ string_freez(lookup_id);
103
104
spinlock_lock(&host->rrdctx.hub_queue.spinlock);
105
if(unlikely(!service_running(SERVICE_CONTEXT)))
122
}
123
}
124
else if(do_it) {
122
- rrdcontext_del_from_hub_queue(rc, true);
123
- dropped++;
124
- if(queued > 0)
125
- queued--;
125
+ RRDCONTEXT *rc_at_idx = RRDCONTEXT_QUEUE_GET(&host->rrdctx.hub_queue, idx);
126
+ if(rc_at_idx == rc) {
127
+ rrdcontext_del_from_hub_queue(rc, true);
128
+ dropped++;
129
+ if(queued > 0)
130
+ queued--;
131
+ }
132
}
133
}
134
spinlock_unlock(&host->rrdctx.hub_queue.spinlock);
232
rc = RRDCONTEXT_QUEUE_NEXT(&host->rrdctx.pp_queue, &idx)) {
233
if(unlikely(!service_running(SERVICE_CONTEXT))) break;
234
229
- const DICTIONARY_ITEM *item = dictionary_get_and_acquire_item(host->rrdctx.contexts, string2str(rc->id));
230
- bool do_it = dictionary_acquired_item_value(item) == rc;
235
+ STRING *lookup_id = string_dup(rc->id);
236
+ spinlock_unlock(&host->rrdctx.pp_queue.spinlock);
237
+
238
+ const DICTIONARY_ITEM *item = dictionary_get_and_acquire_item(host->rrdctx.contexts, string2str(lookup_id));
239
+ bool do_it = item && (dictionary_acquired_item_value(item) == rc);
240
+ bool process_it = false;
241
232
- if(do_it)
233
- rrdcontext_del_from_pp_queue(rc, true);
242
+ spinlock_lock(&host->rrdctx.pp_queue.spinlock);
243
+ if(unlikely(!service_running(SERVICE_CONTEXT))) {
244
+ spinlock_unlock(&host->rrdctx.pp_queue.spinlock);
245
+ if(item)
246
+ dictionary_acquired_item_release(host->rrdctx.contexts, item);
247
+ string_freez(lookup_id);
248
+ return;
249
+ }
250
+
251
+ if(do_it) {
252
+ RRDCONTEXT *rc_at_idx = RRDCONTEXT_QUEUE_GET(&host->rrdctx.pp_queue, idx);
253
+ if(rc_at_idx == rc) {
254
+ rrdcontext_del_from_pp_queue(rc, true);
255
+ process_it = true;
256
+ }
257
+ else
258
+ do_it = false;
259
+ }
260
261
spinlock_unlock(&host->rrdctx.pp_queue.spinlock);
262
263
if(item) {
238
- if (do_it)
264
+ if (process_it)
265
rrdcontext_post_process_updates(rc, false, RRD_FLAG_NONE, true);
240
-
266
dictionary_acquired_item_release(host->rrdctx.contexts, item);
267
}
268
+ string_freez(lookup_id);
269
270
spinlock_lock(&host->rrdctx.pp_queue.spinlock);
271
}
307
if(unlikely(messages_added >= MESSAGES_PER_BUNDLE_TO_SEND_TO_HUB_PER_HOST))
308
break;
309
284
- const DICTIONARY_ITEM *item = dictionary_get_and_acquire_item(host->rrdctx.contexts, string2str(rc->id));
285
- bool do_it = dictionary_acquired_item_value(item) == rc;
286
-
310
+ STRING *lookup_id = string_dup(rc->id);
311
spinlock_unlock(&host->rrdctx.hub_queue.spinlock);
312
313
+ const DICTIONARY_ITEM *item = dictionary_get_and_acquire_item(host->rrdctx.contexts, string2str(lookup_id));
314
+ bool do_it = item && (dictionary_acquired_item_value(item) == rc);
315
+ bool dispatch_ready = false;
316
+
317
if(item) {
318
if (do_it) {
291
- worker_is_busy(WORKER_JOB_QUEUED);
292
- usec_t dispatch_ut = rrdcontext_calculate_queued_dispatch_time_ut(rc, now_ut);
293
- if(unlikely(now_ut >= dispatch_ut) && claim_id_is_set(claim_id)) {
319
+ spinlock_lock(&host->rrdctx.hub_queue.spinlock);
320
+
321
+ if(likely(service_running(SERVICE_CONTEXT))) {
322
+ RRDCONTEXT *rc_at_idx = RRDCONTEXT_QUEUE_GET(&host->rrdctx.hub_queue, idx);
323
+ if(rc_at_idx == rc) {
324
+ worker_is_busy(WORKER_JOB_QUEUED);
325
+ usec_t dispatch_ut = rrdcontext_calculate_queued_dispatch_time_ut(rc, now_ut);
326
+ dispatch_ready = unlikely(now_ut >= dispatch_ut) && claim_id_is_set(claim_id);
327
+ }
328
+ else
329
+ do_it = false;
330
+ }
331
+ else
332
+ do_it = false;
333
+
334
+ spinlock_unlock(&host->rrdctx.hub_queue.spinlock);
335
+
336
+ if(dispatch_ready) {
337
worker_is_busy(WORKER_JOB_CHECK);
338
339
rrdcontext_lock(rc);
367
rrdcontext_dequeue_from_post_processing(rc);
368
rrdcontext_delete_from_sql_unsafe(rc);
369
327
- STRING *id = string_dup(rc->id);
370
+ STRING *delete_id = string_dup(rc->id);
371
rrdcontext_unlock(rc);
372
373
// delete it from the master dictionary
331
- if(!dictionary_del(host->rrdctx.contexts, string2str(rc->id)))
374
+ if(!dictionary_del(host->rrdctx.contexts, string2str(delete_id)))
375
netdata_log_error("RRDCONTEXT: '%s' of host '%s' failed to be deleted from rrdcontext dictionary.",
333
- string2str(id), rrdhost_hostname(host));
376
+ string2str(delete_id), rrdhost_hostname(host));
377
335
- string_freez(id);
378
+ string_freez(delete_id);
379
}
380
else
381
rrdcontext_unlock(rc);
386
387
dictionary_acquired_item_release(host->rrdctx.contexts, item);
388
}
389
+ string_freez(lookup_id);
390
+
391
spinlock_lock(&host->rrdctx.hub_queue.spinlock);
392
+ if(unlikely(!service_running(SERVICE_CONTEXT)))
393
+ break;
394
395
if(do_it) {
396
worker_is_busy(WORKER_JOB_DEQUEUE);
350
- rrdcontext_del_from_hub_queue(rc, true);
397
+ RRDCONTEXT *rc_at_idx = RRDCONTEXT_QUEUE_GET(&host->rrdctx.hub_queue, idx);
398
+ if(rc_at_idx == rc)
399
+ rrdcontext_del_from_hub_queue(rc, true);
400
}
401
}
402
spinlock_unlock(&host->rrdctx.hub_queue.spinlock);