Nilorea Library
C utilities for networking, threading, graphics
Loading...
Searching...
No Matches
n_kafka.c
Go to the documentation of this file.
1/*
2 * Nilorea Library
3 * Copyright (C) 2005-2026 Castagnier Mickael
4 *
5 * Licensed under the Apache License, Version 2.0 (the "License");
6 * you may not use this file except in compliance with the License.
7 * You may obtain a copy of the License at
8 *
9 * http://www.apache.org/licenses/LICENSE-2.0
10 *
11 * Unless required by applicable law or agreed to in writing, software
12 * distributed under the License is distributed on an "AS IS" BASIS,
13 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
14 * implied. See the License for the specific language governing
15 * permissions and limitations under the License.
16 *
17 * SPDX-License-Identifier: Apache-2.0
18 */
19
28#include "nilorea/n_kafka.h"
29#include "nilorea/n_common.h"
30#include "nilorea/n_base64.h"
31#include <limits.h>
32#include <sys/stat.h>
33#include <sys/types.h>
34
44static const char* n_kafka_event_fdesc(const N_KAFKA_EVENT* event) {
45 static _Thread_local char fragment[2 * PATH_MAX];
46 if (event && event->event_files_to_delete && event->event_files_to_delete->data) {
47 snprintf(fragment, sizeof(fragment), " (file: %s)", event->event_files_to_delete->data);
48 return fragment;
49 }
50 return "";
51}
52
58static int n_kafka_mkdir_p(const char* path) {
59 __n_assert(path, return FALSE);
60
61 size_t len = strlen(path);
62 if (len == 0)
63 return FALSE;
64
65 char* tmp = strdup(path);
66 __n_assert(tmp, return FALSE);
67
68 /* strip a trailing slash so the final component is created too */
69 if (tmp[len - 1] == '/')
70 tmp[len - 1] = '\0';
71
72 for (char* p = tmp + 1; *p; p++) {
73 if (*p == '/') {
74 *p = '\0';
75 if (mkdir(tmp, 0775) != 0 && errno != EEXIST) {
76 int error = errno;
77 n_log(LOG_ERR, "could not create directory \"%s\": %s", tmp, strerror(error));
78 free(tmp);
79 return FALSE;
80 }
81 *p = '/';
82 }
83 }
84 if (mkdir(tmp, 0775) != 0 && errno != EEXIST) {
85 int error = errno;
86 n_log(LOG_ERR, "could not create directory \"%s\": %s", tmp, strerror(error));
87 free(tmp);
88 return FALSE;
89 }
90 free(tmp);
91 return TRUE;
92}
93
100static int n_kafka_move_file_to_sended(const N_KAFKA* kafka, const char* filepath) {
101 __n_assert(kafka, return FALSE);
102 __n_assert(kafka->sended_dir, return FALSE);
103 __n_assert(filepath, return FALSE);
104
105 /* compute the path relative to the 'to-send' base directory; fall back to
106 * the basename when the file does not live under tosend_dir */
107 const char* relative = NULL;
108 if (kafka->tosend_dir) {
109 size_t base_len = strlen(kafka->tosend_dir);
110 if (strncmp(filepath, kafka->tosend_dir, base_len) == 0) {
111 relative = filepath + base_len;
112 while (*relative == '/')
113 relative++;
114 }
115 }
116 N_STR* dest = NULL;
117 if (relative && *relative) {
118 nstrprintf(dest, "%s/%s", kafka->sended_dir, relative);
119 } else {
120 char* tmp = strdup(filepath);
121 __n_assert(tmp, return FALSE);
122 nstrprintf(dest, "%s/%s", kafka->sended_dir, basename(tmp));
123 free(tmp);
124 }
125 __n_assert(dest, return FALSE);
126
127 /* ensure the destination directory exists */
128 char* dest_copy = strdup(_nstr(dest));
129 if (dest_copy) {
130 n_kafka_mkdir_p(dirname(dest_copy));
131 free(dest_copy);
132 }
133
134 int ret = TRUE;
135 if (rename(filepath, _nstr(dest)) != 0) {
136 int error = errno;
137 n_log(LOG_ERR, "could not move acknowledged file \"%s\" to \"%s\": %s", filepath, _nstr(dest), strerror(error));
138 ret = FALSE;
139 } else {
140 n_log(LOG_DEBUG, "moved acknowledged file \"%s\" to \"%s\"", filepath, _nstr(dest));
141 }
142 free_nstr(&dest);
143 return ret;
144}
145
151int32_t n_kafka_get_schema_from_char(const char* string) {
152 __n_assert(string, return -1);
153
154 uint32_t raw_schema_id = 0;
155 memcpy(&raw_schema_id, string + 1, sizeof(uint32_t));
156
157 return (int32_t)ntohl(raw_schema_id);
158}
159
165int32_t n_kafka_get_schema_from_nstr(const N_STR* string) {
166 __n_assert(string, return -1);
167 __n_assert(string->data, return -1);
168 __n_assert((string->written >= sizeof(int32_t)), return -1);
169 return n_kafka_get_schema_from_char(string->data);
170}
171
178int n_kafka_put_schema_in_char(char* string, int schema_id) {
179 __n_assert(string, return FALSE);
180 uint32_t schema_id_htonl = htonl((uint32_t)schema_id); // cast to unsigned to avoid warning
181 memcpy(string + 1, &schema_id_htonl, sizeof(uint32_t));
182 return TRUE;
183}
184
191int n_kafka_put_schema_in_nstr(N_STR* string, int schema_id) {
192 __n_assert(string, return FALSE);
193 __n_assert(string->data, return FALSE);
194 __n_assert((string->written >= sizeof(int32_t)), return FALSE);
195 return n_kafka_put_schema_in_char(string->data, schema_id);
196}
197
206int n_kafka_get_status(N_KAFKA* kafka, size_t* nb_queued, size_t* nb_waiting, size_t* nb_error) {
207 __n_assert(kafka, return FALSE);
208 __n_assert(nb_waiting, return FALSE);
209 __n_assert(nb_error, return FALSE);
210 read_lock(kafka->rwlock);
211 int status = kafka->polling_thread_status;
212 *nb_queued = kafka->nb_queued;
213 *nb_waiting = kafka->nb_waiting;
214 *nb_error = kafka->nb_error;
215 unlock(kafka->rwlock);
216 return status;
217}
218
225int n_kafka_set_error_timeout(N_KAFKA* kafka, int32_t error_timeout) {
226 __n_assert(kafka, return FALSE);
227 if (error_timeout < 0) {
228 n_log(LOG_ERR, "n_kafka_set_error_timeout: negative timeout %d rejected", error_timeout);
229 return FALSE;
230 }
231 write_lock(kafka->rwlock);
232 kafka->error_timeout = error_timeout;
233 unlock(kafka->rwlock);
234 return TRUE;
235}
236
244int n_kafka_set_sended_dir(N_KAFKA* kafka, const char* tosend_dir, const char* sended_dir) {
245 __n_assert(kafka, return FALSE);
246 write_lock(kafka->rwlock);
247 FreeNoLog(kafka->tosend_dir);
248 FreeNoLog(kafka->sended_dir);
249 if (tosend_dir)
250 kafka->tosend_dir = strdup(tosend_dir);
251 if (sended_dir)
252 kafka->sended_dir = strdup(sended_dir);
253 unlock(kafka->rwlock);
254 return TRUE;
255}
256
266// cppcheck-suppress[constParameterCallback,constParameter] ; callback signature must match rd_kafka typedef
267static void n_kafka_delivery_message_callback(rd_kafka_t* rk, const rd_kafka_message_t* rkmessage, void* opaque) {
268 (void)opaque;
269
270 __n_assert(rk, n_log(LOG_ERR, "rk=NULL is not a valid kafka handle"); return);
271 __n_assert(rkmessage, n_log(LOG_ERR, "rkmessage=NULL is not a valid kafka message"); return);
272
273 N_KAFKA_EVENT* event = (N_KAFKA_EVENT*)rkmessage->_private;
274
275 if (rkmessage->err) {
276 n_log(LOG_ERR, "message delivery failed for event %p%s: %s", (void*)event, n_kafka_event_fdesc(event), rd_kafka_err2str(rkmessage->err));
277 if (!event) {
278 n_log(LOG_ERR, "fatal: event is NULL");
279 return;
280 }
281 if (event->parent_table)
282 write_lock(event->parent_table->rwlock);
283 event->status = N_KAFKA_EVENT_ERROR;
284 event->error_time = time(NULL);
285 /* adjust counters: event was WAITING_ACK, now it's ERROR */
286 if (event->parent_table) {
287 event->parent_table->nb_waiting--;
288 event->parent_table->nb_error++;
289 }
290 if (event->parent_table)
291 unlock(event->parent_table->rwlock);
292 } else {
293 n_log(LOG_DEBUG, "message delivered (%ld bytes, partition %d)", rkmessage->len, rkmessage->partition);
294 if (event) {
295 // lock
296 if (event->parent_table)
297 write_lock(event->parent_table->rwlock);
298 // finalize produce event linked files: move them to the 'sended'
299 // directory when one is configured, otherwise delete them (legacy)
300 if (event->event_files_to_delete) {
301 int move_to_sended = (event->parent_table && event->parent_table->sended_dir) ? 1 : 0;
302 char** files_to_delete = split(_nstr(event->event_files_to_delete), ";", 0);
303 if (files_to_delete) {
304 if (split_count(files_to_delete) > 0) {
305 int index = 0;
306 while (files_to_delete[index]) {
307 if (move_to_sended) {
308 n_kafka_move_file_to_sended(event->parent_table, files_to_delete[index]);
309 } else {
310 int ret = unlink(files_to_delete[index]);
311 int error = errno;
312 if (ret == 0) {
313 n_log(LOG_DEBUG, "deleted on produce ack: %s", files_to_delete[index]);
314 } else {
315 n_log(LOG_ERR, "couldn't delete \"%s\": %s", files_to_delete[index], strerror(error));
316 }
317 }
318 index++;
319 }
320 } else {
321 n_log(LOG_ERR, "split result is empty !");
322 }
323 }
324 free_split_result(&files_to_delete);
325 }
326 // set status
327 event->status = N_KAFKA_EVENT_OK;
328 // unlock
329 if (event->parent_table)
330 unlock(event->parent_table->rwlock);
331 n_log(LOG_INFO, "kafka event %p%s received an ack !", (void*)event, n_kafka_event_fdesc(event));
332 } else {
333 n_log(LOG_ERR, "fatal: event is NULL");
334 }
335 }
336 return;
337 /* The rkmessage is destroyed automatically by librdkafka */
338}
339
345 __n_assert(kafka, return);
346
347 FreeNoLog(kafka->topic);
348 free_split_result(&kafka->topics);
349 FreeNoLog(kafka->event_cmd);
350 FreeNoLog(kafka->groupid);
351 FreeNoLog(kafka->tosend_dir);
352 FreeNoLog(kafka->sended_dir);
353
355
356 /* rd_kafka_topic_destroy must be called before rd_kafka_destroy */
357 if (kafka->rd_kafka_topic)
358 rd_kafka_topic_destroy(kafka->rd_kafka_topic);
359 kafka->rd_kafka_topic = NULL;
360
361 if (kafka->rd_kafka_handle) {
362 if (kafka->mode == RD_KAFKA_CONSUMER) {
363 /* close the consumer. May already be closed by stop_polling_thread,
364 * calling it again is safe (returns an error but does not crash) */
365 rd_kafka_consumer_close(kafka->rd_kafka_handle);
366 if (kafka->subscription)
367 rd_kafka_topic_partition_list_destroy(kafka->subscription);
368 }
369 if (kafka->mode == RD_KAFKA_PRODUCER) {
370 rd_kafka_flush(kafka->rd_kafka_handle, kafka->poll_timeout);
371 }
372 rd_kafka_destroy(kafka->rd_kafka_handle);
373 n_log(LOG_DEBUG, "kafka handle destroyed");
374 }
375
376 if (kafka->rd_kafka_conf) {
377 rd_kafka_conf_destroy(kafka->rd_kafka_conf);
378 }
379
380 if (kafka->configuration)
381 cJSON_Delete(kafka->configuration);
382
383 if (kafka->errstr)
384 free_nstr(&kafka->errstr);
385
386 if (kafka->events_to_send)
388
389 if (kafka->received_events)
391
392 rw_lock_destroy(kafka->rwlock);
393
394 Free(kafka->bootstrap_servers);
395
396 Free(kafka);
397 return;
398}
399
407N_KAFKA* n_kafka_new(int32_t poll_timeout, int32_t poll_interval, size_t errstr_len) {
408 N_KAFKA* kafka = NULL;
409 Malloc(kafka, N_KAFKA, 1);
410 __n_assert(kafka, return NULL);
411
412 kafka->errstr = new_nstr(errstr_len);
413 __n_assert(kafka->errstr, Free(kafka); return NULL);
414
415 kafka->events_to_send = NULL;
416 kafka->received_events = NULL;
417 kafka->rd_kafka_conf = NULL;
418 kafka->rd_kafka_handle = NULL;
419 kafka->configuration = NULL;
420 kafka->groupid = NULL;
421 kafka->topics = NULL;
422 kafka->event_cmd = NULL;
423 kafka->subscription = NULL;
424 kafka->topic = NULL;
425 kafka->mode = -1;
426 kafka->schema_id = -1;
427 kafka->poll_timeout = poll_timeout;
428 kafka->poll_interval = poll_interval;
429 kafka->monitored_directory_interval = 3000; // 3000 msecs as a default
430 kafka->nb_queued = 0;
431 kafka->nb_waiting = 0;
432 kafka->nb_error = 0;
433 kafka->error_timeout = 0;
434 kafka->event_consumption_enabled = TRUE;
435 kafka->event_production_enabled = TRUE;
436 kafka->polling_thread_status = 0;
437 kafka->bootstrap_servers = NULL;
438 kafka->is_transactional = FALSE;
439 kafka->tosend_dir = NULL;
440 kafka->sended_dir = NULL;
441
443 __n_assert(kafka->events_to_send, n_kafka_delete(kafka); return NULL);
444
446 __n_assert(kafka->received_events, n_kafka_delete(kafka); return NULL);
447
448 if (init_lock(kafka->rwlock) != 0) {
449 n_log(LOG_ERR, "could not init kafka rwlock in kafka structure at address %p", kafka);
450 n_kafka_delete(kafka);
451 return NULL;
452 }
453
454 return kafka;
455}
456
464 __n_assert(config_file, return NULL);
465
466 N_KAFKA* kafka = NULL;
467
468 kafka = n_kafka_new(-1, 100, 1024);
469 __n_assert(kafka, return NULL);
470
471 // initialize kafka object
472 kafka->rd_kafka_conf = rd_kafka_conf_new();
473
474 // load config file
475 N_STR* config_string = NULL;
476 config_string = file_to_nstr(config_file);
477 if (!config_string) {
478 n_log(LOG_ERR, "unable to read config from file %s !", config_file);
479 n_kafka_delete(kafka);
480 return NULL;
481 }
482 cJSON* json = NULL;
483 json = cJSON_Parse(_nstrp(config_string));
484 /* config_string is no longer needed after parsing */
485 free_nstr(&config_string);
486 if (!json) {
487 n_log(LOG_ERR, "unable to parse json from file %s", config_file);
488 n_kafka_delete(kafka);
489 return NULL;
490 }
491
492 int jsonIndex;
493 for (jsonIndex = 0; jsonIndex < cJSON_GetArraySize(json); jsonIndex++) {
494 cJSON* entry = cJSON_GetArrayItem(json, jsonIndex);
495
496 if (!entry) continue;
497 __n_assert(entry->string, continue);
498 if (!entry->valuestring) {
499 n_log(LOG_DEBUG, "no valuestring for entry %s", _str(entry->string));
500 continue;
501 }
502
503 if (entry->string[0] != '-') {
504 // if it's not one of the optionnal parameters not managed by kafka, then we can use rd_kafka_conf_set on them
505 if (strcmp("topic", entry->string) != 0 &&
506 strcmp("topics", entry->string) != 0 &&
507 strcmp("event_cmd", entry->string) != 0 &&
508 strcmp("value.schema.id", entry->string) != 0 &&
509 strcmp("value.schema.type", entry->string) != 0 &&
510 strcmp("poll.interval", entry->string) != 0 &&
511 strcmp("poll.timeout", entry->string) != 0 &&
512 strcmp("group.id.autogen", entry->string) != 0 &&
513 strcmp("monitored.directory.interval", entry->string) != 0) {
514 if (!strcmp("group.id", entry->string)) {
515 // exclude group id for producer
516 if (mode == RD_KAFKA_PRODUCER)
517 continue;
518 kafka->groupid = strdup(entry->valuestring);
519 }
520
521 if (rd_kafka_conf_set(kafka->rd_kafka_conf, entry->string, entry->valuestring, _nstr(kafka->errstr), kafka->errstr->length) != RD_KAFKA_CONF_OK) {
522 n_log(LOG_ERR, "kafka config: %s", _nstr(kafka->errstr));
523 } else {
524 n_log(LOG_DEBUG, "kafka config enabled: %s => %s", entry->string, entry->valuestring);
525 }
526 }
527 } else {
528 n_log(LOG_DEBUG, "kafka disabled config: %s => %s", entry->string, entry->valuestring);
529 }
530 }
531
532 // other parameters, not directly managed by kafka API (will cause an error if used along rd_kafka_conf_set )
533 // producer topic
534 cJSON* jstr = NULL;
535 jstr = cJSON_GetObjectItem(json, "topic");
536 if (jstr && jstr->valuestring) {
537 kafka->topic = strdup(jstr->valuestring);
538 n_log(LOG_DEBUG, "kafka producer topic: %s", kafka->topic);
539 } else {
540 if (mode == RD_KAFKA_PRODUCER) {
541 n_log(LOG_ERR, "no topic configured !");
542 cJSON_Delete(json);
543 n_kafka_delete(kafka);
544 return NULL;
545 }
546 }
547 // consumer topics
548 jstr = cJSON_GetObjectItem(json, "topics");
549 if (jstr && jstr->valuestring) {
550 kafka->topics = split(jstr->valuestring, ",", 0);
551 n_log(LOG_DEBUG, "kafka consumer topics: %s", jstr->valuestring);
552 } else {
553 if (mode == RD_KAFKA_CONSUMER) {
554 n_log(LOG_ERR, "no topics configured !");
555 cJSON_Delete(json);
556 n_kafka_delete(kafka);
557 return NULL;
558 }
559 }
560
561 // eventual consumer custom event_cmd
562 jstr = cJSON_GetObjectItem(json, "event_cmd");
563 if (jstr && jstr->valuestring) {
564 kafka->event_cmd = strdup(jstr->valuestring);
565 n_log(LOG_DEBUG, "kafka consumer event_cmd: %s", jstr->valuestring);
566 }
567
568 // schema id if any
569 jstr = cJSON_GetObjectItem(json, "value.schema.id");
570 if (jstr && jstr->valuestring) {
571 int schem_v = atoi(jstr->valuestring);
572 if (schem_v < -1 || schem_v > 9999) {
573 n_log(LOG_ERR, "invalid schema id %d", schem_v);
574 cJSON_Delete(json);
575 n_kafka_delete(kafka);
576 return NULL;
577 }
578 n_log(LOG_DEBUG, "kafka schema id: %d", schem_v);
579 kafka->schema_id = schem_v;
580 }
581
582 // kafka broker poll interval
583 jstr = cJSON_GetObjectItem(json, "poll.interval");
584 if (jstr && jstr->valuestring) {
585 kafka->poll_interval = atoi(jstr->valuestring);
586 n_log(LOG_DEBUG, "kafka poll interval: %d", kafka->poll_interval);
587 }
588
589 // kafka broker poll timeout
590 jstr = cJSON_GetObjectItem(json, "poll.timeout");
591 if (jstr && jstr->valuestring) {
592 kafka->poll_timeout = atoi(jstr->valuestring);
593 n_log(LOG_DEBUG, "kafka poll timeout: %d", kafka->poll_timeout);
594 }
595
596 // local directory poll interval
597 jstr = cJSON_GetObjectItem(json, "monitored.directory.interval");
598 if (jstr && jstr->valuestring) {
599 kafka->monitored_directory_interval = atoi(jstr->valuestring);
600 n_log(LOG_DEBUG, "kafka monitored directory interval: %d", kafka->monitored_directory_interval);
601 }
602
603 // bootstrap servers
604 jstr = cJSON_GetObjectItem(json, "bootstrap.servers");
605 if (jstr && jstr->valuestring) {
606 kafka->bootstrap_servers = strdup(jstr->valuestring);
607 n_log(LOG_DEBUG, "kafka bootstrap server: %s", kafka->bootstrap_servers);
608 }
609
610 if (mode == RD_KAFKA_PRODUCER) {
611 // a configured 'transactional.id' switches the producer to transactional mode.
612 // the key itself is already forwarded to rd_kafka_conf_set by the generic
613 // loop above (it is a regular librdkafka config key), we only detect it here.
614 jstr = cJSON_GetObjectItem(json, "transactional.id");
615 if (jstr && jstr->valuestring && jstr->valuestring[0] != '\0') {
616 kafka->is_transactional = TRUE;
617 n_log(LOG_DEBUG, "kafka producer transactional.id: %s, enabling transactional produces", jstr->valuestring);
618 }
619
620 // set delivery callback
621 rd_kafka_conf_set_dr_msg_cb(kafka->rd_kafka_conf, n_kafka_delivery_message_callback);
622
623 kafka->rd_kafka_handle = rd_kafka_new(RD_KAFKA_PRODUCER, kafka->rd_kafka_conf, _nstr(kafka->errstr), kafka->errstr->length);
624 if (!kafka->rd_kafka_handle) {
625 n_log(LOG_ERR, "failed to create new producer: %s", _nstr(kafka->errstr));
626 cJSON_Delete(json);
627 n_kafka_delete(kafka);
628 return NULL;
629 }
630 // conf is now owned by kafka handle
631 kafka->rd_kafka_conf = NULL;
632 // Create topic object
633 kafka->rd_kafka_topic = rd_kafka_topic_new(kafka->rd_kafka_handle, kafka->topic, NULL);
634 if (!kafka->rd_kafka_topic) {
635 cJSON_Delete(json);
636 n_kafka_delete(kafka);
637 return NULL;
638 }
639 // initialize the transactional producer (acquires the producer id/epoch).
640 // -1 lets librdkafka use 2 * transaction.timeout.ms as the blocking timeout.
641 if (kafka->is_transactional) {
642 rd_kafka_error_t* txn_err = rd_kafka_init_transactions(kafka->rd_kafka_handle, -1);
643 if (txn_err) {
644 n_log(LOG_ERR, "failed to init transactions for producer %p, topic %s: %s", kafka->rd_kafka_handle, kafka->topic, rd_kafka_error_string(txn_err));
645 rd_kafka_error_destroy(txn_err);
646 cJSON_Delete(json);
647 n_kafka_delete(kafka);
648 return NULL;
649 }
650 n_log(LOG_DEBUG, "transactions initialized for producer %p, topic %s", kafka->rd_kafka_handle, kafka->topic);
651 }
652 kafka->mode = RD_KAFKA_PRODUCER;
653 } else if (mode == RD_KAFKA_CONSUMER) {
654 /* If there is no previously committed offset for a partition
655 * the auto.offset.reset strategy will be used to decide where
656 * in the partition to start fetching messages.
657 * By setting this to earliest the consumer will read all messages
658 * in the partition if there was no previously committed offset. */
659 /*if (rd_kafka_conf_set(kafka -> rd_kafka_conf, "auto.offset.reset", "earliest", _nstr( kafka -> errstr ) , kafka -> errstr -> length) != RD_KAFKA_CONF_OK) {
660 n_log( LOG_ERR , "kafka conf set: %s", kafka -> errstr);
661 n_kafka_delete( kafka );
662 return NULL ;
663 }*/
664
665 // if groupid is not set, generate a unique one
666 if (!kafka->groupid) {
667 char computer_name[1024] = "";
668 get_computer_name(computer_name, 1024);
669 N_STR* groupid = new_nstr(1024);
670 char* topics = join(kafka->topics, "_");
671 // generating group id
672 jstr = cJSON_GetObjectItem(json, "group.id.autogen");
673 if (jstr && jstr->valuestring) {
674 if (strcmp(jstr->valuestring, "host-topic-group") == 0) {
675 // nstrprintf( groupid , "%s_%s_%s" , computer_name , topics ,kafka -> bootstrap_servers );
676 nstrprintf(groupid, "%s_%s", computer_name, topics);
677 } else if (strcmp(jstr->valuestring, "unique-group") == 0) {
678 // nstrprintf( groupid , "%s_%s_%s_%d" , computer_name , topics , kafka -> bootstrap_servers , getpid() );
679 nstrprintf(groupid, "%s_%s_%d", computer_name, topics, getpid());
680 }
681 } else // default unique group
682 {
683 // nstrprintf( groupid , "%s_%s_%s_%d" , computer_name , topics , kafka -> bootstrap_servers , getpid() );
684 nstrprintf(groupid, "%s_%s_%d", computer_name, topics, getpid());
685 n_log(LOG_DEBUG, "group.id is not set and group.id.autogen is not set, generated unique group id: %s", _nstr(groupid));
686 }
687 free(topics);
688 kafka->groupid = groupid->data;
689 groupid->data = NULL;
690 free_nstr(&groupid);
691 }
692 if (rd_kafka_conf_set(kafka->rd_kafka_conf, "group.id", kafka->groupid, _nstr(kafka->errstr), kafka->errstr->length) != RD_KAFKA_CONF_OK) {
693 n_log(LOG_ERR, "kafka consumer group.id error: %s", _nstr(kafka->errstr));
694 cJSON_Delete(json);
695 n_kafka_delete(kafka);
696 return NULL;
697 } else {
698 n_log(LOG_DEBUG, "kafka consumer group.id => %s", kafka->groupid);
699 }
700
701 /* Create consumer instance.
702 * NOTE: rd_kafka_new() takes ownership of the conf object
703 * and the application must not reference it again after
704 * this call.
705 */
706 kafka->rd_kafka_handle = rd_kafka_new(RD_KAFKA_CONSUMER, kafka->rd_kafka_conf, _nstr(kafka->errstr), kafka->errstr->length);
707 if (!kafka->rd_kafka_handle) {
708 n_log(LOG_ERR, "%% Failed to create new consumer: %s", kafka->errstr);
709 cJSON_Delete(json);
710 n_kafka_delete(kafka);
711 return NULL;
712 }
713 // conf is now owned by kafka handle
714 kafka->rd_kafka_conf = NULL;
715
716 /* Redirect all messages from per-partition queues to
717 * the main queue so that messages can be consumed with one
718 * call from all assigned partitions.
719 *
720 * The alternative is to poll the main queue (for events)
721 * and each partition queue separately, which requires setting
722 * up a rebalance callback and keeping track of the assignment:
723 * but that is more complex and typically not recommended. */
724 rd_kafka_poll_set_consumer(kafka->rd_kafka_handle);
725
726 /* Convert the list of topics to a format suitable for librdkafka */
727 int topic_cnt = split_count(kafka->topics);
728 kafka->subscription = rd_kafka_topic_partition_list_new(topic_cnt);
729 for (int i = 0; i < topic_cnt; i++)
730 rd_kafka_topic_partition_list_add(kafka->subscription, kafka->topics[i],
731 /* the partition is ignored
732 * by subscribe() */
733 RD_KAFKA_PARTITION_UA);
734 /* Assign the topic. This method is disabled as it does noat allow dynamic partition assignement
735 int err = rd_kafka_assign(kafka -> rd_kafka_handle, kafka -> subscription);
736 if( err )
737 {
738 n_log( LOG_ERR , "kafka consumer: failed to assign %d topics: %s", kafka -> subscription->cnt, rd_kafka_err2str(err));
739 n_kafka_delete( kafka );
740 return NULL;
741 } */
742
743 /* Subscribe to the list of topics */
744 int err = rd_kafka_subscribe(kafka->rd_kafka_handle, kafka->subscription);
745 if (err) {
746 n_log(LOG_ERR, "kafka consumer: failed to subscribe to %d topics: %s", kafka->subscription->cnt, rd_kafka_err2str(err));
747 cJSON_Delete(json);
748 n_kafka_delete(kafka);
749 return NULL;
750 }
751
752 n_log(LOG_DEBUG, "kafka consumer created and subscribed to %d topic(s), waiting for rebalance and messages...", kafka->subscription->cnt);
753
754 kafka->mode = RD_KAFKA_CONSUMER;
755 } else {
756 n_log(LOG_ERR, "invalid mode %d", mode);
757 cJSON_Delete(json);
758 n_kafka_delete(kafka);
759 return NULL;
760 }
761
762 cJSON_Delete(json);
763 return kafka;
764} /* n_kafka_load_config */
765
772 N_KAFKA_EVENT* event = NULL;
773 Malloc(event, N_KAFKA_EVENT, 1);
774 __n_assert(event, return NULL);
775
776 event->event_string = NULL;
777 event->event_files_to_delete = NULL;
778 event->from_topic = NULL;
779 event->rd_kafka_headers = NULL;
780 event->received_headers = NULL;
781 event->schema_id = schema_id;
782 event->status = N_KAFKA_EVENT_CREATED;
783 event->parent_table = NULL;
784 event->error_time = 0;
785
786 return event;
787} /* n_kafka_new_event */
788
795int n_kafka_new_headers(N_KAFKA_EVENT* event, size_t count) {
796 __n_assert(event, return FALSE);
797
798 if (count == 0)
799 count = 1;
800
801 if (!event->rd_kafka_headers) {
802 event->rd_kafka_headers = rd_kafka_headers_new(count);
803 __n_assert(event->rd_kafka_headers, return FALSE);
804 } else {
805 n_log(LOG_ERR, "event headers already allocated for event %p%s", event, n_kafka_event_fdesc(event));
806 return FALSE;
807 }
808 return TRUE;
809}
810
820int n_kafka_add_header_ex(N_KAFKA_EVENT* event, char* key, size_t key_length, char* value, size_t value_length) {
821 __n_assert(event, return FALSE);
822 __n_assert(event->rd_kafka_headers, return FALSE);
823 __n_assert(key, return FALSE);
824 __n_assert(value, return FALSE);
825
826 if (key_length < 1 || key_length > SSIZE_MAX) {
827 n_log(LOG_ERR, "Invalid key length (%zu) for header in event %p%s", key_length, event, n_kafka_event_fdesc(event));
828 return FALSE;
829 }
830
831 if (value_length < 1 || value_length > SSIZE_MAX) {
832 n_log(LOG_ERR, "Invalid value length (%zu) for key '%s' in event %p%s", value_length, key, event, n_kafka_event_fdesc(event));
833 return FALSE;
834 }
835
836 rd_kafka_resp_err_t err = rd_kafka_header_add(event->rd_kafka_headers, key, (ssize_t)key_length, value, (ssize_t)value_length);
837
838 if (err) {
839 n_log(LOG_ERR, "Failed to add header [%s:%zu=%s:%zu] to event %p%s: %s",
840 key, key_length, value, value_length, event, n_kafka_event_fdesc(event), rd_kafka_err2str(err));
841 return FALSE;
842 }
843
844 return TRUE;
845}
846
855 __n_assert(event, return FALSE);
856 __n_assert(key, return FALSE);
857 __n_assert(key->data, return FALSE);
858 __n_assert(value, return FALSE);
859 __n_assert(value->data, return FALSE);
860
861 return n_kafka_add_header_ex(event, key->data, key->written, value->data, value->written);
862}
863
871 __n_assert(kafka, return FALSE);
872 __n_assert(event, return FALSE);
873
874 event->parent_table = kafka;
875
876 write_lock(kafka->rwlock);
877 kafka->nb_queued++;
879 unlock(kafka->rwlock);
880
881 // Success to *enqueue* event
882 n_log(LOG_DEBUG, "successfully enqueued event %p%s in producer %p waitlist, topic: %s", event, n_kafka_event_fdesc(event), kafka->rd_kafka_handle, kafka->topic);
883
884 return TRUE;
885} /* n_kafka_produce */
886
894 char* event_string = NULL;
895 size_t event_length = 0;
896
897 event->parent_table = kafka;
898
899 event_string = event->event_string->data;
900 event_length = event->event_string->written;
901
902 if (event->rd_kafka_headers) {
903 rd_kafka_headers_t* hdrs_copy;
904 hdrs_copy = rd_kafka_headers_copy(event->rd_kafka_headers);
905
906 rd_kafka_resp_err_t err = rd_kafka_producev(
907 kafka->rd_kafka_handle, RD_KAFKA_V_RKT(kafka->rd_kafka_topic),
908 RD_KAFKA_V_PARTITION(RD_KAFKA_PARTITION_UA),
909 RD_KAFKA_V_MSGFLAGS(RD_KAFKA_MSG_F_COPY),
910 RD_KAFKA_V_VALUE(event_string, event_length),
911 RD_KAFKA_V_HEADERS(hdrs_copy),
912 RD_KAFKA_V_OPAQUE((void*)event),
913 RD_KAFKA_V_END);
914
915 if (err) {
916 rd_kafka_headers_destroy(hdrs_copy);
917 event->status = N_KAFKA_EVENT_ERROR;
918 event->error_time = time(NULL);
919 n_log(LOG_ERR, "failed to produce event: %p%s with headers %p, producer: %p, topic: %s, error: %s", event, n_kafka_event_fdesc(event), event->rd_kafka_headers, kafka->rd_kafka_handle, kafka->topic, rd_kafka_err2str(err));
920 return FALSE;
921 }
922 } else {
923 if (rd_kafka_produce(kafka->rd_kafka_topic, RD_KAFKA_PARTITION_UA, RD_KAFKA_MSG_F_COPY, event_string, event_length, NULL, 0, event) == -1) {
924 int error = errno;
925 event->status = N_KAFKA_EVENT_ERROR;
926 event->error_time = time(NULL);
927 n_log(LOG_ERR, "failed to produce event: %p%s, producer: %p, topic: %s, error: %s", event, n_kafka_event_fdesc(event), kafka->rd_kafka_handle, kafka->topic, strerror(error));
928 return FALSE;
929 }
930 }
931 // Success to *enqueue* event
932 n_log(LOG_DEBUG, "successfully enqueued event %p%s in local producer %p : %s", event, n_kafka_event_fdesc(event), kafka->rd_kafka_handle, kafka->topic);
933
934 event->status = N_KAFKA_EVENT_WAITING_ACK;
935
936 return TRUE;
937} /* n_kafka_produce_ex */
938
946N_KAFKA_EVENT* n_kafka_new_event_from_char(const char* string, size_t written, int schema_id) {
947 __n_assert(string, return NULL);
948
949 size_t offset = 0;
950 if (schema_id != -1)
951 offset = 5;
952
953 N_KAFKA_EVENT* event = n_kafka_new_event(schema_id);
954 __n_assert(event, return NULL);
955
956 Malloc(event->event_string, N_STR, 1);
957 __n_assert(event->event_string, free(event); return NULL);
958
959 // allocate the size of the event + (the size of the schema id + magic byte) + one ending \0
960 size_t length = written + offset;
961 Malloc(event->event_string->data, char, length);
962 __n_assert(event->event_string->data, free(event->event_string); free(event); return NULL);
963 event->event_string->length = length;
964
965 // copy incomming body
966 memcpy(event->event_string->data + offset, string, written);
967 event->event_string->written = written + offset;
968
969 if (schema_id != -1)
970 n_kafka_put_schema_in_nstr(event->event_string, schema_id);
971
972 // set status and schema id
973 event->status = N_KAFKA_EVENT_QUEUED;
974
975 return event;
976} /* n_kafka_new_event_from_char */
977
984N_KAFKA_EVENT* n_kafka_new_event_from_string(const N_STR* string, int schema_id) {
985 __n_assert(string, return NULL);
986 __n_assert(string->data, return NULL);
987
988 N_KAFKA_EVENT* event = n_kafka_new_event_from_char(string->data, string->written, schema_id);
989
990 return event;
991} /* n_kafka_new_event_from_string */
992
999N_KAFKA_EVENT* n_kafka_new_event_from_file(char* filename, int schema_id) {
1000 __n_assert(filename, return NULL);
1001
1002 N_STR* from = file_to_nstr(filename);
1003 __n_assert(from, return NULL);
1004
1005 N_KAFKA_EVENT* event = n_kafka_new_event_from_string(from, schema_id);
1006 free_nstr(&from);
1007
1008 return event;
1009} /* n_kafka_new_event_from_file */
1010
1015void n_kafka_event_destroy_ptr(void* event_ptr) {
1016 __n_assert(event_ptr, return);
1017 N_KAFKA_EVENT* event = (N_KAFKA_EVENT*)event_ptr;
1018 __n_assert(event, return);
1019
1020 if (event->event_string)
1021 free_nstr(&event->event_string);
1022
1023 if (event->event_files_to_delete)
1024 free_nstr(&event->event_files_to_delete);
1025
1026 FreeNoLog(event->from_topic);
1027
1028 if (event->rd_kafka_headers)
1029 rd_kafka_headers_destroy(event->rd_kafka_headers);
1030
1031 if (event->received_headers)
1032 list_destroy(&event->received_headers);
1033
1034 free(event);
1035 return;
1036} /* n_kafka_event_destroy_ptr */
1037
1044 __n_assert(event && (*event), return FALSE);
1045 n_kafka_event_destroy_ptr((*event));
1046 (*event) = NULL;
1047 return TRUE;
1048} /* n_kafka_event_destroy */
1049
1056 __n_assert(kafka, return FALSE);
1057
1058 int event_consumption_enabled;
1059 int event_production_enabled;
1060
1061 read_lock(kafka->rwlock);
1062 event_consumption_enabled = kafka->event_consumption_enabled;
1063 event_production_enabled = kafka->event_production_enabled;
1064 unlock(kafka->rwlock);
1065
1066 // wait poll interval msecs for kafka response
1067 if (kafka->mode == RD_KAFKA_PRODUCER) {
1068 if (event_production_enabled) {
1069 int nb_events = rd_kafka_poll(kafka->rd_kafka_handle, kafka->poll_interval);
1070 (void)nb_events;
1071 /* transactional mode: a kafka transaction is opened lazily before the
1072 * first queued event is produced in this cycle, and committed after the
1073 * locked pass below. commit/abort must run outside the rwlock because
1074 * they synchronously fire the delivery callbacks which take that lock. */
1075 int txn_active = 0;
1076 int txn_begin_failed = 0;
1077 write_lock(kafka->rwlock);
1078 // check events status in event table
1079 LIST_NODE* node = kafka->events_to_send->start;
1080 while (node) {
1081 N_KAFKA_EVENT* event = (N_KAFKA_EVENT*)node->ptr;
1082 if (event->status == N_KAFKA_EVENT_OK) {
1083 kafka->nb_waiting--;
1084 n_log(LOG_DEBUG, "removing event OK %p%s", (void*)event, n_kafka_event_fdesc(event));
1085 LIST_NODE* node_to_kill = node;
1086 node = node->next;
1087 N_KAFKA_EVENT* event_to_kill = remove_list_node(kafka->events_to_send, node_to_kill, N_KAFKA_EVENT);
1088 n_kafka_event_destroy(&event_to_kill);
1089 continue;
1090 } else if (event->status == N_KAFKA_EVENT_QUEUED) {
1091 /* in transactional mode, open a transaction before the first produce */
1092 if (kafka->is_transactional && !txn_active && !txn_begin_failed) {
1093 rd_kafka_error_t* txn_err = rd_kafka_begin_transaction(kafka->rd_kafka_handle);
1094 if (txn_err) {
1095 n_log(LOG_ERR, "could not begin transaction on producer %p, topic %s: %s%s", kafka->rd_kafka_handle, kafka->topic, rd_kafka_error_string(txn_err), rd_kafka_error_is_fatal(txn_err) ? " (FATAL)" : "");
1096 rd_kafka_error_destroy(txn_err);
1097 /* leave the queued events untouched, retry on the next poll */
1098 txn_begin_failed = 1;
1099 } else {
1100 txn_active = 1;
1101 }
1102 }
1103 /* if a transaction was required but could not be opened, skip producing this cycle */
1104 if (kafka->is_transactional && !txn_active) {
1105 node = node->next;
1106 continue;
1107 }
1108 if (n_kafka_produce_ex(kafka, event) == FALSE) {
1109 /* produce_ex already set status to ERROR,
1110 * adjust counters: QUEUED -> ERROR */
1111 kafka->nb_queued--;
1112 kafka->nb_error++;
1113 } else {
1114 kafka->nb_waiting++;
1115 kafka->nb_queued--;
1116 }
1117 } else if (event->status == N_KAFKA_EVENT_ERROR) {
1118 /* counters already adjusted at transition time
1119 * (in produce_ex failure above or in delivery callback) */
1120 if (kafka->error_timeout > 0 &&
1121 (time(NULL) - event->error_time) >= (time_t)kafka->error_timeout) {
1122 n_log(LOG_INFO, "retrying errored event %p%s after %d s", (void*)event, n_kafka_event_fdesc(event), kafka->error_timeout);
1123 event->status = N_KAFKA_EVENT_QUEUED;
1124 event->error_time = 0;
1125 kafka->nb_error--;
1126 kafka->nb_queued++;
1127 }
1128 }
1129 node = node->next;
1130 }
1131 unlock(kafka->rwlock);
1132
1133 /* close the transaction outside the lock: commit_transaction (and
1134 * abort_transaction) drive the delivery callbacks that, on success,
1135 * move the produced files to 'sended' and, on failure, flip events
1136 * back to ERROR. Both callbacks take the rwlock themselves.
1137 * Retry the commit on retriable errors; on any other failure (or once
1138 * retries are exhausted) abort to bring the producer back to a clean
1139 * state: the purged messages are reported to the delivery callback as
1140 * errors so their files stay in 'to-send' and are produced again later. */
1141 if (txn_active) {
1142 int committed = 0;
1143 int attempts = 0;
1144 while (!committed && attempts < 5) {
1145 attempts++;
1146 rd_kafka_error_t* commit_err = rd_kafka_commit_transaction(kafka->rd_kafka_handle, -1);
1147 if (!commit_err) {
1148 committed = 1;
1149 break;
1150 }
1151 n_log(LOG_ERR, "could not commit transaction (attempt %d) on producer %p, topic %s: %s", attempts, kafka->rd_kafka_handle, kafka->topic, rd_kafka_error_string(commit_err));
1152 int retriable = rd_kafka_error_is_retriable(commit_err);
1153 int fatal = rd_kafka_error_is_fatal(commit_err);
1154 rd_kafka_error_destroy(commit_err);
1155 if (retriable)
1156 continue; // resume the in-flight commit
1157 if (fatal)
1158 n_log(LOG_ERR, "FATAL transaction error on producer %p, topic %s, producer is no longer usable", kafka->rd_kafka_handle, kafka->topic);
1159 break; // abort-required, fatal or other non-retriable error
1160 }
1161 if (!committed) {
1162 n_log(LOG_ERR, "aborting transaction on producer %p, topic %s", kafka->rd_kafka_handle, kafka->topic);
1163 rd_kafka_error_t* abort_err = rd_kafka_abort_transaction(kafka->rd_kafka_handle, -1);
1164 if (abort_err) {
1165 n_log(LOG_ERR, "could not abort transaction on producer %p, topic %s: %s", kafka->rd_kafka_handle, kafka->topic, rd_kafka_error_string(abort_err));
1166 rd_kafka_error_destroy(abort_err);
1167 }
1168 }
1169 }
1170 } else {
1171 uint64_t sleep_us = (kafka->poll_interval > 0) ? ((uint64_t)(uint32_t)kafka->poll_interval * 1000ULL) : 0ULL;
1172 usleep((unsigned int)(sleep_us > (uint64_t)UINT_MAX ? (uint64_t)UINT_MAX : sleep_us)); // production is disabled, sleep instead of producing
1173 }
1174 } else if (kafka->mode == RD_KAFKA_CONSUMER) {
1175 if (event_consumption_enabled) {
1176 rd_kafka_message_t* rkm = NULL;
1177 while (event_consumption_enabled && (rkm = rd_kafka_consumer_poll(kafka->rd_kafka_handle, kafka->poll_interval))) {
1178 read_lock(kafka->rwlock);
1179 event_consumption_enabled = kafka->event_consumption_enabled;
1180 unlock(kafka->rwlock);
1181 if (rkm->err) {
1182 /* Consumer errors are generally to be considered
1183 * informational as the consumer will automatically
1184 * try to recover from all types of errors. */
1185 n_log(LOG_ERR, "consumer: %s", rd_kafka_message_errstr(rkm));
1186 rd_kafka_message_destroy(rkm);
1187 uint64_t sleep_us = (kafka->poll_interval > 0) ? ((uint64_t)(uint32_t)kafka->poll_interval * 1000ULL) : 0ULL;
1188 usleep((unsigned int)(sleep_us > (uint64_t)UINT_MAX ? (uint64_t)UINT_MAX : sleep_us));
1189 return FALSE;
1190 }
1191 // Reminder of rkm contents
1192 n_log(LOG_DEBUG, "Message on %s [%" PRId32 "] at offset %" PRId64 " (leader epoch %" PRId32 ")", rd_kafka_topic_name(rkm->rkt), rkm->partition, rkm->offset, rd_kafka_message_leader_epoch(rkm));
1193
1194 // Print the message key
1195 if (rkm->key && rkm->key_len > 0)
1196 n_log(LOG_DEBUG, "Key: %.*s", (int)rkm->key_len, (const char*)rkm->key);
1197 else if (rkm->key)
1198 n_log(LOG_DEBUG, "Key: (%d bytes)", (int)rkm->key_len);
1199
1200 if (rkm->payload && rkm->len > 0) {
1201 write_lock(kafka->rwlock);
1202 // make a copy of the event for further processing
1203 N_KAFKA_EVENT* event = NULL;
1204 event = n_kafka_new_event_from_char(rkm->payload, rkm->len, -1); // no schema id because we want a full raw copy here
1205 event->parent_table = kafka;
1206 // test if there are headers, save them
1207 rd_kafka_headers_t* hdrs = NULL;
1208 if (!rd_kafka_message_headers(rkm, &hdrs)) {
1209 size_t idx = 0;
1210 const char* name = NULL;
1211 const void* val = NULL;
1212 size_t size = 0;
1213 event->received_headers = new_generic_list(MAX_LIST_ITEMS);
1214 while (!rd_kafka_header_get_all(hdrs, idx, &name, &val, &size)) {
1215 N_STR* header_entry = NULL;
1216 nstrprintf(header_entry, "%s=%s", _str(name), _str((char*)val));
1217 list_push(event->received_headers, header_entry, &free_nstr_ptr);
1218 idx++;
1219 }
1220 }
1221 // save originating topic (can help sorting if there are multiples one)
1222 event->from_topic = strdup(rd_kafka_topic_name(rkm->rkt));
1223 if (kafka->schema_id != -1)
1224 event->schema_id = n_kafka_get_schema_from_nstr(event->event_string);
1226 n_log(LOG_DEBUG, "Consumer received event of (%d bytes) from topic %s", (int)rkm->len, event->from_topic);
1227 unlock(kafka->rwlock);
1228 }
1229 rd_kafka_message_destroy(rkm);
1230 }
1231 } else {
1232 uint64_t sleep_us = (kafka->poll_interval > 0) ? ((uint64_t)(uint32_t)kafka->poll_interval * 1000ULL) : 0ULL;
1233 usleep((unsigned int)(sleep_us > (uint64_t)UINT_MAX ? (uint64_t)UINT_MAX : sleep_us)); // consumption is disabled, sleep instead of consuming
1234 }
1235 }
1236 // n_log( LOG_DEBUG , "kafka poll for handle %p returned %d elements" , kafka -> rd_kafka_handle , nb_events );
1237 return TRUE;
1238} /* n_kafka_poll */
1239
1245void* n_kafka_polling_thread(void* ptr) {
1246 N_KAFKA* kafka = (N_KAFKA*)ptr;
1247
1248 int status = 1;
1249
1250 N_TIME chrono;
1251
1252 if (kafka->mode == RD_KAFKA_PRODUCER)
1253 n_log(LOG_DEBUG, "starting polling thread for kafka handler %p mode PRODUCER (%d) topic %s", kafka->rd_kafka_handle, RD_KAFKA_PRODUCER, kafka->topic);
1254 if (kafka->mode == RD_KAFKA_CONSUMER) {
1255 char* topiclist = join(kafka->topics, ",");
1256 n_log(LOG_DEBUG, "starting polling thread for kafka handler %p mode CONSUMER (%d) topic %s", kafka->rd_kafka_handle, RD_KAFKA_CONSUMER, _str(topiclist));
1257 FreeNoLog(topiclist);
1258 }
1259
1260 start_HiTimer(&chrono);
1261
1262 int64_t remaining_time = (int64_t)kafka->poll_timeout * 1000;
1263 while (status == 1) {
1264 if (n_kafka_poll(kafka) == FALSE) {
1265 if (kafka->mode == RD_KAFKA_PRODUCER) {
1266 n_log(LOG_ERR, "failed to poll kafka producer handle %p with topic %s", kafka->rd_kafka_handle, rd_kafka_topic_name(kafka->rd_kafka_topic));
1267 } else if (kafka->topics) {
1268 char* topiclist = join(kafka->topics, ",");
1269 n_log(LOG_ERR, "failed to poll kafka consumer handle %p with topic %s", kafka->rd_kafka_handle, _str(topiclist));
1270 FreeNoLog(topiclist);
1271 }
1272 }
1273
1274 read_lock(kafka->rwlock);
1275 status = kafka->polling_thread_status;
1276 unlock(kafka->rwlock);
1277
1278 if (status == 2)
1279 break;
1280
1281 int64_t elapsed_time = get_usec(&chrono);
1282 if (kafka->poll_timeout != -1) {
1283 remaining_time -= elapsed_time;
1284 if (remaining_time < 0) {
1285 if (kafka->mode == RD_KAFKA_PRODUCER) {
1286 n_log(LOG_DEBUG, "timeouted on kafka handle %p", kafka->rd_kafka_handle);
1287 } else if (kafka->mode == RD_KAFKA_CONSUMER) {
1288 n_log(LOG_DEBUG, "timeouted on kafka handle %p", kafka->rd_kafka_handle);
1289 }
1290 break;
1291 }
1292 }
1293 // n_log( LOG_DEBUG , "remaining time: %d on kafka handle %p" , remaining_time , kafka -> rd_kafka_handle );
1294 }
1295
1296 write_lock(kafka->rwlock);
1297 kafka->polling_thread_status = 0;
1298 unlock(kafka->rwlock);
1299
1300 n_log(LOG_DEBUG, "exiting polling thread for kafka handler %p mode %s", kafka->rd_kafka_handle, (kafka->mode == RD_KAFKA_PRODUCER) ? "PRODUCER" : "CONSUMER");
1301 pthread_exit(NULL);
1302 return NULL;
1303} /* n_kafka_polling_thread */
1304
1311 __n_assert(kafka, return FALSE);
1312
1313 read_lock(kafka->rwlock);
1314 int status = kafka->polling_thread_status;
1315 unlock(kafka->rwlock);
1316
1317 if (status == 1) {
1318 n_log(LOG_ERR, "kafka polling thread already started for handle %p", kafka);
1319 return FALSE;
1320 }
1321
1322 write_lock(kafka->rwlock);
1323 kafka->polling_thread_status = 1;
1324
1325 if (pthread_create(&kafka->polling_thread, NULL, n_kafka_polling_thread, (void*)kafka) != 0) {
1326 n_log(LOG_ERR, "unable to create polling_thread for kafka handle %p", kafka);
1327 unlock(kafka->rwlock);
1328 return FALSE;
1329 }
1330 unlock(kafka->rwlock);
1331
1332 n_log(LOG_DEBUG, "pthread_create sucess for kafka handle %p->%p", kafka, kafka->rd_kafka_handle);
1333
1334 return TRUE;
1335} /* n_kafka_start_polling_thread */
1336
1343 __n_assert(kafka, return FALSE);
1344
1345 if (kafka->mode == RD_KAFKA_CONSUMER) {
1347 }
1348
1349 read_lock(kafka->rwlock);
1350 int polling_thread_status = kafka->polling_thread_status;
1351 unlock(kafka->rwlock);
1352
1353 if (polling_thread_status == 0) {
1354 n_log(LOG_DEBUG, "kafka polling thread already stopped for handle %p", kafka);
1355 return FALSE;
1356 }
1357 if (polling_thread_status == 2) {
1358 n_log(LOG_DEBUG, "kafka polling ask for stop thread already done for handle %p", kafka);
1359 return FALSE;
1360 }
1361
1362 write_lock(kafka->rwlock);
1363 kafka->polling_thread_status = 2;
1364 unlock(kafka->rwlock);
1365
1366 if (kafka->rd_kafka_handle && kafka->mode == RD_KAFKA_CONSUMER) {
1367 rd_kafka_consumer_close(kafka->rd_kafka_handle);
1368 }
1369
1370 struct timespec ts;
1371 clock_gettime(CLOCK_REALTIME, &ts);
1372 ts.tv_sec += 10;
1373
1374 int rc = pthread_timedjoin_np(kafka->polling_thread, NULL, &ts);
1375 if (rc != 0) {
1376 n_log(LOG_ERR, "polling thread did not stop in 10s, force stop !");
1377 return FALSE;
1378 }
1379
1380 return TRUE;
1381} /* n_kafka_stop_polling_thread */
1382
1389 __n_assert(kafka, return FALSE);
1390 write_lock(kafka->rwlock);
1391 kafka->event_consumption_enabled = TRUE;
1392 unlock(kafka->rwlock);
1393 return TRUE;
1394} /* n_kafka_set_consumer_event_polling */
1395
1402 __n_assert(kafka, return FALSE);
1403 write_lock(kafka->rwlock);
1404 kafka->event_consumption_enabled = FALSE;
1405 unlock(kafka->rwlock);
1406 return TRUE;
1407} /* n_kafka_set_consumer_event_polling */
1408
1415 __n_assert(kafka, return FALSE);
1416 write_lock(kafka->rwlock);
1417 kafka->event_production_enabled = TRUE;
1418 unlock(kafka->rwlock);
1419 return TRUE;
1420} /* n_kafka_set_consumer_event_polling */
1421
1428 __n_assert(kafka, return FALSE);
1429 write_lock(kafka->rwlock);
1430 kafka->event_production_enabled = FALSE;
1431 unlock(kafka->rwlock);
1432 return TRUE;
1433} /* n_kafka_set_consumer_event_polling */
1434
1441int n_kafka_dump_unprocessed(N_KAFKA* kafka, char* directory) {
1442 __n_assert(kafka, return FALSE);
1443 __n_assert(directory, return FALSE);
1444
1445 int status = 0;
1446 size_t nb_todump = 0;
1447 read_lock(kafka->rwlock);
1448 status = kafka->polling_thread_status;
1449 nb_todump = kafka->nb_queued + kafka->nb_waiting + kafka->nb_error;
1450 if (status != 0) {
1451 n_log(LOG_ERR, "kafka handle %p thread polling func is still running, aborting dump", kafka);
1452 unlock(kafka->rwlock);
1453 return FALSE;
1454 }
1455 if (nb_todump == 0) {
1456 n_log(LOG_DEBUG, "kafka handle %p: nothing to dump, all events processed correctly", kafka);
1457 unlock(kafka->rwlock);
1458 return TRUE;
1459 }
1460
1461 /* use a plain struct to alias event data without allocating a data buffer
1462 * that would leak when overwritten */
1463 N_STR* dumpstr = NULL;
1464 Malloc(dumpstr, N_STR, 1);
1465 __n_assert(dumpstr, unlock(kafka->rwlock); return FALSE);
1466 list_foreach(node, kafka->events_to_send) {
1467 N_KAFKA_EVENT* event = node->ptr;
1468 if (event->status != N_KAFKA_EVENT_OK) {
1469 size_t offset = 0;
1470 if (event->schema_id != -1)
1471 offset = 5;
1472
1473 N_STR* filename = NULL;
1474 nstrprintf(filename, "%s/%s+%p", directory, kafka->topic, event);
1475 n_log(LOG_DEBUG, "Dumping unprocessed event %p%s to %s", (void*)event, n_kafka_event_fdesc(event), _nstr(filename));
1476 // dump event here: alias the data pointer, do not own it
1477 dumpstr->data = event->event_string->data + offset;
1478 dumpstr->written = event->event_string->written - offset;
1479 dumpstr->length = event->event_string->length - offset;
1480 nstr_to_file(dumpstr, _nstr(filename));
1481 free_nstr(&filename);
1482 }
1483 }
1484 unlock(kafka->rwlock);
1485 /* clear aliased pointer before freeing the wrapper struct */
1486 dumpstr->data = NULL;
1487 free_nstr(&dumpstr);
1488 return TRUE;
1489} /* n_kafka_dump_unprocessed */
1490
1497int n_kafka_load_unprocessed(N_KAFKA* kafka, const char* directory) {
1498 __n_assert(kafka, return FALSE);
1499 __n_assert(directory, return FALSE);
1500
1501 write_lock(kafka->rwlock);
1502 /* load events from filename using base64decode( filename ) split( result , '+' )
1503 to get brokersname+topic to check against what's saved with the event
1504 and what's inside kafka's handle conf */
1505 unlock(kafka->rwlock);
1506
1507 return TRUE;
1508} /* n_kafka_load_unprocessed */
1509
1516 __n_assert(kafka, return NULL);
1517
1518 N_KAFKA_EVENT* event = NULL;
1519
1520 write_lock(kafka->rwlock);
1521 if (kafka->received_events->start)
1523 unlock(kafka->rwlock);
1524
1525 return event;
1526} /* n_kafka_get_event */
static int mode
char * event_string
Definition ex_kafka.c:47
char * config_file
Definition ex_kafka.c:46
char * key
#define init_lock(__rwlock_mutex)
Macro for initializing a rwlock.
Definition n_common.h:350
#define FreeNoLog(__ptr)
Free Handler without log.
Definition n_common.h:272
#define Malloc(__ptr, __struct, __size)
Malloc Handler to get errors and set to 0.
Definition n_common.h:204
#define __n_assert(__ptr, __ret)
macro to assert things
Definition n_common.h:279
#define _str(__PTR)
define true
Definition n_common.h:193
int get_computer_name(char *computer_name, size_t len)
get the computer name
Definition n_common.c:70
#define rw_lock_destroy(__rwlock_mutex)
Macro to destroy rwlock mutex.
Definition n_common.h:409
#define unlock(__rwlock_mutex)
Macro for releasing read/write lock a rwlock mutex.
Definition n_common.h:396
#define _nstrp(__PTR)
N_STR or NULL pointer for testing purposes.
Definition n_common.h:201
#define write_lock(__rwlock_mutex)
Macro for acquiring a write lock on a rwlock mutex.
Definition n_common.h:382
#define Free(__ptr)
Free Handler to get errors.
Definition n_common.h:263
#define read_lock(__rwlock_mutex)
Macro for acquiring a read lock on a rwlock mutex.
Definition n_common.h:368
#define _nstr(__PTR)
N_STR or "NULL" string for logging purposes.
Definition n_common.h:199
void * ptr
void pointer to store
Definition n_list.h:46
LIST_NODE * start
pointer to the start of the list
Definition n_list.h:66
struct LIST_NODE * next
pointer to the next node
Definition n_list.h:52
int list_push(LIST *list, void *ptr, void(*destructor)(void *ptr))
Add a pointer to the end of the list.
Definition n_list.c:228
#define list_foreach(__ITEM_, __LIST_)
ForEach macro helper, safe for node removal during iteration.
Definition n_list.h:89
#define remove_list_node(__LIST_, __NODE_, __TYPE_)
Remove macro helper for void pointer casting.
Definition n_list.h:98
int list_destroy(LIST **list)
Empty and Free a list container.
Definition n_list.c:548
LIST * new_generic_list(size_t max_items)
Initialiaze a generic list container to max_items pointers.
Definition n_list.c:37
#define MAX_LIST_ITEMS
flag to pass to new_generic_list for the maximum possible number of item in a list
Definition n_list.h:75
Structure of a generic list node.
Definition n_list.h:44
#define n_log(__LEVEL__,...)
Logging function wrapper to get line and func.
Definition n_log.h:89
#define LOG_DEBUG
debug-level messages
Definition n_log.h:84
#define LOG_ERR
error conditions
Definition n_log.h:76
#define LOG_INFO
informational
Definition n_log.h:82
int32_t event_production_enabled
bool flag to suspend or restart event production
Definition n_kafka.h:141
int32_t poll_interval
poll interval in msecs
Definition n_kafka.h:125
rd_kafka_topic_partition_list_t * subscription
subscribed topics
Definition n_kafka.h:101
int polling_thread_status
polling thread status, 0 => off , 1 => on , 2 => wants to stop, will be turned out to 0 by exiting po...
Definition n_kafka.h:129
cJSON * configuration
kafka json configuration holder
Definition n_kafka.h:121
size_t nb_waiting
number of events waiting for an ack in the waiting list
Definition n_kafka.h:135
int mode
kafka handle mode: RD_KAFKA_CONSUMER or RD_KAFKA_PRODUCER
Definition n_kafka.h:119
char ** topics
list of topics to subscribe to
Definition n_kafka.h:97
rd_kafka_topic_t * rd_kafka_topic
kafka topic handle
Definition n_kafka.h:111
rd_kafka_headers_t * rd_kafka_headers
kafka produce event headers structure handle
Definition n_kafka.h:77
int32_t error_timeout
retry interval in seconds for events stuck in N_KAFKA_EVENT_ERROR.
Definition n_kafka.h:143
int is_transactional
TRUE when a 'transactional.id' is set in the producer config: produces are then wrapped in kafka tran...
Definition n_kafka.h:145
pthread_rwlock_t rwlock
access lock
Definition n_kafka.h:117
N_STR * errstr
kafka error string holder
Definition n_kafka.h:113
rd_kafka_conf_t * rd_kafka_conf
kafka structure handle
Definition n_kafka.h:103
LIST * received_events
list of received N_KAFKA_EVENT
Definition n_kafka.h:93
LIST * events_to_send
list of N_KAFKA_EVENT to send
Definition n_kafka.h:91
int32_t poll_timeout
poll timeout in msecs
Definition n_kafka.h:123
pthread_t polling_thread
polling thread id
Definition n_kafka.h:131
char * sended_dir
base directory produced event files are moved to once acknowledged ('sended').
Definition n_kafka.h:149
rd_kafka_t * rd_kafka_handle
kafka handle (producer or consumer)
Definition n_kafka.h:105
char * event_cmd
eventual custom event_cmd
Definition n_kafka.h:99
int32_t monitored_directory_interval
monitored directory refresh interval in msecs
Definition n_kafka.h:127
char * tosend_dir
base directory the produced event files are read from ('to-send').
Definition n_kafka.h:147
char * groupid
consumer group id
Definition n_kafka.h:95
N_STR * event_files_to_delete
string containing the original event source file name if it is to be deleted when event is produced.
Definition n_kafka.h:71
int32_t event_consumption_enabled
bool flag to suspend or restart event consumption
Definition n_kafka.h:139
int schema_id
kafka schema id in network order
Definition n_kafka.h:115
size_t nb_queued
number of waiting events in the producer waiting list
Definition n_kafka.h:133
size_t nb_error
number of errored events
Definition n_kafka.h:137
char * bootstrap_servers
kafka bootstrap servers string
Definition n_kafka.h:109
char * topic
kafka topic string
Definition n_kafka.h:107
int n_kafka_add_header_ex(N_KAFKA_EVENT *event, char *key, size_t key_length, char *value, size_t value_length)
add a header entry to an event.
Definition n_kafka.c:820
int n_kafka_dump_unprocessed(N_KAFKA *kafka, char *directory)
dump unprocessed/unset events
Definition n_kafka.c:1441
N_KAFKA_EVENT * n_kafka_get_event(N_KAFKA *kafka)
get a received event from the N_KAFKA kafka handle
Definition n_kafka.c:1515
#define N_KAFKA_EVENT_OK
state of an OK event
Definition n_kafka.h:62
int n_kafka_stop_polling_thread(N_KAFKA *kafka)
stop the polling thread of a kafka handle
Definition n_kafka.c:1342
#define N_KAFKA_EVENT_ERROR
state of an errored event
Definition n_kafka.h:60
int32_t n_kafka_get_schema_from_nstr(const N_STR *string)
get a schema from the first 4 bytes of a N_STR *string, returning ntohl of the value
Definition n_kafka.c:165
int n_kafka_new_headers(N_KAFKA_EVENT *event, size_t count)
allocate a headers array for the event
Definition n_kafka.c:795
int n_kafka_put_schema_in_nstr(N_STR *string, int schema_id)
put a htonl schema id into the first 4 bytes of a N_STR *string
Definition n_kafka.c:191
int n_kafka_enable_event_production(N_KAFKA *kafka)
enable event production
Definition n_kafka.c:1414
int n_kafka_disable_event_consumption(N_KAFKA *kafka)
disable event consumption
Definition n_kafka.c:1401
N_KAFKA * n_kafka_load_config(char *config_file, int mode)
load a kafka configuration from a file
Definition n_kafka.c:463
int n_kafka_disable_event_production(N_KAFKA *kafka)
disable event production
Definition n_kafka.c:1427
int n_kafka_get_status(N_KAFKA *kafka, size_t *nb_queued, size_t *nb_waiting, size_t *nb_error)
return the queues status
Definition n_kafka.c:206
#define N_KAFKA_EVENT_WAITING_ACK
state of a sent event waiting for acknowledgement
Definition n_kafka.h:58
int n_kafka_put_schema_in_char(char *string, int schema_id)
put a htonl schema id into the first 4 bytes of a char *string
Definition n_kafka.c:178
int n_kafka_poll(N_KAFKA *kafka)
Poll kafka handle in producer or consumer mode.
Definition n_kafka.c:1055
int n_kafka_start_polling_thread(N_KAFKA *kafka)
start the polling thread of a kafka handle
Definition n_kafka.c:1310
N_KAFKA_EVENT * n_kafka_new_event_from_char(const char *string, size_t written, int schema_id)
make a new event from a char *string
Definition n_kafka.c:946
N_KAFKA_EVENT * n_kafka_new_event_from_string(const N_STR *string, int schema_id)
make a new event from a N_STR *string
Definition n_kafka.c:984
int n_kafka_event_destroy(N_KAFKA_EVENT **event)
destroy a kafka event and set it's pointer to NULL
Definition n_kafka.c:1043
int n_kafka_enable_event_consumption(N_KAFKA *kafka)
enable event consumption
Definition n_kafka.c:1388
void n_kafka_event_destroy_ptr(void *event_ptr)
festroy a kafka event
Definition n_kafka.c:1015
int32_t n_kafka_get_schema_from_char(const char *string)
get a schema from the first 4 bytes of a char *string, returning ntohl of the value
Definition n_kafka.c:151
int n_kafka_add_header(N_KAFKA_EVENT *event, N_STR *key, N_STR *value)
add a header entry to an event.
Definition n_kafka.c:854
int n_kafka_set_sended_dir(N_KAFKA *kafka, const char *tosend_dir, const char *sended_dir)
set the 'to-send' and 'sended' base directories so acknowledged event files are moved instead of dele...
Definition n_kafka.c:244
int n_kafka_load_unprocessed(N_KAFKA *kafka, const char *directory)
load unprocessed/unset events
Definition n_kafka.c:1497
int n_kafka_produce(N_KAFKA *kafka, N_KAFKA_EVENT *event)
put an event in the events_to_send list
Definition n_kafka.c:870
N_KAFKA_EVENT * n_kafka_new_event(int schema_id)
create a new empty event
Definition n_kafka.c:771
int n_kafka_set_error_timeout(N_KAFKA *kafka, int32_t error_timeout)
set the retry interval (in seconds) for events stuck in N_KAFKA_EVENT_ERROR
Definition n_kafka.c:225
#define N_KAFKA_EVENT_CREATED
state of a freshly created event
Definition n_kafka.h:64
N_KAFKA_EVENT * n_kafka_new_event_from_file(char *filename, int schema_id)
make a new event from a N_STR *string
Definition n_kafka.c:999
N_KAFKA * n_kafka_new(int32_t poll_timeout, int32_t poll_interval, size_t errstr_len)
allocate a new kafka handle
Definition n_kafka.c:407
#define N_KAFKA_EVENT_QUEUED
state of a queued event
Definition n_kafka.h:56
void n_kafka_delete(N_KAFKA *kafka)
delete a N_KAFKA handle
Definition n_kafka.c:344
structure of a KAFKA consumer or producer handle
Definition n_kafka.h:89
structure of a KAFKA message
Definition n_kafka.h:67
size_t written
number of meaningful bytes in data, excluding the null terminator; the size including the null termin...
Definition n_str.h:68
char * data
the string
Definition n_str.h:63
size_t length
total allocation (in bytes) of the data buffer, padding included
Definition n_str.h:65
void free_nstr_ptr(void *ptr)
Free a N_STR pointer structure.
Definition n_str.c:70
#define free_nstr(__ptr)
free a N_STR structure and set the pointer to NULL
Definition n_str.h:203
int split_count(char **split_result)
Count split elements.
Definition n_str.c:1000
int nstr_to_file(N_STR *str, char *filename)
Write a N_STR content into a file.
Definition n_str.c:424
char * join(char **splitresult, const char *delim)
join the array into a string
Definition n_str.c:1037
N_STR * new_nstr(NSTRBYTE size)
create a new N_STR string
Definition n_str.c:207
#define nstrprintf(__nstr_var, __format,...)
Macro to quickly allocate and sprintf to N_STR.
Definition n_str.h:117
char ** split(const char *str, const char *delim, int empty)
split the strings into a an array of char *pointer , ended by a NULL one.
Definition n_str.c:920
N_STR * file_to_nstr(char *filename)
Load a whole file into a N_STR.
Definition n_str.c:288
int free_split_result(char ***tab)
Free a split result allocated array.
Definition n_str.c:1016
A box including a string and his lenght.
Definition n_str.h:61
int start_HiTimer(N_TIME *timer)
Initialize or restart from zero any N_TIME HiTimer.
Definition n_time.c:85
time_t get_usec(N_TIME *timer)
Poll any N_TIME HiTimer, returning usec, and moving currentTime to startTime.
Definition n_time.c:107
Timing Structure.
Definition n_time.h:49
Base64 encoding and decoding functions using N_STR.
Common headers and low-level functions & define.
static int n_kafka_move_file_to_sended(const N_KAFKA *kafka, const char *filepath)
move a single acknowledged event file to the 'sended' directory, preserving its subpath relative to '...
Definition n_kafka.c:100
static void n_kafka_delivery_message_callback(rd_kafka_t *rk, const rd_kafka_message_t *rkmessage, void *opaque)
Message delivery report callback.
Definition n_kafka.c:267
static int n_kafka_mkdir_p(const char *path)
create a directory and all its missing parents (mkdir -p)
Definition n_kafka.c:58
void * n_kafka_polling_thread(void *ptr)
kafka produce or consume polling thread function
Definition n_kafka.c:1245
int n_kafka_produce_ex(N_KAFKA *kafka, N_KAFKA_EVENT *event)
produce an event on a N_KAFKA *kafka handle
Definition n_kafka.c:893
static const char * n_kafka_event_fdesc(const N_KAFKA_EVENT *event)
return a " (file: ...)" log fragment naming the source file(s) of an event
Definition n_kafka.c:44
Kafka generic produce and consume event header.
#define _Thread_local
thread-local pre-connection error buffer (DNS, socket creation)
Definition n_network.c:75