diff --git a/pkg/api/playback.go b/pkg/api/playback.go index 275b775bb..1445d37e4 100644 --- a/pkg/api/playback.go +++ b/pkg/api/playback.go @@ -51,12 +51,7 @@ func (a *StreamplaceAPI) HandleWebRTCPlayback(ctx context.Context) httprouter.Ha return } offer := webrtc.SessionDescription{Type: webrtc.SDPTypeOffer, SDP: string(body)} - var answer *webrtc.SessionDescription - if a.CLI.NewWebRTCPlayback { - answer, err = a.MediaManager.WebRTCPlayback2(ctx, user, rendition, &offer, "") - } else { - answer, err = a.MediaManager.WebRTCPlayback(ctx, user, rendition, &offer) - } + answer, err := a.MediaManager.WebRTCPlayback2(ctx, user, rendition, &offer, "") if err != nil { errors.WriteHTTPInternalServerError(w, fmt.Sprintf("error playing back: %s", err.Error()), err) return diff --git a/pkg/config/config.go b/pkg/config/config.go index 1b563b294..9ec684889 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -121,7 +121,6 @@ type CLI struct { ServiceAuthKey jwk.Key dataDirFlags []*string DiscordWebhooks []*discordtypes.Webhook - NewWebRTCPlayback bool AppleTeamID string AndroidCertFingerprint string Labelers []string @@ -613,13 +612,6 @@ func (cli *CLI) NewCommand(name string) *urfavecli.Command { }, Sources: urfavecli.EnvVars("SP_DISCORD_WEBHOOKS"), }, - &urfavecli.BoolFlag{ - Name: "new-webrtc-playback", - Usage: "enable new webrtc playback", - Value: true, - Destination: &cli.NewWebRTCPlayback, - Sources: urfavecli.EnvVars("SP_NEW_WEBRTC_PLAYBACK"), - }, &urfavecli.StringFlag{ Name: "apple-team-id", Usage: "apple team id for deep linking", diff --git a/pkg/media/concat.go b/pkg/media/concat.go deleted file mode 100644 index cde33410f..000000000 --- a/pkg/media/concat.go +++ /dev/null @@ -1,392 +0,0 @@ -package media - -import ( - "bytes" - "context" - "errors" - "fmt" - "io" - "strings" - "sync" - - "github.com/go-gst/go-gst/gst" - "github.com/go-gst/go-gst/gst/app" - "stream.place/streamplace/pkg/bus" - "stream.place/streamplace/pkg/log" -) - -type ConcatStreamer interface { - SubscribeSegment(ctx context.Context, user string, rendition string) *bus.SegChan - UnsubscribeSegment(ctx context.Context, user string, rendition string, ch *bus.SegChan) -} - -// This function remains in scope for the duration of a single users' playback -func ConcatStream(ctx context.Context, pipeline *gst.Pipeline, user string, rendition string, streamer ConcatStreamer) (*gst.Element, <-chan struct{}, error) { - ctx = log.WithLogValues(ctx, "func", "ConcatStream") - ctx, cancel := context.WithCancel(ctx) - defer cancel() - - // make 1000000000000 elements! - - // input multiqueue - inputQueue, err := gst.NewElementWithProperties("multiqueue", map[string]any{}) - if err != nil { - return nil, nil, fmt.Errorf("failed to create multiqueue element: %w", err) - } - err = pipeline.Add(inputQueue) - if err != nil { - return nil, nil, fmt.Errorf("failed to add input multiqueue to pipeline: %w", err) - } - inputQueuePadVideoSink := inputQueue.GetRequestPad("sink_%u") - if inputQueuePadVideoSink == nil { - return nil, nil, fmt.Errorf("failed to get input queue video sink pad") - } - inputQueuePadAudioSink := inputQueue.GetRequestPad("sink_%u") - if inputQueuePadAudioSink == nil { - return nil, nil, fmt.Errorf("failed to get input queue audio sink pad") - } - inputQueuePadVideoSrc := inputQueue.GetStaticPad("src_0") - if inputQueuePadVideoSrc == nil { - return nil, nil, fmt.Errorf("failed to get input queue video src pad") - } - inputQueuePadAudioSrc := inputQueue.GetStaticPad("src_1") - if inputQueuePadAudioSrc == nil { - return nil, nil, fmt.Errorf("failed to get input queue audio src pad") - } - // streamsynchronizer - streamsynchronizer, err := gst.NewElementWithProperties("streamsynchronizer", map[string]any{}) - if err != nil { - return nil, nil, fmt.Errorf("failed to create streamsynchronizer element: %w", err) - } - - err = pipeline.Add(streamsynchronizer) - if err != nil { - return nil, nil, fmt.Errorf("failed to add streamsynchronizer to pipeline: %w", err) - } - syncPadVideoSink := streamsynchronizer.GetRequestPad("sink_%u") - if syncPadVideoSink == nil { - return nil, nil, fmt.Errorf("failed to get sync video sink pad") - } - syncPadAudioSink := streamsynchronizer.GetRequestPad("sink_%u") - if syncPadAudioSink == nil { - return nil, nil, fmt.Errorf("failed to get sync audio sink pad") - } - syncPadVideoSrc := streamsynchronizer.GetStaticPad("src_0") - if syncPadVideoSrc == nil { - return nil, nil, fmt.Errorf("failed to get sync video src pad") - } - syncPadAudioSrc := streamsynchronizer.GetStaticPad("src_1") - if syncPadAudioSrc == nil { - return nil, nil, fmt.Errorf("failed to get sync audio src pad") - } - - // output multiqueue - outputQueue, err := gst.NewElementWithProperties("multiqueue", map[string]any{ - "name": "concat-output-queue", - }) - if err != nil { - return nil, nil, fmt.Errorf("failed to create multiqueue element: %w", err) - } - err = pipeline.Add(outputQueue) - if err != nil { - return nil, nil, fmt.Errorf("failed to add output multiqueue to pipeline: %w", err) - } - outputQueuePadVideoSink := outputQueue.GetRequestPad("sink_%u") - if outputQueuePadVideoSink == nil { - return nil, nil, fmt.Errorf("failed to get output queue video sink pad") - } - outputQueuePadAudioSink := outputQueue.GetRequestPad("sink_%u") - if outputQueuePadAudioSink == nil { - return nil, nil, fmt.Errorf("failed to get output queue audio sink pad") - } - // linking - - // input queue to streamsynchronizer - ret := inputQueuePadVideoSrc.Link(syncPadVideoSink) - if ret != gst.PadLinkOK { - return nil, nil, fmt.Errorf("failed to link multiqueue to streamsynchronizer: %v", ret) - } - ret = inputQueuePadAudioSrc.Link(syncPadAudioSink) - if ret != gst.PadLinkOK { - return nil, nil, fmt.Errorf("failed to link multiqueue to streamsynchronizer: %v", ret) - } - - // streamsynchronizer to output queue - ret = syncPadVideoSrc.Link(outputQueuePadVideoSink) - if ret != gst.PadLinkOK { - return nil, nil, fmt.Errorf("failed to link streamsynchronizer to output queue: %v", ret) - } - ret = syncPadAudioSrc.Link(outputQueuePadAudioSink) - if ret != gst.PadLinkOK { - return nil, nil, fmt.Errorf("failed to link streamsynchronizer to output queue: %v", ret) - } - - // ok now we can start looping over input files - - // this goroutine will read all the files from the segment queue and buffer - // them in a pipe so that we don't miss any in between iterations of the output - allFiles := make(chan []byte, 1024) - go func() { - ch := streamer.SubscribeSegment(ctx, user, rendition) - defer streamer.UnsubscribeSegment(ctx, user, rendition, ch) - for { - select { - case <-ctx.Done(): - log.Debug(ctx, "exiting segment reader") - return - case file := <-ch.C: - log.Debug(ctx, "got segment", "file", file.Filepath) - allFiles <- file.Data - if len(file.Data) == 0 { - log.Warn(ctx, "no more segments, stopping segment reader") - return - } - } - } - }() - - segCount := 0 - - // nextFile is the primary loop that pops off a file, creates new demuxer elements for it, - // and pushes into the pipeline - var nextFile func() - nextFile = func() { - mySegCount := segCount - segCount += 1 - segDone := make(chan struct{}) - log.Debug(ctx, "moving to next file", "segCount", mySegCount) - pr, pw := io.Pipe() - go func() { - select { - case <-ctx.Done(): - pr.Close() - pw.Close() - return - case bs := <-allFiles: - if len(bs) == 0 { - log.Warn(ctx, "no more segments, ending stream") - pr.Close() - pw.Close() - cancel() - return - } - _, err = io.Copy(pw, bytes.NewReader(bs)) - if err != nil { - log.Error(ctx, "failed to copy segment file", "error", err) - cancel() - return - } - return - } - }() - - demux, err := gst.NewElementWithProperties("qtdemux", map[string]any{ - "name": fmt.Sprintf("concat-demux-%d", mySegCount), - }) - if err != nil { - log.Error(ctx, "failed to create demux element", "error", err) - cancel() - return - } - - err = pipeline.Add(demux) - if err != nil { - log.Error(ctx, "failed to add demux to pipeline", "error", err) - cancel() - return - } - - demuxSinkPad := demux.GetStaticPad("sink") - if demuxSinkPad == nil { - log.Error(ctx, "failed to get demux sink pad") - cancel() - return - } - - mu := sync.Mutex{} - count := 0 - _, err = demux.Connect("pad-added", func(self *gst.Element, pad *gst.Pad) { - mu.Lock() - count += 1 - mu.Unlock() - log.Debug(ctx, "demux pad-added", "name", pad.GetName(), "direction", pad.GetDirection()) - var downstreamPad *gst.Pad - if strings.HasPrefix(pad.GetName(), "video_") { - downstreamPad = inputQueuePadVideoSink - } else if strings.HasPrefix(pad.GetName(), "audio_") { - downstreamPad = inputQueuePadAudioSink - } else { - log.Error(ctx, "unknown pad", "name", pad.GetName(), "direction", pad.GetDirection()) - cancel() - return - } - ret := pad.Link(downstreamPad) - if ret != gst.PadLinkOK { - log.Error(ctx, "failed to link demux to downstream pad", "name", pad.GetName(), "direction", pad.GetDirection(), "error", ret) - cancel() - return - } - if pad.GetDirection() == gst.PadDirectionSource { - pad.AddProbe(gst.PadProbeTypeEventBoth, func(pad *gst.Pad, info *gst.PadProbeInfo) gst.PadProbeReturn { - if info.GetEvent().Type() != gst.EventTypeEOS { - return gst.PadProbeOK - } - log.Debug(ctx, "demux EOS", "name", pad.GetName(), "direction", pad.GetDirection()) - pad.Unlink(downstreamPad) - mu.Lock() - defer mu.Unlock() - count -= 1 - - if count == 0 { - // don't keep going if our context is done - if ctx.Err() == nil { - go nextFile() - segDone <- struct{}{} - } - } else { - log.Debug(ctx, "demux has more pads, waiting for them to close") - } - return gst.PadProbeRemove - }) - } - }) - if err != nil { - log.Error(ctx, "failed to connect demux pad-added", "error", err) - cancel() - return - } - - appsrc, err := gst.NewElementWithProperties("appsrc", map[string]any{ - "name": fmt.Sprintf("concat-appsrc-%d", mySegCount), - "is-live": true, - }) - if err != nil { - log.Error(ctx, "failed to get appsrc element from pipeline", "error", err) - cancel() - return - } - - src := app.SrcFromElement(appsrc) - - appSrcPad := appsrc.GetStaticPad("src") - if appSrcPad == nil { - log.Error(ctx, "failed to get appsrc pad") - cancel() - return - } - - done := func() { - // appsrc.Unlink(demux) - pads, err := src.GetPads() - if err != nil { - log.Error(ctx, "failed to get pads", "error", err) - cancel() - return - } - for _, pad := range pads { - log.Debug(ctx, "setting pad-idle", "name", pad.GetName(), "direction", pad.GetDirection()) - - pad.AddProbe(gst.PadProbeTypeIdle, func(pad *gst.Pad, info *gst.PadProbeInfo) gst.PadProbeReturn { - log.Debug(ctx, "pad-idle", "name", pad.GetName(), "direction", pad.GetDirection()) - src.EndStream() - return gst.PadProbeRemove - }) - } - } - - src.SetAutomaticEOS(false) - src.SetCallbacks(&app.SourceCallbacks{ - NeedDataFunc: func(self *app.Source, length uint) { - bs := make([]byte, length) - read, err := pr.Read(bs) - if err != nil { - if errors.Is(err, io.EOF) { - if read > 0 { - log.Debug(ctx, "got data on eof???") - cancel() - return - } - log.Debug(ctx, "EOF, ending segment", "length", read) - done() - return - } else { - log.Debug(ctx, "failed to read data, ending stream", "error", err) - cancel() - return - } - } - toPush := bs - if uint(read) < length { - toPush = bs[:read] - } - buffer := gst.NewBufferWithSize(int64(len(toPush))) - buffer.Map(gst.MapWrite).WriteData(toPush) - defer buffer.Unmap() - self.PushBuffer(buffer) - - if uint(read) < length { - log.Debug(ctx, "short write, ending segment", "length", read) - done() - } - }, - }) - err = pipeline.Add(appsrc) - if err != nil { - log.Error(ctx, "failed to add appsrc to pipeline", "error", err) - cancel() - return - } - - ret := appSrcPad.Link(demuxSinkPad) - if ret != gst.PadLinkOK { - log.Error(ctx, "failed to link appsrc to demux", "error", ret) - cancel() - return - } - - err = demux.SetState(gst.StatePlaying) - if err != nil { - log.Error(ctx, "failed to set demux state", "error", err) - cancel() - return - } - err = appsrc.SetState(gst.StatePlaying) - if err != nil { - log.Error(ctx, "failed to set appsrc state", "error", err) - cancel() - return - } - - select { - case <-ctx.Done(): - return - case <-segDone: - } - - log.Debug(ctx, "ending segment") - if err := demux.SetState(gst.StateNull); err != nil { - log.Error(ctx, "failed to set demux state", "error", err) - return - } - src.SetCallbacks(&app.SourceCallbacks{}) - if err := appsrc.SetState(gst.StateNull); err != nil { - log.Error(ctx, "failed to set appsrc state", "error", err) - return - } - if err := pipeline.Remove(demux); err != nil { - log.Error(ctx, "failed to remove demux from pipleine", "error", err) - return - } - if err := pipeline.Remove(appsrc); err != nil { - log.Error(ctx, "failed to remove appsrc from pipleine", "error", err) - return - } - pr.Close() - pw.Close() - } - - // fire it up! - go nextFile() - - return outputQueue, ctx.Done(), nil -} diff --git a/pkg/media/progressive.go b/pkg/media/progressive.go deleted file mode 100644 index 79684c861..000000000 --- a/pkg/media/progressive.go +++ /dev/null @@ -1,185 +0,0 @@ -package media - -// func (mm *MediaManager) MP4Playback(ctx context.Context, user string, rendition string, w io.Writer) error { -// uu, err := uuid.NewV7() -// if err != nil { -// return err -// } -// ctx = log.WithLogValues(ctx, "playbackID", uu.String()) -// ctx, cancel := context.WithCancel(ctx) - -// ctx = log.WithLogValues(ctx, "mediafunc", "MP4Playback") - -// pipelineSlice := []string{ -// "mp4mux name=muxer fragment-mode=first-moov-then-finalise fragment-duration=1000 streamable=true ! appsink name=mp4sink", -// "h264parse name=videoparse ! muxer.", -// "opusparse name=audioparse ! muxer.", -// } - -// pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) -// if err != nil { -// return fmt.Errorf("failed to create GStreamer pipeline: %w", err) -// } - -// go func() { -// HandleBusMessages(ctx, pipeline) -// cancel() -// }() - -// outputQueue, done, err := ConcatStream(ctx, pipeline, user, rendition, mm) -// if err != nil { -// return fmt.Errorf("failed to get output queue: %w", err) -// } -// go func() { -// select { -// case <-ctx.Done(): -// return -// case <-done: -// cancel() -// } -// }() - -// videoParse, err := pipeline.GetElementByName("videoparse") -// if err != nil { -// return fmt.Errorf("failed to get video sink element from pipeline: %w", err) -// } -// err = outputQueue.Link(videoParse) -// if err != nil { -// return fmt.Errorf("failed to link output queue to video parse: %w", err) -// } - -// audioParse, err := pipeline.GetElementByName("audioparse") -// if err != nil { -// return fmt.Errorf("failed to get audio parse element from pipeline: %w", err) -// } -// err = outputQueue.Link(audioParse) -// if err != nil { -// return fmt.Errorf("failed to link output queue to audio parse: %w", err) -// } - -// go func() { -// ticker := time.NewTicker(time.Second * 1) -// for { -// select { -// case <-ctx.Done(): -// return -// case <-ticker.C: -// state := pipeline.GetCurrentState() -// log.Debug(ctx, "pipeline state", "state", state) -// } -// } -// }() - -// mp4sinkele, err := pipeline.GetElementByName("mp4sink") -// if err != nil { -// return fmt.Errorf("failed to get video sink element from pipeline: %w", err) -// } -// mp4sink := app.SinkFromElement(mp4sinkele) -// mp4sink.SetCallbacks(&app.SinkCallbacks{ -// NewSampleFunc: WriterNewSample(ctx, w), -// EOSFunc: func(sink *app.Sink) { -// log.Warn(ctx, "mp4sink EOSFunc") -// cancel() -// }, -// }) - -// pipeline.SetState(gst.StatePlaying) - -// <-ctx.Done() - -// pipeline.BlockSetState(gst.StateNull) - -// return nil -// } - -// func (mm *MediaManager) MKVPlayback(ctx context.Context, user string, rendition string, w io.Writer) error { -// uu, err := uuid.NewV7() -// if err != nil { -// return err -// } -// ctx = log.WithLogValues(ctx, "playbackID", uu.String()) -// ctx, cancel := context.WithCancel(ctx) - -// ctx = log.WithLogValues(ctx, "mediafunc", "MKVPlayback") - -// pipelineSlice := []string{ -// "matroskamux name=muxer streamable=true ! appsink name=mkvsink", -// "h264parse name=videoparse ! muxer.", -// "opusparse name=audioparse ! muxer.", -// } - -// pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) -// if err != nil { -// return fmt.Errorf("failed to create GStreamer pipeline: %w", err) -// } - -// go func() { -// HandleBusMessages(ctx, pipeline) -// cancel() -// }() - -// outputQueue, done, err := ConcatStream(ctx, pipeline, user, rendition, mm) -// if err != nil { -// return fmt.Errorf("failed to get output queue: %w", err) -// } -// go func() { -// select { -// case <-ctx.Done(): -// return -// case <-done: -// cancel() -// } -// }() - -// videoParse, err := pipeline.GetElementByName("videoparse") -// if err != nil { -// return fmt.Errorf("failed to get video sink element from pipeline: %w", err) -// } -// err = outputQueue.Link(videoParse) -// if err != nil { -// return fmt.Errorf("failed to link output queue to video parse: %w", err) -// } - -// audioParse, err := pipeline.GetElementByName("audioparse") -// if err != nil { -// return fmt.Errorf("failed to get audio parse element from pipeline: %w", err) -// } -// err = outputQueue.Link(audioParse) -// if err != nil { -// return fmt.Errorf("failed to link output queue to audio parse: %w", err) -// } - -// go func() { -// ticker := time.NewTicker(time.Second * 1) -// for { -// select { -// case <-ctx.Done(): -// return -// case <-ticker.C: -// state := pipeline.GetCurrentState() -// log.Debug(ctx, "pipeline state", "state", state) -// } -// } -// }() - -// mkvsinkele, err := pipeline.GetElementByName("mkvsink") -// if err != nil { -// return fmt.Errorf("failed to get video sink element from pipeline: %w", err) -// } -// mkvsink := app.SinkFromElement(mkvsinkele) -// mkvsink.SetCallbacks(&app.SinkCallbacks{ -// NewSampleFunc: WriterNewSample(ctx, w), -// EOSFunc: func(sink *app.Sink) { -// log.Warn(ctx, "mp4sink EOSFunc") -// cancel() -// }, -// }) - -// pipeline.SetState(gst.StatePlaying) - -// <-ctx.Done() - -// pipeline.BlockSetState(gst.StateNull) - -// return nil -// } diff --git a/pkg/media/webrtc_playback.go b/pkg/media/webrtc_playback.go deleted file mode 100644 index 53b2bc9e0..000000000 --- a/pkg/media/webrtc_playback.go +++ /dev/null @@ -1,375 +0,0 @@ -package media - -import ( - "context" - "fmt" - "strings" - "time" - - "github.com/go-gst/go-gst/gst" - "github.com/go-gst/go-gst/gst/app" - "github.com/google/uuid" - "github.com/pion/webrtc/v4" - "github.com/pion/webrtc/v4/pkg/media" - "stream.place/streamplace/pkg/bus" - "stream.place/streamplace/pkg/log" -) - -// we have a bug that prevents us from correctly probing video durations -// a lot of the time. so when we don't have them we use the last duration -// that we had, and when we don't have that we use a default duration -var DefaultDuration = time.Duration(32 * time.Millisecond) - -// This function remains in scope for the duration of a single users' playback -func (mm *MediaManager) WebRTCPlayback(ctx context.Context, user string, rendition string, offer *webrtc.SessionDescription) (*webrtc.SessionDescription, error) { - uu, err := uuid.NewV7() - if err != nil { - return nil, err - } - ctx = log.WithLogValues(ctx, "webrtcID", uu.String()) - ctx = log.WithLogValues(ctx, "mediafunc", "WebRTCPlayback") - ctx, cancel := context.WithCancel(ctx) //nolint:all - - pipelineSlice := []string{ - "h264parse name=videoparse ! video/x-h264,stream-format=byte-stream ! appsink name=videoappsink", - "opusparse name=audioparse ! appsink name=audioappsink", - } - - pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) - if err != nil { - return nil, fmt.Errorf("failed to create GStreamer pipeline: %w", err) //nolint:all - } - - segBuffer := make(chan *bus.Seg, 1024) - go func() { - segChan := mm.bus.SubscribeSegment(ctx, user, rendition) - defer mm.bus.UnsubscribeSegment(ctx, user, rendition, segChan) - for { - select { - case <-ctx.Done(): - log.Debug(ctx, "exiting segment reader") - return - case file := <-segChan.C: - log.Debug(ctx, "got segment", "file", file.Filepath) - segBuffer <- file - } - } - }() - - segCh := make(chan *bus.Seg) - go func() { - for { - select { - case <-ctx.Done(): - log.Debug(ctx, "exiting segment reader") - return - case seg := <-segBuffer: - select { - case <-ctx.Done(): - return - case segCh <- seg: - } - } - } - }() - - concatBin, err := ConcatBin(ctx, segCh, true) - if err != nil { - return nil, fmt.Errorf("failed to create concat bin: %w", err) - } - - err = pipeline.Add(concatBin.Element) - if err != nil { - return nil, fmt.Errorf("failed to add concat bin to pipeline: %w", err) - } - - videoPad := concatBin.GetStaticPad("video_0") - if videoPad == nil { - return nil, fmt.Errorf("video pad not found") - } - - audioPad := concatBin.GetStaticPad("audio_0") - if audioPad == nil { - return nil, fmt.Errorf("audio pad not found") - } - - // queuePadVideo := outputQueue.GetRequestPad("src_%u") - // if queuePadVideo == nil { - // return nil, fmt.Errorf("failed to get queue video pad") - // } - // queuePadAudio := outputQueue.GetRequestPad("src_%u") - // if queuePadAudio == nil { - // return nil, fmt.Errorf("failed to get queue audio pad") - // } - - videoParse, err := pipeline.GetElementByName("videoparse") - if err != nil { - return nil, fmt.Errorf("failed to get video sink element from pipeline: %w", err) - } - videoParsePad := videoParse.GetStaticPad("sink") - if videoParsePad == nil { - return nil, fmt.Errorf("video parse pad not found") - } - linked := videoPad.Link(videoParsePad) - if linked != gst.PadLinkOK { - return nil, fmt.Errorf("failed to link video pad to video parse pad: %v", linked) - } - - audioParse, err := pipeline.GetElementByName("audioparse") - if err != nil { - return nil, fmt.Errorf("failed to get audio parse element from pipeline: %w", err) - } - audioParsePad := audioParse.GetStaticPad("sink") - if audioParsePad == nil { - return nil, fmt.Errorf("audio parse pad not found") - } - linked = audioPad.Link(audioParsePad) - if linked != gst.PadLinkOK { - return nil, fmt.Errorf("failed to link audio pad to audio parse pad: %v", linked) - } - - videoappsinkele, err := pipeline.GetElementByName("videoappsink") - if err != nil { - return nil, fmt.Errorf("failed to get video sink element from pipeline: %w", err) - } - - audioappsinkele, err := pipeline.GetElementByName("audioappsink") - if err != nil { - return nil, fmt.Errorf("failed to get audio sink element from pipeline: %w", err) - } - - // Create a new RTCPeerConnection - peerConnection, err := mm.webrtcAPI.NewPeerConnection(mm.webrtcConfig) - if err != nil { - return nil, fmt.Errorf("failed to create WebRTC peer connection: %w", err) - } - go func() { - <-ctx.Done() - if cErr := peerConnection.Close(); cErr != nil { - log.Log(ctx, "cannot close peerConnection: %v\n", cErr) - } - }() - - videoTrack, err := webrtc.NewTrackLocalStaticSample(webrtc.RTPCodecCapability{MimeType: webrtc.MimeTypeH264}, "video", "pion") - if err != nil { - return nil, fmt.Errorf("failed to create video track: %w", err) - } - videoRTPSender, err := peerConnection.AddTrack(videoTrack) - if err != nil { - return nil, fmt.Errorf("failed to add video track to peer connection: %w", err) - } - - audioTrack, err := webrtc.NewTrackLocalStaticSample(webrtc.RTPCodecCapability{MimeType: webrtc.MimeTypeOpus}, "audio", "pion") - if err != nil { - return nil, fmt.Errorf("failed to create audio track: %w", err) - } - audioRTPSender, err := peerConnection.AddTrack(audioTrack) - if err != nil { - return nil, fmt.Errorf("failed to add audio track to peer connection: %w", err) - } - - // Set the remote SessionDescription - if err = peerConnection.SetRemoteDescription(*offer); err != nil { - return nil, fmt.Errorf("failed to set remote description: %w", err) - } - - // Create answer - answer, err := peerConnection.CreateAnswer(nil) - if err != nil { - return nil, fmt.Errorf("failed to create answer: %w", err) - } - - // Sets the LocalDescription, and starts our UDP listeners - if err = peerConnection.SetLocalDescription(answer); err != nil { - return nil, fmt.Errorf("failed to set local description: %w", err) - } - - // Create channel that is blocked until ICE Gathering is complete - gatherComplete := webrtc.GatheringCompletePromise(peerConnection) - - // Setup complete! Now we boot up streaming in the background while returning the SDP offer to the user. - - go func() { - ticker := time.NewTicker(time.Second * 1) - for { - select { - case <-ctx.Done(): - return - case <-ticker.C: - state := pipeline.GetCurrentState() - log.Debug(ctx, "pipeline state", "state", state) - } - } - }() - - var lastVideoDuration = &DefaultDuration - - go func() { - go func() { - if err := HandleBusMessages(ctx, pipeline); err != nil { - log.Log(ctx, "pipeline error", "error", err) - } - cancel() - }() - - videoappsink := app.SinkFromElement(videoappsinkele) - videoappsink.SetCallbacks(&app.SinkCallbacks{ - NewSampleFunc: func(sink *app.Sink) gst.FlowReturn { - sample := sink.PullSample() - if sample == nil { - return gst.FlowEOS - } - - buffer := sample.GetBuffer() - if buffer == nil { - return gst.FlowError - } - - samples := buffer.Map(gst.MapRead).Bytes() - defer buffer.Unmap() - clockTime := buffer.Duration() - dur := clockTime.AsDuration() - mediaSample := media.Sample{Data: samples} - if dur != nil { - mediaSample.Duration = *dur - lastVideoDuration = dur - } else if lastVideoDuration != nil { - // log.Log(ctx, "no video duration, using last duration", "lastVideoDuration", lastVideoDuration) - mediaSample.Duration = *lastVideoDuration - } else { - log.Log(ctx, "no video duration", "samples", len(samples)) - // cancel() - return gst.FlowOK - } - - if err := videoTrack.WriteSample(mediaSample); err != nil { - log.Log(ctx, "failed to write video sample", "error", err) - cancel() - } - - return gst.FlowOK - }, - EOSFunc: func(sink *app.Sink) { - log.Warn(ctx, "videoappsink EOSFunc") - cancel() - }, - }) - - audioappsink := app.SinkFromElement(audioappsinkele) - audioappsink.SetCallbacks(&app.SinkCallbacks{ - NewSampleFunc: func(sink *app.Sink) gst.FlowReturn { - sample := sink.PullSample() - if sample == nil { - return gst.FlowEOS - } - - buffer := sample.GetBuffer() - if buffer == nil { - return gst.FlowError - } - - samples := buffer.Map(gst.MapRead).Bytes() - defer buffer.Unmap() - - b2 := make([]byte, len(samples)) - copy(b2, samples) - - clockTime := buffer.Duration() - dur := clockTime.AsDuration() - mediaSample := media.Sample{Data: b2} - if dur != nil { - mediaSample.Duration = *dur - } else { - log.Log(ctx, "no audio duration", "samples", len(b2)) - // cancel() - return gst.FlowOK - } - if err := audioTrack.WriteSample(mediaSample); err != nil { - log.Log(ctx, "failed to write audio sample", "error", err) - return gst.FlowOK - } - - return gst.FlowOK - }, - EOSFunc: func(sink *app.Sink) { - log.Warn(ctx, "audioappsink EOSFunc") - cancel() - }, - }) - - // Start the pipeline - err := pipeline.SetState(gst.StatePlaying) - if err != nil { - log.Log(ctx, "failed to set pipeline state to null", "error", err) - } - mm.IncrementViewerCount(user, "webrtc") - defer mm.DecrementViewerCount(user, "webrtc") - - go func() { - rtcpBuf := make([]byte, 1500) - for { - if _, _, rtcpErr := videoRTPSender.Read(rtcpBuf); rtcpErr != nil { - return - } - } - }() - - go func() { - rtcpBuf := make([]byte, 1500) - for { - if _, _, rtcpErr := audioRTPSender.Read(rtcpBuf); rtcpErr != nil { - return - } - } - }() - - // Set the handler for ICE connection state - // This will notify you when the peer has connected/disconnected - peerConnection.OnICEConnectionStateChange(func(connectionState webrtc.ICEConnectionState) { - log.Log(ctx, "Connection State has changed", "state", connectionState.String()) - }) - - // Set the handler for Peer connection state - // This will notify you when the peer has connected/disconnected - peerConnection.OnConnectionStateChange(func(s webrtc.PeerConnectionState) { - log.Log(ctx, "Peer Connection State has changed", "state", s.String()) - - if s == webrtc.PeerConnectionStateFailed || s == webrtc.PeerConnectionStateClosed || s == webrtc.PeerConnectionStateDisconnected { - // Wait until PeerConnection has had no network activity for 30 seconds or another failure. It may be reconnected using an ICE Restart. - // Use webrtc.PeerConnectionStateDisconnected if you are interested in detecting faster timeout. - // Note that the PeerConnection may come back from PeerConnectionStateDisconnected. - log.Log(ctx, "Peer Connection has gone to failed, exiting") - cancel() - } - }) - - <-ctx.Done() - - log.Warn(ctx, "setting playback pipeline state to null") - err = pipeline.BlockSetState(gst.StateNull) - if err != nil { - log.Log(ctx, "failed to set pipeline state to null", "error", err) - } - - videoappsink.SetCallbacks(&app.SinkCallbacks{}) - err = videoappsinkele.SetState(gst.StateNull) - if err != nil { - log.Log(ctx, "failed to set videoappsinkele state to null", "error", err) - } - - audioappsink.SetCallbacks(&app.SinkCallbacks{}) - err = audioappsinkele.SetState(gst.StateNull) - if err != nil { - log.Log(ctx, "failed to set audioappsinkele state to null", "error", err) - } - - log.Warn(ctx, "exiting playback") - - }() - select { - case <-gatherComplete: - return peerConnection.LocalDescription(), nil - case <-ctx.Done(): - return nil, ctx.Err() - } -}