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