From 159d5ea3cab3a02c6a0ec0f5bfffcd474f33935c Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Sat, 6 Jun 2026 16:38:16 -0700 Subject: [PATCH] media: worker-side self-watchdog for the detached/WHIP ingest paths MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Wedge containment had a hole: the supervisor watchdog (MKVIngestIsolated) only covers the fd-4 fallback. On the default detached path and WHIP, the worker is detached (not tied to main's context), so main can't kill a stuck one — and a wedged native pipeline (e.g. the 4-audio MKV that leaves matroskademux pads unlinked, never emits EOS, produces no frames) would run forever as an orphan. That's the exact failure isolation is supposed to contain. Add a worker-side watchdog: a FrameWriter wrapper kicks a timer on every frame the worker emits; going ingestWorkerWatchdog with no output means the pipeline is wedged, so it cancels the worker's context. muxlSignSegmentElem's ctx-done hook closes the signer pipe, so the worker actually returns and the process exits — the fault stays contained to the subprocess. Wired into both RunMKVIngestWorker and ServeWHIPIngestWorkerSocket (which kicks once after the SDP answer, so the first-segment clock starts post-negotiation; a pre-answer hang is already bounded by main's whipAnswerTimeout). The kick is mutex-guarded since a worker emits frames from more than one goroutine. This is additive on the fd-4 path (the supervisor watchdog still fires first there); it's the ONLY wedge containment on the detached/WHIP paths. Test: TestRunMKVIngestWorkerSelfWatchdog feeds the wedging MKV straight to RunMKVIngestWorker (no supervisor) and asserts it self-terminates rather than hanging. Co-Authored-By: Claude Opus 4.8 --- pkg/media/ingest_worker.go | 7 ++++ pkg/media/ingest_worker_test.go | 45 +++++++++++++++++++++ pkg/media/whip_worker.go | 10 ++++- pkg/media/worker_watchdog.go | 69 +++++++++++++++++++++++++++++++++ 4 files changed, 130 insertions(+), 1 deletion(-) create mode 100644 pkg/media/worker_watchdog.go diff --git a/pkg/media/ingest_worker.go b/pkg/media/ingest_worker.go index cc95cf731..09d47263b 100644 --- a/pkg/media/ingest_worker.go +++ b/pkg/media/ingest_worker.go @@ -168,6 +168,13 @@ func RunMKVIngestWorker(ctx context.Context, cfg IngestWorkerConfig, stdin io.Re ctx, cancel := context.WithCancel(ctx) defer cancel() + // Self-watchdog: if the pipeline wedges and stops emitting frames, tear it + // down so the worker exits and the fault stays contained (the only wedge + // containment on the detached path — main can't kill a detached worker). + wd := newWorkerWatchdog(ctx, ingestWorkerWatchdog, cancel, cfg.StreamerDID) + defer wd.stop() + frames = wd.wrap(frames) + // 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. diff --git a/pkg/media/ingest_worker_test.go b/pkg/media/ingest_worker_test.go index ee8ce64f3..618d170a3 100644 --- a/pkg/media/ingest_worker_test.go +++ b/pkg/media/ingest_worker_test.go @@ -187,3 +187,48 @@ func TestRunMKVIngestWorkerRecords(t *testing.T) { return rerr == nil && bytes.Equal(got, mkv) }, 10*time.Second, 25*time.Millisecond, "worker records the ingest media verbatim") } + +// TestRunMKVIngestWorkerSelfWatchdog proves the worker's OWN watchdog contains a +// wedge. This is the only wedge containment on the detached/WHIP paths, where +// main can't kill a detached worker — so the worker has to notice it's stuck and +// exit itself. The 4-audio sample-stream.mkv leaves matroskademux pads unlinked, +// so it wedges with no EOS and emits no frames; the watchdog must tear the +// pipeline down and return rather than hang forever. (The fd-4 path's +// supervisor-side watchdog is covered separately by +// TestMKVIngestIsolatedWedgeContained.) +func TestRunMKVIngestWorkerSelfWatchdog(t *testing.T) { + old := ingestWorkerWatchdog + ingestWorkerWatchdog = 3 * time.Second + defer func() { ingestWorkerWatchdog = old }() + + 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) + cfg := IngestWorkerConfig{ + StreamerDID: ms.Streamer(), + KeyPEM: keyPEM, + CertPEM: ms.Cert, + Manifest: manifest, + BroadcasterHost: "test.example.com", + } + + wedge, err := os.ReadFile(getFixture("sample-stream.mkv")) + require.NoError(t, err) + + start := time.Now() + done := make(chan error, 1) + go func() { + done <- RunMKVIngestWorker(ctx, cfg, bytes.NewReader(wedge), ingestframe.NewWriter(io.Discard)) + }() + select { + case <-done: + elapsed := time.Since(start) + require.Less(t, elapsed, 25*time.Second, "watchdog bounded the wedge") + t.Logf("worker self-terminated on wedge in %s", elapsed.Round(time.Second)) + case <-time.After(30 * time.Second): + t.Fatal("worker-side watchdog did not contain the wedge (RunMKVIngestWorker hung)") + } +} diff --git a/pkg/media/whip_worker.go b/pkg/media/whip_worker.go index ff53f0af2..09bb221e2 100644 --- a/pkg/media/whip_worker.go +++ b/pkg/media/whip_worker.go @@ -76,7 +76,14 @@ func ServeWHIPIngestWorkerSocket(ctx context.Context, cfg IngestWorkerConfig) er return finish(fmt.Errorf("peer connection wrapper: %w", err)) } - onSegment, flush := mm.workerSegmentSink(ctx, cfg, srv) + // Self-watchdog: a connected-but-silent publisher (or a wedged pipeline) that + // stops producing frames tears the worker down so the fault stays contained. + // The clock to the first segment starts at the answer (kick below), after ICE + // negotiation; a pre-answer hang is bounded by main's whipAnswerTimeout. + wd := newWorkerWatchdog(ctx, ingestWorkerWatchdog, cancel, cfg.StreamerDID) + defer wd.stop() + + onSegment, flush := mm.workerSegmentSink(ctx, cfg, wd.wrap(srv)) signerElem, signerDone, err := muxlSignSegmentElem(ctx, mm.cli, workerSignStream(cfg), onSegment) if err != nil { return finish(fmt.Errorf("build signer element: %w", err)) @@ -93,6 +100,7 @@ func ServeWHIPIngestWorkerSocket(ctx context.Context, cfg IngestWorkerConfig) er if aerr := srv.Answer(answer.SDP); aerr != nil { log.Error(ctx, "whip worker: frame answer", "error", aerr) } + wd.kick() // negotiation done + answer sent; start the first-segment clock here // Streaming runs until the peer disconnects / errors; webRTCIngestPipeline // cancels ctx then, which drains the signer. Wait for that, flush the diff --git a/pkg/media/worker_watchdog.go b/pkg/media/worker_watchdog.go new file mode 100644 index 000000000..1abf748e9 --- /dev/null +++ b/pkg/media/worker_watchdog.go @@ -0,0 +1,69 @@ +package media + +import ( + "context" + "sync" + "time" + + "stream.place/streamplace/pkg/log" +) + +// workerWatchdog self-terminates a wedged ingest worker. A healthy stream emits +// an output frame every GoP (~1–2s); going ingestWorkerWatchdog with no frame +// means the native pipeline is wedged (gst can't drain — e.g. a pathological +// stream that leaves demux pads unlinked) and won't recover. The watchdog then +// fires onWedge — the worker's context cancel — tearing the pipeline down so the +// process exits and the fault stays contained to this subprocess. +// +// This is the worker-side counterpart to MKVIngestIsolated's main-side watchdog, +// and the ONLY wedge containment on the detached and WHIP paths: those workers +// are detached (not tied to main's context), so main can't kill a stuck one — +// the worker has to notice and exit itself. +type workerWatchdog struct { + mu sync.Mutex + timer *time.Timer + to time.Duration +} + +func newWorkerWatchdog(ctx context.Context, timeout time.Duration, onWedge func(), streamer string) *workerWatchdog { + return &workerWatchdog{ + to: timeout, + timer: time.AfterFunc(timeout, func() { + log.Warn(ctx, "ingest worker watchdog fired (no output frames); self-terminating", + "streamer", streamer, "timeout", timeout) + onWedge() + }), + } +} + +// kick resets the no-output timer. Safe for concurrent use — a worker emits +// frames from more than one goroutine (the signer event loop and the +// transcoder's completion callback). +func (w *workerWatchdog) kick() { + w.mu.Lock() + defer w.mu.Unlock() + w.timer.Reset(w.to) +} + +// stop halts the watchdog; call it once the pipeline has ended so it can't fire +// during the post-stream drain. +func (w *workerWatchdog) stop() { + w.mu.Lock() + defer w.mu.Unlock() + w.timer.Stop() +} + +// wrap returns a FrameWriter that kicks the watchdog on every frame it writes, +// so any output the worker produces counts as liveness. +func (w *workerWatchdog) wrap(inner FrameWriter) FrameWriter { + return watchdogFrameWriter{FrameWriter: inner, wd: w} +} + +type watchdogFrameWriter struct { + FrameWriter + wd *workerWatchdog +} + +func (w watchdogFrameWriter) Segment(seg []byte) error { w.wd.kick(); return w.FrameWriter.Segment(seg) } +func (w watchdogFrameWriter) End() error { w.wd.kick(); return w.FrameWriter.End() } +func (w watchdogFrameWriter) Error(msg string) error { w.wd.kick(); return w.FrameWriter.Error(msg) } -- 2.51.2