Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
51 changes: 46 additions & 5 deletions src/hackney.erl
Original file line number Diff line number Diff line change
Expand Up @@ -708,17 +708,58 @@ request_with_transport(Method, URL, Headers0, Body, Options0,
_ -> <<Path/binary, "?", Query/binary>>
end,

%% Check for proxy
case maybe_proxy(Transport, URL#hackney_url.scheme, Host, Port, Options) of
Dial = fun() -> maybe_proxy(Transport, URL#hackney_url.scheme, Host, Port, Options) end,
Send = fun(ConnPid, ReqPath, ReqHeaders) ->
do_request(ConnPid, Method, ReqPath, ReqHeaders, Body, Options, URL, Host)
end,
case dial_and_send(Dial, Send, FinalPath, Headers0) of
{{error, {goaway, _}} = Refused, ConnPid} ->
%% The server refused the stream with a GOAWAY, so it never processed the
%% request (RFC 9113 6.8) and the request can go again on a fresh
%% connection whatever its method, as long as the body can be sent twice.
%% Once only: a refusal from the fresh connection is the caller's.
case replayable_body(Body) of
true ->
maybe_stop_unpooled(ConnPid, Options),
{Result, _} = dial_and_send(Dial, Send, FinalPath, Headers0),
Result;
false ->
Refused
end;
{Result, _} ->
Result
end.

%% @private Dial (possibly through a proxy) and send, returning the connection
%% alongside the result so a refused request can be retried elsewhere.
dial_and_send(Dial, Send, FinalPath, Headers0) ->
case Dial() of
{ok, ConnPid} ->
do_request(ConnPid, Method, FinalPath, Headers0, Body, Options, URL, Host);
{Send(ConnPid, FinalPath, Headers0), ConnPid};
{ok, ConnPid, {http_proxy, TargetScheme, TargetHost, TargetPort, ProxyAuth}} ->
%% HTTP proxy mode - use absolute URLs
AbsolutePath = build_absolute_url(TargetScheme, TargetHost, TargetPort, FinalPath),
Headers1 = add_proxy_auth_header(Headers0, ProxyAuth),
do_request(ConnPid, Method, AbsolutePath, Headers1, Body, Options, URL, Host);
{Send(ConnPid, AbsolutePath, Headers1), ConnPid};
Error ->
Error
{Error, undefined}
end.

%% @private A body produced by a fun may already have been consumed by the
%% first attempt. Everything else hackney accepts is re-sendable: binaries,
%% iolists, forms, multipart, files, and `stream', whose chunks only start
%% after the headers were accepted.
replayable_body(Body) when is_function(Body) -> false;
replayable_body({Fun, _State}) when is_function(Fun, 1) -> false;
replayable_body(_Body) -> true.

%% @private A pooled connection a GOAWAY refused has already left its pool and
%% closes on its own. An unpooled one belongs to this caller, who is about to
%% replace it.
maybe_stop_unpooled(ConnPid, Options) ->
case use_pool(Options) of
false -> stop_conn(ConnPid);
_ -> ok
end.

%% @doc Send a request on an existing connection.
Expand Down
12 changes: 9 additions & 3 deletions test/hackney_h2_goaway_server.erl
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,10 @@
%%% Pick(FirstStreamId, SecondStreamId) as last_stream_id
%%% first_data when a stream's first body chunk arrives, with that
%%% stream as last_stream_id
%%% first_stream when the first stream opens, refusing it
%%% and shaped by options:
%%% trigger_on => first | all only the first connection gets the GOAWAY,
%%% or every connection does (default first)
%%% goaway => direct | two_step two_step first sends GOAWAY(2^31-1), the
%%% graceful shutdown of RFC 9113 6.8 (default direct)
%%% answer => boolean() false answers nothing after the GOAWAY, not
Expand All @@ -37,19 +40,20 @@ start(Trigger, Opts) ->
{versions, ['tlsv1.2', 'tlsv1.3']},
{active, false}, {mode, binary}, {reuseaddr, true}]),
{ok, {_, Port}} = ssl:sockname(LSock),
Options = maps:merge(#{goaway => direct, answer => true}, Opts),
Options = maps:merge(#{goaway => direct, answer => true, trigger_on => first}, Opts),
Pid = spawn(fun() -> accept_loop(LSock, Trigger, Options) end),
Url = iolist_to_binary([<<"https://localhost:">>, integer_to_list(Port), <<"/">>]),
{Pid, Url}.

stop(Pid) ->
exit(Pid, kill).

accept_loop(LSock, Trigger, Opts) ->
accept_loop(LSock, Trigger, #{trigger_on := TriggerOn} = Opts) ->
case ssl:transport_accept(LSock, 2000) of
{ok, TSock} ->
spawn(fun() -> serve(TSock, Trigger, Opts) end),
accept_loop(LSock, none, maps:remove(notify, Opts));
Next = case TriggerOn of all -> Trigger; first -> none end,
accept_loop(LSock, Next, maps:remove(notify, Opts));
{error, timeout} -> accept_loop(LSock, Trigger, Opts);
{error, closed} -> ok
end.
Expand Down Expand Up @@ -141,6 +145,8 @@ goaway_for(headers, _Sid, _EndStream, #{trigger := {second_stream, Pick}, seen :
Pick(First, Second);
goaway_for(data, Sid, false, #{trigger := first_data}) ->
Sid;
goaway_for(headers, Sid, _EndStream, #{trigger := first_stream, seen := [Sid]}) ->
max(Sid - 2, 0);
goaway_for(_Kind, _Sid, _EndStream, _St) ->
none.

Expand Down
53 changes: 40 additions & 13 deletions test/hackney_http2_goaway_drain_tests.erl
Original file line number Diff line number Diff line change
@@ -1,10 +1,13 @@
%%% GOAWAY must only fail the streams the peer did not accept.
%%% GOAWAY must only fail the streams the peer did not accept, and those go
%%% again on a fresh connection.
%%%
%%% RFC 9113 6.8: streams up to and including the GOAWAY's last_stream_id may
%%% still be processed, and the peer keeps the connection open to finish them.
%%% hackney used to abort every in-flight stream with {error, {goaway, _}}, so a
%%% request the server went on to complete came back as an error: for a payment
%%% API that is a charge that succeeded and was reported as failed.
%%% API that is a charge that succeeded and was reported as failed. A stream
%%% above last_stream_id was never processed, so request/5 sends it once more
%%% on a fresh connection instead of handing the caller the refusal.
%%%
%%% The server (hackney_h2_goaway_server) holds the first two streams of its
%%% first connection, sends GOAWAY with a chosen last_stream_id, then answers
Expand All @@ -18,7 +21,9 @@
goaway_drain_test_() ->
{foreach, fun setup/0, fun cleanup/1,
[{timeout, 30, fun accepted_streams_complete/0},
{timeout, 30, fun unaccepted_stream_fails_fast/0},
{timeout, 30, fun refused_stream_is_sent_again/0},
{timeout, 30, fun refused_unpooled_request_is_sent_again/0},
{timeout, 30, fun refusal_from_the_fresh_connection_is_returned/0},
{timeout, 30, fun two_step_shutdown/0},
{timeout, 30, fun stalled_drain_ends_with_the_stream/0}]}.

Expand Down Expand Up @@ -53,47 +58,69 @@ accepted_streams_complete() ->
hackney_h2_goaway_server:stop(Server)
end.

%% GOAWAY(last_stream_id = 1) after streams 1 and 3: stream 3 was not accepted
%% and fails at once with the goaway reason, stream 1 still completes. Only
%% the refused stream is reset.
unaccepted_stream_fails_fast() ->
%% GOAWAY(last_stream_id = 1) after streams 1 and 3: stream 3 was not accepted,
%% so it is reset and sent again on a fresh connection, where it is that
%% connection's stream 1. Stream 1 on the first connection still completes.
refused_stream_is_sent_again() ->
{Server, Url} = hackney_h2_goaway_server:start({second_stream, fun(First, _Second) -> First end},
#{notify => self()}),
try
[R1, R2] = concurrent_requests(Url, 2),
?assertEqual({ok, 200, <<"1">>}, R1),
?assertEqual({error, {goaway, no_error}}, R2),
?assertEqual({ok, 200, <<"1">>}, R2),
?assertEqual([3], hackney_h2_goaway_server:rst_streams())
after
hackney_h2_goaway_server:stop(Server)
end.

%% Without a pool the refused connection is the caller's own; it is stopped and
%% the request goes again on a new one.
refused_unpooled_request_is_sent_again() ->
{Server, Url} = hackney_h2_goaway_server:start(first_stream),
try
?assertEqual({ok, 200, <<"1">>}, fetch(Url, [{pool, false}]))
after
hackney_h2_goaway_server:stop(Server)
end.

%% The request is sent again once. When the fresh connection refuses it too,
%% the caller gets that refusal.
refusal_from_the_fresh_connection_is_returned() ->
{Server, Url} = hackney_h2_goaway_server:start(first_stream, #{trigger_on => all}),
try
?assertEqual({error, {goaway, no_error}}, fetch(Url, []))
after
hackney_h2_goaway_server:stop(Server)
end.

%% The graceful shutdown of RFC 9113 6.8: GOAWAY(2^31-1) stops new streams
%% without refusing any, then a second GOAWAY gives the real last_stream_id.
%% The first frame alone must not fail anything.
%% The first frame alone must not fail anything, and the stream the second
%% one refuses goes again on a fresh connection.
two_step_shutdown() ->
{Server, Url} = hackney_h2_goaway_server:start({second_stream, fun(First, _Second) -> First end},
#{goaway => two_step, notify => self()}),
try
[R1, R2] = concurrent_requests(Url, 2),
?assertEqual({ok, 200, <<"1">>}, R1),
?assertEqual({error, {goaway, no_error}}, R2),
?assertEqual({ok, 200, <<"1">>}, R2),
?assertEqual([3], hackney_h2_goaway_server:rst_streams())
after
hackney_h2_goaway_server:stop(Server)
end.

%% A server that accepts a stream and then never answers it must not keep the
%% draining connection around: the stream's own recv_timeout ends it, and with
%% it the connection, so the next request gets a fresh one.
%% it the connection. The refused request went again on a fresh connection,
%% which the request after then shares.
stalled_drain_ends_with_the_stream() ->
{Server, Url} = hackney_h2_goaway_server:start({second_stream, fun(First, _Second) -> First end},
#{answer => false}),
try
[R1, R2] = concurrent_requests(Url, 2, [{recv_timeout, 1000}]),
?assertEqual({error, timeout}, R1),
?assertEqual({error, {goaway, no_error}}, R2),
?assertEqual({ok, 200, <<"1">>}, fetch(Url, []))
?assertEqual({ok, 200, <<"1">>}, R2),
?assertMatch({ok, 200, _}, fetch(Url, []))
after
hackney_h2_goaway_server:stop(Server)
end.
Expand Down
3 changes: 2 additions & 1 deletion test/hackney_http2_goaway_streaming_tests.erl
Original file line number Diff line number Diff line change
Expand Up @@ -96,7 +96,8 @@ accepted_streamed_response_completes() ->
Self ! {self(), R}
end),
timer:sleep(300),
?assertEqual({error, {goaway, no_error}}, fetch(Url)),
%% The plain request is refused and goes again on a fresh connection.
?assertEqual({ok, 200, <<"1">>}, fetch(Url)),
?assertEqual({ok, <<"1">>}, await(P1))
after
hackney_h2_goaway_server:stop(Server)
Expand Down
Loading