Migrate to pion/webrtc v3 (#254)

* Add initial pion/webrtc v3 migration

* Use custom timestamps when sending video packets

* Update dependencies
This commit is contained in:
sergystepanov 2021-01-04 21:09:03 +03:00 committed by GitHub
parent 10db5d6c88
commit 0898fb06ca
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
5 changed files with 240 additions and 235 deletions

111
pkg/webrtc/track.go Normal file
View file

@ -0,0 +1,111 @@
package webrtc
import (
"strings"
"sync"
"github.com/pion/rtp"
"github.com/pion/rtp/codecs"
"github.com/pion/webrtc/v3"
"github.com/pion/webrtc/v3/pkg/media"
)
// CustomTrackSample is used just for adding custom timestamps
// into outgoing packets, since packetizer is not accessible anymore.
// Use webrtc.TrackLocalStaticSample instead if you use constant rate streams.
type CustomTrackSample struct {
packetizer rtp.Packetizer
rtpTrack *webrtc.TrackLocalStaticRTP
clockRate float64
mu sync.RWMutex
}
func NewCustomTrackSample(c webrtc.RTPCodecCapability, id, streamID string) (*CustomTrackSample, error) {
rtpTrack, err := webrtc.NewTrackLocalStaticRTP(c, id, streamID)
if err != nil {
return nil, err
}
return &CustomTrackSample{rtpTrack: rtpTrack}, nil
}
func (s *CustomTrackSample) ID() string { return s.rtpTrack.ID() }
func (s *CustomTrackSample) StreamID() string { return s.rtpTrack.StreamID() }
func (s *CustomTrackSample) Kind() webrtc.RTPCodecType { return s.rtpTrack.Kind() }
func (s *CustomTrackSample) Codec() webrtc.RTPCodecCapability { return s.rtpTrack.Codec() }
func (s *CustomTrackSample) Bind(t webrtc.TrackLocalContext) (webrtc.RTPCodecParameters, error) {
rtpOutboundMTU := 1200
codec, err := s.rtpTrack.Bind(t)
if err != nil {
return codec, err
}
s.mu.Lock()
defer s.mu.Unlock()
if s.packetizer != nil {
return codec, nil
}
payloader, err := payloaderForCodec(codec.RTPCodecCapability)
if err != nil {
return codec, err
}
s.packetizer = rtp.NewPacketizer(
rtpOutboundMTU,
0, // Value is handled when writing
0, // Value is handled when writing
payloader,
rtp.NewRandomSequencer(),
codec.ClockRate,
)
s.clockRate = float64(codec.RTPCodecCapability.ClockRate)
return codec, nil
}
func (s *CustomTrackSample) Unbind(t webrtc.TrackLocalContext) error {
return s.rtpTrack.Unbind(t)
}
func (s *CustomTrackSample) WriteSampleWithTimestamp(sample media.Sample, timestamp uint32) (err error) {
s.mu.RLock()
p := s.packetizer
clockRate := s.clockRate
s.mu.RUnlock()
if p == nil {
return nil
}
samples := sample.Duration.Seconds() * clockRate
packets := p.(rtp.Packetizer).Packetize(sample.Data, uint32(samples))
for _, p := range packets {
p.Timestamp = timestamp
err = s.rtpTrack.WriteRTP(p)
}
return
}
func payloaderForCodec(codec webrtc.RTPCodecCapability) (rtp.Payloader, error) {
switch strings.ToLower(codec.MimeType) {
case strings.ToLower(webrtc.MimeTypeH264):
return &codecs.H264Payloader{}, nil
case strings.ToLower(webrtc.MimeTypeOpus):
return &codecs.OpusPayloader{}, nil
case strings.ToLower(webrtc.MimeTypeVP8):
return &codecs.VP8Payloader{}, nil
case strings.ToLower(webrtc.MimeTypeVP9):
return &codecs.VP9Payloader{}, nil
case strings.ToLower(webrtc.MimeTypeG722):
return &codecs.G722Payloader{}, nil
case strings.ToLower(webrtc.MimeTypePCMU), strings.ToLower(webrtc.MimeTypePCMA):
return &codecs.G711Payloader{}, nil
default:
return nil, webrtc.ErrNoPayloaderForCodec
}
}

View file

@ -6,7 +6,6 @@ import (
"encoding/json"
"fmt"
"log"
"math/rand"
"runtime/debug"
"time"
@ -14,18 +13,13 @@ import (
"github.com/giongto35/cloud-game/v2/pkg/encoder"
"github.com/giongto35/cloud-game/v2/pkg/util"
"github.com/gofrs/uuid"
"github.com/pion/webrtc/v2"
"github.com/pion/webrtc/v2/pkg/media"
"github.com/pion/webrtc/v3"
"github.com/pion/webrtc/v3/pkg/media"
)
// TODO: double check if no need TURN server here
var webrtcconfig = webrtc.Configuration{ICEServers: []webrtc.ICEServer{{URLs: []string{"stun:stun.l.google.com:19302"}}}}
type InputDataPair struct {
data int
time time.Time
}
type WebFrame struct {
Data []byte
Timestamp uint32
@ -116,7 +110,7 @@ func (w *WebRTC) StartClient(isMobile bool, iceCB OnIceCallback) (string, error)
}
}()
var err error
var videoTrack *webrtc.Track
var videoTrack *CustomTrackSample
// reset client
if w.isConnected {
@ -131,11 +125,13 @@ func (w *WebRTC) StartClient(isMobile bool, iceCB OnIceCallback) (string, error)
}
// add video track
var codec webrtc.RTPCodecCapability
if util.GetVideoEncoder(isMobile) == encoder.H264 {
videoTrack, err = w.connection.NewTrack(webrtc.DefaultPayloadTypeH264, rand.Uint32(), "video", "game-video")
codec = webrtc.RTPCodecCapability{MimeType: "video/h264"}
} else {
videoTrack, err = w.connection.NewTrack(webrtc.DefaultPayloadTypeVP8, rand.Uint32(), "video", "game-video")
codec = webrtc.RTPCodecCapability{MimeType: "video/vp8"}
}
videoTrack, err = NewCustomTrackSample(codec, "video", "game-video")
if err != nil {
return "", err
}
@ -147,7 +143,7 @@ func (w *WebRTC) StartClient(isMobile bool, iceCB OnIceCallback) (string, error)
log.Println("Add video track")
// add audio track
opusTrack, err := w.connection.NewTrack(webrtc.DefaultPayloadTypeOpus, rand.Uint32(), "audio", "game-audio")
opusTrack, err := webrtc.NewTrackLocalStaticSample(webrtc.RTPCodecCapability{MimeType: "audio/opus"}, "audio", "game-audio")
if err != nil {
return "", err
}
@ -209,7 +205,7 @@ func (w *WebRTC) StartClient(isMobile bool, iceCB OnIceCallback) (string, error)
})
w.connection.OnTrack(func(remoteTrack *webrtc.Track, receiver *webrtc.RTPReceiver) {
w.connection.OnTrack(func(remoteTrack *webrtc.TrackRemote, receiver *webrtc.RTPReceiver) {
//NOTE: High CPU due to constantly for loop. Turn it off first, Fix it later.
//rtpBuf := make([]byte, 1400)
@ -317,7 +313,7 @@ func (w *WebRTC) IsConnected() bool {
return w.isConnected
}
func (w *WebRTC) startStreaming(vp8Track *webrtc.Track, opusTrack *webrtc.Track) {
func (w *WebRTC) startStreaming(vp8Track *CustomTrackSample, opusTrack *webrtc.TrackLocalStaticSample) {
log.Println("Start streaming")
// receive frame buffer
go func() {
@ -329,14 +325,10 @@ func (w *WebRTC) startStreaming(vp8Track *webrtc.Track, opusTrack *webrtc.Track)
}()
for data := range w.ImageChannel {
packets := vp8Track.Packetizer().Packetize(data.Data, 1)
for _, p := range packets {
p.Header.Timestamp = data.Timestamp
err := vp8Track.WriteRTP(p)
if err != nil {
log.Println("Warn: Err write sample: ", err)
break
}
err := vp8Track.WriteSampleWithTimestamp(media.Sample{Data: data.Data}, data.Timestamp)
if err != nil {
log.Println("Warn: Err write sample: ", err)
break
}
}
}()
@ -350,13 +342,14 @@ func (w *WebRTC) startStreaming(vp8Track *webrtc.Track, opusTrack *webrtc.Track)
}
}()
opusSamples := uint32(w.cfg.Encoder.Audio.GetFrameDuration() / w.cfg.Encoder.Audio.Channels)
//opusSamples := uint32(w.cfg.Encoder.Audio.GetFrameDuration() / w.cfg.Encoder.Audio.Channels)
audioDuration := time.Duration(w.cfg.Encoder.Audio.Frame) * time.Millisecond
for data := range w.AudioChannel {
if !w.isConnected {
return
}
err := opusTrack.WriteSample(media.Sample{Data: data, Samples: opusSamples})
err := opusTrack.WriteSample(media.Sample{Data: data, Duration: audioDuration})
if err != nil {
log.Println("Warn: Err write sample: ", err)
}
@ -376,7 +369,8 @@ func (w *WebRTC) startStreaming(vp8Track *webrtc.Track, opusTrack *webrtc.Track)
if !w.isConnected {
return
}
_, err := opusTrack.Write(data)
// !to pass duration from the input
err := opusTrack.WriteSample(media.Sample{Data: data})
if err != nil {
log.Println("Warn: Err write sample: ", err)
}

View file

@ -157,12 +157,13 @@ func (r *Room) startVideo(width, height int, videoCodec encoder.VideoCodec) {
for data := range eoutput {
// TODO: r.rtcSessions is rarely updated. Lock will hold down perf
for _, webRTC := range r.rtcSessions {
if !webRTC.IsConnected() {
continue
}
// encode frame
// fanout imageChannel
if webRTC.IsConnected() {
// NOTE: can block here
webRTC.ImageChannel <- webrtc.WebFrame{Data: data.Data, Timestamp: data.Timestamp}
}
// NOTE: can block here
webRTC.ImageChannel <- webrtc.WebFrame{Data: data.Data, Timestamp: data.Timestamp}
}
}
}()