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
5 changes: 5 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -15,3 +15,8 @@ jobs:

test:
uses: membraneframework/membrane_actions/.github/workflows/test.yml@main
with:
setup-commands: |
VERSION=0.14.14
curl -fsSL "https://github.com/moq-dev/moq/releases/download/moq-relay-v${VERSION}/moq-relay-v${VERSION}-x86_64-unknown-linux-gnu.tar.gz" \
| tar -xz --strip-components=2 -C /usr/local/bin "moq-relay-v${VERSION}-x86_64-unknown-linux-gnu/bin/moq-relay"
13 changes: 9 additions & 4 deletions examples/dynamic_subscriber.exs
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,8 @@ defmodule Subscriber do
[:current_track, available_tracks: %{}, generation: 0, moq_disconnected?: false]
end

@subscription %ExMoQ.Subscription{latency_ns: 200_000_000}

@impl true
def handle_init(_ctx, opts) do
state = %State{url: opts[:url], broadcast: opts[:broadcast]}
Expand All @@ -57,8 +59,7 @@ defmodule Subscriber do
child(:source, %Membrane.MoQ.Source{
url: state.url,
broadcast: state.broadcast,
disable_tls_verify?: true,
latency: Membrane.Time.milliseconds(200)
disable_tls_verify?: true
})

