Current section

Files

Jump to
pulsar_elixir lib pulsar connection.ex~
Raw

lib/pulsar/connection.ex~

defmodule Pulsar.Socket do
@doc """
This module handles all socket-related operations.
"""
defstruct socket: nil, ssl: false
@type t :: %__MODULE__{
socket: :gen_tcp.socket() | :ssl.sslsocket(),
ssl: boolean()
}
@doc """
Establishes a connection to the provided host.
For more information about the available options, see `:ssl.connect/3` for ssl
or `:gen_tcp.connect/3` for non ssl.
"""
# connect("pulsar+ssl://istio-stagingmig.euw1-turtle.streamnative.g.snio.cloud:6651")
def connect(uri, socket_opts \\ [:binary, nodelay: true, active: false, keepalive: true, verify: :verify_none]) do
{host, port, ssl} = parse_uri(uri)
socket = do_connect(host, port, socket_opts, ssl)
%__MODULE__{socket: socket, ssl: ssl}
end
defp parse_uri(host) do
uri = URI.parse(host)
ssl = (Map.get(uri, :scheme) == "pulsar+ssl")
host = Map.get(uri, :host)
port = Map.get(uri, :port, 6650)
{host, port, ssl}
end
defp do_connect(host, port, socket_opts, ssl) do
socket_module = socket_module(ssl)
host = String.to_charlist(host)
{:ok, socket} = apply(socket_module, :connect, [host, port, socket_opts, 5_000])
socket
end
@doc """
Closes the connection.
"""
@spec close(__MODULE__.t()) :: :ok
def close(%{socket: socket, ssl: ssl}), do: apply(socket_module(ssl), :close, [socket])
defp socket_module(true), do: :ssl
defp socket_module(false), do: :gen_tcp
end