Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
2.1 kB · 99 lines
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100package media
import ( "context" "fmt" "strings" "time"
"github.com/bluenviron/gortsplib/v5/pkg/format" "github.com/go-gst/go-gst/gst" "stream.place/streamplace/pkg/log")
type RTMPH264Data struct { AU [][]byte PTS time.Duration DTS time.Duration}
type RTMPAACData struct { AU []byte PTS time.Duration}
type RTMPSession struct { EventChan chan any VideoTrack *format.H264 AudioTrack *format.MPEG4Audio MediaSigner MediaSigner}
func (mm *MediaManager) RTMPIngest(ctx context.Context, rtmpURL string, ms MediaSigner) error { ctx, cancel := context.WithCancel(ctx) defer cancel() // Mint the source audio: RTMP/FLV audio is already AAC, so pass it through // (aacparse) rather than transcoding to Opus. The validate path completes // each segment to also carry Opus when a consumer (WebRTC) needs it — so // the old RTMP-AAC→Opus→HLS-AAC double-transcode is gone. pipelineSlice := []string{ fmt.Sprintf("rtmp2src location=%s ! flvdemux name=demux", rtmpURL), "demux.audio ! queue ! aacparse name=audioenc", "demux.video ! queue ! h264parse name=parse", } pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) if err != nil { return fmt.Errorf("error creating RTMPIngest pipeline: %w", err) }
signer, err := mm.SegmentAndSignElem(ctx, ms) if err != nil { return err }
parseEle, err := pipeline.GetElementByName("parse") if err != nil { return err }
err = pipeline.Add(signer) if err != nil { return err } err = parseEle.Link(signer) if err != nil { return err } audioenc, err := pipeline.GetElementByName("audioenc") if err != nil { return err } err = audioenc.Link(signer) if err != nil { return err }
busErr := make(chan error) go func() { err := HandleBusMessages(ctx, pipeline) busErr <- err }()
go mm.HandleKeyRevocation(ctx, ms, pipeline)
err = pipeline.SetState(gst.StatePlaying) if err != nil { return err }
defer func() { err := pipeline.SetState(gst.StateNull) if err != nil { log.Error(ctx, "error setting pipeline to null state", "error", err) } }()
err = <-busErr
return err}