From f2c72925145d98aa5c2af9453c7bea36aac60145 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Sat, 6 Jun 2026 16:15:10 -0700 Subject: [PATCH] media: keep debug recording on the isolated ingest paths MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Debug recording (the DebugRecording per-stream setting → a debug-recordings// dump) only worked on the in-process paths and the fd-4-pipe fallback, where main is in the data path and can tee. The zero-downtime detached MKV path and the WHIP worker put main OUT of the data path, so recording silently stopped there. Move the recording into the worker, uniformly across the isolated paths. main still owns the DECISION (shouldRecord needs the DB), evaluated once in buildWorkerConfig and passed as cfg.Record + cfg.DataDir in the handshake; the worker, which owns the data path, carries it out: - RunMKVIngestWorker tees its ingest media to dumpToFile when cfg.Record. - ServeWHIPIngestWorkerSocket passes cfg.Record to NewRecordingPeerConnection. - The now-redundant main-side tee is removed from MKVIngestIsolated, so there's one recording path (no double-record) and main stays out of the data path on the fd-4 path too. The in-process MKVIngest / NewPeerConnection recordings are unchanged. Bonus: because the worker owns the recording, a recording now survives a main restart along with the rest of the detached session. cfg.DataDir is only set when recording (rtcrec/dumpToFile panic on an empty DataDir, but they're only reached when enabled). Test: TestRunMKVIngestWorkerRecords runs the worker with Record + a temp DataDir and asserts debug-recordings//*.rtmp.mkv is written verbatim from the ingested media. Co-Authored-By: Claude Opus 4.8 --- pkg/media/ingest_supervisor.go | 23 ++++++++--------- pkg/media/ingest_worker.go | 34 +++++++++++++++++++++--- pkg/media/ingest_worker_test.go | 46 +++++++++++++++++++++++++++++++++ pkg/media/whip_worker.go | 8 +++--- 4 files changed, 91 insertions(+), 20 deletions(-) diff --git a/pkg/media/ingest_supervisor.go b/pkg/media/ingest_supervisor.go index 8ba69db18..47fa4bc49 100644 --- a/pkg/media/ingest_supervisor.go +++ b/pkg/media/ingest_supervisor.go @@ -47,18 +47,8 @@ func (mm *MediaManager) MKVIngestIsolated(ctx context.Context, input io.Reader, return fmt.Errorf("marshal worker config: %w", err) } - // Optional recording stays in main: tee the raw input before it reaches the - // worker, so the worker needs no data-dir access. - if shouldRecord, rerr := mm.shouldRecord(ctx, ms.Streamer()); rerr == nil && shouldRecord { - log.Log(ctx, "recording RTMP stream to file", "streamer", ms.Streamer()) - pr, pw := io.Pipe() - input = io.TeeReader(input, pw) - go func() { - if derr := mm.dumpToFile(ctx, pr, ms.Streamer(), ".rtmp.mkv"); derr != nil { - log.Error(ctx, "error dumping to file", "error", derr) - } - }() - } + // Debug recording now happens inside the worker (cfg.Record), uniformly across + // the isolated paths — main stays out of the data path here too. exe, err := os.Executable() if err != nil { @@ -243,6 +233,15 @@ func (mm *MediaManager) buildWorkerConfig(ctx context.Context, ms MediaSigner) ( Manifest: manifest, BroadcasterHost: mm.cli.BroadcasterHost, } + // Debug recording: main owns the per-stream setting (it needs the DB); the + // worker carries out the recording (it owns the data path). A lookup failure + // is non-fatal — just don't record. + if rec, rerr := mm.shouldRecord(ctx, ms.Streamer()); rerr != nil { + log.Warn(ctx, "could not read debug-recording setting; not recording", "streamer", ms.Streamer(), "error", rerr) + } else if rec { + cfg.Record = true + cfg.DataDir = mm.cli.DataDir + } // 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. diff --git a/pkg/media/ingest_worker.go b/pkg/media/ingest_worker.go index 5941b691a..cc95cf731 100644 --- a/pkg/media/ingest_worker.go +++ b/pkg/media/ingest_worker.go @@ -62,6 +62,16 @@ type IngestWorkerConfig struct { Prebuf []byte `json:"prebuf,omitempty"` Chunked bool `json:"chunked,omitempty"` + // Record, when true, makes the worker write a debug recording of this session + // (the MKV/RTMP push body, or the WHIP session) under + // DataDir/debug-recordings//. main evaluates the per-stream DebugRecording + // setting (which needs the DB) and the worker carries it out — so debug + // recording keeps working on the isolated paths without main being in the data + // path, and a recording even survives a main restart. DataDir is set (only when + // Record) to the node data dir the worker writes recordings under. + Record bool `json:"record,omitempty"` + DataDir string `json:"data_dir,omitempty"` + // Transport selects the worker's ingest source: "" / "mkv" reads MKV media // (stdin or InputFD); "whip" makes the worker own the WebRTC PeerConnection, // built from OfferSDP — no media fd to pass. @@ -158,16 +168,32 @@ 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}} + // Minimal manager: the broadcaster identity the transcode completion + // (finishTranscodedSegment) stamps into the node-signed AAC track, plus the + // data dir for an optional debug recording. + mm := &MediaManager{cli: &config.CLI{BroadcasterHost: cfg.BroadcasterHost, DataDir: cfg.DataDir}} onSegment, flush := mm.workerSegmentSink(ctx, cfg, frames) + // Debug recording: tee the ingest media to a file before it reaches gst. main + // decided this (cfg.Record) and handed us DataDir; recording here keeps main + // out of the data path and lets the recording survive a main restart. + media := stdin + if cfg.Record { + log.Log(ctx, "recording ingest media to file", "streamer", cfg.StreamerDID) + pr, pw := io.Pipe() + media = io.TeeReader(stdin, pw) + go func() { + if derr := mm.dumpToFile(ctx, pr, cfg.StreamerDID, ".rtmp.mkv"); derr != nil { + log.Error(ctx, "ingest worker: dump recording to file", "error", derr) + } + }() + } + signerElem, done, err := muxlSignSegmentElem(ctx, mm.cli, workerSignStream(cfg), onSegment) if err != nil { return fmt.Errorf("build signer element: %w", err) } - pipeline, err := buildMKVIngestPipeline(ctx, stdin, signerElem) + pipeline, err := buildMKVIngestPipeline(ctx, media, signerElem) if err != nil { return fmt.Errorf("build pipeline: %w", err) } diff --git a/pkg/media/ingest_worker_test.go b/pkg/media/ingest_worker_test.go index c1d732999..ee8ce64f3 100644 --- a/pkg/media/ingest_worker_test.go +++ b/pkg/media/ingest_worker_test.go @@ -6,6 +6,8 @@ import ( "errors" "fmt" "io" + "os" + "path/filepath" "strings" "testing" "time" @@ -141,3 +143,47 @@ func TestRunMKVIngestWorkerProducesValidSignedFrames(t *testing.T) { require.GreaterOrEqual(t, segs, 1, "worker emitted at least one signed dual-codec segment") t.Logf("worker emitted %d valid dual-codec segments", segs) } + +// TestRunMKVIngestWorkerRecords proves debug recording works INSIDE the worker: +// with cfg.Record set and a DataDir handed over, the worker tees its ingest +// media to debug-recordings//.rtmp.mkv. This is what keeps debug +// recording working on the isolated paths where main is out of the data path +// (it can't tee the bytes itself), so main decides and the worker records. +func TestRunMKVIngestWorkerRecords(t *testing.T) { + ctx := context.Background() + ms := newBareSegmentSigner(t) + + keyPEM, err := signers.MarshalES256KPrivateKeyPEM(ms.Signer) + require.NoError(t, err) + manifest, err := ms.buildManifest(ctx, time.Now().UnixMilli()) + require.NoError(t, err) + + dataDir := t.TempDir() + cfg := IngestWorkerConfig{ + StreamerDID: ms.Streamer(), + KeyPEM: keyPEM, + CertPEM: ms.Cert, + Manifest: manifest, + BroadcasterHost: "test.example.com", + Record: true, + DataDir: dataDir, + } + + mkv := makeH264AACMKV(t, ctx, getFixture("5sec.mp4")) + + // Frames are irrelevant here (recording is on the input side); discard them. + require.NoError(t, RunMKVIngestWorker(ctx, cfg, bytes.NewReader(mkv), ingestframe.NewWriter(io.Discard))) + + // The recording lands at debug-recordings//.rtmp.mkv and + // must contain exactly the media the worker ingested. The dump goroutine + // flushes asynchronously, so allow it a moment to finish the last write. + glob := filepath.Join(dataDir, "debug-recordings", "*", "*.rtmp.mkv") + require.Eventually(t, func() bool { + matches, _ := filepath.Glob(glob) + if len(matches) != 1 { + return false + } + got, rerr := os.ReadFile(matches[0]) + return rerr == nil && bytes.Equal(got, mkv) + }, 10*time.Second, 25*time.Millisecond, "worker records the ingest media verbatim") +} diff --git a/pkg/media/whip_worker.go b/pkg/media/whip_worker.go index ffb88592f..ff53f0af2 100644 --- a/pkg/media/whip_worker.go +++ b/pkg/media/whip_worker.go @@ -58,11 +58,11 @@ func ServeWHIPIngestWorkerSocket(ctx context.Context, cfg IngestWorkerConfig) er return runErr } - mm := &MediaManager{cli: &config.CLI{BroadcasterHost: cfg.BroadcasterHost}} + mm := &MediaManager{cli: &config.CLI{BroadcasterHost: cfg.BroadcasterHost, DataDir: cfg.DataDir}} // The worker owns the PeerConnection (its own UDP sockets), built with the - // same codec/interceptor setup as the in-process server. No recording here — - // the worker has no model-backed settings. + // same codec/interceptor setup as the in-process server. Debug recording is + // decided by main (cfg.Record) and written by the worker under cfg.DataDir. api, webrtcConfig, err := newWebRTCAPI() if err != nil { return finish(fmt.Errorf("webrtc api: %w", err)) @@ -71,7 +71,7 @@ func ServeWHIPIngestWorkerSocket(ctx context.Context, cfg IngestWorkerConfig) er if err != nil { return finish(fmt.Errorf("peer connection: %w", err)) } - pc, err := rtcrec.NewRecordingPeerConnection(ctx, *mm.cli, cfg.StreamerDID, pionpc, false) + pc, err := rtcrec.NewRecordingPeerConnection(ctx, *mm.cli, cfg.StreamerDID, pionpc, cfg.Record) if err != nil { return finish(fmt.Errorf("peer connection wrapper: %w", err)) } -- 2.51.2