Skip to content
2 changes: 1 addition & 1 deletion .dialyzer_ignore
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
lib/mint/tunnel_proxy.ex:50
lib/mint/http1.ex:1069
lib/mint/http1.ex:1104
lib/mint/unsafe_proxy.ex:173
lib/mint/unsafe_proxy.ex:198
test/support
Expand Down
222 changes: 143 additions & 79 deletions lib/mint/http1.ex
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,19 @@ defmodule Mint.HTTP1 do
* `:request_body_is_streaming` - when you call `request/5` to send a new
request but another request is already streaming.

* `:request_is_not_streaming` - when you call `stream_request_body/3` for a
request whose body is not being streamed, for example after sending `:eof`.

* `:unknown_request_to_stream` - when you call `stream_request_body/3` with a
request reference that doesn't belong to this connection.

* `:unprocessed` - when a pipelined request gets no response because the server
answered a previous request with a `connection: close` header. The server didn't
process the request, so it can be retried on a new connection. Requests pipelined
behind a response that closes the connection in other ways get a
`Mint.TransportError` with reason `:closed` instead, since the server might have
processed them.

* `{:unexpected_data, data}` - when unexpected data is received from the server.

* `:invalid_status_line` - when the HTTP/1 status line is invalid.
Expand Down Expand Up @@ -350,17 +363,19 @@ defmodule Mint.HTTP1 do
),
:ok <- transport.send(socket, iodata) do
request_ref = make_ref()
request = new_request(request_ref, method, body, encoding)

case request.state do
{:stream_request, _} ->
conn = %{conn | streaming_request: request}
{:ok, conn, request_ref}
conn = enqueue_request(conn, new_request(request_ref, method))

# The request is enqueued right away so that a response the server
# sends before the body is complete (such as 100 Continue or an early
# 413) is parsed rather than treated as unexpected data.
conn =
if body == :stream do
%{conn | streaming_request: %{ref: request_ref, encoding: encoding}}
else
conn
end

_ ->
conn = enqueue_request(conn, request)
{:ok, conn, request_ref}
end
{:ok, conn, request_ref}
else
{:error, %TransportError{reason: :closed} = error} ->
conn = internal_close(conn)
Expand Down Expand Up @@ -405,26 +420,28 @@ defmodule Mint.HTTP1 do
iodata() | :eof | {:eof, trailer_headers :: Types.headers()}
) ::
{:ok, t()} | {:error, t(), Types.error()}
def stream_request_body(%__MODULE__{state: :closed} = conn, _request_ref, _chunk) do
{:error, conn, wrap_error(:closed)}
end

def stream_request_body(
%__MODULE__{streaming_request: %{state: {:stream_request, :identity}, ref: ref}} = conn,
%__MODULE__{streaming_request: %{encoding: :identity, ref: ref}} = conn,
ref,
:eof
) do
request = %{conn.streaming_request | state: :status}
conn = enqueue_request(%{conn | streaming_request: nil}, request)
{:ok, conn}
{:ok, %{conn | streaming_request: nil}}
end

def stream_request_body(
%__MODULE__{streaming_request: %{state: {:stream_request, :identity}, ref: ref}} = conn,
%__MODULE__{streaming_request: %{encoding: :identity, ref: ref}} = conn,
ref,
{:eof, _trailer_headers}
) do
{:error, conn, wrap_error(:trailing_headers_but_not_chunked_encoding)}
end

def stream_request_body(
%__MODULE__{streaming_request: %{state: {:stream_request, :identity}, ref: ref}} = conn,
%__MODULE__{streaming_request: %{encoding: :identity, ref: ref}} = conn,
ref,
body
) do
Expand All @@ -442,25 +459,16 @@ defmodule Mint.HTTP1 do
end

def stream_request_body(
%__MODULE__{streaming_request: %{state: {:stream_request, :chunked}, ref: ref}} = conn,
%__MODULE__{streaming_request: %{encoding: :chunked, ref: ref}} = conn,
ref,
chunk
) do
with {:ok, chunk} <- validate_chunk(conn, chunk),
:ok <- conn.transport.send(conn.socket, Request.encode_chunk(chunk)) do
case chunk do
:eof ->
request = %{conn.streaming_request | state: :status}
conn = enqueue_request(%{conn | streaming_request: nil}, request)
{:ok, conn}

{:eof, _trailer_headers} ->
request = %{conn.streaming_request | state: :status}
conn = enqueue_request(%{conn | streaming_request: nil}, request)
{:ok, conn}

