Page Menu
Home
Phorge
Search
Configure Global Search
Log In
Files
F85628858
No One
Temporary
Actions
View File
Edit File
Delete File
View Transforms
Subscribe
Award Token
Flag For Later
Size
58 KB
Referenced Files
None
Subscribers
None
View Options
diff --git a/changelog.d/media-proxy-truncated-stream.fix b/changelog.d/media-proxy-truncated-stream.fix
new file mode 100644
index 000000000..12413e53e
--- /dev/null
+++ b/changelog.d/media-proxy-truncated-stream.fix
@@ -0,0 +1 @@
+Abort incomplete media proxy responses instead of caching partial bodies.
diff --git a/lib/pleroma/reverse_proxy.ex b/lib/pleroma/reverse_proxy.ex
index e87c3f625..12e48eadd 100644
--- a/lib/pleroma/reverse_proxy.ex
+++ b/lib/pleroma/reverse_proxy.ex
@@ -1,614 +1,713 @@
# Pleroma: A lightweight social networking server
# Copyright © 2017-2022 Pleroma Authors <https://pleroma.social/>
# SPDX-License-Identifier: AGPL-3.0-only
defmodule Pleroma.ReverseProxy do
alias Pleroma.Utils.URIEncoding
@range_headers ~w(range if-range)
@keep_req_headers ~w(accept accept-encoding cache-control if-modified-since) ++
~w(if-unmodified-since if-none-match) ++ @range_headers
@resp_cache_headers ~w(etag date last-modified)
@keep_resp_headers @resp_cache_headers ++
~w(content-length content-type content-disposition content-encoding) ++
~w(content-range accept-ranges vary)
- @default_cache_control_header "public, max-age=1209600, immutable"
+ @default_cache_control_header "public, max-age=1209600, immutable, no-transform"
@valid_resp_codes [200, 206, 304]
@max_read_duration :timer.seconds(30)
@max_body_length :infinity
@failed_request_ttl :timer.seconds(60)
@sniff_bytes 8 * 1024
@methods ~w(GET HEAD)
@allowed_mime_types Pleroma.Config.get([Pleroma.Upload, :allowed_mime_types], [])
@cachex Pleroma.Config.get([:cachex, :provider], Cachex)
+ defmodule StreamError do
+ @moduledoc false
+ defexception [:url, :reason]
+
+ @impl true
+ def message(%{url: url, reason: reason}) do
+ "reverse proxy stream from #{url} failed: #{inspect(reason)}"
+ end
+ end
+
def max_read_duration_default, do: @max_read_duration
def default_cache_control_header, do: @default_cache_control_header
@moduledoc """
A reverse proxy.
Pleroma.ReverseProxy.call(conn, url, options)
It is not meant to be added into a plug pipeline, but to be called from another plug or controller.
Supports `#{inspect(@methods)}` HTTP methods, and only allows `#{inspect(@valid_resp_codes)}` status codes.
Responses are chunked to the client while downloading from the upstream.
Some request / responses headers are preserved:
* request: `#{inspect(@keep_req_headers)}`
* response: `#{inspect(@keep_resp_headers)}`
Options:
* `redirect_on_failure` (default `false`). Redirects the client to the real remote URL if there's any HTTP
errors. Any error during body processing will not be redirected as the response is chunked. This may expose
remote URL, clients IPs, ….
* `max_body_length` (default `#{inspect(@max_body_length)}`): limits the content length to be approximately the
specified length. It is validated with the `content-length` header and also verified when proxying.
* `max_read_duration` (default `#{inspect(@max_read_duration)}` ms): the total time the connection is allowed to
read from the remote upstream.
* `failed_request_ttl` (default `#{inspect(@failed_request_ttl)}` ms): the time the failed request is cached and cannot be retried.
* `inline_content_types`:
* `true` will not alter `content-disposition` (up to the upstream),
* `false` will add `content-disposition: attachment` to any request,
* a list of whitelisted content types for which `content-disposition: inline`
is always set (overriding any upstream header) so the media can be embedded
in pages; the filename is derived from the content type
* `sniff_content_type` (default `false`): detects an image MIME type from the first
response chunk when the upstream type is missing or `application/octet-stream`.
* `req_headers`, `resp_headers` additional headers.
* `http`: options for [hackney](https://github.com/benoitc/hackney) or [gun](https://github.com/ninenines/gun).
"""
@default_options [pool: :media]
@inline_content_types [
"image/apng",
"image/avif",
"image/bmp",
"image/gif",
"image/jpeg",
"image/jpg",
"image/png",
"image/svg+xml",
"image/webp",
"audio/mpeg",
"audio/mp3",
"video/webm",
"video/mp4",
"video/quicktime"
]
require Logger
import Plug.Conn
@type option() ::
{:max_read_duration, non_neg_integer() | :infinity}
| {:max_body_length, non_neg_integer() | :infinity}
| {:failed_request_ttl, non_neg_integer() | :infinity}
| {:http, keyword()}
| {:req_headers, [{String.t(), String.t()}]}
| {:resp_headers, [{String.t(), String.t()}]}
| {:inline_content_types, boolean() | list(String.t())}
| {:sniff_content_type, boolean()}
| {:redirect_on_failure, boolean()}
@spec call(Plug.Conn.t(), String.t(), list(option())) :: Plug.Conn.t()
def call(_conn, _url, _opts \\ [])
def call(conn = %{method: method}, url, opts) when method in @methods do
client_opts = Keyword.merge(@default_options, Keyword.get(opts, :http, []))
req_headers = build_req_headers(conn.req_headers, opts)
opts =
if filename = Pleroma.Web.MediaProxy.filename(url) do
Keyword.put_new(opts, :attachment_name, filename)
else
opts
end
with {:ok, nil} <- @cachex.get(:failed_proxy_url_cache, url),
{:ok, code, headers, client} <-
request_with_constraints(method, url, req_headers, client_opts, opts) do
response(conn, client, url, code, headers, opts)
else
{:ok, true} ->
conn
|> error_or_redirect(url, 500, "Request failed", opts)
|> halt()
{:ok, code, headers} ->
head_response(conn, url, code, headers, opts)
|> halt()
{:error, {:invalid_http_response, code}} ->
Logger.error("#{__MODULE__}: request to #{inspect(url)} failed with HTTP status #{code}")
track_failed_url(url, code, opts)
conn
|> error_or_redirect(
url,
code,
"Request failed: " <> Plug.Conn.Status.reason_phrase(code),
opts
)
|> halt()
{:error, error} ->
Logger.error("#{__MODULE__}: request to #{inspect(url)} failed: #{inspect(error)}")
track_failed_url(url, error, opts)
conn
|> error_or_redirect(url, 500, "Request failed", opts)
|> halt()
end
end
def call(conn, _, _) do
conn
|> send_resp(400, Plug.Conn.Status.reason_phrase(400))
|> halt()
end
defp request(method, url, headers, opts) do
method = method |> String.downcase() |> String.to_existing_atom()
url = maybe_encode_url(url)
Logger.debug("#{__MODULE__} #{method} #{url} #{inspect(headers)}")
case client().request(method, url, headers, "", opts) do
{:ok, code, headers, client} when code in @valid_resp_codes ->
{:ok, code, downcase_headers(headers), client}
{:ok, code, headers} when code in @valid_resp_codes ->
{:ok, code, downcase_headers(headers)}
{:ok, code, _, client} ->
client().close(client)
{:error, {:invalid_http_response, code}}
{:ok, code, _} ->
{:error, {:invalid_http_response, code}}
{:error, error} ->
{:error, error}
end
end
defp request_with_constraints(method, url, headers, client_opts, opts) do
- case request(method, url, headers, client_opts) do
- {:ok, _code, headers, client} = response ->
- case header_length_constraint(
- headers,
- Keyword.get(opts, :max_body_length, @max_body_length)
- ) do
- :ok ->
- response
-
- error ->
- client().close(client)
- error
- end
+ method
+ |> request(url, headers, client_opts)
+ |> constrain_response(Keyword.get(opts, :max_body_length, @max_body_length))
+ end
+
+ defp constrain_response({:ok, code, headers, client}, limit) do
+ case header_length_constraint(headers, limit) do
+ {:ok, headers} ->
+ {:ok, code, headers, client}
+
+ error ->
+ client().close(client)
+ error
+ end
+ end
- response ->
- response
+ defp constrain_response({:ok, code, headers}, limit) do
+ case header_length_constraint(headers, limit) do
+ {:ok, headers} -> {:ok, code, headers}
+ error -> error
end
end
+ defp constrain_response(response, _limit), do: response
+
defp response(conn, client, url, status, headers, opts) do
Logger.debug("#{__MODULE__} #{status} #{url} #{inspect(headers)}")
- {headers, client, prefetched} = maybe_prefetch_for_content_type(headers, client, opts)
+ case maybe_prefetch_for_content_type(headers, client, opts) do
+ {_headers, client, {:error, {:upstream, error}}} ->
+ Logger.warning(
+ "#{__MODULE__} request to #{url} failed before streaming: #{inspect(error)}"
+ )
+ client().close(client)
+ track_failed_url(url, error, opts)
+
+ conn
+ |> error_or_redirect(url, 500, "Request failed", opts)
+ |> halt()
+
+ {headers, client, prefetched} ->
+ stream_response(conn, client, url, status, headers, opts, prefetched)
+ end
+ end
+
+ defp stream_response(conn, client, url, status, headers, opts, prefetched) do
result =
conn
|> put_resp_headers(build_resp_headers(headers, opts))
- |> streaming_compat
|> send_chunked(status)
|> chunk_reply(client, opts, prefetched)
case result do
{:ok, conn} ->
halt(conn)
- {:error, :closed, conn} ->
+ {:error, {:downstream, _error}, conn} ->
client().close(client)
halt(conn)
- {:error, error, conn} ->
+ {:error, {:upstream, error}, conn} ->
Logger.warning(
"#{__MODULE__} request to #{url} failed while reading/chunking: #{inspect(error)}"
)
client().close(client)
- halt(conn)
+ track_failed_url(url, error, opts)
+ raise_stream_error(conn, url, error)
end
end
defp chunk_reply(conn, client, opts, :none) do
chunk_reply(conn, client, opts, 0, 0)
end
defp chunk_reply(conn, _client, _opts, :done), do: {:ok, conn}
defp chunk_reply(conn, _client, _opts, {:error, error}), do: {:error, error, conn}
defp chunk_reply(conn, client, opts, {:ok, data, duration}) do
with :ok <-
body_size_constraint(
byte_size(data),
Keyword.get(opts, :max_body_length, @max_body_length)
),
- {:ok, conn} <- chunk(conn, data) do
+ {:ok, conn} <- write_chunk(conn, data) do
chunk_reply(conn, client, opts, byte_size(data), duration)
else
- {:error, error} -> {:error, error, conn}
+ {:error, {:downstream, error}} -> {:error, {:downstream, error}, conn}
+ {:error, error} -> {:error, {:upstream, error}, conn}
end
end
defp chunk_reply(conn, client, opts, sent_so_far, duration) do
with {:ok, data, client, duration} <- read_chunk(client, duration, opts),
sent_so_far = sent_so_far + byte_size(data),
:ok <-
body_size_constraint(
sent_so_far,
Keyword.get(opts, :max_body_length, @max_body_length)
),
- {:ok, conn} <- chunk(conn, data) do
+ {:ok, conn} <- write_chunk(conn, data) do
chunk_reply(conn, client, opts, sent_so_far, duration)
else
:done -> {:ok, conn}
- {:error, error} -> {:error, error, conn}
+ {:error, {source, error}} -> {:error, {source, error}, conn}
+ {:error, error} -> {:error, {:upstream, error}, conn}
end
end
defp read_chunk(client, duration, opts) do
- with {:ok, timer} <-
- check_read_duration(
- duration,
- Keyword.get(opts, :max_read_duration, @max_read_duration)
- ),
- {:ok, data, client} <- client().stream_body(client),
- {:ok, duration} <- increase_read_duration(timer) do
- {:ok, data, client, duration}
+ result =
+ with {:ok, timer} <-
+ check_read_duration(
+ duration,
+ Keyword.get(opts, :max_read_duration, @max_read_duration)
+ ),
+ {:ok, data, client} <- client().stream_body(client),
+ {:ok, duration} <- increase_read_duration(timer) do
+ {:ok, data, client, duration}
+ end
+
+ case result do
+ {:error, error} -> {:error, {:upstream, error}}
+ result -> result
+ end
+ end
+
+ defp write_chunk(conn, data) do
+ case chunk(conn, data) do
+ {:error, error} -> {:error, {:downstream, error}}
+ result -> result
end
end
defp maybe_prefetch_for_content_type(headers, client, opts) do
if Keyword.get(opts, :sniff_content_type, false) and generic_content_type?(headers) do
case read_chunk(client, 0, opts) do
{:ok, data, client, duration} ->
{maybe_put_image_content_type(headers, data), client, {:ok, data, duration}}
:done ->
{headers, client, :done}
{:error, error} ->
{headers, client, {:error, error}}
end
else
{headers, client, :none}
end
end
defp generic_content_type?(headers) do
headers
|> get_content_type()
|> String.trim()
|> String.downcase()
|> then(&(&1 in ["", "application/octet-stream"]))
end
defp maybe_put_image_content_type(headers, data) do
with false <- data == "",
{:ok, %{mime_type: "image/" <> _ = content_type}} <-
Majic.perform({:bytes, binary_part(data, 0, min(byte_size(data), @sniff_bytes))},
pool: Pleroma.MajicPool
) do
[
{"content-type", content_type}
| Enum.reject(headers, fn {key, _} -> key == "content-type" end)
]
else
_ -> headers
end
rescue
error ->
Logger.debug("#{__MODULE__}: content-type sniffing failed: #{Exception.message(error)}")
headers
catch
kind, reason ->
Logger.debug("#{__MODULE__}: content-type sniffing failed: #{kind}: #{inspect(reason)}")
headers
end
defp head_response(conn, url, code, headers, opts) do
Logger.debug("#{__MODULE__} #{code} #{url} #{inspect(headers)}")
conn
|> put_resp_headers(build_resp_headers(headers, opts))
|> send_resp(code, "")
end
defp error_or_redirect(conn, url, code, body, opts) do
if Keyword.get(opts, :redirect_on_failure, false) do
conn
|> Phoenix.Controller.redirect(external: url)
|> halt()
else
conn
|> send_resp(code, body)
|> halt
end
end
defp downcase_headers(headers) do
Enum.map(headers, fn {k, v} ->
{String.downcase(k), v}
end)
end
defp get_content_type(headers) do
{_, content_type} =
List.keyfind(headers, "content-type", 0, {"content-type", "application/octet-stream"})
[content_type | _] = String.split(content_type, ";")
content_type
end
defp put_resp_headers(conn, headers) do
Enum.reduce(headers, conn, fn {k, v}, conn ->
put_resp_header(conn, k, v)
end)
end
defp build_req_headers(headers, opts) do
headers
|> downcase_headers()
|> Enum.filter(fn {k, _} -> k in @keep_req_headers end)
|> build_req_range_or_encoding_header(opts)
|> build_req_user_agent_header(opts)
|> merge_headers(Keyword.get(opts, :req_headers, []))
|> maybe_force_identity_encoding(opts)
end
defp merge_headers(headers, extra_headers) do
Enum.reduce(extra_headers, headers, fn {key, value}, headers ->
List.keystore(headers, String.downcase(key), 0, {String.downcase(key), value})
end)
end
defp maybe_force_identity_encoding(headers, opts) do
if Keyword.get(opts, :sniff_content_type, false) do
List.keystore(headers, "accept-encoding", 0, {"accept-encoding", "identity"})
else
headers
end
end
# Disable content-encoding if any @range_headers are requested (see #1823).
defp build_req_range_or_encoding_header(headers, _opts) do
range? = Enum.any?(headers, fn {header, _} -> Enum.member?(@range_headers, header) end)
cond do
range? && List.keymember?(headers, "accept-encoding", 0) ->
List.keydelete(headers, "accept-encoding", 0)
true ->
headers
end
end
defp build_req_user_agent_header(headers, _opts) do
List.keystore(
headers,
"user-agent",
0,
{"user-agent", Pleroma.Application.user_agent()}
)
end
defp build_resp_headers(headers, opts) do
headers
|> Enum.filter(fn {k, _} -> k in @keep_resp_headers end)
|> build_resp_cache_headers(opts)
|> sanitise_content_type()
|> build_resp_content_disposition_header(opts)
- |> Keyword.merge(Keyword.get(opts, :resp_headers, []))
+ |> merge_headers(Keyword.get(opts, :resp_headers, []))
+ |> ensure_no_transform()
+ end
+
+ defp ensure_no_transform(headers) do
+ {_, cache_control} =
+ List.keyfind(headers, "cache-control", 0, {"cache-control", @default_cache_control_header})
+
+ directives =
+ cache_control
+ |> String.split(",")
+ |> Enum.map(&(&1 |> String.trim() |> String.downcase()))
+
+ if "no-transform" in directives do
+ headers
+ else
+ replace_header(headers, "cache-control", cache_control <> ", no-transform")
+ end
end
defp sanitise_content_type(headers) do
original_ct = get_content_type(headers)
safe_ct =
Pleroma.Web.Plugs.Utils.get_safe_mime_type(
%{allowed_mime_types: @allowed_mime_types},
original_ct
)
[
{"content-type", safe_ct}
| Enum.filter(headers, fn {k, _v} -> k != "content-type" end)
]
end
defp build_resp_cache_headers(headers, _opts) do
has_cache? = Enum.any?(headers, fn {k, _} -> k in @resp_cache_headers end)
cond do
has_cache? ->
# There's caching header present but no cache-control -- we need to set our own
# as Plug defaults to "max-age=0, private, must-revalidate"
List.keystore(
headers,
"cache-control",
0,
{"cache-control", @default_cache_control_header}
)
true ->
List.keystore(
headers,
"cache-control",
0,
{"cache-control", @default_cache_control_header}
)
end
end
defp build_resp_content_disposition_header(headers, opts) do
opt = Keyword.get(opts, :inline_content_types, @inline_content_types)
content_type = get_content_type(headers)
attachment? =
cond do
is_list(opt) && !Enum.member?(opt, content_type) -> true
opt == false -> true
true -> false
end
if attachment? do
name =
try do
{{"content-disposition", content_disposition_string}, _} =
List.keytake(headers, "content-disposition", 0)
[name | _] =
Regex.run(
~r/filename="((?:[^"\\]|\\.)*)"/u,
content_disposition_string || "",
capture: :all_but_first
)
name
rescue
MatchError -> Keyword.get(opts, :attachment_name, "attachment")
end
disposition = "attachment; filename=\"#{name}\""
replace_header(headers, "content-disposition", disposition)
else
if opt == true do
headers
else
name = inline_filename(content_type)
disposition =
if name do
"inline; filename=\"#{name}\""
else
"inline"
end
replace_header(headers, "content-disposition", disposition)
end
end
end
defp replace_header(headers, key, value) do
[{key, value} | Enum.reject(headers, fn {header, _} -> header == key end)]
end
defp inline_filename(content_type) do
case MIME.extensions(content_type) do
[ext | _] when ext != "" -> "inline.#{ext}"
_ -> nil
end
end
- defp header_length_constraint(headers, limit) when is_integer(limit) and limit > 0 do
- with {_, size} <- List.keyfind(headers, "content-length", 0),
- {size, _} <- Integer.parse(size),
- true <- size <= limit do
- :ok
- else
- false ->
- {:error, :body_too_large}
+ defp header_length_constraint(headers, limit) do
+ lengths =
+ for {"content-length", value} <- headers,
+ do: parse_content_length(value)
+
+ case Enum.uniq(lengths) do
+ [] ->
+ {:ok, headers}
+
+ [size] when is_integer(size) ->
+ case body_size_constraint(size, limit) do
+ :ok -> {:ok, replace_header(headers, "content-length", to_string(size))}
+ error -> error
+ end
_ ->
- :ok
+ {:error, :invalid_content_length}
+ end
+ end
+
+ defp parse_content_length(value) when is_binary(value) do
+ case Integer.parse(value) do
+ {size, ""} when size >= 0 -> size
+ _ -> :invalid
end
end
- defp header_length_constraint(_, _), do: :ok
+ defp parse_content_length(_value), do: :invalid
defp body_size_constraint(size, limit) when is_integer(limit) and limit > 0 and size > limit do
{:error, :body_too_large}
end
defp body_size_constraint(_, _), do: :ok
defp check_read_duration(duration, max)
when is_integer(duration) and is_integer(max) and max > 0 do
if duration > max do
{:error, :read_duration_exceeded}
else
{:ok, {duration, :erlang.system_time(:millisecond)}}
end
end
defp check_read_duration(_, _), do: {:ok, :no_duration_limit}
defp increase_read_duration(:no_duration_limit), do: {:ok, :no_duration_limit}
defp increase_read_duration({previous_duration, started})
when is_integer(previous_duration) and is_integer(started) do
duration = :erlang.system_time(:millisecond) - started
{:ok, previous_duration + duration}
end
defp client, do: Pleroma.ReverseProxy.Client.Wrapper
+ # Neither Plug adapter exposes a public way to abort a committed response.
+ # Cowboy finalizes chunked HTTP/1 responses when the request process exits,
+ # so terminate the connection first to leave the response visibly incomplete.
+ defp raise_stream_error(
+ %Plug.Conn{adapter: {Plug.Cowboy.Conn, %{pid: connection_pid}}} = conn,
+ url,
+ error
+ ) do
+ if Plug.Conn.get_http_protocol(conn) in [:"HTTP/1.0", :"HTTP/1.1"] do
+ Process.exit(connection_pid, :kill)
+ end
+
+ raise StreamError, url: url, reason: error
+ end
+
+ # Bandit turns this exception into RST_STREAM; a regular exception sends an
+ # empty DATA frame with END_STREAM and makes the partial body look complete.
+ defp raise_stream_error(
+ %Plug.Conn{
+ adapter: {Bandit.Adapter, %{transport: %{stream_id: stream_id}}}
+ },
+ url,
+ error
+ ) do
+ raise Bandit.HTTP2.Errors.StreamError,
+ message: StreamError.message(%StreamError{url: url, reason: error}),
+ error_code: Bandit.HTTP2.Errors.internal_error(),
+ stream_id: stream_id
+ end
+
+ defp raise_stream_error(_conn, url, error) do
+ raise StreamError, url: url, reason: error
+ end
+
defp track_failed_url(url, error, opts) do
ttl =
unless error in [:body_too_large, 400, 204] do
Keyword.get(opts, :failed_request_ttl, @failed_request_ttl)
else
nil
end
@cachex.put(:failed_proxy_url_cache, url, true, ttl: ttl)
end
- # When Cowboy handles a chunked response with a content-length header it streams
- # over HTTP 1.1 instead of chunking. Bandit cannot stream over HTTP 1.1 so the header
- # must be stripped or it breaks RFC compliance for Transfer Encoding: Chunked. RFC9112§6.2
- #
- # HTTP2 is always streamed for all adapters.
- defp streaming_compat(conn) do
- with Phoenix.Endpoint.Cowboy2Adapter <- Pleroma.Web.Endpoint.config(:adapter) do
- conn
- else
- _ -> delete_resp_header(conn, "content-length")
- end
- end
-
# Only when Tesla adapter is Hackney or Finch does the URL
# need encoding before Reverse Proxying as both end up
# using the raw Hackney client and cannot leverage our
# EncodeUrl Tesla middleware
# Also do it for test environment
defp maybe_encode_url(url) do
case Application.get_env(:tesla, :adapter) do
Tesla.Adapter.Hackney -> URIEncoding.encode_url(url)
{Tesla.Adapter.Finch, _} -> URIEncoding.encode_url(url)
Tesla.Mock -> URIEncoding.encode_url(url)
_ -> url
end
end
end
diff --git a/test/pleroma/reverse_proxy_test.exs b/test/pleroma/reverse_proxy_test.exs
index 63ab4202c..bf0a809c4 100644
--- a/test/pleroma/reverse_proxy_test.exs
+++ b/test/pleroma/reverse_proxy_test.exs
@@ -1,594 +1,1018 @@
# Pleroma: A lightweight social networking server
# Copyright © 2017-2022 Pleroma Authors <https://pleroma.social/>
# SPDX-License-Identifier: AGPL-3.0-only
defmodule Pleroma.ReverseProxyTest do
use Pleroma.Web.ConnCase
import ExUnit.CaptureLog
import Mox
alias Pleroma.ReverseProxy
alias Pleroma.ReverseProxy.ClientMock
alias Plug.Conn
+ defmodule ProxyPlug do
+ @moduledoc false
+
+ def init(options), do: options
+
+ def call(conn, {url, opts}), do: Pleroma.ReverseProxy.call(conn, url, opts)
+ def call(conn, url), do: Pleroma.ReverseProxy.call(conn, url)
+ end
+
setup_all do
{:ok, _} = Registry.start_link(keys: :unique, name: ClientMock)
:ok
end
setup :verify_on_exit!
+ defp start_bandit(url, opts \\ []) do
+ start_supervised!(
+ {Bandit,
+ plug: {ProxyPlug, {url, opts}},
+ ip: {127, 0, 0, 1},
+ port: 0,
+ http_options: [log_exceptions_with_status_codes: []]}
+ )
+ end
+
+ defp start_cowboy(url, opts \\ []) do
+ ref = {__MODULE__, make_ref()}
+
+ {:ok, _pid} =
+ Plug.Cowboy.http(ProxyPlug, {url, opts},
+ ip: {127, 0, 0, 1},
+ port: 0,
+ ref: ref,
+ protocol_options: [stream_handlers: [:cowboy_stream_h]]
+ )
+
+ on_exit(fn -> Plug.Cowboy.shutdown(ref) end)
+ :ranch.get_port(ref)
+ end
+
+ defp bandit_request(url, headers \\ [], opts \\ []) do
+ pid = start_bandit(url, opts)
+ {:ok, {{127, 0, 0, 1}, port}} = ThousandIsland.listener_info(pid)
+ http1_request(port, headers)
+ end
+
+ defp cowboy_request(url, headers \\ [], opts \\ []) do
+ url
+ |> start_cowboy(opts)
+ |> http1_request(headers)
+ end
+
+ defp http1_request(port, headers) do
+ {:ok, socket} = :gen_tcp.connect({127, 0, 0, 1}, port, [:binary, active: false])
+
+ request_headers =
+ [["host: localhost", "connection: close"] | headers]
+ |> List.flatten()
+ |> Enum.join("\r\n")
+
+ :ok = :gen_tcp.send(socket, "GET / HTTP/1.1\r\n#{request_headers}\r\n\r\n")
+ recv_all(socket, [])
+ end
+
+ defp recv_all(socket, chunks) do
+ case :gen_tcp.recv(socket, 0, 5_000) do
+ {:ok, data} -> recv_all(socket, [data | chunks])
+ {:error, reason} -> {reason, chunks |> Enum.reverse() |> IO.iodata_to_binary()}
+ end
+ end
+
+ defp receive_mint_responses(conn, responses \\ []) do
+ receive do
+ message ->
+ case Mint.HTTP2.stream(conn, message) do
+ {:ok, conn, new_responses} ->
+ responses = responses ++ new_responses
+
+ if Enum.any?(new_responses, fn response ->
+ elem(response, 0) in [:done, :error]
+ end) do
+ responses
+ else
+ receive_mint_responses(conn, responses)
+ end
+
+ {:error, _conn, error, new_responses} ->
+ responses ++ new_responses ++ [{:connection_error, error}]
+
+ :unknown ->
+ receive_mint_responses(conn, responses)
+ end
+ after
+ 2_000 -> responses ++ [:timeout]
+ end
+ end
+
defp request_mock(invokes) do
ClientMock
|> expect(:request, fn :get, url, headers, _body, _opts ->
Registry.register(ClientMock, url, 0)
body = headers |> Enum.into(%{}) |> Jason.encode!()
{:ok, 200,
[
{"content-type", "application/json"},
{"content-length", byte_size(body) |> to_string()}
], %{url: url, body: body}}
end)
|> expect(:stream_body, invokes, fn %{url: url, body: body} = client ->
case Registry.lookup(ClientMock, url) do
[{_, 0}] ->
Registry.update_value(ClientMock, url, &(&1 + 1))
{:ok, body, client}
[{_, 1}] ->
Registry.unregister(ClientMock, url)
:done
end
end)
end
+ defp stream_then_error_mock(url, headers, chunk \\ "partial", error \\ :closed) do
+ ClientMock
+ |> expect(:request, fn :get, ^url, _, _, _ ->
+ Registry.register(ClientMock, url, 0)
+ {:ok, 200, headers, %{url: url}}
+ end)
+ |> expect(:stream_body, 2, fn %{url: ^url} = client ->
+ case Registry.lookup(ClientMock, url) do
+ [{_, 0}] ->
+ Registry.update_value(ClientMock, url, &(&1 + 1))
+ {:ok, chunk, client}
+
+ [{_, 1}] ->
+ Registry.unregister(ClientMock, url)
+ {:error, error}
+ end
+ end)
+ |> expect(:close, fn _ -> :ok end)
+ end
+
+ defp complete_stream_mock(url, headers, body) do
+ ClientMock
+ |> expect(:request, fn :get, ^url, _, _, _ ->
+ Registry.register(ClientMock, url, 0)
+ {:ok, 200, headers, %{url: url}}
+ end)
+ |> expect(:stream_body, 2, fn %{url: ^url} = client ->
+ case Registry.lookup(ClientMock, url) do
+ [{_, 0}] ->
+ Registry.update_value(ClientMock, url, &(&1 + 1))
+ {:ok, body, client}
+
+ [{_, 1}] ->
+ Registry.unregister(ClientMock, url)
+ :done
+ end
+ end)
+ end
+
describe "reverse proxy" do
test "do not track successful request", %{conn: conn} do
request_mock(2)
url = "/success"
conn = ReverseProxy.call(conn, url)
assert conn.status == 200
assert Cachex.get(:failed_proxy_url_cache, url) == {:ok, nil}
end
end
test "use Pleroma's user agent in the request; don't pass the client's", %{conn: conn} do
request_mock(2)
conn =
conn
|> Plug.Conn.put_req_header("user-agent", "fake/1.0")
|> ReverseProxy.call("/user-agent")
# Convert the response to a map without relying on json_response
body = conn.resp_body
assert conn.status == 200
response = Jason.decode!(body)
assert response == %{"user-agent" => Pleroma.Application.user_agent()}
end
- test "closed connection", %{conn: conn} do
- ClientMock
- |> expect(:request, fn :get, "/closed", _, _, _ -> {:ok, 200, [], %{}} end)
- |> expect(:stream_body, fn _ -> {:error, :closed} end)
- |> expect(:close, fn _ -> :ok end)
+ test "aborts a partially streamed response when the upstream closes", %{conn: conn} do
+ stream_then_error_mock("/closed", [{"content-length", "20"}])
- conn = ReverseProxy.call(conn, "/closed")
- assert conn.halted
+ assert_raise ReverseProxy.StreamError, fn ->
+ ReverseProxy.call(conn, "/closed")
+ end
+
+ assert Cachex.get(:failed_proxy_url_cache, "/closed") == {:ok, true}
+ end
+
+ describe "Bandit HTTP/1 streaming" do
+ test "closes an incomplete chunked response without a terminating chunk" do
+ Mox.set_mox_global(ClientMock)
+ stream_then_error_mock("/bandit-closed", [{"content-type", "image/png"}])
+
+ assert {:closed, response} = bandit_request("/bandit-closed")
+ assert response =~ "transfer-encoding: chunked"
+ assert response =~ "partial"
+ refute response =~ "0\r\n\r\n"
+ end
+
+ test "closes an incomplete fixed-length response before the declared length" do
+ Mox.set_mox_global(ClientMock)
+
+ stream_then_error_mock("/bandit-fixed-closed", [
+ {"content-length", "20"},
+ {"content-type", "image/png"}
+ ])
+
+ assert {:closed, response} = bandit_request("/bandit-fixed-closed")
+ assert response =~ "content-length: 20"
+ refute response =~ "transfer-encoding"
+ assert String.ends_with?(response, "partial")
+ end
+
+ test "streams a complete fixed-length response without transformation" do
+ Mox.set_mox_global(ClientMock)
+ body = "complete body"
+
+ complete_stream_mock(
+ "/bandit-complete",
+ [
+ {"content-length", to_string(byte_size(body))},
+ {"content-type", "image/png"}
+ ],
+ body
+ )
+
+ assert {:closed, response} =
+ bandit_request("/bandit-complete", ["accept-encoding: deflate"])
+
+ assert response =~ "content-length: #{byte_size(body)}"
+ assert response =~ "cache-control: public, max-age=1209600, immutable, no-transform"
+ refute response =~ "content-encoding: gzip"
+ refute response =~ "content-encoding: deflate"
+ assert String.ends_with?(response, body)
+ end
+
+ test "keeps no-transform when custom response headers replace cache-control" do
+ Mox.set_mox_global(ClientMock)
+ body = "complete body"
+
+ complete_stream_mock(
+ "/bandit-custom-cache-control",
+ [
+ {"content-length", to_string(byte_size(body))},
+ {"content-type", "image/png"}
+ ],
+ body
+ )
+
+ assert {:closed, response} =
+ bandit_request(
+ "/bandit-custom-cache-control",
+ ["accept-encoding: deflate"],
+ resp_headers: [{"cache-control", "private, max-age=60"}]
+ )
+
+ assert response =~ "content-length: #{byte_size(body)}"
+ assert response =~ "cache-control: private, max-age=60, no-transform"
+ refute response =~ "content-encoding: gzip"
+ refute response =~ "content-encoding: deflate"
+ assert String.ends_with?(response, body)
+ end
+
+ test "does not blacklist the upstream when the downstream disconnects" do
+ Mox.set_mox_global(ClientMock)
+ test_pid = self()
+
+ ClientMock
+ |> expect(:request, fn :get, "/bandit-client-closed", _, _, _ ->
+ {:ok, 200, [{"content-type", "image/png"}], %{url: "/bandit-client-closed"}}
+ end)
+ |> stub(:stream_body, fn client ->
+ send(test_pid, {:stream_requested, self()})
+
+ receive do
+ :continue -> {:ok, :binary.copy("x", 1_000_000), client}
+ after
+ 1_000 -> raise "timed out waiting for downstream disconnect"
+ end
+ end)
+ |> expect(:close, fn _ ->
+ send(test_pid, :upstream_closed)
+ :ok
+ end)
+
+ pid = start_bandit("/bandit-client-closed")
+ {:ok, {{127, 0, 0, 1}, port}} = ThousandIsland.listener_info(pid)
+ {:ok, socket} = :gen_tcp.connect({127, 0, 0, 1}, port, [:binary, active: false])
+ :ok = :gen_tcp.send(socket, "GET / HTTP/1.1\r\nhost: localhost\r\n\r\n")
+
+ assert_receive {:stream_requested, request_pid}
+ :ok = :inet.setopts(socket, linger: {true, 0})
+ :ok = :gen_tcp.close(socket)
+ send(request_pid, :continue)
+
+ assert_receive :upstream_closed
+ assert Cachex.get(:failed_proxy_url_cache, "/bandit-client-closed") == {:ok, nil}
+ end
+ end
+
+ describe "Cowboy streaming" do
+ test "closes an incomplete HTTP/1 chunked response without a terminating chunk" do
+ Mox.set_mox_global(ClientMock)
+ stream_then_error_mock("/cowboy-closed", [{"content-type", "image/png"}])
+
+ assert {:closed, response} = cowboy_request("/cowboy-closed")
+ assert response =~ "transfer-encoding: chunked"
+ assert response =~ "partial"
+ refute response =~ "0\r\n\r\n"
+ end
+
+ test "closes an incomplete HTTP/1 fixed-length response before the declared length" do
+ Mox.set_mox_global(ClientMock)
+
+ stream_then_error_mock("/cowboy-fixed-closed", [
+ {"content-length", "20"},
+ {"content-type", "image/png"}
+ ])
+
+ assert {:closed, response} = cowboy_request("/cowboy-fixed-closed")
+ assert response =~ "content-length: 20"
+ refute response =~ "transfer-encoding"
+ assert String.ends_with?(response, "partial")
+ end
+
+ test "resets an incomplete HTTP/2 stream" do
+ Mox.set_mox_global(ClientMock)
+ stream_then_error_mock("/cowboy-http2-closed", [{"content-type", "image/png"}])
+
+ port = start_cowboy("/cowboy-http2-closed")
+ {:ok, mint} = Mint.HTTP2.connect(:http, "localhost", port)
+ {:ok, mint, request_ref} = Mint.HTTP2.request(mint, "GET", "/", [], nil)
+ responses = receive_mint_responses(mint)
+
+ assert {:status, request_ref, 200} in responses
+ assert {:data, request_ref, "partial"} in responses
+
+ assert Enum.any?(responses, fn
+ {:error, ^request_ref, %Mint.HTTPError{reason: {:server_closed_request, _}}} ->
+ true
+
+ _ ->
+ false
+ end)
+
+ refute {:done, request_ref} in responses
+ refute :timeout in responses
+ end
+ end
+
+ test "resets an incomplete Bandit HTTP/2 stream" do
+ Mox.set_mox_global(ClientMock)
+ stream_then_error_mock("/bandit-http2-closed", [{"content-type", "image/png"}])
+
+ pid = start_bandit("/bandit-http2-closed")
+ {:ok, {{127, 0, 0, 1}, port}} = ThousandIsland.listener_info(pid)
+ {:ok, mint} = Mint.HTTP2.connect(:http, "localhost", port)
+ {:ok, mint, request_ref} = Mint.HTTP2.request(mint, "GET", "/", [], nil)
+ responses = receive_mint_responses(mint)
+
+ assert {:status, request_ref, 200} in responses
+ assert {:data, request_ref, "partial"} in responses
+
+ assert Enum.any?(responses, fn
+ {:error, ^request_ref,
+ %Mint.HTTPError{reason: {:server_closed_request, :internal_error}}} ->
+ true
+
+ _ ->
+ false
+ end)
+
+ refute {:done, request_ref} in responses
+ refute :timeout in responses
end
defp stream_mock(invokes, with_close? \\ false) do
ClientMock
|> expect(:request, fn :get, "/stream-bytes/" <> length, _, _, _ ->
Registry.register(ClientMock, "/stream-bytes/" <> length, 0)
{:ok, 200, [{"content-type", "application/octet-stream"}],
%{url: "/stream-bytes/" <> length}}
end)
|> expect(:stream_body, invokes, fn %{url: "/stream-bytes/" <> length} = client ->
max = String.to_integer(length)
case Registry.lookup(ClientMock, "/stream-bytes/" <> length) do
[{_, current}] when current < max ->
Registry.update_value(
ClientMock,
"/stream-bytes/" <> length,
&(&1 + 10)
)
{:ok, "0123456789", client}
[{_, ^max}] ->
Registry.unregister(ClientMock, "/stream-bytes/" <> length)
:done
end
end)
if with_close? do
expect(ClientMock, :close, fn _ -> :ok end)
end
end
describe "max_body" do
test "length returns error if content-length more than option", %{conn: conn} do
request_mock(0)
expect(ClientMock, :close, fn _ -> :ok end)
assert capture_log(fn ->
ReverseProxy.call(conn, "/huge-file", max_body_length: 4)
end) =~
"[error] Elixir.Pleroma.ReverseProxy: request to \"/huge-file\" failed: :body_too_large"
assert {:ok, true} == Cachex.get(:failed_proxy_url_cache, "/huge-file")
assert capture_log(fn ->
ReverseProxy.call(conn, "/huge-file", max_body_length: 4)
end) == ""
end
test "closes streamed responses with invalid status", %{conn: conn} do
ClientMock
|> expect(:request, fn :get, "/invalid-status", _, _, _ ->
{:ok, 404, [], %{url: "/invalid-status"}}
end)
|> expect(:close, fn %{url: "/invalid-status"} -> :ok end)
ReverseProxy.call(conn, "/invalid-status")
end
test "max_body_length returns error if streaming body more than that option", %{conn: conn} do
stream_mock(3, true)
- assert capture_log(fn ->
- ReverseProxy.call(conn, "/stream-bytes/50", max_body_length: 29)
- end) =~
+ log =
+ capture_log(fn ->
+ assert_raise ReverseProxy.StreamError, fn ->
+ ReverseProxy.call(conn, "/stream-bytes/50", max_body_length: 29)
+ end
+ end)
+
+ assert log =~
"Elixir.Pleroma.ReverseProxy request to /stream-bytes/50 failed while reading/chunking: :body_too_large"
end
test "max_body_length accepts a body exactly at the limit", %{conn: conn} do
stream_mock(4)
conn = ReverseProxy.call(conn, "/stream-bytes/30", max_body_length: 30)
assert byte_size(conn.resp_body) == 30
end
end
describe "HEAD requests" do
test "common", %{conn: conn} do
ClientMock
|> expect(:request, fn :head, "/head", _, _, _ ->
{:ok, 200, [{"content-type", "image/png"}]}
end)
conn = ReverseProxy.call(Map.put(conn, :method, "HEAD"), "/head")
assert conn.status == 200
assert Conn.get_resp_header(conn, "content-type") == ["image/png"]
assert conn.resp_body == ""
end
+
+ test "rejects a malformed content-length", %{conn: conn} do
+ url = "/head-invalid-content-length"
+
+ expect(ClientMock, :request, fn :head, ^url, _, _, _ ->
+ {:ok, 200, [{"content-length", "20junk"}]}
+ end)
+
+ conn = ReverseProxy.call(Map.put(conn, :method, "HEAD"), url)
+
+ assert conn.status == 500
+ assert Cachex.get(:failed_proxy_url_cache, url) == {:ok, true}
+ end
end
defp error_mock(status) when is_integer(status) do
ClientMock
|> expect(:request, fn :get, "/status/" <> _, _, _, _ ->
{:error, status}
end)
end
describe "returns error on" do
test "500", %{conn: conn} do
error_mock(500)
url = "/status/500"
capture_log(fn -> ReverseProxy.call(conn, url) end) =~
"[error] Elixir.Pleroma.ReverseProxy: request to /status/500 failed with HTTP status 500"
assert Cachex.get(:failed_proxy_url_cache, url) == {:ok, true}
{:ok, ttl} = Cachex.ttl(:failed_proxy_url_cache, url)
assert ttl <= 60_000
end
test "400", %{conn: conn} do
error_mock(400)
url = "/status/400"
capture_log(fn -> ReverseProxy.call(conn, url) end) =~
"[error] Elixir.Pleroma.ReverseProxy: request to /status/400 failed with HTTP status 400"
assert Cachex.get(:failed_proxy_url_cache, url) == {:ok, true}
assert Cachex.ttl(:failed_proxy_url_cache, url) == {:ok, nil}
end
test "403", %{conn: conn} do
error_mock(403)
url = "/status/403"
capture_log(fn ->
ReverseProxy.call(conn, url, failed_request_ttl: :timer.seconds(120))
end) =~
"[error] Elixir.Pleroma.ReverseProxy: request to /status/403 failed with HTTP status 403"
{:ok, ttl} = Cachex.ttl(:failed_proxy_url_cache, url)
assert ttl > 100_000
end
test "204", %{conn: conn} do
url = "/status/204"
ClientMock
|> expect(:request, fn :get, _url, _, _, _ -> {:ok, 204, [], %{}} end)
|> expect(:close, fn %{} -> :ok end)
capture_log(fn ->
conn = ReverseProxy.call(conn, url)
assert conn.resp_body == "Request failed: No Content"
assert conn.halted
end) =~
"[error] Elixir.Pleroma.ReverseProxy: request to \"/status/204\" failed with HTTP status 204"
assert Cachex.get(:failed_proxy_url_cache, url) == {:ok, true}
assert Cachex.ttl(:failed_proxy_url_cache, url) == {:ok, nil}
end
end
test "streaming", %{conn: conn} do
stream_mock(21)
conn = ReverseProxy.call(conn, "/stream-bytes/200")
assert conn.state == :chunked
assert byte_size(conn.resp_body) == 200
assert Conn.get_resp_header(conn, "content-type") == ["application/octet-stream"]
end
test "sniffs and preserves an image body with a generic content type", %{conn: conn} do
body = File.read!("test/fixtures/image.jpg")
ClientMock
|> expect(:request, fn :get, "/extensionless", headers, _, _ ->
assert {"accept-encoding", "identity"} in headers
{:ok, 200, [{"content-type", "application/octet-stream"}], %{body: body}}
end)
|> expect(:stream_body, fn %{body: ^body} = client ->
{:ok, body, Map.delete(client, :body)}
end)
|> expect(:stream_body, fn %{} -> :done end)
conn =
ReverseProxy.call(conn, "/extensionless",
sniff_content_type: true,
max_read_duration: :infinity,
req_headers: [{"accept-encoding", "gzip"}]
)
assert conn.resp_body == body
assert Conn.get_resp_header(conn, "content-type") == ["image/jpeg"]
assert Conn.get_resp_header(conn, "content-disposition") == [
"inline; filename=\"inline.jpg\""
]
end
+ test "redirects a content sniffing failure before committing the response", %{conn: conn} do
+ url = "https://example.com/extensionless"
+
+ ClientMock
+ |> expect(:request, fn :get, ^url, _, _, _ ->
+ {:ok, 200, [{"content-type", "application/octet-stream"}], %{url: url}}
+ end)
+ |> expect(:stream_body, fn _ -> {:error, :closed} end)
+ |> expect(:close, fn _ -> :ok end)
+
+ conn =
+ ReverseProxy.call(conn, url,
+ sniff_content_type: true,
+ redirect_on_failure: true
+ )
+
+ assert conn.status == 302
+ assert Conn.get_resp_header(conn, "location") == [url]
+ assert conn.state == :sent
+ assert Cachex.get(:failed_proxy_url_cache, url) == {:ok, true}
+ end
+
test "preserves a generic content type for a non-image body", %{conn: conn} do
body = "not an image"
ClientMock
|> expect(:request, fn :get, "/extensionless", _, _, _ ->
{:ok, 200, [], %{body: body}}
end)
|> expect(:stream_body, fn %{body: ^body} = client ->
{:ok, body, Map.delete(client, :body)}
end)
|> expect(:stream_body, fn %{} -> :done end)
conn = ReverseProxy.call(conn, "/extensionless", sniff_content_type: true)
assert conn.resp_body == body
assert Conn.get_resp_header(conn, "content-type") == ["application/octet-stream"]
assert Conn.get_resp_header(conn, "content-disposition") == [
"attachment; filename=\"extensionless\""
]
end
defp headers_mock(_) do
ClientMock
|> expect(:request, fn :get, "/headers", headers, _, _ ->
Registry.register(ClientMock, "/headers", 0)
{:ok, 200, [{"content-type", "application/json"}], %{url: "/headers", headers: headers}}
end)
|> expect(:stream_body, 2, fn %{url: url, headers: headers} = client ->
case Registry.lookup(ClientMock, url) do
[{_, 0}] ->
Registry.update_value(ClientMock, url, &(&1 + 1))
headers = for {k, v} <- headers, into: %{}, do: {String.capitalize(k), v}
{:ok, Jason.encode!(%{headers: headers}), client}
[{_, 1}] ->
Registry.unregister(ClientMock, url)
:done
end
end)
:ok
end
describe "keep request headers" do
setup [:headers_mock]
test "header passes", %{conn: conn} do
conn =
Conn.put_req_header(
conn,
"accept",
"text/html"
)
|> ReverseProxy.call("/headers")
body = conn.resp_body
assert conn.status == 200
response = Jason.decode!(body)
headers = response["headers"]
assert headers["Accept"] == "text/html"
end
test "header is filtered", %{conn: conn} do
conn =
Conn.put_req_header(
conn,
"accept-language",
"en-US"
)
|> ReverseProxy.call("/headers")
body = conn.resp_body
assert conn.status == 200
response = Jason.decode!(body)
headers = response["headers"]
refute headers["Accept-Language"]
end
end
test "returns 400 on non GET, HEAD requests", %{conn: conn} do
conn = ReverseProxy.call(Map.put(conn, :method, "POST"), "/ip")
assert conn.status == 400
end
describe "cache resp headers" do
test "add cache-control", %{conn: conn} do
ClientMock
|> expect(:request, fn :get, "/cache", _, _, _ ->
{:ok, 200, [{"ETag", "some ETag"}], %{}}
end)
|> expect(:stream_body, fn _ -> :done end)
conn = ReverseProxy.call(conn, "/cache")
- assert {"cache-control", "public, max-age=1209600, immutable"} in conn.resp_headers
+
+ assert {"cache-control", "public, max-age=1209600, immutable, no-transform"} in conn.resp_headers
+ end
+
+ test "preserves a valid upstream content-length", %{conn: conn} do
+ body = "complete body"
+
+ ClientMock
+ |> expect(:request, fn :get, "/content-length", _, _, _ ->
+ {:ok, 200, [{"content-length", to_string(byte_size(body))}], %{body: body}}
+ end)
+ |> expect(:stream_body, fn %{body: ^body} = client ->
+ {:ok, body, Map.delete(client, :body)}
+ end)
+ |> expect(:stream_body, fn %{} -> :done end)
+
+ conn = ReverseProxy.call(conn, "/content-length")
+
+ assert Conn.get_resp_header(conn, "content-length") == [to_string(byte_size(body))]
+ end
+
+ test "rejects malformed upstream content-length values" do
+ Enum.each(["20junk", "-1"], fn content_length ->
+ url = "/invalid-content-length/#{content_length}"
+
+ ClientMock
+ |> expect(:request, fn :get, ^url, _, _, _ ->
+ {:ok, 200, [{"content-length", content_length}], %{url: url}}
+ end)
+ |> expect(:close, fn %{url: ^url} -> :ok end)
+
+ conn = ReverseProxy.call(build_conn(), url)
+
+ assert conn.status == 500
+ assert Cachex.get(:failed_proxy_url_cache, url) == {:ok, true}
+ end)
+ end
+
+ test "rejects conflicting upstream content-length values" do
+ url = "/conflicting-content-length"
+
+ ClientMock
+ |> expect(:request, fn :get, ^url, _, _, _ ->
+ {:ok, 200, [{"content-length", "7"}, {"content-length", "20"}], %{url: url}}
+ end)
+ |> stub(:stream_body, fn _ -> :done end)
+ |> stub(:close, fn _ -> :ok end)
+
+ conn = ReverseProxy.call(build_conn(), url)
+
+ assert conn.status == 500
+ assert Cachex.get(:failed_proxy_url_cache, url) == {:ok, true}
+ end
+
+ test "normalizes equivalent upstream content-length values", %{conn: conn} do
+ body = "partial"
+
+ complete_stream_mock(
+ "/duplicate-content-length",
+ [{"content-length", "7"}, {"content-length", "07"}],
+ body
+ )
+
+ conn = ReverseProxy.call(conn, "/duplicate-content-length")
+
+ assert Conn.get_resp_header(conn, "content-length") == ["7"]
+ assert conn.resp_body == body
end
end
defp disposition_headers_mock(headers) do
ClientMock
|> expect(:request, fn :get, "/disposition", _, _, _ ->
Registry.register(ClientMock, "/disposition", 0)
{:ok, 200, headers, %{url: "/disposition"}}
end)
|> expect(:stream_body, 2, fn %{url: "/disposition"} = client ->
case Registry.lookup(ClientMock, "/disposition") do
[{_, 0}] ->
Registry.update_value(ClientMock, "/disposition", &(&1 + 1))
{:ok, "", client}
[{_, 1}] ->
Registry.unregister(ClientMock, "/disposition")
:done
end
end)
end
describe "response content disposition header" do
test "not attachment", %{conn: conn} do
disposition_headers_mock([
{"content-type", "image/gif"},
{"content-length", "0"}
])
conn = ReverseProxy.call(conn, "/disposition")
assert {"content-type", "image/gif"} in conn.resp_headers
assert {"content-disposition", "inline; filename=\"inline.gif\""} in conn.resp_headers
end
test "forces inline for inline content types overriding upstream attachment", %{
conn: conn
} do
disposition_headers_mock([
{"content-type", "image/png"},
{"content-disposition", "attachment; filename=\"filename.png\""},
{"content-disposition", "attachment; filename=\"duplicate.png\""},
{"content-length", "0"}
])
conn = ReverseProxy.call(conn, "/disposition")
[disposition] = Conn.get_resp_header(conn, "content-disposition")
assert String.starts_with?(disposition, "inline")
refute String.starts_with?(disposition, "attachment")
end
test "with content-disposition header", %{conn: conn} do
disposition_headers_mock([
{"content-disposition", "attachment; filename=\"filename.jpg\""},
{"content-length", "0"}
])
conn = ReverseProxy.call(conn, "/disposition")
assert {"content-disposition", "attachment; filename=\"filename.jpg\""} in conn.resp_headers
end
test "with inline_content_types: true leaves upstream headers untouched", %{
conn: conn
} do
# opt == true: the proxy must not synthesise or rewrite content-disposition.
disposition_headers_mock([
{"content-type", "image/png"},
{"content-disposition", "attachment; filename=\"upstream.png\""},
{"content-length", "0"}
])
conn = ReverseProxy.call(conn, "/disposition", inline_content_types: true)
assert {"content-disposition", "attachment; filename=\"upstream.png\""} in conn.resp_headers
end
test "forces bare inline for a whitelisted type with no MIME extension", %{
conn: conn
} do
# image/x-foo-bar has no entry in MIME's database, so inline_filename/1
# returns nil and the disposition should be the bare token "inline".
disposition_headers_mock([
{"content-type", "image/x-foo-bar"},
{"content-length", "0"}
])
conn = ReverseProxy.call(conn, "/disposition", inline_content_types: ["image/x-foo-bar"])
[disposition] = Conn.get_resp_header(conn, "content-disposition")
assert disposition == "inline"
end
test "serves modern browser image types inline", %{conn: conn} do
disposition_headers_mock([
{"content-type", "image/webp"},
{"content-length", "0"}
])
conn = ReverseProxy.call(conn, "/disposition")
assert Conn.get_resp_header(conn, "content-disposition") == [
"inline; filename=\"inline.webp\""
]
end
end
describe "content-type sanitisation" do
test "preserves allowed image type", %{conn: conn} do
ClientMock
|> expect(:request, fn :get, "/content", _, _, _ ->
{:ok, 200, [{"content-type", "image/png"}], %{url: "/content"}}
end)
|> expect(:stream_body, fn _ -> :done end)
conn = ReverseProxy.call(conn, "/content")
assert conn.status == 200
assert Conn.get_resp_header(conn, "content-type") == ["image/png"]
end
test "preserves allowed video type", %{conn: conn} do
ClientMock
|> expect(:request, fn :get, "/content", _, _, _ ->
{:ok, 200, [{"content-type", "video/mp4"}], %{url: "/content"}}
end)
|> expect(:stream_body, fn _ -> :done end)
conn = ReverseProxy.call(conn, "/content")
assert conn.status == 200
assert Conn.get_resp_header(conn, "content-type") == ["video/mp4"]
end
test "sanitizes ActivityPub content type", %{conn: conn} do
ClientMock
|> expect(:request, fn :get, "/content", _, _, _ ->
{:ok, 200, [{"content-type", "application/activity+json"}], %{url: "/content"}}
end)
|> expect(:stream_body, fn _ -> :done end)
conn = ReverseProxy.call(conn, "/content")
assert conn.status == 200
assert Conn.get_resp_header(conn, "content-type") == ["application/octet-stream"]
end
test "sanitizes LD-JSON content type", %{conn: conn} do
ClientMock
|> expect(:request, fn :get, "/content", _, _, _ ->
{:ok, 200, [{"content-type", "application/ld+json"}], %{url: "/content"}}
end)
|> expect(:stream_body, fn _ -> :done end)
conn = ReverseProxy.call(conn, "/content")
assert conn.status == 200
assert Conn.get_resp_header(conn, "content-type") == ["application/octet-stream"]
end
end
# Hackney is used for Reverse Proxy when Hackney or Finch is the Tesla Adapter
# Gun is able to proxy through Tesla, so it does not need testing as the
# test cases in the Pleroma.HTTPTest module are sufficient
describe "Hackney URL encoding:" do
setup do
ClientMock
|> expect(:request, fn
:get,
"https://example.com/emoji/Pack%201/koronebless.png?foo=bar+baz",
_headers,
_body,
_opts ->
{:ok, 200, [{"content-type", "image/png"}], "It works!"}
:get,
"https://example.com/media/foo/bar%20!$&'()*+,;=/:%20@a%20%5Bbaz%5D.mp4",
_headers,
_body,
_opts ->
{:ok, 200, [{"content-type", "video/mp4"}], "Allowed reserved chars."}
:get, "https://example.com/media/unicode%20%F0%9F%99%82%20.gif", _headers, _body, _opts ->
{:ok, 200, [{"content-type", "image/gif"}], "Unicode emoji in path"}
end)
|> stub(:stream_body, fn _ -> :done end)
|> stub(:close, fn _ -> :ok end)
:ok
end
test "properly encodes URLs with spaces", %{conn: conn} do
url_with_space = "https://example.com/emoji/Pack 1/koronebless.png?foo=bar baz"
result = ReverseProxy.call(conn, url_with_space)
assert result.status == 200
end
test "properly encoded URL should not be altered", %{conn: conn} do
properly_encoded_url = "https://example.com/emoji/Pack%201/koronebless.png?foo=bar+baz"
result = ReverseProxy.call(conn, properly_encoded_url)
assert result.status == 200
end
test "properly encodes URLs with allowed reserved characters", %{conn: conn} do
url_with_reserved_chars = "https://example.com/media/foo/bar !$&'()*+,;=/: @a [baz].mp4"
result = ReverseProxy.call(conn, url_with_reserved_chars)
assert result.status == 200
end
test "properly encodes URLs with unicode in path", %{conn: conn} do
url_with_unicode = "https://example.com/media/unicode 🙂 .gif"
result = ReverseProxy.call(conn, url_with_unicode)
assert result.status == 200
end
end
end
File Metadata
Details
Attached
Mime Type
text/x-diff
Expires
Sat, Aug 8, 2:59 PM (1 d, 20 h)
Storage Engine
blob
Storage Format
Raw Data
Storage Handle
1723632
Default Alt Text
(58 KB)
Attached To
Mode
rPUBE pleroma-upstream
Attached
Detach File
Event Timeline
Log In to Comment