with launch

This commit is contained in:
kaedwen committed 2024-03-29 16:00:35 +01:00
1 parent 555fce9f10
commit 424be67dda
113 files changed
+18959 -109

No files matched your search

+65 -8
View File
@@ -7,6 +7,7 @@ import (
"github.com/go-gst/go-gst/gst/app"
"github.com/kaedwen/webrtc/pkg/common"
"github.com/pion/webrtc/v3/pkg/media"
"go.uber.org/zap"
)
func setCallback(sink *app.Sink, ch chan media.Sample) {
@@ -37,7 +38,25 @@ func setCallback(sink *app.Sink, ch chan media.Sample) {
})
}
func CreateVideoPipelineSink(s StreamElement) (*gst.Pipeline, <-chan media.Sample, error) {
func CreateVideoPipelineSinkWithLaunch(s StreamElement, _ *zap.Logger) (*gst.Pipeline, <-chan media.Sample, error) {
pipeline, err := gst.NewPipelineFromString(`v4l2src device=/dev/video4 ! capsfilter caps="video/x-raw,width=(int)320,height=(int)240" ! videoconvert ! queue ! capsfilter caps="video/x-raw,format=(string)I420" ! x264enc speed-preset=ultrafast tune=zerolatency key-int-max=20 ! capsfilter caps="video/x-h264,stream-format=(string)byte-stream" ! appsink name=appsink`)
if err != nil {
return nil, nil, err
}
elem, err := pipeline.GetElementByName("appsink")
if err != nil {
return nil, nil, err
}
appsink := app.SinkFromElement(elem)
ch := make(chan media.Sample, 100)
setCallback(appsink, ch)
return pipeline, ch, nil
}
func CreateVideoPipelineSink(s StreamElement, lg *zap.Logger) (*gst.Pipeline, <-chan media.Sample, error) {
// Create a pipeline
pipeline, err := gst.NewPipeline("pion-video-pipeline")
if err != nil {
@@ -62,9 +81,14 @@ func CreateVideoPipelineSink(s StreamElement) (*gst.Pipeline, <-chan media.Sampl
if err != nil {
return nil, nil, err
}
filter.Set("caps", s.ToGstCaps())
elems = append(elems, filter)
c := s.Caps.Build()
lg.Info("capsfilter", zap.String("caps", c.String()))
err = filter.SetProperty("caps", c)
if err != nil {
return nil, nil, err
}
}
// just to be on the save side
@@ -98,6 +122,20 @@ func CreateVideoPipelineSink(s StreamElement) (*gst.Pipeline, <-chan media.Sampl
enc.SetProperty("cpu-used", 5)
enc.SetProperty("deadline", 1)
case common.H264:
filterIn, err := gst.NewElement("capsfilter")
if err != nil {
return nil, nil, err
}
elems = append(elems, filterIn)
cIn := gst.NewEmptySimpleCaps("video/x-raw")
cIn.SetValue("format", "I420")
lg.Info("capsfilter", zap.String("caps", cIn.String()))
err = filterIn.SetProperty("caps", cIn)
if err != nil {
return nil, nil, err
}
// Create the enc
enc, err := gst.NewElement("x264enc")
if err != nil {
@@ -107,8 +145,22 @@ func CreateVideoPipelineSink(s StreamElement) (*gst.Pipeline, <-chan media.Sampl
enc.SetProperty("speed-preset", "ultrafast")
enc.SetProperty("tune", "zerolatency")
enc.SetProperty("key-int-max", 2)
enc.SetProperty("bitrate", 300)
enc.SetProperty("key-int-max", 20)
//enc.SetProperty("bitrate", 300)
filterOut, err := gst.NewElement("capsfilter")
if err != nil {
return nil, nil, err
}
elems = append(elems, filterOut)
cOut := gst.NewEmptySimpleCaps("video/x-h264")
cOut.SetValue("stream-format", "byte-stream")
lg.Info("capsfilter", zap.String("caps", cOut.String()))
err = filterOut.SetProperty("caps", cOut)
if err != nil {
return nil, nil, err
}
default:
return nil, nil, fmt.Errorf("unsupported video codec given - %s", s.Codec)
}
@@ -132,7 +184,7 @@ func CreateVideoPipelineSink(s StreamElement) (*gst.Pipeline, <-chan media.Sampl
return pipeline, ch, nil
}
func CreateAudioPipelineSink(s StreamElement) (*gst.Pipeline, <-chan media.Sample, error) {
func CreateAudioPipelineSink(s StreamElement, lg *zap.Logger) (*gst.Pipeline, <-chan media.Sample, error) {
// Create a pipeline
pipeline, err := gst.NewPipeline("pion-audio-pipeline")
if err != nil {
@@ -157,9 +209,14 @@ func CreateAudioPipelineSink(s StreamElement) (*gst.Pipeline, <-chan media.Sampl
if err != nil {
return nil, nil, err
}
filter.Set("caps", s.ToGstCaps())
elems = append(elems, filter)
c := s.Caps.Build()
lg.Info("capsfilter", zap.String("caps", c.String()))
err = filter.SetProperty("caps", c)
if err != nil {
return nil, nil, err
}
}
// just to be on the save side
-50
View File
@@ -1,66 +1,16 @@
package streamer
import (
"fmt"
"reflect"
"strings"
"github.com/go-gst/go-gst/gst"
"github.com/go-gst/go-gst/gst/app"
"github.com/kaedwen/webrtc/pkg/common"
"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
Codec common.StreamCodec
Properties map[string]interface{}
Caps *StreamElementCaps
Queue bool
}
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
+83
View File
@@ -0,0 +1,83 @@
package streamer
import (
"github.com/go-gst/go-gst/gst"
"github.com/kaedwen/webrtc/pkg/common"
)
type StreamElement struct {
Kind string
Codec common.StreamCodec
Properties map[string]interface{}
Caps *CapsBuilder
Queue bool
}
type CapsBuilder struct {
caps Caps
}
type Caps struct {
mime string
values []CapsValue
}
type CapsValue struct {
name string
value any
}
func NewCapsBuilder(mime string) *CapsBuilder {
return &CapsBuilder{
caps: Caps{mime, nil},
}
}
func (b *CapsBuilder) Height(v int) *CapsBuilder {
b.caps.values = append(b.caps.values, CapsValue{
name: "height",
value: v,
})
return b
}
func (b *CapsBuilder) Width(v int) *CapsBuilder {
b.caps.values = append(b.caps.values, CapsValue{
name: "width",
value: v,
})
return b
}
func (b *CapsBuilder) Channels(v int) *CapsBuilder {
b.caps.values = append(b.caps.values, CapsValue{
name: "channels",
value: v,
})
return b
}
func (b *CapsBuilder) Format(v string) *CapsBuilder {
b.caps.values = append(b.caps.values, CapsValue{
name: "format",
value: v,
})
return b
}
func (b *CapsBuilder) Rate(v int) *CapsBuilder {
b.caps.values = append(b.caps.values, CapsValue{
name: "rate",
value: v,
})
return b
}
func (b *CapsBuilder) Build() *gst.Caps {
caps := gst.NewEmptySimpleCaps(b.caps.mime)
for _, v := range b.caps.values {
caps.SetValue(v.name, v.value)
}
return caps
}
+6 -15
View File
@@ -120,18 +120,14 @@ func (wh *WebrtcHandler) handleAudioSamples(ctx context.Context, cfg *common.Con
src := streamer.StreamElement{
Kind: cfg.Source,
Properties: properties,
Caps: &streamer.StreamElementCaps{
Mime: "audio/x-raw",
Channels: cfg.Channels,
Rate: 48000,
},
Queue: cfg.Queue,
Codec: cfg.Codec,
Caps: streamer.NewCapsBuilder("audio/x-raw").Channels(int(cfg.Channels)).Rate(48000),
Queue: cfg.Queue,
Codec: cfg.Codec,
}
var err error
var audioCh <-chan media.Sample
wh.audioPipeline, audioCh, err = streamer.CreateAudioPipelineSink(src)
wh.audioPipeline, audioCh, err = streamer.CreateAudioPipelineSink(src, wh.lg)
if err != nil {
return err
}
@@ -164,19 +160,14 @@ func (wh *WebrtcHandler) handleVideoSamples(ctx context.Context, cfg *common.Con
Properties: map[string]interface{}{
"device": cfg.Device,
},
Caps: &streamer.StreamElementCaps{
Mime: "video/x-raw",
Format: "YUY2",
Width: cfg.Width,
Height: cfg.Height,
},
Caps: streamer.NewCapsBuilder("video/x-raw").Format("YUY2").Width(int(cfg.Width)).Height(int(cfg.Height)),
Queue: cfg.Queue,
Codec: cfg.Codec,
}
var err error
var videoCh <-chan media.Sample
wh.videoPipeline, videoCh, err = streamer.CreateVideoPipelineSink(src)
wh.videoPipeline, videoCh, err = streamer.CreateVideoPipelineSinkWithLaunch(src, wh.lg)
if err != nil {
return err
}