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)) }