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); - -/**@}*/
