mirror of
https://github.com/aler9/gortsplib
synced 2025-10-04 23:02:45 +08:00
client: add range argument to clientconn.Play()
This commit is contained in:
@@ -165,7 +165,7 @@ func (c *Client) DialReadContext(ctx context.Context, address string) (*ClientCo
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
_, err = conn.Play()
|
_, err = conn.Play(nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
conn.Close()
|
conn.Close()
|
||||||
return nil, err
|
return nil, err
|
||||||
|
@@ -1252,7 +1252,7 @@ func TestClientReadPause(t *testing.T) {
|
|||||||
conn.ReadFrames(func(id int, typ StreamType, payload []byte) {
|
conn.ReadFrames(func(id int, typ StreamType, payload []byte) {
|
||||||
})
|
})
|
||||||
|
|
||||||
_, err = conn.Play()
|
_, err = conn.Play(nil)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
firstFrame = int32(0)
|
firstFrame = int32(0)
|
||||||
@@ -1734,3 +1734,139 @@ func TestClientReadIgnoreTCPInvalidTrack(t *testing.T) {
|
|||||||
conn.Close()
|
conn.Close()
|
||||||
<-done
|
<-done
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestClientReadSeek(t *testing.T) {
|
||||||
|
l, err := net.Listen("tcp", "localhost:8554")
|
||||||
|
require.NoError(t, err)
|
||||||
|
defer l.Close()
|
||||||
|
|
||||||
|
serverDone := make(chan struct{})
|
||||||
|
defer func() { <-serverDone }()
|
||||||
|
go func() {
|
||||||
|
defer close(serverDone)
|
||||||
|
|
||||||
|
conn, err := l.Accept()
|
||||||
|
require.NoError(t, err)
|
||||||
|
defer conn.Close()
|
||||||
|
bconn := bufio.NewReadWriter(bufio.NewReader(conn), bufio.NewWriter(conn))
|
||||||
|
|
||||||
|
req, err := readRequest(bconn.Reader)
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.Equal(t, base.Options, req.Method)
|
||||||
|
|
||||||
|
err = base.Response{
|
||||||
|
StatusCode: base.StatusOK,
|
||||||
|
Header: base.Header{
|
||||||
|
"Public": base.HeaderValue{strings.Join([]string{
|
||||||
|
string(base.Describe),
|
||||||
|
string(base.Setup),
|
||||||
|
string(base.Play),
|
||||||
|
}, ", ")},
|
||||||
|
},
|
||||||
|
}.Write(bconn.Writer)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
req, err = readRequest(bconn.Reader)
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.Equal(t, base.Describe, req.Method)
|
||||||
|
|
||||||
|
track, err := NewTrackH264(96, []byte("123456"), []byte("123456"))
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
err = base.Response{
|
||||||
|
StatusCode: base.StatusOK,
|
||||||
|
Header: base.Header{
|
||||||
|
"Content-Type": base.HeaderValue{"application/sdp"},
|
||||||
|
"Content-Base": base.HeaderValue{"rtsp://localhost:8554/teststream/"},
|
||||||
|
},
|
||||||
|
Body: Tracks{track}.Write(),
|
||||||
|
}.Write(bconn.Writer)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
req, err = readRequest(bconn.Reader)
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.Equal(t, base.Setup, req.Method)
|
||||||
|
|
||||||
|
var inTH headers.Transport
|
||||||
|
err = inTH.Read(req.Header["Transport"])
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
th := headers.Transport{
|
||||||
|
Delivery: func() *base.StreamDelivery {
|
||||||
|
v := base.StreamDeliveryUnicast
|
||||||
|
return &v
|
||||||
|
}(),
|
||||||
|
Protocol: StreamProtocolTCP,
|
||||||
|
InterleavedIDs: inTH.InterleavedIDs,
|
||||||
|
}
|
||||||
|
|
||||||
|
err = base.Response{
|
||||||
|
StatusCode: base.StatusOK,
|
||||||
|
Header: base.Header{
|
||||||
|
"Transport": th.Write(),
|
||||||
|
},
|
||||||
|
}.Write(bconn.Writer)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
req, err = readRequest(bconn.Reader)
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.Equal(t, base.Play, req.Method)
|
||||||
|
|
||||||
|
var ra headers.Range
|
||||||
|
err = ra.Read(req.Header["Range"])
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.Equal(t, headers.Range{
|
||||||
|
Value: &headers.RangeNPT{
|
||||||
|
Start: headers.RangeNPTTime(5500 * time.Millisecond),
|
||||||
|
},
|
||||||
|
}, ra)
|
||||||
|
|
||||||
|
err = base.Response{
|
||||||
|
StatusCode: base.StatusOK,
|
||||||
|
}.Write(bconn.Writer)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
req, err = readRequest(bconn.Reader)
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.Equal(t, base.Teardown, req.Method)
|
||||||
|
|
||||||
|
err = base.Response{
|
||||||
|
StatusCode: base.StatusOK,
|
||||||
|
}.Write(bconn.Writer)
|
||||||
|
require.NoError(t, err)
|
||||||
|
}()
|
||||||
|
|
||||||
|
c := &Client{
|
||||||
|
StreamProtocol: func() *StreamProtocol {
|
||||||
|
v := StreamProtocolTCP
|
||||||
|
return &v
|
||||||
|
}(),
|
||||||
|
}
|
||||||
|
|
||||||
|
u, err := base.ParseURL("rtsp://localhost:8554/teststream")
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
conn, err := c.Dial(u.Scheme, u.Host)
|
||||||
|
require.NoError(t, err)
|
||||||
|
defer conn.Close()
|
||||||
|
|
||||||
|
_, err = conn.Options(u)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
tracks, _, err := conn.Describe(u)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
for _, track := range tracks {
|
||||||
|
_, err := conn.Setup(headers.TransportModePlay, track, 0, 0)
|
||||||
|
require.NoError(t, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
_, err = conn.Play(&headers.Range{
|
||||||
|
Value: &headers.RangeNPT{
|
||||||
|
Start: headers.RangeNPTTime(5500 * time.Millisecond),
|
||||||
|
},
|
||||||
|
})
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
// asdasdasd
|
||||||
|
}
|
||||||
|
@@ -96,6 +96,7 @@ type setupReq struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
type playReq struct {
|
type playReq struct {
|
||||||
|
ra *headers.Range
|
||||||
res chan clientRes
|
res chan clientRes
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -131,6 +132,7 @@ type ClientConn struct {
|
|||||||
streamBaseURL *base.URL
|
streamBaseURL *base.URL
|
||||||
streamProtocol *StreamProtocol
|
streamProtocol *StreamProtocol
|
||||||
tracks map[int]clientConnTrack
|
tracks map[int]clientConnTrack
|
||||||
|
lastRange *headers.Range
|
||||||
backgroundRunning bool
|
backgroundRunning bool
|
||||||
backgroundErr error
|
backgroundErr error
|
||||||
tcpFrameBuffer *multibuffer.MultiBuffer // tcp
|
tcpFrameBuffer *multibuffer.MultiBuffer // tcp
|
||||||
@@ -267,7 +269,7 @@ outer:
|
|||||||
req.res <- clientRes{res: res, err: err}
|
req.res <- clientRes{res: res, err: err}
|
||||||
|
|
||||||
case req := <-cc.play:
|
case req := <-cc.play:
|
||||||
res, err := cc.doPlay(false)
|
res, err := cc.doPlay(req.ra, false)
|
||||||
req.res <- clientRes{res: res, err: err}
|
req.res <- clientRes{res: res, err: err}
|
||||||
|
|
||||||
case req := <-cc.record:
|
case req := <-cc.record:
|
||||||
@@ -389,7 +391,7 @@ func (cc *ClientConn) switchProtocolIfTimeout(err error) error {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
_, err = cc.doPlay(true)
|
_, err = cc.doPlay(cc.lastRange, true)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
@@ -1409,7 +1411,7 @@ func (cc *ClientConn) Setup(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (cc *ClientConn) doPlay(isSwitchingProtocol bool) (*base.Response, error) {
|
func (cc *ClientConn) doPlay(ra *headers.Range, isSwitchingProtocol bool) (*base.Response, error) {
|
||||||
err := cc.checkState(map[clientConnState]struct{}{
|
err := cc.checkState(map[clientConnState]struct{}{
|
||||||
clientConnStatePrePlay: {},
|
clientConnStatePrePlay: {},
|
||||||
})
|
})
|
||||||
@@ -1417,9 +1419,16 @@ func (cc *ClientConn) doPlay(isSwitchingProtocol bool) (*base.Response, error) {
|
|||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
header := make(base.Header)
|
||||||
|
|
||||||
|
if ra != nil {
|
||||||
|
header["Range"] = ra.Write()
|
||||||
|
}
|
||||||
|
|
||||||
res, err := cc.do(&base.Request{
|
res, err := cc.do(&base.Request{
|
||||||
Method: base.Play,
|
Method: base.Play,
|
||||||
URL: cc.streamBaseURL,
|
URL: cc.streamBaseURL,
|
||||||
|
Header: header,
|
||||||
}, false)
|
}, false)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -1432,6 +1441,7 @@ func (cc *ClientConn) doPlay(isSwitchingProtocol bool) (*base.Response, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
cc.state = clientConnStatePlay
|
cc.state = clientConnStatePlay
|
||||||
|
cc.lastRange = ra
|
||||||
|
|
||||||
if !isSwitchingProtocol {
|
if !isSwitchingProtocol {
|
||||||
// use a temporary callback that is replaces as soon as
|
// use a temporary callback that is replaces as soon as
|
||||||
@@ -1455,10 +1465,10 @@ func (cc *ClientConn) doPlay(isSwitchingProtocol bool) (*base.Response, error) {
|
|||||||
|
|
||||||
// Play writes a PLAY request and reads a Response.
|
// Play writes a PLAY request and reads a Response.
|
||||||
// This can be called only after Setup().
|
// This can be called only after Setup().
|
||||||
func (cc *ClientConn) Play() (*base.Response, error) {
|
func (cc *ClientConn) Play(ra *headers.Range) (*base.Response, error) {
|
||||||
cres := make(chan clientRes)
|
cres := make(chan clientRes)
|
||||||
select {
|
select {
|
||||||
case cc.play <- playReq{res: cres}:
|
case cc.play <- playReq{ra: ra, res: cres}:
|
||||||
res := <-cres
|
res := <-cres
|
||||||
return res.res, res.err
|
return res.res, res.err
|
||||||
|
|
||||||
|
@@ -46,7 +46,7 @@ func main() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// play setupped tracks
|
// play setupped tracks
|
||||||
_, err = conn.Play()
|
_, err = conn.Play(nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
panic(err)
|
panic(err)
|
||||||
}
|
}
|
||||||
|
@@ -47,7 +47,7 @@ func main() {
|
|||||||
time.Sleep(5 * time.Second)
|
time.Sleep(5 * time.Second)
|
||||||
|
|
||||||
// play again
|
// play again
|
||||||
_, err = conn.Play()
|
_, err = conn.Play(nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
panic(err)
|
panic(err)
|
||||||
}
|
}
|
||||||
|
Reference in New Issue
Block a user