package media import ( "bytes" "context" "encoding/json" "fmt" "io" "os" "runtime" "strings" "time" "github.com/go-gst/go-glib/glib" "github.com/go-gst/go-gst/gst" "github.com/go-gst/go-gst/gst/app" "github.com/skip2/go-qrcode" "golang.org/x/sync/errgroup" "stream.place/streamplace/pkg/aqtime" "stream.place/streamplace/pkg/log" "stream.place/streamplace/test" ) const HLSPlaylist = "stream.m3u8" // Pipe with a mechanism to keep the FDs not garbage collected func SafePipe() (*os.File, *os.File, func(), error) { r, w, err := os.Pipe() if err != nil { return nil, nil, nil, err } return r, w, func() { runtime.KeepAlive(r.Fd()) runtime.KeepAlive(w.Fd()) }, nil } // basic test to make sure gstreamer functionality is working func SelfTest(ctx context.Context) error { ctx = log.WithLogValues(ctx, "mediafunc", "SelfTest") ctx, cancel := context.WithTimeout(ctx, 5*time.Second) defer cancel() f, err := test.Files.Open("fixtures/sample-segment.mp4") if err != nil { return fmt.Errorf("failed to open test file: %w", err) } defer f.Close() bs, err := io.ReadAll(f) if err != nil { return fmt.Errorf("failed to read test file: %w", err) } pipeline, err := gst.NewPipeline("self-test") if err != nil { return fmt.Errorf("failed to create pipeline: %w", err) } srcele, err := gst.NewElementWithProperties("appsrc", map[string]interface{}{ "name": "self-test-src", }) if err != nil { return fmt.Errorf("failed to create appsrc element: %w", err) } err = pipeline.Add(srcele) if err != nil { return fmt.Errorf("failed to add appsrc to pipeline: %w", err) } sinkele, err := gst.NewElementWithProperties("appsink", map[string]interface{}{ "name": "self-test-sink", "sync": false, }) if err != nil { return fmt.Errorf("failed to create appsink element: %w", err) } err = pipeline.Add(sinkele) if err != nil { return fmt.Errorf("failed to add appsink to pipeline: %w", err) } err = srcele.Link(sinkele) if err != nil { return fmt.Errorf("failed to link appsrc to appsink: %w", err) } // pipeline, err := gst.NewPipelineFromString("appsrc name=src ! appsink name=sink") // if err != nil { // return err // } src := app.SrcFromElement(srcele) src.SetCallbacks(&app.SourceCallbacks{ NeedDataFunc: func(self *app.Source, _ uint) { buffer := gst.NewBufferWithSize(int64(len(bs))) buffer.Map(gst.MapWrite).WriteData(bs) defer buffer.Unmap() self.PushBuffer(buffer) log.Debug(ctx, "ending stream") self.EndStream() }, }) output := &bytes.Buffer{} if err != nil { return fmt.Errorf("unexpected error: %w", err) } appsink := app.SinkFromElement(sinkele) appsink.SetCallbacks(&app.SinkCallbacks{ NewSampleFunc: func(sink *app.Sink) gst.FlowReturn { sample := sink.PullSample() if sample == nil { return gst.FlowOK } // defer sample.Unref() // Retrieve the buffer from the sample. buffer := sample.GetBuffer() _, err := io.Copy(output, buffer.Reader()) if err != nil { panic(err) } return gst.FlowOK }, EOSFunc: func(sink *app.Sink) { log.Debug(ctx, "EOSFunc") cancel() }, }) go func() { if err := HandleBusMessages(ctx, pipeline); err != nil { log.Debug(ctx, "handle bus messages failed", "error", err) } cancel() }() // Start the pipeline log.Debug(ctx, "setting pipeline to playing state") err = pipeline.SetState(gst.StatePlaying) if err != nil { return fmt.Errorf("failed to set pipeline to playing state: %w", err) } <-ctx.Done() if len(output.Bytes()) < 1 { return fmt.Errorf("got a zero-byte buffer from SelfTest") } err = pipeline.BlockSetState(gst.StateNull) if err != nil { return fmt.Errorf("failed to set pipeline to null state: %w", err) } return nil } const TestSrcWidth = 1280 const TestSrcHeight = 720 const QRSize = 256 type QRData struct { Now int64 `json:"now"` } func (mm *MediaManager) TestSource(ctx context.Context, ms MediaSigner) error { mainLoop := glib.NewMainLoop(glib.MainContextDefault(), false) pipelineSlice := []string{ "h264parse name=videoparse", "compositor name=comp ! videoconvert ! video/x-raw,format=I420 ! x264enc speed-preset=ultrafast key-int-max=30 ! queue ! videoparse.", fmt.Sprintf(`videotestsrc is-live=true ! video/x-raw,format=AYUV,framerate=30/1,width=%d,height=%d ! comp.`, TestSrcWidth, TestSrcHeight), fmt.Sprintf("videobox border-alpha=0 top=-%d left=-%d name=box ! comp.", (TestSrcHeight/2)-(QRSize/2), (TestSrcWidth/2)-(QRSize/2)), "appsrc name=pngsrc ! pngdec ! videoconvert ! videorate ! video/x-raw,format=AYUV,framerate=1/1 ! box.", "appsrc name=timetext ! pngdec ! videoconvert ! videorate ! video/x-raw,format=AYUV,framerate=1/1 ! comp.", "audiotestsrc ! audioconvert ! opusenc inband-fec=true perfect-timestamp=true bitrate=128000 ! queue ! opusparse name=audioparse", } pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) if err != nil { return fmt.Errorf("error creating TestSource pipeline: %w", err) } videoparse, err := pipeline.GetElementByName("videoparse") if err != nil { return err } audioparse, err := pipeline.GetElementByName("audioparse") if err != nil { return err } signer, err := mm.SegmentAndSignElem(ctx, ms) if err != nil { return err } if err := pipeline.Add(signer); err != nil { return err } err = videoparse.Link(signer) if err != nil { return fmt.Errorf("link to signer failed: %w", err) } err = audioparse.Link(signer) if err != nil { return fmt.Errorf("link to signer failed: %w", err) } pngele, err := pipeline.GetElementByName("pngsrc") if err != nil { return err } src := app.SrcFromElement(pngele) src.SetCallbacks(&app.SourceCallbacks{ NeedDataFunc: func(self *app.Source, _ uint) { now := time.Now().UnixMilli() data := QRData{Now: now} bs, err := json.Marshal(data) if err != nil { panic(err) } png, err := qrcode.Encode(string(bs), qrcode.Medium, 256) if err != nil { panic(err) } buffer := gst.NewBufferWithSize(int64(len(png))) buffer.Map(gst.MapWrite).WriteData(png) defer buffer.Unmap() self.PushBuffer(buffer) }, }) tr, err := NewTextRenderer() if err != nil { return err } timetext, err := pipeline.GetElementByName("timetext") if err != nil { return err } timesrc := app.SrcFromElement(timetext) timesrc.SetCallbacks(&app.SourceCallbacks{ NeedDataFunc: func(self *app.Source, _ uint) { aqt := aqtime.FromTime(time.Now()) png, err := tr.GenerateImage(aqt.String(), "#ffffff", "#000000", 36) if err != nil { panic(err) } buffer := gst.NewBufferWithSize(int64(len(png))) buffer.Map(gst.MapWrite).WriteData(png) defer buffer.Unmap() self.PushBuffer(buffer) }, }) ctx, cancel := context.WithCancel(ctx) defer cancel() go func() { <-ctx.Done() if err := pipeline.BlockSetState(gst.StateNull); err != nil { log.Log(ctx, "failed to set pipeline state", "error", err) } mainLoop.Quit() }() go func() { if err := HandleBusMessages(ctx, pipeline); err != nil { log.Log(ctx, "pipeline error", "error", err) } cancel() }() // Start the pipeline if err := pipeline.SetState(gst.StatePlaying); err != nil { log.Log(ctx, "failed to set pipeline state", "error", err) } g, _ := errgroup.WithContext(ctx) g.Go(func() error { mainLoop.Run() log.Log(ctx, "main loop complete") return nil }) return g.Wait() }