Home
last modified time | relevance | path

Searched refs:rkcg (Results 1 – 25 of 36) sorted by relevance

12

/dports/sysutils/vector/vector-0.10.0/cargo-crates/rdkafka-sys-2.0.0+1.4.2/librdkafka/src/
H A Drdkafka_cgrp.c178 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 Drdkafka_cgrp.h258 #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 Drdkafka_subscription.c40 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 Drdkafka_assignor.c164 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 Drdkafka.c2947 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 Drdkafka_offset.c359 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 Drdkafka_assignor.h111 rd_kafka_assignor_run (struct rd_kafka_cgrp_s *rkcg,
H A Drdkafka_request.c1036 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 Drdkafka_cgrp.c369 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 Drdkafka_cgrp.h296 #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 Drdkafka_subscription.c40 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 Drdkafka_assignor.c189 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 Drdkafka.c3166 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 Drdkafka_offset.c357 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 Drdkafka_metadata.c1026 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 Drdkafka_assignor.h192 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 Drdkafka_cgrp.c307 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 Drdkafka_cgrp.h289 #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 Drdkafka_subscription.c40 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 Drdkafka_assignor.c189 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 Drdkafka.c3165 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 Drdkafka_offset.c357 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 Drdkafka_metadata.c1026 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 Drdkafka_assignor.h192 rd_kafka_assignor_run (struct rd_kafka_cgrp_s *rkcg,
H A Drdkafka_request.c1180 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 …]

12