improve write performance by

* buffering packets before sending them
* removing mutexes
This commit is contained in:
aler9
2021-12-08 17:46:56 +01:00
parent a1de5ffdf9
commit f3096ec102
20 changed files with 969 additions and 952 deletions

View File

@@ -2,6 +2,7 @@ package gortsplib
import (
"bufio"
"bytes"
"crypto/tls"
"net"
"strings"
@@ -48,14 +49,15 @@ func TestClientPublishSerial(t *testing.T) {
conn, err := l.Accept()
require.NoError(t, err)
defer conn.Close()
bconn := bufio.NewReadWriter(bufio.NewReader(conn), bufio.NewWriter(conn))
br := bufio.NewReader(conn)
var bb bytes.Buffer
req, err := readRequest(bconn.Reader)
req, err := readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Options, req.Method)
require.Equal(t, mustParseURL(scheme+"://localhost:8554/teststream"), req.URL)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
Header: base.Header{
"Public": base.HeaderValue{strings.Join([]string{
@@ -64,20 +66,22 @@ func TestClientPublishSerial(t *testing.T) {
string(base.Record),
}, ", ")},
},
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Announce, req.Method)
require.Equal(t, mustParseURL(scheme+"://localhost:8554/teststream"), req.URL)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Setup, req.Method)
require.Equal(t, mustParseURL(scheme+"://localhost:8554/teststream/trackID=0"), req.URL)
@@ -114,22 +118,24 @@ func TestClientPublishSerial(t *testing.T) {
th.InterleavedIDs = inTH.InterleavedIDs
}
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
Header: base.Header{
"Transport": th.Write(),
},
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Record, req.Method)
require.Equal(t, mustParseURL(scheme+"://localhost:8554/teststream"), req.URL)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
// client -> server
@@ -141,7 +147,7 @@ func TestClientPublishSerial(t *testing.T) {
} else {
var f base.InterleavedFrame
f.Payload = make([]byte, 2048)
err = f.Read(bconn.Reader)
err = f.Read(br)
require.NoError(t, err)
require.Equal(t, 0, f.Channel)
require.Equal(t, []byte{0x01, 0x02, 0x03, 0x04}, f.Payload)
@@ -154,21 +160,23 @@ func TestClientPublishSerial(t *testing.T) {
Port: th.ClientPorts[1],
})
} else {
err = base.InterleavedFrame{
base.InterleavedFrame{
Channel: 1,
Payload: []byte{0x05, 0x06, 0x07, 0x08},
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
}
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Teardown, req.Method)
require.Equal(t, mustParseURL(scheme+"://localhost:8554/teststream"), req.URL)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
}()
@@ -254,13 +262,14 @@ func TestClientPublishParallel(t *testing.T) {
conn, err := l.Accept()
require.NoError(t, err)
defer conn.Close()
bconn := bufio.NewReadWriter(bufio.NewReader(conn), bufio.NewWriter(conn))
br := bufio.NewReader(conn)
var bb bytes.Buffer
req, err := readRequest(bconn.Reader)
req, err := readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Options, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
Header: base.Header{
"Public": base.HeaderValue{strings.Join([]string{
@@ -269,19 +278,21 @@ func TestClientPublishParallel(t *testing.T) {
string(base.Record),
}, ", ")},
},
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Announce, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Setup, req.Method)
@@ -305,30 +316,33 @@ func TestClientPublishParallel(t *testing.T) {
th.InterleavedIDs = inTH.InterleavedIDs
}
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
Header: base.Header{
"Transport": th.Write(),
},
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Record, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequestIgnoreFrames(bconn.Reader)
req, err = readRequestIgnoreFrames(br)
require.NoError(t, err)
require.Equal(t, base.Teardown, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
}()
@@ -397,13 +411,14 @@ func TestClientPublishPauseSerial(t *testing.T) {
conn, err := l.Accept()
require.NoError(t, err)
defer conn.Close()
bconn := bufio.NewReadWriter(bufio.NewReader(conn), bufio.NewWriter(conn))
br := bufio.NewReader(conn)
var bb bytes.Buffer
req, err := readRequest(bconn.Reader)
req, err := readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Options, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
Header: base.Header{
"Public": base.HeaderValue{strings.Join([]string{
@@ -413,19 +428,21 @@ func TestClientPublishPauseSerial(t *testing.T) {
string(base.Pause),
}, ", ")},
},
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Announce, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Setup, req.Method)
@@ -449,48 +466,53 @@ func TestClientPublishPauseSerial(t *testing.T) {
th.InterleavedIDs = inTH.InterleavedIDs
}
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
Header: base.Header{
"Transport": th.Write(),
},
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Record, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequestIgnoreFrames(bconn.Reader)
req, err = readRequestIgnoreFrames(br)
require.NoError(t, err)
require.Equal(t, base.Pause, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Record, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequestIgnoreFrames(bconn.Reader)
req, err = readRequestIgnoreFrames(br)
require.NoError(t, err)
require.Equal(t, base.Teardown, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
}()
@@ -550,13 +572,14 @@ func TestClientPublishPauseParallel(t *testing.T) {
conn, err := l.Accept()
require.NoError(t, err)
defer conn.Close()
bconn := bufio.NewReadWriter(bufio.NewReader(conn), bufio.NewWriter(conn))
br := bufio.NewReader(conn)
var bb bytes.Buffer
req, err := readRequest(bconn.Reader)
req, err := readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Options, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
Header: base.Header{
"Public": base.HeaderValue{strings.Join([]string{
@@ -566,19 +589,21 @@ func TestClientPublishPauseParallel(t *testing.T) {
string(base.Pause),
}, ", ")},
},
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Announce, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Setup, req.Method)
@@ -602,30 +627,33 @@ func TestClientPublishPauseParallel(t *testing.T) {
th.InterleavedIDs = inTH.InterleavedIDs
}
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
Header: base.Header{
"Transport": th.Write(),
},
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Record, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequestIgnoreFrames(bconn.Reader)
req, err = readRequestIgnoreFrames(br)
require.NoError(t, err)
require.Equal(t, base.Pause, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
}()
@@ -689,14 +717,15 @@ func TestClientPublishAutomaticProtocol(t *testing.T) {
conn, err := l.Accept()
require.NoError(t, err)
defer conn.Close()
bconn := bufio.NewReadWriter(bufio.NewReader(conn), bufio.NewWriter(conn))
br := bufio.NewReader(conn)
var bb bytes.Buffer
req, err := readRequest(bconn.Reader)
req, err := readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Options, req.Method)
require.Equal(t, mustParseURL("rtsp://localhost:8554/teststream"), req.URL)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
Header: base.Header{
"Public": base.HeaderValue{strings.Join([]string{
@@ -705,29 +734,32 @@ func TestClientPublishAutomaticProtocol(t *testing.T) {
string(base.Record),
}, ", ")},
},
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Announce, req.Method)
require.Equal(t, mustParseURL("rtsp://localhost:8554/teststream"), req.URL)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Setup, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusUnsupportedTransport,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Setup, req.Method)
@@ -745,38 +777,41 @@ func TestClientPublishAutomaticProtocol(t *testing.T) {
InterleavedIDs: &[2]int{0, 1},
}
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
Header: base.Header{
"Transport": th.Write(),
},
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Record, req.Method)
require.Equal(t, mustParseURL("rtsp://localhost:8554/teststream"), req.URL)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
var f base.InterleavedFrame
f.Payload = make([]byte, 2048)
err = f.Read(bconn.Reader)
err = f.Read(br)
require.NoError(t, err)
require.Equal(t, 0, f.Channel)
require.Equal(t, []byte{0x01, 0x02, 0x03, 0x04}, f.Payload)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Teardown, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
}()
@@ -812,13 +847,14 @@ func TestClientPublishRTCPReport(t *testing.T) {
conn, err := l.Accept()
require.NoError(t, err)
defer conn.Close()
bconn := bufio.NewReadWriter(bufio.NewReader(conn), bufio.NewWriter(conn))
br := bufio.NewReader(conn)
var bb bytes.Buffer
req, err := readRequest(bconn.Reader)
req, err := readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Options, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
Header: base.Header{
"Public": base.HeaderValue{strings.Join([]string{
@@ -827,19 +863,21 @@ func TestClientPublishRTCPReport(t *testing.T) {
string(base.Record),
}, ", ")},
},
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Announce, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Setup, req.Method)
@@ -855,7 +893,7 @@ func TestClientPublishRTCPReport(t *testing.T) {
require.NoError(t, err)
defer l2.Close()
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
Header: base.Header{
"Transport": headers.Transport{
@@ -868,16 +906,18 @@ func TestClientPublishRTCPReport(t *testing.T) {
ServerPorts: &[2]int{34556, 34557},
}.Write(),
},
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Record, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
rr := rtcpreceiver.New(nil, 90000)
@@ -905,13 +945,15 @@ func TestClientPublishRTCPReport(t *testing.T) {
close(reportReceived)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Teardown, req.Method)
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
}()
c := &Client{
@@ -958,13 +1000,14 @@ func TestClientPublishIgnoreTCPRTPPackets(t *testing.T) {
conn, err := l.Accept()
require.NoError(t, err)
defer conn.Close()
bconn := bufio.NewReadWriter(bufio.NewReader(conn), bufio.NewWriter(conn))
br := bufio.NewReader(conn)
var bb bytes.Buffer
req, err := readRequest(bconn.Reader)
req, err := readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Options, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
Header: base.Header{
"Public": base.HeaderValue{strings.Join([]string{
@@ -973,19 +1016,21 @@ func TestClientPublishIgnoreTCPRTPPackets(t *testing.T) {
string(base.Record),
}, ", ")},
},
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Announce, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Setup, req.Method)
@@ -1002,42 +1047,48 @@ func TestClientPublishIgnoreTCPRTPPackets(t *testing.T) {
InterleavedIDs: inTH.InterleavedIDs,
}
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
Header: base.Header{
"Transport": th.Write(),
},
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Record, req.Method)
err = base.Response{
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
err = base.InterleavedFrame{
base.InterleavedFrame{
Channel: 0,
Payload: []byte{0x01, 0x02, 0x03, 0x04},
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
err = base.InterleavedFrame{
base.InterleavedFrame{
Channel: 1,
Payload: []byte{0x05, 0x06, 0x07, 0x08},
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
req, err = readRequest(bconn.Reader)
req, err = readRequest(br)
require.NoError(t, err)
require.Equal(t, base.Teardown, req.Method)
base.Response{
StatusCode: base.StatusOK,
}.Write(bconn.Writer)
}.Write(&bb)
_, err = conn.Write(bb.Bytes())
require.NoError(t, err)
}()
rtcpReceived := make(chan struct{})