diff --git a/go.mod b/go.mod index a7a7e6998..0a55b6341 100644 --- a/go.mod +++ b/go.mod @@ -10,6 +10,11 @@ replace github.com/AxisCommunications/go-dpop => github.com/streamplace/go-dpop replace github.com/bluesky-social/indigo => github.com/streamplace/indigo v0.0.0-20260218231908-939cdaf0c507 +// LOCAL DEV: build against the working-tree muxl (dual-codec / track-id work). +// Before merge: commit + push the muxl changes and bump the pinned +// pseudo-version above, then drop this replace (the documented merge-blocker). +replace github.com/streamplace/muxl/go => ../muxl/go + tool github.com/bluesky-social/indigo/cmd/lexgen require ( diff --git a/go.sum b/go.sum index bb3b3be2d..d2eedc836 100644 --- a/go.sum +++ b/go.sum @@ -1375,8 +1375,6 @@ github.com/streamplace/go-dpop v0.0.0-20250510031900-c897158a8ad4 h1:L1fS4HJSaAy github.com/streamplace/go-dpop v0.0.0-20250510031900-c897158a8ad4/go.mod h1:bGUXY9Wd4mnd+XUrOYZr358J2f6z9QO/dLhL1SsiD+0= github.com/streamplace/indigo v0.0.0-20260218231908-939cdaf0c507 h1:e8M3qPLr37NxEjlr18TaAwGP+OVyherVjgUG5VVmgWI= github.com/streamplace/indigo v0.0.0-20260218231908-939cdaf0c507/go.mod h1:Pm2I1+iDXn/hLbF7XCg/DsZi6uDCiOo7hZGWprSM7k0= -github.com/streamplace/muxl/go v0.0.0-20260526201538-22bb7055201b h1:PyMsuvIibqAXi81+nNJPtqW6SedDHyhu22gGffwd/DI= -github.com/streamplace/muxl/go v0.0.0-20260526201538-22bb7055201b/go.mod h1:aCyYTW3o6c1Kush9UJ/Yv6EYMUbj8l8GTD7cHKcSxw8= github.com/streamplace/oatproxy v0.0.0-20260508220721-f8852e8dbf44 h1:b38ToXNQCvKqBlx0SeQYUOhx6uqqnrF5AL6B0Z3zzdk= github.com/streamplace/oatproxy v0.0.0-20260508220721-f8852e8dbf44/go.mod h1:j1+zdhe1IC0+PTE2rIXk2LqYn4Lz2x9SJaV3eHkQyfs= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= diff --git a/pkg/atproto/server_repo.go b/pkg/atproto/server_repo.go index 67747a130..e98479cc7 100644 --- a/pkg/atproto/server_repo.go +++ b/pkg/atproto/server_repo.go @@ -27,8 +27,11 @@ import ( "gorm.io/gorm" "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/crypto/spkey" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/statedb" + + gocrypto "crypto" ) var ServerRepo *atrepo.Repo @@ -40,6 +43,22 @@ var serverRepoLock sync.Mutex var serverCommitDB *gorm.DB var serverRepoSigner func(ctx context.Context, did string, sb []byte) ([]byte, error) +// serverRepoPriv is the node's secp256k1 server-repo private key — the key +// behind its did:web identity — captured by MakeServerRepo. Nil until then. +var serverRepoPriv *atcrypto.PrivateKeyK256 + +// ServerCryptoSigner returns a crypto.Signer for the node's own secp256k1 +// identity (the server-repo key behind its did:web). Used to S2PA-sign +// node-produced artifacts — e.g. a transcoded audio track minted at validate +// time, signed as a c2pa.transcoded derivative of the streamer's segment. +// Errors if the server repo hasn't been initialized yet. +func ServerCryptoSigner() (gocrypto.Signer, error) { + if serverRepoPriv == nil { + return nil, fmt.Errorf("server repo key not initialized") + } + return spkey.KeyToSigner(serverRepoPriv) +} + // serverCommitSubscribers is notified when new commit events are created. var serverCommitSubscribers []chan *ServerCommitEvent var serverCommitSubLock sync.Mutex @@ -136,6 +155,8 @@ func MakeServerRepo(ctx context.Context, cli *config.CLI, state *statedb.Statefu } } + serverRepoPriv = priv + pub, err := priv.PublicKey() if err != nil { return nil, fmt.Errorf("failed to get server repo public key: %w", err) diff --git a/pkg/livehls/livehls.go b/pkg/livehls/livehls.go index 06033d66c..55e818a2e 100644 --- a/pkg/livehls/livehls.go +++ b/pkg/livehls/livehls.go @@ -253,18 +253,35 @@ func (w *Writer) MasterPlaylist(trackURL func(trackID string) string) string { var b strings.Builder b.WriteString("#EXTM3U\n#EXT-X-VERSION:7\n#EXT-X-INDEPENDENT-SEGMENTS\n") + // Pick the primary audio rendition: prefer AAC (broadest HLS support — + // Safari has no Opus) so it's the default; any other codec (Opus) stays a + // selectable alternate. A segment may carry both after codec completion. + primaryAudio := "" + for _, tid := range w.order { + if t := w.tracks[tid]; t.Type == "audio" { + if primaryAudio == "" { + primaryAudio = tid + } + if strings.HasPrefix(t.Codec, "mp4a") { + primaryAudio = tid + break + } + } + } var audioCodec string - haveAudio := false + haveAudio := primaryAudio != "" for _, tid := range w.order { t := w.tracks[tid] - if t.Type == "audio" { - haveAudio = true - if audioCodec == "" { - audioCodec = t.Codec - } - fmt.Fprintf(&b, "#EXT-X-MEDIA:TYPE=AUDIO,GROUP-ID=%q,NAME=%q,DEFAULT=YES,AUTOSELECT=YES,CHANNELS=%q,URI=%q\n", - "audio", t.Codec, strconv.Itoa(int(maxU32(t.Channels, 2))), trackURL(tid)) + if t.Type != "audio" { + continue + } + def := "NO" + if tid == primaryAudio { + def = "YES" + audioCodec = t.Codec } + fmt.Fprintf(&b, "#EXT-X-MEDIA:TYPE=AUDIO,GROUP-ID=%q,NAME=%q,DEFAULT=%s,AUTOSELECT=YES,CHANNELS=%q,URI=%q\n", + "audio", t.Codec, def, strconv.Itoa(int(maxU32(t.Channels, 2))), trackURL(tid)) } for _, tid := range w.order { t := w.tracks[tid] diff --git a/pkg/media/media.go b/pkg/media/media.go index 43d10aea9..34ee1803a 100644 --- a/pkg/media/media.go +++ b/pkg/media/media.go @@ -54,6 +54,15 @@ type MediaManager struct { webrtcAPI *webrtc.API webrtcConfig webrtc.Configuration localDB localdb.LocalDB + + // Node S2PA transcode signer (cert + PKCS#8 key PEM), built once from the + // server-repo key. Used to sign transcode-completed audio tracks under the + // node's own did:web identity, signed in-wasm (the node key is software). + // See transcode.go. + nodeSignerOnce sync.Once + nodeCert []byte + nodeKeyPEM []byte + nodeSignerErr error } type NewSegmentNotification struct { diff --git a/pkg/media/media_data_parser.go b/pkg/media/media_data_parser.go index 9216c483e..7f1fa1203 100644 --- a/pkg/media/media_data_parser.go +++ b/pkg/media/media_data_parser.go @@ -27,14 +27,27 @@ func ParseSegmentMediaData(ctx context.Context, mp4bs []byte) (*localdb.SegmentM ctx, span := otel.Tracer("signer").Start(ctx, "ParseSegmentMediaData") defer span.End() ctx = log.WithLogValues(ctx, "GStreamerFunc", "ParseSegmentMediaData") - ctx, cancel := context.WithCancel(ctx) + // Watchdog: parsing a ~1s segment is sub-second. A stalled qtdemux only + // posts non-fatal warnings, so bound it — a hang surfaces as an error + // rather than wedging the validate/ingest pipeline (which is what a stray + // unresolved delayed link did). + ctx, cancel := context.WithTimeout(ctx, 30*time.Second) defer cancel() + // Codec-agnostic: audio rate/channels come from the qtdemux pad caps (read + // in onPadAdded) and durations from the demuxed buffers, so no per-codec + // parser is needed — the segment may carry AAC or Opus. We link only + // video_0 and audio_0. A dual-codec (completed) segment's extra audio track + // (audio_1) is left unlinked: qtdemux's flow combiner tolerates the + // not-linked pad, and the linked sinks still reach EOS. We must NOT add an + // idle sink for it (or delayed-link a pad that may not exist) — an unlinked + // sink never receives EOS, so the pipeline never posts EOS to the bus and + // the parse stalls until its watchdog fires. pipelineSlice := []string{ "appsrc name=appsrc ! qtdemux name=demux", fmt.Sprintf("demux.video_0 ! %s ! tee name=videotee", constants.Queue2Big), fmt.Sprintf("videotee. ! %s ! h2642json ! appsink sync=false name=jsonappsink", constants.Queue2Big), fmt.Sprintf("videotee. ! %s ! appsink sync=false name=videoappsink", constants.Queue2Big), - fmt.Sprintf("demux.audio_0 ! %s ! opusparse name=audioparse ! appsink sync=false name=audioappsink", constants.Queue2Big), + fmt.Sprintf("demux.audio_0 ! %s ! appsink sync=false name=audioappsink", constants.Queue2Big), } pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) @@ -89,8 +102,8 @@ func ParseSegmentMediaData(ctx context.Context, mp4bs []byte) (*localdb.SegmentM if info.GetEvent().Type() != gst.EventTypeEOS { return gst.PadProbeOK } - if padsAdded != 2 { - err := fmt.Errorf("expected 2 tracks in input, got %d", padsAdded) + if padsAdded < 2 { + err := fmt.Errorf("expected at least 2 tracks (video + audio), got %d", padsAdded) pipeline.Error(err.Error(), err) } padProbe = padProbeEmpty @@ -150,9 +163,11 @@ func ParseSegmentMediaData(ctx context.Context, mp4bs []byte) (*localdb.SegmentM } } - if name[:5] == "audio" { + // Primary audio track (audio_0) is statically linked to audioappsink; + // read its rate/channels for the segment metadata. Extra audio tracks + // (audio_1+, on a dual-codec completed segment) are left unlinked. + if name[:5] == "audio" && pad.GetName() == "audio_0" { audioMetadata = &localdb.SegmentMediadataAudio{} - // Get some common audio properties rateVal, _ := structure.GetValue("rate") channelsVal, _ := structure.GetValue("channels") diff --git a/pkg/media/packetize.go b/pkg/media/packetize.go index f576822e4..dd5906f21 100644 --- a/pkg/media/packetize.go +++ b/pkg/media/packetize.go @@ -14,6 +14,7 @@ import ( "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/constants" "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/muxl" ) // take in a segment and return a bunch of packets suitable for webrtc @@ -27,6 +28,22 @@ func Packetize(ctx context.Context, cli *config.CLI, seg *bus.Seg) (*bus.Packeti defer cancel() ctx = log.WithLogValues(ctx, "func", "Packetize", "uuid", uu.String()) + + // WebRTC playback needs Opus. From a dual-codec segment select video+Opus + // and present it as a flat MP4, so the demux below yields exactly one + // (Opus) audio pad for opusparse — no extra AAC pad to strand. + if len(seg.Muxl) > 0 { + opusM4s, err := filterSegmentToCodec(ctx, seg.Muxl, true) + if err != nil { + return nil, fmt.Errorf("select opus audio: %w", err) + } + var flat bytes.Buffer + if err := muxl.RunMuxlWrap(ctx, bytes.NewReader(opusM4s), "flat", &flat); err != nil { + return nil, fmt.Errorf("wrap opus segment: %w", err) + } + seg = &bus.Seg{Filepath: seg.Filepath, Data: flat.Bytes(), Muxl: opusM4s} + } + cli.DumpDebugSegment(ctx, fmt.Sprintf("packetize-input-%s.mp4", uu.String()), bytes.NewReader(seg.Data)) pipelineSlice := []string{ diff --git a/pkg/media/rtmp_ingest.go b/pkg/media/rtmp_ingest.go index a2d0994bc..e5a12ee4a 100644 --- a/pkg/media/rtmp_ingest.go +++ b/pkg/media/rtmp_ingest.go @@ -32,9 +32,13 @@ type RTMPSession struct { 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 ! fdkaacdec ! audioresample ! opusenc name=audioenc", + "demux.audio ! queue ! aacparse name=audioenc", "demux.video ! queue ! h264parse name=parse", } pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) diff --git a/pkg/media/rtmp_push.go b/pkg/media/rtmp_push.go index 167fa661c..cee48a800 100644 --- a/pkg/media/rtmp_push.go +++ b/pkg/media/rtmp_push.go @@ -33,7 +33,9 @@ func (mm *MediaManager) RTMPPush(ctx context.Context, user string, rendition str "appsrc name=muxlsrc ! qtdemux name=demux", "flvmux name=muxer ! rtmp2sink name=rtmp2sink", fmt.Sprintf("%s name=videoqueue ! h264parse ! muxer.video", constants.Queue2Big), - fmt.Sprintf("%s name=audioqueue ! opusparse ! opusdec ! audioresample ! fdkaacenc ! muxer.audio", constants.Queue2Big), + // Segments carry AAC (we feed only the AAC track below), so pass it + // straight to flvmux — no Opus→AAC transcode. + fmt.Sprintf("%s name=audioqueue ! aacparse ! muxer.audio", constants.Queue2Big), } pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) @@ -150,9 +152,18 @@ func (mm *MediaManager) RTMPPush(ctx context.Context, user string, rendition str log.Warn(ctx, "source segment has no MUXL bytes, skipping", "file", seg.Filepath) continue } + // RTMP wants AAC: select video + the AAC audio track from the + // dual-codec segment and feed only those, so flvmux gets AAC + // with no transcode. + aacSeg, err := filterSegmentToCodec(ctx, seg.Muxl, false) + if err != nil { + log.Error(ctx, "failed to select AAC audio", "error", err) + pw.CloseWithError(err) + return + } if first { var init bytes.Buffer - if err := muxl.RunMuxlWrapInit(ctx, bytes.NewReader(seg.Muxl), &init); err != nil { + if err := muxl.RunMuxlWrapInit(ctx, bytes.NewReader(aacSeg), &init); err != nil { pw.CloseWithError(fmt.Errorf("synthesize init segment: %w", err)) return } @@ -164,8 +175,8 @@ func (mm *MediaManager) RTMPPush(ctx context.Context, user string, rendition str } first = false } - log.Debug(ctx, "writing segment", "size", len(seg.Muxl)) - if _, err := pw.Write(seg.Muxl); err != nil { + log.Debug(ctx, "writing segment", "size", len(aacSeg)) + if _, err := pw.Write(aacSeg); err != nil { log.Error(ctx, "failed to write segment", "error", err) pw.CloseWithError(err) return @@ -193,7 +204,7 @@ func (mm *MediaManager) RTMPPush(ctx context.Context, user string, rendition str } // qtdemux exposes its track pads only after parsing the moov, so link them - // on pad-added: video → h264parse, audio → opusparse. + // on pad-added: video → h264parse, audio → aacparse. demux, err := pipeline.GetElementByName("demux") if err != nil { return fmt.Errorf("failed to get demux element from pipeline: %w", err) diff --git a/pkg/media/transcode.go b/pkg/media/transcode.go new file mode 100644 index 000000000..7517da287 --- /dev/null +++ b/pkg/media/transcode.go @@ -0,0 +1,474 @@ +package media + +import ( + "bytes" + "context" + "fmt" + "sort" + "strconv" + "strings" + "time" + + "github.com/go-gst/go-gst/gst" + "github.com/go-gst/go-gst/gst/app" + "stream.place/streamplace/pkg/atproto" + "stream.place/streamplace/pkg/crypto/signers" + "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/muxl" +) + +// Audio codec identification. muxl catalogs carry the MIME-ish codec string +// (`mp4a.40.2` for AAC, `opus` for Opus). We complete every segment to carry +// one of each so downstream pipelines can pick what they need. +func isAACCodec(codec string) bool { return strings.HasPrefix(codec, "mp4a") } +func isOpusCodec(codec string) bool { return strings.HasPrefix(codec, "opus") } + +// transcodeManifest is the C2PA manifest for a transcode-completed audio +// track: it opens the source segment as an ingredient and records a +// c2pa.transcoded action referencing it, both by TranscodeIngredientLabel. +// muxl-sign stamps the per-segment time into cawg.metadata as it signs. +func transcodeManifest(creator string) []byte { + return []byte(fmt.Sprintf(`{ + "title": "transcoded audio", + "assertions": [ + {"label": "c2pa.actions", "data": {"actions": [ + {"action": "c2pa.opened", "parameters": {"org.cai.ingredientIds": [%q]}}, + {"action": "c2pa.transcoded", "parameters": {"org.cai.ingredientIds": [%q]}} + ]}}, + {"label": "cawg.metadata", "data": { + "@context": {"dc": "http://purl.org/dc/elements/1.1/"}, + "dc:creator": %q + }} + ] + }`, muxl.TranscodeIngredientLabel, muxl.TranscodeIngredientLabel, creator)) +} + +// transcodeSigner returns the node's S2PA cert + PKCS#8 key PEM, built once +// from the server-repo key. The transcode-completed track is signed under the +// node's own did:web identity (it vouches for its transcode), declaring the +// streamer's source track as a parentOf ingredient. The node key is software, +// so it's signed in-wasm via KeyPEM — the same proven path MediaSignerLocal +// uses for ecdsa keys. +func (mm *MediaManager) transcodeSigner() (cert []byte, keyPEM []byte, err error) { + mm.nodeSignerOnce.Do(func() { + signer, e := atproto.ServerCryptoSigner() + if e != nil { + mm.nodeSignerErr = fmt.Errorf("node transcode signer: %w", e) + return + } + c, e := signers.GenerateES256KCert(signer) + if e != nil { + mm.nodeSignerErr = fmt.Errorf("node transcode cert: %w", e) + return + } + k, e := signers.MarshalES256KPrivateKeyPEM(signer) + if e != nil { + mm.nodeSignerErr = fmt.Errorf("node transcode key: %w", e) + return + } + mm.nodeCert = c + mm.nodeKeyPEM = k + }) + return mm.nodeCert, mm.nodeKeyPEM, mm.nodeSignerErr +} + +// completeAudioCodecs ensures a validated segment carries both AAC and Opus +// audio. seg is the bare canonical .m4s for one GoP (all tracks). It inspects +// the embedded catalog; if exactly one audio codec is present it transcodes +// the audio to the other, mints it at a free track id, signs it as a +// c2pa.transcoded derivative of the source audio track (under the node +// identity), and returns seg with the new signed track appended. It is a +// no-op (returns seg unchanged) when both codecs are already present, when +// there is no audio, when there is no video to anchor GoP boundaries, or when +// the audio codec is neither AAC nor Opus. +func (mm *MediaManager) completeAudioCodecs(ctx context.Context, seg []byte) ([]byte, error) { + // 1. Unwrap the source segment: catalog (codecs/track-ids) + per-track + // bytes (the signed source-audio track we'll declare as the parent). + events, err := unwrapMuxlEvents(ctx, seg) + if err != nil { + return nil, fmt.Errorf("unwrap segment: %w", err) + } + cat, tracks := catalogAndTracks(events) + if cat == nil || cat.Audio == nil { + return seg, nil // no audio to complete + } + + var ( + haveAAC, haveOpus bool + srcAudioTID uint32 + srcAudioCodec string + maxTID uint32 + ) + if cat.Video != nil { + for _, v := range cat.Video.Renditions { + maxTID = maxU32(maxTID, v.TrackID()) + } + } + for _, a := range cat.Audio.Renditions { + maxTID = maxU32(maxTID, a.TrackID()) + switch { + case isAACCodec(a.Codec): + haveAAC = true + case isOpusCodec(a.Codec): + haveOpus = true + } + srcAudioTID = a.TrackID() + srcAudioCodec = a.Codec + } + + // 2. Decide whether (and to what) we need to transcode. + if haveAAC && haveOpus { + return seg, nil // already complete + } + var target string + switch { + case haveAAC && !haveOpus: + target = "opus" + case haveOpus && !haveAAC: + target = "aac" + default: + log.Warn(ctx, "segment audio codec not AAC/Opus, skipping codec completion", "codec", srcAudioCodec) + return seg, nil + } + if cat.Video == nil { + // No video reference: GoP boundaries can't be anchored cleanly for the + // transcoded track. Audio-only live is rare; skip for now. + log.Warn(ctx, "audio-only segment, skipping codec completion") + return seg, nil + } + + sourceAudio := tracks[strconv.FormatUint(uint64(srcAudioTID), 10)] + if len(sourceAudio) == 0 { + return nil, fmt.Errorf("source audio track %d missing from segment", srcAudioTID) + } + + cert, keyPEM, err := mm.transcodeSigner() + if err != nil { + // No node signing identity (e.g. server repo not initialized) — leave + // the segment single-codec rather than failing ingest. + log.Warn(ctx, "node transcode signer unavailable, skipping codec completion", "error", err) + return seg, nil + } + + // 3. Wrap the whole segment to a flat MP4 (gstreamer needs video present so + // muxl's keyframe-anchored canonicalization yields one aligned segment). + var flat bytes.Buffer + if err := muxl.RunMuxlWrap(ctx, bytes.NewReader(seg), "flat", &flat); err != nil { + return nil, fmt.Errorf("wrap segment for transcode: %w", err) + } + + // 4. Transcode the audio (video passes through, for GoP alignment only). + transFmp4, err := transcodeAudioSegment(ctx, flat.Bytes(), target) + if err != nil { + return nil, fmt.Errorf("transcode audio to %s: %w", target, err) + } + + // 5. Find the transcoded audio track's id, then canonicalize remapping it + // to a free id so it can join the source tracks without colliding. + transEvents, err := segmentMuxlEvents(ctx, transFmp4) + if err != nil { + return nil, fmt.Errorf("segment transcoded output: %w", err) + } + transCat, _ := catalogAndTracks(transEvents) + if transCat == nil || transCat.Audio == nil { + return nil, fmt.Errorf("transcoded output has no audio track") + } + var transAudioTID uint32 + for _, a := range transCat.Audio.Renditions { + transAudioTID = a.TrackID() + } + freeTID := maxTID + 1 + canon, err := muxl.RunMuxlCanonicalize(ctx, transFmp4, map[uint32]uint32{transAudioTID: freeTID}) + if err != nil { + return nil, fmt.Errorf("canonicalize transcoded output: %w", err) + } + + // 6. Extract the (now-free-id) transcoded audio track as a bare unsigned + // canonical .m4s — the SignTranscode output. + canonEvents, err := unwrapMuxlEvents(ctx, canon) + if err != nil { + return nil, fmt.Errorf("unwrap canonicalized output: %w", err) + } + _, canonTracks := catalogAndTracks(canonEvents) + output := canonTracks[strconv.FormatUint(uint64(freeTID), 10)] + if len(output) == 0 { + return nil, fmt.Errorf("transcoded audio track %d missing after canonicalize", freeTID) + } + + // 7. Sign the transcoded track as a c2pa.transcoded derivative of the + // source audio track, under the node identity. + signed, err := muxl.RunMuxlSignTranscode(ctx, muxl.TranscodeInput{ + Output: output, + Source: sourceAudio, + CertPEM: cert, + KeyPEM: keyPEM, + Manifest: transcodeManifest(mm.cli.BroadcasterDID()), + }) + if err != nil { + return nil, fmt.Errorf("sign transcoded track: %w", err) + } + + // 8. Append the signed track. Its id is the largest, so concatenating it + // last preserves the canonical track-id-ascending segment order. + completed := make([]byte, 0, len(seg)+len(signed)) + completed = append(completed, seg...) + completed = append(completed, signed...) + log.Log(ctx, "completed segment audio codecs", + "source_codec", srcAudioCodec, "added_codec", target, + "source_track", srcAudioTID, "added_track", freeTID, + "in_bytes", len(seg), "out_bytes", len(completed)) + return completed, nil +} + +// transcodeAudioSegment transcodes the audio of a flat MP4 segment to the +// target codec ("opus" or "aac"), passing video through untouched, and returns +// a fragmented MP4 (video + transcoded audio). Video rides along only so muxl +// can anchor the canonical segment on the video keyframe; the caller keeps the +// original signed video and uses only the transcoded audio track. +func transcodeAudioSegment(ctx context.Context, flat []byte, target string) ([]byte, error) { + // Watchdog: a per-segment audio transcode is a sub-second real-time job. A + // stalled gstreamer pipeline only posts non-fatal warnings (not bus errors), + // so bound it explicitly — a hang surfaces as an error instead of wedging + // the whole validate/ingest path. + ctx, cancel := context.WithTimeout(ctx, 30*time.Second) + defer cancel() + + var audioChain string + switch target { + case "opus": // source is AAC + audioChain = "queue name=aq ! aacparse ! fdkaacdec ! audioconvert ! audioresample ! opusenc name=aenc" + case "aac": // source is Opus + audioChain = "queue name=aq ! opusparse ! opusdec ! audioconvert ! audioresample ! fdkaacenc name=aenc" + default: + return nil, fmt.Errorf("unsupported transcode target %q", target) + } + + pipeline, err := gst.NewPipelineFromString(strings.Join([]string{ + "appsrc name=src ! qtdemux name=demux", + "queue name=vq ! h264parse name=vparse", + audioChain, + }, "\n")) + if err != nil { + return nil, fmt.Errorf("create transcode pipeline: %w", err) + } + defer func() { + if e := pipeline.SetState(gst.StateNull); e != nil { + log.Error(ctx, "transcode: set null", "error", e) + } + }() + + // Fragmented muxer + appsink, built programmatically so we can set the + // enum fragment-mode and pre-request pads in video-then-audio order. + mux, err := gst.NewElementWithProperties("mp4mux", map[string]any{ + "name": "mux", + "fragment-mode": 0, + "fragment-duration": 1, + }) + if err != nil { + return nil, fmt.Errorf("create mp4mux: %w", err) + } + sink, err := gst.NewElementWithProperties("appsink", map[string]any{"name": "sink", "sync": false}) + if err != nil { + return nil, fmt.Errorf("create appsink: %w", err) + } + if err := pipeline.AddMany(mux, sink); err != nil { + return nil, fmt.Errorf("add mux/sink: %w", err) + } + if err := mux.Link(sink); err != nil { + return nil, fmt.Errorf("link mux→sink: %w", err) + } + + // Pre-request pads in order so the muxer assigns video=track 1, audio=2. + videoMuxPad := mux.GetRequestPad("video_%u") + audioMuxPad := mux.GetRequestPad("audio_%u") + if videoMuxPad == nil || audioMuxPad == nil { + return nil, fmt.Errorf("failed to request mp4mux pads") + } + + vparse, err := pipeline.GetElementByName("vparse") + if err != nil { + return nil, err + } + aenc, err := pipeline.GetElementByName("aenc") + if err != nil { + return nil, err + } + if r := vparse.GetStaticPad("src").Link(videoMuxPad); r != gst.PadLinkOK { + return nil, fmt.Errorf("link video chain → mux: %v", r) + } + if r := aenc.GetStaticPad("src").Link(audioMuxPad); r != gst.PadLinkOK { + return nil, fmt.Errorf("link audio chain → mux: %v", r) + } + + vq, err := pipeline.GetElementByName("vq") + if err != nil { + return nil, err + } + aq, err := pipeline.GetElementByName("aq") + if err != nil { + return nil, err + } + demux, err := pipeline.GetElementByName("demux") + if err != nil { + return nil, err + } + if _, err := demux.Connect("pad-added", func(self *gst.Element, pad *gst.Pad) { + name := pad.GetName() + var dst *gst.Pad + switch { + case strings.HasPrefix(name, "video"): + dst = vq.GetStaticPad("sink") + case strings.HasPrefix(name, "audio"): + dst = aq.GetStaticPad("sink") + default: + return + } + if r := pad.Link(dst); r != gst.PadLinkOK { + log.Error(ctx, "transcode: link demux pad", "name", name, "result", r) + } + }); err != nil { + return nil, fmt.Errorf("connect demux pad-added: %w", err) + } + + srcEle, err := pipeline.GetElementByName("src") + if err != nil { + return nil, err + } + app.SrcFromElement(srcEle).SetCallbacks(&app.SourceCallbacks{ + NeedDataFunc: ReaderNeedDataIncremental(ctx, bytes.NewReader(flat)), + }) + + var out bytes.Buffer + app.SinkFromElement(sink).SetCallbacks(&app.SinkCallbacks{ + NewSampleFunc: WriterNewSample(ctx, &out), + }) + + errCh := make(chan error, 1) + go func() { errCh <- HandleBusMessages(ctx, pipeline) }() + if err := pipeline.SetState(gst.StatePlaying); err != nil { + return nil, fmt.Errorf("transcode: set playing: %w", err) + } + if err := <-errCh; err != nil { + return nil, fmt.Errorf("transcode pipeline: %w", err) + } + if out.Len() == 0 { + return nil, fmt.Errorf("transcode produced no output") + } + return out.Bytes(), nil +} + +// --- muxl event helpers --- + +func unwrapMuxlEvents(ctx context.Context, seg []byte) ([]*muxl.MuxlEvent, error) { + return collectMuxlEvents(func(ch chan *muxl.MuxlEvent) error { + return muxl.RunMuxlUnwrapEvents(ctx, bytes.NewReader(seg), ch) + }) +} + +func segmentMuxlEvents(ctx context.Context, fmp4 []byte) ([]*muxl.MuxlEvent, error) { + return collectMuxlEvents(func(ch chan *muxl.MuxlEvent) error { + return muxl.RunMuxlSegmenterEvents(ctx, bytes.NewReader(fmp4), ch) + }) +} + +func collectMuxlEvents(run func(chan *muxl.MuxlEvent) error) ([]*muxl.MuxlEvent, error) { + ch := make(chan *muxl.MuxlEvent, 16) + errCh := make(chan error, 1) + go func() { + err := run(ch) + close(ch) + errCh <- err + }() + var out []*muxl.MuxlEvent + for ev := range ch { + out = append(out, ev) + } + return out, <-errCh +} + +// catalogAndTracks pulls the catalog (from the init event) and the per-track +// bytes (from the first segment/signed-segment event) out of a muxl event +// stream. +func catalogAndTracks(events []*muxl.MuxlEvent) (*muxl.MuxlCatalog, map[string][]byte) { + var cat *muxl.MuxlCatalog + var tracks map[string][]byte + for _, ev := range events { + switch ev.Type { + case "init": + if ev.Catalog != nil { + cat = ev.Catalog + } + case "segment", "signed-segment": + if tracks == nil && len(ev.Tracks) > 0 { + tracks = ev.Tracks + } + } + } + return cat, tracks +} + +func maxU32(a, b uint32) uint32 { + if a > b { + return a + } + return b +} + +// filterSegmentToCodec returns a bare canonical .m4s containing every video +// track plus the single audio track matching the requested codec (Opus when +// wantOpus, else AAC) — how an output consumer "asks for the audio it needs" +// from a dual-codec segment. Track bytes are carried verbatim (signatures +// intact) in ascending track-id order. If no audio matches the requested +// codec, any one audio track is kept (degraded but playable); with no audio +// info at all, the segment is returned unchanged. +func filterSegmentToCodec(ctx context.Context, seg []byte, wantOpus bool) ([]byte, error) { + events, err := unwrapMuxlEvents(ctx, seg) + if err != nil { + return nil, fmt.Errorf("unwrap segment for codec filter: %w", err) + } + cat, tracks := catalogAndTracks(events) + if cat == nil { + return seg, nil + } + + keep := map[uint32]bool{} + if cat.Video != nil { + for _, v := range cat.Video.Renditions { + keep[v.TrackID()] = true + } + } + if cat.Audio != nil && len(cat.Audio.Renditions) > 0 { + var chosen uint32 + found := false + for _, a := range cat.Audio.Renditions { + if (wantOpus && isOpusCodec(a.Codec)) || (!wantOpus && isAACCodec(a.Codec)) { + chosen, found = a.TrackID(), true + break + } + } + if !found { // no exact codec match — keep some audio rather than none + for _, a := range cat.Audio.Renditions { + chosen, found = a.TrackID(), true + break + } + } + if found { + keep[chosen] = true + } + } + + ids := make([]uint32, 0, len(keep)) + for id := range keep { + ids = append(ids, id) + } + sort.Slice(ids, func(i, j int) bool { return ids[i] < ids[j] }) + + var out []byte + for _, id := range ids { + out = append(out, tracks[strconv.FormatUint(uint64(id), 10)]...) + } + if len(out) == 0 { + return seg, nil + } + return out, nil +} diff --git a/pkg/media/transcode_test.go b/pkg/media/transcode_test.go new file mode 100644 index 000000000..77c082cff --- /dev/null +++ b/pkg/media/transcode_test.go @@ -0,0 +1,110 @@ +package media + +import ( + "bytes" + "context" + "os" + "sort" + "testing" + + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/crypto/signers" + "stream.place/streamplace/pkg/muxl" +) + +// audioCodecsOf returns the sorted distinct audio codecs in a bare .m4s segment. +func audioCodecsOf(t *testing.T, ctx context.Context, seg []byte) []string { + t.Helper() + events, err := unwrapMuxlEvents(ctx, seg) + require.NoError(t, err) + cat, _ := catalogAndTracks(events) + require.NotNil(t, cat, "segment has a catalog") + var out []string + if cat.Audio != nil { + for _, a := range cat.Audio.Renditions { + out = append(out, a.Codec) + } + } + sort.Strings(out) + return out +} + +// signedBareSegment signs the given fragmented fixture and returns the first +// GoP's bare canonical .m4s (all tracks) — the live ingest shape. +func signedBareSegment(t *testing.T, ctx context.Context, ms *MediaSignerLocal, fragPath string) []byte { + t.Helper() + frag, err := os.ReadFile(fragPath) + require.NoError(t, err) + eventCh := make(chan *muxl.MuxlEvent, 16) + errCh := make(chan error, 1) + go func() { + err := ms.SignSegmentStream(ctx, bytes.NewReader(frag), eventCh) + close(eventCh) + errCh <- err + }() + var seg []byte + for ev := range eventCh { + if ev.Type == "signed-segment" && seg == nil { + seg = concatTracksSorted(ev.Tracks) + } + } + require.NoError(t, <-errCh) + require.NotEmpty(t, seg, "expected at least one signed GoP") + return seg +} + +// TestCompleteAudioCodecs drives the validate-time codec completion directly: +// sign a single-audio (Opus) segment, then run completeAudioCodecs and confirm +// it transcodes + transcode-signs the missing AAC track so the segment carries +// both codecs and still verifies. This is the live-path machinery (gstreamer +// transcode + muxl canonicalize/remap + SignTranscode), exercised without an +// RTMP source so a hang/error surfaces locally with a stack trace. +func TestCompleteAudioCodecs(t *testing.T) { + ctx := context.Background() + ms := newBareSegmentSigner(t) + + seg := signedBareSegment(t, ctx, ms, getFixture("h264-opus-frag.mp4")) + + srcCodecs := audioCodecsOf(t, ctx, seg) + require.Len(t, srcCodecs, 1, "source should be single-audio, got %v", srcCodecs) + + // Build a MediaManager with a node transcode signer (reuse the test key). + mm := &MediaManager{cli: &config.CLI{BroadcasterHost: "test.example.com"}} + keyPEM, err := signers.MarshalES256KPrivateKeyPEM(ms.Signer) + require.NoError(t, err) + mm.nodeCert = ms.Cert + mm.nodeKeyPEM = keyPEM + mm.nodeSignerOnce.Do(func() {}) // mark built so transcodeSigner returns the preset cert/key + + completed, err := mm.completeAudioCodecs(ctx, seg) + require.NoError(t, err) + require.Greater(t, len(completed), len(seg), "completion should append a signed track") + + codecs := audioCodecsOf(t, ctx, completed) + require.Len(t, codecs, 2, "expected both AAC and Opus after completion, got %v", codecs) + hasAAC, hasOpus := false, false + for _, c := range codecs { + if isAACCodec(c) { + hasAAC = true + } + if isOpusCodec(c) { + hasOpus = true + } + } + require.True(t, hasAAC, "expected an AAC track, got %v", codecs) + require.True(t, hasOpus, "expected the original Opus track, got %v", codecs) + + // The completed segment must still verify end-to-end (every track). + out, err := muxl.RunMuxlVerify(ctx, bytes.NewReader(completed)) + require.NoError(t, err) + require.Contains(t, out, "track_id", "verify output should describe each track") + + // And it must media-parse without stalling: the extra (audio_1) track is + // left unlinked in the parse pipeline, which must not block EOS. With the + // -timeout on the test, a regression here shows up as a 30s watchdog stall. + res, err := ValidateMP4Media(ctx, completed) + require.NoError(t, err) + require.NotEmpty(t, res.MediaData.Video, "3-track segment parses video") + require.NotEmpty(t, res.MediaData.Audio, "3-track segment parses audio") +} diff --git a/pkg/media/validate.go b/pkg/media/validate.go index 223b86df5..ced389ee0 100644 --- a/pkg/media/validate.go +++ b/pkg/media/validate.go @@ -121,6 +121,17 @@ func (mm *MediaManager) ValidateMP4(ctx context.Context, input io.Reader, local } } + // Complete the segment's audio so it carries both AAC and Opus (no-op when + // already dual-codec, audio-only, or the codec is neither). The added track + // is signed under the node's own identity as a c2pa.transcoded derivative + // of the source audio track. Everything archived/distributed below uses the + // completed bytes; replicas re-run this and find it already complete. + completed, err := mm.completeAudioCodecs(ctx, buf) + if err != nil { + return fmt.Errorf("complete audio codecs: %w", err) + } + buf = completed + _, fileSpan := tracer.Start(ctx, "ValidateMP4.SegmentArchiveWrite", trace.WithAttributes( attribute.Int("bytes", len(buf)), )) diff --git a/pkg/muxl/muxl.go b/pkg/muxl/muxl.go index 544d76f4a..b98939a60 100644 --- a/pkg/muxl/muxl.go +++ b/pkg/muxl/muxl.go @@ -26,13 +26,23 @@ type ( MuxlAudioConfig = upstream.AudioConfig MuxlContainer = upstream.Container SignerInput = upstream.SignerInput + TranscodeInput = upstream.TranscodeInput ) +// TranscodeIngredientLabel is the C2PA ingredient label SignTranscode assigns +// to the source segment; a TranscodeInput.Manifest references the source by +// listing it in an action's "org.cai.ingredientIds" (see the upstream doc). +const TranscodeIngredientLabel = upstream.TranscodeIngredientLabel + // SignerToCallback adapts a crypto.Signer into the host-sign callback that -// SignerInput.Sign expects (SHA-256 digest, ECDSA DER → fixed-width r‖s). See -// the upstream doc. +// SignerInput.Sign / TranscodeInput.Sign expect (SHA-256 digest, ECDSA DER → +// fixed-width r‖s). See the upstream doc. var SignerToCallback = upstream.SignerToCallback +// RawSignerToCallback adapts a raw-secp256k1 digest signer (e.g. an Ethereum +// keystore's SignHash) into the host-sign callback. See the upstream doc. +var RawSignerToCallback = upstream.RawSignerToCallback + // --- engine singleton ------------------------------------------------------- var ( @@ -131,6 +141,35 @@ func RunMuxlSignSegment(ctx context.Context, input io.Reader, in SignerInput, in return eng.SignSegment(ctx, input, in, initCh, segCh, eventCh) } +// RunMuxlCanonicalize converts a flat or fragmented MP4 (e.g. a transcoder's +// output) into a canonical MUXL fMP4. trackRemap (may be nil) reassigns track +// IDs in the output — used to mint a transcoded rendition at a free id so it +// can join the source's tracks in one multi-track segment without colliding. +func RunMuxlCanonicalize(ctx context.Context, mp4 []byte, trackRemap map[uint32]uint32) ([]byte, error) { + eng, err := getEngine() + if err != nil { + return nil, err + } + var opts []upstream.CanonicalizeOption + if len(trackRemap) > 0 { + opts = append(opts, upstream.WithTrackRemap(trackRemap)) + } + return eng.Canonicalize(ctx, mp4, opts...) +} + +// RunMuxlSignTranscode signs in.Output (an unsigned canonical MUXL segment — +// the transcoded result) as a standalone asset declaring in.Source (the +// canonical segment it was transcoded from) as a c2pa.transcoded parentOf +// ingredient. Returns the signed segment. Exactly one of in.KeyPEM or in.Sign +// must be set. +func RunMuxlSignTranscode(ctx context.Context, in TranscodeInput) ([]byte, error) { + eng, err := getEngine() + if err != nil { + return nil, err + } + return eng.SignTranscode(ctx, in) +} + // --- push-style Concatenator ------------------------------------------------ // Concatenator is a push wrapper over the Engine: Write whole fMP4 archives,