Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
9.5 kB · 308 lines
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309package media
import ( "bytes" "context" "fmt" "io" "os" "strings" "time"
"github.com/go-gst/go-gst/gst" "github.com/go-gst/go-gst/gst/app" "go.opentelemetry.io/otel" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/trace" "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/spmetrics")
// For testing. Normally, We don't want to stop the pipeline upon a// segmentation error because we want to keep the stream alive. Lots// of weird invalid data coming in from WebRTC connections on phones.// Better we drop one weird segment than force the stream to restart.// But for tests, we want (sometimes) to know if there's a problem.var FatalSegmentationErrors = false
// element that takes the input stream, muxes to mp4, and signs the resultfunc SegmentElem(ctx context.Context, cli *config.CLI, streamer string, doH264Parse bool, cb func(ctx context.Context, buf []byte, now int64) error) (*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 }
p := elem.GetRequestPad("video") if p == nil { return nil, fmt.Errorf("failed to get video pad") } p = elem.GetRequestPad("audio_%u") if p == nil { return nil, fmt.Errorf("failed to get audio pad") }
resetTimer := make(chan struct{})
go func() { for { select { case <-ctx.Done(): return case <-resetTimer: continue case <-time.After(time.Second * 30): log.Warn(ctx, "no new segment for 30 seconds") elem.ErrorMessage(gst.DomainCore, gst.CoreErrorFailed, "No new segment for 30 seconds", "No new segment for 30 seconds (debug)") return } } }()
// we didn't need faststart but i'm leaving this commented here in case // you want to change any other muxer properties in the future
_, err = elem.Connect("muxer-added", func(split, muxEle *gst.Element) { err := muxEle.SetProperty("presentation-time", false) if err != nil { panic("error setting presentation-time to false: " + err.Error()) } err = muxEle.SetProperty("interleave-bytes", InterleaveBytes) if err != nil { panic("error setting interleave-bytes" + err.Error()) } err = muxEle.SetProperty("interleave-time", InterleaveTime) if err != nil { panic("error setting interleave-time" + err.Error()) } err = muxEle.SetProperty("faststart", true) if err != nil { panic("error setting faststart" + err.Error()) } err = muxEle.SetProperty("movie-timescale", uint(60000)) if err != nil { panic("error setting movie-timescale" + err.Error()) } err = muxEle.SetProperty("trak-timescale", uint(60000)) if err != nil { panic("error setting trak-timescale" + err.Error()) } }) if err != nil { return nil, fmt.Errorf("failed to connect muxer-added handler: %w", err) }
// channel to make sure data is emitted in order var ch chan struct{}
_, err = elem.Connect("sink-added", func(split, sinkEle *gst.Element) { previousSegCh := ch mySegCh := make(chan struct{}, 1) ch = mySegCh buf := &bytes.Buffer{} err := sinkEle.SetProperty("sync", false) if err != nil { panic("error setting sync to false: " + err.Error()) } appsink := app.SinkFromElement(sinkEle) if appsink == nil { panic("appsink should not be nil") }
appsink.SetCallbacks(&app.SinkCallbacks{ NewSampleFunc: WriterNewSample(ctx, buf), EOSFunc: func(sink *app.Sink) { // ctx, span := otel.Tracer("signer").Start(ctx, "SegmentAndSignElem", trace.WithAttributes( // attribute.String("streamer", ms.Streamer()), // )) // defer span.End() now := time.Now().UnixMilli() bs := buf.Bytes()
if previousSegCh != nil { <-previousSegCh } resetTimer <- struct{}{} convergeAndSign := func() error { log.Debug(ctx, "signing segment", "size", len(bs)) err = cb(ctx, bs, now) if err != nil { return fmt.Errorf("error signing segment: %w", err) } return nil } err := func() error { convergeDone := make(chan error) go func() { convergeDone <- convergeAndSign() }() select { case <-ctx.Done(): return ctx.Err() case err := <-convergeDone: return err case <-time.After(time.Second * 3): go func() { err = cli.DataFileWrite([]string{"debug-recordings", streamer, fmt.Sprintf("converge-timeout-%d.mp4", now)}, bytes.NewReader(bs), true) if err != nil { log.Error(ctx, "error writing debug recording", "error", err) } }() return fmt.Errorf("timeout converging segment") } }() close(mySegCh) if err != nil { log.Error(ctx, "error in segmenter", "error", err) if FatalSegmentationErrors { sink.ErrorMessage(gst.DomainCore, gst.CoreErrorFailed, "error in segmenter", err.Error()) return } } }, }) }) if err != nil { return nil, fmt.Errorf("failed to connect sink-added handler: %w", err) }
return elem, nil}
func (mm *MediaManager) SegmentAndSignElem(ctx context.Context, ms MediaSigner) (*gst.Element, error) { tracer := otel.Tracer("signer") streamer := ms.Streamer()
// Stamp a fresh ingest-session epoch on the context that flows down to every // segment of this session (onSegment → ValidateMP4 → feedStreamTranscoder). A // new live session (RTMP/WHIP (re)connect) restarts the media timeline; the // per-DID continuous transcoder keys on this epoch and rebuilds rather than // feeding the restarted timeline into the previous session's encoder. ctx = withIngestSession(ctx, mm.nextIngestSession())
// muxl path: stream the fMP4 through the per-segment signer. Each GoP // arrives as a bare canonical .m4s, which ValidateMP4 verifies, archives // (as .m4s), and distributes. muxl-sign stamps the signing time into the // segment, so ValidateMP4 derives StartTime from the segment itself — no // per-GoP wall-clock is threaded through here. onSegment := func(ctx context.Context, segment []byte) error { // Top-level span for segment-delivery latency: signed bytes from the // streaming signer go in, a validated/persisted/distributed segment // comes out (ValidateMP4 = in-wasm c2pa verify + media parse + DB). ctx, span := tracer.Start(ctx, "SegmentAndSign", trace.WithAttributes( attribute.String("streamer", streamer), attribute.Int("segment_bytes", len(segment)), )) defer span.End() startTime := time.Now() defer func() { spmetrics.SegmentDeliveryDuration.WithLabelValues(streamer). Observe(float64(time.Since(startTime).Milliseconds())) }() if err := mm.ValidateMP4(ctx, bytes.NewReader(segment), true); err != nil { span.SetAttributes(attribute.String("error", "validate")) mm.cli.DumpDebugSegment(ctx, "just-signed-segment.m4s", bytes.NewReader(segment)) return fmt.Errorf("error validating signed segment: %w", err) } return nil } return MuxlSignSegmentElem(ctx, mm.cli, ms, onSegment)}
func SegmentFileUnsigned(ctx context.Context, cli *config.CLI, streamer string, input string, ch chan *SplitSegment) error { fd, err := os.OpenFile(input, os.O_RDONLY, 0644) log.Log(ctx, "reading file", "file", input) if err != nil { return fmt.Errorf("failed to read file: %w", err) } defer fd.Close() return SegmentUnsigned(ctx, cli, streamer, fd, false, ch)}
func SegmentUnsigned(ctx context.Context, cli *config.CLI, streamer string, input io.Reader, doH264Parse bool, ch chan *SplitSegment) error { ctx, cancel := context.WithCancel(ctx) defer cancel() pipelineSlice := []string{ "appsrc name=appsrc ! qtdemux name=demux", "demux. ! queue ! h264parse name=videoparse disable-passthrough=true config-interval=0", "demux. ! queue ! opusparse name=audioparse", } pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) if err != nil { return fmt.Errorf("error creating SegmentUnsigned pipeline: %w", err) }
srcele, err := pipeline.GetElementByName("appsrc") if err != nil { return err } src := app.SrcFromElement(srcele) src.SetCallbacks(&app.SourceCallbacks{ NeedDataFunc: ReaderNeedDataIncremental(ctx, input), }) videoParseEle, err := pipeline.GetElementByName("videoparse") if err != nil { return err }
segmenter, err := SegmentElem(ctx, cli, streamer, doH264Parse, func(ctx context.Context, buf []byte, now int64) error { ch <- &SplitSegment{ Filename: fmt.Sprintf("%d.mp4", now), Data: buf, } return nil }) if err != nil { return err }
err = pipeline.Add(segmenter) if err != nil { return err } err = videoParseEle.Link(segmenter) if err != nil { return err } audioparse, err := pipeline.GetElementByName("audioparse") if err != nil { return err } err = audioparse.Link(segmenter) if err != nil { return err }
busErr := make(chan error) go func() { err := HandleBusMessages(ctx, pipeline) cancel() busErr <- err }()
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 if err != nil { return err }
return nil}