first working

This commit is contained in:
kaedwen committed 2023-07-04 19:53:26 +02:00
1 parent b7abd95296
commit 957bf18721
2528 files changed
+787399 -155

No files matched your search

+34
View File
@@ -9,6 +9,7 @@ import (
type Config struct {
Logging ConfigLogging `arg:"-"`
Stream ConfigStream `arg:"-"`
Http ConfigHTTP `arg:"-"`
}
@@ -24,12 +25,45 @@ type ConfigHTTP struct {
StaticPath *string `arg:"--http-static,env:HTTP_STATIC"`
}
type ConfigStream struct {
VideoOut ConfigVideoOutputStream `arg:"-"`
AudioOut ConfigAudioOutputStream `arg:"-"`
AudioIn ConfigAudioInputStream `arg:"-"`
}
type ConfigVideoOutputStream struct {
Source string `arg:"--video-src,env:VIDEO_SRC" default:"v4l2src"`
Device string `arg:"--video-device,env:VIDEO_DEVICE" default:"/dev/video0"`
Codec string `arg:"--video-codec,env:VIDEO_CODEC" default:"vp8"`
Height uint `arg:"--video-height,env:VIDEO_HEIGHT" default:"480"`
Width uint `arg:"--video-width,env:VIDEO_WIDTH" default:"640"`
}
type ConfigAudioOutputStream struct {
Source string `arg:"--audio-src,env:AUDIO_SRC" default:"alsasrc"`
DeviceName string `arg:"--audio-device-name,env:AUDIO_DEVICE" default:"default"`
Device *string `arg:"--audio-device,env:AUDIO_DEVICE"`
Codec string `arg:"--audio-codec,env:AUDIO_CODEC" default:"opus"`
Channels uint `arg:"--audio-channels,env:AUDIO_CHANNELS" default:"1"`
}
type ConfigAudioInputStream struct {
Sink string `arg:"--audio-src,env:AUDIO_SINK" default:"alsasink"`
DeviceName string `arg:"--audio-device-name,env:AUDIO_DEVICE" default:"default"`
Device *string `arg:"--audio-device,env:AUDIO_DEVICE"`
Codec string `arg:"--audio-codec,env:AUDIO_CODEC" default:"opus"`
Channels uint `arg:"--audio-channels,env:AUDIO_CHANNELS" default:"1"`
}
func (c *ConfigHTTP) Address() string {
return net.JoinHostPort(c.Host, fmt.Sprint(c.Port))
}
func (c *Config) MustParse() {
arg.MustParse(&c.Logging)
arg.MustParse(&c.Stream.VideoOut)
arg.MustParse(&c.Stream.AudioOut)
arg.MustParse(&c.Stream.AudioIn)
arg.MustParse(&c.Http)
arg.MustParse(c)
}
+1 -1
View File
@@ -11,8 +11,8 @@ import (
"strings"
"time"
"gitea.heinrich.blue/PHI/webrtc-gst/pkg/common"
"github.com/gin-gonic/gin"
"github.com/kaedwen/webrtc/pkg/common"
"github.com/pion/webrtc/v3"
"go.uber.org/zap"
"nhooyr.io/websocket"
+67 -19
View File
@@ -1,4 +1,4 @@
package webrtc
package streamer
import (
"github.com/pion/webrtc/v3/pkg/media"
@@ -6,10 +6,6 @@ import (
"github.com/tinyzimmer/go-gst/gst/app"
)
func init() {
gst.Init(nil)
}
func setCallback(sink *app.Sink, ch chan media.Sample) {
sink.SetCallbacks(&app.SinkCallbacks{
NewSampleFunc: func(sink *app.Sink) gst.FlowReturn {
@@ -38,24 +34,49 @@ func setCallback(sink *app.Sink, ch chan media.Sample) {
})
}
func CreateVideoPipeline(s string) (*gst.Pipeline, <-chan media.Sample, error) {
func CreateVideoPipelineSink(s StreamElement) (*gst.Pipeline, <-chan media.Sample, error) {
// Create a pipeline
pipeline, err := gst.NewPipeline("")
pipeline, err := gst.NewPipeline("pion-video-pipeline")
if err != nil {
return nil, nil, err
}
elems := make([]*gst.Element, 0)
// Create the src
src, err := gst.NewElement(s)
src, err := gst.NewElement(s.Kind)
if err != nil {
return nil, nil, err
}
elems = append(elems, src)
for name, value := range s.Properties {
src.Set(name, value)
}
if s.Caps != nil {
filter, err := gst.NewElement("capsfilter")
if err != nil {
return nil, nil, err
}
filter.Set("caps", s.ToGstCaps())
elems = append(elems, filter)
}
// just to be on the save side
conv, err := gst.NewElement("videoconvert")
if err != nil {
return nil, nil, err
}
elems = append(elems, conv)
// Create the enc
enc, err := gst.NewElement("vp8enc")
if err != nil {
return nil, nil, err
}
elems = append(elems, enc)
enc.SetProperty("error-resilient", "partitions")
enc.SetProperty("keyframe-max-dist", 10)
@@ -64,56 +85,83 @@ func CreateVideoPipeline(s string) (*gst.Pipeline, <-chan media.Sample, error) {
enc.SetProperty("deadline", 1)
// Create the sink
sink, err := app.NewAppSink()
appsink, err := app.NewAppSink()
if err != nil {
return nil, nil, err
}
elems = append(elems, appsink.Element)
ch := make(chan media.Sample, 100)
setCallback(sink, ch)
setCallback(appsink, ch)
// Add the elements to the pipeline
pipeline.AddMany(src, enc, sink.Element)
pipeline.AddMany(elems...)
// link the elements
gst.ElementLinkMany(src, enc, sink.Element)
gst.ElementLinkMany(elems...)
return pipeline, ch, nil
}
func CreateAudioPipeline(s string) (*gst.Pipeline, <-chan media.Sample, error) {
func CreateAudioPipelineSink(s StreamElement) (*gst.Pipeline, <-chan media.Sample, error) {
// Create a pipeline
pipeline, err := gst.NewPipeline("")
pipeline, err := gst.NewPipeline("pion-audio-pipeline")
if err != nil {
return nil, nil, err
}
elems := make([]*gst.Element, 0)
// Create the src
src, err := gst.NewElement(s)
src, err := gst.NewElement(s.Kind)
if err != nil {
return nil, nil, err
}
elems = append(elems, src)
for name, value := range s.Properties {
src.Set(name, value)
}
if s.Caps != nil {
filter, err := gst.NewElement("capsfilter")
if err != nil {
return nil, nil, err
}
filter.Set("caps", s.ToGstCaps())
elems = append(elems, filter)
}
// just to be on the save side
conv, err := gst.NewElement("audioconvert")
if err != nil {
return nil, nil, err
}
elems = append(elems, conv)
// Create the enc
enc, err := gst.NewElement("opusenc")
if err != nil {
return nil, nil, err
}
elems = append(elems, enc)
// Create the sink
sink, err := app.NewAppSink()
appsink, err := app.NewAppSink()
if err != nil {
return nil, nil, err
}
elems = append(elems, appsink.Element)
ch := make(chan media.Sample, 100)
setCallback(sink, ch)
setCallback(appsink, ch)
// Add the elements to the pipeline
pipeline.AddMany(src, enc, sink.Element)
pipeline.AddMany(elems...)
// link the elements
gst.ElementLinkMany(src, enc, sink.Element)
gst.ElementLinkMany(elems...)
return pipeline, ch, nil
}
+79
View File
@@ -0,0 +1,79 @@
package streamer
import (
"fmt"
"github.com/tinyzimmer/go-gst/gst"
"github.com/tinyzimmer/go-gst/gst/app"
)
type SrcPipeline struct {
*gst.Pipeline
src *app.Source
}
func CreateAudioPipelineSrc(dst StreamElement) (*SrcPipeline, error) {
// Create a pipeline
pipeline, err := gst.NewPipeline("")
if err != nil {
return nil, err
}
elems := make([]*gst.Element, 0)
// Create the src
appsrc, err := app.NewAppSrc()
if err != nil {
return nil, err
}
elems = append(elems, appsrc.Element)
appsrc.SetFormat(gst.FormatTime)
appsrc.SetDoTimestamp(true)
appsrc.SetLive(true)
// Create the opus decoder
decCodec, err := gst.NewElement("rtpopusdepay")
if err != nil {
return nil, err
}
elems = append(elems, decCodec)
// Create the bin decoder
decBin, err := gst.NewElement("decodebin")
if err != nil {
return nil, err
}
elems = append(elems, decBin)
// Create the sink
sink, err := gst.NewElement(dst.Kind)
if err != nil {
return nil, err
}
elems = append(elems, sink)
for name, value := range dst.Properties {
sink.SetProperty(name, value)
}
// Add the elements to the pipeline
pipeline.AddMany(elems...)
// link the elements
gst.ElementLinkMany(elems...)
return &SrcPipeline{
Pipeline: pipeline,
src: appsrc,
}, nil
}
func (p *SrcPipeline) Push(data []byte) error {
err := p.src.PushBuffer(gst.NewBufferFromBytes(data))
if err != gst.FlowOK {
return fmt.Errorf("failed to send bytes to gst pipeline - %s", err.String())
}
return nil
}
+88
View File
@@ -0,0 +1,88 @@
package streamer
import (
"fmt"
"reflect"
"strings"
"time"
"github.com/tinyzimmer/go-gst/gst"
"github.com/tinyzimmer/go-gst/gst/app"
"go.uber.org/zap"
)
const capsTagName = "caps"
func init() {
gst.Init(nil)
}
type StreamElementCaps struct {
Mime string `caps:"mime"`
Height uint `caps:"height"`
Width uint `caps:"width"`
Format string `caps:"format"`
Channels uint `caps:"channels"`
Rate uint `caps:"rate"`
}
type StreamElement struct {
Kind string
Properties map[string]interface{}
Caps *StreamElementCaps
}
func (se StreamElement) ToGstCaps() *gst.Caps {
if se.Caps == nil {
return nil
}
t := reflect.ValueOf(*se.Caps)
var mime string
var caps []string
for i := 0; i < t.NumField(); i++ {
valueField := t.Field(i)
tagField := t.Type().Field(i).Tag.Get(capsTagName)
if tagField == "mime" {
mime = valueField.String()
} else {
if !valueField.IsZero() {
caps = append(caps, fmt.Sprintf("%s=%v", tagField, valueField.Interface()))
}
}
}
return gst.NewCapsFromString(strings.Join(append([]string{mime}, caps...), ","))
}
func handleMessage(msg *gst.Message) error {
switch msg.Type() {
case gst.MessageEOS:
return app.ErrEOS
case gst.MessageError:
return msg.ParseError()
}
return nil
}
func LoopBus(lg *zap.Logger, pipeline *gst.Pipeline) {
// Retrieve the bus from the pipeline
bus := pipeline.GetPipelineBus()
// Loop over messsages from the pipeline
go func() {
for {
msg := bus.TimedPop(time.Duration(-1))
if msg == nil {
return
}
if err := handleMessage(msg); err != nil {
lg.Error("failed to handle message", zap.Error(err))
}
}
}()
}
+97 -8
View File
@@ -3,8 +3,12 @@ package webrtc
import (
"context"
"sync"
"time"
"gitea.heinrich.blue/PHI/webrtc-gst/pkg/server"
"github.com/kaedwen/webrtc/pkg/common"
"github.com/kaedwen/webrtc/pkg/server"
"github.com/kaedwen/webrtc/pkg/streamer"
"github.com/pion/rtcp"
"github.com/pion/webrtc/v3"
"github.com/pion/webrtc/v3/pkg/media"
"github.com/tinyzimmer/go-gst/gst"
@@ -26,19 +30,19 @@ type PeerHandle struct {
videoTrack *webrtc.TrackLocalStaticSample
}
func NewWebrtcHandler(ctx context.Context, lg *zap.Logger, ch <-chan *server.SignalingHandle) error {
func NewWebrtcHandler(ctx context.Context, lg *zap.Logger, cfg *common.ConfigStream, ch <-chan *server.SignalingHandle) error {
wh := WebrtcHandler{
lg: lg,
mu: &sync.Mutex{},
peerHandles: make(map[string]*PeerHandle, 0),
}
err := wh.handleAudioSamples(ctx, "audiotestsrc")
err := wh.handleAudioSamples(ctx, &cfg.AudioOut)
if err != nil {
return err
}
err = wh.handleVideoSamples(ctx, "videotestsrc")
err = wh.handleVideoSamples(ctx, &cfg.VideoOut)
if err != nil {
return err
}
@@ -99,14 +103,35 @@ func (wh *WebrtcHandler) stopPipelines() {
}
func (wh *WebrtcHandler) handleAudioSamples(ctx context.Context, src string) error {
func (wh *WebrtcHandler) handleAudioSamples(ctx context.Context, cfg *common.ConfigAudioOutputStream) error {
properties := map[string]interface{}{}
if cfg.Source == "alsasrc" {
if cfg.Device != nil {
properties["device"] = *cfg.Device
} else {
properties["device-name"] = cfg.DeviceName
}
}
src := streamer.StreamElement{
Kind: cfg.Source,
Properties: properties,
Caps: &streamer.StreamElementCaps{
Mime: "audio/x-raw",
Channels: cfg.Channels,
Rate: 48000,
},
}
var err error
var audioCh <-chan media.Sample
wh.audioPipeline, audioCh, err = CreateAudioPipeline(src)
wh.audioPipeline, audioCh, err = streamer.CreateAudioPipelineSink(src)
if err != nil {
return err
}
streamer.LoopBus(wh.lg.With(zap.String("sub-context", "audio")), wh.audioPipeline)
go func() {
wh.lg.Info("wait for audio sample")
for {
@@ -127,14 +152,29 @@ func (wh *WebrtcHandler) handleAudioSamples(ctx context.Context, src string) err
return nil
}
func (wh *WebrtcHandler) handleVideoSamples(ctx context.Context, src string) error {
func (wh *WebrtcHandler) handleVideoSamples(ctx context.Context, cfg *common.ConfigVideoOutputStream) error {
src := streamer.StreamElement{
Kind: cfg.Source,
Properties: map[string]interface{}{
"device": cfg.Device,
},
Caps: &streamer.StreamElementCaps{
Mime: "video/x-raw",
Format: "YUY2",
Width: cfg.Width,
Height: cfg.Height,
},
}
var err error
var videoCh <-chan media.Sample
wh.videoPipeline, videoCh, err = CreateVideoPipeline(src)
wh.videoPipeline, videoCh, err = streamer.CreateVideoPipelineSink(src)
if err != nil {
return err
}
streamer.LoopBus(wh.lg.With(zap.String("sub-context", "video")), wh.videoPipeline)
go func() {
wh.lg.Info("wait for video sample")
for {
@@ -170,6 +210,55 @@ func (wh *WebrtcHandler) createPeerHandle(ctx context.Context, sh *server.Signal
return err
}
// 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) {
wh.lg.Info("received track", zap.String("kind", track.Kind().String()), zap.String("codec", track.Codec().MimeType))
if track.Codec().MimeType != "audio/opus" {
wh.lg.Error("mimetype not supported", zap.String("mime", track.Codec().MimeType))
return
}
// Send a PLI on an interval so that the publisher is pushing a keyframe every rtcpPLIInterval
go func() {
ticker := time.NewTicker(time.Second * 3)
for {
select {
case <-ticker.C:
if err := peerConnection.WriteRTCP([]rtcp.Packet{&rtcp.PictureLossIndication{MediaSSRC: uint32(track.SSRC())}}); err != nil {
wh.lg.Error("failed to send PLI", zap.Error(err))
}
case <-ctx.Done():
return
}
}
}()
pipeline, err := streamer.CreateAudioPipelineSrc(streamer.StreamElement{
Kind: "autoaudiosink",
})
if err != nil {
wh.lg.Error("failed to create src pipeline", zap.Error(err))
}
pipeline.Start()
buf := make([]byte, 1400)
for {
i, _, readErr := track.Read(buf)
if readErr != nil {
wh.lg.Error("read failed", zap.Error(err))
}
if i > 0 {
pipeline.Push(buf[:i])
if err != nil {
wh.lg.Error("push failed", zap.Error(err))
}
}
}
})
// Set the handler for ICE connection state
// This will notify you when the peer has connected/disconnected
peerConnection.OnICEConnectionStateChange(func(connectionState webrtc.ICEConnectionState) {