Packages
membrane_rtc_engine
0.10.3
0.25.0
0.24.0
0.23.0
0.22.0
0.21.0
0.20.0
0.19.0
0.18.0
0.17.1
0.17.0
0.16.0
0.15.1
0.15.0
0.14.2
0.14.1
0.14.0
0.13.0
0.12.1
0.12.0
0.11.0
0.10.3
0.10.2
0.10.1
0.10.0
0.9.1
0.9.0
0.8.2
0.8.1
0.8.0
0.7.0
0.6.0
0.5.1
0.5.0
0.4.1
0.4.0
0.3.2
0.3.1
0.3.0
0.2.0
0.1.0
0.1.0-alpha.2
0.1.0-alpha.1
0.1.0-alpha
Membrane RTC Engine and its client library
Current section
Files
Jump to
Current section
Files
lib/membrane_rtc_engine/endpoints/hls_endpoint.ex
if Enum.all?(
[
Membrane.H264.FFmpeg.Parser,
Membrane.HTTPAdaptiveStream.SinkBin,
Membrane.Opus.Decoder,
Membrane.AAC.Parser,
Membrane.AAC.FDK.Encoder
],
&Code.ensure_loaded?/1
) do
defmodule Membrane.RTC.Engine.Endpoint.HLS do
@moduledoc """
An Endpoint responsible for converting incoming tracks to HLS playlist.
This module requires the following plugins to be present in your `mix.exs` for H264 & OPUS input:
```
[
:membrane_h264_ffmpeg_plugin,
:membrane_http_adaptive_stream_plugin,
]
```
It can perform transcoding (see `Membrane.RTC.Engine.Endpoint.HLS.TranscodingConfig`),
in such case these plugins are also needed:
```
[
:membrane_ffmpeg_swscale_plugin,
:membrane_framerate_converter_plugin
]
```
Plus, optionally it supports OPUS audio input - for that, additional dependencies are needed:
```
[
:membrane_opus_plugin,
:membrane_aac_plugin,
:membrane_aac_fdk_plugin
]
```
"""
use Membrane.Bin
require Membrane.Logger
alias Membrane.RTC.Engine
alias Membrane.RTC.Engine.Endpoint.HLS.TranscodingConfig
alias Membrane.RTC.Engine.Endpoint.WebRTC.TrackReceiver
alias Membrane.RTC.Engine.Track
@transcoding_deps [
Membrane.H264.FFmpeg.Decoder,
Membrane.H264.FFmpeg.Encoder,
Membrane.FFmpeg.SWScale.Scaler,
Membrane.FramerateConverter
]
def_input_pad :input,
demand_unit: :buffers,
caps: :any,
availability: :on_request
def_options rtc_engine: [
spec: pid(),
description: "Pid of parent Engine"
],
output_directory: [
spec: Path.t(),
description: "Path to directory under which HLS output will be saved",
default: "hls_output"
],
owner: [
spec: pid(),
description: """
Pid of parent all notifications will be send to.
These notifications are:
* `{:playlist_playable, content_type, stream_id, origin}`
* `{:cleanup, clean_function, stream_id}`
"""
],
hls_mode: [
spec: :separate_av | :muxed_av,
default: :separate_av,
description: """
Defines output mode for `Membrane.HTTPAdaptiveStream.SinkBin`.
- `:separate_av` - audio and video tracks will be separated
- `:muxed_av` - audio will be attached to every video track
"""
],
target_window_duration: [
type: :time,
spec: Membrane.Time.t() | :infinity,
default: Membrane.Time.seconds(20),
description: """
Max duration of stream that will be stored. Segments that are older than window duration will be removed.
"""
],
target_segment_duration: [
type: :time,
spec: Membrane.Time.t(),
default: Membrane.Time.seconds(5),
description: """
Expected length of each segment. Setting it is not necessary, but
may help players achieve better UX.
"""
],
framerate: [
spec: {integer(), integer()} | nil,
description: """
Framerate of input tracks
""",
default: nil
],
transcoding_config: [
spec: TranscodingConfig.t(),
default: %TranscodingConfig{},
description: """
Transcoding configuration
"""
]
@impl true
def handle_init(opts) do
state = %{
rtc_engine: opts.rtc_engine,
tracks: %{},
stream_ids: MapSet.new(),
output_directory: opts.output_directory,
owner: opts.owner,
hls_mode: opts.hls_mode,
target_window_duration: opts.target_window_duration,
framerate: opts.framerate,
target_segment_duration: opts.target_segment_duration,
transcoding_config: opts.transcoding_config
}
{:ok, state}
end
@impl true
def handle_other({:new_tracks, tracks}, ctx, state) do
{:endpoint, endpoint_id} = ctx.name
state =
Enum.reduce(tracks, state, fn track, state ->
case Engine.subscribe(state.rtc_engine, endpoint_id, track.id) do
:ok ->
put_in(state, [:tracks, track.id], track)
{:error, :invalid_track_id} ->
Membrane.Logger.debug("""
Couldn't subscribe to the track: #{inspect(track.id)}. No such track.
It had to be removed just after publishing it. Ignoring.
""")
state
{:error, reason} ->
raise "Couldn't subscribe to the track: #{inspect(track.id)}. Reason: #{inspect(reason)}"
end
end)
{:ok, state}
end
@impl true
def handle_other(msg, _ctx, state) do
Membrane.Logger.debug("Unexpected message: #{inspect(msg)}. Ignoring.")
{:ok, state}
end
def handle_notification(
{:track_playable, {content_type, track_id}},
{:hls_sink_bin, stream_id},
_ctx,
state
) do
%{origin: origin} = Map.fetch!(state.tracks, track_id)
# notify about playable just when video becomes available
send(state.owner, {:playlist_playable, content_type, stream_id, origin})
{:ok, state}
end
def handle_notification(
{:cleanup, clean_function},
{:hls_sink_bin, stream_id},
_ctx,
state
) do
# notify about possibility to cleanup as the stream is finished.
send(state.owner, {:cleanup, clean_function, stream_id})
{:ok, state}
end
@impl true
def handle_notification(notification, _element, _context, state) do
Membrane.Logger.warn("Unexpected notification: #{inspect(notification)}. Ignoring.")
{:ok, state}
end
@impl true
def handle_pad_removed(Pad.ref(:input, track_id), ctx, state) do
children =
[
:opus_decoder,
:aac_encoder,
:aac_parser,
:video_parser,
:video_parser_out,
:decoder,
:encoder,
:resolution_scaler,
:framerate_converter
]
|> Enum.map(&{&1, track_id})
|> Enum.filter(&Map.has_key?(ctx.children, &1))
{removed_track, tracks} = Map.pop!(state.tracks, track_id)
state = %{state | tracks: tracks}
sink_bin_used? =
Enum.any?(tracks, fn {_id, track} ->
track.stream_id == removed_track.stream_id
end)
children =
if sink_bin_used?,
do: children,
else: [{:hls_sink_bin, removed_track.stream_id} | children]
{{:ok, [remove_child: children]}, state}
end
@impl true
def handle_pad_added(Pad.ref(:input, track_id) = pad, _ctx, state) do
link_builder = link_bin_input(pad)
track = Map.get(state.tracks, track_id)
directory = Path.join(state.output_directory, track.stream_id)
spec = hls_links_and_children(link_builder, track, state)
{spec, state} =
if MapSet.member?(state.stream_ids, track.stream_id) do
{spec, state}
else
# remove directory if it already exists
File.rm_rf(directory)
File.mkdir_p!(directory)
hls_sink_bin = %Membrane.HTTPAdaptiveStream.SinkBin{
manifest_module: Membrane.HTTPAdaptiveStream.HLS,
target_window_duration: state.target_window_duration,
target_segment_duration: state.target_segment_duration,
persist?: false,
storage: %Membrane.HTTPAdaptiveStream.Storages.FileStorage{
directory: directory
},
hls_mode: state.hls_mode
}
new_spec = %{
spec
| children: Map.put(spec.children, {:hls_sink_bin, track.stream_id}, hls_sink_bin)
}
{new_spec, %{state | stream_ids: MapSet.put(state.stream_ids, track.stream_id)}}
end
{{:ok, spec: spec}, state}
end
defp hls_links_and_children(link_builder, %Track{encoding: :OPUS} = track, _state) do
%ParentSpec{
children: %{
{:track_receiver, track.id} => %TrackReceiver{
track: track,
initial_target_variant: :high
},
{:depayloader, track.id} => get_depayloader(track),
{:opus_decoder, track.id} => Membrane.Opus.Decoder,
{:aac_encoder, track.id} => Membrane.AAC.FDK.Encoder,
{:aac_parser, track.id} => %Membrane.AAC.Parser{out_encapsulation: :none}
},
links: [
link_builder
|> to({:track_receiver, track.id})
|> to({:depayloader, track.id})
|> to({:opus_decoder, track.id})
|> to({:aac_encoder, track.id})
|> to({:aac_parser, track.id})
|> via_in(Pad.ref(:input, {:audio, track.id}), options: [encoding: :AAC])
|> to({:hls_sink_bin, track.stream_id})
]
}
end
defp hls_links_and_children(link_builder, %Track{encoding: :H264} = track, state) do
link_to_transcoder = create_transcoder_link(state.transcoding_config, track.id)
%ParentSpec{
children: %{
{:track_receiver, track.id} => %TrackReceiver{
track: track,
initial_target_variant: :high,
keyframe_request_interval: state.target_segment_duration
},
{:depayloader, track.id} => get_depayloader(track),
{:video_parser, track.id} => %Membrane.H264.FFmpeg.Parser{
alignment: :au,
attach_nalus?: true,
framerate: state.framerate
}
},
links: [
link_builder
|> to({:track_receiver, track.id})
|> to({:depayloader, track.id})
|> to({:video_parser, track.id})
|> then(link_to_transcoder)
|> via_in(Pad.ref(:input, {:video, track.id}), options: [encoding: :H264])
|> to({:hls_sink_bin, track.stream_id})
]
}
end
defp get_depayloader(track) do
track
|> Track.get_depayloader()
|> tap(&unless &1, do: raise("Couldn't find depayloader for track #{inspect(track)}"))
end
defp create_transcoder_link(%TranscodingConfig{enabled?: false}, _track_id), do: & &1
if Enum.all?(@transcoding_deps, &Code.ensure_loaded?/1) do
defp create_transcoder_link(transcoding_config, track_id) do
resolution_scaler = %Membrane.FFmpeg.SWScale.Scaler{
output_width: transcoding_config.output_width,
output_height: transcoding_config.output_height
}
framerate_converter = %Membrane.FramerateConverter{
framerate: transcoding_config.output_framerate
}
video_parser_out = %Membrane.H264.FFmpeg.Parser{
alignment: :au,
attach_nalus?: true,
framerate: transcoding_config.output_framerate
}
fn link_builder ->
link_builder
|> to({:decoder, track_id}, Membrane.H264.FFmpeg.Decoder)
|> to({:resolution_scaler, track_id}, resolution_scaler)
|> to({:framerate_converter, track_id}, framerate_converter)
|> to({:encoder, track_id}, Membrane.H264.FFmpeg.Encoder)
|> to({:video_parser_out, track_id}, video_parser_out)
end
end
else
defp create_transcoder_link(_transcoding_config, _track_id) do
raise """
Cannot find some of the modules required to perform transcoding.
Ensure `:membrane_ffmpeg_swscale_plugin` and `membrane_framerate_converter_plugin` are added to the deps.
"""
end
end
end
end