_other ->
{:ok, conn}
:eof -> {:ok, %{conn | streaming_request: nil}}
{:eof, _trailer_headers} -> {:ok, %{conn | streaming_request: nil}}
_other -> {:ok, conn}
end
else
:empty_chunk ->
Expand All @@ -475,6 +483,19 @@ defmodule Mint.HTTP1 do
end
end

def stream_request_body(%__MODULE__{} = conn, request_ref, _chunk)
when is_reference(request_ref) do
known? =
(conn.request != nil and conn.request.ref == request_ref) or
Enum.any?(:queue.to_list(conn.requests), &(&1.ref == request_ref))

if known? do
{:error, conn, wrap_error(:request_is_not_streaming)}
else
{:error, conn, wrap_error(:unknown_request_to_stream)}
end
end

defp validate_chunk(conn, {:eof, trailers}) do
trailers = Headers.from_raw(trailers)

Expand Down Expand Up @@ -554,16 +575,16 @@ defmodule Mint.HTTP1 do
end
end

defp handle_close(%__MODULE__{request: request} = conn) do
conn = internal_close(conn)
conn = request_done(conn)
defp handle_close(%__MODULE__{request: %{body: :until_closed} = request} = conn) do
conn = pop_request(conn)
responses = [{:done, request.ref}]
{conn, responses} = close_after_response(conn, responses, conn.transport.wrap_error(:closed))
{:ok, conn, Enum.reverse(responses)}
end

if request && request.body == :until_closed do
conn = put_in(conn.state, :closed)
{:ok, conn, [{:done, request.ref}]}
else
{:error, conn, conn.transport.wrap_error(:closed), []}
end
defp handle_close(conn) do
conn = conn |> internal_close() |> pop_request()
{:error, conn, conn.transport.wrap_error(:closed), []}
end

defp handle_transport_error(conn, error) do
Expand Down Expand Up @@ -635,10 +656,8 @@ defmodule Mint.HTTP1 do
@spec open_request_count(t()) :: non_neg_integer()
def open_request_count(%__MODULE__{} = conn) do
case conn do
%{request: nil, streaming_request: nil} -> 0
%{request: nil} -> 1
%{streaming_request: nil} -> 1 + :queue.len(conn.requests)
_ -> 2 + :queue.len(conn.requests)
%{request: nil} -> 0
_ -> 1 + :queue.len(conn.requests)
end
end

Expand Down Expand Up @@ -815,11 +834,23 @@ defmodule Mint.HTTP1 do
end
end

defp decode_body(:none, conn, data, request_ref, responses) do
conn = put_in(conn.buffer, data)
conn = request_done(conn)
responses = [{:done, request_ref} | responses]
{:ok, conn, responses}
# A successful CONNECT switches the connection to tunnel mode, so bytes after
# the header section belong to the tunnel rather than to another response.
defp decode_body(
:none,
%{request: %{method: "CONNECT", status: status}} = conn,
data,
_request_ref,
responses
)
when status in 200..299 do
{conn, responses} = request_done(conn, responses)
{:ok, %{conn | buffer: data}, responses}
end

defp decode_body(:none, conn, data, _request_ref, responses) do
{conn, responses} = request_done(conn, responses)
next_request(conn, data, responses)
end

# Informational (1xx) responses have no body and must not finalize the
Expand All @@ -844,10 +875,9 @@ defmodule Mint.HTTP1 do
decode(:status, conn, data, responses)
end

defp decode_body(:single, conn, data, request_ref, responses) do
defp decode_body(:single, conn, data, _request_ref, responses) do
{conn, responses} = add_body(conn, data, responses)
conn = request_done(conn)
responses = [{:done, request_ref} | responses]
{conn, responses} = request_done(conn, responses)
{:ok, conn, responses}
end

Expand All @@ -856,7 +886,7 @@ defmodule Mint.HTTP1 do
{:ok, conn, responses}
end

defp decode_body({:content_length, length}, conn, data, request_ref, responses) do
defp decode_body({:content_length, length}, conn, data, _request_ref, responses) do
cond do
length > byte_size(data) ->
conn = put_in(conn.request.body, {:content_length, length - byte_size(data)})
Expand All @@ -866,8 +896,7 @@ defmodule Mint.HTTP1 do
length <= byte_size(data) ->
{body, rest} = :erlang.split_binary(data, length)
{conn, responses} = add_body(conn, body, responses)
conn = request_done(conn)
responses = [{:done, request_ref} | responses]
{conn, responses} = request_done(conn, responses)
next_request(conn, rest, responses)
end
end
Expand Down Expand Up @@ -966,12 +995,8 @@ defmodule Mint.HTTP1 do
{:ok, _request} ->
headers = Headers.remove_unallowed_trailer(headers)

