http://git-wip-us.apache.org/repos/asf/nifi-minifi-cpp/blob/7528d23e/thirdparty/librdkafka-0.11.1/src/rdkafka_lz4.c
----------------------------------------------------------------------
diff --git a/thirdparty/librdkafka-0.11.1/src/rdkafka_lz4.c 
b/thirdparty/librdkafka-0.11.1/src/rdkafka_lz4.c
deleted file mode 100644
index 9ee50de..0000000
--- a/thirdparty/librdkafka-0.11.1/src/rdkafka_lz4.c
+++ /dev/null
@@ -1,429 +0,0 @@
-/*
- * librdkafka - Apache Kafka C library
- *
- * Copyright (c) 2017 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_lz4.h"
-
-#if WITH_LZ4_EXT
-#include <lz4frame.h>
-#else
-#include "lz4frame.h"
-#endif
-#include "xxhash.h"
-
-#include "rdbuf.h"
-
-/**
- * Fix-up bad LZ4 framing caused by buggy Kafka client / broker.
- * The LZ4F framing format is described in detail here:
- * https://github.com/Cyan4973/lz4/blob/master/lz4_Frame_format.md
- *
- * NOTE: This modifies 'inbuf'.
- *
- * Returns an error on failure to fix (nothing modified), else NO_ERROR.
- */
-static rd_kafka_resp_err_t
-rd_kafka_lz4_decompress_fixup_bad_framing (rd_kafka_broker_t *rkb,
-                                           char *inbuf, size_t inlen) {
-        static const char magic[4] = { 0x04, 0x22, 0x4d, 0x18 };
-        uint8_t FLG, HC, correct_HC;
-        size_t of = 4;
-
-        /* Format is:
-         *    int32_t magic;
-         *    int8_t_ FLG;
-         *    int8_t  BD;
-         *  [ int64_t contentSize; ]
-         *    int8_t  HC;
-         */
-        if (inlen < 4+3 || memcmp(inbuf, magic, 4)) {
-                rd_rkb_dbg(rkb, BROKER,  "LZ4FIXUP",
-                           "Unable to fix-up legacy LZ4 framing "
-                           "(%"PRIusz" bytes): invalid length or magic value",
-                           inlen);
-                return RD_KAFKA_RESP_ERR__BAD_COMPRESSION;
-        }
-
-        of = 4; /* past magic */
-        FLG = inbuf[of++];
-        of++; /* BD */
-
-        if ((FLG >> 3) & 1) /* contentSize */
-                of += 8;
-
-        if (of >= inlen) {
-                rd_rkb_dbg(rkb, BROKER,  "LZ4FIXUP",
-                           "Unable to fix-up legacy LZ4 framing "
-                           "(%"PRIusz" bytes): requires %"PRIusz" bytes",
-                           inlen, of);
-                return RD_KAFKA_RESP_ERR__BAD_COMPRESSION;
-        }
-
-        /* Header hash code */
-        HC = inbuf[of];
-
-        /* Calculate correct header hash code */
-        correct_HC = (XXH32(inbuf+4, of-4, 0) >> 8) & 0xff;
-
-        if (HC != correct_HC)
-                inbuf[of] = correct_HC;
-
-        return RD_KAFKA_RESP_ERR_NO_ERROR;
-}
-
-
-/**
- * Reverse of fix-up: break LZ4 framing caused to be compatbile with with
- * buggy Kafka client / broker.
- *
- * NOTE: This modifies 'outbuf'.
- *
- * Returns an error on failure to recognize format (nothing modified),
- * else NO_ERROR.
- */
-static rd_kafka_resp_err_t
-rd_kafka_lz4_compress_break_framing (rd_kafka_broker_t *rkb,
-                                     char *outbuf, size_t outlen) {
-        static const char magic[4] = { 0x04, 0x22, 0x4d, 0x18 };
-        uint8_t FLG, HC, bad_HC;
-        size_t of = 4;
-
-        /* Format is:
-         *    int32_t magic;
-         *    int8_t_ FLG;
-         *    int8_t  BD;
-         *  [ int64_t contentSize; ]
-         *    int8_t  HC;
-         */
-        if (outlen < 4+3 || memcmp(outbuf, magic, 4)) {
-                rd_rkb_dbg(rkb, BROKER,  "LZ4FIXDOWN",
-                           "Unable to break legacy LZ4 framing "
-                           "(%"PRIusz" bytes): invalid length or magic value",
-                           outlen);
-                return RD_KAFKA_RESP_ERR__BAD_COMPRESSION;
-        }
-
-        of = 4; /* past magic */
-        FLG = outbuf[of++];
-        of++; /* BD */
-
-        if ((FLG >> 3) & 1) /* contentSize */
-                of += 8;
-
-        if (of >= outlen) {
-                rd_rkb_dbg(rkb, BROKER,  "LZ4FIXUP",
-                           "Unable to break legacy LZ4 framing "
-                           "(%"PRIusz" bytes): requires %"PRIusz" bytes",
-                           outlen, of);
-                return RD_KAFKA_RESP_ERR__BAD_COMPRESSION;
-        }
-
-        /* Header hash code */
-        HC = outbuf[of];
-
-        /* Calculate bad header hash code (include magic) */
-        bad_HC = (XXH32(outbuf, of, 0) >> 8) & 0xff;
-
-        if (HC != bad_HC)
-                outbuf[of] = bad_HC;
-
-        return RD_KAFKA_RESP_ERR_NO_ERROR;
-}
-
-
-
-/**
- * @brief Decompress LZ4F (framed) data.
- *        Kafka broker versions <0.10.0.0 (MsgVersion 0) breaks LZ4 framing
- *        checksum, if \p proper_hc we assume the checksum is okay
- *        (broker version >=0.10.0, MsgVersion >= 1) else we fix it up.
- *
- * @remark May modify \p inbuf (if not \p proper_hc)
- */
-rd_kafka_resp_err_t
-rd_kafka_lz4_decompress (rd_kafka_broker_t *rkb, int proper_hc, int64_t Offset,
-                         char *inbuf, size_t inlen,
-                         void **outbuf, size_t *outlenp) {
-        LZ4F_errorCode_t code;
-        LZ4F_decompressionContext_t dctx;
-        LZ4F_frameInfo_t fi;
-        size_t in_sz, out_sz;
-        size_t in_of, out_of;
-        size_t r;
-        size_t estimated_uncompressed_size;
-        size_t outlen;
-        rd_kafka_resp_err_t err = RD_KAFKA_RESP_ERR_NO_ERROR;
-        char *out = NULL;
-
-        *outbuf = NULL;
-
-        code = LZ4F_createDecompressionContext(&dctx, LZ4F_VERSION);
-        if (LZ4F_isError(code)) {
-                rd_rkb_dbg(rkb, BROKER, "LZ4DECOMPR",
-                           "Unable to create LZ4 decompression context: %s",
-                           LZ4F_getErrorName(code));
-                return RD_KAFKA_RESP_ERR__CRIT_SYS_RESOURCE;
-        }
-
-        if (!proper_hc) {
-                /* The original/legacy LZ4 framing in Kafka was buggy and
-                 * calculated the LZ4 framing header hash code (HC) 
incorrectly.
-                 * We do a fix-up of it here. */
-                if ((err = rd_kafka_lz4_decompress_fixup_bad_framing(rkb,
-                                                                     inbuf,
-                                                                     inlen)))
-                        goto done;
-        }
-
-        in_sz = inlen;
-        r = LZ4F_getFrameInfo(dctx, &fi, (const void *)inbuf, &in_sz);
-        if (LZ4F_isError(r)) {
-                rd_rkb_dbg(rkb, BROKER, "LZ4DECOMPR",
-                           "Failed to gather LZ4 frame info: %s",
-                           LZ4F_getErrorName(r));
-                err = RD_KAFKA_RESP_ERR__BAD_COMPRESSION;
-                goto done;
-        }
-
-        /* If uncompressed size is unknown or out of bounds make up a
-         * worst-case uncompressed size
-         * More info on max size: http://stackoverflow.com/a/25751871/1821055 
*/
-        if (fi.contentSize == 0 || fi.contentSize > inlen * 255)
-                estimated_uncompressed_size = inlen * 255;
-        else
-                estimated_uncompressed_size = (size_t)fi.contentSize;
-
-        /* Allocate output buffer, we increase this later if needed,
-         * but hopefully not. */
-        out = rd_malloc(estimated_uncompressed_size);
-        if (!out) {
-                rd_rkb_log(rkb, LOG_WARNING, "LZ4DEC",
-                           "Unable to allocate decompression "
-                           "buffer of %zd bytes: %s",
-                           estimated_uncompressed_size, rd_strerror(errno));
-                err = RD_KAFKA_RESP_ERR__CRIT_SYS_RESOURCE;
-                goto done;
-        }
-
-
-        /* Decompress input buffer to output buffer until input is exhausted. 
*/
-        outlen = estimated_uncompressed_size;
-        in_of = in_sz;
-        out_of = 0;
-        while (in_of < inlen) {
-                out_sz = outlen - out_of;
-                in_sz = inlen - in_of;
-                r = LZ4F_decompress(dctx, out+out_of, &out_sz,
-                                    inbuf+in_of, &in_sz, NULL);
-                if (unlikely(LZ4F_isError(r))) {
-                        rd_rkb_dbg(rkb, MSG, "LZ4DEC",
-                                   "Failed to LZ4 (%s HC) decompress message "
-                                   "(offset %"PRId64") at "
-                                   "payload offset %"PRIusz"/%"PRIusz": %s",
-                                   proper_hc ? "proper":"legacy",
-                                   Offset, in_of, inlen,  
LZ4F_getErrorName(r));
-                        err = RD_KAFKA_RESP_ERR__BAD_COMPRESSION;
-                        goto done;
-                }
-
-                rd_kafka_assert(NULL, out_of + out_sz < outlen &&
-                                in_of + in_sz <= inlen);
-                out_of += out_sz;
-                in_of += in_sz;
-                if (r == 0)
-                        break;
-
-                /* Need to grow output buffer, this shouldn't happen if
-                 * contentSize was properly set. */
-                if (unlikely(r > 0 && out_of == outlen)) {
-                        char *tmp;
-                        size_t extra = (r > 1024 ? r : 1024) * 2;
-
-                        rd_atomic64_add(&rkb->rkb_c.zbuf_grow, 1);
-
-                        if ((tmp = rd_realloc(outbuf, outlen + extra))) {
-                                rd_rkb_log(rkb, LOG_WARNING, "LZ4DEC",
-                                           "Unable to grow decompression "
-                                           "buffer to %zd+%zd bytes: %s",
-                                           outlen, extra,rd_strerror(errno));
-                                err = RD_KAFKA_RESP_ERR__CRIT_SYS_RESOURCE;
-                                goto done;
-                        }
-                        outlen += extra;
-                }
-        }
-
-
-        if (in_of < inlen) {
-                rd_rkb_dbg(rkb, MSG, "LZ4DEC",
-                           "Failed to LZ4 (%s HC) decompress message "
-                           "(offset %"PRId64"): "
-                           "%"PRIusz" (out of %"PRIusz") bytes remaining",
-                           proper_hc ? "proper":"legacy",
-                           Offset, inlen-in_of, inlen);
-                err = RD_KAFKA_RESP_ERR__BAD_MSG;
-                goto done;
-        }
-
-        *outbuf = out;
-        *outlenp = out_of;
-
- done:
-        code = LZ4F_freeDecompressionContext(dctx);
-        if (LZ4F_isError(code)) {
-                rd_rkb_dbg(rkb, BROKER, "LZ4DECOMPR",
-                           "Failed to close LZ4 compression context: %s",
-                           LZ4F_getErrorName(code));
-                err = RD_KAFKA_RESP_ERR__BAD_COMPRESSION;
-        }
-
-        if (err && out)
-                rd_free(out);
-
-        return err;
-}
-
-
-/**
- * Allocate space for \p *outbuf and compress all \p iovlen buffers in \p iov.
- * @param proper_hc generate a proper HC (checksum) (kafka >=0.10.0.0, 
MsgVersion >= 1)
- * @param MessageSetSize indicates (at least) full uncompressed data size,
- *                       possibly including MessageSet fields that will not
- *                       be compressed.
- *
- * @returns allocated buffer in \p *outbuf, length in \p *outlenp.
- */
-rd_kafka_resp_err_t
-rd_kafka_lz4_compress (rd_kafka_broker_t *rkb, int proper_hc,
-                       rd_slice_t *slice, void **outbuf, size_t *outlenp) {
-        LZ4F_compressionContext_t cctx;
-        LZ4F_errorCode_t r;
-        rd_kafka_resp_err_t err = RD_KAFKA_RESP_ERR_NO_ERROR;
-        size_t len = rd_slice_remains(slice);
-        size_t out_sz;
-        size_t out_of = 0;
-        char *out;
-        const void *p;
-        size_t rlen;
-
-        /* Required by Kafka */
-        const LZ4F_preferences_t prefs =
-                { .frameInfo = { .blockMode = LZ4F_blockIndependent } };
-
-        *outbuf = NULL;
-
-        out_sz = LZ4F_compressBound(len, NULL) + 1000;
-        if (LZ4F_isError(out_sz)) {
-                rd_rkb_dbg(rkb, MSG, "LZ4COMPR",
-                           "Unable to query LZ4 compressed size "
-                           "(for %"PRIusz" uncompressed bytes): %s",
-                           len, LZ4F_getErrorName(out_sz));
-                return RD_KAFKA_RESP_ERR__BAD_MSG;
-        }
-
-        out = rd_malloc(out_sz);
-        if (!out) {
-                rd_rkb_dbg(rkb, MSG, "LZ4COMPR",
-                           "Unable to allocate output buffer "
-                           "(%"PRIusz" bytes): %s",
-                           out_sz, rd_strerror(errno));
-                return RD_KAFKA_RESP_ERR__CRIT_SYS_RESOURCE;
-        }
-
-        r = LZ4F_createCompressionContext(&cctx, LZ4F_VERSION);
-        if (LZ4F_isError(r)) {
-                rd_rkb_dbg(rkb, MSG, "LZ4COMPR",
-                           "Unable to create LZ4 compression context: %s",
-                           LZ4F_getErrorName(r));
-                return RD_KAFKA_RESP_ERR__CRIT_SYS_RESOURCE;
-        }
-
-        r = LZ4F_compressBegin(cctx, out, out_sz, &prefs);
-        if (LZ4F_isError(r)) {
-                rd_rkb_dbg(rkb, MSG, "LZ4COMPR",
-                           "Unable to begin LZ4 compression "
-                           "(out buffer is %"PRIusz" bytes): %s",
-                           out_sz, LZ4F_getErrorName(r));
-                err = RD_KAFKA_RESP_ERR__BAD_COMPRESSION;
-                goto done;
-        }
-
-        out_of += r;
-
-        while ((rlen = rd_slice_reader(slice, &p))) {
-                rd_assert(out_of < out_sz);
-                r = LZ4F_compressUpdate(cctx, out+out_of, out_sz-out_of,
-                                        p, rlen, NULL);
-                if (unlikely(LZ4F_isError(r))) {
-                        rd_rkb_dbg(rkb, MSG, "LZ4COMPR",
-                                   "LZ4 compression failed "
-                                   "(at of %"PRIusz" bytes, with "
-                                   "%"PRIusz" bytes remaining in out buffer): "
-                                   "%s",
-                                   rlen, out_sz - out_of,
-                                   LZ4F_getErrorName(r));
-                        err = RD_KAFKA_RESP_ERR__BAD_COMPRESSION;
-                        goto done;
-                }
-
-                out_of += r;
-        }
-
-        rd_assert(rd_slice_remains(slice) == 0);
-
-        r = LZ4F_compressEnd(cctx, out+out_of, out_sz-out_of, NULL);
-        if (unlikely(LZ4F_isError(r))) {
-                rd_rkb_dbg(rkb, MSG, "LZ4COMPR",
-                           "Failed to finalize LZ4 compression "
-                           "of %"PRIusz" bytes: %s",
-                           len, LZ4F_getErrorName(r));
-                err = RD_KAFKA_RESP_ERR__BAD_COMPRESSION;
-                goto done;
-        }
-
-        out_of += r;
-
-        /* For the broken legacy framing we need to mess up the header checksum
-         * so that the Kafka client / broker code accepts it. */
-        if (!proper_hc)
-                if ((err = rd_kafka_lz4_compress_break_framing(rkb,
-                                                               out, out_of)))
-                        goto done;
-
-
-        *outbuf  = out;
-        *outlenp = out_of;
-
- done:
-        LZ4F_freeCompressionContext(cctx);
-
-        if (err)
-                rd_free(out);
-
-        return err;
-
-}

