Current section

Files

Jump to
kvasir_agent_server lib kvasir agent_server command connection.ex
Raw

lib/kvasir/agent_server/command/connection.ex

defmodule Kvasir.AgentServer.Command.Connection do
@transport :gen_tcp
@socket_opts [:binary, active: false]
@latency 1_000
def create(host, port) do
{:ok, pid} = start_link(host, port)
{:ok, GenServer.call(pid, :get_conn)}
end
def send_command(conn, command = %{__meta__: m}) do
{registry, socket} = conn
response = m.wait != :dispatch
if response do
:ets.insert(registry, {m.id, self()})
end
cmd = %{
command
| __meta__: m |> Map.from_struct() |> Enum.reject(&(elem(&1, 1) == nil)) |> Map.new()
}
packed = :erlang.term_to_binary(cmd, minor_version: 2, compressed: 9)
size = byte_size(packed)
@transport.send(socket, [<<size::unsigned-integer-32>>, packed])
if response do
wait_for = m.timeout + @latency
receive do
{:ok, meta} -> {:ok, %{command | __meta__: struct!(Kvasir.Command.Meta, meta)}}
{_, err} -> err
after
wait_for -> {:error, :remote_command_timeout}
end
else
{:ok, command}
end
end
def start_link(host, port) do
GenServer.start_link(__MODULE__, {host, port})
end
@behaviour GenServer
@impl GenServer
def init({host, port}) do
h = if(is_binary(host), do: String.to_charlist(host), else: host)
{:ok, socket} = @transport.connect(h, port, @socket_opts)
registry = :ets.new(:registry, [:set, :public, write_concurrency: true])
reader = spawn_link(fn -> response_loop(registry, h, port, socket) end)
{:ok, %{socket: socket, registry: registry, reader: reader}}
end
defp response_loop(registry, host, port, socket) do
case @transport.recv(socket, 4, :infinity) do
{:ok, <<length::unsigned-integer-32>>} ->
{:ok, data} = @transport.recv(socket, length, :infinity)
response = :erlang.binary_to_term(data)
id =
case response do
{:ok, %{id: id}} -> id
{id, _} -> id
end
[{^id, pid}] = :ets.take(registry, id)
send(pid, response)
response_loop(registry, host, port, socket)
{:error, :closed} ->
connect_loop(registry, host, port)
end
end
defp connect_loop(registry, host, port, attempt \\ 1) do
case @transport.connect(host, port, @socket_opts) do
{:ok, socket} ->
response_loop(registry, host, port, socket)
_ ->
:timer.sleep(attempt * 500)
if attempt <= 5 do
connect_loop(registry, host, port, attempt + 1)
else
raise "Connection lost."
end
end
end
@impl GenServer
def handle_call(:get_conn, _from, state = %{registry: r, socket: s}) do
{:reply, {r, s}, state}
end
end