moq: support draft-16 (#6045)

This commit is contained in:
Alessandro Ros
2026-08-05 12:16:07 +02:00
committed by GitHub
parent f0d2f11525
commit 1b943637a4
12 changed files with 352 additions and 183 deletions
+1
View File
@@ -98,6 +98,7 @@ components:
MoQVersion: MoQVersion:
type: string type: string
enum: enum:
- moqt-16
- moqt-17 - moqt-17
- moqt-18 - moqt-18
- moqt-19 - moqt-19
+1 -1
View File
@@ -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: 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 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 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). - 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).
+1
View File
@@ -28,6 +28,7 @@ type APIMoQVersion string
// protocol versions. // protocol versions.
const ( const (
APIMoQVersionDraft16 APIMoQVersion = "moqt-16"
APIMoQVersionDraft17 APIMoQVersion = "moqt-17" APIMoQVersionDraft17 APIMoQVersion = "moqt-17"
APIMoQVersionDraft18 APIMoQVersion = "moqt-18" APIMoQVersionDraft18 APIMoQVersion = "moqt-18"
APIMoQVersionDraft19 APIMoQVersion = "moqt-19" APIMoQVersionDraft19 APIMoQVersion = "moqt-19"
@@ -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
}
@@ -42,6 +42,10 @@ func Read(r io.Reader) (Message, error) {
switch t { switch t {
case typeSetup: case typeSetup:
m = &Setup{} m = &Setup{}
case typeClientSetup:
m = &ClientSetup{}
case typeServerSetup:
m = &ServerSetup{}
case typeSubscribe: case typeSubscribe:
m = &Subscribe{} m = &Subscribe{}
case typeSubscribeOk: case typeSubscribeOk:
@@ -7,6 +7,7 @@ import (
"github.com/stretchr/testify/require" "github.com/stretchr/testify/require"
"github.com/bluenviron/mediamtx/internal/protocols/moq/controlmessage" "github.com/bluenviron/mediamtx/internal/protocols/moq/controlmessage"
"github.com/bluenviron/mediamtx/internal/protocols/moq/namespace"
"github.com/bluenviron/mediamtx/internal/protocols/moq/parameter" "github.com/bluenviron/mediamtx/internal/protocols/moq/parameter"
) )
@@ -15,6 +16,22 @@ var cases = []struct {
enc []byte enc []byte
dec controlmessage.Message 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", name: "setup",
enc: []byte{ enc: []byte{
@@ -66,7 +83,7 @@ var cases = []struct {
}, },
dec: &controlmessage.Subscribe{ dec: &controlmessage.Subscribe{
RequestID: 1, RequestID: 1,
Namespace: []string{"foo"}, Namespace: namespace.Namespace{"foo"},
TrackName: "bar", TrackName: "bar",
}, },
}, },
@@ -88,7 +105,7 @@ var cases = []struct {
}, },
dec: &controlmessage.Subscribe{ dec: &controlmessage.Subscribe{
RequestID: 1, RequestID: 1,
Namespace: []string{"foo"}, Namespace: namespace.Namespace{"foo"},
TrackName: "bar", TrackName: "bar",
Parameters: []parameter.Parameter{ Parameters: []parameter.Parameter{
&parameter.AuthorizationToken{ &parameter.AuthorizationToken{
@@ -157,7 +174,7 @@ var cases = []struct {
}, },
dec: &controlmessage.Publish{ dec: &controlmessage.Publish{
RequestID: 1, RequestID: 1,
Namespace: []string{"foo"}, Namespace: namespace.Namespace{"foo"},
TrackName: "bar", TrackName: "bar",
TrackAlias: 2, TrackAlias: 2,
}, },
@@ -181,7 +198,7 @@ var cases = []struct {
}, },
dec: &controlmessage.Publish{ dec: &controlmessage.Publish{
RequestID: 1, RequestID: 1,
Namespace: []string{"foo"}, Namespace: namespace.Namespace{"foo"},
TrackName: "bar", TrackName: "bar",
TrackAlias: 2, TrackAlias: 2,
Parameters: []parameter.Parameter{ Parameters: []parameter.Parameter{
@@ -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
}
@@ -74,7 +74,7 @@ func (m *Setup) unmarshal(buf []byte) error {
return nil return nil
} }
func (m Setup) marshalSize() int { func (m Setup) marshalSize(t varint.Varint) int {
payloadSize := 0 payloadSize := 0
var previousType varint.Varint var previousType varint.Varint
@@ -93,13 +93,13 @@ func (m Setup) marshalSize() int {
len(m.Authority) len(m.Authority)
} }
return typeSetup.MarshalSize() + 2 + payloadSize return t.MarshalSize() + 2 + payloadSize
} }
func (m Setup) marshalTo(buf []byte) int { func (m Setup) marshalTo(buf []byte, t varint.Varint) int {
payloadSize := m.marshalSize() - typeSetup.MarshalSize() - 2 payloadSize := m.marshalSize(t) - t.MarshalSize() - 2
pos := typeSetup.MarshalTo(buf) pos := t.MarshalTo(buf)
buf[pos] = byte(payloadSize >> 8) buf[pos] = byte(payloadSize >> 8)
buf[pos+1] = byte(payloadSize) buf[pos+1] = byte(payloadSize)
pos += 2 pos += 2
@@ -126,7 +126,7 @@ func (m Setup) marshalTo(buf []byte) int {
// Marshal implements Message. // Marshal implements Message.
func (m Setup) Marshal() []byte { func (m Setup) Marshal() []byte {
buf := make([]byte, m.marshalSize()) buf := make([]byte, m.marshalSize(typeSetup))
m.marshalTo(buf) m.marshalTo(buf, typeSetup)
return buf return buf
} }
+1
View File
@@ -46,6 +46,7 @@ var supportedMoqtVersions = []defs.APIMoQVersion{
defs.APIMoQVersionDraft19, defs.APIMoQVersionDraft19,
defs.APIMoQVersionDraft18, defs.APIMoQVersionDraft18,
defs.APIMoQVersionDraft17, defs.APIMoQVersionDraft17,
defs.APIMoQVersionDraft16,
} }
type ginUnwrapper interface { type ginUnwrapper interface {
+4
View File
@@ -19,6 +19,7 @@ var supportedMoqtALPNs = []string{
string(defs.APIMoQVersionDraft19), string(defs.APIMoQVersionDraft19),
string(defs.APIMoQVersionDraft18), string(defs.APIMoQVersionDraft18),
string(defs.APIMoQVersionDraft17), string(defs.APIMoQVersionDraft17),
string(defs.APIMoQVersionDraft16),
} }
type nativeListenerParent interface { type nativeListenerParent interface {
@@ -101,6 +102,9 @@ func alpnToVersion(alpn string) defs.APIMoQVersion {
case string(defs.APIMoQVersionDraft17): case string(defs.APIMoQVersionDraft17):
return defs.APIMoQVersionDraft17 return defs.APIMoQVersionDraft17
case string(defs.APIMoQVersionDraft16):
return defs.APIMoQVersionDraft16
} }
return "" return ""
+182 -126
View File
@@ -37,6 +37,76 @@ func (p *serverDummyPath) ExternalCmdEnv() externalcmd.Environment { retur
func (p *serverDummyPath) RemovePublisher(_ defs.PathRemovePublisherReq) {} func (p *serverDummyPath) RemovePublisher(_ defs.PathRemovePublisherReq) {}
func (p *serverDummyPath) RemoveReader(_ defs.PathRemoveReaderReq) {} 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) { func TestAuthError(t *testing.T) {
serverCertFile := test.CreateTempFile(t, test.TLSCertPub) serverCertFile := test.CreateTempFile(t, test.TLSCertPub)
serverKeyFile := test.CreateTempFile(t, test.TLSCertKey) serverKeyFile := test.CreateTempFile(t, test.TLSCertKey)
@@ -166,20 +236,7 @@ func TestAuthError(t *testing.T) {
defer sx.CloseWithError(0, "") //nolint:errcheck defer sx.CloseWithError(0, "") //nolint:errcheck
defer res.Body.Close() //nolint:errcheck defer res.Body.Close() //nolint:errcheck
setupStream, err := sx.AcceptUniStream(ctx) performWTSetup(ctx, t, sx, defs.APIMoQVersionDraft19)
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)
switch ca { switch ca {
case "subscribe": case "subscribe":
@@ -272,6 +329,11 @@ func TestServer(t *testing.T) {
clientProtocols []string clientProtocols []string
expectedVersion defs.APIMoQVersion expectedVersion defs.APIMoQVersion
}{ }{
{
name: "draft-16",
clientProtocols: []string{"moqt-16"},
expectedVersion: defs.APIMoQVersionDraft16,
},
{ {
name: "draft-17", name: "draft-17",
clientProtocols: []string{"moqt-17"}, clientProtocols: []string{"moqt-17"},
@@ -294,7 +356,7 @@ func TestServer(t *testing.T) {
}, },
{ {
name: "highest-preferred", name: "highest-preferred",
clientProtocols: []string{"moqt-17", "moqt-18", "moqt-19"}, clientProtocols: []string{"moqt-16", "moqt-17", "moqt-18", "moqt-19"},
expectedVersion: defs.APIMoQVersionDraft19, expectedVersion: defs.APIMoQVersionDraft19,
}, },
} { } {
@@ -368,18 +430,7 @@ func TestServer(t *testing.T) {
require.Equal(t, `"`+string(ca.expectedVersion)+`"`, res.Header.Get("WT-Protocol")) require.Equal(t, `"`+string(ca.expectedVersion)+`"`, res.Header.Get("WT-Protocol"))
setupStream, err := sx.AcceptUniStream(ctx) performWTSetup(ctx, t, sx, ca.expectedVersion)
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)
catalogBidi, err := sx.OpenStreamSync(ctx) catalogBidi, err := sx.OpenStreamSync(ctx)
require.NoError(t, err) require.NoError(t, err)
@@ -541,106 +592,111 @@ func TestServerUnsupportedVersion(t *testing.T) {
} }
func TestServerNativeQUICSubscribe(t *testing.T) { func TestServerNativeQUICSubscribe(t *testing.T) {
desc := &description.Session{Medias: []*description.Media{test.UniqueMediaH264()}} for _, ca := range []struct {
strm := &stream.Stream{ name string
OrigDesc: desc, version defs.APIMoQVersion
WriteQueueSize: 512, }{
RTPMaxPayloadSize: 1450, {
Parent: test.NilLogger, name: "draft-16",
} version: defs.APIMoQVersionDraft16,
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 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)
} }
+84 -45
View File
@@ -175,9 +175,11 @@ func (s *session) runInner() error {
return s.runBidiStreamAcceptor(errGroup) return s.runBidiStreamAcceptor(errGroup)
}) })
errGroup.Go(func() error { if s.version != defs.APIMoQVersionDraft16 {
return s.runSetupWriter() errGroup.Go(func() error {
}) return s.runSetupWriter()
})
}
select { select {
case <-s.ctx.Done(): case <-s.ctx.Done():
@@ -254,47 +256,11 @@ func (s *session) onUniMessage(r io.Reader) error {
switch m := msg.(type) { switch m := msg.(type) {
case *controlmessage.Setup: case *controlmessage.Setup:
err = func() error { if s.version == defs.APIMoQVersionDraft16 {
s.mutex.Lock() return fmt.Errorf("received SETUP over unidirectional stream with draft-16")
defer s.mutex.Unlock() }
if s.transport == defs.APIMoQSessionTransportWebTransport { err = s.processSetupMessage(m)
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
}
}()
if err != nil { if err != nil {
return err 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 { 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 { select {
case <-s.setupReceived: case <-s.setupReceived:
case <-s.ctx.Done(): case <-s.ctx.Done():
@@ -622,7 +661,7 @@ func (s *session) onPublishCatalog(wstream io.ReadWriteCloser, m *controlmessage
} }
var ackPayload []byte var ackPayload []byte
if s.version == defs.APIMoQVersionDraft17 { if s.version == defs.APIMoQVersionDraft16 || s.version == defs.APIMoQVersionDraft17 {
ackPayload = controlmessage.PublishOk{}.Marshal() ackPayload = controlmessage.PublishOk{}.Marshal()
} else { } else {
ackPayload = controlmessage.RequestOk{}.Marshal() ackPayload = controlmessage.RequestOk{}.Marshal()
@@ -646,7 +685,7 @@ func (s *session) onPublishTrack(wstream io.ReadWriteCloser) error {
s.mutex.Unlock() s.mutex.Unlock()
var ackPayload []byte var ackPayload []byte
if s.version == defs.APIMoQVersionDraft17 { if s.version == defs.APIMoQVersionDraft16 || s.version == defs.APIMoQVersionDraft17 {
ackPayload = controlmessage.PublishOk{}.Marshal() ackPayload = controlmessage.PublishOk{}.Marshal()
} else { } else {
ackPayload = controlmessage.RequestOk{}.Marshal() ackPayload = controlmessage.RequestOk{}.Marshal()