http://git-wip-us.apache.org/repos/asf/nifi-minifi-cpp/blob/7528d23e/thirdparty/librdkafka-0.11.1/src/rdkafka_lz4.h
----------------------------------------------------------------------
diff --git a/thirdparty/librdkafka-0.11.1/src/rdkafka_lz4.h 
b/thirdparty/librdkafka-0.11.1/src/rdkafka_lz4.h
deleted file mode 100644
index fb72f21..0000000
--- a/thirdparty/librdkafka-0.11.1/src/rdkafka_lz4.h
+++ /dev/null
@@ -1,40 +0,0 @@
-/*
- * librdkafka - Apache Kafka C library
- *
- * Copyright (c) 2017 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
-
-
-rd_kafka_resp_err_t
-rd_kafka_lz4_decompress (rd_kafka_broker_t *rkb, int proper_hc, int64_t Offset,
-                         char *inbuf, size_t inlen,
-                         void **outbuf, size_t *outlenp);
-
-rd_kafka_resp_err_t
-rd_kafka_lz4_compress (rd_kafka_broker_t *rkb, int proper_hc,
-                       rd_slice_t *slice, void **outbuf, size_t *outlenp);

http://git-wip-us.apache.org/repos/asf/nifi-minifi-cpp/blob/7528d23e/thirdparty/librdkafka-0.11.1/src/rdkafka_metadata.c
----------------------------------------------------------------------
diff --git a/thirdparty/librdkafka-0.11.1/src/rdkafka_metadata.c 
b/thirdparty/librdkafka-0.11.1/src/rdkafka_metadata.c
deleted file mode 100644
index 135ac84..0000000
--- a/thirdparty/librdkafka-0.11.1/src/rdkafka_metadata.c
+++ /dev/null
@@ -1,1017 +0,0 @@
-/*
- * librdkafka - Apache Kafka C library
- *
- * Copyright (c) 2012-2013, 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 "rd.h"
-#include "rdkafka_int.h"
-#include "rdkafka_topic.h"
-#include "rdkafka_broker.h"
-#include "rdkafka_request.h"
-#include "rdkafka_metadata.h"
-
-#include <string.h>
-
-
-
-rd_kafka_resp_err_t
-rd_kafka_metadata (rd_kafka_t *rk, int all_topics,
-                   rd_kafka_topic_t *only_rkt,
-                   const struct rd_kafka_metadata **metadatap,
-                   int timeout_ms) {
-        rd_kafka_q_t *rkq;
-        rd_kafka_broker_t *rkb;
-        rd_kafka_op_t *rko;
-       rd_ts_t ts_end = rd_timeout_init(timeout_ms);
-        rd_list_t topics;
-
-        /* Query any broker that is up, and if none are up pick the first one,
-         * if we're lucky it will be up before the timeout */
-       rkb = rd_kafka_broker_any_usable(rk, timeout_ms, 1);
-       if (!rkb)
-               return RD_KAFKA_RESP_ERR__TRANSPORT;
-
-        rkq = rd_kafka_q_new(rk);
-
-        rd_list_init(&topics, 0, rd_free);
-        if (!all_topics) {
-                if (only_rkt)
-                        rd_list_add(&topics,
-                                    rd_strdup(rd_kafka_topic_a2i(only_rkt)->
-                                              rkt_topic->str));
-                else
-                        rd_kafka_local_topics_to_list(rkb->rkb_rk, &topics);
-        }
-
-        /* Async: request metadata */
-        rko = rd_kafka_op_new(RD_KAFKA_OP_METADATA);
-        rd_kafka_op_set_replyq(rko, rkq, 0);
-        rko->rko_u.metadata.force = 1; /* Force metadata request regardless
-                                        * of outstanding metadata requests. */
-        rd_kafka_MetadataRequest(rkb, &topics, "application requested", rko);
-
-        rd_list_destroy(&topics);
-        rd_kafka_broker_destroy(rkb);
-
-        /* Wait for reply (or timeout) */
-        rko = rd_kafka_q_pop(rkq, rd_timeout_remains(ts_end), 0);
-
-        rd_kafka_q_destroy(rkq);
-
-        /* Timeout */
-        if (!rko)
-                return RD_KAFKA_RESP_ERR__TIMED_OUT;
-
-        /* Error */
-        if (rko->rko_err) {
-                rd_kafka_resp_err_t err = rko->rko_err;
-                rd_kafka_op_destroy(rko);
-                return err;
-        }
-
-        /* Reply: pass metadata pointer to application who now owns it*/
-        rd_kafka_assert(rk, rko->rko_u.metadata.md);
-        *metadatap = rko->rko_u.metadata.md;
-        rko->rko_u.metadata.md = NULL;
-        rd_kafka_op_destroy(rko);
-
-        return RD_KAFKA_RESP_ERR_NO_ERROR;
-}
-
-
-
-void rd_kafka_metadata_destroy (const struct rd_kafka_metadata *metadata) {
-        rd_free((void *)metadata);
-}
-
-
-/**
- * @returns a newly allocated copy of metadata \p src of size \p size
- */
-struct rd_kafka_metadata *
-rd_kafka_metadata_copy (const struct rd_kafka_metadata *src, size_t size) {
-       struct rd_kafka_metadata *md;
-       rd_tmpabuf_t tbuf;
-       int i;
-
-       /* metadata is stored in one contigious buffer where structs and
-        * and pointed-to fields are layed out in a memory aligned fashion.
-        * rd_tmpabuf_t provides the infrastructure to do this.
-        * Because of this we copy all the structs verbatim but
-        * any pointer fields needs to be copied explicitly to update
-        * the pointer address. */
-       rd_tmpabuf_new(&tbuf, size, 1/*assert on fail*/);
-       md = rd_tmpabuf_write(&tbuf, src, sizeof(*md));
-
-       rd_tmpabuf_write_str(&tbuf, src->orig_broker_name);
-
-
-       /* Copy Brokers */
-       md->brokers = rd_tmpabuf_write(&tbuf, src->brokers,
-                                     md->broker_cnt * sizeof(*md->brokers));
-
-       for (i = 0 ; i < md->broker_cnt ; i++)
-               md->brokers[i].host =
-                       rd_tmpabuf_write_str(&tbuf, src->brokers[i].host);
-
-
-       /* Copy TopicMetadata */
-        md->topics = rd_tmpabuf_write(&tbuf, src->topics,
-                                     md->topic_cnt * sizeof(*md->topics));
-
-       for (i = 0 ; i < md->topic_cnt ; i++) {
-               int j;
-
-               md->topics[i].topic = rd_tmpabuf_write_str(&tbuf,
-                                                          
src->topics[i].topic);
-
-
-               /* Copy partitions */
-               md->topics[i].partitions =
-                       rd_tmpabuf_write(&tbuf, src->topics[i].partitions,
-                                        md->topics[i].partition_cnt *
-                                        sizeof(*md->topics[i].partitions));
-
-               for (j = 0 ; j < md->topics[i].partition_cnt ; j++) {
-                       /* Copy replicas and ISRs */
-                       md->topics[i].partitions[j].replicas =
-                               rd_tmpabuf_write(&tbuf,
-                                                src->topics[i].partitions[j].
-                                                replicas,
-                                                md->topics[i].partitions[j].
-                                                replica_cnt *
-                                                sizeof(*md->topics[i].
-                                                       partitions[j].
-                                                       replicas));
-
-                       md->topics[i].partitions[j].isrs =
-                               rd_tmpabuf_write(&tbuf,
-                                                src->topics[i].partitions[j].
-                                                isrs,
-                                                md->topics[i].partitions[j].
-                                                isr_cnt *
-                                                sizeof(*md->topics[i].
-                                                       partitions[j].
-                                                       isrs));
-
-               }
-       }
-
-       /* Check for tmpabuf errors */
-       if (rd_tmpabuf_failed(&tbuf))
-               rd_kafka_assert(NULL, !*"metadata copy failed");
-
-       /* Delibarely not destroying the tmpabuf since we return
-        * its allocated memory. */
-
-       return md;
-}
-
-
-
-
-/**
- * Handle a Metadata response message.
- *
- * @param topics are the requested topics (may be NULL)
- *
- * The metadata will be marshalled into 'struct rd_kafka_metadata*' structs.
- *
- * Returns the marshalled metadata, or NULL on parse error.
- *
- * @locality rdkafka main thread
- */
-struct rd_kafka_metadata *
-rd_kafka_parse_Metadata (rd_kafka_broker_t *rkb,
-                         rd_kafka_buf_t *request,
-                         rd_kafka_buf_t *rkbuf) {
-        rd_kafka_t *rk = rkb->rkb_rk;
-        int i, j, k;
-        rd_tmpabuf_t tbuf;
-        struct rd_kafka_metadata *md;
-        size_t rkb_namelen;
-        const int log_decode_errors = LOG_ERR;
-        rd_list_t *missing_topics = NULL;
-        const rd_list_t *requested_topics = request->rkbuf_u.Metadata.topics;
-        int all_topics = request->rkbuf_u.Metadata.all_topics;
-        const char *reason = request->rkbuf_u.Metadata.reason ?
-                request->rkbuf_u.Metadata.reason : "(no reason)";
-        int ApiVersion = request->rkbuf_reqhdr.ApiVersion;
-        rd_kafkap_str_t cluster_id = RD_ZERO_INIT;
-        int32_t controller_id = -1;
-
-        rd_kafka_assert(NULL, thrd_is_current(rk->rk_thread));
-
-        /* Remove topics from missing_topics as they are seen in Metadata. */
-        if (requested_topics)
-                missing_topics = rd_list_copy(requested_topics,
-                                              rd_list_string_copy, NULL);
-
-        rd_kafka_broker_lock(rkb);
-        rkb_namelen = strlen(rkb->rkb_name)+1;
-        /* We assume that the marshalled representation is
-         * no more than 4 times larger than the wire representation. */
-        rd_tmpabuf_new(&tbuf,
-                       sizeof(*md) + rkb_namelen + (rkbuf->rkbuf_totlen * 4),
-                       0/*dont assert on fail*/);
-
-        if (!(md = rd_tmpabuf_alloc(&tbuf, sizeof(*md))))
-                goto err;
-        md->orig_broker_id = rkb->rkb_nodeid;
-        md->orig_broker_name = rd_tmpabuf_write(&tbuf,
-                                                rkb->rkb_name, rkb_namelen);
-        rd_kafka_broker_unlock(rkb);
-
-        /* Read Brokers */
-        rd_kafka_buf_read_i32a(rkbuf, md->broker_cnt);
-        if (md->broker_cnt > RD_KAFKAP_BROKERS_MAX)
-                rd_kafka_buf_parse_fail(rkbuf, "Broker_cnt %i > BROKERS_MAX 
%i",
-                                        md->broker_cnt, RD_KAFKAP_BROKERS_MAX);
-
-        if (!(md->brokers = rd_tmpabuf_alloc(&tbuf, md->broker_cnt *
-                                             sizeof(*md->brokers))))
-                rd_kafka_buf_parse_fail(rkbuf,
-                                        "%d brokers: tmpabuf memory shortage",
-                                        md->broker_cnt);
-
-        for (i = 0 ; i < md->broker_cnt ; i++) {
-                rd_kafka_buf_read_i32a(rkbuf, md->brokers[i].id);
-                rd_kafka_buf_read_str_tmpabuf(rkbuf, &tbuf, 
md->brokers[i].host);
-                rd_kafka_buf_read_i32a(rkbuf, md->brokers[i].port);
-
-                if (ApiVersion >= 1) {
-                        rd_kafkap_str_t rack;
-                        rd_kafka_buf_read_str(rkbuf, &rack);
-                }
-        }
-
-        if (ApiVersion >= 2)
-                rd_kafka_buf_read_str(rkbuf, &cluster_id);
-
-        if (ApiVersion >= 1) {
-                rd_kafka_buf_read_i32(rkbuf, &controller_id);
-                rd_rkb_dbg(rkb, METADATA,
-                           "METADATA", "ClusterId: %.*s, ControllerId: 
%"PRId32,
-                           RD_KAFKAP_STR_PR(&cluster_id), controller_id);
-        }
-
-
-
-        /* Read TopicMetadata */
-        rd_kafka_buf_read_i32a(rkbuf, md->topic_cnt);
-        rd_rkb_dbg(rkb, METADATA, "METADATA", "%i brokers, %i topics",
-                   md->broker_cnt, md->topic_cnt);
-
-        if (md->topic_cnt > RD_KAFKAP_TOPICS_MAX)
-                rd_kafka_buf_parse_fail(rkbuf, "TopicMetadata_cnt %"PRId32
-                                        " > TOPICS_MAX %i",
-                                        md->topic_cnt, RD_KAFKAP_TOPICS_MAX);
-
-        if (!(md->topics = rd_tmpabuf_alloc(&tbuf,
-                                            md->topic_cnt *
-                                            sizeof(*md->topics))))
-                rd_kafka_buf_parse_fail(rkbuf,
-                                        "%d topics: tmpabuf memory shortage",
-                                        md->topic_cnt);
-
-        for (i = 0 ; i < md->topic_cnt ; i++) {
-                rd_kafka_buf_read_i16a(rkbuf, md->topics[i].err);
-                rd_kafka_buf_read_str_tmpabuf(rkbuf, &tbuf, 
md->topics[i].topic);
-                if (ApiVersion >= 1) {
-                        int8_t is_internal;
-                        rd_kafka_buf_read_i8(rkbuf, &is_internal);
-                }
-
-                /* PartitionMetadata */
-                rd_kafka_buf_read_i32a(rkbuf, md->topics[i].partition_cnt);
-                if (md->topics[i].partition_cnt > RD_KAFKAP_PARTITIONS_MAX)
-                        rd_kafka_buf_parse_fail(rkbuf,
-                                                "TopicMetadata[%i]."
-                                                "PartitionMetadata_cnt %i "
-                                                "> PARTITIONS_MAX %i",
-                                                i, md->topics[i].partition_cnt,
-                                                RD_KAFKAP_PARTITIONS_MAX);
-
-                if (!(md->topics[i].partitions =
-                      rd_tmpabuf_alloc(&tbuf,
-                                       md->topics[i].partition_cnt *
-                                       sizeof(*md->topics[i].partitions))))
-                        rd_kafka_buf_parse_fail(rkbuf,
-                                                "%s: %d partitions: "
-                                                "tmpabuf memory shortage",
-                                                md->topics[i].topic,
-                                                md->topics[i].partition_cnt);
-
-                for (j = 0 ; j < md->topics[i].partition_cnt ; j++) {
-                        rd_kafka_buf_read_i16a(rkbuf, 
md->topics[i].partitions[j].err);
-                        rd_kafka_buf_read_i32a(rkbuf, 
md->topics[i].partitions[j].id);
-                        rd_kafka_buf_read_i32a(rkbuf, 
md->topics[i].partitions[j].leader);
-
-                        /* Replicas */
-                        rd_kafka_buf_read_i32a(rkbuf, 
md->topics[i].partitions[j].replica_cnt);
-                        if (md->topics[i].partitions[j].replica_cnt >
-                            RD_KAFKAP_BROKERS_MAX)
-                                rd_kafka_buf_parse_fail(rkbuf,
-                                                        "TopicMetadata[%i]."
-                                                        
"PartitionMetadata[%i]."
-                                                        "Replica_cnt "
-                                                        "%i > BROKERS_MAX %i",
-                                                        i, j,
-                                                        md->topics[i].
-                                                        partitions[j].
-                                                        replica_cnt,
-                                                        RD_KAFKAP_BROKERS_MAX);
-
-                        if (!(md->topics[i].partitions[j].replicas =
-                              rd_tmpabuf_alloc(&tbuf,
-                                               md->topics[i].
-                                               partitions[j].replica_cnt *
-                                               sizeof(*md->topics[i].
-                                                      
partitions[j].replicas))))
-                                rd_kafka_buf_parse_fail(
-                                        rkbuf,
-                                        "%s [%"PRId32"]: %d replicas: "
-                                        "tmpabuf memory shortage",
-                                        md->topics[i].topic,
-                                        md->topics[i].partitions[j].id,
-                                        
md->topics[i].partitions[j].replica_cnt);
-
-
-                        for (k = 0 ;
-                             k < md->topics[i].partitions[j].replica_cnt; k++)
-                                rd_kafka_buf_read_i32a(rkbuf, 
md->topics[i].partitions[j].
-                                           replicas[k]);
-
-                        /* Isrs */
-                        rd_kafka_buf_read_i32a(rkbuf, 
md->topics[i].partitions[j].isr_cnt);
-                        if (md->topics[i].partitions[j].isr_cnt >
-                            RD_KAFKAP_BROKERS_MAX)
-                                rd_kafka_buf_parse_fail(rkbuf,
-                                                        "TopicMetadata[%i]."
-                                                        
"PartitionMetadata[%i]."
-                                                        "Isr_cnt "
-                                                        "%i > BROKERS_MAX %i",
-                                                        i, j,
-                                                        md->topics[i].
-                                                        partitions[j].isr_cnt,
-                                                        RD_KAFKAP_BROKERS_MAX);
-
-                        if (!(md->topics[i].partitions[j].isrs =
-                              rd_tmpabuf_alloc(&tbuf,
-                                               md->topics[i].
-                                               partitions[j].isr_cnt *
-                                               sizeof(*md->topics[i].
-                                                      partitions[j].isrs))))
-                                rd_kafka_buf_parse_fail(
-                                        rkbuf,
-                                        "%s [%"PRId32"]: %d isrs: "
-                                        "tmpabuf memory shortage",
-                                        md->topics[i].topic,
-                                        md->topics[i].partitions[j].id,
-                                        md->topics[i].partitions[j].isr_cnt);
-
-
-                        for (k = 0 ;
-                             k < md->topics[i].partitions[j].isr_cnt; k++)
-                                rd_kafka_buf_read_i32a(rkbuf, md->topics[i].
-                                                       partitions[j].isrs[k]);
-
-                }
-        }
-
-        /* Entire Metadata response now parsed without errors:
-         * update our internal state according to the response. */
-
-        /* Avoid metadata updates when we're terminating. */
-        if (rd_kafka_terminating(rkb->rkb_rk))
-                goto done;
-
-        if (md->broker_cnt == 0 && md->topic_cnt == 0) {
-                rd_rkb_dbg(rkb, METADATA, "METADATA",
-                           "No brokers or topics in metadata: retrying");
-                goto err;
-        }
-
-        /* Update our list of brokers. */
-        for (i = 0 ; i < md->broker_cnt ; i++) {
-                rd_rkb_dbg(rkb, METADATA, "METADATA",
-                           "  Broker #%i/%i: %s:%i NodeId %"PRId32,
-                           i, md->broker_cnt,
-                           md->brokers[i].host,
-                           md->brokers[i].port,
-                           md->brokers[i].id);
-                rd_kafka_broker_update(rkb->rkb_rk, rkb->rkb_proto,
-                                       &md->brokers[i]);
-        }
-
-        /* Update partition count and leader for each topic we know about */
-        for (i = 0 ; i < md->topic_cnt ; i++) {
-                rd_kafka_metadata_topic_t *mdt = &md->topics[i];
-                rd_rkb_dbg(rkb, METADATA, "METADATA",
-                           "  Topic #%i/%i: %s with %i partitions%s%s",
-                           i, md->topic_cnt, mdt->topic,
-                           mdt->partition_cnt,
-                           mdt->err ? ": " : "",
-                           mdt->err ? rd_kafka_err2str(mdt->err) : "");
-
-                /* Ignore topics in blacklist */
-                if (rkb->rkb_rk->rk_conf.topic_blacklist &&
-                    
rd_kafka_pattern_match(rkb->rkb_rk->rk_conf.topic_blacklist,
-                                           mdt->topic)) {
-                        rd_rkb_dbg(rkb, TOPIC, "BLACKLIST",
-                                   "Ignoring blacklisted topic \"%s\" "
-                                   "in metadata", mdt->topic);
-                        continue;
-                }
-
-                /* Ignore metadata completely for temporary errors. (issue 
#513)
-                 *   LEADER_NOT_AVAILABLE: Broker is rebalancing
-                 */
-                if (mdt->err == RD_KAFKA_RESP_ERR_LEADER_NOT_AVAILABLE &&
-                    mdt->partition_cnt == 0) {
-                        rd_rkb_dbg(rkb, TOPIC, "METADATA",
-                                   "Temporary error in metadata reply for "
-                                   "topic %s (PartCnt %i): %s: ignoring",
-                                   mdt->topic, mdt->partition_cnt,
-                                   rd_kafka_err2str(mdt->err));
-                        rd_list_free_cb(missing_topics,
-                                        rd_list_remove_cmp(missing_topics,
-                                                           mdt->topic,
-                                                           (void *)strcmp));
-                        continue;
-                }
-
-
-                /* Update local topic & partition state based on metadata */
-                rd_kafka_topic_metadata_update2(rkb, mdt);
-
-                if (requested_topics) {
-                        rd_list_free_cb(missing_topics,
-                                        rd_list_remove_cmp(missing_topics,
-                                                           mdt->topic,
-                                                           (void*)strcmp));
-                        if (!all_topics) {
-                                rd_kafka_wrlock(rk);
-                                rd_kafka_metadata_cache_topic_update(rk, mdt);
-                                rd_kafka_wrunlock(rk);
-                        }
-                }
-        }
-
-
-        /* Requested topics not seen in metadata? Propogate to topic code. */
-        if (missing_topics) {
-                char *topic;
-                rd_rkb_dbg(rkb, TOPIC, "METADATA",
-                           "%d/%d requested topic(s) seen in metadata",
-                           rd_list_cnt(requested_topics) -
-                           rd_list_cnt(missing_topics),
-                           rd_list_cnt(requested_topics));
-                for (i = 0 ; i < rd_list_cnt(missing_topics) ; i++)
-                        rd_rkb_dbg(rkb, TOPIC, "METADATA", "wanted %s",
-                                   (char *)(missing_topics->rl_elems[i]));
-                RD_LIST_FOREACH(topic, missing_topics, i) {
-                        shptr_rd_kafka_itopic_t *s_rkt;
-
-                        s_rkt = rd_kafka_topic_find(rkb->rkb_rk, topic, 
1/*lock*/);
-                        if (s_rkt) {
-                                rd_kafka_topic_metadata_none(
-                                        rd_kafka_topic_s2i(s_rkt));
-                                rd_kafka_topic_destroy0(s_rkt);
-                        }
-                }
-        }
-
-
-        rd_kafka_wrlock(rkb->rkb_rk);
-        rkb->rkb_rk->rk_ts_metadata = rd_clock();
-
-        /* Update cached cluster id. */
-        if (RD_KAFKAP_STR_LEN(&cluster_id) > 0 &&
-            (!rkb->rkb_rk->rk_clusterid ||
-             rd_kafkap_str_cmp_str(&cluster_id, rkb->rkb_rk->rk_clusterid))) {
-                rd_rkb_dbg(rkb, BROKER|RD_KAFKA_DBG_GENERIC, "CLUSTERID",
-                           "ClusterId update \"%s\" -> \"%.*s\"",
-                           rkb->rkb_rk->rk_clusterid ?
-                           rkb->rkb_rk->rk_clusterid : "",
-                           RD_KAFKAP_STR_PR(&cluster_id));
-                if (rkb->rkb_rk->rk_clusterid)
-                        rd_free(rkb->rkb_rk->rk_clusterid);
-                rkb->rkb_rk->rk_clusterid = RD_KAFKAP_STR_DUP(&cluster_id);
-        }
-
-        if (all_topics) {
-                rd_kafka_metadata_cache_update(rkb->rkb_rk,
-                                               md, 1/*abs update*/);
-
-                if (rkb->rkb_rk->rk_full_metadata)
-                        
rd_kafka_metadata_destroy(rkb->rkb_rk->rk_full_metadata);
-                rkb->rkb_rk->rk_full_metadata =
-                        rd_kafka_metadata_copy(md, tbuf.of);
-                rkb->rkb_rk->rk_ts_full_metadata = rkb->rkb_rk->rk_ts_metadata;
-                rd_rkb_dbg(rkb, METADATA, "METADATA",
-                           "Caching full metadata with "
-                           "%d broker(s) and %d topic(s): %s",
-                           md->broker_cnt, md->topic_cnt, reason);
-        } else {
-                rd_kafka_metadata_cache_expiry_start(rk);
-        }
-
-        /* Remove cache hints for the originally requested topics. */
-        if (requested_topics)
-                rd_kafka_metadata_cache_purge_hints(rk, requested_topics);
-
-        rd_kafka_wrunlock(rkb->rkb_rk);
-
-        /* Check if cgrp effective subscription is affected by
-         * new metadata. */
-        if (rkb->rkb_rk->rk_cgrp)
-                rd_kafka_cgrp_metadata_update_check(
-                        rkb->rkb_rk->rk_cgrp, 1/*do join*/);
-
-
-
-done:
-        if (missing_topics)
-                rd_list_destroy(missing_topics);
-
-        /* This metadata request was triggered by someone wanting
-         * the metadata information back as a reply, so send that reply now.
-         * In this case we must not rd_free the metadata memory here,
-         * the requestee will do.
-         * The tbuf is explicitly not destroyed as we return its memory
-         * to the caller. */
-        return md;
-
- err_parse:
-err:
-        if (requested_topics) {
-                /* Failed requests shall purge cache hints for
-                 * the requested topics. */
-                rd_kafka_wrlock(rkb->rkb_rk);
-                rd_kafka_metadata_cache_purge_hints(rk, requested_topics);
-                rd_kafka_wrunlock(rkb->rkb_rk);
-        }
-
-        if (missing_topics)
-                rd_list_destroy(missing_topics);
-
-        rd_tmpabuf_destroy(&tbuf);
-        return NULL;
-}
-
-
-/**
- * @brief Add all topics in current cached full metadata
- *        to \p tinfos (rd_kafka_topic_info_t *)
- *        that matches the topics in \p match
- *
- * @returns the number of topics matched and added to \p list
- *
- * @locks none
- * @locality any
- */
-size_t
-rd_kafka_metadata_topic_match (rd_kafka_t *rk, rd_list_t *tinfos,
-                               const rd_kafka_topic_partition_list_t *match) {
-        int ti;
-        size_t cnt = 0;
-        const struct rd_kafka_metadata *metadata;
-
-
-        rd_kafka_rdlock(rk);
-        metadata = rk->rk_full_metadata;
-        if (!metadata) {
-                rd_kafka_rdunlock(rk);
-                return 0;
-        }
-
-        /* For each topic in the cluster, scan through the match list
-         * to find matching topic. */
-        for (ti = 0 ; ti < metadata->topic_cnt ; ti++) {
-                const char *topic = metadata->topics[ti].topic;
-                int i;
-
-                /* Ignore topics in blacklist */
-                if (rk->rk_conf.topic_blacklist &&
-                    rd_kafka_pattern_match(rk->rk_conf.topic_blacklist, topic))
-                        continue;
-
-                /* Scan for matches */
-                for (i = 0 ; i < match->cnt ; i++) {
-                        if (!rd_kafka_topic_match(rk,
-                                                  match->elems[i].topic, 
topic))
-                                continue;
-
-                        if (metadata->topics[ti].err)
-                                continue; /* Skip errored topics */
-
-                        rd_list_add(tinfos,
-                                    rd_kafka_topic_info_new(
-                                            topic,
-                                            
metadata->topics[ti].partition_cnt));
-                        cnt++;
-                }
-        }
-        rd_kafka_rdunlock(rk);
-
-        return cnt;
-}
-
-
-/**
- * @brief Add all topics in \p match that matches cached metadata.
- * @remark MUST NOT be used with wildcard topics,
- *         see rd_kafka_metadata_topic_match() for that.
- *
- * @returns the number of topics matched and added to \p tinfos
- * @locks none
- */
-size_t
-rd_kafka_metadata_topic_filter (rd_kafka_t *rk, rd_list_t *tinfos,
-                               const rd_kafka_topic_partition_list_t *match) {
-        int i;
-        size_t cnt = 0;
-
-        rd_kafka_rdlock(rk);
-        /* For each topic in match, look up the topic in the cache. */
-        for (i = 0 ; i < match->cnt ; i++) {
-                const char *topic = match->elems[i].topic;
-                const rd_kafka_metadata_topic_t *mtopic;
-
-                /* Ignore topics in blacklist */
-                if (rk->rk_conf.topic_blacklist &&
-                    rd_kafka_pattern_match(rk->rk_conf.topic_blacklist, topic))
-                        continue;
-
-                mtopic = rd_kafka_metadata_cache_topic_get(rk, topic,
-                                                           1/*valid*/);
-                if (mtopic && !mtopic->err) {
-                        rd_list_add(tinfos,
-                                    rd_kafka_topic_info_new(
-                                            topic, mtopic->partition_cnt));
-
-                        cnt++;
-                }
-        }
-        rd_kafka_rdunlock(rk);
-
-        return cnt;
-}
-
-
-void rd_kafka_metadata_log (rd_kafka_t *rk, const char *fac,
-                            const struct rd_kafka_metadata *md) {
-        int i;
-
-        rd_kafka_dbg(rk, METADATA, fac,
-                     "Metadata with %d broker(s) and %d topic(s):",
-                     md->broker_cnt, md->topic_cnt);
-
-        for (i = 0 ; i < md->broker_cnt ; i++) {
-                rd_kafka_dbg(rk, METADATA, fac,
-                             "  Broker #%i/%i: %s:%i NodeId %"PRId32,
-                             i, md->broker_cnt,
-                             md->brokers[i].host,
-                             md->brokers[i].port,
-                             md->brokers[i].id);
-        }
-
-        for (i = 0 ; i < md->topic_cnt ; i++) {
-                rd_kafka_dbg(rk, METADATA, fac,
-                             "  Topic #%i/%i: %s with %i partitions%s%s",
-                             i, md->topic_cnt, md->topics[i].topic,
-                             md->topics[i].partition_cnt,
-                             md->topics[i].err ? ": " : "",
-                             md->topics[i].err ?
-                             rd_kafka_err2str(md->topics[i].err) : "");
-        }
-}
-
-
-
-
-/**
- * @brief Refresh metadata for \p topics
- *
- * @param rk: used to look up usable broker if \p rkb is NULL.
- * @param rkb: use this broker, unless NULL then any usable broker from \p rk
- * @param force: force refresh even if topics are up-to-date in cache
- *
- * @returns an error code
- *
- * @locality any
- * @locks none
- */
-rd_kafka_resp_err_t
-rd_kafka_metadata_refresh_topics (rd_kafka_t *rk, rd_kafka_broker_t *rkb,
-                                  const rd_list_t *topics, int force,
-                                  const char *reason) {
-        rd_list_t q_topics;
-        int destroy_rkb = 0;
-
-        if (!rk)
-                rk = rkb->rkb_rk;
-
-        rd_kafka_wrlock(rk);
-
-        if (!rkb) {
-                if (!(rkb = rd_kafka_broker_any_usable(rk, RD_POLL_NOWAIT, 
0))){
-                        rd_kafka_wrunlock(rk);
-                        rd_kafka_dbg(rk, METADATA, "METADATA",
-                                     "Skipping metadata refresh of %d 
topic(s):"
-                                     " no usable brokers",
-                                     rd_list_cnt(topics));
-                        return RD_KAFKA_RESP_ERR__TRANSPORT;
-                }
-                destroy_rkb = 1;
-        }
-
-        rd_list_init(&q_topics, rd_list_cnt(topics), rd_free);
-
-        if (!force) {
-
-                /* Hint cache of upcoming MetadataRequest and filter
-                 * out any topics that are already being requested.
-                 * q_topics will contain remaining topics to query. */
-                rd_kafka_metadata_cache_hint(rk, topics, &q_topics,
-                                             0/*dont replace*/);
-                rd_kafka_wrunlock(rk);
-
-                if (rd_list_cnt(&q_topics) == 0) {
-                        /* No topics need new query. */
-                        rd_kafka_dbg(rk, METADATA, "METADATA",
-                                     "Skipping metadata refresh of "
-                                     "%d topic(s): %s: "
-                                     "already being requested",
-                                     rd_list_cnt(topics), reason);
-                        rd_list_destroy(&q_topics);
-                        if (destroy_rkb)
-                                rd_kafka_broker_destroy(rkb);
-                        return RD_KAFKA_RESP_ERR_NO_ERROR;
-                }
-
-        } else {
-                rd_kafka_wrunlock(rk);
-                rd_list_copy_to(&q_topics, topics, rd_list_string_copy, NULL);
-        }
-
-        rd_kafka_dbg(rk, METADATA, "METADATA",
-                     "Requesting metadata for %d/%d topics: %s",
-                     rd_list_cnt(&q_topics), rd_list_cnt(topics), reason);
-
-        rd_kafka_MetadataRequest(rkb, &q_topics, reason, NULL);
-
-        rd_list_destroy(&q_topics);
-
-        if (destroy_rkb)
-                rd_kafka_broker_destroy(rkb);
-
-        return RD_KAFKA_RESP_ERR_NO_ERROR;
-}
-
-
-/**
- * @brief Refresh metadata for known topics
- *
- * @param rk: used to look up usable broker if \p rkb is NULL.
- * @param rkb: use this broker, unless NULL then any usable broker from \p rk
- * @param force: refresh even if cache is up-to-date
- *
- * @returns an error code (__UNKNOWN_TOPIC if there are no local topics)
- *
- * @locality any
- * @locks none
- */
-rd_kafka_resp_err_t
-rd_kafka_metadata_refresh_known_topics (rd_kafka_t *rk, rd_kafka_broker_t *rkb,
-                                        int force, const char *reason) {
-        rd_list_t topics;
-        rd_kafka_resp_err_t err;
-
-        if (!rk)
-                rk = rkb->rkb_rk;
-
-        rd_list_init(&topics, 8, rd_free);
-        rd_kafka_local_topics_to_list(rk, &topics);
-
-        if (rd_list_cnt(&topics) == 0)
-                err = RD_KAFKA_RESP_ERR__UNKNOWN_TOPIC;
-        else
-                err = rd_kafka_metadata_refresh_topics(rk, rkb,
-                                                       &topics, force, reason);
-
-        rd_list_destroy(&topics);
-
-        return err;
-}
-
-
-/**
- * @brief Refresh broker list by metadata.
- *
- * Attempts to use sparse metadata request if possible, else falls back
- * on a full metadata request. (NOTE: sparse not implemented, KIP-4)
- *
- * @param rk: used to look up usable broker if \p rkb is NULL.
- * @param rkb: use this broker, unless NULL then any usable broker from \p rk
- *
- * @returns an error code
- *
- * @locality any
- * @locks none
- */
-rd_kafka_resp_err_t
-rd_kafka_metadata_refresh_brokers (rd_kafka_t *rk, rd_kafka_broker_t *rkb,
-                                   const char *reason) {
-        return rd_kafka_metadata_request(rk, rkb, NULL /*brokers only*/,
-                                         reason, NULL);
-}
-
-
-
-/**
- * @brief Refresh metadata for all topics in cluster.
- *        This is a full metadata request which might be taxing on the
- *        broker if the cluster has many topics.
- *
- * @locality any
- * @locks none
- */
-rd_kafka_resp_err_t
-rd_kafka_metadata_refresh_all (rd_kafka_t *rk, rd_kafka_broker_t *rkb,
-                               const char *reason) {
-        int destroy_rkb = 0;
-        rd_list_t topics;
-
-        if (!rk)
-                rk = rkb->rkb_rk;
-
-        if (!rkb) {
-                if (!(rkb = rd_kafka_broker_any_usable(rk, RD_POLL_NOWAIT, 1)))
-                        return RD_KAFKA_RESP_ERR__TRANSPORT;
-                destroy_rkb = 1;
-        }
-
-        rd_list_init(&topics, 0, NULL); /* empty list = all topics */
-        rd_kafka_MetadataRequest(rkb, &topics, reason, NULL);
-        rd_list_destroy(&topics);
-
-        if (destroy_rkb)
-                rd_kafka_broker_destroy(rkb);
-
-        return RD_KAFKA_RESP_ERR_NO_ERROR;
-}
-
-
-/**
-
- * @brief Lower-level Metadata request that takes a callback (with replyq set)
- *        which will be triggered after parsing is complete.
- *
- * @locks none
- * @locality any
- */
-rd_kafka_resp_err_t
-rd_kafka_metadata_request (rd_kafka_t *rk, rd_kafka_broker_t *rkb,
-                           const rd_list_t *topics,
-                           const char *reason, rd_kafka_op_t *rko) {
-        int destroy_rkb = 0;
-
-        if (!rkb) {
-                if (!(rkb = rd_kafka_broker_any_usable(rk, RD_POLL_NOWAIT, 1)))
-                        return RD_KAFKA_RESP_ERR__TRANSPORT;
-                destroy_rkb = 1;
-        }
-
-        rd_kafka_MetadataRequest(rkb, topics, reason, rko);
-
-        if (destroy_rkb)
-                rd_kafka_broker_destroy(rkb);
-
-        return RD_KAFKA_RESP_ERR_NO_ERROR;
-}
-
-
-/**
- * @brief Query timer callback to trigger refresh for topics
- *        that are missing their leaders.
- *
- * @locks none
- * @locality rdkafka main thread
- */
-static void rd_kafka_metadata_leader_query_tmr_cb (rd_kafka_timers_t *rkts,
-                                                   void *arg) {
-        rd_kafka_t *rk = rkts->rkts_rk;
-        rd_kafka_timer_t *rtmr = &rk->rk_metadata_cache.rkmc_query_tmr;
-        rd_kafka_itopic_t *rkt;
-        rd_list_t topics;
-
-        rd_kafka_wrlock(rk);
-        rd_list_init(&topics, rk->rk_topic_cnt, rd_free);
-
-        TAILQ_FOREACH(rkt, &rk->rk_topics, rkt_link) {
-                int i, no_leader = 0;
-                rd_kafka_topic_rdlock(rkt);
-
-                if (rkt->rkt_state == RD_KAFKA_TOPIC_S_NOTEXISTS) {
-                        /* Skip topics that are known to not exist. */
-                        rd_kafka_topic_rdunlock(rkt);
-                        continue;
-                }
-
-                no_leader = rkt->rkt_flags & RD_KAFKA_TOPIC_F_LEADER_UNAVAIL;
-
-                /* Check if any partitions are missing their leaders. */
-                for (i = 0 ; !no_leader && i < rkt->rkt_partition_cnt ; i++) {
-                        rd_kafka_toppar_t *rktp =
-                                rd_kafka_toppar_s2i(rkt->rkt_p[i]);
-                        rd_kafka_toppar_lock(rktp);
-                        no_leader = !rktp->rktp_leader &&
-                                !rktp->rktp_next_leader;
-                        rd_kafka_toppar_unlock(rktp);
-                }
-
-                if (no_leader || rkt->rkt_partition_cnt == 0)
-                        rd_list_add(&topics, rd_strdup(rkt->rkt_topic->str));
-
-                rd_kafka_topic_rdunlock(rkt);
-        }
-
-        rd_kafka_wrunlock(rk);
-
-        if (rd_list_cnt(&topics) == 0) {
-                /* No leader-less topics+partitions, stop the timer. */
-                rd_kafka_timer_stop(rkts, rtmr, 1/*lock*/);
-        } else {
-                rd_kafka_metadata_refresh_topics(rk, NULL, &topics, 1/*force*/,
-                                                 "partition leader query");
-                /* Back off next query exponentially until we reach
-                 * the standard query interval - then stop the timer
-                 * since the intervalled querier will do the job for us. */
-                if (rk->rk_conf.metadata_refresh_interval_ms > 0 &&
-                    rtmr->rtmr_interval * 2 / 1000 >=
-                    rk->rk_conf.metadata_refresh_interval_ms)
-                        rd_kafka_timer_stop(rkts, rtmr, 1/*lock*/);
-                else
-                        rd_kafka_timer_backoff(rkts, rtmr,
-                                               (int)rtmr->rtmr_interval);
-        }
-
-        rd_list_destroy(&topics);
-}
-
-
-
-/**
- * @brief Trigger fast leader query to quickly pick up on leader changes.
- *        The fast leader query is a quick query followed by later queries at
- *        exponentially increased intervals until no topics are missing
- *        leaders.
- *
- * @locks none
- * @locality any
- */
-void rd_kafka_metadata_fast_leader_query (rd_kafka_t *rk) {
-        rd_ts_t next;
-
-        /* Restart the timer if it will speed things up. */
-        next = rd_kafka_timer_next(&rk->rk_timers,
-                                   &rk->rk_metadata_cache.rkmc_query_tmr,
-                                   1/*lock*/);
-        if (next == -1 /* not started */ ||
-            next > rk->rk_conf.metadata_refresh_fast_interval_ms*1000) {
-                rd_kafka_dbg(rk, METADATA|RD_KAFKA_DBG_TOPIC, "FASTQUERY",
-                             "Starting fast leader query");
-                rd_kafka_timer_start(&rk->rk_timers,
-                                     &rk->rk_metadata_cache.rkmc_query_tmr,
-                                     rk->rk_conf.
-                                     metadata_refresh_fast_interval_ms*1000,
-                                     rd_kafka_metadata_leader_query_tmr_cb,
-                                     NULL);
-        }
-}

