diff --git a/pkg/cmd/combine.go b/pkg/cmd/combine.go index 3b7e82f23..53271e5ed 100644 --- a/pkg/cmd/combine.go +++ b/pkg/cmd/combine.go @@ -12,6 +12,7 @@ import ( "stream.place/streamplace/pkg/gstinit" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/media" + "stream.place/streamplace/pkg/muxl" ) func Combine(ctx context.Context, cli *config.CLI, debugDir string, outFile string, inputs []string) error { @@ -23,41 +24,33 @@ func Combine(ctx context.Context, cli *config.CLI, debugDir string, outFile stri return fmt.Errorf("failed to create debug directory: %w", err) } } - log.Debug(context.Background(), "combine command: starting", "outFile", outFile, "inputs", inputs) ctx = log.WithDebugValue(ctx, cli.Debug) - cryptoSigner, err := createSigner(ctx, cli) - if err != nil { - return err - } - ms, err := media.MakeMediaSigner(ctx, cli, "combine", cryptoSigner, nil) - if err != nil { - return err - } - log.Log(ctx, "combining segments", "outFile", outFile, "inputs", inputs) + outFd, err := os.Create(outFile) if err != nil { return err } defer outFd.Close() - inputFds := make([]io.ReadSeeker, len(inputs)) - for i, input := range inputs { + + readers := make([]io.Reader, 0, len(inputs)) + for _, input := range inputs { fd, err := os.Open(input) if err != nil { return err } defer fd.Close() - inputFds[i] = fd + readers = append(readers, fd) } - err = media.CombineSegments(ctx, inputFds, ms, outFd) - if err != nil { - return err - } - err = CheckCombined(ctx, cli, outFd, debugDir) - if err != nil { - return err + + // Inputs are canonical MUXL segments; concatenated they wrap straight into + // a flat MP4 — one synthesized ftyp+moov over every segment, each segment's + // signature preserved verbatim. No remux, no re-signing. + if err := muxl.RunMuxlWrap(ctx, io.MultiReader(readers...), "flat", outFd); err != nil { + return fmt.Errorf("failed to combine segments: %w", err) } - return nil + + return CheckCombined(ctx, cli, outFd, debugDir) } func CheckCombined(ctx context.Context, cli *config.CLI, inFD io.ReadWriteSeeker, debugDir string) error { diff --git a/pkg/media/clip_user.go b/pkg/media/clip_user.go index fe365c448..04a6d2d4a 100644 --- a/pkg/media/clip_user.go +++ b/pkg/media/clip_user.go @@ -11,8 +11,14 @@ import ( "stream.place/streamplace/pkg/aqtime" "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/localdb" + "stream.place/streamplace/pkg/muxl" ) +// ClipUser concatenates a user's stored canonical .m4s segments over the given +// time window into a single flat MP4 written to writer (used for +// moderation/report clips). The segments are blindly concatenatable MUXL, so +// the clip is just their bytes wrapped in a synthesized ftyp+moov — no remux, +// and each segment's C2PA signature passes through verbatim. func ClipUser(ctx context.Context, localDB localdb.LocalDB, cli *config.CLI, user string, writer io.Writer, before *time.Time, after *time.Time) error { segments, err := localDB.LatestSegmentsForUser(user, -1, false, before, after) if err != nil { @@ -21,14 +27,14 @@ func ClipUser(ctx context.Context, localDB localdb.LocalDB, cli *config.CLI, use if len(segments) == 0 { return fmt.Errorf("no segments found") } - // Sort segments by StartTime, oldest first + // Oldest first, so the clip plays in order. sort.Slice(segments, func(i, j int) bool { return segments[i].StartTime.Before(segments[j].StartTime) }) - segmentFiles := []io.ReadSeeker{} + readers := make([]io.Reader, 0, len(segments)) for _, segment := range segments { aqt := aqtime.FromTime(segment.StartTime) - fpath, err := cli.SegmentFilePath(user, fmt.Sprintf("%s.%s", aqt.FileSafeString(), "mp4")) + fpath, err := cli.SegmentFilePath(user, fmt.Sprintf("%s.%s", aqt.FileSafeString(), "m4s")) if err != nil { return fmt.Errorf("unable to get segment file path: %w", err) } @@ -37,10 +43,11 @@ func ClipUser(ctx context.Context, localDB localdb.LocalDB, cli *config.CLI, use return fmt.Errorf("unable to open segment file: %w", err) } defer fd.Close() - segmentFiles = append(segmentFiles, fd) + readers = append(readers, fd) } - err = CombineSegmentsUnsigned(ctx, segmentFiles, writer, false) - if err != nil { + // muxl wrap unwraps every concatenated segment, aggregates their catalogs, + // and synthesizes one flat MP4 over the lot. + if err := muxl.RunMuxlWrap(ctx, io.MultiReader(readers...), "flat", writer); err != nil { return fmt.Errorf("unable to clip segments: %w", err) } return nil diff --git a/pkg/media/media_signer.go b/pkg/media/media_signer.go index 7495a2a9a..f7e45aa55 100644 --- a/pkg/media/media_signer.go +++ b/pkg/media/media_signer.go @@ -6,7 +6,6 @@ import ( "crypto/ecdsa" "crypto/rand" "crypto/sha256" - "encoding/base64" "encoding/json" "fmt" "io" @@ -17,15 +16,12 @@ import ( "go.opentelemetry.io/otel/trace" "stream.place/streamplace/pkg/aqtime" "stream.place/streamplace/pkg/atproto" - c2patypes "stream.place/streamplace/pkg/c2patypes" "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/crypto/aqpub" "stream.place/streamplace/pkg/crypto/signers" - "stream.place/streamplace/pkg/iroh/generated/iroh_streamplace" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/model" "stream.place/streamplace/pkg/muxl" - "stream.place/streamplace/pkg/spmetrics" ) var signerTracer = otel.Tracer("signer") @@ -39,7 +35,6 @@ type MediaSigner interface { Pub() aqpub.Pub Streamer() string DID() string - SignConcatMP4(ctx context.Context, input io.ReadSeeker, ingredients []io.ReadSeeker, output io.ReadWriteSeeker) error } var DoReplay = false @@ -172,78 +167,6 @@ func (ms *MediaSignerLocal) SignSegmentStream(ctx context.Context, input io.Read return muxl.RunMuxlSignSegment(ctx, input, in, nil, nil, eventCh) } -func (ms *MediaSignerLocal) SignConcatMP4(ctx context.Context, input io.ReadSeeker, ingredients []io.ReadSeeker, output io.ReadWriteSeeker) error { - startTime := time.Now() - ctx, span := otel.Tracer("signer").Start(ctx, "SignMP4") - defer span.End() - // for _, ingredient := range ingredients { - // _, err := iroh_streamplace.GetManifestAndCert(c2patypes.NewReader(aqio.NewReadWriteSeeker(ingredient))) - // if err != nil { - // return nil, err - // } - // } - // title := "livestream" - mani := obj{ - "title": "Livestream Clip", - // "assertions": []obj{ - // { - // "label": "c2pa.actions", - // "data": obj{ - // "actions": []obj{ - // {"action": "c2pa.created"}, - // {"action": "c2pa.published"}, - // }, - // }, - // }, - // { - // "label": StreamplaceMetadata, - // "data": obj{ - // "@context": obj{ - // "dc": "http://purl.org/dc/elements/1.1/", - // }, - // "dc:creator": ms.StreamerName, - // "dc:title": []string{title}, - // "dc:date": []string{aqtime.FromMillis(start).String()}, - // }, - // }, - // }, - } - ctx, span = otel.Tracer("signer").Start(ctx, "SignMP4_MarshalManifest") - manifestBs, err := json.Marshal(mani) - if err != nil { - return fmt.Errorf("failed to marshal manifest: %w", err) - } - var manifest c2patypes.ManifestDefinition - err = json.Unmarshal(manifestBs, &manifest) - if err != nil { - return fmt.Errorf("failed to unmarshal manifest: %w", err) - } - span.End() - - ctx, span = otel.Tracer("signer").Start(ctx, "SignMP4_Sign") - rustCallbackSigner := &RustCallbackSigner{ - Signer: ms.Signer, - } - many := c2patypes.NewManyStreams() - for _, ingredient := range ingredients { - many.AddStream(ingredient) - } - err = iroh_streamplace.SignWithIngredients(string(manifestBs), c2patypes.NewReader(input), base64.StdEncoding.EncodeToString(ms.Cert), many, rustCallbackSigner, c2patypes.NewWriter(output)) - if err != nil { - return err - } - span.End() - - ctx, span = otel.Tracer("signer").Start(ctx, "SignMP4_OutputBytes") - defer ctx.Done() - if err != nil { - return fmt.Errorf("failed to get output bytes: %w", err) - } - span.End() - spmetrics.SigningDuration.WithLabelValues(ms.StreamerName).Observe(float64(time.Since(startTime).Milliseconds())) - return nil -} - // don't call externally! this is used as a callback for the rust library func (ms *MediaSignerLocal) Pub() aqpub.Pub { diff --git a/pkg/media/segment_combine.go b/pkg/media/segment_combine.go deleted file mode 100644 index ceebb4482..000000000 --- a/pkg/media/segment_combine.go +++ /dev/null @@ -1,152 +0,0 @@ -package media - -import ( - "bytes" - "context" - "fmt" - "io" - "strings" - - "github.com/go-gst/go-gst/gst" - "github.com/go-gst/go-gst/gst/app" - "stream.place/streamplace/pkg/aqio" - "stream.place/streamplace/pkg/bus" - "stream.place/streamplace/pkg/log" -) - -// CombineSegments combines a list of segments into a single segment that maintains all of the manifests -func CombineSegments(ctx context.Context, inputFds []io.ReadSeeker, ms MediaSigner, output io.ReadWriteSeeker) error { - rws := aqio.NewReadWriteSeeker([]byte{}) - err := CombineSegmentsUnsigned(ctx, inputFds, rws, true) - if err != nil { - return err - } - // rewind all the inputs for the signer - for _, fd := range inputFds { - _, err := fd.Seek(0, io.SeekStart) - if err != nil { - return err - } - } - bs, err := rws.Bytes() - if err != nil { - return err - } - err = ms.SignConcatMP4(ctx, bytes.NewReader(bs), inputFds, output) - if err != nil { - return err - } - return nil -} - -func CombineSegmentsUnsigned(ctx context.Context, sources []io.ReadSeeker, w io.Writer, doH264Parse bool) error { - ctx = log.WithLogValues(ctx, "mediafunc", "CombineSegmentsUnsigned") - ctx, cancel := context.WithCancel(ctx) - defer cancel() - - pipelineSlice := []string{ - fmt.Sprintf("mp4mux name=muxer faststart=true interleave-bytes=%d interleave-time=%d movie-timescale=60000 trak-timescale=60000 ! appsink sync=false name=mp4sink", InterleaveBytes, InterleaveTime), - "capsfilter caps=video/x-h264,parsed=true name=videoqueue ! queue ! muxer.", - "capsfilter caps=audio/x-opus,framed=true name=audioparse ! queue ! muxer.", - } - - pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) - if err != nil { - return fmt.Errorf("failed to create GStreamer pipeline: %w", err) - } - - segCh := make(chan *bus.Seg) - go func() { - for _, source := range sources { - bs, err := io.ReadAll(source) - if err != nil { - err = fmt.Errorf("failed to read file: %w", err) - pipeline.Error(err.Error(), err) - return - } - segCh <- &bus.Seg{ - Filepath: "ignored", - Data: bs, - } - } - close(segCh) - }() - - concatBin, err := ConcatBin(ctx, segCh, doH264Parse) - if err != nil { - return fmt.Errorf("failed to create concat bin: %w", err) - } - - err = pipeline.Add(concatBin.Element) - if err != nil { - return fmt.Errorf("failed to add concat bin to pipeline: %w", err) - } - - videoPad := concatBin.GetStaticPad("video_0") - if videoPad == nil { - return fmt.Errorf("video pad not found") - } - - audioPad := concatBin.GetStaticPad("audio_0") - if audioPad == nil { - return fmt.Errorf("audio pad not found") - } - - // Get the videoparse and audioparse elements from the pipeline - videoQueue, err := pipeline.GetElementByName("videoqueue") - if err != nil { - return fmt.Errorf("failed to get video parse element: %w", err) - } - - audioParse, err := pipeline.GetElementByName("audioparse") - if err != nil { - return fmt.Errorf("failed to get audio parse element: %w", err) - } - - // Link the concat bin pads to the parse element sink pads - linked := videoPad.Link(videoQueue.GetStaticPad("sink")) - if linked != gst.PadLinkOK { - return fmt.Errorf("failed to link video pad to video parse element: %v", linked) - } - - linked = audioPad.Link(audioParse.GetStaticPad("sink")) - if linked != gst.PadLinkOK { - return fmt.Errorf("failed to link audio pad to audio parse element: %v", linked) - } - - // Get the mp4sink element and set up its callback - mp4Sink, err := pipeline.GetElementByName("mp4sink") - if err != nil { - return fmt.Errorf("failed to get mp4sink element: %w", err) - } - - appSink := app.SinkFromElement(mp4Sink) - appSink.SetCallbacks(&app.SinkCallbacks{ - NewSampleFunc: WriterNewSample(ctx, w), - }) - - errCh := make(chan error) - go func() { - err := HandleBusMessages(ctx, pipeline) - 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.BlockSetState(gst.StateNull) - if err != nil { - log.Error(ctx, "failed to set pipeline state to null", "error", err) - } - }() - - err = <-errCh - if err != nil { - return fmt.Errorf("pipeline error: %w", err) - } - - return nil -} diff --git a/pkg/media/segment_combine_test.go b/pkg/media/segment_combine_test.go deleted file mode 100644 index 1e944af18..000000000 --- a/pkg/media/segment_combine_test.go +++ /dev/null @@ -1,53 +0,0 @@ -package media - -import ( - "context" - "fmt" - "io" - "os" - "testing" - - "github.com/stretchr/testify/require" - "golang.org/x/sync/errgroup" - "stream.place/streamplace/pkg/aqio" - "stream.place/streamplace/pkg/log" - "stream.place/streamplace/test/remote" -) - -func TestCombineSegmentsUnsigned(t *testing.T) { - withNoGSTLeaks(t, func() { - g, _ := errgroup.WithContext(context.Background()) - for range streamplaceTestCount { - g.Go(func() error { - return innerTestClip(t) - }) - } - err := g.Wait() - require.NoError(t, err) - }) -} - -func innerTestClip(t *testing.T) error { - ctx := log.WithDebugValue(context.Background(), map[string]map[string]int{"func": {"ConcatDemuxBin": 9, "ConcatBin": 9}}) - dirname := remote.RemoteArchive("c21e9352e72ca0729c66af2fcabec1b8997b509601241e8d38d5728f9687386b/threesegs.tar.gz") - inputFiles := []string{ - fmt.Sprintf("%s/2025-11-15T21-05-00-399Z.mp4", dirname), - fmt.Sprintf("%s/2025-11-15T21-05-01-385Z.mp4", dirname), - fmt.Sprintf("%s/2025-11-15T21-05-02-393Z.mp4", dirname), - } - inputFds := make([]io.ReadSeeker, len(inputFiles)) - for i, fName := range inputFiles { - fd, err := os.Open(fName) - if err != nil { - return fmt.Errorf("unable to open segment file: %w", err) - } - inputFds[i] = fd - } - buf := aqio.NewReadWriteSeeker([]byte{}) - err := CombineSegmentsUnsigned(ctx, inputFds, buf, true) - require.NoError(t, err) - slice, err := buf.Bytes() - require.NoError(t, err) - require.Equal(t, 4725181, len(slice)) - return nil -}