From dfbed9999f1c14d63385402125960293c53031e8 Mon Sep 17 00:00:00 2001 From: vikrambombhi Date: Tue, 4 Dec 2018 11:03:28 -0500 Subject: [PATCH 1/6] handle read error in read method --- lib/active_proxy/proxy.ex | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/lib/active_proxy/proxy.ex b/lib/active_proxy/proxy.ex index 0f74586..a62d435 100644 --- a/lib/active_proxy/proxy.ex +++ b/lib/active_proxy/proxy.ex @@ -50,6 +50,7 @@ defmodule ActiveProxy.Proxy do if {ok, packet} == {:error, :closed} do # If reading from socket closed exit serve loop + # TODO: # Should probably log error # # @@ -58,13 +59,10 @@ defmodule ActiveProxy.Proxy do # TODO: Handle failure to write to application write(upstream_socket, packet) - case read(upstream_socket, timeout) do - {:ok, payload} -> - write(socket, payload) + {status, payload} = read(upstream_socket, timeout) - {:error, _} -> - # shutdown writes to signal that no more data is to be sent and wait for the read side of the socket to be closed - :gen_tcp.shutdown(socket, :write) + if status == :ok do + write(socket, payload) end serve(socket, upstream_socket) @@ -80,8 +78,10 @@ defmodule ActiveProxy.Proxy do {:ok, packet} -> {:ok, packet} - {:error, timeout} -> - {:error, timeout} + {:error, reason} -> + # shutdown writes to signal that no more data is to be sent and wait for the read side of the socket to be closed + :gen_tcp.shutdown(socket, :write) + {:error, reason} end end From b580ceed2481336a60b066659c13de7c008b3f56 Mon Sep 17 00:00:00 2001 From: vikrambombhi Date: Tue, 4 Dec 2018 11:04:21 -0500 Subject: [PATCH 2/6] ignore elixir language server files --- .gitignore | 2 ++ 1 file changed, 2 insertions(+) diff --git a/.gitignore b/.gitignore index 0b142a3..30d4179 100644 --- a/.gitignore +++ b/.gitignore @@ -22,3 +22,5 @@ erl_crash.dump # Ignore package tarball (built via "mix hex.build"). active_proxy-*.tar +# Elixir language server +.elixir_ls From ea682d4e603343bee7fbc65a0ca1a459a802fcd6 Mon Sep 17 00:00:00 2001 From: vikrambombhi Date: Wed, 5 Dec 2018 13:23:13 -0500 Subject: [PATCH 3/6] config file with single upstream node address and ms before timeout from upstream --- config/config.exs | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/config/config.exs b/config/config.exs index 46f253c..ffd0646 100644 --- a/config/config.exs +++ b/config/config.exs @@ -28,3 +28,7 @@ use Mix.Config # here (which is why it is important to import them last). # # import_config "#{Mix.env()}.exs" + +config :active_proxy, + timeout: 100, + node1_address: "159.203.44.11" From 85c0408971181cdd859b8e089bafbec706ed5b45 Mon Sep 17 00:00:00 2001 From: vikrambombhi Date: Wed, 5 Dec 2018 13:24:20 -0500 Subject: [PATCH 4/6] use timeout specified in config file when reading from upstream node --- lib/active_proxy/proxy.ex | 59 +++++++++++++++------------------------ 1 file changed, 23 insertions(+), 36 deletions(-) diff --git a/lib/active_proxy/proxy.ex b/lib/active_proxy/proxy.ex index a62d435..1a31bcc 100644 --- a/lib/active_proxy/proxy.ex +++ b/lib/active_proxy/proxy.ex @@ -22,12 +22,16 @@ defmodule ActiveProxy.Proxy do end defp start_serve_process(socket) do + timeout = Application.get_env(:active_proxy, :timeout) + application_address = Application.get_env(:active_proxy, :node1_address) + Logger.info("Forwarding to application at #{application_address} with timout of #{timeout}ms") + {:ok, pid} = Task.Supervisor.start_child(ActiveProxy.TaskSupervisor, fn -> # TODO: make the host configurable {:ok, upstream_socket} = :gen_tcp.connect( - '159.203.44.11', + to_charlist(application_address), 4000, [ :binary, @@ -38,50 +42,33 @@ defmodule ActiveProxy.Proxy do 1000 ) - serve(socket, upstream_socket) + serve(socket, upstream_socket, timeout) end) :ok = :gen_tcp.controlling_process(socket, pid) end - defp serve(socket, upstream_socket) do - timeout = 10 - {ok, packet} = read(socket) - - if {ok, packet} == {:error, :closed} do - # If reading from socket closed exit serve loop - # TODO: - # Should probably log error - # - # - # Also handle other types of errors other than :closed - else - # TODO: Handle failure to write to application - write(upstream_socket, packet) + defp serve(socket, upstream_socket, timeout) do + case :gen_tcp.recv(socket, 0) do + {:ok, packet} -> + # TODO: Handle failure to write to application + write(upstream_socket, packet) - {status, payload} = read(upstream_socket, timeout) + case :gen_tcp.recv(upstream_socket, 0, timeout) do + {:ok, packet} -> + write(socket, packet) + serve(socket, upstream_socket, timeout) - if status == :ok do - write(socket, payload) - end + {:error, :timeout} -> + # TODO: Consider failing over to different upstream node + nil + end - serve(socket, upstream_socket) - end - end - - defp read(socket) do - :gen_tcp.recv(socket, 0) - end - - defp read(socket, timeout) do - case :gen_tcp.recv(socket, 0, timeout) do - {:ok, packet} -> - {:ok, packet} + {:error, :closed} -> + # In case of socket being closed exit serve loop + nil - {:error, reason} -> - # shutdown writes to signal that no more data is to be sent and wait for the read side of the socket to be closed - :gen_tcp.shutdown(socket, :write) - {:error, reason} + # TODO: handle other types of read error to end end From c74eef2926f4cc77cbb24943f854570e05d60e1b Mon Sep 17 00:00:00 2001 From: vikrambombhi Date: Wed, 5 Dec 2018 14:17:28 -0500 Subject: [PATCH 5/6] cleaner functional flow for read_from_client |> forward_to_upstream |> write_to_client --- lib/active_proxy/proxy.ex | 38 +++++++++++++++++++++++++------------- 1 file changed, 25 insertions(+), 13 deletions(-) diff --git a/lib/active_proxy/proxy.ex b/lib/active_proxy/proxy.ex index 1a31bcc..9bae76f 100644 --- a/lib/active_proxy/proxy.ex +++ b/lib/active_proxy/proxy.ex @@ -49,20 +49,13 @@ defmodule ActiveProxy.Proxy do end defp serve(socket, upstream_socket, timeout) do + read_from_client(socket, upstream_socket, timeout) + end + + defp read_from_client(socket, upstream_socket, timeout) do case :gen_tcp.recv(socket, 0) do {:ok, packet} -> - # TODO: Handle failure to write to application - write(upstream_socket, packet) - - case :gen_tcp.recv(upstream_socket, 0, timeout) do - {:ok, packet} -> - write(socket, packet) - serve(socket, upstream_socket, timeout) - - {:error, :timeout} -> - # TODO: Consider failing over to different upstream node - nil - end + packet |> forward_to_upstream(upstream_socket, socket, timeout) {:error, :closed} -> # In case of socket being closed exit serve loop @@ -72,7 +65,26 @@ defmodule ActiveProxy.Proxy do end end - defp write(socket, packet) do + defp forward_to_upstream(packet, upstream_socket, socket, timeout) do + # TODO: Handle failure to write to application + write(packet, upstream_socket) + + case :gen_tcp.recv(upstream_socket, 0, timeout) do + {:ok, packet} -> + packet |> write_to_client(upstream_socket, socket, timeout) + + {:error, :timeout} -> + # TODO: Consider failing over to different upstream node + nil + end + end + + defp write_to_client(packet, upstream_socket, socket, timeout) do + write(packet, socket) + serve(socket, upstream_socket, timeout) + end + + defp write(packet, socket) do :gen_tcp.send(socket, packet) end end From 75cda5f3260654c7d7bb248872986ded2f4e7ecf Mon Sep 17 00:00:00 2001 From: vikrambombhi Date: Wed, 5 Dec 2018 14:39:05 -0500 Subject: [PATCH 6/6] use monad to handle errors, Thanks @JackyChiu --- lib/active_proxy/proxy.ex | 43 +++++++++------------------------------ 1 file changed, 10 insertions(+), 33 deletions(-) diff --git a/lib/active_proxy/proxy.ex b/lib/active_proxy/proxy.ex index 9bae76f..d76ce32 100644 --- a/lib/active_proxy/proxy.ex +++ b/lib/active_proxy/proxy.ex @@ -48,43 +48,20 @@ defmodule ActiveProxy.Proxy do :ok = :gen_tcp.controlling_process(socket, pid) end - defp serve(socket, upstream_socket, timeout) do - read_from_client(socket, upstream_socket, timeout) - end - - defp read_from_client(socket, upstream_socket, timeout) do - case :gen_tcp.recv(socket, 0) do - {:ok, packet} -> - packet |> forward_to_upstream(upstream_socket, socket, timeout) - + defp serve(client_socket, upstream_socket, timeout) do + with {:ok, packet} <- :gen_tcp.recv(client_socket, 0), + :ok <- :gen_tcp.send(upstream_socket, packet), + {:ok, packet} <- :gen_tcp.recv(upstream_socket, 0, timeout), + :ok <- :gen_tcp.send(client_socket, packet) do + serve(client_socket, upstream_socket, timeout) + else {:error, :closed} -> - # In case of socket being closed exit serve loop nil - # TODO: handle other types of read error to - end - end - - defp forward_to_upstream(packet, upstream_socket, socket, timeout) do - # TODO: Handle failure to write to application - write(packet, upstream_socket) - - case :gen_tcp.recv(upstream_socket, 0, timeout) do - {:ok, packet} -> - packet |> write_to_client(upstream_socket, socket, timeout) - - {:error, :timeout} -> - # TODO: Consider failing over to different upstream node + # TODO: Handle errors + # Consider doing a failover + _ -> nil end end - - defp write_to_client(packet, upstream_socket, socket, timeout) do - write(packet, socket) - serve(socket, upstream_socket, timeout) - end - - defp write(packet, socket) do - :gen_tcp.send(socket, packet) - end end