Packages

kafka_ex_tc

0.12.1-11-debug
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 socket.ex
Raw

lib/kafka_ex/socket.ex

defmodule KafkaEx.Socket do
@moduledoc """
This module handle all socket related operations.
"""
defstruct socket: nil, ssl: false
@type t :: %KafkaEx.Socket{
socket: :gen_tcp.socket() | :ssl.sslsocket(),
ssl: boolean
}
@doc """
Creates a socket.
For more information about the available options, see `:ssl.connect/3` for ssl
or `:gen_tcp.connect/3` for non ssl.
"""
@spec create(:inet.ip_address(), non_neg_integer, [] | [...]) ::
{:ok, KafkaEx.Socket.t()} | {:error, any}
def create(host, port, socket_options \\ [], is_ssl \\ false) do
case create_socket(host, port, is_ssl, socket_options) do
{:ok, socket} -> {:ok, %KafkaEx.Socket{socket: socket, ssl: is_ssl}}
{:error, reason} -> {:error, reason}
end
end
@doc """
Closes the socket.
For more information, see `:ssl.close/1` for ssl or `:gen_tcp.send/1` for non ssl.
"""
@spec close(KafkaEx.Socket.t()) :: :ok
def close(%KafkaEx.Socket{ssl: true} = socket), do: :ssl.close(socket.socket)
def close(socket), do: :gen_tcp.close(socket.socket)
@doc """
Sends data over the socket.
It handles both, SSL and non SSL sockets.
For more information, see `:ssl.send/2` for ssl or `:gen_tcp.send/2` for non ssl.
"""
@spec send(KafkaEx.Socket.t(), iodata) :: :ok | {:error, any}
def send(%KafkaEx.Socket{ssl: true} = socket, data) do
:ssl.send(socket.socket, data)
end
def send(socket, data) do
:gen_tcp.send(socket.socket, data)
end
@doc """
Set options to the socket.
For more information, see `:ssl.setopts/2` for ssl or `:inet.setopts/2` for non ssl.
"""
@spec setopts(KafkaEx.Socket.t(), list) :: :ok | {:error, any}
def setopts(%KafkaEx.Socket{ssl: true} = socket, options) do
:ssl.setopts(socket.socket, options)
end
def setopts(socket, options) do
:inet.setopts(socket.socket, options)
end
@doc """
Receives data from the socket.
For more information, see `:ssl.recv/2` for ssl or `:gen_tcp.recv/2` for non ssl.
"""
@spec recv(KafkaEx.Socket.t(), non_neg_integer) ::
{:ok, String.t() | binary | term} | {:error, any}
def recv(%KafkaEx.Socket{ssl: true} = socket, length) do
:ssl.recv(socket.socket, length)
end
def recv(socket, length) do
:gen_tcp.recv(socket.socket, length)
end
@spec recv(KafkaEx.Socket.t(), non_neg_integer, timeout) ::
{:ok, String.t() | binary | term} | {:error, any}
def recv(%KafkaEx.Socket{ssl: true} = socket, length, timeout) do
:ssl.recv(socket.socket, length, timeout)
end
def recv(socket, length, timeout) do
:gen_tcp.recv(socket.socket, length, timeout)
end
@doc """
Returns true if the socket is open
"""
@spec open?(KafkaEx.Socket.t()) :: boolean
def open?(%KafkaEx.Socket{} = socket) do
info(socket) != nil
end
@doc """
Returns the information about the socket.
For more information, see `Port.info`
"""
@spec info(KafkaEx.Socket.t()) :: list | nil
def info(socket) do
socket
|> extract_port
|> Port.info()
end
defp extract_port(%KafkaEx.Socket{ssl: true} = socket) do
{:sslsocket, {:gen_tcp, port, _, _}, _} = socket.socket
port
end
defp extract_port(socket), do: socket.socket
defp create_socket(host, port, true, socket_options) do
:ssl.connect(host, port, socket_options)
end
defp create_socket(host, port, false, socket_options) do
:gen_tcp.connect(host, port, socket_options)
end
end