Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
7.4 kB · 294 lines
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295package 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 collectedfunc 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 workingfunc 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 = 1280const TestSrcHeight = 720const 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()}