diff --git a/pkg/media/ingest_supervisor.go b/pkg/media/ingest_supervisor.go index 02e74711..52d365eb 100644 --- a/pkg/media/ingest_supervisor.go +++ b/pkg/media/ingest_supervisor.go @@ -222,10 +222,21 @@ func (mm *MediaManager) buildWorkerConfig(ctx context.Context, ms MediaSigner) ( if err != nil { return IngestWorkerConfig{}, fmt.Errorf("build manifest: %w", err) } - return IngestWorkerConfig{ - StreamerDID: ms.Streamer(), - KeyPEM: keyPEM, - CertPEM: local.Cert, - Manifest: manifest, - }, nil + cfg := IngestWorkerConfig{ + StreamerDID: ms.Streamer(), + KeyPEM: keyPEM, + CertPEM: local.Cert, + Manifest: manifest, + BroadcasterHost: mm.cli.BroadcasterHost, + } + // Node transcode signer lets the worker complete to dual-codec itself. If it's + // unavailable, the worker emits single-codec (the node doesn't re-transcode the + // worker's output) — an acceptable, logged fallback rather than a hard failure. + if nodeCert, nodeKeyPEM, serr := mm.transcodeSigner(); serr == nil { + cfg.NodeCertPEM = nodeCert + cfg.NodeKeyPEM = nodeKeyPEM + } else { + log.Warn(ctx, "node transcode signer unavailable; isolated worker will emit single-codec", "error", serr) + } + return cfg, nil } diff --git a/pkg/media/ingest_worker.go b/pkg/media/ingest_worker.go index 22d64f5d..84128fdc 100644 --- a/pkg/media/ingest_worker.go +++ b/pkg/media/ingest_worker.go @@ -33,6 +33,14 @@ type IngestWorkerConfig struct { // pre-live → live transition) don't yet cross the boundary; that needs a // control channel and is tracked as future work. Manifest []byte `json:"manifest"` + + // Node transcode signer + broadcaster identity. When set, the worker completes + // each single-codec source segment to dual-codec (Opus+AAC) itself — the + // transcode runs in this isolated process too — signing the added track under + // the node identity. Empty → the worker emits single-codec source segments. + NodeCertPEM []byte `json:"node_cert_pem,omitempty"` + NodeKeyPEM []byte `json:"node_key_pem,omitempty"` + BroadcasterHost string `json:"broadcaster_host,omitempty"` } // RunMKVIngestWorker is the body of the `ingest-worker` subcommand. It reads an @@ -50,6 +58,10 @@ func RunMKVIngestWorker(ctx context.Context, cfg IngestWorkerConfig, stdin io.Re ctx, cancel := context.WithCancel(ctx) defer cancel() + // Minimal manager: just the broadcaster identity the transcode completion + // (finishTranscodedSegment) stamps into the node-signed AAC track. + mm := &MediaManager{cli: &config.CLI{BroadcasterHost: cfg.BroadcasterHost}} + // The worker signs everything itself: forward the streamer key PEM + cert + // prebuilt manifest straight to muxl-sign. No MediaSigner / model / DB needed. signStream := func(ctx context.Context, input io.Reader, eventCh chan *muxl.MuxlEvent) error { @@ -62,11 +74,34 @@ func RunMKVIngestWorker(ctx context.Context, cfg IngestWorkerConfig, stdin io.Re }, nil, nil, eventCh) } + // With a node transcode key, the worker completes each single-codec source + // segment to dual-codec itself: feed the signed source segment into a + // per-stream transcoder running in THIS process; its completion callback + // frames the finished dual-codec segment. The transcoder runs on a + // non-cancellable context so draining the signer (cancel, below) can't kill it + // before its ~1-GoP tail is flushed by Close. One process == one session, so + // the per-DID transcoder-reuse hazard simply can't arise here. + var transcoder *streamTranscoder onSegment := func(_ context.Context, segment []byte) error { - return frames.Segment(segment) + if len(cfg.NodeKeyPEM) == 0 { + return frames.Segment(segment) // no node signer → single-codec + } + if transcoder == nil { + target, need := mm.audioCompletionTarget(ctx, segment) + if !need { + return frames.Segment(segment) // already dual-codec / no audio track + } + transcoder = mm.newStreamTranscoder(context.WithoutCancel(ctx), target, cfg.NodeCertPEM, cfg.NodeKeyPEM, + func(_ any, completed []byte) { + if ferr := frames.Segment(completed); ferr != nil { + log.Error(ctx, "ingest worker: frame completed segment", "error", ferr) + } + }) + } + return transcoder.Feed(segment, nil) } - signerElem, done, err := muxlSignSegmentElem(ctx, &config.CLI{}, signStream, onSegment) + signerElem, done, err := muxlSignSegmentElem(ctx, mm.cli, signStream, onSegment) if err != nil { return fmt.Errorf("build signer element: %w", err) } @@ -90,10 +125,17 @@ func RunMKVIngestWorker(ctx context.Context, cfg IngestWorkerConfig, stdin io.Re }() // Wait for the pipeline to finish (EOS or error), then drain the signer: - // cancelling unblocks the signer's input pipe so it flushes the final GoP, - // and <-done guarantees every segment frame is written before we return. + // cancelling unblocks the signer's input pipe so it flushes the final GoP, and + // <-done guarantees every source segment has been fed. Then flush the + // transcoder's tail so the last dual-codec completions are framed before we + // return (the caller's End can't race ahead of them). pipeErr := <-busErr cancel() <-done + if transcoder != nil { + if cerr := transcoder.Close(); cerr != nil { + log.Error(ctx, "ingest worker: transcoder close", "error", cerr) + } + } return pipeErr } diff --git a/pkg/media/ingest_worker_test.go b/pkg/media/ingest_worker_test.go index 9eb8d6c1..efb699c9 100644 --- a/pkg/media/ingest_worker_test.go +++ b/pkg/media/ingest_worker_test.go @@ -66,11 +66,16 @@ func TestRunMKVIngestWorkerProducesValidSignedFrames(t *testing.T) { manifest, err := ms.buildManifest(ctx, time.Now().UnixMilli()) require.NoError(t, err) + // Provide the node transcode signer (the same test key serves both roles + // here, as in the transcoder tests) so the worker completes to dual-codec. cfg := IngestWorkerConfig{ - StreamerDID: ms.Streamer(), - KeyPEM: keyPEM, - CertPEM: ms.Cert, - Manifest: manifest, + StreamerDID: ms.Streamer(), + KeyPEM: keyPEM, + CertPEM: ms.Cert, + Manifest: manifest, + NodeCertPEM: ms.Cert, + NodeKeyPEM: keyPEM, + BroadcasterHost: "test.example.com", } mkv := makeH264AACMKV(t, ctx, getFixture("5sec.mp4")) @@ -95,8 +100,23 @@ func TestRunMKVIngestWorkerProducesValidSignedFrames(t *testing.T) { out, err := muxl.RunMuxlVerify(ctx, bytes.NewReader(payload)) require.NoError(t, err, "segment %d verify", segs) require.NotContains(t, out, `"validation_state":"Invalid"`, "segment %d must validate", segs) + + // With a node key the worker completes to dual-codec: every segment must + // carry both the source Opus and a worker-transcoded AAC track. + codecs := audioCodecsOf(t, ctx, payload) + hasOpus, hasAAC := false, false + for _, c := range codecs { + if isOpusCodec(c) { + hasOpus = true + } + if isAACCodec(c) { + hasAAC = true + } + } + require.True(t, hasOpus, "segment %d keeps source Opus (got %v)", segs, codecs) + require.True(t, hasAAC, "segment %d gains worker-transcoded AAC (got %v)", segs, codecs) segs++ } - require.GreaterOrEqual(t, segs, 1, "worker emitted at least one signed segment") - t.Logf("worker emitted %d valid signed segments", segs) + require.GreaterOrEqual(t, segs, 1, "worker emitted at least one signed dual-codec segment") + t.Logf("worker emitted %d valid dual-codec segments", segs) }