http://git-wip-us.apache.org/repos/asf/nifi-minifi-cpp/blob/7528d23e/thirdparty/librdkafka-0.11.1/src/rdkafka_queue.c ---------------------------------------------------------------------- diff --git a/thirdparty/librdkafka-0.11.1/src/rdkafka_queue.c b/thirdparty/librdkafka-0.11.1/src/rdkafka_queue.c deleted file mode 100644 index fab2899..0000000 --- a/thirdparty/librdkafka-0.11.1/src/rdkafka_queue.c +++ /dev/null @@ -1,860 +0,0 @@ -/* - * librdkafka - The Apache Kafka C/C++ library - * - * Copyright (c) 2016 Magnus Edenhill - * All rights reserved. - * - * Redistribution and use in source and binary forms, with or without - * modification, are permitted provided that the following conditions are met: - * - * 1. Redistributions of source code must retain the above copyright notice, - * this list of conditions and the following disclaimer. - * 2. Redistributions in binary form must reproduce the above copyright notice, - * this list of conditions and the following disclaimer in the documentation - * and/or other materials provided with the distribution. - * - * THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS" - * AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE - * IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE - * ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT OWNER OR CONTRIBUTORS BE - * LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR - * CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF - * SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS - * INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN - * CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) - * ARISING IN ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE - * POSSIBILITY OF SUCH DAMAGE. - */ - -#include "rdkafka_int.h" -#include "rdkafka_offset.h" -#include "rdkafka_topic.h" -#include "rdkafka_interceptor.h" - -int RD_TLS rd_kafka_yield_thread = 0; - -void rd_kafka_yield (rd_kafka_t *rk) { - rd_kafka_yield_thread = 1; -} - - -/** - * Destroy a queue. refcnt must be at zero. - */ -void rd_kafka_q_destroy_final (rd_kafka_q_t *rkq) { - - mtx_lock(&rkq->rkq_lock); - if (unlikely(rkq->rkq_qio != NULL)) { - rd_free(rkq->rkq_qio); - rkq->rkq_qio = NULL; - } - rd_kafka_q_fwd_set0(rkq, NULL, 0/*no-lock*/, 0 /*no-fwd-app*/); - rd_kafka_q_disable0(rkq, 0/*no-lock*/); - rd_kafka_q_purge0(rkq, 0/*no-lock*/); - assert(!rkq->rkq_fwdq); - mtx_unlock(&rkq->rkq_lock); - mtx_destroy(&rkq->rkq_lock); - cnd_destroy(&rkq->rkq_cond); - - if (rkq->rkq_flags & RD_KAFKA_Q_F_ALLOCATED) - rd_free(rkq); -} - - - -/** - * Initialize a queue. - */ -void rd_kafka_q_init (rd_kafka_q_t *rkq, rd_kafka_t *rk) { - rd_kafka_q_reset(rkq); - rkq->rkq_fwdq = NULL; - rkq->rkq_refcnt = 1; - rkq->rkq_flags = RD_KAFKA_Q_F_READY; - rkq->rkq_rk = rk; - rkq->rkq_qio = NULL; - rkq->rkq_serve = NULL; - rkq->rkq_opaque = NULL; - mtx_init(&rkq->rkq_lock, mtx_plain); - cnd_init(&rkq->rkq_cond); -} - - -/** - * Allocate a new queue and initialize it. - */ -rd_kafka_q_t *rd_kafka_q_new0 (rd_kafka_t *rk, const char *func, int line) { - rd_kafka_q_t *rkq = rd_malloc(sizeof(*rkq)); - rd_kafka_q_init(rkq, rk); - rkq->rkq_flags |= RD_KAFKA_Q_F_ALLOCATED; -#if ENABLE_DEVEL - rd_snprintf(rkq->rkq_name, sizeof(rkq->rkq_name), "%s:%d", func, line); -#else - rkq->rkq_name = func; -#endif - return rkq; -} - -/** - * Set/clear forward queue. - * Queue forwarding enables message routing inside rdkafka. - * Typical use is to re-route all fetched messages for all partitions - * to one single queue. - * - * All access to rkq_fwdq are protected by rkq_lock. - */ -void rd_kafka_q_fwd_set0 (rd_kafka_q_t *srcq, rd_kafka_q_t *destq, - int do_lock, int fwd_app) { - - if (do_lock) - mtx_lock(&srcq->rkq_lock); - if (fwd_app) - srcq->rkq_flags |= RD_KAFKA_Q_F_FWD_APP; - if (srcq->rkq_fwdq) { - rd_kafka_q_destroy(srcq->rkq_fwdq); - srcq->rkq_fwdq = NULL; - } - if (destq) { - rd_kafka_q_keep(destq); - - /* If rkq has ops in queue, append them to fwdq's queue. - * This is an irreversible operation. */ - if (srcq->rkq_qlen > 0) { - rd_dassert(destq->rkq_flags & RD_KAFKA_Q_F_READY); - rd_kafka_q_concat(destq, srcq); - } - - srcq->rkq_fwdq = destq; - } - if (do_lock) - mtx_unlock(&srcq->rkq_lock); -} - -/** - * Purge all entries from a queue. - */ -int rd_kafka_q_purge0 (rd_kafka_q_t *rkq, int do_lock) { - rd_kafka_op_t *rko, *next; - TAILQ_HEAD(, rd_kafka_op_s) tmpq = TAILQ_HEAD_INITIALIZER(tmpq); - rd_kafka_q_t *fwdq; - int cnt = 0; - - if (do_lock) - mtx_lock(&rkq->rkq_lock); - - if ((fwdq = rd_kafka_q_fwd_get(rkq, 0))) { - if (do_lock) - mtx_unlock(&rkq->rkq_lock); - cnt = rd_kafka_q_purge(fwdq); - rd_kafka_q_destroy(fwdq); - return cnt; - } - - /* Move ops queue to tmpq to avoid lock-order issue - * by locks taken from rd_kafka_op_destroy(). */ - TAILQ_MOVE(&tmpq, &rkq->rkq_q, rko_link); - - /* Zero out queue */ - rd_kafka_q_reset(rkq); - - if (do_lock) - mtx_unlock(&rkq->rkq_lock); - - /* Destroy the ops */ - next = TAILQ_FIRST(&tmpq); - while ((rko = next)) { - next = TAILQ_NEXT(next, rko_link); - rd_kafka_op_destroy(rko); - cnt++; - } - - return cnt; -} - - -/** - * Purge all entries from a queue with a rktp version smaller than `version` - * This shaves off the head of the queue, up until the first rko with - * a non-matching rktp or version. - */ -void rd_kafka_q_purge_toppar_version (rd_kafka_q_t *rkq, - rd_kafka_toppar_t *rktp, int version) { - rd_kafka_op_t *rko, *next; - TAILQ_HEAD(, rd_kafka_op_s) tmpq = TAILQ_HEAD_INITIALIZER(tmpq); - int32_t cnt = 0; - int64_t size = 0; - rd_kafka_q_t *fwdq; - - mtx_lock(&rkq->rkq_lock); - - if ((fwdq = rd_kafka_q_fwd_get(rkq, 0))) { - mtx_unlock(&rkq->rkq_lock); - rd_kafka_q_purge_toppar_version(fwdq, rktp, version); - rd_kafka_q_destroy(fwdq); - return; - } - - /* Move ops to temporary queue and then destroy them from there - * without locks to avoid lock-ordering problems in op_destroy() */ - while ((rko = TAILQ_FIRST(&rkq->rkq_q)) && rko->rko_rktp && - rd_kafka_toppar_s2i(rko->rko_rktp) == rktp && - rko->rko_version < version) { - TAILQ_REMOVE(&rkq->rkq_q, rko, rko_link); - TAILQ_INSERT_TAIL(&tmpq, rko, rko_link); - cnt++; - size += rko->rko_len; - } - - - rkq->rkq_qlen -= cnt; - rkq->rkq_qsize -= size; - mtx_unlock(&rkq->rkq_lock); - - next = TAILQ_FIRST(&tmpq); - while ((rko = next)) { - next = TAILQ_NEXT(next, rko_link); - rd_kafka_op_destroy(rko); - } -} - - -/** - * Move 'cnt' entries from 'srcq' to 'dstq'. - * If 'cnt' == -1 all entries will be moved. - * Returns the number of entries moved. - */ -int rd_kafka_q_move_cnt (rd_kafka_q_t *dstq, rd_kafka_q_t *srcq, - int cnt, int do_locks) { - rd_kafka_op_t *rko; - int mcnt = 0; - - if (do_locks) { - mtx_lock(&srcq->rkq_lock); - mtx_lock(&dstq->rkq_lock); - } - - if (!dstq->rkq_fwdq && !srcq->rkq_fwdq) { - if (cnt > 0 && dstq->rkq_qlen == 0) - rd_kafka_q_io_event(dstq); - - /* Optimization, if 'cnt' is equal/larger than all - * items of 'srcq' we can move the entire queue. */ - if (cnt == -1 || - cnt >= (int)srcq->rkq_qlen) { - mcnt = srcq->rkq_qlen; - rd_kafka_q_concat0(dstq, srcq, 0/*no-lock*/); - } else { - while (mcnt < cnt && - (rko = TAILQ_FIRST(&srcq->rkq_q))) { - TAILQ_REMOVE(&srcq->rkq_q, rko, rko_link); - if (likely(!rko->rko_prio)) - TAILQ_INSERT_TAIL(&dstq->rkq_q, rko, - rko_link); - else - TAILQ_INSERT_SORTED( - &dstq->rkq_q, rko, - rd_kafka_op_t *, rko_link, - rd_kafka_op_cmp_prio); - - srcq->rkq_qlen--; - dstq->rkq_qlen++; - srcq->rkq_qsize -= rko->rko_len; - dstq->rkq_qsize += rko->rko_len; - mcnt++; - } - } - } else - mcnt = rd_kafka_q_move_cnt(dstq->rkq_fwdq ? dstq->rkq_fwdq:dstq, - srcq->rkq_fwdq ? srcq->rkq_fwdq:srcq, - cnt, do_locks); - - if (do_locks) { - mtx_unlock(&dstq->rkq_lock); - mtx_unlock(&srcq->rkq_lock); - } - - return mcnt; -} - - -/** - * Filters out outdated ops. - */ -static RD_INLINE rd_kafka_op_t *rd_kafka_op_filter (rd_kafka_q_t *rkq, - rd_kafka_op_t *rko, - int version) { - if (unlikely(!rko)) - return NULL; - - if (unlikely(rd_kafka_op_version_outdated(rko, version))) { - rd_kafka_q_deq0(rkq, rko); - rd_kafka_op_destroy(rko); - return NULL; - } - - return rko; -} - - - -/** - * Pop an op from a queue. - * - * Locality: any thread. - */ - - -/** - * Serve q like rd_kafka_q_serve() until an op is found that can be returned - * as an event to the application. - * - * @returns the first event:able op, or NULL on timeout. - * - * Locality: any thread - */ -rd_kafka_op_t *rd_kafka_q_pop_serve (rd_kafka_q_t *rkq, int timeout_ms, - int32_t version, - rd_kafka_q_cb_type_t cb_type, - rd_kafka_q_serve_cb_t *callback, - void *opaque) { - rd_kafka_op_t *rko; - rd_kafka_q_t *fwdq; - - rd_dassert(cb_type); - - if (timeout_ms == RD_POLL_INFINITE) - timeout_ms = INT_MAX; - - mtx_lock(&rkq->rkq_lock); - - rd_kafka_yield_thread = 0; - if (!(fwdq = rd_kafka_q_fwd_get(rkq, 0))) { - do { - rd_kafka_op_res_t res; - rd_ts_t pre; - - /* Filter out outdated ops */ - retry: - while ((rko = TAILQ_FIRST(&rkq->rkq_q)) && - !(rko = rd_kafka_op_filter(rkq, rko, version))) - ; - - if (rko) { - /* Proper versioned op */ - rd_kafka_q_deq0(rkq, rko); - - /* Ops with callbacks are considered handled - * and we move on to the next op, if any. - * Ops w/o callbacks are returned immediately */ - res = rd_kafka_op_handle(rkq->rkq_rk, rkq, rko, - cb_type, opaque, - callback); - if (res == RD_KAFKA_OP_RES_HANDLED) - goto retry; /* Next op */ - else if (unlikely(res == - RD_KAFKA_OP_RES_YIELD)) { - /* Callback yielded, unroll */ - mtx_unlock(&rkq->rkq_lock); - return NULL; - } else - break; /* Proper op, handle below. */ - } - - /* No op, wait for one */ - pre = rd_clock(); - if (cnd_timedwait_ms(&rkq->rkq_cond, - &rkq->rkq_lock, - timeout_ms) == - thrd_timedout) { - mtx_unlock(&rkq->rkq_lock); - return NULL; - } - /* Remove spent time */ - timeout_ms -= (int) (rd_clock()-pre) / 1000; - if (timeout_ms < 0) - timeout_ms = RD_POLL_NOWAIT; - - } while (timeout_ms != RD_POLL_NOWAIT); - - mtx_unlock(&rkq->rkq_lock); - - } else { - /* Since the q_pop may block we need to release the parent - * queue's lock. */ - mtx_unlock(&rkq->rkq_lock); - rko = rd_kafka_q_pop_serve(fwdq, timeout_ms, version, - cb_type, callback, opaque); - rd_kafka_q_destroy(fwdq); - } - - - return rko; -} - -rd_kafka_op_t *rd_kafka_q_pop (rd_kafka_q_t *rkq, int timeout_ms, - int32_t version) { - return rd_kafka_q_pop_serve(rkq, timeout_ms, version, - RD_KAFKA_Q_CB_RETURN, - NULL, NULL); -} - - -/** - * Pop all available ops from a queue and call the provided - * callback for each op. - * `max_cnt` limits the number of ops served, 0 = no limit. - * - * Returns the number of ops served. - * - * Locality: any thread. - */ -int rd_kafka_q_serve (rd_kafka_q_t *rkq, int timeout_ms, - int max_cnt, rd_kafka_q_cb_type_t cb_type, - rd_kafka_q_serve_cb_t *callback, void *opaque) { - rd_kafka_t *rk = rkq->rkq_rk; - rd_kafka_op_t *rko; - rd_kafka_q_t localq; - rd_kafka_q_t *fwdq; - int cnt = 0; - - rd_dassert(cb_type); - - mtx_lock(&rkq->rkq_lock); - - rd_dassert(TAILQ_EMPTY(&rkq->rkq_q) || rkq->rkq_qlen > 0); - if ((fwdq = rd_kafka_q_fwd_get(rkq, 0))) { - int ret; - /* Since the q_pop may block we need to release the parent - * queue's lock. */ - mtx_unlock(&rkq->rkq_lock); - ret = rd_kafka_q_serve(fwdq, timeout_ms, max_cnt, - cb_type, callback, opaque); - rd_kafka_q_destroy(fwdq); - return ret; - } - - if (timeout_ms == RD_POLL_INFINITE) - timeout_ms = INT_MAX; - - /* Wait for op */ - while (!(rko = TAILQ_FIRST(&rkq->rkq_q)) && timeout_ms != 0) { - if (cnd_timedwait_ms(&rkq->rkq_cond, - &rkq->rkq_lock, - timeout_ms) != thrd_success) - break; - - timeout_ms = 0; - } - - if (!rko) { - mtx_unlock(&rkq->rkq_lock); - return 0; - } - - /* Move the first `max_cnt` ops. */ - rd_kafka_q_init(&localq, rkq->rkq_rk); - rd_kafka_q_move_cnt(&localq, rkq, max_cnt == 0 ? -1/*all*/ : max_cnt, - 0/*no-locks*/); - - mtx_unlock(&rkq->rkq_lock); - - rd_kafka_yield_thread = 0; - - /* Call callback for each op */ - while ((rko = TAILQ_FIRST(&localq.rkq_q))) { - rd_kafka_op_res_t res; - - rd_kafka_q_deq0(&localq, rko); - res = rd_kafka_op_handle(rk, &localq, rko, cb_type, - opaque, callback); - /* op must have been handled */ - rd_kafka_assert(NULL, res != RD_KAFKA_OP_RES_PASS); - cnt++; - - if (unlikely(res == RD_KAFKA_OP_RES_YIELD || - rd_kafka_yield_thread)) { - /* Callback called rd_kafka_yield(), we must - * stop our callback dispatching and put the - * ops in localq back on the original queue head. */ - if (!TAILQ_EMPTY(&localq.rkq_q)) - rd_kafka_q_prepend(rkq, &localq); - break; - } - } - - rd_kafka_q_destroy(&localq); - - return cnt; -} - - - - - -/** - * Populate 'rkmessages' array with messages from 'rkq'. - * If 'auto_commit' is set, each message's offset will be committed - * to the offset store for that toppar. - * - * Returns the number of messages added. - */ - -int rd_kafka_q_serve_rkmessages (rd_kafka_q_t *rkq, int timeout_ms, - rd_kafka_message_t **rkmessages, - size_t rkmessages_size) { - unsigned int cnt = 0; - TAILQ_HEAD(, rd_kafka_op_s) tmpq = TAILQ_HEAD_INITIALIZER(tmpq); - rd_kafka_op_t *rko, *next; - rd_kafka_t *rk = rkq->rkq_rk; - rd_kafka_q_t *fwdq; - - mtx_lock(&rkq->rkq_lock); - if ((fwdq = rd_kafka_q_fwd_get(rkq, 0))) { - /* Since the q_pop may block we need to release the parent - * queue's lock. */ - mtx_unlock(&rkq->rkq_lock); - cnt = rd_kafka_q_serve_rkmessages(fwdq, timeout_ms, - rkmessages, rkmessages_size); - rd_kafka_q_destroy(fwdq); - return cnt; - } - mtx_unlock(&rkq->rkq_lock); - - rd_kafka_yield_thread = 0; - while (cnt < rkmessages_size) { - rd_kafka_op_res_t res; - - mtx_lock(&rkq->rkq_lock); - - while (!(rko = TAILQ_FIRST(&rkq->rkq_q))) { - if (cnd_timedwait_ms(&rkq->rkq_cond, &rkq->rkq_lock, - timeout_ms) == thrd_timedout) - break; - } - - if (!rko) { - mtx_unlock(&rkq->rkq_lock); - break; /* Timed out */ - } - - rd_kafka_q_deq0(rkq, rko); - - mtx_unlock(&rkq->rkq_lock); - - if (rd_kafka_op_version_outdated(rko, 0)) { - /* Outdated op, put on discard queue */ - TAILQ_INSERT_TAIL(&tmpq, rko, rko_link); - continue; - } - - /* Serve non-FETCH callbacks */ - res = rd_kafka_poll_cb(rk, rkq, rko, - RD_KAFKA_Q_CB_RETURN, NULL); - if (res == RD_KAFKA_OP_RES_HANDLED) { - /* Callback served, rko is destroyed. */ - continue; - } else if (unlikely(res == RD_KAFKA_OP_RES_YIELD || - rd_kafka_yield_thread)) { - /* Yield. */ - break; - } - rd_dassert(res == RD_KAFKA_OP_RES_PASS); - - /* Auto-commit offset, if enabled. */ - if (!rko->rko_err && rko->rko_type == RD_KAFKA_OP_FETCH) { - rd_kafka_toppar_t *rktp; - rktp = rd_kafka_toppar_s2i(rko->rko_rktp); - rd_kafka_toppar_lock(rktp); - rktp->rktp_app_offset = rko->rko_u.fetch.rkm.rkm_offset+1; - if (rktp->rktp_cgrp && - rk->rk_conf.enable_auto_offset_store) - rd_kafka_offset_store0(rktp, - rktp->rktp_app_offset, - 0/* no lock */); - rd_kafka_toppar_unlock(rktp); - } - - /* Get rkmessage from rko and append to array. */ - rkmessages[cnt++] = rd_kafka_message_get(rko); - } - - /* Discard non-desired and already handled ops */ - next = TAILQ_FIRST(&tmpq); - while (next) { - rko = next; - next = TAILQ_NEXT(next, rko_link); - rd_kafka_op_destroy(rko); - } - - - return cnt; -} - - - -void rd_kafka_queue_destroy (rd_kafka_queue_t *rkqu) { - rd_kafka_q_disable(rkqu->rkqu_q); - rd_kafka_q_destroy(rkqu->rkqu_q); - rd_free(rkqu); -} - -rd_kafka_queue_t *rd_kafka_queue_new0 (rd_kafka_t *rk, rd_kafka_q_t *rkq) { - rd_kafka_queue_t *rkqu; - - rkqu = rd_calloc(1, sizeof(*rkqu)); - - rkqu->rkqu_q = rkq; - rd_kafka_q_keep(rkq); - - rkqu->rkqu_rk = rk; - - return rkqu; -} - - -rd_kafka_queue_t *rd_kafka_queue_new (rd_kafka_t *rk) { - rd_kafka_q_t *rkq; - rd_kafka_queue_t *rkqu; - - rkq = rd_kafka_q_new(rk); - rkqu = rd_kafka_queue_new0(rk, rkq); - rd_kafka_q_destroy(rkq); /* Loose refcount from q_new, one is held - * by queue_new0 */ - return rkqu; -} - - -rd_kafka_queue_t *rd_kafka_queue_get_main (rd_kafka_t *rk) { - return rd_kafka_queue_new0(rk, rk->rk_rep); -} - - -rd_kafka_queue_t *rd_kafka_queue_get_consumer (rd_kafka_t *rk) { - if (!rk->rk_cgrp) - return NULL; - return rd_kafka_queue_new0(rk, rk->rk_cgrp->rkcg_q); -} - -rd_kafka_queue_t *rd_kafka_queue_get_partition (rd_kafka_t *rk, - const char *topic, - int32_t partition) { - shptr_rd_kafka_toppar_t *s_rktp; - rd_kafka_toppar_t *rktp; - rd_kafka_queue_t *result; - - if (rk->rk_type == RD_KAFKA_PRODUCER) - return NULL; - - s_rktp = rd_kafka_toppar_get2(rk, topic, - partition, - 0, /* no ua_on_miss */ - 1 /* create_on_miss */); - - if (!s_rktp) - return NULL; - - rktp = rd_kafka_toppar_s2i(s_rktp); - result = rd_kafka_queue_new0(rk, rktp->rktp_fetchq); - rd_kafka_toppar_destroy(s_rktp); - - return result; -} - -rd_kafka_resp_err_t rd_kafka_set_log_queue (rd_kafka_t *rk, - rd_kafka_queue_t *rkqu) { - rd_kafka_q_t *rkq; - if (!rkqu) - rkq = rk->rk_rep; - else - rkq = rkqu->rkqu_q; - rd_kafka_q_fwd_set(rk->rk_logq, rkq); - return RD_KAFKA_RESP_ERR_NO_ERROR; -} - -void rd_kafka_queue_forward (rd_kafka_queue_t *src, rd_kafka_queue_t *dst) { - rd_kafka_q_fwd_set0(src->rkqu_q, dst ? dst->rkqu_q : NULL, - 1, /* do_lock */ - 1 /* fwd_app */); -} - - -size_t rd_kafka_queue_length (rd_kafka_queue_t *rkqu) { - return (size_t)rd_kafka_q_len(rkqu->rkqu_q); -} - -/** - * @brief Enable or disable(fd==-1) fd-based wake-ups for queue - */ -void rd_kafka_q_io_event_enable (rd_kafka_q_t *rkq, int fd, - const void *payload, size_t size) { - struct rd_kafka_q_io *qio = NULL; - - if (fd != -1) { - qio = rd_malloc(sizeof(*qio) + size); - qio->fd = fd; - qio->size = size; - qio->payload = (void *)(qio+1); - memcpy(qio->payload, payload, size); - } - - mtx_lock(&rkq->rkq_lock); - if (rkq->rkq_qio) { - rd_free(rkq->rkq_qio); - rkq->rkq_qio = NULL; - } - - if (fd != -1) { - rkq->rkq_qio = qio; - } - - mtx_unlock(&rkq->rkq_lock); - -} - -void rd_kafka_queue_io_event_enable (rd_kafka_queue_t *rkqu, int fd, - const void *payload, size_t size) { - rd_kafka_q_io_event_enable(rkqu->rkqu_q, fd, payload, size); -} - - -/** - * Helper: wait for single op on 'rkq', and return its error, - * or .._TIMED_OUT on timeout. - */ -rd_kafka_resp_err_t rd_kafka_q_wait_result (rd_kafka_q_t *rkq, int timeout_ms) { - rd_kafka_op_t *rko; - rd_kafka_resp_err_t err; - - rko = rd_kafka_q_pop(rkq, timeout_ms, 0); - if (!rko) - err = RD_KAFKA_RESP_ERR__TIMED_OUT; - else { - err = rko->rko_err; - rd_kafka_op_destroy(rko); - } - - return err; -} - - -/** - * Apply \p callback on each op in queue. - * If the callback wishes to remove the rko it must do so using - * using rd_kafka_op_deq0(). - * - * @returns the sum of \p callback() return values. - * @remark rkq will be locked, callers should take care not to - * interact with \p rkq through other means from the callback to avoid - * deadlocks. - */ -int rd_kafka_q_apply (rd_kafka_q_t *rkq, - int (*callback) (rd_kafka_q_t *rkq, rd_kafka_op_t *rko, - void *opaque), - void *opaque) { - rd_kafka_op_t *rko, *next; - rd_kafka_q_t *fwdq; - int cnt = 0; - - mtx_lock(&rkq->rkq_lock); - if ((fwdq = rd_kafka_q_fwd_get(rkq, 0))) { - mtx_unlock(&rkq->rkq_lock); - cnt = rd_kafka_q_apply(fwdq, callback, opaque); - rd_kafka_q_destroy(fwdq); - return cnt; - } - - next = TAILQ_FIRST(&rkq->rkq_q); - while ((rko = next)) { - next = TAILQ_NEXT(next, rko_link); - cnt += callback(rkq, rko, opaque); - } - mtx_unlock(&rkq->rkq_lock); - - return cnt; -} - -/** - * @brief Convert relative to absolute offsets and also purge any messages - * that are older than \p min_offset. - * @remark Error ops with ERR__NOT_IMPLEMENTED will not be purged since - * they are used to indicate unknnown compression codecs and compressed - * messagesets may have a starting offset lower than what we requested. - * @remark \p rkq locking is not performed (caller's responsibility) - * @remark Must NOT be used on fwdq. - */ -void rd_kafka_q_fix_offsets (rd_kafka_q_t *rkq, int64_t min_offset, - int64_t base_offset) { - rd_kafka_op_t *rko, *next; - int adj_len = 0; - int64_t adj_size = 0; - - rd_kafka_assert(NULL, !rkq->rkq_fwdq); - - next = TAILQ_FIRST(&rkq->rkq_q); - while ((rko = next)) { - next = TAILQ_NEXT(next, rko_link); - - if (unlikely(rko->rko_type != RD_KAFKA_OP_FETCH)) - continue; - - rko->rko_u.fetch.rkm.rkm_offset += base_offset; - - if (rko->rko_u.fetch.rkm.rkm_offset < min_offset && - rko->rko_err != RD_KAFKA_RESP_ERR__NOT_IMPLEMENTED) { - adj_len++; - adj_size += rko->rko_len; - TAILQ_REMOVE(&rkq->rkq_q, rko, rko_link); - rd_kafka_op_destroy(rko); - continue; - } - } - - - rkq->rkq_qlen -= adj_len; - rkq->rkq_qsize -= adj_size; -} - - -/** - * @brief Print information and contents of queue - */ -void rd_kafka_q_dump (FILE *fp, rd_kafka_q_t *rkq) { - mtx_lock(&rkq->rkq_lock); - fprintf(fp, "Queue %p \"%s\" (refcnt %d, flags 0x%x, %d ops, " - "%"PRId64" bytes)\n", - rkq, rkq->rkq_name, rkq->rkq_refcnt, rkq->rkq_flags, - rkq->rkq_qlen, rkq->rkq_qsize); - - if (rkq->rkq_qio) - fprintf(fp, " QIO fd %d\n", rkq->rkq_qio->fd); - if (rkq->rkq_serve) - fprintf(fp, " Serve callback %p, opaque %p\n", - rkq->rkq_serve, rkq->rkq_opaque); - - if (rkq->rkq_fwdq) { - fprintf(fp, " Forwarded ->\n"); - rd_kafka_q_dump(fp, rkq->rkq_fwdq); - } else { - rd_kafka_op_t *rko; - - if (!TAILQ_EMPTY(&rkq->rkq_q)) - fprintf(fp, " Queued ops:\n"); - TAILQ_FOREACH(rko, &rkq->rkq_q, rko_link) { - fprintf(fp, " %p %s (v%"PRId32", flags 0x%x, " - "prio %d, len %"PRId32", source %s, " - "replyq %p)\n", - rko, rd_kafka_op2str(rko->rko_type), - rko->rko_version, rko->rko_flags, - rko->rko_prio, rko->rko_len, - #if ENABLE_DEVEL - rko->rko_source - #else - "-" - #endif - , - rko->rko_replyq.q - ); - } - } - - mtx_unlock(&rkq->rkq_lock); -}
http://git-wip-us.apache.org/repos/asf/nifi-minifi-cpp/blob/7528d23e/thirdparty/librdkafka-0.11.1/src/rdkafka_queue.h ---------------------------------------------------------------------- diff --git a/thirdparty/librdkafka-0.11.1/src/rdkafka_queue.h b/thirdparty/librdkafka-0.11.1/src/rdkafka_queue.h deleted file mode 100644 index 7595b17..0000000 --- a/thirdparty/librdkafka-0.11.1/src/rdkafka_queue.h +++ /dev/null @@ -1,731 +0,0 @@ -/* - * librdkafka - The Apache Kafka C/C++ library - * - * Copyright (c) 2016 Magnus Edenhill - * All rights reserved. - * - * Redistribution and use in source and binary forms, with or without - * modification, are permitted provided that the following conditions are met: - * - * 1. Redistributions of source code must retain the above copyright notice, - * this list of conditions and the following disclaimer. - * 2. Redistributions in binary form must reproduce the above copyright notice, - * this list of conditions and the following disclaimer in the documentation - * and/or other materials provided with the distribution. - * - * THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS" - * AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE - * IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE - * ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT OWNER OR CONTRIBUTORS BE - * LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR - * CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF - * SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS - * INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN - * CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) - * ARISING IN ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE - * POSSIBILITY OF SUCH DAMAGE. - */ - -#pragma once - -#include "rdkafka_op.h" -#include "rdkafka_int.h" - -#ifdef _MSC_VER -#include <io.h> /* for _write() */ -#endif - - -TAILQ_HEAD(rd_kafka_op_tailq, rd_kafka_op_s); - -struct rd_kafka_q_s { - mtx_t rkq_lock; - cnd_t rkq_cond; - struct rd_kafka_q_s *rkq_fwdq; /* Forwarded/Routed queue. - * Used in place of this queue - * for all operations. */ - - struct rd_kafka_op_tailq rkq_q; /* TAILQ_HEAD(, rd_kafka_op_s) */ - int rkq_qlen; /* Number of entries in queue */ - int64_t rkq_qsize; /* Size of all entries in queue */ - int rkq_refcnt; - int rkq_flags; -#define RD_KAFKA_Q_F_ALLOCATED 0x1 /* Allocated: rd_free on destroy */ -#define RD_KAFKA_Q_F_READY 0x2 /* Queue is ready to be used. - * Flag is cleared on destroy */ -#define RD_KAFKA_Q_F_FWD_APP 0x4 /* Queue is being forwarded by a call - * to rd_kafka_queue_forward. */ - - rd_kafka_t *rkq_rk; - struct rd_kafka_q_io *rkq_qio; /* FD-based application signalling */ - - /* Op serve callback (optional). - * Mainly used for forwarded queues to use the original queue's - * serve function from the forwarded position. - * Shall return 1 if op was handled, else 0. */ - rd_kafka_q_serve_cb_t *rkq_serve; - void *rkq_opaque; - -#if ENABLE_DEVEL - char rkq_name[64]; /* Debugging: queue name (FUNC:LINE) */ -#else - const char *rkq_name; /* Debugging: queue name (FUNC) */ -#endif -}; - - -/* FD-based application signalling state holder. */ -struct rd_kafka_q_io { - int fd; - void *payload; - size_t size; -}; - - - -/** - * @return true if queue is ready/enabled, else false. - * @remark queue luck must be held by caller (if applicable) - */ -static RD_INLINE RD_UNUSED -int rd_kafka_q_ready (rd_kafka_q_t *rkq) { - return rkq->rkq_flags & RD_KAFKA_Q_F_READY; -} - - - - -void rd_kafka_q_init (rd_kafka_q_t *rkq, rd_kafka_t *rk); -rd_kafka_q_t *rd_kafka_q_new0 (rd_kafka_t *rk, const char *func, int line); -#define rd_kafka_q_new(rk) rd_kafka_q_new0(rk,__FUNCTION__,__LINE__) -void rd_kafka_q_destroy_final (rd_kafka_q_t *rkq); - -#define rd_kafka_q_lock(rkqu) mtx_lock(&(rkqu)->rkq_lock) -#define rd_kafka_q_unlock(rkqu) mtx_unlock(&(rkqu)->rkq_lock) - -static RD_INLINE RD_UNUSED -rd_kafka_q_t *rd_kafka_q_keep (rd_kafka_q_t *rkq) { - mtx_lock(&rkq->rkq_lock); - rkq->rkq_refcnt++; - mtx_unlock(&rkq->rkq_lock); - return rkq; -} - -static RD_INLINE RD_UNUSED -rd_kafka_q_t *rd_kafka_q_keep_nolock (rd_kafka_q_t *rkq) { - rkq->rkq_refcnt++; - return rkq; -} - - -/** - * @returns the queue's name (used for debugging) - */ -static RD_INLINE RD_UNUSED -const char *rd_kafka_q_name (rd_kafka_q_t *rkq) { - return rkq->rkq_name; -} - -/** - * @returns the final destination queue name (after forwarding) - * @remark rkq MUST NOT be locked - */ -static RD_INLINE RD_UNUSED -const char *rd_kafka_q_dest_name (rd_kafka_q_t *rkq) { - const char *ret; - mtx_lock(&rkq->rkq_lock); - if (rkq->rkq_fwdq) - ret = rd_kafka_q_dest_name(rkq->rkq_fwdq); - else - ret = rd_kafka_q_name(rkq); - mtx_unlock(&rkq->rkq_lock); - return ret; -} - - -static RD_INLINE RD_UNUSED -void rd_kafka_q_destroy (rd_kafka_q_t *rkq) { - int do_delete = 0; - - mtx_lock(&rkq->rkq_lock); - rd_kafka_assert(NULL, rkq->rkq_refcnt > 0); - do_delete = !--rkq->rkq_refcnt; - mtx_unlock(&rkq->rkq_lock); - - if (unlikely(do_delete)) - rd_kafka_q_destroy_final(rkq); -} - - -/** - * Reset a queue. - * WARNING: All messages will be lost and leaked. - * NOTE: No locking is performed. - */ -static RD_INLINE RD_UNUSED -void rd_kafka_q_reset (rd_kafka_q_t *rkq) { - TAILQ_INIT(&rkq->rkq_q); - rd_dassert(TAILQ_EMPTY(&rkq->rkq_q)); - rkq->rkq_qlen = 0; - rkq->rkq_qsize = 0; -} - - -/** - * Disable a queue. - * Attempting to enqueue messages to the queue will destroy them. - */ -static RD_INLINE RD_UNUSED -void rd_kafka_q_disable0 (rd_kafka_q_t *rkq, int do_lock) { - if (do_lock) - mtx_lock(&rkq->rkq_lock); - rkq->rkq_flags &= ~RD_KAFKA_Q_F_READY; - if (do_lock) - mtx_unlock(&rkq->rkq_lock); -} -#define rd_kafka_q_disable(rkq) rd_kafka_q_disable0(rkq, 1/*lock*/) - -/** - * Forward 'srcq' to 'destq' - */ -void rd_kafka_q_fwd_set0 (rd_kafka_q_t *srcq, rd_kafka_q_t *destq, - int do_lock, int fwd_app); -#define rd_kafka_q_fwd_set(S,D) rd_kafka_q_fwd_set0(S,D,1/*lock*/,\ - 0/*no fwd_app*/) - -/** - * @returns the forward queue (if any) with its refcount increased. - * @locks rd_kafka_q_lock(rkq) == !do_lock - */ -static RD_INLINE RD_UNUSED -rd_kafka_q_t *rd_kafka_q_fwd_get (rd_kafka_q_t *rkq, int do_lock) { - rd_kafka_q_t *fwdq; - if (do_lock) - mtx_lock(&rkq->rkq_lock); - - if ((fwdq = rkq->rkq_fwdq)) - rd_kafka_q_keep(fwdq); - - if (do_lock) - mtx_unlock(&rkq->rkq_lock); - - return fwdq; -} - - -/** - * @returns true if queue is forwarded, else false. - * - * @remark Thread-safe. - */ -static RD_INLINE RD_UNUSED int rd_kafka_q_is_fwded (rd_kafka_q_t *rkq) { - int r; - mtx_lock(&rkq->rkq_lock); - r = rkq->rkq_fwdq ? 1 : 0; - mtx_unlock(&rkq->rkq_lock); - return r; -} - - - -/** - * @brief Trigger an IO event for this queue. - * - * @remark Queue MUST be locked - */ -static RD_INLINE RD_UNUSED -void rd_kafka_q_io_event (rd_kafka_q_t *rkq) { - ssize_t r; - - if (likely(!rkq->rkq_qio)) - return; - -#ifdef _MSC_VER - r = _write(rkq->rkq_qio->fd, rkq->rkq_qio->payload, (int)rkq->rkq_qio->size); -#else - r = write(rkq->rkq_qio->fd, rkq->rkq_qio->payload, rkq->rkq_qio->size); -#endif - if (r == -1) { - fprintf(stderr, - "[ERROR:librdkafka:rd_kafka_q_io_event: " - "write(%d,..,%d) failed on queue %p \"%s\": %s: " - "disabling further IO events]\n", - rkq->rkq_qio->fd, (int)rkq->rkq_qio->size, - rkq, rd_kafka_q_name(rkq), rd_strerror(errno)); - /* FIXME: Log this, somehow */ - rd_free(rkq->rkq_qio); - rkq->rkq_qio = NULL; - } -} - - -/** - * @brief rko->rko_prio comparator - * @remark: descending order: higher priority takes preceedence. - */ -static RD_INLINE RD_UNUSED -int rd_kafka_op_cmp_prio (const void *_a, const void *_b) { - const rd_kafka_op_t *a = _a, *b = _b; - - return b->rko_prio - a->rko_prio; -} - - -/** - * @brief Low-level unprotected enqueue that only performs - * the actual queue enqueue and counter updates. - * @remark Will not perform locking, signaling, fwdq, READY checking, etc. - */ -static RD_INLINE RD_UNUSED void -rd_kafka_q_enq0 (rd_kafka_q_t *rkq, rd_kafka_op_t *rko, int at_head) { - if (likely(!rko->rko_prio)) - TAILQ_INSERT_TAIL(&rkq->rkq_q, rko, rko_link); - else if (at_head) - TAILQ_INSERT_HEAD(&rkq->rkq_q, rko, rko_link); - else - TAILQ_INSERT_SORTED(&rkq->rkq_q, rko, rd_kafka_op_t *, - rko_link, rd_kafka_op_cmp_prio); - rkq->rkq_qlen++; - rkq->rkq_qsize += rko->rko_len; -} - - -/** - * @brief Enqueue the 'rko' op at the tail of the queue 'rkq'. - * - * The provided 'rko' is either enqueued or destroyed. - * - * @returns 1 if op was enqueued or 0 if queue is disabled and - * there was no replyq to enqueue on in which case the rko is destroyed. - * - * Locality: any thread. - */ -static RD_INLINE RD_UNUSED -int rd_kafka_q_enq (rd_kafka_q_t *rkq, rd_kafka_op_t *rko) { - rd_kafka_q_t *fwdq; - - mtx_lock(&rkq->rkq_lock); - - rd_dassert(rkq->rkq_refcnt > 0); - - if (unlikely(!(rkq->rkq_flags & RD_KAFKA_Q_F_READY))) { - - /* Queue has been disabled, reply to and fail the rko. */ - mtx_unlock(&rkq->rkq_lock); - - return rd_kafka_op_reply(rko, RD_KAFKA_RESP_ERR__DESTROY); - } - - if (!rko->rko_serve && rkq->rkq_serve) { - /* Store original queue's serve callback and opaque - * prior to forwarding. */ - rko->rko_serve = rkq->rkq_serve; - rko->rko_serve_opaque = rkq->rkq_opaque; - } - - if (!(fwdq = rd_kafka_q_fwd_get(rkq, 0))) { - rd_kafka_q_enq0(rkq, rko, 0); - cnd_signal(&rkq->rkq_cond); - if (rkq->rkq_qlen == 1) - rd_kafka_q_io_event(rkq); - mtx_unlock(&rkq->rkq_lock); - } else { - mtx_unlock(&rkq->rkq_lock); - rd_kafka_q_enq(fwdq, rko); - rd_kafka_q_destroy(fwdq); - } - - return 1; -} - - -/** - * @brief Re-enqueue rko at head of rkq. - * - * The provided 'rko' is either enqueued or destroyed. - * - * @returns 1 if op was enqueued or 0 if queue is disabled and - * there was no replyq to enqueue on in which case the rko is destroyed. - * - * @locks rkq MUST BE LOCKED - * - * Locality: any thread. - */ -static RD_INLINE RD_UNUSED -int rd_kafka_q_reenq (rd_kafka_q_t *rkq, rd_kafka_op_t *rko) { - rd_kafka_q_t *fwdq; - - rd_dassert(rkq->rkq_refcnt > 0); - - if (unlikely(!(rkq->rkq_flags & RD_KAFKA_Q_F_READY))) - return rd_kafka_op_reply(rko, RD_KAFKA_RESP_ERR__DESTROY); - - if (!rko->rko_serve && rkq->rkq_serve) { - /* Store original queue's serve callback and opaque - * prior to forwarding. */ - rko->rko_serve = rkq->rkq_serve; - rko->rko_serve_opaque = rkq->rkq_opaque; - } - - if (!(fwdq = rd_kafka_q_fwd_get(rkq, 0))) { - rd_kafka_q_enq0(rkq, rko, 1/*at_head*/); - cnd_signal(&rkq->rkq_cond); - if (rkq->rkq_qlen == 1) - rd_kafka_q_io_event(rkq); - } else { - rd_kafka_q_enq(fwdq, rko); - rd_kafka_q_destroy(fwdq); - } - - return 1; -} - - -/** - * Dequeue 'rko' from queue 'rkq'. - * - * NOTE: rkq_lock MUST be held - * Locality: any thread - */ -static RD_INLINE RD_UNUSED -void rd_kafka_q_deq0 (rd_kafka_q_t *rkq, rd_kafka_op_t *rko) { - rd_dassert(rkq->rkq_flags & RD_KAFKA_Q_F_READY); - rd_dassert(rkq->rkq_qlen > 0 && - rkq->rkq_qsize >= (int64_t)rko->rko_len); - - TAILQ_REMOVE(&rkq->rkq_q, rko, rko_link); - rkq->rkq_qlen--; - rkq->rkq_qsize -= rko->rko_len; -} - -/** - * Concat all elements of 'srcq' onto tail of 'rkq'. - * 'rkq' will be be locked (if 'do_lock'==1), but 'srcq' will not. - * NOTE: 'srcq' will be reset. - * - * Locality: any thread. - * - * @returns 0 if operation was performed or -1 if rkq is disabled. - */ -static RD_INLINE RD_UNUSED -int rd_kafka_q_concat0 (rd_kafka_q_t *rkq, rd_kafka_q_t *srcq, int do_lock) { - int r = 0; - - while (srcq->rkq_fwdq) /* Resolve source queue */ - srcq = srcq->rkq_fwdq; - if (unlikely(srcq->rkq_qlen == 0)) - return 0; /* Don't do anything if source queue is empty */ - - if (do_lock) - mtx_lock(&rkq->rkq_lock); - if (!rkq->rkq_fwdq) { - rd_kafka_op_t *rko; - - rd_dassert(TAILQ_EMPTY(&srcq->rkq_q) || - srcq->rkq_qlen > 0); - if (unlikely(!(rkq->rkq_flags & RD_KAFKA_Q_F_READY))) { - if (do_lock) - mtx_unlock(&rkq->rkq_lock); - return -1; - } - /* First insert any prioritized ops from srcq - * in the right position in rkq. */ - while ((rko = TAILQ_FIRST(&srcq->rkq_q)) && rko->rko_prio > 0) { - TAILQ_REMOVE(&srcq->rkq_q, rko, rko_link); - TAILQ_INSERT_SORTED(&rkq->rkq_q, rko, - rd_kafka_op_t *, rko_link, - rd_kafka_op_cmp_prio); - } - - TAILQ_CONCAT(&rkq->rkq_q, &srcq->rkq_q, rko_link); - if (rkq->rkq_qlen == 0) - rd_kafka_q_io_event(rkq); - rkq->rkq_qlen += srcq->rkq_qlen; - rkq->rkq_qsize += srcq->rkq_qsize; - cnd_signal(&rkq->rkq_cond); - - rd_kafka_q_reset(srcq); - } else - r = rd_kafka_q_concat0(rkq->rkq_fwdq ? rkq->rkq_fwdq : rkq, - srcq, - rkq->rkq_fwdq ? do_lock : 0); - if (do_lock) - mtx_unlock(&rkq->rkq_lock); - - return r; -} - -#define rd_kafka_q_concat(dstq,srcq) rd_kafka_q_concat0(dstq,srcq,1/*lock*/) - - -/** - * @brief Prepend all elements of 'srcq' onto head of 'rkq'. - * 'rkq' will be be locked (if 'do_lock'==1), but 'srcq' will not. - * 'srcq' will be reset. - * - * @remark Will not respect priority of ops, srcq will be prepended in its - * original form to rkq. - * - * @locality any thread. - */ -static RD_INLINE RD_UNUSED -void rd_kafka_q_prepend0 (rd_kafka_q_t *rkq, rd_kafka_q_t *srcq, - int do_lock) { - if (do_lock) - mtx_lock(&rkq->rkq_lock); - if (!rkq->rkq_fwdq && !srcq->rkq_fwdq) { - /* FIXME: prio-aware */ - /* Concat rkq on srcq */ - TAILQ_CONCAT(&srcq->rkq_q, &rkq->rkq_q, rko_link); - /* Move srcq to rkq */ - TAILQ_MOVE(&rkq->rkq_q, &srcq->rkq_q, rko_link); - if (rkq->rkq_qlen == 0 && srcq->rkq_qlen > 0) - rd_kafka_q_io_event(rkq); - rkq->rkq_qlen += srcq->rkq_qlen; - rkq->rkq_qsize += srcq->rkq_qsize; - - rd_kafka_q_reset(srcq); - } else - rd_kafka_q_prepend0(rkq->rkq_fwdq ? rkq->rkq_fwdq : rkq, - srcq->rkq_fwdq ? srcq->rkq_fwdq : srcq, - rkq->rkq_fwdq ? do_lock : 0); - if (do_lock) - mtx_unlock(&rkq->rkq_lock); -} - -#define rd_kafka_q_prepend(dstq,srcq) rd_kafka_q_prepend0(dstq,srcq,1/*lock*/) - - -/* Returns the number of elements in the queue */ -static RD_INLINE RD_UNUSED -int rd_kafka_q_len (rd_kafka_q_t *rkq) { - int qlen; - rd_kafka_q_t *fwdq; - mtx_lock(&rkq->rkq_lock); - if (!(fwdq = rd_kafka_q_fwd_get(rkq, 0))) { - qlen = rkq->rkq_qlen; - mtx_unlock(&rkq->rkq_lock); - } else { - mtx_unlock(&rkq->rkq_lock); - qlen = rd_kafka_q_len(fwdq); - rd_kafka_q_destroy(fwdq); - } - return qlen; -} - -/* Returns the total size of elements in the queue */ -static RD_INLINE RD_UNUSED -uint64_t rd_kafka_q_size (rd_kafka_q_t *rkq) { - uint64_t sz; - rd_kafka_q_t *fwdq; - mtx_lock(&rkq->rkq_lock); - if (!(fwdq = rd_kafka_q_fwd_get(rkq, 0))) { - sz = rkq->rkq_qsize; - mtx_unlock(&rkq->rkq_lock); - } else { - mtx_unlock(&rkq->rkq_lock); - sz = rd_kafka_q_size(fwdq); - rd_kafka_q_destroy(fwdq); - } - return sz; -} - - -/* Construct temporary on-stack replyq with increased Q refcount and - * optional VERSION. */ -#if ENABLE_DEVEL -#define RD_KAFKA_REPLYQ(Q,VERSION) \ - (rd_kafka_replyq_t){rd_kafka_q_keep(Q), VERSION, \ - rd_strdup(__FUNCTION__) } -#else -#define RD_KAFKA_REPLYQ(Q,VERSION) \ - (rd_kafka_replyq_t){rd_kafka_q_keep(Q), VERSION} -#endif - -/* Construct temporary on-stack replyq for indicating no replyq. */ -#if ENABLE_DEVEL -#define RD_KAFKA_NO_REPLYQ (rd_kafka_replyq_t){NULL, 0, NULL} -#else -#define RD_KAFKA_NO_REPLYQ (rd_kafka_replyq_t){NULL, 0} -#endif - -/** - * Set up replyq. - * Q refcnt is increased. - */ -static RD_INLINE RD_UNUSED void -rd_kafka_set_replyq (rd_kafka_replyq_t *replyq, - rd_kafka_q_t *rkq, int32_t version) { - replyq->q = rkq ? rd_kafka_q_keep(rkq) : NULL; - replyq->version = version; -#if ENABLE_DEVEL - replyq->_id = strdup(__FUNCTION__); -#endif -} - -/** - * Set rko's replyq with an optional version (versionptr != NULL). - * Q refcnt is increased. - */ -static RD_INLINE RD_UNUSED void -rd_kafka_op_set_replyq (rd_kafka_op_t *rko, rd_kafka_q_t *rkq, - rd_atomic32_t *versionptr) { - rd_kafka_set_replyq(&rko->rko_replyq, rkq, - versionptr ? rd_atomic32_get(versionptr) : 0); -} - -/* Set reply rko's version from replyq's version */ -#define rd_kafka_op_get_reply_version(REPLY_RKO, ORIG_RKO) do { \ - (REPLY_RKO)->rko_version = (ORIG_RKO)->rko_replyq.version; \ - } while (0) - - -/* Clear replyq holder without decreasing any .q references. */ -static RD_INLINE RD_UNUSED void -rd_kafka_replyq_clear (rd_kafka_replyq_t *replyq) { - memset(replyq, 0, sizeof(*replyq)); -} - -/** - * @brief Make a copy of \p src in \p dst, with its own queue reference - */ -static RD_INLINE RD_UNUSED void -rd_kafka_replyq_copy (rd_kafka_replyq_t *dst, rd_kafka_replyq_t *src) { - dst->version = src->version; - dst->q = src->q; - if (dst->q) - rd_kafka_q_keep(dst->q); -#if ENABLE_DEVEL - if (src->_id) - dst->_id = rd_strdup(src->_id); - else - dst->_id = NULL; -#endif -} - - -/** - * Clear replyq holder and destroy any .q references. - */ -static RD_INLINE RD_UNUSED void -rd_kafka_replyq_destroy (rd_kafka_replyq_t *replyq) { - if (replyq->q) - rd_kafka_q_destroy(replyq->q); -#if ENABLE_DEVEL - if (replyq->_id) { - rd_free(replyq->_id); - replyq->_id = NULL; - } -#endif - rd_kafka_replyq_clear(replyq); -} - - -/** - * @brief Wrapper for rd_kafka_q_enq() that takes a replyq, - * steals its queue reference, enqueues the op with the replyq version, - * and then destroys the queue reference. - * - * If \p version is non-zero it will be updated, else replyq->version. - * - * @returns Same as rd_kafka_q_enq() - */ -static RD_INLINE RD_UNUSED int -rd_kafka_replyq_enq (rd_kafka_replyq_t *replyq, rd_kafka_op_t *rko, - int version) { - rd_kafka_q_t *rkq = replyq->q; - int r; - - if (version) - rko->rko_version = version; - else - rko->rko_version = replyq->version; - - /* The replyq queue reference is done after we've enqueued the rko - * so clear it here. */ - replyq->q = NULL; - -#if ENABLE_DEVEL - if (replyq->_id) { - rd_free(replyq->_id); - replyq->_id = NULL; - } -#endif - - /* Retain replyq->version since it is used by buf_callback - * when dispatching the callback. */ - - r = rd_kafka_q_enq(rkq, rko); - - rd_kafka_q_destroy(rkq); - - return r; -} - - - -rd_kafka_op_t *rd_kafka_q_pop_serve (rd_kafka_q_t *rkq, int timeout_ms, - int32_t version, - rd_kafka_q_cb_type_t cb_type, - rd_kafka_q_serve_cb_t *callback, - void *opaque); -rd_kafka_op_t *rd_kafka_q_pop (rd_kafka_q_t *rkq, int timeout_ms, - int32_t version); -int rd_kafka_q_serve (rd_kafka_q_t *rkq, int timeout_ms, int max_cnt, - rd_kafka_q_cb_type_t cb_type, - rd_kafka_q_serve_cb_t *callback, - void *opaque); - -int rd_kafka_q_purge0 (rd_kafka_q_t *rkq, int do_lock); -#define rd_kafka_q_purge(rkq) rd_kafka_q_purge0(rkq, 1/*lock*/) -void rd_kafka_q_purge_toppar_version (rd_kafka_q_t *rkq, - rd_kafka_toppar_t *rktp, int version); - -int rd_kafka_q_move_cnt (rd_kafka_q_t *dstq, rd_kafka_q_t *srcq, - int cnt, int do_locks); - -int rd_kafka_q_serve_rkmessages (rd_kafka_q_t *rkq, int timeout_ms, - rd_kafka_message_t **rkmessages, - size_t rkmessages_size); -rd_kafka_resp_err_t rd_kafka_q_wait_result (rd_kafka_q_t *rkq, int timeout_ms); - -int rd_kafka_q_apply (rd_kafka_q_t *rkq, - int (*callback) (rd_kafka_q_t *rkq, rd_kafka_op_t *rko, - void *opaque), - void *opaque); - -void rd_kafka_q_fix_offsets (rd_kafka_q_t *rkq, int64_t min_offset, - int64_t base_offset); - -/** - * @returns the last op in the queue matching \p op_type and \p allow_err (bool) - * @remark The \p rkq must be properly locked before this call, the returned rko - * is not removed from the queue and may thus not be held for longer - * than the lock is held. - */ -static RD_INLINE RD_UNUSED -rd_kafka_op_t *rd_kafka_q_last (rd_kafka_q_t *rkq, rd_kafka_op_type_t op_type, - int allow_err) { - rd_kafka_op_t *rko; - TAILQ_FOREACH_REVERSE(rko, &rkq->rkq_q, rd_kafka_op_tailq, rko_link) { - if (rko->rko_type == op_type && - (allow_err || !rko->rko_err)) - return rko; - } - - return NULL; -} - -void rd_kafka_q_io_event_enable (rd_kafka_q_t *rkq, int fd, - const void *payload, size_t size); - -/* Public interface */ -struct rd_kafka_queue_s { - rd_kafka_q_t *rkqu_q; - rd_kafka_t *rkqu_rk; -}; - - -void rd_kafka_q_dump (FILE *fp, rd_kafka_q_t *rkq); - -extern int RD_TLS rd_kafka_yield_thread; http://git-wip-us.apache.org/repos/asf/nifi-minifi-cpp/blob/7528d23e/thirdparty/librdkafka-0.11.1/src/rdkafka_range_assignor.c ---------------------------------------------------------------------- diff --git a/thirdparty/librdkafka-0.11.1/src/rdkafka_range_assignor.c b/thirdparty/librdkafka-0.11.1/src/rdkafka_range_assignor.c deleted file mode 100644 index dfa9893..0000000 --- a/thirdparty/librdkafka-0.11.1/src/rdkafka_range_assignor.c +++ /dev/null @@ -1,125 +0,0 @@ -/* - * librdkafka - The Apache Kafka C/C++ library - * - * Copyright (c) 2015 Magnus Edenhill - * All rights reserved. - * - * Redistribution and use in source and binary forms, with or without - * modification, are permitted provided that the following conditions are met: - * - * 1. Redistributions of source code must retain the above copyright notice, - * this list of conditions and the following disclaimer. - * 2. Redistributions in binary form must reproduce the above copyright notice, - * this list of conditions and the following disclaimer in the documentation - * and/or other materials provided with the distribution. - * - * THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS" - * AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE - * IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE - * ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT OWNER OR CONTRIBUTORS BE - * LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR - * CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF - * SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS - * INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN - * CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) - * ARISING IN ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE - * POSSIBILITY OF SUCH DAMAGE. - */ -#include "rdkafka_int.h" -#include "rdkafka_assignor.h" - - - - - - -/** - * Source: https://github.com/apache/kafka/blob/trunk/clients/src/main/java/org/apache/kafka/clients/consumer/RangeAssignor.java - * - * The range assignor works on a per-topic basis. For each topic, we lay out the available partitions in numeric order - * and the consumers in lexicographic order. We then divide the number of partitions by the total number of - * consumers to determine the number of partitions to assign to each consumer. If it does not evenly - * divide, then the first few consumers will have one extra partition. - * - * For example, suppose there are two consumers C0 and C1, two topics t0 and t1, and each topic has 3 partitions, - * resulting in partitions t0p0, t0p1, t0p2, t1p0, t1p1, and t1p2. - * - * The assignment will be: - * C0: [t0p0, t0p1, t1p0, t1p1] - * C1: [t0p2, t1p2] - */ - -rd_kafka_resp_err_t -rd_kafka_range_assignor_assign_cb (rd_kafka_t *rk, - const char *member_id, - const char *protocol_name, - const rd_kafka_metadata_t *metadata, - rd_kafka_group_member_t *members, - size_t member_cnt, - rd_kafka_assignor_topic_t **eligible_topics, - size_t eligible_topic_cnt, - char *errstr, size_t errstr_size, - void *opaque) { - unsigned int ti; - int i; - - /* The range assignor works on a per-topic basis. */ - for (ti = 0 ; ti < eligible_topic_cnt ; ti++) { - rd_kafka_assignor_topic_t *eligible_topic = eligible_topics[ti]; - int numPartitionsPerConsumer; - int consumersWithExtraPartition; - - /* For each topic, we lay out the available partitions in - * numeric order and the consumers in lexicographic order. */ - rd_list_sort(&eligible_topic->members, - rd_kafka_group_member_cmp); - - /* We then divide the number of partitions by the total number of - * consumers to determine the number of partitions to assign to - * each consumer. */ - numPartitionsPerConsumer = - eligible_topic->metadata->partition_cnt / - rd_list_cnt(&eligible_topic->members); - - /* If it does not evenly divide, then the first few consumers - * will have one extra partition. */ - consumersWithExtraPartition = - eligible_topic->metadata->partition_cnt % - rd_list_cnt(&eligible_topic->members); - - rd_kafka_dbg(rk, CGRP, "ASSIGN", - "range: Topic %s with %d partition(s) and " - "%d subscribing member(s)", - eligible_topic->metadata->topic, - eligible_topic->metadata->partition_cnt, - rd_list_cnt(&eligible_topic->members)); - - for (i = 0 ; i < rd_list_cnt(&eligible_topic->members) ; i++) { - rd_kafka_group_member_t *rkgm = - rd_list_elem(&eligible_topic->members, i); - int start = numPartitionsPerConsumer * i + - RD_MIN(i, consumersWithExtraPartition); - int length = numPartitionsPerConsumer + - (i + 1 > consumersWithExtraPartition ? 0 : 1); - - if (length == 0) - continue; - - rd_kafka_dbg(rk, CGRP, "ASSIGN", - "range: Member \"%s\": " - "assigned topic %s partitions %d..%d", - rkgm->rkgm_member_id->str, - eligible_topic->metadata->topic, - start, start+length-1); - rd_kafka_topic_partition_list_add_range( - rkgm->rkgm_assignment, - eligible_topic->metadata->topic, - start, start+length-1); - } - } - - return 0; -} - - -