http://git-wip-us.apache.org/repos/asf/nifi-minifi-cpp/blob/7528d23e/thirdparty/librdkafka-0.11.1/src/rdkafka_metadata.h
----------------------------------------------------------------------
diff --git a/thirdparty/librdkafka-0.11.1/src/rdkafka_metadata.h 
b/thirdparty/librdkafka-0.11.1/src/rdkafka_metadata.h
deleted file mode 100644
index a0b77e1..0000000
--- a/thirdparty/librdkafka-0.11.1/src/rdkafka_metadata.h
+++ /dev/null
@@ -1,157 +0,0 @@
-/*
- * librdkafka - Apache Kafka C library
- *
- * Copyright (c) 2012-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.
- */
-
-#pragma once
-
-#include "rdavl.h"
-
-struct rd_kafka_metadata *
-rd_kafka_parse_Metadata (rd_kafka_broker_t *rkb,
-                         rd_kafka_buf_t *request, rd_kafka_buf_t *rkbuf);
-
-struct rd_kafka_metadata *
-rd_kafka_metadata_copy (const struct rd_kafka_metadata *md, size_t size);
-
-size_t
-rd_kafka_metadata_topic_match (rd_kafka_t *rk, rd_list_t *tinfos,
-                               const rd_kafka_topic_partition_list_t *match);
-size_t
-rd_kafka_metadata_topic_filter (rd_kafka_t *rk, rd_list_t *tinfos,
-                                const rd_kafka_topic_partition_list_t *match);
-
-void rd_kafka_metadata_log (rd_kafka_t *rk, const char *fac,
-                            const struct rd_kafka_metadata *md);
-
-
-
-rd_kafka_resp_err_t
-rd_kafka_metadata_refresh_topics (rd_kafka_t *rk, rd_kafka_broker_t *rkb,
-                                  const rd_list_t *topics, int force,
-                                  const char *reason);
-rd_kafka_resp_err_t
-rd_kafka_metadata_refresh_known_topics (rd_kafka_t *rk, rd_kafka_broker_t *rkb,
-                                        int force, const char *reason);
-rd_kafka_resp_err_t
-rd_kafka_metadata_refresh_brokers (rd_kafka_t *rk, rd_kafka_broker_t *rkb,
-                                   const char *reason);
-rd_kafka_resp_err_t
-rd_kafka_metadata_refresh_all (rd_kafka_t *rk, rd_kafka_broker_t *rkb,
-                               const char *reason);
-
-rd_kafka_resp_err_t
-rd_kafka_metadata_request (rd_kafka_t *rk, rd_kafka_broker_t *rkb,
-                           const rd_list_t *topics,
-                           const char *reason, rd_kafka_op_t *rko);
-
-
-/**
- * @{
- *
- * @brief Metadata cache
- */
-
-struct rd_kafka_metadata_cache_entry {
-        rd_avl_node_t rkmce_avlnode;             /* rkmc_avl */
-        TAILQ_ENTRY(rd_kafka_metadata_cache_entry) rkmce_link; /* rkmc_expiry 
*/
-        rd_ts_t rkmce_ts_expires;                /* Expire time */
-        rd_ts_t rkmce_ts_insert;                 /* Insert time */
-        rd_kafka_metadata_topic_t rkmce_mtopic;  /* Cached topic metadata */
-        /* rkmce_partitions memory points here. */
-};
-
-#define RD_KAFKA_METADATA_CACHE_VALID(rkmce) \
-        ((rkmce)->rkmce_mtopic.err != RD_KAFKA_RESP_ERR__WAIT_CACHE)
-
-struct rd_kafka_metadata_cache {
-        rd_avl_t         rkmc_avl;
-        TAILQ_HEAD(, rd_kafka_metadata_cache_entry) rkmc_expiry;
-        rd_kafka_timer_t rkmc_expiry_tmr;
-        int              rkmc_cnt;
-
-        /* Protected by full_lock: */
-        mtx_t            rkmc_full_lock;
-        int              rkmc_full_topics_sent; /* Full MetadataRequest for
-                                                 * all topics has been sent,
-                                                 * awaiting response. */
-        int              rkmc_full_brokers_sent; /* Full MetadataRequest for
-                                                  * all brokers (but not 
topics)
-                                                  * has been sent,
-                                                  * awaiting response. */
-
-        rd_kafka_timer_t rkmc_query_tmr; /* Query timer for topic's without
-                                          * leaders. */
-        cnd_t            rkmc_cnd;       /* cache_wait_change() cond. */
-        mtx_t            rkmc_cnd_lock;  /* lock for rkmc_cnd */
-};
-
-
-
-void rd_kafka_metadata_cache_expiry_start (rd_kafka_t *rk);
-void
-rd_kafka_metadata_cache_topic_update (rd_kafka_t *rk,
-                                      const rd_kafka_metadata_topic_t *mdt);
-void rd_kafka_metadata_cache_update (rd_kafka_t *rk,
-                                     const rd_kafka_metadata_t *md,
-                                     int abs_update);
-struct rd_kafka_metadata_cache_entry *
-rd_kafka_metadata_cache_find (rd_kafka_t *rk, const char *topic, int valid);
-void rd_kafka_metadata_cache_purge_hints (rd_kafka_t *rk,
-                                          const rd_list_t *topics);
-int rd_kafka_metadata_cache_hint (rd_kafka_t *rk,
-                                  const rd_list_t *topics, rd_list_t *dst,
-                                  int replace);
-int rd_kafka_metadata_cache_hint_rktparlist (
-        rd_kafka_t *rk,
-        const rd_kafka_topic_partition_list_t *rktparlist,
-        rd_list_t *dst,
-        int replace);
-
-const rd_kafka_metadata_topic_t *
-rd_kafka_metadata_cache_topic_get (rd_kafka_t *rk, const char *topic,
-                                   int valid);
-int rd_kafka_metadata_cache_topic_partition_get (
-        rd_kafka_t *rk,
-        const rd_kafka_metadata_topic_t **mtopicp,
-        const rd_kafka_metadata_partition_t **mpartp,
-        const char *topic, int32_t partition, int valid);
-
-int rd_kafka_metadata_cache_topics_count_exists (rd_kafka_t *rk,
-                                                 const rd_list_t *topics,
-                                                 int *metadata_agep);
-int rd_kafka_metadata_cache_topics_filter_hinted (rd_kafka_t *rk,
-                                                  rd_list_t *dst,
-                                                  const rd_list_t *src);
-
-void rd_kafka_metadata_fast_leader_query (rd_kafka_t *rk);
-
-void rd_kafka_metadata_cache_init (rd_kafka_t *rk);
-void rd_kafka_metadata_cache_destroy (rd_kafka_t *rk);
-int  rd_kafka_metadata_cache_wait_change (rd_kafka_t *rk, int timeout_ms);
-void rd_kafka_metadata_cache_dump (FILE *fp, rd_kafka_t *rk);
-
-/**@}*/

Reply via email to