2020-02-11 07:12:57 +00:00
|
|
|
# Pleroma: A lightweight social networking server
|
2020-03-03 23:16:24 +00:00
|
|
|
# Copyright © 2017-2020 Pleroma Authors <https://pleroma.social/>
|
2020-02-11 07:12:57 +00:00
|
|
|
# SPDX-License-Identifier: AGPL-3.0-only
|
|
|
|
|
|
|
|
defmodule Pleroma.ReverseProxy.Client.Tesla do
|
|
|
|
@type headers() :: [{String.t(), String.t()}]
|
|
|
|
@type status() :: pos_integer()
|
|
|
|
|
|
|
|
@behaviour Pleroma.ReverseProxy.Client
|
|
|
|
|
|
|
|
@spec request(atom(), String.t(), headers(), String.t(), keyword()) ::
|
|
|
|
{:ok, status(), headers}
|
|
|
|
| {:ok, status(), headers, map()}
|
|
|
|
| {:error, atom() | String.t()}
|
|
|
|
| no_return()
|
|
|
|
|
|
|
|
@impl true
|
|
|
|
def request(method, url, headers, body, opts \\ []) do
|
2020-03-03 09:53:37 +00:00
|
|
|
check_adapter()
|
2020-02-11 07:12:57 +00:00
|
|
|
|
2020-03-03 11:56:49 +00:00
|
|
|
opts = Keyword.merge(opts, body_as: :chunks)
|
|
|
|
|
|
|
|
with {:ok, response} <-
|
2020-02-11 07:12:57 +00:00
|
|
|
Pleroma.HTTP.request(
|
|
|
|
method,
|
|
|
|
url,
|
|
|
|
body,
|
|
|
|
headers,
|
|
|
|
Keyword.put(opts, :adapter, opts)
|
|
|
|
) do
|
|
|
|
if is_map(response.body) and method != :head do
|
|
|
|
{:ok, response.status, response.headers, response.body}
|
|
|
|
else
|
|
|
|
{:ok, response.status, response.headers}
|
|
|
|
end
|
|
|
|
else
|
|
|
|
{:error, error} -> {:error, error}
|
|
|
|
end
|
|
|
|
end
|
|
|
|
|
|
|
|
@impl true
|
|
|
|
@spec stream_body(map()) :: {:ok, binary(), map()} | {:error, atom() | String.t()} | :done
|
|
|
|
def stream_body(%{pid: pid, opts: opts, fin: true}) do
|
2020-03-03 11:56:49 +00:00
|
|
|
# if connection was reused, but in tesla were redirects,
|
|
|
|
# tesla returns new opened connection, which must be closed manually
|
2020-02-11 07:12:57 +00:00
|
|
|
if opts[:old_conn], do: Tesla.Adapter.Gun.close(pid)
|
|
|
|
# if there were redirects we need to checkout old conn
|
|
|
|
conn = opts[:old_conn] || opts[:conn]
|
|
|
|
|
|
|
|
if conn, do: :ok = Pleroma.Pool.Connections.checkout(conn, self(), :gun_connections)
|
|
|
|
|
|
|
|
:done
|
|
|
|
end
|
|
|
|
|
|
|
|
def stream_body(client) do
|
|
|
|
case read_chunk!(client) do
|
|
|
|
{:fin, body} ->
|
|
|
|
{:ok, body, Map.put(client, :fin, true)}
|
|
|
|
|
|
|
|
{:nofin, part} ->
|
|
|
|
{:ok, part, client}
|
|
|
|
|
|
|
|
{:error, error} ->
|
|
|
|
{:error, error}
|
|
|
|
end
|
|
|
|
end
|
|
|
|
|
|
|
|
defp read_chunk!(%{pid: pid, stream: stream, opts: opts}) do
|
|
|
|
adapter = check_adapter()
|
|
|
|
adapter.read_chunk(pid, stream, opts)
|
|
|
|
end
|
|
|
|
|
|
|
|
@impl true
|
|
|
|
@spec close(map) :: :ok | no_return()
|
|
|
|
def close(%{pid: pid}) do
|
|
|
|
adapter = check_adapter()
|
|
|
|
adapter.close(pid)
|
|
|
|
end
|
|
|
|
|
|
|
|
defp check_adapter do
|
|
|
|
adapter = Application.get_env(:tesla, :adapter)
|
|
|
|
|
|
|
|
unless adapter == Tesla.Adapter.Gun do
|
|
|
|
raise "#{adapter} doesn't support reading body in chunks"
|
|
|
|
end
|
|
|
|
|
|
|
|
adapter
|
|
|
|
end
|
|
|
|
end
|