diff --git a/.gitignore b/.gitignore index ee1ff2bf7..999339923 100644 --- a/.gitignore +++ b/.gitignore @@ -14,3 +14,4 @@ description.md lerna-debug.log build-* .ci/build.env +*.log diff --git a/Makefile b/Makefile index b1afed20a..96a3917c6 100644 --- a/Makefile +++ b/Makefile @@ -119,22 +119,33 @@ OPTS = -D "gst-plugins-base:audioresample=enabled" \ -D "gst-plugins-base:playback=enabled" \ -D "gst-plugins-base:opus=enabled" \ -D "gst-plugins-base:gio-typefinder=enabled" \ + -D "gst-plugins-base:videotestsrc=enabled" \ + -D "gst-plugins-base:videoconvertscale=enabled" \ -D "gst-plugins-base:typefind=enabled" \ + -D "gst-plugins-base:compositor=enabled" \ + -D "gst-plugins-base:videorate=enabled" \ + -D "gst-plugins-base:app=enabled" \ -D "gst-plugins-good:matroska=enabled" \ -D "gst-plugins-good:multifile=enabled" \ -D "gst-plugins-bad:fdkaac=enabled" \ -D "gst-plugins-bad:hls=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-bad:videoparsers=enabled" \ -D "gst-plugins-bad:mpegtsmux=enabled" \ + -D "gst-plugins-ugly:x264=enabled" \ + -D "gst-plugins-ugly:gpl=enabled" \ + -D "x264:asm=enabled" \ -D "gstreamer-full:gst-full=enabled" \ - -D "gstreamer-full:gst-full-plugins=libgstaudioresample.a;libgstmatroska.a;libgstmultifile.a;libgstfdkaac.a;libgsthls.a;libgstopus.a;libgstvideoparsersbad.a;libgstaudioparsers.a;libgstmpegtsmux.a;libgstplayback.a;libgsttypefindfunctions.a" \ + -D "gstreamer-full:gst-full-plugins=libgstaudioresample.a;libgstmatroska.a;libgstmultifile.a;libgstfdkaac.a;libgstisomp4.a;libgstapp.a;libgstvideoconvertscale.a;libgstvideobox.a;libgstvideorate.a;libgstpng.a;libgstcompositor.a;libgsthls.a;libgstx264.a;libgstopus.a;libgstvideotestsrc.a;libgstvideoparsersbad.a;libgstaudioparsers.a;libgstmpegtsmux.a;libgstplayback.a;libgsttypefindfunctions.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-target-type=static_library" \ - -D "gstreamer-full:gst-full-elements=coreelements:fdsrc,fdsink,queue,queue2,typefind,tee,filesink" \ + -D "gstreamer-full:gst-full-elements=coreelements:fdsrc,filesrc,fdsink,filesink,queue,queue2,typefind,tee,filesink,capsfilter" \ -D "gstreamer-full:bad=enabled" \ - -D "gstreamer-full:ugly=disabled" \ -D "gstreamer-full:tls=disabled" \ + -D "gstreamer-full:ugly=enabled" \ -D "gstreamer-full:gpl=enabled" \ -D "gstreamer-full:gst-full-typefind-functions=" @@ -188,7 +199,7 @@ windows-amd64: # unbuffer here is a workaround for wine trying to pop up a terminal window and failing .PHONY: windows-amd64-startup-test windows-amd64-startup-test: - bash -c 'set -euo pipefail && unbuffer wine64 ./build-windows-amd64/aquareum.exe --version | cat' + bash -c 'set -euo pipefail && unbuffer wine64 ./build-windows-amd64/aquareum.exe self-test | cat' .PHONY: node-all-platforms-macos node-all-platforms-macos: app @@ -199,6 +210,7 @@ node-all-platforms-macos: app && tar -czvf ../bin/aquareum-$(VERSION)-darwin-arm64.tar.gz ./aquareum \ && cd - ./build-darwin-arm64/aquareum --version + ./build-darwin-arm64/aquareum self-test rustup target add x86_64-apple-darwin meson setup --buildtype debugoptimized --cross-file util/darwin-amd64-apple.ini build-darwin-amd64 $(OPTS) meson compile -C build-darwin-amd64 @@ -207,6 +219,7 @@ node-all-platforms-macos: app && tar -czvf ../bin/aquareum-$(VERSION)-darwin-amd64.tar.gz ./aquareum \ && cd - ./build-darwin-amd64/aquareum --version + ./build-darwin-arm64/aquareum self-test $(MAKE) desktop-macos meson test -C build-darwin-arm64 go-tests diff --git a/docker/release.Dockerfile b/docker/release.Dockerfile index 3bf8556ff..5ce017daf 100644 --- a/docker/release.Dockerfile +++ b/docker/release.Dockerfile @@ -5,5 +5,5 @@ RUN apt update && apt install -y curl ARG AQUAREUM_URL ENV AQUAREUM_URL $AQUAREUM_URL RUN echo "downloading $AQUAREUM_URL" && cd /usr/local/bin && curl -L "$AQUAREUM_URL" | tar xzv -RUN aquareum --version +RUN aquareum self-test CMD aquareum diff --git a/go.mod b/go.mod index edad79c30..82993716a 100644 --- a/go.mod +++ b/go.mod @@ -112,6 +112,7 @@ require ( github.com/sergi/go-diff v1.3.2-0.20230802210424-5b0b94c5c0d3 // indirect github.com/sirupsen/logrus v1.9.3 // indirect github.com/skeema/knownhosts v1.2.2 // indirect + github.com/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e // indirect github.com/stretchr/objx v0.5.2 // indirect github.com/supranational/blst v0.3.11 // indirect github.com/thales-e-security/pool v0.0.2 // indirect diff --git a/go.sum b/go.sum index ae0590af3..e6b4b3e55 100644 --- a/go.sum +++ b/go.sum @@ -306,6 +306,8 @@ github.com/sirupsen/logrus v1.9.3 h1:dueUQJ1C2q9oE3F7wvmSGAaVtTmUizReu6fjN8uqzbQ github.com/sirupsen/logrus v1.9.3/go.mod h1:naHLuLoDiP4jHNo9R0sCBMtWGeIprob74mVsIT4qYEQ= github.com/skeema/knownhosts v1.2.2 h1:Iug2P4fLmDw9f41PB6thxUkNUkJzB5i+1/exaj40L3A= github.com/skeema/knownhosts v1.2.2/go.mod h1:xYbVRSPxqBZFrdmDyMmsOs+uX1UZC3nTN3ThzgDxUwo= +github.com/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e h1:MRM5ITcdelLK2j1vwZ3Je0FKVCfqOLp5zO6trqMLYs0= +github.com/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e/go.mod h1:XV66xRDqSt+GTGFMVlhk3ULuV0y9ZmzeVGR4mloJI3M= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw= github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpEOglKo= diff --git a/meson.build b/meson.build index 75284fe37..b640b3bed 100644 --- a/meson.build +++ b/meson.build @@ -6,7 +6,7 @@ project( 'cpp_std': 'c++11', 'default_library': 'static', 'auto_features': 'disabled', - 'force_fallback_for': 'glib-2.0,gobject-2.0,gio-2.0,gio-unix-2.0,fdk-aac,zlib,libffi,pcre2,intl', + 'force_fallback_for': 'glib-2.0,gobject-2.0,gio-2.0,gio-unix-2.0,fdk-aac,zlib,libffi,pcre2,intl,x264', 'buildtype': 'debug', }, ) @@ -194,6 +194,7 @@ aquareum_deps += [ gst_full_dep, dependency('gio-2.0'), dependency('gstreamer-controller-1.0'), + dependency('gstreamer-app-1.0'), dependency('gstreamer-pbutils-1.0'), dependency('gstplayback'), ] diff --git a/pkg/api/api_internal.go b/pkg/api/api_internal.go index 7be1f77ba..b2b9870bb 100644 --- a/pkg/api/api_internal.go +++ b/pkg/api/api_internal.go @@ -2,24 +2,25 @@ package api import ( "bufio" + "bytes" "context" "encoding/base64" "fmt" "io" "log/slog" "net/http" + "os" "regexp" + "runtime/pprof" "strconv" "strings" "time" "aquareum.tv/aquareum/pkg/errors" "aquareum.tv/aquareum/pkg/log" - "aquareum.tv/aquareum/pkg/media" "aquareum.tv/aquareum/pkg/mist/mistconfig" "aquareum.tv/aquareum/pkg/mist/misttriggers" v0 "aquareum.tv/aquareum/pkg/schema/v0" - "github.com/google/uuid" "github.com/julienschmidt/httprouter" sloghttp "github.com/samber/slog-http" "golang.org/x/sync/errgroup" @@ -38,18 +39,10 @@ func (a *AquareumAPI) ServeInternalHTTP(ctx context.Context) error { } // lightweight way to authenticate push requests to ourself -var secretUUID string var mkvRE *regexp.Regexp func init() { - uu, err := uuid.NewV7() - if err != nil { - panic(err) - } - secretUUID = uu.String() - mkvRE = regexp.MustCompile(`^\d+\.mkv$`) - } func (a *AquareumAPI) InternalHandler(ctx context.Context) (http.Handler, error) { @@ -181,14 +174,16 @@ func (a *AquareumAPI) InternalHandler(ctx context.Context) (http.Handler, error) w.WriteHeader(200) }) + // self-destruct code, useful for dumping goroutines on windows + router.POST("/abort", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { + pprof.Lookup("goroutine").WriteTo(os.Stderr, 2) + log.Log(ctx, "got POST /abort, self-destructing") + os.Exit(1) + }) + // internal route called for each pushed segment from ffmpeg - router.POST("/segment/:uuid/:user/:file", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { + router.POST("/segment/:user/:file", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { ms := time.Now().UnixMilli() - uu := p.ByName("uuid") - if uu != secretUUID { - errors.WriteHTTPForbidden(w, "unable to authenticate internal url", nil) - return - } user := p.ByName("user") if user == "" { log.Log(ctx, "invalid code path: got empty user?") @@ -202,7 +197,9 @@ func (a *AquareumAPI) InternalHandler(ctx context.Context) (http.Handler, error) return } ctx := log.WithLogValues(ctx, "user", user, "file", p.ByName("file"), "time", fmt.Sprintf("%d", ms)) - err := a.MediaManager.SignSegment(ctx, r.Body, ms) + buf := &bytes.Buffer{} + io.Copy(buf, r.Body) + err := a.MediaManager.SignSegment(ctx, bytes.NewReader(buf.Bytes()), ms) if err != nil { log.Log(ctx, "segment error", "error", err) errors.WriteHTTPInternalServerError(w, "segment error", err) @@ -210,16 +207,9 @@ func (a *AquareumAPI) InternalHandler(ctx context.Context) (http.Handler, error) } }) - // route to accept an incoming mkv stream from OBS, segment it, and push the segments back to this HTTP handler - router.POST("/stream/:key", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { + handleIncomingStream := func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { log.Log(ctx, "stream start") - user, err := a.keyToUser(ctx, p.ByName("key")) - if err != nil { - errors.WriteHTTPForbidden(w, "unable to authenticate stream key", err) - return - } - prefix := fmt.Sprintf("%s/segment/%s/%s", a.CLI.OwnInternalURL(), secretUUID, user) - err = media.SegmentToHTTP(ctx, r.Body, prefix) + err := a.MediaManager.IngestStream(ctx, r.Body) if err != nil { log.Log(ctx, "stream error", "error", err) @@ -227,7 +217,11 @@ func (a *AquareumAPI) InternalHandler(ctx context.Context) (http.Handler, error) return } log.Log(ctx, "stream success", "url", r.URL.String()) - }) + } + + // route to accept an incoming mkv stream from OBS, segment it, and push the segments back to this HTTP handler + router.POST("/stream/:key", handleIncomingStream) + router.PUT("/stream/:key", handleIncomingStream) handler := sloghttp.Recovery(router) handler = sloghttp.New(slog.Default())(handler) return handler, nil diff --git a/pkg/cmd/aquareum.go b/pkg/cmd/aquareum.go index 2e94923d3..8e5269462 100644 --- a/pkg/cmd/aquareum.go +++ b/pkg/cmd/aquareum.go @@ -45,6 +45,16 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { return Stream(os.Args[2]) } + if len(os.Args) > 1 && os.Args[1] == "self-test" { + err := media.RunSelfTest(context.Background()) + if err != nil { + fmt.Println(err.Error()) + os.Exit(1) + } + fmt.Println("self-test successful!") + os.Exit(0) + } + fs := flag.NewFlagSet("aquareum", flag.ExitOnError) cli := config.CLI{Build: build} fs.StringVar(&cli.DataDir, "data-dir", config.DefaultDataDir(), "directory for keeping all aquareum data") @@ -73,13 +83,14 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { fs.StringVar(&cli.StreamerName, "streamer-name", "", "name of the person streaming from this aquareum node") cli.AddressSliceFlag(fs, &cli.AllowedStreams, "allowed-streams", "", "comma-separated list of addresses that this node will replicate") cli.StringSliceFlag(fs, &cli.Peers, "peers", "", "other aquareum nodes to replicate to") + fs.BoolVar(&cli.TestStream, "test-stream", false, "run a built-in test stream on boot") fs.Bool("insecure", false, "DEPRECATED, does nothing.") version := fs.Bool("version", false, "print version and exit") if runtime.GOOS == "linux" { - fs.BoolVar(&cli.NoMist, "no-mist", false, "Disable MistServer") + fs.BoolVar(&cli.NoMist, "no-mist", true, "Disable MistServer") fs.IntVar(&cli.MistAdminPort, "mist-admin-port", 14242, "MistServer admin port (internal use only)") fs.IntVar(&cli.MistRTMPPort, "mist-rtmp-port", 11935, "MistServer RTMP port (internal use only)") fs.IntVar(&cli.MistHTTPPort, "mist-http-port", 18080, "MistServer HTTP port (internal use only)") @@ -240,6 +251,12 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { return a.ServeInternalHTTP(ctx) }) + if cli.TestStream { + group.Go(func() error { + return mm.TestSource(ctx) + }) + } + for _, job := range platformJobs { group.Go(func() error { return job(ctx, &cli) diff --git a/pkg/config/config.go b/pkg/config/config.go index 4e4750d95..636074a54 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -72,6 +72,7 @@ type CLI struct { StreamerName string AllowedStreams []aqpub.Pub Peers []string + TestStream bool dataDirFlags []*string } diff --git a/pkg/media/gstreamer.go b/pkg/media/gstreamer.go index 6bf2214b0..8eb758f4e 100644 --- a/pkg/media/gstreamer.go +++ b/pkg/media/gstreamer.go @@ -1,17 +1,24 @@ package media import ( + "bytes" "context" + "encoding/json" + "errors" "fmt" "io" "os" "path/filepath" "runtime" "strings" + "time" "aquareum.tv/aquareum/pkg/log" + "aquareum.tv/aquareum/test" "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" ) @@ -129,6 +136,96 @@ func AddOpusToMKV(ctx context.Context, input io.Reader, output io.Writer) error return g.Wait() } +// basic test to make sure gstreamer functionality is working +func SelfTest(ctx context.Context) error { + ctx, cancel := context.WithTimeout(ctx, 5*time.Second) + defer cancel() + f, err := test.Files.Open("fixtures/sample-segment.mp4") + if err != nil { + return err + } + defer f.Close() + bs, err := io.ReadAll(f) + if err != nil { + return err + } + + pipeline, err := gst.NewPipelineFromString("appsrc name=src ! appsink name=sink") + if err != nil { + return err + } + + srcele, err := pipeline.GetElementByName("src") + if err != nil { + return err + } + if srcele == nil { + return fmt.Errorf("srcele not found") + } + 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) + self.PushBuffer(buffer) + self.EndStream() + }, + }) + + mainLoop := glib.NewMainLoop(glib.MainContextDefault(), false) + + output := &bytes.Buffer{} + sinkele, err := pipeline.GetElementByName("sink") + if err != nil { + return err + } + if sinkele == nil { + return fmt.Errorf("sinkele not found") + } + 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) { + cancel() + }, + }) + + // Start the pipeline + pipeline.SetState(gst.StatePlaying) + + go func() { + <-ctx.Done() + mainLoop.Quit() + }() + + mainLoop.Run() + + if err != nil { + return err + } + if len(output.Bytes()) < 1 { + return fmt.Errorf("got a zero-byte buffer from SelfTest") + } + return nil +} + func ToHLS(ctx context.Context, input io.Reader, dir string) error { ir, iw, idone, err := SafePipe() if err != nil { @@ -197,11 +294,249 @@ func ToHLS(ctx context.Context, input io.Reader, dir string) error { return nil }) - g.Go(func() error { - runtime.GC() - log.Log(ctx, "output copy complete", "error", err) + return g.Wait() +} + +func (mm *MediaManager) IngestStream(ctx context.Context, input io.Reader) error { + pipelineSlice := []string{ + "appsrc name=streamsrc ! matroskademux name=demux ! h264parse name=parse", + } + pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) + if err != nil { + return fmt.Errorf("error creating IngestStream pipeline: %w", err) + } + defer runtime.KeepAlive(pipeline) + srcele, err := pipeline.GetElementByName("streamsrc") + if err != nil { return err + } + // defer runtime.KeepAlive(srcele) + src := app.SrcFromElement(srcele) + src.SetCallbacks(&app.SourceCallbacks{ + NeedDataFunc: func(self *app.Source, length uint) { + bs := make([]byte, length) + read, err := input.Read(bs) + if err != nil { + if errors.Is(err, io.EOF) { + if read > 0 { + panic("got data on eof???") + } + log.Log(ctx, "EOF, ending stream", "length", read) + self.EndStream() + return + } else { + panic(err) + } + } + toPush := bs + if uint(read) < length { + toPush = bs[:read] + } + buffer := gst.NewBufferWithSize(int64(len(toPush))) + buffer.Map(gst.MapWrite).WriteData(toPush) + self.PushBuffer(buffer) + }, + }) + parseEle, err := pipeline.GetElementByName("parse") + if err != nil { + return err + } + // defer runtime.KeepAlive(parseEle) + signer, err := mm.SegmentAndSignElem(ctx) + if err != nil { + return err + } + // defer runtime.KeepAlive(signer) + pipeline.Add(signer) + err = parseEle.Link(signer) + if err != nil { + return err + } + + mainLoop := glib.NewMainLoop(glib.MainContextDefault(), false) + + pipeline.GetPipelineBus().AddWatch(func(msg *gst.Message) bool { + switch msg.Type() { + + case gst.MessageEOS: // When end-of-stream is received flush the pipeling and stop the main loop + mainLoop.Quit() + case gst.MessageError: // Error messages are always fatal + err := msg.ParseError() + log.Log(ctx, "gstreamer error", "error", err.Error()) + if debug := err.DebugString(); debug != "" { + log.Log(ctx, "gstreamer debug", "message", debug) + } + mainLoop.Quit() + // default: + // log.Log(ctx, msg.String()) + } + return true + }) + + pipeline.SetState(gst.StatePlaying) + + mainLoop.Run() + + return nil +} + +const TESTSRC_WIDTH = 1280 +const TESTSRC_HEIGHT = 720 +const QR_SIZE = 256 + +type QRData struct { + Now int64 `json:"now"` +} + +func (mm *MediaManager) TestSource(ctx context.Context) error { + mainLoop := glib.NewMainLoop(glib.MainContextDefault(), false) + + pipelineSlice := []string{ + "h264parse name=parser", + "compositor name=comp ! videoconvert ! x264enc speed-preset=ultrafast key-int-max=30 ! parser.", + fmt.Sprintf(`videotestsrc is-live=true ! video/x-raw,format=AYUV,framerate=30/1,width=%d,height=%d ! comp.`, TESTSRC_WIDTH, TESTSRC_HEIGHT), + fmt.Sprintf("videobox border-alpha=0 top=-%d left=-%d name=box ! comp.", (TESTSRC_HEIGHT/2)-(QR_SIZE/2), (TESTSRC_WIDTH/2)-(QR_SIZE/2)), + "appsrc name=pngsrc ! pngdec ! videoconvert ! videorate ! video/x-raw,format=AYUV,framerate=1/1 ! box.", + } + + pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) + if err != nil { + return fmt.Errorf("error creating TestSource pipeline: %w", err) + } + + pngele, err := pipeline.GetElementByName("pngsrc") + if err != nil { + return err + } + if pngele == nil { + return fmt.Errorf("pngsrc not found") + } + + parseele, err := pipeline.GetElementByName("parser") + if err != nil { + return err + } + if parseele == nil { + return fmt.Errorf("parseele not found") + } + + signer, err := mm.SegmentAndSignElem(ctx) + if err != nil { + return err + } + pipeline.Add(signer) + err = parseele.Link(signer) + if err != nil { + return fmt.Errorf("link to signer failed: %w", 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) + self.PushBuffer(buffer) + }, + }) + ctx, cancel := context.WithCancel(ctx) + defer cancel() + go func() { + <-ctx.Done() + pipeline.BlockSetState(gst.StateNull) + mainLoop.Quit() + }() + + // Add a message handler to the pipeline bus, printing interesting information to the console. + pipeline.GetPipelineBus().AddWatch(func(msg *gst.Message) bool { + switch msg.Type() { + + case gst.MessageEOS: // When end-of-stream is received flush the pipeling and stop the main loop + cancel() + case gst.MessageError: // Error messages are always fatal + err := msg.ParseError() + log.Log(ctx, "gstreamer error", "error", err.Error()) + if debug := err.DebugString(); debug != "" { + log.Log(ctx, "gstreamer debug", "message", debug) + } + cancel() + // default: + // log.Log(ctx, msg.String()) + } + return true + }) + + // Start the pipeline + pipeline.SetState(gst.StatePlaying) + + g, _ := errgroup.WithContext(ctx) + + g.Go(func() error { + mainLoop.Run() + log.Log(ctx, "main loop complete") + return nil }) return g.Wait() } + +// element that takes the input stream, muxes to mp4, and signs the result +func (mm *MediaManager) SegmentAndSignElem(ctx context.Context) (*gst.Element, error) { + // elem, err := gst.NewElement("splitmuxsink name=splitter async-finalize=true sink-factory=appsink muxer-factory=matroskamux max-size-bytes=1") + elem, err := gst.NewElementWithProperties("splitmuxsink", map[string]any{ + "name": "signer", + "async-finalize": true, + "sink-factory": "appsink", + "muxer-factory": "mp4mux", + "max-size-bytes": 1, + }) + if err != nil { + return nil, err + } + + elem.Connect("sink-added", func(split, sinkEle *gst.Element) { + buf := &bytes.Buffer{} + appsink := app.SinkFromElement(sinkEle) + if appsink == nil { + panic("appsink should not be nil") + } + appsink.SetCallbacks(&app.SinkCallbacks{ + NewSampleFunc: func(sink *app.Sink) gst.FlowReturn { + sample := sink.PullSample() + if sample == nil { + return gst.FlowOK + } + sample.Ref() + defer sample.Unref() + + // Retrieve the buffer from the sample. + buffer := sample.GetBuffer() + + _, err := io.Copy(buf, buffer.Reader()) + + if err != nil { + panic(err) + } + + return gst.FlowOK + }, + EOSFunc: func(sink *app.Sink) { + err := mm.SignSegment(ctx, bytes.NewReader(buf.Bytes()), time.Now().UnixMilli()) + if err != nil { + log.Log(ctx, "error signing segment", "error", err) + } + }, + }) + }) + + return elem, nil +} diff --git a/pkg/media/media.go b/pkg/media/media.go index 1ebc85ac4..9ce105385 100644 --- a/pkg/media/media.go +++ b/pkg/media/media.go @@ -53,8 +53,17 @@ type HLSStream struct { Wait func() string } +func RunSelfTest(ctx context.Context) error { + gst.Init(nil) + return SelfTest(ctx) +} + func MakeMediaManager(ctx context.Context, cli *config.CLI, signer crypto.Signer, rep replication.Replicator) (*MediaManager, error) { gst.Init(nil) + err := SelfTest(ctx) + if err != nil { + return nil, fmt.Errorf("error in gstreamer self-test: %w", err) + } hex := signers.HexAddr(signer.Public().(*ecdsa.PublicKey)) exists, err := cli.DataFileExists([]string{hex, CERT_FILE}) if err != nil { @@ -86,16 +95,10 @@ func MakeMediaManager(ctx context.Context, cli *config.CLI, signer crypto.Signer }, nil } -// accept an incoming mkv segment, mux to mp4, and sign it -func (mm *MediaManager) SignSegment(ctx context.Context, input io.Reader, ms int64) error { - buf := bytes.Buffer{} - err := MuxToMP4(ctx, input, &buf) - if err != nil { - return fmt.Errorf("error muxing to mp4: %w", err) - } - reader := bytes.NewReader(buf.Bytes()) +// accept an incoming mkv, and sign it +func (mm *MediaManager) SignSegment(ctx context.Context, input io.ReadSeeker, ms int64) error { rws := &aqio.ReadWriteSeeker{} - err = mm.SignMP4(ctx, reader, rws, ms) + err := mm.SignMP4(ctx, input, rws, ms) if err != nil { return fmt.Errorf("error signing mp4: %w", err) } @@ -191,60 +194,10 @@ func MuxToMP4(ctx context.Context, input io.Reader, output io.Writer) error { return err } of.Close() - status, info, err := ffmpeg.GetCodecInfo(oname) - if err != nil { - return fmt.Errorf("error in GetCodecInfo: %w", err) - } - fmt.Printf("%v %v\n", status, info.DurSecs) // log.Log(ctx, "transmuxing complete", "out-file", oname, "wrote", written) return nil } -func SegmentToHTTP(ctx context.Context, input io.Reader, prefix string) error { - tc := ffmpeg.NewTranscoder() - defer tc.StopTranscoder() - ir, iw, idone, err := SafePipe() - if err != nil { - return fmt.Errorf("error opening pipe: %w", err) - } - defer idone() - out := []ffmpeg.TranscodeOptions{ - { - Oname: fmt.Sprintf("%s/%%d.mkv", prefix), - VideoEncoder: ffmpeg.ComponentOptions{ - Name: "copy", - }, - AudioEncoder: ffmpeg.ComponentOptions{ - Name: "copy", - }, - Profile: ffmpeg.VideoProfile{Format: ffmpeg.FormatNone}, - Muxer: ffmpeg.ComponentOptions{ - Name: "stream_segment", - Opts: map[string]string{ - "segment_time": "0.1", - }, - }, - }, - } - iname := fmt.Sprintf("pipe:%d", ir.Fd()) - in := &ffmpeg.TranscodeOptionsIn{Fname: iname, Transmuxing: true} - g, _ := errgroup.WithContext(ctx) - g.Go(func() error { - _, err := io.Copy(iw, input) - // log.Log(ctx, "input copy done", "error", err) - iw.Close() - return err - }) - g.Go(func() error { - _, err = tc.Transcode(in, out) - // log.Log(ctx, "transcode done", "error", err) - tc.StopTranscoder() - ir.Close() - return err - }) - return g.Wait() -} - func (mm *MediaManager) SegmentToMKVPlusOpus(ctx context.Context, user string, w io.Writer) error { muxer := ffmpeg.ComponentOptions{ Name: "matroska", @@ -566,5 +519,6 @@ func (mm *MediaManager) ValidateMP4(ctx context.Context, input io.Reader) error io.Copy(fd, r) base := filepath.Base(fd.Name()) go mm.PublishSegment(ctx, mm.user, base) + log.Log(ctx, "successfully ingested segment", "user", pub.String(), "timestamp", meta.StartTime) return nil } diff --git a/subprojects/gstreamer-full.wrap b/subprojects/gstreamer-full.wrap index 7abf0ba66..dfb992d42 100644 --- a/subprojects/gstreamer-full.wrap +++ b/subprojects/gstreamer-full.wrap @@ -1,4 +1,4 @@ [wrap-git] -url = https://gitlab.freedesktop.org/gstreamer/gstreamer.git -revision = 95ca7014c88e5538b0b0380f41cd15904a34fc53 +url = https://gitlab.freedesktop.org/amyspark/gstreamer.git +revision = 05f1896f667102d036c30e25a95e3912081c3cd6 depth = 1 diff --git a/test/fixtures.go b/test/fixtures.go new file mode 100644 index 000000000..d678c54ec --- /dev/null +++ b/test/fixtures.go @@ -0,0 +1,6 @@ +package test + +import "embed" + +//go:embed fixtures/** +var Files embed.FS