Packages
bandit
0.7.5
1.12.0
1.11.1
1.11.0
1.10.4
1.10.3
1.10.2
1.10.1
1.10.0
retired
1.9.0
1.8.0
1.7.0
1.6.11
1.6.10
1.6.9
1.6.8
1.6.7
1.6.6
1.6.5
1.6.4
1.6.3
1.6.2
1.6.1
1.6.0
1.5.7
1.5.6
1.5.5
1.5.4
1.5.3
1.5.2
1.5.1
1.5.0
1.4.2
1.4.1
1.4.0
1.3.0
1.2.3
1.2.2
1.2.1
1.2.0
1.1.3
1.1.2
1.1.1
1.1.0
1.0.0
1.0.0-pre.18
1.0.0-pre.17
1.0.0-pre.16
1.0.0-pre.15
1.0.0-pre.14
1.0.0-pre.13
1.0.0-pre.12
1.0.0-pre.11
1.0.0-pre.10
1.0.0-pre.9
1.0.0-pre.8
1.0.0-pre.7
1.0.0-pre.6
1.0.0-pre.5
1.0.0-pre.4
1.0.0-pre.3
1.0.0-pre.2
1.0.0-pre.1
0.7.7
0.7.6
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.11
0.6.10
0.6.9
0.6.8
0.6.7
0.6.6
0.6.5
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
0.5.11
0.5.10
0.5.9
0.5.8
0.5.7
0.5.6
0.5.5
0.5.4
0.5.3
0.5.2
0.5.1
0.5.0
0.4.10
0.4.9
0.4.8
0.4.7
0.4.6
0.4.5
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.3.9
0.3.8
0.3.7
0.3.6
0.3.5
0.3.4
0.3.3
0.3.2
0.2.3
0.2.2
0.2.1
0.2.0
0.1.1
0.1.0
A pure-Elixir HTTP server built for Plug & WebSock apps
Security advisory:
This version has known vulnerabilities.
View advisories
Current section
Files
Jump to
Current section
Files
lib/bandit/http2/connection.ex
defmodule Bandit.HTTP2.Connection do
@moduledoc false
# Represents the state of an HTTP/2 connection
require Logger
alias ThousandIsland.{Handler, Socket}
alias Bandit.HTTP2.{
Connection,
Errors,
FlowControl,
Frame,
Settings,
Stream,
StreamCollection
}
defstruct local_settings: %Settings{},
remote_settings: %Settings{},
fragment_frame: nil,
send_hpack_state: HPAX.new(4096),
recv_hpack_state: HPAX.new(4096),
send_window_size: 65_535,
recv_window_size: 65_535,
streams: %StreamCollection{},
pending_sends: [],
plug: nil,
opts: []
@typedoc "A description of a connection error"
@type error :: {:connection, Errors.error_code(), String.t()}
@typedoc "Encapsulates the state of an HTTP/2 connection"
@type t :: %__MODULE__{
local_settings: Settings.t(),
remote_settings: Settings.t(),
fragment_frame: Frame.Headers.t() | nil,
send_hpack_state: HPAX.Table.t(),
recv_hpack_state: HPAX.Table.t(),
send_window_size: non_neg_integer(),
recv_window_size: non_neg_integer(),
streams: StreamCollection.t(),
pending_sends: [{Stream.stream_id(), iodata(), boolean(), fun()}],
plug: Bandit.plug(),
opts: keyword()
}
@spec init(Socket.t(), Bandit.plug(), keyword()) :: {:ok, t()}
def init(socket, plug, opts) do
connection = %__MODULE__{
plug: plug,
opts: opts,
local_settings: struct!(Settings, Keyword.get(opts, :default_local_settings, []))
}
# Send SETTINGS frame per RFC9113§3.4
%Frame.Settings{ack: false, settings: connection.local_settings}
|> send_frame(socket, connection)
{:ok, connection}
end
#
# Receiving while expecting CONTINUATION frames is a special case (RFC9113§6.10); handle it first
#
@spec handle_frame(Frame.frame(), Socket.t(), t()) :: Handler.handler_result()
def handle_frame(
%Frame.Continuation{end_headers: true, stream_id: stream_id} = frame,
socket,
%Connection{fragment_frame: %Frame.Headers{stream_id: stream_id}} = connection
) do
header_block = connection.fragment_frame.fragment <> frame.fragment
header_frame = %{connection.fragment_frame | end_headers: true, fragment: header_block}
handle_frame(header_frame, socket, %{connection | fragment_frame: nil})
end
def handle_frame(
%Frame.Continuation{end_headers: false, stream_id: stream_id} = frame,
_socket,
%Connection{fragment_frame: %Frame.Headers{stream_id: stream_id}} = connection
) do
header_block = connection.fragment_frame.fragment <> frame.fragment
header_frame = %{connection.fragment_frame | fragment: header_block}
{:continue, %{connection | fragment_frame: header_frame}}
end
def handle_frame(_frame, socket, %Connection{fragment_frame: %Frame.Headers{}} = connection) do
shutdown_connection(
Errors.protocol_error(),
"Expected CONTINUATION frame (RFC9113§6.10)",
socket,
connection
)
end
#
# Connection-level receiving
#
def handle_frame(%Frame.Settings{ack: true}, _socket, connection) do
{:continue, connection}
end
def handle_frame(%Frame.Settings{ack: false} = frame, socket, connection) do
%Frame.Settings{ack: true} |> send_frame(socket, connection)
streams =
connection.streams
|> StreamCollection.update_initial_send_window_size(frame.settings.initial_window_size)
|> StreamCollection.update_max_concurrent_streams(frame.settings.max_concurrent_streams)
send_hpack_state = HPAX.resize(connection.send_hpack_state, frame.settings.header_table_size)
do_pending_sends(socket, %{
connection
| remote_settings: frame.settings,
streams: streams,
send_hpack_state: send_hpack_state
})
end
def handle_frame(%Frame.Ping{ack: true}, _socket, connection) do
{:continue, connection}
end
def handle_frame(%Frame.Ping{ack: false} = frame, socket, connection) do
%Frame.Ping{ack: true, payload: frame.payload} |> send_frame(socket, connection)
{:continue, connection}
end
def handle_frame(%Frame.Goaway{}, socket, connection) do
shutdown_connection(Errors.no_error(), "Received GOAWAY", socket, connection)
end
def handle_frame(%Frame.WindowUpdate{stream_id: 0} = frame, socket, connection) do
case FlowControl.update_send_window(connection.send_window_size, frame.size_increment) do
{:ok, new_window} ->
do_pending_sends(socket, %{connection | send_window_size: new_window})
{:error, error} ->
shutdown_connection(Errors.flow_control_error(), error, socket, connection)
end
end
#
# Stream-level receiving
#
def handle_frame(%Frame.WindowUpdate{} = frame, socket, connection) do
with {:ok, stream} <- StreamCollection.get_stream(connection.streams, frame.stream_id),
{:ok, stream} <- Stream.recv_window_update(stream, frame.size_increment),
{:ok, streams} <- StreamCollection.put_stream(connection.streams, stream) do
do_pending_sends(socket, %{connection | streams: streams})
else
{:error, {:connection, error_code, error_message}} ->
shutdown_connection(error_code, error_message, socket, connection)
{:error, {:stream, stream_id, error_code, error_message}} ->
handle_stream_error(stream_id, error_code, error_message, socket, connection)
{:error, error} ->
shutdown_connection(Errors.internal_error(), error, socket, connection)
end
end
def handle_frame(
%Frame.Headers{stream_id: stream_id, stream_dependency: stream_id},
socket,
connection
) do
# This is no longer mentioned in RFC9113, but error anyway since it's squarely illogical
handle_stream_error(
stream_id,
Errors.protocol_error(),
"Stream cannot list itself as a dependency (RFC7540§5.3.1)",
socket,
connection
)
end
def handle_frame(%Frame.Headers{end_headers: true} = frame, socket, connection) do
with block <- frame.fragment,
end_stream <- frame.end_stream,
{:hpack, {:ok, headers, recv_hpack_state}} <-
{:hpack, HPAX.decode(block, connection.recv_hpack_state)},
{:ok, stream} <- StreamCollection.get_stream(connection.streams, frame.stream_id),
true <- accept_stream?(connection),
true <- accept_headers?(headers, connection.opts, stream),
transport_info <- build_transport_info(socket),
{:ok, stream} <-
Stream.recv_headers(
stream,
transport_info,
headers,
end_stream,
connection.plug,
connection.opts
),
{:ok, stream} <- Stream.recv_end_of_stream(stream, end_stream),
{:ok, streams} <- StreamCollection.put_stream(connection.streams, stream) do
{:continue, %{connection | recv_hpack_state: recv_hpack_state, streams: streams}}
else
{:hpack, _} ->
shutdown_connection(Errors.compression_error(), "Header decode error", socket, connection)
{:error, {:connection, error_code, error_message}} ->
shutdown_connection(error_code, error_message, socket, connection)
{:error, {:stream, stream_id, error_code, error_message}} ->
handle_stream_error(stream_id, error_code, error_message, socket, connection)
{:error, error} ->
shutdown_connection(Errors.internal_error(), error, socket, connection)
end
end
def handle_frame(%Frame.Headers{end_headers: false} = frame, _socket, connection) do
{:continue, %{connection | fragment_frame: frame}}
end
def handle_frame(%Frame.Continuation{}, socket, connection) do
shutdown_connection(
Errors.protocol_error(),
"Received unexpected CONTINUATION frame (RFC9113§6.10)",
socket,
connection
)
end
def handle_frame(%Frame.Data{} = frame, socket, connection) do
{connection_recv_window_size, connection_window_increment} =
FlowControl.compute_recv_window(connection.recv_window_size, byte_size(frame.data))
with {:ok, stream} <- StreamCollection.get_stream(connection.streams, frame.stream_id),
{:ok, stream, stream_window_increment} <- Stream.recv_data(stream, frame.data),
{:ok, stream} <- Stream.recv_end_of_stream(stream, frame.end_stream),
{:ok, streams} <- StreamCollection.put_stream(connection.streams, stream) do
if connection_window_increment > 0 do
%Frame.WindowUpdate{stream_id: 0, size_increment: connection_window_increment}
|> send_frame(socket, connection)
end
if stream_window_increment > 0 do
%Frame.WindowUpdate{stream_id: frame.stream_id, size_increment: stream_window_increment}
|> send_frame(socket, connection)
end
{:continue, %{connection | recv_window_size: connection_recv_window_size, streams: streams}}
else
{:error, {:connection, error_code, error_message}} ->
shutdown_connection(error_code, error_message, socket, connection)
{:error, {:stream, stream_id, error_code, error_message}} ->
# If we're erroring out on a stream error, RFC9113§6.9 stipulates that we MUST take into
# account the sizes of errored frames. As such, ensure that we update our connection
# window to reflect that space taken up by this frame. We needn't worry about the stream's
# window since we're shutting it down anyway
connection = %{connection | recv_window_size: connection_recv_window_size}
handle_stream_error(stream_id, error_code, error_message, socket, connection)
{:error, error} ->
shutdown_connection(Errors.internal_error(), error, socket, connection)
end
end
def handle_frame(
%Frame.Priority{stream_id: stream_id, dependent_stream_id: stream_id},
socket,
connection
) do
# This is no longer mentioned in RFC9113, but error anyway since it's squarely illogical
handle_stream_error(
stream_id,
Errors.protocol_error(),
"Stream cannot list itself as a dependency (RFC7540§5.3.1)",
socket,
connection
)
end
def handle_frame(%Frame.Priority{}, _socket, connection) do
{:continue, connection}
end
def handle_frame(%Frame.RstStream{} = frame, socket, connection) do
with {:ok, stream} <- StreamCollection.get_stream(connection.streams, frame.stream_id),
{:ok, stream} <- Stream.recv_rst_stream(stream, frame.error_code),
{:ok, streams} <- StreamCollection.put_stream(connection.streams, stream) do
{:continue, %{connection | streams: streams}}
else
{:error, {:connection, error_code, error_message}} ->
shutdown_connection(error_code, error_message, socket, connection)
{:error, error} ->
shutdown_connection(Errors.internal_error(), error, socket, connection)
end
end
# Catch-all handler for unknown frame types
def handle_frame(%Frame.Unknown{} = frame, _socket, connection) do
Logger.warning("Unknown frame (#{inspect(Map.from_struct(frame))})")
{:continue, connection}
end
defp accept_stream?(connection) do
max_requests = Keyword.get(connection.opts, :max_requests, 0)
if max_requests != 0 && StreamCollection.stream_count(connection.streams) >= max_requests do
{:error, {:connection, Errors.refused_stream(), "Connection count exceeded"}}
else
true
end
end
defp accept_headers?(headers, opts, stream) do
with true <- valid_header_count?(headers, Keyword.get(opts, :max_header_count, 50), stream),
max_header_key_length <- Keyword.get(opts, :max_header_key_length, 10_000),
max_header_value_length <- Keyword.get(opts, :max_header_value_length, 10_000) do
valid_headers?(headers, max_header_key_length, max_header_value_length, stream)
end
end
defp valid_header_count?(headers, max, _stream) when length(headers) <= max, do: true
defp valid_header_count?(_, _, stream),
do:
{:error,
{:stream, stream.stream_id, Errors.frame_size_error(), "Request contains too many headers"}}
defp valid_headers?(headers, max_key_length, max_value_length, stream) do
if Enum.any?(headers, fn
{_, nil} ->
false
{key, value} ->
String.length(key) > max_key_length || String.length(value) > max_value_length
end) do
{:error,
{:stream, stream.stream_id, Errors.frame_size_error(),
"Request contains overlong header(s)"}}
else
true
end
end
defp build_transport_info(socket) do
{ThousandIsland.Socket.secure?(socket), ThousandIsland.Socket.local_info(socket),
ThousandIsland.Socket.peer_info(socket), ThousandIsland.Socket.telemetry_span(socket)}
end
# Shared logic to send any pending frames upon adjustment of our send window
defp do_pending_sends(socket, connection) do
connection.pending_sends
|> Enum.reverse()
|> Enum.reduce_while({:continue, connection}, fn pending_send, {:continue, connection} ->
connection = connection |> Map.update!(:pending_sends, &List.delete(&1, pending_send))
{stream_id, rest, end_stream, on_unblock} = pending_send
case do_send_data(stream_id, rest, end_stream, on_unblock, socket, connection) do
{:ok, true, connection} ->
on_unblock.()
{:cont, {:continue, connection}}
{:ok, false, connection} ->
{:cont, {:continue, connection}}
{:error, error} ->
{:halt, {:error, error, connection}}
end
end)
end
#
# Error handling on receipt of frames
#
@spec shutdown_connection(Errors.error_code(), term(), Socket.t(), t()) ::
{:close, t()} | {:error, term(), t()}
def shutdown_connection(error_code, reason, socket, connection) do
last_remote_stream_id = StreamCollection.last_remote_stream_id(connection.streams)
%Frame.Goaway{last_stream_id: last_remote_stream_id, error_code: error_code}
|> send_frame(socket, connection)
if error_code == Errors.no_error() do
{:close, connection}
else
{:error, reason, connection}
end
end
defp handle_stream_error(stream_id, error_code, reason, socket, connection) do
{:ok, stream} = StreamCollection.get_stream(connection.streams, stream_id)
Stream.terminate_stream(stream, {:bandit, reason})
%Frame.RstStream{stream_id: stream_id, error_code: error_code}
|> send_frame(socket, connection)
{:continue, connection}
end
#
# Sending logic
#
# All callers of functions below will be from stream tasks
#
#
# Stream-level sending
#
@spec send_headers(Stream.stream_id(), Plug.Conn.headers(), boolean(), Socket.t(), t()) ::
{:ok, t()} | {:error, term()}
def send_headers(stream_id, headers, end_stream, socket, connection) do
with enc_headers <- Enum.map(headers, fn {key, value} -> {:store, key, value} end),
{block, send_hpack_state} <- HPAX.encode(enc_headers, connection.send_hpack_state),
{:ok, stream} <- StreamCollection.get_stream(connection.streams, stream_id),
{:ok, stream} <- Stream.send_headers(stream),
{:ok, stream} <- Stream.send_end_of_stream(stream, end_stream),
{:ok, streams} <- StreamCollection.put_stream(connection.streams, stream) do
%Frame.Headers{
stream_id: stream_id,
end_stream: end_stream,
fragment: block
}
|> send_frame(socket, connection)
{:ok, %{connection | send_hpack_state: send_hpack_state, streams: streams}}
end
end
@spec send_data(Stream.stream_id(), iodata(), boolean(), fun(), Socket.t(), t()) ::
{:ok, boolean(), t()} | {:error, term()}
def send_data(stream_id, data, end_stream, on_unblock, socket, connection) do
do_send_data(stream_id, data, end_stream, on_unblock, socket, connection)
end
defp do_send_data(stream_id, data, end_stream, on_unblock, socket, connection) do
with {:ok, stream} <- StreamCollection.get_stream(connection.streams, stream_id),
stream_window_size <- Stream.get_send_window_size(stream),
connection_window_size <- connection.send_window_size,
max_bytes_to_send <- max(min(stream_window_size, connection_window_size), 0),
{data_to_send, bytes_to_send, rest} <- split_data(data, max_bytes_to_send),
{:ok, stream} <- Stream.send_data(stream, bytes_to_send),
connection <- %{connection | send_window_size: connection_window_size - bytes_to_send},
end_stream_to_send <- end_stream && rest == <<>>,
{:ok, stream} <- Stream.send_end_of_stream(stream, end_stream_to_send),
{:ok, streams} <- StreamCollection.put_stream(connection.streams, stream) do
if end_stream_to_send || IO.iodata_length(data_to_send) > 0 do
%Frame.Data{stream_id: stream_id, end_stream: end_stream_to_send, data: data_to_send}
|> send_frame(socket, connection)
end
if byte_size(rest) == 0 do
{:ok, true, %{connection | streams: streams}}
else
pending_sends = [{stream_id, rest, end_stream, on_unblock} | connection.pending_sends]
{:ok, false, %{connection | streams: streams, pending_sends: pending_sends}}
end
end
end
defp split_data(data, desired_length) do
data_length = IO.iodata_length(data)
if data_length <= desired_length do
{data, data_length, <<>>}
else
<<to_send::binary-size(desired_length), rest::binary>> = IO.iodata_to_binary(data)
{to_send, desired_length, rest}
end
end
@spec stream_terminated(pid(), term(), Socket.t(), t()) :: {:ok, t()} | {:error, term()}
def stream_terminated(pid, reason, socket, connection) do
with {:ok, stream} <- StreamCollection.get_active_stream_by_pid(connection.streams, pid),
{:ok, stream, error_code} <- Stream.stream_terminated(stream, reason),
{:ok, streams} <- StreamCollection.put_stream(connection.streams, stream) do
if !is_nil(error_code) do
%Frame.RstStream{stream_id: stream.stream_id, error_code: error_code}
|> send_frame(socket, connection)
end
{:ok, %{connection | streams: streams}}
end
end
#
# Utility functions
#
defp send_frame(frame, socket, connection) do
Socket.send(socket, Frame.serialize(frame, connection.remote_settings.max_frame_size))
end
end