/dports/sysutils/vector/vector-0.10.0/cargo-crates/rdkafka-sys-2.0.0+1.4.2/librdkafka/src/ |
H A D | rdkafka_cgrp.c | 178 rd_kafka_assert(rkcg->rkcg_rk, !rkcg->rkcg_assignment); in rd_kafka_cgrp_destroy_final() 179 rd_kafka_assert(rkcg->rkcg_rk, !rkcg->rkcg_subscription); in rd_kafka_cgrp_destroy_final() 219 rkcg = rd_calloc(1, sizeof(*rkcg)); in rd_kafka_cgrp_new() 231 rkcg->rkcg_ops->rkq_opaque = rkcg; in rd_kafka_cgrp_new() 724 rkcg, in rd_kafka_rebalance_op() 1479 rkcg->rkcg_assignment ? rkcg->rkcg_assignment->cnt : 0, in rd_kafka_cgrp_handle_Heartbeat() 1936 (rkcg->rkcg_assignment ? rkcg->rkcg_assignment->cnt : 0)); in rd_kafka_cgrp_partitions_fetch_start0() 2299 rkcg->rkcg_coord, rkcg, offsets, in rd_kafka_cgrp_offsets_commit() 2410 rkcg, rkcg->rkcg_assignment, 0); in rd_kafka_cgrp_unassign_done() 3386 rkcg, in rd_kafka_cgrp_serve() [all …]
|
H A D | rdkafka_cgrp.h | 258 #define rd_kafka_cgrp_lock(rkcg) mtx_lock(&(rkcg)->rkcg_lock) argument 259 #define rd_kafka_cgrp_unlock(rkcg) mtx_unlock(&(rkcg)->rkcg_lock) argument 262 #define RD_KAFKA_CGRP_BROKER_IS_COORD(rkcg,rkb) \ argument 263 ((rkcg)->rkcg_coord_id != -1 && \ 264 (rkcg)->rkcg_coord_id == (rkb)->rkb_nodeid) 269 #define RD_KAFKA_CGRP_IS_STATIC_MEMBER(rkcg) \ argument 270 !RD_KAFKAP_STR_IS_NULL((rkcg)->rkcg_group_instance_id) 275 void rd_kafka_cgrp_destroy_final (rd_kafka_cgrp_t *rkcg); 279 void rd_kafka_cgrp_serve (rd_kafka_cgrp_t *rkcg); 297 void rd_kafka_cgrp_handle_SyncGroup (rd_kafka_cgrp_t *rkcg, [all …]
|
H A D | rdkafka_subscription.c | 40 rd_kafka_cgrp_t *rkcg; in rd_kafka_unsubscribe() local 42 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_unsubscribe() 45 return rd_kafka_op_err_destroy(rd_kafka_op_req2(rkcg->rkcg_ops, in rd_kafka_unsubscribe() 76 rd_kafka_cgrp_t *rkcg; in rd_kafka_subscribe() local 78 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_subscribe() 99 rd_kafka_cgrp_t *rkcg; in rd_kafka_assign() local 101 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_assign() 120 rd_kafka_cgrp_t *rkcg; in rd_kafka_assignment() local 122 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_assignment() 148 rd_kafka_cgrp_t *rkcg; in rd_kafka_subscription() local [all …]
|
H A D | rdkafka_assignor.c | 164 rd_kafka_cgrp_t *rkcg, in rd_kafka_member_subscription_match() argument 220 rd_kafka_member_subscriptions_map (rd_kafka_cgrp_t *rkcg, in rd_kafka_member_subscriptions_map() argument 238 if (rkcg->rkcg_rk->rk_conf.topic_blacklist && in rd_kafka_member_subscriptions_map() 239 rd_kafka_pattern_match(rkcg->rkcg_rk->rk_conf. in rd_kafka_member_subscriptions_map() 283 rd_kafka_assignor_run (rd_kafka_cgrp_t *rkcg, in rd_kafka_assignor_run() argument 309 if (rkcg->rkcg_rk->rk_conf.debug & RD_KAFKA_DBG_CGRP) { in rd_kafka_assignor_run() 310 rd_kafka_dbg(rkcg->rkcg_rk, CGRP, "ASSIGN", in rd_kafka_assignor_run() 340 err = rkas->rkas_assign_cb(rkcg->rkcg_rk, in rd_kafka_assignor_run() 341 rkcg->rkcg_member_id->str, in rd_kafka_assignor_run() 351 rd_kafka_dbg(rkcg->rkcg_rk, CGRP, "ASSIGN", in rd_kafka_assignor_run() [all …]
|
H A D | rdkafka.c | 2947 rd_kafka_cgrp_t *rkcg; in rd_kafka_poll_set_consumer() local 2949 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_poll_set_consumer() 2961 rd_kafka_cgrp_t *rkcg; in rd_kafka_consumer_poll() local 2974 rd_kafka_cgrp_t *rkcg; in rd_kafka_consumer_close() local 2979 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_consumer_close() 2996 rd_kafka_q_fwd_set(rkcg->rkcg_q, rkq); in rd_kafka_consumer_close() 3048 rd_kafka_cgrp_t *rkcg; in rd_kafka_committed() local 3054 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_committed() 3834 rkcg->rkcg_flags); in rd_kafka_dump0() 3901 rd_kafka_cgrp_t *rkcg; in rd_kafka_memberid() local [all …]
|
H A D | rdkafka_offset.c | 359 rd_kafka_cgrp_t *rkcg; in rd_kafka_commit0() local 362 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_commit0() 377 rd_kafka_q_enq(rkcg->rkcg_ops, rko); in rd_kafka_commit0() 391 rd_kafka_cgrp_t *rkcg; in rd_kafka_commit() local 396 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_commit()
|
H A D | rdkafka_assignor.h | 111 rd_kafka_assignor_run (struct rd_kafka_cgrp_s *rkcg,
|
H A D | rdkafka_request.c | 1036 rd_kafka_cgrp_t *rkcg, in rd_kafka_OffsetCommitRequest() argument 1063 rd_kafka_buf_write_kstr(rkbuf, rkcg->rkcg_group_id); in rd_kafka_OffsetCommitRequest() 1068 rd_kafka_buf_write_i32(rkbuf, rkcg->rkcg_generation_id); in rd_kafka_OffsetCommitRequest() 1070 rd_kafka_buf_write_kstr(rkbuf, rkcg->rkcg_member_id); in rd_kafka_OffsetCommitRequest() 1075 rd_kafka_buf_write_kstr(rkbuf, rkcg->rkcg_group_instance_id); in rd_kafka_OffsetCommitRequest() 1295 rd_kafka_cgrp_t *rkcg = opaque; in rd_kafka_handle_SyncGroup() local 1301 if (rkcg->rkcg_join_state != RD_KAFKA_CGRP_JOIN_STATE_WAIT_SYNC) { in rd_kafka_handle_SyncGroup() 1305 rd_kafka_cgrp_join_state_names[rkcg-> in rd_kafka_handle_SyncGroup() 1327 rd_kafka_cgrp_op(rkcg, NULL, RD_KAFKA_NO_REPLYQ, in rd_kafka_handle_SyncGroup() 1519 rd_kafka_cgrp_t *rkcg = opaque; in rd_kafka_handle_LeaveGroup() local [all …]
|
/dports/net/librdkafka/librdkafka-1.8.2/src/ |
H A D | rdkafka_cgrp.c | 369 rd_kafka_assert(rkcg->rkcg_rk, !rkcg->rkcg_subscription); in rd_kafka_cgrp_destroy_final() 370 rd_kafka_assert(rkcg->rkcg_rk, !rkcg->rkcg_group_leader.members); in rd_kafka_cgrp_destroy_final() 378 rd_kafka_assert(rkcg->rkcg_rk, TAILQ_EMPTY(&rkcg->rkcg_topics)); in rd_kafka_cgrp_destroy_final() 413 rkcg = rd_calloc(1, sizeof(*rkcg)); in rd_kafka_cgrp_new() 424 rkcg->rkcg_ops->rkq_opaque = rkcg; in rd_kafka_cgrp_new() 426 rkcg->rkcg_wait_coord_q->rkq_serve = rkcg->rkcg_ops->rkq_serve; in rd_kafka_cgrp_new() 463 return rkcg; in rd_kafka_cgrp_new() 2005 if (rkcg->rkcg_assignor && rkcg->rkcg_assignor != rkas) { in rd_kafka_cgrp_handle_JoinGroup() 3274 rkcg->rkcg_coord, rkcg, offsets, in rd_kafka_cgrp_offsets_commit() 5363 rkcg, in rd_kafka_cgrp_serve() [all …]
|
H A D | rdkafka_cgrp.h | 296 #define RD_KAFKA_CGRP_BROKER_IS_COORD(rkcg,rkb) \ argument 297 ((rkcg)->rkcg_coord_id != -1 && \ 298 (rkcg)->rkcg_coord_id == (rkb)->rkb_nodeid) 303 #define RD_KAFKA_CGRP_IS_STATIC_MEMBER(rkcg) \ argument 304 !RD_KAFKAP_STR_IS_NULL((rkcg)->rkcg_group_instance_id) 309 void rd_kafka_cgrp_destroy_final (rd_kafka_cgrp_t *rkcg); 313 void rd_kafka_cgrp_serve (rd_kafka_cgrp_t *rkcg); 334 void rd_kafka_cgrp_coord_query (rd_kafka_cgrp_t *rkcg, 338 void rd_kafka_cgrp_metadata_update_check (rd_kafka_cgrp_t *rkcg, 344 rd_kafka_cgrp_assigned_offsets_commit (rd_kafka_cgrp_t *rkcg, [all …]
|
H A D | rdkafka_subscription.c | 40 rd_kafka_cgrp_t *rkcg; in rd_kafka_unsubscribe() local 42 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_unsubscribe() 76 rd_kafka_cgrp_t *rkcg; in rd_kafka_subscribe() local 79 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_subscribe() 108 rd_kafka_cgrp_t *rkcg; in rd_kafka_assign0() local 110 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_assign0() 176 rd_kafka_cgrp_t *rkcg; in rd_kafka_assignment_lost() local 178 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_assignment_lost() 188 rd_kafka_cgrp_t *rkcg; in rd_kafka_rebalance_protocol() local 217 rd_kafka_cgrp_t *rkcg; in rd_kafka_assignment() local [all …]
|
H A D | rdkafka_assignor.c | 189 rd_kafka_cgrp_t *rkcg, in rd_kafka_member_subscription_match() argument 246 rd_kafka_member_subscriptions_map (rd_kafka_cgrp_t *rkcg, in rd_kafka_member_subscriptions_map() argument 305 rd_kafka_assignor_run (rd_kafka_cgrp_t *rkcg, in rd_kafka_assignor_run() argument 323 if (rkcg->rkcg_rk->rk_conf.debug & in rd_kafka_assignor_run() 330 rkcg->rkcg_group_id->str, in rd_kafka_assignor_run() 353 rd_kafka_dbg(rkcg->rkcg_rk, in rd_kafka_assignor_run() 365 err = rkas->rkas_assign_cb(rkcg->rkcg_rk, rkas, in rd_kafka_assignor_run() 380 rkcg->rkcg_group_id->str, in rd_kafka_assignor_run() 383 } else if (rkcg->rkcg_rk->rk_conf.debug & in rd_kafka_assignor_run() 389 rkcg->rkcg_group_id->str, in rd_kafka_assignor_run() [all …]
|
H A D | rdkafka.c | 3166 rd_kafka_cgrp_t *rkcg; in rd_kafka_poll_set_consumer() local 3168 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_poll_set_consumer() 3180 rd_kafka_cgrp_t *rkcg; in rd_kafka_consumer_poll() local 3193 rd_kafka_cgrp_t *rkcg; in rd_kafka_consumer_close() local 3198 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_consumer_close() 3215 rd_kafka_q_fwd_set(rkcg->rkcg_q, rkq); in rd_kafka_consumer_close() 3267 rd_kafka_cgrp_t *rkcg; in rd_kafka_committed() local 3273 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_committed() 4045 rkcg->rkcg_flags); in rd_kafka_dump0() 4103 rd_kafka_cgrp_t *rkcg; in rd_kafka_memberid() local [all …]
|
H A D | rdkafka_offset.c | 357 rd_kafka_cgrp_t *rkcg; in rd_kafka_commit0() local 360 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_commit0() 375 rd_kafka_q_enq(rkcg->rkcg_ops, rko); in rd_kafka_commit0() 389 rd_kafka_cgrp_t *rkcg; in rd_kafka_commit() local 394 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_commit()
|
H A D | rdkafka_metadata.c | 1026 rd_kafka_cgrp_t *rkcg; in rd_kafka_metadata_refresh_consumer_topics() local 1036 rkcg = rk->rk_cgrp; in rd_kafka_metadata_refresh_consumer_topics() 1037 rd_assert(rkcg != NULL); in rd_kafka_metadata_refresh_consumer_topics() 1039 if (rkcg->rkcg_flags & RD_KAFKA_CGRP_F_WILDCARD_SUBSCRIPTION) { in rd_kafka_metadata_refresh_consumer_topics() 1055 if (rkcg->rkcg_subscription) in rd_kafka_metadata_refresh_consumer_topics() 1057 rkcg->rkcg_subscription, &topics, in rd_kafka_metadata_refresh_consumer_topics()
|
H A D | rdkafka_assignor.h | 192 rd_kafka_assignor_run (struct rd_kafka_cgrp_s *rkcg,
|
/dports/sysutils/fluent-bit/fluent-bit-1.8.11/plugins/out_kafka/librdkafka-1.7.0/src/ |
H A D | rdkafka_cgrp.c | 307 rd_kafka_assert(rkcg->rkcg_rk, !rkcg->rkcg_subscription); in rd_kafka_cgrp_destroy_final() 308 rd_kafka_assert(rkcg->rkcg_rk, !rkcg->rkcg_group_leader.members); in rd_kafka_cgrp_destroy_final() 316 rd_kafka_assert(rkcg->rkcg_rk, TAILQ_EMPTY(&rkcg->rkcg_topics)); in rd_kafka_cgrp_destroy_final() 351 rkcg = rd_calloc(1, sizeof(*rkcg)); in rd_kafka_cgrp_new() 361 rkcg->rkcg_ops->rkq_opaque = rkcg; in rd_kafka_cgrp_new() 363 rkcg->rkcg_wait_coord_q->rkq_serve = rkcg->rkcg_ops->rkq_serve; in rd_kafka_cgrp_new() 400 return rkcg; in rd_kafka_cgrp_new() 1760 if (rkcg->rkcg_assignor && rkcg->rkcg_assignor != rkas) { in rd_kafka_cgrp_handle_JoinGroup() 3016 rkcg->rkcg_coord, rkcg, offsets, in rd_kafka_cgrp_offsets_commit() 5101 rkcg, in rd_kafka_cgrp_serve() [all …]
|
H A D | rdkafka_cgrp.h | 289 #define RD_KAFKA_CGRP_BROKER_IS_COORD(rkcg,rkb) \ argument 290 ((rkcg)->rkcg_coord_id != -1 && \ 291 (rkcg)->rkcg_coord_id == (rkb)->rkb_nodeid) 296 #define RD_KAFKA_CGRP_IS_STATIC_MEMBER(rkcg) \ argument 297 !RD_KAFKAP_STR_IS_NULL((rkcg)->rkcg_group_instance_id) 302 void rd_kafka_cgrp_destroy_final (rd_kafka_cgrp_t *rkcg); 306 void rd_kafka_cgrp_serve (rd_kafka_cgrp_t *rkcg); 324 void rd_kafka_cgrp_handle_SyncGroup (rd_kafka_cgrp_t *rkcg, 331 void rd_kafka_cgrp_coord_query (rd_kafka_cgrp_t *rkcg, 341 rd_kafka_cgrp_assigned_offsets_commit (rd_kafka_cgrp_t *rkcg, [all …]
|
H A D | rdkafka_subscription.c | 40 rd_kafka_cgrp_t *rkcg; in rd_kafka_unsubscribe() local 42 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_unsubscribe() 76 rd_kafka_cgrp_t *rkcg; in rd_kafka_subscribe() local 79 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_subscribe() 108 rd_kafka_cgrp_t *rkcg; in rd_kafka_assign0() local 110 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_assign0() 176 rd_kafka_cgrp_t *rkcg; in rd_kafka_assignment_lost() local 178 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_assignment_lost() 188 rd_kafka_cgrp_t *rkcg; in rd_kafka_rebalance_protocol() local 217 rd_kafka_cgrp_t *rkcg; in rd_kafka_assignment() local [all …]
|
H A D | rdkafka_assignor.c | 189 rd_kafka_cgrp_t *rkcg, in rd_kafka_member_subscription_match() argument 246 rd_kafka_member_subscriptions_map (rd_kafka_cgrp_t *rkcg, in rd_kafka_member_subscriptions_map() argument 305 rd_kafka_assignor_run (rd_kafka_cgrp_t *rkcg, in rd_kafka_assignor_run() argument 323 if (rkcg->rkcg_rk->rk_conf.debug & in rd_kafka_assignor_run() 330 rkcg->rkcg_group_id->str, in rd_kafka_assignor_run() 353 rd_kafka_dbg(rkcg->rkcg_rk, in rd_kafka_assignor_run() 365 err = rkas->rkas_assign_cb(rkcg->rkcg_rk, rkas, in rd_kafka_assignor_run() 380 rkcg->rkcg_group_id->str, in rd_kafka_assignor_run() 383 } else if (rkcg->rkcg_rk->rk_conf.debug & in rd_kafka_assignor_run() 389 rkcg->rkcg_group_id->str, in rd_kafka_assignor_run() [all …]
|
H A D | rdkafka.c | 3165 rd_kafka_cgrp_t *rkcg; 3167 if (!(rkcg = rd_kafka_cgrp_get(rk))) 3179 rd_kafka_cgrp_t *rkcg; 3192 rd_kafka_cgrp_t *rkcg; 3197 if (!(rkcg = rd_kafka_cgrp_get(rk))) 3214 rd_kafka_q_fwd_set(rkcg->rkcg_q, rkq); 3266 rd_kafka_cgrp_t *rkcg; 3272 if (!(rkcg = rd_kafka_cgrp_get(rk))) 4044 rkcg->rkcg_flags); 4102 rd_kafka_cgrp_t *rkcg; [all …]
|
H A D | rdkafka_offset.c | 357 rd_kafka_cgrp_t *rkcg; in rd_kafka_broker_get_state() 360 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_broker_get_state() 375 rd_kafka_q_enq(rkcg->rkcg_ops, rko); in rd_kafka_broker_get_state() 389 rd_kafka_cgrp_t *rkcg; in rd_kafka_broker_supports() 394 if (!(rkcg = rd_kafka_cgrp_get(rk))) in rd_kafka_broker_supports()
|
H A D | rdkafka_metadata.c | 1026 rd_kafka_cgrp_t *rkcg; 1036 rkcg = rk->rk_cgrp; 1037 rd_assert(rkcg != NULL); 1039 if (rkcg->rkcg_flags & RD_KAFKA_CGRP_F_WILDCARD_SUBSCRIPTION) { 1055 if (rkcg->rkcg_subscription) 1057 rkcg->rkcg_subscription, &topics,
|
H A D | rdkafka_assignor.h | 192 rd_kafka_assignor_run (struct rd_kafka_cgrp_s *rkcg,
|
H A D | rdkafka_request.c | 1180 rd_kafka_cgrp_t *rkcg, in rd_kafka_OffsetCommitRequest() argument 1207 rd_kafka_buf_write_kstr(rkbuf, rkcg->rkcg_group_id); in rd_kafka_OffsetCommitRequest() 1212 rd_kafka_buf_write_i32(rkbuf, rkcg->rkcg_generation_id); in rd_kafka_OffsetCommitRequest() 1214 rd_kafka_buf_write_kstr(rkbuf, rkcg->rkcg_member_id); in rd_kafka_OffsetCommitRequest() 1219 rd_kafka_buf_write_kstr(rkbuf, rkcg->rkcg_group_instance_id); in rd_kafka_OffsetCommitRequest() 1479 rd_kafka_cgrp_t *rkcg = opaque; in rd_kafka_handle_SyncGroup() local 1485 if (rkcg->rkcg_join_state != RD_KAFKA_CGRP_JOIN_STATE_WAIT_SYNC) { in rd_kafka_handle_SyncGroup() 1489 rd_kafka_cgrp_join_state_names[rkcg-> in rd_kafka_handle_SyncGroup() 1511 rd_kafka_cgrp_op(rkcg, NULL, RD_KAFKA_NO_REPLYQ, in rd_kafka_handle_SyncGroup() 1703 rd_kafka_cgrp_t *rkcg = opaque; in rd_kafka_handle_LeaveGroup() local [all …]
|