Packages
electric_client
0.2.4-pre-1
0.10.3
0.10.2
0.10.1
0.10.1-beta-1
0.10.0
0.9.5-beta-1
0.9.4
0.9.4-beta-1
0.9.3
0.9.2
0.9.1
0.9.0
0.8.3
0.8.3-beta-1
0.8.2
0.8.1
0.8.0
0.8.0-beta-1
0.7.3
0.7.2
0.7.1
0.7.0
0.6.5
0.6.5-beta-5
0.6.5-beta-4
0.6.5-beta-3
0.6.5-beta-2
0.6.5-beta-1
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
0.5.0
0.5.0-beta-1
0.4.1
0.4.0
0.3.2
0.3.1
0.3.0
0.3.0-beta.4
0.3.0-beta.3
0.3.0-beta.2
0.2.6-pre-1
retired
0.2.6-beta.1
0.2.6-beta.0
0.2.5
0.2.4
0.2.4-pre-8
0.2.4-pre-7
0.2.4-pre-6
0.2.4-pre-5
0.2.4-pre-4
0.2.4-pre-3
0.2.4-pre-2
0.2.4-pre-1
0.2.3
0.2.3-rc-1
0.2.2
0.2.2-rc-1
0.2.1
0.2.1-rc-3
0.2.1-rc-2
0.2.1-rc-1
0.2.0
0.1.2
0.1.1
0.1.0
0.1.0-dev-9
0.1.0-dev-8
0.1.0-dev-7
0.1.0-dev-6
0.1.0-dev-5
0.1.0-dev-4
0.1.0-dev-3
0.1.0-dev-2
0.1.0-dev-17
0.1.0-dev-16
0.1.0-dev-15
0.1.0-dev-14
0.1.0-dev-13
0.1.0-dev-12
0.1.0-dev-11
0.1.0-dev-10
0.1.0-dev
Elixir client for ElectricSQL
Current section
Files
Jump to
Current section
Files
lib/electric/client/fetch/mint/connection.ex
defmodule Electric.Client.Fetch.Mint.Connection do
use GenServer
alias Electric.Client.Fetch
require Logger
def name(stream_id) do
{:via, Registry, {Electric.Client.Registry, {__MODULE__, stream_id}}}
end
def start_link(stream_id) do
GenServer.start_link(__MODULE__, stream_id, name: name(stream_id))
end
def fetch(conn, request, opts) do
GenServer.call(conn, {:request, request, opts}, :infinity)
end
@impl GenServer
def init(stream_id) do
Logger.debug(fn ->
"Starting "
end)
{:ok,
%{stream_id: stream_id, from: nil, conn: nil, ref: nil, resp: nil, timeout: nil, req: nil}}
end
@impl GenServer
def handle_call({:request, request, opts}, from, state) do
state = request({request, opts}, state)
{:noreply, %{state | from: from, req: {request, opts}}}
end
@impl GenServer
def handle_info({:ssl, _socket, _data} = msg, %{conn: conn} = state) do
{:ok, conn, responses} = Mint.HTTP.stream(conn, msg)
handle_responses(responses, %{state | conn: conn})
end
def handle_info({:ssl_closed, _socket}, state) do
Logger.info("[#{state.stream_id}] SSL closed")
{:noreply, state |> maybe_close() |> maybe_retry()}
end
def handle_info({:ssl_error, _socket, reason}, state) do
Logger.error("[#{state.stream_id}] SSL error: #{inspect(reason)}")
{:noreply, state |> maybe_close() |> maybe_retry()}
end
def handle_info({:tcp, _socket, _data} = msg, %{conn: conn} = state) do
{:ok, conn, responses} = Mint.HTTP.stream(conn, msg)
handle_responses(responses, %{state | conn: conn})
end
def handle_info({:tcp_closed, _socket}, state) do
Logger.info("[#{state.stream_id}] TCP closed")
{:noreply, state |> maybe_close() |> maybe_retry()}
end
def handle_info({:timeout, _ref}, state) do
Logger.warning("[#{state.stream_id}] TIMEOUT")
{:noreply, maybe_retry(state)}
end
defp request({request, opts}, state) do
uri = Fetch.Request.uri(request, opts)
state
|> connect(uri, opts)
|> make_request(request, uri, opts)
end
defp connect(%{conn: nil} = state, uri, _opts) do
Logger.debug(
"[#{state.stream_id}] Opening #{uri.scheme} connection to #{uri.host}:#{uri.port}"
)
{:ok, conn} =
Mint.HTTP.connect(String.to_atom(uri.scheme), uri.host, uri.port, transport_opts: [])
%{state | conn: conn}
end
defp connect(state, _request, _opts) do
state
end
defp make_request(%{conn: conn} = state, request, uri, _opts) do
{:ok, conn, request_ref} =
Mint.HTTP.request(
conn,
method(request.method),
uri.path <> "?" <> uri.query,
Enum.to_list(request.headers),
nil
)
ref = Process.send_after(self(), {:timeout, request_ref}, 30_000)
%{state | timeout: ref, conn: conn, ref: request_ref, resp: %Fetch.Response{}}
end
defp method(:get), do: "GET"
defp method(:post), do: "POST"
defp method(:delete), do: "DELETE"
defp handle_responses(responses, %{from: from} = state) do
case Enum.reduce(responses, {:cont, state.resp}, &handle_response/2) do
{:cont, resp} ->
{:noreply, %{state | resp: resp}}
{:done, resp} ->
GenServer.reply(from, resp)
{:noreply, reset(state)}
end
end
defp handle_response({:status, _ref, status}, {:cont, resp}) do
{:cont, %{resp | status: status}}
end
defp handle_response({:headers, _ref, headers}, {:cont, resp}) do
{:cont, %{resp | headers: Enum.reduce(headers, resp.headers, &add_header/2)}}
end
defp handle_response({:data, _ref, data}, {:cont, resp}) do
{:cont, %{resp | body: [resp.body | data]}}
end
defp handle_response({:done, _ref}, {:cont, resp}) do
body =
case IO.iodata_to_binary(resp.body) do
"" -> ""
json -> Jason.decode!(json)
end
{:done, Fetch.Response.decode!(%{resp | body: body})}
end
defp add_header({key, value}, headers) do
Map.update(headers, key, [value], &[value | &1])
end
defp maybe_retry(%{from: nil} = state) do
cancel_timer(state)
end
defp maybe_retry(%{from: _from, req: {request, opts}} = state) do
Logger.debug(fn ->
"Retrying request #{inspect(request)}"
end)
request({request, opts}, cancel_timer(%{state | ref: nil, resp: nil}))
end
defp cancel_timer(%{timeout: timeout} = state) when is_reference(timeout) do
Process.cancel_timer(timeout)
%{state | timeout: nil}
end
defp cancel_timer(state) do
state
end
defp reset(state) do
cancel_timer(%{state | resp: nil, ref: nil, from: nil, req: nil})
end
defp maybe_close(%{conn: conn} = state) do
if conn && Mint.HTTP.open?(conn) do
{:ok, _conn} = Mint.HTTP.close(conn)
%{state | conn: nil}
else
state
end
end
end