add builder

This commit is contained in:
kaedwen committed 2024-04-01 15:40:23 +02:00
1 parent 424be67dda
commit ac2e689131
4 files changed
+142 -117

No files matched your search

+26
View File
@@ -0,0 +1,26 @@
package streamer
import "fmt"
type PipelineBuilder struct {
parts []string
}
func NewPipelineBuilder() *PipelineBuilder {
return &PipelineBuilder{}
}
func (pb *PipelineBuilder) Add(element string) *PipelineBuilder {
pb.parts = append(pb.parts, element)
return pb
}
func (pb *PipelineBuilder) AddWithProperties(element string, properties map[string]any) *PipelineBuilder {
pl := make([]string, 0, len(properties))
for k, v := range properties {
pl = append(pl, fmt.Sprintf("%s=%v", k, v))
}
pb.parts = append(pb.parts, fmt.Sprint(element, pl))
return pb
}
+70 -107
View File
@@ -38,7 +38,10 @@ func setCallback(sink *app.Sink, ch chan media.Sample) {
}) })
} }
func CreateVideoPipelineSinkWithLaunch(s StreamElement, _ *zap.Logger) (*gst.Pipeline, <-chan media.Sample, error) { func CreateVideoPipelineSinkWithLaunch(_ *zap.Logger, s StreamElement) (*gst.Pipeline, <-chan media.Sample, error) {
pb := NewPipelineBuilder()
pb.AddWithProperties("v4l2src", map[string]any{"device", s.})
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`) 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 { if err != nil {
return nil, nil, err return nil, nil, err
@@ -56,130 +59,92 @@ func CreateVideoPipelineSinkWithLaunch(s StreamElement, _ *zap.Logger) (*gst.Pip
return pipeline, ch, nil return pipeline, ch, nil
} }
func CreateVideoPipelineSink(s StreamElement, lg *zap.Logger) (*gst.Pipeline, <-chan media.Sample, error) { func CreateVideoPipelineSink(lg *zap.Logger, s StreamElement) (*gst.Pipeline, <-chan media.Sample, error) {
// Create a pipeline // Create a pipeline
pipeline, err := gst.NewPipeline("pion-video-pipeline") pipeline, err := gst.NewPipeline("pion-video-pipeline")
if err != nil { if err != nil {
return nil, nil, err return nil, nil, err
} }
elems := make([]*gst.Element, 0)
// Create the src // Create the src
src, err := gst.NewElement(s.Kind) src, err := gst.NewElement("v4l2src")
if err != nil { if err != nil {
return nil, nil, err return nil, nil, err
} }
elems = append(elems, src)
for name, value := range s.Properties { src.Set("device", "/dev/video4")
src.Set(name, value)
}
if s.Caps != nil { src_c := gst.NewEmptySimpleCaps("video/x-raw")
filter, err := gst.NewElement("capsfilter") src_c.SetValue("width", int(320))
if err != nil { src_c.SetValue("height", int(240))
return nil, nil, err lg.Info("capsfilter", zap.String("caps", src_c.String()))
}
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 // just to be on the save side
conv, err := gst.NewElement("videoconvert") conv, err := gst.NewElement("videoconvert")
if err != nil { if err != nil {
return nil, nil, err return nil, nil, err
} }
elems = append(elems, conv)
if s.Queue { // add a queue
// add a queue queue, err := gst.NewElement("queue")
queue, err := gst.NewElement("queue") if err != nil {
if err != nil { return nil, nil, err
return nil, nil, err
}
elems = append(elems, queue)
} }
switch s.Codec { enc_in_c := gst.NewEmptySimpleCaps("video/x-raw")
case common.VP8: enc_in_c.SetValue("format", "I420")
// Create the enc lg.Info("capsfilter", zap.String("caps", enc_in_c.String()))
enc, err := gst.NewElement("vp8enc")
if err != nil {
return nil, nil, err
}
elems = append(elems, enc)
enc.SetProperty("error-resilient", "partitions") // Create the enc
enc.SetProperty("keyframe-max-dist", 10) enc, err := gst.NewElement("x264enc")
enc.SetProperty("auto-alt-ref", true) if err != nil {
enc.SetProperty("cpu-used", 5) return nil, nil, err
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 {
return nil, nil, err
}
elems = append(elems, enc)
enc.SetProperty("speed-preset", "ultrafast")
enc.SetProperty("tune", "zerolatency")
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)
} }
enc.SetProperty("speed-preset", "ultrafast")
enc.SetProperty("tune", "zerolatency")
enc.SetProperty("key-int-max", 20)
//enc.SetProperty("bitrate", 300)
enc_out_c := gst.NewEmptySimpleCaps("video/x-h264")
enc_out_c.SetValue("stream-format", "byte-stream")
lg.Info("capsfilter", zap.String("caps", enc_out_c.String()))
// Create the sink // Create the sink
appsink, err := app.NewAppSink() appsink, err := app.NewAppSink()
if err != nil { if err != nil {
return nil, nil, err return nil, nil, err
} }
elems = append(elems, appsink.Element)
ch := make(chan media.Sample, 100) ch := make(chan media.Sample, 100)
setCallback(appsink, ch) setCallback(appsink, ch)
// Add the elements to the pipeline // Add the elements to the pipeline
pipeline.AddMany(elems...) err = pipeline.AddMany(src, conv, queue, enc, appsink.Element)
if err != nil {
return nil, nil, err
}
// link the elements // link the elements
gst.ElementLinkMany(elems...) err = src.LinkFiltered(conv, src_c)
if err != nil {
return nil, nil, err
}
err = conv.Link(queue)
if err != nil {
return nil, nil, err
}
err = queue.LinkFiltered(enc, enc_in_c)
if err != nil {
return nil, nil, err
}
err = enc.LinkFiltered(appsink.Element, enc_out_c)
if err != nil {
return nil, nil, err
}
return pipeline, ch, nil return pipeline, ch, nil
} }
@@ -191,32 +156,24 @@ func CreateAudioPipelineSink(s StreamElement, lg *zap.Logger) (*gst.Pipeline, <-
return nil, nil, err return nil, nil, err
} }
elems := make([]*gst.Element, 0) elems := make(ElementList, 0)
// Create the src // Create the src
src, err := gst.NewElement(s.Kind) src, err := gst.NewElement(s.Kind)
if err != nil { if err != nil {
return nil, nil, err return nil, nil, err
} }
elems = append(elems, src)
for name, value := range s.Properties { for name, value := range s.Properties {
src.Set(name, value) src.Set(name, value)
} }
if s.Caps != nil { if s.Caps != nil {
filter, err := gst.NewElement("capsfilter")
if err != nil {
return nil, nil, err
}
elems = append(elems, filter)
c := s.Caps.Build() c := s.Caps.Build()
lg.Info("capsfilter", zap.String("caps", c.String())) lg.Info("capsfilter", zap.String("caps", c.String()))
err = filter.SetProperty("caps", c) elems = append(elems, Element{src, c})
if err != nil { } else {
return nil, nil, err elems = append(elems, Element{src, nil})
}
} }
// just to be on the save side // just to be on the save side
@@ -224,7 +181,7 @@ func CreateAudioPipelineSink(s StreamElement, lg *zap.Logger) (*gst.Pipeline, <-
if err != nil { if err != nil {
return nil, nil, err return nil, nil, err
} }
elems = append(elems, conv) elems = append(elems, Element{conv, nil})
if s.Queue { if s.Queue {
// add a queue // add a queue
@@ -232,7 +189,7 @@ func CreateAudioPipelineSink(s StreamElement, lg *zap.Logger) (*gst.Pipeline, <-
if err != nil { if err != nil {
return nil, nil, err return nil, nil, err
} }
elems = append(elems, queue) elems = append(elems, Element{queue, nil})
} }
switch s.Codec { switch s.Codec {
@@ -242,7 +199,7 @@ func CreateAudioPipelineSink(s StreamElement, lg *zap.Logger) (*gst.Pipeline, <-
if err != nil { if err != nil {
return nil, nil, err return nil, nil, err
} }
elems = append(elems, enc) elems = append(elems, Element{enc, nil})
default: default:
return nil, nil, fmt.Errorf("unsupported audio codec given - %s", s.Codec) return nil, nil, fmt.Errorf("unsupported audio codec given - %s", s.Codec)
} }
@@ -252,16 +209,22 @@ func CreateAudioPipelineSink(s StreamElement, lg *zap.Logger) (*gst.Pipeline, <-
if err != nil { if err != nil {
return nil, nil, err return nil, nil, err
} }
elems = append(elems, appsink.Element) elems = append(elems, Element{appsink.Element, nil})
ch := make(chan media.Sample, 100) ch := make(chan media.Sample, 100)
setCallback(appsink, ch) setCallback(appsink, ch)
// Add the elements to the pipeline // Add the elements to the pipeline
pipeline.AddMany(elems...) err = pipeline.AddMany(elems.List()...)
if err != nil {
return nil, nil, err
}
// link the elements // link the elements
gst.ElementLinkMany(elems...) err = elems.Link()
if err != nil {
return nil, nil, err
}
return pipeline, ch, nil return pipeline, ch, nil
} }
+36
View File
@@ -10,6 +10,42 @@ func init() {
gst.Init(nil) gst.Init(nil)
} }
type Element struct {
*gst.Element
filter *gst.Caps
}
type ElementList []Element
func (elems ElementList) Link() error {
for idx, elem := range elems {
if idx == 0 {
// skip the first one as the loop always links previous to current
continue
}
pe := elems[idx-1]
if elem.filter != nil {
if err := pe.LinkFiltered(elem.Element, pe.filter); err != nil {
return err
}
} else {
if err := pe.Link(elem.Element); err != nil {
return err
}
}
}
return nil
}
func (elems ElementList) List() []*gst.Element {
l := make([]*gst.Element, 0, len(elems))
for _, e := range elems {
l = append(l, e.Element)
}
return l
}
func handleMessage(msg *gst.Message) error { func handleMessage(msg *gst.Message) error {
switch msg.Type() { switch msg.Type() {
case gst.MessageEOS: case gst.MessageEOS:
+10 -10
View File
@@ -155,19 +155,19 @@ func (wh *WebrtcHandler) handleAudioSamples(ctx context.Context, cfg *common.Con
} }
func (wh *WebrtcHandler) handleVideoSamples(ctx context.Context, cfg *common.ConfigVideoSourceStream) error { func (wh *WebrtcHandler) handleVideoSamples(ctx context.Context, cfg *common.ConfigVideoSourceStream) error {
src := streamer.StreamElement{ // src := streamer.StreamElement{
Kind: cfg.Source, // Kind: cfg.Source,
Properties: map[string]interface{}{ // Properties: map[string]interface{}{
"device": cfg.Device, // "device": cfg.Device,
}, // },
Caps: streamer.NewCapsBuilder("video/x-raw").Format("YUY2").Width(int(cfg.Width)).Height(int(cfg.Height)), // Caps: streamer.NewCapsBuilder("video/x-raw").Format("YUY2").Width(int(cfg.Width)).Height(int(cfg.Height)),
Queue: cfg.Queue, // Queue: cfg.Queue,
Codec: cfg.Codec, // Codec: cfg.Codec,
} // }
var err error var err error
var videoCh <-chan media.Sample var videoCh <-chan media.Sample
wh.videoPipeline, videoCh, err = streamer.CreateVideoPipelineSinkWithLaunch(src, wh.lg) wh.videoPipeline, videoCh, err = streamer.CreateVideoPipelineSinkWithLaunch(wh.lg)
if err != nil { if err != nil {
return err return err
} }