Skip to content
Merged
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
2 changes: 1 addition & 1 deletion rebar.config
Original file line number Diff line number Diff line change
Expand Up @@ -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"},
Expand Down
6 changes: 3 additions & 3 deletions rebar.lock
Original file line number Diff line number Diff line change
@@ -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},
Expand All @@ -10,7 +10,7 @@
[
{pkg_hash,[
{<<"certifi">>, <<"835748414307E15E05B17D0E518190228CE648B08D569A5CC93A85A40F3E5C9B">>},
{<<"h2">>, <<"20E3FD0E384EC6F586E4736ACD409A57EA87B4A56002D8DCF1132514B3D7600A">>},
{<<"h2">>, <<"F8CB769CCC17369D185717135353BEE1D50F2CD193A9E8CD998D88CE6CAFC51D">>},
{<<"idna">>, <<"1067A13043538129602D2F2CE6899D8713125C7D19734AA557CE2E3EA55BD4F1">>},
{<<"mimerl">>, <<"F35ACA6F23242339B3666E0AC0702379E362B469D0AEA167F6CC713547E777ED">>},
{<<"parse_trans">>, <<"C352DDC1A0D5E54F9B1654D45F9C432EEF76F9CEA371C55DDFF769EF688FDB74">>},
Expand All @@ -19,7 +19,7 @@
{<<"webtransport">>, <<"8E0ABD5875DAAB05C7020B8FC0C3B318FE93AA9329302FE5ECDA2A0F953CF688">>}]},
{pkg_hash_ext,[
{<<"certifi">>, <<"8122798A17F0293C80DAADA25D0F81C7F4D708C73FEF782C7C9B1950E26E4D21">>},
{<<"h2">>, <<"996AF98698F7DC68BCC7688D70D97384B53DDD0286BA07E6D4A9AC54F1970D32">>},
{<<"h2">>, <<"238553579D2599E823F1F7B1CACB763362EE6E2DF8E8F3B0C8FC67DA982C0675">>},
{<<"idna">>, <<"6AE959A025BF36DF61A8CAB8508D9654891B5426A84C44D82DEAFFD6DDF8C71F">>},
{<<"mimerl">>, <<"DB648CE065BAE14EA84CA8B5DD123F42F49417CEF693541110BF6F9E9BE9ECC4">>},
{<<"parse_trans">>, <<"4C25347DE3B7C35732D32E69AB43D1CEEE0BEAE3F3B3ADE1B59CBD3DD224D9CA">>},
Expand Down
11 changes: 8 additions & 3 deletions src/hackney_conn.erl
Original file line number Diff line number Diff line change
Expand Up @@ -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}.
Expand Down
29 changes: 27 additions & 2 deletions test/hackney_h2_goaway_server.erl
Original file line number Diff line number Diff line change
Expand Up @@ -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">>).

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down
21 changes: 14 additions & 7 deletions test/hackney_http2_goaway_drain_tests.erl
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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.
Expand Down
14 changes: 10 additions & 4 deletions test/hackney_http2_goaway_streaming_tests.erl
Original file line number Diff line number Diff line change
Expand Up @@ -41,14 +41,15 @@ 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">>),
timer:sleep(300),
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)
Expand All @@ -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),
Expand All @@ -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)
Expand Down Expand Up @@ -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.
Expand Down
Loading