diff --git a/rebar.config b/rebar.config index 6f3ab241..8265bd2b 100644 --- a/rebar.config +++ b/rebar.config @@ -55,7 +55,7 @@ %% Pure Erlang QUIC + HTTP/3 stack {quic, "~>2.0"}, %% Pure Erlang HTTP/2 stack - {h2, "~>0.12.3"}, + {h2, "~>0.12.4"}, %% WebTransport client (HTTP/3 and HTTP/2) - powers the wt_* API {webtransport, "~>0.4.7"}, {idna, "~>7.1.0"}, diff --git a/rebar.lock b/rebar.lock index a6ed1ec7..0b4a3758 100644 --- a/rebar.lock +++ b/rebar.lock @@ -1,6 +1,6 @@ {"1.2.0", [{<<"certifi">>,{pkg,<<"certifi">>,<<"2.17.0">>},0}, - {<<"h2">>,{pkg,<<"h2">>,<<"0.12.3">>},0}, + {<<"h2">>,{pkg,<<"h2">>,<<"0.12.4">>},0}, {<<"idna">>,{pkg,<<"idna">>,<<"7.1.0">>},0}, {<<"mimerl">>,{pkg,<<"mimerl">>,<<"1.5.0">>},0}, {<<"parse_trans">>,{pkg,<<"parse_trans">>,<<"3.4.2">>},0}, @@ -10,7 +10,7 @@ [ {pkg_hash,[ {<<"certifi">>, <<"835748414307E15E05B17D0E518190228CE648B08D569A5CC93A85A40F3E5C9B">>}, - {<<"h2">>, <<"20E3FD0E384EC6F586E4736ACD409A57EA87B4A56002D8DCF1132514B3D7600A">>}, + {<<"h2">>, <<"F8CB769CCC17369D185717135353BEE1D50F2CD193A9E8CD998D88CE6CAFC51D">>}, {<<"idna">>, <<"1067A13043538129602D2F2CE6899D8713125C7D19734AA557CE2E3EA55BD4F1">>}, {<<"mimerl">>, <<"F35ACA6F23242339B3666E0AC0702379E362B469D0AEA167F6CC713547E777ED">>}, {<<"parse_trans">>, <<"C352DDC1A0D5E54F9B1654D45F9C432EEF76F9CEA371C55DDFF769EF688FDB74">>}, @@ -19,7 +19,7 @@ {<<"webtransport">>, <<"8E0ABD5875DAAB05C7020B8FC0C3B318FE93AA9329302FE5ECDA2A0F953CF688">>}]}, {pkg_hash_ext,[ {<<"certifi">>, <<"8122798A17F0293C80DAADA25D0F81C7F4D708C73FEF782C7C9B1950E26E4D21">>}, - {<<"h2">>, <<"996AF98698F7DC68BCC7688D70D97384B53DDD0286BA07E6D4A9AC54F1970D32">>}, + {<<"h2">>, <<"238553579D2599E823F1F7B1CACB763362EE6E2DF8E8F3B0C8FC67DA982C0675">>}, {<<"idna">>, <<"6AE959A025BF36DF61A8CAB8508D9654891B5426A84C44D82DEAFFD6DDF8C71F">>}, {<<"mimerl">>, <<"DB648CE065BAE14EA84CA8B5DD123F42F49417CEF693541110BF6F9E9BE9ECC4">>}, {<<"parse_trans">>, <<"4C25347DE3B7C35732D32E69AB43D1CEEE0BEAE3F3B3ADE1B59CBD3DD224D9CA">>}, diff --git a/src/hackney_conn.erl b/src/hackney_conn.erl index 150d954e..c1ae0de1 100644 --- a/src/hackney_conn.erl +++ b/src/hackney_conn.erl @@ -4085,16 +4085,21 @@ collect_h2_aborts(Err, #conn_data{h2_streams = Streams} = Data) -> %% @private Fail the refused streams and keep the rest of the connection %% running. A request body still being sent has no caller parked on it, so %% its stream is kept as a `refused' entry for the caller's next call to find. -abort_h2_streams(StreamIds, ErrorCode, #conn_data{h2_streams = Streams} = Data) -> +%% Each newly refused stream is reset so the h2 layer drops it too; a marker +%% left by an earlier GOAWAY was reset then. +abort_h2_streams(StreamIds, ErrorCode, #conn_data{h2_streams = Streams, + h2_conn = H2Conn} = Data) -> Refused = maps:with(StreamIds, Streams), Replies = h2_abort_replies({goaway, ErrorCode}, Refused), Data1 = maps:fold(fun + (_StreamId, {_Owner, {stream, refused, _}}, D) -> + D; (StreamId, {Owner, {stream, sending}}, D) -> + _ = cancel_h2_stream(H2Conn, StreamId), D#conn_data{h2_streams = maps:put(StreamId, {Owner, {stream, refused, ErrorCode}}, D#conn_data.h2_streams)}; - (_StreamId, {_Owner, {stream, refused, _}}, D) -> - D; (StreamId, _Entry, D) -> + _ = cancel_h2_stream(H2Conn, StreamId), drop_h2_stream(StreamId, D) end, Data, Refused), {Replies, Data1}. diff --git a/test/hackney_h2_goaway_server.erl b/test/hackney_h2_goaway_server.erl index 1d62bb6c..e5c0de33 100644 --- a/test/hackney_h2_goaway_server.erl +++ b/test/hackney_h2_goaway_server.erl @@ -16,9 +16,12 @@ %%% graceful shutdown of RFC 9113 6.8 (default direct) %%% answer => boolean() false answers nothing after the GOAWAY, not %%% even the accepted streams (default true) +%%% notify => pid() gets {goaway_server, rst_stream, StreamId} for +%%% each RST_STREAM on the first connection, then +%%% {goaway_server, done} when that connection ends -module(hackney_h2_goaway_server). --export([start/1, start/2, stop/1]). +-export([start/1, start/2, stop/1, rst_streams/0]). -define(PREFACE, <<"PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n">>). @@ -46,12 +49,31 @@ accept_loop(LSock, Trigger, Opts) -> case ssl:transport_accept(LSock, 2000) of {ok, TSock} -> spawn(fun() -> serve(TSock, Trigger, Opts) end), - accept_loop(LSock, none, Opts); + accept_loop(LSock, none, maps:remove(notify, Opts)); {error, timeout} -> accept_loop(LSock, Trigger, Opts); {error, closed} -> ok end. serve(TSock, Trigger, Opts) -> + serve_conn(TSock, Trigger, Opts), + notify(done, Opts). + +%% @doc The stream ids the client reset on the first connection, in order, +%% once that connection has ended. Use with the notify option. +rst_streams() -> + rst_streams([]). + +rst_streams(Acc) -> + receive + {goaway_server, rst_stream, StreamId} -> rst_streams([StreamId | Acc]); + {goaway_server, done} -> lists:reverse(Acc) + after 10000 -> {timeout, lists:reverse(Acc)} + end. + +notify(Msg, #{notify := Pid}) -> Pid ! {goaway_server, Msg}; +notify(_Msg, _Opts) -> ok. + +serve_conn(TSock, Trigger, Opts) -> case ssl:handshake(TSock, 5000) of {ok, Sock} -> case recv_preface(Sock, <<>>) of @@ -95,6 +117,9 @@ loop(Sock, Buf, St) -> handle(Sock, {settings, _}, St) -> send(Sock, h2_frame:settings_ack()), {continue, St}; handle(Sock, {ping, D}, St) -> send(Sock, h2_frame:ping_ack(D)), {continue, St}; handle(_Sock, {goaway, _, _, _}, _St) -> stop; +handle(_Sock, {rst_stream, Sid, _Code}, #{notify := Pid} = St) -> + Pid ! {goaway_server, rst_stream, Sid}, + {continue, St}; handle(Sock, {headers, Sid, _Block, EndStream, _EndHeaders}, St) -> {continue, on_request_frame(Sock, headers, Sid, EndStream, St)}; %% decode/1 adds the flow-controlled size as a fifth element. diff --git a/test/hackney_http2_goaway_drain_tests.erl b/test/hackney_http2_goaway_drain_tests.erl index a46b7ef1..f805f062 100644 --- a/test/hackney_http2_goaway_drain_tests.erl +++ b/test/hackney_http2_goaway_drain_tests.erl @@ -40,24 +40,30 @@ stop_pool() -> %% complete. A request made while they drain must not land on the draining %% connection, where the peer would ignore it, but dial a fresh one. accepted_streams_complete() -> - {Server, Url} = hackney_h2_goaway_server:start({second_stream, fun(_First, Second) -> Second end}), + {Server, Url} = hackney_h2_goaway_server:start({second_stream, fun(_First, Second) -> Second end}, + #{notify => self()}), try [R1, R2, R3] = concurrent_requests(Url, 3), ?assertEqual({ok, 200, <<"1">>}, R1), ?assertEqual({ok, 200, <<"3">>}, R2), - ?assertEqual({ok, 200, <<"1">>}, R3) + ?assertEqual({ok, 200, <<"1">>}, R3), + %% Nothing was refused, so nothing is reset. + ?assertEqual([], hackney_h2_goaway_server:rst_streams()) after 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. +%% and fails at once with the goaway reason, stream 1 still completes. Only +%% the refused stream is reset. unaccepted_stream_fails_fast() -> - {Server, Url} = hackney_h2_goaway_server:start({second_stream, fun(First, _Second) -> First end}), + {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({error, {goaway, no_error}}, R2), + ?assertEqual([3], hackney_h2_goaway_server:rst_streams()) after hackney_h2_goaway_server:stop(Server) end. @@ -67,11 +73,12 @@ unaccepted_stream_fails_fast() -> %% The first frame alone must not fail anything. two_step_shutdown() -> {Server, Url} = hackney_h2_goaway_server:start({second_stream, fun(First, _Second) -> First end}, - #{goaway => two_step}), + #{goaway => two_step, notify => self()}), try [R1, R2] = concurrent_requests(Url, 2), ?assertEqual({ok, 200, <<"1">>}, R1), - ?assertEqual({error, {goaway, no_error}}, R2) + ?assertEqual({error, {goaway, no_error}}, R2), + ?assertEqual([3], hackney_h2_goaway_server:rst_streams()) 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 2d7c258b..457e84f7 100644 --- a/test/hackney_http2_goaway_streaming_tests.erl +++ b/test/hackney_http2_goaway_streaming_tests.erl @@ -41,7 +41,7 @@ stop_pool() -> %% body chunk arrives. The rest of the body goes out, the response comes back, %% and the connection is gone for the request after. accepted_upload_completes() -> - {Server, Url} = hackney_h2_goaway_server:start(first_data), + {Server, Url} = hackney_h2_goaway_server:start(first_data, #{notify => self()}), try {ok, Conn} = hackney:request(post, Url, [], stream, opts()), ok = hackney:send_body(Conn, <<"abc">>), @@ -49,6 +49,7 @@ accepted_upload_completes() -> ok = hackney:finish_send_body(Conn), {ok, 200, _Headers, Conn} = hackney:start_response(Conn), ?assertEqual({ok, <<"1">>}, hackney:body(Conn)), + ?assertEqual([], hackney_h2_goaway_server:rst_streams()), ?assertEqual({ok, 200, <<"1">>}, fetch(Url)) after hackney_h2_goaway_server:stop(Server) @@ -59,7 +60,8 @@ accepted_upload_completes() -> %% the one that fails, stream 1 still completes, and the connection is gone for %% the request after. refused_upload_fails_on(Call) -> - {Server, Url} = hackney_h2_goaway_server:start({second_stream, fun(First, _Second) -> First end}), + {Server, Url} = hackney_h2_goaway_server:start({second_stream, fun(First, _Second) -> First end}, + #{notify => self()}), try P1 = spawn_fetch(Url), timer:sleep(300), @@ -72,6 +74,8 @@ refused_upload_fails_on(Call) -> end, ?assertEqual({error, {goaway, no_error}}, Result), ?assertEqual({ok, 200, <<"1">>}, await(P1)), + %% The upload was reset when the GOAWAY refused it, and nothing after. + ?assertEqual([3], hackney_h2_goaway_server:rst_streams()), ?assertEqual({ok, 200, <<"1">>}, fetch(Url)) after hackney_h2_goaway_server:stop(Server) @@ -101,14 +105,16 @@ accepted_streamed_response_completes() -> %% Stream 1 is a plain request the server holds, stream 3 a request whose %% response would be read as a stream, and the GOAWAY covers only stream 1. refused_streamed_response_fails_fast() -> - {Server, Url} = hackney_h2_goaway_server:start({second_stream, fun(First, _Second) -> First end}), + {Server, Url} = hackney_h2_goaway_server:start({second_stream, fun(First, _Second) -> First end}, + #{notify => self()}), try P1 = spawn_fetch(Url), timer:sleep(300), {ok, Conn} = connect(Url), ?assertEqual({error, {goaway, no_error}}, hackney:send_request(Conn, {get, <<"/">>, [], <<>>})), - ?assertEqual({ok, 200, <<"1">>}, await(P1)) + ?assertEqual({ok, 200, <<"1">>}, await(P1)), + ?assertEqual([3], hackney_h2_goaway_server:rst_streams()) after hackney_h2_goaway_server:stop(Server) end.