From 1b943637a4b5778bb929a7af7687b048fecaa03f Mon Sep 17 00:00:00 2001 From: Alessandro Ros Date: Wed, 5 Aug 2026 12:16:07 +0200 Subject: [PATCH] moq: support draft-16 (#6045) --- api/openapi.yaml | 1 + docs/3-publish/01-moq-clients.md | 2 +- internal/defs/api_moq.go | 1 + .../moq/controlmessage/client_setup.go | 23 ++ .../protocols/moq/controlmessage/message.go | 4 + .../moq/controlmessage/message_test.go | 25 +- .../moq/controlmessage/server_setup.go | 23 ++ .../protocols/moq/controlmessage/setup.go | 14 +- internal/servers/moq/http_server.go | 1 + internal/servers/moq/native_listener.go | 4 + internal/servers/moq/server_test.go | 308 +++++++++++------- internal/servers/moq/session.go | 129 +++++--- 12 files changed, 352 insertions(+), 183 deletions(-) create mode 100644 internal/protocols/moq/controlmessage/client_setup.go create mode 100644 internal/protocols/moq/controlmessage/server_setup.go diff --git a/api/openapi.yaml b/api/openapi.yaml index 8e882a57..0e054326 100644 --- a/api/openapi.yaml +++ b/api/openapi.yaml @@ -98,6 +98,7 @@ components: MoQVersion: type: string enum: + - moqt-16 - moqt-17 - moqt-18 - moqt-19 diff --git a/docs/3-publish/01-moq-clients.md b/docs/3-publish/01-moq-clients.md index 6e4ccc5e..c630c5cb 100644 --- a/docs/3-publish/01-moq-clients.md +++ b/docs/3-publish/01-moq-clients.md @@ -9,7 +9,7 @@ Media-over-QUIC is a streaming protocol built upon cutting edge protocols (QUIC, Media-over-QUIC has a wide range of features and variants, most of them in active development. We currently support the following: -- The server supports `draft-19`, `draft-18` and `draft-17` of the [main specification](https://datatracker.ietf.org/doc/html/draft-ietf-moq-transport-19), and prefers `draft-19` when multiple versions are offered during negotiation. +- The server supports `draft-19`, `draft-18`, `draft-17` and `draft-16` of the [main specification](https://datatracker.ietf.org/doc/html/draft-ietf-moq-transport-19), and prefers `draft-19` when multiple versions are offered during negotiation. - We support using Media-over-QUIC through browsers with the WebTransport API and through native QUIC clients. - We support the `PUBLISH` and `SUBSCRIBE` messages only, which are the ones meant to be used with a routing solution like _MediaMTX_. - We use the MOQT Streaming Format (MSF) to advertise tracks, described in [this specification](https://datatracker.ietf.org/doc/html/draft-ietf-moq-msf-00). diff --git a/internal/defs/api_moq.go b/internal/defs/api_moq.go index d229e648..9817ade9 100644 --- a/internal/defs/api_moq.go +++ b/internal/defs/api_moq.go @@ -28,6 +28,7 @@ type APIMoQVersion string // protocol versions. const ( + APIMoQVersionDraft16 APIMoQVersion = "moqt-16" APIMoQVersionDraft17 APIMoQVersion = "moqt-17" APIMoQVersionDraft18 APIMoQVersion = "moqt-18" APIMoQVersionDraft19 APIMoQVersion = "moqt-19" diff --git a/internal/protocols/moq/controlmessage/client_setup.go b/internal/protocols/moq/controlmessage/client_setup.go new file mode 100644 index 00000000..9e89164a --- /dev/null +++ b/internal/protocols/moq/controlmessage/client_setup.go @@ -0,0 +1,23 @@ +package controlmessage + +import "github.com/bluenviron/mediamtx/internal/protocols/moq/varint" + +const typeClientSetup varint.Varint = 0x20 + +// ClientSetup is the CLIENT_SETUP control message. +// spec: draft-16, section 9.3 +type ClientSetup Setup + +func (*ClientSetup) isMessage() {} + +func (m *ClientSetup) unmarshal(buf []byte) error { + return (*Setup)(m).unmarshal(buf) +} + +// Marshal implements Message. +func (m ClientSetup) Marshal() []byte { + s := Setup(m) + buf := make([]byte, s.marshalSize(typeClientSetup)) + s.marshalTo(buf, typeClientSetup) + return buf +} diff --git a/internal/protocols/moq/controlmessage/message.go b/internal/protocols/moq/controlmessage/message.go index b66dc215..a9bb6ab8 100644 --- a/internal/protocols/moq/controlmessage/message.go +++ b/internal/protocols/moq/controlmessage/message.go @@ -42,6 +42,10 @@ func Read(r io.Reader) (Message, error) { switch t { case typeSetup: m = &Setup{} + case typeClientSetup: + m = &ClientSetup{} + case typeServerSetup: + m = &ServerSetup{} case typeSubscribe: m = &Subscribe{} case typeSubscribeOk: diff --git a/internal/protocols/moq/controlmessage/message_test.go b/internal/protocols/moq/controlmessage/message_test.go index 0590cae8..f1125117 100644 --- a/internal/protocols/moq/controlmessage/message_test.go +++ b/internal/protocols/moq/controlmessage/message_test.go @@ -7,6 +7,7 @@ import ( "github.com/stretchr/testify/require" "github.com/bluenviron/mediamtx/internal/protocols/moq/controlmessage" + "github.com/bluenviron/mediamtx/internal/protocols/moq/namespace" "github.com/bluenviron/mediamtx/internal/protocols/moq/parameter" ) @@ -15,6 +16,22 @@ var cases = []struct { enc []byte dec controlmessage.Message }{ + { + name: "draft-16 client setup", + enc: []byte{ + 0x20, // type 0x20 + 0x00, 0x00, // length = 0 + }, + dec: &controlmessage.ClientSetup{}, + }, + { + name: "draft-16 server setup", + enc: []byte{ + 0x21, // type 0x21 + 0x00, 0x00, // length = 0 + }, + dec: &controlmessage.ServerSetup{}, + }, { name: "setup", enc: []byte{ @@ -66,7 +83,7 @@ var cases = []struct { }, dec: &controlmessage.Subscribe{ RequestID: 1, - Namespace: []string{"foo"}, + Namespace: namespace.Namespace{"foo"}, TrackName: "bar", }, }, @@ -88,7 +105,7 @@ var cases = []struct { }, dec: &controlmessage.Subscribe{ RequestID: 1, - Namespace: []string{"foo"}, + Namespace: namespace.Namespace{"foo"}, TrackName: "bar", Parameters: []parameter.Parameter{ ¶meter.AuthorizationToken{ @@ -157,7 +174,7 @@ var cases = []struct { }, dec: &controlmessage.Publish{ RequestID: 1, - Namespace: []string{"foo"}, + Namespace: namespace.Namespace{"foo"}, TrackName: "bar", TrackAlias: 2, }, @@ -181,7 +198,7 @@ var cases = []struct { }, dec: &controlmessage.Publish{ RequestID: 1, - Namespace: []string{"foo"}, + Namespace: namespace.Namespace{"foo"}, TrackName: "bar", TrackAlias: 2, Parameters: []parameter.Parameter{ diff --git a/internal/protocols/moq/controlmessage/server_setup.go b/internal/protocols/moq/controlmessage/server_setup.go new file mode 100644 index 00000000..66740346 --- /dev/null +++ b/internal/protocols/moq/controlmessage/server_setup.go @@ -0,0 +1,23 @@ +package controlmessage + +import "github.com/bluenviron/mediamtx/internal/protocols/moq/varint" + +const typeServerSetup varint.Varint = 0x21 + +// ServerSetup is the SERVER_SETUP control message. +// spec: draft-16, section 9.3 +type ServerSetup Setup + +func (*ServerSetup) isMessage() {} + +func (m *ServerSetup) unmarshal(buf []byte) error { + return (*Setup)(m).unmarshal(buf) +} + +// Marshal implements Message. +func (m ServerSetup) Marshal() []byte { + s := Setup(m) + buf := make([]byte, s.marshalSize(typeServerSetup)) + s.marshalTo(buf, typeServerSetup) + return buf +} diff --git a/internal/protocols/moq/controlmessage/setup.go b/internal/protocols/moq/controlmessage/setup.go index dc33a95a..90614d9f 100644 --- a/internal/protocols/moq/controlmessage/setup.go +++ b/internal/protocols/moq/controlmessage/setup.go @@ -74,7 +74,7 @@ func (m *Setup) unmarshal(buf []byte) error { return nil } -func (m Setup) marshalSize() int { +func (m Setup) marshalSize(t varint.Varint) int { payloadSize := 0 var previousType varint.Varint @@ -93,13 +93,13 @@ func (m Setup) marshalSize() int { len(m.Authority) } - return typeSetup.MarshalSize() + 2 + payloadSize + return t.MarshalSize() + 2 + payloadSize } -func (m Setup) marshalTo(buf []byte) int { - payloadSize := m.marshalSize() - typeSetup.MarshalSize() - 2 +func (m Setup) marshalTo(buf []byte, t varint.Varint) int { + payloadSize := m.marshalSize(t) - t.MarshalSize() - 2 - pos := typeSetup.MarshalTo(buf) + pos := t.MarshalTo(buf) buf[pos] = byte(payloadSize >> 8) buf[pos+1] = byte(payloadSize) pos += 2 @@ -126,7 +126,7 @@ func (m Setup) marshalTo(buf []byte) int { // Marshal implements Message. func (m Setup) Marshal() []byte { - buf := make([]byte, m.marshalSize()) - m.marshalTo(buf) + buf := make([]byte, m.marshalSize(typeSetup)) + m.marshalTo(buf, typeSetup) return buf } diff --git a/internal/servers/moq/http_server.go b/internal/servers/moq/http_server.go index a5cea1e7..f7a7f685 100644 --- a/internal/servers/moq/http_server.go +++ b/internal/servers/moq/http_server.go @@ -46,6 +46,7 @@ var supportedMoqtVersions = []defs.APIMoQVersion{ defs.APIMoQVersionDraft19, defs.APIMoQVersionDraft18, defs.APIMoQVersionDraft17, + defs.APIMoQVersionDraft16, } type ginUnwrapper interface { diff --git a/internal/servers/moq/native_listener.go b/internal/servers/moq/native_listener.go index 3f20849c..b4c20b86 100644 --- a/internal/servers/moq/native_listener.go +++ b/internal/servers/moq/native_listener.go @@ -19,6 +19,7 @@ var supportedMoqtALPNs = []string{ string(defs.APIMoQVersionDraft19), string(defs.APIMoQVersionDraft18), string(defs.APIMoQVersionDraft17), + string(defs.APIMoQVersionDraft16), } type nativeListenerParent interface { @@ -101,6 +102,9 @@ func alpnToVersion(alpn string) defs.APIMoQVersion { case string(defs.APIMoQVersionDraft17): return defs.APIMoQVersionDraft17 + + case string(defs.APIMoQVersionDraft16): + return defs.APIMoQVersionDraft16 } return "" diff --git a/internal/servers/moq/server_test.go b/internal/servers/moq/server_test.go index c41d3d35..da58ca04 100644 --- a/internal/servers/moq/server_test.go +++ b/internal/servers/moq/server_test.go @@ -37,6 +37,76 @@ func (p *serverDummyPath) ExternalCmdEnv() externalcmd.Environment { retur func (p *serverDummyPath) RemovePublisher(_ defs.PathRemovePublisherReq) {} func (p *serverDummyPath) RemoveReader(_ defs.PathRemoveReaderReq) {} +func performWTSetup(ctx context.Context, t *testing.T, sx *webtransport.Session, version defs.APIMoQVersion) { + t.Helper() + + if version == defs.APIMoQVersionDraft16 { + setupBidi, err := sx.OpenStreamSync(ctx) + require.NoError(t, err) + + _, err = setupBidi.Write(controlmessage.ClientSetup(controlmessage.Setup{}).Marshal()) + require.NoError(t, err) + + setupMsg, err := controlmessage.Read(setupBidi) + require.NoError(t, err) + require.Equal(t, &controlmessage.ServerSetup{}, setupMsg) + + setupBidi.Close() //nolint:errcheck + return + } + + setupStream, err := sx.AcceptUniStream(ctx) + require.NoError(t, err) + + setupMsg, err := controlmessage.Read(setupStream) + require.NoError(t, err) + require.Equal(t, &controlmessage.Setup{}, setupMsg) + + clientSetup, err := sx.OpenUniStreamSync(ctx) + require.NoError(t, err) + + _, err = clientSetup.Write(controlmessage.Setup{}.Marshal()) + require.NoError(t, err) +} + +func performNativeQUICSetup( + ctx context.Context, + t *testing.T, + conn *quic.Conn, + version defs.APIMoQVersion, + path string, +) { + t.Helper() + + if version == defs.APIMoQVersionDraft16 { + setupBidi, err := conn.OpenStreamSync(ctx) + require.NoError(t, err) + + _, err = setupBidi.Write(controlmessage.ClientSetup(controlmessage.Setup{Path: path}).Marshal()) + require.NoError(t, err) + + setupMsg, err := controlmessage.Read(setupBidi) + require.NoError(t, err) + require.Equal(t, &controlmessage.ServerSetup{}, setupMsg) + + setupBidi.Close() //nolint:errcheck + return + } + + setupStream, err := conn.AcceptUniStream(ctx) + require.NoError(t, err) + + setupMsg, err := controlmessage.Read(setupStream) + require.NoError(t, err) + require.Equal(t, &controlmessage.Setup{}, setupMsg) + + clientSetup, err := conn.OpenUniStreamSync(ctx) + require.NoError(t, err) + + _, err = clientSetup.Write(controlmessage.Setup{Path: path}.Marshal()) + require.NoError(t, err) +} + func TestAuthError(t *testing.T) { serverCertFile := test.CreateTempFile(t, test.TLSCertPub) serverKeyFile := test.CreateTempFile(t, test.TLSCertKey) @@ -166,20 +236,7 @@ func TestAuthError(t *testing.T) { defer sx.CloseWithError(0, "") //nolint:errcheck defer res.Body.Close() //nolint:errcheck - setupStream, err := sx.AcceptUniStream(ctx) - require.NoError(t, err) - - setupMsg, err := controlmessage.Read(setupStream) - require.NoError(t, err) - - _, ok := setupMsg.(*controlmessage.Setup) - require.True(t, ok) - - clientSetup, err := sx.OpenUniStreamSync(ctx) - require.NoError(t, err) - - _, err = clientSetup.Write(controlmessage.Setup{}.Marshal()) - require.NoError(t, err) + performWTSetup(ctx, t, sx, defs.APIMoQVersionDraft19) switch ca { case "subscribe": @@ -272,6 +329,11 @@ func TestServer(t *testing.T) { clientProtocols []string expectedVersion defs.APIMoQVersion }{ + { + name: "draft-16", + clientProtocols: []string{"moqt-16"}, + expectedVersion: defs.APIMoQVersionDraft16, + }, { name: "draft-17", clientProtocols: []string{"moqt-17"}, @@ -294,7 +356,7 @@ func TestServer(t *testing.T) { }, { name: "highest-preferred", - clientProtocols: []string{"moqt-17", "moqt-18", "moqt-19"}, + clientProtocols: []string{"moqt-16", "moqt-17", "moqt-18", "moqt-19"}, expectedVersion: defs.APIMoQVersionDraft19, }, } { @@ -368,18 +430,7 @@ func TestServer(t *testing.T) { require.Equal(t, `"`+string(ca.expectedVersion)+`"`, res.Header.Get("WT-Protocol")) - setupStream, err := sx.AcceptUniStream(ctx) - require.NoError(t, err) - - setupMsg, err := controlmessage.Read(setupStream) - require.NoError(t, err) - require.Equal(t, &controlmessage.Setup{}, setupMsg) - - clientSetup, err := sx.OpenUniStreamSync(ctx) - require.NoError(t, err) - - _, err = clientSetup.Write(controlmessage.Setup{}.Marshal()) - require.NoError(t, err) + performWTSetup(ctx, t, sx, ca.expectedVersion) catalogBidi, err := sx.OpenStreamSync(ctx) require.NoError(t, err) @@ -541,106 +592,111 @@ func TestServerUnsupportedVersion(t *testing.T) { } func TestServerNativeQUICSubscribe(t *testing.T) { - desc := &description.Session{Medias: []*description.Media{test.UniqueMediaH264()}} - strm := &stream.Stream{ - OrigDesc: desc, - WriteQueueSize: 512, - RTPMaxPayloadSize: 1450, - Parent: test.NilLogger, - } - err := strm.Initialize() - require.NoError(t, err) - defer strm.Close() - - pm := &test.PathManager{ - FindPathConfImpl: func(_ defs.PathFindPathConfReq) (*defs.PathFindPathConfRes, error) { - return &defs.PathFindPathConfRes{Conf: &conf.Path{}}, nil + for _, ca := range []struct { + name string + version defs.APIMoQVersion + }{ + { + name: "draft-16", + version: defs.APIMoQVersionDraft16, }, - AddReaderImpl: func(_ defs.PathAddReaderReq) (*defs.PathAddReaderRes, error) { - return &defs.PathAddReaderRes{Path: &serverDummyPath{}, Stream: strm}, nil + { + name: "draft-19", + version: defs.APIMoQVersionDraft19, }, + } { + t.Run(ca.name, func(t *testing.T) { + desc := &description.Session{Medias: []*description.Media{test.UniqueMediaH264()}} + strm := &stream.Stream{ + OrigDesc: desc, + WriteQueueSize: 512, + RTPMaxPayloadSize: 1450, + Parent: test.NilLogger, + } + err := strm.Initialize() + require.NoError(t, err) + defer strm.Close() + + pm := &test.PathManager{ + FindPathConfImpl: func(_ defs.PathFindPathConfReq) (*defs.PathFindPathConfRes, error) { + return &defs.PathFindPathConfRes{Conf: &conf.Path{}}, nil + }, + AddReaderImpl: func(_ defs.PathAddReaderReq) (*defs.PathAddReaderRes, error) { + return &defs.PathAddReaderRes{Path: &serverDummyPath{}, Stream: strm}, nil + }, + } + + serverCertFile := test.CreateTempFile(t, test.TLSCertPub) + serverKeyFile := test.CreateTempFile(t, test.TLSCertKey) + + s := &moq.Server{ + HTTP2Address: "127.0.0.1:19895", + HTTP3Address: "127.0.0.1:19896", + QUICAddress: "127.0.0.1:19897", + ServerCert: serverCertFile, + ServerKey: serverKeyFile, + AllowOrigins: []string{"*"}, + TrustedProxies: conf.IPNetworks{}, + ReadTimeout: conf.Duration(10 * time.Second), + WriteTimeout: conf.Duration(10 * time.Second), + PathManager: pm, + Parent: test.NilLogger, + } + err = s.Initialize() + require.NoError(t, err) + defer s.Close() + + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + + conn, err := quic.DialAddr(ctx, "127.0.0.1:19897", &tls.Config{ //nolint:gosec + InsecureSkipVerify: true, + NextProtos: []string{string(ca.version)}, + }, &quic.Config{EnableDatagrams: true}) + require.NoError(t, err) + defer conn.CloseWithError(0, "") //nolint:errcheck + + performNativeQUICSetup(ctx, t, conn, ca.version, "/teststream") + + catalogBidi, err := conn.OpenStreamSync(ctx) + require.NoError(t, err) + + _, err = catalogBidi.Write(controlmessage.Subscribe{ + RequestID: 1, + TrackName: ".catalog", + }.Marshal()) + require.NoError(t, err) + + catalogOkMsg, err := controlmessage.Read(catalogBidi) + require.NoError(t, err) + require.Equal(t, &controlmessage.SubscribeOk{TrackAlias: 1}, catalogOkMsg) + + sessions, err := s.APISessionsList() + require.NoError(t, err) + require.Equal(t, 1, len(sessions.Items)) + require.Equal(t, ca.version, sessions.Items[0].Version) + require.Equal(t, defs.APIMoQSessionTransportQUIC, sessions.Items[0].Transport) + + catalogDataStream, err := conn.AcceptUniStream(ctx) + require.NoError(t, err) + + var catalogSG subgroup.SubGroup + err = catalogSG.Read(catalogDataStream) + require.NoError(t, err) + + var cat catalog.Catalog + err = json.Unmarshal(catalogSG.Objects[0].Payload, &cat) + require.NoError(t, err) + + require.Equal(t, catalog.Catalog{ + Version: 1, + Tracks: []catalog.Track{{ + Name: "0", + Packaging: "loc", + IsLive: true, + Codec: "avc3.640028", + }}, + }, cat) + }) } - - serverCertFile := test.CreateTempFile(t, test.TLSCertPub) - serverKeyFile := test.CreateTempFile(t, test.TLSCertKey) - - s := &moq.Server{ - HTTP2Address: "127.0.0.1:19895", - HTTP3Address: "127.0.0.1:19896", - QUICAddress: "127.0.0.1:19897", - ServerCert: serverCertFile, - ServerKey: serverKeyFile, - AllowOrigins: []string{"*"}, - TrustedProxies: conf.IPNetworks{}, - ReadTimeout: conf.Duration(10 * time.Second), - WriteTimeout: conf.Duration(10 * time.Second), - PathManager: pm, - Parent: test.NilLogger, - } - err = s.Initialize() - require.NoError(t, err) - defer s.Close() - - ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) - defer cancel() - - conn, err := quic.DialAddr(ctx, "127.0.0.1:19897", &tls.Config{ //nolint:gosec - InsecureSkipVerify: true, - NextProtos: []string{string(defs.APIMoQVersionDraft19)}, - }, &quic.Config{EnableDatagrams: true}) - require.NoError(t, err) - defer conn.CloseWithError(0, "") //nolint:errcheck - - setupStream, err := conn.AcceptUniStream(ctx) - require.NoError(t, err) - - setupMsg, err := controlmessage.Read(setupStream) - require.NoError(t, err) - require.Equal(t, &controlmessage.Setup{}, setupMsg) - - clientSetup, err := conn.OpenUniStreamSync(ctx) - require.NoError(t, err) - - _, err = clientSetup.Write(controlmessage.Setup{Path: "/teststream"}.Marshal()) - require.NoError(t, err) - - catalogBidi, err := conn.OpenStreamSync(ctx) - require.NoError(t, err) - - _, err = catalogBidi.Write(controlmessage.Subscribe{ - RequestID: 1, - TrackName: ".catalog", - }.Marshal()) - require.NoError(t, err) - - catalogOkMsg, err := controlmessage.Read(catalogBidi) - require.NoError(t, err) - require.Equal(t, &controlmessage.SubscribeOk{TrackAlias: 1}, catalogOkMsg) - - sessions, err := s.APISessionsList() - require.NoError(t, err) - require.Equal(t, 1, len(sessions.Items)) - require.Equal(t, defs.APIMoQVersionDraft19, sessions.Items[0].Version) - require.Equal(t, defs.APIMoQSessionTransportQUIC, sessions.Items[0].Transport) - - catalogDataStream, err := conn.AcceptUniStream(ctx) - require.NoError(t, err) - - var catalogSG subgroup.SubGroup - err = catalogSG.Read(catalogDataStream) - require.NoError(t, err) - - var cat catalog.Catalog - err = json.Unmarshal(catalogSG.Objects[0].Payload, &cat) - require.NoError(t, err) - - require.Equal(t, catalog.Catalog{ - Version: 1, - Tracks: []catalog.Track{{ - Name: "0", - Packaging: "loc", - IsLive: true, - Codec: "avc3.640028", - }}, - }, cat) } diff --git a/internal/servers/moq/session.go b/internal/servers/moq/session.go index c3620283..fd4cd955 100644 --- a/internal/servers/moq/session.go +++ b/internal/servers/moq/session.go @@ -175,9 +175,11 @@ func (s *session) runInner() error { return s.runBidiStreamAcceptor(errGroup) }) - errGroup.Go(func() error { - return s.runSetupWriter() - }) + if s.version != defs.APIMoQVersionDraft16 { + errGroup.Go(func() error { + return s.runSetupWriter() + }) + } select { case <-s.ctx.Done(): @@ -254,47 +256,11 @@ func (s *session) onUniMessage(r io.Reader) error { switch m := msg.(type) { case *controlmessage.Setup: - err = func() error { - s.mutex.Lock() - defer s.mutex.Unlock() + if s.version == defs.APIMoQVersionDraft16 { + return fmt.Errorf("received SETUP over unidirectional stream with draft-16") + } - if s.transport == defs.APIMoQSessionTransportWebTransport { - if m.Path != "" { - return fmt.Errorf("received PATH setup option over WebTransport") - } - if m.Authority != "" { - return fmt.Errorf("received AUTHORITY setup option over WebTransport") - } - } - - if s.transport == defs.APIMoQSessionTransportQUIC && s.pathName == "" { - pathWithQuery := m.Path - if pathWithQuery == "" { - return fmt.Errorf("missing PATH setup option") - } - - u, err2 := url.ParseRequestURI(pathWithQuery) - if err2 != nil { - return fmt.Errorf("invalid PATH setup option: %w", err2) - } - - pathName := strings.Trim(u.Path, "/") - if pathName == "" { - return fmt.Errorf("invalid PATH setup option: empty path") - } - - s.pathName = pathName - s.query = u.RawQuery - } - - select { - case <-s.setupReceived: - return fmt.Errorf("SETUP stream is already present") - default: - close(s.setupReceived) - return nil - } - }() + err = s.processSetupMessage(m) if err != nil { return err } @@ -307,7 +273,80 @@ func (s *session) onUniMessage(r io.Reader) error { } } +func (s *session) processSetupMessage(m *controlmessage.Setup) error { + s.mutex.Lock() + defer s.mutex.Unlock() + + if s.transport == defs.APIMoQSessionTransportWebTransport { + if m.Path != "" { + return fmt.Errorf("received PATH setup option over WebTransport") + } + if m.Authority != "" { + return fmt.Errorf("received AUTHORITY setup option over WebTransport") + } + } + + if s.transport == defs.APIMoQSessionTransportQUIC && s.pathName == "" { + pathWithQuery := m.Path + if pathWithQuery == "" { + return fmt.Errorf("missing PATH setup option") + } + + u, err := url.ParseRequestURI(pathWithQuery) + if err != nil { + return fmt.Errorf("invalid PATH setup option: %w", err) + } + + pathName := strings.Trim(u.Path, "/") + if pathName == "" { + return fmt.Errorf("invalid PATH setup option: empty path") + } + + s.pathName = pathName + s.query = u.RawQuery + } + + select { + case <-s.setupReceived: + return fmt.Errorf("SETUP stream is already present") + default: + close(s.setupReceived) + return nil + } +} + func (s *session) runBidiStream(wstream io.ReadWriteCloser) error { + if s.version == defs.APIMoQVersionDraft16 { + select { + case <-s.setupReceived: + default: + msg, err := controlmessage.Read(wstream) + if err != nil { + return err + } + + setupMsg, ok := msg.(*controlmessage.ClientSetup) + if !ok { + return fmt.Errorf("expected CLIENT_SETUP as first message on draft-16 bidirectional stream") + } + + err = s.processSetupMessage((*controlmessage.Setup)(setupMsg)) + if err != nil { + return err + } + + buf := controlmessage.ServerSetup(controlmessage.Setup{}).Marshal() + + _, err = wstream.Write(buf) + if err != nil { + return err + } + + _, err = io.Copy(io.Discard, wstream) + return err + } + } + select { case <-s.setupReceived: case <-s.ctx.Done(): @@ -622,7 +661,7 @@ func (s *session) onPublishCatalog(wstream io.ReadWriteCloser, m *controlmessage } var ackPayload []byte - if s.version == defs.APIMoQVersionDraft17 { + if s.version == defs.APIMoQVersionDraft16 || s.version == defs.APIMoQVersionDraft17 { ackPayload = controlmessage.PublishOk{}.Marshal() } else { ackPayload = controlmessage.RequestOk{}.Marshal() @@ -646,7 +685,7 @@ func (s *session) onPublishTrack(wstream io.ReadWriteCloser) error { s.mutex.Unlock() var ackPayload []byte - if s.version == defs.APIMoQVersionDraft17 { + if s.version == defs.APIMoQVersionDraft16 || s.version == defs.APIMoQVersionDraft17 { ackPayload = controlmessage.PublishOk{}.Marshal() } else { ackPayload = controlmessage.RequestOk{}.Marshal()