diff --git a/src/hackney.erl b/src/hackney.erl index c3fe472f..c37a3ed5 100644 --- a/src/hackney.erl +++ b/src/hackney.erl @@ -708,17 +708,58 @@ request_with_transport(Method, URL, Headers0, Body, Options0, _ -> <> 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. diff --git a/test/hackney_h2_goaway_server.erl b/test/hackney_h2_goaway_server.erl index e5c0de33..83249b38 100644 --- a/test/hackney_h2_goaway_server.erl +++ b/test/hackney_h2_goaway_server.erl @@ -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 @@ -37,7 +40,7 @@ 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}. @@ -45,11 +48,12 @@ start(Trigger, Opts) -> 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. @@ -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. diff --git a/test/hackney_http2_goaway_drain_tests.erl b/test/hackney_http2_goaway_drain_tests.erl index f805f062..17c56d5b 100644 --- a/test/hackney_http2_goaway_drain_tests.erl +++ b/test/hackney_http2_goaway_drain_tests.erl @@ -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 @@ -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}]}. @@ -53,31 +58,52 @@ 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) @@ -85,15 +111,16 @@ two_step_shutdown() -> %% 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. diff --git a/test/hackney_http2_goaway_streaming_tests.erl b/test/hackney_http2_goaway_streaming_tests.erl index 457e84f7..bb47f4fa 100644 --- a/test/hackney_http2_goaway_streaming_tests.erl +++ b/test/hackney_http2_goaway_streaming_tests.erl @@ -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)