diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 2d214d7..da5b9eb 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -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" diff --git a/examples/dynamic_subscriber.exs b/examples/dynamic_subscriber.exs index 02f6b3e..1e1169e 100644 --- a/examples/dynamic_subscriber.exs +++ b/examples/dynamic_subscriber.exs @@ -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]} @@ -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} @@ -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 @@ -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] + ) |> child({:parser, generation}, %Membrane.H265.Parser{ generate_best_effort_timestamps: %{framerate: framerate || {30, 1}}, output_stream_structure: :annexb diff --git a/examples/publish_and_play.exs b/examples/publish_and_play.exs index 561f12d..e20db62 100644 --- a/examples/publish_and_play.exs +++ b/examples/publish_and_play.exs @@ -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 diff --git a/lib/source.ex b/lib/source.ex index b14b18c..9d17d0c 100644 --- a/lib/source.ex +++ b/lib/source.ex @@ -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()` """ ] ] @@ -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, diff --git a/mix.exs b/mix.exs index 530d312..0251e81 100644 --- a/mix.exs +++ b/mix.exs @@ -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}, diff --git a/mix.lock b/mix.lock index ba10dad..bc7a8e8 100644 --- a/mix.lock +++ b/mix.lock @@ -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"}, diff --git a/test/integration_test.exs b/test/integration_test.exs index ecf656d..5d39ee2 100644 --- a/test/integration_test.exs +++ b/test/integration_test.exs @@ -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 @@ -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) ] ) @@ -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{ @@ -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 diff --git a/test/test_helper.exs b/test/test_helper.exs index 1462514..054a55e 100644 --- a/test/test_helper.exs +++ b/test/test_helper.exs @@ -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 + +""")