Current section

Files

Jump to
kafka_ex lib kafka_ex client client.ex
Raw

lib/kafka_ex/client/client.ex

defmodule KafkaEx.Client do
@moduledoc """
Kayrock-compatible KafkaEx.Server implementation
This implementation attempts to keep as much Kafka 'business logic' as possible
out of the server implementation, with the motivation that this should make
the client easier to maintain as the Kafka protocol evolves.
This implementation does, however, include implementations of all of the
legacy KafkaEx.Server `handle_call` clauses so that it can be compatible with
the legacy KafkaEx API.
"""
alias KafkaEx.Config
alias KafkaEx.Network.NetworkClient
alias KafkaEx.Telemetry
alias KafkaEx.Client.Error
alias KafkaEx.Client.NodeSelector
alias KafkaEx.Client.RequestBudget
alias KafkaEx.Client.RequestBuilder
alias KafkaEx.Client.RequestContext
alias KafkaEx.Client.ResponseParser
alias KafkaEx.Client.State
alias KafkaEx.Cluster.Broker
alias KafkaEx.Cluster.ClusterMetadata
alias KafkaEx.Messages.Fetch
alias KafkaEx.Messages.FindCoordinator, as: FindCoordinatorMsg
alias KafkaEx.Support.OptionalDeps
alias KafkaEx.Support.Retry
@protocol Application.compile_env(:kafka_ex, :protocol, KafkaEx.Protocol.KayrockProtocol)
use GenServer
@type args :: [KafkaEx.worker_setting() | {:allow_auto_topic_creation, boolean}]
@doc """
Start the server in a supervision tree
"""
@spec start_link(args, GenServer.name() | :no_name) :: GenServer.on_start()
def start_link(args, name \\ __MODULE__)
def start_link(args, :no_name) do
GenServer.start_link(__MODULE__, [args, nil])
end
def start_link(args, name) do
GenServer.start_link(__MODULE__, [args, name], name: name)
end
@doc """
Send a protocol request to the appropriate broker
Broker metadata will be updated if necessary
"""
@spec send_request(KafkaEx.API.client(), map, KafkaEx.Client.NodeSelector.t(), pos_integer | nil) ::
{:ok, term} | {:error, term}
def send_request(server, request, node_selector, timeout \\ nil) do
GenServer.call(server, {:network_request, request, node_selector}, timeout_val(timeout))
end
require Logger
# The request retry budget. This is the single source of truth: KafkaEx.API's
# fetch call_timeout budget (@fetch_max_retries) is derived from it via
# retry_count/0, so the two can no longer drift by hand. See #357.
@retry_count 3
# JoinGroup/SyncGroup are sent send-once at the client (one attempt, no
# transport-level socket retry): the broker legitimately holds these responses
# for the rebalance/session window, so retrying the same socket is pointless,
# and ConsumerGroup.Manager owns rejoin (@max_join_retries / @max_sync_retries).
# Matches Java/kafka-python/librdkafka/brod, which all send-once + rejoin higher up.
@coordinator_max_attempts 1
@reconnect_max_retries 3
@reconnect_delay_ms 500
# ApiVersions negotiation retry settings (issue #433)
# Handles intermittent parse errors during initial connection
@api_versions_max_retries 3
@api_versions_retry_base_delay_ms 100
@doc false
# The request retry budget (single source of truth for #357). KafkaEx.API
# derives its fetch call_timeout from this so the two cannot drift.
@spec retry_count() :: pos_integer()
def retry_count, do: @retry_count
@doc false
# Exposed so KafkaEx.API can size the JoinGroup/SyncGroup call budget from it.
@spec coordinator_max_attempts() :: pos_integer()
def coordinator_max_attempts, do: @coordinator_max_attempts
@impl true
def init([args, _name]) do
state = State.static_init(args)
unless KafkaEx.valid_consumer_group?(state.consumer_group_for_auto_commit) do
raise KafkaEx.InvalidConsumerGroupError,
state.consumer_group_for_auto_commit
end
# Crash loudly at boot if the user's config implies an optional
# dep that isn't loaded (e.g. :msk_iam SASL without :aws_signature,
# or :snappy compression without :snappyer). Much friendlier than
# UndefinedFunctionError at first produce/auth.
:ok = OptionalDeps.validate!(state.auth)
# One-time deprecation notice for the legacy :sync_timeout key (→ :request_timeout).
:ok = Config.warn_deprecated_timeout_config()
brokers =
state.bootstrap_uris
|> Enum.with_index()
|> Enum.into(%{}, fn {{host, port}, ix} ->
{ix + 1,
%Broker{
host: host,
port: port,
socket: NetworkClient.create_socket(host, port, state.ssl_options, state.use_ssl, state.auth)
}}
end)
state = %{state | cluster_metadata: %ClusterMetadata{brokers: brokers}}
# Wrap remaining init in try to ensure sockets are closed on failure
try do
check_brokers_sockets!(brokers)
# ApiVersions negotiation with retry (issue #433)
{ok_or_err, api_versions_or_error, state} = get_api_versions_with_retry(state)
case {ok_or_err, api_versions_or_error} do
{:error, :parse_error} ->
sleep_for_reconnect()
raise "ApiVersions negotiation failed: unable to parse broker response after #{@api_versions_max_retries} retries. " <>
"This may indicate network issues, broker instability, or protocol incompatibility."
{:error, :no_broker} ->
sleep_for_reconnect()
raise "ApiVersions negotiation failed: no broker available. " <>
"Check that brokers are reachable at the configured addresses."
{:error, reason} ->
sleep_for_reconnect()
raise "ApiVersions negotiation failed: #{inspect(reason)}. " <>
"Check broker connectivity and network configuration."
{:ok, api_versions} ->
api_versions
end
api_versions = api_versions_or_error
initial_topics = Keyword.get(args, :initial_topics, [])
state = State.ingest_api_versions(state, api_versions)
state =
try do
update_metadata(state, initial_topics)
rescue
e ->
sleep_for_reconnect()
Kernel.reraise(e, __STACKTRACE__)
end
{:ok, timer_ref} = :timer.send_interval(state.metadata_update_interval, :update_metadata)
state = %{state | metadata_timer_ref: timer_ref}
{:ok, state}
rescue
e ->
# Close all sockets before re-raising to prevent resource leak
close_all_sockets(state)
reraise e, __STACKTRACE__
end
end
defp close_all_sockets(state) do
Enum.each(State.brokers(state), fn broker ->
NetworkClient.close_socket(broker, broker.socket, :init_error)
end)
end
@impl true
def handle_call(:cluster_metadata, _from, state) do
{:reply, {:ok, state.cluster_metadata}, state}
end
def handle_call(:correlation_id, _from, state) do
{:reply, {:ok, state.correlation_id}, state}
end
def handle_call(:update_metadata, _from, state) do
updated_state = update_metadata(state)
{:reply, {:ok, updated_state.cluster_metadata}, updated_state}
end
def handle_call(
{:set_consumer_group_for_auto_commit, consumer_group},
_from,
state
) do
if KafkaEx.valid_consumer_group?(consumer_group) do
{:reply, :ok, %{state | consumer_group_for_auto_commit: consumer_group}}
else
{:reply, {:error, :invalid_consumer_group}, state}
end
end
def handle_call({:topic_metadata, topics, allow_topic_creation}, _from, state) do
{topic_metadata, updated_state} = fetch_topics_metadata(state, topics, allow_topic_creation)
{:reply, {:ok, topic_metadata}, updated_state}
end
def handle_call({:api_versions, opts}, _from, state) do
case api_versions_request(opts, state) do
{:error, error} -> {:reply, {:error, error}, state}
{result, updated_state} -> {:reply, result, updated_state}
end
end
def handle_call({:metadata, topics, opts, _api_version}, _from, state) do
case metadata_request(topics, opts, state) do
{:error, error} -> {:reply, {:error, error}, state}
{result, updated_state} -> {:reply, result, updated_state}
end
end
def handle_call({:list_offsets, [{topic, partitions_data}], opts}, _from, state) do
case list_offset_request({topic, partitions_data}, opts, state) do
{:error, error} -> {:reply, {:error, error}, state}
{result, updated_state} -> {:reply, result, updated_state}
end
end
# Backward compatibility, to be deleted once we delete legacy code
def handle_call({:offset, topic, partition, timestamp}, _from, state) do
partition_data = %{partition_num: partition, timestamp: timestamp}
case list_offset_request({topic, [partition_data]}, [], state) do
{:error, error} -> {:reply, {:error, error}, state}
{result, updated_state} -> {:reply, result, updated_state}
end
end
def handle_call({:describe_groups, [consumer_group_name], opts}, _from, state) do
if KafkaEx.valid_consumer_group?(consumer_group_name) do
case describe_group_request(consumer_group_name, opts, state) do
{:error, error} -> {:reply, {:error, error}, state}
{result, updated_state} -> {:reply, result, updated_state}
end
else
{:reply, {:error, :invalid_consumer_group}, state}
end
end
def handle_call({:list_groups, node_id, opts}, _from, state) do
case list_groups_request(node_id, opts, state) do
{:error, error} -> {:reply, {:error, error}, state}
{result, updated_state} -> {:reply, result, updated_state}
end
end
def handle_call({:offset_fetch, consumer_group, [{topic, partitions_data}], opts}, _from, state) do
if KafkaEx.valid_consumer_group?(consumer_group) and is_binary(consumer_group) do
case offset_fetch_request(consumer_group, {topic, partitions_data}, opts, state) do
{:error, error} -> {:reply, {:error, error}, state}
{result, updated_state} -> {:reply, result, updated_state}
end
else
{:reply, {:error, :invalid_consumer_group}, state}
end
end
def handle_call({:offset_commit, consumer_group, [{topic, partitions_data}], opts}, _from, state) do
if KafkaEx.valid_consumer_group?(consumer_group) and is_binary(consumer_group) do
case offset_commit_request(consumer_group, {topic, partitions_data}, opts, state) do
{:error, error} -> {:reply, {:error, error}, state}
{result, updated_state} -> {:reply, result, updated_state}
end
else
{:reply, {:error, :invalid_consumer_group}, state}
end
end
def handle_call({:heartbeat, consumer_group, member_id, generation_id, opts}, _from, state) do
if KafkaEx.valid_consumer_group?(consumer_group) and is_binary(consumer_group) do
case heartbeat_request(consumer_group, member_id, generation_id, opts, state) do
{:error, error} -> {:reply, {:error, error}, state}
{result, updated_state} -> {:reply, result, updated_state}
end
else
{:reply, {:error, :invalid_consumer_group}, state}
end
end
def handle_call({:join_group, consumer_group, member_id, opts}, _from, state) do
if KafkaEx.valid_consumer_group?(consumer_group) and is_binary(consumer_group) do
case join_group_request(consumer_group, member_id, opts, state) do
{:error, error} -> {:reply, {:error, error}, state}
{result, updated_state} -> {:reply, result, updated_state}
end
else
{:reply, {:error, :invalid_consumer_group}, state}
end
end
def handle_call({:leave_group, consumer_group, member_id, opts}, _from, state) do
if KafkaEx.valid_consumer_group?(consumer_group) and is_binary(consumer_group) do
case leave_group_request(consumer_group, member_id, opts, state) do
{:error, error} -> {:reply, {:error, error}, state}
{result, updated_state} -> {:reply, result, updated_state}
end
else
{:reply, {:error, :invalid_consumer_group}, state}
end
end
def handle_call({:sync_group, consumer_group, generation_id, member_id, opts}, _from, state) do
if KafkaEx.valid_consumer_group?(consumer_group) and is_binary(consumer_group) do
case sync_group_request(consumer_group, generation_id, member_id, opts, state) do
{:error, error} -> {:reply, {:error, error}, state}
{result, updated_state} -> {:reply, result, updated_state}
end
else
{:reply, {:error, :invalid_consumer_group}, state}
end
end
def handle_call({:produce, topic, partition, messages, opts}, _from, state) do
case produce_request(topic, partition, messages, opts, state) do
{:error, error} -> {:reply, {:error, error}, state}
{result, updated_state} -> {:reply, result, updated_state}
end
end
def handle_call({:fetch, topic, partition, offset, opts}, _from, state) do
case fetch_request(topic, partition, offset, opts, state) do
{:error, error} -> {:reply, {:error, error}, state}
{result, updated_state} -> {:reply, result, updated_state}
end
end
def handle_call({:find_coordinator, group_id, opts}, _from, state) do
case find_coordinator_request(group_id, opts, state) do
{:error, error} -> {:reply, {:error, error}, state}
{result, updated_state} -> {:reply, result, updated_state}
end
end
def handle_call({:create_topics, topics, timeout, opts}, _from, state) do
case create_topics_request(topics, timeout, opts, state) do
{:error, error} -> {:reply, {:error, error}, state}
{result, updated_state} -> {:reply, result, updated_state}
end
end
def handle_call({:delete_topics, topics, timeout, opts}, _from, state) do
case delete_topics_request(topics, timeout, opts, state) do
{:error, error} -> {:reply, {:error, error}, state}
{result, updated_state} -> {:reply, result, updated_state}
end
end
def handle_call({:network_request, request, node_selector}, _from, state) do
{response, updated_state} = network_request(request, node_selector, state)
{:reply, response, updated_state}
end
def handle_call(:consumer_group, _from, state) do
{:reply, state.consumer_group_for_auto_commit, state}
end
@max_metadata_update_retries 3
@impl true
def handle_info(:update_metadata, state) do
{:noreply, update_metadata_with_retry(state, @max_metadata_update_retries)}
end
def handle_info({:tcp_closed, socket}, state) do
state_out = close_broker_by_socket(state, socket)
{:noreply, state_out}
end
def handle_info({:ssl_closed, socket}, state) do
state_out = close_broker_by_socket(state, socket)
{:noreply, state_out}
end
defp update_metadata_with_retry(_state, retries_left) when retries_left <= 0 do
raise KafkaEx.MetadataUpdateError, attempts: @max_metadata_update_retries
end
defp update_metadata_with_retry(state, retries_left) do
update_metadata(state)
rescue
error ->
Logger.warning("Periodic metadata update failed (#{retries_left - 1} retries left): #{inspect(error)}")
Process.sleep(500)
update_metadata_with_retry(state, retries_left - 1)
end
@impl true
def terminate(reason, state) do
Logger.debug("Shutting down worker #{inspect(worker_identity())}, reason: #{inspect(reason)}")
if state.metadata_timer_ref, do: :timer.cancel(state.metadata_timer_ref), else: nil
Enum.each(State.brokers(state), &NetworkClient.close_socket(&1, &1.socket, :shutdown))
end
defp worker_identity do
case Process.info(self(), :registered_name) do
{:registered_name, []} -> self()
{:registered_name, name} -> name
end
end
defp update_metadata(state, topics \\ []) do
# make sure we update metadata about known topics
known_topics = ClusterMetadata.known_topics(state.cluster_metadata)
topics = Enum.uniq(known_topics ++ topics)
{updated_state, parsed_metadata} = retrieve_metadata(state, config_request_timeout(), topics)
case parsed_metadata do
nil ->
updated_state
%ClusterMetadata{} = new_cluster_metadata ->
{updated_cluster_metadata, brokers_to_close} =
ClusterMetadata.merge_brokers(updated_state.cluster_metadata, new_cluster_metadata)
:ok = Enum.each(brokers_to_close, &NetworkClient.close_socket(&1, &1.socket, :metadata_update))
state_with_meta = %{updated_state | cluster_metadata: updated_cluster_metadata}
updated_state = State.update_brokers(state_with_meta, &maybe_connect_broker(&1, state))
updated_state
end
end
# ----------------------------------------------------------------------------------------------------
defp describe_group_request(consumer_group_name, opts, state) do
node_selector = NodeSelector.consumer_group(consumer_group_name)
req_data = [{:group_names, [consumer_group_name]} | opts]
case RequestBuilder.describe_groups_request(req_data, state) do
{:ok, request} -> handle_describe_group_request(request, node_selector, state)
{:error, error} -> {:error, error}
end
end
defp list_groups_request(node_id, opts, state) do
node_selector = NodeSelector.node_id(node_id)
case RequestBuilder.list_groups_request(opts, state) do
{:ok, request} -> handle_list_groups_request(request, node_selector, state)
{:error, error} -> {:error, error}
end
end
defp list_offset_request({topic, partitions_data}, opts, state) do
[%{partition_num: partition_num}] = partitions_data
node_selector = NodeSelector.topic_partition(topic, partition_num)
req_data = [{:topics, [{topic, partitions_data}]} | opts]
case RequestBuilder.lists_offset_request(req_data, state) do
{:ok, request} -> handle_lists_offsets_request(request, node_selector, state)
{:error, error} -> {:error, error}
end
end
defp offset_fetch_request(consumer_group, {topic, partitions_data}, opts, state) do
node_selector = NodeSelector.consumer_group(consumer_group)
req_data = [{:group_id, consumer_group}, {:topics, [{topic, partitions_data}]} | opts]
case RequestBuilder.offset_fetch_request(req_data, state) do
{:ok, request} -> handle_offset_fetch_request(request, node_selector, state)
{:error, error} -> {:error, error}
end
end
defp offset_commit_request(consumer_group, {topic, partitions_data}, opts, state) do
client_id = Config.client_id()
partition_count = length(partitions_data)
metadata = Telemetry.commit_metadata(consumer_group, client_id, topic, partition_count)
Telemetry.span([:kafka_ex, :consumer, :commit], metadata, fn ->
do_offset_commit_request(consumer_group, {topic, partitions_data}, opts, state, metadata)
end)
end
defp do_offset_commit_request(consumer_group, {topic, partitions_data}, opts, state, metadata) do
node_selector = NodeSelector.consumer_group(consumer_group)
req_data = [{:group_id, consumer_group}, {:topics, [{topic, partitions_data}]} | opts]
with {:ok, request} <- RequestBuilder.offset_commit_request(req_data, state),
{result, updated_state} <- handle_offset_commit_request(request, node_selector, state) do
{{result, updated_state}, metadata}
else
{:error, error} -> {{:error, error}, metadata}
end
end
defp heartbeat_request(consumer_group, member_id, generation_id, opts, state) do
metadata = Telemetry.heartbeat_metadata(consumer_group, member_id, generation_id)
Telemetry.span([:kafka_ex, :consumer, :heartbeat], metadata, fn ->
do_heartbeat_request(consumer_group, member_id, generation_id, opts, state, metadata)
end)
end
defp do_heartbeat_request(consumer_group, member_id, generation_id, opts, state, metadata) do
node_selector = NodeSelector.consumer_group(consumer_group)
{network_timeout, req_opts} = Keyword.pop(opts, :network_timeout)
req_data = [{:group_id, consumer_group}, {:member_id, member_id}, {:generation_id, generation_id} | req_opts]
with {:ok, request} <- RequestBuilder.heartbeat_request(req_data, state),
{result, updated_state} <- handle_heartbeat_request(request, node_selector, state, network_timeout) do
{{result, updated_state}, metadata}
else
{:error, error} -> {{:error, error}, metadata}
end
end
defp join_group_request(consumer_group, member_id, opts, state) do
topics = Keyword.get(opts, :topics, [])
metadata = Telemetry.join_group_metadata(consumer_group, member_id || "", topics)
Telemetry.span([:kafka_ex, :consumer, :join], metadata, fn ->
do_join_group_request(consumer_group, member_id, opts, state, metadata)
end)
end
defp do_join_group_request(consumer_group, member_id, opts, state, metadata) do
node_selector = NodeSelector.consumer_group(consumer_group)
# :network_timeout is a control opt (derived in KafkaEx.API.join_group from
# rebalance_timeout), not a protocol field — pop it before building req_data.
{network_timeout, req_opts} = Keyword.pop(opts, :network_timeout)
req_data = [{:group_id, consumer_group}, {:member_id, member_id} | req_opts]
with {:ok, request} <- RequestBuilder.join_group_request(req_data, state),
{{:ok, result}, updated_state} <- handle_join_group_request(request, node_selector, state, network_timeout) do
stop_metadata = add_join_group_result_metadata(metadata, result)
{{{:ok, result}, updated_state}, stop_metadata}
else
# KIP-394: JoinGroup V4+ requires two-step join.
# First attempt with empty member_id returns :member_id_required
# with an assigned member_id. Retry with that member_id.
{{:error, %Error{error: :member_id_required, metadata: %{member_id: assigned_id}}}, updated_state} ->
Logger.info("JoinGroup returned member_id_required (KIP-394), retrying with assigned member_id")
do_join_group_request(consumer_group, assigned_id, opts, updated_state, metadata)
{:error, error} ->
{{:error, error}, metadata}
error_result ->
{error_result, metadata}
end
end
defp add_join_group_result_metadata(metadata, result) do
is_leader = result.leader_id == result.member_id
metadata
|> Map.put(:generation_id, result.generation_id)
|> Map.put(:is_leader, is_leader)
end
defp leave_group_request(consumer_group, member_id, opts, state) do
metadata = Telemetry.leave_group_metadata(consumer_group, member_id)
Telemetry.span([:kafka_ex, :consumer, :leave], metadata, fn ->
do_leave_group_request(consumer_group, member_id, opts, state, metadata)
end)
end
defp do_leave_group_request(consumer_group, member_id, opts, state, metadata) do
node_selector = NodeSelector.consumer_group(consumer_group)
req_data = [{:group_id, consumer_group}, {:member_id, member_id} | opts]
with {:ok, request} <- RequestBuilder.leave_group_request(req_data, state),
{result, updated_state} <- handle_leave_group_request(request, node_selector, state) do
{{result, updated_state}, metadata}
else
{:error, error} -> {{:error, error}, metadata}
end
end
defp sync_group_request(consumer_group, generation_id, member_id, opts, state) do
# Leader sends non-empty group_assignment, followers send empty
group_assignment = Keyword.get(opts, :group_assignment, [])
is_leader = group_assignment != []
metadata = Telemetry.sync_group_metadata(consumer_group, member_id, generation_id, is_leader)
Telemetry.span([:kafka_ex, :consumer, :sync], metadata, fn ->
do_sync_group_request(consumer_group, generation_id, member_id, opts, state, metadata)
end)
end
defp do_sync_group_request(consumer_group, generation_id, member_id, opts, state, metadata) do
node_selector = NodeSelector.consumer_group(consumer_group)
# :network_timeout is a control opt (derived in KafkaEx.API.sync_group), not a
# protocol field — pop it before building req_data.
{network_timeout, req_opts} = Keyword.pop(opts, :network_timeout)
req_data = [{:group_id, consumer_group}, {:generation_id, generation_id}, {:member_id, member_id} | req_opts]
with {:ok, request} <- RequestBuilder.sync_group_request(req_data, state),
{{:ok, result}, updated_state} <- handle_sync_group_request(request, node_selector, state, network_timeout) do
assigned_partitions = count_assigned_partitions(result.partition_assignments)
stop_metadata = Map.put(metadata, :assigned_partitions, assigned_partitions)
{{{:ok, result}, updated_state}, stop_metadata}
else
{:error, error} -> {{:error, error}, metadata}
error_result -> {error_result, metadata}
end
end
defp count_assigned_partitions(partition_assignments) do
Enum.reduce(partition_assignments, 0, fn assignment, acc ->
acc + length(assignment.partitions)
end)
end
defp produce_request(topic, partition, messages, opts, state) do
message_count = length(messages)
required_acks = Keyword.get(opts, :required_acks, 1)
client_id = Config.client_id()
metadata = Telemetry.produce_metadata(topic, partition, client_id, required_acks)
start_measurements = %{message_count: message_count}
Telemetry.span([:kafka_ex, :produce], Map.merge(metadata, start_measurements), fn ->
do_produce_request(topic, partition, messages, opts, state, metadata)
end)
end
defp do_produce_request(topic, partition, messages, opts, state, metadata) do
node_selector = NodeSelector.topic_partition(topic, partition)
req_data = [{:topic, topic}, {:partition, partition}, {:messages, messages} | opts]
with {:ok, request} <- RequestBuilder.produce_request(req_data, state),
{{:ok, result}, updated_state} <- handle_produce_request(request, node_selector, state) do
stop_metadata = add_offset_to_metadata(metadata, result)
{{{:ok, result}, updated_state}, stop_metadata}
else
{:error, error} -> {{:error, error}, metadata}
error_result -> {error_result, metadata}
end
end
defp add_offset_to_metadata(metadata, result) do
case Map.get(result, :base_offset) do
nil -> metadata
offset -> Map.put(metadata, :offset, offset)
end
end
defp fetch_request(topic, partition, offset, opts, state) do
client_id = Config.client_id()
metadata = Telemetry.fetch_metadata(topic, partition, offset, client_id)
Telemetry.span([:kafka_ex, :fetch], metadata, fn ->
do_fetch_request(topic, partition, offset, opts, state, metadata)
end)
end
defp do_fetch_request(topic, partition, offset, opts, state, metadata) do
node_selector = NodeSelector.topic_partition(topic, partition)
{network_timeout, req_opts} = Keyword.pop(opts, :network_timeout)
req_data = [{:topic, topic}, {:partition, partition}, {:offset, offset} | req_opts]
with {:ok, request} <- RequestBuilder.fetch_request(req_data, state),
{{:ok, result}, updated_state} <-
handle_fetch_request(request, node_selector, state, network_timeout) do
filtered_result = Fetch.filter_from_offset(result, offset)
message_count = length(Map.get(filtered_result, :records, []))
stop_metadata = Map.put(metadata, :message_count, message_count)
{{{:ok, filtered_result}, updated_state}, stop_metadata}
else
{:error, error} -> {{:error, error}, metadata}
error_result -> {error_result, metadata}
end
end
defp find_coordinator_request(group_id, opts, state) do
# FindCoordinator can be sent to any broker
node_selector = NodeSelector.first_available()
req_data = [{:group_id, group_id} | opts]
case RequestBuilder.find_coordinator_request(req_data, state) do
{:ok, request} -> handle_find_coordinator_request(request, node_selector, state)
{:error, error} -> {:error, error}
end
end
defp create_topics_request(topics, timeout, opts, state) do
# CreateTopics must be sent to the controller broker
node_selector = NodeSelector.controller()
{network_timeout, req_opts} = Keyword.pop(opts, :network_timeout)
req_data = [{:topics, topics}, {:timeout, timeout} | req_opts]
case RequestBuilder.create_topics_request(req_data, state) do
{:ok, request} -> handle_create_topics_request(request, node_selector, state, network_timeout)
{:error, error} -> {:error, error}
end
end
defp delete_topics_request(topics, timeout, opts, state) do
# DeleteTopics must be sent to the controller broker
node_selector = NodeSelector.controller()
{network_timeout, req_opts} = Keyword.pop(opts, :network_timeout)
req_data = [{:topics, topics}, {:timeout, timeout} | req_opts]
case RequestBuilder.delete_topics_request(req_data, state) do
{:ok, request} -> handle_delete_topics_request(request, node_selector, state, network_timeout)
{:error, error} -> {:error, error}
end
end
defp api_versions_request(opts, state) do
node_selector = NodeSelector.random()
case RequestBuilder.api_versions_request(opts, state) do
{:ok, request} -> handle_api_versions_request(request, node_selector, state)
{:error, error} -> {:error, error}
end
end
defp metadata_request(topics, opts, state) do
client_id = Config.client_id()
topic_list = if is_list(topics), do: topics, else: [topics]
metadata = Telemetry.metadata_update_metadata(client_id, topic_list)
Telemetry.span([:kafka_ex, :metadata, :update], metadata, fn ->
do_metadata_request(topics, opts, state)
end)
end
defp do_metadata_request(topics, opts, state) do
# Metadata can be fetched from any broker, use random selection
node_selector = NodeSelector.random()
req_data = [{:topics, topics} | opts]
case RequestBuilder.metadata_request(req_data, state) do
{:ok, request} ->
case handle_metadata_request(request, node_selector, state) do
{{:ok, cluster_metadata}, _updated_state} = result ->
broker_count = map_size(cluster_metadata.brokers)
topic_count = map_size(cluster_metadata.topics)
{result, %{broker_count: broker_count, topic_count: topic_count}}
error_result ->
{error_result, %{}}
end
{:error, error} ->
{{:error, error}, %{}}
end
end
# ----------------------------------------------------------------------------------------------------
defp handle_api_versions_request(request, node_selector, state) do
%RequestContext{request: request, parser_fn: &ResponseParser.api_versions_response/1, node_selector: node_selector}
|> handle_request_with_retry(state)
end
defp handle_describe_group_request(request, node_selector, state) do
%RequestContext{
request: request,
parser_fn: &ResponseParser.describe_groups_response/1,
node_selector: node_selector
}
|> handle_request_with_retry(state)
end
defp handle_list_groups_request(request, node_selector, state) do
%RequestContext{
request: request,
parser_fn: &ResponseParser.list_groups_response/1,
node_selector: node_selector
}
|> handle_request_with_retry(state)
end
defp handle_lists_offsets_request(request, node_selector, state) do
%RequestContext{request: request, parser_fn: &ResponseParser.list_offsets_response/1, node_selector: node_selector}
|> handle_request_with_retry(state)
end
defp handle_offset_fetch_request(request, node_selector, state) do
%RequestContext{request: request, parser_fn: &ResponseParser.offset_fetch_response/1, node_selector: node_selector}
|> handle_request_with_retry(state)
end
defp handle_offset_commit_request(request, node_selector, state) do
%RequestContext{
request: request,
parser_fn: &ResponseParser.offset_commit_response/1,
node_selector: node_selector,
retryable?: &Retry.commit_retryable?/1
}
|> handle_request_with_retry(state)
end
defp handle_heartbeat_request(request, node_selector, state, network_timeout) do
%RequestContext{
request: request,
parser_fn: &ResponseParser.heartbeat_response/1,
node_selector: node_selector,
retryable?: &Retry.heartbeat_retryable?/1,
network_timeout: network_timeout
}
|> handle_request_with_retry(state)
end
defp handle_join_group_request(request, node_selector, state, network_timeout) do
%RequestContext{
request: request,
parser_fn: &ResponseParser.join_group_response/1,
node_selector: node_selector,
retryable?: &Retry.join_group_retryable?/1,
network_timeout: network_timeout
}
|> handle_request_with_retry(state, @coordinator_max_attempts)
end
defp handle_leave_group_request(request, node_selector, state) do
%RequestContext{request: request, parser_fn: &ResponseParser.leave_group_response/1, node_selector: node_selector}
|> handle_request_with_retry(state)
end
defp handle_sync_group_request(request, node_selector, state, network_timeout) do
%RequestContext{
request: request,
parser_fn: &ResponseParser.sync_group_response/1,
node_selector: node_selector,
retryable?: &Retry.sync_group_retryable?/1,
network_timeout: network_timeout
}
|> handle_request_with_retry(state, @coordinator_max_attempts)
end
defp handle_produce_request(request, node_selector, state) do
# Produce requests should only retry on leadership errors to avoid duplicates.
# Timeout errors are NOT safe to retry because the message may have been written
# but the response was lost. Use Kafka's idempotent producer for true exactly-once.
%RequestContext{
request: request,
parser_fn: &ResponseParser.produce_response/1,
node_selector: node_selector,
retryable?: &Retry.produce_retryable?/1
}
|> handle_request_with_retry(state)
end
defp handle_fetch_request(request, node_selector, state, network_timeout) do
%RequestContext{
request: request,
parser_fn: &ResponseParser.fetch_response/1,
node_selector: node_selector,
network_timeout: network_timeout
}
|> handle_request_with_retry(state)
end
defp handle_find_coordinator_request(request, node_selector, state) do
%RequestContext{
request: request,
parser_fn: &ResponseParser.find_coordinator_response/1,
node_selector: node_selector
}
|> handle_request_with_retry(state)
end
defp handle_create_topics_request(request, node_selector, state, network_timeout) do
%RequestContext{
request: request,
parser_fn: &ResponseParser.create_topics_response/1,
node_selector: node_selector,
network_timeout: network_timeout
}
|> handle_request_with_retry(state)
end
defp handle_delete_topics_request(request, node_selector, state, network_timeout) do
%RequestContext{
request: request,
parser_fn: &ResponseParser.delete_topics_response/1,
node_selector: node_selector,
network_timeout: network_timeout
}
|> handle_request_with_retry(state)
end
defp handle_metadata_request(request, node_selector, state) do
ctx = %RequestContext{
request: request,
parser_fn: &ResponseParser.metadata_response/1,
node_selector: node_selector
}
case handle_request_with_retry(ctx, state) do
{{:ok, cluster_metadata}, updated_state} ->
merged_state = %{updated_state | cluster_metadata: cluster_metadata}
{{:ok, cluster_metadata}, merged_state}
error_result ->
error_result
end
end
# ----------------------------------------------------------------------------------------------------
# Request retry logic with optional retryable filter
#
# The `retryable?` field on the `%RequestContext{}` lets callers specify which
# errors should trigger a retry. This is important for produce requests where
# blind retries can cause duplicate messages. By default
# (`Retry.data_plane_retryable?/1`), all errors are retried.
#
# For produce requests, the context sets `&Retry.produce_retryable?/1` to only retry leadership errors.
# ----------------------------------------------------------------------------------------------------
defp handle_request_with_retry(ctx, state, retry_count \\ @retry_count, last_error \\ nil)
defp handle_request_with_retry(%RequestContext{}, state, 0, last_error) do
{{:error, last_error}, state}
end
defp handle_request_with_retry(%RequestContext{} = ctx, state, retry_count, _last_error) do
case network_request(ctx.request, ctx.node_selector, state, ctx.network_timeout) do
{{:ok, response}, state_out} ->
case ctx.parser_fn.(response) do
{:ok, result} ->
{{:ok, result}, state_out}
{:error, [error | _]} ->
handle_request_error(ctx, state_out, retry_count, error)
{:error, %Error{} = error} ->
handle_request_error(ctx, state_out, retry_count, error)
end
{{:error, reason}, state_out} ->
handle_request_error(ctx, state_out, retry_count, build_transport_error(reason))
{other, state_out} ->
Logger.warning("Unexpected network_request result for #{inspect(ctx.request.__struct__)}: #{inspect(other)}")
handle_request_error(ctx, state_out, retry_count, Error.build(:unknown, %{}))
end
end
defp build_transport_error(reason) when is_atom(reason), do: Error.build(reason, %{})
defp build_transport_error(reason), do: Error.build(:unknown, %{transport_reason: reason})
defp handle_request_error(%RequestContext{} = ctx, state, retry_count, error) do
request_name = ctx.request.__struct__
error_atom = extract_error_atom(error)
Logger.warning("Request #{inspect(request_name)} failed with error #{inspect(error_atom)}")
cond do
# Send-once only for coordinator requests (broker holds them for the rebalance/
# session window; rejoin/commit re-issue higher up). Data-plane retries via ctx.retryable?.
Retry.transport_timeout?(error_atom) and coordinator_request?(ctx.node_selector) ->
Logger.info("Coordinator request #{inspect(request_name)} timed out; not retrying at the client (send-once)")
{{:error, error}, state}
ctx.retryable?.(error_atom) ->
updated_state = refresh_for_error(error_atom, ctx, state)
handle_request_with_retry(ctx, updated_state, retry_count - 1, error)
true ->
Logger.info("Error #{inspect(error_atom)} is not retryable for #{inspect(request_name)}, failing immediately")
{{:error, error}, state}
end
end
defp extract_error_atom(%Error{error: error}), do: error
defp extract_error_atom(error) when is_atom(error), do: error
defp extract_error_atom(_), do: :unknown
defp refresh_for_error(error_atom, ctx, state) do
cond do
Retry.leadership_error?(error_atom) ->
Logger.info("Refreshing metadata due to #{inspect(error_atom)} error")
update_metadata(state)
coordinator_request?(ctx.node_selector) and Retry.coordinator_refresh_error?(error_atom) ->
Logger.info("Re-discovering group coordinator due to #{inspect(error_atom)} error")
refresh_coordinator(state, ctx.node_selector.consumer_group_name)
true ->
state
end
end
defp coordinator_request?(%NodeSelector{strategy: :consumer_group}), do: true
defp coordinator_request?(_), do: false
defp refresh_coordinator(state, consumer_group) do
state
|> State.drop_consumer_group_coordinator(consumer_group)
|> update_consumer_group_coordinator(consumer_group)
end
# ----------------------------------------------------------------------------------------------------
defp maybe_connect_broker(broker, state) do
case Broker.connected?(broker) do
true ->
broker
false ->
NetworkClient.close_socket(broker, broker.socket, :reconnecting)
socket = NetworkClient.create_socket(broker.host, broker.port, state.ssl_options, state.use_ssl, state.auth)
%{broker | socket: socket}
end
end
# Attempts to reconnect a broker with retries.
# Returns updated broker struct (with socket or nil if all retries failed).
defp reconnect_broker(broker, state, retry_count \\ 0)
defp reconnect_broker(broker, state, retry_count) when retry_count < @reconnect_max_retries do
Logger.info("Reconnecting to #{Broker.to_string(broker)}, attempt #{retry_count + 1}/#{@reconnect_max_retries}")
case NetworkClient.create_socket(broker.host, broker.port, state.ssl_options, state.use_ssl, state.auth) do
nil ->
Process.sleep(@reconnect_delay_ms)
reconnect_broker(broker, state, retry_count + 1)
socket ->
Logger.info("Successfully reconnected to #{Broker.to_string(broker)}")
%{broker | socket: socket}
end
end
defp reconnect_broker(broker, _state, _retry_count) do
Logger.warning("Failed to reconnect to #{Broker.to_string(broker)} after #{@reconnect_max_retries} attempts")
broker
end
defp retrieve_metadata(state, request_timeout, topics) do
retrieve_metadata(state, request_timeout, topics, @retry_count)
end
defp retrieve_metadata(state, _request_timeout, topics, 0) do
Logger.log(:error, "Metadata request for topics #{inspect(topics)} failed: leader_not_available after retries")
{state, nil}
end
defp retrieve_metadata(state, request_timeout, topics, retry) do
req_opts = [
topics: topics,
allow_auto_topic_creation: state.allow_auto_topic_creation
]
with {:ok, request} <- RequestBuilder.metadata_request(req_opts, state),
{{:ok, response}, state_out} <-
network_request(request, NodeSelector.first_available(), state) do
parse_metadata_response(response, state, state_out, request_timeout, topics, retry)
else
{:error, _} ->
Logger.error("Unable to build metadata request.")
{state, nil}
{{:error, _}, state_out} ->
Logger.error("Unable to fetch metadata from any brokers. Timeout is #{request_timeout}.")
{state_out, nil}
end
end
defp parse_metadata_response(response, state, state_out, request_timeout, topics, retry) do
case ResponseParser.metadata_response(response) do
{:ok, cluster_metadata} ->
# Check if any requested topics are missing (leader_not_available).
# The parser filters out topics with errors, so missing = still propagating.
if topics != [] and has_missing_topics?(topics, cluster_metadata) do
:timer.sleep(300)
retrieve_metadata(state, request_timeout, topics, retry - 1)
else
{state_out, cluster_metadata}
end
{:error, _reason} ->
Logger.error("Unable to parse metadata response from brokers.")
{state_out, nil}
end
end
defp has_missing_topics?(requested_topics, %ClusterMetadata{topics: known_topics}) do
Enum.any?(requested_topics, fn topic -> not Map.has_key?(known_topics, topic) end)
end
defp sleep_for_reconnect do
Process.sleep(Application.get_env(:kafka_ex, :sleep_for_reconnect, 400))
end
defp check_brokers_sockets!(brokers) do
any_socket_opened =
brokers
|> Enum.map(fn {_, %Broker{socket: socket}} -> !is_nil(socket) end)
|> Enum.reduce(&(&1 || &2))
if !any_socket_opened do
sleep_for_reconnect()
raise "Brokers sockets are not opened"
end
end
defp client_request(request, state) do
%{request | client_id: Config.client_id(), correlation_id: state.correlation_id}
end
# select a broker, updating state if necessary (e.g., metadata or consumer group)
# returns {broker, maybe_updated_state} - broker will be nil in case of
# failure. Ensures the returned broker is connected, attempting reconnection if needed.
defp select_broker_with_update(state, selector, state_updater) do
case State.select_broker(state, selector) do
{:error, _} ->
updated_state = state_updater.(state)
case State.select_broker(updated_state, selector) do
{:error, _} -> {nil, updated_state}
{:ok, broker} -> ensure_broker_connected(broker, updated_state)
end
{:ok, broker} ->
ensure_broker_connected(broker, state)
end
end
# Ensures broker is connected, reconnecting if necessary.
# Returns {broker, updated_state} where broker may have a new socket,
# or nil if reconnection failed.
defp ensure_broker_connected(broker, state) do
if Broker.connected?(broker) do
{broker, state}
else
reconnected_broker = reconnect_broker(broker, state)
if Broker.connected?(reconnected_broker) do
updated_state = update_broker_in_state(state, reconnected_broker)
{reconnected_broker, updated_state}
else
{nil, state}
end
end
end
defp update_broker_in_state(state, broker) do
State.update_brokers(state, fn b ->
if b.node_id == broker.node_id, do: broker, else: b
end)
end
@max_reconnect_attempts 3
defp ensure_any_broker_connected(state) do
brokers = State.brokers(state)
connected = Enum.filter(brokers, &Broker.connected?/1)
if Enum.empty?(connected) do
brokers_to_try = Enum.take(brokers, @max_reconnect_attempts)
reconnect_any_broker(brokers_to_try, state)
else
{state, connected}
end
end
defp reconnect_any_broker([], state), do: {state, []}
defp reconnect_any_broker([broker | rest], state) do
reconnected = reconnect_broker(broker, state)
if Broker.connected?(reconnected) do
updated_state = update_broker_in_state(state, reconnected)
{updated_state, [reconnected]}
else
reconnect_any_broker(rest, state)
end
end
defp broker_for_partition_with_update(state, topic, partition) do
node = NodeSelector.topic_partition(topic, partition)
select_broker_with_update(state, node, &update_metadata(&1, [topic]))
end
defp broker_for_consumer_group_with_update(state, consumer_group) do
node = NodeSelector.consumer_group(consumer_group)
select_broker_with_update(state, node, &update_consumer_group_coordinator(&1, consumer_group))
end
defp update_consumer_group_coordinator(state, consumer_group) do
req_opts = [group_id: consumer_group]
with {:ok, request} <- RequestBuilder.find_coordinator_request(req_opts, state),
{{:ok, raw_response}, updated_state} <-
network_request(request, NodeSelector.first_available(), state) do
parse_coordinator_response(raw_response, updated_state, consumer_group)
else
{:error, error} ->
Logger.warning(
"Unable to build find_coordinator request for #{inspect(consumer_group)}: Error #{inspect(error)}"
)
state
{{:error, error}, updated_state} ->
Logger.warning(
"Unable to find consumer group coordinator for #{inspect(consumer_group)}: Error #{inspect(error)}"
)
updated_state
end
end
defp parse_coordinator_response(raw_response, updated_state, consumer_group) do
case ResponseParser.find_coordinator_response(raw_response) do
{:ok, %FindCoordinatorMsg{} = result} ->
node_id = FindCoordinatorMsg.coordinator_node_id(result)
State.put_consumer_group_coordinator(updated_state, consumer_group, node_id)
{:error, error} ->
Logger.warning(
"Unable to find consumer group coordinator for #{inspect(consumer_group)}: Error #{inspect(error)}"
)
updated_state
end
end
defp first_broker_response(request, brokers, timeout) do
first_broker_response(request, Enum.shuffle(brokers), timeout, nil)
end
defp first_broker_response(_request, [], _timeout, last_error) do
{last_error || {:error, :no_connected_broker}, nil}
end
defp first_broker_response(request, [broker | rest], timeout, _last_error) do
case try_broker(broker, request, timeout) do
nil -> first_broker_response(request, rest, timeout, {:error, :broker_failed})
{:error, _} = error -> first_broker_response(request, rest, timeout, error)
response -> {response, broker}
end
end
defp try_broker(broker, request, timeout) do
case NetworkClient.send_sync_request(broker, request, timeout) do
{:error, :not_connected} ->
Logger.debug("#{Broker.to_string(broker)} not connected, skipping")
nil
{:error, error} ->
Logger.warning("Network call to #{Broker.to_string(broker)} failed: #{inspect(error)}")
nil
response ->
response
end
end
# send_request/4 is send-once; budget from :request_timeout, never the implicit
# 5000 default (a larger :request_timeout would exit the caller; #357).
defp timeout_val(nil), do: RequestBudget.call_budget(Config.request_timeout())
defp timeout_val(timeout) when is_integer(timeout), do: timeout
# The per-attempt socket-recv deadline for synchronous requests that do not
# set their own `network_timeout` (metadata, offset, heartbeat, produce ack,
# …). Resolves via `Config.request_timeout/0` (`:request_timeout` > deprecated
# `:sync_timeout` > default). JoinGroup/SyncGroup pass an explicit timeout.
defp config_request_timeout(timeout \\ nil) do
timeout || Config.request_timeout()
end
defp get_api_versions(state) do
{:ok, request} = RequestBuilder.api_versions_request([api_version: 0], state)
{{ok_or_error, response}, state_out} = network_request(request, NodeSelector.first_available(), state)
case {ok_or_error, response} do
{:ok, raw_response} ->
case ResponseParser.api_versions_response(raw_response) do
{:ok, api_versions} -> {:ok, api_versions, state_out}
{:error, error} -> {:error, error, state_out}
end
{:error, error} ->
{:error, error, state_out}
end
end
# ApiVersions negotiation with retry and exponential backoff (issue #433)
# Handles intermittent parse errors that can occur during initial connection,
# especially in high-latency environments or during broker restarts.
# Note: Uses manual retry loop instead of Retry.with_retry because we need to
# thread state changes through retries (the state is updated after each request).
defp get_api_versions_with_retry(state, retries_left \\ @api_versions_max_retries)
defp get_api_versions_with_retry(state, retries_left) do
case get_api_versions(state) do
{:ok, api_versions, updated_state} ->
{:ok, api_versions, updated_state}
{:error, error, updated_state} when retries_left > 0 ->
if Retry.transient_error?(error) do
attempt = @api_versions_max_retries - retries_left
delay = Retry.backoff_delay(attempt, @api_versions_retry_base_delay_ms)
Logger.warning(
"ApiVersions negotiation failed with #{inspect(error)}, " <>
"retrying in #{delay}ms (#{retries_left - 1} retries left)"
)
Process.sleep(delay)
get_api_versions_with_retry(updated_state, retries_left - 1)
else
Logger.error("ApiVersions negotiation failed with non-retryable error: #{inspect(error)}")
{:error, error, updated_state}
end
{:error, error, updated_state} ->
Logger.error("ApiVersions negotiation failed after #{@api_versions_max_retries} attempts: #{inspect(error)}")
{:error, error, updated_state}
end
end
defp network_request(request, node_selector, state, network_timeout \\ nil) do
synchronous = if Map.get(request, :acks) == 0, do: false, else: true
network_timeout = config_request_timeout(network_timeout)
{send_request, updated_state} = get_send_request_function(node_selector, state, network_timeout, synchronous)
case send_request do
:no_broker ->
{{:error, :no_broker}, updated_state}
{:error, _} = error ->
{error, updated_state}
_ ->
request = client_request(request, updated_state)
response = run_client_request(request, send_request, synchronous)
{response, State.increment_correlation_id(updated_state)}
end
end
defp run_client_request(
%{client_id: client_id, correlation_id: correlation_id} = client_request,
send_request,
synchronous
)
when not is_nil(client_id) and not is_nil(correlation_id) do
# Start with empty broker info - will be populated after send
start_metadata = Telemetry.request_metadata(client_request, %{})
Telemetry.span([:kafka_ex, :request], start_metadata, fn ->
do_run_client_request(client_request, send_request, synchronous, start_metadata)
end)
end
defp do_run_client_request(client_request, send_request, synchronous, start_metadata) do
wire_request = @protocol.serialize_request(client_request)
bytes_sent = IO.iodata_length(wire_request)
{result, bytes_received, broker_info} =
case send_request.(wire_request) do
{{:error, reason}, broker} ->
{{:error, reason}, 0, broker_to_telemetry_info(broker)}
{data, broker} when synchronous ->
{deserialize(data, client_request), byte_size(data), broker_to_telemetry_info(broker)}
{data, broker} ->
{data, byte_size(data), broker_to_telemetry_info(broker)}
end
additional_stop_metadata = %{bytes_sent: bytes_sent, bytes_received: bytes_received, broker: broker_info}
stop_metadata = Map.merge(start_metadata, additional_stop_metadata)
{result, stop_metadata}
end
defp get_send_request_function(%NodeSelector{strategy: :first_available}, state, network_timeout, _synchronous) do
{updated_state, connected_brokers} = ensure_any_broker_connected(state)
if Enum.empty?(connected_brokers) do
{:no_broker, updated_state}
else
{fn wire_request -> first_broker_response(wire_request, connected_brokers, network_timeout) end, updated_state}
end
end
defp get_send_request_function(
%NodeSelector{strategy: :topic_partition, topic: topic, partition: partition},
state,
network_timeout,
synchronous
) do
{broker, updated_state} = broker_for_partition_with_update(state, topic, partition)
if broker do
if synchronous do
{send_sync_request_fn(broker, network_timeout), updated_state}
else
{send_async_request_fn(broker), updated_state}
end
else
{:no_broker, updated_state}
end
end
defp get_send_request_function(
%NodeSelector{
strategy: :consumer_group,
consumer_group_name: consumer_group
},
state,
network_timeout,
_synchronous
) do
{broker, updated_state} = broker_for_consumer_group_with_update(state, consumer_group)
if broker do
{send_sync_request_fn(broker, network_timeout), updated_state}
else
{:no_broker, updated_state}
end
end
defp get_send_request_function(%NodeSelector{} = node_selector, state, network_timeout, _synchronous) do
case State.select_broker(state, node_selector) do
{:ok, broker} ->
{connected_broker, updated_state} = ensure_broker_connected(broker, state)
if connected_broker do
{send_sync_request_fn(connected_broker, network_timeout), updated_state}
else
{:no_broker, updated_state}
end
{:error, _} = error ->
{error, state}
end
end
defp send_sync_request_fn(broker, network_timeout) do
fn wire_request ->
{NetworkClient.send_sync_request(broker, wire_request, network_timeout), broker}
end
end
defp send_async_request_fn(broker) do
fn wire_request ->
{NetworkClient.send_async_request(broker, wire_request), broker}
end
end
defp broker_to_telemetry_info(nil), do: %{}
defp broker_to_telemetry_info(broker), do: %{node_id: broker.node_id, host: broker.host, port: broker.port}
defp deserialize(data, request) do
{resp, _} = @protocol.response_deserializer(request).(data)
{:ok, resp}
rescue
error ->
Logger.error(
"Failed to parse a response from the server: " <>
inspect(data, limit: :infinity) <>
" for request #{inspect(request, limit: :infinity)} " <>
"error: #{inspect(error)}"
)
{:error, :parse_error}
end
defp fetch_topics_metadata(state, topics, allow_topic_creation) do
allow_auto_topic_creation = state.allow_auto_topic_creation
updated_state = update_metadata(%{state | allow_auto_topic_creation: allow_topic_creation}, topics)
topic_metadata = State.topics_metadata(updated_state, topics)
{topic_metadata, %{updated_state | allow_auto_topic_creation: allow_auto_topic_creation}}
end
defp close_broker_by_socket(state, socket, reason \\ :remote_closed) do
State.update_brokers(state, fn broker ->
if Broker.has_socket?(broker, socket) do
Logger.debug("#{Broker.to_string(broker)} closed connection")
# Socket is already closed (received :tcp_closed/:ssl_closed), just emit telemetry
NetworkClient.close_socket(broker, socket, reason)
Broker.put_socket(broker, nil)
else
broker
end
end)
end
end