Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
10 kB · 348 lines
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349package media
import ( "context" "fmt" "io" "strings"
"github.com/go-gst/go-gst/gst" "github.com/go-gst/go-gst/gst/app" "golang.org/x/sync/errgroup" "stream.place/streamplace/pkg/log")
// Splits out video into MPEG-TS and audio into MP4 (to be recombined after transcoding)func MP4ToMPEGTSVideoMP4Audio(ctx context.Context, input io.Reader, videoOutput io.Writer, audioOutput io.Writer) error { ctx = log.WithLogValues(ctx, "func", "MP4ToMPEGTSVideoMP4Audio") pipelineStr := strings.Join([]string{ "appsrc name=appsrc ! qtdemux name=demux", "mpegtsmux name=videomux ! appsink name=videoappsink sync=false", "mp4mux name=audiomux ! appsink name=audioappsink sync=false", // config-interval=-1: SPS/PPS before every IDR, taken from the MP4's // avcC. A transcoder decodes each pushed segment with a fresh // decoder, so a segment without them ("Could not find codec // parameters") kills its session — encoders that don't repeat the // headers in-band (most) produced exactly that. // Unlimited queues: a muxer holds its first video buffers until the // audio pad has data too, and a 1s 1080p60 segment with B-frame // reordering fills the default 1s queue before the demuxer gets to // the audio — the pipeline deadlocked and produced nothing. "demux.video_0 ! h264parse config-interval=-1 ! video/x-h264,stream-format=byte-stream ! queue name=videoqueue max-size-time=0 max-size-buffers=0 max-size-bytes=0", // parsebin plugs the parser the audio actually needs (opusparse for a // WHIP/Opus ingest, aacparse for RTMP/AAC); a hard-wired opusparse // never linked for AAC and the pipeline hung on it, wedging every // transcode behind the in-flight guard. "demux.audio_0 ! parsebin ! queue name=audioqueue max-size-time=0 max-size-buffers=0 max-size-bytes=0", }, " ")
pipeline, err := gst.NewPipelineFromString(pipelineStr) if err != nil { return err }
videomux, err := pipeline.GetElementByName("videomux") if err != nil { return err } muxVideoSinkPad := videomux.GetRequestPad("sink_%d") if muxVideoSinkPad == nil { return fmt.Errorf("failed to get video sink pad") }
audiomux, err := pipeline.GetElementByName("audiomux") if err != nil { return err } muxAudioSinkPad := audiomux.GetRequestPad("audio_%u") if muxAudioSinkPad == nil { return fmt.Errorf("failed to get audio sink pad") }
videoQueue, err := pipeline.GetElementByName("videoqueue") if err != nil { return err } audioQueue, err := pipeline.GetElementByName("audioqueue") if err != nil { return err }
videoQueueSrcPad := videoQueue.GetStaticPad("src") if videoQueueSrcPad == nil { return fmt.Errorf("failed to get video queue source pad") } audioQueueSrcPad := audioQueue.GetStaticPad("src") if audioQueueSrcPad == nil { return fmt.Errorf("failed to get audio queue source pad") }
ok := videoQueueSrcPad.Link(muxVideoSinkPad) if ok != gst.PadLinkOK { return fmt.Errorf("failed to link video queue source pad to mux video sink pad: %v", ok) } ok = audioQueueSrcPad.Link(muxAudioSinkPad) if ok != gst.PadLinkOK { return fmt.Errorf("failed to link audio queue source pad to mux audio sink pad: %v", ok) }
// Get elements appsrc, err := pipeline.GetElementByName("appsrc") if err != nil { return err } videoappsink, err := pipeline.GetElementByName("videoappsink") if err != nil { return err } audioappsink, err := pipeline.GetElementByName("audioappsink") if err != nil { return err }
source := app.SrcFromElement(appsrc) videoSink := app.SinkFromElement(videoappsink) audioSink := app.SinkFromElement(audioappsink)
// Set up source callbacks source.SetCallbacks(&app.SourceCallbacks{ NeedDataFunc: ReaderNeedDataIncremental(ctx, input), EnoughDataFunc: func(self *app.Source) { // Nothing to do here }, SeekDataFunc: func(self *app.Source, offset uint64) bool { return false // We don't support seeking }, })
// Set up sink callbacks videoSink.SetCallbacks(&app.SinkCallbacks{ NewSampleFunc: WriterNewSample(ctx, videoOutput), NewPrerollFunc: func(self *app.Sink) gst.FlowReturn { return gst.FlowOK }, })
// Set up sink callbacks audioSink.SetCallbacks(&app.SinkCallbacks{ NewSampleFunc: WriterNewSample(ctx, audioOutput), NewPrerollFunc: func(self *app.Sink) gst.FlowReturn { return gst.FlowOK }, })
ctx, cancel := context.WithCancel(ctx) defer cancel()
// Handle bus messages in a separate goroutine g, ctx := errgroup.WithContext(ctx) g.Go(func() error { err = HandleBusMessages(ctx, pipeline) cancel() return err })
// Start the pipeline err = pipeline.SetState(gst.StatePlaying) if err != nil { return fmt.Errorf("failed to set pipeline state to playing: %w", err) }
// Wait for the pipeline to finish or context to be canceled <-ctx.Done()
// Clean up err = pipeline.SetState(gst.StateNull) if err != nil { return fmt.Errorf("failed to set pipeline state to null: %w", err) }
return g.Wait()}
// Joins video and audio back together from MPEG-TS and MP4 (from transcoding)func MPEGTSVideoMP4AudioToMP4(ctx context.Context, videoInput io.Reader, audioInput io.Reader, output io.Writer) error { pipelineStr := strings.Join([]string{ "appsrc name=videoappsrc ! tsdemux name=videodemux", "appsrc name=audioappsrc ! qtdemux name=audiodemux", "mp4mux name=mux ! appsink name=appsink sync=false", "h264parse name=videoparse ! video/x-h264,stream-format=avc ! queue name=videoqueue max-size-time=0 max-size-buffers=0 max-size-bytes=0", "audiodemux.audio_0 ! parsebin ! queue name=audioqueue max-size-time=0 max-size-buffers=0 max-size-bytes=0", }, " ")
pipeline, err := gst.NewPipelineFromString(pipelineStr) if err != nil { return err }
mux, err := pipeline.GetElementByName("mux") if err != nil { return err } muxVideoSinkPad := mux.GetRequestPad("video_%u") if muxVideoSinkPad == nil { return fmt.Errorf("failed to get video sink pad") } muxAudioSinkPad := mux.GetRequestPad("audio_%u") if muxAudioSinkPad == nil { return fmt.Errorf("failed to get audio sink pad") }
videoQueue, err := pipeline.GetElementByName("videoqueue") if err != nil { return err } audioQueue, err := pipeline.GetElementByName("audioqueue") if err != nil { return err }
videoQueueSrcPad := videoQueue.GetStaticPad("src") if videoQueueSrcPad == nil { return fmt.Errorf("failed to get video queue source pad") } audioQueueSrcPad := audioQueue.GetStaticPad("src") if audioQueueSrcPad == nil { return fmt.Errorf("failed to get audio queue source pad") }
ok := videoQueueSrcPad.Link(muxVideoSinkPad) if ok != gst.PadLinkOK { return fmt.Errorf("failed to link video queue source pad to mux video sink pad: %v", ok) } ok = audioQueueSrcPad.Link(muxAudioSinkPad) if ok != gst.PadLinkOK { return fmt.Errorf("failed to link audio queue source pad to mux audio sink pad: %v", ok) }
videodemux, err := pipeline.GetElementByName("videodemux") if err != nil { return err } videoparse, err := pipeline.GetElementByName("videoparse") if err != nil { return err } videoParseSinkPad := videoparse.GetStaticPad("sink") if videoParseSinkPad == nil { return fmt.Errorf("failed to get video parse sink pad") }
// Get elements videoappsrc, err := pipeline.GetElementByName("videoappsrc") if err != nil { return err } audioappsrc, err := pipeline.GetElementByName("audioappsrc") if err != nil { return err } appsink, err := pipeline.GetElementByName("appsink") if err != nil { return err }
videoSource := app.SrcFromElement(videoappsrc) audioSource := app.SrcFromElement(audioappsrc) sink := app.SinkFromElement(appsink)
// Set up source callbacks videoSource.SetCallbacks(&app.SourceCallbacks{ NeedDataFunc: ReaderNeedDataIncremental(ctx, videoInput), EnoughDataFunc: func(self *app.Source) { // Nothing to do here }, SeekDataFunc: func(self *app.Source, offset uint64) bool { return false // We don't support seeking }, })
audioSource.SetCallbacks(&app.SourceCallbacks{ NeedDataFunc: ReaderNeedDataIncremental(ctx, audioInput), EnoughDataFunc: func(self *app.Source) { // Nothing to do here }, SeekDataFunc: func(self *app.Source, offset uint64) bool { return false // We don't support seeking }, })
wroteAnything := false
// Set up sink callbacks sink.SetCallbacks(&app.SinkCallbacks{ NewSampleFunc: func(sink *app.Sink) gst.FlowReturn { sample := sink.PullSample() if sample == nil { return gst.FlowOK }
// Retrieve the buffer from the sample. buffer := sample.GetBuffer() bs := buffer.Map(gst.MapRead).Bytes() defer buffer.Unmap()
_, err := output.Write(bs)
if err != nil { log.Error(ctx, "error writing to output", "error", err) return gst.FlowError }
wroteAnything = true
return gst.FlowOK }, NewPrerollFunc: func(self *app.Sink) gst.FlowReturn { return gst.FlowOK }, })
ctx, cancel := context.WithCancel(ctx) defer cancel()
onPadAdded := func(element *gst.Element, pad *gst.Pad) { if pad.GetDirection() == gst.PadDirectionSource { ok := pad.Link(videoParseSinkPad) if ok != gst.PadLinkOK { log.Error(ctx, "failed to link video parse sink pad to video demux pad", "error", ok) cancel() } } } if _, err := videodemux.Connect("pad-added", onPadAdded); err != nil { return fmt.Errorf("failed connect pad-added handler: %w", err) }
errCh := make(chan error) go func() { err = HandleBusMessages(ctx, pipeline) cancel() errCh <- err }()
// Start the pipeline err = pipeline.SetState(gst.StatePlaying) if err != nil { return fmt.Errorf("failed to set pipeline state to playing: %w", err) }
defer func() { err = pipeline.SetState(gst.StateNull) if err != nil { log.Error(ctx, "failed to set pipeline state to null", "error", err) } videoParseSinkPad = nil }()
err = <-errCh if err != nil { return fmt.Errorf("pipeline error: %w", err) }
if !wroteAnything { return fmt.Errorf("no data written to output") }
return nil}