Packages
membrane_rtp_plugin
0.23.2
0.31.5
0.31.4
0.31.3
0.31.2
0.31.1
0.31.0
0.30.0
0.29.1
0.29.0
0.28.0
0.27.1
0.27.0
0.26.0
0.25.0
0.24.1
0.24.0
0.23.2
0.23.1
0.23.0
0.22.1
0.22.0
0.21.0
0.20.0
0.19.1
0.19.0
0.18.0
0.17.1
0.17.0
0.16.0
0.15.0
0.15.0-rc.1
0.14.0
0.13.0
0.12.2
0.12.1
0.12.0
0.11.0
0.10.0
0.9.0
0.8.2
0.8.1
0.8.0
0.7.1-alpha.3
0.7.1-alpha.2
0.7.0-alpha.2
0.7.0-alpha.1
0.7.0-alpha
0.6.1
0.6.0
0.5.1
0.5.0
0.4.0-alpha
Membrane Multimedia Framework plugin for RTP
Current section
Files
Jump to
Current section
Files
lib/membrane/rtp/twcc_sender/receiver_rate.ex
defmodule Membrane.RTP.TWCCSender.ReceiverRate do
@moduledoc false
# module responsible for calculating bitrate
# received by a receiver in last `window` time
# (referred to as R_hat in the GCC draft, sec. 5.5)
alias Membrane.Time
@type t() :: %__MODULE__{
# estimated bitrate in bps
value: float() | nil,
# time window for measuring the received bitrate, between [0.5, 1]s (reffered to as "T" in the draft)
window: Time.t(),
# accumulator for packets and their timestamps that have been received in last `window` time
packets_received: Qex.t({Time.t(), pos_integer()})
}
@enforce_keys [:window]
defstruct @enforce_keys ++ [:value, packets_received: Qex.new()]
@spec new(Time.t()) :: t()
def new(window), do: %__MODULE__{window: window}
@spec update(t(), Time.t(), [Time.t() | :not_received], [pos_integer()]) :: t()
def update(%__MODULE__{value: nil} = rr, reference_time, receive_deltas, packet_sizes) do
packets_received = resolve_receive_deltas(receive_deltas, reference_time, packet_sizes)
packets_received = Qex.join(rr.packets_received, packets_received)
{first_packet_timestamp, _first_packet_size} = Qex.first!(packets_received)
{last_packet_timestamp, _last_packet_size} = Qex.last!(packets_received)
if last_packet_timestamp - first_packet_timestamp >= rr.window do
do_update(rr, packets_received)
else
%__MODULE__{rr | packets_received: packets_received}
end
end
def update(%__MODULE__{} = rr, reference_time, receive_deltas, packet_sizes) do
packets_received = resolve_receive_deltas(receive_deltas, reference_time, packet_sizes)
do_update(rr, packets_received)
end
defp do_update(rr, packets_received) do
{last_packet_timestamp, _last_packet_size} = Qex.last!(packets_received)
threshold = last_packet_timestamp - rr.window
packets_received =
rr.packets_received
|> Qex.join(packets_received)
|> Enum.drop_while(fn {timestamp, _size} -> timestamp < threshold end)
|> Qex.new()
received_sizes_sum =
Enum.reduce(packets_received, 0, fn {_timestamp, size}, acc -> acc + size end)
value = 1 / (Time.as_milliseconds(rr.window) / 1000) * received_sizes_sum
%__MODULE__{rr | value: value, packets_received: packets_received}
end
defp resolve_receive_deltas(receive_deltas, reference_time, packet_sizes) do
receive_deltas
|> Enum.zip(packet_sizes)
|> Enum.reject(fn {delta, _size} -> delta == :not_received end)
|> Enum.map_reduce(reference_time, fn {recv_delta, size}, prev_timestamp ->
receive_timestamp = prev_timestamp + recv_delta
{{receive_timestamp, size}, receive_timestamp}
end)
# take the packets_received
|> elem(0)
|> Qex.new()
end
end