diff --git a/pkg/media/ingest_subprocess_test.go b/pkg/media/ingest_subprocess_test.go index c852325c..27b791f0 100644 --- a/pkg/media/ingest_subprocess_test.go +++ b/pkg/media/ingest_subprocess_test.go @@ -124,3 +124,30 @@ func TestIngestWorkerSubprocess(t *testing.T) { require.True(t, sawEnd, "worker subprocess emitted a clean End frame") t.Logf("worker subprocess emitted %d valid signed segments + clean End", segs) } + +// TestMKVIngestIsolatedWedgeContained is the isolation guarantee: sample-stream.mkv +// carries four audio tracks, so the single-audio ingest pipeline leaves three +// matroskademux pads unlinked and wedges with no EOS — exactly the kind of native +// wedge that would hang (or, with a runaway buffer, OOM-kill) an in-process +// ingest and take the node with it. Run in a worker, it must be contained: the +// watchdog kills the worker and MKVIngestIsolated returns an error, bounded in +// time, with THIS process — the node — still running to assert it. +func TestMKVIngestIsolatedWedgeContained(t *testing.T) { + old := ingestWorkerWatchdog + ingestWorkerWatchdog = 6 * time.Second + defer func() { ingestWorkerWatchdog = old }() + + mm, _ := getStaticTestMediaManager(t) + ms := newBareSegmentSigner(t) + + wedge, err := os.ReadFile(getFixture("sample-stream.mkv")) + require.NoError(t, err) + + start := time.Now() + err = mm.MKVIngestIsolated(context.Background(), bytes.NewReader(wedge), ms) + elapsed := time.Since(start) + + require.Error(t, err, "a wedged worker must surface as an error, not a hang") + require.Less(t, elapsed, 25*time.Second, "watchdog bounded the wedge") + t.Logf("wedged worker contained in %s: %v", elapsed.Round(time.Second), err) +} diff --git a/pkg/media/ingest_supervisor.go b/pkg/media/ingest_supervisor.go index ce558d7e..02e74711 100644 --- a/pkg/media/ingest_supervisor.go +++ b/pkg/media/ingest_supervisor.go @@ -19,6 +19,13 @@ import ( "stream.place/streamplace/pkg/log" ) +// ingestWorkerWatchdog bounds how long an isolated worker may go without +// producing a segment frame before it's presumed wedged (a native pipeline gst +// can't drain — e.g. a pathological stream) and killed. A healthy stream emits a +// segment every GoP (~1–2s), so this is generous. Var (not const) so tests can +// shorten it. +var ingestWorkerWatchdog = 30 * time.Second + // MKVIngestIsolated is the process-isolated counterpart to MKVIngest. Instead of // running the demux + sign pipeline in this process — where a native gst fault, // OOM, or deadlock would take the whole node down — it spawns a dedicated @@ -124,8 +131,20 @@ func (mm *MediaManager) MKVIngestIsolated(ctx context.Context, input io.Reader, go func() { defer logsWG.Done(); streamWorkerLogs(ctx, stdout, ms.Streamer()) }() go func() { defer logsWG.Done(); streamWorkerLogs(ctx, stderr, ms.Streamer()) }() + // Watchdog: a worker that stops producing frames is presumed wedged and + // killed (cancel → CommandContext SIGKILLs it), so a single bad stream can't + // hang its session forever. Reset on every frame. + watchdog := time.AfterFunc(ingestWorkerWatchdog, func() { + log.Warn(ctx, "ingest worker watchdog fired (no frames); killing worker", + "streamer", ms.Streamer(), "timeout", ingestWorkerWatchdog) + cancel() + }) + defer watchdog.Stop() + // Read signed-segment frames and feed each into the normal chokepoint. - sawEnd, readErr := mm.consumeWorkerFrames(ctx, framesR, ms.Streamer()) + sawEnd, readErr := mm.consumeWorkerFrames(ctx, framesR, ms.Streamer(), func() { + watchdog.Reset(ingestWorkerWatchdog) + }) logsWG.Wait() werr := cmd.Wait() @@ -146,7 +165,7 @@ func (mm *MediaManager) MKVIngestIsolated(ctx context.Context, input io.Reader, // over each. It returns whether a clean End frame was seen and the terminal read // error: nil on a clean close (End then EOF), or io.ErrUnexpectedEOF / a desync // error when the worker died mid-frame. -func (mm *MediaManager) consumeWorkerFrames(ctx context.Context, stdout io.Reader, streamer string) (sawEnd bool, _ error) { +func (mm *MediaManager) consumeWorkerFrames(ctx context.Context, stdout io.Reader, streamer string, onProgress func()) (sawEnd bool, _ error) { fr := ingestframe.NewReader(stdout) for { typ, payload, err := fr.ReadFrame() @@ -156,6 +175,9 @@ func (mm *MediaManager) consumeWorkerFrames(ctx context.Context, stdout io.Reade } return sawEnd, err } + if onProgress != nil { + onProgress() + } switch typ { case ingestframe.Segment: if verr := mm.ValidateMP4(ctx, bytes.NewReader(payload), true); verr != nil {