Current section
Files
Jump to
Current section
Files
lib/kafka_ex/network/socket.ex
defmodule KafkaEx.Network.Socket do
@moduledoc """
This module handle all socket related operations.
"""
defstruct socket: nil, ssl: false
@type t :: %KafkaEx.Network.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.
"""
# Bounds :gen_tcp.connect / :ssl.connect, which otherwise default to
# :infinity and block the calling process for the full OS TCP timeout
# (tens of seconds to minutes) against an unreachable/black-holed broker.
@default_connect_timeout 10_000
@spec create(:inet.ip_address(), non_neg_integer, [] | [...], boolean, timeout) ::
{:ok, KafkaEx.Network.Socket.t()} | {:error, any}
def create(host, port, socket_options \\ [], is_ssl \\ false, connect_timeout \\ @default_connect_timeout) do
case create_socket(host, port, is_ssl, socket_options, connect_timeout) do
{:ok, socket} -> {:ok, %KafkaEx.Network.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.
Handles both Socket structs and raw ports (for tcp_closed/ssl_closed handlers).
"""
@spec close(KafkaEx.Network.Socket.t() | port() | reference() | nil) :: :ok
def close(nil), do: :ok
def close(%KafkaEx.Network.Socket{ssl: true} = socket), do: :ssl.close(socket.socket)
def close(%KafkaEx.Network.Socket{} = socket), do: :gen_tcp.close(socket.socket)
# Handle raw ports (from tcp_closed messages)
# Socket may already be closed, ignore errors
def close(port) when is_port(port) do
:gen_tcp.close(port)
catch
_, _ -> :ok
end
# Handle raw SSL socket references (from ssl_closed messages)
def close(ssl_socket) when is_reference(ssl_socket) or is_tuple(ssl_socket) do
:ssl.close(ssl_socket)
catch
_, _ -> :ok
end
@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.Network.Socket.t(), iodata) :: :ok | {:error, any}
def send(%KafkaEx.Network.Socket{ssl: true} = socket, data), do: :ssl.send(socket.socket, data)
def send(socket, data), do: :gen_tcp.send(socket.socket, data)
@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.Network.Socket.t(), list) :: :ok | {:error, any}
def setopts(%KafkaEx.Network.Socket{ssl: true} = socket, options), do: :ssl.setopts(socket.socket, options)
def setopts(socket, options), do: :inet.setopts(socket.socket, options)
@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.Network.Socket.t(), non_neg_integer) :: {:ok, String.t() | binary | term} | {:error, any}
def recv(%KafkaEx.Network.Socket{ssl: true} = socket, length), do: :ssl.recv(socket.socket, length)
def recv(socket, length), do: :gen_tcp.recv(socket.socket, length)
@spec recv(KafkaEx.Network.Socket.t(), non_neg_integer, timeout) :: {:ok, String.t() | binary | term} | {:error, any}
def recv(%KafkaEx.Network.Socket{ssl: true} = socket, length, timeout), do: :ssl.recv(socket.socket, length, timeout)
def recv(socket, length, timeout), do: :gen_tcp.recv(socket.socket, length, timeout)
@doc """
Returns true if the socket is open
"""
@spec open?(KafkaEx.Network.Socket.t()) :: boolean
def open?(%KafkaEx.Network.Socket{} = socket), do: !is_nil(info(socket))
@doc """
Returns the information about the socket.
For more information, see `Port.info`
"""
@spec info(KafkaEx.Network.Socket.t()) :: list | nil
def info(socket) do
socket |> extract_port() |> Port.info()
end
defp extract_port(%KafkaEx.Network.Socket{ssl: true} = socket) do
case socket.socket do
# OTP pre 28
{:sslsocket, {:gen_tcp, port, _, _}, _} -> port
# OTP 28
{:sslsocket, port, _, _, :gen_tcp, _, _, _} -> port
end
end
defp extract_port(socket), do: socket.socket
defp create_socket(host, port, true, socket_options, timeout),
do: :ssl.connect(host, port, socket_options, timeout)
defp create_socket(host, port, false, socket_options, timeout),
do: :gen_tcp.connect(host, port, socket_options, timeout)
end