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;
-}
-
-
-

Reply via email to