defmodule Phoenix.Sync.Electric.ClientAdapter do @moduledoc false defstruct [:client, :shape_definition] defimpl Phoenix.Sync.Adapter.PlugApi do alias Electric.Client alias Phoenix.Sync.PredefinedShape def predefined_shape(sync_client, %PredefinedShape{} = predefined_shape) do shape_client = PredefinedShape.client(sync_client.client, predefined_shape) {:ok, %Phoenix.Sync.Electric.ClientAdapter{ client: shape_client, shape_definition: predefined_shape }} end # this is the server-defined shape route, so we want to only pass on the # per-request/stream position params leaving the shape-definition params # from the configured client. def call(%{shape_definition: %PredefinedShape{} = shape} = sync_client, conn, params) do request = Client.request( sync_client.client, method: :get, offset: params["offset"], shape_handle: params["handle"], live: live?(params["live"]), next_cursor: params["cursor"] ) fetch_upstream(sync_client, conn, request, shape) end # this version is the pure client-defined shape version def call(sync_client, %{method: method} = conn, params) do request = Client.request( sync_client.client, method: normalise_method(method), params: params ) fetch_upstream(sync_client, conn, request, nil) end defp normalise_method(method), do: method |> String.downcase() |> String.to_atom() defp live?(live), do: live == "true" defp fetch_upstream(sync_client, conn, request, shape) do request = put_req_headers(request, conn.req_headers) response = case Client.Fetch.request(sync_client.client, request) do %Client.Fetch.Response{} = response -> response {:error, %Client.Fetch.Response{} = response} -> response end body = if response.status == 200 do Phoenix.Sync.Electric.map_response_body( response.body, PredefinedShape.transform_fun(shape) ) else response.body end conn |> put_resp_headers(response.headers) |> Plug.Conn.send_resp(response.status, body) end defp put_req_headers(request, headers) do merged_headers = Enum.reduce(headers, request.headers, fn {header, value}, acc -> Map.update(acc, header, [value], fn existing -> [value | List.wrap(existing)] end) end) |> expand_headers() %{request | headers: merged_headers} end defp put_resp_headers(conn, headers) do resp_headers = headers |> Map.delete("transfer-encoding") |> expand_headers() Plug.Conn.merge_resp_headers(conn, resp_headers) end # turn headers into a list which is more compatible than a map # representation as it preserves multiple values for a header. defp expand_headers(headers) when is_map(headers) do Enum.flat_map(headers, fn {k, v} -> Enum.map(List.wrap(v), &{k, &1}) end) end end end