Packages

kafka_ex_tc

0.12.1-21
0.13.0 0.12.1 0.12.1-51-1 0.12.1-50-3 0.12.1-50-2 0.12.1-50-1 0.12.1-49-2 0.12.1-49-1 0.12.1-49-0 0.12.1-48-2 0.12.1-48-1 0.12.1-47-1 0.12.1-46-2 0.12.1-45-2 0.12.1-45-1 0.12.1-44-1 0.12.1-43-1 0.12.1-42-8 0.12.1-42-6 0.12.1-42-5 0.12.1-42-4 0.12.1-42-1 0.12.1-41-2 0.12.1-41-1 0.12.1-40-9 0.12.1-40-8 0.12.1-40-7 0.12.1-40-6 0.12.1-40-5 0.12.1-40-4 0.12.1-40-3 0.12.1-40-2 0.12.1-40-10 0.12.1-40-1 0.12.1-39-2 0.12.1-39-1 0.12.1-36-3 0.12.1-36-2 0.12.1-36-1 0.12.1-34-b 0.12.1-34-a 0.12.1-33-1 0.12.1-32-a1 0.12.1-29-s-3 0.12.1-29-s-2 0.12.1-29-s-1 0.12.1-28-rumble-1 0.12.1-27-2 0.12.1-27-1 0.12.1-26-rumble-2 0.12.1-26-rumble-1 0.12.1-26-dev-1 0.12.1-26-dev.1 0.12.1-26-3 0.12.1-26-2 0.12.1-26-1 0.12.1-25-s-8 0.12.1-25-s-7 0.12.1-25-s-6 0.12.1-25-s-5 0.12.1-25-s-4 0.12.1-25-s-3 0.12.1-25-s-2 0.12.1-25-s-1 0.12.1-25-dev-22 0.12.1-25-dev 0.12.1-24-2 0.12.1-24-1 0.12.1-22-poc-b-7 0.12.1-20-poc-b-6 0.12.1-20-poc-b-5 0.12.1-20-poc-b-4 0.12.1-20-poc-b-3 0.12.1-19-poc-b-3 0.12.1-19-poc-b-2 0.12.1-19-poc-b-1 0.12.1-19-poc-7 0.12.1-19-poc-6 0.12.1-19-poc-4 0.12.1-19-poc-3 0.12.1-19-poc-2 0.12.1-19-poc-1 0.12.1-19-poc 0.12.1-11-debug-1 0.12.1-11-debug 0.12.1-39 0.12.1-38.a 0.12.1-38 0.12.1-37.c 0.12.1-37.b 0.12.1-37.a 0.12.1-37.1 0.12.1-37 0.12.1-35.b 0.12.1-35.a 0.12.1-34 0.12.1-33 0.12.1-32 0.12.1-31 0.12.1-30 0.12.1-29 0.12.1-28 0.12.1-27 0.12.1-26 0.12.1-25 0.12.1-24 0.12.1-23 0.12.1-22.1 0.12.1-22 0.12.1-21 0.12.1-20 0.12.1-19 0.12.1-18 0.12.1-17 0.12.1-16 0.12.1-15 0.12.1-14 0.12.1-13 0.12.1-12 0.12.1-11 0.12.1-10 0.12.1-9 0.12.1-8 0.12.1-7 0.12.1-6 0.12.1-5 0.12.1-4 0.12.1-3 0.12.1-2 0.12.1-1

Kafka client for Elixir/Erlang.

Current section

Files

Jump to
kafka_ex_tc lib kafka_ex network_client.ex
Raw

lib/kafka_ex/network_client.ex

defmodule KafkaEx.NetworkClient do
alias KafkaEx.New
alias KafkaEx.Protocol.Metadata.Broker
alias KafkaEx.Socket
alias KafkaEx.Utils.Logger
@moduledoc false
@spec create_socket(
binary,
non_neg_integer,
KafkaEx.ssl_options(),
boolean,
boolean
) ::
nil | Socket.t()
def create_socket(
host,
port,
ssl_options \\ [],
use_ssl \\ false,
retry_until_success \\ false
) do
case Socket.create(
format_host(host),
port,
build_socket_options(ssl_options),
use_ssl
) do
{:ok, socket} ->
Logger.info(
"Successfully connected to broker #{inspect(host)}:#{inspect(port)}"
)
socket
err ->
Logger.error(
"Could not connect to broker #{inspect(host)}:#{inspect(port)} because of error #{
inspect(err)
}"
)
case retry_until_success do
true ->
sleep_for_reconnect()
create_socket(host, port, ssl_options, use_ssl, retry_until_success)
false ->
nil
end
end
end
@spec close_socket(nil | Socket.t()) :: :ok
def close_socket(nil), do: :ok
def close_socket(socket), do: Socket.close(socket)
@spec send_async_request(Broker.t() | New.Broker.t(), iodata) ::
:ok | {:error, :closed | :inet.posix()}
def send_async_request(broker, data) do
socket = broker.socket
case Socket.send(socket, data) do
:ok ->
:ok
{_, reason} ->
Logger.error(
"Asynchronously sending data to broker #{inspect(broker.host)}:#{
inspect(broker.port)
} failed with #{inspect(reason)}"
)
reason
end
end
@spec send_sync_request(Broker.t() | New.Broker.t(), iodata, timeout) ::
iodata | {:error, any()}
def send_sync_request(%{:socket => socket} = broker, data, timeout) do
case Socket.setopts(socket, [:binary, {:packet, 4}, {:active, false}]) do
:ok ->
case Socket.send(socket, data) do
:ok ->
# credo:disable-for-next-line Credo.Check.Refactor.Nesting
case Socket.recv(socket, 0, timeout) do
{:ok, data} ->
:ok =
Socket.setopts(socket, [
:binary,
{:packet, 4},
{:active, true}
])
data
{:error, reason} ->
Logger.error(
"Receiving data from broker #{inspect(broker.host)}:#{
inspect(broker.port)
} failed with #{inspect(reason)}"
)
Socket.close(socket)
{:error, reason}
end
{_, reason} ->
Logger.error(
"Sending data to broker #{inspect(broker.host)}:#{
inspect(broker.port)
} failed with #{inspect(reason)}"
)
Socket.close(socket)
{:error, reason}
end
{:error, reason} ->
Logger.error(
"Setting opts active false to broker #{inspect(broker.host)}:#{
inspect(broker.port)
} failed with #{inspect(reason)}"
)
Socket.close(socket)
{:error, reason}
end
end
def send_sync_request(nil, _, _) do
{:error, :no_broker}
end
@spec format_host(binary) :: [char] | :inet.ip_address()
def format_host(host) do
case Regex.scan(~r/^(\d{1,3})\.(\d{1,3})\.(\d{1,3})\.(\d{1,3})$/, host) do
[match_data] = [[_, _, _, _, _]] ->
match_data
|> tl
|> List.flatten()
|> Enum.map(&String.to_integer/1)
|> List.to_tuple()
# to_char_list is deprecated from Elixir 1.3 onward
_ ->
apply(String, :to_char_list, [host])
end
end
defp build_socket_options([]) do
[:binary, {:packet, 4}]
end
defp build_socket_options(ssl_options) do
build_socket_options([]) ++ ssl_options
end
defp sleep_for_reconnect() do
Process.sleep(Application.get_env(:kafka_ex, :sleep_for_reconnect, 400))
end
end