diff --git a/Makefile b/Makefile index 380253f7..306f7ba8 100644 --- a/Makefile +++ b/Makefile @@ -170,12 +170,14 @@ BASE_OPTS = \ -D "gst-plugins-good:multifile=enabled" \ -D "gst-plugins-good:rtp=enabled" \ -D "gst-plugins-bad:fdkaac=enabled" \ + -D "gst-plugins-bad:rtmp2=enabled" \ -D "gst-plugins-good:audioparsers=enabled" \ -D "gst-plugins-good:isomp4=enabled" \ -D "gst-plugins-good:png=enabled" \ -D "gst-plugins-good:videobox=enabled" \ -D "gst-plugins-good:jpeg=enabled" \ -D "gst-plugins-good:audioparsers=enabled" \ + -D "gst-plugins-good:flv=enabled" \ -D "gst-plugins-bad:videoparsers=enabled" \ -D "gst-plugins-bad:mpegtsmux=enabled" \ -D "gst-plugins-bad:mpegtsdemux=enabled" \ @@ -185,7 +187,7 @@ BASE_OPTS = \ -D "gst-plugins-ugly:gpl=enabled" \ -D "x264:asm=enabled" \ -D "gstreamer-full:gst-full=enabled" \ - -D "gstreamer-full:gst-full-plugins=libgstopusparse.a;libgstcodectimestamper.a;libgstrtp.a;libgstaudioresample.a;libgstlibav.a;libgstmatroska.a;libgstmultifile.a;libgstjpeg.a;libgstaudiotestsrc.a;libgstaudioconvert.a;libgstaudioparsers.a;libgstfdkaac.a;libgstisomp4.a;libgstapp.a;libgstvideoconvertscale.a;libgstvideobox.a;libgstvideorate.a;libgstpng.a;libgstcompositor.a;libgstaudiorate.a;libgstx264.a;libgstopus.a;libgstvideotestsrc.a;libgstvideoparsersbad.a;libgstaudioparsers.a;libgstmpegtsmux.a;libgstmpegtsdemux.a;libgstplayback.a;libgsttypefindfunctions.a;libgstcoretracers.a;libgstcodec2json.a" \ + -D "gstreamer-full:gst-full-plugins=libgstflv.a;libgstrtmp2.a;libgstopusparse.a;libgstcodectimestamper.a;libgstrtp.a;libgstaudioresample.a;libgstlibav.a;libgstmatroska.a;libgstmultifile.a;libgstjpeg.a;libgstaudiotestsrc.a;libgstaudioconvert.a;libgstaudioparsers.a;libgstfdkaac.a;libgstisomp4.a;libgstapp.a;libgstvideoconvertscale.a;libgstvideobox.a;libgstvideorate.a;libgstpng.a;libgstcompositor.a;libgstaudiorate.a;libgstx264.a;libgstopus.a;libgstvideotestsrc.a;libgstvideoparsersbad.a;libgstaudioparsers.a;libgstmpegtsmux.a;libgstmpegtsdemux.a;libgstplayback.a;libgsttypefindfunctions.a;libgstcoretracers.a;libgstcodec2json.a" \ -D "gstreamer-full:gst-full-libraries=gstreamer-controller-1.0,gstreamer-plugins-base-1.0,gstreamer-pbutils-1.0" \ -D "gstreamer-full:gst-full-elements=coreelements:concat,filesrc,filesink,queue,queue2,multiqueue,typefind,tee,capsfilter,fakesink,identity" \ -D "gstreamer-full:bad=enabled" \ diff --git a/go.mod b/go.mod index d54f2209..ff196f7d 100644 --- a/go.mod +++ b/go.mod @@ -19,6 +19,7 @@ require ( github.com/acarl005/stripansi v0.0.0-20180116102854-5a71ef0e047d github.com/bluenviron/gortmplib v0.1.2 github.com/bluenviron/gortsplib/v5 v5.2.1 + github.com/bluenviron/mediacommon/v2 v2.5.2 github.com/bluesky-social/indigo v0.0.0-20251206005924-d49b45419635 github.com/cenkalti/backoff v2.2.1+incompatible github.com/cenkalti/backoff/v5 v5.0.2 @@ -162,7 +163,6 @@ require ( github.com/bkielbasa/cyclop v1.2.3 // indirect github.com/blizzy78/varnamelen v0.8.0 // indirect github.com/bluenviron/gortsplib/v4 v4.12.3 // indirect - github.com/bluenviron/mediacommon/v2 v2.5.2 // indirect github.com/bombsimon/wsl/v4 v4.7.0 // indirect github.com/breml/bidichk v0.3.3 // indirect github.com/breml/errchkjson v0.4.1 // indirect diff --git a/pkg/api/api.go b/pkg/api/api.go index 58c18ca2..92ae289c 100644 --- a/pkg/api/api.go +++ b/pkg/api/api.go @@ -79,6 +79,10 @@ type StreamplaceAPI struct { HTTPRedirectTLSPort *int sessions map[string]map[string]time.Time sessionsLock sync.RWMutex + + rtmpSessions map[string]*media.RTMPSession + rtmpSessionsLock sync.Mutex + rtmpInternalPlaybackAddr string } type WebsocketTracker struct { @@ -109,6 +113,8 @@ func MakeStreamplaceAPI(cli *config.CLI, mod model.Model, statefulDB *statedb.St op: op, sessions: make(map[string]map[string]time.Time), sessionsLock: sync.RWMutex{}, + rtmpSessions: make(map[string]*media.RTMPSession), + rtmpSessionsLock: sync.Mutex{}, } a.Mimes, err = updater.GetMimes() if err != nil { diff --git a/pkg/api/rtmp_server.go b/pkg/api/rtmp_server.go index 858f42ac..c58f7bbb 100644 --- a/pkg/api/rtmp_server.go +++ b/pkg/api/rtmp_server.go @@ -27,10 +27,15 @@ import ( // readers []*gortmplib.Writer // ) +var RTMPTimeout = 10 * time.Second + const RTMPPrefix = "/live/" func (a *StreamplaceAPI) HandleRTMPPublisher(ctx context.Context, sc *gortmplib.ServerConn) error { - sc.RW.(net.Conn).SetReadDeadline(time.Now().Add(10 * time.Second)) + err := sc.RW.(net.Conn).SetReadDeadline(time.Now().Add(RTMPTimeout)) + if err != nil { + return err + } if !strings.HasPrefix(sc.URL.Path, RTMPPrefix) { return fmt.Errorf("RTMP publisher is not allowed to publish to %s (must start with %s)", sc.URL.String(), RTMPPrefix) @@ -41,12 +46,27 @@ func (a *StreamplaceAPI) HandleRTMPPublisher(ctx context.Context, sc *gortmplib. return fmt.Errorf("failed to make media signer: %w", err) } - ctx = log.WithLogValues(ctx, "streamer", mediaSigner.Streamer()) + streamer := mediaSigner.Streamer() + ctx = log.WithLogValues(ctx, "streamer", streamer) + session := &media.RTMPSession{ + EventChan: make(chan any, 1024), + MediaSigner: mediaSigner, + } + a.rtmpSessionsLock.Lock() + a.rtmpSessions[streamer] = session + a.rtmpSessionsLock.Unlock() - videoInput := make(chan *media.RTMPH264Data, 1024) - defer close(videoInput) - audioInput := make(chan *media.RTMPAACData, 1024) - defer close(audioInput) + defer func() { + a.rtmpSessionsLock.Lock() + delete(a.rtmpSessions, streamer) + a.rtmpSessionsLock.Unlock() + close(session.EventChan) + }() + + // videoInput := make(chan *media.RTMPH264Data, 1024) + // defer close(videoInput) + // audioInput := make(chan *media.RTMPAACData, 1024) + // defer close(audioInput) r := &gortmplib.Reader{ Conn: sc, @@ -61,18 +81,21 @@ func (a *StreamplaceAPI) HandleRTMPPublisher(ctx context.Context, sc *gortmplib. switch track := track.(type) { case *format.H264: + session.VideoTrack = track r.OnDataH264(track, func(pts time.Duration, dts time.Duration, au [][]byte) { - log.Log(ctx, "got H264", "len", len(au), "pts", pts, "dts", dts) - videoInput <- &media.RTMPH264Data{ + // log.Log(ctx, "got H264", "len", len(au), "pts", pts, "dts", dts) + session.EventChan <- &media.RTMPH264Data{ AU: au, PTS: pts, + DTS: dts, } }) case *format.MPEG4Audio: + session.AudioTrack = track r.OnDataMPEG4Audio(track, func(pts time.Duration, au []byte) { - log.Log(ctx, "got MPEG4Au", "len", len(au), "pts", pts) - audioInput <- &media.RTMPAACData{ + // log.Log(ctx, "got MPEG4Au", "len", len(au), "pts", pts) + session.EventChan <- &media.RTMPAACData{ AU: au, PTS: pts, } @@ -89,7 +112,10 @@ func (a *StreamplaceAPI) HandleRTMPPublisher(ctx context.Context, sc *gortmplib. if ctx.Err() != nil { return ctx.Err() } - sc.RW.(net.Conn).SetReadDeadline(time.Now().Add(10 * time.Second)) + err = sc.RW.(net.Conn).SetReadDeadline(time.Now().Add(RTMPTimeout)) + if err != nil { + return err + } err = r.Read() if err != nil { return err @@ -98,19 +124,68 @@ func (a *StreamplaceAPI) HandleRTMPPublisher(ctx context.Context, sc *gortmplib. }) g.Go(func() error { - return a.MediaManager.RTMPIngest(ctx, videoInput, audioInput, mediaSigner) + return a.MediaManager.RTMPIngest(ctx, fmt.Sprintf("rtmp://%s/live/%s", a.rtmpInternalPlaybackAddr, streamer), mediaSigner) }) return g.Wait() } -func (a *StreamplaceAPI) HandleRTMPConnInner(ctx context.Context, conn net.Conn) error { - conn.SetReadDeadline(time.Now().Add(10 * time.Second)) +func (a *StreamplaceAPI) HandleRTMPPlayback(ctx context.Context, sc *gortmplib.ServerConn) error { + if !strings.HasPrefix(sc.URL.Path, RTMPPrefix) { + return fmt.Errorf("RTMP publisher is not allowed to publish to %s (must start with %s)", sc.URL.String(), RTMPPrefix) + } + streamer := strings.TrimPrefix(sc.URL.Path, RTMPPrefix) + a.rtmpSessionsLock.Lock() + session, ok := a.rtmpSessions[streamer] + a.rtmpSessionsLock.Unlock() + if !ok { + return fmt.Errorf("RTMP session not found for streamer %s", streamer) + } + + w := &gortmplib.Writer{ + Conn: sc, + Tracks: []format.Format{session.VideoTrack, session.AudioTrack}, + } + err := w.Initialize() + if err != nil { + return err + } + for { + select { + case <-ctx.Done(): + return ctx.Err() + case event := <-session.EventChan: + if event == nil { + return fmt.Errorf("RTMP session closed") + } + switch event := event.(type) { + case *media.RTMPH264Data: + err := w.WriteH264(session.VideoTrack, event.PTS, event.DTS, event.AU) + if err != nil { + return fmt.Errorf("error writing H264: %w", err) + } + case *media.RTMPAACData: + err := w.WriteMPEG4Audio(session.AudioTrack, event.PTS, event.AU) + if err != nil { + return fmt.Errorf("error writing MPEG4Audio: %w", err) + } + default: + return fmt.Errorf("unsupported event type: %T", event) + } + } + } +} + +func (a *StreamplaceAPI) HandleRTMPPublishConn(ctx context.Context, conn net.Conn) error { + err := conn.SetReadDeadline(time.Now().Add(RTMPTimeout)) + if err != nil { + return err + } sc := &gortmplib.ServerConn{ RW: conn, } - err := sc.Initialize() + err = sc.Initialize() if err != nil { return err } @@ -123,35 +198,114 @@ func (a *StreamplaceAPI) HandleRTMPConnInner(ctx context.Context, conn net.Conn) if sc.Publish { return a.HandleRTMPPublisher(ctx, sc) } - return fmt.Errorf("RTMP playback is not supported") + return fmt.Errorf("RTMP playback is not allowed") } -func (a *StreamplaceAPI) HandleRTMPConn(ctx context.Context, conn net.Conn) { - defer conn.Close() +func (a *StreamplaceAPI) HandleRTMPPlaybackConn(ctx context.Context, conn net.Conn) error { + err := conn.SetReadDeadline(time.Now().Add(RTMPTimeout)) + if err != nil { + return err + } + + sc := &gortmplib.ServerConn{ + RW: conn, + } + err = sc.Initialize() + if err != nil { + return err + } - log.Log(ctx, "connection opened", "remoteAddr", conn.RemoteAddr()) - err := a.HandleRTMPConnInner(ctx, conn) - log.Log(ctx, "connection closed", "remoteAddr", conn.RemoteAddr(), "error", err) + err = sc.Accept() + if err != nil { + return err + } + + if !sc.Publish { + return a.HandleRTMPPlayback(ctx, sc) + } + return fmt.Errorf("RTMP playback is not allowed") } -func (a *StreamplaceAPI) StartRTMPServer(ctx context.Context) error { - ln, err := net.Listen("tcp", ":1935") +func (a *StreamplaceAPI) ServeRTMP(ctx context.Context) error { + ln, err := net.Listen("tcp", a.CLI.RTMPAddr) if err != nil { return fmt.Errorf("failed to listen: %w", err) } defer ln.Close() - log.Log(ctx, "listening on :1935") + log.Log(ctx, "rtmp server starting", "addr", a.CLI.RTMPAddr) + + g, ctx := errgroup.WithContext(ctx) + g.Go(func() error { + return a.ServeRTMPInternalPlayback(ctx) + }) + g.Go(func() error { + for { + if ctx.Err() != nil { + return ctx.Err() + } + conn, err := ln.Accept() + if err != nil { + return fmt.Errorf("error accepting RTMP connection: %w", err) + } + go func() { + err := a.HandleRTMPPublishConn(ctx, conn) + if err != nil { + log.Error(ctx, "error handling RTMP publish connection", "error", err) + } + }() + } + }) + + <-ctx.Done() + + err = ln.Close() + if err != nil { + return fmt.Errorf("failed to close RTMP listener: %w", err) + } + + return g.Wait() +} + +// Serve RTMP internal playback server for gstreamer to pull from +func (a *StreamplaceAPI) ServeRTMPInternalPlayback(ctx context.Context) error { + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + return fmt.Errorf("failed to listen: %w", err) + } + addr := ln.Addr().String() + defer ln.Close() + + _, port, err := net.SplitHostPort(addr) + if err != nil { + return fmt.Errorf("failed to split host and port: %w", err) + } + + a.rtmpInternalPlaybackAddr = fmt.Sprintf("127.0.0.1:%s", port) + + log.Log(ctx, "rtmp internal playback server starting", "addr", a.rtmpInternalPlaybackAddr) // Accept loop in a goroutine so we can select on context.Done go func() { for { + if ctx.Err() != nil { + return + } conn, err := ln.Accept() if err != nil { + if ctx.Err() != nil { + return + } log.Error(ctx, "error accepting RTMP connection", "error", err) + continue } - go a.HandleRTMPConn(ctx, conn) + go func() { + err := a.HandleRTMPPlaybackConn(ctx, conn) + if err != nil { + log.Error(ctx, "error handling RTMP internal playback connection", "error", err) + } + }() } }() diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index 1bfea088..b18b1bb4 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -411,13 +411,16 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { }) if cli.RTMPServerAddon != "" { group.Go(func() error { - return rtmps.ServeRTMPS(ctx, &cli) + return rtmps.ServeRTMPSAddon(ctx, &cli) }) } } else { group.Go(func() error { return a.ServeHTTP(ctx) }) + group.Go(func() error { + return a.ServeRTMP(ctx) + }) } group.Go(func() error { @@ -447,10 +450,6 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { return mod.StartSegmentCleaner(ctx) }) - group.Go(func() error { - return a.StartRTMPServer(ctx) - }) - group.Go(func() error { return replicator.Start(ctx, &cli) }) diff --git a/pkg/config/config.go b/pkg/config/config.go index ecde4150..84106e94 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -65,7 +65,8 @@ type CLI struct { HTTPAddr string HTTPInternalAddr string HTTPSAddr string - RtmpsAddr string + RTMPAddr string + RTMPSAddr string Secure bool NoMist bool MistAdminPort int @@ -208,7 +209,8 @@ func (cli *CLI) NewFlagSet(name string) *flag.FlagSet { fs.IntVar(&cli.RateLimitBurst, "rate-limit-burst", 0, "rate limit burst for requests per ip") fs.IntVar(&cli.RateLimitWebsocket, "rate-limit-websocket", 10, "number of concurrent websocket connections allowed per ip") fs.StringVar(&cli.RTMPServerAddon, "rtmp-server-addon", "", "address of external RTMP server to forward streams to") - fs.StringVar(&cli.RtmpsAddr, "rtmps-addr", ":1935", "address to listen for RTMPS connections") + fs.StringVar(&cli.RTMPSAddr, "rtmps-addr", ":1935", "address to listen for RTMPS connections (when --secure=true)") + fs.StringVar(&cli.RTMPAddr, "rtmp-addr", ":1935", "address to listen for RTMP connections (when --secure=false)") cli.JSONFlag(fs, &cli.DiscordWebhooks, "discord-webhooks", "[]", "JSON array of Discord webhooks to send notifications to") fs.BoolVar(&cli.NewWebRTCPlayback, "new-webrtc-playback", true, "enable new webrtc playback") fs.StringVar(&cli.AppleTeamID, "apple-team-id", "", "apple team id for deep linking") diff --git a/pkg/media/rtmp_ingest.go b/pkg/media/rtmp_ingest.go index d101f3b5..adad5959 100644 --- a/pkg/media/rtmp_ingest.go +++ b/pkg/media/rtmp_ingest.go @@ -6,15 +6,15 @@ import ( "strings" "time" - "github.com/bluenviron/mediacommon/v2/pkg/codecs/h264" + "github.com/bluenviron/gortsplib/v5/pkg/format" "github.com/go-gst/go-gst/gst" - "github.com/go-gst/go-gst/gst/app" "stream.place/streamplace/pkg/log" ) type RTMPH264Data struct { AU [][]byte PTS time.Duration + DTS time.Duration } type RTMPAACData struct { @@ -22,102 +22,36 @@ type RTMPAACData struct { PTS time.Duration } -// ingest a H264+AAC RTMP stream -func (mm *MediaManager) RTMPIngest(ctx context.Context, videoInput chan *RTMPH264Data, audioInput chan *RTMPAACData, ms MediaSigner) error { +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() pipelineSlice := []string{ - "appsrc name=videosrc ! queue ! h264parse name=parse", - "appsrc name=audiosrc ! queue ! fdkaacdec ! audioresample ! opusenc name=audioenc", + fmt.Sprintf("rtmp2src location=%s ! flvdemux name=demux", rtmpURL), + "demux.audio ! queue ! fdkaacdec ! audioresample ! opusenc 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) } - videosrcEle, err := pipeline.GetElementByName("videosrc") - if err != nil { - return err - } - // defer runtime.KeepAlive(srcele) - videosrc := app.SrcFromElement(videosrcEle) - videosrc.SetCaps(gst.NewCapsFromString("video/x-h264,stream-format=byte-stream")) - videosrc.SetCallbacks(&app.SourceCallbacks{ - NeedDataFunc: func(self *app.Source, length uint) { - if ctx.Err() != nil { - self.EndStream() - return - } - - packet := <-videoInput - if packet == nil { - log.Debug(ctx, "video input closed, ending stream") - self.EndStream() - return - } - - // allBytes := bytes.Buffer{} - // for _, au := range packet.AU { - // allBytes.Write(au) - // } - - avcc, err := h264.AnnexB(packet.AU).Marshal() - if err != nil { - log.Error(ctx, "failed to marshal AVCC", "error", err) - self.Error("failed to marshal AVCC", fmt.Errorf("failed to marshal AVCC: %w", err)) - return - } - - buf := gst.NewBufferFromBytes(avcc) - buf.SetPresentationTimestamp(gst.ClockTime(uint64(packet.PTS.Nanoseconds()))) - ret := self.PushBuffer(buf) - if ret != gst.FlowOK { - log.Error(ctx, "failed to push video buffer", "error", ret.String()) - self.Error("failed to push video buffer", fmt.Errorf("failed to push video buffer: %s", ret.String())) - return - } - }, - }) - - audiosrcEle, err := pipeline.GetElementByName("videosrc") + signer, err := mm.SegmentAndSignElem(ctx, ms) if err != nil { return err } - // defer runtime.KeepAlive(srcele) - audiosrc := app.SrcFromElement(audiosrcEle) - audiosrc.SetCallbacks(&app.SourceCallbacks{ - NeedDataFunc: func(self *app.Source, length uint) { - if ctx.Err() != nil { - self.EndStream() - return - } - packet := <-audioInput - if packet == nil { - log.Debug(ctx, "audio input closed, ending stream") - self.EndStream() - return - } - buf := gst.NewBufferFromBytes(packet.AU) - buf.SetPresentationTimestamp(gst.ClockTime(uint64(packet.PTS.Nanoseconds()))) - ret := self.PushBuffer(buf) - if ret != gst.FlowOK { - log.Error(ctx, "failed to push audio buffer", "error", ret.String()) - self.Error("failed to push audio buffer", fmt.Errorf("failed to push audio buffer: %s", ret.String())) - return - } - }, - }) parseEle, err := pipeline.GetElementByName("parse") if err != nil { return err } - signer, err := mm.SegmentAndSignElem(ctx, ms) - if err != nil { - return err - } - err = pipeline.Add(signer) if err != nil { return err @@ -159,3 +93,157 @@ func (mm *MediaManager) RTMPIngest(ctx context.Context, videoInput chan *RTMPH26 return err } + +// // ingest a H264+AAC RTMP stream +// func (mm *MediaManager) RTMPIngest(ctx context.Context, videoInput chan *RTMPH264Data, audioInput chan *RTMPAACData, ms MediaSigner) error { +// ctx, cancel := context.WithCancel(ctx) +// defer cancel() +// pipelineSlice := []string{ +// "appsrc name=videosrc ! queue ! h264parse name=parse", +// "appsrc name=audiosrc ! queue ! fdkaacdec ! audioresample ! opusenc name=audioenc", +// } +// pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) +// if err != nil { +// return fmt.Errorf("error creating RTMPIngest pipeline: %w", err) +// } + +// videosrcEle, err := pipeline.GetElementByName("videosrc") +// if err != nil { +// return err +// } +// first := true +// // defer runtime.KeepAlive(srcele) +// videosrc := app.SrcFromElement(videosrcEle) +// videosrc.SetCaps(gst.NewCapsFromString("video/x-h264,stream-format=avc3")) +// videosrc.SetCallbacks(&app.SourceCallbacks{ +// NeedDataFunc: func(self *app.Source, length uint) { +// if ctx.Err() != nil { +// self.EndStream() +// return +// } + +// packet := <-videoInput +// if packet == nil { +// log.Debug(ctx, "video input closed, ending stream") +// self.EndStream() +// return +// } + +// // allBytes := bytes.Buffer{} +// // for _, au := range packet.AU { +// // allBytes.Write(au) +// // } + +// var avc []byte +// if first { +// c := h264conf.Conf{ +// SPS: packet.AU[0], +// PPS: packet.AU[1], +// } +// avc, err = c.Marshal() +// if err != nil { +// log.Error(ctx, "failed to marshal H264 config", "error", err) +// self.Error("failed to marshal H264 config", fmt.Errorf("failed to marshal H264 config: %w", err)) +// return +// } +// first = false +// } else { +// avc, err = h264.AVCC(packet.AU).Marshal() +// if err != nil { +// log.Error(ctx, "failed to marshal AnnexB", "error", err) +// self.Error("failed to marshal AnnexB", fmt.Errorf("failed to marshal AnnexB: %w", err)) +// return +// } +// } + +// buf := gst.NewBufferFromBytes(avc) +// buf.SetPresentationTimestamp(gst.ClockTime(uint64(packet.PTS.Nanoseconds()))) +// ret := self.PushBuffer(buf) +// if ret != gst.FlowOK { +// log.Error(ctx, "failed to push video buffer", "error", ret.String()) +// self.Error("failed to push video buffer", fmt.Errorf("failed to push video buffer: %s", ret.String())) +// return +// } +// }, +// }) + +// audiosrcEle, err := pipeline.GetElementByName("videosrc") +// if err != nil { +// return err +// } +// // defer runtime.KeepAlive(srcele) +// audiosrc := app.SrcFromElement(audiosrcEle) +// audiosrc.SetCallbacks(&app.SourceCallbacks{ +// NeedDataFunc: func(self *app.Source, length uint) { +// if ctx.Err() != nil { +// self.EndStream() +// return +// } +// packet := <-audioInput +// if packet == nil { +// log.Debug(ctx, "audio input closed, ending stream") +// self.EndStream() +// return +// } +// buf := gst.NewBufferFromBytes(packet.AU) +// buf.SetPresentationTimestamp(gst.ClockTime(uint64(packet.PTS.Nanoseconds()))) +// ret := self.PushBuffer(buf) +// if ret != gst.FlowOK { +// log.Error(ctx, "failed to push audio buffer", "error", ret.String()) +// self.Error("failed to push audio buffer", fmt.Errorf("failed to push audio buffer: %s", ret.String())) +// return +// } +// }, +// }) + +// parseEle, err := pipeline.GetElementByName("parse") +// if err != nil { +// return err +// } + +// signer, err := mm.SegmentAndSignElem(ctx, ms) +// 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 +// } diff --git a/pkg/notifications/firebase.go b/pkg/notifications/firebase.go index a501f062..2dd16740 100644 --- a/pkg/notifications/firebase.go +++ b/pkg/notifications/firebase.go @@ -5,9 +5,10 @@ import ( "encoding/json" "fmt" + "context" + firebase "firebase.google.com/go/v4" "firebase.google.com/go/v4/messaging" - "golang.org/x/net/context" "google.golang.org/api/option" "stream.place/streamplace/pkg/log" ) diff --git a/pkg/rtmps/rtmps.go b/pkg/rtmps/rtmps.go index a5bc324e..3e1d379b 100644 --- a/pkg/rtmps/rtmps.go +++ b/pkg/rtmps/rtmps.go @@ -14,7 +14,7 @@ import ( ) // passthrough RTMPS TLS terminator to external RTMP server -func ServeRTMPS(ctx context.Context, cli *config.CLI) error { +func ServeRTMPSAddon(ctx context.Context, cli *config.CLI) error { if cli.RTMPServerAddon == "" { return fmt.Errorf("RTMP server address not configured") } @@ -29,13 +29,13 @@ func ServeRTMPS(ctx context.Context, cli *config.CLI) error { MinVersion: tls.VersionTLS12, } - listener, err := tls.Listen("tcp", cli.RtmpsAddr, tlsConfig) + listener, err := tls.Listen("tcp", cli.RTMPAddr, tlsConfig) if err != nil { return fmt.Errorf("failed to create RTMPS listener: %w", err) } log.Log(ctx, "rtmps server starting", - "addr", cli.RtmpsAddr, + "addr", cli.RTMPAddr, "forwarding_to", cli.RTMPServerAddon) go func() {