defmodule Kvasir.AgentServer.Control.Protocol do def start_link(ref, transport, config) do {:ok, spawn_link(__MODULE__, :init, [ref, transport, config])} end def init(ref, transport, {server, protocol, custom_commands}) do # Apply to not warn for client only {:ok, socket} = apply(:ranch, :handshake, [ref]) transport.send(socket, ["HELLO ", to_string(System.system_time(:second)), ?\n]) protocol.init(socket, transport, server) loop(socket, transport, server, protocol, custom_commands, []) end defp loop(socket, transport, server, protocol, custom_commands, buffer) do case transport.recv(socket, 0, :infinity) do {:ok, data} -> loop( socket, transport, server, protocol, custom_commands, process_data(socket, transport, server, protocol, custom_commands, data, buffer) ) _err -> :ok = transport.close(socket) protocol.close(server) end end defp process_data(socket, transport, server, protocol, custom_commands, data, buffer) do case String.split(data, ~r/\r?\n/) do [b] -> [b | buffer] process -> handle_commands(socket, transport, server, protocol, custom_commands, process, buffer) end end defp handle_commands(socket, transport, server, protocol, custom_commands, received, []), do: handle_loop(socket, transport, server, protocol, custom_commands, received) defp handle_commands( socket, transport, server, protocol, custom_commands, [r | received], buffer ) do b = :erlang.iolist_to_binary(:lists.reverse([r | buffer])) handle_commands(socket, transport, server, protocol, custom_commands, [b | received], []) end defp handle_loop(_socket, _transport, _server, _protocol, _custom_commands, [""]), do: [] defp handle_loop(_socket, _transport, _server, _protocol, _custom_commands, [b]), do: [b] defp handle_loop(socket, transport, server, protocol, custom_commands, [cmd | received]) do case String.split(cmd, ~r/\s+/, trim: true) do [] -> :ok ["PING"] -> transport.send(socket, "PONG\n") c -> case protocol.handle_command(c, server, custom_commands) do {:reply, d} -> transport.send(socket, d) {:async, reply, call} -> transport.send(socket, reply) spawn(fn -> call.(&transport.send(socket, &1)) end) _ -> :ok end end handle_loop(socket, transport, server, protocol, custom_commands, received) end end