responses = [
{:done, conn.request.ref}
| add_trailer_headers(headers, conn.request.ref, responses)
]

conn = request_done(conn)
responses = add_trailer_headers(headers, conn.request.ref, responses)
{conn, responses} = request_done(conn, responses)
next_request(conn, rest, responses)

{:error, reason} ->
Expand Down Expand Up @@ -1033,10 +1058,20 @@ defmodule Mint.HTTP1 do
end
end

# A response that closes the connection is the last one the server sends on
# it, so anything after it can't be a response to a queued request.
defp next_request(%{state: :closed} = conn, _data, responses) do
{:ok, %{conn | buffer: ""}, responses}
end

defp next_request(%{request: nil} = conn, "", responses) do
{:ok, %{conn | buffer: ""}, responses}
end

# Bytes left over after the last in-flight response would otherwise be
# delivered as the response to whichever request is issued next.
defp next_request(%{request: nil} = conn, data, responses) do
# TODO: Figure out if we should keep buffering even though there are no
# requests in flight
{:ok, %{conn | buffer: data}, responses}
{:error, conn, wrap_error({:unexpected_data, data}), responses}
end

defp next_request(conn, data, responses) do
Expand Down Expand Up @@ -1137,23 +1172,46 @@ defmodule Mint.HTTP1 do
# lifetime of the underlying socket. In particular, HTTP/1.0 responses are
# otherwise treated as non-persistent and would close the newly-established
# tunnel before the caller can use it.
defp request_done(%{request: %{method: "CONNECT", status: status}} = conn)
defp request_done(%{request: %{method: "CONNECT", status: status} = request} = conn, responses)
when status in 200..299 do
pop_request(conn)
{pop_request(conn), [{:done, request.ref} | responses]}
end

defp request_done(%{request: request} = conn) do
defp request_done(%{request: request} = conn, responses) do
conn = pop_request(conn)
responses = [{:done, request.ref} | responses]

cond do
!request -> conn
"close" in request.connection -> internal_close(conn)
request.version >= {1, 1} -> conn
"keep-alive" in request.connection -> conn
true -> internal_close(conn)
# The server doesn't process any requests after one it answers with
# "Connection: close" (RFC 9112 section 9.6), so the queued requests can
# be retried.
"close" in request.connection ->
close_after_response(conn, responses, wrap_error(:unprocessed))

request.version >= {1, 1} ->
{conn, responses}

"keep-alive" in request.connection ->
{conn, responses}

true ->
close_after_response(conn, responses, conn.transport.wrap_error(:closed))
end
end

# Requests pipelined behind a response that closes the connection never get a
# response of their own.
defp close_after_response(conn, responses, error) do
requests = if conn.request, do: [conn.request | :queue.to_list(conn.requests)], else: []

responses =
Enum.reduce(requests, responses, fn request, responses ->
[{:error, request.ref, error} | responses]
end)

{internal_close(%{conn | request: nil, requests: :queue.new()}), responses}
end

defp pop_request(conn) do
case :queue.out(conn.requests) do
{{:value, request}, requests} ->
Expand Down Expand Up @@ -1254,17 +1312,10 @@ defmodule Mint.HTTP1 do
:ok
end

defp new_request(ref, method, body, encoding) do
state =
if body == :stream do
{:stream_request, encoding}
else
:status
end

defp new_request(ref, method) do
%{
ref: ref,
state: state,
state: :status,
method: method,
version: nil,
status: nil,
Expand Down Expand Up @@ -1344,10 +1395,23 @@ defmodule Mint.HTTP1 do
"the connection is closed"
end

def format_error(:unprocessed) do
"request was not processed because the server answered a previous request with " <>
"\"connection: close\", so it's safe to retry on a new connection"
end

def format_error(:request_body_is_streaming) do
"a request body is currently streaming, so no new requests can be issued"
end

def format_error(:request_is_not_streaming) do
"can't send more data on a request that is not streaming its body"
end

def format_error(:unknown_request_to_stream) do
"can't stream the request body because the request is not known to this connection"
end

def format_error({:unexpected_data, data}) do
"received unexpected data: " <> inspect(data)
end
Expand Down
Loading
Loading