Current section

Files

Jump to
membrane_rtc_engine_recording lib edge_timestamp_saver.ex
Raw

lib/edge_timestamp_saver.ex

defmodule Membrane.RTC.Engine.Endpoint.Recording.EdgeTimestampSaver do
@moduledoc false
use Membrane.Filter
alias Membrane.Buffer
alias Membrane.RTC.Engine.Endpoint.Recording
def_input_pad :input, accepted_format: _accepted_format
def_output_pad :output, accepted_format: _accepted_format
def_options reporter: [
spec: pid(),
description: """
Pid of the recording reporter
"""
]
@packets_interval 200
@impl true
def handle_init(_ctx, options) do
{[],
%{
first_buffer_timestamp: nil,
last_buffer_timestamp: nil,
counter: 0,
reporter: options.reporter
}}
end
@impl true
def handle_buffer(_pad, %Buffer{} = buffer, ctx, %{first_buffer_timestamp: nil} = state) do
Recording.Reporter.start_timestamp(
state.reporter,
track_id(ctx),
buffer.metadata.rtp.timestamp
)
{[buffer: {:output, buffer}],
%{state | first_buffer_timestamp: buffer.metadata.rtp.timestamp}}
end
@impl true
def handle_buffer(_pad, %Buffer{} = buffer, ctx, state) do
counter = rem(state.counter + 1, @packets_interval)
if counter == 0,
do:
Recording.Reporter.end_timestamp(
state.reporter,
track_id(ctx),
buffer.metadata.rtp.timestamp
)
{[buffer: {:output, buffer}],
%{state | counter: counter, last_buffer_timestamp: buffer.metadata.rtp.timestamp}}
end
@impl true
def handle_end_of_stream(_pad, ctx, state) do
Recording.Reporter.end_timestamp(state.reporter, track_id(ctx), state.last_buffer_timestamp)
{[end_of_stream: :output], state}
end
defp track_id(%{name: {:edge_timestamp_saver, track_id}}), do: track_id
end