fix offer/answer handling

This commit is contained in:
kaedwen committed 2024-03-25 20:37:24 +01:00
1 parent 2a70fd9331
commit 555fce9f10
9 files changed
+274 -71

No files matched your search

+11 -12
View File
@@ -10,7 +10,6 @@ import (
"github.com/gin-gonic/gin"
"github.com/kaedwen/webrtc/pkg/common"
"github.com/kaedwen/webrtc/static"
"github.com/pion/webrtc/v3"
"go.uber.org/zap"
"nhooyr.io/websocket"
"nhooyr.io/websocket/wsjson"
@@ -18,8 +17,8 @@ import (
type SignalingHandle struct {
Id string
Recv chan webrtc.SessionDescription
Trcv chan webrtc.SessionDescription
Recv chan *IncomingSignalingMessage
Trcv chan *OutgoingSignalingMessage
}
type HttpServer struct {
@@ -32,8 +31,8 @@ type HttpServer struct {
func NewSignalingHandle(id string) SignalingHandle {
return SignalingHandle{
Id: id,
Recv: make(chan webrtc.SessionDescription, 10),
Trcv: make(chan webrtc.SessionDescription, 10),
Recv: make(chan *IncomingSignalingMessage, 10),
Trcv: make(chan *OutgoingSignalingMessage, 10),
}
}
@@ -143,8 +142,8 @@ func (h *HttpServer) signalingHandler(c *gin.Context) {
go func() {
for {
// wait and read/parse message
var sdp webrtc.SessionDescription
err := wsjson.Read(ctx, conn, &sdp)
var m IncomingSignalingMessage
err := wsjson.Read(ctx, conn, &m)
if err != nil {
status := websocket.CloseStatus(err)
if status == websocket.StatusNormalClosure || status == websocket.StatusGoingAway {
@@ -161,18 +160,18 @@ func (h *HttpServer) signalingHandler(c *gin.Context) {
break
}
h.lg.Info("received message", zap.String("type", sdp.Type.String()))
h.lg.Info("received message", zap.String("type", m.Type))
// forward in channel
hndl.Recv <- sdp
hndl.Recv <- &m
}
cancel()
}()
// loop write
for sdp := range hndl.Trcv {
err := wsjson.Write(ctx, conn, sdp)
for m := range hndl.Trcv {
err := wsjson.Write(ctx, conn, m)
if err != nil {
status := websocket.CloseStatus(err)
if status == websocket.StatusNormalClosure || status == websocket.StatusGoingAway {
@@ -184,7 +183,7 @@ func (h *HttpServer) signalingHandler(c *gin.Context) {
break
}
h.lg.Info("tranceived message", zap.String("type", sdp.Type.String()))
h.lg.Info("tranceived message", zap.String("type", m.Type))
}
}
+101
View File
@@ -0,0 +1,101 @@
package server
import (
"encoding/json"
"github.com/pion/webrtc/v3"
)
type SignalingMessageType = string
const (
MessageTypeIceCandidate SignalingMessageType = "new-ice-candidate"
MessageTypeAnswer SignalingMessageType = "answer"
MessageTypeOffer SignalingMessageType = "offer"
)
// INCOMING
type IncomingSignalingMessage struct {
Type SignalingMessageType `json:"type"`
Data json.RawMessage `json:"data"`
}
type IceCandidateMessage struct {
*IncomingSignalingMessage
Candidate webrtc.ICECandidateInit
}
type OfferMessage struct {
*IncomingSignalingMessage
Offer webrtc.SessionDescription
}
type AnswerMessage struct {
*IncomingSignalingMessage
Answer webrtc.SessionDescription
}
func (m *IncomingSignalingMessage) IsIceCandidateMessage() bool {
return m.Type == MessageTypeIceCandidate
}
func (m *IncomingSignalingMessage) IsAnswerMessage() bool {
return m.Type == MessageTypeAnswer
}
func (m *IncomingSignalingMessage) IsOfferMessage() bool {
return m.Type == MessageTypeOffer
}
func (m *IncomingSignalingMessage) ToIceCandidateMessage() (*IceCandidateMessage, error) {
nm := IceCandidateMessage{
IncomingSignalingMessage: m,
}
return &nm, json.Unmarshal(m.Data, &nm.Candidate)
}
func (m *IncomingSignalingMessage) ToAnswerMessage() (*AnswerMessage, error) {
nm := AnswerMessage{
IncomingSignalingMessage: m,
}
return &nm, json.Unmarshal(m.Data, &nm.Answer)
}
func (m *IncomingSignalingMessage) ToOfferMessage() (*OfferMessage, error) {
nm := OfferMessage{
IncomingSignalingMessage: m,
}
return &nm, json.Unmarshal(m.Data, &nm.Offer)
}
// OUTGOING
type OutgoingSignalingMessage struct {
Type SignalingMessageType `json:"type"`
Data any `json:"data"`
}
func NewIceCandidateMessage(candidate webrtc.ICECandidate) *OutgoingSignalingMessage {
return &OutgoingSignalingMessage{
Type: MessageTypeIceCandidate,
Data: candidate,
}
}
func NewAnswerMessage(answer *webrtc.SessionDescription) *OutgoingSignalingMessage {
return &OutgoingSignalingMessage{
Type: MessageTypeAnswer,
Data: answer,
}
}
func NewOfferMessage(offer *webrtc.SessionDescription) *OutgoingSignalingMessage {
return &OutgoingSignalingMessage{
Type: MessageTypeOffer,
Data: offer,
}
}
+1 -1
View File
@@ -47,7 +47,7 @@ func CreateAudioPipelineSrc(dst StreamElement) (*SrcPipeline, error) {
sink := elems[len(elems)-1]
for name, value := range dst.Properties {
sink.SetProperty(name, value)
sink.Set(name, value)
}
// Add the elements to the pipeline and link them
+47 -19
View File
@@ -221,6 +221,13 @@ func (wh *WebrtcHandler) createPeerHandle(rctx context.Context, sh *server.Signa
// create a context for this handle
hctx, hcancel := context.WithCancel(rctx)
// sent the candidate out when where is one
peerConnection.OnICECandidate(func(i *webrtc.ICECandidate) {
if i != nil {
sh.Trcv <- server.NewIceCandidateMessage(*i)
}
})
// Set a handler for when a new remote track starts, this handler creates a gstreamer pipeline
// for the given codec
peerConnection.OnTrack(func(track *webrtc.TrackRemote, receiver *webrtc.RTPReceiver) {
@@ -263,11 +270,9 @@ func (wh *WebrtcHandler) createPeerHandle(rctx context.Context, sh *server.Signa
}()
pipeline, err := streamer.CreateAudioPipelineSrc(streamer.StreamElement{
Kind: cfg.Sink,
Properties: map[string]interface{}{
"device": cfg.Device,
},
Queue: cfg.Queue,
Kind: cfg.Sink,
Properties: properties,
Queue: cfg.Queue,
})
if err != nil {
wh.lg.Error("failed to create src pipeline", zap.Error(err))
@@ -358,7 +363,7 @@ func (wh *WebrtcHandler) createPeerHandle(rctx context.Context, sh *server.Signa
}
// Create channel that is blocked until ICE Gathering is complete
gatherComplete := webrtc.GatheringCompletePromise(peerConnection)
//gatherComplete := webrtc.GatheringCompletePromise(peerConnection)
// Sets the LocalDescription, and starts our UDP listeners
err = peerConnection.SetLocalDescription(answer)
@@ -366,23 +371,23 @@ func (wh *WebrtcHandler) createPeerHandle(rctx context.Context, sh *server.Signa
return err
}
<-gatherComplete
//<-gatherComplete
// Send the answer
sh.Trcv <- *peerConnection.LocalDescription()
sh.Trcv <- server.NewAnswerMessage(peerConnection.LocalDescription())
return nil
}
hndl := PeerHandle{
audioTrack: audioTrack,
videoTrack: videoTrack,
}
hndl := PeerHandle{audioTrack, videoTrack}
// add handle to list
wh.peerHandles[sh.Id] = &hndl
go func() {
for {
select {
case offer, ok := <-sh.Recv:
case m, ok := <-sh.Recv:
// return when channel is closed
if !ok {
wh.mu.Lock()
@@ -399,9 +404,35 @@ func (wh *WebrtcHandler) createPeerHandle(rctx context.Context, sh *server.Signa
return
}
err := onOfferReceived(offer)
if err != nil {
wh.lg.Error("failed to handle offer", zap.Error(err))
switch true {
case m.IsIceCandidateMessage():
pm, err := m.ToIceCandidateMessage()
if err != nil {
wh.lg.Error("failed to parse candidate", zap.Error(err))
continue
}
wh.lg.Info("received new ice candidate", zap.Any("candidate", pm.Candidate))
err = peerConnection.AddICECandidate(pm.Candidate)
if err != nil {
wh.lg.Error("failed to add ice candidate", zap.Error(err))
}
case m.IsAnswerMessage():
wh.lg.Info("received answer", zap.Any("data", m.Data))
case m.IsOfferMessage():
pm, err := m.ToOfferMessage()
if err != nil {
wh.lg.Error("failed to parse offer", zap.Error(err))
continue
}
wh.lg.Info("received offer", zap.Any("data", pm.Offer))
err = onOfferReceived(pm.Offer)
if err != nil {
wh.lg.Error("failed to handle offer", zap.Error(err))
}
}
case <-hctx.Done():
return
@@ -409,8 +440,5 @@ func (wh *WebrtcHandler) createPeerHandle(rctx context.Context, sh *server.Signa
}
}()
// add handle to list
wh.peerHandles[sh.Id] = &hndl
return nil
}