{[spec: source_spec], state}
Expand Down Expand Up @@ -147,7 +148,9 @@ defmodule Subscriber do
defp track_spec(name, generation, %Membrane.H264{framerate: framerate}),
do:
get_child(:source)
|> via_out(Pad.ref(:output, generation), options: [track: name])
|> via_out(Pad.ref(:output, generation),
options: [track: name, subscription: @subscription]
)
|> child({:parser, generation}, %Membrane.H264.Parser{
generate_best_effort_timestamps: %{framerate: framerate || {30, 1}},
output_stream_structure: :annexb
Expand All @@ -159,7 +162,9 @@ defmodule Subscriber do
defp track_spec(name, generation, %Membrane.H265{framerate: framerate}),
do:
get_child(:source)
|> via_out(Pad.ref(:output, generation), options: [track: name])
|> via_out(Pad.ref(:output, generation),
options: [track: name, subscription: @subscription]

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Does it make sense, to require specifying both track name and subscription? Maybe sbd who has some expertise in MoQ could look at it too

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't think it's too important, the struct is just a bundle of parameters, we may just as well put :track inside it. It would be closer to moq-lite's structure of SUBSCRIBE, see https://datatracker.ietf.org/doc/draft-lcurley-moq-lite/ §7.7

for comparison, IETF MoQ splits the track name and parameters: https://www.ietf.org/archive/id/draft-ietf-moq-transport-21.html#section-9.6

It makes sense for me to have the current shape with the track name as the "most important" parameter, but it could just as well be a required value for the struct, your call

)
|> child({:parser, generation}, %Membrane.H265.Parser{
generate_best_effort_timestamps: %{framerate: framerate || {30, 1}},
output_stream_structure: :annexb
Expand Down
4 changes: 3 additions & 1 deletion examples/publish_and_play.exs
Original file line number Diff line number Diff line change
Expand Up @@ -75,7 +75,9 @@ defmodule Player do
broadcast: opts[:broadcast],
disable_tls_verify?: true
})
|> via_out(Pad.ref(:output, :video), options: [track: opts[:track]])
|> via_out(Pad.ref(:output, :video),
options: [track: opts[:track], subscription: %ExMoQ.Subscription{group_start: 0}]
)
# `MoQ.Source` emits frames in decode order carrying their presentation
# PTS but no DTS, so the decoder has no monotonic decode timeline to
# reorder against and would emit frames still in decode order — which the
Expand Down
15 changes: 7 additions & 8 deletions lib/source.ex
Original file line number Diff line number Diff line change
Expand Up @@ -52,13 +52,12 @@ defmodule Membrane.MoQ.Source do
see `Track` at https://doc.moq.dev/concept/layer/moq-lite.html#terminology
"""
],
priority: [
spec: 0..255 | nil,
default: nil,
subscription: [
spec: ExMoQ.Subscription.t(),
default: %ExMoQ.Subscription{},
description: """
Delivery priority of this subscription.
Under congestion, tracks with a higher value are sent first.
When nil, hang defaults for the track's media kind are used.
Parameters configuring the subscriber for this track.
For more info, see `ExMoQ.Subscription.t()`
"""
]
]
Expand Down Expand Up @@ -367,11 +366,11 @@ defmodule Membrane.MoQ.Source do
@spec subscribe_pad(Membrane.Pad.ref(), Membrane.Element.CallbackContext.t(), State.t()) ::
{[Membrane.Element.Action.t()], State.t()}
defp subscribe_pad(pad, ctx, state) do
%{track: track, priority: priority} = ctx.pads[pad].options
%{track: track, subscription: subscription} = ctx.pads[pad].options

with format when format != nil <- Catalog.rendition(state.catalog, track),
token = state.next_token,
:ok <- Native.subscribe_track(state.consumer, track, token, priority) do
:ok <- Native.subscribe_track(state.consumer, track, token, subscription) do
state = %{
state
| next_token: token + 1,
Expand Down
3 changes: 2 additions & 1 deletion mix.exs
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,8 @@ defmodule Membrane.MoQ.Mixfile do
{:membrane_h265_format, "~> 0.2.0"},
{:membrane_aac_format, "~> 0.8.0"},
{:membrane_opus_format, "~> 0.3.0"},
{:ex_moq, "~> 0.1.0"},
# TODO: switch to released ex_moq once it gets merged
{:ex_moq, github: "membraneframework/ex_moq", branch: "kidq330/add_sub_params"},
{:muontrap, "~> 1.8", only: :test},
{:membrane_aac_plugin, "~> 0.19", only: :test},
{:membrane_file_plugin, "~> 0.17", only: :test},
Expand Down
2 changes: 1 addition & 1 deletion mix.lock
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@
"elixir_make": {:hex, :elixir_make, "0.10.0", "16577e2583a79bb79237bbff349619ef5d80afffc07eac6e4faf0d00e2ddaf7d", [:mix], [], "hexpm", "dc1f09fb7fa68866b886abd5f0f3c83553b1a19a52359a899e92af1bb3b31982"},
"erlex": {:hex, :erlex, "0.2.9", "7debbbaa9f4f368b8cd648983e0f1d7963028508e9c59e9d4ed504e94ef52a55", [:mix], [], "hexpm", "8cfffc0ec7159e6d73de2ab28a588064de80f88b2798d5cbe4482cbbc200178b"},
"ex_doc": {:hex, :ex_doc, "0.40.4", "66f2e42bf588594d5a8aab31cad87f2ddad09d0da1b1a2f379340ec2c2e497cb", [:mix], [{:earmark_parser, "~> 1.4.46", [hex: :earmark_parser, repo: "hexpm", optional: false]}, {:makeup_c, ">= 0.1.0", [hex: :makeup_c, repo: "hexpm", optional: true]}, {:makeup_elixir, "~> 0.14 or ~> 1.0", [hex: :makeup_elixir, repo: "hexpm", optional: false]}, {:makeup_erlang, "~> 0.1 or ~> 1.0", [hex: :makeup_erlang, repo: "hexpm", optional: false]}, {:makeup_html, ">= 0.1.0", [hex: :makeup_html, repo: "hexpm", optional: true]}], "hexpm", "6222b9e423d76584ee34df2c82a5ed72c2d53dc153f7f483ad28b378694186cc"},
"ex_moq": {:hex, :ex_moq, "0.1.0", "5ad97ad6eb9d7a66e70bdad5b5deca9839e808fc2c1da21b7127b0351ce04cb8", [:mix], [{:muontrap, "~> 1.8", [hex: :muontrap, repo: "hexpm", optional: true]}, {:rustler, "~> 0.38.0", [hex: :rustler, repo: "hexpm", optional: false]}], "hexpm", "4bfbef0832aa858e2c49844ab0c3b0cdaa71b93acd8f8ad5e4b8c54dca1713bc"},
"ex_moq": {:git, "https://github.com/membraneframework/ex_moq.git", "7190e76f827cb8e40f459d3e0451431fdec21f53", [branch: "kidq330/add_sub_params"]},
"file_system": {:hex, :file_system, "1.1.1", "31864f4685b0148f25bd3fbef2b1228457c0c89024ad67f7a81a3ffbc0bbad3a", [:mix], [], "hexpm", "7a15ff97dfe526aeefb090a7a9d3d03aa907e100e262a0f8f7746b78f8f87a5d"},
"jason": {:hex, :jason, "1.4.5", "2e3a008590b0b8d7388c20293e9dcc9cf3e5d642fd2a114e4cbbb52e595d940a", [:mix], [{:decimal, "~> 1.0 or ~> 2.0 or ~> 3.0", [hex: :decimal, repo: "hexpm", optional: true]}], "hexpm", "b0c823996102bcd0239b3c2444eb00409b72f6a140c1950bc8b457d836b30684"},
"logger_backends": {:hex, :logger_backends, "1.0.1", "08a9bd09eed271eac41ff35dd7603b6abf9e0f8b6ffe40359846cb0bfd37704f", [:mix], [], "hexpm", "d7903244fe8d84eeb3d3d42f8d7de080fd0409b613040df54c78b3591672d652"},
Expand Down
16 changes: 12 additions & 4 deletions test/integration_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@ defmodule Membrane.MoQ.IntegrationTest do
@track "video"
@audio_track "audio"

@subscription %ExMoQ.Subscription{group_start: 0, latency_ns: 5_000_000_000}

defmodule EndOfStreamSource do
use Membrane.Source

Expand Down Expand Up @@ -194,10 +196,14 @@ defmodule Membrane.MoQ.IntegrationTest do
broadcast: broadcast,
disable_tls_verify?: relay.disable_tls_verify?
})
|> via_out(Pad.ref(:output, :video), options: [track: @track])
|> via_out(Pad.ref(:output, :video),
options: [track: @track, subscription: @subscription]
)
|> child(:video_sink, Testing.Sink),
get_child(:source)
|> via_out(Pad.ref(:output, :audio), options: [track: @audio_track])
|> via_out(Pad.ref(:output, :audio),
options: [track: @audio_track, subscription: @subscription]
)
|> child(:audio_sink, Testing.Sink)
]
)
Expand Down Expand Up @@ -397,7 +403,9 @@ defmodule Membrane.MoQ.IntegrationTest do
broadcast: broadcast,
disable_tls_verify?: relay.disable_tls_verify?
})
|> via_out(Pad.ref(:output, :video), options: [track: @track])
|> via_out(Pad.ref(:output, :video),
options: [track: @track, subscription: @subscription]
)
|> child(:tee, Membrane.Tee)
|> via_in(Pad.ref(:input, :video), options: [track: @track])
|> child(:moq_sink, %Membrane.MoQ.Sink{
Expand Down Expand Up @@ -457,7 +465,7 @@ defmodule Membrane.MoQ.IntegrationTest do
broadcast: broadcast,
disable_tls_verify?: relay.disable_tls_verify?
})
|> via_out(Pad.ref(:output, track), options: [track: track])
|> via_out(Pad.ref(:output, track), options: [track: track, subscription: @subscription])
|> child(:sink, Testing.Sink)
)
end
Expand Down
10 changes: 7 additions & 3 deletions test/test_helper.exs
Original file line number Diff line number Diff line change
@@ -1,5 +1,9 @@
ExUnit.start(capture_log: true)

# Integration tests talk to a real MoQ relay.
# Opt in with `mix test --include integration`.
ExUnit.configure(exclude: [:integration])
IO.puts("""
Some tests assume a moq-relay binary is available in your PATH.
To disable them, run:

mix test --exclude integration

""")
Loading